Skip to main content

moirai_core/channel/mpmc/
channel.rs

1#![expect(
2    clippy::unwrap_used,
3    reason = "ratchet MOIRAI-UNWRAP-1: pre-existing debt"
4)]
5
6use super::recv::MpmcReceiver;
7use super::send::MpmcSender;
8use super::{MPMC_BLOCK_SPINS, MpmcState, block::backoff_step};
9use crate::channel::CHANNEL_STORE_LOAD_ORDER;
10use crate::channel::error::{Channel, ChannelError, Result};
11use moirai_utils::queue::LockFreeQueue;
12use std::collections::VecDeque;
13use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering, fence};
14use std::sync::{Arc, Condvar, Mutex};
15
16mod roles;
17
18/// Multi-Producer Multi-Consumer channel with bounded capacity
19/// Uses mutex-based implementation for simplicity and correctness.
20///
21/// The shared state lives *directly* in this struct rather than behind a field
22/// per `Arc`: [`MpmcSender`]/[`MpmcReceiver`] already share one handle through
23/// `Arc<MpmcChannel<T>>`, so a channel costs one allocation plus the bounded
24/// ring's slot array. Every operation reaches the mutex, condvars, both waiter
25/// counters and the ring through that single handle, with no second
26/// indirection on the send/receive paths.
27pub struct MpmcChannel<T> {
28    pub(super) state: (Mutex<MpmcState<T>>, Condvar, Condvar),
29    pub(super) bounded: Option<LockFreeQueue<T>>,
30    pub(super) closed: AtomicBool,
31    pub(super) sender_waiter_count: AtomicUsize,
32    pub(super) receiver_waiter_count: AtomicUsize,
33}
34
35/// Slots preallocated for an unbounded channel's mutex-guarded deque.
36///
37/// The unbounded path is the only one that stores items in `MpmcState::queue`,
38/// and it has no capacity to size that allocation from. One `VecDeque` growth
39/// step reallocates and moves every element while the channel mutex is held,
40/// so the first few sends are paid for up front; 16 covers the short-burst
41/// case without committing memory a mostly-idle channel never uses.
42const UNBOUNDED_PREALLOCATED_SLOTS: usize = 16;
43
44impl<T> MpmcChannel<T> {
45    /// Create a new MPMC channel with optional capacity
46    pub fn new(capacity: Option<usize>) -> Self {
47        let state = MpmcState {
48            // Deliberately asymmetric, and not the inversion it looks like: a
49            // bounded channel stores nothing here. Every `Channel` method
50            // dispatches on `self.bounded` first, and `bounded` is `Some`
51            // exactly when `capacity` is, so for a bounded channel this deque
52            // stays empty for life and preallocating `capacity` slots for it
53            // would reserve a second copy of the ring that is never written.
54            // Items live in the lock-free `LockFreeQueue`, which allocates
55            // its ring in `LockFreeQueue::with_capacity` below.
56            queue: if capacity.is_some() {
57                VecDeque::new()
58            } else {
59                VecDeque::with_capacity(UNBOUNDED_PREALLOCATED_SLOTS)
60            },
61            capacity,
62            closed: false,
63            sender_count: 0,
64            receiver_count: 0,
65        };
66
67        let bounded = capacity.map(LockFreeQueue::with_capacity);
68
69        Self {
70            state: (Mutex::new(state), Condvar::new(), Condvar::new()),
71            bounded,
72            closed: AtomicBool::new(false),
73            sender_waiter_count: AtomicUsize::new(0),
74            receiver_waiter_count: AtomicUsize::new(0),
75        }
76    }
77
78    /// Create an unbounded channel
79    pub fn unbounded() -> Self {
80        Self::new(None)
81    }
82
83    /// Create a bounded channel with given capacity
84    pub fn bounded(capacity: usize) -> Self {
85        Self::new(Some(capacity))
86    }
87
88    /// Create a channel pair for ergonomic usage
89    pub fn channel(capacity: Option<usize>) -> (MpmcSender<T>, MpmcReceiver<T>) {
90        let channel = Arc::new(Self::new(capacity));
91        let (mutex, _, _) = &channel.state;
92
93        {
94            let mut state = mutex.lock().unwrap();
95            state.sender_count = 1;
96            state.receiver_count = 1;
97        }
98
99        (
100            MpmcSender {
101                channel: channel.clone(),
102            },
103            MpmcReceiver { channel },
104        )
105    }
106
107    fn send_bounded(&self, queue: &LockFreeQueue<T>, mut value: T) -> Result<()>
108    where
109        T: Send,
110    {
111        let mut spin_count = 0;
112
113        loop {
114            if self.closed.load(Ordering::Acquire) {
115                return Err(ChannelError::Closed);
116            }
117
118            match queue.try_enqueue(value) {
119                Ok(()) => {
120                    self.wake_receiver_after_push();
121                    return Ok(());
122                }
123                Err(returned) => {
124                    value = returned;
125                }
126            }
127
128            if spin_count < MPMC_BLOCK_SPINS {
129                // The ring paths spend a round without releasing anything: they
130                // hold no lock, so unlike the mutex paths there is nothing to
131                // re-acquire afterwards.
132                backoff_step(&mut spin_count);
133                continue;
134            }
135
136            // Fallback to condvar wait to prevent CPU contention and busy-looping
137            let (mutex, not_full, _) = &self.state;
138            let mut guard = mutex.lock().unwrap();
139
140            if self.closed.load(Ordering::Acquire) || guard.closed {
141                return Err(ChannelError::Closed);
142            }
143
144            // Register *before* the re-check, never after: the counter must be
145            // visible to any receiver that frees a slot from here on, or that
146            // receiver reads zero and skips the notify while this thread goes
147            // on to park. SeqCst, load-bearing: this is the waiter half of the
148            // Dekker pair described above. The explicit fence is the Store→Load
149            // barrier separating the registration from the `try_enqueue` below:
150            // the C++/Rust memory model orders an SC fence against the
151            // notifier's SC fence, but gives a bare SC read-modify-write no
152            // such force over a later non-SC load. x86-64 supplies it in
153            // hardware; the fence makes the guarantee portable, and it costs
154            // nothing on the hot path because this branch runs only after the
155            // spin budget is spent.
156            self.sender_waiter_count
157                .fetch_add(1, CHANNEL_STORE_LOAD_ORDER);
158            fence(CHANNEL_STORE_LOAD_ORDER);
159
160            match queue.try_enqueue(value) {
161                Ok(()) => {
162                    // No fence here: this push holds the channel mutex, and a
163                    // receiver registers (and re-checks the queue) under that
164                    // same mutex, so the mutex orders the two sides.
165                    //
166                    // Relaxed: deregistration. No happens-before edge is
167                    // needed — the counter only ever gates a `notify_one`, so
168                    // a receiver still reading the pre-decrement value takes
169                    // the mutex and signals a condvar nobody waits on. A
170                    // spurious notify is free; a missed one is a hang, and
171                    // only the increment above can be missed.
172                    self.sender_waiter_count.fetch_sub(1, Ordering::Relaxed);
173                    drop(guard);
174                    // Relaxed: unlike the fast path above, this push happened
175                    // while holding the channel mutex, and a receiver
176                    // registers (and re-checks the queue) while holding that
177                    // same mutex. Either it registered before this thread took
178                    // the lock — then its `fetch_add` happens-before this load
179                    // through the mutex and the load observes it — or it takes
180                    // the lock afterwards, and its own re-check finds the item
181                    // this thread just pushed. Neither branch parks.
182                    if self.receiver_waiter_count.load(Ordering::Relaxed) > 0 {
183                        let (_, _, not_empty) = &self.state;
184                        not_empty.notify_one();
185                    }
186                    return Ok(());
187                }
188                Err(returned) => {
189                    value = returned;
190                }
191            }
192
193            guard = not_full.wait(guard).unwrap();
194            // Relaxed: deregistration, as above.
195            self.sender_waiter_count.fetch_sub(1, Ordering::Relaxed);
196        }
197    }
198
199    /// Wake one parked receiver after a lock-free push.
200    ///
201    /// SeqCst, load-bearing: this is the notifier half of a store-buffer
202    /// (Dekker) pair with `recv_bounded`'s registration. Here the queue write
203    /// precedes the counter read; there the counter write precedes the queue
204    /// read. If either side could reorder Store→Load, this side reads "no
205    /// waiters" while that side reads "still empty", and a receiver parks
206    /// forever on an item that is already queued. Acquire is insufficient — it
207    /// orders Load→Load and Load→Store, never Store→Load. The queue is
208    /// lock-free, so the channel mutex orders neither side. The waiter half gets
209    /// that barrier from its `SeqCst` RMW; this half does not — `try_enqueue`
210    /// ends in a plain release store and a `SeqCst` load is an ordinary `mov` on
211    /// x86-64 — so the fence is explicit.
212    ///
213    /// The fence and the counter read are taken on every push, never gated on
214    /// the ring having been empty: occupancy observed before the slot is
215    /// published says nothing about occupancy when a receiver re-checks. A
216    /// receiver can drain the items ahead of this push in that window, find this
217    /// slot still unpublished, and park; skipping the notify then strands it,
218    /// and a sender that fills the ring behind it parks on `not_full` too.
219    /// `tests/loom_mpmc_waiter.rs` enumerates both interleavings
220    /// (`notifier_without_the_store_load_barrier_loses_the_wakeup`,
221    /// `occupancy_read_before_publish_loses_the_wakeup`).
222    fn wake_receiver_after_push(&self) {
223        fence(CHANNEL_STORE_LOAD_ORDER);
224        if self.receiver_waiter_count.load(CHANNEL_STORE_LOAD_ORDER) > 0 {
225            let (mutex, _, not_empty) = &self.state;
226            let _guard = mutex.lock().unwrap();
227            not_empty.notify_one();
228        }
229    }
230
231    /// Wake one parked sender after a successful pop.
232    ///
233    /// The mirror of [`Self::wake_receiver_after_push`], and the notifier half
234    /// of the Dekker pair with `send_bounded`'s registration.
235    fn wake_sender_after_pop(&self) {
236        fence(CHANNEL_STORE_LOAD_ORDER);
237        if self.sender_waiter_count.load(CHANNEL_STORE_LOAD_ORDER) > 0 {
238            let (mutex, not_full, _) = &self.state;
239            let _guard = mutex.lock().unwrap();
240            not_full.notify_one();
241        }
242    }
243
244    fn recv_bounded(&self, queue: &LockFreeQueue<T>) -> Result<T>
245    where
246        T: Send,
247    {
248        let mut spin_count = 0;
249
250        loop {
251            if let Some(value) = queue.try_dequeue() {
252                self.wake_sender_after_pop();
253                return Ok(value);
254            }
255
256            if self.closed.load(Ordering::Acquire) {
257                if queue.is_empty() {
258                    return Err(ChannelError::Closed);
259                }
260                std::hint::spin_loop();
261                continue;
262            }
263
264            if spin_count < MPMC_BLOCK_SPINS {
265                // The ring paths spend a round without releasing anything: they
266                // hold no lock, so unlike the mutex paths there is nothing to
267                // re-acquire afterwards.
268                backoff_step(&mut spin_count);
269                continue;
270            }
271
272            // Fallback to condvar wait to prevent CPU contention and busy-looping
273            let (mutex, _, not_empty) = &self.state;
274            let mut guard = mutex.lock().unwrap();
275
276            // Register *before* the re-check below. Previously the order was
277            // inverted here (re-check, then register) while `send_bounded`
278            // registered first, and the asymmetry was a lost wakeup: a
279            // producer could push and read `receiver_waiter_count == 0`
280            // between this thread's failed `try_dequeue` and its `fetch_add`,
281            // skip the notify, and leave this thread parked on a queue that
282            // already holds its item. Registering first makes the producer's
283            // counter read and this thread's queue read a Dekker pair that
284            // `SeqCst` closes.
285            //
286            // SeqCst, load-bearing: waiter half of the pair. The fence is the
287            // Store→Load barrier before the `try_dequeue` that follows (see
288            // `send_bounded` for why the registration alone is not one).
289            self.receiver_waiter_count
290                .fetch_add(1, CHANNEL_STORE_LOAD_ORDER);
291            fence(CHANNEL_STORE_LOAD_ORDER);
292
293            if let Some(value) = queue.try_dequeue() {
294                // Relaxed: deregistration (see `send_bounded`).
295                self.receiver_waiter_count.fetch_sub(1, Ordering::Relaxed);
296                // Relaxed: this pop and the matching sender registration both
297                // happen under the channel mutex, which supplies the edge —
298                // a sender that registered earlier is visible through the
299                // lock, and one that registers later re-checks the queue slot
300                // this pop just freed.
301                if self.sender_waiter_count.load(Ordering::Relaxed) > 0 {
302                    let (_, not_full, _) = &self.state;
303                    not_full.notify_one();
304                }
305                drop(guard);
306                return Ok(value);
307            }
308
309            if self.closed.load(Ordering::Acquire) || guard.closed {
310                // Relaxed: deregistration on the error exit.
311                self.receiver_waiter_count.fetch_sub(1, Ordering::Relaxed);
312                return Err(ChannelError::Closed);
313            }
314
315            guard = not_empty.wait(guard).unwrap();
316            // Relaxed: deregistration (see `send_bounded`).
317            self.receiver_waiter_count.fetch_sub(1, Ordering::Relaxed);
318        }
319    }
320}
321
322impl<T: Send> Channel<T> for MpmcChannel<T> {
323    fn send(&self, value: T) -> Result<()> {
324        if let Some(queue) = &self.bounded {
325            return self.send_bounded(queue, value);
326        }
327        self.send_unbounded(value)
328    }
329
330    fn try_send(&self, value: T) -> Result<()> {
331        if let Some(queue) = &self.bounded {
332            if self.closed.load(Ordering::Acquire) {
333                return Err(ChannelError::Closed);
334            }
335            queue.try_enqueue(value).map_err(|_| ChannelError::Full)?;
336            // The mutex is not held here, so nothing else orders this
337            // Store->Load; the helper's fence supplies it.
338            self.wake_receiver_after_push();
339            return Ok(());
340        }
341        self.try_send_unbounded(value)
342    }
343
344    fn recv(&self) -> Result<T> {
345        if let Some(queue) = &self.bounded {
346            return self.recv_bounded(queue);
347        }
348        self.recv_unbounded()
349    }
350
351    fn try_recv(&self) -> Result<T> {
352        if let Some(queue) = &self.bounded {
353            if let Some(value) = queue.try_dequeue() {
354                self.wake_sender_after_pop();
355                return Ok(value);
356            }
357            if self.closed.load(Ordering::Acquire) {
358                return Err(ChannelError::Closed);
359            }
360            return Err(ChannelError::Empty);
361        }
362        self.try_recv_unbounded()
363    }
364
365    fn is_empty(&self) -> bool {
366        if let Some(queue) = &self.bounded {
367            return queue.is_empty();
368        }
369        self.is_empty_unbounded()
370    }
371
372    fn is_full(&self) -> bool {
373        if let Some(queue) = &self.bounded {
374            return queue.is_full();
375        }
376        self.is_full_unbounded()
377    }
378
379    fn capacity(&self) -> Option<usize> {
380        if let Some(queue) = &self.bounded {
381            return Some(queue.capacity());
382        }
383        self.capacity_unbounded()
384    }
385}