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