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, 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
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
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    /// Construct the concrete supervisor behavior hidden by `Compose`.
126    ///
127    /// # Panics
128    /// Panics when configured child nonces are not unique. Such a topology
129    /// would violate creator-local child routing and creation freshness before
130    /// the behavior could return an initialization result.
131    #[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    /// Replace the pure reaction used for typed supervision failures.
175    pub fn with_failure_reaction(mut self, reaction: SupervisionFailureReaction<B>) -> Self {
176        self.on_failure = reaction;
177        self
178    }
179
180    #[must_use]
181    /// Report whether a known supervised proxy is alive.
182    ///
183    /// # Panics
184    /// Panics when `nonce` is not part of this supervisor topology.
185    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                // Worker realization does not change the stable proxy's
364                // liveness. The typed result remains distinct from a proxy
365                // terminal observation.
366                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}