1use std::time::Duration;
4
5use super::super::domain::{Fleet, RestartBudget};
6use super::super::policy::{
7 RestartPolicy, Strategy, SupervisionFailure, SupervisionFailureReaction,
8 retire_on_supervision_failure,
9};
10use super::super::protocol::{ProxyCommand, SupervisionEvent};
11use super::proxy::Proxy;
12use crate::behavior::{
13 Actions, Address, Behavior, Births, Create, Delivery, Recipient, SendAlgebra, ServiceSends,
14};
15use crate::next::{Never, Step};
16use crate::protocol::{
17 ChildEvent, CreationEvent, ObserveChild, WorkerCreationEvent, WorkerStopped,
18};
19use crate::{Become, Exit, SupervisionFailureReason};
20use crate::{Inner, Own, SendInput};
21
22pub struct SupervisorSends<A, Sends, C>
24where
25 A: Address,
26 A::Nonce: From<u64>,
27 C: Behavior<Addr = A, Ph = Never>,
28{
29 pub behavior: Sends,
30 pub child_observations: ServiceSends<ObserveChild<A::Nonce>>,
31 pub replacement_commands: Vec<Delivery<Proxy<C>>>,
32}
33
34impl<A, Sends, C> SendAlgebra for SupervisorSends<A, Sends, C>
35where
36 A: Address,
37 A::Nonce: From<u64>,
38 Sends: SendAlgebra,
39 C: Behavior<Addr = A, Ph = Never>,
40{
41 fn empty() -> Self {
42 Self {
43 behavior: Sends::empty(),
44 child_observations: ServiceSends::empty(),
45 replacement_commands: Vec::new(),
46 }
47 }
48
49 fn append(&mut self, other: Self) {
50 self.behavior.append(other.behavior);
51 self.child_observations.append(other.child_observations);
52 self.replacement_commands.extend(other.replacement_commands);
53 }
54}
55
56impl<A, Sends, C> SendInput<ObserveChild<A::Nonce>, Own> for SupervisorSends<A, Sends, C>
57where
58 A: Address,
59 A::Nonce: From<u64>,
60 C: Behavior<Addr = A, Ph = Never>,
61{
62 fn emit(&mut self, input: ObserveChild<A::Nonce>) {
63 self.child_observations.send(input);
64 }
65}
66
67impl<A, Sends, C> SendInput<Delivery<Proxy<C>>, Own> for SupervisorSends<A, Sends, C>
68where
69 A: Address,
70 A::Nonce: From<u64>,
71 C: Behavior<Addr = A, Ph = Never>,
72{
73 fn emit(&mut self, input: Delivery<Proxy<C>>) {
74 self.replacement_commands.push(input);
75 }
76}
77
78impl<A, Sends, C, Input, Path> SendInput<Input, Inner<Path>> for SupervisorSends<A, Sends, C>
79where
80 A: Address,
81 A::Nonce: From<u64>,
82 C: Behavior<Addr = A, Ph = Never>,
83 Sends: SendInput<Input, Path>,
84{
85 fn emit(&mut self, input: Input) {
86 <Sends as SendInput<Input, Path>>::emit(&mut self.behavior, input);
87 }
88}
89
90pub type SupervisorActions<B, C> = Actions<
91 <B as Behavior>::Addr,
92 <B as Behavior>::Ph,
93 SupervisorSends<<B as Behavior>::Addr, <B as Behavior>::Sends, C>,
94 Births<Proxy<C>>,
95>;
96
97enum ReplacementDecision<A, C>
98where
99 A: Address,
100 A::Nonce: From<u64>,
101 C: Behavior<Addr = A, Ph = Never>,
102{
103 Retire,
104 Replace(Vec<Delivery<Proxy<C>>>),
105 Failed(SupervisionFailure<A>),
106}
107
108pub struct Supervisor<B: Behavior, C: Behavior<Ph = Never, Addr = B::Addr>> {
109 inner: B,
110 fleet: Fleet<<B::Addr as Address>::Nonce>,
111 build: fn(usize) -> C,
112 strategy: Strategy,
113 policy: RestartPolicy,
114 budget: RestartBudget,
115 on_failure: SupervisionFailureReaction<B>,
116}
117
118impl<B, C> Supervisor<B, C>
119where
120 B: Behavior<Birth = Births<C>>,
121 <B::Addr as Address>::Nonce: From<u64>,
122 C: Behavior<Ph = Never, Addr = B::Addr>,
123{
124 #[allow(clippy::too_many_arguments, reason = "hidden by Compose")]
125 #[must_use]
132 pub fn new(
133 inner: B,
134 nonces: fn(usize) -> <B::Addr as Address>::Nonce,
135 count: usize,
136 build: fn(usize) -> C,
137 strategy: Strategy,
138 policy: RestartPolicy,
139 max_restarts: u32,
140 window: Duration,
141 ) -> Self {
142 let fleet = Fleet::configured((0..count).map(nonces))
143 .unwrap_or_else(|_| panic!("configured child nonces must be fresh"));
144 Self {
145 inner,
146 fleet,
147 build,
148 strategy,
149 policy,
150 budget: RestartBudget::new(max_restarts, window),
151 on_failure: retire_on_supervision_failure::<B>,
152 }
153 }
154
155 #[must_use]
156 pub fn with_strategy(mut self, strategy: Strategy) -> Self {
157 self.strategy = strategy;
158 self
159 }
160
161 #[must_use]
162 pub fn with_policy(mut self, policy: RestartPolicy) -> Self {
163 self.policy = policy;
164 self
165 }
166
167 #[must_use]
168 pub fn with_budget(mut self, max: u32, window: Duration) -> Self {
169 self.budget = RestartBudget::new(max, window);
170 self
171 }
172
173 #[must_use]
174 pub fn with_failure_reaction(mut self, reaction: SupervisionFailureReaction<B>) -> Self {
176 self.on_failure = reaction;
177 self
178 }
179
180 #[must_use]
181 pub fn is_alive(&self, nonce: <B::Addr as Address>::Nonce) -> bool {
186 self.fleet
187 .is_available(nonce)
188 .unwrap_or_else(|_| panic!("unknown supervised nonce"))
189 }
190
191 #[must_use]
192 pub fn child_count(&self) -> usize {
193 self.fleet.len()
194 }
195
196 #[must_use]
197 pub fn restarts_in_window(&self) -> usize {
198 self.budget.admitted()
199 }
200
201 fn replacement_decision(
202 &mut self,
203 event: &WorkerStopped<B::Addr>,
204 ) -> ReplacementDecision<B::Addr, C> {
205 let eligible = match self.policy {
206 RestartPolicy::Permanent => true,
207 RestartPolicy::Transient => {
208 !matches!(&event.outcome, Ok(Exit::Normal | Exit::Collected))
209 }
210 RestartPolicy::Temporary => false,
211 };
212 if !eligible {
213 self.fleet
214 .retire(event.proxy)
215 .unwrap_or_else(|_| panic!("unknown supervised nonce"));
216 return ReplacementDecision::Retire;
217 }
218 let candidates = self
219 .fleet
220 .replacements(event.proxy, self.strategy)
221 .unwrap_or_else(|_| panic!("unknown supervised nonce"));
222 if let Err(reason) = self.budget.admit(event.at, candidates.len()) {
223 self.fleet
224 .retire(event.proxy)
225 .unwrap_or_else(|_| panic!("unknown supervised nonce"));
226 return ReplacementDecision::Failed(SupervisionFailure::new(
227 event.proxy,
228 event.outcome,
229 SupervisionFailureReason::RestartDenied(reason),
230 ));
231 }
232 for candidate in &candidates {
233 self.fleet
234 .replacement_requested(candidate.nonce)
235 .unwrap_or_else(|_| unreachable!("candidate belongs to fleet"));
236 }
237 ReplacementDecision::Replace(
238 candidates
239 .into_iter()
240 .map(|candidate| {
241 Delivery::new(
242 Recipient::child(candidate.nonce),
243 ProxyCommand::Replace((self.build)(candidate.index)),
244 )
245 })
246 .collect(),
247 )
248 }
249
250 fn react_to_failure(
251 &mut self,
252 failure: &SupervisionFailure<B::Addr>,
253 ) -> Result<Become<B::Addr, B::Ph>, B::Error> {
254 Ok(match (self.on_failure)(&mut self.inner, failure)? {
255 Step::Continue => Step::Continue,
256 Step::Goto(never) => match never {},
257 Step::Stop(exit) => Step::Stop(exit),
258 })
259 }
260
261 fn wrap(
262 &mut self,
263 actions: Actions<B::Addr, B::Ph, B::Sends, Births<C>>,
264 ) -> SupervisorActions<B, C> {
265 let born: Vec<_> = actions.creates.iter().map(|create| create.nonce).collect();
266 for create in &actions.creates {
267 self.fleet
268 .register(create.nonce)
269 .unwrap_or_else(|_| panic!("a child birth nonce must be fresh"));
270 }
271 Actions::new(
272 SupervisorSends {
273 behavior: actions.sends,
274 child_observations: ServiceSends::new(
275 born.into_iter().map(ObserveChild::new).collect(),
276 ),
277 replacement_commands: Vec::new(),
278 },
279 actions
280 .creates
281 .into_iter()
282 .map(|create| Create::new(create.nonce, Proxy::new(create.child), create.kind))
283 .collect(),
284 actions.become_,
285 )
286 }
287}
288
289impl<B, C, A, Ph, Sends> Behavior for Supervisor<B, C>
290where
291 A: Address,
292 Sends: SendAlgebra,
293 B: Behavior<Addr = A, Ph = Ph, Sends = Sends, Birth = Births<C>>,
294 B::Event: ChildEvent + CreationEvent + WorkerCreationEvent,
295 A::Nonce: From<u64>,
296 C: Behavior<Ph = Never, Addr = B::Addr>,
297{
298 type Addr = A;
299 type Msg = B::Msg;
300 type Event = SupervisionEvent<B::Event>;
301 type Sends = SupervisorSends<A, Sends, C>;
302 type Ph = Ph;
303 type Error = B::Error;
304 type Birth = Births<Proxy<C>>;
305
306 fn init(&mut self) -> Result<SupervisorActions<B, C>, B::Error> {
307 let actions = self.inner.init()?;
308 let mut actions = self.wrap(actions);
309 actions.creates.extend(
310 self.fleet
311 .configured_nonces()
312 .enumerate()
313 .map(|(index, nonce)| Create::birth(nonce, Proxy::new((self.build)(index)))),
314 );
315 actions
316 .sends
317 .child_observations
318 .extend(self.fleet.configured_nonces().map(ObserveChild::new));
319 Ok(actions)
320 }
321
322 fn transition(&mut self, event: Self::Event) -> Result<SupervisorActions<B, C>, B::Error> {
323 match event {
324 SupervisionEvent::WorkerStopped(event) => {
325 let decision = self.replacement_decision(&event);
326 match decision {
327 ReplacementDecision::Retire => Ok(Actions::cont()),
328 ReplacementDecision::Replace(replacements) => Ok(Actions::new(
329 SupervisorSends {
330 behavior: B::Sends::empty(),
331 child_observations: ServiceSends::empty(),
332 replacement_commands: replacements,
333 },
334 Vec::new(),
335 Step::Continue,
336 )),
337 ReplacementDecision::Failed(failure) => {
338 Ok(Actions::just(self.react_to_failure(&failure)?))
339 }
340 }
341 }
342 SupervisionEvent::ChildStopped(event) => {
343 self.fleet
344 .retire(event.nonce)
345 .unwrap_or_else(|_| panic!("unknown supervised nonce"));
346 let failure = SupervisionFailure::new(
347 event.nonce,
348 event.outcome,
349 SupervisionFailureReason::StableChildStopped,
350 );
351 Ok(Actions::just(self.react_to_failure(&failure)?))
352 }
353 SupervisionEvent::CreationResolved(event) => {
354 self.fleet.resolve_creation(event.nonce, event.result);
355 if let Some(event) = B::Event::creation_resolved(event) {
356 let actions = self.inner.transition(event)?;
357 Ok(self.wrap(actions))
358 } else {
359 Ok(Actions::cont())
360 }
361 }
362 SupervisionEvent::WorkerCreationResolved(event) => {
363 if let Some(event) = B::Event::worker_creation_resolved(event) {
367 let actions = self.inner.transition(event)?;
368 Ok(self.wrap(actions))
369 } else {
370 Ok(Actions::cont())
371 }
372 }
373 SupervisionEvent::Inner(event) => {
374 let actions = self.inner.transition(event)?;
375 Ok(self.wrap(actions))
376 }
377 }
378 }
379}