1use behavior::{
4 Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
5 SendEffects, User,
6};
7use thiserror::Error;
8
9use crate::DeliveryRoute;
10
11pub struct TopicMembership<K, Route> {
13 pub topic: K,
15 pub subscribers: Vec<Route>,
17}
18
19pub enum PubSubMessage<K, P, Route> {
21 Subscribe {
23 topic: K,
25 subscriber: Route,
27 },
28 Unsubscribe {
30 topic: K,
32 subscriber: Route,
34 },
35 Publish {
37 topic: K,
39 value: P,
41 },
42}
43
44#[derive(Error, PartialEq, Eq)]
46pub enum PubSubError<K, P, Route> {
47 #[error("pub-sub topic is unknown")]
49 UnknownTopic {
50 topic: K,
52 subscriber: Route,
54 },
55 #[error("recipient is not subscribed to the pub-sub topic")]
57 NotSubscribed {
58 topic: K,
60 subscriber: Route,
62 },
63 #[error("pub-sub topic has no subscribers")]
65 NoSubscribers {
66 topic: K,
68 value: P,
70 },
71}
72
73impl<K: core::fmt::Debug, P: core::fmt::Debug, Route> core::fmt::Debug
74 for PubSubError<K, P, Route>
75{
76 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
77 match self {
78 Self::UnknownTopic { topic, .. } => formatter
79 .debug_struct("UnknownTopic")
80 .field("topic", topic)
81 .field("subscriber", &"<typed recipient>")
82 .finish(),
83 Self::NotSubscribed { topic, .. } => formatter
84 .debug_struct("NotSubscribed")
85 .field("topic", topic)
86 .field("subscriber", &"<typed recipient>")
87 .finish(),
88 Self::NoSubscribers { topic, value } => formatter
89 .debug_struct("NoSubscribers")
90 .field("topic", topic)
91 .field("value", value)
92 .finish(),
93 }
94 }
95}
96
97pub struct PubSub<A: Address, K, P, Route> {
111 topics: Vec<TopicMembership<K, Route>>,
112 marker: core::marker::PhantomData<fn() -> (A, P)>,
113}
114
115impl<A, K, P, Route> PubSub<A, K, P, Route>
116where
117 A: Address,
118 K: Eq,
119 Route: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = P>> + PartialEq,
120{
121 #[must_use]
123 pub const fn new() -> Self {
124 Self {
125 topics: Vec::new(),
126 marker: core::marker::PhantomData,
127 }
128 }
129
130 #[must_use]
132 pub fn topics(&self) -> &[TopicMembership<K, Route>] {
133 &self.topics
134 }
135
136 fn subscribe(&mut self, topic: K, subscriber: Route) {
137 if let Some(membership) = self.topics.iter_mut().find(|entry| entry.topic == topic) {
138 if !membership.subscribers.contains(&subscriber) {
139 membership.subscribers.push(subscriber);
140 }
141 } else {
142 self.topics.push(TopicMembership {
143 topic,
144 subscribers: vec![subscriber],
145 });
146 }
147 }
148
149 fn unsubscribe(&mut self, topic: K, subscriber: Route) -> Result<(), PubSubError<K, P, Route>> {
150 let Some(membership) = self.topics.iter_mut().find(|entry| entry.topic == topic) else {
151 return Err(PubSubError::UnknownTopic { topic, subscriber });
152 };
153 let Some(index) = membership
154 .subscribers
155 .iter()
156 .position(|member| member == &subscriber)
157 else {
158 return Err(PubSubError::NotSubscribed { topic, subscriber });
159 };
160 membership.subscribers.remove(index);
161 Ok(())
162 }
163}
164
165impl<A, K, P, Route> Default for PubSub<A, K, P, Route>
166where
167 A: Address,
168 K: Eq,
169 Route: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = P>> + PartialEq,
170{
171 fn default() -> Self {
172 Self::new()
173 }
174}
175
176impl<A, K, P, Route> BehaviorBase for PubSub<A, K, P, Route>
177where
178 A: Address,
179{
180 type Base = Self;
181 fn base(&self) -> &Self {
182 self
183 }
184}
185
186impl<A: Address, K, P, Route> behavior::Protocol for PubSub<A, K, P, Route> {
187 type Addr = A;
188 type Msg = PubSubMessage<K, P, Route>;
189}
190
191impl<A, K, P, Route> Behavior for PubSub<A, K, P, Route>
192where
193 A: Address,
194 K: Clone + Eq,
195 P: Clone,
196 Route: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = P>> + Clone + PartialEq,
197 Route::Sends: behavior::SendsFor<User<A, PubSubMessage<K, P, Route>>>,
198{
199 type Protocol = Self;
200 type Event = User<A, behavior::BehaviorMessage<Self>>;
201 type Sends = Route::Sends;
202 type Ph = Never;
203 type Error = PubSubError<K, P, Route>;
204 type Birth = NoBirths;
205 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
206 match event.message {
207 PubSubMessage::Subscribe { topic, subscriber } => {
208 self.subscribe(topic, subscriber);
209 Ok(Actions::cont())
210 }
211 PubSubMessage::Unsubscribe { topic, subscriber } => {
212 self.unsubscribe(topic, subscriber)?;
213 Ok(Actions::cont())
214 }
215 PubSubMessage::Publish { topic, value } => {
216 let Some(membership) = self.topics.iter().find(|entry| entry.topic == topic) else {
217 return Err(PubSubError::NoSubscribers { topic, value });
218 };
219 if membership.subscribers.is_empty() {
220 return Err(PubSubError::NoSubscribers { topic, value });
221 }
222 let mut sends = Route::Sends::empty();
223 for subscriber in &membership.subscribers {
224 sends.append(subscriber.clone().deliver(value.clone()));
225 }
226 Ok(Actions::send(sends))
227 }
228 }
229 }
230}
231
232#[cfg(test)]
233mod tests {
234 use super::*;
235 use crate::Activate as _;
236 use behavior::Recipient;
237 use behavior::{Delivery, MailAddr};
238 struct Destination;
239 impl behavior::Protocol for Destination {
240 type Addr = MailAddr;
241 type Msg = u8;
242 }
243
244 type Subject = PubSub<MailAddr, u8, u8, Recipient<Destination>>;
245 #[test]
246 fn topics_and_subscribers_preserve_first_order() {
247 let one = Recipient::global(MailAddr(1));
248 let two = Recipient::global(MailAddr(2));
249 let mut s = (Subject::new()).initialize().unwrap().behavior;
250 for subscriber in [one, two, one] {
251 let subscribed = s
252 .receive(
253 MailAddr(9),
254 PubSubMessage::Subscribe {
255 topic: 7,
256 subscriber,
257 },
258 )
259 .unwrap();
260 assert!(subscribed.sends.is_empty());
261 assert!(subscribed.creates.is_empty());
262 assert_eq!(subscribed.become_, behavior::Step::Continue);
263 }
264 let a = s
265 .receive(MailAddr(9), PubSubMessage::Publish { topic: 7, value: 4 })
266 .unwrap();
267 assert!(a.sends == vec![Delivery::new(one, 4), Delivery::new(two, 4)]);
268 }
269 #[test]
270 fn empty_known_and_unknown_topics_return_publication() {
271 let one = Recipient::global(MailAddr(1));
272 let mut s = (Subject::new()).initialize().unwrap().behavior;
273 let subscribed = s
274 .receive(
275 MailAddr(9),
276 PubSubMessage::Subscribe {
277 topic: 1,
278 subscriber: one,
279 },
280 )
281 .unwrap();
282 assert!(subscribed.sends.is_empty());
283 assert!(subscribed.creates.is_empty());
284 assert_eq!(subscribed.become_, behavior::Step::Continue);
285 let unsubscribed = s
286 .receive(
287 MailAddr(9),
288 PubSubMessage::Unsubscribe {
289 topic: 1,
290 subscriber: one,
291 },
292 )
293 .unwrap();
294 assert!(unsubscribed.sends.is_empty());
295 assert!(unsubscribed.creates.is_empty());
296 assert_eq!(unsubscribed.become_, behavior::Step::Continue);
297 for topic in [1, 2] {
298 let rejection = s.receive(MailAddr(9), PubSubMessage::Publish { topic, value: 8 });
299 assert!(
300 matches!(rejection,Err(PubSubError::NoSubscribers{topic:returned,value:8}) if returned==topic)
301 );
302 }
303 }
304}