use crate::communication::RingBuffer;
use std::marker::PhantomData;
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize};
use std::sync::{Arc, Mutex};
mod future;
mod recv;
mod send;
#[cfg(test)]
mod tests;
pub use self::future::RecvFuture;
pub use self::recv::HybridReceiver;
pub use self::send::HybridSender;
pub struct HybridChannel<T>(PhantomData<fn() -> T>);
impl<T: Send> HybridChannel<T> {
pub fn new(capacity: usize) -> (HybridSender<T>, HybridReceiver<T>) {
let ring = Arc::new(RingBuffer::new(capacity));
let parker = Arc::new(Mutex::new(Vec::new()));
let async_wakers = Arc::new(Mutex::new(Vec::new()));
let parked_count = Arc::new(AtomicUsize::new(0));
let waker_count = Arc::new(AtomicUsize::new(0));
let closed = Arc::new(AtomicBool::new(false));
let next_id = Arc::new(AtomicU64::new(0));
let sender = HybridSender {
ring: ring.clone(),
parker: parker.clone(),
async_wakers: async_wakers.clone(),
parked_count: parked_count.clone(),
waker_count: waker_count.clone(),
closed: closed.clone(),
_marker: PhantomData,
};
let receiver = HybridReceiver {
ring,
parker,
async_wakers,
parked_count,
waker_count,
closed,
next_id,
_marker: PhantomData,
};
(sender, receiver)
}
}
unsafe impl<T: Send> Send for HybridSender<T> {}
unsafe impl<T: Send> Send for HybridReceiver<T> {}
const _: () = {
use core::marker::PhantomData;
struct Probe<T>(PhantomData<T>);
trait CloneProbe {
const IS_CLONE: bool = false;
}
impl<T> CloneProbe for Probe<T> {}
impl<T: Clone> Probe<T> {
const IS_CLONE: bool = true;
}
struct NeverClone;
assert!(
Probe::<u32>::IS_CLONE,
"Clone detector failed positive control"
);
assert!(
!Probe::<NeverClone>::IS_CLONE,
"Clone detector failed negative control"
);
assert!(
!Probe::<HybridSender<u8>>::IS_CLONE,
"HybridSender must not be Clone: a second producer would race the SPSC ring"
);
assert!(
!Probe::<HybridReceiver<u8>>::IS_CLONE,
"HybridReceiver must not be Clone: a second consumer would race the SPSC ring"
);
};