Skip to main content

bombay/mailbox/
communication.rs

1//! Bombay Communication-backed two-lane actor mailboxes.
2
3use communication::{
4    Config, Consumer, ControlClosed, ControlSender, Received, UserAnchor, UserClosed, UserSender,
5    channel,
6};
7
8use crate::{EventSender, EventSource};
9
10/// Configuration for bounded actor user mailboxes.
11///
12/// Bombay deliberately selects Bombay Communication's zero-aging mode:
13/// every waiting control event precedes every waiting user event. Sustained
14/// control traffic may therefore starve user traffic. A second priority
15/// protocol or fairness consumer is required before bombay can justify an
16/// aging policy.
17///
18/// User events are FIFO in Communication's admission-ticket order. Clones of
19/// one sender have no additional cross-producer ordering contract: concurrent
20/// producers may acquire admission tickets in either order, while each
21/// producer's sequential sends retain their order.
22///
23/// The user lane is bounded. Async sends apply backpressure while retaining
24/// their event and return that exact event if retirement closes the lane;
25/// bombay never silently drops an accepted event.
26#[derive(Debug, Clone, Copy)]
27pub struct MailboxConfig {
28    config: Config,
29}
30
31impl MailboxConfig {
32    /// Construct a bounded user input.
33    #[must_use]
34    pub const fn bounded(user_capacity: usize) -> Self {
35        Self {
36            config: Config::new(user_capacity),
37        }
38    }
39
40    /// Create one concrete sender/source pair with this configuration.
41    #[must_use]
42    pub(crate) fn create<E: Send>(&self) -> (MailboxSender<E>, MailboxReceiver<E>) {
43        let (control, user, consumer) = channel::<E, E>(self.config);
44        (MailboxSender { control, user }, MailboxReceiver(consumer))
45    }
46}
47
48/// Counting edge handle for bounded user delivery.
49#[doc(hidden)]
50pub struct MailboxSender<E> {
51    control: ControlSender<E>,
52    user: UserSender<E>,
53}
54
55impl<E> Clone for MailboxSender<E> {
56    fn clone(&self) -> Self {
57        Self {
58            control: self.control.clone(),
59            user: self.user.clone(),
60        }
61    }
62}
63
64impl<E> MailboxSender<E> {
65    /// Derive a non-owning user endpoint suitable for address registration.
66    #[must_use]
67    pub(crate) fn anchor(&self) -> MailboxAnchor<E> {
68        MailboxAnchor(self.user.anchor())
69    }
70
71    /// Submit one already-formed priority event without user backpressure.
72    pub(crate) fn send_control(&self, event: E) -> Result<(), ControlClosed<E>> {
73        self.control.send(event)
74    }
75}
76
77impl<E: Send> EventSender for MailboxSender<E> {
78    type Event = E;
79    type Error = UserClosed<E>;
80
81    async fn send(&self, event: E) -> Result<(), Self::Error> {
82        self.user.send(event).await
83    }
84}
85
86/// Non-owning address-table endpoint; it cannot keep the user lane alive.
87#[doc(hidden)]
88pub struct MailboxAnchor<E>(UserAnchor<E>);
89
90impl<E> Clone for MailboxAnchor<E> {
91    fn clone(&self) -> Self {
92        Self(self.0.clone())
93    }
94}
95
96impl<E: Send> EventSender for MailboxAnchor<E> {
97    type Event = E;
98    type Error = UserClosed<E>;
99
100    async fn send(&self, event: E) -> Result<(), Self::Error> {
101        self.0.send(event).await
102    }
103}
104
105/// Actor-owned bounded user event source.
106#[doc(hidden)]
107pub struct MailboxReceiver<E>(Consumer<E, E>);
108
109impl<E: Send> EventSource for MailboxReceiver<E> {
110    type Event = E;
111
112    async fn next(&mut self) -> Option<Self::Event> {
113        match self.0.recv().await {
114            Some(Received::Control(event) | Received::User(event)) => Some(event),
115            Some(Received::UserLaneClosed) | None => None,
116        }
117    }
118}
119
120#[cfg(test)]
121mod tests {
122    use crate::{EventSender, EventSource};
123    use tokio::sync::Barrier;
124    use tokio::task::yield_now;
125
126    use std::sync::Arc;
127
128    use super::MailboxConfig;
129
130    #[tokio::test]
131    async fn registry_anchor_does_not_keep_mailbox_alive() {
132        let (sender, mut source) = MailboxConfig::bounded(2).create::<u64>();
133        let anchor = sender.anchor();
134        drop(sender);
135        assert_eq!(source.next().await, None);
136        assert!(anchor.send(7).await.is_err());
137    }
138
139    #[tokio::test]
140    async fn queued_user_events_drain_before_lane_closure() {
141        let (sender, mut source) = MailboxConfig::bounded(2).create::<u64>();
142        sender.send(1).await.unwrap();
143        drop(sender);
144        assert_eq!(source.next().await, Some(1));
145        assert_eq!(source.next().await, None);
146    }
147
148    #[tokio::test]
149    async fn blocked_producer_recovers_payload_when_receiver_retires() {
150        let (sender, receiver) = MailboxConfig::bounded(1).create::<u64>();
151        sender.send(1).await.unwrap();
152        sender.send(2).await.unwrap();
153        let barrier = Arc::new(Barrier::new(2));
154        let blocked = tokio::spawn({
155            let sender = sender.clone();
156            let barrier = barrier.clone();
157            async move {
158                barrier.wait().await;
159                sender.send(3).await
160            }
161        });
162        barrier.wait().await;
163        yield_now().await;
164        assert!(!blocked.is_finished());
165
166        drop(receiver);
167        let error = blocked
168            .await
169            .unwrap()
170            .expect_err("receiver retirement rejects send");
171        assert_eq!(error.0, 3);
172    }
173
174    #[tokio::test]
175    async fn control_event_precedes_queued_user_events() {
176        let (sender, mut source) = MailboxConfig::bounded(2).create::<u64>();
177        sender.send(1).await.unwrap();
178        sender.send(2).await.unwrap();
179        sender.send_control(9).unwrap();
180
181        assert_eq!(source.next().await, Some(9));
182        assert_eq!(source.next().await, Some(1));
183        assert_eq!(source.next().await, Some(2));
184    }
185
186    #[tokio::test]
187    async fn zero_aging_drains_complete_control_backlog_before_waiting_user() {
188        let (sender, mut source) = MailboxConfig::bounded(1).create::<u64>();
189        sender.send(1).await.unwrap();
190        for control in 10..18 {
191            sender.send_control(control).unwrap();
192        }
193
194        for control in 10..18 {
195            assert_eq!(source.next().await, Some(control));
196        }
197        assert_eq!(source.next().await, Some(1));
198    }
199}