Skip to main content

moirai_core/channel/hybrid/
mod.rs

1//! Zero-copy hybrid channel for async/sync interop.
2//!
3//! Uses a lock-free ring buffer with memory barriers to ensure safe
4//! zero-copy communication between async and sync contexts.
5
6use crate::communication::RingBuffer;
7use std::marker::PhantomData;
8use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize};
9use std::sync::{Arc, Mutex};
10
11mod future;
12mod notify;
13mod recv;
14mod send;
15#[cfg(test)]
16mod tests;
17
18pub use self::future::RecvFuture;
19pub use self::recv::HybridReceiver;
20pub use self::send::HybridSender;
21
22/// Zero-sized factory for a zero-copy hybrid channel endpoint pair.
23///
24/// The sender and receiver own the shared ring and wake registries. The factory
25/// carries only its scalar parameter, so constructing endpoints never creates
26/// an unused channel-owner allocation or duplicate synchronization state.
27pub struct HybridChannel<T>(PhantomData<fn() -> T>);
28
29impl<T: Send> HybridChannel<T> {
30    /// Create a new hybrid channel with specified capacity
31    pub fn new(capacity: usize) -> (HybridSender<T>, HybridReceiver<T>) {
32        let ring = Arc::new(RingBuffer::new(capacity));
33        let parker = Arc::new(Mutex::new(Vec::new()));
34        let async_wakers = Arc::new(Mutex::new(Vec::new()));
35        let parked_count = Arc::new(AtomicUsize::new(0));
36        let waker_count = Arc::new(AtomicUsize::new(0));
37        let closed = Arc::new(AtomicBool::new(false));
38        let next_id = Arc::new(AtomicU64::new(0));
39
40        let sender = HybridSender {
41            ring: ring.clone(),
42            parker: parker.clone(),
43            async_wakers: async_wakers.clone(),
44            parked_count: parked_count.clone(),
45            waker_count: waker_count.clone(),
46            closed: closed.clone(),
47            _marker: PhantomData,
48        };
49
50        let receiver = HybridReceiver {
51            ring,
52            parker,
53            async_wakers,
54            parked_count,
55            waker_count,
56            closed,
57            next_id,
58            _marker: PhantomData,
59        };
60
61        (sender, receiver)
62    }
63}
64
65// SAFETY: `HybridSender`/`HybridReceiver` each hold an `Arc<RingBuffer<T>>`,
66// which is not auto-`Send` because `RingBuffer` is `!Sync` (a single-producer /
67// single-consumer structure). These manual impls assert the SPSC contract: the
68// channel hands out exactly one sender and one receiver, each touching a
69// disjoint end of the ring (producer_seq vs consumer_seq), so the single
70// producer and single consumer may live on different threads without ever
71// aliasing the same slot or sequence. Soundness therefore requires that neither
72// half is `Clone` — a second sender or receiver would create two producers or
73// two consumers racing the same end. The `spsc_send_invariant` guard below turns
74// any future `Clone` impl on a half into a compile error so this contract cannot
75// be silently broken.
76unsafe impl<T: Send> Send for HybridSender<T> {}
77unsafe impl<T: Send> Send for HybridReceiver<T> {}
78
79/// Compile-time guard protecting the SPSC `Send` soundness contract above.
80///
81/// The manual `unsafe impl Send` on each half is sound only while each half is
82/// the unique owner of its endpoint. Adding `Clone` to either would allow two
83/// producers/consumers on the `!Sync` ring — latent UB with no other compile
84/// error. This module statically asserts neither half is `Clone`, and includes a
85/// positive/negative control proving the detector itself is wired correctly.
86const _: () = {
87    use core::marker::PhantomData;
88
89    struct Probe<T>(PhantomData<T>);
90    trait CloneProbe {
91        const IS_CLONE: bool = false;
92    }
93    impl<T> CloneProbe for Probe<T> {}
94    impl<T: Clone> Probe<T> {
95        const IS_CLONE: bool = true;
96    }
97
98    struct NeverClone;
99
100    assert!(
101        Probe::<u32>::IS_CLONE,
102        "Clone detector failed positive control"
103    );
104    assert!(
105        !Probe::<NeverClone>::IS_CLONE,
106        "Clone detector failed negative control"
107    );
108
109    assert!(
110        !Probe::<HybridSender<u8>>::IS_CLONE,
111        "HybridSender must not be Clone: a second producer would race the SPSC ring"
112    );
113    assert!(
114        !Probe::<HybridReceiver<u8>>::IS_CLONE,
115        "HybridReceiver must not be Clone: a second consumer would race the SPSC ring"
116    );
117};