1use std::time::Duration;
4
5use super::super::domain::{Fleet, FleetError, 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 ChildStopped, CreationResolved, ObserveChild, WorkerCreationResolved, WorkerStopped,
18};
19use crate::{Become, Exit, SupervisionFailureReason};
20use crate::{Own, RouteInput, 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
78pub(crate) type SupervisorActions<B, C> = Actions<
79 <B as Behavior>::Addr,
80 <B as Behavior>::Ph,
81 SupervisorSends<<B as Behavior>::Addr, <B as Behavior>::Sends, C>,
82 Births<Proxy<C>>,
83>;
84
85#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
87pub enum SupervisorError<E, N> {
88 #[error("supervised behavior rejected the transition")]
90 Behavior(E),
91 #[error(transparent)]
93 Fleet(#[from] FleetError<N>),
94 #[error("worker factory rejected configured fleet index {index}")]
96 FactoryIndex { index: usize },
97}
98
99enum ReplacementDecision<A, C>
100where
101 A: Address,
102 A::Nonce: From<u64>,
103 C: Behavior<Addr = A, Ph = Never>,
104{
105 Retire,
106 Replace(Vec<Delivery<Proxy<C>>>),
107 Failed(SupervisionFailure<A>),
108}
109
110pub struct Supervisor<B: Behavior, C: Behavior<Ph = Never, Addr = B::Addr>> {
111 inner: B,
112 fleet: Fleet<<B::Addr as Address>::Nonce>,
113 build: fn(usize) -> Option<C>,
114 strategy: Strategy,
115 policy: RestartPolicy,
116 budget: RestartBudget,
117 on_failure: SupervisionFailureReaction<B>,
118}
119
120impl<B, C> crate::BehaviorBase for Supervisor<B, C>
121where
122 B: Behavior<Birth = Births<C>> + crate::BehaviorBase,
123 <B::Addr as Address>::Nonce: From<u64>,
124 C: Behavior<Ph = Never, Addr = B::Addr>,
125{
126 type Base = B::Base;
127
128 fn base(&self) -> &Self::Base {
129 self.inner.base()
130 }
131}
132
133impl<B, C> crate::StashStatus for Supervisor<B, C>
134where
135 B: Behavior<Birth = Births<C>> + crate::StashStatus,
136 <B::Addr as Address>::Nonce: From<u64>,
137 C: Behavior<Ph = Never, Addr = B::Addr>,
138{
139 fn stashed_messages(&self) -> usize {
140 self.inner.stashed_messages()
141 }
142}
143
144impl<B, C> Supervisor<B, C>
145where
146 B: Behavior<Birth = Births<C>>,
147 <B::Addr as Address>::Nonce: From<u64>,
148 C: Behavior<Ph = Never, Addr = B::Addr>,
149{
150 #[allow(clippy::too_many_arguments, reason = "hidden by Compose")]
151 pub fn new(
158 inner: B,
159 nonces: fn(usize) -> <B::Addr as Address>::Nonce,
160 count: usize,
161 build: fn(usize) -> Option<C>,
162 strategy: Strategy,
163 policy: RestartPolicy,
164 max_restarts: u32,
165 window: Duration,
166 ) -> Result<Self, FleetError<<B::Addr as Address>::Nonce>> {
167 let fleet = Fleet::configured((0..count).map(nonces))?;
168 Ok(Self {
169 inner,
170 fleet,
171 build,
172 strategy,
173 policy,
174 budget: RestartBudget::new(max_restarts, window),
175 on_failure: retire_on_supervision_failure::<B>,
176 })
177 }
178
179 #[must_use]
180 pub fn with_strategy(mut self, strategy: Strategy) -> Self {
181 self.strategy = strategy;
182 self
183 }
184
185 #[must_use]
186 pub fn with_policy(mut self, policy: RestartPolicy) -> Self {
187 self.policy = policy;
188 self
189 }
190
191 #[must_use]
192 pub fn with_budget(mut self, max: u32, window: Duration) -> Self {
193 self.budget = RestartBudget::new(max, window);
194 self
195 }
196
197 #[must_use]
198 pub fn with_failure_reaction(mut self, reaction: SupervisionFailureReaction<B>) -> Self {
200 self.on_failure = reaction;
201 self
202 }
203
204 #[must_use]
205 pub fn is_alive(
210 &self,
211 nonce: <B::Addr as Address>::Nonce,
212 ) -> Result<bool, SupervisorError<core::convert::Infallible, <B::Addr as Address>::Nonce>> {
213 Ok(self.fleet.is_available(nonce)?)
214 }
215
216 #[must_use]
217 pub fn child_count(&self) -> usize {
218 self.fleet.len()
219 }
220
221 #[must_use]
222 pub fn restarts_in_window(&self) -> usize {
223 self.budget.admitted()
224 }
225
226 fn replacement_decision(
227 &mut self,
228 event: &WorkerStopped<B::Addr>,
229 ) -> Result<
230 ReplacementDecision<B::Addr, C>,
231 SupervisorError<B::Error, <B::Addr as Address>::Nonce>,
232 > {
233 let policy = self.policy;
234 let strategy = self.strategy;
235 let eligible = match policy {
236 RestartPolicy::Permanent => true,
237 RestartPolicy::Transient => {
238 !matches!(&event.outcome, Ok(Exit::Normal | Exit::Collected))
239 }
240 RestartPolicy::Temporary => false,
241 };
242 if !eligible {
243 self.fleet.retire(event.proxy)?;
244 return Ok(ReplacementDecision::Retire);
245 }
246 let candidates = self.fleet.replacements(event.proxy, strategy)?;
247 let replacements = candidates
248 .iter()
249 .map(|candidate| {
250 (self.build)(candidate.index)
251 .map(|child| (candidate.nonce, child))
252 .ok_or(SupervisorError::FactoryIndex {
253 index: candidate.index,
254 })
255 })
256 .collect::<Result<Vec<_>, _>>()?;
257 if let Err(reason) = self.budget.admit(event.at, candidates.len()) {
258 self.fleet.retire(event.proxy)?;
259 return Ok(ReplacementDecision::Failed(SupervisionFailure::new(
260 event.proxy,
261 event.outcome,
262 SupervisionFailureReason::RestartDenied(reason),
263 )));
264 }
265 for candidate in &candidates {
266 self.fleet.replacement_requested(candidate.nonce)?;
267 }
268 Ok(ReplacementDecision::Replace(
269 replacements
270 .into_iter()
271 .map(|(nonce, child)| {
272 Delivery::new(Recipient::child(nonce), ProxyCommand::Replace(child))
273 })
274 .collect(),
275 ))
276 }
277
278 fn react_to_failure(
279 &mut self,
280 failure: &SupervisionFailure<B::Addr>,
281 ) -> Result<Become<B::Addr, B::Ph>, B::Error> {
282 Ok(match (self.on_failure)(&mut self.inner, failure)? {
283 Step::Continue => Step::Continue,
284 Step::Goto(never) => match never {},
285 Step::Stop(exit) => Step::Stop(exit),
286 })
287 }
288
289 fn wrap(
290 &mut self,
291 actions: Actions<B::Addr, B::Ph, B::Sends, Births<C>>,
292 ) -> Result<SupervisorActions<B, C>, SupervisorError<B::Error, <B::Addr as Address>::Nonce>>
293 {
294 let fleet = &mut self.fleet;
295 let born: Vec<_> = actions.creates.iter().map(|create| create.nonce).collect();
296 for create in &actions.creates {
297 fleet.register(create.nonce)?;
298 }
299 Ok(Actions::new(
300 SupervisorSends {
301 behavior: actions.sends,
302 child_observations: ServiceSends::new(
303 born.into_iter().map(ObserveChild::new).collect(),
304 ),
305 replacement_commands: Vec::new(),
306 },
307 actions
308 .creates
309 .into_iter()
310 .map(|create| Create::new(create.nonce, Proxy::new(create.child), create.kind))
311 .collect(),
312 actions.become_,
313 ))
314 }
315}
316
317impl<B, C, A, Ph, Sends> Behavior for Supervisor<B, C>
318where
319 A: Address,
320 Sends: SendAlgebra,
321 B: Behavior<Addr = A, Ph = Ph, Sends = Sends, Birth = Births<C>>,
322 B::Event: crate::RouteInput<ChildStopped<A>>
323 + crate::RouteInput<CreationResolved<A::Nonce>>
324 + crate::RouteInput<WorkerCreationResolved<A::Nonce>>,
325 A::Nonce: From<u64>,
326 C: Behavior<Ph = Never, Addr = B::Addr>,
327{
328 type Addr = A;
329 type Msg = B::Msg;
330 type Event = SupervisionEvent<B::Event>;
331 type Sends = SupervisorSends<A, Sends, C>;
332 type Ph = Ph;
333 type Error = SupervisorError<B::Error, A::Nonce>;
334 type Birth = Births<Proxy<C>>;
335
336 fn init(
337 &mut self,
338 _: crate::InitializationTurn,
339 ) -> Result<SupervisorActions<B, C>, Self::Error> {
340 let configured: Vec<_> = self.fleet.configured_nonces().collect();
341 let workers = configured
342 .iter()
343 .copied()
344 .enumerate()
345 .map(|(index, nonce)| {
346 (self.build)(index)
347 .map(|worker| (nonce, worker))
348 .ok_or(SupervisorError::FactoryIndex { index })
349 })
350 .collect::<Result<Vec<_>, _>>()?;
351 let actions =
352 crate::calculus::initialize(&mut self.inner).map_err(SupervisorError::Behavior)?;
353 let mut actions = self.wrap(actions)?;
354 actions.creates.extend(
355 workers
356 .into_iter()
357 .map(|(nonce, worker)| Create::birth(nonce, Proxy::new(worker))),
358 );
359 actions
360 .sends
361 .child_observations
362 .extend(configured.into_iter().map(ObserveChild::new));
363 Ok(actions)
364 }
365
366 fn transition(
367 &mut self,
368 _: crate::ActiveTurn,
369 event: Self::Event,
370 ) -> Result<SupervisorActions<B, C>, Self::Error> {
371 match event {
372 SupervisionEvent::WorkerStopped(event) => {
373 let decision = self.replacement_decision(&event)?;
374 match decision {
375 ReplacementDecision::Retire => Ok(Actions::cont()),
376 ReplacementDecision::Replace(replacements) => Ok(Actions::new(
377 SupervisorSends {
378 behavior: B::Sends::empty(),
379 child_observations: ServiceSends::empty(),
380 replacement_commands: replacements,
381 },
382 Vec::new(),
383 Step::Continue,
384 )),
385 ReplacementDecision::Failed(failure) => Ok(Actions::just(
386 self.react_to_failure(&failure)
387 .map_err(SupervisorError::Behavior)?,
388 )),
389 }
390 }
391 SupervisionEvent::ChildStopped(event) => {
392 self.fleet.retire(event.nonce)?;
393 let failure = SupervisionFailure::new(
394 event.nonce,
395 event.outcome,
396 SupervisionFailureReason::StableChildStopped,
397 );
398 Ok(Actions::just(
399 self.react_to_failure(&failure)
400 .map_err(SupervisorError::Behavior)?,
401 ))
402 }
403 SupervisionEvent::CreationResolved(event) => {
404 self.fleet.resolve_creation(event.nonce, event.result);
405 if let Ok(event) = B::Event::route(event) {
406 let actions = crate::calculus::delegate_transition(&mut self.inner, event)
407 .map_err(SupervisorError::Behavior)?;
408 self.wrap(actions)
409 } else {
410 Ok(Actions::cont())
411 }
412 }
413 SupervisionEvent::WorkerCreationResolved(event) => {
414 if let Ok(event) = B::Event::route(event) {
418 let actions = crate::calculus::delegate_transition(&mut self.inner, event)
419 .map_err(SupervisorError::Behavior)?;
420 self.wrap(actions)
421 } else {
422 Ok(Actions::cont())
423 }
424 }
425 SupervisionEvent::Behavior(event) => {
426 let actions = crate::calculus::delegate_transition(&mut self.inner, event)
427 .map_err(SupervisorError::Behavior)?;
428 self.wrap(actions)
429 }
430 }
431 }
432}