Skip to main content

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}