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