saluki_core/runtime/
supervisor.rs

1use std::{
2    future::Future,
3    pin::Pin,
4    sync::{
5        atomic::{AtomicU64, AtomicUsize, Ordering},
6        Arc, Mutex,
7    },
8    time::Duration,
9};
10
11use saluki_common::sync::shutdown::ShutdownHandle;
12use saluki_error::GenericError;
13use serde::{Deserialize, Serialize};
14use snafu::{OptionExt as _, Snafu};
15use tokio::{pin, runtime::Handle, select, sync::mpsc};
16use tracing::{debug, error, warn};
17
18use super::{
19    dedicated::{spawn_dedicated_runtime, RuntimeConfiguration, RuntimeMode},
20    restart::{RestartAction, RestartMode, RestartState, RestartStrategy, RestartType},
21    supervisable::{InitializationError, ShutdownStrategy, Supervisable},
22    tree::{ChildFacts, ChildKey, NodeConfig, Roster, SupervisionTreeHandle, SupervisorNode},
23    worker_state::WorkerState,
24};
25use crate::runtime::{
26    process::{Process, ProcessExt as _},
27    state::DataspaceRegistry,
28};
29
30/// Process name segment used for a child whose own name can't be turned into a valid process name.
31///
32/// See [`SupervisedChild::create_process`].
33const UNNAMED_CHILD: &str = "unnamed";
34
35/// A `Future` that represents the full lifecycle of a worker, including initialization.
36///
37/// Unlike [`SupervisorFuture`], which only represents the runtime phase, this future first performs async
38/// initialization and then runs the worker. This allows initialization to happen concurrently when multiple workers are
39/// spawned, and keeps the supervisor loop responsive to shutdown signals during initialization.
40pub(super) type WorkerFuture = Pin<Box<dyn Future<Output = Result<(), WorkerError>> + Send>>;
41
42/// Worker lifecycle errors.
43///
44/// Distinguishes between initialization failures (which shouldn't trigger restart logic) and runtime failures (which
45/// are eligible for restart).
46#[derive(Debug)]
47pub(super) enum WorkerError {
48    /// The worker failed during async initialization.
49    ///
50    /// The optional `child_name` carries the name of the original failing child when the error originates from a
51    /// nested supervisor. This allows the parent to include it in its own `FailedToInitialize` error for better
52    /// diagnostics across supervision tree levels.
53    Initialization {
54        child_name: Option<String>,
55        source: InitializationError,
56    },
57
58    /// The worker failed during runtime execution.
59    Runtime(GenericError),
60
61    /// The worker was a nested supervisor that completed a requested shutdown after forcefully aborting one or more of
62    /// its own workers.
63    ///
64    /// Carried as a distinct variant (rather than collapsed into [`Runtime`][WorkerError::Runtime]) so the parent's
65    /// shutdown drain can recover the structured count and merge it into its own tally, aggregating forced aborts up
66    /// the supervision tree.
67    ShutdownTimedOut {
68        /// The number of workers the nested supervisor forcefully aborted, summed across its own supervision tree.
69        aborted: usize,
70    },
71}
72
73impl From<SupervisorError> for WorkerError {
74    fn from(err: SupervisorError) -> Self {
75        match err {
76            // Propagate initialization failures so the parent supervisor does NOT attempt to restart.
77            // Preserve the original child name so the parent can include it in diagnostics.
78            SupervisorError::FailedToInitialize { child_name, source } => WorkerError::Initialization {
79                child_name: Some(child_name),
80                source,
81            },
82            // Preserve the structured abort count so the parent can merge it into its own shutdown tally.
83            SupervisorError::ShutdownTimedOut { aborted } => WorkerError::ShutdownTimedOut { aborted },
84            // All other supervisor errors (shutdown, no children, invalid name) are runtime-level.
85            other => WorkerError::Runtime(other.into()),
86        }
87    }
88}
89
90/// Process errors.
91#[derive(Debug, Snafu)]
92pub enum ProcessError {
93    /// The child process was aborted by the supervisor.
94    #[snafu(display("Child process was aborted by the supervisor."))]
95    Aborted,
96
97    /// The child process panicked.
98    #[snafu(display("Child process panicked."))]
99    Panicked,
100
101    /// The child process terminated with an error.
102    #[snafu(display("Child process terminated with an error: {}", source))]
103    Terminated {
104        /// The error that caused the termination.
105        source: GenericError,
106    },
107}
108
109/// Policy for automatically shutting a supervisor down based on the termination of its _significant_ children.
110///
111/// A significant child (see [`ChildBuilder::with_significant`][crate::runtime::ChildBuilder::with_significant]) is one whose termination -- when it isn't restarted -- can
112/// drive the supervisor to shut down. This mirrors Erlang/OTP's `auto_shutdown` supervisor flag, and is how an
113/// unexpected (or intentional) child exit cascades into the supervisor stopping, and thus propagating up the tree,
114/// without that child being restarted.
115#[derive(Clone, Copy, Debug, PartialEq, Eq, Default, Deserialize, Serialize)]
116#[serde(rename_all = "snake_case")]
117pub enum AutoShutdown {
118    /// Never shut down automatically; significant children have no special effect. This is the default.
119    #[default]
120    Never,
121
122    /// Shut down as soon as _any_ significant child terminates without being restarted.
123    AnySignificant,
124
125    /// Shut down once _all_ significant children have terminated without being restarted.
126    AllSignificant,
127}
128
129/// Supervisor errors.
130#[derive(Debug, Snafu)]
131#[snafu(context(suffix(false)))]
132pub enum SupervisorError {
133    /// Supervisor or worker name is invalid.
134    #[snafu(display("Invalid name for supervisor or worker: '{}'", name))]
135    InvalidName {
136        /// The name of the supervisor is invalid.
137        name: String,
138    },
139
140    /// A child process failed to initialize.
141    ///
142    /// This error indicates that a child couldn't complete its async initialization. This is distinct from runtime
143    /// failures and doesn't trigger restart logic.
144    #[snafu(display("Child process '{}' failed to initialize: {}", child_name, source))]
145    FailedToInitialize {
146        /// The name of the child that failed to initialize.
147        child_name: String,
148
149        /// The underlying initialization error.
150        source: InitializationError,
151    },
152
153    /// The supervisor exceeded its restart limits and was forced to shutdown.
154    #[snafu(display("Supervisor has exceeded restart limits and was forced to shutdown."))]
155    Shutdown,
156
157    /// The supervisor shut down because a significant child terminated.
158    ///
159    /// See [`AutoShutdown`] and [`ChildBuilder::with_significant`][crate::runtime::ChildBuilder::with_significant]. The supervisor stopped, and drained its remaining
160    /// children, because a child marked significant terminated without being restarted.
161    #[snafu(display("Supervisor shut down after a significant child terminated."))]
162    SignificantChildExited,
163
164    /// The supervisor completed a requested shutdown, but one or more workers ignored graceful shutdown and had to be
165    /// forcefully aborted after exceeding their shutdown timeout.
166    ///
167    /// The shutdown itself was requested and otherwise orderly; this variant exists so that having to forcefully stop a
168    /// worker is surfaced as a failure rather than reported as a clean shutdown. The count aggregates forced aborts
169    /// across the entire supervision tree: a parent merges in the counts reported by any child supervisors that also
170    /// timed out, so the value observed at the root supervisor is the total number of workers that had to be
171    /// force-stopped tree-wide.
172    #[snafu(display(
173        "Shutdown completed uncleanly: {} worker(s) were forcefully aborted after exceeding their shutdown timeout.",
174        aborted
175    ))]
176    ShutdownTimedOut {
177        /// The number of workers that had to be forcefully aborted.
178        aborted: usize,
179    },
180}
181
182/// A specification for a process to be added to a [`Supervisor`].
183///
184/// A child specification describes how the supervisor should create and manage a child: the underlying future that
185/// represents the process, along with metadata such as its name and shutdown strategy. All processes in a supervisor,
186/// whether a worker or a (nested) supervisor, are represented by a [`ChildSpecification`].
187///
188/// A specification is a description, not a control surface: it carries no public methods of its own. There are two
189/// ways to obtain one, matching the two levels of control:
190///
191/// - Pass a worker or supervisor directly to [`add_worker`][Supervisor::add_worker], [`spawn`][crate::runtime::spawn],
192///   or [`SupervisorHandle::spawn`], all of which accept a [`Supervisor`] or any [`Supervisable`] and convert it for
193///   you, applying the defaults appropriate to that kind of child -- including the shutdown strategy that lets a
194///   nested supervisor drain its whole subtree.
195/// - Configure one with [`ChildBuilder`][crate::runtime::ChildBuilder], which is the only way to set a restart policy,
196///   significance, placement, or a shutdown deadline. The builder exposes only the settings that make sense for the
197///   kind of child being described, and [`build`][crate::runtime::ChildBuilder::build] hands the result to
198///   [`add_worker`][Supervisor::add_worker].
199///
200/// Supervisors have no per-child settings of their own, so there is nothing to configure for a nested supervisor.
201pub struct ChildSpecification<S = WorkerSpec> {
202    spec_inner: S,
203}
204
205/// Child specification state for a worker.
206pub struct WorkerSpec {
207    worker: Arc<dyn Supervisable>,
208    options: ChildOptions,
209}
210
211/// Child specification state for a supervisor.
212pub struct SupervisorSpec {
213    supervisor: Supervisor,
214    options: ChildOptions,
215}
216
217// The configuration surface below is deliberately crate-internal: `ChildBuilder` is the public front end for all of
218// it, and is what decides which settings are offered for which kind of child. Keeping these methods off the public API
219// means a combination the builder refuses to express -- a permanent child marked significant, say -- can't be reached
220// by going around it. In-crate callers use them directly where the builder would be a layering inversion: this
221// module's own tests, which exercise the lowering these methods feed.
222impl ChildSpecification<WorkerSpec> {
223    /// Creates a specification for the given worker.
224    pub(crate) fn worker<T: Supervisable + 'static>(worker: T) -> Self {
225        Self {
226            spec_inner: WorkerSpec {
227                worker: Arc::new(worker),
228                options: ChildOptions::default(),
229            },
230        }
231    }
232
233    /// Creates a specification for a worker that can only run once.
234    ///
235    /// This function is shorthand for calling [`worker`][Self::worker] followed by
236    /// [`with_restart_type`][Self::with_restart_type] set to [`RestartType::Temporary`][RestartType::Temporary].
237    pub(crate) fn one_shot_worker<T: Supervisable + 'static>(worker: T) -> Self {
238        Self::worker(worker).with_restart_type(RestartType::Temporary)
239    }
240
241    /// Sets the restart policy for this worker.
242    ///
243    /// When left unset, the policy depends on how the child is registered: a child added up front with
244    /// [`Supervisor::add_worker`] defaults to [`RestartType::Permanent`], while one spawned dynamically with
245    /// [`SupervisorHandle::spawn`] defaults to [`RestartType::Temporary`].
246    #[must_use]
247    pub(crate) fn with_restart_type(mut self, restart_type: RestartType) -> Self {
248        self.spec_inner.options.restart = Some(restart_type);
249        self
250    }
251
252    /// Sets whether this worker is _significant_.
253    ///
254    /// A significant worker's termination (when it isn't restarted) can drive the supervisor to shut down, per the
255    /// supervisor's [`AutoShutdown`] policy. Only meaningful for non-permanent workers, since a permanent worker is
256    /// always restarted and so never terminates without being restarted.
257    #[must_use]
258    pub(crate) fn with_significant(mut self, significant: bool) -> Self {
259        self.spec_inner.options.significant = significant;
260        self
261    }
262
263    /// Runs this worker on the given Tokio runtime rather than the supervisor's own runtime.
264    ///
265    /// By default, a worker runs on whatever runtime its supervisor runs on. Use this for compute-heavy workers that
266    /// shouldn't contend with the supervisor's runtime -- for example, a topology component offloading encoding work
267    /// onto a shared worker pool.
268    ///
269    /// Note that this only affects where the worker's task is spawned. Supervision itself -- shutdown signalling,
270    /// restart handling, and abort-on-timeout -- is unchanged, and is still driven from the supervisor's runtime.
271    #[must_use]
272    pub(crate) fn with_runtime(mut self, handle: Handle) -> Self {
273        self.spec_inner.options.runtime = Some(handle);
274        self
275    }
276
277    /// Overrides the shutdown strategy for this worker.
278    ///
279    /// By default, a worker's strategy comes from [`Supervisable::shutdown_strategy`], which itself defaults to
280    /// `Graceful(5s)`. Use this when the grace period depends on where the worker is used rather than on the worker
281    /// type: a worker that a component drains during its own shutdown needs at least as long as the component itself,
282    /// otherwise it is forcefully aborted while the component is still waiting on it.
283    #[must_use]
284    pub(crate) fn with_shutdown_strategy(mut self, strategy: ShutdownStrategy) -> Self {
285        self.spec_inner.options.shutdown = ChildShutdown::Explicit(strategy);
286        self
287    }
288
289    /// Gives this worker no shutdown deadline of its own, leaving it bounded solely by its supervisor's shutdown
290    /// budget.
291    ///
292    /// Use this for a worker whose acceptable drain time is a property of the subtree it belongs to rather than of the
293    /// worker itself -- a task spawned by a topology component, for instance, where what matters is that the component
294    /// as a whole stops in time. See [`Supervisor::with_shutdown_budget`].
295    ///
296    /// A supervisor with no budget has nothing to bound the worker with, so in that case the worker falls back to the
297    /// strategy it reports through [`Supervisable::shutdown_strategy`] rather than being left to stall the drain
298    /// indefinitely.
299    #[must_use]
300    pub(crate) fn with_budget_bounded_shutdown(mut self) -> Self {
301        self.spec_inner.options.shutdown = ChildShutdown::BudgetBounded;
302        self
303    }
304}
305
306// Crate-internal for the same reason as the worker surface above: `NestedSupervisorBuilder` is the public front end,
307// and going around it would allow combinations the builder refuses to express.
308//
309// Deliberately narrower than the worker surface, though. A nested supervisor bounds its own subtree through its
310// children's deadlines: `SupervisedChild::shutdown_strategy` reports `Graceful(Duration::MAX)` for one, and
311// `WorkerState::add_worker` exempts it from the parent's budget. Offering a shutdown setting here would let a caller
312// truncate a drain the subtree is already responsible for, so there isn't one. Placement is likewise absent: a nested
313// supervisor runs wherever its parent does, and its children carry their own placement.
314impl ChildSpecification<SupervisorSpec> {
315    /// Sets the restart policy for this nested supervisor.
316    ///
317    /// When left unset, the policy depends on how the child is registered: a child added up front with
318    /// [`Supervisor::add_worker`] defaults to [`RestartType::Permanent`], while one spawned dynamically with
319    /// [`SupervisorHandle::spawn`] defaults to [`RestartType::Temporary`].
320    #[must_use]
321    pub(crate) fn with_restart_type(mut self, restart_type: RestartType) -> Self {
322        self.spec_inner.options.restart = Some(restart_type);
323        self
324    }
325
326    /// Sets whether this nested supervisor is _significant_.
327    ///
328    /// A significant child's termination (when it isn't restarted) can drive the parent supervisor to shut down, per
329    /// the parent's [`AutoShutdown`] policy.
330    #[must_use]
331    pub(crate) fn with_significant(mut self, significant: bool) -> Self {
332        self.spec_inner.options.significant = significant;
333        self
334    }
335}
336
337impl<T> From<T> for ChildSpecification<WorkerSpec>
338where
339    T: Supervisable + 'static,
340{
341    fn from(worker: T) -> Self {
342        Self::worker(worker)
343    }
344}
345
346impl From<Supervisor> for ChildSpecification<SupervisorSpec> {
347    fn from(supervisor: Supervisor) -> Self {
348        Self {
349            spec_inner: SupervisorSpec {
350                supervisor,
351                options: ChildOptions::default(),
352            },
353        }
354    }
355}
356
357mod sealed {
358    pub trait Sealed {}
359}
360
361impl sealed::Sealed for WorkerSpec {}
362impl sealed::Sealed for SupervisorSpec {}
363
364/// Child specification state.
365///
366/// This trait is sealed -- it cannot be implemented outside of this crate -- and is implemented only for
367/// [`WorkerSpec`] and [`SupervisorSpec`]. It exists so that [`Supervisor::add_worker`] and
368/// [`SupervisorHandle::spawn`] can both accept a [`ChildSpecification`] in either state (as well as bare workers and
369/// supervisors) while lowering each into the supervisor's internal representation.
370pub trait ChildState: sealed::Sealed + Sized {
371    /// Lowers a specification into the supervisor's internal representation of a child.
372    ///
373    /// `default_restart` supplies the restart policy for a specification that didn't set one, which differs by
374    /// registration path: children added up front are permanent, dynamically spawned children are temporary.
375    #[doc(hidden)]
376    fn into_child_parts(spec: ChildSpecification<Self>, default_restart: RestartType) -> LoweredChild;
377}
378
379/// A child specification lowered into the supervisor's internal representation.
380///
381/// Opaque to callers: it exists only to carry the output of [`ChildState::into_child_parts`] to the supervisor that
382/// registers the child, and is public only because [`ChildState`] is.
383pub struct LoweredChild {
384    spec: SupervisedChild,
385    config: ChildConfig,
386}
387
388impl ChildState for WorkerSpec {
389    fn into_child_parts(spec: ChildSpecification<Self>, default_restart: RestartType) -> LoweredChild {
390        let WorkerSpec { worker, options } = spec.spec_inner;
391        LoweredChild {
392            spec: SupervisedChild::Worker(worker),
393            config: options.resolve(default_restart),
394        }
395    }
396}
397
398impl ChildState for SupervisorSpec {
399    fn into_child_parts(spec: ChildSpecification<Self>, default_restart: RestartType) -> LoweredChild {
400        let SupervisorSpec { supervisor, options } = spec.spec_inner;
401        LoweredChild {
402            spec: SupervisedChild::Supervisor(supervisor),
403            config: options.resolve(default_restart),
404        }
405    }
406}
407
408/// The type-erased, runnable form of a child: either a worker or a nested supervisor.
409pub(super) enum SupervisedChild {
410    Worker(Arc<dyn Supervisable>),
411    Supervisor(Supervisor),
412}
413
414impl SupervisedChild {
415    /// Returns whether this child is a nested supervisor rather than a leaf worker.
416    pub(super) fn is_supervisor(&self) -> bool {
417        matches!(self, Self::Supervisor(_))
418    }
419
420    /// This is the link that makes the tree walkable: a parent records its child's node alongside its own, so an
421    /// observer holding the parent's node can descend into the child's.
422    pub(super) fn node(&self) -> Option<Arc<SupervisorNode>> {
423        match self {
424            Self::Worker(_) => None,
425            Self::Supervisor(supervisor) => Some(Arc::clone(&supervisor.node)),
426        }
427    }
428
429    fn process_type(&self) -> &'static str {
430        match self {
431            Self::Worker(_) => "worker",
432            Self::Supervisor(_) => "supervisor",
433        }
434    }
435
436    fn name(&self) -> &str {
437        match self {
438            Self::Worker(worker) => worker.name(),
439            Self::Supervisor(supervisor) => &supervisor.supervisor_id,
440        }
441    }
442
443    /// Returns whether this child observes the shutdown signal it is given.
444    ///
445    /// Always true for a nested supervisor: the signal is how it learns to drain its own subtree, and a supervisor on
446    /// a dedicated runtime receives it across the thread boundary through `spawn_dedicated_runtime`, where aborting
447    /// the awaiting future wouldn't stop the runtime thread anyway.
448    pub(super) fn wants_shutdown_signal(&self) -> bool {
449        match self {
450            Self::Worker(worker) => worker.wants_shutdown_signal(),
451            Self::Supervisor(_) => true,
452        }
453    }
454
455    pub(super) fn shutdown_strategy(&self) -> ShutdownStrategy {
456        match self {
457            Self::Worker(worker) => worker.shutdown_strategy(),
458
459            // Supervisors should always be given as much time as necessary shutdown down gracefully to ensure that the
460            // entire supervision subtree can be shutdown cleanly.
461            Self::Supervisor(_) => ShutdownStrategy::Graceful(Duration::MAX),
462        }
463    }
464
465    /// Creates the process for this child under `parent_process`.
466    ///
467    /// A name that sanitizes to nothing at all (an empty string, or one made up entirely of separators) can't be used
468    /// as a process name. Rather than refuse to start the child -- which for a dynamically spawned child would mean
469    /// silently losing work that the caller was told had been accepted -- the child runs under
470    /// [`UNNAMED_CHILD`] instead, and the substitution is logged.
471    pub(super) fn create_process(&self, parent_process: &Process) -> Process {
472        let name = self.name();
473        let process = match self {
474            Self::Worker(_) => Process::worker(name, parent_process),
475            Self::Supervisor(_) => Process::supervisor(name, Some(parent_process)),
476        };
477
478        process.unwrap_or_else(|| {
479            warn!(
480                parent_process = parent_process.name(),
481                child_name = name,
482                "Child process name is not usable as a process name; falling back to '{}'.",
483                UNNAMED_CHILD
484            );
485
486            match self {
487                Self::Worker(_) => Process::worker(UNNAMED_CHILD, parent_process),
488                Self::Supervisor(_) => Process::supervisor(UNNAMED_CHILD, Some(parent_process)),
489            }
490            .expect("placeholder child name is always a valid process name")
491        })
492    }
493
494    pub(super) fn create_worker_future(
495        &self, process: Process, process_shutdown: ShutdownHandle,
496    ) -> Result<WorkerFuture, SupervisorError> {
497        match self {
498            Self::Worker(worker) => {
499                let worker = Arc::clone(worker);
500                Ok(Box::pin(async move {
501                    let run_future =
502                        worker
503                            .initialize(process_shutdown)
504                            .await
505                            .map_err(|source| WorkerError::Initialization {
506                                child_name: None,
507                                source,
508                            })?;
509                    run_future.await.map_err(WorkerError::Runtime)
510                }))
511            }
512            Self::Supervisor(sup) => {
513                match sup.runtime_mode() {
514                    RuntimeMode::Ambient => {
515                        // Run on the parent's ambient runtime.
516                        Ok(sup.as_nested_process(process, process_shutdown))
517                    }
518                    RuntimeMode::Dedicated(config) => {
519                        // Spawn in a dedicated runtime on a new OS thread, passing the parent's
520                        // dataspace so the nested supervisor inherits it across the thread boundary.
521                        //
522                        // TODO: Only the dataspace is carried across, so the supervisor re-roots its own process name
523                        // when it starts (`run_with_shutdown_inner` passes no parent) rather than staying scoped
524                        // under us. That also leaves the process we build here registered as a resource group that
525                        // nothing ever enters, so it reads zero forever. Threading this process through instead would
526                        // fix both, at the cost of renaming the affected resource groups and their metric labels.
527                        let child_name = sup.supervisor_id.to_string();
528                        let dataspace = process.dataspace().clone();
529                        let handle =
530                            spawn_dedicated_runtime(sup.inner_clone(), config.clone(), process_shutdown, dataspace)
531                                .map_err(|e| SupervisorError::FailedToInitialize {
532                                    child_name,
533                                    source: e.into(),
534                                })?;
535
536                        Ok(Box::pin(async move { handle.await.map_err(WorkerError::from) }))
537                    }
538                }
539            }
540        }
541    }
542}
543
544impl Clone for SupervisedChild {
545    fn clone(&self) -> Self {
546        match self {
547            Self::Worker(worker) => Self::Worker(Arc::clone(worker)),
548            Self::Supervisor(supervisor) => Self::Supervisor(supervisor.inner_clone()),
549        }
550    }
551}
552
553/// How a child's shutdown strategy is determined.
554#[derive(Clone, Copy, Debug, Default)]
555pub(super) enum ChildShutdown {
556    /// Use whatever the worker reports through [`Supervisable::shutdown_strategy`]. This is the default.
557    #[default]
558    Worker,
559
560    /// Use this strategy, overriding whatever the worker reports.
561    Explicit(ShutdownStrategy),
562
563    /// The child carries no deadline of its own and is bounded solely by its supervisor's shutdown budget.
564    ///
565    /// A supervisor with no budget has nothing to bound the child with, so this falls back to the worker's own
566    /// strategy rather than leaving the child free to stall the drain indefinitely.
567    BudgetBounded,
568}
569
570/// Per-child settings as configured on a [`ChildSpecification`], before they are resolved for a specific registration
571/// path.
572///
573/// Separate from [`ChildConfig`] because the restart policy has no single default: a child registered up front with
574/// [`Supervisor::add_worker`] is permanent, while one spawned dynamically with [`SupervisorHandle::spawn`] is
575/// temporary. Leaving the policy unset here is what lets both paths share one specification type.
576#[derive(Clone, Debug, Default)]
577pub(super) struct ChildOptions {
578    restart: Option<RestartType>,
579    significant: bool,
580
581    /// Runtime to spawn the child on. `None` means the supervisor's own runtime.
582    runtime: Option<Handle>,
583
584    shutdown: ChildShutdown,
585}
586
587impl ChildOptions {
588    /// Resolves these options into a concrete configuration, applying `default_restart` if no policy was set.
589    fn resolve(self, default_restart: RestartType) -> ChildConfig {
590        ChildConfig {
591            restart: self.restart.unwrap_or(default_restart),
592            significant: self.significant,
593            runtime: self.runtime,
594            shutdown: self.shutdown,
595        }
596    }
597}
598
599/// Per-child configuration: its [`RestartType`], whether it is _significant_ (see [`AutoShutdown`]), the runtime it
600/// runs on, and how its shutdown strategy is decided.
601#[derive(Clone, Debug)]
602pub(super) struct ChildConfig {
603    restart: RestartType,
604    significant: bool,
605    runtime: Option<Handle>,
606    shutdown: ChildShutdown,
607}
608
609impl ChildConfig {
610    /// Returns the runtime the child should be spawned on, if it isn't the supervisor's own.
611    pub(super) fn runtime(&self) -> Option<&Handle> {
612        self.runtime.as_ref()
613    }
614
615    /// Returns how the child's shutdown strategy should be determined.
616    pub(super) fn shutdown(&self) -> ChildShutdown {
617        self.shutdown
618    }
619}
620
621/// A registered child: its specification together with the configuration chosen at registration time.
622#[derive(Clone)]
623struct ChildEntry {
624    spec: SupervisedChild,
625    config: ChildConfig,
626    /// Whether this child was added dynamically (via [`SupervisorHandle`]) rather than statically before the run. Used
627    /// to maintain the dynamic-children gauge.
628    dynamic: bool,
629}
630
631/// Identifier for a child managed by a [`Supervisor`].
632///
633/// Returned by [`SupervisorHandle::spawn`] for dynamically spawned children. Unique within a single process for the
634/// lifetime of a supervisor run.
635#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
636pub struct ChildId(u64);
637
638impl ChildId {
639    /// Returns the raw numeric value of this identifier.
640    pub const fn as_u64(self) -> u64 {
641        self.0
642    }
643}
644
645/// A dynamic spawn request handed from a [`SupervisorHandle`] to the running supervisor.
646struct PendingSpawn {
647    id: u64,
648    spec: SupervisedChild,
649    config: ChildConfig,
650}
651
652/// Number of queued spawn requests the supervisor takes in one go before returning to its loop.
653///
654/// Draining in batches keeps a burst of spawns to a single wake-up rather than one per child, while still bounding how
655/// long the supervisor can spend registering children before it re-checks the rest of its loop (most importantly,
656/// shutdown).
657const SPAWN_DRAIN_BATCH: usize = 64;
658
659/// A handle for spawning dynamic children on a running [`Supervisor`].
660///
661/// Obtained from [`Supervisor::handle`]. Handles are cheap to clone and can be shared across tasks.
662///
663/// Spawning is synchronous and infallible, in the spirit of [`tokio::spawn`]: the child is queued for the running
664/// supervisor and the call returns immediately with the child's [`ChildId`]. Also as with [`tokio::spawn`], being
665/// accepted is not a promise of being run -- if the supervisor isn't running, or shuts down before it gets to the
666/// queued child, the child is never started at all.
667///
668/// # Ambient spawning
669///
670/// Code running under supervision usually doesn't need a handle at all: [`spawn`][crate::runtime::spawn] targets the
671/// supervisor of whatever process is currently running. Use a handle when spawning from outside supervision, or when
672/// targeting a supervisor other than the ambient one. [`scope`][Self::scope] bridges the two by making a handle the
673/// ambient supervisor for a future.
674#[derive(Clone)]
675pub struct SupervisorHandle {
676    name: Arc<str>,
677    // The currently running supervisor publishes its spawn queue here so handles can reach the live run; it's cleared
678    // when no run is active, at which point spawns are accepted and dropped.
679    current_tx: Arc<Mutex<Option<mpsc::UnboundedSender<PendingSpawn>>>>,
680    id_counter: Arc<AtomicU64>,
681    active: Arc<AtomicUsize>,
682}
683
684impl SupervisorHandle {
685    /// Returns the name of the supervisor this handle refers to.
686    pub fn name(&self) -> &str {
687        &self.name
688    }
689
690    /// Spawns a new dynamic child.
691    ///
692    /// Accepts anything [`Supervisor::add_worker`] accepts: a bare [`Supervisable`], a [`Supervisor`] to run as a
693    /// nested supervision subtree, or a [`ChildSpecification`] configured in detail.
694    ///
695    /// Unless [`ChildBuilder`][crate::runtime::ChildBuilder] says otherwise, dynamic children are
696    /// [`temporary`][RestartType::Temporary]: they
697    /// aren't restarted when they die, and they aren't restored when the supervisor itself restarts. That suits
698    /// short-lived, non-critical work that still wants structured concurrency -- the child is stopped when the
699    /// supervisor is restarted or terminated.
700    ///
701    /// The returned [`ChildId`] identifies the child for the lifetime of the supervisor run. The child is queued
702    /// rather than started synchronously, so it may not have begun running by the time this returns; if the supervisor
703    /// isn't running, or shuts down before reaching the child, it never runs at all.
704    pub fn spawn<S, T>(&self, child: T) -> ChildId
705    where
706        S: ChildState,
707        T: Into<ChildSpecification<S>>,
708    {
709        let LoweredChild { spec, config } = S::into_child_parts(child.into(), RestartType::Temporary);
710
711        // Take the id before we try to enqueue: the caller gets a stable identifier either way, and ids are only
712        // meaningful within a run.
713        let id = self.id_counter.fetch_add(1, Ordering::Relaxed);
714        let pending = PendingSpawn { id, spec, config };
715
716        // Clone the sender out from under the lock rather than sending while holding the guard.
717        //
718        // The queue behind it is unbounded on purpose: a queued child is a child that will be started, so the only
719        // thing a depth limit could buy is discarding work the caller was told had been accepted. A backlog only forms
720        // while the supervisor can't drain -- mid-restart, or mid-drain -- and holding it until it can is the whole
721        // point.
722        let tx = self.current_tx.lock().unwrap().clone();
723        match tx {
724            // Racing a teardown is normal rather than exceptional -- a source that spawns a child per connection will
725            // do it every time it is shut down mid-accept -- so this stays at debug level.
726            Some(tx) => {
727                if let Err(e) = tx.send(pending) {
728                    debug!(
729                        supervisor_id = %self.name,
730                        child_name = e.0.spec.name(),
731                        "Supervisor is shutting down; dynamic child will not be started."
732                    );
733                }
734            }
735            // Spawning against a supervisor that never ran, on the other hand, is a wiring mistake: nothing about the
736            // program's normal operation produces it, and the child is silently lost.
737            None => warn!(
738                supervisor_id = %self.name,
739                child_name = pending.spec.name(),
740                "Supervisor is not running; dynamic child will not be started."
741            ),
742        }
743
744        ChildId(id)
745    }
746
747    /// Returns whether the supervisor is currently running.
748    pub fn is_running(&self) -> bool {
749        self.current_tx.lock().unwrap().is_some()
750    }
751
752    /// Returns the number of dynamic children currently running under the supervisor.
753    ///
754    /// Counts children the supervisor has actually started, so a child that has been spawned but not yet picked up
755    /// isn't included yet.
756    pub fn active_children(&self) -> usize {
757        self.active.load(Ordering::Relaxed)
758    }
759}
760
761/// Supervises a set of workers.
762///
763/// # Workers
764///
765/// All workers are defined through implementation of the [`Supervisable`] trait, which provides the logic for both
766/// creating the underlying worker future that's spawned, as well as other metadata, such as the worker's name, how the
767/// worker should be shutdown, and so on.
768///
769/// Supervisors also (indirectly) implement the [`Supervisable`] trait, allowing them to be supervised by other
770/// supervisors in order to construct _supervision trees_.
771///
772/// # Instrumentation
773///
774/// Supervisors automatically create their own allocation group
775/// ([`TrackingAllocator`][saluki_common::resource_tracking::TrackingAllocator]), which is used to track both the memory
776/// usage of the supervisor itself and its children. Additionally, individual worker processes are wrapped in a
777/// dedicated [`tracing::Span`] to allow tracing the causal relationship between arbitrary code and the worker executing
778/// it, and statistics about task polls (poll count, poll duration) are collected.
779///
780/// # Restart Strategies
781///
782/// As the main purpose of a supervisor, restart behavior is fully configurable. A number of restart strategies are
783/// available, which generally relate to the purpose of the supervisor: whether the workers being managed are
784/// independent or interdependent.
785///
786/// All restart strategies are configured through [`RestartStrategy`], which has more information on the available
787/// strategies and configuration settings.
788pub struct Supervisor {
789    supervisor_id: Arc<str>,
790    child_specs: Vec<ChildEntry>,
791    runtime_mode: RuntimeMode,
792    // Shared across clones (a nested supervisor is cloned each time it runs) and across all handles. While a run is
793    // active it holds that run's spawn queue so handles can reach the live supervisor; it's `None` whenever no run is
794    // active, at which point spawned children are dropped rather than queued. Doubles as the `is_running` signal.
795    current_tx: Arc<Mutex<Option<mpsc::UnboundedSender<PendingSpawn>>>>,
796    id_counter: Arc<AtomicU64>,
797    // Number of dynamic children currently running, shared with handles so it can be surfaced as a gauge.
798    active: Arc<AtomicUsize>,
799    // Observable bookkeeping for this supervisor, shared with clones (so every generation writes to the same place),
800    // with this supervisor's parent (so the tree can be walked downward), and with any tree handle.
801    node: Arc<SupervisorNode>,
802}
803
804impl Supervisor {
805    /// Creates an empty `Supervisor` with the default restart strategy.
806    pub fn new<S: AsRef<str>>(supervisor_id: S) -> Result<Self, SupervisorError> {
807        // We try to throw an error about invalid names as early as possible. This is a manual check, so we might still
808        // encounter an error later when actually running the supervisor, but this is a good first step to catch the
809        // bulk of invalid names.
810        if supervisor_id.as_ref().is_empty() {
811            return Err(SupervisorError::InvalidName {
812                name: supervisor_id.as_ref().to_string(),
813            });
814        }
815
816        let supervisor_id: Arc<str> = supervisor_id.as_ref().into();
817
818        Ok(Self {
819            node: Arc::new(SupervisorNode::new(Arc::clone(&supervisor_id))),
820            supervisor_id,
821            child_specs: Vec::new(),
822            runtime_mode: RuntimeMode::default(),
823            current_tx: Arc::new(Mutex::new(None)),
824            id_counter: Arc::new(AtomicU64::new(0)),
825            active: Arc::new(AtomicUsize::new(0)),
826        })
827    }
828
829    /// Returns the supervisor's ID.
830    pub fn id(&self) -> &str {
831        &self.supervisor_id
832    }
833
834    /// Sets the restart strategy for the supervisor.
835    pub fn with_restart_strategy(self, strategy: RestartStrategy) -> Self {
836        self.node.update_config(|config| config.restart_strategy = strategy);
837        self
838    }
839
840    /// Sets the supervisor's automatic-shutdown policy.
841    ///
842    /// Controls whether the termination of _significant_ children (see [`ChildBuilder::with_significant`][crate::runtime::ChildBuilder::with_significant]) drives the
843    /// supervisor to shut down. Defaults to [`AutoShutdown::Never`].
844    pub fn with_auto_shutdown(self, auto_shutdown: AutoShutdown) -> Self {
845        self.node.update_config(|config| config.auto_shutdown = auto_shutdown);
846        self
847    }
848
849    /// Bounds how long this supervisor waits for its worker children during shutdown.
850    ///
851    /// Without a budget, a supervisor waits as long as each child's own [`ShutdownStrategy`] allows, and waits
852    /// indefinitely for any child that has no finite deadline of its own. A budget makes the supervisor responsible for
853    /// the deadline instead: children need no individual timeouts, and whatever is still running when the budget
854    /// elapses is forcefully aborted -- each one named in the logs, and counted in the resulting
855    /// [`SupervisorError::ShutdownTimedOut`].
856    ///
857    /// The budget is a ceiling, not a replacement: a child that also carries its own finite deadline is still held to
858    /// whichever elapses first.
859    ///
860    /// Since children are always drained concurrently, the budget bounds the drain as a whole rather than accruing
861    /// per child: it is measured from the moment shutdown begins, and every child is held to it simultaneously.
862    ///
863    /// Two kinds of child are outside it. A nested supervisor is never cut off by its parent's budget -- it bounds its
864    /// own subtree, and aborting it would both truncate that drain and, for a supervisor running on a dedicated
865    /// runtime, fail to stop it at all. A [`ShutdownStrategy::Brutal`] child is aborted up front and never waited on.
866    /// Neither can a budget bound work that ignores cancellation, since an abort only takes effect at an await point.
867    ///
868    /// Use this where one deadline for a whole subtree is more meaningful than a guess per worker -- a topology
869    /// component and its background tasks, for instance, where what matters is that the component as a whole stops in
870    /// time.
871    #[must_use]
872    pub fn with_shutdown_budget(self, budget: Duration) -> Self {
873        self.node.update_config(|config| config.shutdown_budget = Some(budget));
874        self
875    }
876
877    /// Returns a handle for spawning dynamic children on this supervisor while it runs.
878    ///
879    /// The handle can be created before the supervisor starts and cloned freely. Spawning through it always succeeds,
880    /// but a child is only ever started while the supervisor is actually running: one spawned before the supervisor
881    /// starts, or after it has shut down, is accepted and then dropped.
882    pub fn handle(&self) -> SupervisorHandle {
883        SupervisorHandle {
884            name: Arc::clone(&self.supervisor_id),
885            current_tx: Arc::clone(&self.current_tx),
886            id_counter: Arc::clone(&self.id_counter),
887            active: Arc::clone(&self.active),
888        }
889    }
890
891    /// Returns a read-only handle for taking snapshots of this supervisor and the subtree beneath it.
892    ///
893    /// The handle can be created before the supervisor starts and remains valid across every restart of it. Unlike
894    /// [`handle`][Self::handle], it grants no ability to affect the supervisor -- only to observe it.
895    pub fn tree_handle(&self) -> SupervisionTreeHandle {
896        SupervisionTreeHandle::new(Arc::clone(&self.node))
897    }
898
899    /// Configures this supervisor to run in a dedicated runtime.
900    ///
901    /// When this supervisor is added as a child to another supervisor, it will spawn its own OS threads and Tokio
902    /// runtime instead of running on the parent's ambient runtime.
903    ///
904    /// This provides runtime isolation, which can be useful for:
905    /// - CPU-bound work that shouldn't block the parent's runtime
906    /// - Isolating failures in one part of the system
907    /// - Using different runtime configurations (for example, single-threaded vs multi-threaded)
908    pub fn with_dedicated_runtime(mut self, config: RuntimeConfiguration) -> Self {
909        let worker_threads = config.worker_threads();
910        self.node
911            .update_config(|node_config| node_config.dedicated_threads = Some(worker_threads));
912        self.runtime_mode = RuntimeMode::Dedicated(config);
913        self
914    }
915
916    /// Returns the runtime mode for this supervisor.
917    pub(crate) fn runtime_mode(&self) -> &RuntimeMode {
918        &self.runtime_mode
919    }
920
921    /// Adds a worker (or nested supervisor) to the supervisor.
922    ///
923    /// A worker can be anything that implements the [`Supervisable`] trait. A [`Supervisor`] can also be added as a
924    /// worker and managed in a nested fashion, known as a supervision tree.
925    ///
926    /// Anything that needs configuring -- a restart policy, significance, placement, a shutdown deadline -- is
927    /// described with [`ChildBuilder`][crate::runtime::ChildBuilder] and handed over via
928    /// [`build`][crate::runtime::ChildBuilder::build]. See [`ChildSpecification`] for how children are represented
929    /// internally.
930    pub fn add_worker<S, T>(&mut self, child: T)
931    where
932        S: ChildState,
933        T: Into<ChildSpecification<S>>,
934    {
935        let LoweredChild { spec, config } = S::into_child_parts(child.into(), RestartType::Permanent);
936        self.push_child(ChildEntry {
937            spec,
938            config,
939            dynamic: false,
940        });
941    }
942
943    /// Warns when a child was marked significant but nothing will act on it.
944    ///
945    /// Significance only has an effect for a child that can terminate without being restarted, under a supervisor
946    /// whose [`AutoShutdown`] policy isn't [`Never`][AutoShutdown::Never]. Either mismatch makes the flag inert, which
947    /// is worth saying out loud: a caller who marked a child significant is asserting that its termination matters,
948    /// and silently ignoring that is how a supervisor ends up outliving something it can't work without.
949    ///
950    /// Called as children are started rather than as they are registered, because the policy half of the question
951    /// isn't answerable any earlier: [`with_auto_shutdown`][Self::with_auto_shutdown] consumes the supervisor while
952    /// [`add_worker`][Self::add_worker] borrows it, so a caller is free to add children first and set the policy
953    /// afterwards. Checking at registration time would flag that -- entirely correct -- ordering as a mistake.
954    ///
955    /// Warn-only: the child still starts, since an inert flag is useless rather than unsafe.
956    fn warn_if_significance_is_inert(&self, config: &ChildConfig, child_name: &str, auto_shutdown: AutoShutdown) {
957        if !config.significant {
958            return;
959        }
960
961        if config.restart == RestartType::Permanent {
962            warn!(
963                supervisor_id = %self.supervisor_id,
964                child_name,
965                "Child is marked significant but is permanent, so it is always restarted and the flag has no effect."
966            );
967        }
968
969        if auto_shutdown == AutoShutdown::Never {
970            warn!(
971                supervisor_id = %self.supervisor_id,
972                child_name,
973                "Child is marked significant but the supervisor's auto-shutdown policy is `Never`, so the flag has \
974                 no effect."
975            );
976        }
977    }
978
979    fn push_child(&mut self, entry: ChildEntry) {
980        debug!(
981            supervisor_id = %self.supervisor_id,
982            "Adding new static child process #{}. ({}, {}, {:?})",
983            self.child_specs.len(),
984            entry.spec.process_type(),
985            entry.spec.name(),
986            entry.config,
987        );
988
989        // The policy half of the inert-significance check has to wait until the supervisor runs (see
990        // `warn_if_significance_is_inert`), but this half doesn't depend on anything but the child itself, and here we
991        // are still in the caller's frame. `ChildBuilder` makes the combination unreachable from outside the crate, so
992        // this guards in-crate construction.
993        debug_assert!(
994            !(entry.config.significant && entry.config.restart == RestartType::Permanent),
995            "child '{}' was marked significant but is permanent, so it is always restarted and its termination can \
996             never drive auto-shutdown",
997            entry.spec.name()
998        );
999
1000        self.child_specs.push(entry);
1001    }
1002
1003    /// Assembled here rather than inside the roster because only the supervisor can see a child's specification and
1004    /// resolved configuration.
1005    fn child_facts(entry: &ChildEntry, key: ChildKey) -> ChildFacts {
1006        ChildFacts {
1007            key,
1008            name: entry.spec.name().into(),
1009            node: entry.spec.node(),
1010            restart: entry.config.restart,
1011            significant: entry.config.significant,
1012        }
1013    }
1014
1015    fn spawn_static_children(
1016        &self, roster: &mut Roster<ChildEntry>, worker_state: &mut WorkerState, auto_shutdown: AutoShutdown,
1017    ) -> Result<(), SupervisorError> {
1018        debug!(supervisor_id = %self.supervisor_id, "Spawning all static child processes.");
1019        for (index, entry) in self.child_specs.iter().enumerate() {
1020            self.warn_if_significance_is_inert(&entry.config, entry.spec.name(), auto_shutdown);
1021
1022            let id = self.id_counter.fetch_add(1, Ordering::Relaxed);
1023            let started = worker_state.add_worker(id, &entry.spec, &entry.config)?;
1024            let facts = Self::child_facts(entry, ChildKey::Static(index));
1025            roster.insert(id, entry.clone(), facts, started);
1026        }
1027
1028        Ok(())
1029    }
1030
1031    /// Respawns children after a one-for-all restart, honoring each child's [`RestartType`].
1032    ///
1033    /// Every child except [`RestartType::Temporary`] is restarted, matching Erlang/OTP: a group restart restarts all
1034    /// permanent and transient children -- regardless of how they last exited, including a transient child that had
1035    /// already exited cleanly -- but never temporary children, which are shut down with the group and not brought back.
1036    /// A transient child's "restart only on abnormal exit" rule governs its _own_ termination, not a group restart
1037    /// driven by a sibling. Dynamic children are not restored (they are lost on a supervisor-level restart).
1038    fn respawn_children_one_for_all(
1039        &self, roster: &mut Roster<ChildEntry>, worker_state: &mut WorkerState,
1040    ) -> Result<(), SupervisorError> {
1041        debug!(supervisor_id = %self.supervisor_id, "Restarting all eligible static child processes.");
1042        for (index, entry) in self.child_specs.iter().enumerate() {
1043            // Temporary children are never restarted by a group restart (matching OTP): they are shut down with the
1044            // group but not brought back.
1045            if entry.config.restart == RestartType::Temporary {
1046                continue;
1047            }
1048            let id = self.id_counter.fetch_add(1, Ordering::Relaxed);
1049            let started = worker_state.add_worker(id, &entry.spec, &entry.config)?;
1050            // Keyed by position in the static child list rather than by roster id: a group restart hands out fresh
1051            // ids, and the child's creation time and restart count have to survive that.
1052            let facts = Self::child_facts(entry, ChildKey::Static(index));
1053            roster.insert(id, entry.clone(), facts, started);
1054        }
1055
1056        Ok(())
1057    }
1058
1059    /// Registers one dynamic child into the running supervisor's worker set and roster.
1060    fn spawn_dynamic_child(
1061        &self, spawn: PendingSpawn, worker_state: &mut WorkerState, roster: &mut Roster<ChildEntry>,
1062        significant_remaining: &mut usize, auto_shutdown: AutoShutdown,
1063    ) {
1064        let PendingSpawn { id, spec, config } = spawn;
1065        let entry = ChildEntry {
1066            spec,
1067            config,
1068            dynamic: true,
1069        };
1070        self.warn_if_significance_is_inert(&entry.config, entry.spec.name(), auto_shutdown);
1071
1072        match worker_state.add_worker(id, &entry.spec, &entry.config) {
1073            Ok(started) => {
1074                if entry.config.significant {
1075                    *significant_remaining += 1;
1076                }
1077                self.active.fetch_add(1, Ordering::Relaxed);
1078                let facts = Self::child_facts(&entry, ChildKey::Dynamic(id));
1079                roster.insert(id, entry, facts, started);
1080            }
1081            Err(e) => {
1082                // The only way registration fails now that child names always resolve is a nested supervisor on a
1083                // dedicated runtime failing to get an OS thread. There's no caller left to report it to -- spawning is
1084                // infallible -- so the child is dropped and the failure is logged here.
1085                error!(
1086                    supervisor_id = %self.supervisor_id,
1087                    child_name = entry.spec.name(),
1088                    error = %e,
1089                    "Failed to start dynamic child."
1090                );
1091            }
1092        }
1093    }
1094
1095    async fn run_inner(&self, process: Process, process_shutdown: ShutdownHandle) -> Result<(), SupervisorError> {
1096        // Publish a fresh spawn queue for this run so handles can spawn dynamic children into it; while it's set,
1097        // handles observe us as running.
1098        let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
1099        *self.current_tx.lock().unwrap() = Some(cmd_tx);
1100
1101        // Record the process we're actually running under. It has to come from here rather than from whatever our
1102        // parent created for us, because a supervisor on a dedicated runtime re-roots its process name when it
1103        // starts, and only the resulting name matches the resource group its allocations land in.
1104        self.node.begin_run(&process);
1105
1106        let result = self.supervise(process, process_shutdown, cmd_rx).await;
1107
1108        // The run is over. Clear the sender so later spawns are dropped rather than queued, and reset the
1109        // dynamic-children gauge. Dropping the receiver (owned by `supervise`) already discarded anything in flight.
1110        *self.current_tx.lock().unwrap() = None;
1111        self.active.store(0, Ordering::Relaxed);
1112
1113        // Report as stopped rather than presenting the children of a generation that has ended. Every way out of
1114        // `supervise` -- clean shutdown, a child failing to initialize, the restart limit being exceeded, a
1115        // significant child exiting -- returns through here, so this covers all of them.
1116        self.node.end_run();
1117
1118        result
1119    }
1120
1121    async fn supervise(
1122        &self, process: Process, process_shutdown: ShutdownHandle, mut cmd_rx: mpsc::UnboundedReceiver<PendingSpawn>,
1123    ) -> Result<(), SupervisorError> {
1124        // Read once: configuration can't change while a run is in flight, and this keeps the reap and spawn arms
1125        // below off the node's lock entirely.
1126        let NodeConfig {
1127            restart_strategy,
1128            auto_shutdown,
1129            shutdown_budget,
1130            ..
1131        } = self.node.config();
1132
1133        let mut restart_state = RestartState::new(restart_strategy);
1134        let mut worker_state = WorkerState::new(process, self.handle(), shutdown_budget, Arc::clone(&self.node));
1135
1136        // The live roster of children -- both static (seeded below) and dynamic (added via the handle) -- keyed by a
1137        // stable id. A restart re-runs a child by id; a child that isn't restarted is removed from the roster.
1138        //
1139        // Every change also lands in this supervisor's supervision-tree bookkeeping, which is what makes the tree
1140        // observable. Going through the roster for all of it is what keeps the two from drifting apart.
1141        let mut roster = Roster::new(Arc::clone(&self.node));
1142
1143        // Spawn the static children. Initialization is folded into each worker's task, so this returns immediately --
1144        // children initialize concurrently in the background.
1145        self.spawn_static_children(&mut roster, &mut worker_state, auto_shutdown)?;
1146
1147        // Track how many significant children are still running, for `AutoShutdown` evaluation.
1148        let mut significant_remaining = roster.values().filter(|entry| entry.config.significant).count();
1149
1150        // Scratch space reused across every batched drain of the spawn queue.
1151        let mut spawn_batch = Vec::with_capacity(SPAWN_DRAIN_BATCH);
1152
1153        // Now we supervise.
1154        pin!(process_shutdown);
1155
1156        let outcome = loop {
1157            select! {
1158                // Shutdown first, then reaping, then taking on new work -- so neither a flood of spawns nor a stream
1159                // of exiting children can starve anything ranked above it.
1160                biased;
1161
1162                // Shutdown has been triggered; break out of the loop with a clean outcome and tear down below. (We
1163                // can't touch `cmd_rx` in any arm's handler -- the `recv_many` arm below borrows it for the whole
1164                // `select!` -- so all teardown happens after the loop.)
1165                _ = &mut process_shutdown => break Ok(()),
1166
1167                // Reaping outranks taking on new work: it is the only place children leave the join set, the roster
1168                // and the `active` gauge, so a steady stream of spawns must not be able to starve it. It parks
1169                // whenever there are no children, so it can't starve spawning in return.
1170                (child_id, worker_result) = worker_state.wait_for_next_worker() => {
1171                    // Pull out what we need from the roster before we mutate it.
1172                    let (child_name, config, dynamic) = {
1173                        let entry = roster.get(child_id).expect("completed worker must be present in the roster");
1174                        (entry.spec.name().to_string(), entry.config.clone(), entry.dynamic)
1175                    };
1176
1177                    // Initialization failures are not eligible for restart -- they propagate immediately.
1178                    if let Err(WorkerError::Initialization { child_name: inner, source }) = worker_result {
1179                        // If the error came from a nested supervisor, include the original child name to make the error
1180                        // chain more informative (e.g., "ctrl-pln/privileged-api").
1181                        let full_name = match inner {
1182                            Some(inner) => format!("{}/{}", child_name, inner),
1183                            None => child_name.clone(),
1184                        };
1185
1186                        error!(supervisor_id = %self.supervisor_id, worker_name = full_name, "Child process failed to initialize: {}", source);
1187                        break Err(SupervisorError::FailedToInitialize { child_name: full_name, source });
1188                    }
1189
1190                    // A worker exited abnormally if it returned an error, panicked, or was aborted; a clean exit is
1191                    // `Ok(())`. Together with the worker's restart policy, this determines whether we restart it.
1192                    let abnormal = worker_result.is_err();
1193                    let worker_result = worker_result.map_err(|e| match e {
1194                        WorkerError::Runtime(e) => ProcessError::Terminated { source: e },
1195                        WorkerError::Initialization { .. } => unreachable!("handled above"),
1196                        // A nested supervisor only reports `ShutdownTimedOut` while draining, which is driven by its own
1197                        // `process_shutdown` -- and that fires only when *this* supervisor is itself draining it, i.e.
1198                        // from `shutdown_workers` below, never from this main-loop arm. Treat it as a runtime
1199                        // termination defensively rather than asserting unreachable.
1200                        WorkerError::ShutdownTimedOut { aborted } => ProcessError::Terminated {
1201                            source: SupervisorError::ShutdownTimedOut { aborted }.into(),
1202                        },
1203                    });
1204
1205                    if !config.restart.should_restart(abnormal) {
1206                        // Not eligible for restart given how it exited. Drop it from the roster, and free its slot/gauge
1207                        // if it was dynamic. Crucially, we do NOT consult `evaluate_restart` here: non-restarts must not
1208                        // consume the restart-intensity budget, otherwise a steady stream of terminating temporary
1209                        // children would eventually trip the limit and tear the supervisor (and its siblings) down.
1210                        //
1211                        // An abnormal exit is reported at `warn` rather than `debug`: a child that isn't restarted --
1212                        // every dynamically-spawned child, in practice -- has no other path back to its owner, so this
1213                        // is the only place its failure is surfaced.
1214                        if abnormal {
1215                            warn!(supervisor_id = %self.supervisor_id, worker_name = %child_name, restart = ?config.restart, ?worker_result, "Child process exited with an error and is not eligible for restart.");
1216                        } else {
1217                            debug!(supervisor_id = %self.supervisor_id, worker_name = %child_name, restart = ?config.restart, "Child process exited and is not eligible for restart.");
1218                        }
1219                        roster.remove(child_id);
1220                        if dynamic {
1221                            self.active.fetch_sub(1, Ordering::Relaxed);
1222                        }
1223
1224                        // A significant child terminating without restart can drive the supervisor to shut down, per its
1225                        // `AutoShutdown` policy -- cascading an unexpected (or intentional) child exit into the
1226                        // supervisor stopping and propagating up the tree.
1227                        if config.significant {
1228                            significant_remaining = significant_remaining.saturating_sub(1);
1229                            let auto_shutdown = match auto_shutdown {
1230                                AutoShutdown::Never => false,
1231                                AutoShutdown::AnySignificant => true,
1232                                AutoShutdown::AllSignificant => significant_remaining == 0,
1233                            };
1234                            if auto_shutdown {
1235                                warn!(supervisor_id = %self.supervisor_id, worker_name = %child_name, ?worker_result, "Significant child terminated; shutting down supervisor.");
1236                                break Err(SupervisorError::SignificantChildExited);
1237                            }
1238                        }
1239                    } else {
1240                        match restart_state.evaluate_restart() {
1241                            RestartAction::Restart(mode) => match mode {
1242                                RestartMode::OneForOne => {
1243                                    warn!(supervisor_id = %self.supervisor_id, worker_name = %child_name, ?worker_result, "Child process terminated, restarting.");
1244                                    let spec = roster.get(child_id).expect("present for restart").spec.clone();
1245                                    match worker_state.add_worker(child_id, &spec, &config) {
1246                                        Ok(started) => roster.restart_in_place(child_id, started),
1247                                        Err(e) => break Err(e),
1248                                    }
1249                                }
1250                                RestartMode::OneForAll => {
1251                                    warn!(supervisor_id = %self.supervisor_id, worker_name = %child_name, ?worker_result, "Child process terminated, restarting all processes.");
1252                                    // This drain is part of a restart, not a shutdown: any forced aborts here are
1253                                    // already logged per-worker, and the supervisor keeps running, so the count does
1254                                    // not feed the unclean-shutdown signal.
1255                                    let _ = worker_state.shutdown_workers().await;
1256                                    // A one-for-all restart resets to the static roster; dynamic children are not
1257                                    // restored (they're lost on a supervisor-level restart, matching Erlang/OTP), and
1258                                    // temporary children are not restarted.
1259                                    roster.clear_for_group_restart();
1260                                    self.active.store(0, Ordering::Relaxed);
1261                                    let respawn = self.respawn_children_one_for_all(&mut roster, &mut worker_state);
1262                                    if let Err(e) = respawn {
1263                                        break Err(e);
1264                                    }
1265                                    significant_remaining =
1266                                        roster.values().filter(|entry| entry.config.significant).count();
1267                                }
1268                            },
1269                            RestartAction::Shutdown => {
1270                                error!(supervisor_id = %self.supervisor_id, worker_name = %child_name, ?worker_result, "Supervisor shutting down due to restart limits.");
1271                                break Err(SupervisorError::Shutdown);
1272                            }
1273                        }
1274                    }
1275                }
1276
1277                // A handle asked us to spawn one or more dynamic children. The published sender keeps the queue open
1278                // for the whole run, so this only yields zero once we close it during teardown. Draining in batches
1279                // keeps a burst of spawns to a single wake-up.
1280                _ = cmd_rx.recv_many(&mut spawn_batch, SPAWN_DRAIN_BATCH) => {
1281                    for spawn in spawn_batch.drain(..) {
1282                        self.spawn_dynamic_child(
1283                            spawn,
1284                            &mut worker_state,
1285                            &mut roster,
1286                            &mut significant_remaining,
1287                            auto_shutdown,
1288                        );
1289                    }
1290                }
1291            }
1292        };
1293
1294        // The run is ending -- either cleanly (shutdown was signalled) or with an error (a child failed to initialize
1295        // or restart, the restart limit was exceeded, or a significant child exited). On every path: stop accepting
1296        // spawns and discard anything still queued, rather than starting children only to tear them down immediately,
1297        // and then shut down all children.
1298        cmd_rx.close();
1299        let mut discarded = 0;
1300        while cmd_rx.try_recv().is_ok() {
1301            discarded += 1;
1302        }
1303        if discarded > 0 {
1304            debug!(
1305                supervisor_id = %self.supervisor_id,
1306                discarded,
1307                "Discarded queued dynamic children during shutdown."
1308            );
1309        }
1310        let aborted = worker_state.shutdown_workers().await;
1311
1312        // A requested shutdown that nonetheless had to forcefully abort one or more workers (here or anywhere in the
1313        // subtree below us) is surfaced as an unclean shutdown so it propagates up the tree rather than being reported
1314        // as success. An outcome that already carries an error (initialization, restart limit, significant child)
1315        // takes precedence -- that's the root cause -- and the forced aborts are left to the per-worker warnings.
1316        match outcome {
1317            Ok(()) if aborted > 0 => {
1318                warn!(supervisor_id = %self.supervisor_id, aborted, "Shutdown completed uncleanly; workers were forcefully aborted.");
1319                Err(SupervisorError::ShutdownTimedOut { aborted })
1320            }
1321            outcome => outcome,
1322        }
1323    }
1324
1325    fn as_nested_process(&self, process: Process, process_shutdown: ShutdownHandle) -> WorkerFuture {
1326        // Simple wrapper around `run_inner` to satisfy the return type signature needed when running the supervisor as
1327        // a nested child process in another supervisor.
1328        debug!(supervisor_id = %self.supervisor_id, "Nested supervisor starting.");
1329
1330        // Create a standalone clone of ourselves so we can fulfill the future signature.
1331        let sup = self.inner_clone();
1332
1333        Box::pin(async move {
1334            sup.run_inner(process, process_shutdown)
1335                .await
1336                .map_err(WorkerError::from)
1337        })
1338    }
1339
1340    /// Runs the supervisor forever.
1341    ///
1342    /// # Errors
1343    ///
1344    /// If the supervisor exceeds its restart limits, or fails to initialize a child process, an error is returned.
1345    pub async fn run(&mut self) -> Result<(), SupervisorError> {
1346        // Create a no-op `ShutdownHandle` to satisfy the `run_inner` function. This is never used since we want to run
1347        // forever, but we need to satisfy the signature.
1348        let process_shutdown = ShutdownHandle::noop();
1349        let process = Process::supervisor(&self.supervisor_id, None).context(InvalidName {
1350            name: self.supervisor_id.to_string(),
1351        })?;
1352
1353        debug!(supervisor_id = %self.supervisor_id, "Supervisor starting.");
1354        self.run_inner(process.clone(), process_shutdown)
1355            .into_process_future(process)
1356            .await
1357    }
1358
1359    /// Runs the supervisor until shutdown is triggered.
1360    ///
1361    /// When `shutdown` resolves, the supervisor will shutdown all child processes according to their shutdown strategy,
1362    /// and then return.
1363    ///
1364    /// # Errors
1365    ///
1366    /// If the supervisor exceeds its restart limits, or fails to initialize a child process, an error is returned.
1367    pub async fn run_with_shutdown<F: Future + Send + 'static>(&mut self, shutdown: F) -> Result<(), SupervisorError> {
1368        // Drive the caller-provided shutdown future into a trigger so the supervisor can begin shutting down its
1369        // children once `shutdown` resolves. The trigger fires at most once (guarded), and otherwise fires on drop if
1370        // the supervisor returns on its own first.
1371        let (shutdown_coordinator, shutdown_handle) = ShutdownHandle::paired();
1372        let run = self.run_with_shutdown_inner(shutdown_handle, None);
1373        pin!(run, shutdown);
1374
1375        let mut shutdown_coordinator = Some(shutdown_coordinator);
1376        loop {
1377            select! {
1378                result = &mut run => return result,
1379                _ = &mut shutdown, if shutdown_coordinator.is_some() => {
1380                    shutdown_coordinator.take().expect("coordinator present per select guard").shutdown();
1381                }
1382            }
1383        }
1384    }
1385
1386    /// Runs the supervisor until the given `ShutdownHandle` signal is received.
1387    ///
1388    /// This is an internal variant of `run_with_shutdown` that takes a `ShutdownHandle` directly, used when spawning
1389    /// supervisors in dedicated runtimes where the shutdown signal is already wrapped in a `ShutdownHandle`.
1390    ///
1391    /// If `dataspace` is provided, the supervisor will use it instead of creating a new one. This is used to propagate
1392    /// the parent's dataspace across OS thread boundaries for dedicated runtimes.
1393    ///
1394    /// # Errors
1395    ///
1396    /// If the supervisor exceeds its restart limits, or fails to initialize a child process, an error is returned.
1397    pub(crate) async fn run_with_shutdown_inner(
1398        &mut self, process_shutdown: ShutdownHandle, dataspace: Option<DataspaceRegistry>,
1399    ) -> Result<(), SupervisorError> {
1400        let process =
1401            Process::supervisor_with_dataspace(&self.supervisor_id, None, dataspace).context(InvalidName {
1402                name: self.supervisor_id.to_string(),
1403            })?;
1404
1405        debug!(supervisor_id = %self.supervisor_id, "Supervisor starting.");
1406        self.run_inner(process.clone(), process_shutdown)
1407            .into_process_future(process)
1408            .await
1409    }
1410
1411    fn inner_clone(&self) -> Self {
1412        // This is no different than if we just implemented `Clone` directly, but it allows us to avoid exposing a
1413        // _public_ implementation of `Clone`, which we don't want normal users to be able to do. We only need this
1414        // internally to support nested supervisors.
1415        Self {
1416            supervisor_id: Arc::clone(&self.supervisor_id),
1417            child_specs: self.child_specs.clone(),
1418            runtime_mode: self.runtime_mode.clone(),
1419            current_tx: Arc::clone(&self.current_tx),
1420            id_counter: Arc::clone(&self.id_counter),
1421            active: Arc::clone(&self.active),
1422            node: Arc::clone(&self.node),
1423        }
1424    }
1425}
1426
1427#[cfg(test)]
1428mod tests {
1429    use std::{
1430        future::pending,
1431        sync::atomic::{AtomicBool, AtomicUsize, Ordering},
1432    };
1433
1434    use async_trait::async_trait;
1435    use saluki_common::sync::shutdown::ShutdownCoordinator;
1436    use saluki_metrics::test::TestRecorder;
1437    use tokio::{
1438        sync::oneshot,
1439        task::JoinHandle,
1440        time::{sleep, timeout},
1441    };
1442
1443    use super::*;
1444    use crate::runtime::{self, FnWorker, NodeKind, NodeSnapshot, NodeState, SupervisorFuture};
1445    use crate::test_support::wait_until;
1446
1447    /// Behavior for a mock worker during initialization.
1448    #[derive(Clone)]
1449    enum InitBehavior {
1450        /// Initialization succeeds immediately.
1451        Instant,
1452
1453        /// Initialization takes the given duration before succeeding.
1454        Slow(Duration),
1455
1456        /// Initialization fails with the given message.
1457        Fail(&'static str),
1458    }
1459
1460    /// Behavior for a mock worker during runtime (after initialization).
1461    #[derive(Clone)]
1462    enum RunBehavior {
1463        /// Runs until shutdown is received.
1464        UntilShutdown,
1465
1466        /// Fails with the given error message after the given delay.
1467        FailAfter(Duration, &'static str),
1468
1469        /// Completes successfully after the given delay.
1470        CompleteAfter(Duration),
1471
1472        /// On shutdown, sleeps for the given duration before exiting (to exercise concurrent draining).
1473        SlowShutdown(Duration),
1474
1475        /// Ignores shutdown entirely and runs forever (to exercise abort-at-deadline).
1476        IgnoreShutdown,
1477
1478        /// Panics after the given delay, unless shutdown arrives first.
1479        PanicAfter(Duration),
1480    }
1481
1482    /// A configurable mock worker for testing supervisor behavior.
1483    struct MockWorker {
1484        name: &'static str,
1485        init_behavior: InitBehavior,
1486        run_behavior: RunBehavior,
1487        start_count: Arc<AtomicUsize>,
1488        finish_count: Arc<AtomicUsize>,
1489        brutal_shutdown: bool,
1490        graceful_timeout: Duration,
1491    }
1492
1493    impl MockWorker {
1494        /// Creates a worker that runs until shutdown.
1495        fn long_running(name: &'static str) -> Self {
1496            Self {
1497                name,
1498                init_behavior: InitBehavior::Instant,
1499                run_behavior: RunBehavior::UntilShutdown,
1500                start_count: Arc::new(AtomicUsize::new(0)),
1501                finish_count: Arc::new(AtomicUsize::new(0)),
1502                brutal_shutdown: false,
1503                graceful_timeout: Duration::from_millis(500),
1504            }
1505        }
1506
1507        /// Creates a worker that fails after the given delay.
1508        fn failing(name: &'static str, delay: Duration) -> Self {
1509            Self {
1510                name,
1511                init_behavior: InitBehavior::Instant,
1512                run_behavior: RunBehavior::FailAfter(delay, "worker failed"),
1513                start_count: Arc::new(AtomicUsize::new(0)),
1514                finish_count: Arc::new(AtomicUsize::new(0)),
1515                brutal_shutdown: false,
1516                graceful_timeout: Duration::from_millis(500),
1517            }
1518        }
1519
1520        /// Creates a worker that completes successfully after the given delay.
1521        fn completing(name: &'static str, delay: Duration) -> Self {
1522            Self {
1523                name,
1524                init_behavior: InitBehavior::Instant,
1525                run_behavior: RunBehavior::CompleteAfter(delay),
1526                start_count: Arc::new(AtomicUsize::new(0)),
1527                finish_count: Arc::new(AtomicUsize::new(0)),
1528                brutal_shutdown: false,
1529                graceful_timeout: Duration::from_millis(500),
1530            }
1531        }
1532
1533        /// Creates a worker that sleeps for `delay` after observing shutdown before exiting.
1534        fn slow_shutdown(name: &'static str, delay: Duration) -> Self {
1535            Self {
1536                name,
1537                init_behavior: InitBehavior::Instant,
1538                run_behavior: RunBehavior::SlowShutdown(delay),
1539                start_count: Arc::new(AtomicUsize::new(0)),
1540                finish_count: Arc::new(AtomicUsize::new(0)),
1541                brutal_shutdown: false,
1542                graceful_timeout: Duration::from_millis(500),
1543            }
1544        }
1545
1546        /// Creates a worker that never reacts to shutdown.
1547        fn ignore_shutdown(name: &'static str) -> Self {
1548            Self {
1549                name,
1550                init_behavior: InitBehavior::Instant,
1551                run_behavior: RunBehavior::IgnoreShutdown,
1552                start_count: Arc::new(AtomicUsize::new(0)),
1553                finish_count: Arc::new(AtomicUsize::new(0)),
1554                brutal_shutdown: false,
1555                graceful_timeout: Duration::from_millis(500),
1556            }
1557        }
1558
1559        /// Creates a worker that panics after the given delay.
1560        fn panicking(name: &'static str, delay: Duration) -> Self {
1561            Self {
1562                name,
1563                init_behavior: InitBehavior::Instant,
1564                run_behavior: RunBehavior::PanicAfter(delay),
1565                start_count: Arc::new(AtomicUsize::new(0)),
1566                finish_count: Arc::new(AtomicUsize::new(0)),
1567                brutal_shutdown: false,
1568                graceful_timeout: Duration::from_millis(500),
1569            }
1570        }
1571
1572        /// Creates a worker that fails during initialization.
1573        fn init_failure(name: &'static str) -> Self {
1574            Self {
1575                name,
1576                init_behavior: InitBehavior::Fail("init failed"),
1577                run_behavior: RunBehavior::UntilShutdown,
1578                start_count: Arc::new(AtomicUsize::new(0)),
1579                finish_count: Arc::new(AtomicUsize::new(0)),
1580                brutal_shutdown: false,
1581                graceful_timeout: Duration::from_millis(500),
1582            }
1583        }
1584
1585        /// Creates a worker with slow initialization.
1586        fn slow_init(name: &'static str, init_delay: Duration) -> Self {
1587            Self {
1588                name,
1589                init_behavior: InitBehavior::Slow(init_delay),
1590                run_behavior: RunBehavior::UntilShutdown,
1591                start_count: Arc::new(AtomicUsize::new(0)),
1592                finish_count: Arc::new(AtomicUsize::new(0)),
1593                brutal_shutdown: false,
1594                graceful_timeout: Duration::from_millis(500),
1595            }
1596        }
1597
1598        /// Returns a shared handle to the start count for this worker.
1599        ///
1600        /// The start count ticks up the instant the worker's run future begins executing, which is *before* any
1601        /// programmed delay elapses. It records that the worker started (or was restarted), not that it ran to any
1602        /// particular outcome.
1603        fn start_count(&self) -> Arc<AtomicUsize> {
1604            Arc::clone(&self.start_count)
1605        }
1606
1607        /// Returns a shared handle to the finish count for this worker.
1608        ///
1609        /// The finish count ticks up only when the worker runs to its *own* programmed terminal state -- a
1610        /// [`RunBehavior::FailAfter`] failure, a [`RunBehavior::CompleteAfter`] completion, or a
1611        /// [`RunBehavior::SlowShutdown`] drain that finished -- and not when it is cut short by an abort. Tests use it
1612        /// to wait for a worker to actually fail or complete (rather than merely start) before asserting on restart
1613        /// behavior, so the failure/completion path is genuinely exercised.
1614        fn finish_count(&self) -> Arc<AtomicUsize> {
1615            Arc::clone(&self.finish_count)
1616        }
1617
1618        /// Configures this worker to use a `Brutal` shutdown strategy (immediate abort, no graceful wait).
1619        fn with_brutal_shutdown(mut self) -> Self {
1620            self.brutal_shutdown = true;
1621            self
1622        }
1623
1624        /// Overrides the worker's graceful shutdown timeout (defaults to 500 milliseconds).
1625        fn with_graceful_timeout(mut self, timeout: Duration) -> Self {
1626            self.graceful_timeout = timeout;
1627            self
1628        }
1629    }
1630
1631    #[async_trait]
1632    impl Supervisable for MockWorker {
1633        fn name(&self) -> &str {
1634            self.name
1635        }
1636
1637        fn shutdown_strategy(&self) -> ShutdownStrategy {
1638            if self.brutal_shutdown {
1639                ShutdownStrategy::Brutal
1640            } else {
1641                ShutdownStrategy::Graceful(self.graceful_timeout)
1642            }
1643        }
1644
1645        async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
1646            match &self.init_behavior {
1647                InitBehavior::Instant => {}
1648                InitBehavior::Slow(delay) => {
1649                    sleep(*delay).await;
1650                }
1651                InitBehavior::Fail(msg) => {
1652                    return Err(InitializationError::Failed {
1653                        source: GenericError::msg(*msg),
1654                    });
1655                }
1656            }
1657
1658            let start_count = Arc::clone(&self.start_count);
1659            let finish_count = Arc::clone(&self.finish_count);
1660            let run_behavior = self.run_behavior.clone();
1661
1662            Ok(Box::pin(async move {
1663                start_count.fetch_add(1, Ordering::SeqCst);
1664
1665                match run_behavior {
1666                    RunBehavior::UntilShutdown => {
1667                        process_shutdown.await;
1668                        Ok(())
1669                    }
1670                    RunBehavior::FailAfter(delay, msg) => {
1671                        select! {
1672                            _ = sleep(delay) => {
1673                                // Ran to our own programmed failure rather than being cut short by shutdown; record
1674                                // it so tests can wait for the failure to actually happen before asserting.
1675                                finish_count.fetch_add(1, Ordering::SeqCst);
1676                                Err(GenericError::msg(msg))
1677                            }
1678                            _ = process_shutdown => {
1679                                Ok(())
1680                            }
1681                        }
1682                    }
1683                    RunBehavior::CompleteAfter(delay) => {
1684                        select! {
1685                            _ = sleep(delay) => {
1686                                // Ran to our own programmed completion rather than being cut short by shutdown.
1687                                finish_count.fetch_add(1, Ordering::SeqCst);
1688                                Ok(())
1689                            }
1690                            _ = process_shutdown => Ok(()),
1691                        }
1692                    }
1693                    RunBehavior::SlowShutdown(delay) => {
1694                        process_shutdown.await;
1695                        sleep(delay).await;
1696                        // Finished draining rather than being aborted partway through it.
1697                        finish_count.fetch_add(1, Ordering::SeqCst);
1698                        Ok(())
1699                    }
1700                    RunBehavior::IgnoreShutdown => {
1701                        // Hold the handle (so the supervisor counts us as outstanding) but never react to it.
1702                        let _hold = process_shutdown;
1703                        pending().await
1704                    }
1705                    RunBehavior::PanicAfter(delay) => {
1706                        select! {
1707                            _ = sleep(delay) => panic!("worker panicked"),
1708                            _ = process_shutdown => Ok(()),
1709                        }
1710                    }
1711                }
1712            }))
1713        }
1714    }
1715
1716    /// Helper: run a supervisor with a oneshot-based shutdown trigger.
1717    ///
1718    /// Returns the shutdown sender and a join handle for the run. The supervisor is polled to a running state (its
1719    /// static children spawned) via readiness polling rather than a blind startup sleep, so callers can rely on it
1720    /// being live on return.
1721    async fn run_supervisor_with_trigger(
1722        supervisor: Supervisor,
1723    ) -> (oneshot::Sender<()>, JoinHandle<Result<(), SupervisorError>>) {
1724        // Grab a handle before moving the supervisor into the run task so we can observe when it actually starts.
1725        let sup_handle = supervisor.handle();
1726        let mut supervisor = supervisor;
1727
1728        let (tx, rx) = oneshot::channel();
1729        let handle = tokio::spawn(async move { supervisor.run_with_shutdown(rx).await });
1730
1731        wait_until("supervisor is running", || sup_handle.is_running()).await;
1732        (tx, handle)
1733    }
1734
1735    /// Helper: awaits a spawned supervisor run to completion under a bounded timeout, unwrapping the join.
1736    ///
1737    /// Collapses the `timeout(..).await.unwrap().unwrap()` suffix repeated across the restart/shutdown tests into one
1738    /// call with useful panic messages.
1739    async fn join_supervisor(handle: JoinHandle<Result<(), SupervisorError>>) -> Result<(), SupervisorError> {
1740        timeout(Duration::from_secs(2), handle)
1741            .await
1742            .expect("supervisor should exit promptly")
1743            .expect("supervisor task should not panic")
1744    }
1745
1746    // -- Supervisor run mode tests ---------------------------------------------------------
1747
1748    #[tokio::test]
1749    async fn standalone_supervisor_shuts_down_cleanly() {
1750        let mut sup = Supervisor::new("test-sup").unwrap();
1751        sup.add_worker(MockWorker::long_running("worker1"));
1752        sup.add_worker(MockWorker::long_running("worker2"));
1753
1754        let (tx, handle) = run_supervisor_with_trigger(sup).await;
1755        tx.send(()).unwrap();
1756
1757        let result = join_supervisor(handle).await;
1758        assert!(result.is_ok());
1759    }
1760
1761    #[tokio::test]
1762    async fn nested_supervisor_shuts_down_cleanly() {
1763        let mut child_sup = Supervisor::new("child-sup").unwrap();
1764        child_sup.add_worker(MockWorker::long_running("inner-worker"));
1765
1766        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
1767        parent_sup.add_worker(MockWorker::long_running("outer-worker"));
1768        parent_sup.add_worker(child_sup);
1769
1770        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
1771        tx.send(()).unwrap();
1772
1773        let result = join_supervisor(handle).await;
1774        assert!(result.is_ok());
1775    }
1776
1777    #[tokio::test]
1778    async fn empty_supervisor_idles_until_shutdown() {
1779        // A supervisor with no static children is valid: it idles, waiting for dynamic children, and shuts down
1780        // cleanly when signalled. (Before dynamic children were folded in, this returned a `NoChildren` error.)
1781        let sup = Supervisor::new("empty-sup").unwrap();
1782
1783        let (tx, handle) = run_supervisor_with_trigger(sup).await;
1784        assert!(!handle.is_finished(), "an empty supervisor must idle rather than exit");
1785
1786        tx.send(()).unwrap();
1787        let result = join_supervisor(handle).await;
1788        assert!(result.is_ok());
1789    }
1790
1791    // -- Child restart behavior tests ------------------------------------------------------
1792
1793    #[tokio::test]
1794    async fn one_for_one_restarts_only_failed_child() {
1795        let failing = MockWorker::failing("failing-worker", Duration::from_millis(50));
1796        let failing_count = failing.start_count();
1797
1798        let stable = MockWorker::long_running("stable-worker");
1799        let stable_count = stable.start_count();
1800
1801        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1802            RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
1803        );
1804        sup.add_worker(stable);
1805        sup.add_worker(failing);
1806
1807        let (tx, handle) = run_supervisor_with_trigger(sup).await;
1808
1809        // Wait until the failing worker has actually been restarted (its second start), then shut down.
1810        wait_until("the failing worker has been restarted", || {
1811            failing_count.load(Ordering::SeqCst) >= 2
1812        })
1813        .await;
1814        let _ = tx.send(());
1815
1816        let result = join_supervisor(handle).await;
1817        assert!(result.is_ok());
1818
1819        // The failing worker should have been started multiple times.
1820        assert!(
1821            failing_count.load(Ordering::SeqCst) >= 2,
1822            "failing worker should have been restarted"
1823        );
1824        // The stable worker should only have been started once (never restarted).
1825        assert_eq!(
1826            stable_count.load(Ordering::SeqCst),
1827            1,
1828            "stable worker should not have been restarted"
1829        );
1830    }
1831
1832    #[tokio::test]
1833    async fn one_for_all_restarts_all_children() {
1834        let failing = MockWorker::failing("failing-worker", Duration::from_millis(50));
1835        let failing_count = failing.start_count();
1836
1837        let stable = MockWorker::long_running("stable-worker");
1838        let stable_count = stable.start_count();
1839
1840        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1841            RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
1842        );
1843        sup.add_worker(stable);
1844        sup.add_worker(failing);
1845
1846        let (tx, handle) = run_supervisor_with_trigger(sup).await;
1847
1848        // Wait until a one-for-all cycle has restarted both workers (each on its second start), then shut down.
1849        wait_until("both workers have been restarted", || {
1850            failing_count.load(Ordering::SeqCst) >= 2 && stable_count.load(Ordering::SeqCst) >= 2
1851        })
1852        .await;
1853        let _ = tx.send(());
1854
1855        let result = join_supervisor(handle).await;
1856        assert!(result.is_ok());
1857
1858        // Both workers should have been started multiple times.
1859        assert!(
1860            failing_count.load(Ordering::SeqCst) >= 2,
1861            "failing worker should have been restarted"
1862        );
1863        assert!(
1864            stable_count.load(Ordering::SeqCst) >= 2,
1865            "stable worker should also have been restarted"
1866        );
1867    }
1868
1869    #[tokio::test]
1870    async fn one_for_all_does_not_restart_temporary_children() {
1871        // A permanent worker that fails repeatedly drives one-for-all restarts; a temporary sibling is shut down with
1872        // the group on each cycle but, per OTP semantics, must never be brought back.
1873        let failing = MockWorker::failing("failing-worker", Duration::from_millis(50));
1874        let failing_count = failing.start_count();
1875
1876        let temp = MockWorker::long_running("temp-worker");
1877        let temp_count = temp.start_count();
1878
1879        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1880            RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
1881        );
1882        sup.add_worker(ChildSpecification::worker(temp).with_restart_type(RestartType::Temporary));
1883        sup.add_worker(failing);
1884
1885        let (tx, handle) = run_supervisor_with_trigger(sup).await;
1886
1887        // Wait until the permanent worker has driven at least one one-for-all restart, then shut down.
1888        wait_until("the permanent worker has been restarted", || {
1889            failing_count.load(Ordering::SeqCst) >= 2
1890        })
1891        .await;
1892        let _ = tx.send(());
1893
1894        let result = join_supervisor(handle).await;
1895        assert!(result.is_ok());
1896        assert!(
1897            failing_count.load(Ordering::SeqCst) >= 2,
1898            "permanent worker should have been restarted by one-for-all"
1899        );
1900        assert_eq!(
1901            temp_count.load(Ordering::SeqCst),
1902            1,
1903            "temporary child must not be restarted by a one-for-all group restart"
1904        );
1905    }
1906
1907    #[tokio::test]
1908    async fn one_for_all_restarts_transient_children() {
1909        // A transient child that exits cleanly is not restarted on its own, but a one-for-all restart triggered by a
1910        // sibling restarts it anyway -- matching OTP, where only temporary children are exempt from group restarts.
1911        let transient = MockWorker::completing("transient-worker", Duration::from_millis(30));
1912        let transient_count = transient.start_count();
1913
1914        // Fails after the transient has already exited cleanly, so the group restart is what brings the transient back.
1915        let failing = MockWorker::failing("failing-worker", Duration::from_millis(80));
1916
1917        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1918            RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
1919        );
1920        sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
1921        sup.add_worker(failing);
1922
1923        let (tx, handle) = run_supervisor_with_trigger(sup).await;
1924
1925        wait_until("the transient worker has been restarted by the group", || {
1926            transient_count.load(Ordering::SeqCst) >= 2
1927        })
1928        .await;
1929        let _ = tx.send(());
1930
1931        let result = join_supervisor(handle).await;
1932        assert!(result.is_ok());
1933        assert!(
1934            transient_count.load(Ordering::SeqCst) >= 2,
1935            "transient child must be restarted by a one-for-all group restart, even after a clean exit"
1936        );
1937    }
1938
1939    #[tokio::test]
1940    async fn transient_abnormal_exit_triggers_one_for_all() {
1941        // A transient child's *own* abnormal exit is restartable, so under one-for-all it triggers a whole-group
1942        // restart -- the sibling is restarted too, not just the transient.
1943        let transient = MockWorker::failing("transient-worker", Duration::from_millis(50));
1944        let transient_count = transient.start_count();
1945
1946        let stable = MockWorker::long_running("stable-worker");
1947        let stable_count = stable.start_count();
1948
1949        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1950            RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
1951        );
1952        sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
1953        sup.add_worker(stable);
1954
1955        let (tx, handle) = run_supervisor_with_trigger(sup).await;
1956
1957        wait_until("the abnormal exit has restarted both workers", || {
1958            transient_count.load(Ordering::SeqCst) >= 2 && stable_count.load(Ordering::SeqCst) >= 2
1959        })
1960        .await;
1961        let _ = tx.send(());
1962
1963        let result = join_supervisor(handle).await;
1964        assert!(result.is_ok());
1965        assert!(
1966            transient_count.load(Ordering::SeqCst) >= 2,
1967            "transient worker must be restarted after its own abnormal exit"
1968        );
1969        assert!(
1970            stable_count.load(Ordering::SeqCst) >= 2,
1971            "the transient's abnormal exit must trigger a one-for-all that also restarts the sibling"
1972        );
1973    }
1974
1975    #[tokio::test]
1976    async fn restart_limit_exceeded_shuts_down_supervisor() {
1977        let mut sup = Supervisor::new("test-sup")
1978            .unwrap()
1979            .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(1, Duration::from_secs(10)));
1980        // This worker fails immediately, which will exhaust the restart budget quickly.
1981        sup.add_worker(MockWorker::failing("fast-fail", Duration::ZERO));
1982
1983        let (tx, rx) = oneshot::channel::<()>();
1984        let handle = tokio::spawn(async move { sup.run_with_shutdown(rx).await });
1985
1986        let result = join_supervisor(handle).await;
1987        drop(tx);
1988
1989        assert!(matches!(result, Err(SupervisorError::Shutdown)));
1990    }
1991
1992    // -- Restart type tests ----------------------------------------------------------------
1993
1994    #[tokio::test]
1995    async fn temporary_child_is_not_restarted() {
1996        // A temporary worker that fails quickly, alongside a long-running worker that keeps the supervisor alive.
1997        let temp = MockWorker::failing("temp-worker", Duration::from_millis(50));
1998        let temp_started = temp.start_count();
1999        let temp_failed = temp.finish_count();
2000
2001        let stable = MockWorker::long_running("stable-worker");
2002
2003        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2004            RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2005        );
2006        sup.add_worker(stable);
2007        sup.add_worker(ChildSpecification::worker(temp).with_restart_type(RestartType::Temporary));
2008
2009        let (tx, handle) = run_supervisor_with_trigger(sup).await;
2010
2011        // Wait for the worker to *actually fail*, not merely start. `start_count` ticks up the instant the worker
2012        // begins running -- well before its 50ms failure -- so shutting down as soon as it reached 1 would tear the
2013        // supervisor down before the failure -> no-restart path ever ran, hiding a regression that restarted a
2014        // temporary child (or charged the failure against restart intensity). `finish_count` ticks only once the
2015        // worker runs to its own failure, so waiting on it genuinely exercises that path before we shut down.
2016        wait_until("the temporary worker has failed once", || {
2017            temp_failed.load(Ordering::SeqCst) == 1
2018        })
2019        .await;
2020        let _ = tx.send(());
2021
2022        let result = join_supervisor(handle).await;
2023        assert!(result.is_ok());
2024        assert_eq!(
2025            temp_started.load(Ordering::SeqCst),
2026            1,
2027            "temporary worker must not be restarted after it fails"
2028        );
2029    }
2030
2031    #[tokio::test]
2032    async fn transient_child_is_not_restarted_on_clean_exit() {
2033        let transient = MockWorker::completing("transient-worker", Duration::from_millis(50));
2034        let transient_started = transient.start_count();
2035        let transient_finished = transient.finish_count();
2036
2037        let stable = MockWorker::long_running("stable-worker");
2038
2039        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2040            RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2041        );
2042        sup.add_worker(stable);
2043        sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
2044
2045        let (tx, handle) = run_supervisor_with_trigger(sup).await;
2046
2047        // Wait for the worker to *actually complete*, not merely start: `start_count` ticks the instant it begins
2048        // running, so shutting down as soon as it reached 1 would drive the supervisor's teardown before the clean
2049        // exit -> no-restart path ran, hiding a regression that restarted a transient child after a clean exit.
2050        // `finish_count` ticks only once the worker runs to its own completion.
2051        wait_until("the transient worker has completed once", || {
2052            transient_finished.load(Ordering::SeqCst) == 1
2053        })
2054        .await;
2055        let _ = tx.send(());
2056
2057        let result = join_supervisor(handle).await;
2058        assert!(result.is_ok());
2059        assert_eq!(
2060            transient_started.load(Ordering::SeqCst),
2061            1,
2062            "transient worker must not be restarted after a clean exit"
2063        );
2064    }
2065
2066    #[tokio::test]
2067    async fn transient_child_is_restarted_on_failure() {
2068        let transient = MockWorker::failing("transient-worker", Duration::from_millis(50));
2069        let transient_count = transient.start_count();
2070
2071        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2072            RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2073        );
2074        sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
2075
2076        let (tx, handle) = run_supervisor_with_trigger(sup).await;
2077
2078        wait_until("the transient worker has been restarted", || {
2079            transient_count.load(Ordering::SeqCst) >= 2
2080        })
2081        .await;
2082        let _ = tx.send(());
2083
2084        let result = join_supervisor(handle).await;
2085        assert!(result.is_ok());
2086        assert!(
2087            transient_count.load(Ordering::SeqCst) >= 2,
2088            "transient worker must be restarted after an abnormal exit"
2089        );
2090    }
2091
2092    #[tokio::test]
2093    async fn transient_child_is_restarted_on_panic() {
2094        // The other half of what makes `transient` the right policy for a worker serving an object with a lifetime of
2095        // its own -- a cache's expiration loop, say. A panic is a bug worth recovering from, so the worker comes back;
2096        // the clean exit that follows the served object being dropped is taken at face value (see
2097        // `transient_child_is_not_restarted_on_clean_exit`), so it stays stopped. Under `Permanent` that clean exit
2098        // would be restarted straight into another clean exit, over and over, until the restart intensity failed the
2099        // whole supervisor.
2100        let transient = MockWorker::panicking("transient-worker", Duration::from_millis(50));
2101        let transient_count = transient.start_count();
2102
2103        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2104            RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2105        );
2106        sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
2107
2108        let (tx, handle) = run_supervisor_with_trigger(sup).await;
2109
2110        wait_until("the panicking transient worker has been restarted", || {
2111            transient_count.load(Ordering::SeqCst) >= 2
2112        })
2113        .await;
2114        let _ = tx.send(());
2115
2116        let result = join_supervisor(handle).await;
2117        assert!(result.is_ok());
2118        assert!(
2119            transient_count.load(Ordering::SeqCst) >= 2,
2120            "transient worker must be restarted after a panic"
2121        );
2122    }
2123
2124    #[tokio::test]
2125    async fn permanent_child_is_restarted_on_clean_exit() {
2126        // A permanent worker that completes cleanly must still be restarted -- this is what distinguishes
2127        // `Permanent` from `Transient`, which is left stopped after a clean exit.
2128        let permanent = MockWorker::completing("permanent-worker", Duration::from_millis(50));
2129        let permanent_count = permanent.start_count();
2130
2131        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2132            RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2133        );
2134        // Added with the default restart policy, which is `Permanent`.
2135        sup.add_worker(permanent);
2136
2137        let (tx, handle) = run_supervisor_with_trigger(sup).await;
2138
2139        wait_until("the permanent worker has been restarted", || {
2140            permanent_count.load(Ordering::SeqCst) >= 2
2141        })
2142        .await;
2143        let _ = tx.send(());
2144
2145        let result = join_supervisor(handle).await;
2146        assert!(result.is_ok());
2147        assert!(
2148            permanent_count.load(Ordering::SeqCst) >= 2,
2149            "permanent worker must be restarted even after a clean exit"
2150        );
2151    }
2152
2153    #[tokio::test]
2154    async fn temporary_failures_do_not_consume_restart_intensity() {
2155        // With intensity=1, two *restartable* failures within the period would shut the supervisor down. Here several
2156        // temporary workers all fail quickly. Because temporary exits aren't eligible for restart, they must not consume
2157        // the restart-intensity budget, and the supervisor must stay up.
2158        let mut sup = Supervisor::new("test-sup")
2159            .unwrap()
2160            .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(1, Duration::from_secs(10)));
2161
2162        let workers = [
2163            MockWorker::failing("temp-0", Duration::from_millis(20)),
2164            MockWorker::failing("temp-1", Duration::from_millis(20)),
2165            MockWorker::failing("temp-2", Duration::from_millis(20)),
2166            MockWorker::failing("temp-3", Duration::from_millis(20)),
2167            MockWorker::failing("temp-4", Duration::from_millis(20)),
2168        ];
2169        let started: Vec<_> = workers.iter().map(|w| w.start_count()).collect();
2170        let failed: Vec<_> = workers.iter().map(|w| w.finish_count()).collect();
2171        for worker in workers {
2172            sup.add_worker(ChildSpecification::worker(worker).with_restart_type(RestartType::Temporary));
2173        }
2174        // A long-running worker so the supervisor doesn't simply idle once the temporaries are gone.
2175        sup.add_worker(MockWorker::long_running("stable-worker"));
2176
2177        let (tx, handle) = run_supervisor_with_trigger(sup).await;
2178        // Wait for every temporary worker to *actually fail* on its own. Keying off `start_count` would let shutdown
2179        // cut them short before their failures ran, so the supervisor would never get the chance to (mis)charge those
2180        // failures against its intensity=1 budget -- hiding the very regression this guards against.
2181        wait_until("every temporary worker has failed once", || {
2182            failed.iter().all(|c| c.load(Ordering::SeqCst) == 1)
2183        })
2184        .await;
2185        let _ = tx.send(());
2186
2187        let result = join_supervisor(handle).await;
2188        assert!(
2189            result.is_ok(),
2190            "supervisor must not trip its restart limit on temporary exits"
2191        );
2192        for count in started {
2193            assert_eq!(
2194                count.load(Ordering::SeqCst),
2195                1,
2196                "each temporary worker runs exactly once"
2197            );
2198        }
2199    }
2200
2201    #[tokio::test]
2202    async fn transient_clean_exits_do_not_consume_restart_intensity() {
2203        // With intensity=1, two *restartable* exits within the period would shut the supervisor down. Here several
2204        // transient workers all complete cleanly. A transient child's clean exit isn't eligible for restart, so it
2205        // must not consume the restart-intensity budget, and the supervisor must stay up.
2206        let mut sup = Supervisor::new("test-sup")
2207            .unwrap()
2208            .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(1, Duration::from_secs(10)));
2209
2210        let workers = [
2211            MockWorker::completing("transient-0", Duration::from_millis(20)),
2212            MockWorker::completing("transient-1", Duration::from_millis(20)),
2213            MockWorker::completing("transient-2", Duration::from_millis(20)),
2214            MockWorker::completing("transient-3", Duration::from_millis(20)),
2215            MockWorker::completing("transient-4", Duration::from_millis(20)),
2216        ];
2217        let started: Vec<_> = workers.iter().map(|w| w.start_count()).collect();
2218        let finished: Vec<_> = workers.iter().map(|w| w.finish_count()).collect();
2219        for worker in workers {
2220            sup.add_worker(ChildSpecification::worker(worker).with_restart_type(RestartType::Transient));
2221        }
2222        // A long-running worker so the supervisor doesn't simply idle once the transients have completed.
2223        sup.add_worker(MockWorker::long_running("stable-worker"));
2224
2225        let (tx, handle) = run_supervisor_with_trigger(sup).await;
2226        // Wait for every transient to *actually complete* on its own. Keying off `start_count` would let shutdown cut
2227        // the workers short before their clean exits ran, so the supervisor would never get the chance to (mis)charge
2228        // those exits against its intensity=1 budget -- hiding the very regression this guards against.
2229        wait_until("every transient worker has completed once", || {
2230            finished.iter().all(|c| c.load(Ordering::SeqCst) == 1)
2231        })
2232        .await;
2233        let _ = tx.send(());
2234
2235        let result = join_supervisor(handle).await;
2236        assert!(
2237            result.is_ok(),
2238            "supervisor must not trip its restart limit on clean transient exits"
2239        );
2240        for count in started {
2241            assert_eq!(
2242                count.load(Ordering::SeqCst),
2243                1,
2244                "each transient worker runs exactly once"
2245            );
2246        }
2247    }
2248
2249    #[tokio::test]
2250    async fn supervisor_idles_when_all_temporary_children_exit() {
2251        // When every static child is temporary and they all exit, the worker set drains. The supervisor must not panic
2252        // or exit on its own; it must keep running and remain able to accept new (dynamic) work until shutdown is
2253        // triggered.
2254        let temp_a = MockWorker::completing("temp-a", Duration::from_millis(10));
2255        let a_finished = temp_a.finish_count();
2256        let temp_b = MockWorker::completing("temp-b", Duration::from_millis(10));
2257        let b_finished = temp_b.finish_count();
2258
2259        let mut sup = Supervisor::new("test-sup").unwrap();
2260        let handle = sup.handle();
2261        sup.add_worker(ChildSpecification::worker(temp_a).with_restart_type(RestartType::Temporary));
2262        sup.add_worker(ChildSpecification::worker(temp_b).with_restart_type(RestartType::Temporary));
2263
2264        let (tx, run) = run_supervisor_with_trigger(sup).await;
2265
2266        // Wait for both temporary children to actually complete -- draining the worker set to empty -- before probing.
2267        // Keying off `start_count` could spawn the probe child before the set ever emptied, letting a supervisor that
2268        // (wrongly) exited once its last child left slip through.
2269        wait_until("both temporary children have completed", || {
2270            a_finished.load(Ordering::SeqCst) == 1 && b_finished.load(Ordering::SeqCst) == 1
2271        })
2272        .await;
2273
2274        // The supervisor must still be alive after its worker set empties: spawning a new dynamic child succeeds and
2275        // runs, which is only possible if the supervise loop kept running rather than exiting when the last child left.
2276        let dynamic = MockWorker::long_running("late-comer");
2277        let dynamic_count = dynamic.start_count();
2278        handle.spawn(dynamic);
2279        wait_until("the late dynamic child has started", || {
2280            dynamic_count.load(Ordering::SeqCst) == 1
2281        })
2282        .await;
2283        assert!(
2284            handle.is_running(),
2285            "supervisor must keep running after all temporary children exit"
2286        );
2287
2288        tx.send(()).unwrap();
2289        let result = join_supervisor(run).await;
2290        assert!(result.is_ok());
2291    }
2292
2293    // -- Significant child / auto-shutdown tests -------------------------------------------
2294
2295    #[tokio::test]
2296    async fn significant_child_drives_auto_shutdown() {
2297        // With `AnySignificant`, a significant child terminating (even cleanly, and without being restarted) must
2298        // shut the supervisor down, surfacing the significant-exit error.
2299        let mut sup = Supervisor::new("test-sup")
2300            .unwrap()
2301            .with_auto_shutdown(AutoShutdown::AnySignificant);
2302        sup.add_worker(MockWorker::long_running("stable"));
2303        sup.add_worker(
2304            ChildSpecification::worker(MockWorker::completing("significant", Duration::from_millis(50)))
2305                .with_restart_type(RestartType::Temporary)
2306                .with_significant(true),
2307        );
2308
2309        // Hold the shutdown sender so the only thing that can stop the supervisor is the significant child.
2310        let (_tx, rx) = oneshot::channel::<()>();
2311        let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2312            .await
2313            .unwrap();
2314        assert!(matches!(result, Err(SupervisorError::SignificantChildExited)));
2315    }
2316
2317    #[tokio::test]
2318    async fn significant_child_added_before_the_auto_shutdown_policy_still_drives_it() {
2319        // `with_auto_shutdown` consumes the supervisor while `add_worker` borrows it, so adding children first and
2320        // setting the policy afterwards is a perfectly good way to build one up. Nothing about registration may
2321        // assume the policy is already final -- an earlier version of the inert-significance check read it at
2322        // registration time and flagged this ordering as a mistake.
2323        let mut sup = Supervisor::new("test-sup").unwrap();
2324        sup.add_worker(MockWorker::long_running("stable"));
2325        sup.add_worker(
2326            runtime::supervisable(MockWorker::completing("significant", Duration::from_millis(50)))
2327                .temporary()
2328                .with_significant(true)
2329                .build(),
2330        );
2331        let mut sup = sup.with_auto_shutdown(AutoShutdown::AnySignificant);
2332
2333        let (_tx, rx) = oneshot::channel::<()>();
2334        let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2335            .await
2336            .unwrap();
2337        assert!(
2338            matches!(result, Err(SupervisorError::SignificantChildExited)),
2339            "the policy set after registration should still have applied, got {result:?}"
2340        );
2341    }
2342
2343    #[tokio::test]
2344    async fn non_significant_exit_does_not_auto_shutdown() {
2345        // Even with `AnySignificant` set, a non-significant child exiting must not shut the supervisor down.
2346        let plain = MockWorker::completing("plain", Duration::from_millis(10));
2347        let plain_finished = plain.finish_count();
2348
2349        let mut sup = Supervisor::new("test-sup")
2350            .unwrap()
2351            .with_auto_shutdown(AutoShutdown::AnySignificant);
2352        let handle = sup.handle();
2353        sup.add_worker(MockWorker::long_running("stable"));
2354        sup.add_worker(ChildSpecification::worker(plain).with_restart_type(RestartType::Temporary));
2355
2356        let (tx, run) = run_supervisor_with_trigger(sup).await;
2357
2358        // Let the non-significant child actually run to completion -- not merely start. Its completion is what could
2359        // (wrongly) trip `AnySignificant`, so we must observe the real exit before probing liveness; keying off
2360        // `start_count` could assert before the completion was ever processed.
2361        wait_until("the non-significant child has completed", || {
2362            plain_finished.load(Ordering::SeqCst) == 1
2363        })
2364        .await;
2365
2366        // The supervisor must still be alive after the non-significant child exits (had it been treated as
2367        // significant, `AnySignificant` would have torn the supervisor down). Spawning a dynamic child and observing
2368        // it start proves the supervise loop is still running.
2369        let dynamic = MockWorker::long_running("late-comer");
2370        let dynamic_count = dynamic.start_count();
2371        handle.spawn(dynamic);
2372        wait_until("the late dynamic child has started", || {
2373            dynamic_count.load(Ordering::SeqCst) == 1
2374        })
2375        .await;
2376        assert!(
2377            handle.is_running(),
2378            "a non-significant child exiting must not trigger auto-shutdown"
2379        );
2380
2381        tx.send(()).unwrap();
2382        let result = join_supervisor(run).await;
2383        assert!(result.is_ok());
2384    }
2385
2386    #[tokio::test]
2387    async fn all_significant_waits_for_last() {
2388        // With `AllSignificant`, the supervisor shuts down only once *all* significant children have terminated.
2389        let mut sup = Supervisor::new("test-sup")
2390            .unwrap()
2391            .with_auto_shutdown(AutoShutdown::AllSignificant);
2392        sup.add_worker(
2393            ChildSpecification::worker(MockWorker::completing("sig-a", Duration::from_millis(50)))
2394                .with_restart_type(RestartType::Temporary)
2395                .with_significant(true),
2396        );
2397        sup.add_worker(
2398            ChildSpecification::worker(MockWorker::completing("sig-b", Duration::from_millis(250)))
2399                .with_restart_type(RestartType::Temporary)
2400                .with_significant(true),
2401        );
2402
2403        let (_tx, rx) = oneshot::channel::<()>();
2404        let start = std::time::Instant::now();
2405        let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2406            .await
2407            .unwrap();
2408        let elapsed = start.elapsed();
2409
2410        assert!(matches!(result, Err(SupervisorError::SignificantChildExited)));
2411        // The first significant child exits at ~50ms but must NOT trigger shutdown; only the second (~250ms) does.
2412        assert!(
2413            elapsed >= Duration::from_millis(200),
2414            "auto-shutdown must wait for all significant children (took {elapsed:?})"
2415        );
2416    }
2417
2418    // -- Initialization failure tests ------------------------------------------------------
2419
2420    #[tokio::test]
2421    async fn init_failure_propagates_with_child_name() {
2422        let mut sup = Supervisor::new("test-sup").unwrap();
2423        sup.add_worker(MockWorker::long_running("good-worker"));
2424        sup.add_worker(MockWorker::init_failure("bad-worker"));
2425
2426        let (_tx, rx) = oneshot::channel::<()>();
2427        let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2428            .await
2429            .unwrap();
2430
2431        match result {
2432            Err(SupervisorError::FailedToInitialize { child_name, .. }) => {
2433                assert_eq!(child_name, "bad-worker");
2434            }
2435            other => panic!("expected FailedToInitialize, got: {:?}", other),
2436        }
2437    }
2438
2439    #[tokio::test]
2440    async fn init_failure_does_not_trigger_restart() {
2441        let init_fail = MockWorker::init_failure("bad-worker");
2442        let start_count = init_fail.start_count();
2443
2444        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2445            RestartStrategy::one_to_one().with_intensity_and_period(10, Duration::from_secs(10)),
2446        );
2447        sup.add_worker(init_fail);
2448
2449        let (_tx, rx) = oneshot::channel::<()>();
2450        let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2451            .await
2452            .unwrap();
2453
2454        assert!(matches!(result, Err(SupervisorError::FailedToInitialize { .. })));
2455        // The worker never got past init, so start_count should be 0.
2456        assert_eq!(start_count.load(Ordering::SeqCst), 0);
2457    }
2458
2459    // -- Shutdown responsiveness tests -----------------------------------------------------
2460
2461    #[tokio::test]
2462    async fn shutdown_completes_promptly_in_steady_state() {
2463        let mut sup = Supervisor::new("test-sup").unwrap();
2464        sup.add_worker(MockWorker::long_running("worker1"));
2465        sup.add_worker(MockWorker::long_running("worker2"));
2466
2467        let (tx, handle) = run_supervisor_with_trigger(sup).await;
2468        tx.send(()).unwrap();
2469
2470        // Shutdown should complete well within 1 second (workers respond to shutdown signal immediately).
2471        let result = timeout(Duration::from_secs(1), handle).await;
2472        assert!(result.is_ok(), "shutdown should complete promptly");
2473    }
2474
2475    #[tokio::test]
2476    async fn shutdown_during_slow_init_completes_promptly() {
2477        let mut sup = Supervisor::new("test-sup").unwrap();
2478        // This worker takes 30 seconds to initialize — but we'll trigger shutdown immediately.
2479        sup.add_worker(MockWorker::slow_init("slow-worker", Duration::from_secs(30)));
2480
2481        let (tx, rx) = oneshot::channel();
2482        let handle = tokio::spawn(async move { sup.run_with_shutdown(rx).await });
2483
2484        // Give the supervisor just enough time to spawn the task, then trigger shutdown.
2485        sleep(Duration::from_millis(20)).await;
2486        tx.send(()).unwrap();
2487
2488        // Shutdown should complete quickly even though the worker hasn't finished initializing.
2489        // The supervisor loop sees the shutdown signal and aborts the still-initializing task.
2490        let result = timeout(Duration::from_secs(2), handle).await;
2491        assert!(result.is_ok(), "shutdown during slow init should complete promptly");
2492    }
2493
2494    // -- Dynamic children tests ------------------------------------------------------------
2495
2496    #[tokio::test]
2497    async fn dynamic_children_spawn_after_start() {
2498        let sup = Supervisor::new("dyn-sup").unwrap();
2499        let handle = sup.handle();
2500        let (tx, run) = run_supervisor_with_trigger(sup).await;
2501        wait_until("supervisor is running", || handle.is_running()).await;
2502
2503        let c1 = MockWorker::long_running("c1");
2504        let c2 = MockWorker::long_running("c2");
2505        let c1_count = c1.start_count();
2506        let c2_count = c2.start_count();
2507        handle.spawn(c1);
2508        handle.spawn(c2);
2509
2510        wait_until("both dynamic children have started", || {
2511            c1_count.load(Ordering::SeqCst) == 1 && c2_count.load(Ordering::SeqCst) == 1
2512        })
2513        .await;
2514        assert_eq!(handle.active_children(), 2);
2515
2516        tx.send(()).unwrap();
2517        let result = join_supervisor(run).await;
2518        assert!(result.is_ok());
2519        assert_eq!(
2520            handle.active_children(),
2521            0,
2522            "all dynamic children must be drained on shutdown"
2523        );
2524    }
2525
2526    #[tokio::test]
2527    async fn temporary_dynamic_child_failure_is_isolated() {
2528        // A dynamic child added with the default config (temporary, not significant) is fault-isolated: its failure is
2529        // reaped and removed without restarting it or disturbing the supervisor or its siblings.
2530        let sup = Supervisor::new("dyn-sup").unwrap();
2531        let handle = sup.handle();
2532        let (tx, run) = run_supervisor_with_trigger(sup).await;
2533        wait_until("supervisor is running", || handle.is_running()).await;
2534
2535        let failing = MockWorker::failing("boom", Duration::from_millis(20));
2536        let failing_count = failing.start_count();
2537        handle.spawn(failing);
2538        wait_until("the failing dynamic child has run once", || {
2539            failing_count.load(Ordering::SeqCst) == 1
2540        })
2541        .await;
2542        wait_until("all dynamic children have drained", || handle.active_children() == 0).await;
2543
2544        sleep(Duration::from_millis(50)).await;
2545        assert!(
2546            handle.is_running(),
2547            "supervisor stays up after an isolated child failure"
2548        );
2549        assert_eq!(
2550            failing_count.load(Ordering::SeqCst),
2551            1,
2552            "a temporary child is never restarted"
2553        );
2554
2555        // It still accepts new children.
2556        handle.spawn(MockWorker::long_running("c2"));
2557        wait_until("one dynamic child is running", || handle.active_children() == 1).await;
2558
2559        tx.send(()).unwrap();
2560        let result = join_supervisor(run).await;
2561        assert!(result.is_ok());
2562    }
2563
2564    #[tokio::test]
2565    async fn temporary_dynamic_child_panic_is_isolated() {
2566        // A panicking temporary, non-significant child is isolated exactly like an error exit.
2567        let sup = Supervisor::new("dyn-sup").unwrap();
2568        let handle = sup.handle();
2569        let (tx, run) = run_supervisor_with_trigger(sup).await;
2570        wait_until("supervisor is running", || handle.is_running()).await;
2571
2572        handle.spawn(MockWorker::panicking("boom", Duration::from_millis(20)));
2573        wait_until("all dynamic children have drained", || handle.active_children() == 0).await;
2574
2575        sleep(Duration::from_millis(50)).await;
2576        assert!(handle.is_running(), "supervisor stays up after an isolated child panic");
2577
2578        tx.send(()).unwrap();
2579        let result = join_supervisor(run).await;
2580        assert!(result.is_ok());
2581    }
2582
2583    #[tokio::test]
2584    async fn significant_dynamic_child_failure_shuts_down_supervisor() {
2585        // A dynamic child added as significant, under `AutoShutdown::AnySignificant`, drives the supervisor to shut
2586        // down when it terminates -- the opt-in mechanism that replaces the old escalate-on-error behavior.
2587        let sup = Supervisor::new("dyn-sup")
2588            .unwrap()
2589            .with_auto_shutdown(AutoShutdown::AnySignificant);
2590        let handle = sup.handle();
2591        let (_tx, run) = run_supervisor_with_trigger(sup).await;
2592        wait_until("supervisor is running", || handle.is_running()).await;
2593
2594        handle.spawn(
2595            ChildSpecification::worker(MockWorker::failing("boom", Duration::from_millis(20))).with_significant(true),
2596        );
2597
2598        let result = join_supervisor(run).await;
2599        assert!(matches!(result, Err(SupervisorError::SignificantChildExited)));
2600    }
2601
2602    #[tokio::test]
2603    async fn dynamic_spawn_outside_a_run_is_accepted_and_dropped() {
2604        // Spawning is infallible in the same sense `tokio::spawn` is: the child is always accepted, but a child handed
2605        // to a supervisor that isn't running is never started. That holds both before a run and after one, and a child
2606        // spawned before the run must not be held over and started by it -- children belong to a run, not to the
2607        // supervisor across runs.
2608        let sup = Supervisor::new("dyn-sup").unwrap();
2609        let handle = sup.handle();
2610
2611        assert!(!handle.is_running());
2612        let before = MockWorker::long_running("before-start");
2613        let before_count = before.start_count();
2614        handle.spawn(before);
2615
2616        // Once it's running, spawns do start children.
2617        let (tx, run) = run_supervisor_with_trigger(sup).await;
2618        wait_until("supervisor is running", || handle.is_running()).await;
2619        let worker = MockWorker::long_running("after-start");
2620        let started = worker.start_count();
2621        handle.spawn(worker);
2622        wait_until("the dynamic child has started", || started.load(Ordering::SeqCst) == 1).await;
2623        assert_eq!(
2624            before_count.load(Ordering::SeqCst),
2625            0,
2626            "a child spawned before the run must not be started by it"
2627        );
2628
2629        tx.send(()).unwrap();
2630        let result = join_supervisor(run).await;
2631        assert!(result.is_ok());
2632
2633        // And once it has shut down there is nothing left to start children either.
2634        wait_until("the supervisor has stopped", || !handle.is_running()).await;
2635        let after = MockWorker::long_running("after-shutdown");
2636        let after_count = after.start_count();
2637        handle.spawn(after);
2638
2639        sleep(Duration::from_millis(50)).await;
2640        assert_eq!(
2641            after_count.load(Ordering::SeqCst),
2642            0,
2643            "a child spawned after shutdown must never start"
2644        );
2645    }
2646
2647    #[tokio::test]
2648    async fn dynamic_spawn_allocates_an_id_eagerly() {
2649        // The id comes back synchronously, before the supervisor has picked the child up, so it can't depend on
2650        // registration having happened. Ids come from the supervisor's shared counter, so with no static children the
2651        // first dynamic child takes id 0.
2652        let sup = Supervisor::new("dyn-sup").unwrap();
2653        let handle = sup.handle();
2654        let (tx, run) = run_supervisor_with_trigger(sup).await;
2655        wait_until("supervisor is running", || handle.is_running()).await;
2656
2657        let worker = MockWorker::long_running("c");
2658        let started = worker.start_count();
2659        let id = handle.spawn(worker);
2660        assert_eq!(id.as_u64(), 0);
2661        wait_until("the dynamic child has started", || started.load(Ordering::SeqCst) == 1).await;
2662
2663        tx.send(()).unwrap();
2664        let result = join_supervisor(run).await;
2665        assert!(result.is_ok());
2666    }
2667
2668    #[tokio::test]
2669    async fn dynamic_child_with_an_unusable_name_runs_under_a_placeholder() {
2670        // A name that sanitizes to nothing can't be used as a process name. Spawning is infallible, so rather than
2671        // quietly discarding work the caller was told had been accepted, the child runs under a placeholder segment.
2672        // The poll metric is what proves the substituted name is what the child actually ran as, rather than the child
2673        // merely having started somehow.
2674        let recorder = TestRecorder::default();
2675        let _guard = metrics::set_default_local_recorder(&recorder);
2676
2677        let sup = Supervisor::new("dyn-sup").unwrap();
2678        let handle = sup.handle();
2679        let (tx, run) = run_supervisor_with_trigger(sup).await;
2680        wait_until("supervisor is running", || handle.is_running()).await;
2681
2682        let worker = MockWorker::long_running("");
2683        let started = worker.start_count();
2684        handle.spawn(worker);
2685        wait_until("the unnamed dynamic child has started", || {
2686            started.load(Ordering::SeqCst) == 1
2687        })
2688        .await;
2689
2690        // The supervisor stays up and still accepts normally-named children.
2691        assert!(handle.is_running());
2692        handle.spawn(MockWorker::long_running("ok"));
2693        wait_until("both dynamic children are running", || handle.active_children() == 2).await;
2694
2695        tx.send(()).unwrap();
2696        let result = join_supervisor(run).await;
2697        assert!(result.is_ok());
2698
2699        let polls = recorder.counter(("runtime_task_poll_count", &[("task_name", "dyn_sup.unnamed")]));
2700        assert!(
2701            polls.is_some_and(|polls| polls > 0),
2702            "the child should have run under the placeholder name, got {polls:?}"
2703        );
2704    }
2705
2706    #[tokio::test]
2707    async fn dynamic_spawns_are_not_capped() {
2708        // Spawn requests are queued without a bound, so a burst well past what any fixed-capacity channel would hold
2709        // still starts every child. Being accepted means being started -- a depth limit could only deliver that by
2710        // discarding work the caller was already told had been taken.
2711        const CHILDREN: usize = 2048;
2712
2713        let sup = Supervisor::new("dyn-sup").unwrap();
2714        let handle = sup.handle();
2715        let (tx, run) = run_supervisor_with_trigger(sup).await;
2716        wait_until("supervisor is running", || handle.is_running()).await;
2717
2718        for _ in 0..CHILDREN {
2719            // A generous deadline: this test is about every child starting, not about how fast a couple of thousand
2720            // of them can be reaped, which is slow enough in a debug build to trip a short per-child timeout.
2721            handle.spawn(MockWorker::long_running("burst").with_graceful_timeout(Duration::from_secs(30)));
2722        }
2723
2724        wait_until("every child in the burst has started", || {
2725            handle.active_children() == CHILDREN
2726        })
2727        .await;
2728
2729        tx.send(()).unwrap();
2730        let result = timeout(Duration::from_secs(30), run)
2731            .await
2732            .expect("supervisor should stop")
2733            .expect("supervisor task should not panic");
2734        assert!(result.is_ok(), "the burst should have drained cleanly: {result:?}");
2735    }
2736
2737    #[tokio::test]
2738    async fn dynamically_spawned_supervisor_runs_and_drains() {
2739        // A dynamic child can be a whole supervision subtree, not just a worker: spawning a `Supervisor` runs it
2740        // nested, and shutting the parent down drains it along with everything under it.
2741        let child_worker = MockWorker::long_running("nested-child");
2742        let child_started = child_worker.start_count();
2743        let mut nested = Supervisor::new("nested-sup").unwrap();
2744        nested.add_worker(child_worker);
2745
2746        let sup = Supervisor::new("dyn-sup").unwrap();
2747        let handle = sup.handle();
2748        let (tx, run) = run_supervisor_with_trigger(sup).await;
2749        wait_until("supervisor is running", || handle.is_running()).await;
2750
2751        handle.spawn(nested);
2752        wait_until("the nested supervisor's own child has started", || {
2753            child_started.load(Ordering::SeqCst) == 1
2754        })
2755        .await;
2756
2757        tx.send(()).unwrap();
2758        let result = join_supervisor(run).await;
2759        assert!(
2760            result.is_ok(),
2761            "the nested subtree should have drained cleanly: {result:?}"
2762        );
2763    }
2764
2765    /// Counts how many times a nested subtree is started, by counting starts of its sole child.
2766    ///
2767    /// The child fails immediately and the subtree has a restart intensity of zero, so the subtree gives up the first
2768    /// time it fails. That makes the child's start count equal to the number of times the *subtree* ran, which is what
2769    /// these tests are actually asserting on -- without the zero intensity the subtree's own one-for-one restart would
2770    /// be indistinguishable from the parent restarting the subtree.
2771    fn failing_subtree(name: &'static str) -> (Supervisor, Arc<AtomicUsize>) {
2772        let worker = MockWorker::failing("nested-child", Duration::from_millis(5));
2773        let started = worker.start_count();
2774        let mut nested = Supervisor::new(name)
2775            .unwrap()
2776            .with_restart_strategy(RestartStrategy::new(RestartMode::OneForOne, 0, Duration::from_secs(30)));
2777        nested.add_worker(worker);
2778
2779        (nested, started)
2780    }
2781
2782    #[tokio::test]
2783    async fn dynamic_nested_supervisor_defaults_to_temporary() {
2784        // Spawning a bare `Supervisor` takes the dynamic default, so a subtree that gives up stays gone. For a
2785        // listener that means it silently disappears while whatever owns it keeps reporting healthy -- which is the
2786        // reason `nested_supervisor` exists.
2787        let (nested, started) = failing_subtree("nested-temp");
2788
2789        let sup = Supervisor::new("dyn-temp-sup").unwrap();
2790        let handle = sup.handle();
2791        let (tx, run) = run_supervisor_with_trigger(sup).await;
2792
2793        handle.spawn(nested);
2794        wait_until("the subtree has started once", || started.load(Ordering::SeqCst) == 1).await;
2795
2796        // Give the subtree time to fail and be reaped. Nothing brings it back.
2797        sleep(Duration::from_millis(200)).await;
2798        assert_eq!(
2799            started.load(Ordering::SeqCst),
2800            1,
2801            "a temporary subtree must not be restarted"
2802        );
2803
2804        tx.send(()).unwrap();
2805        assert!(join_supervisor(run).await.is_ok());
2806    }
2807
2808    #[tokio::test]
2809    async fn dynamic_nested_supervisor_can_be_made_permanent() {
2810        // What `nested_supervisor` buys: the same subtree, spawned permanent, is brought back when it terminates.
2811        let (nested, started) = failing_subtree("nested-perm");
2812
2813        let sup = Supervisor::new("dyn-perm-sup")
2814            .unwrap()
2815            .with_restart_strategy(RestartStrategy::new(
2816                RestartMode::OneForOne,
2817                100,
2818                Duration::from_secs(30),
2819            ));
2820        let handle = sup.handle();
2821        let (tx, run) = run_supervisor_with_trigger(sup).await;
2822
2823        handle.nested_supervisor(nested).spawn();
2824        wait_until("the subtree has been restarted", || started.load(Ordering::SeqCst) >= 2).await;
2825
2826        tx.send(()).unwrap();
2827        assert!(join_supervisor(run).await.is_ok());
2828    }
2829
2830    #[tokio::test]
2831    async fn dynamic_nested_supervisor_can_be_significant() {
2832        // The other half: a subtree its parent can't function without takes the parent with it when it terminates,
2833        // rather than leaving it running with nothing behind it.
2834        let (nested, _started) = failing_subtree("nested-sig");
2835
2836        let sup = Supervisor::new("dyn-sig-sup")
2837            .unwrap()
2838            .with_auto_shutdown(AutoShutdown::AnySignificant);
2839        let handle = sup.handle();
2840        let (tx, run) = run_supervisor_with_trigger(sup).await;
2841
2842        handle
2843            .nested_supervisor(nested)
2844            .temporary()
2845            .with_significant(true)
2846            .spawn();
2847
2848        let result = timeout(Duration::from_secs(5), run)
2849            .await
2850            .expect("supervisor should stop once the significant subtree terminates")
2851            .expect("supervisor task should not panic");
2852        assert!(
2853            matches!(result, Err(SupervisorError::SignificantChildExited)),
2854            "the significant subtree's termination should have stopped the parent, got {result:?}"
2855        );
2856
2857        // The run already ended; the trigger is redundant but keeps the sender alive to the end of the test.
2858        let _ = tx.send(());
2859    }
2860
2861    #[tokio::test]
2862    async fn budget_of_duration_max_does_not_leave_a_budget_bounded_child_unbounded() {
2863        // `Duration::MAX` is the natural spelling of "no ceiling", and a budget too large to become a deadline bounds
2864        // nothing at all. A budget-bounded child under one must therefore fall back to its own deadline: without that,
2865        // setting `MAX` would be strictly *worse* than setting no budget, since the child would never be abandoned and
2866        // the drain would hang.
2867        let mut sup = Supervisor::new("test-sup").unwrap().with_shutdown_budget(Duration::MAX);
2868        sup.add_worker(
2869            ChildSpecification::one_shot_worker(
2870                MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_millis(100)),
2871            )
2872            .with_budget_bounded_shutdown(),
2873        );
2874
2875        let (tx, run) = run_supervisor_with_trigger(sup).await;
2876        tx.send(()).unwrap();
2877
2878        let started = tokio::time::Instant::now();
2879        let result = join_supervisor(run).await;
2880        let elapsed = started.elapsed();
2881
2882        assert!(
2883            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
2884            "the worker's own deadline should have aborted it, got {result:?}"
2885        );
2886        assert!(
2887            elapsed < Duration::from_secs(1),
2888            "the child should have been bounded by its own 100ms deadline; took {elapsed:?}"
2889        );
2890    }
2891
2892    #[tokio::test]
2893    async fn budget_bounded_child_falls_back_to_its_own_deadline_without_a_budget() {
2894        // A one-shot child asks to be bounded by its supervisor's budget rather than carrying a deadline of its own.
2895        // On a supervisor with no budget there'd be nothing bounding it at all, so it falls back to the strategy the
2896        // worker reports -- here a short one, which is what lets this test finish rather than hang.
2897        let mut sup = Supervisor::new("test-sup").unwrap();
2898        sup.add_worker(
2899            ChildSpecification::one_shot_worker(
2900                MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_millis(100)),
2901            )
2902            .with_budget_bounded_shutdown(),
2903        );
2904
2905        let (tx, run) = run_supervisor_with_trigger(sup).await;
2906        tx.send(()).unwrap();
2907
2908        let started = tokio::time::Instant::now();
2909        let result = join_supervisor(run).await;
2910        let elapsed = started.elapsed();
2911
2912        assert!(
2913            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
2914            "the worker's own deadline should have aborted it, got {result:?}"
2915        );
2916        assert!(
2917            elapsed < Duration::from_secs(1),
2918            "the child should have been bounded by its own 100ms deadline; took {elapsed:?}"
2919        );
2920    }
2921
2922    #[tokio::test]
2923    async fn concurrent_shutdown_drains_many_children_quickly() {
2924        const CHILDREN: usize = 500;
2925        const SHUTDOWN_DELAY: Duration = Duration::from_millis(50);
2926
2927        let sup = Supervisor::new("dyn-sup").unwrap();
2928        let handle = sup.handle();
2929        let (tx, run) = run_supervisor_with_trigger(sup).await;
2930        wait_until("supervisor is running", || handle.is_running()).await;
2931
2932        for _ in 0..CHILDREN {
2933            handle.spawn(MockWorker::slow_shutdown("conn", SHUTDOWN_DELAY));
2934        }
2935        wait_until("all dynamic children are running", || {
2936            handle.active_children() == CHILDREN
2937        })
2938        .await;
2939
2940        // Each child sleeps after observing shutdown. Concurrent shutdown drains them all in roughly one delay; an
2941        // ordered shutdown would take CHILDREN * delay (25s here). Assert it finishes well under that.
2942        let start = std::time::Instant::now();
2943        tx.send(()).unwrap();
2944        let result = timeout(Duration::from_secs(5), run).await.unwrap().unwrap();
2945        let elapsed = start.elapsed();
2946
2947        assert!(result.is_ok());
2948        assert_eq!(handle.active_children(), 0, "active count must return to zero");
2949        assert!(
2950            elapsed < Duration::from_secs(2),
2951            "shutdown must be concurrent (took {elapsed:?})"
2952        );
2953    }
2954
2955    #[tokio::test]
2956    async fn concurrent_shutdown_aborts_unresponsive_children() {
2957        let sup = Supervisor::new("dyn-sup").unwrap();
2958        let handle = sup.handle();
2959        let (tx, run) = run_supervisor_with_trigger(sup).await;
2960        wait_until("supervisor is running", || handle.is_running()).await;
2961
2962        handle.spawn(MockWorker::ignore_shutdown("stuck"));
2963        wait_until("one dynamic child is running", || handle.active_children() == 1).await;
2964
2965        // The child never reacts to shutdown, so it must be aborted once its graceful deadline (500ms) elapses rather
2966        // than hanging the supervisor.
2967        let start = std::time::Instant::now();
2968        tx.send(()).unwrap();
2969        let result = join_supervisor(run).await;
2970        let elapsed = start.elapsed();
2971
2972        // Forcefully aborting an unresponsive child is surfaced as an unclean shutdown rather than reported as success.
2973        assert!(
2974            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
2975            "aborting a stuck child must surface as an unclean shutdown, got {result:?}"
2976        );
2977        assert_eq!(handle.active_children(), 0);
2978        assert!(
2979            elapsed < Duration::from_secs(1),
2980            "stuck child must be aborted at the deadline (took {elapsed:?})"
2981        );
2982    }
2983
2984    #[tokio::test]
2985    async fn concurrent_shutdown_honors_per_child_deadline() {
2986        // Each child must be aborted at its OWN graceful deadline, not a single shared one. A responsive child with an
2987        // effectively-infinite timeout (modeling a nested supervisor, which uses `Graceful(Duration::MAX)`) coexists
2988        // with an unresponsive child with a short timeout. Under a shared `max` deadline the short-timeout child would
2989        // never be aborted (the shared deadline would be `MAX`) and shutdown would hang.
2990        let sup = Supervisor::new("dyn-sup").unwrap();
2991        let handle = sup.handle();
2992        let (tx, run) = run_supervisor_with_trigger(sup).await;
2993        wait_until("supervisor is running", || handle.is_running()).await;
2994
2995        // Responds to shutdown promptly, but its deadline is effectively infinite.
2996        handle.spawn(MockWorker::long_running("responsive").with_graceful_timeout(Duration::MAX));
2997        // Never responds; must be aborted at its own short deadline.
2998        handle.spawn(MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_millis(200)));
2999        wait_until("both dynamic children are running", || handle.active_children() == 2).await;
3000
3001        let start = std::time::Instant::now();
3002        tx.send(()).unwrap();
3003        let result = join_supervisor(run).await;
3004        let elapsed = start.elapsed();
3005
3006        // Only the stuck child is aborted (the responsive one exits cleanly), so the unclean-shutdown tally is exactly 1.
3007        assert!(
3008            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3009            "aborting the stuck child must surface as an unclean shutdown with a count of 1, got {result:?}"
3010        );
3011        assert_eq!(handle.active_children(), 0);
3012        assert!(
3013            elapsed < Duration::from_secs(1),
3014            "stuck child must be aborted at its own deadline despite an infinite-timeout sibling (took {elapsed:?})"
3015        );
3016    }
3017
3018    #[tokio::test]
3019    async fn unresponsive_child_is_aborted_at_its_deadline() {
3020        // A child that never reacts to shutdown must be aborted once its graceful deadline (500ms) elapses, rather
3021        // than hanging the supervisor indefinitely.
3022        let mut sup = Supervisor::new("test-sup").unwrap();
3023        sup.add_worker(MockWorker::ignore_shutdown("stuck"));
3024
3025        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3026
3027        let start = std::time::Instant::now();
3028        tx.send(()).unwrap();
3029        let result = join_supervisor(handle).await;
3030        let elapsed = start.elapsed();
3031
3032        assert!(
3033            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3034            "aborting a stuck child must surface as an unclean shutdown, got {result:?}"
3035        );
3036        assert!(
3037            elapsed < Duration::from_secs(1),
3038            "unresponsive child must be aborted at its deadline (took {elapsed:?})"
3039        );
3040    }
3041
3042    #[tokio::test]
3043    async fn brutal_shutdown_aborts_child_immediately() {
3044        // A child with a `Brutal` shutdown strategy is aborted at once on shutdown, with no graceful wait -- so even a
3045        // child that ignores shutdown is torn down promptly rather than after the graceful deadline.
3046        let mut sup = Supervisor::new("test-sup").unwrap();
3047        sup.add_worker(MockWorker::ignore_shutdown("brutal-stuck").with_brutal_shutdown());
3048
3049        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3050
3051        let start = std::time::Instant::now();
3052        tx.send(()).unwrap();
3053        let result = join_supervisor(handle).await;
3054        let elapsed = start.elapsed();
3055
3056        // A brutal abort is the configured, expected way to stop this child -- not a graceful-timeout overrun -- so it
3057        // is NOT counted toward the unclean-shutdown tally, and the shutdown reports success.
3058        assert!(result.is_ok());
3059        assert!(
3060            elapsed < Duration::from_millis(200),
3061            "brutal-shutdown child must be aborted immediately, not after a graceful wait (took {elapsed:?})"
3062        );
3063    }
3064
3065    #[tokio::test]
3066    async fn shutdown_timeout_aborts_aggregate_to_root() {
3067        // Forced aborts must surface as an unclean shutdown and aggregate up the tree: a supervisor adds the workers it
3068        // aborts directly to the counts reported by any child supervisors that also timed out. Here the parent aborts
3069        // one direct child and a nested supervisor aborts one of its own, so the root observes a total of 2.
3070        let mut child_sup = Supervisor::new("child-sup").unwrap();
3071        child_sup
3072            .add_worker(MockWorker::ignore_shutdown("child-stuck").with_graceful_timeout(Duration::from_millis(200)));
3073
3074        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3075        parent_sup
3076            .add_worker(MockWorker::ignore_shutdown("parent-stuck").with_graceful_timeout(Duration::from_millis(200)));
3077        parent_sup.add_worker(MockWorker::long_running("parent-clean"));
3078        parent_sup.add_worker(child_sup);
3079
3080        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3081        tx.send(()).unwrap();
3082
3083        let result = join_supervisor(handle).await;
3084        assert!(
3085            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 2 })),
3086            "forced aborts must aggregate across the tree (1 direct + 1 nested), got {result:?}"
3087        );
3088    }
3089
3090    // -- Restart-policy edge cases ---------------------------------------------------------
3091
3092    #[tokio::test]
3093    async fn restart_intensity_zero_shuts_down_on_first_failure() {
3094        // A restart intensity of zero means the supervisor gives up the moment any restartable child fails: it shuts
3095        // down on the very first failure without ever restarting the worker. (See `RestartState::evaluate_restart`,
3096        // which short-circuits to `Shutdown` when intensity is zero.)
3097        let worker = MockWorker::failing("boom", Duration::from_millis(20));
3098        let start_count = worker.start_count();
3099
3100        let mut sup = Supervisor::new("test-sup")
3101            .unwrap()
3102            .with_restart_strategy(RestartStrategy::new(RestartMode::OneForOne, 0, Duration::from_secs(5)));
3103        sup.add_worker(worker);
3104
3105        let (_tx, rx) = oneshot::channel::<()>();
3106        let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
3107            .await
3108            .unwrap();
3109
3110        assert!(matches!(result, Err(SupervisorError::Shutdown)));
3111        assert_eq!(
3112            start_count.load(Ordering::SeqCst),
3113            1,
3114            "with intensity zero the worker must run exactly once and never be restarted"
3115        );
3116    }
3117
3118    #[tokio::test]
3119    async fn one_for_all_restart_loses_dynamic_children() {
3120        // Documented one-for-all semantics: a group restart resets to the static roster only -- dynamic children are
3121        // NOT restored (they're lost on a supervisor-level restart, matching Erlang/OTP). A permanent static worker
3122        // that keeps failing drives repeated one-for-all restarts; a dynamic child spawned before the first restart
3123        // must be torn down and never brought back.
3124        let failing = MockWorker::failing("failing-static", Duration::from_millis(50));
3125        let failing_count = failing.start_count();
3126
3127        let sup = Supervisor::new("dyn-sup").unwrap().with_restart_strategy(
3128            RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
3129        );
3130        let handle = sup.handle();
3131        let mut sup = sup;
3132        sup.add_worker(failing);
3133
3134        let (tx, run) = run_supervisor_with_trigger(sup).await;
3135
3136        // Spawn a long-running dynamic child and wait for it to be running.
3137        let dynamic = MockWorker::long_running("dynamic");
3138        let dynamic_count = dynamic.start_count();
3139        handle.spawn(dynamic);
3140        wait_until("the dynamic child is running", || handle.active_children() == 1).await;
3141
3142        // Let the static worker drive at least one one-for-all restart (its second start).
3143        wait_until("the static worker has been restarted", || {
3144            failing_count.load(Ordering::SeqCst) >= 2
3145        })
3146        .await;
3147
3148        // The one-for-all restart must have discarded the dynamic child: the active count returns to zero, and the
3149        // dynamic child ran exactly once (it was never restored).
3150        wait_until("the dynamic child has been discarded", || handle.active_children() == 0).await;
3151        assert_eq!(
3152            dynamic_count.load(Ordering::SeqCst),
3153            1,
3154            "a dynamic child must be lost -- not restored -- across a one-for-all restart"
3155        );
3156
3157        tx.send(()).unwrap();
3158        let result = join_supervisor(run).await;
3159        assert!(result.is_ok());
3160    }
3161
3162    // -- Dedicated-runtime tests -----------------------------------------------------------
3163
3164    #[tokio::test]
3165    async fn dedicated_single_threaded_runtime_runs_nested_worker_and_shuts_down_cleanly() {
3166        // A nested supervisor configured with a dedicated single-threaded runtime spawns its own OS thread and Tokio
3167        // runtime (via `spawn_dedicated_runtime`). Its worker must run there, and a shutdown signalled by the parent
3168        // must propagate across the thread boundary and tear it down cleanly.
3169        let worker = MockWorker::long_running("dedicated-worker");
3170        let worker_count = worker.start_count();
3171
3172        let mut child_sup = Supervisor::new("child-sup")
3173            .unwrap()
3174            .with_dedicated_runtime(RuntimeConfiguration::single_threaded());
3175        child_sup.add_worker(worker);
3176
3177        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3178        parent_sup.add_worker(child_sup);
3179
3180        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3181
3182        // The worker starts on the dedicated runtime's own thread.
3183        wait_until("the dedicated worker has started", || {
3184            worker_count.load(Ordering::SeqCst) == 1
3185        })
3186        .await;
3187
3188        tx.send(()).unwrap();
3189        let result = join_supervisor(handle).await;
3190        assert!(
3191            result.is_ok(),
3192            "dedicated-runtime supervisor should shut down cleanly, got {result:?}"
3193        );
3194    }
3195
3196    #[tokio::test]
3197    async fn dedicated_multi_threaded_runtime_runs_nested_worker() {
3198        // The same nested-dedicated flow, but exercising the multi-threaded dedicated runtime builder path.
3199        let worker = MockWorker::long_running("dedicated-worker");
3200        let worker_count = worker.start_count();
3201
3202        let mut child_sup = Supervisor::new("child-sup")
3203            .unwrap()
3204            .with_dedicated_runtime(RuntimeConfiguration::multi_threaded(2));
3205        child_sup.add_worker(worker);
3206
3207        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3208        parent_sup.add_worker(child_sup);
3209
3210        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3211        wait_until("the dedicated worker has started", || {
3212            worker_count.load(Ordering::SeqCst) == 1
3213        })
3214        .await;
3215
3216        tx.send(()).unwrap();
3217        let result = join_supervisor(handle).await;
3218        assert!(
3219            result.is_ok(),
3220            "multi-threaded dedicated-runtime supervisor should shut down cleanly, got {result:?}"
3221        );
3222    }
3223
3224    #[tokio::test]
3225    async fn dedicated_runtime_forced_abort_aggregates_to_root() {
3226        // A worker inside a dedicated-runtime nested supervisor that ignores shutdown must be forcefully aborted at its
3227        // deadline, and that abort tally must survive the OS-thread boundary (`DedicatedRuntimeHandle` -> `WorkerError`)
3228        // and be observed by the root supervisor as `ShutdownTimedOut`.
3229        let stuck = MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_millis(200));
3230        let stuck_count = stuck.start_count();
3231
3232        let mut child_sup = Supervisor::new("child-sup")
3233            .unwrap()
3234            .with_dedicated_runtime(RuntimeConfiguration::single_threaded());
3235        child_sup.add_worker(stuck);
3236
3237        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3238        parent_sup.add_worker(child_sup);
3239
3240        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3241
3242        // Make sure the stuck worker is actually running on the dedicated runtime before signalling shutdown, so the
3243        // forced-abort path (rather than an early exit) is what we exercise.
3244        wait_until("the stuck worker has started", || {
3245            stuck_count.load(Ordering::SeqCst) == 1
3246        })
3247        .await;
3248
3249        tx.send(()).unwrap();
3250        let result = join_supervisor(handle).await;
3251        assert!(
3252            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3253            "a stuck worker in a dedicated runtime must surface as an unclean shutdown aggregated to the root, got {result:?}"
3254        );
3255    }
3256
3257    // -- Per-child override tests ----------------------------------------------------------
3258
3259    #[tokio::test]
3260    async fn child_with_runtime_override_runs_on_that_runtime() {
3261        // `with_runtime` places an individual child's task on a caller-provided runtime instead of the supervisor's
3262        // own. The child reports the name of the thread it's actually running on, which must belong to that runtime.
3263        let child_runtime = tokio::runtime::Builder::new_multi_thread()
3264            .worker_threads(1)
3265            .thread_name("child-rt-test")
3266            .enable_all()
3267            .build()
3268            .expect("should build child runtime");
3269
3270        let (thread_tx, thread_rx) = oneshot::channel();
3271        let worker = FnWorker::new("placed", async move {
3272            let thread_name = std::thread::current().name().unwrap_or_default().to_string();
3273            let _ = thread_tx.send(thread_name);
3274            pending::<()>().await;
3275        });
3276
3277        // The worker only has to stay alive long enough to report where it ran, so it's aborted at shutdown rather
3278        // than given a terminal condition to reach.
3279        let mut sup = Supervisor::new("test-sup").unwrap();
3280        sup.add_worker(
3281            ChildSpecification::worker(worker)
3282                .with_restart_type(RestartType::Temporary)
3283                .with_runtime(child_runtime.handle().clone())
3284                .with_shutdown_strategy(ShutdownStrategy::Brutal),
3285        );
3286
3287        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3288
3289        let thread_name = timeout(Duration::from_secs(2), thread_rx)
3290            .await
3291            .expect("child should report its thread promptly")
3292            .expect("child should not be dropped before reporting");
3293        assert!(
3294            thread_name.starts_with("child-rt-test"),
3295            "child must run on the runtime given to `with_runtime`, but ran on thread {thread_name:?}"
3296        );
3297
3298        tx.send(()).unwrap();
3299        assert!(join_supervisor(handle).await.is_ok());
3300
3301        // Dropping a `Runtime` from within an async context panics, so tear it down without blocking.
3302        child_runtime.shutdown_background();
3303    }
3304
3305    #[tokio::test]
3306    async fn child_shutdown_strategy_override_takes_precedence_over_worker() {
3307        // `with_shutdown_strategy` overrides what the worker reports for itself. The worker below asks for a 30-second
3308        // grace period and then ignores shutdown entirely; the override cuts that to 50ms, so the supervisor must
3309        // abort it and report an unclean shutdown well inside `join_supervisor`'s two-second bound. Without the
3310        // override taking precedence, this test times out.
3311        let worker = MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_secs(30));
3312
3313        let mut sup = Supervisor::new("test-sup").unwrap();
3314        sup.add_worker(
3315            ChildSpecification::worker(worker)
3316                .with_restart_type(RestartType::Temporary)
3317                .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::from_millis(50))),
3318        );
3319
3320        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3321        tx.send(()).unwrap();
3322
3323        let result = join_supervisor(handle).await;
3324        assert!(
3325            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3326            "the overridden 50ms deadline should have aborted the stuck child, got {result:?}"
3327        );
3328    }
3329
3330    /// A worker that waits for shutdown and then drains a [`ShutdownCoordinator`] before exiting.
3331    ///
3332    /// Stands in for a component that owns background work and waits for it during its own shutdown. Unlike a
3333    /// closure-based worker it genuinely needs the shutdown signal, which is exactly the case [`Supervisable`] exists
3334    /// for.
3335    struct DrainWaiter {
3336        coordinator: Mutex<Option<ShutdownCoordinator>>,
3337        finished: Arc<AtomicBool>,
3338    }
3339
3340    #[async_trait]
3341    impl Supervisable for DrainWaiter {
3342        fn name(&self) -> &str {
3343            "waiter"
3344        }
3345
3346        async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
3347            let coordinator = self
3348                .coordinator
3349                .lock()
3350                .expect("drain waiter mutex poisoned")
3351                .take()
3352                .expect("drain waiter runs once");
3353            let finished = Arc::clone(&self.finished);
3354
3355            Ok(Box::pin(async move {
3356                process_shutdown.await;
3357                coordinator.shutdown_and_wait().await;
3358                finished.store(true, Ordering::SeqCst);
3359                Ok(())
3360            }))
3361        }
3362    }
3363
3364    /// Builds two children modelling an owner that drains a background task during shutdown.
3365    ///
3366    /// `stuck` holds a shutdown handle and never releases it voluntarily, so the only way it goes away is a forced
3367    /// abort. `waiter` blocks on that handle being dropped, standing in for a component's `shutdown_and_wait`. The
3368    /// returned flag records whether `waiter` ran to completion rather than being aborted itself.
3369    fn build_drain_pair(
3370        stuck_strategy: ShutdownStrategy, waiter_strategy: ShutdownStrategy,
3371    ) -> (Supervisor, Arc<AtomicBool>) {
3372        let mut coordinator = ShutdownCoordinator::default();
3373        let held_handle = coordinator.register();
3374
3375        let stuck = FnWorker::new("stuck", async move {
3376            // Hold the handle for as long as this future lives, and ignore shutdown entirely.
3377            let _held = held_handle;
3378            pending::<()>().await;
3379        });
3380
3381        let waiter_finished = Arc::new(AtomicBool::new(false));
3382        let waiter = DrainWaiter {
3383            coordinator: Mutex::new(Some(coordinator)),
3384            finished: Arc::clone(&waiter_finished),
3385        };
3386
3387        let mut sup = Supervisor::new("test-sup").unwrap();
3388        sup.add_worker(
3389            ChildSpecification::worker(stuck)
3390                .with_restart_type(RestartType::Temporary)
3391                .with_shutdown_strategy(stuck_strategy),
3392        );
3393        sup.add_worker(
3394            ChildSpecification::worker(waiter)
3395                .with_restart_type(RestartType::Temporary)
3396                .with_shutdown_strategy(waiter_strategy),
3397        );
3398
3399        (sup, waiter_finished)
3400    }
3401
3402    #[tokio::test]
3403    async fn shorter_child_deadline_releases_a_waiting_sibling() {
3404        // Aborting a stuck child drops the shutdown handle it was holding, which is what releases anything waiting on
3405        // it. A child bounded more tightly than its waiter therefore stays recoverable: the child is aborted, the
3406        // waiter unblocks and finishes cleanly, and only one abort is reported.
3407        let (sup, waiter_finished) = build_drain_pair(
3408            ShutdownStrategy::Graceful(Duration::from_millis(100)),
3409            ShutdownStrategy::Graceful(Duration::from_secs(1)),
3410        );
3411
3412        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3413        tx.send(()).unwrap();
3414
3415        let result = join_supervisor(handle).await;
3416        assert!(
3417            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3418            "only the stuck child should have been aborted, got {result:?}"
3419        );
3420        assert!(
3421            waiter_finished.load(Ordering::SeqCst),
3422            "the waiter should have been released by the stuck child's abort and run to completion"
3423        );
3424    }
3425
3426    #[tokio::test]
3427    async fn equal_child_deadlines_abort_the_waiter_too() {
3428        // The counterpart: concurrent shutdown computes every deadline from one shared instant, so identical timeouts
3429        // elapse in the same pass and the waiter is aborted alongside the child it was waiting on. This is also what a
3430        // shutdown budget does to a whole subtree, which is why the budget is set at a level where losing the entire
3431        // group at once is the intended outcome.
3432        let (sup, waiter_finished) = build_drain_pair(
3433            ShutdownStrategy::Graceful(Duration::from_millis(100)),
3434            ShutdownStrategy::Graceful(Duration::from_millis(100)),
3435        );
3436
3437        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3438        tx.send(()).unwrap();
3439
3440        let result = join_supervisor(handle).await;
3441        assert!(
3442            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 2 })),
3443            "both children should have been aborted together, got {result:?}"
3444        );
3445        assert!(
3446            !waiter_finished.load(Ordering::SeqCst),
3447            "the waiter should have been aborted mid-wait, not completed"
3448        );
3449    }
3450
3451    // -- Shutdown budget tests -------------------------------------------------------------
3452
3453    #[tokio::test]
3454    async fn budget_bounds_children_that_have_no_deadline_of_their_own() {
3455        // Without a budget, `Graceful(Duration::MAX)` children are waited on indefinitely and a stuck one hangs
3456        // shutdown forever. A budget makes the supervisor responsible for the deadline instead, and each child it has
3457        // to abort is still counted individually -- so an overrun says how many tasks were responsible, not merely
3458        // that the group as a whole overran.
3459        let mut sup = Supervisor::new("test-sup")
3460            .unwrap()
3461            .with_shutdown_budget(Duration::from_millis(100));
3462
3463        for name in ["stuck_one", "stuck_two"] {
3464            sup.add_worker(
3465                ChildSpecification::worker(FnWorker::new(name, pending::<()>()))
3466                    .with_restart_type(RestartType::Temporary)
3467                    .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::MAX)),
3468            );
3469        }
3470
3471        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3472        tx.send(()).unwrap();
3473
3474        let result = join_supervisor(handle).await;
3475        assert!(
3476            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 2 })),
3477            "the budget should have aborted both deadline-less children, got {result:?}"
3478        );
3479    }
3480
3481    #[tokio::test]
3482    async fn budget_does_not_delay_children_that_stop_on_their_own() {
3483        // The budget is a ceiling, not a wait. Asserting on elapsed time rather than merely on success is what makes
3484        // this meaningful: a supervisor that waited out its budget regardless would still report `Ok`.
3485        let mut sup = Supervisor::new("test-sup")
3486            .unwrap()
3487            .with_shutdown_budget(Duration::from_secs(30));
3488        sup.add_worker(
3489            ChildSpecification::one_shot_worker(MockWorker::long_running("prompt"))
3490                .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::MAX)),
3491        );
3492
3493        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3494        let started = tokio::time::Instant::now();
3495        tx.send(()).unwrap();
3496
3497        assert!(join_supervisor(handle).await.is_ok());
3498        let elapsed = started.elapsed();
3499        assert!(
3500            elapsed < Duration::from_millis(500),
3501            "shutdown should finish as soon as the child does, not burn the budget; took {elapsed:?}"
3502        );
3503    }
3504
3505    #[tokio::test]
3506    async fn child_deadline_shorter_than_budget_still_wins() {
3507        // A child that carries its own finite deadline is held to whichever elapses first, so a component can still
3508        // bound one particular task more tightly than the budget covering the rest.
3509        let mut sup = Supervisor::new("test-sup")
3510            .unwrap()
3511            .with_shutdown_budget(Duration::from_secs(30));
3512        sup.add_worker(
3513            ChildSpecification::worker(FnWorker::new("stuck", pending::<()>()))
3514                .with_restart_type(RestartType::Temporary)
3515                .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::from_millis(100))),
3516        );
3517
3518        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3519        tx.send(()).unwrap();
3520
3521        // Again bounded at two seconds: if the 30-second budget had won, this would time out instead.
3522        let result = join_supervisor(handle).await;
3523        assert!(
3524            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3525            "the child's own 100ms deadline should have won over the budget, got {result:?}"
3526        );
3527    }
3528
3529    #[tokio::test]
3530    async fn every_worker_records_poll_metrics() {
3531        // Poll timing is a property of being supervised, not something a child opts into, so a plain statically
3532        // registered worker gets it too -- tagged with its fully qualified process name.
3533        let recorder = TestRecorder::default();
3534        let _guard = metrics::set_default_local_recorder(&recorder);
3535
3536        // The recorder has to be installed before the worker spawns: its metric handles are resolved once, at spawn.
3537        let mut sup = Supervisor::new("metrics_sup").unwrap();
3538        sup.add_worker(ChildSpecification::one_shot_worker(MockWorker::long_running("timed")));
3539
3540        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3541        tx.send(()).unwrap();
3542        assert!(join_supervisor(handle).await.is_ok());
3543
3544        let polls = recorder.counter(("runtime_task_poll_count", &[("task_name", "metrics_sup.timed")]));
3545        assert!(
3546            polls.is_some_and(|polls| polls > 0),
3547            "a supervised worker should have recorded poll metrics, got {polls:?}"
3548        );
3549    }
3550
3551    #[tokio::test]
3552    async fn budget_bounds_the_whole_drain_rather_than_each_child() {
3553        // The budget is measured once, from the start of the drain, and every child is held to that same instant --
3554        // it does not reset per child. Three children that each ignore a 10-second deadline must all be aborted at
3555        // the shared 150ms budget, not 30 seconds later.
3556        let mut sup = Supervisor::new("test-sup")
3557            .unwrap()
3558            .with_shutdown_budget(Duration::from_millis(150));
3559
3560        for name in ["stuck_one", "stuck_two", "stuck_three"] {
3561            sup.add_worker(
3562                ChildSpecification::one_shot_worker(FnWorker::new(name, pending::<()>()))
3563                    .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::from_secs(10))),
3564            );
3565        }
3566
3567        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3568        tx.send(()).unwrap();
3569
3570        let result = join_supervisor(handle).await;
3571        assert!(
3572            matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 3 })),
3573            "the budget should have bounded the whole drain, got {result:?}"
3574        );
3575    }
3576
3577    #[tokio::test]
3578    async fn budget_of_duration_max_is_treated_as_no_budget() {
3579        // `Duration::MAX` is the natural spelling of "no ceiling" and used to panic the supervisor task on an instant
3580        // overflow.
3581        let mut sup = Supervisor::new("test-sup").unwrap().with_shutdown_budget(Duration::MAX);
3582        sup.add_worker(ChildSpecification::one_shot_worker(MockWorker::long_running("prompt")));
3583
3584        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3585        tx.send(()).unwrap();
3586        assert!(join_supervisor(handle).await.is_ok());
3587    }
3588
3589    #[tokio::test]
3590    async fn near_max_child_timeout_does_not_panic() {
3591        // Resolving a graceful timeout to an instant makes anything just under `Duration::MAX` overflow unless it is
3592        // added with `checked_add`. See `resolve_abort_deadline`.
3593        let mut sup = Supervisor::new("test-sup").unwrap();
3594        sup.add_worker(
3595            ChildSpecification::one_shot_worker(MockWorker::long_running("prompt"))
3596                .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::MAX - Duration::from_nanos(1))),
3597        );
3598
3599        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3600        tx.send(()).unwrap();
3601        assert!(join_supervisor(handle).await.is_ok());
3602    }
3603
3604    #[tokio::test]
3605    async fn budget_does_not_cut_off_a_nested_supervisor_mid_drain() {
3606        // A nested supervisor bounds its own subtree, so a parent's budget must not abort it: doing so truncates the
3607        // subtree's drain, discards its abort tally, and -- for a supervisor on a dedicated runtime, whose work is on
3608        // another OS thread -- reports it as stopped without actually stopping it.
3609        let slow = MockWorker::slow_shutdown("slow", Duration::from_millis(300));
3610        let drained = slow.finish_count();
3611
3612        let mut nested = Supervisor::new("nested").unwrap();
3613        nested.add_worker(
3614            ChildSpecification::one_shot_worker(slow)
3615                .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::from_secs(10))),
3616        );
3617
3618        let mut parent = Supervisor::new("parent")
3619            .unwrap()
3620            .with_shutdown_budget(Duration::from_millis(50));
3621        parent.add_worker(nested);
3622
3623        let (tx, handle) = run_supervisor_with_trigger(parent).await;
3624        tx.send(()).unwrap();
3625
3626        let result = join_supervisor(handle).await;
3627        assert!(
3628            drained.load(Ordering::SeqCst) == 1,
3629            "the nested subtree should have drained rather than being cut off by the parent's budget: {result:?}"
3630        );
3631        assert!(
3632            result.is_ok(),
3633            "the nested drain finished in time, so shutdown was clean: {result:?}"
3634        );
3635    }
3636    // -- Supervision tree snapshot tests ---------------------------------------------------
3637
3638    fn find_node<'a>(nodes: &'a [NodeSnapshot], name: &str) -> &'a NodeSnapshot {
3639        nodes.iter().find(|node| node.name == name).unwrap_or_else(|| {
3640            panic!(
3641                "no node named '{name}' among {:?}",
3642                nodes.iter().map(|n| &n.name).collect::<Vec<_>>()
3643            )
3644        })
3645    }
3646
3647    #[tokio::test]
3648    async fn snapshot_before_run_reports_a_registered_root() {
3649        let mut sup = Supervisor::new("test-sup")
3650            .unwrap()
3651            .with_restart_strategy(RestartStrategy::one_for_all().with_intensity_and_period(7, Duration::from_secs(11)))
3652            .with_auto_shutdown(AutoShutdown::AnySignificant)
3653            .with_shutdown_budget(Duration::from_secs(3));
3654        sup.add_worker(MockWorker::long_running("worker1"));
3655
3656        // A tree that hasn't started still has a declared shape, and configuration is known from the moment it is
3657        // set, so all of it should be readable without running anything.
3658        let snapshot = sup.tree_handle().snapshot();
3659
3660        assert_eq!(snapshot.root.name, "test-sup");
3661        assert_eq!(snapshot.root.kind, NodeKind::Supervisor);
3662        assert_eq!(snapshot.root.state, NodeState::Registered);
3663        assert_eq!(snapshot.root.process_id, None);
3664        assert_eq!(snapshot.root.process_name, None);
3665        assert_eq!(snapshot.root.started_at, None);
3666        assert_eq!(snapshot.root.uptime_ms, None);
3667        assert!(snapshot.root.children.is_empty(), "no children have been started yet");
3668
3669        let supervision = snapshot.root.supervision.expect("a supervisor reports its settings");
3670        assert_eq!(supervision.restart_mode, RestartMode::OneForAll);
3671        assert_eq!(supervision.restart_intensity, 7);
3672        assert_eq!(supervision.restart_period_ms, 11_000);
3673        assert_eq!(supervision.auto_shutdown, AutoShutdown::AnySignificant);
3674        assert_eq!(supervision.shutdown_budget_ms, Some(3_000));
3675        assert_eq!(supervision.dedicated_threads, None);
3676        assert_eq!(supervision.generation, 0);
3677
3678        assert_eq!(snapshot.totals.supervisors, 1);
3679        assert_eq!(snapshot.totals.workers, 0);
3680        assert_eq!(snapshot.totals.registered, 1);
3681        assert_eq!(snapshot.totals.max_depth, 1);
3682    }
3683
3684    #[tokio::test]
3685    async fn snapshot_lists_static_children_in_declaration_order() {
3686        let mut sup = Supervisor::new("test-sup").unwrap();
3687        sup.add_worker(MockWorker::long_running("worker1"));
3688        sup.add_worker(MockWorker::long_running("worker2"));
3689
3690        let tree = sup.tree_handle();
3691        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3692        wait_until("both children appear in the snapshot", || {
3693            tree.snapshot().root.children.len() == 2
3694        })
3695        .await;
3696
3697        let snapshot = tree.snapshot();
3698        assert_eq!(snapshot.root.state, NodeState::Running);
3699        assert!(snapshot.root.process_id.is_some());
3700        assert_eq!(snapshot.root.process_name.as_deref(), Some("test_sup"));
3701        assert_eq!(snapshot.root.supervision.expect("supervisor settings").generation, 1);
3702
3703        // Declaration order, not hash order: it is how the tree is written, so it is how it should read.
3704        let names: Vec<&str> = snapshot.root.children.iter().map(|c| c.name.as_str()).collect();
3705        assert_eq!(names, vec!["worker1", "worker2"]);
3706
3707        for (child, expected_process) in snapshot
3708            .root
3709            .children
3710            .iter()
3711            .zip(["test_sup.worker1", "test_sup.worker2"])
3712        {
3713            assert_eq!(child.kind, NodeKind::Worker);
3714            assert_eq!(child.state, NodeState::Running);
3715            assert_eq!(child.restart, RestartType::Permanent);
3716            assert_eq!(child.restart_count, 0);
3717            assert_eq!(child.process_name.as_deref(), Some(expected_process));
3718            assert!(child.process_id.is_some(), "a running child has a process");
3719            assert!(child.children.is_empty(), "a worker has no children");
3720            assert!(child.supervision.is_none(), "a worker has no supervision settings");
3721            assert!(
3722                child.created_at <= child.started_at.expect("a running child has a start time"),
3723                "a child cannot start before it is created"
3724            );
3725        }
3726
3727        let first = snapshot.root.children[0].process_id;
3728        let second = snapshot.root.children[1].process_id;
3729        assert_ne!(first, second, "each child runs as its own process");
3730
3731        assert_eq!(snapshot.totals.supervisors, 1);
3732        assert_eq!(snapshot.totals.workers, 2);
3733        assert_eq!(snapshot.totals.running, 3);
3734        assert_eq!(snapshot.totals.max_depth, 2);
3735
3736        tx.send(()).unwrap();
3737        let _ = join_supervisor(handle).await;
3738    }
3739
3740    #[tokio::test]
3741    async fn one_for_one_restart_increments_only_the_failing_child() {
3742        let mut sup = Supervisor::new("test-sup")
3743            .unwrap()
3744            .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(5, Duration::from_secs(5)));
3745        sup.add_worker(MockWorker::failing("flapper", Duration::from_millis(20)));
3746        sup.add_worker(MockWorker::long_running("stable"));
3747
3748        let tree = sup.tree_handle();
3749        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3750
3751        wait_until("the failing child has appeared", || {
3752            tree.snapshot().root.children.len() == 2
3753        })
3754        .await;
3755        let before = tree.snapshot();
3756        let created_before = find_node(&before.root.children, "flapper").created_at;
3757        let process_before = find_node(&before.root.children, "flapper").process_id;
3758
3759        wait_until("the failing child has been restarted", || {
3760            find_node(&tree.snapshot().root.children, "flapper").restart_count >= 1
3761        })
3762        .await;
3763
3764        let after = tree.snapshot();
3765        let flapper = find_node(&after.root.children, "flapper");
3766        let stable = find_node(&after.root.children, "stable");
3767
3768        assert_eq!(stable.restart_count, 0, "a one-for-one restart leaves siblings alone");
3769        assert_ne!(flapper.process_id, process_before, "a restart is a new process");
3770        assert_eq!(
3771            flapper.created_at, created_before,
3772            "creation time is a property of the child, not of the incarnation"
3773        );
3774
3775        // A single child flapping is distinguishable from the supervisor restarting its whole group.
3776        assert!(after.root.supervision.expect("supervisor settings").restarts_performed >= 1);
3777
3778        tx.send(()).unwrap();
3779        let _ = join_supervisor(handle).await;
3780    }
3781
3782    #[tokio::test]
3783    async fn one_for_all_restart_increments_every_restarted_child() {
3784        let mut sup = Supervisor::new("test-sup")
3785            .unwrap()
3786            .with_restart_strategy(RestartStrategy::one_for_all().with_intensity_and_period(5, Duration::from_secs(5)));
3787        sup.add_worker(MockWorker::failing("flapper", Duration::from_millis(20)));
3788        sup.add_worker(MockWorker::long_running("sibling"));
3789
3790        let tree = sup.tree_handle();
3791        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3792
3793        wait_until("both children have appeared", || {
3794            tree.snapshot().root.children.len() == 2
3795        })
3796        .await;
3797        let before = tree.snapshot();
3798        let sibling_created = find_node(&before.root.children, "sibling").created_at;
3799        let sibling_process = find_node(&before.root.children, "sibling").process_id;
3800
3801        // The whole group is brought back, so every child's count moves -- which is what keeps each child's count
3802        // consistent with the new process and start time in the same record.
3803        wait_until("the group has been restarted", || {
3804            let snapshot = tree.snapshot();
3805            snapshot.root.children.len() == 2 && snapshot.root.children.iter().all(|child| child.restart_count >= 1)
3806        })
3807        .await;
3808
3809        let after = tree.snapshot();
3810        let sibling = find_node(&after.root.children, "sibling");
3811        assert_ne!(sibling.process_id, sibling_process, "a group restart is a new process");
3812        assert_eq!(
3813            sibling.created_at, sibling_created,
3814            "creation time survives a group restart"
3815        );
3816
3817        tx.send(()).unwrap();
3818        let _ = join_supervisor(handle).await;
3819    }
3820
3821    #[tokio::test]
3822    async fn temporary_child_survives_a_group_restart_as_a_tombstone() {
3823        let mut sup = Supervisor::new("test-sup")
3824            .unwrap()
3825            .with_restart_strategy(RestartStrategy::one_for_all().with_intensity_and_period(5, Duration::from_secs(5)));
3826        sup.add_worker(
3827            runtime::supervisable(MockWorker::completing("one-shot", Duration::from_millis(10)))
3828                .temporary()
3829                .build(),
3830        );
3831        sup.add_worker(MockWorker::failing("flapper", Duration::from_millis(30)));
3832
3833        let tree = sup.tree_handle();
3834        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3835
3836        // A temporary child is not brought back by a group restart, but it was declared, so it stays visible as a
3837        // tombstone rather than vanishing -- "declared, ran, stopped for good" is the useful answer.
3838        wait_until("the temporary child has exited and the group has restarted", || {
3839            let snapshot = tree.snapshot();
3840            snapshot
3841                .root
3842                .children
3843                .iter()
3844                .any(|child| child.name == "one-shot" && child.state == NodeState::Exited)
3845                && snapshot
3846                    .root
3847                    .children
3848                    .iter()
3849                    .any(|child| child.name == "flapper" && child.restart_count >= 1)
3850        })
3851        .await;
3852
3853        let snapshot = tree.snapshot();
3854        let one_shot = find_node(&snapshot.root.children, "one-shot");
3855        assert_eq!(one_shot.state, NodeState::Exited);
3856        assert_eq!(one_shot.restart, RestartType::Temporary);
3857        assert_eq!(one_shot.restart_count, 0, "a temporary child is never brought back");
3858        assert!(one_shot.exited_at.is_some(), "a tombstone records when it exited");
3859        assert_eq!(one_shot.uptime_ms, None, "a node that isn't running has no uptime");
3860        assert!(snapshot.totals.exited >= 1);
3861
3862        tx.send(()).unwrap();
3863        let _ = join_supervisor(handle).await;
3864    }
3865
3866    #[tokio::test]
3867    async fn dynamic_children_leave_no_tombstones() {
3868        let mut sup = Supervisor::new("test-sup").unwrap();
3869        sup.add_worker(MockWorker::long_running("static-worker"));
3870
3871        let tree = sup.tree_handle();
3872        let sup_handle = sup.handle();
3873        let (tx, handle) = run_supervisor_with_trigger(sup).await;
3874        wait_until("the static child is running", || {
3875            tree.snapshot().root.children.len() == 1
3876        })
3877        .await;
3878
3879        // The canonical dynamic child is one short-lived task per unit of work, so retaining every one that has ever
3880        // finished would grow without bound. Spawn far more than a tree should ever hold and require the count to
3881        // come back to the static baseline.
3882        for _ in 0..500 {
3883            sup_handle.spawn(MockWorker::completing("ephemeral", Duration::from_millis(1)));
3884        }
3885
3886        wait_until("every dynamic child has been reaped", || {
3887            sup_handle.active_children() == 0 && tree.snapshot().root.children.len() == 1
3888        })
3889        .await;
3890
3891        let snapshot = tree.snapshot();
3892        assert_eq!(snapshot.root.children.len(), 1, "only the static child remains");
3893        assert_eq!(snapshot.root.children[0].name, "static-worker");
3894
3895        tx.send(()).unwrap();
3896        let _ = join_supervisor(handle).await;
3897    }
3898
3899    #[tokio::test]
3900    async fn snapshot_nests_child_supervisors() {
3901        let mut child_sup = Supervisor::new("child-sup").unwrap();
3902        child_sup.add_worker(MockWorker::long_running("inner-worker"));
3903
3904        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3905        parent_sup.add_worker(MockWorker::long_running("outer-worker"));
3906        parent_sup.add_worker(child_sup);
3907
3908        let tree = parent_sup.tree_handle();
3909        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3910
3911        wait_until("the nested subtree is populated", || {
3912            let snapshot = tree.snapshot();
3913            snapshot
3914                .root
3915                .children
3916                .iter()
3917                .any(|child| child.name == "child-sup" && !child.children.is_empty())
3918        })
3919        .await;
3920
3921        let snapshot = tree.snapshot();
3922        let nested = find_node(&snapshot.root.children, "child-sup");
3923        assert_eq!(nested.kind, NodeKind::Supervisor);
3924        assert!(nested.supervision.is_some(), "a nested supervisor reports its settings");
3925        assert_eq!(nested.process_name.as_deref(), Some("parent_sup.child_sup"));
3926        assert_eq!(nested.children.len(), 1);
3927
3928        let grandchild = &nested.children[0];
3929        assert_eq!(grandchild.name, "inner-worker");
3930        assert_eq!(grandchild.kind, NodeKind::Worker);
3931        assert_eq!(
3932            grandchild.process_name.as_deref(),
3933            Some("parent_sup.child_sup.inner_worker")
3934        );
3935
3936        assert_eq!(snapshot.totals.supervisors, 2);
3937        assert_eq!(snapshot.totals.workers, 2);
3938        assert_eq!(snapshot.totals.max_depth, 3);
3939
3940        tx.send(()).unwrap();
3941        let _ = join_supervisor(handle).await;
3942    }
3943
3944    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3945    async fn dedicated_runtime_subtree_is_visible_across_the_thread_boundary() {
3946        let mut child_sup = Supervisor::new("child-sup")
3947            .unwrap()
3948            .with_dedicated_runtime(RuntimeConfiguration::single_threaded());
3949        child_sup.add_worker(MockWorker::long_running("inner-worker"));
3950
3951        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3952        parent_sup.add_worker(child_sup);
3953
3954        // Taken from the parent, before the child is running on an OS thread of its own.
3955        let tree = parent_sup.tree_handle();
3956        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3957
3958        wait_until("the dedicated-runtime subtree is visible", || {
3959            tree.snapshot()
3960                .root
3961                .children
3962                .iter()
3963                .any(|child| child.name == "child-sup" && !child.children.is_empty())
3964        })
3965        .await;
3966
3967        // Read the tree from a thread that is neither the parent's nor the dedicated runtime's.
3968        let probe = tree.clone();
3969        let snapshot = tokio::task::spawn_blocking(move || probe.snapshot())
3970            .await
3971            .expect("probe should not panic");
3972
3973        let nested = find_node(&snapshot.root.children, "child-sup");
3974        assert_eq!(nested.state, NodeState::Running);
3975        assert_eq!(
3976            nested.supervision.expect("supervisor settings").dedicated_threads,
3977            Some(1)
3978        );
3979        assert_eq!(nested.children.len(), 1);
3980
3981        // A supervisor on a dedicated runtime re-roots its process name when it starts, so it is *not* scoped under
3982        // its parent. That is a known defect (see the `Dedicated` branch of `create_worker_future`), not a design
3983        // choice -- but it is also why a node's name has to come from its own run rather than from the process its
3984        // parent created for it, since only the former names the resource group its allocations land in. Pinned here
3985        // so that fixing the defect is a deliberate act rather than a silent regression.
3986        assert_eq!(nested.process_name.as_deref(), Some("child_sup"));
3987        assert_eq!(
3988            nested.children[0].process_name.as_deref(),
3989            Some("child_sup.inner_worker")
3990        );
3991
3992        tx.send(()).unwrap();
3993        let _ = join_supervisor(handle).await;
3994    }
3995
3996    #[tokio::test]
3997    async fn a_running_node_always_has_a_process() {
3998        // A nested supervisor being restarted is briefly stopped, and a snapshot taken in that window must not
3999        // present the previous generation's processes as if they were alive. Sampling exactly inside the window is
4000        // inherently racy, so assert the invariant that has to hold at every instant instead.
4001        let mut child_sup = Supervisor::new("child-sup")
4002            .unwrap()
4003            .with_restart_strategy(RestartStrategy::new(RestartMode::OneForOne, 0, Duration::from_secs(5)));
4004        child_sup.add_worker(MockWorker::failing("inner", Duration::from_millis(10)));
4005
4006        let mut parent_sup = Supervisor::new("parent-sup")
4007            .unwrap()
4008            .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(50, Duration::from_secs(5)));
4009        parent_sup.add_worker(child_sup);
4010
4011        let tree = parent_sup.tree_handle();
4012        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
4013
4014        for _ in 0..60 {
4015            let snapshot = tree.snapshot();
4016            assert_running_nodes_have_processes(&snapshot.root);
4017            tokio::time::sleep(Duration::from_millis(5)).await;
4018        }
4019
4020        wait_until("the nested supervisor has been restarted", || {
4021            find_node(&tree.snapshot().root.children, "child-sup").restart_count >= 1
4022        })
4023        .await;
4024
4025        tx.send(()).unwrap();
4026        let _ = join_supervisor(handle).await;
4027    }
4028
4029    fn assert_running_nodes_have_processes(node: &NodeSnapshot) {
4030        if node.state == NodeState::Running {
4031            assert!(
4032                node.process_id.is_some() && node.process_name.is_some(),
4033                "node '{}' reports as running but names no process",
4034                node.name
4035            );
4036            assert!(
4037                node.uptime_ms.is_some(),
4038                "node '{}' is running but has no uptime",
4039                node.name
4040            );
4041        } else {
4042            assert!(
4043                node.uptime_ms.is_none(),
4044                "node '{}' is not running but reports an uptime",
4045                node.name
4046            );
4047        }
4048
4049        for child in &node.children {
4050            assert_running_nodes_have_processes(child);
4051        }
4052    }
4053
4054    #[tokio::test]
4055    async fn snapshot_after_the_run_reports_a_stopped_tree() {
4056        fn assert_send_sync<T: Send + Sync + 'static>() {}
4057        assert_send_sync::<SupervisionTreeHandle>();
4058
4059        let mut sup = Supervisor::new("test-sup").unwrap();
4060        sup.add_worker(MockWorker::long_running("worker1"));
4061
4062        let tree = sup.tree_handle();
4063        let (tx, handle) = run_supervisor_with_trigger(sup).await;
4064        wait_until("the child is running", || tree.snapshot().root.children.len() == 1).await;
4065
4066        tx.send(()).unwrap();
4067        join_supervisor(handle).await.unwrap();
4068
4069        // A stopped supervisor reports as stopped rather than presenting a subtree of processes that no longer exist.
4070        let snapshot = tree.snapshot();
4071        assert_eq!(snapshot.root.state, NodeState::Registered);
4072        assert!(snapshot.root.children.is_empty());
4073        assert_eq!(
4074            snapshot.root.supervision.expect("supervisor settings").generation,
4075            1,
4076            "the generation count outlives the run"
4077        );
4078
4079        // Not running is not the same as never having run: the node keeps the identity of the run that just ended, so
4080        // that a stopped tree still says what it was. `state` is what distinguishes the two.
4081        assert!(snapshot.root.process_id.is_some(), "a stopped node says what it ran as");
4082        assert_eq!(snapshot.root.process_name.as_deref(), Some("test_sup"));
4083        assert!(snapshot.root.started_at.is_some());
4084        assert_eq!(snapshot.root.uptime_ms, None, "a node that isn't running has no uptime");
4085    }
4086
4087    #[tokio::test]
4088    async fn an_exited_nested_supervisor_says_what_it_ran_as() {
4089        // A nested supervisor that stops for good is kept by its parent as a tombstone, and what it ran as is what a
4090        // postmortem wants from it -- which is what its own workers have always reported in the same situation.
4091        let mut child_sup = Supervisor::new("child-sup")
4092            .unwrap()
4093            .with_auto_shutdown(AutoShutdown::AnySignificant);
4094        child_sup.add_worker(
4095            runtime::supervisable(MockWorker::completing("one-shot", Duration::from_millis(10)))
4096                .temporary()
4097                .with_significant(true)
4098                .build(),
4099        );
4100
4101        // Temporary, so the parent does not bring it back once it stops.
4102        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
4103        parent_sup.add_worker(ChildSpecification::from(child_sup).with_restart_type(RestartType::Temporary));
4104
4105        let tree = parent_sup.tree_handle();
4106        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
4107
4108        // Record the identity while it is still running, so the tombstone has to report the same one rather than
4109        // merely reporting something.
4110        let mut running = None;
4111        wait_until("the nested supervisor is running", || {
4112            let snapshot = tree.snapshot();
4113            let nested = find_node(&snapshot.root.children, "child-sup");
4114            if nested.state == NodeState::Running {
4115                running = Some((nested.process_id, nested.process_name.clone(), nested.started_at));
4116                return true;
4117            }
4118            false
4119        })
4120        .await;
4121        let (process_id, process_name, started_at) = running.expect("observed running");
4122        assert!(process_id.is_some(), "a running supervisor has a process");
4123
4124        wait_until("the nested supervisor has exited", || {
4125            find_node(&tree.snapshot().root.children, "child-sup").state == NodeState::Exited
4126        })
4127        .await;
4128
4129        let snapshot = tree.snapshot();
4130        let nested = find_node(&snapshot.root.children, "child-sup");
4131        assert_eq!(nested.process_id, process_id, "an exited supervisor keeps its identity");
4132        assert_eq!(nested.process_name, process_name);
4133        assert_eq!(nested.started_at, started_at);
4134        assert!(nested.exited_at.is_some(), "a tombstone records when it exited");
4135        assert_eq!(nested.uptime_ms, None, "a node that isn't running has no uptime");
4136        assert!(
4137            nested.resource_group.is_some(),
4138            "and the group its allocations were attributed to"
4139        );
4140
4141        tx.send(()).unwrap();
4142        let _ = join_supervisor(handle).await;
4143    }
4144
4145    #[tokio::test]
4146    async fn initialization_failure_leaves_nothing_running() {
4147        let mut sup = Supervisor::new("test-sup").unwrap();
4148        sup.add_worker(MockWorker::init_failure("broken"));
4149
4150        let tree = sup.tree_handle();
4151
4152        // Run directly rather than through the readiness barrier: this supervisor fails almost immediately, so it may
4153        // never be observed running at all.
4154        let mut sup = sup;
4155        let (_tx, rx) = oneshot::channel::<()>();
4156        let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
4157            .await
4158            .expect("supervisor should exit promptly");
4159        assert!(matches!(result, Err(SupervisorError::FailedToInitialize { .. })));
4160
4161        // The failure path returns through `run_inner` like every other, so the tree is cleaned up either way.
4162        let snapshot = tree.snapshot();
4163        assert_eq!(snapshot.root.state, NodeState::Registered);
4164        assert_eq!(snapshot.totals.running, 0);
4165    }
4166
4167    #[tokio::test]
4168    async fn resources_are_attributed_to_supervisors_not_workers() {
4169        let mut child_sup = Supervisor::new("child-sup").unwrap();
4170        child_sup.add_worker(MockWorker::long_running("inner-worker"));
4171
4172        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
4173        parent_sup.add_worker(child_sup);
4174
4175        let tree = parent_sup.tree_handle();
4176        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
4177        wait_until("the nested subtree is populated", || {
4178            tree.snapshot()
4179                .root
4180                .children
4181                .iter()
4182                .any(|child| !child.children.is_empty())
4183        })
4184        .await;
4185
4186        let snapshot = tree.snapshot();
4187
4188        // The tracking allocator is a process-wide facility that the test binary doesn't install, so every byte count
4189        // here reads zero. That is exactly why the snapshot reports whether tracking is on at all: without it, zero
4190        // bytes is indistinguishable from nothing being measured.
4191        assert!(!snapshot.resource_tracking_enabled);
4192
4193        let nested = find_node(&snapshot.root.children, "child-sup");
4194        assert!(
4195            nested.resources.is_some(),
4196            "a supervisor owns the resource group named by its own process"
4197        );
4198        assert_eq!(nested.resource_group.as_deref(), nested.process_name.as_deref());
4199
4200        // A worker inherits its supervisor's group rather than owning one, so it names the group but carries no
4201        // figures of its own -- reporting the supervisor's totals against each of its workers would count them twice.
4202        let worker = &nested.children[0];
4203        assert!(worker.resources.is_none(), "a worker owns no resource group");
4204        assert_eq!(worker.resource_group.as_deref(), nested.process_name.as_deref());
4205
4206        tx.send(()).unwrap();
4207        let _ = join_supervisor(handle).await;
4208    }
4209
4210    #[tokio::test]
4211    async fn snapshot_serializes_to_the_expected_shape() {
4212        let mut child_sup = Supervisor::new("child-sup").unwrap();
4213        child_sup.add_worker(MockWorker::long_running("inner-worker"));
4214
4215        let mut parent_sup = Supervisor::new("parent-sup").unwrap();
4216        parent_sup.add_worker(child_sup);
4217
4218        let tree = parent_sup.tree_handle();
4219        let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
4220        wait_until("the nested subtree is populated", || {
4221            tree.snapshot()
4222                .root
4223                .children
4224                .iter()
4225                .any(|child| !child.children.is_empty())
4226        })
4227        .await;
4228
4229        let value = serde_json::to_value(tree.snapshot()).expect("snapshot should serialize");
4230
4231        // This is the wire contract an HTTP consumer reads, so pin its shape rather than just its serializability.
4232        assert!(value["captured_at"].is_u64(), "a timestamp is a plain integer");
4233        assert_eq!(value["root"]["kind"], "supervisor");
4234        assert_eq!(value["root"]["state"], "running");
4235        assert_eq!(value["root"]["restart"], "permanent");
4236        assert_eq!(value["root"]["supervision"]["restart_mode"], "one_for_one");
4237        assert_eq!(value["root"]["supervision"]["auto_shutdown"], "never");
4238
4239        let nested = &value["root"]["children"][0];
4240        assert_eq!(nested["kind"], "supervisor");
4241        let worker = &nested["children"][0];
4242        assert_eq!(worker["kind"], "worker");
4243        assert!(
4244            worker["children"]
4245                .as_array()
4246                .expect("children is always an array")
4247                .is_empty(),
4248            "children is present and empty for a worker, so consumers can recurse unconditionally"
4249        );
4250        assert!(worker.get("resources").is_none(), "a worker reports no resource usage");
4251        assert!(
4252            worker.get("supervision").is_none(),
4253            "a worker reports no supervision settings"
4254        );
4255        assert!(
4256            worker.get("children_truncated").is_none(),
4257            "an untruncated node omits the truncation flag"
4258        );
4259
4260        tx.send(()).unwrap();
4261        let _ = join_supervisor(handle).await;
4262    }
4263
4264    #[tokio::test]
4265    async fn a_worker_that_drives_its_own_supervisor_is_shown_as_its_parent() {
4266        // Some workers own a supervisor rather than being one: they build it after initialization and run it inside
4267        // their own future, so their parent has no supervisor value to record and the subtree would otherwise be
4268        // invisible. This is how the largest subtrees in a real process are shaped.
4269        struct SupervisorDrivingWorker;
4270
4271        #[async_trait]
4272        impl Supervisable for SupervisorDrivingWorker {
4273            fn name(&self) -> &str {
4274                "driver"
4275            }
4276
4277            fn shutdown_strategy(&self) -> ShutdownStrategy {
4278                ShutdownStrategy::Graceful(Duration::MAX)
4279            }
4280
4281            async fn initialize(
4282                &self, process_shutdown: ShutdownHandle,
4283            ) -> Result<SupervisorFuture, InitializationError> {
4284                Ok(Box::pin(async move {
4285                    let mut inner = Supervisor::new("inner-sup").expect("valid name");
4286                    inner.add_worker(MockWorker::long_running("inner-worker"));
4287                    inner
4288                        .run_with_shutdown_inner(process_shutdown, None)
4289                        .await
4290                        .map_err(Into::into)
4291                }))
4292            }
4293        }
4294
4295        let mut sup = Supervisor::new("test-sup").unwrap();
4296        sup.add_worker(SupervisorDrivingWorker);
4297
4298        let tree = sup.tree_handle();
4299        let (tx, handle) = run_supervisor_with_trigger(sup).await;
4300
4301        wait_until("the worker's own supervisor has attached itself", || {
4302            let snapshot = tree.snapshot();
4303            snapshot
4304                .root
4305                .children
4306                .first()
4307                .is_some_and(|child| !child.children.is_empty())
4308        })
4309        .await;
4310
4311        let snapshot = tree.snapshot();
4312        let driver = find_node(&snapshot.root.children, "driver");
4313        assert_eq!(
4314            driver.kind,
4315            NodeKind::Worker,
4316            "the worker is what its parent supervises"
4317        );
4318
4319        let inner_sup = find_node(&driver.children, "inner-sup");
4320        assert_eq!(inner_sup.kind, NodeKind::Supervisor);
4321        assert_eq!(inner_sup.children.len(), 1);
4322        assert_eq!(inner_sup.children[0].name, "inner-worker");
4323        assert_eq!(snapshot.totals.max_depth, 4);
4324
4325        tx.send(()).unwrap();
4326        let _ = join_supervisor(handle).await;
4327    }
4328
4329    #[tokio::test]
4330    async fn mixed_child_traffic_keeps_the_roster_consistent() {
4331        // Exercises every path that mutates the roster in one run -- static spawn, dynamic spawn, restart in place,
4332        // and removal without restart. It asserts little itself: the drift check inside the roster is what is under
4333        // test, and it runs on every mutation.
4334        let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
4335            RestartStrategy::one_to_one().with_intensity_and_period(100, Duration::from_secs(5)),
4336        );
4337        sup.add_worker(MockWorker::long_running("stable"));
4338        sup.add_worker(MockWorker::failing("flapper", Duration::from_millis(5)));
4339
4340        let tree = sup.tree_handle();
4341        let sup_handle = sup.handle();
4342        let (tx, handle) = run_supervisor_with_trigger(sup).await;
4343
4344        for _ in 0..25 {
4345            sup_handle.spawn(MockWorker::completing("ephemeral", Duration::from_millis(2)));
4346            sup_handle.spawn(MockWorker::long_running("lingering"));
4347            let snapshot = tree.snapshot();
4348            assert_running_nodes_have_processes(&snapshot.root);
4349            tokio::time::sleep(Duration::from_millis(2)).await;
4350        }
4351
4352        wait_until("the flapping child has restarted several times", || {
4353            find_node(&tree.snapshot().root.children, "flapper").restart_count >= 3
4354        })
4355        .await;
4356
4357        tx.send(()).unwrap();
4358        let _ = join_supervisor(handle).await;
4359    }
4360}