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_terminal_journal(mut self, path: PathBuf, daemon_incarnation: String) -> Self {
2074 self.spawn_events
2078 .configure_incarnation(daemon_incarnation.clone());
2079 self.terminal_journal = Some(Arc::new(crate::terminal_journal::TerminalJournal::open(
2080 path,
2081 daemon_incarnation,
2082 )));
2083 self
2084 }
2085
2086 pub fn with_forwarding(mut self, forwarding: Arc<ForwardingTable>) -> Self {
2087 self.forwarding = Some(forwarding);
2088 self
2089 }
2090
2091 pub fn with_handle(mut self, supervisor_handle: SupervisorHandle) -> Self {
2092 self.spawn_events = supervisor_handle.spawn_events.clone();
2093 self.supervisor_handle = Some(supervisor_handle);
2094 self
2095 }
2096
2097 pub fn with_health_config(mut self, health: HealthConfig) -> Self {
2098 self.health = health;
2099 self
2100 }
2101
2102 #[cfg(target_os = "linux")]
2103 pub fn with_cgroup_placement(
2104 mut self,
2105 cgroup_placement: Option<subc_cgroup::Placement>,
2106 ) -> Self {
2107 self.cgroup_placement = cgroup_placement;
2108 self
2109 }
2110
2111 pub fn spawn(&self, spec: ModuleSpec) -> Result<SupervisedModule, SuperviseError> {
2117 validate_spec(&spec)?;
2118
2119 let runtime = self.runtime_config();
2120 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2121 let child = spawn_child(
2122 &spec,
2123 runtime.connection_file_path.as_deref(),
2124 self.supervisor_handle.as_ref(),
2125 &runtime.stderr_ring,
2126 runtime.capture_logs_dir.as_deref(),
2127 &runtime.child_roster,
2128 #[cfg(target_os = "linux")]
2129 runtime.cgroup_placement.as_ref(),
2130 )?;
2131 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2132 self.process_liveness
2133 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2134
2135 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2136 }
2137
2138 pub fn supervise_configured(
2144 &self,
2145 spec: ModuleSpec,
2146 enabled: bool,
2147 ) -> Result<SupervisedModule, SuperviseError> {
2148 validate_spec(&spec)?;
2149
2150 let runtime = self.runtime_config();
2151 if !enabled {
2152 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2153 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2154 }
2155
2156 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2157 match spawn_child(
2158 &spec,
2159 runtime.connection_file_path.as_deref(),
2160 self.supervisor_handle.as_ref(),
2161 &runtime.stderr_ring,
2162 runtime.capture_logs_dir.as_deref(),
2163 &runtime.child_roster,
2164 #[cfg(target_os = "linux")]
2165 runtime.cgroup_placement.as_ref(),
2166 ) {
2167 Ok(child) => {
2168 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2169 self.process_liveness
2170 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2171 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2172 }
2173 Err(err) => {
2174 error!(
2175 module_id = %spec.module_id,
2176 program = %spec.program.display(),
2177 error = %err,
2178 "configured module failed to spawn; marking failed and continuing"
2179 );
2180 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2181 Ok(self.supervised_module(spec, runtime, snapshot, None))
2182 }
2183 }
2184 }
2185
2186 pub fn supervise_configured_with_health(
2192 &self,
2193 spec: ModuleSpec,
2194 enabled: bool,
2195 health: HealthConfig,
2196 drain_timeout_ms: Option<u64>,
2197 restart_policy: RestartPolicy,
2198 ) -> Result<SupervisedModule, SuperviseError> {
2199 validate_spec(&spec)?;
2200
2201 let mut runtime = self.runtime_config();
2202 runtime.health = health;
2203 runtime.restart_policy = restart_policy;
2204 if let Some(ms) = drain_timeout_ms {
2205 runtime.drain_timeout = Duration::from_millis(ms);
2206 *runtime
2207 .effective_drain_timeout
2208 .lock()
2209 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
2210 }
2211 if !enabled {
2212 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::disabled()));
2213 return Ok(self.supervised_module(spec, runtime, snapshot, None));
2214 }
2215
2216 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
2217 match spawn_child(
2218 &spec,
2219 runtime.connection_file_path.as_deref(),
2220 self.supervisor_handle.as_ref(),
2221 &runtime.stderr_ring,
2222 runtime.capture_logs_dir.as_deref(),
2223 &runtime.child_roster,
2224 #[cfg(target_os = "linux")]
2225 runtime.cgroup_placement.as_ref(),
2226 ) {
2227 Ok(child) => {
2228 set_running(&snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
2229 self.process_liveness
2230 .track(spec.module_id.clone(), Arc::clone(&snapshot));
2231 Ok(self.supervised_module(spec, runtime, snapshot, Some(child)))
2232 }
2233 Err(err) => {
2234 if health.critical {
2235 error!(
2236 module_id = %spec.module_id,
2237 program = %spec.program.display(),
2238 error = %err,
2239 "critical configured module failed to spawn; marking failed and alerting"
2240 );
2241 } else {
2242 error!(
2243 module_id = %spec.module_id,
2244 program = %spec.program.display(),
2245 error = %err,
2246 "configured module failed to spawn; marking failed and continuing"
2247 );
2248 }
2249 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::failed()));
2250 Ok(self.supervised_module(spec, runtime, snapshot, None))
2251 }
2252 }
2253 }
2254
2255 fn runtime_config(&self) -> SupervisorRuntimeConfig {
2256 let effective_drain_timeout = Arc::new(Mutex::new(self.drain_timeout));
2257 SupervisorRuntimeConfig {
2258 restart_policy: self.restart_policy,
2259 drain_timeout: self.drain_timeout,
2260 child_roster: self
2263 .child_roster
2264 .for_module(Arc::clone(&effective_drain_timeout)),
2265 effective_drain_timeout,
2266 default_drain_timeout: self.drain_timeout,
2267 health: self.health,
2268 connection_file_path: self.connection_file_path.clone(),
2269 capture_logs_dir: self.capture_logs_dir.clone(),
2270 forwarding: self.forwarding.clone(),
2271 supervisor_handle: self.supervisor_handle.clone(),
2272 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
2273 terminal_ring: Arc::new(Mutex::new(
2274 TerminalRing::new(
2275 TerminalRingConfig::default(),
2276 self.daemon_start_clock.started_at_ms(),
2277 )
2278 .with_start_clock(self.daemon_start_clock)
2279 .with_journal(self.terminal_journal.clone()),
2280 )),
2281 spawn_events: self.spawn_events.clone(),
2282 #[cfg(target_os = "linux")]
2283 cgroup_placement: self.cgroup_placement.clone(),
2284 #[cfg(test)]
2285 test_seed_stale_facts_before_enable_spawn: false,
2286 }
2287 }
2288
2289 fn supervised_module(
2290 &self,
2291 spec: ModuleSpec,
2292 runtime: SupervisorRuntimeConfig,
2293 snapshot: SharedSnapshot,
2294 child: Option<SupervisedChild>,
2295 ) -> SupervisedModule {
2296 let configuration = Arc::new(Mutex::new(SupervisedConfiguration {
2297 spec: spec.clone(),
2298 health: runtime.health,
2299 }));
2300 let stderr_ring = Arc::clone(&runtime.stderr_ring);
2301 let terminal_ring = Arc::clone(&runtime.terminal_ring);
2302 let restart_policy = runtime.restart_policy;
2306 let effective_drain_timeout = Arc::clone(&runtime.effective_drain_timeout);
2307 let (tx, rx) = mpsc::channel(4);
2308 let monitor = tokio::spawn(supervise_loop(
2309 spec.clone(),
2310 runtime,
2311 Arc::clone(&self.registry),
2312 Arc::clone(&self.process_liveness),
2313 Arc::clone(&snapshot),
2314 child,
2315 rx,
2316 ));
2317
2318 let module_id = spec.module_id.clone();
2319 let module = SupervisedModule {
2320 inner: Arc::new(SupervisedModuleInner {
2321 module_id: module_id.clone(),
2322 registry: Arc::clone(&self.registry),
2323 snapshot,
2324 configuration,
2325 stderr_ring,
2326 terminal_ring,
2327 commands: tx,
2328 monitor: Mutex::new(Some(monitor)),
2329 restart_policy,
2330 effective_drain_timeout,
2331 provenance_probe: self.provenance_probe.clone(),
2332 }),
2333 };
2334 if let Some(supervisor_handle) = &self.supervisor_handle {
2335 supervisor_handle.apply_identity_configuration(&spec);
2336 supervisor_handle.insert(module.clone());
2337 }
2338 module
2339 }
2340}
2341
2342impl Default for Supervisor {
2343 fn default() -> Self {
2344 Self::new(Arc::new(Registry::default()), RestartPolicy::default())
2345 }
2346}
2347
2348#[derive(Clone)]
2350pub struct SupervisedModule {
2351 inner: Arc<SupervisedModuleInner>,
2352}
2353
2354struct SupervisedModuleInner {
2355 module_id: String,
2356 registry: Arc<Registry>,
2357 snapshot: SharedSnapshot,
2358 configuration: Arc<Mutex<SupervisedConfiguration>>,
2359 stderr_ring: Arc<Mutex<StderrRing>>,
2360 terminal_ring: Arc<Mutex<TerminalRing>>,
2361 commands: mpsc::Sender<SupervisorCommand>,
2362 monitor: Mutex<Option<JoinHandle<()>>>,
2363 restart_policy: RestartPolicy,
2367 effective_drain_timeout: Arc<Mutex<Duration>>,
2368 provenance_probe: ExecutableIdentityProbe,
2369}
2370
2371impl fmt::Debug for SupervisedModule {
2372 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2373 f.debug_struct("SupervisedModule")
2374 .field("module_id", &self.inner.module_id)
2375 .field("status", &self.status())
2376 .finish_non_exhaustive()
2377 }
2378}
2379
2380impl SupervisedModule {
2381 pub fn module_id(&self) -> &str {
2382 &self.inner.module_id
2383 }
2384
2385 #[cfg(test)]
2389 pub(crate) fn record_health_probe_failure_for_test(
2390 &self,
2391 detail: &str,
2392 ) -> Result<(), SuperviseError> {
2393 update_snapshot(&self.inner.snapshot, Some(&self.inner.module_id), |state| {
2394 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
2395 state.health.detail = Some(detail.to_string());
2396 })
2397 }
2398
2399 pub fn state(&self) -> Result<ModuleState, SuperviseError> {
2400 Ok(lock_snapshot(&self.inner.snapshot)?.state)
2401 }
2402
2403 pub fn stderr_tail(
2410 &self,
2411 max_lines: Option<usize>,
2412 max_bytes: Option<usize>,
2413 ) -> StderrTailSnapshot {
2414 self.inner
2415 .stderr_ring
2416 .lock()
2417 .unwrap_or_else(|poisoned| poisoned.into_inner())
2418 .snapshot(max_lines, max_bytes)
2419 }
2420
2421 pub fn terminal_history(&self) -> TerminalHistorySnapshot {
2426 self.inner
2427 .terminal_ring
2428 .lock()
2429 .unwrap_or_else(|poisoned| poisoned.into_inner())
2430 .snapshot()
2431 }
2432
2433 pub fn durable_terminal_history(&self) -> subc_control::TerminalHistory {
2438 durable_terminal_history_of(&self.inner.terminal_ring, &self.inner.module_id)
2439 }
2440
2441 pub(crate) async fn read_durable_terminal_history(
2446 &self,
2447 ) -> Result<subc_control::TerminalHistory, tokio::task::JoinError> {
2448 let terminal_ring = Arc::clone(&self.inner.terminal_ring);
2449 let module_id = self.inner.module_id.clone();
2450 tokio::task::spawn_blocking(move || durable_terminal_history_of(&terminal_ring, &module_id))
2451 .await
2452 }
2453
2454 pub fn status(&self) -> Result<ModuleStatus, SuperviseError> {
2455 self.status_with_snapshot_lock(&self.inner.snapshot, None)
2456 }
2457
2458 pub(crate) fn record_deliberate_severance(
2459 &self,
2460 identity: ProcessIdentity,
2461 ) -> Result<bool, SuperviseError> {
2462 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2463 if snapshot.pid != Some(identity.pid)
2464 || snapshot.process_start_time != Some(identity.start_time)
2465 {
2466 return Ok(false);
2467 }
2468 snapshot.deliberate_severance = Some(identity);
2469 Ok(true)
2470 }
2471
2472 pub(crate) fn status_for_control(
2477 &self,
2478 caller: &'static str,
2479 ) -> Result<ModuleStatus, SuperviseError> {
2480 self.status_with_snapshot_lock(&self.inner.snapshot, Some(caller))
2481 }
2482
2483 fn status_with_snapshot_lock(
2484 &self,
2485 snapshot: &SharedSnapshot,
2486 caller: Option<&'static str>,
2487 ) -> Result<ModuleStatus, SuperviseError> {
2488 let mut guard = match caller {
2489 Some(caller) => lock_snapshot_for_control(snapshot, &self.inner.module_id, caller)?,
2490 None => lock_snapshot(snapshot)?,
2491 };
2492 let restart_count =
2495 guard.crash_restarts_in_window(self.inner.restart_policy.window, Instant::now());
2496 let snapshot = guard.clone();
2497 drop(guard);
2498 let drain_timeout = *self.inner.effective_drain_timeout.lock().map_err(|_| {
2499 SuperviseError::StatePoisoned {
2500 module_id: Some(self.inner.module_id.clone()),
2501 }
2502 })?;
2503 let registration_active = self
2504 .inner
2505 .registry
2506 .get_module(&self.inner.module_id)
2507 .map_err(SuperviseError::Registry)?
2508 .is_some();
2509 let protocol = self.declared_protocol()?;
2510 let running_process =
2511 snapshot.enabled && snapshot.state == ModuleState::Running && snapshot.process_alive;
2512 let live = match protocol {
2518 ModuleProtocol::Subc => running_process && registration_active,
2519 ModuleProtocol::None => running_process,
2520 };
2521
2522 Ok(ModuleStatus {
2523 module_id: self.inner.module_id.clone(),
2524 state: snapshot.state,
2525 enabled: snapshot.enabled,
2526 process_alive: snapshot.process_alive,
2527 registration_active,
2528 protocol,
2529 live,
2530 restart_count,
2531 lifetime_restarts: snapshot.lifetime_restarts,
2532 spawn_generation: snapshot.spawn_generation,
2533 max_restarts: self.inner.restart_policy.max_restarts,
2534 restart_window: self.inner.restart_policy.window,
2535 drain_timeout,
2536 restart_backoff: self.inner.restart_policy.backoff,
2537 restart_max_backoff: self.inner.restart_policy.max_backoff,
2538 pid: snapshot.pid,
2539 spawned_at_ms: snapshot.spawned_at_ms,
2540 spawned_from: snapshot.spawned_from,
2541 process_start_time: snapshot.process_start_time,
2542 last_exit: snapshot.last_exit,
2543 health: snapshot.health,
2544 })
2545 }
2546
2547 #[cfg(test)]
2548 pub(crate) fn hold_snapshot_for_test(
2549 &self,
2550 acquired: std::sync::mpsc::Sender<()>,
2551 hold: Duration,
2552 ) -> std::thread::JoinHandle<()> {
2553 let snapshot = Arc::clone(&self.inner.snapshot);
2554 std::thread::spawn(move || {
2555 let _guard = snapshot.lock().expect("test snapshot lock is not poisoned");
2556 acquired
2557 .send(())
2558 .expect("test receiver waits for snapshot lock");
2559 std::thread::sleep(hold);
2560 })
2561 }
2562
2563 pub(crate) async fn running_image_agreement(&self) -> subc_control::RunningImageAgreement {
2564 let snapshot = match lock_snapshot(&self.inner.snapshot) {
2565 Ok(snapshot) => snapshot.clone(),
2566 Err(_) => {
2567 return subc_control::RunningImageAgreement::Unavailable {
2568 reason: subc_control::RunningImageUnavailableReason::NotRunning,
2569 };
2570 }
2571 };
2572 self.inner
2573 .provenance_probe
2574 .observe(
2575 snapshot.pid,
2576 snapshot.spawned_from.as_deref(),
2577 snapshot.spawned_file_identity,
2578 snapshot.process_start_time,
2579 )
2580 .await
2581 }
2582
2583 pub(crate) fn will_recover_after_connection_loss(&self) -> Result<bool, SuperviseError> {
2584 let mut snapshot = lock_snapshot(&self.inner.snapshot)?;
2585 Ok(match snapshot.state {
2586 ModuleState::Restarting => true,
2587 ModuleState::Failed | ModuleState::Disabled => false,
2588 _ => daemon_will_restart(&mut snapshot, &self.inner.restart_policy, Instant::now()),
2589 })
2590 }
2591
2592 #[cfg(test)]
2593 pub(crate) fn is_warming(&self) -> Result<bool, SuperviseError> {
2594 self.is_warming_with_snapshot_lock(None)
2595 }
2596
2597 pub(crate) fn is_warming_for_control(
2598 &self,
2599 caller: &'static str,
2600 ) -> Result<bool, SuperviseError> {
2601 self.is_warming_with_snapshot_lock(Some(caller))
2602 }
2603
2604 fn is_warming_with_snapshot_lock(
2605 &self,
2606 caller: Option<&'static str>,
2607 ) -> Result<bool, SuperviseError> {
2608 let snapshot = match caller {
2609 Some(caller) => {
2610 lock_snapshot_for_control(&self.inner.snapshot, &self.inner.module_id, caller)?
2611 }
2612 None => lock_snapshot(&self.inner.snapshot)?,
2613 }
2614 .clone();
2615 Ok(matches!(
2616 snapshot.state,
2617 ModuleState::Starting | ModuleState::Running | ModuleState::Restarting
2618 ))
2619 }
2620
2621 pub async fn drain(&self) -> Result<(), SuperviseError> {
2623 self.stop().await
2624 }
2625
2626 pub(crate) async fn retire(&self) -> Result<(), SuperviseError> {
2627 match self.state()? {
2628 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2629 ModuleState::Starting
2630 | ModuleState::Running
2631 | ModuleState::Unresponsive
2632 | ModuleState::Restarting
2633 | ModuleState::Draining
2634 | ModuleState::Disabled => {}
2635 }
2636
2637 let (reply_tx, reply_rx) = oneshot::channel();
2638 self.inner
2639 .commands
2640 .send(SupervisorCommand::Retire { reply: reply_tx })
2641 .await
2642 .map_err(|_| SuperviseError::CommandClosed {
2643 module_id: self.inner.module_id.clone(),
2644 })?;
2645 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2646 module_id: self.inner.module_id.clone(),
2647 })?
2648 }
2649
2650 pub async fn stop(&self) -> Result<(), SuperviseError> {
2651 match self.state()? {
2652 ModuleState::Stopped | ModuleState::Failed => return Ok(()),
2653 ModuleState::Starting
2654 | ModuleState::Running
2655 | ModuleState::Unresponsive
2656 | ModuleState::Restarting
2657 | ModuleState::Draining
2658 | ModuleState::Disabled => {}
2659 }
2660
2661 let (reply_tx, reply_rx) = oneshot::channel();
2662 self.inner
2663 .commands
2664 .send(SupervisorCommand::Drain { reply: reply_tx })
2665 .await
2666 .map_err(|_| SuperviseError::CommandClosed {
2667 module_id: self.inner.module_id.clone(),
2668 })?;
2669 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2670 module_id: self.inner.module_id.clone(),
2671 })?
2672 }
2673
2674 pub async fn restart(&self, drain_timeout_ms: Option<u64>) -> Result<(), SuperviseError> {
2675 let (reply_tx, reply_rx) = oneshot::channel();
2676 self.inner
2677 .commands
2678 .send(SupervisorCommand::Restart {
2679 drain_timeout_ms,
2680 reply: reply_tx,
2681 })
2682 .await
2683 .map_err(|_| SuperviseError::CommandClosed {
2684 module_id: self.inner.module_id.clone(),
2685 })?;
2686 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2687 module_id: self.inner.module_id.clone(),
2688 })?
2689 }
2690
2691 pub async fn swap(&self, ready_timeout: Option<Duration>) -> Result<(), SuperviseError> {
2696 let (reply_tx, reply_rx) = oneshot::channel();
2697 self.inner
2698 .commands
2699 .send(SupervisorCommand::Swap {
2700 ready_timeout,
2701 reply: reply_tx,
2702 })
2703 .await
2704 .map_err(|_| SuperviseError::CommandClosed {
2705 module_id: self.inner.module_id.clone(),
2706 })?;
2707 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2708 module_id: self.inner.module_id.clone(),
2709 })?
2710 }
2711
2712 pub async fn reload(&self) -> Result<(), SuperviseError> {
2713 let (reply_tx, reply_rx) = oneshot::channel();
2714 self.inner
2715 .commands
2716 .send(SupervisorCommand::Reload { reply: reply_tx })
2717 .await
2718 .map_err(|_| SuperviseError::CommandClosed {
2719 module_id: self.inner.module_id.clone(),
2720 })?;
2721 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2722 module_id: self.inner.module_id.clone(),
2723 })?
2724 }
2725
2726 pub async fn set_enabled(&self, enabled: bool) -> Result<bool, SuperviseError> {
2727 let (reply_tx, reply_rx) = oneshot::channel();
2728 self.inner
2729 .commands
2730 .send(SupervisorCommand::SetEnabled {
2731 enabled,
2732 reply: reply_tx,
2733 })
2734 .await
2735 .map_err(|_| SuperviseError::CommandClosed {
2736 module_id: self.inner.module_id.clone(),
2737 })?;
2738 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2739 module_id: self.inner.module_id.clone(),
2740 })?
2741 }
2742
2743 pub(crate) fn declared_protocol(&self) -> Result<ModuleProtocol, SuperviseError> {
2748 Ok(self
2749 .inner
2750 .configuration
2751 .lock()
2752 .map_err(|_| SuperviseError::StatePoisoned {
2753 module_id: Some(self.inner.module_id.clone()),
2754 })?
2755 .spec
2756 .protocol)
2757 }
2758
2759 pub(crate) fn configuration(&self) -> Result<(ModuleSpec, HealthConfig), SuperviseError> {
2760 let configuration =
2761 self.inner
2762 .configuration
2763 .lock()
2764 .map_err(|_| SuperviseError::StatePoisoned {
2765 module_id: Some(self.inner.module_id.clone()),
2766 })?;
2767 Ok((configuration.spec.clone(), configuration.health))
2768 }
2769
2770 #[cfg(any(test, feature = "test-support"))]
2774 pub async fn update_spec_for_test(&self, spec: ModuleSpec) -> Result<(), SuperviseError> {
2775 let (_, health) = self.configuration()?;
2776 let drain_timeout_ms = u64::try_from(
2777 self.inner
2778 .effective_drain_timeout
2779 .lock()
2780 .unwrap_or_else(|poisoned| poisoned.into_inner())
2781 .as_millis(),
2782 )
2783 .ok();
2784 self.update_configuration(spec, health, drain_timeout_ms)
2785 .await
2786 }
2787
2788 pub(crate) async fn update_configuration(
2789 &self,
2790 spec: ModuleSpec,
2791 health: HealthConfig,
2792 drain_timeout_ms: Option<u64>,
2793 ) -> Result<(), SuperviseError> {
2794 if spec.module_id != self.inner.module_id {
2795 return Err(SuperviseError::InvalidSpec {
2796 reason: "a supervised module's module_id cannot be changed".to_string(),
2797 });
2798 }
2799 validate_spec(&spec)?;
2800 let (reply_tx, reply_rx) = oneshot::channel();
2801 self.inner
2802 .commands
2803 .send(SupervisorCommand::UpdateConfiguration {
2804 spec: spec.clone(),
2805 health,
2806 drain_timeout_ms,
2807 reply: reply_tx,
2808 })
2809 .await
2810 .map_err(|_| SuperviseError::CommandClosed {
2811 module_id: self.inner.module_id.clone(),
2812 })?;
2813 reply_rx.await.map_err(|_| SuperviseError::CommandClosed {
2814 module_id: self.inner.module_id.clone(),
2815 })?;
2816 let mut configuration =
2817 self.inner
2818 .configuration
2819 .lock()
2820 .map_err(|_| SuperviseError::StatePoisoned {
2821 module_id: Some(self.inner.module_id.clone()),
2822 })?;
2823 configuration.spec = spec;
2824 configuration.health = health;
2825 Ok(())
2826 }
2827}
2828
2829impl Drop for SupervisedModuleInner {
2830 fn drop(&mut self) {
2831 let Ok(mut monitor) = self.monitor.lock() else {
2832 return;
2833 };
2834 if let Some(monitor) = monitor.as_ref().filter(|monitor| !monitor.is_finished()) {
2835 let _ = update_snapshot(&self.snapshot, Some(&self.module_id), |state| {
2836 state.state = ModuleState::Stopped;
2837 clear_current_process_facts(state);
2838 });
2839 monitor.abort();
2840 }
2841 let _ = monitor.take();
2842 }
2843}
2844
2845#[derive(Debug)]
2846enum SupervisorCommand {
2847 Drain {
2848 reply: oneshot::Sender<Result<(), SuperviseError>>,
2849 },
2850 Retire {
2851 reply: oneshot::Sender<Result<(), SuperviseError>>,
2852 },
2853 Restart {
2854 drain_timeout_ms: Option<u64>,
2859 reply: oneshot::Sender<Result<(), SuperviseError>>,
2860 },
2861 Reload {
2862 reply: oneshot::Sender<Result<(), SuperviseError>>,
2863 },
2864 SetEnabled {
2865 enabled: bool,
2866 reply: oneshot::Sender<Result<bool, SuperviseError>>,
2867 },
2868 UpdateConfiguration {
2869 spec: ModuleSpec,
2870 health: HealthConfig,
2871 drain_timeout_ms: Option<u64>,
2874 reply: oneshot::Sender<()>,
2875 },
2876 Swap {
2877 ready_timeout: Option<Duration>,
2880 reply: oneshot::Sender<Result<(), SuperviseError>>,
2882 },
2883}
2884
2885#[derive(Debug)]
2886pub enum SuperviseError {
2887 InvalidSpec {
2888 reason: String,
2889 },
2890 Spawn {
2891 program: PathBuf,
2892 source: io::Error,
2893 cgroup_path: Option<PathBuf>,
2894 },
2895 Cgroup {
2896 module_id: String,
2897 source: io::Error,
2898 },
2899 LaunchNonce {
2902 reason: String,
2903 },
2904 Wait {
2905 module_id: String,
2906 source: io::Error,
2907 },
2908 Kill {
2909 module_id: String,
2910 source: io::Error,
2911 },
2912 Forwarding(ForwardingError),
2913 Registry(RegistryError),
2914 ReloadUnavailable {
2915 module_id: String,
2916 reason: String,
2917 },
2918 Disabled {
2923 module_id: String,
2924 },
2925 ReloadFailed {
2926 module_id: String,
2927 reason: String,
2928 },
2929 RegistrationStillActive {
2930 module_id: String,
2931 waited: Duration,
2932 },
2933 StatePoisoned {
2934 module_id: Option<String>,
2935 },
2936 CommandClosed {
2937 module_id: String,
2938 },
2939 SwapInProgress {
2943 module_id: String,
2944 },
2945 SwapRefused {
2947 module_id: String,
2948 reason: SwapRefusal,
2949 },
2950 SwapFailed {
2954 module_id: String,
2955 arm: SwapFailureArm,
2956 detail: String,
2957 candidate_exit: Option<ExitReport>,
2960 },
2961}
2962
2963#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2965pub enum SwapRefusal {
2966 OverlapExclusive,
2968 NotRegistered,
2971 ProtocolNone,
2974 NotConfigured,
2977 AlreadySwapping,
2979}
2980
2981impl SwapRefusal {
2982 pub fn as_str(self) -> &'static str {
2983 match self {
2984 Self::OverlapExclusive => "overlap_exclusive",
2985 Self::NotRegistered => "not_registered",
2986 Self::ProtocolNone => "protocol_none",
2987 Self::NotConfigured => "not_configured",
2988 Self::AlreadySwapping => "already_swapping",
2989 }
2990 }
2991}
2992
2993#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2996pub enum SwapFailureArm {
2997 SpawnFailed,
2999 NeverRegistered,
3001 NeverReady,
3003 CandidateExited,
3005 CandidateUnhealthy,
3007 Interrupted,
3011 CutoverLost,
3016}
3017
3018impl SwapFailureArm {
3019 pub fn as_str(self) -> &'static str {
3020 match self {
3021 Self::SpawnFailed => "spawn_failed",
3022 Self::NeverRegistered => "never_registered",
3023 Self::NeverReady => "never_ready",
3024 Self::CandidateExited => "candidate_exited",
3025 Self::CandidateUnhealthy => "candidate_unhealthy",
3026 Self::Interrupted => "interrupted",
3027 Self::CutoverLost => "cutover_lost",
3028 }
3029 }
3030}
3031
3032impl fmt::Display for SuperviseError {
3033 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3034 match self {
3035 Self::InvalidSpec { reason } => write!(f, "invalid module spec: {reason}"),
3036 Self::Spawn {
3037 program,
3038 source,
3039 cgroup_path: Some(cgroup_path),
3040 } => write!(
3041 f,
3042 "failed to place module in cgroup '{}' while spawning '{}': {source}",
3043 cgroup_path.display(),
3044 program.display()
3045 ),
3046 Self::Spawn {
3047 program,
3048 source,
3049 cgroup_path: None,
3050 } => write!(
3051 f,
3052 "failed to spawn module '{}': {source}",
3053 program.display()
3054 ),
3055 Self::Cgroup { module_id, source } => {
3056 write!(
3057 f,
3058 "failed to prepare cgroup for module '{module_id}': {source}"
3059 )
3060 }
3061 Self::LaunchNonce { reason } => {
3062 write!(
3063 f,
3064 "failed to generate reserved-module launch nonce: {reason}"
3065 )
3066 }
3067 Self::Wait { module_id, source } => {
3068 write!(f, "failed to wait for module '{module_id}': {source}")
3069 }
3070 Self::Kill { module_id, source } => {
3071 write!(f, "failed to kill module '{module_id}': {source}")
3072 }
3073 Self::Forwarding(err) => write!(f, "forwarding error: {err}"),
3074 Self::Registry(err) => write!(f, "registry error: {err}"),
3075 Self::ReloadUnavailable { module_id, reason } => {
3076 write!(f, "reload unavailable for module '{module_id}': {reason}")
3077 }
3078 Self::Disabled { module_id } => {
3079 write!(
3080 f,
3081 "module '{module_id}' is disabled; enable it before restart or reload"
3082 )
3083 }
3084 Self::ReloadFailed { module_id, reason } => {
3085 write!(f, "reload failed for module '{module_id}': {reason}")
3086 }
3087 Self::RegistrationStillActive { module_id, waited } => write!(
3088 f,
3089 "module '{module_id}' registration remained active after waiting {waited:?}"
3090 ),
3091 Self::StatePoisoned { module_id } => match module_id {
3092 Some(module_id) => {
3093 write!(f, "supervisor state for module '{module_id}' was poisoned")
3094 }
3095 None => write!(f, "supervisor state was poisoned"),
3096 },
3097 Self::CommandClosed { module_id } => {
3098 write!(
3099 f,
3100 "supervisor command channel for module '{module_id}' is closed"
3101 )
3102 }
3103 Self::SwapInProgress { module_id } => write!(
3104 f,
3105 "module '{module_id}' is being swapped; retry once the swap has cut over or failed, or stop the module to abort the swap"
3106 ),
3107 Self::SwapRefused { module_id, reason } => match reason {
3108 SwapRefusal::OverlapExclusive => write!(
3109 f,
3110 "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"
3111 ),
3112 SwapRefusal::NotRegistered => write!(
3113 f,
3114 "module '{module_id}' is not registered, so there is no serving process to keep while a replacement warms; use a plain restart"
3115 ),
3116 SwapRefusal::ProtocolNone => write!(
3117 f,
3118 "module '{module_id}' is protocol: \"none\" and never registers, so a swap could never see its replacement become ready; use a plain restart"
3119 ),
3120 SwapRefusal::NotConfigured => write!(
3121 f,
3122 "module '{module_id}' cannot be swapped: the supervisor was built without the forwarding table or shared handle a swap needs"
3123 ),
3124 SwapRefusal::AlreadySwapping => {
3125 write!(f, "module '{module_id}' is already being swapped")
3126 }
3127 },
3128 Self::SwapFailed {
3129 module_id,
3130 arm,
3131 detail,
3132 ..
3133 } => write!(
3134 f,
3135 "swap of module '{module_id}' failed ({}): {detail}; the running process was left serving",
3136 arm.as_str()
3137 ),
3138 }
3139 }
3140}
3141
3142impl Error for SuperviseError {
3143 fn source(&self) -> Option<&(dyn Error + 'static)> {
3144 match self {
3145 Self::Spawn { source, .. }
3146 | Self::Cgroup { source, .. }
3147 | Self::Wait { source, .. }
3148 | Self::Kill { source, .. } => Some(source),
3149 Self::Forwarding(err) => Some(err),
3150 Self::Registry(err) => Some(err),
3151 Self::LaunchNonce { .. }
3152 | Self::InvalidSpec { .. }
3153 | Self::ReloadUnavailable { .. }
3154 | Self::Disabled { .. }
3155 | Self::ReloadFailed { .. }
3156 | Self::RegistrationStillActive { .. }
3157 | Self::StatePoisoned { .. }
3158 | Self::CommandClosed { .. }
3159 | Self::SwapInProgress { .. }
3160 | Self::SwapRefused { .. }
3161 | Self::SwapFailed { .. } => None,
3162 }
3163 }
3164}
3165
3166pub(crate) fn validate_spec(spec: &ModuleSpec) -> Result<(), SuperviseError> {
3167 if spec.module_id.trim().is_empty() {
3168 return Err(SuperviseError::InvalidSpec {
3169 reason: "module_id must not be empty".to_string(),
3170 });
3171 }
3172
3173 Ok(())
3174}
3175
3176#[derive(Debug, Default)]
3177struct HealthProbeRuntime {
3178 registered_connection: Option<crate::ConnectionId>,
3179 advertised: bool,
3180 next_probe_at: Option<Instant>,
3181 probe_index: u64,
3182}
3183
3184impl HealthProbeRuntime {
3185 fn refresh_registration(
3186 &mut self,
3187 spec: &ModuleSpec,
3188 runtime: &SupervisorRuntimeConfig,
3189 registry: &Registry,
3190 snapshot: &SharedSnapshot,
3191 ) {
3192 if spec.protocol == ModuleProtocol::None {
3204 self.registered_connection = None;
3205 self.advertised = false;
3206 self.next_probe_at = None;
3207 return;
3208 }
3209
3210 let registration = match registry.get_module(&spec.module_id) {
3211 Ok(registration) => registration,
3212 Err(err) => {
3213 warn!(module_id = %spec.module_id, error = %err, "health prober could not read registry");
3214 self.advertised = false;
3215 self.next_probe_at = None;
3216 return;
3217 }
3218 };
3219
3220 let Some(registration) = registration else {
3221 self.registered_connection = None;
3222 self.advertised = false;
3223 self.next_probe_at = None;
3224 return;
3225 };
3226
3227 let advertised = registration
3228 .control_ops
3229 .iter()
3230 .any(|op| op == MODULE_CONTROL_OP_HEALTH_CHECK);
3231 if !advertised {
3232 self.registered_connection = Some(registration.connection_id);
3233 self.advertised = false;
3234 self.next_probe_at = None;
3235 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3236 state.health.status = SupervisorHealthStatus::Unknown;
3237 state.health.consecutive_failures = 0;
3238 state.health.last_probe_ms = None;
3239 state.health.detail = None;
3240 state.health.metrics = None;
3241 });
3242 return;
3243 }
3244
3245 let reregistered = self.registered_connection != Some(registration.connection_id);
3246 self.registered_connection = Some(registration.connection_id);
3247 self.advertised = true;
3248 if reregistered || self.next_probe_at.is_none() {
3249 self.probe_index = 0;
3250 self.next_probe_at = Some(
3251 Instant::now() + jittered_health_delay(&spec.module_id, 0, runtime.health.cadence),
3252 );
3253 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3254 state.health.status = SupervisorHealthStatus::Unknown;
3255 state.health.consecutive_failures = 0;
3256 state.health.detail = None;
3257 state.health.metrics = None;
3258 });
3259 }
3260 }
3261
3262 fn wake_after(&self) -> Duration {
3263 if !self.advertised {
3264 return REGISTRY_RELEASE_POLL;
3265 }
3266 self.next_probe_at
3267 .map(|next| next.saturating_duration_since(Instant::now()))
3268 .unwrap_or(REGISTRY_RELEASE_POLL)
3269 }
3270
3271 fn due(&self) -> bool {
3272 self.advertised
3273 && self
3274 .next_probe_at
3275 .is_some_and(|next| Instant::now() >= next)
3276 }
3277
3278 fn schedule_next(&mut self, spec: &ModuleSpec, cadence: Duration) {
3279 self.probe_index = self.probe_index.wrapping_add(1);
3280 self.next_probe_at = Some(
3281 Instant::now() + jittered_health_delay(&spec.module_id, self.probe_index, cadence),
3282 );
3283 }
3284}
3285
3286#[derive(Debug)]
3321enum HealthProbeEvidence {
3322 LaneDead,
3324 NoAnswer,
3326 BadAnswer,
3328 Misconfigured,
3330}
3331
3332#[derive(Debug)]
3333struct HealthProbeError {
3334 evidence: HealthProbeEvidence,
3335 message: String,
3336}
3337
3338impl HealthProbeError {
3339 fn lane_dead(message: impl Into<String>) -> Self {
3340 Self::with(HealthProbeEvidence::LaneDead, message)
3341 }
3342
3343 fn no_answer(message: impl Into<String>) -> Self {
3344 Self::with(HealthProbeEvidence::NoAnswer, message)
3345 }
3346
3347 fn bad_answer(message: impl Into<String>) -> Self {
3348 Self::with(HealthProbeEvidence::BadAnswer, message)
3349 }
3350
3351 fn misconfigured(message: impl Into<String>) -> Self {
3352 Self::with(HealthProbeEvidence::Misconfigured, message)
3353 }
3354
3355 fn with(evidence: HealthProbeEvidence, message: impl Into<String>) -> Self {
3356 Self {
3357 evidence,
3358 message: message.into(),
3359 }
3360 }
3361
3362 #[allow(dead_code)]
3376 fn is_proof_of_death(&self) -> bool {
3377 matches!(self.evidence, HealthProbeEvidence::LaneDead)
3378 }
3379
3380 fn label(&self) -> &'static str {
3388 match self.evidence {
3389 HealthProbeEvidence::LaneDead => "lane-dead",
3390 HealthProbeEvidence::NoAnswer => "no-answer",
3391 HealthProbeEvidence::BadAnswer => "bad-answer",
3392 HealthProbeEvidence::Misconfigured => "daemon-misconfigured",
3393 }
3394 }
3395}
3396
3397impl fmt::Display for HealthProbeError {
3398 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
3399 f.write_str(&self.message)
3400 }
3401}
3402
3403async fn run_health_probe_cycle(
3404 spec: &ModuleSpec,
3405 runtime: &SupervisorRuntimeConfig,
3406 registry: &Registry,
3407 process_liveness: &SupervisorProcessLiveness,
3408 snapshot: &SharedSnapshot,
3409 child: &mut Option<SupervisedChild>,
3410) {
3411 let now_ms = unix_ms_now();
3412 match probe_module_health(&spec.module_id, runtime, None).await {
3413 Ok(report) => {
3414 handle_health_report(
3415 spec,
3416 runtime,
3417 registry,
3418 process_liveness,
3419 snapshot,
3420 child,
3421 report,
3422 now_ms,
3423 )
3424 .await;
3425 }
3426 Err(err) => {
3427 handle_health_probe_failure(
3428 spec,
3429 runtime,
3430 registry,
3431 process_liveness,
3432 snapshot,
3433 child,
3434 err,
3435 now_ms,
3436 )
3437 .await;
3438 }
3439 }
3440}
3441
3442async fn probe_module_health(
3443 module_id: &str,
3444 runtime: &SupervisorRuntimeConfig,
3445 drain_deadline: Option<Instant>,
3446) -> Result<HealthReport, HealthProbeError> {
3447 let Some(forwarding) = runtime.forwarding.as_ref() else {
3448 return Err(HealthProbeError::misconfigured(
3449 "supervisor was not configured with a forwarding table",
3450 ));
3451 };
3452 let probe_started_at = Instant::now();
3453 let mut deadline = probe_started_at + runtime.health.deadline;
3454 if let Some(drain_deadline) = drain_deadline {
3455 deadline = deadline.min(drain_deadline);
3456 }
3457 let pending = if drain_deadline.is_some() {
3458 forwarding.begin_drain_health_probe_rpc_for(
3459 module_id,
3460 MODULE_CONTROL_OP_HEALTH_CHECK,
3461 probe_started_at,
3462 deadline,
3463 )
3464 } else {
3465 forwarding.begin_health_probe_rpc_for(
3466 module_id,
3467 MODULE_CONTROL_OP_HEALTH_CHECK,
3468 probe_started_at,
3469 deadline,
3470 )
3471 }
3472 .map_err(|err| {
3473 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3476 })?;
3477 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3478}
3479
3480async fn probe_endpoint_health(
3487 endpoint: crate::ModuleEndpointId,
3488 runtime: &SupervisorRuntimeConfig,
3489 deadline_cap: Option<Instant>,
3490) -> Result<HealthReport, HealthProbeError> {
3491 let Some(forwarding) = runtime.forwarding.as_ref() else {
3492 return Err(HealthProbeError::misconfigured(
3493 "supervisor was not configured with a forwarding table",
3494 ));
3495 };
3496 let probe_started_at = Instant::now();
3497 let mut deadline = probe_started_at + runtime.health.deadline;
3498 if let Some(cap) = deadline_cap {
3499 deadline = deadline.min(cap);
3500 }
3501 let pending = forwarding
3502 .begin_endpoint_health_probe_rpc_for(
3503 endpoint,
3504 MODULE_CONTROL_OP_HEALTH_CHECK,
3505 probe_started_at,
3506 deadline,
3507 )
3508 .map_err(|err| {
3509 HealthProbeError::lane_dead(format!("failed to begin health.check RPC: {err}"))
3510 })?;
3511 await_health_probe(forwarding, pending, deadline, runtime.health.deadline).await
3512}
3513
3514async fn await_health_probe(
3516 forwarding: &ForwardingTable,
3517 pending: PendingModuleControlRpc,
3518 deadline: Instant,
3519 probe_budget: Duration,
3520) -> Result<HealthReport, HealthProbeError> {
3521 let PendingModuleControlRpc {
3522 endpoint,
3523 module_sink,
3524 negotiated_ver,
3525 corr,
3526 receiver,
3527 } = pending;
3528 let body = serde_json::to_vec(&ModuleControlRequest::HealthCheck {}).map_err(|err| {
3529 HealthProbeError::misconfigured(format!("failed to encode health.check: {err}"))
3530 })?;
3531 let frame = Frame::build_with_version(
3532 negotiated_ver,
3533 FrameType::Request,
3534 control_flags(),
3535 0,
3536 0,
3537 corr,
3538 body,
3539 )
3540 .map_err(|err| {
3541 HealthProbeError::misconfigured(format!("failed to build health.check frame: {err}"))
3542 })?;
3543
3544 match timeout_at(deadline, module_sink.send(frame)).await {
3550 Ok(Ok(())) => {}
3551 Ok(Err(err)) => {
3552 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3553 return Err(HealthProbeError::lane_dead(format!(
3556 "failed to send health.check: {err}"
3557 )));
3558 }
3559 Err(_elapsed) => {
3560 let _ = forwarding.cancel_module_control_rpc(endpoint, corr);
3561 return Err(HealthProbeError::no_answer(
3565 "health.check send timed out before enqueue (module egress full)",
3566 ));
3567 }
3568 }
3569
3570 match timeout_at(deadline, receiver).await {
3571 Ok(Ok(ModuleControlRpcOutcome::Response(response))) => {
3575 response.health_report().ok_or_else(|| {
3576 HealthProbeError::bad_answer("health.check RPC returned a non-health response")
3577 })
3578 }
3579 Ok(Ok(ModuleControlRpcOutcome::Rejected(body))) => Err(HealthProbeError::bad_answer(
3580 format!("health.check rejected: {}", body.message),
3581 )),
3582 Ok(Ok(ModuleControlRpcOutcome::ModuleGone(message))) => {
3583 Err(HealthProbeError::lane_dead(message))
3584 }
3585 Ok(Ok(ModuleControlRpcOutcome::MalformedResponse(message))) => {
3586 Err(HealthProbeError::bad_answer(message))
3587 }
3588 Ok(Ok(ModuleControlRpcOutcome::UnexpectedOp { expected, actual })) => {
3589 Err(HealthProbeError::bad_answer(format!(
3590 "expected module-control op '{expected}', got '{actual}'"
3591 )))
3592 }
3593 Ok(Ok(ModuleControlRpcOutcome::DeadlineElapsed)) => Err(HealthProbeError::bad_answer(
3597 "module answered health.check after its daemon deadline",
3598 )),
3599 Ok(Err(_)) => Err(HealthProbeError::misconfigured(
3600 "health.check waiter was canceled before the module responded",
3601 )),
3602 Err(_) => {
3603 let _ = forwarding.tombstone_health_probe_rpc(endpoint, corr);
3604 Err(HealthProbeError::no_answer(format!(
3605 "module did not answer health.check within {probe_budget:?}"
3606 )))
3607 }
3608 }
3609}
3610
3611#[allow(clippy::too_many_arguments)]
3612async fn handle_health_report(
3613 spec: &ModuleSpec,
3614 runtime: &SupervisorRuntimeConfig,
3615 registry: &Registry,
3616 process_liveness: &SupervisorProcessLiveness,
3617 snapshot: &SharedSnapshot,
3618 child: &mut Option<SupervisedChild>,
3619 report: HealthReport,
3620 now_ms: u64,
3621) {
3622 let status = supervisor_health_status(report.status);
3623 let detail = report.detail.clone();
3624 let metrics = truncate_health_metrics(report.metrics);
3625 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3626 state.health.status = status;
3627 state.health.last_probe_ms = Some(now_ms);
3628 state.health.detail = detail.clone();
3629 state.health.metrics = metrics.clone();
3630 state.health.consecutive_failures = 0;
3631 });
3632
3633 let action = match report.status {
3634 HealthStatus::Ok => return,
3635 HealthStatus::Degraded => runtime.health.on_degraded,
3636 HealthStatus::Failing => runtime.health.on_failing,
3637 };
3638 apply_l3_health_action(
3639 spec,
3640 runtime,
3641 registry,
3642 process_liveness,
3643 snapshot,
3644 child,
3645 status,
3646 detail.as_deref(),
3647 action,
3648 now_ms,
3649 )
3650 .await;
3651}
3652
3653#[allow(clippy::too_many_arguments)]
3654async fn handle_health_probe_failure(
3655 spec: &ModuleSpec,
3656 runtime: &SupervisorRuntimeConfig,
3657 registry: &Registry,
3658 process_liveness: &SupervisorProcessLiveness,
3659 snapshot: &SharedSnapshot,
3660 child: &mut Option<SupervisedChild>,
3661 err: HealthProbeError,
3662 now_ms: u64,
3663) {
3664 let threshold = runtime.health.failure_threshold.max(1);
3665 let mut failures = 0;
3666 let detail = format!("[{}] {err}", err.label());
3671 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3672 state.health.last_probe_ms = Some(now_ms);
3673 state.health.consecutive_failures = state.health.consecutive_failures.saturating_add(1);
3674 state.health.detail = Some(detail.clone());
3675 state.health.metrics = None;
3676 failures = state.health.consecutive_failures;
3677 });
3678
3679 if failures < threshold {
3680 warn!(
3681 module_id = %spec.module_id,
3682 consecutive_failures = failures,
3683 threshold,
3684 evidence = err.label(),
3685 detail = %detail,
3686 "health.check probe failed"
3687 );
3688 return;
3689 }
3690
3691 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
3692 state.state = ModuleState::Unresponsive;
3693 state.health.status = SupervisorHealthStatus::Unresponsive;
3694 });
3695 if runtime.health.critical {
3699 error!(
3700 module_id = %spec.module_id,
3701 status = "unresponsive",
3702 evidence = err.label(),
3703 detail = %detail,
3704 "critical module health alert"
3705 );
3706 } else {
3707 warn!(
3708 module_id = %spec.module_id,
3709 status = "unresponsive",
3710 evidence = err.label(),
3711 detail = %detail,
3712 "module health threshold breached"
3713 );
3714 }
3715 if let Err(err) = health_restart_child(
3716 spec,
3717 runtime,
3718 registry,
3719 process_liveness,
3720 snapshot,
3721 child,
3722 SupervisorHealthStatus::Unresponsive,
3723 Some(&detail),
3724 now_ms,
3725 )
3726 .await
3727 {
3728 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3729 }
3730}
3731
3732#[allow(clippy::too_many_arguments)]
3733async fn apply_l3_health_action(
3734 spec: &ModuleSpec,
3735 runtime: &SupervisorRuntimeConfig,
3736 registry: &Registry,
3737 process_liveness: &SupervisorProcessLiveness,
3738 snapshot: &SharedSnapshot,
3739 child: &mut Option<SupervisedChild>,
3740 status: SupervisorHealthStatus,
3741 detail: Option<&str>,
3742 action: HealthAction,
3743 now_ms: u64,
3744) {
3745 record_health_action(snapshot, &spec.module_id, action.to_string(), now_ms);
3746 match action {
3747 HealthAction::Report => {
3748 info!(
3749 module_id = %spec.module_id,
3750 status = ?status,
3751 detail,
3752 "module reported non-ok health"
3753 );
3754 }
3755 HealthAction::Alert => {
3756 error!(
3757 module_id = %spec.module_id,
3758 status = ?status,
3759 detail,
3760 "module health alert"
3761 );
3762 }
3763 HealthAction::Restart => {
3764 if let Err(err) = health_restart_child(
3765 spec,
3766 runtime,
3767 registry,
3768 process_liveness,
3769 snapshot,
3770 child,
3771 status,
3772 detail,
3773 now_ms,
3774 )
3775 .await
3776 {
3777 error!(module_id = %spec.module_id, error = %err, "health-triggered restart failed");
3778 }
3779 }
3780 }
3781}
3782
3783#[allow(clippy::too_many_arguments)]
3784async fn health_restart_child(
3785 spec: &ModuleSpec,
3786 runtime: &SupervisorRuntimeConfig,
3787 registry: &Registry,
3788 process_liveness: &SupervisorProcessLiveness,
3789 snapshot: &SharedSnapshot,
3790 child: &mut Option<SupervisedChild>,
3791 status: SupervisorHealthStatus,
3792 detail: Option<&str>,
3793 now_ms: u64,
3794) -> Result<(), SuperviseError> {
3795 let (enabled, schedule) = {
3796 let mut state = lock_snapshot(snapshot)?;
3797 let enabled = state.enabled;
3798 let schedule = if enabled {
3799 state.next_crash_restart(&runtime.restart_policy, Instant::now())
3800 } else {
3801 None
3802 };
3803 (enabled, schedule)
3804 };
3805
3806 if !enabled {
3807 return Err(SuperviseError::Disabled {
3808 module_id: spec.module_id.clone(),
3809 });
3810 }
3811
3812 if schedule.is_none() {
3813 record_health_action(snapshot, &spec.module_id, "disabled".to_string(), now_ms);
3814 error!(
3815 module_id = %spec.module_id,
3816 status = ?status,
3817 detail,
3818 max_restarts = runtime.restart_policy.max_restarts,
3819 window_secs = runtime.restart_policy.window.as_secs(),
3820 "health restart budget exhausted; disabling module"
3821 );
3822 begin_forwarding_drain_if_configured(
3823 spec,
3824 runtime,
3825 registry,
3826 snapshot,
3827 Some(false),
3828 RouteCloseReason::Disable,
3829 )
3830 .await?;
3831 drain_optional_child(
3832 &spec.module_id,
3833 spec.protocol,
3834 registry,
3835 snapshot,
3836 &runtime.terminal_ring,
3837 &runtime.spawn_events,
3838 child,
3839 runtime.drain_timeout,
3840 ModuleState::Disabled,
3841 Some(false),
3842 )
3843 .await?;
3844 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3845 return Ok(());
3846 }
3847
3848 let schedule = schedule.expect("a health restart must have a crash-restart schedule");
3849 let mut restart_count = 0;
3850 update_snapshot(snapshot, Some(&spec.module_id), |state| {
3851 restart_count = state.crash_restarts.len();
3852 state.state = ModuleState::Unresponsive;
3853 state.health.status = status;
3854 state.health.last_action = Some(HealthAction::Restart.to_string());
3855 state.health.last_action_ms = Some(now_ms);
3856 })?;
3857 warn!(
3858 module_id = %spec.module_id,
3859 status = ?status,
3860 detail,
3861 restart_count,
3862 restart_in_window = schedule.restart_in_window,
3863 delay_ms = schedule.delay.as_millis() as u64,
3864 "health-triggered module restart"
3865 );
3866
3867 begin_forwarding_drain_if_configured(
3868 spec,
3869 runtime,
3870 registry,
3871 snapshot,
3872 Some(true),
3873 RouteCloseReason::Restart,
3874 )
3875 .await?;
3876 drain_optional_child(
3877 &spec.module_id,
3878 spec.protocol,
3879 registry,
3880 snapshot,
3881 &runtime.terminal_ring,
3882 &runtime.spawn_events,
3883 child,
3884 runtime.drain_timeout,
3885 ModuleState::Restarting,
3886 Some(true),
3887 )
3888 .await?;
3889 sleep(schedule.delay).await;
3890 if !respawn_still_pending(snapshot) {
3894 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3895 return Ok(());
3896 }
3897 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
3898 match spawn_and_mark_running(spec, runtime, snapshot) {
3899 Ok(next_child) => {
3900 *child = Some(next_child);
3901 Ok(())
3902 }
3903 Err(err) => {
3904 fail_snapshot(snapshot, Some(&spec.module_id), None);
3905 process_liveness.untrack_if_current(&spec.module_id, snapshot);
3906 *child = None;
3907 Err(err)
3908 }
3909 }
3910}
3911
3912fn record_health_action(snapshot: &SharedSnapshot, module_id: &str, action: String, now_ms: u64) {
3913 let _ = update_snapshot(snapshot, Some(module_id), |state| {
3914 state.health.last_action = Some(action);
3915 state.health.last_action_ms = Some(now_ms);
3916 });
3917}
3918
3919fn supervisor_health_status(status: HealthStatus) -> SupervisorHealthStatus {
3920 match status {
3921 HealthStatus::Ok => SupervisorHealthStatus::Ok,
3922 HealthStatus::Degraded => SupervisorHealthStatus::Degraded,
3923 HealthStatus::Failing => SupervisorHealthStatus::Failing,
3924 }
3925}
3926
3927fn truncate_health_metrics(metrics: Option<Value>) -> Option<Value> {
3939 let metrics = metrics?;
3940 match serde_json::to_vec(&metrics) {
3941 Ok(encoded) if encoded.len() > MAX_HEALTH_METRICS_BYTES => Some(serde_json::json!({
3942 "truncated": true,
3943 "original_bytes": encoded.len(),
3944 })),
3945 Ok(_) | Err(_) => Some(metrics),
3946 }
3947}
3948
3949fn jittered_health_delay(module_id: &str, probe_index: u64, cadence: Duration) -> Duration {
3955 if cadence.is_zero() {
3956 return Duration::ZERO;
3957 }
3958 let cadence_ms = cadence.as_millis() as u64;
3959 if cadence_ms == 0 {
3975 return cadence;
3976 }
3977 let jitter_span = (cadence_ms / 10).max(1);
3992 let hash = module_id.as_bytes().iter().fold(
3993 probe_index.wrapping_mul(0x9E37_79B9_7F4A_7C15),
3994 |acc, byte| {
3995 acc.wrapping_mul(1099511628211)
3996 .wrapping_add(u64::from(*byte))
3997 },
3998 );
3999 cadence + Duration::from_millis(hash % jitter_span)
4000}
4001
4002#[cfg(test)]
4003mod tests {
4004 use super::*;
4005
4006 #[test]
4007 fn readding_a_module_clears_its_rescan_removal_tombstone() {
4008 let handle = SupervisorHandle::new();
4009 let module_id = "readded-tombstone";
4010 handle.record_rescan_removal(module_id);
4011 assert!(handle.removal_tombstone_age_ms(module_id).is_some());
4012
4013 handle.apply_identity_configuration(&ModuleSpec {
4014 module_id: module_id.to_string(),
4015 program: PathBuf::from("/test/module"),
4016 args: Vec::new(),
4017 env: Vec::new(),
4018 reserved: false,
4019 reserved_prefixes: Vec::new(),
4020 protocol: ModuleProtocol::Subc,
4021 overlap: Default::default(),
4022 });
4023
4024 assert!(
4025 handle.removal_tombstone_age_ms(module_id).is_none(),
4026 "a re-added module must not retain a stale removal tombstone"
4027 );
4028 }
4029
4030 fn stale_process_snapshot(state: ModuleState, enabled: bool) -> SharedSnapshot {
4031 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::new(state, enabled)));
4032 update_snapshot(&snapshot, Some("stale-process-facts"), |snapshot| {
4033 snapshot.process_alive = true;
4034 snapshot.pid = Some(41);
4035 snapshot.spawned_at_ms = Some(42);
4036 snapshot.spawned_from = Some(PathBuf::from("/spawned/module"));
4037 snapshot.spawned_file_identity = Some(SpawnedFileIdentity {
4038 device: 43,
4039 inode: 44,
4040 });
4041 })
4042 .unwrap();
4043 snapshot
4044 }
4045
4046 fn assert_snapshot_process_facts_cleared(snapshot: &SharedSnapshot) {
4047 let snapshot = lock_snapshot(snapshot).unwrap();
4048 assert!(!snapshot.process_alive);
4049 assert_eq!(snapshot.pid, None);
4050 assert_eq!(snapshot.spawned_at_ms, None);
4051 assert_eq!(snapshot.spawned_from, None);
4052 assert_eq!(snapshot.spawned_file_identity, None);
4053 }
4054
4055 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4056 async fn failed_enable_spawn_clears_preexisting_current_process_facts() {
4057 let supervisor = Supervisor::default();
4058 let mut runtime = supervisor.runtime_config();
4059 runtime.test_seed_stale_facts_before_enable_spawn = true;
4060 let snapshot = stale_process_snapshot(ModuleState::Disabled, false);
4061 let mut child = None;
4062 let spec = ModuleSpec {
4063 module_id: "failed-enable-clears-facts".to_string(),
4064 program: PathBuf::from("/definitely/missing/failed-enable-module"),
4065 args: Vec::new(),
4066 env: Vec::new(),
4067 reserved: false,
4068 reserved_prefixes: Vec::new(),
4069 protocol: ModuleProtocol::Subc,
4070 overlap: Default::default(),
4071 };
4072
4073 let result = set_child_enabled(
4074 &spec,
4075 &runtime,
4076 &supervisor.registry,
4077 &supervisor.process_liveness,
4078 &snapshot,
4079 &mut child,
4080 true,
4081 )
4082 .await;
4083
4084 assert!(matches!(result, Err(SuperviseError::Spawn { .. })));
4085 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4086 assert_snapshot_process_facts_cleared(&snapshot);
4087 }
4088
4089 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4090 async fn failed_reload_spawn_clears_current_process_facts() {
4091 let supervisor = Supervisor::default();
4092 let mut runtime = supervisor.runtime_config();
4093 runtime.restart_policy = RestartPolicy::new(0, Duration::ZERO);
4094 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4095 let mut child = None;
4096 let spec = ModuleSpec {
4097 module_id: "failed-reload-clears-facts".to_string(),
4098 program: PathBuf::from("/unused/failed-reload-module"),
4099 args: Vec::new(),
4100 env: Vec::new(),
4101 reserved: false,
4102 reserved_prefixes: Vec::new(),
4103 protocol: ModuleProtocol::Subc,
4104 overlap: Default::default(),
4105 };
4106
4107 let result = handle_reload_spawn_failure(
4108 &spec,
4109 &runtime,
4110 &supervisor.process_liveness,
4111 &snapshot,
4112 &mut child,
4113 "forced reload spawn failure".to_string(),
4114 )
4115 .await;
4116
4117 assert!(matches!(result, Err(SuperviseError::ReloadFailed { .. })));
4118 assert_eq!(lock_snapshot(&snapshot).unwrap().state, ModuleState::Failed);
4119 assert_snapshot_process_facts_cleared(&snapshot);
4120 }
4121
4122 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4123 async fn dropping_a_module_with_an_active_monitor_clears_current_process_facts() {
4124 let supervisor = Supervisor::default();
4125 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4126 let module = supervisor.supervised_module(
4127 ModuleSpec {
4128 module_id: "drop-clears-facts".to_string(),
4129 program: PathBuf::from("/unused/drop-module"),
4130 args: Vec::new(),
4131 env: Vec::new(),
4132 reserved: false,
4133 reserved_prefixes: Vec::new(),
4134 protocol: ModuleProtocol::Subc,
4135 overlap: Default::default(),
4136 },
4137 supervisor.runtime_config(),
4138 Arc::clone(&snapshot),
4139 None,
4140 );
4141 assert!(!module
4142 .inner
4143 .monitor
4144 .lock()
4145 .unwrap()
4146 .as_ref()
4147 .unwrap()
4148 .is_finished());
4149
4150 drop(module);
4151
4152 assert_eq!(
4153 lock_snapshot(&snapshot).unwrap().state,
4154 ModuleState::Stopped
4155 );
4156 assert_snapshot_process_facts_cleared(&snapshot);
4157 }
4158
4159 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4160 async fn configuration_update_does_not_replace_captured_running_process_facts() {
4161 let supervisor = Supervisor::default();
4162 let snapshot = stale_process_snapshot(ModuleState::Running, true);
4163 let initial = ModuleSpec {
4164 module_id: "rescan-preserves-spawn-facts".to_string(),
4165 program: PathBuf::from("/spawned/module"),
4166 args: Vec::new(),
4167 env: Vec::new(),
4168 reserved: false,
4169 reserved_prefixes: Vec::new(),
4170 protocol: ModuleProtocol::Subc,
4171 overlap: Default::default(),
4172 };
4173 let module = supervisor.supervised_module(
4174 initial.clone(),
4175 supervisor.runtime_config(),
4176 snapshot,
4177 None,
4178 );
4179 let before = module.status().unwrap();
4180 let mut replacement = initial;
4181 replacement.program = PathBuf::from("/rescanned/replacement-module");
4182
4183 module
4184 .update_configuration(replacement, HealthConfig::default(), None)
4185 .await
4186 .unwrap();
4187
4188 let after = module.status().unwrap();
4189 assert_eq!(after.pid, before.pid);
4190 assert_eq!(after.spawned_at_ms, before.spawned_at_ms);
4191 assert_eq!(after.spawned_from, before.spawned_from);
4192 drop(module);
4193 }
4194}
4195
4196fn unix_ms_now() -> u64 {
4197 SystemTime::now()
4198 .duration_since(UNIX_EPOCH)
4199 .map(|duration| duration.as_millis().min(u128::from(u64::MAX)) as u64)
4200 .unwrap_or(0)
4201}
4202
4203async fn supervise_loop(
4204 mut spec: ModuleSpec,
4205 mut runtime: SupervisorRuntimeConfig,
4206 registry: Arc<Registry>,
4207 process_liveness: Arc<SupervisorProcessLiveness>,
4208 snapshot: SharedSnapshot,
4209 mut child: Option<SupervisedChild>,
4210 mut commands: mpsc::Receiver<SupervisorCommand>,
4211) {
4212 let mut health_probe = HealthProbeRuntime::default();
4213 let mut pending_respawn: Option<Instant> = None;
4217 let mut requeued: VecDeque<SupervisorCommand> = VecDeque::new();
4220 loop {
4221 if let Some(command) = requeued.pop_front() {
4222 if !handle_supervisor_command(
4223 command,
4224 &mut spec,
4225 &mut runtime,
4226 ®istry,
4227 &process_liveness,
4228 &snapshot,
4229 &mut child,
4230 &mut commands,
4231 &mut requeued,
4232 )
4233 .await
4234 {
4235 return;
4236 }
4237 if child.is_some() || !respawn_still_pending(&snapshot) {
4238 pending_respawn = None;
4239 }
4240 continue;
4241 }
4242 if child.is_some() {
4243 health_probe.refresh_registration(&spec, &runtime, ®istry, &snapshot);
4244 let probe_sleep = sleep(health_probe.wake_after());
4245 tokio::pin!(probe_sleep);
4246 let active_child = child.as_mut().expect("child checked above");
4247 tokio::select! {
4248 wait_result = active_child.wait() => {
4249 let exit_report = match wait_result {
4258 Ok(status) => classify_reaped_child_exit(&snapshot, active_child, &status),
4259 Err(err) => {
4260 active_child.drain_stderr(&spec.module_id).await;
4261 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4262 record_wait_error_terminal(
4268 &spec.module_id,
4269 &runtime.terminal_ring,
4270 &runtime.spawn_events,
4271 );
4272 untrack_if_registration_released(
4273 &process_liveness,
4274 ®istry,
4275 &spec.module_id,
4276 &snapshot,
4277 );
4278 error!(module_id = %spec.module_id, error = %err, "failed to wait for supervised module");
4279 child = None;
4280 continue;
4281 }
4282 };
4283 active_child.drain_stderr(&spec.module_id).await;
4284
4285 match on_child_exit(
4286 &spec,
4287 runtime.restart_policy,
4288 ®istry,
4289 &snapshot,
4290 &runtime.terminal_ring,
4291 &runtime.spawn_events,
4292 exit_report,
4293 ).await {
4294 NextAction::Stop { registration_released } => {
4295 if registration_released {
4296 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4297 }
4298 child = None;
4299 }
4300 NextAction::Restart { schedule } => {
4301 let delay = schedule.map_or(
4302 runtime.restart_policy.delay_for_restart(0),
4303 |schedule| schedule.delay,
4304 );
4305 if let Some(schedule) = schedule {
4306 log_crash_respawn(&spec.module_id, schedule);
4307 }
4308 child = None;
4316 pending_respawn = Some(Instant::now() + delay);
4317 }
4318 }
4319 }
4320 command = commands.recv() => {
4321 let Some(command) = command else {
4322 return;
4323 };
4324 if !handle_supervisor_command(
4325 command,
4326 &mut spec,
4327 &mut runtime,
4328 ®istry,
4329 &process_liveness,
4330 &snapshot,
4331 &mut child,
4332 &mut commands,
4333 &mut requeued,
4334 ).await {
4335 return;
4336 }
4337 }
4338 _ = &mut probe_sleep => {
4339 if health_probe.due() {
4340 run_health_probe_cycle(
4341 &spec,
4342 &runtime,
4343 ®istry,
4344 &process_liveness,
4345 &snapshot,
4346 &mut child,
4347 ).await;
4348 if child.is_some() {
4349 health_probe.schedule_next(&spec, runtime.health.cadence);
4350 }
4351 }
4352 }
4353 }
4354 } else if let Some(deadline) = pending_respawn {
4355 tokio::select! {
4356 _ = sleep_until(deadline) => {
4357 pending_respawn = None;
4358 if !respawn_still_pending(&snapshot) {
4362 continue;
4363 }
4364 if let Err(err) = wait_for_registration_release(
4365 ®istry,
4366 &spec.module_id,
4367 REGISTRY_RELEASE_TIMEOUT,
4368 ).await {
4369 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4370 error!(module_id = %spec.module_id, error = %err, "registration did not release before restart");
4371 continue;
4372 }
4373
4374 match spawn_and_mark_running(&spec, &runtime, &snapshot) {
4375 Ok(next_child) => {
4376 child = Some(next_child);
4377 debug!(module_id = %spec.module_id, "supervised module restarted after crash");
4378 }
4379 Err(err) => {
4380 fail_snapshot(&snapshot, Some(&spec.module_id), None);
4381 process_liveness.untrack_if_current(&spec.module_id, &snapshot);
4382 error!(module_id = %spec.module_id, error = %err, "failed to restart supervised module");
4383 }
4384 }
4385 }
4386 command = commands.recv() => {
4387 let Some(command) = command else {
4388 return;
4389 };
4390 if !handle_supervisor_command(
4391 command,
4392 &mut spec,
4393 &mut runtime,
4394 ®istry,
4395 &process_liveness,
4396 &snapshot,
4397 &mut child,
4398 &mut commands,
4399 &mut requeued,
4400 ).await {
4401 return;
4402 }
4403 if child.is_some() || !respawn_still_pending(&snapshot) {
4408 pending_respawn = None;
4409 }
4410 }
4411 }
4412 } else {
4413 let Some(command) = commands.recv().await else {
4414 return;
4415 };
4416 if !handle_supervisor_command(
4417 command,
4418 &mut spec,
4419 &mut runtime,
4420 ®istry,
4421 &process_liveness,
4422 &snapshot,
4423 &mut child,
4424 &mut commands,
4425 &mut requeued,
4426 )
4427 .await
4428 {
4429 return;
4430 }
4431 }
4432 }
4433}
4434
4435fn log_crash_respawn(module_id: &str, schedule: CrashRestartSchedule) {
4436 info!(
4437 module_id,
4438 restart_in_window = schedule.restart_in_window,
4439 delay_ms = schedule.delay.as_millis() as u64,
4440 "respawning after crash"
4441 );
4442}
4443
4444fn respawn_still_pending(snapshot: &SharedSnapshot) -> bool {
4450 matches!(
4451 lock_snapshot(snapshot),
4452 Ok(state) if state.enabled && state.state == ModuleState::Restarting
4453 )
4454}
4455
4456enum NextAction {
4457 Stop {
4458 registration_released: bool,
4459 },
4460 Restart {
4461 schedule: Option<CrashRestartSchedule>,
4462 },
4463}
4464
4465#[allow(clippy::too_many_arguments)]
4466async fn handle_supervisor_command(
4467 command: SupervisorCommand,
4468 spec: &mut ModuleSpec,
4469 runtime: &mut SupervisorRuntimeConfig,
4470 registry: &Registry,
4471 process_liveness: &SupervisorProcessLiveness,
4472 snapshot: &SharedSnapshot,
4473 child: &mut Option<SupervisedChild>,
4474 commands: &mut mpsc::Receiver<SupervisorCommand>,
4475 requeued: &mut VecDeque<SupervisorCommand>,
4476) -> bool {
4477 match command {
4478 SupervisorCommand::Drain { reply } => {
4479 let result = drain_optional_child(
4480 &spec.module_id,
4481 spec.protocol,
4482 registry,
4483 snapshot,
4484 &runtime.terminal_ring,
4485 &runtime.spawn_events,
4486 child,
4487 runtime.drain_timeout,
4488 ModuleState::Stopped,
4489 None,
4490 )
4491 .await;
4492 let registration_released = result.is_ok();
4493 let _ = reply.send(result);
4494 if registration_released {
4495 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4496 }
4497 false
4498 }
4499 SupervisorCommand::Retire { reply } => {
4500 let result = async {
4501 begin_forwarding_drain_if_configured(
4502 spec,
4503 runtime,
4504 registry,
4505 snapshot,
4506 None,
4507 RouteCloseReason::Disable,
4508 )
4509 .await?;
4510 drain_optional_child(
4511 &spec.module_id,
4512 spec.protocol,
4513 registry,
4514 snapshot,
4515 &runtime.terminal_ring,
4516 &runtime.spawn_events,
4517 child,
4518 runtime.drain_timeout,
4519 ModuleState::Stopped,
4520 None,
4521 )
4522 .await
4523 }
4524 .await;
4525 let registration_released = result.is_ok();
4526 let _ = reply.send(result);
4527 if registration_released {
4528 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4529 }
4530 false
4531 }
4532 SupervisorCommand::Restart {
4533 drain_timeout_ms,
4534 reply,
4535 } => {
4536 let validation = match lock_snapshot(snapshot) {
4548 Ok(state) if !state.enabled => Err(SuperviseError::Disabled {
4549 module_id: spec.module_id.clone(),
4550 }),
4551 Ok(_) => Ok(()),
4552 Err(err) => Err(err),
4553 };
4554 let initiated = validation.is_ok();
4555 let _ = reply.send(validation);
4556 if initiated {
4557 let drain_timeout = drain_timeout_ms
4560 .map(Duration::from_millis)
4561 .unwrap_or(runtime.drain_timeout);
4562 if let Err(err) = restart_child(
4563 spec,
4564 runtime,
4565 registry,
4566 process_liveness,
4567 snapshot,
4568 child,
4569 drain_timeout,
4570 )
4571 .await
4572 {
4573 warn!(
4574 module_id = %spec.module_id,
4575 error = %err,
4576 "operator restart failed after initiation ack; module state carries the outcome"
4577 );
4578 let _ = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4579 state.state = ModuleState::Failed;
4580 clear_current_process_facts(state);
4581 });
4582 }
4583 }
4584 true
4585 }
4586 SupervisorCommand::Reload { reply } => {
4587 let result =
4588 reload_child(spec, runtime, registry, process_liveness, snapshot, child).await;
4589 let _ = reply.send(result);
4590 true
4591 }
4592 SupervisorCommand::SetEnabled { enabled, reply } => {
4593 let result = set_child_enabled(
4594 spec,
4595 runtime,
4596 registry,
4597 process_liveness,
4598 snapshot,
4599 child,
4600 enabled,
4601 )
4602 .await;
4603 let _ = reply.send(result);
4604 true
4605 }
4606 SupervisorCommand::UpdateConfiguration {
4607 spec: next_spec,
4608 health,
4609 drain_timeout_ms,
4610 reply,
4611 } => {
4612 if let Some(handle) = &runtime.supervisor_handle {
4613 handle.apply_identity_configuration(&next_spec);
4614 }
4615 *spec = next_spec;
4616 runtime.health = health;
4617 runtime.drain_timeout = drain_timeout_ms
4618 .map(Duration::from_millis)
4619 .unwrap_or(runtime.default_drain_timeout);
4620 *runtime
4621 .effective_drain_timeout
4622 .lock()
4623 .unwrap_or_else(|poisoned| poisoned.into_inner()) = runtime.drain_timeout;
4624 let _ = reply.send(());
4625 true
4626 }
4627 SupervisorCommand::Swap {
4628 ready_timeout,
4629 reply,
4630 } => {
4631 let end = swap::run_swap(
4632 spec,
4633 runtime,
4634 registry,
4635 process_liveness,
4636 snapshot,
4637 child,
4638 commands,
4639 ready_timeout.unwrap_or(DEFAULT_SWAP_READY_TIMEOUT),
4640 reply,
4641 )
4642 .await;
4643 requeued.extend(end.requeue);
4644 true
4645 }
4646 }
4647}
4648
4649async fn restart_child(
4650 spec: &ModuleSpec,
4651 runtime: &SupervisorRuntimeConfig,
4652 registry: &Registry,
4653 process_liveness: &SupervisorProcessLiveness,
4654 snapshot: &SharedSnapshot,
4655 child: &mut Option<SupervisedChild>,
4656 drain_timeout: Duration,
4657) -> Result<(), SuperviseError> {
4658 if !lock_snapshot(snapshot)?.enabled {
4660 return Err(SuperviseError::Disabled {
4661 module_id: spec.module_id.clone(),
4662 });
4663 }
4664 begin_forwarding_drain_with_timeout(
4665 spec,
4666 runtime,
4667 registry,
4668 snapshot,
4669 None,
4670 RouteCloseReason::Restart,
4671 drain_timeout,
4672 )
4673 .await?;
4674
4675 if child.is_some() {
4676 drain_optional_child(
4677 &spec.module_id,
4678 spec.protocol,
4679 registry,
4680 snapshot,
4681 &runtime.terminal_ring,
4682 &runtime.spawn_events,
4683 child,
4684 drain_timeout,
4685 ModuleState::Restarting,
4686 Some(true),
4687 )
4688 .await?;
4689 } else {
4690 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4691 state.enabled = true;
4692 state.state = ModuleState::Restarting;
4693 clear_current_process_facts(state);
4694 })?;
4695 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4696 }
4697
4698 reset_restart_count(snapshot, &spec.module_id)?;
4699 sleep(runtime.restart_policy.backoff).await;
4700 if !respawn_still_pending(snapshot) {
4703 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4704 return Ok(());
4705 }
4706 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4707 match spawn_and_mark_running(spec, runtime, snapshot) {
4713 Ok(next_child) => {
4714 *child = Some(next_child);
4715 debug!(module_id = %spec.module_id, "supervised module restarted by operator request");
4716 Ok(())
4717 }
4718 Err(err) => {
4719 fail_snapshot(snapshot, Some(&spec.module_id), None);
4720 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4721 *child = None;
4722 Err(err)
4723 }
4724 }
4725}
4726
4727async fn reload_child(
4728 spec: &ModuleSpec,
4729 runtime: &SupervisorRuntimeConfig,
4730 registry: &Registry,
4731 process_liveness: &SupervisorProcessLiveness,
4732 snapshot: &SharedSnapshot,
4733 child: &mut Option<SupervisedChild>,
4734) -> Result<(), SuperviseError> {
4735 if !lock_snapshot(snapshot)?.enabled {
4737 return Err(SuperviseError::Disabled {
4738 module_id: spec.module_id.clone(),
4739 });
4740 }
4741 begin_forwarding_drain(
4742 spec,
4743 runtime,
4744 registry,
4745 snapshot,
4746 Some(true),
4747 RouteCloseReason::Reload,
4748 )
4749 .await?;
4750
4751 if child.is_some() {
4752 drain_optional_child(
4753 &spec.module_id,
4754 spec.protocol,
4755 registry,
4756 snapshot,
4757 &runtime.terminal_ring,
4758 &runtime.spawn_events,
4759 child,
4760 runtime.drain_timeout,
4761 ModuleState::Restarting,
4762 Some(true),
4763 )
4764 .await?;
4765 } else {
4766 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4767 state.enabled = true;
4768 state.state = ModuleState::Restarting;
4769 clear_current_process_facts(state);
4770 })?;
4771 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4772 }
4773
4774 reset_restart_count(snapshot, &spec.module_id)?;
4775 sleep(runtime.restart_policy.backoff).await;
4776 if !respawn_still_pending(snapshot) {
4779 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4780 return Ok(());
4781 }
4782 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4783 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
4784 Ok(next_child) => next_child,
4785 Err(err) => {
4786 return handle_reload_spawn_failure(
4787 spec,
4788 runtime,
4789 process_liveness,
4790 snapshot,
4791 child,
4792 format!("new child failed to spawn: {err}"),
4793 )
4794 .await;
4795 }
4796 };
4797 *child = Some(next_child);
4798
4799 let wait_outcome = {
4800 let active_child = child.as_mut().expect("new reload child was just stored");
4801 wait_for_registration_after_reload(
4802 registry,
4803 &spec.module_id,
4804 snapshot,
4805 active_child,
4806 REGISTRY_RELEASE_TIMEOUT,
4807 )
4808 .await?
4809 };
4810
4811 match wait_outcome {
4812 RegistrationWaitOutcome::Registered => {
4813 debug!(module_id = %spec.module_id, "supervised module reloaded and registered");
4814 Ok(())
4815 }
4816 RegistrationWaitOutcome::Exited(exit_report) => {
4817 if let Some(active_child) = child.as_mut() {
4818 active_child.drain_stderr(&spec.module_id).await;
4819 }
4820 *child = None;
4821 handle_reload_child_registration_failure(
4822 spec,
4823 runtime,
4824 registry,
4825 process_liveness,
4826 snapshot,
4827 child,
4828 ReloadRegistrationFailure {
4829 exit_report: registration_failure_exit_report(exit_report),
4830 reason: "new child exited before registering".to_string(),
4831 },
4832 )
4833 .await
4834 }
4835 RegistrationWaitOutcome::TimedOut => {
4836 let mut timed_out_child = child
4837 .take()
4838 .expect("timed-out reload child is still running");
4839 timed_out_child
4840 .start_kill()
4841 .map_err(|source| SuperviseError::Kill {
4842 module_id: spec.module_id.clone(),
4843 source,
4844 })?;
4845 let status = timed_out_child
4846 .wait()
4847 .await
4848 .map_err(|source| SuperviseError::Wait {
4849 module_id: spec.module_id.clone(),
4850 source,
4851 })?;
4852 timed_out_child.drain_stderr(&spec.module_id).await;
4853 handle_reload_child_registration_failure(
4854 spec,
4855 runtime,
4856 registry,
4857 process_liveness,
4858 snapshot,
4859 child,
4860 ReloadRegistrationFailure {
4861 exit_report: registration_failure_exit_report(classify_reaped_child_exit(
4862 snapshot,
4863 &timed_out_child,
4864 &status,
4865 )),
4866 reason: format!(
4867 "new child did not register within {:?}",
4868 REGISTRY_RELEASE_TIMEOUT
4869 ),
4870 },
4871 )
4872 .await
4873 }
4874 }
4875}
4876
4877async fn set_child_enabled(
4878 spec: &ModuleSpec,
4879 runtime: &SupervisorRuntimeConfig,
4880 registry: &Registry,
4881 process_liveness: &SupervisorProcessLiveness,
4882 snapshot: &SharedSnapshot,
4883 child: &mut Option<SupervisedChild>,
4884 enabled: bool,
4885) -> Result<bool, SuperviseError> {
4886 let (current_enabled, current_state) = {
4887 let state = lock_snapshot(snapshot)?;
4888 (state.enabled, state.state)
4889 };
4890 let revive_terminal = enabled
4898 && current_enabled
4899 && child.is_none()
4900 && matches!(current_state, ModuleState::Failed | ModuleState::Stopped);
4901 if current_enabled == enabled && !revive_terminal {
4902 return Ok(false);
4903 }
4904
4905 if enabled {
4906 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4907 state.enabled = true;
4908 state.state = ModuleState::Starting;
4909 clear_current_process_facts(state);
4910 })?;
4911 #[cfg(test)]
4912 if runtime.test_seed_stale_facts_before_enable_spawn {
4913 update_snapshot(snapshot, Some(&spec.module_id), |state| {
4914 state.process_alive = true;
4915 state.pid = Some(41);
4916 state.spawned_at_ms = Some(42);
4917 state.spawned_from = Some(PathBuf::from("/spawned/module"));
4918 state.spawned_file_identity = Some(SpawnedFileIdentity {
4919 device: 43,
4920 inode: 44,
4921 });
4922 })?;
4923 }
4924 wait_for_registration_release(registry, &spec.module_id, REGISTRY_RELEASE_TIMEOUT).await?;
4925 reset_restart_count(snapshot, &spec.module_id)?;
4926 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
4927 let next_child = match spawn_and_mark_running(spec, runtime, snapshot) {
4928 Ok(next_child) => next_child,
4929 Err(err) => {
4930 if let Err(state_err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4931 state.state = ModuleState::Failed;
4932 clear_current_process_facts(state);
4933 }) {
4934 error!(module_id = %spec.module_id, error = %state_err, "failed to record enable spawn failure");
4935 }
4936 process_liveness.untrack_if_current(&spec.module_id, snapshot);
4937 return Err(err);
4938 }
4939 };
4940 *child = Some(next_child);
4941 debug!(module_id = %spec.module_id, "supervised module enabled");
4942 Ok(true)
4943 } else {
4944 begin_forwarding_drain_if_configured(
4945 spec,
4946 runtime,
4947 registry,
4948 snapshot,
4949 Some(false),
4950 RouteCloseReason::Disable,
4951 )
4952 .await?;
4953 drain_optional_child(
4954 &spec.module_id,
4955 spec.protocol,
4956 registry,
4957 snapshot,
4958 &runtime.terminal_ring,
4959 &runtime.spawn_events,
4960 child,
4961 runtime.drain_timeout,
4962 ModuleState::Disabled,
4963 Some(false),
4964 )
4965 .await?;
4966 debug!(module_id = %spec.module_id, "supervised module disabled");
4967 Ok(true)
4968 }
4969}
4970
4971async fn on_child_exit(
4972 spec: &ModuleSpec,
4973 policy: RestartPolicy,
4974 registry: &Registry,
4975 snapshot: &SharedSnapshot,
4976 terminal_ring: &Arc<Mutex<TerminalRing>>,
4977 spawn_events: &SpawnEventFeed,
4978 exit_report: ExitReport,
4979) -> NextAction {
4980 match exit_report.kind {
4981 ExitKind::Clean => {
4982 info!(
4983 module_id = %spec.module_id,
4984 exit_code = ?exit_report.code,
4985 exit_signal = ?exit_report.signal,
4986 "supervised module exited cleanly"
4987 );
4988 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
4989 state.state = ModuleState::Stopped;
4990 clear_current_process_facts(state);
4991 state.last_exit = Some(exit_report.clone());
4992 }) {
4993 error!(module_id = %spec.module_id, error = %err, "failed to record clean module exit");
4994 }
4995 record_terminal(
4996 &spec.module_id,
4997 terminal_ring,
4998 spawn_events,
4999 &exit_report,
5000 TerminalDisposition::Stopped,
5001 );
5002 let registration_released = match wait_for_registration_release(
5003 registry,
5004 &spec.module_id,
5005 REGISTRY_RELEASE_TIMEOUT,
5006 )
5007 .await
5008 {
5009 Ok(()) => true,
5010 Err(err) => {
5011 warn!(module_id = %spec.module_id, error = %err, "registration still active after clean exit");
5012 false
5013 }
5014 };
5015 NextAction::Stop {
5016 registration_released,
5017 }
5018 }
5019 ExitKind::Crash => {
5020 warn!(
5021 module_id = %spec.module_id,
5022 exit_code = ?exit_report.code,
5023 exit_signal = ?exit_report.signal,
5024 "supervised module exited abnormally (crash)"
5025 );
5026 let mut restart_schedule = None;
5027 let mut disposition = TerminalDisposition::Disabled;
5028 let mut disposition_detail = None;
5032 let now = Instant::now();
5033 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5034 clear_current_process_facts(state);
5035 state.last_exit = Some(exit_report.clone());
5036 if state.enabled {
5037 if let Some(schedule) = state.next_crash_restart(&policy, now) {
5038 state.state = ModuleState::Restarting;
5039 restart_schedule = Some(schedule);
5040 disposition = TerminalDisposition::Restarting;
5041 } else {
5042 state.state = ModuleState::Failed;
5043 disposition = TerminalDisposition::Failed;
5044 disposition_detail = Some(policy.budget_exhausted_detail());
5045 }
5046 } else {
5047 state.state = ModuleState::Disabled;
5048 disposition = TerminalDisposition::Disabled;
5049 }
5050 }) {
5051 error!(module_id = %spec.module_id, error = %err, "failed to record crashed module exit");
5052 return NextAction::Stop {
5053 registration_released: false,
5054 };
5055 }
5056 if disposition_detail.is_some() {
5057 error!(
5062 module_id = %spec.module_id,
5063 max_restarts = policy.max_restarts,
5064 window_secs = policy.window.as_secs(),
5065 "module stopped: {}",
5066 policy.budget_exhausted_detail()
5067 );
5068 }
5069 record_terminal_with_detail(
5070 &spec.module_id,
5071 terminal_ring,
5072 spawn_events,
5073 &exit_report,
5074 disposition,
5075 disposition_detail,
5076 );
5077
5078 if let Some(schedule) = restart_schedule {
5079 NextAction::Restart {
5080 schedule: Some(schedule),
5081 }
5082 } else {
5083 let registration_released = match wait_for_registration_release(
5084 registry,
5085 &spec.module_id,
5086 REGISTRY_RELEASE_TIMEOUT,
5087 )
5088 .await
5089 {
5090 Ok(()) => true,
5091 Err(err) => {
5092 warn!(module_id = %spec.module_id, error = %err, "registration still active after failed module");
5093 false
5094 }
5095 };
5096 NextAction::Stop {
5097 registration_released,
5098 }
5099 }
5100 }
5101 ExitKind::DeliberateSeverance => {
5102 warn!(
5103 module_id = %spec.module_id,
5104 exit_code = ?exit_report.code,
5105 exit_signal = ?exit_report.signal,
5106 "supervised module exited after deliberate connection severance"
5107 );
5108 let mut should_restart = false;
5109 let mut disposition = TerminalDisposition::Disabled;
5110 if let Err(err) = update_snapshot(snapshot, Some(&spec.module_id), |state| {
5111 clear_current_process_facts(state);
5112 state.last_exit = Some(exit_report.clone());
5113 state.lifetime_restarts += 1;
5114 if state.enabled {
5115 state.state = ModuleState::Restarting;
5116 should_restart = true;
5117 disposition = TerminalDisposition::Restarting;
5118 } else {
5119 state.state = ModuleState::Disabled;
5120 }
5121 }) {
5122 error!(module_id = %spec.module_id, error = %err, "failed to record deliberately severed module exit");
5123 return NextAction::Stop {
5124 registration_released: false,
5125 };
5126 }
5127 record_terminal(
5128 &spec.module_id,
5129 terminal_ring,
5130 spawn_events,
5131 &exit_report,
5132 disposition,
5133 );
5134
5135 if should_restart {
5136 NextAction::Restart { schedule: None }
5137 } else {
5138 let registration_released = match wait_for_registration_release(
5139 registry,
5140 &spec.module_id,
5141 REGISTRY_RELEASE_TIMEOUT,
5142 )
5143 .await
5144 {
5145 Ok(()) => true,
5146 Err(err) => {
5147 warn!(module_id = %spec.module_id, error = %err, "registration still active after deliberately severed module exit");
5148 false
5149 }
5150 };
5151 NextAction::Stop {
5152 registration_released,
5153 }
5154 }
5155 }
5156 }
5157}
5158
5159fn record_wait_error_terminal(
5160 module_id: &str,
5161 terminal_ring: &Arc<Mutex<TerminalRing>>,
5162 spawn_events: &SpawnEventFeed,
5163) {
5164 record_terminal(
5165 module_id,
5166 terminal_ring,
5167 spawn_events,
5168 &wait_error_exit_report(),
5169 TerminalDisposition::Failed,
5170 );
5171}
5172
5173fn record_terminal(
5174 module_id: &str,
5175 terminal_ring: &Arc<Mutex<TerminalRing>>,
5176 spawn_events: &SpawnEventFeed,
5177 exit_report: &ExitReport,
5178 disposition: TerminalDisposition,
5179) {
5180 record_terminal_with_detail(
5181 module_id,
5182 terminal_ring,
5183 spawn_events,
5184 exit_report,
5185 disposition,
5186 None,
5187 );
5188}
5189
5190fn durable_terminal_history_of(
5194 terminal_ring: &Mutex<TerminalRing>,
5195 module_id: &str,
5196) -> subc_control::TerminalHistory {
5197 let read = terminal_ring
5198 .lock()
5199 .unwrap_or_else(|p| p.into_inner())
5200 .capture_durable_history();
5201 read.read(module_id)
5202}
5203
5204fn record_terminal_with_detail(
5205 module_id: &str,
5206 terminal_ring: &Arc<Mutex<TerminalRing>>,
5207 spawn_events: &SpawnEventFeed,
5208 exit_report: &ExitReport,
5209 disposition: TerminalDisposition,
5210 disposition_detail: Option<String>,
5211) {
5212 spawn_events.emit_exited(module_id, exit_report.code, exit_report.signal);
5213 let mut ring = terminal_ring
5214 .lock()
5215 .unwrap_or_else(|poisoned| poisoned.into_inner());
5216 let record = TerminalRecord {
5217 exit_code: exit_report.code,
5218 exit_signal: exit_report.signal,
5219 at_ms: exit_report.at_ms,
5220 disposition,
5221 exit_kind: exit_report.kind.into(),
5222 disposition_detail,
5223 };
5224 ring.append_journal(module_id, &record);
5225 ring.push(record);
5226}
5227
5228fn untrack_if_registration_released(
5229 process_liveness: &SupervisorProcessLiveness,
5230 registry: &Registry,
5231 module_id: &str,
5232 snapshot: &SharedSnapshot,
5233) {
5234 match registry.get_module(module_id) {
5235 Ok(None) => process_liveness.untrack_if_current(module_id, snapshot),
5236 Ok(Some(_)) => {}
5237 Err(err) => {
5238 warn!(module_id, error = %err, "could not determine whether supervisor liveness can be untracked");
5239 }
5240 }
5241}
5242
5243#[cfg(test)]
5257fn apply_wire_spawn_args(
5258 command: &mut Command,
5259 spec: &ModuleSpec,
5260 connection_file_path: Option<&std::path::Path>,
5261 handle: Option<&SupervisorHandle>,
5262) -> Result<(), SuperviseError> {
5263 apply_wire_spawn_args_for_role(
5264 command,
5265 spec,
5266 connection_file_path,
5267 handle,
5268 SpawnRole::Plain,
5269 )
5270}
5271
5272fn apply_wire_spawn_args_for_role(
5281 command: &mut Command,
5282 spec: &ModuleSpec,
5283 connection_file_path: Option<&std::path::Path>,
5284 handle: Option<&SupervisorHandle>,
5285 role: SpawnRole,
5286) -> Result<(), SuperviseError> {
5287 command.env(SUBC_MODULE_ID_ENV, &spec.module_id);
5288 if spec.protocol == ModuleProtocol::None {
5289 return Ok(());
5290 }
5291 if let Some(connection_file_path) = connection_file_path {
5292 command.arg(SUBC_ARG).arg(connection_file_path);
5293 }
5294
5295 let nonce = generate_launch_nonce()?;
5299 if let Some(handle) = handle {
5300 match role {
5301 SpawnRole::Plain => {
5302 handle.set_spawn_nonce(&spec.module_id, nonce.clone());
5303 if spec.reserved {
5304 handle.set_reserved_nonce(&spec.module_id, nonce.clone());
5305 }
5306 }
5307 SpawnRole::SwapCandidate => handle.open_swap(&spec.module_id, nonce.clone()),
5308 }
5309 }
5310 command.env(SUBC_LAUNCH_NONCE_ENV, nonce);
5311 Ok(())
5312}
5313
5314fn apply_child_env(command: &mut Command, spec: &ModuleSpec) {
5315 command.env_remove(CK_LOG_ENV);
5316 command.env_remove(SUBC_SPAWN_ROLE_ENV);
5323 for (key, value) in &spec.env {
5324 if matches!(
5328 key.as_str(),
5329 CAPTURE_MAX_FILE_MB_ENV | CAPTURE_KEEP_ENV | CAPTURE_MAX_AGE_DAYS_ENV
5330 ) || key == SUBC_SPAWN_ROLE_ENV
5331 {
5332 continue;
5333 }
5334 command.env(key, value);
5335 }
5336}
5337
5338#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5341enum SpawnRole {
5342 Plain,
5343 SwapCandidate,
5344}
5345
5346fn apply_spawn_role(command: &mut Command, role: SpawnRole) {
5349 if role == SpawnRole::SwapCandidate {
5350 command.env(SUBC_SPAWN_ROLE_ENV, SPAWN_ROLE_SWAP_CANDIDATE);
5351 }
5352}
5353
5354fn spawn_child(
5355 spec: &ModuleSpec,
5356 connection_file_path: Option<&std::path::Path>,
5357 handle: Option<&SupervisorHandle>,
5358 ring: &Arc<Mutex<StderrRing>>,
5359 capture_logs_dir: Option<&std::path::Path>,
5360 roster: &ChildRoster,
5361 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5362) -> Result<SupervisedChild, SuperviseError> {
5363 spawn_child_in_slot(
5364 spec,
5365 connection_file_path,
5366 handle,
5367 ring,
5368 capture_logs_dir,
5369 roster,
5370 #[cfg(target_os = "linux")]
5371 cgroup_placement,
5372 SpawnRole::Plain,
5373 false,
5374 )
5375}
5376
5377#[allow(clippy::too_many_arguments)]
5390fn spawn_child_in_slot(
5391 spec: &ModuleSpec,
5392 connection_file_path: Option<&std::path::Path>,
5393 handle: Option<&SupervisorHandle>,
5394 ring: &Arc<Mutex<StderrRing>>,
5395 capture_logs_dir: Option<&std::path::Path>,
5396 roster: &ChildRoster,
5397 #[cfg(target_os = "linux")] cgroup_placement: Option<&subc_cgroup::Placement>,
5398 role: SpawnRole,
5399 alternate_slot: bool,
5400) -> Result<SupervisedChild, SuperviseError> {
5401 if roster.is_closed() {
5402 return Err(SuperviseError::Spawn {
5403 program: spec.program.clone(),
5404 source: io::Error::other("the daemon is shutting down; not starting a new process"),
5405 cgroup_path: None,
5406 });
5407 }
5408 #[cfg(target_os = "linux")]
5409 let cgroup_name = swap::cgroup_name(&spec.module_id, alternate_slot);
5410 #[cfg(not(target_os = "linux"))]
5411 let _ = alternate_slot;
5412 let mut command = Command::new(&spec.program);
5413 command.args(&spec.args);
5414 apply_child_env(&mut command, spec);
5444 apply_spawn_role(&mut command, role);
5445 apply_wire_spawn_args_for_role(&mut command, spec, connection_file_path, handle, role)?;
5446
5447 #[cfg(target_os = "linux")]
5448 let cgroup_path = cgroup_placement
5449 .map(|placement| placement.module_path(&cgroup_name))
5450 .transpose()
5451 .map_err(|source| SuperviseError::Cgroup {
5452 module_id: spec.module_id.clone(),
5453 source,
5454 })?;
5455 #[cfg(not(target_os = "linux"))]
5456 let cgroup_path: Option<PathBuf> = None;
5457 #[cfg(target_os = "linux")]
5458 if let Some(path) = &cgroup_path {
5459 if let Err(error) = apply_cgroup_placement(&mut command, spec, path) {
5460 if let Some(placement) = cgroup_placement {
5461 remove_module_cgroup(placement, &cgroup_name);
5462 }
5463 return Err(error);
5464 }
5465 }
5466
5467 let output_sink = if let Some(logs_dir) = capture_logs_dir {
5468 let path = logs_dir.join(format!("{}.stderr.log", spec.module_id));
5469 match ChildOutputSink::open(&path, capture_retention(spec)) {
5470 Ok(sink) => sink,
5471 Err(error) => {
5472 warn!(
5473 module_id = %spec.module_id,
5474 path = %path.display(),
5475 error = %error,
5476 "could not open child output capture file; forwarding to stderr"
5477 );
5478 ChildOutputSink::Stderr
5479 }
5480 }
5481 } else {
5482 ChildOutputSink::Stderr
5483 };
5484
5485 command.stdout(Stdio::piped());
5486 command.stderr(Stdio::piped());
5487 command.kill_on_drop(true);
5488 #[cfg(unix)]
5505 command.process_group(0);
5506 command.stdin(Stdio::null());
5507 let mut child = match command.spawn() {
5508 Ok(child) => child,
5509 Err(source) => {
5510 #[cfg(target_os = "linux")]
5511 if let Some(placement) = cgroup_placement {
5512 remove_module_cgroup(placement, &cgroup_name);
5513 }
5514 return Err(SuperviseError::Spawn {
5515 program: spec.program.clone(),
5516 source,
5517 cgroup_path,
5518 });
5519 }
5520 };
5521 let spawned_at_ms = unix_ms_now();
5522 let spawned_from = spec.program.clone();
5523 let spawned_file_identity = spawned_file_identity(&spawned_from);
5524 let pid = child.id().ok_or_else(|| SuperviseError::Spawn {
5525 program: spec.program.clone(),
5526 source: io::Error::other("spawned child exposed no live pid"),
5527 cgroup_path: cgroup_path.clone(),
5528 })?;
5529 let process_start_time = crate::provenance::process_start_time(pid);
5530 let process_identity = process_start_time.map(|start_time| ProcessIdentity { pid, start_time });
5531 let roster_guard = roster.admit(
5532 spec.module_id.clone(),
5533 pid,
5534 spec.protocol,
5535 process_start_time,
5536 );
5537
5538 let stdout_pump = match child.stdout.take() {
5539 Some(stdout) => Some(tokio::spawn(pump_stdout_to(stdout, output_sink.clone()))),
5540 None => {
5541 warn!(
5542 module_id = %spec.module_id,
5543 "spawned child exposed no stdout pipe; file capture will be incomplete"
5544 );
5545 None
5546 }
5547 };
5548 let stderr_pump = match child.stderr.take() {
5549 Some(stderr) => {
5550 ring.lock()
5551 .unwrap_or_else(|poisoned| poisoned.into_inner())
5552 .push_process_start();
5553 Some(tokio::spawn(pump_stderr_to(
5554 stderr,
5555 Arc::clone(ring),
5556 output_sink,
5557 )))
5558 }
5559 None => {
5560 ring.lock()
5564 .unwrap_or_else(|poisoned| poisoned.into_inner())
5565 .mark_not_captured("stderr pipe was not available on spawn");
5566 warn!(
5567 module_id = %spec.module_id,
5568 "spawned child exposed no stderr pipe; tail will be unavailable"
5569 );
5570 None
5571 }
5572 };
5573
5574 Ok(SupervisedChild {
5575 child,
5576 #[cfg(target_os = "linux")]
5577 module_id: cgroup_name,
5578 #[cfg(target_os = "linux")]
5579 cgroup_placement: cgroup_placement.cloned(),
5580 stdout_pump,
5581 stderr_pump,
5582 stderr_ring: Arc::clone(ring),
5583 spawned_at_ms,
5584 spawned_from,
5585 spawned_file_identity,
5586 process_start_time,
5587 process_identity,
5588 pid,
5589 roster_guard: Some(roster_guard),
5590 })
5591}
5592
5593#[cfg(target_os = "linux")]
5594fn remove_module_cgroup(placement: &subc_cgroup::Placement, module_id: &str) {
5595 match placement.remove_module(module_id) {
5596 Ok(()) => debug!(module_id, "removed module cgroup after process exit"),
5597 Err(error) => warn!(
5598 module_id,
5599 error = %error,
5600 "could not remove module cgroup after process exit; continuing teardown"
5601 ),
5602 }
5603}
5604
5605#[cfg(target_os = "linux")]
5606fn apply_cgroup_placement(
5607 command: &mut Command,
5608 spec: &ModuleSpec,
5609 path: &std::path::Path,
5610) -> Result<(), SuperviseError> {
5611 subc_cgroup::apply(command, path).map_err(|source| SuperviseError::Cgroup {
5612 module_id: spec.module_id.clone(),
5613 source,
5614 })
5615}
5616
5617fn capture_retention(spec: &ModuleSpec) -> Retention {
5618 let defaults = Retention::default();
5619 let value = |name: &str| {
5620 spec.env
5621 .iter()
5622 .rev()
5623 .find_map(|(key, value)| (key == name).then_some(value.as_str()))
5624 };
5625 Retention {
5626 max_file_mb: value(CAPTURE_MAX_FILE_MB_ENV)
5627 .and_then(|value| value.parse().ok())
5628 .unwrap_or(defaults.max_file_mb),
5629 keep: value(CAPTURE_KEEP_ENV)
5630 .and_then(|value| value.parse().ok())
5631 .unwrap_or(defaults.keep),
5632 max_age_days: value(CAPTURE_MAX_AGE_DAYS_ENV)
5633 .and_then(|value| value.parse().ok())
5634 .unwrap_or(defaults.max_age_days),
5635 }
5636}
5637
5638fn generate_launch_nonce() -> Result<String, SuperviseError> {
5641 let mut bytes = [0u8; 32];
5642 getrandom::getrandom(&mut bytes).map_err(|source| SuperviseError::LaunchNonce {
5643 reason: source.to_string(),
5644 })?;
5645 let mut hex = String::with_capacity(64);
5646 for b in bytes {
5647 use std::fmt::Write;
5648 let _ = write!(hex, "{b:02x}");
5649 }
5650 Ok(hex)
5651}
5652
5653fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
5656 if a.len() != b.len() {
5657 return false;
5658 }
5659 let mut diff = 0u8;
5660 for (x, y) in a.iter().zip(b.iter()) {
5661 diff |= x ^ y;
5662 }
5663 diff == 0
5664}
5665
5666fn spawn_and_mark_running(
5667 spec: &ModuleSpec,
5668 runtime: &SupervisorRuntimeConfig,
5669 snapshot: &SharedSnapshot,
5670) -> Result<SupervisedChild, SuperviseError> {
5671 let child = spawn_child(
5672 spec,
5673 runtime.connection_file_path.as_deref(),
5674 runtime.supervisor_handle.as_ref(),
5675 &runtime.stderr_ring,
5676 runtime.capture_logs_dir.as_deref(),
5677 &runtime.child_roster,
5678 #[cfg(target_os = "linux")]
5679 runtime.cgroup_placement.as_ref(),
5680 )?;
5681 set_running(snapshot, &child, &spec.module_id, &runtime.spawn_events)?;
5682 Ok(child)
5683}
5684
5685enum RegistrationWaitOutcome {
5686 Registered,
5687 Exited(ExitReport),
5688 TimedOut,
5689}
5690
5691struct ReloadRegistrationFailure {
5692 exit_report: ExitReport,
5693 reason: String,
5694}
5695
5696#[derive(Debug, Clone, Copy, PartialEq, Eq)]
5697enum BusyGaugeObservation {
5698 Quiescent,
5699 Busy,
5700 Omitted,
5701}
5702
5703fn busy_gauge_observation(metrics: Option<&Value>, gauges: &[String]) -> BusyGaugeObservation {
5704 let Some(metrics) = metrics.and_then(Value::as_object) else {
5705 return BusyGaugeObservation::Omitted;
5706 };
5707 let mut sum = 0u128;
5708 for gauge in gauges {
5709 let Some(value) = metrics.get(gauge) else {
5710 return BusyGaugeObservation::Omitted;
5711 };
5712 let Some(value) = value.as_u64() else {
5713 return BusyGaugeObservation::Busy;
5714 };
5715 sum = sum.saturating_add(u128::from(value));
5716 }
5717 if sum == 0 {
5718 BusyGaugeObservation::Quiescent
5719 } else {
5720 BusyGaugeObservation::Busy
5721 }
5722}
5723
5724fn declared_busy_gauges(
5725 registry: &Registry,
5726 module_id: &str,
5727) -> Result<Vec<String>, SuperviseError> {
5728 busy_gauges_of(
5729 registry
5730 .get_module(module_id)
5731 .map_err(SuperviseError::Registry)?,
5732 )
5733}
5734
5735fn declared_busy_gauges_for_connection(
5739 registry: &Registry,
5740 connection_id: ConnectionId,
5741) -> Result<Vec<String>, SuperviseError> {
5742 busy_gauges_of(
5743 registry
5744 .get_module_by_connection(connection_id)
5745 .map_err(SuperviseError::Registry)?,
5746 )
5747}
5748
5749fn busy_gauges_of(
5750 registration: Option<crate::registry::ModuleRegistration>,
5751) -> Result<Vec<String>, SuperviseError> {
5752 let Some(registration) = registration else {
5753 return Ok(Vec::new());
5754 };
5755 let Some(self_signals) = registration.manifest.self_signals else {
5756 return Ok(Vec::new());
5757 };
5758
5759 let mut gauges = Vec::new();
5760 for declaration in self_signals {
5761 if declaration.kind != SelfSignalKind::Busy {
5762 continue;
5763 }
5764 match declaration.anchored_to {
5765 SignalAnchor::HealthGauges { gauges: declared } if !declared.is_empty() => {
5766 gauges.extend(declared)
5767 }
5768 _ => {
5769 gauges.push(String::new());
5772 }
5773 }
5774 }
5775 Ok(gauges)
5776}
5777
5778async fn wait_for_forwarding_quiescence(
5783 forwarding: &ForwardingTable,
5784 module_id: &str,
5785 runtime: &SupervisorRuntimeConfig,
5786 endpoint: crate::ModuleEndpointId,
5787 deadline: Instant,
5788 busy_gauges: &[String],
5789 scope: DrainScope,
5790) -> Result<bool, SuperviseError> {
5791 let mut gauges_quiescent = busy_gauges.is_empty();
5792 let mut next_probe_at = Instant::now();
5793 let mut omission_counted = false;
5794
5795 loop {
5796 let now = Instant::now();
5797 if !busy_gauges.is_empty() && now >= next_probe_at && now < deadline {
5798 let report = match scope {
5799 DrainScope::Active => probe_module_health(module_id, runtime, Some(deadline)).await,
5800 DrainScope::Endpoint(endpoint) => {
5801 probe_endpoint_health(endpoint, runtime, Some(deadline)).await
5802 }
5803 };
5804 gauges_quiescent = match report {
5805 Ok(report) => match busy_gauge_observation(report.metrics.as_ref(), busy_gauges) {
5806 BusyGaugeObservation::Quiescent => true,
5807 BusyGaugeObservation::Busy => false,
5808 BusyGaugeObservation::Omitted => {
5809 if !omission_counted {
5810 forwarding
5811 .counters()
5812 .increment_drains_with_undeclared_gauge();
5813 omission_counted = true;
5814 }
5815 false
5816 }
5817 },
5818 Err(err) => {
5819 warn!(
5820 module_id,
5821 error = %err,
5822 "drain health.check did not produce declared busy gauges; treating module as busy"
5823 );
5824 false
5825 }
5826 };
5827 next_probe_at = Instant::now() + runtime.health.cadence.max(REGISTRY_RELEASE_POLL);
5828 }
5829
5830 let in_flight = forwarding
5831 .endpoint_in_flight_count(endpoint)
5832 .map_err(SuperviseError::Forwarding)?;
5833 if in_flight == 0 && gauges_quiescent {
5834 return Ok(true);
5835 }
5836
5837 let now = Instant::now();
5838 if now >= deadline {
5839 return Ok(false);
5840 }
5841 let mut wait = deadline
5842 .saturating_duration_since(now)
5843 .min(REGISTRY_RELEASE_POLL);
5844 if !busy_gauges.is_empty() {
5845 wait = wait.min(next_probe_at.saturating_duration_since(now));
5846 }
5847 sleep(wait).await;
5848 }
5849}
5850
5851fn drained_after_quiescence_wait(wait_result: &Result<bool, SuperviseError>) -> bool {
5859 match wait_result {
5860 Ok(drained) => *drained,
5861 Err(_) => false,
5862 }
5863}
5864
5865fn send_route_goodbyes(forwarding: &ForwardingTable, released_routes: Vec<GoodbyeTarget>) {
5866 for released in released_routes {
5867 let frame = match Frame::build_with_version(
5868 released.negotiated_ver,
5869 FrameType::Goodbye,
5870 control_flags(),
5871 released.channel,
5872 released.epoch,
5873 0,
5874 Vec::new(),
5875 ) {
5876 Ok(frame) => frame,
5877 Err(err) => {
5878 warn!(
5879 route_channel = released.channel,
5880 error = %err,
5881 "failed to build supervisor drain route GOODBYE frame"
5882 );
5883 continue;
5884 }
5885 };
5886 if let Err(err) = released.sink.try_send(frame) {
5887 if released.close_on_delivery_failure() {
5888 warn!(
5889 target_connection_id = released.connection_id.get(),
5890 route_channel = released.channel,
5891 error = %err,
5892 "supervisor drain route GOODBYE was not delivered to client; closing target connection"
5893 );
5894 let _ = forwarding.escalate_client_delivery_failure(
5895 released.connection_id,
5896 released.channel,
5897 released.epoch,
5898 CloseReason::new(
5899 "route_goodbye_delivery_failed",
5900 format!(
5901 "failed to enqueue supervisor drain route GOODBYE for channel {}: {err}",
5902 released.channel
5903 ),
5904 ),
5905 crate::forwarding::UndeliveredFrame {
5906 module_id: released.module_id.as_deref(),
5907 sink: &released.sink,
5908 },
5909 );
5910 } else {
5911 warn!(
5912 target_connection_id = released.connection_id.get(),
5913 route_channel = released.channel,
5914 error = %err,
5915 "supervisor drain route GOODBYE to module dropped under backpressure; not closing shared module connection"
5916 );
5917 }
5918 }
5919 }
5920}
5921
5922fn send_module_draining(
5923 module_id: &str,
5924 reason: RouteCloseReason,
5925 deadline_ms: u64,
5926 target: &ModuleDrainTarget,
5927) {
5928 let body = match serde_json::to_vec(&ModuleControlCommand::Draining {
5929 reason,
5930 deadline_ms,
5931 }) {
5932 Ok(body) => body,
5933 Err(err) => {
5934 warn!(
5935 module_id,
5936 error = %err,
5937 "failed to encode module draining command"
5938 );
5939 return;
5940 }
5941 };
5942 let frame = match Frame::build_with_version(
5943 target.negotiated_ver,
5944 FrameType::Push,
5945 control_flags(),
5946 0,
5947 0,
5948 0,
5949 body,
5950 ) {
5951 Ok(frame) => frame,
5952 Err(err) => {
5953 warn!(
5954 module_id,
5955 error = %err,
5956 "failed to build module draining command frame"
5957 );
5958 return;
5959 }
5960 };
5961 if let Err(err) = target.sink.try_send(frame) {
5962 warn!(
5963 module_id,
5964 target_connection_id = target.endpoint.connection_id.get(),
5965 error = %err,
5966 "module draining command was not delivered to peer"
5967 );
5968 }
5969}
5970
5971fn send_module_goodbye(module_id: &str, forwarding: &ForwardingTable, target: &ModuleDrainTarget) {
5972 let frame = match Frame::build_with_version(
5973 target.negotiated_ver,
5974 FrameType::Goodbye,
5975 control_flags(),
5976 0,
5977 0,
5978 0,
5979 Vec::new(),
5980 ) {
5981 Ok(frame) => frame,
5982 Err(err) => {
5983 warn!(
5984 module_id,
5985 error = %err,
5986 "failed to build supervisor drain module GOODBYE frame"
5987 );
5988 return;
5989 }
5990 };
5991 if let Err(err) = target.sink.try_send(frame) {
5992 warn!(
5993 module_id,
5994 target_connection_id = target.endpoint.connection_id.get(),
5995 error = %err,
5996 "supervisor drain module GOODBYE was not delivered to peer; closing module connection"
5997 );
5998 forwarding.request_connection_close(
5999 target.endpoint.connection_id,
6000 CloseReason::new(
6001 "module_goodbye_delivery_failed",
6002 format!("failed to enqueue supervisor drain module GOODBYE for module '{module_id}': {err}"),
6003 ),
6004 );
6005 }
6006}
6007
6008#[derive(Clone, Copy)]
6009struct ForwardingDrainContext<'a> {
6010 spec: &'a ModuleSpec,
6011 runtime: &'a SupervisorRuntimeConfig,
6012 registry: &'a Registry,
6013 scope: DrainScope,
6014}
6015
6016#[derive(Debug, Clone, Copy, PartialEq, Eq)]
6018enum DrainScope {
6019 Active,
6022 Endpoint(crate::ModuleEndpointId),
6027}
6028
6029async fn begin_forwarding_drain(
6030 spec: &ModuleSpec,
6031 runtime: &SupervisorRuntimeConfig,
6032 registry: &Registry,
6033 snapshot: &SharedSnapshot,
6034 enabled: Option<bool>,
6035 reason: RouteCloseReason,
6036) -> Result<(), SuperviseError> {
6037 let Some(forwarding) = runtime.forwarding.as_ref() else {
6038 return Err(SuperviseError::ReloadUnavailable {
6039 module_id: spec.module_id.clone(),
6040 reason: "supervisor was not configured with a forwarding table".to_string(),
6041 });
6042 };
6043
6044 begin_forwarding_drain_with(
6045 forwarding,
6046 ForwardingDrainContext {
6047 spec,
6048 runtime,
6049 registry,
6050 scope: DrainScope::Active,
6051 },
6052 snapshot,
6053 enabled,
6054 reason,
6055 runtime.drain_timeout,
6056 )
6057 .await
6058}
6059
6060async fn begin_forwarding_drain_if_configured(
6061 spec: &ModuleSpec,
6062 runtime: &SupervisorRuntimeConfig,
6063 registry: &Registry,
6064 snapshot: &SharedSnapshot,
6065 enabled: Option<bool>,
6066 reason: RouteCloseReason,
6067) -> Result<(), SuperviseError> {
6068 begin_forwarding_drain_with_timeout(
6069 spec,
6070 runtime,
6071 registry,
6072 snapshot,
6073 enabled,
6074 reason,
6075 runtime.drain_timeout,
6076 )
6077 .await
6078}
6079
6080async fn begin_forwarding_drain_with_timeout(
6084 spec: &ModuleSpec,
6085 runtime: &SupervisorRuntimeConfig,
6086 registry: &Registry,
6087 snapshot: &SharedSnapshot,
6088 enabled: Option<bool>,
6089 reason: RouteCloseReason,
6090 drain_timeout: Duration,
6091) -> Result<(), SuperviseError> {
6092 let Some(forwarding) = runtime.forwarding.as_ref() else {
6093 return Ok(());
6094 };
6095
6096 begin_forwarding_drain_with(
6097 forwarding,
6098 ForwardingDrainContext {
6099 spec,
6100 runtime,
6101 registry,
6102 scope: DrainScope::Active,
6103 },
6104 snapshot,
6105 enabled,
6106 reason,
6107 drain_timeout,
6108 )
6109 .await
6110}
6111
6112async fn begin_forwarding_drain_with(
6113 forwarding: &ForwardingTable,
6114 context: ForwardingDrainContext<'_>,
6115 snapshot: &SharedSnapshot,
6116 enabled: Option<bool>,
6117 reason: RouteCloseReason,
6118 drain_timeout: Duration,
6119) -> Result<(), SuperviseError> {
6120 let ForwardingDrainContext {
6121 spec,
6122 runtime,
6123 registry,
6124 scope,
6125 } = context;
6126 debug_assert_ne!(reason, RouteCloseReason::Crash);
6127 let terminal = matches!(reason, RouteCloseReason::Disable);
6128 let drain_started_at = Instant::now();
6129 let drain_deadline = drain_started_at + drain_timeout;
6130 let deadline_ms =
6131 unix_ms_now().saturating_add(u64::try_from(drain_timeout.as_millis()).unwrap_or(u64::MAX));
6132 let busy_gauges = match scope {
6133 DrainScope::Active => declared_busy_gauges(registry, &spec.module_id)?,
6134 DrainScope::Endpoint(endpoint) => {
6135 declared_busy_gauges_for_connection(registry, endpoint.connection_id)?
6136 }
6137 };
6138
6139 let drain_target = match scope {
6142 DrainScope::Active => forwarding.begin_module_drain(&spec.module_id, reason),
6143 DrainScope::Endpoint(endpoint) => forwarding.begin_endpoint_drain(endpoint, reason),
6144 }
6145 .map_err(SuperviseError::Forwarding)?;
6146 if scope == DrainScope::Active {
6147 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6148 state.state = ModuleState::Draining;
6149 if let Some(enabled) = enabled {
6150 state.enabled = enabled;
6151 }
6152 })?;
6153 }
6154
6155 if let Some(target) = drain_target.as_ref() {
6156 send_module_draining(&spec.module_id, reason, deadline_ms, target);
6157 let routes = forwarding
6158 .endpoint_routes(target.endpoint)
6159 .map_err(SuperviseError::Forwarding)?;
6160 let routes_notified = routes.len();
6161 crate::control::send_route_control_pushes(
6162 forwarding,
6163 routes.clone(),
6164 ClientControlPush::RouteClosing {
6165 module_id: spec.module_id.clone(),
6166 reason,
6167 },
6168 );
6169 send_route_goodbyes(forwarding, target.abandoned_bindings.clone());
6170
6171 let wait_result = wait_for_forwarding_quiescence(
6177 forwarding,
6178 &spec.module_id,
6179 runtime,
6180 target.endpoint,
6181 drain_deadline,
6182 &busy_gauges,
6183 scope,
6184 )
6185 .await;
6186 let drained = drained_after_quiescence_wait(&wait_result);
6187 if let Err(err) = &wait_result {
6188 error!(
6189 module_id = %spec.module_id,
6190 ?reason,
6191 error = %err,
6192 "forwarding quiescence wait failed after route.closing; forcing route.closed(drained: false) so the client is not left waiting on an unfulfilled promise"
6193 );
6194 } else if !drained {
6195 let holdouts = forwarding
6201 .endpoint_drain_holdouts(target.endpoint)
6202 .unwrap_or_default();
6203 warn!(
6204 module_id = %spec.module_id,
6205 waited = ?drain_timeout,
6206 ?reason,
6207 held_requests = holdouts.requests,
6208 held_routes = holdouts.routes,
6209 total_routes = holdouts.total_routes,
6210 top_connections = ?holdouts.top_connections,
6211 held = %holdouts
6214 .held
6215 .iter()
6216 .map(|(channel, corr)| format!("{channel}:{corr}"))
6217 .collect::<Vec<_>>()
6218 .join(","),
6219 "route drain timed out before request quiescence; forcing teardown"
6220 );
6221 }
6222 crate::control::send_route_control_pushes(
6223 forwarding,
6224 routes,
6225 ClientControlPush::RouteClosed {
6226 module_id: spec.module_id.clone(),
6227 reason,
6228 drained,
6229 abandoned: target.abandoned_bindings.len() as u32,
6230 excluded_subscriptions: target.excluded_subscriptions,
6231 terminal: Some(terminal),
6232 },
6233 );
6234 wait_result?;
6235
6236 let released_routes = match forwarding.release_module_endpoint_routes(target.endpoint) {
6242 Ok(routes) => routes,
6243 Err(err) => {
6244 warn!(
6245 module_id = %spec.module_id,
6246 ?reason,
6247 error = %err,
6248 "failed to release module endpoint routes after route.closed; module GOODBYE will still be sent"
6249 );
6250 send_module_goodbye(&spec.module_id, forwarding, target);
6251 return Err(SuperviseError::Forwarding(err));
6252 }
6253 };
6254 let route_goodbye_count = released_routes.len();
6255 send_route_goodbyes(forwarding, released_routes);
6256 send_module_goodbye(&spec.module_id, forwarding, target);
6257
6258 info!(
6264 module_id = %spec.module_id,
6265 ?reason,
6266 routes_notified,
6267 route_goodbyes = route_goodbye_count,
6268 abandoned_reservations = target.abandoned_bindings.len(),
6269 excluded_subscriptions = target.excluded_subscriptions,
6270 drained,
6271 "module drain complete; consumers notified via route.closing/route.closed pushes and per-route GOODBYE frames"
6272 );
6273 }
6274
6275 Ok(())
6276}
6277
6278async fn wait_for_registration_after_reload(
6281 registry: &Registry,
6282 module_id: &str,
6283 snapshot: &SharedSnapshot,
6284 child: &mut SupervisedChild,
6285 wait: Duration,
6286) -> Result<RegistrationWaitOutcome, SuperviseError> {
6287 wait_for_slot_registration(
6288 registry,
6289 crate::registry::RegistrationSlot::Active(module_id),
6290 module_id,
6291 snapshot,
6292 child,
6293 wait,
6294 )
6295 .await
6296}
6297
6298async fn wait_for_slot_registration(
6306 registry: &Registry,
6307 slot: crate::registry::RegistrationSlot<'_>,
6308 module_id: &str,
6309 snapshot: &SharedSnapshot,
6310 child: &mut SupervisedChild,
6311 wait: Duration,
6312) -> Result<RegistrationWaitOutcome, SuperviseError> {
6313 let deadline = Instant::now() + wait;
6314 loop {
6315 if registry
6316 .registration(slot)
6317 .map_err(SuperviseError::Registry)?
6318 .is_some()
6319 {
6320 return Ok(RegistrationWaitOutcome::Registered);
6321 }
6322
6323 let now = Instant::now();
6324 if now >= deadline {
6325 return Ok(RegistrationWaitOutcome::TimedOut);
6326 }
6327 let remaining = deadline.saturating_duration_since(now);
6328 let poll = remaining.min(REGISTRY_RELEASE_POLL);
6329
6330 tokio::select! {
6331 wait_result = child.wait() => {
6332 let status = wait_result.map_err(|source| SuperviseError::Wait {
6333 module_id: module_id.to_string(),
6334 source,
6335 })?;
6336 return Ok(RegistrationWaitOutcome::Exited(classify_reaped_child_exit(
6337 snapshot,
6338 child,
6339 &status,
6340 )));
6341 }
6342 _ = sleep(poll) => {}
6343 }
6344 }
6345}
6346
6347fn registration_failure_exit_report(mut exit_report: ExitReport) -> ExitReport {
6348 if exit_report.kind != ExitKind::DeliberateSeverance {
6351 exit_report.kind = ExitKind::Crash;
6352 }
6353 exit_report
6354}
6355
6356async fn handle_reload_child_registration_failure(
6357 spec: &ModuleSpec,
6358 runtime: &SupervisorRuntimeConfig,
6359 registry: &Registry,
6360 process_liveness: &SupervisorProcessLiveness,
6361 snapshot: &SharedSnapshot,
6362 child: &mut Option<SupervisedChild>,
6363 failure: ReloadRegistrationFailure,
6364) -> Result<(), SuperviseError> {
6365 let ReloadRegistrationFailure {
6366 exit_report,
6367 reason,
6368 } = failure;
6369 match on_child_exit(
6370 spec,
6371 runtime.restart_policy,
6372 registry,
6373 snapshot,
6374 &runtime.terminal_ring,
6375 &runtime.spawn_events,
6376 exit_report,
6377 )
6378 .await
6379 {
6380 NextAction::Stop {
6381 registration_released,
6382 } => {
6383 if registration_released {
6384 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6385 }
6386 }
6387 NextAction::Restart { schedule } => {
6388 let delay = schedule.map_or(runtime.restart_policy.delay_for_restart(0), |schedule| {
6389 schedule.delay
6390 });
6391 if let Some(schedule) = schedule {
6392 log_crash_respawn(&spec.module_id, schedule);
6393 }
6394 sleep(delay).await;
6395 if respawn_still_pending(snapshot) {
6399 if let Err(err) = wait_for_registration_release(
6400 registry,
6401 &spec.module_id,
6402 REGISTRY_RELEASE_TIMEOUT,
6403 )
6404 .await
6405 {
6406 fail_snapshot(snapshot, Some(&spec.module_id), None);
6407 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6408 return Err(SuperviseError::ReloadFailed {
6409 module_id: spec.module_id.clone(),
6410 reason: format!(
6411 "{reason}; registration did not release before policy retry: {err}"
6412 ),
6413 });
6414 }
6415 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6416 match spawn_and_mark_running(spec, runtime, snapshot) {
6417 Ok(next_child) => {
6418 *child = Some(next_child);
6419 }
6420 Err(err) => {
6421 fail_snapshot(snapshot, Some(&spec.module_id), None);
6422 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6423 return Err(SuperviseError::ReloadFailed {
6424 module_id: spec.module_id.clone(),
6425 reason: format!("{reason}; policy retry spawn failed: {err}"),
6426 });
6427 }
6428 }
6429 }
6430 }
6431 }
6432
6433 Err(SuperviseError::ReloadFailed {
6434 module_id: spec.module_id.clone(),
6435 reason,
6436 })
6437}
6438
6439async fn handle_reload_spawn_failure(
6440 spec: &ModuleSpec,
6441 runtime: &SupervisorRuntimeConfig,
6442 process_liveness: &SupervisorProcessLiveness,
6443 snapshot: &SharedSnapshot,
6444 child: &mut Option<SupervisedChild>,
6445 reason: String,
6446) -> Result<(), SuperviseError> {
6447 let mut should_retry = false;
6448 let now = Instant::now();
6449 update_snapshot(snapshot, Some(&spec.module_id), |state| {
6450 clear_current_process_facts(state);
6451 if daemon_will_restart(state, &runtime.restart_policy, now) {
6452 state.record_crash_restart(&runtime.restart_policy, now);
6453 state.state = ModuleState::Restarting;
6454 should_retry = true;
6455 } else if state.enabled {
6456 state.state = ModuleState::Failed;
6457 } else {
6458 state.state = ModuleState::Disabled;
6459 }
6460 })?;
6461
6462 if should_retry {
6463 sleep(runtime.restart_policy.backoff).await;
6464 if respawn_still_pending(snapshot) {
6468 process_liveness.track(spec.module_id.clone(), Arc::clone(snapshot));
6469 match spawn_and_mark_running(spec, runtime, snapshot) {
6470 Ok(next_child) => {
6471 *child = Some(next_child);
6472 }
6473 Err(err) => {
6474 fail_snapshot(snapshot, Some(&spec.module_id), None);
6475 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6476 return Err(SuperviseError::ReloadFailed {
6477 module_id: spec.module_id.clone(),
6478 reason: format!("{reason}; policy retry spawn failed: {err}"),
6479 });
6480 }
6481 }
6482 }
6483 } else {
6484 process_liveness.untrack_if_current(&spec.module_id, snapshot);
6485 }
6486
6487 Err(SuperviseError::ReloadFailed {
6488 module_id: spec.module_id.clone(),
6489 reason,
6490 })
6491}
6492
6493fn control_flags() -> Flags {
6494 Flags::new(false, Priority::Passive, false)
6495}
6496
6497#[allow(clippy::too_many_arguments)]
6498async fn drain_optional_child(
6499 module_id: &str,
6500 protocol: ModuleProtocol,
6501 registry: &Registry,
6502 snapshot: &SharedSnapshot,
6503 terminal_ring: &Arc<Mutex<TerminalRing>>,
6504 spawn_events: &SpawnEventFeed,
6505 child: &mut Option<SupervisedChild>,
6506 drain_timeout: Duration,
6507 final_state: ModuleState,
6508 enabled: Option<bool>,
6509) -> Result<(), SuperviseError> {
6510 if let Some(child) = child.take() {
6511 drain_child_to_state(
6512 module_id,
6513 protocol,
6514 registry,
6515 snapshot,
6516 terminal_ring,
6517 spawn_events,
6518 child,
6519 drain_timeout,
6520 final_state,
6521 enabled,
6522 )
6523 .await
6524 } else {
6525 update_snapshot(snapshot, Some(module_id), |state| {
6526 state.state = final_state;
6527 if let Some(enabled) = enabled {
6528 state.enabled = enabled;
6529 }
6530 clear_current_process_facts(state);
6531 })?;
6532 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
6533 }
6534}
6535
6536#[allow(clippy::too_many_arguments)]
6537async fn drain_child_to_state(
6538 module_id: &str,
6539 protocol: ModuleProtocol,
6540 registry: &Registry,
6541 snapshot: &SharedSnapshot,
6542 terminal_ring: &Arc<Mutex<TerminalRing>>,
6543 spawn_events: &SpawnEventFeed,
6544 mut child: SupervisedChild,
6545 drain_timeout: Duration,
6546 final_state: ModuleState,
6547 enabled: Option<bool>,
6548) -> Result<(), SuperviseError> {
6549 update_snapshot(snapshot, Some(module_id), |state| {
6550 state.state = ModuleState::Draining;
6551 if let Some(enabled) = enabled {
6552 state.enabled = enabled;
6553 }
6554 })?;
6555
6556 if protocol == ModuleProtocol::None {
6562 request_graceful_stop(module_id, &child);
6563 }
6564
6565 let exit_report = match timeout(drain_timeout, child.wait()).await {
6566 Ok(Ok(status)) => classify_reaped_child_exit(snapshot, &child, &status),
6567 Ok(Err(source)) => {
6568 fail_snapshot(snapshot, Some(module_id), None);
6569 return Err(SuperviseError::Wait {
6570 module_id: module_id.to_string(),
6571 source,
6572 });
6573 }
6574 Err(_) => {
6575 child.start_kill().map_err(|source| {
6584 fail_snapshot(snapshot, Some(module_id), None);
6585 SuperviseError::Kill {
6586 module_id: module_id.to_string(),
6587 source,
6588 }
6589 })?;
6590 let status = child.wait().await.map_err(|source| {
6591 fail_snapshot(snapshot, Some(module_id), None);
6592 SuperviseError::Wait {
6593 module_id: module_id.to_string(),
6594 source,
6595 }
6596 })?;
6597 classify_reaped_child_exit(snapshot, &child, &status)
6598 }
6599 };
6600
6601 update_snapshot(snapshot, Some(module_id), |state| {
6602 state.state = final_state;
6603 if let Some(enabled) = enabled {
6604 state.enabled = enabled;
6605 }
6606 clear_current_process_facts(state);
6607 state.last_exit = Some(exit_report.clone());
6608 if exit_report.kind == ExitKind::DeliberateSeverance {
6609 state.lifetime_restarts += 1;
6610 }
6611 })?;
6612 record_terminal(
6613 module_id,
6614 terminal_ring,
6615 spawn_events,
6616 &exit_report,
6617 terminal_disposition(final_state),
6618 );
6619 child.drain_stderr(module_id).await;
6620
6621 wait_for_registration_release(registry, module_id, REGISTRY_RELEASE_TIMEOUT).await
6622}
6623
6624#[cfg(unix)]
6644fn request_graceful_stop(module_id: &str, child: &SupervisedChild) {
6645 let Some(pid) = child
6646 .id()
6647 .and_then(|pid| i32::try_from(pid).ok())
6648 .and_then(rustix::process::Pid::from_raw)
6649 else {
6650 debug!(
6651 module_id,
6652 "no pid to signal for protocol: none teardown; falling through to the drain wait"
6653 );
6654 return;
6655 };
6656 match rustix::process::kill_process(pid, rustix::process::Signal::TERM) {
6657 Ok(()) => debug!(module_id, "sent SIGTERM to protocol: none module"),
6658 Err(err) => debug!(
6659 module_id,
6660 error = %err,
6661 "SIGTERM to protocol: none module failed; the drain wait and kill still apply"
6662 ),
6663 }
6664}
6665
6666#[cfg(not(unix))]
6674fn request_graceful_stop(module_id: &str, _child: &SupervisedChild) {
6675 debug!(
6676 module_id,
6677 "no graceful stop signal exists on this platform; protocol: none teardown waits, then kills"
6678 );
6679}
6680
6681fn terminal_disposition(final_state: ModuleState) -> TerminalDisposition {
6682 match final_state {
6683 ModuleState::Stopped => TerminalDisposition::Stopped,
6684 ModuleState::Disabled => TerminalDisposition::Disabled,
6685 ModuleState::Restarting => TerminalDisposition::Restarting,
6686 ModuleState::Failed => TerminalDisposition::Failed,
6687 ModuleState::Starting
6688 | ModuleState::Running
6689 | ModuleState::Unresponsive
6690 | ModuleState::Draining => {
6691 unreachable!("terminal exits only finish in terminal or restarting states")
6692 }
6693 }
6694}
6695
6696async fn wait_for_registration_release(
6699 registry: &Registry,
6700 module_id: &str,
6701 wait: Duration,
6702) -> Result<(), SuperviseError> {
6703 wait_for_slot_registration_release(
6704 registry,
6705 crate::registry::RegistrationSlot::Active(module_id),
6706 wait,
6707 )
6708 .await
6709}
6710
6711async fn wait_for_slot_registration_release(
6719 registry: &Registry,
6720 slot: crate::registry::RegistrationSlot<'_>,
6721 wait: Duration,
6722) -> Result<(), SuperviseError> {
6723 let deadline = Instant::now() + wait;
6724 let mut release_events = registration_release_events().subscribe();
6725 let still_active = |registration: &crate::registry::ModuleRegistration| {
6726 SuperviseError::RegistrationStillActive {
6727 module_id: registration.manifest.module_id.clone(),
6728 waited: wait,
6729 }
6730 };
6731 loop {
6732 let _observed_generation = *release_events.borrow_and_update();
6733 let Some(registration) = registry
6734 .registration(slot)
6735 .map_err(SuperviseError::Registry)?
6736 else {
6737 return Ok(());
6738 };
6739
6740 let now = Instant::now();
6741 if now >= deadline {
6742 return Err(still_active(®istration));
6743 }
6744
6745 let remaining = deadline.saturating_duration_since(now);
6746 match timeout(remaining, release_events.changed()).await {
6747 Ok(Ok(())) | Ok(Err(_)) => {}
6748 Err(_) => return Err(still_active(®istration)),
6749 }
6750 }
6751}
6752
6753#[cfg(test)]
6754mod slot_registration_wait_tests {
6755 use super::*;
6756 use crate::registry::{ConnectionId, RegistrationSlot};
6757 use subc_protocol::manifest::ModuleManifest;
6758
6759 const INCUMBENT: u64 = 1;
6760 const CANDIDATE: u64 = 2;
6761
6762 fn swapped_registry() -> Arc<Registry> {
6763 let registry = Arc::new(Registry::default());
6764 let manifest = ModuleManifest::builder("m", "0.1.0").build();
6765 registry
6766 .register_with_control_ops(
6767 manifest.clone(),
6768 1,
6769 ConnectionId::new(INCUMBENT),
6770 Vec::new(),
6771 )
6772 .unwrap();
6773 registry
6774 .register_candidate_with_control_ops(
6775 manifest,
6776 1,
6777 ConnectionId::new(CANDIDATE),
6778 Vec::new(),
6779 )
6780 .unwrap();
6781 registry
6782 }
6783
6784 #[tokio::test]
6788 async fn incumbent_release_is_awaited_by_connection_not_by_module_id() {
6789 let registry = swapped_registry();
6790 registry.promote_candidate("m").unwrap().unwrap();
6791
6792 assert!(matches!(
6793 wait_for_registration_release(®istry, "m", Duration::from_millis(50)).await,
6794 Err(SuperviseError::RegistrationStillActive { .. })
6795 ));
6796
6797 assert!(matches!(
6799 wait_for_slot_registration_release(
6800 ®istry,
6801 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
6802 Duration::from_millis(50),
6803 )
6804 .await,
6805 Err(SuperviseError::RegistrationStillActive { .. })
6806 ));
6807
6808 let releaser = Arc::clone(®istry);
6809 let release = tokio::spawn(async move {
6810 sleep(Duration::from_millis(20)).await;
6811 releaser
6812 .deregister_connection(ConnectionId::new(INCUMBENT))
6813 .unwrap();
6814 notify_registration_release();
6815 });
6816 wait_for_slot_registration_release(
6817 ®istry,
6818 RegistrationSlot::Connection(ConnectionId::new(INCUMBENT)),
6819 Duration::from_secs(5),
6820 )
6821 .await
6822 .expect("the incumbent's own registration is released");
6823 release.await.unwrap();
6824 assert!(registry.get_module("m").unwrap().is_some());
6825 }
6826
6827 #[tokio::test]
6830 async fn candidate_slot_wait_ignores_the_incumbents_registration() {
6831 let registry = swapped_registry();
6832 assert!(matches!(
6833 wait_for_slot_registration_release(
6834 ®istry,
6835 RegistrationSlot::Candidate("m"),
6836 Duration::from_millis(50),
6837 )
6838 .await,
6839 Err(SuperviseError::RegistrationStillActive { .. })
6840 ));
6841 registry
6842 .deregister_connection(ConnectionId::new(CANDIDATE))
6843 .unwrap();
6844 wait_for_slot_registration_release(
6845 ®istry,
6846 RegistrationSlot::Candidate("m"),
6847 Duration::from_millis(50),
6848 )
6849 .await
6850 .expect("a candidate slot with no candidate is released");
6851 assert!(registry
6852 .registration(RegistrationSlot::Active("m"))
6853 .unwrap()
6854 .is_some());
6855 }
6856}
6857
6858fn classify_exit(status: &ExitStatus) -> ExitReport {
6859 ExitReport {
6860 kind: if status.success() {
6861 ExitKind::Clean
6862 } else {
6863 ExitKind::Crash
6864 },
6865 code: status.code(),
6866 signal: exit_signal(status),
6867 at_ms: unix_ms_now(),
6868 }
6869}
6870
6871fn wait_error_exit_report() -> ExitReport {
6877 ExitReport {
6878 kind: ExitKind::Crash,
6879 code: None,
6880 signal: None,
6881 at_ms: unix_ms_now(),
6882 }
6883}
6884
6885#[cfg(unix)]
6886fn exit_signal(status: &ExitStatus) -> Option<i32> {
6887 use std::os::unix::process::ExitStatusExt;
6888
6889 status.signal()
6890}
6891
6892#[cfg(not(unix))]
6893fn exit_signal(_status: &ExitStatus) -> Option<i32> {
6894 None
6895}
6896
6897fn reset_restart_count(snapshot: &SharedSnapshot, module_id: &str) -> Result<(), SuperviseError> {
6903 update_snapshot(snapshot, Some(module_id), |state| {
6904 state.clear_crash_restarts();
6905 })
6906}
6907
6908fn set_running(
6909 snapshot: &SharedSnapshot,
6910 child: &SupervisedChild,
6911 module_id: &str,
6912 spawn_events: &SpawnEventFeed,
6913) -> Result<(), SuperviseError> {
6914 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
6915 module_id: Some(module_id.to_string()),
6916 })?;
6917 state.spawn_generation = spawn_events.emit_spawned(module_id, child.pid, child.spawned_at_ms);
6918 state.in_alternate_slot = false;
6921 state.state = ModuleState::Running;
6922 state.enabled = true;
6923 state.process_alive = true;
6924 state.pid = child.id();
6925 state.spawned_at_ms = Some(child.spawned_at_ms);
6926 state.spawned_from = Some(child.spawned_from.clone());
6927 state.spawned_file_identity = child.spawned_file_identity;
6928 state.process_start_time = child.process_start_time;
6929 Ok(())
6930}
6931
6932fn clear_current_process_facts(state: &mut SupervisorSnapshot) {
6933 state.process_alive = false;
6934 state.pid = None;
6935 state.spawned_at_ms = None;
6936 state.spawned_from = None;
6937 state.spawned_file_identity = None;
6938 state.process_start_time = None;
6939 state.deliberate_severance = None;
6940}
6941
6942#[cfg(test)]
6943fn record_deliberate_severance(
6944 snapshot: &SharedSnapshot,
6945 identity: ProcessIdentity,
6946) -> Result<(), SuperviseError> {
6947 update_snapshot(snapshot, None, |state| {
6948 state.deliberate_severance = Some(identity);
6949 })
6950}
6951
6952fn apply_deliberate_severance_marker(
6953 snapshot: &SharedSnapshot,
6954 exited_identity: Option<ProcessIdentity>,
6955 mut exit_report: ExitReport,
6956) -> ExitReport {
6957 let marker = lock_snapshot(snapshot)
6958 .ok()
6959 .and_then(|mut state| state.deliberate_severance.take());
6960 if marker.is_some() && marker == exited_identity {
6961 exit_report.kind = ExitKind::DeliberateSeverance;
6962 }
6963 exit_report
6964}
6965
6966fn classify_reaped_child_exit(
6967 snapshot: &SharedSnapshot,
6968 child: &SupervisedChild,
6969 status: &ExitStatus,
6970) -> ExitReport {
6971 apply_deliberate_severance_marker(snapshot, child.process_identity(), classify_exit(status))
6972}
6973
6974fn fail_snapshot(
6975 snapshot: &SharedSnapshot,
6976 module_id: Option<&str>,
6977 last_exit: Option<ExitReport>,
6978) {
6979 if let Err(err) = update_snapshot(snapshot, module_id, |state| {
6980 state.state = ModuleState::Failed;
6981 clear_current_process_facts(state);
6982 if let Some(last_exit) = last_exit {
6983 state.last_exit = Some(last_exit);
6984 }
6985 }) {
6986 error!(error = %err, "failed to mark supervisor state failed");
6987 }
6988}
6989
6990fn update_snapshot(
6991 snapshot: &SharedSnapshot,
6992 module_id: Option<&str>,
6993 update: impl FnOnce(&mut SupervisorSnapshot),
6994) -> Result<(), SuperviseError> {
6995 let mut state = snapshot.lock().map_err(|_| SuperviseError::StatePoisoned {
6996 module_id: module_id.map(ToOwned::to_owned),
6997 })?;
6998 update(&mut state);
6999 Ok(())
7000}
7001
7002const SLOW_SNAPSHOT_LOCK_THRESHOLD: Duration = Duration::from_millis(250);
7003
7004fn lock_snapshot_for_control<'a>(
7005 snapshot: &'a SharedSnapshot,
7006 module_id: &str,
7007 caller: &'static str,
7008) -> Result<std::sync::MutexGuard<'a, SupervisorSnapshot>, SuperviseError> {
7009 let started_at = Instant::now();
7010 let guard = lock_snapshot(snapshot)?;
7011 let waited = started_at.elapsed();
7012 if waited >= SLOW_SNAPSHOT_LOCK_THRESHOLD {
7013 warn!(
7014 module_id = %module_id,
7015 waited_ms = waited.as_millis() as u64,
7016 caller = %caller,
7017 "slow snapshot lock"
7018 );
7019 }
7020 Ok(guard)
7021}
7022
7023fn lock_snapshot(
7024 snapshot: &SharedSnapshot,
7025) -> Result<std::sync::MutexGuard<'_, SupervisorSnapshot>, SuperviseError> {
7026 snapshot
7027 .lock()
7028 .map_err(|_| SuperviseError::StatePoisoned { module_id: None })
7029}
7030
7031#[cfg(test)]
7032mod terminal_history_tests {
7033 use std::{
7034 path::PathBuf,
7035 sync::Arc,
7036 time::{Duration, Instant},
7037 };
7038
7039 use tokio::time::sleep;
7040
7041 use super::{
7042 apply_deliberate_severance_marker, daemon_will_restart, drain_child_to_state,
7043 drained_after_quiescence_wait, handle_reload_spawn_failure, health_restart_child,
7044 lock_snapshot, on_child_exit, record_deliberate_severance, record_wait_error_terminal,
7045 reset_restart_count, spawn_and_mark_running, update_snapshot, wait_error_exit_report,
7046 ExitKind, ExitReport, ModuleProtocol, ModuleSpec, ModuleState, NextAction, ProcessIdentity,
7047 RestartPolicy, SpawnEventKind, SuperviseError, SupervisedModule, Supervisor,
7048 SupervisorHandle, SupervisorHealthStatus, SupervisorSnapshot,
7049 };
7050 use super::Instant as ClockInstant;
7055 use crate::{
7056 registry::Registry,
7057 terminal_ring::{TerminalRing, TerminalRingConfig},
7058 };
7059 use std::sync::Mutex;
7060 use subc_control::TerminalDisposition;
7061
7062 fn fake_aft_stub_path() -> PathBuf {
7067 let mut path = std::env::current_exe().expect("current_exe available in tests");
7068 path.pop();
7069 path.pop();
7070 path.push(if cfg!(windows) {
7071 "fake-aft-stub.exe"
7072 } else {
7073 "fake-aft-stub"
7074 });
7075 assert!(
7076 path.exists(),
7077 "fake-aft-stub not built at {}: run `cargo test -p subc-core` (which builds \
7078 [[bin]] targets) rather than `cargo test -p subc-core --lib` (which does not)",
7079 path.display()
7080 );
7081 path
7082 }
7083
7084 #[test]
7085 fn reserved_never_spawned_refuses_every_hello() {
7086 let supervisor = SupervisorHandle::default();
7091 supervisor.apply_identity_configuration(&ModuleSpec {
7092 module_id: "never-spawned".to_string(),
7093 program: PathBuf::from("/usr/bin/false"),
7094 args: Vec::new(),
7095 env: Vec::new(),
7096 reserved: true,
7097 reserved_prefixes: Vec::new(),
7098 protocol: ModuleProtocol::Subc,
7099 overlap: Default::default(),
7100 });
7101 assert!(
7102 supervisor
7103 .reserved_hello_rejection("never-spawned", Some("any-forged-nonce"))
7104 .is_some(),
7105 "forged nonce must refuse on a reserved never-spawned id"
7106 );
7107 assert!(
7108 supervisor
7109 .reserved_hello_rejection("never-spawned", None)
7110 .is_some(),
7111 "absent nonce must refuse on a reserved never-spawned id"
7112 );
7113 supervisor.set_spawn_nonce("never-spawned", "minted".to_string());
7115 supervisor.apply_identity_configuration(&ModuleSpec {
7116 module_id: "never-spawned".to_string(),
7117 program: PathBuf::from("/usr/bin/false"),
7118 args: Vec::new(),
7119 env: Vec::new(),
7120 reserved: true,
7121 reserved_prefixes: Vec::new(),
7122 protocol: ModuleProtocol::Subc,
7123 overlap: Default::default(),
7124 });
7125 assert!(supervisor
7126 .reserved_hello_rejection("never-spawned", Some("minted"))
7127 .is_none());
7128 assert!(supervisor
7129 .reserved_hello_rejection("never-spawned", Some("forged"))
7130 .is_some());
7131 }
7132
7133 fn seed_crash_restarts(state: &mut SupervisorSnapshot, count: u32) {
7136 let now = ClockInstant::now();
7137 for _ in 0..count {
7138 state.crash_restarts.push_back(now);
7139 }
7140 }
7141
7142 fn age_oldest_crash_restart_out_of_window(state: &mut SupervisorSnapshot, window: Duration) {
7146 let aged = state
7147 .crash_restarts
7148 .front()
7149 .expect("a crash restart must be recorded before it can be aged")
7150 .checked_sub(window + Duration::from_secs(1))
7151 .expect("the test clock is far enough from its origin to age an instant");
7152 state.crash_restarts[0] = aged;
7153 }
7154
7155 fn snapshot_with_restarts(enabled: bool, count: u32) -> SupervisorSnapshot {
7156 let mut state = SupervisorSnapshot::new(ModuleState::Running, enabled);
7157 seed_crash_restarts(&mut state, count);
7158 state
7159 }
7160
7161 #[test]
7162 fn daemon_owned_recovery_predicate_uses_the_pre_increment_budget() {
7163 let policy = RestartPolicy::new(3, Duration::ZERO);
7164 let now = ClockInstant::now();
7165 assert!(daemon_will_restart(
7166 &mut snapshot_with_restarts(true, 2),
7167 &policy,
7168 now
7169 ));
7170 assert!(!daemon_will_restart(
7171 &mut snapshot_with_restarts(true, 3),
7172 &policy,
7173 now
7174 ));
7175 assert!(!daemon_will_restart(
7176 &mut snapshot_with_restarts(false, 0),
7177 &policy,
7178 now
7179 ));
7180 }
7181
7182 #[test]
7183 fn crash_restart_backoff_escalates_with_in_window_count() {
7184 let policy = RestartPolicy::new(4, Duration::from_millis(100))
7185 .with_max_backoff(Duration::from_secs(30));
7186 let now = ClockInstant::now();
7187 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7188 let schedules = (0..4)
7189 .map(|_| {
7190 state
7191 .next_crash_restart(&policy, now)
7192 .expect("the test policy allows four crash restarts")
7193 })
7194 .collect::<Vec<_>>();
7195
7196 assert_eq!(
7197 schedules
7198 .iter()
7199 .map(|schedule| schedule.restart_in_window)
7200 .collect::<Vec<_>>(),
7201 vec![0, 1, 2, 3]
7202 );
7203 assert_eq!(
7204 schedules
7205 .iter()
7206 .map(|schedule| schedule.delay)
7207 .collect::<Vec<_>>(),
7208 vec![
7209 Duration::from_millis(100),
7210 Duration::from_secs(1),
7211 Duration::from_secs(10),
7212 Duration::from_secs(30),
7213 ]
7214 );
7215 }
7216
7217 #[test]
7218 fn crash_restart_backoff_resets_after_ring_clear() {
7219 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7220 let now = ClockInstant::now();
7221 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7222 assert_eq!(
7223 state.next_crash_restart(&policy, now).unwrap().delay,
7224 Duration::from_millis(100)
7225 );
7226 assert_eq!(
7227 state.next_crash_restart(&policy, now).unwrap().delay,
7228 Duration::from_secs(1)
7229 );
7230
7231 state.clear_crash_restarts();
7232 let schedule = state
7233 .next_crash_restart(&policy, now)
7234 .expect("a cleared ring must allow another restart");
7235 assert_eq!(schedule.restart_in_window, 0);
7236 assert_eq!(schedule.delay, Duration::from_millis(100));
7237 }
7238
7239 #[test]
7240 fn crash_restart_backoff_ignores_aged_restarts() {
7241 let policy = RestartPolicy::new(3, Duration::from_millis(100));
7242 let now = ClockInstant::now();
7243 let mut state = SupervisorSnapshot::new(ModuleState::Running, true);
7244 state
7245 .next_crash_restart(&policy, now)
7246 .expect("the first restart is allowed");
7247 state
7248 .next_crash_restart(&policy, now)
7249 .expect("the second restart is allowed");
7250 state.crash_restarts[0] = now
7251 .checked_sub(policy.window + Duration::from_secs(1))
7252 .expect("the fake clock can age a restart past the window");
7253
7254 let schedule = state
7255 .next_crash_restart(&policy, now)
7256 .expect("an aged restart must release its slot");
7257 assert_eq!(schedule.restart_in_window, 1);
7258 assert_eq!(schedule.delay, Duration::from_secs(1));
7259 assert_eq!(state.crash_restarts.len(), 2);
7260 }
7261
7262 #[test]
7266 fn a_budget_spent_before_the_window_no_longer_refuses() {
7267 let policy = RestartPolicy::new(3, Duration::ZERO);
7268 let mut state = snapshot_with_restarts(true, 3);
7269 let now = ClockInstant::now();
7270 assert!(!daemon_will_restart(&mut state, &policy, now));
7271
7272 assert!(daemon_will_restart(
7273 &mut state,
7274 &policy,
7275 now + policy.window + Duration::from_secs(1)
7276 ));
7277 assert!(
7278 state.crash_restarts.is_empty(),
7279 "reading the budget must drop the instants that left the window"
7280 );
7281 }
7282
7283 fn module_with_recovery_snapshot(
7284 state: ModuleState,
7285 enabled: bool,
7286 restart_count: u32,
7287 ) -> SupervisedModule {
7288 let registry = Arc::new(Registry::default());
7289 let supervisor =
7290 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(3, Duration::ZERO));
7291 let module = supervisor
7292 .spawn(ModuleSpec {
7293 module_id: "recovery-snapshot".to_string(),
7294 program: fake_aft_stub_path(),
7295 args: Vec::new(),
7296 env: Vec::new(),
7297 reserved: false,
7298 reserved_prefixes: Vec::new(),
7299 protocol: ModuleProtocol::Subc,
7300 overlap: Default::default(),
7301 })
7302 .unwrap();
7303 update_snapshot(
7304 &module.inner.snapshot,
7305 Some("recovery-snapshot"),
7306 |snapshot| {
7307 snapshot.state = state;
7308 snapshot.enabled = enabled;
7309 seed_crash_restarts(snapshot, restart_count);
7310 },
7311 )
7312 .unwrap();
7313 module
7314 }
7315
7316 #[cfg(target_os = "linux")]
7317 #[tokio::test]
7318 async fn no_cgroup_placement_does_not_block_fake_aft_stub_spawn() {
7319 let supervisor = Supervisor::new(Arc::new(Registry::default()), RestartPolicy::default())
7320 .with_cgroup_placement(None);
7321 let result = supervisor.spawn(ModuleSpec {
7322 module_id: "no-cgroup-placement".to_string(),
7323 program: fake_aft_stub_path(),
7324 args: Vec::new(),
7325 env: Vec::new(),
7326 reserved: false,
7327 reserved_prefixes: Vec::new(),
7328 protocol: ModuleProtocol::Subc,
7329 overlap: Default::default(),
7330 });
7331
7332 assert!(
7333 result.is_ok(),
7334 "no delegation must not turn an otherwise valid spawn into a failure: {result:?}"
7335 );
7336 }
7337
7338 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7339 async fn undecided_snapshot_uses_shared_restart_predicate() {
7340 assert!(module_with_recovery_snapshot(ModuleState::Running, true, 2)
7341 .will_recover_after_connection_loss()
7342 .unwrap());
7343 assert!(
7344 !module_with_recovery_snapshot(ModuleState::Running, true, 3)
7345 .will_recover_after_connection_loss()
7346 .unwrap()
7347 );
7348 }
7349
7350 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7351 async fn restarting_snapshot_at_exhausted_budget_is_non_terminal() {
7352 assert!(
7353 module_with_recovery_snapshot(ModuleState::Restarting, true, 3)
7354 .will_recover_after_connection_loss()
7355 .unwrap()
7356 );
7357 }
7358
7359 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7360 async fn terminal_phase_snapshots_are_terminal_before_budget_exhaustion() {
7361 assert!(!module_with_recovery_snapshot(ModuleState::Failed, true, 0)
7362 .will_recover_after_connection_loss()
7363 .unwrap());
7364 assert!(
7365 !module_with_recovery_snapshot(ModuleState::Disabled, true, 0)
7366 .will_recover_after_connection_loss()
7367 .unwrap()
7368 );
7369 }
7370
7371 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7372 async fn warming_snapshot_is_limited_to_startup_phases() {
7373 for state in [
7374 ModuleState::Starting,
7375 ModuleState::Running,
7376 ModuleState::Restarting,
7377 ] {
7378 assert!(
7379 module_with_recovery_snapshot(state, true, 0)
7380 .is_warming()
7381 .unwrap(),
7382 "{state:?} should be warming"
7383 );
7384 }
7385 for state in [
7386 ModuleState::Unresponsive,
7387 ModuleState::Draining,
7388 ModuleState::Stopped,
7389 ModuleState::Failed,
7390 ModuleState::Disabled,
7391 ] {
7392 assert!(
7393 !module_with_recovery_snapshot(state, true, 0)
7394 .is_warming()
7395 .unwrap(),
7396 "{state:?} should not be warming"
7397 );
7398 }
7399 }
7400
7401 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7402 async fn terminal_history_survives_respawn_and_keeps_both_crashes_in_order() {
7403 let registry = Arc::new(Registry::default());
7404 let supervisor =
7405 Supervisor::new(Arc::clone(®istry), RestartPolicy::new(1, Duration::ZERO));
7406 let module = supervisor
7407 .spawn(ModuleSpec {
7408 module_id: "terminal-history".to_string(),
7409 program: fake_aft_stub_path(),
7410 args: Vec::new(),
7411 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7412 reserved: false,
7413 reserved_prefixes: Vec::new(),
7414 protocol: ModuleProtocol::Subc,
7415 overlap: Default::default(),
7416 })
7417 .unwrap();
7418
7419 let deadline = Instant::now() + Duration::from_secs(5);
7420 loop {
7421 let history = module.terminal_history();
7422 if history.entries.len() == 2 {
7423 assert_eq!(module.status().unwrap().state, ModuleState::Failed);
7424 assert_eq!(history.dropped, 0);
7425 assert_eq!(
7426 history
7427 .entries
7428 .iter()
7429 .map(|entry| entry.exit_code)
7430 .collect::<Vec<_>>(),
7431 vec![Some(23), Some(23)]
7432 );
7433 assert!(history.entries[0].at_ms <= history.entries[1].at_ms);
7434 return;
7435 }
7436 assert!(
7437 Instant::now() < deadline,
7438 "module did not retain two terminal exits: {history:?}"
7439 );
7440 sleep(Duration::from_millis(10)).await;
7441 }
7442 }
7443
7444 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7448 async fn disable_during_crash_backoff_cancels_pending_respawn() {
7449 let backoff = Duration::from_secs(2);
7450 let supervisor = Supervisor::new(
7451 Arc::new(Registry::default()),
7452 RestartPolicy::new(10, backoff),
7453 );
7454 let module = supervisor
7455 .spawn(ModuleSpec {
7456 module_id: "disable-during-backoff".to_string(),
7457 program: fake_aft_stub_path(),
7458 args: Vec::new(),
7459 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7460 reserved: false,
7461 reserved_prefixes: Vec::new(),
7462 protocol: ModuleProtocol::Subc,
7463 overlap: Default::default(),
7464 })
7465 .unwrap();
7466
7467 let deadline = Instant::now() + Duration::from_secs(5);
7469 loop {
7470 if module.status().unwrap().state == ModuleState::Restarting {
7471 break;
7472 }
7473 assert!(
7474 Instant::now() < deadline,
7475 "module never entered the crash backoff"
7476 );
7477 sleep(Duration::from_millis(10)).await;
7478 }
7479
7480 let started = Instant::now();
7481 module.set_enabled(false).await.unwrap();
7482 let waited = started.elapsed();
7483
7484 assert!(
7485 waited < backoff / 2,
7486 "disable waited {waited:?} behind the {backoff:?} crash backoff; the operator command must preempt the pending respawn"
7487 );
7488 assert_eq!(module.status().unwrap().state, ModuleState::Disabled);
7489
7490 sleep(backoff + Duration::from_millis(500)).await;
7492 let status = module.status().unwrap();
7493 assert_eq!(status.state, ModuleState::Disabled);
7494 assert_eq!(
7495 status.spawn_generation, 1,
7496 "module respawned after the operator disabled it"
7497 );
7498 }
7499
7500 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7504 async fn every_restart_increment_path_advances_lifetime_count() {
7505 let supervisor = Supervisor::new(
7506 Arc::new(Registry::default()),
7507 RestartPolicy::new(1, Duration::ZERO),
7508 );
7509 let runtime = supervisor.runtime_config();
7510 let spec = ModuleSpec {
7511 module_id: "lifetime-increment-path".to_string(),
7512 program: PathBuf::from("/unused/lifetime-increment-path"),
7513 args: Vec::new(),
7514 env: Vec::new(),
7515 reserved: false,
7516 reserved_prefixes: Vec::new(),
7517 protocol: ModuleProtocol::Subc,
7518 overlap: Default::default(),
7519 };
7520
7521 let crash_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7522 assert!(matches!(
7523 on_child_exit(
7524 &spec,
7525 runtime.restart_policy,
7526 &supervisor.registry,
7527 &crash_snapshot,
7528 &runtime.terminal_ring,
7529 &runtime.spawn_events,
7530 ExitReport {
7531 kind: ExitKind::Crash,
7532 code: Some(1),
7533 signal: None,
7534 at_ms: 1,
7535 },
7536 )
7537 .await,
7538 NextAction::Restart { schedule: _ }
7539 ));
7540 let (crash_restarts, crash_lifetime) = {
7541 let state = lock_snapshot(&crash_snapshot).unwrap();
7542 (state.crash_restarts.len(), state.lifetime_restarts)
7543 };
7544 assert_eq!(crash_restarts, 1);
7545 assert_eq!(crash_lifetime, 1);
7546
7547 let health_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7548 let mut health_child = None;
7549 assert!(matches!(
7550 health_restart_child(
7551 &spec,
7552 &runtime,
7553 &supervisor.registry,
7554 &supervisor.process_liveness,
7555 &health_snapshot,
7556 &mut health_child,
7557 SupervisorHealthStatus::Failing,
7558 None,
7559 2,
7560 )
7561 .await,
7562 Err(SuperviseError::Spawn { .. })
7563 ));
7564 let (health_restarts, health_lifetime) = {
7565 let state = lock_snapshot(&health_snapshot).unwrap();
7566 (state.crash_restarts.len(), state.lifetime_restarts)
7567 };
7568 assert_eq!(health_restarts, 1);
7569 assert_eq!(health_lifetime, 1);
7570
7571 let reload_snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7572 let mut reload_child = None;
7573 assert!(matches!(
7574 handle_reload_spawn_failure(
7575 &spec,
7576 &runtime,
7577 &supervisor.process_liveness,
7578 &reload_snapshot,
7579 &mut reload_child,
7580 "forced reload spawn failure".to_string(),
7581 )
7582 .await,
7583 Err(SuperviseError::ReloadFailed { .. })
7584 ));
7585 let (reload_restarts, reload_lifetime) = {
7586 let state = lock_snapshot(&reload_snapshot).unwrap();
7587 (state.crash_restarts.len(), state.lifetime_restarts)
7588 };
7589 assert_eq!(reload_restarts, 1);
7590 assert_eq!(reload_lifetime, 1);
7591 }
7592
7593 #[tokio::test]
7594 async fn deliberately_severed_live_child_records_lifetime_without_spending_restart_budget() {
7595 let supervisor = Supervisor::new(
7596 Arc::new(Registry::default()),
7597 RestartPolicy::new(3, Duration::ZERO),
7598 );
7599 let runtime = supervisor.runtime_config();
7600 let spec = ModuleSpec {
7601 module_id: "deliberately-severed".to_string(),
7602 program: PathBuf::from("/unused/deliberately-severed"),
7603 args: Vec::new(),
7604 env: Vec::new(),
7605 reserved: false,
7606 reserved_prefixes: Vec::new(),
7607 protocol: ModuleProtocol::Subc,
7608 overlap: Default::default(),
7609 };
7610 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7611 let process = ProcessIdentity {
7612 pid: 41,
7613 start_time: 101,
7614 };
7615 record_deliberate_severance(&snapshot, process).unwrap();
7616 let exit_report = apply_deliberate_severance_marker(
7617 &snapshot,
7618 Some(process),
7619 ExitReport {
7620 kind: ExitKind::Crash,
7621 code: Some(1),
7622 signal: None,
7623 at_ms: 1,
7624 },
7625 );
7626 assert_eq!(exit_report.kind, ExitKind::DeliberateSeverance);
7627
7628 assert!(matches!(
7629 on_child_exit(
7630 &spec,
7631 runtime.restart_policy,
7632 &supervisor.registry,
7633 &snapshot,
7634 &runtime.terminal_ring,
7635 &runtime.spawn_events,
7636 exit_report,
7637 )
7638 .await,
7639 NextAction::Restart { schedule: _ }
7640 ));
7641 let state = lock_snapshot(&snapshot).unwrap();
7642 assert_eq!(state.lifetime_restarts, 1);
7643 assert_eq!(state.crash_restarts.len(), 0);
7644 }
7645
7646 #[tokio::test]
7647 async fn genuine_crash_spends_restart_budget_and_records_lifetime() {
7648 let supervisor = Supervisor::new(
7649 Arc::new(Registry::default()),
7650 RestartPolicy::new(3, Duration::ZERO),
7651 );
7652 let runtime = supervisor.runtime_config();
7653 let spec = ModuleSpec {
7654 module_id: "genuine-crash".to_string(),
7655 program: PathBuf::from("/unused/genuine-crash"),
7656 args: Vec::new(),
7657 env: Vec::new(),
7658 reserved: false,
7659 reserved_prefixes: Vec::new(),
7660 protocol: ModuleProtocol::Subc,
7661 overlap: Default::default(),
7662 };
7663 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7664
7665 assert!(matches!(
7666 on_child_exit(
7667 &spec,
7668 runtime.restart_policy,
7669 &supervisor.registry,
7670 &snapshot,
7671 &runtime.terminal_ring,
7672 &runtime.spawn_events,
7673 ExitReport {
7674 kind: ExitKind::Crash,
7675 code: Some(1),
7676 signal: None,
7677 at_ms: 1,
7678 },
7679 )
7680 .await,
7681 NextAction::Restart { schedule: _ }
7682 ));
7683 let state = lock_snapshot(&snapshot).unwrap();
7684 assert_eq!(state.lifetime_restarts, 1);
7685 assert_eq!(state.crash_restarts.len(), 1);
7686 }
7687
7688 fn crash_exit_report(at_ms: u64) -> ExitReport {
7689 ExitReport {
7690 kind: ExitKind::Crash,
7691 code: Some(1),
7692 signal: None,
7693 at_ms,
7694 }
7695 }
7696
7697 fn windowed_crash_spec(module_id: &str) -> ModuleSpec {
7698 ModuleSpec {
7699 module_id: module_id.to_string(),
7700 program: PathBuf::from("/unused").join(module_id),
7701 args: Vec::new(),
7702 env: Vec::new(),
7703 reserved: false,
7704 reserved_prefixes: Vec::new(),
7705 protocol: ModuleProtocol::Subc,
7706 overlap: Default::default(),
7707 }
7708 }
7709
7710 #[tokio::test]
7716 async fn three_crashes_inside_the_window_stop_the_module_and_name_the_window() {
7717 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::ERROR);
7718 let supervisor = Supervisor::new(
7719 Arc::new(Registry::default()),
7720 RestartPolicy::new(2, Duration::ZERO),
7721 );
7722 let runtime = supervisor.runtime_config();
7723 let spec = windowed_crash_spec("crash-loop-in-window");
7724 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7725
7726 for attempt in 1..=2 {
7727 assert!(
7728 matches!(
7729 on_child_exit(
7730 &spec,
7731 runtime.restart_policy,
7732 &supervisor.registry,
7733 &snapshot,
7734 &runtime.terminal_ring,
7735 &runtime.spawn_events,
7736 crash_exit_report(attempt),
7737 )
7738 .await,
7739 NextAction::Restart { schedule: _ }
7740 ),
7741 "crash {attempt} is inside the budget and must respawn"
7742 );
7743 }
7744
7745 assert!(matches!(
7746 on_child_exit(
7747 &spec,
7748 runtime.restart_policy,
7749 &supervisor.registry,
7750 &snapshot,
7751 &runtime.terminal_ring,
7752 &runtime.spawn_events,
7753 crash_exit_report(3),
7754 )
7755 .await,
7756 NextAction::Stop { .. }
7757 ));
7758
7759 {
7760 let state = lock_snapshot(&snapshot).unwrap();
7761 assert_eq!(state.state, ModuleState::Failed);
7762 assert_eq!(state.crash_restarts.len(), 2);
7763 assert_eq!(state.lifetime_restarts, 2);
7764 }
7765
7766 let history = runtime
7767 .terminal_ring
7768 .lock()
7769 .expect("terminal ring is not poisoned")
7770 .snapshot();
7771 let last = history
7772 .entries
7773 .last()
7774 .expect("the refused crash is retained");
7775 assert_eq!(last.disposition, TerminalDisposition::Failed);
7776 assert_eq!(
7777 last.disposition_detail.as_deref(),
7778 Some("crash budget exhausted: max_restarts=2 within window_secs=600")
7779 );
7780
7781 let captured = crate::router::test_log::captured_logs(&logs);
7782 assert!(
7783 captured.contains("crash budget exhausted: max_restarts=2 within window_secs=600"),
7784 "the stop must be logged with its window: {captured}"
7785 );
7786 }
7787
7788 #[tokio::test]
7796 async fn a_crash_older_than_the_window_frees_its_slot_for_a_later_crash() {
7797 let supervisor = Supervisor::new(
7798 Arc::new(Registry::default()),
7799 RestartPolicy::new(2, Duration::ZERO),
7800 );
7801 let runtime = supervisor.runtime_config();
7802 let spec = windowed_crash_spec("crash-across-windows");
7803 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7804
7805 for attempt in 1..=2 {
7806 assert!(matches!(
7807 on_child_exit(
7808 &spec,
7809 runtime.restart_policy,
7810 &supervisor.registry,
7811 &snapshot,
7812 &runtime.terminal_ring,
7813 &runtime.spawn_events,
7814 crash_exit_report(attempt),
7815 )
7816 .await,
7817 NextAction::Restart { schedule: _ }
7818 ));
7819 }
7820
7821 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
7824 age_oldest_crash_restart_out_of_window(state, runtime.restart_policy.window);
7825 })
7826 .unwrap();
7827
7828 assert!(
7829 matches!(
7830 on_child_exit(
7831 &spec,
7832 runtime.restart_policy,
7833 &supervisor.registry,
7834 &snapshot,
7835 &runtime.terminal_ring,
7836 &runtime.spawn_events,
7837 crash_exit_report(3),
7838 )
7839 .await,
7840 NextAction::Restart { schedule: _ }
7841 ),
7842 "a crash older than the window must not hold a budget slot"
7843 );
7844
7845 let state = lock_snapshot(&snapshot).unwrap();
7846 assert_eq!(state.state, ModuleState::Restarting);
7847 assert_eq!(
7848 state.crash_restarts.len(),
7849 2,
7850 "the aged instant is dropped and the new one takes its place"
7851 );
7852 assert_eq!(
7853 state.lifetime_restarts, 3,
7854 "the ledger counts every restart, including the ones the window forgot"
7855 );
7856 }
7857
7858 #[tokio::test]
7863 async fn an_operator_restart_clears_the_ring_and_leaves_the_ledger_alone() {
7864 let supervisor = Supervisor::new(
7865 Arc::new(Registry::default()),
7866 RestartPolicy::new(2, Duration::ZERO),
7867 );
7868 let runtime = supervisor.runtime_config();
7869 let spec = windowed_crash_spec("operator-cleared-budget");
7870 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7871
7872 for attempt in 1..=2 {
7873 assert!(matches!(
7874 on_child_exit(
7875 &spec,
7876 runtime.restart_policy,
7877 &supervisor.registry,
7878 &snapshot,
7879 &runtime.terminal_ring,
7880 &runtime.spawn_events,
7881 crash_exit_report(attempt),
7882 )
7883 .await,
7884 NextAction::Restart { schedule: _ }
7885 ));
7886 }
7887
7888 reset_restart_count(&snapshot, &spec.module_id).unwrap();
7889 {
7890 let state = lock_snapshot(&snapshot).unwrap();
7891 assert!(
7892 state.crash_restarts.is_empty(),
7893 "an operator restart returns the full budget"
7894 );
7895 assert_eq!(
7896 state.lifetime_restarts, 2,
7897 "clearing the budget must not unmake the crashes"
7898 );
7899 }
7900
7901 assert!(
7902 matches!(
7903 on_child_exit(
7904 &spec,
7905 runtime.restart_policy,
7906 &supervisor.registry,
7907 &snapshot,
7908 &runtime.terminal_ring,
7909 &runtime.spawn_events,
7910 crash_exit_report(3),
7911 )
7912 .await,
7913 NextAction::Restart { schedule: _ }
7914 ),
7915 "the cleared budget must be spendable again"
7916 );
7917 let state = lock_snapshot(&snapshot).unwrap();
7918 assert_eq!(state.crash_restarts.len(), 1);
7919 assert_eq!(state.lifetime_restarts, 3);
7920 }
7921
7922 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7923 async fn severance_marker_for_a_dead_child_does_not_label_its_successor() {
7924 let severed = ProcessIdentity {
7925 pid: 41,
7926 start_time: 101,
7927 };
7928 let successor = ProcessIdentity {
7929 pid: 41,
7930 start_time: 202,
7931 };
7932 let module = module_with_recovery_snapshot(ModuleState::Running, true, 0);
7933 update_snapshot(&module.inner.snapshot, Some("recovery-snapshot"), |state| {
7934 state.pid = Some(successor.pid);
7935 state.process_start_time = Some(successor.start_time);
7936 })
7937 .unwrap();
7938 assert!(!module.record_deliberate_severance(severed).unwrap());
7939
7940 let exit_report = apply_deliberate_severance_marker(
7941 &module.inner.snapshot,
7942 Some(successor),
7943 ExitReport {
7944 kind: ExitKind::Crash,
7945 code: Some(1),
7946 signal: None,
7947 at_ms: 1,
7948 },
7949 );
7950
7951 assert_eq!(exit_report.kind, ExitKind::Crash);
7952 }
7953
7954 #[tokio::test]
7955 async fn drain_reap_marks_deliberate_severance_and_records_lifetime_without_budget() {
7956 let registry = Registry::default();
7957 let supervisor = Supervisor::new(
7958 Arc::new(Registry::default()),
7959 RestartPolicy::new(3, Duration::ZERO),
7960 );
7961 let runtime = supervisor.runtime_config();
7962 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
7963 let spec = ModuleSpec {
7964 module_id: "drain-deliberate-severance".to_string(),
7965 program: fake_aft_stub_path(),
7966 args: Vec::new(),
7967 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
7968 reserved: false,
7969 reserved_prefixes: Vec::new(),
7970 protocol: ModuleProtocol::Subc,
7971 overlap: Default::default(),
7972 };
7973 let mut child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
7974 let process = ProcessIdentity {
7975 pid: 41,
7976 start_time: 101,
7977 };
7978 child.process_identity = Some(process);
7979 update_snapshot(&snapshot, Some(&spec.module_id), |state| {
7980 state.pid = Some(process.pid);
7981 state.process_start_time = Some(process.start_time);
7982 })
7983 .unwrap();
7984 record_deliberate_severance(&snapshot, process).unwrap();
7985
7986 drain_child_to_state(
7987 &spec.module_id,
7988 spec.protocol,
7989 ®istry,
7990 &snapshot,
7991 &runtime.terminal_ring,
7992 &runtime.spawn_events,
7993 child,
7994 Duration::from_secs(1),
7995 ModuleState::Stopped,
7996 Some(false),
7997 )
7998 .await
7999 .unwrap();
8000
8001 let state = lock_snapshot(&snapshot).unwrap();
8002 assert_eq!(
8003 state.last_exit.as_ref().map(|exit| exit.kind),
8004 Some(ExitKind::DeliberateSeverance)
8005 );
8006 assert_eq!(state.lifetime_restarts, 1);
8007 assert_eq!(state.crash_restarts.len(), 0);
8008 drop(state);
8009 let history = runtime.terminal_ring.lock().unwrap().snapshot();
8010 assert_eq!(
8011 history.entries[0].exit_kind,
8012 subc_control::TerminalExitKind::DeliberateSeverance
8013 );
8014 }
8015
8016 #[tokio::test]
8017 async fn ordinary_drain_reap_does_not_record_a_lifetime_restart() {
8018 let registry = Registry::default();
8019 let supervisor = Supervisor::new(
8020 Arc::new(Registry::default()),
8021 RestartPolicy::new(3, Duration::ZERO),
8022 );
8023 let runtime = supervisor.runtime_config();
8024 let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting()));
8025 let spec = ModuleSpec {
8026 module_id: "ordinary-drain".to_string(),
8027 program: fake_aft_stub_path(),
8028 args: Vec::new(),
8029 env: vec![("FAKE_AFT_EXIT_CODE".to_string(), "23".to_string())],
8030 reserved: false,
8031 reserved_prefixes: Vec::new(),
8032 protocol: ModuleProtocol::Subc,
8033 overlap: Default::default(),
8034 };
8035 let child = spawn_and_mark_running(&spec, &runtime, &snapshot).unwrap();
8036
8037 drain_child_to_state(
8038 &spec.module_id,
8039 spec.protocol,
8040 ®istry,
8041 &snapshot,
8042 &runtime.terminal_ring,
8043 &runtime.spawn_events,
8044 child,
8045 Duration::from_secs(1),
8046 ModuleState::Stopped,
8047 Some(false),
8048 )
8049 .await
8050 .unwrap();
8051
8052 let state = lock_snapshot(&snapshot).unwrap();
8053 assert_eq!(
8054 state.last_exit.as_ref().map(|exit| exit.kind),
8055 Some(ExitKind::Crash)
8056 );
8057 assert_eq!(state.lifetime_restarts, 0);
8058 assert_eq!(state.crash_restarts.len(), 0);
8059 }
8060
8061 #[test]
8062 fn fatal_connection_teardown_cannot_arm_a_marker_for_a_surviving_process() {
8063 assert!(!include_str!("server.rs")
8069 .contains("router.record_deliberate_connection_severance(ctx.connection_id)"));
8070 }
8071
8072 #[test]
8079 fn drained_after_quiescence_wait_passes_ok_through_and_forces_false_on_err() {
8080 assert!(drained_after_quiescence_wait(&Ok(true)));
8081 assert!(!drained_after_quiescence_wait(&Ok(false)));
8082 assert!(!drained_after_quiescence_wait(&Err(
8083 SuperviseError::StatePoisoned { module_id: None }
8084 )));
8085 }
8086
8087 #[test]
8096 fn wait_error_exit_report_records_a_failed_terminal_with_no_code_or_signal() {
8097 let ring = Arc::new(Mutex::new(TerminalRing::new(
8098 TerminalRingConfig::default(),
8099 0,
8100 )));
8101 record_wait_error_terminal("wait-error", &ring, &super::SpawnEventFeed::default());
8102
8103 let snapshot = ring.lock().unwrap().snapshot();
8104 assert_eq!(snapshot.entries.len(), 1);
8105 let entry = &snapshot.entries[0];
8106 assert_eq!(entry.exit_code, None);
8107 assert_eq!(entry.exit_signal, None);
8108 assert_eq!(entry.disposition, TerminalDisposition::Failed);
8109 }
8110
8111 #[test]
8112 fn wait_error_exit_path_preserves_spawn_event_density() {
8113 let feed = super::SpawnEventFeed::default();
8114 feed.configure_incarnation("wait-error-density".to_string());
8115 feed.emit_spawned("wait-error", 41, 1);
8116 let ring = Arc::new(Mutex::new(TerminalRing::new(
8117 TerminalRingConfig::default(),
8118 0,
8119 )));
8120
8121 record_wait_error_terminal("wait-error", &ring, &feed);
8122 feed.emit_spawned("after-wait-error", 42, 2);
8123
8124 let state = feed.0.lock().unwrap();
8125 let sequences = state
8126 .events
8127 .iter()
8128 .map(|event| event.cursor.seq)
8129 .collect::<Vec<_>>();
8130 assert_eq!(sequences, vec![1, 2, 3]);
8131 assert_eq!(state.events[1].kind, SpawnEventKind::Exited);
8132 assert_eq!(state.events[1].exit_code, None);
8133 assert_eq!(state.events[1].exit_signal, None);
8134 }
8135
8136 #[test]
8140 fn wait_error_exit_report_is_classified_as_a_crash() {
8141 assert_eq!(wait_error_exit_report().kind, ExitKind::Crash);
8142 }
8143}
8144
8145#[cfg(test)]
8146mod health_evidence_tests {
8147 use super::{HealthProbeError, HealthProbeEvidence};
8148 use std::collections::HashSet;
8149
8150 #[test]
8158 fn only_a_dead_lane_is_proof_of_death() {
8159 assert!(HealthProbeError::lane_dead("gone").is_proof_of_death());
8160 assert!(!HealthProbeError::no_answer("timed out").is_proof_of_death());
8164 assert!(!HealthProbeError::bad_answer("garbage").is_proof_of_death());
8165 assert!(!HealthProbeError::misconfigured("no table").is_proof_of_death());
8166 }
8167
8168 #[test]
8174 fn every_evidence_class_has_a_distinct_label() {
8175 let labels = [
8176 HealthProbeError::lane_dead("").label(),
8177 HealthProbeError::no_answer("").label(),
8178 HealthProbeError::bad_answer("").label(),
8179 HealthProbeError::misconfigured("").label(),
8180 ];
8181 let unique: HashSet<_> = labels.iter().collect();
8182 assert_eq!(unique.len(), labels.len(), "labels collided: {labels:?}");
8183 }
8184
8185 #[test]
8191 fn classification_preserves_the_original_message() {
8192 let err = HealthProbeError::no_answer("module did not answer within 5s");
8193 assert_eq!(err.to_string(), "module did not answer within 5s");
8194 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8195 }
8196}
8197
8198#[cfg(test)]
8199mod health_tombstone_tests {
8200 use std::{path::PathBuf, sync::Arc, time::Duration};
8201
8202 use subc_protocol::{
8203 manifest::Concurrency,
8204 session::{HealthStatus, ModuleControlResponse},
8205 };
8206 use tokio::sync::mpsc;
8207
8208 use super::{
8209 probe_module_health, HealthAction, HealthConfig, HealthProbeEvidence, ModuleProtocol,
8210 ModuleSpec, RestartPolicy, Supervisor, SupervisorRuntimeConfig,
8211 };
8212 use crate::{
8213 control::ControlHandler,
8214 forwarding::{ForwardingTable, ModuleControlRpcCompletion, ModuleControlRpcOutcome},
8215 registry::{ConnectionId, Registry},
8216 router::FrameSink,
8217 };
8218
8219 struct ProbeHarness {
8220 spec: ModuleSpec,
8221 runtime: SupervisorRuntimeConfig,
8222 forwarding: Arc<ForwardingTable>,
8223 module_connection: ConnectionId,
8224 module_rx: mpsc::Receiver<crate::router::OutboundFrame>,
8225 handler: ControlHandler,
8226 module: super::SupervisedModule,
8227 }
8228
8229 fn probe_harness() -> ProbeHarness {
8230 let registry = Arc::new(Registry::default());
8231 let forwarding = Arc::new(ForwardingTable::default());
8232 let supervisor_handle = super::SupervisorHandle::new();
8233 let health = HealthConfig {
8234 cadence: Duration::from_secs(30),
8235 deadline: Duration::from_secs(5),
8236 failure_threshold: 3,
8237 on_degraded: HealthAction::Report,
8238 on_failing: HealthAction::Report,
8239 critical: false,
8240 };
8241 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
8242 .with_forwarding(Arc::clone(&forwarding))
8243 .with_handle(supervisor_handle.clone())
8244 .with_health_config(health);
8245 let spec = ModuleSpec {
8246 module_id: "late-health-module".to_string(),
8247 program: PathBuf::from("disabled-module"),
8248 args: Vec::new(),
8249 env: Vec::new(),
8250 reserved: false,
8251 reserved_prefixes: Vec::new(),
8252 protocol: ModuleProtocol::Subc,
8253 overlap: Default::default(),
8254 };
8255 let module = supervisor
8256 .supervise_configured(spec.clone(), false)
8257 .unwrap();
8258 let runtime = supervisor.runtime_config();
8259 let handler = ControlHandler::with_forwarding(registry, Arc::clone(&forwarding))
8260 .with_supervisor(supervisor_handle);
8261 let module_connection = ConnectionId::new(700);
8262 let (module_tx, module_rx) = mpsc::channel(8);
8263 forwarding
8264 .register_module_connection(
8265 module_connection,
8266 spec.module_id.clone(),
8267 subc_protocol::PROTOCOL_VERSION,
8268 Concurrency::ModuleManaged,
8269 FrameSink::new(module_tx),
8270 )
8271 .unwrap();
8272
8273 ProbeHarness {
8274 spec,
8275 runtime,
8276 forwarding,
8277 module_connection,
8278 module_rx,
8279 handler,
8280 module,
8281 }
8282 }
8283
8284 async fn finish_after(
8285 harness: &mut ProbeHarness,
8286 stall: Duration,
8287 ) -> ModuleControlRpcCompletion {
8288 assert!(stall > harness.runtime.health.deadline);
8289 let deadline = harness.runtime.health.deadline;
8290 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8291 let answer = async {
8292 let frame = harness.module_rx.recv().await.expect("health.check frame");
8293 tokio::time::advance(deadline).await;
8294 tokio::task::yield_now().await;
8295 tokio::time::advance(stall - deadline).await;
8296 harness
8297 .forwarding
8298 .complete_module_control_rpc(
8299 harness.module_connection,
8300 frame.header.corr,
8301 Some("health.check"),
8302 ModuleControlRpcOutcome::Response(ModuleControlResponse::HealthCheck {
8303 status: HealthStatus::Ok,
8304 detail: None,
8305 metrics: None,
8306 }),
8307 )
8308 .unwrap()
8309 };
8310 let (probe_result, completion) = tokio::join!(probe, answer);
8311 let err = probe_result.expect_err("probe must miss its deadline");
8312 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8313 completion
8314 }
8315
8316 async fn time_out_without_answer(harness: &mut ProbeHarness) {
8317 let deadline = harness.runtime.health.deadline;
8318 let probe = probe_module_health(&harness.spec.module_id, &harness.runtime, None);
8319 let exhaust_deadline = async {
8320 let _frame = harness.module_rx.recv().await.expect("health.check frame");
8321 tokio::time::advance(deadline).await;
8322 tokio::task::yield_now().await;
8323 };
8324 let (probe_result, ()) = tokio::join!(probe, exhaust_deadline);
8325 let err = probe_result.expect_err("probe must miss its deadline");
8326 assert!(matches!(err.evidence, HealthProbeEvidence::NoAnswer));
8327 }
8328
8329 #[tokio::test(start_paused = true)]
8330 async fn late_health_answers_record_start_anchored_latency_for_two_stalls() {
8331 let mut harness = probe_harness();
8332
8333 let first = finish_after(&mut harness, Duration::from_secs(8)).await;
8334 let first_latency = match &first {
8335 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8336 other => panic!("late answer was not retained: {other:?}"),
8337 };
8338 assert!(harness.handler.observe_module_control_completion(first));
8339
8340 let second = finish_after(&mut harness, Duration::from_secs(11)).await;
8341 let second_latency = match &second {
8342 ModuleControlRpcCompletion::LateHealthAnswer { latency, .. } => *latency,
8343 other => panic!("late answer was not retained: {other:?}"),
8344 };
8345 assert!(harness.handler.observe_module_control_completion(second));
8346
8347 assert_eq!(first_latency, Duration::from_secs(8));
8348 assert_eq!(
8349 second_latency - first_latency,
8350 Duration::from_secs(3),
8351 "latency must grow linearly with the additional stall"
8352 );
8353 let health = harness.module.status().unwrap().health;
8354 assert_eq!(health.late_answer_count, 2);
8355 assert_eq!(health.last_late_answer_latency_ms, Some(11_000));
8356 }
8357
8358 #[tokio::test(start_paused = true)]
8366 async fn late_answer_clears_the_consecutive_failure_streak() {
8367 let mut harness = probe_harness();
8368
8369 time_out_without_answer(&mut harness).await;
8371 harness
8372 .module
8373 .record_health_probe_failure_for_test("[no-answer] test miss")
8374 .unwrap();
8375 assert_eq!(
8376 harness.module.status().unwrap().health.consecutive_failures,
8377 1,
8378 "precondition: the miss must be on the streak before the late answer"
8379 );
8380
8381 let late = finish_after(&mut harness, Duration::from_secs(9)).await;
8383 assert!(matches!(
8384 late,
8385 ModuleControlRpcCompletion::LateHealthAnswer { .. }
8386 ));
8387 assert!(harness.handler.observe_module_control_completion(late));
8388
8389 let health = harness.module.status().unwrap().health;
8390 assert_eq!(
8391 health.consecutive_failures, 0,
8392 "a late answer is an answer: the streak must reset"
8393 );
8394 assert_eq!(health.late_answer_count, 1);
8395 }
8396
8397 #[tokio::test(start_paused = true)]
8398 async fn repeated_serial_probe_cycles_keep_one_tombstone_per_endpoint() {
8399 let mut harness = probe_harness();
8400
8401 for _ in 0..20 {
8402 time_out_without_answer(&mut harness).await;
8403 assert_eq!(
8404 harness.forwarding.health_probe_tombstone_count().unwrap(),
8405 1
8406 );
8407 }
8408 }
8409}
8410
8411#[cfg(test)]
8412mod child_env_tests {
8413 use super::{
8414 apply_child_env, apply_spawn_role, apply_wire_spawn_args, ModuleProtocol, ModuleSpec,
8415 SpawnRole, SupervisorHandle, SPAWN_ROLE_SWAP_CANDIDATE, SUBC_ARG, SUBC_LAUNCH_NONCE_ENV,
8416 SUBC_MODULE_ID_ENV, SUBC_SPAWN_ROLE_ENV,
8417 };
8418 use std::{ffi::OsStr, path::PathBuf};
8419 use tokio::process::Command;
8420
8421 fn spec(env: Vec<(String, String)>) -> ModuleSpec {
8422 ModuleSpec {
8423 module_id: "env-plan".to_string(),
8424 program: PathBuf::from("/nonexistent"),
8425 args: Vec::new(),
8426 env,
8427 reserved: false,
8428 reserved_prefixes: Vec::new(),
8429 protocol: ModuleProtocol::Subc,
8430 overlap: Default::default(),
8431 }
8432 }
8433
8434 #[test]
8448 fn ambient_ck_log_is_removed_and_a_configured_one_survives() {
8449 let mut command = Command::new("/nonexistent");
8450 apply_child_env(&mut command, &spec(Vec::new()));
8451 let removed = command
8452 .as_std()
8453 .get_envs()
8454 .any(|(key, value)| key == OsStr::new("CK_LOG") && value.is_none());
8455 assert!(
8456 removed,
8457 "ambient CK_LOG must be explicitly removed for an unconfigured module"
8458 );
8459
8460 let mut configured = Command::new("/nonexistent");
8461 apply_child_env(
8462 &mut configured,
8463 &spec(vec![("CK_LOG".to_string(), "debug".to_string())]),
8464 );
8465 let effective = configured
8466 .as_std()
8467 .get_envs()
8468 .filter(|(key, _)| *key == OsStr::new("CK_LOG"))
8469 .last()
8470 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()));
8471 assert_eq!(
8472 effective,
8473 Some(Some("debug".to_string())),
8474 "a module's configured CK_LOG must survive the ambient removal"
8475 );
8476 }
8477
8478 #[test]
8487 fn protocol_none_spawn_carries_no_subc_argument_and_no_nonce() {
8488 let connection_file = std::path::Path::new("/run/subc-connection.json");
8489 let handle = SupervisorHandle::new();
8490
8491 let mut none_spec = spec(Vec::new());
8492 none_spec.protocol = ModuleProtocol::None;
8493 let mut none = Command::new("/nonexistent");
8494 apply_wire_spawn_args(&mut none, &none_spec, Some(connection_file), Some(&handle))
8495 .expect("protocol-none spawn args apply");
8496 let none_args: Vec<String> = none
8497 .as_std()
8498 .get_args()
8499 .map(|a| a.to_string_lossy().into_owned())
8500 .collect();
8501 assert!(
8502 !none_args.iter().any(|a| a == SUBC_ARG),
8503 "protocol:none argv must not carry --subc; got {none_args:?}"
8504 );
8505 let none_has_nonce = none
8506 .as_std()
8507 .get_envs()
8508 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some());
8509 assert!(
8510 !none_has_nonce,
8511 "protocol:none spawn must not receive a launch nonce"
8512 );
8513 let none_has_module_id = none
8514 .as_std()
8515 .get_envs()
8516 .any(|(key, value)| key == OsStr::new(SUBC_MODULE_ID_ENV) && value.is_some());
8517 assert!(
8518 none_has_module_id,
8519 "SUBC_MODULE_ID is inert and stays on every path"
8520 );
8521 assert!(
8522 handle.spawn_nonce(&none_spec.module_id).is_none(),
8523 "no nonce record for a process that will never present one"
8524 );
8525
8526 let wire_spec = spec(Vec::new());
8528 let mut wire = Command::new("/nonexistent");
8529 apply_wire_spawn_args(&mut wire, &wire_spec, Some(connection_file), Some(&handle))
8530 .expect("subc-wire spawn args apply");
8531 let wire_args: Vec<String> = wire
8532 .as_std()
8533 .get_args()
8534 .map(|a| a.to_string_lossy().into_owned())
8535 .collect();
8536 assert_eq!(
8537 wire_args,
8538 vec![
8539 SUBC_ARG.to_string(),
8540 connection_file.to_string_lossy().into_owned()
8541 ],
8542 "a subc-wire spawn still carries --subc <path>"
8543 );
8544 assert!(wire
8545 .as_std()
8546 .get_envs()
8547 .any(|(key, value)| key == OsStr::new(SUBC_LAUNCH_NONCE_ENV) && value.is_some()));
8548 assert!(handle.spawn_nonce(&wire_spec.module_id).is_some());
8549 }
8550
8551 #[test]
8561 fn plain_spawn_removes_the_spawn_role_even_when_the_spec_sets_it() {
8562 let role = |command: &Command| {
8563 command
8564 .as_std()
8565 .get_envs()
8566 .filter(|(key, _)| *key == OsStr::new(SUBC_SPAWN_ROLE_ENV))
8567 .last()
8568 .map(|(_, value)| value.map(|v| v.to_string_lossy().into_owned()))
8569 };
8570 let forged = spec(vec![(
8571 SUBC_SPAWN_ROLE_ENV.to_string(),
8572 SPAWN_ROLE_SWAP_CANDIDATE.to_string(),
8573 )]);
8574
8575 let mut plain = Command::new("/nonexistent");
8576 apply_child_env(&mut plain, &forged);
8577 apply_spawn_role(&mut plain, SpawnRole::Plain);
8578 assert_eq!(
8579 role(&plain),
8580 Some(None),
8581 "a plain spawn must remove SUBC_SPAWN_ROLE, whatever the spec says"
8582 );
8583
8584 let mut candidate = Command::new("/nonexistent");
8585 apply_child_env(&mut candidate, &spec(Vec::new()));
8586 apply_spawn_role(&mut candidate, SpawnRole::SwapCandidate);
8587 assert_eq!(
8588 role(&candidate),
8589 Some(Some(SPAWN_ROLE_SWAP_CANDIDATE.to_string()))
8590 );
8591 }
8592
8593 #[test]
8599 fn daemon_private_capture_keys_are_not_passed_to_the_child() {
8600 let mut command = Command::new("/nonexistent");
8601 apply_child_env(
8602 &mut command,
8603 &spec(vec![
8604 (super::CAPTURE_KEEP_ENV.to_string(), "5".to_string()),
8605 ("KEPT".to_string(), "yes".to_string()),
8606 ]),
8607 );
8608 let keys: Vec<String> = command
8609 .as_std()
8610 .get_envs()
8611 .filter(|(_, value)| value.is_some())
8612 .map(|(key, _)| key.to_string_lossy().into_owned())
8613 .collect();
8614 assert!(keys.contains(&"KEPT".to_string()), "got {keys:?}");
8615 assert!(
8616 !keys.contains(&super::CAPTURE_KEEP_ENV.to_string()),
8617 "daemon-private capture key leaked to the child: {keys:?}"
8618 );
8619 }
8620}
8621
8622#[cfg(test)]
8623mod jitter_tests {
8624 use super::jittered_health_delay;
8625 use std::{collections::HashSet, time::Duration};
8626
8627 const FLEET: [&str; 14] = [
8636 "aft",
8637 "alfonso-core",
8638 "magic-context",
8639 "broca",
8640 "thalamus",
8641 "quota",
8642 "engram",
8643 "plexus",
8644 "cerebellum",
8645 "astrocyte",
8646 "synapse",
8647 "subc-mcp",
8648 "cortexkit-credentials",
8649 "subc-federation",
8650 ];
8651
8652 #[test]
8660 fn probe_delays_disperse_across_the_fleet() {
8661 let cadence = Duration::from_secs(30);
8662 let delays: HashSet<Duration> = FLEET
8663 .iter()
8664 .map(|id| jittered_health_delay(id, 0, cadence))
8665 .collect();
8666 assert_eq!(
8667 delays.len(),
8668 FLEET.len(),
8669 "every supervised module must land on its own probe offset"
8670 );
8671 }
8672
8673 #[test]
8679 fn jitter_only_delays_and_stays_within_one_tenth_of_cadence() {
8680 let cadence = Duration::from_secs(30);
8681 let span = cadence / 10;
8682 for id in FLEET {
8683 for probe_index in 0..8 {
8684 let delay = jittered_health_delay(id, probe_index, cadence);
8685 assert!(
8686 delay >= cadence,
8687 "{id}#{probe_index}: jitter must not shorten the cadence"
8688 );
8689 assert!(
8690 delay < cadence + span,
8691 "{id}#{probe_index}: jitter must stay inside one tenth of the cadence"
8692 );
8693 }
8694 }
8695 }
8696
8697 #[test]
8703 fn a_module_offset_is_stable_across_restarts() {
8704 let cadence = Duration::from_secs(30);
8705 for id in FLEET {
8706 assert_eq!(
8707 jittered_health_delay(id, 0, cadence),
8708 jittered_health_delay(id, 0, cadence),
8709 "{id}: the same module and probe index must produce the same offset"
8710 );
8711 }
8712 }
8713
8714 #[test]
8716 fn zero_cadence_yields_zero_delay() {
8717 assert_eq!(
8718 jittered_health_delay("aft", 0, Duration::ZERO),
8719 Duration::ZERO
8720 );
8721 }
8722}
8723
8724#[cfg(all(test, target_os = "linux"))]
8725mod cgroup_placement_tests {
8726 use super::{
8727 apply_cgroup_placement, remove_module_cgroup, ModuleProtocol, ModuleSpec, SuperviseError,
8728 SupervisedChild,
8729 };
8730 use crate::{
8731 stderr_tail::{StderrRing, StderrTailConfig},
8732 test_support::TestTempDir,
8733 };
8734 use std::{
8735 fs, io,
8736 path::{Path, PathBuf},
8737 sync::{Arc, Mutex},
8738 };
8739 use tokio::process::Command;
8740
8741 #[test]
8742 fn failed_parent_cgroup_open_is_a_cgroup_supervision_error() {
8743 let path = Path::new("/definitely-missing-subc-cgroup");
8744 let mut command = Command::new("true");
8745 let error = apply_cgroup_placement(
8746 &mut command,
8747 &ModuleSpec {
8748 module_id: "broken-cgroup".to_string(),
8749 program: PathBuf::from("true"),
8750 args: Vec::new(),
8751 env: Vec::new(),
8752 reserved: false,
8753 reserved_prefixes: Vec::new(),
8754 protocol: ModuleProtocol::Subc,
8755 overlap: Default::default(),
8756 },
8757 path,
8758 )
8759 .expect_err("a parent cgroup open failure must reject the supervised spawn");
8760 let reason = error.to_string();
8761
8762 assert!(
8763 matches!(error, SuperviseError::Cgroup { .. }),
8764 "parent cgroup open must be reported as a cgroup supervision error: {reason}"
8765 );
8766 assert!(
8767 reason.contains("/definitely-missing-subc-cgroup/cgroup.procs"),
8768 "parent cgroup open failure must name cgroup.procs: {reason}"
8769 );
8770 }
8771
8772 #[tokio::test]
8773 async fn reaping_a_child_removes_its_empty_module_cgroup() {
8774 let root = TestTempDir::new("supervisor-reap-cgroup");
8775 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
8776 let placement = subc_cgroup::prepare_at(&root)
8777 .expect("prepare scratch cgroup root")
8778 .expect("scratch root has a cgroup.procs marker");
8779 let module_id = "reaped-module";
8780 let module = placement
8781 .module_path(module_id)
8782 .expect("create scratch module cgroup");
8783 let child = Command::new("true")
8784 .spawn()
8785 .expect("spawn short-lived child");
8786 let pid = child.id().expect("spawned child has pid");
8787 let mut child = SupervisedChild {
8788 child,
8789 module_id: module_id.to_string(),
8790 cgroup_placement: Some(placement),
8791 stdout_pump: None,
8792 stderr_pump: None,
8793 stderr_ring: Arc::new(Mutex::new(StderrRing::new(StderrTailConfig::default()))),
8794 spawned_at_ms: 0,
8795 spawned_from: PathBuf::from("true"),
8796 spawned_file_identity: None,
8797 process_start_time: None,
8798 process_identity: None,
8799 pid,
8800 roster_guard: None,
8801 };
8802
8803 child.wait().await.expect("reap short-lived child");
8804
8805 assert!(
8806 !module.exists(),
8807 "reaping the supervised child must remove its empty cgroup"
8808 );
8809 }
8810
8811 #[test]
8812 fn non_empty_cgroup_removal_is_reported_without_blocking_teardown() {
8813 let root = TestTempDir::new("supervisor-non-empty-cgroup");
8814 fs::write(root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
8815 let placement = subc_cgroup::prepare_at(&root)
8816 .expect("prepare scratch cgroup root")
8817 .expect("scratch root has a cgroup.procs marker");
8818 let module = placement
8819 .module_path("surviving-module")
8820 .expect("create scratch module cgroup");
8821 fs::write(module.join("surviving-process"), b"still present")
8822 .expect("make scratch cgroup non-empty");
8823 let (logs, _guard) = crate::router::test_log::log_capture(tracing::Level::WARN);
8824
8825 remove_module_cgroup(&placement, "surviving-module");
8826
8827 let logs = crate::router::test_log::captured_logs(&logs);
8828 assert!(
8829 module.exists(),
8830 "failed removal must leave the cgroup intact"
8831 );
8832 assert!(
8833 logs.contains("could not remove module cgroup after process exit; continuing teardown")
8834 && logs.contains("surviving-module"),
8835 "best-effort removal must report the failure without returning it: {logs}"
8836 );
8837 }
8838
8839 #[test]
8840 fn cgroup_pre_exec_spawn_failure_names_the_cgroup_path() {
8841 let cgroup_path = PathBuf::from("/sys/fs/cgroup/subc-modules/broken-module");
8842 let reason = SuperviseError::Spawn {
8843 program: PathBuf::from("/bin/true"),
8844 source: io::Error::from_raw_os_error(13),
8845 cgroup_path: Some(cgroup_path.clone()),
8846 }
8847 .to_string();
8848
8849 assert!(
8850 reason.contains(&cgroup_path.display().to_string()),
8851 "a pre_exec spawn failure must name the cgroup path: {reason}"
8852 );
8853 }
8854}
8855
8856#[cfg(test)]
8857mod spawn_subscriber_lag_tests {
8858 use super::*;
8859
8860 #[tokio::test]
8865 async fn lagged_spawn_subscriber_receives_a_terminal_lagged_error_after_its_queued_frames() {
8866 let feed = SpawnEventFeed::default();
8867 feed.configure_incarnation("lag-incarnation".to_string());
8868 let (tx, mut rx) = mpsc::channel(1);
8871 feed.subscribe(ConnectionId::new(1), 7, 1, None, FrameSink::new(tx))
8872 .expect("subscribe");
8873 let emitted = SPAWN_SUBSCRIBER_BUFFER + 16;
8874 for index in 0..emitted {
8875 feed.emit_spawned(&format!("lag-module-{index}"), 1000, 0);
8876 tokio::task::yield_now().await;
8879 }
8880 assert_eq!(
8881 feed.subscriber_count(),
8882 0,
8883 "the lagged subscriber must be removed"
8884 );
8885
8886 let mut data = Vec::new();
8887 let mut last = None;
8888 loop {
8889 let next = tokio::time::timeout(Duration::from_secs(5), rx.recv())
8890 .await
8891 .expect("the forwarder must finish once the subscriber is dropped");
8892 let Some(outbound) = next else { break };
8893 let frame = outbound.frame;
8894 if frame.header.ty == FrameType::StreamData {
8895 assert!(last.is_none(), "no data may follow the terminal frame");
8896 let event: SpawnEvent = serde_json::from_slice(&frame.body).unwrap();
8897 data.push(event.cursor.seq);
8898 } else {
8899 assert!(last.is_none(), "exactly one terminal frame");
8900 last = Some(frame);
8901 }
8902 }
8903 assert!(!data.is_empty(), "queued frames drain before the terminal");
8904 for pair in data.windows(2) {
8905 assert_eq!(
8906 pair[1],
8907 pair[0] + 1,
8908 "queued frames arrive dense and in order"
8909 );
8910 }
8911 let terminal = last.expect("a lagged subscriber must receive a terminal frame");
8912 assert_eq!(terminal.header.ty, FrameType::Error);
8913 assert_eq!(terminal.header.corr, 7);
8914 let body: subc_protocol::ErrorBody = serde_json::from_slice(&terminal.body).unwrap();
8915 assert_eq!(body.code, SPAWN_SUBSCRIBER_LAGGED_CODE);
8916 let detail = body.detail.expect("lagged error carries detail");
8917 assert_eq!(
8918 detail["first_undelivered_cursor"]["seq"],
8919 data.last().unwrap() + 1,
8920 "the named cursor is the first event the subscriber did not receive"
8921 );
8922 assert_eq!(
8923 detail["first_undelivered_cursor"]["daemon_incarnation"],
8924 "lag-incarnation"
8925 );
8926 }
8927}
8928
8929#[cfg(test)]
8930mod terminal_history_read_concurrency_tests {
8931 use super::*;
8932 use crate::{terminal_journal::read_pause, test_support::TestTempDir};
8933 use std::sync::mpsc as std_mpsc;
8934
8935 fn journaled_ring(
8936 journal: &Arc<crate::terminal_journal::TerminalJournal>,
8937 ) -> Arc<Mutex<TerminalRing>> {
8938 Arc::new(Mutex::new(
8939 TerminalRing::new(TerminalRingConfig::default(), 1)
8940 .with_journal(Some(Arc::clone(journal))),
8941 ))
8942 }
8943
8944 fn crash(at_ms: u64) -> ExitReport {
8945 ExitReport {
8946 kind: ExitKind::Crash,
8947 code: Some(1),
8948 signal: None,
8949 at_ms,
8950 }
8951 }
8952
8953 fn record_within(
8956 module_id: &'static str,
8957 ring: &Arc<Mutex<TerminalRing>>,
8958 at_ms: u64,
8959 bound: Duration,
8960 ) -> bool {
8961 let ring = Arc::clone(ring);
8962 let (done, done_rx) = std_mpsc::channel();
8963 std::thread::spawn(move || {
8964 record_terminal(
8965 module_id,
8966 &ring,
8967 &SpawnEventFeed::default(),
8968 &crash(at_ms),
8969 TerminalDisposition::Restarting,
8970 );
8971 let _ = done.send(());
8972 });
8973 done_rx.recv_timeout(bound).is_ok()
8974 }
8975
8976 #[test]
8981 fn exits_recorded_during_a_paused_history_read_are_not_blocked_or_half_merged() {
8982 let dir = TestTempDir::new("terminal-history-concurrent-read");
8983 let path = dir.join("terminals.jsonl");
8984 let journal = Arc::new(crate::terminal_journal::TerminalJournal::open(
8985 path.clone(),
8986 "daemon".into(),
8987 ));
8988 let reader_ring = journaled_ring(&journal);
8989 let other_ring = journaled_ring(&journal);
8990 assert!(record_within(
8991 "reader-module",
8992 &reader_ring,
8993 10,
8994 Duration::from_secs(5)
8995 ));
8996
8997 let (started, release) = read_pause::install(&path);
8998 let reading = {
8999 let ring = Arc::clone(&reader_ring);
9000 std::thread::spawn(move || durable_terminal_history_of(&ring, "reader-module"))
9001 };
9002 started
9003 .recv_timeout(Duration::from_secs(5))
9004 .expect("the history read reached its pause");
9005
9006 let bound = Duration::from_secs(1);
9007 assert!(
9008 record_within("other-module", &other_ring, 20, bound),
9009 "another module's exit waited on a history read (journal writer held)"
9010 );
9011 assert!(
9012 record_within("reader-module", &reader_ring, 30, bound),
9013 "the read module's own exit waited on its history read (ring held)"
9014 );
9015
9016 drop(release);
9017 let paused = reading.join().unwrap();
9018 assert_eq!(
9019 paused.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9020 vec![10],
9021 "an exit recorded after the read began lands in neither half of it"
9022 );
9023 assert_eq!(paused.journal_skipped_lines, 0);
9024 assert_eq!(paused.journal_read_errors, 0);
9025
9026 let after = durable_terminal_history_of(&reader_ring, "reader-module");
9027 assert_eq!(
9028 after.entries.iter().map(|e| e.at_ms).collect::<Vec<_>>(),
9029 vec![10, 30],
9030 "the next read merges ring and journal with no duplicate"
9031 );
9032 assert_eq!(after.journal_skipped_lines, 0);
9033 }
9034}