bombay/mailbox/
communication.rs1use communication::{
4 Config, Consumer, ControlClosed, ControlSender, Received, UserAnchor, UserClosed, UserSender,
5 channel,
6};
7
8use crate::{EventSender, EventSource};
9
10#[derive(Debug, Clone, Copy)]
27pub struct MailboxConfig {
28 config: Config,
29}
30
31impl MailboxConfig {
32 #[must_use]
34 pub const fn bounded(user_capacity: usize) -> Self {
35 Self {
36 config: Config::new(user_capacity),
37 }
38 }
39
40 #[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#[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 #[must_use]
67 pub(crate) fn anchor(&self) -> MailboxAnchor<E> {
68 MailboxAnchor(self.user.anchor())
69 }
70
71 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#[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#[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}