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}