Skip to main content

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}