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};