Skip to main content

moirai_core/communication/
ring_buffer.rs

1use moirai_utils::cache::CacheAligned;
2use std::cell::{Cell, UnsafeCell};
3use std::mem::MaybeUninit;
4use std::sync::atomic::{AtomicUsize, Ordering};
5
6/// Zero-copy ring buffer for high-throughput streaming
7///
8/// # Safety
9///
10/// This structure uses `MaybeUninit` for zero-copy performance:
11/// - Values are written with `write()` before incrementing `producer_seq`
12/// - The `assume_init_read()` in `try_consume()` is safe because we check
13///   that `producer_seq` > current, ensuring data was written
14///
15/// # Why this ring also backs the SPSC channel
16///
17/// [`SpscRing`](crate::channel::SpscRing) and `channel::spsc`'s halves are backed
18/// by this type. There is one Lamport protocol here — a masked slot array, a
19/// `Relaxed` load of the owner's own cursor, an `Acquire` load of the peer's, a
20/// `Release` store back, and the element written in between — reached through two
21/// access disciplines:
22///
23/// - This type is public, exposes the uncached entry points, and stays `!Sync`
24///   (see the `Send` impl below). Its `&self` methods mutate through
25///   `UnsafeCell`, so a shared `&RingBuffer` would let two safe threads race one
26///   end of the ring.
27/// - `channel::spsc::SpscChannel` is crate-private and *is* `Sync`, because its
28///   `Arc` must be `Send` to back `'static` halves. Its safety argument is the
29///   non-`Clone` halves plus crate-private reach (ADR-024), which a public type
30///   cannot invoke. It adds the cached-index layer, a `closed` flag, and the
31///   spin-then-yield policy.
32///
33/// The *bound* is not shared: granting `Sync` here would be unsound for
34/// downstream users, and a public type cannot invoke the channel's argument. The
35/// *code* is shared instead. The cached primitives live on this type as
36/// `pub(crate)` methods, so the discipline stays on the wrapper that can enforce
37/// it while the storage, the cursors, and the publication algebra exist exactly
38/// once (ADR-016 item 3).
39pub struct RingBuffer<T> {
40    /// Buffer storage
41    buffer: Box<[UnsafeCell<MaybeUninit<T>>]>,
42    /// Capacity mask for fast modulo
43    mask: usize,
44    /// Producer sequence number
45    producer_seq: CacheAligned<AtomicUsize>,
46    /// Consumer sequence number
47    consumer_seq: CacheAligned<AtomicUsize>,
48}
49
50// SAFETY: the ring owns its `T` values inside `UnsafeCell<MaybeUninit<T>>`, so it
51// may move between threads exactly when `T: Send`. It is deliberately NOT `Sync`:
52// concurrent shared access is only sound under the single-producer/single-consumer
53// discipline (producer touches `producer_seq` + tail slots, consumer touches
54// `consumer_seq` + head slots, never the same slot), which is enforced by the
55// non-`Clone` `HybridSender`/`HybridReceiver` halves rather than by the type
56// system here. Granting `Sync` would permit two producers (or two consumers) to
57// race the same end, so it is intentionally withheld.
58unsafe impl<T: Send> Send for RingBuffer<T> {}
59
60impl<T> RingBuffer<T> {
61    /// Create a new ring buffer with given capacity
62    pub fn new(capacity: usize) -> Self {
63        let capacity = capacity.next_power_of_two();
64        let buffer = (0..capacity)
65            .map(|_| UnsafeCell::new(MaybeUninit::uninit()))
66            .collect::<Vec<_>>()
67            .into_boxed_slice();
68
69        Self {
70            buffer,
71            mask: capacity - 1,
72            producer_seq: CacheAligned::new(AtomicUsize::new(0)),
73            consumer_seq: CacheAligned::new(AtomicUsize::new(0)),
74        }
75    }
76
77    /// Try to produce a value, handing it back when the ring is full.
78    pub fn try_produce(&self, value: T) -> Result<(), T> {
79        let current = self.producer_relaxed();
80        let consumer = self.consumer_acquire();
81
82        // Check if full
83        if current.wrapping_sub(consumer) >= self.buffer.len() {
84            return Err(value);
85        }
86
87        // SAFETY: the capacity check above keeps `current` outside the consumer
88        // window, and the withheld `Sync` makes this thread the ring's sole
89        // producer.
90        unsafe { self.produce_at(current, value) };
91        Ok(())
92    }
93
94    /// Try to consume a value
95    pub fn try_consume(&self) -> Option<T> {
96        let current = self.consumer_relaxed();
97        let producer = self.producer_acquire();
98
99        if current == producer {
100            return None;
101        }
102
103        // SAFETY: the acquire load above proves the producer published this slot,
104        // and the withheld `Sync` makes this thread the ring's sole consumer.
105        Some(unsafe { self.consume_at(current) })
106    }
107
108    /// Get the capacity of the ring buffer
109    pub fn capacity(&self) -> usize {
110        self.buffer.len()
111    }
112
113    /// Check if the ring buffer is empty
114    pub fn is_empty(&self) -> bool {
115        self.consumer_acquire() == self.producer_acquire()
116    }
117
118    /// Check if the ring buffer is full
119    pub fn is_full(&self) -> bool {
120        self.producer_acquire()
121            .wrapping_sub(self.consumer_acquire())
122            >= self.buffer.len()
123    }
124
125    /// Get the number of items currently in the buffer
126    pub fn len(&self) -> usize {
127        self.producer_acquire()
128            .wrapping_sub(self.consumer_acquire())
129    }
130}
131
132/// The protocol's atoms, shared with the SPSC channel.
133///
134/// These are the primitives `channel::spsc::SpscChannel` composes: that wrapper
135/// owns the *policy* (the cached-index layer, the `closed` flag, the
136/// spin-then-yield schedule) and drives these for the *mechanism*. They are
137/// `pub(crate)` rather than public because they take `&self` and mutate through
138/// the ring's cells, so reaching them from outside the crate would let any number
139/// of threads drive one end — the discipline `channel::spsc`'s non-`Clone` halves
140/// enforce and a public type cannot.
141impl<T> RingBuffer<T> {
142    /// The producer and consumer cursors, in that order, read without
143    /// synchronization.
144    ///
145    /// `Relaxed` is correct only because every caller holds the ring
146    /// exclusively: [`SpscRing`](crate::channel::SpscRing)'s methods take `&self`
147    /// or `&mut self`, and its halves borrow it, so no half can exist — and
148    /// therefore no other thread can be advancing either cursor — while this
149    /// runs. Reading these from a live half needs the acquire loads the send and
150    /// receive paths use.
151    pub(crate) fn indices(&self) -> (usize, usize) {
152        (
153            self.producer_seq.0.load(Ordering::Relaxed),
154            self.consumer_seq.0.load(Ordering::Relaxed),
155        )
156    }
157
158    /// The producer's own cursor, `Relaxed`: only the producing half advances it,
159    /// and it only moves forward.
160    pub(crate) fn producer_relaxed(&self) -> usize {
161        self.producer_seq.0.load(Ordering::Relaxed)
162    }
163
164    /// The producer's cursor, `Acquire`: the `Release` store that publishes
165    /// closure must be visible before the emptiness re-check that reads it.
166    pub(crate) fn producer_acquire(&self) -> usize {
167        self.producer_seq.0.load(Ordering::Acquire)
168    }
169
170    /// The consumer's own cursor, `Relaxed`, for the same reason as
171    /// [`Self::producer_relaxed`].
172    pub(crate) fn consumer_relaxed(&self) -> usize {
173        self.consumer_seq.0.load(Ordering::Relaxed)
174    }
175
176    /// The consumer's cursor, `Acquire`: the `Release` store that publishes an
177    /// element must be visible before that element is read.
178    pub(crate) fn consumer_acquire(&self) -> usize {
179        self.consumer_seq.0.load(Ordering::Acquire)
180    }
181
182    /// Room for one more value, consulting `cached_consumer` before the
183    /// consumer's real cursor.
184    ///
185    /// The cached cursor is always at or behind the true one, because only the
186    /// consumer advances it and it only moves forward. A stale value therefore
187    /// makes the ring look *fuller* than it is, never emptier, so this may take
188    /// the slow path unnecessarily but can never report room that does not exist.
189    /// That one-sidedness is what makes the cache sound.
190    pub(crate) fn has_room(&self, producer: usize, cached_consumer: &Cell<usize>) -> bool {
191        if producer.wrapping_sub(cached_consumer.get()) < self.buffer.len() {
192            return true;
193        }
194        // The cache says full; consult the consumer and try once more. This is
195        // the only load that touches the consumer's cache line.
196        let consumer = self.consumer_seq.0.load(Ordering::Acquire);
197        cached_consumer.set(consumer);
198        producer.wrapping_sub(consumer) < self.buffer.len()
199    }
200
201    /// A value is available, consulting `cached_producer` before the producer's
202    /// real cursor. Mirrors [`Self::has_room`]: a stale cache understates what is
203    /// queued, so it can cost an extra load but never invent an element.
204    pub(crate) fn has_value(&self, consumer: usize, cached_producer: &Cell<usize>) -> bool {
205        if consumer != cached_producer.get() {
206            return true;
207        }
208        let producer = self.producer_seq.0.load(Ordering::Acquire);
209        cached_producer.set(producer);
210        consumer != producer
211    }
212
213    /// Write `value` into the slot at `producer` and publish the cursor.
214    ///
215    /// # Safety
216    ///
217    /// The caller must have established room for `producer` — through
218    /// [`Self::has_room`], or the uncached full check [`Self::try_produce`]
219    /// performs — and must be the sole producer: this writes a slot the consumer
220    /// may reach as soon as the release store below lands.
221    pub(crate) unsafe fn produce_at(&self, producer: usize, value: T) {
222        // SAFETY: `producer` is past the consumer's cursor, so this slot is not
223        // one the consumer may read until the release store publishes it, and
224        // only the producing half writes slots.
225        unsafe {
226            let slot = &mut *self.buffer[producer & self.mask].get();
227            slot.write(value);
228        }
229        self.producer_seq
230            .0
231            .store(producer.wrapping_add(1), Ordering::Release);
232    }
233
234    /// Take the value in the slot at `consumer` and publish the cursor.
235    ///
236    /// # Safety
237    ///
238    /// The caller must have established that `consumer` is behind the published
239    /// producer cursor — through [`Self::has_value`], or the uncached emptiness
240    /// check [`Self::try_consume`] performs — and must be the sole consumer: the
241    /// slot is read once here and never again.
242    pub(crate) unsafe fn consume_at(&self, consumer: usize) -> T {
243        // SAFETY: `consumer` is behind the published producer cursor, so this
244        // slot was written and released by the producer. It has not been read
245        // before — the cursor advances once per value, and only the consuming
246        // half advances it.
247        let value = unsafe {
248            let slot = &*self.buffer[consumer & self.mask].get();
249            slot.assume_init_read()
250        };
251        self.consumer_seq
252            .0
253            .store(consumer.wrapping_add(1), Ordering::Release);
254        value
255    }
256}
257
258impl<T> Drop for RingBuffer<T> {
259    fn drop(&mut self) {
260        let consumer = *self.consumer_seq.0.get_mut();
261        let producer = *self.producer_seq.0.get_mut();
262        let len = producer.wrapping_sub(consumer);
263        for i in 0..len {
264            let idx = (consumer.wrapping_add(i)) & self.mask;
265            // SAFETY: exclusive `&mut self` in drop; every live index in
266            // `consumer..producer` was written by produce and not yet read,
267            // so dropping it here discharges each value exactly once.
268            unsafe {
269                let slot = &mut *self.buffer[idx].get();
270                slot.assume_init_drop();
271            }
272        }
273    }
274}
275
276#[cfg(test)]
277mod tests {
278    use super::*;
279
280    #[test]
281    fn test_wrapping_drop_correctness() {
282        use std::sync::atomic::{AtomicUsize, Ordering};
283
284        static DROP_COUNT: AtomicUsize = AtomicUsize::new(0);
285        struct TrackDrop;
286        impl Drop for TrackDrop {
287            fn drop(&mut self) {
288                DROP_COUNT.fetch_add(1, Ordering::SeqCst);
289            }
290        }
291
292        {
293            let mut rb = RingBuffer::<TrackDrop>::new(4);
294            let mask = rb.mask;
295            unsafe {
296                let slot1 = &mut *rb.buffer[(usize::MAX - 1) & mask].get();
297                slot1.write(TrackDrop);
298                let slot2 = &mut *rb.buffer[usize::MAX & mask].get();
299                slot2.write(TrackDrop);
300                let slot3 = &mut *rb.buffer[0].get();
301                slot3.write(TrackDrop);
302            }
303
304            *rb.consumer_seq.0.get_mut() = usize::MAX - 1;
305            *rb.producer_seq.0.get_mut() = 1;
306        }
307
308        assert_eq!(DROP_COUNT.load(Ordering::SeqCst), 3);
309    }
310}