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;
146 if result.is_ok() {
147 self.roster_guard = None;
150 }
151 #[cfg(target_os = "linux")]
152 if result.is_ok() {
153 if let Some(placement) = self.cgroup_placement.take() {
154 remove_module_cgroup(&placement, &self.module_id);
155 }
156 }
157 result
158 }
159
160 fn start_kill(&mut self) -> io::Result<()> {
161 self.child.start_kill()
162 }
163
164 async fn drain_stderr(&mut self, module_id: &str) {
165 if let Some(mut pump) = self.stdout_pump.take() {
166 match timeout(STDERR_PUMP_DRAIN_TIMEOUT, &mut pump).await {
167 Ok(Ok(())) => {}
168 Ok(Err(error)) => {
169 warn!(module_id, error = %error, "stdout pump ended unexpectedly");
170 }
171 Err(_) => {
172 pump.abort();
173 warn!(
174 module_id,
175 waited = ?STDERR_PUMP_DRAIN_TIMEOUT,
176 "stdout pump did not drain before restart; stopped it before the next process"
177 );
178 }
179 }
180 }
181
182 let Some(pump) = self.stderr_pump.take() else {
183 return;
184 };
185 settle_stderr_pump(
186 module_id,
187 &self.stderr_ring,
188 pump,
189 STDERR_PUMP_DRAIN_TIMEOUT,
190 )
191 .await;
192 }
193}
194
195struct StderrPump {
198 task: JoinHandle<()>,
199 generation: u64,
200}
201
202async fn settle_stderr_pump(
208 module_id: &str,
209 ring: &Arc<Mutex<StderrRing>>,
210 pump: StderrPump,
211 bound: Duration,
212) {
213 let lock = || ring.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
214 let StderrPump {
215 mut task,
216 generation,
217 } = pump;
218 lock().retire_pump(generation);
219 match timeout(bound, &mut task).await {
220 Ok(Ok(())) => {}
221 Ok(Err(err)) => {
222 let mut ring = lock();
223 ring.mark_incomplete(format!("stderr pump ended unexpectedly: {err}"));
224 ring.finish_pump(generation);
225 warn!(module_id, error = %err, "stderr pump ended before clean EOF");
226 }
227 Err(_) => {
228 drop(task);
230 lock().mark_pump_late(
231 generation,
232 format!(
233 "stderr of the exited process had not reached EOF {bound:?} after it was \
234 retired (a descendant may still hold the pipe open); lines it still \
235 writes are kept in that process's section"
236 ),
237 );
238 warn!(
239 module_id,
240 waited = ?bound,
241 "stderr pipe of the exited process is still open; its reader keeps running without delaying the restart"
242 );
243 }
244 }
245}
246
247fn registration_release_events() -> &'static watch::Sender<u64> {
248 static EVENTS: OnceLock<watch::Sender<u64>> = OnceLock::new();
249 EVENTS.get_or_init(|| {
250 let (sender, _receiver) = watch::channel(0);
251 sender
252 })
253}
254
255pub(crate) fn notify_registration_release() {
256 let events = registration_release_events();
257 let next_generation = (*events.borrow()).wrapping_add(1);
258 events.send_replace(next_generation);
259}
260
261#[derive(Debug, Clone, PartialEq, Eq)]
263pub struct ModuleSpec {
264 pub module_id: String,
265 pub program: PathBuf,
266 pub args: Vec<String>,
267 pub env: Vec<(String, String)>,
268 pub reserved: bool,
273 pub reserved_prefixes: Vec<String>,
278 pub protocol: ModuleProtocol,
294 pub overlap: ModuleOverlap,
299}
300
301#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
308pub enum ModuleOverlap {
309 #[default]
311 Exclusive,
312 Safe,
324}
325
326impl ModuleOverlap {
327 pub fn as_str(self) -> &'static str {
328 match self {
329 Self::Exclusive => "exclusive",
330 Self::Safe => "safe",
331 }
332 }
333}
334
335pub const SUBC_SPAWN_ROLE_ENV: &str = "SUBC_SPAWN_ROLE";
345pub const SPAWN_ROLE_SWAP_CANDIDATE: &str = "swap_candidate";
347pub const DEFAULT_SWAP_READY_TIMEOUT: Duration = Duration::from_secs(100);
352
353#[derive(Debug, Clone, Copy, PartialEq, Eq)]
371pub struct RestartPolicy {
372 pub max_restarts: u32,
373 pub backoff: Duration,
376 pub max_backoff: Duration,
378 pub window: Duration,
382}
383
384impl RestartPolicy {
385 pub fn new(max_restarts: u32, backoff: Duration) -> Self {
389 Self {
390 max_restarts,
391 backoff,
392 max_backoff: DEFAULT_MAX_BACKOFF,
393 window: DEFAULT_RESTART_WINDOW,
394 }
395 }
396
397 pub fn with_max_backoff(mut self, max_backoff: Duration) -> Self {
398 self.max_backoff = max_backoff;
399 self
400 }
401
402 pub fn with_window(mut self, window: Duration) -> Self {
403 self.window = window;
404 self
405 }
406
407 fn delay_for_restart(&self, restart_in_window: u32) -> Duration {
412 if self.backoff.is_zero() || self.max_backoff.is_zero() {
413 return Duration::ZERO;
414 }
415
416 let mut delay = self.backoff;
417 for _ in 0..restart_in_window {
418 if delay >= self.max_backoff {
419 return self.max_backoff;
420 }
421 delay = delay
422 .checked_mul(10)
423 .unwrap_or(self.max_backoff)
424 .min(self.max_backoff);
425 }
426 delay.min(self.max_backoff)
427 }
428
429 fn budget_exhausted_detail(&self) -> String {
434 format!(
435 "crash budget exhausted: max_restarts={} within window_secs={}",
436 self.max_restarts,
437 self.window.as_secs()
438 )
439 }
440}
441
442impl Default for RestartPolicy {
443 fn default() -> Self {
444 Self {
445 max_restarts: DEFAULT_MAX_RESTARTS,
446 backoff: DEFAULT_BACKOFF,
447 max_backoff: DEFAULT_MAX_BACKOFF,
448 window: DEFAULT_RESTART_WINDOW,
449 }
450 }
451}
452
453#[derive(Debug, Clone, Copy, PartialEq, Eq)]
454struct CrashRestartSchedule {
455 restart_in_window: u32,
456 delay: Duration,
457}
458
459fn daemon_will_restart(
466 state: &mut SupervisorSnapshot,
467 policy: &RestartPolicy,
468 now: Instant,
469) -> bool {
470 state.enabled && state.crash_restarts_in_window(policy.window, now) < policy.max_restarts
471}
472
473const DEFAULT_HEALTH_CADENCE: Duration = Duration::from_secs(30);
474const DEFAULT_HEALTH_DEADLINE: Duration = Duration::from_secs(5);
475const DEFAULT_HEALTH_FAILURE_THRESHOLD: u32 = 3;
476const MAX_HEALTH_METRICS_BYTES: usize = 16 * 1024;
477
478#[derive(Debug, Clone, Copy, PartialEq, Eq)]
479pub enum HealthAction {
480 Report,
481 Restart,
482 Alert,
483}
484
485impl fmt::Display for HealthAction {
486 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
487 f.write_str(match self {
488 Self::Report => "report",
489 Self::Restart => "restart",
490 Self::Alert => "alert",
491 })
492 }
493}
494
495#[derive(Debug, Clone, Copy, PartialEq, Eq)]
496pub struct HealthConfig {
497 pub cadence: Duration,
498 pub deadline: Duration,
499 pub failure_threshold: u32,
500 pub on_degraded: HealthAction,
501 pub on_failing: HealthAction,
502 pub critical: bool,
503}
504
505impl Default for HealthConfig {
506 fn default() -> Self {
507 Self {
508 cadence: DEFAULT_HEALTH_CADENCE,
509 deadline: DEFAULT_HEALTH_DEADLINE,
510 failure_threshold: DEFAULT_HEALTH_FAILURE_THRESHOLD,
511 on_degraded: HealthAction::Report,
512 on_failing: HealthAction::Report,
513 critical: false,
514 }
515 }
516}
517
518#[derive(Debug, Clone, PartialEq)]
536pub struct ModuleHealthStatus {
537 pub status: SupervisorHealthStatus,
538 pub last_probe_ms: Option<u64>,
539 pub detail: Option<String>,
540 pub metrics: Option<Value>,
541 pub consecutive_failures: u32,
542 pub late_answer_count: u64,
545 pub last_late_answer_latency_ms: Option<u64>,
547 pub last_action: Option<String>,
548 pub last_action_ms: Option<u64>,
552}
553
554impl Default for ModuleHealthStatus {
555 fn default() -> Self {
556 Self {
557 status: SupervisorHealthStatus::Unknown,
558 last_probe_ms: None,
559 detail: None,
560 metrics: None,
561 consecutive_failures: 0,
562 late_answer_count: 0,
563 last_late_answer_latency_ms: None,
564 last_action: None,
565 last_action_ms: None,
566 }
567 }
568}
569
570#[derive(Debug, Clone, Copy, PartialEq, Eq)]
572pub enum ModuleState {
573 Starting,
574 Running,
575 Unresponsive,
576 Restarting,
577 Draining,
578 Stopped,
579 Failed,
580 Disabled,
581}
582
583impl fmt::Display for ModuleState {
584 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
585 f.write_str(match self {
586 Self::Starting => "starting",
587 Self::Running => "running",
588 Self::Unresponsive => "unresponsive",
589 Self::Restarting => "restarting",
590 Self::Draining => "draining",
591 Self::Stopped => "stopped",
592 Self::Failed => "failed",
593 Self::Disabled => "disabled",
594 })
595 }
596}
597
598#[derive(Debug, Clone, Copy, PartialEq, Eq)]
600pub enum ExitKind {
601 Clean,
602 Crash,
603 DeliberateSeverance,
604}
605
606impl From<ExitKind> for TerminalExitKind {
607 fn from(kind: ExitKind) -> Self {
608 match kind {
609 ExitKind::Clean => Self::Clean,
610 ExitKind::Crash => Self::Crash,
611 ExitKind::DeliberateSeverance => Self::DeliberateSeverance,
612 }
613 }
614}
615
616#[derive(Debug, Clone, Copy, PartialEq, Eq)]
619pub(crate) struct ProcessIdentity {
620 pub(crate) pid: u32,
621 pub(crate) start_time: u64,
622}
623
624#[derive(Debug, Clone, PartialEq, Eq)]
626pub struct ExitReport {
627 pub kind: ExitKind,
628 pub code: Option<i32>,
629 pub signal: Option<i32>,
630 pub at_ms: u64,
631}
632
633#[derive(Debug, Clone, PartialEq)]
636pub struct ModuleStatus {
637 pub module_id: String,
638 pub state: ModuleState,
639 pub enabled: bool,
640 pub process_alive: bool,
641 pub registration_active: bool,
642 pub protocol: ModuleProtocol,
645 pub live: bool,
656 pub restart_count: u32,
660 pub lifetime_restarts: u32,
664 pub spawn_generation: u64,
665 pub max_restarts: u32,
670 pub restart_window: Duration,
674 pub drain_timeout: Duration,
678 pub restart_backoff: Duration,
679 pub restart_max_backoff: Duration,
680 pub pid: Option<u32>,
681 pub spawned_at_ms: Option<u64>,
682 pub spawned_from: Option<PathBuf>,
683 pub process_start_time: Option<u64>,
684 pub last_exit: Option<ExitReport>,
685 pub health: ModuleHealthStatus,
686}
687
688#[derive(Debug, Clone, PartialEq)]
689struct SupervisorSnapshot {
690 state: ModuleState,
691 enabled: bool,
692 process_alive: bool,
693 crash_restarts: VecDeque<Instant>,
699 lifetime_restarts: u32,
700 spawn_generation: u64,
709 pid: Option<u32>,
710 spawned_at_ms: Option<u64>,
711 spawned_from: Option<PathBuf>,
712 spawned_file_identity: Option<SpawnedFileIdentity>,
713 process_start_time: Option<u64>,
714 deliberate_severance: Option<ProcessIdentity>,
715 last_exit: Option<ExitReport>,
716 health: ModuleHealthStatus,
717 in_alternate_slot: bool,
722}
723
724impl SupervisorSnapshot {
725 fn starting() -> Self {
726 Self::new(ModuleState::Starting, true)
727 }
728
729 fn disabled() -> Self {
730 Self::new(ModuleState::Disabled, false)
731 }
732
733 fn failed() -> Self {
734 Self::new(ModuleState::Failed, true)
735 }
736
737 fn crash_restarts_in_window(&mut self, window: Duration, now: Instant) -> u32 {
741 while let Some(oldest) = self.crash_restarts.front() {
742 if now.duration_since(*oldest) > window {
743 self.crash_restarts.pop_front();
744 } else {
745 break;
746 }
747 }
748 u32::try_from(self.crash_restarts.len()).unwrap_or(u32::MAX)
749 }
750
751 fn record_crash_restart(&mut self, policy: &RestartPolicy, now: Instant) {
757 self.crash_restarts.push_back(now);
758 while self.crash_restarts.len() > policy.max_restarts as usize {
759 self.crash_restarts.pop_front();
760 }
761 self.lifetime_restarts += 1;
762 }
763
764 fn next_crash_restart(
768 &mut self,
769 policy: &RestartPolicy,
770 now: Instant,
771 ) -> Option<CrashRestartSchedule> {
772 let restart_in_window = self.crash_restarts_in_window(policy.window, now);
773 if restart_in_window >= policy.max_restarts {
774 return None;
775 }
776 self.record_crash_restart(policy, now);
777 Some(CrashRestartSchedule {
778 restart_in_window,
779 delay: policy.delay_for_restart(restart_in_window),
780 })
781 }
782
783 fn clear_crash_restarts(&mut self) {
788 self.crash_restarts.clear();
789 }
790
791 fn new(state: ModuleState, enabled: bool) -> Self {
792 Self {
793 state,
794 enabled,
795 process_alive: false,
796 crash_restarts: VecDeque::new(),
797 lifetime_restarts: 0,
798 spawn_generation: 0,
799 pid: None,
800 spawned_at_ms: None,
801 spawned_from: None,
802 spawned_file_identity: None,
803 process_start_time: None,
804 deliberate_severance: None,
805 last_exit: None,
806 health: ModuleHealthStatus::default(),
807 in_alternate_slot: false,
808 }
809 }
810}
811
812type SharedSnapshot = Arc<Mutex<SupervisorSnapshot>>;
813
814type SpawnSubscriberKey = (ConnectionId, u64);
815
816#[derive(Debug)]
817struct SpawnSubscriber {
818 version: u8,
819 frames: mpsc::Sender<Frame>,
820 lagged: Option<oneshot::Sender<SpawnCursor>>,
824}
825
826#[derive(Debug)]
827struct SpawnEventState {
828 daemon_incarnation: String,
829 seq: u64,
830 capacity: usize,
831 live: HashMap<String, LiveSpawn>,
832 generations: HashMap<String, u64>,
833 events: VecDeque<SpawnEvent>,
834 subscribers: HashMap<SpawnSubscriberKey, SpawnSubscriber>,
835}
836
837impl Default for SpawnEventState {
838 fn default() -> Self {
839 Self {
840 daemon_incarnation: "unconfigured".to_string(),
841 seq: 0,
842 capacity: SPAWN_EVENT_RING_CAPACITY,
843 live: HashMap::new(),
844 generations: HashMap::new(),
845 events: VecDeque::new(),
846 subscribers: HashMap::new(),
847 }
848 }
849}
850
851#[derive(Debug, Clone, Default)]
852struct SpawnEventFeed(Arc<Mutex<SpawnEventState>>);
853
854#[derive(Debug, Clone, PartialEq, Eq)]
855pub(crate) enum SpawnSubscribeRefusal {
856 ForeignIncarnation { current: String },
857 TooOld { oldest: SpawnCursor },
858 Frame(String),
859}
860
861impl SpawnEventFeed {
862 fn configure_incarnation(&self, daemon_incarnation: String) {
863 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
864 state.daemon_incarnation = daemon_incarnation;
865 state.seq = 0;
866 state.live.clear();
867 state.generations.clear();
868 state.events.clear();
869 state.subscribers.clear();
870 }
871
872 fn cursor(state: &SpawnEventState) -> SpawnCursor {
873 SpawnCursor {
874 daemon_incarnation: state.daemon_incarnation.clone(),
875 seq: state.seq,
876 }
877 }
878
879 fn snapshot(&self) -> SpawnSnapshot {
880 let state = self.0.lock().unwrap_or_else(|p| p.into_inner());
881 let mut live = state.live.values().cloned().collect::<Vec<_>>();
882 live.sort_by(|left, right| left.module_id.cmp(&right.module_id));
883 SpawnSnapshot {
884 cursor: Self::cursor(&state),
885 ring_bound: state.capacity as u64,
886 live,
887 }
888 }
889
890 fn emit_spawned(&self, module_id: &str, pid: u32, spawned_at_ms: u64) -> u64 {
891 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
892 let generation = state
893 .generations
894 .get(module_id)
895 .copied()
896 .unwrap_or(0)
897 .checked_add(1)
898 .expect("spawn generation exhausted");
899 state.generations.insert(module_id.to_string(), generation);
900 let live = LiveSpawn {
901 module_id: module_id.to_string(),
902 spawn_generation: generation,
903 pid,
904 spawned_at_ms,
905 };
906 state.live.insert(module_id.to_string(), live);
907 Self::emit_locked(
908 &mut state,
909 SpawnEventKind::Spawned,
910 module_id.to_string(),
911 generation,
912 pid,
913 None,
914 None,
915 );
916 generation
917 }
918
919 fn emit_exited(&self, module_id: &str, exit_code: Option<i32>, exit_signal: Option<i32>) {
920 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
921 let Some(live) = state.live.remove(module_id) else {
922 warn!(
923 module_id,
924 "terminal record had no live spawn event identity"
925 );
926 return;
927 };
928 Self::emit_locked(
929 &mut state,
930 SpawnEventKind::Exited,
931 module_id.to_string(),
932 live.spawn_generation,
933 live.pid,
934 exit_code,
935 exit_signal,
936 );
937 }
938
939 fn emit_superseded_exited(
946 &self,
947 module_id: &str,
948 spawn_generation: u64,
949 pid: u32,
950 exit_code: Option<i32>,
951 exit_signal: Option<i32>,
952 ) {
953 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
954 if state
955 .live
956 .get(module_id)
957 .is_some_and(|live| live.spawn_generation == spawn_generation)
958 {
959 state.live.remove(module_id);
960 }
961 Self::emit_locked(
962 &mut state,
963 SpawnEventKind::Exited,
964 module_id.to_string(),
965 spawn_generation,
966 pid,
967 exit_code,
968 exit_signal,
969 );
970 }
971
972 #[allow(clippy::too_many_arguments)]
973 fn emit_locked(
974 state: &mut SpawnEventState,
975 kind: SpawnEventKind,
976 module_id: String,
977 spawn_generation: u64,
978 pid: u32,
979 exit_code: Option<i32>,
980 exit_signal: Option<i32>,
981 ) {
982 state.seq = state
983 .seq
984 .checked_add(1)
985 .expect("spawn event sequence exhausted");
986 let event = SpawnEvent {
987 cursor: Self::cursor(state),
988 kind,
989 module_id,
990 spawn_generation,
991 pid,
992 exit_code,
993 exit_signal,
994 };
995 state.events.push_back(event.clone());
996 while state.events.len() > state.capacity {
997 state.events.pop_front();
998 }
999 let body = match serde_json::to_vec(&event) {
1000 Ok(body) => body,
1001 Err(error) => {
1002 error!(%error, "failed to serialize supervisor spawn event");
1003 return;
1004 }
1005 };
1006 state.subscribers.retain(|(connection_id, corr), subscriber| {
1007 let frame = Frame::build_with_version(
1008 subscriber.version,
1009 FrameType::StreamData,
1010 control_flags(),
1011 0,
1012 0,
1013 *corr,
1014 body.clone(),
1015 );
1016 match frame {
1017 Ok(frame) => {
1018 if subscriber.frames.try_send(frame).is_ok() {
1019 true
1020 } else {
1021 warn!(connection_id = connection_id.get(), corr, "dropping lagged supervisor spawn subscriber");
1022 if let Some(lagged) = subscriber.lagged.take() {
1023 let _ = lagged.send(event.cursor.clone());
1024 }
1025 false
1026 }
1027 }
1028 Err(error) => {
1029 warn!(connection_id = connection_id.get(), corr, %error, "dropping supervisor spawn subscriber after frame build failure");
1030 false
1031 }
1032 }
1033 });
1034 }
1035
1036 fn subscribe(
1037 &self,
1038 connection_id: ConnectionId,
1039 corr: u64,
1040 version: u8,
1041 since: Option<SpawnCursor>,
1042 sink: FrameSink,
1043 ) -> Result<(), SpawnSubscribeRefusal> {
1044 let (frames, mut receiver) = mpsc::channel(SPAWN_SUBSCRIBER_BUFFER);
1045 let (lagged, mut lagged_rx) = oneshot::channel::<SpawnCursor>();
1046 {
1047 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
1048 let replay = if let Some(since) = since {
1049 if since.daemon_incarnation != state.daemon_incarnation {
1050 return Err(SpawnSubscribeRefusal::ForeignIncarnation {
1051 current: state.daemon_incarnation.clone(),
1052 });
1053 }
1054 if let Some(oldest) = state.events.front().map(|event| event.cursor.clone()) {
1055 if since.seq < oldest.seq.saturating_sub(1) {
1056 return Err(SpawnSubscribeRefusal::TooOld { oldest });
1057 }
1058 }
1059 state
1060 .events
1061 .iter()
1062 .filter(|event| event.cursor.seq > since.seq)
1063 .cloned()
1064 .collect::<Vec<_>>()
1065 } else {
1066 Vec::new()
1067 };
1068 for event in replay {
1069 let body = serde_json::to_vec(&event)
1070 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1071 let frame = Frame::build_with_version(
1072 version,
1073 FrameType::StreamData,
1074 control_flags(),
1075 0,
1076 0,
1077 corr,
1078 body,
1079 )
1080 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1081 frames
1082 .try_send(frame)
1083 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1084 }
1085 state.subscribers.insert(
1086 (connection_id, corr),
1087 SpawnSubscriber {
1088 version,
1089 frames,
1090 lagged: Some(lagged),
1091 },
1092 );
1093 }
1094 tokio::spawn(async move {
1105 while let Some(frame) = receiver.recv().await {
1106 if sink.send(frame).await.is_err() {
1107 return;
1108 }
1109 }
1110 let Ok(first_undelivered) = lagged_rx.try_recv() else {
1111 return;
1112 };
1113 match spawn_subscriber_lagged_frame(version, corr, first_undelivered) {
1114 Ok(frame) => {
1115 let _ = sink.send(frame).await;
1116 }
1117 Err(error) => {
1118 error!(%error, corr, "failed to build lagged spawn subscriber terminal frame");
1119 }
1120 }
1121 });
1122 Ok(())
1123 }
1124
1125 fn cancel(&self, connection_id: ConnectionId, corr: u64) -> bool {
1126 let Some(subscriber) = self
1127 .0
1128 .lock()
1129 .unwrap_or_else(|p| p.into_inner())
1130 .subscribers
1131 .remove(&(connection_id, corr))
1132 else {
1133 return false;
1134 };
1135 if let Ok(frame) = Frame::build_with_version(
1136 subscriber.version,
1137 FrameType::StreamEnd,
1138 control_flags(),
1139 0,
1140 0,
1141 corr,
1142 Vec::new(),
1143 ) {
1144 tokio::spawn(async move {
1145 let _ = subscriber.frames.send(frame).await;
1146 });
1147 }
1148 true
1149 }
1150
1151 fn remove_connection(&self, connection_id: ConnectionId) {
1152 self.0
1153 .lock()
1154 .unwrap_or_else(|p| p.into_inner())
1155 .subscribers
1156 .retain(|(subscriber_connection, _), _| *subscriber_connection != connection_id);
1157 }
1158
1159 #[cfg(any(test, feature = "test-support"))]
1160 fn set_capacity(&self, capacity: usize) {
1161 self.0.lock().unwrap_or_else(|p| p.into_inner()).capacity = capacity;
1162 }
1163
1164 #[cfg(any(test, feature = "test-support"))]
1165 fn subscriber_count(&self) -> usize {
1166 self.0
1167 .lock()
1168 .unwrap_or_else(|p| p.into_inner())
1169 .subscribers
1170 .len()
1171 }
1172}
1173
1174fn spawn_subscriber_lagged_frame(
1177 version: u8,
1178 corr: u64,
1179 first_undelivered: SpawnCursor,
1180) -> Result<Frame, String> {
1181 let body = serde_json::to_vec(&subc_protocol::ErrorBody {
1182 code: SPAWN_SUBSCRIBER_LAGGED_CODE.to_string(),
1183 message: "spawn subscriber fell behind and was dropped; resubscribe from the last cursor received"
1184 .to_string(),
1185 detail: Some(serde_json::json!({
1186 "first_undelivered_cursor": first_undelivered
1187 })),
1188 })
1189 .map_err(|error| error.to_string())?;
1190 Frame::build_with_version(version, FrameType::Error, control_flags(), 0, 0, corr, body)
1191 .map_err(|error| error.to_string())
1192}
1193
1194pub trait ModuleProcessLiveness: Send + Sync {
1195 fn process_live(&self, module_id: &str) -> Option<bool>;
1196}
1197
1198#[derive(Debug, Clone, Default)]
1200pub struct SupervisorProcessLiveness {
1201 snapshots: Arc<Mutex<HashMap<String, SharedSnapshot>>>,
1202}
1203
1204impl SupervisorProcessLiveness {
1205 pub fn new() -> Self {
1206 Self::default()
1207 }
1208
1209 fn track(&self, module_id: String, snapshot: SharedSnapshot) {
1210 let mut snapshots = self
1211 .snapshots
1212 .lock()
1213 .unwrap_or_else(|poisoned| poisoned.into_inner());
1214 snapshots.insert(module_id, snapshot);
1215 }
1216
1217 fn untrack_if_current(&self, module_id: &str, snapshot: &SharedSnapshot) {
1218 let mut snapshots = self
1219 .snapshots
1220 .lock()
1221 .unwrap_or_else(|poisoned| poisoned.into_inner());
1222 let is_current = snapshots
1223 .get(module_id)
1224 .map(|tracked| Arc::ptr_eq(tracked, snapshot))
1225 .unwrap_or(false);
1226 if is_current {
1227 snapshots.remove(module_id);
1228 }
1229 }
1230}
1231
1232impl ModuleProcessLiveness for SupervisorProcessLiveness {
1233 fn process_live(&self, module_id: &str) -> Option<bool> {
1234 let snapshot = {
1235 let snapshots = self
1236 .snapshots
1237 .lock()
1238 .unwrap_or_else(|poisoned| poisoned.into_inner());
1239 snapshots.get(module_id).cloned()
1240 }?;
1241 let snapshot = snapshot
1242 .lock()
1243 .unwrap_or_else(|poisoned| poisoned.into_inner());
1244 Some(snapshot.state == ModuleState::Running && snapshot.process_alive)
1245 }
1246}
1247
1248#[derive(Debug, Clone)]
1249struct SupervisorRuntimeConfig {
1250 restart_policy: RestartPolicy,
1251 drain_timeout: Duration,
1254 effective_drain_timeout: Arc<Mutex<Duration>>,
1257 default_drain_timeout: Duration,
1260 health: HealthConfig,
1261 connection_file_path: Option<PathBuf>,
1262 capture_logs_dir: Option<PathBuf>,
1263 forwarding: Option<Arc<ForwardingTable>>,
1264 supervisor_handle: Option<SupervisorHandle>,
1267 stderr_ring: Arc<Mutex<StderrRing>>,
1274 terminal_ring: Arc<Mutex<TerminalRing>>,
1275 spawn_events: SpawnEventFeed,
1276 child_roster: ChildRoster,
1277 #[cfg(target_os = "linux")]
1278 cgroup_placement: Option<subc_cgroup::Placement>,
1279 #[cfg(test)]
1280 test_seed_stale_facts_before_enable_spawn: bool,
1281}
1282
1283#[derive(Debug, Clone, PartialEq, Eq)]
1284struct SupervisedConfiguration {
1285 spec: ModuleSpec,
1286 health: HealthConfig,
1287}
1288
1289#[derive(Debug, Clone, Default)]
1295pub struct SupervisorHandle {
1296 modules: Arc<Mutex<HashMap<String, SupervisedModule>>>,
1297 spawn_events: SpawnEventFeed,
1298 reserved_nonces: Arc<Mutex<HashMap<String, Option<String>>>>,
1309 removal_tombstones: Arc<Mutex<HashMap<String, u64>>>,
1315 spawn_nonces: Arc<Mutex<HashMap<String, String>>>,
1319 reserved_prefix_owners: Arc<Mutex<HashMap<String, String>>>,
1327 swaps: Arc<Mutex<HashMap<String, OpenSwap>>>,
1333 promotion_observer: PromotionObserverSlot,
1335 operation_lock: Arc<AsyncMutex<()>>,
1339}
1340
1341pub(crate) trait SwapPromotionObserver: Send + Sync {
1350 fn swap_promoted(&self, registration: &crate::registry::ModuleRegistration);
1351}
1352
1353#[derive(Clone, Default)]
1357struct PromotionObserverSlot(Arc<Mutex<Option<std::sync::Weak<dyn SwapPromotionObserver>>>>);
1358
1359impl fmt::Debug for PromotionObserverSlot {
1360 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1361 f.write_str("PromotionObserverSlot")
1362 }
1363}
1364
1365#[derive(Debug, Clone)]
1367struct OpenSwap {
1368 candidate_nonce: String,
1371 incumbent_nonce: Option<String>,
1376 candidate_admitted: bool,
1380}
1381
1382#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1385pub(crate) enum SwapHelloAdmission {
1386 NotSwapping,
1389 Candidate,
1391 Refused,
1394}
1395
1396#[derive(Debug, Clone, PartialEq, Eq)]
1397pub(crate) enum ReservedHelloRejection {
1398 Exact {
1399 module_id: String,
1400 },
1401 Prefix {
1402 prefix: String,
1403 owner_module_id: String,
1404 },
1405}
1406
1407impl SupervisorHandle {
1408 pub fn new() -> Self {
1409 Self::default()
1410 }
1411
1412 pub(crate) fn spawn_snapshot(&self) -> SpawnSnapshot {
1413 self.spawn_events.snapshot()
1414 }
1415
1416 pub(crate) fn subscribe_spawns(
1417 &self,
1418 connection_id: ConnectionId,
1419 corr: u64,
1420 version: u8,
1421 since: Option<SpawnCursor>,
1422 sink: FrameSink,
1423 ) -> Result<(), SpawnSubscribeRefusal> {
1424 self.spawn_events
1425 .subscribe(connection_id, corr, version, since, sink)
1426 }
1427
1428 pub(crate) fn cancel_spawn_subscription(&self, connection_id: ConnectionId, corr: u64) -> bool {
1429 self.spawn_events.cancel(connection_id, corr)
1430 }
1431
1432 pub(crate) fn remove_spawn_subscribers(&self, connection_id: ConnectionId) {
1433 self.spawn_events.remove_connection(connection_id);
1434 }
1435
1436 #[cfg(any(test, feature = "test-support"))]
1437 pub fn set_spawn_event_capacity_for_test(&self, capacity: usize) {
1438 assert!(capacity > 0, "spawn event capacity must be non-zero");
1439 self.spawn_events.set_capacity(capacity);
1440 }
1441
1442 #[cfg(any(test, feature = "test-support"))]
1443 pub fn spawn_subscriber_count_for_test(&self) -> usize {
1444 self.spawn_events.subscriber_count()
1445 }
1446
1447 pub fn set_spawn_nonce(&self, module_id: &str, nonce: String) {
1450 self.spawn_nonces
1451 .lock()
1452 .unwrap_or_else(|poisoned| poisoned.into_inner())
1453 .insert(module_id.to_string(), nonce);
1454 }
1455
1456 pub fn set_reserved_nonce(&self, module_id: &str, nonce: String) {
1459 self.reserved_nonces
1460 .lock()
1461 .unwrap_or_else(|poisoned| poisoned.into_inner())
1462 .insert(module_id.to_string(), Some(nonce));
1463 }
1464
1465 pub fn set_reserved_prefixes(&self, owner_module_id: &str, prefixes: &[String]) {
1467 let mut owners = self
1468 .reserved_prefix_owners
1469 .lock()
1470 .unwrap_or_else(|poisoned| poisoned.into_inner());
1471 owners.retain(|_, owner| owner != owner_module_id);
1472 for prefix in prefixes {
1473 owners.insert(prefix.clone(), owner_module_id.to_string());
1474 }
1475 }
1476
1477 #[cfg(test)]
1479 pub(crate) fn spawn_nonce(&self, module_id: &str) -> Option<String> {
1480 self.spawn_nonces
1481 .lock()
1482 .unwrap_or_else(|poisoned| poisoned.into_inner())
1483 .get(module_id)
1484 .cloned()
1485 }
1486
1487 fn apply_identity_configuration(&self, spec: &ModuleSpec) {
1488 self.set_reserved_prefixes(&spec.module_id, &spec.reserved_prefixes);
1489 let spawn_nonce = self
1490 .spawn_nonces
1491 .lock()
1492 .unwrap_or_else(|poisoned| poisoned.into_inner())
1493 .get(&spec.module_id)
1494 .cloned();
1495 let mut reserved_nonces = self
1496 .reserved_nonces
1497 .lock()
1498 .unwrap_or_else(|poisoned| poisoned.into_inner());
1499 if spec.reserved {
1500 reserved_nonces.insert(spec.module_id.clone(), spawn_nonce);
1505 }
1506 drop(reserved_nonces);
1507 self.removal_tombstones
1511 .lock()
1512 .unwrap_or_else(|poisoned| poisoned.into_inner())
1513 .remove(&spec.module_id);
1514 }
1515
1516 pub fn reserved_hello_authorized(&self, module_id: &str, presented: Option<&str>) -> bool {
1521 self.reserved_hello_rejection(module_id, presented)
1522 .is_none()
1523 }
1524
1525 pub(crate) fn reserved_hello_rejection(
1526 &self,
1527 module_id: &str,
1528 presented: Option<&str>,
1529 ) -> Option<ReservedHelloRejection> {
1530 let nonces = self
1531 .reserved_nonces
1532 .lock()
1533 .unwrap_or_else(|poisoned| poisoned.into_inner());
1534 if let Some(expected) = nonces.get(module_id) {
1535 let authorized = match expected {
1539 Some(expected) => {
1540 presented.is_some_and(|p| constant_time_eq(expected.as_bytes(), p.as_bytes()))
1541 }
1542 None => false,
1543 };
1544 if authorized {
1545 return None;
1546 }
1547 return Some(ReservedHelloRejection::Exact {
1548 module_id: module_id.to_string(),
1549 });
1550 }
1551 drop(nonces);
1552
1553 let matched_prefix = self
1554 .reserved_prefix_owners
1555 .lock()
1556 .unwrap_or_else(|poisoned| poisoned.into_inner())
1557 .iter()
1558 .filter(|(prefix, _)| module_id.starts_with(prefix.as_str()))
1559 .max_by_key(|(prefix, _)| prefix.len())
1560 .map(|(prefix, owner)| (prefix.clone(), owner.clone()));
1561 let (prefix, owner_module_id) = matched_prefix?;
1562
1563 let authorized = presented.is_some_and(|presented| {
1564 self.spawn_nonces
1565 .lock()
1566 .unwrap_or_else(|poisoned| poisoned.into_inner())
1567 .get(&owner_module_id)
1568 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()))
1569 || self.swap_nonce_matches(&owner_module_id, presented)
1572 });
1573 if authorized {
1574 None
1575 } else {
1576 Some(ReservedHelloRejection::Prefix {
1577 prefix,
1578 owner_module_id,
1579 })
1580 }
1581 }
1582
1583 pub fn spawned_consumer_authorized(&self, module_id: &str, presented: &str) -> bool {
1588 if presented.is_empty() {
1589 return false;
1590 }
1591 let nonces = self
1592 .spawn_nonces
1593 .lock()
1594 .unwrap_or_else(|poisoned| poisoned.into_inner());
1595 let current = nonces
1596 .get(module_id)
1597 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()));
1598 drop(nonces);
1599 current || self.swap_nonce_matches(module_id, presented)
1604 }
1605
1606 fn swap_nonce_matches(&self, module_id: &str, presented: &str) -> bool {
1608 let swaps = self
1609 .swaps
1610 .lock()
1611 .unwrap_or_else(|poisoned| poisoned.into_inner());
1612 swaps.get(module_id).is_some_and(|swap| {
1613 constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes())
1614 || swap.incumbent_nonce.as_deref().is_some_and(|incumbent| {
1615 constant_time_eq(incumbent.as_bytes(), presented.as_bytes())
1616 })
1617 })
1618 }
1619
1620 pub(crate) fn open_swap(&self, module_id: &str, candidate_nonce: String) {
1623 let incumbent_nonce = self
1624 .spawn_nonces
1625 .lock()
1626 .unwrap_or_else(|poisoned| poisoned.into_inner())
1627 .get(module_id)
1628 .cloned();
1629 self.swaps
1630 .lock()
1631 .unwrap_or_else(|poisoned| poisoned.into_inner())
1632 .insert(
1633 module_id.to_string(),
1634 OpenSwap {
1635 candidate_nonce,
1636 incumbent_nonce,
1637 candidate_admitted: false,
1638 },
1639 );
1640 }
1641
1642 pub(crate) fn close_swap(&self, module_id: &str) {
1645 self.swaps
1646 .lock()
1647 .unwrap_or_else(|poisoned| poisoned.into_inner())
1648 .remove(module_id);
1649 }
1650
1651 pub(crate) fn set_swap_promotion_observer(
1654 &self,
1655 observer: std::sync::Weak<dyn SwapPromotionObserver>,
1656 ) {
1657 *self
1658 .promotion_observer
1659 .0
1660 .lock()
1661 .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(observer);
1662 }
1663
1664 fn notify_swap_promoted(&self, registration: &crate::registry::ModuleRegistration) {
1667 let observer = self
1668 .promotion_observer
1669 .0
1670 .lock()
1671 .unwrap_or_else(|poisoned| poisoned.into_inner())
1672 .as_ref()
1673 .and_then(std::sync::Weak::upgrade);
1674 if let Some(observer) = observer {
1675 observer.swap_promoted(registration);
1676 }
1677 }
1678
1679 pub(crate) fn swap_open(&self, module_id: &str) -> bool {
1681 self.swaps
1682 .lock()
1683 .unwrap_or_else(|poisoned| poisoned.into_inner())
1684 .contains_key(module_id)
1685 }
1686
1687 fn promote_swap_nonce(&self, module_id: &str, reserved: bool) {
1692 let candidate_nonce = self
1693 .swaps
1694 .lock()
1695 .unwrap_or_else(|poisoned| poisoned.into_inner())
1696 .get(module_id)
1697 .map(|swap| swap.candidate_nonce.clone());
1698 let Some(nonce) = candidate_nonce else {
1699 return;
1700 };
1701 self.set_spawn_nonce(module_id, nonce.clone());
1702 if reserved {
1703 self.set_reserved_nonce(module_id, nonce);
1704 }
1705 }
1706
1707 pub(crate) fn swap_hello_admission(
1722 &self,
1723 module_id: &str,
1724 presented: Option<&str>,
1725 ) -> SwapHelloAdmission {
1726 let swaps = self
1727 .swaps
1728 .lock()
1729 .unwrap_or_else(|poisoned| poisoned.into_inner());
1730 let Some(swap) = swaps.get(module_id) else {
1731 return SwapHelloAdmission::NotSwapping;
1732 };
1733 let Some(presented) = presented else {
1734 return SwapHelloAdmission::Refused;
1735 };
1736 if constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes()) {
1737 return if swap.candidate_admitted {
1738 SwapHelloAdmission::Refused
1739 } else {
1740 SwapHelloAdmission::Candidate
1741 };
1742 }
1743 if swap
1744 .incumbent_nonce
1745 .as_deref()
1746 .is_some_and(|incumbent| constant_time_eq(incumbent.as_bytes(), presented.as_bytes()))
1747 {
1748 return SwapHelloAdmission::NotSwapping;
1749 }
1750 SwapHelloAdmission::Refused
1751 }
1752
1753 pub(crate) fn mark_swap_candidate_admitted(&self, module_id: &str) {
1756 if let Some(swap) = self
1757 .swaps
1758 .lock()
1759 .unwrap_or_else(|poisoned| poisoned.into_inner())
1760 .get_mut(module_id)
1761 {
1762 swap.candidate_admitted = true;
1763 }
1764 }
1765
1766 pub fn spawn_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1768 self.spawn_nonces
1769 .lock()
1770 .unwrap_or_else(|poisoned| poisoned.into_inner())
1771 .get(module_id)
1772 .cloned()
1773 }
1774
1775 pub fn reserved_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1777 self.reserved_nonces
1778 .lock()
1779 .unwrap_or_else(|poisoned| poisoned.into_inner())
1780 .get(module_id)
1781 .cloned()
1782 .flatten()
1783 }
1784
1785 pub fn insert(&self, module: SupervisedModule) -> Option<SupervisedModule> {
1786 let mut modules = self
1787 .modules
1788 .lock()
1789 .unwrap_or_else(|poisoned| poisoned.into_inner());
1790 modules.insert(module.module_id().to_string(), module)
1791 }
1792
1793 pub fn get(&self, module_id: &str) -> Option<SupervisedModule> {
1794 let modules = self
1795 .modules
1796 .lock()
1797 .unwrap_or_else(|poisoned| poisoned.into_inner());
1798 modules.get(module_id).cloned()
1799 }
1800
1801 pub(crate) fn record_late_health_answer(
1802 &self,
1803 module_id: &str,
1804 latency_ms: u64,
1805 ) -> Result<bool, SuperviseError> {
1806 let Some(module) = self.get(module_id) else {
1807 return Ok(false);
1808 };
1809 update_snapshot(&module.inner.snapshot, Some(module_id), |state| {
1810 state.health.late_answer_count = state.health.late_answer_count.saturating_add(1);
1811 state.health.last_late_answer_latency_ms = Some(latency_ms);
1812 state.health.consecutive_failures = 0;
1820 })?;
1821 Ok(true)
1822 }
1823
1824 pub fn record_deliberate_severance(&self, module_id: &str) -> Result<bool, SuperviseError> {
1830 let Some(module) = self.get(module_id) else {
1831 return Ok(false);
1832 };
1833 let status = module.status()?;
1834 let Some((pid, start_time)) = status.pid.zip(status.process_start_time) else {
1835 return Ok(false);
1836 };
1837 module.record_deliberate_severance(ProcessIdentity { pid, start_time })
1838 }
1839
1840 pub fn list(&self) -> Vec<SupervisedModule> {
1841 let modules = self
1842 .modules
1843 .lock()
1844 .unwrap_or_else(|poisoned| poisoned.into_inner());
1845 let mut modules = modules.values().cloned().collect::<Vec<_>>();
1846 modules.sort_by(|left, right| left.module_id().cmp(right.module_id()));
1847 modules
1848 }
1849
1850 pub(crate) fn retire(&self, module_id: &str) -> Option<SupervisedModule> {
1851 self.spawn_nonces
1852 .lock()
1853 .unwrap_or_else(|poisoned| poisoned.into_inner())
1854 .remove(module_id);
1855 self.close_swap(module_id);
1856 let mut reserved_nonces = self
1857 .reserved_nonces
1858 .lock()
1859 .unwrap_or_else(|poisoned| poisoned.into_inner());
1860 if reserved_nonces.contains_key(module_id) {
1861 reserved_nonces.insert(module_id.to_string(), None);
1864 }
1865 drop(reserved_nonces);
1866 self.reserved_prefix_owners
1867 .lock()
1868 .unwrap_or_else(|poisoned| poisoned.into_inner())
1869 .retain(|_, owner| owner != module_id);
1870 self.modules
1871 .lock()
1872 .unwrap_or_else(|poisoned| poisoned.into_inner())
1873 .remove(module_id)
1874 }
1875
1876 pub(crate) fn record_rescan_removal(&self, module_id: &str) {
1879 self.removal_tombstones
1880 .lock()
1881 .unwrap_or_else(|poisoned| poisoned.into_inner())
1882 .insert(module_id.to_string(), unix_ms_now());
1883 }
1884
1885 pub(crate) fn removal_tombstone_age_ms(&self, module_id: &str) -> Option<u64> {
1887 self.removal_tombstones
1888 .lock()
1889 .unwrap_or_else(|poisoned| poisoned.into_inner())
1890 .get(module_id)
1891 .copied()
1892 .map(|removed_at_ms| unix_ms_now().saturating_sub(removed_at_ms))
1893 }
1894
1895 pub(crate) fn release_retained_reserved_gate(&self, module_id: &str) -> bool {
1900 if self.get(module_id).is_some() {
1901 return false;
1902 }
1903 let mut reserved_nonces = self
1904 .reserved_nonces
1905 .lock()
1906 .unwrap_or_else(|poisoned| poisoned.into_inner());
1907 if !matches!(reserved_nonces.get(module_id), Some(None)) {
1908 return false;
1909 }
1910 reserved_nonces.remove(module_id);
1911 true
1912 }
1913
1914 pub(crate) fn operation_lock(&self) -> Arc<AsyncMutex<()>> {
1915 Arc::clone(&self.operation_lock)
1916 }
1917}
1918
1919#[derive(Debug, Clone)]
1921pub struct Supervisor {
1922 registry: Arc<Registry>,
1923 restart_policy: RestartPolicy,
1924 drain_timeout: Duration,
1925 connection_file_path: Option<PathBuf>,
1926 capture_logs_dir: Option<PathBuf>,
1927 forwarding: Option<Arc<ForwardingTable>>,
1928 process_liveness: Arc<SupervisorProcessLiveness>,
1929 supervisor_handle: Option<SupervisorHandle>,
1930 health: HealthConfig,
1931 daemon_start_clock: crate::clock::StartClock,
1932 terminal_journal: Option<Arc<crate::terminal_journal::TerminalJournal>>,
1933 spawn_events: SpawnEventFeed,
1934 provenance_probe: ExecutableIdentityProbe,
1935 child_roster: ChildRoster,
1938 #[cfg(target_os = "linux")]
1939 cgroup_placement: Option<subc_cgroup::Placement>,
1940}
1941
1942impl Supervisor {
1943 #[cfg(unix)]
1954 pub(crate) fn begin_daemon_shutdown(&self) {
1955 self.child_roster.close();
1956 if let Some(journal) = &self.terminal_journal {
1957 journal.stamp_shutdown();
1958 }
1959 }
1960
1961 #[cfg(unix)]
1965 pub(crate) async fn drain_for_daemon_shutdown(&self) -> Result<(), SuperviseError> {
1966 const NOTICE_BUDGET: Duration = Duration::from_millis(500);
1967 const DRAIN_BUDGET: Duration = Duration::from_secs(2);
1968 let Some(forwarding) = &self.forwarding else {
1969 return Ok(());
1970 };
1971 let module_ids = forwarding
1972 .begin_daemon_drain()
1973 .map_err(SuperviseError::Forwarding)?;
1974 let deadline_ms =
1975 unix_ms_now().saturating_add((NOTICE_BUDGET + DRAIN_BUDGET).as_millis() as u64);
1976 let mut notices = tokio::task::JoinSet::new();
1977 let mut drains = Vec::new();
1978 for module_id in module_ids {
1979 let Some(target) = forwarding
1980 .begin_module_drain(&module_id, RouteCloseReason::Restart)
1981 .map_err(SuperviseError::Forwarding)?
1982 else {
1983 continue;
1984 };
1985 let routes = forwarding
1986 .endpoint_routes(target.endpoint)
1987 .map_err(SuperviseError::Forwarding)?;
1988 let command = serde_json::to_vec(&ModuleControlCommand::Draining {
1994 reason: RouteCloseReason::Restart,
1995 deadline_ms,
1996 })
1997 .expect("module draining serializes");
1998 let closing = serde_json::to_vec(&ClientControlPush::RouteClosing {
1999 module_id: module_id.clone(),
2000 reason: RouteCloseReason::Restart,
2001 })
2002 .expect("route closing serializes");
2003 let mut recipients = vec![(target.sink.clone(), target.negotiated_ver, command)];
2004 let mut seen = std::collections::HashSet::new();
2005 for route in routes {
2006 let client = route.goodbye_target;
2007 if seen.insert(client.connection_id) {
2008 recipients.push((client.sink, client.negotiated_ver, closing.clone()));
2009 }
2010 }
2011 for (sink, version, body) in recipients {
2012 notices.spawn(async move {
2013 let frame = Frame::build_with_version(
2014 version,
2015 FrameType::Push,
2016 control_flags(),
2017 0,
2018 0,
2019 0,
2020 body,
2021 )
2022 .expect("bounded lifecycle notice frame builds");
2023 sink.send_flushed(frame).await
2024 });
2025 }
2026 let gauges = declared_busy_gauges(&self.registry, &module_id)?;
2027 drains.push((module_id, target.endpoint, gauges));
2028 }
2029 let notice_deadline = Instant::now() + NOTICE_BUDGET;
2032 while let Ok(Some(result)) = timeout_at(notice_deadline, notices.join_next()).await {
2033 if !matches!(result, Ok(Ok(()))) {
2034 warn!(?result, "daemon shutdown notice delivery failed");
2035 }
2036 }
2037 notices.abort_all();
2038 let deadline = Instant::now() + DRAIN_BUDGET;
2039 let mut waits = tokio::task::JoinSet::new();
2040 for (module_id, endpoint, gauges) in drains {
2041 let forwarding = Arc::clone(forwarding);
2042 let mut runtime = self.runtime_config();
2043 runtime.health.cadence = Duration::from_millis(100);
2044 waits.spawn(async move {
2045 wait_for_forwarding_quiescence(
2046 &forwarding,
2047 &module_id,
2048 &runtime,
2049 endpoint,
2050 deadline,
2051 &gauges,
2052 DrainScope::Active,
2053 )
2054 .await
2055 });
2056 }
2057 while let Ok(Some(result)) = timeout_at(deadline, waits.join_next()).await {
2058 if !matches!(result, Ok(Ok(true))) {
2059 warn!(?result, "daemon shutdown drain did not reach quiescence");
2060 }
2061 }
2062 Ok(())
2063 }
2064
2065 #[cfg(unix)]
2075 pub(crate) async fn end_children_for_daemon_shutdown(
2076 &self,
2077 already_escalated: bool,
2078 escalate: impl std::future::Future<Output = ()>,
2079 ) {
2080 if let Some(forwarding) = &self.forwarding {
2081 let closed = forwarding.close_all_connections(&CloseReason::new(
2082 "daemon_shutdown",
2083 "the daemon is exiting after its shutdown notice and drain",
2084 ));
2085 debug!(closed, "closed established connections for daemon shutdown");
2086 }
2087 crate::child_roster::end_children_for_daemon_shutdown(
2088 &self.child_roster,
2089 already_escalated,
2090 escalate,
2091 )
2092 .await;
2093 }
2094
2095 pub fn new(registry: Arc<Registry>, restart_policy: RestartPolicy) -> Self {
2096 Self {
2097 registry,
2098 restart_policy,
2099 drain_timeout: DEFAULT_DRAIN_TIMEOUT,
2100 connection_file_path: None,
2101 capture_logs_dir: None,
2102 forwarding: None,
2103 process_liveness: Arc::new(SupervisorProcessLiveness::default()),
2104 supervisor_handle: None,
2105 health: HealthConfig::default(),
2106 daemon_start_clock: crate::clock::StartClock::capture(),
2107 terminal_journal: None,
2108 spawn_events: SpawnEventFeed::default(),
2109 provenance_probe: ExecutableIdentityProbe::default(),
2110 child_roster: ChildRoster::default(),
2111 #[cfg(target_os = "linux")]
2112 cgroup_placement: None,
2113 }
2114 }
2115
2116 pub fn with_drain_timeout(mut self, drain_timeout: Duration) -> Self {
2117 self.drain_timeout = drain_timeout;
2118 self
2119 }
2120
2121 pub fn with_process_liveness(
2122 mut self,
2123 process_liveness: Arc<SupervisorProcessLiveness>,
2124 ) -> Self {
2125 self.process_liveness = process_liveness;
2126 self
2127 }
2128
2129 pub fn with_connection_file_path(mut self, connection_file_path: impl Into<PathBuf>) -> Self {
2130 self.connection_file_path = Some(connection_file_path.into());
2131 self
2132 }
2133
2134 pub fn with_capture_logs_dir(mut self, logs_dir: impl Into<PathBuf>) -> Self {
2136 self.capture_logs_dir = Some(logs_dir.into());
2137 self
2138 }
2139
2140 pub fn with_daemon_incarnation(self, daemon_incarnation: String) -> Self {
2143 self.spawn_events.configure_incarnation(daemon_incarnation);
2147 self
2148 }
2149
2150 pub fn with_terminal_journal(self, path: PathBuf, daemon_incarnation: String) -> Self {
2153 let mut this = self.with_daemon_incarnation(daemon_incarnation.clone());
2154 this.terminal_journal = Some(Arc::new(crate::terminal_journal::TerminalJournal::open(
2155 path,
2156 daemon_incarnation,
2157 )));
2158 this
2159 }
2160
2161 pub fn with_forwarding(mut self, forwarding: Arc<ForwardingTable>) -> Self {
2162 self.forwarding = Some(forwarding);
2163 self
2164 }
2165
2166 pub fn with_handle(mut self, supervisor_handle: SupervisorHandle) -> Self {
2167 self.spawn_events = supervisor_handle.spawn_events.clone();
2168 self.supervisor_handle = Some(supervisor_handle);
2169 self
2170 }
2171
2172 pub fn with_health_config(mut self, health: HealthConfig) -> Self {
2173 self.health = health;
2174 self
2175 }
2176
2177 #[cfg(target_os = "linux")]
2178 pub fn with_cgroup_placement(
2179 mut self,
2180 cgroup_placement: Option<subc_cgroup::Placement>,
2181 ) -> Self {
2182 self.cgroup_placement = cgroup_placement;
2183 self
2184 }
2185
2186 pub fn spawn(&self, spec: ModuleSpec) -> Result<SupervisedModule, SuperviseError> {
2192 validate_spec(&spec)?;
2193
2194 let runtime = self.runtime_config();
2195 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2196 let child = spawn_child(
2197 &spec,
2198 runtime.connection_file_path.as_deref(),
2199 self.supervisor_handle.as_ref(),
2200 &runtime.stderr_ring,
2201 runtime.capture_logs_dir.as_deref(),
2202 &runtime.child_roster,
2203 #[cfg(target_os = "linux")]
2204 runtime.cgroup_placement.as_ref(),
2205 )?;
2206 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2207 self.process_liveness
2208 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2209
2210 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2211 }
2212
2213 pub fn supervise_configured(
2219 &self,
2220 spec: ModuleSpec,
2221 enabled: bool,
2222 ) -> Result<SupervisedModule, SuperviseError> {
2223 validate_spec(&spec)?;
2224
2225 let runtime = self.runtime_config();
2226 if !enabled {
2227 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2228 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2229 }
2230
2231 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2232 match spawn_child(
2233 &spec,
2234 runtime.connection_file_path.as_deref(),
2235 self.supervisor_handle.as_ref(),
2236 &runtime.stderr_ring,
2237 runtime.capture_logs_dir.as_deref(),
2238 &runtime.child_roster,
2239 #[cfg(target_os = "linux")]
2240 runtime.cgroup_placement.as_ref(),
2241 ) {
2242 Ok(child) => {
2243 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2244 self.process_liveness
2245 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2246 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2247 }
2248 Err(err) => {
2249 error!(
2250 module_id = %spec.module_id,
2251 program = %spec.program.display(),
2252 error = %err,
2253 "configured module failed to spawn; marking failed and continuing"
2254 );
2255 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2256 Ok(self.supervised_module(spec, runtime, snapshot, None))
2257 }
2258 }
2259 }
2260
2261 pub fn supervise_configured_with_health(
2267 &self,
2268 spec: ModuleSpec,
2269 enabled: bool,
2270 health: HealthConfig,
2271 drain_timeout_ms: Option<u64>,
2272 restart_policy: RestartPolicy,
2273 ) -> Result<SupervisedModule, SuperviseError> {
2274 validate_spec(&spec)?;
2275
2276 let mut runtime = self.runtime_config();
2277 runtime.health = health;
2278 runtime.restart_policy = restart_policy;
2279 if let Some(ms) = drain_timeout_ms {
2280 runtime.drain_timeout = Duration::from_millis(ms);
2281 *runtime
2282 .effective_drain_timeout
2283 .lock()
2284 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
2285 }
2286 if !enabled {
2287 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2288 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2289 }
2290
2291 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2292 match spawn_child(
2293 &spec,
2294 runtime.connection_file_path.as_deref(),
2295 self.supervisor_handle.as_ref(),
2296 &runtime.stderr_ring,
2297 runtime.capture_logs_dir.as_deref(),
2298 &runtime.child_roster,
2299 #[cfg(target_os = "linux")]
2300 runtime.cgroup_placement.as_ref(),
2301 ) {
2302 Ok(child) => {
2303 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2304 self.process_liveness
2305 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2306 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2307 }
2308 Err(err) => {
2309 if health.critical {
2310 error!(
2311 module_id = %spec.module_id,
2312 program = %spec.program.display(),
2313 error = %err,
2314 "critical configured module failed to spawn; marking failed and alerting"
2315 );
2316 } else {
2317 error!(
2318 module_id = %spec.module_id,
2319 program = %spec.program.display(),
2320 error = %err,
2321 "configured module failed to spawn; marking failed and continuing"
2322 );
2323 }
2324 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2325 Ok(self.supervised_module(spec, runtime, snapshot, None))
2326 }
2327 }
2328 }
2329
2330 fn runtime_config(&self) -> SupervisorRuntimeConfig {
2331 let effective_drain_timeout = Arc::new(Mutex::new(self.drain_timeout));
2332 SupervisorRuntimeConfig {
2333 restart_policy: self.restart_policy,
2334 drain_timeout: self.drain_timeout,
2335 child_roster: self
2338 .child_roster
2339 .for_module(Arc::clone(&effective_drain_timeout)),
2340 effective_drain_timeout,
2341 default_drain_timeout: self.drain_timeout,
2342 health: self.health,
2343 connection_file_path: self.connection_file_path.clone(),
2344 capture_logs_dir: self.capture_logs_dir.clone(),
2345 forwarding: self.forwarding.clone(),
2346 supervisor_handle: self.supervisor_handle.clone(),
2347 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
2348 terminal_ring: Arc::new(Mutex::new(
2349 TerminalRing::new(
2350 TerminalRingConfig::default(),
2351 self.daemon_start_clock.started_at_ms(),
2352 )
2353 .with_start_clock(self.daemon_start_clock)
2354 .with_journal(self.terminal_journal.clone())
2355 .with_daemon_shutdown(self.child_roster.shutdown_flag()),
2356 )),
2357 spawn_events: self.spawn_events.clone(),
2358 #[cfg(target_os = "linux")]
2359 cgroup_placement: self.cgroup_placement.clone(),
2360 #[cfg(test)]
2361 test_seed_stale_facts_before_enable_spawn: false,
2362 }
2363 }
2364
2365 fn supervised_module(
2366 &self,
2367 spec: ModuleSpec,
2368 runtime: SupervisorRuntimeConfig,
2369 snapshot: SharedSnapshot,
2370 child: Option<SupervisedChild>,
2371 ) -> SupervisedModule {
2372 let configuration = Arc::new(Mutex::new(SupervisedConfiguration {
2373 spec: spec.clone(),
2374 health: runtime.health,
2375 }));
2376 let stderr_ring = Arc::clone(&runtime.stderr_ring);
2377 let terminal_ring = Arc::clone(&runtime.terminal_ring);
2378 let restart_policy = runtime.restart_policy;
2382 let effective_drain_timeout = Arc::clone(&runtime.effective_drain_timeout);
2383 let (tx, rx) = mpsc::channel(4);
2384 let monitor = tokio::spawn(supervise_loop(
2385 spec.clone(),
2386 runtime,
2387 Arc::clone(&self.registry),
2388 Arc::clone(&self.process_liveness),
2389 Arc::clone(&snapshot),
2390 child,
2391 rx,
2392 ));
2393
2394 let module_id = spec.module_id.clone();
2395 let module = SupervisedModule {
2396 inner: Arc::new(SupervisedModuleInner {
2397 module_id: module_id.clone(),
2398 registry: Arc::clone(&self.registry),
2399 snapshot,
2400 configuration,
2401 stderr_ring,
2402 terminal_ring,
2403 commands: tx,
2404 monitor: Mutex::new(Some(monitor)),
2405 restart_policy,
2406 effective_drain_timeout,
2407 provenance_probe: self.provenance_probe.clone(),
2408 }),
2409 };
2410 if let Some(supervisor_handle) = &self.supervisor_handle {
2411 supervisor_handle.apply_identity_configuration(&spec);
2412 supervisor_handle.insert(module.clone());
2413 }
2414 module
2415 }
2416}
2417
2418impl Default for Supervisor {
2419 fn default() -> Self {
2420 Self::new(Arc::new(Registry::default()), RestartPolicy::default())
2421 }
2422}
2423
2424#[derive(Clone)]
2426pub struct SupervisedModule {
2427 inner: Arc<SupervisedModuleInner>,
2428}
2429
2430struct SupervisedModuleInner {
2431 module_id: String,
2432 registry: Arc<Registry>,
2433 snapshot: SharedSnapshot,
2434 configuration: Arc<Mutex<SupervisedConfiguration>>,
2435 stderr_ring: Arc<Mutex<StderrRing>>,
2436 terminal_ring: Arc<Mutex<TerminalRing>>,
2437 commands: mpsc::Sender<SupervisorCommand>,
2438 monitor: Mutex<Option<JoinHandle<()>>>,
2439 restart_policy: RestartPolicy,
2443 effective_drain_timeout: Arc<Mutex<Duration>>,
2444 provenance_probe: ExecutableIdentityProbe,
2445}
2446
2447impl fmt::Debug for SupervisedModule {
2448 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2449 f.debug_struct("SupervisedModule")
2450 .field("module_id", &self.inner.module_id)
2451 .field("status", &self.status())
2452 .finish_non_exhaustive()
2453 }
2454}
2455
2456impl SupervisedModule {
2457 pub fn module_id(&self) -> &str {
2458 &self.inner.module_id
2459 }
2460
2461 #[cfg(test)]
2465 pub(crate) fn record_health_probe_failure_for_test(
2466 &self,
2467 detail: &str,
2468 ) -> Result<(), SuperviseError> {
2469 update_snapshot(&self.inner.snapshot, Some(&self.inner.module_id), |state| {
2470 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
2471 state.health.detail = Some(detail.to_string());
2472 })
2473 }
2474
2475 pub fn state(&self) -> Result<ModuleState, SuperviseError> {
2476 Ok(lock_snapshot(&self.inner.snapshot)?.state)
2477 }
2478
2479 pub fn stderr_tail(
2486 &self,
2487 max_lines: Option<usize>,
2488 max_bytes: Option<usize>,
2489 ) -> StderrTailSnapshot {
2490 self.inner
2491 .stderr_ring
2492 .lock()
2493 .unwrap_or_else(|poisoned| poisoned.into_inner())
2494 .snapshot(max_lines, max_bytes)
2495 }
2496
2497 pub fn terminal_history(&self) -> TerminalHistorySnapshot {
2502 self.inner
2503 .terminal_ring
2504 .lock()
2505 .unwrap_or_else(|poisoned| poisoned.into_inner())
2506 .snapshot()
2507 }
2508
2509 pub fn durable_terminal_history(&self) -> subc_control::TerminalHistory {
2514 durable_terminal_history_of(&self.inner.terminal_ring, &self.inner.module_id)
2515 }
2516
2517 pub(crate) async fn read_durable_terminal_history(
2522 &self,
2523 ) -> Result<subc_control::TerminalHistory, tokio::task::JoinError> {
2524 let terminal_ring = Arc::clone(&self.inner.terminal_ring);
2525 let module_id = self.inner.module_id.clone();
2526 tokio::task::spawn_blocking(move || durable_terminal_history_of(&terminal_ring, &module_id))
2527 .await
2528 }
2529
2530 pub fn status(&self) -> Result<ModuleStatus, SuperviseError> {
2531 self.status_with_snapshot_lock(&self.inner.snapshot, None)
2532 }
2533
2534 pub(crate) fn record_deliberate_severance(
2535 &self,
2536 identity: ProcessIdentity,
2537 ) -> Result<bool, SuperviseError> {
2538 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2539 if snapshot.pid != Some(identity.pid)
2540 || snapshot.process_start_time != Some(identity.start_time)
2541 {
2542 return Ok(false);
2543 }
2544 snapshot.deliberate_severance = Some(identity);
2545 Ok(true)
2546 }
2547
2548 pub(crate) fn status_for_control(
2553 &self,
2554 caller: &'static str,
2555 ) -> Result<ModuleStatus, SuperviseError> {
2556 self.status_with_snapshot_lock(&self.inner.snapshot, Some(caller))
2557 }
2558
2559 fn status_with_snapshot_lock(
2560 &self,
2561 snapshot: &SharedSnapshot,
2562 caller: Option<&'static str>,
2563 ) -> Result<ModuleStatus, SuperviseError> {
2564 let mut guard = match caller {
2565 Some(caller) => lock_snapshot_for_control(snapshot, &self.inner.module_id, caller)?,
2566 None => lock_snapshot(snapshot)?,
2567 };
2568 let restart_count =
2571 guard.crash_restarts_in_window(self.inner.restart_policy.window, Instant::now());
2572 let snapshot = guard.clone();
2573 drop(guard);
2574 let drain_timeout = *self.inner.effective_drain_timeout.lock().map_err(|_| {
2575 SuperviseError::StatePoisoned {
2576 module_id: Some(self.inner.module_id.clone()),
2577 }
2578 })?;
2579 let registration_active = self
2580 .inner
2581 .registry
2582 .get_module(&self.inner.module_id)
2583 .map_err(SuperviseError::Registry)?
2584 .is_some();
2585 let protocol = self.declared_protocol()?;
2586 let running_process =
2587 snapshot.enabled && snapshot.state == ModuleState::Running && snapshot.process_alive;
2588 let live = match protocol {
2594 ModuleProtocol::Subc => running_process && registration_active,
2595 ModuleProtocol::None => running_process,
2596 };
2597
2598 Ok(ModuleStatus {
2599 module_id: self.inner.module_id.clone(),
2600 state: snapshot.state,
2601 enabled: snapshot.enabled,
2602 process_alive: snapshot.process_alive,
2603 registration_active,
2604 protocol,
2605 live,
2606 restart_count,
2607 lifetime_restarts: snapshot.lifetime_restarts,
2608 spawn_generation: snapshot.spawn_generation,
2609 max_restarts: self.inner.restart_policy.max_restarts,
2610 restart_window: self.inner.restart_policy.window,
2611 drain_timeout,
2612 restart_backoff: self.inner.restart_policy.backoff,
2613 restart_max_backoff: self.inner.restart_policy.max_backoff,
2614 pid: snapshot.pid,
2615 spawned_at_ms: snapshot.spawned_at_ms,
2616 spawned_from: snapshot.spawned_from,
2617 process_start_time: snapshot.process_start_time,
2618 last_exit: snapshot.last_exit,
2619 health: snapshot.health,
2620 })
2621 }
2622
2623 #[cfg(test)]
2624 pub(crate) fn hold_snapshot_for_test(
2625 &self,
2626 acquired: std::sync::mpsc::Sender<()>,
2627 hold: Duration,
2628 ) -> std::thread::JoinHandle<()> {
2629 let snapshot = Arc::clone(&self.inner.snapshot);
2630 std::thread::spawn(move || {
2631 let _guard = snapshot.lock().expect("test snapshot lock is not poisoned");
2632 acquired
2633 .send(())
2634 .expect("test receiver waits for snapshot lock");
2635 std::thread::sleep(hold);
2636 })
2637 }
2638
2639 pub(crate) async fn running_image_agreement(&self) -> subc_control::RunningImageAgreement {
2640 let snapshot = match lock_snapshot(&self.inner.snapshot) {
2641 Ok(snapshot) => snapshot.clone(),
2642 Err(_) => {
2643 return subc_control::RunningImageAgreement::Unavailable {
2644 reason: subc_control::RunningImageUnavailableReason::NotRunning,
2645 };
2646 }
2647 };
2648 self.inner
2649 .provenance_probe
2650 .observe(
2651 snapshot.pid,
2652 snapshot.spawned_from.as_deref(),
2653 snapshot.spawned_file_identity,
2654 snapshot.process_start_time,
2655 )
2656 .await
2657 }
2658
2659 pub(crate) fn will_recover_after_connection_loss(&self) -> Result<bool, SuperviseError> {
2660 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2661 Ok(match snapshot.state {
2662 ModuleState::Restarting => true,
2663 ModuleState::Failed | ModuleState::Disabled => false,
2664 _ => daemon_will_restart(&mut snapshot, &self.inner.restart_policy, Instant::now()),
2665 })
2666 }
2667
2668 #[cfg(test)]
2669 pub(crate) fn is_warming(&self) -> Result<bool, SuperviseError> {
2670 self.is_warming_with_snapshot_lock(None)
2671 }
2672
2673 pub(crate) fn is_warming_for_control(
2674 &self,
2675 caller: &'static str,
2676 ) -> Result<bool, SuperviseError> {
2677 self.is_warming_with_snapshot_lock(Some(caller))
2678 }
2679
2680 fn is_warming_with_snapshot_lock(
2681 &self,
2682 caller: Option<&'static str>,
2683 ) -> Result<bool, SuperviseError> {
2684 let snapshot = match caller {
2685 Some(caller) => {
2686 lock_snapshot_for_control(&self.inner.snapshot, &self.inner.module_id, caller)?
2687 }
2688 None => lock_snapshot(&self.inner.snapshot)?,
2689 }
2690 .clone();
2691 Ok(matches!(
2692 snapshot.state,
2693 ModuleState::Starting | ModuleState::Running | ModuleState::Restarting
2694 ))
2695 }
2696
2697 pub async fn drain(&self) -> Result<(), SuperviseError> {
2699 self.stop().await
2700 }
2701
2702 pub(crate) async fn retire(&self) -> Result<(), SuperviseError> {
2703 match self.state()? {
2704 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2705 ModuleState::Starting
2706 | ModuleState::Running
2707 | ModuleState::Unresponsive
2708 | ModuleState::Restarting
2709 | ModuleState::Draining
2710 | ModuleState::Disabled => {}
2711 }
2712
2713 let (reply_tx, reply_rx) = oneshot::channel();
2714 self.inner
2715 .commands
2716 .send(SupervisorCommand::Retire { reply: reply_tx })
2717 .await
2718 .map_err(|_| SuperviseError::CommandClosed {
2719 module_id: self.inner.module_id.clone(),
2720 })?;
2721 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2722 module_id: self.inner.module_id.clone(),
2723 })?
2724 }
2725
2726 pub async fn stop(&self) -> Result<(), SuperviseError> {
2727 match self.state()? {
2728 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2729 ModuleState::Starting
2730 | ModuleState::Running
2731 | ModuleState::Unresponsive
2732 | ModuleState::Restarting
2733 | ModuleState::Draining
2734 | ModuleState::Disabled => {}
2735 }
2736
2737 let (reply_tx, reply_rx) = oneshot::channel();
2738 self.inner
2739 .commands
2740 .send(SupervisorCommand::Drain { reply: reply_tx })
2741 .await
2742 .map_err(|_| SuperviseError::CommandClosed {
2743 module_id: self.inner.module_id.clone(),
2744 })?;
2745 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2746 module_id: self.inner.module_id.clone(),
2747 })?
2748 }
2749
2750 pub async fn restart(&self, drain_timeout_ms: Option<u64>) -> Result<(), SuperviseError> {
2751 let (reply_tx, reply_rx) = oneshot::channel();
2752 self.inner
2753 .commands
2754 .send(SupervisorCommand::Restart {
2755 drain_timeout_ms,
2756 reply: reply_tx,
2757 })
2758 .await
2759 .map_err(|_| SuperviseError::CommandClosed {
2760 module_id: self.inner.module_id.clone(),
2761 })?;
2762 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2763 module_id: self.inner.module_id.clone(),
2764 })?
2765 }
2766
2767 pub async fn swap(&self, ready_timeout: Option<Duration>) -> Result<(), SuperviseError> {
2772 let (reply_tx, reply_rx) = oneshot::channel();
2773 self.inner
2774 .commands
2775 .send(SupervisorCommand::Swap {
2776 ready_timeout,
2777 reply: reply_tx,
2778 })
2779 .await
2780 .map_err(|_| SuperviseError::CommandClosed {
2781 module_id: self.inner.module_id.clone(),
2782 })?;
2783 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2784 module_id: self.inner.module_id.clone(),
2785 })?
2786 }
2787
2788 pub async fn reload(&self) -> Result<(), SuperviseError> {
2789 let (reply_tx, reply_rx) = oneshot::channel();
2790 self.inner
2791 .commands
2792 .send(SupervisorCommand::Reload { reply: reply_tx })
2793 .await
2794 .map_err(|_| SuperviseError::CommandClosed {
2795 module_id: self.inner.module_id.clone(),
2796 })?;
2797 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2798 module_id: self.inner.module_id.clone(),
2799 })?
2800 }
2801
2802 pub async fn set_enabled(&self, enabled: bool) -> Result<bool, SuperviseError> {
2803 let (reply_tx, reply_rx) = oneshot::channel();
2804 self.inner
2805 .commands
2806 .send(SupervisorCommand::SetEnabled {
2807 enabled,
2808 reply: reply_tx,
2809 })
2810 .await
2811 .map_err(|_| SuperviseError::CommandClosed {
2812 module_id: self.inner.module_id.clone(),
2813 })?;
2814 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2815 module_id: self.inner.module_id.clone(),
2816 })?
2817 }
2818
2819 pub(crate) fn declared_protocol(&self) -> Result<ModuleProtocol, SuperviseError> {
2824 Ok(self
2825 .inner
2826 .configuration
2827 .lock()
2828 .map_err(|_| SuperviseError::StatePoisoned {
2829 module_id: Some(self.inner.module_id.clone()),
2830 })?
2831 .spec
2832 .protocol)
2833 }
2834
2835 pub(crate) fn configuration(&self) -> Result<(ModuleSpec, HealthConfig), SuperviseError> {
2836 let configuration =
2837 self.inner
2838 .configuration
2839 .lock()
2840 .map_err(|_| SuperviseError::StatePoisoned {
2841 module_id: Some(self.inner.module_id.clone()),
2842 })?;
2843 Ok((configuration.spec.clone(), configuration.health))
2844 }
2845
2846 #[cfg(any(test, feature = "test-support"))]
2850 pub async fn update_spec_for_test(&self, spec: ModuleSpec) -> Result<(), SuperviseError> {
2851 let (_, health) = self.configuration()?;
2852 let drain_timeout_ms = u64::try_from(
2853 self.inner
2854 .effective_drain_timeout
2855 .lock()
2856 .unwrap_or_else(|poisoned| poisoned.into_inner())
2857 .as_millis(),
2858 )
2859 .ok();
2860 self.update_configuration(spec, health, drain_timeout_ms)
2861 .await
2862 }
2863
2864 pub(crate) async fn update_configuration(
2865 &self,
2866 spec: ModuleSpec,
2867 health: HealthConfig,
2868 drain_timeout_ms: Option<u64>,
2869 ) -> Result<(), SuperviseError> {
2870 if spec.module_id != self.inner.module_id {
2871 return Err(SuperviseError::InvalidSpec {
2872 reason: "a supervised module's module_id cannot be changed".to_string(),
2873 });
2874 }
2875 validate_spec(&spec)?;
2876 let (reply_tx, reply_rx) = oneshot::channel();
2877 self.inner
2878 .commands
2879 .send(SupervisorCommand::UpdateConfiguration {
2880 spec: spec.clone(),
2881 health,
2882 drain_timeout_ms,
2883 reply: reply_tx,
2884 })
2885 .await
2886 .map_err(|_| SuperviseError::CommandClosed {
2887 module_id: self.inner.module_id.clone(),
2888 })?;
2889 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2890 module_id: self.inner.module_id.clone(),
2891 })?;
2892 let mut configuration =
2893 self.inner
2894 .configuration
2895 .lock()
2896 .map_err(|_| SuperviseError::StatePoisoned {
2897 module_id: Some(self.inner.module_id.clone()),
2898 })?;
2899 configuration.spec = spec;
2900 configuration.health = health;
2901 Ok(())
2902 }
2903}
2904
2905impl Drop for SupervisedModuleInner {
2906 fn drop(&mut self) {
2907 let Ok(mut monitor) = self.monitor.lock() else {
2908 return;
2909 };
2910 if let Some(monitor) = monitor.as_ref().filter(|monitor| !monitor.is_finished()) {
2911 let _ = update_snapshot(&self.snapshot, Some(&self.module_id), |state| {
2912 state.state = ModuleState::Stopped;
2913 clear_current_process_facts(state);
2914 });
2915 monitor.abort();
2916 }
2917 let _ = monitor.take();
2918 }
2919}
2920
2921#[derive(Debug)]
2922enum SupervisorCommand {
2923 Drain {
2924 reply: oneshot::Sender<Result<(), SuperviseError>>,
2925 },
2926 Retire {
2927 reply: oneshot::Sender<Result<(), SuperviseError>>,
2928 },
2929 Restart {
2930 drain_timeout_ms: Option<u64>,
2935 reply: oneshot::Sender<Result<(), SuperviseError>>,
2936 },
2937 Reload {
2938 reply: oneshot::Sender<Result<(), SuperviseError>>,
2939 },
2940 SetEnabled {
2941 enabled: bool,
2942 reply: oneshot::Sender<Result<bool, SuperviseError>>,
2943 },
2944 UpdateConfiguration {
2945 spec: ModuleSpec,
2946 health: HealthConfig,
2947 drain_timeout_ms: Option<u64>,
2950 reply: oneshot::Sender<()>,
2951 },
2952 Swap {
2953 ready_timeout: Option<Duration>,
2956 reply: oneshot::Sender<Result<(), SuperviseError>>,
2958 },
2959}
2960
2961#[derive(Debug)]
2962pub enum SuperviseError {
2963 InvalidSpec {
2964 reason: String,
2965 },
2966 Spawn {
2967 program: PathBuf,
2968 source: io::Error,
2969 cgroup_path: Option<PathBuf>,
2970 },
2971 Cgroup {
2972 module_id: String,
2973 source: io::Error,
2974 },
2975 LaunchNonce {
2978 reason: String,
2979 },
2980 Wait {
2981 module_id: String,
2982 source: io::Error,
2983 },
2984 Kill {
2985 module_id: String,
2986 source: io::Error,
2987 },
2988 Forwarding(ForwardingError),
2989 Registry(RegistryError),
2990 ReloadUnavailable {
2991 module_id: String,
2992 reason: String,
2993 },
2994 Disabled {
2999 module_id: String,
3000 },
3001 ReloadFailed {
3002 module_id: String,
3003 reason: String,
3004 },
3005 RegistrationStillActive {
3006 module_id: String,
3007 waited: Duration,
3008 },
3009 StatePoisoned {
3010 module_id: Option<String>,
3011 },
3012 CommandClosed {
3013 module_id: String,
3014 },
3015 SwapInProgress {
3019 module_id: String,
3020 },
3021 SwapRefused {
3023 module_id: String,
3024 reason: SwapRefusal,
3025 },
3026 SwapFailed {
3030 module_id: String,
3031 arm: SwapFailureArm,
3032 detail: String,
3033 candidate_exit: Option<ExitReport>,
3036 },
3037}
3038
3039#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3041pub enum SwapRefusal {
3042 OverlapExclusive,
3044 NotRegistered,
3047 ProtocolNone,
3050 NotConfigured,
3053 AlreadySwapping,
3055}
3056
3057impl SwapRefusal {
3058 pub fn as_str(self) -> &'static str {
3059 match self {
3060 Self::OverlapExclusive => "overlap_exclusive",
3061 Self::NotRegistered => "not_registered",
3062 Self::ProtocolNone => "protocol_none",
3063 Self::NotConfigured => "not_configured",
3064 Self::AlreadySwapping => "already_swapping",
3065 }
3066 }
3067}
3068
3069#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3072pub enum SwapFailureArm {
3073 SpawnFailed,
3075 NeverRegistered,
3077 NeverReady,
3079 CandidateExited,
3081 CandidateUnhealthy,
3083 Interrupted,
3087 CutoverLost,
3092}
3093
3094impl SwapFailureArm {
3095 pub fn as_str(self) -> &'static str {
3096 match self {
3097 Self::SpawnFailed => "spawn_failed",
3098 Self::NeverRegistered => "never_registered",
3099 Self::NeverReady => "never_ready",
3100 Self::CandidateExited => "candidate_exited",
3101 Self::CandidateUnhealthy => "candidate_unhealthy",
3102 Self::Interrupted => "interrupted",
3103 Self::CutoverLost => "cutover_lost",
3104 }
3105 }
3106}
3107
3108impl fmt::Display for SuperviseError {
3109 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3110 match self {
3111 Self::InvalidSpec { reason } => write!(f, "invalid module spec: {reason}"),
3112 Self::Spawn {
3113 program,
3114 source,
3115 cgroup_path: Some(cgroup_path),
3116 } => write!(
3117 f,
3118 "failed to place module in cgroup '{}' while spawning '{}': {source}",
3119 cgroup_path.display(),
3120 program.display()
3121 ),
3122 Self::Spawn {
3123 program,
3124 source,
3125 cgroup_path: None,
3126 } => write!(
3127 f,
3128 "failed to spawn module '{}': {source}",
3129 program.display()
3130 ),
3131 Self::Cgroup { module_id, source } => {
3132 write!(
3133 f,
3134 "failed to prepare cgroup for module '{module_id}': {source}"
3135 )
3136 }
3137 Self::LaunchNonce { reason } => {
3138 write!(
3139 f,
3140 "failed to generate reserved-module launch nonce: {reason}"
3141 )
3142 }
3143 Self::Wait { module_id, source } => {
3144 write!(f, "failed to wait for module '{module_id}': {source}")
3145 }
3146 Self::Kill { module_id, source } => {
3147 write!(f, "failed to kill module '{module_id}': {source}")
3148 }
3149 Self::Forwarding(err) => write!(f, "forwarding error: {err}"),
3150 Self::Registry(err) => write!(f, "registry error: {err}"),
3151 Self::ReloadUnavailable { module_id, reason } => {
3152 write!(f, "reload unavailable for module '{module_id}': {reason}")
3153 }
3154 Self::Disabled { module_id } => {
3155 write!(
3156 f,
3157 "module '{module_id}' is disabled; enable it before restart or reload"
3158 )
3159 }
3160 Self::ReloadFailed { module_id, reason } => {
3161 write!(f, "reload failed for module '{module_id}': {reason}")
3162 }
3163 Self::RegistrationStillActive { module_id, waited } => write!(
3164 f,
3165 "module '{module_id}' registration remained active after waiting {waited:?}"
3166 ),
3167 Self::StatePoisoned { module_id } => match module_id {
3168 Some(module_id) => {
3169 write!(f, "supervisor state for module '{module_id}' was poisoned")
3170 }
3171 None => write!(f, "supervisor state was poisoned"),
3172 },
3173 Self::CommandClosed { module_id } => {
3174 write!(
3175 f,
3176 "supervisor command channel for module '{module_id}' is closed"
3177 )
3178 }
3179 Self::SwapInProgress { module_id } => write!(
3180 f,
3181 "module '{module_id}' is being swapped; retry once the swap has cut over or failed, or stop the module to abort the swap"
3182 ),
3183 Self::SwapRefused { module_id, reason } => match reason {
3184 SwapRefusal::OverlapExclusive => write!(
3185 f,
3186 "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"
3187 ),
3188 SwapRefusal::NotRegistered => write!(
3189 f,
3190 "module '{module_id}' is not registered, so there is no serving process to keep while a replacement warms; use a plain restart"
3191 ),
3192 SwapRefusal::ProtocolNone => write!(
3193 f,
3194 "module '{module_id}' is protocol: \"none\" and never registers, so a swap could never see its replacement become ready; use a plain restart"
3195 ),
3196 SwapRefusal::NotConfigured => write!(
3197 f,
3198 "module '{module_id}' cannot be swapped: the supervisor was built without the forwarding table or shared handle a swap needs"
3199 ),
3200 SwapRefusal::AlreadySwapping => {
3201 write!(f, "module '{module_id}' is already being swapped")
3202 }
3203 },
3204 Self::SwapFailed {
3205 module_id,
3206 arm,
3207 detail,
3208 ..
3209 } => write!(
3210 f,
3211 "swap of module '{module_id}' failed ({}): {detail}; the running process was left serving",
3212 arm.as_str()
3213 ),
3214 }
3215 }
3216}
3217
3218impl Error for SuperviseError {
3219 fn source(&self) -> Option<&(dyn Error + 'static)> {
3220 match self {
3221 Self::Spawn { source, .. }
3222 | Self::Cgroup { source, .. }
3223 | Self::Wait { source, .. }
3224 | Self::Kill { source, .. } => Some(source),
3225 Self::Forwarding(err) => Some(err),
3226 Self::Registry(err) => Some(err),
3227 Self::LaunchNonce { .. }
3228 | Self::InvalidSpec { .. }
3229 | Self::ReloadUnavailable { .. }
3230 | Self::Disabled { .. }
3231 | Self::ReloadFailed { .. }
3232 | Self::RegistrationStillActive { .. }
3233 | Self::StatePoisoned { .. }
3234 | Self::CommandClosed { .. }
3235 | Self::SwapInProgress { .. }
3236 | Self::SwapRefused { .. }
3237 | Self::SwapFailed { .. } => None,
3238 }
3239 }
3240}
3241
3242pub(crate) fn validate_spec(spec: &ModuleSpec) -> Result<(), SuperviseError> {
3243 if spec.module_id.trim().is_empty() {
3244 return Err(SuperviseError::InvalidSpec {
3245 reason: "module_id must not be empty".to_string(),
3246 });
3247 }
3248
3249 Ok(())
3250}
3251
3252#[derive(Debug, Default)]
3253struct HealthProbeRuntime {
3254 registered_connection: Option<crate::ConnectionId>,
3255 advertised: bool,
3256 next_probe_at: Option<Instant>,
3257 probe_index: u64,
3258}
3259
3260impl HealthProbeRuntime {
3261 fn refresh_registration(
3262 &mut self,
3263 spec: &ModuleSpec,
3264 runtime: &SupervisorRuntimeConfig,
3265 registry: &Registry,
3266 snapshot: &SharedSnapshot,
3267 ) {
3268 if spec.protocol == ModuleProtocol::None {
3280 self.registered_connection = None;
3281 self.advertised = false;
3282 self.next_probe_at = None;
3283 return;
3284 }
3285
3286 let registration = match registry.get_module(&spec.module_id) {
3287 Ok(registration) => registration,
3288 Err(err) => {
3289 warn!(module_id = %spec.module_id, error = %err, "health prober could not read registry");
3290 self.advertised = false;
3291 self.next_probe_at = None;
3292 return;
3293 }
3294 };
3295
3296 let Some(registration) = registration else {
3297 self.registered_connection = None;
3298 self.advertised = false;
3299 self.next_probe_at = None;
3300 return;
3301 };
3302
3303 let advertised = registration
3304 .control_ops
3305 .iter()
3306 .any(|op| op == MODULE_CONTROL_OP_HEALTH_CHECK);
3307 if !advertised {
3308 self.registered_connection = Some(registration.connection_id);
3309 self.advertised = false;
3310 self.next_probe_at = None;
3311 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3312 state.health.status = SupervisorHealthStatus::Unknown;
3313 state.health.consecutive_failures = 0;
3314 state.health.last_probe_ms = None;
3315 state.health.detail = None;
3316 state.health.metrics = None;
3317 });
3318 return;
3319 }
3320
3321 let reregistered = self.registered_connection != Some(registration.connection_id);
3322 self.registered_connection = Some(registration.connection_id);
3323 self.advertised = true;
3324 if reregistered || self.next_probe_at.is_none() {
3325 self.probe_index = 0;
3326 self.next_probe_at = Some(
3327 Instant::now() + jittered_health_delay(&spec.module_id, 0, runtime.health.cadence),
3328 );
3329 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3330 state.health.status = SupervisorHealthStatus::Unknown;
3331 state.health.consecutive_failures = 0;
3332 state.health.detail = None;
3333 state.health.metrics = None;
3334 });
3335 }
3336 }
3337
3338 fn wake_after(&self) -> Duration {
3339 if !self.advertised {
3340 return REGISTRY_RELEASE_POLL;
3341 }
3342 self.next_probe_at
3343 .map(|next| next.saturating_duration_since(Instant::now()))
3344 .unwrap_or(REGISTRY_RELEASE_POLL)
3345 }
3346
3347 fn due(&self) -> bool {
3348 self.advertised
3349 && self
3350 .next_probe_at
3351 .is_some_and(|next| Instant::now() >= next)
3352 }
3353
3354 fn schedule_next(&mut self, spec: &ModuleSpec, cadence: Duration) {
3355 self.probe_index = self.probe_index.wrapping_add(1);
3356 self.next_probe_at = Some(
3357 Instant::now() + jittered_health_delay(&spec.module_id, self.probe_index, cadence),
3358 );
3359 }
3360}
3361
3362#[derive(Debug)]
3397enum HealthProbeEvidence {
3398 LaneDead,
3400 NoAnswer,
3402 BadAnswer,
3404 Misconfigured,
3406}
3407
3408#[derive(Debug)]
3409struct HealthProbeError {
3410 evidence: HealthProbeEvidence,
3411 message: String,
3412}
3413
3414impl HealthProbeError {
3415 fn lane_dead(message: impl Into<String>) -> Self {
3416 Self::with(HealthProbeEvidence::LaneDead, message)
3417 }
3418
3419 fn no_answer(message: impl Into<String>) -> Self {
3420 Self::with(HealthProbeEvidence::NoAnswer, message)
3421 }
3422
3423 fn bad_answer(message: impl Into<String>) -> Self {
3424 Self::with(HealthProbeEvidence::BadAnswer, message)
3425 }
3426
3427 fn misconfigured(message: impl Into<String>) -> Self {
3428 Self::with(HealthProbeEvidence::Misconfigured, message)
3429 }
3430
3431 fn with(evidence: HealthProbeEvidence, message: impl Into<String>) -> Self {
3432 Self {
3433 evidence,
3434 message: message.into(),
3435 }
3436 }
3437
3438 #[allow(dead_code)]
3452 fn is_proof_of_death(&self) -> bool {
3453 matches!(self.evidence, HealthProbeEvidence::LaneDead)
3454 }
3455
3456 fn label(&self) -> &'static str {
3464 match self.evidence {
3465 HealthProbeEvidence::LaneDead => "lane-dead",
3466 HealthProbeEvidence::NoAnswer => "no-answer",
3467 HealthProbeEvidence::BadAnswer => "bad-answer",
3468 HealthProbeEvidence::Misconfigured => "daemon-misconfigured",
3469 }
3470 }
3471}
3472
3473impl fmt::Display for HealthProbeError {
3474 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3475 f.write_str(&self.message)
3476 }
3477}
3478
3479async fn run_health_probe_cycle(
3480 spec: &ModuleSpec,
3481 runtime: &SupervisorRuntimeConfig,
3482 registry: &Registry,
3483 process_liveness: &SupervisorProcessLiveness,
3484 snapshot: &SharedSnapshot,
3485 child: &mut Option<SupervisedChild>,
3486) {
3487 let now_ms = unix_ms_now();
3488 match probe_module_health(&spec.module_id, runtime, None).await {
3489 Ok(report) => {
3490 handle_health_report(
3491 spec,
3492 runtime,
3493 registry,
3494 process_liveness,
3495 snapshot,
3496 child,
3497 report,
3498 now_ms,
3499 )
3500 .await;
3501 }
3502 Err(err) => {
3503 handle_health_probe_failure(
3504 spec,
3505 runtime,
3506 registry,
3507 process_liveness,
3508 snapshot,
3509 child,
3510 err,
3511 now_ms,
3512 )
3513 .await;
3514 }
3515 }
3516}
3517
3518async fn probe_module_health(
3519 module_id: &str,
3520 runtime: &SupervisorRuntimeConfig,
3521 drain_deadline: Option<Instant>,
3522) -> Result<HealthReport, HealthProbeError> {
3523 let Some(forwarding) = runtime.forwarding.as_ref() else {
3524 return Err(HealthProbeError::misconfigured(
3525 "supervisor was not configured with a forwarding table",
3526 ));
3527 };
3528 let probe_started_at = Instant::now();
3529 let mut deadline = probe_started_at + runtime.health.deadline;
3530 if let Some(drain_deadline) = drain_deadline {
3531 deadline = deadline.min(drain_deadline);
3532 }
3533 let pending = if drain_deadline.is_some() {
3534 forwarding.begin_drain_health_probe_rpc_for(
3535 module_id,
3536 MODULE_CONTROL_OP_HEALTH_CHECK,
3537 probe_started_at,
3538 deadline,
3539 )
3540 } else {
3541 forwarding.begin_health_probe_rpc_for(
3542 module_id,
3543 MODULE_CONTROL_OP_HEALTH_CHECK,
3544 probe_started_at,
3545 deadline,
3546 )
3547 }
3548 .map_err(|err| {
3549 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3552 })?;
3553 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3554}
3555
3556async fn probe_endpoint_health(
3563 endpoint: crate::ModuleEndpointId,
3564 runtime: &SupervisorRuntimeConfig,
3565 deadline_cap: Option<Instant>,
3566) -> Result<HealthReport, HealthProbeError> {
3567 let Some(forwarding) = runtime.forwarding.as_ref() else {
3568 return Err(HealthProbeError::misconfigured(
3569 "supervisor was not configured with a forwarding table",
3570 ));
3571 };
3572 let probe_started_at = Instant::now();
3573 let mut deadline = probe_started_at + runtime.health.deadline;
3574 if let Some(cap) = deadline_cap {
3575 deadline = deadline.min(cap);
3576 }
3577 let pending = forwarding
3578 .begin_endpoint_health_probe_rpc_for(
3579 endpoint,
3580 MODULE_CONTROL_OP_HEALTH_CHECK,
3581 probe_started_at,
3582 deadline,
3583 )
3584 .map_err(|err| {
3585 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3586 })?;
3587 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3588}
3589
3590async fn await_health_probe(
3592 forwarding: &ForwardingTable,
3593 pending: PendingModuleControlRpc,
3594 deadline: Instant,
3595 probe_budget: Duration,
3596) -> Result<HealthReport, HealthProbeError> {
3597 let PendingModuleControlRpc {
3598 endpoint,
3599 module_sink,
3600 negotiated_ver,
3601 corr,
3602 receiver,
3603 } = pending;
3604 let body = serde_json::to_vec(&ModuleControlRequest::HealthCheck {}).map_err(|err| {
3605 HealthProbeError::misconfigured(format!("failed to encode health.check: {err}"))
3606 })?;
3607 let frame = Frame::build_with_version(
3608 negotiated_ver,
3609 FrameType::Request,
3610 control_flags(),
3611 0,
3612 0,
3613 corr,
3614 body,
3615 )
3616 .map_err(|err| {
3617 HealthProbeError::misconfigured(format!("failed to build health.check frame: {err}"))
3618 })?;
3619
3620 match timeout_at(deadline, module_sink.send(frame)).await {
3626 Ok(Ok(())) => {}
3627 Ok(Err(err)) => {
3628 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3629 return Err(HealthProbeError::lane_dead(format!(
3632 "failed to send health.check: {err}"
3633 )));
3634 }
3635 Err(_elapsed) => {
3636 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3637 return Err(HealthProbeError::no_answer(
3641 "health.check send timed out before enqueue (module egress full)",
3642 ));
3643 }
3644 }
3645
3646 match timeout_at(deadline, receiver).await {
3647 Ok(Ok(ModuleControlRpcOutcome::Response(response))) => {
3651 response.health_report().ok_or_else(|| {
3652 HealthProbeError::bad_answer("health.check RPC returned a non-health response")
3653 })
3654 }
3655 Ok(Ok(ModuleControlRpcOutcome::Rejected(body))) => Err(HealthProbeError::bad_answer(
3656 format!("health.check rejected: {}", body.message),
3657 )),
3658 Ok(Ok(ModuleControlRpcOutcome::ModuleGone(message))) => {
3659 Err(HealthProbeError::lane_dead(message))
3660 }
3661 Ok(Ok(ModuleControlRpcOutcome::MalformedResponse(message))) => {
3662 Err(HealthProbeError::bad_answer(message))
3663 }
3664 Ok(Ok(ModuleControlRpcOutcome::UnexpectedOp { expected, actual })) => {
3665 Err(HealthProbeError::bad_answer(format!(
3666 "expected module-control op '{expected}', got '{actual}'"
3667 )))
3668 }
3669 Ok(Ok(ModuleControlRpcOutcome::DeadlineElapsed)) => Err(HealthProbeError::bad_answer(
3673 "module answered health.check after its daemon deadline",
3674 )),
3675 Ok(Err(_)) => Err(HealthProbeError::misconfigured(
3676 "health.check waiter was canceled before the module responded",
3677 )),
3678 Err(_) => {
3679 let _ = forwarding.tombstone_health_probe_rpc(endpoint, corr);
3680 Err(HealthProbeError::no_answer(format!(
3681 "module did not answer health.check within {probe_budget:?}"
3682 )))
3683 }
3684 }
3685}
3686
3687#[allow(clippy::too_many_arguments)]
3688async fn handle_health_report(
3689 spec: &ModuleSpec,
3690 runtime: &SupervisorRuntimeConfig,
3691 registry: &Registry,
3692 process_liveness: &SupervisorProcessLiveness,
3693 snapshot: &SharedSnapshot,
3694 child: &mut Option<SupervisedChild>,
3695 report: HealthReport,
3696 now_ms: u64,
3697) {
3698 let status = supervisor_health_status(report.status);
3699 let detail = report.detail.clone();
3700 let metrics = truncate_health_metrics(report.metrics);
3701 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3702 state.health.status = status;
3703 state.health.last_probe_ms = Some(now_ms);
3704 state.health.detail = detail.clone();
3705 state.health.metrics = metrics.clone();
3706 state.health.consecutive_failures = 0;
3707 });
3708
3709 let action = match report.status {
3710 HealthStatus::Ok => return,
3711 HealthStatus::Degraded => runtime.health.on_degraded,
3712 HealthStatus::Failing => runtime.health.on_failing,
3713 };
3714 apply_l3_health_action(
3715 spec,
3716 runtime,
3717 registry,
3718 process_liveness,
3719 snapshot,
3720 child,
3721 status,
3722 detail.as_deref(),
3723 action,
3724 now_ms,
3725 )
3726 .await;
3727}
3728
3729#[allow(clippy::too_many_arguments)]
3730async fn handle_health_probe_failure(
3731 spec: &ModuleSpec,
3732 runtime: &SupervisorRuntimeConfig,
3733 registry: &Registry,
3734 process_liveness: &SupervisorProcessLiveness,
3735 snapshot: &SharedSnapshot,
3736 child: &mut Option<SupervisedChild>,
3737 err: HealthProbeError,
3738 now_ms: u64,
3739) {
3740 let threshold = runtime.health.failure_threshold.max(1);
3741 let mut failures = 0;
3742 let detail = format!("[{}] {err}", err.label());
3747 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3748 state.health.last_probe_ms = Some(now_ms);
3749 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
3750 state.health.detail = Some(detail.clone());
3751 state.health.metrics = None;
3752 failures = state.health.consecutive_failures;
3753 });
3754
3755 if failures < threshold {
3756 warn!(
3757 module_id = %spec.module_id,
3758 consecutive_failures = failures,
3759 threshold,
3760 evidence = err.label(),
3761 detail = %detail,
3762 "health.check probe failed"
3763 );
3764 return;
3765 }
3766
3767 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3768 state.state = ModuleState::Unresponsive;
3769 state.health.status = SupervisorHealthStatus::Unresponsive;
3770 });
3771 if runtime.health.critical {
3775 error!(
3776 module_id = %spec.module_id,
3777 status = "unresponsive",
3778 evidence = err.label(),
3779 detail = %detail,
3780 "critical module health alert"
3781 );
3782 } else {
3783 warn!(
3784 module_id = %spec.module_id,
3785 status = "unresponsive",
3786 evidence = err.label(),
3787 detail = %detail,
3788 "module health threshold breached"
3789 );
3790 }
3791 if let Err(err) = health_restart_child(
3792 spec,
3793 runtime,
3794 registry,
3795 process_liveness,
3796 snapshot,
3797 child,
3798 SupervisorHealthStatus::Unresponsive,
3799 Some(&detail),
3800 now_ms,
3801 )
3802 .await
3803 {
3804 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3805 }
3806}
3807
3808#[allow(clippy::too_many_arguments)]
3809async fn apply_l3_health_action(
3810 spec: &ModuleSpec,
3811 runtime: &SupervisorRuntimeConfig,
3812 registry: &Registry,
3813 process_liveness: &SupervisorProcessLiveness,
3814 snapshot: &SharedSnapshot,
3815 child: &mut Option<SupervisedChild>,
3816 status: SupervisorHealthStatus,
3817 detail: Option<&str>,
3818 action: HealthAction,
3819 now_ms: u64,
3820) {
3821 record_health_action(snapshot, &spec.module_id, action.to_string(), now_ms);
3822 match action {
3823 HealthAction::Report => {
3824 info!(
3825 module_id = %spec.module_id,
3826 status = ?status,
3827 detail,
3828 "module reported non-ok health"
3829 );
3830 }
3831 HealthAction::Alert => {
3832 error!(
3833 module_id = %spec.module_id,
3834 status = ?status,
3835 detail,
3836 "module health alert"
3837 );
3838 }
3839 HealthAction::Restart => {
3840 if let Err(err) = health_restart_child(
3841 spec,
3842 runtime,
3843 registry,
3844 process_liveness,
3845 snapshot,
3846 child,
3847 status,
3848 detail,
3849 now_ms,
3850 )
3851 .await
3852 {
3853 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3854 }
3855 }
3856 }
3857}
3858
3859#[allow(clippy::too_many_arguments)]
3860async fn health_restart_child(
3861 spec: &ModuleSpec,
3862 runtime: &SupervisorRuntimeConfig,
3863 registry: &Registry,
3864 process_liveness: &SupervisorProcessLiveness,
3865 snapshot: &SharedSnapshot,
3866 child: &mut Option<SupervisedChild>,
3867 status: SupervisorHealthStatus,
3868 detail: Option<&str>,
3869 now_ms: u64,
3870) -> Result<(), SuperviseError> {
3871 let (enabled, schedule) = {
3872 let mut state = lock_snapshot(snapshot)?;
3873 let enabled = state.enabled;
3874 let schedule = if enabled {
3875 state.next_crash_restart(&runtime.restart_policy, Instant::now())
3876 } else {
3877 None
3878 };
3879 (enabled, schedule)
3880 };
3881
3882 if !enabled {
3883 return Err(SuperviseError::Disabled {
3884 module_id: spec.module_id.clone(),
3885 });
3886 }
3887
3888 if schedule.is_none() {
3889 record_health_action(snapshot, &spec.module_id, "disabled".to_string(), now_ms);
3890 error!(
3891 module_id = %spec.module_id,
3892 status = ?status,
3893 detail,
3894 max_restarts = runtime.restart_policy.max_restarts,
3895 window_secs = runtime.restart_policy.window.as_secs(),
3896 "health restart budget exhausted; disabling module"
3897 );
3898 begin_forwarding_drain_if_configured(
3899 spec,
3900 runtime,
3901 registry,
3902 snapshot,
3903 Some(false),
3904 RouteCloseReason::Disable,
3905 )
3906 .await?;
3907 drain_optional_child(
3908 &spec.module_id,
3909 spec.protocol,
3910 registry,
3911 snapshot,
3912 &runtime.terminal_ring,
3913 &runtime.spawn_events,
3914 child,
3915 runtime.drain_timeout,
3916 ModuleState::Disabled,
3917 Some(false),
3918 )
3919 .await?;
3920 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3921 return Ok(());
3922 }
3923
3924 let schedule = schedule.expect("a health restart must have a crash-restart schedule");
3925 let mut restart_count = 0;
3926 update_snapshot(snapshot, Some(&spec.module_id), |state| {
3927 restart_count = state.crash_restarts.len();
3928 state.state = ModuleState::Unresponsive;
3929 state.health.status = status;
3930 state.health.last_action = Some(HealthAction::Restart.to_string());
3931 state.health.last_action_ms = Some(now_ms);
3932 })?;
3933 warn!(
3934 module_id = %spec.module_id,
3935 status = ?status,
3936 detail,
3937 restart_count,
3938 restart_in_window = schedule.restart_in_window,
3939 delay_ms = schedule.delay.as_millis() as u64,
3940 "health-triggered module restart"
3941 );
3942
3943 begin_forwarding_drain_if_configured(
3944 spec,
3945 runtime,
3946 registry,
3947 snapshot,
3948 Some(true),
3949 RouteCloseReason::Restart,
3950 )
3951 .await?;
3952 drain_optional_child(
3953 &spec.module_id,
3954 spec.protocol,
3955 registry,
3956 snapshot,
3957 &runtime.terminal_ring,
3958 &runtime.spawn_events,
3959 child,
3960 runtime.drain_timeout,
3961 ModuleState::Restarting,
3962 Some(true),
3963 )
3964 .await?;
3965 sleep(schedule.delay).await;
3966 if !respawn_still_pending(snapshot) {
3970 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3971 return Ok(());
3972 }
3973 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
3974 match spawn_and_mark_running(spec, runtime, snapshot) {
3975 Ok(next_child) => {
3976 *child = Some(next_child);
3977 Ok(())
3978 }
3979 Err(err) => {
3980 fail_snapshot(snapshot, Some(&spec.module_id), None);
3981 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3982 *child = None;
3983 Err(err)
3984 }
3985 }
3986}
3987
3988fn record_health_action(snapshot: &SharedSnapshot, module_id: &str, action: String, now_ms: u64) {
3989 let _ = update_snapshot(snapshot, Some(module_id), |state| {
3990 state.health.last_action = Some(action);
3991 state.health.last_action_ms = Some(now_ms);
3992 });
3993}
3994
3995fn supervisor_health_status(status: HealthStatus) -> SupervisorHealthStatus {
3996 match status {
3997 HealthStatus::Ok => SupervisorHealthStatus::Ok,
3998 HealthStatus::Degraded => SupervisorHealthStatus::Degraded,
3999 HealthStatus::Failing => SupervisorHealthStatus::Failing,
4000 }
4001}
4002
4003fn truncate_health_metrics(metrics: Option<Value>) -> Option<Value> {
4015 let metrics = metrics?;
4016 match serde_json::to_vec(&metrics) {
4017 Ok(encoded) if encoded.len() > MAX_HEALTH_METRICS_BYTES => Some(serde_json::json!({
4018 "truncated": true,
4019 "original_bytes": encoded.len(),
4020 })),
4021 Ok(_) | Err(_) => Some(metrics),
4022 }
4023}
4024
4025fn jittered_health_delay(module_id: &str, probe_index: u64, cadence: Duration) -> Duration {
4031 if cadence.is_zero() {
4032 return Duration::ZERO;
4033 }
4034 let cadence_ms = cadence.as_millis() as u64;
4035 if cadence_ms == 0 {
4051 return cadence;
4052 }
4053 let jitter_span = (cadence_ms / 10).max(1);
4068 let hash = module_id.as_bytes().iter().fold(
4069 probe_index.wrapping_mul(0x9E37_79B9_7F4A_7C15),
4070 |acc, byte| {
4071 acc.wrapping_mul(1099511628211)
4072 .wrapping_add(u64::from(*byte))
4073 },
4074 );
4075 cadence + Duration::from_millis(hash % jitter_span)
4076}
4077
4078#[cfg(test)]
4079mod tests {
4080 use super::*;
4081
4082 #[test]
4083 fn readding_a_module_clears_its_rescan_removal_tombstone() {
4084 let handle = SupervisorHandle::new();
4085 let module_id = "readded-tombstone";
4086 handle.record_rescan_removal(module_id);
4087 assert!(handle.removal_tombstone_age_ms(module_id).is_some());
4088
4089 handle.apply_identity_configuration(&ModuleSpec {
4090 module_id: module_id.to_string(),
4091 program: PathBuf::from("/test/module"),
4092 args: Vec::new(),
4093 env: Vec::new(),
4094 reserved: false,
4095 reserved_prefixes: Vec::new(),
4096 protocol: ModuleProtocol::Subc,
4097 overlap: Default::default(),
4098 });
4099
4100 assert!(
4101 handle.removal_tombstone_age_ms(module_id).is_none(),
4102 "a re-added module must not retain a stale removal tombstone"
4103 );
4104 }
4105
4106 fn stale_process_snapshot(state: ModuleState, enabled: bool) -> SharedSnapshot {
4107 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::new(state, enabled)));
4108 update_snapshot(&snapshot, Some("stale-process-facts"), |snapshot| {
4109 snapshot.process_alive = true;
4110 snapshot.pid = Some(41);
4111 snapshot.spawned_at_ms = Some(42);
4112 snapshot.spawned_from = Some(PathBuf::from("/spawned/module"));
4113 snapshot.spawned_file_identity = Some(SpawnedFileIdentity {
4114 device: 43,
4115 inode: 44,
4116 });
4117 })
4118 .unwrap();
4119 snapshot
4120 }
4121
4122 fn assert_snapshot_process_facts_cleared(snapshot: &SharedSnapshot) {
4123 let snapshot = lock_snapshot(snapshot).unwrap();
4124 assert!(!snapshot.process_alive);
4125 assert_eq!(snapshot.pid, None);
4126 assert_eq!(snapshot.spawned_at_ms, None);
4127 assert_eq!(snapshot.spawned_from, None);
4128 assert_eq!(snapshot.spawned_file_identity, None);
4129 }
4130
4131 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4132 async fn failed_enable_spawn_clears_preexisting_current_process_facts() {
4133 let supervisor = Supervisor::default();
4134 let mut runtime = supervisor.runtime_config();
4135 runtime.test_seed_stale_facts_before_enable_spawn = true;
4136 let snapshot = stale_process_snapshot(ModuleState::Disabled, false);
4137 let mut child = None;
4138 let spec = ModuleSpec {
4139 module_id: "failed-enable-clears-facts".to_string(),
4140 program: PathBuf::from("/definitely/missing/failed-enable-module"),
4141 args: Vec::new(),
4142 env: Vec::new(),
4143 reserved: false,
4144 reserved_prefixes: Vec::new(),
4145 protocol: ModuleProtocol::Subc,
4146 overlap: Default::default(),
4147 };
4148
4149 let result = set_child_enabled(
4150 &spec,
4151 &runtime,
4152 &supervisor.registry,
4153 &supervisor.process_liveness,
4154 &snapshot,
4155 &mut child,
4156 true,
4157 )
4158 .await;
4159
4160 assert!(matches!(result, Err(SuperviseError::Spawn { .. })));
4161 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4162 assert_snapshot_process_facts_cleared(&snapshot);
4163 }
4164
4165 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4166 async fn failed_reload_spawn_clears_current_process_facts() {
4167 let supervisor = Supervisor::default();
4168 let mut runtime = supervisor.runtime_config();
4169 runtime.restart_policy = RestartPolicy::new(0, Duration::ZERO);
4170 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4171 let mut child = None;
4172 let spec = ModuleSpec {
4173 module_id: "failed-reload-clears-facts".to_string(),
4174 program: PathBuf::from("/unused/failed-reload-module"),
4175 args: Vec::new(),
4176 env: Vec::new(),
4177 reserved: false,
4178 reserved_prefixes: Vec::new(),
4179 protocol: ModuleProtocol::Subc,
4180 overlap: Default::default(),
4181 };
4182
4183 let result = handle_reload_spawn_failure(
4184 &spec,
4185 &runtime,
4186 &supervisor.process_liveness,
4187 &snapshot,
4188 &mut child,
4189 "forced reload spawn failure".to_string(),
4190 )
4191 .await;
4192
4193 assert!(matches!(result, Err(SuperviseError::ReloadFailed { .. })));
4194 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4195 assert_snapshot_process_facts_cleared(&snapshot);
4196 }
4197
4198 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4199 async fn dropping_a_module_with_an_active_monitor_clears_current_process_facts() {
4200 let supervisor = Supervisor::default();
4201 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4202 let module = supervisor.supervised_module(
4203 ModuleSpec {
4204 module_id: "drop-clears-facts".to_string(),
4205 program: PathBuf::from("/unused/drop-module"),
4206 args: Vec::new(),
4207 env: Vec::new(),
4208 reserved: false,
4209 reserved_prefixes: Vec::new(),
4210 protocol: ModuleProtocol::Subc,
4211 overlap: Default::default(),
4212 },
4213 supervisor.runtime_config(),
4214 Arc::clone(&snapshot),
4215 None,
4216 );
4217 assert!(!module
4218 .inner
4219 .monitor
4220 .lock()
4221 .unwrap()
4222 .as_ref()
4223 .unwrap()
4224 .is_finished());
4225
4226 drop(module);
4227
4228 assert_eq!(
4229 lock_snapshot(&snapshot).unwrap().state,
4230 ModuleState::Stopped
4231 );
4232 assert_snapshot_process_facts_cleared(&snapshot);
4233 }
4234
4235 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4236 async fn configuration_update_does_not_replace_captured_running_process_facts() {
4237 let supervisor = Supervisor::default();
4238 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4239 let initial = ModuleSpec {
4240 module_id: "rescan-preserves-spawn-facts".to_string(),
4241 program: PathBuf::from("/spawned/module"),
4242 args: Vec::new(),
4243 env: Vec::new(),
4244 reserved: false,
4245 reserved_prefixes: Vec::new(),
4246 protocol: ModuleProtocol::Subc,
4247 overlap: Default::default(),
4248 };
4249 let module = supervisor.supervised_module(
4250 initial.clone(),
4251 supervisor.runtime_config(),
4252 snapshot,
4253 None,
4254 );
4255 let before = module.status().unwrap();
4256 let mut replacement = initial;
4257 replacement.program = PathBuf::from("/rescanned/replacement-module");
4258
4259 module
4260 .update_configuration(replacement, HealthConfig::default(), None)
4261 .await
4262 .unwrap();
4263
4264 let after = module.status().unwrap();
4265 assert_eq!(after.pid, before.pid);
4266 assert_eq!(after.spawned_at_ms, before.spawned_at_ms);
4267 assert_eq!(after.spawned_from, before.spawned_from);
4268 drop(module);
4269 }
4270}
4271
4272fn unix_ms_now() -> u64 {
4273 SystemTime::now()
4274 .duration_since(UNIX_EPOCH)
4275 .map(|duration| duration.as_millis().min(u128::from(u64::MAX)) as u64)
4276 .unwrap_or(0)
4277}
4278
4279async fn supervise_loop(
4280 mut spec: ModuleSpec,
4281 mut runtime: SupervisorRuntimeConfig,
4282 registry: Arc<Registry>,
4283 process_liveness: Arc<SupervisorProcessLiveness>,
4284 snapshot: SharedSnapshot,
4285 mut child: Option<SupervisedChild>,
4286 mut commands: mpsc::Receiver<SupervisorCommand>,
4287) {
4288 let mut health_probe = HealthProbeRuntime::default();
4289 let mut pending_respawn: Option<Instant> = None;
4293 let mut requeued: VecDeque<SupervisorCommand> = VecDeque::new();
4296 loop {
4297 if let Some(command) = requeued.pop_front() {
4298 if !handle_supervisor_command(
4299 command,
4300 &mut spec,
4301 &mut runtime,
4302 ®istry,
4303 &process_liveness,
4304 &snapshot,
4305 &mut child,
4306 &mut commands,
4307 &mut requeued,
4308 )
4309 .await
4310 {
4311 return;
4312 }
4313 if child.is_some() || !respawn_still_pending(&snapshot) {
4314 pending_respawn = None;
4315 }
4316 continue;
4317 }
4318 if child.is_some() {
4319 health_probe.refresh_registration(&spec, &runtime, ®istry, &snapshot);
4320 let probe_sleep = sleep(health_probe.wake_after());
4321 tokio::pin!(probe_sleep);
4322 let active_child = child.as_mut().expect("child checked above");
4323 tokio::select! {
4324 wait_result = active_child.wait() => {
4325 let exit_report = match wait_result {
4334 Ok(status) => classify_reaped_child_exit(&snapshot, active_child, &status),
4335 Err(err) => {
4336 active_child.drain_stderr(&spec.module_id).await;
4337 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4338 record_wait_error_terminal(
4344 &spec.module_id,
4345 &runtime.terminal_ring,
4346 &runtime.spawn_events,
4347 );
4348 untrack_if_registration_released(
4349 &process_liveness,
4350 ®istry,
4351 &spec.module_id,
4352 &snapshot,
4353 );
4354 error!(module_id = %spec.module_id, error = %err, "failed to wait for supervised module");
4355 child = None;
4356 continue;
4357 }
4358 };
4359 active_child.drain_stderr(&spec.module_id).await;
4360
4361 match on_child_exit(
4362 &spec,
4363 runtime.restart_policy,
4364 ®istry,
4365 &snapshot,
4366 &runtime.terminal_ring,
4367 &runtime.spawn_events,
4368 &runtime.child_roster,
4369 exit_report,
4370 ).await {
4371 NextAction::Stop { registration_released } => {
4372 if registration_released {
4373 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4374 }
4375 child = None;
4376 }
4377 NextAction::Restart { schedule } => {
4378 let delay = schedule.map_or(
4379 runtime.restart_policy.delay_for_restart(0),
4380 |schedule| schedule.delay,
4381 );
4382 if let Some(schedule) = schedule {
4383 log_crash_respawn(&spec.module_id, schedule);
4384 }
4385 child = None;
4393 pending_respawn = Some(Instant::now() + delay);
4394 }
4395 }
4396 }
4397 command = commands.recv() => {
4398 let Some(command) = command else {
4399 return;
4400 };
4401 if !handle_supervisor_command(
4402 command,
4403 &mut spec,
4404 &mut runtime,
4405 ®istry,
4406 &process_liveness,
4407 &snapshot,
4408 &mut child,
4409 &mut commands,
4410 &mut requeued,
4411 ).await {
4412 return;
4413 }
4414 }
4415 _ = &mut probe_sleep => {
4416 if health_probe.due() {
4417 run_health_probe_cycle(
4418 &spec,
4419 &runtime,
4420 ®istry,
4421 &process_liveness,
4422 &snapshot,
4423 &mut child,
4424 ).await;
4425 if child.is_some() {
4426 health_probe.schedule_next(&spec, runtime.health.cadence);
4427 }
4428 }
4429 }
4430 }
4431 } else if let Some(deadline) = pending_respawn {
4432 tokio::select! {
4433 _ = sleep_until(deadline) => {
4434 pending_respawn = None;
4435 if !respawn_still_pending(&snapshot) {
4439 continue;
4440 }
4441 if runtime.child_roster.is_closed() {
4446 let _ = update_snapshot(&snapshot, Some(&spec.module_id), |state| {
4447 state.state = ModuleState::Stopped;
4448 });
4449 debug!(module_id = %spec.module_id, "crash respawn cancelled by daemon shutdown");
4450 continue;
4451 }
4452 if let Err(err) = wait_for_registration_release(
4453 ®istry,
4454 &spec.module_id,
4455 REGISTRY_RELEASE_TIMEOUT,
4456 ).await {
4457 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4458 error!(module_id = %spec.module_id, error = %err, "registration did not release before restart");
4459 continue;
4460 }
4461
4462 match spawn_and_mark_running(&spec, &runtime, &snapshot) {
4463 Ok(next_child) => {
4464 child = Some(next_child);
4465 debug!(module_id = %spec.module_id, "supervised module restarted after crash");
4466 }
4467 Err(err) => {
4468 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4469 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4470 error!(module_id = %spec.module_id, error = %err, "failed to restart supervised module");
4471 }
4472 }
4473 }
4474 command = commands.recv() => {
4475 let Some(command) = command else {
4476 return;
4477 };
4478 if !handle_supervisor_command(
4479 command,
4480 &mut spec,
4481 &mut runtime,
4482 ®istry,
4483 &process_liveness,
4484 &snapshot,
4485 &mut child,
4486 &mut commands,
4487 &mut requeued,
4488 ).await {
4489 return;
4490 }
4491 if child.is_some() || !respawn_still_pending(&snapshot) {
4496 pending_respawn = None;
4497 }
4498 }
4499 }
4500 } else {
4501 let Some(command) = commands.recv().await else {
4502 return;
4503 };
4504 if !handle_supervisor_command(
4505 command,
4506 &mut spec,
4507 &mut runtime,
4508 ®istry,
4509 &process_liveness,
4510 &snapshot,
4511 &mut child,
4512 &mut commands,
4513 &mut requeued,
4514 )
4515 .await
4516 {
4517 return;
4518 }
4519 }
4520 }
4521}
4522
4523fn log_crash_respawn(module_id: &str, schedule: CrashRestartSchedule) {
4524 info!(
4525 module_id,
4526 restart_in_window = schedule.restart_in_window,
4527 delay_ms = schedule.delay.as_millis() as u64,
4528 "respawning after crash"
4529 );
4530}
4531
4532fn respawn_still_pending(snapshot: &SharedSnapshot) -> bool {
4538 matches!(
4539 lock_snapshot(snapshot),
4540 Ok(state) if state.enabled && state.state == ModuleState::Restarting
4541 )
4542}
4543
4544enum NextAction {
4545 Stop {
4546 registration_released: bool,
4547 },
4548 Restart {
4549 schedule: Option<CrashRestartSchedule>,
4550 },
4551}
4552
4553#[allow(clippy::too_many_arguments)]
4554async fn handle_supervisor_command(
4555 command: SupervisorCommand,
4556 spec: &mut ModuleSpec,
4557 runtime: &mut SupervisorRuntimeConfig,
4558 registry: &Registry,
4559 process_liveness: &SupervisorProcessLiveness,
4560 snapshot: &SharedSnapshot,
4561 child: &mut Option<SupervisedChild>,
4562 commands: &mut mpsc::Receiver<SupervisorCommand>,
4563 requeued: &mut VecDeque<SupervisorCommand>,
4564) -> bool {
4565 match command {
4566 SupervisorCommand::Drain { reply } => {
4567 let result = drain_optional_child(
4568 &spec.module_id,
4569 spec.protocol,
4570 registry,
4571 snapshot,
4572 &runtime.terminal_ring,
4573 &runtime.spawn_events,
4574 child,
4575 runtime.drain_timeout,
4576 ModuleState::Stopped,
4577 None,
4578 )
4579 .await;
4580 let registration_released = result.is_ok();
4581 let _ = reply.send(result);
4582 if registration_released {
4583 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4584 }
4585 false
4586 }
4587 SupervisorCommand::Retire { reply } => {
4588 let result = async {
4589 begin_forwarding_drain_if_configured(
4590 spec,
4591 runtime,
4592 registry,
4593 snapshot,
4594 None,
4595 RouteCloseReason::Disable,
4596 )
4597 .await?;
4598 drain_optional_child(
4599 &spec.module_id,
4600 spec.protocol,
4601 registry,
4602 snapshot,
4603 &runtime.terminal_ring,
4604 &runtime.spawn_events,
4605 child,
4606 runtime.drain_timeout,
4607 ModuleState::Stopped,
4608 None,
4609 )
4610 .await
4611 }
4612 .await;
4613 let registration_released = result.is_ok();
4614 let _ = reply.send(result);
4615 if registration_released {
4616 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4617 }
4618 false
4619 }
4620 SupervisorCommand::Restart {
4621 drain_timeout_ms,
4622 reply,
4623 } => {
4624 let validation = match lock_snapshot(snapshot) {
4636 Ok(state) if !state.enabled => Err(SuperviseError::Disabled {
4637 module_id: spec.module_id.clone(),
4638 }),
4639 Ok(_) => Ok(()),
4640 Err(err) => Err(err),
4641 };
4642 let initiated = validation.is_ok();
4643 let _ = reply.send(validation);
4644 if initiated {
4645 let drain_timeout = drain_timeout_ms
4648 .map(Duration::from_millis)
4649 .unwrap_or(runtime.drain_timeout);
4650 if let Err(err) = restart_child(
4651 spec,
4652 runtime,
4653 registry,
4654 process_liveness,
4655 snapshot,
4656 child,
4657 drain_timeout,
4658 )
4659 .await
4660 {
4661 warn!(
4662 module_id = %spec.module_id,
4663 error = %err,
4664 "operator restart failed after initiation ack; module state carries the outcome"
4665 );
4666 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4667 state.state = ModuleState::Failed;
4668 clear_current_process_facts(state);
4669 });
4670 }
4671 }
4672 true
4673 }
4674 SupervisorCommand::Reload { reply } => {
4675 let result =
4676 reload_child(spec, runtime, registry, process_liveness, snapshot, child).await;
4677 let _ = reply.send(result);
4678 true
4679 }
4680 SupervisorCommand::SetEnabled { enabled, reply } => {
4681 let result = set_child_enabled(
4682 spec,
4683 runtime,
4684 registry,
4685 process_liveness,
4686 snapshot,
4687 child,
4688 enabled,
4689 )
4690 .await;
4691 let _ = reply.send(result);
4692 true
4693 }
4694 SupervisorCommand::UpdateConfiguration {
4695 spec: next_spec,
4696 health,
4697 drain_timeout_ms,
4698 reply,
4699 } => {
4700 if let Some(handle) = &runtime.supervisor_handle {
4701 handle.apply_identity_configuration(&next_spec);
4702 }
4703 *spec = next_spec;
4704 runtime.health = health;
4705 runtime.drain_timeout = drain_timeout_ms
4706 .map(Duration::from_millis)
4707 .unwrap_or(runtime.default_drain_timeout);
4708 *runtime
4709 .effective_drain_timeout
4710 .lock()
4711 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
4712 let _ = reply.send(());
4713 true
4714 }
4715 SupervisorCommand::Swap {
4716 ready_timeout,
4717 reply,
4718 } => {
4719 let end = swap::run_swap(
4720 spec,
4721 runtime,
4722 registry,
4723 process_liveness,
4724 snapshot,
4725 child,
4726 commands,
4727 ready_timeout.unwrap_or(DEFAULT_SWAP_READY_TIMEOUT),
4728 reply,
4729 )
4730 .await;
4731 requeued.extend(end.requeue);
4732 true
4733 }
4734 }
4735}
4736
4737async fn restart_child(
4738 spec: &ModuleSpec,
4739 runtime: &SupervisorRuntimeConfig,
4740 registry: &Registry,
4741 process_liveness: &SupervisorProcessLiveness,
4742 snapshot: &SharedSnapshot,
4743 child: &mut Option<SupervisedChild>,
4744 drain_timeout: Duration,
4745) -> Result<(), SuperviseError> {
4746 if !lock_snapshot(snapshot)?.enabled {
4748 return Err(SuperviseError::Disabled {
4749 module_id: spec.module_id.clone(),
4750 });
4751 }
4752 begin_forwarding_drain_with_timeout(
4753 spec,
4754 runtime,
4755 registry,
4756 snapshot,
4757 None,
4758 RouteCloseReason::Restart,
4759 drain_timeout,
4760 )
4761 .await?;
4762
4763 if child.is_some() {
4764 drain_optional_child(
4765 &spec.module_id,
4766 spec.protocol,
4767 registry,
4768 snapshot,
4769 &runtime.terminal_ring,
4770 &runtime.spawn_events,
4771 child,
4772 drain_timeout,
4773 ModuleState::Restarting,
4774 Some(true),
4775 )
4776 .await?;
4777 } else {
4778 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4779 state.enabled = true;
4780 state.state = ModuleState::Restarting;
4781 clear_current_process_facts(state);
4782 })?;
4783 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4784 }
4785
4786 reset_restart_count(snapshot, &spec.module_id)?;
4787 sleep(runtime.restart_policy.backoff).await;
4788 if !respawn_still_pending(snapshot) {
4791 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4792 return Ok(());
4793 }
4794 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4795 match spawn_and_mark_running(spec, runtime, snapshot) {
4801 Ok(next_child) => {
4802 *child = Some(next_child);
4803 debug!(module_id = %spec.module_id, "supervised module restarted by operator request");
4804 Ok(())
4805 }
4806 Err(err) => {
4807 fail_snapshot(snapshot, Some(&spec.module_id), None);
4808 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4809 *child = None;
4810 Err(err)
4811 }
4812 }
4813}
4814
4815async fn reload_child(
4816 spec: &ModuleSpec,
4817 runtime: &SupervisorRuntimeConfig,
4818 registry: &Registry,
4819 process_liveness: &SupervisorProcessLiveness,
4820 snapshot: &SharedSnapshot,
4821 child: &mut Option<SupervisedChild>,
4822) -> Result<(), SuperviseError> {
4823 if !lock_snapshot(snapshot)?.enabled {
4825 return Err(SuperviseError::Disabled {
4826 module_id: spec.module_id.clone(),
4827 });
4828 }
4829 begin_forwarding_drain(
4830 spec,
4831 runtime,
4832 registry,
4833 snapshot,
4834 Some(true),
4835 RouteCloseReason::Reload,
4836 )
4837 .await?;
4838
4839 if child.is_some() {
4840 drain_optional_child(
4841 &spec.module_id,
4842 spec.protocol,
4843 registry,
4844 snapshot,
4845 &runtime.terminal_ring,
4846 &runtime.spawn_events,
4847 child,
4848 runtime.drain_timeout,
4849 ModuleState::Restarting,
4850 Some(true),
4851 )
4852 .await?;
4853 } else {
4854 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4855 state.enabled = true;
4856 state.state = ModuleState::Restarting;
4857 clear_current_process_facts(state);
4858 })?;
4859 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4860 }
4861
4862 reset_restart_count(snapshot, &spec.module_id)?;
4863 sleep(runtime.restart_policy.backoff).await;
4864 if !respawn_still_pending(snapshot) {
4867 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4868 return Ok(());
4869 }
4870 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4871 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
4872 Ok(next_child) => next_child,
4873 Err(err) => {
4874 return handle_reload_spawn_failure(
4875 spec,
4876 runtime,
4877 process_liveness,
4878 snapshot,
4879 child,
4880 format!("new child failed to spawn: {err}"),
4881 )
4882 .await;
4883 }
4884 };
4885 *child = Some(next_child);
4886
4887 let wait_outcome = {
4888 let active_child = child.as_mut().expect("new reload child was just stored");
4889 wait_for_registration_after_reload(
4890 registry,
4891 &spec.module_id,
4892 snapshot,
4893 active_child,
4894 REGISTRY_RELEASE_TIMEOUT,
4895 )
4896 .await?
4897 };
4898
4899 match wait_outcome {
4900 RegistrationWaitOutcome::Registered => {
4901 debug!(module_id = %spec.module_id, "supervised module reloaded and registered");
4902 Ok(())
4903 }
4904 RegistrationWaitOutcome::Exited(exit_report) => {
4905 if let Some(active_child) = child.as_mut() {
4906 active_child.drain_stderr(&spec.module_id).await;
4907 }
4908 *child = None;
4909 handle_reload_child_registration_failure(
4910 spec,
4911 runtime,
4912 registry,
4913 process_liveness,
4914 snapshot,
4915 child,
4916 ReloadRegistrationFailure {
4917 exit_report: registration_failure_exit_report(exit_report),
4918 reason: "new child exited before registering".to_string(),
4919 },
4920 )
4921 .await
4922 }
4923 RegistrationWaitOutcome::TimedOut => {
4924 let mut timed_out_child = child
4925 .take()
4926 .expect("timed-out reload child is still running");
4927 timed_out_child
4928 .start_kill()
4929 .map_err(|source| SuperviseError::Kill {
4930 module_id: spec.module_id.clone(),
4931 source,
4932 })?;
4933 let status = timed_out_child
4934 .wait()
4935 .await
4936 .map_err(|source| SuperviseError::Wait {
4937 module_id: spec.module_id.clone(),
4938 source,
4939 })?;
4940 timed_out_child.drain_stderr(&spec.module_id).await;
4941 handle_reload_child_registration_failure(
4942 spec,
4943 runtime,
4944 registry,
4945 process_liveness,
4946 snapshot,
4947 child,
4948 ReloadRegistrationFailure {
4949 exit_report: registration_failure_exit_report(classify_reaped_child_exit(
4950 snapshot,
4951 &timed_out_child,
4952 &status,
4953 )),
4954 reason: format!(
4955 "new child did not register within {:?}",
4956 REGISTRY_RELEASE_TIMEOUT
4957 ),
4958 },
4959 )
4960 .await
4961 }
4962 }
4963}
4964
4965async fn set_child_enabled(
4966 spec: &ModuleSpec,
4967 runtime: &SupervisorRuntimeConfig,
4968 registry: &Registry,
4969 process_liveness: &SupervisorProcessLiveness,
4970 snapshot: &SharedSnapshot,
4971 child: &mut Option<SupervisedChild>,
4972 enabled: bool,
4973) -> Result<bool, SuperviseError> {
4974 let (current_enabled, current_state) = {
4975 let state = lock_snapshot(snapshot)?;
4976 (state.enabled, state.state)
4977 };
4978 let revive_terminal = enabled
4986 && current_enabled
4987 && child.is_none()
4988 && matches!(current_state, ModuleState::Failed | ModuleState::Stopped);
4989 if current_enabled == enabled && !revive_terminal {
4990 return Ok(false);
4991 }
4992
4993 if enabled {
4994 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4995 state.enabled = true;
4996 state.state = ModuleState::Starting;
4997 clear_current_process_facts(state);
4998 })?;
4999 #[cfg(test)]
5000 if runtime.test_seed_stale_facts_before_enable_spawn {
5001 update_snapshot(snapshot, Some(&spec.module_id), |state| {
5002 state.process_alive = true;
5003 state.pid = Some(41);
5004 state.spawned_at_ms = Some(42);
5005 state.spawned_from = Some(PathBuf::from("/spawned/module"));
5006 state.spawned_file_identity = Some(SpawnedFileIdentity {
5007 device: 43,
5008 inode: 44,
5009 });
5010 })?;
5011 }
5012 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
5013 reset_restart_count(snapshot, &spec.module_id)?;
5014 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
5015 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
5016 Ok(next_child) => next_child,
5017 Err(err) => {
5018 if let Err(state_err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5019 state.state = ModuleState::Failed;
5020 clear_current_process_facts(state);
5021 }) {
5022 error!(module_id = %spec.module_id, error = %state_err, "failed to record enable spawn failure");
5023 }
5024 process_liveness.untrack_if_current(&spec.module_id, snapshot);
5025 return Err(err);
5026 }
5027 };
5028 *child = Some(next_child);
5029 debug!(module_id = %spec.module_id, "supervised module enabled");
5030 Ok(true)
5031 } else {
5032 begin_forwarding_drain_if_configured(
5033 spec,
5034 runtime,
5035 registry,
5036 snapshot,
5037 Some(false),
5038 RouteCloseReason::Disable,
5039 )
5040 .await?;
5041 drain_optional_child(
5042 &spec.module_id,
5043 spec.protocol,
5044 registry,
5045 snapshot,
5046 &runtime.terminal_ring,
5047 &runtime.spawn_events,
5048 child,
5049 runtime.drain_timeout,
5050 ModuleState::Disabled,
5051 Some(false),
5052 )
5053 .await?;
5054 debug!(module_id = %spec.module_id, "supervised module disabled");
5055 Ok(true)
5056 }
5057}
5058
5059#[allow(clippy::too_many_arguments)]
5060async fn on_child_exit(
5061 spec: &ModuleSpec,
5062 policy: RestartPolicy,
5063 registry: &Registry,
5064 snapshot: &SharedSnapshot,
5065 terminal_ring: &Arc<Mutex<TerminalRing>>,
5066 spawn_events: &SpawnEventFeed,
5067 roster: &ChildRoster,
5068 exit_report: ExitReport,
5069) -> NextAction {
5070 if roster.is_closed() {
5076 return on_child_exit_during_daemon_shutdown(
5077 spec,
5078 registry,
5079 snapshot,
5080 terminal_ring,
5081 spawn_events,
5082 exit_report,
5083 )
5084 .await;
5085 }
5086 match exit_report.kind {
5087 ExitKind::Clean => {
5088 info!(
5089 module_id = %spec.module_id,
5090 exit_code = ?exit_report.code,
5091 exit_signal = ?exit_report.signal,
5092 "supervised module exited cleanly"
5093 );
5094 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5095 state.state = ModuleState::Stopped;
5096 clear_current_process_facts(state);
5097 state.last_exit = Some(exit_report.clone());
5098 }) {
5099 error!(module_id = %spec.module_id, error = %err, "failed to record clean module exit");
5100 }
5101 record_terminal(
5102 &spec.module_id,
5103 terminal_ring,
5104 spawn_events,
5105 &exit_report,
5106 TerminalDisposition::Stopped,
5107 );
5108 let registration_released = match wait_for_registration_release(
5109 registry,
5110 &spec.module_id,
5111 REGISTRY_RELEASE_TIMEOUT,
5112 )
5113 .await
5114 {
5115 Ok(()) => true,
5116 Err(err) => {
5117 warn!(module_id = %spec.module_id, error = %err, "registration still active after clean exit");
5118 false
5119 }
5120 };
5121 NextAction::Stop {
5122 registration_released,
5123 }
5124 }
5125 ExitKind::Crash => {
5126 warn!(
5127 module_id = %spec.module_id,
5128 exit_code = ?exit_report.code,
5129 exit_signal = ?exit_report.signal,
5130 "supervised module exited abnormally (crash)"
5131 );
5132 let mut restart_schedule = None;
5133 let mut disposition = TerminalDisposition::Disabled;
5134 let mut disposition_detail = None;
5138 let now = Instant::now();
5139 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5140 clear_current_process_facts(state);
5141 state.last_exit = Some(exit_report.clone());
5142 if state.enabled {
5143 if let Some(schedule) = state.next_crash_restart(&policy, now) {
5144 state.state = ModuleState::Restarting;
5145 restart_schedule = Some(schedule);
5146 disposition = TerminalDisposition::Restarting;
5147 } else {
5148 state.state = ModuleState::Failed;
5149 disposition = TerminalDisposition::Failed;
5150 disposition_detail = Some(policy.budget_exhausted_detail());
5151 }
5152 } else {
5153 state.state = ModuleState::Disabled;
5154 disposition = TerminalDisposition::Disabled;
5155 }
5156 }) {
5157 error!(module_id = %spec.module_id, error = %err, "failed to record crashed module exit");
5158 return NextAction::Stop {
5159 registration_released: false,
5160 };
5161 }
5162 if disposition_detail.is_some() {
5163 error!(
5168 module_id = %spec.module_id,
5169 max_restarts = policy.max_restarts,
5170 window_secs = policy.window.as_secs(),
5171 "module stopped: {}",
5172 policy.budget_exhausted_detail()
5173 );
5174 }
5175 record_terminal_with_detail(
5176 &spec.module_id,
5177 terminal_ring,
5178 spawn_events,
5179 &exit_report,
5180 disposition,
5181 disposition_detail,
5182 );
5183
5184 if let Some(schedule) = restart_schedule {
5185 NextAction::Restart {
5186 schedule: Some(schedule),
5187 }
5188 } else {
5189 let registration_released = match wait_for_registration_release(
5190 registry,
5191 &spec.module_id,
5192 REGISTRY_RELEASE_TIMEOUT,
5193 )
5194 .await
5195 {
5196 Ok(()) => true,
5197 Err(err) => {
5198 warn!(module_id = %spec.module_id, error = %err, "registration still active after failed module");
5199 false
5200 }
5201 };
5202 NextAction::Stop {
5203 registration_released,
5204 }
5205 }
5206 }
5207 ExitKind::DeliberateSeverance => {
5208 warn!(
5209 module_id = %spec.module_id,
5210 exit_code = ?exit_report.code,
5211 exit_signal = ?exit_report.signal,
5212 "supervised module exited after deliberate connection severance"
5213 );
5214 let mut should_restart = false;
5215 let mut disposition = TerminalDisposition::Disabled;
5216 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5217 clear_current_process_facts(state);
5218 state.last_exit = Some(exit_report.clone());
5219 state.lifetime_restarts += 1;
5220 if state.enabled {
5221 state.state = ModuleState::Restarting;
5222 should_restart = true;
5223 disposition = TerminalDisposition::Restarting;
5224 } else {
5225 state.state = ModuleState::Disabled;
5226 }
5227 }) {
5228 error!(module_id = %spec.module_id, error = %err, "failed to record deliberately severed module exit");
5229 return NextAction::Stop {
5230 registration_released: false,
5231 };
5232 }
5233 record_terminal(
5234 &spec.module_id,
5235 terminal_ring,
5236 spawn_events,
5237 &exit_report,
5238 disposition,
5239 );
5240
5241 if should_restart {
5242 NextAction::Restart { schedule: None }
5243 } else {
5244 let registration_released = match wait_for_registration_release(
5245 registry,
5246 &spec.module_id,
5247 REGISTRY_RELEASE_TIMEOUT,
5248 )
5249 .await
5250 {
5251 Ok(()) => true,
5252 Err(err) => {
5253 warn!(module_id = %spec.module_id, error = %err, "registration still active after deliberately severed module exit");
5254 false
5255 }
5256 };
5257 NextAction::Stop {
5258 registration_released,
5259 }
5260 }
5261 }
5262 }
5263}
5264
5265async fn on_child_exit_during_daemon_shutdown(
5266 spec: &ModuleSpec,
5267 registry: &Registry,
5268 snapshot: &SharedSnapshot,
5269 terminal_ring: &Arc<Mutex<TerminalRing>>,
5270 spawn_events: &SpawnEventFeed,
5271 exit_report: ExitReport,
5272) -> NextAction {
5273 info!(
5274 module_id = %spec.module_id,
5275 exit_code = ?exit_report.code,
5276 exit_signal = ?exit_report.signal,
5277 exit_kind = ?exit_report.kind,
5278 "supervised module exited during daemon shutdown; not restarting it"
5279 );
5280 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5281 state.state = ModuleState::Stopped;
5282 clear_current_process_facts(state);
5283 state.last_exit = Some(exit_report.clone());
5284 }) {
5285 error!(module_id = %spec.module_id, error = %err, "failed to record module exit during daemon shutdown");
5286 }
5287 record_terminal(
5288 &spec.module_id,
5289 terminal_ring,
5290 spawn_events,
5291 &exit_report,
5292 TerminalDisposition::DaemonShutdown,
5293 );
5294 let registration_released =
5295 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT)
5296 .await
5297 .is_ok();
5298 NextAction::Stop {
5299 registration_released,
5300 }
5301}
5302
5303fn record_wait_error_terminal(
5304 module_id: &str,
5305 terminal_ring: &Arc<Mutex<TerminalRing>>,
5306 spawn_events: &SpawnEventFeed,
5307) {
5308 record_terminal(
5309 module_id,
5310 terminal_ring,
5311 spawn_events,
5312 &wait_error_exit_report(),
5313 TerminalDisposition::Failed,
5314 );
5315}
5316
5317fn record_terminal(
5318 module_id: &str,
5319 terminal_ring: &Arc<Mutex<TerminalRing>>,
5320 spawn_events: &SpawnEventFeed,
5321 exit_report: &ExitReport,
5322 disposition: TerminalDisposition,
5323) {
5324 record_terminal_with_detail(
5325 module_id,
5326 terminal_ring,
5327 spawn_events,
5328 exit_report,
5329 disposition,
5330 None,
5331 );
5332}
5333
5334fn durable_terminal_history_of(
5338 terminal_ring: &Mutex<TerminalRing>,
5339 module_id: &str,
5340) -> subc_control::TerminalHistory {
5341 let read = terminal_ring
5342 .lock()
5343 .unwrap_or_else(|p| p.into_inner())
5344 .capture_durable_history();
5345 read.read(module_id)
5346}
5347
5348fn record_terminal_with_detail(
5349 module_id: &str,
5350 terminal_ring: &Arc<Mutex<TerminalRing>>,
5351 spawn_events: &SpawnEventFeed,
5352 exit_report: &ExitReport,
5353 disposition: TerminalDisposition,
5354 disposition_detail: Option<String>,
5355) {
5356 spawn_events.emit_exited(module_id, exit_report.code, exit_report.signal);
5357 let record = TerminalRecord {
5358 exit_code: exit_report.code,
5359 exit_signal: exit_report.signal,
5360 at_ms: exit_report.at_ms,
5361 disposition,
5362 exit_kind: exit_report.kind.into(),
5363 disposition_detail,
5364 };
5365 terminal_ring
5366 .lock()
5367 .unwrap_or_else(|poisoned| poisoned.into_inner())
5368 .record_exit(module_id, record);
5369}
5370
5371fn untrack_if_registration_released(
5372 process_liveness: &SupervisorProcessLiveness,
5373 registry: &Registry,
5374 module_id: &str,
5375 snapshot: &SharedSnapshot,
5376) {
5377 match registry.get_module(module_id) {
5378 Ok(None) => process_liveness.untrack_if_current(module_id, snapshot),
5379 Ok(Some(_)) => {}
5380 Err(err) => {
5381 warn!(module_id, error = %err, "could not determine whether supervisor liveness can be untracked");
5382 }
5383 }
5384}
5385
5386#[cfg(test)]
5400fn apply_wire_spawn_args(
5401 command: &mut Command,
5402 spec: &ModuleSpec,
5403 connection_file_path: Option<&std::path::Path>,
5404 handle: Option<&SupervisorHandle>,
5405) -> Result<(), SuperviseError> {
5406 apply_wire_spawn_args_for_role(
5407 command,
5408 spec,
5409 connection_file_path,
5410 handle,
5411 SpawnRole::Plain,
5412 )
5413}
5414
5415fn apply_wire_spawn_args_for_role(
5424 command: &mut Command,
5425 spec: &ModuleSpec,
5426 connection_file_path: Option<&std::path::Path>,
5427 handle: Option<&SupervisorHandle>,
5428 role: SpawnRole,
5429) -> Result<(), SuperviseError> {
5430 command.env(SUBC_MODULE_ID_ENV, &spec.module_id);
5431 if spec.protocol == ModuleProtocol::None {
5432 return Ok(());
5433 }
5434 if let Some(connection_file_path) = connection_file_path {
5435 command.arg(SUBC_ARG).arg(connection_file_path);
5436 }
5437
5438 let nonce = generate_launch_nonce()?;
5442 if let Some(handle) = handle {
5443 match role {
5444 SpawnRole::Plain => {
5445 handle.set_spawn_nonce(&spec.module_id, nonce.clone());
5446 if spec.reserved {
5447 handle.set_reserved_nonce(&spec.module_id, nonce.clone());
5448 }
5449 }
5450 SpawnRole::SwapCandidate => handle.open_swap(&spec.module_id, nonce.clone()),
5451 }
5452 }
5453 command.env(SUBC_LAUNCH_NONCE_ENV, nonce);
5454 Ok(())
5455}
5456
5457fn apply_child_env(command: &mut Command, spec: &ModuleSpec) {
5458 command.env_remove(CK_LOG_ENV);
5459 command.env_remove(SUBC_SPAWN_ROLE_ENV);
5466 for (key, value) in &spec.env {
5467 if matches!(
5471 key.as_str(),
5472 CAPTURE_MAX_FILE_MB_ENV | CAPTURE_KEEP_ENV | CAPTURE_MAX_AGE_DAYS_ENV
5473 ) || key == SUBC_SPAWN_ROLE_ENV
5474 {
5475 continue;
5476 }
5477 command.env(key, value);
5478 }
5479}
5480
5481#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5484enum SpawnRole {
5485 Plain,
5486 SwapCandidate,
5487}
5488
5489fn apply_spawn_role(command: &mut Command, role: SpawnRole) {
5492 if role == SpawnRole::SwapCandidate {
5493 command.env(SUBC_SPAWN_ROLE_ENV, SPAWN_ROLE_SWAP_CANDIDATE);
5494 }
5495}
5496
5497fn spawn_child(
5498 spec: &ModuleSpec,
5499 connection_file_path: Option<&std::path::Path>,
5500 handle: Option<&SupervisorHandle>,
5501 ring: &Arc<Mutex<StderrRing>>,
5502 capture_logs_dir: Option<&std::path::Path>,
5503 roster: &ChildRoster,
5504 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5505) -> Result<SupervisedChild, SuperviseError> {
5506 spawn_child_in_slot(
5507 spec,
5508 connection_file_path,
5509 handle,
5510 ring,
5511 capture_logs_dir,
5512 roster,
5513 #[cfg(target_os = "linux")]
5514 cgroup_placement,
5515 SpawnRole::Plain,
5516 false,
5517 )
5518}
5519
5520#[allow(clippy::too_many_arguments)]
5533fn spawn_child_in_slot(
5534 spec: &ModuleSpec,
5535 connection_file_path: Option<&std::path::Path>,
5536 handle: Option<&SupervisorHandle>,
5537 ring: &Arc<Mutex<StderrRing>>,
5538 capture_logs_dir: Option<&std::path::Path>,
5539 roster: &ChildRoster,
5540 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5541 role: SpawnRole,
5542 alternate_slot: bool,
5543) -> Result<SupervisedChild, SuperviseError> {
5544 if roster.is_closed() {
5545 return Err(SuperviseError::Spawn {
5546 program: spec.program.clone(),
5547 source: io::Error::other("the daemon is shutting down; not starting a new process"),
5548 cgroup_path: None,
5549 });
5550 }
5551 #[cfg(target_os = "linux")]
5552 let cgroup_name = swap::cgroup_name(&spec.module_id, alternate_slot);
5553 #[cfg(not(target_os = "linux"))]
5554 let _ = alternate_slot;
5555 let mut command = Command::new(&spec.program);
5556 command.args(&spec.args);
5557 apply_child_env(&mut command, spec);
5587 apply_spawn_role(&mut command, role);
5588 apply_wire_spawn_args_for_role(&mut command, spec, connection_file_path, handle, role)?;
5589
5590 #[cfg(target_os = "linux")]
5591 let cgroup_path = cgroup_placement
5592 .map(|placement| placement.module_path(&cgroup_name))
5593 .transpose()
5594 .map_err(|source| SuperviseError::Cgroup {
5595 module_id: spec.module_id.clone(),
5596 source,
5597 })?;
5598 #[cfg(not(target_os = "linux"))]
5599 let cgroup_path: Option<PathBuf> = None;
5600 #[cfg(target_os = "linux")]
5601 if let Some(path) = &cgroup_path {
5602 if let Err(error) = apply_cgroup_placement(&mut command, spec, path) {
5603 if let Some(placement) = cgroup_placement {
5604 remove_module_cgroup(placement, &cgroup_name);
5605 }
5606 return Err(error);
5607 }
5608 }
5609
5610 let output_sink = if let Some(logs_dir) = capture_logs_dir {
5611 let path = logs_dir.join(format!("{}.stderr.log", spec.module_id));
5612 match ChildOutputSink::open(&path, capture_retention(spec)) {
5613 Ok(sink) => sink,
5614 Err(error) => {
5615 warn!(
5616 module_id = %spec.module_id,
5617 path = %path.display(),
5618 error = %error,
5619 "could not open child output capture file; forwarding to stderr"
5620 );
5621 ChildOutputSink::Stderr
5622 }
5623 }
5624 } else {
5625 ChildOutputSink::Stderr
5626 };
5627
5628 command.stdout(Stdio::piped());
5629 command.stderr(Stdio::piped());
5630 command.kill_on_drop(true);
5631 #[cfg(unix)]
5648 command.process_group(0);
5649 command.stdin(Stdio::null());
5650 let mut child = match command.spawn() {
5651 Ok(child) => child,
5652 Err(source) => {
5653 #[cfg(target_os = "linux")]
5654 if let Some(placement) = cgroup_placement {
5655 remove_module_cgroup(placement, &cgroup_name);
5656 }
5657 return Err(SuperviseError::Spawn {
5658 program: spec.program.clone(),
5659 source,
5660 cgroup_path,
5661 });
5662 }
5663 };
5664 let spawned_at_ms = unix_ms_now();
5665 let spawned_from = spec.program.clone();
5666 let spawned_file_identity = spawned_file_identity(&spawned_from);
5667 let pid = child.id().ok_or_else(|| SuperviseError::Spawn {
5668 program: spec.program.clone(),
5669 source: io::Error::other("spawned child exposed no live pid"),
5670 cgroup_path: cgroup_path.clone(),
5671 })?;
5672 let process_start_time = crate::provenance::process_start_time(pid);
5673 let process_identity = process_start_time.map(|start_time| ProcessIdentity { pid, start_time });
5674 let roster_guard = roster.admit(
5675 spec.module_id.clone(),
5676 pid,
5677 spec.protocol,
5678 process_start_time,
5679 );
5680 if roster.is_closed() {
5689 if let Err(error) = child.start_kill() {
5690 debug!(module_id = %spec.module_id, pid, %error, "kill of a process spawned during daemon shutdown failed; it may already have exited");
5691 }
5692 drop(roster_guard);
5693 return Err(SuperviseError::Spawn {
5694 program: spec.program.clone(),
5695 source: io::Error::other(
5696 "the daemon began shutting down while this process was starting; ended it",
5697 ),
5698 cgroup_path,
5699 });
5700 }
5701
5702 let stdout_pump = match child.stdout.take() {
5703 Some(stdout) => Some(tokio::spawn(pump_stdout_to(stdout, output_sink.clone()))),
5704 None => {
5705 warn!(
5706 module_id = %spec.module_id,
5707 "spawned child exposed no stdout pipe; file capture will be incomplete"
5708 );
5709 None
5710 }
5711 };
5712 let stderr_pump = match child.stderr.take() {
5713 Some(stderr) => {
5714 let generation = ring
5715 .lock()
5716 .unwrap_or_else(|poisoned| poisoned.into_inner())
5717 .begin_process();
5718 Some(StderrPump {
5719 task: tokio::spawn(pump_stderr_to(
5720 stderr,
5721 Arc::clone(ring),
5722 generation,
5723 output_sink,
5724 )),
5725 generation,
5726 })
5727 }
5728 None => {
5729 ring.lock()
5733 .unwrap_or_else(|poisoned| poisoned.into_inner())
5734 .mark_not_captured("stderr pipe was not available on spawn");
5735 warn!(
5736 module_id = %spec.module_id,
5737 "spawned child exposed no stderr pipe; tail will be unavailable"
5738 );
5739 None
5740 }
5741 };
5742
5743 Ok(SupervisedChild {
5744 child,
5745 #[cfg(target_os = "linux")]
5746 module_id: cgroup_name,
5747 #[cfg(target_os = "linux")]
5748 cgroup_placement: cgroup_placement.cloned(),
5749 stdout_pump,
5750 stderr_pump,
5751 stderr_ring: Arc::clone(ring),
5752 spawned_at_ms,
5753 spawned_from,
5754 spawned_file_identity,
5755 process_start_time,
5756 process_identity,
5757 pid,
5758 roster_guard: Some(roster_guard),
5759 })
5760}
5761
5762#[cfg(target_os = "linux")]
5763fn remove_module_cgroup(placement: &subc_cgroup::Placement, module_id: &str) {
5764 match placement.remove_module(module_id) {
5765 Ok(()) => debug!(module_id, "removed module cgroup after process exit"),
5766 Err(error) => warn!(
5767 module_id,
5768 error = %error,
5769 "could not remove module cgroup after process exit; continuing teardown"
5770 ),
5771 }
5772}
5773
5774#[cfg(target_os = "linux")]
5775fn apply_cgroup_placement(
5776 command: &mut Command,
5777 spec: &ModuleSpec,
5778 path: &std::path::Path,
5779) -> Result<(), SuperviseError> {
5780 subc_cgroup::apply(command, path).map_err(|source| SuperviseError::Cgroup {
5781 module_id: spec.module_id.clone(),
5782 source,
5783 })
5784}
5785
5786fn capture_retention(spec: &ModuleSpec) -> Retention {
5787 let defaults = Retention::default();
5788 let value = |name: &str| {
5789 spec.env
5790 .iter()
5791 .rev()
5792 .find_map(|(key, value)| (key == name).then_some(value.as_str()))
5793 };
5794 Retention {
5795 max_file_mb: value(CAPTURE_MAX_FILE_MB_ENV)
5796 .and_then(|value| value.parse().ok())
5797 .unwrap_or(defaults.max_file_mb),
5798 keep: value(CAPTURE_KEEP_ENV)
5799 .and_then(|value| value.parse().ok())
5800 .unwrap_or(defaults.keep),
5801 max_age_days: value(CAPTURE_MAX_AGE_DAYS_ENV)
5802 .and_then(|value| value.parse().ok())
5803 .unwrap_or(defaults.max_age_days),
5804 }
5805}
5806
5807fn generate_launch_nonce() -> Result<String, SuperviseError> {
5810 let mut bytes = [0u8; 32];
5811 getrandom::getrandom(&mut bytes).map_err(|source| SuperviseError::LaunchNonce {
5812 reason: source.to_string(),
5813 })?;
5814 let mut hex = String::with_capacity(64);
5815 for b in bytes {
5816 use std::fmt::Write;
5817 let _ = write!(hex, "{b:02x}");
5818 }
5819 Ok(hex)
5820}
5821
5822fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
5825 if a.len() != b.len() {
5826 return false;
5827 }
5828 let mut diff = 0u8;
5829 for (x, y) in a.iter().zip(b.iter()) {
5830 diff |= x ^ y;
5831 }
5832 diff == 0
5833}
5834
5835fn spawn_and_mark_running(
5836 spec: &ModuleSpec,
5837 runtime: &SupervisorRuntimeConfig,
5838 snapshot: &SharedSnapshot,
5839) -> Result<SupervisedChild, SuperviseError> {
5840 let child = spawn_child(
5841 spec,
5842 runtime.connection_file_path.as_deref(),
5843 runtime.supervisor_handle.as_ref(),
5844 &runtime.stderr_ring,
5845 runtime.capture_logs_dir.as_deref(),
5846 &runtime.child_roster,
5847 #[cfg(target_os = "linux")]
5848 runtime.cgroup_placement.as_ref(),
5849 )?;
5850 set_running(snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
5851 Ok(child)
5852}
5853
5854enum RegistrationWaitOutcome {
5855 Registered,
5856 Exited(ExitReport),
5857 TimedOut,
5858}
5859
5860struct ReloadRegistrationFailure {
5861 exit_report: ExitReport,
5862 reason: String,
5863}
5864
5865#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5866enum BusyGaugeObservation {
5867 Quiescent,
5868 Busy,
5869 Omitted,
5870}
5871
5872fn busy_gauge_observation(metrics: Option<&Value>, gauges: &[String]) -> BusyGaugeObservation {
5873 let Some(metrics) = metrics.and_then(Value::as_object) else {
5874 return BusyGaugeObservation::Omitted;
5875 };
5876 let mut sum = 0u128;
5877 for gauge in gauges {
5878 let Some(value) = metrics.get(gauge) else {
5879 return BusyGaugeObservation::Omitted;
5880 };
5881 let Some(value) = value.as_u64() else {
5882 return BusyGaugeObservation::Busy;
5883 };
5884 sum = sum.saturating_add(u128::from(value));
5885 }
5886 if sum == 0 {
5887 BusyGaugeObservation::Quiescent
5888 } else {
5889 BusyGaugeObservation::Busy
5890 }
5891}
5892
5893fn declared_busy_gauges(
5894 registry: &Registry,
5895 module_id: &str,
5896) -> Result<Vec<String>, SuperviseError> {
5897 busy_gauges_of(
5898 registry
5899 .get_module(module_id)
5900 .map_err(SuperviseError::Registry)?,
5901 )
5902}
5903
5904fn declared_busy_gauges_for_connection(
5908 registry: &Registry,
5909 connection_id: ConnectionId,
5910) -> Result<Vec<String>, SuperviseError> {
5911 busy_gauges_of(
5912 registry
5913 .get_module_by_connection(connection_id)
5914 .map_err(SuperviseError::Registry)?,
5915 )
5916}
5917
5918fn busy_gauges_of(
5919 registration: Option<crate::registry::ModuleRegistration>,
5920) -> Result<Vec<String>, SuperviseError> {
5921 let Some(registration) = registration else {
5922 return Ok(Vec::new());
5923 };
5924 let Some(self_signals) = registration.manifest.self_signals else {
5925 return Ok(Vec::new());
5926 };
5927
5928 let mut gauges = Vec::new();
5929 for declaration in self_signals {
5930 if declaration.kind != SelfSignalKind::Busy {
5931 continue;
5932 }
5933 match declaration.anchored_to {
5934 SignalAnchor::HealthGauges { gauges: declared } if !declared.is_empty() => {
5935 gauges.extend(declared)
5936 }
5937 _ => {
5938 gauges.push(String::new());
5941 }
5942 }
5943 }
5944 Ok(gauges)
5945}
5946
5947async fn wait_for_forwarding_quiescence(
5952 forwarding: &ForwardingTable,
5953 module_id: &str,
5954 runtime: &SupervisorRuntimeConfig,
5955 endpoint: crate::ModuleEndpointId,
5956 deadline: Instant,
5957 busy_gauges: &[String],
5958 scope: DrainScope,
5959) -> Result<bool, SuperviseError> {
5960 let mut gauges_quiescent = busy_gauges.is_empty();
5961 let mut next_probe_at = Instant::now();
5962 let mut omission_counted = false;
5963
5964 loop {
5965 let now = Instant::now();
5966 if !busy_gauges.is_empty() && now >= next_probe_at && now < deadline {
5967 let report = match scope {
5968 DrainScope::Active => probe_module_health(module_id, runtime, Some(deadline)).await,
5969 DrainScope::Endpoint(endpoint) => {
5970 probe_endpoint_health(endpoint, runtime, Some(deadline)).await
5971 }
5972 };
5973 gauges_quiescent = match report {
5974 Ok(report) => match busy_gauge_observation(report.metrics.as_ref(), busy_gauges) {
5975 BusyGaugeObservation::Quiescent => true,
5976 BusyGaugeObservation::Busy => false,
5977 BusyGaugeObservation::Omitted => {
5978 if !omission_counted {
5979 forwarding
5980 .counters()
5981 .increment_drains_with_undeclared_gauge();
5982 omission_counted = true;
5983 }
5984 false
5985 }
5986 },
5987 Err(err) => {
5988 warn!(
5989 module_id,
5990 error = %err,
5991 "drain health.check did not produce declared busy gauges; treating module as busy"
5992 );
5993 false
5994 }
5995 };
5996 next_probe_at = Instant::now() + runtime.health.cadence.max(REGISTRY_RELEASE_POLL);
5997 }
5998
5999 let in_flight = forwarding
6000 .endpoint_in_flight_count(endpoint)
6001 .map_err(SuperviseError::Forwarding)?;
6002 if in_flight == 0 && gauges_quiescent {
6003 return Ok(true);
6004 }
6005
6006 let now = Instant::now();
6007 if now >= deadline {
6008 return Ok(false);
6009 }
6010 let mut wait = deadline
6011 .saturating_duration_since(now)
6012 .min(REGISTRY_RELEASE_POLL);
6013 if !busy_gauges.is_empty() {
6014 wait = wait.min(next_probe_at.saturating_duration_since(now));
6015 }
6016 sleep(wait).await;
6017 }
6018}
6019
6020fn drained_after_quiescence_wait(wait_result: &Result<bool, SuperviseError>) -> bool {
6028 match wait_result {
6029 Ok(drained) => *drained,
6030 Err(_) => false,
6031 }
6032}
6033
6034fn send_route_goodbyes(forwarding: &ForwardingTable, released_routes: Vec<GoodbyeTarget>) {
6035 for released in released_routes {
6036 let frame = match Frame::build_with_version(
6037 released.negotiated_ver,
6038 FrameType::Goodbye,
6039 control_flags(),
6040 released.channel,
6041 released.epoch,
6042 0,
6043 Vec::new(),
6044 ) {
6045 Ok(frame) => frame,
6046 Err(err) => {
6047 warn!(
6048 route_channel = released.channel,
6049 error = %err,
6050 "failed to build supervisor drain route GOODBYE frame"
6051 );
6052 continue;
6053 }
6054 };
6055 if !released.close_on_delivery_failure() {
6056 crate::forwarding::send_module_route_goodbye(
6057 &forwarding.counters(),
6058 &released.sink,
6059 frame,
6060 released.module_id.as_deref(),
6061 "supervisor drain",
6062 );
6063 continue;
6064 }
6065 if let Err(err) = released.sink.try_send(frame) {
6066 warn!(
6067 target_connection_id = released.connection_id.get(),
6068 route_channel = released.channel,
6069 error = %err,
6070 "supervisor drain route GOODBYE was not delivered to client; closing target connection"
6071 );
6072 let _ = forwarding.escalate_client_delivery_failure(
6073 released.connection_id,
6074 released.channel,
6075 released.epoch,
6076 CloseReason::new(
6077 "route_goodbye_delivery_failed",
6078 format!(
6079 "failed to enqueue supervisor drain route GOODBYE for channel {}: {err}",
6080 released.channel
6081 ),
6082 ),
6083 crate::forwarding::UndeliveredFrame {
6084 module_id: released.module_id.as_deref(),
6085 sink: &released.sink,
6086 },
6087 );
6088 }
6089 }
6090}
6091
6092fn send_module_draining(
6093 module_id: &str,
6094 reason: RouteCloseReason,
6095 deadline_ms: u64,
6096 target: &ModuleDrainTarget,
6097) {
6098 let body = match serde_json::to_vec(&ModuleControlCommand::Draining {
6099 reason,
6100 deadline_ms,
6101 }) {
6102 Ok(body) => body,
6103 Err(err) => {
6104 warn!(
6105 module_id,
6106 error = %err,
6107 "failed to encode module draining command"
6108 );
6109 return;
6110 }
6111 };
6112 let frame = match Frame::build_with_version(
6113 target.negotiated_ver,
6114 FrameType::Push,
6115 control_flags(),
6116 0,
6117 0,
6118 0,
6119 body,
6120 ) {
6121 Ok(frame) => frame,
6122 Err(err) => {
6123 warn!(
6124 module_id,
6125 error = %err,
6126 "failed to build module draining command frame"
6127 );
6128 return;
6129 }
6130 };
6131 if let Err(err) = target.sink.try_send(frame) {
6132 warn!(
6133 module_id,
6134 target_connection_id = target.endpoint.connection_id.get(),
6135 error = %err,
6136 "module draining command was not delivered to peer"
6137 );
6138 }
6139}
6140
6141fn send_module_goodbye(module_id: &str, forwarding: &ForwardingTable, target: &ModuleDrainTarget) {
6142 let frame = match Frame::build_with_version(
6143 target.negotiated_ver,
6144 FrameType::Goodbye,
6145 control_flags(),
6146 0,
6147 0,
6148 0,
6149 Vec::new(),
6150 ) {
6151 Ok(frame) => frame,
6152 Err(err) => {
6153 warn!(
6154 module_id,
6155 error = %err,
6156 "failed to build supervisor drain module GOODBYE frame"
6157 );
6158 return;
6159 }
6160 };
6161 if let Err(err) = target.sink.try_send(frame) {
6162 warn!(
6163 module_id,
6164 target_connection_id = target.endpoint.connection_id.get(),
6165 error = %err,
6166 "supervisor drain module GOODBYE was not delivered to peer; closing module connection"
6167 );
6168 forwarding.request_connection_close(
6169 target.endpoint.connection_id,
6170 CloseReason::new(
6171 "module_goodbye_delivery_failed",
6172 format!("failed to enqueue supervisor drain module GOODBYE for module '{module_id}': {err}"),
6173 ),
6174 );
6175 }
6176}
6177
6178#[derive(Clone, Copy)]
6179struct ForwardingDrainContext<'a> {
6180 spec: &'a ModuleSpec,
6181 runtime: &'a SupervisorRuntimeConfig,
6182 registry: &'a Registry,
6183 scope: DrainScope,
6184}
6185
6186#[derive(Debug, Clone, Copy, PartialEq, Eq)]
6188enum DrainScope {
6189 Active,
6192 Endpoint(crate::ModuleEndpointId),
6197}
6198
6199async fn begin_forwarding_drain(
6200 spec: &ModuleSpec,
6201 runtime: &SupervisorRuntimeConfig,
6202 registry: &Registry,
6203 snapshot: &SharedSnapshot,
6204 enabled: Option<bool>,
6205 reason: RouteCloseReason,
6206) -> Result<(), SuperviseError> {
6207 let Some(forwarding) = runtime.forwarding.as_ref() else {
6208 return Err(SuperviseError::ReloadUnavailable {
6209 module_id: spec.module_id.clone(),
6210 reason: "supervisor was not configured with a forwarding table".to_string(),
6211 });
6212 };
6213
6214 begin_forwarding_drain_with(
6215 forwarding,
6216 ForwardingDrainContext {
6217 spec,
6218 runtime,
6219 registry,
6220 scope: DrainScope::Active,
6221 },
6222 snapshot,
6223 enabled,
6224 reason,
6225 runtime.drain_timeout,
6226 )
6227 .await
6228}
6229
6230async fn begin_forwarding_drain_if_configured(
6231 spec: &ModuleSpec,
6232 runtime: &SupervisorRuntimeConfig,
6233 registry: &Registry,
6234 snapshot: &SharedSnapshot,
6235 enabled: Option<bool>,
6236 reason: RouteCloseReason,
6237) -> Result<(), SuperviseError> {
6238 begin_forwarding_drain_with_timeout(
6239 spec,
6240 runtime,
6241 registry,
6242 snapshot,
6243 enabled,
6244 reason,
6245 runtime.drain_timeout,
6246 )
6247 .await
6248}
6249
6250async fn begin_forwarding_drain_with_timeout(
6254 spec: &ModuleSpec,
6255 runtime: &SupervisorRuntimeConfig,
6256 registry: &Registry,
6257 snapshot: &SharedSnapshot,
6258 enabled: Option<bool>,
6259 reason: RouteCloseReason,
6260 drain_timeout: Duration,
6261) -> Result<(), SuperviseError> {
6262 let Some(forwarding) = runtime.forwarding.as_ref() else {
6263 return Ok(());
6264 };
6265
6266 begin_forwarding_drain_with(
6267 forwarding,
6268 ForwardingDrainContext {
6269 spec,
6270 runtime,
6271 registry,
6272 scope: DrainScope::Active,
6273 },
6274 snapshot,
6275 enabled,
6276 reason,
6277 drain_timeout,
6278 )
6279 .await
6280}
6281
6282async fn begin_forwarding_drain_with(
6283 forwarding: &ForwardingTable,
6284 context: ForwardingDrainContext<'_>,
6285 snapshot: &SharedSnapshot,
6286 enabled: Option<bool>,
6287 reason: RouteCloseReason,
6288 drain_timeout: Duration,
6289) -> Result<(), SuperviseError> {
6290 let ForwardingDrainContext {
6291 spec,
6292 runtime,
6293 registry,
6294 scope,
6295 } = context;
6296 debug_assert_ne!(reason, RouteCloseReason::Crash);
6297 let terminal = matches!(reason, RouteCloseReason::Disable);
6298 let drain_started_at = Instant::now();
6299 let drain_deadline = drain_started_at + drain_timeout;
6300 let deadline_ms =
6301 unix_ms_now().saturating_add(u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX));
6302 let busy_gauges = match scope {
6303 DrainScope::Active => declared_busy_gauges(registry, &spec.module_id)?,
6304 DrainScope::Endpoint(endpoint) => {
6305 declared_busy_gauges_for_connection(registry, endpoint.connection_id)?
6306 }
6307 };
6308
6309 let drain_target = match scope {
6312 DrainScope::Active => forwarding.begin_module_drain(&spec.module_id, reason),
6313 DrainScope::Endpoint(endpoint) => forwarding.begin_endpoint_drain(endpoint, reason),
6314 }
6315 .map_err(SuperviseError::Forwarding)?;
6316 if scope == DrainScope::Active {
6317 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6318 state.state = ModuleState::Draining;
6319 if let Some(enabled) = enabled {
6320 state.enabled = enabled;
6321 }
6322 })?;
6323 }
6324
6325 if let Some(target) = drain_target.as_ref() {
6326 send_module_draining(&spec.module_id, reason, deadline_ms, target);
6327 let routes = forwarding
6328 .endpoint_routes(target.endpoint)
6329 .map_err(SuperviseError::Forwarding)?;
6330 let routes_notified = routes.len();
6331 crate::control::send_route_control_pushes(
6332 forwarding,
6333 routes.clone(),
6334 ClientControlPush::RouteClosing {
6335 module_id: spec.module_id.clone(),
6336 reason,
6337 },
6338 );
6339 send_route_goodbyes(forwarding, target.abandoned_bindings.clone());
6340
6341 let wait_result = wait_for_forwarding_quiescence(
6347 forwarding,
6348 &spec.module_id,
6349 runtime,
6350 target.endpoint,
6351 drain_deadline,
6352 &busy_gauges,
6353 scope,
6354 )
6355 .await;
6356 let drained = drained_after_quiescence_wait(&wait_result);
6357 if let Err(err) = &wait_result {
6358 error!(
6359 module_id = %spec.module_id,
6360 ?reason,
6361 error = %err,
6362 "forwarding quiescence wait failed after route.closing; forcing route.closed(drained: false) so the client is not left waiting on an unfulfilled promise"
6363 );
6364 } else if !drained {
6365 let holdouts = forwarding
6371 .endpoint_drain_holdouts(target.endpoint)
6372 .unwrap_or_default();
6373 warn!(
6374 module_id = %spec.module_id,
6375 waited = ?drain_timeout,
6376 ?reason,
6377 held_requests = holdouts.requests,
6378 held_routes = holdouts.routes,
6379 total_routes = holdouts.total_routes,
6380 top_connections = ?holdouts.top_connections,
6381 held = %holdouts
6384 .held
6385 .iter()
6386 .map(|(channel, corr)| format!("{channel}:{corr}"))
6387 .collect::<Vec<_>>()
6388 .join(","),
6389 "route drain timed out before request quiescence; forcing teardown"
6390 );
6391 }
6392 crate::control::send_route_control_pushes(
6393 forwarding,
6394 routes,
6395 ClientControlPush::RouteClosed {
6396 module_id: spec.module_id.clone(),
6397 reason,
6398 drained,
6399 abandoned: target.abandoned_bindings.len() as u32,
6400 excluded_subscriptions: target.excluded_subscriptions,
6401 terminal: Some(terminal),
6402 },
6403 );
6404 wait_result?;
6405
6406 let released_routes = match forwarding.release_module_endpoint_routes(target.endpoint) {
6412 Ok(routes) => routes,
6413 Err(err) => {
6414 warn!(
6415 module_id = %spec.module_id,
6416 ?reason,
6417 error = %err,
6418 "failed to release module endpoint routes after route.closed; module GOODBYE will still be sent"
6419 );
6420 send_module_goodbye(&spec.module_id, forwarding, target);
6421 return Err(SuperviseError::Forwarding(err));
6422 }
6423 };
6424 let route_goodbye_count = released_routes.len();
6425 send_route_goodbyes(forwarding, released_routes);
6426 send_module_goodbye(&spec.module_id, forwarding, target);
6427
6428 info!(
6434 module_id = %spec.module_id,
6435 ?reason,
6436 routes_notified,
6437 route_goodbyes = route_goodbye_count,
6438 abandoned_reservations = target.abandoned_bindings.len(),
6439 excluded_subscriptions = target.excluded_subscriptions,
6440 drained,
6441 "module drain complete; consumers notified via route.closing/route.closed pushes and per-route GOODBYE frames"
6442 );
6443 }
6444
6445 Ok(())
6446}
6447
6448async fn wait_for_registration_after_reload(
6451 registry: &Registry,
6452 module_id: &str,
6453 snapshot: &SharedSnapshot,
6454 child: &mut SupervisedChild,
6455 wait: Duration,
6456) -> Result<RegistrationWaitOutcome, SuperviseError> {
6457 wait_for_slot_registration(
6458 registry,
6459 crate::registry::RegistrationSlot::Active(module_id),
6460 module_id,
6461 snapshot,
6462 child,
6463 wait,
6464 )
6465 .await
6466}
6467
6468async fn wait_for_slot_registration(
6476 registry: &Registry,
6477 slot: crate::registry::RegistrationSlot<'_>,
6478 module_id: &str,
6479 snapshot: &SharedSnapshot,
6480 child: &mut SupervisedChild,
6481 wait: Duration,
6482) -> Result<RegistrationWaitOutcome, SuperviseError> {
6483 let deadline = Instant::now() + wait;
6484 loop {
6485 if registry
6486 .registration(slot)
6487 .map_err(SuperviseError::Registry)?
6488 .is_some()
6489 {
6490 return Ok(RegistrationWaitOutcome::Registered);
6491 }
6492
6493 let now = Instant::now();
6494 if now >= deadline {
6495 return Ok(RegistrationWaitOutcome::TimedOut);
6496 }
6497 let remaining = deadline.saturating_duration_since(now);
6498 let poll = remaining.min(REGISTRY_RELEASE_POLL);
6499
6500 tokio::select! {
6501 wait_result = child.wait() => {
6502 let status = wait_result.map_err(|source| SuperviseError::Wait {
6503 module_id: module_id.to_string(),
6504 source,
6505 })?;
6506 return Ok(RegistrationWaitOutcome::Exited(classify_reaped_child_exit(
6507 snapshot,
6508 child,
6509 &status,
6510 )));
6511 }
6512 _ = sleep(poll) => {}
6513 }
6514 }
6515}
6516
6517fn registration_failure_exit_report(mut exit_report: ExitReport) -> ExitReport {
6518 if exit_report.kind != ExitKind::DeliberateSeverance {
6521 exit_report.kind = ExitKind::Crash;
6522 }
6523 exit_report
6524}
6525
6526async fn handle_reload_child_registration_failure(
6527 spec: &ModuleSpec,
6528 runtime: &SupervisorRuntimeConfig,
6529 registry: &Registry,
6530 process_liveness: &SupervisorProcessLiveness,
6531 snapshot: &SharedSnapshot,
6532 child: &mut Option<SupervisedChild>,
6533 failure: ReloadRegistrationFailure,
6534) -> Result<(), SuperviseError> {
6535 let ReloadRegistrationFailure {
6536 exit_report,
6537 reason,
6538 } = failure;
6539 match on_child_exit(
6540 spec,
6541 runtime.restart_policy,
6542 registry,
6543 snapshot,
6544 &runtime.terminal_ring,
6545 &runtime.spawn_events,
6546 &runtime.child_roster,
6547 exit_report,
6548 )
6549 .await
6550 {
6551 NextAction::Stop {
6552 registration_released,
6553 } => {
6554 if registration_released {
6555 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6556 }
6557 }
6558 NextAction::Restart { schedule } => {
6559 let delay = schedule.map_or(runtime.restart_policy.delay_for_restart(0), |schedule| {
6560 schedule.delay
6561 });
6562 if let Some(schedule) = schedule {
6563 log_crash_respawn(&spec.module_id, schedule);
6564 }
6565 sleep(delay).await;
6566 if respawn_still_pending(snapshot) {
6570 if let Err(err) = wait_for_registration_release(
6571 registry,
6572 &spec.module_id,
6573 REGISTRY_RELEASE_TIMEOUT,
6574 )
6575 .await
6576 {
6577 fail_snapshot(snapshot, Some(&spec.module_id), None);
6578 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6579 return Err(SuperviseError::ReloadFailed {
6580 module_id: spec.module_id.clone(),
6581 reason: format!(
6582 "{reason}; registration did not release before policy retry: {err}"
6583 ),
6584 });
6585 }
6586 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6587 match spawn_and_mark_running(spec, runtime, snapshot) {
6588 Ok(next_child) => {
6589 *child = Some(next_child);
6590 }
6591 Err(err) => {
6592 fail_snapshot(snapshot, Some(&spec.module_id), None);
6593 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6594 return Err(SuperviseError::ReloadFailed {
6595 module_id: spec.module_id.clone(),
6596 reason: format!("{reason}; policy retry spawn failed: {err}"),
6597 });
6598 }
6599 }
6600 }
6601 }
6602 }
6603
6604 Err(SuperviseError::ReloadFailed {
6605 module_id: spec.module_id.clone(),
6606 reason,
6607 })
6608}
6609
6610async fn handle_reload_spawn_failure(
6611 spec: &ModuleSpec,
6612 runtime: &SupervisorRuntimeConfig,
6613 process_liveness: &SupervisorProcessLiveness,
6614 snapshot: &SharedSnapshot,
6615 child: &mut Option<SupervisedChild>,
6616 reason: String,
6617) -> Result<(), SuperviseError> {
6618 let mut should_retry = false;
6619 let now = Instant::now();
6620 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6621 clear_current_process_facts(state);
6622 if daemon_will_restart(state, &runtime.restart_policy, now) {
6623 state.record_crash_restart(&runtime.restart_policy, now);
6624 state.state = ModuleState::Restarting;
6625 should_retry = true;
6626 } else if state.enabled {
6627 state.state = ModuleState::Failed;
6628 } else {
6629 state.state = ModuleState::Disabled;
6630 }
6631 })?;
6632
6633 if should_retry {
6634 sleep(runtime.restart_policy.backoff).await;
6635 if respawn_still_pending(snapshot) {
6639 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6640 match spawn_and_mark_running(spec, runtime, snapshot) {
6641 Ok(next_child) => {
6642 *child = Some(next_child);
6643 }
6644 Err(err) => {
6645 fail_snapshot(snapshot, Some(&spec.module_id), None);
6646 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6647 return Err(SuperviseError::ReloadFailed {
6648 module_id: spec.module_id.clone(),
6649 reason: format!("{reason}; policy retry spawn failed: {err}"),
6650 });
6651 }
6652 }
6653 }
6654 } else {
6655 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6656 }
6657
6658 Err(SuperviseError::ReloadFailed {
6659 module_id: spec.module_id.clone(),
6660 reason,
6661 })
6662}
6663
6664fn control_flags() -> Flags {
6665 Flags::new(false, Priority::Passive, false)
6666}
6667
6668#[allow(clippy::too_many_arguments)]
6669async fn drain_optional_child(
6670 module_id: &str,
6671 protocol: ModuleProtocol,
6672 registry: &Registry,
6673 snapshot: &SharedSnapshot,
6674 terminal_ring: &Arc<Mutex<TerminalRing>>,
6675 spawn_events: &SpawnEventFeed,
6676 child: &mut Option<SupervisedChild>,
6677 drain_timeout: Duration,
6678 final_state: ModuleState,
6679 enabled: Option<bool>,
6680) -> Result<(), SuperviseError> {
6681 if let Some(child) = child.take() {
6682 drain_child_to_state(
6683 module_id,
6684 protocol,
6685 registry,
6686 snapshot,
6687 terminal_ring,
6688 spawn_events,
6689 child,
6690 drain_timeout,
6691 final_state,
6692 enabled,
6693 )
6694 .await
6695 } else {
6696 update_snapshot(snapshot, Some(module_id), |state| {
6697 state.state = final_state;
6698 if let Some(enabled) = enabled {
6699 state.enabled = enabled;
6700 }
6701 clear_current_process_facts(state);
6702 })?;
6703 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
6704 }
6705}
6706
6707#[allow(clippy::too_many_arguments)]
6708async fn drain_child_to_state(
6709 module_id: &str,
6710 protocol: ModuleProtocol,
6711 registry: &Registry,
6712 snapshot: &SharedSnapshot,
6713 terminal_ring: &Arc<Mutex<TerminalRing>>,
6714 spawn_events: &SpawnEventFeed,
6715 mut child: SupervisedChild,
6716 drain_timeout: Duration,
6717 final_state: ModuleState,
6718 enabled: Option<bool>,
6719) -> Result<(), SuperviseError> {
6720 update_snapshot(snapshot, Some(module_id), |state| {
6721 state.state = ModuleState::Draining;
6722 if let Some(enabled) = enabled {
6723 state.enabled = enabled;
6724 }
6725 })?;
6726
6727 if protocol == ModuleProtocol::None {
6733 request_graceful_stop(module_id, &child);
6734 }
6735
6736 let exit_report = match timeout(drain_timeout, child.wait()).await {
6737 Ok(Ok(status)) => classify_reaped_child_exit(snapshot, &child, &status),
6738 Ok(Err(source)) => {
6739 fail_snapshot(snapshot, Some(module_id), None);
6740 return Err(SuperviseError::Wait {
6741 module_id: module_id.to_string(),
6742 source,
6743 });
6744 }
6745 Err(_) => {
6746 child.start_kill().map_err(|source| {
6755 fail_snapshot(snapshot, Some(module_id), None);
6756 SuperviseError::Kill {
6757 module_id: module_id.to_string(),
6758 source,
6759 }
6760 })?;
6761 let status = child.wait().await.map_err(|source| {
6762 fail_snapshot(snapshot, Some(module_id), None);
6763 SuperviseError::Wait {
6764 module_id: module_id.to_string(),
6765 source,
6766 }
6767 })?;
6768 classify_reaped_child_exit(snapshot, &child, &status)
6769 }
6770 };
6771
6772 update_snapshot(snapshot, Some(module_id), |state| {
6773 state.state = final_state;
6774 if let Some(enabled) = enabled {
6775 state.enabled = enabled;
6776 }
6777 clear_current_process_facts(state);
6778 state.last_exit = Some(exit_report.clone());
6779 if exit_report.kind == ExitKind::DeliberateSeverance {
6780 state.lifetime_restarts += 1;
6781 }
6782 })?;
6783 record_terminal(
6784 module_id,
6785 terminal_ring,
6786 spawn_events,
6787 &exit_report,
6788 terminal_disposition(final_state),
6789 );
6790 child.drain_stderr(module_id).await;
6791
6792 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
6793}
6794
6795#[cfg(unix)]
6815fn request_graceful_stop(module_id: &str, child: &SupervisedChild) {
6816 let Some(pid) = child
6817 .id()
6818 .and_then(|pid| i32::try_from(pid).ok())
6819 .and_then(rustix::process::Pid::from_raw)
6820 else {
6821 debug!(
6822 module_id,
6823 "no pid to signal for protocol: none teardown; falling through to the drain wait"
6824 );
6825 return;
6826 };
6827 match rustix::process::kill_process(pid, rustix::process::Signal::TERM) {
6828 Ok(()) => debug!(module_id, "sent SIGTERM to protocol: none module"),
6829 Err(err) => debug!(
6830 module_id,
6831 error = %err,
6832 "SIGTERM to protocol: none module failed; the drain wait and kill still apply"
6833 ),
6834 }
6835}
6836
6837#[cfg(not(unix))]
6845fn request_graceful_stop(module_id: &str, _child: &SupervisedChild) {
6846 debug!(
6847 module_id,
6848 "no graceful stop signal exists on this platform; protocol: none teardown waits, then kills"
6849 );
6850}
6851
6852fn terminal_disposition(final_state: ModuleState) -> TerminalDisposition {
6853 match final_state {
6854 ModuleState::Stopped => TerminalDisposition::Stopped,
6855 ModuleState::Disabled => TerminalDisposition::Disabled,
6856 ModuleState::Restarting => TerminalDisposition::Restarting,
6857 ModuleState::Failed => TerminalDisposition::Failed,
6858 ModuleState::Starting
6859 | ModuleState::Running
6860 | ModuleState::Unresponsive
6861 | ModuleState::Draining => {
6862 unreachable!("terminal exits only finish in terminal or restarting states")
6863 }
6864 }
6865}
6866
6867async fn wait_for_registration_release(
6870 registry: &Registry,
6871 module_id: &str,
6872 wait: Duration,
6873) -> Result<(), SuperviseError> {
6874 wait_for_slot_registration_release(
6875 registry,
6876 crate::registry::RegistrationSlot::Active(module_id),
6877 wait,
6878 )
6879 .await
6880}
6881
6882async fn wait_for_slot_registration_release(
6890 registry: &Registry,
6891 slot: crate::registry::RegistrationSlot<'_>,
6892 wait: Duration,
6893) -> Result<(), SuperviseError> {
6894 let deadline = Instant::now() + wait;
6895 let mut release_events = registration_release_events().subscribe();
6896 let still_active = |registration: &crate::registry::ModuleRegistration| {
6897 SuperviseError::RegistrationStillActive {
6898 module_id: registration.manifest.module_id.clone(),
6899 waited: wait,
6900 }
6901 };
6902 loop {
6903 let _observed_generation = *release_events.borrow_and_update();
6904 let Some(registration) = registry
6905 .registration(slot)
6906 .map_err(SuperviseError::Registry)?
6907 else {
6908 return Ok(());
6909 };
6910
6911 let now = Instant::now();
6912 if now >= deadline {
6913 return Err(still_active(®istration));
6914 }
6915
6916 let remaining = deadline.saturating_duration_since(now);
6917 match timeout(remaining, release_events.changed()).await {
6918 Ok(Ok(())) | Ok(Err(_)) => {}
6919 Err(_) => return Err(still_active(®istration)),
6920 }
6921 }
6922}
6923
6924#[cfg(test)]
6925mod slot_registration_wait_tests {
6926 use super::*;
6927 use crate::registry::{ConnectionId, RegistrationSlot};
6928 use subc_protocol::manifest::ModuleManifest;
6929
6930 const INCUMBENT: u64 = 1;
6931 const CANDIDATE: u64 = 2;
6932
6933 fn swapped_registry() -> Arc<Registry> {
6934 let registry = Arc::new(Registry::default());
6935 let manifest = ModuleManifest::builder("m", "0.1.0").build();
6936 registry
6937 .register_with_control_ops(
6938 manifest.clone(),
6939 1,
6940 ConnectionId::new(INCUMBENT),
6941 Vec::new(),
6942 )
6943 .unwrap();
6944 registry
6945 .register_candidate_with_control_ops(
6946 manifest,
6947 1,
6948 ConnectionId::new(CANDIDATE),
6949 Vec::new(),
6950 )
6951 .unwrap();
6952 registry
6953 }
6954
6955 #[tokio::test]
6959 async fn incumbent_release_is_awaited_by_connection_not_by_module_id() {
6960 let registry = swapped_registry();
6961 registry.promote_candidate("m").unwrap().unwrap();
6962
6963 assert!(matches!(
6964 wait_for_registration_release(®istry, "m", Duration::from_millis(50)).await,
6965 Err(SuperviseError::RegistrationStillActive { .. })
6966 ));
6967
6968 assert!(matches!(
6970 wait_for_slot_registration_release(
6971 ®istry,
6972 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
6973 Duration::from_millis(50),
6974 )
6975 .await,
6976 Err(SuperviseError::RegistrationStillActive { .. })
6977 ));
6978
6979 let releaser = Arc::clone(®istry);
6980 let release = tokio::spawn(async move {
6981 sleep(Duration::from_millis(20)).await;
6982 releaser
6983 .deregister_connection(ConnectionId::new(INCUMBENT))
6984 .unwrap();
6985 notify_registration_release();
6986 });
6987 wait_for_slot_registration_release(
6988 ®istry,
6989 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
6990 Duration::from_secs(5),
6991 )
6992 .await
6993 .expect("the incumbent's own registration is released");
6994 release.await.unwrap();
6995 assert!(registry.get_module("m").unwrap().is_some());
6996 }
6997
6998 #[tokio::test]
7001 async fn candidate_slot_wait_ignores_the_incumbents_registration() {
7002 let registry = swapped_registry();
7003 assert!(matches!(
7004 wait_for_slot_registration_release(
7005 ®istry,
7006 RegistrationSlot::Candidate("m"),
7007 Duration::from_millis(50),
7008 )
7009 .await,
7010 Err(SuperviseError::RegistrationStillActive { .. })
7011 ));
7012 registry
7013 .deregister_connection(ConnectionId::new(CANDIDATE))
7014 .unwrap();
7015 wait_for_slot_registration_release(
7016 ®istry,
7017 RegistrationSlot::Candidate("m"),
7018 Duration::from_millis(50),
7019 )
7020 .await
7021 .expect("a candidate slot with no candidate is released");
7022 assert!(registry
7023 .registration(RegistrationSlot::Active("m"))
7024 .unwrap()
7025 .is_some());
7026 }
7027}
7028
7029fn classify_exit(status: &ExitStatus) -> ExitReport {
7030 ExitReport {
7031 kind: if status.success() {
7032 ExitKind::Clean
7033 } else {
7034 ExitKind::Crash
7035 },
7036 code: status.code(),
7037 signal: exit_signal(status),
7038 at_ms: unix_ms_now(),
7039 }
7040}
7041
7042fn wait_error_exit_report() -> ExitReport {
7048 ExitReport {
7049 kind: ExitKind::Crash,
7050 code: None,
7051 signal: None,
7052 at_ms: unix_ms_now(),
7053 }
7054}
7055
7056#[cfg(unix)]
7057fn exit_signal(status: &ExitStatus) -> Option<i32> {
7058 use std::os::unix::process::ExitStatusExt;
7059
7060 status.signal()
7061}
7062
7063#[cfg(not(unix))]
7064fn exit_signal(_status: &ExitStatus) -> Option<i32> {
7065 None
7066}
7067
7068fn reset_restart_count(snapshot: &SharedSnapshot, module_id: &str) -> Result<(), SuperviseError> {
7074 update_snapshot(snapshot, Some(module_id), |state| {
7075 state.clear_crash_restarts();
7076 })
7077}
7078
7079fn set_running(
7080 snapshot: &SharedSnapshot,
7081 child: &SupervisedChild,
7082 module_id: &str,
7083 spawn_events: &SpawnEventFeed,
7084) -> Result<(), SuperviseError> {
7085 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
7086 module_id: Some(module_id.to_string()),
7087 })?;
7088 state.spawn_generation = spawn_events.emit_spawned(module_id, child.pid, child.spawned_at_ms);
7089 state.in_alternate_slot = false;
7092 state.state = ModuleState::Running;
7093 state.enabled = true;
7094 state.process_alive = true;
7095 state.pid = child.id();
7096 state.spawned_at_ms = Some(child.spawned_at_ms);
7097 state.spawned_from = Some(child.spawned_from.clone());
7098 state.spawned_file_identity = child.spawned_file_identity;
7099 state.process_start_time = child.process_start_time;
7100 Ok(())
7101}
7102
7103fn clear_current_process_facts(state: &mut SupervisorSnapshot) {
7104 state.process_alive = false;
7105 state.pid = None;
7106 state.spawned_at_ms = None;
7107 state.spawned_from = None;
7108 state.spawned_file_identity = None;
7109 state.process_start_time = None;
7110 state.deliberate_severance = None;
7111}
7112
7113#[cfg(test)]
7114fn record_deliberate_severance(
7115 snapshot: &SharedSnapshot,
7116 identity: ProcessIdentity,
7117) -> Result<(), SuperviseError> {
7118 update_snapshot(snapshot, None, |state| {
7119 state.deliberate_severance = Some(identity);
7120 })
7121}
7122
7123fn apply_deliberate_severance_marker(
7124 snapshot: &SharedSnapshot,
7125 exited_identity: Option<ProcessIdentity>,
7126 mut exit_report: ExitReport,
7127) -> ExitReport {
7128 let marker = lock_snapshot(snapshot)
7129 .ok()
7130 .and_then(|mut state| state.deliberate_severance.take());
7131 if marker.is_some() && marker == exited_identity {
7132 exit_report.kind = ExitKind::DeliberateSeverance;
7133 }
7134 exit_report
7135}
7136
7137fn classify_reaped_child_exit(
7138 snapshot: &SharedSnapshot,
7139 child: &SupervisedChild,
7140 status: &ExitStatus,
7141) -> ExitReport {
7142 apply_deliberate_severance_marker(snapshot, child.process_identity(), classify_exit(status))
7143}
7144
7145fn fail_snapshot(
7146 snapshot: &SharedSnapshot,
7147 module_id: Option<&str>,
7148 last_exit: Option<ExitReport>,
7149) {
7150 if let Err(err) = update_snapshot(snapshot, module_id, |state| {
7151 state.state = ModuleState::Failed;
7152 clear_current_process_facts(state);
7153 if let Some(last_exit) = last_exit {
7154 state.last_exit = Some(last_exit);
7155 }
7156 }) {
7157 error!(error = %err, "failed to mark supervisor state failed");
7158 }
7159}
7160
7161fn update_snapshot(
7162 snapshot: &SharedSnapshot,
7163 module_id: Option<&str>,
7164 update: impl FnOnce(&mut SupervisorSnapshot),
7165) -> Result<(), SuperviseError> {
7166 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
7167 module_id: module_id.map(ToOwned::to_owned),
7168 })?;
7169 update(&mut state);
7170 Ok(())
7171}
7172
7173const SLOW_SNAPSHOT_LOCK_THRESHOLD: Duration = Duration::from_millis(250);
7174
7175fn lock_snapshot_for_control<'a>(
7176 snapshot: &'a SharedSnapshot,
7177 module_id: &str,
7178 caller: &'static str,
7179) -> Result<std::sync::MutexGuard<'a, SupervisorSnapshot>, SuperviseError> {
7180 let started_at = Instant::now();
7181 let guard = lock_snapshot(snapshot)?;
7182 let waited = started_at.elapsed();
7183 if waited >= SLOW_SNAPSHOT_LOCK_THRESHOLD {
7184 warn!(
7185 module_id = %module_id,
7186 waited_ms = waited.as_millis() as u64,
7187 caller = %caller,
7188 "slow snapshot lock"
7189 );
7190 }
7191 Ok(guard)
7192}
7193
7194fn lock_snapshot(
7195 snapshot: &SharedSnapshot,
7196) -> Result<std::sync::MutexGuard<'_, SupervisorSnapshot>, SuperviseError> {
7197 snapshot
7198 .lock()
7199 .map_err(|_| SuperviseError::StatePoisoned { module_id: None })
7200}
7201
7202#[cfg(test)]
7203mod terminal_history_tests {
7204 use std::{
7205 path::PathBuf,
7206 sync::Arc,
7207 time::{Duration, Instant},
7208 };
7209
7210 use tokio::time::sleep;
7211
7212 use super::{
7213 apply_deliberate_severance_marker, daemon_will_restart, drain_child_to_state,
7214 drained_after_quiescence_wait, handle_reload_spawn_failure, health_restart_child,
7215 lock_snapshot, on_child_exit, record_deliberate_severance, record_wait_error_terminal,
7216 reset_restart_count, spawn_and_mark_running, update_snapshot, wait_error_exit_report,
7217 ExitKind, ExitReport, ModuleProtocol, ModuleSpec, ModuleState, NextAction, ProcessIdentity,
7218 RestartPolicy, SpawnEventKind, SuperviseError, SupervisedModule, Supervisor,
7219 SupervisorHandle, SupervisorHealthStatus, SupervisorSnapshot,
7220 };
7221 use super::Instant as ClockInstant;
7226 use crate::{
7227 registry::Registry,
7228 terminal_ring::{TerminalRing, TerminalRingConfig},
7229 };
7230 use std::sync::Mutex;
7231 use subc_control::TerminalDisposition;
7232
7233 fn fake_aft_stub_path() -> PathBuf {
7238 let mut path = std::env::current_exe().expect("current_exe available in tests");
7239 path.pop();
7240 path.pop();
7241 path.push(if cfg!(windows) {
7242 "fake-aft-stub.exe"
7243 } else {
7244 "fake-aft-stub"
7245 });
7246 assert!(
7247 path.exists(),
7248 "fake-aft-stub not built at {}: run `cargo test -p subc-core` (which builds \
7249 [[bin]] targets) rather than `cargo test -p subc-core --lib` (which does not)",
7250 path.display()
7251 );
7252 path
7253 }
7254
7255 #[test]
7256 fn reserved_never_spawned_refuses_every_hello() {
7257 let supervisor = SupervisorHandle::default();
7262 supervisor.apply_identity_configuration(&ModuleSpec {
7263 module_id: "never-spawned".to_string(),
7264 program: PathBuf::from("/usr/bin/false"),
7265 args: Vec::new(),
7266 env: Vec::new(),
7267 reserved: true,
7268 reserved_prefixes: Vec::new(),
7269 protocol: ModuleProtocol::Subc,
7270 overlap: Default::default(),
7271 });
7272 assert!(
7273 supervisor
7274 .reserved_hello_rejection("never-spawned", Some("any-forged-nonce"))
7275 .is_some(),
7276 "forged nonce must refuse on a reserved never-spawned id"
7277 );
7278 assert!(
7279 supervisor
7280 .reserved_hello_rejection("never-spawned", None)
7281 .is_some(),
7282 "absent nonce must refuse on a reserved never-spawned id"
7283 );
7284 supervisor.set_spawn_nonce("never-spawned", "minted".to_string());
7286 supervisor.apply_identity_configuration(&ModuleSpec {
7287 module_id: "never-spawned".to_string(),
7288 program: PathBuf::from("/usr/bin/false"),
7289 args: Vec::new(),
7290 env: Vec::new(),
7291 reserved: true,
7292 reserved_prefixes: Vec::new(),
7293 protocol: ModuleProtocol::Subc,
7294 overlap: Default::default(),
7295 });
7296 assert!(supervisor
7297 .reserved_hello_rejection("never-spawned", Some("minted"))
7298 .is_none());
7299 assert!(supervisor
7300 .reserved_hello_rejection("never-spawned", Some("forged"))
7301 .is_some());
7302 }
7303
7304 fn seed_crash_restarts(state: &mut SupervisorSnapshot, count: u32) {
7307 let now = ClockInstant::now();
7308 for _ in 0..count {
7309 state.crash_restarts.push_back(now);
7310 }
7311 }
7312
7313 fn age_oldest_crash_restart_out_of_window(state: &mut SupervisorSnapshot, window: Duration) {
7317 let aged = state
7318 .crash_restarts
7319 .front()
7320 .expect("a crash restart must be recorded before it can be aged")
7321 .checked_sub(window + Duration::from_secs(1))
7322 .expect("the test clock is far enough from its origin to age an instant");
7323 state.crash_restarts[0] = aged;
7324 }
7325
7326 fn snapshot_with_restarts(enabled: bool, count: u32) -> SupervisorSnapshot {
7327 let mut state = SupervisorSnapshot::new(ModuleState::Running, enabled);
7328 seed_crash_restarts(&mut state, count);
7329 state
7330 }
7331
7332 #[test]
7333 fn daemon_owned_recovery_predicate_uses_the_pre_increment_budget() {
7334 let policy = RestartPolicy::new(3, Duration::ZERO);
7335 let now = ClockInstant::now();
7336 assert!(daemon_will_restart(
7337 &mut snapshot_with_restarts(true, 2),
7338 &policy,
7339 now
7340 ));
7341 assert!(!daemon_will_restart(
7342 &mut snapshot_with_restarts(true, 3),
7343 &policy,
7344 now
7345 ));
7346 assert!(!daemon_will_restart(
7347 &mut snapshot_with_restarts(false, 0),
7348 &policy,
7349 now
7350 ));
7351 }
7352
7353 #[test]
7354 fn crash_restart_backoff_escalates_with_in_window_count() {
7355 let policy = RestartPolicy::new(4, Duration::from_millis(100))
7356 .with_max_backoff(Duration::from_secs(30));
7357 let now = ClockInstant::now();
7358 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7359 let schedules = (0..4)
7360 .map(|_| {
7361 state
7362 .next_crash_restart(&policy, now)
7363 .expect("the test policy allows four crash restarts")
7364 })
7365 .collect::<Vec<_>>();
7366
7367 assert_eq!(
7368 schedules
7369 .iter()
7370 .map(|schedule| schedule.restart_in_window)
7371 .collect::<Vec<_>>(),
7372 vec![0, 1, 2, 3]
7373 );
7374 assert_eq!(
7375 schedules
7376 .iter()
7377 .map(|schedule| schedule.delay)
7378 .collect::<Vec<_>>(),
7379 vec![
7380 Duration::from_millis(100),
7381 Duration::from_secs(1),
7382 Duration::from_secs(10),
7383 Duration::from_secs(30),
7384 ]
7385 );
7386 }
7387
7388 #[test]
7389 fn crash_restart_backoff_resets_after_ring_clear() {
7390 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7391 let now = ClockInstant::now();
7392 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7393 assert_eq!(
7394 state.next_crash_restart(&policy, now).unwrap().delay,
7395 Duration::from_millis(100)
7396 );
7397 assert_eq!(
7398 state.next_crash_restart(&policy, now).unwrap().delay,
7399 Duration::from_secs(1)
7400 );
7401
7402 state.clear_crash_restarts();
7403 let schedule = state
7404 .next_crash_restart(&policy, now)
7405 .expect("a cleared ring must allow another restart");
7406 assert_eq!(schedule.restart_in_window, 0);
7407 assert_eq!(schedule.delay, Duration::from_millis(100));
7408 }
7409
7410 #[test]
7411 fn crash_restart_backoff_ignores_aged_restarts() {
7412 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7413 let now = ClockInstant::now();
7414 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7415 state
7416 .next_crash_restart(&policy, now)
7417 .expect("the first restart is allowed");
7418 state
7419 .next_crash_restart(&policy, now)
7420 .expect("the second restart is allowed");
7421 state.crash_restarts[0] = now
7422 .checked_sub(policy.window + Duration::from_secs(1))
7423 .expect("the fake clock can age a restart past the window");
7424
7425 let schedule = state
7426 .next_crash_restart(&policy, now)
7427 .expect("an aged restart must release its slot");
7428 assert_eq!(schedule.restart_in_window, 1);
7429 assert_eq!(schedule.delay, Duration::from_secs(1));
7430 assert_eq!(state.crash_restarts.len(), 2);
7431 }
7432
7433 #[test]
7437 fn a_budget_spent_before_the_window_no_longer_refuses() {
7438 let policy = RestartPolicy::new(3, Duration::ZERO);
7439 let mut state = snapshot_with_restarts(true, 3);
7440 let now = ClockInstant::now();
7441 assert!(!daemon_will_restart(&mut state, &policy, now));
7442
7443 assert!(daemon_will_restart(
7444 &mut state,
7445 &policy,
7446 now + policy.window + Duration::from_secs(1)
7447 ));
7448 assert!(
7449 state.crash_restarts.is_empty(),
7450 "reading the budget must drop the instants that left the window"
7451 );
7452 }
7453
7454 fn module_with_recovery_snapshot(
7455 state: ModuleState,
7456 enabled: bool,
7457 restart_count: u32,
7458 ) -> SupervisedModule {
7459 let registry = Arc::new(Registry::default());
7460 let supervisor =
7461 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(3, Duration::ZERO));
7462 let module = supervisor
7463 .spawn(ModuleSpec {
7464 module_id: "recovery-snapshot".to_string(),
7465 program: fake_aft_stub_path(),
7466 args: Vec::new(),
7467 env: Vec::new(),
7468 reserved: false,
7469 reserved_prefixes: Vec::new(),
7470 protocol: ModuleProtocol::Subc,
7471 overlap: Default::default(),
7472 })
7473 .unwrap();
7474 update_snapshot(
7475 &module.inner.snapshot,
7476 Some("recovery-snapshot"),
7477 |snapshot| {
7478 snapshot.state = state;
7479 snapshot.enabled = enabled;
7480 seed_crash_restarts(snapshot, restart_count);
7481 },
7482 )
7483 .unwrap();
7484 module
7485 }
7486
7487 #[cfg(target_os = "linux")]
7488 #[tokio::test]
7489 async fn no_cgroup_placement_does_not_block_fake_aft_stub_spawn() {
7490 let supervisor = Supervisor::new(Arc::new(Registry::default()), RestartPolicy::default())
7491 .with_cgroup_placement(None);
7492 let result = supervisor.spawn(ModuleSpec {
7493 module_id: "no-cgroup-placement".to_string(),
7494 program: fake_aft_stub_path(),
7495 args: Vec::new(),
7496 env: Vec::new(),
7497 reserved: false,
7498 reserved_prefixes: Vec::new(),
7499 protocol: ModuleProtocol::Subc,
7500 overlap: Default::default(),
7501 });
7502
7503 assert!(
7504 result.is_ok(),
7505 "no delegation must not turn an otherwise valid spawn into a failure: {result:?}"
7506 );
7507 }
7508
7509 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7510 async fn undecided_snapshot_uses_shared_restart_predicate() {
7511 assert!(module_with_recovery_snapshot(ModuleState::Running, true, 2)
7512 .will_recover_after_connection_loss()
7513 .unwrap());
7514 assert!(
7515 !module_with_recovery_snapshot(ModuleState::Running, true, 3)
7516 .will_recover_after_connection_loss()
7517 .unwrap()
7518 );
7519 }
7520
7521 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7522 async fn restarting_snapshot_at_exhausted_budget_is_non_terminal() {
7523 assert!(
7524 module_with_recovery_snapshot(ModuleState::Restarting, true, 3)
7525 .will_recover_after_connection_loss()
7526 .unwrap()
7527 );
7528 }
7529
7530 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7531 async fn terminal_phase_snapshots_are_terminal_before_budget_exhaustion() {
7532 assert!(!module_with_recovery_snapshot(ModuleState::Failed, true, 0)
7533 .will_recover_after_connection_loss()
7534 .unwrap());
7535 assert!(
7536 !module_with_recovery_snapshot(ModuleState::Disabled, true, 0)
7537 .will_recover_after_connection_loss()
7538 .unwrap()
7539 );
7540 }
7541
7542 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7543 async fn warming_snapshot_is_limited_to_startup_phases() {
7544 for state in [
7545 ModuleState::Starting,
7546 ModuleState::Running,
7547 ModuleState::Restarting,
7548 ] {
7549 assert!(
7550 module_with_recovery_snapshot(state, true, 0)
7551 .is_warming()
7552 .unwrap(),
7553 "{state:?} should be warming"
7554 );
7555 }
7556 for state in [
7557 ModuleState::Unresponsive,
7558 ModuleState::Draining,
7559 ModuleState::Stopped,
7560 ModuleState::Failed,
7561 ModuleState::Disabled,
7562 ] {
7563 assert!(
7564 !module_with_recovery_snapshot(state, true, 0)
7565 .is_warming()
7566 .unwrap(),
7567 "{state:?} should not be warming"
7568 );
7569 }
7570 }
7571
7572 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7573 async fn terminal_history_survives_respawn_and_keeps_both_crashes_in_order() {
7574 let registry = Arc::new(Registry::default());
7575 let supervisor =
7576 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(1, Duration::ZERO));
7577 let module = supervisor
7578 .spawn(ModuleSpec {
7579 module_id: "terminal-history".to_string(),
7580 program: fake_aft_stub_path(),
7581 args: Vec::new(),
7582 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7583 reserved: false,
7584 reserved_prefixes: Vec::new(),
7585 protocol: ModuleProtocol::Subc,
7586 overlap: Default::default(),
7587 })
7588 .unwrap();
7589
7590 let deadline = Instant::now() + Duration::from_secs(5);
7591 loop {
7592 let history = module.terminal_history();
7593 if history.entries.len() == 2 {
7594 assert_eq!(module.status().unwrap().state, ModuleState::Failed);
7595 assert_eq!(history.dropped, 0);
7596 assert_eq!(
7597 history
7598 .entries
7599 .iter()
7600 .map(|entry| entry.exit_code)
7601 .collect::<Vec<_>>(),
7602 vec![Some(23), Some(23)]
7603 );
7604 assert!(history.entries[0].at_ms <= history.entries[1].at_ms);
7605 return;
7606 }
7607 assert!(
7608 Instant::now() < deadline,
7609 "module did not retain two terminal exits: {history:?}"
7610 );
7611 sleep(Duration::from_millis(10)).await;
7612 }
7613 }
7614
7615 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7619 async fn disable_during_crash_backoff_cancels_pending_respawn() {
7620 let backoff = Duration::from_secs(2);
7621 let supervisor = Supervisor::new(
7622 Arc::new(Registry::default()),
7623 RestartPolicy::new(10, backoff),
7624 );
7625 let module = supervisor
7626 .spawn(ModuleSpec {
7627 module_id: "disable-during-backoff".to_string(),
7628 program: fake_aft_stub_path(),
7629 args: Vec::new(),
7630 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7631 reserved: false,
7632 reserved_prefixes: Vec::new(),
7633 protocol: ModuleProtocol::Subc,
7634 overlap: Default::default(),
7635 })
7636 .unwrap();
7637
7638 let deadline = Instant::now() + Duration::from_secs(5);
7640 loop {
7641 if module.status().unwrap().state == ModuleState::Restarting {
7642 break;
7643 }
7644 assert!(
7645 Instant::now() < deadline,
7646 "module never entered the crash backoff"
7647 );
7648 sleep(Duration::from_millis(10)).await;
7649 }
7650
7651 let started = Instant::now();
7652 module.set_enabled(false).await.unwrap();
7653 let waited = started.elapsed();
7654
7655 assert!(
7656 waited < backoff / 2,
7657 "disable waited {waited:?} behind the {backoff:?} crash backoff; the operator command must preempt the pending respawn"
7658 );
7659 assert_eq!(module.status().unwrap().state, ModuleState::Disabled);
7660
7661 sleep(backoff + Duration::from_millis(500)).await;
7663 let status = module.status().unwrap();
7664 assert_eq!(status.state, ModuleState::Disabled);
7665 assert_eq!(
7666 status.spawn_generation, 1,
7667 "module respawned after the operator disabled it"
7668 );
7669 }
7670
7671 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7675 async fn every_restart_increment_path_advances_lifetime_count() {
7676 let supervisor = Supervisor::new(
7677 Arc::new(Registry::default()),
7678 RestartPolicy::new(1, Duration::ZERO),
7679 );
7680 let runtime = supervisor.runtime_config();
7681 let spec = ModuleSpec {
7682 module_id: "lifetime-increment-path".to_string(),
7683 program: PathBuf::from("/unused/lifetime-increment-path"),
7684 args: Vec::new(),
7685 env: Vec::new(),
7686 reserved: false,
7687 reserved_prefixes: Vec::new(),
7688 protocol: ModuleProtocol::Subc,
7689 overlap: Default::default(),
7690 };
7691
7692 let crash_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7693 assert!(matches!(
7694 on_child_exit(
7695 &spec,
7696 runtime.restart_policy,
7697 &supervisor.registry,
7698 &crash_snapshot,
7699 &runtime.terminal_ring,
7700 &runtime.spawn_events,
7701 &runtime.child_roster,
7702 ExitReport {
7703 kind: ExitKind::Crash,
7704 code: Some(1),
7705 signal: None,
7706 at_ms: 1,
7707 },
7708 )
7709 .await,
7710 NextAction::Restart { schedule: _ }
7711 ));
7712 let (crash_restarts, crash_lifetime) = {
7713 let state = lock_snapshot(&crash_snapshot).unwrap();
7714 (state.crash_restarts.len(), state.lifetime_restarts)
7715 };
7716 assert_eq!(crash_restarts, 1);
7717 assert_eq!(crash_lifetime, 1);
7718
7719 let health_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7720 let mut health_child = None;
7721 assert!(matches!(
7722 health_restart_child(
7723 &spec,
7724 &runtime,
7725 &supervisor.registry,
7726 &supervisor.process_liveness,
7727 &health_snapshot,
7728 &mut health_child,
7729 SupervisorHealthStatus::Failing,
7730 None,
7731 2,
7732 )
7733 .await,
7734 Err(SuperviseError::Spawn { .. })
7735 ));
7736 let (health_restarts, health_lifetime) = {
7737 let state = lock_snapshot(&health_snapshot).unwrap();
7738 (state.crash_restarts.len(), state.lifetime_restarts)
7739 };
7740 assert_eq!(health_restarts, 1);
7741 assert_eq!(health_lifetime, 1);
7742
7743 let reload_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7744 let mut reload_child = None;
7745 assert!(matches!(
7746 handle_reload_spawn_failure(
7747 &spec,
7748 &runtime,
7749 &supervisor.process_liveness,
7750 &reload_snapshot,
7751 &mut reload_child,
7752 "forced reload spawn failure".to_string(),
7753 )
7754 .await,
7755 Err(SuperviseError::ReloadFailed { .. })
7756 ));
7757 let (reload_restarts, reload_lifetime) = {
7758 let state = lock_snapshot(&reload_snapshot).unwrap();
7759 (state.crash_restarts.len(), state.lifetime_restarts)
7760 };
7761 assert_eq!(reload_restarts, 1);
7762 assert_eq!(reload_lifetime, 1);
7763 }
7764
7765 #[tokio::test]
7766 async fn deliberately_severed_live_child_records_lifetime_without_spending_restart_budget() {
7767 let supervisor = Supervisor::new(
7768 Arc::new(Registry::default()),
7769 RestartPolicy::new(3, Duration::ZERO),
7770 );
7771 let runtime = supervisor.runtime_config();
7772 let spec = ModuleSpec {
7773 module_id: "deliberately-severed".to_string(),
7774 program: PathBuf::from("/unused/deliberately-severed"),
7775 args: Vec::new(),
7776 env: Vec::new(),
7777 reserved: false,
7778 reserved_prefixes: Vec::new(),
7779 protocol: ModuleProtocol::Subc,
7780 overlap: Default::default(),
7781 };
7782 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7783 let process = ProcessIdentity {
7784 pid: 41,
7785 start_time: 101,
7786 };
7787 record_deliberate_severance(&snapshot, process).unwrap();
7788 let exit_report = apply_deliberate_severance_marker(
7789 &snapshot,
7790 Some(process),
7791 ExitReport {
7792 kind: ExitKind::Crash,
7793 code: Some(1),
7794 signal: None,
7795 at_ms: 1,
7796 },
7797 );
7798 assert_eq!(exit_report.kind, ExitKind::DeliberateSeverance);
7799
7800 assert!(matches!(
7801 on_child_exit(
7802 &spec,
7803 runtime.restart_policy,
7804 &supervisor.registry,
7805 &snapshot,
7806 &runtime.terminal_ring,
7807 &runtime.spawn_events,
7808 &runtime.child_roster,
7809 exit_report,
7810 )
7811 .await,
7812 NextAction::Restart { schedule: _ }
7813 ));
7814 let state = lock_snapshot(&snapshot).unwrap();
7815 assert_eq!(state.lifetime_restarts, 1);
7816 assert_eq!(state.crash_restarts.len(), 0);
7817 }
7818
7819 #[tokio::test]
7820 async fn genuine_crash_spends_restart_budget_and_records_lifetime() {
7821 let supervisor = Supervisor::new(
7822 Arc::new(Registry::default()),
7823 RestartPolicy::new(3, Duration::ZERO),
7824 );
7825 let runtime = supervisor.runtime_config();
7826 let spec = ModuleSpec {
7827 module_id: "genuine-crash".to_string(),
7828 program: PathBuf::from("/unused/genuine-crash"),
7829 args: Vec::new(),
7830 env: Vec::new(),
7831 reserved: false,
7832 reserved_prefixes: Vec::new(),
7833 protocol: ModuleProtocol::Subc,
7834 overlap: Default::default(),
7835 };
7836 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7837
7838 assert!(matches!(
7839 on_child_exit(
7840 &spec,
7841 runtime.restart_policy,
7842 &supervisor.registry,
7843 &snapshot,
7844 &runtime.terminal_ring,
7845 &runtime.spawn_events,
7846 &runtime.child_roster,
7847 ExitReport {
7848 kind: ExitKind::Crash,
7849 code: Some(1),
7850 signal: None,
7851 at_ms: 1,
7852 },
7853 )
7854 .await,
7855 NextAction::Restart { schedule: _ }
7856 ));
7857 let state = lock_snapshot(&snapshot).unwrap();
7858 assert_eq!(state.lifetime_restarts, 1);
7859 assert_eq!(state.crash_restarts.len(), 1);
7860 }
7861
7862 fn crash_exit_report(at_ms: u64) -> ExitReport {
7863 ExitReport {
7864 kind: ExitKind::Crash,
7865 code: Some(1),
7866 signal: None,
7867 at_ms,
7868 }
7869 }
7870
7871 fn windowed_crash_spec(module_id: &str) -> ModuleSpec {
7872 ModuleSpec {
7873 module_id: module_id.to_string(),
7874 program: PathBuf::from("/unused").join(module_id),
7875 args: Vec::new(),
7876 env: Vec::new(),
7877 reserved: false,
7878 reserved_prefixes: Vec::new(),
7879 protocol: ModuleProtocol::Subc,
7880 overlap: Default::default(),
7881 }
7882 }
7883
7884 #[tokio::test]
7890 async fn three_crashes_inside_the_window_stop_the_module_and_name_the_window() {
7891 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::ERROR);
7892 let supervisor = Supervisor::new(
7893 Arc::new(Registry::default()),
7894 RestartPolicy::new(2, Duration::ZERO),
7895 );
7896 let runtime = supervisor.runtime_config();
7897 let spec = windowed_crash_spec("crash-loop-in-window");
7898 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7899
7900 for attempt in 1..=2 {
7901 assert!(
7902 matches!(
7903 on_child_exit(
7904 &spec,
7905 runtime.restart_policy,
7906 &supervisor.registry,
7907 &snapshot,
7908 &runtime.terminal_ring,
7909 &runtime.spawn_events,
7910 &runtime.child_roster,
7911 crash_exit_report(attempt),
7912 )
7913 .await,
7914 NextAction::Restart { schedule: _ }
7915 ),
7916 "crash {attempt} is inside the budget and must respawn"
7917 );
7918 }
7919
7920 assert!(matches!(
7921 on_child_exit(
7922 &spec,
7923 runtime.restart_policy,
7924 &supervisor.registry,
7925 &snapshot,
7926 &runtime.terminal_ring,
7927 &runtime.spawn_events,
7928 &runtime.child_roster,
7929 crash_exit_report(3),
7930 )
7931 .await,
7932 NextAction::Stop { .. }
7933 ));
7934
7935 {
7936 let state = lock_snapshot(&snapshot).unwrap();
7937 assert_eq!(state.state, ModuleState::Failed);
7938 assert_eq!(state.crash_restarts.len(), 2);
7939 assert_eq!(state.lifetime_restarts, 2);
7940 }
7941
7942 let history = runtime
7943 .terminal_ring
7944 .lock()
7945 .expect("terminal ring is not poisoned")
7946 .snapshot();
7947 let last = history
7948 .entries
7949 .last()
7950 .expect("the refused crash is retained");
7951 assert_eq!(last.disposition, TerminalDisposition::Failed);
7952 assert_eq!(
7953 last.disposition_detail.as_deref(),
7954 Some("crash budget exhausted: max_restarts=2 within window_secs=600")
7955 );
7956
7957 let captured = crate::router::test_log::captured_logs(&logs);
7958 assert!(
7959 captured.contains("crash budget exhausted: max_restarts=2 within window_secs=600"),
7960 "the stop must be logged with its window: {captured}"
7961 );
7962 }
7963
7964 #[tokio::test]
7972 async fn a_crash_older_than_the_window_frees_its_slot_for_a_later_crash() {
7973 let supervisor = Supervisor::new(
7974 Arc::new(Registry::default()),
7975 RestartPolicy::new(2, Duration::ZERO),
7976 );
7977 let runtime = supervisor.runtime_config();
7978 let spec = windowed_crash_spec("crash-across-windows");
7979 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7980
7981 for attempt in 1..=2 {
7982 assert!(matches!(
7983 on_child_exit(
7984 &spec,
7985 runtime.restart_policy,
7986 &supervisor.registry,
7987 &snapshot,
7988 &runtime.terminal_ring,
7989 &runtime.spawn_events,
7990 &runtime.child_roster,
7991 crash_exit_report(attempt),
7992 )
7993 .await,
7994 NextAction::Restart { schedule: _ }
7995 ));
7996 }
7997
7998 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
8001 age_oldest_crash_restart_out_of_window(state, runtime.restart_policy.window);
8002 })
8003 .unwrap();
8004
8005 assert!(
8006 matches!(
8007 on_child_exit(
8008 &spec,
8009 runtime.restart_policy,
8010 &supervisor.registry,
8011 &snapshot,
8012 &runtime.terminal_ring,
8013 &runtime.spawn_events,
8014 &runtime.child_roster,
8015 crash_exit_report(3),
8016 )
8017 .await,
8018 NextAction::Restart { schedule: _ }
8019 ),
8020 "a crash older than the window must not hold a budget slot"
8021 );
8022
8023 let state = lock_snapshot(&snapshot).unwrap();
8024 assert_eq!(state.state, ModuleState::Restarting);
8025 assert_eq!(
8026 state.crash_restarts.len(),
8027 2,
8028 "the aged instant is dropped and the new one takes its place"
8029 );
8030 assert_eq!(
8031 state.lifetime_restarts, 3,
8032 "the ledger counts every restart, including the ones the window forgot"
8033 );
8034 }
8035
8036 #[tokio::test]
8041 async fn an_operator_restart_clears_the_ring_and_leaves_the_ledger_alone() {
8042 let supervisor = Supervisor::new(
8043 Arc::new(Registry::default()),
8044 RestartPolicy::new(2, Duration::ZERO),
8045 );
8046 let runtime = supervisor.runtime_config();
8047 let spec = windowed_crash_spec("operator-cleared-budget");
8048 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8049
8050 for attempt in 1..=2 {
8051 assert!(matches!(
8052 on_child_exit(
8053 &spec,
8054 runtime.restart_policy,
8055 &supervisor.registry,
8056 &snapshot,
8057 &runtime.terminal_ring,
8058 &runtime.spawn_events,
8059 &runtime.child_roster,
8060 crash_exit_report(attempt),
8061 )
8062 .await,
8063 NextAction::Restart { schedule: _ }
8064 ));
8065 }
8066
8067 reset_restart_count(&snapshot, &spec.module_id).unwrap();
8068 {
8069 let state = lock_snapshot(&snapshot).unwrap();
8070 assert!(
8071 state.crash_restarts.is_empty(),
8072 "an operator restart returns the full budget"
8073 );
8074 assert_eq!(
8075 state.lifetime_restarts, 2,
8076 "clearing the budget must not unmake the crashes"
8077 );
8078 }
8079
8080 assert!(
8081 matches!(
8082 on_child_exit(
8083 &spec,
8084 runtime.restart_policy,
8085 &supervisor.registry,
8086 &snapshot,
8087 &runtime.terminal_ring,
8088 &runtime.spawn_events,
8089 &runtime.child_roster,
8090 crash_exit_report(3),
8091 )
8092 .await,
8093 NextAction::Restart { schedule: _ }
8094 ),
8095 "the cleared budget must be spendable again"
8096 );
8097 let state = lock_snapshot(&snapshot).unwrap();
8098 assert_eq!(state.crash_restarts.len(), 1);
8099 assert_eq!(state.lifetime_restarts, 3);
8100 }
8101
8102 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
8103 async fn severance_marker_for_a_dead_child_does_not_label_its_successor() {
8104 let severed = ProcessIdentity {
8105 pid: 41,
8106 start_time: 101,
8107 };
8108 let successor = ProcessIdentity {
8109 pid: 41,
8110 start_time: 202,
8111 };
8112 let module = module_with_recovery_snapshot(ModuleState::Running, true, 0);
8113 update_snapshot(&module.inner.snapshot, Some("recovery-snapshot"), |state| {
8114 state.pid = Some(successor.pid);
8115 state.process_start_time = Some(successor.start_time);
8116 })
8117 .unwrap();
8118 assert!(!module.record_deliberate_severance(severed).unwrap());
8119
8120 let exit_report = apply_deliberate_severance_marker(
8121 &module.inner.snapshot,
8122 Some(successor),
8123 ExitReport {
8124 kind: ExitKind::Crash,
8125 code: Some(1),
8126 signal: None,
8127 at_ms: 1,
8128 },
8129 );
8130
8131 assert_eq!(exit_report.kind, ExitKind::Crash);
8132 }
8133
8134 #[tokio::test]
8135 async fn drain_reap_marks_deliberate_severance_and_records_lifetime_without_budget() {
8136 let registry = Registry::default();
8137 let supervisor = Supervisor::new(
8138 Arc::new(Registry::default()),
8139 RestartPolicy::new(3, Duration::ZERO),
8140 );
8141 let runtime = supervisor.runtime_config();
8142 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8143 let spec = ModuleSpec {
8144 module_id: "drain-deliberate-severance".to_string(),
8145 program: fake_aft_stub_path(),
8146 args: Vec::new(),
8147 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
8148 reserved: false,
8149 reserved_prefixes: Vec::new(),
8150 protocol: ModuleProtocol::Subc,
8151 overlap: Default::default(),
8152 };
8153 let mut child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
8154 let process = ProcessIdentity {
8155 pid: 41,
8156 start_time: 101,
8157 };
8158 child.process_identity = Some(process);
8159 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
8160 state.pid = Some(process.pid);
8161 state.process_start_time = Some(process.start_time);
8162 })
8163 .unwrap();
8164 record_deliberate_severance(&snapshot, process).unwrap();
8165
8166 drain_child_to_state(
8167 &spec.module_id,
8168 spec.protocol,
8169 ®istry,
8170 &snapshot,
8171 &runtime.terminal_ring,
8172 &runtime.spawn_events,
8173 child,
8174 Duration::from_secs(1),
8175 ModuleState::Stopped,
8176 Some(false),
8177 )
8178 .await
8179 .unwrap();
8180
8181 let state = lock_snapshot(&snapshot).unwrap();
8182 assert_eq!(
8183 state.last_exit.as_ref().map(|exit| exit.kind),
8184 Some(ExitKind::DeliberateSeverance)
8185 );
8186 assert_eq!(state.lifetime_restarts, 1);
8187 assert_eq!(state.crash_restarts.len(), 0);
8188 drop(state);
8189 let history = runtime.terminal_ring.lock().unwrap().snapshot();
8190 assert_eq!(
8191 history.entries[0].exit_kind,
8192 subc_control::TerminalExitKind::DeliberateSeverance
8193 );
8194 }
8195
8196 #[tokio::test]
8197 async fn ordinary_drain_reap_does_not_record_a_lifetime_restart() {
8198 let registry = Registry::default();
8199 let supervisor = Supervisor::new(
8200 Arc::new(Registry::default()),
8201 RestartPolicy::new(3, Duration::ZERO),
8202 );
8203 let runtime = supervisor.runtime_config();
8204 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8205 let spec = ModuleSpec {
8206 module_id: "ordinary-drain".to_string(),
8207 program: fake_aft_stub_path(),
8208 args: Vec::new(),
8209 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
8210 reserved: false,
8211 reserved_prefixes: Vec::new(),
8212 protocol: ModuleProtocol::Subc,
8213 overlap: Default::default(),
8214 };
8215 let child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
8216
8217 drain_child_to_state(
8218 &spec.module_id,
8219 spec.protocol,
8220 ®istry,
8221 &snapshot,
8222 &runtime.terminal_ring,
8223 &runtime.spawn_events,
8224 child,
8225 Duration::from_secs(1),
8226 ModuleState::Stopped,
8227 Some(false),
8228 )
8229 .await
8230 .unwrap();
8231
8232 let state = lock_snapshot(&snapshot).unwrap();
8233 assert_eq!(
8234 state.last_exit.as_ref().map(|exit| exit.kind),
8235 Some(ExitKind::Crash)
8236 );
8237 assert_eq!(state.lifetime_restarts, 0);
8238 assert_eq!(state.crash_restarts.len(), 0);
8239 }
8240
8241 #[test]
8242 fn fatal_connection_teardown_cannot_arm_a_marker_for_a_surviving_process() {
8243 assert!(!include_str!("server.rs")
8249 .contains("router.record_deliberate_connection_severance(ctx.connection_id)"));
8250 }
8251
8252 #[test]
8259 fn drained_after_quiescence_wait_passes_ok_through_and_forces_false_on_err() {
8260 assert!(drained_after_quiescence_wait(&Ok(true)));
8261 assert!(!drained_after_quiescence_wait(&Ok(false)));
8262 assert!(!drained_after_quiescence_wait(&Err(
8263 SuperviseError::StatePoisoned { module_id: None }
8264 )));
8265 }
8266
8267 #[test]
8276 fn wait_error_exit_report_records_a_failed_terminal_with_no_code_or_signal() {
8277 let ring = Arc::new(Mutex::new(TerminalRing::new(
8278 TerminalRingConfig::default(),
8279 0,
8280 )));
8281 record_wait_error_terminal("wait-error", &ring, &super::SpawnEventFeed::default());
8282
8283 let snapshot = ring.lock().unwrap().snapshot();
8284 assert_eq!(snapshot.entries.len(), 1);
8285 let entry = &snapshot.entries[0];
8286 assert_eq!(entry.exit_code, None);
8287 assert_eq!(entry.exit_signal, None);
8288 assert_eq!(entry.disposition, TerminalDisposition::Failed);
8289 }
8290
8291 #[test]
8292 fn wait_error_exit_path_preserves_spawn_event_density() {
8293 let feed = super::SpawnEventFeed::default();
8294 feed.configure_incarnation("wait-error-density".to_string());
8295 feed.emit_spawned("wait-error", 41, 1);
8296 let ring = Arc::new(Mutex::new(TerminalRing::new(
8297 TerminalRingConfig::default(),
8298 0,
8299 )));
8300
8301 record_wait_error_terminal("wait-error", &ring, &feed);
8302 feed.emit_spawned("after-wait-error", 42, 2);
8303
8304 let state = feed.0.lock().unwrap();
8305 let sequences = state
8306 .events
8307 .iter()
8308 .map(|event| event.cursor.seq)
8309 .collect::<Vec<_>>();
8310 assert_eq!(sequences, vec![1, 2, 3]);
8311 assert_eq!(state.events[1].kind, SpawnEventKind::Exited);
8312 assert_eq!(state.events[1].exit_code, None);
8313 assert_eq!(state.events[1].exit_signal, None);
8314 }
8315
8316 #[test]
8320 fn wait_error_exit_report_is_classified_as_a_crash() {
8321 assert_eq!(wait_error_exit_report().kind, ExitKind::Crash);
8322 }
8323}
8324
8325#[cfg(test)]
8326mod health_evidence_tests {
8327 use super::{HealthProbeError, HealthProbeEvidence};
8328 use std::collections::HashSet;
8329
8330 #[test]
8338 fn only_a_dead_lane_is_proof_of_death() {
8339 assert!(HealthProbeError::lane_dead("gone").is_proof_of_death());
8340 assert!(!HealthProbeError::no_answer("timed out").is_proof_of_death());
8344 assert!(!HealthProbeError::bad_answer("garbage").is_proof_of_death());
8345 assert!(!HealthProbeError::misconfigured("no table").is_proof_of_death());
8346 }
8347
8348 #[test]
8354 fn every_evidence_class_has_a_distinct_label() {
8355 let labels = [
8356 HealthProbeError::lane_dead("").label(),
8357 HealthProbeError::no_answer("").label(),
8358 HealthProbeError::bad_answer("").label(),
8359 HealthProbeError::misconfigured("").label(),
8360 ];
8361 let unique: HashSet<_> = labels.iter().collect();
8362 assert_eq!(unique.len(), labels.len(), "labels collided: {labels:?}");
8363 }
8364
8365 #[test]
8371 fn classification_preserves_the_original_message() {
8372 let err = HealthProbeError::no_answer("module did not answer within 5s");
8373 assert_eq!(err.to_string(), "module did not answer within 5s");
8374 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8375 }
8376}
8377
8378#[cfg(test)]
8379mod health_tombstone_tests {
8380 use std::{path::PathBuf, sync::Arc, time::Duration};
8381
8382 use subc_protocol::{
8383 manifest::Concurrency,
8384 session::{HealthStatus, ModuleControlResponse},
8385 };
8386 use tokio::sync::mpsc;
8387
8388 use super::{
8389 probe_module_health, HealthAction, HealthConfig, HealthProbeEvidence, ModuleProtocol,
8390 ModuleSpec, RestartPolicy, Supervisor, SupervisorRuntimeConfig,
8391 };
8392 use crate::{
8393 control::ControlHandler,
8394 forwarding::{ForwardingTable, ModuleControlRpcCompletion, ModuleControlRpcOutcome},
8395 registry::{ConnectionId, Registry},
8396 router::FrameSink,
8397 };
8398
8399 struct ProbeHarness {
8400 spec: ModuleSpec,
8401 runtime: SupervisorRuntimeConfig,
8402 forwarding: Arc<ForwardingTable>,
8403 module_connection: ConnectionId,
8404 module_rx: mpsc::Receiver<crate::router::OutboundFrame>,
8405 handler: ControlHandler,
8406 module: super::SupervisedModule,
8407 }
8408
8409 fn probe_harness() -> ProbeHarness {
8410 let registry = Arc::new(Registry::default());
8411 let forwarding = Arc::new(ForwardingTable::default());
8412 let supervisor_handle = super::SupervisorHandle::new();
8413 let health = HealthConfig {
8414 cadence: Duration::from_secs(30),
8415 deadline: Duration::from_secs(5),
8416 failure_threshold: 3,
8417 on_degraded: HealthAction::Report,
8418 on_failing: HealthAction::Report,
8419 critical: false,
8420 };
8421 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
8422 .with_forwarding(Arc::clone(&forwarding))
8423 .with_handle(supervisor_handle.clone())
8424 .with_health_config(health);
8425 let spec = ModuleSpec {
8426 module_id: "late-health-module".to_string(),
8427 program: PathBuf::from("disabled-module"),
8428 args: Vec::new(),
8429 env: Vec::new(),
8430 reserved: false,
8431 reserved_prefixes: Vec::new(),
8432 protocol: ModuleProtocol::Subc,
8433 overlap: Default::default(),
8434 };
8435 let module = supervisor
8436 .supervise_configured(spec.clone(), false)
8437 .unwrap();
8438 let runtime = supervisor.runtime_config();
8439 let handler = ControlHandler::with_forwarding(registry, Arc::clone(&forwarding))
8440 .with_supervisor(supervisor_handle);
8441 let module_connection = ConnectionId::new(700);
8442 let (module_tx, module_rx) = mpsc::channel(8);
8443 forwarding
8444 .register_module_connection(
8445 module_connection,
8446 spec.module_id.clone(),
8447 subc_protocol::PROTOCOL_VERSION,
8448 Concurrency::ModuleManaged,
8449 FrameSink::new(module_tx),
8450 )
8451 .unwrap();
8452
8453 ProbeHarness {
8454 spec,
8455 runtime,
8456 forwarding,
8457 module_connection,
8458 module_rx,
8459 handler,
8460 module,
8461 }
8462 }
8463
8464 async fn finish_after(
8465 harness: &mut ProbeHarness,
8466 stall: Duration,
8467 ) -> ModuleControlRpcCompletion {
8468 assert!(stall > harness.runtime.health.deadline);
8469 let deadline = harness.runtime.health.deadline;
8470 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8471 let answer = async {
8472 let frame = harness.module_rx.recv().await.expect("health.check frame");
8473 tokio::time::advance(deadline).await;
8474 tokio::task::yield_now().await;
8475 tokio::time::advance(stall - deadline).await;
8476 harness
8477 .forwarding
8478 .complete_module_control_rpc(
8479 harness.module_connection,
8480 frame.header.corr,
8481 Some("health.check"),
8482 ModuleControlRpcOutcome::Response(ModuleControlResponse::HealthCheck {
8483 status: HealthStatus::Ok,
8484 detail: None,
8485 metrics: None,
8486 }),
8487 )
8488 .unwrap()
8489 };
8490 let (probe_result, completion) = tokio::join!(probe, answer);
8491 let err = probe_result.expect_err("probe must miss its deadline");
8492 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8493 completion
8494 }
8495
8496 async fn time_out_without_answer(harness: &mut ProbeHarness) {
8497 let deadline = harness.runtime.health.deadline;
8498 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8499 let exhaust_deadline = async {
8500 let _frame = harness.module_rx.recv().await.expect("health.check frame");
8501 tokio::time::advance(deadline).await;
8502 tokio::task::yield_now().await;
8503 };
8504 let (probe_result, ()) = tokio::join!(probe, exhaust_deadline);
8505 let err = probe_result.expect_err("probe must miss its deadline");
8506 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8507 }
8508
8509 #[tokio::test(start_paused = true)]
8510 async fn late_health_answers_record_start_anchored_latency_for_two_stalls() {
8511 let mut harness = probe_harness();
8512
8513 let first = finish_after(&mut harness, Duration::from_secs(8)).await;
8514 let first_latency = match &first {
8515 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8516 other => panic!("late answer was not retained: {other:?}"),
8517 };
8518 assert!(harness.handler.observe_module_control_completion(first));
8519
8520 let second = finish_after(&mut harness, Duration::from_secs(11)).await;
8521 let second_latency = match &second {
8522 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8523 other => panic!("late answer was not retained: {other:?}"),
8524 };
8525 assert!(harness.handler.observe_module_control_completion(second));
8526
8527 assert_eq!(first_latency, Duration::from_secs(8));
8528 assert_eq!(
8529 second_latency - first_latency,
8530 Duration::from_secs(3),
8531 "latency must grow linearly with the additional stall"
8532 );
8533 let health = harness.module.status().unwrap().health;
8534 assert_eq!(health.late_answer_count, 2);
8535 assert_eq!(health.last_late_answer_latency_ms, Some(11_000));
8536 }
8537
8538 #[tokio::test(start_paused = true)]
8546 async fn late_answer_clears_the_consecutive_failure_streak() {
8547 let mut harness = probe_harness();
8548
8549 time_out_without_answer(&mut harness).await;
8551 harness
8552 .module
8553 .record_health_probe_failure_for_test("[no-answer] test miss")
8554 .unwrap();
8555 assert_eq!(
8556 harness.module.status().unwrap().health.consecutive_failures,
8557 1,
8558 "precondition: the miss must be on the streak before the late answer"
8559 );
8560
8561 let late = finish_after(&mut harness, Duration::from_secs(9)).await;
8563 assert!(matches!(
8564 late,
8565 ModuleControlRpcCompletion::LateHealthAnswer { .. }
8566 ));
8567 assert!(harness.handler.observe_module_control_completion(late));
8568
8569 let health = harness.module.status().unwrap().health;
8570 assert_eq!(
8571 health.consecutive_failures, 0,
8572 "a late answer is an answer: the streak must reset"
8573 );
8574 assert_eq!(health.late_answer_count, 1);
8575 }
8576
8577 #[tokio::test(start_paused = true)]
8578 async fn repeated_serial_probe_cycles_keep_one_tombstone_per_endpoint() {
8579 let mut harness = probe_harness();
8580
8581 for _ in 0..20 {
8582 time_out_without_answer(&mut harness).await;
8583 assert_eq!(
8584 harness.forwarding.health_probe_tombstone_count().unwrap(),
8585 1
8586 );
8587 }
8588 }
8589}
8590
8591#[cfg(test)]
8592mod child_env_tests {
8593 use super::{
8594 apply_child_env, apply_spawn_role, apply_wire_spawn_args, ModuleProtocol, ModuleSpec,
8595 SpawnRole, SupervisorHandle, SPAWN_ROLE_SWAP_CANDIDATE, SUBC_ARG, SUBC_LAUNCH_NONCE_ENV,
8596 SUBC_MODULE_ID_ENV, SUBC_SPAWN_ROLE_ENV,
8597 };
8598 use std::{ffi::OsStr, path::PathBuf};
8599 use tokio::process::Command;
8600
8601 fn spec(env: Vec<(String, String)>) -> ModuleSpec {
8602 ModuleSpec {
8603 module_id: "env-plan".to_string(),
8604 program: PathBuf::from("/nonexistent"),
8605 args: Vec::new(),
8606 env,
8607 reserved: false,
8608 reserved_prefixes: Vec::new(),
8609 protocol: ModuleProtocol::Subc,
8610 overlap: Default::default(),
8611 }
8612 }
8613
8614 #[test]
8628 fn ambient_ck_log_is_removed_and_a_configured_one_survives() {
8629 let mut command = Command::new("/nonexistent");
8630 apply_child_env(&mut command, &spec(Vec::new()));
8631 let removed = command
8632 .as_std()
8633 .get_envs()
8634 .any(|(key, value)| key == OsStr::new("CK_LOG") && value.is_none());
8635 assert!(
8636 removed,
8637 "ambient CK_LOG must be explicitly removed for an unconfigured module"
8638 );
8639
8640 let mut configured = Command::new("/nonexistent");
8641 apply_child_env(
8642 &mut configured,
8643 &spec(vec![("CK_LOG".to_string(), "debug".to_string())]),
8644 );
8645 let effective = configured
8646 .as_std()
8647 .get_envs()
8648 .filter(|(key, _)| *key == OsStr::new("CK_LOG"))
8649 .last()
8650 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()));
8651 assert_eq!(
8652 effective,
8653 Some(Some("debug".to_string())),
8654 "a module's configured CK_LOG must survive the ambient removal"
8655 );
8656 }
8657
8658 #[test]
8667 fn protocol_none_spawn_carries_no_subc_argument_and_no_nonce() {
8668 let connection_file = std::path::Path::new("/run/subc-connection.json");
8669 let handle = SupervisorHandle::new();
8670
8671 let mut none_spec = spec(Vec::new());
8672 none_spec.protocol = ModuleProtocol::None;
8673 let mut none = Command::new("/nonexistent");
8674 apply_wire_spawn_args(&mut none, &none_spec, Some(connection_file), Some(&handle))
8675 .expect("protocol-none spawn args apply");
8676 let none_args: Vec<String> = none
8677 .as_std()
8678 .get_args()
8679 .map(|a| a.to_string_lossy().into_owned())
8680 .collect();
8681 assert!(
8682 !none_args.iter().any(|a| a == SUBC_ARG),
8683 "protocol:none argv must not carry --subc; got {none_args:?}"
8684 );
8685 let none_has_nonce = none
8686 .as_std()
8687 .get_envs()
8688 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some());
8689 assert!(
8690 !none_has_nonce,
8691 "protocol:none spawn must not receive a launch nonce"
8692 );
8693 let none_has_module_id = none
8694 .as_std()
8695 .get_envs()
8696 .any(|(key, value)| key == OsStr::new(SUBC_MODULE_ID_ENV) && value.is_some());
8697 assert!(
8698 none_has_module_id,
8699 "SUBC_MODULE_ID is inert and stays on every path"
8700 );
8701 assert!(
8702 handle.spawn_nonce(&none_spec.module_id).is_none(),
8703 "no nonce record for a process that will never present one"
8704 );
8705
8706 let wire_spec = spec(Vec::new());
8708 let mut wire = Command::new("/nonexistent");
8709 apply_wire_spawn_args(&mut wire, &wire_spec, Some(connection_file), Some(&handle))
8710 .expect("subc-wire spawn args apply");
8711 let wire_args: Vec<String> = wire
8712 .as_std()
8713 .get_args()
8714 .map(|a| a.to_string_lossy().into_owned())
8715 .collect();
8716 assert_eq!(
8717 wire_args,
8718 vec![
8719 SUBC_ARG.to_string(),
8720 connection_file.to_string_lossy().into_owned()
8721 ],
8722 "a subc-wire spawn still carries --subc <path>"
8723 );
8724 assert!(wire
8725 .as_std()
8726 .get_envs()
8727 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some()));
8728 assert!(handle.spawn_nonce(&wire_spec.module_id).is_some());
8729 }
8730
8731 #[test]
8741 fn plain_spawn_removes_the_spawn_role_even_when_the_spec_sets_it() {
8742 let role = |command: &Command| {
8743 command
8744 .as_std()
8745 .get_envs()
8746 .filter(|(key, _)| *key == OsStr::new(SUBC_SPAWN_ROLE_ENV))
8747 .last()
8748 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()))
8749 };
8750 let forged = spec(vec![(
8751 SUBC_SPAWN_ROLE_ENV.to_string(),
8752 SPAWN_ROLE_SWAP_CANDIDATE.to_string(),
8753 )]);
8754
8755 let mut plain = Command::new("/nonexistent");
8756 apply_child_env(&mut plain, &forged);
8757 apply_spawn_role(&mut plain, SpawnRole::Plain);
8758 assert_eq!(
8759 role(&plain),
8760 Some(None),
8761 "a plain spawn must remove SUBC_SPAWN_ROLE, whatever the spec says"
8762 );
8763
8764 let mut candidate = Command::new("/nonexistent");
8765 apply_child_env(&mut candidate, &spec(Vec::new()));
8766 apply_spawn_role(&mut candidate, SpawnRole::SwapCandidate);
8767 assert_eq!(
8768 role(&candidate),
8769 Some(Some(SPAWN_ROLE_SWAP_CANDIDATE.to_string()))
8770 );
8771 }
8772
8773 #[test]
8779 fn daemon_private_capture_keys_are_not_passed_to_the_child() {
8780 let mut command = Command::new("/nonexistent");
8781 apply_child_env(
8782 &mut command,
8783 &spec(vec![
8784 (super::CAPTURE_KEEP_ENV.to_string(), "5".to_string()),
8785 ("KEPT".to_string(), "yes".to_string()),
8786 ]),
8787 );
8788 let keys: Vec<String> = command
8789 .as_std()
8790 .get_envs()
8791 .filter(|(_, value)| value.is_some())
8792 .map(|(key, _)| key.to_string_lossy().into_owned())
8793 .collect();
8794 assert!(keys.contains(&"KEPT".to_string()), "got {keys:?}");
8795 assert!(
8796 !keys.contains(&super::CAPTURE_KEEP_ENV.to_string()),
8797 "daemon-private capture key leaked to the child: {keys:?}"
8798 );
8799 }
8800}
8801
8802#[cfg(test)]
8803mod jitter_tests {
8804 use super::jittered_health_delay;
8805 use std::{collections::HashSet, time::Duration};
8806
8807 const FLEET: [&str; 14] = [
8816 "aft",
8817 "alfonso-core",
8818 "magic-context",
8819 "broca",
8820 "thalamus",
8821 "quota",
8822 "engram",
8823 "plexus",
8824 "cerebellum",
8825 "astrocyte",
8826 "synapse",
8827 "subc-mcp",
8828 "cortexkit-credentials",
8829 "subc-federation",
8830 ];
8831
8832 #[test]
8840 fn probe_delays_disperse_across_the_fleet() {
8841 let cadence = Duration::from_secs(30);
8842 let delays: HashSet<Duration> = FLEET
8843 .iter()
8844 .map(|id| jittered_health_delay(id, 0, cadence))
8845 .collect();
8846 assert_eq!(
8847 delays.len(),
8848 FLEET.len(),
8849 "every supervised module must land on its own probe offset"
8850 );
8851 }
8852
8853 #[test]
8859 fn jitter_only_delays_and_stays_within_one_tenth_of_cadence() {
8860 let cadence = Duration::from_secs(30);
8861 let span = cadence / 10;
8862 for id in FLEET {
8863 for probe_index in 0..8 {
8864 let delay = jittered_health_delay(id, probe_index, cadence);
8865 assert!(
8866 delay >= cadence,
8867 "{id}#{probe_index}: jitter must not shorten the cadence"
8868 );
8869 assert!(
8870 delay < cadence + span,
8871 "{id}#{probe_index}: jitter must stay inside one tenth of the cadence"
8872 );
8873 }
8874 }
8875 }
8876
8877 #[test]
8883 fn a_module_offset_is_stable_across_restarts() {
8884 let cadence = Duration::from_secs(30);
8885 for id in FLEET {
8886 assert_eq!(
8887 jittered_health_delay(id, 0, cadence),
8888 jittered_health_delay(id, 0, cadence),
8889 "{id}: the same module and probe index must produce the same offset"
8890 );
8891 }
8892 }
8893
8894 #[test]
8896 fn zero_cadence_yields_zero_delay() {
8897 assert_eq!(
8898 jittered_health_delay("aft", 0, Duration::ZERO),
8899 Duration::ZERO
8900 );
8901 }
8902}
8903
8904#[cfg(all(test, target_os = "linux"))]
8905mod cgroup_placement_tests {
8906 use super::{
8907 apply_cgroup_placement, remove_module_cgroup, ModuleProtocol, ModuleSpec, SuperviseError,
8908 SupervisedChild,
8909 };
8910 use crate::{
8911 stderr_tail::{StderrRing, StderrTailConfig},
8912 test_support::TestTempDir,
8913 };
8914 use std::{
8915 fs, io,
8916 path::{Path, PathBuf},
8917 sync::{Arc, Mutex},
8918 };
8919 use tokio::process::Command;
8920
8921 #[test]
8922 fn failed_parent_cgroup_open_is_a_cgroup_supervision_error() {
8923 let path = Path::new("/definitely-missing-subc-cgroup");
8924 let mut command = Command::new("true");
8925 let error = apply_cgroup_placement(
8926 &mut command,
8927 &ModuleSpec {
8928 module_id: "broken-cgroup".to_string(),
8929 program: PathBuf::from("true"),
8930 args: Vec::new(),
8931 env: Vec::new(),
8932 reserved: false,
8933 reserved_prefixes: Vec::new(),
8934 protocol: ModuleProtocol::Subc,
8935 overlap: Default::default(),
8936 },
8937 path,
8938 )
8939 .expect_err("a parent cgroup open failure must reject the supervised spawn");
8940 let reason = error.to_string();
8941
8942 assert!(
8943 matches!(error, SuperviseError::Cgroup { .. }),
8944 "parent cgroup open must be reported as a cgroup supervision error: {reason}"
8945 );
8946 assert!(
8947 reason.contains("/definitely-missing-subc-cgroup/cgroup.procs"),
8948 "parent cgroup open failure must name cgroup.procs: {reason}"
8949 );
8950 }
8951
8952 #[tokio::test]
8953 async fn reaping_a_child_removes_its_empty_module_cgroup() {
8954 let root = TestTempDir::new("supervisor-reap-cgroup");
8955 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
8956 let placement = subc_cgroup::prepare_at(&root)
8957 .expect("prepare scratch cgroup root")
8958 .expect("scratch root has a cgroup.procs marker");
8959 let module_id = "reaped-module";
8960 let module = placement
8961 .module_path(module_id)
8962 .expect("create scratch module cgroup");
8963 let child = Command::new("true")
8964 .spawn()
8965 .expect("spawn short-lived child");
8966 let pid = child.id().expect("spawned child has pid");
8967 let mut child = SupervisedChild {
8968 child,
8969 module_id: module_id.to_string(),
8970 cgroup_placement: Some(placement),
8971 stdout_pump: None,
8972 stderr_pump: None,
8973 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
8974 spawned_at_ms: 0,
8975 spawned_from: PathBuf::from("true"),
8976 spawned_file_identity: None,
8977 process_start_time: None,
8978 process_identity: None,
8979 pid,
8980 roster_guard: None,
8981 };
8982
8983 child.wait().await.expect("reap short-lived child");
8984
8985 assert!(
8986 !module.exists(),
8987 "reaping the supervised child must remove its empty cgroup"
8988 );
8989 }
8990
8991 #[test]
8992 fn non_empty_cgroup_removal_is_reported_without_blocking_teardown() {
8993 let root = TestTempDir::new("supervisor-non-empty-cgroup");
8994 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
8995 let placement = subc_cgroup::prepare_at(&root)
8996 .expect("prepare scratch cgroup root")
8997 .expect("scratch root has a cgroup.procs marker");
8998 let module = placement
8999 .module_path("surviving-module")
9000 .expect("create scratch module cgroup");
9001 fs::write(module.join("surviving-process"), b"still present")
9002 .expect("make scratch cgroup non-empty");
9003 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::WARN);
9004
9005 remove_module_cgroup(&placement, "surviving-module");
9006
9007 let logs = crate::router::test_log::captured_logs(&logs);
9008 assert!(
9009 module.exists(),
9010 "failed removal must leave the cgroup intact"
9011 );
9012 assert!(
9013 logs.contains("could not remove module cgroup after process exit; continuing teardown")
9014 && logs.contains("surviving-module"),
9015 "best-effort removal must report the failure without returning it: {logs}"
9016 );
9017 }
9018
9019 #[test]
9020 fn cgroup_pre_exec_spawn_failure_names_the_cgroup_path() {
9021 let cgroup_path = PathBuf::from("/sys/fs/cgroup/subc-modules/broken-module");
9022 let reason = SuperviseError::Spawn {
9023 program: PathBuf::from("/bin/true"),
9024 source: io::Error::from_raw_os_error(13),
9025 cgroup_path: Some(cgroup_path.clone()),
9026 }
9027 .to_string();
9028
9029 assert!(
9030 reason.contains(&cgroup_path.display().to_string()),
9031 "a pre_exec spawn failure must name the cgroup path: {reason}"
9032 );
9033 }
9034}
9035
9036#[cfg(test)]
9037mod spawn_subscriber_lag_tests {
9038 use super::*;
9039
9040 #[tokio::test]
9045 async fn lagged_spawn_subscriber_receives_a_terminal_lagged_error_after_its_queued_frames() {
9046 let feed = SpawnEventFeed::default();
9047 feed.configure_incarnation("lag-incarnation".to_string());
9048 let (tx, mut rx) = mpsc::channel(1);
9051 feed.subscribe(ConnectionId::new(1), 7, 1, None, FrameSink::new(tx))
9052 .expect("subscribe");
9053 let emitted = SPAWN_SUBSCRIBER_BUFFER + 16;
9054 for index in 0..emitted {
9055 feed.emit_spawned(&format!("lag-module-{index}"), 1000, 0);
9056 tokio::task::yield_now().await;
9059 }
9060 assert_eq!(
9061 feed.subscriber_count(),
9062 0,
9063 "the lagged subscriber must be removed"
9064 );
9065
9066 let mut data = Vec::new();
9067 let mut last = None;
9068 loop {
9069 let next = tokio::time::timeout(Duration::from_secs(5), rx.recv())
9070 .await
9071 .expect("the forwarder must finish once the subscriber is dropped");
9072 let Some(outbound) = next else { break };
9073 let frame = outbound.frame;
9074 if frame.header.ty == FrameType::StreamData {
9075 assert!(last.is_none(), "no data may follow the terminal frame");
9076 let event: SpawnEvent = serde_json::from_slice(&frame.body).unwrap();
9077 data.push(event.cursor.seq);
9078 } else {
9079 assert!(last.is_none(), "exactly one terminal frame");
9080 last = Some(frame);
9081 }
9082 }
9083 assert!(!data.is_empty(), "queued frames drain before the terminal");
9084 for pair in data.windows(2) {
9085 assert_eq!(
9086 pair[1],
9087 pair[0] + 1,
9088 "queued frames arrive dense and in order"
9089 );
9090 }
9091 let terminal = last.expect("a lagged subscriber must receive a terminal frame");
9092 assert_eq!(terminal.header.ty, FrameType::Error);
9093 assert_eq!(terminal.header.corr, 7);
9094 let body: subc_protocol::ErrorBody = serde_json::from_slice(&terminal.body).unwrap();
9095 assert_eq!(body.code, SPAWN_SUBSCRIBER_LAGGED_CODE);
9096 let detail = body.detail.expect("lagged error carries detail");
9097 assert_eq!(
9098 detail["first_undelivered_cursor"]["seq"],
9099 data.last().unwrap() + 1,
9100 "the named cursor is the first event the subscriber did not receive"
9101 );
9102 assert_eq!(
9103 detail["first_undelivered_cursor"]["daemon_incarnation"],
9104 "lag-incarnation"
9105 );
9106 }
9107}
9108
9109#[cfg(test)]
9110mod terminal_history_read_concurrency_tests {
9111 use super::*;
9112 use crate::{terminal_journal::read_pause, test_support::TestTempDir};
9113 use std::sync::mpsc as std_mpsc;
9114
9115 fn journaled_ring(
9116 journal: &Arc<crate::terminal_journal::TerminalJournal>,
9117 ) -> Arc<Mutex<TerminalRing>> {
9118 Arc::new(Mutex::new(
9119 TerminalRing::new(TerminalRingConfig::default(), 1)
9120 .with_journal(Some(Arc::clone(journal))),
9121 ))
9122 }
9123
9124 fn crash(at_ms: u64) -> ExitReport {
9125 ExitReport {
9126 kind: ExitKind::Crash,
9127 code: Some(1),
9128 signal: None,
9129 at_ms,
9130 }
9131 }
9132
9133 fn record_within(
9136 module_id: &'static str,
9137 ring: &Arc<Mutex<TerminalRing>>,
9138 at_ms: u64,
9139 bound: Duration,
9140 ) -> bool {
9141 let ring = Arc::clone(ring);
9142 let (done, done_rx) = std_mpsc::channel();
9143 std::thread::spawn(move || {
9144 record_terminal(
9145 module_id,
9146 &ring,
9147 &SpawnEventFeed::default(),
9148 &crash(at_ms),
9149 TerminalDisposition::Restarting,
9150 );
9151 let _ = done.send(());
9152 });
9153 done_rx.recv_timeout(bound).is_ok()
9154 }
9155
9156 #[test]
9161 fn exits_recorded_during_a_paused_history_read_are_not_blocked_or_half_merged() {
9162 let dir = TestTempDir::new("terminal-history-concurrent-read");
9163 let path = dir.join("terminals.jsonl");
9164 let journal = Arc::new(crate::terminal_journal::TerminalJournal::open(
9165 path.clone(),
9166 "daemon".into(),
9167 ));
9168 let reader_ring = journaled_ring(&journal);
9169 let other_ring = journaled_ring(&journal);
9170 assert!(record_within(
9171 "reader-module",
9172 &reader_ring,
9173 10,
9174 Duration::from_secs(5)
9175 ));
9176
9177 let (started, release) = read_pause::install(&path);
9178 let reading = {
9179 let ring = Arc::clone(&reader_ring);
9180 std::thread::spawn(move || durable_terminal_history_of(&ring, "reader-module"))
9181 };
9182 started
9183 .recv_timeout(Duration::from_secs(5))
9184 .expect("the history read reached its pause");
9185
9186 let bound = Duration::from_secs(1);
9187 assert!(
9188 record_within("other-module", &other_ring, 20, bound),
9189 "another module's exit waited on a history read (journal writer held)"
9190 );
9191 assert!(
9192 record_within("reader-module", &reader_ring, 30, bound),
9193 "the read module's own exit waited on its history read (ring held)"
9194 );
9195
9196 drop(release);
9197 let paused = reading.join().unwrap();
9198 assert_eq!(
9199 paused.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9200 vec![10],
9201 "an exit recorded after the read began lands in neither half of it"
9202 );
9203 assert_eq!(paused.journal_skipped_lines, 0);
9204 assert_eq!(paused.journal_read_errors, 0);
9205
9206 let after = durable_terminal_history_of(&reader_ring, "reader-module");
9207 assert_eq!(
9208 after.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9209 vec![10, 30],
9210 "the next read merges ring and journal with no duplicate"
9211 );
9212 assert_eq!(after.journal_skipped_lines, 0);
9213 }
9214}
9215
9216#[cfg(test)]
9221mod stderr_settle_tests {
9222 use std::{
9223 future::Future,
9224 io,
9225 pin::Pin,
9226 sync::{Arc, Mutex},
9227 task::{Context, Poll},
9228 time::Duration,
9229 };
9230
9231 use tokio::{
9232 io::{AsyncRead, ReadBuf},
9233 sync::oneshot,
9234 time::Instant,
9235 };
9236
9237 use super::{settle_stderr_pump, StderrPump};
9238 use crate::stderr_tail::{
9239 pump_stderr_to, CaptureState, OutputSink, StderrRing, StderrTailConfig, TailEntry,
9240 };
9241
9242 const BOUND: Duration = Duration::from_millis(250);
9243
9244 struct HeldReader {
9248 before: Option<Vec<u8>>,
9249 gate: Option<oneshot::Receiver<()>>,
9250 after: io::Cursor<Vec<u8>>,
9251 }
9252
9253 impl AsyncRead for HeldReader {
9254 fn poll_read(
9255 mut self: Pin<&mut Self>,
9256 cx: &mut Context<'_>,
9257 buf: &mut ReadBuf<'_>,
9258 ) -> Poll<io::Result<()>> {
9259 if let Some(bytes) = self.before.take() {
9260 buf.put_slice(&bytes);
9261 return Poll::Ready(Ok(()));
9262 }
9263 if let Some(gate) = self.gate.as_mut() {
9264 match Pin::new(gate).poll(cx) {
9265 Poll::Pending => return Poll::Pending,
9266 Poll::Ready(_) => self.gate = None,
9267 }
9268 }
9269 Pin::new(&mut self.after).poll_read(cx, buf)
9270 }
9271 }
9272
9273 struct DiscardSink;
9274
9275 impl OutputSink for DiscardSink {
9276 fn write_line(&mut self, _line: &[u8]) {}
9277 }
9278
9279 fn line(text: &str) -> TailEntry {
9280 TailEntry::Line {
9281 text: text.to_string(),
9282 truncated: false,
9283 }
9284 }
9285
9286 fn lock(ring: &Arc<Mutex<StderrRing>>) -> std::sync::MutexGuard<'_, StderrRing> {
9287 ring.lock().unwrap()
9288 }
9289
9290 fn held_pump(
9294 ring: &Arc<Mutex<StderrRing>>,
9295 before: &str,
9296 after: &str,
9297 ) -> (StderrPump, oneshot::Sender<()>) {
9298 let generation = lock(ring).begin_process();
9299 let (release, gate) = oneshot::channel();
9300 let reader = HeldReader {
9301 before: Some(before.as_bytes().to_vec()),
9302 gate: Some(gate),
9303 after: io::Cursor::new(after.as_bytes().to_vec()),
9304 };
9305 let task = tokio::spawn(pump_stderr_to(
9306 reader,
9307 Arc::clone(ring),
9308 generation,
9309 DiscardSink,
9310 ));
9311 (StderrPump { task, generation }, release)
9312 }
9313
9314 async fn wait_until(ring: &Arc<Mutex<StderrRing>>, done: impl Fn(&StderrRing) -> bool) {
9315 for _ in 0..1000 {
9316 if done(&lock(ring)) {
9317 return;
9318 }
9319 tokio::time::sleep(Duration::from_millis(1)).await;
9320 }
9321 panic!(
9322 "ring never reached the expected state: {:?}",
9323 lock(ring).snapshot(None, None)
9324 );
9325 }
9326
9327 #[tokio::test(start_paused = true)]
9328 async fn a_crash_line_the_reader_had_not_reached_by_the_bound_is_kept_before_the_restart() {
9329 let ring = Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default())));
9330 let (pump, release) = held_pump(&ring, "booting\n", "config error: missing storage\n");
9331
9332 settle_stderr_pump("crasher", &ring, pump, BOUND).await;
9333 let before_release = lock(&ring).snapshot(None, None);
9334 assert!(
9335 matches!(before_release.capture, CaptureState::Incomplete { .. }),
9336 "a reader that has not reached EOF cannot claim a whole tail: {before_release:?}"
9337 );
9338
9339 let next = lock(&ring).begin_process();
9342 lock(&ring).push_line_from(next, "next process booting");
9343 release.send(()).unwrap();
9344 wait_until(&ring, |ring| {
9345 ring.snapshot(None, None).capture == CaptureState::Captured
9346 })
9347 .await;
9348
9349 assert_eq!(
9350 lock(&ring).snapshot(None, None).entries,
9351 vec![
9352 line("booting"),
9353 line("config error: missing storage"),
9354 TailEntry::ProcessStart,
9355 line("next process booting"),
9356 ],
9357 "the crash's last line must survive a slow reader and stay in the crashed process's section"
9358 );
9359 }
9360
9361 #[tokio::test(start_paused = true)]
9362 async fn a_pipe_held_open_by_a_descendant_reads_incomplete_without_delaying_the_restart_past_the_bound(
9363 ) {
9364 let ring = Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default())));
9365 let (pump, _held) = held_pump(&ring, "parent exiting\n", "");
9368
9369 let started = Instant::now();
9370 settle_stderr_pump("orphaning", &ring, pump, BOUND).await;
9371 assert_eq!(
9372 started.elapsed(),
9373 BOUND,
9374 "the restart must wait exactly the bound for a pipe that stays open, no longer"
9375 );
9376
9377 let next = lock(&ring).begin_process();
9378 lock(&ring).push_line_from(next, "next process booting");
9379 tokio::time::sleep(Duration::from_secs(60)).await;
9380
9381 let snapshot = lock(&ring).snapshot(None, None);
9382 match &snapshot.capture {
9383 CaptureState::Incomplete { reason } => assert!(
9384 reason.contains("had not reached EOF") && reason.contains("250ms"),
9385 "the reason must say what is missing and after how long: {reason}"
9386 ),
9387 other => panic!("expected Incomplete while the pipe is held open, got {other:?}"),
9388 }
9389 assert_eq!(
9390 snapshot.entries,
9391 vec![
9392 line("parent exiting"),
9393 TailEntry::ProcessStart,
9394 line("next process booting"),
9395 ]
9396 );
9397 }
9398
9399 #[tokio::test(start_paused = true)]
9400 async fn a_reader_that_reaches_eof_within_the_bound_leaves_the_tail_captured() {
9401 let ring = Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default())));
9402 let (pump, release) = held_pump(&ring, "one\n", "two\n");
9403 release.send(()).unwrap();
9404
9405 settle_stderr_pump("clean", &ring, pump, BOUND).await;
9406
9407 let snapshot = lock(&ring).snapshot(None, None);
9408 assert_eq!(snapshot.capture, CaptureState::Captured);
9409 assert_eq!(snapshot.entries, vec![line("one"), line("two")]);
9410 }
9411}