1use std::{
2 collections::{HashMap, VecDeque},
3 error::Error,
4 fmt, io,
5 path::PathBuf,
6 process::{ExitStatus, Stdio},
7 sync::{Arc, Mutex, OnceLock},
8 time::{Duration, SystemTime, UNIX_EPOCH},
9};
10
11use cortexkit_log::Retention;
12use serde_json::Value;
13use subc_control::{
14 ClientControlPush, LiveSpawn, ModuleProtocol, RouteCloseReason, SpawnCursor, SpawnEvent,
15 SpawnEventKind, SpawnSnapshot, SupervisorHealthStatus, TerminalDisposition, TerminalExitKind,
16};
17use subc_protocol::{
18 manifest::{SelfSignalKind, SignalAnchor},
19 session::{
20 HealthReport, HealthStatus, ModuleControlCommand, ModuleControlRequest,
21 MODULE_CONTROL_OP_HEALTH_CHECK,
22 },
23 Flags, FrameType, Priority, SUBC_LAUNCH_NONCE_ENV, SUBC_MODULE_ID_ENV,
24};
25use tokio::{
26 process::{Child, Command},
27 sync::{mpsc, oneshot, watch, Mutex as AsyncMutex},
28 task::JoinHandle,
29 time::{sleep, sleep_until, timeout, timeout_at, Instant},
30};
31use tracing::{debug, error, info, warn};
32
33use crate::{
34 child_roster::ChildRoster,
35 daemon_config::{
36 CAPTURE_KEEP_ENV, CAPTURE_MAX_AGE_DAYS_ENV, CAPTURE_MAX_FILE_MB_ENV, CK_LOG_ENV,
37 },
38 forwarding::{
39 CloseReason, ForwardingError, ForwardingTable, GoodbyeTarget, ModuleControlRpcOutcome,
40 ModuleDrainTarget, PendingModuleControlRpc,
41 },
42 provenance::{spawned_file_identity, ExecutableIdentityProbe, SpawnedFileIdentity},
43 registry::{ConnectionId, RegistryError},
44 stderr_tail::{
45 pump_stderr_to, pump_stdout_to, ChildOutputSink, StderrRing, StderrTailConfig,
46 StderrTailSnapshot,
47 },
48 terminal_ring::{TerminalHistorySnapshot, TerminalRecord, TerminalRing, TerminalRingConfig},
49 Frame, FrameSink, Registry,
50};
51
52#[path = "supervise_swap.rs"]
53mod swap;
54
55pub const SUBC_ARG: &str = "--subc";
61
62const DEFAULT_MAX_RESTARTS: u32 = 3;
63const DEFAULT_BACKOFF: Duration = Duration::from_millis(100);
64const DEFAULT_MAX_BACKOFF: Duration = Duration::from_secs(30);
65const DEFAULT_RESTART_WINDOW: Duration = Duration::from_secs(600);
69pub const DEFAULT_DRAIN_TIMEOUT: Duration = Duration::from_secs(30);
80const REGISTRY_RELEASE_TIMEOUT: Duration = Duration::from_secs(1);
81const REGISTRY_RELEASE_POLL: Duration = Duration::from_millis(10);
82const STDERR_PUMP_DRAIN_TIMEOUT: Duration = Duration::from_millis(250);
83pub const SPAWN_EVENT_RING_CAPACITY: usize = 4096;
85const SPAWN_SUBSCRIBER_BUFFER: usize = SPAWN_EVENT_RING_CAPACITY + 1;
86
87struct SupervisedChild {
88 child: Child,
89 #[cfg(target_os = "linux")]
92 module_id: String,
93 #[cfg(target_os = "linux")]
94 cgroup_placement: Option<subc_cgroup::Placement>,
95 stdout_pump: Option<JoinHandle<()>>,
96 stderr_pump: Option<JoinHandle<()>>,
97 stderr_ring: Arc<Mutex<StderrRing>>,
98 spawned_at_ms: u64,
99 spawned_from: PathBuf,
100 spawned_file_identity: Option<SpawnedFileIdentity>,
101 process_start_time: Option<u64>,
102 process_identity: Option<ProcessIdentity>,
103 pid: u32,
104 roster_guard: Option<crate::child_roster::RosterGuard>,
107}
108
109impl SupervisedChild {
110 fn id(&self) -> Option<u32> {
111 Some(self.pid)
112 }
113
114 fn process_identity(&self) -> Option<ProcessIdentity> {
115 self.process_identity
116 }
117
118 async fn wait(&mut self) -> io::Result<ExitStatus> {
119 let result = self.child.wait().await;
120 if result.is_ok() {
121 self.roster_guard = None;
124 }
125 #[cfg(target_os = "linux")]
126 if result.is_ok() {
127 if let Some(placement) = self.cgroup_placement.take() {
128 remove_module_cgroup(&placement, &self.module_id);
129 }
130 }
131 result
132 }
133
134 fn start_kill(&mut self) -> io::Result<()> {
135 self.child.start_kill()
136 }
137
138 async fn drain_stderr(&mut self, module_id: &str) {
139 if let Some(mut pump) = self.stdout_pump.take() {
140 match timeout(STDERR_PUMP_DRAIN_TIMEOUT, &mut pump).await {
141 Ok(Ok(())) => {}
142 Ok(Err(error)) => {
143 warn!(module_id, error = %error, "stdout pump ended unexpectedly");
144 }
145 Err(_) => {
146 pump.abort();
147 warn!(
148 module_id,
149 waited = ?STDERR_PUMP_DRAIN_TIMEOUT,
150 "stdout pump did not drain before restart; stopped it before the next process"
151 );
152 }
153 }
154 }
155
156 let Some(mut pump) = self.stderr_pump.take() else {
157 return;
158 };
159 match timeout(STDERR_PUMP_DRAIN_TIMEOUT, &mut pump).await {
160 Ok(Ok(())) => {}
161 Ok(Err(err)) => {
162 self.stderr_ring
163 .lock()
164 .unwrap_or_else(|poisoned| poisoned.into_inner())
165 .mark_incomplete(format!("stderr pump ended unexpectedly: {err}"));
166 warn!(module_id, error = %err, "stderr pump ended before clean EOF");
167 }
168 Err(_) => {
169 pump.abort();
170 self.stderr_ring
171 .lock()
172 .unwrap_or_else(|poisoned| poisoned.into_inner())
173 .mark_incomplete(format!(
174 "stderr pump did not reach EOF within {:?} before restart",
175 STDERR_PUMP_DRAIN_TIMEOUT
176 ));
177 warn!(
178 module_id,
179 waited = ?STDERR_PUMP_DRAIN_TIMEOUT,
180 "stderr pump did not drain before restart; stopped it before marking the new process"
181 );
182 }
183 }
184 }
185}
186
187fn registration_release_events() -> &'static watch::Sender<u64> {
188 static EVENTS: OnceLock<watch::Sender<u64>> = OnceLock::new();
189 EVENTS.get_or_init(|| {
190 let (sender, _receiver) = watch::channel(0);
191 sender
192 })
193}
194
195pub(crate) fn notify_registration_release() {
196 let events = registration_release_events();
197 let next_generation = (*events.borrow()).wrapping_add(1);
198 events.send_replace(next_generation);
199}
200
201#[derive(Debug, Clone, PartialEq, Eq)]
203pub struct ModuleSpec {
204 pub module_id: String,
205 pub program: PathBuf,
206 pub args: Vec<String>,
207 pub env: Vec<(String, String)>,
208 pub reserved: bool,
213 pub reserved_prefixes: Vec<String>,
218 pub protocol: ModuleProtocol,
234 pub overlap: ModuleOverlap,
239}
240
241#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
248pub enum ModuleOverlap {
249 #[default]
251 Exclusive,
252 Safe,
264}
265
266impl ModuleOverlap {
267 pub fn as_str(self) -> &'static str {
268 match self {
269 Self::Exclusive => "exclusive",
270 Self::Safe => "safe",
271 }
272 }
273}
274
275pub const SUBC_SPAWN_ROLE_ENV: &str = "SUBC_SPAWN_ROLE";
285pub const SPAWN_ROLE_SWAP_CANDIDATE: &str = "swap_candidate";
287pub const DEFAULT_SWAP_READY_TIMEOUT: Duration = Duration::from_secs(100);
292
293#[derive(Debug, Clone, Copy, PartialEq, Eq)]
311pub struct RestartPolicy {
312 pub max_restarts: u32,
313 pub backoff: Duration,
316 pub max_backoff: Duration,
318 pub window: Duration,
322}
323
324impl RestartPolicy {
325 pub fn new(max_restarts: u32, backoff: Duration) -> Self {
329 Self {
330 max_restarts,
331 backoff,
332 max_backoff: DEFAULT_MAX_BACKOFF,
333 window: DEFAULT_RESTART_WINDOW,
334 }
335 }
336
337 pub fn with_max_backoff(mut self, max_backoff: Duration) -> Self {
338 self.max_backoff = max_backoff;
339 self
340 }
341
342 pub fn with_window(mut self, window: Duration) -> Self {
343 self.window = window;
344 self
345 }
346
347 fn delay_for_restart(&self, restart_in_window: u32) -> Duration {
352 if self.backoff.is_zero() || self.max_backoff.is_zero() {
353 return Duration::ZERO;
354 }
355
356 let mut delay = self.backoff;
357 for _ in 0..restart_in_window {
358 if delay >= self.max_backoff {
359 return self.max_backoff;
360 }
361 delay = delay
362 .checked_mul(10)
363 .unwrap_or(self.max_backoff)
364 .min(self.max_backoff);
365 }
366 delay.min(self.max_backoff)
367 }
368
369 fn budget_exhausted_detail(&self) -> String {
374 format!(
375 "crash budget exhausted: max_restarts={} within window_secs={}",
376 self.max_restarts,
377 self.window.as_secs()
378 )
379 }
380}
381
382impl Default for RestartPolicy {
383 fn default() -> Self {
384 Self {
385 max_restarts: DEFAULT_MAX_RESTARTS,
386 backoff: DEFAULT_BACKOFF,
387 max_backoff: DEFAULT_MAX_BACKOFF,
388 window: DEFAULT_RESTART_WINDOW,
389 }
390 }
391}
392
393#[derive(Debug, Clone, Copy, PartialEq, Eq)]
394struct CrashRestartSchedule {
395 restart_in_window: u32,
396 delay: Duration,
397}
398
399fn daemon_will_restart(
406 state: &mut SupervisorSnapshot,
407 policy: &RestartPolicy,
408 now: Instant,
409) -> bool {
410 state.enabled && state.crash_restarts_in_window(policy.window, now) < policy.max_restarts
411}
412
413const DEFAULT_HEALTH_CADENCE: Duration = Duration::from_secs(30);
414const DEFAULT_HEALTH_DEADLINE: Duration = Duration::from_secs(5);
415const DEFAULT_HEALTH_FAILURE_THRESHOLD: u32 = 3;
416const MAX_HEALTH_METRICS_BYTES: usize = 16 * 1024;
417
418#[derive(Debug, Clone, Copy, PartialEq, Eq)]
419pub enum HealthAction {
420 Report,
421 Restart,
422 Alert,
423}
424
425impl fmt::Display for HealthAction {
426 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
427 f.write_str(match self {
428 Self::Report => "report",
429 Self::Restart => "restart",
430 Self::Alert => "alert",
431 })
432 }
433}
434
435#[derive(Debug, Clone, Copy, PartialEq, Eq)]
436pub struct HealthConfig {
437 pub cadence: Duration,
438 pub deadline: Duration,
439 pub failure_threshold: u32,
440 pub on_degraded: HealthAction,
441 pub on_failing: HealthAction,
442 pub critical: bool,
443}
444
445impl Default for HealthConfig {
446 fn default() -> Self {
447 Self {
448 cadence: DEFAULT_HEALTH_CADENCE,
449 deadline: DEFAULT_HEALTH_DEADLINE,
450 failure_threshold: DEFAULT_HEALTH_FAILURE_THRESHOLD,
451 on_degraded: HealthAction::Report,
452 on_failing: HealthAction::Report,
453 critical: false,
454 }
455 }
456}
457
458#[derive(Debug, Clone, PartialEq)]
476pub struct ModuleHealthStatus {
477 pub status: SupervisorHealthStatus,
478 pub last_probe_ms: Option<u64>,
479 pub detail: Option<String>,
480 pub metrics: Option<Value>,
481 pub consecutive_failures: u32,
482 pub late_answer_count: u64,
485 pub last_late_answer_latency_ms: Option<u64>,
487 pub last_action: Option<String>,
488 pub last_action_ms: Option<u64>,
492}
493
494impl Default for ModuleHealthStatus {
495 fn default() -> Self {
496 Self {
497 status: SupervisorHealthStatus::Unknown,
498 last_probe_ms: None,
499 detail: None,
500 metrics: None,
501 consecutive_failures: 0,
502 late_answer_count: 0,
503 last_late_answer_latency_ms: None,
504 last_action: None,
505 last_action_ms: None,
506 }
507 }
508}
509
510#[derive(Debug, Clone, Copy, PartialEq, Eq)]
512pub enum ModuleState {
513 Starting,
514 Running,
515 Unresponsive,
516 Restarting,
517 Draining,
518 Stopped,
519 Failed,
520 Disabled,
521}
522
523impl fmt::Display for ModuleState {
524 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
525 f.write_str(match self {
526 Self::Starting => "starting",
527 Self::Running => "running",
528 Self::Unresponsive => "unresponsive",
529 Self::Restarting => "restarting",
530 Self::Draining => "draining",
531 Self::Stopped => "stopped",
532 Self::Failed => "failed",
533 Self::Disabled => "disabled",
534 })
535 }
536}
537
538#[derive(Debug, Clone, Copy, PartialEq, Eq)]
540pub enum ExitKind {
541 Clean,
542 Crash,
543 DeliberateSeverance,
544}
545
546impl From<ExitKind> for TerminalExitKind {
547 fn from(kind: ExitKind) -> Self {
548 match kind {
549 ExitKind::Clean => Self::Clean,
550 ExitKind::Crash => Self::Crash,
551 ExitKind::DeliberateSeverance => Self::DeliberateSeverance,
552 }
553 }
554}
555
556#[derive(Debug, Clone, Copy, PartialEq, Eq)]
559pub(crate) struct ProcessIdentity {
560 pub(crate) pid: u32,
561 pub(crate) start_time: u64,
562}
563
564#[derive(Debug, Clone, PartialEq, Eq)]
566pub struct ExitReport {
567 pub kind: ExitKind,
568 pub code: Option<i32>,
569 pub signal: Option<i32>,
570 pub at_ms: u64,
571}
572
573#[derive(Debug, Clone, PartialEq)]
576pub struct ModuleStatus {
577 pub module_id: String,
578 pub state: ModuleState,
579 pub enabled: bool,
580 pub process_alive: bool,
581 pub registration_active: bool,
582 pub protocol: ModuleProtocol,
585 pub live: bool,
596 pub restart_count: u32,
600 pub lifetime_restarts: u32,
604 pub spawn_generation: u64,
605 pub max_restarts: u32,
610 pub restart_window: Duration,
614 pub drain_timeout: Duration,
618 pub restart_backoff: Duration,
619 pub restart_max_backoff: Duration,
620 pub pid: Option<u32>,
621 pub spawned_at_ms: Option<u64>,
622 pub spawned_from: Option<PathBuf>,
623 pub process_start_time: Option<u64>,
624 pub last_exit: Option<ExitReport>,
625 pub health: ModuleHealthStatus,
626}
627
628#[derive(Debug, Clone, PartialEq)]
629struct SupervisorSnapshot {
630 state: ModuleState,
631 enabled: bool,
632 process_alive: bool,
633 crash_restarts: VecDeque<Instant>,
639 lifetime_restarts: u32,
640 spawn_generation: u64,
649 pid: Option<u32>,
650 spawned_at_ms: Option<u64>,
651 spawned_from: Option<PathBuf>,
652 spawned_file_identity: Option<SpawnedFileIdentity>,
653 process_start_time: Option<u64>,
654 deliberate_severance: Option<ProcessIdentity>,
655 last_exit: Option<ExitReport>,
656 health: ModuleHealthStatus,
657 in_alternate_slot: bool,
662}
663
664impl SupervisorSnapshot {
665 fn starting() -> Self {
666 Self::new(ModuleState::Starting, true)
667 }
668
669 fn disabled() -> Self {
670 Self::new(ModuleState::Disabled, false)
671 }
672
673 fn failed() -> Self {
674 Self::new(ModuleState::Failed, true)
675 }
676
677 fn crash_restarts_in_window(&mut self, window: Duration, now: Instant) -> u32 {
681 while let Some(oldest) = self.crash_restarts.front() {
682 if now.duration_since(*oldest) > window {
683 self.crash_restarts.pop_front();
684 } else {
685 break;
686 }
687 }
688 u32::try_from(self.crash_restarts.len()).unwrap_or(u32::MAX)
689 }
690
691 fn record_crash_restart(&mut self, policy: &RestartPolicy, now: Instant) {
697 self.crash_restarts.push_back(now);
698 while self.crash_restarts.len() > policy.max_restarts as usize {
699 self.crash_restarts.pop_front();
700 }
701 self.lifetime_restarts += 1;
702 }
703
704 fn next_crash_restart(
708 &mut self,
709 policy: &RestartPolicy,
710 now: Instant,
711 ) -> Option<CrashRestartSchedule> {
712 let restart_in_window = self.crash_restarts_in_window(policy.window, now);
713 if restart_in_window >= policy.max_restarts {
714 return None;
715 }
716 self.record_crash_restart(policy, now);
717 Some(CrashRestartSchedule {
718 restart_in_window,
719 delay: policy.delay_for_restart(restart_in_window),
720 })
721 }
722
723 fn clear_crash_restarts(&mut self) {
728 self.crash_restarts.clear();
729 }
730
731 fn new(state: ModuleState, enabled: bool) -> Self {
732 Self {
733 state,
734 enabled,
735 process_alive: false,
736 crash_restarts: VecDeque::new(),
737 lifetime_restarts: 0,
738 spawn_generation: 0,
739 pid: None,
740 spawned_at_ms: None,
741 spawned_from: None,
742 spawned_file_identity: None,
743 process_start_time: None,
744 deliberate_severance: None,
745 last_exit: None,
746 health: ModuleHealthStatus::default(),
747 in_alternate_slot: false,
748 }
749 }
750}
751
752type SharedSnapshot = Arc<Mutex<SupervisorSnapshot>>;
753
754type SpawnSubscriberKey = (ConnectionId, u64);
755
756#[derive(Debug)]
757struct SpawnSubscriber {
758 version: u8,
759 frames: mpsc::Sender<Frame>,
760}
761
762#[derive(Debug)]
763struct SpawnEventState {
764 daemon_incarnation: String,
765 seq: u64,
766 capacity: usize,
767 live: HashMap<String, LiveSpawn>,
768 generations: HashMap<String, u64>,
769 events: VecDeque<SpawnEvent>,
770 subscribers: HashMap<SpawnSubscriberKey, SpawnSubscriber>,
771}
772
773impl Default for SpawnEventState {
774 fn default() -> Self {
775 Self {
776 daemon_incarnation: "unconfigured".to_string(),
777 seq: 0,
778 capacity: SPAWN_EVENT_RING_CAPACITY,
779 live: HashMap::new(),
780 generations: HashMap::new(),
781 events: VecDeque::new(),
782 subscribers: HashMap::new(),
783 }
784 }
785}
786
787#[derive(Debug, Clone, Default)]
788struct SpawnEventFeed(Arc<Mutex<SpawnEventState>>);
789
790#[derive(Debug, Clone, PartialEq, Eq)]
791pub(crate) enum SpawnSubscribeRefusal {
792 ForeignIncarnation { current: String },
793 TooOld { oldest: SpawnCursor },
794 Frame(String),
795}
796
797impl SpawnEventFeed {
798 fn configure_incarnation(&self, daemon_incarnation: String) {
799 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
800 state.daemon_incarnation = daemon_incarnation;
801 state.seq = 0;
802 state.live.clear();
803 state.generations.clear();
804 state.events.clear();
805 state.subscribers.clear();
806 }
807
808 fn cursor(state: &SpawnEventState) -> SpawnCursor {
809 SpawnCursor {
810 daemon_incarnation: state.daemon_incarnation.clone(),
811 seq: state.seq,
812 }
813 }
814
815 fn snapshot(&self) -> SpawnSnapshot {
816 let state = self.0.lock().unwrap_or_else(|p| p.into_inner());
817 let mut live = state.live.values().cloned().collect::<Vec<_>>();
818 live.sort_by(|left, right| left.module_id.cmp(&right.module_id));
819 SpawnSnapshot {
820 cursor: Self::cursor(&state),
821 ring_bound: state.capacity as u64,
822 live,
823 }
824 }
825
826 fn emit_spawned(&self, module_id: &str, pid: u32, spawned_at_ms: u64) -> u64 {
827 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
828 let generation = state
829 .generations
830 .get(module_id)
831 .copied()
832 .unwrap_or(0)
833 .checked_add(1)
834 .expect("spawn generation exhausted");
835 state.generations.insert(module_id.to_string(), generation);
836 let live = LiveSpawn {
837 module_id: module_id.to_string(),
838 spawn_generation: generation,
839 pid,
840 spawned_at_ms,
841 };
842 state.live.insert(module_id.to_string(), live);
843 Self::emit_locked(
844 &mut state,
845 SpawnEventKind::Spawned,
846 module_id.to_string(),
847 generation,
848 pid,
849 None,
850 None,
851 );
852 generation
853 }
854
855 fn emit_exited(&self, module_id: &str, exit_code: Option<i32>, exit_signal: Option<i32>) {
856 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
857 let Some(live) = state.live.remove(module_id) else {
858 warn!(
859 module_id,
860 "terminal record had no live spawn event identity"
861 );
862 return;
863 };
864 Self::emit_locked(
865 &mut state,
866 SpawnEventKind::Exited,
867 module_id.to_string(),
868 live.spawn_generation,
869 live.pid,
870 exit_code,
871 exit_signal,
872 );
873 }
874
875 fn emit_superseded_exited(
882 &self,
883 module_id: &str,
884 spawn_generation: u64,
885 pid: u32,
886 exit_code: Option<i32>,
887 exit_signal: Option<i32>,
888 ) {
889 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
890 if state
891 .live
892 .get(module_id)
893 .is_some_and(|live| live.spawn_generation == spawn_generation)
894 {
895 state.live.remove(module_id);
896 }
897 Self::emit_locked(
898 &mut state,
899 SpawnEventKind::Exited,
900 module_id.to_string(),
901 spawn_generation,
902 pid,
903 exit_code,
904 exit_signal,
905 );
906 }
907
908 #[allow(clippy::too_many_arguments)]
909 fn emit_locked(
910 state: &mut SpawnEventState,
911 kind: SpawnEventKind,
912 module_id: String,
913 spawn_generation: u64,
914 pid: u32,
915 exit_code: Option<i32>,
916 exit_signal: Option<i32>,
917 ) {
918 state.seq = state
919 .seq
920 .checked_add(1)
921 .expect("spawn event sequence exhausted");
922 let event = SpawnEvent {
923 cursor: Self::cursor(state),
924 kind,
925 module_id,
926 spawn_generation,
927 pid,
928 exit_code,
929 exit_signal,
930 };
931 state.events.push_back(event.clone());
932 while state.events.len() > state.capacity {
933 state.events.pop_front();
934 }
935 let body = match serde_json::to_vec(&event) {
936 Ok(body) => body,
937 Err(error) => {
938 error!(%error, "failed to serialize supervisor spawn event");
939 return;
940 }
941 };
942 state.subscribers.retain(|(connection_id, corr), subscriber| {
943 let frame = Frame::build_with_version(
944 subscriber.version,
945 FrameType::StreamData,
946 control_flags(),
947 0,
948 0,
949 *corr,
950 body.clone(),
951 );
952 match frame {
953 Ok(frame) => {
954 if subscriber.frames.try_send(frame).is_ok() {
955 true
956 } else {
957 warn!(connection_id = connection_id.get(), corr, "dropping lagged supervisor spawn subscriber");
958 false
959 }
960 }
961 Err(error) => {
962 warn!(connection_id = connection_id.get(), corr, %error, "dropping supervisor spawn subscriber after frame build failure");
963 false
964 }
965 }
966 });
967 }
968
969 fn subscribe(
970 &self,
971 connection_id: ConnectionId,
972 corr: u64,
973 version: u8,
974 since: Option<SpawnCursor>,
975 sink: FrameSink,
976 ) -> Result<(), SpawnSubscribeRefusal> {
977 let (frames, mut receiver) = mpsc::channel(SPAWN_SUBSCRIBER_BUFFER);
978 {
979 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
980 let replay = if let Some(since) = since {
981 if since.daemon_incarnation != state.daemon_incarnation {
982 return Err(SpawnSubscribeRefusal::ForeignIncarnation {
983 current: state.daemon_incarnation.clone(),
984 });
985 }
986 if let Some(oldest) = state.events.front().map(|event| event.cursor.clone()) {
987 if since.seq < oldest.seq.saturating_sub(1) {
988 return Err(SpawnSubscribeRefusal::TooOld { oldest });
989 }
990 }
991 state
992 .events
993 .iter()
994 .filter(|event| event.cursor.seq > since.seq)
995 .cloned()
996 .collect::<Vec<_>>()
997 } else {
998 Vec::new()
999 };
1000 for event in replay {
1001 let body = serde_json::to_vec(&event)
1002 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1003 let frame = Frame::build_with_version(
1004 version,
1005 FrameType::StreamData,
1006 control_flags(),
1007 0,
1008 0,
1009 corr,
1010 body,
1011 )
1012 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1013 frames
1014 .try_send(frame)
1015 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1016 }
1017 state.subscribers.insert(
1018 (connection_id, corr),
1019 SpawnSubscriber {
1020 version,
1021 frames: frames.clone(),
1022 },
1023 );
1024 }
1025 tokio::spawn(async move {
1026 while let Some(frame) = receiver.recv().await {
1027 if sink.send(frame).await.is_err() {
1028 break;
1029 }
1030 }
1031 });
1032 Ok(())
1033 }
1034
1035 fn cancel(&self, connection_id: ConnectionId, corr: u64) -> bool {
1036 let Some(subscriber) = self
1037 .0
1038 .lock()
1039 .unwrap_or_else(|p| p.into_inner())
1040 .subscribers
1041 .remove(&(connection_id, corr))
1042 else {
1043 return false;
1044 };
1045 if let Ok(frame) = Frame::build_with_version(
1046 subscriber.version,
1047 FrameType::StreamEnd,
1048 control_flags(),
1049 0,
1050 0,
1051 corr,
1052 Vec::new(),
1053 ) {
1054 tokio::spawn(async move {
1055 let _ = subscriber.frames.send(frame).await;
1056 });
1057 }
1058 true
1059 }
1060
1061 fn remove_connection(&self, connection_id: ConnectionId) {
1062 self.0
1063 .lock()
1064 .unwrap_or_else(|p| p.into_inner())
1065 .subscribers
1066 .retain(|(subscriber_connection, _), _| *subscriber_connection != connection_id);
1067 }
1068
1069 #[cfg(any(test, feature = "test-support"))]
1070 fn set_capacity(&self, capacity: usize) {
1071 self.0.lock().unwrap_or_else(|p| p.into_inner()).capacity = capacity;
1072 }
1073
1074 #[cfg(any(test, feature = "test-support"))]
1075 fn subscriber_count(&self) -> usize {
1076 self.0
1077 .lock()
1078 .unwrap_or_else(|p| p.into_inner())
1079 .subscribers
1080 .len()
1081 }
1082}
1083
1084pub trait ModuleProcessLiveness: Send + Sync {
1086 fn process_live(&self, module_id: &str) -> Option<bool>;
1087}
1088
1089#[derive(Debug, Clone, Default)]
1091pub struct SupervisorProcessLiveness {
1092 snapshots: Arc<Mutex<HashMap<String, SharedSnapshot>>>,
1093}
1094
1095impl SupervisorProcessLiveness {
1096 pub fn new() -> Self {
1097 Self::default()
1098 }
1099
1100 fn track(&self, module_id: String, snapshot: SharedSnapshot) {
1101 let mut snapshots = self
1102 .snapshots
1103 .lock()
1104 .unwrap_or_else(|poisoned| poisoned.into_inner());
1105 snapshots.insert(module_id, snapshot);
1106 }
1107
1108 fn untrack_if_current(&self, module_id: &str, snapshot: &SharedSnapshot) {
1109 let mut snapshots = self
1110 .snapshots
1111 .lock()
1112 .unwrap_or_else(|poisoned| poisoned.into_inner());
1113 let is_current = snapshots
1114 .get(module_id)
1115 .map(|tracked| Arc::ptr_eq(tracked, snapshot))
1116 .unwrap_or(false);
1117 if is_current {
1118 snapshots.remove(module_id);
1119 }
1120 }
1121}
1122
1123impl ModuleProcessLiveness for SupervisorProcessLiveness {
1124 fn process_live(&self, module_id: &str) -> Option<bool> {
1125 let snapshot = {
1126 let snapshots = self
1127 .snapshots
1128 .lock()
1129 .unwrap_or_else(|poisoned| poisoned.into_inner());
1130 snapshots.get(module_id).cloned()
1131 }?;
1132 let snapshot = snapshot
1133 .lock()
1134 .unwrap_or_else(|poisoned| poisoned.into_inner());
1135 Some(snapshot.state == ModuleState::Running && snapshot.process_alive)
1136 }
1137}
1138
1139#[derive(Debug, Clone)]
1140struct SupervisorRuntimeConfig {
1141 restart_policy: RestartPolicy,
1142 drain_timeout: Duration,
1145 effective_drain_timeout: Arc<Mutex<Duration>>,
1148 default_drain_timeout: Duration,
1151 health: HealthConfig,
1152 connection_file_path: Option<PathBuf>,
1153 capture_logs_dir: Option<PathBuf>,
1154 forwarding: Option<Arc<ForwardingTable>>,
1155 supervisor_handle: Option<SupervisorHandle>,
1158 stderr_ring: Arc<Mutex<StderrRing>>,
1165 terminal_ring: Arc<Mutex<TerminalRing>>,
1166 spawn_events: SpawnEventFeed,
1167 child_roster: ChildRoster,
1168 #[cfg(target_os = "linux")]
1169 cgroup_placement: Option<subc_cgroup::Placement>,
1170 #[cfg(test)]
1171 test_seed_stale_facts_before_enable_spawn: bool,
1172}
1173
1174#[derive(Debug, Clone, PartialEq, Eq)]
1175struct SupervisedConfiguration {
1176 spec: ModuleSpec,
1177 health: HealthConfig,
1178}
1179
1180#[derive(Debug, Clone, Default)]
1186pub struct SupervisorHandle {
1187 modules: Arc<Mutex<HashMap<String, SupervisedModule>>>,
1188 spawn_events: SpawnEventFeed,
1189 reserved_nonces: Arc<Mutex<HashMap<String, Option<String>>>>,
1200 removal_tombstones: Arc<Mutex<HashMap<String, u64>>>,
1206 spawn_nonces: Arc<Mutex<HashMap<String, String>>>,
1210 reserved_prefix_owners: Arc<Mutex<HashMap<String, String>>>,
1218 swaps: Arc<Mutex<HashMap<String, OpenSwap>>>,
1224 promotion_observer: PromotionObserverSlot,
1226 operation_lock: Arc<AsyncMutex<()>>,
1230}
1231
1232pub(crate) trait SwapPromotionObserver: Send + Sync {
1241 fn swap_promoted(&self, registration: &crate::registry::ModuleRegistration);
1242}
1243
1244#[derive(Clone, Default)]
1248struct PromotionObserverSlot(Arc<Mutex<Option<std::sync::Weak<dyn SwapPromotionObserver>>>>);
1249
1250impl fmt::Debug for PromotionObserverSlot {
1251 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1252 f.write_str("PromotionObserverSlot")
1253 }
1254}
1255
1256#[derive(Debug, Clone)]
1258struct OpenSwap {
1259 candidate_nonce: String,
1262 incumbent_nonce: Option<String>,
1267 candidate_admitted: bool,
1271}
1272
1273#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1276pub(crate) enum SwapHelloAdmission {
1277 NotSwapping,
1280 Candidate,
1282 Refused,
1285}
1286
1287#[derive(Debug, Clone, PartialEq, Eq)]
1288pub(crate) enum ReservedHelloRejection {
1289 Exact {
1290 module_id: String,
1291 },
1292 Prefix {
1293 prefix: String,
1294 owner_module_id: String,
1295 },
1296}
1297
1298impl SupervisorHandle {
1299 pub fn new() -> Self {
1300 Self::default()
1301 }
1302
1303 pub(crate) fn spawn_snapshot(&self) -> SpawnSnapshot {
1304 self.spawn_events.snapshot()
1305 }
1306
1307 pub(crate) fn subscribe_spawns(
1308 &self,
1309 connection_id: ConnectionId,
1310 corr: u64,
1311 version: u8,
1312 since: Option<SpawnCursor>,
1313 sink: FrameSink,
1314 ) -> Result<(), SpawnSubscribeRefusal> {
1315 self.spawn_events
1316 .subscribe(connection_id, corr, version, since, sink)
1317 }
1318
1319 pub(crate) fn cancel_spawn_subscription(&self, connection_id: ConnectionId, corr: u64) -> bool {
1320 self.spawn_events.cancel(connection_id, corr)
1321 }
1322
1323 pub(crate) fn remove_spawn_subscribers(&self, connection_id: ConnectionId) {
1324 self.spawn_events.remove_connection(connection_id);
1325 }
1326
1327 #[cfg(any(test, feature = "test-support"))]
1328 pub fn set_spawn_event_capacity_for_test(&self, capacity: usize) {
1329 assert!(capacity > 0, "spawn event capacity must be non-zero");
1330 self.spawn_events.set_capacity(capacity);
1331 }
1332
1333 #[cfg(any(test, feature = "test-support"))]
1334 pub fn spawn_subscriber_count_for_test(&self) -> usize {
1335 self.spawn_events.subscriber_count()
1336 }
1337
1338 pub fn set_spawn_nonce(&self, module_id: &str, nonce: String) {
1341 self.spawn_nonces
1342 .lock()
1343 .unwrap_or_else(|poisoned| poisoned.into_inner())
1344 .insert(module_id.to_string(), nonce);
1345 }
1346
1347 pub fn set_reserved_nonce(&self, module_id: &str, nonce: String) {
1350 self.reserved_nonces
1351 .lock()
1352 .unwrap_or_else(|poisoned| poisoned.into_inner())
1353 .insert(module_id.to_string(), Some(nonce));
1354 }
1355
1356 pub fn set_reserved_prefixes(&self, owner_module_id: &str, prefixes: &[String]) {
1358 let mut owners = self
1359 .reserved_prefix_owners
1360 .lock()
1361 .unwrap_or_else(|poisoned| poisoned.into_inner());
1362 owners.retain(|_, owner| owner != owner_module_id);
1363 for prefix in prefixes {
1364 owners.insert(prefix.clone(), owner_module_id.to_string());
1365 }
1366 }
1367
1368 #[cfg(test)]
1370 pub(crate) fn spawn_nonce(&self, module_id: &str) -> Option<String> {
1371 self.spawn_nonces
1372 .lock()
1373 .unwrap_or_else(|poisoned| poisoned.into_inner())
1374 .get(module_id)
1375 .cloned()
1376 }
1377
1378 fn apply_identity_configuration(&self, spec: &ModuleSpec) {
1379 self.set_reserved_prefixes(&spec.module_id, &spec.reserved_prefixes);
1380 let spawn_nonce = self
1381 .spawn_nonces
1382 .lock()
1383 .unwrap_or_else(|poisoned| poisoned.into_inner())
1384 .get(&spec.module_id)
1385 .cloned();
1386 let mut reserved_nonces = self
1387 .reserved_nonces
1388 .lock()
1389 .unwrap_or_else(|poisoned| poisoned.into_inner());
1390 if spec.reserved {
1391 reserved_nonces.insert(spec.module_id.clone(), spawn_nonce);
1396 }
1397 drop(reserved_nonces);
1398 self.removal_tombstones
1402 .lock()
1403 .unwrap_or_else(|poisoned| poisoned.into_inner())
1404 .remove(&spec.module_id);
1405 }
1406
1407 pub fn reserved_hello_authorized(&self, module_id: &str, presented: Option<&str>) -> bool {
1412 self.reserved_hello_rejection(module_id, presented)
1413 .is_none()
1414 }
1415
1416 pub(crate) fn reserved_hello_rejection(
1417 &self,
1418 module_id: &str,
1419 presented: Option<&str>,
1420 ) -> Option<ReservedHelloRejection> {
1421 let nonces = self
1422 .reserved_nonces
1423 .lock()
1424 .unwrap_or_else(|poisoned| poisoned.into_inner());
1425 if let Some(expected) = nonces.get(module_id) {
1426 let authorized = match expected {
1430 Some(expected) => {
1431 presented.is_some_and(|p| constant_time_eq(expected.as_bytes(), p.as_bytes()))
1432 }
1433 None => false,
1434 };
1435 if authorized {
1436 return None;
1437 }
1438 return Some(ReservedHelloRejection::Exact {
1439 module_id: module_id.to_string(),
1440 });
1441 }
1442 drop(nonces);
1443
1444 let matched_prefix = self
1445 .reserved_prefix_owners
1446 .lock()
1447 .unwrap_or_else(|poisoned| poisoned.into_inner())
1448 .iter()
1449 .filter(|(prefix, _)| module_id.starts_with(prefix.as_str()))
1450 .max_by_key(|(prefix, _)| prefix.len())
1451 .map(|(prefix, owner)| (prefix.clone(), owner.clone()));
1452 let (prefix, owner_module_id) = matched_prefix?;
1453
1454 let authorized = presented.is_some_and(|presented| {
1455 self.spawn_nonces
1456 .lock()
1457 .unwrap_or_else(|poisoned| poisoned.into_inner())
1458 .get(&owner_module_id)
1459 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()))
1460 || self.swap_nonce_matches(&owner_module_id, presented)
1463 });
1464 if authorized {
1465 None
1466 } else {
1467 Some(ReservedHelloRejection::Prefix {
1468 prefix,
1469 owner_module_id,
1470 })
1471 }
1472 }
1473
1474 pub fn spawned_consumer_authorized(&self, module_id: &str, presented: &str) -> bool {
1479 if presented.is_empty() {
1480 return false;
1481 }
1482 let nonces = self
1483 .spawn_nonces
1484 .lock()
1485 .unwrap_or_else(|poisoned| poisoned.into_inner());
1486 let current = nonces
1487 .get(module_id)
1488 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()));
1489 drop(nonces);
1490 current || self.swap_nonce_matches(module_id, presented)
1495 }
1496
1497 fn swap_nonce_matches(&self, module_id: &str, presented: &str) -> bool {
1499 let swaps = self
1500 .swaps
1501 .lock()
1502 .unwrap_or_else(|poisoned| poisoned.into_inner());
1503 swaps.get(module_id).is_some_and(|swap| {
1504 constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes())
1505 || swap.incumbent_nonce.as_deref().is_some_and(|incumbent| {
1506 constant_time_eq(incumbent.as_bytes(), presented.as_bytes())
1507 })
1508 })
1509 }
1510
1511 pub(crate) fn open_swap(&self, module_id: &str, candidate_nonce: String) {
1514 let incumbent_nonce = self
1515 .spawn_nonces
1516 .lock()
1517 .unwrap_or_else(|poisoned| poisoned.into_inner())
1518 .get(module_id)
1519 .cloned();
1520 self.swaps
1521 .lock()
1522 .unwrap_or_else(|poisoned| poisoned.into_inner())
1523 .insert(
1524 module_id.to_string(),
1525 OpenSwap {
1526 candidate_nonce,
1527 incumbent_nonce,
1528 candidate_admitted: false,
1529 },
1530 );
1531 }
1532
1533 pub(crate) fn close_swap(&self, module_id: &str) {
1536 self.swaps
1537 .lock()
1538 .unwrap_or_else(|poisoned| poisoned.into_inner())
1539 .remove(module_id);
1540 }
1541
1542 pub(crate) fn set_swap_promotion_observer(
1545 &self,
1546 observer: std::sync::Weak<dyn SwapPromotionObserver>,
1547 ) {
1548 *self
1549 .promotion_observer
1550 .0
1551 .lock()
1552 .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(observer);
1553 }
1554
1555 fn notify_swap_promoted(&self, registration: &crate::registry::ModuleRegistration) {
1558 let observer = self
1559 .promotion_observer
1560 .0
1561 .lock()
1562 .unwrap_or_else(|poisoned| poisoned.into_inner())
1563 .as_ref()
1564 .and_then(std::sync::Weak::upgrade);
1565 if let Some(observer) = observer {
1566 observer.swap_promoted(registration);
1567 }
1568 }
1569
1570 pub(crate) fn swap_open(&self, module_id: &str) -> bool {
1572 self.swaps
1573 .lock()
1574 .unwrap_or_else(|poisoned| poisoned.into_inner())
1575 .contains_key(module_id)
1576 }
1577
1578 fn promote_swap_nonce(&self, module_id: &str, reserved: bool) {
1583 let candidate_nonce = self
1584 .swaps
1585 .lock()
1586 .unwrap_or_else(|poisoned| poisoned.into_inner())
1587 .get(module_id)
1588 .map(|swap| swap.candidate_nonce.clone());
1589 let Some(nonce) = candidate_nonce else {
1590 return;
1591 };
1592 self.set_spawn_nonce(module_id, nonce.clone());
1593 if reserved {
1594 self.set_reserved_nonce(module_id, nonce);
1595 }
1596 }
1597
1598 pub(crate) fn swap_hello_admission(
1613 &self,
1614 module_id: &str,
1615 presented: Option<&str>,
1616 ) -> SwapHelloAdmission {
1617 let swaps = self
1618 .swaps
1619 .lock()
1620 .unwrap_or_else(|poisoned| poisoned.into_inner());
1621 let Some(swap) = swaps.get(module_id) else {
1622 return SwapHelloAdmission::NotSwapping;
1623 };
1624 let Some(presented) = presented else {
1625 return SwapHelloAdmission::Refused;
1626 };
1627 if constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes()) {
1628 return if swap.candidate_admitted {
1629 SwapHelloAdmission::Refused
1630 } else {
1631 SwapHelloAdmission::Candidate
1632 };
1633 }
1634 if swap
1635 .incumbent_nonce
1636 .as_deref()
1637 .is_some_and(|incumbent| constant_time_eq(incumbent.as_bytes(), presented.as_bytes()))
1638 {
1639 return SwapHelloAdmission::NotSwapping;
1640 }
1641 SwapHelloAdmission::Refused
1642 }
1643
1644 pub(crate) fn mark_swap_candidate_admitted(&self, module_id: &str) {
1647 if let Some(swap) = self
1648 .swaps
1649 .lock()
1650 .unwrap_or_else(|poisoned| poisoned.into_inner())
1651 .get_mut(module_id)
1652 {
1653 swap.candidate_admitted = true;
1654 }
1655 }
1656
1657 pub fn spawn_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1659 self.spawn_nonces
1660 .lock()
1661 .unwrap_or_else(|poisoned| poisoned.into_inner())
1662 .get(module_id)
1663 .cloned()
1664 }
1665
1666 pub fn reserved_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1668 self.reserved_nonces
1669 .lock()
1670 .unwrap_or_else(|poisoned| poisoned.into_inner())
1671 .get(module_id)
1672 .cloned()
1673 .flatten()
1674 }
1675
1676 pub fn insert(&self, module: SupervisedModule) -> Option<SupervisedModule> {
1677 let mut modules = self
1678 .modules
1679 .lock()
1680 .unwrap_or_else(|poisoned| poisoned.into_inner());
1681 modules.insert(module.module_id().to_string(), module)
1682 }
1683
1684 pub fn get(&self, module_id: &str) -> Option<SupervisedModule> {
1685 let modules = self
1686 .modules
1687 .lock()
1688 .unwrap_or_else(|poisoned| poisoned.into_inner());
1689 modules.get(module_id).cloned()
1690 }
1691
1692 pub(crate) fn record_late_health_answer(
1693 &self,
1694 module_id: &str,
1695 latency_ms: u64,
1696 ) -> Result<bool, SuperviseError> {
1697 let Some(module) = self.get(module_id) else {
1698 return Ok(false);
1699 };
1700 update_snapshot(&module.inner.snapshot, Some(module_id), |state| {
1701 state.health.late_answer_count = state.health.late_answer_count.saturating_add(1);
1702 state.health.last_late_answer_latency_ms = Some(latency_ms);
1703 state.health.consecutive_failures = 0;
1711 })?;
1712 Ok(true)
1713 }
1714
1715 pub fn record_deliberate_severance(&self, module_id: &str) -> Result<bool, SuperviseError> {
1721 let Some(module) = self.get(module_id) else {
1722 return Ok(false);
1723 };
1724 let status = module.status()?;
1725 let Some((pid, start_time)) = status.pid.zip(status.process_start_time) else {
1726 return Ok(false);
1727 };
1728 module.record_deliberate_severance(ProcessIdentity { pid, start_time })
1729 }
1730
1731 pub fn list(&self) -> Vec<SupervisedModule> {
1732 let modules = self
1733 .modules
1734 .lock()
1735 .unwrap_or_else(|poisoned| poisoned.into_inner());
1736 let mut modules = modules.values().cloned().collect::<Vec<_>>();
1737 modules.sort_by(|left, right| left.module_id().cmp(right.module_id()));
1738 modules
1739 }
1740
1741 pub(crate) fn retire(&self, module_id: &str) -> Option<SupervisedModule> {
1742 self.spawn_nonces
1743 .lock()
1744 .unwrap_or_else(|poisoned| poisoned.into_inner())
1745 .remove(module_id);
1746 self.close_swap(module_id);
1747 let mut reserved_nonces = self
1748 .reserved_nonces
1749 .lock()
1750 .unwrap_or_else(|poisoned| poisoned.into_inner());
1751 if reserved_nonces.contains_key(module_id) {
1752 reserved_nonces.insert(module_id.to_string(), None);
1755 }
1756 drop(reserved_nonces);
1757 self.reserved_prefix_owners
1758 .lock()
1759 .unwrap_or_else(|poisoned| poisoned.into_inner())
1760 .retain(|_, owner| owner != module_id);
1761 self.modules
1762 .lock()
1763 .unwrap_or_else(|poisoned| poisoned.into_inner())
1764 .remove(module_id)
1765 }
1766
1767 pub(crate) fn record_rescan_removal(&self, module_id: &str) {
1770 self.removal_tombstones
1771 .lock()
1772 .unwrap_or_else(|poisoned| poisoned.into_inner())
1773 .insert(module_id.to_string(), unix_ms_now());
1774 }
1775
1776 pub(crate) fn removal_tombstone_age_ms(&self, module_id: &str) -> Option<u64> {
1778 self.removal_tombstones
1779 .lock()
1780 .unwrap_or_else(|poisoned| poisoned.into_inner())
1781 .get(module_id)
1782 .copied()
1783 .map(|removed_at_ms| unix_ms_now().saturating_sub(removed_at_ms))
1784 }
1785
1786 pub(crate) fn release_retained_reserved_gate(&self, module_id: &str) -> bool {
1791 if self.get(module_id).is_some() {
1792 return false;
1793 }
1794 let mut reserved_nonces = self
1795 .reserved_nonces
1796 .lock()
1797 .unwrap_or_else(|poisoned| poisoned.into_inner());
1798 if !matches!(reserved_nonces.get(module_id), Some(None)) {
1799 return false;
1800 }
1801 reserved_nonces.remove(module_id);
1802 true
1803 }
1804
1805 pub(crate) fn operation_lock(&self) -> Arc<AsyncMutex<()>> {
1806 Arc::clone(&self.operation_lock)
1807 }
1808}
1809
1810#[derive(Debug, Clone)]
1812pub struct Supervisor {
1813 registry: Arc<Registry>,
1814 restart_policy: RestartPolicy,
1815 drain_timeout: Duration,
1816 connection_file_path: Option<PathBuf>,
1817 capture_logs_dir: Option<PathBuf>,
1818 forwarding: Option<Arc<ForwardingTable>>,
1819 process_liveness: Arc<SupervisorProcessLiveness>,
1820 supervisor_handle: Option<SupervisorHandle>,
1821 health: HealthConfig,
1822 daemon_start_clock: crate::clock::StartClock,
1823 terminal_journal: Option<Arc<crate::terminal_journal::TerminalJournal>>,
1824 spawn_events: SpawnEventFeed,
1825 provenance_probe: ExecutableIdentityProbe,
1826 child_roster: ChildRoster,
1829 #[cfg(target_os = "linux")]
1830 cgroup_placement: Option<subc_cgroup::Placement>,
1831}
1832
1833impl Supervisor {
1834 #[cfg(unix)]
1835 pub(crate) fn stamp_shutdown(&self) {
1836 if let Some(journal) = &self.terminal_journal {
1837 journal.stamp_shutdown();
1838 }
1839 }
1840
1841 #[cfg(unix)]
1845 pub(crate) async fn drain_for_daemon_shutdown(&self) -> Result<(), SuperviseError> {
1846 const NOTICE_BUDGET: Duration = Duration::from_millis(500);
1847 const DRAIN_BUDGET: Duration = Duration::from_secs(2);
1848 let Some(forwarding) = &self.forwarding else {
1849 return Ok(());
1850 };
1851 let module_ids = forwarding
1852 .begin_daemon_drain()
1853 .map_err(SuperviseError::Forwarding)?;
1854 let deadline_ms =
1855 unix_ms_now().saturating_add((NOTICE_BUDGET + DRAIN_BUDGET).as_millis() as u64);
1856 let mut notices = tokio::task::JoinSet::new();
1857 let mut drains = Vec::new();
1858 for module_id in module_ids {
1859 let Some(target) = forwarding
1860 .begin_module_drain(&module_id, RouteCloseReason::Restart)
1861 .map_err(SuperviseError::Forwarding)?
1862 else {
1863 continue;
1864 };
1865 let routes = forwarding
1866 .endpoint_routes(target.endpoint)
1867 .map_err(SuperviseError::Forwarding)?;
1868 let command = serde_json::to_vec(&ModuleControlCommand::Draining {
1872 reason: RouteCloseReason::Restart,
1873 deadline_ms,
1874 })
1875 .expect("module draining serializes");
1876 let closing = serde_json::to_vec(&ClientControlPush::RouteClosing {
1877 module_id: module_id.clone(),
1878 reason: RouteCloseReason::Restart,
1879 })
1880 .expect("route closing serializes");
1881 let mut recipients = vec![(target.sink.clone(), target.negotiated_ver, command)];
1882 let mut seen = std::collections::HashSet::new();
1883 for route in routes {
1884 let client = route.goodbye_target;
1885 if seen.insert(client.connection_id) {
1886 recipients.push((client.sink, client.negotiated_ver, closing.clone()));
1887 }
1888 }
1889 for (sink, version, body) in recipients {
1890 notices.spawn(async move {
1891 let frame = Frame::build_with_version(
1892 version,
1893 FrameType::Push,
1894 control_flags(),
1895 0,
1896 0,
1897 0,
1898 body,
1899 )
1900 .expect("bounded lifecycle notice frame builds");
1901 sink.send_flushed(frame).await
1902 });
1903 }
1904 let gauges = declared_busy_gauges(&self.registry, &module_id)?;
1905 drains.push((module_id, target.endpoint, gauges));
1906 }
1907 let notice_deadline = Instant::now() + NOTICE_BUDGET;
1910 while let Ok(Some(result)) = timeout_at(notice_deadline, notices.join_next()).await {
1911 if !matches!(result, Ok(Ok(()))) {
1912 warn!(?result, "daemon shutdown notice delivery failed");
1913 }
1914 }
1915 notices.abort_all();
1916 let deadline = Instant::now() + DRAIN_BUDGET;
1917 let mut waits = tokio::task::JoinSet::new();
1918 for (module_id, endpoint, gauges) in drains {
1919 let forwarding = Arc::clone(forwarding);
1920 let mut runtime = self.runtime_config();
1921 runtime.health.cadence = Duration::from_millis(100);
1922 waits.spawn(async move {
1923 wait_for_forwarding_quiescence(
1924 &forwarding,
1925 &module_id,
1926 &runtime,
1927 endpoint,
1928 deadline,
1929 &gauges,
1930 DrainScope::Active,
1931 )
1932 .await
1933 });
1934 }
1935 while let Ok(Some(result)) = timeout_at(deadline, waits.join_next()).await {
1936 if !matches!(result, Ok(Ok(true))) {
1937 warn!(?result, "daemon shutdown drain did not reach quiescence");
1938 }
1939 }
1940 Ok(())
1941 }
1942
1943 #[cfg(unix)]
1953 pub(crate) async fn end_children_for_daemon_shutdown(
1954 &self,
1955 already_escalated: bool,
1956 escalate: impl std::future::Future<Output = ()>,
1957 ) {
1958 if let Some(forwarding) = &self.forwarding {
1959 let closed = forwarding.close_all_connections(&CloseReason::new(
1960 "daemon_shutdown",
1961 "the daemon is exiting after its shutdown notice and drain",
1962 ));
1963 debug!(closed, "closed established connections for daemon shutdown");
1964 }
1965 crate::child_roster::end_children_for_daemon_shutdown(
1966 &self.child_roster,
1967 already_escalated,
1968 escalate,
1969 )
1970 .await;
1971 }
1972
1973 pub fn new(registry: Arc<Registry>, restart_policy: RestartPolicy) -> Self {
1974 Self {
1975 registry,
1976 restart_policy,
1977 drain_timeout: DEFAULT_DRAIN_TIMEOUT,
1978 connection_file_path: None,
1979 capture_logs_dir: None,
1980 forwarding: None,
1981 process_liveness: Arc::new(SupervisorProcessLiveness::default()),
1982 supervisor_handle: None,
1983 health: HealthConfig::default(),
1984 daemon_start_clock: crate::clock::StartClock::capture(),
1985 terminal_journal: None,
1986 spawn_events: SpawnEventFeed::default(),
1987 provenance_probe: ExecutableIdentityProbe::default(),
1988 child_roster: ChildRoster::default(),
1989 #[cfg(target_os = "linux")]
1990 cgroup_placement: None,
1991 }
1992 }
1993
1994 pub fn with_drain_timeout(mut self, drain_timeout: Duration) -> Self {
1995 self.drain_timeout = drain_timeout;
1996 self
1997 }
1998
1999 pub fn with_process_liveness(
2000 mut self,
2001 process_liveness: Arc<SupervisorProcessLiveness>,
2002 ) -> Self {
2003 self.process_liveness = process_liveness;
2004 self
2005 }
2006
2007 pub fn with_connection_file_path(mut self, connection_file_path: impl Into<PathBuf>) -> Self {
2008 self.connection_file_path = Some(connection_file_path.into());
2009 self
2010 }
2011
2012 pub fn with_capture_logs_dir(mut self, logs_dir: impl Into<PathBuf>) -> Self {
2014 self.capture_logs_dir = Some(logs_dir.into());
2015 self
2016 }
2017
2018 pub fn with_terminal_journal(mut self, path: PathBuf, daemon_incarnation: String) -> Self {
2020 self.spawn_events
2024 .configure_incarnation(daemon_incarnation.clone());
2025 self.terminal_journal = Some(Arc::new(crate::terminal_journal::TerminalJournal::open(
2026 path,
2027 daemon_incarnation,
2028 )));
2029 self
2030 }
2031
2032 pub fn with_forwarding(mut self, forwarding: Arc<ForwardingTable>) -> Self {
2033 self.forwarding = Some(forwarding);
2034 self
2035 }
2036
2037 pub fn with_handle(mut self, supervisor_handle: SupervisorHandle) -> Self {
2038 self.spawn_events = supervisor_handle.spawn_events.clone();
2039 self.supervisor_handle = Some(supervisor_handle);
2040 self
2041 }
2042
2043 pub fn with_health_config(mut self, health: HealthConfig) -> Self {
2044 self.health = health;
2045 self
2046 }
2047
2048 #[cfg(target_os = "linux")]
2049 pub fn with_cgroup_placement(
2050 mut self,
2051 cgroup_placement: Option<subc_cgroup::Placement>,
2052 ) -> Self {
2053 self.cgroup_placement = cgroup_placement;
2054 self
2055 }
2056
2057 pub fn spawn(&self, spec: ModuleSpec) -> Result<SupervisedModule, SuperviseError> {
2063 validate_spec(&spec)?;
2064
2065 let runtime = self.runtime_config();
2066 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2067 let child = spawn_child(
2068 &spec,
2069 runtime.connection_file_path.as_deref(),
2070 self.supervisor_handle.as_ref(),
2071 &runtime.stderr_ring,
2072 runtime.capture_logs_dir.as_deref(),
2073 &runtime.child_roster,
2074 #[cfg(target_os = "linux")]
2075 runtime.cgroup_placement.as_ref(),
2076 )?;
2077 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2078 self.process_liveness
2079 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2080
2081 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2082 }
2083
2084 pub fn supervise_configured(
2090 &self,
2091 spec: ModuleSpec,
2092 enabled: bool,
2093 ) -> Result<SupervisedModule, SuperviseError> {
2094 validate_spec(&spec)?;
2095
2096 let runtime = self.runtime_config();
2097 if !enabled {
2098 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2099 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2100 }
2101
2102 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2103 match spawn_child(
2104 &spec,
2105 runtime.connection_file_path.as_deref(),
2106 self.supervisor_handle.as_ref(),
2107 &runtime.stderr_ring,
2108 runtime.capture_logs_dir.as_deref(),
2109 &runtime.child_roster,
2110 #[cfg(target_os = "linux")]
2111 runtime.cgroup_placement.as_ref(),
2112 ) {
2113 Ok(child) => {
2114 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2115 self.process_liveness
2116 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2117 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2118 }
2119 Err(err) => {
2120 error!(
2121 module_id = %spec.module_id,
2122 program = %spec.program.display(),
2123 error = %err,
2124 "configured module failed to spawn; marking failed and continuing"
2125 );
2126 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2127 Ok(self.supervised_module(spec, runtime, snapshot, None))
2128 }
2129 }
2130 }
2131
2132 pub fn supervise_configured_with_health(
2138 &self,
2139 spec: ModuleSpec,
2140 enabled: bool,
2141 health: HealthConfig,
2142 drain_timeout_ms: Option<u64>,
2143 restart_policy: RestartPolicy,
2144 ) -> Result<SupervisedModule, SuperviseError> {
2145 validate_spec(&spec)?;
2146
2147 let mut runtime = self.runtime_config();
2148 runtime.health = health;
2149 runtime.restart_policy = restart_policy;
2150 if let Some(ms) = drain_timeout_ms {
2151 runtime.drain_timeout = Duration::from_millis(ms);
2152 *runtime
2153 .effective_drain_timeout
2154 .lock()
2155 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
2156 }
2157 if !enabled {
2158 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2159 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2160 }
2161
2162 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2163 match spawn_child(
2164 &spec,
2165 runtime.connection_file_path.as_deref(),
2166 self.supervisor_handle.as_ref(),
2167 &runtime.stderr_ring,
2168 runtime.capture_logs_dir.as_deref(),
2169 &runtime.child_roster,
2170 #[cfg(target_os = "linux")]
2171 runtime.cgroup_placement.as_ref(),
2172 ) {
2173 Ok(child) => {
2174 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2175 self.process_liveness
2176 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2177 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2178 }
2179 Err(err) => {
2180 if health.critical {
2181 error!(
2182 module_id = %spec.module_id,
2183 program = %spec.program.display(),
2184 error = %err,
2185 "critical configured module failed to spawn; marking failed and alerting"
2186 );
2187 } else {
2188 error!(
2189 module_id = %spec.module_id,
2190 program = %spec.program.display(),
2191 error = %err,
2192 "configured module failed to spawn; marking failed and continuing"
2193 );
2194 }
2195 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2196 Ok(self.supervised_module(spec, runtime, snapshot, None))
2197 }
2198 }
2199 }
2200
2201 fn runtime_config(&self) -> SupervisorRuntimeConfig {
2202 let effective_drain_timeout = Arc::new(Mutex::new(self.drain_timeout));
2203 SupervisorRuntimeConfig {
2204 restart_policy: self.restart_policy,
2205 drain_timeout: self.drain_timeout,
2206 child_roster: self
2209 .child_roster
2210 .for_module(Arc::clone(&effective_drain_timeout)),
2211 effective_drain_timeout,
2212 default_drain_timeout: self.drain_timeout,
2213 health: self.health,
2214 connection_file_path: self.connection_file_path.clone(),
2215 capture_logs_dir: self.capture_logs_dir.clone(),
2216 forwarding: self.forwarding.clone(),
2217 supervisor_handle: self.supervisor_handle.clone(),
2218 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
2219 terminal_ring: Arc::new(Mutex::new(
2220 TerminalRing::new(
2221 TerminalRingConfig::default(),
2222 self.daemon_start_clock.started_at_ms(),
2223 )
2224 .with_start_clock(self.daemon_start_clock)
2225 .with_journal(self.terminal_journal.clone()),
2226 )),
2227 spawn_events: self.spawn_events.clone(),
2228 #[cfg(target_os = "linux")]
2229 cgroup_placement: self.cgroup_placement.clone(),
2230 #[cfg(test)]
2231 test_seed_stale_facts_before_enable_spawn: false,
2232 }
2233 }
2234
2235 fn supervised_module(
2236 &self,
2237 spec: ModuleSpec,
2238 runtime: SupervisorRuntimeConfig,
2239 snapshot: SharedSnapshot,
2240 child: Option<SupervisedChild>,
2241 ) -> SupervisedModule {
2242 let configuration = Arc::new(Mutex::new(SupervisedConfiguration {
2243 spec: spec.clone(),
2244 health: runtime.health,
2245 }));
2246 let stderr_ring = Arc::clone(&runtime.stderr_ring);
2247 let terminal_ring = Arc::clone(&runtime.terminal_ring);
2248 let restart_policy = runtime.restart_policy;
2252 let effective_drain_timeout = Arc::clone(&runtime.effective_drain_timeout);
2253 let (tx, rx) = mpsc::channel(4);
2254 let monitor = tokio::spawn(supervise_loop(
2255 spec.clone(),
2256 runtime,
2257 Arc::clone(&self.registry),
2258 Arc::clone(&self.process_liveness),
2259 Arc::clone(&snapshot),
2260 child,
2261 rx,
2262 ));
2263
2264 let module_id = spec.module_id.clone();
2265 let module = SupervisedModule {
2266 inner: Arc::new(SupervisedModuleInner {
2267 module_id: module_id.clone(),
2268 registry: Arc::clone(&self.registry),
2269 snapshot,
2270 configuration,
2271 stderr_ring,
2272 terminal_ring,
2273 commands: tx,
2274 monitor: Mutex::new(Some(monitor)),
2275 restart_policy,
2276 effective_drain_timeout,
2277 provenance_probe: self.provenance_probe.clone(),
2278 }),
2279 };
2280 if let Some(supervisor_handle) = &self.supervisor_handle {
2281 supervisor_handle.apply_identity_configuration(&spec);
2282 supervisor_handle.insert(module.clone());
2283 }
2284 module
2285 }
2286}
2287
2288impl Default for Supervisor {
2289 fn default() -> Self {
2290 Self::new(Arc::new(Registry::default()), RestartPolicy::default())
2291 }
2292}
2293
2294#[derive(Clone)]
2296pub struct SupervisedModule {
2297 inner: Arc<SupervisedModuleInner>,
2298}
2299
2300struct SupervisedModuleInner {
2301 module_id: String,
2302 registry: Arc<Registry>,
2303 snapshot: SharedSnapshot,
2304 configuration: Arc<Mutex<SupervisedConfiguration>>,
2305 stderr_ring: Arc<Mutex<StderrRing>>,
2306 terminal_ring: Arc<Mutex<TerminalRing>>,
2307 commands: mpsc::Sender<SupervisorCommand>,
2308 monitor: Mutex<Option<JoinHandle<()>>>,
2309 restart_policy: RestartPolicy,
2313 effective_drain_timeout: Arc<Mutex<Duration>>,
2314 provenance_probe: ExecutableIdentityProbe,
2315}
2316
2317impl fmt::Debug for SupervisedModule {
2318 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2319 f.debug_struct("SupervisedModule")
2320 .field("module_id", &self.inner.module_id)
2321 .field("status", &self.status())
2322 .finish_non_exhaustive()
2323 }
2324}
2325
2326impl SupervisedModule {
2327 pub fn module_id(&self) -> &str {
2328 &self.inner.module_id
2329 }
2330
2331 #[cfg(test)]
2335 pub(crate) fn record_health_probe_failure_for_test(
2336 &self,
2337 detail: &str,
2338 ) -> Result<(), SuperviseError> {
2339 update_snapshot(&self.inner.snapshot, Some(&self.inner.module_id), |state| {
2340 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
2341 state.health.detail = Some(detail.to_string());
2342 })
2343 }
2344
2345 pub fn state(&self) -> Result<ModuleState, SuperviseError> {
2346 Ok(lock_snapshot(&self.inner.snapshot)?.state)
2347 }
2348
2349 pub fn stderr_tail(
2356 &self,
2357 max_lines: Option<usize>,
2358 max_bytes: Option<usize>,
2359 ) -> StderrTailSnapshot {
2360 self.inner
2361 .stderr_ring
2362 .lock()
2363 .unwrap_or_else(|poisoned| poisoned.into_inner())
2364 .snapshot(max_lines, max_bytes)
2365 }
2366
2367 pub fn terminal_history(&self) -> TerminalHistorySnapshot {
2372 self.inner
2373 .terminal_ring
2374 .lock()
2375 .unwrap_or_else(|poisoned| poisoned.into_inner())
2376 .snapshot()
2377 }
2378
2379 pub fn durable_terminal_history(&self) -> subc_control::TerminalHistory {
2381 self.inner
2382 .terminal_ring
2383 .lock()
2384 .unwrap_or_else(|p| p.into_inner())
2385 .durable_history(&self.inner.module_id)
2386 }
2387
2388 pub fn status(&self) -> Result<ModuleStatus, SuperviseError> {
2389 self.status_with_snapshot_lock(&self.inner.snapshot, None)
2390 }
2391
2392 pub(crate) fn record_deliberate_severance(
2393 &self,
2394 identity: ProcessIdentity,
2395 ) -> Result<bool, SuperviseError> {
2396 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2397 if snapshot.pid != Some(identity.pid)
2398 || snapshot.process_start_time != Some(identity.start_time)
2399 {
2400 return Ok(false);
2401 }
2402 snapshot.deliberate_severance = Some(identity);
2403 Ok(true)
2404 }
2405
2406 pub(crate) fn status_for_control(
2411 &self,
2412 caller: &'static str,
2413 ) -> Result<ModuleStatus, SuperviseError> {
2414 self.status_with_snapshot_lock(&self.inner.snapshot, Some(caller))
2415 }
2416
2417 fn status_with_snapshot_lock(
2418 &self,
2419 snapshot: &SharedSnapshot,
2420 caller: Option<&'static str>,
2421 ) -> Result<ModuleStatus, SuperviseError> {
2422 let mut guard = match caller {
2423 Some(caller) => lock_snapshot_for_control(snapshot, &self.inner.module_id, caller)?,
2424 None => lock_snapshot(snapshot)?,
2425 };
2426 let restart_count =
2429 guard.crash_restarts_in_window(self.inner.restart_policy.window, Instant::now());
2430 let snapshot = guard.clone();
2431 drop(guard);
2432 let drain_timeout = *self.inner.effective_drain_timeout.lock().map_err(|_| {
2433 SuperviseError::StatePoisoned {
2434 module_id: Some(self.inner.module_id.clone()),
2435 }
2436 })?;
2437 let registration_active = self
2438 .inner
2439 .registry
2440 .get_module(&self.inner.module_id)
2441 .map_err(SuperviseError::Registry)?
2442 .is_some();
2443 let protocol = self.declared_protocol()?;
2444 let running_process =
2445 snapshot.enabled && snapshot.state == ModuleState::Running && snapshot.process_alive;
2446 let live = match protocol {
2452 ModuleProtocol::Subc => running_process && registration_active,
2453 ModuleProtocol::None => running_process,
2454 };
2455
2456 Ok(ModuleStatus {
2457 module_id: self.inner.module_id.clone(),
2458 state: snapshot.state,
2459 enabled: snapshot.enabled,
2460 process_alive: snapshot.process_alive,
2461 registration_active,
2462 protocol,
2463 live,
2464 restart_count,
2465 lifetime_restarts: snapshot.lifetime_restarts,
2466 spawn_generation: snapshot.spawn_generation,
2467 max_restarts: self.inner.restart_policy.max_restarts,
2468 restart_window: self.inner.restart_policy.window,
2469 drain_timeout,
2470 restart_backoff: self.inner.restart_policy.backoff,
2471 restart_max_backoff: self.inner.restart_policy.max_backoff,
2472 pid: snapshot.pid,
2473 spawned_at_ms: snapshot.spawned_at_ms,
2474 spawned_from: snapshot.spawned_from,
2475 process_start_time: snapshot.process_start_time,
2476 last_exit: snapshot.last_exit,
2477 health: snapshot.health,
2478 })
2479 }
2480
2481 #[cfg(test)]
2482 pub(crate) fn hold_snapshot_for_test(
2483 &self,
2484 acquired: std::sync::mpsc::Sender<()>,
2485 hold: Duration,
2486 ) -> std::thread::JoinHandle<()> {
2487 let snapshot = Arc::clone(&self.inner.snapshot);
2488 std::thread::spawn(move || {
2489 let _guard = snapshot.lock().expect("test snapshot lock is not poisoned");
2490 acquired
2491 .send(())
2492 .expect("test receiver waits for snapshot lock");
2493 std::thread::sleep(hold);
2494 })
2495 }
2496
2497 pub(crate) async fn running_image_agreement(&self) -> subc_control::RunningImageAgreement {
2498 let snapshot = match lock_snapshot(&self.inner.snapshot) {
2499 Ok(snapshot) => snapshot.clone(),
2500 Err(_) => {
2501 return subc_control::RunningImageAgreement::Unavailable {
2502 reason: subc_control::RunningImageUnavailableReason::NotRunning,
2503 };
2504 }
2505 };
2506 self.inner
2507 .provenance_probe
2508 .observe(
2509 snapshot.pid,
2510 snapshot.spawned_from.as_deref(),
2511 snapshot.spawned_file_identity,
2512 snapshot.process_start_time,
2513 )
2514 .await
2515 }
2516
2517 pub(crate) fn will_recover_after_connection_loss(&self) -> Result<bool, SuperviseError> {
2518 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2519 Ok(match snapshot.state {
2520 ModuleState::Restarting => true,
2521 ModuleState::Failed | ModuleState::Disabled => false,
2522 _ => daemon_will_restart(&mut snapshot, &self.inner.restart_policy, Instant::now()),
2523 })
2524 }
2525
2526 #[cfg(test)]
2527 pub(crate) fn is_warming(&self) -> Result<bool, SuperviseError> {
2528 self.is_warming_with_snapshot_lock(None)
2529 }
2530
2531 pub(crate) fn is_warming_for_control(
2532 &self,
2533 caller: &'static str,
2534 ) -> Result<bool, SuperviseError> {
2535 self.is_warming_with_snapshot_lock(Some(caller))
2536 }
2537
2538 fn is_warming_with_snapshot_lock(
2539 &self,
2540 caller: Option<&'static str>,
2541 ) -> Result<bool, SuperviseError> {
2542 let snapshot = match caller {
2543 Some(caller) => {
2544 lock_snapshot_for_control(&self.inner.snapshot, &self.inner.module_id, caller)?
2545 }
2546 None => lock_snapshot(&self.inner.snapshot)?,
2547 }
2548 .clone();
2549 Ok(matches!(
2550 snapshot.state,
2551 ModuleState::Starting | ModuleState::Running | ModuleState::Restarting
2552 ))
2553 }
2554
2555 pub async fn drain(&self) -> Result<(), SuperviseError> {
2557 self.stop().await
2558 }
2559
2560 pub(crate) async fn retire(&self) -> Result<(), SuperviseError> {
2561 match self.state()? {
2562 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2563 ModuleState::Starting
2564 | ModuleState::Running
2565 | ModuleState::Unresponsive
2566 | ModuleState::Restarting
2567 | ModuleState::Draining
2568 | ModuleState::Disabled => {}
2569 }
2570
2571 let (reply_tx, reply_rx) = oneshot::channel();
2572 self.inner
2573 .commands
2574 .send(SupervisorCommand::Retire { reply: reply_tx })
2575 .await
2576 .map_err(|_| SuperviseError::CommandClosed {
2577 module_id: self.inner.module_id.clone(),
2578 })?;
2579 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2580 module_id: self.inner.module_id.clone(),
2581 })?
2582 }
2583
2584 pub async fn stop(&self) -> Result<(), SuperviseError> {
2585 match self.state()? {
2586 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2587 ModuleState::Starting
2588 | ModuleState::Running
2589 | ModuleState::Unresponsive
2590 | ModuleState::Restarting
2591 | ModuleState::Draining
2592 | ModuleState::Disabled => {}
2593 }
2594
2595 let (reply_tx, reply_rx) = oneshot::channel();
2596 self.inner
2597 .commands
2598 .send(SupervisorCommand::Drain { reply: reply_tx })
2599 .await
2600 .map_err(|_| SuperviseError::CommandClosed {
2601 module_id: self.inner.module_id.clone(),
2602 })?;
2603 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2604 module_id: self.inner.module_id.clone(),
2605 })?
2606 }
2607
2608 pub async fn restart(&self, drain_timeout_ms: Option<u64>) -> Result<(), SuperviseError> {
2609 let (reply_tx, reply_rx) = oneshot::channel();
2610 self.inner
2611 .commands
2612 .send(SupervisorCommand::Restart {
2613 drain_timeout_ms,
2614 reply: reply_tx,
2615 })
2616 .await
2617 .map_err(|_| SuperviseError::CommandClosed {
2618 module_id: self.inner.module_id.clone(),
2619 })?;
2620 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2621 module_id: self.inner.module_id.clone(),
2622 })?
2623 }
2624
2625 pub async fn swap(&self, ready_timeout: Option<Duration>) -> Result<(), SuperviseError> {
2630 let (reply_tx, reply_rx) = oneshot::channel();
2631 self.inner
2632 .commands
2633 .send(SupervisorCommand::Swap {
2634 ready_timeout,
2635 reply: reply_tx,
2636 })
2637 .await
2638 .map_err(|_| SuperviseError::CommandClosed {
2639 module_id: self.inner.module_id.clone(),
2640 })?;
2641 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2642 module_id: self.inner.module_id.clone(),
2643 })?
2644 }
2645
2646 pub async fn reload(&self) -> Result<(), SuperviseError> {
2647 let (reply_tx, reply_rx) = oneshot::channel();
2648 self.inner
2649 .commands
2650 .send(SupervisorCommand::Reload { reply: reply_tx })
2651 .await
2652 .map_err(|_| SuperviseError::CommandClosed {
2653 module_id: self.inner.module_id.clone(),
2654 })?;
2655 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2656 module_id: self.inner.module_id.clone(),
2657 })?
2658 }
2659
2660 pub async fn set_enabled(&self, enabled: bool) -> Result<bool, SuperviseError> {
2661 let (reply_tx, reply_rx) = oneshot::channel();
2662 self.inner
2663 .commands
2664 .send(SupervisorCommand::SetEnabled {
2665 enabled,
2666 reply: reply_tx,
2667 })
2668 .await
2669 .map_err(|_| SuperviseError::CommandClosed {
2670 module_id: self.inner.module_id.clone(),
2671 })?;
2672 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2673 module_id: self.inner.module_id.clone(),
2674 })?
2675 }
2676
2677 pub(crate) fn declared_protocol(&self) -> Result<ModuleProtocol, SuperviseError> {
2682 Ok(self
2683 .inner
2684 .configuration
2685 .lock()
2686 .map_err(|_| SuperviseError::StatePoisoned {
2687 module_id: Some(self.inner.module_id.clone()),
2688 })?
2689 .spec
2690 .protocol)
2691 }
2692
2693 pub(crate) fn configuration(&self) -> Result<(ModuleSpec, HealthConfig), SuperviseError> {
2694 let configuration =
2695 self.inner
2696 .configuration
2697 .lock()
2698 .map_err(|_| SuperviseError::StatePoisoned {
2699 module_id: Some(self.inner.module_id.clone()),
2700 })?;
2701 Ok((configuration.spec.clone(), configuration.health))
2702 }
2703
2704 #[cfg(any(test, feature = "test-support"))]
2708 pub async fn update_spec_for_test(&self, spec: ModuleSpec) -> Result<(), SuperviseError> {
2709 let (_, health) = self.configuration()?;
2710 let drain_timeout_ms = u64::try_from(
2711 self.inner
2712 .effective_drain_timeout
2713 .lock()
2714 .unwrap_or_else(|poisoned| poisoned.into_inner())
2715 .as_millis(),
2716 )
2717 .ok();
2718 self.update_configuration(spec, health, drain_timeout_ms)
2719 .await
2720 }
2721
2722 pub(crate) async fn update_configuration(
2723 &self,
2724 spec: ModuleSpec,
2725 health: HealthConfig,
2726 drain_timeout_ms: Option<u64>,
2727 ) -> Result<(), SuperviseError> {
2728 if spec.module_id != self.inner.module_id {
2729 return Err(SuperviseError::InvalidSpec {
2730 reason: "a supervised module's module_id cannot be changed".to_string(),
2731 });
2732 }
2733 validate_spec(&spec)?;
2734 let (reply_tx, reply_rx) = oneshot::channel();
2735 self.inner
2736 .commands
2737 .send(SupervisorCommand::UpdateConfiguration {
2738 spec: spec.clone(),
2739 health,
2740 drain_timeout_ms,
2741 reply: reply_tx,
2742 })
2743 .await
2744 .map_err(|_| SuperviseError::CommandClosed {
2745 module_id: self.inner.module_id.clone(),
2746 })?;
2747 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2748 module_id: self.inner.module_id.clone(),
2749 })?;
2750 let mut configuration =
2751 self.inner
2752 .configuration
2753 .lock()
2754 .map_err(|_| SuperviseError::StatePoisoned {
2755 module_id: Some(self.inner.module_id.clone()),
2756 })?;
2757 configuration.spec = spec;
2758 configuration.health = health;
2759 Ok(())
2760 }
2761}
2762
2763impl Drop for SupervisedModuleInner {
2764 fn drop(&mut self) {
2765 let Ok(mut monitor) = self.monitor.lock() else {
2766 return;
2767 };
2768 if let Some(monitor) = monitor.as_ref().filter(|monitor| !monitor.is_finished()) {
2769 let _ = update_snapshot(&self.snapshot, Some(&self.module_id), |state| {
2770 state.state = ModuleState::Stopped;
2771 clear_current_process_facts(state);
2772 });
2773 monitor.abort();
2774 }
2775 let _ = monitor.take();
2776 }
2777}
2778
2779#[derive(Debug)]
2780enum SupervisorCommand {
2781 Drain {
2782 reply: oneshot::Sender<Result<(), SuperviseError>>,
2783 },
2784 Retire {
2785 reply: oneshot::Sender<Result<(), SuperviseError>>,
2786 },
2787 Restart {
2788 drain_timeout_ms: Option<u64>,
2793 reply: oneshot::Sender<Result<(), SuperviseError>>,
2794 },
2795 Reload {
2796 reply: oneshot::Sender<Result<(), SuperviseError>>,
2797 },
2798 SetEnabled {
2799 enabled: bool,
2800 reply: oneshot::Sender<Result<bool, SuperviseError>>,
2801 },
2802 UpdateConfiguration {
2803 spec: ModuleSpec,
2804 health: HealthConfig,
2805 drain_timeout_ms: Option<u64>,
2808 reply: oneshot::Sender<()>,
2809 },
2810 Swap {
2811 ready_timeout: Option<Duration>,
2814 reply: oneshot::Sender<Result<(), SuperviseError>>,
2816 },
2817}
2818
2819#[derive(Debug)]
2820pub enum SuperviseError {
2821 InvalidSpec {
2822 reason: String,
2823 },
2824 Spawn {
2825 program: PathBuf,
2826 source: io::Error,
2827 cgroup_path: Option<PathBuf>,
2828 },
2829 Cgroup {
2830 module_id: String,
2831 source: io::Error,
2832 },
2833 LaunchNonce {
2836 reason: String,
2837 },
2838 Wait {
2839 module_id: String,
2840 source: io::Error,
2841 },
2842 Kill {
2843 module_id: String,
2844 source: io::Error,
2845 },
2846 Forwarding(ForwardingError),
2847 Registry(RegistryError),
2848 ReloadUnavailable {
2849 module_id: String,
2850 reason: String,
2851 },
2852 Disabled {
2857 module_id: String,
2858 },
2859 ReloadFailed {
2860 module_id: String,
2861 reason: String,
2862 },
2863 RegistrationStillActive {
2864 module_id: String,
2865 waited: Duration,
2866 },
2867 StatePoisoned {
2868 module_id: Option<String>,
2869 },
2870 CommandClosed {
2871 module_id: String,
2872 },
2873 SwapInProgress {
2877 module_id: String,
2878 },
2879 SwapRefused {
2881 module_id: String,
2882 reason: SwapRefusal,
2883 },
2884 SwapFailed {
2888 module_id: String,
2889 arm: SwapFailureArm,
2890 detail: String,
2891 candidate_exit: Option<ExitReport>,
2894 },
2895}
2896
2897#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2899pub enum SwapRefusal {
2900 OverlapExclusive,
2902 NotRegistered,
2905 ProtocolNone,
2908 NotConfigured,
2911 AlreadySwapping,
2913}
2914
2915impl SwapRefusal {
2916 pub fn as_str(self) -> &'static str {
2917 match self {
2918 Self::OverlapExclusive => "overlap_exclusive",
2919 Self::NotRegistered => "not_registered",
2920 Self::ProtocolNone => "protocol_none",
2921 Self::NotConfigured => "not_configured",
2922 Self::AlreadySwapping => "already_swapping",
2923 }
2924 }
2925}
2926
2927#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2930pub enum SwapFailureArm {
2931 SpawnFailed,
2933 NeverRegistered,
2935 NeverReady,
2937 CandidateExited,
2939 CandidateUnhealthy,
2941 Interrupted,
2945 CutoverLost,
2950}
2951
2952impl SwapFailureArm {
2953 pub fn as_str(self) -> &'static str {
2954 match self {
2955 Self::SpawnFailed => "spawn_failed",
2956 Self::NeverRegistered => "never_registered",
2957 Self::NeverReady => "never_ready",
2958 Self::CandidateExited => "candidate_exited",
2959 Self::CandidateUnhealthy => "candidate_unhealthy",
2960 Self::Interrupted => "interrupted",
2961 Self::CutoverLost => "cutover_lost",
2962 }
2963 }
2964}
2965
2966impl fmt::Display for SuperviseError {
2967 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2968 match self {
2969 Self::InvalidSpec { reason } => write!(f, "invalid module spec: {reason}"),
2970 Self::Spawn {
2971 program,
2972 source,
2973 cgroup_path: Some(cgroup_path),
2974 } => write!(
2975 f,
2976 "failed to place module in cgroup '{}' while spawning '{}': {source}",
2977 cgroup_path.display(),
2978 program.display()
2979 ),
2980 Self::Spawn {
2981 program,
2982 source,
2983 cgroup_path: None,
2984 } => write!(
2985 f,
2986 "failed to spawn module '{}': {source}",
2987 program.display()
2988 ),
2989 Self::Cgroup { module_id, source } => {
2990 write!(
2991 f,
2992 "failed to prepare cgroup for module '{module_id}': {source}"
2993 )
2994 }
2995 Self::LaunchNonce { reason } => {
2996 write!(
2997 f,
2998 "failed to generate reserved-module launch nonce: {reason}"
2999 )
3000 }
3001 Self::Wait { module_id, source } => {
3002 write!(f, "failed to wait for module '{module_id}': {source}")
3003 }
3004 Self::Kill { module_id, source } => {
3005 write!(f, "failed to kill module '{module_id}': {source}")
3006 }
3007 Self::Forwarding(err) => write!(f, "forwarding error: {err}"),
3008 Self::Registry(err) => write!(f, "registry error: {err}"),
3009 Self::ReloadUnavailable { module_id, reason } => {
3010 write!(f, "reload unavailable for module '{module_id}': {reason}")
3011 }
3012 Self::Disabled { module_id } => {
3013 write!(
3014 f,
3015 "module '{module_id}' is disabled; enable it before restart or reload"
3016 )
3017 }
3018 Self::ReloadFailed { module_id, reason } => {
3019 write!(f, "reload failed for module '{module_id}': {reason}")
3020 }
3021 Self::RegistrationStillActive { module_id, waited } => write!(
3022 f,
3023 "module '{module_id}' registration remained active after waiting {waited:?}"
3024 ),
3025 Self::StatePoisoned { module_id } => match module_id {
3026 Some(module_id) => {
3027 write!(f, "supervisor state for module '{module_id}' was poisoned")
3028 }
3029 None => write!(f, "supervisor state was poisoned"),
3030 },
3031 Self::CommandClosed { module_id } => {
3032 write!(
3033 f,
3034 "supervisor command channel for module '{module_id}' is closed"
3035 )
3036 }
3037 Self::SwapInProgress { module_id } => write!(
3038 f,
3039 "module '{module_id}' is being swapped; retry once the swap has cut over or failed, or stop the module to abort the swap"
3040 ),
3041 Self::SwapRefused { module_id, reason } => match reason {
3042 SwapRefusal::OverlapExclusive => write!(
3043 f,
3044 "module '{module_id}' is declared overlap: \"exclusive\" (the default): two processes of it must not run at once, so it cannot be swapped; use a plain restart, or declare overlap: \"safe\" in its config if it really tolerates a second process"
3045 ),
3046 SwapRefusal::NotRegistered => write!(
3047 f,
3048 "module '{module_id}' is not registered, so there is no serving process to keep while a replacement warms; use a plain restart"
3049 ),
3050 SwapRefusal::ProtocolNone => write!(
3051 f,
3052 "module '{module_id}' is protocol: \"none\" and never registers, so a swap could never see its replacement become ready; use a plain restart"
3053 ),
3054 SwapRefusal::NotConfigured => write!(
3055 f,
3056 "module '{module_id}' cannot be swapped: the supervisor was built without the forwarding table or shared handle a swap needs"
3057 ),
3058 SwapRefusal::AlreadySwapping => {
3059 write!(f, "module '{module_id}' is already being swapped")
3060 }
3061 },
3062 Self::SwapFailed {
3063 module_id,
3064 arm,
3065 detail,
3066 ..
3067 } => write!(
3068 f,
3069 "swap of module '{module_id}' failed ({}): {detail}; the running process was left serving",
3070 arm.as_str()
3071 ),
3072 }
3073 }
3074}
3075
3076impl Error for SuperviseError {
3077 fn source(&self) -> Option<&(dyn Error + 'static)> {
3078 match self {
3079 Self::Spawn { source, .. }
3080 | Self::Cgroup { source, .. }
3081 | Self::Wait { source, .. }
3082 | Self::Kill { source, .. } => Some(source),
3083 Self::Forwarding(err) => Some(err),
3084 Self::Registry(err) => Some(err),
3085 Self::LaunchNonce { .. }
3086 | Self::InvalidSpec { .. }
3087 | Self::ReloadUnavailable { .. }
3088 | Self::Disabled { .. }
3089 | Self::ReloadFailed { .. }
3090 | Self::RegistrationStillActive { .. }
3091 | Self::StatePoisoned { .. }
3092 | Self::CommandClosed { .. }
3093 | Self::SwapInProgress { .. }
3094 | Self::SwapRefused { .. }
3095 | Self::SwapFailed { .. } => None,
3096 }
3097 }
3098}
3099
3100pub(crate) fn validate_spec(spec: &ModuleSpec) -> Result<(), SuperviseError> {
3101 if spec.module_id.trim().is_empty() {
3102 return Err(SuperviseError::InvalidSpec {
3103 reason: "module_id must not be empty".to_string(),
3104 });
3105 }
3106
3107 Ok(())
3108}
3109
3110#[derive(Debug, Default)]
3111struct HealthProbeRuntime {
3112 registered_connection: Option<crate::ConnectionId>,
3113 advertised: bool,
3114 next_probe_at: Option<Instant>,
3115 probe_index: u64,
3116}
3117
3118impl HealthProbeRuntime {
3119 fn refresh_registration(
3120 &mut self,
3121 spec: &ModuleSpec,
3122 runtime: &SupervisorRuntimeConfig,
3123 registry: &Registry,
3124 snapshot: &SharedSnapshot,
3125 ) {
3126 if spec.protocol == ModuleProtocol::None {
3138 self.registered_connection = None;
3139 self.advertised = false;
3140 self.next_probe_at = None;
3141 return;
3142 }
3143
3144 let registration = match registry.get_module(&spec.module_id) {
3145 Ok(registration) => registration,
3146 Err(err) => {
3147 warn!(module_id = %spec.module_id, error = %err, "health prober could not read registry");
3148 self.advertised = false;
3149 self.next_probe_at = None;
3150 return;
3151 }
3152 };
3153
3154 let Some(registration) = registration else {
3155 self.registered_connection = None;
3156 self.advertised = false;
3157 self.next_probe_at = None;
3158 return;
3159 };
3160
3161 let advertised = registration
3162 .control_ops
3163 .iter()
3164 .any(|op| op == MODULE_CONTROL_OP_HEALTH_CHECK);
3165 if !advertised {
3166 self.registered_connection = Some(registration.connection_id);
3167 self.advertised = false;
3168 self.next_probe_at = None;
3169 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3170 state.health.status = SupervisorHealthStatus::Unknown;
3171 state.health.consecutive_failures = 0;
3172 state.health.last_probe_ms = None;
3173 state.health.detail = None;
3174 state.health.metrics = None;
3175 });
3176 return;
3177 }
3178
3179 let reregistered = self.registered_connection != Some(registration.connection_id);
3180 self.registered_connection = Some(registration.connection_id);
3181 self.advertised = true;
3182 if reregistered || self.next_probe_at.is_none() {
3183 self.probe_index = 0;
3184 self.next_probe_at = Some(
3185 Instant::now() + jittered_health_delay(&spec.module_id, 0, runtime.health.cadence),
3186 );
3187 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3188 state.health.status = SupervisorHealthStatus::Unknown;
3189 state.health.consecutive_failures = 0;
3190 state.health.detail = None;
3191 state.health.metrics = None;
3192 });
3193 }
3194 }
3195
3196 fn wake_after(&self) -> Duration {
3197 if !self.advertised {
3198 return REGISTRY_RELEASE_POLL;
3199 }
3200 self.next_probe_at
3201 .map(|next| next.saturating_duration_since(Instant::now()))
3202 .unwrap_or(REGISTRY_RELEASE_POLL)
3203 }
3204
3205 fn due(&self) -> bool {
3206 self.advertised
3207 && self
3208 .next_probe_at
3209 .is_some_and(|next| Instant::now() >= next)
3210 }
3211
3212 fn schedule_next(&mut self, spec: &ModuleSpec, cadence: Duration) {
3213 self.probe_index = self.probe_index.wrapping_add(1);
3214 self.next_probe_at = Some(
3215 Instant::now() + jittered_health_delay(&spec.module_id, self.probe_index, cadence),
3216 );
3217 }
3218}
3219
3220#[derive(Debug)]
3255enum HealthProbeEvidence {
3256 LaneDead,
3258 NoAnswer,
3260 BadAnswer,
3262 Misconfigured,
3264}
3265
3266#[derive(Debug)]
3267struct HealthProbeError {
3268 evidence: HealthProbeEvidence,
3269 message: String,
3270}
3271
3272impl HealthProbeError {
3273 fn lane_dead(message: impl Into<String>) -> Self {
3274 Self::with(HealthProbeEvidence::LaneDead, message)
3275 }
3276
3277 fn no_answer(message: impl Into<String>) -> Self {
3278 Self::with(HealthProbeEvidence::NoAnswer, message)
3279 }
3280
3281 fn bad_answer(message: impl Into<String>) -> Self {
3282 Self::with(HealthProbeEvidence::BadAnswer, message)
3283 }
3284
3285 fn misconfigured(message: impl Into<String>) -> Self {
3286 Self::with(HealthProbeEvidence::Misconfigured, message)
3287 }
3288
3289 fn with(evidence: HealthProbeEvidence, message: impl Into<String>) -> Self {
3290 Self {
3291 evidence,
3292 message: message.into(),
3293 }
3294 }
3295
3296 #[allow(dead_code)]
3310 fn is_proof_of_death(&self) -> bool {
3311 matches!(self.evidence, HealthProbeEvidence::LaneDead)
3312 }
3313
3314 fn label(&self) -> &'static str {
3322 match self.evidence {
3323 HealthProbeEvidence::LaneDead => "lane-dead",
3324 HealthProbeEvidence::NoAnswer => "no-answer",
3325 HealthProbeEvidence::BadAnswer => "bad-answer",
3326 HealthProbeEvidence::Misconfigured => "daemon-misconfigured",
3327 }
3328 }
3329}
3330
3331impl fmt::Display for HealthProbeError {
3332 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3333 f.write_str(&self.message)
3334 }
3335}
3336
3337async fn run_health_probe_cycle(
3338 spec: &ModuleSpec,
3339 runtime: &SupervisorRuntimeConfig,
3340 registry: &Registry,
3341 process_liveness: &SupervisorProcessLiveness,
3342 snapshot: &SharedSnapshot,
3343 child: &mut Option<SupervisedChild>,
3344) {
3345 let now_ms = unix_ms_now();
3346 match probe_module_health(&spec.module_id, runtime, None).await {
3347 Ok(report) => {
3348 handle_health_report(
3349 spec,
3350 runtime,
3351 registry,
3352 process_liveness,
3353 snapshot,
3354 child,
3355 report,
3356 now_ms,
3357 )
3358 .await;
3359 }
3360 Err(err) => {
3361 handle_health_probe_failure(
3362 spec,
3363 runtime,
3364 registry,
3365 process_liveness,
3366 snapshot,
3367 child,
3368 err,
3369 now_ms,
3370 )
3371 .await;
3372 }
3373 }
3374}
3375
3376async fn probe_module_health(
3377 module_id: &str,
3378 runtime: &SupervisorRuntimeConfig,
3379 drain_deadline: Option<Instant>,
3380) -> Result<HealthReport, HealthProbeError> {
3381 let Some(forwarding) = runtime.forwarding.as_ref() else {
3382 return Err(HealthProbeError::misconfigured(
3383 "supervisor was not configured with a forwarding table",
3384 ));
3385 };
3386 let probe_started_at = Instant::now();
3387 let mut deadline = probe_started_at + runtime.health.deadline;
3388 if let Some(drain_deadline) = drain_deadline {
3389 deadline = deadline.min(drain_deadline);
3390 }
3391 let pending = if drain_deadline.is_some() {
3392 forwarding.begin_drain_health_probe_rpc_for(
3393 module_id,
3394 MODULE_CONTROL_OP_HEALTH_CHECK,
3395 probe_started_at,
3396 deadline,
3397 )
3398 } else {
3399 forwarding.begin_health_probe_rpc_for(
3400 module_id,
3401 MODULE_CONTROL_OP_HEALTH_CHECK,
3402 probe_started_at,
3403 deadline,
3404 )
3405 }
3406 .map_err(|err| {
3407 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3410 })?;
3411 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3412}
3413
3414async fn probe_endpoint_health(
3421 endpoint: crate::ModuleEndpointId,
3422 runtime: &SupervisorRuntimeConfig,
3423 deadline_cap: Option<Instant>,
3424) -> Result<HealthReport, HealthProbeError> {
3425 let Some(forwarding) = runtime.forwarding.as_ref() else {
3426 return Err(HealthProbeError::misconfigured(
3427 "supervisor was not configured with a forwarding table",
3428 ));
3429 };
3430 let probe_started_at = Instant::now();
3431 let mut deadline = probe_started_at + runtime.health.deadline;
3432 if let Some(cap) = deadline_cap {
3433 deadline = deadline.min(cap);
3434 }
3435 let pending = forwarding
3436 .begin_endpoint_health_probe_rpc_for(
3437 endpoint,
3438 MODULE_CONTROL_OP_HEALTH_CHECK,
3439 probe_started_at,
3440 deadline,
3441 )
3442 .map_err(|err| {
3443 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3444 })?;
3445 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3446}
3447
3448async fn await_health_probe(
3450 forwarding: &ForwardingTable,
3451 pending: PendingModuleControlRpc,
3452 deadline: Instant,
3453 probe_budget: Duration,
3454) -> Result<HealthReport, HealthProbeError> {
3455 let PendingModuleControlRpc {
3456 endpoint,
3457 module_sink,
3458 negotiated_ver,
3459 corr,
3460 receiver,
3461 } = pending;
3462 let body = serde_json::to_vec(&ModuleControlRequest::HealthCheck {}).map_err(|err| {
3463 HealthProbeError::misconfigured(format!("failed to encode health.check: {err}"))
3464 })?;
3465 let frame = Frame::build_with_version(
3466 negotiated_ver,
3467 FrameType::Request,
3468 control_flags(),
3469 0,
3470 0,
3471 corr,
3472 body,
3473 )
3474 .map_err(|err| {
3475 HealthProbeError::misconfigured(format!("failed to build health.check frame: {err}"))
3476 })?;
3477
3478 match timeout_at(deadline, module_sink.send(frame)).await {
3484 Ok(Ok(())) => {}
3485 Ok(Err(err)) => {
3486 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3487 return Err(HealthProbeError::lane_dead(format!(
3490 "failed to send health.check: {err}"
3491 )));
3492 }
3493 Err(_elapsed) => {
3494 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3495 return Err(HealthProbeError::no_answer(
3499 "health.check send timed out before enqueue (module egress full)",
3500 ));
3501 }
3502 }
3503
3504 match timeout_at(deadline, receiver).await {
3505 Ok(Ok(ModuleControlRpcOutcome::Response(response))) => {
3509 response.health_report().ok_or_else(|| {
3510 HealthProbeError::bad_answer("health.check RPC returned a non-health response")
3511 })
3512 }
3513 Ok(Ok(ModuleControlRpcOutcome::Rejected(body))) => Err(HealthProbeError::bad_answer(
3514 format!("health.check rejected: {}", body.message),
3515 )),
3516 Ok(Ok(ModuleControlRpcOutcome::ModuleGone(message))) => {
3517 Err(HealthProbeError::lane_dead(message))
3518 }
3519 Ok(Ok(ModuleControlRpcOutcome::MalformedResponse(message))) => {
3520 Err(HealthProbeError::bad_answer(message))
3521 }
3522 Ok(Ok(ModuleControlRpcOutcome::UnexpectedOp { expected, actual })) => {
3523 Err(HealthProbeError::bad_answer(format!(
3524 "expected module-control op '{expected}', got '{actual}'"
3525 )))
3526 }
3527 Ok(Ok(ModuleControlRpcOutcome::DeadlineElapsed)) => Err(HealthProbeError::bad_answer(
3531 "module answered health.check after its daemon deadline",
3532 )),
3533 Ok(Err(_)) => Err(HealthProbeError::misconfigured(
3534 "health.check waiter was canceled before the module responded",
3535 )),
3536 Err(_) => {
3537 let _ = forwarding.tombstone_health_probe_rpc(endpoint, corr);
3538 Err(HealthProbeError::no_answer(format!(
3539 "module did not answer health.check within {probe_budget:?}"
3540 )))
3541 }
3542 }
3543}
3544
3545#[allow(clippy::too_many_arguments)]
3546async fn handle_health_report(
3547 spec: &ModuleSpec,
3548 runtime: &SupervisorRuntimeConfig,
3549 registry: &Registry,
3550 process_liveness: &SupervisorProcessLiveness,
3551 snapshot: &SharedSnapshot,
3552 child: &mut Option<SupervisedChild>,
3553 report: HealthReport,
3554 now_ms: u64,
3555) {
3556 let status = supervisor_health_status(report.status);
3557 let detail = report.detail.clone();
3558 let metrics = truncate_health_metrics(report.metrics);
3559 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3560 state.health.status = status;
3561 state.health.last_probe_ms = Some(now_ms);
3562 state.health.detail = detail.clone();
3563 state.health.metrics = metrics.clone();
3564 state.health.consecutive_failures = 0;
3565 });
3566
3567 let action = match report.status {
3568 HealthStatus::Ok => return,
3569 HealthStatus::Degraded => runtime.health.on_degraded,
3570 HealthStatus::Failing => runtime.health.on_failing,
3571 };
3572 apply_l3_health_action(
3573 spec,
3574 runtime,
3575 registry,
3576 process_liveness,
3577 snapshot,
3578 child,
3579 status,
3580 detail.as_deref(),
3581 action,
3582 now_ms,
3583 )
3584 .await;
3585}
3586
3587#[allow(clippy::too_many_arguments)]
3588async fn handle_health_probe_failure(
3589 spec: &ModuleSpec,
3590 runtime: &SupervisorRuntimeConfig,
3591 registry: &Registry,
3592 process_liveness: &SupervisorProcessLiveness,
3593 snapshot: &SharedSnapshot,
3594 child: &mut Option<SupervisedChild>,
3595 err: HealthProbeError,
3596 now_ms: u64,
3597) {
3598 let threshold = runtime.health.failure_threshold.max(1);
3599 let mut failures = 0;
3600 let detail = format!("[{}] {err}", err.label());
3605 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3606 state.health.last_probe_ms = Some(now_ms);
3607 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
3608 state.health.detail = Some(detail.clone());
3609 state.health.metrics = None;
3610 failures = state.health.consecutive_failures;
3611 });
3612
3613 if failures < threshold {
3614 warn!(
3615 module_id = %spec.module_id,
3616 consecutive_failures = failures,
3617 threshold,
3618 evidence = err.label(),
3619 detail = %detail,
3620 "health.check probe failed"
3621 );
3622 return;
3623 }
3624
3625 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3626 state.state = ModuleState::Unresponsive;
3627 state.health.status = SupervisorHealthStatus::Unresponsive;
3628 });
3629 if runtime.health.critical {
3633 error!(
3634 module_id = %spec.module_id,
3635 status = "unresponsive",
3636 evidence = err.label(),
3637 detail = %detail,
3638 "critical module health alert"
3639 );
3640 } else {
3641 warn!(
3642 module_id = %spec.module_id,
3643 status = "unresponsive",
3644 evidence = err.label(),
3645 detail = %detail,
3646 "module health threshold breached"
3647 );
3648 }
3649 if let Err(err) = health_restart_child(
3650 spec,
3651 runtime,
3652 registry,
3653 process_liveness,
3654 snapshot,
3655 child,
3656 SupervisorHealthStatus::Unresponsive,
3657 Some(&detail),
3658 now_ms,
3659 )
3660 .await
3661 {
3662 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3663 }
3664}
3665
3666#[allow(clippy::too_many_arguments)]
3667async fn apply_l3_health_action(
3668 spec: &ModuleSpec,
3669 runtime: &SupervisorRuntimeConfig,
3670 registry: &Registry,
3671 process_liveness: &SupervisorProcessLiveness,
3672 snapshot: &SharedSnapshot,
3673 child: &mut Option<SupervisedChild>,
3674 status: SupervisorHealthStatus,
3675 detail: Option<&str>,
3676 action: HealthAction,
3677 now_ms: u64,
3678) {
3679 record_health_action(snapshot, &spec.module_id, action.to_string(), now_ms);
3680 match action {
3681 HealthAction::Report => {
3682 info!(
3683 module_id = %spec.module_id,
3684 status = ?status,
3685 detail,
3686 "module reported non-ok health"
3687 );
3688 }
3689 HealthAction::Alert => {
3690 error!(
3691 module_id = %spec.module_id,
3692 status = ?status,
3693 detail,
3694 "module health alert"
3695 );
3696 }
3697 HealthAction::Restart => {
3698 if let Err(err) = health_restart_child(
3699 spec,
3700 runtime,
3701 registry,
3702 process_liveness,
3703 snapshot,
3704 child,
3705 status,
3706 detail,
3707 now_ms,
3708 )
3709 .await
3710 {
3711 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3712 }
3713 }
3714 }
3715}
3716
3717#[allow(clippy::too_many_arguments)]
3718async fn health_restart_child(
3719 spec: &ModuleSpec,
3720 runtime: &SupervisorRuntimeConfig,
3721 registry: &Registry,
3722 process_liveness: &SupervisorProcessLiveness,
3723 snapshot: &SharedSnapshot,
3724 child: &mut Option<SupervisedChild>,
3725 status: SupervisorHealthStatus,
3726 detail: Option<&str>,
3727 now_ms: u64,
3728) -> Result<(), SuperviseError> {
3729 let (enabled, schedule) = {
3730 let mut state = lock_snapshot(snapshot)?;
3731 let enabled = state.enabled;
3732 let schedule = if enabled {
3733 state.next_crash_restart(&runtime.restart_policy, Instant::now())
3734 } else {
3735 None
3736 };
3737 (enabled, schedule)
3738 };
3739
3740 if !enabled {
3741 return Err(SuperviseError::Disabled {
3742 module_id: spec.module_id.clone(),
3743 });
3744 }
3745
3746 if schedule.is_none() {
3747 record_health_action(snapshot, &spec.module_id, "disabled".to_string(), now_ms);
3748 error!(
3749 module_id = %spec.module_id,
3750 status = ?status,
3751 detail,
3752 max_restarts = runtime.restart_policy.max_restarts,
3753 window_secs = runtime.restart_policy.window.as_secs(),
3754 "health restart budget exhausted; disabling module"
3755 );
3756 begin_forwarding_drain_if_configured(
3757 spec,
3758 runtime,
3759 registry,
3760 snapshot,
3761 Some(false),
3762 RouteCloseReason::Disable,
3763 )
3764 .await?;
3765 drain_optional_child(
3766 &spec.module_id,
3767 spec.protocol,
3768 registry,
3769 snapshot,
3770 &runtime.terminal_ring,
3771 &runtime.spawn_events,
3772 child,
3773 runtime.drain_timeout,
3774 ModuleState::Disabled,
3775 Some(false),
3776 )
3777 .await?;
3778 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3779 return Ok(());
3780 }
3781
3782 let schedule = schedule.expect("a health restart must have a crash-restart schedule");
3783 let mut restart_count = 0;
3784 update_snapshot(snapshot, Some(&spec.module_id), |state| {
3785 restart_count = state.crash_restarts.len();
3786 state.state = ModuleState::Unresponsive;
3787 state.health.status = status;
3788 state.health.last_action = Some(HealthAction::Restart.to_string());
3789 state.health.last_action_ms = Some(now_ms);
3790 })?;
3791 warn!(
3792 module_id = %spec.module_id,
3793 status = ?status,
3794 detail,
3795 restart_count,
3796 restart_in_window = schedule.restart_in_window,
3797 delay_ms = schedule.delay.as_millis() as u64,
3798 "health-triggered module restart"
3799 );
3800
3801 begin_forwarding_drain_if_configured(
3802 spec,
3803 runtime,
3804 registry,
3805 snapshot,
3806 Some(true),
3807 RouteCloseReason::Restart,
3808 )
3809 .await?;
3810 drain_optional_child(
3811 &spec.module_id,
3812 spec.protocol,
3813 registry,
3814 snapshot,
3815 &runtime.terminal_ring,
3816 &runtime.spawn_events,
3817 child,
3818 runtime.drain_timeout,
3819 ModuleState::Restarting,
3820 Some(true),
3821 )
3822 .await?;
3823 sleep(schedule.delay).await;
3824 if !respawn_still_pending(snapshot) {
3828 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3829 return Ok(());
3830 }
3831 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
3832 match spawn_and_mark_running(spec, runtime, snapshot) {
3833 Ok(next_child) => {
3834 *child = Some(next_child);
3835 Ok(())
3836 }
3837 Err(err) => {
3838 fail_snapshot(snapshot, Some(&spec.module_id), None);
3839 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3840 *child = None;
3841 Err(err)
3842 }
3843 }
3844}
3845
3846fn record_health_action(snapshot: &SharedSnapshot, module_id: &str, action: String, now_ms: u64) {
3847 let _ = update_snapshot(snapshot, Some(module_id), |state| {
3848 state.health.last_action = Some(action);
3849 state.health.last_action_ms = Some(now_ms);
3850 });
3851}
3852
3853fn supervisor_health_status(status: HealthStatus) -> SupervisorHealthStatus {
3854 match status {
3855 HealthStatus::Ok => SupervisorHealthStatus::Ok,
3856 HealthStatus::Degraded => SupervisorHealthStatus::Degraded,
3857 HealthStatus::Failing => SupervisorHealthStatus::Failing,
3858 }
3859}
3860
3861fn truncate_health_metrics(metrics: Option<Value>) -> Option<Value> {
3873 let metrics = metrics?;
3874 match serde_json::to_vec(&metrics) {
3875 Ok(encoded) if encoded.len() > MAX_HEALTH_METRICS_BYTES => Some(serde_json::json!({
3876 "truncated": true,
3877 "original_bytes": encoded.len(),
3878 })),
3879 Ok(_) | Err(_) => Some(metrics),
3880 }
3881}
3882
3883fn jittered_health_delay(module_id: &str, probe_index: u64, cadence: Duration) -> Duration {
3889 if cadence.is_zero() {
3890 return Duration::ZERO;
3891 }
3892 let cadence_ms = cadence.as_millis() as u64;
3893 if cadence_ms == 0 {
3909 return cadence;
3910 }
3911 let jitter_span = (cadence_ms / 10).max(1);
3926 let hash = module_id.as_bytes().iter().fold(
3927 probe_index.wrapping_mul(0x9E37_79B9_7F4A_7C15),
3928 |acc, byte| {
3929 acc.wrapping_mul(1099511628211)
3930 .wrapping_add(u64::from(*byte))
3931 },
3932 );
3933 cadence + Duration::from_millis(hash % jitter_span)
3934}
3935
3936#[cfg(test)]
3937mod tests {
3938 use super::*;
3939
3940 #[test]
3941 fn readding_a_module_clears_its_rescan_removal_tombstone() {
3942 let handle = SupervisorHandle::new();
3943 let module_id = "readded-tombstone";
3944 handle.record_rescan_removal(module_id);
3945 assert!(handle.removal_tombstone_age_ms(module_id).is_some());
3946
3947 handle.apply_identity_configuration(&ModuleSpec {
3948 module_id: module_id.to_string(),
3949 program: PathBuf::from("/test/module"),
3950 args: Vec::new(),
3951 env: Vec::new(),
3952 reserved: false,
3953 reserved_prefixes: Vec::new(),
3954 protocol: ModuleProtocol::Subc,
3955 overlap: Default::default(),
3956 });
3957
3958 assert!(
3959 handle.removal_tombstone_age_ms(module_id).is_none(),
3960 "a re-added module must not retain a stale removal tombstone"
3961 );
3962 }
3963
3964 fn stale_process_snapshot(state: ModuleState, enabled: bool) -> SharedSnapshot {
3965 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::new(state, enabled)));
3966 update_snapshot(&snapshot, Some("stale-process-facts"), |snapshot| {
3967 snapshot.process_alive = true;
3968 snapshot.pid = Some(41);
3969 snapshot.spawned_at_ms = Some(42);
3970 snapshot.spawned_from = Some(PathBuf::from("/spawned/module"));
3971 snapshot.spawned_file_identity = Some(SpawnedFileIdentity {
3972 device: 43,
3973 inode: 44,
3974 });
3975 })
3976 .unwrap();
3977 snapshot
3978 }
3979
3980 fn assert_snapshot_process_facts_cleared(snapshot: &SharedSnapshot) {
3981 let snapshot = lock_snapshot(snapshot).unwrap();
3982 assert!(!snapshot.process_alive);
3983 assert_eq!(snapshot.pid, None);
3984 assert_eq!(snapshot.spawned_at_ms, None);
3985 assert_eq!(snapshot.spawned_from, None);
3986 assert_eq!(snapshot.spawned_file_identity, None);
3987 }
3988
3989 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3990 async fn failed_enable_spawn_clears_preexisting_current_process_facts() {
3991 let supervisor = Supervisor::default();
3992 let mut runtime = supervisor.runtime_config();
3993 runtime.test_seed_stale_facts_before_enable_spawn = true;
3994 let snapshot = stale_process_snapshot(ModuleState::Disabled, false);
3995 let mut child = None;
3996 let spec = ModuleSpec {
3997 module_id: "failed-enable-clears-facts".to_string(),
3998 program: PathBuf::from("/definitely/missing/failed-enable-module"),
3999 args: Vec::new(),
4000 env: Vec::new(),
4001 reserved: false,
4002 reserved_prefixes: Vec::new(),
4003 protocol: ModuleProtocol::Subc,
4004 overlap: Default::default(),
4005 };
4006
4007 let result = set_child_enabled(
4008 &spec,
4009 &runtime,
4010 &supervisor.registry,
4011 &supervisor.process_liveness,
4012 &snapshot,
4013 &mut child,
4014 true,
4015 )
4016 .await;
4017
4018 assert!(matches!(result, Err(SuperviseError::Spawn { .. })));
4019 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4020 assert_snapshot_process_facts_cleared(&snapshot);
4021 }
4022
4023 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4024 async fn failed_reload_spawn_clears_current_process_facts() {
4025 let supervisor = Supervisor::default();
4026 let mut runtime = supervisor.runtime_config();
4027 runtime.restart_policy = RestartPolicy::new(0, Duration::ZERO);
4028 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4029 let mut child = None;
4030 let spec = ModuleSpec {
4031 module_id: "failed-reload-clears-facts".to_string(),
4032 program: PathBuf::from("/unused/failed-reload-module"),
4033 args: Vec::new(),
4034 env: Vec::new(),
4035 reserved: false,
4036 reserved_prefixes: Vec::new(),
4037 protocol: ModuleProtocol::Subc,
4038 overlap: Default::default(),
4039 };
4040
4041 let result = handle_reload_spawn_failure(
4042 &spec,
4043 &runtime,
4044 &supervisor.process_liveness,
4045 &snapshot,
4046 &mut child,
4047 "forced reload spawn failure".to_string(),
4048 )
4049 .await;
4050
4051 assert!(matches!(result, Err(SuperviseError::ReloadFailed { .. })));
4052 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4053 assert_snapshot_process_facts_cleared(&snapshot);
4054 }
4055
4056 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4057 async fn dropping_a_module_with_an_active_monitor_clears_current_process_facts() {
4058 let supervisor = Supervisor::default();
4059 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4060 let module = supervisor.supervised_module(
4061 ModuleSpec {
4062 module_id: "drop-clears-facts".to_string(),
4063 program: PathBuf::from("/unused/drop-module"),
4064 args: Vec::new(),
4065 env: Vec::new(),
4066 reserved: false,
4067 reserved_prefixes: Vec::new(),
4068 protocol: ModuleProtocol::Subc,
4069 overlap: Default::default(),
4070 },
4071 supervisor.runtime_config(),
4072 Arc::clone(&snapshot),
4073 None,
4074 );
4075 assert!(!module
4076 .inner
4077 .monitor
4078 .lock()
4079 .unwrap()
4080 .as_ref()
4081 .unwrap()
4082 .is_finished());
4083
4084 drop(module);
4085
4086 assert_eq!(
4087 lock_snapshot(&snapshot).unwrap().state,
4088 ModuleState::Stopped
4089 );
4090 assert_snapshot_process_facts_cleared(&snapshot);
4091 }
4092
4093 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4094 async fn configuration_update_does_not_replace_captured_running_process_facts() {
4095 let supervisor = Supervisor::default();
4096 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4097 let initial = ModuleSpec {
4098 module_id: "rescan-preserves-spawn-facts".to_string(),
4099 program: PathBuf::from("/spawned/module"),
4100 args: Vec::new(),
4101 env: Vec::new(),
4102 reserved: false,
4103 reserved_prefixes: Vec::new(),
4104 protocol: ModuleProtocol::Subc,
4105 overlap: Default::default(),
4106 };
4107 let module = supervisor.supervised_module(
4108 initial.clone(),
4109 supervisor.runtime_config(),
4110 snapshot,
4111 None,
4112 );
4113 let before = module.status().unwrap();
4114 let mut replacement = initial;
4115 replacement.program = PathBuf::from("/rescanned/replacement-module");
4116
4117 module
4118 .update_configuration(replacement, HealthConfig::default(), None)
4119 .await
4120 .unwrap();
4121
4122 let after = module.status().unwrap();
4123 assert_eq!(after.pid, before.pid);
4124 assert_eq!(after.spawned_at_ms, before.spawned_at_ms);
4125 assert_eq!(after.spawned_from, before.spawned_from);
4126 drop(module);
4127 }
4128}
4129
4130fn unix_ms_now() -> u64 {
4131 SystemTime::now()
4132 .duration_since(UNIX_EPOCH)
4133 .map(|duration| duration.as_millis().min(u128::from(u64::MAX)) as u64)
4134 .unwrap_or(0)
4135}
4136
4137async fn supervise_loop(
4138 mut spec: ModuleSpec,
4139 mut runtime: SupervisorRuntimeConfig,
4140 registry: Arc<Registry>,
4141 process_liveness: Arc<SupervisorProcessLiveness>,
4142 snapshot: SharedSnapshot,
4143 mut child: Option<SupervisedChild>,
4144 mut commands: mpsc::Receiver<SupervisorCommand>,
4145) {
4146 let mut health_probe = HealthProbeRuntime::default();
4147 let mut pending_respawn: Option<Instant> = None;
4151 let mut requeued: VecDeque<SupervisorCommand> = VecDeque::new();
4154 loop {
4155 if let Some(command) = requeued.pop_front() {
4156 if !handle_supervisor_command(
4157 command,
4158 &mut spec,
4159 &mut runtime,
4160 ®istry,
4161 &process_liveness,
4162 &snapshot,
4163 &mut child,
4164 &mut commands,
4165 &mut requeued,
4166 )
4167 .await
4168 {
4169 return;
4170 }
4171 if child.is_some() || !respawn_still_pending(&snapshot) {
4172 pending_respawn = None;
4173 }
4174 continue;
4175 }
4176 if child.is_some() {
4177 health_probe.refresh_registration(&spec, &runtime, ®istry, &snapshot);
4178 let probe_sleep = sleep(health_probe.wake_after());
4179 tokio::pin!(probe_sleep);
4180 let active_child = child.as_mut().expect("child checked above");
4181 tokio::select! {
4182 wait_result = active_child.wait() => {
4183 let exit_report = match wait_result {
4192 Ok(status) => classify_reaped_child_exit(&snapshot, active_child, &status),
4193 Err(err) => {
4194 active_child.drain_stderr(&spec.module_id).await;
4195 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4196 record_wait_error_terminal(
4202 &spec.module_id,
4203 &runtime.terminal_ring,
4204 &runtime.spawn_events,
4205 );
4206 untrack_if_registration_released(
4207 &process_liveness,
4208 ®istry,
4209 &spec.module_id,
4210 &snapshot,
4211 );
4212 error!(module_id = %spec.module_id, error = %err, "failed to wait for supervised module");
4213 child = None;
4214 continue;
4215 }
4216 };
4217 active_child.drain_stderr(&spec.module_id).await;
4218
4219 match on_child_exit(
4220 &spec,
4221 runtime.restart_policy,
4222 ®istry,
4223 &snapshot,
4224 &runtime.terminal_ring,
4225 &runtime.spawn_events,
4226 exit_report,
4227 ).await {
4228 NextAction::Stop { registration_released } => {
4229 if registration_released {
4230 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4231 }
4232 child = None;
4233 }
4234 NextAction::Restart { schedule } => {
4235 let delay = schedule.map_or(
4236 runtime.restart_policy.delay_for_restart(0),
4237 |schedule| schedule.delay,
4238 );
4239 if let Some(schedule) = schedule {
4240 log_crash_respawn(&spec.module_id, schedule);
4241 }
4242 child = None;
4250 pending_respawn = Some(Instant::now() + delay);
4251 }
4252 }
4253 }
4254 command = commands.recv() => {
4255 let Some(command) = command else {
4256 return;
4257 };
4258 if !handle_supervisor_command(
4259 command,
4260 &mut spec,
4261 &mut runtime,
4262 ®istry,
4263 &process_liveness,
4264 &snapshot,
4265 &mut child,
4266 &mut commands,
4267 &mut requeued,
4268 ).await {
4269 return;
4270 }
4271 }
4272 _ = &mut probe_sleep => {
4273 if health_probe.due() {
4274 run_health_probe_cycle(
4275 &spec,
4276 &runtime,
4277 ®istry,
4278 &process_liveness,
4279 &snapshot,
4280 &mut child,
4281 ).await;
4282 if child.is_some() {
4283 health_probe.schedule_next(&spec, runtime.health.cadence);
4284 }
4285 }
4286 }
4287 }
4288 } else if let Some(deadline) = pending_respawn {
4289 tokio::select! {
4290 _ = sleep_until(deadline) => {
4291 pending_respawn = None;
4292 if !respawn_still_pending(&snapshot) {
4296 continue;
4297 }
4298 if let Err(err) = wait_for_registration_release(
4299 ®istry,
4300 &spec.module_id,
4301 REGISTRY_RELEASE_TIMEOUT,
4302 ).await {
4303 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4304 error!(module_id = %spec.module_id, error = %err, "registration did not release before restart");
4305 continue;
4306 }
4307
4308 match spawn_and_mark_running(&spec, &runtime, &snapshot) {
4309 Ok(next_child) => {
4310 child = Some(next_child);
4311 debug!(module_id = %spec.module_id, "supervised module restarted after crash");
4312 }
4313 Err(err) => {
4314 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4315 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4316 error!(module_id = %spec.module_id, error = %err, "failed to restart supervised module");
4317 }
4318 }
4319 }
4320 command = commands.recv() => {
4321 let Some(command) = command else {
4322 return;
4323 };
4324 if !handle_supervisor_command(
4325 command,
4326 &mut spec,
4327 &mut runtime,
4328 ®istry,
4329 &process_liveness,
4330 &snapshot,
4331 &mut child,
4332 &mut commands,
4333 &mut requeued,
4334 ).await {
4335 return;
4336 }
4337 if child.is_some() || !respawn_still_pending(&snapshot) {
4342 pending_respawn = None;
4343 }
4344 }
4345 }
4346 } else {
4347 let Some(command) = commands.recv().await else {
4348 return;
4349 };
4350 if !handle_supervisor_command(
4351 command,
4352 &mut spec,
4353 &mut runtime,
4354 ®istry,
4355 &process_liveness,
4356 &snapshot,
4357 &mut child,
4358 &mut commands,
4359 &mut requeued,
4360 )
4361 .await
4362 {
4363 return;
4364 }
4365 }
4366 }
4367}
4368
4369fn log_crash_respawn(module_id: &str, schedule: CrashRestartSchedule) {
4370 info!(
4371 module_id,
4372 restart_in_window = schedule.restart_in_window,
4373 delay_ms = schedule.delay.as_millis() as u64,
4374 "respawning after crash"
4375 );
4376}
4377
4378fn respawn_still_pending(snapshot: &SharedSnapshot) -> bool {
4384 matches!(
4385 lock_snapshot(snapshot),
4386 Ok(state) if state.enabled && state.state == ModuleState::Restarting
4387 )
4388}
4389
4390enum NextAction {
4391 Stop {
4392 registration_released: bool,
4393 },
4394 Restart {
4395 schedule: Option<CrashRestartSchedule>,
4396 },
4397}
4398
4399#[allow(clippy::too_many_arguments)]
4400async fn handle_supervisor_command(
4401 command: SupervisorCommand,
4402 spec: &mut ModuleSpec,
4403 runtime: &mut SupervisorRuntimeConfig,
4404 registry: &Registry,
4405 process_liveness: &SupervisorProcessLiveness,
4406 snapshot: &SharedSnapshot,
4407 child: &mut Option<SupervisedChild>,
4408 commands: &mut mpsc::Receiver<SupervisorCommand>,
4409 requeued: &mut VecDeque<SupervisorCommand>,
4410) -> bool {
4411 match command {
4412 SupervisorCommand::Drain { reply } => {
4413 let result = drain_optional_child(
4414 &spec.module_id,
4415 spec.protocol,
4416 registry,
4417 snapshot,
4418 &runtime.terminal_ring,
4419 &runtime.spawn_events,
4420 child,
4421 runtime.drain_timeout,
4422 ModuleState::Stopped,
4423 None,
4424 )
4425 .await;
4426 let registration_released = result.is_ok();
4427 let _ = reply.send(result);
4428 if registration_released {
4429 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4430 }
4431 false
4432 }
4433 SupervisorCommand::Retire { reply } => {
4434 let result = async {
4435 begin_forwarding_drain_if_configured(
4436 spec,
4437 runtime,
4438 registry,
4439 snapshot,
4440 None,
4441 RouteCloseReason::Disable,
4442 )
4443 .await?;
4444 drain_optional_child(
4445 &spec.module_id,
4446 spec.protocol,
4447 registry,
4448 snapshot,
4449 &runtime.terminal_ring,
4450 &runtime.spawn_events,
4451 child,
4452 runtime.drain_timeout,
4453 ModuleState::Stopped,
4454 None,
4455 )
4456 .await
4457 }
4458 .await;
4459 let registration_released = result.is_ok();
4460 let _ = reply.send(result);
4461 if registration_released {
4462 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4463 }
4464 false
4465 }
4466 SupervisorCommand::Restart {
4467 drain_timeout_ms,
4468 reply,
4469 } => {
4470 let validation = match lock_snapshot(snapshot) {
4482 Ok(state) if !state.enabled => Err(SuperviseError::Disabled {
4483 module_id: spec.module_id.clone(),
4484 }),
4485 Ok(_) => Ok(()),
4486 Err(err) => Err(err),
4487 };
4488 let initiated = validation.is_ok();
4489 let _ = reply.send(validation);
4490 if initiated {
4491 let drain_timeout = drain_timeout_ms
4494 .map(Duration::from_millis)
4495 .unwrap_or(runtime.drain_timeout);
4496 if let Err(err) = restart_child(
4497 spec,
4498 runtime,
4499 registry,
4500 process_liveness,
4501 snapshot,
4502 child,
4503 drain_timeout,
4504 )
4505 .await
4506 {
4507 warn!(
4508 module_id = %spec.module_id,
4509 error = %err,
4510 "operator restart failed after initiation ack; module state carries the outcome"
4511 );
4512 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4513 state.state = ModuleState::Failed;
4514 clear_current_process_facts(state);
4515 });
4516 }
4517 }
4518 true
4519 }
4520 SupervisorCommand::Reload { reply } => {
4521 let result =
4522 reload_child(spec, runtime, registry, process_liveness, snapshot, child).await;
4523 let _ = reply.send(result);
4524 true
4525 }
4526 SupervisorCommand::SetEnabled { enabled, reply } => {
4527 let result = set_child_enabled(
4528 spec,
4529 runtime,
4530 registry,
4531 process_liveness,
4532 snapshot,
4533 child,
4534 enabled,
4535 )
4536 .await;
4537 let _ = reply.send(result);
4538 true
4539 }
4540 SupervisorCommand::UpdateConfiguration {
4541 spec: next_spec,
4542 health,
4543 drain_timeout_ms,
4544 reply,
4545 } => {
4546 if let Some(handle) = &runtime.supervisor_handle {
4547 handle.apply_identity_configuration(&next_spec);
4548 }
4549 *spec = next_spec;
4550 runtime.health = health;
4551 runtime.drain_timeout = drain_timeout_ms
4552 .map(Duration::from_millis)
4553 .unwrap_or(runtime.default_drain_timeout);
4554 *runtime
4555 .effective_drain_timeout
4556 .lock()
4557 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
4558 let _ = reply.send(());
4559 true
4560 }
4561 SupervisorCommand::Swap {
4562 ready_timeout,
4563 reply,
4564 } => {
4565 let end = swap::run_swap(
4566 spec,
4567 runtime,
4568 registry,
4569 process_liveness,
4570 snapshot,
4571 child,
4572 commands,
4573 ready_timeout.unwrap_or(DEFAULT_SWAP_READY_TIMEOUT),
4574 reply,
4575 )
4576 .await;
4577 requeued.extend(end.requeue);
4578 true
4579 }
4580 }
4581}
4582
4583async fn restart_child(
4584 spec: &ModuleSpec,
4585 runtime: &SupervisorRuntimeConfig,
4586 registry: &Registry,
4587 process_liveness: &SupervisorProcessLiveness,
4588 snapshot: &SharedSnapshot,
4589 child: &mut Option<SupervisedChild>,
4590 drain_timeout: Duration,
4591) -> Result<(), SuperviseError> {
4592 if !lock_snapshot(snapshot)?.enabled {
4594 return Err(SuperviseError::Disabled {
4595 module_id: spec.module_id.clone(),
4596 });
4597 }
4598 begin_forwarding_drain_with_timeout(
4599 spec,
4600 runtime,
4601 registry,
4602 snapshot,
4603 None,
4604 RouteCloseReason::Restart,
4605 drain_timeout,
4606 )
4607 .await?;
4608
4609 if child.is_some() {
4610 drain_optional_child(
4611 &spec.module_id,
4612 spec.protocol,
4613 registry,
4614 snapshot,
4615 &runtime.terminal_ring,
4616 &runtime.spawn_events,
4617 child,
4618 drain_timeout,
4619 ModuleState::Restarting,
4620 Some(true),
4621 )
4622 .await?;
4623 } else {
4624 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4625 state.enabled = true;
4626 state.state = ModuleState::Restarting;
4627 clear_current_process_facts(state);
4628 })?;
4629 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4630 }
4631
4632 reset_restart_count(snapshot, &spec.module_id)?;
4633 sleep(runtime.restart_policy.backoff).await;
4634 if !respawn_still_pending(snapshot) {
4637 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4638 return Ok(());
4639 }
4640 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4641 match spawn_and_mark_running(spec, runtime, snapshot) {
4647 Ok(next_child) => {
4648 *child = Some(next_child);
4649 debug!(module_id = %spec.module_id, "supervised module restarted by operator request");
4650 Ok(())
4651 }
4652 Err(err) => {
4653 fail_snapshot(snapshot, Some(&spec.module_id), None);
4654 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4655 *child = None;
4656 Err(err)
4657 }
4658 }
4659}
4660
4661async fn reload_child(
4662 spec: &ModuleSpec,
4663 runtime: &SupervisorRuntimeConfig,
4664 registry: &Registry,
4665 process_liveness: &SupervisorProcessLiveness,
4666 snapshot: &SharedSnapshot,
4667 child: &mut Option<SupervisedChild>,
4668) -> Result<(), SuperviseError> {
4669 if !lock_snapshot(snapshot)?.enabled {
4671 return Err(SuperviseError::Disabled {
4672 module_id: spec.module_id.clone(),
4673 });
4674 }
4675 begin_forwarding_drain(
4676 spec,
4677 runtime,
4678 registry,
4679 snapshot,
4680 Some(true),
4681 RouteCloseReason::Reload,
4682 )
4683 .await?;
4684
4685 if child.is_some() {
4686 drain_optional_child(
4687 &spec.module_id,
4688 spec.protocol,
4689 registry,
4690 snapshot,
4691 &runtime.terminal_ring,
4692 &runtime.spawn_events,
4693 child,
4694 runtime.drain_timeout,
4695 ModuleState::Restarting,
4696 Some(true),
4697 )
4698 .await?;
4699 } else {
4700 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4701 state.enabled = true;
4702 state.state = ModuleState::Restarting;
4703 clear_current_process_facts(state);
4704 })?;
4705 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4706 }
4707
4708 reset_restart_count(snapshot, &spec.module_id)?;
4709 sleep(runtime.restart_policy.backoff).await;
4710 if !respawn_still_pending(snapshot) {
4713 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4714 return Ok(());
4715 }
4716 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4717 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
4718 Ok(next_child) => next_child,
4719 Err(err) => {
4720 return handle_reload_spawn_failure(
4721 spec,
4722 runtime,
4723 process_liveness,
4724 snapshot,
4725 child,
4726 format!("new child failed to spawn: {err}"),
4727 )
4728 .await;
4729 }
4730 };
4731 *child = Some(next_child);
4732
4733 let wait_outcome = {
4734 let active_child = child.as_mut().expect("new reload child was just stored");
4735 wait_for_registration_after_reload(
4736 registry,
4737 &spec.module_id,
4738 snapshot,
4739 active_child,
4740 REGISTRY_RELEASE_TIMEOUT,
4741 )
4742 .await?
4743 };
4744
4745 match wait_outcome {
4746 RegistrationWaitOutcome::Registered => {
4747 debug!(module_id = %spec.module_id, "supervised module reloaded and registered");
4748 Ok(())
4749 }
4750 RegistrationWaitOutcome::Exited(exit_report) => {
4751 if let Some(active_child) = child.as_mut() {
4752 active_child.drain_stderr(&spec.module_id).await;
4753 }
4754 *child = None;
4755 handle_reload_child_registration_failure(
4756 spec,
4757 runtime,
4758 registry,
4759 process_liveness,
4760 snapshot,
4761 child,
4762 ReloadRegistrationFailure {
4763 exit_report: registration_failure_exit_report(exit_report),
4764 reason: "new child exited before registering".to_string(),
4765 },
4766 )
4767 .await
4768 }
4769 RegistrationWaitOutcome::TimedOut => {
4770 let mut timed_out_child = child
4771 .take()
4772 .expect("timed-out reload child is still running");
4773 timed_out_child
4774 .start_kill()
4775 .map_err(|source| SuperviseError::Kill {
4776 module_id: spec.module_id.clone(),
4777 source,
4778 })?;
4779 let status = timed_out_child
4780 .wait()
4781 .await
4782 .map_err(|source| SuperviseError::Wait {
4783 module_id: spec.module_id.clone(),
4784 source,
4785 })?;
4786 timed_out_child.drain_stderr(&spec.module_id).await;
4787 handle_reload_child_registration_failure(
4788 spec,
4789 runtime,
4790 registry,
4791 process_liveness,
4792 snapshot,
4793 child,
4794 ReloadRegistrationFailure {
4795 exit_report: registration_failure_exit_report(classify_reaped_child_exit(
4796 snapshot,
4797 &timed_out_child,
4798 &status,
4799 )),
4800 reason: format!(
4801 "new child did not register within {:?}",
4802 REGISTRY_RELEASE_TIMEOUT
4803 ),
4804 },
4805 )
4806 .await
4807 }
4808 }
4809}
4810
4811async fn set_child_enabled(
4812 spec: &ModuleSpec,
4813 runtime: &SupervisorRuntimeConfig,
4814 registry: &Registry,
4815 process_liveness: &SupervisorProcessLiveness,
4816 snapshot: &SharedSnapshot,
4817 child: &mut Option<SupervisedChild>,
4818 enabled: bool,
4819) -> Result<bool, SuperviseError> {
4820 let (current_enabled, current_state) = {
4821 let state = lock_snapshot(snapshot)?;
4822 (state.enabled, state.state)
4823 };
4824 let revive_terminal = enabled
4832 && current_enabled
4833 && child.is_none()
4834 && matches!(current_state, ModuleState::Failed | ModuleState::Stopped);
4835 if current_enabled == enabled && !revive_terminal {
4836 return Ok(false);
4837 }
4838
4839 if enabled {
4840 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4841 state.enabled = true;
4842 state.state = ModuleState::Starting;
4843 clear_current_process_facts(state);
4844 })?;
4845 #[cfg(test)]
4846 if runtime.test_seed_stale_facts_before_enable_spawn {
4847 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4848 state.process_alive = true;
4849 state.pid = Some(41);
4850 state.spawned_at_ms = Some(42);
4851 state.spawned_from = Some(PathBuf::from("/spawned/module"));
4852 state.spawned_file_identity = Some(SpawnedFileIdentity {
4853 device: 43,
4854 inode: 44,
4855 });
4856 })?;
4857 }
4858 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4859 reset_restart_count(snapshot, &spec.module_id)?;
4860 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4861 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
4862 Ok(next_child) => next_child,
4863 Err(err) => {
4864 if let Err(state_err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4865 state.state = ModuleState::Failed;
4866 clear_current_process_facts(state);
4867 }) {
4868 error!(module_id = %spec.module_id, error = %state_err, "failed to record enable spawn failure");
4869 }
4870 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4871 return Err(err);
4872 }
4873 };
4874 *child = Some(next_child);
4875 debug!(module_id = %spec.module_id, "supervised module enabled");
4876 Ok(true)
4877 } else {
4878 begin_forwarding_drain_if_configured(
4879 spec,
4880 runtime,
4881 registry,
4882 snapshot,
4883 Some(false),
4884 RouteCloseReason::Disable,
4885 )
4886 .await?;
4887 drain_optional_child(
4888 &spec.module_id,
4889 spec.protocol,
4890 registry,
4891 snapshot,
4892 &runtime.terminal_ring,
4893 &runtime.spawn_events,
4894 child,
4895 runtime.drain_timeout,
4896 ModuleState::Disabled,
4897 Some(false),
4898 )
4899 .await?;
4900 debug!(module_id = %spec.module_id, "supervised module disabled");
4901 Ok(true)
4902 }
4903}
4904
4905async fn on_child_exit(
4906 spec: &ModuleSpec,
4907 policy: RestartPolicy,
4908 registry: &Registry,
4909 snapshot: &SharedSnapshot,
4910 terminal_ring: &Arc<Mutex<TerminalRing>>,
4911 spawn_events: &SpawnEventFeed,
4912 exit_report: ExitReport,
4913) -> NextAction {
4914 match exit_report.kind {
4915 ExitKind::Clean => {
4916 info!(
4917 module_id = %spec.module_id,
4918 exit_code = ?exit_report.code,
4919 exit_signal = ?exit_report.signal,
4920 "supervised module exited cleanly"
4921 );
4922 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4923 state.state = ModuleState::Stopped;
4924 clear_current_process_facts(state);
4925 state.last_exit = Some(exit_report.clone());
4926 }) {
4927 error!(module_id = %spec.module_id, error = %err, "failed to record clean module exit");
4928 }
4929 record_terminal(
4930 &spec.module_id,
4931 terminal_ring,
4932 spawn_events,
4933 &exit_report,
4934 TerminalDisposition::Stopped,
4935 );
4936 let registration_released = match wait_for_registration_release(
4937 registry,
4938 &spec.module_id,
4939 REGISTRY_RELEASE_TIMEOUT,
4940 )
4941 .await
4942 {
4943 Ok(()) => true,
4944 Err(err) => {
4945 warn!(module_id = %spec.module_id, error = %err, "registration still active after clean exit");
4946 false
4947 }
4948 };
4949 NextAction::Stop {
4950 registration_released,
4951 }
4952 }
4953 ExitKind::Crash => {
4954 warn!(
4955 module_id = %spec.module_id,
4956 exit_code = ?exit_report.code,
4957 exit_signal = ?exit_report.signal,
4958 "supervised module exited abnormally (crash)"
4959 );
4960 let mut restart_schedule = None;
4961 let mut disposition = TerminalDisposition::Disabled;
4962 let mut disposition_detail = None;
4966 let now = Instant::now();
4967 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4968 clear_current_process_facts(state);
4969 state.last_exit = Some(exit_report.clone());
4970 if state.enabled {
4971 if let Some(schedule) = state.next_crash_restart(&policy, now) {
4972 state.state = ModuleState::Restarting;
4973 restart_schedule = Some(schedule);
4974 disposition = TerminalDisposition::Restarting;
4975 } else {
4976 state.state = ModuleState::Failed;
4977 disposition = TerminalDisposition::Failed;
4978 disposition_detail = Some(policy.budget_exhausted_detail());
4979 }
4980 } else {
4981 state.state = ModuleState::Disabled;
4982 disposition = TerminalDisposition::Disabled;
4983 }
4984 }) {
4985 error!(module_id = %spec.module_id, error = %err, "failed to record crashed module exit");
4986 return NextAction::Stop {
4987 registration_released: false,
4988 };
4989 }
4990 if disposition_detail.is_some() {
4991 error!(
4996 module_id = %spec.module_id,
4997 max_restarts = policy.max_restarts,
4998 window_secs = policy.window.as_secs(),
4999 "module stopped: {}",
5000 policy.budget_exhausted_detail()
5001 );
5002 }
5003 record_terminal_with_detail(
5004 &spec.module_id,
5005 terminal_ring,
5006 spawn_events,
5007 &exit_report,
5008 disposition,
5009 disposition_detail,
5010 );
5011
5012 if let Some(schedule) = restart_schedule {
5013 NextAction::Restart {
5014 schedule: Some(schedule),
5015 }
5016 } else {
5017 let registration_released = match wait_for_registration_release(
5018 registry,
5019 &spec.module_id,
5020 REGISTRY_RELEASE_TIMEOUT,
5021 )
5022 .await
5023 {
5024 Ok(()) => true,
5025 Err(err) => {
5026 warn!(module_id = %spec.module_id, error = %err, "registration still active after failed module");
5027 false
5028 }
5029 };
5030 NextAction::Stop {
5031 registration_released,
5032 }
5033 }
5034 }
5035 ExitKind::DeliberateSeverance => {
5036 warn!(
5037 module_id = %spec.module_id,
5038 exit_code = ?exit_report.code,
5039 exit_signal = ?exit_report.signal,
5040 "supervised module exited after deliberate connection severance"
5041 );
5042 let mut should_restart = false;
5043 let mut disposition = TerminalDisposition::Disabled;
5044 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5045 clear_current_process_facts(state);
5046 state.last_exit = Some(exit_report.clone());
5047 state.lifetime_restarts += 1;
5048 if state.enabled {
5049 state.state = ModuleState::Restarting;
5050 should_restart = true;
5051 disposition = TerminalDisposition::Restarting;
5052 } else {
5053 state.state = ModuleState::Disabled;
5054 }
5055 }) {
5056 error!(module_id = %spec.module_id, error = %err, "failed to record deliberately severed module exit");
5057 return NextAction::Stop {
5058 registration_released: false,
5059 };
5060 }
5061 record_terminal(
5062 &spec.module_id,
5063 terminal_ring,
5064 spawn_events,
5065 &exit_report,
5066 disposition,
5067 );
5068
5069 if should_restart {
5070 NextAction::Restart { schedule: None }
5071 } else {
5072 let registration_released = match wait_for_registration_release(
5073 registry,
5074 &spec.module_id,
5075 REGISTRY_RELEASE_TIMEOUT,
5076 )
5077 .await
5078 {
5079 Ok(()) => true,
5080 Err(err) => {
5081 warn!(module_id = %spec.module_id, error = %err, "registration still active after deliberately severed module exit");
5082 false
5083 }
5084 };
5085 NextAction::Stop {
5086 registration_released,
5087 }
5088 }
5089 }
5090 }
5091}
5092
5093fn record_wait_error_terminal(
5094 module_id: &str,
5095 terminal_ring: &Arc<Mutex<TerminalRing>>,
5096 spawn_events: &SpawnEventFeed,
5097) {
5098 record_terminal(
5099 module_id,
5100 terminal_ring,
5101 spawn_events,
5102 &wait_error_exit_report(),
5103 TerminalDisposition::Failed,
5104 );
5105}
5106
5107fn record_terminal(
5108 module_id: &str,
5109 terminal_ring: &Arc<Mutex<TerminalRing>>,
5110 spawn_events: &SpawnEventFeed,
5111 exit_report: &ExitReport,
5112 disposition: TerminalDisposition,
5113) {
5114 record_terminal_with_detail(
5115 module_id,
5116 terminal_ring,
5117 spawn_events,
5118 exit_report,
5119 disposition,
5120 None,
5121 );
5122}
5123
5124fn record_terminal_with_detail(
5125 module_id: &str,
5126 terminal_ring: &Arc<Mutex<TerminalRing>>,
5127 spawn_events: &SpawnEventFeed,
5128 exit_report: &ExitReport,
5129 disposition: TerminalDisposition,
5130 disposition_detail: Option<String>,
5131) {
5132 spawn_events.emit_exited(module_id, exit_report.code, exit_report.signal);
5133 let mut ring = terminal_ring
5134 .lock()
5135 .unwrap_or_else(|poisoned| poisoned.into_inner());
5136 let record = TerminalRecord {
5137 exit_code: exit_report.code,
5138 exit_signal: exit_report.signal,
5139 at_ms: exit_report.at_ms,
5140 disposition,
5141 exit_kind: exit_report.kind.into(),
5142 disposition_detail,
5143 };
5144 ring.append_journal(module_id, &record);
5145 ring.push(record);
5146}
5147
5148fn untrack_if_registration_released(
5149 process_liveness: &SupervisorProcessLiveness,
5150 registry: &Registry,
5151 module_id: &str,
5152 snapshot: &SharedSnapshot,
5153) {
5154 match registry.get_module(module_id) {
5155 Ok(None) => process_liveness.untrack_if_current(module_id, snapshot),
5156 Ok(Some(_)) => {}
5157 Err(err) => {
5158 warn!(module_id, error = %err, "could not determine whether supervisor liveness can be untracked");
5159 }
5160 }
5161}
5162
5163#[cfg(test)]
5177fn apply_wire_spawn_args(
5178 command: &mut Command,
5179 spec: &ModuleSpec,
5180 connection_file_path: Option<&std::path::Path>,
5181 handle: Option<&SupervisorHandle>,
5182) -> Result<(), SuperviseError> {
5183 apply_wire_spawn_args_for_role(
5184 command,
5185 spec,
5186 connection_file_path,
5187 handle,
5188 SpawnRole::Plain,
5189 )
5190}
5191
5192fn apply_wire_spawn_args_for_role(
5201 command: &mut Command,
5202 spec: &ModuleSpec,
5203 connection_file_path: Option<&std::path::Path>,
5204 handle: Option<&SupervisorHandle>,
5205 role: SpawnRole,
5206) -> Result<(), SuperviseError> {
5207 command.env(SUBC_MODULE_ID_ENV, &spec.module_id);
5208 if spec.protocol == ModuleProtocol::None {
5209 return Ok(());
5210 }
5211 if let Some(connection_file_path) = connection_file_path {
5212 command.arg(SUBC_ARG).arg(connection_file_path);
5213 }
5214
5215 let nonce = generate_launch_nonce()?;
5219 if let Some(handle) = handle {
5220 match role {
5221 SpawnRole::Plain => {
5222 handle.set_spawn_nonce(&spec.module_id, nonce.clone());
5223 if spec.reserved {
5224 handle.set_reserved_nonce(&spec.module_id, nonce.clone());
5225 }
5226 }
5227 SpawnRole::SwapCandidate => handle.open_swap(&spec.module_id, nonce.clone()),
5228 }
5229 }
5230 command.env(SUBC_LAUNCH_NONCE_ENV, nonce);
5231 Ok(())
5232}
5233
5234fn apply_child_env(command: &mut Command, spec: &ModuleSpec) {
5235 command.env_remove(CK_LOG_ENV);
5236 command.env_remove(SUBC_SPAWN_ROLE_ENV);
5243 for (key, value) in &spec.env {
5244 if matches!(
5248 key.as_str(),
5249 CAPTURE_MAX_FILE_MB_ENV | CAPTURE_KEEP_ENV | CAPTURE_MAX_AGE_DAYS_ENV
5250 ) || key == SUBC_SPAWN_ROLE_ENV
5251 {
5252 continue;
5253 }
5254 command.env(key, value);
5255 }
5256}
5257
5258#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5261enum SpawnRole {
5262 Plain,
5263 SwapCandidate,
5264}
5265
5266fn apply_spawn_role(command: &mut Command, role: SpawnRole) {
5269 if role == SpawnRole::SwapCandidate {
5270 command.env(SUBC_SPAWN_ROLE_ENV, SPAWN_ROLE_SWAP_CANDIDATE);
5271 }
5272}
5273
5274fn spawn_child(
5275 spec: &ModuleSpec,
5276 connection_file_path: Option<&std::path::Path>,
5277 handle: Option<&SupervisorHandle>,
5278 ring: &Arc<Mutex<StderrRing>>,
5279 capture_logs_dir: Option<&std::path::Path>,
5280 roster: &ChildRoster,
5281 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5282) -> Result<SupervisedChild, SuperviseError> {
5283 spawn_child_in_slot(
5284 spec,
5285 connection_file_path,
5286 handle,
5287 ring,
5288 capture_logs_dir,
5289 roster,
5290 #[cfg(target_os = "linux")]
5291 cgroup_placement,
5292 SpawnRole::Plain,
5293 false,
5294 )
5295}
5296
5297#[allow(clippy::too_many_arguments)]
5310fn spawn_child_in_slot(
5311 spec: &ModuleSpec,
5312 connection_file_path: Option<&std::path::Path>,
5313 handle: Option<&SupervisorHandle>,
5314 ring: &Arc<Mutex<StderrRing>>,
5315 capture_logs_dir: Option<&std::path::Path>,
5316 roster: &ChildRoster,
5317 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5318 role: SpawnRole,
5319 alternate_slot: bool,
5320) -> Result<SupervisedChild, SuperviseError> {
5321 if roster.is_closed() {
5322 return Err(SuperviseError::Spawn {
5323 program: spec.program.clone(),
5324 source: io::Error::other("the daemon is shutting down; not starting a new process"),
5325 cgroup_path: None,
5326 });
5327 }
5328 #[cfg(target_os = "linux")]
5329 let cgroup_name = swap::cgroup_name(&spec.module_id, alternate_slot);
5330 #[cfg(not(target_os = "linux"))]
5331 let _ = alternate_slot;
5332 let mut command = Command::new(&spec.program);
5333 command.args(&spec.args);
5334 apply_child_env(&mut command, spec);
5364 apply_spawn_role(&mut command, role);
5365 apply_wire_spawn_args_for_role(&mut command, spec, connection_file_path, handle, role)?;
5366
5367 #[cfg(target_os = "linux")]
5368 let cgroup_path = cgroup_placement
5369 .map(|placement| placement.module_path(&cgroup_name))
5370 .transpose()
5371 .map_err(|source| SuperviseError::Cgroup {
5372 module_id: spec.module_id.clone(),
5373 source,
5374 })?;
5375 #[cfg(not(target_os = "linux"))]
5376 let cgroup_path: Option<PathBuf> = None;
5377 #[cfg(target_os = "linux")]
5378 if let Some(path) = &cgroup_path {
5379 if let Err(error) = apply_cgroup_placement(&mut command, spec, path) {
5380 if let Some(placement) = cgroup_placement {
5381 remove_module_cgroup(placement, &cgroup_name);
5382 }
5383 return Err(error);
5384 }
5385 }
5386
5387 let output_sink = if let Some(logs_dir) = capture_logs_dir {
5388 let path = logs_dir.join(format!("{}.stderr.log", spec.module_id));
5389 match ChildOutputSink::open(&path, capture_retention(spec)) {
5390 Ok(sink) => sink,
5391 Err(error) => {
5392 warn!(
5393 module_id = %spec.module_id,
5394 path = %path.display(),
5395 error = %error,
5396 "could not open child output capture file; forwarding to stderr"
5397 );
5398 ChildOutputSink::Stderr
5399 }
5400 }
5401 } else {
5402 ChildOutputSink::Stderr
5403 };
5404
5405 command.stdout(Stdio::piped());
5406 command.stderr(Stdio::piped());
5407 command.kill_on_drop(true);
5408 #[cfg(unix)]
5425 command.process_group(0);
5426 command.stdin(Stdio::null());
5427 let mut child = match command.spawn() {
5428 Ok(child) => child,
5429 Err(source) => {
5430 #[cfg(target_os = "linux")]
5431 if let Some(placement) = cgroup_placement {
5432 remove_module_cgroup(placement, &cgroup_name);
5433 }
5434 return Err(SuperviseError::Spawn {
5435 program: spec.program.clone(),
5436 source,
5437 cgroup_path,
5438 });
5439 }
5440 };
5441 let spawned_at_ms = unix_ms_now();
5442 let spawned_from = spec.program.clone();
5443 let spawned_file_identity = spawned_file_identity(&spawned_from);
5444 let pid = child.id().ok_or_else(|| SuperviseError::Spawn {
5445 program: spec.program.clone(),
5446 source: io::Error::other("spawned child exposed no live pid"),
5447 cgroup_path: cgroup_path.clone(),
5448 })?;
5449 let process_start_time = crate::provenance::process_start_time(pid);
5450 let process_identity = process_start_time.map(|start_time| ProcessIdentity { pid, start_time });
5451 let roster_guard = roster.admit(
5452 spec.module_id.clone(),
5453 pid,
5454 spec.protocol,
5455 process_start_time,
5456 );
5457
5458 let stdout_pump = match child.stdout.take() {
5459 Some(stdout) => Some(tokio::spawn(pump_stdout_to(stdout, output_sink.clone()))),
5460 None => {
5461 warn!(
5462 module_id = %spec.module_id,
5463 "spawned child exposed no stdout pipe; file capture will be incomplete"
5464 );
5465 None
5466 }
5467 };
5468 let stderr_pump = match child.stderr.take() {
5469 Some(stderr) => {
5470 ring.lock()
5471 .unwrap_or_else(|poisoned| poisoned.into_inner())
5472 .push_process_start();
5473 Some(tokio::spawn(pump_stderr_to(
5474 stderr,
5475 Arc::clone(ring),
5476 output_sink,
5477 )))
5478 }
5479 None => {
5480 ring.lock()
5484 .unwrap_or_else(|poisoned| poisoned.into_inner())
5485 .mark_not_captured("stderr pipe was not available on spawn");
5486 warn!(
5487 module_id = %spec.module_id,
5488 "spawned child exposed no stderr pipe; tail will be unavailable"
5489 );
5490 None
5491 }
5492 };
5493
5494 Ok(SupervisedChild {
5495 child,
5496 #[cfg(target_os = "linux")]
5497 module_id: cgroup_name,
5498 #[cfg(target_os = "linux")]
5499 cgroup_placement: cgroup_placement.cloned(),
5500 stdout_pump,
5501 stderr_pump,
5502 stderr_ring: Arc::clone(ring),
5503 spawned_at_ms,
5504 spawned_from,
5505 spawned_file_identity,
5506 process_start_time,
5507 process_identity,
5508 pid,
5509 roster_guard: Some(roster_guard),
5510 })
5511}
5512
5513#[cfg(target_os = "linux")]
5514fn remove_module_cgroup(placement: &subc_cgroup::Placement, module_id: &str) {
5515 match placement.remove_module(module_id) {
5516 Ok(()) => debug!(module_id, "removed module cgroup after process exit"),
5517 Err(error) => warn!(
5518 module_id,
5519 error = %error,
5520 "could not remove module cgroup after process exit; continuing teardown"
5521 ),
5522 }
5523}
5524
5525#[cfg(target_os = "linux")]
5526fn apply_cgroup_placement(
5527 command: &mut Command,
5528 spec: &ModuleSpec,
5529 path: &std::path::Path,
5530) -> Result<(), SuperviseError> {
5531 subc_cgroup::apply(command, path).map_err(|source| SuperviseError::Cgroup {
5532 module_id: spec.module_id.clone(),
5533 source,
5534 })
5535}
5536
5537fn capture_retention(spec: &ModuleSpec) -> Retention {
5538 let defaults = Retention::default();
5539 let value = |name: &str| {
5540 spec.env
5541 .iter()
5542 .rev()
5543 .find_map(|(key, value)| (key == name).then_some(value.as_str()))
5544 };
5545 Retention {
5546 max_file_mb: value(CAPTURE_MAX_FILE_MB_ENV)
5547 .and_then(|value| value.parse().ok())
5548 .unwrap_or(defaults.max_file_mb),
5549 keep: value(CAPTURE_KEEP_ENV)
5550 .and_then(|value| value.parse().ok())
5551 .unwrap_or(defaults.keep),
5552 max_age_days: value(CAPTURE_MAX_AGE_DAYS_ENV)
5553 .and_then(|value| value.parse().ok())
5554 .unwrap_or(defaults.max_age_days),
5555 }
5556}
5557
5558fn generate_launch_nonce() -> Result<String, SuperviseError> {
5561 let mut bytes = [0u8; 32];
5562 getrandom::getrandom(&mut bytes).map_err(|source| SuperviseError::LaunchNonce {
5563 reason: source.to_string(),
5564 })?;
5565 let mut hex = String::with_capacity(64);
5566 for b in bytes {
5567 use std::fmt::Write;
5568 let _ = write!(hex, "{b:02x}");
5569 }
5570 Ok(hex)
5571}
5572
5573fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
5576 if a.len() != b.len() {
5577 return false;
5578 }
5579 let mut diff = 0u8;
5580 for (x, y) in a.iter().zip(b.iter()) {
5581 diff |= x ^ y;
5582 }
5583 diff == 0
5584}
5585
5586fn spawn_and_mark_running(
5587 spec: &ModuleSpec,
5588 runtime: &SupervisorRuntimeConfig,
5589 snapshot: &SharedSnapshot,
5590) -> Result<SupervisedChild, SuperviseError> {
5591 let child = spawn_child(
5592 spec,
5593 runtime.connection_file_path.as_deref(),
5594 runtime.supervisor_handle.as_ref(),
5595 &runtime.stderr_ring,
5596 runtime.capture_logs_dir.as_deref(),
5597 &runtime.child_roster,
5598 #[cfg(target_os = "linux")]
5599 runtime.cgroup_placement.as_ref(),
5600 )?;
5601 set_running(snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
5602 Ok(child)
5603}
5604
5605enum RegistrationWaitOutcome {
5606 Registered,
5607 Exited(ExitReport),
5608 TimedOut,
5609}
5610
5611struct ReloadRegistrationFailure {
5612 exit_report: ExitReport,
5613 reason: String,
5614}
5615
5616#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5617enum BusyGaugeObservation {
5618 Quiescent,
5619 Busy,
5620 Omitted,
5621}
5622
5623fn busy_gauge_observation(metrics: Option<&Value>, gauges: &[String]) -> BusyGaugeObservation {
5624 let Some(metrics) = metrics.and_then(Value::as_object) else {
5625 return BusyGaugeObservation::Omitted;
5626 };
5627 let mut sum = 0u128;
5628 for gauge in gauges {
5629 let Some(value) = metrics.get(gauge) else {
5630 return BusyGaugeObservation::Omitted;
5631 };
5632 let Some(value) = value.as_u64() else {
5633 return BusyGaugeObservation::Busy;
5634 };
5635 sum = sum.saturating_add(u128::from(value));
5636 }
5637 if sum == 0 {
5638 BusyGaugeObservation::Quiescent
5639 } else {
5640 BusyGaugeObservation::Busy
5641 }
5642}
5643
5644fn declared_busy_gauges(
5645 registry: &Registry,
5646 module_id: &str,
5647) -> Result<Vec<String>, SuperviseError> {
5648 busy_gauges_of(
5649 registry
5650 .get_module(module_id)
5651 .map_err(SuperviseError::Registry)?,
5652 )
5653}
5654
5655fn declared_busy_gauges_for_connection(
5659 registry: &Registry,
5660 connection_id: ConnectionId,
5661) -> Result<Vec<String>, SuperviseError> {
5662 busy_gauges_of(
5663 registry
5664 .get_module_by_connection(connection_id)
5665 .map_err(SuperviseError::Registry)?,
5666 )
5667}
5668
5669fn busy_gauges_of(
5670 registration: Option<crate::registry::ModuleRegistration>,
5671) -> Result<Vec<String>, SuperviseError> {
5672 let Some(registration) = registration else {
5673 return Ok(Vec::new());
5674 };
5675 let Some(self_signals) = registration.manifest.self_signals else {
5676 return Ok(Vec::new());
5677 };
5678
5679 let mut gauges = Vec::new();
5680 for declaration in self_signals {
5681 if declaration.kind != SelfSignalKind::Busy {
5682 continue;
5683 }
5684 match declaration.anchored_to {
5685 SignalAnchor::HealthGauges { gauges: declared } if !declared.is_empty() => {
5686 gauges.extend(declared)
5687 }
5688 _ => {
5689 gauges.push(String::new());
5692 }
5693 }
5694 }
5695 Ok(gauges)
5696}
5697
5698async fn wait_for_forwarding_quiescence(
5703 forwarding: &ForwardingTable,
5704 module_id: &str,
5705 runtime: &SupervisorRuntimeConfig,
5706 endpoint: crate::ModuleEndpointId,
5707 deadline: Instant,
5708 busy_gauges: &[String],
5709 scope: DrainScope,
5710) -> Result<bool, SuperviseError> {
5711 let mut gauges_quiescent = busy_gauges.is_empty();
5712 let mut next_probe_at = Instant::now();
5713 let mut omission_counted = false;
5714
5715 loop {
5716 let now = Instant::now();
5717 if !busy_gauges.is_empty() && now >= next_probe_at && now < deadline {
5718 let report = match scope {
5719 DrainScope::Active => probe_module_health(module_id, runtime, Some(deadline)).await,
5720 DrainScope::Endpoint(endpoint) => {
5721 probe_endpoint_health(endpoint, runtime, Some(deadline)).await
5722 }
5723 };
5724 gauges_quiescent = match report {
5725 Ok(report) => match busy_gauge_observation(report.metrics.as_ref(), busy_gauges) {
5726 BusyGaugeObservation::Quiescent => true,
5727 BusyGaugeObservation::Busy => false,
5728 BusyGaugeObservation::Omitted => {
5729 if !omission_counted {
5730 forwarding
5731 .counters()
5732 .increment_drains_with_undeclared_gauge();
5733 omission_counted = true;
5734 }
5735 false
5736 }
5737 },
5738 Err(err) => {
5739 warn!(
5740 module_id,
5741 error = %err,
5742 "drain health.check did not produce declared busy gauges; treating module as busy"
5743 );
5744 false
5745 }
5746 };
5747 next_probe_at = Instant::now() + runtime.health.cadence.max(REGISTRY_RELEASE_POLL);
5748 }
5749
5750 let in_flight = forwarding
5751 .endpoint_in_flight_count(endpoint)
5752 .map_err(SuperviseError::Forwarding)?;
5753 if in_flight == 0 && gauges_quiescent {
5754 return Ok(true);
5755 }
5756
5757 let now = Instant::now();
5758 if now >= deadline {
5759 return Ok(false);
5760 }
5761 let mut wait = deadline
5762 .saturating_duration_since(now)
5763 .min(REGISTRY_RELEASE_POLL);
5764 if !busy_gauges.is_empty() {
5765 wait = wait.min(next_probe_at.saturating_duration_since(now));
5766 }
5767 sleep(wait).await;
5768 }
5769}
5770
5771fn drained_after_quiescence_wait(wait_result: &Result<bool, SuperviseError>) -> bool {
5779 match wait_result {
5780 Ok(drained) => *drained,
5781 Err(_) => false,
5782 }
5783}
5784
5785fn send_route_goodbyes(forwarding: &ForwardingTable, released_routes: Vec<GoodbyeTarget>) {
5786 for released in released_routes {
5787 let frame = match Frame::build_with_version(
5788 released.negotiated_ver,
5789 FrameType::Goodbye,
5790 control_flags(),
5791 released.channel,
5792 released.epoch,
5793 0,
5794 Vec::new(),
5795 ) {
5796 Ok(frame) => frame,
5797 Err(err) => {
5798 warn!(
5799 route_channel = released.channel,
5800 error = %err,
5801 "failed to build supervisor drain route GOODBYE frame"
5802 );
5803 continue;
5804 }
5805 };
5806 if let Err(err) = released.sink.try_send(frame) {
5807 if released.close_on_delivery_failure() {
5808 warn!(
5809 target_connection_id = released.connection_id.get(),
5810 route_channel = released.channel,
5811 error = %err,
5812 "supervisor drain route GOODBYE was not delivered to client; closing target connection"
5813 );
5814 let _ = forwarding.escalate_client_delivery_failure(
5815 released.connection_id,
5816 released.channel,
5817 released.epoch,
5818 CloseReason::new(
5819 "route_goodbye_delivery_failed",
5820 format!(
5821 "failed to enqueue supervisor drain route GOODBYE for channel {}: {err}",
5822 released.channel
5823 ),
5824 ),
5825 );
5826 } else {
5827 warn!(
5828 target_connection_id = released.connection_id.get(),
5829 route_channel = released.channel,
5830 error = %err,
5831 "supervisor drain route GOODBYE to module dropped under backpressure; not closing shared module connection"
5832 );
5833 }
5834 }
5835 }
5836}
5837
5838fn send_module_draining(
5839 module_id: &str,
5840 reason: RouteCloseReason,
5841 deadline_ms: u64,
5842 target: &ModuleDrainTarget,
5843) {
5844 let body = match serde_json::to_vec(&ModuleControlCommand::Draining {
5845 reason,
5846 deadline_ms,
5847 }) {
5848 Ok(body) => body,
5849 Err(err) => {
5850 warn!(
5851 module_id,
5852 error = %err,
5853 "failed to encode module draining command"
5854 );
5855 return;
5856 }
5857 };
5858 let frame = match Frame::build_with_version(
5859 target.negotiated_ver,
5860 FrameType::Push,
5861 control_flags(),
5862 0,
5863 0,
5864 0,
5865 body,
5866 ) {
5867 Ok(frame) => frame,
5868 Err(err) => {
5869 warn!(
5870 module_id,
5871 error = %err,
5872 "failed to build module draining command frame"
5873 );
5874 return;
5875 }
5876 };
5877 if let Err(err) = target.sink.try_send(frame) {
5878 warn!(
5879 module_id,
5880 target_connection_id = target.endpoint.connection_id.get(),
5881 error = %err,
5882 "module draining command was not delivered to peer"
5883 );
5884 }
5885}
5886
5887fn send_module_goodbye(module_id: &str, forwarding: &ForwardingTable, target: &ModuleDrainTarget) {
5888 let frame = match Frame::build_with_version(
5889 target.negotiated_ver,
5890 FrameType::Goodbye,
5891 control_flags(),
5892 0,
5893 0,
5894 0,
5895 Vec::new(),
5896 ) {
5897 Ok(frame) => frame,
5898 Err(err) => {
5899 warn!(
5900 module_id,
5901 error = %err,
5902 "failed to build supervisor drain module GOODBYE frame"
5903 );
5904 return;
5905 }
5906 };
5907 if let Err(err) = target.sink.try_send(frame) {
5908 warn!(
5909 module_id,
5910 target_connection_id = target.endpoint.connection_id.get(),
5911 error = %err,
5912 "supervisor drain module GOODBYE was not delivered to peer; closing module connection"
5913 );
5914 forwarding.request_connection_close(
5915 target.endpoint.connection_id,
5916 CloseReason::new(
5917 "module_goodbye_delivery_failed",
5918 format!("failed to enqueue supervisor drain module GOODBYE for module '{module_id}': {err}"),
5919 ),
5920 );
5921 }
5922}
5923
5924#[derive(Clone, Copy)]
5925struct ForwardingDrainContext<'a> {
5926 spec: &'a ModuleSpec,
5927 runtime: &'a SupervisorRuntimeConfig,
5928 registry: &'a Registry,
5929 scope: DrainScope,
5930}
5931
5932#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5934enum DrainScope {
5935 Active,
5938 Endpoint(crate::ModuleEndpointId),
5943}
5944
5945async fn begin_forwarding_drain(
5946 spec: &ModuleSpec,
5947 runtime: &SupervisorRuntimeConfig,
5948 registry: &Registry,
5949 snapshot: &SharedSnapshot,
5950 enabled: Option<bool>,
5951 reason: RouteCloseReason,
5952) -> Result<(), SuperviseError> {
5953 let Some(forwarding) = runtime.forwarding.as_ref() else {
5954 return Err(SuperviseError::ReloadUnavailable {
5955 module_id: spec.module_id.clone(),
5956 reason: "supervisor was not configured with a forwarding table".to_string(),
5957 });
5958 };
5959
5960 begin_forwarding_drain_with(
5961 forwarding,
5962 ForwardingDrainContext {
5963 spec,
5964 runtime,
5965 registry,
5966 scope: DrainScope::Active,
5967 },
5968 snapshot,
5969 enabled,
5970 reason,
5971 runtime.drain_timeout,
5972 )
5973 .await
5974}
5975
5976async fn begin_forwarding_drain_if_configured(
5977 spec: &ModuleSpec,
5978 runtime: &SupervisorRuntimeConfig,
5979 registry: &Registry,
5980 snapshot: &SharedSnapshot,
5981 enabled: Option<bool>,
5982 reason: RouteCloseReason,
5983) -> Result<(), SuperviseError> {
5984 begin_forwarding_drain_with_timeout(
5985 spec,
5986 runtime,
5987 registry,
5988 snapshot,
5989 enabled,
5990 reason,
5991 runtime.drain_timeout,
5992 )
5993 .await
5994}
5995
5996async fn begin_forwarding_drain_with_timeout(
6000 spec: &ModuleSpec,
6001 runtime: &SupervisorRuntimeConfig,
6002 registry: &Registry,
6003 snapshot: &SharedSnapshot,
6004 enabled: Option<bool>,
6005 reason: RouteCloseReason,
6006 drain_timeout: Duration,
6007) -> Result<(), SuperviseError> {
6008 let Some(forwarding) = runtime.forwarding.as_ref() else {
6009 return Ok(());
6010 };
6011
6012 begin_forwarding_drain_with(
6013 forwarding,
6014 ForwardingDrainContext {
6015 spec,
6016 runtime,
6017 registry,
6018 scope: DrainScope::Active,
6019 },
6020 snapshot,
6021 enabled,
6022 reason,
6023 drain_timeout,
6024 )
6025 .await
6026}
6027
6028async fn begin_forwarding_drain_with(
6029 forwarding: &ForwardingTable,
6030 context: ForwardingDrainContext<'_>,
6031 snapshot: &SharedSnapshot,
6032 enabled: Option<bool>,
6033 reason: RouteCloseReason,
6034 drain_timeout: Duration,
6035) -> Result<(), SuperviseError> {
6036 let ForwardingDrainContext {
6037 spec,
6038 runtime,
6039 registry,
6040 scope,
6041 } = context;
6042 debug_assert_ne!(reason, RouteCloseReason::Crash);
6043 let terminal = matches!(reason, RouteCloseReason::Disable);
6044 let drain_started_at = Instant::now();
6045 let drain_deadline = drain_started_at + drain_timeout;
6046 let deadline_ms =
6047 unix_ms_now().saturating_add(u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX));
6048 let busy_gauges = match scope {
6049 DrainScope::Active => declared_busy_gauges(registry, &spec.module_id)?,
6050 DrainScope::Endpoint(endpoint) => {
6051 declared_busy_gauges_for_connection(registry, endpoint.connection_id)?
6052 }
6053 };
6054
6055 let drain_target = match scope {
6058 DrainScope::Active => forwarding.begin_module_drain(&spec.module_id, reason),
6059 DrainScope::Endpoint(endpoint) => forwarding.begin_endpoint_drain(endpoint, reason),
6060 }
6061 .map_err(SuperviseError::Forwarding)?;
6062 if scope == DrainScope::Active {
6063 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6064 state.state = ModuleState::Draining;
6065 if let Some(enabled) = enabled {
6066 state.enabled = enabled;
6067 }
6068 })?;
6069 }
6070
6071 if let Some(target) = drain_target.as_ref() {
6072 send_module_draining(&spec.module_id, reason, deadline_ms, target);
6073 let routes = forwarding
6074 .endpoint_routes(target.endpoint)
6075 .map_err(SuperviseError::Forwarding)?;
6076 let routes_notified = routes.len();
6077 crate::control::send_route_control_pushes(
6078 forwarding,
6079 routes.clone(),
6080 ClientControlPush::RouteClosing {
6081 module_id: spec.module_id.clone(),
6082 reason,
6083 },
6084 );
6085 send_route_goodbyes(forwarding, target.abandoned_bindings.clone());
6086
6087 let wait_result = wait_for_forwarding_quiescence(
6093 forwarding,
6094 &spec.module_id,
6095 runtime,
6096 target.endpoint,
6097 drain_deadline,
6098 &busy_gauges,
6099 scope,
6100 )
6101 .await;
6102 let drained = drained_after_quiescence_wait(&wait_result);
6103 if let Err(err) = &wait_result {
6104 error!(
6105 module_id = %spec.module_id,
6106 ?reason,
6107 error = %err,
6108 "forwarding quiescence wait failed after route.closing; forcing route.closed(drained: false) so the client is not left waiting on an unfulfilled promise"
6109 );
6110 } else if !drained {
6111 let holdouts = forwarding
6117 .endpoint_drain_holdouts(target.endpoint)
6118 .unwrap_or_default();
6119 warn!(
6120 module_id = %spec.module_id,
6121 waited = ?drain_timeout,
6122 ?reason,
6123 held_requests = holdouts.requests,
6124 held_routes = holdouts.routes,
6125 total_routes = holdouts.total_routes,
6126 top_connections = ?holdouts.top_connections,
6127 held = %holdouts
6130 .held
6131 .iter()
6132 .map(|(channel, corr)| format!("{channel}:{corr}"))
6133 .collect::<Vec<_>>()
6134 .join(","),
6135 "route drain timed out before request quiescence; forcing teardown"
6136 );
6137 }
6138 crate::control::send_route_control_pushes(
6139 forwarding,
6140 routes,
6141 ClientControlPush::RouteClosed {
6142 module_id: spec.module_id.clone(),
6143 reason,
6144 drained,
6145 abandoned: target.abandoned_bindings.len() as u32,
6146 excluded_subscriptions: target.excluded_subscriptions,
6147 terminal: Some(terminal),
6148 },
6149 );
6150 wait_result?;
6151
6152 let released_routes = match forwarding.release_module_endpoint_routes(target.endpoint) {
6158 Ok(routes) => routes,
6159 Err(err) => {
6160 warn!(
6161 module_id = %spec.module_id,
6162 ?reason,
6163 error = %err,
6164 "failed to release module endpoint routes after route.closed; module GOODBYE will still be sent"
6165 );
6166 send_module_goodbye(&spec.module_id, forwarding, target);
6167 return Err(SuperviseError::Forwarding(err));
6168 }
6169 };
6170 let route_goodbye_count = released_routes.len();
6171 send_route_goodbyes(forwarding, released_routes);
6172 send_module_goodbye(&spec.module_id, forwarding, target);
6173
6174 info!(
6180 module_id = %spec.module_id,
6181 ?reason,
6182 routes_notified,
6183 route_goodbyes = route_goodbye_count,
6184 abandoned_reservations = target.abandoned_bindings.len(),
6185 excluded_subscriptions = target.excluded_subscriptions,
6186 drained,
6187 "module drain complete; consumers notified via route.closing/route.closed pushes and per-route GOODBYE frames"
6188 );
6189 }
6190
6191 Ok(())
6192}
6193
6194async fn wait_for_registration_after_reload(
6197 registry: &Registry,
6198 module_id: &str,
6199 snapshot: &SharedSnapshot,
6200 child: &mut SupervisedChild,
6201 wait: Duration,
6202) -> Result<RegistrationWaitOutcome, SuperviseError> {
6203 wait_for_slot_registration(
6204 registry,
6205 crate::registry::RegistrationSlot::Active(module_id),
6206 module_id,
6207 snapshot,
6208 child,
6209 wait,
6210 )
6211 .await
6212}
6213
6214async fn wait_for_slot_registration(
6222 registry: &Registry,
6223 slot: crate::registry::RegistrationSlot<'_>,
6224 module_id: &str,
6225 snapshot: &SharedSnapshot,
6226 child: &mut SupervisedChild,
6227 wait: Duration,
6228) -> Result<RegistrationWaitOutcome, SuperviseError> {
6229 let deadline = Instant::now() + wait;
6230 loop {
6231 if registry
6232 .registration(slot)
6233 .map_err(SuperviseError::Registry)?
6234 .is_some()
6235 {
6236 return Ok(RegistrationWaitOutcome::Registered);
6237 }
6238
6239 let now = Instant::now();
6240 if now >= deadline {
6241 return Ok(RegistrationWaitOutcome::TimedOut);
6242 }
6243 let remaining = deadline.saturating_duration_since(now);
6244 let poll = remaining.min(REGISTRY_RELEASE_POLL);
6245
6246 tokio::select! {
6247 wait_result = child.wait() => {
6248 let status = wait_result.map_err(|source| SuperviseError::Wait {
6249 module_id: module_id.to_string(),
6250 source,
6251 })?;
6252 return Ok(RegistrationWaitOutcome::Exited(classify_reaped_child_exit(
6253 snapshot,
6254 child,
6255 &status,
6256 )));
6257 }
6258 _ = sleep(poll) => {}
6259 }
6260 }
6261}
6262
6263fn registration_failure_exit_report(mut exit_report: ExitReport) -> ExitReport {
6264 if exit_report.kind != ExitKind::DeliberateSeverance {
6267 exit_report.kind = ExitKind::Crash;
6268 }
6269 exit_report
6270}
6271
6272async fn handle_reload_child_registration_failure(
6273 spec: &ModuleSpec,
6274 runtime: &SupervisorRuntimeConfig,
6275 registry: &Registry,
6276 process_liveness: &SupervisorProcessLiveness,
6277 snapshot: &SharedSnapshot,
6278 child: &mut Option<SupervisedChild>,
6279 failure: ReloadRegistrationFailure,
6280) -> Result<(), SuperviseError> {
6281 let ReloadRegistrationFailure {
6282 exit_report,
6283 reason,
6284 } = failure;
6285 match on_child_exit(
6286 spec,
6287 runtime.restart_policy,
6288 registry,
6289 snapshot,
6290 &runtime.terminal_ring,
6291 &runtime.spawn_events,
6292 exit_report,
6293 )
6294 .await
6295 {
6296 NextAction::Stop {
6297 registration_released,
6298 } => {
6299 if registration_released {
6300 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6301 }
6302 }
6303 NextAction::Restart { schedule } => {
6304 let delay = schedule.map_or(runtime.restart_policy.delay_for_restart(0), |schedule| {
6305 schedule.delay
6306 });
6307 if let Some(schedule) = schedule {
6308 log_crash_respawn(&spec.module_id, schedule);
6309 }
6310 sleep(delay).await;
6311 if respawn_still_pending(snapshot) {
6315 if let Err(err) = wait_for_registration_release(
6316 registry,
6317 &spec.module_id,
6318 REGISTRY_RELEASE_TIMEOUT,
6319 )
6320 .await
6321 {
6322 fail_snapshot(snapshot, Some(&spec.module_id), None);
6323 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6324 return Err(SuperviseError::ReloadFailed {
6325 module_id: spec.module_id.clone(),
6326 reason: format!(
6327 "{reason}; registration did not release before policy retry: {err}"
6328 ),
6329 });
6330 }
6331 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6332 match spawn_and_mark_running(spec, runtime, snapshot) {
6333 Ok(next_child) => {
6334 *child = Some(next_child);
6335 }
6336 Err(err) => {
6337 fail_snapshot(snapshot, Some(&spec.module_id), None);
6338 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6339 return Err(SuperviseError::ReloadFailed {
6340 module_id: spec.module_id.clone(),
6341 reason: format!("{reason}; policy retry spawn failed: {err}"),
6342 });
6343 }
6344 }
6345 }
6346 }
6347 }
6348
6349 Err(SuperviseError::ReloadFailed {
6350 module_id: spec.module_id.clone(),
6351 reason,
6352 })
6353}
6354
6355async fn handle_reload_spawn_failure(
6356 spec: &ModuleSpec,
6357 runtime: &SupervisorRuntimeConfig,
6358 process_liveness: &SupervisorProcessLiveness,
6359 snapshot: &SharedSnapshot,
6360 child: &mut Option<SupervisedChild>,
6361 reason: String,
6362) -> Result<(), SuperviseError> {
6363 let mut should_retry = false;
6364 let now = Instant::now();
6365 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6366 clear_current_process_facts(state);
6367 if daemon_will_restart(state, &runtime.restart_policy, now) {
6368 state.record_crash_restart(&runtime.restart_policy, now);
6369 state.state = ModuleState::Restarting;
6370 should_retry = true;
6371 } else if state.enabled {
6372 state.state = ModuleState::Failed;
6373 } else {
6374 state.state = ModuleState::Disabled;
6375 }
6376 })?;
6377
6378 if should_retry {
6379 sleep(runtime.restart_policy.backoff).await;
6380 if respawn_still_pending(snapshot) {
6384 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6385 match spawn_and_mark_running(spec, runtime, snapshot) {
6386 Ok(next_child) => {
6387 *child = Some(next_child);
6388 }
6389 Err(err) => {
6390 fail_snapshot(snapshot, Some(&spec.module_id), None);
6391 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6392 return Err(SuperviseError::ReloadFailed {
6393 module_id: spec.module_id.clone(),
6394 reason: format!("{reason}; policy retry spawn failed: {err}"),
6395 });
6396 }
6397 }
6398 }
6399 } else {
6400 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6401 }
6402
6403 Err(SuperviseError::ReloadFailed {
6404 module_id: spec.module_id.clone(),
6405 reason,
6406 })
6407}
6408
6409fn control_flags() -> Flags {
6410 Flags::new(false, Priority::Passive, false)
6411}
6412
6413#[allow(clippy::too_many_arguments)]
6414async fn drain_optional_child(
6415 module_id: &str,
6416 protocol: ModuleProtocol,
6417 registry: &Registry,
6418 snapshot: &SharedSnapshot,
6419 terminal_ring: &Arc<Mutex<TerminalRing>>,
6420 spawn_events: &SpawnEventFeed,
6421 child: &mut Option<SupervisedChild>,
6422 drain_timeout: Duration,
6423 final_state: ModuleState,
6424 enabled: Option<bool>,
6425) -> Result<(), SuperviseError> {
6426 if let Some(child) = child.take() {
6427 drain_child_to_state(
6428 module_id,
6429 protocol,
6430 registry,
6431 snapshot,
6432 terminal_ring,
6433 spawn_events,
6434 child,
6435 drain_timeout,
6436 final_state,
6437 enabled,
6438 )
6439 .await
6440 } else {
6441 update_snapshot(snapshot, Some(module_id), |state| {
6442 state.state = final_state;
6443 if let Some(enabled) = enabled {
6444 state.enabled = enabled;
6445 }
6446 clear_current_process_facts(state);
6447 })?;
6448 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
6449 }
6450}
6451
6452#[allow(clippy::too_many_arguments)]
6453async fn drain_child_to_state(
6454 module_id: &str,
6455 protocol: ModuleProtocol,
6456 registry: &Registry,
6457 snapshot: &SharedSnapshot,
6458 terminal_ring: &Arc<Mutex<TerminalRing>>,
6459 spawn_events: &SpawnEventFeed,
6460 mut child: SupervisedChild,
6461 drain_timeout: Duration,
6462 final_state: ModuleState,
6463 enabled: Option<bool>,
6464) -> Result<(), SuperviseError> {
6465 update_snapshot(snapshot, Some(module_id), |state| {
6466 state.state = ModuleState::Draining;
6467 if let Some(enabled) = enabled {
6468 state.enabled = enabled;
6469 }
6470 })?;
6471
6472 if protocol == ModuleProtocol::None {
6478 request_graceful_stop(module_id, &child);
6479 }
6480
6481 let exit_report = match timeout(drain_timeout, child.wait()).await {
6482 Ok(Ok(status)) => classify_reaped_child_exit(snapshot, &child, &status),
6483 Ok(Err(source)) => {
6484 fail_snapshot(snapshot, Some(module_id), None);
6485 return Err(SuperviseError::Wait {
6486 module_id: module_id.to_string(),
6487 source,
6488 });
6489 }
6490 Err(_) => {
6491 child.start_kill().map_err(|source| {
6500 fail_snapshot(snapshot, Some(module_id), None);
6501 SuperviseError::Kill {
6502 module_id: module_id.to_string(),
6503 source,
6504 }
6505 })?;
6506 let status = child.wait().await.map_err(|source| {
6507 fail_snapshot(snapshot, Some(module_id), None);
6508 SuperviseError::Wait {
6509 module_id: module_id.to_string(),
6510 source,
6511 }
6512 })?;
6513 classify_reaped_child_exit(snapshot, &child, &status)
6514 }
6515 };
6516
6517 update_snapshot(snapshot, Some(module_id), |state| {
6518 state.state = final_state;
6519 if let Some(enabled) = enabled {
6520 state.enabled = enabled;
6521 }
6522 clear_current_process_facts(state);
6523 state.last_exit = Some(exit_report.clone());
6524 if exit_report.kind == ExitKind::DeliberateSeverance {
6525 state.lifetime_restarts += 1;
6526 }
6527 })?;
6528 record_terminal(
6529 module_id,
6530 terminal_ring,
6531 spawn_events,
6532 &exit_report,
6533 terminal_disposition(final_state),
6534 );
6535 child.drain_stderr(module_id).await;
6536
6537 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
6538}
6539
6540#[cfg(unix)]
6560fn request_graceful_stop(module_id: &str, child: &SupervisedChild) {
6561 let Some(pid) = child
6562 .id()
6563 .and_then(|pid| i32::try_from(pid).ok())
6564 .and_then(rustix::process::Pid::from_raw)
6565 else {
6566 debug!(
6567 module_id,
6568 "no pid to signal for protocol: none teardown; falling through to the drain wait"
6569 );
6570 return;
6571 };
6572 match rustix::process::kill_process(pid, rustix::process::Signal::TERM) {
6573 Ok(()) => debug!(module_id, "sent SIGTERM to protocol: none module"),
6574 Err(err) => debug!(
6575 module_id,
6576 error = %err,
6577 "SIGTERM to protocol: none module failed; the drain wait and kill still apply"
6578 ),
6579 }
6580}
6581
6582#[cfg(not(unix))]
6590fn request_graceful_stop(module_id: &str, _child: &SupervisedChild) {
6591 debug!(
6592 module_id,
6593 "no graceful stop signal exists on this platform; protocol: none teardown waits, then kills"
6594 );
6595}
6596
6597fn terminal_disposition(final_state: ModuleState) -> TerminalDisposition {
6598 match final_state {
6599 ModuleState::Stopped => TerminalDisposition::Stopped,
6600 ModuleState::Disabled => TerminalDisposition::Disabled,
6601 ModuleState::Restarting => TerminalDisposition::Restarting,
6602 ModuleState::Failed => TerminalDisposition::Failed,
6603 ModuleState::Starting
6604 | ModuleState::Running
6605 | ModuleState::Unresponsive
6606 | ModuleState::Draining => {
6607 unreachable!("terminal exits only finish in terminal or restarting states")
6608 }
6609 }
6610}
6611
6612async fn wait_for_registration_release(
6615 registry: &Registry,
6616 module_id: &str,
6617 wait: Duration,
6618) -> Result<(), SuperviseError> {
6619 wait_for_slot_registration_release(
6620 registry,
6621 crate::registry::RegistrationSlot::Active(module_id),
6622 wait,
6623 )
6624 .await
6625}
6626
6627async fn wait_for_slot_registration_release(
6635 registry: &Registry,
6636 slot: crate::registry::RegistrationSlot<'_>,
6637 wait: Duration,
6638) -> Result<(), SuperviseError> {
6639 let deadline = Instant::now() + wait;
6640 let mut release_events = registration_release_events().subscribe();
6641 let still_active = |registration: &crate::registry::ModuleRegistration| {
6642 SuperviseError::RegistrationStillActive {
6643 module_id: registration.manifest.module_id.clone(),
6644 waited: wait,
6645 }
6646 };
6647 loop {
6648 let _observed_generation = *release_events.borrow_and_update();
6649 let Some(registration) = registry
6650 .registration(slot)
6651 .map_err(SuperviseError::Registry)?
6652 else {
6653 return Ok(());
6654 };
6655
6656 let now = Instant::now();
6657 if now >= deadline {
6658 return Err(still_active(®istration));
6659 }
6660
6661 let remaining = deadline.saturating_duration_since(now);
6662 match timeout(remaining, release_events.changed()).await {
6663 Ok(Ok(())) | Ok(Err(_)) => {}
6664 Err(_) => return Err(still_active(®istration)),
6665 }
6666 }
6667}
6668
6669#[cfg(test)]
6670mod slot_registration_wait_tests {
6671 use super::*;
6672 use crate::registry::{ConnectionId, RegistrationSlot};
6673 use subc_protocol::manifest::ModuleManifest;
6674
6675 const INCUMBENT: u64 = 1;
6676 const CANDIDATE: u64 = 2;
6677
6678 fn swapped_registry() -> Arc<Registry> {
6679 let registry = Arc::new(Registry::default());
6680 let manifest = ModuleManifest::builder("m", "0.1.0").build();
6681 registry
6682 .register_with_control_ops(
6683 manifest.clone(),
6684 1,
6685 ConnectionId::new(INCUMBENT),
6686 Vec::new(),
6687 )
6688 .unwrap();
6689 registry
6690 .register_candidate_with_control_ops(
6691 manifest,
6692 1,
6693 ConnectionId::new(CANDIDATE),
6694 Vec::new(),
6695 )
6696 .unwrap();
6697 registry
6698 }
6699
6700 #[tokio::test]
6704 async fn incumbent_release_is_awaited_by_connection_not_by_module_id() {
6705 let registry = swapped_registry();
6706 registry.promote_candidate("m").unwrap().unwrap();
6707
6708 assert!(matches!(
6709 wait_for_registration_release(®istry, "m", Duration::from_millis(50)).await,
6710 Err(SuperviseError::RegistrationStillActive { .. })
6711 ));
6712
6713 assert!(matches!(
6715 wait_for_slot_registration_release(
6716 ®istry,
6717 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
6718 Duration::from_millis(50),
6719 )
6720 .await,
6721 Err(SuperviseError::RegistrationStillActive { .. })
6722 ));
6723
6724 let releaser = Arc::clone(®istry);
6725 let release = tokio::spawn(async move {
6726 sleep(Duration::from_millis(20)).await;
6727 releaser
6728 .deregister_connection(ConnectionId::new(INCUMBENT))
6729 .unwrap();
6730 notify_registration_release();
6731 });
6732 wait_for_slot_registration_release(
6733 ®istry,
6734 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
6735 Duration::from_secs(5),
6736 )
6737 .await
6738 .expect("the incumbent's own registration is released");
6739 release.await.unwrap();
6740 assert!(registry.get_module("m").unwrap().is_some());
6741 }
6742
6743 #[tokio::test]
6746 async fn candidate_slot_wait_ignores_the_incumbents_registration() {
6747 let registry = swapped_registry();
6748 assert!(matches!(
6749 wait_for_slot_registration_release(
6750 ®istry,
6751 RegistrationSlot::Candidate("m"),
6752 Duration::from_millis(50),
6753 )
6754 .await,
6755 Err(SuperviseError::RegistrationStillActive { .. })
6756 ));
6757 registry
6758 .deregister_connection(ConnectionId::new(CANDIDATE))
6759 .unwrap();
6760 wait_for_slot_registration_release(
6761 ®istry,
6762 RegistrationSlot::Candidate("m"),
6763 Duration::from_millis(50),
6764 )
6765 .await
6766 .expect("a candidate slot with no candidate is released");
6767 assert!(registry
6768 .registration(RegistrationSlot::Active("m"))
6769 .unwrap()
6770 .is_some());
6771 }
6772}
6773
6774fn classify_exit(status: &ExitStatus) -> ExitReport {
6775 ExitReport {
6776 kind: if status.success() {
6777 ExitKind::Clean
6778 } else {
6779 ExitKind::Crash
6780 },
6781 code: status.code(),
6782 signal: exit_signal(status),
6783 at_ms: unix_ms_now(),
6784 }
6785}
6786
6787fn wait_error_exit_report() -> ExitReport {
6793 ExitReport {
6794 kind: ExitKind::Crash,
6795 code: None,
6796 signal: None,
6797 at_ms: unix_ms_now(),
6798 }
6799}
6800
6801#[cfg(unix)]
6802fn exit_signal(status: &ExitStatus) -> Option<i32> {
6803 use std::os::unix::process::ExitStatusExt;
6804
6805 status.signal()
6806}
6807
6808#[cfg(not(unix))]
6809fn exit_signal(_status: &ExitStatus) -> Option<i32> {
6810 None
6811}
6812
6813fn reset_restart_count(snapshot: &SharedSnapshot, module_id: &str) -> Result<(), SuperviseError> {
6819 update_snapshot(snapshot, Some(module_id), |state| {
6820 state.clear_crash_restarts();
6821 })
6822}
6823
6824fn set_running(
6825 snapshot: &SharedSnapshot,
6826 child: &SupervisedChild,
6827 module_id: &str,
6828 spawn_events: &SpawnEventFeed,
6829) -> Result<(), SuperviseError> {
6830 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
6831 module_id: Some(module_id.to_string()),
6832 })?;
6833 state.spawn_generation = spawn_events.emit_spawned(module_id, child.pid, child.spawned_at_ms);
6834 state.in_alternate_slot = false;
6837 state.state = ModuleState::Running;
6838 state.enabled = true;
6839 state.process_alive = true;
6840 state.pid = child.id();
6841 state.spawned_at_ms = Some(child.spawned_at_ms);
6842 state.spawned_from = Some(child.spawned_from.clone());
6843 state.spawned_file_identity = child.spawned_file_identity;
6844 state.process_start_time = child.process_start_time;
6845 Ok(())
6846}
6847
6848fn clear_current_process_facts(state: &mut SupervisorSnapshot) {
6849 state.process_alive = false;
6850 state.pid = None;
6851 state.spawned_at_ms = None;
6852 state.spawned_from = None;
6853 state.spawned_file_identity = None;
6854 state.process_start_time = None;
6855 state.deliberate_severance = None;
6856}
6857
6858#[cfg(test)]
6859fn record_deliberate_severance(
6860 snapshot: &SharedSnapshot,
6861 identity: ProcessIdentity,
6862) -> Result<(), SuperviseError> {
6863 update_snapshot(snapshot, None, |state| {
6864 state.deliberate_severance = Some(identity);
6865 })
6866}
6867
6868fn apply_deliberate_severance_marker(
6869 snapshot: &SharedSnapshot,
6870 exited_identity: Option<ProcessIdentity>,
6871 mut exit_report: ExitReport,
6872) -> ExitReport {
6873 let marker = lock_snapshot(snapshot)
6874 .ok()
6875 .and_then(|mut state| state.deliberate_severance.take());
6876 if marker.is_some() && marker == exited_identity {
6877 exit_report.kind = ExitKind::DeliberateSeverance;
6878 }
6879 exit_report
6880}
6881
6882fn classify_reaped_child_exit(
6883 snapshot: &SharedSnapshot,
6884 child: &SupervisedChild,
6885 status: &ExitStatus,
6886) -> ExitReport {
6887 apply_deliberate_severance_marker(snapshot, child.process_identity(), classify_exit(status))
6888}
6889
6890fn fail_snapshot(
6891 snapshot: &SharedSnapshot,
6892 module_id: Option<&str>,
6893 last_exit: Option<ExitReport>,
6894) {
6895 if let Err(err) = update_snapshot(snapshot, module_id, |state| {
6896 state.state = ModuleState::Failed;
6897 clear_current_process_facts(state);
6898 if let Some(last_exit) = last_exit {
6899 state.last_exit = Some(last_exit);
6900 }
6901 }) {
6902 error!(error = %err, "failed to mark supervisor state failed");
6903 }
6904}
6905
6906fn update_snapshot(
6907 snapshot: &SharedSnapshot,
6908 module_id: Option<&str>,
6909 update: impl FnOnce(&mut SupervisorSnapshot),
6910) -> Result<(), SuperviseError> {
6911 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
6912 module_id: module_id.map(ToOwned::to_owned),
6913 })?;
6914 update(&mut state);
6915 Ok(())
6916}
6917
6918const SLOW_SNAPSHOT_LOCK_THRESHOLD: Duration = Duration::from_millis(250);
6919
6920fn lock_snapshot_for_control<'a>(
6921 snapshot: &'a SharedSnapshot,
6922 module_id: &str,
6923 caller: &'static str,
6924) -> Result<std::sync::MutexGuard<'a, SupervisorSnapshot>, SuperviseError> {
6925 let started_at = Instant::now();
6926 let guard = lock_snapshot(snapshot)?;
6927 let waited = started_at.elapsed();
6928 if waited >= SLOW_SNAPSHOT_LOCK_THRESHOLD {
6929 warn!(
6930 module_id = %module_id,
6931 waited_ms = waited.as_millis() as u64,
6932 caller = %caller,
6933 "slow snapshot lock"
6934 );
6935 }
6936 Ok(guard)
6937}
6938
6939fn lock_snapshot(
6940 snapshot: &SharedSnapshot,
6941) -> Result<std::sync::MutexGuard<'_, SupervisorSnapshot>, SuperviseError> {
6942 snapshot
6943 .lock()
6944 .map_err(|_| SuperviseError::StatePoisoned { module_id: None })
6945}
6946
6947#[cfg(test)]
6948mod terminal_history_tests {
6949 use std::{
6950 path::PathBuf,
6951 sync::Arc,
6952 time::{Duration, Instant},
6953 };
6954
6955 use tokio::time::sleep;
6956
6957 use super::{
6958 apply_deliberate_severance_marker, daemon_will_restart, drain_child_to_state,
6959 drained_after_quiescence_wait, handle_reload_spawn_failure, health_restart_child,
6960 lock_snapshot, on_child_exit, record_deliberate_severance, record_wait_error_terminal,
6961 reset_restart_count, spawn_and_mark_running, update_snapshot, wait_error_exit_report,
6962 ExitKind, ExitReport, ModuleProtocol, ModuleSpec, ModuleState, NextAction, ProcessIdentity,
6963 RestartPolicy, SpawnEventKind, SuperviseError, SupervisedModule, Supervisor,
6964 SupervisorHandle, SupervisorHealthStatus, SupervisorSnapshot,
6965 };
6966 use super::Instant as ClockInstant;
6971 use crate::{
6972 registry::Registry,
6973 terminal_ring::{TerminalRing, TerminalRingConfig},
6974 };
6975 use std::sync::Mutex;
6976 use subc_control::TerminalDisposition;
6977
6978 fn fake_aft_stub_path() -> PathBuf {
6983 let mut path = std::env::current_exe().expect("current_exe available in tests");
6984 path.pop();
6985 path.pop();
6986 path.push(if cfg!(windows) {
6987 "fake-aft-stub.exe"
6988 } else {
6989 "fake-aft-stub"
6990 });
6991 assert!(
6992 path.exists(),
6993 "fake-aft-stub not built at {}: run `cargo test -p subc-core` (which builds \
6994 [[bin]] targets) rather than `cargo test -p subc-core --lib` (which does not)",
6995 path.display()
6996 );
6997 path
6998 }
6999
7000 #[test]
7001 fn reserved_never_spawned_refuses_every_hello() {
7002 let supervisor = SupervisorHandle::default();
7007 supervisor.apply_identity_configuration(&ModuleSpec {
7008 module_id: "never-spawned".to_string(),
7009 program: PathBuf::from("/usr/bin/false"),
7010 args: Vec::new(),
7011 env: Vec::new(),
7012 reserved: true,
7013 reserved_prefixes: Vec::new(),
7014 protocol: ModuleProtocol::Subc,
7015 overlap: Default::default(),
7016 });
7017 assert!(
7018 supervisor
7019 .reserved_hello_rejection("never-spawned", Some("any-forged-nonce"))
7020 .is_some(),
7021 "forged nonce must refuse on a reserved never-spawned id"
7022 );
7023 assert!(
7024 supervisor
7025 .reserved_hello_rejection("never-spawned", None)
7026 .is_some(),
7027 "absent nonce must refuse on a reserved never-spawned id"
7028 );
7029 supervisor.set_spawn_nonce("never-spawned", "minted".to_string());
7031 supervisor.apply_identity_configuration(&ModuleSpec {
7032 module_id: "never-spawned".to_string(),
7033 program: PathBuf::from("/usr/bin/false"),
7034 args: Vec::new(),
7035 env: Vec::new(),
7036 reserved: true,
7037 reserved_prefixes: Vec::new(),
7038 protocol: ModuleProtocol::Subc,
7039 overlap: Default::default(),
7040 });
7041 assert!(supervisor
7042 .reserved_hello_rejection("never-spawned", Some("minted"))
7043 .is_none());
7044 assert!(supervisor
7045 .reserved_hello_rejection("never-spawned", Some("forged"))
7046 .is_some());
7047 }
7048
7049 fn seed_crash_restarts(state: &mut SupervisorSnapshot, count: u32) {
7052 let now = ClockInstant::now();
7053 for _ in 0..count {
7054 state.crash_restarts.push_back(now);
7055 }
7056 }
7057
7058 fn age_oldest_crash_restart_out_of_window(state: &mut SupervisorSnapshot, window: Duration) {
7062 let aged = state
7063 .crash_restarts
7064 .front()
7065 .expect("a crash restart must be recorded before it can be aged")
7066 .checked_sub(window + Duration::from_secs(1))
7067 .expect("the test clock is far enough from its origin to age an instant");
7068 state.crash_restarts[0] = aged;
7069 }
7070
7071 fn snapshot_with_restarts(enabled: bool, count: u32) -> SupervisorSnapshot {
7072 let mut state = SupervisorSnapshot::new(ModuleState::Running, enabled);
7073 seed_crash_restarts(&mut state, count);
7074 state
7075 }
7076
7077 #[test]
7078 fn daemon_owned_recovery_predicate_uses_the_pre_increment_budget() {
7079 let policy = RestartPolicy::new(3, Duration::ZERO);
7080 let now = ClockInstant::now();
7081 assert!(daemon_will_restart(
7082 &mut snapshot_with_restarts(true, 2),
7083 &policy,
7084 now
7085 ));
7086 assert!(!daemon_will_restart(
7087 &mut snapshot_with_restarts(true, 3),
7088 &policy,
7089 now
7090 ));
7091 assert!(!daemon_will_restart(
7092 &mut snapshot_with_restarts(false, 0),
7093 &policy,
7094 now
7095 ));
7096 }
7097
7098 #[test]
7099 fn crash_restart_backoff_escalates_with_in_window_count() {
7100 let policy = RestartPolicy::new(4, Duration::from_millis(100))
7101 .with_max_backoff(Duration::from_secs(30));
7102 let now = ClockInstant::now();
7103 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7104 let schedules = (0..4)
7105 .map(|_| {
7106 state
7107 .next_crash_restart(&policy, now)
7108 .expect("the test policy allows four crash restarts")
7109 })
7110 .collect::<Vec<_>>();
7111
7112 assert_eq!(
7113 schedules
7114 .iter()
7115 .map(|schedule| schedule.restart_in_window)
7116 .collect::<Vec<_>>(),
7117 vec![0, 1, 2, 3]
7118 );
7119 assert_eq!(
7120 schedules
7121 .iter()
7122 .map(|schedule| schedule.delay)
7123 .collect::<Vec<_>>(),
7124 vec![
7125 Duration::from_millis(100),
7126 Duration::from_secs(1),
7127 Duration::from_secs(10),
7128 Duration::from_secs(30),
7129 ]
7130 );
7131 }
7132
7133 #[test]
7134 fn crash_restart_backoff_resets_after_ring_clear() {
7135 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7136 let now = ClockInstant::now();
7137 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7138 assert_eq!(
7139 state.next_crash_restart(&policy, now).unwrap().delay,
7140 Duration::from_millis(100)
7141 );
7142 assert_eq!(
7143 state.next_crash_restart(&policy, now).unwrap().delay,
7144 Duration::from_secs(1)
7145 );
7146
7147 state.clear_crash_restarts();
7148 let schedule = state
7149 .next_crash_restart(&policy, now)
7150 .expect("a cleared ring must allow another restart");
7151 assert_eq!(schedule.restart_in_window, 0);
7152 assert_eq!(schedule.delay, Duration::from_millis(100));
7153 }
7154
7155 #[test]
7156 fn crash_restart_backoff_ignores_aged_restarts() {
7157 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7158 let now = ClockInstant::now();
7159 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7160 state
7161 .next_crash_restart(&policy, now)
7162 .expect("the first restart is allowed");
7163 state
7164 .next_crash_restart(&policy, now)
7165 .expect("the second restart is allowed");
7166 state.crash_restarts[0] = now
7167 .checked_sub(policy.window + Duration::from_secs(1))
7168 .expect("the fake clock can age a restart past the window");
7169
7170 let schedule = state
7171 .next_crash_restart(&policy, now)
7172 .expect("an aged restart must release its slot");
7173 assert_eq!(schedule.restart_in_window, 1);
7174 assert_eq!(schedule.delay, Duration::from_secs(1));
7175 assert_eq!(state.crash_restarts.len(), 2);
7176 }
7177
7178 #[test]
7182 fn a_budget_spent_before_the_window_no_longer_refuses() {
7183 let policy = RestartPolicy::new(3, Duration::ZERO);
7184 let mut state = snapshot_with_restarts(true, 3);
7185 let now = ClockInstant::now();
7186 assert!(!daemon_will_restart(&mut state, &policy, now));
7187
7188 assert!(daemon_will_restart(
7189 &mut state,
7190 &policy,
7191 now + policy.window + Duration::from_secs(1)
7192 ));
7193 assert!(
7194 state.crash_restarts.is_empty(),
7195 "reading the budget must drop the instants that left the window"
7196 );
7197 }
7198
7199 fn module_with_recovery_snapshot(
7200 state: ModuleState,
7201 enabled: bool,
7202 restart_count: u32,
7203 ) -> SupervisedModule {
7204 let registry = Arc::new(Registry::default());
7205 let supervisor =
7206 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(3, Duration::ZERO));
7207 let module = supervisor
7208 .spawn(ModuleSpec {
7209 module_id: "recovery-snapshot".to_string(),
7210 program: fake_aft_stub_path(),
7211 args: Vec::new(),
7212 env: Vec::new(),
7213 reserved: false,
7214 reserved_prefixes: Vec::new(),
7215 protocol: ModuleProtocol::Subc,
7216 overlap: Default::default(),
7217 })
7218 .unwrap();
7219 update_snapshot(
7220 &module.inner.snapshot,
7221 Some("recovery-snapshot"),
7222 |snapshot| {
7223 snapshot.state = state;
7224 snapshot.enabled = enabled;
7225 seed_crash_restarts(snapshot, restart_count);
7226 },
7227 )
7228 .unwrap();
7229 module
7230 }
7231
7232 #[cfg(target_os = "linux")]
7233 #[tokio::test]
7234 async fn no_cgroup_placement_does_not_block_fake_aft_stub_spawn() {
7235 let supervisor = Supervisor::new(Arc::new(Registry::default()), RestartPolicy::default())
7236 .with_cgroup_placement(None);
7237 let result = supervisor.spawn(ModuleSpec {
7238 module_id: "no-cgroup-placement".to_string(),
7239 program: fake_aft_stub_path(),
7240 args: Vec::new(),
7241 env: Vec::new(),
7242 reserved: false,
7243 reserved_prefixes: Vec::new(),
7244 protocol: ModuleProtocol::Subc,
7245 overlap: Default::default(),
7246 });
7247
7248 assert!(
7249 result.is_ok(),
7250 "no delegation must not turn an otherwise valid spawn into a failure: {result:?}"
7251 );
7252 }
7253
7254 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7255 async fn undecided_snapshot_uses_shared_restart_predicate() {
7256 assert!(module_with_recovery_snapshot(ModuleState::Running, true, 2)
7257 .will_recover_after_connection_loss()
7258 .unwrap());
7259 assert!(
7260 !module_with_recovery_snapshot(ModuleState::Running, true, 3)
7261 .will_recover_after_connection_loss()
7262 .unwrap()
7263 );
7264 }
7265
7266 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7267 async fn restarting_snapshot_at_exhausted_budget_is_non_terminal() {
7268 assert!(
7269 module_with_recovery_snapshot(ModuleState::Restarting, true, 3)
7270 .will_recover_after_connection_loss()
7271 .unwrap()
7272 );
7273 }
7274
7275 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7276 async fn terminal_phase_snapshots_are_terminal_before_budget_exhaustion() {
7277 assert!(!module_with_recovery_snapshot(ModuleState::Failed, true, 0)
7278 .will_recover_after_connection_loss()
7279 .unwrap());
7280 assert!(
7281 !module_with_recovery_snapshot(ModuleState::Disabled, true, 0)
7282 .will_recover_after_connection_loss()
7283 .unwrap()
7284 );
7285 }
7286
7287 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7288 async fn warming_snapshot_is_limited_to_startup_phases() {
7289 for state in [
7290 ModuleState::Starting,
7291 ModuleState::Running,
7292 ModuleState::Restarting,
7293 ] {
7294 assert!(
7295 module_with_recovery_snapshot(state, true, 0)
7296 .is_warming()
7297 .unwrap(),
7298 "{state:?} should be warming"
7299 );
7300 }
7301 for state in [
7302 ModuleState::Unresponsive,
7303 ModuleState::Draining,
7304 ModuleState::Stopped,
7305 ModuleState::Failed,
7306 ModuleState::Disabled,
7307 ] {
7308 assert!(
7309 !module_with_recovery_snapshot(state, true, 0)
7310 .is_warming()
7311 .unwrap(),
7312 "{state:?} should not be warming"
7313 );
7314 }
7315 }
7316
7317 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7318 async fn terminal_history_survives_respawn_and_keeps_both_crashes_in_order() {
7319 let registry = Arc::new(Registry::default());
7320 let supervisor =
7321 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(1, Duration::ZERO));
7322 let module = supervisor
7323 .spawn(ModuleSpec {
7324 module_id: "terminal-history".to_string(),
7325 program: fake_aft_stub_path(),
7326 args: Vec::new(),
7327 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7328 reserved: false,
7329 reserved_prefixes: Vec::new(),
7330 protocol: ModuleProtocol::Subc,
7331 overlap: Default::default(),
7332 })
7333 .unwrap();
7334
7335 let deadline = Instant::now() + Duration::from_secs(5);
7336 loop {
7337 let history = module.terminal_history();
7338 if history.entries.len() == 2 {
7339 assert_eq!(module.status().unwrap().state, ModuleState::Failed);
7340 assert_eq!(history.dropped, 0);
7341 assert_eq!(
7342 history
7343 .entries
7344 .iter()
7345 .map(|entry| entry.exit_code)
7346 .collect::<Vec<_>>(),
7347 vec![Some(23), Some(23)]
7348 );
7349 assert!(history.entries[0].at_ms <= history.entries[1].at_ms);
7350 return;
7351 }
7352 assert!(
7353 Instant::now() < deadline,
7354 "module did not retain two terminal exits: {history:?}"
7355 );
7356 sleep(Duration::from_millis(10)).await;
7357 }
7358 }
7359
7360 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7364 async fn disable_during_crash_backoff_cancels_pending_respawn() {
7365 let backoff = Duration::from_secs(2);
7366 let supervisor = Supervisor::new(
7367 Arc::new(Registry::default()),
7368 RestartPolicy::new(10, backoff),
7369 );
7370 let module = supervisor
7371 .spawn(ModuleSpec {
7372 module_id: "disable-during-backoff".to_string(),
7373 program: fake_aft_stub_path(),
7374 args: Vec::new(),
7375 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7376 reserved: false,
7377 reserved_prefixes: Vec::new(),
7378 protocol: ModuleProtocol::Subc,
7379 overlap: Default::default(),
7380 })
7381 .unwrap();
7382
7383 let deadline = Instant::now() + Duration::from_secs(5);
7385 loop {
7386 if module.status().unwrap().state == ModuleState::Restarting {
7387 break;
7388 }
7389 assert!(
7390 Instant::now() < deadline,
7391 "module never entered the crash backoff"
7392 );
7393 sleep(Duration::from_millis(10)).await;
7394 }
7395
7396 let started = Instant::now();
7397 module.set_enabled(false).await.unwrap();
7398 let waited = started.elapsed();
7399
7400 assert!(
7401 waited < backoff / 2,
7402 "disable waited {waited:?} behind the {backoff:?} crash backoff; the operator command must preempt the pending respawn"
7403 );
7404 assert_eq!(module.status().unwrap().state, ModuleState::Disabled);
7405
7406 sleep(backoff + Duration::from_millis(500)).await;
7408 let status = module.status().unwrap();
7409 assert_eq!(status.state, ModuleState::Disabled);
7410 assert_eq!(
7411 status.spawn_generation, 1,
7412 "module respawned after the operator disabled it"
7413 );
7414 }
7415
7416 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7420 async fn every_restart_increment_path_advances_lifetime_count() {
7421 let supervisor = Supervisor::new(
7422 Arc::new(Registry::default()),
7423 RestartPolicy::new(1, Duration::ZERO),
7424 );
7425 let runtime = supervisor.runtime_config();
7426 let spec = ModuleSpec {
7427 module_id: "lifetime-increment-path".to_string(),
7428 program: PathBuf::from("/unused/lifetime-increment-path"),
7429 args: Vec::new(),
7430 env: Vec::new(),
7431 reserved: false,
7432 reserved_prefixes: Vec::new(),
7433 protocol: ModuleProtocol::Subc,
7434 overlap: Default::default(),
7435 };
7436
7437 let crash_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7438 assert!(matches!(
7439 on_child_exit(
7440 &spec,
7441 runtime.restart_policy,
7442 &supervisor.registry,
7443 &crash_snapshot,
7444 &runtime.terminal_ring,
7445 &runtime.spawn_events,
7446 ExitReport {
7447 kind: ExitKind::Crash,
7448 code: Some(1),
7449 signal: None,
7450 at_ms: 1,
7451 },
7452 )
7453 .await,
7454 NextAction::Restart { schedule: _ }
7455 ));
7456 let (crash_restarts, crash_lifetime) = {
7457 let state = lock_snapshot(&crash_snapshot).unwrap();
7458 (state.crash_restarts.len(), state.lifetime_restarts)
7459 };
7460 assert_eq!(crash_restarts, 1);
7461 assert_eq!(crash_lifetime, 1);
7462
7463 let health_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7464 let mut health_child = None;
7465 assert!(matches!(
7466 health_restart_child(
7467 &spec,
7468 &runtime,
7469 &supervisor.registry,
7470 &supervisor.process_liveness,
7471 &health_snapshot,
7472 &mut health_child,
7473 SupervisorHealthStatus::Failing,
7474 None,
7475 2,
7476 )
7477 .await,
7478 Err(SuperviseError::Spawn { .. })
7479 ));
7480 let (health_restarts, health_lifetime) = {
7481 let state = lock_snapshot(&health_snapshot).unwrap();
7482 (state.crash_restarts.len(), state.lifetime_restarts)
7483 };
7484 assert_eq!(health_restarts, 1);
7485 assert_eq!(health_lifetime, 1);
7486
7487 let reload_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7488 let mut reload_child = None;
7489 assert!(matches!(
7490 handle_reload_spawn_failure(
7491 &spec,
7492 &runtime,
7493 &supervisor.process_liveness,
7494 &reload_snapshot,
7495 &mut reload_child,
7496 "forced reload spawn failure".to_string(),
7497 )
7498 .await,
7499 Err(SuperviseError::ReloadFailed { .. })
7500 ));
7501 let (reload_restarts, reload_lifetime) = {
7502 let state = lock_snapshot(&reload_snapshot).unwrap();
7503 (state.crash_restarts.len(), state.lifetime_restarts)
7504 };
7505 assert_eq!(reload_restarts, 1);
7506 assert_eq!(reload_lifetime, 1);
7507 }
7508
7509 #[tokio::test]
7510 async fn deliberately_severed_live_child_records_lifetime_without_spending_restart_budget() {
7511 let supervisor = Supervisor::new(
7512 Arc::new(Registry::default()),
7513 RestartPolicy::new(3, Duration::ZERO),
7514 );
7515 let runtime = supervisor.runtime_config();
7516 let spec = ModuleSpec {
7517 module_id: "deliberately-severed".to_string(),
7518 program: PathBuf::from("/unused/deliberately-severed"),
7519 args: Vec::new(),
7520 env: Vec::new(),
7521 reserved: false,
7522 reserved_prefixes: Vec::new(),
7523 protocol: ModuleProtocol::Subc,
7524 overlap: Default::default(),
7525 };
7526 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7527 let process = ProcessIdentity {
7528 pid: 41,
7529 start_time: 101,
7530 };
7531 record_deliberate_severance(&snapshot, process).unwrap();
7532 let exit_report = apply_deliberate_severance_marker(
7533 &snapshot,
7534 Some(process),
7535 ExitReport {
7536 kind: ExitKind::Crash,
7537 code: Some(1),
7538 signal: None,
7539 at_ms: 1,
7540 },
7541 );
7542 assert_eq!(exit_report.kind, ExitKind::DeliberateSeverance);
7543
7544 assert!(matches!(
7545 on_child_exit(
7546 &spec,
7547 runtime.restart_policy,
7548 &supervisor.registry,
7549 &snapshot,
7550 &runtime.terminal_ring,
7551 &runtime.spawn_events,
7552 exit_report,
7553 )
7554 .await,
7555 NextAction::Restart { schedule: _ }
7556 ));
7557 let state = lock_snapshot(&snapshot).unwrap();
7558 assert_eq!(state.lifetime_restarts, 1);
7559 assert_eq!(state.crash_restarts.len(), 0);
7560 }
7561
7562 #[tokio::test]
7563 async fn genuine_crash_spends_restart_budget_and_records_lifetime() {
7564 let supervisor = Supervisor::new(
7565 Arc::new(Registry::default()),
7566 RestartPolicy::new(3, Duration::ZERO),
7567 );
7568 let runtime = supervisor.runtime_config();
7569 let spec = ModuleSpec {
7570 module_id: "genuine-crash".to_string(),
7571 program: PathBuf::from("/unused/genuine-crash"),
7572 args: Vec::new(),
7573 env: Vec::new(),
7574 reserved: false,
7575 reserved_prefixes: Vec::new(),
7576 protocol: ModuleProtocol::Subc,
7577 overlap: Default::default(),
7578 };
7579 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7580
7581 assert!(matches!(
7582 on_child_exit(
7583 &spec,
7584 runtime.restart_policy,
7585 &supervisor.registry,
7586 &snapshot,
7587 &runtime.terminal_ring,
7588 &runtime.spawn_events,
7589 ExitReport {
7590 kind: ExitKind::Crash,
7591 code: Some(1),
7592 signal: None,
7593 at_ms: 1,
7594 },
7595 )
7596 .await,
7597 NextAction::Restart { schedule: _ }
7598 ));
7599 let state = lock_snapshot(&snapshot).unwrap();
7600 assert_eq!(state.lifetime_restarts, 1);
7601 assert_eq!(state.crash_restarts.len(), 1);
7602 }
7603
7604 fn crash_exit_report(at_ms: u64) -> ExitReport {
7605 ExitReport {
7606 kind: ExitKind::Crash,
7607 code: Some(1),
7608 signal: None,
7609 at_ms,
7610 }
7611 }
7612
7613 fn windowed_crash_spec(module_id: &str) -> ModuleSpec {
7614 ModuleSpec {
7615 module_id: module_id.to_string(),
7616 program: PathBuf::from("/unused").join(module_id),
7617 args: Vec::new(),
7618 env: Vec::new(),
7619 reserved: false,
7620 reserved_prefixes: Vec::new(),
7621 protocol: ModuleProtocol::Subc,
7622 overlap: Default::default(),
7623 }
7624 }
7625
7626 #[tokio::test]
7632 async fn three_crashes_inside_the_window_stop_the_module_and_name_the_window() {
7633 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::ERROR);
7634 let supervisor = Supervisor::new(
7635 Arc::new(Registry::default()),
7636 RestartPolicy::new(2, Duration::ZERO),
7637 );
7638 let runtime = supervisor.runtime_config();
7639 let spec = windowed_crash_spec("crash-loop-in-window");
7640 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7641
7642 for attempt in 1..=2 {
7643 assert!(
7644 matches!(
7645 on_child_exit(
7646 &spec,
7647 runtime.restart_policy,
7648 &supervisor.registry,
7649 &snapshot,
7650 &runtime.terminal_ring,
7651 &runtime.spawn_events,
7652 crash_exit_report(attempt),
7653 )
7654 .await,
7655 NextAction::Restart { schedule: _ }
7656 ),
7657 "crash {attempt} is inside the budget and must respawn"
7658 );
7659 }
7660
7661 assert!(matches!(
7662 on_child_exit(
7663 &spec,
7664 runtime.restart_policy,
7665 &supervisor.registry,
7666 &snapshot,
7667 &runtime.terminal_ring,
7668 &runtime.spawn_events,
7669 crash_exit_report(3),
7670 )
7671 .await,
7672 NextAction::Stop { .. }
7673 ));
7674
7675 {
7676 let state = lock_snapshot(&snapshot).unwrap();
7677 assert_eq!(state.state, ModuleState::Failed);
7678 assert_eq!(state.crash_restarts.len(), 2);
7679 assert_eq!(state.lifetime_restarts, 2);
7680 }
7681
7682 let history = runtime
7683 .terminal_ring
7684 .lock()
7685 .expect("terminal ring is not poisoned")
7686 .snapshot();
7687 let last = history
7688 .entries
7689 .last()
7690 .expect("the refused crash is retained");
7691 assert_eq!(last.disposition, TerminalDisposition::Failed);
7692 assert_eq!(
7693 last.disposition_detail.as_deref(),
7694 Some("crash budget exhausted: max_restarts=2 within window_secs=600")
7695 );
7696
7697 let captured = crate::router::test_log::captured_logs(&logs);
7698 assert!(
7699 captured.contains("crash budget exhausted: max_restarts=2 within window_secs=600"),
7700 "the stop must be logged with its window: {captured}"
7701 );
7702 }
7703
7704 #[tokio::test]
7712 async fn a_crash_older_than_the_window_frees_its_slot_for_a_later_crash() {
7713 let supervisor = Supervisor::new(
7714 Arc::new(Registry::default()),
7715 RestartPolicy::new(2, Duration::ZERO),
7716 );
7717 let runtime = supervisor.runtime_config();
7718 let spec = windowed_crash_spec("crash-across-windows");
7719 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7720
7721 for attempt in 1..=2 {
7722 assert!(matches!(
7723 on_child_exit(
7724 &spec,
7725 runtime.restart_policy,
7726 &supervisor.registry,
7727 &snapshot,
7728 &runtime.terminal_ring,
7729 &runtime.spawn_events,
7730 crash_exit_report(attempt),
7731 )
7732 .await,
7733 NextAction::Restart { schedule: _ }
7734 ));
7735 }
7736
7737 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
7740 age_oldest_crash_restart_out_of_window(state, runtime.restart_policy.window);
7741 })
7742 .unwrap();
7743
7744 assert!(
7745 matches!(
7746 on_child_exit(
7747 &spec,
7748 runtime.restart_policy,
7749 &supervisor.registry,
7750 &snapshot,
7751 &runtime.terminal_ring,
7752 &runtime.spawn_events,
7753 crash_exit_report(3),
7754 )
7755 .await,
7756 NextAction::Restart { schedule: _ }
7757 ),
7758 "a crash older than the window must not hold a budget slot"
7759 );
7760
7761 let state = lock_snapshot(&snapshot).unwrap();
7762 assert_eq!(state.state, ModuleState::Restarting);
7763 assert_eq!(
7764 state.crash_restarts.len(),
7765 2,
7766 "the aged instant is dropped and the new one takes its place"
7767 );
7768 assert_eq!(
7769 state.lifetime_restarts, 3,
7770 "the ledger counts every restart, including the ones the window forgot"
7771 );
7772 }
7773
7774 #[tokio::test]
7779 async fn an_operator_restart_clears_the_ring_and_leaves_the_ledger_alone() {
7780 let supervisor = Supervisor::new(
7781 Arc::new(Registry::default()),
7782 RestartPolicy::new(2, Duration::ZERO),
7783 );
7784 let runtime = supervisor.runtime_config();
7785 let spec = windowed_crash_spec("operator-cleared-budget");
7786 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7787
7788 for attempt in 1..=2 {
7789 assert!(matches!(
7790 on_child_exit(
7791 &spec,
7792 runtime.restart_policy,
7793 &supervisor.registry,
7794 &snapshot,
7795 &runtime.terminal_ring,
7796 &runtime.spawn_events,
7797 crash_exit_report(attempt),
7798 )
7799 .await,
7800 NextAction::Restart { schedule: _ }
7801 ));
7802 }
7803
7804 reset_restart_count(&snapshot, &spec.module_id).unwrap();
7805 {
7806 let state = lock_snapshot(&snapshot).unwrap();
7807 assert!(
7808 state.crash_restarts.is_empty(),
7809 "an operator restart returns the full budget"
7810 );
7811 assert_eq!(
7812 state.lifetime_restarts, 2,
7813 "clearing the budget must not unmake the crashes"
7814 );
7815 }
7816
7817 assert!(
7818 matches!(
7819 on_child_exit(
7820 &spec,
7821 runtime.restart_policy,
7822 &supervisor.registry,
7823 &snapshot,
7824 &runtime.terminal_ring,
7825 &runtime.spawn_events,
7826 crash_exit_report(3),
7827 )
7828 .await,
7829 NextAction::Restart { schedule: _ }
7830 ),
7831 "the cleared budget must be spendable again"
7832 );
7833 let state = lock_snapshot(&snapshot).unwrap();
7834 assert_eq!(state.crash_restarts.len(), 1);
7835 assert_eq!(state.lifetime_restarts, 3);
7836 }
7837
7838 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7839 async fn severance_marker_for_a_dead_child_does_not_label_its_successor() {
7840 let severed = ProcessIdentity {
7841 pid: 41,
7842 start_time: 101,
7843 };
7844 let successor = ProcessIdentity {
7845 pid: 41,
7846 start_time: 202,
7847 };
7848 let module = module_with_recovery_snapshot(ModuleState::Running, true, 0);
7849 update_snapshot(&module.inner.snapshot, Some("recovery-snapshot"), |state| {
7850 state.pid = Some(successor.pid);
7851 state.process_start_time = Some(successor.start_time);
7852 })
7853 .unwrap();
7854 assert!(!module.record_deliberate_severance(severed).unwrap());
7855
7856 let exit_report = apply_deliberate_severance_marker(
7857 &module.inner.snapshot,
7858 Some(successor),
7859 ExitReport {
7860 kind: ExitKind::Crash,
7861 code: Some(1),
7862 signal: None,
7863 at_ms: 1,
7864 },
7865 );
7866
7867 assert_eq!(exit_report.kind, ExitKind::Crash);
7868 }
7869
7870 #[tokio::test]
7871 async fn drain_reap_marks_deliberate_severance_and_records_lifetime_without_budget() {
7872 let registry = Registry::default();
7873 let supervisor = Supervisor::new(
7874 Arc::new(Registry::default()),
7875 RestartPolicy::new(3, Duration::ZERO),
7876 );
7877 let runtime = supervisor.runtime_config();
7878 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7879 let spec = ModuleSpec {
7880 module_id: "drain-deliberate-severance".to_string(),
7881 program: fake_aft_stub_path(),
7882 args: Vec::new(),
7883 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7884 reserved: false,
7885 reserved_prefixes: Vec::new(),
7886 protocol: ModuleProtocol::Subc,
7887 overlap: Default::default(),
7888 };
7889 let mut child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
7890 let process = ProcessIdentity {
7891 pid: 41,
7892 start_time: 101,
7893 };
7894 child.process_identity = Some(process);
7895 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
7896 state.pid = Some(process.pid);
7897 state.process_start_time = Some(process.start_time);
7898 })
7899 .unwrap();
7900 record_deliberate_severance(&snapshot, process).unwrap();
7901
7902 drain_child_to_state(
7903 &spec.module_id,
7904 spec.protocol,
7905 ®istry,
7906 &snapshot,
7907 &runtime.terminal_ring,
7908 &runtime.spawn_events,
7909 child,
7910 Duration::from_secs(1),
7911 ModuleState::Stopped,
7912 Some(false),
7913 )
7914 .await
7915 .unwrap();
7916
7917 let state = lock_snapshot(&snapshot).unwrap();
7918 assert_eq!(
7919 state.last_exit.as_ref().map(|exit| exit.kind),
7920 Some(ExitKind::DeliberateSeverance)
7921 );
7922 assert_eq!(state.lifetime_restarts, 1);
7923 assert_eq!(state.crash_restarts.len(), 0);
7924 drop(state);
7925 let history = runtime.terminal_ring.lock().unwrap().snapshot();
7926 assert_eq!(
7927 history.entries[0].exit_kind,
7928 subc_control::TerminalExitKind::DeliberateSeverance
7929 );
7930 }
7931
7932 #[tokio::test]
7933 async fn ordinary_drain_reap_does_not_record_a_lifetime_restart() {
7934 let registry = Registry::default();
7935 let supervisor = Supervisor::new(
7936 Arc::new(Registry::default()),
7937 RestartPolicy::new(3, Duration::ZERO),
7938 );
7939 let runtime = supervisor.runtime_config();
7940 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7941 let spec = ModuleSpec {
7942 module_id: "ordinary-drain".to_string(),
7943 program: fake_aft_stub_path(),
7944 args: Vec::new(),
7945 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7946 reserved: false,
7947 reserved_prefixes: Vec::new(),
7948 protocol: ModuleProtocol::Subc,
7949 overlap: Default::default(),
7950 };
7951 let child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
7952
7953 drain_child_to_state(
7954 &spec.module_id,
7955 spec.protocol,
7956 ®istry,
7957 &snapshot,
7958 &runtime.terminal_ring,
7959 &runtime.spawn_events,
7960 child,
7961 Duration::from_secs(1),
7962 ModuleState::Stopped,
7963 Some(false),
7964 )
7965 .await
7966 .unwrap();
7967
7968 let state = lock_snapshot(&snapshot).unwrap();
7969 assert_eq!(
7970 state.last_exit.as_ref().map(|exit| exit.kind),
7971 Some(ExitKind::Crash)
7972 );
7973 assert_eq!(state.lifetime_restarts, 0);
7974 assert_eq!(state.crash_restarts.len(), 0);
7975 }
7976
7977 #[test]
7978 fn fatal_connection_teardown_cannot_arm_a_marker_for_a_surviving_process() {
7979 assert!(!include_str!("server.rs")
7985 .contains("router.record_deliberate_connection_severance(ctx.connection_id)"));
7986 }
7987
7988 #[test]
7995 fn drained_after_quiescence_wait_passes_ok_through_and_forces_false_on_err() {
7996 assert!(drained_after_quiescence_wait(&Ok(true)));
7997 assert!(!drained_after_quiescence_wait(&Ok(false)));
7998 assert!(!drained_after_quiescence_wait(&Err(
7999 SuperviseError::StatePoisoned { module_id: None }
8000 )));
8001 }
8002
8003 #[test]
8012 fn wait_error_exit_report_records_a_failed_terminal_with_no_code_or_signal() {
8013 let ring = Arc::new(Mutex::new(TerminalRing::new(
8014 TerminalRingConfig::default(),
8015 0,
8016 )));
8017 record_wait_error_terminal("wait-error", &ring, &super::SpawnEventFeed::default());
8018
8019 let snapshot = ring.lock().unwrap().snapshot();
8020 assert_eq!(snapshot.entries.len(), 1);
8021 let entry = &snapshot.entries[0];
8022 assert_eq!(entry.exit_code, None);
8023 assert_eq!(entry.exit_signal, None);
8024 assert_eq!(entry.disposition, TerminalDisposition::Failed);
8025 }
8026
8027 #[test]
8028 fn wait_error_exit_path_preserves_spawn_event_density() {
8029 let feed = super::SpawnEventFeed::default();
8030 feed.configure_incarnation("wait-error-density".to_string());
8031 feed.emit_spawned("wait-error", 41, 1);
8032 let ring = Arc::new(Mutex::new(TerminalRing::new(
8033 TerminalRingConfig::default(),
8034 0,
8035 )));
8036
8037 record_wait_error_terminal("wait-error", &ring, &feed);
8038 feed.emit_spawned("after-wait-error", 42, 2);
8039
8040 let state = feed.0.lock().unwrap();
8041 let sequences = state
8042 .events
8043 .iter()
8044 .map(|event| event.cursor.seq)
8045 .collect::<Vec<_>>();
8046 assert_eq!(sequences, vec![1, 2, 3]);
8047 assert_eq!(state.events[1].kind, SpawnEventKind::Exited);
8048 assert_eq!(state.events[1].exit_code, None);
8049 assert_eq!(state.events[1].exit_signal, None);
8050 }
8051
8052 #[test]
8056 fn wait_error_exit_report_is_classified_as_a_crash() {
8057 assert_eq!(wait_error_exit_report().kind, ExitKind::Crash);
8058 }
8059}
8060
8061#[cfg(test)]
8062mod health_evidence_tests {
8063 use super::{HealthProbeError, HealthProbeEvidence};
8064 use std::collections::HashSet;
8065
8066 #[test]
8074 fn only_a_dead_lane_is_proof_of_death() {
8075 assert!(HealthProbeError::lane_dead("gone").is_proof_of_death());
8076 assert!(!HealthProbeError::no_answer("timed out").is_proof_of_death());
8080 assert!(!HealthProbeError::bad_answer("garbage").is_proof_of_death());
8081 assert!(!HealthProbeError::misconfigured("no table").is_proof_of_death());
8082 }
8083
8084 #[test]
8090 fn every_evidence_class_has_a_distinct_label() {
8091 let labels = [
8092 HealthProbeError::lane_dead("").label(),
8093 HealthProbeError::no_answer("").label(),
8094 HealthProbeError::bad_answer("").label(),
8095 HealthProbeError::misconfigured("").label(),
8096 ];
8097 let unique: HashSet<_> = labels.iter().collect();
8098 assert_eq!(unique.len(), labels.len(), "labels collided: {labels:?}");
8099 }
8100
8101 #[test]
8107 fn classification_preserves_the_original_message() {
8108 let err = HealthProbeError::no_answer("module did not answer within 5s");
8109 assert_eq!(err.to_string(), "module did not answer within 5s");
8110 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8111 }
8112}
8113
8114#[cfg(test)]
8115mod health_tombstone_tests {
8116 use std::{path::PathBuf, sync::Arc, time::Duration};
8117
8118 use subc_protocol::{
8119 manifest::Concurrency,
8120 session::{HealthStatus, ModuleControlResponse},
8121 };
8122 use tokio::sync::mpsc;
8123
8124 use super::{
8125 probe_module_health, HealthAction, HealthConfig, HealthProbeEvidence, ModuleProtocol,
8126 ModuleSpec, RestartPolicy, Supervisor, SupervisorRuntimeConfig,
8127 };
8128 use crate::{
8129 control::ControlHandler,
8130 forwarding::{ForwardingTable, ModuleControlRpcCompletion, ModuleControlRpcOutcome},
8131 registry::{ConnectionId, Registry},
8132 router::FrameSink,
8133 };
8134
8135 struct ProbeHarness {
8136 spec: ModuleSpec,
8137 runtime: SupervisorRuntimeConfig,
8138 forwarding: Arc<ForwardingTable>,
8139 module_connection: ConnectionId,
8140 module_rx: mpsc::Receiver<crate::router::OutboundFrame>,
8141 handler: ControlHandler,
8142 module: super::SupervisedModule,
8143 }
8144
8145 fn probe_harness() -> ProbeHarness {
8146 let registry = Arc::new(Registry::default());
8147 let forwarding = Arc::new(ForwardingTable::default());
8148 let supervisor_handle = super::SupervisorHandle::new();
8149 let health = HealthConfig {
8150 cadence: Duration::from_secs(30),
8151 deadline: Duration::from_secs(5),
8152 failure_threshold: 3,
8153 on_degraded: HealthAction::Report,
8154 on_failing: HealthAction::Report,
8155 critical: false,
8156 };
8157 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
8158 .with_forwarding(Arc::clone(&forwarding))
8159 .with_handle(supervisor_handle.clone())
8160 .with_health_config(health);
8161 let spec = ModuleSpec {
8162 module_id: "late-health-module".to_string(),
8163 program: PathBuf::from("disabled-module"),
8164 args: Vec::new(),
8165 env: Vec::new(),
8166 reserved: false,
8167 reserved_prefixes: Vec::new(),
8168 protocol: ModuleProtocol::Subc,
8169 overlap: Default::default(),
8170 };
8171 let module = supervisor
8172 .supervise_configured(spec.clone(), false)
8173 .unwrap();
8174 let runtime = supervisor.runtime_config();
8175 let handler = ControlHandler::with_forwarding(registry, Arc::clone(&forwarding))
8176 .with_supervisor(supervisor_handle);
8177 let module_connection = ConnectionId::new(700);
8178 let (module_tx, module_rx) = mpsc::channel(8);
8179 forwarding
8180 .register_module_connection(
8181 module_connection,
8182 spec.module_id.clone(),
8183 subc_protocol::PROTOCOL_VERSION,
8184 Concurrency::ModuleManaged,
8185 FrameSink::new(module_tx),
8186 )
8187 .unwrap();
8188
8189 ProbeHarness {
8190 spec,
8191 runtime,
8192 forwarding,
8193 module_connection,
8194 module_rx,
8195 handler,
8196 module,
8197 }
8198 }
8199
8200 async fn finish_after(
8201 harness: &mut ProbeHarness,
8202 stall: Duration,
8203 ) -> ModuleControlRpcCompletion {
8204 assert!(stall > harness.runtime.health.deadline);
8205 let deadline = harness.runtime.health.deadline;
8206 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8207 let answer = async {
8208 let frame = harness.module_rx.recv().await.expect("health.check frame");
8209 tokio::time::advance(deadline).await;
8210 tokio::task::yield_now().await;
8211 tokio::time::advance(stall - deadline).await;
8212 harness
8213 .forwarding
8214 .complete_module_control_rpc(
8215 harness.module_connection,
8216 frame.header.corr,
8217 Some("health.check"),
8218 ModuleControlRpcOutcome::Response(ModuleControlResponse::HealthCheck {
8219 status: HealthStatus::Ok,
8220 detail: None,
8221 metrics: None,
8222 }),
8223 )
8224 .unwrap()
8225 };
8226 let (probe_result, completion) = tokio::join!(probe, answer);
8227 let err = probe_result.expect_err("probe must miss its deadline");
8228 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8229 completion
8230 }
8231
8232 async fn time_out_without_answer(harness: &mut ProbeHarness) {
8233 let deadline = harness.runtime.health.deadline;
8234 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8235 let exhaust_deadline = async {
8236 let _frame = harness.module_rx.recv().await.expect("health.check frame");
8237 tokio::time::advance(deadline).await;
8238 tokio::task::yield_now().await;
8239 };
8240 let (probe_result, ()) = tokio::join!(probe, exhaust_deadline);
8241 let err = probe_result.expect_err("probe must miss its deadline");
8242 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8243 }
8244
8245 #[tokio::test(start_paused = true)]
8246 async fn late_health_answers_record_start_anchored_latency_for_two_stalls() {
8247 let mut harness = probe_harness();
8248
8249 let first = finish_after(&mut harness, Duration::from_secs(8)).await;
8250 let first_latency = match &first {
8251 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8252 other => panic!("late answer was not retained: {other:?}"),
8253 };
8254 assert!(harness.handler.observe_module_control_completion(first));
8255
8256 let second = finish_after(&mut harness, Duration::from_secs(11)).await;
8257 let second_latency = match &second {
8258 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8259 other => panic!("late answer was not retained: {other:?}"),
8260 };
8261 assert!(harness.handler.observe_module_control_completion(second));
8262
8263 assert_eq!(first_latency, Duration::from_secs(8));
8264 assert_eq!(
8265 second_latency - first_latency,
8266 Duration::from_secs(3),
8267 "latency must grow linearly with the additional stall"
8268 );
8269 let health = harness.module.status().unwrap().health;
8270 assert_eq!(health.late_answer_count, 2);
8271 assert_eq!(health.last_late_answer_latency_ms, Some(11_000));
8272 }
8273
8274 #[tokio::test(start_paused = true)]
8282 async fn late_answer_clears_the_consecutive_failure_streak() {
8283 let mut harness = probe_harness();
8284
8285 time_out_without_answer(&mut harness).await;
8287 harness
8288 .module
8289 .record_health_probe_failure_for_test("[no-answer] test miss")
8290 .unwrap();
8291 assert_eq!(
8292 harness.module.status().unwrap().health.consecutive_failures,
8293 1,
8294 "precondition: the miss must be on the streak before the late answer"
8295 );
8296
8297 let late = finish_after(&mut harness, Duration::from_secs(9)).await;
8299 assert!(matches!(
8300 late,
8301 ModuleControlRpcCompletion::LateHealthAnswer { .. }
8302 ));
8303 assert!(harness.handler.observe_module_control_completion(late));
8304
8305 let health = harness.module.status().unwrap().health;
8306 assert_eq!(
8307 health.consecutive_failures, 0,
8308 "a late answer is an answer: the streak must reset"
8309 );
8310 assert_eq!(health.late_answer_count, 1);
8311 }
8312
8313 #[tokio::test(start_paused = true)]
8314 async fn repeated_serial_probe_cycles_keep_one_tombstone_per_endpoint() {
8315 let mut harness = probe_harness();
8316
8317 for _ in 0..20 {
8318 time_out_without_answer(&mut harness).await;
8319 assert_eq!(
8320 harness.forwarding.health_probe_tombstone_count().unwrap(),
8321 1
8322 );
8323 }
8324 }
8325}
8326
8327#[cfg(test)]
8328mod child_env_tests {
8329 use super::{
8330 apply_child_env, apply_spawn_role, apply_wire_spawn_args, ModuleProtocol, ModuleSpec,
8331 SpawnRole, SupervisorHandle, SPAWN_ROLE_SWAP_CANDIDATE, SUBC_ARG, SUBC_LAUNCH_NONCE_ENV,
8332 SUBC_MODULE_ID_ENV, SUBC_SPAWN_ROLE_ENV,
8333 };
8334 use std::{ffi::OsStr, path::PathBuf};
8335 use tokio::process::Command;
8336
8337 fn spec(env: Vec<(String, String)>) -> ModuleSpec {
8338 ModuleSpec {
8339 module_id: "env-plan".to_string(),
8340 program: PathBuf::from("/nonexistent"),
8341 args: Vec::new(),
8342 env,
8343 reserved: false,
8344 reserved_prefixes: Vec::new(),
8345 protocol: ModuleProtocol::Subc,
8346 overlap: Default::default(),
8347 }
8348 }
8349
8350 #[test]
8364 fn ambient_ck_log_is_removed_and_a_configured_one_survives() {
8365 let mut command = Command::new("/nonexistent");
8366 apply_child_env(&mut command, &spec(Vec::new()));
8367 let removed = command
8368 .as_std()
8369 .get_envs()
8370 .any(|(key, value)| key == OsStr::new("CK_LOG") && value.is_none());
8371 assert!(
8372 removed,
8373 "ambient CK_LOG must be explicitly removed for an unconfigured module"
8374 );
8375
8376 let mut configured = Command::new("/nonexistent");
8377 apply_child_env(
8378 &mut configured,
8379 &spec(vec![("CK_LOG".to_string(), "debug".to_string())]),
8380 );
8381 let effective = configured
8382 .as_std()
8383 .get_envs()
8384 .filter(|(key, _)| *key == OsStr::new("CK_LOG"))
8385 .last()
8386 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()));
8387 assert_eq!(
8388 effective,
8389 Some(Some("debug".to_string())),
8390 "a module's configured CK_LOG must survive the ambient removal"
8391 );
8392 }
8393
8394 #[test]
8403 fn protocol_none_spawn_carries_no_subc_argument_and_no_nonce() {
8404 let connection_file = std::path::Path::new("/run/subc-connection.json");
8405 let handle = SupervisorHandle::new();
8406
8407 let mut none_spec = spec(Vec::new());
8408 none_spec.protocol = ModuleProtocol::None;
8409 let mut none = Command::new("/nonexistent");
8410 apply_wire_spawn_args(&mut none, &none_spec, Some(connection_file), Some(&handle))
8411 .expect("protocol-none spawn args apply");
8412 let none_args: Vec<String> = none
8413 .as_std()
8414 .get_args()
8415 .map(|a| a.to_string_lossy().into_owned())
8416 .collect();
8417 assert!(
8418 !none_args.iter().any(|a| a == SUBC_ARG),
8419 "protocol:none argv must not carry --subc; got {none_args:?}"
8420 );
8421 let none_has_nonce = none
8422 .as_std()
8423 .get_envs()
8424 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some());
8425 assert!(
8426 !none_has_nonce,
8427 "protocol:none spawn must not receive a launch nonce"
8428 );
8429 let none_has_module_id = none
8430 .as_std()
8431 .get_envs()
8432 .any(|(key, value)| key == OsStr::new(SUBC_MODULE_ID_ENV) && value.is_some());
8433 assert!(
8434 none_has_module_id,
8435 "SUBC_MODULE_ID is inert and stays on every path"
8436 );
8437 assert!(
8438 handle.spawn_nonce(&none_spec.module_id).is_none(),
8439 "no nonce record for a process that will never present one"
8440 );
8441
8442 let wire_spec = spec(Vec::new());
8444 let mut wire = Command::new("/nonexistent");
8445 apply_wire_spawn_args(&mut wire, &wire_spec, Some(connection_file), Some(&handle))
8446 .expect("subc-wire spawn args apply");
8447 let wire_args: Vec<String> = wire
8448 .as_std()
8449 .get_args()
8450 .map(|a| a.to_string_lossy().into_owned())
8451 .collect();
8452 assert_eq!(
8453 wire_args,
8454 vec![
8455 SUBC_ARG.to_string(),
8456 connection_file.to_string_lossy().into_owned()
8457 ],
8458 "a subc-wire spawn still carries --subc <path>"
8459 );
8460 assert!(wire
8461 .as_std()
8462 .get_envs()
8463 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some()));
8464 assert!(handle.spawn_nonce(&wire_spec.module_id).is_some());
8465 }
8466
8467 #[test]
8477 fn plain_spawn_removes_the_spawn_role_even_when_the_spec_sets_it() {
8478 let role = |command: &Command| {
8479 command
8480 .as_std()
8481 .get_envs()
8482 .filter(|(key, _)| *key == OsStr::new(SUBC_SPAWN_ROLE_ENV))
8483 .last()
8484 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()))
8485 };
8486 let forged = spec(vec![(
8487 SUBC_SPAWN_ROLE_ENV.to_string(),
8488 SPAWN_ROLE_SWAP_CANDIDATE.to_string(),
8489 )]);
8490
8491 let mut plain = Command::new("/nonexistent");
8492 apply_child_env(&mut plain, &forged);
8493 apply_spawn_role(&mut plain, SpawnRole::Plain);
8494 assert_eq!(
8495 role(&plain),
8496 Some(None),
8497 "a plain spawn must remove SUBC_SPAWN_ROLE, whatever the spec says"
8498 );
8499
8500 let mut candidate = Command::new("/nonexistent");
8501 apply_child_env(&mut candidate, &spec(Vec::new()));
8502 apply_spawn_role(&mut candidate, SpawnRole::SwapCandidate);
8503 assert_eq!(
8504 role(&candidate),
8505 Some(Some(SPAWN_ROLE_SWAP_CANDIDATE.to_string()))
8506 );
8507 }
8508
8509 #[test]
8515 fn daemon_private_capture_keys_are_not_passed_to_the_child() {
8516 let mut command = Command::new("/nonexistent");
8517 apply_child_env(
8518 &mut command,
8519 &spec(vec![
8520 (super::CAPTURE_KEEP_ENV.to_string(), "5".to_string()),
8521 ("KEPT".to_string(), "yes".to_string()),
8522 ]),
8523 );
8524 let keys: Vec<String> = command
8525 .as_std()
8526 .get_envs()
8527 .filter(|(_, value)| value.is_some())
8528 .map(|(key, _)| key.to_string_lossy().into_owned())
8529 .collect();
8530 assert!(keys.contains(&"KEPT".to_string()), "got {keys:?}");
8531 assert!(
8532 !keys.contains(&super::CAPTURE_KEEP_ENV.to_string()),
8533 "daemon-private capture key leaked to the child: {keys:?}"
8534 );
8535 }
8536}
8537
8538#[cfg(test)]
8539mod jitter_tests {
8540 use super::jittered_health_delay;
8541 use std::{collections::HashSet, time::Duration};
8542
8543 const FLEET: [&str; 14] = [
8552 "aft",
8553 "alfonso-core",
8554 "magic-context",
8555 "broca",
8556 "thalamus",
8557 "quota",
8558 "engram",
8559 "plexus",
8560 "cerebellum",
8561 "astrocyte",
8562 "synapse",
8563 "subc-mcp",
8564 "cortexkit-credentials",
8565 "subc-federation",
8566 ];
8567
8568 #[test]
8576 fn probe_delays_disperse_across_the_fleet() {
8577 let cadence = Duration::from_secs(30);
8578 let delays: HashSet<Duration> = FLEET
8579 .iter()
8580 .map(|id| jittered_health_delay(id, 0, cadence))
8581 .collect();
8582 assert_eq!(
8583 delays.len(),
8584 FLEET.len(),
8585 "every supervised module must land on its own probe offset"
8586 );
8587 }
8588
8589 #[test]
8595 fn jitter_only_delays_and_stays_within_one_tenth_of_cadence() {
8596 let cadence = Duration::from_secs(30);
8597 let span = cadence / 10;
8598 for id in FLEET {
8599 for probe_index in 0..8 {
8600 let delay = jittered_health_delay(id, probe_index, cadence);
8601 assert!(
8602 delay >= cadence,
8603 "{id}#{probe_index}: jitter must not shorten the cadence"
8604 );
8605 assert!(
8606 delay < cadence + span,
8607 "{id}#{probe_index}: jitter must stay inside one tenth of the cadence"
8608 );
8609 }
8610 }
8611 }
8612
8613 #[test]
8619 fn a_module_offset_is_stable_across_restarts() {
8620 let cadence = Duration::from_secs(30);
8621 for id in FLEET {
8622 assert_eq!(
8623 jittered_health_delay(id, 0, cadence),
8624 jittered_health_delay(id, 0, cadence),
8625 "{id}: the same module and probe index must produce the same offset"
8626 );
8627 }
8628 }
8629
8630 #[test]
8632 fn zero_cadence_yields_zero_delay() {
8633 assert_eq!(
8634 jittered_health_delay("aft", 0, Duration::ZERO),
8635 Duration::ZERO
8636 );
8637 }
8638}
8639
8640#[cfg(all(test, target_os = "linux"))]
8641mod cgroup_placement_tests {
8642 use super::{
8643 apply_cgroup_placement, remove_module_cgroup, ModuleProtocol, ModuleSpec, SuperviseError,
8644 SupervisedChild,
8645 };
8646 use crate::{
8647 stderr_tail::{StderrRing, StderrTailConfig},
8648 test_support::TestTempDir,
8649 };
8650 use std::{
8651 fs, io,
8652 path::{Path, PathBuf},
8653 sync::{Arc, Mutex},
8654 };
8655 use tokio::process::Command;
8656
8657 #[test]
8658 fn failed_parent_cgroup_open_is_a_cgroup_supervision_error() {
8659 let path = Path::new("/definitely-missing-subc-cgroup");
8660 let mut command = Command::new("true");
8661 let error = apply_cgroup_placement(
8662 &mut command,
8663 &ModuleSpec {
8664 module_id: "broken-cgroup".to_string(),
8665 program: PathBuf::from("true"),
8666 args: Vec::new(),
8667 env: Vec::new(),
8668 reserved: false,
8669 reserved_prefixes: Vec::new(),
8670 protocol: ModuleProtocol::Subc,
8671 overlap: Default::default(),
8672 },
8673 path,
8674 )
8675 .expect_err("a parent cgroup open failure must reject the supervised spawn");
8676 let reason = error.to_string();
8677
8678 assert!(
8679 matches!(error, SuperviseError::Cgroup { .. }),
8680 "parent cgroup open must be reported as a cgroup supervision error: {reason}"
8681 );
8682 assert!(
8683 reason.contains("/definitely-missing-subc-cgroup/cgroup.procs"),
8684 "parent cgroup open failure must name cgroup.procs: {reason}"
8685 );
8686 }
8687
8688 #[tokio::test]
8689 async fn reaping_a_child_removes_its_empty_module_cgroup() {
8690 let root = TestTempDir::new("supervisor-reap-cgroup");
8691 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
8692 let placement = subc_cgroup::prepare_at(&root)
8693 .expect("prepare scratch cgroup root")
8694 .expect("scratch root has a cgroup.procs marker");
8695 let module_id = "reaped-module";
8696 let module = placement
8697 .module_path(module_id)
8698 .expect("create scratch module cgroup");
8699 let child = Command::new("true")
8700 .spawn()
8701 .expect("spawn short-lived child");
8702 let pid = child.id().expect("spawned child has pid");
8703 let mut child = SupervisedChild {
8704 child,
8705 module_id: module_id.to_string(),
8706 cgroup_placement: Some(placement),
8707 stdout_pump: None,
8708 stderr_pump: None,
8709 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
8710 spawned_at_ms: 0,
8711 spawned_from: PathBuf::from("true"),
8712 spawned_file_identity: None,
8713 process_start_time: None,
8714 process_identity: None,
8715 pid,
8716 roster_guard: None,
8717 };
8718
8719 child.wait().await.expect("reap short-lived child");
8720
8721 assert!(
8722 !module.exists(),
8723 "reaping the supervised child must remove its empty cgroup"
8724 );
8725 }
8726
8727 #[test]
8728 fn non_empty_cgroup_removal_is_reported_without_blocking_teardown() {
8729 let root = TestTempDir::new("supervisor-non-empty-cgroup");
8730 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
8731 let placement = subc_cgroup::prepare_at(&root)
8732 .expect("prepare scratch cgroup root")
8733 .expect("scratch root has a cgroup.procs marker");
8734 let module = placement
8735 .module_path("surviving-module")
8736 .expect("create scratch module cgroup");
8737 fs::write(module.join("surviving-process"), b"still present")
8738 .expect("make scratch cgroup non-empty");
8739 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::WARN);
8740
8741 remove_module_cgroup(&placement, "surviving-module");
8742
8743 let logs = crate::router::test_log::captured_logs(&logs);
8744 assert!(
8745 module.exists(),
8746 "failed removal must leave the cgroup intact"
8747 );
8748 assert!(
8749 logs.contains("could not remove module cgroup after process exit; continuing teardown")
8750 && logs.contains("surviving-module"),
8751 "best-effort removal must report the failure without returning it: {logs}"
8752 );
8753 }
8754
8755 #[test]
8756 fn cgroup_pre_exec_spawn_failure_names_the_cgroup_path() {
8757 let cgroup_path = PathBuf::from("/sys/fs/cgroup/subc-modules/broken-module");
8758 let reason = SuperviseError::Spawn {
8759 program: PathBuf::from("/bin/true"),
8760 source: io::Error::from_raw_os_error(13),
8761 cgroup_path: Some(cgroup_path.clone()),
8762 }
8763 .to_string();
8764
8765 assert!(
8766 reason.contains(&cgroup_path.display().to_string()),
8767 "a pre_exec spawn failure must name the cgroup path: {reason}"
8768 );
8769 }
8770}