moirai_core/channel/unified/
receiver.rs1use std::marker::PhantomData;
4use std::sync::Arc;
5
6use super::core::UnifiedChannel;
7use crate::channel::error::ChannelError;
8use crate::channel::stats::ChannelStatistics;
9
10pub struct UnifiedReceiver<T> {
12 pub(crate) channel: Arc<UnifiedChannel<T>>,
13 pub(crate) _phantom: PhantomData<T>,
14}
15
16impl<T> UnifiedReceiver<T> {
17 pub fn recv(&self) -> Result<T, ChannelError> {
19 self.channel.recv()
20 }
21
22 pub fn try_recv(&self) -> Result<T, ChannelError> {
27 self.channel.recv()
28 }
29
30 pub fn recv_batch(&self, max_count: usize) -> Vec<T> {
32 self.channel.recv_batch(max_count)
33 }
34
35 pub fn is_closed(&self) -> bool {
37 self.channel.is_closed()
38 }
39
40 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
55unsafe impl<T: Send> Send for UnifiedReceiver<T> {}
57unsafe impl<T: Send> Sync for UnifiedReceiver<T> {}