1#[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
13pub enum RegistryResult<K, D: Protocol> {
15 Found { key: K, recipient: Recipient<D> },
17 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
55pub enum RegistryMessage<K, D: Protocol, Route> {
60 Bind { key: K, recipient: Recipient<D> },
62 Unbind { key: K, recipient: Recipient<D> },
64 Lookup { key: K, reply_to: Route },
66}
67
68#[derive(Error, Clone, PartialEq, Eq)]
70pub enum RegistryError<K, D: Protocol> {
71 #[error("registry key is already bound")]
73 AlreadyBound {
74 key: K,
76 recipient: Recipient<D>,
78 current: Recipient<D>,
80 },
81 #[error("registry key is not bound")]
83 NotBound { key: K, recipient: Recipient<D> },
84 #[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
117pub 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 #[must_use]
145 pub const fn new() -> Self {
146 Self {
147 bindings: Vec::new(),
148 address: core::marker::PhantomData,
149 }
150 }
151
152 #[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}