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}