Skip to main content

behavior_actors/
watch.rs

1//! Recurring logical-name peer observation.
2
3use crate::protocol::{ObservePeer, PeerStopped};
4use crate::{Crash, Exit, TerminationMonitorError, TerminationObservation};
5use behavior::Step;
6use behavior::{
7    Actions, Address, Become, Behavior, BehaviorActed, BirthMode, EventLayer, InterpreterRequests,
8    SendEffects, SendLayer, UserEvent,
9};
10
11/// Complete event sum accepted by a logical peer watch.
12pub type WatchEvent<E, Report = PeerStopped<<E as UserEvent>::Addr>> = EventLayer<Report, E>;
13
14/// Infallible reaction to one matching logical peer stop.
15pub type LinkReaction<B> = fn(
16    &mut B,
17    behavior::BehaviorAddr<B>,
18    &Result<Exit<behavior::BehaviorAddr<B>>, Crash>,
19) -> Become;
20
21/// Observe a logical peer name across any number of later incarnations.
22///
23/// A logical name can denote another incarnation after a stop, so every
24/// matching report is accepted and the watch remains observing. This recurrence
25/// law is distinct from [`crate::TerminationMonitor`], which consumes one
26/// correlated terminal relationship. Initialization emits one typed
27/// [`ObservePeer`] request after preserving the inner initialization effects.
28pub struct Watch<B: Behavior> {
29    inner: B,
30    peer: behavior::BehaviorAddr<B>,
31    on_stopped: LinkReaction<B>,
32}
33
34type WatchActions<B> = Actions<
35    behavior::BehaviorAddr<B>,
36    <B as Behavior>::Ph,
37    SendLayer<InterpreterRequests<ObservePeer<behavior::BehaviorAddr<B>>>, <B as Behavior>::Sends>,
38    <B as Behavior>::Birth,
39>;
40
41impl<B: Behavior> Watch<B> {
42    /// Wrap `inner` with recurring observation of one logical peer name.
43    #[must_use]
44    pub const fn new(
45        inner: B,
46        peer: behavior::BehaviorAddr<B>,
47        on_stopped: LinkReaction<B>,
48    ) -> Self {
49        Self {
50            inner,
51            peer,
52            on_stopped,
53        }
54    }
55
56    fn wrap(
57        actions: Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
58        observations: InterpreterRequests<ObservePeer<behavior::BehaviorAddr<B>>>,
59    ) -> WatchActions<B> {
60        actions.map_sends(|inner| SendLayer::new(observations, inner))
61    }
62}
63
64impl<B> behavior::BehaviorBase for Watch<B>
65where
66    B: Behavior + behavior::BehaviorBase,
67{
68    type Base = B::Base;
69
70    fn base(&self) -> &Self::Base {
71        self.inner.base()
72    }
73}
74
75impl<B> crate::StashStatus for Watch<B>
76where
77    B: Behavior + crate::StashStatus,
78{
79    fn stashed_messages(&self) -> usize {
80        self.inner.stashed_messages()
81    }
82}
83
84impl<B, A, Ph, Sends, Br> Behavior for Watch<B>
85where
86    A: Address,
87    Sends: SendEffects + behavior::SendsFor<B::Event>,
88    Br: BirthMode,
89    B: Behavior<Ph = Ph, Sends = Sends, Birth = Br>,
90    B::Protocol: behavior::Protocol<Addr = A>,
91{
92    type Protocol = B::Protocol;
93    type Event = WatchEvent<B::Event>;
94    type Sends = SendLayer<InterpreterRequests<ObservePeer<A>>, Sends>;
95    type Ph = Ph;
96    type Error = TerminationMonitorError<B::Error, PeerStopped<A>>;
97    type Birth = Br;
98
99    fn init(&mut self, _: behavior::InitializationTurn) -> BehaviorActed<Self> {
100        let actions =
101            behavior::initialize(&mut self.inner).map_err(TerminationMonitorError::Inner)?;
102        Ok(Self::wrap(
103            actions,
104            InterpreterRequests::one(ObservePeer::new(self.peer)),
105        ))
106    }
107
108    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
109        match event {
110            EventLayer::Owned(report) if report.peer == self.peer => {
111                let become_ = match (self.on_stopped)(&mut self.inner, report.peer, &report.outcome)
112                {
113                    Step::Continue => Step::Continue,
114                    Step::Goto(never) => match never {},
115                    Step::Stop(stopped) => Step::Stop(stopped),
116                };
117                Ok(Self::wrap(
118                    Actions::just(become_),
119                    InterpreterRequests::empty(),
120                ))
121            }
122            EventLayer::Owned(report) => Err(TerminationMonitorError::UnexpectedReport {
123                observation: TerminationObservation::Observing,
124                report,
125            }),
126            EventLayer::Inner(event) => behavior::delegate_transition(&mut self.inner, event)
127                .map(|actions| Self::wrap(actions, InterpreterRequests::empty()))
128                .map_err(TerminationMonitorError::Inner),
129        }
130    }
131}
132
133/// Stop when a logical watch reports an abnormal outcome.
134pub fn stop_on_abnormal_death<B: Behavior>(
135    _behavior: &mut B,
136    _peer: behavior::BehaviorAddr<B>,
137    outcome: &Result<Exit<behavior::BehaviorAddr<B>>, Crash>,
138) -> Become {
139    if let Ok(Exit::Normal | Exit::Collected) = outcome {
140        Step::Continue
141    } else {
142        Step::Stop(behavior::Stopped)
143    }
144}