saluki_core/runtime/tree/
mod.rs

1//! Supervision-tree introspection.
2//!
3//! A supervisor knows the shape of the subtree below it -- which children it has, whether each is a worker or a
4//! nested supervisor, how each is configured, and how each has fared -- but that knowledge lives inside the future
5//! that drives the supervisor and is unreachable from anywhere else. This module makes it observable.
6//!
7//! # How it works
8//!
9//! Every [`Supervisor`][super::Supervisor] owns an [`Arc<SupervisorNode>`][SupervisorNode]: a small piece of shared
10//! bookkeeping that the supervisor's run loop writes to as children start, exit, and restart. A parent records a clone
11//! of each child supervisor's node, so the nodes form a graph mirroring the supervision tree, and walking it from any
12//! node yields the subtree rooted there -- which is what [`SupervisionTreeHandle::snapshot`] does. Nodes are shared
13//! rather than copied, so a handle taken before a supervisor starts observes every subsequent generation of it,
14//! including one running on a dedicated runtime on another OS thread.
15//!
16//! Note that this module owns more than an observer's view: [`Roster`] holds the supervisor's own live child roster,
17//! not a mirror of it. That is deliberate -- it is what makes the two impossible to drift apart -- but it means
18//! changes here affect supervision itself and not only what is reported about it.
19
20use std::{
21    sync::{Arc, Mutex},
22    time::{Duration, Instant, SystemTime, UNIX_EPOCH},
23};
24
25use saluki_common::{
26    collections::{FastHashMap, FastHashSet, FastIndexMap},
27    resource_tracking::{ResourceGroupRegistry, ResourceStatsSnapshot},
28};
29use tracing::warn;
30
31use super::{
32    process::Process,
33    restart::{RestartStrategy, RestartType},
34    supervisor::AutoShutdown,
35    ProcessId,
36};
37
38mod api;
39pub use self::api::{SupervisionTreeAPIHandler, SupervisionTreeState};
40
41mod worker;
42pub use self::worker::SupervisionTreeWorker;
43
44mod snapshot;
45pub use self::snapshot::{
46    NodeKind, NodeSnapshot, NodeState, ResourceUsage, SupervisionSettings, TreeSnapshot, TreeTotals, UnixMillis,
47};
48
49/// API route serving a snapshot of a supervision tree.
50///
51/// Exported so that a client can address the route without restating the path.
52pub const SUPERVISION_TREE_ROUTE: &str = "/runtime/processes";
53
54/// Maximum depth the snapshot walk descends before it stops and reports the subtree as truncated.
55///
56/// The node graph cannot contain a cycle -- a supervisor can't be its own ancestor, since
57/// [`add_worker`][super::Supervisor::add_worker] takes its child by value -- so this is not a termination condition
58/// but a stack guard. The walk is recursive and runs wherever the caller asks for a snapshot, which may be a
59/// request-serving task with a modest stack, and a diagnostics endpoint should not be able to abort the process even
60/// if a future change does manage to build a pathological tree.
61const MAX_TREE_DEPTH: usize = 64;
62
63/// A wall-clock and monotonic reading of the same instant.
64///
65/// Both are needed. Wall clock is what an operator wants to see (and what correlates with logs), but it can step
66/// backwards, so ages computed from it can come out negative or absurd. The monotonic reading can't, so every
67/// duration is derived from it and every displayed timestamp from the wall clock.
68///
69/// Uses [`std::time::Instant`] rather than [`tokio::time::Instant`] deliberately: a snapshot may be taken from a
70/// thread with no Tokio runtime, and it should report real elapsed time rather than a test's virtual clock.
71#[derive(Clone, Copy)]
72pub(super) struct Stamp {
73    wall: SystemTime,
74    mono: Instant,
75}
76
77impl Stamp {
78    pub(super) fn now() -> Self {
79        Self {
80            wall: SystemTime::now(),
81            mono: Instant::now(),
82        }
83    }
84
85    fn wall_millis(&self) -> UnixMillis {
86        // A pre-epoch clock makes `duration_since` fail rather than return a negative duration; report it as the
87        // epoch rather than panicking inside a diagnostics path.
88        UnixMillis(duration_millis(
89            self.wall.duration_since(UNIX_EPOCH).unwrap_or_default(),
90        ))
91    }
92
93    fn elapsed_millis(&self) -> u64 {
94        duration_millis(self.mono.elapsed())
95    }
96}
97
98/// Identity of a child that survives the id churn of a restart.
99///
100/// A supervisor keys its live roster by an id drawn from a monotonic counter, and a restart -- whether a
101/// [`OneForAll`][super::RestartMode::OneForAll] restart of the group or a restart of the supervisor itself --
102/// discards the whole roster and re-registers every eligible child under a *fresh* id. Facts that must outlive that
103/// -- when the child was first created, how many times it has been restarted -- therefore can't be keyed by roster
104/// id.
105#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
106pub(super) enum ChildKey {
107    /// A child declared before the run, identified by its index in the supervisor's static child list.
108    ///
109    /// That list is fixed once the supervisor starts, so the index is a stable identity across generations.
110    Static(usize),
111
112    /// A dynamically spawned child, identified by its roster id.
113    ///
114    /// A dynamic child is never restored across generations, so its roster id is as stable an identity as it needs.
115    /// Ids come from a monotonic counter, so ordering by one is ordering by spawn.
116    Dynamic(u64),
117}
118
119impl ChildKey {
120    fn is_dynamic(&self) -> bool {
121        matches!(self, Self::Dynamic(_))
122    }
123}
124
125/// The process a child was started under.
126///
127/// Produced by the supervisor's worker bookkeeping at the moment a child's task is spawned, which is the only place
128/// the child's [`Process`] exists and therefore the only place its identity can be captured.
129#[derive(Clone)]
130pub(super) struct StartedChild {
131    process_id: ProcessId,
132    process_name: Arc<str>,
133    at: Stamp,
134}
135
136impl StartedChild {
137    pub(super) fn new(process: &Process, process_name: Arc<str>) -> Self {
138        Self {
139            process_id: *process.id(),
140            process_name,
141            at: Stamp::now(),
142        }
143    }
144}
145
146/// Facts about a child that its supervisor knows at registration time.
147///
148/// Assembled by the supervisor (which can see its own child specifications and configuration) and handed to
149/// [`Roster::insert`], so that this module needs no visibility into either.
150pub(super) struct ChildFacts {
151    /// Identity that survives a group restart.
152    pub(super) key: ChildKey,
153
154    pub(super) name: Arc<str>,
155
156    /// The child's own node, if the child is a nested supervisor.
157    ///
158    /// This is what links a parent's bookkeeping to its child's, and so what makes the tree walkable.
159    pub(super) node: Option<Arc<SupervisorNode>>,
160
161    pub(super) restart: RestartType,
162
163    /// Whether the child's termination can drive its supervisor to shut down.
164    pub(super) significant: bool,
165}
166
167/// What a child is, from its parent's point of view.
168enum RecordKind {
169    Worker,
170
171    Supervisor(Arc<SupervisorNode>),
172}
173
174/// A supervisor's record of one child in the current generation.
175struct ChildRecord {
176    key: ChildKey,
177    name: Arc<str>,
178    kind: RecordKind,
179    restart: RestartType,
180    significant: bool,
181    /// The process the child is currently running under, or last ran under. `None` only for a child that never
182    /// successfully started.
183    start: Option<StartedChild>,
184    /// When the child exited for good, if it has. `None` while it is running.
185    exited: Option<Stamp>,
186}
187
188/// Cumulative facts about a child that outlive any single generation of it.
189struct ChildHistory {
190    /// When the child first entered its supervisor's roster.
191    ///
192    /// Deliberately not its registration time: registration happens once, during program construction, and never
193    /// again on a later generation, so it measures how long ago a builder method was called rather than anything
194    /// about the running tree.
195    created: Stamp,
196    restarts: u32,
197}
198
199/// A supervisor's configuration.
200///
201/// Owned by the supervisor's node rather than by the [`Supervisor`][super::Supervisor] itself, so that there is one
202/// copy rather than two that have to be kept in step, and so that it can be read before the supervisor has ever run.
203#[derive(Clone, Copy)]
204pub(super) struct NodeConfig {
205    pub(super) restart_strategy: RestartStrategy,
206    pub(super) auto_shutdown: AutoShutdown,
207    pub(super) shutdown_budget: Option<Duration>,
208    pub(super) dedicated_threads: Option<usize>,
209}
210
211/// The mutable half of a [`SupervisorNode`].
212struct NodeInner {
213    config: NodeConfig,
214    /// The supervisor's most recent run. `None` only before its first one.
215    ///
216    /// Retained after a run ends rather than cleared, so that a stopped supervisor still reports what it last ran as
217    /// -- which is what a postmortem asks of it, and what its own workers have always answered. Whether this run is
218    /// still in progress is carried by `running`, not by the presence of an identity.
219    ///
220    /// The process name here is recorded by the supervisor itself rather than by its parent, because the two can
221    /// differ: a supervisor on a dedicated runtime re-roots its process name when it starts (a known defect), so only
222    /// the supervisor knows what it actually runs as -- and only that name matches the resource group its allocations
223    /// are attributed to.
224    run: Option<StartedChild>,
225    /// Whether `run` is in progress rather than ended.
226    running: bool,
227    generation: u64,
228    restarts_performed: u64,
229    /// Children of the current generation, in declaration order.
230    ///
231    /// Ordered rather than hashed so the rendered tree and the serialized JSON come out in the order children were
232    /// declared, which is how the tree is written and therefore how it reads.
233    children: FastIndexMap<u64, ChildRecord>,
234    /// Per-child facts that outlive a generation, keyed by an identity that survives a restart.
235    ///
236    /// Only statically declared children keep an entry across a restart of the supervisor itself, since only they
237    /// have an identity that a later generation can look up again.
238    history: FastHashMap<ChildKey, ChildHistory>,
239    /// Supervisors that a worker child drives internally, keyed by that worker's roster id.
240    ///
241    /// Some workers own a supervisor rather than being one -- they build it after initialization and run it inside
242    /// their own future -- so the parent has no supervisor value to record at registration time and the subtree would
243    /// otherwise be invisible. Such a supervisor registers itself here when it starts.
244    adopted: FastHashMap<u64, Arc<SupervisorNode>>,
245}
246
247impl NodeInner {
248    /// Discards everything specific to a single generation of the supervisor.
249    ///
250    /// Used when a run begins or ends, which discards the generation wholesale. A group restart within a run keeps
251    /// the tombstones it won't bring back, and so prunes the roster itself.
252    ///
253    /// A dynamic child's history belongs to the generation and goes with it: its key is a roster id drawn from a
254    /// counter that is never reset, so once the generation is gone nothing can look the key up again -- no later
255    /// generation restores the child, and no later spawn reuses the id. A statically declared child's history is
256    /// kept, since the next generation re-registers it under the same identity.
257    fn clear_generation(&mut self) {
258        self.children.clear();
259        self.adopted.clear();
260        self.history.retain(|key, _| !key.is_dynamic());
261    }
262
263    fn running_children(&self) -> usize {
264        self.children.values().filter(|r| r.exited.is_none()).count()
265    }
266}
267
268/// The supervised child slot a worker is running as.
269///
270/// Installed around every supervised worker task so that a supervisor a worker builds and runs inside its own future
271/// can attach itself to the tree, which it could not otherwise do: the parent has no supervisor value to record at
272/// registration time, because the supervisor does not exist until the worker is already running.
273#[derive(Clone)]
274pub(super) struct TreeParent {
275    node: Arc<SupervisorNode>,
276    child_id: u64,
277}
278
279impl TreeParent {
280    pub(super) fn new(node: Arc<SupervisorNode>, child_id: u64) -> Self {
281        Self { node, child_id }
282    }
283}
284
285tokio::task_local! {
286    /// The child slot the currently running worker occupies in its supervisor.
287    pub(super) static CURRENT_TREE_PARENT: TreeParent;
288}
289
290/// Shared bookkeeping for one supervisor.
291///
292/// Held by the supervisor itself, by every internal clone of it (so all generations write to the same place), by its
293/// parent's record of it, and by any [`SupervisionTreeHandle`] pointed at it.
294pub(super) struct SupervisorNode {
295    id: Arc<str>,
296    created: Stamp,
297    state: Mutex<NodeInner>,
298}
299
300impl SupervisorNode {
301    pub(super) fn new(id: Arc<str>) -> Self {
302        Self {
303            id,
304            created: Stamp::now(),
305            state: Mutex::new(NodeInner {
306                config: NodeConfig {
307                    restart_strategy: RestartStrategy::default(),
308                    auto_shutdown: AutoShutdown::default(),
309                    shutdown_budget: None,
310                    dedicated_threads: None,
311                },
312                run: None,
313                running: false,
314                generation: 0,
315                restarts_performed: 0,
316                children: FastIndexMap::default(),
317                history: FastHashMap::default(),
318                adopted: FastHashMap::default(),
319            }),
320        }
321    }
322
323    fn state(&self) -> std::sync::MutexGuard<'_, NodeInner> {
324        self.state.lock().expect("supervision-tree node lock poisoned")
325    }
326
327    pub(super) fn config(&self) -> NodeConfig {
328        self.state().config
329    }
330
331    pub(super) fn update_config(&self, f: impl FnOnce(&mut NodeConfig)) {
332        f(&mut self.state().config);
333    }
334
335    /// Records that the supervisor has begun a run under `process`.
336    ///
337    /// Resets everything specific to a generation, so a run that ended abruptly -- a panicked task, or a dedicated
338    /// runtime thread that unwound without reaching its cleanup -- can't leave stale children behind for the next
339    /// generation to be confused with. Counts the new generation as a restart of every statically declared child,
340    /// since every one of them is about to be re-registered under a fresh process.
341    pub(super) fn begin_run(self: &Arc<Self>, process: &Process) {
342        let mut state = self.state();
343
344        if state.running {
345            // Two overlapping runs share one node, one spawn queue, one id counter and one child gauge, so the second
346            // silently clobbers the first. This is unreachable from outside the crate -- `add_worker` takes its child
347            // by value and `Supervisor` has no public `Clone` -- so what this really guards is a restart that begins
348            // before the previous generation has finished stopping.
349            debug_assert!(
350                false,
351                "supervisor '{}' began a run while another was still active",
352                self.id
353            );
354            warn!(
355                supervisor_id = %self.id,
356                "Supervisor began a run while another was still active; a supervisor must not be supervised by two \
357                 parents."
358            );
359        }
360
361        state.run = Some(StartedChild::new(process, process.name().into()));
362        state.running = true;
363        state.generation += 1;
364        state.clear_generation();
365
366        // A statically declared child is identified by its position in a list fixed before the run, so this
367        // generation re-registers it under the same identity and its history carries over. Every one of them comes
368        // back -- `spawn_static_children` filters on nothing -- and each comes back under a fresh process, so this
369        // generation is a restart of each of them.
370        //
371        // Only static children are left to count here: `clear_generation` has just discarded the rest. There is
372        // nothing at all to count on the first run, which has no history to carry over.
373        for history in state.history.values_mut() {
374            history.restarts += 1;
375        }
376        drop(state);
377
378        // If we were started inside a worker's task, attach ourselves to the tree under that worker. This is how a
379        // supervisor built and run inside a worker's future -- rather than handed to a parent as a child -- becomes
380        // visible. Absent when the supervisor is a root, or is running on a dedicated runtime (whose thread has no
381        // task-locals), and in the latter case the parent already recorded us directly.
382        let _ = CURRENT_TREE_PARENT.try_with(|parent| {
383            parent.node.state().adopted.insert(parent.child_id, Arc::clone(self));
384        });
385    }
386
387    /// Records that the supervisor's run has ended.
388    ///
389    /// Keeps everything cumulative -- generation count, restarts performed, static children's history, and the
390    /// identity of the run itself -- and discards the generation's children, so a stopped supervisor reports as
391    /// stopped rather than presenting a subtree of processes that no longer exist.
392    pub(super) fn end_run(&self) {
393        let mut state = self.state();
394        state.running = false;
395        state.clear_generation();
396        drop(state);
397
398        let _ = CURRENT_TREE_PARENT.try_with(|parent| {
399            parent.node.state().adopted.remove(&parent.child_id);
400        });
401    }
402}
403
404/// A supervisor's children, in both the form the supervisor runs from and the form an observer reads.
405///
406/// The supervisor's own roster is a plain map it owns exclusively; the observable bookkeeping lives behind a lock and
407/// is shared. Both are updated through this one type so that they cannot drift: there is no way to add, restart or
408/// remove a child in one without doing so in the other.
409pub(super) struct Roster<E> {
410    live: FastHashMap<u64, E>,
411    node: Arc<SupervisorNode>,
412}
413
414impl<E> Roster<E> {
415    pub(super) fn new(node: Arc<SupervisorNode>) -> Self {
416        Self {
417            live: FastHashMap::default(),
418            node,
419        }
420    }
421
422    pub(super) fn get(&self, id: u64) -> Option<&E> {
423        self.live.get(&id)
424    }
425
426    pub(super) fn values(&self) -> impl Iterator<Item = &E> {
427        self.live.values()
428    }
429
430    /// Registers a newly started child.
431    ///
432    /// `facts` describes the child as its supervisor configured it, and `started` identifies the process it was
433    /// started under. A child re-registered under a fresh id after a group restart keeps the creation time and
434    /// restart count recorded against its [`ChildKey`].
435    pub(super) fn insert(&mut self, id: u64, entry: E, facts: ChildFacts, started: StartedChild) {
436        self.live.insert(id, entry);
437
438        let mut state = self.node.state();
439        state.history.entry(facts.key).or_insert_with(|| ChildHistory {
440            created: started.at,
441            restarts: 0,
442        });
443
444        let kind = match facts.node {
445            Some(node) => RecordKind::Supervisor(node),
446            None => RecordKind::Worker,
447        };
448
449        state.children.insert(
450            id,
451            ChildRecord {
452                key: facts.key,
453                name: facts.name,
454                kind,
455                restart: facts.restart,
456                significant: facts.significant,
457                start: Some(started),
458                exited: None,
459            },
460        );
461        drop(state);
462
463        self.debug_assert_consistent();
464    }
465
466    pub(super) fn restart_in_place(&mut self, id: u64, started: StartedChild) {
467        let mut state = self.node.state();
468        state.restarts_performed += 1;
469
470        // Whatever the previous incarnation adopted refers to a supervisor that has since stopped. The new
471        // incarnation re-registers if it drives one of its own.
472        state.adopted.remove(&id);
473
474        if let Some(record) = state.children.get_mut(&id) {
475            let key = record.key;
476            record.start = Some(started);
477            record.exited = None;
478            if let Some(history) = state.history.get_mut(&key) {
479                history.restarts += 1;
480            }
481        }
482        drop(state);
483
484        self.debug_assert_consistent();
485    }
486
487    /// Removes the child under `id`, which has exited and will not be restarted.
488    ///
489    /// A statically declared child is kept as a tombstone rather than dropped: it remains part of the tree's declared
490    /// shape, and "declared, ran, and then stopped for good" is precisely what an operator needs to see. A
491    /// dynamically spawned child is dropped outright -- the canonical dynamic child is one short-lived task per unit
492    /// of work, so retaining every one that has ever finished would grow without bound.
493    pub(super) fn remove(&mut self, id: u64) -> Option<E> {
494        let entry = self.live.remove(&id);
495
496        let mut state = self.node.state();
497        state.adopted.remove(&id);
498        match state.children.get(&id).map(|record| record.key) {
499            // Order is restored by sorting on `ChildKey` when a snapshot is taken, so the roster itself doesn't need
500            // to preserve it -- which lets a dynamic child, of which there may be one per unit of work, leave in
501            // constant time rather than shifting every entry behind it.
502            Some(key) if key.is_dynamic() => {
503                state.children.swap_remove(&id);
504                state.history.remove(&key);
505            }
506            Some(_) => {
507                if let Some(record) = state.children.get_mut(&id) {
508                    record.exited = Some(Stamp::now());
509                }
510            }
511            None => {}
512        }
513        drop(state);
514
515        self.debug_assert_consistent();
516
517        entry
518    }
519
520    /// Clears the roster ahead of a group restart, counting it as a restart of every child it will bring back.
521    ///
522    /// The children about to be re-registered each get a new process and a new start time, so counting the restart
523    /// against every one of them is what keeps their restart counts consistent with the rest of their own record.
524    /// Children that a group restart doesn't bring back keep their tombstones.
525    pub(super) fn clear_for_group_restart(&mut self) {
526        self.live.clear();
527
528        let mut state = self.node.state();
529        state.restarts_performed += 1;
530
531        let stopped_at = Stamp::now();
532        let NodeInner { children, history, .. } = &mut *state;
533        for record in children.values() {
534            if record.key.is_dynamic() {
535                // A dynamic child is not restored, and its roster id is never handed out again, so anything keyed to
536                // it is unreachable from here on rather than merely stale.
537                history.remove(&record.key);
538                continue;
539            }
540
541            // Exactly the children `respawn_children_one_for_all` brings back: every statically declared child that
542            // isn't temporary, whatever state it was last in. That includes a transient child which had already
543            // exited cleanly, since a group restart starts it afresh rather than honoring the rule that governs its
544            // own termination.
545            if record.restart != RestartType::Temporary {
546                if let Some(history) = history.get_mut(&record.key) {
547                    history.restarts += 1;
548                }
549            }
550        }
551
552        // Keep the tombstones of statically declared children a group restart won't bring back, since they remain
553        // part of the tree's declared shape. Everything else is about to be re-registered under a fresh id, and
554        // keeping the old record too would show each of those children twice.
555        //
556        // A temporary child still running is shut down along with the group and never brought back, so this is where
557        // it stops for good: record it as exited rather than dropping it, or it would vanish from the tree entirely.
558        state.children.retain(|_, record| {
559            if record.key.is_dynamic() || record.restart != RestartType::Temporary {
560                return false;
561            }
562
563            if record.exited.is_none() {
564                record.exited = Some(stopped_at);
565            }
566            true
567        });
568        state.adopted.clear();
569        drop(state);
570
571        self.debug_assert_consistent();
572    }
573
574    /// Checks that the supervisor's roster and the observable bookkeeping still agree.
575    ///
576    /// Every mutating method ends here, so the crate's existing supervision tests double as drift tests for this
577    /// module without needing to know it exists.
578    fn debug_assert_consistent(&self) {
579        debug_assert_eq!(
580            self.live.len(),
581            self.node.state().running_children(),
582            "supervisor '{}' roster and supervision-tree bookkeeping diverged",
583            self.node.id
584        );
585
586        saluki_antithesis::always_or_unreachable!(
587            self.live.len() == self.node.state().running_children(),
588            "supervision-tree bookkeeping tracks the supervisor's live roster"
589        );
590    }
591}
592
593/// A read-only handle for taking snapshots of a supervision tree.
594///
595/// Obtained from [`Supervisor::tree_handle`][super::Supervisor::tree_handle]. Cheap to clone, safe to share across
596/// tasks and threads, and deliberately incapable of anything but observation -- it cannot spawn children or
597/// otherwise affect the tree it reports on.
598///
599/// A handle can be taken before the supervisor starts, and remains valid across every restart of it.
600#[derive(Clone)]
601pub struct SupervisionTreeHandle {
602    node: Arc<SupervisorNode>,
603}
604
605impl SupervisionTreeHandle {
606    pub(super) fn new(node: Arc<SupervisorNode>) -> Self {
607        Self { node }
608    }
609
610    /// Returns the identifier of the supervisor at the root of this handle's view.
611    pub fn name(&self) -> &str {
612        &self.node.id
613    }
614
615    /// Renders the current state of the supervision tree as pretty-printed JSON.
616    ///
617    /// Returns a JSON object describing the failure if the tree can't be serialized, so that a caller writing a
618    /// diagnostic artifact always has something to write.
619    pub fn snapshot_json(&self) -> String {
620        match serde_json::to_string_pretty(&self.snapshot()) {
621            Ok(json) => json,
622            Err(e) => {
623                warn!(error = %e, "Failed to serialize supervision tree.");
624                String::from(r#"{"error": "failed to serialize supervision tree"}"#)
625            }
626        }
627    }
628
629    /// Creates a [`SupervisionTreeAPIHandler`] serving snapshots of this tree.
630    pub fn api_handler(&self) -> SupervisionTreeAPIHandler {
631        SupervisionTreeAPIHandler::from_handle(self.clone())
632    }
633
634    /// Creates a [`SupervisionTreeWorker`] that publishes this tree over the control plane.
635    pub fn worker(&self) -> SupervisionTreeWorker {
636        SupervisionTreeWorker::new(self.clone())
637    }
638
639    /// Captures the current state of the supervision tree.
640    ///
641    /// Descends from this handle's supervisor through every child, nesting each child's own children beneath it.
642    ///
643    /// # Consistency
644    ///
645    /// The tree is read one supervisor at a time, so a snapshot is not a globally atomic view: a child may start or
646    /// exit while the walk is in progress, and the result then mixes observations from slightly different instants.
647    /// This is deliberate. Reading the whole tree atomically would mean holding every supervisor's bookkeeping at
648    /// once, which would let an operator-facing diagnostics call stall supervision across the entire process -- a far
649    /// worse property than a snapshot whose nodes are microseconds apart.
650    pub fn snapshot(&self) -> TreeSnapshot {
651        // Read every resource group once, before touching any node. Doing it per node would take the global resource
652        // registry's lock once per node, each acquisition contending with the group registration that happens
653        // whenever any supervisor anywhere in the process starts a child.
654        let groups = collect_resource_groups();
655
656        let mut walk = Walk {
657            groups,
658            totals: TreeTotals::default(),
659            counted_groups: FastHashSet::default(),
660        };
661        let root = walk.supervisor(&self.node, ParentFacts::root(&self.node), 1);
662
663        TreeSnapshot {
664            captured_at: Stamp::now().wall_millis(),
665            resource_tracking_enabled: ResourceGroupRegistry::allocator_installed(),
666            totals: walk.totals,
667            root,
668        }
669    }
670}
671
672fn collect_resource_groups() -> FastHashMap<Arc<str>, ResourceUsage> {
673    let mut groups = FastHashMap::default();
674
675    ResourceGroupRegistry::global().visit_resource_groups(|name, stats| {
676        groups.insert(
677            Arc::<str>::from(name),
678            ResourceUsage::from(&stats.snapshot_delta(&ResourceStatsSnapshot::empty())),
679        );
680    });
681
682    groups
683}
684
685/// What a node's parent contributes to its snapshot.
686///
687/// A node's own bookkeeping knows what it is doing; its parent knows how it was configured, when it was created, and
688/// how many times it has been restarted. Both halves are needed, and only the parent has the second one -- except at
689/// the root, which has no parent.
690#[derive(Clone)]
691struct ParentFacts {
692    name: Arc<str>,
693    restart: RestartType,
694    significant: bool,
695    created: Stamp,
696    restarts: u32,
697    exited: Option<Stamp>,
698}
699
700impl ParentFacts {
701    /// Facts for a node with no parent: the root of the walk, or a supervisor a worker attached to the tree itself.
702    fn root(node: &SupervisorNode) -> Self {
703        Self {
704            name: Arc::clone(&node.id),
705            restart: RestartType::default(),
706            significant: false,
707            created: node.created,
708            restarts: 0,
709            exited: None,
710        }
711    }
712}
713
714struct RunFacts {
715    process_id: Option<ProcessId>,
716    process_name: Option<Arc<str>>,
717    started: Option<Stamp>,
718    state: NodeState,
719    /// The group this node's allocations are attributed to, which for a worker is its supervisor's rather than its
720    /// own.
721    resource_group: Option<Arc<str>>,
722    resources: Option<ResourceUsage>,
723}
724
725/// A node its parent has seen exit is stopped whatever else is true; otherwise `running` decides.
726///
727/// Each caller supplies `running` from what it actually knows, which is not the same thing in both cases. A
728/// supervisor's own bookkeeping says whether its run is in progress, and must be asked: it keeps the identity of its
729/// most recent run after that run ends, so reading the identity would say "running" forever. A worker has no such
730/// bookkeeping, but its record is created at the moment it starts and keeps an exit stamp once it stops, so for a
731/// worker having an identity and no exit stamp is the same statement.
732fn node_state(exited: Option<Stamp>, running: bool) -> NodeState {
733    match (exited, running) {
734        (Some(_), _) => NodeState::Exited,
735        (None, true) => NodeState::Running,
736        (None, false) => NodeState::Registered,
737    }
738}
739
740struct Walk {
741    groups: FastHashMap<Arc<str>, ResourceUsage>,
742    totals: TreeTotals,
743    /// Resource groups already folded into the totals.
744    ///
745    /// Group registration is idempotent by name, so two supervisors whose fully qualified names coincide share one
746    /// group. Totalling per node would then count the same bytes twice.
747    counted_groups: FastHashSet<Arc<str>>,
748}
749
750impl Walk {
751    fn supervisor(&mut self, node: &Arc<SupervisorNode>, parent: ParentFacts, depth: usize) -> NodeSnapshot {
752        // Copy out everything needed, then release the lock before descending: holding a parent's lock while taking a
753        // child's is the only way this walk could ever deadlock against a running supervisor, and not doing it is
754        // simpler to guarantee than any ordering rule.
755        let (run, running, config, generation, restarts_performed, children) = {
756            let state = node.state();
757
758            let mut children = state
759                .children
760                .iter()
761                .map(|(id, record)| {
762                    let history = state.history.get(&record.key);
763                    PendingChild {
764                        key: record.key,
765                        facts: ParentFacts {
766                            name: Arc::clone(&record.name),
767                            restart: record.restart,
768                            significant: record.significant,
769                            created: history
770                                .map(|h| h.created)
771                                .or_else(|| record.start.as_ref().map(|start| start.at))
772                                .unwrap_or_else(Stamp::now),
773                            restarts: history.map(|h| h.restarts).unwrap_or(0),
774                            exited: record.exited,
775                        },
776                        kind: match &record.kind {
777                            RecordKind::Supervisor(node) => PendingKind::Supervisor(Arc::clone(node)),
778                            RecordKind::Worker => match state.adopted.get(id) {
779                                Some(adopted) => PendingKind::WorkerDriving(Arc::clone(adopted)),
780                                None => PendingKind::Worker,
781                            },
782                        },
783                        start: record.start.clone(),
784                    }
785                })
786                .collect::<Vec<_>>();
787
788            // A restart re-registers a child under a fresh id, so insertion order stops matching declaration order as
789            // soon as anything restarts. `ChildKey` orders statics by declaration and dynamics by spawn, which keeps
790            // successive snapshots diffable.
791            children.sort_by_key(|child| child.key);
792
793            (
794                state.run.clone(),
795                state.running,
796                state.config,
797                state.generation,
798                state.restarts_performed,
799                children,
800            )
801        };
802
803        let (process_id, process_name, started) = split_run(run.as_ref());
804
805        // Taken from `running` rather than from whether an identity is present: the identity is that of the most
806        // recent run and outlives it, so deriving the state from it would have every supervisor that has ever run
807        // report itself as running forever.
808        let state = node_state(parent.exited, running);
809        let resources = process_name.as_ref().and_then(|name| self.groups.get(name).cloned());
810
811        // A worker inherits its supervisor's resource group rather than owning one, so this is what the children below
812        // are attributed to.
813        let worker_group = process_name.clone();
814
815        let mut snapshot = self.node_snapshot(
816            NodeKind::Supervisor,
817            parent,
818            RunFacts {
819                process_id,
820                // A supervisor owns its resource group, and it is named by the process it actually runs as.
821                resource_group: process_name.clone(),
822                process_name,
823                started,
824                state,
825                resources,
826            },
827            depth,
828        );
829
830        snapshot.supervision = Some(SupervisionSettings {
831            restart_mode: config.restart_strategy.mode(),
832            restart_intensity: config.restart_strategy.intensity(),
833            restart_period_ms: duration_millis(config.restart_strategy.period()),
834            auto_shutdown: config.auto_shutdown,
835            shutdown_budget_ms: config.shutdown_budget.map(duration_millis),
836            dedicated_threads: config.dedicated_threads,
837            restarts_performed,
838            generation,
839        });
840
841        if children.is_empty() || !self.may_descend(&node.id, depth) {
842            return snapshot;
843        }
844
845        snapshot.children.reserve(children.len());
846        for child in children {
847            snapshot
848                .children
849                .push(self.child(child, worker_group.as_ref(), depth + 1));
850        }
851
852        snapshot
853    }
854
855    fn child(&mut self, child: PendingChild, worker_group: Option<&Arc<str>>, depth: usize) -> NodeSnapshot {
856        match child.kind {
857            // A nested supervisor reports its own process and its own children; its parent only contributes how it
858            // was configured and how it has fared.
859            PendingKind::Supervisor(node) => self.supervisor(&node, child.facts, depth),
860
861            // A worker that turned out to be driving a supervisor of its own. The worker is what the parent
862            // supervises, so it stays the node; the supervisor it drives hangs beneath it.
863            PendingKind::WorkerDriving(adopted) => {
864                let mut snapshot = self.worker(child.facts, child.start, worker_group, depth);
865                if self.may_descend(&adopted.id, depth) {
866                    let nested = self.supervisor(&adopted, ParentFacts::root(&adopted), depth + 1);
867                    snapshot.children.push(nested);
868                }
869                snapshot
870            }
871
872            PendingKind::Worker => self.worker(child.facts, child.start, worker_group, depth),
873        }
874    }
875
876    fn worker(
877        &mut self, facts: ParentFacts, start: Option<StartedChild>, worker_group: Option<&Arc<str>>, depth: usize,
878    ) -> NodeSnapshot {
879        let (process_id, process_name, started) = split_run(start.as_ref());
880        let state = node_state(facts.exited, process_id.is_some());
881
882        self.node_snapshot(
883            NodeKind::Worker,
884            facts,
885            RunFacts {
886                process_id,
887                process_name,
888                started,
889                state,
890                resource_group: worker_group.cloned(),
891                // A worker owns no resource group: its allocations are counted against the supervisor named above.
892                // Reporting the supervisor's totals against each of its workers would count them many times over.
893                resources: None,
894            },
895            depth,
896        )
897    }
898
899    /// Assembles a node from its two halves, and folds it into the tree totals.
900    fn node_snapshot(&mut self, kind: NodeKind, parent: ParentFacts, run: RunFacts, depth: usize) -> NodeSnapshot {
901        self.totals.max_depth = self.totals.max_depth.max(depth);
902
903        match kind {
904            NodeKind::Supervisor => self.totals.supervisors += 1,
905            NodeKind::Worker => self.totals.workers += 1,
906        }
907        match run.state {
908            NodeState::Running => self.totals.running += 1,
909            NodeState::Exited => self.totals.exited += 1,
910            NodeState::Registered => self.totals.registered += 1,
911        }
912        self.totals.restarts += u64::from(parent.restarts);
913
914        // Several nodes can share one resource group, so fold each group in once rather than once per node.
915        if let (Some(usage), Some(group)) = (&run.resources, &run.process_name) {
916            if self.counted_groups.insert(Arc::clone(group)) {
917                self.totals.live_bytes += usage.live_bytes;
918                self.totals.cpu_time_nanos += usage.cpu_time_nanos;
919            }
920        }
921
922        NodeSnapshot {
923            name: parent.name.to_string(),
924            kind,
925            process_name: run.process_name.as_deref().map(str::to_string),
926            process_id: run.process_id.map(|id| id.as_usize() as u64),
927            state: run.state,
928            restart: parent.restart,
929            significant: parent.significant,
930            created_at: parent.created.wall_millis(),
931            started_at: run.started.map(|started| started.wall_millis()),
932            uptime_ms: match run.state {
933                NodeState::Running => run.started.map(|started| started.elapsed_millis()),
934                _ => None,
935            },
936            restart_count: parent.restarts,
937            exited_at: parent.exited.map(|exited| exited.wall_millis()),
938            resource_group: run.resource_group.as_deref().map(str::to_string),
939            resources: run.resources,
940            supervision: None,
941            children: Vec::new(),
942        }
943    }
944
945    fn may_descend(&self, supervisor_id: &str, depth: usize) -> bool {
946        if depth < MAX_TREE_DEPTH {
947            return true;
948        }
949
950        warn!(
951            supervisor_id,
952            depth, "Supervision tree is deeper than the snapshot walk descends; subtree omitted."
953        );
954        false
955    }
956}
957
958fn split_run(started: Option<&StartedChild>) -> (Option<ProcessId>, Option<Arc<str>>, Option<Stamp>) {
959    match started {
960        Some(start) => (
961            Some(start.process_id),
962            Some(Arc::clone(&start.process_name)),
963            Some(start.at),
964        ),
965        None => (None, None, None),
966    }
967}
968
969/// A child copied out from under its supervisor's lock, ready to be walked.
970struct PendingChild {
971    key: ChildKey,
972    facts: ParentFacts,
973    kind: PendingKind,
974    start: Option<StartedChild>,
975}
976
977enum PendingKind {
978    Worker,
979    /// A worker that drives a supervisor of its own, which hangs beneath it.
980    WorkerDriving(Arc<SupervisorNode>),
981    Supervisor(Arc<SupervisorNode>),
982}
983
984/// Converts a duration to whole milliseconds, saturating rather than overflowing.
985///
986/// Both a shutdown budget and a restart period come from configuration and can be arbitrarily large, up to and
987/// including [`Duration::MAX`], which is used elsewhere to mean `no deadline`.
988fn duration_millis(duration: Duration) -> u64 {
989    duration.as_millis().min(u128::from(u64::MAX)) as u64
990}
991
992#[cfg(test)]
993mod tests {
994    use super::*;
995
996    fn test_node(id: &str) -> Arc<SupervisorNode> {
997        Arc::new(SupervisorNode::new(Arc::from(id)))
998    }
999
1000    fn test_process(name: &str) -> Process {
1001        Process::supervisor(name, None).expect("valid process name")
1002    }
1003
1004    fn test_started(name: &str) -> StartedChild {
1005        StartedChild::new(&test_process(name), Arc::from(name))
1006    }
1007
1008    fn test_facts(key: ChildKey, name: &str, restart: RestartType) -> ChildFacts {
1009        ChildFacts {
1010            key,
1011            name: Arc::from(name),
1012            node: None,
1013            restart,
1014            significant: false,
1015        }
1016    }
1017
1018    fn restarts_of(node: &SupervisorNode, key: ChildKey) -> Option<u32> {
1019        node.state().history.get(&key).map(|history| history.restarts)
1020    }
1021
1022    /// Whether `key` still has a record, and if so whether it has been recorded as exited.
1023    fn record_of(node: &SupervisorNode, key: ChildKey) -> Option<bool> {
1024        node.state()
1025            .children
1026            .values()
1027            .find(|record| record.key == key)
1028            .map(|record| record.exited.is_some())
1029    }
1030
1031    #[test]
1032    fn group_restart_tombstones_a_temporary_child_that_was_still_running() {
1033        let node = test_node("sup");
1034        let mut roster = Roster::new(Arc::clone(&node));
1035        roster.insert(
1036            0,
1037            (),
1038            test_facts(ChildKey::Static(0), "one-shot", RestartType::Temporary),
1039            test_started("sup.one_shot"),
1040        );
1041        roster.insert(
1042            1,
1043            (),
1044            test_facts(ChildKey::Static(1), "flapper", RestartType::Permanent),
1045            test_started("sup.flapper"),
1046        );
1047
1048        roster.clear_for_group_restart();
1049
1050        // A group restart shuts the temporary child down and never brings it back, so this is where it stops for
1051        // good. It was declared, so it stays part of the tree's shape as a tombstone rather than vanishing.
1052        assert_eq!(
1053            record_of(&node, ChildKey::Static(0)),
1054            Some(true),
1055            "a temporary child stopped by a group restart is retained, and records that it stopped"
1056        );
1057        assert_eq!(
1058            restarts_of(&node, ChildKey::Static(0)),
1059            Some(0),
1060            "a child that is not brought back has not been restarted"
1061        );
1062
1063        // The permanent child is about to be re-registered under a fresh id, so keeping the old record too would
1064        // show it twice -- and it really is being restarted.
1065        assert_eq!(record_of(&node, ChildKey::Static(1)), None);
1066        assert_eq!(restarts_of(&node, ChildKey::Static(1)), Some(1));
1067    }
1068
1069    #[test]
1070    fn group_restart_counts_an_exited_transient_child_it_brings_back() {
1071        let node = test_node("sup");
1072        let mut roster = Roster::new(Arc::clone(&node));
1073        roster.insert(
1074            0,
1075            (),
1076            test_facts(ChildKey::Static(0), "transient", RestartType::Transient),
1077            test_started("sup.transient"),
1078        );
1079
1080        // A clean exit leaves a transient child stopped, as a tombstone.
1081        roster.remove(0);
1082        assert_eq!(record_of(&node, ChildKey::Static(0)), Some(true));
1083
1084        roster.clear_for_group_restart();
1085
1086        // A group restart starts every non-temporary static child afresh, including one that had already exited
1087        // cleanly, so the restart is counted against it and its tombstone gives way to the new registration.
1088        assert_eq!(restarts_of(&node, ChildKey::Static(0)), Some(1));
1089        assert_eq!(record_of(&node, ChildKey::Static(0)), None);
1090    }
1091
1092    #[test]
1093    fn group_restart_forgets_dynamic_children_entirely() {
1094        let node = test_node("sup");
1095        let mut roster = Roster::new(Arc::clone(&node));
1096        roster.insert(
1097            0,
1098            (),
1099            test_facts(ChildKey::Static(0), "static", RestartType::Permanent),
1100            test_started("sup.static"),
1101        );
1102        roster.insert(
1103            7,
1104            (),
1105            test_facts(ChildKey::Dynamic(7), "ephemeral", RestartType::Temporary),
1106            test_started("sup.ephemeral"),
1107        );
1108
1109        roster.clear_for_group_restart();
1110
1111        // A dynamic child is not restored, and its roster id is never handed out again, so nothing keyed to it can
1112        // ever be looked up again. Keeping it would strand one entry per dynamic child alive at every restart.
1113        let state = node.state();
1114        assert!(!state.history.contains_key(&ChildKey::Dynamic(7)));
1115        assert!(state.history.contains_key(&ChildKey::Static(0)));
1116    }
1117
1118    #[test]
1119    fn a_new_generation_counts_as_a_restart_of_every_static_child() {
1120        let node = test_node("sup");
1121        node.begin_run(&test_process("sup"));
1122
1123        let mut roster = Roster::new(Arc::clone(&node));
1124        roster.insert(
1125            0,
1126            (),
1127            test_facts(ChildKey::Static(0), "worker", RestartType::Permanent),
1128            test_started("sup.worker"),
1129        );
1130        roster.insert(
1131            5,
1132            (),
1133            test_facts(ChildKey::Dynamic(5), "ephemeral", RestartType::Temporary),
1134            test_started("sup.ephemeral"),
1135        );
1136        assert_eq!(restarts_of(&node, ChildKey::Static(0)), Some(0));
1137        drop(roster);
1138
1139        // The supervisor itself stops and is restarted from above. Its children go down with it, and every static one
1140        // is re-registered under a fresh process by the next generation, which is a restart of each of them.
1141        node.end_run();
1142        node.begin_run(&test_process("sup"));
1143
1144        assert_eq!(restarts_of(&node, ChildKey::Static(0)), Some(1));
1145        assert!(
1146            !node.state().history.contains_key(&ChildKey::Dynamic(5)),
1147            "a dynamic child is never restored, so its history cannot be looked up again"
1148        );
1149    }
1150
1151    #[test]
1152    fn a_supervisor_keeps_the_identity_of_its_most_recent_run() {
1153        let node = test_node("sup");
1154        let handle = SupervisionTreeHandle::new(Arc::clone(&node));
1155
1156        // Never run: there is no identity to report.
1157        let before = handle.snapshot();
1158        assert_eq!(before.root.state, NodeState::Registered);
1159        assert_eq!(before.root.process_id, None);
1160        assert_eq!(before.root.process_name, None);
1161        assert_eq!(before.root.started_at, None);
1162
1163        node.begin_run(&test_process("sup"));
1164        let running = handle.snapshot();
1165        assert_eq!(running.root.state, NodeState::Running);
1166        assert!(running.root.process_id.is_some());
1167        assert!(running.root.uptime_ms.is_some());
1168
1169        // Stopped, but it still says what it ran as -- as an exited worker always has. `state` carries the difference
1170        // between "not running" and "never ran", so the identity does not have to.
1171        node.end_run();
1172        let stopped = handle.snapshot();
1173        assert_eq!(stopped.root.state, NodeState::Registered);
1174        assert_eq!(stopped.root.process_id, running.root.process_id);
1175        assert_eq!(stopped.root.process_name, running.root.process_name);
1176        assert_eq!(stopped.root.started_at, running.root.started_at);
1177        assert_eq!(stopped.root.uptime_ms, None, "a node that isn't running has no uptime");
1178    }
1179}