1use std::{
2 future::Future,
3 pin::Pin,
4 sync::{
5 atomic::{AtomicU64, AtomicUsize, Ordering},
6 Arc, Mutex,
7 },
8 time::Duration,
9};
10
11use saluki_common::sync::shutdown::ShutdownHandle;
12use saluki_error::GenericError;
13use serde::{Deserialize, Serialize};
14use snafu::{OptionExt as _, Snafu};
15use tokio::{pin, runtime::Handle, select, sync::mpsc};
16use tracing::{debug, error, warn};
17
18use super::{
19 dedicated::{spawn_dedicated_runtime, RuntimeConfiguration, RuntimeMode},
20 restart::{RestartAction, RestartMode, RestartState, RestartStrategy, RestartType},
21 supervisable::{InitializationError, ShutdownStrategy, Supervisable},
22 tree::{ChildFacts, ChildKey, NodeConfig, Roster, SupervisionTreeHandle, SupervisorNode},
23 worker_state::WorkerState,
24};
25use crate::runtime::{
26 process::{Process, ProcessExt as _},
27 state::DataspaceRegistry,
28};
29
30const UNNAMED_CHILD: &str = "unnamed";
34
35pub(super) type WorkerFuture = Pin<Box<dyn Future<Output = Result<(), WorkerError>> + Send>>;
41
42#[derive(Debug)]
47pub(super) enum WorkerError {
48 Initialization {
54 child_name: Option<String>,
55 source: InitializationError,
56 },
57
58 Runtime(GenericError),
60
61 ShutdownTimedOut {
68 aborted: usize,
70 },
71}
72
73impl From<SupervisorError> for WorkerError {
74 fn from(err: SupervisorError) -> Self {
75 match err {
76 SupervisorError::FailedToInitialize { child_name, source } => WorkerError::Initialization {
79 child_name: Some(child_name),
80 source,
81 },
82 SupervisorError::ShutdownTimedOut { aborted } => WorkerError::ShutdownTimedOut { aborted },
84 other => WorkerError::Runtime(other.into()),
86 }
87 }
88}
89
90#[derive(Debug, Snafu)]
92pub enum ProcessError {
93 #[snafu(display("Child process was aborted by the supervisor."))]
95 Aborted,
96
97 #[snafu(display("Child process panicked."))]
99 Panicked,
100
101 #[snafu(display("Child process terminated with an error: {}", source))]
103 Terminated {
104 source: GenericError,
106 },
107}
108
109#[derive(Clone, Copy, Debug, PartialEq, Eq, Default, Deserialize, Serialize)]
116#[serde(rename_all = "snake_case")]
117pub enum AutoShutdown {
118 #[default]
120 Never,
121
122 AnySignificant,
124
125 AllSignificant,
127}
128
129#[derive(Debug, Snafu)]
131#[snafu(context(suffix(false)))]
132pub enum SupervisorError {
133 #[snafu(display("Invalid name for supervisor or worker: '{}'", name))]
135 InvalidName {
136 name: String,
138 },
139
140 #[snafu(display("Child process '{}' failed to initialize: {}", child_name, source))]
145 FailedToInitialize {
146 child_name: String,
148
149 source: InitializationError,
151 },
152
153 #[snafu(display("Supervisor has exceeded restart limits and was forced to shutdown."))]
155 Shutdown,
156
157 #[snafu(display("Supervisor shut down after a significant child terminated."))]
162 SignificantChildExited,
163
164 #[snafu(display(
173 "Shutdown completed uncleanly: {} worker(s) were forcefully aborted after exceeding their shutdown timeout.",
174 aborted
175 ))]
176 ShutdownTimedOut {
177 aborted: usize,
179 },
180}
181
182pub struct ChildSpecification<S = WorkerSpec> {
202 spec_inner: S,
203}
204
205pub struct WorkerSpec {
207 worker: Arc<dyn Supervisable>,
208 options: ChildOptions,
209}
210
211pub struct SupervisorSpec {
213 supervisor: Supervisor,
214 options: ChildOptions,
215}
216
217impl ChildSpecification<WorkerSpec> {
223 pub(crate) fn worker<T: Supervisable + 'static>(worker: T) -> Self {
225 Self {
226 spec_inner: WorkerSpec {
227 worker: Arc::new(worker),
228 options: ChildOptions::default(),
229 },
230 }
231 }
232
233 pub(crate) fn one_shot_worker<T: Supervisable + 'static>(worker: T) -> Self {
238 Self::worker(worker).with_restart_type(RestartType::Temporary)
239 }
240
241 #[must_use]
247 pub(crate) fn with_restart_type(mut self, restart_type: RestartType) -> Self {
248 self.spec_inner.options.restart = Some(restart_type);
249 self
250 }
251
252 #[must_use]
258 pub(crate) fn with_significant(mut self, significant: bool) -> Self {
259 self.spec_inner.options.significant = significant;
260 self
261 }
262
263 #[must_use]
272 pub(crate) fn with_runtime(mut self, handle: Handle) -> Self {
273 self.spec_inner.options.runtime = Some(handle);
274 self
275 }
276
277 #[must_use]
284 pub(crate) fn with_shutdown_strategy(mut self, strategy: ShutdownStrategy) -> Self {
285 self.spec_inner.options.shutdown = ChildShutdown::Explicit(strategy);
286 self
287 }
288
289 #[must_use]
300 pub(crate) fn with_budget_bounded_shutdown(mut self) -> Self {
301 self.spec_inner.options.shutdown = ChildShutdown::BudgetBounded;
302 self
303 }
304}
305
306impl ChildSpecification<SupervisorSpec> {
315 #[must_use]
321 pub(crate) fn with_restart_type(mut self, restart_type: RestartType) -> Self {
322 self.spec_inner.options.restart = Some(restart_type);
323 self
324 }
325
326 #[must_use]
331 pub(crate) fn with_significant(mut self, significant: bool) -> Self {
332 self.spec_inner.options.significant = significant;
333 self
334 }
335}
336
337impl<T> From<T> for ChildSpecification<WorkerSpec>
338where
339 T: Supervisable + 'static,
340{
341 fn from(worker: T) -> Self {
342 Self::worker(worker)
343 }
344}
345
346impl From<Supervisor> for ChildSpecification<SupervisorSpec> {
347 fn from(supervisor: Supervisor) -> Self {
348 Self {
349 spec_inner: SupervisorSpec {
350 supervisor,
351 options: ChildOptions::default(),
352 },
353 }
354 }
355}
356
357mod sealed {
358 pub trait Sealed {}
359}
360
361impl sealed::Sealed for WorkerSpec {}
362impl sealed::Sealed for SupervisorSpec {}
363
364pub trait ChildState: sealed::Sealed + Sized {
371 #[doc(hidden)]
376 fn into_child_parts(spec: ChildSpecification<Self>, default_restart: RestartType) -> LoweredChild;
377}
378
379pub struct LoweredChild {
384 spec: SupervisedChild,
385 config: ChildConfig,
386}
387
388impl ChildState for WorkerSpec {
389 fn into_child_parts(spec: ChildSpecification<Self>, default_restart: RestartType) -> LoweredChild {
390 let WorkerSpec { worker, options } = spec.spec_inner;
391 LoweredChild {
392 spec: SupervisedChild::Worker(worker),
393 config: options.resolve(default_restart),
394 }
395 }
396}
397
398impl ChildState for SupervisorSpec {
399 fn into_child_parts(spec: ChildSpecification<Self>, default_restart: RestartType) -> LoweredChild {
400 let SupervisorSpec { supervisor, options } = spec.spec_inner;
401 LoweredChild {
402 spec: SupervisedChild::Supervisor(supervisor),
403 config: options.resolve(default_restart),
404 }
405 }
406}
407
408pub(super) enum SupervisedChild {
410 Worker(Arc<dyn Supervisable>),
411 Supervisor(Supervisor),
412}
413
414impl SupervisedChild {
415 pub(super) fn is_supervisor(&self) -> bool {
417 matches!(self, Self::Supervisor(_))
418 }
419
420 pub(super) fn node(&self) -> Option<Arc<SupervisorNode>> {
423 match self {
424 Self::Worker(_) => None,
425 Self::Supervisor(supervisor) => Some(Arc::clone(&supervisor.node)),
426 }
427 }
428
429 fn process_type(&self) -> &'static str {
430 match self {
431 Self::Worker(_) => "worker",
432 Self::Supervisor(_) => "supervisor",
433 }
434 }
435
436 fn name(&self) -> &str {
437 match self {
438 Self::Worker(worker) => worker.name(),
439 Self::Supervisor(supervisor) => &supervisor.supervisor_id,
440 }
441 }
442
443 pub(super) fn wants_shutdown_signal(&self) -> bool {
449 match self {
450 Self::Worker(worker) => worker.wants_shutdown_signal(),
451 Self::Supervisor(_) => true,
452 }
453 }
454
455 pub(super) fn shutdown_strategy(&self) -> ShutdownStrategy {
456 match self {
457 Self::Worker(worker) => worker.shutdown_strategy(),
458
459 Self::Supervisor(_) => ShutdownStrategy::Graceful(Duration::MAX),
462 }
463 }
464
465 pub(super) fn create_process(&self, parent_process: &Process) -> Process {
472 let name = self.name();
473 let process = match self {
474 Self::Worker(_) => Process::worker(name, parent_process),
475 Self::Supervisor(_) => Process::supervisor(name, Some(parent_process)),
476 };
477
478 process.unwrap_or_else(|| {
479 warn!(
480 parent_process = parent_process.name(),
481 child_name = name,
482 "Child process name is not usable as a process name; falling back to '{}'.",
483 UNNAMED_CHILD
484 );
485
486 match self {
487 Self::Worker(_) => Process::worker(UNNAMED_CHILD, parent_process),
488 Self::Supervisor(_) => Process::supervisor(UNNAMED_CHILD, Some(parent_process)),
489 }
490 .expect("placeholder child name is always a valid process name")
491 })
492 }
493
494 pub(super) fn create_worker_future(
495 &self, process: Process, process_shutdown: ShutdownHandle,
496 ) -> Result<WorkerFuture, SupervisorError> {
497 match self {
498 Self::Worker(worker) => {
499 let worker = Arc::clone(worker);
500 Ok(Box::pin(async move {
501 let run_future =
502 worker
503 .initialize(process_shutdown)
504 .await
505 .map_err(|source| WorkerError::Initialization {
506 child_name: None,
507 source,
508 })?;
509 run_future.await.map_err(WorkerError::Runtime)
510 }))
511 }
512 Self::Supervisor(sup) => {
513 match sup.runtime_mode() {
514 RuntimeMode::Ambient => {
515 Ok(sup.as_nested_process(process, process_shutdown))
517 }
518 RuntimeMode::Dedicated(config) => {
519 let child_name = sup.supervisor_id.to_string();
528 let dataspace = process.dataspace().clone();
529 let handle =
530 spawn_dedicated_runtime(sup.inner_clone(), config.clone(), process_shutdown, dataspace)
531 .map_err(|e| SupervisorError::FailedToInitialize {
532 child_name,
533 source: e.into(),
534 })?;
535
536 Ok(Box::pin(async move { handle.await.map_err(WorkerError::from) }))
537 }
538 }
539 }
540 }
541 }
542}
543
544impl Clone for SupervisedChild {
545 fn clone(&self) -> Self {
546 match self {
547 Self::Worker(worker) => Self::Worker(Arc::clone(worker)),
548 Self::Supervisor(supervisor) => Self::Supervisor(supervisor.inner_clone()),
549 }
550 }
551}
552
553#[derive(Clone, Copy, Debug, Default)]
555pub(super) enum ChildShutdown {
556 #[default]
558 Worker,
559
560 Explicit(ShutdownStrategy),
562
563 BudgetBounded,
568}
569
570#[derive(Clone, Debug, Default)]
577pub(super) struct ChildOptions {
578 restart: Option<RestartType>,
579 significant: bool,
580
581 runtime: Option<Handle>,
583
584 shutdown: ChildShutdown,
585}
586
587impl ChildOptions {
588 fn resolve(self, default_restart: RestartType) -> ChildConfig {
590 ChildConfig {
591 restart: self.restart.unwrap_or(default_restart),
592 significant: self.significant,
593 runtime: self.runtime,
594 shutdown: self.shutdown,
595 }
596 }
597}
598
599#[derive(Clone, Debug)]
602pub(super) struct ChildConfig {
603 restart: RestartType,
604 significant: bool,
605 runtime: Option<Handle>,
606 shutdown: ChildShutdown,
607}
608
609impl ChildConfig {
610 pub(super) fn runtime(&self) -> Option<&Handle> {
612 self.runtime.as_ref()
613 }
614
615 pub(super) fn shutdown(&self) -> ChildShutdown {
617 self.shutdown
618 }
619}
620
621#[derive(Clone)]
623struct ChildEntry {
624 spec: SupervisedChild,
625 config: ChildConfig,
626 dynamic: bool,
629}
630
631#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
636pub struct ChildId(u64);
637
638impl ChildId {
639 pub const fn as_u64(self) -> u64 {
641 self.0
642 }
643}
644
645struct PendingSpawn {
647 id: u64,
648 spec: SupervisedChild,
649 config: ChildConfig,
650}
651
652const SPAWN_DRAIN_BATCH: usize = 64;
658
659#[derive(Clone)]
675pub struct SupervisorHandle {
676 name: Arc<str>,
677 current_tx: Arc<Mutex<Option<mpsc::UnboundedSender<PendingSpawn>>>>,
680 id_counter: Arc<AtomicU64>,
681 active: Arc<AtomicUsize>,
682}
683
684impl SupervisorHandle {
685 pub fn name(&self) -> &str {
687 &self.name
688 }
689
690 pub fn spawn<S, T>(&self, child: T) -> ChildId
705 where
706 S: ChildState,
707 T: Into<ChildSpecification<S>>,
708 {
709 let LoweredChild { spec, config } = S::into_child_parts(child.into(), RestartType::Temporary);
710
711 let id = self.id_counter.fetch_add(1, Ordering::Relaxed);
714 let pending = PendingSpawn { id, spec, config };
715
716 let tx = self.current_tx.lock().unwrap().clone();
723 match tx {
724 Some(tx) => {
727 if let Err(e) = tx.send(pending) {
728 debug!(
729 supervisor_id = %self.name,
730 child_name = e.0.spec.name(),
731 "Supervisor is shutting down; dynamic child will not be started."
732 );
733 }
734 }
735 None => warn!(
738 supervisor_id = %self.name,
739 child_name = pending.spec.name(),
740 "Supervisor is not running; dynamic child will not be started."
741 ),
742 }
743
744 ChildId(id)
745 }
746
747 pub fn is_running(&self) -> bool {
749 self.current_tx.lock().unwrap().is_some()
750 }
751
752 pub fn active_children(&self) -> usize {
757 self.active.load(Ordering::Relaxed)
758 }
759}
760
761pub struct Supervisor {
789 supervisor_id: Arc<str>,
790 child_specs: Vec<ChildEntry>,
791 runtime_mode: RuntimeMode,
792 current_tx: Arc<Mutex<Option<mpsc::UnboundedSender<PendingSpawn>>>>,
796 id_counter: Arc<AtomicU64>,
797 active: Arc<AtomicUsize>,
799 node: Arc<SupervisorNode>,
802}
803
804impl Supervisor {
805 pub fn new<S: AsRef<str>>(supervisor_id: S) -> Result<Self, SupervisorError> {
807 if supervisor_id.as_ref().is_empty() {
811 return Err(SupervisorError::InvalidName {
812 name: supervisor_id.as_ref().to_string(),
813 });
814 }
815
816 let supervisor_id: Arc<str> = supervisor_id.as_ref().into();
817
818 Ok(Self {
819 node: Arc::new(SupervisorNode::new(Arc::clone(&supervisor_id))),
820 supervisor_id,
821 child_specs: Vec::new(),
822 runtime_mode: RuntimeMode::default(),
823 current_tx: Arc::new(Mutex::new(None)),
824 id_counter: Arc::new(AtomicU64::new(0)),
825 active: Arc::new(AtomicUsize::new(0)),
826 })
827 }
828
829 pub fn id(&self) -> &str {
831 &self.supervisor_id
832 }
833
834 pub fn with_restart_strategy(self, strategy: RestartStrategy) -> Self {
836 self.node.update_config(|config| config.restart_strategy = strategy);
837 self
838 }
839
840 pub fn with_auto_shutdown(self, auto_shutdown: AutoShutdown) -> Self {
845 self.node.update_config(|config| config.auto_shutdown = auto_shutdown);
846 self
847 }
848
849 #[must_use]
872 pub fn with_shutdown_budget(self, budget: Duration) -> Self {
873 self.node.update_config(|config| config.shutdown_budget = Some(budget));
874 self
875 }
876
877 pub fn handle(&self) -> SupervisorHandle {
883 SupervisorHandle {
884 name: Arc::clone(&self.supervisor_id),
885 current_tx: Arc::clone(&self.current_tx),
886 id_counter: Arc::clone(&self.id_counter),
887 active: Arc::clone(&self.active),
888 }
889 }
890
891 pub fn tree_handle(&self) -> SupervisionTreeHandle {
896 SupervisionTreeHandle::new(Arc::clone(&self.node))
897 }
898
899 pub fn with_dedicated_runtime(mut self, config: RuntimeConfiguration) -> Self {
909 let worker_threads = config.worker_threads();
910 self.node
911 .update_config(|node_config| node_config.dedicated_threads = Some(worker_threads));
912 self.runtime_mode = RuntimeMode::Dedicated(config);
913 self
914 }
915
916 pub(crate) fn runtime_mode(&self) -> &RuntimeMode {
918 &self.runtime_mode
919 }
920
921 pub fn add_worker<S, T>(&mut self, child: T)
931 where
932 S: ChildState,
933 T: Into<ChildSpecification<S>>,
934 {
935 let LoweredChild { spec, config } = S::into_child_parts(child.into(), RestartType::Permanent);
936 self.push_child(ChildEntry {
937 spec,
938 config,
939 dynamic: false,
940 });
941 }
942
943 fn warn_if_significance_is_inert(&self, config: &ChildConfig, child_name: &str, auto_shutdown: AutoShutdown) {
957 if !config.significant {
958 return;
959 }
960
961 if config.restart == RestartType::Permanent {
962 warn!(
963 supervisor_id = %self.supervisor_id,
964 child_name,
965 "Child is marked significant but is permanent, so it is always restarted and the flag has no effect."
966 );
967 }
968
969 if auto_shutdown == AutoShutdown::Never {
970 warn!(
971 supervisor_id = %self.supervisor_id,
972 child_name,
973 "Child is marked significant but the supervisor's auto-shutdown policy is `Never`, so the flag has \
974 no effect."
975 );
976 }
977 }
978
979 fn push_child(&mut self, entry: ChildEntry) {
980 debug!(
981 supervisor_id = %self.supervisor_id,
982 "Adding new static child process #{}. ({}, {}, {:?})",
983 self.child_specs.len(),
984 entry.spec.process_type(),
985 entry.spec.name(),
986 entry.config,
987 );
988
989 debug_assert!(
994 !(entry.config.significant && entry.config.restart == RestartType::Permanent),
995 "child '{}' was marked significant but is permanent, so it is always restarted and its termination can \
996 never drive auto-shutdown",
997 entry.spec.name()
998 );
999
1000 self.child_specs.push(entry);
1001 }
1002
1003 fn child_facts(entry: &ChildEntry, key: ChildKey) -> ChildFacts {
1006 ChildFacts {
1007 key,
1008 name: entry.spec.name().into(),
1009 node: entry.spec.node(),
1010 restart: entry.config.restart,
1011 significant: entry.config.significant,
1012 }
1013 }
1014
1015 fn spawn_static_children(
1016 &self, roster: &mut Roster<ChildEntry>, worker_state: &mut WorkerState, auto_shutdown: AutoShutdown,
1017 ) -> Result<(), SupervisorError> {
1018 debug!(supervisor_id = %self.supervisor_id, "Spawning all static child processes.");
1019 for (index, entry) in self.child_specs.iter().enumerate() {
1020 self.warn_if_significance_is_inert(&entry.config, entry.spec.name(), auto_shutdown);
1021
1022 let id = self.id_counter.fetch_add(1, Ordering::Relaxed);
1023 let started = worker_state.add_worker(id, &entry.spec, &entry.config)?;
1024 let facts = Self::child_facts(entry, ChildKey::Static(index));
1025 roster.insert(id, entry.clone(), facts, started);
1026 }
1027
1028 Ok(())
1029 }
1030
1031 fn respawn_children_one_for_all(
1039 &self, roster: &mut Roster<ChildEntry>, worker_state: &mut WorkerState,
1040 ) -> Result<(), SupervisorError> {
1041 debug!(supervisor_id = %self.supervisor_id, "Restarting all eligible static child processes.");
1042 for (index, entry) in self.child_specs.iter().enumerate() {
1043 if entry.config.restart == RestartType::Temporary {
1046 continue;
1047 }
1048 let id = self.id_counter.fetch_add(1, Ordering::Relaxed);
1049 let started = worker_state.add_worker(id, &entry.spec, &entry.config)?;
1050 let facts = Self::child_facts(entry, ChildKey::Static(index));
1053 roster.insert(id, entry.clone(), facts, started);
1054 }
1055
1056 Ok(())
1057 }
1058
1059 fn spawn_dynamic_child(
1061 &self, spawn: PendingSpawn, worker_state: &mut WorkerState, roster: &mut Roster<ChildEntry>,
1062 significant_remaining: &mut usize, auto_shutdown: AutoShutdown,
1063 ) {
1064 let PendingSpawn { id, spec, config } = spawn;
1065 let entry = ChildEntry {
1066 spec,
1067 config,
1068 dynamic: true,
1069 };
1070 self.warn_if_significance_is_inert(&entry.config, entry.spec.name(), auto_shutdown);
1071
1072 match worker_state.add_worker(id, &entry.spec, &entry.config) {
1073 Ok(started) => {
1074 if entry.config.significant {
1075 *significant_remaining += 1;
1076 }
1077 self.active.fetch_add(1, Ordering::Relaxed);
1078 let facts = Self::child_facts(&entry, ChildKey::Dynamic(id));
1079 roster.insert(id, entry, facts, started);
1080 }
1081 Err(e) => {
1082 error!(
1086 supervisor_id = %self.supervisor_id,
1087 child_name = entry.spec.name(),
1088 error = %e,
1089 "Failed to start dynamic child."
1090 );
1091 }
1092 }
1093 }
1094
1095 async fn run_inner(&self, process: Process, process_shutdown: ShutdownHandle) -> Result<(), SupervisorError> {
1096 let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
1099 *self.current_tx.lock().unwrap() = Some(cmd_tx);
1100
1101 self.node.begin_run(&process);
1105
1106 let result = self.supervise(process, process_shutdown, cmd_rx).await;
1107
1108 *self.current_tx.lock().unwrap() = None;
1111 self.active.store(0, Ordering::Relaxed);
1112
1113 self.node.end_run();
1117
1118 result
1119 }
1120
1121 async fn supervise(
1122 &self, process: Process, process_shutdown: ShutdownHandle, mut cmd_rx: mpsc::UnboundedReceiver<PendingSpawn>,
1123 ) -> Result<(), SupervisorError> {
1124 let NodeConfig {
1127 restart_strategy,
1128 auto_shutdown,
1129 shutdown_budget,
1130 ..
1131 } = self.node.config();
1132
1133 let mut restart_state = RestartState::new(restart_strategy);
1134 let mut worker_state = WorkerState::new(process, self.handle(), shutdown_budget, Arc::clone(&self.node));
1135
1136 let mut roster = Roster::new(Arc::clone(&self.node));
1142
1143 self.spawn_static_children(&mut roster, &mut worker_state, auto_shutdown)?;
1146
1147 let mut significant_remaining = roster.values().filter(|entry| entry.config.significant).count();
1149
1150 let mut spawn_batch = Vec::with_capacity(SPAWN_DRAIN_BATCH);
1152
1153 pin!(process_shutdown);
1155
1156 let outcome = loop {
1157 select! {
1158 biased;
1161
1162 _ = &mut process_shutdown => break Ok(()),
1166
1167 (child_id, worker_result) = worker_state.wait_for_next_worker() => {
1171 let (child_name, config, dynamic) = {
1173 let entry = roster.get(child_id).expect("completed worker must be present in the roster");
1174 (entry.spec.name().to_string(), entry.config.clone(), entry.dynamic)
1175 };
1176
1177 if let Err(WorkerError::Initialization { child_name: inner, source }) = worker_result {
1179 let full_name = match inner {
1182 Some(inner) => format!("{}/{}", child_name, inner),
1183 None => child_name.clone(),
1184 };
1185
1186 error!(supervisor_id = %self.supervisor_id, worker_name = full_name, "Child process failed to initialize: {}", source);
1187 break Err(SupervisorError::FailedToInitialize { child_name: full_name, source });
1188 }
1189
1190 let abnormal = worker_result.is_err();
1193 let worker_result = worker_result.map_err(|e| match e {
1194 WorkerError::Runtime(e) => ProcessError::Terminated { source: e },
1195 WorkerError::Initialization { .. } => unreachable!("handled above"),
1196 WorkerError::ShutdownTimedOut { aborted } => ProcessError::Terminated {
1201 source: SupervisorError::ShutdownTimedOut { aborted }.into(),
1202 },
1203 });
1204
1205 if !config.restart.should_restart(abnormal) {
1206 if abnormal {
1215 warn!(supervisor_id = %self.supervisor_id, worker_name = %child_name, restart = ?config.restart, ?worker_result, "Child process exited with an error and is not eligible for restart.");
1216 } else {
1217 debug!(supervisor_id = %self.supervisor_id, worker_name = %child_name, restart = ?config.restart, "Child process exited and is not eligible for restart.");
1218 }
1219 roster.remove(child_id);
1220 if dynamic {
1221 self.active.fetch_sub(1, Ordering::Relaxed);
1222 }
1223
1224 if config.significant {
1228 significant_remaining = significant_remaining.saturating_sub(1);
1229 let auto_shutdown = match auto_shutdown {
1230 AutoShutdown::Never => false,
1231 AutoShutdown::AnySignificant => true,
1232 AutoShutdown::AllSignificant => significant_remaining == 0,
1233 };
1234 if auto_shutdown {
1235 warn!(supervisor_id = %self.supervisor_id, worker_name = %child_name, ?worker_result, "Significant child terminated; shutting down supervisor.");
1236 break Err(SupervisorError::SignificantChildExited);
1237 }
1238 }
1239 } else {
1240 match restart_state.evaluate_restart() {
1241 RestartAction::Restart(mode) => match mode {
1242 RestartMode::OneForOne => {
1243 warn!(supervisor_id = %self.supervisor_id, worker_name = %child_name, ?worker_result, "Child process terminated, restarting.");
1244 let spec = roster.get(child_id).expect("present for restart").spec.clone();
1245 match worker_state.add_worker(child_id, &spec, &config) {
1246 Ok(started) => roster.restart_in_place(child_id, started),
1247 Err(e) => break Err(e),
1248 }
1249 }
1250 RestartMode::OneForAll => {
1251 warn!(supervisor_id = %self.supervisor_id, worker_name = %child_name, ?worker_result, "Child process terminated, restarting all processes.");
1252 let _ = worker_state.shutdown_workers().await;
1256 roster.clear_for_group_restart();
1260 self.active.store(0, Ordering::Relaxed);
1261 let respawn = self.respawn_children_one_for_all(&mut roster, &mut worker_state);
1262 if let Err(e) = respawn {
1263 break Err(e);
1264 }
1265 significant_remaining =
1266 roster.values().filter(|entry| entry.config.significant).count();
1267 }
1268 },
1269 RestartAction::Shutdown => {
1270 error!(supervisor_id = %self.supervisor_id, worker_name = %child_name, ?worker_result, "Supervisor shutting down due to restart limits.");
1271 break Err(SupervisorError::Shutdown);
1272 }
1273 }
1274 }
1275 }
1276
1277 _ = cmd_rx.recv_many(&mut spawn_batch, SPAWN_DRAIN_BATCH) => {
1281 for spawn in spawn_batch.drain(..) {
1282 self.spawn_dynamic_child(
1283 spawn,
1284 &mut worker_state,
1285 &mut roster,
1286 &mut significant_remaining,
1287 auto_shutdown,
1288 );
1289 }
1290 }
1291 }
1292 };
1293
1294 cmd_rx.close();
1299 let mut discarded = 0;
1300 while cmd_rx.try_recv().is_ok() {
1301 discarded += 1;
1302 }
1303 if discarded > 0 {
1304 debug!(
1305 supervisor_id = %self.supervisor_id,
1306 discarded,
1307 "Discarded queued dynamic children during shutdown."
1308 );
1309 }
1310 let aborted = worker_state.shutdown_workers().await;
1311
1312 match outcome {
1317 Ok(()) if aborted > 0 => {
1318 warn!(supervisor_id = %self.supervisor_id, aborted, "Shutdown completed uncleanly; workers were forcefully aborted.");
1319 Err(SupervisorError::ShutdownTimedOut { aborted })
1320 }
1321 outcome => outcome,
1322 }
1323 }
1324
1325 fn as_nested_process(&self, process: Process, process_shutdown: ShutdownHandle) -> WorkerFuture {
1326 debug!(supervisor_id = %self.supervisor_id, "Nested supervisor starting.");
1329
1330 let sup = self.inner_clone();
1332
1333 Box::pin(async move {
1334 sup.run_inner(process, process_shutdown)
1335 .await
1336 .map_err(WorkerError::from)
1337 })
1338 }
1339
1340 pub async fn run(&mut self) -> Result<(), SupervisorError> {
1346 let process_shutdown = ShutdownHandle::noop();
1349 let process = Process::supervisor(&self.supervisor_id, None).context(InvalidName {
1350 name: self.supervisor_id.to_string(),
1351 })?;
1352
1353 debug!(supervisor_id = %self.supervisor_id, "Supervisor starting.");
1354 self.run_inner(process.clone(), process_shutdown)
1355 .into_process_future(process)
1356 .await
1357 }
1358
1359 pub async fn run_with_shutdown<F: Future + Send + 'static>(&mut self, shutdown: F) -> Result<(), SupervisorError> {
1368 let (shutdown_coordinator, shutdown_handle) = ShutdownHandle::paired();
1372 let run = self.run_with_shutdown_inner(shutdown_handle, None);
1373 pin!(run, shutdown);
1374
1375 let mut shutdown_coordinator = Some(shutdown_coordinator);
1376 loop {
1377 select! {
1378 result = &mut run => return result,
1379 _ = &mut shutdown, if shutdown_coordinator.is_some() => {
1380 shutdown_coordinator.take().expect("coordinator present per select guard").shutdown();
1381 }
1382 }
1383 }
1384 }
1385
1386 pub(crate) async fn run_with_shutdown_inner(
1398 &mut self, process_shutdown: ShutdownHandle, dataspace: Option<DataspaceRegistry>,
1399 ) -> Result<(), SupervisorError> {
1400 let process =
1401 Process::supervisor_with_dataspace(&self.supervisor_id, None, dataspace).context(InvalidName {
1402 name: self.supervisor_id.to_string(),
1403 })?;
1404
1405 debug!(supervisor_id = %self.supervisor_id, "Supervisor starting.");
1406 self.run_inner(process.clone(), process_shutdown)
1407 .into_process_future(process)
1408 .await
1409 }
1410
1411 fn inner_clone(&self) -> Self {
1412 Self {
1416 supervisor_id: Arc::clone(&self.supervisor_id),
1417 child_specs: self.child_specs.clone(),
1418 runtime_mode: self.runtime_mode.clone(),
1419 current_tx: Arc::clone(&self.current_tx),
1420 id_counter: Arc::clone(&self.id_counter),
1421 active: Arc::clone(&self.active),
1422 node: Arc::clone(&self.node),
1423 }
1424 }
1425}
1426
1427#[cfg(test)]
1428mod tests {
1429 use std::{
1430 future::pending,
1431 sync::atomic::{AtomicBool, AtomicUsize, Ordering},
1432 };
1433
1434 use async_trait::async_trait;
1435 use saluki_common::sync::shutdown::ShutdownCoordinator;
1436 use saluki_metrics::test::TestRecorder;
1437 use tokio::{
1438 sync::oneshot,
1439 task::JoinHandle,
1440 time::{sleep, timeout},
1441 };
1442
1443 use super::*;
1444 use crate::runtime::{self, FnWorker, NodeKind, NodeSnapshot, NodeState, SupervisorFuture};
1445 use crate::test_support::wait_until;
1446
1447 #[derive(Clone)]
1449 enum InitBehavior {
1450 Instant,
1452
1453 Slow(Duration),
1455
1456 Fail(&'static str),
1458 }
1459
1460 #[derive(Clone)]
1462 enum RunBehavior {
1463 UntilShutdown,
1465
1466 FailAfter(Duration, &'static str),
1468
1469 CompleteAfter(Duration),
1471
1472 SlowShutdown(Duration),
1474
1475 IgnoreShutdown,
1477
1478 PanicAfter(Duration),
1480 }
1481
1482 struct MockWorker {
1484 name: &'static str,
1485 init_behavior: InitBehavior,
1486 run_behavior: RunBehavior,
1487 start_count: Arc<AtomicUsize>,
1488 finish_count: Arc<AtomicUsize>,
1489 brutal_shutdown: bool,
1490 graceful_timeout: Duration,
1491 }
1492
1493 impl MockWorker {
1494 fn long_running(name: &'static str) -> Self {
1496 Self {
1497 name,
1498 init_behavior: InitBehavior::Instant,
1499 run_behavior: RunBehavior::UntilShutdown,
1500 start_count: Arc::new(AtomicUsize::new(0)),
1501 finish_count: Arc::new(AtomicUsize::new(0)),
1502 brutal_shutdown: false,
1503 graceful_timeout: Duration::from_millis(500),
1504 }
1505 }
1506
1507 fn failing(name: &'static str, delay: Duration) -> Self {
1509 Self {
1510 name,
1511 init_behavior: InitBehavior::Instant,
1512 run_behavior: RunBehavior::FailAfter(delay, "worker failed"),
1513 start_count: Arc::new(AtomicUsize::new(0)),
1514 finish_count: Arc::new(AtomicUsize::new(0)),
1515 brutal_shutdown: false,
1516 graceful_timeout: Duration::from_millis(500),
1517 }
1518 }
1519
1520 fn completing(name: &'static str, delay: Duration) -> Self {
1522 Self {
1523 name,
1524 init_behavior: InitBehavior::Instant,
1525 run_behavior: RunBehavior::CompleteAfter(delay),
1526 start_count: Arc::new(AtomicUsize::new(0)),
1527 finish_count: Arc::new(AtomicUsize::new(0)),
1528 brutal_shutdown: false,
1529 graceful_timeout: Duration::from_millis(500),
1530 }
1531 }
1532
1533 fn slow_shutdown(name: &'static str, delay: Duration) -> Self {
1535 Self {
1536 name,
1537 init_behavior: InitBehavior::Instant,
1538 run_behavior: RunBehavior::SlowShutdown(delay),
1539 start_count: Arc::new(AtomicUsize::new(0)),
1540 finish_count: Arc::new(AtomicUsize::new(0)),
1541 brutal_shutdown: false,
1542 graceful_timeout: Duration::from_millis(500),
1543 }
1544 }
1545
1546 fn ignore_shutdown(name: &'static str) -> Self {
1548 Self {
1549 name,
1550 init_behavior: InitBehavior::Instant,
1551 run_behavior: RunBehavior::IgnoreShutdown,
1552 start_count: Arc::new(AtomicUsize::new(0)),
1553 finish_count: Arc::new(AtomicUsize::new(0)),
1554 brutal_shutdown: false,
1555 graceful_timeout: Duration::from_millis(500),
1556 }
1557 }
1558
1559 fn panicking(name: &'static str, delay: Duration) -> Self {
1561 Self {
1562 name,
1563 init_behavior: InitBehavior::Instant,
1564 run_behavior: RunBehavior::PanicAfter(delay),
1565 start_count: Arc::new(AtomicUsize::new(0)),
1566 finish_count: Arc::new(AtomicUsize::new(0)),
1567 brutal_shutdown: false,
1568 graceful_timeout: Duration::from_millis(500),
1569 }
1570 }
1571
1572 fn init_failure(name: &'static str) -> Self {
1574 Self {
1575 name,
1576 init_behavior: InitBehavior::Fail("init failed"),
1577 run_behavior: RunBehavior::UntilShutdown,
1578 start_count: Arc::new(AtomicUsize::new(0)),
1579 finish_count: Arc::new(AtomicUsize::new(0)),
1580 brutal_shutdown: false,
1581 graceful_timeout: Duration::from_millis(500),
1582 }
1583 }
1584
1585 fn slow_init(name: &'static str, init_delay: Duration) -> Self {
1587 Self {
1588 name,
1589 init_behavior: InitBehavior::Slow(init_delay),
1590 run_behavior: RunBehavior::UntilShutdown,
1591 start_count: Arc::new(AtomicUsize::new(0)),
1592 finish_count: Arc::new(AtomicUsize::new(0)),
1593 brutal_shutdown: false,
1594 graceful_timeout: Duration::from_millis(500),
1595 }
1596 }
1597
1598 fn start_count(&self) -> Arc<AtomicUsize> {
1604 Arc::clone(&self.start_count)
1605 }
1606
1607 fn finish_count(&self) -> Arc<AtomicUsize> {
1615 Arc::clone(&self.finish_count)
1616 }
1617
1618 fn with_brutal_shutdown(mut self) -> Self {
1620 self.brutal_shutdown = true;
1621 self
1622 }
1623
1624 fn with_graceful_timeout(mut self, timeout: Duration) -> Self {
1626 self.graceful_timeout = timeout;
1627 self
1628 }
1629 }
1630
1631 #[async_trait]
1632 impl Supervisable for MockWorker {
1633 fn name(&self) -> &str {
1634 self.name
1635 }
1636
1637 fn shutdown_strategy(&self) -> ShutdownStrategy {
1638 if self.brutal_shutdown {
1639 ShutdownStrategy::Brutal
1640 } else {
1641 ShutdownStrategy::Graceful(self.graceful_timeout)
1642 }
1643 }
1644
1645 async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
1646 match &self.init_behavior {
1647 InitBehavior::Instant => {}
1648 InitBehavior::Slow(delay) => {
1649 sleep(*delay).await;
1650 }
1651 InitBehavior::Fail(msg) => {
1652 return Err(InitializationError::Failed {
1653 source: GenericError::msg(*msg),
1654 });
1655 }
1656 }
1657
1658 let start_count = Arc::clone(&self.start_count);
1659 let finish_count = Arc::clone(&self.finish_count);
1660 let run_behavior = self.run_behavior.clone();
1661
1662 Ok(Box::pin(async move {
1663 start_count.fetch_add(1, Ordering::SeqCst);
1664
1665 match run_behavior {
1666 RunBehavior::UntilShutdown => {
1667 process_shutdown.await;
1668 Ok(())
1669 }
1670 RunBehavior::FailAfter(delay, msg) => {
1671 select! {
1672 _ = sleep(delay) => {
1673 finish_count.fetch_add(1, Ordering::SeqCst);
1676 Err(GenericError::msg(msg))
1677 }
1678 _ = process_shutdown => {
1679 Ok(())
1680 }
1681 }
1682 }
1683 RunBehavior::CompleteAfter(delay) => {
1684 select! {
1685 _ = sleep(delay) => {
1686 finish_count.fetch_add(1, Ordering::SeqCst);
1688 Ok(())
1689 }
1690 _ = process_shutdown => Ok(()),
1691 }
1692 }
1693 RunBehavior::SlowShutdown(delay) => {
1694 process_shutdown.await;
1695 sleep(delay).await;
1696 finish_count.fetch_add(1, Ordering::SeqCst);
1698 Ok(())
1699 }
1700 RunBehavior::IgnoreShutdown => {
1701 let _hold = process_shutdown;
1703 pending().await
1704 }
1705 RunBehavior::PanicAfter(delay) => {
1706 select! {
1707 _ = sleep(delay) => panic!("worker panicked"),
1708 _ = process_shutdown => Ok(()),
1709 }
1710 }
1711 }
1712 }))
1713 }
1714 }
1715
1716 async fn run_supervisor_with_trigger(
1722 supervisor: Supervisor,
1723 ) -> (oneshot::Sender<()>, JoinHandle<Result<(), SupervisorError>>) {
1724 let sup_handle = supervisor.handle();
1726 let mut supervisor = supervisor;
1727
1728 let (tx, rx) = oneshot::channel();
1729 let handle = tokio::spawn(async move { supervisor.run_with_shutdown(rx).await });
1730
1731 wait_until("supervisor is running", || sup_handle.is_running()).await;
1732 (tx, handle)
1733 }
1734
1735 async fn join_supervisor(handle: JoinHandle<Result<(), SupervisorError>>) -> Result<(), SupervisorError> {
1740 timeout(Duration::from_secs(2), handle)
1741 .await
1742 .expect("supervisor should exit promptly")
1743 .expect("supervisor task should not panic")
1744 }
1745
1746 #[tokio::test]
1749 async fn standalone_supervisor_shuts_down_cleanly() {
1750 let mut sup = Supervisor::new("test-sup").unwrap();
1751 sup.add_worker(MockWorker::long_running("worker1"));
1752 sup.add_worker(MockWorker::long_running("worker2"));
1753
1754 let (tx, handle) = run_supervisor_with_trigger(sup).await;
1755 tx.send(()).unwrap();
1756
1757 let result = join_supervisor(handle).await;
1758 assert!(result.is_ok());
1759 }
1760
1761 #[tokio::test]
1762 async fn nested_supervisor_shuts_down_cleanly() {
1763 let mut child_sup = Supervisor::new("child-sup").unwrap();
1764 child_sup.add_worker(MockWorker::long_running("inner-worker"));
1765
1766 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
1767 parent_sup.add_worker(MockWorker::long_running("outer-worker"));
1768 parent_sup.add_worker(child_sup);
1769
1770 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
1771 tx.send(()).unwrap();
1772
1773 let result = join_supervisor(handle).await;
1774 assert!(result.is_ok());
1775 }
1776
1777 #[tokio::test]
1778 async fn empty_supervisor_idles_until_shutdown() {
1779 let sup = Supervisor::new("empty-sup").unwrap();
1782
1783 let (tx, handle) = run_supervisor_with_trigger(sup).await;
1784 assert!(!handle.is_finished(), "an empty supervisor must idle rather than exit");
1785
1786 tx.send(()).unwrap();
1787 let result = join_supervisor(handle).await;
1788 assert!(result.is_ok());
1789 }
1790
1791 #[tokio::test]
1794 async fn one_for_one_restarts_only_failed_child() {
1795 let failing = MockWorker::failing("failing-worker", Duration::from_millis(50));
1796 let failing_count = failing.start_count();
1797
1798 let stable = MockWorker::long_running("stable-worker");
1799 let stable_count = stable.start_count();
1800
1801 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1802 RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
1803 );
1804 sup.add_worker(stable);
1805 sup.add_worker(failing);
1806
1807 let (tx, handle) = run_supervisor_with_trigger(sup).await;
1808
1809 wait_until("the failing worker has been restarted", || {
1811 failing_count.load(Ordering::SeqCst) >= 2
1812 })
1813 .await;
1814 let _ = tx.send(());
1815
1816 let result = join_supervisor(handle).await;
1817 assert!(result.is_ok());
1818
1819 assert!(
1821 failing_count.load(Ordering::SeqCst) >= 2,
1822 "failing worker should have been restarted"
1823 );
1824 assert_eq!(
1826 stable_count.load(Ordering::SeqCst),
1827 1,
1828 "stable worker should not have been restarted"
1829 );
1830 }
1831
1832 #[tokio::test]
1833 async fn one_for_all_restarts_all_children() {
1834 let failing = MockWorker::failing("failing-worker", Duration::from_millis(50));
1835 let failing_count = failing.start_count();
1836
1837 let stable = MockWorker::long_running("stable-worker");
1838 let stable_count = stable.start_count();
1839
1840 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1841 RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
1842 );
1843 sup.add_worker(stable);
1844 sup.add_worker(failing);
1845
1846 let (tx, handle) = run_supervisor_with_trigger(sup).await;
1847
1848 wait_until("both workers have been restarted", || {
1850 failing_count.load(Ordering::SeqCst) >= 2 && stable_count.load(Ordering::SeqCst) >= 2
1851 })
1852 .await;
1853 let _ = tx.send(());
1854
1855 let result = join_supervisor(handle).await;
1856 assert!(result.is_ok());
1857
1858 assert!(
1860 failing_count.load(Ordering::SeqCst) >= 2,
1861 "failing worker should have been restarted"
1862 );
1863 assert!(
1864 stable_count.load(Ordering::SeqCst) >= 2,
1865 "stable worker should also have been restarted"
1866 );
1867 }
1868
1869 #[tokio::test]
1870 async fn one_for_all_does_not_restart_temporary_children() {
1871 let failing = MockWorker::failing("failing-worker", Duration::from_millis(50));
1874 let failing_count = failing.start_count();
1875
1876 let temp = MockWorker::long_running("temp-worker");
1877 let temp_count = temp.start_count();
1878
1879 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1880 RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
1881 );
1882 sup.add_worker(ChildSpecification::worker(temp).with_restart_type(RestartType::Temporary));
1883 sup.add_worker(failing);
1884
1885 let (tx, handle) = run_supervisor_with_trigger(sup).await;
1886
1887 wait_until("the permanent worker has been restarted", || {
1889 failing_count.load(Ordering::SeqCst) >= 2
1890 })
1891 .await;
1892 let _ = tx.send(());
1893
1894 let result = join_supervisor(handle).await;
1895 assert!(result.is_ok());
1896 assert!(
1897 failing_count.load(Ordering::SeqCst) >= 2,
1898 "permanent worker should have been restarted by one-for-all"
1899 );
1900 assert_eq!(
1901 temp_count.load(Ordering::SeqCst),
1902 1,
1903 "temporary child must not be restarted by a one-for-all group restart"
1904 );
1905 }
1906
1907 #[tokio::test]
1908 async fn one_for_all_restarts_transient_children() {
1909 let transient = MockWorker::completing("transient-worker", Duration::from_millis(30));
1912 let transient_count = transient.start_count();
1913
1914 let failing = MockWorker::failing("failing-worker", Duration::from_millis(80));
1916
1917 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1918 RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
1919 );
1920 sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
1921 sup.add_worker(failing);
1922
1923 let (tx, handle) = run_supervisor_with_trigger(sup).await;
1924
1925 wait_until("the transient worker has been restarted by the group", || {
1926 transient_count.load(Ordering::SeqCst) >= 2
1927 })
1928 .await;
1929 let _ = tx.send(());
1930
1931 let result = join_supervisor(handle).await;
1932 assert!(result.is_ok());
1933 assert!(
1934 transient_count.load(Ordering::SeqCst) >= 2,
1935 "transient child must be restarted by a one-for-all group restart, even after a clean exit"
1936 );
1937 }
1938
1939 #[tokio::test]
1940 async fn transient_abnormal_exit_triggers_one_for_all() {
1941 let transient = MockWorker::failing("transient-worker", Duration::from_millis(50));
1944 let transient_count = transient.start_count();
1945
1946 let stable = MockWorker::long_running("stable-worker");
1947 let stable_count = stable.start_count();
1948
1949 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
1950 RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
1951 );
1952 sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
1953 sup.add_worker(stable);
1954
1955 let (tx, handle) = run_supervisor_with_trigger(sup).await;
1956
1957 wait_until("the abnormal exit has restarted both workers", || {
1958 transient_count.load(Ordering::SeqCst) >= 2 && stable_count.load(Ordering::SeqCst) >= 2
1959 })
1960 .await;
1961 let _ = tx.send(());
1962
1963 let result = join_supervisor(handle).await;
1964 assert!(result.is_ok());
1965 assert!(
1966 transient_count.load(Ordering::SeqCst) >= 2,
1967 "transient worker must be restarted after its own abnormal exit"
1968 );
1969 assert!(
1970 stable_count.load(Ordering::SeqCst) >= 2,
1971 "the transient's abnormal exit must trigger a one-for-all that also restarts the sibling"
1972 );
1973 }
1974
1975 #[tokio::test]
1976 async fn restart_limit_exceeded_shuts_down_supervisor() {
1977 let mut sup = Supervisor::new("test-sup")
1978 .unwrap()
1979 .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(1, Duration::from_secs(10)));
1980 sup.add_worker(MockWorker::failing("fast-fail", Duration::ZERO));
1982
1983 let (tx, rx) = oneshot::channel::<()>();
1984 let handle = tokio::spawn(async move { sup.run_with_shutdown(rx).await });
1985
1986 let result = join_supervisor(handle).await;
1987 drop(tx);
1988
1989 assert!(matches!(result, Err(SupervisorError::Shutdown)));
1990 }
1991
1992 #[tokio::test]
1995 async fn temporary_child_is_not_restarted() {
1996 let temp = MockWorker::failing("temp-worker", Duration::from_millis(50));
1998 let temp_started = temp.start_count();
1999 let temp_failed = temp.finish_count();
2000
2001 let stable = MockWorker::long_running("stable-worker");
2002
2003 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2004 RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2005 );
2006 sup.add_worker(stable);
2007 sup.add_worker(ChildSpecification::worker(temp).with_restart_type(RestartType::Temporary));
2008
2009 let (tx, handle) = run_supervisor_with_trigger(sup).await;
2010
2011 wait_until("the temporary worker has failed once", || {
2017 temp_failed.load(Ordering::SeqCst) == 1
2018 })
2019 .await;
2020 let _ = tx.send(());
2021
2022 let result = join_supervisor(handle).await;
2023 assert!(result.is_ok());
2024 assert_eq!(
2025 temp_started.load(Ordering::SeqCst),
2026 1,
2027 "temporary worker must not be restarted after it fails"
2028 );
2029 }
2030
2031 #[tokio::test]
2032 async fn transient_child_is_not_restarted_on_clean_exit() {
2033 let transient = MockWorker::completing("transient-worker", Duration::from_millis(50));
2034 let transient_started = transient.start_count();
2035 let transient_finished = transient.finish_count();
2036
2037 let stable = MockWorker::long_running("stable-worker");
2038
2039 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2040 RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2041 );
2042 sup.add_worker(stable);
2043 sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
2044
2045 let (tx, handle) = run_supervisor_with_trigger(sup).await;
2046
2047 wait_until("the transient worker has completed once", || {
2052 transient_finished.load(Ordering::SeqCst) == 1
2053 })
2054 .await;
2055 let _ = tx.send(());
2056
2057 let result = join_supervisor(handle).await;
2058 assert!(result.is_ok());
2059 assert_eq!(
2060 transient_started.load(Ordering::SeqCst),
2061 1,
2062 "transient worker must not be restarted after a clean exit"
2063 );
2064 }
2065
2066 #[tokio::test]
2067 async fn transient_child_is_restarted_on_failure() {
2068 let transient = MockWorker::failing("transient-worker", Duration::from_millis(50));
2069 let transient_count = transient.start_count();
2070
2071 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2072 RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2073 );
2074 sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
2075
2076 let (tx, handle) = run_supervisor_with_trigger(sup).await;
2077
2078 wait_until("the transient worker has been restarted", || {
2079 transient_count.load(Ordering::SeqCst) >= 2
2080 })
2081 .await;
2082 let _ = tx.send(());
2083
2084 let result = join_supervisor(handle).await;
2085 assert!(result.is_ok());
2086 assert!(
2087 transient_count.load(Ordering::SeqCst) >= 2,
2088 "transient worker must be restarted after an abnormal exit"
2089 );
2090 }
2091
2092 #[tokio::test]
2093 async fn transient_child_is_restarted_on_panic() {
2094 let transient = MockWorker::panicking("transient-worker", Duration::from_millis(50));
2101 let transient_count = transient.start_count();
2102
2103 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2104 RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2105 );
2106 sup.add_worker(ChildSpecification::worker(transient).with_restart_type(RestartType::Transient));
2107
2108 let (tx, handle) = run_supervisor_with_trigger(sup).await;
2109
2110 wait_until("the panicking transient worker has been restarted", || {
2111 transient_count.load(Ordering::SeqCst) >= 2
2112 })
2113 .await;
2114 let _ = tx.send(());
2115
2116 let result = join_supervisor(handle).await;
2117 assert!(result.is_ok());
2118 assert!(
2119 transient_count.load(Ordering::SeqCst) >= 2,
2120 "transient worker must be restarted after a panic"
2121 );
2122 }
2123
2124 #[tokio::test]
2125 async fn permanent_child_is_restarted_on_clean_exit() {
2126 let permanent = MockWorker::completing("permanent-worker", Duration::from_millis(50));
2129 let permanent_count = permanent.start_count();
2130
2131 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2132 RestartStrategy::one_to_one().with_intensity_and_period(20, Duration::from_secs(10)),
2133 );
2134 sup.add_worker(permanent);
2136
2137 let (tx, handle) = run_supervisor_with_trigger(sup).await;
2138
2139 wait_until("the permanent worker has been restarted", || {
2140 permanent_count.load(Ordering::SeqCst) >= 2
2141 })
2142 .await;
2143 let _ = tx.send(());
2144
2145 let result = join_supervisor(handle).await;
2146 assert!(result.is_ok());
2147 assert!(
2148 permanent_count.load(Ordering::SeqCst) >= 2,
2149 "permanent worker must be restarted even after a clean exit"
2150 );
2151 }
2152
2153 #[tokio::test]
2154 async fn temporary_failures_do_not_consume_restart_intensity() {
2155 let mut sup = Supervisor::new("test-sup")
2159 .unwrap()
2160 .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(1, Duration::from_secs(10)));
2161
2162 let workers = [
2163 MockWorker::failing("temp-0", Duration::from_millis(20)),
2164 MockWorker::failing("temp-1", Duration::from_millis(20)),
2165 MockWorker::failing("temp-2", Duration::from_millis(20)),
2166 MockWorker::failing("temp-3", Duration::from_millis(20)),
2167 MockWorker::failing("temp-4", Duration::from_millis(20)),
2168 ];
2169 let started: Vec<_> = workers.iter().map(|w| w.start_count()).collect();
2170 let failed: Vec<_> = workers.iter().map(|w| w.finish_count()).collect();
2171 for worker in workers {
2172 sup.add_worker(ChildSpecification::worker(worker).with_restart_type(RestartType::Temporary));
2173 }
2174 sup.add_worker(MockWorker::long_running("stable-worker"));
2176
2177 let (tx, handle) = run_supervisor_with_trigger(sup).await;
2178 wait_until("every temporary worker has failed once", || {
2182 failed.iter().all(|c| c.load(Ordering::SeqCst) == 1)
2183 })
2184 .await;
2185 let _ = tx.send(());
2186
2187 let result = join_supervisor(handle).await;
2188 assert!(
2189 result.is_ok(),
2190 "supervisor must not trip its restart limit on temporary exits"
2191 );
2192 for count in started {
2193 assert_eq!(
2194 count.load(Ordering::SeqCst),
2195 1,
2196 "each temporary worker runs exactly once"
2197 );
2198 }
2199 }
2200
2201 #[tokio::test]
2202 async fn transient_clean_exits_do_not_consume_restart_intensity() {
2203 let mut sup = Supervisor::new("test-sup")
2207 .unwrap()
2208 .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(1, Duration::from_secs(10)));
2209
2210 let workers = [
2211 MockWorker::completing("transient-0", Duration::from_millis(20)),
2212 MockWorker::completing("transient-1", Duration::from_millis(20)),
2213 MockWorker::completing("transient-2", Duration::from_millis(20)),
2214 MockWorker::completing("transient-3", Duration::from_millis(20)),
2215 MockWorker::completing("transient-4", Duration::from_millis(20)),
2216 ];
2217 let started: Vec<_> = workers.iter().map(|w| w.start_count()).collect();
2218 let finished: Vec<_> = workers.iter().map(|w| w.finish_count()).collect();
2219 for worker in workers {
2220 sup.add_worker(ChildSpecification::worker(worker).with_restart_type(RestartType::Transient));
2221 }
2222 sup.add_worker(MockWorker::long_running("stable-worker"));
2224
2225 let (tx, handle) = run_supervisor_with_trigger(sup).await;
2226 wait_until("every transient worker has completed once", || {
2230 finished.iter().all(|c| c.load(Ordering::SeqCst) == 1)
2231 })
2232 .await;
2233 let _ = tx.send(());
2234
2235 let result = join_supervisor(handle).await;
2236 assert!(
2237 result.is_ok(),
2238 "supervisor must not trip its restart limit on clean transient exits"
2239 );
2240 for count in started {
2241 assert_eq!(
2242 count.load(Ordering::SeqCst),
2243 1,
2244 "each transient worker runs exactly once"
2245 );
2246 }
2247 }
2248
2249 #[tokio::test]
2250 async fn supervisor_idles_when_all_temporary_children_exit() {
2251 let temp_a = MockWorker::completing("temp-a", Duration::from_millis(10));
2255 let a_finished = temp_a.finish_count();
2256 let temp_b = MockWorker::completing("temp-b", Duration::from_millis(10));
2257 let b_finished = temp_b.finish_count();
2258
2259 let mut sup = Supervisor::new("test-sup").unwrap();
2260 let handle = sup.handle();
2261 sup.add_worker(ChildSpecification::worker(temp_a).with_restart_type(RestartType::Temporary));
2262 sup.add_worker(ChildSpecification::worker(temp_b).with_restart_type(RestartType::Temporary));
2263
2264 let (tx, run) = run_supervisor_with_trigger(sup).await;
2265
2266 wait_until("both temporary children have completed", || {
2270 a_finished.load(Ordering::SeqCst) == 1 && b_finished.load(Ordering::SeqCst) == 1
2271 })
2272 .await;
2273
2274 let dynamic = MockWorker::long_running("late-comer");
2277 let dynamic_count = dynamic.start_count();
2278 handle.spawn(dynamic);
2279 wait_until("the late dynamic child has started", || {
2280 dynamic_count.load(Ordering::SeqCst) == 1
2281 })
2282 .await;
2283 assert!(
2284 handle.is_running(),
2285 "supervisor must keep running after all temporary children exit"
2286 );
2287
2288 tx.send(()).unwrap();
2289 let result = join_supervisor(run).await;
2290 assert!(result.is_ok());
2291 }
2292
2293 #[tokio::test]
2296 async fn significant_child_drives_auto_shutdown() {
2297 let mut sup = Supervisor::new("test-sup")
2300 .unwrap()
2301 .with_auto_shutdown(AutoShutdown::AnySignificant);
2302 sup.add_worker(MockWorker::long_running("stable"));
2303 sup.add_worker(
2304 ChildSpecification::worker(MockWorker::completing("significant", Duration::from_millis(50)))
2305 .with_restart_type(RestartType::Temporary)
2306 .with_significant(true),
2307 );
2308
2309 let (_tx, rx) = oneshot::channel::<()>();
2311 let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2312 .await
2313 .unwrap();
2314 assert!(matches!(result, Err(SupervisorError::SignificantChildExited)));
2315 }
2316
2317 #[tokio::test]
2318 async fn significant_child_added_before_the_auto_shutdown_policy_still_drives_it() {
2319 let mut sup = Supervisor::new("test-sup").unwrap();
2324 sup.add_worker(MockWorker::long_running("stable"));
2325 sup.add_worker(
2326 runtime::supervisable(MockWorker::completing("significant", Duration::from_millis(50)))
2327 .temporary()
2328 .with_significant(true)
2329 .build(),
2330 );
2331 let mut sup = sup.with_auto_shutdown(AutoShutdown::AnySignificant);
2332
2333 let (_tx, rx) = oneshot::channel::<()>();
2334 let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2335 .await
2336 .unwrap();
2337 assert!(
2338 matches!(result, Err(SupervisorError::SignificantChildExited)),
2339 "the policy set after registration should still have applied, got {result:?}"
2340 );
2341 }
2342
2343 #[tokio::test]
2344 async fn non_significant_exit_does_not_auto_shutdown() {
2345 let plain = MockWorker::completing("plain", Duration::from_millis(10));
2347 let plain_finished = plain.finish_count();
2348
2349 let mut sup = Supervisor::new("test-sup")
2350 .unwrap()
2351 .with_auto_shutdown(AutoShutdown::AnySignificant);
2352 let handle = sup.handle();
2353 sup.add_worker(MockWorker::long_running("stable"));
2354 sup.add_worker(ChildSpecification::worker(plain).with_restart_type(RestartType::Temporary));
2355
2356 let (tx, run) = run_supervisor_with_trigger(sup).await;
2357
2358 wait_until("the non-significant child has completed", || {
2362 plain_finished.load(Ordering::SeqCst) == 1
2363 })
2364 .await;
2365
2366 let dynamic = MockWorker::long_running("late-comer");
2370 let dynamic_count = dynamic.start_count();
2371 handle.spawn(dynamic);
2372 wait_until("the late dynamic child has started", || {
2373 dynamic_count.load(Ordering::SeqCst) == 1
2374 })
2375 .await;
2376 assert!(
2377 handle.is_running(),
2378 "a non-significant child exiting must not trigger auto-shutdown"
2379 );
2380
2381 tx.send(()).unwrap();
2382 let result = join_supervisor(run).await;
2383 assert!(result.is_ok());
2384 }
2385
2386 #[tokio::test]
2387 async fn all_significant_waits_for_last() {
2388 let mut sup = Supervisor::new("test-sup")
2390 .unwrap()
2391 .with_auto_shutdown(AutoShutdown::AllSignificant);
2392 sup.add_worker(
2393 ChildSpecification::worker(MockWorker::completing("sig-a", Duration::from_millis(50)))
2394 .with_restart_type(RestartType::Temporary)
2395 .with_significant(true),
2396 );
2397 sup.add_worker(
2398 ChildSpecification::worker(MockWorker::completing("sig-b", Duration::from_millis(250)))
2399 .with_restart_type(RestartType::Temporary)
2400 .with_significant(true),
2401 );
2402
2403 let (_tx, rx) = oneshot::channel::<()>();
2404 let start = std::time::Instant::now();
2405 let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2406 .await
2407 .unwrap();
2408 let elapsed = start.elapsed();
2409
2410 assert!(matches!(result, Err(SupervisorError::SignificantChildExited)));
2411 assert!(
2413 elapsed >= Duration::from_millis(200),
2414 "auto-shutdown must wait for all significant children (took {elapsed:?})"
2415 );
2416 }
2417
2418 #[tokio::test]
2421 async fn init_failure_propagates_with_child_name() {
2422 let mut sup = Supervisor::new("test-sup").unwrap();
2423 sup.add_worker(MockWorker::long_running("good-worker"));
2424 sup.add_worker(MockWorker::init_failure("bad-worker"));
2425
2426 let (_tx, rx) = oneshot::channel::<()>();
2427 let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2428 .await
2429 .unwrap();
2430
2431 match result {
2432 Err(SupervisorError::FailedToInitialize { child_name, .. }) => {
2433 assert_eq!(child_name, "bad-worker");
2434 }
2435 other => panic!("expected FailedToInitialize, got: {:?}", other),
2436 }
2437 }
2438
2439 #[tokio::test]
2440 async fn init_failure_does_not_trigger_restart() {
2441 let init_fail = MockWorker::init_failure("bad-worker");
2442 let start_count = init_fail.start_count();
2443
2444 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
2445 RestartStrategy::one_to_one().with_intensity_and_period(10, Duration::from_secs(10)),
2446 );
2447 sup.add_worker(init_fail);
2448
2449 let (_tx, rx) = oneshot::channel::<()>();
2450 let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
2451 .await
2452 .unwrap();
2453
2454 assert!(matches!(result, Err(SupervisorError::FailedToInitialize { .. })));
2455 assert_eq!(start_count.load(Ordering::SeqCst), 0);
2457 }
2458
2459 #[tokio::test]
2462 async fn shutdown_completes_promptly_in_steady_state() {
2463 let mut sup = Supervisor::new("test-sup").unwrap();
2464 sup.add_worker(MockWorker::long_running("worker1"));
2465 sup.add_worker(MockWorker::long_running("worker2"));
2466
2467 let (tx, handle) = run_supervisor_with_trigger(sup).await;
2468 tx.send(()).unwrap();
2469
2470 let result = timeout(Duration::from_secs(1), handle).await;
2472 assert!(result.is_ok(), "shutdown should complete promptly");
2473 }
2474
2475 #[tokio::test]
2476 async fn shutdown_during_slow_init_completes_promptly() {
2477 let mut sup = Supervisor::new("test-sup").unwrap();
2478 sup.add_worker(MockWorker::slow_init("slow-worker", Duration::from_secs(30)));
2480
2481 let (tx, rx) = oneshot::channel();
2482 let handle = tokio::spawn(async move { sup.run_with_shutdown(rx).await });
2483
2484 sleep(Duration::from_millis(20)).await;
2486 tx.send(()).unwrap();
2487
2488 let result = timeout(Duration::from_secs(2), handle).await;
2491 assert!(result.is_ok(), "shutdown during slow init should complete promptly");
2492 }
2493
2494 #[tokio::test]
2497 async fn dynamic_children_spawn_after_start() {
2498 let sup = Supervisor::new("dyn-sup").unwrap();
2499 let handle = sup.handle();
2500 let (tx, run) = run_supervisor_with_trigger(sup).await;
2501 wait_until("supervisor is running", || handle.is_running()).await;
2502
2503 let c1 = MockWorker::long_running("c1");
2504 let c2 = MockWorker::long_running("c2");
2505 let c1_count = c1.start_count();
2506 let c2_count = c2.start_count();
2507 handle.spawn(c1);
2508 handle.spawn(c2);
2509
2510 wait_until("both dynamic children have started", || {
2511 c1_count.load(Ordering::SeqCst) == 1 && c2_count.load(Ordering::SeqCst) == 1
2512 })
2513 .await;
2514 assert_eq!(handle.active_children(), 2);
2515
2516 tx.send(()).unwrap();
2517 let result = join_supervisor(run).await;
2518 assert!(result.is_ok());
2519 assert_eq!(
2520 handle.active_children(),
2521 0,
2522 "all dynamic children must be drained on shutdown"
2523 );
2524 }
2525
2526 #[tokio::test]
2527 async fn temporary_dynamic_child_failure_is_isolated() {
2528 let sup = Supervisor::new("dyn-sup").unwrap();
2531 let handle = sup.handle();
2532 let (tx, run) = run_supervisor_with_trigger(sup).await;
2533 wait_until("supervisor is running", || handle.is_running()).await;
2534
2535 let failing = MockWorker::failing("boom", Duration::from_millis(20));
2536 let failing_count = failing.start_count();
2537 handle.spawn(failing);
2538 wait_until("the failing dynamic child has run once", || {
2539 failing_count.load(Ordering::SeqCst) == 1
2540 })
2541 .await;
2542 wait_until("all dynamic children have drained", || handle.active_children() == 0).await;
2543
2544 sleep(Duration::from_millis(50)).await;
2545 assert!(
2546 handle.is_running(),
2547 "supervisor stays up after an isolated child failure"
2548 );
2549 assert_eq!(
2550 failing_count.load(Ordering::SeqCst),
2551 1,
2552 "a temporary child is never restarted"
2553 );
2554
2555 handle.spawn(MockWorker::long_running("c2"));
2557 wait_until("one dynamic child is running", || handle.active_children() == 1).await;
2558
2559 tx.send(()).unwrap();
2560 let result = join_supervisor(run).await;
2561 assert!(result.is_ok());
2562 }
2563
2564 #[tokio::test]
2565 async fn temporary_dynamic_child_panic_is_isolated() {
2566 let sup = Supervisor::new("dyn-sup").unwrap();
2568 let handle = sup.handle();
2569 let (tx, run) = run_supervisor_with_trigger(sup).await;
2570 wait_until("supervisor is running", || handle.is_running()).await;
2571
2572 handle.spawn(MockWorker::panicking("boom", Duration::from_millis(20)));
2573 wait_until("all dynamic children have drained", || handle.active_children() == 0).await;
2574
2575 sleep(Duration::from_millis(50)).await;
2576 assert!(handle.is_running(), "supervisor stays up after an isolated child panic");
2577
2578 tx.send(()).unwrap();
2579 let result = join_supervisor(run).await;
2580 assert!(result.is_ok());
2581 }
2582
2583 #[tokio::test]
2584 async fn significant_dynamic_child_failure_shuts_down_supervisor() {
2585 let sup = Supervisor::new("dyn-sup")
2588 .unwrap()
2589 .with_auto_shutdown(AutoShutdown::AnySignificant);
2590 let handle = sup.handle();
2591 let (_tx, run) = run_supervisor_with_trigger(sup).await;
2592 wait_until("supervisor is running", || handle.is_running()).await;
2593
2594 handle.spawn(
2595 ChildSpecification::worker(MockWorker::failing("boom", Duration::from_millis(20))).with_significant(true),
2596 );
2597
2598 let result = join_supervisor(run).await;
2599 assert!(matches!(result, Err(SupervisorError::SignificantChildExited)));
2600 }
2601
2602 #[tokio::test]
2603 async fn dynamic_spawn_outside_a_run_is_accepted_and_dropped() {
2604 let sup = Supervisor::new("dyn-sup").unwrap();
2609 let handle = sup.handle();
2610
2611 assert!(!handle.is_running());
2612 let before = MockWorker::long_running("before-start");
2613 let before_count = before.start_count();
2614 handle.spawn(before);
2615
2616 let (tx, run) = run_supervisor_with_trigger(sup).await;
2618 wait_until("supervisor is running", || handle.is_running()).await;
2619 let worker = MockWorker::long_running("after-start");
2620 let started = worker.start_count();
2621 handle.spawn(worker);
2622 wait_until("the dynamic child has started", || started.load(Ordering::SeqCst) == 1).await;
2623 assert_eq!(
2624 before_count.load(Ordering::SeqCst),
2625 0,
2626 "a child spawned before the run must not be started by it"
2627 );
2628
2629 tx.send(()).unwrap();
2630 let result = join_supervisor(run).await;
2631 assert!(result.is_ok());
2632
2633 wait_until("the supervisor has stopped", || !handle.is_running()).await;
2635 let after = MockWorker::long_running("after-shutdown");
2636 let after_count = after.start_count();
2637 handle.spawn(after);
2638
2639 sleep(Duration::from_millis(50)).await;
2640 assert_eq!(
2641 after_count.load(Ordering::SeqCst),
2642 0,
2643 "a child spawned after shutdown must never start"
2644 );
2645 }
2646
2647 #[tokio::test]
2648 async fn dynamic_spawn_allocates_an_id_eagerly() {
2649 let sup = Supervisor::new("dyn-sup").unwrap();
2653 let handle = sup.handle();
2654 let (tx, run) = run_supervisor_with_trigger(sup).await;
2655 wait_until("supervisor is running", || handle.is_running()).await;
2656
2657 let worker = MockWorker::long_running("c");
2658 let started = worker.start_count();
2659 let id = handle.spawn(worker);
2660 assert_eq!(id.as_u64(), 0);
2661 wait_until("the dynamic child has started", || started.load(Ordering::SeqCst) == 1).await;
2662
2663 tx.send(()).unwrap();
2664 let result = join_supervisor(run).await;
2665 assert!(result.is_ok());
2666 }
2667
2668 #[tokio::test]
2669 async fn dynamic_child_with_an_unusable_name_runs_under_a_placeholder() {
2670 let recorder = TestRecorder::default();
2675 let _guard = metrics::set_default_local_recorder(&recorder);
2676
2677 let sup = Supervisor::new("dyn-sup").unwrap();
2678 let handle = sup.handle();
2679 let (tx, run) = run_supervisor_with_trigger(sup).await;
2680 wait_until("supervisor is running", || handle.is_running()).await;
2681
2682 let worker = MockWorker::long_running("");
2683 let started = worker.start_count();
2684 handle.spawn(worker);
2685 wait_until("the unnamed dynamic child has started", || {
2686 started.load(Ordering::SeqCst) == 1
2687 })
2688 .await;
2689
2690 assert!(handle.is_running());
2692 handle.spawn(MockWorker::long_running("ok"));
2693 wait_until("both dynamic children are running", || handle.active_children() == 2).await;
2694
2695 tx.send(()).unwrap();
2696 let result = join_supervisor(run).await;
2697 assert!(result.is_ok());
2698
2699 let polls = recorder.counter(("runtime_task_poll_count", &[("task_name", "dyn_sup.unnamed")]));
2700 assert!(
2701 polls.is_some_and(|polls| polls > 0),
2702 "the child should have run under the placeholder name, got {polls:?}"
2703 );
2704 }
2705
2706 #[tokio::test]
2707 async fn dynamic_spawns_are_not_capped() {
2708 const CHILDREN: usize = 2048;
2712
2713 let sup = Supervisor::new("dyn-sup").unwrap();
2714 let handle = sup.handle();
2715 let (tx, run) = run_supervisor_with_trigger(sup).await;
2716 wait_until("supervisor is running", || handle.is_running()).await;
2717
2718 for _ in 0..CHILDREN {
2719 handle.spawn(MockWorker::long_running("burst").with_graceful_timeout(Duration::from_secs(30)));
2722 }
2723
2724 wait_until("every child in the burst has started", || {
2725 handle.active_children() == CHILDREN
2726 })
2727 .await;
2728
2729 tx.send(()).unwrap();
2730 let result = timeout(Duration::from_secs(30), run)
2731 .await
2732 .expect("supervisor should stop")
2733 .expect("supervisor task should not panic");
2734 assert!(result.is_ok(), "the burst should have drained cleanly: {result:?}");
2735 }
2736
2737 #[tokio::test]
2738 async fn dynamically_spawned_supervisor_runs_and_drains() {
2739 let child_worker = MockWorker::long_running("nested-child");
2742 let child_started = child_worker.start_count();
2743 let mut nested = Supervisor::new("nested-sup").unwrap();
2744 nested.add_worker(child_worker);
2745
2746 let sup = Supervisor::new("dyn-sup").unwrap();
2747 let handle = sup.handle();
2748 let (tx, run) = run_supervisor_with_trigger(sup).await;
2749 wait_until("supervisor is running", || handle.is_running()).await;
2750
2751 handle.spawn(nested);
2752 wait_until("the nested supervisor's own child has started", || {
2753 child_started.load(Ordering::SeqCst) == 1
2754 })
2755 .await;
2756
2757 tx.send(()).unwrap();
2758 let result = join_supervisor(run).await;
2759 assert!(
2760 result.is_ok(),
2761 "the nested subtree should have drained cleanly: {result:?}"
2762 );
2763 }
2764
2765 fn failing_subtree(name: &'static str) -> (Supervisor, Arc<AtomicUsize>) {
2772 let worker = MockWorker::failing("nested-child", Duration::from_millis(5));
2773 let started = worker.start_count();
2774 let mut nested = Supervisor::new(name)
2775 .unwrap()
2776 .with_restart_strategy(RestartStrategy::new(RestartMode::OneForOne, 0, Duration::from_secs(30)));
2777 nested.add_worker(worker);
2778
2779 (nested, started)
2780 }
2781
2782 #[tokio::test]
2783 async fn dynamic_nested_supervisor_defaults_to_temporary() {
2784 let (nested, started) = failing_subtree("nested-temp");
2788
2789 let sup = Supervisor::new("dyn-temp-sup").unwrap();
2790 let handle = sup.handle();
2791 let (tx, run) = run_supervisor_with_trigger(sup).await;
2792
2793 handle.spawn(nested);
2794 wait_until("the subtree has started once", || started.load(Ordering::SeqCst) == 1).await;
2795
2796 sleep(Duration::from_millis(200)).await;
2798 assert_eq!(
2799 started.load(Ordering::SeqCst),
2800 1,
2801 "a temporary subtree must not be restarted"
2802 );
2803
2804 tx.send(()).unwrap();
2805 assert!(join_supervisor(run).await.is_ok());
2806 }
2807
2808 #[tokio::test]
2809 async fn dynamic_nested_supervisor_can_be_made_permanent() {
2810 let (nested, started) = failing_subtree("nested-perm");
2812
2813 let sup = Supervisor::new("dyn-perm-sup")
2814 .unwrap()
2815 .with_restart_strategy(RestartStrategy::new(
2816 RestartMode::OneForOne,
2817 100,
2818 Duration::from_secs(30),
2819 ));
2820 let handle = sup.handle();
2821 let (tx, run) = run_supervisor_with_trigger(sup).await;
2822
2823 handle.nested_supervisor(nested).spawn();
2824 wait_until("the subtree has been restarted", || started.load(Ordering::SeqCst) >= 2).await;
2825
2826 tx.send(()).unwrap();
2827 assert!(join_supervisor(run).await.is_ok());
2828 }
2829
2830 #[tokio::test]
2831 async fn dynamic_nested_supervisor_can_be_significant() {
2832 let (nested, _started) = failing_subtree("nested-sig");
2835
2836 let sup = Supervisor::new("dyn-sig-sup")
2837 .unwrap()
2838 .with_auto_shutdown(AutoShutdown::AnySignificant);
2839 let handle = sup.handle();
2840 let (tx, run) = run_supervisor_with_trigger(sup).await;
2841
2842 handle
2843 .nested_supervisor(nested)
2844 .temporary()
2845 .with_significant(true)
2846 .spawn();
2847
2848 let result = timeout(Duration::from_secs(5), run)
2849 .await
2850 .expect("supervisor should stop once the significant subtree terminates")
2851 .expect("supervisor task should not panic");
2852 assert!(
2853 matches!(result, Err(SupervisorError::SignificantChildExited)),
2854 "the significant subtree's termination should have stopped the parent, got {result:?}"
2855 );
2856
2857 let _ = tx.send(());
2859 }
2860
2861 #[tokio::test]
2862 async fn budget_of_duration_max_does_not_leave_a_budget_bounded_child_unbounded() {
2863 let mut sup = Supervisor::new("test-sup").unwrap().with_shutdown_budget(Duration::MAX);
2868 sup.add_worker(
2869 ChildSpecification::one_shot_worker(
2870 MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_millis(100)),
2871 )
2872 .with_budget_bounded_shutdown(),
2873 );
2874
2875 let (tx, run) = run_supervisor_with_trigger(sup).await;
2876 tx.send(()).unwrap();
2877
2878 let started = tokio::time::Instant::now();
2879 let result = join_supervisor(run).await;
2880 let elapsed = started.elapsed();
2881
2882 assert!(
2883 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
2884 "the worker's own deadline should have aborted it, got {result:?}"
2885 );
2886 assert!(
2887 elapsed < Duration::from_secs(1),
2888 "the child should have been bounded by its own 100ms deadline; took {elapsed:?}"
2889 );
2890 }
2891
2892 #[tokio::test]
2893 async fn budget_bounded_child_falls_back_to_its_own_deadline_without_a_budget() {
2894 let mut sup = Supervisor::new("test-sup").unwrap();
2898 sup.add_worker(
2899 ChildSpecification::one_shot_worker(
2900 MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_millis(100)),
2901 )
2902 .with_budget_bounded_shutdown(),
2903 );
2904
2905 let (tx, run) = run_supervisor_with_trigger(sup).await;
2906 tx.send(()).unwrap();
2907
2908 let started = tokio::time::Instant::now();
2909 let result = join_supervisor(run).await;
2910 let elapsed = started.elapsed();
2911
2912 assert!(
2913 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
2914 "the worker's own deadline should have aborted it, got {result:?}"
2915 );
2916 assert!(
2917 elapsed < Duration::from_secs(1),
2918 "the child should have been bounded by its own 100ms deadline; took {elapsed:?}"
2919 );
2920 }
2921
2922 #[tokio::test]
2923 async fn concurrent_shutdown_drains_many_children_quickly() {
2924 const CHILDREN: usize = 500;
2925 const SHUTDOWN_DELAY: Duration = Duration::from_millis(50);
2926
2927 let sup = Supervisor::new("dyn-sup").unwrap();
2928 let handle = sup.handle();
2929 let (tx, run) = run_supervisor_with_trigger(sup).await;
2930 wait_until("supervisor is running", || handle.is_running()).await;
2931
2932 for _ in 0..CHILDREN {
2933 handle.spawn(MockWorker::slow_shutdown("conn", SHUTDOWN_DELAY));
2934 }
2935 wait_until("all dynamic children are running", || {
2936 handle.active_children() == CHILDREN
2937 })
2938 .await;
2939
2940 let start = std::time::Instant::now();
2943 tx.send(()).unwrap();
2944 let result = timeout(Duration::from_secs(5), run).await.unwrap().unwrap();
2945 let elapsed = start.elapsed();
2946
2947 assert!(result.is_ok());
2948 assert_eq!(handle.active_children(), 0, "active count must return to zero");
2949 assert!(
2950 elapsed < Duration::from_secs(2),
2951 "shutdown must be concurrent (took {elapsed:?})"
2952 );
2953 }
2954
2955 #[tokio::test]
2956 async fn concurrent_shutdown_aborts_unresponsive_children() {
2957 let sup = Supervisor::new("dyn-sup").unwrap();
2958 let handle = sup.handle();
2959 let (tx, run) = run_supervisor_with_trigger(sup).await;
2960 wait_until("supervisor is running", || handle.is_running()).await;
2961
2962 handle.spawn(MockWorker::ignore_shutdown("stuck"));
2963 wait_until("one dynamic child is running", || handle.active_children() == 1).await;
2964
2965 let start = std::time::Instant::now();
2968 tx.send(()).unwrap();
2969 let result = join_supervisor(run).await;
2970 let elapsed = start.elapsed();
2971
2972 assert!(
2974 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
2975 "aborting a stuck child must surface as an unclean shutdown, got {result:?}"
2976 );
2977 assert_eq!(handle.active_children(), 0);
2978 assert!(
2979 elapsed < Duration::from_secs(1),
2980 "stuck child must be aborted at the deadline (took {elapsed:?})"
2981 );
2982 }
2983
2984 #[tokio::test]
2985 async fn concurrent_shutdown_honors_per_child_deadline() {
2986 let sup = Supervisor::new("dyn-sup").unwrap();
2991 let handle = sup.handle();
2992 let (tx, run) = run_supervisor_with_trigger(sup).await;
2993 wait_until("supervisor is running", || handle.is_running()).await;
2994
2995 handle.spawn(MockWorker::long_running("responsive").with_graceful_timeout(Duration::MAX));
2997 handle.spawn(MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_millis(200)));
2999 wait_until("both dynamic children are running", || handle.active_children() == 2).await;
3000
3001 let start = std::time::Instant::now();
3002 tx.send(()).unwrap();
3003 let result = join_supervisor(run).await;
3004 let elapsed = start.elapsed();
3005
3006 assert!(
3008 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3009 "aborting the stuck child must surface as an unclean shutdown with a count of 1, got {result:?}"
3010 );
3011 assert_eq!(handle.active_children(), 0);
3012 assert!(
3013 elapsed < Duration::from_secs(1),
3014 "stuck child must be aborted at its own deadline despite an infinite-timeout sibling (took {elapsed:?})"
3015 );
3016 }
3017
3018 #[tokio::test]
3019 async fn unresponsive_child_is_aborted_at_its_deadline() {
3020 let mut sup = Supervisor::new("test-sup").unwrap();
3023 sup.add_worker(MockWorker::ignore_shutdown("stuck"));
3024
3025 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3026
3027 let start = std::time::Instant::now();
3028 tx.send(()).unwrap();
3029 let result = join_supervisor(handle).await;
3030 let elapsed = start.elapsed();
3031
3032 assert!(
3033 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3034 "aborting a stuck child must surface as an unclean shutdown, got {result:?}"
3035 );
3036 assert!(
3037 elapsed < Duration::from_secs(1),
3038 "unresponsive child must be aborted at its deadline (took {elapsed:?})"
3039 );
3040 }
3041
3042 #[tokio::test]
3043 async fn brutal_shutdown_aborts_child_immediately() {
3044 let mut sup = Supervisor::new("test-sup").unwrap();
3047 sup.add_worker(MockWorker::ignore_shutdown("brutal-stuck").with_brutal_shutdown());
3048
3049 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3050
3051 let start = std::time::Instant::now();
3052 tx.send(()).unwrap();
3053 let result = join_supervisor(handle).await;
3054 let elapsed = start.elapsed();
3055
3056 assert!(result.is_ok());
3059 assert!(
3060 elapsed < Duration::from_millis(200),
3061 "brutal-shutdown child must be aborted immediately, not after a graceful wait (took {elapsed:?})"
3062 );
3063 }
3064
3065 #[tokio::test]
3066 async fn shutdown_timeout_aborts_aggregate_to_root() {
3067 let mut child_sup = Supervisor::new("child-sup").unwrap();
3071 child_sup
3072 .add_worker(MockWorker::ignore_shutdown("child-stuck").with_graceful_timeout(Duration::from_millis(200)));
3073
3074 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3075 parent_sup
3076 .add_worker(MockWorker::ignore_shutdown("parent-stuck").with_graceful_timeout(Duration::from_millis(200)));
3077 parent_sup.add_worker(MockWorker::long_running("parent-clean"));
3078 parent_sup.add_worker(child_sup);
3079
3080 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3081 tx.send(()).unwrap();
3082
3083 let result = join_supervisor(handle).await;
3084 assert!(
3085 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 2 })),
3086 "forced aborts must aggregate across the tree (1 direct + 1 nested), got {result:?}"
3087 );
3088 }
3089
3090 #[tokio::test]
3093 async fn restart_intensity_zero_shuts_down_on_first_failure() {
3094 let worker = MockWorker::failing("boom", Duration::from_millis(20));
3098 let start_count = worker.start_count();
3099
3100 let mut sup = Supervisor::new("test-sup")
3101 .unwrap()
3102 .with_restart_strategy(RestartStrategy::new(RestartMode::OneForOne, 0, Duration::from_secs(5)));
3103 sup.add_worker(worker);
3104
3105 let (_tx, rx) = oneshot::channel::<()>();
3106 let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
3107 .await
3108 .unwrap();
3109
3110 assert!(matches!(result, Err(SupervisorError::Shutdown)));
3111 assert_eq!(
3112 start_count.load(Ordering::SeqCst),
3113 1,
3114 "with intensity zero the worker must run exactly once and never be restarted"
3115 );
3116 }
3117
3118 #[tokio::test]
3119 async fn one_for_all_restart_loses_dynamic_children() {
3120 let failing = MockWorker::failing("failing-static", Duration::from_millis(50));
3125 let failing_count = failing.start_count();
3126
3127 let sup = Supervisor::new("dyn-sup").unwrap().with_restart_strategy(
3128 RestartStrategy::one_for_all().with_intensity_and_period(20, Duration::from_secs(10)),
3129 );
3130 let handle = sup.handle();
3131 let mut sup = sup;
3132 sup.add_worker(failing);
3133
3134 let (tx, run) = run_supervisor_with_trigger(sup).await;
3135
3136 let dynamic = MockWorker::long_running("dynamic");
3138 let dynamic_count = dynamic.start_count();
3139 handle.spawn(dynamic);
3140 wait_until("the dynamic child is running", || handle.active_children() == 1).await;
3141
3142 wait_until("the static worker has been restarted", || {
3144 failing_count.load(Ordering::SeqCst) >= 2
3145 })
3146 .await;
3147
3148 wait_until("the dynamic child has been discarded", || handle.active_children() == 0).await;
3151 assert_eq!(
3152 dynamic_count.load(Ordering::SeqCst),
3153 1,
3154 "a dynamic child must be lost -- not restored -- across a one-for-all restart"
3155 );
3156
3157 tx.send(()).unwrap();
3158 let result = join_supervisor(run).await;
3159 assert!(result.is_ok());
3160 }
3161
3162 #[tokio::test]
3165 async fn dedicated_single_threaded_runtime_runs_nested_worker_and_shuts_down_cleanly() {
3166 let worker = MockWorker::long_running("dedicated-worker");
3170 let worker_count = worker.start_count();
3171
3172 let mut child_sup = Supervisor::new("child-sup")
3173 .unwrap()
3174 .with_dedicated_runtime(RuntimeConfiguration::single_threaded());
3175 child_sup.add_worker(worker);
3176
3177 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3178 parent_sup.add_worker(child_sup);
3179
3180 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3181
3182 wait_until("the dedicated worker has started", || {
3184 worker_count.load(Ordering::SeqCst) == 1
3185 })
3186 .await;
3187
3188 tx.send(()).unwrap();
3189 let result = join_supervisor(handle).await;
3190 assert!(
3191 result.is_ok(),
3192 "dedicated-runtime supervisor should shut down cleanly, got {result:?}"
3193 );
3194 }
3195
3196 #[tokio::test]
3197 async fn dedicated_multi_threaded_runtime_runs_nested_worker() {
3198 let worker = MockWorker::long_running("dedicated-worker");
3200 let worker_count = worker.start_count();
3201
3202 let mut child_sup = Supervisor::new("child-sup")
3203 .unwrap()
3204 .with_dedicated_runtime(RuntimeConfiguration::multi_threaded(2));
3205 child_sup.add_worker(worker);
3206
3207 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3208 parent_sup.add_worker(child_sup);
3209
3210 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3211 wait_until("the dedicated worker has started", || {
3212 worker_count.load(Ordering::SeqCst) == 1
3213 })
3214 .await;
3215
3216 tx.send(()).unwrap();
3217 let result = join_supervisor(handle).await;
3218 assert!(
3219 result.is_ok(),
3220 "multi-threaded dedicated-runtime supervisor should shut down cleanly, got {result:?}"
3221 );
3222 }
3223
3224 #[tokio::test]
3225 async fn dedicated_runtime_forced_abort_aggregates_to_root() {
3226 let stuck = MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_millis(200));
3230 let stuck_count = stuck.start_count();
3231
3232 let mut child_sup = Supervisor::new("child-sup")
3233 .unwrap()
3234 .with_dedicated_runtime(RuntimeConfiguration::single_threaded());
3235 child_sup.add_worker(stuck);
3236
3237 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3238 parent_sup.add_worker(child_sup);
3239
3240 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3241
3242 wait_until("the stuck worker has started", || {
3245 stuck_count.load(Ordering::SeqCst) == 1
3246 })
3247 .await;
3248
3249 tx.send(()).unwrap();
3250 let result = join_supervisor(handle).await;
3251 assert!(
3252 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3253 "a stuck worker in a dedicated runtime must surface as an unclean shutdown aggregated to the root, got {result:?}"
3254 );
3255 }
3256
3257 #[tokio::test]
3260 async fn child_with_runtime_override_runs_on_that_runtime() {
3261 let child_runtime = tokio::runtime::Builder::new_multi_thread()
3264 .worker_threads(1)
3265 .thread_name("child-rt-test")
3266 .enable_all()
3267 .build()
3268 .expect("should build child runtime");
3269
3270 let (thread_tx, thread_rx) = oneshot::channel();
3271 let worker = FnWorker::new("placed", async move {
3272 let thread_name = std::thread::current().name().unwrap_or_default().to_string();
3273 let _ = thread_tx.send(thread_name);
3274 pending::<()>().await;
3275 });
3276
3277 let mut sup = Supervisor::new("test-sup").unwrap();
3280 sup.add_worker(
3281 ChildSpecification::worker(worker)
3282 .with_restart_type(RestartType::Temporary)
3283 .with_runtime(child_runtime.handle().clone())
3284 .with_shutdown_strategy(ShutdownStrategy::Brutal),
3285 );
3286
3287 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3288
3289 let thread_name = timeout(Duration::from_secs(2), thread_rx)
3290 .await
3291 .expect("child should report its thread promptly")
3292 .expect("child should not be dropped before reporting");
3293 assert!(
3294 thread_name.starts_with("child-rt-test"),
3295 "child must run on the runtime given to `with_runtime`, but ran on thread {thread_name:?}"
3296 );
3297
3298 tx.send(()).unwrap();
3299 assert!(join_supervisor(handle).await.is_ok());
3300
3301 child_runtime.shutdown_background();
3303 }
3304
3305 #[tokio::test]
3306 async fn child_shutdown_strategy_override_takes_precedence_over_worker() {
3307 let worker = MockWorker::ignore_shutdown("stuck").with_graceful_timeout(Duration::from_secs(30));
3312
3313 let mut sup = Supervisor::new("test-sup").unwrap();
3314 sup.add_worker(
3315 ChildSpecification::worker(worker)
3316 .with_restart_type(RestartType::Temporary)
3317 .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::from_millis(50))),
3318 );
3319
3320 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3321 tx.send(()).unwrap();
3322
3323 let result = join_supervisor(handle).await;
3324 assert!(
3325 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3326 "the overridden 50ms deadline should have aborted the stuck child, got {result:?}"
3327 );
3328 }
3329
3330 struct DrainWaiter {
3336 coordinator: Mutex<Option<ShutdownCoordinator>>,
3337 finished: Arc<AtomicBool>,
3338 }
3339
3340 #[async_trait]
3341 impl Supervisable for DrainWaiter {
3342 fn name(&self) -> &str {
3343 "waiter"
3344 }
3345
3346 async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
3347 let coordinator = self
3348 .coordinator
3349 .lock()
3350 .expect("drain waiter mutex poisoned")
3351 .take()
3352 .expect("drain waiter runs once");
3353 let finished = Arc::clone(&self.finished);
3354
3355 Ok(Box::pin(async move {
3356 process_shutdown.await;
3357 coordinator.shutdown_and_wait().await;
3358 finished.store(true, Ordering::SeqCst);
3359 Ok(())
3360 }))
3361 }
3362 }
3363
3364 fn build_drain_pair(
3370 stuck_strategy: ShutdownStrategy, waiter_strategy: ShutdownStrategy,
3371 ) -> (Supervisor, Arc<AtomicBool>) {
3372 let mut coordinator = ShutdownCoordinator::default();
3373 let held_handle = coordinator.register();
3374
3375 let stuck = FnWorker::new("stuck", async move {
3376 let _held = held_handle;
3378 pending::<()>().await;
3379 });
3380
3381 let waiter_finished = Arc::new(AtomicBool::new(false));
3382 let waiter = DrainWaiter {
3383 coordinator: Mutex::new(Some(coordinator)),
3384 finished: Arc::clone(&waiter_finished),
3385 };
3386
3387 let mut sup = Supervisor::new("test-sup").unwrap();
3388 sup.add_worker(
3389 ChildSpecification::worker(stuck)
3390 .with_restart_type(RestartType::Temporary)
3391 .with_shutdown_strategy(stuck_strategy),
3392 );
3393 sup.add_worker(
3394 ChildSpecification::worker(waiter)
3395 .with_restart_type(RestartType::Temporary)
3396 .with_shutdown_strategy(waiter_strategy),
3397 );
3398
3399 (sup, waiter_finished)
3400 }
3401
3402 #[tokio::test]
3403 async fn shorter_child_deadline_releases_a_waiting_sibling() {
3404 let (sup, waiter_finished) = build_drain_pair(
3408 ShutdownStrategy::Graceful(Duration::from_millis(100)),
3409 ShutdownStrategy::Graceful(Duration::from_secs(1)),
3410 );
3411
3412 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3413 tx.send(()).unwrap();
3414
3415 let result = join_supervisor(handle).await;
3416 assert!(
3417 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3418 "only the stuck child should have been aborted, got {result:?}"
3419 );
3420 assert!(
3421 waiter_finished.load(Ordering::SeqCst),
3422 "the waiter should have been released by the stuck child's abort and run to completion"
3423 );
3424 }
3425
3426 #[tokio::test]
3427 async fn equal_child_deadlines_abort_the_waiter_too() {
3428 let (sup, waiter_finished) = build_drain_pair(
3433 ShutdownStrategy::Graceful(Duration::from_millis(100)),
3434 ShutdownStrategy::Graceful(Duration::from_millis(100)),
3435 );
3436
3437 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3438 tx.send(()).unwrap();
3439
3440 let result = join_supervisor(handle).await;
3441 assert!(
3442 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 2 })),
3443 "both children should have been aborted together, got {result:?}"
3444 );
3445 assert!(
3446 !waiter_finished.load(Ordering::SeqCst),
3447 "the waiter should have been aborted mid-wait, not completed"
3448 );
3449 }
3450
3451 #[tokio::test]
3454 async fn budget_bounds_children_that_have_no_deadline_of_their_own() {
3455 let mut sup = Supervisor::new("test-sup")
3460 .unwrap()
3461 .with_shutdown_budget(Duration::from_millis(100));
3462
3463 for name in ["stuck_one", "stuck_two"] {
3464 sup.add_worker(
3465 ChildSpecification::worker(FnWorker::new(name, pending::<()>()))
3466 .with_restart_type(RestartType::Temporary)
3467 .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::MAX)),
3468 );
3469 }
3470
3471 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3472 tx.send(()).unwrap();
3473
3474 let result = join_supervisor(handle).await;
3475 assert!(
3476 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 2 })),
3477 "the budget should have aborted both deadline-less children, got {result:?}"
3478 );
3479 }
3480
3481 #[tokio::test]
3482 async fn budget_does_not_delay_children_that_stop_on_their_own() {
3483 let mut sup = Supervisor::new("test-sup")
3486 .unwrap()
3487 .with_shutdown_budget(Duration::from_secs(30));
3488 sup.add_worker(
3489 ChildSpecification::one_shot_worker(MockWorker::long_running("prompt"))
3490 .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::MAX)),
3491 );
3492
3493 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3494 let started = tokio::time::Instant::now();
3495 tx.send(()).unwrap();
3496
3497 assert!(join_supervisor(handle).await.is_ok());
3498 let elapsed = started.elapsed();
3499 assert!(
3500 elapsed < Duration::from_millis(500),
3501 "shutdown should finish as soon as the child does, not burn the budget; took {elapsed:?}"
3502 );
3503 }
3504
3505 #[tokio::test]
3506 async fn child_deadline_shorter_than_budget_still_wins() {
3507 let mut sup = Supervisor::new("test-sup")
3510 .unwrap()
3511 .with_shutdown_budget(Duration::from_secs(30));
3512 sup.add_worker(
3513 ChildSpecification::worker(FnWorker::new("stuck", pending::<()>()))
3514 .with_restart_type(RestartType::Temporary)
3515 .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::from_millis(100))),
3516 );
3517
3518 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3519 tx.send(()).unwrap();
3520
3521 let result = join_supervisor(handle).await;
3523 assert!(
3524 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 1 })),
3525 "the child's own 100ms deadline should have won over the budget, got {result:?}"
3526 );
3527 }
3528
3529 #[tokio::test]
3530 async fn every_worker_records_poll_metrics() {
3531 let recorder = TestRecorder::default();
3534 let _guard = metrics::set_default_local_recorder(&recorder);
3535
3536 let mut sup = Supervisor::new("metrics_sup").unwrap();
3538 sup.add_worker(ChildSpecification::one_shot_worker(MockWorker::long_running("timed")));
3539
3540 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3541 tx.send(()).unwrap();
3542 assert!(join_supervisor(handle).await.is_ok());
3543
3544 let polls = recorder.counter(("runtime_task_poll_count", &[("task_name", "metrics_sup.timed")]));
3545 assert!(
3546 polls.is_some_and(|polls| polls > 0),
3547 "a supervised worker should have recorded poll metrics, got {polls:?}"
3548 );
3549 }
3550
3551 #[tokio::test]
3552 async fn budget_bounds_the_whole_drain_rather_than_each_child() {
3553 let mut sup = Supervisor::new("test-sup")
3557 .unwrap()
3558 .with_shutdown_budget(Duration::from_millis(150));
3559
3560 for name in ["stuck_one", "stuck_two", "stuck_three"] {
3561 sup.add_worker(
3562 ChildSpecification::one_shot_worker(FnWorker::new(name, pending::<()>()))
3563 .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::from_secs(10))),
3564 );
3565 }
3566
3567 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3568 tx.send(()).unwrap();
3569
3570 let result = join_supervisor(handle).await;
3571 assert!(
3572 matches!(result, Err(SupervisorError::ShutdownTimedOut { aborted: 3 })),
3573 "the budget should have bounded the whole drain, got {result:?}"
3574 );
3575 }
3576
3577 #[tokio::test]
3578 async fn budget_of_duration_max_is_treated_as_no_budget() {
3579 let mut sup = Supervisor::new("test-sup").unwrap().with_shutdown_budget(Duration::MAX);
3582 sup.add_worker(ChildSpecification::one_shot_worker(MockWorker::long_running("prompt")));
3583
3584 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3585 tx.send(()).unwrap();
3586 assert!(join_supervisor(handle).await.is_ok());
3587 }
3588
3589 #[tokio::test]
3590 async fn near_max_child_timeout_does_not_panic() {
3591 let mut sup = Supervisor::new("test-sup").unwrap();
3594 sup.add_worker(
3595 ChildSpecification::one_shot_worker(MockWorker::long_running("prompt"))
3596 .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::MAX - Duration::from_nanos(1))),
3597 );
3598
3599 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3600 tx.send(()).unwrap();
3601 assert!(join_supervisor(handle).await.is_ok());
3602 }
3603
3604 #[tokio::test]
3605 async fn budget_does_not_cut_off_a_nested_supervisor_mid_drain() {
3606 let slow = MockWorker::slow_shutdown("slow", Duration::from_millis(300));
3610 let drained = slow.finish_count();
3611
3612 let mut nested = Supervisor::new("nested").unwrap();
3613 nested.add_worker(
3614 ChildSpecification::one_shot_worker(slow)
3615 .with_shutdown_strategy(ShutdownStrategy::Graceful(Duration::from_secs(10))),
3616 );
3617
3618 let mut parent = Supervisor::new("parent")
3619 .unwrap()
3620 .with_shutdown_budget(Duration::from_millis(50));
3621 parent.add_worker(nested);
3622
3623 let (tx, handle) = run_supervisor_with_trigger(parent).await;
3624 tx.send(()).unwrap();
3625
3626 let result = join_supervisor(handle).await;
3627 assert!(
3628 drained.load(Ordering::SeqCst) == 1,
3629 "the nested subtree should have drained rather than being cut off by the parent's budget: {result:?}"
3630 );
3631 assert!(
3632 result.is_ok(),
3633 "the nested drain finished in time, so shutdown was clean: {result:?}"
3634 );
3635 }
3636 fn find_node<'a>(nodes: &'a [NodeSnapshot], name: &str) -> &'a NodeSnapshot {
3639 nodes.iter().find(|node| node.name == name).unwrap_or_else(|| {
3640 panic!(
3641 "no node named '{name}' among {:?}",
3642 nodes.iter().map(|n| &n.name).collect::<Vec<_>>()
3643 )
3644 })
3645 }
3646
3647 #[tokio::test]
3648 async fn snapshot_before_run_reports_a_registered_root() {
3649 let mut sup = Supervisor::new("test-sup")
3650 .unwrap()
3651 .with_restart_strategy(RestartStrategy::one_for_all().with_intensity_and_period(7, Duration::from_secs(11)))
3652 .with_auto_shutdown(AutoShutdown::AnySignificant)
3653 .with_shutdown_budget(Duration::from_secs(3));
3654 sup.add_worker(MockWorker::long_running("worker1"));
3655
3656 let snapshot = sup.tree_handle().snapshot();
3659
3660 assert_eq!(snapshot.root.name, "test-sup");
3661 assert_eq!(snapshot.root.kind, NodeKind::Supervisor);
3662 assert_eq!(snapshot.root.state, NodeState::Registered);
3663 assert_eq!(snapshot.root.process_id, None);
3664 assert_eq!(snapshot.root.process_name, None);
3665 assert_eq!(snapshot.root.started_at, None);
3666 assert_eq!(snapshot.root.uptime_ms, None);
3667 assert!(snapshot.root.children.is_empty(), "no children have been started yet");
3668
3669 let supervision = snapshot.root.supervision.expect("a supervisor reports its settings");
3670 assert_eq!(supervision.restart_mode, RestartMode::OneForAll);
3671 assert_eq!(supervision.restart_intensity, 7);
3672 assert_eq!(supervision.restart_period_ms, 11_000);
3673 assert_eq!(supervision.auto_shutdown, AutoShutdown::AnySignificant);
3674 assert_eq!(supervision.shutdown_budget_ms, Some(3_000));
3675 assert_eq!(supervision.dedicated_threads, None);
3676 assert_eq!(supervision.generation, 0);
3677
3678 assert_eq!(snapshot.totals.supervisors, 1);
3679 assert_eq!(snapshot.totals.workers, 0);
3680 assert_eq!(snapshot.totals.registered, 1);
3681 assert_eq!(snapshot.totals.max_depth, 1);
3682 }
3683
3684 #[tokio::test]
3685 async fn snapshot_lists_static_children_in_declaration_order() {
3686 let mut sup = Supervisor::new("test-sup").unwrap();
3687 sup.add_worker(MockWorker::long_running("worker1"));
3688 sup.add_worker(MockWorker::long_running("worker2"));
3689
3690 let tree = sup.tree_handle();
3691 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3692 wait_until("both children appear in the snapshot", || {
3693 tree.snapshot().root.children.len() == 2
3694 })
3695 .await;
3696
3697 let snapshot = tree.snapshot();
3698 assert_eq!(snapshot.root.state, NodeState::Running);
3699 assert!(snapshot.root.process_id.is_some());
3700 assert_eq!(snapshot.root.process_name.as_deref(), Some("test_sup"));
3701 assert_eq!(snapshot.root.supervision.expect("supervisor settings").generation, 1);
3702
3703 let names: Vec<&str> = snapshot.root.children.iter().map(|c| c.name.as_str()).collect();
3705 assert_eq!(names, vec!["worker1", "worker2"]);
3706
3707 for (child, expected_process) in snapshot
3708 .root
3709 .children
3710 .iter()
3711 .zip(["test_sup.worker1", "test_sup.worker2"])
3712 {
3713 assert_eq!(child.kind, NodeKind::Worker);
3714 assert_eq!(child.state, NodeState::Running);
3715 assert_eq!(child.restart, RestartType::Permanent);
3716 assert_eq!(child.restart_count, 0);
3717 assert_eq!(child.process_name.as_deref(), Some(expected_process));
3718 assert!(child.process_id.is_some(), "a running child has a process");
3719 assert!(child.children.is_empty(), "a worker has no children");
3720 assert!(child.supervision.is_none(), "a worker has no supervision settings");
3721 assert!(
3722 child.created_at <= child.started_at.expect("a running child has a start time"),
3723 "a child cannot start before it is created"
3724 );
3725 }
3726
3727 let first = snapshot.root.children[0].process_id;
3728 let second = snapshot.root.children[1].process_id;
3729 assert_ne!(first, second, "each child runs as its own process");
3730
3731 assert_eq!(snapshot.totals.supervisors, 1);
3732 assert_eq!(snapshot.totals.workers, 2);
3733 assert_eq!(snapshot.totals.running, 3);
3734 assert_eq!(snapshot.totals.max_depth, 2);
3735
3736 tx.send(()).unwrap();
3737 let _ = join_supervisor(handle).await;
3738 }
3739
3740 #[tokio::test]
3741 async fn one_for_one_restart_increments_only_the_failing_child() {
3742 let mut sup = Supervisor::new("test-sup")
3743 .unwrap()
3744 .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(5, Duration::from_secs(5)));
3745 sup.add_worker(MockWorker::failing("flapper", Duration::from_millis(20)));
3746 sup.add_worker(MockWorker::long_running("stable"));
3747
3748 let tree = sup.tree_handle();
3749 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3750
3751 wait_until("the failing child has appeared", || {
3752 tree.snapshot().root.children.len() == 2
3753 })
3754 .await;
3755 let before = tree.snapshot();
3756 let created_before = find_node(&before.root.children, "flapper").created_at;
3757 let process_before = find_node(&before.root.children, "flapper").process_id;
3758
3759 wait_until("the failing child has been restarted", || {
3760 find_node(&tree.snapshot().root.children, "flapper").restart_count >= 1
3761 })
3762 .await;
3763
3764 let after = tree.snapshot();
3765 let flapper = find_node(&after.root.children, "flapper");
3766 let stable = find_node(&after.root.children, "stable");
3767
3768 assert_eq!(stable.restart_count, 0, "a one-for-one restart leaves siblings alone");
3769 assert_ne!(flapper.process_id, process_before, "a restart is a new process");
3770 assert_eq!(
3771 flapper.created_at, created_before,
3772 "creation time is a property of the child, not of the incarnation"
3773 );
3774
3775 assert!(after.root.supervision.expect("supervisor settings").restarts_performed >= 1);
3777
3778 tx.send(()).unwrap();
3779 let _ = join_supervisor(handle).await;
3780 }
3781
3782 #[tokio::test]
3783 async fn one_for_all_restart_increments_every_restarted_child() {
3784 let mut sup = Supervisor::new("test-sup")
3785 .unwrap()
3786 .with_restart_strategy(RestartStrategy::one_for_all().with_intensity_and_period(5, Duration::from_secs(5)));
3787 sup.add_worker(MockWorker::failing("flapper", Duration::from_millis(20)));
3788 sup.add_worker(MockWorker::long_running("sibling"));
3789
3790 let tree = sup.tree_handle();
3791 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3792
3793 wait_until("both children have appeared", || {
3794 tree.snapshot().root.children.len() == 2
3795 })
3796 .await;
3797 let before = tree.snapshot();
3798 let sibling_created = find_node(&before.root.children, "sibling").created_at;
3799 let sibling_process = find_node(&before.root.children, "sibling").process_id;
3800
3801 wait_until("the group has been restarted", || {
3804 let snapshot = tree.snapshot();
3805 snapshot.root.children.len() == 2 && snapshot.root.children.iter().all(|child| child.restart_count >= 1)
3806 })
3807 .await;
3808
3809 let after = tree.snapshot();
3810 let sibling = find_node(&after.root.children, "sibling");
3811 assert_ne!(sibling.process_id, sibling_process, "a group restart is a new process");
3812 assert_eq!(
3813 sibling.created_at, sibling_created,
3814 "creation time survives a group restart"
3815 );
3816
3817 tx.send(()).unwrap();
3818 let _ = join_supervisor(handle).await;
3819 }
3820
3821 #[tokio::test]
3822 async fn temporary_child_survives_a_group_restart_as_a_tombstone() {
3823 let mut sup = Supervisor::new("test-sup")
3824 .unwrap()
3825 .with_restart_strategy(RestartStrategy::one_for_all().with_intensity_and_period(5, Duration::from_secs(5)));
3826 sup.add_worker(
3827 runtime::supervisable(MockWorker::completing("one-shot", Duration::from_millis(10)))
3828 .temporary()
3829 .build(),
3830 );
3831 sup.add_worker(MockWorker::failing("flapper", Duration::from_millis(30)));
3832
3833 let tree = sup.tree_handle();
3834 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3835
3836 wait_until("the temporary child has exited and the group has restarted", || {
3839 let snapshot = tree.snapshot();
3840 snapshot
3841 .root
3842 .children
3843 .iter()
3844 .any(|child| child.name == "one-shot" && child.state == NodeState::Exited)
3845 && snapshot
3846 .root
3847 .children
3848 .iter()
3849 .any(|child| child.name == "flapper" && child.restart_count >= 1)
3850 })
3851 .await;
3852
3853 let snapshot = tree.snapshot();
3854 let one_shot = find_node(&snapshot.root.children, "one-shot");
3855 assert_eq!(one_shot.state, NodeState::Exited);
3856 assert_eq!(one_shot.restart, RestartType::Temporary);
3857 assert_eq!(one_shot.restart_count, 0, "a temporary child is never brought back");
3858 assert!(one_shot.exited_at.is_some(), "a tombstone records when it exited");
3859 assert_eq!(one_shot.uptime_ms, None, "a node that isn't running has no uptime");
3860 assert!(snapshot.totals.exited >= 1);
3861
3862 tx.send(()).unwrap();
3863 let _ = join_supervisor(handle).await;
3864 }
3865
3866 #[tokio::test]
3867 async fn dynamic_children_leave_no_tombstones() {
3868 let mut sup = Supervisor::new("test-sup").unwrap();
3869 sup.add_worker(MockWorker::long_running("static-worker"));
3870
3871 let tree = sup.tree_handle();
3872 let sup_handle = sup.handle();
3873 let (tx, handle) = run_supervisor_with_trigger(sup).await;
3874 wait_until("the static child is running", || {
3875 tree.snapshot().root.children.len() == 1
3876 })
3877 .await;
3878
3879 for _ in 0..500 {
3883 sup_handle.spawn(MockWorker::completing("ephemeral", Duration::from_millis(1)));
3884 }
3885
3886 wait_until("every dynamic child has been reaped", || {
3887 sup_handle.active_children() == 0 && tree.snapshot().root.children.len() == 1
3888 })
3889 .await;
3890
3891 let snapshot = tree.snapshot();
3892 assert_eq!(snapshot.root.children.len(), 1, "only the static child remains");
3893 assert_eq!(snapshot.root.children[0].name, "static-worker");
3894
3895 tx.send(()).unwrap();
3896 let _ = join_supervisor(handle).await;
3897 }
3898
3899 #[tokio::test]
3900 async fn snapshot_nests_child_supervisors() {
3901 let mut child_sup = Supervisor::new("child-sup").unwrap();
3902 child_sup.add_worker(MockWorker::long_running("inner-worker"));
3903
3904 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3905 parent_sup.add_worker(MockWorker::long_running("outer-worker"));
3906 parent_sup.add_worker(child_sup);
3907
3908 let tree = parent_sup.tree_handle();
3909 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3910
3911 wait_until("the nested subtree is populated", || {
3912 let snapshot = tree.snapshot();
3913 snapshot
3914 .root
3915 .children
3916 .iter()
3917 .any(|child| child.name == "child-sup" && !child.children.is_empty())
3918 })
3919 .await;
3920
3921 let snapshot = tree.snapshot();
3922 let nested = find_node(&snapshot.root.children, "child-sup");
3923 assert_eq!(nested.kind, NodeKind::Supervisor);
3924 assert!(nested.supervision.is_some(), "a nested supervisor reports its settings");
3925 assert_eq!(nested.process_name.as_deref(), Some("parent_sup.child_sup"));
3926 assert_eq!(nested.children.len(), 1);
3927
3928 let grandchild = &nested.children[0];
3929 assert_eq!(grandchild.name, "inner-worker");
3930 assert_eq!(grandchild.kind, NodeKind::Worker);
3931 assert_eq!(
3932 grandchild.process_name.as_deref(),
3933 Some("parent_sup.child_sup.inner_worker")
3934 );
3935
3936 assert_eq!(snapshot.totals.supervisors, 2);
3937 assert_eq!(snapshot.totals.workers, 2);
3938 assert_eq!(snapshot.totals.max_depth, 3);
3939
3940 tx.send(()).unwrap();
3941 let _ = join_supervisor(handle).await;
3942 }
3943
3944 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3945 async fn dedicated_runtime_subtree_is_visible_across_the_thread_boundary() {
3946 let mut child_sup = Supervisor::new("child-sup")
3947 .unwrap()
3948 .with_dedicated_runtime(RuntimeConfiguration::single_threaded());
3949 child_sup.add_worker(MockWorker::long_running("inner-worker"));
3950
3951 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
3952 parent_sup.add_worker(child_sup);
3953
3954 let tree = parent_sup.tree_handle();
3956 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
3957
3958 wait_until("the dedicated-runtime subtree is visible", || {
3959 tree.snapshot()
3960 .root
3961 .children
3962 .iter()
3963 .any(|child| child.name == "child-sup" && !child.children.is_empty())
3964 })
3965 .await;
3966
3967 let probe = tree.clone();
3969 let snapshot = tokio::task::spawn_blocking(move || probe.snapshot())
3970 .await
3971 .expect("probe should not panic");
3972
3973 let nested = find_node(&snapshot.root.children, "child-sup");
3974 assert_eq!(nested.state, NodeState::Running);
3975 assert_eq!(
3976 nested.supervision.expect("supervisor settings").dedicated_threads,
3977 Some(1)
3978 );
3979 assert_eq!(nested.children.len(), 1);
3980
3981 assert_eq!(nested.process_name.as_deref(), Some("child_sup"));
3987 assert_eq!(
3988 nested.children[0].process_name.as_deref(),
3989 Some("child_sup.inner_worker")
3990 );
3991
3992 tx.send(()).unwrap();
3993 let _ = join_supervisor(handle).await;
3994 }
3995
3996 #[tokio::test]
3997 async fn a_running_node_always_has_a_process() {
3998 let mut child_sup = Supervisor::new("child-sup")
4002 .unwrap()
4003 .with_restart_strategy(RestartStrategy::new(RestartMode::OneForOne, 0, Duration::from_secs(5)));
4004 child_sup.add_worker(MockWorker::failing("inner", Duration::from_millis(10)));
4005
4006 let mut parent_sup = Supervisor::new("parent-sup")
4007 .unwrap()
4008 .with_restart_strategy(RestartStrategy::one_to_one().with_intensity_and_period(50, Duration::from_secs(5)));
4009 parent_sup.add_worker(child_sup);
4010
4011 let tree = parent_sup.tree_handle();
4012 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
4013
4014 for _ in 0..60 {
4015 let snapshot = tree.snapshot();
4016 assert_running_nodes_have_processes(&snapshot.root);
4017 tokio::time::sleep(Duration::from_millis(5)).await;
4018 }
4019
4020 wait_until("the nested supervisor has been restarted", || {
4021 find_node(&tree.snapshot().root.children, "child-sup").restart_count >= 1
4022 })
4023 .await;
4024
4025 tx.send(()).unwrap();
4026 let _ = join_supervisor(handle).await;
4027 }
4028
4029 fn assert_running_nodes_have_processes(node: &NodeSnapshot) {
4030 if node.state == NodeState::Running {
4031 assert!(
4032 node.process_id.is_some() && node.process_name.is_some(),
4033 "node '{}' reports as running but names no process",
4034 node.name
4035 );
4036 assert!(
4037 node.uptime_ms.is_some(),
4038 "node '{}' is running but has no uptime",
4039 node.name
4040 );
4041 } else {
4042 assert!(
4043 node.uptime_ms.is_none(),
4044 "node '{}' is not running but reports an uptime",
4045 node.name
4046 );
4047 }
4048
4049 for child in &node.children {
4050 assert_running_nodes_have_processes(child);
4051 }
4052 }
4053
4054 #[tokio::test]
4055 async fn snapshot_after_the_run_reports_a_stopped_tree() {
4056 fn assert_send_sync<T: Send + Sync + 'static>() {}
4057 assert_send_sync::<SupervisionTreeHandle>();
4058
4059 let mut sup = Supervisor::new("test-sup").unwrap();
4060 sup.add_worker(MockWorker::long_running("worker1"));
4061
4062 let tree = sup.tree_handle();
4063 let (tx, handle) = run_supervisor_with_trigger(sup).await;
4064 wait_until("the child is running", || tree.snapshot().root.children.len() == 1).await;
4065
4066 tx.send(()).unwrap();
4067 join_supervisor(handle).await.unwrap();
4068
4069 let snapshot = tree.snapshot();
4071 assert_eq!(snapshot.root.state, NodeState::Registered);
4072 assert!(snapshot.root.children.is_empty());
4073 assert_eq!(
4074 snapshot.root.supervision.expect("supervisor settings").generation,
4075 1,
4076 "the generation count outlives the run"
4077 );
4078
4079 assert!(snapshot.root.process_id.is_some(), "a stopped node says what it ran as");
4082 assert_eq!(snapshot.root.process_name.as_deref(), Some("test_sup"));
4083 assert!(snapshot.root.started_at.is_some());
4084 assert_eq!(snapshot.root.uptime_ms, None, "a node that isn't running has no uptime");
4085 }
4086
4087 #[tokio::test]
4088 async fn an_exited_nested_supervisor_says_what_it_ran_as() {
4089 let mut child_sup = Supervisor::new("child-sup")
4092 .unwrap()
4093 .with_auto_shutdown(AutoShutdown::AnySignificant);
4094 child_sup.add_worker(
4095 runtime::supervisable(MockWorker::completing("one-shot", Duration::from_millis(10)))
4096 .temporary()
4097 .with_significant(true)
4098 .build(),
4099 );
4100
4101 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
4103 parent_sup.add_worker(ChildSpecification::from(child_sup).with_restart_type(RestartType::Temporary));
4104
4105 let tree = parent_sup.tree_handle();
4106 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
4107
4108 let mut running = None;
4111 wait_until("the nested supervisor is running", || {
4112 let snapshot = tree.snapshot();
4113 let nested = find_node(&snapshot.root.children, "child-sup");
4114 if nested.state == NodeState::Running {
4115 running = Some((nested.process_id, nested.process_name.clone(), nested.started_at));
4116 return true;
4117 }
4118 false
4119 })
4120 .await;
4121 let (process_id, process_name, started_at) = running.expect("observed running");
4122 assert!(process_id.is_some(), "a running supervisor has a process");
4123
4124 wait_until("the nested supervisor has exited", || {
4125 find_node(&tree.snapshot().root.children, "child-sup").state == NodeState::Exited
4126 })
4127 .await;
4128
4129 let snapshot = tree.snapshot();
4130 let nested = find_node(&snapshot.root.children, "child-sup");
4131 assert_eq!(nested.process_id, process_id, "an exited supervisor keeps its identity");
4132 assert_eq!(nested.process_name, process_name);
4133 assert_eq!(nested.started_at, started_at);
4134 assert!(nested.exited_at.is_some(), "a tombstone records when it exited");
4135 assert_eq!(nested.uptime_ms, None, "a node that isn't running has no uptime");
4136 assert!(
4137 nested.resource_group.is_some(),
4138 "and the group its allocations were attributed to"
4139 );
4140
4141 tx.send(()).unwrap();
4142 let _ = join_supervisor(handle).await;
4143 }
4144
4145 #[tokio::test]
4146 async fn initialization_failure_leaves_nothing_running() {
4147 let mut sup = Supervisor::new("test-sup").unwrap();
4148 sup.add_worker(MockWorker::init_failure("broken"));
4149
4150 let tree = sup.tree_handle();
4151
4152 let mut sup = sup;
4155 let (_tx, rx) = oneshot::channel::<()>();
4156 let result = timeout(Duration::from_secs(2), sup.run_with_shutdown(rx))
4157 .await
4158 .expect("supervisor should exit promptly");
4159 assert!(matches!(result, Err(SupervisorError::FailedToInitialize { .. })));
4160
4161 let snapshot = tree.snapshot();
4163 assert_eq!(snapshot.root.state, NodeState::Registered);
4164 assert_eq!(snapshot.totals.running, 0);
4165 }
4166
4167 #[tokio::test]
4168 async fn resources_are_attributed_to_supervisors_not_workers() {
4169 let mut child_sup = Supervisor::new("child-sup").unwrap();
4170 child_sup.add_worker(MockWorker::long_running("inner-worker"));
4171
4172 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
4173 parent_sup.add_worker(child_sup);
4174
4175 let tree = parent_sup.tree_handle();
4176 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
4177 wait_until("the nested subtree is populated", || {
4178 tree.snapshot()
4179 .root
4180 .children
4181 .iter()
4182 .any(|child| !child.children.is_empty())
4183 })
4184 .await;
4185
4186 let snapshot = tree.snapshot();
4187
4188 assert!(!snapshot.resource_tracking_enabled);
4192
4193 let nested = find_node(&snapshot.root.children, "child-sup");
4194 assert!(
4195 nested.resources.is_some(),
4196 "a supervisor owns the resource group named by its own process"
4197 );
4198 assert_eq!(nested.resource_group.as_deref(), nested.process_name.as_deref());
4199
4200 let worker = &nested.children[0];
4203 assert!(worker.resources.is_none(), "a worker owns no resource group");
4204 assert_eq!(worker.resource_group.as_deref(), nested.process_name.as_deref());
4205
4206 tx.send(()).unwrap();
4207 let _ = join_supervisor(handle).await;
4208 }
4209
4210 #[tokio::test]
4211 async fn snapshot_serializes_to_the_expected_shape() {
4212 let mut child_sup = Supervisor::new("child-sup").unwrap();
4213 child_sup.add_worker(MockWorker::long_running("inner-worker"));
4214
4215 let mut parent_sup = Supervisor::new("parent-sup").unwrap();
4216 parent_sup.add_worker(child_sup);
4217
4218 let tree = parent_sup.tree_handle();
4219 let (tx, handle) = run_supervisor_with_trigger(parent_sup).await;
4220 wait_until("the nested subtree is populated", || {
4221 tree.snapshot()
4222 .root
4223 .children
4224 .iter()
4225 .any(|child| !child.children.is_empty())
4226 })
4227 .await;
4228
4229 let value = serde_json::to_value(tree.snapshot()).expect("snapshot should serialize");
4230
4231 assert!(value["captured_at"].is_u64(), "a timestamp is a plain integer");
4233 assert_eq!(value["root"]["kind"], "supervisor");
4234 assert_eq!(value["root"]["state"], "running");
4235 assert_eq!(value["root"]["restart"], "permanent");
4236 assert_eq!(value["root"]["supervision"]["restart_mode"], "one_for_one");
4237 assert_eq!(value["root"]["supervision"]["auto_shutdown"], "never");
4238
4239 let nested = &value["root"]["children"][0];
4240 assert_eq!(nested["kind"], "supervisor");
4241 let worker = &nested["children"][0];
4242 assert_eq!(worker["kind"], "worker");
4243 assert!(
4244 worker["children"]
4245 .as_array()
4246 .expect("children is always an array")
4247 .is_empty(),
4248 "children is present and empty for a worker, so consumers can recurse unconditionally"
4249 );
4250 assert!(worker.get("resources").is_none(), "a worker reports no resource usage");
4251 assert!(
4252 worker.get("supervision").is_none(),
4253 "a worker reports no supervision settings"
4254 );
4255 assert!(
4256 worker.get("children_truncated").is_none(),
4257 "an untruncated node omits the truncation flag"
4258 );
4259
4260 tx.send(()).unwrap();
4261 let _ = join_supervisor(handle).await;
4262 }
4263
4264 #[tokio::test]
4265 async fn a_worker_that_drives_its_own_supervisor_is_shown_as_its_parent() {
4266 struct SupervisorDrivingWorker;
4270
4271 #[async_trait]
4272 impl Supervisable for SupervisorDrivingWorker {
4273 fn name(&self) -> &str {
4274 "driver"
4275 }
4276
4277 fn shutdown_strategy(&self) -> ShutdownStrategy {
4278 ShutdownStrategy::Graceful(Duration::MAX)
4279 }
4280
4281 async fn initialize(
4282 &self, process_shutdown: ShutdownHandle,
4283 ) -> Result<SupervisorFuture, InitializationError> {
4284 Ok(Box::pin(async move {
4285 let mut inner = Supervisor::new("inner-sup").expect("valid name");
4286 inner.add_worker(MockWorker::long_running("inner-worker"));
4287 inner
4288 .run_with_shutdown_inner(process_shutdown, None)
4289 .await
4290 .map_err(Into::into)
4291 }))
4292 }
4293 }
4294
4295 let mut sup = Supervisor::new("test-sup").unwrap();
4296 sup.add_worker(SupervisorDrivingWorker);
4297
4298 let tree = sup.tree_handle();
4299 let (tx, handle) = run_supervisor_with_trigger(sup).await;
4300
4301 wait_until("the worker's own supervisor has attached itself", || {
4302 let snapshot = tree.snapshot();
4303 snapshot
4304 .root
4305 .children
4306 .first()
4307 .is_some_and(|child| !child.children.is_empty())
4308 })
4309 .await;
4310
4311 let snapshot = tree.snapshot();
4312 let driver = find_node(&snapshot.root.children, "driver");
4313 assert_eq!(
4314 driver.kind,
4315 NodeKind::Worker,
4316 "the worker is what its parent supervises"
4317 );
4318
4319 let inner_sup = find_node(&driver.children, "inner-sup");
4320 assert_eq!(inner_sup.kind, NodeKind::Supervisor);
4321 assert_eq!(inner_sup.children.len(), 1);
4322 assert_eq!(inner_sup.children[0].name, "inner-worker");
4323 assert_eq!(snapshot.totals.max_depth, 4);
4324
4325 tx.send(()).unwrap();
4326 let _ = join_supervisor(handle).await;
4327 }
4328
4329 #[tokio::test]
4330 async fn mixed_child_traffic_keeps_the_roster_consistent() {
4331 let mut sup = Supervisor::new("test-sup").unwrap().with_restart_strategy(
4335 RestartStrategy::one_to_one().with_intensity_and_period(100, Duration::from_secs(5)),
4336 );
4337 sup.add_worker(MockWorker::long_running("stable"));
4338 sup.add_worker(MockWorker::failing("flapper", Duration::from_millis(5)));
4339
4340 let tree = sup.tree_handle();
4341 let sup_handle = sup.handle();
4342 let (tx, handle) = run_supervisor_with_trigger(sup).await;
4343
4344 for _ in 0..25 {
4345 sup_handle.spawn(MockWorker::completing("ephemeral", Duration::from_millis(2)));
4346 sup_handle.spawn(MockWorker::long_running("lingering"));
4347 let snapshot = tree.snapshot();
4348 assert_running_nodes_have_processes(&snapshot.root);
4349 tokio::time::sleep(Duration::from_millis(2)).await;
4350 }
4351
4352 wait_until("the flapping child has restarted several times", || {
4353 find_node(&tree.snapshot().root.children, "flapper").restart_count >= 3
4354 })
4355 .await;
4356
4357 tx.send(()).unwrap();
4358 let _ = join_supervisor(handle).await;
4359 }
4360}