pub use std::sync::mpsc::{RecvError, RecvTimeoutError, SendError, TryRecvError};
use crate::common::time::Instant;
#[must_use]
pub fn channel<T>() -> (Sender<T>, Receiver<T>) {
let (sender, receiver) = ::loom::sync::mpsc::channel();
(Sender(sender), Receiver(receiver))
}
#[derive_where::derive_where(Clone)]
pub struct Sender<T>(::loom::sync::mpsc::Sender<T>);
impl<T> Sender<T> {
pub fn send(&self, value: T) -> Result<(), SendError<T>> {
self.0.send(value)
}
}
pub struct Receiver<T>(::loom::sync::mpsc::Receiver<T>);
impl<T> Receiver<T> {
pub fn recv(&self) -> Result<T, RecvError> {
crate::no_block::forbid("mpsc::recv");
self.0.recv()
}
pub fn recv_timeout(&self, _deadline: Instant) -> Result<T, RecvTimeoutError> {
panic!("timed channel receives require the flash backend when loom is enabled");
}
pub fn try_recv(&self) -> Result<T, TryRecvError> {
self.0.try_recv()
}
delegate::delegate! {
to self {
#[expr(core::iter::from_fn(move || $.ok()))]
#[call(recv)]
pub fn iter(&self) -> impl Iterator<Item = T> + '_;
#[expr(core::iter::from_fn(move || $.ok()))]
#[call(try_recv)]
pub fn try_iter(&self) -> impl Iterator<Item = T> + '_;
}
}
}