Skip to main content

moirai_utils/queue/
ring.rs

1//! The bounded MPMC ring itself: storage, the two cursors, and the capacity
2//! contract.
3
4use core::cell::UnsafeCell;
5use core::cmp::Ordering as CmpOrdering;
6use core::mem::MaybeUninit;
7use core::sync::atomic::{AtomicUsize, Ordering};
8
9use crate::cache::CacheAligned;
10
11#[cfg(feature = "std")]
12use std::boxed::Box;
13
14#[cfg(not(feature = "std"))]
15use alloc::boxed::Box;
16
17/// Default capacity for [`LockFreeQueue`]. Large enough to avoid backpressure
18/// under normal scheduling load while bounding memory under adversarial
19/// producer rates per the bounded-resource policy.
20const DEFAULT_QUEUE_CAPACITY: usize = 65536;
21
22/// Spins a producer makes while a dequeue reopens a slot, before it yields its
23/// time slice. Reopening is one store after the item is moved out, so a spin or
24/// two covers it; only a consumer preempted in between needs the yield.
25const SPINS_BEFORE_YIELD: u32 = 64;
26
27/// Times a producer found a claimed slot not yet reopened, so a test can
28/// reopen the slot once a producer is observed waiting on it.
29#[cfg(test)]
30pub(super) static REOPEN_WAITS: AtomicUsize = AtomicUsize::new(0);
31
32/// A single slot in the bounded MPMC queue.
33pub(super) struct Slot<T> {
34    /// Monotonic sequence number that distinguishes empty, full, and stale
35    /// states without an ABA hazard.
36    sequence: AtomicUsize,
37    /// The slot's data. It is uninitialized exactly while the sequence number
38    /// reports the slot empty, so occupancy needs no separate flag and no
39    /// `Option` discriminant: the slot costs one machine word per payload.
40    data: UnsafeCell<MaybeUninit<T>>,
41}
42
43/// A bounded, genuinely lock-free multi-producer multi-consumer queue.
44///
45/// This is an array-based MPMC queue using per-slot sequence numbers (the
46/// Vyukov algorithm). Producers and consumers operate through independent
47/// atomic head/tail cursors and never acquire a mutex or spinlock. The
48/// sequence-number protocol eliminates the ABA problem without tagged
49/// pointers or epoch-based reclamation: slots are reused in place, so no
50/// node allocation or deallocation occurs during enqueue/dequeue.
51///
52/// # Capacity
53///
54/// The queue is bounded. [`LockFreeQueue::new`] creates a queue with
55/// `DEFAULT_QUEUE_CAPACITY` usable slots. [`LockFreeQueue::with_capacity`]
56/// accepts any request of one slot or more: the ring is sized to the next power
57/// of two at least two, while [`capacity`](LockFreeQueue::capacity) keeps
58/// reporting the request, so a non-power-of-two capacity bounds the queue
59/// exactly and only the ring's unused tail is wasted. When the queue is full,
60/// [`enqueue`] retries with exponential backoff (preserving the
61/// unblocked-sender contract of the previous API), while [`try_enqueue`]
62/// returns `Err(item)` for callers that prefer explicit backpressure. Full
63/// means `capacity` items queued: a slot a dequeue has emptied but not yet
64/// reopened is waited for, never reported as fullness.
65///
66/// # Memory safety
67///
68/// Each slot's `MaybeUninit<T>` is written by the producer — the only writer,
69/// between `sequence == pos` and `sequence == pos + 1` — and moved out by the
70/// consumer, which is the only reader, between `sequence == pos + 1` and
71/// `sequence == pos + ring_len`. The sequence-number protocol therefore
72/// guarantees that only one thread ever touches a slot's payload.
73///
74/// [`enqueue`]: LockFreeQueue::enqueue
75/// [`try_enqueue`]: LockFreeQueue::try_enqueue
76// No struct-level `repr(align)`: `head`/`tail` are `CacheAligned`, so the
77// struct's alignment already equals `DESTRUCTIVE_INTERFERENCE_SIZE` and tracks
78// the per-target table in `cache.rs` instead of pinning a second literal here.
79pub struct LockFreeQueue<T> {
80    buffer: Box<[Slot<T>]>,
81    mask: usize,
82    /// Slots in the ring: a power of two, at least two, so that a slot's empty
83    /// and full generations never alias. Private because the sequence protocol
84    /// addresses slots with it; callers reason in `capacity`.
85    ring_len: usize,
86    /// Usable slots: what `try_enqueue` accepts before reporting full.
87    capacity: usize,
88    head: CacheAligned<AtomicUsize>,
89    tail: CacheAligned<AtomicUsize>,
90}
91
92// Safety: The sequence-number protocol ensures that each slot's data is
93// accessed by at most one thread at a time: a producer writes between
94// sequence == pos and sequence == pos+1; a consumer takes between
95// sequence == pos+1 and sequence == pos+ring_len. The head and tail atomics
96// are independently advanced via CAS, so no global lock is needed. T: Send
97// is sufficient because ownership of the value transfers between threads
98// through the slot, never shared concurrently.
99unsafe impl<T: Send> Send for LockFreeQueue<T> {}
100unsafe impl<T: Send> Sync for LockFreeQueue<T> {}
101
102impl<T> LockFreeQueue<T> {
103    /// Create a new queue with the default capacity.
104    pub fn new() -> Self {
105        Self::with_capacity(DEFAULT_QUEUE_CAPACITY)
106    }
107
108    /// Create a new queue holding up to `capacity` items.
109    ///
110    /// The ring is the next power of two at least two, which the sequence
111    /// protocol needs to tell a slot's empty generation from its full one; a
112    /// one-slot request therefore gets a two-slot ring and still bounds the
113    /// queue at one item.
114    #[track_caller]
115    pub fn with_capacity(capacity: usize) -> Self {
116        let capacity = capacity.max(1);
117        let ring_len = capacity.next_power_of_two().max(2);
118
119        #[cfg(feature = "std")]
120        let buffer: Box<[Slot<T>]> = (0..ring_len)
121            .map(|i| Slot {
122                sequence: AtomicUsize::new(i),
123                data: UnsafeCell::new(MaybeUninit::uninit()),
124            })
125            .collect::<std::vec::Vec<_>>()
126            .into_boxed_slice();
127
128        #[cfg(not(feature = "std"))]
129        let buffer: Box<[Slot<T>]> = (0..ring_len)
130            .map(|i| Slot {
131                sequence: AtomicUsize::new(i),
132                data: UnsafeCell::new(MaybeUninit::uninit()),
133            })
134            .collect::<alloc::vec::Vec<_>>()
135            .into_boxed_slice();
136
137        Self {
138            buffer,
139            mask: ring_len - 1,
140            ring_len,
141            capacity,
142            head: CacheAligned::new(AtomicUsize::new(0)),
143            tail: CacheAligned::new(AtomicUsize::new(0)),
144        }
145    }
146
147    /// Try to enqueue an item without waiting for space.
148    ///
149    /// Returns `Ok(())` if the item was enqueued, or `Err(item)` if the queue
150    /// holds [`capacity`](LockFreeQueue::capacity) items. It takes no lock.
151    /// When the queue has room but the item's slot still belongs to a dequeue
152    /// that has moved its item out and not yet reopened the slot, it waits for
153    /// that reopening instead of reporting a full queue; the wait is one store
154    /// unless the dequeuing thread was preempted.
155    #[inline]
156    pub fn try_enqueue(&self, item: T) -> Result<(), T> {
157        let mut pos = self.tail.load(Ordering::Relaxed);
158        let mut waits = 0;
159        loop {
160            let slot = &self.buffer[pos & self.mask];
161            let seq = slot.sequence.load(Ordering::Acquire);
162            #[allow(clippy::cast_possible_wrap)]
163            let diff = seq.wrapping_sub(pos) as isize;
164
165            match diff.cmp(&0) {
166                CmpOrdering::Equal => {
167                    // The ring can be larger than the requested capacity, so
168                    // fullness is the request, not the ring size.
169                    if self.holds_capacity(pos) {
170                        return Err(item);
171                    }
172
173                    match self.tail.compare_exchange_weak(
174                        pos,
175                        pos.wrapping_add(1),
176                        Ordering::Relaxed,
177                        Ordering::Relaxed,
178                    ) {
179                        Ok(_) => {
180                            // SAFETY: winning the tail CAS grants exclusive
181                            // right to fill this slot's sequence generation; its
182                            // payload cell is uninitialized (fresh or drained)
183                            // until this write publishes it.
184                            unsafe {
185                                (*slot.data.get()).write(item);
186                            }
187                            slot.sequence.store(pos.wrapping_add(1), Ordering::Release);
188                            return Ok(());
189                        }
190                        Err(actual) => pos = actual,
191                    }
192                }
193                // The slot still holds the generation written one lap ago.
194                // Below capacity, a dequeue has claimed that item and not yet
195                // reopened the slot: wait for it rather than report a fullness
196                // that does not exist, as crossbeam's `ArrayQueue::push` does.
197                CmpOrdering::Less => {
198                    if self.holds_capacity(pos) {
199                        return Err(item);
200                    }
201                    wait_for_reopen(&mut waits);
202                    pos = self.tail.load(Ordering::Relaxed);
203                }
204                // Another producer advanced tail before us: reload and retry.
205                CmpOrdering::Greater => pos = self.tail.load(Ordering::Relaxed),
206            }
207        }
208    }
209
210    /// Whether the queue holds `capacity` items once the tail reaches `pos`.
211    ///
212    /// `pos` may be stale: once other producers fill that position and
213    /// consumers drain it, the head passes `pos` and the wrapped difference is
214    /// huge. No real occupancy exceeds the ring, so a difference past
215    /// `ring_len` is a stale position the caller retries, not a full queue.
216    #[inline]
217    pub(super) fn holds_capacity(&self, pos: usize) -> bool {
218        let queued = pos.wrapping_sub(self.head.load(Ordering::Acquire));
219        (self.capacity..=self.ring_len).contains(&queued)
220    }
221
222    /// Enqueue an item, retrying with exponential backoff if the queue is full.
223    ///
224    /// This preserves the unblocked-sender contract of the previous API: the
225    /// call always eventually succeeds (assuming consumers make progress).
226    /// The backoff path uses `core::hint::spin_loop` and, on std targets,
227    /// `std::thread::yield_now` after heavy contention, but never acquires a
228    /// global lock, so multiple producers can enqueue concurrently.
229    #[inline]
230    pub fn enqueue(&self, item: T) {
231        let mut backoff: usize = 1;
232        let mut item = Some(item);
233        loop {
234            match self.try_enqueue(item.take().expect("invariant: item present")) {
235                Ok(()) => return,
236                Err(returned) => {
237                    item = Some(returned);
238                    for _ in 0..backoff {
239                        core::hint::spin_loop();
240                    }
241                    if backoff < 64 {
242                        backoff = backoff.saturating_mul(2);
243                    } else {
244                        #[cfg(feature = "std")]
245                        {
246                            std::thread::yield_now();
247                        }
248                        backoff = 1;
249                    }
250                }
251            }
252        }
253    }
254
255    /// Try to dequeue an item from the front of the queue.
256    /// Returns `None` if the queue is empty.
257    ///
258    /// This is the lock-free fast path: no spinlock, no mutex.
259    #[inline]
260    pub fn try_dequeue(&self) -> Option<T> {
261        let (pos, item) = self.claim_front()?;
262        // SAFETY: `pos` was claimed just above and is reopened once.
263        unsafe { self.reopen(pos) };
264        Some(item)
265    }
266
267    /// Claims the front item: advances the head past it and moves it out,
268    /// leaving its slot closed to producers until [`Self::reopen`].
269    #[inline]
270    pub(super) fn claim_front(&self) -> Option<(usize, T)> {
271        let mut pos = self.head.load(Ordering::Relaxed);
272        loop {
273            let slot = &self.buffer[pos & self.mask];
274            let seq = slot.sequence.load(Ordering::Acquire);
275            #[allow(clippy::cast_possible_wrap)]
276            let diff = seq.wrapping_sub(pos.wrapping_add(1)) as isize;
277
278            match diff.cmp(&0) {
279                CmpOrdering::Equal => {
280                    // Slot has data: try to claim it by advancing head.
281                    match self.head.compare_exchange_weak(
282                        pos,
283                        pos.wrapping_add(1),
284                        Ordering::Relaxed,
285                        Ordering::Relaxed,
286                    ) {
287                        Ok(_) => {
288                            // SAFETY: winning the head CAS means no other
289                            // consumer can claim this slot, and the sequence
290                            // == pos+1 invariant means the producer finished
291                            // writing; no producer writes it again until
292                            // `reopen` opens the next generation.
293                            let item = unsafe { (*slot.data.get()).assume_init_read() };
294                            return Some((pos, item));
295                        }
296                        Err(actual) => pos = actual,
297                    }
298                }
299                // Queue is empty: sequence has not advanced past pos+1.
300                CmpOrdering::Less => return None,
301                // Another consumer advanced head before us: reload and retry.
302                CmpOrdering::Greater => pos = self.head.load(Ordering::Relaxed),
303            }
304        }
305    }
306
307    /// Reopens the slot of the item claimed at `pos` for the next lap's
308    /// producer.
309    ///
310    /// # Safety
311    /// `pos` must come from [`Self::claim_front`] on this queue and be
312    /// reopened exactly once: reopening any other slot lets a producer write
313    /// a payload a consumer may still be reading.
314    #[inline]
315    pub(super) unsafe fn reopen(&self, pos: usize) {
316        self.buffer[pos & self.mask]
317            .sequence
318            .store(pos.wrapping_add(self.ring_len), Ordering::Release);
319    }
320
321    /// Check if the queue is empty.
322    ///
323    /// This is a best-effort check: the queue may have items added or removed
324    /// between this call and the next operation. It is safe to call
325    /// concurrently with enqueue/dequeue.
326    pub fn is_empty(&self) -> bool {
327        self.tail.load(Ordering::Acquire) == self.head.load(Ordering::Acquire)
328    }
329
330    /// Check if the queue holds `capacity` items.
331    ///
332    /// Best-effort in the same sense as [`is_empty`](LockFreeQueue::is_empty).
333    pub fn is_full(&self) -> bool {
334        self.tail
335            .load(Ordering::Acquire)
336            .wrapping_sub(self.head.load(Ordering::Acquire))
337            >= self.capacity
338    }
339
340    /// Number of items the queue accepts before `try_enqueue` reports full.
341    pub const fn capacity(&self) -> usize {
342        self.capacity
343    }
344
345    /// Number of items currently queued.
346    ///
347    /// Best-effort in the same sense as [`is_empty`](LockFreeQueue::is_empty):
348    /// it is the cursor difference, so a push that has reserved its position but
349    /// not yet published its slot is already counted, and the value can change
350    /// under the caller immediately after the read.
351    pub fn len(&self) -> usize {
352        self.tail
353            .load(Ordering::Acquire)
354            .wrapping_sub(self.head.load(Ordering::Acquire))
355    }
356}
357
358/// Pauses a producer waiting for a dequeue to reopen a slot: spin briefly, then
359/// yield where the platform can, so a preempted consumer gets to run.
360#[inline]
361fn wait_for_reopen(waits: &mut u32) {
362    #[cfg(test)]
363    REOPEN_WAITS.fetch_add(1, Ordering::Relaxed);
364    if *waits < SPINS_BEFORE_YIELD {
365        *waits += 1;
366        core::hint::spin_loop();
367    } else {
368        #[cfg(feature = "std")]
369        std::thread::yield_now();
370        #[cfg(not(feature = "std"))]
371        core::hint::spin_loop();
372    }
373}
374
375impl<T> Default for LockFreeQueue<T> {
376    fn default() -> Self {
377        Self::new()
378    }
379}
380
381impl<T> Drop for LockFreeQueue<T> {
382    fn drop(&mut self) {
383        // Exclusive access: walk the published range and drop what nobody took,
384        // reading the sequence numbers without atomics.
385        let head = *self.head.0.get_mut();
386        let tail = *self.tail.0.get_mut();
387
388        for pos in head..tail {
389            let slot = &mut self.buffer[pos & self.mask];
390            if *slot.sequence.get_mut() == pos.wrapping_add(1) {
391                // SAFETY: the published sequence marks the payload
392                // initialized, and each position is visited once.
393                unsafe {
394                    (*slot.data.get()).assume_init_drop();
395                }
396            }
397        }
398    }
399}