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}