nexus_acto_rs/actor/dispatch/
bounded.rs1use 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}