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);
104pub const SPAWN_EVENT_RING_CAPACITY: usize = 4096;
106const SPAWN_SUBSCRIBER_BUFFER: usize = SPAWN_EVENT_RING_CAPACITY + 1;
107pub(crate) const SPAWN_SUBSCRIBER_LAGGED_CODE: &str = "spawn_subscriber_lagged";
112
113struct SupervisedChild {
114 child: Child,
115 #[cfg(target_os = "linux")]
118 module_id: String,
119 #[cfg(target_os = "linux")]
120 cgroup_placement: Option<subc_cgroup::Placement>,
121 #[cfg(windows)]
143 job: Option<subc_jobobject::JobObject>,
144 stdout_pump: Option<JoinHandle<()>>,
145 stderr_pump: Option<StderrPump>,
146 stderr_ring: Arc<Mutex<StderrRing>>,
147 spawned_at_ms: u64,
148 spawned_from: PathBuf,
149 spawned_file_identity: Option<SpawnedFileIdentity>,
150 process_start_time: Option<u64>,
151 process_identity: Option<ProcessIdentity>,
152 pid: u32,
153 roster_guard: Option<crate::child_roster::RosterGuard>,
156}
157
158impl SupervisedChild {
159 fn id(&self) -> Option<u32> {
160 Some(self.pid)
161 }
162
163 fn process_identity(&self) -> Option<ProcessIdentity> {
164 self.process_identity
165 }
166
167 async fn wait(&mut self) -> io::Result<ExitStatus> {
168 let result = self.child.wait().await;
176 #[cfg(target_os = "linux")]
177 if result.is_ok() {
178 if let Some(placement) = self.cgroup_placement.take() {
179 remove_module_cgroup(&placement, &self.module_id);
180 }
181 }
182 result
183 }
184
185 fn release_roster(&mut self) {
189 self.roster_guard = None;
190 }
191
192 fn start_kill(&mut self) -> io::Result<()> {
205 #[cfg(windows)]
206 if let Some(job) = &self.job {
207 if let Err(error) = job.terminate() {
208 debug!(
209 error = %error,
210 "job termination failed; the direct-child kill still owns the outcome"
211 );
212 }
213 }
214 self.child.start_kill()
215 }
216
217 async fn drain_stderr(&mut self, module_id: &str) {
218 if let Some(mut pump) = self.stdout_pump.take() {
219 match timeout(STDERR_PUMP_DRAIN_TIMEOUT, &mut pump).await {
220 Ok(Ok(())) => {}
221 Ok(Err(error)) => {
222 warn!(module_id, error = %error, "stdout pump ended unexpectedly");
223 }
224 Err(_) => {
225 pump.abort();
226 warn!(
227 module_id,
228 waited = ?STDERR_PUMP_DRAIN_TIMEOUT,
229 "stdout pump did not drain before restart; stopped it before the next process"
230 );
231 }
232 }
233 }
234
235 let Some(pump) = self.stderr_pump.take() else {
236 return;
237 };
238 settle_stderr_pump(
239 module_id,
240 &self.stderr_ring,
241 pump,
242 STDERR_PUMP_DRAIN_TIMEOUT,
243 )
244 .await;
245 }
246}
247
248struct StderrPump {
251 task: JoinHandle<()>,
252 generation: u64,
253}
254
255async fn settle_stderr_pump(
261 module_id: &str,
262 ring: &Arc<Mutex<StderrRing>>,
263 pump: StderrPump,
264 bound: Duration,
265) {
266 let lock = || ring.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
267 let StderrPump {
268 mut task,
269 generation,
270 } = pump;
271 lock().retire_pump(generation);
272 match timeout(bound, &mut task).await {
273 Ok(Ok(())) => {}
274 Ok(Err(err)) => {
275 let mut ring = lock();
276 ring.mark_incomplete(format!("stderr pump ended unexpectedly: {err}"));
277 ring.finish_pump(generation);
278 warn!(module_id, error = %err, "stderr pump ended before clean EOF");
279 }
280 Err(_) => {
281 drop(task);
283 lock().mark_pump_late(
284 generation,
285 format!(
286 "stderr of the exited process had not reached EOF {bound:?} after it was \
287 retired (a descendant may still hold the pipe open); lines it still \
288 writes are kept in that process's section"
289 ),
290 );
291 warn!(
292 module_id,
293 waited = ?bound,
294 "stderr pipe of the exited process is still open; its reader keeps running without delaying the restart"
295 );
296 }
297 }
298}
299
300fn registration_release_events() -> &'static watch::Sender<u64> {
301 static EVENTS: OnceLock<watch::Sender<u64>> = OnceLock::new();
302 EVENTS.get_or_init(|| {
303 let (sender, _receiver) = watch::channel(0);
304 sender
305 })
306}
307
308pub(crate) fn notify_registration_release() {
309 let events = registration_release_events();
310 let next_generation = (*events.borrow()).wrapping_add(1);
311 events.send_replace(next_generation);
312}
313
314#[derive(Debug, Clone, PartialEq, Eq)]
316pub struct ModuleSpec {
317 pub module_id: String,
318 pub program: PathBuf,
319 pub args: Vec<String>,
320 pub env: Vec<(String, String)>,
321 pub reserved: bool,
326 pub reserved_prefixes: Vec<String>,
331 pub protocol: ModuleProtocol,
347 pub overlap: ModuleOverlap,
352}
353
354#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
361pub enum ModuleOverlap {
362 #[default]
364 Exclusive,
365 Safe,
377}
378
379impl ModuleOverlap {
380 pub fn as_str(self) -> &'static str {
381 match self {
382 Self::Exclusive => "exclusive",
383 Self::Safe => "safe",
384 }
385 }
386}
387
388pub const SUBC_SPAWN_ROLE_ENV: &str = "SUBC_SPAWN_ROLE";
398pub const SPAWN_ROLE_SWAP_CANDIDATE: &str = "swap_candidate";
400pub const DEFAULT_SWAP_READY_TIMEOUT: Duration = Duration::from_secs(100);
405
406#[derive(Debug, Clone, Copy, PartialEq, Eq)]
424pub struct RestartPolicy {
425 pub max_restarts: u32,
426 pub backoff: Duration,
429 pub max_backoff: Duration,
431 pub window: Duration,
435}
436
437impl RestartPolicy {
438 pub fn new(max_restarts: u32, backoff: Duration) -> Self {
442 Self {
443 max_restarts,
444 backoff,
445 max_backoff: DEFAULT_MAX_BACKOFF,
446 window: DEFAULT_RESTART_WINDOW,
447 }
448 }
449
450 pub fn with_max_backoff(mut self, max_backoff: Duration) -> Self {
451 self.max_backoff = max_backoff;
452 self
453 }
454
455 pub fn with_window(mut self, window: Duration) -> Self {
456 self.window = window;
457 self
458 }
459
460 fn delay_for_restart(&self, restart_in_window: u32) -> Duration {
465 if self.backoff.is_zero() || self.max_backoff.is_zero() {
466 return Duration::ZERO;
467 }
468
469 let mut delay = self.backoff;
470 for _ in 0..restart_in_window {
471 if delay >= self.max_backoff {
472 return self.max_backoff;
473 }
474 delay = delay
475 .checked_mul(10)
476 .unwrap_or(self.max_backoff)
477 .min(self.max_backoff);
478 }
479 delay.min(self.max_backoff)
480 }
481
482 fn budget_exhausted_detail(&self) -> String {
487 format!(
488 "crash budget exhausted: max_restarts={} within window_secs={}",
489 self.max_restarts,
490 self.window.as_secs()
491 )
492 }
493}
494
495impl Default for RestartPolicy {
496 fn default() -> Self {
497 Self {
498 max_restarts: DEFAULT_MAX_RESTARTS,
499 backoff: DEFAULT_BACKOFF,
500 max_backoff: DEFAULT_MAX_BACKOFF,
501 window: DEFAULT_RESTART_WINDOW,
502 }
503 }
504}
505
506#[derive(Debug, Clone, Copy, PartialEq, Eq)]
507struct CrashRestartSchedule {
508 restart_in_window: u32,
509 delay: Duration,
510}
511
512fn daemon_will_restart(
519 state: &mut SupervisorSnapshot,
520 policy: &RestartPolicy,
521 now: Instant,
522) -> bool {
523 state.enabled && state.crash_restarts_in_window(policy.window, now) < policy.max_restarts
524}
525
526const DEFAULT_HEALTH_CADENCE: Duration = Duration::from_secs(30);
527const DEFAULT_HEALTH_DEADLINE: Duration = Duration::from_secs(5);
528const DEFAULT_HEALTH_FAILURE_THRESHOLD: u32 = 3;
529const MAX_HEALTH_METRICS_BYTES: usize = 16 * 1024;
530
531#[derive(Debug, Clone, Copy, PartialEq, Eq)]
532pub enum HealthAction {
533 Report,
534 Restart,
535 Alert,
536}
537
538impl fmt::Display for HealthAction {
539 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
540 f.write_str(match self {
541 Self::Report => "report",
542 Self::Restart => "restart",
543 Self::Alert => "alert",
544 })
545 }
546}
547
548#[derive(Debug, Clone, Copy, PartialEq, Eq)]
549pub struct HealthConfig {
550 pub cadence: Duration,
551 pub deadline: Duration,
552 pub failure_threshold: u32,
553 pub on_degraded: HealthAction,
554 pub on_failing: HealthAction,
555 pub critical: bool,
556}
557
558impl Default for HealthConfig {
559 fn default() -> Self {
560 Self {
561 cadence: DEFAULT_HEALTH_CADENCE,
562 deadline: DEFAULT_HEALTH_DEADLINE,
563 failure_threshold: DEFAULT_HEALTH_FAILURE_THRESHOLD,
564 on_degraded: HealthAction::Report,
565 on_failing: HealthAction::Report,
566 critical: false,
567 }
568 }
569}
570
571#[derive(Debug, Clone, PartialEq)]
589pub struct ModuleHealthStatus {
590 pub status: SupervisorHealthStatus,
591 pub last_probe_ms: Option<u64>,
592 pub detail: Option<String>,
593 pub metrics: Option<Value>,
594 pub consecutive_failures: u32,
595 pub late_answer_count: u64,
598 pub last_late_answer_latency_ms: Option<u64>,
600 pub last_action: Option<String>,
601 pub last_action_ms: Option<u64>,
605}
606
607impl Default for ModuleHealthStatus {
608 fn default() -> Self {
609 Self {
610 status: SupervisorHealthStatus::Unknown,
611 last_probe_ms: None,
612 detail: None,
613 metrics: None,
614 consecutive_failures: 0,
615 late_answer_count: 0,
616 last_late_answer_latency_ms: None,
617 last_action: None,
618 last_action_ms: None,
619 }
620 }
621}
622
623#[derive(Debug, Clone, Copy, PartialEq, Eq)]
625pub enum ModuleState {
626 Starting,
627 Running,
628 Unresponsive,
629 Restarting,
630 Draining,
631 Stopped,
632 Failed,
633 Disabled,
634}
635
636impl fmt::Display for ModuleState {
637 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
638 f.write_str(match self {
639 Self::Starting => "starting",
640 Self::Running => "running",
641 Self::Unresponsive => "unresponsive",
642 Self::Restarting => "restarting",
643 Self::Draining => "draining",
644 Self::Stopped => "stopped",
645 Self::Failed => "failed",
646 Self::Disabled => "disabled",
647 })
648 }
649}
650
651#[derive(Debug, Clone, Copy, PartialEq, Eq)]
653pub enum ExitKind {
654 Clean,
655 Crash,
656 DeliberateSeverance,
657}
658
659impl From<ExitKind> for TerminalExitKind {
660 fn from(kind: ExitKind) -> Self {
661 match kind {
662 ExitKind::Clean => Self::Clean,
663 ExitKind::Crash => Self::Crash,
664 ExitKind::DeliberateSeverance => Self::DeliberateSeverance,
665 }
666 }
667}
668
669#[derive(Debug, Clone, Copy, PartialEq, Eq)]
672pub(crate) struct ProcessIdentity {
673 pub(crate) pid: u32,
674 pub(crate) start_time: u64,
675}
676
677#[derive(Debug, Clone, PartialEq, Eq)]
679pub struct ExitReport {
680 pub kind: ExitKind,
681 pub code: Option<i32>,
682 pub signal: Option<i32>,
683 pub at_ms: u64,
684}
685
686#[derive(Debug, Clone, PartialEq)]
689pub struct ModuleStatus {
690 pub module_id: String,
691 pub state: ModuleState,
692 pub enabled: bool,
693 pub process_alive: bool,
694 pub registration_active: bool,
695 pub protocol: ModuleProtocol,
698 pub live: bool,
709 pub restart_count: u32,
713 pub lifetime_restarts: u32,
717 pub spawn_generation: u64,
718 pub max_restarts: u32,
723 pub restart_window: Duration,
727 pub drain_timeout: Duration,
731 pub restart_backoff: Duration,
732 pub restart_max_backoff: Duration,
733 pub pid: Option<u32>,
734 pub spawned_at_ms: Option<u64>,
735 pub spawned_from: Option<PathBuf>,
736 pub process_start_time: Option<u64>,
737 pub last_exit: Option<ExitReport>,
738 pub health: ModuleHealthStatus,
739}
740
741#[derive(Debug, Clone, PartialEq)]
742struct SupervisorSnapshot {
743 state: ModuleState,
744 enabled: bool,
745 process_alive: bool,
746 crash_restarts: VecDeque<Instant>,
752 lifetime_restarts: u32,
753 spawn_generation: u64,
762 pid: Option<u32>,
763 spawned_at_ms: Option<u64>,
764 spawned_from: Option<PathBuf>,
765 spawned_file_identity: Option<SpawnedFileIdentity>,
766 process_start_time: Option<u64>,
767 deliberate_severance: Option<ProcessIdentity>,
768 last_exit: Option<ExitReport>,
769 health: ModuleHealthStatus,
770 in_alternate_slot: bool,
775 draining_to_replace: bool,
782 configuration_updated_since_spawn: bool,
788}
789
790impl SupervisorSnapshot {
791 fn starting() -> Self {
792 Self::new(ModuleState::Starting, true)
793 }
794
795 fn disabled() -> Self {
796 Self::new(ModuleState::Disabled, false)
797 }
798
799 fn failed() -> Self {
800 Self::new(ModuleState::Failed, true)
801 }
802
803 fn crash_restarts_in_window(&mut self, window: Duration, now: Instant) -> u32 {
807 while let Some(oldest) = self.crash_restarts.front() {
808 if now.duration_since(*oldest) > window {
809 self.crash_restarts.pop_front();
810 } else {
811 break;
812 }
813 }
814 u32::try_from(self.crash_restarts.len()).unwrap_or(u32::MAX)
815 }
816
817 fn record_crash_restart(&mut self, policy: &RestartPolicy, now: Instant) {
823 self.crash_restarts.push_back(now);
824 while self.crash_restarts.len() > policy.max_restarts as usize {
825 self.crash_restarts.pop_front();
826 }
827 self.lifetime_restarts += 1;
828 }
829
830 fn next_crash_restart(
834 &mut self,
835 policy: &RestartPolicy,
836 now: Instant,
837 ) -> Option<CrashRestartSchedule> {
838 let restart_in_window = self.crash_restarts_in_window(policy.window, now);
839 if restart_in_window >= policy.max_restarts {
840 return None;
841 }
842 self.record_crash_restart(policy, now);
843 Some(CrashRestartSchedule {
844 restart_in_window,
845 delay: policy.delay_for_restart(restart_in_window),
846 })
847 }
848
849 fn clear_crash_restarts(&mut self) {
854 self.crash_restarts.clear();
855 }
856
857 fn new(state: ModuleState, enabled: bool) -> Self {
858 Self {
859 state,
860 enabled,
861 process_alive: false,
862 crash_restarts: VecDeque::new(),
863 lifetime_restarts: 0,
864 spawn_generation: 0,
865 pid: None,
866 spawned_at_ms: None,
867 spawned_from: None,
868 spawned_file_identity: None,
869 process_start_time: None,
870 deliberate_severance: None,
871 last_exit: None,
872 health: ModuleHealthStatus::default(),
873 in_alternate_slot: false,
874 draining_to_replace: false,
875 configuration_updated_since_spawn: false,
876 }
877 }
878}
879
880type SharedSnapshot = Arc<Mutex<SupervisorSnapshot>>;
881
882type SpawnSubscriberKey = (ConnectionId, u64);
883
884#[derive(Debug)]
885struct SpawnSubscriber {
886 version: u8,
887 frames: mpsc::Sender<Frame>,
888 lagged: Option<oneshot::Sender<SpawnCursor>>,
892}
893
894#[derive(Debug)]
895struct SpawnEventState {
896 daemon_incarnation: String,
897 seq: u64,
898 capacity: usize,
899 live: HashMap<String, LiveSpawn>,
900 generations: HashMap<String, u64>,
901 events: VecDeque<SpawnEvent>,
902 subscribers: HashMap<SpawnSubscriberKey, SpawnSubscriber>,
903}
904
905impl Default for SpawnEventState {
906 fn default() -> Self {
907 Self {
908 daemon_incarnation: "unconfigured".to_string(),
909 seq: 0,
910 capacity: SPAWN_EVENT_RING_CAPACITY,
911 live: HashMap::new(),
912 generations: HashMap::new(),
913 events: VecDeque::new(),
914 subscribers: HashMap::new(),
915 }
916 }
917}
918
919#[derive(Debug, Clone, Default)]
920struct SpawnEventFeed(Arc<Mutex<SpawnEventState>>);
921
922#[derive(Debug, Clone, PartialEq, Eq)]
923pub(crate) enum SpawnSubscribeRefusal {
924 ForeignIncarnation { current: String },
925 TooOld { oldest: SpawnCursor },
926 Frame(String),
927}
928
929impl SpawnEventFeed {
930 fn configure_incarnation(&self, daemon_incarnation: String) {
931 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
932 state.daemon_incarnation = daemon_incarnation;
933 state.seq = 0;
934 state.live.clear();
935 state.generations.clear();
936 state.events.clear();
937 state.subscribers.clear();
938 }
939
940 fn cursor(state: &SpawnEventState) -> SpawnCursor {
941 SpawnCursor {
942 daemon_incarnation: state.daemon_incarnation.clone(),
943 seq: state.seq,
944 }
945 }
946
947 fn snapshot(&self) -> SpawnSnapshot {
948 let state = self.0.lock().unwrap_or_else(|p| p.into_inner());
949 let mut live = state.live.values().cloned().collect::<Vec<_>>();
950 live.sort_by(|left, right| left.module_id.cmp(&right.module_id));
951 SpawnSnapshot {
952 cursor: Self::cursor(&state),
953 ring_bound: state.capacity as u64,
954 live,
955 }
956 }
957
958 fn emit_spawned(&self, module_id: &str, pid: u32, spawned_at_ms: u64) -> u64 {
959 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
960 let generation = state
961 .generations
962 .get(module_id)
963 .copied()
964 .unwrap_or(0)
965 .checked_add(1)
966 .expect("spawn generation exhausted");
967 state.generations.insert(module_id.to_string(), generation);
968 let live = LiveSpawn {
969 module_id: module_id.to_string(),
970 spawn_generation: generation,
971 pid,
972 spawned_at_ms,
973 };
974 state.live.insert(module_id.to_string(), live);
975 Self::emit_locked(
976 &mut state,
977 SpawnEventKind::Spawned,
978 module_id.to_string(),
979 generation,
980 pid,
981 None,
982 None,
983 );
984 generation
985 }
986
987 fn emit_exited(&self, module_id: &str, exit_code: Option<i32>, exit_signal: Option<i32>) {
988 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
989 let Some(live) = state.live.remove(module_id) else {
990 warn!(
991 module_id,
992 "terminal record had no live spawn event identity"
993 );
994 return;
995 };
996 Self::emit_locked(
997 &mut state,
998 SpawnEventKind::Exited,
999 module_id.to_string(),
1000 live.spawn_generation,
1001 live.pid,
1002 exit_code,
1003 exit_signal,
1004 );
1005 }
1006
1007 fn emit_superseded_exited(
1014 &self,
1015 module_id: &str,
1016 spawn_generation: u64,
1017 pid: u32,
1018 exit_code: Option<i32>,
1019 exit_signal: Option<i32>,
1020 ) {
1021 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
1022 if state
1023 .live
1024 .get(module_id)
1025 .is_some_and(|live| live.spawn_generation == spawn_generation)
1026 {
1027 state.live.remove(module_id);
1028 }
1029 Self::emit_locked(
1030 &mut state,
1031 SpawnEventKind::Exited,
1032 module_id.to_string(),
1033 spawn_generation,
1034 pid,
1035 exit_code,
1036 exit_signal,
1037 );
1038 }
1039
1040 #[allow(clippy::too_many_arguments)]
1041 fn emit_locked(
1042 state: &mut SpawnEventState,
1043 kind: SpawnEventKind,
1044 module_id: String,
1045 spawn_generation: u64,
1046 pid: u32,
1047 exit_code: Option<i32>,
1048 exit_signal: Option<i32>,
1049 ) {
1050 state.seq = state
1051 .seq
1052 .checked_add(1)
1053 .expect("spawn event sequence exhausted");
1054 let event = SpawnEvent {
1055 cursor: Self::cursor(state),
1056 kind,
1057 module_id,
1058 spawn_generation,
1059 pid,
1060 exit_code,
1061 exit_signal,
1062 };
1063 state.events.push_back(event.clone());
1064 while state.events.len() > state.capacity {
1065 state.events.pop_front();
1066 }
1067 let body = match serde_json::to_vec(&event) {
1068 Ok(body) => body,
1069 Err(error) => {
1070 error!(%error, "failed to serialize supervisor spawn event");
1071 return;
1072 }
1073 };
1074 state.subscribers.retain(|(connection_id, corr), subscriber| {
1075 let frame = Frame::build_with_version(
1076 subscriber.version,
1077 FrameType::StreamData,
1078 control_flags(),
1079 0,
1080 0,
1081 *corr,
1082 body.clone(),
1083 );
1084 match frame {
1085 Ok(frame) => {
1086 if subscriber.frames.try_send(frame).is_ok() {
1087 true
1088 } else {
1089 warn!(connection_id = connection_id.get(), corr, "dropping lagged supervisor spawn subscriber");
1090 if let Some(lagged) = subscriber.lagged.take() {
1091 let _ = lagged.send(event.cursor.clone());
1092 }
1093 false
1094 }
1095 }
1096 Err(error) => {
1097 warn!(connection_id = connection_id.get(), corr, %error, "dropping supervisor spawn subscriber after frame build failure");
1098 false
1099 }
1100 }
1101 });
1102 }
1103
1104 fn subscribe(
1105 &self,
1106 connection_id: ConnectionId,
1107 corr: u64,
1108 version: u8,
1109 since: Option<SpawnCursor>,
1110 sink: FrameSink,
1111 ) -> Result<(), SpawnSubscribeRefusal> {
1112 let (frames, mut receiver) = mpsc::channel(SPAWN_SUBSCRIBER_BUFFER);
1113 let (lagged, mut lagged_rx) = oneshot::channel::<SpawnCursor>();
1114 {
1115 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
1116 let replay = if let Some(since) = since {
1117 if since.daemon_incarnation != state.daemon_incarnation {
1118 return Err(SpawnSubscribeRefusal::ForeignIncarnation {
1119 current: state.daemon_incarnation.clone(),
1120 });
1121 }
1122 if let Some(oldest) = state.events.front().map(|event| event.cursor.clone()) {
1123 if since.seq < oldest.seq.saturating_sub(1) {
1124 return Err(SpawnSubscribeRefusal::TooOld { oldest });
1125 }
1126 }
1127 state
1128 .events
1129 .iter()
1130 .filter(|event| event.cursor.seq > since.seq)
1131 .cloned()
1132 .collect::<Vec<_>>()
1133 } else {
1134 Vec::new()
1135 };
1136 for event in replay {
1137 let body = serde_json::to_vec(&event)
1138 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1139 let frame = Frame::build_with_version(
1140 version,
1141 FrameType::StreamData,
1142 control_flags(),
1143 0,
1144 0,
1145 corr,
1146 body,
1147 )
1148 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1149 frames
1150 .try_send(frame)
1151 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1152 }
1153 state.subscribers.insert(
1154 (connection_id, corr),
1155 SpawnSubscriber {
1156 version,
1157 frames,
1158 lagged: Some(lagged),
1159 },
1160 );
1161 }
1162 tokio::spawn(async move {
1173 while let Some(frame) = receiver.recv().await {
1174 if sink.send(frame).await.is_err() {
1175 return;
1176 }
1177 }
1178 let Ok(first_undelivered) = lagged_rx.try_recv() else {
1179 return;
1180 };
1181 match spawn_subscriber_lagged_frame(version, corr, first_undelivered) {
1182 Ok(frame) => {
1183 let _ = sink.send(frame).await;
1184 }
1185 Err(error) => {
1186 error!(%error, corr, "failed to build lagged spawn subscriber terminal frame");
1187 }
1188 }
1189 });
1190 Ok(())
1191 }
1192
1193 fn cancel(&self, connection_id: ConnectionId, corr: u64) -> bool {
1194 let Some(subscriber) = self
1195 .0
1196 .lock()
1197 .unwrap_or_else(|p| p.into_inner())
1198 .subscribers
1199 .remove(&(connection_id, corr))
1200 else {
1201 return false;
1202 };
1203 if let Ok(frame) = Frame::build_with_version(
1204 subscriber.version,
1205 FrameType::StreamEnd,
1206 control_flags(),
1207 0,
1208 0,
1209 corr,
1210 Vec::new(),
1211 ) {
1212 tokio::spawn(async move {
1213 let _ = subscriber.frames.send(frame).await;
1214 });
1215 }
1216 true
1217 }
1218
1219 fn remove_connection(&self, connection_id: ConnectionId) {
1220 self.0
1221 .lock()
1222 .unwrap_or_else(|p| p.into_inner())
1223 .subscribers
1224 .retain(|(subscriber_connection, _), _| *subscriber_connection != connection_id);
1225 }
1226
1227 #[cfg(any(test, feature = "test-support"))]
1228 fn set_capacity(&self, capacity: usize) {
1229 self.0.lock().unwrap_or_else(|p| p.into_inner()).capacity = capacity;
1230 }
1231
1232 #[cfg(any(test, feature = "test-support"))]
1233 fn subscriber_count(&self) -> usize {
1234 self.0
1235 .lock()
1236 .unwrap_or_else(|p| p.into_inner())
1237 .subscribers
1238 .len()
1239 }
1240}
1241
1242fn spawn_subscriber_lagged_frame(
1245 version: u8,
1246 corr: u64,
1247 first_undelivered: SpawnCursor,
1248) -> Result<Frame, String> {
1249 let body = serde_json::to_vec(&subc_protocol::ErrorBody {
1250 code: SPAWN_SUBSCRIBER_LAGGED_CODE.to_string(),
1251 message: "spawn subscriber fell behind and was dropped; resubscribe from the last cursor received"
1252 .to_string(),
1253 detail: Some(serde_json::json!({
1254 "first_undelivered_cursor": first_undelivered
1255 })),
1256 })
1257 .map_err(|error| error.to_string())?;
1258 Frame::build_with_version(version, FrameType::Error, control_flags(), 0, 0, corr, body)
1259 .map_err(|error| error.to_string())
1260}
1261
1262pub trait ModuleProcessLiveness: Send + Sync {
1263 fn process_live(&self, module_id: &str) -> Option<bool>;
1264
1265 fn process_replacing(&self, _module_id: &str) -> bool {
1271 false
1272 }
1273}
1274
1275#[derive(Debug, Clone, Default)]
1277pub struct SupervisorProcessLiveness {
1278 snapshots: Arc<Mutex<HashMap<String, SharedSnapshot>>>,
1279}
1280
1281impl SupervisorProcessLiveness {
1282 pub fn new() -> Self {
1283 Self::default()
1284 }
1285
1286 fn track(&self, module_id: String, snapshot: SharedSnapshot) {
1287 let mut snapshots = self
1288 .snapshots
1289 .lock()
1290 .unwrap_or_else(|poisoned| poisoned.into_inner());
1291 snapshots.insert(module_id, snapshot);
1292 }
1293
1294 fn untrack_if_current(&self, module_id: &str, snapshot: &SharedSnapshot) {
1295 let mut snapshots = self
1296 .snapshots
1297 .lock()
1298 .unwrap_or_else(|poisoned| poisoned.into_inner());
1299 let is_current = snapshots
1300 .get(module_id)
1301 .map(|tracked| Arc::ptr_eq(tracked, snapshot))
1302 .unwrap_or(false);
1303 if is_current {
1304 snapshots.remove(module_id);
1305 }
1306 }
1307}
1308
1309impl ModuleProcessLiveness for SupervisorProcessLiveness {
1310 fn process_live(&self, module_id: &str) -> Option<bool> {
1311 let snapshot = {
1312 let snapshots = self
1313 .snapshots
1314 .lock()
1315 .unwrap_or_else(|poisoned| poisoned.into_inner());
1316 snapshots.get(module_id).cloned()
1317 }?;
1318 let snapshot = snapshot
1319 .lock()
1320 .unwrap_or_else(|poisoned| poisoned.into_inner());
1321 Some(snapshot.state == ModuleState::Running && snapshot.process_alive)
1322 }
1323
1324 fn process_replacing(&self, module_id: &str) -> bool {
1325 let Some(snapshot) = self
1326 .snapshots
1327 .lock()
1328 .unwrap_or_else(|poisoned| poisoned.into_inner())
1329 .get(module_id)
1330 .cloned()
1331 else {
1332 return false;
1333 };
1334 let snapshot = snapshot
1335 .lock()
1336 .unwrap_or_else(|poisoned| poisoned.into_inner());
1337 snapshot.enabled
1338 && match snapshot.state {
1339 ModuleState::Restarting => true,
1340 ModuleState::Draining => snapshot.draining_to_replace,
1341 ModuleState::Starting
1342 | ModuleState::Running
1343 | ModuleState::Unresponsive
1344 | ModuleState::Stopped
1345 | ModuleState::Failed
1346 | ModuleState::Disabled => false,
1347 }
1348 }
1349}
1350
1351#[derive(Debug, Clone)]
1352struct SupervisorRuntimeConfig {
1353 restart_policy: RestartPolicy,
1354 drain_timeout: Duration,
1357 effective_drain_timeout: Arc<Mutex<Duration>>,
1360 default_drain_timeout: Duration,
1363 health: HealthConfig,
1364 connection_file_path: Option<PathBuf>,
1365 capture_logs_dir: Option<PathBuf>,
1366 forwarding: Option<Arc<ForwardingTable>>,
1367 supervisor_handle: Option<SupervisorHandle>,
1370 stderr_ring: Arc<Mutex<StderrRing>>,
1377 terminal_ring: Arc<Mutex<TerminalRing>>,
1378 spawn_events: SpawnEventFeed,
1379 child_roster: ChildRoster,
1380 #[cfg(target_os = "linux")]
1381 cgroup_placement: Option<subc_cgroup::Placement>,
1382 #[cfg(test)]
1383 test_seed_stale_facts_before_enable_spawn: bool,
1384}
1385
1386#[derive(Debug, Clone, PartialEq, Eq)]
1387struct SupervisedConfiguration {
1388 spec: ModuleSpec,
1389 health: HealthConfig,
1390}
1391
1392#[derive(Debug, Clone, Default)]
1398pub struct SupervisorHandle {
1399 modules: Arc<Mutex<HashMap<String, SupervisedModule>>>,
1400 spawn_events: SpawnEventFeed,
1401 reserved_nonces: Arc<Mutex<HashMap<String, Option<String>>>>,
1412 removal_tombstones: Arc<Mutex<HashMap<String, u64>>>,
1418 spawn_nonces: Arc<Mutex<HashMap<String, String>>>,
1422 reserved_prefix_owners: Arc<Mutex<HashMap<String, String>>>,
1430 swaps: Arc<Mutex<HashMap<String, OpenSwap>>>,
1436 promotion_observer: PromotionObserverSlot,
1438 operation_lock: Arc<AsyncMutex<()>>,
1442}
1443
1444pub(crate) trait SwapPromotionObserver: Send + Sync {
1453 fn swap_promoted(&self, registration: &crate::registry::ModuleRegistration);
1454}
1455
1456#[derive(Clone, Default)]
1460struct PromotionObserverSlot(Arc<Mutex<Option<std::sync::Weak<dyn SwapPromotionObserver>>>>);
1461
1462impl fmt::Debug for PromotionObserverSlot {
1463 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1464 f.write_str("PromotionObserverSlot")
1465 }
1466}
1467
1468#[derive(Debug, Clone)]
1470struct OpenSwap {
1471 candidate_nonce: String,
1474 incumbent_nonce: Option<String>,
1479 candidate_admitted: bool,
1483}
1484
1485#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1488pub(crate) enum SwapHelloAdmission {
1489 NotSwapping,
1492 Candidate,
1494 Refused,
1497}
1498
1499#[derive(Debug, Clone, PartialEq, Eq)]
1500pub(crate) enum ReservedHelloRejection {
1501 Exact {
1502 module_id: String,
1503 },
1504 Prefix {
1505 prefix: String,
1506 owner_module_id: String,
1507 },
1508}
1509
1510impl SupervisorHandle {
1511 pub fn new() -> Self {
1512 Self::default()
1513 }
1514
1515 pub(crate) fn spawn_snapshot(&self) -> SpawnSnapshot {
1516 self.spawn_events.snapshot()
1517 }
1518
1519 pub(crate) fn subscribe_spawns(
1520 &self,
1521 connection_id: ConnectionId,
1522 corr: u64,
1523 version: u8,
1524 since: Option<SpawnCursor>,
1525 sink: FrameSink,
1526 ) -> Result<(), SpawnSubscribeRefusal> {
1527 self.spawn_events
1528 .subscribe(connection_id, corr, version, since, sink)
1529 }
1530
1531 pub(crate) fn cancel_spawn_subscription(&self, connection_id: ConnectionId, corr: u64) -> bool {
1532 self.spawn_events.cancel(connection_id, corr)
1533 }
1534
1535 pub(crate) fn remove_spawn_subscribers(&self, connection_id: ConnectionId) {
1536 self.spawn_events.remove_connection(connection_id);
1537 }
1538
1539 #[cfg(any(test, feature = "test-support"))]
1540 pub fn set_spawn_event_capacity_for_test(&self, capacity: usize) {
1541 assert!(capacity > 0, "spawn event capacity must be non-zero");
1542 self.spawn_events.set_capacity(capacity);
1543 }
1544
1545 #[cfg(any(test, feature = "test-support"))]
1546 pub fn spawn_subscriber_count_for_test(&self) -> usize {
1547 self.spawn_events.subscriber_count()
1548 }
1549
1550 pub fn set_spawn_nonce(&self, module_id: &str, nonce: String) {
1553 self.spawn_nonces
1554 .lock()
1555 .unwrap_or_else(|poisoned| poisoned.into_inner())
1556 .insert(module_id.to_string(), nonce);
1557 }
1558
1559 pub fn set_reserved_nonce(&self, module_id: &str, nonce: String) {
1562 self.reserved_nonces
1563 .lock()
1564 .unwrap_or_else(|poisoned| poisoned.into_inner())
1565 .insert(module_id.to_string(), Some(nonce));
1566 }
1567
1568 pub fn set_reserved_prefixes(&self, owner_module_id: &str, prefixes: &[String]) {
1570 let mut owners = self
1571 .reserved_prefix_owners
1572 .lock()
1573 .unwrap_or_else(|poisoned| poisoned.into_inner());
1574 owners.retain(|_, owner| owner != owner_module_id);
1575 for prefix in prefixes {
1576 owners.insert(prefix.clone(), owner_module_id.to_string());
1577 }
1578 }
1579
1580 #[cfg(test)]
1582 pub(crate) fn spawn_nonce(&self, module_id: &str) -> Option<String> {
1583 self.spawn_nonces
1584 .lock()
1585 .unwrap_or_else(|poisoned| poisoned.into_inner())
1586 .get(module_id)
1587 .cloned()
1588 }
1589
1590 fn apply_identity_configuration(&self, spec: &ModuleSpec) {
1591 self.set_reserved_prefixes(&spec.module_id, &spec.reserved_prefixes);
1592 let spawn_nonce = self
1593 .spawn_nonces
1594 .lock()
1595 .unwrap_or_else(|poisoned| poisoned.into_inner())
1596 .get(&spec.module_id)
1597 .cloned();
1598 let mut reserved_nonces = self
1599 .reserved_nonces
1600 .lock()
1601 .unwrap_or_else(|poisoned| poisoned.into_inner());
1602 if spec.reserved {
1603 reserved_nonces.insert(spec.module_id.clone(), spawn_nonce);
1608 }
1609 drop(reserved_nonces);
1610 self.removal_tombstones
1614 .lock()
1615 .unwrap_or_else(|poisoned| poisoned.into_inner())
1616 .remove(&spec.module_id);
1617 }
1618
1619 pub fn reserved_hello_authorized(&self, module_id: &str, presented: Option<&str>) -> bool {
1624 self.reserved_hello_rejection(module_id, presented)
1625 .is_none()
1626 }
1627
1628 pub(crate) fn reserved_hello_rejection(
1629 &self,
1630 module_id: &str,
1631 presented: Option<&str>,
1632 ) -> Option<ReservedHelloRejection> {
1633 let nonces = self
1634 .reserved_nonces
1635 .lock()
1636 .unwrap_or_else(|poisoned| poisoned.into_inner());
1637 if let Some(expected) = nonces.get(module_id) {
1638 let authorized = match expected {
1642 Some(expected) => {
1643 presented.is_some_and(|p| constant_time_eq(expected.as_bytes(), p.as_bytes()))
1644 }
1645 None => false,
1646 };
1647 if authorized {
1648 return None;
1649 }
1650 return Some(ReservedHelloRejection::Exact {
1651 module_id: module_id.to_string(),
1652 });
1653 }
1654 drop(nonces);
1655
1656 let matched_prefix = self
1657 .reserved_prefix_owners
1658 .lock()
1659 .unwrap_or_else(|poisoned| poisoned.into_inner())
1660 .iter()
1661 .filter(|(prefix, _)| module_id.starts_with(prefix.as_str()))
1662 .max_by_key(|(prefix, _)| prefix.len())
1663 .map(|(prefix, owner)| (prefix.clone(), owner.clone()));
1664 let (prefix, owner_module_id) = matched_prefix?;
1665
1666 let authorized = presented.is_some_and(|presented| {
1667 self.spawn_nonces
1668 .lock()
1669 .unwrap_or_else(|poisoned| poisoned.into_inner())
1670 .get(&owner_module_id)
1671 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()))
1672 || self.swap_nonce_matches(&owner_module_id, presented)
1675 });
1676 if authorized {
1677 None
1678 } else {
1679 Some(ReservedHelloRejection::Prefix {
1680 prefix,
1681 owner_module_id,
1682 })
1683 }
1684 }
1685
1686 pub fn spawned_consumer_authorized(&self, module_id: &str, presented: &str) -> bool {
1691 if presented.is_empty() {
1692 return false;
1693 }
1694 let nonces = self
1695 .spawn_nonces
1696 .lock()
1697 .unwrap_or_else(|poisoned| poisoned.into_inner());
1698 let current = nonces
1699 .get(module_id)
1700 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()));
1701 drop(nonces);
1702 current || self.swap_nonce_matches(module_id, presented)
1707 }
1708
1709 fn swap_nonce_matches(&self, module_id: &str, presented: &str) -> bool {
1711 let swaps = self
1712 .swaps
1713 .lock()
1714 .unwrap_or_else(|poisoned| poisoned.into_inner());
1715 swaps.get(module_id).is_some_and(|swap| {
1716 constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes())
1717 || swap.incumbent_nonce.as_deref().is_some_and(|incumbent| {
1718 constant_time_eq(incumbent.as_bytes(), presented.as_bytes())
1719 })
1720 })
1721 }
1722
1723 pub(crate) fn open_swap(&self, module_id: &str, candidate_nonce: String) {
1726 let incumbent_nonce = self
1727 .spawn_nonces
1728 .lock()
1729 .unwrap_or_else(|poisoned| poisoned.into_inner())
1730 .get(module_id)
1731 .cloned();
1732 self.swaps
1733 .lock()
1734 .unwrap_or_else(|poisoned| poisoned.into_inner())
1735 .insert(
1736 module_id.to_string(),
1737 OpenSwap {
1738 candidate_nonce,
1739 incumbent_nonce,
1740 candidate_admitted: false,
1741 },
1742 );
1743 }
1744
1745 pub(crate) fn close_swap(&self, module_id: &str) {
1748 self.swaps
1749 .lock()
1750 .unwrap_or_else(|poisoned| poisoned.into_inner())
1751 .remove(module_id);
1752 }
1753
1754 pub(crate) fn set_swap_promotion_observer(
1757 &self,
1758 observer: std::sync::Weak<dyn SwapPromotionObserver>,
1759 ) {
1760 *self
1761 .promotion_observer
1762 .0
1763 .lock()
1764 .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(observer);
1765 }
1766
1767 fn notify_swap_promoted(&self, registration: &crate::registry::ModuleRegistration) {
1770 let observer = self
1771 .promotion_observer
1772 .0
1773 .lock()
1774 .unwrap_or_else(|poisoned| poisoned.into_inner())
1775 .as_ref()
1776 .and_then(std::sync::Weak::upgrade);
1777 if let Some(observer) = observer {
1778 observer.swap_promoted(registration);
1779 }
1780 }
1781
1782 pub(crate) fn swap_open(&self, module_id: &str) -> bool {
1784 self.swaps
1785 .lock()
1786 .unwrap_or_else(|poisoned| poisoned.into_inner())
1787 .contains_key(module_id)
1788 }
1789
1790 fn promote_swap_nonce(&self, module_id: &str, reserved: bool) {
1795 let candidate_nonce = self
1796 .swaps
1797 .lock()
1798 .unwrap_or_else(|poisoned| poisoned.into_inner())
1799 .get(module_id)
1800 .map(|swap| swap.candidate_nonce.clone());
1801 let Some(nonce) = candidate_nonce else {
1802 return;
1803 };
1804 self.set_spawn_nonce(module_id, nonce.clone());
1805 if reserved {
1806 self.set_reserved_nonce(module_id, nonce);
1807 }
1808 }
1809
1810 pub(crate) fn swap_hello_admission(
1825 &self,
1826 module_id: &str,
1827 presented: Option<&str>,
1828 ) -> SwapHelloAdmission {
1829 let swaps = self
1830 .swaps
1831 .lock()
1832 .unwrap_or_else(|poisoned| poisoned.into_inner());
1833 let Some(swap) = swaps.get(module_id) else {
1834 return SwapHelloAdmission::NotSwapping;
1835 };
1836 let Some(presented) = presented else {
1837 return SwapHelloAdmission::Refused;
1838 };
1839 if constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes()) {
1840 return if swap.candidate_admitted {
1841 SwapHelloAdmission::Refused
1842 } else {
1843 SwapHelloAdmission::Candidate
1844 };
1845 }
1846 if swap
1847 .incumbent_nonce
1848 .as_deref()
1849 .is_some_and(|incumbent| constant_time_eq(incumbent.as_bytes(), presented.as_bytes()))
1850 {
1851 return SwapHelloAdmission::NotSwapping;
1852 }
1853 SwapHelloAdmission::Refused
1854 }
1855
1856 pub(crate) fn mark_swap_candidate_admitted(&self, module_id: &str) {
1859 if let Some(swap) = self
1860 .swaps
1861 .lock()
1862 .unwrap_or_else(|poisoned| poisoned.into_inner())
1863 .get_mut(module_id)
1864 {
1865 swap.candidate_admitted = true;
1866 }
1867 }
1868
1869 pub fn spawn_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1871 self.spawn_nonces
1872 .lock()
1873 .unwrap_or_else(|poisoned| poisoned.into_inner())
1874 .get(module_id)
1875 .cloned()
1876 }
1877
1878 pub fn reserved_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1880 self.reserved_nonces
1881 .lock()
1882 .unwrap_or_else(|poisoned| poisoned.into_inner())
1883 .get(module_id)
1884 .cloned()
1885 .flatten()
1886 }
1887
1888 pub fn insert(&self, module: SupervisedModule) -> Option<SupervisedModule> {
1889 let mut modules = self
1890 .modules
1891 .lock()
1892 .unwrap_or_else(|poisoned| poisoned.into_inner());
1893 modules.insert(module.module_id().to_string(), module)
1894 }
1895
1896 pub fn get(&self, module_id: &str) -> Option<SupervisedModule> {
1897 let modules = self
1898 .modules
1899 .lock()
1900 .unwrap_or_else(|poisoned| poisoned.into_inner());
1901 modules.get(module_id).cloned()
1902 }
1903
1904 pub(crate) fn record_late_health_answer(
1905 &self,
1906 module_id: &str,
1907 latency_ms: u64,
1908 ) -> Result<bool, SuperviseError> {
1909 let Some(module) = self.get(module_id) else {
1910 return Ok(false);
1911 };
1912 update_snapshot(&module.inner.snapshot, Some(module_id), |state| {
1913 state.health.late_answer_count = state.health.late_answer_count.saturating_add(1);
1914 state.health.last_late_answer_latency_ms = Some(latency_ms);
1915 state.health.consecutive_failures = 0;
1923 })?;
1924 Ok(true)
1925 }
1926
1927 pub fn record_deliberate_severance(&self, module_id: &str) -> Result<bool, SuperviseError> {
1933 let Some(module) = self.get(module_id) else {
1934 return Ok(false);
1935 };
1936 let status = module.status()?;
1937 let Some((pid, start_time)) = status.pid.zip(status.process_start_time) else {
1938 return Ok(false);
1939 };
1940 module.record_deliberate_severance(ProcessIdentity { pid, start_time })
1941 }
1942
1943 pub fn list(&self) -> Vec<SupervisedModule> {
1944 let modules = self
1945 .modules
1946 .lock()
1947 .unwrap_or_else(|poisoned| poisoned.into_inner());
1948 let mut modules = modules.values().cloned().collect::<Vec<_>>();
1949 modules.sort_by(|left, right| left.module_id().cmp(right.module_id()));
1950 modules
1951 }
1952
1953 pub(crate) fn retire(&self, module_id: &str) -> Option<SupervisedModule> {
1954 self.spawn_nonces
1955 .lock()
1956 .unwrap_or_else(|poisoned| poisoned.into_inner())
1957 .remove(module_id);
1958 self.close_swap(module_id);
1959 let mut reserved_nonces = self
1960 .reserved_nonces
1961 .lock()
1962 .unwrap_or_else(|poisoned| poisoned.into_inner());
1963 if reserved_nonces.contains_key(module_id) {
1964 reserved_nonces.insert(module_id.to_string(), None);
1967 }
1968 drop(reserved_nonces);
1969 self.reserved_prefix_owners
1970 .lock()
1971 .unwrap_or_else(|poisoned| poisoned.into_inner())
1972 .retain(|_, owner| owner != module_id);
1973 self.modules
1974 .lock()
1975 .unwrap_or_else(|poisoned| poisoned.into_inner())
1976 .remove(module_id)
1977 }
1978
1979 pub(crate) fn record_rescan_removal(&self, module_id: &str) {
1982 self.removal_tombstones
1983 .lock()
1984 .unwrap_or_else(|poisoned| poisoned.into_inner())
1985 .insert(module_id.to_string(), unix_ms_now());
1986 }
1987
1988 pub(crate) fn removal_tombstone_age_ms(&self, module_id: &str) -> Option<u64> {
1990 self.removal_tombstones
1991 .lock()
1992 .unwrap_or_else(|poisoned| poisoned.into_inner())
1993 .get(module_id)
1994 .copied()
1995 .map(|removed_at_ms| unix_ms_now().saturating_sub(removed_at_ms))
1996 }
1997
1998 pub(crate) fn release_retained_reserved_gate(&self, module_id: &str) -> bool {
2003 if self.get(module_id).is_some() {
2004 return false;
2005 }
2006 let mut reserved_nonces = self
2007 .reserved_nonces
2008 .lock()
2009 .unwrap_or_else(|poisoned| poisoned.into_inner());
2010 if !matches!(reserved_nonces.get(module_id), Some(None)) {
2011 return false;
2012 }
2013 reserved_nonces.remove(module_id);
2014 true
2015 }
2016
2017 pub(crate) fn operation_lock(&self) -> Arc<AsyncMutex<()>> {
2018 Arc::clone(&self.operation_lock)
2019 }
2020}
2021
2022#[derive(Debug, Clone)]
2024pub struct Supervisor {
2025 registry: Arc<Registry>,
2026 restart_policy: RestartPolicy,
2027 drain_timeout: Duration,
2028 connection_file_path: Option<PathBuf>,
2029 capture_logs_dir: Option<PathBuf>,
2030 forwarding: Option<Arc<ForwardingTable>>,
2031 process_liveness: Arc<SupervisorProcessLiveness>,
2032 supervisor_handle: Option<SupervisorHandle>,
2033 health: HealthConfig,
2034 daemon_start_clock: crate::clock::StartClock,
2035 terminal_journal: Option<Arc<crate::terminal_journal::TerminalJournal>>,
2036 spawn_events: SpawnEventFeed,
2037 provenance_probe: ExecutableIdentityProbe,
2038 child_roster: ChildRoster,
2041 #[cfg(target_os = "linux")]
2042 cgroup_placement: Option<subc_cgroup::Placement>,
2043}
2044
2045impl Supervisor {
2046 #[cfg(unix)]
2057 pub(crate) fn begin_daemon_shutdown(&self) {
2058 self.child_roster.close();
2059 if let Some(journal) = &self.terminal_journal {
2060 journal.stamp_shutdown();
2061 }
2062 }
2063
2064 #[cfg(unix)]
2068 pub(crate) async fn drain_for_daemon_shutdown(&self) -> Result<(), SuperviseError> {
2069 const NOTICE_BUDGET: Duration = Duration::from_millis(500);
2070 const DRAIN_BUDGET: Duration = Duration::from_secs(2);
2071 let Some(forwarding) = &self.forwarding else {
2072 return Ok(());
2073 };
2074 let module_ids = forwarding
2075 .begin_daemon_drain()
2076 .map_err(SuperviseError::Forwarding)?;
2077 let deadline_ms =
2078 unix_ms_now().saturating_add((NOTICE_BUDGET + DRAIN_BUDGET).as_millis() as u64);
2079 let mut notices = tokio::task::JoinSet::new();
2080 let mut drains = Vec::new();
2081 for module_id in module_ids {
2082 let Some(target) = forwarding
2083 .begin_module_drain(&module_id, RouteCloseReason::Restart)
2084 .map_err(SuperviseError::Forwarding)?
2085 else {
2086 continue;
2087 };
2088 let routes = forwarding
2089 .endpoint_routes(target.endpoint)
2090 .map_err(SuperviseError::Forwarding)?;
2091 let command = serde_json::to_vec(&ModuleControlCommand::Draining {
2097 reason: RouteCloseReason::Restart,
2098 deadline_ms,
2099 })
2100 .expect("module draining serializes");
2101 let closing = serde_json::to_vec(&ClientControlPush::RouteClosing {
2102 module_id: module_id.clone(),
2103 reason: RouteCloseReason::Restart,
2104 })
2105 .expect("route closing serializes");
2106 let mut recipients = vec![(target.sink.clone(), target.negotiated_ver, command)];
2107 let mut seen = std::collections::HashSet::new();
2108 for route in routes {
2109 let client = route.goodbye_target;
2110 if seen.insert(client.connection_id) {
2111 recipients.push((client.sink, client.negotiated_ver, closing.clone()));
2112 }
2113 }
2114 for (sink, version, body) in recipients {
2115 notices.spawn(async move {
2116 let frame = Frame::build_with_version(
2117 version,
2118 FrameType::Push,
2119 control_flags(),
2120 0,
2121 0,
2122 0,
2123 body,
2124 )
2125 .expect("bounded lifecycle notice frame builds");
2126 sink.send_flushed(frame).await
2127 });
2128 }
2129 let gauges = declared_busy_gauges(&self.registry, &module_id)?;
2130 drains.push((module_id, target.endpoint, gauges));
2131 }
2132 let notice_deadline = Instant::now() + NOTICE_BUDGET;
2135 while let Ok(Some(result)) = timeout_at(notice_deadline, notices.join_next()).await {
2136 if !matches!(result, Ok(Ok(()))) {
2137 warn!(?result, "daemon shutdown notice delivery failed");
2138 }
2139 }
2140 notices.abort_all();
2141 let deadline = Instant::now() + DRAIN_BUDGET;
2142 let mut waits = tokio::task::JoinSet::new();
2143 for (module_id, endpoint, gauges) in drains {
2144 let forwarding = Arc::clone(forwarding);
2145 let mut runtime = self.runtime_config();
2146 runtime.health.cadence = Duration::from_millis(100);
2147 waits.spawn(async move {
2148 wait_for_forwarding_quiescence(
2149 &forwarding,
2150 &module_id,
2151 &runtime,
2152 endpoint,
2153 deadline,
2154 &gauges,
2155 DrainScope::Active,
2156 )
2157 .await
2158 });
2159 }
2160 while let Ok(Some(result)) = timeout_at(deadline, waits.join_next()).await {
2161 if !matches!(result, Ok(Ok(true))) {
2162 warn!(?result, "daemon shutdown drain did not reach quiescence");
2163 }
2164 }
2165 Ok(())
2166 }
2167
2168 #[cfg(unix)]
2178 pub(crate) async fn end_children_for_daemon_shutdown(
2179 &self,
2180 already_escalated: bool,
2181 escalate: impl std::future::Future<Output = ()>,
2182 ) {
2183 if let Some(forwarding) = &self.forwarding {
2184 let closed = forwarding.close_all_connections(&CloseReason::new(
2185 "daemon_shutdown",
2186 "the daemon is exiting after its shutdown notice and drain",
2187 ));
2188 debug!(closed, "closed established connections for daemon shutdown");
2189 }
2190 crate::child_roster::end_children_for_daemon_shutdown(
2191 &self.child_roster,
2192 already_escalated,
2193 escalate,
2194 )
2195 .await;
2196 }
2197
2198 pub fn new(registry: Arc<Registry>, restart_policy: RestartPolicy) -> Self {
2199 Self {
2200 registry,
2201 restart_policy,
2202 drain_timeout: DEFAULT_DRAIN_TIMEOUT,
2203 connection_file_path: None,
2204 capture_logs_dir: None,
2205 forwarding: None,
2206 process_liveness: Arc::new(SupervisorProcessLiveness::default()),
2207 supervisor_handle: None,
2208 health: HealthConfig::default(),
2209 daemon_start_clock: crate::clock::StartClock::capture(),
2210 terminal_journal: None,
2211 spawn_events: SpawnEventFeed::default(),
2212 provenance_probe: ExecutableIdentityProbe::default(),
2213 child_roster: ChildRoster::default(),
2214 #[cfg(target_os = "linux")]
2215 cgroup_placement: None,
2216 }
2217 }
2218
2219 pub fn with_drain_timeout(mut self, drain_timeout: Duration) -> Self {
2220 self.drain_timeout = drain_timeout;
2221 self
2222 }
2223
2224 pub fn with_process_liveness(
2225 mut self,
2226 process_liveness: Arc<SupervisorProcessLiveness>,
2227 ) -> Self {
2228 self.process_liveness = process_liveness;
2229 self
2230 }
2231
2232 pub fn with_connection_file_path(mut self, connection_file_path: impl Into<PathBuf>) -> Self {
2233 self.connection_file_path = Some(connection_file_path.into());
2234 self
2235 }
2236
2237 pub fn with_capture_logs_dir(mut self, logs_dir: impl Into<PathBuf>) -> Self {
2239 self.capture_logs_dir = Some(logs_dir.into());
2240 self
2241 }
2242
2243 pub fn with_daemon_incarnation(self, daemon_incarnation: String) -> Self {
2246 self.spawn_events.configure_incarnation(daemon_incarnation);
2250 self
2251 }
2252
2253 pub fn with_terminal_journal(self, path: PathBuf, daemon_incarnation: String) -> Self {
2256 let mut this = self.with_daemon_incarnation(daemon_incarnation.clone());
2257 this.terminal_journal = Some(Arc::new(crate::terminal_journal::TerminalJournal::open(
2258 path,
2259 daemon_incarnation,
2260 )));
2261 this
2262 }
2263
2264 pub fn with_forwarding(mut self, forwarding: Arc<ForwardingTable>) -> Self {
2265 self.forwarding = Some(forwarding);
2266 self
2267 }
2268
2269 pub fn with_handle(mut self, supervisor_handle: SupervisorHandle) -> Self {
2270 self.spawn_events = supervisor_handle.spawn_events.clone();
2271 self.supervisor_handle = Some(supervisor_handle);
2272 self
2273 }
2274
2275 pub fn with_health_config(mut self, health: HealthConfig) -> Self {
2276 self.health = health;
2277 self
2278 }
2279
2280 pub fn with_live_children_record(self, path: impl Into<PathBuf>) -> Self {
2284 self.child_roster.record_to(path.into());
2285 self
2286 }
2287
2288 #[cfg(target_os = "linux")]
2289 pub fn with_cgroup_placement(
2290 mut self,
2291 cgroup_placement: Option<subc_cgroup::Placement>,
2292 ) -> Self {
2293 self.cgroup_placement = cgroup_placement;
2294 self
2295 }
2296
2297 pub fn spawn(&self, spec: ModuleSpec) -> Result<SupervisedModule, SuperviseError> {
2303 validate_spec(&spec)?;
2304
2305 let runtime = self.runtime_config();
2306 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2307 let child = spawn_child(
2308 &spec,
2309 runtime.connection_file_path.as_deref(),
2310 self.supervisor_handle.as_ref(),
2311 &runtime.stderr_ring,
2312 runtime.capture_logs_dir.as_deref(),
2313 &runtime.child_roster,
2314 #[cfg(target_os = "linux")]
2315 runtime.cgroup_placement.as_ref(),
2316 )?;
2317 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2318 self.process_liveness
2319 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2320
2321 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2322 }
2323
2324 pub fn supervise_configured(
2330 &self,
2331 spec: ModuleSpec,
2332 enabled: bool,
2333 ) -> Result<SupervisedModule, SuperviseError> {
2334 validate_spec(&spec)?;
2335
2336 let runtime = self.runtime_config();
2337 if !enabled {
2338 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2339 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2340 }
2341
2342 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2343 match spawn_child(
2344 &spec,
2345 runtime.connection_file_path.as_deref(),
2346 self.supervisor_handle.as_ref(),
2347 &runtime.stderr_ring,
2348 runtime.capture_logs_dir.as_deref(),
2349 &runtime.child_roster,
2350 #[cfg(target_os = "linux")]
2351 runtime.cgroup_placement.as_ref(),
2352 ) {
2353 Ok(child) => {
2354 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2355 self.process_liveness
2356 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2357 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2358 }
2359 Err(err) => {
2360 error!(
2361 module_id = %spec.module_id,
2362 program = %spec.program.display(),
2363 error = %err,
2364 "configured module failed to spawn; marking failed and continuing"
2365 );
2366 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2367 Ok(self.supervised_module(spec, runtime, snapshot, None))
2368 }
2369 }
2370 }
2371
2372 pub fn supervise_configured_with_health(
2378 &self,
2379 spec: ModuleSpec,
2380 enabled: bool,
2381 health: HealthConfig,
2382 drain_timeout_ms: Option<u64>,
2383 restart_policy: RestartPolicy,
2384 ) -> Result<SupervisedModule, SuperviseError> {
2385 validate_spec(&spec)?;
2386
2387 let mut runtime = self.runtime_config();
2388 runtime.health = health;
2389 runtime.restart_policy = restart_policy;
2390 if let Some(ms) = drain_timeout_ms {
2391 runtime.drain_timeout = Duration::from_millis(ms);
2392 *runtime
2393 .effective_drain_timeout
2394 .lock()
2395 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
2396 }
2397 if !enabled {
2398 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2399 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2400 }
2401
2402 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2403 match spawn_child(
2404 &spec,
2405 runtime.connection_file_path.as_deref(),
2406 self.supervisor_handle.as_ref(),
2407 &runtime.stderr_ring,
2408 runtime.capture_logs_dir.as_deref(),
2409 &runtime.child_roster,
2410 #[cfg(target_os = "linux")]
2411 runtime.cgroup_placement.as_ref(),
2412 ) {
2413 Ok(child) => {
2414 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2415 self.process_liveness
2416 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2417 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2418 }
2419 Err(err) => {
2420 if health.critical {
2421 error!(
2422 module_id = %spec.module_id,
2423 program = %spec.program.display(),
2424 error = %err,
2425 "critical configured module failed to spawn; marking failed and alerting"
2426 );
2427 } else {
2428 error!(
2429 module_id = %spec.module_id,
2430 program = %spec.program.display(),
2431 error = %err,
2432 "configured module failed to spawn; marking failed and continuing"
2433 );
2434 }
2435 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2436 Ok(self.supervised_module(spec, runtime, snapshot, None))
2437 }
2438 }
2439 }
2440
2441 fn runtime_config(&self) -> SupervisorRuntimeConfig {
2442 let effective_drain_timeout = Arc::new(Mutex::new(self.drain_timeout));
2443 SupervisorRuntimeConfig {
2444 restart_policy: self.restart_policy,
2445 drain_timeout: self.drain_timeout,
2446 child_roster: self
2449 .child_roster
2450 .for_module(Arc::clone(&effective_drain_timeout)),
2451 effective_drain_timeout,
2452 default_drain_timeout: self.drain_timeout,
2453 health: self.health,
2454 connection_file_path: self.connection_file_path.clone(),
2455 capture_logs_dir: self.capture_logs_dir.clone(),
2456 forwarding: self.forwarding.clone(),
2457 supervisor_handle: self.supervisor_handle.clone(),
2458 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
2459 terminal_ring: Arc::new(Mutex::new(
2460 TerminalRing::new(
2461 TerminalRingConfig::default(),
2462 self.daemon_start_clock.started_at_ms(),
2463 )
2464 .with_start_clock(self.daemon_start_clock)
2465 .with_journal(self.terminal_journal.clone())
2466 .with_daemon_shutdown(self.child_roster.shutdown_flag()),
2467 )),
2468 spawn_events: self.spawn_events.clone(),
2469 #[cfg(target_os = "linux")]
2470 cgroup_placement: self.cgroup_placement.clone(),
2471 #[cfg(test)]
2472 test_seed_stale_facts_before_enable_spawn: false,
2473 }
2474 }
2475
2476 fn supervised_module(
2477 &self,
2478 spec: ModuleSpec,
2479 runtime: SupervisorRuntimeConfig,
2480 snapshot: SharedSnapshot,
2481 child: Option<SupervisedChild>,
2482 ) -> SupervisedModule {
2483 let configuration = Arc::new(Mutex::new(SupervisedConfiguration {
2484 spec: spec.clone(),
2485 health: runtime.health,
2486 }));
2487 let stderr_ring = Arc::clone(&runtime.stderr_ring);
2488 let terminal_ring = Arc::clone(&runtime.terminal_ring);
2489 let restart_policy = runtime.restart_policy;
2493 let effective_drain_timeout = Arc::clone(&runtime.effective_drain_timeout);
2494 let (tx, rx) = mpsc::channel(4);
2495 let monitor = tokio::spawn(supervise_loop(
2496 spec.clone(),
2497 runtime,
2498 Arc::clone(&self.registry),
2499 Arc::clone(&self.process_liveness),
2500 Arc::clone(&snapshot),
2501 child,
2502 rx,
2503 ));
2504
2505 let module_id = spec.module_id.clone();
2506 let module = SupervisedModule {
2507 inner: Arc::new(SupervisedModuleInner {
2508 module_id: module_id.clone(),
2509 registry: Arc::clone(&self.registry),
2510 snapshot,
2511 configuration,
2512 stderr_ring,
2513 terminal_ring,
2514 commands: tx,
2515 monitor: Mutex::new(Some(monitor)),
2516 restart_policy,
2517 effective_drain_timeout,
2518 provenance_probe: self.provenance_probe.clone(),
2519 }),
2520 };
2521 if let Some(supervisor_handle) = &self.supervisor_handle {
2522 supervisor_handle.apply_identity_configuration(&spec);
2523 supervisor_handle.insert(module.clone());
2524 }
2525 module
2526 }
2527}
2528
2529impl Default for Supervisor {
2530 fn default() -> Self {
2531 Self::new(Arc::new(Registry::default()), RestartPolicy::default())
2532 }
2533}
2534
2535#[derive(Clone)]
2537pub struct SupervisedModule {
2538 inner: Arc<SupervisedModuleInner>,
2539}
2540
2541struct SupervisedModuleInner {
2542 module_id: String,
2543 registry: Arc<Registry>,
2544 snapshot: SharedSnapshot,
2545 configuration: Arc<Mutex<SupervisedConfiguration>>,
2546 stderr_ring: Arc<Mutex<StderrRing>>,
2547 terminal_ring: Arc<Mutex<TerminalRing>>,
2548 commands: mpsc::Sender<SupervisorCommand>,
2549 monitor: Mutex<Option<JoinHandle<()>>>,
2550 restart_policy: RestartPolicy,
2554 effective_drain_timeout: Arc<Mutex<Duration>>,
2555 provenance_probe: ExecutableIdentityProbe,
2556}
2557
2558impl fmt::Debug for SupervisedModule {
2559 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2560 f.debug_struct("SupervisedModule")
2561 .field("module_id", &self.inner.module_id)
2562 .field("status", &self.status())
2563 .finish_non_exhaustive()
2564 }
2565}
2566
2567impl SupervisedModule {
2568 pub fn module_id(&self) -> &str {
2569 &self.inner.module_id
2570 }
2571
2572 #[cfg(test)]
2576 pub(crate) fn record_health_probe_failure_for_test(
2577 &self,
2578 detail: &str,
2579 ) -> Result<(), SuperviseError> {
2580 update_snapshot(&self.inner.snapshot, Some(&self.inner.module_id), |state| {
2581 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
2582 state.health.detail = Some(detail.to_string());
2583 })
2584 }
2585
2586 pub fn state(&self) -> Result<ModuleState, SuperviseError> {
2587 Ok(lock_snapshot(&self.inner.snapshot)?.state)
2588 }
2589
2590 pub fn stderr_tail(
2597 &self,
2598 max_lines: Option<usize>,
2599 max_bytes: Option<usize>,
2600 ) -> StderrTailSnapshot {
2601 self.inner
2602 .stderr_ring
2603 .lock()
2604 .unwrap_or_else(|poisoned| poisoned.into_inner())
2605 .snapshot(max_lines, max_bytes)
2606 }
2607
2608 pub fn terminal_history(&self) -> TerminalHistorySnapshot {
2613 self.inner
2614 .terminal_ring
2615 .lock()
2616 .unwrap_or_else(|poisoned| poisoned.into_inner())
2617 .snapshot()
2618 }
2619
2620 pub fn durable_terminal_history(&self) -> subc_control::TerminalHistory {
2625 durable_terminal_history_of(&self.inner.terminal_ring, &self.inner.module_id)
2626 }
2627
2628 pub(crate) async fn read_durable_terminal_history(
2633 &self,
2634 ) -> Result<subc_control::TerminalHistory, tokio::task::JoinError> {
2635 let terminal_ring = Arc::clone(&self.inner.terminal_ring);
2636 let module_id = self.inner.module_id.clone();
2637 tokio::task::spawn_blocking(move || durable_terminal_history_of(&terminal_ring, &module_id))
2638 .await
2639 }
2640
2641 pub fn status(&self) -> Result<ModuleStatus, SuperviseError> {
2642 self.status_with_snapshot_lock(&self.inner.snapshot, None)
2643 }
2644
2645 pub(crate) fn record_deliberate_severance(
2646 &self,
2647 identity: ProcessIdentity,
2648 ) -> Result<bool, SuperviseError> {
2649 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2650 if snapshot.pid != Some(identity.pid)
2651 || snapshot.process_start_time != Some(identity.start_time)
2652 {
2653 return Ok(false);
2654 }
2655 snapshot.deliberate_severance = Some(identity);
2656 Ok(true)
2657 }
2658
2659 pub(crate) fn status_for_control(
2664 &self,
2665 caller: &'static str,
2666 ) -> Result<ModuleStatus, SuperviseError> {
2667 self.status_with_snapshot_lock(&self.inner.snapshot, Some(caller))
2668 }
2669
2670 fn status_with_snapshot_lock(
2671 &self,
2672 snapshot: &SharedSnapshot,
2673 caller: Option<&'static str>,
2674 ) -> Result<ModuleStatus, SuperviseError> {
2675 let mut guard = match caller {
2676 Some(caller) => lock_snapshot_for_control(snapshot, &self.inner.module_id, caller)?,
2677 None => lock_snapshot(snapshot)?,
2678 };
2679 let restart_count =
2682 guard.crash_restarts_in_window(self.inner.restart_policy.window, Instant::now());
2683 let snapshot = guard.clone();
2684 drop(guard);
2685 let drain_timeout = *self.inner.effective_drain_timeout.lock().map_err(|_| {
2686 SuperviseError::StatePoisoned {
2687 module_id: Some(self.inner.module_id.clone()),
2688 }
2689 })?;
2690 let registration_active = self
2691 .inner
2692 .registry
2693 .get_module(&self.inner.module_id)
2694 .map_err(SuperviseError::Registry)?
2695 .is_some();
2696 let protocol = self.declared_protocol()?;
2697 let running_process =
2698 snapshot.enabled && snapshot.state == ModuleState::Running && snapshot.process_alive;
2699 let live = match protocol {
2705 ModuleProtocol::Subc => running_process && registration_active,
2706 ModuleProtocol::None => running_process,
2707 };
2708
2709 Ok(ModuleStatus {
2710 module_id: self.inner.module_id.clone(),
2711 state: snapshot.state,
2712 enabled: snapshot.enabled,
2713 process_alive: snapshot.process_alive,
2714 registration_active,
2715 protocol,
2716 live,
2717 restart_count,
2718 lifetime_restarts: snapshot.lifetime_restarts,
2719 spawn_generation: snapshot.spawn_generation,
2720 max_restarts: self.inner.restart_policy.max_restarts,
2721 restart_window: self.inner.restart_policy.window,
2722 drain_timeout,
2723 restart_backoff: self.inner.restart_policy.backoff,
2724 restart_max_backoff: self.inner.restart_policy.max_backoff,
2725 pid: snapshot.pid,
2726 spawned_at_ms: snapshot.spawned_at_ms,
2727 spawned_from: snapshot.spawned_from,
2728 process_start_time: snapshot.process_start_time,
2729 last_exit: snapshot.last_exit,
2730 health: snapshot.health,
2731 })
2732 }
2733
2734 #[cfg(test)]
2735 pub(crate) fn hold_snapshot_for_test(
2736 &self,
2737 acquired: std::sync::mpsc::Sender<()>,
2738 hold: Duration,
2739 ) -> std::thread::JoinHandle<()> {
2740 let snapshot = Arc::clone(&self.inner.snapshot);
2741 std::thread::spawn(move || {
2742 let _guard = snapshot.lock().expect("test snapshot lock is not poisoned");
2743 acquired
2744 .send(())
2745 .expect("test receiver waits for snapshot lock");
2746 std::thread::sleep(hold);
2747 })
2748 }
2749
2750 pub(crate) async fn running_image_agreement(&self) -> subc_control::RunningImageAgreement {
2751 let snapshot = match lock_snapshot(&self.inner.snapshot) {
2752 Ok(snapshot) => snapshot.clone(),
2753 Err(_) => {
2754 return subc_control::RunningImageAgreement::Unavailable {
2755 reason: subc_control::RunningImageUnavailableReason::NotRunning,
2756 };
2757 }
2758 };
2759 self.inner
2760 .provenance_probe
2761 .observe(
2762 snapshot.pid,
2763 snapshot.spawned_from.as_deref(),
2764 snapshot.spawned_file_identity,
2765 snapshot.process_start_time,
2766 )
2767 .await
2768 }
2769
2770 pub(crate) fn child_resource_usage(&self) -> subc_control::ChildResourceUsage {
2773 let (pid, start_time) = match lock_snapshot(&self.inner.snapshot) {
2774 Ok(snapshot) => (snapshot.pid, snapshot.process_start_time),
2775 Err(_) => {
2776 return subc_control::ChildResourceUsage::Unavailable {
2777 reason: subc_control::ChildResourceUnavailableReason::Unreadable,
2778 }
2779 }
2780 };
2781 crate::child_resources::read(pid, start_time)
2782 }
2783
2784 pub(crate) fn will_recover_after_connection_loss(&self) -> Result<bool, SuperviseError> {
2785 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2786 Ok(match snapshot.state {
2787 ModuleState::Restarting => true,
2788 ModuleState::Failed | ModuleState::Disabled => false,
2789 _ => daemon_will_restart(&mut snapshot, &self.inner.restart_policy, Instant::now()),
2790 })
2791 }
2792
2793 #[cfg(test)]
2794 pub(crate) fn is_warming(&self) -> Result<bool, SuperviseError> {
2795 self.is_warming_with_snapshot_lock(None)
2796 }
2797
2798 pub(crate) fn is_warming_for_control(
2799 &self,
2800 caller: &'static str,
2801 ) -> Result<bool, SuperviseError> {
2802 self.is_warming_with_snapshot_lock(Some(caller))
2803 }
2804
2805 fn is_warming_with_snapshot_lock(
2806 &self,
2807 caller: Option<&'static str>,
2808 ) -> Result<bool, SuperviseError> {
2809 let snapshot = match caller {
2810 Some(caller) => {
2811 lock_snapshot_for_control(&self.inner.snapshot, &self.inner.module_id, caller)?
2812 }
2813 None => lock_snapshot(&self.inner.snapshot)?,
2814 }
2815 .clone();
2816 Ok(matches!(
2817 snapshot.state,
2818 ModuleState::Starting | ModuleState::Running | ModuleState::Restarting
2819 ))
2820 }
2821
2822 pub async fn drain(&self) -> Result<(), SuperviseError> {
2824 self.stop().await
2825 }
2826
2827 pub(crate) async fn retire(&self) -> Result<(), SuperviseError> {
2828 match self.state()? {
2829 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2830 ModuleState::Starting
2831 | ModuleState::Running
2832 | ModuleState::Unresponsive
2833 | ModuleState::Restarting
2834 | ModuleState::Draining
2835 | ModuleState::Disabled => {}
2836 }
2837
2838 let (reply_tx, reply_rx) = oneshot::channel();
2839 self.inner
2840 .commands
2841 .send(SupervisorCommand::Retire { reply: reply_tx })
2842 .await
2843 .map_err(|_| SuperviseError::CommandClosed {
2844 module_id: self.inner.module_id.clone(),
2845 })?;
2846 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2847 module_id: self.inner.module_id.clone(),
2848 })?
2849 }
2850
2851 pub async fn stop(&self) -> Result<(), SuperviseError> {
2852 match self.state()? {
2853 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2854 ModuleState::Starting
2855 | ModuleState::Running
2856 | ModuleState::Unresponsive
2857 | ModuleState::Restarting
2858 | ModuleState::Draining
2859 | ModuleState::Disabled => {}
2860 }
2861
2862 let (reply_tx, reply_rx) = oneshot::channel();
2863 self.inner
2864 .commands
2865 .send(SupervisorCommand::Drain { reply: reply_tx })
2866 .await
2867 .map_err(|_| SuperviseError::CommandClosed {
2868 module_id: self.inner.module_id.clone(),
2869 })?;
2870 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2871 module_id: self.inner.module_id.clone(),
2872 })?
2873 }
2874
2875 pub async fn restart(&self, drain_timeout_ms: Option<u64>) -> Result<(), SuperviseError> {
2876 let received_at_generation = lock_snapshot(&self.inner.snapshot)?.spawn_generation;
2877 let (reply_tx, reply_rx) = oneshot::channel();
2878 self.inner
2879 .commands
2880 .send(SupervisorCommand::Restart {
2881 drain_timeout_ms,
2882 received_at_generation,
2883 queued_at: Instant::now(),
2884 reply: reply_tx,
2885 })
2886 .await
2887 .map_err(|_| SuperviseError::CommandClosed {
2888 module_id: self.inner.module_id.clone(),
2889 })?;
2890 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2891 module_id: self.inner.module_id.clone(),
2892 })?
2893 }
2894
2895 pub async fn swap(&self, ready_timeout: Option<Duration>) -> Result<(), SuperviseError> {
2900 let (reply_tx, reply_rx) = oneshot::channel();
2901 self.inner
2902 .commands
2903 .send(SupervisorCommand::Swap {
2904 ready_timeout,
2905 reply: reply_tx,
2906 })
2907 .await
2908 .map_err(|_| SuperviseError::CommandClosed {
2909 module_id: self.inner.module_id.clone(),
2910 })?;
2911 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2912 module_id: self.inner.module_id.clone(),
2913 })?
2914 }
2915
2916 pub async fn reload(&self) -> Result<(), SuperviseError> {
2917 let (reply_tx, reply_rx) = oneshot::channel();
2918 self.inner
2919 .commands
2920 .send(SupervisorCommand::Reload { reply: reply_tx })
2921 .await
2922 .map_err(|_| SuperviseError::CommandClosed {
2923 module_id: self.inner.module_id.clone(),
2924 })?;
2925 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2926 module_id: self.inner.module_id.clone(),
2927 })?
2928 }
2929
2930 pub async fn set_enabled(&self, enabled: bool) -> Result<bool, SuperviseError> {
2931 let (reply_tx, reply_rx) = oneshot::channel();
2932 self.inner
2933 .commands
2934 .send(SupervisorCommand::SetEnabled {
2935 enabled,
2936 reply: reply_tx,
2937 })
2938 .await
2939 .map_err(|_| SuperviseError::CommandClosed {
2940 module_id: self.inner.module_id.clone(),
2941 })?;
2942 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2943 module_id: self.inner.module_id.clone(),
2944 })?
2945 }
2946
2947 pub(crate) fn declared_protocol(&self) -> Result<ModuleProtocol, SuperviseError> {
2952 Ok(self
2953 .inner
2954 .configuration
2955 .lock()
2956 .map_err(|_| SuperviseError::StatePoisoned {
2957 module_id: Some(self.inner.module_id.clone()),
2958 })?
2959 .spec
2960 .protocol)
2961 }
2962
2963 pub(crate) fn configuration(&self) -> Result<(ModuleSpec, HealthConfig), SuperviseError> {
2964 let configuration =
2965 self.inner
2966 .configuration
2967 .lock()
2968 .map_err(|_| SuperviseError::StatePoisoned {
2969 module_id: Some(self.inner.module_id.clone()),
2970 })?;
2971 Ok((configuration.spec.clone(), configuration.health))
2972 }
2973
2974 #[cfg(any(test, feature = "test-support"))]
2978 pub async fn update_spec_for_test(&self, spec: ModuleSpec) -> Result<(), SuperviseError> {
2979 let (_, health) = self.configuration()?;
2980 let drain_timeout_ms = u64::try_from(
2981 self.inner
2982 .effective_drain_timeout
2983 .lock()
2984 .unwrap_or_else(|poisoned| poisoned.into_inner())
2985 .as_millis(),
2986 )
2987 .ok();
2988 self.update_configuration(spec, health, drain_timeout_ms)
2989 .await
2990 }
2991
2992 pub(crate) async fn update_configuration(
2993 &self,
2994 spec: ModuleSpec,
2995 health: HealthConfig,
2996 drain_timeout_ms: Option<u64>,
2997 ) -> Result<(), SuperviseError> {
2998 if spec.module_id != self.inner.module_id {
2999 return Err(SuperviseError::InvalidSpec {
3000 reason: "a supervised module's module_id cannot be changed".to_string(),
3001 });
3002 }
3003 validate_spec(&spec)?;
3004 let (reply_tx, reply_rx) = oneshot::channel();
3005 self.inner
3006 .commands
3007 .send(SupervisorCommand::UpdateConfiguration {
3008 spec: spec.clone(),
3009 health,
3010 drain_timeout_ms,
3011 reply: reply_tx,
3012 })
3013 .await
3014 .map_err(|_| SuperviseError::CommandClosed {
3015 module_id: self.inner.module_id.clone(),
3016 })?;
3017 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
3018 module_id: self.inner.module_id.clone(),
3019 })?;
3020 let mut configuration =
3021 self.inner
3022 .configuration
3023 .lock()
3024 .map_err(|_| SuperviseError::StatePoisoned {
3025 module_id: Some(self.inner.module_id.clone()),
3026 })?;
3027 configuration.spec = spec;
3028 configuration.health = health;
3029 Ok(())
3030 }
3031}
3032
3033impl Drop for SupervisedModuleInner {
3034 fn drop(&mut self) {
3035 let Ok(mut monitor) = self.monitor.lock() else {
3036 return;
3037 };
3038 if let Some(monitor) = monitor.as_ref().filter(|monitor| !monitor.is_finished()) {
3039 let _ = update_snapshot(&self.snapshot, Some(&self.module_id), |state| {
3040 state.state = ModuleState::Stopped;
3041 clear_current_process_facts(state);
3042 });
3043 monitor.abort();
3044 }
3045 let _ = monitor.take();
3046 }
3047}
3048
3049#[derive(Debug)]
3050enum SupervisorCommand {
3051 Drain {
3052 reply: oneshot::Sender<Result<(), SuperviseError>>,
3053 },
3054 Retire {
3055 reply: oneshot::Sender<Result<(), SuperviseError>>,
3056 },
3057 Restart {
3058 drain_timeout_ms: Option<u64>,
3063 received_at_generation: u64,
3067 queued_at: Instant,
3070 reply: oneshot::Sender<Result<(), SuperviseError>>,
3071 },
3072 Reload {
3073 reply: oneshot::Sender<Result<(), SuperviseError>>,
3074 },
3075 SetEnabled {
3076 enabled: bool,
3077 reply: oneshot::Sender<Result<bool, SuperviseError>>,
3078 },
3079 UpdateConfiguration {
3080 spec: ModuleSpec,
3081 health: HealthConfig,
3082 drain_timeout_ms: Option<u64>,
3085 reply: oneshot::Sender<()>,
3086 },
3087 Swap {
3088 ready_timeout: Option<Duration>,
3091 reply: oneshot::Sender<Result<(), SuperviseError>>,
3093 },
3094}
3095
3096#[derive(Debug)]
3097pub enum SuperviseError {
3098 InvalidSpec {
3099 reason: String,
3100 },
3101 Spawn {
3102 program: PathBuf,
3103 source: io::Error,
3104 cgroup_path: Option<PathBuf>,
3105 },
3106 Cgroup {
3107 module_id: String,
3108 source: io::Error,
3109 },
3110 LaunchNonce {
3113 reason: String,
3114 },
3115 Wait {
3116 module_id: String,
3117 source: io::Error,
3118 },
3119 Kill {
3120 module_id: String,
3121 source: io::Error,
3122 },
3123 Forwarding(ForwardingError),
3124 Registry(RegistryError),
3125 ReloadUnavailable {
3126 module_id: String,
3127 reason: String,
3128 },
3129 Disabled {
3134 module_id: String,
3135 },
3136 ReloadFailed {
3137 module_id: String,
3138 reason: String,
3139 },
3140 RegistrationStillActive {
3141 module_id: String,
3142 waited: Duration,
3143 },
3144 StatePoisoned {
3145 module_id: Option<String>,
3146 },
3147 CommandClosed {
3148 module_id: String,
3149 },
3150 SwapInProgress {
3154 module_id: String,
3155 },
3156 SwapRefused {
3158 module_id: String,
3159 reason: SwapRefusal,
3160 },
3161 SwapFailed {
3165 module_id: String,
3166 arm: SwapFailureArm,
3167 detail: String,
3168 candidate_exit: Option<ExitReport>,
3171 },
3172}
3173
3174#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3176pub enum SwapRefusal {
3177 OverlapExclusive,
3179 NotRegistered,
3182 ProtocolNone,
3185 NotConfigured,
3188 AlreadySwapping,
3190}
3191
3192impl SwapRefusal {
3193 pub fn as_str(self) -> &'static str {
3194 match self {
3195 Self::OverlapExclusive => "overlap_exclusive",
3196 Self::NotRegistered => "not_registered",
3197 Self::ProtocolNone => "protocol_none",
3198 Self::NotConfigured => "not_configured",
3199 Self::AlreadySwapping => "already_swapping",
3200 }
3201 }
3202}
3203
3204#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3207pub enum SwapFailureArm {
3208 SpawnFailed,
3210 NeverRegistered,
3212 NeverReady,
3214 CandidateExited,
3216 CandidateUnhealthy,
3218 Interrupted,
3222 CutoverLost,
3227}
3228
3229impl SwapFailureArm {
3230 pub fn as_str(self) -> &'static str {
3231 match self {
3232 Self::SpawnFailed => "spawn_failed",
3233 Self::NeverRegistered => "never_registered",
3234 Self::NeverReady => "never_ready",
3235 Self::CandidateExited => "candidate_exited",
3236 Self::CandidateUnhealthy => "candidate_unhealthy",
3237 Self::Interrupted => "interrupted",
3238 Self::CutoverLost => "cutover_lost",
3239 }
3240 }
3241}
3242
3243impl fmt::Display for SuperviseError {
3244 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3245 match self {
3246 Self::InvalidSpec { reason } => write!(f, "invalid module spec: {reason}"),
3247 Self::Spawn {
3248 program,
3249 source,
3250 cgroup_path: Some(cgroup_path),
3251 } => write!(
3252 f,
3253 "failed to place module in cgroup '{}' while spawning '{}': {source}",
3254 cgroup_path.display(),
3255 program.display()
3256 ),
3257 Self::Spawn {
3258 program,
3259 source,
3260 cgroup_path: None,
3261 } => write!(
3262 f,
3263 "failed to spawn module '{}': {source}",
3264 program.display()
3265 ),
3266 Self::Cgroup { module_id, source } => {
3267 write!(
3268 f,
3269 "failed to prepare cgroup for module '{module_id}': {source}"
3270 )
3271 }
3272 Self::LaunchNonce { reason } => {
3273 write!(
3274 f,
3275 "failed to generate reserved-module launch nonce: {reason}"
3276 )
3277 }
3278 Self::Wait { module_id, source } => {
3279 write!(f, "failed to wait for module '{module_id}': {source}")
3280 }
3281 Self::Kill { module_id, source } => {
3282 write!(f, "failed to kill module '{module_id}': {source}")
3283 }
3284 Self::Forwarding(err) => write!(f, "forwarding error: {err}"),
3285 Self::Registry(err) => write!(f, "registry error: {err}"),
3286 Self::ReloadUnavailable { module_id, reason } => {
3287 write!(f, "reload unavailable for module '{module_id}': {reason}")
3288 }
3289 Self::Disabled { module_id } => {
3290 write!(
3291 f,
3292 "module '{module_id}' is disabled; enable it before restart or reload"
3293 )
3294 }
3295 Self::ReloadFailed { module_id, reason } => {
3296 write!(f, "reload failed for module '{module_id}': {reason}")
3297 }
3298 Self::RegistrationStillActive { module_id, waited } => write!(
3299 f,
3300 "module '{module_id}' registration remained active after waiting {waited:?}"
3301 ),
3302 Self::StatePoisoned { module_id } => match module_id {
3303 Some(module_id) => {
3304 write!(f, "supervisor state for module '{module_id}' was poisoned")
3305 }
3306 None => write!(f, "supervisor state was poisoned"),
3307 },
3308 Self::CommandClosed { module_id } => {
3309 write!(
3310 f,
3311 "supervisor command channel for module '{module_id}' is closed"
3312 )
3313 }
3314 Self::SwapInProgress { module_id } => write!(
3315 f,
3316 "module '{module_id}' is being swapped; retry once the swap has cut over or failed, or stop the module to abort the swap"
3317 ),
3318 Self::SwapRefused { module_id, reason } => match reason {
3319 SwapRefusal::OverlapExclusive => write!(
3320 f,
3321 "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"
3322 ),
3323 SwapRefusal::NotRegistered => write!(
3324 f,
3325 "module '{module_id}' is not registered, so there is no serving process to keep while a replacement warms; use a plain restart"
3326 ),
3327 SwapRefusal::ProtocolNone => write!(
3328 f,
3329 "module '{module_id}' is protocol: \"none\" and never registers, so a swap could never see its replacement become ready; use a plain restart"
3330 ),
3331 SwapRefusal::NotConfigured => write!(
3332 f,
3333 "module '{module_id}' cannot be swapped: the supervisor was built without the forwarding table or shared handle a swap needs"
3334 ),
3335 SwapRefusal::AlreadySwapping => {
3336 write!(f, "module '{module_id}' is already being swapped")
3337 }
3338 },
3339 Self::SwapFailed {
3340 module_id,
3341 arm,
3342 detail,
3343 ..
3344 } => write!(
3345 f,
3346 "swap of module '{module_id}' failed ({}): {detail}; the running process was left serving",
3347 arm.as_str()
3348 ),
3349 }
3350 }
3351}
3352
3353impl Error for SuperviseError {
3354 fn source(&self) -> Option<&(dyn Error + 'static)> {
3355 match self {
3356 Self::Spawn { source, .. }
3357 | Self::Cgroup { source, .. }
3358 | Self::Wait { source, .. }
3359 | Self::Kill { source, .. } => Some(source),
3360 Self::Forwarding(err) => Some(err),
3361 Self::Registry(err) => Some(err),
3362 Self::LaunchNonce { .. }
3363 | Self::InvalidSpec { .. }
3364 | Self::ReloadUnavailable { .. }
3365 | Self::Disabled { .. }
3366 | Self::ReloadFailed { .. }
3367 | Self::RegistrationStillActive { .. }
3368 | Self::StatePoisoned { .. }
3369 | Self::CommandClosed { .. }
3370 | Self::SwapInProgress { .. }
3371 | Self::SwapRefused { .. }
3372 | Self::SwapFailed { .. } => None,
3373 }
3374 }
3375}
3376
3377pub(crate) fn validate_spec(spec: &ModuleSpec) -> Result<(), SuperviseError> {
3378 if spec.module_id.trim().is_empty() {
3379 return Err(SuperviseError::InvalidSpec {
3380 reason: "module_id must not be empty".to_string(),
3381 });
3382 }
3383
3384 Ok(())
3385}
3386
3387#[derive(Debug, Default)]
3388struct HealthProbeRuntime {
3389 registered_connection: Option<crate::ConnectionId>,
3390 advertised: bool,
3391 next_probe_at: Option<Instant>,
3392 probe_index: u64,
3393}
3394
3395impl HealthProbeRuntime {
3396 fn refresh_registration(
3397 &mut self,
3398 spec: &ModuleSpec,
3399 runtime: &SupervisorRuntimeConfig,
3400 registry: &Registry,
3401 snapshot: &SharedSnapshot,
3402 ) {
3403 if spec.protocol == ModuleProtocol::None {
3415 self.registered_connection = None;
3416 self.advertised = false;
3417 self.next_probe_at = None;
3418 return;
3419 }
3420
3421 let registration = match registry.get_module(&spec.module_id) {
3422 Ok(registration) => registration,
3423 Err(err) => {
3424 warn!(module_id = %spec.module_id, error = %err, "health prober could not read registry");
3425 self.advertised = false;
3426 self.next_probe_at = None;
3427 return;
3428 }
3429 };
3430
3431 let Some(registration) = registration else {
3432 self.registered_connection = None;
3433 self.advertised = false;
3434 self.next_probe_at = None;
3435 return;
3436 };
3437
3438 let advertised = registration
3439 .control_ops
3440 .iter()
3441 .any(|op| op == MODULE_CONTROL_OP_HEALTH_CHECK);
3442 if !advertised {
3443 self.registered_connection = Some(registration.connection_id);
3444 self.advertised = false;
3445 self.next_probe_at = None;
3446 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3447 state.health.status = SupervisorHealthStatus::Unknown;
3448 state.health.consecutive_failures = 0;
3449 state.health.last_probe_ms = None;
3450 state.health.detail = None;
3451 state.health.metrics = None;
3452 });
3453 return;
3454 }
3455
3456 let reregistered = self.registered_connection != Some(registration.connection_id);
3457 self.registered_connection = Some(registration.connection_id);
3458 self.advertised = true;
3459 if reregistered || self.next_probe_at.is_none() {
3460 self.probe_index = 0;
3461 self.next_probe_at = Some(
3462 Instant::now() + jittered_health_delay(&spec.module_id, 0, runtime.health.cadence),
3463 );
3464 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3465 state.health.status = SupervisorHealthStatus::Unknown;
3466 state.health.consecutive_failures = 0;
3467 state.health.detail = None;
3468 state.health.metrics = None;
3469 });
3470 }
3471 }
3472
3473 fn wake_after(&self) -> Duration {
3474 if !self.advertised {
3475 return REGISTRY_RELEASE_POLL;
3476 }
3477 self.next_probe_at
3478 .map(|next| next.saturating_duration_since(Instant::now()))
3479 .unwrap_or(REGISTRY_RELEASE_POLL)
3480 }
3481
3482 fn due(&self) -> bool {
3483 self.advertised
3484 && self
3485 .next_probe_at
3486 .is_some_and(|next| Instant::now() >= next)
3487 }
3488
3489 fn schedule_next(&mut self, spec: &ModuleSpec, cadence: Duration) {
3490 self.probe_index = self.probe_index.wrapping_add(1);
3491 self.next_probe_at = Some(
3492 Instant::now() + jittered_health_delay(&spec.module_id, self.probe_index, cadence),
3493 );
3494 }
3495}
3496
3497#[derive(Debug)]
3532enum HealthProbeEvidence {
3533 LaneDead,
3535 NoAnswer,
3537 BadAnswer,
3539 Misconfigured,
3541}
3542
3543#[derive(Debug)]
3544struct HealthProbeError {
3545 evidence: HealthProbeEvidence,
3546 message: String,
3547}
3548
3549impl HealthProbeError {
3550 fn lane_dead(message: impl Into<String>) -> Self {
3551 Self::with(HealthProbeEvidence::LaneDead, message)
3552 }
3553
3554 fn no_answer(message: impl Into<String>) -> Self {
3555 Self::with(HealthProbeEvidence::NoAnswer, message)
3556 }
3557
3558 fn bad_answer(message: impl Into<String>) -> Self {
3559 Self::with(HealthProbeEvidence::BadAnswer, message)
3560 }
3561
3562 fn misconfigured(message: impl Into<String>) -> Self {
3563 Self::with(HealthProbeEvidence::Misconfigured, message)
3564 }
3565
3566 fn with(evidence: HealthProbeEvidence, message: impl Into<String>) -> Self {
3567 Self {
3568 evidence,
3569 message: message.into(),
3570 }
3571 }
3572
3573 #[allow(dead_code)]
3587 fn is_proof_of_death(&self) -> bool {
3588 matches!(self.evidence, HealthProbeEvidence::LaneDead)
3589 }
3590
3591 fn label(&self) -> &'static str {
3599 match self.evidence {
3600 HealthProbeEvidence::LaneDead => "lane-dead",
3601 HealthProbeEvidence::NoAnswer => "no-answer",
3602 HealthProbeEvidence::BadAnswer => "bad-answer",
3603 HealthProbeEvidence::Misconfigured => "daemon-misconfigured",
3604 }
3605 }
3606}
3607
3608impl fmt::Display for HealthProbeError {
3609 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3610 f.write_str(&self.message)
3611 }
3612}
3613
3614async fn run_health_probe_cycle(
3615 spec: &ModuleSpec,
3616 runtime: &SupervisorRuntimeConfig,
3617 registry: &Registry,
3618 process_liveness: &SupervisorProcessLiveness,
3619 snapshot: &SharedSnapshot,
3620 child: &mut Option<SupervisedChild>,
3621) {
3622 let now_ms = unix_ms_now();
3623 match probe_module_health(&spec.module_id, runtime, None).await {
3624 Ok(report) => {
3625 handle_health_report(
3626 spec,
3627 runtime,
3628 registry,
3629 process_liveness,
3630 snapshot,
3631 child,
3632 report,
3633 now_ms,
3634 )
3635 .await;
3636 }
3637 Err(err) => {
3638 handle_health_probe_failure(
3639 spec,
3640 runtime,
3641 registry,
3642 process_liveness,
3643 snapshot,
3644 child,
3645 err,
3646 now_ms,
3647 )
3648 .await;
3649 }
3650 }
3651}
3652
3653async fn probe_module_health(
3654 module_id: &str,
3655 runtime: &SupervisorRuntimeConfig,
3656 drain_deadline: Option<Instant>,
3657) -> Result<HealthReport, HealthProbeError> {
3658 let Some(forwarding) = runtime.forwarding.as_ref() else {
3659 return Err(HealthProbeError::misconfigured(
3660 "supervisor was not configured with a forwarding table",
3661 ));
3662 };
3663 let probe_started_at = Instant::now();
3664 let mut deadline = probe_started_at + runtime.health.deadline;
3665 if let Some(drain_deadline) = drain_deadline {
3666 deadline = deadline.min(drain_deadline);
3667 }
3668 let pending = if drain_deadline.is_some() {
3669 forwarding.begin_drain_health_probe_rpc_for(
3670 module_id,
3671 MODULE_CONTROL_OP_HEALTH_CHECK,
3672 probe_started_at,
3673 deadline,
3674 )
3675 } else {
3676 forwarding.begin_health_probe_rpc_for(
3677 module_id,
3678 MODULE_CONTROL_OP_HEALTH_CHECK,
3679 probe_started_at,
3680 deadline,
3681 )
3682 }
3683 .map_err(|err| {
3684 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3687 })?;
3688 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3689}
3690
3691async fn probe_endpoint_health(
3698 endpoint: crate::ModuleEndpointId,
3699 runtime: &SupervisorRuntimeConfig,
3700 deadline_cap: Option<Instant>,
3701) -> Result<HealthReport, HealthProbeError> {
3702 let Some(forwarding) = runtime.forwarding.as_ref() else {
3703 return Err(HealthProbeError::misconfigured(
3704 "supervisor was not configured with a forwarding table",
3705 ));
3706 };
3707 let probe_started_at = Instant::now();
3708 let mut deadline = probe_started_at + runtime.health.deadline;
3709 if let Some(cap) = deadline_cap {
3710 deadline = deadline.min(cap);
3711 }
3712 let pending = forwarding
3713 .begin_endpoint_health_probe_rpc_for(
3714 endpoint,
3715 MODULE_CONTROL_OP_HEALTH_CHECK,
3716 probe_started_at,
3717 deadline,
3718 )
3719 .map_err(|err| {
3720 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3721 })?;
3722 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3723}
3724
3725async fn await_health_probe(
3727 forwarding: &ForwardingTable,
3728 pending: PendingModuleControlRpc,
3729 deadline: Instant,
3730 probe_budget: Duration,
3731) -> Result<HealthReport, HealthProbeError> {
3732 let PendingModuleControlRpc {
3733 endpoint,
3734 module_sink,
3735 negotiated_ver,
3736 corr,
3737 receiver,
3738 } = pending;
3739 let body = serde_json::to_vec(&ModuleControlRequest::HealthCheck {}).map_err(|err| {
3740 HealthProbeError::misconfigured(format!("failed to encode health.check: {err}"))
3741 })?;
3742 let frame = Frame::build_with_version(
3743 negotiated_ver,
3744 FrameType::Request,
3745 control_flags(),
3746 0,
3747 0,
3748 corr,
3749 body,
3750 )
3751 .map_err(|err| {
3752 HealthProbeError::misconfigured(format!("failed to build health.check frame: {err}"))
3753 })?;
3754
3755 match timeout_at(deadline, module_sink.send(frame)).await {
3761 Ok(Ok(())) => {}
3762 Ok(Err(err)) => {
3763 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3764 return Err(HealthProbeError::lane_dead(format!(
3767 "failed to send health.check: {err}"
3768 )));
3769 }
3770 Err(_elapsed) => {
3771 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3772 return Err(HealthProbeError::no_answer(
3776 "health.check send timed out before enqueue (module egress full)",
3777 ));
3778 }
3779 }
3780
3781 match timeout_at(deadline, receiver).await {
3782 Ok(Ok(ModuleControlRpcOutcome::Response(response))) => {
3786 response.health_report().ok_or_else(|| {
3787 HealthProbeError::bad_answer("health.check RPC returned a non-health response")
3788 })
3789 }
3790 Ok(Ok(ModuleControlRpcOutcome::Rejected(body))) => Err(HealthProbeError::bad_answer(
3791 format!("health.check rejected: {}", body.message),
3792 )),
3793 Ok(Ok(ModuleControlRpcOutcome::ModuleGone(message))) => {
3794 Err(HealthProbeError::lane_dead(message))
3795 }
3796 Ok(Ok(ModuleControlRpcOutcome::MalformedResponse(message))) => {
3797 Err(HealthProbeError::bad_answer(message))
3798 }
3799 Ok(Ok(ModuleControlRpcOutcome::UnexpectedOp { expected, actual })) => {
3800 Err(HealthProbeError::bad_answer(format!(
3801 "expected module-control op '{expected}', got '{actual}'"
3802 )))
3803 }
3804 Ok(Ok(ModuleControlRpcOutcome::DeadlineElapsed)) => Err(HealthProbeError::bad_answer(
3808 "module answered health.check after its daemon deadline",
3809 )),
3810 Ok(Err(_)) => Err(HealthProbeError::misconfigured(
3811 "health.check waiter was canceled before the module responded",
3812 )),
3813 Err(_) => {
3814 let _ = forwarding.tombstone_health_probe_rpc(endpoint, corr);
3815 Err(HealthProbeError::no_answer(format!(
3816 "module did not answer health.check within {probe_budget:?}"
3817 )))
3818 }
3819 }
3820}
3821
3822#[allow(clippy::too_many_arguments)]
3823async fn handle_health_report(
3824 spec: &ModuleSpec,
3825 runtime: &SupervisorRuntimeConfig,
3826 registry: &Registry,
3827 process_liveness: &SupervisorProcessLiveness,
3828 snapshot: &SharedSnapshot,
3829 child: &mut Option<SupervisedChild>,
3830 report: HealthReport,
3831 now_ms: u64,
3832) {
3833 let status = supervisor_health_status(report.status);
3834 let detail = report.detail.clone();
3835 let metrics = truncate_health_metrics(report.metrics);
3836 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3837 state.health.status = status;
3838 state.health.last_probe_ms = Some(now_ms);
3839 state.health.detail = detail.clone();
3840 state.health.metrics = metrics.clone();
3841 state.health.consecutive_failures = 0;
3842 });
3843
3844 let action = match report.status {
3845 HealthStatus::Ok => return,
3846 HealthStatus::Degraded => runtime.health.on_degraded,
3847 HealthStatus::Failing => runtime.health.on_failing,
3848 };
3849 apply_l3_health_action(
3850 spec,
3851 runtime,
3852 registry,
3853 process_liveness,
3854 snapshot,
3855 child,
3856 status,
3857 detail.as_deref(),
3858 action,
3859 now_ms,
3860 )
3861 .await;
3862}
3863
3864#[allow(clippy::too_many_arguments)]
3865async fn handle_health_probe_failure(
3866 spec: &ModuleSpec,
3867 runtime: &SupervisorRuntimeConfig,
3868 registry: &Registry,
3869 process_liveness: &SupervisorProcessLiveness,
3870 snapshot: &SharedSnapshot,
3871 child: &mut Option<SupervisedChild>,
3872 err: HealthProbeError,
3873 now_ms: u64,
3874) {
3875 let threshold = runtime.health.failure_threshold.max(1);
3876 let mut failures = 0;
3877 let detail = format!("[{}] {err}", err.label());
3882 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3883 state.health.last_probe_ms = Some(now_ms);
3884 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
3885 state.health.detail = Some(detail.clone());
3886 state.health.metrics = None;
3887 failures = state.health.consecutive_failures;
3888 });
3889
3890 if failures < threshold {
3891 warn!(
3892 module_id = %spec.module_id,
3893 consecutive_failures = failures,
3894 threshold,
3895 evidence = err.label(),
3896 detail = %detail,
3897 "health.check probe failed"
3898 );
3899 return;
3900 }
3901
3902 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3903 state.state = ModuleState::Unresponsive;
3904 state.health.status = SupervisorHealthStatus::Unresponsive;
3905 });
3906 if runtime.health.critical {
3910 error!(
3911 module_id = %spec.module_id,
3912 status = "unresponsive",
3913 evidence = err.label(),
3914 detail = %detail,
3915 "critical module health alert"
3916 );
3917 } else {
3918 warn!(
3919 module_id = %spec.module_id,
3920 status = "unresponsive",
3921 evidence = err.label(),
3922 detail = %detail,
3923 "module health threshold breached"
3924 );
3925 }
3926 if let Err(err) = health_restart_child(
3927 spec,
3928 runtime,
3929 registry,
3930 process_liveness,
3931 snapshot,
3932 child,
3933 SupervisorHealthStatus::Unresponsive,
3934 Some(&detail),
3935 now_ms,
3936 )
3937 .await
3938 {
3939 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3940 }
3941}
3942
3943#[allow(clippy::too_many_arguments)]
3944async fn apply_l3_health_action(
3945 spec: &ModuleSpec,
3946 runtime: &SupervisorRuntimeConfig,
3947 registry: &Registry,
3948 process_liveness: &SupervisorProcessLiveness,
3949 snapshot: &SharedSnapshot,
3950 child: &mut Option<SupervisedChild>,
3951 status: SupervisorHealthStatus,
3952 detail: Option<&str>,
3953 action: HealthAction,
3954 now_ms: u64,
3955) {
3956 record_health_action(snapshot, &spec.module_id, action.to_string(), now_ms);
3957 match action {
3958 HealthAction::Report => {
3959 info!(
3960 module_id = %spec.module_id,
3961 status = ?status,
3962 detail,
3963 "module reported non-ok health"
3964 );
3965 }
3966 HealthAction::Alert => {
3967 error!(
3968 module_id = %spec.module_id,
3969 status = ?status,
3970 detail,
3971 "module health alert"
3972 );
3973 }
3974 HealthAction::Restart => {
3975 if let Err(err) = health_restart_child(
3976 spec,
3977 runtime,
3978 registry,
3979 process_liveness,
3980 snapshot,
3981 child,
3982 status,
3983 detail,
3984 now_ms,
3985 )
3986 .await
3987 {
3988 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3989 }
3990 }
3991 }
3992}
3993
3994#[allow(clippy::too_many_arguments)]
3995async fn health_restart_child(
3996 spec: &ModuleSpec,
3997 runtime: &SupervisorRuntimeConfig,
3998 registry: &Registry,
3999 process_liveness: &SupervisorProcessLiveness,
4000 snapshot: &SharedSnapshot,
4001 child: &mut Option<SupervisedChild>,
4002 status: SupervisorHealthStatus,
4003 detail: Option<&str>,
4004 now_ms: u64,
4005) -> Result<(), SuperviseError> {
4006 let (enabled, schedule) = {
4007 let mut state = lock_snapshot(snapshot)?;
4008 let enabled = state.enabled;
4009 let schedule = if enabled {
4010 state.next_crash_restart(&runtime.restart_policy, Instant::now())
4011 } else {
4012 None
4013 };
4014 (enabled, schedule)
4015 };
4016
4017 if !enabled {
4018 return Err(SuperviseError::Disabled {
4019 module_id: spec.module_id.clone(),
4020 });
4021 }
4022
4023 if schedule.is_none() {
4024 record_health_action(snapshot, &spec.module_id, "disabled".to_string(), now_ms);
4025 error!(
4026 module_id = %spec.module_id,
4027 status = ?status,
4028 detail,
4029 max_restarts = runtime.restart_policy.max_restarts,
4030 window_secs = runtime.restart_policy.window.as_secs(),
4031 "health restart budget exhausted; disabling module"
4032 );
4033 let stop_notice = begin_forwarding_drain_if_configured(
4034 spec,
4035 runtime,
4036 registry,
4037 snapshot,
4038 Some(false),
4039 RouteCloseReason::Disable,
4040 )
4041 .await?;
4042 drain_optional_child(
4043 &spec.module_id,
4044 spec.protocol,
4045 stop_notice,
4046 registry,
4047 snapshot,
4048 &runtime.terminal_ring,
4049 &runtime.spawn_events,
4050 child,
4051 runtime.drain_timeout,
4052 ModuleState::Disabled,
4053 Some(false),
4054 )
4055 .await?;
4056 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4057 return Ok(());
4058 }
4059
4060 let schedule = schedule.expect("a health restart must have a crash-restart schedule");
4061 let mut restart_count = 0;
4062 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4063 restart_count = state.crash_restarts.len();
4064 state.state = ModuleState::Unresponsive;
4065 state.health.status = status;
4066 state.health.last_action = Some(HealthAction::Restart.to_string());
4067 state.health.last_action_ms = Some(now_ms);
4068 })?;
4069 warn!(
4070 module_id = %spec.module_id,
4071 status = ?status,
4072 detail,
4073 restart_count,
4074 restart_in_window = schedule.restart_in_window,
4075 delay_ms = schedule.delay.as_millis() as u64,
4076 "health-triggered module restart"
4077 );
4078
4079 let stop_notice = begin_forwarding_drain_if_configured(
4080 spec,
4081 runtime,
4082 registry,
4083 snapshot,
4084 Some(true),
4085 RouteCloseReason::Restart,
4086 )
4087 .await?;
4088 drain_optional_child(
4089 &spec.module_id,
4090 spec.protocol,
4091 stop_notice,
4092 registry,
4093 snapshot,
4094 &runtime.terminal_ring,
4095 &runtime.spawn_events,
4096 child,
4097 runtime.drain_timeout,
4098 ModuleState::Restarting,
4099 Some(true),
4100 )
4101 .await?;
4102 sleep(schedule.delay).await;
4103 if !respawn_still_pending(snapshot) {
4107 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4108 return Ok(());
4109 }
4110 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4111 match spawn_and_mark_running(spec, runtime, snapshot) {
4112 Ok(next_child) => {
4113 *child = Some(next_child);
4114 Ok(())
4115 }
4116 Err(err) => {
4117 fail_snapshot(snapshot, Some(&spec.module_id), None);
4118 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4119 *child = None;
4120 Err(err)
4121 }
4122 }
4123}
4124
4125fn record_health_action(snapshot: &SharedSnapshot, module_id: &str, action: String, now_ms: u64) {
4126 let _ = update_snapshot(snapshot, Some(module_id), |state| {
4127 state.health.last_action = Some(action);
4128 state.health.last_action_ms = Some(now_ms);
4129 });
4130}
4131
4132fn supervisor_health_status(status: HealthStatus) -> SupervisorHealthStatus {
4133 match status {
4134 HealthStatus::Ok => SupervisorHealthStatus::Ok,
4135 HealthStatus::Degraded => SupervisorHealthStatus::Degraded,
4136 HealthStatus::Failing => SupervisorHealthStatus::Failing,
4137 }
4138}
4139
4140fn truncate_health_metrics(metrics: Option<Value>) -> Option<Value> {
4152 let metrics = metrics?;
4153 match serde_json::to_vec(&metrics) {
4154 Ok(encoded) if encoded.len() > MAX_HEALTH_METRICS_BYTES => Some(serde_json::json!({
4155 "truncated": true,
4156 "original_bytes": encoded.len(),
4157 })),
4158 Ok(_) | Err(_) => Some(metrics),
4159 }
4160}
4161
4162fn jittered_health_delay(module_id: &str, probe_index: u64, cadence: Duration) -> Duration {
4168 if cadence.is_zero() {
4169 return Duration::ZERO;
4170 }
4171 let cadence_ms = cadence.as_millis() as u64;
4172 if cadence_ms == 0 {
4188 return cadence;
4189 }
4190 let jitter_span = (cadence_ms / 10).max(1);
4205 let hash = module_id.as_bytes().iter().fold(
4206 probe_index.wrapping_mul(0x9E37_79B9_7F4A_7C15),
4207 |acc, byte| {
4208 acc.wrapping_mul(1099511628211)
4209 .wrapping_add(u64::from(*byte))
4210 },
4211 );
4212 cadence + Duration::from_millis(hash % jitter_span)
4213}
4214
4215#[cfg(test)]
4216mod tests {
4217 use super::*;
4218
4219 #[test]
4220 fn readding_a_module_clears_its_rescan_removal_tombstone() {
4221 let handle = SupervisorHandle::new();
4222 let module_id = "readded-tombstone";
4223 handle.record_rescan_removal(module_id);
4224 assert!(handle.removal_tombstone_age_ms(module_id).is_some());
4225
4226 handle.apply_identity_configuration(&ModuleSpec {
4227 module_id: module_id.to_string(),
4228 program: PathBuf::from("/test/module"),
4229 args: Vec::new(),
4230 env: Vec::new(),
4231 reserved: false,
4232 reserved_prefixes: Vec::new(),
4233 protocol: ModuleProtocol::Subc,
4234 overlap: Default::default(),
4235 });
4236
4237 assert!(
4238 handle.removal_tombstone_age_ms(module_id).is_none(),
4239 "a re-added module must not retain a stale removal tombstone"
4240 );
4241 }
4242
4243 fn stale_process_snapshot(state: ModuleState, enabled: bool) -> SharedSnapshot {
4244 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::new(state, enabled)));
4245 update_snapshot(&snapshot, Some("stale-process-facts"), |snapshot| {
4246 snapshot.process_alive = true;
4247 snapshot.pid = Some(41);
4248 snapshot.spawned_at_ms = Some(42);
4249 snapshot.spawned_from = Some(PathBuf::from("/spawned/module"));
4250 snapshot.spawned_file_identity = Some(SpawnedFileIdentity {
4251 device: 43,
4252 inode: 44,
4253 });
4254 })
4255 .unwrap();
4256 snapshot
4257 }
4258
4259 fn assert_snapshot_process_facts_cleared(snapshot: &SharedSnapshot) {
4260 let snapshot = lock_snapshot(snapshot).unwrap();
4261 assert!(!snapshot.process_alive);
4262 assert_eq!(snapshot.pid, None);
4263 assert_eq!(snapshot.spawned_at_ms, None);
4264 assert_eq!(snapshot.spawned_from, None);
4265 assert_eq!(snapshot.spawned_file_identity, None);
4266 }
4267
4268 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4269 async fn failed_enable_spawn_clears_preexisting_current_process_facts() {
4270 let supervisor = Supervisor::default();
4271 let mut runtime = supervisor.runtime_config();
4272 runtime.test_seed_stale_facts_before_enable_spawn = true;
4273 let snapshot = stale_process_snapshot(ModuleState::Disabled, false);
4274 let mut child = None;
4275 let spec = ModuleSpec {
4276 module_id: "failed-enable-clears-facts".to_string(),
4277 program: PathBuf::from("/definitely/missing/failed-enable-module"),
4278 args: Vec::new(),
4279 env: Vec::new(),
4280 reserved: false,
4281 reserved_prefixes: Vec::new(),
4282 protocol: ModuleProtocol::Subc,
4283 overlap: Default::default(),
4284 };
4285
4286 let result = set_child_enabled(
4287 &spec,
4288 &runtime,
4289 &supervisor.registry,
4290 &supervisor.process_liveness,
4291 &snapshot,
4292 &mut child,
4293 true,
4294 )
4295 .await;
4296
4297 assert!(matches!(result, Err(SuperviseError::Spawn { .. })));
4298 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4299 assert_snapshot_process_facts_cleared(&snapshot);
4300 }
4301
4302 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4303 async fn failed_reload_spawn_clears_current_process_facts() {
4304 let supervisor = Supervisor::default();
4305 let mut runtime = supervisor.runtime_config();
4306 runtime.restart_policy = RestartPolicy::new(0, Duration::ZERO);
4307 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4308 let mut child = None;
4309 let spec = ModuleSpec {
4310 module_id: "failed-reload-clears-facts".to_string(),
4311 program: PathBuf::from("/unused/failed-reload-module"),
4312 args: Vec::new(),
4313 env: Vec::new(),
4314 reserved: false,
4315 reserved_prefixes: Vec::new(),
4316 protocol: ModuleProtocol::Subc,
4317 overlap: Default::default(),
4318 };
4319
4320 let result = handle_reload_spawn_failure(
4321 &spec,
4322 &runtime,
4323 &supervisor.process_liveness,
4324 &snapshot,
4325 &mut child,
4326 "forced reload spawn failure".to_string(),
4327 )
4328 .await;
4329
4330 assert!(matches!(result, Err(SuperviseError::ReloadFailed { .. })));
4331 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4332 assert_snapshot_process_facts_cleared(&snapshot);
4333 }
4334
4335 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4336 async fn dropping_a_module_with_an_active_monitor_clears_current_process_facts() {
4337 let supervisor = Supervisor::default();
4338 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4339 let module = supervisor.supervised_module(
4340 ModuleSpec {
4341 module_id: "drop-clears-facts".to_string(),
4342 program: PathBuf::from("/unused/drop-module"),
4343 args: Vec::new(),
4344 env: Vec::new(),
4345 reserved: false,
4346 reserved_prefixes: Vec::new(),
4347 protocol: ModuleProtocol::Subc,
4348 overlap: Default::default(),
4349 },
4350 supervisor.runtime_config(),
4351 Arc::clone(&snapshot),
4352 None,
4353 );
4354 assert!(!module
4355 .inner
4356 .monitor
4357 .lock()
4358 .unwrap()
4359 .as_ref()
4360 .unwrap()
4361 .is_finished());
4362
4363 drop(module);
4364
4365 assert_eq!(
4366 lock_snapshot(&snapshot).unwrap().state,
4367 ModuleState::Stopped
4368 );
4369 assert_snapshot_process_facts_cleared(&snapshot);
4370 }
4371
4372 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4373 async fn configuration_update_does_not_replace_captured_running_process_facts() {
4374 let supervisor = Supervisor::default();
4375 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4376 let initial = ModuleSpec {
4377 module_id: "rescan-preserves-spawn-facts".to_string(),
4378 program: PathBuf::from("/spawned/module"),
4379 args: Vec::new(),
4380 env: Vec::new(),
4381 reserved: false,
4382 reserved_prefixes: Vec::new(),
4383 protocol: ModuleProtocol::Subc,
4384 overlap: Default::default(),
4385 };
4386 let module = supervisor.supervised_module(
4387 initial.clone(),
4388 supervisor.runtime_config(),
4389 snapshot,
4390 None,
4391 );
4392 let before = module.status().unwrap();
4393 let mut replacement = initial;
4394 replacement.program = PathBuf::from("/rescanned/replacement-module");
4395
4396 module
4397 .update_configuration(replacement, HealthConfig::default(), None)
4398 .await
4399 .unwrap();
4400
4401 let after = module.status().unwrap();
4402 assert_eq!(after.pid, before.pid);
4403 assert_eq!(after.spawned_at_ms, before.spawned_at_ms);
4404 assert_eq!(after.spawned_from, before.spawned_from);
4405 drop(module);
4406 }
4407}
4408
4409fn unix_ms_now() -> u64 {
4410 SystemTime::now()
4411 .duration_since(UNIX_EPOCH)
4412 .map(|duration| duration.as_millis().min(u128::from(u64::MAX)) as u64)
4413 .unwrap_or(0)
4414}
4415
4416async fn supervise_loop(
4417 mut spec: ModuleSpec,
4418 mut runtime: SupervisorRuntimeConfig,
4419 registry: Arc<Registry>,
4420 process_liveness: Arc<SupervisorProcessLiveness>,
4421 snapshot: SharedSnapshot,
4422 mut child: Option<SupervisedChild>,
4423 mut commands: mpsc::Receiver<SupervisorCommand>,
4424) {
4425 let mut health_probe = HealthProbeRuntime::default();
4426 let mut pending_respawn: Option<Instant> = None;
4430 let mut requeued: VecDeque<SupervisorCommand> = VecDeque::new();
4433 loop {
4434 if let Some(command) = requeued.pop_front() {
4435 if !handle_supervisor_command(
4436 command,
4437 &mut spec,
4438 &mut runtime,
4439 ®istry,
4440 &process_liveness,
4441 &snapshot,
4442 &mut child,
4443 &mut commands,
4444 &mut requeued,
4445 )
4446 .await
4447 {
4448 return;
4449 }
4450 if child.is_some() || !respawn_still_pending(&snapshot) {
4451 pending_respawn = None;
4452 }
4453 continue;
4454 }
4455 if child.is_some() {
4456 health_probe.refresh_registration(&spec, &runtime, ®istry, &snapshot);
4457 let probe_sleep = sleep(health_probe.wake_after());
4458 tokio::pin!(probe_sleep);
4459 let active_child = child.as_mut().expect("child checked above");
4460 tokio::select! {
4461 wait_result = active_child.wait() => {
4462 let exit_report = match wait_result {
4471 Ok(status) => classify_reaped_child_exit(&snapshot, active_child, &status),
4472 Err(err) => {
4473 active_child.drain_stderr(&spec.module_id).await;
4474 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4475 record_wait_error_terminal(
4481 &spec.module_id,
4482 &runtime.terminal_ring,
4483 &runtime.spawn_events,
4484 );
4485 untrack_if_registration_released(
4486 &process_liveness,
4487 ®istry,
4488 &spec.module_id,
4489 &snapshot,
4490 );
4491 error!(module_id = %spec.module_id, error = %err, "failed to wait for supervised module");
4492 child = None;
4493 continue;
4494 }
4495 };
4496 active_child.drain_stderr(&spec.module_id).await;
4497
4498 let next = on_child_exit(
4499 &spec,
4500 runtime.restart_policy,
4501 ®istry,
4502 &snapshot,
4503 &runtime.terminal_ring,
4504 &runtime.spawn_events,
4505 &runtime.child_roster,
4506 exit_report,
4507 ).await;
4508 active_child.release_roster();
4511 match next {
4512 NextAction::Stop { registration_released } => {
4513 if registration_released {
4514 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4515 }
4516 child = None;
4517 }
4518 NextAction::Restart { schedule } => {
4519 let delay = schedule.map_or(
4520 runtime.restart_policy.delay_for_restart(0),
4521 |schedule| schedule.delay,
4522 );
4523 if let Some(schedule) = schedule {
4524 log_crash_respawn(&spec.module_id, schedule);
4525 }
4526 child = None;
4534 pending_respawn = Some(Instant::now() + delay);
4535 }
4536 }
4537 }
4538 command = commands.recv() => {
4539 let Some(command) = command else {
4540 return;
4541 };
4542 if !handle_supervisor_command(
4543 command,
4544 &mut spec,
4545 &mut runtime,
4546 ®istry,
4547 &process_liveness,
4548 &snapshot,
4549 &mut child,
4550 &mut commands,
4551 &mut requeued,
4552 ).await {
4553 return;
4554 }
4555 }
4556 _ = &mut probe_sleep => {
4557 if health_probe.due() {
4558 run_health_probe_cycle(
4559 &spec,
4560 &runtime,
4561 ®istry,
4562 &process_liveness,
4563 &snapshot,
4564 &mut child,
4565 ).await;
4566 if child.is_some() {
4567 health_probe.schedule_next(&spec, runtime.health.cadence);
4568 }
4569 }
4570 }
4571 }
4572 } else if let Some(deadline) = pending_respawn {
4573 tokio::select! {
4574 _ = sleep_until(deadline) => {
4575 pending_respawn = None;
4576 if !respawn_still_pending(&snapshot) {
4580 continue;
4581 }
4582 if runtime.child_roster.is_closed() {
4587 let _ = update_snapshot(&snapshot, Some(&spec.module_id), |state| {
4588 state.state = ModuleState::Stopped;
4589 });
4590 debug!(module_id = %spec.module_id, "crash respawn cancelled by daemon shutdown");
4591 continue;
4592 }
4593 if let Err(err) = wait_for_registration_release(
4594 ®istry,
4595 &spec.module_id,
4596 REGISTRY_RELEASE_TIMEOUT,
4597 ).await {
4598 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4599 error!(module_id = %spec.module_id, error = %err, "registration did not release before restart");
4600 continue;
4601 }
4602
4603 match spawn_and_mark_running(&spec, &runtime, &snapshot) {
4604 Ok(next_child) => {
4605 child = Some(next_child);
4606 debug!(module_id = %spec.module_id, "supervised module restarted after crash");
4607 }
4608 Err(err) => {
4609 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4610 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4611 error!(module_id = %spec.module_id, error = %err, "failed to restart supervised module");
4612 }
4613 }
4614 }
4615 command = commands.recv() => {
4616 let Some(command) = command else {
4617 return;
4618 };
4619 if !handle_supervisor_command(
4620 command,
4621 &mut spec,
4622 &mut runtime,
4623 ®istry,
4624 &process_liveness,
4625 &snapshot,
4626 &mut child,
4627 &mut commands,
4628 &mut requeued,
4629 ).await {
4630 return;
4631 }
4632 if child.is_some() || !respawn_still_pending(&snapshot) {
4637 pending_respawn = None;
4638 }
4639 }
4640 }
4641 } else {
4642 let Some(command) = commands.recv().await else {
4643 return;
4644 };
4645 if !handle_supervisor_command(
4646 command,
4647 &mut spec,
4648 &mut runtime,
4649 ®istry,
4650 &process_liveness,
4651 &snapshot,
4652 &mut child,
4653 &mut commands,
4654 &mut requeued,
4655 )
4656 .await
4657 {
4658 return;
4659 }
4660 }
4661 }
4662}
4663
4664fn log_crash_respawn(module_id: &str, schedule: CrashRestartSchedule) {
4665 info!(
4666 module_id,
4667 restart_in_window = schedule.restart_in_window,
4668 delay_ms = schedule.delay.as_millis() as u64,
4669 "respawning after crash"
4670 );
4671}
4672
4673fn respawn_still_pending(snapshot: &SharedSnapshot) -> bool {
4679 matches!(
4680 lock_snapshot(snapshot),
4681 Ok(state) if state.enabled && state.state == ModuleState::Restarting
4682 )
4683}
4684
4685enum NextAction {
4686 Stop {
4687 registration_released: bool,
4688 },
4689 Restart {
4690 schedule: Option<CrashRestartSchedule>,
4691 },
4692}
4693
4694#[allow(clippy::too_many_arguments)]
4695async fn handle_supervisor_command(
4696 command: SupervisorCommand,
4697 spec: &mut ModuleSpec,
4698 runtime: &mut SupervisorRuntimeConfig,
4699 registry: &Registry,
4700 process_liveness: &SupervisorProcessLiveness,
4701 snapshot: &SharedSnapshot,
4702 child: &mut Option<SupervisedChild>,
4703 commands: &mut mpsc::Receiver<SupervisorCommand>,
4704 requeued: &mut VecDeque<SupervisorCommand>,
4705) -> bool {
4706 match command {
4707 SupervisorCommand::Drain { reply } => {
4708 let result = drain_optional_child(
4711 &spec.module_id,
4712 spec.protocol,
4713 StopNotice::NotSent,
4714 registry,
4715 snapshot,
4716 &runtime.terminal_ring,
4717 &runtime.spawn_events,
4718 child,
4719 runtime.drain_timeout,
4720 ModuleState::Stopped,
4721 None,
4722 )
4723 .await;
4724 let registration_released = result.is_ok();
4725 let _ = reply.send(result);
4726 if registration_released {
4727 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4728 }
4729 false
4730 }
4731 SupervisorCommand::Retire { reply } => {
4732 let result = async {
4733 let stop_notice = begin_forwarding_drain_if_configured(
4734 spec,
4735 runtime,
4736 registry,
4737 snapshot,
4738 None,
4739 RouteCloseReason::Disable,
4740 )
4741 .await?;
4742 drain_optional_child(
4743 &spec.module_id,
4744 spec.protocol,
4745 stop_notice,
4746 registry,
4747 snapshot,
4748 &runtime.terminal_ring,
4749 &runtime.spawn_events,
4750 child,
4751 runtime.drain_timeout,
4752 ModuleState::Stopped,
4753 None,
4754 )
4755 .await
4756 }
4757 .await;
4758 let registration_released = result.is_ok();
4759 let _ = reply.send(result);
4760 if registration_released {
4761 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4762 }
4763 false
4764 }
4765 SupervisorCommand::Restart {
4766 drain_timeout_ms,
4767 received_at_generation,
4768 queued_at,
4769 reply,
4770 } => {
4771 info!(
4775 module_id = %spec.module_id,
4776 queued_ms = u64::try_from(queued_at.elapsed().as_millis()).unwrap_or(u64::MAX),
4777 "restart command dequeued"
4778 );
4779 let validation = match lock_snapshot(snapshot) {
4791 Ok(state) if !state.enabled => Err(SuperviseError::Disabled {
4792 module_id: spec.module_id.clone(),
4793 }),
4794 Ok(_) => Ok(()),
4795 Err(err) => Err(err),
4796 };
4797 let initiated = validation.is_ok();
4798 let _ = reply.send(validation);
4799 let satisfied_by_generation = if initiated && child.is_some() {
4810 lock_snapshot(snapshot).ok().and_then(|state| {
4811 (state.spawn_generation > received_at_generation
4812 && !state.configuration_updated_since_spawn)
4813 .then_some(state.spawn_generation)
4814 })
4815 } else {
4816 None
4817 };
4818 if let Some(generation) = satisfied_by_generation {
4819 info!(
4820 module_id = %spec.module_id,
4821 received_at_generation,
4822 "restart already satisfied by generation {generation}; not restarting again"
4823 );
4824 } else if initiated {
4825 let drain_timeout = drain_timeout_ms
4828 .map(Duration::from_millis)
4829 .unwrap_or(runtime.drain_timeout);
4830 if let Err(err) = restart_child(
4831 spec,
4832 runtime,
4833 registry,
4834 process_liveness,
4835 snapshot,
4836 child,
4837 drain_timeout,
4838 )
4839 .await
4840 {
4841 warn!(
4842 module_id = %spec.module_id,
4843 error = %err,
4844 "operator restart failed after initiation ack; module state carries the outcome"
4845 );
4846 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4847 state.state = ModuleState::Failed;
4848 clear_current_process_facts(state);
4849 });
4850 }
4851 }
4852 true
4853 }
4854 SupervisorCommand::Reload { reply } => {
4855 let result =
4856 reload_child(spec, runtime, registry, process_liveness, snapshot, child).await;
4857 let _ = reply.send(result);
4858 true
4859 }
4860 SupervisorCommand::SetEnabled { enabled, reply } => {
4861 let result = set_child_enabled(
4862 spec,
4863 runtime,
4864 registry,
4865 process_liveness,
4866 snapshot,
4867 child,
4868 enabled,
4869 )
4870 .await;
4871 let _ = reply.send(result);
4872 true
4873 }
4874 SupervisorCommand::UpdateConfiguration {
4875 spec: next_spec,
4876 health,
4877 drain_timeout_ms,
4878 reply,
4879 } => {
4880 if let Some(handle) = &runtime.supervisor_handle {
4881 handle.apply_identity_configuration(&next_spec);
4882 }
4883 *spec = next_spec;
4884 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4885 state.configuration_updated_since_spawn = true;
4886 });
4887 runtime.health = health;
4888 runtime.drain_timeout = drain_timeout_ms
4889 .map(Duration::from_millis)
4890 .unwrap_or(runtime.default_drain_timeout);
4891 *runtime
4892 .effective_drain_timeout
4893 .lock()
4894 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
4895 let _ = reply.send(());
4896 true
4897 }
4898 SupervisorCommand::Swap {
4899 ready_timeout,
4900 reply,
4901 } => {
4902 let end = swap::run_swap(
4903 spec,
4904 runtime,
4905 registry,
4906 process_liveness,
4907 snapshot,
4908 child,
4909 commands,
4910 ready_timeout.unwrap_or(DEFAULT_SWAP_READY_TIMEOUT),
4911 reply,
4912 )
4913 .await;
4914 requeued.extend(end.requeue);
4915 true
4916 }
4917 }
4918}
4919
4920async fn restart_child(
4921 spec: &ModuleSpec,
4922 runtime: &SupervisorRuntimeConfig,
4923 registry: &Registry,
4924 process_liveness: &SupervisorProcessLiveness,
4925 snapshot: &SharedSnapshot,
4926 child: &mut Option<SupervisedChild>,
4927 drain_timeout: Duration,
4928) -> Result<(), SuperviseError> {
4929 if !lock_snapshot(snapshot)?.enabled {
4931 return Err(SuperviseError::Disabled {
4932 module_id: spec.module_id.clone(),
4933 });
4934 }
4935 let stop_notice = begin_forwarding_drain_with_timeout(
4936 spec,
4937 runtime,
4938 registry,
4939 snapshot,
4940 None,
4941 RouteCloseReason::Restart,
4942 drain_timeout,
4943 )
4944 .await?;
4945
4946 if child.is_some() {
4947 drain_optional_child(
4948 &spec.module_id,
4949 spec.protocol,
4950 stop_notice,
4951 registry,
4952 snapshot,
4953 &runtime.terminal_ring,
4954 &runtime.spawn_events,
4955 child,
4956 drain_timeout,
4957 ModuleState::Restarting,
4958 Some(true),
4959 )
4960 .await?;
4961 } else {
4962 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4963 state.enabled = true;
4964 state.state = ModuleState::Restarting;
4965 clear_current_process_facts(state);
4966 })?;
4967 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4968 }
4969
4970 reset_restart_count(snapshot, &spec.module_id)?;
4971 sleep(runtime.restart_policy.backoff).await;
4972 if !respawn_still_pending(snapshot) {
4975 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4976 return Ok(());
4977 }
4978 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4979 match spawn_and_mark_running(spec, runtime, snapshot) {
4985 Ok(next_child) => {
4986 *child = Some(next_child);
4987 debug!(module_id = %spec.module_id, "supervised module restarted by operator request");
4988 Ok(())
4989 }
4990 Err(err) => {
4991 fail_snapshot(snapshot, Some(&spec.module_id), None);
4992 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4993 *child = None;
4994 Err(err)
4995 }
4996 }
4997}
4998
4999async fn reload_child(
5000 spec: &ModuleSpec,
5001 runtime: &SupervisorRuntimeConfig,
5002 registry: &Registry,
5003 process_liveness: &SupervisorProcessLiveness,
5004 snapshot: &SharedSnapshot,
5005 child: &mut Option<SupervisedChild>,
5006) -> Result<(), SuperviseError> {
5007 if !lock_snapshot(snapshot)?.enabled {
5009 return Err(SuperviseError::Disabled {
5010 module_id: spec.module_id.clone(),
5011 });
5012 }
5013 let stop_notice = begin_forwarding_drain(
5014 spec,
5015 runtime,
5016 registry,
5017 snapshot,
5018 Some(true),
5019 RouteCloseReason::Reload,
5020 )
5021 .await?;
5022
5023 if child.is_some() {
5024 drain_optional_child(
5025 &spec.module_id,
5026 spec.protocol,
5027 stop_notice,
5028 registry,
5029 snapshot,
5030 &runtime.terminal_ring,
5031 &runtime.spawn_events,
5032 child,
5033 runtime.drain_timeout,
5034 ModuleState::Restarting,
5035 Some(true),
5036 )
5037 .await?;
5038 } else {
5039 update_snapshot(snapshot, Some(&spec.module_id), |state| {
5040 state.enabled = true;
5041 state.state = ModuleState::Restarting;
5042 clear_current_process_facts(state);
5043 })?;
5044 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
5045 }
5046
5047 reset_restart_count(snapshot, &spec.module_id)?;
5048 sleep(runtime.restart_policy.backoff).await;
5049 if !respawn_still_pending(snapshot) {
5052 process_liveness.untrack_if_current(&spec.module_id, snapshot);
5053 return Ok(());
5054 }
5055 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
5056 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
5057 Ok(next_child) => next_child,
5058 Err(err) => {
5059 return handle_reload_spawn_failure(
5060 spec,
5061 runtime,
5062 process_liveness,
5063 snapshot,
5064 child,
5065 format!("new child failed to spawn: {err}"),
5066 )
5067 .await;
5068 }
5069 };
5070 *child = Some(next_child);
5071
5072 let wait_outcome = {
5073 let active_child = child.as_mut().expect("new reload child was just stored");
5074 wait_for_registration_after_reload(
5075 registry,
5076 &spec.module_id,
5077 snapshot,
5078 active_child,
5079 REGISTRY_RELEASE_TIMEOUT,
5080 )
5081 .await?
5082 };
5083
5084 match wait_outcome {
5085 RegistrationWaitOutcome::Registered => {
5086 debug!(module_id = %spec.module_id, "supervised module reloaded and registered");
5087 Ok(())
5088 }
5089 RegistrationWaitOutcome::Exited(exit_report) => {
5090 if let Some(active_child) = child.as_mut() {
5091 active_child.drain_stderr(&spec.module_id).await;
5092 }
5093 *child = None;
5094 handle_reload_child_registration_failure(
5095 spec,
5096 runtime,
5097 registry,
5098 process_liveness,
5099 snapshot,
5100 child,
5101 ReloadRegistrationFailure {
5102 exit_report: registration_failure_exit_report(exit_report),
5103 reason: "new child exited before registering".to_string(),
5104 },
5105 )
5106 .await
5107 }
5108 RegistrationWaitOutcome::TimedOut => {
5109 let mut timed_out_child = child
5110 .take()
5111 .expect("timed-out reload child is still running");
5112 timed_out_child
5113 .start_kill()
5114 .map_err(|source| SuperviseError::Kill {
5115 module_id: spec.module_id.clone(),
5116 source,
5117 })?;
5118 let status = timed_out_child
5119 .wait()
5120 .await
5121 .map_err(|source| SuperviseError::Wait {
5122 module_id: spec.module_id.clone(),
5123 source,
5124 })?;
5125 timed_out_child.drain_stderr(&spec.module_id).await;
5126 handle_reload_child_registration_failure(
5127 spec,
5128 runtime,
5129 registry,
5130 process_liveness,
5131 snapshot,
5132 child,
5133 ReloadRegistrationFailure {
5134 exit_report: registration_failure_exit_report(classify_reaped_child_exit(
5135 snapshot,
5136 &timed_out_child,
5137 &status,
5138 )),
5139 reason: format!(
5140 "new child did not register within {:?}",
5141 REGISTRY_RELEASE_TIMEOUT
5142 ),
5143 },
5144 )
5145 .await
5146 }
5147 }
5148}
5149
5150async fn set_child_enabled(
5151 spec: &ModuleSpec,
5152 runtime: &SupervisorRuntimeConfig,
5153 registry: &Registry,
5154 process_liveness: &SupervisorProcessLiveness,
5155 snapshot: &SharedSnapshot,
5156 child: &mut Option<SupervisedChild>,
5157 enabled: bool,
5158) -> Result<bool, SuperviseError> {
5159 let (current_enabled, current_state) = {
5160 let state = lock_snapshot(snapshot)?;
5161 (state.enabled, state.state)
5162 };
5163 let revive_terminal = enabled
5171 && current_enabled
5172 && child.is_none()
5173 && matches!(current_state, ModuleState::Failed | ModuleState::Stopped);
5174 if current_enabled == enabled && !revive_terminal {
5175 return Ok(false);
5176 }
5177
5178 if enabled {
5179 update_snapshot(snapshot, Some(&spec.module_id), |state| {
5180 state.enabled = true;
5181 state.state = ModuleState::Starting;
5182 clear_current_process_facts(state);
5183 })?;
5184 #[cfg(test)]
5185 if runtime.test_seed_stale_facts_before_enable_spawn {
5186 update_snapshot(snapshot, Some(&spec.module_id), |state| {
5187 state.process_alive = true;
5188 state.pid = Some(41);
5189 state.spawned_at_ms = Some(42);
5190 state.spawned_from = Some(PathBuf::from("/spawned/module"));
5191 state.spawned_file_identity = Some(SpawnedFileIdentity {
5192 device: 43,
5193 inode: 44,
5194 });
5195 })?;
5196 }
5197 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
5198 reset_restart_count(snapshot, &spec.module_id)?;
5199 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
5200 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
5201 Ok(next_child) => next_child,
5202 Err(err) => {
5203 if let Err(state_err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5204 state.state = ModuleState::Failed;
5205 clear_current_process_facts(state);
5206 }) {
5207 error!(module_id = %spec.module_id, error = %state_err, "failed to record enable spawn failure");
5208 }
5209 process_liveness.untrack_if_current(&spec.module_id, snapshot);
5210 return Err(err);
5211 }
5212 };
5213 *child = Some(next_child);
5214 debug!(module_id = %spec.module_id, "supervised module enabled");
5215 Ok(true)
5216 } else {
5217 let stop_notice = begin_forwarding_drain_if_configured(
5218 spec,
5219 runtime,
5220 registry,
5221 snapshot,
5222 Some(false),
5223 RouteCloseReason::Disable,
5224 )
5225 .await?;
5226 drain_optional_child(
5227 &spec.module_id,
5228 spec.protocol,
5229 stop_notice,
5230 registry,
5231 snapshot,
5232 &runtime.terminal_ring,
5233 &runtime.spawn_events,
5234 child,
5235 runtime.drain_timeout,
5236 ModuleState::Disabled,
5237 Some(false),
5238 )
5239 .await?;
5240 debug!(module_id = %spec.module_id, "supervised module disabled");
5241 Ok(true)
5242 }
5243}
5244
5245#[allow(clippy::too_many_arguments)]
5246async fn on_child_exit(
5247 spec: &ModuleSpec,
5248 policy: RestartPolicy,
5249 registry: &Registry,
5250 snapshot: &SharedSnapshot,
5251 terminal_ring: &Arc<Mutex<TerminalRing>>,
5252 spawn_events: &SpawnEventFeed,
5253 roster: &ChildRoster,
5254 exit_report: ExitReport,
5255) -> NextAction {
5256 if roster.is_closed() {
5262 return on_child_exit_during_daemon_shutdown(
5263 spec,
5264 registry,
5265 snapshot,
5266 terminal_ring,
5267 spawn_events,
5268 exit_report,
5269 )
5270 .await;
5271 }
5272 match exit_report.kind {
5273 ExitKind::Clean => {
5274 info!(
5275 module_id = %spec.module_id,
5276 exit_code = ?exit_report.code,
5277 exit_signal = ?exit_report.signal,
5278 "supervised module exited cleanly"
5279 );
5280 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5281 state.state = ModuleState::Stopped;
5282 clear_current_process_facts(state);
5283 state.last_exit = Some(exit_report.clone());
5284 }) {
5285 error!(module_id = %spec.module_id, error = %err, "failed to record clean module exit");
5286 }
5287 record_terminal(
5288 &spec.module_id,
5289 terminal_ring,
5290 spawn_events,
5291 &exit_report,
5292 TerminalDisposition::Stopped,
5293 );
5294 let registration_released = match wait_for_registration_release(
5295 registry,
5296 &spec.module_id,
5297 REGISTRY_RELEASE_TIMEOUT,
5298 )
5299 .await
5300 {
5301 Ok(()) => true,
5302 Err(err) => {
5303 warn!(module_id = %spec.module_id, error = %err, "registration still active after clean exit");
5304 false
5305 }
5306 };
5307 NextAction::Stop {
5308 registration_released,
5309 }
5310 }
5311 ExitKind::Crash => {
5312 warn!(
5313 module_id = %spec.module_id,
5314 exit_code = ?exit_report.code,
5315 exit_signal = ?exit_report.signal,
5316 "supervised module exited abnormally (crash)"
5317 );
5318 let mut restart_schedule = None;
5319 let mut disposition = TerminalDisposition::Disabled;
5320 let mut disposition_detail = None;
5324 let now = Instant::now();
5325 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5326 clear_current_process_facts(state);
5327 state.last_exit = Some(exit_report.clone());
5328 if state.enabled {
5329 if let Some(schedule) = state.next_crash_restart(&policy, now) {
5330 state.state = ModuleState::Restarting;
5331 restart_schedule = Some(schedule);
5332 disposition = TerminalDisposition::Restarting;
5333 } else {
5334 state.state = ModuleState::Failed;
5335 disposition = TerminalDisposition::Failed;
5336 disposition_detail = Some(policy.budget_exhausted_detail());
5337 }
5338 } else {
5339 state.state = ModuleState::Disabled;
5340 disposition = TerminalDisposition::Disabled;
5341 }
5342 }) {
5343 error!(module_id = %spec.module_id, error = %err, "failed to record crashed module exit");
5344 return NextAction::Stop {
5345 registration_released: false,
5346 };
5347 }
5348 if disposition_detail.is_some() {
5349 error!(
5354 module_id = %spec.module_id,
5355 max_restarts = policy.max_restarts,
5356 window_secs = policy.window.as_secs(),
5357 "module stopped: {}",
5358 policy.budget_exhausted_detail()
5359 );
5360 }
5361 record_terminal_with_detail(
5362 &spec.module_id,
5363 terminal_ring,
5364 spawn_events,
5365 &exit_report,
5366 disposition,
5367 disposition_detail,
5368 );
5369
5370 if let Some(schedule) = restart_schedule {
5371 NextAction::Restart {
5372 schedule: Some(schedule),
5373 }
5374 } else {
5375 let registration_released = match wait_for_registration_release(
5376 registry,
5377 &spec.module_id,
5378 REGISTRY_RELEASE_TIMEOUT,
5379 )
5380 .await
5381 {
5382 Ok(()) => true,
5383 Err(err) => {
5384 warn!(module_id = %spec.module_id, error = %err, "registration still active after failed module");
5385 false
5386 }
5387 };
5388 NextAction::Stop {
5389 registration_released,
5390 }
5391 }
5392 }
5393 ExitKind::DeliberateSeverance => {
5394 warn!(
5395 module_id = %spec.module_id,
5396 exit_code = ?exit_report.code,
5397 exit_signal = ?exit_report.signal,
5398 "supervised module exited after deliberate connection severance"
5399 );
5400 let mut should_restart = false;
5401 let mut disposition = TerminalDisposition::Disabled;
5402 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5403 clear_current_process_facts(state);
5404 state.last_exit = Some(exit_report.clone());
5405 state.lifetime_restarts += 1;
5406 if state.enabled {
5407 state.state = ModuleState::Restarting;
5408 should_restart = true;
5409 disposition = TerminalDisposition::Restarting;
5410 } else {
5411 state.state = ModuleState::Disabled;
5412 }
5413 }) {
5414 error!(module_id = %spec.module_id, error = %err, "failed to record deliberately severed module exit");
5415 return NextAction::Stop {
5416 registration_released: false,
5417 };
5418 }
5419 record_terminal(
5420 &spec.module_id,
5421 terminal_ring,
5422 spawn_events,
5423 &exit_report,
5424 disposition,
5425 );
5426
5427 if should_restart {
5428 NextAction::Restart { schedule: None }
5429 } else {
5430 let registration_released = match wait_for_registration_release(
5431 registry,
5432 &spec.module_id,
5433 REGISTRY_RELEASE_TIMEOUT,
5434 )
5435 .await
5436 {
5437 Ok(()) => true,
5438 Err(err) => {
5439 warn!(module_id = %spec.module_id, error = %err, "registration still active after deliberately severed module exit");
5440 false
5441 }
5442 };
5443 NextAction::Stop {
5444 registration_released,
5445 }
5446 }
5447 }
5448 }
5449}
5450
5451async fn on_child_exit_during_daemon_shutdown(
5452 spec: &ModuleSpec,
5453 registry: &Registry,
5454 snapshot: &SharedSnapshot,
5455 terminal_ring: &Arc<Mutex<TerminalRing>>,
5456 spawn_events: &SpawnEventFeed,
5457 exit_report: ExitReport,
5458) -> NextAction {
5459 info!(
5460 module_id = %spec.module_id,
5461 exit_code = ?exit_report.code,
5462 exit_signal = ?exit_report.signal,
5463 exit_kind = ?exit_report.kind,
5464 "supervised module exited during daemon shutdown; not restarting it"
5465 );
5466 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5467 state.state = ModuleState::Stopped;
5468 clear_current_process_facts(state);
5469 state.last_exit = Some(exit_report.clone());
5470 }) {
5471 error!(module_id = %spec.module_id, error = %err, "failed to record module exit during daemon shutdown");
5472 }
5473 record_terminal(
5474 &spec.module_id,
5475 terminal_ring,
5476 spawn_events,
5477 &exit_report,
5478 TerminalDisposition::DaemonShutdown,
5479 );
5480 let registration_released =
5481 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT)
5482 .await
5483 .is_ok();
5484 NextAction::Stop {
5485 registration_released,
5486 }
5487}
5488
5489fn record_wait_error_terminal(
5490 module_id: &str,
5491 terminal_ring: &Arc<Mutex<TerminalRing>>,
5492 spawn_events: &SpawnEventFeed,
5493) {
5494 record_terminal(
5495 module_id,
5496 terminal_ring,
5497 spawn_events,
5498 &wait_error_exit_report(),
5499 TerminalDisposition::Failed,
5500 );
5501}
5502
5503fn record_terminal(
5504 module_id: &str,
5505 terminal_ring: &Arc<Mutex<TerminalRing>>,
5506 spawn_events: &SpawnEventFeed,
5507 exit_report: &ExitReport,
5508 disposition: TerminalDisposition,
5509) {
5510 record_terminal_with_detail(
5511 module_id,
5512 terminal_ring,
5513 spawn_events,
5514 exit_report,
5515 disposition,
5516 None,
5517 );
5518}
5519
5520fn durable_terminal_history_of(
5524 terminal_ring: &Mutex<TerminalRing>,
5525 module_id: &str,
5526) -> subc_control::TerminalHistory {
5527 let read = terminal_ring
5528 .lock()
5529 .unwrap_or_else(|p| p.into_inner())
5530 .capture_durable_history();
5531 read.read(module_id)
5532}
5533
5534fn record_terminal_with_detail(
5535 module_id: &str,
5536 terminal_ring: &Arc<Mutex<TerminalRing>>,
5537 spawn_events: &SpawnEventFeed,
5538 exit_report: &ExitReport,
5539 disposition: TerminalDisposition,
5540 disposition_detail: Option<String>,
5541) {
5542 spawn_events.emit_exited(module_id, exit_report.code, exit_report.signal);
5543 let record = TerminalRecord {
5544 exit_code: exit_report.code,
5545 exit_signal: exit_report.signal,
5546 at_ms: exit_report.at_ms,
5547 disposition,
5548 exit_kind: exit_report.kind.into(),
5549 disposition_detail,
5550 };
5551 terminal_ring
5552 .lock()
5553 .unwrap_or_else(|poisoned| poisoned.into_inner())
5554 .record_exit(module_id, record);
5555}
5556
5557fn untrack_if_registration_released(
5558 process_liveness: &SupervisorProcessLiveness,
5559 registry: &Registry,
5560 module_id: &str,
5561 snapshot: &SharedSnapshot,
5562) {
5563 match registry.get_module(module_id) {
5564 Ok(None) => process_liveness.untrack_if_current(module_id, snapshot),
5565 Ok(Some(_)) => {}
5566 Err(err) => {
5567 warn!(module_id, error = %err, "could not determine whether supervisor liveness can be untracked");
5568 }
5569 }
5570}
5571
5572#[cfg(test)]
5586fn apply_wire_spawn_args(
5587 command: &mut Command,
5588 spec: &ModuleSpec,
5589 connection_file_path: Option<&std::path::Path>,
5590 handle: Option<&SupervisorHandle>,
5591) -> Result<(), SuperviseError> {
5592 apply_wire_spawn_args_for_role(
5593 command,
5594 spec,
5595 connection_file_path,
5596 handle,
5597 SpawnRole::Plain,
5598 )
5599}
5600
5601fn apply_wire_spawn_args_for_role(
5610 command: &mut Command,
5611 spec: &ModuleSpec,
5612 connection_file_path: Option<&std::path::Path>,
5613 handle: Option<&SupervisorHandle>,
5614 role: SpawnRole,
5615) -> Result<(), SuperviseError> {
5616 command.env(SUBC_MODULE_ID_ENV, &spec.module_id);
5617 if spec.protocol == ModuleProtocol::None {
5618 return Ok(());
5619 }
5620 if let Some(connection_file_path) = connection_file_path {
5621 command.arg(SUBC_ARG).arg(connection_file_path);
5622 }
5623
5624 let nonce = generate_launch_nonce()?;
5628 if let Some(handle) = handle {
5629 match role {
5630 SpawnRole::Plain => {
5631 handle.set_spawn_nonce(&spec.module_id, nonce.clone());
5632 if spec.reserved {
5633 handle.set_reserved_nonce(&spec.module_id, nonce.clone());
5634 }
5635 }
5636 SpawnRole::SwapCandidate => handle.open_swap(&spec.module_id, nonce.clone()),
5637 }
5638 }
5639 command.env(SUBC_LAUNCH_NONCE_ENV, nonce);
5640 Ok(())
5641}
5642
5643fn apply_child_env(command: &mut Command, spec: &ModuleSpec) {
5644 command.env_remove(CK_LOG_ENV);
5645 command.env_remove(SUBC_SPAWN_ROLE_ENV);
5652 for (key, value) in &spec.env {
5653 if matches!(
5657 key.as_str(),
5658 CAPTURE_MAX_FILE_MB_ENV | CAPTURE_KEEP_ENV | CAPTURE_MAX_AGE_DAYS_ENV
5659 ) || key == SUBC_SPAWN_ROLE_ENV
5660 {
5661 continue;
5662 }
5663 command.env(key, value);
5664 }
5665}
5666
5667#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5670enum SpawnRole {
5671 Plain,
5672 SwapCandidate,
5673}
5674
5675fn apply_spawn_role(command: &mut Command, role: SpawnRole) {
5678 if role == SpawnRole::SwapCandidate {
5679 command.env(SUBC_SPAWN_ROLE_ENV, SPAWN_ROLE_SWAP_CANDIDATE);
5680 }
5681}
5682
5683fn spawn_child(
5684 spec: &ModuleSpec,
5685 connection_file_path: Option<&std::path::Path>,
5686 handle: Option<&SupervisorHandle>,
5687 ring: &Arc<Mutex<StderrRing>>,
5688 capture_logs_dir: Option<&std::path::Path>,
5689 roster: &ChildRoster,
5690 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5691) -> Result<SupervisedChild, SuperviseError> {
5692 spawn_child_in_slot(
5693 spec,
5694 connection_file_path,
5695 handle,
5696 ring,
5697 capture_logs_dir,
5698 roster,
5699 #[cfg(target_os = "linux")]
5700 cgroup_placement,
5701 SpawnRole::Plain,
5702 false,
5703 )
5704}
5705
5706#[allow(clippy::too_many_arguments)]
5719fn spawn_child_in_slot(
5720 spec: &ModuleSpec,
5721 connection_file_path: Option<&std::path::Path>,
5722 handle: Option<&SupervisorHandle>,
5723 ring: &Arc<Mutex<StderrRing>>,
5724 capture_logs_dir: Option<&std::path::Path>,
5725 roster: &ChildRoster,
5726 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5727 role: SpawnRole,
5728 alternate_slot: bool,
5729) -> Result<SupervisedChild, SuperviseError> {
5730 if roster.is_closed() {
5731 return Err(SuperviseError::Spawn {
5732 program: spec.program.clone(),
5733 source: io::Error::other("the daemon is shutting down; not starting a new process"),
5734 cgroup_path: None,
5735 });
5736 }
5737 #[cfg(target_os = "linux")]
5738 let cgroup_name = swap::cgroup_name(&spec.module_id, alternate_slot);
5739 #[cfg(not(target_os = "linux"))]
5740 let _ = alternate_slot;
5741 let mut command = Command::new(&spec.program);
5742 command.args(&spec.args);
5743 apply_child_env(&mut command, spec);
5773 apply_spawn_role(&mut command, role);
5774 apply_wire_spawn_args_for_role(&mut command, spec, connection_file_path, handle, role)?;
5775
5776 #[cfg(target_os = "linux")]
5777 let cgroup_path = cgroup_placement
5778 .map(|placement| placement.module_path(&cgroup_name))
5779 .transpose()
5780 .map_err(|source| SuperviseError::Cgroup {
5781 module_id: spec.module_id.clone(),
5782 source,
5783 })?;
5784 #[cfg(not(target_os = "linux"))]
5785 let cgroup_path: Option<PathBuf> = None;
5786 #[cfg(target_os = "linux")]
5787 if let Some(path) = &cgroup_path {
5788 if let Err(error) = apply_cgroup_placement(&mut command, spec, path) {
5789 if let Some(placement) = cgroup_placement {
5790 remove_module_cgroup(placement, &cgroup_name);
5791 }
5792 return Err(error);
5793 }
5794 }
5795
5796 let output_sink = if let Some(logs_dir) = capture_logs_dir {
5797 let path = logs_dir.join(format!("{}.stderr.log", spec.module_id));
5798 match ChildOutputSink::open(&path, capture_retention(spec)) {
5799 Ok(sink) => sink,
5800 Err(error) => {
5801 warn!(
5802 module_id = %spec.module_id,
5803 path = %path.display(),
5804 error = %error,
5805 "could not open child output capture file; forwarding to stderr"
5806 );
5807 ChildOutputSink::Stderr
5808 }
5809 }
5810 } else {
5811 ChildOutputSink::Stderr
5812 };
5813
5814 command.stdout(Stdio::piped());
5815 command.stderr(Stdio::piped());
5816 command.kill_on_drop(true);
5817 #[cfg(unix)]
5834 command.process_group(0);
5835 command.stdin(Stdio::null());
5836
5837 #[cfg(windows)]
5842 subc_jobobject::suspend_on_create_async(&mut command);
5843 let mut child = match command.spawn() {
5844 Ok(child) => child,
5845 Err(source) => {
5846 #[cfg(target_os = "linux")]
5847 if let Some(placement) = cgroup_placement {
5848 remove_module_cgroup(placement, &cgroup_name);
5849 }
5850 return Err(SuperviseError::Spawn {
5851 program: spec.program.clone(),
5852 source,
5853 cgroup_path,
5854 });
5855 }
5856 };
5857
5858 #[cfg(windows)]
5860 let job = contain_spawned_child(&child, spec)?;
5861 let spawned_at_ms = unix_ms_now();
5862 let spawned_from = spec.program.clone();
5863 let spawned_file_identity = spawned_file_identity(&spawned_from);
5864 let pid = child.id().ok_or_else(|| SuperviseError::Spawn {
5865 program: spec.program.clone(),
5866 source: io::Error::other("spawned child exposed no live pid"),
5867 cgroup_path: cgroup_path.clone(),
5868 })?;
5869 let process_start_time = crate::provenance::process_start_time(pid);
5870 let process_identity = process_start_time.map(|start_time| ProcessIdentity { pid, start_time });
5871 #[cfg(target_os = "linux")]
5875 let recorded_cgroup_name = cgroup_path.as_ref().map(|_| cgroup_name.clone());
5876 #[cfg(not(target_os = "linux"))]
5877 let recorded_cgroup_name = None;
5878 let roster_guard = roster.admit(
5879 spec.module_id.clone(),
5880 pid,
5881 spec.protocol,
5882 process_start_time,
5883 crate::child_roster::RecordedIdentity {
5884 start_time: subc_os::start_time(pid),
5885 executable: spawned_file_identity.map(|identity| {
5886 crate::live_children::ExecutableIdentity {
5887 device: identity.device,
5888 inode: identity.inode,
5889 }
5890 }),
5891 cgroup_name: recorded_cgroup_name,
5892 },
5893 );
5894 if roster.is_closed() {
5903 if let Err(error) = child.start_kill() {
5904 debug!(module_id = %spec.module_id, pid, %error, "kill of a process spawned during daemon shutdown failed; it may already have exited");
5905 }
5906 drop(roster_guard);
5907 return Err(SuperviseError::Spawn {
5908 program: spec.program.clone(),
5909 source: io::Error::other(
5910 "the daemon began shutting down while this process was starting; ended it",
5911 ),
5912 cgroup_path,
5913 });
5914 }
5915
5916 let stdout_pump = match child.stdout.take() {
5917 Some(stdout) => Some(tokio::spawn(pump_stdout_to(stdout, output_sink.clone()))),
5918 None => {
5919 warn!(
5920 module_id = %spec.module_id,
5921 "spawned child exposed no stdout pipe; file capture will be incomplete"
5922 );
5923 None
5924 }
5925 };
5926 let stderr_pump = match child.stderr.take() {
5927 Some(stderr) => {
5928 let generation = ring
5929 .lock()
5930 .unwrap_or_else(|poisoned| poisoned.into_inner())
5931 .begin_process();
5932 Some(StderrPump {
5933 task: tokio::spawn(pump_stderr_to(
5934 stderr,
5935 Arc::clone(ring),
5936 generation,
5937 output_sink,
5938 )),
5939 generation,
5940 })
5941 }
5942 None => {
5943 ring.lock()
5947 .unwrap_or_else(|poisoned| poisoned.into_inner())
5948 .mark_not_captured("stderr pipe was not available on spawn");
5949 warn!(
5950 module_id = %spec.module_id,
5951 "spawned child exposed no stderr pipe; tail will be unavailable"
5952 );
5953 None
5954 }
5955 };
5956
5957 Ok(SupervisedChild {
5958 child,
5959 #[cfg(target_os = "linux")]
5960 module_id: cgroup_name,
5961 #[cfg(target_os = "linux")]
5962 cgroup_placement: cgroup_placement.cloned(),
5963 #[cfg(windows)]
5964 job,
5965 stdout_pump,
5966 stderr_pump,
5967 stderr_ring: Arc::clone(ring),
5968 spawned_at_ms,
5969 spawned_from,
5970 spawned_file_identity,
5971 process_start_time,
5972 process_identity,
5973 pid,
5974 roster_guard: Some(roster_guard),
5975 })
5976}
5977
5978#[cfg(windows)]
5992fn contain_spawned_child(
5993 child: &Child,
5994 spec: &ModuleSpec,
5995) -> Result<Option<subc_jobobject::JobObject>, SuperviseError> {
5996 let module_id = spec.module_id.as_str();
5997 let Some(pid) = child.id() else {
5998 warn!(
6001 module_id,
6002 "spawned child had already exited before containment; no job object attached"
6003 );
6004 return Ok(None);
6005 };
6006
6007 let job = match subc_jobobject::JobObject::new() {
6008 Ok(job) => job,
6009 Err(source) => {
6010 warn!(
6011 module_id,
6012 error = %source,
6013 "could not create a job object; this module's helper processes will not be \
6014 reaped on teardown"
6015 );
6016 resume_suspended_child(pid, spec)?;
6019 return Ok(None);
6020 }
6021 };
6022
6023 if let Err(source) = job.assign(child) {
6024 warn!(
6025 module_id,
6026 error = %source,
6027 "could not assign the child to its job object; this module's helper processes \
6028 will not be reaped on teardown"
6029 );
6030 resume_suspended_child(pid, spec)?;
6031 return Ok(None);
6032 }
6033
6034 resume_suspended_child(pid, spec)?;
6035 Ok(Some(job))
6036}
6037
6038#[cfg(windows)]
6043fn resume_suspended_child(pid: u32, spec: &ModuleSpec) -> Result<(), SuperviseError> {
6044 if let Err(source) = subc_jobobject::resume_main_thread(pid) {
6045 let _ = std::process::Command::new("taskkill.exe")
6049 .args(["/PID", &pid.to_string(), "/T", "/F"])
6050 .stdin(Stdio::null())
6051 .stdout(Stdio::null())
6052 .stderr(Stdio::null())
6053 .status();
6054 return Err(SuperviseError::Spawn {
6055 program: spec.program.clone(),
6056 source,
6057 cgroup_path: None,
6058 });
6059 }
6060 Ok(())
6061}
6062
6063#[cfg(target_os = "linux")]
6064fn remove_module_cgroup(placement: &subc_cgroup::Placement, module_id: &str) {
6065 match placement.remove_module(module_id) {
6066 Ok(()) => debug!(module_id, "removed module cgroup after process exit"),
6067 Err(error) => warn!(
6068 module_id,
6069 error = %error,
6070 "could not remove module cgroup after process exit; continuing teardown"
6071 ),
6072 }
6073}
6074
6075#[cfg(target_os = "linux")]
6076fn apply_cgroup_placement(
6077 command: &mut Command,
6078 spec: &ModuleSpec,
6079 path: &std::path::Path,
6080) -> Result<(), SuperviseError> {
6081 subc_cgroup::apply(command, path).map_err(|source| SuperviseError::Cgroup {
6082 module_id: spec.module_id.clone(),
6083 source,
6084 })
6085}
6086
6087fn capture_retention(spec: &ModuleSpec) -> Retention {
6088 let defaults = Retention::default();
6089 let value = |name: &str| {
6090 spec.env
6091 .iter()
6092 .rev()
6093 .find_map(|(key, value)| (key == name).then_some(value.as_str()))
6094 };
6095 Retention {
6096 max_file_mb: value(CAPTURE_MAX_FILE_MB_ENV)
6097 .and_then(|value| value.parse().ok())
6098 .unwrap_or(defaults.max_file_mb),
6099 keep: value(CAPTURE_KEEP_ENV)
6100 .and_then(|value| value.parse().ok())
6101 .unwrap_or(defaults.keep),
6102 max_age_days: value(CAPTURE_MAX_AGE_DAYS_ENV)
6103 .and_then(|value| value.parse().ok())
6104 .unwrap_or(defaults.max_age_days),
6105 }
6106}
6107
6108fn generate_launch_nonce() -> Result<String, SuperviseError> {
6111 let mut bytes = [0u8; 32];
6112 getrandom::getrandom(&mut bytes).map_err(|source| SuperviseError::LaunchNonce {
6113 reason: source.to_string(),
6114 })?;
6115 let mut hex = String::with_capacity(64);
6116 for b in bytes {
6117 use std::fmt::Write;
6118 let _ = write!(hex, "{b:02x}");
6119 }
6120 Ok(hex)
6121}
6122
6123fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
6126 if a.len() != b.len() {
6127 return false;
6128 }
6129 let mut diff = 0u8;
6130 for (x, y) in a.iter().zip(b.iter()) {
6131 diff |= x ^ y;
6132 }
6133 diff == 0
6134}
6135
6136fn spawn_and_mark_running(
6137 spec: &ModuleSpec,
6138 runtime: &SupervisorRuntimeConfig,
6139 snapshot: &SharedSnapshot,
6140) -> Result<SupervisedChild, SuperviseError> {
6141 let child = spawn_child(
6142 spec,
6143 runtime.connection_file_path.as_deref(),
6144 runtime.supervisor_handle.as_ref(),
6145 &runtime.stderr_ring,
6146 runtime.capture_logs_dir.as_deref(),
6147 &runtime.child_roster,
6148 #[cfg(target_os = "linux")]
6149 runtime.cgroup_placement.as_ref(),
6150 )?;
6151 set_running(snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
6152 Ok(child)
6153}
6154
6155enum RegistrationWaitOutcome {
6156 Registered,
6157 Exited(ExitReport),
6158 TimedOut,
6159}
6160
6161struct ReloadRegistrationFailure {
6162 exit_report: ExitReport,
6163 reason: String,
6164}
6165
6166#[derive(Debug, Clone, Copy, PartialEq, Eq)]
6167enum BusyGaugeObservation {
6168 Quiescent,
6169 Busy,
6170 Omitted,
6171}
6172
6173fn busy_gauge_observation(metrics: Option<&Value>, gauges: &[String]) -> BusyGaugeObservation {
6174 let Some(metrics) = metrics.and_then(Value::as_object) else {
6175 return BusyGaugeObservation::Omitted;
6176 };
6177 let mut sum = 0u128;
6178 for gauge in gauges {
6179 let Some(value) = metrics.get(gauge) else {
6180 return BusyGaugeObservation::Omitted;
6181 };
6182 let Some(value) = value.as_u64() else {
6183 return BusyGaugeObservation::Busy;
6184 };
6185 sum = sum.saturating_add(u128::from(value));
6186 }
6187 if sum == 0 {
6188 BusyGaugeObservation::Quiescent
6189 } else {
6190 BusyGaugeObservation::Busy
6191 }
6192}
6193
6194fn declared_busy_gauges(
6195 registry: &Registry,
6196 module_id: &str,
6197) -> Result<Vec<String>, SuperviseError> {
6198 busy_gauges_of(
6199 registry
6200 .get_module(module_id)
6201 .map_err(SuperviseError::Registry)?,
6202 )
6203}
6204
6205fn declared_busy_gauges_for_connection(
6209 registry: &Registry,
6210 connection_id: ConnectionId,
6211) -> Result<Vec<String>, SuperviseError> {
6212 busy_gauges_of(
6213 registry
6214 .get_module_by_connection(connection_id)
6215 .map_err(SuperviseError::Registry)?,
6216 )
6217}
6218
6219fn busy_gauges_of(
6220 registration: Option<crate::registry::ModuleRegistration>,
6221) -> Result<Vec<String>, SuperviseError> {
6222 let Some(registration) = registration else {
6223 return Ok(Vec::new());
6224 };
6225 let Some(self_signals) = registration.manifest.self_signals else {
6226 return Ok(Vec::new());
6227 };
6228
6229 let mut gauges = Vec::new();
6230 for declaration in self_signals {
6231 if declaration.kind != SelfSignalKind::Busy {
6232 continue;
6233 }
6234 match declaration.anchored_to {
6235 SignalAnchor::HealthGauges { gauges: declared } if !declared.is_empty() => {
6236 gauges.extend(declared)
6237 }
6238 _ => {
6239 gauges.push(String::new());
6242 }
6243 }
6244 }
6245 Ok(gauges)
6246}
6247
6248async fn wait_for_forwarding_quiescence(
6253 forwarding: &ForwardingTable,
6254 module_id: &str,
6255 runtime: &SupervisorRuntimeConfig,
6256 endpoint: crate::ModuleEndpointId,
6257 deadline: Instant,
6258 busy_gauges: &[String],
6259 scope: DrainScope,
6260) -> Result<bool, SuperviseError> {
6261 let mut gauges_quiescent = busy_gauges.is_empty();
6262 let mut next_probe_at = Instant::now();
6263 let mut omission_counted = false;
6264
6265 loop {
6266 let now = Instant::now();
6267 if !busy_gauges.is_empty() && now >= next_probe_at && now < deadline {
6268 let report = match scope {
6269 DrainScope::Active => probe_module_health(module_id, runtime, Some(deadline)).await,
6270 DrainScope::Endpoint(endpoint) => {
6271 probe_endpoint_health(endpoint, runtime, Some(deadline)).await
6272 }
6273 };
6274 gauges_quiescent = match report {
6275 Ok(report) => match busy_gauge_observation(report.metrics.as_ref(), busy_gauges) {
6276 BusyGaugeObservation::Quiescent => true,
6277 BusyGaugeObservation::Busy => false,
6278 BusyGaugeObservation::Omitted => {
6279 if !omission_counted {
6280 forwarding
6281 .counters()
6282 .increment_drains_with_undeclared_gauge();
6283 omission_counted = true;
6284 }
6285 false
6286 }
6287 },
6288 Err(err) => {
6289 warn!(
6290 module_id,
6291 error = %err,
6292 "drain health.check did not produce declared busy gauges; treating module as busy"
6293 );
6294 false
6295 }
6296 };
6297 next_probe_at = Instant::now() + runtime.health.cadence.max(REGISTRY_RELEASE_POLL);
6298 }
6299
6300 let in_flight = forwarding
6301 .endpoint_in_flight_count(endpoint)
6302 .map_err(SuperviseError::Forwarding)?;
6303 if in_flight == 0 && gauges_quiescent {
6304 return Ok(true);
6305 }
6306
6307 let now = Instant::now();
6308 if now >= deadline {
6309 return Ok(false);
6310 }
6311 let mut wait = deadline
6312 .saturating_duration_since(now)
6313 .min(REGISTRY_RELEASE_POLL);
6314 if !busy_gauges.is_empty() {
6315 wait = wait.min(next_probe_at.saturating_duration_since(now));
6316 }
6317 sleep(wait).await;
6318 }
6319}
6320
6321fn drained_after_quiescence_wait(wait_result: &Result<bool, SuperviseError>) -> bool {
6329 match wait_result {
6330 Ok(drained) => *drained,
6331 Err(_) => false,
6332 }
6333}
6334
6335fn send_route_goodbyes(forwarding: &ForwardingTable, released_routes: Vec<GoodbyeTarget>) {
6336 for released in released_routes {
6337 let frame = match Frame::build_with_version(
6338 released.negotiated_ver,
6339 FrameType::Goodbye,
6340 control_flags(),
6341 released.channel,
6342 released.epoch,
6343 0,
6344 Vec::new(),
6345 ) {
6346 Ok(frame) => frame,
6347 Err(err) => {
6348 warn!(
6349 route_channel = released.channel,
6350 error = %err,
6351 "failed to build supervisor drain route GOODBYE frame"
6352 );
6353 continue;
6354 }
6355 };
6356 if !released.close_on_delivery_failure() {
6357 crate::forwarding::send_module_route_goodbye(
6358 &forwarding.counters(),
6359 &released.sink,
6360 frame,
6361 released.module_id.as_deref(),
6362 "supervisor drain",
6363 );
6364 continue;
6365 }
6366 if let Err(err) = released.sink.try_send(frame) {
6367 warn!(
6368 target_connection_id = released.connection_id.get(),
6369 route_channel = released.channel,
6370 error = %err,
6371 "supervisor drain route GOODBYE was not delivered to client; closing target connection"
6372 );
6373 let _ = forwarding.escalate_client_delivery_failure(
6374 released.connection_id,
6375 released.channel,
6376 released.epoch,
6377 CloseReason::new(
6378 "route_goodbye_delivery_failed",
6379 format!(
6380 "failed to enqueue supervisor drain route GOODBYE for channel {}: {err}",
6381 released.channel
6382 ),
6383 ),
6384 crate::forwarding::UndeliveredFrame {
6385 module_id: released.module_id.as_deref(),
6386 sink: &released.sink,
6387 },
6388 );
6389 }
6390 }
6391}
6392
6393fn send_module_draining(
6394 module_id: &str,
6395 reason: RouteCloseReason,
6396 deadline_ms: u64,
6397 target: &ModuleDrainTarget,
6398) {
6399 let body = match serde_json::to_vec(&ModuleControlCommand::Draining {
6400 reason,
6401 deadline_ms,
6402 }) {
6403 Ok(body) => body,
6404 Err(err) => {
6405 warn!(
6406 module_id,
6407 error = %err,
6408 "failed to encode module draining command"
6409 );
6410 return;
6411 }
6412 };
6413 let frame = match Frame::build_with_version(
6414 target.negotiated_ver,
6415 FrameType::Push,
6416 control_flags(),
6417 0,
6418 0,
6419 0,
6420 body,
6421 ) {
6422 Ok(frame) => frame,
6423 Err(err) => {
6424 warn!(
6425 module_id,
6426 error = %err,
6427 "failed to build module draining command frame"
6428 );
6429 return;
6430 }
6431 };
6432 if let Err(err) = target.sink.try_send(frame) {
6433 warn!(
6434 module_id,
6435 target_connection_id = target.endpoint.connection_id.get(),
6436 error = %err,
6437 "module draining command was not delivered to peer"
6438 );
6439 }
6440}
6441
6442fn send_module_goodbye(module_id: &str, forwarding: &ForwardingTable, target: &ModuleDrainTarget) {
6443 let frame = match Frame::build_with_version(
6444 target.negotiated_ver,
6445 FrameType::Goodbye,
6446 control_flags(),
6447 0,
6448 0,
6449 0,
6450 Vec::new(),
6451 ) {
6452 Ok(frame) => frame,
6453 Err(err) => {
6454 warn!(
6455 module_id,
6456 error = %err,
6457 "failed to build supervisor drain module GOODBYE frame"
6458 );
6459 return;
6460 }
6461 };
6462 if let Err(err) = target.sink.try_send(frame) {
6463 warn!(
6464 module_id,
6465 target_connection_id = target.endpoint.connection_id.get(),
6466 error = %err,
6467 "supervisor drain module GOODBYE was not delivered to peer; closing module connection"
6468 );
6469 forwarding.request_connection_close(
6470 target.endpoint.connection_id,
6471 CloseReason::new(
6472 "module_goodbye_delivery_failed",
6473 format!("failed to enqueue supervisor drain module GOODBYE for module '{module_id}': {err}"),
6474 ),
6475 );
6476 }
6477}
6478
6479#[derive(Clone, Copy)]
6480struct ForwardingDrainContext<'a> {
6481 spec: &'a ModuleSpec,
6482 runtime: &'a SupervisorRuntimeConfig,
6483 registry: &'a Registry,
6484 scope: DrainScope,
6485}
6486
6487#[derive(Debug, Clone, Copy, PartialEq, Eq)]
6489enum DrainScope {
6490 Active,
6493 Endpoint(crate::ModuleEndpointId),
6498}
6499
6500#[derive(Debug, Clone, Copy, PartialEq, Eq)]
6508enum StopNotice {
6509 SentOverConnection,
6512 NoConnection,
6516 NotSent,
6520}
6521
6522async fn begin_forwarding_drain(
6523 spec: &ModuleSpec,
6524 runtime: &SupervisorRuntimeConfig,
6525 registry: &Registry,
6526 snapshot: &SharedSnapshot,
6527 enabled: Option<bool>,
6528 reason: RouteCloseReason,
6529) -> Result<StopNotice, SuperviseError> {
6530 let Some(forwarding) = runtime.forwarding.as_ref() else {
6531 return Err(SuperviseError::ReloadUnavailable {
6532 module_id: spec.module_id.clone(),
6533 reason: "supervisor was not configured with a forwarding table".to_string(),
6534 });
6535 };
6536
6537 begin_forwarding_drain_with(
6538 forwarding,
6539 ForwardingDrainContext {
6540 spec,
6541 runtime,
6542 registry,
6543 scope: DrainScope::Active,
6544 },
6545 snapshot,
6546 enabled,
6547 reason,
6548 runtime.drain_timeout,
6549 )
6550 .await
6551}
6552
6553async fn begin_forwarding_drain_if_configured(
6554 spec: &ModuleSpec,
6555 runtime: &SupervisorRuntimeConfig,
6556 registry: &Registry,
6557 snapshot: &SharedSnapshot,
6558 enabled: Option<bool>,
6559 reason: RouteCloseReason,
6560) -> Result<StopNotice, SuperviseError> {
6561 begin_forwarding_drain_with_timeout(
6562 spec,
6563 runtime,
6564 registry,
6565 snapshot,
6566 enabled,
6567 reason,
6568 runtime.drain_timeout,
6569 )
6570 .await
6571}
6572
6573async fn begin_forwarding_drain_with_timeout(
6577 spec: &ModuleSpec,
6578 runtime: &SupervisorRuntimeConfig,
6579 registry: &Registry,
6580 snapshot: &SharedSnapshot,
6581 enabled: Option<bool>,
6582 reason: RouteCloseReason,
6583 drain_timeout: Duration,
6584) -> Result<StopNotice, SuperviseError> {
6585 let Some(forwarding) = runtime.forwarding.as_ref() else {
6586 return Ok(StopNotice::NotSent);
6587 };
6588
6589 begin_forwarding_drain_with(
6590 forwarding,
6591 ForwardingDrainContext {
6592 spec,
6593 runtime,
6594 registry,
6595 scope: DrainScope::Active,
6596 },
6597 snapshot,
6598 enabled,
6599 reason,
6600 drain_timeout,
6601 )
6602 .await
6603}
6604
6605async fn begin_forwarding_drain_with(
6606 forwarding: &ForwardingTable,
6607 context: ForwardingDrainContext<'_>,
6608 snapshot: &SharedSnapshot,
6609 enabled: Option<bool>,
6610 reason: RouteCloseReason,
6611 drain_timeout: Duration,
6612) -> Result<StopNotice, SuperviseError> {
6613 let ForwardingDrainContext {
6614 spec,
6615 runtime,
6616 registry,
6617 scope,
6618 } = context;
6619 debug_assert_ne!(reason, RouteCloseReason::Crash);
6620 let terminal = matches!(reason, RouteCloseReason::Disable);
6621 let drain_started_at = Instant::now();
6622 let drain_deadline = drain_started_at + drain_timeout;
6623 let deadline_ms =
6624 unix_ms_now().saturating_add(u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX));
6625 let busy_gauges = match scope {
6626 DrainScope::Active => declared_busy_gauges(registry, &spec.module_id)?,
6627 DrainScope::Endpoint(endpoint) => {
6628 declared_busy_gauges_for_connection(registry, endpoint.connection_id)?
6629 }
6630 };
6631
6632 let gate_started = Instant::now();
6635 let drain_target = match scope {
6636 DrainScope::Active => forwarding.begin_module_drain(&spec.module_id, reason),
6637 DrainScope::Endpoint(endpoint) => forwarding.begin_endpoint_drain(endpoint, reason),
6638 }
6639 .map_err(SuperviseError::Forwarding)?;
6640 info!(
6645 module_id = %spec.module_id,
6646 ?reason,
6647 gate_ms = u64::try_from(gate_started.elapsed().as_millis()).unwrap_or(u64::MAX),
6648 connected = drain_target.is_some(),
6649 "module drain began; route admission closed"
6650 );
6651 if scope == DrainScope::Active {
6652 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6653 state.state = ModuleState::Draining;
6654 state.draining_to_replace =
6655 matches!(reason, RouteCloseReason::Restart | RouteCloseReason::Reload);
6656 if let Some(enabled) = enabled {
6657 state.enabled = enabled;
6658 }
6659 })?;
6660 }
6661
6662 let Some(target) = drain_target.as_ref() else {
6663 return Ok(StopNotice::NoConnection);
6667 };
6668 {
6669 send_module_draining(&spec.module_id, reason, deadline_ms, target);
6670 let routes = forwarding
6671 .endpoint_routes(target.endpoint)
6672 .map_err(SuperviseError::Forwarding)?;
6673 let routes_notified = routes.len();
6674 crate::control::send_route_control_pushes(
6675 forwarding,
6676 routes.clone(),
6677 ClientControlPush::RouteClosing {
6678 module_id: spec.module_id.clone(),
6679 reason,
6680 },
6681 );
6682 send_route_goodbyes(forwarding, target.abandoned_bindings.clone());
6683
6684 let wait_result = wait_for_forwarding_quiescence(
6690 forwarding,
6691 &spec.module_id,
6692 runtime,
6693 target.endpoint,
6694 drain_deadline,
6695 &busy_gauges,
6696 scope,
6697 )
6698 .await;
6699 let drained = drained_after_quiescence_wait(&wait_result);
6700 if let Err(err) = &wait_result {
6701 error!(
6702 module_id = %spec.module_id,
6703 ?reason,
6704 error = %err,
6705 "forwarding quiescence wait failed after route.closing; forcing route.closed(drained: false) so the client is not left waiting on an unfulfilled promise"
6706 );
6707 } else if !drained {
6708 let holdouts = forwarding
6714 .endpoint_drain_holdouts(target.endpoint)
6715 .unwrap_or_default();
6716 warn!(
6717 module_id = %spec.module_id,
6718 waited = ?drain_timeout,
6719 ?reason,
6720 held_requests = holdouts.requests,
6721 held_routes = holdouts.routes,
6722 total_routes = holdouts.total_routes,
6723 top_connections = ?holdouts.top_connections,
6724 held = %holdouts
6727 .held
6728 .iter()
6729 .map(|(channel, corr)| format!("{channel}:{corr}"))
6730 .collect::<Vec<_>>()
6731 .join(","),
6732 "route drain timed out before request quiescence; forcing teardown"
6733 );
6734 }
6735 crate::control::send_route_control_pushes(
6736 forwarding,
6737 routes,
6738 ClientControlPush::RouteClosed {
6739 module_id: spec.module_id.clone(),
6740 reason,
6741 drained,
6742 abandoned: target.abandoned_bindings.len() as u32,
6743 excluded_subscriptions: target.excluded_subscriptions,
6744 terminal: Some(terminal),
6745 },
6746 );
6747 wait_result?;
6748
6749 let released_routes = match forwarding.release_module_endpoint_routes(target.endpoint) {
6755 Ok(routes) => routes,
6756 Err(err) => {
6757 warn!(
6758 module_id = %spec.module_id,
6759 ?reason,
6760 error = %err,
6761 "failed to release module endpoint routes after route.closed; module GOODBYE will still be sent"
6762 );
6763 send_module_goodbye(&spec.module_id, forwarding, target);
6764 return Err(SuperviseError::Forwarding(err));
6765 }
6766 };
6767 let route_goodbye_count = released_routes.len();
6768 send_route_goodbyes(forwarding, released_routes);
6769 send_module_goodbye(&spec.module_id, forwarding, target);
6770
6771 info!(
6777 module_id = %spec.module_id,
6778 ?reason,
6779 routes_notified,
6780 route_goodbyes = route_goodbye_count,
6781 abandoned_reservations = target.abandoned_bindings.len(),
6782 excluded_subscriptions = target.excluded_subscriptions,
6783 drained,
6784 "module drain complete; consumers notified via route.closing/route.closed pushes and per-route GOODBYE frames"
6785 );
6786 }
6787
6788 Ok(StopNotice::SentOverConnection)
6789}
6790
6791async fn wait_for_registration_after_reload(
6794 registry: &Registry,
6795 module_id: &str,
6796 snapshot: &SharedSnapshot,
6797 child: &mut SupervisedChild,
6798 wait: Duration,
6799) -> Result<RegistrationWaitOutcome, SuperviseError> {
6800 wait_for_slot_registration(
6801 registry,
6802 crate::registry::RegistrationSlot::Active(module_id),
6803 module_id,
6804 snapshot,
6805 child,
6806 wait,
6807 )
6808 .await
6809}
6810
6811async fn wait_for_slot_registration(
6819 registry: &Registry,
6820 slot: crate::registry::RegistrationSlot<'_>,
6821 module_id: &str,
6822 snapshot: &SharedSnapshot,
6823 child: &mut SupervisedChild,
6824 wait: Duration,
6825) -> Result<RegistrationWaitOutcome, SuperviseError> {
6826 let deadline = Instant::now() + wait;
6827 loop {
6828 if registry
6829 .registration(slot)
6830 .map_err(SuperviseError::Registry)?
6831 .is_some()
6832 {
6833 return Ok(RegistrationWaitOutcome::Registered);
6834 }
6835
6836 let now = Instant::now();
6837 if now >= deadline {
6838 return Ok(RegistrationWaitOutcome::TimedOut);
6839 }
6840 let remaining = deadline.saturating_duration_since(now);
6841 let poll = remaining.min(REGISTRY_RELEASE_POLL);
6842
6843 tokio::select! {
6844 wait_result = child.wait() => {
6845 let status = wait_result.map_err(|source| SuperviseError::Wait {
6846 module_id: module_id.to_string(),
6847 source,
6848 })?;
6849 return Ok(RegistrationWaitOutcome::Exited(classify_reaped_child_exit(
6850 snapshot,
6851 child,
6852 &status,
6853 )));
6854 }
6855 _ = sleep(poll) => {}
6856 }
6857 }
6858}
6859
6860fn registration_failure_exit_report(mut exit_report: ExitReport) -> ExitReport {
6861 if exit_report.kind != ExitKind::DeliberateSeverance {
6864 exit_report.kind = ExitKind::Crash;
6865 }
6866 exit_report
6867}
6868
6869async fn handle_reload_child_registration_failure(
6870 spec: &ModuleSpec,
6871 runtime: &SupervisorRuntimeConfig,
6872 registry: &Registry,
6873 process_liveness: &SupervisorProcessLiveness,
6874 snapshot: &SharedSnapshot,
6875 child: &mut Option<SupervisedChild>,
6876 failure: ReloadRegistrationFailure,
6877) -> Result<(), SuperviseError> {
6878 let ReloadRegistrationFailure {
6879 exit_report,
6880 reason,
6881 } = failure;
6882 match on_child_exit(
6883 spec,
6884 runtime.restart_policy,
6885 registry,
6886 snapshot,
6887 &runtime.terminal_ring,
6888 &runtime.spawn_events,
6889 &runtime.child_roster,
6890 exit_report,
6891 )
6892 .await
6893 {
6894 NextAction::Stop {
6895 registration_released,
6896 } => {
6897 if registration_released {
6898 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6899 }
6900 }
6901 NextAction::Restart { schedule } => {
6902 let delay = schedule.map_or(runtime.restart_policy.delay_for_restart(0), |schedule| {
6903 schedule.delay
6904 });
6905 if let Some(schedule) = schedule {
6906 log_crash_respawn(&spec.module_id, schedule);
6907 }
6908 sleep(delay).await;
6909 if respawn_still_pending(snapshot) {
6913 if let Err(err) = wait_for_registration_release(
6914 registry,
6915 &spec.module_id,
6916 REGISTRY_RELEASE_TIMEOUT,
6917 )
6918 .await
6919 {
6920 fail_snapshot(snapshot, Some(&spec.module_id), None);
6921 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6922 return Err(SuperviseError::ReloadFailed {
6923 module_id: spec.module_id.clone(),
6924 reason: format!(
6925 "{reason}; registration did not release before policy retry: {err}"
6926 ),
6927 });
6928 }
6929 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6930 match spawn_and_mark_running(spec, runtime, snapshot) {
6931 Ok(next_child) => {
6932 *child = Some(next_child);
6933 }
6934 Err(err) => {
6935 fail_snapshot(snapshot, Some(&spec.module_id), None);
6936 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6937 return Err(SuperviseError::ReloadFailed {
6938 module_id: spec.module_id.clone(),
6939 reason: format!("{reason}; policy retry spawn failed: {err}"),
6940 });
6941 }
6942 }
6943 }
6944 }
6945 }
6946
6947 Err(SuperviseError::ReloadFailed {
6948 module_id: spec.module_id.clone(),
6949 reason,
6950 })
6951}
6952
6953async fn handle_reload_spawn_failure(
6954 spec: &ModuleSpec,
6955 runtime: &SupervisorRuntimeConfig,
6956 process_liveness: &SupervisorProcessLiveness,
6957 snapshot: &SharedSnapshot,
6958 child: &mut Option<SupervisedChild>,
6959 reason: String,
6960) -> Result<(), SuperviseError> {
6961 let mut should_retry = false;
6962 let now = Instant::now();
6963 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6964 clear_current_process_facts(state);
6965 if daemon_will_restart(state, &runtime.restart_policy, now) {
6966 state.record_crash_restart(&runtime.restart_policy, now);
6967 state.state = ModuleState::Restarting;
6968 should_retry = true;
6969 } else if state.enabled {
6970 state.state = ModuleState::Failed;
6971 } else {
6972 state.state = ModuleState::Disabled;
6973 }
6974 })?;
6975
6976 if should_retry {
6977 sleep(runtime.restart_policy.backoff).await;
6978 if respawn_still_pending(snapshot) {
6982 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6983 match spawn_and_mark_running(spec, runtime, snapshot) {
6984 Ok(next_child) => {
6985 *child = Some(next_child);
6986 }
6987 Err(err) => {
6988 fail_snapshot(snapshot, Some(&spec.module_id), None);
6989 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6990 return Err(SuperviseError::ReloadFailed {
6991 module_id: spec.module_id.clone(),
6992 reason: format!("{reason}; policy retry spawn failed: {err}"),
6993 });
6994 }
6995 }
6996 }
6997 } else {
6998 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6999 }
7000
7001 Err(SuperviseError::ReloadFailed {
7002 module_id: spec.module_id.clone(),
7003 reason,
7004 })
7005}
7006
7007fn control_flags() -> Flags {
7008 Flags::new(false, Priority::Passive, false)
7009}
7010
7011#[allow(clippy::too_many_arguments)]
7012async fn drain_optional_child(
7013 module_id: &str,
7014 protocol: ModuleProtocol,
7015 stop_notice: StopNotice,
7016 registry: &Registry,
7017 snapshot: &SharedSnapshot,
7018 terminal_ring: &Arc<Mutex<TerminalRing>>,
7019 spawn_events: &SpawnEventFeed,
7020 child: &mut Option<SupervisedChild>,
7021 drain_timeout: Duration,
7022 final_state: ModuleState,
7023 enabled: Option<bool>,
7024) -> Result<(), SuperviseError> {
7025 if let Some(child) = child.take() {
7026 drain_child_to_state(
7027 module_id,
7028 protocol,
7029 stop_notice,
7030 registry,
7031 snapshot,
7032 terminal_ring,
7033 spawn_events,
7034 child,
7035 drain_timeout,
7036 final_state,
7037 enabled,
7038 )
7039 .await
7040 } else {
7041 update_snapshot(snapshot, Some(module_id), |state| {
7042 state.state = final_state;
7043 if let Some(enabled) = enabled {
7044 state.enabled = enabled;
7045 }
7046 clear_current_process_facts(state);
7047 })?;
7048 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
7049 }
7050}
7051
7052#[allow(clippy::too_many_arguments)]
7053async fn drain_child_to_state(
7054 module_id: &str,
7055 protocol: ModuleProtocol,
7056 stop_notice: StopNotice,
7057 registry: &Registry,
7058 snapshot: &SharedSnapshot,
7059 terminal_ring: &Arc<Mutex<TerminalRing>>,
7060 spawn_events: &SpawnEventFeed,
7061 mut child: SupervisedChild,
7062 drain_timeout: Duration,
7063 final_state: ModuleState,
7064 enabled: Option<bool>,
7065) -> Result<(), SuperviseError> {
7066 update_snapshot(snapshot, Some(module_id), |state| {
7067 state.state = ModuleState::Draining;
7068 state.draining_to_replace = final_state == ModuleState::Restarting;
7069 if let Some(enabled) = enabled {
7070 state.enabled = enabled;
7071 }
7072 })?;
7073
7074 if stop_notice != StopNotice::SentOverConnection {
7085 if protocol == ModuleProtocol::Subc && stop_notice == StopNotice::NoConnection {
7086 info!(
7087 module_id,
7088 pid = child.pid,
7089 budget_ms = u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX),
7090 "module has no connection yet; requesting stop by signal"
7091 );
7092 }
7093 request_graceful_stop(module_id, &child);
7094 }
7095
7096 let exit_report = match timeout(drain_timeout, child.wait()).await {
7097 Ok(Ok(status)) => classify_reaped_child_exit(snapshot, &child, &status),
7098 Ok(Err(source)) => {
7099 fail_snapshot(snapshot, Some(module_id), None);
7100 return Err(SuperviseError::Wait {
7101 module_id: module_id.to_string(),
7102 source,
7103 });
7104 }
7105 Err(_) => {
7106 warn!(
7119 module_id,
7120 pid = child.pid,
7121 budget_ms = u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX),
7122 reason = ?final_state,
7123 ?stop_notice,
7124 "drain budget expired before the module exited; killing it"
7125 );
7126 child.start_kill().map_err(|source| {
7127 fail_snapshot(snapshot, Some(module_id), None);
7128 SuperviseError::Kill {
7129 module_id: module_id.to_string(),
7130 source,
7131 }
7132 })?;
7133 let status = child.wait().await.map_err(|source| {
7134 fail_snapshot(snapshot, Some(module_id), None);
7135 SuperviseError::Wait {
7136 module_id: module_id.to_string(),
7137 source,
7138 }
7139 })?;
7140 classify_reaped_child_exit(snapshot, &child, &status)
7141 }
7142 };
7143
7144 update_snapshot(snapshot, Some(module_id), |state| {
7145 state.state = final_state;
7146 if let Some(enabled) = enabled {
7147 state.enabled = enabled;
7148 }
7149 clear_current_process_facts(state);
7150 state.last_exit = Some(exit_report.clone());
7151 if exit_report.kind == ExitKind::DeliberateSeverance {
7152 state.lifetime_restarts += 1;
7153 }
7154 })?;
7155 record_terminal(
7156 module_id,
7157 terminal_ring,
7158 spawn_events,
7159 &exit_report,
7160 terminal_disposition(final_state),
7161 );
7162 child.drain_stderr(module_id).await;
7163
7164 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
7165}
7166
7167#[cfg(unix)]
7187fn request_graceful_stop(module_id: &str, child: &SupervisedChild) {
7188 let Some(pid) = child
7189 .id()
7190 .and_then(|pid| i32::try_from(pid).ok())
7191 .and_then(rustix::process::Pid::from_raw)
7192 else {
7193 debug!(
7194 module_id,
7195 "no pid to signal for teardown; falling through to the drain wait"
7196 );
7197 return;
7198 };
7199 match rustix::process::kill_process(pid, rustix::process::Signal::TERM) {
7200 Ok(()) => debug!(
7201 module_id,
7202 "sent SIGTERM to a module nothing else asked to stop"
7203 ),
7204 Err(err) => debug!(
7205 module_id,
7206 error = %err,
7207 "SIGTERM to module failed; the drain wait and kill still apply"
7208 ),
7209 }
7210}
7211
7212#[cfg(not(unix))]
7220fn request_graceful_stop(module_id: &str, _child: &SupervisedChild) {
7221 debug!(
7222 module_id,
7223 "no graceful stop signal exists on this platform; teardown of a module nothing asked to stop waits, then kills"
7224 );
7225}
7226
7227fn terminal_disposition(final_state: ModuleState) -> TerminalDisposition {
7228 match final_state {
7229 ModuleState::Stopped => TerminalDisposition::Stopped,
7230 ModuleState::Disabled => TerminalDisposition::Disabled,
7231 ModuleState::Restarting => TerminalDisposition::Restarting,
7232 ModuleState::Failed => TerminalDisposition::Failed,
7233 ModuleState::Starting
7234 | ModuleState::Running
7235 | ModuleState::Unresponsive
7236 | ModuleState::Draining => {
7237 unreachable!("terminal exits only finish in terminal or restarting states")
7238 }
7239 }
7240}
7241
7242async fn wait_for_registration_release(
7245 registry: &Registry,
7246 module_id: &str,
7247 wait: Duration,
7248) -> Result<(), SuperviseError> {
7249 wait_for_slot_registration_release(
7250 registry,
7251 crate::registry::RegistrationSlot::Active(module_id),
7252 wait,
7253 )
7254 .await
7255}
7256
7257async fn wait_for_slot_registration_release(
7265 registry: &Registry,
7266 slot: crate::registry::RegistrationSlot<'_>,
7267 wait: Duration,
7268) -> Result<(), SuperviseError> {
7269 let deadline = Instant::now() + wait;
7270 let mut release_events = registration_release_events().subscribe();
7271 let still_active = |registration: &crate::registry::ModuleRegistration| {
7272 SuperviseError::RegistrationStillActive {
7273 module_id: registration.manifest.module_id.clone(),
7274 waited: wait,
7275 }
7276 };
7277 loop {
7278 let _observed_generation = *release_events.borrow_and_update();
7279 let Some(registration) = registry
7280 .registration(slot)
7281 .map_err(SuperviseError::Registry)?
7282 else {
7283 return Ok(());
7284 };
7285
7286 let now = Instant::now();
7287 if now >= deadline {
7288 return Err(still_active(®istration));
7289 }
7290
7291 let remaining = deadline.saturating_duration_since(now);
7292 match timeout(remaining, release_events.changed()).await {
7293 Ok(Ok(())) | Ok(Err(_)) => {}
7294 Err(_) => return Err(still_active(®istration)),
7295 }
7296 }
7297}
7298
7299#[cfg(test)]
7300mod slot_registration_wait_tests {
7301 use super::*;
7302 use crate::registry::{ConnectionId, RegistrationSlot};
7303 use subc_protocol::manifest::ModuleManifest;
7304
7305 const INCUMBENT: u64 = 1;
7306 const CANDIDATE: u64 = 2;
7307
7308 fn swapped_registry() -> Arc<Registry> {
7309 let registry = Arc::new(Registry::default());
7310 let manifest = ModuleManifest::builder("m", "0.1.0").build();
7311 registry
7312 .register_with_control_ops(
7313 manifest.clone(),
7314 1,
7315 ConnectionId::new(INCUMBENT),
7316 Vec::new(),
7317 )
7318 .unwrap();
7319 registry
7320 .register_candidate_with_control_ops(
7321 manifest,
7322 1,
7323 ConnectionId::new(CANDIDATE),
7324 Vec::new(),
7325 )
7326 .unwrap();
7327 registry
7328 }
7329
7330 #[tokio::test]
7334 async fn incumbent_release_is_awaited_by_connection_not_by_module_id() {
7335 let registry = swapped_registry();
7336 registry.promote_candidate("m").unwrap().unwrap();
7337
7338 assert!(matches!(
7339 wait_for_registration_release(®istry, "m", Duration::from_millis(50)).await,
7340 Err(SuperviseError::RegistrationStillActive { .. })
7341 ));
7342
7343 assert!(matches!(
7345 wait_for_slot_registration_release(
7346 ®istry,
7347 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
7348 Duration::from_millis(50),
7349 )
7350 .await,
7351 Err(SuperviseError::RegistrationStillActive { .. })
7352 ));
7353
7354 let releaser = Arc::clone(®istry);
7355 let release = tokio::spawn(async move {
7356 sleep(Duration::from_millis(20)).await;
7357 releaser
7358 .deregister_connection(ConnectionId::new(INCUMBENT))
7359 .unwrap();
7360 notify_registration_release();
7361 });
7362 wait_for_slot_registration_release(
7363 ®istry,
7364 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
7365 Duration::from_secs(5),
7366 )
7367 .await
7368 .expect("the incumbent's own registration is released");
7369 release.await.unwrap();
7370 assert!(registry.get_module("m").unwrap().is_some());
7371 }
7372
7373 #[tokio::test]
7376 async fn candidate_slot_wait_ignores_the_incumbents_registration() {
7377 let registry = swapped_registry();
7378 assert!(matches!(
7379 wait_for_slot_registration_release(
7380 ®istry,
7381 RegistrationSlot::Candidate("m"),
7382 Duration::from_millis(50),
7383 )
7384 .await,
7385 Err(SuperviseError::RegistrationStillActive { .. })
7386 ));
7387 registry
7388 .deregister_connection(ConnectionId::new(CANDIDATE))
7389 .unwrap();
7390 wait_for_slot_registration_release(
7391 ®istry,
7392 RegistrationSlot::Candidate("m"),
7393 Duration::from_millis(50),
7394 )
7395 .await
7396 .expect("a candidate slot with no candidate is released");
7397 assert!(registry
7398 .registration(RegistrationSlot::Active("m"))
7399 .unwrap()
7400 .is_some());
7401 }
7402}
7403
7404fn classify_exit(status: &ExitStatus) -> ExitReport {
7405 ExitReport {
7406 kind: if status.success() {
7407 ExitKind::Clean
7408 } else {
7409 ExitKind::Crash
7410 },
7411 code: status.code(),
7412 signal: exit_signal(status),
7413 at_ms: unix_ms_now(),
7414 }
7415}
7416
7417fn wait_error_exit_report() -> ExitReport {
7423 ExitReport {
7424 kind: ExitKind::Crash,
7425 code: None,
7426 signal: None,
7427 at_ms: unix_ms_now(),
7428 }
7429}
7430
7431#[cfg(unix)]
7432fn exit_signal(status: &ExitStatus) -> Option<i32> {
7433 use std::os::unix::process::ExitStatusExt;
7434
7435 status.signal()
7436}
7437
7438#[cfg(not(unix))]
7439fn exit_signal(_status: &ExitStatus) -> Option<i32> {
7440 None
7441}
7442
7443fn reset_restart_count(snapshot: &SharedSnapshot, module_id: &str) -> Result<(), SuperviseError> {
7449 update_snapshot(snapshot, Some(module_id), |state| {
7450 state.clear_crash_restarts();
7451 })
7452}
7453
7454fn set_running(
7455 snapshot: &SharedSnapshot,
7456 child: &SupervisedChild,
7457 module_id: &str,
7458 spawn_events: &SpawnEventFeed,
7459) -> Result<(), SuperviseError> {
7460 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
7461 module_id: Some(module_id.to_string()),
7462 })?;
7463 state.spawn_generation = spawn_events.emit_spawned(module_id, child.pid, child.spawned_at_ms);
7464 state.in_alternate_slot = false;
7467 state.configuration_updated_since_spawn = false;
7468 state.state = ModuleState::Running;
7469 state.enabled = true;
7470 state.process_alive = true;
7471 state.pid = child.id();
7472 state.spawned_at_ms = Some(child.spawned_at_ms);
7473 state.spawned_from = Some(child.spawned_from.clone());
7474 state.spawned_file_identity = child.spawned_file_identity;
7475 state.process_start_time = child.process_start_time;
7476 Ok(())
7477}
7478
7479fn clear_current_process_facts(state: &mut SupervisorSnapshot) {
7480 state.process_alive = false;
7481 state.pid = None;
7482 state.spawned_at_ms = None;
7483 state.spawned_from = None;
7484 state.spawned_file_identity = None;
7485 state.process_start_time = None;
7486 state.deliberate_severance = None;
7487}
7488
7489#[cfg(test)]
7490fn record_deliberate_severance(
7491 snapshot: &SharedSnapshot,
7492 identity: ProcessIdentity,
7493) -> Result<(), SuperviseError> {
7494 update_snapshot(snapshot, None, |state| {
7495 state.deliberate_severance = Some(identity);
7496 })
7497}
7498
7499fn apply_deliberate_severance_marker(
7500 snapshot: &SharedSnapshot,
7501 exited_identity: Option<ProcessIdentity>,
7502 mut exit_report: ExitReport,
7503) -> ExitReport {
7504 let marker = lock_snapshot(snapshot)
7505 .ok()
7506 .and_then(|mut state| state.deliberate_severance.take());
7507 if marker.is_some() && marker == exited_identity {
7508 exit_report.kind = ExitKind::DeliberateSeverance;
7509 }
7510 exit_report
7511}
7512
7513fn classify_reaped_child_exit(
7514 snapshot: &SharedSnapshot,
7515 child: &SupervisedChild,
7516 status: &ExitStatus,
7517) -> ExitReport {
7518 apply_deliberate_severance_marker(snapshot, child.process_identity(), classify_exit(status))
7519}
7520
7521fn fail_snapshot(
7522 snapshot: &SharedSnapshot,
7523 module_id: Option<&str>,
7524 last_exit: Option<ExitReport>,
7525) {
7526 if let Err(err) = update_snapshot(snapshot, module_id, |state| {
7527 state.state = ModuleState::Failed;
7528 clear_current_process_facts(state);
7529 if let Some(last_exit) = last_exit {
7530 state.last_exit = Some(last_exit);
7531 }
7532 }) {
7533 error!(error = %err, "failed to mark supervisor state failed");
7534 }
7535}
7536
7537fn update_snapshot(
7538 snapshot: &SharedSnapshot,
7539 module_id: Option<&str>,
7540 update: impl FnOnce(&mut SupervisorSnapshot),
7541) -> Result<(), SuperviseError> {
7542 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
7543 module_id: module_id.map(ToOwned::to_owned),
7544 })?;
7545 update(&mut state);
7546 Ok(())
7547}
7548
7549const SLOW_SNAPSHOT_LOCK_THRESHOLD: Duration = Duration::from_millis(250);
7550
7551fn lock_snapshot_for_control<'a>(
7552 snapshot: &'a SharedSnapshot,
7553 module_id: &str,
7554 caller: &'static str,
7555) -> Result<std::sync::MutexGuard<'a, SupervisorSnapshot>, SuperviseError> {
7556 let started_at = Instant::now();
7557 let guard = lock_snapshot(snapshot)?;
7558 let waited = started_at.elapsed();
7559 if waited >= SLOW_SNAPSHOT_LOCK_THRESHOLD {
7560 warn!(
7561 module_id = %module_id,
7562 waited_ms = waited.as_millis() as u64,
7563 caller = %caller,
7564 "slow snapshot lock"
7565 );
7566 }
7567 Ok(guard)
7568}
7569
7570fn lock_snapshot(
7571 snapshot: &SharedSnapshot,
7572) -> Result<std::sync::MutexGuard<'_, SupervisorSnapshot>, SuperviseError> {
7573 snapshot
7574 .lock()
7575 .map_err(|_| SuperviseError::StatePoisoned { module_id: None })
7576}
7577
7578#[cfg(test)]
7579mod terminal_history_tests {
7580 use std::{
7581 path::PathBuf,
7582 sync::Arc,
7583 time::{Duration, Instant},
7584 };
7585
7586 use tokio::time::sleep;
7587
7588 use super::{
7589 apply_deliberate_severance_marker, daemon_will_restart, drain_child_to_state,
7590 drained_after_quiescence_wait, handle_reload_spawn_failure, health_restart_child,
7591 lock_snapshot, on_child_exit, record_deliberate_severance, record_wait_error_terminal,
7592 reset_restart_count, spawn_and_mark_running, update_snapshot, wait_error_exit_report,
7593 ExitKind, ExitReport, ModuleProtocol, ModuleSpec, ModuleState, NextAction, ProcessIdentity,
7594 RestartPolicy, SpawnEventKind, StopNotice, SuperviseError, SupervisedModule, Supervisor,
7595 SupervisorHandle, SupervisorHealthStatus, SupervisorSnapshot,
7596 };
7597 use super::Instant as ClockInstant;
7602 use crate::{
7603 registry::Registry,
7604 terminal_ring::{TerminalRing, TerminalRingConfig},
7605 };
7606 use std::sync::Mutex;
7607 use subc_control::TerminalDisposition;
7608
7609 fn fake_aft_stub_path() -> PathBuf {
7614 let mut path = std::env::current_exe().expect("current_exe available in tests");
7615 path.pop();
7616 path.pop();
7617 path.push(if cfg!(windows) {
7618 "fake-aft-stub.exe"
7619 } else {
7620 "fake-aft-stub"
7621 });
7622 assert!(
7623 path.exists(),
7624 "fake-aft-stub not built at {}: run `cargo test -p subc-core` (which builds \
7625 [[bin]] targets) rather than `cargo test -p subc-core --lib` (which does not)",
7626 path.display()
7627 );
7628 path
7629 }
7630
7631 #[test]
7632 fn reserved_never_spawned_refuses_every_hello() {
7633 let supervisor = SupervisorHandle::default();
7638 supervisor.apply_identity_configuration(&ModuleSpec {
7639 module_id: "never-spawned".to_string(),
7640 program: PathBuf::from("/usr/bin/false"),
7641 args: Vec::new(),
7642 env: Vec::new(),
7643 reserved: true,
7644 reserved_prefixes: Vec::new(),
7645 protocol: ModuleProtocol::Subc,
7646 overlap: Default::default(),
7647 });
7648 assert!(
7649 supervisor
7650 .reserved_hello_rejection("never-spawned", Some("any-forged-nonce"))
7651 .is_some(),
7652 "forged nonce must refuse on a reserved never-spawned id"
7653 );
7654 assert!(
7655 supervisor
7656 .reserved_hello_rejection("never-spawned", None)
7657 .is_some(),
7658 "absent nonce must refuse on a reserved never-spawned id"
7659 );
7660 supervisor.set_spawn_nonce("never-spawned", "minted".to_string());
7662 supervisor.apply_identity_configuration(&ModuleSpec {
7663 module_id: "never-spawned".to_string(),
7664 program: PathBuf::from("/usr/bin/false"),
7665 args: Vec::new(),
7666 env: Vec::new(),
7667 reserved: true,
7668 reserved_prefixes: Vec::new(),
7669 protocol: ModuleProtocol::Subc,
7670 overlap: Default::default(),
7671 });
7672 assert!(supervisor
7673 .reserved_hello_rejection("never-spawned", Some("minted"))
7674 .is_none());
7675 assert!(supervisor
7676 .reserved_hello_rejection("never-spawned", Some("forged"))
7677 .is_some());
7678 }
7679
7680 fn seed_crash_restarts(state: &mut SupervisorSnapshot, count: u32) {
7683 let now = ClockInstant::now();
7684 for _ in 0..count {
7685 state.crash_restarts.push_back(now);
7686 }
7687 }
7688
7689 fn age_oldest_crash_restart_out_of_window(state: &mut SupervisorSnapshot, window: Duration) {
7693 let aged = state
7694 .crash_restarts
7695 .front()
7696 .expect("a crash restart must be recorded before it can be aged")
7697 .checked_sub(window + Duration::from_secs(1))
7698 .expect("the test clock is far enough from its origin to age an instant");
7699 state.crash_restarts[0] = aged;
7700 }
7701
7702 fn snapshot_with_restarts(enabled: bool, count: u32) -> SupervisorSnapshot {
7703 let mut state = SupervisorSnapshot::new(ModuleState::Running, enabled);
7704 seed_crash_restarts(&mut state, count);
7705 state
7706 }
7707
7708 #[test]
7709 fn daemon_owned_recovery_predicate_uses_the_pre_increment_budget() {
7710 let policy = RestartPolicy::new(3, Duration::ZERO);
7711 let now = ClockInstant::now();
7712 assert!(daemon_will_restart(
7713 &mut snapshot_with_restarts(true, 2),
7714 &policy,
7715 now
7716 ));
7717 assert!(!daemon_will_restart(
7718 &mut snapshot_with_restarts(true, 3),
7719 &policy,
7720 now
7721 ));
7722 assert!(!daemon_will_restart(
7723 &mut snapshot_with_restarts(false, 0),
7724 &policy,
7725 now
7726 ));
7727 }
7728
7729 #[test]
7730 fn crash_restart_backoff_escalates_with_in_window_count() {
7731 let policy = RestartPolicy::new(4, Duration::from_millis(100))
7732 .with_max_backoff(Duration::from_secs(30));
7733 let now = ClockInstant::now();
7734 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7735 let schedules = (0..4)
7736 .map(|_| {
7737 state
7738 .next_crash_restart(&policy, now)
7739 .expect("the test policy allows four crash restarts")
7740 })
7741 .collect::<Vec<_>>();
7742
7743 assert_eq!(
7744 schedules
7745 .iter()
7746 .map(|schedule| schedule.restart_in_window)
7747 .collect::<Vec<_>>(),
7748 vec![0, 1, 2, 3]
7749 );
7750 assert_eq!(
7751 schedules
7752 .iter()
7753 .map(|schedule| schedule.delay)
7754 .collect::<Vec<_>>(),
7755 vec![
7756 Duration::from_millis(100),
7757 Duration::from_secs(1),
7758 Duration::from_secs(10),
7759 Duration::from_secs(30),
7760 ]
7761 );
7762 }
7763
7764 #[test]
7765 fn crash_restart_backoff_resets_after_ring_clear() {
7766 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7767 let now = ClockInstant::now();
7768 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7769 assert_eq!(
7770 state.next_crash_restart(&policy, now).unwrap().delay,
7771 Duration::from_millis(100)
7772 );
7773 assert_eq!(
7774 state.next_crash_restart(&policy, now).unwrap().delay,
7775 Duration::from_secs(1)
7776 );
7777
7778 state.clear_crash_restarts();
7779 let schedule = state
7780 .next_crash_restart(&policy, now)
7781 .expect("a cleared ring must allow another restart");
7782 assert_eq!(schedule.restart_in_window, 0);
7783 assert_eq!(schedule.delay, Duration::from_millis(100));
7784 }
7785
7786 #[test]
7787 fn crash_restart_backoff_ignores_aged_restarts() {
7788 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7789 let now = ClockInstant::now();
7790 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7791 state
7792 .next_crash_restart(&policy, now)
7793 .expect("the first restart is allowed");
7794 state
7795 .next_crash_restart(&policy, now)
7796 .expect("the second restart is allowed");
7797 state.crash_restarts[0] = now
7798 .checked_sub(policy.window + Duration::from_secs(1))
7799 .expect("the fake clock can age a restart past the window");
7800
7801 let schedule = state
7802 .next_crash_restart(&policy, now)
7803 .expect("an aged restart must release its slot");
7804 assert_eq!(schedule.restart_in_window, 1);
7805 assert_eq!(schedule.delay, Duration::from_secs(1));
7806 assert_eq!(state.crash_restarts.len(), 2);
7807 }
7808
7809 #[test]
7813 fn a_budget_spent_before_the_window_no_longer_refuses() {
7814 let policy = RestartPolicy::new(3, Duration::ZERO);
7815 let mut state = snapshot_with_restarts(true, 3);
7816 let now = ClockInstant::now();
7817 assert!(!daemon_will_restart(&mut state, &policy, now));
7818
7819 assert!(daemon_will_restart(
7820 &mut state,
7821 &policy,
7822 now + policy.window + Duration::from_secs(1)
7823 ));
7824 assert!(
7825 state.crash_restarts.is_empty(),
7826 "reading the budget must drop the instants that left the window"
7827 );
7828 }
7829
7830 fn module_with_recovery_snapshot(
7831 state: ModuleState,
7832 enabled: bool,
7833 restart_count: u32,
7834 ) -> SupervisedModule {
7835 let registry = Arc::new(Registry::default());
7836 let supervisor =
7837 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(3, Duration::ZERO));
7838 let module = supervisor
7839 .spawn(ModuleSpec {
7840 module_id: "recovery-snapshot".to_string(),
7841 program: fake_aft_stub_path(),
7842 args: Vec::new(),
7843 env: Vec::new(),
7844 reserved: false,
7845 reserved_prefixes: Vec::new(),
7846 protocol: ModuleProtocol::Subc,
7847 overlap: Default::default(),
7848 })
7849 .unwrap();
7850 update_snapshot(
7851 &module.inner.snapshot,
7852 Some("recovery-snapshot"),
7853 |snapshot| {
7854 snapshot.state = state;
7855 snapshot.enabled = enabled;
7856 seed_crash_restarts(snapshot, restart_count);
7857 },
7858 )
7859 .unwrap();
7860 module
7861 }
7862
7863 #[cfg(target_os = "linux")]
7864 #[tokio::test]
7865 async fn no_cgroup_placement_does_not_block_fake_aft_stub_spawn() {
7866 let supervisor = Supervisor::new(Arc::new(Registry::default()), RestartPolicy::default())
7867 .with_cgroup_placement(None);
7868 let result = supervisor.spawn(ModuleSpec {
7869 module_id: "no-cgroup-placement".to_string(),
7870 program: fake_aft_stub_path(),
7871 args: Vec::new(),
7872 env: Vec::new(),
7873 reserved: false,
7874 reserved_prefixes: Vec::new(),
7875 protocol: ModuleProtocol::Subc,
7876 overlap: Default::default(),
7877 });
7878
7879 assert!(
7880 result.is_ok(),
7881 "no delegation must not turn an otherwise valid spawn into a failure: {result:?}"
7882 );
7883 }
7884
7885 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7886 async fn undecided_snapshot_uses_shared_restart_predicate() {
7887 assert!(module_with_recovery_snapshot(ModuleState::Running, true, 2)
7888 .will_recover_after_connection_loss()
7889 .unwrap());
7890 assert!(
7891 !module_with_recovery_snapshot(ModuleState::Running, true, 3)
7892 .will_recover_after_connection_loss()
7893 .unwrap()
7894 );
7895 }
7896
7897 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7898 async fn restarting_snapshot_at_exhausted_budget_is_non_terminal() {
7899 assert!(
7900 module_with_recovery_snapshot(ModuleState::Restarting, true, 3)
7901 .will_recover_after_connection_loss()
7902 .unwrap()
7903 );
7904 }
7905
7906 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7907 async fn terminal_phase_snapshots_are_terminal_before_budget_exhaustion() {
7908 assert!(!module_with_recovery_snapshot(ModuleState::Failed, true, 0)
7909 .will_recover_after_connection_loss()
7910 .unwrap());
7911 assert!(
7912 !module_with_recovery_snapshot(ModuleState::Disabled, true, 0)
7913 .will_recover_after_connection_loss()
7914 .unwrap()
7915 );
7916 }
7917
7918 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7919 async fn warming_snapshot_is_limited_to_startup_phases() {
7920 for state in [
7921 ModuleState::Starting,
7922 ModuleState::Running,
7923 ModuleState::Restarting,
7924 ] {
7925 assert!(
7926 module_with_recovery_snapshot(state, true, 0)
7927 .is_warming()
7928 .unwrap(),
7929 "{state:?} should be warming"
7930 );
7931 }
7932 for state in [
7933 ModuleState::Unresponsive,
7934 ModuleState::Draining,
7935 ModuleState::Stopped,
7936 ModuleState::Failed,
7937 ModuleState::Disabled,
7938 ] {
7939 assert!(
7940 !module_with_recovery_snapshot(state, true, 0)
7941 .is_warming()
7942 .unwrap(),
7943 "{state:?} should not be warming"
7944 );
7945 }
7946 }
7947
7948 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7949 async fn terminal_history_survives_respawn_and_keeps_both_crashes_in_order() {
7950 let registry = Arc::new(Registry::default());
7951 let supervisor =
7952 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(1, Duration::ZERO));
7953 let module = supervisor
7954 .spawn(ModuleSpec {
7955 module_id: "terminal-history".to_string(),
7956 program: fake_aft_stub_path(),
7957 args: Vec::new(),
7958 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7959 reserved: false,
7960 reserved_prefixes: Vec::new(),
7961 protocol: ModuleProtocol::Subc,
7962 overlap: Default::default(),
7963 })
7964 .unwrap();
7965
7966 let deadline = Instant::now() + Duration::from_secs(5);
7967 loop {
7968 let history = module.terminal_history();
7969 if history.entries.len() == 2 {
7970 assert_eq!(module.status().unwrap().state, ModuleState::Failed);
7971 assert_eq!(history.dropped, 0);
7972 assert_eq!(
7973 history
7974 .entries
7975 .iter()
7976 .map(|entry| entry.exit_code)
7977 .collect::<Vec<_>>(),
7978 vec![Some(23), Some(23)]
7979 );
7980 assert!(history.entries[0].at_ms <= history.entries[1].at_ms);
7981 return;
7982 }
7983 assert!(
7984 Instant::now() < deadline,
7985 "module did not retain two terminal exits: {history:?}"
7986 );
7987 sleep(Duration::from_millis(10)).await;
7988 }
7989 }
7990
7991 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7995 async fn disable_during_crash_backoff_cancels_pending_respawn() {
7996 let backoff = Duration::from_secs(2);
7997 let supervisor = Supervisor::new(
7998 Arc::new(Registry::default()),
7999 RestartPolicy::new(10, backoff),
8000 );
8001 let module = supervisor
8002 .spawn(ModuleSpec {
8003 module_id: "disable-during-backoff".to_string(),
8004 program: fake_aft_stub_path(),
8005 args: Vec::new(),
8006 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
8007 reserved: false,
8008 reserved_prefixes: Vec::new(),
8009 protocol: ModuleProtocol::Subc,
8010 overlap: Default::default(),
8011 })
8012 .unwrap();
8013
8014 let deadline = Instant::now() + Duration::from_secs(5);
8016 loop {
8017 if module.status().unwrap().state == ModuleState::Restarting {
8018 break;
8019 }
8020 assert!(
8021 Instant::now() < deadline,
8022 "module never entered the crash backoff"
8023 );
8024 sleep(Duration::from_millis(10)).await;
8025 }
8026
8027 let started = Instant::now();
8028 module.set_enabled(false).await.unwrap();
8029 let waited = started.elapsed();
8030
8031 assert!(
8032 waited < backoff / 2,
8033 "disable waited {waited:?} behind the {backoff:?} crash backoff; the operator command must preempt the pending respawn"
8034 );
8035 assert_eq!(module.status().unwrap().state, ModuleState::Disabled);
8036
8037 sleep(backoff + Duration::from_millis(500)).await;
8039 let status = module.status().unwrap();
8040 assert_eq!(status.state, ModuleState::Disabled);
8041 assert_eq!(
8042 status.spawn_generation, 1,
8043 "module respawned after the operator disabled it"
8044 );
8045 }
8046
8047 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
8051 async fn every_restart_increment_path_advances_lifetime_count() {
8052 let supervisor = Supervisor::new(
8053 Arc::new(Registry::default()),
8054 RestartPolicy::new(1, Duration::ZERO),
8055 );
8056 let runtime = supervisor.runtime_config();
8057 let spec = ModuleSpec {
8058 module_id: "lifetime-increment-path".to_string(),
8059 program: PathBuf::from("/unused/lifetime-increment-path"),
8060 args: Vec::new(),
8061 env: Vec::new(),
8062 reserved: false,
8063 reserved_prefixes: Vec::new(),
8064 protocol: ModuleProtocol::Subc,
8065 overlap: Default::default(),
8066 };
8067
8068 let crash_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8069 assert!(matches!(
8070 on_child_exit(
8071 &spec,
8072 runtime.restart_policy,
8073 &supervisor.registry,
8074 &crash_snapshot,
8075 &runtime.terminal_ring,
8076 &runtime.spawn_events,
8077 &runtime.child_roster,
8078 ExitReport {
8079 kind: ExitKind::Crash,
8080 code: Some(1),
8081 signal: None,
8082 at_ms: 1,
8083 },
8084 )
8085 .await,
8086 NextAction::Restart { schedule: _ }
8087 ));
8088 let (crash_restarts, crash_lifetime) = {
8089 let state = lock_snapshot(&crash_snapshot).unwrap();
8090 (state.crash_restarts.len(), state.lifetime_restarts)
8091 };
8092 assert_eq!(crash_restarts, 1);
8093 assert_eq!(crash_lifetime, 1);
8094
8095 let health_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8096 let mut health_child = None;
8097 assert!(matches!(
8098 health_restart_child(
8099 &spec,
8100 &runtime,
8101 &supervisor.registry,
8102 &supervisor.process_liveness,
8103 &health_snapshot,
8104 &mut health_child,
8105 SupervisorHealthStatus::Failing,
8106 None,
8107 2,
8108 )
8109 .await,
8110 Err(SuperviseError::Spawn { .. })
8111 ));
8112 let (health_restarts, health_lifetime) = {
8113 let state = lock_snapshot(&health_snapshot).unwrap();
8114 (state.crash_restarts.len(), state.lifetime_restarts)
8115 };
8116 assert_eq!(health_restarts, 1);
8117 assert_eq!(health_lifetime, 1);
8118
8119 let reload_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8120 let mut reload_child = None;
8121 assert!(matches!(
8122 handle_reload_spawn_failure(
8123 &spec,
8124 &runtime,
8125 &supervisor.process_liveness,
8126 &reload_snapshot,
8127 &mut reload_child,
8128 "forced reload spawn failure".to_string(),
8129 )
8130 .await,
8131 Err(SuperviseError::ReloadFailed { .. })
8132 ));
8133 let (reload_restarts, reload_lifetime) = {
8134 let state = lock_snapshot(&reload_snapshot).unwrap();
8135 (state.crash_restarts.len(), state.lifetime_restarts)
8136 };
8137 assert_eq!(reload_restarts, 1);
8138 assert_eq!(reload_lifetime, 1);
8139 }
8140
8141 #[tokio::test]
8142 async fn deliberately_severed_live_child_records_lifetime_without_spending_restart_budget() {
8143 let supervisor = Supervisor::new(
8144 Arc::new(Registry::default()),
8145 RestartPolicy::new(3, Duration::ZERO),
8146 );
8147 let runtime = supervisor.runtime_config();
8148 let spec = ModuleSpec {
8149 module_id: "deliberately-severed".to_string(),
8150 program: PathBuf::from("/unused/deliberately-severed"),
8151 args: Vec::new(),
8152 env: Vec::new(),
8153 reserved: false,
8154 reserved_prefixes: Vec::new(),
8155 protocol: ModuleProtocol::Subc,
8156 overlap: Default::default(),
8157 };
8158 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8159 let process = ProcessIdentity {
8160 pid: 41,
8161 start_time: 101,
8162 };
8163 record_deliberate_severance(&snapshot, process).unwrap();
8164 let exit_report = apply_deliberate_severance_marker(
8165 &snapshot,
8166 Some(process),
8167 ExitReport {
8168 kind: ExitKind::Crash,
8169 code: Some(1),
8170 signal: None,
8171 at_ms: 1,
8172 },
8173 );
8174 assert_eq!(exit_report.kind, ExitKind::DeliberateSeverance);
8175
8176 assert!(matches!(
8177 on_child_exit(
8178 &spec,
8179 runtime.restart_policy,
8180 &supervisor.registry,
8181 &snapshot,
8182 &runtime.terminal_ring,
8183 &runtime.spawn_events,
8184 &runtime.child_roster,
8185 exit_report,
8186 )
8187 .await,
8188 NextAction::Restart { schedule: _ }
8189 ));
8190 let state = lock_snapshot(&snapshot).unwrap();
8191 assert_eq!(state.lifetime_restarts, 1);
8192 assert_eq!(state.crash_restarts.len(), 0);
8193 }
8194
8195 #[tokio::test]
8196 async fn genuine_crash_spends_restart_budget_and_records_lifetime() {
8197 let supervisor = Supervisor::new(
8198 Arc::new(Registry::default()),
8199 RestartPolicy::new(3, Duration::ZERO),
8200 );
8201 let runtime = supervisor.runtime_config();
8202 let spec = ModuleSpec {
8203 module_id: "genuine-crash".to_string(),
8204 program: PathBuf::from("/unused/genuine-crash"),
8205 args: Vec::new(),
8206 env: Vec::new(),
8207 reserved: false,
8208 reserved_prefixes: Vec::new(),
8209 protocol: ModuleProtocol::Subc,
8210 overlap: Default::default(),
8211 };
8212 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8213
8214 assert!(matches!(
8215 on_child_exit(
8216 &spec,
8217 runtime.restart_policy,
8218 &supervisor.registry,
8219 &snapshot,
8220 &runtime.terminal_ring,
8221 &runtime.spawn_events,
8222 &runtime.child_roster,
8223 ExitReport {
8224 kind: ExitKind::Crash,
8225 code: Some(1),
8226 signal: None,
8227 at_ms: 1,
8228 },
8229 )
8230 .await,
8231 NextAction::Restart { schedule: _ }
8232 ));
8233 let state = lock_snapshot(&snapshot).unwrap();
8234 assert_eq!(state.lifetime_restarts, 1);
8235 assert_eq!(state.crash_restarts.len(), 1);
8236 }
8237
8238 fn crash_exit_report(at_ms: u64) -> ExitReport {
8239 ExitReport {
8240 kind: ExitKind::Crash,
8241 code: Some(1),
8242 signal: None,
8243 at_ms,
8244 }
8245 }
8246
8247 fn windowed_crash_spec(module_id: &str) -> ModuleSpec {
8248 ModuleSpec {
8249 module_id: module_id.to_string(),
8250 program: PathBuf::from("/unused").join(module_id),
8251 args: Vec::new(),
8252 env: Vec::new(),
8253 reserved: false,
8254 reserved_prefixes: Vec::new(),
8255 protocol: ModuleProtocol::Subc,
8256 overlap: Default::default(),
8257 }
8258 }
8259
8260 #[tokio::test]
8266 async fn three_crashes_inside_the_window_stop_the_module_and_name_the_window() {
8267 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::ERROR);
8268 let supervisor = Supervisor::new(
8269 Arc::new(Registry::default()),
8270 RestartPolicy::new(2, Duration::ZERO),
8271 );
8272 let runtime = supervisor.runtime_config();
8273 let spec = windowed_crash_spec("crash-loop-in-window");
8274 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8275
8276 for attempt in 1..=2 {
8277 assert!(
8278 matches!(
8279 on_child_exit(
8280 &spec,
8281 runtime.restart_policy,
8282 &supervisor.registry,
8283 &snapshot,
8284 &runtime.terminal_ring,
8285 &runtime.spawn_events,
8286 &runtime.child_roster,
8287 crash_exit_report(attempt),
8288 )
8289 .await,
8290 NextAction::Restart { schedule: _ }
8291 ),
8292 "crash {attempt} is inside the budget and must respawn"
8293 );
8294 }
8295
8296 assert!(matches!(
8297 on_child_exit(
8298 &spec,
8299 runtime.restart_policy,
8300 &supervisor.registry,
8301 &snapshot,
8302 &runtime.terminal_ring,
8303 &runtime.spawn_events,
8304 &runtime.child_roster,
8305 crash_exit_report(3),
8306 )
8307 .await,
8308 NextAction::Stop { .. }
8309 ));
8310
8311 {
8312 let state = lock_snapshot(&snapshot).unwrap();
8313 assert_eq!(state.state, ModuleState::Failed);
8314 assert_eq!(state.crash_restarts.len(), 2);
8315 assert_eq!(state.lifetime_restarts, 2);
8316 }
8317
8318 let history = runtime
8319 .terminal_ring
8320 .lock()
8321 .expect("terminal ring is not poisoned")
8322 .snapshot();
8323 let last = history
8324 .entries
8325 .last()
8326 .expect("the refused crash is retained");
8327 assert_eq!(last.disposition, TerminalDisposition::Failed);
8328 assert_eq!(
8329 last.disposition_detail.as_deref(),
8330 Some("crash budget exhausted: max_restarts=2 within window_secs=600")
8331 );
8332
8333 let captured = crate::router::test_log::captured_logs(&logs);
8334 assert!(
8335 captured.contains("crash budget exhausted: max_restarts=2 within window_secs=600"),
8336 "the stop must be logged with its window: {captured}"
8337 );
8338 }
8339
8340 #[tokio::test]
8348 async fn a_crash_older_than_the_window_frees_its_slot_for_a_later_crash() {
8349 let supervisor = Supervisor::new(
8350 Arc::new(Registry::default()),
8351 RestartPolicy::new(2, Duration::ZERO),
8352 );
8353 let runtime = supervisor.runtime_config();
8354 let spec = windowed_crash_spec("crash-across-windows");
8355 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8356
8357 for attempt in 1..=2 {
8358 assert!(matches!(
8359 on_child_exit(
8360 &spec,
8361 runtime.restart_policy,
8362 &supervisor.registry,
8363 &snapshot,
8364 &runtime.terminal_ring,
8365 &runtime.spawn_events,
8366 &runtime.child_roster,
8367 crash_exit_report(attempt),
8368 )
8369 .await,
8370 NextAction::Restart { schedule: _ }
8371 ));
8372 }
8373
8374 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
8377 age_oldest_crash_restart_out_of_window(state, runtime.restart_policy.window);
8378 })
8379 .unwrap();
8380
8381 assert!(
8382 matches!(
8383 on_child_exit(
8384 &spec,
8385 runtime.restart_policy,
8386 &supervisor.registry,
8387 &snapshot,
8388 &runtime.terminal_ring,
8389 &runtime.spawn_events,
8390 &runtime.child_roster,
8391 crash_exit_report(3),
8392 )
8393 .await,
8394 NextAction::Restart { schedule: _ }
8395 ),
8396 "a crash older than the window must not hold a budget slot"
8397 );
8398
8399 let state = lock_snapshot(&snapshot).unwrap();
8400 assert_eq!(state.state, ModuleState::Restarting);
8401 assert_eq!(
8402 state.crash_restarts.len(),
8403 2,
8404 "the aged instant is dropped and the new one takes its place"
8405 );
8406 assert_eq!(
8407 state.lifetime_restarts, 3,
8408 "the ledger counts every restart, including the ones the window forgot"
8409 );
8410 }
8411
8412 #[tokio::test]
8417 async fn an_operator_restart_clears_the_ring_and_leaves_the_ledger_alone() {
8418 let supervisor = Supervisor::new(
8419 Arc::new(Registry::default()),
8420 RestartPolicy::new(2, Duration::ZERO),
8421 );
8422 let runtime = supervisor.runtime_config();
8423 let spec = windowed_crash_spec("operator-cleared-budget");
8424 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8425
8426 for attempt in 1..=2 {
8427 assert!(matches!(
8428 on_child_exit(
8429 &spec,
8430 runtime.restart_policy,
8431 &supervisor.registry,
8432 &snapshot,
8433 &runtime.terminal_ring,
8434 &runtime.spawn_events,
8435 &runtime.child_roster,
8436 crash_exit_report(attempt),
8437 )
8438 .await,
8439 NextAction::Restart { schedule: _ }
8440 ));
8441 }
8442
8443 reset_restart_count(&snapshot, &spec.module_id).unwrap();
8444 {
8445 let state = lock_snapshot(&snapshot).unwrap();
8446 assert!(
8447 state.crash_restarts.is_empty(),
8448 "an operator restart returns the full budget"
8449 );
8450 assert_eq!(
8451 state.lifetime_restarts, 2,
8452 "clearing the budget must not unmake the crashes"
8453 );
8454 }
8455
8456 assert!(
8457 matches!(
8458 on_child_exit(
8459 &spec,
8460 runtime.restart_policy,
8461 &supervisor.registry,
8462 &snapshot,
8463 &runtime.terminal_ring,
8464 &runtime.spawn_events,
8465 &runtime.child_roster,
8466 crash_exit_report(3),
8467 )
8468 .await,
8469 NextAction::Restart { schedule: _ }
8470 ),
8471 "the cleared budget must be spendable again"
8472 );
8473 let state = lock_snapshot(&snapshot).unwrap();
8474 assert_eq!(state.crash_restarts.len(), 1);
8475 assert_eq!(state.lifetime_restarts, 3);
8476 }
8477
8478 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
8479 async fn severance_marker_for_a_dead_child_does_not_label_its_successor() {
8480 let severed = ProcessIdentity {
8481 pid: 41,
8482 start_time: 101,
8483 };
8484 let successor = ProcessIdentity {
8485 pid: 41,
8486 start_time: 202,
8487 };
8488 let module = module_with_recovery_snapshot(ModuleState::Running, true, 0);
8489 update_snapshot(&module.inner.snapshot, Some("recovery-snapshot"), |state| {
8490 state.pid = Some(successor.pid);
8491 state.process_start_time = Some(successor.start_time);
8492 })
8493 .unwrap();
8494 assert!(!module.record_deliberate_severance(severed).unwrap());
8495
8496 let exit_report = apply_deliberate_severance_marker(
8497 &module.inner.snapshot,
8498 Some(successor),
8499 ExitReport {
8500 kind: ExitKind::Crash,
8501 code: Some(1),
8502 signal: None,
8503 at_ms: 1,
8504 },
8505 );
8506
8507 assert_eq!(exit_report.kind, ExitKind::Crash);
8508 }
8509
8510 #[tokio::test]
8511 async fn drain_reap_marks_deliberate_severance_and_records_lifetime_without_budget() {
8512 let registry = Registry::default();
8513 let supervisor = Supervisor::new(
8514 Arc::new(Registry::default()),
8515 RestartPolicy::new(3, Duration::ZERO),
8516 );
8517 let runtime = supervisor.runtime_config();
8518 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8519 let spec = ModuleSpec {
8520 module_id: "drain-deliberate-severance".to_string(),
8521 program: fake_aft_stub_path(),
8522 args: Vec::new(),
8523 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
8524 reserved: false,
8525 reserved_prefixes: Vec::new(),
8526 protocol: ModuleProtocol::Subc,
8527 overlap: Default::default(),
8528 };
8529 let mut child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
8530 let process = ProcessIdentity {
8531 pid: 41,
8532 start_time: 101,
8533 };
8534 child.process_identity = Some(process);
8535 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
8536 state.pid = Some(process.pid);
8537 state.process_start_time = Some(process.start_time);
8538 })
8539 .unwrap();
8540 record_deliberate_severance(&snapshot, process).unwrap();
8541
8542 drain_child_to_state(
8543 &spec.module_id,
8544 spec.protocol,
8545 StopNotice::SentOverConnection,
8548 ®istry,
8549 &snapshot,
8550 &runtime.terminal_ring,
8551 &runtime.spawn_events,
8552 child,
8553 Duration::from_secs(1),
8554 ModuleState::Stopped,
8555 Some(false),
8556 )
8557 .await
8558 .unwrap();
8559
8560 let state = lock_snapshot(&snapshot).unwrap();
8561 assert_eq!(
8562 state.last_exit.as_ref().map(|exit| exit.kind),
8563 Some(ExitKind::DeliberateSeverance)
8564 );
8565 assert_eq!(state.lifetime_restarts, 1);
8566 assert_eq!(state.crash_restarts.len(), 0);
8567 drop(state);
8568 let history = runtime.terminal_ring.lock().unwrap().snapshot();
8569 assert_eq!(
8570 history.entries[0].exit_kind,
8571 subc_control::TerminalExitKind::DeliberateSeverance
8572 );
8573 }
8574
8575 #[tokio::test]
8576 async fn ordinary_drain_reap_does_not_record_a_lifetime_restart() {
8577 let registry = Registry::default();
8578 let supervisor = Supervisor::new(
8579 Arc::new(Registry::default()),
8580 RestartPolicy::new(3, Duration::ZERO),
8581 );
8582 let runtime = supervisor.runtime_config();
8583 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8584 let spec = ModuleSpec {
8585 module_id: "ordinary-drain".to_string(),
8586 program: fake_aft_stub_path(),
8587 args: Vec::new(),
8588 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
8589 reserved: false,
8590 reserved_prefixes: Vec::new(),
8591 protocol: ModuleProtocol::Subc,
8592 overlap: Default::default(),
8593 };
8594 let child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
8595
8596 drain_child_to_state(
8597 &spec.module_id,
8598 spec.protocol,
8599 StopNotice::SentOverConnection,
8602 ®istry,
8603 &snapshot,
8604 &runtime.terminal_ring,
8605 &runtime.spawn_events,
8606 child,
8607 Duration::from_secs(1),
8608 ModuleState::Stopped,
8609 Some(false),
8610 )
8611 .await
8612 .unwrap();
8613
8614 let state = lock_snapshot(&snapshot).unwrap();
8615 assert_eq!(
8616 state.last_exit.as_ref().map(|exit| exit.kind),
8617 Some(ExitKind::Crash)
8618 );
8619 assert_eq!(state.lifetime_restarts, 0);
8620 assert_eq!(state.crash_restarts.len(), 0);
8621 }
8622
8623 #[test]
8624 fn fatal_connection_teardown_cannot_arm_a_marker_for_a_surviving_process() {
8625 assert!(!include_str!("server.rs")
8631 .contains("router.record_deliberate_connection_severance(ctx.connection_id)"));
8632 }
8633
8634 #[test]
8641 fn drained_after_quiescence_wait_passes_ok_through_and_forces_false_on_err() {
8642 assert!(drained_after_quiescence_wait(&Ok(true)));
8643 assert!(!drained_after_quiescence_wait(&Ok(false)));
8644 assert!(!drained_after_quiescence_wait(&Err(
8645 SuperviseError::StatePoisoned { module_id: None }
8646 )));
8647 }
8648
8649 #[test]
8658 fn wait_error_exit_report_records_a_failed_terminal_with_no_code_or_signal() {
8659 let ring = Arc::new(Mutex::new(TerminalRing::new(
8660 TerminalRingConfig::default(),
8661 0,
8662 )));
8663 record_wait_error_terminal("wait-error", &ring, &super::SpawnEventFeed::default());
8664
8665 let snapshot = ring.lock().unwrap().snapshot();
8666 assert_eq!(snapshot.entries.len(), 1);
8667 let entry = &snapshot.entries[0];
8668 assert_eq!(entry.exit_code, None);
8669 assert_eq!(entry.exit_signal, None);
8670 assert_eq!(entry.disposition, TerminalDisposition::Failed);
8671 }
8672
8673 #[test]
8674 fn wait_error_exit_path_preserves_spawn_event_density() {
8675 let feed = super::SpawnEventFeed::default();
8676 feed.configure_incarnation("wait-error-density".to_string());
8677 feed.emit_spawned("wait-error", 41, 1);
8678 let ring = Arc::new(Mutex::new(TerminalRing::new(
8679 TerminalRingConfig::default(),
8680 0,
8681 )));
8682
8683 record_wait_error_terminal("wait-error", &ring, &feed);
8684 feed.emit_spawned("after-wait-error", 42, 2);
8685
8686 let state = feed.0.lock().unwrap();
8687 let sequences = state
8688 .events
8689 .iter()
8690 .map(|event| event.cursor.seq)
8691 .collect::<Vec<_>>();
8692 assert_eq!(sequences, vec![1, 2, 3]);
8693 assert_eq!(state.events[1].kind, SpawnEventKind::Exited);
8694 assert_eq!(state.events[1].exit_code, None);
8695 assert_eq!(state.events[1].exit_signal, None);
8696 }
8697
8698 #[test]
8702 fn wait_error_exit_report_is_classified_as_a_crash() {
8703 assert_eq!(wait_error_exit_report().kind, ExitKind::Crash);
8704 }
8705}
8706
8707#[cfg(test)]
8708mod health_evidence_tests {
8709 use super::{HealthProbeError, HealthProbeEvidence};
8710 use std::collections::HashSet;
8711
8712 #[test]
8720 fn only_a_dead_lane_is_proof_of_death() {
8721 assert!(HealthProbeError::lane_dead("gone").is_proof_of_death());
8722 assert!(!HealthProbeError::no_answer("timed out").is_proof_of_death());
8726 assert!(!HealthProbeError::bad_answer("garbage").is_proof_of_death());
8727 assert!(!HealthProbeError::misconfigured("no table").is_proof_of_death());
8728 }
8729
8730 #[test]
8736 fn every_evidence_class_has_a_distinct_label() {
8737 let labels = [
8738 HealthProbeError::lane_dead("").label(),
8739 HealthProbeError::no_answer("").label(),
8740 HealthProbeError::bad_answer("").label(),
8741 HealthProbeError::misconfigured("").label(),
8742 ];
8743 let unique: HashSet<_> = labels.iter().collect();
8744 assert_eq!(unique.len(), labels.len(), "labels collided: {labels:?}");
8745 }
8746
8747 #[test]
8753 fn classification_preserves_the_original_message() {
8754 let err = HealthProbeError::no_answer("module did not answer within 5s");
8755 assert_eq!(err.to_string(), "module did not answer within 5s");
8756 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8757 }
8758}
8759
8760#[cfg(test)]
8761mod health_tombstone_tests {
8762 use std::{path::PathBuf, sync::Arc, time::Duration};
8763
8764 use subc_protocol::{
8765 manifest::Concurrency,
8766 session::{HealthStatus, ModuleControlResponse},
8767 };
8768 use tokio::sync::mpsc;
8769
8770 use super::{
8771 probe_module_health, HealthAction, HealthConfig, HealthProbeEvidence, ModuleProtocol,
8772 ModuleSpec, RestartPolicy, Supervisor, SupervisorRuntimeConfig,
8773 };
8774 use crate::{
8775 control::ControlHandler,
8776 forwarding::{ForwardingTable, ModuleControlRpcCompletion, ModuleControlRpcOutcome},
8777 registry::{ConnectionId, Registry},
8778 router::FrameSink,
8779 };
8780
8781 struct ProbeHarness {
8782 spec: ModuleSpec,
8783 runtime: SupervisorRuntimeConfig,
8784 forwarding: Arc<ForwardingTable>,
8785 module_connection: ConnectionId,
8786 module_rx: mpsc::Receiver<crate::router::OutboundFrame>,
8787 handler: ControlHandler,
8788 module: super::SupervisedModule,
8789 }
8790
8791 fn probe_harness() -> ProbeHarness {
8792 let registry = Arc::new(Registry::default());
8793 let forwarding = Arc::new(ForwardingTable::default());
8794 let supervisor_handle = super::SupervisorHandle::new();
8795 let health = HealthConfig {
8796 cadence: Duration::from_secs(30),
8797 deadline: Duration::from_secs(5),
8798 failure_threshold: 3,
8799 on_degraded: HealthAction::Report,
8800 on_failing: HealthAction::Report,
8801 critical: false,
8802 };
8803 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
8804 .with_forwarding(Arc::clone(&forwarding))
8805 .with_handle(supervisor_handle.clone())
8806 .with_health_config(health);
8807 let spec = ModuleSpec {
8808 module_id: "late-health-module".to_string(),
8809 program: PathBuf::from("disabled-module"),
8810 args: Vec::new(),
8811 env: Vec::new(),
8812 reserved: false,
8813 reserved_prefixes: Vec::new(),
8814 protocol: ModuleProtocol::Subc,
8815 overlap: Default::default(),
8816 };
8817 let module = supervisor
8818 .supervise_configured(spec.clone(), false)
8819 .unwrap();
8820 let runtime = supervisor.runtime_config();
8821 let handler = ControlHandler::with_forwarding(registry, Arc::clone(&forwarding))
8822 .with_supervisor(supervisor_handle);
8823 let module_connection = ConnectionId::new(700);
8824 let (module_tx, module_rx) = mpsc::channel(8);
8825 forwarding
8826 .register_module_connection(
8827 module_connection,
8828 spec.module_id.clone(),
8829 subc_protocol::PROTOCOL_VERSION,
8830 Concurrency::ModuleManaged,
8831 FrameSink::new(module_tx),
8832 )
8833 .unwrap();
8834
8835 ProbeHarness {
8836 spec,
8837 runtime,
8838 forwarding,
8839 module_connection,
8840 module_rx,
8841 handler,
8842 module,
8843 }
8844 }
8845
8846 async fn finish_after(
8847 harness: &mut ProbeHarness,
8848 stall: Duration,
8849 ) -> ModuleControlRpcCompletion {
8850 assert!(stall > harness.runtime.health.deadline);
8851 let deadline = harness.runtime.health.deadline;
8852 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8853 let answer = async {
8854 let frame = harness.module_rx.recv().await.expect("health.check frame");
8855 tokio::time::advance(deadline).await;
8856 tokio::task::yield_now().await;
8857 tokio::time::advance(stall - deadline).await;
8858 harness
8859 .forwarding
8860 .complete_module_control_rpc(
8861 harness.module_connection,
8862 frame.header.corr,
8863 Some("health.check"),
8864 ModuleControlRpcOutcome::Response(ModuleControlResponse::HealthCheck {
8865 status: HealthStatus::Ok,
8866 detail: None,
8867 metrics: None,
8868 }),
8869 )
8870 .unwrap()
8871 };
8872 let (probe_result, completion) = tokio::join!(probe, answer);
8873 let err = probe_result.expect_err("probe must miss its deadline");
8874 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8875 completion
8876 }
8877
8878 async fn time_out_without_answer(harness: &mut ProbeHarness) {
8879 let deadline = harness.runtime.health.deadline;
8880 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8881 let exhaust_deadline = async {
8882 let _frame = harness.module_rx.recv().await.expect("health.check frame");
8883 tokio::time::advance(deadline).await;
8884 tokio::task::yield_now().await;
8885 };
8886 let (probe_result, ()) = tokio::join!(probe, exhaust_deadline);
8887 let err = probe_result.expect_err("probe must miss its deadline");
8888 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8889 }
8890
8891 #[tokio::test(start_paused = true)]
8892 async fn late_health_answers_record_start_anchored_latency_for_two_stalls() {
8893 let mut harness = probe_harness();
8894
8895 let first = finish_after(&mut harness, Duration::from_secs(8)).await;
8896 let first_latency = match &first {
8897 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8898 other => panic!("late answer was not retained: {other:?}"),
8899 };
8900 assert!(harness.handler.observe_module_control_completion(first));
8901
8902 let second = finish_after(&mut harness, Duration::from_secs(11)).await;
8903 let second_latency = match &second {
8904 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8905 other => panic!("late answer was not retained: {other:?}"),
8906 };
8907 assert!(harness.handler.observe_module_control_completion(second));
8908
8909 assert_eq!(first_latency, Duration::from_secs(8));
8910 assert_eq!(
8911 second_latency - first_latency,
8912 Duration::from_secs(3),
8913 "latency must grow linearly with the additional stall"
8914 );
8915 let health = harness.module.status().unwrap().health;
8916 assert_eq!(health.late_answer_count, 2);
8917 assert_eq!(health.last_late_answer_latency_ms, Some(11_000));
8918 }
8919
8920 #[tokio::test(start_paused = true)]
8928 async fn late_answer_clears_the_consecutive_failure_streak() {
8929 let mut harness = probe_harness();
8930
8931 time_out_without_answer(&mut harness).await;
8933 harness
8934 .module
8935 .record_health_probe_failure_for_test("[no-answer] test miss")
8936 .unwrap();
8937 assert_eq!(
8938 harness.module.status().unwrap().health.consecutive_failures,
8939 1,
8940 "precondition: the miss must be on the streak before the late answer"
8941 );
8942
8943 let late = finish_after(&mut harness, Duration::from_secs(9)).await;
8945 assert!(matches!(
8946 late,
8947 ModuleControlRpcCompletion::LateHealthAnswer { .. }
8948 ));
8949 assert!(harness.handler.observe_module_control_completion(late));
8950
8951 let health = harness.module.status().unwrap().health;
8952 assert_eq!(
8953 health.consecutive_failures, 0,
8954 "a late answer is an answer: the streak must reset"
8955 );
8956 assert_eq!(health.late_answer_count, 1);
8957 }
8958
8959 #[tokio::test(start_paused = true)]
8960 async fn repeated_serial_probe_cycles_keep_one_tombstone_per_endpoint() {
8961 let mut harness = probe_harness();
8962
8963 for _ in 0..20 {
8964 time_out_without_answer(&mut harness).await;
8965 assert_eq!(
8966 harness.forwarding.health_probe_tombstone_count().unwrap(),
8967 1
8968 );
8969 }
8970 }
8971}
8972
8973#[cfg(test)]
8974mod child_env_tests {
8975 use super::{
8976 apply_child_env, apply_spawn_role, apply_wire_spawn_args, ModuleProtocol, ModuleSpec,
8977 SpawnRole, SupervisorHandle, SPAWN_ROLE_SWAP_CANDIDATE, SUBC_ARG, SUBC_LAUNCH_NONCE_ENV,
8978 SUBC_MODULE_ID_ENV, SUBC_SPAWN_ROLE_ENV,
8979 };
8980 use std::{ffi::OsStr, path::PathBuf};
8981 use tokio::process::Command;
8982
8983 fn spec(env: Vec<(String, String)>) -> ModuleSpec {
8984 ModuleSpec {
8985 module_id: "env-plan".to_string(),
8986 program: PathBuf::from("/nonexistent"),
8987 args: Vec::new(),
8988 env,
8989 reserved: false,
8990 reserved_prefixes: Vec::new(),
8991 protocol: ModuleProtocol::Subc,
8992 overlap: Default::default(),
8993 }
8994 }
8995
8996 #[test]
9010 fn ambient_ck_log_is_removed_and_a_configured_one_survives() {
9011 let mut command = Command::new("/nonexistent");
9012 apply_child_env(&mut command, &spec(Vec::new()));
9013 let removed = command
9014 .as_std()
9015 .get_envs()
9016 .any(|(key, value)| key == OsStr::new("CK_LOG") && value.is_none());
9017 assert!(
9018 removed,
9019 "ambient CK_LOG must be explicitly removed for an unconfigured module"
9020 );
9021
9022 let mut configured = Command::new("/nonexistent");
9023 apply_child_env(
9024 &mut configured,
9025 &spec(vec![("CK_LOG".to_string(), "debug".to_string())]),
9026 );
9027 let effective = configured
9028 .as_std()
9029 .get_envs()
9030 .filter(|(key, _)| *key == OsStr::new("CK_LOG"))
9031 .last()
9032 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()));
9033 assert_eq!(
9034 effective,
9035 Some(Some("debug".to_string())),
9036 "a module's configured CK_LOG must survive the ambient removal"
9037 );
9038 }
9039
9040 #[test]
9049 fn protocol_none_spawn_carries_no_subc_argument_and_no_nonce() {
9050 let connection_file = std::path::Path::new("/run/subc-connection.json");
9051 let handle = SupervisorHandle::new();
9052
9053 let mut none_spec = spec(Vec::new());
9054 none_spec.protocol = ModuleProtocol::None;
9055 let mut none = Command::new("/nonexistent");
9056 apply_wire_spawn_args(&mut none, &none_spec, Some(connection_file), Some(&handle))
9057 .expect("protocol-none spawn args apply");
9058 let none_args: Vec<String> = none
9059 .as_std()
9060 .get_args()
9061 .map(|a| a.to_string_lossy().into_owned())
9062 .collect();
9063 assert!(
9064 !none_args.iter().any(|a| a == SUBC_ARG),
9065 "protocol:none argv must not carry --subc; got {none_args:?}"
9066 );
9067 let none_has_nonce = none
9068 .as_std()
9069 .get_envs()
9070 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some());
9071 assert!(
9072 !none_has_nonce,
9073 "protocol:none spawn must not receive a launch nonce"
9074 );
9075 let none_has_module_id = none
9076 .as_std()
9077 .get_envs()
9078 .any(|(key, value)| key == OsStr::new(SUBC_MODULE_ID_ENV) && value.is_some());
9079 assert!(
9080 none_has_module_id,
9081 "SUBC_MODULE_ID is inert and stays on every path"
9082 );
9083 assert!(
9084 handle.spawn_nonce(&none_spec.module_id).is_none(),
9085 "no nonce record for a process that will never present one"
9086 );
9087
9088 let wire_spec = spec(Vec::new());
9090 let mut wire = Command::new("/nonexistent");
9091 apply_wire_spawn_args(&mut wire, &wire_spec, Some(connection_file), Some(&handle))
9092 .expect("subc-wire spawn args apply");
9093 let wire_args: Vec<String> = wire
9094 .as_std()
9095 .get_args()
9096 .map(|a| a.to_string_lossy().into_owned())
9097 .collect();
9098 assert_eq!(
9099 wire_args,
9100 vec![
9101 SUBC_ARG.to_string(),
9102 connection_file.to_string_lossy().into_owned()
9103 ],
9104 "a subc-wire spawn still carries --subc <path>"
9105 );
9106 assert!(wire
9107 .as_std()
9108 .get_envs()
9109 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some()));
9110 assert!(handle.spawn_nonce(&wire_spec.module_id).is_some());
9111 }
9112
9113 #[test]
9123 fn plain_spawn_removes_the_spawn_role_even_when_the_spec_sets_it() {
9124 let role = |command: &Command| {
9125 command
9126 .as_std()
9127 .get_envs()
9128 .filter(|(key, _)| *key == OsStr::new(SUBC_SPAWN_ROLE_ENV))
9129 .last()
9130 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()))
9131 };
9132 let forged = spec(vec![(
9133 SUBC_SPAWN_ROLE_ENV.to_string(),
9134 SPAWN_ROLE_SWAP_CANDIDATE.to_string(),
9135 )]);
9136
9137 let mut plain = Command::new("/nonexistent");
9138 apply_child_env(&mut plain, &forged);
9139 apply_spawn_role(&mut plain, SpawnRole::Plain);
9140 assert_eq!(
9141 role(&plain),
9142 Some(None),
9143 "a plain spawn must remove SUBC_SPAWN_ROLE, whatever the spec says"
9144 );
9145
9146 let mut candidate = Command::new("/nonexistent");
9147 apply_child_env(&mut candidate, &spec(Vec::new()));
9148 apply_spawn_role(&mut candidate, SpawnRole::SwapCandidate);
9149 assert_eq!(
9150 role(&candidate),
9151 Some(Some(SPAWN_ROLE_SWAP_CANDIDATE.to_string()))
9152 );
9153 }
9154
9155 #[test]
9161 fn daemon_private_capture_keys_are_not_passed_to_the_child() {
9162 let mut command = Command::new("/nonexistent");
9163 apply_child_env(
9164 &mut command,
9165 &spec(vec![
9166 (super::CAPTURE_KEEP_ENV.to_string(), "5".to_string()),
9167 ("KEPT".to_string(), "yes".to_string()),
9168 ]),
9169 );
9170 let keys: Vec<String> = command
9171 .as_std()
9172 .get_envs()
9173 .filter(|(_, value)| value.is_some())
9174 .map(|(key, _)| key.to_string_lossy().into_owned())
9175 .collect();
9176 assert!(keys.contains(&"KEPT".to_string()), "got {keys:?}");
9177 assert!(
9178 !keys.contains(&super::CAPTURE_KEEP_ENV.to_string()),
9179 "daemon-private capture key leaked to the child: {keys:?}"
9180 );
9181 }
9182}
9183
9184#[cfg(test)]
9185mod jitter_tests {
9186 use super::jittered_health_delay;
9187 use std::{collections::HashSet, time::Duration};
9188
9189 const FLEET: [&str; 14] = [
9198 "aft",
9199 "alfonso-core",
9200 "magic-context",
9201 "broca",
9202 "thalamus",
9203 "quota",
9204 "engram",
9205 "plexus",
9206 "cerebellum",
9207 "astrocyte",
9208 "synapse",
9209 "subc-mcp",
9210 "cortexkit-credentials",
9211 "subc-federation",
9212 ];
9213
9214 #[test]
9222 fn probe_delays_disperse_across_the_fleet() {
9223 let cadence = Duration::from_secs(30);
9224 let delays: HashSet<Duration> = FLEET
9225 .iter()
9226 .map(|id| jittered_health_delay(id, 0, cadence))
9227 .collect();
9228 assert_eq!(
9229 delays.len(),
9230 FLEET.len(),
9231 "every supervised module must land on its own probe offset"
9232 );
9233 }
9234
9235 #[test]
9241 fn jitter_only_delays_and_stays_within_one_tenth_of_cadence() {
9242 let cadence = Duration::from_secs(30);
9243 let span = cadence / 10;
9244 for id in FLEET {
9245 for probe_index in 0..8 {
9246 let delay = jittered_health_delay(id, probe_index, cadence);
9247 assert!(
9248 delay >= cadence,
9249 "{id}#{probe_index}: jitter must not shorten the cadence"
9250 );
9251 assert!(
9252 delay < cadence + span,
9253 "{id}#{probe_index}: jitter must stay inside one tenth of the cadence"
9254 );
9255 }
9256 }
9257 }
9258
9259 #[test]
9265 fn a_module_offset_is_stable_across_restarts() {
9266 let cadence = Duration::from_secs(30);
9267 for id in FLEET {
9268 assert_eq!(
9269 jittered_health_delay(id, 0, cadence),
9270 jittered_health_delay(id, 0, cadence),
9271 "{id}: the same module and probe index must produce the same offset"
9272 );
9273 }
9274 }
9275
9276 #[test]
9278 fn zero_cadence_yields_zero_delay() {
9279 assert_eq!(
9280 jittered_health_delay("aft", 0, Duration::ZERO),
9281 Duration::ZERO
9282 );
9283 }
9284}
9285
9286#[cfg(all(test, target_os = "linux"))]
9287mod cgroup_placement_tests {
9288 use super::{
9289 apply_cgroup_placement, remove_module_cgroup, ModuleProtocol, ModuleSpec, SuperviseError,
9290 SupervisedChild,
9291 };
9292 use crate::stderr_tail::{StderrRing, StderrTailConfig};
9293 use std::{
9294 fs, io,
9295 path::{Path, PathBuf},
9296 sync::{Arc, Mutex},
9297 };
9298 use subc_test_support::TestTempDir;
9299 use tokio::process::Command;
9300
9301 #[test]
9302 fn failed_parent_cgroup_open_is_a_cgroup_supervision_error() {
9303 let path = Path::new("/definitely-missing-subc-cgroup");
9304 let mut command = Command::new("true");
9305 let error = apply_cgroup_placement(
9306 &mut command,
9307 &ModuleSpec {
9308 module_id: "broken-cgroup".to_string(),
9309 program: PathBuf::from("true"),
9310 args: Vec::new(),
9311 env: Vec::new(),
9312 reserved: false,
9313 reserved_prefixes: Vec::new(),
9314 protocol: ModuleProtocol::Subc,
9315 overlap: Default::default(),
9316 },
9317 path,
9318 )
9319 .expect_err("a parent cgroup open failure must reject the supervised spawn");
9320 let reason = error.to_string();
9321
9322 assert!(
9323 matches!(error, SuperviseError::Cgroup { .. }),
9324 "parent cgroup open must be reported as a cgroup supervision error: {reason}"
9325 );
9326 assert!(
9327 reason.contains("/definitely-missing-subc-cgroup/cgroup.procs"),
9328 "parent cgroup open failure must name cgroup.procs: {reason}"
9329 );
9330 }
9331
9332 #[tokio::test]
9333 async fn reaping_a_child_removes_its_empty_module_cgroup() {
9334 let root = TestTempDir::new("supervisor-reap-cgroup");
9335 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
9336 let placement = subc_cgroup::prepare_at(&root)
9337 .expect("prepare scratch cgroup root")
9338 .expect("scratch root has a cgroup.procs marker");
9339 let module_id = "reaped-module";
9340 let module = placement
9341 .module_path(module_id)
9342 .expect("create scratch module cgroup");
9343 let child = Command::new("true")
9344 .spawn()
9345 .expect("spawn short-lived child");
9346 let pid = child.id().expect("spawned child has pid");
9347 let mut child = SupervisedChild {
9348 child,
9349 module_id: module_id.to_string(),
9350 cgroup_placement: Some(placement),
9351 stdout_pump: None,
9352 stderr_pump: None,
9353 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
9354 spawned_at_ms: 0,
9355 spawned_from: PathBuf::from("true"),
9356 spawned_file_identity: None,
9357 process_start_time: None,
9358 process_identity: None,
9359 pid,
9360 roster_guard: None,
9361 };
9362
9363 child.wait().await.expect("reap short-lived child");
9364
9365 assert!(
9366 !module.exists(),
9367 "reaping the supervised child must remove its empty cgroup"
9368 );
9369 }
9370
9371 #[test]
9372 fn non_empty_cgroup_removal_is_reported_without_blocking_teardown() {
9373 let root = TestTempDir::new("supervisor-non-empty-cgroup");
9374 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
9375 let placement = subc_cgroup::prepare_at(&root)
9376 .expect("prepare scratch cgroup root")
9377 .expect("scratch root has a cgroup.procs marker");
9378 let module = placement
9379 .module_path("surviving-module")
9380 .expect("create scratch module cgroup");
9381 fs::write(module.join("surviving-process"), b"still present")
9382 .expect("make scratch cgroup non-empty");
9383 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::WARN);
9384
9385 remove_module_cgroup(&placement, "surviving-module");
9386
9387 let logs = crate::router::test_log::captured_logs(&logs);
9388 assert!(
9389 module.exists(),
9390 "failed removal must leave the cgroup intact"
9391 );
9392 assert!(
9393 logs.contains("could not remove module cgroup after process exit; continuing teardown")
9394 && logs.contains("surviving-module"),
9395 "best-effort removal must report the failure without returning it: {logs}"
9396 );
9397 }
9398
9399 #[test]
9400 fn cgroup_pre_exec_spawn_failure_names_the_cgroup_path() {
9401 let cgroup_path = PathBuf::from("/sys/fs/cgroup/subc-modules/broken-module");
9402 let reason = SuperviseError::Spawn {
9403 program: PathBuf::from("/bin/true"),
9404 source: io::Error::from_raw_os_error(13),
9405 cgroup_path: Some(cgroup_path.clone()),
9406 }
9407 .to_string();
9408
9409 assert!(
9410 reason.contains(&cgroup_path.display().to_string()),
9411 "a pre_exec spawn failure must name the cgroup path: {reason}"
9412 );
9413 }
9414}
9415
9416#[cfg(test)]
9417mod spawn_subscriber_lag_tests {
9418 use super::*;
9419
9420 #[tokio::test]
9425 async fn lagged_spawn_subscriber_receives_a_terminal_lagged_error_after_its_queued_frames() {
9426 let feed = SpawnEventFeed::default();
9427 feed.configure_incarnation("lag-incarnation".to_string());
9428 let (tx, mut rx) = mpsc::channel(1);
9431 feed.subscribe(ConnectionId::new(1), 7, 1, None, FrameSink::new(tx))
9432 .expect("subscribe");
9433 let emitted = SPAWN_SUBSCRIBER_BUFFER + 16;
9434 for index in 0..emitted {
9435 feed.emit_spawned(&format!("lag-module-{index}"), 1000, 0);
9436 tokio::task::yield_now().await;
9439 }
9440 assert_eq!(
9441 feed.subscriber_count(),
9442 0,
9443 "the lagged subscriber must be removed"
9444 );
9445
9446 let mut data = Vec::new();
9447 let mut last = None;
9448 loop {
9449 let next = tokio::time::timeout(Duration::from_secs(5), rx.recv())
9450 .await
9451 .expect("the forwarder must finish once the subscriber is dropped");
9452 let Some(outbound) = next else { break };
9453 let frame = outbound.frame;
9454 if frame.header.ty == FrameType::StreamData {
9455 assert!(last.is_none(), "no data may follow the terminal frame");
9456 let event: SpawnEvent = serde_json::from_slice(&frame.body).unwrap();
9457 data.push(event.cursor.seq);
9458 } else {
9459 assert!(last.is_none(), "exactly one terminal frame");
9460 last = Some(frame);
9461 }
9462 }
9463 assert!(!data.is_empty(), "queued frames drain before the terminal");
9464 for pair in data.windows(2) {
9465 assert_eq!(
9466 pair[1],
9467 pair[0] + 1,
9468 "queued frames arrive dense and in order"
9469 );
9470 }
9471 let terminal = last.expect("a lagged subscriber must receive a terminal frame");
9472 assert_eq!(terminal.header.ty, FrameType::Error);
9473 assert_eq!(terminal.header.corr, 7);
9474 let body: subc_protocol::ErrorBody = serde_json::from_slice(&terminal.body).unwrap();
9475 assert_eq!(body.code, SPAWN_SUBSCRIBER_LAGGED_CODE);
9476 let detail = body.detail.expect("lagged error carries detail");
9477 assert_eq!(
9478 detail["first_undelivered_cursor"]["seq"],
9479 data.last().unwrap() + 1,
9480 "the named cursor is the first event the subscriber did not receive"
9481 );
9482 assert_eq!(
9483 detail["first_undelivered_cursor"]["daemon_incarnation"],
9484 "lag-incarnation"
9485 );
9486 }
9487}
9488
9489#[cfg(test)]
9490mod terminal_history_read_concurrency_tests {
9491 use super::*;
9492 use crate::terminal_journal::read_pause;
9493 use std::sync::mpsc as std_mpsc;
9494 use subc_test_support::TestTempDir;
9495
9496 fn journaled_ring(
9497 journal: &Arc<crate::terminal_journal::TerminalJournal>,
9498 ) -> Arc<Mutex<TerminalRing>> {
9499 Arc::new(Mutex::new(
9500 TerminalRing::new(TerminalRingConfig::default(), 1)
9501 .with_journal(Some(Arc::clone(journal))),
9502 ))
9503 }
9504
9505 fn crash(at_ms: u64) -> ExitReport {
9506 ExitReport {
9507 kind: ExitKind::Crash,
9508 code: Some(1),
9509 signal: None,
9510 at_ms,
9511 }
9512 }
9513
9514 fn record_within(
9517 module_id: &'static str,
9518 ring: &Arc<Mutex<TerminalRing>>,
9519 at_ms: u64,
9520 bound: Duration,
9521 ) -> bool {
9522 let ring = Arc::clone(ring);
9523 let (done, done_rx) = std_mpsc::channel();
9524 std::thread::spawn(move || {
9525 record_terminal(
9526 module_id,
9527 &ring,
9528 &SpawnEventFeed::default(),
9529 &crash(at_ms),
9530 TerminalDisposition::Restarting,
9531 );
9532 let _ = done.send(());
9533 });
9534 done_rx.recv_timeout(bound).is_ok()
9535 }
9536
9537 #[test]
9542 fn exits_recorded_during_a_paused_history_read_are_not_blocked_or_half_merged() {
9543 let dir = TestTempDir::new("terminal-history-concurrent-read");
9544 let path = dir.join("terminals.jsonl");
9545 let journal = Arc::new(crate::terminal_journal::TerminalJournal::open(
9546 path.clone(),
9547 "daemon".into(),
9548 ));
9549 let reader_ring = journaled_ring(&journal);
9550 let other_ring = journaled_ring(&journal);
9551 assert!(record_within(
9552 "reader-module",
9553 &reader_ring,
9554 10,
9555 Duration::from_secs(5)
9556 ));
9557
9558 let (started, release) = read_pause::install(&path);
9559 let reading = {
9560 let ring = Arc::clone(&reader_ring);
9561 std::thread::spawn(move || durable_terminal_history_of(&ring, "reader-module"))
9562 };
9563 started
9564 .recv_timeout(Duration::from_secs(5))
9565 .expect("the history read reached its pause");
9566
9567 let bound = Duration::from_secs(1);
9568 assert!(
9569 record_within("other-module", &other_ring, 20, bound),
9570 "another module's exit waited on a history read (journal writer held)"
9571 );
9572 assert!(
9573 record_within("reader-module", &reader_ring, 30, bound),
9574 "the read module's own exit waited on its history read (ring held)"
9575 );
9576
9577 drop(release);
9578 let paused = reading.join().unwrap();
9579 assert_eq!(
9580 paused.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9581 vec![10],
9582 "an exit recorded after the read began lands in neither half of it"
9583 );
9584 assert_eq!(paused.journal_skipped_lines, 0);
9585 assert_eq!(paused.journal_read_errors, 0);
9586
9587 let after = durable_terminal_history_of(&reader_ring, "reader-module");
9588 assert_eq!(
9589 after.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9590 vec![10, 30],
9591 "the next read merges ring and journal with no duplicate"
9592 );
9593 assert_eq!(after.journal_skipped_lines, 0);
9594 }
9595}
9596
9597#[cfg(test)]
9602mod stderr_settle_tests {
9603 use std::{
9604 future::Future,
9605 io,
9606 pin::Pin,
9607 sync::{Arc, Mutex},
9608 task::{Context, Poll},
9609 time::Duration,
9610 };
9611
9612 use tokio::{
9613 io::{AsyncRead, ReadBuf},
9614 sync::oneshot,
9615 time::Instant,
9616 };
9617
9618 use super::{settle_stderr_pump, StderrPump};
9619 use crate::stderr_tail::{
9620 pump_stderr_to, CaptureState, OutputSink, StderrRing, StderrTailConfig, TailEntry,
9621 };
9622
9623 const BOUND: Duration = Duration::from_millis(250);
9624
9625 struct HeldReader {
9629 before: Option<Vec<u8>>,
9630 gate: Option<oneshot::Receiver<()>>,
9631 after: io::Cursor<Vec<u8>>,
9632 }
9633
9634 impl AsyncRead for HeldReader {
9635 fn poll_read(
9636 mut self: Pin<&mut Self>,
9637 cx: &mut Context<'_>,
9638 buf: &mut ReadBuf<'_>,
9639 ) -> Poll<io::Result<()>> {
9640 if let Some(bytes) = self.before.take() {
9641 buf.put_slice(&bytes);
9642 return Poll::Ready(Ok(()));
9643 }
9644 if let Some(gate) = self.gate.as_mut() {
9645 match Pin::new(gate).poll(cx) {
9646 Poll::Pending => return Poll::Pending,
9647 Poll::Ready(_) => self.gate = None,
9648 }
9649 }
9650 Pin::new(&mut self.after).poll_read(cx, buf)
9651 }
9652 }
9653
9654 struct DiscardSink;
9655
9656 impl OutputSink for DiscardSink {
9657 fn write_line(&mut self, _line: &[u8]) {}
9658 }
9659
9660 fn line(text: &str) -> TailEntry {
9661 TailEntry::Line {
9662 text: text.to_string(),
9663 truncated: false,
9664 }
9665 }
9666
9667 fn lock(ring: &Arc<Mutex<StderrRing>>) -> std::sync::MutexGuard<'_, StderrRing> {
9668 ring.lock().unwrap()
9669 }
9670
9671 fn held_pump(
9675 ring: &Arc<Mutex<StderrRing>>,
9676 before: &str,
9677 after: &str,
9678 ) -> (StderrPump, oneshot::Sender<()>) {
9679 let generation = lock(ring).begin_process();
9680 let (release, gate) = oneshot::channel();
9681 let reader = HeldReader {
9682 before: Some(before.as_bytes().to_vec()),
9683 gate: Some(gate),
9684 after: io::Cursor::new(after.as_bytes().to_vec()),
9685 };
9686 let task = tokio::spawn(pump_stderr_to(
9687 reader,
9688 Arc::clone(ring),
9689 generation,
9690 DiscardSink,
9691 ));
9692 (StderrPump { task, generation }, release)
9693 }
9694
9695 async fn wait_until(ring: &Arc<Mutex<StderrRing>>, done: impl Fn(&StderrRing) -> bool) {
9696 for _ in 0..1000 {
9697 if done(&lock(ring)) {
9698 return;
9699 }
9700 tokio::time::sleep(Duration::from_millis(1)).await;
9701 }
9702 panic!(
9703 "ring never reached the expected state: {:?}",
9704 lock(ring).snapshot(None, None)
9705 );
9706 }
9707
9708 #[tokio::test(start_paused = true)]
9709 async fn a_crash_line_the_reader_had_not_reached_by_the_bound_is_kept_before_the_restart() {
9710 let ring = Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default())));
9711 let (pump, release) = held_pump(&ring, "booting\n", "config error: missing storage\n");
9712
9713 settle_stderr_pump("crasher", &ring, pump, BOUND).await;
9714 let before_release = lock(&ring).snapshot(None, None);
9715 assert!(
9716 matches!(before_release.capture, CaptureState::Incomplete { .. }),
9717 "a reader that has not reached EOF cannot claim a whole tail: {before_release:?}"
9718 );
9719
9720 let next = lock(&ring).begin_process();
9723 lock(&ring).push_line_from(next, "next process booting");
9724 release.send(()).unwrap();
9725 wait_until(&ring, |ring| {
9726 ring.snapshot(None, None).capture == CaptureState::Captured
9727 })
9728 .await;
9729
9730 assert_eq!(
9731 lock(&ring).snapshot(None, None).entries,
9732 vec![
9733 line("booting"),
9734 line("config error: missing storage"),
9735 TailEntry::ProcessStart,
9736 line("next process booting"),
9737 ],
9738 "the crash's last line must survive a slow reader and stay in the crashed process's section"
9739 );
9740 }
9741
9742 #[tokio::test(start_paused = true)]
9743 async fn a_pipe_held_open_by_a_descendant_reads_incomplete_without_delaying_the_restart_past_the_bound(
9744 ) {
9745 let ring = Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default())));
9746 let (pump, _held) = held_pump(&ring, "parent exiting\n", "");
9749
9750 let started = Instant::now();
9751 settle_stderr_pump("orphaning", &ring, pump, BOUND).await;
9752 assert_eq!(
9753 started.elapsed(),
9754 BOUND,
9755 "the restart must wait exactly the bound for a pipe that stays open, no longer"
9756 );
9757
9758 let next = lock(&ring).begin_process();
9759 lock(&ring).push_line_from(next, "next process booting");
9760 tokio::time::sleep(Duration::from_secs(60)).await;
9761
9762 let snapshot = lock(&ring).snapshot(None, None);
9763 match &snapshot.capture {
9764 CaptureState::Incomplete { reason } => assert!(
9765 reason.contains("had not reached EOF") && reason.contains("250ms"),
9766 "the reason must say what is missing and after how long: {reason}"
9767 ),
9768 other => panic!("expected Incomplete while the pipe is held open, got {other:?}"),
9769 }
9770 assert_eq!(
9771 snapshot.entries,
9772 vec![
9773 line("parent exiting"),
9774 TailEntry::ProcessStart,
9775 line("next process booting"),
9776 ]
9777 );
9778 }
9779
9780 #[tokio::test(start_paused = true)]
9781 async fn a_reader_that_reaches_eof_within_the_bound_leaves_the_tail_captured() {
9782 let ring = Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default())));
9783 let (pump, release) = held_pump(&ring, "one\n", "two\n");
9784 release.send(()).unwrap();
9785
9786 settle_stderr_pump("clean", &ring, pump, BOUND).await;
9787
9788 let snapshot = lock(&ring).snapshot(None, None);
9789 assert_eq!(snapshot.capture, CaptureState::Captured);
9790 assert_eq!(snapshot.entries, vec![line("one"), line("two")]);
9791 }
9792}
9793
9794#[cfg(all(test, windows))]
9808mod job_containment_tests {
9809 use super::*;
9810 use std::{
9811 path::{Path, PathBuf},
9812 sync::{Arc, Mutex},
9813 time::{Duration, Instant},
9814 };
9815 use subc_test_support::TestTempDir;
9816
9817 fn stub_path() -> PathBuf {
9823 let mut path = std::env::current_exe().expect("current_exe available in tests");
9824 path.pop();
9825 path.pop();
9826 path.push("fake-aft-stub.exe");
9827 assert!(
9828 path.exists(),
9829 "fake-aft-stub not built at {}: run `cargo test -p subc-core` (which builds \
9830 [[bin]] targets) rather than `cargo test -p subc-core --lib` (which does not)",
9831 path.display()
9832 );
9833 path
9834 }
9835
9836 fn read_grandchild_pid(path: &Path) -> u32 {
9838 let deadline = Instant::now() + Duration::from_secs(10);
9839 loop {
9840 if let Ok(contents) = std::fs::read_to_string(path) {
9841 if let Ok(pid) = contents.trim().parse() {
9842 return pid;
9843 }
9844 }
9845 assert!(
9846 Instant::now() < deadline,
9847 "the stub never recorded a grandchild pid at {}",
9848 path.display()
9849 );
9850 std::thread::sleep(Duration::from_millis(10));
9851 }
9852 }
9853
9854 struct Fixture {
9857 _dir: TestTempDir,
9858 module_id: String,
9859 grandchild: u32,
9860 child: Option<SupervisedChild>,
9861 registry: Arc<Registry>,
9862 snapshot: Arc<Mutex<SupervisorSnapshot>>,
9863 terminal_ring: Arc<Mutex<TerminalRing>>,
9864 spawn_events: SpawnEventFeed,
9865 }
9866
9867 fn fixture(label: &str, module_id: &str) -> Fixture {
9868 let dir = TestTempDir::new(label);
9869 let pid_file = dir.join("grandchild.pid");
9870 let supervisor = Supervisor::new(
9871 Arc::new(Registry::default()),
9872 RestartPolicy::new(3, Duration::ZERO),
9873 );
9874 let runtime = supervisor.runtime_config();
9875 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
9876 let spec = ModuleSpec {
9877 module_id: module_id.to_string(),
9878 program: stub_path(),
9879 args: Vec::new(),
9883 env: vec![
9884 ("FAKE_AFT_NEVER_CONNECT".to_string(), "1".to_string()),
9885 (
9886 "FAKE_AFT_GRANDCHILD_PID_FILE".to_string(),
9887 pid_file.display().to_string(),
9888 ),
9889 ],
9890 reserved: false,
9891 reserved_prefixes: Vec::new(),
9892 protocol: ModuleProtocol::Subc,
9893 overlap: Default::default(),
9894 };
9895 let child = spawn_and_mark_running(&spec, &runtime, &snapshot)
9896 .expect("spawn the supervised fixture");
9897 let grandchild = read_grandchild_pid(&pid_file);
9898 Fixture {
9899 _dir: dir,
9900 module_id: module_id.to_string(),
9901 grandchild,
9902 child: Some(child),
9903 registry: Arc::new(Registry::default()),
9904 snapshot,
9905 terminal_ring: Arc::clone(&runtime.terminal_ring),
9906 spawn_events: SpawnEventFeed::default(),
9907 }
9908 }
9909
9910 impl Fixture {
9911 async fn drain(&mut self) {
9913 let child = self
9914 .child
9915 .take()
9916 .expect("the fixture child is still present");
9917 drain_child_to_state(
9918 &self.module_id,
9919 ModuleProtocol::Subc,
9920 StopNotice::NotSent,
9923 &self.registry,
9924 &self.snapshot,
9925 &self.terminal_ring,
9926 &self.spawn_events,
9927 child,
9928 Duration::from_millis(500),
9929 ModuleState::Stopped,
9930 Some(false),
9931 )
9932 .await
9933 .expect("drain the supervised fixture");
9934 }
9935 }
9936
9937 #[tokio::test]
9943 async fn teardown_reaps_the_grandchild() {
9944 let mut fixture = fixture("teardown-grandchild", "tree-teardown");
9945 let grandchild = fixture.grandchild;
9946
9947 assert!(
9948 subc_jobobject::process_exists(grandchild),
9949 "grandchild {grandchild} must be alive before teardown, or this proves nothing"
9950 );
9951
9952 fixture.drain().await;
9953
9954 assert!(
9955 subc_jobobject::wait_for_process_exit(grandchild, Duration::from_secs(10)),
9956 "grandchild {grandchild} outlived module teardown: the tree was not contained"
9957 );
9958 }
9959
9960 #[test]
9973 fn an_uncontained_grandchild_survives_a_direct_child_kill() {
9974 let dir = TestTempDir::new("teardown-uncontained");
9975 let pid_file = dir.join("grandchild.pid");
9976 let mut child = std::process::Command::new(stub_path())
9977 .env("FAKE_AFT_NEVER_CONNECT", "1")
9978 .env(
9979 "FAKE_AFT_GRANDCHILD_PID_FILE",
9980 pid_file.display().to_string(),
9981 )
9982 .stdin(std::process::Stdio::null())
9983 .stdout(std::process::Stdio::null())
9984 .stderr(std::process::Stdio::null())
9985 .spawn()
9986 .expect("spawn the uncontained fixture");
9987 let grandchild = read_grandchild_pid(&pid_file);
9988
9989 child.kill().expect("kill the direct child");
9991 let _ = child.wait();
9992
9993 assert!(
9994 subc_jobobject::process_exists(grandchild),
9995 "grandchild {grandchild} died with the direct child, so this control no longer \
9996 distinguishes contained from uncontained teardown and the regression test is \
9997 passing vacuously"
9998 );
9999
10000 kill_tree(grandchild);
10003 }
10004
10005 #[tokio::test]
10014 async fn dropping_containment_reaps_the_grandchild() {
10015 let mut fixture = fixture("drop-containment", "tree-drop");
10016 let grandchild = fixture.grandchild;
10017
10018 assert!(subc_jobobject::process_exists(grandchild));
10019
10020 fixture.child.as_mut().expect("child present").job = None;
10022
10023 assert!(
10024 subc_jobobject::wait_for_process_exit(grandchild, Duration::from_secs(10)),
10025 "grandchild {grandchild} survived the containment handle closing, so a daemon \
10026 crash would leave the tree behind"
10027 );
10028 }
10029
10030 fn kill_tree(pid: u32) {
10032 let _ = std::process::Command::new("taskkill.exe")
10033 .args(["/PID", &pid.to_string(), "/T", "/F"])
10034 .stdin(std::process::Stdio::null())
10035 .stdout(std::process::Stdio::null())
10036 .stderr(std::process::Stdio::null())
10037 .status();
10038 assert!(
10039 subc_jobobject::wait_for_process_exit(pid, Duration::from_secs(10)),
10040 "could not clean up grandchild {pid}"
10041 );
10042 }
10043}