- viewer_stream->total_index_received = stream->total_index_received;
-
- /*
- * If we never received an index for the current stream, delay
- * the opening of the index, otherwise open it right now.
- */
- if (viewer_stream->tracefile_count_current ==
- stream->tracefile_count_current &&
- viewer_stream->total_index_received == 0) {
- viewer_stream->index_read_fd = -1;
- } else {
- ret = open_index(viewer_stream);
- if (ret < 0) {
- goto error;
- }
- }
-
- if (seek_last && viewer_stream->index_read_fd > 0) {
- ret = lseek(viewer_stream->index_read_fd,
- viewer_stream->total_index_received *
- sizeof(struct ctf_packet_index),
- SEEK_CUR);
- if (ret < 0) {
- goto error;
- }
- viewer_stream->last_sent_index =
- viewer_stream->total_index_received;
- }
-
- ret = 0;
-
-error:
- return ret;
-}
-
-/*
- * Rotate a stream to the next tracefile.
- *
- * Returns 0 on success, 1 on EOF, a negative value on error.
- */
-static
-int rotate_viewer_stream(struct relay_viewer_stream *viewer_stream,
- struct relay_stream *stream)
-{
- int ret;
- uint64_t tracefile_id;
-
- assert(viewer_stream);
-
- tracefile_id = (viewer_stream->tracefile_count_current + 1) %
- viewer_stream->tracefile_count;
- /*
- * Detect the last tracefile to open.
- */
- if (viewer_stream->tracefile_count_last != -1ULL &&
- viewer_stream->tracefile_count_last ==
- viewer_stream->tracefile_count_current) {
- ret = 1;
- goto end;
- }
-
- if (stream) {
- pthread_mutex_lock(&stream->viewer_stream_rotation_lock);
- }
- /*
- * The writer and the reader are not working in the same
- * tracefile, we can read up to EOF, we don't care about the
- * total_index_received.
- */
- if (!stream || (stream->tracefile_count_current != tracefile_id)) {
- viewer_stream->close_write_flag = 1;
- } else {
- /*
- * We are opening a file that is still open in write, make
- * sure we limit our reading to the number of indexes
- * received.
- */
- viewer_stream->close_write_flag = 0;
- if (stream) {
- viewer_stream->total_index_received =
- stream->total_index_received;
- }
- }
- viewer_stream->tracefile_count_current = tracefile_id;
-
- ret = close(viewer_stream->index_read_fd);
- if (ret < 0) {
- PERROR("close index file %d",
- viewer_stream->index_read_fd);
- }
- viewer_stream->index_read_fd = -1;
- ret = close(viewer_stream->read_fd);
- if (ret < 0) {
- PERROR("close tracefile %d",
- viewer_stream->read_fd);
- }
- viewer_stream->read_fd = -1;
-
- pthread_mutex_lock(&viewer_stream->overwrite_lock);
- viewer_stream->abort_flag = 0;
- pthread_mutex_unlock(&viewer_stream->overwrite_lock);
-
- viewer_stream->index_read_fd = -1;
- viewer_stream->read_fd = -1;
-
- if (stream) {
- pthread_mutex_unlock(&stream->viewer_stream_rotation_lock);
- }
- ret = open_index(viewer_stream);
- if (ret < 0) {
- goto error;
- }
-
- ret = 0;
-
-end:
-error:
- return ret;
-}
-
-/*
- * Send the viewer the list of current sessions.
- */
-static
-int viewer_get_new_streams(struct relay_command *cmd,
- struct lttng_ht *sessions_ht)
-{
- int ret, send_streams = 0;
- uint32_t nb_streams = 0;
- struct lttng_viewer_new_streams_request request;
- struct lttng_viewer_new_streams_response response;
- struct lttng_viewer_stream send_stream;
- struct relay_stream *stream;
- struct relay_viewer_stream *viewer_stream;
- struct lttng_ht_node_ulong *node;
- struct lttng_ht_iter iter;
- struct relay_session *session;
-
- assert(cmd);
- assert(sessions_ht);
-
- DBG("Get new streams received");
-
- if (cmd->version_check_done == 0) {
- ERR("Trying to get streams before version check");
- ret = -1;
- goto end_no_session;
- }
-
- health_code_update();
-
- ret = cmd->sock->ops->recvmsg(cmd->sock, &request, sizeof(request), 0);
- if (ret < 0 || ret != sizeof(request)) {
- if (ret == 0) {
- /* Orderly shutdown. Not necessary to print an error. */
- DBG("Socket %d did an orderly shutdown", cmd->sock->fd);
- } else {
- ERR("Relay failed to receive the command parameters.");
- }
- ret = -1;
- goto error;
- }
-
- health_code_update();
-
- rcu_read_lock();
- lttng_ht_lookup(sessions_ht,
- (void *)((unsigned long) be64toh(request.session_id)), &iter);
- node = lttng_ht_iter_get_node_ulong(&iter);
- if (node == NULL) {
- DBG("Relay session %" PRIu64 " not found",
- be64toh(request.session_id));
- response.status = htobe32(VIEWER_NEW_STREAMS_ERR);
- goto send_reply;
- }
-
- session = caa_container_of(node, struct relay_session, session_n);
- if (cmd->session_id == session->id) {
- /* We confirmed the viewer is asking for the same session. */
- send_streams = 1;
- response.status = htobe32(VIEWER_NEW_STREAMS_OK);
- } else {
- send_streams = 0;
- response.status = htobe32(VIEWER_NEW_STREAMS_ERR);
- goto send_reply;
- }
-
- /*
- * Fill the viewer_streams_ht to count the number of streams ready to be
- * sent and avoid concurrency issues on the relay_streams_ht and don't rely
- * on a total session stream count.
- */
- pthread_mutex_lock(&session->viewer_ready_lock);
- cds_lfht_for_each_entry(relay_streams_ht->ht, &iter.iter, stream,
- stream_n.node) {
- struct relay_viewer_stream *vstream;
-
- health_code_update();
-
- /*
- * Don't send stream if no ctf_trace, wrong session or if the stream is
- * not ready for the viewer.
- */
- if (stream->session != cmd->session ||
- !stream->ctf_trace || !stream->viewer_ready) {
- continue;
- }
-
- vstream = live_find_viewer_stream_by_id(stream->stream_handle);
- if (!vstream) {
- ret = init_viewer_stream(stream, 0);
- if (ret < 0) {
- pthread_mutex_unlock(&session->viewer_ready_lock);
- goto end_unlock;
- }
- nb_streams++;
- } else if (!vstream->sent_flag) {
- nb_streams++;
- }
- }
- pthread_mutex_unlock(&session->viewer_ready_lock);
-
- response.streams_count = htobe32(nb_streams);