Skip to main content

moirai_core/channel/unified/
receiver.rs

1//! Receiver half of a unified channel.
2
3use std::marker::PhantomData;
4use std::sync::Arc;
5
6use super::core::UnifiedChannel;
7use crate::channel::error::ChannelError;
8use crate::channel::stats::ChannelStatistics;
9
10/// Receiver half of a unified channel
11pub struct UnifiedReceiver<T> {
12    pub(crate) channel: Arc<UnifiedChannel<T>>,
13    pub(crate) _phantom: PhantomData<T>,
14}
15
16impl<T> UnifiedReceiver<T> {
17    /// Receive a message (non-blocking; returns `Err(Empty)` when none is available).
18    pub fn recv(&self) -> Result<T, ChannelError> {
19        self.channel.recv()
20    }
21
22    /// Try to receive without blocking.
23    ///
24    /// Identical to [`Self::recv`] — this channel has no blocking receive
25    /// path. Retained for consumers written against the `try_recv` name.
26    pub fn try_recv(&self) -> Result<T, ChannelError> {
27        self.channel.recv()
28    }
29
30    /// Receive batch of messages
31    pub fn recv_batch(&self, max_count: usize) -> Vec<T> {
32        self.channel.recv_batch(max_count)
33    }
34
35    /// Check if channel is closed
36    pub fn is_closed(&self) -> bool {
37        self.channel.is_closed()
38    }
39
40    /// Get channel statistics
41    pub fn stats(&self) -> ChannelStatistics {
42        self.channel.stats()
43    }
44}
45
46impl<T> Clone for UnifiedReceiver<T> {
47    fn clone(&self) -> Self {
48        Self {
49            channel: self.channel.clone(),
50            _phantom: PhantomData,
51        }
52    }
53}
54
55// Safety: UnifiedReceiver is safe to send and share between threads
56unsafe impl<T: Send> Send for UnifiedReceiver<T> {}
57unsafe impl<T: Send> Sync for UnifiedReceiver<T> {}