Skip to main content

behavior_actors/discovery/
registry.rs

1//! Typed recipient bindings and lookup replies.
2
3#[cfg(test)]
4use behavior::Delivery;
5use behavior::{
6    Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol, Recipient,
7    User,
8};
9use thiserror::Error;
10
11use crate::DeliveryRoute;
12
13/// The complete lookup result returned by [`Registry`].
14pub enum RegistryResult<K, D: Protocol> {
15    /// The key was bound to this typed recipient when the lookup was folded.
16    Found { key: K, recipient: Recipient<D> },
17    /// No binding existed when the lookup was folded.
18    Missing { key: K },
19}
20
21impl<K: Clone, D: Protocol> Clone for RegistryResult<K, D> {
22    fn clone(&self) -> Self {
23        match self {
24            Self::Found { key, recipient } => Self::Found {
25                key: key.clone(),
26                recipient: *recipient,
27            },
28            Self::Missing { key } => Self::Missing { key: key.clone() },
29        }
30    }
31}
32
33impl<K: PartialEq, D: Protocol> PartialEq for RegistryResult<K, D> {
34    fn eq(&self, other: &Self) -> bool {
35        match (self, other) {
36            (
37                Self::Found {
38                    key: left_key,
39                    recipient: left_recipient,
40                },
41                Self::Found {
42                    key: right_key,
43                    recipient: right_recipient,
44                },
45            ) => left_key == right_key && left_recipient == right_recipient,
46            (Self::Missing { key: left }, Self::Missing { key: right }) => left == right,
47            (Self::Found { .. }, Self::Missing { .. })
48            | (Self::Missing { .. }, Self::Found { .. }) => false,
49        }
50    }
51}
52
53impl<K: Eq, D: Protocol> Eq for RegistryResult<K, D> {}
54
55/// Commands accepted by a [`Registry`].
56///
57/// A lookup owns a typed reply recipient; no runtime registry or reply-channel
58/// discovery is performed.
59pub enum RegistryMessage<K, D: Protocol, Route> {
60    /// Establish a previously absent binding.
61    Bind { key: K, recipient: Recipient<D> },
62    /// Remove a binding only if it still names the supplied recipient.
63    Unbind { key: K, recipient: Recipient<D> },
64    /// Return the exact current result to `reply_to`.
65    Lookup { key: K, reply_to: Route },
66}
67
68/// A rejected registry mutation.
69#[derive(Error, Clone, PartialEq, Eq)]
70pub enum RegistryError<K, D: Protocol> {
71    /// A different recipient already owns the key.
72    #[error("registry key is already bound")]
73    AlreadyBound {
74        /// Rejected key.
75        key: K,
76        /// Exact recipient from the rejected bind command.
77        recipient: Recipient<D>,
78        /// Recipient that currently owns the key.
79        current: Recipient<D>,
80    },
81    /// Unbinding named an absent key.
82    #[error("registry key is not bound")]
83    NotBound { key: K, recipient: Recipient<D> },
84    /// Unbinding named a recipient other than the current binding.
85    #[error("registry unbind is stale")]
86    StaleBinding {
87        key: K,
88        recipient: Recipient<D>,
89        current: Recipient<D>,
90    },
91}
92
93impl<K: core::fmt::Debug, D: Protocol> core::fmt::Debug for RegistryError<K, D> {
94    fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
95        match self {
96            Self::AlreadyBound { key, .. } => formatter
97                .debug_struct("AlreadyBound")
98                .field("key", key)
99                .field("recipient", &"<typed recipient>")
100                .field("current", &"<typed recipient>")
101                .finish(),
102            Self::NotBound { key, .. } => formatter
103                .debug_struct("NotBound")
104                .field("key", key)
105                .field("recipient", &"<typed recipient>")
106                .finish(),
107            Self::StaleBinding { key, .. } => formatter
108                .debug_struct("StaleBinding")
109                .field("key", key)
110                .field("recipient", &"<typed recipient>")
111                .field("current", &"<typed recipient>")
112                .finish(),
113        }
114    }
115}
116
117/// Insertion-ordered typed recipient registry.
118///
119/// State is a sequence of unique key/recipient bindings. Inputs are
120/// [`RegistryMessage`] values and lookup outputs are one typed
121/// [`behavior::Delivery`]. Initialization is empty. Duplicate bind, absent
122/// unbind, and stale unbind are explicit errors and leave state unchanged.
123/// Lookups never fail: absence is a factual [`RegistryResult::Missing`]. The
124/// actor does not terminate by policy. Ordering and conflict behavior are
125/// deliberate Bombay policy; endpoint generation and delivery are interpreted
126/// by Bombay Address and Communication.
127pub struct Registry<A, K, D, Route>
128where
129    A: Address,
130    D: Protocol<Addr = A>,
131    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = RegistryResult<K, D>>>,
132{
133    bindings: Vec<(K, Recipient<D>)>,
134    address: core::marker::PhantomData<fn() -> (A, Route)>,
135}
136
137impl<A, K, D, Route> Registry<A, K, D, Route>
138where
139    A: Address,
140    D: Protocol<Addr = A>,
141    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = RegistryResult<K, D>>>,
142{
143    /// Construct an empty registry definition.
144    #[must_use]
145    pub const fn new() -> Self {
146        Self {
147            bindings: Vec::new(),
148            address: core::marker::PhantomData,
149        }
150    }
151
152    /// Current bindings in establishment order.
153    #[must_use]
154    pub fn bindings(&self) -> &[(K, Recipient<D>)] {
155        &self.bindings
156    }
157}
158
159impl<A, K, D, Route> Default for Registry<A, K, D, Route>
160where
161    A: Address,
162    D: Protocol<Addr = A>,
163    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = RegistryResult<K, D>>>,
164{
165    fn default() -> Self {
166        Self::new()
167    }
168}
169
170impl<A, K, D, Route> BehaviorBase for Registry<A, K, D, Route>
171where
172    A: Address,
173    D: Protocol<Addr = A>,
174    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = RegistryResult<K, D>>>,
175{
176    type Base = Self;
177
178    fn base(&self) -> &Self::Base {
179        self
180    }
181}
182
183impl<A, K, D, Route> behavior::Protocol for Registry<A, K, D, Route>
184where
185    A: Address,
186    D: Protocol<Addr = A>,
187    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = RegistryResult<K, D>>>,
188{
189    type Addr = A;
190    type Msg = RegistryMessage<K, D, Route>;
191}
192
193impl<A, K, D, Route> Behavior for Registry<A, K, D, Route>
194where
195    A: Address,
196    K: Clone + Eq,
197    D: Protocol<Addr = A>,
198    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = RegistryResult<K, D>>>,
199    Route::Sends: behavior::SendsFor<User<A, RegistryMessage<K, D, Route>>>,
200{
201    type Protocol = Self;
202    type Event = User<A, behavior::BehaviorMessage<Self>>;
203    type Sends = Route::Sends;
204    type Ph = Never;
205    type Error = RegistryError<K, D>;
206    type Birth = NoBirths;
207
208    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
209        match event.message {
210            RegistryMessage::Bind { key, recipient } => {
211                if let Some((_, current)) = self.bindings.iter().find(|(bound, _)| *bound == key) {
212                    return Err(RegistryError::AlreadyBound {
213                        key,
214                        recipient,
215                        current: *current,
216                    });
217                }
218                self.bindings.push((key, recipient));
219                Ok(Actions::cont())
220            }
221            RegistryMessage::Unbind { key, recipient } => {
222                let Some(index) = self.bindings.iter().position(|(bound, _)| *bound == key) else {
223                    return Err(RegistryError::NotBound { key, recipient });
224                };
225                if self.bindings[index].1 != recipient {
226                    return Err(RegistryError::StaleBinding {
227                        key,
228                        recipient,
229                        current: self.bindings[index].1,
230                    });
231                }
232                self.bindings.remove(index);
233                Ok(Actions::cont())
234            }
235            RegistryMessage::Lookup { key, reply_to } => {
236                let result = self
237                    .bindings
238                    .iter()
239                    .find(|(bound, _)| *bound == key)
240                    .map_or_else(
241                        || RegistryResult::Missing { key: key.clone() },
242                        |(_, recipient)| RegistryResult::Found {
243                            key: key.clone(),
244                            recipient: *recipient,
245                        },
246                    );
247                Ok(Actions::send(reply_to.deliver(result)))
248            }
249        }
250    }
251}
252
253#[cfg(test)]
254mod tests {
255    use super::*;
256    use crate::Activate as _;
257    use behavior::{MailAddr, Step};
258
259    #[derive(PartialEq, Eq)]
260    struct Destination;
261    struct Reply;
262
263    impl behavior::Protocol for Destination {
264        type Addr = MailAddr;
265        type Msg = u8;
266    }
267
268    impl Behavior for Destination {
269        type Protocol = Self;
270        type Event = User<MailAddr, u8>;
271        type Sends = Vec<Never>;
272        type Ph = Never;
273        type Error = Never;
274        type Birth = NoBirths;
275
276        fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
277            Ok(Actions::cont())
278        }
279    }
280
281    impl behavior::Protocol for Reply {
282        type Addr = MailAddr;
283        type Msg = RegistryResult<u8, Destination>;
284    }
285
286    impl Behavior for Reply {
287        type Protocol = Self;
288        type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
289        type Sends = Vec<Never>;
290        type Ph = Never;
291        type Error = Never;
292        type Birth = NoBirths;
293
294        fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
295            Ok(Actions::cont())
296        }
297    }
298
299    type TestRegistry = Registry<MailAddr, u8, Destination, Recipient<Reply>>;
300
301    #[test]
302    fn mutations_are_atomic_and_stale_unbind_is_typed() {
303        let one = Recipient::<Destination>::global(MailAddr(1));
304        let two = Recipient::<Destination>::global(MailAddr(2));
305        let mut registry = (TestRegistry::new()).initialize().unwrap().behavior;
306
307        let bound = registry
308            .receive(
309                MailAddr(9),
310                RegistryMessage::Bind {
311                    key: 4,
312                    recipient: one,
313                },
314            )
315            .unwrap();
316        assert!(bound.sends.is_empty());
317        assert!(bound.creates.is_empty());
318        assert!(matches!(bound.become_, Step::Continue));
319
320        let rejection = registry.receive(
321            MailAddr(9),
322            RegistryMessage::Bind {
323                key: 4,
324                recipient: two,
325            },
326        );
327        assert!(matches!(
328            rejection,
329            Err(RegistryError::AlreadyBound { key: 4, recipient, current })
330                if recipient == two && current == one
331        ));
332        let rejection = registry.receive(
333            MailAddr(9),
334            RegistryMessage::Unbind {
335                key: 4,
336                recipient: two,
337            },
338        );
339        assert!(matches!(
340            rejection,
341            Err(RegistryError::StaleBinding {
342                key: 4,
343                recipient,
344                current,
345            }) if recipient == two && current == one
346        ));
347        assert!(registry.bindings() == [(4, one)]);
348
349        let unbound = registry
350            .receive(
351                MailAddr(9),
352                RegistryMessage::Unbind {
353                    key: 4,
354                    recipient: one,
355                },
356            )
357            .unwrap();
358        assert!(unbound.sends.is_empty());
359        assert!(unbound.creates.is_empty());
360        assert_eq!(unbound.become_, behavior::Step::Continue);
361        assert!(registry.bindings().is_empty());
362    }
363
364    #[test]
365    fn lookups_reply_with_found_and_missing_results() {
366        let destination = Recipient::<Destination>::global(MailAddr(1));
367        let reply = Recipient::<Reply>::global(MailAddr(8));
368        let mut registry = (TestRegistry::new()).initialize().unwrap().behavior;
369        let bound = registry
370            .receive(
371                MailAddr(9),
372                RegistryMessage::Bind {
373                    key: 4,
374                    recipient: destination,
375                },
376            )
377            .unwrap();
378        assert!(bound.sends.is_empty());
379        assert!(bound.creates.is_empty());
380        assert_eq!(bound.become_, behavior::Step::Continue);
381
382        let found = registry
383            .receive(
384                MailAddr(9),
385                RegistryMessage::Lookup {
386                    key: 4,
387                    reply_to: reply,
388                },
389            )
390            .unwrap();
391        assert!(
392            found.sends
393                == vec![Delivery::new(
394                    reply,
395                    RegistryResult::Found {
396                        key: 4,
397                        recipient: destination,
398                    },
399                )]
400        );
401        assert!(found.creates.is_empty());
402        assert!(matches!(found.become_, Step::Continue));
403
404        let missing = registry
405            .receive(
406                MailAddr(9),
407                RegistryMessage::Lookup {
408                    key: 5,
409                    reply_to: reply,
410                },
411            )
412            .unwrap();
413        assert!(missing.sends == vec![Delivery::new(reply, RegistryResult::Missing { key: 5 })]);
414    }
415}