Skip to main content

behavior/
watch.rs

1//! Pure peer-observation composition over an ordinary monitor actor protocol.
2
3use crate::behavior::{
4    Actions, Address, Become, Behavior, BirthMode, SendAlgebra, ServiceSends, User, UserEvent,
5};
6use crate::protocol::forward::forward_event_lane;
7use crate::protocol::{ObservePeer, PeerStopped};
8use crate::{Crash, Exit, Step};
9use crate::{Own, RouteInput, SendInput};
10
11#[derive(Debug, Clone, PartialEq, Eq)]
12pub enum WatchEvent<E: UserEvent> {
13    Behavior(E),
14    PeerStopped(PeerStopped<E::Addr>),
15}
16
17impl<E: UserEvent> crate::RouteInput<PeerStopped<E::Addr>> for WatchEvent<E> {
18    fn route(event: PeerStopped<E::Addr>) -> Result<Self, PeerStopped<E::Addr>> {
19        Ok(Self::PeerStopped(event))
20    }
21}
22
23impl<E: UserEvent> crate::EventInput<PeerStopped<E::Addr>> for WatchEvent<E> {
24    fn inject(event: PeerStopped<E::Addr>) -> Self {
25        Self::PeerStopped(event)
26    }
27}
28
29impl<E: UserEvent> UserEvent for WatchEvent<E> {
30    type Addr = E::Addr;
31    type Message = E::Message;
32
33    fn user(from: Self::Addr, message: Self::Message) -> Self {
34        Self::Behavior(E::user(from, message))
35    }
36
37    fn into_user(self) -> Result<User<Self::Addr, Self::Message>, Self> {
38        match self {
39            Self::Behavior(event) => event.into_user().map_err(Self::Behavior),
40            stopped @ Self::PeerStopped(_) => Err(stopped),
41        }
42    }
43}
44
45forward_event_lane!(WatchEvent, crate::TimerElapsed);
46forward_event_lane!(WatchEvent, crate::ChildStopped<E::Addr>);
47forward_event_lane!(WatchEvent, crate::WorkerStopped<E::Addr>);
48forward_event_lane!(
49    WatchEvent,
50    crate::CreationResolved<<E::Addr as crate::Address>::Nonce>
51);
52forward_event_lane!(
53    WatchEvent,
54    crate::WorkerCreationResolved<<E::Addr as crate::Address>::Nonce>
55);
56forward_event_lane!(WatchEvent, crate::ShutdownRequested);
57
58pub type LinkReaction<B> = fn(
59    &mut B,
60    <B as Behavior>::Addr,
61    &Result<Exit<<B as Behavior>::Addr>, Crash>,
62) -> Result<Become<<B as Behavior>::Addr>, <B as Behavior>::Error>;
63
64/// Named effect lanes added by [`Watch`].
65#[derive(Debug, Clone, PartialEq, Eq)]
66pub struct WatchSends<A: Address, Sends> {
67    pub behavior: Sends,
68    pub observations: ServiceSends<ObservePeer<A>>,
69}
70
71impl<A: Address, Sends: SendAlgebra> SendAlgebra for WatchSends<A, Sends> {
72    fn empty() -> Self {
73        Self {
74            behavior: Sends::empty(),
75            observations: ServiceSends::empty(),
76        }
77    }
78
79    fn append(&mut self, other: Self) {
80        self.behavior.append(other.behavior);
81        self.observations.append(other.observations);
82    }
83}
84
85impl<A: Address, Sends> SendInput<ObservePeer<A>, Own> for WatchSends<A, Sends> {
86    fn emit(&mut self, input: ObservePeer<A>) {
87        self.observations.send(input);
88    }
89}
90
91pub(crate) type WatchActions<B> = Actions<
92    <B as Behavior>::Addr,
93    <B as Behavior>::Ph,
94    WatchSends<<B as Behavior>::Addr, <B as Behavior>::Sends>,
95    <B as Behavior>::Birth,
96>;
97
98/// A pure peer-observation transformation.
99///
100/// Initialization emits exactly one [`ObservePeer`] request after preserving
101/// the inner initialization effects. A matching [`PeerStopped`] result invokes
102/// the configured reaction whether the interpreter produced it immediately
103/// from authoritative retained termination or after observing a live
104/// incarnation. The transformation retains no runtime observation handle or
105/// lifecycle flag; exact-incarnation selection belongs to the interpreter.
106pub struct Watch<B: Behavior> {
107    inner: B,
108    peer: B::Addr,
109    on_stopped: LinkReaction<B>,
110}
111
112impl<B: Behavior> Watch<B> {
113    #[must_use]
114    pub(crate) fn new(inner: B, peer: B::Addr, on_stopped: LinkReaction<B>) -> Self {
115        Self {
116            inner,
117            peer,
118            on_stopped,
119        }
120    }
121}
122
123impl<B: Behavior + crate::BehaviorBase> crate::BehaviorBase for Watch<B> {
124    type Base = B::Base;
125
126    fn base(&self) -> &Self::Base {
127        self.inner.base()
128    }
129}
130
131impl<B> crate::StashStatus for Watch<B>
132where
133    B: Behavior + crate::StashStatus,
134{
135    fn stashed_messages(&self) -> usize {
136        self.inner.stashed_messages()
137    }
138}
139
140impl<B, A, Ph, Sends, Br> Behavior for Watch<B>
141where
142    A: Address,
143    Sends: SendAlgebra,
144    Br: BirthMode,
145    B: Behavior<Addr = A, Ph = Ph, Sends = Sends, Birth = Br>,
146    B::Event: crate::RouteInput<PeerStopped<A>>,
147{
148    type Addr = A;
149    type Msg = B::Msg;
150    type Event = WatchEvent<B::Event>;
151    type Sends = WatchSends<A, Sends>;
152    type Ph = Ph;
153    type Error = B::Error;
154    type Birth = Br;
155
156    fn init(&mut self, _: crate::InitializationTurn) -> Result<WatchActions<B>, B::Error> {
157        let actions = crate::calculus::initialize(&mut self.inner)?;
158        Ok(Self::wrap(actions, ServiceSends::one(self.peer.into())))
159    }
160
161    fn transition(
162        &mut self,
163        _: crate::ActiveTurn,
164        event: Self::Event,
165    ) -> Result<WatchActions<B>, B::Error> {
166        match event {
167            WatchEvent::PeerStopped(event) if event.peer == self.peer => {
168                let become_ = match (self.on_stopped)(&mut self.inner, event.peer, &event.outcome)?
169                {
170                    Step::Continue => Step::Continue,
171                    Step::Goto(never) => match never {},
172                    Step::Stop(exit) => Step::Stop(exit),
173                };
174                Ok(Actions::new(Self::Sends::empty(), Vec::new(), become_))
175            }
176            WatchEvent::PeerStopped(event) => match B::Event::route(event) {
177                Ok(inner) => crate::calculus::delegate_transition(&mut self.inner, inner)
178                    .map(|actions| Self::wrap(actions, ServiceSends::empty())),
179                Err(_) => Ok(Actions::cont()),
180            },
181            WatchEvent::Behavior(event) => {
182                crate::calculus::delegate_transition(&mut self.inner, event)
183                    .map(|actions| Self::wrap(actions, ServiceSends::empty()))
184            }
185        }
186    }
187}
188
189impl<B: Behavior> Watch<B> {
190    fn wrap(
191        actions: Actions<B::Addr, B::Ph, B::Sends, B::Birth>,
192        own: ServiceSends<ObservePeer<B::Addr>>,
193    ) -> WatchActions<B> {
194        actions.map_sends(|behavior| WatchSends {
195            behavior,
196            observations: own,
197        })
198    }
199}
200
201/// Stop when the monitor reports an abnormal outcome.
202///
203/// # Errors
204/// This supplied policy never creates a controlled error.
205pub fn stop_on_abnormal_death<B: Behavior>(
206    _behavior: &mut B,
207    peer: B::Addr,
208    outcome: &Result<Exit<B::Addr>, Crash>,
209) -> Result<Become<B::Addr>, B::Error> {
210    Ok(match outcome {
211        Ok(Exit::Normal | Exit::Collected) => Step::Continue,
212        _ => Step::Stop(Exit::LinkDied(peer)),
213    })
214}