Skip to main content

MpmcChannel

Struct MpmcChannel 

Source
pub struct MpmcChannel<T> { /* private fields */ }
Expand description

Multi-Producer Multi-Consumer channel with bounded capacity Uses mutex-based implementation for simplicity and correctness.

The shared state lives directly in this struct rather than behind a field per Arc: MpmcSender/MpmcReceiver already share one handle through Arc<MpmcChannel<T>>, so a channel costs one allocation plus the bounded ring’s slot array. Every operation reaches the mutex, condvars, both waiter counters and the ring through that single handle, with no second indirection on the send/receive paths.

Implementations§

Source§

impl<T> MpmcChannel<T>

Source

pub fn new(capacity: Option<usize>) -> MpmcChannel<T>

Create a new MPMC channel with optional capacity

Source

pub fn unbounded() -> MpmcChannel<T>

Create an unbounded channel

Source

pub fn bounded(capacity: usize) -> MpmcChannel<T>

Create a bounded channel with given capacity

Source

pub fn channel(capacity: Option<usize>) -> (MpmcSender<T>, MpmcReceiver<T>)

Create a channel pair for ergonomic usage

Trait Implementations§

Source§

impl<T> Channel<T> for MpmcChannel<T>
where T: Send,

Source§

fn send(&self, value: T) -> Result<(), ChannelError>

Send a value, blocking if necessary
Source§

fn try_send(&self, value: T) -> Result<(), ChannelError>

Try to send without blocking
Source§

fn recv(&self) -> Result<T, ChannelError>

Receive a value, blocking if necessary
Source§

fn try_recv(&self) -> Result<T, ChannelError>

Try to receive without blocking
Source§

fn is_empty(&self) -> bool

Check if channel is empty
Source§

fn is_full(&self) -> bool

Check if channel is full
Source§

fn capacity(&self) -> Option<usize>

Get the capacity of the channel
Source§

fn send_batch(&self, values: Vec<T>) -> Result<usize, ChannelError>

Send multiple values in batch. Default sends each individually.
Source§

fn recv_batch(&self, max_count: usize) -> Vec<T>

Receive up to max_count values in batch. Default receives individually.
Source§

fn close(&self)

Close the channel. Default: no-op (channels that support close override).
Source§

fn is_closed(&self) -> bool

Check if channel is closed. Default: false.
Source§

fn len(&self) -> usize

Current number of buffered items. Default: 0.
Source§

fn stats(&self) -> Option<ChannelStatistics>

Return statistics if the channel tracks them.
Source§

impl<T> Consumer<T> for MpmcChannel<T>
where T: Send,

Source§

fn recv(&self) -> Result<T, ChannelError>

Receive a value, blocking until one arrives.
Source§

fn try_recv(&self) -> Result<T, ChannelError>

Receive a value, or report ChannelError::Empty rather than waiting.
Source§

fn is_empty(&self) -> bool

Whether a receive would have to wait.
Source§

impl<T> Producer<T> for MpmcChannel<T>
where T: Send,

Source§

fn send(&self, value: T) -> Result<(), ChannelError>

Send a value, blocking until there is room.
Source§

fn try_send(&self, value: T) -> Result<(), ChannelError>

Send a value, or report ChannelError::Full rather than waiting.
Source§

fn is_full(&self) -> bool

Whether a further send would have to wait.
Source§

fn capacity(&self) -> Option<usize>

Bounded capacity, or None when the channel is unbounded.

Auto Trait Implementations§

§

impl<T> !Freeze for MpmcChannel<T>

§

impl<T> !RefUnwindSafe for MpmcChannel<T>

§

impl<T> Send for MpmcChannel<T>
where (Mutex<MpmcState<T>>, Condvar, Condvar): Send, Option<LockFreeQueue<T>>: Send,

§

impl<T> Sync for MpmcChannel<T>
where (Mutex<MpmcState<T>>, Condvar, Condvar): Sync, Option<LockFreeQueue<T>>: Sync,

§

impl<T> Unpin for MpmcChannel<T>
where (Mutex<MpmcState<T>>, Condvar, Condvar): Unpin, Option<LockFreeQueue<T>>: Unpin,

§

impl<T> UnsafeUnpin for MpmcChannel<T>

§

impl<T> UnwindSafe for MpmcChannel<T>

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.