communication/lib.rs
1//! Two-lane priority merge (bombay card #225) — a ring + control-sideband design.
2//!
3//! Merge a **control** lane and a **user** lane into one [`Consumer`] so that
4//! control signals are served ahead of a user backlog, while keeping FIFO
5//! per lane, no loss, no lost wakeups, a clean teardown drain, a zero-alloc
6//! steady-state send, and no user starvation under a control flood.
7//!
8//! # Mechanism
9//!
10//! - **User lane** — a preallocated bounded Vyukov MPSC ring (power-of-two
11//! slots, one sequence atomic per slot, claim via a single `fetch_add`).
12//! Backpressure lives IN the ring: a producer that finds the head slot
13//! occupied parks on a `send_notify` eventcount and is woken one-per-pop by
14//! the consumer. A steady-state `try_send` is a handful of atomics and
15//! never allocates (P7).
16//! - **Control lane** — an unbounded lock-free MPSC chain of single-use
17//! 64-slot blocks: producers claim a global ticket with one `fetch_add`
18//! and publish into their slot; the single consumer pops in ticket order.
19//! Consumed blocks are reclaimed by the consumer once no producer can
20//! hold a block hint into them (an `in_flight` registration brackets each
21//! push's hint window), so the lane does not leak (P8). Control never
22//! shares the user ring, so a full ring can never delay it.
23//! - **Wakeup** — one shared [`tokio::sync::Notify`] gated by a `parked`
24//! flag. A producer pays a single atomic load when the consumer is active;
25//! a parked consumer is woken by `notify_one` after registering with
26//! enable-then-recheck, so the lost-wakeup class is impossible by
27//! construction and `recv` stays cancel-safe (no item is moved before the
28//! consumer is known to be awake to take it). The flag protocol is
29//! model-checked under `cfg(loom)` (`tests/loom.rs`).
30//!
31//! The public API below is FIXED — the property suite in `communication-testkit`
32//! depends on these exact names and signatures.
33
34use std::fmt;
35use std::mem::{MaybeUninit, replace};
36
37#[cfg(not(loom))]
38use std::sync::atomic::{AtomicBool, AtomicPtr, AtomicU32, AtomicUsize, Ordering};
39#[cfg(not(loom))]
40use std::sync::{Arc, Mutex, Weak};
41
42#[cfg(loom)]
43use loom::sync::atomic::{AtomicBool, AtomicPtr, AtomicU32, AtomicUsize, Ordering};
44#[cfg(loom)]
45use loom::sync::{Arc, Mutex};
46
47#[cfg(not(loom))]
48use tokio::sync::Notify;
49
50/// Synchronous stand-in for `tokio::sync::Notify`, used ONLY under
51/// `cfg(loom)`: models tokio's registration semantics explicitly so the
52/// `parked`/`waiting` flag protocol exercised by the loom model matches the
53/// async build. `enable()` arms a registration (the async twins call
54/// `Notified::enable` at the same points); `notify_one` wakes one
55/// registered waiter or stores a single permit; `notify_waiters` wakes all
56/// registered waiters without storing one. Wake credits are pooled
57/// (`woken`), so a registration abandoned by an early return can only
58/// cause a spurious extra loop turn, never a lost wakeup — matching
59/// tokio's `Notify`.
60#[cfg(loom)]
61mod sync_notify {
62 use loom::sync::atomic::Ordering;
63
64 pub struct Notify {
65 state: loom::sync::Mutex<State>,
66 cv: loom::sync::Condvar,
67 }
68
69 struct State {
70 /// Stored permit (notify with no registered waiter), at most one.
71 permit: bool,
72 /// Armed registrations not yet consumed by a `wait`.
73 enabled: usize,
74 /// Wake credits granted by `notify_one`/`notify_waiters`.
75 woken: usize,
76 }
77
78 impl Notify {
79 pub fn new() -> Self {
80 Self {
81 state: loom::sync::Mutex::new(State {
82 permit: false,
83 enabled: 0,
84 woken: 0,
85 }),
86 cv: loom::sync::Condvar::new(),
87 }
88 }
89
90 /// Arm a registration, mirroring `Notified::enable`. ALSO the
91 /// synchronization point of the loom protocol: the caller announces
92 /// its flag (`parked`/`waiting`) BEFORE this call, so the mutex
93 /// release here publishes the announcement to any subsequent
94 /// `notify_if*` lock.
95 pub fn enable(&self) {
96 self.state.lock().unwrap().enabled += 1;
97 }
98
99 /// Wake one waiter iff `flag` is set. The flag check happens UNDER
100 /// the lock: if this lock follows the waiter's `enable` in the
101 /// mutex order, the check synchronizes-with the waiter's
102 /// announcement and must observe it; otherwise the waiter's
103 /// post-enable re-check observes whatever this caller published
104 /// before locking. A plain atomic check outside the lock has no
105 /// happens-before edge to the announcement, and loom (correctly
106 /// over-approximating hardware) explores the lost wakeup.
107 pub fn notify_if(&self, flag: &loom::sync::atomic::AtomicBool) {
108 let mut g = self.state.lock().unwrap();
109 if !flag.load(Ordering::SeqCst) {
110 return;
111 }
112 if g.enabled > 0 {
113 g.enabled -= 1;
114 g.woken += 1;
115 self.cv.notify_one();
116 } else {
117 g.permit = true;
118 }
119 }
120
121 /// Wake one waiter iff `n` is nonzero — see `notify_if`.
122 pub fn notify_if_nonzero(&self, n: &loom::sync::atomic::AtomicUsize) {
123 let mut g = self.state.lock().unwrap();
124 if n.load(Ordering::SeqCst) == 0 {
125 return;
126 }
127 if g.enabled > 0 {
128 g.enabled -= 1;
129 g.woken += 1;
130 self.cv.notify_one();
131 } else {
132 g.permit = true;
133 }
134 }
135
136 pub fn notify_one(&self) {
137 let mut g = self.state.lock().unwrap();
138 if g.enabled > 0 {
139 g.enabled -= 1;
140 g.woken += 1;
141 self.cv.notify_one();
142 } else {
143 g.permit = true;
144 }
145 }
146
147 pub fn notify_waiters(&self) {
148 let mut g = self.state.lock().unwrap();
149 g.woken += g.enabled;
150 g.enabled = 0;
151 self.cv.notify_all();
152 }
153
154 /// Park until a wake credit or stored permit is available, then
155 /// consume it.
156 pub fn wait(&self) {
157 let mut g = self.state.lock().unwrap();
158 if g.permit {
159 g.permit = false;
160 return;
161 }
162 while g.woken == 0 {
163 g = self.cv.wait(g).unwrap();
164 }
165 g.woken -= 1;
166 }
167
168 /// Lock+unlock with no other effect: a happens-before edge to every
169 /// `notify_if*` that ran before it in the mutex order. The closed
170 /// sweep uses this to observe a departing sender's publishes — the
171 /// sender's pre-decrement `wake_consumer` (a `notify_if` lock) is
172 /// program-ordered before the decrement the consumer observed.
173 pub fn sync_point(&self) {
174 let _g = self.state.lock().unwrap();
175 }
176 }
177}
178
179#[cfg(loom)]
180use sync_notify::Notify;
181
182/// Consumer configuration.
183#[derive(Debug, Clone, Copy)]
184pub struct Config {
185 user_capacity: usize,
186 aging_cap: usize,
187}
188
189impl Config {
190 /// A config with the given bounded user-lane capacity and no aging.
191 ///
192 /// The capacity is a MINIMUM: the ring rounds it up to a power of two
193 /// (floor 2), and the rounded size is the effective capacity.
194 #[must_use]
195 pub const fn new(user_capacity: usize) -> Self {
196 Self {
197 user_capacity,
198 aging_cap: 0,
199 }
200 }
201
202 /// Force one waiting user through after `k` consecutive control dequeues.
203 /// `0` disables aging.
204 #[must_use]
205 pub const fn with_aging_cap(mut self, k: usize) -> Self {
206 self.aging_cap = k;
207 self
208 }
209}
210
211/// The item handed back by [`Consumer::recv`], tagged by its lane.
212#[derive(Debug, Clone, Copy, PartialEq, Eq)]
213pub enum Received<C, U> {
214 /// A control-lane item.
215 Control(C),
216 /// A user-lane item.
217 User(U),
218 /// The user lane is TERMINALLY closed: every `UserSender` is gone and
219 /// the ring is drained. Delivered EXACTLY ONCE, in the user stream's
220 /// FIFO position (after the last user item), subject to the same
221 /// control-first priority and aging as a user item. Terminal by
222 /// construction: no `UserSender` source exists besides `channel()`
223 /// and `Clone`, so a zero count can never rise again.
224 UserLaneClosed,
225}
226
227/// Items still queued when the consumer tears down (see [`Consumer::drain`]).
228#[derive(Debug)]
229pub struct Drained<C, U> {
230 /// Queued control items, in FIFO order.
231 pub control: Vec<C>,
232 /// Queued user items, in FIFO order.
233 pub user: Vec<U>,
234}
235
236/// The control lane closed because the [`Consumer`] was dropped.
237#[derive(Debug, thiserror::Error)]
238#[error("control lane closed: consumer dropped")]
239pub struct ControlClosed<T>(pub T);
240
241/// The user lane or mailbox admission was closed.
242/// The consumer can remain alive and drain accepted work.
243#[derive(Debug, thiserror::Error)]
244#[error("user lane closed")]
245pub struct UserClosed<T>(pub T);
246
247/// A non-blocking user send could not be enqueued.
248#[derive(Debug, thiserror::Error)]
249pub enum TrySendError<T> {
250 /// The bounded user lane is at capacity.
251 #[error("user lane full")]
252 Full(T),
253 /// The user lane or mailbox admission was closed.
254 /// The consumer can remain alive and drain accepted work.
255 #[error("user lane closed")]
256 Closed(T),
257}
258
259/// State shared by both lanes and the consumer: the single wakeup point and
260/// the consumer-liveness flag.
261struct Core {
262 /// Consumer wakeup. Single waiter (the consumer), `notify_one` suffices.
263 notify: Notify,
264 /// Set while the consumer is parked in `notify`; gates the wake call so a
265 /// producer with an active consumer pays one atomic load, nothing more.
266 parked: AtomicBool,
267 consumer_gone: AtomicBool,
268}
269
270impl Core {
271 /// Wake the consumer iff it is parked. The check is an RMW, not a plain
272 /// load: absent a happens-before edge, a load may legally miss the
273 /// consumer's earlier `parked` store (store buffering — loom exhibits
274 /// the lost wakeup), while an RMW always observes the newest value in
275 /// the modification order. The consumer's post-store re-check covers
276 /// the other Dekker direction.
277 #[cfg(not(loom))]
278 #[inline]
279 fn wake_consumer(&self) {
280 if self.parked.load(Ordering::SeqCst) {
281 self.notify.notify_one();
282 }
283 }
284
285 /// Wake the consumer iff it is parked — loom variant: the flag check is
286 /// taken under the `Notify` stand-in's lock so it synchronizes with the
287 /// consumer's announcement (see `sync_notify::Notify::notify_if`).
288 #[cfg(loom)]
289 fn wake_consumer(&self) {
290 self.notify.notify_if(&self.parked);
291 }
292}
293
294/// One slot of the user ring. `seq` is the Vyukov ticket, stored as `u32`
295/// (all ticket arithmetic is mod 2^32; the live-ticket window is bounded by
296/// `ring.len() < 2^31`, so truncation is sound): it equals the slot's
297/// position when free and `position + 1` once published. The narrow ticket
298/// halves slot memory for small `U`.
299struct Slot<U> {
300 seq: AtomicU32,
301 val: std::cell::UnsafeCell<MaybeUninit<U>>,
302}
303
304impl<U> fmt::Debug for Slot<U> {
305 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
306 f.debug_struct("Slot").finish_non_exhaustive()
307 }
308}
309
310/// The bounded user lane: a Vyukov MPSC ring with in-ring backpressure.
311struct UserLane<U> {
312 /// Ring slots, `ring.len()` is the capacity (min 1) rounded up to a
313 /// power of two.
314 ring: Box<[Slot<U>]>,
315 /// Next claim ticket for producers (multi-producer CAS).
316 tail: AtomicUsize,
317 /// Next pop ticket; written only by the single consumer.
318 head: AtomicUsize,
319 /// Parked-producer wakeup; gated by `waiting`.
320 send_notify: Notify,
321 /// Number of producers parked on a full ring (guards `send_notify`).
322 waiting: AtomicUsize,
323 /// Live [`UserSender`] handles; lane is closed when this reaches 0.
324 senders: AtomicUsize,
325 core: Arc<Core>,
326}
327
328// SAFETY: the ring protocol is the standard bounded-queue ticket discipline —
329// a slot's value is written before its `seq` is published (Release) and read
330// only after observing the published `seq` (Acquire); the consumer alone
331// advances `head`, producers alone advance `tail`, and a slot is reused only
332// after the consumer re-tickets it. `U` crosses threads only through that
333// protocol, so `U: Send` suffices.
334unsafe impl<U: Send> Send for UserLane<U> {}
335unsafe impl<U: Send> Sync for UserLane<U> {}
336
337impl<U> UserLane<U> {
338 #[inline]
339 fn mask(&self) -> usize {
340 self.ring.len() - 1
341 }
342
343 /// Try to enqueue `v`; on failure hands `v` back (ring full).
344 #[inline]
345 #[allow(
346 clippy::cast_possible_truncation,
347 clippy::cast_possible_wrap,
348 clippy::comparison_chain,
349 reason = "ticket arithmetic intentionally operates modulo 2^32; channel() bounds the live window"
350 )]
351 fn try_push(&self, v: U) -> Result<(), U> {
352 let mask = self.mask();
353 let mut pos = self.tail.load(Ordering::Acquire);
354 loop {
355 // SAFETY: `pos & mask` is always in bounds.
356 let slot = unsafe { self.ring.get_unchecked(pos & mask) };
357 let seq = slot.seq.load(Ordering::Acquire);
358 let pos32 = pos as u32;
359 let diff = seq.wrapping_sub(pos32) as i32;
360 if diff == 0 {
361 match self.tail.compare_exchange_weak(
362 pos,
363 pos.wrapping_add(1),
364 Ordering::AcqRel,
365 Ordering::Acquire,
366 ) {
367 Ok(_) => {
368 // SAFETY: we claimed cycle `pos` of this slot; the
369 // free `seq` ticket proves the previous occupant was
370 // consumed, so we hold exclusive access.
371 unsafe { slot.val.get().write(MaybeUninit::new(v)) };
372 slot.seq.store(pos32.wrapping_add(1), Ordering::Release);
373 return Ok(());
374 }
375 Err(actual) => pos = actual,
376 }
377 } else if diff < 0 {
378 return Err(v); // slot still occupied → full
379 } else {
380 pos = self.tail.load(Ordering::Acquire);
381 }
382 }
383 }
384
385 /// Dequeue the head item, if one is published. Consumer-only.
386 #[inline]
387 #[allow(
388 clippy::cast_possible_truncation,
389 reason = "ticket arithmetic intentionally operates modulo 2^32; channel() bounds the live window"
390 )]
391 fn pop(&self) -> Option<U> {
392 let head = self.head.load(Ordering::Relaxed);
393 // SAFETY: `head & mask` is always in bounds.
394 let slot = unsafe { self.ring.get_unchecked(head & self.mask()) };
395 let head32 = head as u32;
396 if slot.seq.load(Ordering::Acquire) != head32.wrapping_add(1) {
397 return None;
398 }
399 // SAFETY: the published `seq` (Acquire) orders the producer's write
400 // before this read; only the consumer pops, and the slot is not
401 // reused until re-ticketed below.
402 let v = unsafe { slot.val.get().read().assume_init() };
403 slot.seq.store(
404 head32.wrapping_add(self.ring.len() as u32),
405 Ordering::Release,
406 );
407 self.head.store(head.wrapping_add(1), Ordering::Relaxed);
408 // Wake one parked producer, if any — checked every 4th pop only;
409 // a parked producer is also released by the consumer's pre-park
410 // check, so it can never sleep past an empty ring.
411 if head.is_multiple_of(8) {
412 self.release_one_waiter();
413 }
414 Some(v)
415 }
416
417 /// Release one parked producer, if any is waiting.
418 #[cfg(not(loom))]
419 #[inline]
420 fn release_one_waiter(&self) {
421 if self.waiting.load(Ordering::SeqCst) != 0 {
422 self.send_notify.notify_one();
423 }
424 }
425
426 /// Release one parked producer, if any — loom variant: the `waiting`
427 /// check is taken under the `Notify` stand-in's lock (see
428 /// `sync_notify::Notify::notify_if`).
429 #[cfg(loom)]
430 fn release_one_waiter(&self) {
431 self.send_notify.notify_if_nonzero(&self.waiting);
432 }
433
434 /// True once every [`UserSender`] is gone.
435 #[inline]
436 fn closed(&self) -> bool {
437 self.senders.load(Ordering::Acquire) == 0
438 }
439}
440
441impl<U> Drop for UserLane<U> {
442 #[allow(
443 clippy::cast_possible_truncation,
444 reason = "ticket arithmetic intentionally operates modulo 2^32; channel() bounds the live window"
445 )]
446 fn drop(&mut self) {
447 // Quiescent: all senders and the consumer are gone, so `head..tail`
448 // are exactly the published-but-unconsumed items.
449 let head = self.head.load(Ordering::Relaxed); // quiescent: &mut self
450 let tail = self.tail.load(Ordering::Relaxed); // quiescent: &mut self
451 for pos in head..tail {
452 let slot = &self.ring[pos & self.mask()];
453 if slot.seq.load(Ordering::Relaxed) == (pos as u32).wrapping_add(1) {
454 // SAFETY: published and never consumed; exclusive `&mut self`.
455 // `assume_init_drop` drops the CONTENTS — a plain
456 // `drop_in_place` on the `MaybeUninit` drops only the
457 // wrapper, which has no drop glue (payload leak).
458 unsafe { (*slot.val.get()).assume_init_drop() };
459 }
460 }
461 }
462}
463
464/// Slots per control block. Blocks are single-use (never reused), so a slot
465/// needs no Vyukov ticket — just a publish flag. Under `cfg(loom)` the
466/// block is shrunk so the tiny model crosses block boundaries (linking,
467/// hint advance, and reclamation are all exercised).
468#[cfg(not(loom))]
469const CBLOCK: usize = 64;
470/// `CBLOCK` under loom — see above.
471#[cfg(loom)]
472const CBLOCK: usize = 2;
473
474/// One control slot: written once, published with `ready`.
475struct CSlot<C> {
476 ready: AtomicBool,
477 val: std::cell::UnsafeCell<MaybeUninit<C>>,
478}
479
480/// A linked block of the control queue.
481struct CBlock<C> {
482 /// Global block index (`block_base = idx * CBLOCK`).
483 idx: usize,
484 slots: [CSlot<C>; CBLOCK],
485 next: AtomicPtr<CBlock<C>>,
486}
487
488impl<C> CBlock<C> {
489 #[allow(
490 clippy::unnecessary_box_returns,
491 reason = "control blocks are immediately converted to stable raw pointers for the linked queue"
492 )]
493 fn new(idx: usize) -> Box<Self> {
494 Box::new(Self {
495 idx,
496 slots: std::array::from_fn(|_| CSlot {
497 ready: AtomicBool::new(false),
498 val: std::cell::UnsafeCell::new(MaybeUninit::uninit()),
499 }),
500 next: AtomicPtr::new(std::ptr::null_mut()),
501 })
502 }
503}
504
505/// The unbounded control sideband: a lock-free MPSC chain of single-use
506/// blocks. Producers claim a global ticket with one `fetch_add` and publish
507/// into their slot; the single consumer pops in ticket order lock-free.
508/// Control never shares the user ring, so a full ring can never delay it.
509struct ControlLane<C> {
510 /// Next claim ticket (multi-producer `fetch_add`; unbounded).
511 tail: AtomicUsize,
512 /// Highest linked block (best-effort hint; producers walk forward from
513 /// it). Never points behind `first_block` (reclaim's `min` bound), so
514 /// any value read here is a live block.
515 tail_block: AtomicPtr<CBlock<C>>,
516 /// Live frontier: the first block a push may walk from. Monotonic,
517 /// consumer-only writer, published BEFORE the corresponding reclamation
518 /// check so a push that registers in-flight afterwards is guaranteed
519 /// (`SeqCst` chain) to read a frontier above the freed prefix.
520 first_block: AtomicPtr<CBlock<C>>,
521 /// First unfreed block. Blocks in `[free_anchor, first_block)` are
522 /// retired (fully consumed) but not yet freed — freeing waits for an
523 /// `in_flight == 0` crossing. Consumer-only writer, read at lane drop.
524 free_anchor: AtomicPtr<CBlock<C>>,
525 /// Slow (block-crossing) pushes between their in-flight registration
526 /// and slot publish. The consumer frees retired blocks only while this
527 /// is zero, which proves no push can be walking the freed prefix.
528 in_flight: AtomicUsize,
529 /// Tickets consumed, published by the consumer ONLY at teardown so the
530 /// lane's `Drop` reclaims exactly the published-but-unconsumed items.
531 consumed: AtomicUsize,
532 /// Live [`ControlSender`] handles; lane is closed when this reaches 0.
533 senders: AtomicUsize,
534 core: Arc<Core>,
535}
536
537// SAFETY: a slot's value is written before `ready` is set (Release) and read
538// only after observing `ready` (Acquire); each ticket is claimed by exactly
539// one producer; blocks are linked before any of their slots are published
540// and freed only when all senders and the consumer are gone. `C` crosses
541// threads only through that protocol, so `C: Send` suffices.
542unsafe impl<C: Send> Send for ControlLane<C> {}
543unsafe impl<C: Send> Sync for ControlLane<C> {}
544
545impl<C> ControlLane<C> {
546 /// Claim the next ticket and publish `item` into its slot.
547 fn push(&self, item: C) {
548 // Load the hint BEFORE claiming the ticket. `tail_block` always
549 // points at a block containing an already-claimed ticket `t`, and
550 // any ticket claimed after this load satisfies `pos >= tail > t`,
551 // so `hint.idx <= t / CBLOCK <= pos / CBLOCK == want`: the hint's
552 // block is never BEYOND the ticket's block. (Claim-first ordering
553 // lets a preempted push resume with a hint beyond its ticket's
554 // block and publish into the wrong slot.)
555 let hint = self.tail_block.load(Ordering::SeqCst);
556 let pos = self.tail.fetch_add(1, Ordering::AcqRel);
557 let want = pos / CBLOCK;
558 // SAFETY: `hint` is a live block: it contains a claimed ticket and
559 // `hint.idx <= want` by the ordering above.
560 if unsafe { &*hint }.idx == want {
561 // FAST PATH (steady state): the ticket lives in the hinted
562 // block, no walk. The block cannot be reclaimed before our
563 // publish: `reclaim` only frees blocks the consumer has fully
564 // crossed, and our ticket in it is not even published yet, so
565 // the consumer's cursor cannot pass it.
566 let slot = unsafe { &(*hint).slots[pos % CBLOCK] };
567 unsafe { slot.val.get().write(MaybeUninit::new(item)) };
568 slot.ready.store(true, Ordering::Release);
569 return;
570 }
571 // SLOW PATH (block crossing): the walk passes through intermediate
572 // blocks, which the consumer may reclaim once it has crossed them —
573 // register in-flight so `reclaim` provably holds off (see
574 // `reclaim`), THEN take a provably-live hint. The inc is ordered
575 // before the hint loads below, so any hint this push walks from is
576 // covered by the in_flight guard; the initial `hint` may be stale
577 // and is never dereferenced here.
578 self.in_flight.fetch_add(1, Ordering::SeqCst);
579 let mut b = self.tail_block.load(Ordering::SeqCst);
580 // SAFETY: `b` is live: `tail_block` never points behind the
581 // consumer's reclaim frontier, and `in_flight >= 1` blocks new
582 // frees for the rest of this push.
583 if unsafe { &*b }.idx > want {
584 // Another producer advanced `tail_block` past our ticket's
585 // block; walk from the reclaim frontier instead (also live,
586 // via the SeqCst chain documented in `reclaim`).
587 b = self.first_block.load(Ordering::SeqCst);
588 }
589 // Advance the block chain to the block containing `pos`, linking
590 // fresh blocks as needed. Racing linkers: exactly one CAS wins per
591 // `next` pointer; losers free their spare block.
592 while unsafe { &*b }.idx < want {
593 let next = unsafe { &*b }.next.load(Ordering::Acquire);
594 if next.is_null() {
595 let fresh = Box::into_raw(CBlock::new(unsafe { &*b }.idx + 1));
596 match unsafe { &*b }.next.compare_exchange(
597 std::ptr::null_mut(),
598 fresh,
599 Ordering::AcqRel,
600 Ordering::Acquire,
601 ) {
602 Ok(_) => {
603 let _ = self.tail_block.compare_exchange(
604 b,
605 fresh,
606 Ordering::SeqCst,
607 Ordering::SeqCst,
608 );
609 }
610 Err(_) => drop(unsafe { Box::from_raw(fresh) }),
611 }
612 continue;
613 }
614 let _ = self
615 .tail_block
616 .compare_exchange(b, next, Ordering::SeqCst, Ordering::SeqCst);
617 b = next;
618 }
619 // SAFETY: `b` is the block containing `pos` (the walk terminates
620 // with `b.idx == want`); `b` and every block the walk touched are
621 // live under the in_flight guard; the ticket is unique to this
622 // push, so we hold exclusive access to the slot.
623 let slot = unsafe { &(*b).slots[pos % CBLOCK] };
624 unsafe { slot.val.get().write(MaybeUninit::new(item)) };
625 slot.ready.store(true, Ordering::Release);
626 self.in_flight.fetch_sub(1, Ordering::Relaxed);
627 }
628
629 /// Free fully-consumed blocks behind `cursor` (the consumer's current
630 /// block). Consumer-only. Two-phase:
631 ///
632 /// 1. PUBLISH the live frontier (`first_block`) to the first block at
633 /// or after `min(cursor, tail_block)`. A block below `cursor` was
634 /// fully consumed (values moved out); a block below `tail_block`
635 /// can never be a push's walk start. The `min` is load-bearing: a
636 /// freshly LINKED block briefly lags `tail_block` while its
637 /// linker's advance CAS is in flight.
638 /// 2. FREE the retired prefix `[free_anchor, frontier)` — only if
639 /// `in_flight == 0`; otherwise freeing catches up at a later
640 /// crossing (no leak: `free_anchor` only advances on real frees).
641 ///
642 /// Ordering (Dekker over SeqCst): the frontier store (p1) precedes the
643 /// `in_flight` read (p2) in program order. A slow push's increment
644 /// (p3) precedes its frontier/hint loads (p4). If p2 observes zero,
645 /// any push still walking has p3 AFTER p2 in the `SeqCst` order, hence
646 /// p4 > p3 > p2 > p1: it reads a frontier >= the one published here,
647 /// and `tail_block` >= `first_block` always, so no push ever walks a
648 /// block this call frees. FAST pushes (`hint.idx == want`) write into
649 /// a block containing their own unpublished ticket, which `cursor`
650 /// can never have crossed — they need no guard.
651 fn reclaim(&self, cursor: *const CBlock<C>) {
652 // Under loom the block chain is never reclaimed: loom's visibility
653 // model (deliberately) allows a load to observe a cell's INITIAL
654 // value absent a happens-before edge, so the SeqCst Dekker chain
655 // that makes this safe on hardware cannot be model-checked — a
656 // stale hint read would manufacture a use-after-free the real
657 // protocol forbids. The wakeup/FIFO/no-loss properties the loom
658 // lane exists to prove are unaffected; the leak gate (`tests/
659 // leak.rs`) covers reclamation on hardware.
660 if cfg!(loom) {
661 return;
662 }
663 // Amortize: reclaim every 4th block crossing — the retired prefix
664 // stays bounded by a few blocks (and is reclaimed wholesale at lane
665 // drop), while the drain hot path pays the check once per 4
666 // crossings instead of every crossing.
667 if unsafe { &*cursor }.idx % 8 != 0 {
668 return;
669 }
670 let tb = self.tail_block.load(Ordering::SeqCst);
671 // SAFETY: `tb` and `cursor` are live blocks (reclamation only ever
672 // frees strictly behind both).
673 let bound = unsafe { &*tb }.idx.min(unsafe { &*cursor }.idx);
674 // Advance the frontier to the first block at or after `bound`.
675 let mut stop = self.first_block.load(Ordering::Relaxed);
676 while !stop.is_null() && unsafe { &*stop }.idx < bound {
677 stop = unsafe { &*stop }.next.load(Ordering::Acquire);
678 }
679 self.first_block.store(stop, Ordering::SeqCst);
680 if self.in_flight.load(Ordering::SeqCst) != 0 {
681 return;
682 }
683 let mut b = self.free_anchor.load(Ordering::Relaxed);
684 // Amortize: free in batches — the retired prefix keeps the leak
685 // bounded by a few blocks, and bulk-freeing avoids per-crossing
686 // dealloc churn in the drain hot path.
687 if unsafe { &*stop }.idx - unsafe { &*b }.idx < 4 {
688 return;
689 }
690 while !b.is_null() && b != stop {
691 let next = unsafe { &*b }.next.load(Ordering::Acquire);
692 // SAFETY: the consumer crossed every block below `stop`, so
693 // all their slots were consumed (values moved out, nothing to
694 // drop); `in_flight == 0` proves no push is walking them;
695 // behind `tail_block`, so no new push reaches them.
696 drop(unsafe { Box::from_raw(b) });
697 b = next;
698 }
699 self.free_anchor.store(b, Ordering::Relaxed);
700 }
701
702 /// True once every [`ControlSender`] is gone.
703 #[inline]
704 fn closed(&self) -> bool {
705 self.senders.load(Ordering::Acquire) == 0
706 }
707}
708
709impl<C> Drop for ControlLane<C> {
710 fn drop(&mut self) {
711 // Quiescent: all senders and the consumer are gone. Reclaim every
712 // published-but-unconsumed slot, then free the remaining chain
713 // (blocks the consumer already freed are behind `free_anchor`).
714 let consumed = self.consumed.load(Ordering::Relaxed); // quiescent: &mut self
715 let tail = self.tail.load(Ordering::Relaxed); // quiescent: &mut self
716 let mut b = self.free_anchor.load(Ordering::Relaxed); // quiescent: &mut self
717 while !b.is_null() {
718 let block = unsafe { Box::from_raw(b) };
719 let base = block.idx * CBLOCK;
720 for pos in base.max(consumed)..(base + CBLOCK).min(tail) {
721 let slot = &block.slots[pos % CBLOCK];
722 if slot.ready.load(Ordering::Relaxed) {
723 // SAFETY: published and never consumed; exclusive ownership.
724 // `assume_init_drop` drops the CONTENTS — see
725 // `UserLane::drop`.
726 unsafe { (*slot.val.get()).assume_init_drop() };
727 }
728 }
729 b = block.next.load(Ordering::Relaxed); // quiescent: &mut self
730 }
731 }
732}
733
734/// Cloneable handle to the control lane. `send` never blocks.
735pub struct ControlSender<C> {
736 lane: Arc<ControlLane<C>>,
737}
738
739impl<C> Clone for ControlSender<C> {
740 fn clone(&self) -> Self {
741 self.lane.senders.fetch_add(1, Ordering::Relaxed);
742 Self {
743 lane: self.lane.clone(),
744 }
745 }
746}
747
748impl<C> Drop for ControlSender<C> {
749 fn drop(&mut self) {
750 // Wake BEFORE decrementing: under loom this mutex-synchronizes the
751 // departing sender's publishes with a consumer that later observes
752 // the lane closed (see `sync_notify`); in production it is a cheap
753 // gated check and a harmless spurious wake.
754 self.lane.core.wake_consumer();
755 if self.lane.senders.fetch_sub(1, Ordering::Release) == 1 {
756 // Last sender: wake a parked consumer so `recv` can observe the
757 // closure and return `None` once both lanes are drained.
758 self.lane.core.wake_consumer();
759 }
760 }
761}
762
763impl<C> ControlSender<C> {
764 /// Enqueue a control signal. Never blocks (the lane is unbounded).
765 ///
766 /// # Errors
767 /// Returns [`ControlClosed`] carrying `item` if the consumer is gone.
768 pub fn send(&self, item: C) -> Result<(), ControlClosed<C>> {
769 if self.lane.core.consumer_gone.load(Ordering::Acquire) {
770 return Err(ControlClosed(item));
771 }
772 self.lane.push(item);
773 self.lane.core.wake_consumer();
774 Ok(())
775 }
776}
777
778/// Cloneable handle to the bounded user lane.
779pub struct UserSender<U> {
780 lane: Arc<UserLane<U>>,
781}
782
783impl<U> Clone for UserSender<U> {
784 fn clone(&self) -> Self {
785 self.lane.senders.fetch_add(1, Ordering::Relaxed);
786 Self {
787 lane: self.lane.clone(),
788 }
789 }
790}
791
792impl<U> Drop for UserSender<U> {
793 fn drop(&mut self) {
794 // Wake BEFORE decrementing — see `ControlSender::drop`.
795 self.lane.core.wake_consumer();
796 if self.lane.senders.fetch_sub(1, Ordering::Release) == 1 {
797 self.lane.core.wake_consumer();
798 }
799 }
800}
801
802impl<U> UserSender<U> {
803 /// Enqueue a user message, awaiting capacity (backpressure).
804 ///
805 /// The send linearizes at the ring-slot publish. If the consumer is
806 /// dropped while this send is blocked on a full ring, teardown releases
807 /// it with its payload — a blocked send never reports success for an
808 /// item that was not enqueued.
809 ///
810 /// # Errors
811 /// Returns [`UserClosed`] carrying `item` if the consumer is gone —
812 /// including when teardown releases a send blocked on a full ring.
813 #[cfg(not(loom))]
814 pub async fn send(&self, item: U) -> Result<(), UserClosed<U>> {
815 let lane = &self.lane;
816 let mut item = item;
817 loop {
818 if lane.core.consumer_gone.load(Ordering::Acquire) {
819 return Err(UserClosed(item));
820 }
821 match lane.try_push(item) {
822 Ok(()) => {
823 lane.core.wake_consumer();
824 return Ok(());
825 }
826 Err(v) => item = v,
827 }
828 // Ring full: register, announce, re-check, then park. The
829 // announce (SeqCst) pairs with the consumer's SeqCst `waiting`
830 // load, so a pop racing this park either wakes us or is seen by
831 // the re-check below.
832 let notified = lane.send_notify.notified();
833 tokio::pin!(notified);
834 notified.as_mut().enable();
835 lane.waiting.fetch_add(1, Ordering::SeqCst);
836 match lane.try_push(item) {
837 Ok(()) => {
838 lane.waiting.fetch_sub(1, Ordering::Relaxed);
839 lane.core.wake_consumer();
840 return Ok(());
841 }
842 Err(v) => item = v,
843 }
844 // Park-gating check (RMW — a plain load may miss the consumer's
845 // earlier teardown store; see `wake_consumer`).
846 if lane.core.consumer_gone.load(Ordering::Acquire) {
847 lane.waiting.fetch_sub(1, Ordering::Relaxed);
848 // Teardown raced the registration seam: the item never
849 // linearized (every push attempt handed it back), so it
850 // returns to the caller, exactly once.
851 return Err(UserClosed(item));
852 }
853 notified.await;
854 lane.waiting.fetch_sub(1, Ordering::Relaxed);
855 if lane.core.consumer_gone.load(Ordering::Acquire) {
856 // Teardown released this parked waiter before the item
857 // linearized: return it to the caller, exactly once.
858 return Err(UserClosed(item));
859 }
860 }
861 }
862
863 /// Blocking twin of `send`, compiled ONLY under `cfg(loom)`: the same
864 /// announce/re-check/park protocol over `waiting`, with the sync
865 /// `Notify` stand-in's `wait()` in place of the `Notified` future.
866 /// Drives the loom model in `tests/loom.rs`.
867 #[cfg(loom)]
868 #[doc(hidden)]
869 pub fn send_blocking(&self, item: U) -> Result<(), UserClosed<U>> {
870 let lane = &self.lane;
871 let mut item = item;
872 loop {
873 if lane.core.consumer_gone.load(Ordering::Acquire) {
874 return Err(UserClosed(item));
875 }
876 match lane.try_push(item) {
877 Ok(()) => {
878 lane.core.wake_consumer();
879 return Ok(());
880 }
881 Err(v) => item = v,
882 }
883 // Ring full: announce (SeqCst), re-check, then park. The
884 // stand-in's `enable` lock publishes the announcement; the
885 // SECOND re-check after it catches any pop the first missed
886 // (the consumer's `notify_if_nonzero` lock covers the other
887 // direction).
888 lane.waiting.fetch_add(1, Ordering::SeqCst);
889 match lane.try_push(item) {
890 Ok(()) => {
891 lane.waiting.fetch_sub(1, Ordering::Relaxed);
892 lane.core.wake_consumer();
893 return Ok(());
894 }
895 Err(v) => item = v,
896 }
897 // Park-gating check (RMW — see `wake_consumer`).
898 if lane.core.consumer_gone.load(Ordering::Acquire) {
899 lane.waiting.fetch_sub(1, Ordering::Relaxed);
900 // Teardown raced the registration seam: the item never
901 // linearized, so it returns to the caller, exactly once.
902 return Err(UserClosed(item));
903 }
904 lane.send_notify.enable();
905 match lane.try_push(item) {
906 Ok(()) => {
907 lane.waiting.fetch_sub(1, Ordering::Relaxed);
908 lane.core.wake_consumer();
909 return Ok(());
910 }
911 Err(v) => item = v,
912 }
913 if lane.core.consumer_gone.load(Ordering::Acquire) {
914 lane.waiting.fetch_sub(1, Ordering::Relaxed);
915 // Teardown raced the post-registration re-check: the item
916 // never linearized, so it returns to the caller.
917 return Err(UserClosed(item));
918 }
919 lane.send_notify.wait();
920 lane.waiting.fetch_sub(1, Ordering::Relaxed);
921 if lane.core.consumer_gone.load(Ordering::Acquire) {
922 // Teardown released this parked waiter before the item
923 // linearized: return it to the caller, exactly once.
924 return Err(UserClosed(item));
925 }
926 }
927 }
928
929 /// Enqueue a user message without blocking.
930 ///
931 /// # Errors
932 /// [`TrySendError::Full`] if the lane is at capacity, [`TrySendError::Closed`]
933 /// if the consumer is gone — each carrying `item` back.
934 #[inline]
935 pub fn try_send(&self, item: U) -> Result<(), TrySendError<U>> {
936 if self.lane.core.consumer_gone.load(Ordering::Acquire) {
937 return Err(TrySendError::Closed(item));
938 }
939 match self.lane.try_push(item) {
940 Ok(()) => {
941 self.lane.core.wake_consumer();
942 Ok(())
943 }
944 Err(v) => Err(TrySendError::Full(v)),
945 }
946 }
947
948 /// Derive a non-owning anchor from this sender (address-table
949 /// endpoint). The anchor holds no liveness: the lane closes when the
950 /// last counting `UserSender` drops even while anchors are held.
951 #[must_use]
952 pub fn anchor(&self) -> UserAnchor<U> {
953 #[cfg(not(loom))]
954 let lane = Arc::downgrade(&self.lane);
955 #[cfg(loom)]
956 let lane = self.lane.clone();
957 UserAnchor { lane }
958 }
959}
960
961/// Non-owning weak capability to the user lane (address-table
962/// endpoint). Holds no liveness: the lane closes when the last counting
963/// [`UserSender`] drops even while anchors are held. Every delivery first
964/// atomically acquires a temporary live [`UserSender`], so a delivery racing
965/// the last sender drop either linearizes entirely before `UserLaneClosed`
966/// or fails with its payload entirely after.
967///
968/// Under `cfg(loom)` the anchor holds a strong `Arc` instead of a `Weak`
969/// (loom's `Arc` has no `downgrade`): the semantic liveness is governed by
970/// the `senders` counter, NOT the strong count, so the model stays faithful
971/// — the strong-ref distinction is pure memory reclamation, covered on
972/// hardware by `tests/leak.rs`.
973pub struct UserAnchor<U> {
974 #[cfg(not(loom))]
975 lane: Weak<UserLane<U>>,
976 #[cfg(loom)]
977 lane: Arc<UserLane<U>>,
978}
979
980/// Affine authority over user-message admission for an actor mailbox.
981///
982/// Unlike [`UserSender`], this capability is deliberately not cloneable and
983/// does not expose a counting sender. Cloneable addresses are obtained as
984/// [`MailboxRef`]s instead. Consuming this owner with
985/// [`MailboxOwner::close_admission`] prevents every stale reference from
986/// acquiring a new delivery permit while allowing permits acquired before
987/// the close to finish. The consumer can then receive every accepted message
988/// followed by [`Received::UserLaneClosed`].
989///
990/// Dropping the owner has the same admission-closing effect. The explicit
991/// method is preferred at graceful-retirement sites because it documents the
992/// lifecycle transition.
993#[must_use = "dropping the mailbox owner closes user-message admission"]
994pub struct MailboxOwner<U> {
995 admission: Arc<Mutex<MailboxAdmission<U>>>,
996}
997
998/// Admission is independent of the sender count: already acquired permits
999/// retain their sender while this owner closes every subsequent acquisition.
1000/// Critical sections contain only enum replacement and sender cloning: no
1001/// user code, payload destruction, or await can unwind while holding the lock.
1002enum MailboxAdmission<U> {
1003 Open(UserSender<U>),
1004 Closed,
1005}
1006
1007impl<U> MailboxOwner<U> {
1008 /// Create another non-owning, cloneable address for this mailbox.
1009 #[must_use]
1010 pub fn actor_ref(&self) -> MailboxRef<U> {
1011 #[cfg(not(loom))]
1012 let admission = Arc::downgrade(&self.admission);
1013 #[cfg(loom)]
1014 let admission = self.admission.clone();
1015 MailboxRef { admission }
1016 }
1017
1018 /// Atomically close admission with respect to deliveries through every
1019 /// [`MailboxRef`], including stale clones.
1020 ///
1021 /// A racing delivery has exactly one of two outcomes: it acquired a
1022 /// temporary sender before this close and must be received before the
1023 /// terminal marker, or it receives its original payload back as closed.
1024 /// This operation does not destroy the consumer or discard queued work.
1025 pub fn close_admission(self) {
1026 drop(self);
1027 }
1028}
1029
1030impl<U> Drop for MailboxOwner<U> {
1031 fn drop(&mut self) {
1032 let retired_admission = {
1033 let mut admission = self.admission.lock().expect("mailbox admission poisoned");
1034 replace(&mut *admission, MailboxAdmission::Closed)
1035 };
1036 // Final lane destruction can drop queued user payloads; release the
1037 // admission lock before dropping the durable sender and its lane.
1038 drop(retired_admission);
1039 }
1040}
1041
1042/// Cloneable, non-owning actor address governed by a [`MailboxOwner`].
1043///
1044/// Each send acquires exactly one temporary admission permit. Unlike
1045/// [`UserAnchor`], this restricted address deliberately does not expose that
1046/// permit as a [`UserSender`], so code holding stale references cannot retain
1047/// or clone admission beyond an individual delivery operation.
1048pub struct MailboxRef<U> {
1049 #[cfg(not(loom))]
1050 admission: Weak<Mutex<MailboxAdmission<U>>>,
1051 #[cfg(loom)]
1052 admission: Arc<Mutex<MailboxAdmission<U>>>,
1053}
1054
1055impl<U> Clone for MailboxRef<U> {
1056 fn clone(&self) -> Self {
1057 Self {
1058 admission: self.admission.clone(),
1059 }
1060 }
1061}
1062
1063impl<U> MailboxRef<U> {
1064 fn acquire_sender(&self) -> Option<UserSender<U>> {
1065 #[cfg(not(loom))]
1066 let admission = self.admission.upgrade()?;
1067 #[cfg(loom)]
1068 let admission = self.admission.clone();
1069 let sender = match &*admission.lock().expect("mailbox admission poisoned") {
1070 MailboxAdmission::Open(sender) => Some(sender.clone()),
1071 MailboxAdmission::Closed => None,
1072 };
1073 // The lock and promoted admission Arc are gone before any delivery
1074 // awaits capacity; only the exact acquired UserSender remains live.
1075 sender
1076 }
1077
1078 /// Deliver one message, awaiting bounded-lane capacity.
1079 ///
1080 /// Returns the original message if admission was already closed or the
1081 /// consumer disappeared before the message was accepted.
1082 #[cfg(not(loom))]
1083 pub async fn send(&self, item: U) -> Result<(), UserClosed<U>> {
1084 let Some(sender) = self.acquire_sender() else {
1085 return Err(UserClosed(item));
1086 };
1087 sender.send(item).await
1088 }
1089
1090 /// Try to deliver one message without waiting for capacity.
1091 ///
1092 /// A close racing this call either follows the accepted message in the
1093 /// receive stream or wins and returns the original message as closed.
1094 pub fn try_send(&self, item: U) -> Result<(), TrySendError<U>> {
1095 let Some(sender) = self.acquire_sender() else {
1096 return Err(TrySendError::Closed(item));
1097 };
1098 sender.try_send(item)
1099 }
1100}
1101
1102impl<U> Clone for UserAnchor<U> {
1103 fn clone(&self) -> Self {
1104 Self {
1105 lane: self.lane.clone(),
1106 }
1107 }
1108}
1109
1110impl<U> UserAnchor<U> {
1111 /// Acquire a temporary live [`UserSender`] iff the lane is still open.
1112 ///
1113 /// The liveness increment is CONDITIONAL and atomic — one read-modify-
1114 /// write loop over the live-sender count that succeeds only while the
1115 /// count is nonzero. It never increments zero, so it can never resurrect
1116 /// a closed lane; a last-sender drop racing this upgrade is caught by
1117 /// the RMW reading zero, and the upgrade then fails with no side effect.
1118 /// The returned sender counts as live until dropped — held across an
1119 /// async send, released on success, error, or cancellation.
1120 #[must_use]
1121 pub fn upgrade(&self) -> Option<UserSender<U>> {
1122 // Strong ref first: the consumer or another sender keeps the lane
1123 // alive. Then the conditional RMW gates on the live-sender count —
1124 // the `Weak` step only fails once the lane is fully reclaimed.
1125 #[cfg(not(loom))]
1126 let lane = self.lane.upgrade()?;
1127 #[cfg(loom)]
1128 let lane = self.lane.clone();
1129 lane.senders
1130 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| {
1131 // Conditional: succeed only while the count is nonzero — a
1132 // zero count is terminal (the lane is closed) and must never
1133 // be incremented, or the anchor would resurrect the lane.
1134 if n == 0 { None } else { n.checked_add(1) }
1135 })
1136 .ok()?;
1137 Some(UserSender { lane })
1138 }
1139
1140 /// Enqueue a user message through a temporary live sender, awaiting
1141 /// capacity (backpressure). The temporary sender stays alive across the
1142 /// await — the lane cannot close while the delivery is in flight — and
1143 /// is dropped on success, error, or future cancellation, releasing its
1144 /// temporary liveness.
1145 ///
1146 /// # Errors
1147 /// Returns [`UserClosed`] carrying `item` if the lane is closed (no
1148 /// live counting sender, or the consumer is gone).
1149 #[cfg(not(loom))]
1150 pub async fn send(&self, item: U) -> Result<(), UserClosed<U>> {
1151 let Some(sender) = self.upgrade() else {
1152 return Err(UserClosed(item));
1153 };
1154 sender.send(item).await
1155 }
1156
1157 /// Blocking twin of `send`, compiled ONLY under `cfg(loom)` — drives the
1158 /// anchor loom models in `tests/loom.rs`.
1159 #[cfg(loom)]
1160 #[doc(hidden)]
1161 pub fn send_blocking(&self, item: U) -> Result<(), UserClosed<U>> {
1162 let Some(sender) = self.upgrade() else {
1163 return Err(UserClosed(item));
1164 };
1165 sender.send_blocking(item)
1166 }
1167
1168 /// Enqueue a user message without blocking, through a temporary live
1169 /// sender.
1170 ///
1171 /// # Errors
1172 /// [`TrySendError::Closed`] carrying `item` if the lane is closed (no
1173 /// live counting sender, or the consumer is gone);
1174 /// [`TrySendError::Full`] if the lane is at capacity.
1175 pub fn try_send(&self, item: U) -> Result<(), TrySendError<U>> {
1176 let Some(sender) = self.upgrade() else {
1177 return Err(TrySendError::Closed(item));
1178 };
1179 sender.try_send(item)
1180 }
1181}
1182
1183/// Pop the next published control in ticket order, advancing the consumer's
1184/// block cursor and reclaiming consumed blocks on each crossing. Expands to
1185/// a block expression evaluating to `Option<C>`; takes the consumer
1186/// identifier so it can run under `recv`'s split borrows.
1187macro_rules! pop_control {
1188 ($slf:ident) => {{
1189 let mut block = $slf.ctl_block;
1190 // SAFETY: `ctl_block` is always a live block (reclamation frees
1191 // only strictly behind it; the chain anchor lives until lane drop,
1192 // and the consumer holds an Arc on the lane).
1193 let b = unsafe { &*block };
1194 if $slf.ctl_head == (b.idx + 1) * CBLOCK {
1195 // Crossed a block boundary: the next block is linked before any
1196 // of its slots are published, so `null` means the lane is empty.
1197 let next = b.next.load(Ordering::Acquire);
1198 if next.is_null() {
1199 None
1200 } else {
1201 $slf.ctl_block = next;
1202 // The block just left is fully consumed; reclaim it (and
1203 // anything older) once no push can hold a hint into it.
1204 $slf.ctl.reclaim(next);
1205 block = next;
1206 // SAFETY: as above; `ctl_head % CBLOCK` is in bounds.
1207 let slot = unsafe { &(*block).slots[$slf.ctl_head % CBLOCK] };
1208 if slot.ready.load(Ordering::Acquire) {
1209 // SAFETY: `ready` (Acquire) orders the producer's write
1210 // before this read; only the consumer pops, and each
1211 // slot is published once.
1212 let v = unsafe { slot.val.get().read().assume_init() };
1213 $slf.ctl_head = $slf.ctl_head.wrapping_add(1);
1214 Some(v)
1215 } else {
1216 None
1217 }
1218 }
1219 } else {
1220 // SAFETY: as above; `ctl_head % CBLOCK` is in bounds.
1221 let slot = unsafe { &(*block).slots[$slf.ctl_head % CBLOCK] };
1222 let __r = slot.ready.load(Ordering::Acquire);
1223 if __r {
1224 // SAFETY: as above.
1225 let v = unsafe { slot.val.get().read().assume_init() };
1226 $slf.ctl_head = $slf.ctl_head.wrapping_add(1);
1227 Some(v)
1228 } else {
1229 None
1230 }
1231 }
1232 }};
1233}
1234
1235/// The single consumer that merges the two lanes.
1236pub struct Consumer<C, U> {
1237 ctl: Arc<ControlLane<C>>,
1238 usr: Arc<UserLane<U>>,
1239 core: Arc<Core>,
1240 aging_cap: usize,
1241 /// Consecutive control dequeues since the last user dequeue (the aging
1242 /// streak; never exceeds `aging_cap`).
1243 consec_control: usize,
1244 /// Consumer-side control pop ticket (plain counter, no atomic).
1245 ctl_head: usize,
1246 /// Block containing `ctl_head`; consumer-only. Blocks strictly behind
1247 /// it are reclaimed by `ControlLane::reclaim` as the cursor crosses
1248 /// boundaries; the cursor block itself is freed at lane drop.
1249 ctl_block: *mut CBlock<C>,
1250 /// One-shot latch: the `Received::UserLaneClosed` end-marker has been
1251 /// delivered. Plain consumer-side state — the consumer is single-owner,
1252 /// so no atomics.
1253 usr_closed_reported: bool,
1254}
1255
1256// SAFETY: `ctl_block` is dereferenced only by the consumer that owns it, and
1257// the pointee outlives the consumer (freed in `ControlLane::drop`, and the
1258// consumer holds an `Arc<ControlLane<C>>`).
1259unsafe impl<C: Send, U: Send> Send for Consumer<C, U> {}
1260
1261impl<C, U> Consumer<C, U> {
1262 /// Pop the next published control in ticket order. Consumer-only.
1263 #[inline]
1264 fn pop_control(&mut self) -> Option<C> {
1265 pop_control!(self)
1266 }
1267
1268 /// Receive the next item.
1269 ///
1270 /// Returns `None` once both lanes are closed and empty.
1271 ///
1272 /// POLICY: strict control-first priority with an anti-starvation aging
1273 /// cap. Control is dequeued before any waiting user (P1, overtake); after
1274 /// `aging_cap` consecutive control dequeues one waiting user is forced
1275 /// through (P3 under flood). Users reset the streak; a cap of 0 disables
1276 /// aging entirely.
1277 ///
1278 /// Once every `UserSender` is gone AND the ring is drained, the user
1279 /// stream's one-shot end-marker ([`Received::UserLaneClosed`]) is
1280 /// delivered through the same path a user item takes — control keeps its
1281 /// priority over it, aging forces it through exactly as for users, and
1282 /// it resets the streak (it IS the user stream's end-marker).
1283 #[cfg(not(loom))]
1284 pub async fn recv(&mut self) -> Option<Received<C, U>> {
1285 // Split borrows: lanes + streak + latch + control cursor are disjoint fields.
1286 let usr = &self.usr;
1287 let core = &self.core;
1288 let aging_cap = self.aging_cap;
1289 let streak = &mut self.consec_control;
1290 let leg_reported = &mut self.usr_closed_reported;
1291
1292 macro_rules! take_control {
1293 ($c:expr) => {{
1294 // Guarded increment: streak never exceeds the cap, so it
1295 // cannot overflow.
1296 if aging_cap != 0 && *streak < aging_cap {
1297 *streak += 1;
1298 }
1299 return Some(Received::Control($c));
1300 }};
1301 }
1302 macro_rules! take_user {
1303 ($u:expr) => {{
1304 *streak = 0;
1305 return Some(Received::User($u));
1306 }};
1307 }
1308 macro_rules! take_leg {
1309 () => {{
1310 // The user stream's end-marker takes the user slot: it
1311 // resets the aging streak exactly as a user item would.
1312 *leg_reported = true;
1313 *streak = 0;
1314 return Some(Received::UserLaneClosed);
1315 }};
1316 }
1317 // User-pop site, leg-aware: a drained ring whose lane is TERMINALLY
1318 // closed yields the one-shot end-marker instead of a plain miss.
1319 // `closed()` is monotone (a zero sender count can never rise again),
1320 // so a stale read can only delay the leg a loop turn, never fire it
1321 // early.
1322 macro_rules! pop_user_or_leg {
1323 () => {
1324 match usr.pop() {
1325 Some(u) => take_user!(u),
1326 None if !*leg_reported && usr.closed() => {
1327 // The `closed()` read (Acquire) synchronizes with
1328 // the last sender's `fetch_sub(Release)`, so every
1329 // item published before that decrement is visible
1330 // now — re-pop once before declaring the end-marker,
1331 // so a delivery racing the last drop cannot publish
1332 // an item AFTER it. A zero sender count is terminal
1333 // (the anchor's conditional increment refuses it), so
1334 // one re-pop exhausts the possibilities.
1335 match usr.pop() {
1336 Some(u) => take_user!(u),
1337 None => take_leg!(),
1338 }
1339 }
1340 None => {}
1341 }
1342 };
1343 }
1344
1345 loop {
1346 // Aging safety net: the streak reached the cap and a user is
1347 // waiting — serve it before any more control. The end-marker is
1348 // forced through here exactly as a user item.
1349 if aging_cap != 0 && *streak >= aging_cap {
1350 pop_user_or_leg!();
1351 }
1352 if let Some(c) = pop_control!(self) {
1353 take_control!(c);
1354 }
1355 pop_user_or_leg!();
1356
1357 let ctl_closed = self.ctl.closed();
1358 let usr_closed = usr.closed();
1359 if ctl_closed && usr_closed {
1360 // Both lanes closed: no item can arrive anymore. Sweep TWICE:
1361 // items pushed by departing senders are visible here via the
1362 // `senders` RMW release chain; under loom's (deliberately
1363 // over-approximating) visibility model a single read may
1364 // still come back stale, and a stale read cannot repeat
1365 // within one activation. The second pass is free on
1366 // hardware — it happens once per channel, at teardown.
1367 // The user pops are leg-aware: the end-marker (delivered
1368 // already if the user lane died first) is not re-fired, and
1369 // a stale `closed()` read that delayed it is caught here.
1370 if let Some(c) = pop_control!(self) {
1371 take_control!(c);
1372 }
1373 pop_user_or_leg!();
1374 if let Some(c) = pop_control!(self) {
1375 take_control!(c);
1376 }
1377 pop_user_or_leg!();
1378 return None;
1379 }
1380
1381 // Park. Register FIRST, announce parked (SeqCst — pairs with the
1382 // producers' SeqCst `parked` load), then re-check: any push or
1383 // last-sender drop after the re-check either wakes this
1384 // registration or is seen by the re-check itself.
1385 let notified = core.notify.notified();
1386 tokio::pin!(notified);
1387 notified.as_mut().enable();
1388 core.parked.store(true, Ordering::SeqCst);
1389
1390 if let Some(c) = pop_control!(self) {
1391 core.parked.store(false, Ordering::Relaxed);
1392 take_control!(c);
1393 }
1394 match usr.pop() {
1395 Some(u) => {
1396 core.parked.store(false, Ordering::Relaxed);
1397 take_user!(u);
1398 }
1399 // The leg must be observable from the PARKED path too: the
1400 // last-sender drop wakes this registration, and the
1401 // re-check delivers the end-marker instead of parking on a
1402 // dead user lane. Same re-pop discipline as
1403 // `pop_user_or_leg!`: the `closed()` read (Acquire)
1404 // synchronizes with the last sender's `fetch_sub(Release)`,
1405 // so an item published before that decrement is visible now
1406 // — re-pop once so a delivery racing the drop cannot land
1407 // an item AFTER the end-marker.
1408 None if !*leg_reported && usr.closed() => {
1409 core.parked.store(false, Ordering::Relaxed);
1410 if let Some(u) = usr.pop() {
1411 take_user!(u);
1412 }
1413 take_leg!();
1414 }
1415 None => {}
1416 }
1417 if self.ctl.closed() && usr.closed() {
1418 core.parked.store(false, Ordering::Relaxed);
1419 continue; // loop around to the sweep-and-`None` path
1420 }
1421 // Release one parked producer before sleeping: with the strided
1422 // in-pop `waiting` check, this is what guarantees a producer
1423 // parked on a just-drained ring is always woken (its push then
1424 // wakes us via `parked`).
1425 usr.release_one_waiter();
1426 notified.await;
1427 core.parked.store(false, Ordering::Relaxed);
1428 }
1429 }
1430
1431 /// Blocking twin of `recv`, compiled ONLY under `cfg(loom)`: the same
1432 /// policy loop and the same `parked`-flag Dekker protocol, with the
1433 /// sync `Notify` stand-in's `wait()` in place of the `Notified`
1434 /// future. Drives the loom model in `tests/loom.rs`.
1435 #[cfg(loom)]
1436 #[doc(hidden)]
1437 pub fn recv_blocking(&mut self) -> Option<Received<C, U>> {
1438 let usr = &self.usr;
1439 let core = &self.core;
1440 let aging_cap = self.aging_cap;
1441 let streak = &mut self.consec_control;
1442 let leg_reported = &mut self.usr_closed_reported;
1443
1444 macro_rules! take_control {
1445 ($c:expr) => {{
1446 if aging_cap != 0 && *streak < aging_cap {
1447 *streak += 1;
1448 }
1449 return Some(Received::Control($c));
1450 }};
1451 }
1452 macro_rules! take_user {
1453 ($u:expr) => {{
1454 *streak = 0;
1455 return Some(Received::User($u));
1456 }};
1457 }
1458 macro_rules! take_leg {
1459 () => {{
1460 *leg_reported = true;
1461 *streak = 0;
1462 return Some(Received::UserLaneClosed);
1463 }};
1464 }
1465 // See the async `recv`: a drained, terminally closed user lane
1466 // yields its one-shot end-marker at every user-pop site.
1467 macro_rules! pop_user_or_leg {
1468 () => {
1469 match usr.pop() {
1470 Some(u) => take_user!(u),
1471 None if !*leg_reported && usr.closed() => {
1472 // The `closed()` read (Acquire) synchronizes with
1473 // the last sender's `fetch_sub(Release)`, so every
1474 // item published before that decrement is visible
1475 // now — re-pop once before declaring the end-marker,
1476 // so a delivery racing the last drop cannot publish
1477 // an item AFTER it. A zero sender count is terminal
1478 // (the anchor's conditional increment refuses it), so
1479 // one re-pop exhausts the possibilities.
1480 match usr.pop() {
1481 Some(u) => take_user!(u),
1482 None => take_leg!(),
1483 }
1484 }
1485 None => {}
1486 }
1487 };
1488 }
1489
1490 loop {
1491 if aging_cap != 0 && *streak >= aging_cap {
1492 pop_user_or_leg!();
1493 }
1494 if let Some(c) = pop_control!(self) {
1495 take_control!(c);
1496 }
1497 pop_user_or_leg!();
1498
1499 let ctl_closed = self.ctl.closed();
1500 let usr_closed = usr.closed();
1501 if ctl_closed && usr_closed {
1502 // Both lanes closed. Synchronize with the departing
1503 // senders' `notify_if` locks so their publishes are visible
1504 // to the sweep (loom's visibility model does not honor the
1505 // atomic release chain alone); then sweep TWICE (see the
1506 // async `recv` for why). The sync point also publishes the
1507 // final `senders` decrement, so the leg-aware user pops
1508 // below observe the closure.
1509 core.notify.sync_point();
1510 if let Some(c) = pop_control!(self) {
1511 take_control!(c);
1512 }
1513 pop_user_or_leg!();
1514 if let Some(c) = pop_control!(self) {
1515 take_control!(c);
1516 }
1517 pop_user_or_leg!();
1518 return None;
1519 }
1520
1521 // Park: announce parked FIRST (SeqCst), then register — the
1522 // stand-in's `enable` lock is the synchronization point that
1523 // publishes the announcement to `notify_if` (and vice versa
1524 // for the re-checks below).
1525 core.parked.store(true, Ordering::SeqCst);
1526 core.notify.enable();
1527
1528 if let Some(c) = pop_control!(self) {
1529 core.parked.store(false, Ordering::Relaxed);
1530 take_control!(c);
1531 }
1532 match usr.pop() {
1533 Some(u) => {
1534 core.parked.store(false, Ordering::Relaxed);
1535 take_user!(u);
1536 }
1537 // The leg must be observable from the PARKED path too: the
1538 // last-sender drop wakes this registration, and the
1539 // re-check delivers the end-marker instead of parking on a
1540 // dead user lane. Same re-pop discipline as
1541 // `pop_user_or_leg!` (see the async `recv`).
1542 None if !*leg_reported && usr.closed() => {
1543 core.parked.store(false, Ordering::Relaxed);
1544 if let Some(u) = usr.pop() {
1545 take_user!(u);
1546 }
1547 take_leg!();
1548 }
1549 None => {}
1550 }
1551 if self.ctl.closed() && usr.closed() {
1552 core.parked.store(false, Ordering::Relaxed);
1553 continue;
1554 }
1555 usr.release_one_waiter();
1556 core.notify.wait();
1557 core.parked.store(false, Ordering::Relaxed);
1558 }
1559 }
1560
1561 /// Receive the next CONTROL item only, never consuming the user lane.
1562 /// `None` once the control lane is closed and drained. Cancel-safe under
1563 /// the same registration protocol as `recv`; a user-lane push may cause a
1564 /// spurious wake (one extra loop turn), never a lost wakeup and never a
1565 /// consumed user item.
1566 ///
1567 /// This path consumes no user items, so it must NOT call
1568 /// `usr.release_one_waiter()`: releasing a producer parked on a full ring
1569 /// would just re-park it — and the semantics WANT user producers to stay
1570 /// parked while the process is parked-for-restart (that is the
1571 /// backpressure). The aging streak (`consec_control`) is untouched.
1572 #[cfg(not(loom))]
1573 pub async fn recv_control(&mut self) -> Option<C> {
1574 let core = &self.core;
1575 loop {
1576 if let Some(c) = pop_control!(self) {
1577 return Some(c);
1578 }
1579 if self.ctl.closed() {
1580 // Control lane closed: no item can arrive anymore. Sweep
1581 // TWICE — the same discipline `recv` documents for the
1582 // both-closed path (departing control senders' publishes).
1583 if let Some(c) = pop_control!(self) {
1584 return Some(c);
1585 }
1586 if let Some(c) = pop_control!(self) {
1587 return Some(c);
1588 }
1589 return None;
1590 }
1591
1592 // Park with the identical enable-then-recheck `parked` protocol
1593 // as `recv`, minus all user-lane interaction.
1594 let notified = core.notify.notified();
1595 tokio::pin!(notified);
1596 notified.as_mut().enable();
1597 core.parked.store(true, Ordering::SeqCst);
1598
1599 if let Some(c) = pop_control!(self) {
1600 core.parked.store(false, Ordering::Relaxed);
1601 return Some(c);
1602 }
1603 if self.ctl.closed() {
1604 core.parked.store(false, Ordering::Relaxed);
1605 continue; // loop around to the sweep-and-`None` path
1606 }
1607 notified.await;
1608 core.parked.store(false, Ordering::Relaxed);
1609 }
1610 }
1611
1612 /// The `cfg(loom)` blocking twin (drives the loom model).
1613 #[cfg(loom)]
1614 #[doc(hidden)]
1615 pub fn recv_control_blocking(&mut self) -> Option<C> {
1616 let core = &self.core;
1617 loop {
1618 if let Some(c) = pop_control!(self) {
1619 return Some(c);
1620 }
1621 if self.ctl.closed() {
1622 // Same double-sweep discipline as `recv_blocking`'s
1623 // both-closed path: synchronize with the departing control
1624 // senders' `notify_if` locks, then sweep TWICE.
1625 core.notify.sync_point();
1626 if let Some(c) = pop_control!(self) {
1627 return Some(c);
1628 }
1629 if let Some(c) = pop_control!(self) {
1630 return Some(c);
1631 }
1632 return None;
1633 }
1634
1635 // The identical announce-then-register `parked` protocol as
1636 // `recv_blocking`, minus all user-lane interaction.
1637 core.parked.store(true, Ordering::SeqCst);
1638 core.notify.enable();
1639
1640 if let Some(c) = pop_control!(self) {
1641 core.parked.store(false, Ordering::Relaxed);
1642 return Some(c);
1643 }
1644 if self.ctl.closed() {
1645 core.parked.store(false, Ordering::Relaxed);
1646 continue;
1647 }
1648 core.notify.wait();
1649 core.parked.store(false, Ordering::Relaxed);
1650 }
1651 }
1652
1653 /// Consume the consumer and return everything still queued on both lanes,
1654 /// in FIFO order (P6 teardown seam).
1655 ///
1656 /// # Teardown race
1657 ///
1658 /// A send's linearization point is the ring-slot publish inside
1659 /// `UserLane::try_push`. A `UserSender::send` parked on a full lane at
1660 /// this moment is RELEASED with its payload: it never linearized, so it
1661 /// resolves `Err(UserClosed(item))` and the item comes back to the
1662 /// caller exactly once. A send whose publish completed before it
1663 /// observed teardown resolves `Ok(())`; its item is owned by the lane —
1664 /// collected here if still queued, or dropped with the lane otherwise.
1665 /// Pinned by `tests/teardown_oracle.rs` and
1666 /// `edge_cases::drain_teardown_race_releases_blocked_sender_with_payload`.
1667 #[must_use]
1668 #[allow(
1669 clippy::cast_possible_truncation,
1670 reason = "ticket arithmetic intentionally operates modulo 2^32; channel() bounds the live window"
1671 )]
1672 pub fn drain(mut self) -> Drained<C, U> {
1673 self.teardown();
1674 let mut control: Vec<C> = Vec::new();
1675 while let Some(c) = self.pop_control() {
1676 control.push(c);
1677 }
1678 let mut user = Vec::new();
1679 let mut head = self.usr.head.load(Ordering::Relaxed);
1680 let tail = self.usr.tail.load(Ordering::Acquire);
1681 while head < tail {
1682 // SAFETY: `head & mask` is always in bounds.
1683 let slot = unsafe { self.usr.ring.get_unchecked(head & self.usr.mask()) };
1684 if slot.seq.load(Ordering::Acquire) != (head as u32).wrapping_add(1) {
1685 break;
1686 }
1687 // SAFETY: published and unconsumed; the consumer is being torn
1688 // down, so no other popper exists.
1689 user.push(unsafe { slot.val.get().read().assume_init() });
1690 head = head.wrapping_add(1);
1691 }
1692 self.usr.head.store(head, Ordering::Relaxed);
1693 // Republish the consumed ticket past everything collected above, so
1694 // the lane's `Drop` does not reclaim items we just moved out.
1695 self.ctl.consumed.store(self.ctl_head, Ordering::Release);
1696 Drained { control, user }
1697 }
1698
1699 /// Idempotent teardown: flag the consumer gone, publish the consumed
1700 /// control ticket so the lane's `Drop` reclaims exactly the unconsumed
1701 /// tail, then release parked user senders (their pending park resolves
1702 /// `Err(UserClosed(item))` — the item never linearized, so it returns
1703 /// to the caller) and wake any parked `recv`.
1704 fn teardown(&self) {
1705 self.ctl.consumed.store(self.ctl_head, Ordering::Release);
1706 if self.core.consumer_gone.swap(true, Ordering::AcqRel) {
1707 return;
1708 }
1709 self.usr.send_notify.notify_waiters();
1710 self.core.notify.notify_one();
1711 }
1712}
1713
1714impl<C, U> Drop for Consumer<C, U> {
1715 fn drop(&mut self) {
1716 self.teardown();
1717 }
1718}
1719
1720/// Build a two-lane priority channel: an unbounded control lane and a bounded
1721/// user lane feeding one [`Consumer`].
1722///
1723/// # Panics
1724///
1725/// Panics when the rounded user capacity exceeds the live `u32` ticket window.
1726#[must_use]
1727#[allow(
1728 clippy::cast_possible_truncation,
1729 reason = "capacity is asserted below the u32 ticket window before slot tickets are initialized"
1730)]
1731pub fn channel<C, U>(cfg: Config) -> (ControlSender<C>, UserSender<U>, Consumer<C, U>) {
1732 // A capacity-0 (rendezvous) config is served by buffer slots: the FIFO
1733 // and no-loss semantics are identical, and the send side still paces on
1734 // the consumer taking items. The ring rounds capacity UP to a power of
1735 // two with a floor of 2 — the Vyukov ticket for "free next cycle"
1736 // (`pos + ring_len`) must be distinct from "published" (`pos + 1`),
1737 // which collapses at ring_len == 1 — and the rounded size IS the
1738 // effective capacity (in-ring backpressure).
1739 let capacity = cfg.user_capacity.max(2).next_power_of_two();
1740 assert!(
1741 capacity < (1 << 31),
1742 "user_capacity {capacity} exceeds the u32 ticket window"
1743 );
1744 let ring_len = capacity;
1745 let ring = (0..ring_len)
1746 .map(|i| Slot {
1747 seq: AtomicU32::new(i as u32),
1748 val: std::cell::UnsafeCell::new(MaybeUninit::uninit()),
1749 })
1750 .collect::<Vec<_>>()
1751 .into_boxed_slice();
1752 let core = Arc::new(Core {
1753 notify: Notify::new(),
1754 parked: AtomicBool::new(false),
1755 consumer_gone: AtomicBool::new(false),
1756 });
1757 let first_block = Box::into_raw(CBlock::new(0));
1758 let ctl = Arc::new(ControlLane {
1759 tail: AtomicUsize::new(0),
1760 tail_block: AtomicPtr::new(first_block),
1761 first_block: AtomicPtr::new(first_block),
1762 free_anchor: AtomicPtr::new(first_block),
1763 in_flight: AtomicUsize::new(0),
1764 consumed: AtomicUsize::new(0),
1765 senders: AtomicUsize::new(1),
1766 core: core.clone(),
1767 });
1768 let usr = Arc::new(UserLane {
1769 ring,
1770 tail: AtomicUsize::new(0),
1771 head: AtomicUsize::new(0),
1772 send_notify: Notify::new(),
1773 waiting: AtomicUsize::new(0),
1774 senders: AtomicUsize::new(1),
1775 core: core.clone(),
1776 });
1777 (
1778 ControlSender { lane: ctl.clone() },
1779 UserSender { lane: usr.clone() },
1780 Consumer {
1781 ctl,
1782 usr,
1783 core,
1784 aging_cap: cfg.aging_cap,
1785 consec_control: 0,
1786 ctl_head: 0,
1787 ctl_block: first_block,
1788 usr_closed_reported: false,
1789 },
1790 )
1791}
1792
1793/// Build an actor-oriented channel with affine admission ownership.
1794///
1795/// The returned [`MailboxOwner`] is the sole durable user-lane sender.
1796/// External addresses are restricted [`MailboxRef`]s and therefore neither
1797/// keep the mailbox open, retain admission permits, nor resurrect the mailbox
1798/// after retirement. Call
1799/// [`MailboxOwner::close_admission`], continue receiving until
1800/// [`Received::UserLaneClosed`], and finally drop the control sender when no
1801/// further lifecycle traffic can arrive; [`Consumer::recv`] then drains to
1802/// `None`.
1803pub fn mailbox_channel<C, U>(
1804 cfg: Config,
1805) -> (
1806 ControlSender<C>,
1807 MailboxOwner<U>,
1808 MailboxRef<U>,
1809 Consumer<C, U>,
1810) {
1811 let (control, sender, consumer) = channel(cfg);
1812 let owner = MailboxOwner {
1813 admission: Arc::new(Mutex::new(MailboxAdmission::Open(sender))),
1814 };
1815 let actor_ref = owner.actor_ref();
1816 (control, owner, actor_ref, consumer)
1817}
1818
1819#[cfg(test)]
1820mod mailbox_admission_tests {
1821 use super::{Config, TrySendError, mailbox_channel};
1822
1823 #[test]
1824 fn promoted_reference_cannot_extend_owner_admission() {
1825 let observe_closed_admission = || {
1826 let (control, owner, mailbox, receiver) = mailbox_channel::<(), u32>(Config::new(2));
1827 // Model a reference paused after promoting allocation liveness,
1828 // before taking the admission lock. Promotion is not authority.
1829 let promoted = owner.admission.clone();
1830 owner.close_admission();
1831 let rejected = mailbox.try_send(23);
1832 match rejected {
1833 Err(TrySendError::Closed(payload)) => assert_eq!(payload, 23),
1834 Err(TrySendError::Full(payload)) => {
1835 panic!("closed admission reported full: {payload}")
1836 }
1837 Ok(()) => panic!("promoted allocation retained admission after owner close"),
1838 }
1839 drop((promoted, control, receiver));
1840 };
1841 #[cfg(loom)]
1842 loom::model(observe_closed_admission);
1843 #[cfg(not(loom))]
1844 observe_closed_admission();
1845 }
1846}