Skip to main content

behavior/supervision/adapter/
supervisor.rs

1//! Fleet coordination for supervised stable proxy actors.
2
3use 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
22/// Named effect lanes emitted by a supervised behavior.
23pub 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/// A controlled supervisor-fold failure.
86#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
87pub enum SupervisorError<E, N> {
88    /// The supervised behavior rejected its fold.
89    #[error("supervised behavior rejected the transition")]
90    Behavior(E),
91    /// The supervisor's child topology rejected the operation.
92    #[error(transparent)]
93    Fleet(#[from] FleetError<N>),
94    /// The configured worker factory did not define a requested fleet index.
95    #[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    /// Construct the concrete supervisor behavior hidden by `Compose`.
152    ///
153    /// Invalid configured routes reject construction before a behavior exists.
154    ///
155    /// # Errors
156    /// Returns the first typed topology rejection.
157    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    /// Replace the pure reaction used for typed supervision failures.
199    pub fn with_failure_reaction(mut self, reaction: SupervisionFailureReaction<B>) -> Self {
200        self.on_failure = reaction;
201        self
202    }
203
204    #[must_use]
205    /// Report whether a known supervised proxy is alive.
206    ///
207    /// # Errors
208    /// Returns the unknown nonce when it is not part of this topology.
209    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                // Worker realization does not change the stable proxy's
415                // liveness. The typed result remains distinct from a proxy
416                // terminal observation.
417                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}