Skip to main content

behavior_actors/discovery/
topic.rs

1//! Typed subscription membership and publication.
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/// Commands accepted by a [`Topic`].
12pub enum TopicMessage<P, Route> {
13    /// Add a subscriber if absent.
14    Subscribe(Route),
15    /// Remove a subscriber if present.
16    Unsubscribe(Route),
17    /// Publish one owned value to the current membership snapshot.
18    Publish(P),
19}
20
21/// Publication rejection preserving the unaccepted value.
22#[derive(Debug, Error, Clone, PartialEq, Eq)]
23pub enum TopicError<P> {
24    /// No subscriber was present when publication was folded.
25    #[error("topic publication rejected because there are no subscribers")]
26    NoSubscribers(P),
27}
28
29/// Insertion-ordered typed publication behavior.
30///
31/// Subscription and unsubscription are idempotent and preserve survivor
32/// order. Publication snapshots current membership and emits exactly one typed
33/// delivery per subscriber in that order. With no subscribers it returns
34/// [`TopicError::NoSubscribers`] containing the owned value; it never silently
35/// drops a publication. Initialization is empty and the topic never terminates
36/// itself. Membership order and empty-publication rejection are Bombay policy;
37/// endpoint resolution and delivery are Address/Communication capabilities.
38/// No method has a semantic panic condition.
39pub struct Topic<A: Address, P, Route> {
40    subscribers: Vec<Route>,
41    marker: core::marker::PhantomData<fn() -> (A, P)>,
42}
43
44impl<A: Address, P, Route> Topic<A, P, Route> {
45    /// Construct an empty topic definition.
46    #[must_use]
47    pub const fn new() -> Self {
48        Self {
49            subscribers: Vec::new(),
50            marker: core::marker::PhantomData,
51        }
52    }
53
54    /// Borrow subscribers in publication order.
55    #[must_use]
56    pub fn subscribers(&self) -> &[Route] {
57        &self.subscribers
58    }
59}
60
61impl<A: Address, P, Route> Default for Topic<A, P, Route> {
62    fn default() -> Self {
63        Self::new()
64    }
65}
66
67impl<A, P, Route> BehaviorBase for Topic<A, P, Route>
68where
69    A: Address,
70{
71    type Base = Self;
72
73    fn base(&self) -> &Self {
74        self
75    }
76}
77
78impl<A: Address, P, Route> behavior::Protocol for Topic<A, P, Route> {
79    type Addr = A;
80    type Msg = TopicMessage<P, Route>;
81}
82
83impl<A, P, Route> Behavior for Topic<A, P, Route>
84where
85    A: Address,
86    P: Clone,
87    Route: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = P>> + Clone + PartialEq,
88    Route::Sends: behavior::SendsFor<User<A, TopicMessage<P, Route>>>,
89{
90    type Protocol = Self;
91    type Event = User<A, behavior::BehaviorMessage<Self>>;
92    type Sends = Route::Sends;
93    type Ph = Never;
94    type Error = TopicError<P>;
95    type Birth = NoBirths;
96
97    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
98        match event.message {
99            TopicMessage::Subscribe(subscriber) => {
100                if !self.subscribers.contains(&subscriber) {
101                    self.subscribers.push(subscriber);
102                }
103                Ok(Actions::cont())
104            }
105            TopicMessage::Unsubscribe(subscriber) => {
106                if let Some(index) = self
107                    .subscribers
108                    .iter()
109                    .position(|current| current == &subscriber)
110                {
111                    self.subscribers.remove(index);
112                }
113                Ok(Actions::cont())
114            }
115            TopicMessage::Publish(value) => {
116                if self.subscribers.is_empty() {
117                    return Err(TopicError::NoSubscribers(value));
118                }
119                let mut sends = Route::Sends::empty();
120                for subscriber in &self.subscribers {
121                    sends.append(subscriber.clone().deliver(value.clone()));
122                }
123                Ok(Actions::send(sends))
124            }
125        }
126    }
127}
128
129#[cfg(test)]
130mod tests {
131    use super::*;
132    use crate::Activate as _;
133    use behavior::Recipient;
134    use behavior::{Delivery, MailAddr, MessageProtocol};
135
136    #[test]
137    fn membership_is_idempotent_and_publication_ordered() {
138        let one = Recipient::from(MailAddr(1));
139        let two = Recipient::from(MailAddr(2));
140        let mut topic = Topic::<MailAddr, u8, Recipient<MessageProtocol<MailAddr, u8>>>::new()
141            .initialize()
142            .unwrap()
143            .behavior;
144        for subscriber in [one, two, one] {
145            let subscribed = topic
146                .receive(MailAddr(9), TopicMessage::Subscribe(subscriber))
147                .unwrap();
148            assert!(subscribed.sends.is_empty());
149            assert!(subscribed.creates.is_empty());
150            assert_eq!(subscribed.become_, behavior::Step::Continue);
151        }
152        assert!(topic.subscribers() == [one, two]);
153        let published = topic
154            .receive(MailAddr(9), TopicMessage::Publish(7))
155            .unwrap();
156        assert!(published.sends == vec![Delivery::new(one, 7), Delivery::new(two, 7)]);
157        let unsubscribed = topic
158            .receive(MailAddr(9), TopicMessage::Unsubscribe(one))
159            .unwrap();
160        assert!(unsubscribed.sends.is_empty());
161        assert!(unsubscribed.creates.is_empty());
162        assert_eq!(unsubscribed.become_, behavior::Step::Continue);
163        assert!(topic.subscribers() == [two]);
164    }
165
166    #[test]
167    fn empty_publication_returns_owned_value() {
168        let mut topic = Topic::<MailAddr, u8, Recipient<MessageProtocol<MailAddr, u8>>>::new()
169            .initialize()
170            .unwrap()
171            .behavior;
172        let rejection = topic.receive(MailAddr(9), TopicMessage::Publish(7));
173        assert!(matches!(rejection, Err(TopicError::NoSubscribers(7))));
174    }
175}