* 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA.
*/
-#define _GNU_SOURCE
#define _LGPL_SOURCE
#include <assert.h>
#include <lttng/ust-ctl.h>
#include <common/relayd/relayd.h>
#include <common/compat/fcntl.h>
#include <common/compat/endian.h>
-#include <common/consumer-metadata-cache.h>
-#include <common/consumer-stream.h>
-#include <common/consumer-timer.h>
+#include <common/consumer/consumer-metadata-cache.h>
+#include <common/consumer/consumer-stream.h>
+#include <common/consumer/consumer-timer.h>
#include <common/utils.h>
#include <common/index/index.h>
#include "ust-consumer.h"
+#define INT_MAX_STR_LEN 12 /* includes \0 */
+
extern struct lttng_consumer_global_data consumer_data;
extern int consumer_poll_timeout;
extern volatile int consumer_quit;
*/
if (channel->uchan) {
lttng_ustconsumer_del_channel(channel);
+ lttng_ustconsumer_free_channel(channel);
}
free(channel);
}
uint64_t tracefile_size, uint64_t tracefile_count,
uint64_t session_id_per_pid, unsigned int monitor,
unsigned int live_timer_interval,
- const char *shm_path)
+ const char *root_shm_path, const char *shm_path)
{
assert(pathname);
assert(name);
return consumer_allocate_channel(key, session_id, pathname, name, uid,
gid, relayd_id, output, tracefile_size,
tracefile_count, session_id_per_pid, monitor,
- live_timer_interval, shm_path);
+ live_timer_interval, root_shm_path, shm_path);
}
/*
return ret;
}
+static
+int get_stream_shm_path(char *stream_shm_path, const char *shm_path, int cpu)
+{
+ char cpu_nr[INT_MAX_STR_LEN]; /* int max len */
+ int ret;
+
+ strncpy(stream_shm_path, shm_path, PATH_MAX);
+ stream_shm_path[PATH_MAX - 1] = '\0';
+ ret = snprintf(cpu_nr, INT_MAX_STR_LEN, "%i", cpu);
+ if (ret < 0) {
+ PERROR("snprintf");
+ goto end;
+ }
+ strncat(stream_shm_path, cpu_nr,
+ PATH_MAX - strlen(stream_shm_path) - 1);
+ ret = 0;
+end:
+ return ret;
+}
+
/*
* Create streams for the given channel using liblttng-ust-ctl.
*
/* Keep stream reference when creating metadata. */
if (channel->type == CONSUMER_CHANNEL_TYPE_METADATA) {
channel->metadata_stream = stream;
- stream->ust_metadata_poll_pipe[0] = ust_metadata_pipe[0];
- stream->ust_metadata_poll_pipe[1] = ust_metadata_pipe[1];
+ if (channel->monitor) {
+ /* Set metadata poll pipe if we created one */
+ memcpy(stream->ust_metadata_poll_pipe,
+ ust_metadata_pipe,
+ sizeof(ust_metadata_pipe));
+ }
}
}
return ret;
}
+/*
+ * create_posix_shm is never called concurrently within a process.
+ */
+static
+int create_posix_shm(void)
+{
+ char tmp_name[NAME_MAX];
+ int shmfd, ret;
+
+ ret = snprintf(tmp_name, NAME_MAX, "/ust-shm-consumer-%d", getpid());
+ if (ret < 0) {
+ PERROR("snprintf");
+ return -1;
+ }
+ /*
+ * Allocate shm, and immediately unlink its shm oject, keeping
+ * only the file descriptor as a reference to the object.
+ * We specifically do _not_ use the / at the beginning of the
+ * pathname so that some OS implementations can keep it local to
+ * the process (POSIX leaves this implementation-defined).
+ */
+ shmfd = shm_open(tmp_name, O_CREAT | O_EXCL | O_RDWR, 0700);
+ if (shmfd < 0) {
+ PERROR("shm_open");
+ goto error_shm_open;
+ }
+ ret = shm_unlink(tmp_name);
+ if (ret < 0 && errno != ENOENT) {
+ PERROR("shm_unlink");
+ goto error_shm_release;
+ }
+ return shmfd;
+
+error_shm_release:
+ ret = close(shmfd);
+ if (ret) {
+ PERROR("close");
+ }
+error_shm_open:
+ return -1;
+}
+
+static int open_ust_stream_fd(struct lttng_consumer_channel *channel,
+ struct ustctl_consumer_channel_attr *attr,
+ int cpu)
+{
+ char shm_path[PATH_MAX];
+ int ret;
+
+ if (!channel->shm_path[0]) {
+ return create_posix_shm();
+ }
+ ret = get_stream_shm_path(shm_path, channel->shm_path, cpu);
+ if (ret) {
+ goto error_shm_path;
+ }
+ return run_as_open(shm_path,
+ O_RDWR | O_CREAT | O_EXCL, S_IRUSR | S_IWUSR,
+ channel->uid, channel->gid);
+
+error_shm_path:
+ return -1;
+}
+
/*
* Create an UST channel with the given attributes and send it to the session
* daemon using the ust ctl API.
*
* Return 0 on success or else a negative value.
*/
-static int create_ust_channel(struct ustctl_consumer_channel_attr *attr,
- struct ustctl_consumer_channel **chanp)
+static int create_ust_channel(struct lttng_consumer_channel *channel,
+ struct ustctl_consumer_channel_attr *attr,
+ struct ustctl_consumer_channel **ust_chanp)
{
- int ret;
- struct ustctl_consumer_channel *channel;
+ int ret, nr_stream_fds, i, j;
+ int *stream_fds;
+ struct ustctl_consumer_channel *ust_channel;
+ assert(channel);
assert(attr);
- assert(chanp);
+ assert(ust_chanp);
DBG3("Creating channel to ustctl with attr: [overwrite: %d, "
"subbuf_size: %" PRIu64 ", num_subbuf: %" PRIu64 ", "
attr->num_subbuf, attr->switch_timer_interval,
attr->read_timer_interval, attr->output, attr->type);
- channel = ustctl_create_channel(attr);
- if (!channel) {
+ if (channel->type == CONSUMER_CHANNEL_TYPE_METADATA)
+ nr_stream_fds = 1;
+ else
+ nr_stream_fds = ustctl_get_nr_stream_per_channel();
+ stream_fds = zmalloc(nr_stream_fds * sizeof(*stream_fds));
+ if (!stream_fds) {
+ ret = -1;
+ goto error_alloc;
+ }
+ for (i = 0; i < nr_stream_fds; i++) {
+ stream_fds[i] = open_ust_stream_fd(channel, attr, i);
+ if (stream_fds[i] < 0) {
+ ret = -1;
+ goto error_open;
+ }
+ }
+ ust_channel = ustctl_create_channel(attr, stream_fds, nr_stream_fds);
+ if (!ust_channel) {
ret = -1;
goto error_create;
}
-
- *chanp = channel;
+ channel->nr_stream_fds = nr_stream_fds;
+ channel->stream_fds = stream_fds;
+ *ust_chanp = ust_channel;
return 0;
error_create:
+error_open:
+ for (j = i - 1; j >= 0; j--) {
+ int closeret;
+
+ closeret = close(stream_fds[j]);
+ if (closeret) {
+ PERROR("close");
+ }
+ if (channel->shm_path[0]) {
+ char shm_path[PATH_MAX];
+
+ closeret = get_stream_shm_path(shm_path,
+ channel->shm_path, j);
+ if (closeret) {
+ ERR("Cannot get stream shm path");
+ }
+ closeret = run_as_unlink(shm_path,
+ channel->uid, channel->gid);
+ if (closeret) {
+ PERROR("unlink %s", shm_path);
+ }
+ }
+ }
+ /* Try to rmdir all directories under shm_path root. */
+ if (channel->root_shm_path[0]) {
+ (void) run_as_recursive_rmdir(channel->root_shm_path,
+ channel->uid, channel->gid);
+ }
+ free(stream_fds);
+error_alloc:
return ret;
}
channel->nb_init_stream_left = 0;
/* The reply msg status is handled in the following call. */
- ret = create_ust_channel(attr, &channel->uchan);
+ ret = create_ust_channel(channel, attr, &channel->uchan);
if (ret < 0) {
goto end;
}
}
/*
- * Receive the metadata updates from the sessiond.
+ * Receive the metadata updates from the sessiond. Supports receiving
+ * overlapping metadata, but is needs to always belong to a contiguous
+ * range starting from 0.
+ * Be careful about the locks held when calling this function: it needs
+ * the metadata cache flush to concurrently progress in order to
+ * complete.
*/
int lttng_ustconsumer_recv_metadata(int sock, uint64_t key, uint64_t offset,
uint64_t len, struct lttng_consumer_channel *channel,
msg.u.ask_channel.session_id_per_pid,
msg.u.ask_channel.monitor,
msg.u.ask_channel.live_timer_interval,
+ msg.u.ask_channel.root_shm_path,
msg.u.ask_channel.shm_path);
if (!channel) {
goto end_channel_error;
attr.read_timer_interval = msg.u.ask_channel.read_timer_interval;
attr.chan_id = msg.u.ask_channel.chan_id;
memcpy(attr.uuid, msg.u.ask_channel.uuid, sizeof(attr.uuid));
- strncpy(attr.shm_path, channel->shm_path,
- sizeof(attr.shm_path));
- attr.shm_path[sizeof(attr.shm_path) - 1] = '\0';
/* Match channel buffer type to the UST abi. */
switch (msg.u.ask_channel.output) {
health_code_update();
+ if (!len) {
+ /*
+ * There is nothing to receive. We have simply
+ * checked whether the channel can be found.
+ */
+ ret_code = LTTCOMM_CONSUMERD_SUCCESS;
+ goto end_msg_sessiond;
+ }
+
/* Tell session daemon we are ready to receive the metadata. */
ret = consumer_send_status_msg(sock, LTTCOMM_CONSUMERD_SUCCESS);
if (ret < 0) {
health_code_update();
break;
}
+ case LTTNG_CONSUMER_DISCARDED_EVENTS:
+ {
+ uint64_t ret;
+ struct lttng_ht_iter iter;
+ struct lttng_ht *ht;
+ struct lttng_consumer_stream *stream;
+ uint64_t id = msg.u.discarded_events.session_id;
+ uint64_t key = msg.u.discarded_events.channel_key;
+
+ DBG("UST consumer discarded events command for session id %"
+ PRIu64, id);
+ rcu_read_lock();
+ pthread_mutex_lock(&consumer_data.lock);
+
+ ht = consumer_data.stream_list_ht;
+
+ /*
+ * We only need a reference to the channel, but they are not
+ * directly indexed, so we just use the first matching stream
+ * to extract the information we need, we default to 0 if not
+ * found (no events are dropped if the channel is not yet in
+ * use).
+ */
+ ret = 0;
+ cds_lfht_for_each_entry_duplicate(ht->ht,
+ ht->hash_fct(&id, lttng_ht_seed),
+ ht->match_fct, &id,
+ &iter.iter, stream, node_session_id.node) {
+ if (stream->chan->key == key) {
+ ret = stream->chan->discarded_events;
+ break;
+ }
+ }
+ pthread_mutex_unlock(&consumer_data.lock);
+ rcu_read_unlock();
+
+ DBG("UST consumer discarded events command for session id %"
+ PRIu64 ", channel key %" PRIu64, id, key);
+
+ health_code_update();
+
+ /* Send back returned value to session daemon */
+ ret = lttcomm_send_unix_sock(sock, &ret, sizeof(ret));
+ if (ret < 0) {
+ PERROR("send discarded events");
+ goto error_fatal;
+ }
+
+ break;
+ }
+ case LTTNG_CONSUMER_LOST_PACKETS:
+ {
+ uint64_t ret;
+ struct lttng_ht_iter iter;
+ struct lttng_ht *ht;
+ struct lttng_consumer_stream *stream;
+ uint64_t id = msg.u.lost_packets.session_id;
+ uint64_t key = msg.u.lost_packets.channel_key;
+
+ DBG("UST consumer lost packets command for session id %"
+ PRIu64, id);
+ rcu_read_lock();
+ pthread_mutex_lock(&consumer_data.lock);
+
+ ht = consumer_data.stream_list_ht;
+
+ /*
+ * We only need a reference to the channel, but they are not
+ * directly indexed, so we just use the first matching stream
+ * to extract the information we need, we default to 0 if not
+ * found (no packets lost if the channel is not yet in use).
+ */
+ ret = 0;
+ cds_lfht_for_each_entry_duplicate(ht->ht,
+ ht->hash_fct(&id, lttng_ht_seed),
+ ht->match_fct, &id,
+ &iter.iter, stream, node_session_id.node) {
+ if (stream->chan->key == key) {
+ ret = stream->chan->lost_packets;
+ break;
+ }
+ }
+ pthread_mutex_unlock(&consumer_data.lock);
+ rcu_read_unlock();
+
+ DBG("UST consumer lost packets command for session id %"
+ PRIu64 ", channel key %" PRIu64, id, key);
+
+ health_code_update();
+
+ /* Send back returned value to session daemon */
+ ret = lttcomm_send_unix_sock(sock, &ret, sizeof(ret));
+ if (ret < 0) {
+ PERROR("send lost packets");
+ goto error_fatal;
+ }
+
+ break;
+ }
default:
break;
}
return ustctl_get_current_timestamp(stream->ustream, ts);
}
+int lttng_ustconsumer_get_sequence_number(
+ struct lttng_consumer_stream *stream, uint64_t *seq)
+{
+ assert(stream);
+ assert(stream->ustream);
+ assert(seq);
+
+ return ustctl_get_sequence_number(stream->ustream, seq);
+}
+
/*
* Called when the stream signal the consumer that it has hang up.
*/
void lttng_ustconsumer_del_channel(struct lttng_consumer_channel *chan)
{
+ int i;
+
assert(chan);
assert(chan->uchan);
if (chan->switch_timer_enabled == 1) {
consumer_timer_switch_stop(chan);
}
+ for (i = 0; i < chan->nr_stream_fds; i++) {
+ int ret;
+
+ ret = close(chan->stream_fds[i]);
+ if (ret) {
+ PERROR("close");
+ }
+ if (chan->shm_path[0]) {
+ char shm_path[PATH_MAX];
+
+ ret = get_stream_shm_path(shm_path, chan->shm_path, i);
+ if (ret) {
+ ERR("Cannot get stream shm path");
+ }
+ ret = run_as_unlink(shm_path, chan->uid, chan->gid);
+ if (ret) {
+ PERROR("unlink %s", shm_path);
+ }
+ }
+ }
+ /* Try to rmdir all directories under shm_path root. */
+ if (chan->root_shm_path[0]) {
+ (void) run_as_recursive_rmdir(chan->root_shm_path,
+ chan->uid, chan->gid);
+ }
+}
+
+void lttng_ustconsumer_free_channel(struct lttng_consumer_channel *chan)
+{
+ assert(chan);
+ assert(chan->uchan);
+
consumer_metadata_cache_destroy(chan);
ustctl_destroy_channel(chan->uchan);
+ free(chan->stream_fds);
}
void lttng_ustconsumer_del_stream(struct lttng_consumer_stream *stream)
int ret;
pthread_mutex_lock(&stream->chan->metadata_cache->lock);
- if (stream->chan->metadata_cache->contiguous
+ if (stream->chan->metadata_cache->max_offset
== stream->ust_metadata_pushed) {
ret = 0;
goto end;
write_len = ustctl_write_one_packet_to_channel(stream->chan->uchan,
&stream->chan->metadata_cache->data[stream->ust_metadata_pushed],
- stream->chan->metadata_cache->contiguous
+ stream->chan->metadata_cache->max_offset
- stream->ust_metadata_pushed);
assert(write_len != 0);
if (write_len < 0) {
}
stream->ust_metadata_pushed += write_len;
- assert(stream->chan->metadata_cache->contiguous >=
+ assert(stream->chan->metadata_cache->max_offset >=
stream->ust_metadata_pushed);
ret = write_len;
* Sync metadata meaning request them to the session daemon and snapshot to the
* metadata thread can consumer them.
*
- * Metadata stream lock MUST be acquired.
+ * Metadata stream lock is held here, but we need to release it when
+ * interacting with sessiond, else we cause a deadlock with live
+ * awaiting on metadata to be pushed out.
*
* Return 0 if new metadatda is available, EAGAIN if the metadata stream
* is empty or a negative value on error.
assert(ctx);
assert(metadata);
+ pthread_mutex_unlock(&metadata->lock);
/*
* Request metadata from the sessiond, but don't wait for the flush
* because we locked the metadata thread.
if (ret < 0) {
goto end;
}
+ pthread_mutex_lock(&metadata->lock);
ret = commit_one_metadata_packet(metadata);
if (ret <= 0) {
return ret;
}
+static
+int update_stream_stats(struct lttng_consumer_stream *stream)
+{
+ int ret;
+ uint64_t seq, discarded;
+
+ ret = ustctl_get_sequence_number(stream->ustream, &seq);
+ if (ret < 0) {
+ PERROR("ustctl_get_sequence_number");
+ goto end;
+ }
+ /*
+ * Start the sequence when we extract the first packet in case we don't
+ * start at 0 (for example if a consumer is not connected to the
+ * session immediately after the beginning).
+ */
+ if (stream->last_sequence_number == -1ULL) {
+ stream->last_sequence_number = seq;
+ } else if (seq > stream->last_sequence_number) {
+ stream->chan->lost_packets += seq -
+ stream->last_sequence_number - 1;
+ } else {
+ /* seq <= last_sequence_number */
+ ERR("Sequence number inconsistent : prev = %" PRIu64
+ ", current = %" PRIu64,
+ stream->last_sequence_number, seq);
+ ret = -1;
+ goto end;
+ }
+ stream->last_sequence_number = seq;
+
+ ret = ustctl_get_events_discarded(stream->ustream, &discarded);
+ if (ret < 0) {
+ PERROR("kernctl_get_events_discarded");
+ goto end;
+ }
+ if (discarded < stream->last_discarded_events) {
+ /*
+ * Overflow has occured. We assume only one wrap-around
+ * has occured.
+ */
+ stream->chan->discarded_events +=
+ (1ULL << (CAA_BITS_PER_LONG - 1)) -
+ stream->last_discarded_events + discarded;
+ } else {
+ stream->chan->discarded_events += discarded -
+ stream->last_discarded_events;
+ }
+ stream->last_discarded_events = discarded;
+ ret = 0;
+
+end:
+ return ret;
+}
+
/*
* Read subbuffer from the given stream.
*
if (ret < 0) {
goto end;
}
+
+ /* Update the stream's sequence and discarded events count. */
+ ret = update_stream_stats(stream);
+ if (ret < 0) {
+ PERROR("kernctl_get_events_discarded");
+ goto end;
+ }
} else {
write_index = 0;
}
/*
* In live, block until all the metadata is sent.
*/
+ pthread_mutex_lock(&stream->metadata_timer_lock);
+ assert(!stream->missed_metadata_flush);
+ stream->waiting_on_metadata = true;
+ pthread_mutex_unlock(&stream->metadata_timer_lock);
+
err = consumer_stream_sync_metadata(ctx, stream->session_id);
+
+ pthread_mutex_lock(&stream->metadata_timer_lock);
+ stream->waiting_on_metadata = false;
+ if (stream->missed_metadata_flush) {
+ stream->missed_metadata_flush = false;
+ pthread_mutex_unlock(&stream->metadata_timer_lock);
+ (void) consumer_flush_ust_index(stream);
+ } else {
+ pthread_mutex_unlock(&stream->metadata_timer_lock);
+ }
+
if (err < 0) {
goto end;
}
uint64_t contiguous, pushed;
/* Ease our life a bit. */
- contiguous = stream->chan->metadata_cache->contiguous;
+ contiguous = stream->chan->metadata_cache->max_offset;
pushed = stream->ust_metadata_pushed;
/*
* function or any of its callees. Timers have a very strict locking
* semantic with respect to teardown. Failure to respect this semantic
* introduces deadlocks.
+ *
+ * DON'T hold the metadata lock when calling this function, else this
+ * can cause deadlock involving consumer awaiting for metadata to be
+ * pushed out due to concurrent interaction with the session daemon.
*/
int lttng_ustconsumer_request_metadata(struct lttng_consumer_local_data *ctx,
struct lttng_consumer_channel *channel, int timer, int wait)