- goto end;
- }
-
- /*
- * Rewind the current tracefile to the position at which the rotation
- * should have occured.
- */
- lseek_ret = lseek(stream->stream_fd->fd,
- stream->pos_after_last_complete_data_index, SEEK_SET);
- if (lseek_ret < 0) {
- PERROR("seek truncate stream");
- ret = -1;
- goto end;
- }
-
- /* Move data from the old file to the new file. */
- while (pos < diff) {
- uint64_t count, bytes_left;
- ssize_t io_ret;
-
- bytes_left = diff - pos;
- count = bytes_left > sizeof(buf) ? sizeof(buf) : bytes_left;
- assert(count <= SIZE_MAX);
-
- io_ret = lttng_read(stream->stream_fd->fd, buf, count);
- if (io_ret < (ssize_t) count) {
- char error_string[256];
-
- snprintf(error_string, sizeof(error_string),
- "Failed to read %" PRIu64 " bytes from fd %i in rotate_truncate_stream(), returned %zi",
- count, stream->stream_fd->fd, io_ret);
- if (io_ret == -1) {
- PERROR("%s", error_string);
- } else {
- ERR("%s", error_string);
- }
- ret = -1;
- goto end;
- }
-
- io_ret = lttng_write(new_fd, buf, count);
- if (io_ret < (ssize_t) count) {
- char error_string[256];
-
- snprintf(error_string, sizeof(error_string),
- "Failed to write %" PRIu64 " bytes from fd %i in rotate_truncate_stream(), returned %zi",
- count, new_fd, io_ret);
- if (io_ret == -1) {
- PERROR("%s", error_string);
- } else {
- ERR("%s", error_string);
- }
- ret = -1;
- goto end;
- }
-
- pos += count;
- }
-
- /* Truncate the file to get rid of the excess data. */
- ret = ftruncate(stream->stream_fd->fd,
- stream->pos_after_last_complete_data_index);
- if (ret) {
- PERROR("ftruncate");
- goto end;
- }
-
- ret = close(stream->stream_fd->fd);
- if (ret < 0) {
- PERROR("Closing tracefile");
- goto end;
- }
-
- ret = create_rotate_index_file(stream);
- if (ret < 0) {
- ERR("Rotate stream index file");
- goto end;
- }
-
- /*
- * Update the offset and FD of all the eventual indexes created by the
- * data connection before the rotation command arrived.
- */
- ret = relay_index_switch_all_files(stream);
- if (ret < 0) {
- ERR("Failed to rotate index file");
- goto end;
- }
-
- stream->stream_fd->fd = new_fd;
- stream->tracefile_size_current = diff;
- stream->pos_after_last_complete_data_index = 0;
- stream->rotate_at_seq_num = -1ULL;
-
- ret = 0;
-
-end:
- return ret;
-}
-
-/*
- * Check if a stream should perform a rotation (for session rotation).
- * Must be called with the stream lock held.
- *
- * Return 0 on success, a negative value on error.
- */
-static
-int try_rotate_stream(struct relay_stream *stream)
-{
- int ret = 0;
-
- /* No rotation expected. */
- if (stream->rotate_at_seq_num == -1ULL) {
- goto end;
- }
-
- if (stream->prev_seq < stream->rotate_at_seq_num ||
- stream->prev_seq == -1ULL) {
- DBG("Stream %" PRIu64 " no yet ready for rotation",
- stream->stream_handle);
- goto end;
- } else if (stream->prev_seq > stream->rotate_at_seq_num) {
- DBG("Rotation after too much data has been written in tracefile "
- "for stream %" PRIu64 ", need to truncate before "
- "rotating", stream->stream_handle);
- ret = rotate_truncate_stream(stream);
- if (ret) {
- ERR("Failed to truncate stream");
- goto end;
- }
- } else {
- /* stream->prev_seq == stream->rotate_at_seq_num */
- DBG("Stream %" PRIu64 " ready for rotation",
- stream->stream_handle);
- ret = do_rotate_stream(stream);
- }
-
-end:
- return ret;
-}
-
-/*
- * relay_recv_metadata: receive the metadata for the session.
- */
-static int relay_recv_metadata(const struct lttcomm_relayd_hdr *recv_hdr,
- struct relay_connection *conn,
- const struct lttng_buffer_view *payload)
-{
- int ret = 0;
- ssize_t size_ret;
- struct relay_session *session = conn->session;
- struct lttcomm_relayd_metadata_payload metadata_payload_header;
- struct relay_stream *metadata_stream;
- uint64_t metadata_payload_size;
-
- if (!session) {
- ERR("Metadata sent before version check");
- ret = -1;
- goto end;
- }
-
- if (recv_hdr->data_size < sizeof(struct lttcomm_relayd_metadata_payload)) {
- ERR("Incorrect data size");
- ret = -1;
- goto end;
- }
- metadata_payload_size = recv_hdr->data_size -
- sizeof(struct lttcomm_relayd_metadata_payload);
-
- memcpy(&metadata_payload_header, payload->data,
- sizeof(metadata_payload_header));
- metadata_payload_header.stream_id = be64toh(
- metadata_payload_header.stream_id);
- metadata_payload_header.padding_size = be32toh(
- metadata_payload_header.padding_size);
-
- metadata_stream = stream_get_by_id(metadata_payload_header.stream_id);
- if (!metadata_stream) {
- ret = -1;
- goto end;
- }
-
- pthread_mutex_lock(&metadata_stream->lock);
-
- size_ret = lttng_write(metadata_stream->stream_fd->fd,
- payload->data + sizeof(metadata_payload_header),
- metadata_payload_size);
- if (size_ret < metadata_payload_size) {
- ERR("Relay error writing metadata on file");
- ret = -1;
- goto end_put;
- }
-
- size_ret = write_padding_to_file(metadata_stream->stream_fd->fd,
- metadata_payload_header.padding_size);
- if (size_ret < (int64_t) metadata_payload_header.padding_size) {
- ret = -1;
- goto end_put;
- }
-
- metadata_stream->metadata_received +=
- metadata_payload_size + metadata_payload_header.padding_size;
- DBG2("Relay metadata written. Updated metadata_received %" PRIu64,
- metadata_stream->metadata_received);
-
- ret = try_rotate_stream(metadata_stream);
- if (ret < 0) {
- goto end_put;
- }
-
-end_put:
- pthread_mutex_unlock(&metadata_stream->lock);
- stream_put(metadata_stream);
-end:
- return ret;
-}
-
-/*
- * relay_send_version: send relayd version number
- */
-static int relay_send_version(const struct lttcomm_relayd_hdr *recv_hdr,
- struct relay_connection *conn,
- const struct lttng_buffer_view *payload)
-{
- int ret;
- ssize_t send_ret;
- struct lttcomm_relayd_version reply, msg;
- bool compatible = true;
-
- conn->version_check_done = true;
-
- /* Get version from the other side. */
- if (payload->size < sizeof(msg)) {
- ERR("Unexpected payload size in \"relay_send_version\": expected >= %zu bytes, got %zu bytes",
- sizeof(msg), payload->size);
- ret = -1;
- goto end;
- }
-
- memcpy(&msg, payload->data, sizeof(msg));
- msg.major = be32toh(msg.major);
- msg.minor = be32toh(msg.minor);
-
- memset(&reply, 0, sizeof(reply));
- reply.major = RELAYD_VERSION_COMM_MAJOR;
- reply.minor = RELAYD_VERSION_COMM_MINOR;
-
- /* Major versions must be the same */
- if (reply.major != msg.major) {
- DBG("Incompatible major versions (%u vs %u), deleting session",
- reply.major, msg.major);
- compatible = false;
- }
-
- conn->major = reply.major;
- /* We adapt to the lowest compatible version */
- if (reply.minor <= msg.minor) {
- conn->minor = reply.minor;
- } else {
- conn->minor = msg.minor;
- }
-
- reply.major = htobe32(reply.major);
- reply.minor = htobe32(reply.minor);
- send_ret = conn->sock->ops->sendmsg(conn->sock, &reply,
- sizeof(reply), 0);
- if (send_ret < (ssize_t) sizeof(reply)) {
- ERR("Failed to send \"send version\" command reply (ret = %zd)",
- send_ret);
- ret = -1;
- goto end;
- } else {
- ret = 0;
- }
-
- if (!compatible) {
- ret = -1;
- goto end;
- }
-
- DBG("Version check done using protocol %u.%u", conn->major,
- conn->minor);
-
-end:
- return ret;
-}
-
-/*
- * Check for data pending for a given stream id from the session daemon.
- */
-static int relay_data_pending(const struct lttcomm_relayd_hdr *recv_hdr,
- struct relay_connection *conn,
- const struct lttng_buffer_view *payload)
-{
- struct relay_session *session = conn->session;
- struct lttcomm_relayd_data_pending msg;
- struct lttcomm_relayd_generic_reply reply;
- struct relay_stream *stream;
- ssize_t send_ret;
- int ret;
-
- DBG("Data pending command received");
-
- if (!session || !conn->version_check_done) {
- ERR("Trying to check for data before version check");
- ret = -1;
- goto end_no_session;