Skip to main content

nexus_acto_rs/actor/dispatch/
bounded.rs

1use crate::actor::dispatch::default_mailbox::DefaultMailbox;
2use crate::actor::dispatch::mailbox_handle::MailboxHandle;
3use crate::actor::dispatch::mailbox_middleware::MailboxMiddlewareHandle;
4use crate::actor::dispatch::mailbox_producer::MailboxProducer;
5use crate::actor::dispatch::unbounded::UnboundedMailboxQueue;
6use crate::actor::message::MessageHandle;
7use crate::util::queue::MpscUnboundedChannelQueue;
8use crate::util::queue::RingQueue;
9use crate::util::queue::{QueueBase, QueueError, QueueReader, QueueSize, QueueWriter};
10use async_trait::async_trait;
11use std::fmt::Debug;
12
13#[derive(Debug, Clone)]
14pub struct BoundedMailboxQueue {
15  user_mailbox: RingQueue<MessageHandle>,
16  initial_capacity: usize,
17  dropping: bool,
18}
19
20impl BoundedMailboxQueue {
21  pub(crate) fn new(user_mailbox: RingQueue<MessageHandle>, initial_capacity: usize, dropping: bool) -> Self {
22    BoundedMailboxQueue {
23      user_mailbox,
24      initial_capacity,
25      dropping,
26    }
27  }
28}
29
30#[async_trait]
31impl QueueBase<MessageHandle> for BoundedMailboxQueue {
32  async fn len(&self) -> QueueSize {
33    self.user_mailbox.len().await
34  }
35
36  async fn capacity(&self) -> QueueSize {
37    self.user_mailbox.capacity().await
38  }
39}
40
41#[async_trait]
42impl QueueWriter<MessageHandle> for BoundedMailboxQueue {
43  async fn offer(&mut self, element: MessageHandle) -> Result<(), QueueError<MessageHandle>> {
44    let len = self.user_mailbox.len().await;
45    if self.dropping && len == QueueSize::Limited(self.initial_capacity) {
46      let _ = self.user_mailbox.poll().await;
47    }
48    self.user_mailbox.offer(element).await
49  }
50}
51
52#[async_trait]
53impl QueueReader<MessageHandle> for BoundedMailboxQueue {
54  async fn poll(&mut self) -> Result<Option<MessageHandle>, QueueError<MessageHandle>> {
55    self.user_mailbox.poll().await
56  }
57
58  async fn clean_up(&mut self) {
59    self.user_mailbox.clean_up().await
60  }
61}
62
63pub fn bounded_mailbox_creator_with_opts(
64  size: usize,
65  dropping: bool,
66  mailbox_stats: impl IntoIterator<Item = MailboxMiddlewareHandle> + Send + Sync,
67) -> MailboxProducer {
68  let cloned_mailbox_stats = mailbox_stats.into_iter().collect::<Vec<_>>();
69  MailboxProducer::new(move || {
70    let cloned_mailbox_stats = cloned_mailbox_stats.clone();
71    async move {
72      let user_queue = BoundedMailboxQueue::new(RingQueue::new(size), size, dropping);
73      let system_queue = UnboundedMailboxQueue::new(MpscUnboundedChannelQueue::new());
74      MailboxHandle::new(
75        DefaultMailbox::new(user_queue, system_queue)
76          .with_middlewares(cloned_mailbox_stats.clone())
77          .await,
78      )
79    }
80  })
81}
82
83pub fn bounded_mailbox_creator(size: usize, dropping: bool) -> MailboxProducer {
84  bounded_mailbox_creator_with_opts(size, dropping, [])
85}