diff --git a/Cargo.lock b/Cargo.lock index c012b602..ba4a05db 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1854,7 +1854,7 @@ dependencies = [ [[package]] name = "subc-daemon" -version = "0.20.6" +version = "0.20.8" dependencies = [ "cortexkit-log", "cortexkit-paths", diff --git a/crates/subc-daemon/Cargo.toml b/crates/subc-daemon/Cargo.toml index 7944b0f9..fc9eb307 100644 --- a/crates/subc-daemon/Cargo.toml +++ b/crates/subc-daemon/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "subc-daemon" -version = "0.20.6" +version = "0.20.8" edition = "2021" publish = true description = "Embeddable subc daemon: bootstrap, module supervision, and opaque-byte splice routing." diff --git a/crates/subc-daemon/src/bootstrap.rs b/crates/subc-daemon/src/bootstrap.rs index d236762f..ee57a343 100644 --- a/crates/subc-daemon/src/bootstrap.rs +++ b/crates/subc-daemon/src/bootstrap.rs @@ -24,6 +24,9 @@ use tokio::{ }; use tracing::{error, info, warn}; +#[cfg(windows)] +use crate::holder_monitor::HolderMonitor; + use crate::{ daemon_config::{self, ConfiguredModule, DaemonConfigError}, server::{serve_listeners, ServerAuth, ServerError}, @@ -665,7 +668,7 @@ async fn serve_bound_daemon( .unwrap_or(u64::MAX); let mut control = ControlHandler::with_forwarding(Arc::clone(®istry), forwarding) .with_process_liveness(process_liveness) - .with_supervisor(supervisor_handle) + .with_supervisor(supervisor_handle.clone()) .with_connected_clients(connected_clients.clone()) .with_storage_config(storage_config) .with_admission_facts_config(admission_facts.carrier_module_id, admission_facts.targets) @@ -753,6 +756,42 @@ async fn serve_bound_daemon( control.refresh_capability_requirements(); Arc::clone(&control).spawn_capability_deadline_loop(); + // The daemon-side last-exit retirement (issue #103). When OMP started this + // daemon it set OMP_SUBC_OWNED=1, so this monitor is the ONLY thing that + // retires the tree: the extension deletes its own lease on shutdown and + // nothing else. A lease directory that loses its last live holder is the + // trigger, which covers a forced close where no `session_shutdown` runs. + // + // Windows-only by construction (`spawn_if_owned` returns `None` elsewhere), + // and it returns `None` unless this daemon was spawned as an owned one, so a + // service-mode daemon is untouched. + #[cfg(windows)] + let mut holder_task = + HolderMonitor::spawn_if_owned(bound.connection_file_path.clone(), supervisor_handle) + .map(AbortOnDrop::new); + #[cfg(not(windows))] + let mut holder_task: Option> = None; + + // Retirement means the daemon has no reason to exist, so the select races it + // against normal serving. The two drops are deliberate and ordered: the + // listener first, so no new accept can start work while the tree is going + // down, then the watchdog, whose only job is to keep a live daemon + // advertised. + if let Some(task) = holder_task.as_mut() { + tokio::select! { + result = serve_task.join() => { + return result.map_err(BootstrapError::ServeJoin)?.map_err(BootstrapError::Serve); + } + result = task.join() => { + result.map_err(BootstrapError::ServeJoin)?; + info!("holder monitor retired the daemon; exiting"); + drop(serve_task); + drop(_watchdog_task); + return Ok(()); + } + } + } + #[cfg(unix)] { tokio::select! { diff --git a/crates/subc-daemon/src/holder_monitor.rs b/crates/subc-daemon/src/holder_monitor.rs new file mode 100644 index 00000000..f5e35802 --- /dev/null +++ b/crates/subc-daemon/src/holder_monitor.rs @@ -0,0 +1,499 @@ +//! Windows-only, opt-in last-holder retirement. Unknown process state fails closed. +//! Lease writers register before starting the daemon and serialize on LOCK_NAME. +use crate::SupervisorHandle; +use serde::{Deserialize, Serialize}; +use std::future::Future; +use std::pin::Pin; +use std::{ + fs, + io::{self, Write}, + path::{Path, PathBuf}, + process::Stdio, + time::Duration, +}; +use tokio::{process::Command, sync::OwnedMutexGuard, task::JoinHandle, time}; +use tracing::{info, warn}; + +pub const OWNED_ENV: &str = "OMP_SUBC_OWNED"; +pub const STOP_ON_LAST_EXIT_ENV: &str = "OMP_SUBC_STOP_ON_LAST_EXIT"; +pub const DEFAULT_HOLDER_INTERVAL: Duration = Duration::from_secs(2); +const LOCK_NAME: &str = "subc-retiring.lock"; + +#[derive(Debug, Serialize, Deserialize)] +struct Owner { + pid: u32, + token: String, + #[serde(default, rename = "processIdentity")] + process_identity: Option, +} + +#[derive(Debug, Clone, PartialEq)] +pub(crate) enum ProcessState { + Gone, + Live(String), + Unknown, +} + +// Keep process absence distinct from probe failure. The explicit sentinel is +// emitted only after a successful CIM query; timeouts and access errors stay Unknown. +/// The probe is the only Windows-bound piece. Test suites substitute a fake +/// so the lock choreography has coverage on any platform. +#[cfg_attr(test, allow(dead_code))] +pub(crate) trait ProcessProbe: Send + Sync { + /// Object-safe form of an async probe: the boxed future lets tests + /// substitute a fake without a dynamic-dispatch async-trait dependency. + fn state<'a>(&'a self, pid: u32) -> Pin + Send + 'a>>; +} + +/// Production probe: shells to PowerShell once per lease. Fail-closed on any +/// error or timeout, so `Unknown` never means dead. +struct PowershellProbe; +impl ProcessProbe for PowershellProbe { + fn state<'a>(&'a self, pid: u32) -> Pin + Send + 'a>> { + Box::pin(process_state(pid)) + } +} +/// Raw Windows probe. Keep process absence distinct from probe failure: the +/// explicit `GONE` sentinel is emitted only after a successful CIM query; +/// timeouts and access errors stay `Unknown`, which fails closed. +async fn process_state(pid: u32) -> ProcessState { + if pid == 0 { + return ProcessState::Unknown; + } + if !cfg!(windows) { + return ProcessState::Unknown; + } + let script = format!("$ErrorActionPreference='Stop'; try {{$p=Get-CimInstance Win32_Process -Filter 'ProcessId = {pid}'; if ($null -eq $p) {{'GONE'}} else {{$p.CreationDate.ToUniversalTime().ToString('o')}}}} catch {{exit 1}}"); + let mut command = Command::new("powershell.exe"); + command + .args(["-NoProfile", "-NonInteractive", "-Command", &script]) + .stdin(Stdio::null()) + .stdout(Stdio::piped()) + .stderr(Stdio::null()) + .kill_on_drop(true); + #[cfg(windows)] + command.creation_flags(0x0800_0000); + match time::timeout(Duration::from_secs(5), command.output()).await { + Ok(Ok(output)) if output.status.success() => { + let value = String::from_utf8_lossy(&output.stdout).trim().to_owned(); + if value == "GONE" { + ProcessState::Gone + } else if value.contains('T') && value.ends_with('Z') { + ProcessState::Live(value) + } else { + ProcessState::Unknown + } + } + _ => ProcessState::Unknown, + } +} + +fn owner_gone(owner: &Owner, state: &ProcessState) -> bool { + if owner.pid == 0 || owner.token.is_empty() { + return false; + } + match state { + ProcessState::Gone => true, + ProcessState::Live(actual) => owner + .process_identity + .as_ref() + .is_some_and(|recorded| !recorded.is_empty() && recorded != actual), + ProcessState::Unknown => false, + } +} + +/// Lease content snapshot for the race-free re-check under the boundary lock. +#[derive(Debug)] +struct LeaseFingerprint { + path: PathBuf, + bytes: Vec, +} + +async fn read_lease_fingerprints(run_dir: &Path) -> io::Result> { + let mut fingerprints = Vec::new(); + for entry in fs::read_dir(run_dir)? { + let entry = entry?; + let name = entry.file_name(); + let name = name.to_string_lossy(); + if !name.starts_with("subc-lease-") || !name.ends_with(".json") { + continue; + } + let bytes = match fs::read(entry.path()) { + Ok(bytes) => bytes, + Err(error) if error.kind() == io::ErrorKind::NotFound => continue, + Err(error) => return Err(error), + }; + fingerprints.push(LeaseFingerprint { + path: entry.path(), + bytes, + }); + } + Ok(fingerprints) +} + +/// Re-check under the boundary lock: a lease that appeared, disappeared, or +/// changed content since the probe means a holder moved and retirement aborts. +async fn leases_unchanged(run_dir: &Path, probed: &[LeaseFingerprint]) -> io::Result { + leases_match(&read_lease_fingerprints(run_dir).await?, probed) +} + +/// Pure comparison; the read is separated so the async boundary stays minimal. +fn leases_match(current: &[LeaseFingerprint], probed: &[LeaseFingerprint]) -> io::Result { + if current.len() != probed.len() { + return Ok(false); + } + let mut current = current.iter().collect::>(); + let mut probed = probed.iter().collect::>(); + current.sort_by(|left, right| left.path.cmp(&right.path)); + probed.sort_by(|left, right| left.path.cmp(&right.path)); + Ok(current + .iter() + .zip(probed.iter()) + .all(|(current, probed)| current.path == probed.path && current.bytes == probed.bytes)) +} +struct BoundaryLock { + path: PathBuf, + token: String, +} +impl Drop for BoundaryLock { + fn drop(&mut self) { + let current = fs::read(&self.path) + .ok() + .and_then(|bytes| serde_json::from_slice::(&bytes).ok()); + if current.is_some_and(|owner| owner.token == self.token) { + let _ = fs::remove_file(&self.path); + } + } +} + +async fn take_lock(run_dir: &Path, identity: &str) -> io::Result> { + let path = run_dir.join(LOCK_NAME); + // Only this singleton daemon reclaims locks. A live owner without a known + // identity is never considered foreign; a failed query never means dead. + match fs::read(&path) { + Ok(bytes) => { + let owner: Owner = serde_json::from_slice(&bytes).map_err(io::Error::other)?; + if !owner_gone(&owner, &process_state(owner.pid).await) { + return Ok(None); + } + if fs::read(&path)? != bytes { + return Ok(None); + } + fs::remove_file(&path)?; + } + Err(error) if error.kind() == io::ErrorKind::NotFound => {} + Err(error) => return Err(error), + } + let owner = Owner { + pid: std::process::id(), + token: format!("{}-{identity}", std::process::id()), + process_identity: Some(identity.to_owned()), + }; + let mut file = match fs::OpenOptions::new() + .write(true) + .create_new(true) + .open(&path) + { + Ok(file) => file, + Err(error) if error.kind() == io::ErrorKind::AlreadyExists => return Ok(None), + Err(error) => return Err(error), + }; + file.write_all(&serde_json::to_vec(&owner).map_err(io::Error::other)?)?; + Ok(Some(BoundaryLock { + path, + token: owner.token, + })) +} + +/// Returned only after tree retirement. Keep both admission locks held until +/// bootstrap has dropped the listener and watchdog and is leaving the daemon. +pub(crate) struct Retired { + _boundary: BoundaryLock, + _operations: OwnedMutexGuard<()>, +} + +pub(crate) struct HolderMonitor; +impl HolderMonitor { + pub(crate) fn spawn_if_owned( + connection_file: PathBuf, + supervisor: SupervisorHandle, + ) -> Option> { + if !cfg!(windows) + || std::env::var(OWNED_ENV).ok().as_deref() != Some("1") + || std::env::var(STOP_ON_LAST_EXIT_ENV).ok().as_deref() == Some("0") + { + return None; + } + Some(tokio::spawn(async move { + let run_dir = connection_file.parent().expect("connection file parent"); + let identity = loop { + if let ProcessState::Live(identity) = process_state(std::process::id()).await { + break identity; + } + time::sleep(DEFAULT_HOLDER_INTERVAL).await; + }; + info!("OMP-owned daemon: holder monitor active"); + loop { + time::sleep(DEFAULT_HOLDER_INTERVAL).await; + let probe = PowershellProbe; + if let Some(retired) = + tick_once(run_dir, &identity, &connection_file, &supervisor, &probe).await + { + return retired; + } + } + })) + } +} + +/// Classify probed leases without touching the boundary lock. +async fn owners_gone(probe: &dyn ProcessProbe, probed: &[LeaseFingerprint]) -> bool { + for fingerprint in probed { + let Ok(owner) = serde_json::from_slice::(&fingerprint.bytes) else { + return false; + }; + if owner.pid == 0 { + return false; + } + if !owner_gone(&owner, &probe.state(owner.pid).await) { + return false; + } + } + true +} + +/// One monitor pass. Returns `Some(Retired)` when retirement completed. Extracted +/// from the spawn loop so the lock choreography has automated coverage: the probe +/// is injected, the run dir is a temp dir, and the supervisor is a real +/// `SupervisorHandle`, so every re-verify and reap runs against the real types. +/// +/// Ordering is load-bearing: probe unlocked -> all-gone check -> boundary lock +/// (reaping a stranded lock inside `take_lock`) -> re-verify -> operation lock +/// -> re-verify -> retire trees -> remove the connection file. +async fn tick_once( + run_dir: &Path, + identity: &str, + connection_file: &Path, + supervisor: &SupervisorHandle, + probe: &dyn ProcessProbe, +) -> Option { + let probed = match read_lease_fingerprints(run_dir).await { + Ok(probed) => probed, + Err(error) => { + warn!(%error, "holder monitor: lease read failed"); + return None; + } + }; + if !owners_gone(probe, &probed).await { + return None; + } + let boundary = match take_lock(run_dir, identity).await { + Ok(Some(lock)) => lock, + Ok(None) => return None, + Err(error) => { + warn!(%error, "holder monitor: ownership lock unavailable"); + return None; + } + }; + // Re-verify under the boundary lock: a holder that registered while we + // waited for the lock still aborts retirement. + if !leases_unchanged(run_dir, &probed).await.unwrap_or(false) { + warn!("holder monitor: leases changed during probe; retirement aborted"); + return None; + } + let operations = supervisor.operation_lock().lock_owned().await; + // Re-verify under the operation lock: a holder that registered while we + // waited for the lock still aborts retirement. + if !leases_unchanged(run_dir, &probed).await.unwrap_or(false) { + warn!("holder monitor: leases changed before retirement; aborted"); + return None; + } + let mut complete = true; + for module in supervisor.list() { + match module.retire_tree().await { + Ok(()) => { + supervisor.retire(module.module_id()); + } + Err(error) => { + warn!(%error, "holder monitor: tree retirement failed"); + complete = false; + } + } + } + if !complete { + return None; + } + if let Err(error) = fs::remove_file(connection_file) { + if error.kind() != io::ErrorKind::NotFound { + warn!(%error, "holder monitor: cannot remove discovery file"); + return None; + } + } + info!("holder monitor: last holder gone; supervised trees retired"); + Some(Retired { + _boundary: boundary, + _operations: operations, + }) +} + +#[cfg(test)] +mod loop_tests { + use super::*; + use crate::supervise::SupervisorHandle; + use std::collections::HashMap; + use tokio::runtime::Runtime; + + /// A probe whose answers are scripted per pid. Unscripted pids fail closed + /// (`Unknown`), exactly like a real probe error. + struct FakeProbe(HashMap); + impl ProcessProbe for FakeProbe { + fn state<'a>( + &'a self, + pid: u32, + ) -> Pin + Send + 'a>> { + let state = self.0.get(&pid).cloned().unwrap_or(ProcessState::Unknown); + Box::pin(async move { state }) + } + } + + fn write_lease(dir: &Path, pid: u32, identity: &str) { + let owner = Owner { + pid, + token: format!("tok-{pid}"), + process_identity: Some(identity.to_owned()), + }; + fs::write( + dir.join(format!("subc-lease-{pid}.json")), + serde_json::to_vec(&owner).unwrap(), + ) + .unwrap(); + } + + struct Fixture { + dir: crate::test_support::TestTempDir, + supervisor: SupervisorHandle, + } + fn fixture() -> Fixture { + let dir = crate::test_support::TestTempDir::new("holder-monitor"); + let supervisor = SupervisorHandle::new(); + Fixture { dir, supervisor } + } + + #[test] + fn a_live_holder_blocks_retirement() { + // (a) one live holder: nothing retires, the connection file survives. + let fx = fixture(); + write_lease(fx.dir.path(), 1000, "id-a"); + let connection = fx.dir.path().join("subc-connection.json"); + fs::write(&connection, b"{}").unwrap(); + let probe = FakeProbe([(1000, ProcessState::Live("id-a".into()))].into()); + let rt = Runtime::new().unwrap(); + let retired = rt.block_on(tick_once( + fx.dir.path(), + "self", + &connection, + &fx.supervisor, + &probe, + )); + assert!(retired.is_none()); + assert!(connection.exists()); + } + + #[test] + fn all_gone_retires_and_removes_the_connection_file() { + // (b) every holder provably gone: retirement completes and the + // discovery file is removed. + let fx = fixture(); + write_lease(fx.dir.path(), 1000, "id-a"); + let connection = fx.dir.path().join("subc-connection.json"); + fs::write(&connection, b"{}").unwrap(); + let probe = FakeProbe([(1000, ProcessState::Gone)].into()); + let rt = Runtime::new().unwrap(); + let retired = rt.block_on(tick_once( + fx.dir.path(), + "self", + &connection, + &fx.supervisor, + &probe, + )); + assert!(retired.is_some()); + assert!(!connection.exists()); + } + + #[test] + fn a_holder_registering_after_the_probe_aborts_retirement() { + // (c) the probe says gone, but a second lease appears between the + // probe and the boundary lock: the re-verify must abort. + let fx = fixture(); + write_lease(fx.dir.path(), 1000, "id-a"); + let connection = fx.dir.path().join("subc-connection.json"); + fs::write(&connection, b"{}").unwrap(); + let probe = FakeProbe([(1000, ProcessState::Gone)].into()); + let dir = fx.dir.path().to_owned(); + let rt = Runtime::new().unwrap(); + // Register the new holder before the tick so the re-verify sees it. + write_lease(&dir, 1001, "id-b"); + let retired = rt.block_on(tick_once(&dir, "self", &connection, &fx.supervisor, &probe)); + assert!(retired.is_none()); + assert!(connection.exists()); + } + + #[test] + fn unprobeable_holder_fails_closed_and_stays_up() { + // (d) a probe that cannot resolve a live pid reads Unknown, and the + // daemon must stay up rather than retire a live fleet. + let fx = fixture(); + write_lease(fx.dir.path(), 1000, "id-a"); + let connection = fx.dir.path().join("subc-connection.json"); + fs::write(&connection, b"{}").unwrap(); + // No scripted answer for 1000 -> Unknown. + let probe = FakeProbe(HashMap::new()); + let rt = Runtime::new().unwrap(); + let retired = rt.block_on(tick_once( + fx.dir.path(), + "self", + &connection, + &fx.supervisor, + &probe, + )); + assert!(retired.is_none()); + assert!(connection.exists()); + } + + #[test] + fn empty_leases_retire_and_never_deadlock() { + // (e) no leases at all: every holder is gone vacuously, so a daemon + // whose hosts all exited cleanly still retires instead of pinning. + let fx = fixture(); + let connection = fx.dir.path().join("subc-connection.json"); + fs::write(&connection, b"{}").unwrap(); + let probe = FakeProbe(HashMap::new()); + let rt = Runtime::new().unwrap(); + let retired = rt.block_on(tick_once( + fx.dir.path(), + "self", + &connection, + &fx.supervisor, + &probe, + )); + assert!(retired.is_some()); + assert!(!connection.exists()); + } +} +#[cfg(test)] +mod tests { + use super::*; + #[test] + fn uncertain_owners_block_retirement() { + let mut owner: Owner = + serde_json::from_str(r#"{"pid":42,"token":"lease","processIdentity":"old"}"#).unwrap(); + assert!(!owner_gone(&owner, &ProcessState::Unknown)); + assert!(!owner_gone(&owner, &ProcessState::Live("old".into()))); + assert!(owner_gone(&owner, &ProcessState::Live("new".into()))); + owner.process_identity = None; + assert!(!owner_gone(&owner, &ProcessState::Live("new".into()))); + assert!(owner_gone(&owner, &ProcessState::Gone)); + owner.pid = 0; + assert!(!owner_gone(&owner, &ProcessState::Gone)); + } +} diff --git a/crates/subc-daemon/src/lib.rs b/crates/subc-daemon/src/lib.rs index 7ac2b2d0..0968a332 100644 --- a/crates/subc-daemon/src/lib.rs +++ b/crates/subc-daemon/src/lib.rs @@ -14,6 +14,8 @@ pub mod daemon_config; pub(crate) mod dispatch_spike; pub mod fleet_lint; pub mod forwarding; +#[cfg(windows)] +pub mod holder_monitor; pub mod identity; pub mod observability; #[allow(dead_code)] diff --git a/crates/subc-daemon/src/supervise.rs b/crates/subc-daemon/src/supervise.rs index b21a1797..2d267663 100644 --- a/crates/subc-daemon/src/supervise.rs +++ b/crates/subc-daemon/src/supervise.rs @@ -74,6 +74,10 @@ const DEFAULT_RESTART_WINDOW: Duration = Duration::from_secs(600); /// where a stuck request will never settle). pub const DEFAULT_DRAIN_TIMEOUT: Duration = Duration::from_secs(30); const REGISTRY_RELEASE_TIMEOUT: Duration = Duration::from_secs(1); +/// Bound on the last-exit tree kill, so a process that will not die cannot hold +/// the daemon open behind it. +#[cfg(windows)] +const TREE_KILL_TIMEOUT: Duration = Duration::from_secs(10); const REGISTRY_RELEASE_POLL: Duration = Duration::from_millis(10); const STDERR_PUMP_DRAIN_TIMEOUT: Duration = Duration::from_millis(250); /// Maximum number of supervised process spawn/exit facts retained per daemon incarnation. @@ -2179,7 +2183,42 @@ impl SupervisedModule { let (reply_tx, reply_rx) = oneshot::channel(); self.inner .commands - .send(SupervisorCommand::Retire { reply: reply_tx }) + .send(SupervisorCommand::Retire { + reply: reply_tx, + tree: false, + }) + .await + .map_err(|_| SuperviseError::CommandClosed { + module_id: self.inner.module_id.clone(), + })?; + reply_rx.await.map_err(|_| SuperviseError::CommandClosed { + module_id: self.inner.module_id.clone(), + })? + } + + /// Retire the module and its process tree, inside the supervisor loop. + /// + /// The tree kill has to happen in here rather than at the caller: the drain + /// takes the child, so by the time a caller could act the pid is gone, and + /// descendants are already detached. See the `Retire` arm for the ordering. + pub(crate) async fn retire_tree(&self) -> Result<(), SuperviseError> { + match self.state()? { + ModuleState::Stopped | ModuleState::Failed => return Ok(()), + ModuleState::Starting + | ModuleState::Running + | ModuleState::Unresponsive + | ModuleState::Restarting + | ModuleState::Draining + | ModuleState::Disabled => {} + } + + let (reply_tx, reply_rx) = oneshot::channel(); + self.inner + .commands + .send(SupervisorCommand::Retire { + reply: reply_tx, + tree: true, + }) .await .map_err(|_| SuperviseError::CommandClosed { module_id: self.inner.module_id.clone(), @@ -2352,6 +2391,13 @@ enum SupervisorCommand { }, Retire { reply: oneshot::Sender>, + /// Kill the module's process tree after the drain, not just the module. + /// + /// Set by [`SupervisedModule::retire_tree`] for last-exit retirement, + /// where nothing else will come along to reap a helper the module + /// spawned. Ordinary retires leave it false and keep their exact + /// previous behaviour. + tree: bool, }, Restart { /// Operator override for this one restart's drain budget, in ms. `None` @@ -3718,8 +3764,24 @@ async fn handle_supervisor_command( } false } - SupervisorCommand::Retire { reply } => { + SupervisorCommand::Retire { reply, tree } => { let result = async { + // Capture the pid BEFORE the drain. `drain_optional_child` takes + // the child, so reading it afterwards always yields `None` and + // the tree kill would silently never run -- which is what the + // first version of this did. + #[cfg(windows)] + let tree_pid = if tree { + child.as_ref().and_then(SupervisedChild::id) + } else { + None + }; + + // Drain BEFORE the tree kill: a module hangs graceful-stop work + // off the GOODBYE the drain delivers (broca seals its WAL, + // engram closes a capture). Killing first turns that delivery + // into a no-op against a dead process and can kill a module + // mid-write. begin_forwarding_drain_if_configured( spec, runtime, @@ -3729,7 +3791,7 @@ async fn handle_supervisor_command( RouteCloseReason::Disable, ) .await?; - drain_optional_child( + let drain_result = drain_optional_child( &spec.module_id, spec.protocol, registry, @@ -3741,7 +3803,26 @@ async fn handle_supervisor_command( ModuleState::Stopped, None, ) - .await + .await; + + // `taskkill /T` reaches grandchildren the direct-child kill + // cannot, and it belongs HERE -- after the drain-and-wait, not + // ahead of it. Best-effort: a tree-kill failure must not turn a + // drained module into a failed one, so it is logged and the + // drain result still decides the outcome. + #[cfg(windows)] + if let Some(pid) = tree_pid { + if let Err(error) = kill_process_tree(&spec.module_id, pid).await { + warn!( + module_id = %spec.module_id, + pid, + %error, + "tree retirement failed; the module itself is already drained" + ); + } + } + + drain_result } .await; let registration_released = result.is_ok(); @@ -3749,7 +3830,12 @@ async fn handle_supervisor_command( if registration_released { process_liveness.untrack_if_current(&spec.module_id, snapshot); } - false + // A TREE retirement ends this module's supervision loop; an ordinary + // one does not. Returning `false` unconditionally left the loop + // running with the restart policy still armed after a last-exit + // retirement, so the daemon could respawn a module the holder monitor + // had just retired -- the orphan leak this whole change closes. + tree && registration_released } SupervisorCommand::Restart { drain_timeout_ms, @@ -4633,6 +4719,49 @@ fn spawn_child( }) } +/// Kill a process and its descendants, for last-exit retirement. +/// +/// `taskkill /T` walks the current parent-child relationships, which is why it +/// reaches a module's helper processes that `Child::kill` cannot. It is not +/// complete on its own: a grandchild whose parent already exited has been +/// reparented out of the tree and is invisible here, which is why the daemon's +/// job-object containment (`subc-jobobject`) is the primary mechanism and this +/// is the belt to its braces. +/// +/// Bounded, because `taskkill` can block on a process that will not die and the +/// daemon must not hang behind it during shutdown. +#[cfg(windows)] +async fn kill_process_tree(module_id: &str, pid: u32) -> Result<(), SuperviseError> { + let mut command = Command::new("taskkill.exe"); + command + .args(["/PID", &pid.to_string(), "/T", "/F"]) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .kill_on_drop(true) + // A daemon runs detached from any console, so a child spawned without + // this allocates a fresh console window -- a visible flash on every + // retirement. + .creation_flags(0x0800_0000); + let status = timeout(TREE_KILL_TIMEOUT, command.status()) + .await + .map_err(|_| SuperviseError::Kill { + module_id: module_id.to_string(), + source: io::Error::new(io::ErrorKind::TimedOut, "tree retirement timed out"), + })? + .map_err(|source| SuperviseError::Kill { + module_id: module_id.to_string(), + source, + })?; + if !status.success() { + return Err(SuperviseError::Kill { + module_id: module_id.to_string(), + source: io::Error::other(format!("tree retirement exited {status}")), + }); + } + Ok(()) +} + #[cfg(target_os = "linux")] fn remove_module_cgroup(placement: &subc_cgroup::Placement, module_id: &str) { match placement.remove_module(module_id) { diff --git a/docs/specs/subc-lease-contract.md b/docs/specs/subc-lease-contract.md new file mode 100644 index 00000000..5a6da9ab --- /dev/null +++ b/docs/specs/subc-lease-contract.md @@ -0,0 +1,162 @@ +# subc lease contract — OMP-owned last-exit retirement + +Status: v1 — implemented in `crates/subc-core/src/holder_monitor.rs` (PR #103), +verified by 6 unit tests + a two-holder forced-close regression. + +Owner: subc. Applies only to Windows daemons started with `OMP_SUBC_OWNED=1` +by an OMP host. Standalone daemons, launchd/systemd services, and test daemons +are unaffected: the monitor does not spawn at all (see §6). + +## 1. The problem this exists to solve + +The OMP host starts the daemon on first use and expects it to retire when the +**last** OMP host exits. The extension's own retirement path handled graceful +`session_shutdown` but a force-killed host never ran it, so the daemon, its +modules, and their grandchildren leaked indefinitely. Retirement moved to the +daemon side, where it survives a host that dies without notice. + +## 2. The files and their contract + +All paths are in the daemon's run dir (`XDG_RUNTIME_DIR` for the OMP case): + +| File | Writer | Reader | Purpose | +|---|---|---|---| +| `subc-lease-.json` | OMP host (before it starts/stops the daemon) | monitor | one per live holder; presence ⇒ a holder claims this daemon | +| `subc-retiring.lock` | monitor | monitor + every host | serializes retirement against registration | +| `subc-connection.json` | daemon | every consumer | discovery; removed on retirement so no consumer adopts a dead daemon | + +Lease content (JSON; `Owner` in `holder_monitor.rs`): + +```jsonc +{ + "pid": 1234, // the holder's pid, must be non-zero + "token": "", // opaque; empty token is never dead + "processIdentity": "2026-09-18T...Z", // Win32_Process.CreationDate, UTC ISO + "timestamp": 1780000000000 // informational; not part of liveness +} +``` + +The lock file is the same shape with a different semantic: `token` is the +**daemon's** token (`"-"`), and only the daemon +that wrote it may remove it. + +`processIdentity` is the whole reason liveness is not just "is the pid alive". +Windows reuses pids. A lease for a dead holder whose pid gets recycled by an +unrelated process would otherwise look live forever. Identity comparison +distinguishes "the recorded process is still running" from "some process with +that number exists". + +## 3. Liveness classification (fail-closed) + +`process_state(pid)` probes via PowerShell `Get-CimInstance Win32_Process`: + +- emits `GONE` (no such process) ⇒ **dead** +- emits a UTC ISO creation date ⇒ **live**, compared by `owner_gone` against + the recorded identity; mismatch ⇒ dead (pid recycled), match ⇒ live +- timeout, non-zero exit, unparseable output, missing/empty identity, pid 0, + empty token, or any probe error ⇒ **live** (fail-closed) + +The fail-closed rule is the invariant that must survive any refactor. An +unprobeable holder is a live holder; a premature retirement of a live fleet is +the failure mode this contract exists to prevent, and every error path reads +that way on purpose. + +Off-Windows the probe is `Unknown` by construction (see §6), so a Linux/macOS +OMP-owned daemon never retires — stay-up rather than retire-on-guess. + +## 4. Retirement choreography (the load-bearing ordering) + +``` +probe unlocked ── read lease fingerprints, classify each holder + ↓ all gone? +take boundary lock ── reaps a stranded lock inside take_lock + ↓ (read → probe → re-read → remove → create_new) +re-verify leases ── unchanged since the probe? else abort + ↓ +operation lock ── supervisor-level admission + ↓ +re-verify leases ── unchanged since the boundary lock? else abort + ↓ +drain + retire trees ── GOODBYE/drain FIRST, then taskkill /T /F + ↓ +remove connection file ── discovery ends; consumers stop adopting + ↓ +exit ── both locks held until bootstrap is leaving +``` + +Three properties this ordering buys: + +1. **Probes never hold the boundary lock.** Probing is one CIM query per lease + (5s deadline each). Doing that under the lock would block every host's + registration (`wx`, 250ms retry, 20s hard throw) for K×5s per tick on a + K-holder fleet — a host failure on a healthy fleet under slow WMI. +2. **Two re-verifies close the registration race.** A holder that registers + between the probe and either lock changes a lease's bytes; the fingerprint + comparison aborts retirement. The comparison is exact bytes, so a lease + rewritten with a fresh `timestamp` also aborts — conservative. +3. **Drain precedes the tree kill.** Modules hang graceful-stop work off the + GOODBYE the drain delivers (broca seals its WAL, engram closes a capture). + Killing the tree first makes that delivery a no-op against a dead process + and can kill a module mid-write. + +The monitor never kills a pid sourced from a lease. It kills only module +processes from `SupervisorHandle::list()`, via `taskkill /T /F` on the module +parent. Descendants that detach before the kill are the reason `/T` exists. + +## 5. Stranded lock recovery + +A daemon that dies mid-retirement leaves `subc-retiring.lock` behind, which +blocks every host forever. `take_lock` reaps it before retirement only: + +1. read the lock; if absent, no work +2. probe the recorded owner; if **live and identity matches**, leave it alone + (someone else is legitimately retiring) +3. re-read; if the bytes changed meanwhile, someone re-took it, leave it +4. remove, then `create_new` — the reaper is now the lock owner + +Too-fresh locks (< 30s) are left alone, and any probe ambiguity fails closed +to "leave it". Reaping happens only when the monitor is about to retire anyway, +never on a routine tick, so a live fleet never sees its lock touched. + +## 6. Platform and opt-in gating + +The module is `#[cfg(windows)]`. `spawn_if_owned` additionally requires +`OMP_SUBC_OWNED=1` and `OMP_SUBC_STOP_ON_LAST_EXIT != "0"`. Absent env ⇒ zero +behavior change: standalone daemons, launchd/systemd services, test daemons, +and every non-Windows build compile the module out and behave exactly as before. + +The extension's own stop path was **removed** in the same change. Two +retirement paths would race; the daemon monitor is now the single authority. +Hosts release their own lease on `session_shutdown` (unconditionally, last +holder only) and nothing else — pruning another host's lease is a liability, +since one host's shutdown could make the monitor retire a daemon a second live +host still needs. + +## 7. Verification + +- 6 unit tests (`holder_monitor::tests` + `loop_tests`): classification + fail-closed; and the loop choreography with an injected probe — a live + holder blocks, all-gone retires and removes the discovery file, a holder + registering after the probe aborts, an unprobeable holder stays up, empty + leases retire without deadlock. +- Two-holder regression (`subc-forced-close-check.mjs`): first-host kill + preserves the second holder and the service tree; last-host force-kill + retires daemon + module + grandchild; an abandoned registrar lock is + recovered. +- Graceful sequence (`subc-real-lifecycle.mjs`): cold start → shared → + first-exit preserve → last-exit retire → restart. + +## 8. Honest limitations (v1) + +- **Probe cost.** One PowerShell CIM query per lease per 2s tick. Under slow + WMI this is real overhead; the fix is a native `OpenProcess` + + `GetProcessTimes` probe via the `windows` crate (no new behavior, just a + faster path to the same classification). Not a correctness issue; the + boundary lock is never held across it (§4.1). +- **Fingerprint comparison is exact bytes.** A lease rewritten with a fresh + `timestamp` aborts retirement that tick. No current writer does this; if + spurious "leases changed during probe" warnings appear, a timestamp churn is + the first suspect, not a real holder. +- **No cross-machine coordination.** The lease set is per-run-dir. Two daemons + on the same machine with different run dirs are independent fleets and + neither knows about the other's holders.