Skip to main content

behavior_actors/discovery/
pub_sub.rs

1//! Keyed typed publication and subscription membership.
2
3use behavior::{
4    Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
5    SendEffects, User,
6};
7use thiserror::Error;
8
9use crate::DeliveryRoute;
10
11/// One topic and its subscribers in delivery order.
12pub struct TopicMembership<K, Route> {
13    /// Application-defined topic identity.
14    pub topic: K,
15    /// Unique subscribers in subscription order.
16    pub subscribers: Vec<Route>,
17}
18
19/// Operations accepted by [`PubSub`].
20pub enum PubSubMessage<K, P, Route> {
21    /// Idempotently subscribe one typed recipient.
22    Subscribe {
23        /// Topic identity.
24        topic: K,
25        /// Subscriber.
26        subscriber: Route,
27    },
28    /// Remove one existing subscription.
29    Unsubscribe {
30        /// Topic identity.
31        topic: K,
32        /// Subscriber.
33        subscriber: Route,
34    },
35    /// Publish one value to a point-in-time membership snapshot.
36    Publish {
37        /// Topic identity.
38        topic: K,
39        /// Owned publication.
40        value: P,
41    },
42}
43
44/// Typed keyed-publication rejection.
45#[derive(Error, PartialEq, Eq)]
46pub enum PubSubError<K, P, Route> {
47    /// Unsubscription named an unknown topic.
48    #[error("pub-sub topic is unknown")]
49    UnknownTopic {
50        /// Rejected topic.
51        topic: K,
52        /// Exact subscriber from the rejected command.
53        subscriber: Route,
54    },
55    /// Recipient is not subscribed to the named topic.
56    #[error("recipient is not subscribed to the pub-sub topic")]
57    NotSubscribed {
58        /// Rejected topic.
59        topic: K,
60        /// Exact subscriber from the rejected command.
61        subscriber: Route,
62    },
63    /// Publication had no recipients; ownership is returned.
64    #[error("pub-sub topic has no subscribers")]
65    NoSubscribers {
66        /// Topic with no recipients.
67        topic: K,
68        /// Undelivered owned publication.
69        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
97/// Deterministic keyed typed publish/subscribe behavior.
98///
99/// Topics are introduced by the first subscription and retained after their
100/// final unsubscription so an unknown topic remains distinguishable from a
101/// known empty one. Subscription is idempotent. Unsubscription from an unknown
102/// topic or absent membership is rejected atomically. Publication clones the
103/// value once per point-in-time subscriber in subscription order; an unknown
104/// or empty topic returns the original value in [`PubSubError::NoSubscribers`].
105/// Initialization is empty, the behavior creates no actors, and it never
106/// terminates by policy. Membership retention and ordered snapshot delivery
107/// are Bombay policy; physical delivery remains a Communication capability.
108/// Cloning a publication can panic only if the application-provided `Clone`
109/// implementation panics, before any template state is changed.
110pub 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    /// Construct an empty keyed fabric.
122    #[must_use]
123    pub const fn new() -> Self {
124        Self {
125            topics: Vec::new(),
126            marker: core::marker::PhantomData,
127        }
128    }
129
130    /// Borrow complete topic membership in introduction order.
131    #[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}