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