Skip to main content

nexus_acto_rs/actor/dispatch/
unbounded.rs

1use async_trait::async_trait;
2
3use crate::actor::dispatch::default_mailbox::DefaultMailbox;
4use crate::actor::dispatch::mailbox_handle::MailboxHandle;
5use crate::actor::dispatch::mailbox_middleware::MailboxMiddlewareHandle;
6use crate::actor::dispatch::mailbox_producer::MailboxProducer;
7use crate::actor::message::MessageHandle;
8use crate::util::queue::MpscUnboundedChannelQueue;
9use crate::util::queue::PriorityQueue;
10use crate::util::queue::RingQueue;
11use crate::util::queue::{QueueBase, QueueError, QueueReader, QueueSize, QueueWriter};
12
13#[derive(Debug, Clone)]
14pub struct UnboundedMailboxQueue<Q: QueueReader<MessageHandle> + QueueWriter<MessageHandle>> {
15  user_mailbox: Q,
16}
17
18impl<Q: QueueReader<MessageHandle> + QueueWriter<MessageHandle>> UnboundedMailboxQueue<Q> {
19  pub fn new(user_mailbox: Q) -> Self {
20    UnboundedMailboxQueue { user_mailbox }
21  }
22}
23
24#[async_trait]
25impl<Q: QueueReader<MessageHandle> + QueueWriter<MessageHandle>> QueueBase<MessageHandle> for UnboundedMailboxQueue<Q> {
26  async fn len(&self) -> QueueSize {
27    self.user_mailbox.len().await
28  }
29
30  async fn capacity(&self) -> QueueSize {
31    self.user_mailbox.capacity().await
32  }
33}
34
35#[async_trait]
36impl<Q: QueueReader<MessageHandle> + QueueWriter<MessageHandle>> QueueReader<MessageHandle>
37  for UnboundedMailboxQueue<Q>
38{
39  async fn poll(&mut self) -> Result<Option<MessageHandle>, QueueError<MessageHandle>> {
40    self.user_mailbox.poll().await
41  }
42
43  async fn clean_up(&mut self) {
44    self.user_mailbox.clean_up().await
45  }
46}
47
48#[async_trait]
49impl<Q: QueueReader<MessageHandle> + QueueWriter<MessageHandle>> QueueWriter<MessageHandle>
50  for UnboundedMailboxQueue<Q>
51{
52  async fn offer(&mut self, element: MessageHandle) -> Result<(), QueueError<MessageHandle>> {
53    self.user_mailbox.offer(element).await
54  }
55}
56
57pub fn unbounded_mailbox_creator_with_opts(
58  mailbox_stats: impl IntoIterator<Item = MailboxMiddlewareHandle> + Send + Sync,
59) -> MailboxProducer {
60  let cloned_mailbox_stats = mailbox_stats.into_iter().collect::<Vec<_>>();
61  MailboxProducer::new(move || {
62    let cloned_mailbox_stats = cloned_mailbox_stats.clone();
63    async move {
64      let user_queue = UnboundedMailboxQueue::new(RingQueue::new(10));
65      let system_queue = UnboundedMailboxQueue::new(MpscUnboundedChannelQueue::new());
66      MailboxHandle::new(
67        DefaultMailbox::new(user_queue, system_queue)
68          .with_middlewares(cloned_mailbox_stats.clone())
69          .await,
70      )
71    }
72  })
73}
74
75pub fn unbounded_mailbox_creator() -> MailboxProducer {
76  unbounded_mailbox_creator_with_opts([])
77}
78
79pub fn unbounded_priority_mailbox_creator_with_opts(
80  mailbox_stats: impl IntoIterator<Item = MailboxMiddlewareHandle> + Send + Sync,
81) -> MailboxProducer {
82  let cloned_mailbox_stats = mailbox_stats.into_iter().collect::<Vec<_>>();
83  MailboxProducer::new(move || {
84    let cloned_mailbox_stats = cloned_mailbox_stats.clone();
85    async move {
86      let user_queue = UnboundedMailboxQueue::new(PriorityQueue::new(|| RingQueue::new(10)));
87      let system_queue = UnboundedMailboxQueue::new(MpscUnboundedChannelQueue::new());
88      MailboxHandle::new(
89        DefaultMailbox::new(user_queue, system_queue)
90          .with_middlewares(cloned_mailbox_stats.clone())
91          .await,
92      )
93    }
94  })
95}
96
97pub fn unbounded_priority_mailbox_creator() -> MailboxProducer {
98  unbounded_priority_mailbox_creator_with_opts([])
99}
100
101pub fn unbounded_mpsc_mailbox_creator_with_opts(
102  mailbox_stats: impl IntoIterator<Item = MailboxMiddlewareHandle> + Send + Sync,
103) -> MailboxProducer {
104  let cloned_mailbox_stats = mailbox_stats.into_iter().collect::<Vec<_>>();
105  MailboxProducer::new(move || {
106    let cloned_mailbox_stats = cloned_mailbox_stats.clone();
107    async move {
108      let user_queue = UnboundedMailboxQueue::new(MpscUnboundedChannelQueue::new());
109      let system_queue = UnboundedMailboxQueue::new(MpscUnboundedChannelQueue::new());
110      MailboxHandle::new(
111        DefaultMailbox::new(user_queue, system_queue)
112          .with_middlewares(cloned_mailbox_stats.clone())
113          .await,
114      )
115    }
116  })
117}
118
119pub fn unbounded_mpsc_mailbox_creator() -> MailboxProducer {
120  unbounded_mpsc_mailbox_creator_with_opts([])
121}