Skip to main content

behavior/
watching.rs

1//! Pure peer-observation composition over an ordinary monitor actor protocol.
2
3use crate::behavior::{
4    Actions, Address, Become, Behavior, BirthMode, SendAlgebra, SendProduct, ServiceSends, User,
5    UserEvent,
6};
7use crate::protocol::{
8    ChildEvent, ChildStopped, ObservePeer, PeerEvent, PeerStopped, ShutdownEvent,
9    ShutdownRequested, TimeEvent, TimerElapsed, WorkerEvent, WorkerStopped,
10};
11use crate::{Crash, Exit, Step};
12
13#[derive(Debug, Clone, PartialEq, Eq)]
14pub enum WatchEvent<E, A: Address> {
15    Inner(E),
16    PeerStopped(PeerStopped<A>),
17}
18
19impl<E, A: Address> PeerEvent<A> for WatchEvent<E, A> {
20    fn peer_stopped(event: PeerStopped<A>) -> Option<Self> {
21        Some(Self::PeerStopped(event))
22    }
23}
24
25impl<E: UserEvent, A: Address> UserEvent for WatchEvent<E, A> {
26    type Addr = E::Addr;
27    type Message = E::Message;
28
29    fn user(from: Self::Addr, message: Self::Message) -> Self {
30        Self::Inner(E::user(from, message))
31    }
32
33    fn into_user(self) -> Result<User<Self::Addr, Self::Message>, Self> {
34        match self {
35            Self::Inner(event) => event.into_user().map_err(Self::Inner),
36            stopped @ Self::PeerStopped(_) => Err(stopped),
37        }
38    }
39}
40
41impl<E: TimeEvent, A: Address> TimeEvent for WatchEvent<E, A> {
42    fn time_reached(event: TimerElapsed) -> Option<Self> {
43        E::time_reached(event).map(Self::Inner)
44    }
45}
46
47impl<E: ChildEvent<A>, A: Address> ChildEvent<A> for WatchEvent<E, A> {
48    fn child_stopped(event: ChildStopped<A>) -> Option<Self> {
49        E::child_stopped(event).map(Self::Inner)
50    }
51}
52
53impl<E: WorkerEvent<A>, A: Address> WorkerEvent<A> for WatchEvent<E, A> {
54    fn worker_stopped(event: WorkerStopped<A>) -> Option<Self> {
55        E::worker_stopped(event).map(Self::Inner)
56    }
57}
58
59impl<E: ShutdownEvent, A: Address> ShutdownEvent for WatchEvent<E, A> {
60    fn shutdown_requested(event: ShutdownRequested) -> Option<Self> {
61        E::shutdown_requested(event).map(Self::Inner)
62    }
63}
64
65pub type LinkReaction<B> = fn(
66    &mut B,
67    <B as Behavior>::Addr,
68    &Result<Exit<<B as Behavior>::Addr>, Crash>,
69) -> Result<Become<<B as Behavior>::Addr>, <B as Behavior>::Error>;
70
71pub type WatchSends<B> =
72    SendProduct<<B as Behavior>::Sends, ServiceSends<ObservePeer<<B as Behavior>::Addr>>>;
73
74pub type WatchActions<B> =
75    Actions<<B as Behavior>::Addr, <B as Behavior>::Ph, WatchSends<B>, <B as Behavior>::Birth>;
76
77pub struct Watching<B: Behavior> {
78    inner: B,
79    peer: B::Addr,
80    on_stopped: LinkReaction<B>,
81}
82
83impl<B: Behavior> Watching<B> {
84    #[must_use]
85    pub fn new(inner: B, peer: B::Addr, on_stopped: LinkReaction<B>) -> Self {
86        Self {
87            inner,
88            peer,
89            on_stopped,
90        }
91    }
92
93    #[must_use]
94    pub fn inner(&self) -> &B {
95        &self.inner
96    }
97}
98
99impl<B, A, Ph, Sends, Br> Behavior for Watching<B>
100where
101    A: Address + Send,
102    Sends: SendAlgebra,
103    Br: BirthMode,
104    B: Behavior<Addr = A, Ph = Ph, Sends = Sends, Birth = Br> + Send,
105    B::Event: PeerEvent<B::Addr> + Send,
106    B::Msg: Send,
107{
108    type Addr = A;
109    type Msg = B::Msg;
110    type Event = WatchEvent<B::Event, B::Addr>;
111    type Sends = SendProduct<Sends, ServiceSends<ObservePeer<A>>>;
112    type Ph = Ph;
113    type Error = B::Error;
114    type Birth = Br;
115
116    async fn init(&mut self) -> Result<WatchActions<B>, B::Error> {
117        let actions = self.inner.init().await?;
118        Ok(Self::wrap(
119            actions,
120            ServiceSends::one(ObservePeer { peer: self.peer }),
121        ))
122    }
123
124    async fn step(&mut self, event: Self::Event) -> Result<WatchActions<B>, B::Error> {
125        match event {
126            WatchEvent::PeerStopped(event) if event.peer == self.peer => {
127                let become_ = match (self.on_stopped)(&mut self.inner, event.peer, &event.outcome)?
128                {
129                    Step::Continue => Step::Continue,
130                    Step::Goto(never) => match never {},
131                    Step::Stop(exit) => Step::Stop(exit),
132                };
133                Ok(Actions {
134                    sends: Self::Sends::empty(),
135                    creates: Vec::new(),
136                    become_,
137                })
138            }
139            WatchEvent::PeerStopped(event) => match B::Event::peer_stopped(event) {
140                Some(inner) => self
141                    .inner
142                    .step(inner)
143                    .await
144                    .map(|actions| Self::wrap(actions, ServiceSends::empty())),
145                None => Ok(Actions::cont()),
146            },
147            WatchEvent::Inner(event) => self
148                .inner
149                .step(event)
150                .await
151                .map(|actions| Self::wrap(actions, ServiceSends::empty())),
152        }
153    }
154}
155
156impl<B: Behavior> Watching<B> {
157    fn wrap(
158        actions: Actions<B::Addr, B::Ph, B::Sends, B::Birth>,
159        own: ServiceSends<ObservePeer<B::Addr>>,
160    ) -> WatchActions<B> {
161        Actions {
162            sends: SendProduct {
163                inner: actions.sends,
164                own,
165            },
166            creates: actions.creates,
167            become_: actions.become_,
168        }
169    }
170}
171
172/// Stop when the monitor reports an abnormal outcome.
173///
174/// # Errors
175/// This supplied policy never creates a controlled error.
176pub fn stop_on_abnormal_death<B: Behavior>(
177    _behavior: &mut B,
178    peer: B::Addr,
179    outcome: &Result<Exit<B::Addr>, Crash>,
180) -> Result<Become<B::Addr>, B::Error> {
181    Ok(match outcome {
182        Ok(Exit::Normal | Exit::Collected) => Step::Continue,
183        _ => Step::Stop(Exit::LinkDied(peer)),
184    })
185}