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