+
+static
+int consumer_flush_buffer(struct lttng_consumer_stream *stream, int producer_active)
+{
+ int ret = 0;
+
+ switch (consumer_data.type) {
+ case LTTNG_CONSUMER_KERNEL:
+ ret = kernctl_buffer_flush(stream->wait_fd);
+ if (ret < 0) {
+ ERR("Failed to flush kernel stream");
+ goto end;
+ }
+ break;
+ case LTTNG_CONSUMER32_UST:
+ case LTTNG_CONSUMER64_UST:
+ lttng_ustctl_flush_buffer(stream, producer_active);
+ break;
+ default:
+ ERR("Unknown consumer_data type");
+ abort();
+ }
+
+end:
+ return ret;
+}
+
+/*
+ * Sample the rotate position for all the streams of a channel. If a stream
+ * is already at the rotate position (produced == consumed), we flag it as
+ * ready for rotation. The rotation of ready streams occurs after we have
+ * replied to the session daemon that we have finished sampling the positions.
+ *
+ * Returns 0 on success, < 0 on error
+ */
+int lttng_consumer_rotate_channel(uint64_t key, const char *path,
+ uint64_t relayd_id, uint32_t metadata, uint64_t new_chunk_id,
+ struct lttng_consumer_local_data *ctx)
+{
+ int ret;
+ struct lttng_consumer_channel *channel;
+ struct lttng_consumer_stream *stream;
+ struct lttng_ht_iter iter;
+ struct lttng_ht *ht = consumer_data.stream_per_chan_id_ht;
+
+ DBG("Consumer sample rotate position for channel %" PRIu64, key);
+
+ rcu_read_lock();
+
+ channel = consumer_find_channel(key);
+ if (!channel) {
+ ERR("No channel found for key %" PRIu64, key);
+ ret = -1;
+ goto end;
+ }
+
+ pthread_mutex_lock(&channel->lock);
+ channel->current_chunk_id = new_chunk_id;
+
+ ret = lttng_strncpy(channel->pathname, path, sizeof(channel->pathname));
+ if (ret) {
+ ERR("Failed to copy new path to channel during channel rotation");
+ ret = -1;
+ goto end_unlock_channel;
+ }
+
+ if (relayd_id == -1ULL) {
+ /*
+ * The domain path (/ust or /kernel) has been created before, we
+ * now need to create the last part of the path: the application/user
+ * specific section (uid/1000/64-bit).
+ */
+ ret = utils_mkdir_recursive(channel->pathname, S_IRWXU | S_IRWXG,
+ channel->uid, channel->gid);
+ if (ret < 0) {
+ ERR("Failed to create trace directory at %s during rotation",
+ channel->pathname);
+ ret = -1;
+ goto end_unlock_channel;
+ }
+ }
+
+ cds_lfht_for_each_entry_duplicate(ht->ht,
+ ht->hash_fct(&channel->key, lttng_ht_seed),
+ ht->match_fct, &channel->key, &iter.iter,
+ stream, node_channel_id.node) {
+ unsigned long consumed_pos;
+
+ health_code_update();
+
+ /*
+ * Lock stream because we are about to change its state.
+ */
+ pthread_mutex_lock(&stream->lock);
+
+ ret = lttng_strncpy(stream->channel_read_only_attributes.path,
+ channel->pathname,
+ sizeof(stream->channel_read_only_attributes.path));
+ if (ret) {
+ ERR("Failed to sample channel path name during channel rotation");
+ goto end_unlock_stream;
+ }
+ ret = lttng_consumer_sample_snapshot_positions(stream);
+ if (ret < 0) {
+ ERR("Failed to sample snapshot position during channel rotation");
+ goto end_unlock_stream;
+ }
+
+ ret = lttng_consumer_get_produced_snapshot(stream,
+ &stream->rotate_position);
+ if (ret < 0) {
+ ERR("Failed to sample produced position during channel rotation");
+ goto end_unlock_stream;
+ }
+
+ lttng_consumer_get_consumed_snapshot(stream,
+ &consumed_pos);
+ if (consumed_pos == stream->rotate_position) {
+ stream->rotate_ready = true;
+ }
+ channel->nr_stream_rotate_pending++;
+
+ ret = consumer_flush_buffer(stream, 1);
+ if (ret < 0) {
+ ERR("Failed to flush stream %" PRIu64 " during channel rotation",
+ stream->key);
+ goto end_unlock_stream;
+ }
+
+ pthread_mutex_unlock(&stream->lock);
+ }
+ pthread_mutex_unlock(&channel->lock);
+
+ ret = 0;
+ goto end;
+
+end_unlock_stream:
+ pthread_mutex_unlock(&stream->lock);
+end_unlock_channel:
+ pthread_mutex_unlock(&channel->lock);
+end:
+ rcu_read_unlock();
+ return ret;
+}
+
+/*
+ * Check if a stream is ready to be rotated after extracting it.
+ *
+ * Return 1 if it is ready for rotation, 0 if it is not, a negative value on
+ * error. Stream lock must be held.
+ */
+int lttng_consumer_stream_is_rotate_ready(struct lttng_consumer_stream *stream)
+{
+ int ret;
+ unsigned long consumed_pos;
+
+ if (!stream->rotate_position && !stream->rotate_ready) {
+ ret = 0;
+ goto end;
+ }
+
+ if (stream->rotate_ready) {
+ ret = 1;
+ goto end;
+ }
+
+ /*
+ * If we don't have the rotate_ready flag, check the consumed position
+ * to determine if we need to rotate.
+ */
+ ret = lttng_consumer_sample_snapshot_positions(stream);
+ if (ret < 0) {
+ ERR("Taking snapshot positions");
+ goto end;
+ }
+
+ ret = lttng_consumer_get_consumed_snapshot(stream, &consumed_pos);
+ if (ret < 0) {
+ ERR("Consumed snapshot position");
+ goto end;
+ }
+
+ /* Rotate position not reached yet (with check for overflow). */
+ if ((long) (consumed_pos - stream->rotate_position) < 0) {
+ ret = 0;
+ goto end;
+ }
+ ret = 1;
+
+end:
+ return ret;
+}
+
+/*
+ * Reset the state for a stream after a rotation occurred.
+ */
+void lttng_consumer_reset_stream_rotate_state(struct lttng_consumer_stream *stream)
+{
+ stream->rotate_position = 0;
+ stream->rotate_ready = false;
+}
+
+/*
+ * Perform the rotation a local stream file.
+ */
+int rotate_local_stream(struct lttng_consumer_local_data *ctx,
+ struct lttng_consumer_stream *stream)
+{
+ int ret;
+
+ DBG("Rotate local stream: stream key %" PRIu64 ", channel key %" PRIu64 " at path %s",
+ stream->key,
+ stream->chan->key,
+ stream->channel_read_only_attributes.path);
+
+ ret = close(stream->out_fd);
+ if (ret < 0) {
+ PERROR("Closing trace file (fd %d), stream %" PRIu64,
+ stream->out_fd, stream->key);
+ assert(0);
+ goto error;
+ }
+
+ ret = utils_create_stream_file(
+ stream->channel_read_only_attributes.path,
+ stream->name,
+ stream->channel_read_only_attributes.tracefile_size,
+ stream->tracefile_count_current,
+ stream->uid, stream->gid, NULL);
+ if (ret < 0) {
+ ERR("Rotate create stream file");
+ goto error;
+ }
+ stream->out_fd = ret;
+ stream->tracefile_size_current = 0;
+
+ if (!stream->metadata_flag) {
+ struct lttng_index_file *index_file;
+
+ lttng_index_file_put(stream->index_file);
+
+ index_file = lttng_index_file_create(
+ stream->channel_read_only_attributes.path,
+ stream->name, stream->uid, stream->gid,
+ stream->channel_read_only_attributes.tracefile_size,
+ stream->tracefile_count_current,
+ CTF_INDEX_MAJOR, CTF_INDEX_MINOR);
+ if (!index_file) {
+ ERR("Create index file during rotation");
+ goto error;
+ }
+ stream->index_file = index_file;
+ stream->out_fd_offset = 0;
+ }
+
+ ret = 0;
+ goto end;
+
+error:
+ ret = -1;
+end:
+ return ret;
+
+}
+
+/*
+ * Perform the rotation a stream file on the relay.
+ */
+int rotate_relay_stream(struct lttng_consumer_local_data *ctx,
+ struct lttng_consumer_stream *stream)
+{
+ int ret;
+ struct consumer_relayd_sock_pair *relayd;
+
+ DBG("Rotate relay stream");
+ relayd = consumer_find_relayd(stream->net_seq_idx);
+ if (!relayd) {
+ ERR("Failed to find relayd");
+ ret = -1;
+ goto end;
+ }
+
+ pthread_mutex_lock(&relayd->ctrl_sock_mutex);
+ ret = relayd_rotate_stream(&relayd->control_sock,
+ stream->relayd_stream_id,
+ stream->channel_read_only_attributes.path,
+ stream->chan->current_chunk_id,
+ stream->last_sequence_number);
+ pthread_mutex_unlock(&relayd->ctrl_sock_mutex);
+ if (ret < 0) {
+ ERR("Relayd rotate stream failed. Cleaning up relayd %" PRIu64".", relayd->net_seq_idx);
+ lttng_consumer_cleanup_relayd(relayd);
+ }
+ if (ret) {
+ ERR("Rotate relay stream");
+ }
+
+end:
+ return ret;
+}
+
+/*
+ * Performs the stream rotation for the rotate session feature if needed.
+ * It must be called with the stream lock held.
+ *
+ * Return 0 on success, a negative number of error.
+ */
+int lttng_consumer_rotate_stream(struct lttng_consumer_local_data *ctx,
+ struct lttng_consumer_stream *stream, bool *rotated)
+{
+ int ret;
+
+ DBG("Consumer rotate stream %" PRIu64, stream->key);
+
+ if (stream->net_seq_idx != (uint64_t) -1ULL) {
+ ret = rotate_relay_stream(ctx, stream);
+ } else {
+ ret = rotate_local_stream(ctx, stream);
+ }
+ if (ret < 0) {
+ ERR("Rotate stream");
+ goto error;
+ }
+
+ if (stream->metadata_flag) {
+ switch (consumer_data.type) {
+ case LTTNG_CONSUMER_KERNEL:
+ /*
+ * Reset the position of what has been read from the metadata
+ * cache to 0 so we can dump it again.
+ */
+ ret = kernctl_metadata_cache_dump(stream->wait_fd);
+ if (ret < 0) {
+ ERR("Failed to dump the kernel metadata cache after rotation");
+ goto error;
+ }
+ break;
+ case LTTNG_CONSUMER32_UST:
+ case LTTNG_CONSUMER64_UST:
+ /*
+ * Reset the position pushed from the metadata cache so it
+ * will write from the beginning on the next push.
+ */
+ stream->ust_metadata_pushed = 0;
+ break;
+ default:
+ ERR("Unknown consumer_data type");
+ abort();
+ }
+ }
+ lttng_consumer_reset_stream_rotate_state(stream);
+
+ if (rotated) {
+ *rotated = true;
+ }
+
+ ret = 0;
+
+error:
+ return ret;
+}
+
+/*
+ * Rotate all the ready streams now.
+ *
+ * This is especially important for low throughput streams that have already
+ * been consumed, we cannot wait for their next packet to perform the
+ * rotation.
+ *
+ * Returns 0 on success, < 0 on error
+ */
+int lttng_consumer_rotate_ready_streams(uint64_t key,
+ struct lttng_consumer_local_data *ctx)
+{
+ int ret;
+ struct lttng_consumer_channel *channel;
+ struct lttng_consumer_stream *stream;
+ struct lttng_ht_iter iter;
+ struct lttng_ht *ht = consumer_data.stream_per_chan_id_ht;
+
+ rcu_read_lock();
+
+ DBG("Consumer rotate ready streams in channel %" PRIu64, key);
+
+ channel = consumer_find_channel(key);
+ if (!channel) {
+ ERR("No channel found for key %" PRIu64, key);
+ ret = -1;
+ goto end;
+ }
+
+ cds_lfht_for_each_entry_duplicate(ht->ht,
+ ht->hash_fct(&channel->key, lttng_ht_seed),
+ ht->match_fct, &channel->key, &iter.iter,
+ stream, node_channel_id.node) {
+ health_code_update();
+
+ pthread_mutex_lock(&stream->lock);
+
+ if (!stream->rotate_ready) {
+ pthread_mutex_unlock(&stream->lock);
+ continue;
+ }
+ DBG("Consumer rotate ready stream %" PRIu64, stream->key);
+
+ ret = lttng_consumer_rotate_stream(ctx, stream, NULL);
+ pthread_mutex_unlock(&stream->lock);
+ if (ret) {
+ goto end;
+ }
+
+ ret = consumer_post_rotation(stream, ctx);
+ if (ret) {
+ goto end;
+ }
+ }
+
+ ret = 0;
+
+end:
+ rcu_read_unlock();
+ return ret;
+}
+
+static
+int rotate_rename_local(const char *old_path, const char *new_path,
+ uid_t uid, gid_t gid)
+{
+ int ret;
+
+ assert(old_path);
+ assert(new_path);
+
+ ret = utils_mkdir_recursive(new_path, S_IRWXU | S_IRWXG, uid, gid);
+ if (ret < 0) {
+ ERR("Create directory on rotate");
+ goto end;
+ }
+
+ ret = rename(old_path, new_path);
+ if (ret < 0 && errno != ENOENT) {
+ PERROR("Rename completed rotation chunk");
+ goto end;
+ }
+
+ ret = 0;
+end:
+ return ret;
+}
+
+static
+int rotate_rename_relay(const char *old_path, const char *new_path,
+ uint64_t relayd_id)
+{
+ int ret;
+ struct consumer_relayd_sock_pair *relayd;
+
+ relayd = consumer_find_relayd(relayd_id);
+ if (!relayd) {
+ ERR("Failed to find relayd while running rotate_rename_relay command");
+ ret = -1;
+ goto end;
+ }
+
+ pthread_mutex_lock(&relayd->ctrl_sock_mutex);
+ ret = relayd_rotate_rename(&relayd->control_sock, old_path, new_path);
+ if (ret < 0) {
+ ERR("Relayd rotate rename failed. Cleaning up relayd %" PRIu64".", relayd->net_seq_idx);
+ lttng_consumer_cleanup_relayd(relayd);
+ }
+ pthread_mutex_unlock(&relayd->ctrl_sock_mutex);
+end:
+ return ret;
+}
+
+int lttng_consumer_rotate_rename(const char *old_path, const char *new_path,
+ uid_t uid, gid_t gid, uint64_t relayd_id)
+{
+ if (relayd_id != -1ULL) {
+ return rotate_rename_relay(old_path, new_path, relayd_id);
+ } else {
+ return rotate_rename_local(old_path, new_path, uid, gid);
+ }
+}
+
+int lttng_consumer_rotate_pending_relay(uint64_t session_id,
+ uint64_t relayd_id, uint64_t chunk_id)
+{
+ int ret;
+ struct consumer_relayd_sock_pair *relayd;
+
+ relayd = consumer_find_relayd(relayd_id);
+ if (!relayd) {
+ ERR("Failed to find relayd");
+ ret = -1;
+ goto end;
+ }
+
+ pthread_mutex_lock(&relayd->ctrl_sock_mutex);
+ ret = relayd_rotate_pending(&relayd->control_sock, chunk_id);
+ if (ret < 0) {
+ ERR("Relayd rotate pending failed. Cleaning up relayd %" PRIu64".", relayd->net_seq_idx);
+ lttng_consumer_cleanup_relayd(relayd);
+ }
+ pthread_mutex_unlock(&relayd->ctrl_sock_mutex);
+
+end:
+ return ret;
+}
+
+static
+int mkdir_local(const char *path, uid_t uid, gid_t gid)
+{
+ int ret;
+
+ ret = utils_mkdir_recursive(path, S_IRWXU | S_IRWXG, uid, gid);
+ if (ret < 0) {
+ /* utils_mkdir_recursive logs an error. */
+ goto end;
+ }
+
+ ret = 0;
+end:
+ return ret;
+}
+
+static
+int mkdir_relay(const char *path, uint64_t relayd_id)
+{
+ int ret;
+ struct consumer_relayd_sock_pair *relayd;
+
+ relayd = consumer_find_relayd(relayd_id);
+ if (!relayd) {
+ ERR("Failed to find relayd");
+ ret = -1;
+ goto end;
+ }
+
+ pthread_mutex_lock(&relayd->ctrl_sock_mutex);
+ ret = relayd_mkdir(&relayd->control_sock, path);
+ if (ret < 0) {
+ ERR("Relayd mkdir failed. Cleaning up relayd %" PRIu64".", relayd->net_seq_idx);
+ lttng_consumer_cleanup_relayd(relayd);
+ }
+ pthread_mutex_unlock(&relayd->ctrl_sock_mutex);
+
+end:
+ return ret;
+
+}
+
+int lttng_consumer_mkdir(const char *path, uid_t uid, gid_t gid,
+ uint64_t relayd_id)
+{
+ if (relayd_id != -1ULL) {
+ return mkdir_relay(path, relayd_id);
+ } else {
+ return mkdir_local(path, uid, gid);
+ }
+}