hotpath_drain/lib.rs
1//! Lock-free event transport between instrumented threads and background workers.
2//!
3//! Each producing thread owns a chunked SPSC queue: it appends events into
4//! fixed-size chunks with plain stores and publishes them with a single
5//! `Release` store of the chunk's `len` - no mutex, no RMW atomic on the hot
6//! path. Chunks form a linked list, so the queue is unbounded.
7//!
8//! A single consumer per registry (the subsystem's background worker) sweeps
9//! all registered queues periodically: it `Acquire`-loads each chunk's `len`,
10//! reads the published prefix, and frees fully consumed chunks. Producer and
11//! consumer touch disjoint slot ranges by construction, so a queue is safe to
12//! drain at any moment - including queues of threads parked at shutdown,
13//! which is what guarantees a complete final report.
14//!
15//! Safety invariants:
16//! - Only the owning thread writes slots and the `len`/`next` of its tail chunk.
17//! - The consumer only reads slots below the published `len` and only frees a
18//! chunk after fully consuming it and observing its `next` pointer.
19//! - Single consumer per registry, enforced by the registry's internal mutex
20//! (locked only at thread registration and during sweeps - never on the
21//! event hot path).
22
23use std::cell::{Cell, UnsafeCell};
24use std::mem::MaybeUninit;
25use std::ptr;
26use std::sync::atomic::{AtomicBool, AtomicPtr, AtomicUsize, Ordering};
27use std::sync::{Arc, Mutex};
28
29pub const CHUNK_SIZE: usize = 64;
30
31/// Caps how many chunks one sweep drains from a single queue, so a producer
32/// outpacing the consumer cannot pin the sweep in an endless tail chase.
33/// Leftovers are picked up on the next tick.
34const MAX_CHUNKS_PER_SWEEP: usize = 1024;
35
36struct Chunk<M> {
37 slots: [UnsafeCell<MaybeUninit<M>>; CHUNK_SIZE],
38 /// Number of initialized slots; the producer publishes with `Release`.
39 len: AtomicUsize,
40 /// Set once by the producer when the chunk is full; never unset.
41 next: AtomicPtr<Chunk<M>>,
42}
43
44impl<M> Chunk<M> {
45 fn new_raw() -> *mut Chunk<M> {
46 Box::into_raw(Box::new(Chunk {
47 slots: [const { UnsafeCell::new(MaybeUninit::uninit()) }; CHUNK_SIZE],
48 len: AtomicUsize::new(0),
49 next: AtomicPtr::new(ptr::null_mut()),
50 }))
51 }
52}
53
54/// One thread's event queue: a linked list of chunks. The owning thread
55/// appends via its [`EventProducer`]; the registry's consumer drains.
56pub struct EventQueue<M> {
57 /// Oldest chunk with unconsumed events. Consumer-owned after creation.
58 head: AtomicPtr<Chunk<M>>,
59 /// Consumed slot count within `head`. Consumer-only.
60 consumed: AtomicUsize,
61 /// Set by the producer's TLS drop so the consumer can drain the remainder
62 /// and release the queue.
63 closed: AtomicBool,
64}
65
66// SAFETY: the auto impls are lost to the raw chunk pointers. Sending or
67// sharing the queue only moves `M` values across threads (hence `M: Send`);
68// concurrent access is sound because producer and consumer touch disjoint
69// slot ranges, synchronized by the `Release`/`Acquire` handoff on `len`
70// (see module-level safety invariants).
71unsafe impl<M: Send> Send for EventQueue<M> {}
72// SAFETY: same reasoning as `Send` above.
73unsafe impl<M: Send> Sync for EventQueue<M> {}
74
75impl<M> EventQueue<M> {
76 /// Drains published events into `out`, walking at most `max_chunks` full
77 /// chunks. Returns `true` when the drain reached the queue's tail, `false`
78 /// when it stopped at the cap with chunks still pending. Must only be
79 /// called by the single consumer (see registry).
80 #[cfg_attr(
81 feature = "hotpath-meta",
82 hotpath_meta::measure(impl_type = "EventQueue")
83 )]
84 fn drain_into(&self, out: &mut Vec<M>, max_chunks: usize) -> bool {
85 let mut chunk_ptr = self.head.load(Ordering::Relaxed);
86 let mut consumed = self.consumed.load(Ordering::Relaxed);
87 let mut chunks_walked = 0;
88 let mut reached_tail = true;
89 loop {
90 // SAFETY: `chunk_ptr` is either `head` (never null, freed only by
91 // this single consumer after advancing past it) or a non-null
92 // `next` observed below; the chunk stays alive until this loop
93 // frees it.
94 let chunk = unsafe { &*chunk_ptr };
95 let len = chunk.len.load(Ordering::Acquire);
96 for i in consumed..len {
97 // SAFETY: the `Acquire` load of `len` synchronizes with the
98 // producer's `Release` publish, so slots `..len` are
99 // initialized; slots below `consumed` were already read out
100 // and are never touched twice (consumer-only cursor).
101 out.push(unsafe { (*chunk.slots[i].get()).assume_init_read() });
102 }
103 consumed = len;
104 if len == CHUNK_SIZE {
105 let next = chunk.next.load(Ordering::Acquire);
106 if !next.is_null() {
107 if chunks_walked >= max_chunks {
108 reached_tail = false;
109 break;
110 }
111 // SAFETY: `chunk_ptr` came from `Box::into_raw` in
112 // `Chunk::new_raw`. The chunk is full (`len == CHUNK_SIZE`)
113 // and fully consumed, and the producer moved on to `next`,
114 // so neither side will touch it again; all `M` values were
115 // moved out above, so dropping the box frees only
116 // `MaybeUninit` storage.
117 unsafe { drop(Box::from_raw(chunk_ptr)) };
118 chunk_ptr = next;
119 consumed = 0;
120 chunks_walked += 1;
121 continue;
122 }
123 }
124 break;
125 }
126 self.head.store(chunk_ptr, Ordering::Relaxed);
127 self.consumed.store(consumed, Ordering::Relaxed);
128 reached_tail
129 }
130}
131
132impl<M> Drop for EventQueue<M> {
133 fn drop(&mut self) {
134 // Reached only after the producer is gone (closed) and the registry
135 // released its Arc, so exclusive access is guaranteed.
136 let mut chunk_ptr = *self.head.get_mut();
137 let mut consumed = *self.consumed.get_mut();
138 while !chunk_ptr.is_null() {
139 // SAFETY: `&mut self` proves exclusive access; every chunk from
140 // `head` onward is live and owned by this queue.
141 let chunk = unsafe { &mut *chunk_ptr };
142 let len = *chunk.len.get_mut();
143 for i in consumed..len {
144 // SAFETY: slots `consumed..len` were initialized by the
145 // producer and never read out, so each holds a live `M` that
146 // is dropped exactly once here.
147 unsafe { (*chunk.slots[i].get()).assume_init_drop() };
148 }
149 consumed = 0;
150 let next = *chunk.next.get_mut();
151 // SAFETY: `chunk_ptr` came from `Box::into_raw` in
152 // `Chunk::new_raw` and nothing can reference it after this drop
153 // (exclusive access via `&mut self`).
154 unsafe { drop(Box::from_raw(chunk_ptr)) };
155 chunk_ptr = next;
156 }
157 }
158}
159
160/// The owning thread's write handle, stored in a `thread_local`.
161pub struct EventProducer<M> {
162 queue: Arc<EventQueue<M>>,
163 tail: Cell<*mut Chunk<M>>,
164 len: Cell<usize>,
165}
166
167impl<M> EventProducer<M> {
168 /// Appends one event: a plain slot store plus a `Release` publish of the
169 /// new length. Allocates a fresh chunk every `CHUNK_SIZE` events.
170 #[inline]
171 #[cfg_attr(
172 feature = "hotpath-meta",
173 hotpath_meta::measure(impl_type = "EventProducer")
174 )]
175 pub fn push(&self, m: M) {
176 let tail = self.tail.get();
177 let i = self.len.get();
178 // SAFETY: `tail` is the producer-owned live tail chunk (the consumer
179 // never frees a chunk whose `next` it has not observed, and `next` is
180 // set only after this chunk is full). Slot `i` is above the published
181 // `len`, so the consumer cannot be reading it; the `Release` store of
182 // `len` publishes the write.
183 unsafe {
184 (*tail).slots[i].get().write(MaybeUninit::new(m));
185 (*tail).len.store(i + 1, Ordering::Release);
186 }
187 if i + 1 == CHUNK_SIZE {
188 let new = Chunk::new_raw();
189 // SAFETY: `tail` is still live (see above); only the producer
190 // writes `next`, and only once, when the chunk is full.
191 unsafe { (*tail).next.store(new, Ordering::Release) };
192 self.tail.set(new);
193 self.len.set(0);
194 } else {
195 self.len.set(i + 1);
196 }
197 }
198}
199
200impl<M> Drop for EventProducer<M> {
201 fn drop(&mut self) {
202 self.queue.closed.store(true, Ordering::Release);
203 }
204}
205
206/// All live per-thread queues for one event type, plus the gate that tells
207/// producers whether a worker is consuming.
208pub struct EventQueueRegistry<M> {
209 active: AtomicBool,
210 queues: Mutex<Vec<Arc<EventQueue<M>>>>,
211}
212
213impl<M: Send> Default for EventQueueRegistry<M> {
214 fn default() -> Self {
215 Self::new()
216 }
217}
218
219impl<M: Send> EventQueueRegistry<M> {
220 pub const fn new() -> Self {
221 Self {
222 active: AtomicBool::new(false),
223 queues: Mutex::new(Vec::new()),
224 }
225 }
226
227 /// Whether a worker is consuming. Producers check this before pushing so
228 /// events cannot pile up unbounded when no one will ever drain them.
229 #[inline]
230 pub fn is_active(&self) -> bool {
231 self.active.load(Ordering::Relaxed)
232 }
233
234 pub fn set_active(&self, active: bool) {
235 self.active.store(active, Ordering::Release);
236 }
237
238 /// Creates and registers a queue for the calling thread.
239 #[cfg_attr(
240 feature = "hotpath-meta",
241 hotpath_meta::measure(impl_type = "EventQueueRegistry")
242 )]
243 pub fn register(&self) -> EventProducer<M> {
244 let first = Chunk::new_raw();
245 let queue = Arc::new(EventQueue {
246 head: AtomicPtr::new(first),
247 consumed: AtomicUsize::new(0),
248 closed: AtomicBool::new(false),
249 });
250 if let Ok(mut queues) = self.queues.lock() {
251 queues.push(Arc::clone(&queue));
252 }
253 EventProducer {
254 queue,
255 tail: Cell::new(first),
256 len: Cell::new(0),
257 }
258 }
259
260 /// Drains every registered queue into `out` and releases queues whose
261 /// producer thread has exited. Holding the registry lock for the whole
262 /// sweep is what makes this the single consumer. Each queue is capped at
263 /// `MAX_CHUNKS_PER_SWEEP`; leftovers are picked up on the next tick.
264 #[cfg_attr(
265 feature = "hotpath-meta",
266 hotpath_meta::measure(impl_type = "EventQueueRegistry")
267 )]
268 pub fn sweep(&self, out: &mut Vec<M>) {
269 self.sweep_inner(out, MAX_CHUNKS_PER_SWEEP);
270 }
271
272 /// Uncapped variant for the worker's final sweep at shutdown: there is no
273 /// next tick to pick up leftovers, so every queue is drained to its tail.
274 /// Terminates because producers are deactivated (`set_active(false)`)
275 /// before shutdown is signalled, so queues can no longer grow.
276 #[cfg_attr(
277 feature = "hotpath-meta",
278 hotpath_meta::measure(impl_type = "EventQueueRegistry")
279 )]
280 pub fn drain_all(&self, out: &mut Vec<M>) {
281 self.sweep_inner(out, usize::MAX);
282 }
283
284 fn sweep_inner(&self, out: &mut Vec<M>, max_chunks: usize) {
285 if let Ok(mut queues) = self.queues.lock() {
286 queues.retain(|queue| {
287 // Read `closed` before draining: the flag is set after the
288 // producer's last push, so observing it guarantees the drain
289 // below sees every event.
290 let closed = queue.closed.load(Ordering::Acquire);
291 let reached_tail = queue.drain_into(out, max_chunks);
292 // A closed queue is released only once fully drained; a
293 // cap-truncated drain keeps it for the next sweep so its tail
294 // events are not freed unprocessed by `EventQueue::drop`.
295 !(closed && reached_tail)
296 });
297 }
298 }
299}