//! Unbounded MPSC channels. //! //! Inner state is `Arc>>` so channels can be sent across OS //! threads (required for the multi-scheduler runtime where a sender and //! receiver may run on different scheduler threads simultaneously). //! //! Semantics: //! - Senders are clonable; the last sender drop closes the channel. //! - `Receiver::recv` on an empty open channel parks the receiver. //! - `Receiver::recv` on an empty closed channel returns `Err(RecvError)`. //! - `Sender::send` on an open channel always succeeds. //! - `Sender::send` on a closed channel (receiver dropped) returns //! `Err(SendError(value))`. //! - When a send pushes to a previously empty queue and a receiver is //! parked, the receiver is unparked. use crate::pid::Pid; use std::collections::VecDeque; use std::sync::{Arc, Mutex}; pub fn channel() -> (Sender, Receiver) { let inner = Arc::new(Mutex::new(Inner { queue: VecDeque::new(), parked_receiver: None, senders: 1, receiver_alive: true, })); (Sender { inner: inner.clone() }, Receiver { inner }) } struct Inner { queue: VecDeque, parked_receiver: Option, senders: usize, receiver_alive: bool, } pub struct Sender { inner: Arc>>, } pub struct Receiver { inner: Arc>>, } #[derive(Debug, PartialEq, Eq)] pub struct SendError(pub T); #[derive(Debug, PartialEq, Eq, Clone, Copy)] pub struct RecvError; impl std::fmt::Display for RecvError { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { write!(f, "channel closed") } } impl std::error::Error for RecvError {} impl Clone for Sender { fn clone(&self) -> Self { self.inner.lock().unwrap().senders += 1; Sender { inner: self.inner.clone() } } } impl Drop for Sender { fn drop(&mut self) { let unpark = { let mut g = self.inner.lock().unwrap(); g.senders -= 1; // Wake the parked receiver on the last sender drop regardless of // whether the queue is empty. A plain `recv` only ever parks on an // empty queue (so this is unchanged for it), but a selective // `recv_match` may be parked on a *non-empty* queue holding only // non-matching messages — it must wake to observe closure and // return Err rather than sleep forever. if g.senders == 0 { g.parked_receiver.take() } else { None } }; if let Some(pid) = unpark { crate::scheduler::unpark(pid); } } } impl Drop for Receiver { fn drop(&mut self) { self.inner.lock().unwrap().receiver_alive = false; } } impl Sender { pub fn send(&self, value: T) -> Result<(), SendError> { let unpark = { let mut g = self.inner.lock().unwrap(); if !g.receiver_alive { return Err(SendError(value)); } g.queue.push_back(value); g.parked_receiver.take() }; if let Some(pid) = unpark { crate::te!(crate::trace::Event::Send { sender: crate::actor::current_pid().unwrap_or(crate::pid::Pid::new(u32::MAX, u32::MAX)), receiver: Some(pid) }); crate::scheduler::unpark(pid); } else { crate::te!(crate::trace::Event::Send { sender: crate::actor::current_pid().unwrap_or(crate::pid::Pid::new(u32::MAX, u32::MAX)), receiver: None }); } Ok(()) } } impl Receiver { pub fn recv(&self) -> Result { loop { { let mut g = self.inner.lock().unwrap(); if let Some(v) = g.queue.pop_front() { return Ok(v); } if g.senders == 0 { return Err(RecvError); } let me = crate::actor::current_pid() .expect("recv() called outside an actor"); debug_assert!( g.parked_receiver.is_none(), "channel has more than one receiver" ); g.parked_receiver = Some(me); crate::te!(crate::trace::Event::RecvPark(me)); } // Release the lock before parking — the unparker will need it. crate::scheduler::park_current(); // Woken up — record it before looping to check the queue. crate::te!(crate::trace::Event::RecvWake(crate::actor::current_pid().unwrap())); } } /// Selective receive: remove and return the first queued message for which /// `pred` holds, leaving the rest in arrival order. If no queued message /// matches, parks and re-scans on every send (a selective receiver may park /// on a *non-empty* queue). Returns `Err(RecvError)` only once the channel /// is closed and no queued message matches. /// /// `pred` is run while the channel lock is held: keep it cheap and pure, /// and do not call back into this channel from inside it. It is modelled as /// `Fn` (not `FnMut`) deliberately — it is re-run from scratch on every /// scan, so a stateful predicate would observe surprising re-counting. pub fn recv_match(&self, pred: F) -> Result where F: Fn(&T) -> bool, { loop { { let mut g = self.inner.lock().unwrap(); if let Some(i) = g.queue.iter().position(|v| pred(v)) { // position() found it, so remove() returns Some. return Ok(g.queue.remove(i).unwrap()); } if g.senders == 0 { // Closed and nothing queued can ever match. return Err(RecvError); } let me = crate::actor::current_pid() .expect("recv_match() called outside an actor"); debug_assert!( g.parked_receiver.is_none(), "channel has more than one receiver" ); g.parked_receiver = Some(me); crate::te!(crate::trace::Event::RecvPark(me)); } // Release the lock before parking — the unparker will need it. crate::scheduler::park_current(); crate::te!(crate::trace::Event::RecvWake(crate::actor::current_pid().unwrap())); } } /// Non-blocking selective receive. `Ok(Some(v))` if a queued message /// matched `pred` (removed, rest left in order), `Ok(None)` if the channel /// is open but nothing matched, `Err(RecvError)` if closed and nothing /// matched. Same predicate contract as [`recv_match`](Self::recv_match). pub fn try_recv_match(&self, pred: F) -> Result, RecvError> where F: Fn(&T) -> bool, { let mut g = self.inner.lock().unwrap(); if let Some(i) = g.queue.iter().position(|v| pred(v)) { return Ok(Some(g.queue.remove(i).unwrap())); } if g.senders == 0 { return Err(RecvError); } Ok(None) } /// Non-blocking. `Ok(Some(v))` if a message was available, `Ok(None)` if /// the channel is empty but open, `Err(RecvError)` if closed and drained. pub fn try_recv(&self) -> Result, RecvError> { let mut g = self.inner.lock().unwrap(); if let Some(v) = g.queue.pop_front() { return Ok(Some(v)); } if g.senders == 0 { return Err(RecvError); } Ok(None) } }