Skip to main content

bombay/routing/
delivery.rs

1//! Shared outbound delivery routing.
2
3use 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
15/// Normalized terminal outcome exposed to Bombay Behavior's peer protocol.
16pub(crate) type PeerOutcome<A> = Result<Exit<A>, Crash>;
17
18/// A registered delivery endpoint paired with its exact incarnation outcome.
19#[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/// Failure to capture one currently registered peer incarnation.
60#[doc(hidden)]
61#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
62pub enum PeerObservationError<A> {
63    /// No live incarnation is registered at the observed address.
64    #[error("no live incarnation is registered at address {0:?}")]
65    Unknown(A),
66}
67
68/// Resolves and captures exact-generation peer completion without liveness.
69#[doc(hidden)]
70pub trait PeerObserver<A: Address> {
71    fn observe_peer(&self, peer: A)
72    -> Result<Observation<PeerOutcome<A>>, PeerObservationError<A>>;
73}
74
75/// A resolved destination for one statically typed message protocol.
76///
77/// This port performs no route interpretation or address lookup; those belong
78/// to [`DeliveryRouter`].
79#[doc(hidden)]
80pub trait DeliveryEndpoint<A, M> {
81    /// Delivery failure.
82    type Error;
83
84    /// Deliver `message` on behalf of `from`.
85    fn deliver(
86        &self,
87        from: A,
88        message: M,
89    ) -> impl Future<Output = Result<(), RejectedDelivery<M, Self::Error>>> + Send;
90}
91
92/// One endpoint rejection with ownership of the unchanged message.
93#[derive(Debug, Clone, PartialEq, Eq)]
94pub struct RejectedDelivery<M, E> {
95    /// The message the endpoint did not accept.
96    pub message: M,
97    /// The endpoint-specific rejection reason.
98    pub error: E,
99}
100
101impl<M, E> RejectedDelivery<M, E> {
102    /// Preserve one rejected message with its typed reason.
103    #[must_use]
104    pub const fn new(message: M, error: E) -> Self {
105        Self { message, error }
106    }
107}
108
109/// Registers resolved typed endpoints for later address routing.
110pub trait EndpointRegistry<A, M, D> {
111    /// Registration failure.
112    type Error;
113    /// Ownership token that keeps the exact registration generation live.
114    type Registration;
115
116    /// Register `endpoint` at `address`.
117    ///
118    /// # Errors
119    ///
120    /// Returns the registry's collision or storage failure.
121    fn register(&self, address: A, endpoint: D) -> Result<Self::Registration, Self::Error>;
122}
123
124/// Failure while resolving or invoking a registered endpoint.
125#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
126pub enum RoutingError<A: Address, M, E> {
127    /// No endpoint exists for this resolved address.
128    #[error("no endpoint exists for the resolved address {address:?}")]
129    UnknownAddress {
130        /// The concrete address obtained from the typed route.
131        address: A,
132        /// The unchanged message that was not accepted.
133        message: M,
134    },
135    /// The resolved endpoint rejected delivery.
136    #[error("the endpoint at {address:?} rejected delivery: {rejected:?}")]
137    Endpoint {
138        /// The concrete endpoint address selected by the route.
139        address: A,
140        /// The unchanged rejected message and typed endpoint reason.
141        rejected: RejectedDelivery<M, E>,
142    },
143}
144
145/// A simple shared routing table for one typed message protocol.
146pub 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
223/// Shared capability that resolves and delivers one typed behavior send.
224///
225/// Resolution of global, child-relative, and service routes is distinct from
226/// invoking the resulting [`DeliveryEndpoint`].
227pub trait DeliveryRouter<A: Address, M> {
228    /// Delivery failure.
229    type Error;
230
231    /// Deliver `delivery` on behalf of `from`.
232    fn deliver(
233        &self,
234        from: A,
235        delivery: Delivery<A, M>,
236    ) -> impl Future<Output = Result<(), Self::Error>> + Send;
237}
238
239/// Statically reveals whether a send algebra observes one creation nonce.
240///
241/// The interpreter consults this probe before routing a transition's sends:
242/// a creation failure at an observed nonce becomes a committed
243/// [`behavior::CreationResolved`] rejection, while an unobserved failure
244/// retains its exact runtime error. Custom send algebras must implement this
245/// alongside [`RouteSends`] so the two agree.
246pub trait ObservesCreations<N> {
247    /// Whether this algebra contains an `ObserveCreation` for `nonce`.
248    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
374/// Interprets a composed behavior send algebra through shared routing.
375///
376/// [`behavior::SendAlgebra`] only supplies pure `empty` and `append`
377/// composition. This runtime port consumes that accumulated value in fold
378/// order and reports the exact interpreter leg that failed.
379pub trait RouteSends<A: Address, R>: Sized {
380    /// Delivery failure.
381    type Error;
382
383    /// Route every send in fold order.
384    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/// Failure from one statically selected leg of a composed send algebra.
406#[doc(hidden)]
407#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
408pub enum SendProductError<L, R> {
409    /// The inner protocol leg failed.
410    #[error("inner protocol leg failed: {0:?}")]
411    Inner(L),
412    /// The wrapper-owned protocol leg failed.
413    #[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/// Failure from one named lane of a [`WatchSends`] composition.
439#[doc(hidden)]
440#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
441pub enum WatchSendsError<B, O> {
442    /// The wrapped behavior's own sends failed.
443    #[error("watch behavior sends failed: {0:?}")]
444    Behavior(B),
445    /// The watch observation lane failed.
446    #[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/// Failure from one named lane of a [`DeadlineSends`] composition.
475#[doc(hidden)]
476#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
477pub enum DeadlineSendsError<B, S> {
478    /// The wrapped behavior's own sends failed.
479    #[error("deadline behavior sends failed: {0:?}")]
480    Behavior(B),
481    /// The absolute-schedule lane failed.
482    #[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/// Failure from one named lane of a [`ReceiveTimeoutSends`] composition.
511#[doc(hidden)]
512#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
513pub enum ReceiveTimeoutSendsError<B, S> {
514    /// The wrapped behavior's own sends failed.
515    #[error("receive-timeout behavior sends failed: {0:?}")]
516    Behavior(B),
517    /// The relative-schedule lane failed.
518    #[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/// Failure from one named lane of a [`SupervisorSends`] composition.
547#[doc(hidden)]
548#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
549pub enum SupervisorSendsError<B, O, C> {
550    /// The supervised behavior's own sends failed.
551    #[error("supervised behavior sends failed: {0:?}")]
552    Behavior(B),
553    /// The child-observation lane failed.
554    #[error("supervisor child-observation lane failed: {0:?}")]
555    ChildObservations(O),
556    /// The replacement-command lane failed.
557    #[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/// Failure from one named lane of a [`ProxySends`] composition.
594#[doc(hidden)]
595#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
596pub enum ProxySendsError<D, C, CR, S, RP> {
597    /// The proxied user-delivery lane failed.
598    #[error("proxy delivery lane failed: {0:?}")]
599    Deliveries(D),
600    /// The child-observation lane failed.
601    #[error("proxy child-observation lane failed: {0:?}")]
602    ChildObservations(C),
603    /// The creation-observation lane failed.
604    #[error("proxy creation-observation lane failed: {0:?}")]
605    CreationObservations(CR),
606    /// The worker-stopped report lane failed.
607    #[error("proxy worker-stopped report lane failed: {0:?}")]
608    StoppedReports(S),
609    /// The worker-creation report lane failed.
610    #[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}