Skip to main content

moirai_core/channel/
select.rs

1//! `Select` helper and top-level channel constructor functions.
2
3use super::error::Result;
4use super::mpmc::{MpmcChannel, MpmcReceiver, MpmcSender};
5use super::spsc::{SpscChannel, SpscReceiver, SpscSender};
6
7/// Non-blocking poll over multiple receive closures.
8///
9/// This is a single-pass poll, not a Go-style blocking `select`: each closure
10/// is tried once in order and the first successful receive wins; when every
11/// closure fails the call returns `None` immediately without waiting or
12/// registering wakeups.
13pub struct Select;
14
15impl Select {
16    /// Poll each receiver closure once in order, returning the first
17    /// available value with its index, or `None` if none is ready.
18    pub fn try_recv<T>(receivers: &mut [&mut dyn FnMut() -> Result<T>]) -> Option<(usize, T)> {
19        for (idx, recv) in receivers.iter_mut().enumerate() {
20            if let Ok(value) = recv() {
21                return Some((idx, value));
22            }
23        }
24        None
25    }
26}
27
28/// Create a new SPSC channel pair with the given capacity.
29pub fn spsc<T>(capacity: usize) -> (SpscSender<T>, SpscReceiver<T>) {
30    SpscChannel::channel(capacity)
31}
32
33/// Create a new bounded MPMC channel with the given capacity.
34///
35/// This is the constructor to reach for. Callers with no rate information of
36/// their own should pass [`DEFAULT_CHANNEL_CAPACITY`].
37///
38/// [`DEFAULT_CHANNEL_CAPACITY`]: super::config::DEFAULT_CHANNEL_CAPACITY
39pub fn mpmc<T>(capacity: usize) -> (MpmcSender<T>, MpmcReceiver<T>) {
40    MpmcChannel::channel(Some(capacity))
41}
42
43/// Create a new MPMC channel whose queue grows without limit.
44///
45/// # Memory
46///
47/// This channel applies **no backpressure**: a producer that outruns its
48/// consumer grows the queue until allocation fails. It is deliberately absent
49/// from [`crate::prelude`] and from the crate root so that reaching for it is
50/// a decision rather than the path of least resistance — prefer [`mpmc`], and
51/// use this only where the total number of sends is bounded by
52/// construction (a fixed fan-in, a drained one-shot batch) and blocking a
53/// producer would deadlock.
54pub fn unbounded<T>() -> (MpmcSender<T>, MpmcReceiver<T>) {
55    MpmcChannel::channel(None)
56}