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 stdout_pump: Option<JoinHandle<()>>,
122 stderr_pump: Option<StderrPump>,
123 stderr_ring: Arc<Mutex<StderrRing>>,
124 spawned_at_ms: u64,
125 spawned_from: PathBuf,
126 spawned_file_identity: Option<SpawnedFileIdentity>,
127 process_start_time: Option<u64>,
128 process_identity: Option<ProcessIdentity>,
129 pid: u32,
130 roster_guard: Option<crate::child_roster::RosterGuard>,
133}
134
135impl SupervisedChild {
136 fn id(&self) -> Option<u32> {
137 Some(self.pid)
138 }
139
140 fn process_identity(&self) -> Option<ProcessIdentity> {
141 self.process_identity
142 }
143
144 async fn wait(&mut self) -> io::Result<ExitStatus> {
145 let result = self.child.wait().await;
153 #[cfg(target_os = "linux")]
154 if result.is_ok() {
155 if let Some(placement) = self.cgroup_placement.take() {
156 remove_module_cgroup(&placement, &self.module_id);
157 }
158 }
159 result
160 }
161
162 fn release_roster(&mut self) {
166 self.roster_guard = None;
167 }
168
169 fn start_kill(&mut self) -> io::Result<()> {
170 self.child.start_kill()
171 }
172
173 async fn drain_stderr(&mut self, module_id: &str) {
174 if let Some(mut pump) = self.stdout_pump.take() {
175 match timeout(STDERR_PUMP_DRAIN_TIMEOUT, &mut pump).await {
176 Ok(Ok(())) => {}
177 Ok(Err(error)) => {
178 warn!(module_id, error = %error, "stdout pump ended unexpectedly");
179 }
180 Err(_) => {
181 pump.abort();
182 warn!(
183 module_id,
184 waited = ?STDERR_PUMP_DRAIN_TIMEOUT,
185 "stdout pump did not drain before restart; stopped it before the next process"
186 );
187 }
188 }
189 }
190
191 let Some(pump) = self.stderr_pump.take() else {
192 return;
193 };
194 settle_stderr_pump(
195 module_id,
196 &self.stderr_ring,
197 pump,
198 STDERR_PUMP_DRAIN_TIMEOUT,
199 )
200 .await;
201 }
202}
203
204struct StderrPump {
207 task: JoinHandle<()>,
208 generation: u64,
209}
210
211async fn settle_stderr_pump(
217 module_id: &str,
218 ring: &Arc<Mutex<StderrRing>>,
219 pump: StderrPump,
220 bound: Duration,
221) {
222 let lock = || ring.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
223 let StderrPump {
224 mut task,
225 generation,
226 } = pump;
227 lock().retire_pump(generation);
228 match timeout(bound, &mut task).await {
229 Ok(Ok(())) => {}
230 Ok(Err(err)) => {
231 let mut ring = lock();
232 ring.mark_incomplete(format!("stderr pump ended unexpectedly: {err}"));
233 ring.finish_pump(generation);
234 warn!(module_id, error = %err, "stderr pump ended before clean EOF");
235 }
236 Err(_) => {
237 drop(task);
239 lock().mark_pump_late(
240 generation,
241 format!(
242 "stderr of the exited process had not reached EOF {bound:?} after it was \
243 retired (a descendant may still hold the pipe open); lines it still \
244 writes are kept in that process's section"
245 ),
246 );
247 warn!(
248 module_id,
249 waited = ?bound,
250 "stderr pipe of the exited process is still open; its reader keeps running without delaying the restart"
251 );
252 }
253 }
254}
255
256fn registration_release_events() -> &'static watch::Sender<u64> {
257 static EVENTS: OnceLock<watch::Sender<u64>> = OnceLock::new();
258 EVENTS.get_or_init(|| {
259 let (sender, _receiver) = watch::channel(0);
260 sender
261 })
262}
263
264pub(crate) fn notify_registration_release() {
265 let events = registration_release_events();
266 let next_generation = (*events.borrow()).wrapping_add(1);
267 events.send_replace(next_generation);
268}
269
270#[derive(Debug, Clone, PartialEq, Eq)]
272pub struct ModuleSpec {
273 pub module_id: String,
274 pub program: PathBuf,
275 pub args: Vec<String>,
276 pub env: Vec<(String, String)>,
277 pub reserved: bool,
282 pub reserved_prefixes: Vec<String>,
287 pub protocol: ModuleProtocol,
303 pub overlap: ModuleOverlap,
308}
309
310#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
317pub enum ModuleOverlap {
318 #[default]
320 Exclusive,
321 Safe,
333}
334
335impl ModuleOverlap {
336 pub fn as_str(self) -> &'static str {
337 match self {
338 Self::Exclusive => "exclusive",
339 Self::Safe => "safe",
340 }
341 }
342}
343
344pub const SUBC_SPAWN_ROLE_ENV: &str = "SUBC_SPAWN_ROLE";
354pub const SPAWN_ROLE_SWAP_CANDIDATE: &str = "swap_candidate";
356pub const DEFAULT_SWAP_READY_TIMEOUT: Duration = Duration::from_secs(100);
361
362#[derive(Debug, Clone, Copy, PartialEq, Eq)]
380pub struct RestartPolicy {
381 pub max_restarts: u32,
382 pub backoff: Duration,
385 pub max_backoff: Duration,
387 pub window: Duration,
391}
392
393impl RestartPolicy {
394 pub fn new(max_restarts: u32, backoff: Duration) -> Self {
398 Self {
399 max_restarts,
400 backoff,
401 max_backoff: DEFAULT_MAX_BACKOFF,
402 window: DEFAULT_RESTART_WINDOW,
403 }
404 }
405
406 pub fn with_max_backoff(mut self, max_backoff: Duration) -> Self {
407 self.max_backoff = max_backoff;
408 self
409 }
410
411 pub fn with_window(mut self, window: Duration) -> Self {
412 self.window = window;
413 self
414 }
415
416 fn delay_for_restart(&self, restart_in_window: u32) -> Duration {
421 if self.backoff.is_zero() || self.max_backoff.is_zero() {
422 return Duration::ZERO;
423 }
424
425 let mut delay = self.backoff;
426 for _ in 0..restart_in_window {
427 if delay >= self.max_backoff {
428 return self.max_backoff;
429 }
430 delay = delay
431 .checked_mul(10)
432 .unwrap_or(self.max_backoff)
433 .min(self.max_backoff);
434 }
435 delay.min(self.max_backoff)
436 }
437
438 fn budget_exhausted_detail(&self) -> String {
443 format!(
444 "crash budget exhausted: max_restarts={} within window_secs={}",
445 self.max_restarts,
446 self.window.as_secs()
447 )
448 }
449}
450
451impl Default for RestartPolicy {
452 fn default() -> Self {
453 Self {
454 max_restarts: DEFAULT_MAX_RESTARTS,
455 backoff: DEFAULT_BACKOFF,
456 max_backoff: DEFAULT_MAX_BACKOFF,
457 window: DEFAULT_RESTART_WINDOW,
458 }
459 }
460}
461
462#[derive(Debug, Clone, Copy, PartialEq, Eq)]
463struct CrashRestartSchedule {
464 restart_in_window: u32,
465 delay: Duration,
466}
467
468fn daemon_will_restart(
475 state: &mut SupervisorSnapshot,
476 policy: &RestartPolicy,
477 now: Instant,
478) -> bool {
479 state.enabled && state.crash_restarts_in_window(policy.window, now) < policy.max_restarts
480}
481
482const DEFAULT_HEALTH_CADENCE: Duration = Duration::from_secs(30);
483const DEFAULT_HEALTH_DEADLINE: Duration = Duration::from_secs(5);
484const DEFAULT_HEALTH_FAILURE_THRESHOLD: u32 = 3;
485const MAX_HEALTH_METRICS_BYTES: usize = 16 * 1024;
486
487#[derive(Debug, Clone, Copy, PartialEq, Eq)]
488pub enum HealthAction {
489 Report,
490 Restart,
491 Alert,
492}
493
494impl fmt::Display for HealthAction {
495 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
496 f.write_str(match self {
497 Self::Report => "report",
498 Self::Restart => "restart",
499 Self::Alert => "alert",
500 })
501 }
502}
503
504#[derive(Debug, Clone, Copy, PartialEq, Eq)]
505pub struct HealthConfig {
506 pub cadence: Duration,
507 pub deadline: Duration,
508 pub failure_threshold: u32,
509 pub on_degraded: HealthAction,
510 pub on_failing: HealthAction,
511 pub critical: bool,
512}
513
514impl Default for HealthConfig {
515 fn default() -> Self {
516 Self {
517 cadence: DEFAULT_HEALTH_CADENCE,
518 deadline: DEFAULT_HEALTH_DEADLINE,
519 failure_threshold: DEFAULT_HEALTH_FAILURE_THRESHOLD,
520 on_degraded: HealthAction::Report,
521 on_failing: HealthAction::Report,
522 critical: false,
523 }
524 }
525}
526
527#[derive(Debug, Clone, PartialEq)]
545pub struct ModuleHealthStatus {
546 pub status: SupervisorHealthStatus,
547 pub last_probe_ms: Option<u64>,
548 pub detail: Option<String>,
549 pub metrics: Option<Value>,
550 pub consecutive_failures: u32,
551 pub late_answer_count: u64,
554 pub last_late_answer_latency_ms: Option<u64>,
556 pub last_action: Option<String>,
557 pub last_action_ms: Option<u64>,
561}
562
563impl Default for ModuleHealthStatus {
564 fn default() -> Self {
565 Self {
566 status: SupervisorHealthStatus::Unknown,
567 last_probe_ms: None,
568 detail: None,
569 metrics: None,
570 consecutive_failures: 0,
571 late_answer_count: 0,
572 last_late_answer_latency_ms: None,
573 last_action: None,
574 last_action_ms: None,
575 }
576 }
577}
578
579#[derive(Debug, Clone, Copy, PartialEq, Eq)]
581pub enum ModuleState {
582 Starting,
583 Running,
584 Unresponsive,
585 Restarting,
586 Draining,
587 Stopped,
588 Failed,
589 Disabled,
590}
591
592impl fmt::Display for ModuleState {
593 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
594 f.write_str(match self {
595 Self::Starting => "starting",
596 Self::Running => "running",
597 Self::Unresponsive => "unresponsive",
598 Self::Restarting => "restarting",
599 Self::Draining => "draining",
600 Self::Stopped => "stopped",
601 Self::Failed => "failed",
602 Self::Disabled => "disabled",
603 })
604 }
605}
606
607#[derive(Debug, Clone, Copy, PartialEq, Eq)]
609pub enum ExitKind {
610 Clean,
611 Crash,
612 DeliberateSeverance,
613}
614
615impl From<ExitKind> for TerminalExitKind {
616 fn from(kind: ExitKind) -> Self {
617 match kind {
618 ExitKind::Clean => Self::Clean,
619 ExitKind::Crash => Self::Crash,
620 ExitKind::DeliberateSeverance => Self::DeliberateSeverance,
621 }
622 }
623}
624
625#[derive(Debug, Clone, Copy, PartialEq, Eq)]
628pub(crate) struct ProcessIdentity {
629 pub(crate) pid: u32,
630 pub(crate) start_time: u64,
631}
632
633#[derive(Debug, Clone, PartialEq, Eq)]
635pub struct ExitReport {
636 pub kind: ExitKind,
637 pub code: Option<i32>,
638 pub signal: Option<i32>,
639 pub at_ms: u64,
640}
641
642#[derive(Debug, Clone, PartialEq)]
645pub struct ModuleStatus {
646 pub module_id: String,
647 pub state: ModuleState,
648 pub enabled: bool,
649 pub process_alive: bool,
650 pub registration_active: bool,
651 pub protocol: ModuleProtocol,
654 pub live: bool,
665 pub restart_count: u32,
669 pub lifetime_restarts: u32,
673 pub spawn_generation: u64,
674 pub max_restarts: u32,
679 pub restart_window: Duration,
683 pub drain_timeout: Duration,
687 pub restart_backoff: Duration,
688 pub restart_max_backoff: Duration,
689 pub pid: Option<u32>,
690 pub spawned_at_ms: Option<u64>,
691 pub spawned_from: Option<PathBuf>,
692 pub process_start_time: Option<u64>,
693 pub last_exit: Option<ExitReport>,
694 pub health: ModuleHealthStatus,
695}
696
697#[derive(Debug, Clone, PartialEq)]
698struct SupervisorSnapshot {
699 state: ModuleState,
700 enabled: bool,
701 process_alive: bool,
702 crash_restarts: VecDeque<Instant>,
708 lifetime_restarts: u32,
709 spawn_generation: u64,
718 pid: Option<u32>,
719 spawned_at_ms: Option<u64>,
720 spawned_from: Option<PathBuf>,
721 spawned_file_identity: Option<SpawnedFileIdentity>,
722 process_start_time: Option<u64>,
723 deliberate_severance: Option<ProcessIdentity>,
724 last_exit: Option<ExitReport>,
725 health: ModuleHealthStatus,
726 in_alternate_slot: bool,
731 draining_to_replace: bool,
738 configuration_updated_since_spawn: bool,
744}
745
746impl SupervisorSnapshot {
747 fn starting() -> Self {
748 Self::new(ModuleState::Starting, true)
749 }
750
751 fn disabled() -> Self {
752 Self::new(ModuleState::Disabled, false)
753 }
754
755 fn failed() -> Self {
756 Self::new(ModuleState::Failed, true)
757 }
758
759 fn crash_restarts_in_window(&mut self, window: Duration, now: Instant) -> u32 {
763 while let Some(oldest) = self.crash_restarts.front() {
764 if now.duration_since(*oldest) > window {
765 self.crash_restarts.pop_front();
766 } else {
767 break;
768 }
769 }
770 u32::try_from(self.crash_restarts.len()).unwrap_or(u32::MAX)
771 }
772
773 fn record_crash_restart(&mut self, policy: &RestartPolicy, now: Instant) {
779 self.crash_restarts.push_back(now);
780 while self.crash_restarts.len() > policy.max_restarts as usize {
781 self.crash_restarts.pop_front();
782 }
783 self.lifetime_restarts += 1;
784 }
785
786 fn next_crash_restart(
790 &mut self,
791 policy: &RestartPolicy,
792 now: Instant,
793 ) -> Option<CrashRestartSchedule> {
794 let restart_in_window = self.crash_restarts_in_window(policy.window, now);
795 if restart_in_window >= policy.max_restarts {
796 return None;
797 }
798 self.record_crash_restart(policy, now);
799 Some(CrashRestartSchedule {
800 restart_in_window,
801 delay: policy.delay_for_restart(restart_in_window),
802 })
803 }
804
805 fn clear_crash_restarts(&mut self) {
810 self.crash_restarts.clear();
811 }
812
813 fn new(state: ModuleState, enabled: bool) -> Self {
814 Self {
815 state,
816 enabled,
817 process_alive: false,
818 crash_restarts: VecDeque::new(),
819 lifetime_restarts: 0,
820 spawn_generation: 0,
821 pid: None,
822 spawned_at_ms: None,
823 spawned_from: None,
824 spawned_file_identity: None,
825 process_start_time: None,
826 deliberate_severance: None,
827 last_exit: None,
828 health: ModuleHealthStatus::default(),
829 in_alternate_slot: false,
830 draining_to_replace: false,
831 configuration_updated_since_spawn: false,
832 }
833 }
834}
835
836type SharedSnapshot = Arc<Mutex<SupervisorSnapshot>>;
837
838type SpawnSubscriberKey = (ConnectionId, u64);
839
840#[derive(Debug)]
841struct SpawnSubscriber {
842 version: u8,
843 frames: mpsc::Sender<Frame>,
844 lagged: Option<oneshot::Sender<SpawnCursor>>,
848}
849
850#[derive(Debug)]
851struct SpawnEventState {
852 daemon_incarnation: String,
853 seq: u64,
854 capacity: usize,
855 live: HashMap<String, LiveSpawn>,
856 generations: HashMap<String, u64>,
857 events: VecDeque<SpawnEvent>,
858 subscribers: HashMap<SpawnSubscriberKey, SpawnSubscriber>,
859}
860
861impl Default for SpawnEventState {
862 fn default() -> Self {
863 Self {
864 daemon_incarnation: "unconfigured".to_string(),
865 seq: 0,
866 capacity: SPAWN_EVENT_RING_CAPACITY,
867 live: HashMap::new(),
868 generations: HashMap::new(),
869 events: VecDeque::new(),
870 subscribers: HashMap::new(),
871 }
872 }
873}
874
875#[derive(Debug, Clone, Default)]
876struct SpawnEventFeed(Arc<Mutex<SpawnEventState>>);
877
878#[derive(Debug, Clone, PartialEq, Eq)]
879pub(crate) enum SpawnSubscribeRefusal {
880 ForeignIncarnation { current: String },
881 TooOld { oldest: SpawnCursor },
882 Frame(String),
883}
884
885impl SpawnEventFeed {
886 fn configure_incarnation(&self, daemon_incarnation: String) {
887 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
888 state.daemon_incarnation = daemon_incarnation;
889 state.seq = 0;
890 state.live.clear();
891 state.generations.clear();
892 state.events.clear();
893 state.subscribers.clear();
894 }
895
896 fn cursor(state: &SpawnEventState) -> SpawnCursor {
897 SpawnCursor {
898 daemon_incarnation: state.daemon_incarnation.clone(),
899 seq: state.seq,
900 }
901 }
902
903 fn snapshot(&self) -> SpawnSnapshot {
904 let state = self.0.lock().unwrap_or_else(|p| p.into_inner());
905 let mut live = state.live.values().cloned().collect::<Vec<_>>();
906 live.sort_by(|left, right| left.module_id.cmp(&right.module_id));
907 SpawnSnapshot {
908 cursor: Self::cursor(&state),
909 ring_bound: state.capacity as u64,
910 live,
911 }
912 }
913
914 fn emit_spawned(&self, module_id: &str, pid: u32, spawned_at_ms: u64) -> u64 {
915 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
916 let generation = state
917 .generations
918 .get(module_id)
919 .copied()
920 .unwrap_or(0)
921 .checked_add(1)
922 .expect("spawn generation exhausted");
923 state.generations.insert(module_id.to_string(), generation);
924 let live = LiveSpawn {
925 module_id: module_id.to_string(),
926 spawn_generation: generation,
927 pid,
928 spawned_at_ms,
929 };
930 state.live.insert(module_id.to_string(), live);
931 Self::emit_locked(
932 &mut state,
933 SpawnEventKind::Spawned,
934 module_id.to_string(),
935 generation,
936 pid,
937 None,
938 None,
939 );
940 generation
941 }
942
943 fn emit_exited(&self, module_id: &str, exit_code: Option<i32>, exit_signal: Option<i32>) {
944 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
945 let Some(live) = state.live.remove(module_id) else {
946 warn!(
947 module_id,
948 "terminal record had no live spawn event identity"
949 );
950 return;
951 };
952 Self::emit_locked(
953 &mut state,
954 SpawnEventKind::Exited,
955 module_id.to_string(),
956 live.spawn_generation,
957 live.pid,
958 exit_code,
959 exit_signal,
960 );
961 }
962
963 fn emit_superseded_exited(
970 &self,
971 module_id: &str,
972 spawn_generation: u64,
973 pid: u32,
974 exit_code: Option<i32>,
975 exit_signal: Option<i32>,
976 ) {
977 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
978 if state
979 .live
980 .get(module_id)
981 .is_some_and(|live| live.spawn_generation == spawn_generation)
982 {
983 state.live.remove(module_id);
984 }
985 Self::emit_locked(
986 &mut state,
987 SpawnEventKind::Exited,
988 module_id.to_string(),
989 spawn_generation,
990 pid,
991 exit_code,
992 exit_signal,
993 );
994 }
995
996 #[allow(clippy::too_many_arguments)]
997 fn emit_locked(
998 state: &mut SpawnEventState,
999 kind: SpawnEventKind,
1000 module_id: String,
1001 spawn_generation: u64,
1002 pid: u32,
1003 exit_code: Option<i32>,
1004 exit_signal: Option<i32>,
1005 ) {
1006 state.seq = state
1007 .seq
1008 .checked_add(1)
1009 .expect("spawn event sequence exhausted");
1010 let event = SpawnEvent {
1011 cursor: Self::cursor(state),
1012 kind,
1013 module_id,
1014 spawn_generation,
1015 pid,
1016 exit_code,
1017 exit_signal,
1018 };
1019 state.events.push_back(event.clone());
1020 while state.events.len() > state.capacity {
1021 state.events.pop_front();
1022 }
1023 let body = match serde_json::to_vec(&event) {
1024 Ok(body) => body,
1025 Err(error) => {
1026 error!(%error, "failed to serialize supervisor spawn event");
1027 return;
1028 }
1029 };
1030 state.subscribers.retain(|(connection_id, corr), subscriber| {
1031 let frame = Frame::build_with_version(
1032 subscriber.version,
1033 FrameType::StreamData,
1034 control_flags(),
1035 0,
1036 0,
1037 *corr,
1038 body.clone(),
1039 );
1040 match frame {
1041 Ok(frame) => {
1042 if subscriber.frames.try_send(frame).is_ok() {
1043 true
1044 } else {
1045 warn!(connection_id = connection_id.get(), corr, "dropping lagged supervisor spawn subscriber");
1046 if let Some(lagged) = subscriber.lagged.take() {
1047 let _ = lagged.send(event.cursor.clone());
1048 }
1049 false
1050 }
1051 }
1052 Err(error) => {
1053 warn!(connection_id = connection_id.get(), corr, %error, "dropping supervisor spawn subscriber after frame build failure");
1054 false
1055 }
1056 }
1057 });
1058 }
1059
1060 fn subscribe(
1061 &self,
1062 connection_id: ConnectionId,
1063 corr: u64,
1064 version: u8,
1065 since: Option<SpawnCursor>,
1066 sink: FrameSink,
1067 ) -> Result<(), SpawnSubscribeRefusal> {
1068 let (frames, mut receiver) = mpsc::channel(SPAWN_SUBSCRIBER_BUFFER);
1069 let (lagged, mut lagged_rx) = oneshot::channel::<SpawnCursor>();
1070 {
1071 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
1072 let replay = if let Some(since) = since {
1073 if since.daemon_incarnation != state.daemon_incarnation {
1074 return Err(SpawnSubscribeRefusal::ForeignIncarnation {
1075 current: state.daemon_incarnation.clone(),
1076 });
1077 }
1078 if let Some(oldest) = state.events.front().map(|event| event.cursor.clone()) {
1079 if since.seq < oldest.seq.saturating_sub(1) {
1080 return Err(SpawnSubscribeRefusal::TooOld { oldest });
1081 }
1082 }
1083 state
1084 .events
1085 .iter()
1086 .filter(|event| event.cursor.seq > since.seq)
1087 .cloned()
1088 .collect::<Vec<_>>()
1089 } else {
1090 Vec::new()
1091 };
1092 for event in replay {
1093 let body = serde_json::to_vec(&event)
1094 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1095 let frame = Frame::build_with_version(
1096 version,
1097 FrameType::StreamData,
1098 control_flags(),
1099 0,
1100 0,
1101 corr,
1102 body,
1103 )
1104 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1105 frames
1106 .try_send(frame)
1107 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1108 }
1109 state.subscribers.insert(
1110 (connection_id, corr),
1111 SpawnSubscriber {
1112 version,
1113 frames,
1114 lagged: Some(lagged),
1115 },
1116 );
1117 }
1118 tokio::spawn(async move {
1129 while let Some(frame) = receiver.recv().await {
1130 if sink.send(frame).await.is_err() {
1131 return;
1132 }
1133 }
1134 let Ok(first_undelivered) = lagged_rx.try_recv() else {
1135 return;
1136 };
1137 match spawn_subscriber_lagged_frame(version, corr, first_undelivered) {
1138 Ok(frame) => {
1139 let _ = sink.send(frame).await;
1140 }
1141 Err(error) => {
1142 error!(%error, corr, "failed to build lagged spawn subscriber terminal frame");
1143 }
1144 }
1145 });
1146 Ok(())
1147 }
1148
1149 fn cancel(&self, connection_id: ConnectionId, corr: u64) -> bool {
1150 let Some(subscriber) = self
1151 .0
1152 .lock()
1153 .unwrap_or_else(|p| p.into_inner())
1154 .subscribers
1155 .remove(&(connection_id, corr))
1156 else {
1157 return false;
1158 };
1159 if let Ok(frame) = Frame::build_with_version(
1160 subscriber.version,
1161 FrameType::StreamEnd,
1162 control_flags(),
1163 0,
1164 0,
1165 corr,
1166 Vec::new(),
1167 ) {
1168 tokio::spawn(async move {
1169 let _ = subscriber.frames.send(frame).await;
1170 });
1171 }
1172 true
1173 }
1174
1175 fn remove_connection(&self, connection_id: ConnectionId) {
1176 self.0
1177 .lock()
1178 .unwrap_or_else(|p| p.into_inner())
1179 .subscribers
1180 .retain(|(subscriber_connection, _), _| *subscriber_connection != connection_id);
1181 }
1182
1183 #[cfg(any(test, feature = "test-support"))]
1184 fn set_capacity(&self, capacity: usize) {
1185 self.0.lock().unwrap_or_else(|p| p.into_inner()).capacity = capacity;
1186 }
1187
1188 #[cfg(any(test, feature = "test-support"))]
1189 fn subscriber_count(&self) -> usize {
1190 self.0
1191 .lock()
1192 .unwrap_or_else(|p| p.into_inner())
1193 .subscribers
1194 .len()
1195 }
1196}
1197
1198fn spawn_subscriber_lagged_frame(
1201 version: u8,
1202 corr: u64,
1203 first_undelivered: SpawnCursor,
1204) -> Result<Frame, String> {
1205 let body = serde_json::to_vec(&subc_protocol::ErrorBody {
1206 code: SPAWN_SUBSCRIBER_LAGGED_CODE.to_string(),
1207 message: "spawn subscriber fell behind and was dropped; resubscribe from the last cursor received"
1208 .to_string(),
1209 detail: Some(serde_json::json!({
1210 "first_undelivered_cursor": first_undelivered
1211 })),
1212 })
1213 .map_err(|error| error.to_string())?;
1214 Frame::build_with_version(version, FrameType::Error, control_flags(), 0, 0, corr, body)
1215 .map_err(|error| error.to_string())
1216}
1217
1218pub trait ModuleProcessLiveness: Send + Sync {
1219 fn process_live(&self, module_id: &str) -> Option<bool>;
1220
1221 fn process_replacing(&self, _module_id: &str) -> bool {
1227 false
1228 }
1229}
1230
1231#[derive(Debug, Clone, Default)]
1233pub struct SupervisorProcessLiveness {
1234 snapshots: Arc<Mutex<HashMap<String, SharedSnapshot>>>,
1235}
1236
1237impl SupervisorProcessLiveness {
1238 pub fn new() -> Self {
1239 Self::default()
1240 }
1241
1242 fn track(&self, module_id: String, snapshot: SharedSnapshot) {
1243 let mut snapshots = self
1244 .snapshots
1245 .lock()
1246 .unwrap_or_else(|poisoned| poisoned.into_inner());
1247 snapshots.insert(module_id, snapshot);
1248 }
1249
1250 fn untrack_if_current(&self, module_id: &str, snapshot: &SharedSnapshot) {
1251 let mut snapshots = self
1252 .snapshots
1253 .lock()
1254 .unwrap_or_else(|poisoned| poisoned.into_inner());
1255 let is_current = snapshots
1256 .get(module_id)
1257 .map(|tracked| Arc::ptr_eq(tracked, snapshot))
1258 .unwrap_or(false);
1259 if is_current {
1260 snapshots.remove(module_id);
1261 }
1262 }
1263}
1264
1265impl ModuleProcessLiveness for SupervisorProcessLiveness {
1266 fn process_live(&self, module_id: &str) -> Option<bool> {
1267 let snapshot = {
1268 let snapshots = self
1269 .snapshots
1270 .lock()
1271 .unwrap_or_else(|poisoned| poisoned.into_inner());
1272 snapshots.get(module_id).cloned()
1273 }?;
1274 let snapshot = snapshot
1275 .lock()
1276 .unwrap_or_else(|poisoned| poisoned.into_inner());
1277 Some(snapshot.state == ModuleState::Running && snapshot.process_alive)
1278 }
1279
1280 fn process_replacing(&self, module_id: &str) -> bool {
1281 let Some(snapshot) = self
1282 .snapshots
1283 .lock()
1284 .unwrap_or_else(|poisoned| poisoned.into_inner())
1285 .get(module_id)
1286 .cloned()
1287 else {
1288 return false;
1289 };
1290 let snapshot = snapshot
1291 .lock()
1292 .unwrap_or_else(|poisoned| poisoned.into_inner());
1293 snapshot.enabled
1294 && match snapshot.state {
1295 ModuleState::Restarting => true,
1296 ModuleState::Draining => snapshot.draining_to_replace,
1297 ModuleState::Starting
1298 | ModuleState::Running
1299 | ModuleState::Unresponsive
1300 | ModuleState::Stopped
1301 | ModuleState::Failed
1302 | ModuleState::Disabled => false,
1303 }
1304 }
1305}
1306
1307#[derive(Debug, Clone)]
1308struct SupervisorRuntimeConfig {
1309 restart_policy: RestartPolicy,
1310 drain_timeout: Duration,
1313 effective_drain_timeout: Arc<Mutex<Duration>>,
1316 default_drain_timeout: Duration,
1319 health: HealthConfig,
1320 connection_file_path: Option<PathBuf>,
1321 capture_logs_dir: Option<PathBuf>,
1322 forwarding: Option<Arc<ForwardingTable>>,
1323 supervisor_handle: Option<SupervisorHandle>,
1326 stderr_ring: Arc<Mutex<StderrRing>>,
1333 terminal_ring: Arc<Mutex<TerminalRing>>,
1334 spawn_events: SpawnEventFeed,
1335 child_roster: ChildRoster,
1336 #[cfg(target_os = "linux")]
1337 cgroup_placement: Option<subc_cgroup::Placement>,
1338 #[cfg(test)]
1339 test_seed_stale_facts_before_enable_spawn: bool,
1340}
1341
1342#[derive(Debug, Clone, PartialEq, Eq)]
1343struct SupervisedConfiguration {
1344 spec: ModuleSpec,
1345 health: HealthConfig,
1346}
1347
1348#[derive(Debug, Clone, Default)]
1354pub struct SupervisorHandle {
1355 modules: Arc<Mutex<HashMap<String, SupervisedModule>>>,
1356 spawn_events: SpawnEventFeed,
1357 reserved_nonces: Arc<Mutex<HashMap<String, Option<String>>>>,
1368 removal_tombstones: Arc<Mutex<HashMap<String, u64>>>,
1374 spawn_nonces: Arc<Mutex<HashMap<String, String>>>,
1378 reserved_prefix_owners: Arc<Mutex<HashMap<String, String>>>,
1386 swaps: Arc<Mutex<HashMap<String, OpenSwap>>>,
1392 promotion_observer: PromotionObserverSlot,
1394 operation_lock: Arc<AsyncMutex<()>>,
1398}
1399
1400pub(crate) trait SwapPromotionObserver: Send + Sync {
1409 fn swap_promoted(&self, registration: &crate::registry::ModuleRegistration);
1410}
1411
1412#[derive(Clone, Default)]
1416struct PromotionObserverSlot(Arc<Mutex<Option<std::sync::Weak<dyn SwapPromotionObserver>>>>);
1417
1418impl fmt::Debug for PromotionObserverSlot {
1419 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1420 f.write_str("PromotionObserverSlot")
1421 }
1422}
1423
1424#[derive(Debug, Clone)]
1426struct OpenSwap {
1427 candidate_nonce: String,
1430 incumbent_nonce: Option<String>,
1435 candidate_admitted: bool,
1439}
1440
1441#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1444pub(crate) enum SwapHelloAdmission {
1445 NotSwapping,
1448 Candidate,
1450 Refused,
1453}
1454
1455#[derive(Debug, Clone, PartialEq, Eq)]
1456pub(crate) enum ReservedHelloRejection {
1457 Exact {
1458 module_id: String,
1459 },
1460 Prefix {
1461 prefix: String,
1462 owner_module_id: String,
1463 },
1464}
1465
1466impl SupervisorHandle {
1467 pub fn new() -> Self {
1468 Self::default()
1469 }
1470
1471 pub(crate) fn spawn_snapshot(&self) -> SpawnSnapshot {
1472 self.spawn_events.snapshot()
1473 }
1474
1475 pub(crate) fn subscribe_spawns(
1476 &self,
1477 connection_id: ConnectionId,
1478 corr: u64,
1479 version: u8,
1480 since: Option<SpawnCursor>,
1481 sink: FrameSink,
1482 ) -> Result<(), SpawnSubscribeRefusal> {
1483 self.spawn_events
1484 .subscribe(connection_id, corr, version, since, sink)
1485 }
1486
1487 pub(crate) fn cancel_spawn_subscription(&self, connection_id: ConnectionId, corr: u64) -> bool {
1488 self.spawn_events.cancel(connection_id, corr)
1489 }
1490
1491 pub(crate) fn remove_spawn_subscribers(&self, connection_id: ConnectionId) {
1492 self.spawn_events.remove_connection(connection_id);
1493 }
1494
1495 #[cfg(any(test, feature = "test-support"))]
1496 pub fn set_spawn_event_capacity_for_test(&self, capacity: usize) {
1497 assert!(capacity > 0, "spawn event capacity must be non-zero");
1498 self.spawn_events.set_capacity(capacity);
1499 }
1500
1501 #[cfg(any(test, feature = "test-support"))]
1502 pub fn spawn_subscriber_count_for_test(&self) -> usize {
1503 self.spawn_events.subscriber_count()
1504 }
1505
1506 pub fn set_spawn_nonce(&self, module_id: &str, nonce: String) {
1509 self.spawn_nonces
1510 .lock()
1511 .unwrap_or_else(|poisoned| poisoned.into_inner())
1512 .insert(module_id.to_string(), nonce);
1513 }
1514
1515 pub fn set_reserved_nonce(&self, module_id: &str, nonce: String) {
1518 self.reserved_nonces
1519 .lock()
1520 .unwrap_or_else(|poisoned| poisoned.into_inner())
1521 .insert(module_id.to_string(), Some(nonce));
1522 }
1523
1524 pub fn set_reserved_prefixes(&self, owner_module_id: &str, prefixes: &[String]) {
1526 let mut owners = self
1527 .reserved_prefix_owners
1528 .lock()
1529 .unwrap_or_else(|poisoned| poisoned.into_inner());
1530 owners.retain(|_, owner| owner != owner_module_id);
1531 for prefix in prefixes {
1532 owners.insert(prefix.clone(), owner_module_id.to_string());
1533 }
1534 }
1535
1536 #[cfg(test)]
1538 pub(crate) fn spawn_nonce(&self, module_id: &str) -> Option<String> {
1539 self.spawn_nonces
1540 .lock()
1541 .unwrap_or_else(|poisoned| poisoned.into_inner())
1542 .get(module_id)
1543 .cloned()
1544 }
1545
1546 fn apply_identity_configuration(&self, spec: &ModuleSpec) {
1547 self.set_reserved_prefixes(&spec.module_id, &spec.reserved_prefixes);
1548 let spawn_nonce = self
1549 .spawn_nonces
1550 .lock()
1551 .unwrap_or_else(|poisoned| poisoned.into_inner())
1552 .get(&spec.module_id)
1553 .cloned();
1554 let mut reserved_nonces = self
1555 .reserved_nonces
1556 .lock()
1557 .unwrap_or_else(|poisoned| poisoned.into_inner());
1558 if spec.reserved {
1559 reserved_nonces.insert(spec.module_id.clone(), spawn_nonce);
1564 }
1565 drop(reserved_nonces);
1566 self.removal_tombstones
1570 .lock()
1571 .unwrap_or_else(|poisoned| poisoned.into_inner())
1572 .remove(&spec.module_id);
1573 }
1574
1575 pub fn reserved_hello_authorized(&self, module_id: &str, presented: Option<&str>) -> bool {
1580 self.reserved_hello_rejection(module_id, presented)
1581 .is_none()
1582 }
1583
1584 pub(crate) fn reserved_hello_rejection(
1585 &self,
1586 module_id: &str,
1587 presented: Option<&str>,
1588 ) -> Option<ReservedHelloRejection> {
1589 let nonces = self
1590 .reserved_nonces
1591 .lock()
1592 .unwrap_or_else(|poisoned| poisoned.into_inner());
1593 if let Some(expected) = nonces.get(module_id) {
1594 let authorized = match expected {
1598 Some(expected) => {
1599 presented.is_some_and(|p| constant_time_eq(expected.as_bytes(), p.as_bytes()))
1600 }
1601 None => false,
1602 };
1603 if authorized {
1604 return None;
1605 }
1606 return Some(ReservedHelloRejection::Exact {
1607 module_id: module_id.to_string(),
1608 });
1609 }
1610 drop(nonces);
1611
1612 let matched_prefix = self
1613 .reserved_prefix_owners
1614 .lock()
1615 .unwrap_or_else(|poisoned| poisoned.into_inner())
1616 .iter()
1617 .filter(|(prefix, _)| module_id.starts_with(prefix.as_str()))
1618 .max_by_key(|(prefix, _)| prefix.len())
1619 .map(|(prefix, owner)| (prefix.clone(), owner.clone()));
1620 let (prefix, owner_module_id) = matched_prefix?;
1621
1622 let authorized = presented.is_some_and(|presented| {
1623 self.spawn_nonces
1624 .lock()
1625 .unwrap_or_else(|poisoned| poisoned.into_inner())
1626 .get(&owner_module_id)
1627 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()))
1628 || self.swap_nonce_matches(&owner_module_id, presented)
1631 });
1632 if authorized {
1633 None
1634 } else {
1635 Some(ReservedHelloRejection::Prefix {
1636 prefix,
1637 owner_module_id,
1638 })
1639 }
1640 }
1641
1642 pub fn spawned_consumer_authorized(&self, module_id: &str, presented: &str) -> bool {
1647 if presented.is_empty() {
1648 return false;
1649 }
1650 let nonces = self
1651 .spawn_nonces
1652 .lock()
1653 .unwrap_or_else(|poisoned| poisoned.into_inner());
1654 let current = nonces
1655 .get(module_id)
1656 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()));
1657 drop(nonces);
1658 current || self.swap_nonce_matches(module_id, presented)
1663 }
1664
1665 fn swap_nonce_matches(&self, module_id: &str, presented: &str) -> bool {
1667 let swaps = self
1668 .swaps
1669 .lock()
1670 .unwrap_or_else(|poisoned| poisoned.into_inner());
1671 swaps.get(module_id).is_some_and(|swap| {
1672 constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes())
1673 || swap.incumbent_nonce.as_deref().is_some_and(|incumbent| {
1674 constant_time_eq(incumbent.as_bytes(), presented.as_bytes())
1675 })
1676 })
1677 }
1678
1679 pub(crate) fn open_swap(&self, module_id: &str, candidate_nonce: String) {
1682 let incumbent_nonce = self
1683 .spawn_nonces
1684 .lock()
1685 .unwrap_or_else(|poisoned| poisoned.into_inner())
1686 .get(module_id)
1687 .cloned();
1688 self.swaps
1689 .lock()
1690 .unwrap_or_else(|poisoned| poisoned.into_inner())
1691 .insert(
1692 module_id.to_string(),
1693 OpenSwap {
1694 candidate_nonce,
1695 incumbent_nonce,
1696 candidate_admitted: false,
1697 },
1698 );
1699 }
1700
1701 pub(crate) fn close_swap(&self, module_id: &str) {
1704 self.swaps
1705 .lock()
1706 .unwrap_or_else(|poisoned| poisoned.into_inner())
1707 .remove(module_id);
1708 }
1709
1710 pub(crate) fn set_swap_promotion_observer(
1713 &self,
1714 observer: std::sync::Weak<dyn SwapPromotionObserver>,
1715 ) {
1716 *self
1717 .promotion_observer
1718 .0
1719 .lock()
1720 .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(observer);
1721 }
1722
1723 fn notify_swap_promoted(&self, registration: &crate::registry::ModuleRegistration) {
1726 let observer = self
1727 .promotion_observer
1728 .0
1729 .lock()
1730 .unwrap_or_else(|poisoned| poisoned.into_inner())
1731 .as_ref()
1732 .and_then(std::sync::Weak::upgrade);
1733 if let Some(observer) = observer {
1734 observer.swap_promoted(registration);
1735 }
1736 }
1737
1738 pub(crate) fn swap_open(&self, module_id: &str) -> bool {
1740 self.swaps
1741 .lock()
1742 .unwrap_or_else(|poisoned| poisoned.into_inner())
1743 .contains_key(module_id)
1744 }
1745
1746 fn promote_swap_nonce(&self, module_id: &str, reserved: bool) {
1751 let candidate_nonce = self
1752 .swaps
1753 .lock()
1754 .unwrap_or_else(|poisoned| poisoned.into_inner())
1755 .get(module_id)
1756 .map(|swap| swap.candidate_nonce.clone());
1757 let Some(nonce) = candidate_nonce else {
1758 return;
1759 };
1760 self.set_spawn_nonce(module_id, nonce.clone());
1761 if reserved {
1762 self.set_reserved_nonce(module_id, nonce);
1763 }
1764 }
1765
1766 pub(crate) fn swap_hello_admission(
1781 &self,
1782 module_id: &str,
1783 presented: Option<&str>,
1784 ) -> SwapHelloAdmission {
1785 let swaps = self
1786 .swaps
1787 .lock()
1788 .unwrap_or_else(|poisoned| poisoned.into_inner());
1789 let Some(swap) = swaps.get(module_id) else {
1790 return SwapHelloAdmission::NotSwapping;
1791 };
1792 let Some(presented) = presented else {
1793 return SwapHelloAdmission::Refused;
1794 };
1795 if constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes()) {
1796 return if swap.candidate_admitted {
1797 SwapHelloAdmission::Refused
1798 } else {
1799 SwapHelloAdmission::Candidate
1800 };
1801 }
1802 if swap
1803 .incumbent_nonce
1804 .as_deref()
1805 .is_some_and(|incumbent| constant_time_eq(incumbent.as_bytes(), presented.as_bytes()))
1806 {
1807 return SwapHelloAdmission::NotSwapping;
1808 }
1809 SwapHelloAdmission::Refused
1810 }
1811
1812 pub(crate) fn mark_swap_candidate_admitted(&self, module_id: &str) {
1815 if let Some(swap) = self
1816 .swaps
1817 .lock()
1818 .unwrap_or_else(|poisoned| poisoned.into_inner())
1819 .get_mut(module_id)
1820 {
1821 swap.candidate_admitted = true;
1822 }
1823 }
1824
1825 pub fn spawn_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1827 self.spawn_nonces
1828 .lock()
1829 .unwrap_or_else(|poisoned| poisoned.into_inner())
1830 .get(module_id)
1831 .cloned()
1832 }
1833
1834 pub fn reserved_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1836 self.reserved_nonces
1837 .lock()
1838 .unwrap_or_else(|poisoned| poisoned.into_inner())
1839 .get(module_id)
1840 .cloned()
1841 .flatten()
1842 }
1843
1844 pub fn insert(&self, module: SupervisedModule) -> Option<SupervisedModule> {
1845 let mut modules = self
1846 .modules
1847 .lock()
1848 .unwrap_or_else(|poisoned| poisoned.into_inner());
1849 modules.insert(module.module_id().to_string(), module)
1850 }
1851
1852 pub fn get(&self, module_id: &str) -> Option<SupervisedModule> {
1853 let modules = self
1854 .modules
1855 .lock()
1856 .unwrap_or_else(|poisoned| poisoned.into_inner());
1857 modules.get(module_id).cloned()
1858 }
1859
1860 pub(crate) fn record_late_health_answer(
1861 &self,
1862 module_id: &str,
1863 latency_ms: u64,
1864 ) -> Result<bool, SuperviseError> {
1865 let Some(module) = self.get(module_id) else {
1866 return Ok(false);
1867 };
1868 update_snapshot(&module.inner.snapshot, Some(module_id), |state| {
1869 state.health.late_answer_count = state.health.late_answer_count.saturating_add(1);
1870 state.health.last_late_answer_latency_ms = Some(latency_ms);
1871 state.health.consecutive_failures = 0;
1879 })?;
1880 Ok(true)
1881 }
1882
1883 pub fn record_deliberate_severance(&self, module_id: &str) -> Result<bool, SuperviseError> {
1889 let Some(module) = self.get(module_id) else {
1890 return Ok(false);
1891 };
1892 let status = module.status()?;
1893 let Some((pid, start_time)) = status.pid.zip(status.process_start_time) else {
1894 return Ok(false);
1895 };
1896 module.record_deliberate_severance(ProcessIdentity { pid, start_time })
1897 }
1898
1899 pub fn list(&self) -> Vec<SupervisedModule> {
1900 let modules = self
1901 .modules
1902 .lock()
1903 .unwrap_or_else(|poisoned| poisoned.into_inner());
1904 let mut modules = modules.values().cloned().collect::<Vec<_>>();
1905 modules.sort_by(|left, right| left.module_id().cmp(right.module_id()));
1906 modules
1907 }
1908
1909 pub(crate) fn retire(&self, module_id: &str) -> Option<SupervisedModule> {
1910 self.spawn_nonces
1911 .lock()
1912 .unwrap_or_else(|poisoned| poisoned.into_inner())
1913 .remove(module_id);
1914 self.close_swap(module_id);
1915 let mut reserved_nonces = self
1916 .reserved_nonces
1917 .lock()
1918 .unwrap_or_else(|poisoned| poisoned.into_inner());
1919 if reserved_nonces.contains_key(module_id) {
1920 reserved_nonces.insert(module_id.to_string(), None);
1923 }
1924 drop(reserved_nonces);
1925 self.reserved_prefix_owners
1926 .lock()
1927 .unwrap_or_else(|poisoned| poisoned.into_inner())
1928 .retain(|_, owner| owner != module_id);
1929 self.modules
1930 .lock()
1931 .unwrap_or_else(|poisoned| poisoned.into_inner())
1932 .remove(module_id)
1933 }
1934
1935 pub(crate) fn record_rescan_removal(&self, module_id: &str) {
1938 self.removal_tombstones
1939 .lock()
1940 .unwrap_or_else(|poisoned| poisoned.into_inner())
1941 .insert(module_id.to_string(), unix_ms_now());
1942 }
1943
1944 pub(crate) fn removal_tombstone_age_ms(&self, module_id: &str) -> Option<u64> {
1946 self.removal_tombstones
1947 .lock()
1948 .unwrap_or_else(|poisoned| poisoned.into_inner())
1949 .get(module_id)
1950 .copied()
1951 .map(|removed_at_ms| unix_ms_now().saturating_sub(removed_at_ms))
1952 }
1953
1954 pub(crate) fn release_retained_reserved_gate(&self, module_id: &str) -> bool {
1959 if self.get(module_id).is_some() {
1960 return false;
1961 }
1962 let mut reserved_nonces = self
1963 .reserved_nonces
1964 .lock()
1965 .unwrap_or_else(|poisoned| poisoned.into_inner());
1966 if !matches!(reserved_nonces.get(module_id), Some(None)) {
1967 return false;
1968 }
1969 reserved_nonces.remove(module_id);
1970 true
1971 }
1972
1973 pub(crate) fn operation_lock(&self) -> Arc<AsyncMutex<()>> {
1974 Arc::clone(&self.operation_lock)
1975 }
1976}
1977
1978#[derive(Debug, Clone)]
1980pub struct Supervisor {
1981 registry: Arc<Registry>,
1982 restart_policy: RestartPolicy,
1983 drain_timeout: Duration,
1984 connection_file_path: Option<PathBuf>,
1985 capture_logs_dir: Option<PathBuf>,
1986 forwarding: Option<Arc<ForwardingTable>>,
1987 process_liveness: Arc<SupervisorProcessLiveness>,
1988 supervisor_handle: Option<SupervisorHandle>,
1989 health: HealthConfig,
1990 daemon_start_clock: crate::clock::StartClock,
1991 terminal_journal: Option<Arc<crate::terminal_journal::TerminalJournal>>,
1992 spawn_events: SpawnEventFeed,
1993 provenance_probe: ExecutableIdentityProbe,
1994 child_roster: ChildRoster,
1997 #[cfg(target_os = "linux")]
1998 cgroup_placement: Option<subc_cgroup::Placement>,
1999}
2000
2001impl Supervisor {
2002 #[cfg(unix)]
2013 pub(crate) fn begin_daemon_shutdown(&self) {
2014 self.child_roster.close();
2015 if let Some(journal) = &self.terminal_journal {
2016 journal.stamp_shutdown();
2017 }
2018 }
2019
2020 #[cfg(unix)]
2024 pub(crate) async fn drain_for_daemon_shutdown(&self) -> Result<(), SuperviseError> {
2025 const NOTICE_BUDGET: Duration = Duration::from_millis(500);
2026 const DRAIN_BUDGET: Duration = Duration::from_secs(2);
2027 let Some(forwarding) = &self.forwarding else {
2028 return Ok(());
2029 };
2030 let module_ids = forwarding
2031 .begin_daemon_drain()
2032 .map_err(SuperviseError::Forwarding)?;
2033 let deadline_ms =
2034 unix_ms_now().saturating_add((NOTICE_BUDGET + DRAIN_BUDGET).as_millis() as u64);
2035 let mut notices = tokio::task::JoinSet::new();
2036 let mut drains = Vec::new();
2037 for module_id in module_ids {
2038 let Some(target) = forwarding
2039 .begin_module_drain(&module_id, RouteCloseReason::Restart)
2040 .map_err(SuperviseError::Forwarding)?
2041 else {
2042 continue;
2043 };
2044 let routes = forwarding
2045 .endpoint_routes(target.endpoint)
2046 .map_err(SuperviseError::Forwarding)?;
2047 let command = serde_json::to_vec(&ModuleControlCommand::Draining {
2053 reason: RouteCloseReason::Restart,
2054 deadline_ms,
2055 })
2056 .expect("module draining serializes");
2057 let closing = serde_json::to_vec(&ClientControlPush::RouteClosing {
2058 module_id: module_id.clone(),
2059 reason: RouteCloseReason::Restart,
2060 })
2061 .expect("route closing serializes");
2062 let mut recipients = vec![(target.sink.clone(), target.negotiated_ver, command)];
2063 let mut seen = std::collections::HashSet::new();
2064 for route in routes {
2065 let client = route.goodbye_target;
2066 if seen.insert(client.connection_id) {
2067 recipients.push((client.sink, client.negotiated_ver, closing.clone()));
2068 }
2069 }
2070 for (sink, version, body) in recipients {
2071 notices.spawn(async move {
2072 let frame = Frame::build_with_version(
2073 version,
2074 FrameType::Push,
2075 control_flags(),
2076 0,
2077 0,
2078 0,
2079 body,
2080 )
2081 .expect("bounded lifecycle notice frame builds");
2082 sink.send_flushed(frame).await
2083 });
2084 }
2085 let gauges = declared_busy_gauges(&self.registry, &module_id)?;
2086 drains.push((module_id, target.endpoint, gauges));
2087 }
2088 let notice_deadline = Instant::now() + NOTICE_BUDGET;
2091 while let Ok(Some(result)) = timeout_at(notice_deadline, notices.join_next()).await {
2092 if !matches!(result, Ok(Ok(()))) {
2093 warn!(?result, "daemon shutdown notice delivery failed");
2094 }
2095 }
2096 notices.abort_all();
2097 let deadline = Instant::now() + DRAIN_BUDGET;
2098 let mut waits = tokio::task::JoinSet::new();
2099 for (module_id, endpoint, gauges) in drains {
2100 let forwarding = Arc::clone(forwarding);
2101 let mut runtime = self.runtime_config();
2102 runtime.health.cadence = Duration::from_millis(100);
2103 waits.spawn(async move {
2104 wait_for_forwarding_quiescence(
2105 &forwarding,
2106 &module_id,
2107 &runtime,
2108 endpoint,
2109 deadline,
2110 &gauges,
2111 DrainScope::Active,
2112 )
2113 .await
2114 });
2115 }
2116 while let Ok(Some(result)) = timeout_at(deadline, waits.join_next()).await {
2117 if !matches!(result, Ok(Ok(true))) {
2118 warn!(?result, "daemon shutdown drain did not reach quiescence");
2119 }
2120 }
2121 Ok(())
2122 }
2123
2124 #[cfg(unix)]
2134 pub(crate) async fn end_children_for_daemon_shutdown(
2135 &self,
2136 already_escalated: bool,
2137 escalate: impl std::future::Future<Output = ()>,
2138 ) {
2139 if let Some(forwarding) = &self.forwarding {
2140 let closed = forwarding.close_all_connections(&CloseReason::new(
2141 "daemon_shutdown",
2142 "the daemon is exiting after its shutdown notice and drain",
2143 ));
2144 debug!(closed, "closed established connections for daemon shutdown");
2145 }
2146 crate::child_roster::end_children_for_daemon_shutdown(
2147 &self.child_roster,
2148 already_escalated,
2149 escalate,
2150 )
2151 .await;
2152 }
2153
2154 pub fn new(registry: Arc<Registry>, restart_policy: RestartPolicy) -> Self {
2155 Self {
2156 registry,
2157 restart_policy,
2158 drain_timeout: DEFAULT_DRAIN_TIMEOUT,
2159 connection_file_path: None,
2160 capture_logs_dir: None,
2161 forwarding: None,
2162 process_liveness: Arc::new(SupervisorProcessLiveness::default()),
2163 supervisor_handle: None,
2164 health: HealthConfig::default(),
2165 daemon_start_clock: crate::clock::StartClock::capture(),
2166 terminal_journal: None,
2167 spawn_events: SpawnEventFeed::default(),
2168 provenance_probe: ExecutableIdentityProbe::default(),
2169 child_roster: ChildRoster::default(),
2170 #[cfg(target_os = "linux")]
2171 cgroup_placement: None,
2172 }
2173 }
2174
2175 pub fn with_drain_timeout(mut self, drain_timeout: Duration) -> Self {
2176 self.drain_timeout = drain_timeout;
2177 self
2178 }
2179
2180 pub fn with_process_liveness(
2181 mut self,
2182 process_liveness: Arc<SupervisorProcessLiveness>,
2183 ) -> Self {
2184 self.process_liveness = process_liveness;
2185 self
2186 }
2187
2188 pub fn with_connection_file_path(mut self, connection_file_path: impl Into<PathBuf>) -> Self {
2189 self.connection_file_path = Some(connection_file_path.into());
2190 self
2191 }
2192
2193 pub fn with_capture_logs_dir(mut self, logs_dir: impl Into<PathBuf>) -> Self {
2195 self.capture_logs_dir = Some(logs_dir.into());
2196 self
2197 }
2198
2199 pub fn with_daemon_incarnation(self, daemon_incarnation: String) -> Self {
2202 self.spawn_events.configure_incarnation(daemon_incarnation);
2206 self
2207 }
2208
2209 pub fn with_terminal_journal(self, path: PathBuf, daemon_incarnation: String) -> Self {
2212 let mut this = self.with_daemon_incarnation(daemon_incarnation.clone());
2213 this.terminal_journal = Some(Arc::new(crate::terminal_journal::TerminalJournal::open(
2214 path,
2215 daemon_incarnation,
2216 )));
2217 this
2218 }
2219
2220 pub fn with_forwarding(mut self, forwarding: Arc<ForwardingTable>) -> Self {
2221 self.forwarding = Some(forwarding);
2222 self
2223 }
2224
2225 pub fn with_handle(mut self, supervisor_handle: SupervisorHandle) -> Self {
2226 self.spawn_events = supervisor_handle.spawn_events.clone();
2227 self.supervisor_handle = Some(supervisor_handle);
2228 self
2229 }
2230
2231 pub fn with_health_config(mut self, health: HealthConfig) -> Self {
2232 self.health = health;
2233 self
2234 }
2235
2236 pub fn with_live_children_record(self, path: impl Into<PathBuf>) -> Self {
2240 self.child_roster.record_to(path.into());
2241 self
2242 }
2243
2244 #[cfg(target_os = "linux")]
2245 pub fn with_cgroup_placement(
2246 mut self,
2247 cgroup_placement: Option<subc_cgroup::Placement>,
2248 ) -> Self {
2249 self.cgroup_placement = cgroup_placement;
2250 self
2251 }
2252
2253 pub fn spawn(&self, spec: ModuleSpec) -> Result<SupervisedModule, SuperviseError> {
2259 validate_spec(&spec)?;
2260
2261 let runtime = self.runtime_config();
2262 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2263 let child = spawn_child(
2264 &spec,
2265 runtime.connection_file_path.as_deref(),
2266 self.supervisor_handle.as_ref(),
2267 &runtime.stderr_ring,
2268 runtime.capture_logs_dir.as_deref(),
2269 &runtime.child_roster,
2270 #[cfg(target_os = "linux")]
2271 runtime.cgroup_placement.as_ref(),
2272 )?;
2273 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2274 self.process_liveness
2275 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2276
2277 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2278 }
2279
2280 pub fn supervise_configured(
2286 &self,
2287 spec: ModuleSpec,
2288 enabled: bool,
2289 ) -> Result<SupervisedModule, SuperviseError> {
2290 validate_spec(&spec)?;
2291
2292 let runtime = self.runtime_config();
2293 if !enabled {
2294 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2295 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2296 }
2297
2298 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2299 match spawn_child(
2300 &spec,
2301 runtime.connection_file_path.as_deref(),
2302 self.supervisor_handle.as_ref(),
2303 &runtime.stderr_ring,
2304 runtime.capture_logs_dir.as_deref(),
2305 &runtime.child_roster,
2306 #[cfg(target_os = "linux")]
2307 runtime.cgroup_placement.as_ref(),
2308 ) {
2309 Ok(child) => {
2310 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2311 self.process_liveness
2312 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2313 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2314 }
2315 Err(err) => {
2316 error!(
2317 module_id = %spec.module_id,
2318 program = %spec.program.display(),
2319 error = %err,
2320 "configured module failed to spawn; marking failed and continuing"
2321 );
2322 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2323 Ok(self.supervised_module(spec, runtime, snapshot, None))
2324 }
2325 }
2326 }
2327
2328 pub fn supervise_configured_with_health(
2334 &self,
2335 spec: ModuleSpec,
2336 enabled: bool,
2337 health: HealthConfig,
2338 drain_timeout_ms: Option<u64>,
2339 restart_policy: RestartPolicy,
2340 ) -> Result<SupervisedModule, SuperviseError> {
2341 validate_spec(&spec)?;
2342
2343 let mut runtime = self.runtime_config();
2344 runtime.health = health;
2345 runtime.restart_policy = restart_policy;
2346 if let Some(ms) = drain_timeout_ms {
2347 runtime.drain_timeout = Duration::from_millis(ms);
2348 *runtime
2349 .effective_drain_timeout
2350 .lock()
2351 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
2352 }
2353 if !enabled {
2354 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2355 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2356 }
2357
2358 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2359 match spawn_child(
2360 &spec,
2361 runtime.connection_file_path.as_deref(),
2362 self.supervisor_handle.as_ref(),
2363 &runtime.stderr_ring,
2364 runtime.capture_logs_dir.as_deref(),
2365 &runtime.child_roster,
2366 #[cfg(target_os = "linux")]
2367 runtime.cgroup_placement.as_ref(),
2368 ) {
2369 Ok(child) => {
2370 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2371 self.process_liveness
2372 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2373 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2374 }
2375 Err(err) => {
2376 if health.critical {
2377 error!(
2378 module_id = %spec.module_id,
2379 program = %spec.program.display(),
2380 error = %err,
2381 "critical configured module failed to spawn; marking failed and alerting"
2382 );
2383 } else {
2384 error!(
2385 module_id = %spec.module_id,
2386 program = %spec.program.display(),
2387 error = %err,
2388 "configured module failed to spawn; marking failed and continuing"
2389 );
2390 }
2391 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2392 Ok(self.supervised_module(spec, runtime, snapshot, None))
2393 }
2394 }
2395 }
2396
2397 fn runtime_config(&self) -> SupervisorRuntimeConfig {
2398 let effective_drain_timeout = Arc::new(Mutex::new(self.drain_timeout));
2399 SupervisorRuntimeConfig {
2400 restart_policy: self.restart_policy,
2401 drain_timeout: self.drain_timeout,
2402 child_roster: self
2405 .child_roster
2406 .for_module(Arc::clone(&effective_drain_timeout)),
2407 effective_drain_timeout,
2408 default_drain_timeout: self.drain_timeout,
2409 health: self.health,
2410 connection_file_path: self.connection_file_path.clone(),
2411 capture_logs_dir: self.capture_logs_dir.clone(),
2412 forwarding: self.forwarding.clone(),
2413 supervisor_handle: self.supervisor_handle.clone(),
2414 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
2415 terminal_ring: Arc::new(Mutex::new(
2416 TerminalRing::new(
2417 TerminalRingConfig::default(),
2418 self.daemon_start_clock.started_at_ms(),
2419 )
2420 .with_start_clock(self.daemon_start_clock)
2421 .with_journal(self.terminal_journal.clone())
2422 .with_daemon_shutdown(self.child_roster.shutdown_flag()),
2423 )),
2424 spawn_events: self.spawn_events.clone(),
2425 #[cfg(target_os = "linux")]
2426 cgroup_placement: self.cgroup_placement.clone(),
2427 #[cfg(test)]
2428 test_seed_stale_facts_before_enable_spawn: false,
2429 }
2430 }
2431
2432 fn supervised_module(
2433 &self,
2434 spec: ModuleSpec,
2435 runtime: SupervisorRuntimeConfig,
2436 snapshot: SharedSnapshot,
2437 child: Option<SupervisedChild>,
2438 ) -> SupervisedModule {
2439 let configuration = Arc::new(Mutex::new(SupervisedConfiguration {
2440 spec: spec.clone(),
2441 health: runtime.health,
2442 }));
2443 let stderr_ring = Arc::clone(&runtime.stderr_ring);
2444 let terminal_ring = Arc::clone(&runtime.terminal_ring);
2445 let restart_policy = runtime.restart_policy;
2449 let effective_drain_timeout = Arc::clone(&runtime.effective_drain_timeout);
2450 let (tx, rx) = mpsc::channel(4);
2451 let monitor = tokio::spawn(supervise_loop(
2452 spec.clone(),
2453 runtime,
2454 Arc::clone(&self.registry),
2455 Arc::clone(&self.process_liveness),
2456 Arc::clone(&snapshot),
2457 child,
2458 rx,
2459 ));
2460
2461 let module_id = spec.module_id.clone();
2462 let module = SupervisedModule {
2463 inner: Arc::new(SupervisedModuleInner {
2464 module_id: module_id.clone(),
2465 registry: Arc::clone(&self.registry),
2466 snapshot,
2467 configuration,
2468 stderr_ring,
2469 terminal_ring,
2470 commands: tx,
2471 monitor: Mutex::new(Some(monitor)),
2472 restart_policy,
2473 effective_drain_timeout,
2474 provenance_probe: self.provenance_probe.clone(),
2475 }),
2476 };
2477 if let Some(supervisor_handle) = &self.supervisor_handle {
2478 supervisor_handle.apply_identity_configuration(&spec);
2479 supervisor_handle.insert(module.clone());
2480 }
2481 module
2482 }
2483}
2484
2485impl Default for Supervisor {
2486 fn default() -> Self {
2487 Self::new(Arc::new(Registry::default()), RestartPolicy::default())
2488 }
2489}
2490
2491#[derive(Clone)]
2493pub struct SupervisedModule {
2494 inner: Arc<SupervisedModuleInner>,
2495}
2496
2497struct SupervisedModuleInner {
2498 module_id: String,
2499 registry: Arc<Registry>,
2500 snapshot: SharedSnapshot,
2501 configuration: Arc<Mutex<SupervisedConfiguration>>,
2502 stderr_ring: Arc<Mutex<StderrRing>>,
2503 terminal_ring: Arc<Mutex<TerminalRing>>,
2504 commands: mpsc::Sender<SupervisorCommand>,
2505 monitor: Mutex<Option<JoinHandle<()>>>,
2506 restart_policy: RestartPolicy,
2510 effective_drain_timeout: Arc<Mutex<Duration>>,
2511 provenance_probe: ExecutableIdentityProbe,
2512}
2513
2514impl fmt::Debug for SupervisedModule {
2515 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2516 f.debug_struct("SupervisedModule")
2517 .field("module_id", &self.inner.module_id)
2518 .field("status", &self.status())
2519 .finish_non_exhaustive()
2520 }
2521}
2522
2523impl SupervisedModule {
2524 pub fn module_id(&self) -> &str {
2525 &self.inner.module_id
2526 }
2527
2528 #[cfg(test)]
2532 pub(crate) fn record_health_probe_failure_for_test(
2533 &self,
2534 detail: &str,
2535 ) -> Result<(), SuperviseError> {
2536 update_snapshot(&self.inner.snapshot, Some(&self.inner.module_id), |state| {
2537 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
2538 state.health.detail = Some(detail.to_string());
2539 })
2540 }
2541
2542 pub fn state(&self) -> Result<ModuleState, SuperviseError> {
2543 Ok(lock_snapshot(&self.inner.snapshot)?.state)
2544 }
2545
2546 pub fn stderr_tail(
2553 &self,
2554 max_lines: Option<usize>,
2555 max_bytes: Option<usize>,
2556 ) -> StderrTailSnapshot {
2557 self.inner
2558 .stderr_ring
2559 .lock()
2560 .unwrap_or_else(|poisoned| poisoned.into_inner())
2561 .snapshot(max_lines, max_bytes)
2562 }
2563
2564 pub fn terminal_history(&self) -> TerminalHistorySnapshot {
2569 self.inner
2570 .terminal_ring
2571 .lock()
2572 .unwrap_or_else(|poisoned| poisoned.into_inner())
2573 .snapshot()
2574 }
2575
2576 pub fn durable_terminal_history(&self) -> subc_control::TerminalHistory {
2581 durable_terminal_history_of(&self.inner.terminal_ring, &self.inner.module_id)
2582 }
2583
2584 pub(crate) async fn read_durable_terminal_history(
2589 &self,
2590 ) -> Result<subc_control::TerminalHistory, tokio::task::JoinError> {
2591 let terminal_ring = Arc::clone(&self.inner.terminal_ring);
2592 let module_id = self.inner.module_id.clone();
2593 tokio::task::spawn_blocking(move || durable_terminal_history_of(&terminal_ring, &module_id))
2594 .await
2595 }
2596
2597 pub fn status(&self) -> Result<ModuleStatus, SuperviseError> {
2598 self.status_with_snapshot_lock(&self.inner.snapshot, None)
2599 }
2600
2601 pub(crate) fn record_deliberate_severance(
2602 &self,
2603 identity: ProcessIdentity,
2604 ) -> Result<bool, SuperviseError> {
2605 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2606 if snapshot.pid != Some(identity.pid)
2607 || snapshot.process_start_time != Some(identity.start_time)
2608 {
2609 return Ok(false);
2610 }
2611 snapshot.deliberate_severance = Some(identity);
2612 Ok(true)
2613 }
2614
2615 pub(crate) fn status_for_control(
2620 &self,
2621 caller: &'static str,
2622 ) -> Result<ModuleStatus, SuperviseError> {
2623 self.status_with_snapshot_lock(&self.inner.snapshot, Some(caller))
2624 }
2625
2626 fn status_with_snapshot_lock(
2627 &self,
2628 snapshot: &SharedSnapshot,
2629 caller: Option<&'static str>,
2630 ) -> Result<ModuleStatus, SuperviseError> {
2631 let mut guard = match caller {
2632 Some(caller) => lock_snapshot_for_control(snapshot, &self.inner.module_id, caller)?,
2633 None => lock_snapshot(snapshot)?,
2634 };
2635 let restart_count =
2638 guard.crash_restarts_in_window(self.inner.restart_policy.window, Instant::now());
2639 let snapshot = guard.clone();
2640 drop(guard);
2641 let drain_timeout = *self.inner.effective_drain_timeout.lock().map_err(|_| {
2642 SuperviseError::StatePoisoned {
2643 module_id: Some(self.inner.module_id.clone()),
2644 }
2645 })?;
2646 let registration_active = self
2647 .inner
2648 .registry
2649 .get_module(&self.inner.module_id)
2650 .map_err(SuperviseError::Registry)?
2651 .is_some();
2652 let protocol = self.declared_protocol()?;
2653 let running_process =
2654 snapshot.enabled && snapshot.state == ModuleState::Running && snapshot.process_alive;
2655 let live = match protocol {
2661 ModuleProtocol::Subc => running_process && registration_active,
2662 ModuleProtocol::None => running_process,
2663 };
2664
2665 Ok(ModuleStatus {
2666 module_id: self.inner.module_id.clone(),
2667 state: snapshot.state,
2668 enabled: snapshot.enabled,
2669 process_alive: snapshot.process_alive,
2670 registration_active,
2671 protocol,
2672 live,
2673 restart_count,
2674 lifetime_restarts: snapshot.lifetime_restarts,
2675 spawn_generation: snapshot.spawn_generation,
2676 max_restarts: self.inner.restart_policy.max_restarts,
2677 restart_window: self.inner.restart_policy.window,
2678 drain_timeout,
2679 restart_backoff: self.inner.restart_policy.backoff,
2680 restart_max_backoff: self.inner.restart_policy.max_backoff,
2681 pid: snapshot.pid,
2682 spawned_at_ms: snapshot.spawned_at_ms,
2683 spawned_from: snapshot.spawned_from,
2684 process_start_time: snapshot.process_start_time,
2685 last_exit: snapshot.last_exit,
2686 health: snapshot.health,
2687 })
2688 }
2689
2690 #[cfg(test)]
2691 pub(crate) fn hold_snapshot_for_test(
2692 &self,
2693 acquired: std::sync::mpsc::Sender<()>,
2694 hold: Duration,
2695 ) -> std::thread::JoinHandle<()> {
2696 let snapshot = Arc::clone(&self.inner.snapshot);
2697 std::thread::spawn(move || {
2698 let _guard = snapshot.lock().expect("test snapshot lock is not poisoned");
2699 acquired
2700 .send(())
2701 .expect("test receiver waits for snapshot lock");
2702 std::thread::sleep(hold);
2703 })
2704 }
2705
2706 pub(crate) async fn running_image_agreement(&self) -> subc_control::RunningImageAgreement {
2707 let snapshot = match lock_snapshot(&self.inner.snapshot) {
2708 Ok(snapshot) => snapshot.clone(),
2709 Err(_) => {
2710 return subc_control::RunningImageAgreement::Unavailable {
2711 reason: subc_control::RunningImageUnavailableReason::NotRunning,
2712 };
2713 }
2714 };
2715 self.inner
2716 .provenance_probe
2717 .observe(
2718 snapshot.pid,
2719 snapshot.spawned_from.as_deref(),
2720 snapshot.spawned_file_identity,
2721 snapshot.process_start_time,
2722 )
2723 .await
2724 }
2725
2726 pub(crate) fn will_recover_after_connection_loss(&self) -> Result<bool, SuperviseError> {
2727 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2728 Ok(match snapshot.state {
2729 ModuleState::Restarting => true,
2730 ModuleState::Failed | ModuleState::Disabled => false,
2731 _ => daemon_will_restart(&mut snapshot, &self.inner.restart_policy, Instant::now()),
2732 })
2733 }
2734
2735 #[cfg(test)]
2736 pub(crate) fn is_warming(&self) -> Result<bool, SuperviseError> {
2737 self.is_warming_with_snapshot_lock(None)
2738 }
2739
2740 pub(crate) fn is_warming_for_control(
2741 &self,
2742 caller: &'static str,
2743 ) -> Result<bool, SuperviseError> {
2744 self.is_warming_with_snapshot_lock(Some(caller))
2745 }
2746
2747 fn is_warming_with_snapshot_lock(
2748 &self,
2749 caller: Option<&'static str>,
2750 ) -> Result<bool, SuperviseError> {
2751 let snapshot = match caller {
2752 Some(caller) => {
2753 lock_snapshot_for_control(&self.inner.snapshot, &self.inner.module_id, caller)?
2754 }
2755 None => lock_snapshot(&self.inner.snapshot)?,
2756 }
2757 .clone();
2758 Ok(matches!(
2759 snapshot.state,
2760 ModuleState::Starting | ModuleState::Running | ModuleState::Restarting
2761 ))
2762 }
2763
2764 pub async fn drain(&self) -> Result<(), SuperviseError> {
2766 self.stop().await
2767 }
2768
2769 pub(crate) async fn retire(&self) -> Result<(), SuperviseError> {
2770 match self.state()? {
2771 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2772 ModuleState::Starting
2773 | ModuleState::Running
2774 | ModuleState::Unresponsive
2775 | ModuleState::Restarting
2776 | ModuleState::Draining
2777 | ModuleState::Disabled => {}
2778 }
2779
2780 let (reply_tx, reply_rx) = oneshot::channel();
2781 self.inner
2782 .commands
2783 .send(SupervisorCommand::Retire { reply: reply_tx })
2784 .await
2785 .map_err(|_| SuperviseError::CommandClosed {
2786 module_id: self.inner.module_id.clone(),
2787 })?;
2788 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2789 module_id: self.inner.module_id.clone(),
2790 })?
2791 }
2792
2793 pub async fn stop(&self) -> Result<(), SuperviseError> {
2794 match self.state()? {
2795 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2796 ModuleState::Starting
2797 | ModuleState::Running
2798 | ModuleState::Unresponsive
2799 | ModuleState::Restarting
2800 | ModuleState::Draining
2801 | ModuleState::Disabled => {}
2802 }
2803
2804 let (reply_tx, reply_rx) = oneshot::channel();
2805 self.inner
2806 .commands
2807 .send(SupervisorCommand::Drain { reply: reply_tx })
2808 .await
2809 .map_err(|_| SuperviseError::CommandClosed {
2810 module_id: self.inner.module_id.clone(),
2811 })?;
2812 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2813 module_id: self.inner.module_id.clone(),
2814 })?
2815 }
2816
2817 pub async fn restart(&self, drain_timeout_ms: Option<u64>) -> Result<(), SuperviseError> {
2818 let received_at_generation = lock_snapshot(&self.inner.snapshot)?.spawn_generation;
2819 let (reply_tx, reply_rx) = oneshot::channel();
2820 self.inner
2821 .commands
2822 .send(SupervisorCommand::Restart {
2823 drain_timeout_ms,
2824 received_at_generation,
2825 queued_at: Instant::now(),
2826 reply: reply_tx,
2827 })
2828 .await
2829 .map_err(|_| SuperviseError::CommandClosed {
2830 module_id: self.inner.module_id.clone(),
2831 })?;
2832 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2833 module_id: self.inner.module_id.clone(),
2834 })?
2835 }
2836
2837 pub async fn swap(&self, ready_timeout: Option<Duration>) -> Result<(), SuperviseError> {
2842 let (reply_tx, reply_rx) = oneshot::channel();
2843 self.inner
2844 .commands
2845 .send(SupervisorCommand::Swap {
2846 ready_timeout,
2847 reply: reply_tx,
2848 })
2849 .await
2850 .map_err(|_| SuperviseError::CommandClosed {
2851 module_id: self.inner.module_id.clone(),
2852 })?;
2853 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2854 module_id: self.inner.module_id.clone(),
2855 })?
2856 }
2857
2858 pub async fn reload(&self) -> Result<(), SuperviseError> {
2859 let (reply_tx, reply_rx) = oneshot::channel();
2860 self.inner
2861 .commands
2862 .send(SupervisorCommand::Reload { reply: reply_tx })
2863 .await
2864 .map_err(|_| SuperviseError::CommandClosed {
2865 module_id: self.inner.module_id.clone(),
2866 })?;
2867 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2868 module_id: self.inner.module_id.clone(),
2869 })?
2870 }
2871
2872 pub async fn set_enabled(&self, enabled: bool) -> Result<bool, SuperviseError> {
2873 let (reply_tx, reply_rx) = oneshot::channel();
2874 self.inner
2875 .commands
2876 .send(SupervisorCommand::SetEnabled {
2877 enabled,
2878 reply: reply_tx,
2879 })
2880 .await
2881 .map_err(|_| SuperviseError::CommandClosed {
2882 module_id: self.inner.module_id.clone(),
2883 })?;
2884 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2885 module_id: self.inner.module_id.clone(),
2886 })?
2887 }
2888
2889 pub(crate) fn declared_protocol(&self) -> Result<ModuleProtocol, SuperviseError> {
2894 Ok(self
2895 .inner
2896 .configuration
2897 .lock()
2898 .map_err(|_| SuperviseError::StatePoisoned {
2899 module_id: Some(self.inner.module_id.clone()),
2900 })?
2901 .spec
2902 .protocol)
2903 }
2904
2905 pub(crate) fn configuration(&self) -> Result<(ModuleSpec, HealthConfig), SuperviseError> {
2906 let configuration =
2907 self.inner
2908 .configuration
2909 .lock()
2910 .map_err(|_| SuperviseError::StatePoisoned {
2911 module_id: Some(self.inner.module_id.clone()),
2912 })?;
2913 Ok((configuration.spec.clone(), configuration.health))
2914 }
2915
2916 #[cfg(any(test, feature = "test-support"))]
2920 pub async fn update_spec_for_test(&self, spec: ModuleSpec) -> Result<(), SuperviseError> {
2921 let (_, health) = self.configuration()?;
2922 let drain_timeout_ms = u64::try_from(
2923 self.inner
2924 .effective_drain_timeout
2925 .lock()
2926 .unwrap_or_else(|poisoned| poisoned.into_inner())
2927 .as_millis(),
2928 )
2929 .ok();
2930 self.update_configuration(spec, health, drain_timeout_ms)
2931 .await
2932 }
2933
2934 pub(crate) async fn update_configuration(
2935 &self,
2936 spec: ModuleSpec,
2937 health: HealthConfig,
2938 drain_timeout_ms: Option<u64>,
2939 ) -> Result<(), SuperviseError> {
2940 if spec.module_id != self.inner.module_id {
2941 return Err(SuperviseError::InvalidSpec {
2942 reason: "a supervised module's module_id cannot be changed".to_string(),
2943 });
2944 }
2945 validate_spec(&spec)?;
2946 let (reply_tx, reply_rx) = oneshot::channel();
2947 self.inner
2948 .commands
2949 .send(SupervisorCommand::UpdateConfiguration {
2950 spec: spec.clone(),
2951 health,
2952 drain_timeout_ms,
2953 reply: reply_tx,
2954 })
2955 .await
2956 .map_err(|_| SuperviseError::CommandClosed {
2957 module_id: self.inner.module_id.clone(),
2958 })?;
2959 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2960 module_id: self.inner.module_id.clone(),
2961 })?;
2962 let mut configuration =
2963 self.inner
2964 .configuration
2965 .lock()
2966 .map_err(|_| SuperviseError::StatePoisoned {
2967 module_id: Some(self.inner.module_id.clone()),
2968 })?;
2969 configuration.spec = spec;
2970 configuration.health = health;
2971 Ok(())
2972 }
2973}
2974
2975impl Drop for SupervisedModuleInner {
2976 fn drop(&mut self) {
2977 let Ok(mut monitor) = self.monitor.lock() else {
2978 return;
2979 };
2980 if let Some(monitor) = monitor.as_ref().filter(|monitor| !monitor.is_finished()) {
2981 let _ = update_snapshot(&self.snapshot, Some(&self.module_id), |state| {
2982 state.state = ModuleState::Stopped;
2983 clear_current_process_facts(state);
2984 });
2985 monitor.abort();
2986 }
2987 let _ = monitor.take();
2988 }
2989}
2990
2991#[derive(Debug)]
2992enum SupervisorCommand {
2993 Drain {
2994 reply: oneshot::Sender<Result<(), SuperviseError>>,
2995 },
2996 Retire {
2997 reply: oneshot::Sender<Result<(), SuperviseError>>,
2998 },
2999 Restart {
3000 drain_timeout_ms: Option<u64>,
3005 received_at_generation: u64,
3009 queued_at: Instant,
3012 reply: oneshot::Sender<Result<(), SuperviseError>>,
3013 },
3014 Reload {
3015 reply: oneshot::Sender<Result<(), SuperviseError>>,
3016 },
3017 SetEnabled {
3018 enabled: bool,
3019 reply: oneshot::Sender<Result<bool, SuperviseError>>,
3020 },
3021 UpdateConfiguration {
3022 spec: ModuleSpec,
3023 health: HealthConfig,
3024 drain_timeout_ms: Option<u64>,
3027 reply: oneshot::Sender<()>,
3028 },
3029 Swap {
3030 ready_timeout: Option<Duration>,
3033 reply: oneshot::Sender<Result<(), SuperviseError>>,
3035 },
3036}
3037
3038#[derive(Debug)]
3039pub enum SuperviseError {
3040 InvalidSpec {
3041 reason: String,
3042 },
3043 Spawn {
3044 program: PathBuf,
3045 source: io::Error,
3046 cgroup_path: Option<PathBuf>,
3047 },
3048 Cgroup {
3049 module_id: String,
3050 source: io::Error,
3051 },
3052 LaunchNonce {
3055 reason: String,
3056 },
3057 Wait {
3058 module_id: String,
3059 source: io::Error,
3060 },
3061 Kill {
3062 module_id: String,
3063 source: io::Error,
3064 },
3065 Forwarding(ForwardingError),
3066 Registry(RegistryError),
3067 ReloadUnavailable {
3068 module_id: String,
3069 reason: String,
3070 },
3071 Disabled {
3076 module_id: String,
3077 },
3078 ReloadFailed {
3079 module_id: String,
3080 reason: String,
3081 },
3082 RegistrationStillActive {
3083 module_id: String,
3084 waited: Duration,
3085 },
3086 StatePoisoned {
3087 module_id: Option<String>,
3088 },
3089 CommandClosed {
3090 module_id: String,
3091 },
3092 SwapInProgress {
3096 module_id: String,
3097 },
3098 SwapRefused {
3100 module_id: String,
3101 reason: SwapRefusal,
3102 },
3103 SwapFailed {
3107 module_id: String,
3108 arm: SwapFailureArm,
3109 detail: String,
3110 candidate_exit: Option<ExitReport>,
3113 },
3114}
3115
3116#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3118pub enum SwapRefusal {
3119 OverlapExclusive,
3121 NotRegistered,
3124 ProtocolNone,
3127 NotConfigured,
3130 AlreadySwapping,
3132}
3133
3134impl SwapRefusal {
3135 pub fn as_str(self) -> &'static str {
3136 match self {
3137 Self::OverlapExclusive => "overlap_exclusive",
3138 Self::NotRegistered => "not_registered",
3139 Self::ProtocolNone => "protocol_none",
3140 Self::NotConfigured => "not_configured",
3141 Self::AlreadySwapping => "already_swapping",
3142 }
3143 }
3144}
3145
3146#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3149pub enum SwapFailureArm {
3150 SpawnFailed,
3152 NeverRegistered,
3154 NeverReady,
3156 CandidateExited,
3158 CandidateUnhealthy,
3160 Interrupted,
3164 CutoverLost,
3169}
3170
3171impl SwapFailureArm {
3172 pub fn as_str(self) -> &'static str {
3173 match self {
3174 Self::SpawnFailed => "spawn_failed",
3175 Self::NeverRegistered => "never_registered",
3176 Self::NeverReady => "never_ready",
3177 Self::CandidateExited => "candidate_exited",
3178 Self::CandidateUnhealthy => "candidate_unhealthy",
3179 Self::Interrupted => "interrupted",
3180 Self::CutoverLost => "cutover_lost",
3181 }
3182 }
3183}
3184
3185impl fmt::Display for SuperviseError {
3186 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3187 match self {
3188 Self::InvalidSpec { reason } => write!(f, "invalid module spec: {reason}"),
3189 Self::Spawn {
3190 program,
3191 source,
3192 cgroup_path: Some(cgroup_path),
3193 } => write!(
3194 f,
3195 "failed to place module in cgroup '{}' while spawning '{}': {source}",
3196 cgroup_path.display(),
3197 program.display()
3198 ),
3199 Self::Spawn {
3200 program,
3201 source,
3202 cgroup_path: None,
3203 } => write!(
3204 f,
3205 "failed to spawn module '{}': {source}",
3206 program.display()
3207 ),
3208 Self::Cgroup { module_id, source } => {
3209 write!(
3210 f,
3211 "failed to prepare cgroup for module '{module_id}': {source}"
3212 )
3213 }
3214 Self::LaunchNonce { reason } => {
3215 write!(
3216 f,
3217 "failed to generate reserved-module launch nonce: {reason}"
3218 )
3219 }
3220 Self::Wait { module_id, source } => {
3221 write!(f, "failed to wait for module '{module_id}': {source}")
3222 }
3223 Self::Kill { module_id, source } => {
3224 write!(f, "failed to kill module '{module_id}': {source}")
3225 }
3226 Self::Forwarding(err) => write!(f, "forwarding error: {err}"),
3227 Self::Registry(err) => write!(f, "registry error: {err}"),
3228 Self::ReloadUnavailable { module_id, reason } => {
3229 write!(f, "reload unavailable for module '{module_id}': {reason}")
3230 }
3231 Self::Disabled { module_id } => {
3232 write!(
3233 f,
3234 "module '{module_id}' is disabled; enable it before restart or reload"
3235 )
3236 }
3237 Self::ReloadFailed { module_id, reason } => {
3238 write!(f, "reload failed for module '{module_id}': {reason}")
3239 }
3240 Self::RegistrationStillActive { module_id, waited } => write!(
3241 f,
3242 "module '{module_id}' registration remained active after waiting {waited:?}"
3243 ),
3244 Self::StatePoisoned { module_id } => match module_id {
3245 Some(module_id) => {
3246 write!(f, "supervisor state for module '{module_id}' was poisoned")
3247 }
3248 None => write!(f, "supervisor state was poisoned"),
3249 },
3250 Self::CommandClosed { module_id } => {
3251 write!(
3252 f,
3253 "supervisor command channel for module '{module_id}' is closed"
3254 )
3255 }
3256 Self::SwapInProgress { module_id } => write!(
3257 f,
3258 "module '{module_id}' is being swapped; retry once the swap has cut over or failed, or stop the module to abort the swap"
3259 ),
3260 Self::SwapRefused { module_id, reason } => match reason {
3261 SwapRefusal::OverlapExclusive => write!(
3262 f,
3263 "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"
3264 ),
3265 SwapRefusal::NotRegistered => write!(
3266 f,
3267 "module '{module_id}' is not registered, so there is no serving process to keep while a replacement warms; use a plain restart"
3268 ),
3269 SwapRefusal::ProtocolNone => write!(
3270 f,
3271 "module '{module_id}' is protocol: \"none\" and never registers, so a swap could never see its replacement become ready; use a plain restart"
3272 ),
3273 SwapRefusal::NotConfigured => write!(
3274 f,
3275 "module '{module_id}' cannot be swapped: the supervisor was built without the forwarding table or shared handle a swap needs"
3276 ),
3277 SwapRefusal::AlreadySwapping => {
3278 write!(f, "module '{module_id}' is already being swapped")
3279 }
3280 },
3281 Self::SwapFailed {
3282 module_id,
3283 arm,
3284 detail,
3285 ..
3286 } => write!(
3287 f,
3288 "swap of module '{module_id}' failed ({}): {detail}; the running process was left serving",
3289 arm.as_str()
3290 ),
3291 }
3292 }
3293}
3294
3295impl Error for SuperviseError {
3296 fn source(&self) -> Option<&(dyn Error + 'static)> {
3297 match self {
3298 Self::Spawn { source, .. }
3299 | Self::Cgroup { source, .. }
3300 | Self::Wait { source, .. }
3301 | Self::Kill { source, .. } => Some(source),
3302 Self::Forwarding(err) => Some(err),
3303 Self::Registry(err) => Some(err),
3304 Self::LaunchNonce { .. }
3305 | Self::InvalidSpec { .. }
3306 | Self::ReloadUnavailable { .. }
3307 | Self::Disabled { .. }
3308 | Self::ReloadFailed { .. }
3309 | Self::RegistrationStillActive { .. }
3310 | Self::StatePoisoned { .. }
3311 | Self::CommandClosed { .. }
3312 | Self::SwapInProgress { .. }
3313 | Self::SwapRefused { .. }
3314 | Self::SwapFailed { .. } => None,
3315 }
3316 }
3317}
3318
3319pub(crate) fn validate_spec(spec: &ModuleSpec) -> Result<(), SuperviseError> {
3320 if spec.module_id.trim().is_empty() {
3321 return Err(SuperviseError::InvalidSpec {
3322 reason: "module_id must not be empty".to_string(),
3323 });
3324 }
3325
3326 Ok(())
3327}
3328
3329#[derive(Debug, Default)]
3330struct HealthProbeRuntime {
3331 registered_connection: Option<crate::ConnectionId>,
3332 advertised: bool,
3333 next_probe_at: Option<Instant>,
3334 probe_index: u64,
3335}
3336
3337impl HealthProbeRuntime {
3338 fn refresh_registration(
3339 &mut self,
3340 spec: &ModuleSpec,
3341 runtime: &SupervisorRuntimeConfig,
3342 registry: &Registry,
3343 snapshot: &SharedSnapshot,
3344 ) {
3345 if spec.protocol == ModuleProtocol::None {
3357 self.registered_connection = None;
3358 self.advertised = false;
3359 self.next_probe_at = None;
3360 return;
3361 }
3362
3363 let registration = match registry.get_module(&spec.module_id) {
3364 Ok(registration) => registration,
3365 Err(err) => {
3366 warn!(module_id = %spec.module_id, error = %err, "health prober could not read registry");
3367 self.advertised = false;
3368 self.next_probe_at = None;
3369 return;
3370 }
3371 };
3372
3373 let Some(registration) = registration else {
3374 self.registered_connection = None;
3375 self.advertised = false;
3376 self.next_probe_at = None;
3377 return;
3378 };
3379
3380 let advertised = registration
3381 .control_ops
3382 .iter()
3383 .any(|op| op == MODULE_CONTROL_OP_HEALTH_CHECK);
3384 if !advertised {
3385 self.registered_connection = Some(registration.connection_id);
3386 self.advertised = false;
3387 self.next_probe_at = None;
3388 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3389 state.health.status = SupervisorHealthStatus::Unknown;
3390 state.health.consecutive_failures = 0;
3391 state.health.last_probe_ms = None;
3392 state.health.detail = None;
3393 state.health.metrics = None;
3394 });
3395 return;
3396 }
3397
3398 let reregistered = self.registered_connection != Some(registration.connection_id);
3399 self.registered_connection = Some(registration.connection_id);
3400 self.advertised = true;
3401 if reregistered || self.next_probe_at.is_none() {
3402 self.probe_index = 0;
3403 self.next_probe_at = Some(
3404 Instant::now() + jittered_health_delay(&spec.module_id, 0, runtime.health.cadence),
3405 );
3406 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3407 state.health.status = SupervisorHealthStatus::Unknown;
3408 state.health.consecutive_failures = 0;
3409 state.health.detail = None;
3410 state.health.metrics = None;
3411 });
3412 }
3413 }
3414
3415 fn wake_after(&self) -> Duration {
3416 if !self.advertised {
3417 return REGISTRY_RELEASE_POLL;
3418 }
3419 self.next_probe_at
3420 .map(|next| next.saturating_duration_since(Instant::now()))
3421 .unwrap_or(REGISTRY_RELEASE_POLL)
3422 }
3423
3424 fn due(&self) -> bool {
3425 self.advertised
3426 && self
3427 .next_probe_at
3428 .is_some_and(|next| Instant::now() >= next)
3429 }
3430
3431 fn schedule_next(&mut self, spec: &ModuleSpec, cadence: Duration) {
3432 self.probe_index = self.probe_index.wrapping_add(1);
3433 self.next_probe_at = Some(
3434 Instant::now() + jittered_health_delay(&spec.module_id, self.probe_index, cadence),
3435 );
3436 }
3437}
3438
3439#[derive(Debug)]
3474enum HealthProbeEvidence {
3475 LaneDead,
3477 NoAnswer,
3479 BadAnswer,
3481 Misconfigured,
3483}
3484
3485#[derive(Debug)]
3486struct HealthProbeError {
3487 evidence: HealthProbeEvidence,
3488 message: String,
3489}
3490
3491impl HealthProbeError {
3492 fn lane_dead(message: impl Into<String>) -> Self {
3493 Self::with(HealthProbeEvidence::LaneDead, message)
3494 }
3495
3496 fn no_answer(message: impl Into<String>) -> Self {
3497 Self::with(HealthProbeEvidence::NoAnswer, message)
3498 }
3499
3500 fn bad_answer(message: impl Into<String>) -> Self {
3501 Self::with(HealthProbeEvidence::BadAnswer, message)
3502 }
3503
3504 fn misconfigured(message: impl Into<String>) -> Self {
3505 Self::with(HealthProbeEvidence::Misconfigured, message)
3506 }
3507
3508 fn with(evidence: HealthProbeEvidence, message: impl Into<String>) -> Self {
3509 Self {
3510 evidence,
3511 message: message.into(),
3512 }
3513 }
3514
3515 #[allow(dead_code)]
3529 fn is_proof_of_death(&self) -> bool {
3530 matches!(self.evidence, HealthProbeEvidence::LaneDead)
3531 }
3532
3533 fn label(&self) -> &'static str {
3541 match self.evidence {
3542 HealthProbeEvidence::LaneDead => "lane-dead",
3543 HealthProbeEvidence::NoAnswer => "no-answer",
3544 HealthProbeEvidence::BadAnswer => "bad-answer",
3545 HealthProbeEvidence::Misconfigured => "daemon-misconfigured",
3546 }
3547 }
3548}
3549
3550impl fmt::Display for HealthProbeError {
3551 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3552 f.write_str(&self.message)
3553 }
3554}
3555
3556async fn run_health_probe_cycle(
3557 spec: &ModuleSpec,
3558 runtime: &SupervisorRuntimeConfig,
3559 registry: &Registry,
3560 process_liveness: &SupervisorProcessLiveness,
3561 snapshot: &SharedSnapshot,
3562 child: &mut Option<SupervisedChild>,
3563) {
3564 let now_ms = unix_ms_now();
3565 match probe_module_health(&spec.module_id, runtime, None).await {
3566 Ok(report) => {
3567 handle_health_report(
3568 spec,
3569 runtime,
3570 registry,
3571 process_liveness,
3572 snapshot,
3573 child,
3574 report,
3575 now_ms,
3576 )
3577 .await;
3578 }
3579 Err(err) => {
3580 handle_health_probe_failure(
3581 spec,
3582 runtime,
3583 registry,
3584 process_liveness,
3585 snapshot,
3586 child,
3587 err,
3588 now_ms,
3589 )
3590 .await;
3591 }
3592 }
3593}
3594
3595async fn probe_module_health(
3596 module_id: &str,
3597 runtime: &SupervisorRuntimeConfig,
3598 drain_deadline: Option<Instant>,
3599) -> Result<HealthReport, HealthProbeError> {
3600 let Some(forwarding) = runtime.forwarding.as_ref() else {
3601 return Err(HealthProbeError::misconfigured(
3602 "supervisor was not configured with a forwarding table",
3603 ));
3604 };
3605 let probe_started_at = Instant::now();
3606 let mut deadline = probe_started_at + runtime.health.deadline;
3607 if let Some(drain_deadline) = drain_deadline {
3608 deadline = deadline.min(drain_deadline);
3609 }
3610 let pending = if drain_deadline.is_some() {
3611 forwarding.begin_drain_health_probe_rpc_for(
3612 module_id,
3613 MODULE_CONTROL_OP_HEALTH_CHECK,
3614 probe_started_at,
3615 deadline,
3616 )
3617 } else {
3618 forwarding.begin_health_probe_rpc_for(
3619 module_id,
3620 MODULE_CONTROL_OP_HEALTH_CHECK,
3621 probe_started_at,
3622 deadline,
3623 )
3624 }
3625 .map_err(|err| {
3626 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3629 })?;
3630 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3631}
3632
3633async fn probe_endpoint_health(
3640 endpoint: crate::ModuleEndpointId,
3641 runtime: &SupervisorRuntimeConfig,
3642 deadline_cap: Option<Instant>,
3643) -> Result<HealthReport, HealthProbeError> {
3644 let Some(forwarding) = runtime.forwarding.as_ref() else {
3645 return Err(HealthProbeError::misconfigured(
3646 "supervisor was not configured with a forwarding table",
3647 ));
3648 };
3649 let probe_started_at = Instant::now();
3650 let mut deadline = probe_started_at + runtime.health.deadline;
3651 if let Some(cap) = deadline_cap {
3652 deadline = deadline.min(cap);
3653 }
3654 let pending = forwarding
3655 .begin_endpoint_health_probe_rpc_for(
3656 endpoint,
3657 MODULE_CONTROL_OP_HEALTH_CHECK,
3658 probe_started_at,
3659 deadline,
3660 )
3661 .map_err(|err| {
3662 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3663 })?;
3664 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3665}
3666
3667async fn await_health_probe(
3669 forwarding: &ForwardingTable,
3670 pending: PendingModuleControlRpc,
3671 deadline: Instant,
3672 probe_budget: Duration,
3673) -> Result<HealthReport, HealthProbeError> {
3674 let PendingModuleControlRpc {
3675 endpoint,
3676 module_sink,
3677 negotiated_ver,
3678 corr,
3679 receiver,
3680 } = pending;
3681 let body = serde_json::to_vec(&ModuleControlRequest::HealthCheck {}).map_err(|err| {
3682 HealthProbeError::misconfigured(format!("failed to encode health.check: {err}"))
3683 })?;
3684 let frame = Frame::build_with_version(
3685 negotiated_ver,
3686 FrameType::Request,
3687 control_flags(),
3688 0,
3689 0,
3690 corr,
3691 body,
3692 )
3693 .map_err(|err| {
3694 HealthProbeError::misconfigured(format!("failed to build health.check frame: {err}"))
3695 })?;
3696
3697 match timeout_at(deadline, module_sink.send(frame)).await {
3703 Ok(Ok(())) => {}
3704 Ok(Err(err)) => {
3705 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3706 return Err(HealthProbeError::lane_dead(format!(
3709 "failed to send health.check: {err}"
3710 )));
3711 }
3712 Err(_elapsed) => {
3713 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3714 return Err(HealthProbeError::no_answer(
3718 "health.check send timed out before enqueue (module egress full)",
3719 ));
3720 }
3721 }
3722
3723 match timeout_at(deadline, receiver).await {
3724 Ok(Ok(ModuleControlRpcOutcome::Response(response))) => {
3728 response.health_report().ok_or_else(|| {
3729 HealthProbeError::bad_answer("health.check RPC returned a non-health response")
3730 })
3731 }
3732 Ok(Ok(ModuleControlRpcOutcome::Rejected(body))) => Err(HealthProbeError::bad_answer(
3733 format!("health.check rejected: {}", body.message),
3734 )),
3735 Ok(Ok(ModuleControlRpcOutcome::ModuleGone(message))) => {
3736 Err(HealthProbeError::lane_dead(message))
3737 }
3738 Ok(Ok(ModuleControlRpcOutcome::MalformedResponse(message))) => {
3739 Err(HealthProbeError::bad_answer(message))
3740 }
3741 Ok(Ok(ModuleControlRpcOutcome::UnexpectedOp { expected, actual })) => {
3742 Err(HealthProbeError::bad_answer(format!(
3743 "expected module-control op '{expected}', got '{actual}'"
3744 )))
3745 }
3746 Ok(Ok(ModuleControlRpcOutcome::DeadlineElapsed)) => Err(HealthProbeError::bad_answer(
3750 "module answered health.check after its daemon deadline",
3751 )),
3752 Ok(Err(_)) => Err(HealthProbeError::misconfigured(
3753 "health.check waiter was canceled before the module responded",
3754 )),
3755 Err(_) => {
3756 let _ = forwarding.tombstone_health_probe_rpc(endpoint, corr);
3757 Err(HealthProbeError::no_answer(format!(
3758 "module did not answer health.check within {probe_budget:?}"
3759 )))
3760 }
3761 }
3762}
3763
3764#[allow(clippy::too_many_arguments)]
3765async fn handle_health_report(
3766 spec: &ModuleSpec,
3767 runtime: &SupervisorRuntimeConfig,
3768 registry: &Registry,
3769 process_liveness: &SupervisorProcessLiveness,
3770 snapshot: &SharedSnapshot,
3771 child: &mut Option<SupervisedChild>,
3772 report: HealthReport,
3773 now_ms: u64,
3774) {
3775 let status = supervisor_health_status(report.status);
3776 let detail = report.detail.clone();
3777 let metrics = truncate_health_metrics(report.metrics);
3778 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3779 state.health.status = status;
3780 state.health.last_probe_ms = Some(now_ms);
3781 state.health.detail = detail.clone();
3782 state.health.metrics = metrics.clone();
3783 state.health.consecutive_failures = 0;
3784 });
3785
3786 let action = match report.status {
3787 HealthStatus::Ok => return,
3788 HealthStatus::Degraded => runtime.health.on_degraded,
3789 HealthStatus::Failing => runtime.health.on_failing,
3790 };
3791 apply_l3_health_action(
3792 spec,
3793 runtime,
3794 registry,
3795 process_liveness,
3796 snapshot,
3797 child,
3798 status,
3799 detail.as_deref(),
3800 action,
3801 now_ms,
3802 )
3803 .await;
3804}
3805
3806#[allow(clippy::too_many_arguments)]
3807async fn handle_health_probe_failure(
3808 spec: &ModuleSpec,
3809 runtime: &SupervisorRuntimeConfig,
3810 registry: &Registry,
3811 process_liveness: &SupervisorProcessLiveness,
3812 snapshot: &SharedSnapshot,
3813 child: &mut Option<SupervisedChild>,
3814 err: HealthProbeError,
3815 now_ms: u64,
3816) {
3817 let threshold = runtime.health.failure_threshold.max(1);
3818 let mut failures = 0;
3819 let detail = format!("[{}] {err}", err.label());
3824 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3825 state.health.last_probe_ms = Some(now_ms);
3826 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
3827 state.health.detail = Some(detail.clone());
3828 state.health.metrics = None;
3829 failures = state.health.consecutive_failures;
3830 });
3831
3832 if failures < threshold {
3833 warn!(
3834 module_id = %spec.module_id,
3835 consecutive_failures = failures,
3836 threshold,
3837 evidence = err.label(),
3838 detail = %detail,
3839 "health.check probe failed"
3840 );
3841 return;
3842 }
3843
3844 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3845 state.state = ModuleState::Unresponsive;
3846 state.health.status = SupervisorHealthStatus::Unresponsive;
3847 });
3848 if runtime.health.critical {
3852 error!(
3853 module_id = %spec.module_id,
3854 status = "unresponsive",
3855 evidence = err.label(),
3856 detail = %detail,
3857 "critical module health alert"
3858 );
3859 } else {
3860 warn!(
3861 module_id = %spec.module_id,
3862 status = "unresponsive",
3863 evidence = err.label(),
3864 detail = %detail,
3865 "module health threshold breached"
3866 );
3867 }
3868 if let Err(err) = health_restart_child(
3869 spec,
3870 runtime,
3871 registry,
3872 process_liveness,
3873 snapshot,
3874 child,
3875 SupervisorHealthStatus::Unresponsive,
3876 Some(&detail),
3877 now_ms,
3878 )
3879 .await
3880 {
3881 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3882 }
3883}
3884
3885#[allow(clippy::too_many_arguments)]
3886async fn apply_l3_health_action(
3887 spec: &ModuleSpec,
3888 runtime: &SupervisorRuntimeConfig,
3889 registry: &Registry,
3890 process_liveness: &SupervisorProcessLiveness,
3891 snapshot: &SharedSnapshot,
3892 child: &mut Option<SupervisedChild>,
3893 status: SupervisorHealthStatus,
3894 detail: Option<&str>,
3895 action: HealthAction,
3896 now_ms: u64,
3897) {
3898 record_health_action(snapshot, &spec.module_id, action.to_string(), now_ms);
3899 match action {
3900 HealthAction::Report => {
3901 info!(
3902 module_id = %spec.module_id,
3903 status = ?status,
3904 detail,
3905 "module reported non-ok health"
3906 );
3907 }
3908 HealthAction::Alert => {
3909 error!(
3910 module_id = %spec.module_id,
3911 status = ?status,
3912 detail,
3913 "module health alert"
3914 );
3915 }
3916 HealthAction::Restart => {
3917 if let Err(err) = health_restart_child(
3918 spec,
3919 runtime,
3920 registry,
3921 process_liveness,
3922 snapshot,
3923 child,
3924 status,
3925 detail,
3926 now_ms,
3927 )
3928 .await
3929 {
3930 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3931 }
3932 }
3933 }
3934}
3935
3936#[allow(clippy::too_many_arguments)]
3937async fn health_restart_child(
3938 spec: &ModuleSpec,
3939 runtime: &SupervisorRuntimeConfig,
3940 registry: &Registry,
3941 process_liveness: &SupervisorProcessLiveness,
3942 snapshot: &SharedSnapshot,
3943 child: &mut Option<SupervisedChild>,
3944 status: SupervisorHealthStatus,
3945 detail: Option<&str>,
3946 now_ms: u64,
3947) -> Result<(), SuperviseError> {
3948 let (enabled, schedule) = {
3949 let mut state = lock_snapshot(snapshot)?;
3950 let enabled = state.enabled;
3951 let schedule = if enabled {
3952 state.next_crash_restart(&runtime.restart_policy, Instant::now())
3953 } else {
3954 None
3955 };
3956 (enabled, schedule)
3957 };
3958
3959 if !enabled {
3960 return Err(SuperviseError::Disabled {
3961 module_id: spec.module_id.clone(),
3962 });
3963 }
3964
3965 if schedule.is_none() {
3966 record_health_action(snapshot, &spec.module_id, "disabled".to_string(), now_ms);
3967 error!(
3968 module_id = %spec.module_id,
3969 status = ?status,
3970 detail,
3971 max_restarts = runtime.restart_policy.max_restarts,
3972 window_secs = runtime.restart_policy.window.as_secs(),
3973 "health restart budget exhausted; disabling module"
3974 );
3975 let stop_notice = begin_forwarding_drain_if_configured(
3976 spec,
3977 runtime,
3978 registry,
3979 snapshot,
3980 Some(false),
3981 RouteCloseReason::Disable,
3982 )
3983 .await?;
3984 drain_optional_child(
3985 &spec.module_id,
3986 spec.protocol,
3987 stop_notice,
3988 registry,
3989 snapshot,
3990 &runtime.terminal_ring,
3991 &runtime.spawn_events,
3992 child,
3993 runtime.drain_timeout,
3994 ModuleState::Disabled,
3995 Some(false),
3996 )
3997 .await?;
3998 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3999 return Ok(());
4000 }
4001
4002 let schedule = schedule.expect("a health restart must have a crash-restart schedule");
4003 let mut restart_count = 0;
4004 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4005 restart_count = state.crash_restarts.len();
4006 state.state = ModuleState::Unresponsive;
4007 state.health.status = status;
4008 state.health.last_action = Some(HealthAction::Restart.to_string());
4009 state.health.last_action_ms = Some(now_ms);
4010 })?;
4011 warn!(
4012 module_id = %spec.module_id,
4013 status = ?status,
4014 detail,
4015 restart_count,
4016 restart_in_window = schedule.restart_in_window,
4017 delay_ms = schedule.delay.as_millis() as u64,
4018 "health-triggered module restart"
4019 );
4020
4021 let stop_notice = begin_forwarding_drain_if_configured(
4022 spec,
4023 runtime,
4024 registry,
4025 snapshot,
4026 Some(true),
4027 RouteCloseReason::Restart,
4028 )
4029 .await?;
4030 drain_optional_child(
4031 &spec.module_id,
4032 spec.protocol,
4033 stop_notice,
4034 registry,
4035 snapshot,
4036 &runtime.terminal_ring,
4037 &runtime.spawn_events,
4038 child,
4039 runtime.drain_timeout,
4040 ModuleState::Restarting,
4041 Some(true),
4042 )
4043 .await?;
4044 sleep(schedule.delay).await;
4045 if !respawn_still_pending(snapshot) {
4049 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4050 return Ok(());
4051 }
4052 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4053 match spawn_and_mark_running(spec, runtime, snapshot) {
4054 Ok(next_child) => {
4055 *child = Some(next_child);
4056 Ok(())
4057 }
4058 Err(err) => {
4059 fail_snapshot(snapshot, Some(&spec.module_id), None);
4060 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4061 *child = None;
4062 Err(err)
4063 }
4064 }
4065}
4066
4067fn record_health_action(snapshot: &SharedSnapshot, module_id: &str, action: String, now_ms: u64) {
4068 let _ = update_snapshot(snapshot, Some(module_id), |state| {
4069 state.health.last_action = Some(action);
4070 state.health.last_action_ms = Some(now_ms);
4071 });
4072}
4073
4074fn supervisor_health_status(status: HealthStatus) -> SupervisorHealthStatus {
4075 match status {
4076 HealthStatus::Ok => SupervisorHealthStatus::Ok,
4077 HealthStatus::Degraded => SupervisorHealthStatus::Degraded,
4078 HealthStatus::Failing => SupervisorHealthStatus::Failing,
4079 }
4080}
4081
4082fn truncate_health_metrics(metrics: Option<Value>) -> Option<Value> {
4094 let metrics = metrics?;
4095 match serde_json::to_vec(&metrics) {
4096 Ok(encoded) if encoded.len() > MAX_HEALTH_METRICS_BYTES => Some(serde_json::json!({
4097 "truncated": true,
4098 "original_bytes": encoded.len(),
4099 })),
4100 Ok(_) | Err(_) => Some(metrics),
4101 }
4102}
4103
4104fn jittered_health_delay(module_id: &str, probe_index: u64, cadence: Duration) -> Duration {
4110 if cadence.is_zero() {
4111 return Duration::ZERO;
4112 }
4113 let cadence_ms = cadence.as_millis() as u64;
4114 if cadence_ms == 0 {
4130 return cadence;
4131 }
4132 let jitter_span = (cadence_ms / 10).max(1);
4147 let hash = module_id.as_bytes().iter().fold(
4148 probe_index.wrapping_mul(0x9E37_79B9_7F4A_7C15),
4149 |acc, byte| {
4150 acc.wrapping_mul(1099511628211)
4151 .wrapping_add(u64::from(*byte))
4152 },
4153 );
4154 cadence + Duration::from_millis(hash % jitter_span)
4155}
4156
4157#[cfg(test)]
4158mod tests {
4159 use super::*;
4160
4161 #[test]
4162 fn readding_a_module_clears_its_rescan_removal_tombstone() {
4163 let handle = SupervisorHandle::new();
4164 let module_id = "readded-tombstone";
4165 handle.record_rescan_removal(module_id);
4166 assert!(handle.removal_tombstone_age_ms(module_id).is_some());
4167
4168 handle.apply_identity_configuration(&ModuleSpec {
4169 module_id: module_id.to_string(),
4170 program: PathBuf::from("/test/module"),
4171 args: Vec::new(),
4172 env: Vec::new(),
4173 reserved: false,
4174 reserved_prefixes: Vec::new(),
4175 protocol: ModuleProtocol::Subc,
4176 overlap: Default::default(),
4177 });
4178
4179 assert!(
4180 handle.removal_tombstone_age_ms(module_id).is_none(),
4181 "a re-added module must not retain a stale removal tombstone"
4182 );
4183 }
4184
4185 fn stale_process_snapshot(state: ModuleState, enabled: bool) -> SharedSnapshot {
4186 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::new(state, enabled)));
4187 update_snapshot(&snapshot, Some("stale-process-facts"), |snapshot| {
4188 snapshot.process_alive = true;
4189 snapshot.pid = Some(41);
4190 snapshot.spawned_at_ms = Some(42);
4191 snapshot.spawned_from = Some(PathBuf::from("/spawned/module"));
4192 snapshot.spawned_file_identity = Some(SpawnedFileIdentity {
4193 device: 43,
4194 inode: 44,
4195 });
4196 })
4197 .unwrap();
4198 snapshot
4199 }
4200
4201 fn assert_snapshot_process_facts_cleared(snapshot: &SharedSnapshot) {
4202 let snapshot = lock_snapshot(snapshot).unwrap();
4203 assert!(!snapshot.process_alive);
4204 assert_eq!(snapshot.pid, None);
4205 assert_eq!(snapshot.spawned_at_ms, None);
4206 assert_eq!(snapshot.spawned_from, None);
4207 assert_eq!(snapshot.spawned_file_identity, None);
4208 }
4209
4210 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4211 async fn failed_enable_spawn_clears_preexisting_current_process_facts() {
4212 let supervisor = Supervisor::default();
4213 let mut runtime = supervisor.runtime_config();
4214 runtime.test_seed_stale_facts_before_enable_spawn = true;
4215 let snapshot = stale_process_snapshot(ModuleState::Disabled, false);
4216 let mut child = None;
4217 let spec = ModuleSpec {
4218 module_id: "failed-enable-clears-facts".to_string(),
4219 program: PathBuf::from("/definitely/missing/failed-enable-module"),
4220 args: Vec::new(),
4221 env: Vec::new(),
4222 reserved: false,
4223 reserved_prefixes: Vec::new(),
4224 protocol: ModuleProtocol::Subc,
4225 overlap: Default::default(),
4226 };
4227
4228 let result = set_child_enabled(
4229 &spec,
4230 &runtime,
4231 &supervisor.registry,
4232 &supervisor.process_liveness,
4233 &snapshot,
4234 &mut child,
4235 true,
4236 )
4237 .await;
4238
4239 assert!(matches!(result, Err(SuperviseError::Spawn { .. })));
4240 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4241 assert_snapshot_process_facts_cleared(&snapshot);
4242 }
4243
4244 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4245 async fn failed_reload_spawn_clears_current_process_facts() {
4246 let supervisor = Supervisor::default();
4247 let mut runtime = supervisor.runtime_config();
4248 runtime.restart_policy = RestartPolicy::new(0, Duration::ZERO);
4249 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4250 let mut child = None;
4251 let spec = ModuleSpec {
4252 module_id: "failed-reload-clears-facts".to_string(),
4253 program: PathBuf::from("/unused/failed-reload-module"),
4254 args: Vec::new(),
4255 env: Vec::new(),
4256 reserved: false,
4257 reserved_prefixes: Vec::new(),
4258 protocol: ModuleProtocol::Subc,
4259 overlap: Default::default(),
4260 };
4261
4262 let result = handle_reload_spawn_failure(
4263 &spec,
4264 &runtime,
4265 &supervisor.process_liveness,
4266 &snapshot,
4267 &mut child,
4268 "forced reload spawn failure".to_string(),
4269 )
4270 .await;
4271
4272 assert!(matches!(result, Err(SuperviseError::ReloadFailed { .. })));
4273 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4274 assert_snapshot_process_facts_cleared(&snapshot);
4275 }
4276
4277 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4278 async fn dropping_a_module_with_an_active_monitor_clears_current_process_facts() {
4279 let supervisor = Supervisor::default();
4280 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4281 let module = supervisor.supervised_module(
4282 ModuleSpec {
4283 module_id: "drop-clears-facts".to_string(),
4284 program: PathBuf::from("/unused/drop-module"),
4285 args: Vec::new(),
4286 env: Vec::new(),
4287 reserved: false,
4288 reserved_prefixes: Vec::new(),
4289 protocol: ModuleProtocol::Subc,
4290 overlap: Default::default(),
4291 },
4292 supervisor.runtime_config(),
4293 Arc::clone(&snapshot),
4294 None,
4295 );
4296 assert!(!module
4297 .inner
4298 .monitor
4299 .lock()
4300 .unwrap()
4301 .as_ref()
4302 .unwrap()
4303 .is_finished());
4304
4305 drop(module);
4306
4307 assert_eq!(
4308 lock_snapshot(&snapshot).unwrap().state,
4309 ModuleState::Stopped
4310 );
4311 assert_snapshot_process_facts_cleared(&snapshot);
4312 }
4313
4314 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4315 async fn configuration_update_does_not_replace_captured_running_process_facts() {
4316 let supervisor = Supervisor::default();
4317 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4318 let initial = ModuleSpec {
4319 module_id: "rescan-preserves-spawn-facts".to_string(),
4320 program: PathBuf::from("/spawned/module"),
4321 args: Vec::new(),
4322 env: Vec::new(),
4323 reserved: false,
4324 reserved_prefixes: Vec::new(),
4325 protocol: ModuleProtocol::Subc,
4326 overlap: Default::default(),
4327 };
4328 let module = supervisor.supervised_module(
4329 initial.clone(),
4330 supervisor.runtime_config(),
4331 snapshot,
4332 None,
4333 );
4334 let before = module.status().unwrap();
4335 let mut replacement = initial;
4336 replacement.program = PathBuf::from("/rescanned/replacement-module");
4337
4338 module
4339 .update_configuration(replacement, HealthConfig::default(), None)
4340 .await
4341 .unwrap();
4342
4343 let after = module.status().unwrap();
4344 assert_eq!(after.pid, before.pid);
4345 assert_eq!(after.spawned_at_ms, before.spawned_at_ms);
4346 assert_eq!(after.spawned_from, before.spawned_from);
4347 drop(module);
4348 }
4349}
4350
4351fn unix_ms_now() -> u64 {
4352 SystemTime::now()
4353 .duration_since(UNIX_EPOCH)
4354 .map(|duration| duration.as_millis().min(u128::from(u64::MAX)) as u64)
4355 .unwrap_or(0)
4356}
4357
4358async fn supervise_loop(
4359 mut spec: ModuleSpec,
4360 mut runtime: SupervisorRuntimeConfig,
4361 registry: Arc<Registry>,
4362 process_liveness: Arc<SupervisorProcessLiveness>,
4363 snapshot: SharedSnapshot,
4364 mut child: Option<SupervisedChild>,
4365 mut commands: mpsc::Receiver<SupervisorCommand>,
4366) {
4367 let mut health_probe = HealthProbeRuntime::default();
4368 let mut pending_respawn: Option<Instant> = None;
4372 let mut requeued: VecDeque<SupervisorCommand> = VecDeque::new();
4375 loop {
4376 if let Some(command) = requeued.pop_front() {
4377 if !handle_supervisor_command(
4378 command,
4379 &mut spec,
4380 &mut runtime,
4381 ®istry,
4382 &process_liveness,
4383 &snapshot,
4384 &mut child,
4385 &mut commands,
4386 &mut requeued,
4387 )
4388 .await
4389 {
4390 return;
4391 }
4392 if child.is_some() || !respawn_still_pending(&snapshot) {
4393 pending_respawn = None;
4394 }
4395 continue;
4396 }
4397 if child.is_some() {
4398 health_probe.refresh_registration(&spec, &runtime, ®istry, &snapshot);
4399 let probe_sleep = sleep(health_probe.wake_after());
4400 tokio::pin!(probe_sleep);
4401 let active_child = child.as_mut().expect("child checked above");
4402 tokio::select! {
4403 wait_result = active_child.wait() => {
4404 let exit_report = match wait_result {
4413 Ok(status) => classify_reaped_child_exit(&snapshot, active_child, &status),
4414 Err(err) => {
4415 active_child.drain_stderr(&spec.module_id).await;
4416 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4417 record_wait_error_terminal(
4423 &spec.module_id,
4424 &runtime.terminal_ring,
4425 &runtime.spawn_events,
4426 );
4427 untrack_if_registration_released(
4428 &process_liveness,
4429 ®istry,
4430 &spec.module_id,
4431 &snapshot,
4432 );
4433 error!(module_id = %spec.module_id, error = %err, "failed to wait for supervised module");
4434 child = None;
4435 continue;
4436 }
4437 };
4438 active_child.drain_stderr(&spec.module_id).await;
4439
4440 let next = on_child_exit(
4441 &spec,
4442 runtime.restart_policy,
4443 ®istry,
4444 &snapshot,
4445 &runtime.terminal_ring,
4446 &runtime.spawn_events,
4447 &runtime.child_roster,
4448 exit_report,
4449 ).await;
4450 active_child.release_roster();
4453 match next {
4454 NextAction::Stop { registration_released } => {
4455 if registration_released {
4456 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4457 }
4458 child = None;
4459 }
4460 NextAction::Restart { schedule } => {
4461 let delay = schedule.map_or(
4462 runtime.restart_policy.delay_for_restart(0),
4463 |schedule| schedule.delay,
4464 );
4465 if let Some(schedule) = schedule {
4466 log_crash_respawn(&spec.module_id, schedule);
4467 }
4468 child = None;
4476 pending_respawn = Some(Instant::now() + delay);
4477 }
4478 }
4479 }
4480 command = commands.recv() => {
4481 let Some(command) = command else {
4482 return;
4483 };
4484 if !handle_supervisor_command(
4485 command,
4486 &mut spec,
4487 &mut runtime,
4488 ®istry,
4489 &process_liveness,
4490 &snapshot,
4491 &mut child,
4492 &mut commands,
4493 &mut requeued,
4494 ).await {
4495 return;
4496 }
4497 }
4498 _ = &mut probe_sleep => {
4499 if health_probe.due() {
4500 run_health_probe_cycle(
4501 &spec,
4502 &runtime,
4503 ®istry,
4504 &process_liveness,
4505 &snapshot,
4506 &mut child,
4507 ).await;
4508 if child.is_some() {
4509 health_probe.schedule_next(&spec, runtime.health.cadence);
4510 }
4511 }
4512 }
4513 }
4514 } else if let Some(deadline) = pending_respawn {
4515 tokio::select! {
4516 _ = sleep_until(deadline) => {
4517 pending_respawn = None;
4518 if !respawn_still_pending(&snapshot) {
4522 continue;
4523 }
4524 if runtime.child_roster.is_closed() {
4529 let _ = update_snapshot(&snapshot, Some(&spec.module_id), |state| {
4530 state.state = ModuleState::Stopped;
4531 });
4532 debug!(module_id = %spec.module_id, "crash respawn cancelled by daemon shutdown");
4533 continue;
4534 }
4535 if let Err(err) = wait_for_registration_release(
4536 ®istry,
4537 &spec.module_id,
4538 REGISTRY_RELEASE_TIMEOUT,
4539 ).await {
4540 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4541 error!(module_id = %spec.module_id, error = %err, "registration did not release before restart");
4542 continue;
4543 }
4544
4545 match spawn_and_mark_running(&spec, &runtime, &snapshot) {
4546 Ok(next_child) => {
4547 child = Some(next_child);
4548 debug!(module_id = %spec.module_id, "supervised module restarted after crash");
4549 }
4550 Err(err) => {
4551 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4552 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4553 error!(module_id = %spec.module_id, error = %err, "failed to restart supervised module");
4554 }
4555 }
4556 }
4557 command = commands.recv() => {
4558 let Some(command) = command else {
4559 return;
4560 };
4561 if !handle_supervisor_command(
4562 command,
4563 &mut spec,
4564 &mut runtime,
4565 ®istry,
4566 &process_liveness,
4567 &snapshot,
4568 &mut child,
4569 &mut commands,
4570 &mut requeued,
4571 ).await {
4572 return;
4573 }
4574 if child.is_some() || !respawn_still_pending(&snapshot) {
4579 pending_respawn = None;
4580 }
4581 }
4582 }
4583 } else {
4584 let Some(command) = commands.recv().await else {
4585 return;
4586 };
4587 if !handle_supervisor_command(
4588 command,
4589 &mut spec,
4590 &mut runtime,
4591 ®istry,
4592 &process_liveness,
4593 &snapshot,
4594 &mut child,
4595 &mut commands,
4596 &mut requeued,
4597 )
4598 .await
4599 {
4600 return;
4601 }
4602 }
4603 }
4604}
4605
4606fn log_crash_respawn(module_id: &str, schedule: CrashRestartSchedule) {
4607 info!(
4608 module_id,
4609 restart_in_window = schedule.restart_in_window,
4610 delay_ms = schedule.delay.as_millis() as u64,
4611 "respawning after crash"
4612 );
4613}
4614
4615fn respawn_still_pending(snapshot: &SharedSnapshot) -> bool {
4621 matches!(
4622 lock_snapshot(snapshot),
4623 Ok(state) if state.enabled && state.state == ModuleState::Restarting
4624 )
4625}
4626
4627enum NextAction {
4628 Stop {
4629 registration_released: bool,
4630 },
4631 Restart {
4632 schedule: Option<CrashRestartSchedule>,
4633 },
4634}
4635
4636#[allow(clippy::too_many_arguments)]
4637async fn handle_supervisor_command(
4638 command: SupervisorCommand,
4639 spec: &mut ModuleSpec,
4640 runtime: &mut SupervisorRuntimeConfig,
4641 registry: &Registry,
4642 process_liveness: &SupervisorProcessLiveness,
4643 snapshot: &SharedSnapshot,
4644 child: &mut Option<SupervisedChild>,
4645 commands: &mut mpsc::Receiver<SupervisorCommand>,
4646 requeued: &mut VecDeque<SupervisorCommand>,
4647) -> bool {
4648 match command {
4649 SupervisorCommand::Drain { reply } => {
4650 let result = drain_optional_child(
4653 &spec.module_id,
4654 spec.protocol,
4655 StopNotice::NotSent,
4656 registry,
4657 snapshot,
4658 &runtime.terminal_ring,
4659 &runtime.spawn_events,
4660 child,
4661 runtime.drain_timeout,
4662 ModuleState::Stopped,
4663 None,
4664 )
4665 .await;
4666 let registration_released = result.is_ok();
4667 let _ = reply.send(result);
4668 if registration_released {
4669 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4670 }
4671 false
4672 }
4673 SupervisorCommand::Retire { reply } => {
4674 let result = async {
4675 let stop_notice = begin_forwarding_drain_if_configured(
4676 spec,
4677 runtime,
4678 registry,
4679 snapshot,
4680 None,
4681 RouteCloseReason::Disable,
4682 )
4683 .await?;
4684 drain_optional_child(
4685 &spec.module_id,
4686 spec.protocol,
4687 stop_notice,
4688 registry,
4689 snapshot,
4690 &runtime.terminal_ring,
4691 &runtime.spawn_events,
4692 child,
4693 runtime.drain_timeout,
4694 ModuleState::Stopped,
4695 None,
4696 )
4697 .await
4698 }
4699 .await;
4700 let registration_released = result.is_ok();
4701 let _ = reply.send(result);
4702 if registration_released {
4703 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4704 }
4705 false
4706 }
4707 SupervisorCommand::Restart {
4708 drain_timeout_ms,
4709 received_at_generation,
4710 queued_at,
4711 reply,
4712 } => {
4713 info!(
4717 module_id = %spec.module_id,
4718 queued_ms = u64::try_from(queued_at.elapsed().as_millis()).unwrap_or(u64::MAX),
4719 "restart command dequeued"
4720 );
4721 let validation = match lock_snapshot(snapshot) {
4733 Ok(state) if !state.enabled => Err(SuperviseError::Disabled {
4734 module_id: spec.module_id.clone(),
4735 }),
4736 Ok(_) => Ok(()),
4737 Err(err) => Err(err),
4738 };
4739 let initiated = validation.is_ok();
4740 let _ = reply.send(validation);
4741 let satisfied_by_generation = if initiated && child.is_some() {
4752 lock_snapshot(snapshot).ok().and_then(|state| {
4753 (state.spawn_generation > received_at_generation
4754 && !state.configuration_updated_since_spawn)
4755 .then_some(state.spawn_generation)
4756 })
4757 } else {
4758 None
4759 };
4760 if let Some(generation) = satisfied_by_generation {
4761 info!(
4762 module_id = %spec.module_id,
4763 received_at_generation,
4764 "restart already satisfied by generation {generation}; not restarting again"
4765 );
4766 } else if initiated {
4767 let drain_timeout = drain_timeout_ms
4770 .map(Duration::from_millis)
4771 .unwrap_or(runtime.drain_timeout);
4772 if let Err(err) = restart_child(
4773 spec,
4774 runtime,
4775 registry,
4776 process_liveness,
4777 snapshot,
4778 child,
4779 drain_timeout,
4780 )
4781 .await
4782 {
4783 warn!(
4784 module_id = %spec.module_id,
4785 error = %err,
4786 "operator restart failed after initiation ack; module state carries the outcome"
4787 );
4788 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4789 state.state = ModuleState::Failed;
4790 clear_current_process_facts(state);
4791 });
4792 }
4793 }
4794 true
4795 }
4796 SupervisorCommand::Reload { reply } => {
4797 let result =
4798 reload_child(spec, runtime, registry, process_liveness, snapshot, child).await;
4799 let _ = reply.send(result);
4800 true
4801 }
4802 SupervisorCommand::SetEnabled { enabled, reply } => {
4803 let result = set_child_enabled(
4804 spec,
4805 runtime,
4806 registry,
4807 process_liveness,
4808 snapshot,
4809 child,
4810 enabled,
4811 )
4812 .await;
4813 let _ = reply.send(result);
4814 true
4815 }
4816 SupervisorCommand::UpdateConfiguration {
4817 spec: next_spec,
4818 health,
4819 drain_timeout_ms,
4820 reply,
4821 } => {
4822 if let Some(handle) = &runtime.supervisor_handle {
4823 handle.apply_identity_configuration(&next_spec);
4824 }
4825 *spec = next_spec;
4826 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4827 state.configuration_updated_since_spawn = true;
4828 });
4829 runtime.health = health;
4830 runtime.drain_timeout = drain_timeout_ms
4831 .map(Duration::from_millis)
4832 .unwrap_or(runtime.default_drain_timeout);
4833 *runtime
4834 .effective_drain_timeout
4835 .lock()
4836 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
4837 let _ = reply.send(());
4838 true
4839 }
4840 SupervisorCommand::Swap {
4841 ready_timeout,
4842 reply,
4843 } => {
4844 let end = swap::run_swap(
4845 spec,
4846 runtime,
4847 registry,
4848 process_liveness,
4849 snapshot,
4850 child,
4851 commands,
4852 ready_timeout.unwrap_or(DEFAULT_SWAP_READY_TIMEOUT),
4853 reply,
4854 )
4855 .await;
4856 requeued.extend(end.requeue);
4857 true
4858 }
4859 }
4860}
4861
4862async fn restart_child(
4863 spec: &ModuleSpec,
4864 runtime: &SupervisorRuntimeConfig,
4865 registry: &Registry,
4866 process_liveness: &SupervisorProcessLiveness,
4867 snapshot: &SharedSnapshot,
4868 child: &mut Option<SupervisedChild>,
4869 drain_timeout: Duration,
4870) -> Result<(), SuperviseError> {
4871 if !lock_snapshot(snapshot)?.enabled {
4873 return Err(SuperviseError::Disabled {
4874 module_id: spec.module_id.clone(),
4875 });
4876 }
4877 let stop_notice = begin_forwarding_drain_with_timeout(
4878 spec,
4879 runtime,
4880 registry,
4881 snapshot,
4882 None,
4883 RouteCloseReason::Restart,
4884 drain_timeout,
4885 )
4886 .await?;
4887
4888 if child.is_some() {
4889 drain_optional_child(
4890 &spec.module_id,
4891 spec.protocol,
4892 stop_notice,
4893 registry,
4894 snapshot,
4895 &runtime.terminal_ring,
4896 &runtime.spawn_events,
4897 child,
4898 drain_timeout,
4899 ModuleState::Restarting,
4900 Some(true),
4901 )
4902 .await?;
4903 } else {
4904 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4905 state.enabled = true;
4906 state.state = ModuleState::Restarting;
4907 clear_current_process_facts(state);
4908 })?;
4909 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4910 }
4911
4912 reset_restart_count(snapshot, &spec.module_id)?;
4913 sleep(runtime.restart_policy.backoff).await;
4914 if !respawn_still_pending(snapshot) {
4917 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4918 return Ok(());
4919 }
4920 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4921 match spawn_and_mark_running(spec, runtime, snapshot) {
4927 Ok(next_child) => {
4928 *child = Some(next_child);
4929 debug!(module_id = %spec.module_id, "supervised module restarted by operator request");
4930 Ok(())
4931 }
4932 Err(err) => {
4933 fail_snapshot(snapshot, Some(&spec.module_id), None);
4934 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4935 *child = None;
4936 Err(err)
4937 }
4938 }
4939}
4940
4941async fn reload_child(
4942 spec: &ModuleSpec,
4943 runtime: &SupervisorRuntimeConfig,
4944 registry: &Registry,
4945 process_liveness: &SupervisorProcessLiveness,
4946 snapshot: &SharedSnapshot,
4947 child: &mut Option<SupervisedChild>,
4948) -> Result<(), SuperviseError> {
4949 if !lock_snapshot(snapshot)?.enabled {
4951 return Err(SuperviseError::Disabled {
4952 module_id: spec.module_id.clone(),
4953 });
4954 }
4955 let stop_notice = begin_forwarding_drain(
4956 spec,
4957 runtime,
4958 registry,
4959 snapshot,
4960 Some(true),
4961 RouteCloseReason::Reload,
4962 )
4963 .await?;
4964
4965 if child.is_some() {
4966 drain_optional_child(
4967 &spec.module_id,
4968 spec.protocol,
4969 stop_notice,
4970 registry,
4971 snapshot,
4972 &runtime.terminal_ring,
4973 &runtime.spawn_events,
4974 child,
4975 runtime.drain_timeout,
4976 ModuleState::Restarting,
4977 Some(true),
4978 )
4979 .await?;
4980 } else {
4981 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4982 state.enabled = true;
4983 state.state = ModuleState::Restarting;
4984 clear_current_process_facts(state);
4985 })?;
4986 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4987 }
4988
4989 reset_restart_count(snapshot, &spec.module_id)?;
4990 sleep(runtime.restart_policy.backoff).await;
4991 if !respawn_still_pending(snapshot) {
4994 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4995 return Ok(());
4996 }
4997 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4998 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
4999 Ok(next_child) => next_child,
5000 Err(err) => {
5001 return handle_reload_spawn_failure(
5002 spec,
5003 runtime,
5004 process_liveness,
5005 snapshot,
5006 child,
5007 format!("new child failed to spawn: {err}"),
5008 )
5009 .await;
5010 }
5011 };
5012 *child = Some(next_child);
5013
5014 let wait_outcome = {
5015 let active_child = child.as_mut().expect("new reload child was just stored");
5016 wait_for_registration_after_reload(
5017 registry,
5018 &spec.module_id,
5019 snapshot,
5020 active_child,
5021 REGISTRY_RELEASE_TIMEOUT,
5022 )
5023 .await?
5024 };
5025
5026 match wait_outcome {
5027 RegistrationWaitOutcome::Registered => {
5028 debug!(module_id = %spec.module_id, "supervised module reloaded and registered");
5029 Ok(())
5030 }
5031 RegistrationWaitOutcome::Exited(exit_report) => {
5032 if let Some(active_child) = child.as_mut() {
5033 active_child.drain_stderr(&spec.module_id).await;
5034 }
5035 *child = None;
5036 handle_reload_child_registration_failure(
5037 spec,
5038 runtime,
5039 registry,
5040 process_liveness,
5041 snapshot,
5042 child,
5043 ReloadRegistrationFailure {
5044 exit_report: registration_failure_exit_report(exit_report),
5045 reason: "new child exited before registering".to_string(),
5046 },
5047 )
5048 .await
5049 }
5050 RegistrationWaitOutcome::TimedOut => {
5051 let mut timed_out_child = child
5052 .take()
5053 .expect("timed-out reload child is still running");
5054 timed_out_child
5055 .start_kill()
5056 .map_err(|source| SuperviseError::Kill {
5057 module_id: spec.module_id.clone(),
5058 source,
5059 })?;
5060 let status = timed_out_child
5061 .wait()
5062 .await
5063 .map_err(|source| SuperviseError::Wait {
5064 module_id: spec.module_id.clone(),
5065 source,
5066 })?;
5067 timed_out_child.drain_stderr(&spec.module_id).await;
5068 handle_reload_child_registration_failure(
5069 spec,
5070 runtime,
5071 registry,
5072 process_liveness,
5073 snapshot,
5074 child,
5075 ReloadRegistrationFailure {
5076 exit_report: registration_failure_exit_report(classify_reaped_child_exit(
5077 snapshot,
5078 &timed_out_child,
5079 &status,
5080 )),
5081 reason: format!(
5082 "new child did not register within {:?}",
5083 REGISTRY_RELEASE_TIMEOUT
5084 ),
5085 },
5086 )
5087 .await
5088 }
5089 }
5090}
5091
5092async fn set_child_enabled(
5093 spec: &ModuleSpec,
5094 runtime: &SupervisorRuntimeConfig,
5095 registry: &Registry,
5096 process_liveness: &SupervisorProcessLiveness,
5097 snapshot: &SharedSnapshot,
5098 child: &mut Option<SupervisedChild>,
5099 enabled: bool,
5100) -> Result<bool, SuperviseError> {
5101 let (current_enabled, current_state) = {
5102 let state = lock_snapshot(snapshot)?;
5103 (state.enabled, state.state)
5104 };
5105 let revive_terminal = enabled
5113 && current_enabled
5114 && child.is_none()
5115 && matches!(current_state, ModuleState::Failed | ModuleState::Stopped);
5116 if current_enabled == enabled && !revive_terminal {
5117 return Ok(false);
5118 }
5119
5120 if enabled {
5121 update_snapshot(snapshot, Some(&spec.module_id), |state| {
5122 state.enabled = true;
5123 state.state = ModuleState::Starting;
5124 clear_current_process_facts(state);
5125 })?;
5126 #[cfg(test)]
5127 if runtime.test_seed_stale_facts_before_enable_spawn {
5128 update_snapshot(snapshot, Some(&spec.module_id), |state| {
5129 state.process_alive = true;
5130 state.pid = Some(41);
5131 state.spawned_at_ms = Some(42);
5132 state.spawned_from = Some(PathBuf::from("/spawned/module"));
5133 state.spawned_file_identity = Some(SpawnedFileIdentity {
5134 device: 43,
5135 inode: 44,
5136 });
5137 })?;
5138 }
5139 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
5140 reset_restart_count(snapshot, &spec.module_id)?;
5141 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
5142 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
5143 Ok(next_child) => next_child,
5144 Err(err) => {
5145 if let Err(state_err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5146 state.state = ModuleState::Failed;
5147 clear_current_process_facts(state);
5148 }) {
5149 error!(module_id = %spec.module_id, error = %state_err, "failed to record enable spawn failure");
5150 }
5151 process_liveness.untrack_if_current(&spec.module_id, snapshot);
5152 return Err(err);
5153 }
5154 };
5155 *child = Some(next_child);
5156 debug!(module_id = %spec.module_id, "supervised module enabled");
5157 Ok(true)
5158 } else {
5159 let stop_notice = begin_forwarding_drain_if_configured(
5160 spec,
5161 runtime,
5162 registry,
5163 snapshot,
5164 Some(false),
5165 RouteCloseReason::Disable,
5166 )
5167 .await?;
5168 drain_optional_child(
5169 &spec.module_id,
5170 spec.protocol,
5171 stop_notice,
5172 registry,
5173 snapshot,
5174 &runtime.terminal_ring,
5175 &runtime.spawn_events,
5176 child,
5177 runtime.drain_timeout,
5178 ModuleState::Disabled,
5179 Some(false),
5180 )
5181 .await?;
5182 debug!(module_id = %spec.module_id, "supervised module disabled");
5183 Ok(true)
5184 }
5185}
5186
5187#[allow(clippy::too_many_arguments)]
5188async fn on_child_exit(
5189 spec: &ModuleSpec,
5190 policy: RestartPolicy,
5191 registry: &Registry,
5192 snapshot: &SharedSnapshot,
5193 terminal_ring: &Arc<Mutex<TerminalRing>>,
5194 spawn_events: &SpawnEventFeed,
5195 roster: &ChildRoster,
5196 exit_report: ExitReport,
5197) -> NextAction {
5198 if roster.is_closed() {
5204 return on_child_exit_during_daemon_shutdown(
5205 spec,
5206 registry,
5207 snapshot,
5208 terminal_ring,
5209 spawn_events,
5210 exit_report,
5211 )
5212 .await;
5213 }
5214 match exit_report.kind {
5215 ExitKind::Clean => {
5216 info!(
5217 module_id = %spec.module_id,
5218 exit_code = ?exit_report.code,
5219 exit_signal = ?exit_report.signal,
5220 "supervised module exited cleanly"
5221 );
5222 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5223 state.state = ModuleState::Stopped;
5224 clear_current_process_facts(state);
5225 state.last_exit = Some(exit_report.clone());
5226 }) {
5227 error!(module_id = %spec.module_id, error = %err, "failed to record clean module exit");
5228 }
5229 record_terminal(
5230 &spec.module_id,
5231 terminal_ring,
5232 spawn_events,
5233 &exit_report,
5234 TerminalDisposition::Stopped,
5235 );
5236 let registration_released = match wait_for_registration_release(
5237 registry,
5238 &spec.module_id,
5239 REGISTRY_RELEASE_TIMEOUT,
5240 )
5241 .await
5242 {
5243 Ok(()) => true,
5244 Err(err) => {
5245 warn!(module_id = %spec.module_id, error = %err, "registration still active after clean exit");
5246 false
5247 }
5248 };
5249 NextAction::Stop {
5250 registration_released,
5251 }
5252 }
5253 ExitKind::Crash => {
5254 warn!(
5255 module_id = %spec.module_id,
5256 exit_code = ?exit_report.code,
5257 exit_signal = ?exit_report.signal,
5258 "supervised module exited abnormally (crash)"
5259 );
5260 let mut restart_schedule = None;
5261 let mut disposition = TerminalDisposition::Disabled;
5262 let mut disposition_detail = None;
5266 let now = Instant::now();
5267 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5268 clear_current_process_facts(state);
5269 state.last_exit = Some(exit_report.clone());
5270 if state.enabled {
5271 if let Some(schedule) = state.next_crash_restart(&policy, now) {
5272 state.state = ModuleState::Restarting;
5273 restart_schedule = Some(schedule);
5274 disposition = TerminalDisposition::Restarting;
5275 } else {
5276 state.state = ModuleState::Failed;
5277 disposition = TerminalDisposition::Failed;
5278 disposition_detail = Some(policy.budget_exhausted_detail());
5279 }
5280 } else {
5281 state.state = ModuleState::Disabled;
5282 disposition = TerminalDisposition::Disabled;
5283 }
5284 }) {
5285 error!(module_id = %spec.module_id, error = %err, "failed to record crashed module exit");
5286 return NextAction::Stop {
5287 registration_released: false,
5288 };
5289 }
5290 if disposition_detail.is_some() {
5291 error!(
5296 module_id = %spec.module_id,
5297 max_restarts = policy.max_restarts,
5298 window_secs = policy.window.as_secs(),
5299 "module stopped: {}",
5300 policy.budget_exhausted_detail()
5301 );
5302 }
5303 record_terminal_with_detail(
5304 &spec.module_id,
5305 terminal_ring,
5306 spawn_events,
5307 &exit_report,
5308 disposition,
5309 disposition_detail,
5310 );
5311
5312 if let Some(schedule) = restart_schedule {
5313 NextAction::Restart {
5314 schedule: Some(schedule),
5315 }
5316 } else {
5317 let registration_released = match wait_for_registration_release(
5318 registry,
5319 &spec.module_id,
5320 REGISTRY_RELEASE_TIMEOUT,
5321 )
5322 .await
5323 {
5324 Ok(()) => true,
5325 Err(err) => {
5326 warn!(module_id = %spec.module_id, error = %err, "registration still active after failed module");
5327 false
5328 }
5329 };
5330 NextAction::Stop {
5331 registration_released,
5332 }
5333 }
5334 }
5335 ExitKind::DeliberateSeverance => {
5336 warn!(
5337 module_id = %spec.module_id,
5338 exit_code = ?exit_report.code,
5339 exit_signal = ?exit_report.signal,
5340 "supervised module exited after deliberate connection severance"
5341 );
5342 let mut should_restart = false;
5343 let mut disposition = TerminalDisposition::Disabled;
5344 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5345 clear_current_process_facts(state);
5346 state.last_exit = Some(exit_report.clone());
5347 state.lifetime_restarts += 1;
5348 if state.enabled {
5349 state.state = ModuleState::Restarting;
5350 should_restart = true;
5351 disposition = TerminalDisposition::Restarting;
5352 } else {
5353 state.state = ModuleState::Disabled;
5354 }
5355 }) {
5356 error!(module_id = %spec.module_id, error = %err, "failed to record deliberately severed module exit");
5357 return NextAction::Stop {
5358 registration_released: false,
5359 };
5360 }
5361 record_terminal(
5362 &spec.module_id,
5363 terminal_ring,
5364 spawn_events,
5365 &exit_report,
5366 disposition,
5367 );
5368
5369 if should_restart {
5370 NextAction::Restart { schedule: None }
5371 } else {
5372 let registration_released = match wait_for_registration_release(
5373 registry,
5374 &spec.module_id,
5375 REGISTRY_RELEASE_TIMEOUT,
5376 )
5377 .await
5378 {
5379 Ok(()) => true,
5380 Err(err) => {
5381 warn!(module_id = %spec.module_id, error = %err, "registration still active after deliberately severed module exit");
5382 false
5383 }
5384 };
5385 NextAction::Stop {
5386 registration_released,
5387 }
5388 }
5389 }
5390 }
5391}
5392
5393async fn on_child_exit_during_daemon_shutdown(
5394 spec: &ModuleSpec,
5395 registry: &Registry,
5396 snapshot: &SharedSnapshot,
5397 terminal_ring: &Arc<Mutex<TerminalRing>>,
5398 spawn_events: &SpawnEventFeed,
5399 exit_report: ExitReport,
5400) -> NextAction {
5401 info!(
5402 module_id = %spec.module_id,
5403 exit_code = ?exit_report.code,
5404 exit_signal = ?exit_report.signal,
5405 exit_kind = ?exit_report.kind,
5406 "supervised module exited during daemon shutdown; not restarting it"
5407 );
5408 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5409 state.state = ModuleState::Stopped;
5410 clear_current_process_facts(state);
5411 state.last_exit = Some(exit_report.clone());
5412 }) {
5413 error!(module_id = %spec.module_id, error = %err, "failed to record module exit during daemon shutdown");
5414 }
5415 record_terminal(
5416 &spec.module_id,
5417 terminal_ring,
5418 spawn_events,
5419 &exit_report,
5420 TerminalDisposition::DaemonShutdown,
5421 );
5422 let registration_released =
5423 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT)
5424 .await
5425 .is_ok();
5426 NextAction::Stop {
5427 registration_released,
5428 }
5429}
5430
5431fn record_wait_error_terminal(
5432 module_id: &str,
5433 terminal_ring: &Arc<Mutex<TerminalRing>>,
5434 spawn_events: &SpawnEventFeed,
5435) {
5436 record_terminal(
5437 module_id,
5438 terminal_ring,
5439 spawn_events,
5440 &wait_error_exit_report(),
5441 TerminalDisposition::Failed,
5442 );
5443}
5444
5445fn record_terminal(
5446 module_id: &str,
5447 terminal_ring: &Arc<Mutex<TerminalRing>>,
5448 spawn_events: &SpawnEventFeed,
5449 exit_report: &ExitReport,
5450 disposition: TerminalDisposition,
5451) {
5452 record_terminal_with_detail(
5453 module_id,
5454 terminal_ring,
5455 spawn_events,
5456 exit_report,
5457 disposition,
5458 None,
5459 );
5460}
5461
5462fn durable_terminal_history_of(
5466 terminal_ring: &Mutex<TerminalRing>,
5467 module_id: &str,
5468) -> subc_control::TerminalHistory {
5469 let read = terminal_ring
5470 .lock()
5471 .unwrap_or_else(|p| p.into_inner())
5472 .capture_durable_history();
5473 read.read(module_id)
5474}
5475
5476fn record_terminal_with_detail(
5477 module_id: &str,
5478 terminal_ring: &Arc<Mutex<TerminalRing>>,
5479 spawn_events: &SpawnEventFeed,
5480 exit_report: &ExitReport,
5481 disposition: TerminalDisposition,
5482 disposition_detail: Option<String>,
5483) {
5484 spawn_events.emit_exited(module_id, exit_report.code, exit_report.signal);
5485 let record = TerminalRecord {
5486 exit_code: exit_report.code,
5487 exit_signal: exit_report.signal,
5488 at_ms: exit_report.at_ms,
5489 disposition,
5490 exit_kind: exit_report.kind.into(),
5491 disposition_detail,
5492 };
5493 terminal_ring
5494 .lock()
5495 .unwrap_or_else(|poisoned| poisoned.into_inner())
5496 .record_exit(module_id, record);
5497}
5498
5499fn untrack_if_registration_released(
5500 process_liveness: &SupervisorProcessLiveness,
5501 registry: &Registry,
5502 module_id: &str,
5503 snapshot: &SharedSnapshot,
5504) {
5505 match registry.get_module(module_id) {
5506 Ok(None) => process_liveness.untrack_if_current(module_id, snapshot),
5507 Ok(Some(_)) => {}
5508 Err(err) => {
5509 warn!(module_id, error = %err, "could not determine whether supervisor liveness can be untracked");
5510 }
5511 }
5512}
5513
5514#[cfg(test)]
5528fn apply_wire_spawn_args(
5529 command: &mut Command,
5530 spec: &ModuleSpec,
5531 connection_file_path: Option<&std::path::Path>,
5532 handle: Option<&SupervisorHandle>,
5533) -> Result<(), SuperviseError> {
5534 apply_wire_spawn_args_for_role(
5535 command,
5536 spec,
5537 connection_file_path,
5538 handle,
5539 SpawnRole::Plain,
5540 )
5541}
5542
5543fn apply_wire_spawn_args_for_role(
5552 command: &mut Command,
5553 spec: &ModuleSpec,
5554 connection_file_path: Option<&std::path::Path>,
5555 handle: Option<&SupervisorHandle>,
5556 role: SpawnRole,
5557) -> Result<(), SuperviseError> {
5558 command.env(SUBC_MODULE_ID_ENV, &spec.module_id);
5559 if spec.protocol == ModuleProtocol::None {
5560 return Ok(());
5561 }
5562 if let Some(connection_file_path) = connection_file_path {
5563 command.arg(SUBC_ARG).arg(connection_file_path);
5564 }
5565
5566 let nonce = generate_launch_nonce()?;
5570 if let Some(handle) = handle {
5571 match role {
5572 SpawnRole::Plain => {
5573 handle.set_spawn_nonce(&spec.module_id, nonce.clone());
5574 if spec.reserved {
5575 handle.set_reserved_nonce(&spec.module_id, nonce.clone());
5576 }
5577 }
5578 SpawnRole::SwapCandidate => handle.open_swap(&spec.module_id, nonce.clone()),
5579 }
5580 }
5581 command.env(SUBC_LAUNCH_NONCE_ENV, nonce);
5582 Ok(())
5583}
5584
5585fn apply_child_env(command: &mut Command, spec: &ModuleSpec) {
5586 command.env_remove(CK_LOG_ENV);
5587 command.env_remove(SUBC_SPAWN_ROLE_ENV);
5594 for (key, value) in &spec.env {
5595 if matches!(
5599 key.as_str(),
5600 CAPTURE_MAX_FILE_MB_ENV | CAPTURE_KEEP_ENV | CAPTURE_MAX_AGE_DAYS_ENV
5601 ) || key == SUBC_SPAWN_ROLE_ENV
5602 {
5603 continue;
5604 }
5605 command.env(key, value);
5606 }
5607}
5608
5609#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5612enum SpawnRole {
5613 Plain,
5614 SwapCandidate,
5615}
5616
5617fn apply_spawn_role(command: &mut Command, role: SpawnRole) {
5620 if role == SpawnRole::SwapCandidate {
5621 command.env(SUBC_SPAWN_ROLE_ENV, SPAWN_ROLE_SWAP_CANDIDATE);
5622 }
5623}
5624
5625fn spawn_child(
5626 spec: &ModuleSpec,
5627 connection_file_path: Option<&std::path::Path>,
5628 handle: Option<&SupervisorHandle>,
5629 ring: &Arc<Mutex<StderrRing>>,
5630 capture_logs_dir: Option<&std::path::Path>,
5631 roster: &ChildRoster,
5632 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5633) -> Result<SupervisedChild, SuperviseError> {
5634 spawn_child_in_slot(
5635 spec,
5636 connection_file_path,
5637 handle,
5638 ring,
5639 capture_logs_dir,
5640 roster,
5641 #[cfg(target_os = "linux")]
5642 cgroup_placement,
5643 SpawnRole::Plain,
5644 false,
5645 )
5646}
5647
5648#[allow(clippy::too_many_arguments)]
5661fn spawn_child_in_slot(
5662 spec: &ModuleSpec,
5663 connection_file_path: Option<&std::path::Path>,
5664 handle: Option<&SupervisorHandle>,
5665 ring: &Arc<Mutex<StderrRing>>,
5666 capture_logs_dir: Option<&std::path::Path>,
5667 roster: &ChildRoster,
5668 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5669 role: SpawnRole,
5670 alternate_slot: bool,
5671) -> Result<SupervisedChild, SuperviseError> {
5672 if roster.is_closed() {
5673 return Err(SuperviseError::Spawn {
5674 program: spec.program.clone(),
5675 source: io::Error::other("the daemon is shutting down; not starting a new process"),
5676 cgroup_path: None,
5677 });
5678 }
5679 #[cfg(target_os = "linux")]
5680 let cgroup_name = swap::cgroup_name(&spec.module_id, alternate_slot);
5681 #[cfg(not(target_os = "linux"))]
5682 let _ = alternate_slot;
5683 let mut command = Command::new(&spec.program);
5684 command.args(&spec.args);
5685 apply_child_env(&mut command, spec);
5715 apply_spawn_role(&mut command, role);
5716 apply_wire_spawn_args_for_role(&mut command, spec, connection_file_path, handle, role)?;
5717
5718 #[cfg(target_os = "linux")]
5719 let cgroup_path = cgroup_placement
5720 .map(|placement| placement.module_path(&cgroup_name))
5721 .transpose()
5722 .map_err(|source| SuperviseError::Cgroup {
5723 module_id: spec.module_id.clone(),
5724 source,
5725 })?;
5726 #[cfg(not(target_os = "linux"))]
5727 let cgroup_path: Option<PathBuf> = None;
5728 #[cfg(target_os = "linux")]
5729 if let Some(path) = &cgroup_path {
5730 if let Err(error) = apply_cgroup_placement(&mut command, spec, path) {
5731 if let Some(placement) = cgroup_placement {
5732 remove_module_cgroup(placement, &cgroup_name);
5733 }
5734 return Err(error);
5735 }
5736 }
5737
5738 let output_sink = if let Some(logs_dir) = capture_logs_dir {
5739 let path = logs_dir.join(format!("{}.stderr.log", spec.module_id));
5740 match ChildOutputSink::open(&path, capture_retention(spec)) {
5741 Ok(sink) => sink,
5742 Err(error) => {
5743 warn!(
5744 module_id = %spec.module_id,
5745 path = %path.display(),
5746 error = %error,
5747 "could not open child output capture file; forwarding to stderr"
5748 );
5749 ChildOutputSink::Stderr
5750 }
5751 }
5752 } else {
5753 ChildOutputSink::Stderr
5754 };
5755
5756 command.stdout(Stdio::piped());
5757 command.stderr(Stdio::piped());
5758 command.kill_on_drop(true);
5759 #[cfg(unix)]
5776 command.process_group(0);
5777 command.stdin(Stdio::null());
5778 let mut child = match command.spawn() {
5779 Ok(child) => child,
5780 Err(source) => {
5781 #[cfg(target_os = "linux")]
5782 if let Some(placement) = cgroup_placement {
5783 remove_module_cgroup(placement, &cgroup_name);
5784 }
5785 return Err(SuperviseError::Spawn {
5786 program: spec.program.clone(),
5787 source,
5788 cgroup_path,
5789 });
5790 }
5791 };
5792 let spawned_at_ms = unix_ms_now();
5793 let spawned_from = spec.program.clone();
5794 let spawned_file_identity = spawned_file_identity(&spawned_from);
5795 let pid = child.id().ok_or_else(|| SuperviseError::Spawn {
5796 program: spec.program.clone(),
5797 source: io::Error::other("spawned child exposed no live pid"),
5798 cgroup_path: cgroup_path.clone(),
5799 })?;
5800 let process_start_time = crate::provenance::process_start_time(pid);
5801 let process_identity = process_start_time.map(|start_time| ProcessIdentity { pid, start_time });
5802 #[cfg(target_os = "linux")]
5806 let recorded_cgroup_name = cgroup_path.as_ref().map(|_| cgroup_name.clone());
5807 #[cfg(not(target_os = "linux"))]
5808 let recorded_cgroup_name = None;
5809 let roster_guard = roster.admit(
5810 spec.module_id.clone(),
5811 pid,
5812 spec.protocol,
5813 process_start_time,
5814 crate::child_roster::RecordedIdentity {
5815 start_time: subc_os::start_time(pid),
5816 executable: spawned_file_identity.map(|identity| {
5817 crate::live_children::ExecutableIdentity {
5818 device: identity.device,
5819 inode: identity.inode,
5820 }
5821 }),
5822 cgroup_name: recorded_cgroup_name,
5823 },
5824 );
5825 if roster.is_closed() {
5834 if let Err(error) = child.start_kill() {
5835 debug!(module_id = %spec.module_id, pid, %error, "kill of a process spawned during daemon shutdown failed; it may already have exited");
5836 }
5837 drop(roster_guard);
5838 return Err(SuperviseError::Spawn {
5839 program: spec.program.clone(),
5840 source: io::Error::other(
5841 "the daemon began shutting down while this process was starting; ended it",
5842 ),
5843 cgroup_path,
5844 });
5845 }
5846
5847 let stdout_pump = match child.stdout.take() {
5848 Some(stdout) => Some(tokio::spawn(pump_stdout_to(stdout, output_sink.clone()))),
5849 None => {
5850 warn!(
5851 module_id = %spec.module_id,
5852 "spawned child exposed no stdout pipe; file capture will be incomplete"
5853 );
5854 None
5855 }
5856 };
5857 let stderr_pump = match child.stderr.take() {
5858 Some(stderr) => {
5859 let generation = ring
5860 .lock()
5861 .unwrap_or_else(|poisoned| poisoned.into_inner())
5862 .begin_process();
5863 Some(StderrPump {
5864 task: tokio::spawn(pump_stderr_to(
5865 stderr,
5866 Arc::clone(ring),
5867 generation,
5868 output_sink,
5869 )),
5870 generation,
5871 })
5872 }
5873 None => {
5874 ring.lock()
5878 .unwrap_or_else(|poisoned| poisoned.into_inner())
5879 .mark_not_captured("stderr pipe was not available on spawn");
5880 warn!(
5881 module_id = %spec.module_id,
5882 "spawned child exposed no stderr pipe; tail will be unavailable"
5883 );
5884 None
5885 }
5886 };
5887
5888 Ok(SupervisedChild {
5889 child,
5890 #[cfg(target_os = "linux")]
5891 module_id: cgroup_name,
5892 #[cfg(target_os = "linux")]
5893 cgroup_placement: cgroup_placement.cloned(),
5894 stdout_pump,
5895 stderr_pump,
5896 stderr_ring: Arc::clone(ring),
5897 spawned_at_ms,
5898 spawned_from,
5899 spawned_file_identity,
5900 process_start_time,
5901 process_identity,
5902 pid,
5903 roster_guard: Some(roster_guard),
5904 })
5905}
5906
5907#[cfg(target_os = "linux")]
5908fn remove_module_cgroup(placement: &subc_cgroup::Placement, module_id: &str) {
5909 match placement.remove_module(module_id) {
5910 Ok(()) => debug!(module_id, "removed module cgroup after process exit"),
5911 Err(error) => warn!(
5912 module_id,
5913 error = %error,
5914 "could not remove module cgroup after process exit; continuing teardown"
5915 ),
5916 }
5917}
5918
5919#[cfg(target_os = "linux")]
5920fn apply_cgroup_placement(
5921 command: &mut Command,
5922 spec: &ModuleSpec,
5923 path: &std::path::Path,
5924) -> Result<(), SuperviseError> {
5925 subc_cgroup::apply(command, path).map_err(|source| SuperviseError::Cgroup {
5926 module_id: spec.module_id.clone(),
5927 source,
5928 })
5929}
5930
5931fn capture_retention(spec: &ModuleSpec) -> Retention {
5932 let defaults = Retention::default();
5933 let value = |name: &str| {
5934 spec.env
5935 .iter()
5936 .rev()
5937 .find_map(|(key, value)| (key == name).then_some(value.as_str()))
5938 };
5939 Retention {
5940 max_file_mb: value(CAPTURE_MAX_FILE_MB_ENV)
5941 .and_then(|value| value.parse().ok())
5942 .unwrap_or(defaults.max_file_mb),
5943 keep: value(CAPTURE_KEEP_ENV)
5944 .and_then(|value| value.parse().ok())
5945 .unwrap_or(defaults.keep),
5946 max_age_days: value(CAPTURE_MAX_AGE_DAYS_ENV)
5947 .and_then(|value| value.parse().ok())
5948 .unwrap_or(defaults.max_age_days),
5949 }
5950}
5951
5952fn generate_launch_nonce() -> Result<String, SuperviseError> {
5955 let mut bytes = [0u8; 32];
5956 getrandom::getrandom(&mut bytes).map_err(|source| SuperviseError::LaunchNonce {
5957 reason: source.to_string(),
5958 })?;
5959 let mut hex = String::with_capacity(64);
5960 for b in bytes {
5961 use std::fmt::Write;
5962 let _ = write!(hex, "{b:02x}");
5963 }
5964 Ok(hex)
5965}
5966
5967fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
5970 if a.len() != b.len() {
5971 return false;
5972 }
5973 let mut diff = 0u8;
5974 for (x, y) in a.iter().zip(b.iter()) {
5975 diff |= x ^ y;
5976 }
5977 diff == 0
5978}
5979
5980fn spawn_and_mark_running(
5981 spec: &ModuleSpec,
5982 runtime: &SupervisorRuntimeConfig,
5983 snapshot: &SharedSnapshot,
5984) -> Result<SupervisedChild, SuperviseError> {
5985 let child = spawn_child(
5986 spec,
5987 runtime.connection_file_path.as_deref(),
5988 runtime.supervisor_handle.as_ref(),
5989 &runtime.stderr_ring,
5990 runtime.capture_logs_dir.as_deref(),
5991 &runtime.child_roster,
5992 #[cfg(target_os = "linux")]
5993 runtime.cgroup_placement.as_ref(),
5994 )?;
5995 set_running(snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
5996 Ok(child)
5997}
5998
5999enum RegistrationWaitOutcome {
6000 Registered,
6001 Exited(ExitReport),
6002 TimedOut,
6003}
6004
6005struct ReloadRegistrationFailure {
6006 exit_report: ExitReport,
6007 reason: String,
6008}
6009
6010#[derive(Debug, Clone, Copy, PartialEq, Eq)]
6011enum BusyGaugeObservation {
6012 Quiescent,
6013 Busy,
6014 Omitted,
6015}
6016
6017fn busy_gauge_observation(metrics: Option<&Value>, gauges: &[String]) -> BusyGaugeObservation {
6018 let Some(metrics) = metrics.and_then(Value::as_object) else {
6019 return BusyGaugeObservation::Omitted;
6020 };
6021 let mut sum = 0u128;
6022 for gauge in gauges {
6023 let Some(value) = metrics.get(gauge) else {
6024 return BusyGaugeObservation::Omitted;
6025 };
6026 let Some(value) = value.as_u64() else {
6027 return BusyGaugeObservation::Busy;
6028 };
6029 sum = sum.saturating_add(u128::from(value));
6030 }
6031 if sum == 0 {
6032 BusyGaugeObservation::Quiescent
6033 } else {
6034 BusyGaugeObservation::Busy
6035 }
6036}
6037
6038fn declared_busy_gauges(
6039 registry: &Registry,
6040 module_id: &str,
6041) -> Result<Vec<String>, SuperviseError> {
6042 busy_gauges_of(
6043 registry
6044 .get_module(module_id)
6045 .map_err(SuperviseError::Registry)?,
6046 )
6047}
6048
6049fn declared_busy_gauges_for_connection(
6053 registry: &Registry,
6054 connection_id: ConnectionId,
6055) -> Result<Vec<String>, SuperviseError> {
6056 busy_gauges_of(
6057 registry
6058 .get_module_by_connection(connection_id)
6059 .map_err(SuperviseError::Registry)?,
6060 )
6061}
6062
6063fn busy_gauges_of(
6064 registration: Option<crate::registry::ModuleRegistration>,
6065) -> Result<Vec<String>, SuperviseError> {
6066 let Some(registration) = registration else {
6067 return Ok(Vec::new());
6068 };
6069 let Some(self_signals) = registration.manifest.self_signals else {
6070 return Ok(Vec::new());
6071 };
6072
6073 let mut gauges = Vec::new();
6074 for declaration in self_signals {
6075 if declaration.kind != SelfSignalKind::Busy {
6076 continue;
6077 }
6078 match declaration.anchored_to {
6079 SignalAnchor::HealthGauges { gauges: declared } if !declared.is_empty() => {
6080 gauges.extend(declared)
6081 }
6082 _ => {
6083 gauges.push(String::new());
6086 }
6087 }
6088 }
6089 Ok(gauges)
6090}
6091
6092async fn wait_for_forwarding_quiescence(
6097 forwarding: &ForwardingTable,
6098 module_id: &str,
6099 runtime: &SupervisorRuntimeConfig,
6100 endpoint: crate::ModuleEndpointId,
6101 deadline: Instant,
6102 busy_gauges: &[String],
6103 scope: DrainScope,
6104) -> Result<bool, SuperviseError> {
6105 let mut gauges_quiescent = busy_gauges.is_empty();
6106 let mut next_probe_at = Instant::now();
6107 let mut omission_counted = false;
6108
6109 loop {
6110 let now = Instant::now();
6111 if !busy_gauges.is_empty() && now >= next_probe_at && now < deadline {
6112 let report = match scope {
6113 DrainScope::Active => probe_module_health(module_id, runtime, Some(deadline)).await,
6114 DrainScope::Endpoint(endpoint) => {
6115 probe_endpoint_health(endpoint, runtime, Some(deadline)).await
6116 }
6117 };
6118 gauges_quiescent = match report {
6119 Ok(report) => match busy_gauge_observation(report.metrics.as_ref(), busy_gauges) {
6120 BusyGaugeObservation::Quiescent => true,
6121 BusyGaugeObservation::Busy => false,
6122 BusyGaugeObservation::Omitted => {
6123 if !omission_counted {
6124 forwarding
6125 .counters()
6126 .increment_drains_with_undeclared_gauge();
6127 omission_counted = true;
6128 }
6129 false
6130 }
6131 },
6132 Err(err) => {
6133 warn!(
6134 module_id,
6135 error = %err,
6136 "drain health.check did not produce declared busy gauges; treating module as busy"
6137 );
6138 false
6139 }
6140 };
6141 next_probe_at = Instant::now() + runtime.health.cadence.max(REGISTRY_RELEASE_POLL);
6142 }
6143
6144 let in_flight = forwarding
6145 .endpoint_in_flight_count(endpoint)
6146 .map_err(SuperviseError::Forwarding)?;
6147 if in_flight == 0 && gauges_quiescent {
6148 return Ok(true);
6149 }
6150
6151 let now = Instant::now();
6152 if now >= deadline {
6153 return Ok(false);
6154 }
6155 let mut wait = deadline
6156 .saturating_duration_since(now)
6157 .min(REGISTRY_RELEASE_POLL);
6158 if !busy_gauges.is_empty() {
6159 wait = wait.min(next_probe_at.saturating_duration_since(now));
6160 }
6161 sleep(wait).await;
6162 }
6163}
6164
6165fn drained_after_quiescence_wait(wait_result: &Result<bool, SuperviseError>) -> bool {
6173 match wait_result {
6174 Ok(drained) => *drained,
6175 Err(_) => false,
6176 }
6177}
6178
6179fn send_route_goodbyes(forwarding: &ForwardingTable, released_routes: Vec<GoodbyeTarget>) {
6180 for released in released_routes {
6181 let frame = match Frame::build_with_version(
6182 released.negotiated_ver,
6183 FrameType::Goodbye,
6184 control_flags(),
6185 released.channel,
6186 released.epoch,
6187 0,
6188 Vec::new(),
6189 ) {
6190 Ok(frame) => frame,
6191 Err(err) => {
6192 warn!(
6193 route_channel = released.channel,
6194 error = %err,
6195 "failed to build supervisor drain route GOODBYE frame"
6196 );
6197 continue;
6198 }
6199 };
6200 if !released.close_on_delivery_failure() {
6201 crate::forwarding::send_module_route_goodbye(
6202 &forwarding.counters(),
6203 &released.sink,
6204 frame,
6205 released.module_id.as_deref(),
6206 "supervisor drain",
6207 );
6208 continue;
6209 }
6210 if let Err(err) = released.sink.try_send(frame) {
6211 warn!(
6212 target_connection_id = released.connection_id.get(),
6213 route_channel = released.channel,
6214 error = %err,
6215 "supervisor drain route GOODBYE was not delivered to client; closing target connection"
6216 );
6217 let _ = forwarding.escalate_client_delivery_failure(
6218 released.connection_id,
6219 released.channel,
6220 released.epoch,
6221 CloseReason::new(
6222 "route_goodbye_delivery_failed",
6223 format!(
6224 "failed to enqueue supervisor drain route GOODBYE for channel {}: {err}",
6225 released.channel
6226 ),
6227 ),
6228 crate::forwarding::UndeliveredFrame {
6229 module_id: released.module_id.as_deref(),
6230 sink: &released.sink,
6231 },
6232 );
6233 }
6234 }
6235}
6236
6237fn send_module_draining(
6238 module_id: &str,
6239 reason: RouteCloseReason,
6240 deadline_ms: u64,
6241 target: &ModuleDrainTarget,
6242) {
6243 let body = match serde_json::to_vec(&ModuleControlCommand::Draining {
6244 reason,
6245 deadline_ms,
6246 }) {
6247 Ok(body) => body,
6248 Err(err) => {
6249 warn!(
6250 module_id,
6251 error = %err,
6252 "failed to encode module draining command"
6253 );
6254 return;
6255 }
6256 };
6257 let frame = match Frame::build_with_version(
6258 target.negotiated_ver,
6259 FrameType::Push,
6260 control_flags(),
6261 0,
6262 0,
6263 0,
6264 body,
6265 ) {
6266 Ok(frame) => frame,
6267 Err(err) => {
6268 warn!(
6269 module_id,
6270 error = %err,
6271 "failed to build module draining command frame"
6272 );
6273 return;
6274 }
6275 };
6276 if let Err(err) = target.sink.try_send(frame) {
6277 warn!(
6278 module_id,
6279 target_connection_id = target.endpoint.connection_id.get(),
6280 error = %err,
6281 "module draining command was not delivered to peer"
6282 );
6283 }
6284}
6285
6286fn send_module_goodbye(module_id: &str, forwarding: &ForwardingTable, target: &ModuleDrainTarget) {
6287 let frame = match Frame::build_with_version(
6288 target.negotiated_ver,
6289 FrameType::Goodbye,
6290 control_flags(),
6291 0,
6292 0,
6293 0,
6294 Vec::new(),
6295 ) {
6296 Ok(frame) => frame,
6297 Err(err) => {
6298 warn!(
6299 module_id,
6300 error = %err,
6301 "failed to build supervisor drain module GOODBYE frame"
6302 );
6303 return;
6304 }
6305 };
6306 if let Err(err) = target.sink.try_send(frame) {
6307 warn!(
6308 module_id,
6309 target_connection_id = target.endpoint.connection_id.get(),
6310 error = %err,
6311 "supervisor drain module GOODBYE was not delivered to peer; closing module connection"
6312 );
6313 forwarding.request_connection_close(
6314 target.endpoint.connection_id,
6315 CloseReason::new(
6316 "module_goodbye_delivery_failed",
6317 format!("failed to enqueue supervisor drain module GOODBYE for module '{module_id}': {err}"),
6318 ),
6319 );
6320 }
6321}
6322
6323#[derive(Clone, Copy)]
6324struct ForwardingDrainContext<'a> {
6325 spec: &'a ModuleSpec,
6326 runtime: &'a SupervisorRuntimeConfig,
6327 registry: &'a Registry,
6328 scope: DrainScope,
6329}
6330
6331#[derive(Debug, Clone, Copy, PartialEq, Eq)]
6333enum DrainScope {
6334 Active,
6337 Endpoint(crate::ModuleEndpointId),
6342}
6343
6344#[derive(Debug, Clone, Copy, PartialEq, Eq)]
6352enum StopNotice {
6353 SentOverConnection,
6356 NoConnection,
6360 NotSent,
6364}
6365
6366async fn begin_forwarding_drain(
6367 spec: &ModuleSpec,
6368 runtime: &SupervisorRuntimeConfig,
6369 registry: &Registry,
6370 snapshot: &SharedSnapshot,
6371 enabled: Option<bool>,
6372 reason: RouteCloseReason,
6373) -> Result<StopNotice, SuperviseError> {
6374 let Some(forwarding) = runtime.forwarding.as_ref() else {
6375 return Err(SuperviseError::ReloadUnavailable {
6376 module_id: spec.module_id.clone(),
6377 reason: "supervisor was not configured with a forwarding table".to_string(),
6378 });
6379 };
6380
6381 begin_forwarding_drain_with(
6382 forwarding,
6383 ForwardingDrainContext {
6384 spec,
6385 runtime,
6386 registry,
6387 scope: DrainScope::Active,
6388 },
6389 snapshot,
6390 enabled,
6391 reason,
6392 runtime.drain_timeout,
6393 )
6394 .await
6395}
6396
6397async fn begin_forwarding_drain_if_configured(
6398 spec: &ModuleSpec,
6399 runtime: &SupervisorRuntimeConfig,
6400 registry: &Registry,
6401 snapshot: &SharedSnapshot,
6402 enabled: Option<bool>,
6403 reason: RouteCloseReason,
6404) -> Result<StopNotice, SuperviseError> {
6405 begin_forwarding_drain_with_timeout(
6406 spec,
6407 runtime,
6408 registry,
6409 snapshot,
6410 enabled,
6411 reason,
6412 runtime.drain_timeout,
6413 )
6414 .await
6415}
6416
6417async fn begin_forwarding_drain_with_timeout(
6421 spec: &ModuleSpec,
6422 runtime: &SupervisorRuntimeConfig,
6423 registry: &Registry,
6424 snapshot: &SharedSnapshot,
6425 enabled: Option<bool>,
6426 reason: RouteCloseReason,
6427 drain_timeout: Duration,
6428) -> Result<StopNotice, SuperviseError> {
6429 let Some(forwarding) = runtime.forwarding.as_ref() else {
6430 return Ok(StopNotice::NotSent);
6431 };
6432
6433 begin_forwarding_drain_with(
6434 forwarding,
6435 ForwardingDrainContext {
6436 spec,
6437 runtime,
6438 registry,
6439 scope: DrainScope::Active,
6440 },
6441 snapshot,
6442 enabled,
6443 reason,
6444 drain_timeout,
6445 )
6446 .await
6447}
6448
6449async fn begin_forwarding_drain_with(
6450 forwarding: &ForwardingTable,
6451 context: ForwardingDrainContext<'_>,
6452 snapshot: &SharedSnapshot,
6453 enabled: Option<bool>,
6454 reason: RouteCloseReason,
6455 drain_timeout: Duration,
6456) -> Result<StopNotice, SuperviseError> {
6457 let ForwardingDrainContext {
6458 spec,
6459 runtime,
6460 registry,
6461 scope,
6462 } = context;
6463 debug_assert_ne!(reason, RouteCloseReason::Crash);
6464 let terminal = matches!(reason, RouteCloseReason::Disable);
6465 let drain_started_at = Instant::now();
6466 let drain_deadline = drain_started_at + drain_timeout;
6467 let deadline_ms =
6468 unix_ms_now().saturating_add(u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX));
6469 let busy_gauges = match scope {
6470 DrainScope::Active => declared_busy_gauges(registry, &spec.module_id)?,
6471 DrainScope::Endpoint(endpoint) => {
6472 declared_busy_gauges_for_connection(registry, endpoint.connection_id)?
6473 }
6474 };
6475
6476 let gate_started = Instant::now();
6479 let drain_target = match scope {
6480 DrainScope::Active => forwarding.begin_module_drain(&spec.module_id, reason),
6481 DrainScope::Endpoint(endpoint) => forwarding.begin_endpoint_drain(endpoint, reason),
6482 }
6483 .map_err(SuperviseError::Forwarding)?;
6484 info!(
6489 module_id = %spec.module_id,
6490 ?reason,
6491 gate_ms = u64::try_from(gate_started.elapsed().as_millis()).unwrap_or(u64::MAX),
6492 connected = drain_target.is_some(),
6493 "module drain began; route admission closed"
6494 );
6495 if scope == DrainScope::Active {
6496 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6497 state.state = ModuleState::Draining;
6498 state.draining_to_replace =
6499 matches!(reason, RouteCloseReason::Restart | RouteCloseReason::Reload);
6500 if let Some(enabled) = enabled {
6501 state.enabled = enabled;
6502 }
6503 })?;
6504 }
6505
6506 let Some(target) = drain_target.as_ref() else {
6507 return Ok(StopNotice::NoConnection);
6511 };
6512 {
6513 send_module_draining(&spec.module_id, reason, deadline_ms, target);
6514 let routes = forwarding
6515 .endpoint_routes(target.endpoint)
6516 .map_err(SuperviseError::Forwarding)?;
6517 let routes_notified = routes.len();
6518 crate::control::send_route_control_pushes(
6519 forwarding,
6520 routes.clone(),
6521 ClientControlPush::RouteClosing {
6522 module_id: spec.module_id.clone(),
6523 reason,
6524 },
6525 );
6526 send_route_goodbyes(forwarding, target.abandoned_bindings.clone());
6527
6528 let wait_result = wait_for_forwarding_quiescence(
6534 forwarding,
6535 &spec.module_id,
6536 runtime,
6537 target.endpoint,
6538 drain_deadline,
6539 &busy_gauges,
6540 scope,
6541 )
6542 .await;
6543 let drained = drained_after_quiescence_wait(&wait_result);
6544 if let Err(err) = &wait_result {
6545 error!(
6546 module_id = %spec.module_id,
6547 ?reason,
6548 error = %err,
6549 "forwarding quiescence wait failed after route.closing; forcing route.closed(drained: false) so the client is not left waiting on an unfulfilled promise"
6550 );
6551 } else if !drained {
6552 let holdouts = forwarding
6558 .endpoint_drain_holdouts(target.endpoint)
6559 .unwrap_or_default();
6560 warn!(
6561 module_id = %spec.module_id,
6562 waited = ?drain_timeout,
6563 ?reason,
6564 held_requests = holdouts.requests,
6565 held_routes = holdouts.routes,
6566 total_routes = holdouts.total_routes,
6567 top_connections = ?holdouts.top_connections,
6568 held = %holdouts
6571 .held
6572 .iter()
6573 .map(|(channel, corr)| format!("{channel}:{corr}"))
6574 .collect::<Vec<_>>()
6575 .join(","),
6576 "route drain timed out before request quiescence; forcing teardown"
6577 );
6578 }
6579 crate::control::send_route_control_pushes(
6580 forwarding,
6581 routes,
6582 ClientControlPush::RouteClosed {
6583 module_id: spec.module_id.clone(),
6584 reason,
6585 drained,
6586 abandoned: target.abandoned_bindings.len() as u32,
6587 excluded_subscriptions: target.excluded_subscriptions,
6588 terminal: Some(terminal),
6589 },
6590 );
6591 wait_result?;
6592
6593 let released_routes = match forwarding.release_module_endpoint_routes(target.endpoint) {
6599 Ok(routes) => routes,
6600 Err(err) => {
6601 warn!(
6602 module_id = %spec.module_id,
6603 ?reason,
6604 error = %err,
6605 "failed to release module endpoint routes after route.closed; module GOODBYE will still be sent"
6606 );
6607 send_module_goodbye(&spec.module_id, forwarding, target);
6608 return Err(SuperviseError::Forwarding(err));
6609 }
6610 };
6611 let route_goodbye_count = released_routes.len();
6612 send_route_goodbyes(forwarding, released_routes);
6613 send_module_goodbye(&spec.module_id, forwarding, target);
6614
6615 info!(
6621 module_id = %spec.module_id,
6622 ?reason,
6623 routes_notified,
6624 route_goodbyes = route_goodbye_count,
6625 abandoned_reservations = target.abandoned_bindings.len(),
6626 excluded_subscriptions = target.excluded_subscriptions,
6627 drained,
6628 "module drain complete; consumers notified via route.closing/route.closed pushes and per-route GOODBYE frames"
6629 );
6630 }
6631
6632 Ok(StopNotice::SentOverConnection)
6633}
6634
6635async fn wait_for_registration_after_reload(
6638 registry: &Registry,
6639 module_id: &str,
6640 snapshot: &SharedSnapshot,
6641 child: &mut SupervisedChild,
6642 wait: Duration,
6643) -> Result<RegistrationWaitOutcome, SuperviseError> {
6644 wait_for_slot_registration(
6645 registry,
6646 crate::registry::RegistrationSlot::Active(module_id),
6647 module_id,
6648 snapshot,
6649 child,
6650 wait,
6651 )
6652 .await
6653}
6654
6655async fn wait_for_slot_registration(
6663 registry: &Registry,
6664 slot: crate::registry::RegistrationSlot<'_>,
6665 module_id: &str,
6666 snapshot: &SharedSnapshot,
6667 child: &mut SupervisedChild,
6668 wait: Duration,
6669) -> Result<RegistrationWaitOutcome, SuperviseError> {
6670 let deadline = Instant::now() + wait;
6671 loop {
6672 if registry
6673 .registration(slot)
6674 .map_err(SuperviseError::Registry)?
6675 .is_some()
6676 {
6677 return Ok(RegistrationWaitOutcome::Registered);
6678 }
6679
6680 let now = Instant::now();
6681 if now >= deadline {
6682 return Ok(RegistrationWaitOutcome::TimedOut);
6683 }
6684 let remaining = deadline.saturating_duration_since(now);
6685 let poll = remaining.min(REGISTRY_RELEASE_POLL);
6686
6687 tokio::select! {
6688 wait_result = child.wait() => {
6689 let status = wait_result.map_err(|source| SuperviseError::Wait {
6690 module_id: module_id.to_string(),
6691 source,
6692 })?;
6693 return Ok(RegistrationWaitOutcome::Exited(classify_reaped_child_exit(
6694 snapshot,
6695 child,
6696 &status,
6697 )));
6698 }
6699 _ = sleep(poll) => {}
6700 }
6701 }
6702}
6703
6704fn registration_failure_exit_report(mut exit_report: ExitReport) -> ExitReport {
6705 if exit_report.kind != ExitKind::DeliberateSeverance {
6708 exit_report.kind = ExitKind::Crash;
6709 }
6710 exit_report
6711}
6712
6713async fn handle_reload_child_registration_failure(
6714 spec: &ModuleSpec,
6715 runtime: &SupervisorRuntimeConfig,
6716 registry: &Registry,
6717 process_liveness: &SupervisorProcessLiveness,
6718 snapshot: &SharedSnapshot,
6719 child: &mut Option<SupervisedChild>,
6720 failure: ReloadRegistrationFailure,
6721) -> Result<(), SuperviseError> {
6722 let ReloadRegistrationFailure {
6723 exit_report,
6724 reason,
6725 } = failure;
6726 match on_child_exit(
6727 spec,
6728 runtime.restart_policy,
6729 registry,
6730 snapshot,
6731 &runtime.terminal_ring,
6732 &runtime.spawn_events,
6733 &runtime.child_roster,
6734 exit_report,
6735 )
6736 .await
6737 {
6738 NextAction::Stop {
6739 registration_released,
6740 } => {
6741 if registration_released {
6742 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6743 }
6744 }
6745 NextAction::Restart { schedule } => {
6746 let delay = schedule.map_or(runtime.restart_policy.delay_for_restart(0), |schedule| {
6747 schedule.delay
6748 });
6749 if let Some(schedule) = schedule {
6750 log_crash_respawn(&spec.module_id, schedule);
6751 }
6752 sleep(delay).await;
6753 if respawn_still_pending(snapshot) {
6757 if let Err(err) = wait_for_registration_release(
6758 registry,
6759 &spec.module_id,
6760 REGISTRY_RELEASE_TIMEOUT,
6761 )
6762 .await
6763 {
6764 fail_snapshot(snapshot, Some(&spec.module_id), None);
6765 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6766 return Err(SuperviseError::ReloadFailed {
6767 module_id: spec.module_id.clone(),
6768 reason: format!(
6769 "{reason}; registration did not release before policy retry: {err}"
6770 ),
6771 });
6772 }
6773 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6774 match spawn_and_mark_running(spec, runtime, snapshot) {
6775 Ok(next_child) => {
6776 *child = Some(next_child);
6777 }
6778 Err(err) => {
6779 fail_snapshot(snapshot, Some(&spec.module_id), None);
6780 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6781 return Err(SuperviseError::ReloadFailed {
6782 module_id: spec.module_id.clone(),
6783 reason: format!("{reason}; policy retry spawn failed: {err}"),
6784 });
6785 }
6786 }
6787 }
6788 }
6789 }
6790
6791 Err(SuperviseError::ReloadFailed {
6792 module_id: spec.module_id.clone(),
6793 reason,
6794 })
6795}
6796
6797async fn handle_reload_spawn_failure(
6798 spec: &ModuleSpec,
6799 runtime: &SupervisorRuntimeConfig,
6800 process_liveness: &SupervisorProcessLiveness,
6801 snapshot: &SharedSnapshot,
6802 child: &mut Option<SupervisedChild>,
6803 reason: String,
6804) -> Result<(), SuperviseError> {
6805 let mut should_retry = false;
6806 let now = Instant::now();
6807 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6808 clear_current_process_facts(state);
6809 if daemon_will_restart(state, &runtime.restart_policy, now) {
6810 state.record_crash_restart(&runtime.restart_policy, now);
6811 state.state = ModuleState::Restarting;
6812 should_retry = true;
6813 } else if state.enabled {
6814 state.state = ModuleState::Failed;
6815 } else {
6816 state.state = ModuleState::Disabled;
6817 }
6818 })?;
6819
6820 if should_retry {
6821 sleep(runtime.restart_policy.backoff).await;
6822 if respawn_still_pending(snapshot) {
6826 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6827 match spawn_and_mark_running(spec, runtime, snapshot) {
6828 Ok(next_child) => {
6829 *child = Some(next_child);
6830 }
6831 Err(err) => {
6832 fail_snapshot(snapshot, Some(&spec.module_id), None);
6833 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6834 return Err(SuperviseError::ReloadFailed {
6835 module_id: spec.module_id.clone(),
6836 reason: format!("{reason}; policy retry spawn failed: {err}"),
6837 });
6838 }
6839 }
6840 }
6841 } else {
6842 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6843 }
6844
6845 Err(SuperviseError::ReloadFailed {
6846 module_id: spec.module_id.clone(),
6847 reason,
6848 })
6849}
6850
6851fn control_flags() -> Flags {
6852 Flags::new(false, Priority::Passive, false)
6853}
6854
6855#[allow(clippy::too_many_arguments)]
6856async fn drain_optional_child(
6857 module_id: &str,
6858 protocol: ModuleProtocol,
6859 stop_notice: StopNotice,
6860 registry: &Registry,
6861 snapshot: &SharedSnapshot,
6862 terminal_ring: &Arc<Mutex<TerminalRing>>,
6863 spawn_events: &SpawnEventFeed,
6864 child: &mut Option<SupervisedChild>,
6865 drain_timeout: Duration,
6866 final_state: ModuleState,
6867 enabled: Option<bool>,
6868) -> Result<(), SuperviseError> {
6869 if let Some(child) = child.take() {
6870 drain_child_to_state(
6871 module_id,
6872 protocol,
6873 stop_notice,
6874 registry,
6875 snapshot,
6876 terminal_ring,
6877 spawn_events,
6878 child,
6879 drain_timeout,
6880 final_state,
6881 enabled,
6882 )
6883 .await
6884 } else {
6885 update_snapshot(snapshot, Some(module_id), |state| {
6886 state.state = final_state;
6887 if let Some(enabled) = enabled {
6888 state.enabled = enabled;
6889 }
6890 clear_current_process_facts(state);
6891 })?;
6892 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
6893 }
6894}
6895
6896#[allow(clippy::too_many_arguments)]
6897async fn drain_child_to_state(
6898 module_id: &str,
6899 protocol: ModuleProtocol,
6900 stop_notice: StopNotice,
6901 registry: &Registry,
6902 snapshot: &SharedSnapshot,
6903 terminal_ring: &Arc<Mutex<TerminalRing>>,
6904 spawn_events: &SpawnEventFeed,
6905 mut child: SupervisedChild,
6906 drain_timeout: Duration,
6907 final_state: ModuleState,
6908 enabled: Option<bool>,
6909) -> Result<(), SuperviseError> {
6910 update_snapshot(snapshot, Some(module_id), |state| {
6911 state.state = ModuleState::Draining;
6912 state.draining_to_replace = final_state == ModuleState::Restarting;
6913 if let Some(enabled) = enabled {
6914 state.enabled = enabled;
6915 }
6916 })?;
6917
6918 if stop_notice != StopNotice::SentOverConnection {
6929 if protocol == ModuleProtocol::Subc && stop_notice == StopNotice::NoConnection {
6930 info!(
6931 module_id,
6932 pid = child.pid,
6933 budget_ms = u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX),
6934 "module has no connection yet; requesting stop by signal"
6935 );
6936 }
6937 request_graceful_stop(module_id, &child);
6938 }
6939
6940 let exit_report = match timeout(drain_timeout, child.wait()).await {
6941 Ok(Ok(status)) => classify_reaped_child_exit(snapshot, &child, &status),
6942 Ok(Err(source)) => {
6943 fail_snapshot(snapshot, Some(module_id), None);
6944 return Err(SuperviseError::Wait {
6945 module_id: module_id.to_string(),
6946 source,
6947 });
6948 }
6949 Err(_) => {
6950 warn!(
6963 module_id,
6964 pid = child.pid,
6965 budget_ms = u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX),
6966 reason = ?final_state,
6967 ?stop_notice,
6968 "drain budget expired before the module exited; killing it"
6969 );
6970 child.start_kill().map_err(|source| {
6971 fail_snapshot(snapshot, Some(module_id), None);
6972 SuperviseError::Kill {
6973 module_id: module_id.to_string(),
6974 source,
6975 }
6976 })?;
6977 let status = child.wait().await.map_err(|source| {
6978 fail_snapshot(snapshot, Some(module_id), None);
6979 SuperviseError::Wait {
6980 module_id: module_id.to_string(),
6981 source,
6982 }
6983 })?;
6984 classify_reaped_child_exit(snapshot, &child, &status)
6985 }
6986 };
6987
6988 update_snapshot(snapshot, Some(module_id), |state| {
6989 state.state = final_state;
6990 if let Some(enabled) = enabled {
6991 state.enabled = enabled;
6992 }
6993 clear_current_process_facts(state);
6994 state.last_exit = Some(exit_report.clone());
6995 if exit_report.kind == ExitKind::DeliberateSeverance {
6996 state.lifetime_restarts += 1;
6997 }
6998 })?;
6999 record_terminal(
7000 module_id,
7001 terminal_ring,
7002 spawn_events,
7003 &exit_report,
7004 terminal_disposition(final_state),
7005 );
7006 child.drain_stderr(module_id).await;
7007
7008 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
7009}
7010
7011#[cfg(unix)]
7031fn request_graceful_stop(module_id: &str, child: &SupervisedChild) {
7032 let Some(pid) = child
7033 .id()
7034 .and_then(|pid| i32::try_from(pid).ok())
7035 .and_then(rustix::process::Pid::from_raw)
7036 else {
7037 debug!(
7038 module_id,
7039 "no pid to signal for teardown; falling through to the drain wait"
7040 );
7041 return;
7042 };
7043 match rustix::process::kill_process(pid, rustix::process::Signal::TERM) {
7044 Ok(()) => debug!(
7045 module_id,
7046 "sent SIGTERM to a module nothing else asked to stop"
7047 ),
7048 Err(err) => debug!(
7049 module_id,
7050 error = %err,
7051 "SIGTERM to module failed; the drain wait and kill still apply"
7052 ),
7053 }
7054}
7055
7056#[cfg(not(unix))]
7064fn request_graceful_stop(module_id: &str, _child: &SupervisedChild) {
7065 debug!(
7066 module_id,
7067 "no graceful stop signal exists on this platform; teardown of a module nothing asked to stop waits, then kills"
7068 );
7069}
7070
7071fn terminal_disposition(final_state: ModuleState) -> TerminalDisposition {
7072 match final_state {
7073 ModuleState::Stopped => TerminalDisposition::Stopped,
7074 ModuleState::Disabled => TerminalDisposition::Disabled,
7075 ModuleState::Restarting => TerminalDisposition::Restarting,
7076 ModuleState::Failed => TerminalDisposition::Failed,
7077 ModuleState::Starting
7078 | ModuleState::Running
7079 | ModuleState::Unresponsive
7080 | ModuleState::Draining => {
7081 unreachable!("terminal exits only finish in terminal or restarting states")
7082 }
7083 }
7084}
7085
7086async fn wait_for_registration_release(
7089 registry: &Registry,
7090 module_id: &str,
7091 wait: Duration,
7092) -> Result<(), SuperviseError> {
7093 wait_for_slot_registration_release(
7094 registry,
7095 crate::registry::RegistrationSlot::Active(module_id),
7096 wait,
7097 )
7098 .await
7099}
7100
7101async fn wait_for_slot_registration_release(
7109 registry: &Registry,
7110 slot: crate::registry::RegistrationSlot<'_>,
7111 wait: Duration,
7112) -> Result<(), SuperviseError> {
7113 let deadline = Instant::now() + wait;
7114 let mut release_events = registration_release_events().subscribe();
7115 let still_active = |registration: &crate::registry::ModuleRegistration| {
7116 SuperviseError::RegistrationStillActive {
7117 module_id: registration.manifest.module_id.clone(),
7118 waited: wait,
7119 }
7120 };
7121 loop {
7122 let _observed_generation = *release_events.borrow_and_update();
7123 let Some(registration) = registry
7124 .registration(slot)
7125 .map_err(SuperviseError::Registry)?
7126 else {
7127 return Ok(());
7128 };
7129
7130 let now = Instant::now();
7131 if now >= deadline {
7132 return Err(still_active(®istration));
7133 }
7134
7135 let remaining = deadline.saturating_duration_since(now);
7136 match timeout(remaining, release_events.changed()).await {
7137 Ok(Ok(())) | Ok(Err(_)) => {}
7138 Err(_) => return Err(still_active(®istration)),
7139 }
7140 }
7141}
7142
7143#[cfg(test)]
7144mod slot_registration_wait_tests {
7145 use super::*;
7146 use crate::registry::{ConnectionId, RegistrationSlot};
7147 use subc_protocol::manifest::ModuleManifest;
7148
7149 const INCUMBENT: u64 = 1;
7150 const CANDIDATE: u64 = 2;
7151
7152 fn swapped_registry() -> Arc<Registry> {
7153 let registry = Arc::new(Registry::default());
7154 let manifest = ModuleManifest::builder("m", "0.1.0").build();
7155 registry
7156 .register_with_control_ops(
7157 manifest.clone(),
7158 1,
7159 ConnectionId::new(INCUMBENT),
7160 Vec::new(),
7161 )
7162 .unwrap();
7163 registry
7164 .register_candidate_with_control_ops(
7165 manifest,
7166 1,
7167 ConnectionId::new(CANDIDATE),
7168 Vec::new(),
7169 )
7170 .unwrap();
7171 registry
7172 }
7173
7174 #[tokio::test]
7178 async fn incumbent_release_is_awaited_by_connection_not_by_module_id() {
7179 let registry = swapped_registry();
7180 registry.promote_candidate("m").unwrap().unwrap();
7181
7182 assert!(matches!(
7183 wait_for_registration_release(®istry, "m", Duration::from_millis(50)).await,
7184 Err(SuperviseError::RegistrationStillActive { .. })
7185 ));
7186
7187 assert!(matches!(
7189 wait_for_slot_registration_release(
7190 ®istry,
7191 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
7192 Duration::from_millis(50),
7193 )
7194 .await,
7195 Err(SuperviseError::RegistrationStillActive { .. })
7196 ));
7197
7198 let releaser = Arc::clone(®istry);
7199 let release = tokio::spawn(async move {
7200 sleep(Duration::from_millis(20)).await;
7201 releaser
7202 .deregister_connection(ConnectionId::new(INCUMBENT))
7203 .unwrap();
7204 notify_registration_release();
7205 });
7206 wait_for_slot_registration_release(
7207 ®istry,
7208 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
7209 Duration::from_secs(5),
7210 )
7211 .await
7212 .expect("the incumbent's own registration is released");
7213 release.await.unwrap();
7214 assert!(registry.get_module("m").unwrap().is_some());
7215 }
7216
7217 #[tokio::test]
7220 async fn candidate_slot_wait_ignores_the_incumbents_registration() {
7221 let registry = swapped_registry();
7222 assert!(matches!(
7223 wait_for_slot_registration_release(
7224 ®istry,
7225 RegistrationSlot::Candidate("m"),
7226 Duration::from_millis(50),
7227 )
7228 .await,
7229 Err(SuperviseError::RegistrationStillActive { .. })
7230 ));
7231 registry
7232 .deregister_connection(ConnectionId::new(CANDIDATE))
7233 .unwrap();
7234 wait_for_slot_registration_release(
7235 ®istry,
7236 RegistrationSlot::Candidate("m"),
7237 Duration::from_millis(50),
7238 )
7239 .await
7240 .expect("a candidate slot with no candidate is released");
7241 assert!(registry
7242 .registration(RegistrationSlot::Active("m"))
7243 .unwrap()
7244 .is_some());
7245 }
7246}
7247
7248fn classify_exit(status: &ExitStatus) -> ExitReport {
7249 ExitReport {
7250 kind: if status.success() {
7251 ExitKind::Clean
7252 } else {
7253 ExitKind::Crash
7254 },
7255 code: status.code(),
7256 signal: exit_signal(status),
7257 at_ms: unix_ms_now(),
7258 }
7259}
7260
7261fn wait_error_exit_report() -> ExitReport {
7267 ExitReport {
7268 kind: ExitKind::Crash,
7269 code: None,
7270 signal: None,
7271 at_ms: unix_ms_now(),
7272 }
7273}
7274
7275#[cfg(unix)]
7276fn exit_signal(status: &ExitStatus) -> Option<i32> {
7277 use std::os::unix::process::ExitStatusExt;
7278
7279 status.signal()
7280}
7281
7282#[cfg(not(unix))]
7283fn exit_signal(_status: &ExitStatus) -> Option<i32> {
7284 None
7285}
7286
7287fn reset_restart_count(snapshot: &SharedSnapshot, module_id: &str) -> Result<(), SuperviseError> {
7293 update_snapshot(snapshot, Some(module_id), |state| {
7294 state.clear_crash_restarts();
7295 })
7296}
7297
7298fn set_running(
7299 snapshot: &SharedSnapshot,
7300 child: &SupervisedChild,
7301 module_id: &str,
7302 spawn_events: &SpawnEventFeed,
7303) -> Result<(), SuperviseError> {
7304 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
7305 module_id: Some(module_id.to_string()),
7306 })?;
7307 state.spawn_generation = spawn_events.emit_spawned(module_id, child.pid, child.spawned_at_ms);
7308 state.in_alternate_slot = false;
7311 state.configuration_updated_since_spawn = false;
7312 state.state = ModuleState::Running;
7313 state.enabled = true;
7314 state.process_alive = true;
7315 state.pid = child.id();
7316 state.spawned_at_ms = Some(child.spawned_at_ms);
7317 state.spawned_from = Some(child.spawned_from.clone());
7318 state.spawned_file_identity = child.spawned_file_identity;
7319 state.process_start_time = child.process_start_time;
7320 Ok(())
7321}
7322
7323fn clear_current_process_facts(state: &mut SupervisorSnapshot) {
7324 state.process_alive = false;
7325 state.pid = None;
7326 state.spawned_at_ms = None;
7327 state.spawned_from = None;
7328 state.spawned_file_identity = None;
7329 state.process_start_time = None;
7330 state.deliberate_severance = None;
7331}
7332
7333#[cfg(test)]
7334fn record_deliberate_severance(
7335 snapshot: &SharedSnapshot,
7336 identity: ProcessIdentity,
7337) -> Result<(), SuperviseError> {
7338 update_snapshot(snapshot, None, |state| {
7339 state.deliberate_severance = Some(identity);
7340 })
7341}
7342
7343fn apply_deliberate_severance_marker(
7344 snapshot: &SharedSnapshot,
7345 exited_identity: Option<ProcessIdentity>,
7346 mut exit_report: ExitReport,
7347) -> ExitReport {
7348 let marker = lock_snapshot(snapshot)
7349 .ok()
7350 .and_then(|mut state| state.deliberate_severance.take());
7351 if marker.is_some() && marker == exited_identity {
7352 exit_report.kind = ExitKind::DeliberateSeverance;
7353 }
7354 exit_report
7355}
7356
7357fn classify_reaped_child_exit(
7358 snapshot: &SharedSnapshot,
7359 child: &SupervisedChild,
7360 status: &ExitStatus,
7361) -> ExitReport {
7362 apply_deliberate_severance_marker(snapshot, child.process_identity(), classify_exit(status))
7363}
7364
7365fn fail_snapshot(
7366 snapshot: &SharedSnapshot,
7367 module_id: Option<&str>,
7368 last_exit: Option<ExitReport>,
7369) {
7370 if let Err(err) = update_snapshot(snapshot, module_id, |state| {
7371 state.state = ModuleState::Failed;
7372 clear_current_process_facts(state);
7373 if let Some(last_exit) = last_exit {
7374 state.last_exit = Some(last_exit);
7375 }
7376 }) {
7377 error!(error = %err, "failed to mark supervisor state failed");
7378 }
7379}
7380
7381fn update_snapshot(
7382 snapshot: &SharedSnapshot,
7383 module_id: Option<&str>,
7384 update: impl FnOnce(&mut SupervisorSnapshot),
7385) -> Result<(), SuperviseError> {
7386 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
7387 module_id: module_id.map(ToOwned::to_owned),
7388 })?;
7389 update(&mut state);
7390 Ok(())
7391}
7392
7393const SLOW_SNAPSHOT_LOCK_THRESHOLD: Duration = Duration::from_millis(250);
7394
7395fn lock_snapshot_for_control<'a>(
7396 snapshot: &'a SharedSnapshot,
7397 module_id: &str,
7398 caller: &'static str,
7399) -> Result<std::sync::MutexGuard<'a, SupervisorSnapshot>, SuperviseError> {
7400 let started_at = Instant::now();
7401 let guard = lock_snapshot(snapshot)?;
7402 let waited = started_at.elapsed();
7403 if waited >= SLOW_SNAPSHOT_LOCK_THRESHOLD {
7404 warn!(
7405 module_id = %module_id,
7406 waited_ms = waited.as_millis() as u64,
7407 caller = %caller,
7408 "slow snapshot lock"
7409 );
7410 }
7411 Ok(guard)
7412}
7413
7414fn lock_snapshot(
7415 snapshot: &SharedSnapshot,
7416) -> Result<std::sync::MutexGuard<'_, SupervisorSnapshot>, SuperviseError> {
7417 snapshot
7418 .lock()
7419 .map_err(|_| SuperviseError::StatePoisoned { module_id: None })
7420}
7421
7422#[cfg(test)]
7423mod terminal_history_tests {
7424 use std::{
7425 path::PathBuf,
7426 sync::Arc,
7427 time::{Duration, Instant},
7428 };
7429
7430 use tokio::time::sleep;
7431
7432 use super::{
7433 apply_deliberate_severance_marker, daemon_will_restart, drain_child_to_state,
7434 drained_after_quiescence_wait, handle_reload_spawn_failure, health_restart_child,
7435 lock_snapshot, on_child_exit, record_deliberate_severance, record_wait_error_terminal,
7436 reset_restart_count, spawn_and_mark_running, update_snapshot, wait_error_exit_report,
7437 ExitKind, ExitReport, ModuleProtocol, ModuleSpec, ModuleState, NextAction, ProcessIdentity,
7438 RestartPolicy, SpawnEventKind, StopNotice, SuperviseError, SupervisedModule, Supervisor,
7439 SupervisorHandle, SupervisorHealthStatus, SupervisorSnapshot,
7440 };
7441 use super::Instant as ClockInstant;
7446 use crate::{
7447 registry::Registry,
7448 terminal_ring::{TerminalRing, TerminalRingConfig},
7449 };
7450 use std::sync::Mutex;
7451 use subc_control::TerminalDisposition;
7452
7453 fn fake_aft_stub_path() -> PathBuf {
7458 let mut path = std::env::current_exe().expect("current_exe available in tests");
7459 path.pop();
7460 path.pop();
7461 path.push(if cfg!(windows) {
7462 "fake-aft-stub.exe"
7463 } else {
7464 "fake-aft-stub"
7465 });
7466 assert!(
7467 path.exists(),
7468 "fake-aft-stub not built at {}: run `cargo test -p subc-core` (which builds \
7469 [[bin]] targets) rather than `cargo test -p subc-core --lib` (which does not)",
7470 path.display()
7471 );
7472 path
7473 }
7474
7475 #[test]
7476 fn reserved_never_spawned_refuses_every_hello() {
7477 let supervisor = SupervisorHandle::default();
7482 supervisor.apply_identity_configuration(&ModuleSpec {
7483 module_id: "never-spawned".to_string(),
7484 program: PathBuf::from("/usr/bin/false"),
7485 args: Vec::new(),
7486 env: Vec::new(),
7487 reserved: true,
7488 reserved_prefixes: Vec::new(),
7489 protocol: ModuleProtocol::Subc,
7490 overlap: Default::default(),
7491 });
7492 assert!(
7493 supervisor
7494 .reserved_hello_rejection("never-spawned", Some("any-forged-nonce"))
7495 .is_some(),
7496 "forged nonce must refuse on a reserved never-spawned id"
7497 );
7498 assert!(
7499 supervisor
7500 .reserved_hello_rejection("never-spawned", None)
7501 .is_some(),
7502 "absent nonce must refuse on a reserved never-spawned id"
7503 );
7504 supervisor.set_spawn_nonce("never-spawned", "minted".to_string());
7506 supervisor.apply_identity_configuration(&ModuleSpec {
7507 module_id: "never-spawned".to_string(),
7508 program: PathBuf::from("/usr/bin/false"),
7509 args: Vec::new(),
7510 env: Vec::new(),
7511 reserved: true,
7512 reserved_prefixes: Vec::new(),
7513 protocol: ModuleProtocol::Subc,
7514 overlap: Default::default(),
7515 });
7516 assert!(supervisor
7517 .reserved_hello_rejection("never-spawned", Some("minted"))
7518 .is_none());
7519 assert!(supervisor
7520 .reserved_hello_rejection("never-spawned", Some("forged"))
7521 .is_some());
7522 }
7523
7524 fn seed_crash_restarts(state: &mut SupervisorSnapshot, count: u32) {
7527 let now = ClockInstant::now();
7528 for _ in 0..count {
7529 state.crash_restarts.push_back(now);
7530 }
7531 }
7532
7533 fn age_oldest_crash_restart_out_of_window(state: &mut SupervisorSnapshot, window: Duration) {
7537 let aged = state
7538 .crash_restarts
7539 .front()
7540 .expect("a crash restart must be recorded before it can be aged")
7541 .checked_sub(window + Duration::from_secs(1))
7542 .expect("the test clock is far enough from its origin to age an instant");
7543 state.crash_restarts[0] = aged;
7544 }
7545
7546 fn snapshot_with_restarts(enabled: bool, count: u32) -> SupervisorSnapshot {
7547 let mut state = SupervisorSnapshot::new(ModuleState::Running, enabled);
7548 seed_crash_restarts(&mut state, count);
7549 state
7550 }
7551
7552 #[test]
7553 fn daemon_owned_recovery_predicate_uses_the_pre_increment_budget() {
7554 let policy = RestartPolicy::new(3, Duration::ZERO);
7555 let now = ClockInstant::now();
7556 assert!(daemon_will_restart(
7557 &mut snapshot_with_restarts(true, 2),
7558 &policy,
7559 now
7560 ));
7561 assert!(!daemon_will_restart(
7562 &mut snapshot_with_restarts(true, 3),
7563 &policy,
7564 now
7565 ));
7566 assert!(!daemon_will_restart(
7567 &mut snapshot_with_restarts(false, 0),
7568 &policy,
7569 now
7570 ));
7571 }
7572
7573 #[test]
7574 fn crash_restart_backoff_escalates_with_in_window_count() {
7575 let policy = RestartPolicy::new(4, Duration::from_millis(100))
7576 .with_max_backoff(Duration::from_secs(30));
7577 let now = ClockInstant::now();
7578 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7579 let schedules = (0..4)
7580 .map(|_| {
7581 state
7582 .next_crash_restart(&policy, now)
7583 .expect("the test policy allows four crash restarts")
7584 })
7585 .collect::<Vec<_>>();
7586
7587 assert_eq!(
7588 schedules
7589 .iter()
7590 .map(|schedule| schedule.restart_in_window)
7591 .collect::<Vec<_>>(),
7592 vec![0, 1, 2, 3]
7593 );
7594 assert_eq!(
7595 schedules
7596 .iter()
7597 .map(|schedule| schedule.delay)
7598 .collect::<Vec<_>>(),
7599 vec![
7600 Duration::from_millis(100),
7601 Duration::from_secs(1),
7602 Duration::from_secs(10),
7603 Duration::from_secs(30),
7604 ]
7605 );
7606 }
7607
7608 #[test]
7609 fn crash_restart_backoff_resets_after_ring_clear() {
7610 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7611 let now = ClockInstant::now();
7612 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7613 assert_eq!(
7614 state.next_crash_restart(&policy, now).unwrap().delay,
7615 Duration::from_millis(100)
7616 );
7617 assert_eq!(
7618 state.next_crash_restart(&policy, now).unwrap().delay,
7619 Duration::from_secs(1)
7620 );
7621
7622 state.clear_crash_restarts();
7623 let schedule = state
7624 .next_crash_restart(&policy, now)
7625 .expect("a cleared ring must allow another restart");
7626 assert_eq!(schedule.restart_in_window, 0);
7627 assert_eq!(schedule.delay, Duration::from_millis(100));
7628 }
7629
7630 #[test]
7631 fn crash_restart_backoff_ignores_aged_restarts() {
7632 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7633 let now = ClockInstant::now();
7634 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7635 state
7636 .next_crash_restart(&policy, now)
7637 .expect("the first restart is allowed");
7638 state
7639 .next_crash_restart(&policy, now)
7640 .expect("the second restart is allowed");
7641 state.crash_restarts[0] = now
7642 .checked_sub(policy.window + Duration::from_secs(1))
7643 .expect("the fake clock can age a restart past the window");
7644
7645 let schedule = state
7646 .next_crash_restart(&policy, now)
7647 .expect("an aged restart must release its slot");
7648 assert_eq!(schedule.restart_in_window, 1);
7649 assert_eq!(schedule.delay, Duration::from_secs(1));
7650 assert_eq!(state.crash_restarts.len(), 2);
7651 }
7652
7653 #[test]
7657 fn a_budget_spent_before_the_window_no_longer_refuses() {
7658 let policy = RestartPolicy::new(3, Duration::ZERO);
7659 let mut state = snapshot_with_restarts(true, 3);
7660 let now = ClockInstant::now();
7661 assert!(!daemon_will_restart(&mut state, &policy, now));
7662
7663 assert!(daemon_will_restart(
7664 &mut state,
7665 &policy,
7666 now + policy.window + Duration::from_secs(1)
7667 ));
7668 assert!(
7669 state.crash_restarts.is_empty(),
7670 "reading the budget must drop the instants that left the window"
7671 );
7672 }
7673
7674 fn module_with_recovery_snapshot(
7675 state: ModuleState,
7676 enabled: bool,
7677 restart_count: u32,
7678 ) -> SupervisedModule {
7679 let registry = Arc::new(Registry::default());
7680 let supervisor =
7681 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(3, Duration::ZERO));
7682 let module = supervisor
7683 .spawn(ModuleSpec {
7684 module_id: "recovery-snapshot".to_string(),
7685 program: fake_aft_stub_path(),
7686 args: Vec::new(),
7687 env: Vec::new(),
7688 reserved: false,
7689 reserved_prefixes: Vec::new(),
7690 protocol: ModuleProtocol::Subc,
7691 overlap: Default::default(),
7692 })
7693 .unwrap();
7694 update_snapshot(
7695 &module.inner.snapshot,
7696 Some("recovery-snapshot"),
7697 |snapshot| {
7698 snapshot.state = state;
7699 snapshot.enabled = enabled;
7700 seed_crash_restarts(snapshot, restart_count);
7701 },
7702 )
7703 .unwrap();
7704 module
7705 }
7706
7707 #[cfg(target_os = "linux")]
7708 #[tokio::test]
7709 async fn no_cgroup_placement_does_not_block_fake_aft_stub_spawn() {
7710 let supervisor = Supervisor::new(Arc::new(Registry::default()), RestartPolicy::default())
7711 .with_cgroup_placement(None);
7712 let result = supervisor.spawn(ModuleSpec {
7713 module_id: "no-cgroup-placement".to_string(),
7714 program: fake_aft_stub_path(),
7715 args: Vec::new(),
7716 env: Vec::new(),
7717 reserved: false,
7718 reserved_prefixes: Vec::new(),
7719 protocol: ModuleProtocol::Subc,
7720 overlap: Default::default(),
7721 });
7722
7723 assert!(
7724 result.is_ok(),
7725 "no delegation must not turn an otherwise valid spawn into a failure: {result:?}"
7726 );
7727 }
7728
7729 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7730 async fn undecided_snapshot_uses_shared_restart_predicate() {
7731 assert!(module_with_recovery_snapshot(ModuleState::Running, true, 2)
7732 .will_recover_after_connection_loss()
7733 .unwrap());
7734 assert!(
7735 !module_with_recovery_snapshot(ModuleState::Running, true, 3)
7736 .will_recover_after_connection_loss()
7737 .unwrap()
7738 );
7739 }
7740
7741 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7742 async fn restarting_snapshot_at_exhausted_budget_is_non_terminal() {
7743 assert!(
7744 module_with_recovery_snapshot(ModuleState::Restarting, true, 3)
7745 .will_recover_after_connection_loss()
7746 .unwrap()
7747 );
7748 }
7749
7750 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7751 async fn terminal_phase_snapshots_are_terminal_before_budget_exhaustion() {
7752 assert!(!module_with_recovery_snapshot(ModuleState::Failed, true, 0)
7753 .will_recover_after_connection_loss()
7754 .unwrap());
7755 assert!(
7756 !module_with_recovery_snapshot(ModuleState::Disabled, true, 0)
7757 .will_recover_after_connection_loss()
7758 .unwrap()
7759 );
7760 }
7761
7762 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7763 async fn warming_snapshot_is_limited_to_startup_phases() {
7764 for state in [
7765 ModuleState::Starting,
7766 ModuleState::Running,
7767 ModuleState::Restarting,
7768 ] {
7769 assert!(
7770 module_with_recovery_snapshot(state, true, 0)
7771 .is_warming()
7772 .unwrap(),
7773 "{state:?} should be warming"
7774 );
7775 }
7776 for state in [
7777 ModuleState::Unresponsive,
7778 ModuleState::Draining,
7779 ModuleState::Stopped,
7780 ModuleState::Failed,
7781 ModuleState::Disabled,
7782 ] {
7783 assert!(
7784 !module_with_recovery_snapshot(state, true, 0)
7785 .is_warming()
7786 .unwrap(),
7787 "{state:?} should not be warming"
7788 );
7789 }
7790 }
7791
7792 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7793 async fn terminal_history_survives_respawn_and_keeps_both_crashes_in_order() {
7794 let registry = Arc::new(Registry::default());
7795 let supervisor =
7796 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(1, Duration::ZERO));
7797 let module = supervisor
7798 .spawn(ModuleSpec {
7799 module_id: "terminal-history".to_string(),
7800 program: fake_aft_stub_path(),
7801 args: Vec::new(),
7802 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7803 reserved: false,
7804 reserved_prefixes: Vec::new(),
7805 protocol: ModuleProtocol::Subc,
7806 overlap: Default::default(),
7807 })
7808 .unwrap();
7809
7810 let deadline = Instant::now() + Duration::from_secs(5);
7811 loop {
7812 let history = module.terminal_history();
7813 if history.entries.len() == 2 {
7814 assert_eq!(module.status().unwrap().state, ModuleState::Failed);
7815 assert_eq!(history.dropped, 0);
7816 assert_eq!(
7817 history
7818 .entries
7819 .iter()
7820 .map(|entry| entry.exit_code)
7821 .collect::<Vec<_>>(),
7822 vec![Some(23), Some(23)]
7823 );
7824 assert!(history.entries[0].at_ms <= history.entries[1].at_ms);
7825 return;
7826 }
7827 assert!(
7828 Instant::now() < deadline,
7829 "module did not retain two terminal exits: {history:?}"
7830 );
7831 sleep(Duration::from_millis(10)).await;
7832 }
7833 }
7834
7835 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7839 async fn disable_during_crash_backoff_cancels_pending_respawn() {
7840 let backoff = Duration::from_secs(2);
7841 let supervisor = Supervisor::new(
7842 Arc::new(Registry::default()),
7843 RestartPolicy::new(10, backoff),
7844 );
7845 let module = supervisor
7846 .spawn(ModuleSpec {
7847 module_id: "disable-during-backoff".to_string(),
7848 program: fake_aft_stub_path(),
7849 args: Vec::new(),
7850 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7851 reserved: false,
7852 reserved_prefixes: Vec::new(),
7853 protocol: ModuleProtocol::Subc,
7854 overlap: Default::default(),
7855 })
7856 .unwrap();
7857
7858 let deadline = Instant::now() + Duration::from_secs(5);
7860 loop {
7861 if module.status().unwrap().state == ModuleState::Restarting {
7862 break;
7863 }
7864 assert!(
7865 Instant::now() < deadline,
7866 "module never entered the crash backoff"
7867 );
7868 sleep(Duration::from_millis(10)).await;
7869 }
7870
7871 let started = Instant::now();
7872 module.set_enabled(false).await.unwrap();
7873 let waited = started.elapsed();
7874
7875 assert!(
7876 waited < backoff / 2,
7877 "disable waited {waited:?} behind the {backoff:?} crash backoff; the operator command must preempt the pending respawn"
7878 );
7879 assert_eq!(module.status().unwrap().state, ModuleState::Disabled);
7880
7881 sleep(backoff + Duration::from_millis(500)).await;
7883 let status = module.status().unwrap();
7884 assert_eq!(status.state, ModuleState::Disabled);
7885 assert_eq!(
7886 status.spawn_generation, 1,
7887 "module respawned after the operator disabled it"
7888 );
7889 }
7890
7891 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7895 async fn every_restart_increment_path_advances_lifetime_count() {
7896 let supervisor = Supervisor::new(
7897 Arc::new(Registry::default()),
7898 RestartPolicy::new(1, Duration::ZERO),
7899 );
7900 let runtime = supervisor.runtime_config();
7901 let spec = ModuleSpec {
7902 module_id: "lifetime-increment-path".to_string(),
7903 program: PathBuf::from("/unused/lifetime-increment-path"),
7904 args: Vec::new(),
7905 env: Vec::new(),
7906 reserved: false,
7907 reserved_prefixes: Vec::new(),
7908 protocol: ModuleProtocol::Subc,
7909 overlap: Default::default(),
7910 };
7911
7912 let crash_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7913 assert!(matches!(
7914 on_child_exit(
7915 &spec,
7916 runtime.restart_policy,
7917 &supervisor.registry,
7918 &crash_snapshot,
7919 &runtime.terminal_ring,
7920 &runtime.spawn_events,
7921 &runtime.child_roster,
7922 ExitReport {
7923 kind: ExitKind::Crash,
7924 code: Some(1),
7925 signal: None,
7926 at_ms: 1,
7927 },
7928 )
7929 .await,
7930 NextAction::Restart { schedule: _ }
7931 ));
7932 let (crash_restarts, crash_lifetime) = {
7933 let state = lock_snapshot(&crash_snapshot).unwrap();
7934 (state.crash_restarts.len(), state.lifetime_restarts)
7935 };
7936 assert_eq!(crash_restarts, 1);
7937 assert_eq!(crash_lifetime, 1);
7938
7939 let health_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7940 let mut health_child = None;
7941 assert!(matches!(
7942 health_restart_child(
7943 &spec,
7944 &runtime,
7945 &supervisor.registry,
7946 &supervisor.process_liveness,
7947 &health_snapshot,
7948 &mut health_child,
7949 SupervisorHealthStatus::Failing,
7950 None,
7951 2,
7952 )
7953 .await,
7954 Err(SuperviseError::Spawn { .. })
7955 ));
7956 let (health_restarts, health_lifetime) = {
7957 let state = lock_snapshot(&health_snapshot).unwrap();
7958 (state.crash_restarts.len(), state.lifetime_restarts)
7959 };
7960 assert_eq!(health_restarts, 1);
7961 assert_eq!(health_lifetime, 1);
7962
7963 let reload_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7964 let mut reload_child = None;
7965 assert!(matches!(
7966 handle_reload_spawn_failure(
7967 &spec,
7968 &runtime,
7969 &supervisor.process_liveness,
7970 &reload_snapshot,
7971 &mut reload_child,
7972 "forced reload spawn failure".to_string(),
7973 )
7974 .await,
7975 Err(SuperviseError::ReloadFailed { .. })
7976 ));
7977 let (reload_restarts, reload_lifetime) = {
7978 let state = lock_snapshot(&reload_snapshot).unwrap();
7979 (state.crash_restarts.len(), state.lifetime_restarts)
7980 };
7981 assert_eq!(reload_restarts, 1);
7982 assert_eq!(reload_lifetime, 1);
7983 }
7984
7985 #[tokio::test]
7986 async fn deliberately_severed_live_child_records_lifetime_without_spending_restart_budget() {
7987 let supervisor = Supervisor::new(
7988 Arc::new(Registry::default()),
7989 RestartPolicy::new(3, Duration::ZERO),
7990 );
7991 let runtime = supervisor.runtime_config();
7992 let spec = ModuleSpec {
7993 module_id: "deliberately-severed".to_string(),
7994 program: PathBuf::from("/unused/deliberately-severed"),
7995 args: Vec::new(),
7996 env: Vec::new(),
7997 reserved: false,
7998 reserved_prefixes: Vec::new(),
7999 protocol: ModuleProtocol::Subc,
8000 overlap: Default::default(),
8001 };
8002 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8003 let process = ProcessIdentity {
8004 pid: 41,
8005 start_time: 101,
8006 };
8007 record_deliberate_severance(&snapshot, process).unwrap();
8008 let exit_report = apply_deliberate_severance_marker(
8009 &snapshot,
8010 Some(process),
8011 ExitReport {
8012 kind: ExitKind::Crash,
8013 code: Some(1),
8014 signal: None,
8015 at_ms: 1,
8016 },
8017 );
8018 assert_eq!(exit_report.kind, ExitKind::DeliberateSeverance);
8019
8020 assert!(matches!(
8021 on_child_exit(
8022 &spec,
8023 runtime.restart_policy,
8024 &supervisor.registry,
8025 &snapshot,
8026 &runtime.terminal_ring,
8027 &runtime.spawn_events,
8028 &runtime.child_roster,
8029 exit_report,
8030 )
8031 .await,
8032 NextAction::Restart { schedule: _ }
8033 ));
8034 let state = lock_snapshot(&snapshot).unwrap();
8035 assert_eq!(state.lifetime_restarts, 1);
8036 assert_eq!(state.crash_restarts.len(), 0);
8037 }
8038
8039 #[tokio::test]
8040 async fn genuine_crash_spends_restart_budget_and_records_lifetime() {
8041 let supervisor = Supervisor::new(
8042 Arc::new(Registry::default()),
8043 RestartPolicy::new(3, Duration::ZERO),
8044 );
8045 let runtime = supervisor.runtime_config();
8046 let spec = ModuleSpec {
8047 module_id: "genuine-crash".to_string(),
8048 program: PathBuf::from("/unused/genuine-crash"),
8049 args: Vec::new(),
8050 env: Vec::new(),
8051 reserved: false,
8052 reserved_prefixes: Vec::new(),
8053 protocol: ModuleProtocol::Subc,
8054 overlap: Default::default(),
8055 };
8056 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8057
8058 assert!(matches!(
8059 on_child_exit(
8060 &spec,
8061 runtime.restart_policy,
8062 &supervisor.registry,
8063 &snapshot,
8064 &runtime.terminal_ring,
8065 &runtime.spawn_events,
8066 &runtime.child_roster,
8067 ExitReport {
8068 kind: ExitKind::Crash,
8069 code: Some(1),
8070 signal: None,
8071 at_ms: 1,
8072 },
8073 )
8074 .await,
8075 NextAction::Restart { schedule: _ }
8076 ));
8077 let state = lock_snapshot(&snapshot).unwrap();
8078 assert_eq!(state.lifetime_restarts, 1);
8079 assert_eq!(state.crash_restarts.len(), 1);
8080 }
8081
8082 fn crash_exit_report(at_ms: u64) -> ExitReport {
8083 ExitReport {
8084 kind: ExitKind::Crash,
8085 code: Some(1),
8086 signal: None,
8087 at_ms,
8088 }
8089 }
8090
8091 fn windowed_crash_spec(module_id: &str) -> ModuleSpec {
8092 ModuleSpec {
8093 module_id: module_id.to_string(),
8094 program: PathBuf::from("/unused").join(module_id),
8095 args: Vec::new(),
8096 env: Vec::new(),
8097 reserved: false,
8098 reserved_prefixes: Vec::new(),
8099 protocol: ModuleProtocol::Subc,
8100 overlap: Default::default(),
8101 }
8102 }
8103
8104 #[tokio::test]
8110 async fn three_crashes_inside_the_window_stop_the_module_and_name_the_window() {
8111 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::ERROR);
8112 let supervisor = Supervisor::new(
8113 Arc::new(Registry::default()),
8114 RestartPolicy::new(2, Duration::ZERO),
8115 );
8116 let runtime = supervisor.runtime_config();
8117 let spec = windowed_crash_spec("crash-loop-in-window");
8118 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8119
8120 for attempt in 1..=2 {
8121 assert!(
8122 matches!(
8123 on_child_exit(
8124 &spec,
8125 runtime.restart_policy,
8126 &supervisor.registry,
8127 &snapshot,
8128 &runtime.terminal_ring,
8129 &runtime.spawn_events,
8130 &runtime.child_roster,
8131 crash_exit_report(attempt),
8132 )
8133 .await,
8134 NextAction::Restart { schedule: _ }
8135 ),
8136 "crash {attempt} is inside the budget and must respawn"
8137 );
8138 }
8139
8140 assert!(matches!(
8141 on_child_exit(
8142 &spec,
8143 runtime.restart_policy,
8144 &supervisor.registry,
8145 &snapshot,
8146 &runtime.terminal_ring,
8147 &runtime.spawn_events,
8148 &runtime.child_roster,
8149 crash_exit_report(3),
8150 )
8151 .await,
8152 NextAction::Stop { .. }
8153 ));
8154
8155 {
8156 let state = lock_snapshot(&snapshot).unwrap();
8157 assert_eq!(state.state, ModuleState::Failed);
8158 assert_eq!(state.crash_restarts.len(), 2);
8159 assert_eq!(state.lifetime_restarts, 2);
8160 }
8161
8162 let history = runtime
8163 .terminal_ring
8164 .lock()
8165 .expect("terminal ring is not poisoned")
8166 .snapshot();
8167 let last = history
8168 .entries
8169 .last()
8170 .expect("the refused crash is retained");
8171 assert_eq!(last.disposition, TerminalDisposition::Failed);
8172 assert_eq!(
8173 last.disposition_detail.as_deref(),
8174 Some("crash budget exhausted: max_restarts=2 within window_secs=600")
8175 );
8176
8177 let captured = crate::router::test_log::captured_logs(&logs);
8178 assert!(
8179 captured.contains("crash budget exhausted: max_restarts=2 within window_secs=600"),
8180 "the stop must be logged with its window: {captured}"
8181 );
8182 }
8183
8184 #[tokio::test]
8192 async fn a_crash_older_than_the_window_frees_its_slot_for_a_later_crash() {
8193 let supervisor = Supervisor::new(
8194 Arc::new(Registry::default()),
8195 RestartPolicy::new(2, Duration::ZERO),
8196 );
8197 let runtime = supervisor.runtime_config();
8198 let spec = windowed_crash_spec("crash-across-windows");
8199 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8200
8201 for attempt in 1..=2 {
8202 assert!(matches!(
8203 on_child_exit(
8204 &spec,
8205 runtime.restart_policy,
8206 &supervisor.registry,
8207 &snapshot,
8208 &runtime.terminal_ring,
8209 &runtime.spawn_events,
8210 &runtime.child_roster,
8211 crash_exit_report(attempt),
8212 )
8213 .await,
8214 NextAction::Restart { schedule: _ }
8215 ));
8216 }
8217
8218 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
8221 age_oldest_crash_restart_out_of_window(state, runtime.restart_policy.window);
8222 })
8223 .unwrap();
8224
8225 assert!(
8226 matches!(
8227 on_child_exit(
8228 &spec,
8229 runtime.restart_policy,
8230 &supervisor.registry,
8231 &snapshot,
8232 &runtime.terminal_ring,
8233 &runtime.spawn_events,
8234 &runtime.child_roster,
8235 crash_exit_report(3),
8236 )
8237 .await,
8238 NextAction::Restart { schedule: _ }
8239 ),
8240 "a crash older than the window must not hold a budget slot"
8241 );
8242
8243 let state = lock_snapshot(&snapshot).unwrap();
8244 assert_eq!(state.state, ModuleState::Restarting);
8245 assert_eq!(
8246 state.crash_restarts.len(),
8247 2,
8248 "the aged instant is dropped and the new one takes its place"
8249 );
8250 assert_eq!(
8251 state.lifetime_restarts, 3,
8252 "the ledger counts every restart, including the ones the window forgot"
8253 );
8254 }
8255
8256 #[tokio::test]
8261 async fn an_operator_restart_clears_the_ring_and_leaves_the_ledger_alone() {
8262 let supervisor = Supervisor::new(
8263 Arc::new(Registry::default()),
8264 RestartPolicy::new(2, Duration::ZERO),
8265 );
8266 let runtime = supervisor.runtime_config();
8267 let spec = windowed_crash_spec("operator-cleared-budget");
8268 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8269
8270 for attempt in 1..=2 {
8271 assert!(matches!(
8272 on_child_exit(
8273 &spec,
8274 runtime.restart_policy,
8275 &supervisor.registry,
8276 &snapshot,
8277 &runtime.terminal_ring,
8278 &runtime.spawn_events,
8279 &runtime.child_roster,
8280 crash_exit_report(attempt),
8281 )
8282 .await,
8283 NextAction::Restart { schedule: _ }
8284 ));
8285 }
8286
8287 reset_restart_count(&snapshot, &spec.module_id).unwrap();
8288 {
8289 let state = lock_snapshot(&snapshot).unwrap();
8290 assert!(
8291 state.crash_restarts.is_empty(),
8292 "an operator restart returns the full budget"
8293 );
8294 assert_eq!(
8295 state.lifetime_restarts, 2,
8296 "clearing the budget must not unmake the crashes"
8297 );
8298 }
8299
8300 assert!(
8301 matches!(
8302 on_child_exit(
8303 &spec,
8304 runtime.restart_policy,
8305 &supervisor.registry,
8306 &snapshot,
8307 &runtime.terminal_ring,
8308 &runtime.spawn_events,
8309 &runtime.child_roster,
8310 crash_exit_report(3),
8311 )
8312 .await,
8313 NextAction::Restart { schedule: _ }
8314 ),
8315 "the cleared budget must be spendable again"
8316 );
8317 let state = lock_snapshot(&snapshot).unwrap();
8318 assert_eq!(state.crash_restarts.len(), 1);
8319 assert_eq!(state.lifetime_restarts, 3);
8320 }
8321
8322 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
8323 async fn severance_marker_for_a_dead_child_does_not_label_its_successor() {
8324 let severed = ProcessIdentity {
8325 pid: 41,
8326 start_time: 101,
8327 };
8328 let successor = ProcessIdentity {
8329 pid: 41,
8330 start_time: 202,
8331 };
8332 let module = module_with_recovery_snapshot(ModuleState::Running, true, 0);
8333 update_snapshot(&module.inner.snapshot, Some("recovery-snapshot"), |state| {
8334 state.pid = Some(successor.pid);
8335 state.process_start_time = Some(successor.start_time);
8336 })
8337 .unwrap();
8338 assert!(!module.record_deliberate_severance(severed).unwrap());
8339
8340 let exit_report = apply_deliberate_severance_marker(
8341 &module.inner.snapshot,
8342 Some(successor),
8343 ExitReport {
8344 kind: ExitKind::Crash,
8345 code: Some(1),
8346 signal: None,
8347 at_ms: 1,
8348 },
8349 );
8350
8351 assert_eq!(exit_report.kind, ExitKind::Crash);
8352 }
8353
8354 #[tokio::test]
8355 async fn drain_reap_marks_deliberate_severance_and_records_lifetime_without_budget() {
8356 let registry = Registry::default();
8357 let supervisor = Supervisor::new(
8358 Arc::new(Registry::default()),
8359 RestartPolicy::new(3, Duration::ZERO),
8360 );
8361 let runtime = supervisor.runtime_config();
8362 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8363 let spec = ModuleSpec {
8364 module_id: "drain-deliberate-severance".to_string(),
8365 program: fake_aft_stub_path(),
8366 args: Vec::new(),
8367 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
8368 reserved: false,
8369 reserved_prefixes: Vec::new(),
8370 protocol: ModuleProtocol::Subc,
8371 overlap: Default::default(),
8372 };
8373 let mut child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
8374 let process = ProcessIdentity {
8375 pid: 41,
8376 start_time: 101,
8377 };
8378 child.process_identity = Some(process);
8379 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
8380 state.pid = Some(process.pid);
8381 state.process_start_time = Some(process.start_time);
8382 })
8383 .unwrap();
8384 record_deliberate_severance(&snapshot, process).unwrap();
8385
8386 drain_child_to_state(
8387 &spec.module_id,
8388 spec.protocol,
8389 StopNotice::SentOverConnection,
8392 ®istry,
8393 &snapshot,
8394 &runtime.terminal_ring,
8395 &runtime.spawn_events,
8396 child,
8397 Duration::from_secs(1),
8398 ModuleState::Stopped,
8399 Some(false),
8400 )
8401 .await
8402 .unwrap();
8403
8404 let state = lock_snapshot(&snapshot).unwrap();
8405 assert_eq!(
8406 state.last_exit.as_ref().map(|exit| exit.kind),
8407 Some(ExitKind::DeliberateSeverance)
8408 );
8409 assert_eq!(state.lifetime_restarts, 1);
8410 assert_eq!(state.crash_restarts.len(), 0);
8411 drop(state);
8412 let history = runtime.terminal_ring.lock().unwrap().snapshot();
8413 assert_eq!(
8414 history.entries[0].exit_kind,
8415 subc_control::TerminalExitKind::DeliberateSeverance
8416 );
8417 }
8418
8419 #[tokio::test]
8420 async fn ordinary_drain_reap_does_not_record_a_lifetime_restart() {
8421 let registry = Registry::default();
8422 let supervisor = Supervisor::new(
8423 Arc::new(Registry::default()),
8424 RestartPolicy::new(3, Duration::ZERO),
8425 );
8426 let runtime = supervisor.runtime_config();
8427 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8428 let spec = ModuleSpec {
8429 module_id: "ordinary-drain".to_string(),
8430 program: fake_aft_stub_path(),
8431 args: Vec::new(),
8432 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
8433 reserved: false,
8434 reserved_prefixes: Vec::new(),
8435 protocol: ModuleProtocol::Subc,
8436 overlap: Default::default(),
8437 };
8438 let child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
8439
8440 drain_child_to_state(
8441 &spec.module_id,
8442 spec.protocol,
8443 StopNotice::SentOverConnection,
8446 ®istry,
8447 &snapshot,
8448 &runtime.terminal_ring,
8449 &runtime.spawn_events,
8450 child,
8451 Duration::from_secs(1),
8452 ModuleState::Stopped,
8453 Some(false),
8454 )
8455 .await
8456 .unwrap();
8457
8458 let state = lock_snapshot(&snapshot).unwrap();
8459 assert_eq!(
8460 state.last_exit.as_ref().map(|exit| exit.kind),
8461 Some(ExitKind::Crash)
8462 );
8463 assert_eq!(state.lifetime_restarts, 0);
8464 assert_eq!(state.crash_restarts.len(), 0);
8465 }
8466
8467 #[test]
8468 fn fatal_connection_teardown_cannot_arm_a_marker_for_a_surviving_process() {
8469 assert!(!include_str!("server.rs")
8475 .contains("router.record_deliberate_connection_severance(ctx.connection_id)"));
8476 }
8477
8478 #[test]
8485 fn drained_after_quiescence_wait_passes_ok_through_and_forces_false_on_err() {
8486 assert!(drained_after_quiescence_wait(&Ok(true)));
8487 assert!(!drained_after_quiescence_wait(&Ok(false)));
8488 assert!(!drained_after_quiescence_wait(&Err(
8489 SuperviseError::StatePoisoned { module_id: None }
8490 )));
8491 }
8492
8493 #[test]
8502 fn wait_error_exit_report_records_a_failed_terminal_with_no_code_or_signal() {
8503 let ring = Arc::new(Mutex::new(TerminalRing::new(
8504 TerminalRingConfig::default(),
8505 0,
8506 )));
8507 record_wait_error_terminal("wait-error", &ring, &super::SpawnEventFeed::default());
8508
8509 let snapshot = ring.lock().unwrap().snapshot();
8510 assert_eq!(snapshot.entries.len(), 1);
8511 let entry = &snapshot.entries[0];
8512 assert_eq!(entry.exit_code, None);
8513 assert_eq!(entry.exit_signal, None);
8514 assert_eq!(entry.disposition, TerminalDisposition::Failed);
8515 }
8516
8517 #[test]
8518 fn wait_error_exit_path_preserves_spawn_event_density() {
8519 let feed = super::SpawnEventFeed::default();
8520 feed.configure_incarnation("wait-error-density".to_string());
8521 feed.emit_spawned("wait-error", 41, 1);
8522 let ring = Arc::new(Mutex::new(TerminalRing::new(
8523 TerminalRingConfig::default(),
8524 0,
8525 )));
8526
8527 record_wait_error_terminal("wait-error", &ring, &feed);
8528 feed.emit_spawned("after-wait-error", 42, 2);
8529
8530 let state = feed.0.lock().unwrap();
8531 let sequences = state
8532 .events
8533 .iter()
8534 .map(|event| event.cursor.seq)
8535 .collect::<Vec<_>>();
8536 assert_eq!(sequences, vec![1, 2, 3]);
8537 assert_eq!(state.events[1].kind, SpawnEventKind::Exited);
8538 assert_eq!(state.events[1].exit_code, None);
8539 assert_eq!(state.events[1].exit_signal, None);
8540 }
8541
8542 #[test]
8546 fn wait_error_exit_report_is_classified_as_a_crash() {
8547 assert_eq!(wait_error_exit_report().kind, ExitKind::Crash);
8548 }
8549}
8550
8551#[cfg(test)]
8552mod health_evidence_tests {
8553 use super::{HealthProbeError, HealthProbeEvidence};
8554 use std::collections::HashSet;
8555
8556 #[test]
8564 fn only_a_dead_lane_is_proof_of_death() {
8565 assert!(HealthProbeError::lane_dead("gone").is_proof_of_death());
8566 assert!(!HealthProbeError::no_answer("timed out").is_proof_of_death());
8570 assert!(!HealthProbeError::bad_answer("garbage").is_proof_of_death());
8571 assert!(!HealthProbeError::misconfigured("no table").is_proof_of_death());
8572 }
8573
8574 #[test]
8580 fn every_evidence_class_has_a_distinct_label() {
8581 let labels = [
8582 HealthProbeError::lane_dead("").label(),
8583 HealthProbeError::no_answer("").label(),
8584 HealthProbeError::bad_answer("").label(),
8585 HealthProbeError::misconfigured("").label(),
8586 ];
8587 let unique: HashSet<_> = labels.iter().collect();
8588 assert_eq!(unique.len(), labels.len(), "labels collided: {labels:?}");
8589 }
8590
8591 #[test]
8597 fn classification_preserves_the_original_message() {
8598 let err = HealthProbeError::no_answer("module did not answer within 5s");
8599 assert_eq!(err.to_string(), "module did not answer within 5s");
8600 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8601 }
8602}
8603
8604#[cfg(test)]
8605mod health_tombstone_tests {
8606 use std::{path::PathBuf, sync::Arc, time::Duration};
8607
8608 use subc_protocol::{
8609 manifest::Concurrency,
8610 session::{HealthStatus, ModuleControlResponse},
8611 };
8612 use tokio::sync::mpsc;
8613
8614 use super::{
8615 probe_module_health, HealthAction, HealthConfig, HealthProbeEvidence, ModuleProtocol,
8616 ModuleSpec, RestartPolicy, Supervisor, SupervisorRuntimeConfig,
8617 };
8618 use crate::{
8619 control::ControlHandler,
8620 forwarding::{ForwardingTable, ModuleControlRpcCompletion, ModuleControlRpcOutcome},
8621 registry::{ConnectionId, Registry},
8622 router::FrameSink,
8623 };
8624
8625 struct ProbeHarness {
8626 spec: ModuleSpec,
8627 runtime: SupervisorRuntimeConfig,
8628 forwarding: Arc<ForwardingTable>,
8629 module_connection: ConnectionId,
8630 module_rx: mpsc::Receiver<crate::router::OutboundFrame>,
8631 handler: ControlHandler,
8632 module: super::SupervisedModule,
8633 }
8634
8635 fn probe_harness() -> ProbeHarness {
8636 let registry = Arc::new(Registry::default());
8637 let forwarding = Arc::new(ForwardingTable::default());
8638 let supervisor_handle = super::SupervisorHandle::new();
8639 let health = HealthConfig {
8640 cadence: Duration::from_secs(30),
8641 deadline: Duration::from_secs(5),
8642 failure_threshold: 3,
8643 on_degraded: HealthAction::Report,
8644 on_failing: HealthAction::Report,
8645 critical: false,
8646 };
8647 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
8648 .with_forwarding(Arc::clone(&forwarding))
8649 .with_handle(supervisor_handle.clone())
8650 .with_health_config(health);
8651 let spec = ModuleSpec {
8652 module_id: "late-health-module".to_string(),
8653 program: PathBuf::from("disabled-module"),
8654 args: Vec::new(),
8655 env: Vec::new(),
8656 reserved: false,
8657 reserved_prefixes: Vec::new(),
8658 protocol: ModuleProtocol::Subc,
8659 overlap: Default::default(),
8660 };
8661 let module = supervisor
8662 .supervise_configured(spec.clone(), false)
8663 .unwrap();
8664 let runtime = supervisor.runtime_config();
8665 let handler = ControlHandler::with_forwarding(registry, Arc::clone(&forwarding))
8666 .with_supervisor(supervisor_handle);
8667 let module_connection = ConnectionId::new(700);
8668 let (module_tx, module_rx) = mpsc::channel(8);
8669 forwarding
8670 .register_module_connection(
8671 module_connection,
8672 spec.module_id.clone(),
8673 subc_protocol::PROTOCOL_VERSION,
8674 Concurrency::ModuleManaged,
8675 FrameSink::new(module_tx),
8676 )
8677 .unwrap();
8678
8679 ProbeHarness {
8680 spec,
8681 runtime,
8682 forwarding,
8683 module_connection,
8684 module_rx,
8685 handler,
8686 module,
8687 }
8688 }
8689
8690 async fn finish_after(
8691 harness: &mut ProbeHarness,
8692 stall: Duration,
8693 ) -> ModuleControlRpcCompletion {
8694 assert!(stall > harness.runtime.health.deadline);
8695 let deadline = harness.runtime.health.deadline;
8696 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8697 let answer = async {
8698 let frame = harness.module_rx.recv().await.expect("health.check frame");
8699 tokio::time::advance(deadline).await;
8700 tokio::task::yield_now().await;
8701 tokio::time::advance(stall - deadline).await;
8702 harness
8703 .forwarding
8704 .complete_module_control_rpc(
8705 harness.module_connection,
8706 frame.header.corr,
8707 Some("health.check"),
8708 ModuleControlRpcOutcome::Response(ModuleControlResponse::HealthCheck {
8709 status: HealthStatus::Ok,
8710 detail: None,
8711 metrics: None,
8712 }),
8713 )
8714 .unwrap()
8715 };
8716 let (probe_result, completion) = tokio::join!(probe, answer);
8717 let err = probe_result.expect_err("probe must miss its deadline");
8718 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8719 completion
8720 }
8721
8722 async fn time_out_without_answer(harness: &mut ProbeHarness) {
8723 let deadline = harness.runtime.health.deadline;
8724 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8725 let exhaust_deadline = async {
8726 let _frame = harness.module_rx.recv().await.expect("health.check frame");
8727 tokio::time::advance(deadline).await;
8728 tokio::task::yield_now().await;
8729 };
8730 let (probe_result, ()) = tokio::join!(probe, exhaust_deadline);
8731 let err = probe_result.expect_err("probe must miss its deadline");
8732 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8733 }
8734
8735 #[tokio::test(start_paused = true)]
8736 async fn late_health_answers_record_start_anchored_latency_for_two_stalls() {
8737 let mut harness = probe_harness();
8738
8739 let first = finish_after(&mut harness, Duration::from_secs(8)).await;
8740 let first_latency = match &first {
8741 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8742 other => panic!("late answer was not retained: {other:?}"),
8743 };
8744 assert!(harness.handler.observe_module_control_completion(first));
8745
8746 let second = finish_after(&mut harness, Duration::from_secs(11)).await;
8747 let second_latency = match &second {
8748 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8749 other => panic!("late answer was not retained: {other:?}"),
8750 };
8751 assert!(harness.handler.observe_module_control_completion(second));
8752
8753 assert_eq!(first_latency, Duration::from_secs(8));
8754 assert_eq!(
8755 second_latency - first_latency,
8756 Duration::from_secs(3),
8757 "latency must grow linearly with the additional stall"
8758 );
8759 let health = harness.module.status().unwrap().health;
8760 assert_eq!(health.late_answer_count, 2);
8761 assert_eq!(health.last_late_answer_latency_ms, Some(11_000));
8762 }
8763
8764 #[tokio::test(start_paused = true)]
8772 async fn late_answer_clears_the_consecutive_failure_streak() {
8773 let mut harness = probe_harness();
8774
8775 time_out_without_answer(&mut harness).await;
8777 harness
8778 .module
8779 .record_health_probe_failure_for_test("[no-answer] test miss")
8780 .unwrap();
8781 assert_eq!(
8782 harness.module.status().unwrap().health.consecutive_failures,
8783 1,
8784 "precondition: the miss must be on the streak before the late answer"
8785 );
8786
8787 let late = finish_after(&mut harness, Duration::from_secs(9)).await;
8789 assert!(matches!(
8790 late,
8791 ModuleControlRpcCompletion::LateHealthAnswer { .. }
8792 ));
8793 assert!(harness.handler.observe_module_control_completion(late));
8794
8795 let health = harness.module.status().unwrap().health;
8796 assert_eq!(
8797 health.consecutive_failures, 0,
8798 "a late answer is an answer: the streak must reset"
8799 );
8800 assert_eq!(health.late_answer_count, 1);
8801 }
8802
8803 #[tokio::test(start_paused = true)]
8804 async fn repeated_serial_probe_cycles_keep_one_tombstone_per_endpoint() {
8805 let mut harness = probe_harness();
8806
8807 for _ in 0..20 {
8808 time_out_without_answer(&mut harness).await;
8809 assert_eq!(
8810 harness.forwarding.health_probe_tombstone_count().unwrap(),
8811 1
8812 );
8813 }
8814 }
8815}
8816
8817#[cfg(test)]
8818mod child_env_tests {
8819 use super::{
8820 apply_child_env, apply_spawn_role, apply_wire_spawn_args, ModuleProtocol, ModuleSpec,
8821 SpawnRole, SupervisorHandle, SPAWN_ROLE_SWAP_CANDIDATE, SUBC_ARG, SUBC_LAUNCH_NONCE_ENV,
8822 SUBC_MODULE_ID_ENV, SUBC_SPAWN_ROLE_ENV,
8823 };
8824 use std::{ffi::OsStr, path::PathBuf};
8825 use tokio::process::Command;
8826
8827 fn spec(env: Vec<(String, String)>) -> ModuleSpec {
8828 ModuleSpec {
8829 module_id: "env-plan".to_string(),
8830 program: PathBuf::from("/nonexistent"),
8831 args: Vec::new(),
8832 env,
8833 reserved: false,
8834 reserved_prefixes: Vec::new(),
8835 protocol: ModuleProtocol::Subc,
8836 overlap: Default::default(),
8837 }
8838 }
8839
8840 #[test]
8854 fn ambient_ck_log_is_removed_and_a_configured_one_survives() {
8855 let mut command = Command::new("/nonexistent");
8856 apply_child_env(&mut command, &spec(Vec::new()));
8857 let removed = command
8858 .as_std()
8859 .get_envs()
8860 .any(|(key, value)| key == OsStr::new("CK_LOG") && value.is_none());
8861 assert!(
8862 removed,
8863 "ambient CK_LOG must be explicitly removed for an unconfigured module"
8864 );
8865
8866 let mut configured = Command::new("/nonexistent");
8867 apply_child_env(
8868 &mut configured,
8869 &spec(vec![("CK_LOG".to_string(), "debug".to_string())]),
8870 );
8871 let effective = configured
8872 .as_std()
8873 .get_envs()
8874 .filter(|(key, _)| *key == OsStr::new("CK_LOG"))
8875 .last()
8876 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()));
8877 assert_eq!(
8878 effective,
8879 Some(Some("debug".to_string())),
8880 "a module's configured CK_LOG must survive the ambient removal"
8881 );
8882 }
8883
8884 #[test]
8893 fn protocol_none_spawn_carries_no_subc_argument_and_no_nonce() {
8894 let connection_file = std::path::Path::new("/run/subc-connection.json");
8895 let handle = SupervisorHandle::new();
8896
8897 let mut none_spec = spec(Vec::new());
8898 none_spec.protocol = ModuleProtocol::None;
8899 let mut none = Command::new("/nonexistent");
8900 apply_wire_spawn_args(&mut none, &none_spec, Some(connection_file), Some(&handle))
8901 .expect("protocol-none spawn args apply");
8902 let none_args: Vec<String> = none
8903 .as_std()
8904 .get_args()
8905 .map(|a| a.to_string_lossy().into_owned())
8906 .collect();
8907 assert!(
8908 !none_args.iter().any(|a| a == SUBC_ARG),
8909 "protocol:none argv must not carry --subc; got {none_args:?}"
8910 );
8911 let none_has_nonce = none
8912 .as_std()
8913 .get_envs()
8914 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some());
8915 assert!(
8916 !none_has_nonce,
8917 "protocol:none spawn must not receive a launch nonce"
8918 );
8919 let none_has_module_id = none
8920 .as_std()
8921 .get_envs()
8922 .any(|(key, value)| key == OsStr::new(SUBC_MODULE_ID_ENV) && value.is_some());
8923 assert!(
8924 none_has_module_id,
8925 "SUBC_MODULE_ID is inert and stays on every path"
8926 );
8927 assert!(
8928 handle.spawn_nonce(&none_spec.module_id).is_none(),
8929 "no nonce record for a process that will never present one"
8930 );
8931
8932 let wire_spec = spec(Vec::new());
8934 let mut wire = Command::new("/nonexistent");
8935 apply_wire_spawn_args(&mut wire, &wire_spec, Some(connection_file), Some(&handle))
8936 .expect("subc-wire spawn args apply");
8937 let wire_args: Vec<String> = wire
8938 .as_std()
8939 .get_args()
8940 .map(|a| a.to_string_lossy().into_owned())
8941 .collect();
8942 assert_eq!(
8943 wire_args,
8944 vec![
8945 SUBC_ARG.to_string(),
8946 connection_file.to_string_lossy().into_owned()
8947 ],
8948 "a subc-wire spawn still carries --subc <path>"
8949 );
8950 assert!(wire
8951 .as_std()
8952 .get_envs()
8953 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some()));
8954 assert!(handle.spawn_nonce(&wire_spec.module_id).is_some());
8955 }
8956
8957 #[test]
8967 fn plain_spawn_removes_the_spawn_role_even_when_the_spec_sets_it() {
8968 let role = |command: &Command| {
8969 command
8970 .as_std()
8971 .get_envs()
8972 .filter(|(key, _)| *key == OsStr::new(SUBC_SPAWN_ROLE_ENV))
8973 .last()
8974 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()))
8975 };
8976 let forged = spec(vec![(
8977 SUBC_SPAWN_ROLE_ENV.to_string(),
8978 SPAWN_ROLE_SWAP_CANDIDATE.to_string(),
8979 )]);
8980
8981 let mut plain = Command::new("/nonexistent");
8982 apply_child_env(&mut plain, &forged);
8983 apply_spawn_role(&mut plain, SpawnRole::Plain);
8984 assert_eq!(
8985 role(&plain),
8986 Some(None),
8987 "a plain spawn must remove SUBC_SPAWN_ROLE, whatever the spec says"
8988 );
8989
8990 let mut candidate = Command::new("/nonexistent");
8991 apply_child_env(&mut candidate, &spec(Vec::new()));
8992 apply_spawn_role(&mut candidate, SpawnRole::SwapCandidate);
8993 assert_eq!(
8994 role(&candidate),
8995 Some(Some(SPAWN_ROLE_SWAP_CANDIDATE.to_string()))
8996 );
8997 }
8998
8999 #[test]
9005 fn daemon_private_capture_keys_are_not_passed_to_the_child() {
9006 let mut command = Command::new("/nonexistent");
9007 apply_child_env(
9008 &mut command,
9009 &spec(vec![
9010 (super::CAPTURE_KEEP_ENV.to_string(), "5".to_string()),
9011 ("KEPT".to_string(), "yes".to_string()),
9012 ]),
9013 );
9014 let keys: Vec<String> = command
9015 .as_std()
9016 .get_envs()
9017 .filter(|(_, value)| value.is_some())
9018 .map(|(key, _)| key.to_string_lossy().into_owned())
9019 .collect();
9020 assert!(keys.contains(&"KEPT".to_string()), "got {keys:?}");
9021 assert!(
9022 !keys.contains(&super::CAPTURE_KEEP_ENV.to_string()),
9023 "daemon-private capture key leaked to the child: {keys:?}"
9024 );
9025 }
9026}
9027
9028#[cfg(test)]
9029mod jitter_tests {
9030 use super::jittered_health_delay;
9031 use std::{collections::HashSet, time::Duration};
9032
9033 const FLEET: [&str; 14] = [
9042 "aft",
9043 "alfonso-core",
9044 "magic-context",
9045 "broca",
9046 "thalamus",
9047 "quota",
9048 "engram",
9049 "plexus",
9050 "cerebellum",
9051 "astrocyte",
9052 "synapse",
9053 "subc-mcp",
9054 "cortexkit-credentials",
9055 "subc-federation",
9056 ];
9057
9058 #[test]
9066 fn probe_delays_disperse_across_the_fleet() {
9067 let cadence = Duration::from_secs(30);
9068 let delays: HashSet<Duration> = FLEET
9069 .iter()
9070 .map(|id| jittered_health_delay(id, 0, cadence))
9071 .collect();
9072 assert_eq!(
9073 delays.len(),
9074 FLEET.len(),
9075 "every supervised module must land on its own probe offset"
9076 );
9077 }
9078
9079 #[test]
9085 fn jitter_only_delays_and_stays_within_one_tenth_of_cadence() {
9086 let cadence = Duration::from_secs(30);
9087 let span = cadence / 10;
9088 for id in FLEET {
9089 for probe_index in 0..8 {
9090 let delay = jittered_health_delay(id, probe_index, cadence);
9091 assert!(
9092 delay >= cadence,
9093 "{id}#{probe_index}: jitter must not shorten the cadence"
9094 );
9095 assert!(
9096 delay < cadence + span,
9097 "{id}#{probe_index}: jitter must stay inside one tenth of the cadence"
9098 );
9099 }
9100 }
9101 }
9102
9103 #[test]
9109 fn a_module_offset_is_stable_across_restarts() {
9110 let cadence = Duration::from_secs(30);
9111 for id in FLEET {
9112 assert_eq!(
9113 jittered_health_delay(id, 0, cadence),
9114 jittered_health_delay(id, 0, cadence),
9115 "{id}: the same module and probe index must produce the same offset"
9116 );
9117 }
9118 }
9119
9120 #[test]
9122 fn zero_cadence_yields_zero_delay() {
9123 assert_eq!(
9124 jittered_health_delay("aft", 0, Duration::ZERO),
9125 Duration::ZERO
9126 );
9127 }
9128}
9129
9130#[cfg(all(test, target_os = "linux"))]
9131mod cgroup_placement_tests {
9132 use super::{
9133 apply_cgroup_placement, remove_module_cgroup, ModuleProtocol, ModuleSpec, SuperviseError,
9134 SupervisedChild,
9135 };
9136 use crate::{
9137 stderr_tail::{StderrRing, StderrTailConfig},
9138 test_support::TestTempDir,
9139 };
9140 use std::{
9141 fs, io,
9142 path::{Path, PathBuf},
9143 sync::{Arc, Mutex},
9144 };
9145 use tokio::process::Command;
9146
9147 #[test]
9148 fn failed_parent_cgroup_open_is_a_cgroup_supervision_error() {
9149 let path = Path::new("/definitely-missing-subc-cgroup");
9150 let mut command = Command::new("true");
9151 let error = apply_cgroup_placement(
9152 &mut command,
9153 &ModuleSpec {
9154 module_id: "broken-cgroup".to_string(),
9155 program: PathBuf::from("true"),
9156 args: Vec::new(),
9157 env: Vec::new(),
9158 reserved: false,
9159 reserved_prefixes: Vec::new(),
9160 protocol: ModuleProtocol::Subc,
9161 overlap: Default::default(),
9162 },
9163 path,
9164 )
9165 .expect_err("a parent cgroup open failure must reject the supervised spawn");
9166 let reason = error.to_string();
9167
9168 assert!(
9169 matches!(error, SuperviseError::Cgroup { .. }),
9170 "parent cgroup open must be reported as a cgroup supervision error: {reason}"
9171 );
9172 assert!(
9173 reason.contains("/definitely-missing-subc-cgroup/cgroup.procs"),
9174 "parent cgroup open failure must name cgroup.procs: {reason}"
9175 );
9176 }
9177
9178 #[tokio::test]
9179 async fn reaping_a_child_removes_its_empty_module_cgroup() {
9180 let root = TestTempDir::new("supervisor-reap-cgroup");
9181 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
9182 let placement = subc_cgroup::prepare_at(&root)
9183 .expect("prepare scratch cgroup root")
9184 .expect("scratch root has a cgroup.procs marker");
9185 let module_id = "reaped-module";
9186 let module = placement
9187 .module_path(module_id)
9188 .expect("create scratch module cgroup");
9189 let child = Command::new("true")
9190 .spawn()
9191 .expect("spawn short-lived child");
9192 let pid = child.id().expect("spawned child has pid");
9193 let mut child = SupervisedChild {
9194 child,
9195 module_id: module_id.to_string(),
9196 cgroup_placement: Some(placement),
9197 stdout_pump: None,
9198 stderr_pump: None,
9199 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
9200 spawned_at_ms: 0,
9201 spawned_from: PathBuf::from("true"),
9202 spawned_file_identity: None,
9203 process_start_time: None,
9204 process_identity: None,
9205 pid,
9206 roster_guard: None,
9207 };
9208
9209 child.wait().await.expect("reap short-lived child");
9210
9211 assert!(
9212 !module.exists(),
9213 "reaping the supervised child must remove its empty cgroup"
9214 );
9215 }
9216
9217 #[test]
9218 fn non_empty_cgroup_removal_is_reported_without_blocking_teardown() {
9219 let root = TestTempDir::new("supervisor-non-empty-cgroup");
9220 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
9221 let placement = subc_cgroup::prepare_at(&root)
9222 .expect("prepare scratch cgroup root")
9223 .expect("scratch root has a cgroup.procs marker");
9224 let module = placement
9225 .module_path("surviving-module")
9226 .expect("create scratch module cgroup");
9227 fs::write(module.join("surviving-process"), b"still present")
9228 .expect("make scratch cgroup non-empty");
9229 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::WARN);
9230
9231 remove_module_cgroup(&placement, "surviving-module");
9232
9233 let logs = crate::router::test_log::captured_logs(&logs);
9234 assert!(
9235 module.exists(),
9236 "failed removal must leave the cgroup intact"
9237 );
9238 assert!(
9239 logs.contains("could not remove module cgroup after process exit; continuing teardown")
9240 && logs.contains("surviving-module"),
9241 "best-effort removal must report the failure without returning it: {logs}"
9242 );
9243 }
9244
9245 #[test]
9246 fn cgroup_pre_exec_spawn_failure_names_the_cgroup_path() {
9247 let cgroup_path = PathBuf::from("/sys/fs/cgroup/subc-modules/broken-module");
9248 let reason = SuperviseError::Spawn {
9249 program: PathBuf::from("/bin/true"),
9250 source: io::Error::from_raw_os_error(13),
9251 cgroup_path: Some(cgroup_path.clone()),
9252 }
9253 .to_string();
9254
9255 assert!(
9256 reason.contains(&cgroup_path.display().to_string()),
9257 "a pre_exec spawn failure must name the cgroup path: {reason}"
9258 );
9259 }
9260}
9261
9262#[cfg(test)]
9263mod spawn_subscriber_lag_tests {
9264 use super::*;
9265
9266 #[tokio::test]
9271 async fn lagged_spawn_subscriber_receives_a_terminal_lagged_error_after_its_queued_frames() {
9272 let feed = SpawnEventFeed::default();
9273 feed.configure_incarnation("lag-incarnation".to_string());
9274 let (tx, mut rx) = mpsc::channel(1);
9277 feed.subscribe(ConnectionId::new(1), 7, 1, None, FrameSink::new(tx))
9278 .expect("subscribe");
9279 let emitted = SPAWN_SUBSCRIBER_BUFFER + 16;
9280 for index in 0..emitted {
9281 feed.emit_spawned(&format!("lag-module-{index}"), 1000, 0);
9282 tokio::task::yield_now().await;
9285 }
9286 assert_eq!(
9287 feed.subscriber_count(),
9288 0,
9289 "the lagged subscriber must be removed"
9290 );
9291
9292 let mut data = Vec::new();
9293 let mut last = None;
9294 loop {
9295 let next = tokio::time::timeout(Duration::from_secs(5), rx.recv())
9296 .await
9297 .expect("the forwarder must finish once the subscriber is dropped");
9298 let Some(outbound) = next else { break };
9299 let frame = outbound.frame;
9300 if frame.header.ty == FrameType::StreamData {
9301 assert!(last.is_none(), "no data may follow the terminal frame");
9302 let event: SpawnEvent = serde_json::from_slice(&frame.body).unwrap();
9303 data.push(event.cursor.seq);
9304 } else {
9305 assert!(last.is_none(), "exactly one terminal frame");
9306 last = Some(frame);
9307 }
9308 }
9309 assert!(!data.is_empty(), "queued frames drain before the terminal");
9310 for pair in data.windows(2) {
9311 assert_eq!(
9312 pair[1],
9313 pair[0] + 1,
9314 "queued frames arrive dense and in order"
9315 );
9316 }
9317 let terminal = last.expect("a lagged subscriber must receive a terminal frame");
9318 assert_eq!(terminal.header.ty, FrameType::Error);
9319 assert_eq!(terminal.header.corr, 7);
9320 let body: subc_protocol::ErrorBody = serde_json::from_slice(&terminal.body).unwrap();
9321 assert_eq!(body.code, SPAWN_SUBSCRIBER_LAGGED_CODE);
9322 let detail = body.detail.expect("lagged error carries detail");
9323 assert_eq!(
9324 detail["first_undelivered_cursor"]["seq"],
9325 data.last().unwrap() + 1,
9326 "the named cursor is the first event the subscriber did not receive"
9327 );
9328 assert_eq!(
9329 detail["first_undelivered_cursor"]["daemon_incarnation"],
9330 "lag-incarnation"
9331 );
9332 }
9333}
9334
9335#[cfg(test)]
9336mod terminal_history_read_concurrency_tests {
9337 use super::*;
9338 use crate::{terminal_journal::read_pause, test_support::TestTempDir};
9339 use std::sync::mpsc as std_mpsc;
9340
9341 fn journaled_ring(
9342 journal: &Arc<crate::terminal_journal::TerminalJournal>,
9343 ) -> Arc<Mutex<TerminalRing>> {
9344 Arc::new(Mutex::new(
9345 TerminalRing::new(TerminalRingConfig::default(), 1)
9346 .with_journal(Some(Arc::clone(journal))),
9347 ))
9348 }
9349
9350 fn crash(at_ms: u64) -> ExitReport {
9351 ExitReport {
9352 kind: ExitKind::Crash,
9353 code: Some(1),
9354 signal: None,
9355 at_ms,
9356 }
9357 }
9358
9359 fn record_within(
9362 module_id: &'static str,
9363 ring: &Arc<Mutex<TerminalRing>>,
9364 at_ms: u64,
9365 bound: Duration,
9366 ) -> bool {
9367 let ring = Arc::clone(ring);
9368 let (done, done_rx) = std_mpsc::channel();
9369 std::thread::spawn(move || {
9370 record_terminal(
9371 module_id,
9372 &ring,
9373 &SpawnEventFeed::default(),
9374 &crash(at_ms),
9375 TerminalDisposition::Restarting,
9376 );
9377 let _ = done.send(());
9378 });
9379 done_rx.recv_timeout(bound).is_ok()
9380 }
9381
9382 #[test]
9387 fn exits_recorded_during_a_paused_history_read_are_not_blocked_or_half_merged() {
9388 let dir = TestTempDir::new("terminal-history-concurrent-read");
9389 let path = dir.join("terminals.jsonl");
9390 let journal = Arc::new(crate::terminal_journal::TerminalJournal::open(
9391 path.clone(),
9392 "daemon".into(),
9393 ));
9394 let reader_ring = journaled_ring(&journal);
9395 let other_ring = journaled_ring(&journal);
9396 assert!(record_within(
9397 "reader-module",
9398 &reader_ring,
9399 10,
9400 Duration::from_secs(5)
9401 ));
9402
9403 let (started, release) = read_pause::install(&path);
9404 let reading = {
9405 let ring = Arc::clone(&reader_ring);
9406 std::thread::spawn(move || durable_terminal_history_of(&ring, "reader-module"))
9407 };
9408 started
9409 .recv_timeout(Duration::from_secs(5))
9410 .expect("the history read reached its pause");
9411
9412 let bound = Duration::from_secs(1);
9413 assert!(
9414 record_within("other-module", &other_ring, 20, bound),
9415 "another module's exit waited on a history read (journal writer held)"
9416 );
9417 assert!(
9418 record_within("reader-module", &reader_ring, 30, bound),
9419 "the read module's own exit waited on its history read (ring held)"
9420 );
9421
9422 drop(release);
9423 let paused = reading.join().unwrap();
9424 assert_eq!(
9425 paused.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9426 vec![10],
9427 "an exit recorded after the read began lands in neither half of it"
9428 );
9429 assert_eq!(paused.journal_skipped_lines, 0);
9430 assert_eq!(paused.journal_read_errors, 0);
9431
9432 let after = durable_terminal_history_of(&reader_ring, "reader-module");
9433 assert_eq!(
9434 after.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9435 vec![10, 30],
9436 "the next read merges ring and journal with no duplicate"
9437 );
9438 assert_eq!(after.journal_skipped_lines, 0);
9439 }
9440}
9441
9442#[cfg(test)]
9447mod stderr_settle_tests {
9448 use std::{
9449 future::Future,
9450 io,
9451 pin::Pin,
9452 sync::{Arc, Mutex},
9453 task::{Context, Poll},
9454 time::Duration,
9455 };
9456
9457 use tokio::{
9458 io::{AsyncRead, ReadBuf},
9459 sync::oneshot,
9460 time::Instant,
9461 };
9462
9463 use super::{settle_stderr_pump, StderrPump};
9464 use crate::stderr_tail::{
9465 pump_stderr_to, CaptureState, OutputSink, StderrRing, StderrTailConfig, TailEntry,
9466 };
9467
9468 const BOUND: Duration = Duration::from_millis(250);
9469
9470 struct HeldReader {
9474 before: Option<Vec<u8>>,
9475 gate: Option<oneshot::Receiver<()>>,
9476 after: io::Cursor<Vec<u8>>,
9477 }
9478
9479 impl AsyncRead for HeldReader {
9480 fn poll_read(
9481 mut self: Pin<&mut Self>,
9482 cx: &mut Context<'_>,
9483 buf: &mut ReadBuf<'_>,
9484 ) -> Poll<io::Result<()>> {
9485 if let Some(bytes) = self.before.take() {
9486 buf.put_slice(&bytes);
9487 return Poll::Ready(Ok(()));
9488 }
9489 if let Some(gate) = self.gate.as_mut() {
9490 match Pin::new(gate).poll(cx) {
9491 Poll::Pending => return Poll::Pending,
9492 Poll::Ready(_) => self.gate = None,
9493 }
9494 }
9495 Pin::new(&mut self.after).poll_read(cx, buf)
9496 }
9497 }
9498
9499 struct DiscardSink;
9500
9501 impl OutputSink for DiscardSink {
9502 fn write_line(&mut self, _line: &[u8]) {}
9503 }
9504
9505 fn line(text: &str) -> TailEntry {
9506 TailEntry::Line {
9507 text: text.to_string(),
9508 truncated: false,
9509 }
9510 }
9511
9512 fn lock(ring: &Arc<Mutex<StderrRing>>) -> std::sync::MutexGuard<'_, StderrRing> {
9513 ring.lock().unwrap()
9514 }
9515
9516 fn held_pump(
9520 ring: &Arc<Mutex<StderrRing>>,
9521 before: &str,
9522 after: &str,
9523 ) -> (StderrPump, oneshot::Sender<()>) {
9524 let generation = lock(ring).begin_process();
9525 let (release, gate) = oneshot::channel();
9526 let reader = HeldReader {
9527 before: Some(before.as_bytes().to_vec()),
9528 gate: Some(gate),
9529 after: io::Cursor::new(after.as_bytes().to_vec()),
9530 };
9531 let task = tokio::spawn(pump_stderr_to(
9532 reader,
9533 Arc::clone(ring),
9534 generation,
9535 DiscardSink,
9536 ));
9537 (StderrPump { task, generation }, release)
9538 }
9539
9540 async fn wait_until(ring: &Arc<Mutex<StderrRing>>, done: impl Fn(&StderrRing) -> bool) {
9541 for _ in 0..1000 {
9542 if done(&lock(ring)) {
9543 return;
9544 }
9545 tokio::time::sleep(Duration::from_millis(1)).await;
9546 }
9547 panic!(
9548 "ring never reached the expected state: {:?}",
9549 lock(ring).snapshot(None, None)
9550 );
9551 }
9552
9553 #[tokio::test(start_paused = true)]
9554 async fn a_crash_line_the_reader_had_not_reached_by_the_bound_is_kept_before_the_restart() {
9555 let ring = Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default())));
9556 let (pump, release) = held_pump(&ring, "booting\n", "config error: missing storage\n");
9557
9558 settle_stderr_pump("crasher", &ring, pump, BOUND).await;
9559 let before_release = lock(&ring).snapshot(None, None);
9560 assert!(
9561 matches!(before_release.capture, CaptureState::Incomplete { .. }),
9562 "a reader that has not reached EOF cannot claim a whole tail: {before_release:?}"
9563 );
9564
9565 let next = lock(&ring).begin_process();
9568 lock(&ring).push_line_from(next, "next process booting");
9569 release.send(()).unwrap();
9570 wait_until(&ring, |ring| {
9571 ring.snapshot(None, None).capture == CaptureState::Captured
9572 })
9573 .await;
9574
9575 assert_eq!(
9576 lock(&ring).snapshot(None, None).entries,
9577 vec![
9578 line("booting"),
9579 line("config error: missing storage"),
9580 TailEntry::ProcessStart,
9581 line("next process booting"),
9582 ],
9583 "the crash's last line must survive a slow reader and stay in the crashed process's section"
9584 );
9585 }
9586
9587 #[tokio::test(start_paused = true)]
9588 async fn a_pipe_held_open_by_a_descendant_reads_incomplete_without_delaying_the_restart_past_the_bound(
9589 ) {
9590 let ring = Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default())));
9591 let (pump, _held) = held_pump(&ring, "parent exiting\n", "");
9594
9595 let started = Instant::now();
9596 settle_stderr_pump("orphaning", &ring, pump, BOUND).await;
9597 assert_eq!(
9598 started.elapsed(),
9599 BOUND,
9600 "the restart must wait exactly the bound for a pipe that stays open, no longer"
9601 );
9602
9603 let next = lock(&ring).begin_process();
9604 lock(&ring).push_line_from(next, "next process booting");
9605 tokio::time::sleep(Duration::from_secs(60)).await;
9606
9607 let snapshot = lock(&ring).snapshot(None, None);
9608 match &snapshot.capture {
9609 CaptureState::Incomplete { reason } => assert!(
9610 reason.contains("had not reached EOF") && reason.contains("250ms"),
9611 "the reason must say what is missing and after how long: {reason}"
9612 ),
9613 other => panic!("expected Incomplete while the pipe is held open, got {other:?}"),
9614 }
9615 assert_eq!(
9616 snapshot.entries,
9617 vec![
9618 line("parent exiting"),
9619 TailEntry::ProcessStart,
9620 line("next process booting"),
9621 ]
9622 );
9623 }
9624
9625 #[tokio::test(start_paused = true)]
9626 async fn a_reader_that_reaches_eof_within_the_bound_leaves_the_tail_captured() {
9627 let ring = Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default())));
9628 let (pump, release) = held_pump(&ring, "one\n", "two\n");
9629 release.send(()).unwrap();
9630
9631 settle_stderr_pump("clean", &ring, pump, BOUND).await;
9632
9633 let snapshot = lock(&ring).snapshot(None, None);
9634 assert_eq!(snapshot.capture, CaptureState::Captured);
9635 assert_eq!(snapshot.entries, vec![line("one"), line("two")]);
9636 }
9637}