1use std::{
2 collections::{HashMap, VecDeque},
3 error::Error,
4 fmt, io,
5 path::PathBuf,
6 process::{ExitStatus, Stdio},
7 sync::{Arc, Mutex, OnceLock},
8 time::{Duration, SystemTime, UNIX_EPOCH},
9};
10
11use cortexkit_log::Retention;
12use serde_json::Value;
13use subc_control::{
14 ClientControlPush, LiveSpawn, ModuleProtocol, RouteCloseReason, SpawnCursor, SpawnEvent,
15 SpawnEventKind, SpawnSnapshot, SupervisorHealthStatus, TerminalDisposition, TerminalExitKind,
16};
17use subc_protocol::{
18 manifest::{SelfSignalKind, SignalAnchor},
19 session::{
20 HealthReport, HealthStatus, ModuleControlCommand, ModuleControlRequest,
21 MODULE_CONTROL_OP_HEALTH_CHECK,
22 },
23 Flags, FrameType, Priority, SUBC_LAUNCH_NONCE_ENV, SUBC_MODULE_ID_ENV,
24};
25use tokio::{
26 process::{Child, Command},
27 sync::{mpsc, oneshot, watch, Mutex as AsyncMutex},
28 task::JoinHandle,
29 time::{sleep, sleep_until, timeout, timeout_at, Instant},
30};
31use tracing::{debug, error, info, warn};
32
33use crate::{
34 child_roster::ChildRoster,
35 daemon_config::{
36 CAPTURE_KEEP_ENV, CAPTURE_MAX_AGE_DAYS_ENV, CAPTURE_MAX_FILE_MB_ENV, CK_LOG_ENV,
37 },
38 forwarding::{
39 CloseReason, ForwardingError, ForwardingTable, GoodbyeTarget, ModuleControlRpcOutcome,
40 ModuleDrainTarget, PendingModuleControlRpc,
41 },
42 provenance::{spawned_file_identity, ExecutableIdentityProbe, SpawnedFileIdentity},
43 registry::{ConnectionId, RegistryError},
44 stderr_tail::{
45 pump_stderr_to, pump_stdout_to, ChildOutputSink, StderrRing, StderrTailConfig,
46 StderrTailSnapshot,
47 },
48 terminal_ring::{TerminalHistorySnapshot, TerminalRecord, TerminalRing, TerminalRingConfig},
49 Frame, FrameSink, Registry,
50};
51
52#[path = "supervise_swap.rs"]
53mod swap;
54
55pub const SUBC_ARG: &str = "--subc";
61
62const DEFAULT_MAX_RESTARTS: u32 = 3;
63const DEFAULT_BACKOFF: Duration = Duration::from_millis(100);
64const DEFAULT_MAX_BACKOFF: Duration = Duration::from_secs(30);
65const DEFAULT_RESTART_WINDOW: Duration = Duration::from_secs(600);
69pub const DEFAULT_DRAIN_TIMEOUT: Duration = Duration::from_secs(30);
80const REGISTRY_RELEASE_TIMEOUT: Duration = Duration::from_secs(1);
81const REGISTRY_RELEASE_POLL: Duration = Duration::from_millis(10);
82const STDERR_PUMP_DRAIN_TIMEOUT: Duration = Duration::from_millis(250);
83pub const SPAWN_EVENT_RING_CAPACITY: usize = 4096;
85const SPAWN_SUBSCRIBER_BUFFER: usize = SPAWN_EVENT_RING_CAPACITY + 1;
86pub(crate) const SPAWN_SUBSCRIBER_LAGGED_CODE: &str = "spawn_subscriber_lagged";
91
92struct SupervisedChild {
93 child: Child,
94 #[cfg(target_os = "linux")]
97 module_id: String,
98 #[cfg(target_os = "linux")]
99 cgroup_placement: Option<subc_cgroup::Placement>,
100 stdout_pump: Option<JoinHandle<()>>,
101 stderr_pump: Option<JoinHandle<()>>,
102 stderr_ring: Arc<Mutex<StderrRing>>,
103 spawned_at_ms: u64,
104 spawned_from: PathBuf,
105 spawned_file_identity: Option<SpawnedFileIdentity>,
106 process_start_time: Option<u64>,
107 process_identity: Option<ProcessIdentity>,
108 pid: u32,
109 roster_guard: Option<crate::child_roster::RosterGuard>,
112}
113
114impl SupervisedChild {
115 fn id(&self) -> Option<u32> {
116 Some(self.pid)
117 }
118
119 fn process_identity(&self) -> Option<ProcessIdentity> {
120 self.process_identity
121 }
122
123 async fn wait(&mut self) -> io::Result<ExitStatus> {
124 let result = self.child.wait().await;
125 if result.is_ok() {
126 self.roster_guard = None;
129 }
130 #[cfg(target_os = "linux")]
131 if result.is_ok() {
132 if let Some(placement) = self.cgroup_placement.take() {
133 remove_module_cgroup(&placement, &self.module_id);
134 }
135 }
136 result
137 }
138
139 fn start_kill(&mut self) -> io::Result<()> {
140 self.child.start_kill()
141 }
142
143 async fn drain_stderr(&mut self, module_id: &str) {
144 if let Some(mut pump) = self.stdout_pump.take() {
145 match timeout(STDERR_PUMP_DRAIN_TIMEOUT, &mut pump).await {
146 Ok(Ok(())) => {}
147 Ok(Err(error)) => {
148 warn!(module_id, error = %error, "stdout pump ended unexpectedly");
149 }
150 Err(_) => {
151 pump.abort();
152 warn!(
153 module_id,
154 waited = ?STDERR_PUMP_DRAIN_TIMEOUT,
155 "stdout pump did not drain before restart; stopped it before the next process"
156 );
157 }
158 }
159 }
160
161 let Some(mut pump) = self.stderr_pump.take() else {
162 return;
163 };
164 match timeout(STDERR_PUMP_DRAIN_TIMEOUT, &mut pump).await {
165 Ok(Ok(())) => {}
166 Ok(Err(err)) => {
167 self.stderr_ring
168 .lock()
169 .unwrap_or_else(|poisoned| poisoned.into_inner())
170 .mark_incomplete(format!("stderr pump ended unexpectedly: {err}"));
171 warn!(module_id, error = %err, "stderr pump ended before clean EOF");
172 }
173 Err(_) => {
174 pump.abort();
175 self.stderr_ring
176 .lock()
177 .unwrap_or_else(|poisoned| poisoned.into_inner())
178 .mark_incomplete(format!(
179 "stderr pump did not reach EOF within {:?} before restart",
180 STDERR_PUMP_DRAIN_TIMEOUT
181 ));
182 warn!(
183 module_id,
184 waited = ?STDERR_PUMP_DRAIN_TIMEOUT,
185 "stderr pump did not drain before restart; stopped it before marking the new process"
186 );
187 }
188 }
189 }
190}
191
192fn registration_release_events() -> &'static watch::Sender<u64> {
193 static EVENTS: OnceLock<watch::Sender<u64>> = OnceLock::new();
194 EVENTS.get_or_init(|| {
195 let (sender, _receiver) = watch::channel(0);
196 sender
197 })
198}
199
200pub(crate) fn notify_registration_release() {
201 let events = registration_release_events();
202 let next_generation = (*events.borrow()).wrapping_add(1);
203 events.send_replace(next_generation);
204}
205
206#[derive(Debug, Clone, PartialEq, Eq)]
208pub struct ModuleSpec {
209 pub module_id: String,
210 pub program: PathBuf,
211 pub args: Vec<String>,
212 pub env: Vec<(String, String)>,
213 pub reserved: bool,
218 pub reserved_prefixes: Vec<String>,
223 pub protocol: ModuleProtocol,
239 pub overlap: ModuleOverlap,
244}
245
246#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
253pub enum ModuleOverlap {
254 #[default]
256 Exclusive,
257 Safe,
269}
270
271impl ModuleOverlap {
272 pub fn as_str(self) -> &'static str {
273 match self {
274 Self::Exclusive => "exclusive",
275 Self::Safe => "safe",
276 }
277 }
278}
279
280pub const SUBC_SPAWN_ROLE_ENV: &str = "SUBC_SPAWN_ROLE";
290pub const SPAWN_ROLE_SWAP_CANDIDATE: &str = "swap_candidate";
292pub const DEFAULT_SWAP_READY_TIMEOUT: Duration = Duration::from_secs(100);
297
298#[derive(Debug, Clone, Copy, PartialEq, Eq)]
316pub struct RestartPolicy {
317 pub max_restarts: u32,
318 pub backoff: Duration,
321 pub max_backoff: Duration,
323 pub window: Duration,
327}
328
329impl RestartPolicy {
330 pub fn new(max_restarts: u32, backoff: Duration) -> Self {
334 Self {
335 max_restarts,
336 backoff,
337 max_backoff: DEFAULT_MAX_BACKOFF,
338 window: DEFAULT_RESTART_WINDOW,
339 }
340 }
341
342 pub fn with_max_backoff(mut self, max_backoff: Duration) -> Self {
343 self.max_backoff = max_backoff;
344 self
345 }
346
347 pub fn with_window(mut self, window: Duration) -> Self {
348 self.window = window;
349 self
350 }
351
352 fn delay_for_restart(&self, restart_in_window: u32) -> Duration {
357 if self.backoff.is_zero() || self.max_backoff.is_zero() {
358 return Duration::ZERO;
359 }
360
361 let mut delay = self.backoff;
362 for _ in 0..restart_in_window {
363 if delay >= self.max_backoff {
364 return self.max_backoff;
365 }
366 delay = delay
367 .checked_mul(10)
368 .unwrap_or(self.max_backoff)
369 .min(self.max_backoff);
370 }
371 delay.min(self.max_backoff)
372 }
373
374 fn budget_exhausted_detail(&self) -> String {
379 format!(
380 "crash budget exhausted: max_restarts={} within window_secs={}",
381 self.max_restarts,
382 self.window.as_secs()
383 )
384 }
385}
386
387impl Default for RestartPolicy {
388 fn default() -> Self {
389 Self {
390 max_restarts: DEFAULT_MAX_RESTARTS,
391 backoff: DEFAULT_BACKOFF,
392 max_backoff: DEFAULT_MAX_BACKOFF,
393 window: DEFAULT_RESTART_WINDOW,
394 }
395 }
396}
397
398#[derive(Debug, Clone, Copy, PartialEq, Eq)]
399struct CrashRestartSchedule {
400 restart_in_window: u32,
401 delay: Duration,
402}
403
404fn daemon_will_restart(
411 state: &mut SupervisorSnapshot,
412 policy: &RestartPolicy,
413 now: Instant,
414) -> bool {
415 state.enabled && state.crash_restarts_in_window(policy.window, now) < policy.max_restarts
416}
417
418const DEFAULT_HEALTH_CADENCE: Duration = Duration::from_secs(30);
419const DEFAULT_HEALTH_DEADLINE: Duration = Duration::from_secs(5);
420const DEFAULT_HEALTH_FAILURE_THRESHOLD: u32 = 3;
421const MAX_HEALTH_METRICS_BYTES: usize = 16 * 1024;
422
423#[derive(Debug, Clone, Copy, PartialEq, Eq)]
424pub enum HealthAction {
425 Report,
426 Restart,
427 Alert,
428}
429
430impl fmt::Display for HealthAction {
431 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
432 f.write_str(match self {
433 Self::Report => "report",
434 Self::Restart => "restart",
435 Self::Alert => "alert",
436 })
437 }
438}
439
440#[derive(Debug, Clone, Copy, PartialEq, Eq)]
441pub struct HealthConfig {
442 pub cadence: Duration,
443 pub deadline: Duration,
444 pub failure_threshold: u32,
445 pub on_degraded: HealthAction,
446 pub on_failing: HealthAction,
447 pub critical: bool,
448}
449
450impl Default for HealthConfig {
451 fn default() -> Self {
452 Self {
453 cadence: DEFAULT_HEALTH_CADENCE,
454 deadline: DEFAULT_HEALTH_DEADLINE,
455 failure_threshold: DEFAULT_HEALTH_FAILURE_THRESHOLD,
456 on_degraded: HealthAction::Report,
457 on_failing: HealthAction::Report,
458 critical: false,
459 }
460 }
461}
462
463#[derive(Debug, Clone, PartialEq)]
481pub struct ModuleHealthStatus {
482 pub status: SupervisorHealthStatus,
483 pub last_probe_ms: Option<u64>,
484 pub detail: Option<String>,
485 pub metrics: Option<Value>,
486 pub consecutive_failures: u32,
487 pub late_answer_count: u64,
490 pub last_late_answer_latency_ms: Option<u64>,
492 pub last_action: Option<String>,
493 pub last_action_ms: Option<u64>,
497}
498
499impl Default for ModuleHealthStatus {
500 fn default() -> Self {
501 Self {
502 status: SupervisorHealthStatus::Unknown,
503 last_probe_ms: None,
504 detail: None,
505 metrics: None,
506 consecutive_failures: 0,
507 late_answer_count: 0,
508 last_late_answer_latency_ms: None,
509 last_action: None,
510 last_action_ms: None,
511 }
512 }
513}
514
515#[derive(Debug, Clone, Copy, PartialEq, Eq)]
517pub enum ModuleState {
518 Starting,
519 Running,
520 Unresponsive,
521 Restarting,
522 Draining,
523 Stopped,
524 Failed,
525 Disabled,
526}
527
528impl fmt::Display for ModuleState {
529 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
530 f.write_str(match self {
531 Self::Starting => "starting",
532 Self::Running => "running",
533 Self::Unresponsive => "unresponsive",
534 Self::Restarting => "restarting",
535 Self::Draining => "draining",
536 Self::Stopped => "stopped",
537 Self::Failed => "failed",
538 Self::Disabled => "disabled",
539 })
540 }
541}
542
543#[derive(Debug, Clone, Copy, PartialEq, Eq)]
545pub enum ExitKind {
546 Clean,
547 Crash,
548 DeliberateSeverance,
549}
550
551impl From<ExitKind> for TerminalExitKind {
552 fn from(kind: ExitKind) -> Self {
553 match kind {
554 ExitKind::Clean => Self::Clean,
555 ExitKind::Crash => Self::Crash,
556 ExitKind::DeliberateSeverance => Self::DeliberateSeverance,
557 }
558 }
559}
560
561#[derive(Debug, Clone, Copy, PartialEq, Eq)]
564pub(crate) struct ProcessIdentity {
565 pub(crate) pid: u32,
566 pub(crate) start_time: u64,
567}
568
569#[derive(Debug, Clone, PartialEq, Eq)]
571pub struct ExitReport {
572 pub kind: ExitKind,
573 pub code: Option<i32>,
574 pub signal: Option<i32>,
575 pub at_ms: u64,
576}
577
578#[derive(Debug, Clone, PartialEq)]
581pub struct ModuleStatus {
582 pub module_id: String,
583 pub state: ModuleState,
584 pub enabled: bool,
585 pub process_alive: bool,
586 pub registration_active: bool,
587 pub protocol: ModuleProtocol,
590 pub live: bool,
601 pub restart_count: u32,
605 pub lifetime_restarts: u32,
609 pub spawn_generation: u64,
610 pub max_restarts: u32,
615 pub restart_window: Duration,
619 pub drain_timeout: Duration,
623 pub restart_backoff: Duration,
624 pub restart_max_backoff: Duration,
625 pub pid: Option<u32>,
626 pub spawned_at_ms: Option<u64>,
627 pub spawned_from: Option<PathBuf>,
628 pub process_start_time: Option<u64>,
629 pub last_exit: Option<ExitReport>,
630 pub health: ModuleHealthStatus,
631}
632
633#[derive(Debug, Clone, PartialEq)]
634struct SupervisorSnapshot {
635 state: ModuleState,
636 enabled: bool,
637 process_alive: bool,
638 crash_restarts: VecDeque<Instant>,
644 lifetime_restarts: u32,
645 spawn_generation: u64,
654 pid: Option<u32>,
655 spawned_at_ms: Option<u64>,
656 spawned_from: Option<PathBuf>,
657 spawned_file_identity: Option<SpawnedFileIdentity>,
658 process_start_time: Option<u64>,
659 deliberate_severance: Option<ProcessIdentity>,
660 last_exit: Option<ExitReport>,
661 health: ModuleHealthStatus,
662 in_alternate_slot: bool,
667}
668
669impl SupervisorSnapshot {
670 fn starting() -> Self {
671 Self::new(ModuleState::Starting, true)
672 }
673
674 fn disabled() -> Self {
675 Self::new(ModuleState::Disabled, false)
676 }
677
678 fn failed() -> Self {
679 Self::new(ModuleState::Failed, true)
680 }
681
682 fn crash_restarts_in_window(&mut self, window: Duration, now: Instant) -> u32 {
686 while let Some(oldest) = self.crash_restarts.front() {
687 if now.duration_since(*oldest) > window {
688 self.crash_restarts.pop_front();
689 } else {
690 break;
691 }
692 }
693 u32::try_from(self.crash_restarts.len()).unwrap_or(u32::MAX)
694 }
695
696 fn record_crash_restart(&mut self, policy: &RestartPolicy, now: Instant) {
702 self.crash_restarts.push_back(now);
703 while self.crash_restarts.len() > policy.max_restarts as usize {
704 self.crash_restarts.pop_front();
705 }
706 self.lifetime_restarts += 1;
707 }
708
709 fn next_crash_restart(
713 &mut self,
714 policy: &RestartPolicy,
715 now: Instant,
716 ) -> Option<CrashRestartSchedule> {
717 let restart_in_window = self.crash_restarts_in_window(policy.window, now);
718 if restart_in_window >= policy.max_restarts {
719 return None;
720 }
721 self.record_crash_restart(policy, now);
722 Some(CrashRestartSchedule {
723 restart_in_window,
724 delay: policy.delay_for_restart(restart_in_window),
725 })
726 }
727
728 fn clear_crash_restarts(&mut self) {
733 self.crash_restarts.clear();
734 }
735
736 fn new(state: ModuleState, enabled: bool) -> Self {
737 Self {
738 state,
739 enabled,
740 process_alive: false,
741 crash_restarts: VecDeque::new(),
742 lifetime_restarts: 0,
743 spawn_generation: 0,
744 pid: None,
745 spawned_at_ms: None,
746 spawned_from: None,
747 spawned_file_identity: None,
748 process_start_time: None,
749 deliberate_severance: None,
750 last_exit: None,
751 health: ModuleHealthStatus::default(),
752 in_alternate_slot: false,
753 }
754 }
755}
756
757type SharedSnapshot = Arc<Mutex<SupervisorSnapshot>>;
758
759type SpawnSubscriberKey = (ConnectionId, u64);
760
761#[derive(Debug)]
762struct SpawnSubscriber {
763 version: u8,
764 frames: mpsc::Sender<Frame>,
765 lagged: Option<oneshot::Sender<SpawnCursor>>,
769}
770
771#[derive(Debug)]
772struct SpawnEventState {
773 daemon_incarnation: String,
774 seq: u64,
775 capacity: usize,
776 live: HashMap<String, LiveSpawn>,
777 generations: HashMap<String, u64>,
778 events: VecDeque<SpawnEvent>,
779 subscribers: HashMap<SpawnSubscriberKey, SpawnSubscriber>,
780}
781
782impl Default for SpawnEventState {
783 fn default() -> Self {
784 Self {
785 daemon_incarnation: "unconfigured".to_string(),
786 seq: 0,
787 capacity: SPAWN_EVENT_RING_CAPACITY,
788 live: HashMap::new(),
789 generations: HashMap::new(),
790 events: VecDeque::new(),
791 subscribers: HashMap::new(),
792 }
793 }
794}
795
796#[derive(Debug, Clone, Default)]
797struct SpawnEventFeed(Arc<Mutex<SpawnEventState>>);
798
799#[derive(Debug, Clone, PartialEq, Eq)]
800pub(crate) enum SpawnSubscribeRefusal {
801 ForeignIncarnation { current: String },
802 TooOld { oldest: SpawnCursor },
803 Frame(String),
804}
805
806impl SpawnEventFeed {
807 fn configure_incarnation(&self, daemon_incarnation: String) {
808 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
809 state.daemon_incarnation = daemon_incarnation;
810 state.seq = 0;
811 state.live.clear();
812 state.generations.clear();
813 state.events.clear();
814 state.subscribers.clear();
815 }
816
817 fn cursor(state: &SpawnEventState) -> SpawnCursor {
818 SpawnCursor {
819 daemon_incarnation: state.daemon_incarnation.clone(),
820 seq: state.seq,
821 }
822 }
823
824 fn snapshot(&self) -> SpawnSnapshot {
825 let state = self.0.lock().unwrap_or_else(|p| p.into_inner());
826 let mut live = state.live.values().cloned().collect::<Vec<_>>();
827 live.sort_by(|left, right| left.module_id.cmp(&right.module_id));
828 SpawnSnapshot {
829 cursor: Self::cursor(&state),
830 ring_bound: state.capacity as u64,
831 live,
832 }
833 }
834
835 fn emit_spawned(&self, module_id: &str, pid: u32, spawned_at_ms: u64) -> u64 {
836 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
837 let generation = state
838 .generations
839 .get(module_id)
840 .copied()
841 .unwrap_or(0)
842 .checked_add(1)
843 .expect("spawn generation exhausted");
844 state.generations.insert(module_id.to_string(), generation);
845 let live = LiveSpawn {
846 module_id: module_id.to_string(),
847 spawn_generation: generation,
848 pid,
849 spawned_at_ms,
850 };
851 state.live.insert(module_id.to_string(), live);
852 Self::emit_locked(
853 &mut state,
854 SpawnEventKind::Spawned,
855 module_id.to_string(),
856 generation,
857 pid,
858 None,
859 None,
860 );
861 generation
862 }
863
864 fn emit_exited(&self, module_id: &str, exit_code: Option<i32>, exit_signal: Option<i32>) {
865 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
866 let Some(live) = state.live.remove(module_id) else {
867 warn!(
868 module_id,
869 "terminal record had no live spawn event identity"
870 );
871 return;
872 };
873 Self::emit_locked(
874 &mut state,
875 SpawnEventKind::Exited,
876 module_id.to_string(),
877 live.spawn_generation,
878 live.pid,
879 exit_code,
880 exit_signal,
881 );
882 }
883
884 fn emit_superseded_exited(
891 &self,
892 module_id: &str,
893 spawn_generation: u64,
894 pid: u32,
895 exit_code: Option<i32>,
896 exit_signal: Option<i32>,
897 ) {
898 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
899 if state
900 .live
901 .get(module_id)
902 .is_some_and(|live| live.spawn_generation == spawn_generation)
903 {
904 state.live.remove(module_id);
905 }
906 Self::emit_locked(
907 &mut state,
908 SpawnEventKind::Exited,
909 module_id.to_string(),
910 spawn_generation,
911 pid,
912 exit_code,
913 exit_signal,
914 );
915 }
916
917 #[allow(clippy::too_many_arguments)]
918 fn emit_locked(
919 state: &mut SpawnEventState,
920 kind: SpawnEventKind,
921 module_id: String,
922 spawn_generation: u64,
923 pid: u32,
924 exit_code: Option<i32>,
925 exit_signal: Option<i32>,
926 ) {
927 state.seq = state
928 .seq
929 .checked_add(1)
930 .expect("spawn event sequence exhausted");
931 let event = SpawnEvent {
932 cursor: Self::cursor(state),
933 kind,
934 module_id,
935 spawn_generation,
936 pid,
937 exit_code,
938 exit_signal,
939 };
940 state.events.push_back(event.clone());
941 while state.events.len() > state.capacity {
942 state.events.pop_front();
943 }
944 let body = match serde_json::to_vec(&event) {
945 Ok(body) => body,
946 Err(error) => {
947 error!(%error, "failed to serialize supervisor spawn event");
948 return;
949 }
950 };
951 state.subscribers.retain(|(connection_id, corr), subscriber| {
952 let frame = Frame::build_with_version(
953 subscriber.version,
954 FrameType::StreamData,
955 control_flags(),
956 0,
957 0,
958 *corr,
959 body.clone(),
960 );
961 match frame {
962 Ok(frame) => {
963 if subscriber.frames.try_send(frame).is_ok() {
964 true
965 } else {
966 warn!(connection_id = connection_id.get(), corr, "dropping lagged supervisor spawn subscriber");
967 if let Some(lagged) = subscriber.lagged.take() {
968 let _ = lagged.send(event.cursor.clone());
969 }
970 false
971 }
972 }
973 Err(error) => {
974 warn!(connection_id = connection_id.get(), corr, %error, "dropping supervisor spawn subscriber after frame build failure");
975 false
976 }
977 }
978 });
979 }
980
981 fn subscribe(
982 &self,
983 connection_id: ConnectionId,
984 corr: u64,
985 version: u8,
986 since: Option<SpawnCursor>,
987 sink: FrameSink,
988 ) -> Result<(), SpawnSubscribeRefusal> {
989 let (frames, mut receiver) = mpsc::channel(SPAWN_SUBSCRIBER_BUFFER);
990 let (lagged, mut lagged_rx) = oneshot::channel::<SpawnCursor>();
991 {
992 let mut state = self.0.lock().unwrap_or_else(|p| p.into_inner());
993 let replay = if let Some(since) = since {
994 if since.daemon_incarnation != state.daemon_incarnation {
995 return Err(SpawnSubscribeRefusal::ForeignIncarnation {
996 current: state.daemon_incarnation.clone(),
997 });
998 }
999 if let Some(oldest) = state.events.front().map(|event| event.cursor.clone()) {
1000 if since.seq < oldest.seq.saturating_sub(1) {
1001 return Err(SpawnSubscribeRefusal::TooOld { oldest });
1002 }
1003 }
1004 state
1005 .events
1006 .iter()
1007 .filter(|event| event.cursor.seq > since.seq)
1008 .cloned()
1009 .collect::<Vec<_>>()
1010 } else {
1011 Vec::new()
1012 };
1013 for event in replay {
1014 let body = serde_json::to_vec(&event)
1015 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1016 let frame = Frame::build_with_version(
1017 version,
1018 FrameType::StreamData,
1019 control_flags(),
1020 0,
1021 0,
1022 corr,
1023 body,
1024 )
1025 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1026 frames
1027 .try_send(frame)
1028 .map_err(|error| SpawnSubscribeRefusal::Frame(error.to_string()))?;
1029 }
1030 state.subscribers.insert(
1031 (connection_id, corr),
1032 SpawnSubscriber {
1033 version,
1034 frames,
1035 lagged: Some(lagged),
1036 },
1037 );
1038 }
1039 tokio::spawn(async move {
1050 while let Some(frame) = receiver.recv().await {
1051 if sink.send(frame).await.is_err() {
1052 return;
1053 }
1054 }
1055 let Ok(first_undelivered) = lagged_rx.try_recv() else {
1056 return;
1057 };
1058 match spawn_subscriber_lagged_frame(version, corr, first_undelivered) {
1059 Ok(frame) => {
1060 let _ = sink.send(frame).await;
1061 }
1062 Err(error) => {
1063 error!(%error, corr, "failed to build lagged spawn subscriber terminal frame");
1064 }
1065 }
1066 });
1067 Ok(())
1068 }
1069
1070 fn cancel(&self, connection_id: ConnectionId, corr: u64) -> bool {
1071 let Some(subscriber) = self
1072 .0
1073 .lock()
1074 .unwrap_or_else(|p| p.into_inner())
1075 .subscribers
1076 .remove(&(connection_id, corr))
1077 else {
1078 return false;
1079 };
1080 if let Ok(frame) = Frame::build_with_version(
1081 subscriber.version,
1082 FrameType::StreamEnd,
1083 control_flags(),
1084 0,
1085 0,
1086 corr,
1087 Vec::new(),
1088 ) {
1089 tokio::spawn(async move {
1090 let _ = subscriber.frames.send(frame).await;
1091 });
1092 }
1093 true
1094 }
1095
1096 fn remove_connection(&self, connection_id: ConnectionId) {
1097 self.0
1098 .lock()
1099 .unwrap_or_else(|p| p.into_inner())
1100 .subscribers
1101 .retain(|(subscriber_connection, _), _| *subscriber_connection != connection_id);
1102 }
1103
1104 #[cfg(any(test, feature = "test-support"))]
1105 fn set_capacity(&self, capacity: usize) {
1106 self.0.lock().unwrap_or_else(|p| p.into_inner()).capacity = capacity;
1107 }
1108
1109 #[cfg(any(test, feature = "test-support"))]
1110 fn subscriber_count(&self) -> usize {
1111 self.0
1112 .lock()
1113 .unwrap_or_else(|p| p.into_inner())
1114 .subscribers
1115 .len()
1116 }
1117}
1118
1119fn spawn_subscriber_lagged_frame(
1122 version: u8,
1123 corr: u64,
1124 first_undelivered: SpawnCursor,
1125) -> Result<Frame, String> {
1126 let body = serde_json::to_vec(&subc_protocol::ErrorBody {
1127 code: SPAWN_SUBSCRIBER_LAGGED_CODE.to_string(),
1128 message: "spawn subscriber fell behind and was dropped; resubscribe from the last cursor received"
1129 .to_string(),
1130 detail: Some(serde_json::json!({
1131 "first_undelivered_cursor": first_undelivered
1132 })),
1133 })
1134 .map_err(|error| error.to_string())?;
1135 Frame::build_with_version(version, FrameType::Error, control_flags(), 0, 0, corr, body)
1136 .map_err(|error| error.to_string())
1137}
1138
1139pub trait ModuleProcessLiveness: Send + Sync {
1140 fn process_live(&self, module_id: &str) -> Option<bool>;
1141}
1142
1143#[derive(Debug, Clone, Default)]
1145pub struct SupervisorProcessLiveness {
1146 snapshots: Arc<Mutex<HashMap<String, SharedSnapshot>>>,
1147}
1148
1149impl SupervisorProcessLiveness {
1150 pub fn new() -> Self {
1151 Self::default()
1152 }
1153
1154 fn track(&self, module_id: String, snapshot: SharedSnapshot) {
1155 let mut snapshots = self
1156 .snapshots
1157 .lock()
1158 .unwrap_or_else(|poisoned| poisoned.into_inner());
1159 snapshots.insert(module_id, snapshot);
1160 }
1161
1162 fn untrack_if_current(&self, module_id: &str, snapshot: &SharedSnapshot) {
1163 let mut snapshots = self
1164 .snapshots
1165 .lock()
1166 .unwrap_or_else(|poisoned| poisoned.into_inner());
1167 let is_current = snapshots
1168 .get(module_id)
1169 .map(|tracked| Arc::ptr_eq(tracked, snapshot))
1170 .unwrap_or(false);
1171 if is_current {
1172 snapshots.remove(module_id);
1173 }
1174 }
1175}
1176
1177impl ModuleProcessLiveness for SupervisorProcessLiveness {
1178 fn process_live(&self, module_id: &str) -> Option<bool> {
1179 let snapshot = {
1180 let snapshots = self
1181 .snapshots
1182 .lock()
1183 .unwrap_or_else(|poisoned| poisoned.into_inner());
1184 snapshots.get(module_id).cloned()
1185 }?;
1186 let snapshot = snapshot
1187 .lock()
1188 .unwrap_or_else(|poisoned| poisoned.into_inner());
1189 Some(snapshot.state == ModuleState::Running && snapshot.process_alive)
1190 }
1191}
1192
1193#[derive(Debug, Clone)]
1194struct SupervisorRuntimeConfig {
1195 restart_policy: RestartPolicy,
1196 drain_timeout: Duration,
1199 effective_drain_timeout: Arc<Mutex<Duration>>,
1202 default_drain_timeout: Duration,
1205 health: HealthConfig,
1206 connection_file_path: Option<PathBuf>,
1207 capture_logs_dir: Option<PathBuf>,
1208 forwarding: Option<Arc<ForwardingTable>>,
1209 supervisor_handle: Option<SupervisorHandle>,
1212 stderr_ring: Arc<Mutex<StderrRing>>,
1219 terminal_ring: Arc<Mutex<TerminalRing>>,
1220 spawn_events: SpawnEventFeed,
1221 child_roster: ChildRoster,
1222 #[cfg(target_os = "linux")]
1223 cgroup_placement: Option<subc_cgroup::Placement>,
1224 #[cfg(test)]
1225 test_seed_stale_facts_before_enable_spawn: bool,
1226}
1227
1228#[derive(Debug, Clone, PartialEq, Eq)]
1229struct SupervisedConfiguration {
1230 spec: ModuleSpec,
1231 health: HealthConfig,
1232}
1233
1234#[derive(Debug, Clone, Default)]
1240pub struct SupervisorHandle {
1241 modules: Arc<Mutex<HashMap<String, SupervisedModule>>>,
1242 spawn_events: SpawnEventFeed,
1243 reserved_nonces: Arc<Mutex<HashMap<String, Option<String>>>>,
1254 removal_tombstones: Arc<Mutex<HashMap<String, u64>>>,
1260 spawn_nonces: Arc<Mutex<HashMap<String, String>>>,
1264 reserved_prefix_owners: Arc<Mutex<HashMap<String, String>>>,
1272 swaps: Arc<Mutex<HashMap<String, OpenSwap>>>,
1278 promotion_observer: PromotionObserverSlot,
1280 operation_lock: Arc<AsyncMutex<()>>,
1284}
1285
1286pub(crate) trait SwapPromotionObserver: Send + Sync {
1295 fn swap_promoted(&self, registration: &crate::registry::ModuleRegistration);
1296}
1297
1298#[derive(Clone, Default)]
1302struct PromotionObserverSlot(Arc<Mutex<Option<std::sync::Weak<dyn SwapPromotionObserver>>>>);
1303
1304impl fmt::Debug for PromotionObserverSlot {
1305 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1306 f.write_str("PromotionObserverSlot")
1307 }
1308}
1309
1310#[derive(Debug, Clone)]
1312struct OpenSwap {
1313 candidate_nonce: String,
1316 incumbent_nonce: Option<String>,
1321 candidate_admitted: bool,
1325}
1326
1327#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1330pub(crate) enum SwapHelloAdmission {
1331 NotSwapping,
1334 Candidate,
1336 Refused,
1339}
1340
1341#[derive(Debug, Clone, PartialEq, Eq)]
1342pub(crate) enum ReservedHelloRejection {
1343 Exact {
1344 module_id: String,
1345 },
1346 Prefix {
1347 prefix: String,
1348 owner_module_id: String,
1349 },
1350}
1351
1352impl SupervisorHandle {
1353 pub fn new() -> Self {
1354 Self::default()
1355 }
1356
1357 pub(crate) fn spawn_snapshot(&self) -> SpawnSnapshot {
1358 self.spawn_events.snapshot()
1359 }
1360
1361 pub(crate) fn subscribe_spawns(
1362 &self,
1363 connection_id: ConnectionId,
1364 corr: u64,
1365 version: u8,
1366 since: Option<SpawnCursor>,
1367 sink: FrameSink,
1368 ) -> Result<(), SpawnSubscribeRefusal> {
1369 self.spawn_events
1370 .subscribe(connection_id, corr, version, since, sink)
1371 }
1372
1373 pub(crate) fn cancel_spawn_subscription(&self, connection_id: ConnectionId, corr: u64) -> bool {
1374 self.spawn_events.cancel(connection_id, corr)
1375 }
1376
1377 pub(crate) fn remove_spawn_subscribers(&self, connection_id: ConnectionId) {
1378 self.spawn_events.remove_connection(connection_id);
1379 }
1380
1381 #[cfg(any(test, feature = "test-support"))]
1382 pub fn set_spawn_event_capacity_for_test(&self, capacity: usize) {
1383 assert!(capacity > 0, "spawn event capacity must be non-zero");
1384 self.spawn_events.set_capacity(capacity);
1385 }
1386
1387 #[cfg(any(test, feature = "test-support"))]
1388 pub fn spawn_subscriber_count_for_test(&self) -> usize {
1389 self.spawn_events.subscriber_count()
1390 }
1391
1392 pub fn set_spawn_nonce(&self, module_id: &str, nonce: String) {
1395 self.spawn_nonces
1396 .lock()
1397 .unwrap_or_else(|poisoned| poisoned.into_inner())
1398 .insert(module_id.to_string(), nonce);
1399 }
1400
1401 pub fn set_reserved_nonce(&self, module_id: &str, nonce: String) {
1404 self.reserved_nonces
1405 .lock()
1406 .unwrap_or_else(|poisoned| poisoned.into_inner())
1407 .insert(module_id.to_string(), Some(nonce));
1408 }
1409
1410 pub fn set_reserved_prefixes(&self, owner_module_id: &str, prefixes: &[String]) {
1412 let mut owners = self
1413 .reserved_prefix_owners
1414 .lock()
1415 .unwrap_or_else(|poisoned| poisoned.into_inner());
1416 owners.retain(|_, owner| owner != owner_module_id);
1417 for prefix in prefixes {
1418 owners.insert(prefix.clone(), owner_module_id.to_string());
1419 }
1420 }
1421
1422 #[cfg(test)]
1424 pub(crate) fn spawn_nonce(&self, module_id: &str) -> Option<String> {
1425 self.spawn_nonces
1426 .lock()
1427 .unwrap_or_else(|poisoned| poisoned.into_inner())
1428 .get(module_id)
1429 .cloned()
1430 }
1431
1432 fn apply_identity_configuration(&self, spec: &ModuleSpec) {
1433 self.set_reserved_prefixes(&spec.module_id, &spec.reserved_prefixes);
1434 let spawn_nonce = self
1435 .spawn_nonces
1436 .lock()
1437 .unwrap_or_else(|poisoned| poisoned.into_inner())
1438 .get(&spec.module_id)
1439 .cloned();
1440 let mut reserved_nonces = self
1441 .reserved_nonces
1442 .lock()
1443 .unwrap_or_else(|poisoned| poisoned.into_inner());
1444 if spec.reserved {
1445 reserved_nonces.insert(spec.module_id.clone(), spawn_nonce);
1450 }
1451 drop(reserved_nonces);
1452 self.removal_tombstones
1456 .lock()
1457 .unwrap_or_else(|poisoned| poisoned.into_inner())
1458 .remove(&spec.module_id);
1459 }
1460
1461 pub fn reserved_hello_authorized(&self, module_id: &str, presented: Option<&str>) -> bool {
1466 self.reserved_hello_rejection(module_id, presented)
1467 .is_none()
1468 }
1469
1470 pub(crate) fn reserved_hello_rejection(
1471 &self,
1472 module_id: &str,
1473 presented: Option<&str>,
1474 ) -> Option<ReservedHelloRejection> {
1475 let nonces = self
1476 .reserved_nonces
1477 .lock()
1478 .unwrap_or_else(|poisoned| poisoned.into_inner());
1479 if let Some(expected) = nonces.get(module_id) {
1480 let authorized = match expected {
1484 Some(expected) => {
1485 presented.is_some_and(|p| constant_time_eq(expected.as_bytes(), p.as_bytes()))
1486 }
1487 None => false,
1488 };
1489 if authorized {
1490 return None;
1491 }
1492 return Some(ReservedHelloRejection::Exact {
1493 module_id: module_id.to_string(),
1494 });
1495 }
1496 drop(nonces);
1497
1498 let matched_prefix = self
1499 .reserved_prefix_owners
1500 .lock()
1501 .unwrap_or_else(|poisoned| poisoned.into_inner())
1502 .iter()
1503 .filter(|(prefix, _)| module_id.starts_with(prefix.as_str()))
1504 .max_by_key(|(prefix, _)| prefix.len())
1505 .map(|(prefix, owner)| (prefix.clone(), owner.clone()));
1506 let (prefix, owner_module_id) = matched_prefix?;
1507
1508 let authorized = presented.is_some_and(|presented| {
1509 self.spawn_nonces
1510 .lock()
1511 .unwrap_or_else(|poisoned| poisoned.into_inner())
1512 .get(&owner_module_id)
1513 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()))
1514 || self.swap_nonce_matches(&owner_module_id, presented)
1517 });
1518 if authorized {
1519 None
1520 } else {
1521 Some(ReservedHelloRejection::Prefix {
1522 prefix,
1523 owner_module_id,
1524 })
1525 }
1526 }
1527
1528 pub fn spawned_consumer_authorized(&self, module_id: &str, presented: &str) -> bool {
1533 if presented.is_empty() {
1534 return false;
1535 }
1536 let nonces = self
1537 .spawn_nonces
1538 .lock()
1539 .unwrap_or_else(|poisoned| poisoned.into_inner());
1540 let current = nonces
1541 .get(module_id)
1542 .is_some_and(|expected| constant_time_eq(expected.as_bytes(), presented.as_bytes()));
1543 drop(nonces);
1544 current || self.swap_nonce_matches(module_id, presented)
1549 }
1550
1551 fn swap_nonce_matches(&self, module_id: &str, presented: &str) -> bool {
1553 let swaps = self
1554 .swaps
1555 .lock()
1556 .unwrap_or_else(|poisoned| poisoned.into_inner());
1557 swaps.get(module_id).is_some_and(|swap| {
1558 constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes())
1559 || swap.incumbent_nonce.as_deref().is_some_and(|incumbent| {
1560 constant_time_eq(incumbent.as_bytes(), presented.as_bytes())
1561 })
1562 })
1563 }
1564
1565 pub(crate) fn open_swap(&self, module_id: &str, candidate_nonce: String) {
1568 let incumbent_nonce = self
1569 .spawn_nonces
1570 .lock()
1571 .unwrap_or_else(|poisoned| poisoned.into_inner())
1572 .get(module_id)
1573 .cloned();
1574 self.swaps
1575 .lock()
1576 .unwrap_or_else(|poisoned| poisoned.into_inner())
1577 .insert(
1578 module_id.to_string(),
1579 OpenSwap {
1580 candidate_nonce,
1581 incumbent_nonce,
1582 candidate_admitted: false,
1583 },
1584 );
1585 }
1586
1587 pub(crate) fn close_swap(&self, module_id: &str) {
1590 self.swaps
1591 .lock()
1592 .unwrap_or_else(|poisoned| poisoned.into_inner())
1593 .remove(module_id);
1594 }
1595
1596 pub(crate) fn set_swap_promotion_observer(
1599 &self,
1600 observer: std::sync::Weak<dyn SwapPromotionObserver>,
1601 ) {
1602 *self
1603 .promotion_observer
1604 .0
1605 .lock()
1606 .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(observer);
1607 }
1608
1609 fn notify_swap_promoted(&self, registration: &crate::registry::ModuleRegistration) {
1612 let observer = self
1613 .promotion_observer
1614 .0
1615 .lock()
1616 .unwrap_or_else(|poisoned| poisoned.into_inner())
1617 .as_ref()
1618 .and_then(std::sync::Weak::upgrade);
1619 if let Some(observer) = observer {
1620 observer.swap_promoted(registration);
1621 }
1622 }
1623
1624 pub(crate) fn swap_open(&self, module_id: &str) -> bool {
1626 self.swaps
1627 .lock()
1628 .unwrap_or_else(|poisoned| poisoned.into_inner())
1629 .contains_key(module_id)
1630 }
1631
1632 fn promote_swap_nonce(&self, module_id: &str, reserved: bool) {
1637 let candidate_nonce = self
1638 .swaps
1639 .lock()
1640 .unwrap_or_else(|poisoned| poisoned.into_inner())
1641 .get(module_id)
1642 .map(|swap| swap.candidate_nonce.clone());
1643 let Some(nonce) = candidate_nonce else {
1644 return;
1645 };
1646 self.set_spawn_nonce(module_id, nonce.clone());
1647 if reserved {
1648 self.set_reserved_nonce(module_id, nonce);
1649 }
1650 }
1651
1652 pub(crate) fn swap_hello_admission(
1667 &self,
1668 module_id: &str,
1669 presented: Option<&str>,
1670 ) -> SwapHelloAdmission {
1671 let swaps = self
1672 .swaps
1673 .lock()
1674 .unwrap_or_else(|poisoned| poisoned.into_inner());
1675 let Some(swap) = swaps.get(module_id) else {
1676 return SwapHelloAdmission::NotSwapping;
1677 };
1678 let Some(presented) = presented else {
1679 return SwapHelloAdmission::Refused;
1680 };
1681 if constant_time_eq(swap.candidate_nonce.as_bytes(), presented.as_bytes()) {
1682 return if swap.candidate_admitted {
1683 SwapHelloAdmission::Refused
1684 } else {
1685 SwapHelloAdmission::Candidate
1686 };
1687 }
1688 if swap
1689 .incumbent_nonce
1690 .as_deref()
1691 .is_some_and(|incumbent| constant_time_eq(incumbent.as_bytes(), presented.as_bytes()))
1692 {
1693 return SwapHelloAdmission::NotSwapping;
1694 }
1695 SwapHelloAdmission::Refused
1696 }
1697
1698 pub(crate) fn mark_swap_candidate_admitted(&self, module_id: &str) {
1701 if let Some(swap) = self
1702 .swaps
1703 .lock()
1704 .unwrap_or_else(|poisoned| poisoned.into_inner())
1705 .get_mut(module_id)
1706 {
1707 swap.candidate_admitted = true;
1708 }
1709 }
1710
1711 pub fn spawn_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1713 self.spawn_nonces
1714 .lock()
1715 .unwrap_or_else(|poisoned| poisoned.into_inner())
1716 .get(module_id)
1717 .cloned()
1718 }
1719
1720 pub fn reserved_launch_nonce_for(&self, module_id: &str) -> Option<String> {
1722 self.reserved_nonces
1723 .lock()
1724 .unwrap_or_else(|poisoned| poisoned.into_inner())
1725 .get(module_id)
1726 .cloned()
1727 .flatten()
1728 }
1729
1730 pub fn insert(&self, module: SupervisedModule) -> Option<SupervisedModule> {
1731 let mut modules = self
1732 .modules
1733 .lock()
1734 .unwrap_or_else(|poisoned| poisoned.into_inner());
1735 modules.insert(module.module_id().to_string(), module)
1736 }
1737
1738 pub fn get(&self, module_id: &str) -> Option<SupervisedModule> {
1739 let modules = self
1740 .modules
1741 .lock()
1742 .unwrap_or_else(|poisoned| poisoned.into_inner());
1743 modules.get(module_id).cloned()
1744 }
1745
1746 pub(crate) fn record_late_health_answer(
1747 &self,
1748 module_id: &str,
1749 latency_ms: u64,
1750 ) -> Result<bool, SuperviseError> {
1751 let Some(module) = self.get(module_id) else {
1752 return Ok(false);
1753 };
1754 update_snapshot(&module.inner.snapshot, Some(module_id), |state| {
1755 state.health.late_answer_count = state.health.late_answer_count.saturating_add(1);
1756 state.health.last_late_answer_latency_ms = Some(latency_ms);
1757 state.health.consecutive_failures = 0;
1765 })?;
1766 Ok(true)
1767 }
1768
1769 pub fn record_deliberate_severance(&self, module_id: &str) -> Result<bool, SuperviseError> {
1775 let Some(module) = self.get(module_id) else {
1776 return Ok(false);
1777 };
1778 let status = module.status()?;
1779 let Some((pid, start_time)) = status.pid.zip(status.process_start_time) else {
1780 return Ok(false);
1781 };
1782 module.record_deliberate_severance(ProcessIdentity { pid, start_time })
1783 }
1784
1785 pub fn list(&self) -> Vec<SupervisedModule> {
1786 let modules = self
1787 .modules
1788 .lock()
1789 .unwrap_or_else(|poisoned| poisoned.into_inner());
1790 let mut modules = modules.values().cloned().collect::<Vec<_>>();
1791 modules.sort_by(|left, right| left.module_id().cmp(right.module_id()));
1792 modules
1793 }
1794
1795 pub(crate) fn retire(&self, module_id: &str) -> Option<SupervisedModule> {
1796 self.spawn_nonces
1797 .lock()
1798 .unwrap_or_else(|poisoned| poisoned.into_inner())
1799 .remove(module_id);
1800 self.close_swap(module_id);
1801 let mut reserved_nonces = self
1802 .reserved_nonces
1803 .lock()
1804 .unwrap_or_else(|poisoned| poisoned.into_inner());
1805 if reserved_nonces.contains_key(module_id) {
1806 reserved_nonces.insert(module_id.to_string(), None);
1809 }
1810 drop(reserved_nonces);
1811 self.reserved_prefix_owners
1812 .lock()
1813 .unwrap_or_else(|poisoned| poisoned.into_inner())
1814 .retain(|_, owner| owner != module_id);
1815 self.modules
1816 .lock()
1817 .unwrap_or_else(|poisoned| poisoned.into_inner())
1818 .remove(module_id)
1819 }
1820
1821 pub(crate) fn record_rescan_removal(&self, module_id: &str) {
1824 self.removal_tombstones
1825 .lock()
1826 .unwrap_or_else(|poisoned| poisoned.into_inner())
1827 .insert(module_id.to_string(), unix_ms_now());
1828 }
1829
1830 pub(crate) fn removal_tombstone_age_ms(&self, module_id: &str) -> Option<u64> {
1832 self.removal_tombstones
1833 .lock()
1834 .unwrap_or_else(|poisoned| poisoned.into_inner())
1835 .get(module_id)
1836 .copied()
1837 .map(|removed_at_ms| unix_ms_now().saturating_sub(removed_at_ms))
1838 }
1839
1840 pub(crate) fn release_retained_reserved_gate(&self, module_id: &str) -> bool {
1845 if self.get(module_id).is_some() {
1846 return false;
1847 }
1848 let mut reserved_nonces = self
1849 .reserved_nonces
1850 .lock()
1851 .unwrap_or_else(|poisoned| poisoned.into_inner());
1852 if !matches!(reserved_nonces.get(module_id), Some(None)) {
1853 return false;
1854 }
1855 reserved_nonces.remove(module_id);
1856 true
1857 }
1858
1859 pub(crate) fn operation_lock(&self) -> Arc<AsyncMutex<()>> {
1860 Arc::clone(&self.operation_lock)
1861 }
1862}
1863
1864#[derive(Debug, Clone)]
1866pub struct Supervisor {
1867 registry: Arc<Registry>,
1868 restart_policy: RestartPolicy,
1869 drain_timeout: Duration,
1870 connection_file_path: Option<PathBuf>,
1871 capture_logs_dir: Option<PathBuf>,
1872 forwarding: Option<Arc<ForwardingTable>>,
1873 process_liveness: Arc<SupervisorProcessLiveness>,
1874 supervisor_handle: Option<SupervisorHandle>,
1875 health: HealthConfig,
1876 daemon_start_clock: crate::clock::StartClock,
1877 terminal_journal: Option<Arc<crate::terminal_journal::TerminalJournal>>,
1878 spawn_events: SpawnEventFeed,
1879 provenance_probe: ExecutableIdentityProbe,
1880 child_roster: ChildRoster,
1883 #[cfg(target_os = "linux")]
1884 cgroup_placement: Option<subc_cgroup::Placement>,
1885}
1886
1887impl Supervisor {
1888 #[cfg(unix)]
1889 pub(crate) fn stamp_shutdown(&self) {
1890 if let Some(journal) = &self.terminal_journal {
1891 journal.stamp_shutdown();
1892 }
1893 }
1894
1895 #[cfg(unix)]
1899 pub(crate) async fn drain_for_daemon_shutdown(&self) -> Result<(), SuperviseError> {
1900 const NOTICE_BUDGET: Duration = Duration::from_millis(500);
1901 const DRAIN_BUDGET: Duration = Duration::from_secs(2);
1902 let Some(forwarding) = &self.forwarding else {
1903 return Ok(());
1904 };
1905 let module_ids = forwarding
1906 .begin_daemon_drain()
1907 .map_err(SuperviseError::Forwarding)?;
1908 let deadline_ms =
1909 unix_ms_now().saturating_add((NOTICE_BUDGET + DRAIN_BUDGET).as_millis() as u64);
1910 let mut notices = tokio::task::JoinSet::new();
1911 let mut drains = Vec::new();
1912 for module_id in module_ids {
1913 let Some(target) = forwarding
1914 .begin_module_drain(&module_id, RouteCloseReason::Restart)
1915 .map_err(SuperviseError::Forwarding)?
1916 else {
1917 continue;
1918 };
1919 let routes = forwarding
1920 .endpoint_routes(target.endpoint)
1921 .map_err(SuperviseError::Forwarding)?;
1922 let command = serde_json::to_vec(&ModuleControlCommand::Draining {
1926 reason: RouteCloseReason::Restart,
1927 deadline_ms,
1928 })
1929 .expect("module draining serializes");
1930 let closing = serde_json::to_vec(&ClientControlPush::RouteClosing {
1931 module_id: module_id.clone(),
1932 reason: RouteCloseReason::Restart,
1933 })
1934 .expect("route closing serializes");
1935 let mut recipients = vec![(target.sink.clone(), target.negotiated_ver, command)];
1936 let mut seen = std::collections::HashSet::new();
1937 for route in routes {
1938 let client = route.goodbye_target;
1939 if seen.insert(client.connection_id) {
1940 recipients.push((client.sink, client.negotiated_ver, closing.clone()));
1941 }
1942 }
1943 for (sink, version, body) in recipients {
1944 notices.spawn(async move {
1945 let frame = Frame::build_with_version(
1946 version,
1947 FrameType::Push,
1948 control_flags(),
1949 0,
1950 0,
1951 0,
1952 body,
1953 )
1954 .expect("bounded lifecycle notice frame builds");
1955 sink.send_flushed(frame).await
1956 });
1957 }
1958 let gauges = declared_busy_gauges(&self.registry, &module_id)?;
1959 drains.push((module_id, target.endpoint, gauges));
1960 }
1961 let notice_deadline = Instant::now() + NOTICE_BUDGET;
1964 while let Ok(Some(result)) = timeout_at(notice_deadline, notices.join_next()).await {
1965 if !matches!(result, Ok(Ok(()))) {
1966 warn!(?result, "daemon shutdown notice delivery failed");
1967 }
1968 }
1969 notices.abort_all();
1970 let deadline = Instant::now() + DRAIN_BUDGET;
1971 let mut waits = tokio::task::JoinSet::new();
1972 for (module_id, endpoint, gauges) in drains {
1973 let forwarding = Arc::clone(forwarding);
1974 let mut runtime = self.runtime_config();
1975 runtime.health.cadence = Duration::from_millis(100);
1976 waits.spawn(async move {
1977 wait_for_forwarding_quiescence(
1978 &forwarding,
1979 &module_id,
1980 &runtime,
1981 endpoint,
1982 deadline,
1983 &gauges,
1984 DrainScope::Active,
1985 )
1986 .await
1987 });
1988 }
1989 while let Ok(Some(result)) = timeout_at(deadline, waits.join_next()).await {
1990 if !matches!(result, Ok(Ok(true))) {
1991 warn!(?result, "daemon shutdown drain did not reach quiescence");
1992 }
1993 }
1994 Ok(())
1995 }
1996
1997 #[cfg(unix)]
2007 pub(crate) async fn end_children_for_daemon_shutdown(
2008 &self,
2009 already_escalated: bool,
2010 escalate: impl std::future::Future<Output = ()>,
2011 ) {
2012 if let Some(forwarding) = &self.forwarding {
2013 let closed = forwarding.close_all_connections(&CloseReason::new(
2014 "daemon_shutdown",
2015 "the daemon is exiting after its shutdown notice and drain",
2016 ));
2017 debug!(closed, "closed established connections for daemon shutdown");
2018 }
2019 crate::child_roster::end_children_for_daemon_shutdown(
2020 &self.child_roster,
2021 already_escalated,
2022 escalate,
2023 )
2024 .await;
2025 }
2026
2027 pub fn new(registry: Arc<Registry>, restart_policy: RestartPolicy) -> Self {
2028 Self {
2029 registry,
2030 restart_policy,
2031 drain_timeout: DEFAULT_DRAIN_TIMEOUT,
2032 connection_file_path: None,
2033 capture_logs_dir: None,
2034 forwarding: None,
2035 process_liveness: Arc::new(SupervisorProcessLiveness::default()),
2036 supervisor_handle: None,
2037 health: HealthConfig::default(),
2038 daemon_start_clock: crate::clock::StartClock::capture(),
2039 terminal_journal: None,
2040 spawn_events: SpawnEventFeed::default(),
2041 provenance_probe: ExecutableIdentityProbe::default(),
2042 child_roster: ChildRoster::default(),
2043 #[cfg(target_os = "linux")]
2044 cgroup_placement: None,
2045 }
2046 }
2047
2048 pub fn with_drain_timeout(mut self, drain_timeout: Duration) -> Self {
2049 self.drain_timeout = drain_timeout;
2050 self
2051 }
2052
2053 pub fn with_process_liveness(
2054 mut self,
2055 process_liveness: Arc<SupervisorProcessLiveness>,
2056 ) -> Self {
2057 self.process_liveness = process_liveness;
2058 self
2059 }
2060
2061 pub fn with_connection_file_path(mut self, connection_file_path: impl Into<PathBuf>) -> Self {
2062 self.connection_file_path = Some(connection_file_path.into());
2063 self
2064 }
2065
2066 pub fn with_capture_logs_dir(mut self, logs_dir: impl Into<PathBuf>) -> Self {
2068 self.capture_logs_dir = Some(logs_dir.into());
2069 self
2070 }
2071
2072 pub fn with_daemon_incarnation(self, daemon_incarnation: String) -> Self {
2075 self.spawn_events.configure_incarnation(daemon_incarnation);
2079 self
2080 }
2081
2082 pub fn with_terminal_journal(self, path: PathBuf, daemon_incarnation: String) -> Self {
2085 let mut this = self.with_daemon_incarnation(daemon_incarnation.clone());
2086 this.terminal_journal = Some(Arc::new(crate::terminal_journal::TerminalJournal::open(
2087 path,
2088 daemon_incarnation,
2089 )));
2090 this
2091 }
2092
2093 pub fn with_forwarding(mut self, forwarding: Arc<ForwardingTable>) -> Self {
2094 self.forwarding = Some(forwarding);
2095 self
2096 }
2097
2098 pub fn with_handle(mut self, supervisor_handle: SupervisorHandle) -> Self {
2099 self.spawn_events = supervisor_handle.spawn_events.clone();
2100 self.supervisor_handle = Some(supervisor_handle);
2101 self
2102 }
2103
2104 pub fn with_health_config(mut self, health: HealthConfig) -> Self {
2105 self.health = health;
2106 self
2107 }
2108
2109 #[cfg(target_os = "linux")]
2110 pub fn with_cgroup_placement(
2111 mut self,
2112 cgroup_placement: Option<subc_cgroup::Placement>,
2113 ) -> Self {
2114 self.cgroup_placement = cgroup_placement;
2115 self
2116 }
2117
2118 pub fn spawn(&self, spec: ModuleSpec) -> Result<SupervisedModule, SuperviseError> {
2124 validate_spec(&spec)?;
2125
2126 let runtime = self.runtime_config();
2127 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2128 let child = spawn_child(
2129 &spec,
2130 runtime.connection_file_path.as_deref(),
2131 self.supervisor_handle.as_ref(),
2132 &runtime.stderr_ring,
2133 runtime.capture_logs_dir.as_deref(),
2134 &runtime.child_roster,
2135 #[cfg(target_os = "linux")]
2136 runtime.cgroup_placement.as_ref(),
2137 )?;
2138 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2139 self.process_liveness
2140 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2141
2142 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2143 }
2144
2145 pub fn supervise_configured(
2151 &self,
2152 spec: ModuleSpec,
2153 enabled: bool,
2154 ) -> Result<SupervisedModule, SuperviseError> {
2155 validate_spec(&spec)?;
2156
2157 let runtime = self.runtime_config();
2158 if !enabled {
2159 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2160 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2161 }
2162
2163 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2164 match spawn_child(
2165 &spec,
2166 runtime.connection_file_path.as_deref(),
2167 self.supervisor_handle.as_ref(),
2168 &runtime.stderr_ring,
2169 runtime.capture_logs_dir.as_deref(),
2170 &runtime.child_roster,
2171 #[cfg(target_os = "linux")]
2172 runtime.cgroup_placement.as_ref(),
2173 ) {
2174 Ok(child) => {
2175 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2176 self.process_liveness
2177 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2178 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2179 }
2180 Err(err) => {
2181 error!(
2182 module_id = %spec.module_id,
2183 program = %spec.program.display(),
2184 error = %err,
2185 "configured module failed to spawn; marking failed and continuing"
2186 );
2187 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2188 Ok(self.supervised_module(spec, runtime, snapshot, None))
2189 }
2190 }
2191 }
2192
2193 pub fn supervise_configured_with_health(
2199 &self,
2200 spec: ModuleSpec,
2201 enabled: bool,
2202 health: HealthConfig,
2203 drain_timeout_ms: Option<u64>,
2204 restart_policy: RestartPolicy,
2205 ) -> Result<SupervisedModule, SuperviseError> {
2206 validate_spec(&spec)?;
2207
2208 let mut runtime = self.runtime_config();
2209 runtime.health = health;
2210 runtime.restart_policy = restart_policy;
2211 if let Some(ms) = drain_timeout_ms {
2212 runtime.drain_timeout = Duration::from_millis(ms);
2213 *runtime
2214 .effective_drain_timeout
2215 .lock()
2216 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
2217 }
2218 if !enabled {
2219 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2220 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2221 }
2222
2223 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2224 match spawn_child(
2225 &spec,
2226 runtime.connection_file_path.as_deref(),
2227 self.supervisor_handle.as_ref(),
2228 &runtime.stderr_ring,
2229 runtime.capture_logs_dir.as_deref(),
2230 &runtime.child_roster,
2231 #[cfg(target_os = "linux")]
2232 runtime.cgroup_placement.as_ref(),
2233 ) {
2234 Ok(child) => {
2235 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2236 self.process_liveness
2237 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2238 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2239 }
2240 Err(err) => {
2241 if health.critical {
2242 error!(
2243 module_id = %spec.module_id,
2244 program = %spec.program.display(),
2245 error = %err,
2246 "critical configured module failed to spawn; marking failed and alerting"
2247 );
2248 } else {
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 }
2256 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2257 Ok(self.supervised_module(spec, runtime, snapshot, None))
2258 }
2259 }
2260 }
2261
2262 fn runtime_config(&self) -> SupervisorRuntimeConfig {
2263 let effective_drain_timeout = Arc::new(Mutex::new(self.drain_timeout));
2264 SupervisorRuntimeConfig {
2265 restart_policy: self.restart_policy,
2266 drain_timeout: self.drain_timeout,
2267 child_roster: self
2270 .child_roster
2271 .for_module(Arc::clone(&effective_drain_timeout)),
2272 effective_drain_timeout,
2273 default_drain_timeout: self.drain_timeout,
2274 health: self.health,
2275 connection_file_path: self.connection_file_path.clone(),
2276 capture_logs_dir: self.capture_logs_dir.clone(),
2277 forwarding: self.forwarding.clone(),
2278 supervisor_handle: self.supervisor_handle.clone(),
2279 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
2280 terminal_ring: Arc::new(Mutex::new(
2281 TerminalRing::new(
2282 TerminalRingConfig::default(),
2283 self.daemon_start_clock.started_at_ms(),
2284 )
2285 .with_start_clock(self.daemon_start_clock)
2286 .with_journal(self.terminal_journal.clone()),
2287 )),
2288 spawn_events: self.spawn_events.clone(),
2289 #[cfg(target_os = "linux")]
2290 cgroup_placement: self.cgroup_placement.clone(),
2291 #[cfg(test)]
2292 test_seed_stale_facts_before_enable_spawn: false,
2293 }
2294 }
2295
2296 fn supervised_module(
2297 &self,
2298 spec: ModuleSpec,
2299 runtime: SupervisorRuntimeConfig,
2300 snapshot: SharedSnapshot,
2301 child: Option<SupervisedChild>,
2302 ) -> SupervisedModule {
2303 let configuration = Arc::new(Mutex::new(SupervisedConfiguration {
2304 spec: spec.clone(),
2305 health: runtime.health,
2306 }));
2307 let stderr_ring = Arc::clone(&runtime.stderr_ring);
2308 let terminal_ring = Arc::clone(&runtime.terminal_ring);
2309 let restart_policy = runtime.restart_policy;
2313 let effective_drain_timeout = Arc::clone(&runtime.effective_drain_timeout);
2314 let (tx, rx) = mpsc::channel(4);
2315 let monitor = tokio::spawn(supervise_loop(
2316 spec.clone(),
2317 runtime,
2318 Arc::clone(&self.registry),
2319 Arc::clone(&self.process_liveness),
2320 Arc::clone(&snapshot),
2321 child,
2322 rx,
2323 ));
2324
2325 let module_id = spec.module_id.clone();
2326 let module = SupervisedModule {
2327 inner: Arc::new(SupervisedModuleInner {
2328 module_id: module_id.clone(),
2329 registry: Arc::clone(&self.registry),
2330 snapshot,
2331 configuration,
2332 stderr_ring,
2333 terminal_ring,
2334 commands: tx,
2335 monitor: Mutex::new(Some(monitor)),
2336 restart_policy,
2337 effective_drain_timeout,
2338 provenance_probe: self.provenance_probe.clone(),
2339 }),
2340 };
2341 if let Some(supervisor_handle) = &self.supervisor_handle {
2342 supervisor_handle.apply_identity_configuration(&spec);
2343 supervisor_handle.insert(module.clone());
2344 }
2345 module
2346 }
2347}
2348
2349impl Default for Supervisor {
2350 fn default() -> Self {
2351 Self::new(Arc::new(Registry::default()), RestartPolicy::default())
2352 }
2353}
2354
2355#[derive(Clone)]
2357pub struct SupervisedModule {
2358 inner: Arc<SupervisedModuleInner>,
2359}
2360
2361struct SupervisedModuleInner {
2362 module_id: String,
2363 registry: Arc<Registry>,
2364 snapshot: SharedSnapshot,
2365 configuration: Arc<Mutex<SupervisedConfiguration>>,
2366 stderr_ring: Arc<Mutex<StderrRing>>,
2367 terminal_ring: Arc<Mutex<TerminalRing>>,
2368 commands: mpsc::Sender<SupervisorCommand>,
2369 monitor: Mutex<Option<JoinHandle<()>>>,
2370 restart_policy: RestartPolicy,
2374 effective_drain_timeout: Arc<Mutex<Duration>>,
2375 provenance_probe: ExecutableIdentityProbe,
2376}
2377
2378impl fmt::Debug for SupervisedModule {
2379 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2380 f.debug_struct("SupervisedModule")
2381 .field("module_id", &self.inner.module_id)
2382 .field("status", &self.status())
2383 .finish_non_exhaustive()
2384 }
2385}
2386
2387impl SupervisedModule {
2388 pub fn module_id(&self) -> &str {
2389 &self.inner.module_id
2390 }
2391
2392 #[cfg(test)]
2396 pub(crate) fn record_health_probe_failure_for_test(
2397 &self,
2398 detail: &str,
2399 ) -> Result<(), SuperviseError> {
2400 update_snapshot(&self.inner.snapshot, Some(&self.inner.module_id), |state| {
2401 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
2402 state.health.detail = Some(detail.to_string());
2403 })
2404 }
2405
2406 pub fn state(&self) -> Result<ModuleState, SuperviseError> {
2407 Ok(lock_snapshot(&self.inner.snapshot)?.state)
2408 }
2409
2410 pub fn stderr_tail(
2417 &self,
2418 max_lines: Option<usize>,
2419 max_bytes: Option<usize>,
2420 ) -> StderrTailSnapshot {
2421 self.inner
2422 .stderr_ring
2423 .lock()
2424 .unwrap_or_else(|poisoned| poisoned.into_inner())
2425 .snapshot(max_lines, max_bytes)
2426 }
2427
2428 pub fn terminal_history(&self) -> TerminalHistorySnapshot {
2433 self.inner
2434 .terminal_ring
2435 .lock()
2436 .unwrap_or_else(|poisoned| poisoned.into_inner())
2437 .snapshot()
2438 }
2439
2440 pub fn durable_terminal_history(&self) -> subc_control::TerminalHistory {
2445 durable_terminal_history_of(&self.inner.terminal_ring, &self.inner.module_id)
2446 }
2447
2448 pub(crate) async fn read_durable_terminal_history(
2453 &self,
2454 ) -> Result<subc_control::TerminalHistory, tokio::task::JoinError> {
2455 let terminal_ring = Arc::clone(&self.inner.terminal_ring);
2456 let module_id = self.inner.module_id.clone();
2457 tokio::task::spawn_blocking(move || durable_terminal_history_of(&terminal_ring, &module_id))
2458 .await
2459 }
2460
2461 pub fn status(&self) -> Result<ModuleStatus, SuperviseError> {
2462 self.status_with_snapshot_lock(&self.inner.snapshot, None)
2463 }
2464
2465 pub(crate) fn record_deliberate_severance(
2466 &self,
2467 identity: ProcessIdentity,
2468 ) -> Result<bool, SuperviseError> {
2469 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2470 if snapshot.pid != Some(identity.pid)
2471 || snapshot.process_start_time != Some(identity.start_time)
2472 {
2473 return Ok(false);
2474 }
2475 snapshot.deliberate_severance = Some(identity);
2476 Ok(true)
2477 }
2478
2479 pub(crate) fn status_for_control(
2484 &self,
2485 caller: &'static str,
2486 ) -> Result<ModuleStatus, SuperviseError> {
2487 self.status_with_snapshot_lock(&self.inner.snapshot, Some(caller))
2488 }
2489
2490 fn status_with_snapshot_lock(
2491 &self,
2492 snapshot: &SharedSnapshot,
2493 caller: Option<&'static str>,
2494 ) -> Result<ModuleStatus, SuperviseError> {
2495 let mut guard = match caller {
2496 Some(caller) => lock_snapshot_for_control(snapshot, &self.inner.module_id, caller)?,
2497 None => lock_snapshot(snapshot)?,
2498 };
2499 let restart_count =
2502 guard.crash_restarts_in_window(self.inner.restart_policy.window, Instant::now());
2503 let snapshot = guard.clone();
2504 drop(guard);
2505 let drain_timeout = *self.inner.effective_drain_timeout.lock().map_err(|_| {
2506 SuperviseError::StatePoisoned {
2507 module_id: Some(self.inner.module_id.clone()),
2508 }
2509 })?;
2510 let registration_active = self
2511 .inner
2512 .registry
2513 .get_module(&self.inner.module_id)
2514 .map_err(SuperviseError::Registry)?
2515 .is_some();
2516 let protocol = self.declared_protocol()?;
2517 let running_process =
2518 snapshot.enabled && snapshot.state == ModuleState::Running && snapshot.process_alive;
2519 let live = match protocol {
2525 ModuleProtocol::Subc => running_process && registration_active,
2526 ModuleProtocol::None => running_process,
2527 };
2528
2529 Ok(ModuleStatus {
2530 module_id: self.inner.module_id.clone(),
2531 state: snapshot.state,
2532 enabled: snapshot.enabled,
2533 process_alive: snapshot.process_alive,
2534 registration_active,
2535 protocol,
2536 live,
2537 restart_count,
2538 lifetime_restarts: snapshot.lifetime_restarts,
2539 spawn_generation: snapshot.spawn_generation,
2540 max_restarts: self.inner.restart_policy.max_restarts,
2541 restart_window: self.inner.restart_policy.window,
2542 drain_timeout,
2543 restart_backoff: self.inner.restart_policy.backoff,
2544 restart_max_backoff: self.inner.restart_policy.max_backoff,
2545 pid: snapshot.pid,
2546 spawned_at_ms: snapshot.spawned_at_ms,
2547 spawned_from: snapshot.spawned_from,
2548 process_start_time: snapshot.process_start_time,
2549 last_exit: snapshot.last_exit,
2550 health: snapshot.health,
2551 })
2552 }
2553
2554 #[cfg(test)]
2555 pub(crate) fn hold_snapshot_for_test(
2556 &self,
2557 acquired: std::sync::mpsc::Sender<()>,
2558 hold: Duration,
2559 ) -> std::thread::JoinHandle<()> {
2560 let snapshot = Arc::clone(&self.inner.snapshot);
2561 std::thread::spawn(move || {
2562 let _guard = snapshot.lock().expect("test snapshot lock is not poisoned");
2563 acquired
2564 .send(())
2565 .expect("test receiver waits for snapshot lock");
2566 std::thread::sleep(hold);
2567 })
2568 }
2569
2570 pub(crate) async fn running_image_agreement(&self) -> subc_control::RunningImageAgreement {
2571 let snapshot = match lock_snapshot(&self.inner.snapshot) {
2572 Ok(snapshot) => snapshot.clone(),
2573 Err(_) => {
2574 return subc_control::RunningImageAgreement::Unavailable {
2575 reason: subc_control::RunningImageUnavailableReason::NotRunning,
2576 };
2577 }
2578 };
2579 self.inner
2580 .provenance_probe
2581 .observe(
2582 snapshot.pid,
2583 snapshot.spawned_from.as_deref(),
2584 snapshot.spawned_file_identity,
2585 snapshot.process_start_time,
2586 )
2587 .await
2588 }
2589
2590 pub(crate) fn will_recover_after_connection_loss(&self) -> Result<bool, SuperviseError> {
2591 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2592 Ok(match snapshot.state {
2593 ModuleState::Restarting => true,
2594 ModuleState::Failed | ModuleState::Disabled => false,
2595 _ => daemon_will_restart(&mut snapshot, &self.inner.restart_policy, Instant::now()),
2596 })
2597 }
2598
2599 #[cfg(test)]
2600 pub(crate) fn is_warming(&self) -> Result<bool, SuperviseError> {
2601 self.is_warming_with_snapshot_lock(None)
2602 }
2603
2604 pub(crate) fn is_warming_for_control(
2605 &self,
2606 caller: &'static str,
2607 ) -> Result<bool, SuperviseError> {
2608 self.is_warming_with_snapshot_lock(Some(caller))
2609 }
2610
2611 fn is_warming_with_snapshot_lock(
2612 &self,
2613 caller: Option<&'static str>,
2614 ) -> Result<bool, SuperviseError> {
2615 let snapshot = match caller {
2616 Some(caller) => {
2617 lock_snapshot_for_control(&self.inner.snapshot, &self.inner.module_id, caller)?
2618 }
2619 None => lock_snapshot(&self.inner.snapshot)?,
2620 }
2621 .clone();
2622 Ok(matches!(
2623 snapshot.state,
2624 ModuleState::Starting | ModuleState::Running | ModuleState::Restarting
2625 ))
2626 }
2627
2628 pub async fn drain(&self) -> Result<(), SuperviseError> {
2630 self.stop().await
2631 }
2632
2633 pub(crate) async fn retire(&self) -> Result<(), SuperviseError> {
2634 match self.state()? {
2635 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2636 ModuleState::Starting
2637 | ModuleState::Running
2638 | ModuleState::Unresponsive
2639 | ModuleState::Restarting
2640 | ModuleState::Draining
2641 | ModuleState::Disabled => {}
2642 }
2643
2644 let (reply_tx, reply_rx) = oneshot::channel();
2645 self.inner
2646 .commands
2647 .send(SupervisorCommand::Retire { reply: reply_tx })
2648 .await
2649 .map_err(|_| SuperviseError::CommandClosed {
2650 module_id: self.inner.module_id.clone(),
2651 })?;
2652 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2653 module_id: self.inner.module_id.clone(),
2654 })?
2655 }
2656
2657 pub async fn stop(&self) -> Result<(), SuperviseError> {
2658 match self.state()? {
2659 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2660 ModuleState::Starting
2661 | ModuleState::Running
2662 | ModuleState::Unresponsive
2663 | ModuleState::Restarting
2664 | ModuleState::Draining
2665 | ModuleState::Disabled => {}
2666 }
2667
2668 let (reply_tx, reply_rx) = oneshot::channel();
2669 self.inner
2670 .commands
2671 .send(SupervisorCommand::Drain { reply: reply_tx })
2672 .await
2673 .map_err(|_| SuperviseError::CommandClosed {
2674 module_id: self.inner.module_id.clone(),
2675 })?;
2676 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2677 module_id: self.inner.module_id.clone(),
2678 })?
2679 }
2680
2681 pub async fn restart(&self, drain_timeout_ms: Option<u64>) -> Result<(), SuperviseError> {
2682 let (reply_tx, reply_rx) = oneshot::channel();
2683 self.inner
2684 .commands
2685 .send(SupervisorCommand::Restart {
2686 drain_timeout_ms,
2687 reply: reply_tx,
2688 })
2689 .await
2690 .map_err(|_| SuperviseError::CommandClosed {
2691 module_id: self.inner.module_id.clone(),
2692 })?;
2693 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2694 module_id: self.inner.module_id.clone(),
2695 })?
2696 }
2697
2698 pub async fn swap(&self, ready_timeout: Option<Duration>) -> Result<(), SuperviseError> {
2703 let (reply_tx, reply_rx) = oneshot::channel();
2704 self.inner
2705 .commands
2706 .send(SupervisorCommand::Swap {
2707 ready_timeout,
2708 reply: reply_tx,
2709 })
2710 .await
2711 .map_err(|_| SuperviseError::CommandClosed {
2712 module_id: self.inner.module_id.clone(),
2713 })?;
2714 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2715 module_id: self.inner.module_id.clone(),
2716 })?
2717 }
2718
2719 pub async fn reload(&self) -> Result<(), SuperviseError> {
2720 let (reply_tx, reply_rx) = oneshot::channel();
2721 self.inner
2722 .commands
2723 .send(SupervisorCommand::Reload { reply: reply_tx })
2724 .await
2725 .map_err(|_| SuperviseError::CommandClosed {
2726 module_id: self.inner.module_id.clone(),
2727 })?;
2728 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2729 module_id: self.inner.module_id.clone(),
2730 })?
2731 }
2732
2733 pub async fn set_enabled(&self, enabled: bool) -> Result<bool, SuperviseError> {
2734 let (reply_tx, reply_rx) = oneshot::channel();
2735 self.inner
2736 .commands
2737 .send(SupervisorCommand::SetEnabled {
2738 enabled,
2739 reply: reply_tx,
2740 })
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(crate) fn declared_protocol(&self) -> Result<ModuleProtocol, SuperviseError> {
2755 Ok(self
2756 .inner
2757 .configuration
2758 .lock()
2759 .map_err(|_| SuperviseError::StatePoisoned {
2760 module_id: Some(self.inner.module_id.clone()),
2761 })?
2762 .spec
2763 .protocol)
2764 }
2765
2766 pub(crate) fn configuration(&self) -> Result<(ModuleSpec, HealthConfig), SuperviseError> {
2767 let configuration =
2768 self.inner
2769 .configuration
2770 .lock()
2771 .map_err(|_| SuperviseError::StatePoisoned {
2772 module_id: Some(self.inner.module_id.clone()),
2773 })?;
2774 Ok((configuration.spec.clone(), configuration.health))
2775 }
2776
2777 #[cfg(any(test, feature = "test-support"))]
2781 pub async fn update_spec_for_test(&self, spec: ModuleSpec) -> Result<(), SuperviseError> {
2782 let (_, health) = self.configuration()?;
2783 let drain_timeout_ms = u64::try_from(
2784 self.inner
2785 .effective_drain_timeout
2786 .lock()
2787 .unwrap_or_else(|poisoned| poisoned.into_inner())
2788 .as_millis(),
2789 )
2790 .ok();
2791 self.update_configuration(spec, health, drain_timeout_ms)
2792 .await
2793 }
2794
2795 pub(crate) async fn update_configuration(
2796 &self,
2797 spec: ModuleSpec,
2798 health: HealthConfig,
2799 drain_timeout_ms: Option<u64>,
2800 ) -> Result<(), SuperviseError> {
2801 if spec.module_id != self.inner.module_id {
2802 return Err(SuperviseError::InvalidSpec {
2803 reason: "a supervised module's module_id cannot be changed".to_string(),
2804 });
2805 }
2806 validate_spec(&spec)?;
2807 let (reply_tx, reply_rx) = oneshot::channel();
2808 self.inner
2809 .commands
2810 .send(SupervisorCommand::UpdateConfiguration {
2811 spec: spec.clone(),
2812 health,
2813 drain_timeout_ms,
2814 reply: reply_tx,
2815 })
2816 .await
2817 .map_err(|_| SuperviseError::CommandClosed {
2818 module_id: self.inner.module_id.clone(),
2819 })?;
2820 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2821 module_id: self.inner.module_id.clone(),
2822 })?;
2823 let mut configuration =
2824 self.inner
2825 .configuration
2826 .lock()
2827 .map_err(|_| SuperviseError::StatePoisoned {
2828 module_id: Some(self.inner.module_id.clone()),
2829 })?;
2830 configuration.spec = spec;
2831 configuration.health = health;
2832 Ok(())
2833 }
2834}
2835
2836impl Drop for SupervisedModuleInner {
2837 fn drop(&mut self) {
2838 let Ok(mut monitor) = self.monitor.lock() else {
2839 return;
2840 };
2841 if let Some(monitor) = monitor.as_ref().filter(|monitor| !monitor.is_finished()) {
2842 let _ = update_snapshot(&self.snapshot, Some(&self.module_id), |state| {
2843 state.state = ModuleState::Stopped;
2844 clear_current_process_facts(state);
2845 });
2846 monitor.abort();
2847 }
2848 let _ = monitor.take();
2849 }
2850}
2851
2852#[derive(Debug)]
2853enum SupervisorCommand {
2854 Drain {
2855 reply: oneshot::Sender<Result<(), SuperviseError>>,
2856 },
2857 Retire {
2858 reply: oneshot::Sender<Result<(), SuperviseError>>,
2859 },
2860 Restart {
2861 drain_timeout_ms: Option<u64>,
2866 reply: oneshot::Sender<Result<(), SuperviseError>>,
2867 },
2868 Reload {
2869 reply: oneshot::Sender<Result<(), SuperviseError>>,
2870 },
2871 SetEnabled {
2872 enabled: bool,
2873 reply: oneshot::Sender<Result<bool, SuperviseError>>,
2874 },
2875 UpdateConfiguration {
2876 spec: ModuleSpec,
2877 health: HealthConfig,
2878 drain_timeout_ms: Option<u64>,
2881 reply: oneshot::Sender<()>,
2882 },
2883 Swap {
2884 ready_timeout: Option<Duration>,
2887 reply: oneshot::Sender<Result<(), SuperviseError>>,
2889 },
2890}
2891
2892#[derive(Debug)]
2893pub enum SuperviseError {
2894 InvalidSpec {
2895 reason: String,
2896 },
2897 Spawn {
2898 program: PathBuf,
2899 source: io::Error,
2900 cgroup_path: Option<PathBuf>,
2901 },
2902 Cgroup {
2903 module_id: String,
2904 source: io::Error,
2905 },
2906 LaunchNonce {
2909 reason: String,
2910 },
2911 Wait {
2912 module_id: String,
2913 source: io::Error,
2914 },
2915 Kill {
2916 module_id: String,
2917 source: io::Error,
2918 },
2919 Forwarding(ForwardingError),
2920 Registry(RegistryError),
2921 ReloadUnavailable {
2922 module_id: String,
2923 reason: String,
2924 },
2925 Disabled {
2930 module_id: String,
2931 },
2932 ReloadFailed {
2933 module_id: String,
2934 reason: String,
2935 },
2936 RegistrationStillActive {
2937 module_id: String,
2938 waited: Duration,
2939 },
2940 StatePoisoned {
2941 module_id: Option<String>,
2942 },
2943 CommandClosed {
2944 module_id: String,
2945 },
2946 SwapInProgress {
2950 module_id: String,
2951 },
2952 SwapRefused {
2954 module_id: String,
2955 reason: SwapRefusal,
2956 },
2957 SwapFailed {
2961 module_id: String,
2962 arm: SwapFailureArm,
2963 detail: String,
2964 candidate_exit: Option<ExitReport>,
2967 },
2968}
2969
2970#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2972pub enum SwapRefusal {
2973 OverlapExclusive,
2975 NotRegistered,
2978 ProtocolNone,
2981 NotConfigured,
2984 AlreadySwapping,
2986}
2987
2988impl SwapRefusal {
2989 pub fn as_str(self) -> &'static str {
2990 match self {
2991 Self::OverlapExclusive => "overlap_exclusive",
2992 Self::NotRegistered => "not_registered",
2993 Self::ProtocolNone => "protocol_none",
2994 Self::NotConfigured => "not_configured",
2995 Self::AlreadySwapping => "already_swapping",
2996 }
2997 }
2998}
2999
3000#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3003pub enum SwapFailureArm {
3004 SpawnFailed,
3006 NeverRegistered,
3008 NeverReady,
3010 CandidateExited,
3012 CandidateUnhealthy,
3014 Interrupted,
3018 CutoverLost,
3023}
3024
3025impl SwapFailureArm {
3026 pub fn as_str(self) -> &'static str {
3027 match self {
3028 Self::SpawnFailed => "spawn_failed",
3029 Self::NeverRegistered => "never_registered",
3030 Self::NeverReady => "never_ready",
3031 Self::CandidateExited => "candidate_exited",
3032 Self::CandidateUnhealthy => "candidate_unhealthy",
3033 Self::Interrupted => "interrupted",
3034 Self::CutoverLost => "cutover_lost",
3035 }
3036 }
3037}
3038
3039impl fmt::Display for SuperviseError {
3040 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3041 match self {
3042 Self::InvalidSpec { reason } => write!(f, "invalid module spec: {reason}"),
3043 Self::Spawn {
3044 program,
3045 source,
3046 cgroup_path: Some(cgroup_path),
3047 } => write!(
3048 f,
3049 "failed to place module in cgroup '{}' while spawning '{}': {source}",
3050 cgroup_path.display(),
3051 program.display()
3052 ),
3053 Self::Spawn {
3054 program,
3055 source,
3056 cgroup_path: None,
3057 } => write!(
3058 f,
3059 "failed to spawn module '{}': {source}",
3060 program.display()
3061 ),
3062 Self::Cgroup { module_id, source } => {
3063 write!(
3064 f,
3065 "failed to prepare cgroup for module '{module_id}': {source}"
3066 )
3067 }
3068 Self::LaunchNonce { reason } => {
3069 write!(
3070 f,
3071 "failed to generate reserved-module launch nonce: {reason}"
3072 )
3073 }
3074 Self::Wait { module_id, source } => {
3075 write!(f, "failed to wait for module '{module_id}': {source}")
3076 }
3077 Self::Kill { module_id, source } => {
3078 write!(f, "failed to kill module '{module_id}': {source}")
3079 }
3080 Self::Forwarding(err) => write!(f, "forwarding error: {err}"),
3081 Self::Registry(err) => write!(f, "registry error: {err}"),
3082 Self::ReloadUnavailable { module_id, reason } => {
3083 write!(f, "reload unavailable for module '{module_id}': {reason}")
3084 }
3085 Self::Disabled { module_id } => {
3086 write!(
3087 f,
3088 "module '{module_id}' is disabled; enable it before restart or reload"
3089 )
3090 }
3091 Self::ReloadFailed { module_id, reason } => {
3092 write!(f, "reload failed for module '{module_id}': {reason}")
3093 }
3094 Self::RegistrationStillActive { module_id, waited } => write!(
3095 f,
3096 "module '{module_id}' registration remained active after waiting {waited:?}"
3097 ),
3098 Self::StatePoisoned { module_id } => match module_id {
3099 Some(module_id) => {
3100 write!(f, "supervisor state for module '{module_id}' was poisoned")
3101 }
3102 None => write!(f, "supervisor state was poisoned"),
3103 },
3104 Self::CommandClosed { module_id } => {
3105 write!(
3106 f,
3107 "supervisor command channel for module '{module_id}' is closed"
3108 )
3109 }
3110 Self::SwapInProgress { module_id } => write!(
3111 f,
3112 "module '{module_id}' is being swapped; retry once the swap has cut over or failed, or stop the module to abort the swap"
3113 ),
3114 Self::SwapRefused { module_id, reason } => match reason {
3115 SwapRefusal::OverlapExclusive => write!(
3116 f,
3117 "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"
3118 ),
3119 SwapRefusal::NotRegistered => write!(
3120 f,
3121 "module '{module_id}' is not registered, so there is no serving process to keep while a replacement warms; use a plain restart"
3122 ),
3123 SwapRefusal::ProtocolNone => write!(
3124 f,
3125 "module '{module_id}' is protocol: \"none\" and never registers, so a swap could never see its replacement become ready; use a plain restart"
3126 ),
3127 SwapRefusal::NotConfigured => write!(
3128 f,
3129 "module '{module_id}' cannot be swapped: the supervisor was built without the forwarding table or shared handle a swap needs"
3130 ),
3131 SwapRefusal::AlreadySwapping => {
3132 write!(f, "module '{module_id}' is already being swapped")
3133 }
3134 },
3135 Self::SwapFailed {
3136 module_id,
3137 arm,
3138 detail,
3139 ..
3140 } => write!(
3141 f,
3142 "swap of module '{module_id}' failed ({}): {detail}; the running process was left serving",
3143 arm.as_str()
3144 ),
3145 }
3146 }
3147}
3148
3149impl Error for SuperviseError {
3150 fn source(&self) -> Option<&(dyn Error + 'static)> {
3151 match self {
3152 Self::Spawn { source, .. }
3153 | Self::Cgroup { source, .. }
3154 | Self::Wait { source, .. }
3155 | Self::Kill { source, .. } => Some(source),
3156 Self::Forwarding(err) => Some(err),
3157 Self::Registry(err) => Some(err),
3158 Self::LaunchNonce { .. }
3159 | Self::InvalidSpec { .. }
3160 | Self::ReloadUnavailable { .. }
3161 | Self::Disabled { .. }
3162 | Self::ReloadFailed { .. }
3163 | Self::RegistrationStillActive { .. }
3164 | Self::StatePoisoned { .. }
3165 | Self::CommandClosed { .. }
3166 | Self::SwapInProgress { .. }
3167 | Self::SwapRefused { .. }
3168 | Self::SwapFailed { .. } => None,
3169 }
3170 }
3171}
3172
3173pub(crate) fn validate_spec(spec: &ModuleSpec) -> Result<(), SuperviseError> {
3174 if spec.module_id.trim().is_empty() {
3175 return Err(SuperviseError::InvalidSpec {
3176 reason: "module_id must not be empty".to_string(),
3177 });
3178 }
3179
3180 Ok(())
3181}
3182
3183#[derive(Debug, Default)]
3184struct HealthProbeRuntime {
3185 registered_connection: Option<crate::ConnectionId>,
3186 advertised: bool,
3187 next_probe_at: Option<Instant>,
3188 probe_index: u64,
3189}
3190
3191impl HealthProbeRuntime {
3192 fn refresh_registration(
3193 &mut self,
3194 spec: &ModuleSpec,
3195 runtime: &SupervisorRuntimeConfig,
3196 registry: &Registry,
3197 snapshot: &SharedSnapshot,
3198 ) {
3199 if spec.protocol == ModuleProtocol::None {
3211 self.registered_connection = None;
3212 self.advertised = false;
3213 self.next_probe_at = None;
3214 return;
3215 }
3216
3217 let registration = match registry.get_module(&spec.module_id) {
3218 Ok(registration) => registration,
3219 Err(err) => {
3220 warn!(module_id = %spec.module_id, error = %err, "health prober could not read registry");
3221 self.advertised = false;
3222 self.next_probe_at = None;
3223 return;
3224 }
3225 };
3226
3227 let Some(registration) = registration else {
3228 self.registered_connection = None;
3229 self.advertised = false;
3230 self.next_probe_at = None;
3231 return;
3232 };
3233
3234 let advertised = registration
3235 .control_ops
3236 .iter()
3237 .any(|op| op == MODULE_CONTROL_OP_HEALTH_CHECK);
3238 if !advertised {
3239 self.registered_connection = Some(registration.connection_id);
3240 self.advertised = false;
3241 self.next_probe_at = None;
3242 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3243 state.health.status = SupervisorHealthStatus::Unknown;
3244 state.health.consecutive_failures = 0;
3245 state.health.last_probe_ms = None;
3246 state.health.detail = None;
3247 state.health.metrics = None;
3248 });
3249 return;
3250 }
3251
3252 let reregistered = self.registered_connection != Some(registration.connection_id);
3253 self.registered_connection = Some(registration.connection_id);
3254 self.advertised = true;
3255 if reregistered || self.next_probe_at.is_none() {
3256 self.probe_index = 0;
3257 self.next_probe_at = Some(
3258 Instant::now() + jittered_health_delay(&spec.module_id, 0, runtime.health.cadence),
3259 );
3260 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3261 state.health.status = SupervisorHealthStatus::Unknown;
3262 state.health.consecutive_failures = 0;
3263 state.health.detail = None;
3264 state.health.metrics = None;
3265 });
3266 }
3267 }
3268
3269 fn wake_after(&self) -> Duration {
3270 if !self.advertised {
3271 return REGISTRY_RELEASE_POLL;
3272 }
3273 self.next_probe_at
3274 .map(|next| next.saturating_duration_since(Instant::now()))
3275 .unwrap_or(REGISTRY_RELEASE_POLL)
3276 }
3277
3278 fn due(&self) -> bool {
3279 self.advertised
3280 && self
3281 .next_probe_at
3282 .is_some_and(|next| Instant::now() >= next)
3283 }
3284
3285 fn schedule_next(&mut self, spec: &ModuleSpec, cadence: Duration) {
3286 self.probe_index = self.probe_index.wrapping_add(1);
3287 self.next_probe_at = Some(
3288 Instant::now() + jittered_health_delay(&spec.module_id, self.probe_index, cadence),
3289 );
3290 }
3291}
3292
3293#[derive(Debug)]
3328enum HealthProbeEvidence {
3329 LaneDead,
3331 NoAnswer,
3333 BadAnswer,
3335 Misconfigured,
3337}
3338
3339#[derive(Debug)]
3340struct HealthProbeError {
3341 evidence: HealthProbeEvidence,
3342 message: String,
3343}
3344
3345impl HealthProbeError {
3346 fn lane_dead(message: impl Into<String>) -> Self {
3347 Self::with(HealthProbeEvidence::LaneDead, message)
3348 }
3349
3350 fn no_answer(message: impl Into<String>) -> Self {
3351 Self::with(HealthProbeEvidence::NoAnswer, message)
3352 }
3353
3354 fn bad_answer(message: impl Into<String>) -> Self {
3355 Self::with(HealthProbeEvidence::BadAnswer, message)
3356 }
3357
3358 fn misconfigured(message: impl Into<String>) -> Self {
3359 Self::with(HealthProbeEvidence::Misconfigured, message)
3360 }
3361
3362 fn with(evidence: HealthProbeEvidence, message: impl Into<String>) -> Self {
3363 Self {
3364 evidence,
3365 message: message.into(),
3366 }
3367 }
3368
3369 #[allow(dead_code)]
3383 fn is_proof_of_death(&self) -> bool {
3384 matches!(self.evidence, HealthProbeEvidence::LaneDead)
3385 }
3386
3387 fn label(&self) -> &'static str {
3395 match self.evidence {
3396 HealthProbeEvidence::LaneDead => "lane-dead",
3397 HealthProbeEvidence::NoAnswer => "no-answer",
3398 HealthProbeEvidence::BadAnswer => "bad-answer",
3399 HealthProbeEvidence::Misconfigured => "daemon-misconfigured",
3400 }
3401 }
3402}
3403
3404impl fmt::Display for HealthProbeError {
3405 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3406 f.write_str(&self.message)
3407 }
3408}
3409
3410async fn run_health_probe_cycle(
3411 spec: &ModuleSpec,
3412 runtime: &SupervisorRuntimeConfig,
3413 registry: &Registry,
3414 process_liveness: &SupervisorProcessLiveness,
3415 snapshot: &SharedSnapshot,
3416 child: &mut Option<SupervisedChild>,
3417) {
3418 let now_ms = unix_ms_now();
3419 match probe_module_health(&spec.module_id, runtime, None).await {
3420 Ok(report) => {
3421 handle_health_report(
3422 spec,
3423 runtime,
3424 registry,
3425 process_liveness,
3426 snapshot,
3427 child,
3428 report,
3429 now_ms,
3430 )
3431 .await;
3432 }
3433 Err(err) => {
3434 handle_health_probe_failure(
3435 spec,
3436 runtime,
3437 registry,
3438 process_liveness,
3439 snapshot,
3440 child,
3441 err,
3442 now_ms,
3443 )
3444 .await;
3445 }
3446 }
3447}
3448
3449async fn probe_module_health(
3450 module_id: &str,
3451 runtime: &SupervisorRuntimeConfig,
3452 drain_deadline: Option<Instant>,
3453) -> Result<HealthReport, HealthProbeError> {
3454 let Some(forwarding) = runtime.forwarding.as_ref() else {
3455 return Err(HealthProbeError::misconfigured(
3456 "supervisor was not configured with a forwarding table",
3457 ));
3458 };
3459 let probe_started_at = Instant::now();
3460 let mut deadline = probe_started_at + runtime.health.deadline;
3461 if let Some(drain_deadline) = drain_deadline {
3462 deadline = deadline.min(drain_deadline);
3463 }
3464 let pending = if drain_deadline.is_some() {
3465 forwarding.begin_drain_health_probe_rpc_for(
3466 module_id,
3467 MODULE_CONTROL_OP_HEALTH_CHECK,
3468 probe_started_at,
3469 deadline,
3470 )
3471 } else {
3472 forwarding.begin_health_probe_rpc_for(
3473 module_id,
3474 MODULE_CONTROL_OP_HEALTH_CHECK,
3475 probe_started_at,
3476 deadline,
3477 )
3478 }
3479 .map_err(|err| {
3480 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3483 })?;
3484 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3485}
3486
3487async fn probe_endpoint_health(
3494 endpoint: crate::ModuleEndpointId,
3495 runtime: &SupervisorRuntimeConfig,
3496 deadline_cap: Option<Instant>,
3497) -> Result<HealthReport, HealthProbeError> {
3498 let Some(forwarding) = runtime.forwarding.as_ref() else {
3499 return Err(HealthProbeError::misconfigured(
3500 "supervisor was not configured with a forwarding table",
3501 ));
3502 };
3503 let probe_started_at = Instant::now();
3504 let mut deadline = probe_started_at + runtime.health.deadline;
3505 if let Some(cap) = deadline_cap {
3506 deadline = deadline.min(cap);
3507 }
3508 let pending = forwarding
3509 .begin_endpoint_health_probe_rpc_for(
3510 endpoint,
3511 MODULE_CONTROL_OP_HEALTH_CHECK,
3512 probe_started_at,
3513 deadline,
3514 )
3515 .map_err(|err| {
3516 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3517 })?;
3518 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3519}
3520
3521async fn await_health_probe(
3523 forwarding: &ForwardingTable,
3524 pending: PendingModuleControlRpc,
3525 deadline: Instant,
3526 probe_budget: Duration,
3527) -> Result<HealthReport, HealthProbeError> {
3528 let PendingModuleControlRpc {
3529 endpoint,
3530 module_sink,
3531 negotiated_ver,
3532 corr,
3533 receiver,
3534 } = pending;
3535 let body = serde_json::to_vec(&ModuleControlRequest::HealthCheck {}).map_err(|err| {
3536 HealthProbeError::misconfigured(format!("failed to encode health.check: {err}"))
3537 })?;
3538 let frame = Frame::build_with_version(
3539 negotiated_ver,
3540 FrameType::Request,
3541 control_flags(),
3542 0,
3543 0,
3544 corr,
3545 body,
3546 )
3547 .map_err(|err| {
3548 HealthProbeError::misconfigured(format!("failed to build health.check frame: {err}"))
3549 })?;
3550
3551 match timeout_at(deadline, module_sink.send(frame)).await {
3557 Ok(Ok(())) => {}
3558 Ok(Err(err)) => {
3559 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3560 return Err(HealthProbeError::lane_dead(format!(
3563 "failed to send health.check: {err}"
3564 )));
3565 }
3566 Err(_elapsed) => {
3567 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3568 return Err(HealthProbeError::no_answer(
3572 "health.check send timed out before enqueue (module egress full)",
3573 ));
3574 }
3575 }
3576
3577 match timeout_at(deadline, receiver).await {
3578 Ok(Ok(ModuleControlRpcOutcome::Response(response))) => {
3582 response.health_report().ok_or_else(|| {
3583 HealthProbeError::bad_answer("health.check RPC returned a non-health response")
3584 })
3585 }
3586 Ok(Ok(ModuleControlRpcOutcome::Rejected(body))) => Err(HealthProbeError::bad_answer(
3587 format!("health.check rejected: {}", body.message),
3588 )),
3589 Ok(Ok(ModuleControlRpcOutcome::ModuleGone(message))) => {
3590 Err(HealthProbeError::lane_dead(message))
3591 }
3592 Ok(Ok(ModuleControlRpcOutcome::MalformedResponse(message))) => {
3593 Err(HealthProbeError::bad_answer(message))
3594 }
3595 Ok(Ok(ModuleControlRpcOutcome::UnexpectedOp { expected, actual })) => {
3596 Err(HealthProbeError::bad_answer(format!(
3597 "expected module-control op '{expected}', got '{actual}'"
3598 )))
3599 }
3600 Ok(Ok(ModuleControlRpcOutcome::DeadlineElapsed)) => Err(HealthProbeError::bad_answer(
3604 "module answered health.check after its daemon deadline",
3605 )),
3606 Ok(Err(_)) => Err(HealthProbeError::misconfigured(
3607 "health.check waiter was canceled before the module responded",
3608 )),
3609 Err(_) => {
3610 let _ = forwarding.tombstone_health_probe_rpc(endpoint, corr);
3611 Err(HealthProbeError::no_answer(format!(
3612 "module did not answer health.check within {probe_budget:?}"
3613 )))
3614 }
3615 }
3616}
3617
3618#[allow(clippy::too_many_arguments)]
3619async fn handle_health_report(
3620 spec: &ModuleSpec,
3621 runtime: &SupervisorRuntimeConfig,
3622 registry: &Registry,
3623 process_liveness: &SupervisorProcessLiveness,
3624 snapshot: &SharedSnapshot,
3625 child: &mut Option<SupervisedChild>,
3626 report: HealthReport,
3627 now_ms: u64,
3628) {
3629 let status = supervisor_health_status(report.status);
3630 let detail = report.detail.clone();
3631 let metrics = truncate_health_metrics(report.metrics);
3632 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3633 state.health.status = status;
3634 state.health.last_probe_ms = Some(now_ms);
3635 state.health.detail = detail.clone();
3636 state.health.metrics = metrics.clone();
3637 state.health.consecutive_failures = 0;
3638 });
3639
3640 let action = match report.status {
3641 HealthStatus::Ok => return,
3642 HealthStatus::Degraded => runtime.health.on_degraded,
3643 HealthStatus::Failing => runtime.health.on_failing,
3644 };
3645 apply_l3_health_action(
3646 spec,
3647 runtime,
3648 registry,
3649 process_liveness,
3650 snapshot,
3651 child,
3652 status,
3653 detail.as_deref(),
3654 action,
3655 now_ms,
3656 )
3657 .await;
3658}
3659
3660#[allow(clippy::too_many_arguments)]
3661async fn handle_health_probe_failure(
3662 spec: &ModuleSpec,
3663 runtime: &SupervisorRuntimeConfig,
3664 registry: &Registry,
3665 process_liveness: &SupervisorProcessLiveness,
3666 snapshot: &SharedSnapshot,
3667 child: &mut Option<SupervisedChild>,
3668 err: HealthProbeError,
3669 now_ms: u64,
3670) {
3671 let threshold = runtime.health.failure_threshold.max(1);
3672 let mut failures = 0;
3673 let detail = format!("[{}] {err}", err.label());
3678 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3679 state.health.last_probe_ms = Some(now_ms);
3680 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
3681 state.health.detail = Some(detail.clone());
3682 state.health.metrics = None;
3683 failures = state.health.consecutive_failures;
3684 });
3685
3686 if failures < threshold {
3687 warn!(
3688 module_id = %spec.module_id,
3689 consecutive_failures = failures,
3690 threshold,
3691 evidence = err.label(),
3692 detail = %detail,
3693 "health.check probe failed"
3694 );
3695 return;
3696 }
3697
3698 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3699 state.state = ModuleState::Unresponsive;
3700 state.health.status = SupervisorHealthStatus::Unresponsive;
3701 });
3702 if runtime.health.critical {
3706 error!(
3707 module_id = %spec.module_id,
3708 status = "unresponsive",
3709 evidence = err.label(),
3710 detail = %detail,
3711 "critical module health alert"
3712 );
3713 } else {
3714 warn!(
3715 module_id = %spec.module_id,
3716 status = "unresponsive",
3717 evidence = err.label(),
3718 detail = %detail,
3719 "module health threshold breached"
3720 );
3721 }
3722 if let Err(err) = health_restart_child(
3723 spec,
3724 runtime,
3725 registry,
3726 process_liveness,
3727 snapshot,
3728 child,
3729 SupervisorHealthStatus::Unresponsive,
3730 Some(&detail),
3731 now_ms,
3732 )
3733 .await
3734 {
3735 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3736 }
3737}
3738
3739#[allow(clippy::too_many_arguments)]
3740async fn apply_l3_health_action(
3741 spec: &ModuleSpec,
3742 runtime: &SupervisorRuntimeConfig,
3743 registry: &Registry,
3744 process_liveness: &SupervisorProcessLiveness,
3745 snapshot: &SharedSnapshot,
3746 child: &mut Option<SupervisedChild>,
3747 status: SupervisorHealthStatus,
3748 detail: Option<&str>,
3749 action: HealthAction,
3750 now_ms: u64,
3751) {
3752 record_health_action(snapshot, &spec.module_id, action.to_string(), now_ms);
3753 match action {
3754 HealthAction::Report => {
3755 info!(
3756 module_id = %spec.module_id,
3757 status = ?status,
3758 detail,
3759 "module reported non-ok health"
3760 );
3761 }
3762 HealthAction::Alert => {
3763 error!(
3764 module_id = %spec.module_id,
3765 status = ?status,
3766 detail,
3767 "module health alert"
3768 );
3769 }
3770 HealthAction::Restart => {
3771 if let Err(err) = health_restart_child(
3772 spec,
3773 runtime,
3774 registry,
3775 process_liveness,
3776 snapshot,
3777 child,
3778 status,
3779 detail,
3780 now_ms,
3781 )
3782 .await
3783 {
3784 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3785 }
3786 }
3787 }
3788}
3789
3790#[allow(clippy::too_many_arguments)]
3791async fn health_restart_child(
3792 spec: &ModuleSpec,
3793 runtime: &SupervisorRuntimeConfig,
3794 registry: &Registry,
3795 process_liveness: &SupervisorProcessLiveness,
3796 snapshot: &SharedSnapshot,
3797 child: &mut Option<SupervisedChild>,
3798 status: SupervisorHealthStatus,
3799 detail: Option<&str>,
3800 now_ms: u64,
3801) -> Result<(), SuperviseError> {
3802 let (enabled, schedule) = {
3803 let mut state = lock_snapshot(snapshot)?;
3804 let enabled = state.enabled;
3805 let schedule = if enabled {
3806 state.next_crash_restart(&runtime.restart_policy, Instant::now())
3807 } else {
3808 None
3809 };
3810 (enabled, schedule)
3811 };
3812
3813 if !enabled {
3814 return Err(SuperviseError::Disabled {
3815 module_id: spec.module_id.clone(),
3816 });
3817 }
3818
3819 if schedule.is_none() {
3820 record_health_action(snapshot, &spec.module_id, "disabled".to_string(), now_ms);
3821 error!(
3822 module_id = %spec.module_id,
3823 status = ?status,
3824 detail,
3825 max_restarts = runtime.restart_policy.max_restarts,
3826 window_secs = runtime.restart_policy.window.as_secs(),
3827 "health restart budget exhausted; disabling module"
3828 );
3829 begin_forwarding_drain_if_configured(
3830 spec,
3831 runtime,
3832 registry,
3833 snapshot,
3834 Some(false),
3835 RouteCloseReason::Disable,
3836 )
3837 .await?;
3838 drain_optional_child(
3839 &spec.module_id,
3840 spec.protocol,
3841 registry,
3842 snapshot,
3843 &runtime.terminal_ring,
3844 &runtime.spawn_events,
3845 child,
3846 runtime.drain_timeout,
3847 ModuleState::Disabled,
3848 Some(false),
3849 )
3850 .await?;
3851 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3852 return Ok(());
3853 }
3854
3855 let schedule = schedule.expect("a health restart must have a crash-restart schedule");
3856 let mut restart_count = 0;
3857 update_snapshot(snapshot, Some(&spec.module_id), |state| {
3858 restart_count = state.crash_restarts.len();
3859 state.state = ModuleState::Unresponsive;
3860 state.health.status = status;
3861 state.health.last_action = Some(HealthAction::Restart.to_string());
3862 state.health.last_action_ms = Some(now_ms);
3863 })?;
3864 warn!(
3865 module_id = %spec.module_id,
3866 status = ?status,
3867 detail,
3868 restart_count,
3869 restart_in_window = schedule.restart_in_window,
3870 delay_ms = schedule.delay.as_millis() as u64,
3871 "health-triggered module restart"
3872 );
3873
3874 begin_forwarding_drain_if_configured(
3875 spec,
3876 runtime,
3877 registry,
3878 snapshot,
3879 Some(true),
3880 RouteCloseReason::Restart,
3881 )
3882 .await?;
3883 drain_optional_child(
3884 &spec.module_id,
3885 spec.protocol,
3886 registry,
3887 snapshot,
3888 &runtime.terminal_ring,
3889 &runtime.spawn_events,
3890 child,
3891 runtime.drain_timeout,
3892 ModuleState::Restarting,
3893 Some(true),
3894 )
3895 .await?;
3896 sleep(schedule.delay).await;
3897 if !respawn_still_pending(snapshot) {
3901 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3902 return Ok(());
3903 }
3904 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
3905 match spawn_and_mark_running(spec, runtime, snapshot) {
3906 Ok(next_child) => {
3907 *child = Some(next_child);
3908 Ok(())
3909 }
3910 Err(err) => {
3911 fail_snapshot(snapshot, Some(&spec.module_id), None);
3912 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3913 *child = None;
3914 Err(err)
3915 }
3916 }
3917}
3918
3919fn record_health_action(snapshot: &SharedSnapshot, module_id: &str, action: String, now_ms: u64) {
3920 let _ = update_snapshot(snapshot, Some(module_id), |state| {
3921 state.health.last_action = Some(action);
3922 state.health.last_action_ms = Some(now_ms);
3923 });
3924}
3925
3926fn supervisor_health_status(status: HealthStatus) -> SupervisorHealthStatus {
3927 match status {
3928 HealthStatus::Ok => SupervisorHealthStatus::Ok,
3929 HealthStatus::Degraded => SupervisorHealthStatus::Degraded,
3930 HealthStatus::Failing => SupervisorHealthStatus::Failing,
3931 }
3932}
3933
3934fn truncate_health_metrics(metrics: Option<Value>) -> Option<Value> {
3946 let metrics = metrics?;
3947 match serde_json::to_vec(&metrics) {
3948 Ok(encoded) if encoded.len() > MAX_HEALTH_METRICS_BYTES => Some(serde_json::json!({
3949 "truncated": true,
3950 "original_bytes": encoded.len(),
3951 })),
3952 Ok(_) | Err(_) => Some(metrics),
3953 }
3954}
3955
3956fn jittered_health_delay(module_id: &str, probe_index: u64, cadence: Duration) -> Duration {
3962 if cadence.is_zero() {
3963 return Duration::ZERO;
3964 }
3965 let cadence_ms = cadence.as_millis() as u64;
3966 if cadence_ms == 0 {
3982 return cadence;
3983 }
3984 let jitter_span = (cadence_ms / 10).max(1);
3999 let hash = module_id.as_bytes().iter().fold(
4000 probe_index.wrapping_mul(0x9E37_79B9_7F4A_7C15),
4001 |acc, byte| {
4002 acc.wrapping_mul(1099511628211)
4003 .wrapping_add(u64::from(*byte))
4004 },
4005 );
4006 cadence + Duration::from_millis(hash % jitter_span)
4007}
4008
4009#[cfg(test)]
4010mod tests {
4011 use super::*;
4012
4013 #[test]
4014 fn readding_a_module_clears_its_rescan_removal_tombstone() {
4015 let handle = SupervisorHandle::new();
4016 let module_id = "readded-tombstone";
4017 handle.record_rescan_removal(module_id);
4018 assert!(handle.removal_tombstone_age_ms(module_id).is_some());
4019
4020 handle.apply_identity_configuration(&ModuleSpec {
4021 module_id: module_id.to_string(),
4022 program: PathBuf::from("/test/module"),
4023 args: Vec::new(),
4024 env: Vec::new(),
4025 reserved: false,
4026 reserved_prefixes: Vec::new(),
4027 protocol: ModuleProtocol::Subc,
4028 overlap: Default::default(),
4029 });
4030
4031 assert!(
4032 handle.removal_tombstone_age_ms(module_id).is_none(),
4033 "a re-added module must not retain a stale removal tombstone"
4034 );
4035 }
4036
4037 fn stale_process_snapshot(state: ModuleState, enabled: bool) -> SharedSnapshot {
4038 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::new(state, enabled)));
4039 update_snapshot(&snapshot, Some("stale-process-facts"), |snapshot| {
4040 snapshot.process_alive = true;
4041 snapshot.pid = Some(41);
4042 snapshot.spawned_at_ms = Some(42);
4043 snapshot.spawned_from = Some(PathBuf::from("/spawned/module"));
4044 snapshot.spawned_file_identity = Some(SpawnedFileIdentity {
4045 device: 43,
4046 inode: 44,
4047 });
4048 })
4049 .unwrap();
4050 snapshot
4051 }
4052
4053 fn assert_snapshot_process_facts_cleared(snapshot: &SharedSnapshot) {
4054 let snapshot = lock_snapshot(snapshot).unwrap();
4055 assert!(!snapshot.process_alive);
4056 assert_eq!(snapshot.pid, None);
4057 assert_eq!(snapshot.spawned_at_ms, None);
4058 assert_eq!(snapshot.spawned_from, None);
4059 assert_eq!(snapshot.spawned_file_identity, None);
4060 }
4061
4062 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4063 async fn failed_enable_spawn_clears_preexisting_current_process_facts() {
4064 let supervisor = Supervisor::default();
4065 let mut runtime = supervisor.runtime_config();
4066 runtime.test_seed_stale_facts_before_enable_spawn = true;
4067 let snapshot = stale_process_snapshot(ModuleState::Disabled, false);
4068 let mut child = None;
4069 let spec = ModuleSpec {
4070 module_id: "failed-enable-clears-facts".to_string(),
4071 program: PathBuf::from("/definitely/missing/failed-enable-module"),
4072 args: Vec::new(),
4073 env: Vec::new(),
4074 reserved: false,
4075 reserved_prefixes: Vec::new(),
4076 protocol: ModuleProtocol::Subc,
4077 overlap: Default::default(),
4078 };
4079
4080 let result = set_child_enabled(
4081 &spec,
4082 &runtime,
4083 &supervisor.registry,
4084 &supervisor.process_liveness,
4085 &snapshot,
4086 &mut child,
4087 true,
4088 )
4089 .await;
4090
4091 assert!(matches!(result, Err(SuperviseError::Spawn { .. })));
4092 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4093 assert_snapshot_process_facts_cleared(&snapshot);
4094 }
4095
4096 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4097 async fn failed_reload_spawn_clears_current_process_facts() {
4098 let supervisor = Supervisor::default();
4099 let mut runtime = supervisor.runtime_config();
4100 runtime.restart_policy = RestartPolicy::new(0, Duration::ZERO);
4101 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4102 let mut child = None;
4103 let spec = ModuleSpec {
4104 module_id: "failed-reload-clears-facts".to_string(),
4105 program: PathBuf::from("/unused/failed-reload-module"),
4106 args: Vec::new(),
4107 env: Vec::new(),
4108 reserved: false,
4109 reserved_prefixes: Vec::new(),
4110 protocol: ModuleProtocol::Subc,
4111 overlap: Default::default(),
4112 };
4113
4114 let result = handle_reload_spawn_failure(
4115 &spec,
4116 &runtime,
4117 &supervisor.process_liveness,
4118 &snapshot,
4119 &mut child,
4120 "forced reload spawn failure".to_string(),
4121 )
4122 .await;
4123
4124 assert!(matches!(result, Err(SuperviseError::ReloadFailed { .. })));
4125 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4126 assert_snapshot_process_facts_cleared(&snapshot);
4127 }
4128
4129 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4130 async fn dropping_a_module_with_an_active_monitor_clears_current_process_facts() {
4131 let supervisor = Supervisor::default();
4132 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4133 let module = supervisor.supervised_module(
4134 ModuleSpec {
4135 module_id: "drop-clears-facts".to_string(),
4136 program: PathBuf::from("/unused/drop-module"),
4137 args: Vec::new(),
4138 env: Vec::new(),
4139 reserved: false,
4140 reserved_prefixes: Vec::new(),
4141 protocol: ModuleProtocol::Subc,
4142 overlap: Default::default(),
4143 },
4144 supervisor.runtime_config(),
4145 Arc::clone(&snapshot),
4146 None,
4147 );
4148 assert!(!module
4149 .inner
4150 .monitor
4151 .lock()
4152 .unwrap()
4153 .as_ref()
4154 .unwrap()
4155 .is_finished());
4156
4157 drop(module);
4158
4159 assert_eq!(
4160 lock_snapshot(&snapshot).unwrap().state,
4161 ModuleState::Stopped
4162 );
4163 assert_snapshot_process_facts_cleared(&snapshot);
4164 }
4165
4166 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4167 async fn configuration_update_does_not_replace_captured_running_process_facts() {
4168 let supervisor = Supervisor::default();
4169 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4170 let initial = ModuleSpec {
4171 module_id: "rescan-preserves-spawn-facts".to_string(),
4172 program: PathBuf::from("/spawned/module"),
4173 args: Vec::new(),
4174 env: Vec::new(),
4175 reserved: false,
4176 reserved_prefixes: Vec::new(),
4177 protocol: ModuleProtocol::Subc,
4178 overlap: Default::default(),
4179 };
4180 let module = supervisor.supervised_module(
4181 initial.clone(),
4182 supervisor.runtime_config(),
4183 snapshot,
4184 None,
4185 );
4186 let before = module.status().unwrap();
4187 let mut replacement = initial;
4188 replacement.program = PathBuf::from("/rescanned/replacement-module");
4189
4190 module
4191 .update_configuration(replacement, HealthConfig::default(), None)
4192 .await
4193 .unwrap();
4194
4195 let after = module.status().unwrap();
4196 assert_eq!(after.pid, before.pid);
4197 assert_eq!(after.spawned_at_ms, before.spawned_at_ms);
4198 assert_eq!(after.spawned_from, before.spawned_from);
4199 drop(module);
4200 }
4201}
4202
4203fn unix_ms_now() -> u64 {
4204 SystemTime::now()
4205 .duration_since(UNIX_EPOCH)
4206 .map(|duration| duration.as_millis().min(u128::from(u64::MAX)) as u64)
4207 .unwrap_or(0)
4208}
4209
4210async fn supervise_loop(
4211 mut spec: ModuleSpec,
4212 mut runtime: SupervisorRuntimeConfig,
4213 registry: Arc<Registry>,
4214 process_liveness: Arc<SupervisorProcessLiveness>,
4215 snapshot: SharedSnapshot,
4216 mut child: Option<SupervisedChild>,
4217 mut commands: mpsc::Receiver<SupervisorCommand>,
4218) {
4219 let mut health_probe = HealthProbeRuntime::default();
4220 let mut pending_respawn: Option<Instant> = None;
4224 let mut requeued: VecDeque<SupervisorCommand> = VecDeque::new();
4227 loop {
4228 if let Some(command) = requeued.pop_front() {
4229 if !handle_supervisor_command(
4230 command,
4231 &mut spec,
4232 &mut runtime,
4233 ®istry,
4234 &process_liveness,
4235 &snapshot,
4236 &mut child,
4237 &mut commands,
4238 &mut requeued,
4239 )
4240 .await
4241 {
4242 return;
4243 }
4244 if child.is_some() || !respawn_still_pending(&snapshot) {
4245 pending_respawn = None;
4246 }
4247 continue;
4248 }
4249 if child.is_some() {
4250 health_probe.refresh_registration(&spec, &runtime, ®istry, &snapshot);
4251 let probe_sleep = sleep(health_probe.wake_after());
4252 tokio::pin!(probe_sleep);
4253 let active_child = child.as_mut().expect("child checked above");
4254 tokio::select! {
4255 wait_result = active_child.wait() => {
4256 let exit_report = match wait_result {
4265 Ok(status) => classify_reaped_child_exit(&snapshot, active_child, &status),
4266 Err(err) => {
4267 active_child.drain_stderr(&spec.module_id).await;
4268 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4269 record_wait_error_terminal(
4275 &spec.module_id,
4276 &runtime.terminal_ring,
4277 &runtime.spawn_events,
4278 );
4279 untrack_if_registration_released(
4280 &process_liveness,
4281 ®istry,
4282 &spec.module_id,
4283 &snapshot,
4284 );
4285 error!(module_id = %spec.module_id, error = %err, "failed to wait for supervised module");
4286 child = None;
4287 continue;
4288 }
4289 };
4290 active_child.drain_stderr(&spec.module_id).await;
4291
4292 match on_child_exit(
4293 &spec,
4294 runtime.restart_policy,
4295 ®istry,
4296 &snapshot,
4297 &runtime.terminal_ring,
4298 &runtime.spawn_events,
4299 exit_report,
4300 ).await {
4301 NextAction::Stop { registration_released } => {
4302 if registration_released {
4303 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4304 }
4305 child = None;
4306 }
4307 NextAction::Restart { schedule } => {
4308 let delay = schedule.map_or(
4309 runtime.restart_policy.delay_for_restart(0),
4310 |schedule| schedule.delay,
4311 );
4312 if let Some(schedule) = schedule {
4313 log_crash_respawn(&spec.module_id, schedule);
4314 }
4315 child = None;
4323 pending_respawn = Some(Instant::now() + delay);
4324 }
4325 }
4326 }
4327 command = commands.recv() => {
4328 let Some(command) = command else {
4329 return;
4330 };
4331 if !handle_supervisor_command(
4332 command,
4333 &mut spec,
4334 &mut runtime,
4335 ®istry,
4336 &process_liveness,
4337 &snapshot,
4338 &mut child,
4339 &mut commands,
4340 &mut requeued,
4341 ).await {
4342 return;
4343 }
4344 }
4345 _ = &mut probe_sleep => {
4346 if health_probe.due() {
4347 run_health_probe_cycle(
4348 &spec,
4349 &runtime,
4350 ®istry,
4351 &process_liveness,
4352 &snapshot,
4353 &mut child,
4354 ).await;
4355 if child.is_some() {
4356 health_probe.schedule_next(&spec, runtime.health.cadence);
4357 }
4358 }
4359 }
4360 }
4361 } else if let Some(deadline) = pending_respawn {
4362 tokio::select! {
4363 _ = sleep_until(deadline) => {
4364 pending_respawn = None;
4365 if !respawn_still_pending(&snapshot) {
4369 continue;
4370 }
4371 if let Err(err) = wait_for_registration_release(
4372 ®istry,
4373 &spec.module_id,
4374 REGISTRY_RELEASE_TIMEOUT,
4375 ).await {
4376 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4377 error!(module_id = %spec.module_id, error = %err, "registration did not release before restart");
4378 continue;
4379 }
4380
4381 match spawn_and_mark_running(&spec, &runtime, &snapshot) {
4382 Ok(next_child) => {
4383 child = Some(next_child);
4384 debug!(module_id = %spec.module_id, "supervised module restarted after crash");
4385 }
4386 Err(err) => {
4387 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4388 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4389 error!(module_id = %spec.module_id, error = %err, "failed to restart supervised module");
4390 }
4391 }
4392 }
4393 command = commands.recv() => {
4394 let Some(command) = command else {
4395 return;
4396 };
4397 if !handle_supervisor_command(
4398 command,
4399 &mut spec,
4400 &mut runtime,
4401 ®istry,
4402 &process_liveness,
4403 &snapshot,
4404 &mut child,
4405 &mut commands,
4406 &mut requeued,
4407 ).await {
4408 return;
4409 }
4410 if child.is_some() || !respawn_still_pending(&snapshot) {
4415 pending_respawn = None;
4416 }
4417 }
4418 }
4419 } else {
4420 let Some(command) = commands.recv().await else {
4421 return;
4422 };
4423 if !handle_supervisor_command(
4424 command,
4425 &mut spec,
4426 &mut runtime,
4427 ®istry,
4428 &process_liveness,
4429 &snapshot,
4430 &mut child,
4431 &mut commands,
4432 &mut requeued,
4433 )
4434 .await
4435 {
4436 return;
4437 }
4438 }
4439 }
4440}
4441
4442fn log_crash_respawn(module_id: &str, schedule: CrashRestartSchedule) {
4443 info!(
4444 module_id,
4445 restart_in_window = schedule.restart_in_window,
4446 delay_ms = schedule.delay.as_millis() as u64,
4447 "respawning after crash"
4448 );
4449}
4450
4451fn respawn_still_pending(snapshot: &SharedSnapshot) -> bool {
4457 matches!(
4458 lock_snapshot(snapshot),
4459 Ok(state) if state.enabled && state.state == ModuleState::Restarting
4460 )
4461}
4462
4463enum NextAction {
4464 Stop {
4465 registration_released: bool,
4466 },
4467 Restart {
4468 schedule: Option<CrashRestartSchedule>,
4469 },
4470}
4471
4472#[allow(clippy::too_many_arguments)]
4473async fn handle_supervisor_command(
4474 command: SupervisorCommand,
4475 spec: &mut ModuleSpec,
4476 runtime: &mut SupervisorRuntimeConfig,
4477 registry: &Registry,
4478 process_liveness: &SupervisorProcessLiveness,
4479 snapshot: &SharedSnapshot,
4480 child: &mut Option<SupervisedChild>,
4481 commands: &mut mpsc::Receiver<SupervisorCommand>,
4482 requeued: &mut VecDeque<SupervisorCommand>,
4483) -> bool {
4484 match command {
4485 SupervisorCommand::Drain { reply } => {
4486 let result = drain_optional_child(
4487 &spec.module_id,
4488 spec.protocol,
4489 registry,
4490 snapshot,
4491 &runtime.terminal_ring,
4492 &runtime.spawn_events,
4493 child,
4494 runtime.drain_timeout,
4495 ModuleState::Stopped,
4496 None,
4497 )
4498 .await;
4499 let registration_released = result.is_ok();
4500 let _ = reply.send(result);
4501 if registration_released {
4502 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4503 }
4504 false
4505 }
4506 SupervisorCommand::Retire { reply } => {
4507 let result = async {
4508 begin_forwarding_drain_if_configured(
4509 spec,
4510 runtime,
4511 registry,
4512 snapshot,
4513 None,
4514 RouteCloseReason::Disable,
4515 )
4516 .await?;
4517 drain_optional_child(
4518 &spec.module_id,
4519 spec.protocol,
4520 registry,
4521 snapshot,
4522 &runtime.terminal_ring,
4523 &runtime.spawn_events,
4524 child,
4525 runtime.drain_timeout,
4526 ModuleState::Stopped,
4527 None,
4528 )
4529 .await
4530 }
4531 .await;
4532 let registration_released = result.is_ok();
4533 let _ = reply.send(result);
4534 if registration_released {
4535 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4536 }
4537 false
4538 }
4539 SupervisorCommand::Restart {
4540 drain_timeout_ms,
4541 reply,
4542 } => {
4543 let validation = match lock_snapshot(snapshot) {
4555 Ok(state) if !state.enabled => Err(SuperviseError::Disabled {
4556 module_id: spec.module_id.clone(),
4557 }),
4558 Ok(_) => Ok(()),
4559 Err(err) => Err(err),
4560 };
4561 let initiated = validation.is_ok();
4562 let _ = reply.send(validation);
4563 if initiated {
4564 let drain_timeout = drain_timeout_ms
4567 .map(Duration::from_millis)
4568 .unwrap_or(runtime.drain_timeout);
4569 if let Err(err) = restart_child(
4570 spec,
4571 runtime,
4572 registry,
4573 process_liveness,
4574 snapshot,
4575 child,
4576 drain_timeout,
4577 )
4578 .await
4579 {
4580 warn!(
4581 module_id = %spec.module_id,
4582 error = %err,
4583 "operator restart failed after initiation ack; module state carries the outcome"
4584 );
4585 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4586 state.state = ModuleState::Failed;
4587 clear_current_process_facts(state);
4588 });
4589 }
4590 }
4591 true
4592 }
4593 SupervisorCommand::Reload { reply } => {
4594 let result =
4595 reload_child(spec, runtime, registry, process_liveness, snapshot, child).await;
4596 let _ = reply.send(result);
4597 true
4598 }
4599 SupervisorCommand::SetEnabled { enabled, reply } => {
4600 let result = set_child_enabled(
4601 spec,
4602 runtime,
4603 registry,
4604 process_liveness,
4605 snapshot,
4606 child,
4607 enabled,
4608 )
4609 .await;
4610 let _ = reply.send(result);
4611 true
4612 }
4613 SupervisorCommand::UpdateConfiguration {
4614 spec: next_spec,
4615 health,
4616 drain_timeout_ms,
4617 reply,
4618 } => {
4619 if let Some(handle) = &runtime.supervisor_handle {
4620 handle.apply_identity_configuration(&next_spec);
4621 }
4622 *spec = next_spec;
4623 runtime.health = health;
4624 runtime.drain_timeout = drain_timeout_ms
4625 .map(Duration::from_millis)
4626 .unwrap_or(runtime.default_drain_timeout);
4627 *runtime
4628 .effective_drain_timeout
4629 .lock()
4630 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
4631 let _ = reply.send(());
4632 true
4633 }
4634 SupervisorCommand::Swap {
4635 ready_timeout,
4636 reply,
4637 } => {
4638 let end = swap::run_swap(
4639 spec,
4640 runtime,
4641 registry,
4642 process_liveness,
4643 snapshot,
4644 child,
4645 commands,
4646 ready_timeout.unwrap_or(DEFAULT_SWAP_READY_TIMEOUT),
4647 reply,
4648 )
4649 .await;
4650 requeued.extend(end.requeue);
4651 true
4652 }
4653 }
4654}
4655
4656async fn restart_child(
4657 spec: &ModuleSpec,
4658 runtime: &SupervisorRuntimeConfig,
4659 registry: &Registry,
4660 process_liveness: &SupervisorProcessLiveness,
4661 snapshot: &SharedSnapshot,
4662 child: &mut Option<SupervisedChild>,
4663 drain_timeout: Duration,
4664) -> Result<(), SuperviseError> {
4665 if !lock_snapshot(snapshot)?.enabled {
4667 return Err(SuperviseError::Disabled {
4668 module_id: spec.module_id.clone(),
4669 });
4670 }
4671 begin_forwarding_drain_with_timeout(
4672 spec,
4673 runtime,
4674 registry,
4675 snapshot,
4676 None,
4677 RouteCloseReason::Restart,
4678 drain_timeout,
4679 )
4680 .await?;
4681
4682 if child.is_some() {
4683 drain_optional_child(
4684 &spec.module_id,
4685 spec.protocol,
4686 registry,
4687 snapshot,
4688 &runtime.terminal_ring,
4689 &runtime.spawn_events,
4690 child,
4691 drain_timeout,
4692 ModuleState::Restarting,
4693 Some(true),
4694 )
4695 .await?;
4696 } else {
4697 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4698 state.enabled = true;
4699 state.state = ModuleState::Restarting;
4700 clear_current_process_facts(state);
4701 })?;
4702 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4703 }
4704
4705 reset_restart_count(snapshot, &spec.module_id)?;
4706 sleep(runtime.restart_policy.backoff).await;
4707 if !respawn_still_pending(snapshot) {
4710 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4711 return Ok(());
4712 }
4713 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4714 match spawn_and_mark_running(spec, runtime, snapshot) {
4720 Ok(next_child) => {
4721 *child = Some(next_child);
4722 debug!(module_id = %spec.module_id, "supervised module restarted by operator request");
4723 Ok(())
4724 }
4725 Err(err) => {
4726 fail_snapshot(snapshot, Some(&spec.module_id), None);
4727 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4728 *child = None;
4729 Err(err)
4730 }
4731 }
4732}
4733
4734async fn reload_child(
4735 spec: &ModuleSpec,
4736 runtime: &SupervisorRuntimeConfig,
4737 registry: &Registry,
4738 process_liveness: &SupervisorProcessLiveness,
4739 snapshot: &SharedSnapshot,
4740 child: &mut Option<SupervisedChild>,
4741) -> Result<(), SuperviseError> {
4742 if !lock_snapshot(snapshot)?.enabled {
4744 return Err(SuperviseError::Disabled {
4745 module_id: spec.module_id.clone(),
4746 });
4747 }
4748 begin_forwarding_drain(
4749 spec,
4750 runtime,
4751 registry,
4752 snapshot,
4753 Some(true),
4754 RouteCloseReason::Reload,
4755 )
4756 .await?;
4757
4758 if child.is_some() {
4759 drain_optional_child(
4760 &spec.module_id,
4761 spec.protocol,
4762 registry,
4763 snapshot,
4764 &runtime.terminal_ring,
4765 &runtime.spawn_events,
4766 child,
4767 runtime.drain_timeout,
4768 ModuleState::Restarting,
4769 Some(true),
4770 )
4771 .await?;
4772 } else {
4773 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4774 state.enabled = true;
4775 state.state = ModuleState::Restarting;
4776 clear_current_process_facts(state);
4777 })?;
4778 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4779 }
4780
4781 reset_restart_count(snapshot, &spec.module_id)?;
4782 sleep(runtime.restart_policy.backoff).await;
4783 if !respawn_still_pending(snapshot) {
4786 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4787 return Ok(());
4788 }
4789 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4790 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
4791 Ok(next_child) => next_child,
4792 Err(err) => {
4793 return handle_reload_spawn_failure(
4794 spec,
4795 runtime,
4796 process_liveness,
4797 snapshot,
4798 child,
4799 format!("new child failed to spawn: {err}"),
4800 )
4801 .await;
4802 }
4803 };
4804 *child = Some(next_child);
4805
4806 let wait_outcome = {
4807 let active_child = child.as_mut().expect("new reload child was just stored");
4808 wait_for_registration_after_reload(
4809 registry,
4810 &spec.module_id,
4811 snapshot,
4812 active_child,
4813 REGISTRY_RELEASE_TIMEOUT,
4814 )
4815 .await?
4816 };
4817
4818 match wait_outcome {
4819 RegistrationWaitOutcome::Registered => {
4820 debug!(module_id = %spec.module_id, "supervised module reloaded and registered");
4821 Ok(())
4822 }
4823 RegistrationWaitOutcome::Exited(exit_report) => {
4824 if let Some(active_child) = child.as_mut() {
4825 active_child.drain_stderr(&spec.module_id).await;
4826 }
4827 *child = None;
4828 handle_reload_child_registration_failure(
4829 spec,
4830 runtime,
4831 registry,
4832 process_liveness,
4833 snapshot,
4834 child,
4835 ReloadRegistrationFailure {
4836 exit_report: registration_failure_exit_report(exit_report),
4837 reason: "new child exited before registering".to_string(),
4838 },
4839 )
4840 .await
4841 }
4842 RegistrationWaitOutcome::TimedOut => {
4843 let mut timed_out_child = child
4844 .take()
4845 .expect("timed-out reload child is still running");
4846 timed_out_child
4847 .start_kill()
4848 .map_err(|source| SuperviseError::Kill {
4849 module_id: spec.module_id.clone(),
4850 source,
4851 })?;
4852 let status = timed_out_child
4853 .wait()
4854 .await
4855 .map_err(|source| SuperviseError::Wait {
4856 module_id: spec.module_id.clone(),
4857 source,
4858 })?;
4859 timed_out_child.drain_stderr(&spec.module_id).await;
4860 handle_reload_child_registration_failure(
4861 spec,
4862 runtime,
4863 registry,
4864 process_liveness,
4865 snapshot,
4866 child,
4867 ReloadRegistrationFailure {
4868 exit_report: registration_failure_exit_report(classify_reaped_child_exit(
4869 snapshot,
4870 &timed_out_child,
4871 &status,
4872 )),
4873 reason: format!(
4874 "new child did not register within {:?}",
4875 REGISTRY_RELEASE_TIMEOUT
4876 ),
4877 },
4878 )
4879 .await
4880 }
4881 }
4882}
4883
4884async fn set_child_enabled(
4885 spec: &ModuleSpec,
4886 runtime: &SupervisorRuntimeConfig,
4887 registry: &Registry,
4888 process_liveness: &SupervisorProcessLiveness,
4889 snapshot: &SharedSnapshot,
4890 child: &mut Option<SupervisedChild>,
4891 enabled: bool,
4892) -> Result<bool, SuperviseError> {
4893 let (current_enabled, current_state) = {
4894 let state = lock_snapshot(snapshot)?;
4895 (state.enabled, state.state)
4896 };
4897 let revive_terminal = enabled
4905 && current_enabled
4906 && child.is_none()
4907 && matches!(current_state, ModuleState::Failed | ModuleState::Stopped);
4908 if current_enabled == enabled && !revive_terminal {
4909 return Ok(false);
4910 }
4911
4912 if enabled {
4913 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4914 state.enabled = true;
4915 state.state = ModuleState::Starting;
4916 clear_current_process_facts(state);
4917 })?;
4918 #[cfg(test)]
4919 if runtime.test_seed_stale_facts_before_enable_spawn {
4920 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4921 state.process_alive = true;
4922 state.pid = Some(41);
4923 state.spawned_at_ms = Some(42);
4924 state.spawned_from = Some(PathBuf::from("/spawned/module"));
4925 state.spawned_file_identity = Some(SpawnedFileIdentity {
4926 device: 43,
4927 inode: 44,
4928 });
4929 })?;
4930 }
4931 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4932 reset_restart_count(snapshot, &spec.module_id)?;
4933 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4934 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
4935 Ok(next_child) => next_child,
4936 Err(err) => {
4937 if let Err(state_err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4938 state.state = ModuleState::Failed;
4939 clear_current_process_facts(state);
4940 }) {
4941 error!(module_id = %spec.module_id, error = %state_err, "failed to record enable spawn failure");
4942 }
4943 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4944 return Err(err);
4945 }
4946 };
4947 *child = Some(next_child);
4948 debug!(module_id = %spec.module_id, "supervised module enabled");
4949 Ok(true)
4950 } else {
4951 begin_forwarding_drain_if_configured(
4952 spec,
4953 runtime,
4954 registry,
4955 snapshot,
4956 Some(false),
4957 RouteCloseReason::Disable,
4958 )
4959 .await?;
4960 drain_optional_child(
4961 &spec.module_id,
4962 spec.protocol,
4963 registry,
4964 snapshot,
4965 &runtime.terminal_ring,
4966 &runtime.spawn_events,
4967 child,
4968 runtime.drain_timeout,
4969 ModuleState::Disabled,
4970 Some(false),
4971 )
4972 .await?;
4973 debug!(module_id = %spec.module_id, "supervised module disabled");
4974 Ok(true)
4975 }
4976}
4977
4978async fn on_child_exit(
4979 spec: &ModuleSpec,
4980 policy: RestartPolicy,
4981 registry: &Registry,
4982 snapshot: &SharedSnapshot,
4983 terminal_ring: &Arc<Mutex<TerminalRing>>,
4984 spawn_events: &SpawnEventFeed,
4985 exit_report: ExitReport,
4986) -> NextAction {
4987 match exit_report.kind {
4988 ExitKind::Clean => {
4989 info!(
4990 module_id = %spec.module_id,
4991 exit_code = ?exit_report.code,
4992 exit_signal = ?exit_report.signal,
4993 "supervised module exited cleanly"
4994 );
4995 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4996 state.state = ModuleState::Stopped;
4997 clear_current_process_facts(state);
4998 state.last_exit = Some(exit_report.clone());
4999 }) {
5000 error!(module_id = %spec.module_id, error = %err, "failed to record clean module exit");
5001 }
5002 record_terminal(
5003 &spec.module_id,
5004 terminal_ring,
5005 spawn_events,
5006 &exit_report,
5007 TerminalDisposition::Stopped,
5008 );
5009 let registration_released = match wait_for_registration_release(
5010 registry,
5011 &spec.module_id,
5012 REGISTRY_RELEASE_TIMEOUT,
5013 )
5014 .await
5015 {
5016 Ok(()) => true,
5017 Err(err) => {
5018 warn!(module_id = %spec.module_id, error = %err, "registration still active after clean exit");
5019 false
5020 }
5021 };
5022 NextAction::Stop {
5023 registration_released,
5024 }
5025 }
5026 ExitKind::Crash => {
5027 warn!(
5028 module_id = %spec.module_id,
5029 exit_code = ?exit_report.code,
5030 exit_signal = ?exit_report.signal,
5031 "supervised module exited abnormally (crash)"
5032 );
5033 let mut restart_schedule = None;
5034 let mut disposition = TerminalDisposition::Disabled;
5035 let mut disposition_detail = None;
5039 let now = Instant::now();
5040 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5041 clear_current_process_facts(state);
5042 state.last_exit = Some(exit_report.clone());
5043 if state.enabled {
5044 if let Some(schedule) = state.next_crash_restart(&policy, now) {
5045 state.state = ModuleState::Restarting;
5046 restart_schedule = Some(schedule);
5047 disposition = TerminalDisposition::Restarting;
5048 } else {
5049 state.state = ModuleState::Failed;
5050 disposition = TerminalDisposition::Failed;
5051 disposition_detail = Some(policy.budget_exhausted_detail());
5052 }
5053 } else {
5054 state.state = ModuleState::Disabled;
5055 disposition = TerminalDisposition::Disabled;
5056 }
5057 }) {
5058 error!(module_id = %spec.module_id, error = %err, "failed to record crashed module exit");
5059 return NextAction::Stop {
5060 registration_released: false,
5061 };
5062 }
5063 if disposition_detail.is_some() {
5064 error!(
5069 module_id = %spec.module_id,
5070 max_restarts = policy.max_restarts,
5071 window_secs = policy.window.as_secs(),
5072 "module stopped: {}",
5073 policy.budget_exhausted_detail()
5074 );
5075 }
5076 record_terminal_with_detail(
5077 &spec.module_id,
5078 terminal_ring,
5079 spawn_events,
5080 &exit_report,
5081 disposition,
5082 disposition_detail,
5083 );
5084
5085 if let Some(schedule) = restart_schedule {
5086 NextAction::Restart {
5087 schedule: Some(schedule),
5088 }
5089 } else {
5090 let registration_released = match wait_for_registration_release(
5091 registry,
5092 &spec.module_id,
5093 REGISTRY_RELEASE_TIMEOUT,
5094 )
5095 .await
5096 {
5097 Ok(()) => true,
5098 Err(err) => {
5099 warn!(module_id = %spec.module_id, error = %err, "registration still active after failed module");
5100 false
5101 }
5102 };
5103 NextAction::Stop {
5104 registration_released,
5105 }
5106 }
5107 }
5108 ExitKind::DeliberateSeverance => {
5109 warn!(
5110 module_id = %spec.module_id,
5111 exit_code = ?exit_report.code,
5112 exit_signal = ?exit_report.signal,
5113 "supervised module exited after deliberate connection severance"
5114 );
5115 let mut should_restart = false;
5116 let mut disposition = TerminalDisposition::Disabled;
5117 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5118 clear_current_process_facts(state);
5119 state.last_exit = Some(exit_report.clone());
5120 state.lifetime_restarts += 1;
5121 if state.enabled {
5122 state.state = ModuleState::Restarting;
5123 should_restart = true;
5124 disposition = TerminalDisposition::Restarting;
5125 } else {
5126 state.state = ModuleState::Disabled;
5127 }
5128 }) {
5129 error!(module_id = %spec.module_id, error = %err, "failed to record deliberately severed module exit");
5130 return NextAction::Stop {
5131 registration_released: false,
5132 };
5133 }
5134 record_terminal(
5135 &spec.module_id,
5136 terminal_ring,
5137 spawn_events,
5138 &exit_report,
5139 disposition,
5140 );
5141
5142 if should_restart {
5143 NextAction::Restart { schedule: None }
5144 } else {
5145 let registration_released = match wait_for_registration_release(
5146 registry,
5147 &spec.module_id,
5148 REGISTRY_RELEASE_TIMEOUT,
5149 )
5150 .await
5151 {
5152 Ok(()) => true,
5153 Err(err) => {
5154 warn!(module_id = %spec.module_id, error = %err, "registration still active after deliberately severed module exit");
5155 false
5156 }
5157 };
5158 NextAction::Stop {
5159 registration_released,
5160 }
5161 }
5162 }
5163 }
5164}
5165
5166fn record_wait_error_terminal(
5167 module_id: &str,
5168 terminal_ring: &Arc<Mutex<TerminalRing>>,
5169 spawn_events: &SpawnEventFeed,
5170) {
5171 record_terminal(
5172 module_id,
5173 terminal_ring,
5174 spawn_events,
5175 &wait_error_exit_report(),
5176 TerminalDisposition::Failed,
5177 );
5178}
5179
5180fn record_terminal(
5181 module_id: &str,
5182 terminal_ring: &Arc<Mutex<TerminalRing>>,
5183 spawn_events: &SpawnEventFeed,
5184 exit_report: &ExitReport,
5185 disposition: TerminalDisposition,
5186) {
5187 record_terminal_with_detail(
5188 module_id,
5189 terminal_ring,
5190 spawn_events,
5191 exit_report,
5192 disposition,
5193 None,
5194 );
5195}
5196
5197fn durable_terminal_history_of(
5201 terminal_ring: &Mutex<TerminalRing>,
5202 module_id: &str,
5203) -> subc_control::TerminalHistory {
5204 let read = terminal_ring
5205 .lock()
5206 .unwrap_or_else(|p| p.into_inner())
5207 .capture_durable_history();
5208 read.read(module_id)
5209}
5210
5211fn record_terminal_with_detail(
5212 module_id: &str,
5213 terminal_ring: &Arc<Mutex<TerminalRing>>,
5214 spawn_events: &SpawnEventFeed,
5215 exit_report: &ExitReport,
5216 disposition: TerminalDisposition,
5217 disposition_detail: Option<String>,
5218) {
5219 spawn_events.emit_exited(module_id, exit_report.code, exit_report.signal);
5220 let mut ring = terminal_ring
5221 .lock()
5222 .unwrap_or_else(|poisoned| poisoned.into_inner());
5223 let record = TerminalRecord {
5224 exit_code: exit_report.code,
5225 exit_signal: exit_report.signal,
5226 at_ms: exit_report.at_ms,
5227 disposition,
5228 exit_kind: exit_report.kind.into(),
5229 disposition_detail,
5230 };
5231 ring.append_journal(module_id, &record);
5232 ring.push(record);
5233}
5234
5235fn untrack_if_registration_released(
5236 process_liveness: &SupervisorProcessLiveness,
5237 registry: &Registry,
5238 module_id: &str,
5239 snapshot: &SharedSnapshot,
5240) {
5241 match registry.get_module(module_id) {
5242 Ok(None) => process_liveness.untrack_if_current(module_id, snapshot),
5243 Ok(Some(_)) => {}
5244 Err(err) => {
5245 warn!(module_id, error = %err, "could not determine whether supervisor liveness can be untracked");
5246 }
5247 }
5248}
5249
5250#[cfg(test)]
5264fn apply_wire_spawn_args(
5265 command: &mut Command,
5266 spec: &ModuleSpec,
5267 connection_file_path: Option<&std::path::Path>,
5268 handle: Option<&SupervisorHandle>,
5269) -> Result<(), SuperviseError> {
5270 apply_wire_spawn_args_for_role(
5271 command,
5272 spec,
5273 connection_file_path,
5274 handle,
5275 SpawnRole::Plain,
5276 )
5277}
5278
5279fn apply_wire_spawn_args_for_role(
5288 command: &mut Command,
5289 spec: &ModuleSpec,
5290 connection_file_path: Option<&std::path::Path>,
5291 handle: Option<&SupervisorHandle>,
5292 role: SpawnRole,
5293) -> Result<(), SuperviseError> {
5294 command.env(SUBC_MODULE_ID_ENV, &spec.module_id);
5295 if spec.protocol == ModuleProtocol::None {
5296 return Ok(());
5297 }
5298 if let Some(connection_file_path) = connection_file_path {
5299 command.arg(SUBC_ARG).arg(connection_file_path);
5300 }
5301
5302 let nonce = generate_launch_nonce()?;
5306 if let Some(handle) = handle {
5307 match role {
5308 SpawnRole::Plain => {
5309 handle.set_spawn_nonce(&spec.module_id, nonce.clone());
5310 if spec.reserved {
5311 handle.set_reserved_nonce(&spec.module_id, nonce.clone());
5312 }
5313 }
5314 SpawnRole::SwapCandidate => handle.open_swap(&spec.module_id, nonce.clone()),
5315 }
5316 }
5317 command.env(SUBC_LAUNCH_NONCE_ENV, nonce);
5318 Ok(())
5319}
5320
5321fn apply_child_env(command: &mut Command, spec: &ModuleSpec) {
5322 command.env_remove(CK_LOG_ENV);
5323 command.env_remove(SUBC_SPAWN_ROLE_ENV);
5330 for (key, value) in &spec.env {
5331 if matches!(
5335 key.as_str(),
5336 CAPTURE_MAX_FILE_MB_ENV | CAPTURE_KEEP_ENV | CAPTURE_MAX_AGE_DAYS_ENV
5337 ) || key == SUBC_SPAWN_ROLE_ENV
5338 {
5339 continue;
5340 }
5341 command.env(key, value);
5342 }
5343}
5344
5345#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5348enum SpawnRole {
5349 Plain,
5350 SwapCandidate,
5351}
5352
5353fn apply_spawn_role(command: &mut Command, role: SpawnRole) {
5356 if role == SpawnRole::SwapCandidate {
5357 command.env(SUBC_SPAWN_ROLE_ENV, SPAWN_ROLE_SWAP_CANDIDATE);
5358 }
5359}
5360
5361fn spawn_child(
5362 spec: &ModuleSpec,
5363 connection_file_path: Option<&std::path::Path>,
5364 handle: Option<&SupervisorHandle>,
5365 ring: &Arc<Mutex<StderrRing>>,
5366 capture_logs_dir: Option<&std::path::Path>,
5367 roster: &ChildRoster,
5368 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5369) -> Result<SupervisedChild, SuperviseError> {
5370 spawn_child_in_slot(
5371 spec,
5372 connection_file_path,
5373 handle,
5374 ring,
5375 capture_logs_dir,
5376 roster,
5377 #[cfg(target_os = "linux")]
5378 cgroup_placement,
5379 SpawnRole::Plain,
5380 false,
5381 )
5382}
5383
5384#[allow(clippy::too_many_arguments)]
5397fn spawn_child_in_slot(
5398 spec: &ModuleSpec,
5399 connection_file_path: Option<&std::path::Path>,
5400 handle: Option<&SupervisorHandle>,
5401 ring: &Arc<Mutex<StderrRing>>,
5402 capture_logs_dir: Option<&std::path::Path>,
5403 roster: &ChildRoster,
5404 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5405 role: SpawnRole,
5406 alternate_slot: bool,
5407) -> Result<SupervisedChild, SuperviseError> {
5408 if roster.is_closed() {
5409 return Err(SuperviseError::Spawn {
5410 program: spec.program.clone(),
5411 source: io::Error::other("the daemon is shutting down; not starting a new process"),
5412 cgroup_path: None,
5413 });
5414 }
5415 #[cfg(target_os = "linux")]
5416 let cgroup_name = swap::cgroup_name(&spec.module_id, alternate_slot);
5417 #[cfg(not(target_os = "linux"))]
5418 let _ = alternate_slot;
5419 let mut command = Command::new(&spec.program);
5420 command.args(&spec.args);
5421 apply_child_env(&mut command, spec);
5451 apply_spawn_role(&mut command, role);
5452 apply_wire_spawn_args_for_role(&mut command, spec, connection_file_path, handle, role)?;
5453
5454 #[cfg(target_os = "linux")]
5455 let cgroup_path = cgroup_placement
5456 .map(|placement| placement.module_path(&cgroup_name))
5457 .transpose()
5458 .map_err(|source| SuperviseError::Cgroup {
5459 module_id: spec.module_id.clone(),
5460 source,
5461 })?;
5462 #[cfg(not(target_os = "linux"))]
5463 let cgroup_path: Option<PathBuf> = None;
5464 #[cfg(target_os = "linux")]
5465 if let Some(path) = &cgroup_path {
5466 if let Err(error) = apply_cgroup_placement(&mut command, spec, path) {
5467 if let Some(placement) = cgroup_placement {
5468 remove_module_cgroup(placement, &cgroup_name);
5469 }
5470 return Err(error);
5471 }
5472 }
5473
5474 let output_sink = if let Some(logs_dir) = capture_logs_dir {
5475 let path = logs_dir.join(format!("{}.stderr.log", spec.module_id));
5476 match ChildOutputSink::open(&path, capture_retention(spec)) {
5477 Ok(sink) => sink,
5478 Err(error) => {
5479 warn!(
5480 module_id = %spec.module_id,
5481 path = %path.display(),
5482 error = %error,
5483 "could not open child output capture file; forwarding to stderr"
5484 );
5485 ChildOutputSink::Stderr
5486 }
5487 }
5488 } else {
5489 ChildOutputSink::Stderr
5490 };
5491
5492 command.stdout(Stdio::piped());
5493 command.stderr(Stdio::piped());
5494 command.kill_on_drop(true);
5495 #[cfg(unix)]
5512 command.process_group(0);
5513 command.stdin(Stdio::null());
5514 let mut child = match command.spawn() {
5515 Ok(child) => child,
5516 Err(source) => {
5517 #[cfg(target_os = "linux")]
5518 if let Some(placement) = cgroup_placement {
5519 remove_module_cgroup(placement, &cgroup_name);
5520 }
5521 return Err(SuperviseError::Spawn {
5522 program: spec.program.clone(),
5523 source,
5524 cgroup_path,
5525 });
5526 }
5527 };
5528 let spawned_at_ms = unix_ms_now();
5529 let spawned_from = spec.program.clone();
5530 let spawned_file_identity = spawned_file_identity(&spawned_from);
5531 let pid = child.id().ok_or_else(|| SuperviseError::Spawn {
5532 program: spec.program.clone(),
5533 source: io::Error::other("spawned child exposed no live pid"),
5534 cgroup_path: cgroup_path.clone(),
5535 })?;
5536 let process_start_time = crate::provenance::process_start_time(pid);
5537 let process_identity = process_start_time.map(|start_time| ProcessIdentity { pid, start_time });
5538 let roster_guard = roster.admit(
5539 spec.module_id.clone(),
5540 pid,
5541 spec.protocol,
5542 process_start_time,
5543 );
5544
5545 let stdout_pump = match child.stdout.take() {
5546 Some(stdout) => Some(tokio::spawn(pump_stdout_to(stdout, output_sink.clone()))),
5547 None => {
5548 warn!(
5549 module_id = %spec.module_id,
5550 "spawned child exposed no stdout pipe; file capture will be incomplete"
5551 );
5552 None
5553 }
5554 };
5555 let stderr_pump = match child.stderr.take() {
5556 Some(stderr) => {
5557 ring.lock()
5558 .unwrap_or_else(|poisoned| poisoned.into_inner())
5559 .push_process_start();
5560 Some(tokio::spawn(pump_stderr_to(
5561 stderr,
5562 Arc::clone(ring),
5563 output_sink,
5564 )))
5565 }
5566 None => {
5567 ring.lock()
5571 .unwrap_or_else(|poisoned| poisoned.into_inner())
5572 .mark_not_captured("stderr pipe was not available on spawn");
5573 warn!(
5574 module_id = %spec.module_id,
5575 "spawned child exposed no stderr pipe; tail will be unavailable"
5576 );
5577 None
5578 }
5579 };
5580
5581 Ok(SupervisedChild {
5582 child,
5583 #[cfg(target_os = "linux")]
5584 module_id: cgroup_name,
5585 #[cfg(target_os = "linux")]
5586 cgroup_placement: cgroup_placement.cloned(),
5587 stdout_pump,
5588 stderr_pump,
5589 stderr_ring: Arc::clone(ring),
5590 spawned_at_ms,
5591 spawned_from,
5592 spawned_file_identity,
5593 process_start_time,
5594 process_identity,
5595 pid,
5596 roster_guard: Some(roster_guard),
5597 })
5598}
5599
5600#[cfg(target_os = "linux")]
5601fn remove_module_cgroup(placement: &subc_cgroup::Placement, module_id: &str) {
5602 match placement.remove_module(module_id) {
5603 Ok(()) => debug!(module_id, "removed module cgroup after process exit"),
5604 Err(error) => warn!(
5605 module_id,
5606 error = %error,
5607 "could not remove module cgroup after process exit; continuing teardown"
5608 ),
5609 }
5610}
5611
5612#[cfg(target_os = "linux")]
5613fn apply_cgroup_placement(
5614 command: &mut Command,
5615 spec: &ModuleSpec,
5616 path: &std::path::Path,
5617) -> Result<(), SuperviseError> {
5618 subc_cgroup::apply(command, path).map_err(|source| SuperviseError::Cgroup {
5619 module_id: spec.module_id.clone(),
5620 source,
5621 })
5622}
5623
5624fn capture_retention(spec: &ModuleSpec) -> Retention {
5625 let defaults = Retention::default();
5626 let value = |name: &str| {
5627 spec.env
5628 .iter()
5629 .rev()
5630 .find_map(|(key, value)| (key == name).then_some(value.as_str()))
5631 };
5632 Retention {
5633 max_file_mb: value(CAPTURE_MAX_FILE_MB_ENV)
5634 .and_then(|value| value.parse().ok())
5635 .unwrap_or(defaults.max_file_mb),
5636 keep: value(CAPTURE_KEEP_ENV)
5637 .and_then(|value| value.parse().ok())
5638 .unwrap_or(defaults.keep),
5639 max_age_days: value(CAPTURE_MAX_AGE_DAYS_ENV)
5640 .and_then(|value| value.parse().ok())
5641 .unwrap_or(defaults.max_age_days),
5642 }
5643}
5644
5645fn generate_launch_nonce() -> Result<String, SuperviseError> {
5648 let mut bytes = [0u8; 32];
5649 getrandom::getrandom(&mut bytes).map_err(|source| SuperviseError::LaunchNonce {
5650 reason: source.to_string(),
5651 })?;
5652 let mut hex = String::with_capacity(64);
5653 for b in bytes {
5654 use std::fmt::Write;
5655 let _ = write!(hex, "{b:02x}");
5656 }
5657 Ok(hex)
5658}
5659
5660fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
5663 if a.len() != b.len() {
5664 return false;
5665 }
5666 let mut diff = 0u8;
5667 for (x, y) in a.iter().zip(b.iter()) {
5668 diff |= x ^ y;
5669 }
5670 diff == 0
5671}
5672
5673fn spawn_and_mark_running(
5674 spec: &ModuleSpec,
5675 runtime: &SupervisorRuntimeConfig,
5676 snapshot: &SharedSnapshot,
5677) -> Result<SupervisedChild, SuperviseError> {
5678 let child = spawn_child(
5679 spec,
5680 runtime.connection_file_path.as_deref(),
5681 runtime.supervisor_handle.as_ref(),
5682 &runtime.stderr_ring,
5683 runtime.capture_logs_dir.as_deref(),
5684 &runtime.child_roster,
5685 #[cfg(target_os = "linux")]
5686 runtime.cgroup_placement.as_ref(),
5687 )?;
5688 set_running(snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
5689 Ok(child)
5690}
5691
5692enum RegistrationWaitOutcome {
5693 Registered,
5694 Exited(ExitReport),
5695 TimedOut,
5696}
5697
5698struct ReloadRegistrationFailure {
5699 exit_report: ExitReport,
5700 reason: String,
5701}
5702
5703#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5704enum BusyGaugeObservation {
5705 Quiescent,
5706 Busy,
5707 Omitted,
5708}
5709
5710fn busy_gauge_observation(metrics: Option<&Value>, gauges: &[String]) -> BusyGaugeObservation {
5711 let Some(metrics) = metrics.and_then(Value::as_object) else {
5712 return BusyGaugeObservation::Omitted;
5713 };
5714 let mut sum = 0u128;
5715 for gauge in gauges {
5716 let Some(value) = metrics.get(gauge) else {
5717 return BusyGaugeObservation::Omitted;
5718 };
5719 let Some(value) = value.as_u64() else {
5720 return BusyGaugeObservation::Busy;
5721 };
5722 sum = sum.saturating_add(u128::from(value));
5723 }
5724 if sum == 0 {
5725 BusyGaugeObservation::Quiescent
5726 } else {
5727 BusyGaugeObservation::Busy
5728 }
5729}
5730
5731fn declared_busy_gauges(
5732 registry: &Registry,
5733 module_id: &str,
5734) -> Result<Vec<String>, SuperviseError> {
5735 busy_gauges_of(
5736 registry
5737 .get_module(module_id)
5738 .map_err(SuperviseError::Registry)?,
5739 )
5740}
5741
5742fn declared_busy_gauges_for_connection(
5746 registry: &Registry,
5747 connection_id: ConnectionId,
5748) -> Result<Vec<String>, SuperviseError> {
5749 busy_gauges_of(
5750 registry
5751 .get_module_by_connection(connection_id)
5752 .map_err(SuperviseError::Registry)?,
5753 )
5754}
5755
5756fn busy_gauges_of(
5757 registration: Option<crate::registry::ModuleRegistration>,
5758) -> Result<Vec<String>, SuperviseError> {
5759 let Some(registration) = registration else {
5760 return Ok(Vec::new());
5761 };
5762 let Some(self_signals) = registration.manifest.self_signals else {
5763 return Ok(Vec::new());
5764 };
5765
5766 let mut gauges = Vec::new();
5767 for declaration in self_signals {
5768 if declaration.kind != SelfSignalKind::Busy {
5769 continue;
5770 }
5771 match declaration.anchored_to {
5772 SignalAnchor::HealthGauges { gauges: declared } if !declared.is_empty() => {
5773 gauges.extend(declared)
5774 }
5775 _ => {
5776 gauges.push(String::new());
5779 }
5780 }
5781 }
5782 Ok(gauges)
5783}
5784
5785async fn wait_for_forwarding_quiescence(
5790 forwarding: &ForwardingTable,
5791 module_id: &str,
5792 runtime: &SupervisorRuntimeConfig,
5793 endpoint: crate::ModuleEndpointId,
5794 deadline: Instant,
5795 busy_gauges: &[String],
5796 scope: DrainScope,
5797) -> Result<bool, SuperviseError> {
5798 let mut gauges_quiescent = busy_gauges.is_empty();
5799 let mut next_probe_at = Instant::now();
5800 let mut omission_counted = false;
5801
5802 loop {
5803 let now = Instant::now();
5804 if !busy_gauges.is_empty() && now >= next_probe_at && now < deadline {
5805 let report = match scope {
5806 DrainScope::Active => probe_module_health(module_id, runtime, Some(deadline)).await,
5807 DrainScope::Endpoint(endpoint) => {
5808 probe_endpoint_health(endpoint, runtime, Some(deadline)).await
5809 }
5810 };
5811 gauges_quiescent = match report {
5812 Ok(report) => match busy_gauge_observation(report.metrics.as_ref(), busy_gauges) {
5813 BusyGaugeObservation::Quiescent => true,
5814 BusyGaugeObservation::Busy => false,
5815 BusyGaugeObservation::Omitted => {
5816 if !omission_counted {
5817 forwarding
5818 .counters()
5819 .increment_drains_with_undeclared_gauge();
5820 omission_counted = true;
5821 }
5822 false
5823 }
5824 },
5825 Err(err) => {
5826 warn!(
5827 module_id,
5828 error = %err,
5829 "drain health.check did not produce declared busy gauges; treating module as busy"
5830 );
5831 false
5832 }
5833 };
5834 next_probe_at = Instant::now() + runtime.health.cadence.max(REGISTRY_RELEASE_POLL);
5835 }
5836
5837 let in_flight = forwarding
5838 .endpoint_in_flight_count(endpoint)
5839 .map_err(SuperviseError::Forwarding)?;
5840 if in_flight == 0 && gauges_quiescent {
5841 return Ok(true);
5842 }
5843
5844 let now = Instant::now();
5845 if now >= deadline {
5846 return Ok(false);
5847 }
5848 let mut wait = deadline
5849 .saturating_duration_since(now)
5850 .min(REGISTRY_RELEASE_POLL);
5851 if !busy_gauges.is_empty() {
5852 wait = wait.min(next_probe_at.saturating_duration_since(now));
5853 }
5854 sleep(wait).await;
5855 }
5856}
5857
5858fn drained_after_quiescence_wait(wait_result: &Result<bool, SuperviseError>) -> bool {
5866 match wait_result {
5867 Ok(drained) => *drained,
5868 Err(_) => false,
5869 }
5870}
5871
5872fn send_route_goodbyes(forwarding: &ForwardingTable, released_routes: Vec<GoodbyeTarget>) {
5873 for released in released_routes {
5874 let frame = match Frame::build_with_version(
5875 released.negotiated_ver,
5876 FrameType::Goodbye,
5877 control_flags(),
5878 released.channel,
5879 released.epoch,
5880 0,
5881 Vec::new(),
5882 ) {
5883 Ok(frame) => frame,
5884 Err(err) => {
5885 warn!(
5886 route_channel = released.channel,
5887 error = %err,
5888 "failed to build supervisor drain route GOODBYE frame"
5889 );
5890 continue;
5891 }
5892 };
5893 if let Err(err) = released.sink.try_send(frame) {
5894 if released.close_on_delivery_failure() {
5895 warn!(
5896 target_connection_id = released.connection_id.get(),
5897 route_channel = released.channel,
5898 error = %err,
5899 "supervisor drain route GOODBYE was not delivered to client; closing target connection"
5900 );
5901 let _ = forwarding.escalate_client_delivery_failure(
5902 released.connection_id,
5903 released.channel,
5904 released.epoch,
5905 CloseReason::new(
5906 "route_goodbye_delivery_failed",
5907 format!(
5908 "failed to enqueue supervisor drain route GOODBYE for channel {}: {err}",
5909 released.channel
5910 ),
5911 ),
5912 crate::forwarding::UndeliveredFrame {
5913 module_id: released.module_id.as_deref(),
5914 sink: &released.sink,
5915 },
5916 );
5917 } else {
5918 warn!(
5919 target_connection_id = released.connection_id.get(),
5920 route_channel = released.channel,
5921 error = %err,
5922 "supervisor drain route GOODBYE to module dropped under backpressure; not closing shared module connection"
5923 );
5924 }
5925 }
5926 }
5927}
5928
5929fn send_module_draining(
5930 module_id: &str,
5931 reason: RouteCloseReason,
5932 deadline_ms: u64,
5933 target: &ModuleDrainTarget,
5934) {
5935 let body = match serde_json::to_vec(&ModuleControlCommand::Draining {
5936 reason,
5937 deadline_ms,
5938 }) {
5939 Ok(body) => body,
5940 Err(err) => {
5941 warn!(
5942 module_id,
5943 error = %err,
5944 "failed to encode module draining command"
5945 );
5946 return;
5947 }
5948 };
5949 let frame = match Frame::build_with_version(
5950 target.negotiated_ver,
5951 FrameType::Push,
5952 control_flags(),
5953 0,
5954 0,
5955 0,
5956 body,
5957 ) {
5958 Ok(frame) => frame,
5959 Err(err) => {
5960 warn!(
5961 module_id,
5962 error = %err,
5963 "failed to build module draining command frame"
5964 );
5965 return;
5966 }
5967 };
5968 if let Err(err) = target.sink.try_send(frame) {
5969 warn!(
5970 module_id,
5971 target_connection_id = target.endpoint.connection_id.get(),
5972 error = %err,
5973 "module draining command was not delivered to peer"
5974 );
5975 }
5976}
5977
5978fn send_module_goodbye(module_id: &str, forwarding: &ForwardingTable, target: &ModuleDrainTarget) {
5979 let frame = match Frame::build_with_version(
5980 target.negotiated_ver,
5981 FrameType::Goodbye,
5982 control_flags(),
5983 0,
5984 0,
5985 0,
5986 Vec::new(),
5987 ) {
5988 Ok(frame) => frame,
5989 Err(err) => {
5990 warn!(
5991 module_id,
5992 error = %err,
5993 "failed to build supervisor drain module GOODBYE frame"
5994 );
5995 return;
5996 }
5997 };
5998 if let Err(err) = target.sink.try_send(frame) {
5999 warn!(
6000 module_id,
6001 target_connection_id = target.endpoint.connection_id.get(),
6002 error = %err,
6003 "supervisor drain module GOODBYE was not delivered to peer; closing module connection"
6004 );
6005 forwarding.request_connection_close(
6006 target.endpoint.connection_id,
6007 CloseReason::new(
6008 "module_goodbye_delivery_failed",
6009 format!("failed to enqueue supervisor drain module GOODBYE for module '{module_id}': {err}"),
6010 ),
6011 );
6012 }
6013}
6014
6015#[derive(Clone, Copy)]
6016struct ForwardingDrainContext<'a> {
6017 spec: &'a ModuleSpec,
6018 runtime: &'a SupervisorRuntimeConfig,
6019 registry: &'a Registry,
6020 scope: DrainScope,
6021}
6022
6023#[derive(Debug, Clone, Copy, PartialEq, Eq)]
6025enum DrainScope {
6026 Active,
6029 Endpoint(crate::ModuleEndpointId),
6034}
6035
6036async fn begin_forwarding_drain(
6037 spec: &ModuleSpec,
6038 runtime: &SupervisorRuntimeConfig,
6039 registry: &Registry,
6040 snapshot: &SharedSnapshot,
6041 enabled: Option<bool>,
6042 reason: RouteCloseReason,
6043) -> Result<(), SuperviseError> {
6044 let Some(forwarding) = runtime.forwarding.as_ref() else {
6045 return Err(SuperviseError::ReloadUnavailable {
6046 module_id: spec.module_id.clone(),
6047 reason: "supervisor was not configured with a forwarding table".to_string(),
6048 });
6049 };
6050
6051 begin_forwarding_drain_with(
6052 forwarding,
6053 ForwardingDrainContext {
6054 spec,
6055 runtime,
6056 registry,
6057 scope: DrainScope::Active,
6058 },
6059 snapshot,
6060 enabled,
6061 reason,
6062 runtime.drain_timeout,
6063 )
6064 .await
6065}
6066
6067async fn begin_forwarding_drain_if_configured(
6068 spec: &ModuleSpec,
6069 runtime: &SupervisorRuntimeConfig,
6070 registry: &Registry,
6071 snapshot: &SharedSnapshot,
6072 enabled: Option<bool>,
6073 reason: RouteCloseReason,
6074) -> Result<(), SuperviseError> {
6075 begin_forwarding_drain_with_timeout(
6076 spec,
6077 runtime,
6078 registry,
6079 snapshot,
6080 enabled,
6081 reason,
6082 runtime.drain_timeout,
6083 )
6084 .await
6085}
6086
6087async fn begin_forwarding_drain_with_timeout(
6091 spec: &ModuleSpec,
6092 runtime: &SupervisorRuntimeConfig,
6093 registry: &Registry,
6094 snapshot: &SharedSnapshot,
6095 enabled: Option<bool>,
6096 reason: RouteCloseReason,
6097 drain_timeout: Duration,
6098) -> Result<(), SuperviseError> {
6099 let Some(forwarding) = runtime.forwarding.as_ref() else {
6100 return Ok(());
6101 };
6102
6103 begin_forwarding_drain_with(
6104 forwarding,
6105 ForwardingDrainContext {
6106 spec,
6107 runtime,
6108 registry,
6109 scope: DrainScope::Active,
6110 },
6111 snapshot,
6112 enabled,
6113 reason,
6114 drain_timeout,
6115 )
6116 .await
6117}
6118
6119async fn begin_forwarding_drain_with(
6120 forwarding: &ForwardingTable,
6121 context: ForwardingDrainContext<'_>,
6122 snapshot: &SharedSnapshot,
6123 enabled: Option<bool>,
6124 reason: RouteCloseReason,
6125 drain_timeout: Duration,
6126) -> Result<(), SuperviseError> {
6127 let ForwardingDrainContext {
6128 spec,
6129 runtime,
6130 registry,
6131 scope,
6132 } = context;
6133 debug_assert_ne!(reason, RouteCloseReason::Crash);
6134 let terminal = matches!(reason, RouteCloseReason::Disable);
6135 let drain_started_at = Instant::now();
6136 let drain_deadline = drain_started_at + drain_timeout;
6137 let deadline_ms =
6138 unix_ms_now().saturating_add(u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX));
6139 let busy_gauges = match scope {
6140 DrainScope::Active => declared_busy_gauges(registry, &spec.module_id)?,
6141 DrainScope::Endpoint(endpoint) => {
6142 declared_busy_gauges_for_connection(registry, endpoint.connection_id)?
6143 }
6144 };
6145
6146 let drain_target = match scope {
6149 DrainScope::Active => forwarding.begin_module_drain(&spec.module_id, reason),
6150 DrainScope::Endpoint(endpoint) => forwarding.begin_endpoint_drain(endpoint, reason),
6151 }
6152 .map_err(SuperviseError::Forwarding)?;
6153 if scope == DrainScope::Active {
6154 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6155 state.state = ModuleState::Draining;
6156 if let Some(enabled) = enabled {
6157 state.enabled = enabled;
6158 }
6159 })?;
6160 }
6161
6162 if let Some(target) = drain_target.as_ref() {
6163 send_module_draining(&spec.module_id, reason, deadline_ms, target);
6164 let routes = forwarding
6165 .endpoint_routes(target.endpoint)
6166 .map_err(SuperviseError::Forwarding)?;
6167 let routes_notified = routes.len();
6168 crate::control::send_route_control_pushes(
6169 forwarding,
6170 routes.clone(),
6171 ClientControlPush::RouteClosing {
6172 module_id: spec.module_id.clone(),
6173 reason,
6174 },
6175 );
6176 send_route_goodbyes(forwarding, target.abandoned_bindings.clone());
6177
6178 let wait_result = wait_for_forwarding_quiescence(
6184 forwarding,
6185 &spec.module_id,
6186 runtime,
6187 target.endpoint,
6188 drain_deadline,
6189 &busy_gauges,
6190 scope,
6191 )
6192 .await;
6193 let drained = drained_after_quiescence_wait(&wait_result);
6194 if let Err(err) = &wait_result {
6195 error!(
6196 module_id = %spec.module_id,
6197 ?reason,
6198 error = %err,
6199 "forwarding quiescence wait failed after route.closing; forcing route.closed(drained: false) so the client is not left waiting on an unfulfilled promise"
6200 );
6201 } else if !drained {
6202 let holdouts = forwarding
6208 .endpoint_drain_holdouts(target.endpoint)
6209 .unwrap_or_default();
6210 warn!(
6211 module_id = %spec.module_id,
6212 waited = ?drain_timeout,
6213 ?reason,
6214 held_requests = holdouts.requests,
6215 held_routes = holdouts.routes,
6216 total_routes = holdouts.total_routes,
6217 top_connections = ?holdouts.top_connections,
6218 held = %holdouts
6221 .held
6222 .iter()
6223 .map(|(channel, corr)| format!("{channel}:{corr}"))
6224 .collect::<Vec<_>>()
6225 .join(","),
6226 "route drain timed out before request quiescence; forcing teardown"
6227 );
6228 }
6229 crate::control::send_route_control_pushes(
6230 forwarding,
6231 routes,
6232 ClientControlPush::RouteClosed {
6233 module_id: spec.module_id.clone(),
6234 reason,
6235 drained,
6236 abandoned: target.abandoned_bindings.len() as u32,
6237 excluded_subscriptions: target.excluded_subscriptions,
6238 terminal: Some(terminal),
6239 },
6240 );
6241 wait_result?;
6242
6243 let released_routes = match forwarding.release_module_endpoint_routes(target.endpoint) {
6249 Ok(routes) => routes,
6250 Err(err) => {
6251 warn!(
6252 module_id = %spec.module_id,
6253 ?reason,
6254 error = %err,
6255 "failed to release module endpoint routes after route.closed; module GOODBYE will still be sent"
6256 );
6257 send_module_goodbye(&spec.module_id, forwarding, target);
6258 return Err(SuperviseError::Forwarding(err));
6259 }
6260 };
6261 let route_goodbye_count = released_routes.len();
6262 send_route_goodbyes(forwarding, released_routes);
6263 send_module_goodbye(&spec.module_id, forwarding, target);
6264
6265 info!(
6271 module_id = %spec.module_id,
6272 ?reason,
6273 routes_notified,
6274 route_goodbyes = route_goodbye_count,
6275 abandoned_reservations = target.abandoned_bindings.len(),
6276 excluded_subscriptions = target.excluded_subscriptions,
6277 drained,
6278 "module drain complete; consumers notified via route.closing/route.closed pushes and per-route GOODBYE frames"
6279 );
6280 }
6281
6282 Ok(())
6283}
6284
6285async fn wait_for_registration_after_reload(
6288 registry: &Registry,
6289 module_id: &str,
6290 snapshot: &SharedSnapshot,
6291 child: &mut SupervisedChild,
6292 wait: Duration,
6293) -> Result<RegistrationWaitOutcome, SuperviseError> {
6294 wait_for_slot_registration(
6295 registry,
6296 crate::registry::RegistrationSlot::Active(module_id),
6297 module_id,
6298 snapshot,
6299 child,
6300 wait,
6301 )
6302 .await
6303}
6304
6305async fn wait_for_slot_registration(
6313 registry: &Registry,
6314 slot: crate::registry::RegistrationSlot<'_>,
6315 module_id: &str,
6316 snapshot: &SharedSnapshot,
6317 child: &mut SupervisedChild,
6318 wait: Duration,
6319) -> Result<RegistrationWaitOutcome, SuperviseError> {
6320 let deadline = Instant::now() + wait;
6321 loop {
6322 if registry
6323 .registration(slot)
6324 .map_err(SuperviseError::Registry)?
6325 .is_some()
6326 {
6327 return Ok(RegistrationWaitOutcome::Registered);
6328 }
6329
6330 let now = Instant::now();
6331 if now >= deadline {
6332 return Ok(RegistrationWaitOutcome::TimedOut);
6333 }
6334 let remaining = deadline.saturating_duration_since(now);
6335 let poll = remaining.min(REGISTRY_RELEASE_POLL);
6336
6337 tokio::select! {
6338 wait_result = child.wait() => {
6339 let status = wait_result.map_err(|source| SuperviseError::Wait {
6340 module_id: module_id.to_string(),
6341 source,
6342 })?;
6343 return Ok(RegistrationWaitOutcome::Exited(classify_reaped_child_exit(
6344 snapshot,
6345 child,
6346 &status,
6347 )));
6348 }
6349 _ = sleep(poll) => {}
6350 }
6351 }
6352}
6353
6354fn registration_failure_exit_report(mut exit_report: ExitReport) -> ExitReport {
6355 if exit_report.kind != ExitKind::DeliberateSeverance {
6358 exit_report.kind = ExitKind::Crash;
6359 }
6360 exit_report
6361}
6362
6363async fn handle_reload_child_registration_failure(
6364 spec: &ModuleSpec,
6365 runtime: &SupervisorRuntimeConfig,
6366 registry: &Registry,
6367 process_liveness: &SupervisorProcessLiveness,
6368 snapshot: &SharedSnapshot,
6369 child: &mut Option<SupervisedChild>,
6370 failure: ReloadRegistrationFailure,
6371) -> Result<(), SuperviseError> {
6372 let ReloadRegistrationFailure {
6373 exit_report,
6374 reason,
6375 } = failure;
6376 match on_child_exit(
6377 spec,
6378 runtime.restart_policy,
6379 registry,
6380 snapshot,
6381 &runtime.terminal_ring,
6382 &runtime.spawn_events,
6383 exit_report,
6384 )
6385 .await
6386 {
6387 NextAction::Stop {
6388 registration_released,
6389 } => {
6390 if registration_released {
6391 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6392 }
6393 }
6394 NextAction::Restart { schedule } => {
6395 let delay = schedule.map_or(runtime.restart_policy.delay_for_restart(0), |schedule| {
6396 schedule.delay
6397 });
6398 if let Some(schedule) = schedule {
6399 log_crash_respawn(&spec.module_id, schedule);
6400 }
6401 sleep(delay).await;
6402 if respawn_still_pending(snapshot) {
6406 if let Err(err) = wait_for_registration_release(
6407 registry,
6408 &spec.module_id,
6409 REGISTRY_RELEASE_TIMEOUT,
6410 )
6411 .await
6412 {
6413 fail_snapshot(snapshot, Some(&spec.module_id), None);
6414 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6415 return Err(SuperviseError::ReloadFailed {
6416 module_id: spec.module_id.clone(),
6417 reason: format!(
6418 "{reason}; registration did not release before policy retry: {err}"
6419 ),
6420 });
6421 }
6422 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6423 match spawn_and_mark_running(spec, runtime, snapshot) {
6424 Ok(next_child) => {
6425 *child = Some(next_child);
6426 }
6427 Err(err) => {
6428 fail_snapshot(snapshot, Some(&spec.module_id), None);
6429 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6430 return Err(SuperviseError::ReloadFailed {
6431 module_id: spec.module_id.clone(),
6432 reason: format!("{reason}; policy retry spawn failed: {err}"),
6433 });
6434 }
6435 }
6436 }
6437 }
6438 }
6439
6440 Err(SuperviseError::ReloadFailed {
6441 module_id: spec.module_id.clone(),
6442 reason,
6443 })
6444}
6445
6446async fn handle_reload_spawn_failure(
6447 spec: &ModuleSpec,
6448 runtime: &SupervisorRuntimeConfig,
6449 process_liveness: &SupervisorProcessLiveness,
6450 snapshot: &SharedSnapshot,
6451 child: &mut Option<SupervisedChild>,
6452 reason: String,
6453) -> Result<(), SuperviseError> {
6454 let mut should_retry = false;
6455 let now = Instant::now();
6456 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6457 clear_current_process_facts(state);
6458 if daemon_will_restart(state, &runtime.restart_policy, now) {
6459 state.record_crash_restart(&runtime.restart_policy, now);
6460 state.state = ModuleState::Restarting;
6461 should_retry = true;
6462 } else if state.enabled {
6463 state.state = ModuleState::Failed;
6464 } else {
6465 state.state = ModuleState::Disabled;
6466 }
6467 })?;
6468
6469 if should_retry {
6470 sleep(runtime.restart_policy.backoff).await;
6471 if respawn_still_pending(snapshot) {
6475 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6476 match spawn_and_mark_running(spec, runtime, snapshot) {
6477 Ok(next_child) => {
6478 *child = Some(next_child);
6479 }
6480 Err(err) => {
6481 fail_snapshot(snapshot, Some(&spec.module_id), None);
6482 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6483 return Err(SuperviseError::ReloadFailed {
6484 module_id: spec.module_id.clone(),
6485 reason: format!("{reason}; policy retry spawn failed: {err}"),
6486 });
6487 }
6488 }
6489 }
6490 } else {
6491 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6492 }
6493
6494 Err(SuperviseError::ReloadFailed {
6495 module_id: spec.module_id.clone(),
6496 reason,
6497 })
6498}
6499
6500fn control_flags() -> Flags {
6501 Flags::new(false, Priority::Passive, false)
6502}
6503
6504#[allow(clippy::too_many_arguments)]
6505async fn drain_optional_child(
6506 module_id: &str,
6507 protocol: ModuleProtocol,
6508 registry: &Registry,
6509 snapshot: &SharedSnapshot,
6510 terminal_ring: &Arc<Mutex<TerminalRing>>,
6511 spawn_events: &SpawnEventFeed,
6512 child: &mut Option<SupervisedChild>,
6513 drain_timeout: Duration,
6514 final_state: ModuleState,
6515 enabled: Option<bool>,
6516) -> Result<(), SuperviseError> {
6517 if let Some(child) = child.take() {
6518 drain_child_to_state(
6519 module_id,
6520 protocol,
6521 registry,
6522 snapshot,
6523 terminal_ring,
6524 spawn_events,
6525 child,
6526 drain_timeout,
6527 final_state,
6528 enabled,
6529 )
6530 .await
6531 } else {
6532 update_snapshot(snapshot, Some(module_id), |state| {
6533 state.state = final_state;
6534 if let Some(enabled) = enabled {
6535 state.enabled = enabled;
6536 }
6537 clear_current_process_facts(state);
6538 })?;
6539 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
6540 }
6541}
6542
6543#[allow(clippy::too_many_arguments)]
6544async fn drain_child_to_state(
6545 module_id: &str,
6546 protocol: ModuleProtocol,
6547 registry: &Registry,
6548 snapshot: &SharedSnapshot,
6549 terminal_ring: &Arc<Mutex<TerminalRing>>,
6550 spawn_events: &SpawnEventFeed,
6551 mut child: SupervisedChild,
6552 drain_timeout: Duration,
6553 final_state: ModuleState,
6554 enabled: Option<bool>,
6555) -> Result<(), SuperviseError> {
6556 update_snapshot(snapshot, Some(module_id), |state| {
6557 state.state = ModuleState::Draining;
6558 if let Some(enabled) = enabled {
6559 state.enabled = enabled;
6560 }
6561 })?;
6562
6563 if protocol == ModuleProtocol::None {
6569 request_graceful_stop(module_id, &child);
6570 }
6571
6572 let exit_report = match timeout(drain_timeout, child.wait()).await {
6573 Ok(Ok(status)) => classify_reaped_child_exit(snapshot, &child, &status),
6574 Ok(Err(source)) => {
6575 fail_snapshot(snapshot, Some(module_id), None);
6576 return Err(SuperviseError::Wait {
6577 module_id: module_id.to_string(),
6578 source,
6579 });
6580 }
6581 Err(_) => {
6582 child.start_kill().map_err(|source| {
6591 fail_snapshot(snapshot, Some(module_id), None);
6592 SuperviseError::Kill {
6593 module_id: module_id.to_string(),
6594 source,
6595 }
6596 })?;
6597 let status = child.wait().await.map_err(|source| {
6598 fail_snapshot(snapshot, Some(module_id), None);
6599 SuperviseError::Wait {
6600 module_id: module_id.to_string(),
6601 source,
6602 }
6603 })?;
6604 classify_reaped_child_exit(snapshot, &child, &status)
6605 }
6606 };
6607
6608 update_snapshot(snapshot, Some(module_id), |state| {
6609 state.state = final_state;
6610 if let Some(enabled) = enabled {
6611 state.enabled = enabled;
6612 }
6613 clear_current_process_facts(state);
6614 state.last_exit = Some(exit_report.clone());
6615 if exit_report.kind == ExitKind::DeliberateSeverance {
6616 state.lifetime_restarts += 1;
6617 }
6618 })?;
6619 record_terminal(
6620 module_id,
6621 terminal_ring,
6622 spawn_events,
6623 &exit_report,
6624 terminal_disposition(final_state),
6625 );
6626 child.drain_stderr(module_id).await;
6627
6628 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
6629}
6630
6631#[cfg(unix)]
6651fn request_graceful_stop(module_id: &str, child: &SupervisedChild) {
6652 let Some(pid) = child
6653 .id()
6654 .and_then(|pid| i32::try_from(pid).ok())
6655 .and_then(rustix::process::Pid::from_raw)
6656 else {
6657 debug!(
6658 module_id,
6659 "no pid to signal for protocol: none teardown; falling through to the drain wait"
6660 );
6661 return;
6662 };
6663 match rustix::process::kill_process(pid, rustix::process::Signal::TERM) {
6664 Ok(()) => debug!(module_id, "sent SIGTERM to protocol: none module"),
6665 Err(err) => debug!(
6666 module_id,
6667 error = %err,
6668 "SIGTERM to protocol: none module failed; the drain wait and kill still apply"
6669 ),
6670 }
6671}
6672
6673#[cfg(not(unix))]
6681fn request_graceful_stop(module_id: &str, _child: &SupervisedChild) {
6682 debug!(
6683 module_id,
6684 "no graceful stop signal exists on this platform; protocol: none teardown waits, then kills"
6685 );
6686}
6687
6688fn terminal_disposition(final_state: ModuleState) -> TerminalDisposition {
6689 match final_state {
6690 ModuleState::Stopped => TerminalDisposition::Stopped,
6691 ModuleState::Disabled => TerminalDisposition::Disabled,
6692 ModuleState::Restarting => TerminalDisposition::Restarting,
6693 ModuleState::Failed => TerminalDisposition::Failed,
6694 ModuleState::Starting
6695 | ModuleState::Running
6696 | ModuleState::Unresponsive
6697 | ModuleState::Draining => {
6698 unreachable!("terminal exits only finish in terminal or restarting states")
6699 }
6700 }
6701}
6702
6703async fn wait_for_registration_release(
6706 registry: &Registry,
6707 module_id: &str,
6708 wait: Duration,
6709) -> Result<(), SuperviseError> {
6710 wait_for_slot_registration_release(
6711 registry,
6712 crate::registry::RegistrationSlot::Active(module_id),
6713 wait,
6714 )
6715 .await
6716}
6717
6718async fn wait_for_slot_registration_release(
6726 registry: &Registry,
6727 slot: crate::registry::RegistrationSlot<'_>,
6728 wait: Duration,
6729) -> Result<(), SuperviseError> {
6730 let deadline = Instant::now() + wait;
6731 let mut release_events = registration_release_events().subscribe();
6732 let still_active = |registration: &crate::registry::ModuleRegistration| {
6733 SuperviseError::RegistrationStillActive {
6734 module_id: registration.manifest.module_id.clone(),
6735 waited: wait,
6736 }
6737 };
6738 loop {
6739 let _observed_generation = *release_events.borrow_and_update();
6740 let Some(registration) = registry
6741 .registration(slot)
6742 .map_err(SuperviseError::Registry)?
6743 else {
6744 return Ok(());
6745 };
6746
6747 let now = Instant::now();
6748 if now >= deadline {
6749 return Err(still_active(®istration));
6750 }
6751
6752 let remaining = deadline.saturating_duration_since(now);
6753 match timeout(remaining, release_events.changed()).await {
6754 Ok(Ok(())) | Ok(Err(_)) => {}
6755 Err(_) => return Err(still_active(®istration)),
6756 }
6757 }
6758}
6759
6760#[cfg(test)]
6761mod slot_registration_wait_tests {
6762 use super::*;
6763 use crate::registry::{ConnectionId, RegistrationSlot};
6764 use subc_protocol::manifest::ModuleManifest;
6765
6766 const INCUMBENT: u64 = 1;
6767 const CANDIDATE: u64 = 2;
6768
6769 fn swapped_registry() -> Arc<Registry> {
6770 let registry = Arc::new(Registry::default());
6771 let manifest = ModuleManifest::builder("m", "0.1.0").build();
6772 registry
6773 .register_with_control_ops(
6774 manifest.clone(),
6775 1,
6776 ConnectionId::new(INCUMBENT),
6777 Vec::new(),
6778 )
6779 .unwrap();
6780 registry
6781 .register_candidate_with_control_ops(
6782 manifest,
6783 1,
6784 ConnectionId::new(CANDIDATE),
6785 Vec::new(),
6786 )
6787 .unwrap();
6788 registry
6789 }
6790
6791 #[tokio::test]
6795 async fn incumbent_release_is_awaited_by_connection_not_by_module_id() {
6796 let registry = swapped_registry();
6797 registry.promote_candidate("m").unwrap().unwrap();
6798
6799 assert!(matches!(
6800 wait_for_registration_release(®istry, "m", Duration::from_millis(50)).await,
6801 Err(SuperviseError::RegistrationStillActive { .. })
6802 ));
6803
6804 assert!(matches!(
6806 wait_for_slot_registration_release(
6807 ®istry,
6808 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
6809 Duration::from_millis(50),
6810 )
6811 .await,
6812 Err(SuperviseError::RegistrationStillActive { .. })
6813 ));
6814
6815 let releaser = Arc::clone(®istry);
6816 let release = tokio::spawn(async move {
6817 sleep(Duration::from_millis(20)).await;
6818 releaser
6819 .deregister_connection(ConnectionId::new(INCUMBENT))
6820 .unwrap();
6821 notify_registration_release();
6822 });
6823 wait_for_slot_registration_release(
6824 ®istry,
6825 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
6826 Duration::from_secs(5),
6827 )
6828 .await
6829 .expect("the incumbent's own registration is released");
6830 release.await.unwrap();
6831 assert!(registry.get_module("m").unwrap().is_some());
6832 }
6833
6834 #[tokio::test]
6837 async fn candidate_slot_wait_ignores_the_incumbents_registration() {
6838 let registry = swapped_registry();
6839 assert!(matches!(
6840 wait_for_slot_registration_release(
6841 ®istry,
6842 RegistrationSlot::Candidate("m"),
6843 Duration::from_millis(50),
6844 )
6845 .await,
6846 Err(SuperviseError::RegistrationStillActive { .. })
6847 ));
6848 registry
6849 .deregister_connection(ConnectionId::new(CANDIDATE))
6850 .unwrap();
6851 wait_for_slot_registration_release(
6852 ®istry,
6853 RegistrationSlot::Candidate("m"),
6854 Duration::from_millis(50),
6855 )
6856 .await
6857 .expect("a candidate slot with no candidate is released");
6858 assert!(registry
6859 .registration(RegistrationSlot::Active("m"))
6860 .unwrap()
6861 .is_some());
6862 }
6863}
6864
6865fn classify_exit(status: &ExitStatus) -> ExitReport {
6866 ExitReport {
6867 kind: if status.success() {
6868 ExitKind::Clean
6869 } else {
6870 ExitKind::Crash
6871 },
6872 code: status.code(),
6873 signal: exit_signal(status),
6874 at_ms: unix_ms_now(),
6875 }
6876}
6877
6878fn wait_error_exit_report() -> ExitReport {
6884 ExitReport {
6885 kind: ExitKind::Crash,
6886 code: None,
6887 signal: None,
6888 at_ms: unix_ms_now(),
6889 }
6890}
6891
6892#[cfg(unix)]
6893fn exit_signal(status: &ExitStatus) -> Option<i32> {
6894 use std::os::unix::process::ExitStatusExt;
6895
6896 status.signal()
6897}
6898
6899#[cfg(not(unix))]
6900fn exit_signal(_status: &ExitStatus) -> Option<i32> {
6901 None
6902}
6903
6904fn reset_restart_count(snapshot: &SharedSnapshot, module_id: &str) -> Result<(), SuperviseError> {
6910 update_snapshot(snapshot, Some(module_id), |state| {
6911 state.clear_crash_restarts();
6912 })
6913}
6914
6915fn set_running(
6916 snapshot: &SharedSnapshot,
6917 child: &SupervisedChild,
6918 module_id: &str,
6919 spawn_events: &SpawnEventFeed,
6920) -> Result<(), SuperviseError> {
6921 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
6922 module_id: Some(module_id.to_string()),
6923 })?;
6924 state.spawn_generation = spawn_events.emit_spawned(module_id, child.pid, child.spawned_at_ms);
6925 state.in_alternate_slot = false;
6928 state.state = ModuleState::Running;
6929 state.enabled = true;
6930 state.process_alive = true;
6931 state.pid = child.id();
6932 state.spawned_at_ms = Some(child.spawned_at_ms);
6933 state.spawned_from = Some(child.spawned_from.clone());
6934 state.spawned_file_identity = child.spawned_file_identity;
6935 state.process_start_time = child.process_start_time;
6936 Ok(())
6937}
6938
6939fn clear_current_process_facts(state: &mut SupervisorSnapshot) {
6940 state.process_alive = false;
6941 state.pid = None;
6942 state.spawned_at_ms = None;
6943 state.spawned_from = None;
6944 state.spawned_file_identity = None;
6945 state.process_start_time = None;
6946 state.deliberate_severance = None;
6947}
6948
6949#[cfg(test)]
6950fn record_deliberate_severance(
6951 snapshot: &SharedSnapshot,
6952 identity: ProcessIdentity,
6953) -> Result<(), SuperviseError> {
6954 update_snapshot(snapshot, None, |state| {
6955 state.deliberate_severance = Some(identity);
6956 })
6957}
6958
6959fn apply_deliberate_severance_marker(
6960 snapshot: &SharedSnapshot,
6961 exited_identity: Option<ProcessIdentity>,
6962 mut exit_report: ExitReport,
6963) -> ExitReport {
6964 let marker = lock_snapshot(snapshot)
6965 .ok()
6966 .and_then(|mut state| state.deliberate_severance.take());
6967 if marker.is_some() && marker == exited_identity {
6968 exit_report.kind = ExitKind::DeliberateSeverance;
6969 }
6970 exit_report
6971}
6972
6973fn classify_reaped_child_exit(
6974 snapshot: &SharedSnapshot,
6975 child: &SupervisedChild,
6976 status: &ExitStatus,
6977) -> ExitReport {
6978 apply_deliberate_severance_marker(snapshot, child.process_identity(), classify_exit(status))
6979}
6980
6981fn fail_snapshot(
6982 snapshot: &SharedSnapshot,
6983 module_id: Option<&str>,
6984 last_exit: Option<ExitReport>,
6985) {
6986 if let Err(err) = update_snapshot(snapshot, module_id, |state| {
6987 state.state = ModuleState::Failed;
6988 clear_current_process_facts(state);
6989 if let Some(last_exit) = last_exit {
6990 state.last_exit = Some(last_exit);
6991 }
6992 }) {
6993 error!(error = %err, "failed to mark supervisor state failed");
6994 }
6995}
6996
6997fn update_snapshot(
6998 snapshot: &SharedSnapshot,
6999 module_id: Option<&str>,
7000 update: impl FnOnce(&mut SupervisorSnapshot),
7001) -> Result<(), SuperviseError> {
7002 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
7003 module_id: module_id.map(ToOwned::to_owned),
7004 })?;
7005 update(&mut state);
7006 Ok(())
7007}
7008
7009const SLOW_SNAPSHOT_LOCK_THRESHOLD: Duration = Duration::from_millis(250);
7010
7011fn lock_snapshot_for_control<'a>(
7012 snapshot: &'a SharedSnapshot,
7013 module_id: &str,
7014 caller: &'static str,
7015) -> Result<std::sync::MutexGuard<'a, SupervisorSnapshot>, SuperviseError> {
7016 let started_at = Instant::now();
7017 let guard = lock_snapshot(snapshot)?;
7018 let waited = started_at.elapsed();
7019 if waited >= SLOW_SNAPSHOT_LOCK_THRESHOLD {
7020 warn!(
7021 module_id = %module_id,
7022 waited_ms = waited.as_millis() as u64,
7023 caller = %caller,
7024 "slow snapshot lock"
7025 );
7026 }
7027 Ok(guard)
7028}
7029
7030fn lock_snapshot(
7031 snapshot: &SharedSnapshot,
7032) -> Result<std::sync::MutexGuard<'_, SupervisorSnapshot>, SuperviseError> {
7033 snapshot
7034 .lock()
7035 .map_err(|_| SuperviseError::StatePoisoned { module_id: None })
7036}
7037
7038#[cfg(test)]
7039mod terminal_history_tests {
7040 use std::{
7041 path::PathBuf,
7042 sync::Arc,
7043 time::{Duration, Instant},
7044 };
7045
7046 use tokio::time::sleep;
7047
7048 use super::{
7049 apply_deliberate_severance_marker, daemon_will_restart, drain_child_to_state,
7050 drained_after_quiescence_wait, handle_reload_spawn_failure, health_restart_child,
7051 lock_snapshot, on_child_exit, record_deliberate_severance, record_wait_error_terminal,
7052 reset_restart_count, spawn_and_mark_running, update_snapshot, wait_error_exit_report,
7053 ExitKind, ExitReport, ModuleProtocol, ModuleSpec, ModuleState, NextAction, ProcessIdentity,
7054 RestartPolicy, SpawnEventKind, SuperviseError, SupervisedModule, Supervisor,
7055 SupervisorHandle, SupervisorHealthStatus, SupervisorSnapshot,
7056 };
7057 use super::Instant as ClockInstant;
7062 use crate::{
7063 registry::Registry,
7064 terminal_ring::{TerminalRing, TerminalRingConfig},
7065 };
7066 use std::sync::Mutex;
7067 use subc_control::TerminalDisposition;
7068
7069 fn fake_aft_stub_path() -> PathBuf {
7074 let mut path = std::env::current_exe().expect("current_exe available in tests");
7075 path.pop();
7076 path.pop();
7077 path.push(if cfg!(windows) {
7078 "fake-aft-stub.exe"
7079 } else {
7080 "fake-aft-stub"
7081 });
7082 assert!(
7083 path.exists(),
7084 "fake-aft-stub not built at {}: run `cargo test -p subc-core` (which builds \
7085 [[bin]] targets) rather than `cargo test -p subc-core --lib` (which does not)",
7086 path.display()
7087 );
7088 path
7089 }
7090
7091 #[test]
7092 fn reserved_never_spawned_refuses_every_hello() {
7093 let supervisor = SupervisorHandle::default();
7098 supervisor.apply_identity_configuration(&ModuleSpec {
7099 module_id: "never-spawned".to_string(),
7100 program: PathBuf::from("/usr/bin/false"),
7101 args: Vec::new(),
7102 env: Vec::new(),
7103 reserved: true,
7104 reserved_prefixes: Vec::new(),
7105 protocol: ModuleProtocol::Subc,
7106 overlap: Default::default(),
7107 });
7108 assert!(
7109 supervisor
7110 .reserved_hello_rejection("never-spawned", Some("any-forged-nonce"))
7111 .is_some(),
7112 "forged nonce must refuse on a reserved never-spawned id"
7113 );
7114 assert!(
7115 supervisor
7116 .reserved_hello_rejection("never-spawned", None)
7117 .is_some(),
7118 "absent nonce must refuse on a reserved never-spawned id"
7119 );
7120 supervisor.set_spawn_nonce("never-spawned", "minted".to_string());
7122 supervisor.apply_identity_configuration(&ModuleSpec {
7123 module_id: "never-spawned".to_string(),
7124 program: PathBuf::from("/usr/bin/false"),
7125 args: Vec::new(),
7126 env: Vec::new(),
7127 reserved: true,
7128 reserved_prefixes: Vec::new(),
7129 protocol: ModuleProtocol::Subc,
7130 overlap: Default::default(),
7131 });
7132 assert!(supervisor
7133 .reserved_hello_rejection("never-spawned", Some("minted"))
7134 .is_none());
7135 assert!(supervisor
7136 .reserved_hello_rejection("never-spawned", Some("forged"))
7137 .is_some());
7138 }
7139
7140 fn seed_crash_restarts(state: &mut SupervisorSnapshot, count: u32) {
7143 let now = ClockInstant::now();
7144 for _ in 0..count {
7145 state.crash_restarts.push_back(now);
7146 }
7147 }
7148
7149 fn age_oldest_crash_restart_out_of_window(state: &mut SupervisorSnapshot, window: Duration) {
7153 let aged = state
7154 .crash_restarts
7155 .front()
7156 .expect("a crash restart must be recorded before it can be aged")
7157 .checked_sub(window + Duration::from_secs(1))
7158 .expect("the test clock is far enough from its origin to age an instant");
7159 state.crash_restarts[0] = aged;
7160 }
7161
7162 fn snapshot_with_restarts(enabled: bool, count: u32) -> SupervisorSnapshot {
7163 let mut state = SupervisorSnapshot::new(ModuleState::Running, enabled);
7164 seed_crash_restarts(&mut state, count);
7165 state
7166 }
7167
7168 #[test]
7169 fn daemon_owned_recovery_predicate_uses_the_pre_increment_budget() {
7170 let policy = RestartPolicy::new(3, Duration::ZERO);
7171 let now = ClockInstant::now();
7172 assert!(daemon_will_restart(
7173 &mut snapshot_with_restarts(true, 2),
7174 &policy,
7175 now
7176 ));
7177 assert!(!daemon_will_restart(
7178 &mut snapshot_with_restarts(true, 3),
7179 &policy,
7180 now
7181 ));
7182 assert!(!daemon_will_restart(
7183 &mut snapshot_with_restarts(false, 0),
7184 &policy,
7185 now
7186 ));
7187 }
7188
7189 #[test]
7190 fn crash_restart_backoff_escalates_with_in_window_count() {
7191 let policy = RestartPolicy::new(4, Duration::from_millis(100))
7192 .with_max_backoff(Duration::from_secs(30));
7193 let now = ClockInstant::now();
7194 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7195 let schedules = (0..4)
7196 .map(|_| {
7197 state
7198 .next_crash_restart(&policy, now)
7199 .expect("the test policy allows four crash restarts")
7200 })
7201 .collect::<Vec<_>>();
7202
7203 assert_eq!(
7204 schedules
7205 .iter()
7206 .map(|schedule| schedule.restart_in_window)
7207 .collect::<Vec<_>>(),
7208 vec![0, 1, 2, 3]
7209 );
7210 assert_eq!(
7211 schedules
7212 .iter()
7213 .map(|schedule| schedule.delay)
7214 .collect::<Vec<_>>(),
7215 vec![
7216 Duration::from_millis(100),
7217 Duration::from_secs(1),
7218 Duration::from_secs(10),
7219 Duration::from_secs(30),
7220 ]
7221 );
7222 }
7223
7224 #[test]
7225 fn crash_restart_backoff_resets_after_ring_clear() {
7226 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7227 let now = ClockInstant::now();
7228 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7229 assert_eq!(
7230 state.next_crash_restart(&policy, now).unwrap().delay,
7231 Duration::from_millis(100)
7232 );
7233 assert_eq!(
7234 state.next_crash_restart(&policy, now).unwrap().delay,
7235 Duration::from_secs(1)
7236 );
7237
7238 state.clear_crash_restarts();
7239 let schedule = state
7240 .next_crash_restart(&policy, now)
7241 .expect("a cleared ring must allow another restart");
7242 assert_eq!(schedule.restart_in_window, 0);
7243 assert_eq!(schedule.delay, Duration::from_millis(100));
7244 }
7245
7246 #[test]
7247 fn crash_restart_backoff_ignores_aged_restarts() {
7248 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7249 let now = ClockInstant::now();
7250 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7251 state
7252 .next_crash_restart(&policy, now)
7253 .expect("the first restart is allowed");
7254 state
7255 .next_crash_restart(&policy, now)
7256 .expect("the second restart is allowed");
7257 state.crash_restarts[0] = now
7258 .checked_sub(policy.window + Duration::from_secs(1))
7259 .expect("the fake clock can age a restart past the window");
7260
7261 let schedule = state
7262 .next_crash_restart(&policy, now)
7263 .expect("an aged restart must release its slot");
7264 assert_eq!(schedule.restart_in_window, 1);
7265 assert_eq!(schedule.delay, Duration::from_secs(1));
7266 assert_eq!(state.crash_restarts.len(), 2);
7267 }
7268
7269 #[test]
7273 fn a_budget_spent_before_the_window_no_longer_refuses() {
7274 let policy = RestartPolicy::new(3, Duration::ZERO);
7275 let mut state = snapshot_with_restarts(true, 3);
7276 let now = ClockInstant::now();
7277 assert!(!daemon_will_restart(&mut state, &policy, now));
7278
7279 assert!(daemon_will_restart(
7280 &mut state,
7281 &policy,
7282 now + policy.window + Duration::from_secs(1)
7283 ));
7284 assert!(
7285 state.crash_restarts.is_empty(),
7286 "reading the budget must drop the instants that left the window"
7287 );
7288 }
7289
7290 fn module_with_recovery_snapshot(
7291 state: ModuleState,
7292 enabled: bool,
7293 restart_count: u32,
7294 ) -> SupervisedModule {
7295 let registry = Arc::new(Registry::default());
7296 let supervisor =
7297 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(3, Duration::ZERO));
7298 let module = supervisor
7299 .spawn(ModuleSpec {
7300 module_id: "recovery-snapshot".to_string(),
7301 program: fake_aft_stub_path(),
7302 args: Vec::new(),
7303 env: Vec::new(),
7304 reserved: false,
7305 reserved_prefixes: Vec::new(),
7306 protocol: ModuleProtocol::Subc,
7307 overlap: Default::default(),
7308 })
7309 .unwrap();
7310 update_snapshot(
7311 &module.inner.snapshot,
7312 Some("recovery-snapshot"),
7313 |snapshot| {
7314 snapshot.state = state;
7315 snapshot.enabled = enabled;
7316 seed_crash_restarts(snapshot, restart_count);
7317 },
7318 )
7319 .unwrap();
7320 module
7321 }
7322
7323 #[cfg(target_os = "linux")]
7324 #[tokio::test]
7325 async fn no_cgroup_placement_does_not_block_fake_aft_stub_spawn() {
7326 let supervisor = Supervisor::new(Arc::new(Registry::default()), RestartPolicy::default())
7327 .with_cgroup_placement(None);
7328 let result = supervisor.spawn(ModuleSpec {
7329 module_id: "no-cgroup-placement".to_string(),
7330 program: fake_aft_stub_path(),
7331 args: Vec::new(),
7332 env: Vec::new(),
7333 reserved: false,
7334 reserved_prefixes: Vec::new(),
7335 protocol: ModuleProtocol::Subc,
7336 overlap: Default::default(),
7337 });
7338
7339 assert!(
7340 result.is_ok(),
7341 "no delegation must not turn an otherwise valid spawn into a failure: {result:?}"
7342 );
7343 }
7344
7345 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7346 async fn undecided_snapshot_uses_shared_restart_predicate() {
7347 assert!(module_with_recovery_snapshot(ModuleState::Running, true, 2)
7348 .will_recover_after_connection_loss()
7349 .unwrap());
7350 assert!(
7351 !module_with_recovery_snapshot(ModuleState::Running, true, 3)
7352 .will_recover_after_connection_loss()
7353 .unwrap()
7354 );
7355 }
7356
7357 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7358 async fn restarting_snapshot_at_exhausted_budget_is_non_terminal() {
7359 assert!(
7360 module_with_recovery_snapshot(ModuleState::Restarting, true, 3)
7361 .will_recover_after_connection_loss()
7362 .unwrap()
7363 );
7364 }
7365
7366 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7367 async fn terminal_phase_snapshots_are_terminal_before_budget_exhaustion() {
7368 assert!(!module_with_recovery_snapshot(ModuleState::Failed, true, 0)
7369 .will_recover_after_connection_loss()
7370 .unwrap());
7371 assert!(
7372 !module_with_recovery_snapshot(ModuleState::Disabled, true, 0)
7373 .will_recover_after_connection_loss()
7374 .unwrap()
7375 );
7376 }
7377
7378 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7379 async fn warming_snapshot_is_limited_to_startup_phases() {
7380 for state in [
7381 ModuleState::Starting,
7382 ModuleState::Running,
7383 ModuleState::Restarting,
7384 ] {
7385 assert!(
7386 module_with_recovery_snapshot(state, true, 0)
7387 .is_warming()
7388 .unwrap(),
7389 "{state:?} should be warming"
7390 );
7391 }
7392 for state in [
7393 ModuleState::Unresponsive,
7394 ModuleState::Draining,
7395 ModuleState::Stopped,
7396 ModuleState::Failed,
7397 ModuleState::Disabled,
7398 ] {
7399 assert!(
7400 !module_with_recovery_snapshot(state, true, 0)
7401 .is_warming()
7402 .unwrap(),
7403 "{state:?} should not be warming"
7404 );
7405 }
7406 }
7407
7408 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7409 async fn terminal_history_survives_respawn_and_keeps_both_crashes_in_order() {
7410 let registry = Arc::new(Registry::default());
7411 let supervisor =
7412 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(1, Duration::ZERO));
7413 let module = supervisor
7414 .spawn(ModuleSpec {
7415 module_id: "terminal-history".to_string(),
7416 program: fake_aft_stub_path(),
7417 args: Vec::new(),
7418 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7419 reserved: false,
7420 reserved_prefixes: Vec::new(),
7421 protocol: ModuleProtocol::Subc,
7422 overlap: Default::default(),
7423 })
7424 .unwrap();
7425
7426 let deadline = Instant::now() + Duration::from_secs(5);
7427 loop {
7428 let history = module.terminal_history();
7429 if history.entries.len() == 2 {
7430 assert_eq!(module.status().unwrap().state, ModuleState::Failed);
7431 assert_eq!(history.dropped, 0);
7432 assert_eq!(
7433 history
7434 .entries
7435 .iter()
7436 .map(|entry| entry.exit_code)
7437 .collect::<Vec<_>>(),
7438 vec![Some(23), Some(23)]
7439 );
7440 assert!(history.entries[0].at_ms <= history.entries[1].at_ms);
7441 return;
7442 }
7443 assert!(
7444 Instant::now() < deadline,
7445 "module did not retain two terminal exits: {history:?}"
7446 );
7447 sleep(Duration::from_millis(10)).await;
7448 }
7449 }
7450
7451 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7455 async fn disable_during_crash_backoff_cancels_pending_respawn() {
7456 let backoff = Duration::from_secs(2);
7457 let supervisor = Supervisor::new(
7458 Arc::new(Registry::default()),
7459 RestartPolicy::new(10, backoff),
7460 );
7461 let module = supervisor
7462 .spawn(ModuleSpec {
7463 module_id: "disable-during-backoff".to_string(),
7464 program: fake_aft_stub_path(),
7465 args: Vec::new(),
7466 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7467 reserved: false,
7468 reserved_prefixes: Vec::new(),
7469 protocol: ModuleProtocol::Subc,
7470 overlap: Default::default(),
7471 })
7472 .unwrap();
7473
7474 let deadline = Instant::now() + Duration::from_secs(5);
7476 loop {
7477 if module.status().unwrap().state == ModuleState::Restarting {
7478 break;
7479 }
7480 assert!(
7481 Instant::now() < deadline,
7482 "module never entered the crash backoff"
7483 );
7484 sleep(Duration::from_millis(10)).await;
7485 }
7486
7487 let started = Instant::now();
7488 module.set_enabled(false).await.unwrap();
7489 let waited = started.elapsed();
7490
7491 assert!(
7492 waited < backoff / 2,
7493 "disable waited {waited:?} behind the {backoff:?} crash backoff; the operator command must preempt the pending respawn"
7494 );
7495 assert_eq!(module.status().unwrap().state, ModuleState::Disabled);
7496
7497 sleep(backoff + Duration::from_millis(500)).await;
7499 let status = module.status().unwrap();
7500 assert_eq!(status.state, ModuleState::Disabled);
7501 assert_eq!(
7502 status.spawn_generation, 1,
7503 "module respawned after the operator disabled it"
7504 );
7505 }
7506
7507 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7511 async fn every_restart_increment_path_advances_lifetime_count() {
7512 let supervisor = Supervisor::new(
7513 Arc::new(Registry::default()),
7514 RestartPolicy::new(1, Duration::ZERO),
7515 );
7516 let runtime = supervisor.runtime_config();
7517 let spec = ModuleSpec {
7518 module_id: "lifetime-increment-path".to_string(),
7519 program: PathBuf::from("/unused/lifetime-increment-path"),
7520 args: Vec::new(),
7521 env: Vec::new(),
7522 reserved: false,
7523 reserved_prefixes: Vec::new(),
7524 protocol: ModuleProtocol::Subc,
7525 overlap: Default::default(),
7526 };
7527
7528 let crash_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7529 assert!(matches!(
7530 on_child_exit(
7531 &spec,
7532 runtime.restart_policy,
7533 &supervisor.registry,
7534 &crash_snapshot,
7535 &runtime.terminal_ring,
7536 &runtime.spawn_events,
7537 ExitReport {
7538 kind: ExitKind::Crash,
7539 code: Some(1),
7540 signal: None,
7541 at_ms: 1,
7542 },
7543 )
7544 .await,
7545 NextAction::Restart { schedule: _ }
7546 ));
7547 let (crash_restarts, crash_lifetime) = {
7548 let state = lock_snapshot(&crash_snapshot).unwrap();
7549 (state.crash_restarts.len(), state.lifetime_restarts)
7550 };
7551 assert_eq!(crash_restarts, 1);
7552 assert_eq!(crash_lifetime, 1);
7553
7554 let health_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7555 let mut health_child = None;
7556 assert!(matches!(
7557 health_restart_child(
7558 &spec,
7559 &runtime,
7560 &supervisor.registry,
7561 &supervisor.process_liveness,
7562 &health_snapshot,
7563 &mut health_child,
7564 SupervisorHealthStatus::Failing,
7565 None,
7566 2,
7567 )
7568 .await,
7569 Err(SuperviseError::Spawn { .. })
7570 ));
7571 let (health_restarts, health_lifetime) = {
7572 let state = lock_snapshot(&health_snapshot).unwrap();
7573 (state.crash_restarts.len(), state.lifetime_restarts)
7574 };
7575 assert_eq!(health_restarts, 1);
7576 assert_eq!(health_lifetime, 1);
7577
7578 let reload_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7579 let mut reload_child = None;
7580 assert!(matches!(
7581 handle_reload_spawn_failure(
7582 &spec,
7583 &runtime,
7584 &supervisor.process_liveness,
7585 &reload_snapshot,
7586 &mut reload_child,
7587 "forced reload spawn failure".to_string(),
7588 )
7589 .await,
7590 Err(SuperviseError::ReloadFailed { .. })
7591 ));
7592 let (reload_restarts, reload_lifetime) = {
7593 let state = lock_snapshot(&reload_snapshot).unwrap();
7594 (state.crash_restarts.len(), state.lifetime_restarts)
7595 };
7596 assert_eq!(reload_restarts, 1);
7597 assert_eq!(reload_lifetime, 1);
7598 }
7599
7600 #[tokio::test]
7601 async fn deliberately_severed_live_child_records_lifetime_without_spending_restart_budget() {
7602 let supervisor = Supervisor::new(
7603 Arc::new(Registry::default()),
7604 RestartPolicy::new(3, Duration::ZERO),
7605 );
7606 let runtime = supervisor.runtime_config();
7607 let spec = ModuleSpec {
7608 module_id: "deliberately-severed".to_string(),
7609 program: PathBuf::from("/unused/deliberately-severed"),
7610 args: Vec::new(),
7611 env: Vec::new(),
7612 reserved: false,
7613 reserved_prefixes: Vec::new(),
7614 protocol: ModuleProtocol::Subc,
7615 overlap: Default::default(),
7616 };
7617 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7618 let process = ProcessIdentity {
7619 pid: 41,
7620 start_time: 101,
7621 };
7622 record_deliberate_severance(&snapshot, process).unwrap();
7623 let exit_report = apply_deliberate_severance_marker(
7624 &snapshot,
7625 Some(process),
7626 ExitReport {
7627 kind: ExitKind::Crash,
7628 code: Some(1),
7629 signal: None,
7630 at_ms: 1,
7631 },
7632 );
7633 assert_eq!(exit_report.kind, ExitKind::DeliberateSeverance);
7634
7635 assert!(matches!(
7636 on_child_exit(
7637 &spec,
7638 runtime.restart_policy,
7639 &supervisor.registry,
7640 &snapshot,
7641 &runtime.terminal_ring,
7642 &runtime.spawn_events,
7643 exit_report,
7644 )
7645 .await,
7646 NextAction::Restart { schedule: _ }
7647 ));
7648 let state = lock_snapshot(&snapshot).unwrap();
7649 assert_eq!(state.lifetime_restarts, 1);
7650 assert_eq!(state.crash_restarts.len(), 0);
7651 }
7652
7653 #[tokio::test]
7654 async fn genuine_crash_spends_restart_budget_and_records_lifetime() {
7655 let supervisor = Supervisor::new(
7656 Arc::new(Registry::default()),
7657 RestartPolicy::new(3, Duration::ZERO),
7658 );
7659 let runtime = supervisor.runtime_config();
7660 let spec = ModuleSpec {
7661 module_id: "genuine-crash".to_string(),
7662 program: PathBuf::from("/unused/genuine-crash"),
7663 args: Vec::new(),
7664 env: Vec::new(),
7665 reserved: false,
7666 reserved_prefixes: Vec::new(),
7667 protocol: ModuleProtocol::Subc,
7668 overlap: Default::default(),
7669 };
7670 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7671
7672 assert!(matches!(
7673 on_child_exit(
7674 &spec,
7675 runtime.restart_policy,
7676 &supervisor.registry,
7677 &snapshot,
7678 &runtime.terminal_ring,
7679 &runtime.spawn_events,
7680 ExitReport {
7681 kind: ExitKind::Crash,
7682 code: Some(1),
7683 signal: None,
7684 at_ms: 1,
7685 },
7686 )
7687 .await,
7688 NextAction::Restart { schedule: _ }
7689 ));
7690 let state = lock_snapshot(&snapshot).unwrap();
7691 assert_eq!(state.lifetime_restarts, 1);
7692 assert_eq!(state.crash_restarts.len(), 1);
7693 }
7694
7695 fn crash_exit_report(at_ms: u64) -> ExitReport {
7696 ExitReport {
7697 kind: ExitKind::Crash,
7698 code: Some(1),
7699 signal: None,
7700 at_ms,
7701 }
7702 }
7703
7704 fn windowed_crash_spec(module_id: &str) -> ModuleSpec {
7705 ModuleSpec {
7706 module_id: module_id.to_string(),
7707 program: PathBuf::from("/unused").join(module_id),
7708 args: Vec::new(),
7709 env: Vec::new(),
7710 reserved: false,
7711 reserved_prefixes: Vec::new(),
7712 protocol: ModuleProtocol::Subc,
7713 overlap: Default::default(),
7714 }
7715 }
7716
7717 #[tokio::test]
7723 async fn three_crashes_inside_the_window_stop_the_module_and_name_the_window() {
7724 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::ERROR);
7725 let supervisor = Supervisor::new(
7726 Arc::new(Registry::default()),
7727 RestartPolicy::new(2, Duration::ZERO),
7728 );
7729 let runtime = supervisor.runtime_config();
7730 let spec = windowed_crash_spec("crash-loop-in-window");
7731 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7732
7733 for attempt in 1..=2 {
7734 assert!(
7735 matches!(
7736 on_child_exit(
7737 &spec,
7738 runtime.restart_policy,
7739 &supervisor.registry,
7740 &snapshot,
7741 &runtime.terminal_ring,
7742 &runtime.spawn_events,
7743 crash_exit_report(attempt),
7744 )
7745 .await,
7746 NextAction::Restart { schedule: _ }
7747 ),
7748 "crash {attempt} is inside the budget and must respawn"
7749 );
7750 }
7751
7752 assert!(matches!(
7753 on_child_exit(
7754 &spec,
7755 runtime.restart_policy,
7756 &supervisor.registry,
7757 &snapshot,
7758 &runtime.terminal_ring,
7759 &runtime.spawn_events,
7760 crash_exit_report(3),
7761 )
7762 .await,
7763 NextAction::Stop { .. }
7764 ));
7765
7766 {
7767 let state = lock_snapshot(&snapshot).unwrap();
7768 assert_eq!(state.state, ModuleState::Failed);
7769 assert_eq!(state.crash_restarts.len(), 2);
7770 assert_eq!(state.lifetime_restarts, 2);
7771 }
7772
7773 let history = runtime
7774 .terminal_ring
7775 .lock()
7776 .expect("terminal ring is not poisoned")
7777 .snapshot();
7778 let last = history
7779 .entries
7780 .last()
7781 .expect("the refused crash is retained");
7782 assert_eq!(last.disposition, TerminalDisposition::Failed);
7783 assert_eq!(
7784 last.disposition_detail.as_deref(),
7785 Some("crash budget exhausted: max_restarts=2 within window_secs=600")
7786 );
7787
7788 let captured = crate::router::test_log::captured_logs(&logs);
7789 assert!(
7790 captured.contains("crash budget exhausted: max_restarts=2 within window_secs=600"),
7791 "the stop must be logged with its window: {captured}"
7792 );
7793 }
7794
7795 #[tokio::test]
7803 async fn a_crash_older_than_the_window_frees_its_slot_for_a_later_crash() {
7804 let supervisor = Supervisor::new(
7805 Arc::new(Registry::default()),
7806 RestartPolicy::new(2, Duration::ZERO),
7807 );
7808 let runtime = supervisor.runtime_config();
7809 let spec = windowed_crash_spec("crash-across-windows");
7810 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7811
7812 for attempt in 1..=2 {
7813 assert!(matches!(
7814 on_child_exit(
7815 &spec,
7816 runtime.restart_policy,
7817 &supervisor.registry,
7818 &snapshot,
7819 &runtime.terminal_ring,
7820 &runtime.spawn_events,
7821 crash_exit_report(attempt),
7822 )
7823 .await,
7824 NextAction::Restart { schedule: _ }
7825 ));
7826 }
7827
7828 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
7831 age_oldest_crash_restart_out_of_window(state, runtime.restart_policy.window);
7832 })
7833 .unwrap();
7834
7835 assert!(
7836 matches!(
7837 on_child_exit(
7838 &spec,
7839 runtime.restart_policy,
7840 &supervisor.registry,
7841 &snapshot,
7842 &runtime.terminal_ring,
7843 &runtime.spawn_events,
7844 crash_exit_report(3),
7845 )
7846 .await,
7847 NextAction::Restart { schedule: _ }
7848 ),
7849 "a crash older than the window must not hold a budget slot"
7850 );
7851
7852 let state = lock_snapshot(&snapshot).unwrap();
7853 assert_eq!(state.state, ModuleState::Restarting);
7854 assert_eq!(
7855 state.crash_restarts.len(),
7856 2,
7857 "the aged instant is dropped and the new one takes its place"
7858 );
7859 assert_eq!(
7860 state.lifetime_restarts, 3,
7861 "the ledger counts every restart, including the ones the window forgot"
7862 );
7863 }
7864
7865 #[tokio::test]
7870 async fn an_operator_restart_clears_the_ring_and_leaves_the_ledger_alone() {
7871 let supervisor = Supervisor::new(
7872 Arc::new(Registry::default()),
7873 RestartPolicy::new(2, Duration::ZERO),
7874 );
7875 let runtime = supervisor.runtime_config();
7876 let spec = windowed_crash_spec("operator-cleared-budget");
7877 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7878
7879 for attempt in 1..=2 {
7880 assert!(matches!(
7881 on_child_exit(
7882 &spec,
7883 runtime.restart_policy,
7884 &supervisor.registry,
7885 &snapshot,
7886 &runtime.terminal_ring,
7887 &runtime.spawn_events,
7888 crash_exit_report(attempt),
7889 )
7890 .await,
7891 NextAction::Restart { schedule: _ }
7892 ));
7893 }
7894
7895 reset_restart_count(&snapshot, &spec.module_id).unwrap();
7896 {
7897 let state = lock_snapshot(&snapshot).unwrap();
7898 assert!(
7899 state.crash_restarts.is_empty(),
7900 "an operator restart returns the full budget"
7901 );
7902 assert_eq!(
7903 state.lifetime_restarts, 2,
7904 "clearing the budget must not unmake the crashes"
7905 );
7906 }
7907
7908 assert!(
7909 matches!(
7910 on_child_exit(
7911 &spec,
7912 runtime.restart_policy,
7913 &supervisor.registry,
7914 &snapshot,
7915 &runtime.terminal_ring,
7916 &runtime.spawn_events,
7917 crash_exit_report(3),
7918 )
7919 .await,
7920 NextAction::Restart { schedule: _ }
7921 ),
7922 "the cleared budget must be spendable again"
7923 );
7924 let state = lock_snapshot(&snapshot).unwrap();
7925 assert_eq!(state.crash_restarts.len(), 1);
7926 assert_eq!(state.lifetime_restarts, 3);
7927 }
7928
7929 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7930 async fn severance_marker_for_a_dead_child_does_not_label_its_successor() {
7931 let severed = ProcessIdentity {
7932 pid: 41,
7933 start_time: 101,
7934 };
7935 let successor = ProcessIdentity {
7936 pid: 41,
7937 start_time: 202,
7938 };
7939 let module = module_with_recovery_snapshot(ModuleState::Running, true, 0);
7940 update_snapshot(&module.inner.snapshot, Some("recovery-snapshot"), |state| {
7941 state.pid = Some(successor.pid);
7942 state.process_start_time = Some(successor.start_time);
7943 })
7944 .unwrap();
7945 assert!(!module.record_deliberate_severance(severed).unwrap());
7946
7947 let exit_report = apply_deliberate_severance_marker(
7948 &module.inner.snapshot,
7949 Some(successor),
7950 ExitReport {
7951 kind: ExitKind::Crash,
7952 code: Some(1),
7953 signal: None,
7954 at_ms: 1,
7955 },
7956 );
7957
7958 assert_eq!(exit_report.kind, ExitKind::Crash);
7959 }
7960
7961 #[tokio::test]
7962 async fn drain_reap_marks_deliberate_severance_and_records_lifetime_without_budget() {
7963 let registry = Registry::default();
7964 let supervisor = Supervisor::new(
7965 Arc::new(Registry::default()),
7966 RestartPolicy::new(3, Duration::ZERO),
7967 );
7968 let runtime = supervisor.runtime_config();
7969 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7970 let spec = ModuleSpec {
7971 module_id: "drain-deliberate-severance".to_string(),
7972 program: fake_aft_stub_path(),
7973 args: Vec::new(),
7974 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7975 reserved: false,
7976 reserved_prefixes: Vec::new(),
7977 protocol: ModuleProtocol::Subc,
7978 overlap: Default::default(),
7979 };
7980 let mut child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
7981 let process = ProcessIdentity {
7982 pid: 41,
7983 start_time: 101,
7984 };
7985 child.process_identity = Some(process);
7986 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
7987 state.pid = Some(process.pid);
7988 state.process_start_time = Some(process.start_time);
7989 })
7990 .unwrap();
7991 record_deliberate_severance(&snapshot, process).unwrap();
7992
7993 drain_child_to_state(
7994 &spec.module_id,
7995 spec.protocol,
7996 ®istry,
7997 &snapshot,
7998 &runtime.terminal_ring,
7999 &runtime.spawn_events,
8000 child,
8001 Duration::from_secs(1),
8002 ModuleState::Stopped,
8003 Some(false),
8004 )
8005 .await
8006 .unwrap();
8007
8008 let state = lock_snapshot(&snapshot).unwrap();
8009 assert_eq!(
8010 state.last_exit.as_ref().map(|exit| exit.kind),
8011 Some(ExitKind::DeliberateSeverance)
8012 );
8013 assert_eq!(state.lifetime_restarts, 1);
8014 assert_eq!(state.crash_restarts.len(), 0);
8015 drop(state);
8016 let history = runtime.terminal_ring.lock().unwrap().snapshot();
8017 assert_eq!(
8018 history.entries[0].exit_kind,
8019 subc_control::TerminalExitKind::DeliberateSeverance
8020 );
8021 }
8022
8023 #[tokio::test]
8024 async fn ordinary_drain_reap_does_not_record_a_lifetime_restart() {
8025 let registry = Registry::default();
8026 let supervisor = Supervisor::new(
8027 Arc::new(Registry::default()),
8028 RestartPolicy::new(3, Duration::ZERO),
8029 );
8030 let runtime = supervisor.runtime_config();
8031 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8032 let spec = ModuleSpec {
8033 module_id: "ordinary-drain".to_string(),
8034 program: fake_aft_stub_path(),
8035 args: Vec::new(),
8036 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
8037 reserved: false,
8038 reserved_prefixes: Vec::new(),
8039 protocol: ModuleProtocol::Subc,
8040 overlap: Default::default(),
8041 };
8042 let child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
8043
8044 drain_child_to_state(
8045 &spec.module_id,
8046 spec.protocol,
8047 ®istry,
8048 &snapshot,
8049 &runtime.terminal_ring,
8050 &runtime.spawn_events,
8051 child,
8052 Duration::from_secs(1),
8053 ModuleState::Stopped,
8054 Some(false),
8055 )
8056 .await
8057 .unwrap();
8058
8059 let state = lock_snapshot(&snapshot).unwrap();
8060 assert_eq!(
8061 state.last_exit.as_ref().map(|exit| exit.kind),
8062 Some(ExitKind::Crash)
8063 );
8064 assert_eq!(state.lifetime_restarts, 0);
8065 assert_eq!(state.crash_restarts.len(), 0);
8066 }
8067
8068 #[test]
8069 fn fatal_connection_teardown_cannot_arm_a_marker_for_a_surviving_process() {
8070 assert!(!include_str!("server.rs")
8076 .contains("router.record_deliberate_connection_severance(ctx.connection_id)"));
8077 }
8078
8079 #[test]
8086 fn drained_after_quiescence_wait_passes_ok_through_and_forces_false_on_err() {
8087 assert!(drained_after_quiescence_wait(&Ok(true)));
8088 assert!(!drained_after_quiescence_wait(&Ok(false)));
8089 assert!(!drained_after_quiescence_wait(&Err(
8090 SuperviseError::StatePoisoned { module_id: None }
8091 )));
8092 }
8093
8094 #[test]
8103 fn wait_error_exit_report_records_a_failed_terminal_with_no_code_or_signal() {
8104 let ring = Arc::new(Mutex::new(TerminalRing::new(
8105 TerminalRingConfig::default(),
8106 0,
8107 )));
8108 record_wait_error_terminal("wait-error", &ring, &super::SpawnEventFeed::default());
8109
8110 let snapshot = ring.lock().unwrap().snapshot();
8111 assert_eq!(snapshot.entries.len(), 1);
8112 let entry = &snapshot.entries[0];
8113 assert_eq!(entry.exit_code, None);
8114 assert_eq!(entry.exit_signal, None);
8115 assert_eq!(entry.disposition, TerminalDisposition::Failed);
8116 }
8117
8118 #[test]
8119 fn wait_error_exit_path_preserves_spawn_event_density() {
8120 let feed = super::SpawnEventFeed::default();
8121 feed.configure_incarnation("wait-error-density".to_string());
8122 feed.emit_spawned("wait-error", 41, 1);
8123 let ring = Arc::new(Mutex::new(TerminalRing::new(
8124 TerminalRingConfig::default(),
8125 0,
8126 )));
8127
8128 record_wait_error_terminal("wait-error", &ring, &feed);
8129 feed.emit_spawned("after-wait-error", 42, 2);
8130
8131 let state = feed.0.lock().unwrap();
8132 let sequences = state
8133 .events
8134 .iter()
8135 .map(|event| event.cursor.seq)
8136 .collect::<Vec<_>>();
8137 assert_eq!(sequences, vec![1, 2, 3]);
8138 assert_eq!(state.events[1].kind, SpawnEventKind::Exited);
8139 assert_eq!(state.events[1].exit_code, None);
8140 assert_eq!(state.events[1].exit_signal, None);
8141 }
8142
8143 #[test]
8147 fn wait_error_exit_report_is_classified_as_a_crash() {
8148 assert_eq!(wait_error_exit_report().kind, ExitKind::Crash);
8149 }
8150}
8151
8152#[cfg(test)]
8153mod health_evidence_tests {
8154 use super::{HealthProbeError, HealthProbeEvidence};
8155 use std::collections::HashSet;
8156
8157 #[test]
8165 fn only_a_dead_lane_is_proof_of_death() {
8166 assert!(HealthProbeError::lane_dead("gone").is_proof_of_death());
8167 assert!(!HealthProbeError::no_answer("timed out").is_proof_of_death());
8171 assert!(!HealthProbeError::bad_answer("garbage").is_proof_of_death());
8172 assert!(!HealthProbeError::misconfigured("no table").is_proof_of_death());
8173 }
8174
8175 #[test]
8181 fn every_evidence_class_has_a_distinct_label() {
8182 let labels = [
8183 HealthProbeError::lane_dead("").label(),
8184 HealthProbeError::no_answer("").label(),
8185 HealthProbeError::bad_answer("").label(),
8186 HealthProbeError::misconfigured("").label(),
8187 ];
8188 let unique: HashSet<_> = labels.iter().collect();
8189 assert_eq!(unique.len(), labels.len(), "labels collided: {labels:?}");
8190 }
8191
8192 #[test]
8198 fn classification_preserves_the_original_message() {
8199 let err = HealthProbeError::no_answer("module did not answer within 5s");
8200 assert_eq!(err.to_string(), "module did not answer within 5s");
8201 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8202 }
8203}
8204
8205#[cfg(test)]
8206mod health_tombstone_tests {
8207 use std::{path::PathBuf, sync::Arc, time::Duration};
8208
8209 use subc_protocol::{
8210 manifest::Concurrency,
8211 session::{HealthStatus, ModuleControlResponse},
8212 };
8213 use tokio::sync::mpsc;
8214
8215 use super::{
8216 probe_module_health, HealthAction, HealthConfig, HealthProbeEvidence, ModuleProtocol,
8217 ModuleSpec, RestartPolicy, Supervisor, SupervisorRuntimeConfig,
8218 };
8219 use crate::{
8220 control::ControlHandler,
8221 forwarding::{ForwardingTable, ModuleControlRpcCompletion, ModuleControlRpcOutcome},
8222 registry::{ConnectionId, Registry},
8223 router::FrameSink,
8224 };
8225
8226 struct ProbeHarness {
8227 spec: ModuleSpec,
8228 runtime: SupervisorRuntimeConfig,
8229 forwarding: Arc<ForwardingTable>,
8230 module_connection: ConnectionId,
8231 module_rx: mpsc::Receiver<crate::router::OutboundFrame>,
8232 handler: ControlHandler,
8233 module: super::SupervisedModule,
8234 }
8235
8236 fn probe_harness() -> ProbeHarness {
8237 let registry = Arc::new(Registry::default());
8238 let forwarding = Arc::new(ForwardingTable::default());
8239 let supervisor_handle = super::SupervisorHandle::new();
8240 let health = HealthConfig {
8241 cadence: Duration::from_secs(30),
8242 deadline: Duration::from_secs(5),
8243 failure_threshold: 3,
8244 on_degraded: HealthAction::Report,
8245 on_failing: HealthAction::Report,
8246 critical: false,
8247 };
8248 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
8249 .with_forwarding(Arc::clone(&forwarding))
8250 .with_handle(supervisor_handle.clone())
8251 .with_health_config(health);
8252 let spec = ModuleSpec {
8253 module_id: "late-health-module".to_string(),
8254 program: PathBuf::from("disabled-module"),
8255 args: Vec::new(),
8256 env: Vec::new(),
8257 reserved: false,
8258 reserved_prefixes: Vec::new(),
8259 protocol: ModuleProtocol::Subc,
8260 overlap: Default::default(),
8261 };
8262 let module = supervisor
8263 .supervise_configured(spec.clone(), false)
8264 .unwrap();
8265 let runtime = supervisor.runtime_config();
8266 let handler = ControlHandler::with_forwarding(registry, Arc::clone(&forwarding))
8267 .with_supervisor(supervisor_handle);
8268 let module_connection = ConnectionId::new(700);
8269 let (module_tx, module_rx) = mpsc::channel(8);
8270 forwarding
8271 .register_module_connection(
8272 module_connection,
8273 spec.module_id.clone(),
8274 subc_protocol::PROTOCOL_VERSION,
8275 Concurrency::ModuleManaged,
8276 FrameSink::new(module_tx),
8277 )
8278 .unwrap();
8279
8280 ProbeHarness {
8281 spec,
8282 runtime,
8283 forwarding,
8284 module_connection,
8285 module_rx,
8286 handler,
8287 module,
8288 }
8289 }
8290
8291 async fn finish_after(
8292 harness: &mut ProbeHarness,
8293 stall: Duration,
8294 ) -> ModuleControlRpcCompletion {
8295 assert!(stall > harness.runtime.health.deadline);
8296 let deadline = harness.runtime.health.deadline;
8297 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8298 let answer = async {
8299 let frame = harness.module_rx.recv().await.expect("health.check frame");
8300 tokio::time::advance(deadline).await;
8301 tokio::task::yield_now().await;
8302 tokio::time::advance(stall - deadline).await;
8303 harness
8304 .forwarding
8305 .complete_module_control_rpc(
8306 harness.module_connection,
8307 frame.header.corr,
8308 Some("health.check"),
8309 ModuleControlRpcOutcome::Response(ModuleControlResponse::HealthCheck {
8310 status: HealthStatus::Ok,
8311 detail: None,
8312 metrics: None,
8313 }),
8314 )
8315 .unwrap()
8316 };
8317 let (probe_result, completion) = tokio::join!(probe, answer);
8318 let err = probe_result.expect_err("probe must miss its deadline");
8319 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8320 completion
8321 }
8322
8323 async fn time_out_without_answer(harness: &mut ProbeHarness) {
8324 let deadline = harness.runtime.health.deadline;
8325 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8326 let exhaust_deadline = async {
8327 let _frame = harness.module_rx.recv().await.expect("health.check frame");
8328 tokio::time::advance(deadline).await;
8329 tokio::task::yield_now().await;
8330 };
8331 let (probe_result, ()) = tokio::join!(probe, exhaust_deadline);
8332 let err = probe_result.expect_err("probe must miss its deadline");
8333 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8334 }
8335
8336 #[tokio::test(start_paused = true)]
8337 async fn late_health_answers_record_start_anchored_latency_for_two_stalls() {
8338 let mut harness = probe_harness();
8339
8340 let first = finish_after(&mut harness, Duration::from_secs(8)).await;
8341 let first_latency = match &first {
8342 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8343 other => panic!("late answer was not retained: {other:?}"),
8344 };
8345 assert!(harness.handler.observe_module_control_completion(first));
8346
8347 let second = finish_after(&mut harness, Duration::from_secs(11)).await;
8348 let second_latency = match &second {
8349 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8350 other => panic!("late answer was not retained: {other:?}"),
8351 };
8352 assert!(harness.handler.observe_module_control_completion(second));
8353
8354 assert_eq!(first_latency, Duration::from_secs(8));
8355 assert_eq!(
8356 second_latency - first_latency,
8357 Duration::from_secs(3),
8358 "latency must grow linearly with the additional stall"
8359 );
8360 let health = harness.module.status().unwrap().health;
8361 assert_eq!(health.late_answer_count, 2);
8362 assert_eq!(health.last_late_answer_latency_ms, Some(11_000));
8363 }
8364
8365 #[tokio::test(start_paused = true)]
8373 async fn late_answer_clears_the_consecutive_failure_streak() {
8374 let mut harness = probe_harness();
8375
8376 time_out_without_answer(&mut harness).await;
8378 harness
8379 .module
8380 .record_health_probe_failure_for_test("[no-answer] test miss")
8381 .unwrap();
8382 assert_eq!(
8383 harness.module.status().unwrap().health.consecutive_failures,
8384 1,
8385 "precondition: the miss must be on the streak before the late answer"
8386 );
8387
8388 let late = finish_after(&mut harness, Duration::from_secs(9)).await;
8390 assert!(matches!(
8391 late,
8392 ModuleControlRpcCompletion::LateHealthAnswer { .. }
8393 ));
8394 assert!(harness.handler.observe_module_control_completion(late));
8395
8396 let health = harness.module.status().unwrap().health;
8397 assert_eq!(
8398 health.consecutive_failures, 0,
8399 "a late answer is an answer: the streak must reset"
8400 );
8401 assert_eq!(health.late_answer_count, 1);
8402 }
8403
8404 #[tokio::test(start_paused = true)]
8405 async fn repeated_serial_probe_cycles_keep_one_tombstone_per_endpoint() {
8406 let mut harness = probe_harness();
8407
8408 for _ in 0..20 {
8409 time_out_without_answer(&mut harness).await;
8410 assert_eq!(
8411 harness.forwarding.health_probe_tombstone_count().unwrap(),
8412 1
8413 );
8414 }
8415 }
8416}
8417
8418#[cfg(test)]
8419mod child_env_tests {
8420 use super::{
8421 apply_child_env, apply_spawn_role, apply_wire_spawn_args, ModuleProtocol, ModuleSpec,
8422 SpawnRole, SupervisorHandle, SPAWN_ROLE_SWAP_CANDIDATE, SUBC_ARG, SUBC_LAUNCH_NONCE_ENV,
8423 SUBC_MODULE_ID_ENV, SUBC_SPAWN_ROLE_ENV,
8424 };
8425 use std::{ffi::OsStr, path::PathBuf};
8426 use tokio::process::Command;
8427
8428 fn spec(env: Vec<(String, String)>) -> ModuleSpec {
8429 ModuleSpec {
8430 module_id: "env-plan".to_string(),
8431 program: PathBuf::from("/nonexistent"),
8432 args: Vec::new(),
8433 env,
8434 reserved: false,
8435 reserved_prefixes: Vec::new(),
8436 protocol: ModuleProtocol::Subc,
8437 overlap: Default::default(),
8438 }
8439 }
8440
8441 #[test]
8455 fn ambient_ck_log_is_removed_and_a_configured_one_survives() {
8456 let mut command = Command::new("/nonexistent");
8457 apply_child_env(&mut command, &spec(Vec::new()));
8458 let removed = command
8459 .as_std()
8460 .get_envs()
8461 .any(|(key, value)| key == OsStr::new("CK_LOG") && value.is_none());
8462 assert!(
8463 removed,
8464 "ambient CK_LOG must be explicitly removed for an unconfigured module"
8465 );
8466
8467 let mut configured = Command::new("/nonexistent");
8468 apply_child_env(
8469 &mut configured,
8470 &spec(vec![("CK_LOG".to_string(), "debug".to_string())]),
8471 );
8472 let effective = configured
8473 .as_std()
8474 .get_envs()
8475 .filter(|(key, _)| *key == OsStr::new("CK_LOG"))
8476 .last()
8477 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()));
8478 assert_eq!(
8479 effective,
8480 Some(Some("debug".to_string())),
8481 "a module's configured CK_LOG must survive the ambient removal"
8482 );
8483 }
8484
8485 #[test]
8494 fn protocol_none_spawn_carries_no_subc_argument_and_no_nonce() {
8495 let connection_file = std::path::Path::new("/run/subc-connection.json");
8496 let handle = SupervisorHandle::new();
8497
8498 let mut none_spec = spec(Vec::new());
8499 none_spec.protocol = ModuleProtocol::None;
8500 let mut none = Command::new("/nonexistent");
8501 apply_wire_spawn_args(&mut none, &none_spec, Some(connection_file), Some(&handle))
8502 .expect("protocol-none spawn args apply");
8503 let none_args: Vec<String> = none
8504 .as_std()
8505 .get_args()
8506 .map(|a| a.to_string_lossy().into_owned())
8507 .collect();
8508 assert!(
8509 !none_args.iter().any(|a| a == SUBC_ARG),
8510 "protocol:none argv must not carry --subc; got {none_args:?}"
8511 );
8512 let none_has_nonce = none
8513 .as_std()
8514 .get_envs()
8515 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some());
8516 assert!(
8517 !none_has_nonce,
8518 "protocol:none spawn must not receive a launch nonce"
8519 );
8520 let none_has_module_id = none
8521 .as_std()
8522 .get_envs()
8523 .any(|(key, value)| key == OsStr::new(SUBC_MODULE_ID_ENV) && value.is_some());
8524 assert!(
8525 none_has_module_id,
8526 "SUBC_MODULE_ID is inert and stays on every path"
8527 );
8528 assert!(
8529 handle.spawn_nonce(&none_spec.module_id).is_none(),
8530 "no nonce record for a process that will never present one"
8531 );
8532
8533 let wire_spec = spec(Vec::new());
8535 let mut wire = Command::new("/nonexistent");
8536 apply_wire_spawn_args(&mut wire, &wire_spec, Some(connection_file), Some(&handle))
8537 .expect("subc-wire spawn args apply");
8538 let wire_args: Vec<String> = wire
8539 .as_std()
8540 .get_args()
8541 .map(|a| a.to_string_lossy().into_owned())
8542 .collect();
8543 assert_eq!(
8544 wire_args,
8545 vec![
8546 SUBC_ARG.to_string(),
8547 connection_file.to_string_lossy().into_owned()
8548 ],
8549 "a subc-wire spawn still carries --subc <path>"
8550 );
8551 assert!(wire
8552 .as_std()
8553 .get_envs()
8554 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some()));
8555 assert!(handle.spawn_nonce(&wire_spec.module_id).is_some());
8556 }
8557
8558 #[test]
8568 fn plain_spawn_removes_the_spawn_role_even_when_the_spec_sets_it() {
8569 let role = |command: &Command| {
8570 command
8571 .as_std()
8572 .get_envs()
8573 .filter(|(key, _)| *key == OsStr::new(SUBC_SPAWN_ROLE_ENV))
8574 .last()
8575 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()))
8576 };
8577 let forged = spec(vec![(
8578 SUBC_SPAWN_ROLE_ENV.to_string(),
8579 SPAWN_ROLE_SWAP_CANDIDATE.to_string(),
8580 )]);
8581
8582 let mut plain = Command::new("/nonexistent");
8583 apply_child_env(&mut plain, &forged);
8584 apply_spawn_role(&mut plain, SpawnRole::Plain);
8585 assert_eq!(
8586 role(&plain),
8587 Some(None),
8588 "a plain spawn must remove SUBC_SPAWN_ROLE, whatever the spec says"
8589 );
8590
8591 let mut candidate = Command::new("/nonexistent");
8592 apply_child_env(&mut candidate, &spec(Vec::new()));
8593 apply_spawn_role(&mut candidate, SpawnRole::SwapCandidate);
8594 assert_eq!(
8595 role(&candidate),
8596 Some(Some(SPAWN_ROLE_SWAP_CANDIDATE.to_string()))
8597 );
8598 }
8599
8600 #[test]
8606 fn daemon_private_capture_keys_are_not_passed_to_the_child() {
8607 let mut command = Command::new("/nonexistent");
8608 apply_child_env(
8609 &mut command,
8610 &spec(vec![
8611 (super::CAPTURE_KEEP_ENV.to_string(), "5".to_string()),
8612 ("KEPT".to_string(), "yes".to_string()),
8613 ]),
8614 );
8615 let keys: Vec<String> = command
8616 .as_std()
8617 .get_envs()
8618 .filter(|(_, value)| value.is_some())
8619 .map(|(key, _)| key.to_string_lossy().into_owned())
8620 .collect();
8621 assert!(keys.contains(&"KEPT".to_string()), "got {keys:?}");
8622 assert!(
8623 !keys.contains(&super::CAPTURE_KEEP_ENV.to_string()),
8624 "daemon-private capture key leaked to the child: {keys:?}"
8625 );
8626 }
8627}
8628
8629#[cfg(test)]
8630mod jitter_tests {
8631 use super::jittered_health_delay;
8632 use std::{collections::HashSet, time::Duration};
8633
8634 const FLEET: [&str; 14] = [
8643 "aft",
8644 "alfonso-core",
8645 "magic-context",
8646 "broca",
8647 "thalamus",
8648 "quota",
8649 "engram",
8650 "plexus",
8651 "cerebellum",
8652 "astrocyte",
8653 "synapse",
8654 "subc-mcp",
8655 "cortexkit-credentials",
8656 "subc-federation",
8657 ];
8658
8659 #[test]
8667 fn probe_delays_disperse_across_the_fleet() {
8668 let cadence = Duration::from_secs(30);
8669 let delays: HashSet<Duration> = FLEET
8670 .iter()
8671 .map(|id| jittered_health_delay(id, 0, cadence))
8672 .collect();
8673 assert_eq!(
8674 delays.len(),
8675 FLEET.len(),
8676 "every supervised module must land on its own probe offset"
8677 );
8678 }
8679
8680 #[test]
8686 fn jitter_only_delays_and_stays_within_one_tenth_of_cadence() {
8687 let cadence = Duration::from_secs(30);
8688 let span = cadence / 10;
8689 for id in FLEET {
8690 for probe_index in 0..8 {
8691 let delay = jittered_health_delay(id, probe_index, cadence);
8692 assert!(
8693 delay >= cadence,
8694 "{id}#{probe_index}: jitter must not shorten the cadence"
8695 );
8696 assert!(
8697 delay < cadence + span,
8698 "{id}#{probe_index}: jitter must stay inside one tenth of the cadence"
8699 );
8700 }
8701 }
8702 }
8703
8704 #[test]
8710 fn a_module_offset_is_stable_across_restarts() {
8711 let cadence = Duration::from_secs(30);
8712 for id in FLEET {
8713 assert_eq!(
8714 jittered_health_delay(id, 0, cadence),
8715 jittered_health_delay(id, 0, cadence),
8716 "{id}: the same module and probe index must produce the same offset"
8717 );
8718 }
8719 }
8720
8721 #[test]
8723 fn zero_cadence_yields_zero_delay() {
8724 assert_eq!(
8725 jittered_health_delay("aft", 0, Duration::ZERO),
8726 Duration::ZERO
8727 );
8728 }
8729}
8730
8731#[cfg(all(test, target_os = "linux"))]
8732mod cgroup_placement_tests {
8733 use super::{
8734 apply_cgroup_placement, remove_module_cgroup, ModuleProtocol, ModuleSpec, SuperviseError,
8735 SupervisedChild,
8736 };
8737 use crate::{
8738 stderr_tail::{StderrRing, StderrTailConfig},
8739 test_support::TestTempDir,
8740 };
8741 use std::{
8742 fs, io,
8743 path::{Path, PathBuf},
8744 sync::{Arc, Mutex},
8745 };
8746 use tokio::process::Command;
8747
8748 #[test]
8749 fn failed_parent_cgroup_open_is_a_cgroup_supervision_error() {
8750 let path = Path::new("/definitely-missing-subc-cgroup");
8751 let mut command = Command::new("true");
8752 let error = apply_cgroup_placement(
8753 &mut command,
8754 &ModuleSpec {
8755 module_id: "broken-cgroup".to_string(),
8756 program: PathBuf::from("true"),
8757 args: Vec::new(),
8758 env: Vec::new(),
8759 reserved: false,
8760 reserved_prefixes: Vec::new(),
8761 protocol: ModuleProtocol::Subc,
8762 overlap: Default::default(),
8763 },
8764 path,
8765 )
8766 .expect_err("a parent cgroup open failure must reject the supervised spawn");
8767 let reason = error.to_string();
8768
8769 assert!(
8770 matches!(error, SuperviseError::Cgroup { .. }),
8771 "parent cgroup open must be reported as a cgroup supervision error: {reason}"
8772 );
8773 assert!(
8774 reason.contains("/definitely-missing-subc-cgroup/cgroup.procs"),
8775 "parent cgroup open failure must name cgroup.procs: {reason}"
8776 );
8777 }
8778
8779 #[tokio::test]
8780 async fn reaping_a_child_removes_its_empty_module_cgroup() {
8781 let root = TestTempDir::new("supervisor-reap-cgroup");
8782 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
8783 let placement = subc_cgroup::prepare_at(&root)
8784 .expect("prepare scratch cgroup root")
8785 .expect("scratch root has a cgroup.procs marker");
8786 let module_id = "reaped-module";
8787 let module = placement
8788 .module_path(module_id)
8789 .expect("create scratch module cgroup");
8790 let child = Command::new("true")
8791 .spawn()
8792 .expect("spawn short-lived child");
8793 let pid = child.id().expect("spawned child has pid");
8794 let mut child = SupervisedChild {
8795 child,
8796 module_id: module_id.to_string(),
8797 cgroup_placement: Some(placement),
8798 stdout_pump: None,
8799 stderr_pump: None,
8800 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
8801 spawned_at_ms: 0,
8802 spawned_from: PathBuf::from("true"),
8803 spawned_file_identity: None,
8804 process_start_time: None,
8805 process_identity: None,
8806 pid,
8807 roster_guard: None,
8808 };
8809
8810 child.wait().await.expect("reap short-lived child");
8811
8812 assert!(
8813 !module.exists(),
8814 "reaping the supervised child must remove its empty cgroup"
8815 );
8816 }
8817
8818 #[test]
8819 fn non_empty_cgroup_removal_is_reported_without_blocking_teardown() {
8820 let root = TestTempDir::new("supervisor-non-empty-cgroup");
8821 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
8822 let placement = subc_cgroup::prepare_at(&root)
8823 .expect("prepare scratch cgroup root")
8824 .expect("scratch root has a cgroup.procs marker");
8825 let module = placement
8826 .module_path("surviving-module")
8827 .expect("create scratch module cgroup");
8828 fs::write(module.join("surviving-process"), b"still present")
8829 .expect("make scratch cgroup non-empty");
8830 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::WARN);
8831
8832 remove_module_cgroup(&placement, "surviving-module");
8833
8834 let logs = crate::router::test_log::captured_logs(&logs);
8835 assert!(
8836 module.exists(),
8837 "failed removal must leave the cgroup intact"
8838 );
8839 assert!(
8840 logs.contains("could not remove module cgroup after process exit; continuing teardown")
8841 && logs.contains("surviving-module"),
8842 "best-effort removal must report the failure without returning it: {logs}"
8843 );
8844 }
8845
8846 #[test]
8847 fn cgroup_pre_exec_spawn_failure_names_the_cgroup_path() {
8848 let cgroup_path = PathBuf::from("/sys/fs/cgroup/subc-modules/broken-module");
8849 let reason = SuperviseError::Spawn {
8850 program: PathBuf::from("/bin/true"),
8851 source: io::Error::from_raw_os_error(13),
8852 cgroup_path: Some(cgroup_path.clone()),
8853 }
8854 .to_string();
8855
8856 assert!(
8857 reason.contains(&cgroup_path.display().to_string()),
8858 "a pre_exec spawn failure must name the cgroup path: {reason}"
8859 );
8860 }
8861}
8862
8863#[cfg(test)]
8864mod spawn_subscriber_lag_tests {
8865 use super::*;
8866
8867 #[tokio::test]
8872 async fn lagged_spawn_subscriber_receives_a_terminal_lagged_error_after_its_queued_frames() {
8873 let feed = SpawnEventFeed::default();
8874 feed.configure_incarnation("lag-incarnation".to_string());
8875 let (tx, mut rx) = mpsc::channel(1);
8878 feed.subscribe(ConnectionId::new(1), 7, 1, None, FrameSink::new(tx))
8879 .expect("subscribe");
8880 let emitted = SPAWN_SUBSCRIBER_BUFFER + 16;
8881 for index in 0..emitted {
8882 feed.emit_spawned(&format!("lag-module-{index}"), 1000, 0);
8883 tokio::task::yield_now().await;
8886 }
8887 assert_eq!(
8888 feed.subscriber_count(),
8889 0,
8890 "the lagged subscriber must be removed"
8891 );
8892
8893 let mut data = Vec::new();
8894 let mut last = None;
8895 loop {
8896 let next = tokio::time::timeout(Duration::from_secs(5), rx.recv())
8897 .await
8898 .expect("the forwarder must finish once the subscriber is dropped");
8899 let Some(outbound) = next else { break };
8900 let frame = outbound.frame;
8901 if frame.header.ty == FrameType::StreamData {
8902 assert!(last.is_none(), "no data may follow the terminal frame");
8903 let event: SpawnEvent = serde_json::from_slice(&frame.body).unwrap();
8904 data.push(event.cursor.seq);
8905 } else {
8906 assert!(last.is_none(), "exactly one terminal frame");
8907 last = Some(frame);
8908 }
8909 }
8910 assert!(!data.is_empty(), "queued frames drain before the terminal");
8911 for pair in data.windows(2) {
8912 assert_eq!(
8913 pair[1],
8914 pair[0] + 1,
8915 "queued frames arrive dense and in order"
8916 );
8917 }
8918 let terminal = last.expect("a lagged subscriber must receive a terminal frame");
8919 assert_eq!(terminal.header.ty, FrameType::Error);
8920 assert_eq!(terminal.header.corr, 7);
8921 let body: subc_protocol::ErrorBody = serde_json::from_slice(&terminal.body).unwrap();
8922 assert_eq!(body.code, SPAWN_SUBSCRIBER_LAGGED_CODE);
8923 let detail = body.detail.expect("lagged error carries detail");
8924 assert_eq!(
8925 detail["first_undelivered_cursor"]["seq"],
8926 data.last().unwrap() + 1,
8927 "the named cursor is the first event the subscriber did not receive"
8928 );
8929 assert_eq!(
8930 detail["first_undelivered_cursor"]["daemon_incarnation"],
8931 "lag-incarnation"
8932 );
8933 }
8934}
8935
8936#[cfg(test)]
8937mod terminal_history_read_concurrency_tests {
8938 use super::*;
8939 use crate::{terminal_journal::read_pause, test_support::TestTempDir};
8940 use std::sync::mpsc as std_mpsc;
8941
8942 fn journaled_ring(
8943 journal: &Arc<crate::terminal_journal::TerminalJournal>,
8944 ) -> Arc<Mutex<TerminalRing>> {
8945 Arc::new(Mutex::new(
8946 TerminalRing::new(TerminalRingConfig::default(), 1)
8947 .with_journal(Some(Arc::clone(journal))),
8948 ))
8949 }
8950
8951 fn crash(at_ms: u64) -> ExitReport {
8952 ExitReport {
8953 kind: ExitKind::Crash,
8954 code: Some(1),
8955 signal: None,
8956 at_ms,
8957 }
8958 }
8959
8960 fn record_within(
8963 module_id: &'static str,
8964 ring: &Arc<Mutex<TerminalRing>>,
8965 at_ms: u64,
8966 bound: Duration,
8967 ) -> bool {
8968 let ring = Arc::clone(ring);
8969 let (done, done_rx) = std_mpsc::channel();
8970 std::thread::spawn(move || {
8971 record_terminal(
8972 module_id,
8973 &ring,
8974 &SpawnEventFeed::default(),
8975 &crash(at_ms),
8976 TerminalDisposition::Restarting,
8977 );
8978 let _ = done.send(());
8979 });
8980 done_rx.recv_timeout(bound).is_ok()
8981 }
8982
8983 #[test]
8988 fn exits_recorded_during_a_paused_history_read_are_not_blocked_or_half_merged() {
8989 let dir = TestTempDir::new("terminal-history-concurrent-read");
8990 let path = dir.join("terminals.jsonl");
8991 let journal = Arc::new(crate::terminal_journal::TerminalJournal::open(
8992 path.clone(),
8993 "daemon".into(),
8994 ));
8995 let reader_ring = journaled_ring(&journal);
8996 let other_ring = journaled_ring(&journal);
8997 assert!(record_within(
8998 "reader-module",
8999 &reader_ring,
9000 10,
9001 Duration::from_secs(5)
9002 ));
9003
9004 let (started, release) = read_pause::install(&path);
9005 let reading = {
9006 let ring = Arc::clone(&reader_ring);
9007 std::thread::spawn(move || durable_terminal_history_of(&ring, "reader-module"))
9008 };
9009 started
9010 .recv_timeout(Duration::from_secs(5))
9011 .expect("the history read reached its pause");
9012
9013 let bound = Duration::from_secs(1);
9014 assert!(
9015 record_within("other-module", &other_ring, 20, bound),
9016 "another module's exit waited on a history read (journal writer held)"
9017 );
9018 assert!(
9019 record_within("reader-module", &reader_ring, 30, bound),
9020 "the read module's own exit waited on its history read (ring held)"
9021 );
9022
9023 drop(release);
9024 let paused = reading.join().unwrap();
9025 assert_eq!(
9026 paused.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9027 vec![10],
9028 "an exit recorded after the read began lands in neither half of it"
9029 );
9030 assert_eq!(paused.journal_skipped_lines, 0);
9031 assert_eq!(paused.journal_read_errors, 0);
9032
9033 let after = durable_terminal_history_of(&reader_ring, "reader-module");
9034 assert_eq!(
9035 after.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9036 vec![10, 30],
9037 "the next read merges ring and journal with no duplicate"
9038 );
9039 assert_eq!(after.journal_skipped_lines, 0);
9040 }
9041}