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}