1use 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
188pub 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}