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};
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
66pub type LinkReaction<B> = fn(
67 &mut B,
68 <B as Behavior>::Addr,
69 &Result<Exit<<B as Behavior>::Addr>, Crash>,
70) -> Result<Become<<B as Behavior>::Addr>, <B as Behavior>::Error>;
71
72pub type WatchSends<B> =
73 SendProduct<<B as Behavior>::Sends, ServiceSends<ObservePeer<<B as Behavior>::Addr>>>;
74
75pub type WatchActions<B> =
76 Actions<<B as Behavior>::Addr, <B as Behavior>::Ph, WatchSends<B>, <B as Behavior>::Birth>;
77
78pub struct Watching<B: Behavior> {
79 inner: B,
80 peer: B::Addr,
81 on_stopped: LinkReaction<B>,
82}
83
84impl<B: Behavior> Watching<B> {
85 #[must_use]
86 pub fn new(inner: B, peer: B::Addr, on_stopped: LinkReaction<B>) -> Self {
87 Self {
88 inner,
89 peer,
90 on_stopped,
91 }
92 }
93
94 #[must_use]
95 pub fn inner(&self) -> &B {
96 &self.inner
97 }
98}
99
100impl<B, A, Ph, Sends, Br> Behavior for Watching<B>
101where
102 A: Address + Send,
103 Sends: SendAlgebra,
104 Br: BirthMode,
105 B: Behavior<
106 Addr = A,
107 Ph = Ph,
108 Sends = Sends,
109 Birth = Br,
110 Effect = Actions<A, Ph, Sends, Br>,
111 Done = Exit<A>,
112 > + Send,
113 B::Event: PeerEvent<B::Addr> + Send,
114 B::Msg: Send,
115{
116 type Addr = A;
117 type Msg = B::Msg;
118 type Event = WatchEvent<B::Event, B::Addr>;
119 type Sends = SendProduct<Sends, ServiceSends<ObservePeer<A>>>;
120 type Ph = Ph;
121 type Error = B::Error;
122 type Birth = Br;
123 type Effect = Actions<A, Ph, Self::Sends, Br>;
124 type Done = Exit<A>;
125
126 async fn init(&mut self) -> Result<Self::Effect, B::Error> {
127 let actions = self.inner.init().await?;
128 Ok(Self::wrap(
129 actions,
130 ServiceSends::one(ObservePeer { peer: self.peer }),
131 ))
132 }
133
134 async fn step(&mut self, event: Self::Event) -> Result<Self::Effect, B::Error> {
135 match event {
136 WatchEvent::PeerStopped(event) if event.peer == self.peer => {
137 let become_ = match (self.on_stopped)(&mut self.inner, event.peer, &event.outcome)?
138 {
139 Step::Continue => Step::Continue,
140 Step::Goto(never) => match never {},
141 Step::Stop(exit) => Step::Stop(exit),
142 };
143 Ok(Actions {
144 sends: Self::Sends::empty(),
145 creates: Vec::new(),
146 become_,
147 })
148 }
149 WatchEvent::PeerStopped(event) => match B::Event::peer_stopped(event) {
150 Some(inner) => self
151 .inner
152 .step(inner)
153 .await
154 .map(|actions| Self::wrap(actions, ServiceSends::empty())),
155 None => Ok(Actions::cont()),
156 },
157 WatchEvent::Inner(event) => self
158 .inner
159 .step(event)
160 .await
161 .map(|actions| Self::wrap(actions, ServiceSends::empty())),
162 }
163 }
164}
165
166impl<B: Behavior> Watching<B> {
167 fn wrap(
168 actions: Actions<B::Addr, B::Ph, B::Sends, B::Birth>,
169 own: ServiceSends<ObservePeer<B::Addr>>,
170 ) -> WatchActions<B> {
171 Actions {
172 sends: SendProduct {
173 inner: actions.sends,
174 own,
175 },
176 creates: actions.creates,
177 become_: actions.become_,
178 }
179 }
180}
181
182pub fn stop_on_abnormal_death<B: Behavior>(
187 _behavior: &mut B,
188 peer: B::Addr,
189 outcome: &Result<Exit<B::Addr>, Crash>,
190) -> Result<Become<B::Addr>, B::Error> {
191 Ok(match outcome {
192 Ok(Exit::Normal | Exit::Collected) => Step::Continue,
193 _ => Step::Stop(Exit::LinkDied(peer)),
194 })
195}