2 * Copyright (C) 2013 Julien Desfossez <jdesfossez@efficios.com>
3 * Copyright (C) 2013 David Goulet <dgoulet@efficios.com>
4 * Copyright (C) 2015 Mathieu Desnoyers <mathieu.desnoyers@efficios.com>
6 * SPDX-License-Identifier: GPL-2.0-only
12 #include "connection.hpp"
14 #include "lttng-relayd.hpp"
17 #include <common/common.hpp>
18 #include <common/compat/endian.hpp>
19 #include <common/urcu.hpp>
20 #include <common/utils.hpp>
23 * Allocate a new relay index object. Pass the stream in which it is
24 * contained as parameter. The sequence number will be used as the hash
27 * Called with stream mutex held.
28 * Return allocated object or else NULL on error.
30 static struct relay_index
*relay_index_create(struct relay_stream
*stream
, uint64_t net_seq_num
)
32 struct relay_index
*index
;
34 DBG2("Creating relay index for stream id %" PRIu64
" and seqnum %" PRIu64
,
35 stream
->stream_handle
,
38 index
= zmalloc
<relay_index
>();
40 PERROR("Relay index zmalloc");
43 if (!stream_get(stream
)) {
44 ERR("Cannot get stream");
49 index
->stream
= stream
;
51 lttng_ht_node_init_u64(&index
->index_n
, net_seq_num
);
52 pthread_mutex_init(&index
->lock
, nullptr);
53 urcu_ref_init(&index
->ref
);
60 * Add unique relay index to the given hash table. In case of a collision, the
61 * already existing object is put in the given _index variable.
63 * RCU read side lock MUST be acquired.
65 static struct relay_index
*relay_index_add_unique(struct relay_stream
*stream
,
66 struct relay_index
*index
)
68 struct cds_lfht_node
*node_ptr
;
69 struct relay_index
*_index
;
71 ASSERT_RCU_READ_LOCKED();
73 DBG2("Adding relay index with stream id %" PRIu64
" and seqnum %" PRIu64
,
74 stream
->stream_handle
,
77 node_ptr
= cds_lfht_add_unique(stream
->indexes_ht
->ht
,
78 stream
->indexes_ht
->hash_fct(&index
->index_n
, lttng_ht_seed
),
79 stream
->indexes_ht
->match_fct
,
81 &index
->index_n
.node
);
82 if (node_ptr
!= &index
->index_n
.node
) {
83 _index
= caa_container_of(node_ptr
, struct relay_index
, index_n
.node
);
91 * Should be called with RCU read-side lock held.
93 static bool relay_index_get(struct relay_index
*index
)
95 ASSERT_RCU_READ_LOCKED();
97 DBG2("index get for stream id %" PRIu64
" and seqnum %" PRIu64
" refcount %d",
98 index
->stream
->stream_handle
,
100 (int) index
->ref
.refcount
);
102 return urcu_ref_get_unless_zero(&index
->ref
);
106 * Get a relayd index in within the given stream, or create it if not
109 * Called with stream mutex held.
110 * Return index object or else NULL on error.
112 struct relay_index
*relay_index_get_by_id_or_create(struct relay_stream
*stream
,
113 uint64_t net_seq_num
)
115 struct lttng_ht_node_u64
*node
;
116 struct lttng_ht_iter iter
;
117 struct relay_index
*index
= nullptr;
119 DBG3("Finding index for stream id %" PRIu64
" and seq_num %" PRIu64
,
120 stream
->stream_handle
,
123 lttng::urcu::read_lock_guard read_lock
;
124 lttng_ht_lookup(stream
->indexes_ht
, &net_seq_num
, &iter
);
125 node
= lttng_ht_iter_get_node_u64(&iter
);
127 index
= lttng::utils::container_of(node
, &relay_index::index_n
);
129 struct relay_index
*oldindex
;
131 index
= relay_index_create(stream
, net_seq_num
);
133 ERR("Cannot create index for stream id %" PRIu64
" and seq_num %" PRIu64
,
134 stream
->stream_handle
,
138 oldindex
= relay_index_add_unique(stream
, index
);
140 /* Added concurrently, keep old. */
141 relay_index_put(index
);
143 if (!relay_index_get(index
)) {
147 stream
->indexes_in_flight
++;
148 index
->in_hash_table
= true;
152 DBG2("Index %sfound or created in HT for stream ID %" PRIu64
" and seqnum %" PRIu64
,
153 (index
== NULL
) ? "NOT " : "",
154 stream
->stream_handle
,
159 int relay_index_set_file(struct relay_index
*index
,
160 struct lttng_index_file
*index_file
,
161 uint64_t data_offset
)
165 pthread_mutex_lock(&index
->lock
);
166 if (index
->index_file
) {
170 lttng_index_file_get(index_file
);
171 index
->index_file
= index_file
;
172 index
->index_data
.offset
= data_offset
;
174 pthread_mutex_unlock(&index
->lock
);
178 int relay_index_set_data(struct relay_index
*index
, const struct ctf_packet_index
*data
)
182 pthread_mutex_lock(&index
->lock
);
183 if (index
->has_index_data
) {
187 /* Set everything except data_offset. */
188 index
->index_data
.packet_size
= data
->packet_size
;
189 index
->index_data
.content_size
= data
->content_size
;
190 index
->index_data
.timestamp_begin
= data
->timestamp_begin
;
191 index
->index_data
.timestamp_end
= data
->timestamp_end
;
192 index
->index_data
.events_discarded
= data
->events_discarded
;
193 index
->index_data
.stream_id
= data
->stream_id
;
194 index
->has_index_data
= true;
196 pthread_mutex_unlock(&index
->lock
);
200 static void index_destroy(struct relay_index
*index
)
205 static void index_destroy_rcu(struct rcu_head
*rcu_head
)
207 struct relay_index
*index
= lttng::utils::container_of(rcu_head
, &relay_index::rcu_node
);
209 index_destroy(index
);
212 /* Stream lock must be held by the caller. */
213 static void index_release(struct urcu_ref
*ref
)
215 struct relay_index
*index
= lttng::utils::container_of(ref
, &relay_index::ref
);
216 struct relay_stream
*stream
= index
->stream
;
218 struct lttng_ht_iter iter
;
220 if (index
->index_file
) {
221 lttng_index_file_put(index
->index_file
);
222 index
->index_file
= nullptr;
224 if (index
->in_hash_table
) {
225 /* Delete index from hash table. */
226 iter
.iter
.node
= &index
->index_n
.node
;
227 ret
= lttng_ht_del(stream
->indexes_ht
, &iter
);
229 stream
->indexes_in_flight
--;
232 stream_put(index
->stream
);
233 index
->stream
= nullptr;
235 call_rcu(&index
->rcu_node
, index_destroy_rcu
);
239 * Called with stream mutex held.
241 * Stream lock must be held by the caller.
243 void relay_index_put(struct relay_index
*index
)
245 DBG2("index put for stream id %" PRIu64
" and seqnum %" PRIu64
" refcount %d",
246 index
->stream
->stream_handle
,
248 (int) index
->ref
.refcount
);
250 * Ensure existence of index->lock for index unlock.
252 lttng::urcu::read_lock_guard read_lock
;
254 * Index lock ensures that concurrent test and update of stream
257 LTTNG_ASSERT(index
->ref
.refcount
!= 0);
258 urcu_ref_put(&index
->ref
, index_release
);
262 * Try to flush index to disk. Releases self-reference to index once
265 * Stream lock must be held by the caller.
266 * Return 0 on successful flush, a negative value on error, or positive
267 * value if no flush was performed.
269 int relay_index_try_flush(struct relay_index
*index
)
272 bool flushed
= false;
274 pthread_mutex_lock(&index
->lock
);
275 if (index
->flushed
) {
278 /* Check if we are ready to flush. */
279 if (!index
->has_index_data
|| !index
->index_file
) {
283 DBG2("Writing index for stream ID %" PRIu64
" and seq num %" PRIu64
,
284 index
->stream
->stream_handle
,
287 index
->flushed
= true;
288 ret
= lttng_index_file_write(index
->index_file
, &index
->index_data
);
290 pthread_mutex_unlock(&index
->lock
);
293 /* Put self-ref from index now that it has been flushed. */
294 relay_index_put(index
);
300 * Close every relay index within a given stream, without flushing
303 void relay_index_close_all(struct relay_stream
*stream
)
305 struct lttng_ht_iter iter
;
306 struct relay_index
*index
;
308 lttng::urcu::read_lock_guard read_lock
;
310 cds_lfht_for_each_entry (stream
->indexes_ht
->ht
, &iter
.iter
, index
, index_n
.node
) {
311 /* Put self-ref from index. */
312 relay_index_put(index
);
316 void relay_index_close_partial_fd(struct relay_stream
*stream
)
318 struct lttng_ht_iter iter
;
319 struct relay_index
*index
;
321 lttng::urcu::read_lock_guard read_lock
;
323 cds_lfht_for_each_entry (stream
->indexes_ht
->ht
, &iter
.iter
, index
, index_n
.node
) {
324 if (!index
->index_file
) {
328 * Partial index has its index_file: we have only
329 * received its info from the data socket.
330 * Put self-ref from index.
332 relay_index_put(index
);
336 uint64_t relay_index_find_last(struct relay_stream
*stream
)
338 struct lttng_ht_iter iter
;
339 struct relay_index
*index
;
340 uint64_t net_seq_num
= -1ULL;
342 lttng::urcu::read_lock_guard read_lock
;
344 cds_lfht_for_each_entry (stream
->indexes_ht
->ht
, &iter
.iter
, index
, index_n
.node
) {
345 if (net_seq_num
== -1ULL || index
->index_n
.key
> net_seq_num
) {
346 net_seq_num
= index
->index_n
.key
;
354 * Update the index file of an already existing relay_index.
355 * Offsets by 'removed_data_count' the offset field of an index.
357 static int relay_index_switch_file(struct relay_index
*index
,
358 struct lttng_index_file
*new_index_file
,
359 uint64_t removed_data_count
)
364 pthread_mutex_lock(&index
->lock
);
365 if (!index
->index_file
) {
366 ERR("No index_file");
371 lttng_index_file_put(index
->index_file
);
372 lttng_index_file_get(new_index_file
);
373 index
->index_file
= new_index_file
;
374 offset
= be64toh(index
->index_data
.offset
);
375 index
->index_data
.offset
= htobe64(offset
- removed_data_count
);
378 pthread_mutex_unlock(&index
->lock
);
383 * Switch the index file of all pending indexes for a stream and update the
384 * data offset by substracting the last safe position.
385 * Stream lock must be held.
387 int relay_index_switch_all_files(struct relay_stream
*stream
)
389 struct lttng_ht_iter iter
;
390 struct relay_index
*index
;
393 lttng::urcu::read_lock_guard read_lock
;
395 cds_lfht_for_each_entry (stream
->indexes_ht
->ht
, &iter
.iter
, index
, index_n
.node
) {
396 ret
= relay_index_switch_file(
397 index
, stream
->index_file
, stream
->pos_after_last_complete_data_index
);
407 * Set index data from the control port to a given index object.
409 int relay_index_set_control_data(struct relay_index
*index
,
410 const struct lttcomm_relayd_index
*data
,
411 unsigned int minor_version
)
413 /* The index on disk is encoded in big endian. */
414 ctf_packet_index index_data
{};
416 index_data
.packet_size
= htobe64(data
->packet_size
);
417 index_data
.content_size
= htobe64(data
->content_size
);
418 index_data
.timestamp_begin
= htobe64(data
->timestamp_begin
);
419 index_data
.timestamp_end
= htobe64(data
->timestamp_end
);
420 index_data
.events_discarded
= htobe64(data
->events_discarded
);
421 index_data
.stream_id
= htobe64(data
->stream_id
);
423 if (minor_version
>= 8) {
424 index
->index_data
.stream_instance_id
= htobe64(data
->stream_instance_id
);
425 index
->index_data
.packet_seq_num
= htobe64(data
->packet_seq_num
);
427 uint64_t unset_value
= -1ULL;
429 index
->index_data
.stream_instance_id
= htobe64(unset_value
);
430 index
->index_data
.packet_seq_num
= htobe64(unset_value
);
433 return relay_index_set_data(index
, &index_data
);