1use core::future::Future;
4use core::hash::Hash;
5
6use behavior::{
7 Address, Behavior, Crash, DeadlineSends, Delivery, Exit, ObserveChild, ObserveCreation,
8 ObservePeer, ProxyCommand, ProxySends, ReceiveTimeoutSends, ReportWorkerCreationResolved,
9 ReportWorkerStopped, Route, SendProduct, ServiceSends, SupervisorSends, WatchSends,
10};
11pub use bombay_address::AddressInUse;
12use bombay_address::{AddressSpace, Lease};
13use observe::{Observation, ObservationSpace};
14
15pub(crate) type PeerOutcome<A> = Result<Exit<A>, Crash>;
17
18#[doc(hidden)]
20pub struct IncarnationEndpoint<A: Address, D> {
21 delivery: D,
22 completion: ObservationSpace<(), PeerOutcome<A>>,
23}
24
25impl<A: Address, D> IncarnationEndpoint<A, D> {
26 pub(crate) const fn new(delivery: D, completion: ObservationSpace<(), PeerOutcome<A>>) -> Self {
27 Self {
28 delivery,
29 completion,
30 }
31 }
32
33 fn observation(&self) -> Observation<PeerOutcome<A>> {
34 self.completion
35 .observe(&())
36 .expect("a registered incarnation must retain its completion subject")
37 }
38}
39
40impl<A: Address, D: Clone> Clone for IncarnationEndpoint<A, D> {
41 fn clone(&self) -> Self {
42 Self::new(self.delivery.clone(), self.completion.clone())
43 }
44}
45
46impl<A, M, D> DeliveryEndpoint<A, M> for IncarnationEndpoint<A, D>
47where
48 A: Address + Send + Sync,
49 M: Send,
50 D: DeliveryEndpoint<A, M> + Sync,
51{
52 type Error = D::Error;
53
54 async fn deliver(&self, from: A, message: M) -> Result<(), RejectedDelivery<M, Self::Error>> {
55 self.delivery.deliver(from, message).await
56 }
57}
58
59#[doc(hidden)]
61#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
62pub enum PeerObservationError<A> {
63 #[error("no live incarnation is registered at address {0:?}")]
65 Unknown(A),
66}
67
68#[doc(hidden)]
70pub trait PeerObserver<A: Address> {
71 fn observe_peer(&self, peer: A)
72 -> Result<Observation<PeerOutcome<A>>, PeerObservationError<A>>;
73}
74
75#[doc(hidden)]
80pub trait DeliveryEndpoint<A, M> {
81 type Error;
83
84 fn deliver(
86 &self,
87 from: A,
88 message: M,
89 ) -> impl Future<Output = Result<(), RejectedDelivery<M, Self::Error>>> + Send;
90}
91
92#[derive(Debug, Clone, PartialEq, Eq)]
94pub struct RejectedDelivery<M, E> {
95 pub message: M,
97 pub error: E,
99}
100
101impl<M, E> RejectedDelivery<M, E> {
102 #[must_use]
104 pub const fn new(message: M, error: E) -> Self {
105 Self { message, error }
106 }
107}
108
109pub trait EndpointRegistry<A, M, D> {
111 type Error;
113 type Registration;
115
116 fn register(&self, address: A, endpoint: D) -> Result<Self::Registration, Self::Error>;
122}
123
124#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
126pub enum RoutingError<A: Address, M, E> {
127 #[error("no endpoint exists for the resolved address {address:?}")]
129 UnknownAddress {
130 address: A,
132 message: M,
134 },
135 #[error("the endpoint at {address:?} rejected delivery: {rejected:?}")]
137 Endpoint {
138 address: A,
140 rejected: RejectedDelivery<M, E>,
142 },
143}
144
145pub struct AddressRouter<A, D> {
147 entries: AddressSpace<A, D>,
148}
149
150impl<A, D> Clone for AddressRouter<A, D> {
151 fn clone(&self) -> Self {
152 Self {
153 entries: self.entries.clone(),
154 }
155 }
156}
157
158impl<A, D> Default for AddressRouter<A, D> {
159 fn default() -> Self {
160 Self {
161 entries: AddressSpace::new(),
162 }
163 }
164}
165
166impl<A: Address + Hash, M, D> EndpointRegistry<A, M, D> for AddressRouter<A, D> {
167 type Error = AddressInUse<A>;
168 type Registration = Lease<A, D>;
169
170 fn register(&self, address: A, endpoint: D) -> Result<Self::Registration, Self::Error> {
171 self.entries.claim(address, endpoint)
172 }
173}
174
175impl<A, M, D> DeliveryRouter<A, M> for AddressRouter<A, D>
176where
177 A: Address + Hash + Send + Sync,
178 A::Nonce: Send,
179 M: Send,
180 D: DeliveryEndpoint<A, M> + Clone + Send + Sync,
181 D::Error: Send,
182{
183 type Error = RoutingError<A, M, D::Error>;
184
185 async fn deliver(&self, from: A, delivery: Delivery<A, M>) -> Result<(), Self::Error> {
186 let address = resolve_route(from, delivery.to.route());
187 let Some(endpoint) = self.entries.resolve(&address) else {
188 return Err(RoutingError::UnknownAddress {
189 address,
190 message: delivery.message,
191 });
192 };
193 endpoint
194 .deliver(from, delivery.message)
195 .await
196 .map_err(|rejected| RoutingError::Endpoint { address, rejected })
197 }
198}
199
200impl<A, D> PeerObserver<A> for AddressRouter<A, IncarnationEndpoint<A, D>>
201where
202 A: Address + Hash,
203 D: Clone,
204{
205 fn observe_peer(
206 &self,
207 peer: A,
208 ) -> Result<Observation<PeerOutcome<A>>, PeerObservationError<A>> {
209 self.entries
210 .resolve(&peer)
211 .map(|endpoint| endpoint.observation())
212 .ok_or(PeerObservationError::Unknown(peer))
213 }
214}
215
216fn resolve_route<A: Address>(from: A, route: Route<A>) -> A {
217 match route {
218 Route::Global(address) => address,
219 Route::Child(nonce) => from.birth(nonce),
220 }
221}
222
223pub trait DeliveryRouter<A: Address, M> {
228 type Error;
230
231 fn deliver(
233 &self,
234 from: A,
235 delivery: Delivery<A, M>,
236 ) -> impl Future<Output = Result<(), Self::Error>> + Send;
237}
238
239pub trait ObservesCreations<N> {
247 fn observes_creation(&self, nonce: N) -> bool;
249}
250
251impl<A: Address, M> ObservesCreations<A::Nonce> for Vec<Delivery<A, M>> {
252 fn observes_creation(&self, _nonce: A::Nonce) -> bool {
253 false
254 }
255}
256
257impl<N> ObservesCreations<N> for ServiceSends<ObserveCreation<N>>
258where
259 N: Copy + Eq,
260{
261 fn observes_creation(&self, nonce: N) -> bool {
262 self.iter().any(|request| request.nonce == nonce)
263 }
264}
265
266macro_rules! observes_no_creations {
267 ($($lane:ty),+ $(,)?) => {
268 $(
269 impl<N> ObservesCreations<N> for $lane {
270 fn observes_creation(&self, _nonce: N) -> bool {
271 false
272 }
273 }
274 )+
275 };
276}
277
278observes_no_creations!(
279 ServiceSends<behavior::ScheduleAt>,
280 ServiceSends<behavior::ScheduleAfter>,
281);
282
283impl<M, N> ObservesCreations<N> for ServiceSends<ObserveChild<M>> {
284 fn observes_creation(&self, _nonce: N) -> bool {
285 false
286 }
287}
288
289impl<A: Address, N> ObservesCreations<N> for ServiceSends<ObservePeer<A>> {
290 fn observes_creation(&self, _nonce: N) -> bool {
291 false
292 }
293}
294
295impl<A, N> ObservesCreations<N> for ServiceSends<behavior::UnwatchPeer<A>> {
296 fn observes_creation(&self, _nonce: N) -> bool {
297 false
298 }
299}
300
301impl<A: Address, N> ObservesCreations<N> for ServiceSends<ReportWorkerStopped<A>> {
302 fn observes_creation(&self, _nonce: N) -> bool {
303 false
304 }
305}
306
307impl<M, N> ObservesCreations<N> for ServiceSends<ReportWorkerCreationResolved<M>> {
308 fn observes_creation(&self, _nonce: N) -> bool {
309 false
310 }
311}
312
313impl<L, R, N> ObservesCreations<N> for SendProduct<L, R>
314where
315 L: ObservesCreations<N>,
316 R: ObservesCreations<N>,
317 N: Copy,
318{
319 fn observes_creation(&self, nonce: N) -> bool {
320 self.inner.observes_creation(nonce) || self.own.observes_creation(nonce)
321 }
322}
323
324impl<A: Address, Sends> ObservesCreations<A::Nonce> for WatchSends<A, Sends>
325where
326 Sends: ObservesCreations<A::Nonce>,
327 A::Nonce: Copy,
328{
329 fn observes_creation(&self, nonce: A::Nonce) -> bool {
330 self.behavior.observes_creation(nonce)
331 }
332}
333
334impl<Sends, N> ObservesCreations<N> for DeadlineSends<Sends>
335where
336 Sends: ObservesCreations<N>,
337 N: Copy,
338{
339 fn observes_creation(&self, nonce: N) -> bool {
340 self.behavior.observes_creation(nonce)
341 }
342}
343
344impl<Sends, N> ObservesCreations<N> for ReceiveTimeoutSends<Sends>
345where
346 Sends: ObservesCreations<N>,
347 N: Copy,
348{
349 fn observes_creation(&self, nonce: N) -> bool {
350 self.behavior.observes_creation(nonce)
351 }
352}
353
354impl<A: Address, Sends, C> ObservesCreations<A::Nonce> for SupervisorSends<A, Sends, C>
355where
356 Sends: ObservesCreations<A::Nonce>,
357 C: Behavior<Addr = A>,
358 A::Nonce: Copy,
359{
360 fn observes_creation(&self, nonce: A::Nonce) -> bool {
361 self.behavior.observes_creation(nonce)
362 }
363}
364
365impl<A: Address, M> ObservesCreations<A::Nonce> for ProxySends<A, M>
366where
367 A::Nonce: Copy + Eq,
368{
369 fn observes_creation(&self, nonce: A::Nonce) -> bool {
370 self.creation_observations.observes_creation(nonce)
371 }
372}
373
374pub trait RouteSends<A: Address, R>: Sized {
380 type Error;
382
383 fn route(self, from: A, router: &mut R)
385 -> impl Future<Output = Result<(), Self::Error>> + Send;
386}
387
388impl<A, M, R> RouteSends<A, R> for Vec<Delivery<A, M>>
389where
390 A: Address + Send,
391 M: Send,
392 R: DeliveryRouter<A, M> + Send + Sync,
393 A::Nonce: Send,
394{
395 type Error = R::Error;
396
397 async fn route(self, from: A, router: &mut R) -> Result<(), Self::Error> {
398 for delivery in self {
399 router.deliver(from, delivery).await?;
400 }
401 Ok(())
402 }
403}
404
405#[doc(hidden)]
407#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
408pub enum SendProductError<L, R> {
409 #[error("inner protocol leg failed: {0:?}")]
411 Inner(L),
412 #[error("wrapper-owned protocol leg failed: {0:?}")]
414 Own(R),
415}
416
417impl<A, R, L, Own> RouteSends<A, R> for SendProduct<L, Own>
418where
419 A: Address + Send,
420 L: RouteSends<A, R> + Send,
421 Own: RouteSends<A, R> + Send,
422 R: Send,
423{
424 type Error = SendProductError<L::Error, Own::Error>;
425
426 async fn route(self, from: A, router: &mut R) -> Result<(), Self::Error> {
427 self.inner
428 .route(from, router)
429 .await
430 .map_err(SendProductError::Inner)?;
431 self.own
432 .route(from, router)
433 .await
434 .map_err(SendProductError::Own)
435 }
436}
437
438#[doc(hidden)]
440#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
441pub enum WatchSendsError<B, O> {
442 #[error("watch behavior sends failed: {0:?}")]
444 Behavior(B),
445 #[error("watch observation lane failed: {0:?}")]
447 Observations(O),
448}
449
450impl<A, R, Sends> RouteSends<A, R> for WatchSends<A, Sends>
451where
452 A: Address + Send,
453 Sends: RouteSends<A, R> + Send,
454 ServiceSends<ObservePeer<A>>: RouteSends<A, R> + Send,
455 R: Send,
456{
457 type Error = WatchSendsError<
458 <Sends as RouteSends<A, R>>::Error,
459 <ServiceSends<ObservePeer<A>> as RouteSends<A, R>>::Error,
460 >;
461
462 async fn route(self, from: A, router: &mut R) -> Result<(), Self::Error> {
463 self.behavior
464 .route(from, router)
465 .await
466 .map_err(WatchSendsError::Behavior)?;
467 self.observations
468 .route(from, router)
469 .await
470 .map_err(WatchSendsError::Observations)
471 }
472}
473
474#[doc(hidden)]
476#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
477pub enum DeadlineSendsError<B, S> {
478 #[error("deadline behavior sends failed: {0:?}")]
480 Behavior(B),
481 #[error("deadline schedule lane failed: {0:?}")]
483 Schedules(S),
484}
485
486impl<A, R, Sends> RouteSends<A, R> for DeadlineSends<Sends>
487where
488 A: Address + Send,
489 Sends: RouteSends<A, R> + Send,
490 ServiceSends<behavior::ScheduleAt>: RouteSends<A, R> + Send,
491 R: Send,
492{
493 type Error = DeadlineSendsError<
494 <Sends as RouteSends<A, R>>::Error,
495 <ServiceSends<behavior::ScheduleAt> as RouteSends<A, R>>::Error,
496 >;
497
498 async fn route(self, from: A, router: &mut R) -> Result<(), Self::Error> {
499 self.behavior
500 .route(from, router)
501 .await
502 .map_err(DeadlineSendsError::Behavior)?;
503 self.schedules
504 .route(from, router)
505 .await
506 .map_err(DeadlineSendsError::Schedules)
507 }
508}
509
510#[doc(hidden)]
512#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
513pub enum ReceiveTimeoutSendsError<B, S> {
514 #[error("receive-timeout behavior sends failed: {0:?}")]
516 Behavior(B),
517 #[error("receive-timeout schedule lane failed: {0:?}")]
519 Schedules(S),
520}
521
522impl<A, R, Sends> RouteSends<A, R> for ReceiveTimeoutSends<Sends>
523where
524 A: Address + Send,
525 Sends: RouteSends<A, R> + Send,
526 ServiceSends<behavior::ScheduleAfter>: RouteSends<A, R> + Send,
527 R: Send,
528{
529 type Error = ReceiveTimeoutSendsError<
530 <Sends as RouteSends<A, R>>::Error,
531 <ServiceSends<behavior::ScheduleAfter> as RouteSends<A, R>>::Error,
532 >;
533
534 async fn route(self, from: A, router: &mut R) -> Result<(), Self::Error> {
535 self.behavior
536 .route(from, router)
537 .await
538 .map_err(ReceiveTimeoutSendsError::Behavior)?;
539 self.schedules
540 .route(from, router)
541 .await
542 .map_err(ReceiveTimeoutSendsError::Schedules)
543 }
544}
545
546#[doc(hidden)]
548#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
549pub enum SupervisorSendsError<B, O, C> {
550 #[error("supervised behavior sends failed: {0:?}")]
552 Behavior(B),
553 #[error("supervisor child-observation lane failed: {0:?}")]
555 ChildObservations(O),
556 #[error("supervisor replacement-command lane failed: {0:?}")]
558 ReplacementCommands(C),
559}
560
561impl<A, R, Sends, C> RouteSends<A, R> for SupervisorSends<A, Sends, C>
562where
563 A: Address + Send,
564 A::Nonce: Send,
565 Sends: RouteSends<A, R> + Send,
566 ServiceSends<ObserveChild<A::Nonce>>: RouteSends<A, R> + Send,
567 Vec<Delivery<A, ProxyCommand<C>>>: RouteSends<A, R> + Send,
568 C: Behavior<Addr = A> + Send,
569 R: Send,
570{
571 type Error = SupervisorSendsError<
572 <Sends as RouteSends<A, R>>::Error,
573 <ServiceSends<ObserveChild<A::Nonce>> as RouteSends<A, R>>::Error,
574 <Vec<Delivery<A, ProxyCommand<C>>> as RouteSends<A, R>>::Error,
575 >;
576
577 async fn route(self, from: A, router: &mut R) -> Result<(), Self::Error> {
578 self.behavior
579 .route(from, router)
580 .await
581 .map_err(SupervisorSendsError::Behavior)?;
582 self.child_observations
583 .route(from, router)
584 .await
585 .map_err(SupervisorSendsError::ChildObservations)?;
586 self.replacement_commands
587 .route(from, router)
588 .await
589 .map_err(SupervisorSendsError::ReplacementCommands)
590 }
591}
592
593#[doc(hidden)]
595#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
596pub enum ProxySendsError<D, C, CR, S, RP> {
597 #[error("proxy delivery lane failed: {0:?}")]
599 Deliveries(D),
600 #[error("proxy child-observation lane failed: {0:?}")]
602 ChildObservations(C),
603 #[error("proxy creation-observation lane failed: {0:?}")]
605 CreationObservations(CR),
606 #[error("proxy worker-stopped report lane failed: {0:?}")]
608 StoppedReports(S),
609 #[error("proxy worker-creation report lane failed: {0:?}")]
611 CreationReports(RP),
612}
613
614impl<A, M, R> RouteSends<A, R> for ProxySends<A, M>
615where
616 A: Address + Send,
617 A::Nonce: Send,
618 Vec<Delivery<A, M>>: RouteSends<A, R> + Send,
619 ServiceSends<ObserveChild<A::Nonce>>: RouteSends<A, R> + Send,
620 ServiceSends<ObserveCreation<A::Nonce>>: RouteSends<A, R> + Send,
621 ServiceSends<ReportWorkerStopped<A>>: RouteSends<A, R> + Send,
622 ServiceSends<ReportWorkerCreationResolved<A::Nonce>>: RouteSends<A, R> + Send,
623 R: Send,
624{
625 type Error = ProxySendsError<
626 <Vec<Delivery<A, M>> as RouteSends<A, R>>::Error,
627 <ServiceSends<ObserveChild<A::Nonce>> as RouteSends<A, R>>::Error,
628 <ServiceSends<ObserveCreation<A::Nonce>> as RouteSends<A, R>>::Error,
629 <ServiceSends<ReportWorkerStopped<A>> as RouteSends<A, R>>::Error,
630 <ServiceSends<ReportWorkerCreationResolved<A::Nonce>> as RouteSends<A, R>>::Error,
631 >;
632
633 async fn route(self, from: A, router: &mut R) -> Result<(), Self::Error> {
634 self.deliveries
635 .route(from, router)
636 .await
637 .map_err(ProxySendsError::Deliveries)?;
638 self.child_observations
639 .route(from, router)
640 .await
641 .map_err(ProxySendsError::ChildObservations)?;
642 self.creation_observations
643 .route(from, router)
644 .await
645 .map_err(ProxySendsError::CreationObservations)?;
646 self.stopped_reports
647 .route(from, router)
648 .await
649 .map_err(ProxySendsError::StoppedReports)?;
650 self.creation_reports
651 .route(from, router)
652 .await
653 .map_err(ProxySendsError::CreationReports)
654 }
655}
656
657#[cfg(test)]
658mod tests {
659 use core::convert::Infallible;
660 use std::sync::{Arc, Mutex};
661
662 use behavior::{
663 Address, Delivery, Exit, MailAddr, ObserveChild, ObserveCreation, ObservePeer, Recipient,
664 ReportWorkerCreationResolved, ReportWorkerStopped, ScheduleAfter, ScheduleAt, SendProduct,
665 ServiceSends, UnwatchPeer, WatchSends,
666 };
667 use observe::ObservationSpace;
668
669 use super::{
670 AddressRouter, DeliveryEndpoint, DeliveryRouter, EndpointRegistry, IncarnationEndpoint,
671 ObservesCreations, PeerObserver, RejectedDelivery, RouteSends, RoutingError,
672 SendProductError,
673 };
674
675 #[test]
676 fn probe_truth_table_distinguishes_creation_lanes_and_nonces() {
677 let creations = ServiceSends::one(ObserveCreation::new(7u64));
678 assert!(creations.observes_creation(7));
679 assert!(
680 !creations.observes_creation(9),
681 "the creation lane matches only its staged nonce"
682 );
683
684 assert!(
685 !Vec::<Delivery<MailAddr, u8>>::new().observes_creation(7u64),
686 "plain deliveries never observe creations"
687 );
688 assert!(!ServiceSends::<ObserveChild<u64>>::new(Vec::new()).observes_creation(7u64));
689 assert!(!ServiceSends::<ObservePeer<MailAddr>>::new(Vec::new()).observes_creation(7u64));
690 assert!(
691 !ServiceSends::<ReportWorkerStopped<MailAddr>>::new(Vec::new()).observes_creation(7u64)
692 );
693 assert!(
694 !ServiceSends::<ReportWorkerCreationResolved<u64>>::new(Vec::new())
695 .observes_creation(7u64)
696 );
697 assert!(!ServiceSends::<ScheduleAt>::new(Vec::new()).observes_creation(7u64));
698 assert!(!ServiceSends::<ScheduleAfter>::new(Vec::new()).observes_creation(7u64));
699
700 let product = SendProduct {
701 inner: Vec::<Delivery<MailAddr, u8>>::new(),
702 own: creations.clone(),
703 };
704 assert!(product.observes_creation(7));
705 assert!(
706 !SendProduct {
707 inner: Vec::<Delivery<MailAddr, u8>>::new(),
708 own: ServiceSends::<ObserveChild<u64>>::new(Vec::new()),
709 }
710 .observes_creation(7u64),
711 "a product observes only through a real creation lane"
712 );
713
714 let wrapped = WatchSends {
715 behavior: creations.clone(),
716 observations: ServiceSends::one(ObservePeer::new(MailAddr(3))),
717 };
718 assert!(wrapped.observes_creation(7));
719 assert!(!wrapped.observes_creation(9));
720 let unwrapped = WatchSends {
721 behavior: Vec::<Delivery<MailAddr, u8>>::new(),
722 observations: ServiceSends::one(ObservePeer::new(MailAddr(3))),
723 };
724 assert!(
725 !unwrapped.observes_creation(7u64),
726 "wrappers delegate to their behavior lane only"
727 );
728
729 let deadline = behavior::DeadlineSends {
730 behavior: creations.clone(),
731 schedules: ServiceSends::<ScheduleAt>::new(Vec::new()),
732 };
733 assert!(deadline.observes_creation(7));
734 let deadline_empty = behavior::DeadlineSends {
735 behavior: Vec::<Delivery<MailAddr, u8>>::new(),
736 schedules: ServiceSends::<ScheduleAt>::new(Vec::new()),
737 };
738 assert!(!deadline_empty.observes_creation(7u64));
739
740 let timeout = behavior::ReceiveTimeoutSends {
741 behavior: creations.clone(),
742 schedules: ServiceSends::<ScheduleAfter>::new(Vec::new()),
743 };
744 assert!(timeout.observes_creation(7));
745 let timeout_empty = behavior::ReceiveTimeoutSends {
746 behavior: Vec::<Delivery<MailAddr, u8>>::new(),
747 schedules: ServiceSends::<ScheduleAfter>::new(Vec::new()),
748 };
749 assert!(!timeout_empty.observes_creation(7u64));
750
751 let supervised = behavior::SupervisorSends::<MailAddr, _, behavior::Pure<AssertChild, u8>> {
752 behavior: creations.clone(),
753 child_observations: ServiceSends::<ObserveChild<u64>>::new(Vec::new()),
754 replacement_commands: Vec::new(),
755 };
756 assert!(supervised.observes_creation(7));
757 let supervised_empty =
758 behavior::SupervisorSends::<MailAddr, _, behavior::Pure<AssertChild, u8>> {
759 behavior: Vec::<Delivery<MailAddr, u8>>::new(),
760 child_observations: ServiceSends::<ObserveChild<u64>>::new(Vec::new()),
761 replacement_commands: Vec::new(),
762 };
763 assert!(!supervised_empty.observes_creation(7u64));
764
765 let proxy = behavior::ProxySends::<MailAddr, u8> {
766 deliveries: Vec::new(),
767 child_observations: ServiceSends::<ObserveChild<u64>>::new(Vec::new()),
768 creation_observations: creations,
769 stopped_reports: ServiceSends::<ReportWorkerStopped<MailAddr>>::new(Vec::new()),
770 creation_reports: ServiceSends::<ReportWorkerCreationResolved<u64>>::new(Vec::new()),
771 };
772 assert!(proxy.observes_creation(7));
773 assert!(!proxy.observes_creation(9));
774 let proxy_empty = behavior::ProxySends::<MailAddr, u8> {
775 deliveries: Vec::new(),
776 child_observations: ServiceSends::<ObserveChild<u64>>::new(Vec::new()),
777 creation_observations: ServiceSends::<ObserveCreation<u64>>::new(Vec::new()),
778 stopped_reports: ServiceSends::<ReportWorkerStopped<MailAddr>>::new(Vec::new()),
779 creation_reports: ServiceSends::<ReportWorkerCreationResolved<u64>>::new(Vec::new()),
780 };
781 assert!(!proxy_empty.observes_creation(7u64));
782 }
783
784 #[test]
785 fn peer_observation_cancellation_never_claims_creation_observation() {
786 assert!(!ServiceSends::<UnwatchPeer<MailAddr>>::new(Vec::new()).observes_creation(7u64));
787 }
788
789 struct AssertChild;
790
791 impl behavior::Handler<u8> for AssertChild {
792 type Addr = MailAddr;
793 type Msg = u8;
794
795 fn receive(
796 &mut self,
797 _from: MailAddr,
798 _message: u8,
799 ) -> behavior::Acted<
800 MailAddr,
801 behavior::Never,
802 Vec<Delivery<MailAddr, u8>>,
803 behavior::NoBirths,
804 behavior::Never,
805 > {
806 Ok(behavior::Actions::cont())
807 }
808 }
809
810 #[derive(Clone, Default)]
811 struct RecordingEndpoint(Arc<Mutex<Vec<(MailAddr, u8)>>>);
812
813 impl DeliveryEndpoint<MailAddr, u8> for RecordingEndpoint {
814 type Error = Infallible;
815
816 async fn deliver(
817 &self,
818 from: MailAddr,
819 message: u8,
820 ) -> Result<(), RejectedDelivery<u8, Self::Error>> {
821 self.0.lock().expect("endpoint lock").push((from, message));
822 Ok(())
823 }
824 }
825
826 #[derive(Clone)]
827 struct FailingEndpoint;
828
829 impl DeliveryEndpoint<MailAddr, u8> for FailingEndpoint {
830 type Error = &'static str;
831
832 async fn deliver(
833 &self,
834 _from: MailAddr,
835 message: u8,
836 ) -> Result<(), RejectedDelivery<u8, Self::Error>> {
837 Err(RejectedDelivery::new(message, "endpoint rejected delivery"))
838 }
839 }
840
841 #[derive(Default)]
842 struct RecordingRouter(Mutex<Vec<(MailAddr, Delivery<MailAddr, u8>)>>);
843
844 impl DeliveryRouter<MailAddr, u8> for RecordingRouter {
845 type Error = &'static str;
846
847 async fn deliver(
848 &self,
849 from: MailAddr,
850 delivery: Delivery<MailAddr, u8>,
851 ) -> Result<(), Self::Error> {
852 if delivery.message == 2 {
853 return Err("delivery failed");
854 }
855 self.0.lock().expect("delivery lock").push((from, delivery));
856 Ok(())
857 }
858 }
859
860 #[tokio::test]
861 async fn router_resolves_global_and_child_routes_before_endpoint_delivery() {
862 let router = AddressRouter::default();
863 let from = MailAddr(7);
864 let global = RecordingEndpoint::default();
865 let child = RecordingEndpoint::default();
866 let _global_lease =
867 EndpointRegistry::<MailAddr, u8, _>::register(&router, MailAddr(9), global.clone())
868 .unwrap();
869 let _child_lease =
870 EndpointRegistry::<MailAddr, u8, _>::register(&router, from.birth(3), child.clone())
871 .unwrap();
872
873 router
874 .deliver(from, Delivery::new(Recipient::global(MailAddr(9)), 11))
875 .await
876 .unwrap();
877 router
878 .deliver(from, Delivery::new(Recipient::child(3), 12))
879 .await
880 .unwrap();
881
882 assert_eq!(
883 *global.0.lock().expect("global endpoint lock"),
884 [(from, 11)]
885 );
886 assert_eq!(*child.0.lock().expect("child endpoint lock"), [(from, 12)]);
887 }
888
889 #[tokio::test]
890 async fn router_distinguishes_lookup_from_endpoint_failure() {
891 let router = AddressRouter::<MailAddr, FailingEndpoint>::default();
892 let from = MailAddr(7);
893
894 assert_eq!(
895 router
896 .deliver(from, Delivery::new(Recipient::global(MailAddr(8)), 1),)
897 .await,
898 Err(RoutingError::UnknownAddress {
899 address: MailAddr(8),
900 message: 1,
901 })
902 );
903
904 let _lease =
905 EndpointRegistry::<MailAddr, u8, _>::register(&router, MailAddr(9), FailingEndpoint)
906 .unwrap();
907 assert_eq!(
908 router
909 .deliver(from, Delivery::new(Recipient::global(MailAddr(9)), 2),)
910 .await,
911 Err(RoutingError::Endpoint {
912 address: MailAddr(9),
913 rejected: RejectedDelivery::new(2, "endpoint rejected delivery"),
914 })
915 );
916 }
917
918 #[tokio::test]
919 async fn incarnation_endpoint_preserves_the_inner_rejection() {
920 let completion = ObservationSpace::new();
921 let _subject = completion.subject(()).unwrap();
922 let endpoint = IncarnationEndpoint::new(FailingEndpoint, completion);
923
924 assert_eq!(
925 endpoint.deliver(MailAddr(7), 23).await,
926 Err(RejectedDelivery::new(23, "endpoint rejected delivery"))
927 );
928 }
929
930 #[test]
931 fn peer_observation_captures_the_resolved_address_generation() {
932 let router = AddressRouter::default();
933 let old_space = ObservationSpace::new();
934 let mut old_subject = old_space.subject(()).unwrap();
935 let old_endpoint = IncarnationEndpoint::new(RecordingEndpoint::default(), old_space);
936 let old_lease =
937 EndpointRegistry::<MailAddr, u8, _>::register(&router, MailAddr(9), old_endpoint)
938 .unwrap();
939 let old_observation = router.observe_peer(MailAddr(9)).unwrap();
940
941 drop(old_lease);
942 let new_space = ObservationSpace::new();
943 let mut new_subject = new_space.subject(()).unwrap();
944 let new_endpoint = IncarnationEndpoint::new(RecordingEndpoint::default(), new_space);
945 let _new_lease =
946 EndpointRegistry::<MailAddr, u8, _>::register(&router, MailAddr(9), new_endpoint)
947 .unwrap();
948 let new_observation = router.observe_peer(MailAddr(9)).unwrap();
949
950 old_subject.complete(Ok(Exit::Normal));
951 new_subject.complete(Ok(Exit::Collected));
952 assert_eq!(old_observation.try_get(), Some(Ok(Exit::Normal)));
953 assert_eq!(new_observation.try_get(), Some(Ok(Exit::Collected)));
954 }
955
956 #[tokio::test]
957 async fn routes_composed_products_in_algebra_order() {
958 let recipient = Recipient::global(MailAddr(9));
959 let sends = SendProduct {
960 inner: vec![Delivery::new(recipient, 3), Delivery::new(recipient, 4)],
961 own: vec![Delivery::new(recipient, 5)],
962 };
963 let mut router = RecordingRouter::default();
964
965 sends.route(MailAddr(7), &mut router).await.unwrap();
966
967 assert_eq!(
968 router
969 .0
970 .lock()
971 .expect("delivery lock")
972 .iter()
973 .map(|entry| entry.1.message)
974 .collect::<Vec<_>>(),
975 [3, 4, 5]
976 );
977 }
978
979 #[tokio::test]
980 async fn stops_at_the_first_delivery_failure() {
981 let recipient = Recipient::global(MailAddr(9));
982 let sends = vec![
983 Delivery::new(recipient, 1),
984 Delivery::new(recipient, 2),
985 Delivery::new(recipient, 3),
986 ];
987 let mut router = RecordingRouter::default();
988
989 assert_eq!(
990 sends.route(MailAddr(7), &mut router).await,
991 Err("delivery failed")
992 );
993 assert_eq!(
994 router
995 .0
996 .lock()
997 .expect("delivery lock")
998 .iter()
999 .map(|entry| entry.1.message)
1000 .collect::<Vec<_>>(),
1001 [1]
1002 );
1003 }
1004
1005 #[tokio::test]
1006 async fn product_failure_preserves_the_exact_interpreter_leg() {
1007 let recipient = Recipient::global(MailAddr(9));
1008 let mut router = RecordingRouter::default();
1009 let inner_failure = SendProduct {
1010 inner: vec![Delivery::new(recipient, 2)],
1011 own: vec![Delivery::new(recipient, 3)],
1012 };
1013 assert_eq!(
1014 inner_failure.route(MailAddr(7), &mut router).await,
1015 Err(SendProductError::Inner("delivery failed"))
1016 );
1017
1018 let own_failure = SendProduct {
1019 inner: vec![Delivery::new(recipient, 1)],
1020 own: vec![Delivery::new(recipient, 2)],
1021 };
1022 assert_eq!(
1023 own_failure.route(MailAddr(7), &mut router).await,
1024 Err(SendProductError::Own("delivery failed"))
1025 );
1026 }
1027}