From d070ce5514c7806e879b1959413e5ad4d4655e85 Mon Sep 17 00:00:00 2001 From: Qiiks Date: Mon, 21 Sep 2026 07:42:08 +0530 Subject: [PATCH 1/2] supc-daemon: job-object containment for module grandchildren (issue #109) A supervised module may spawn helpers of its own -- the Synapse embedding module spawns a CUDA worker that holds the GPU allocation -- and nothing in teardown took them. Child::kill is TerminateProcess scoped to one pid, and Windows has no process group to signal, so a module terminated rather than asked could not close its own pipes and its grandchildren outlived it. A day of restarts accumulated orphans. New leaf crate subc-jobobject follows the subc-cgroup model: this crate forbids unsafe code and the Win32 calls need it. The job carries JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE and breakaway is deliberately not permitted, so a daemon crash reaps the tree that no supervisor code is left alive to kill, and a child cannot escape by CreateProcess with CREATE_BREAKAWAY_FROM_JOB. Spawn is suspend -> assign -> resume, in contain_spawned_child: assignment must happen before the child can run a single instruction, or it could spawn a grandchild that escapes. An assignment failure is NOT fatal -- an uncontained module behaves exactly as before -- and is logged at warn. start_kill terminates the job tree after the drain wait (the maintainer's requirement): the drain and the reaped child run first, the tree kill after. Containment never changes whether a module is reported as stopped. Tests run against the supervisor: teardown reaps the grandchild, an uncontained grandchild survives a direct child kill (the defect reproduction from #109), and dropping containment reaps the tree with no teardown code running (crash durability). The regression test was verified to redden by name under a containment-disabled mutation. The job object is the Windows arm only. Unix containment is the process group plus the child roster, which is already on master; the two compose and use independent process attributes on Linux (setpgid in the child, cgroup.procs via pre_exec). Cargo.lock gains subc-jobobject, subc-daemon takes it as a target-cfg windows dependency, and subc-core's fake-aft-stub gains the grandchild mode the containment tests drive. --- Cargo.lock | 9 + Cargo.toml | 1 + crates/subc-core/src/bin/fake-aft-stub.rs | 45 +++ crates/subc-daemon/Cargo.toml | 5 + crates/subc-daemon/src/supervise.rs | 374 ++++++++++++++++++ crates/subc-jobobject/Cargo.toml | 22 ++ .../src/bin/jobobject-fixture.rs | 66 ++++ crates/subc-jobobject/src/lib.rs | 218 ++++++++++ crates/subc-jobobject/src/sys.rs | 315 +++++++++++++++ crates/subc-jobobject/tests/containment.rs | 298 ++++++++++++++ 10 files changed, 1353 insertions(+) create mode 100644 crates/subc-jobobject/Cargo.toml create mode 100644 crates/subc-jobobject/src/bin/jobobject-fixture.rs create mode 100644 crates/subc-jobobject/src/lib.rs create mode 100644 crates/subc-jobobject/src/sys.rs create mode 100644 crates/subc-jobobject/tests/containment.rs diff --git a/Cargo.lock b/Cargo.lock index 4735f02d..2fc2b1a1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1866,6 +1866,7 @@ dependencies = [ "sha2", "subc-cgroup", "subc-control", + "subc-jobobject", "subc-jsonc", "subc-protocol", "subc-transport", @@ -1886,6 +1887,14 @@ dependencies = [ "tokio", ] +[[package]] +name = "subc-jobobject" +version = "0.1.0" +dependencies = [ + "tokio", + "windows-sys 0.61.2", +] + [[package]] name = "subc-jsonc" version = "0.1.1" diff --git a/Cargo.toml b/Cargo.toml index 38225be2..6115a67b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,6 +8,7 @@ members = [ "crates/subc-jsonc", "crates/subc-cgroup", "crates/subc-uptime", + "crates/subc-jobobject", "crates/subc-core", "crates/subc-daemon", "crates/subc-daemon/tests/consumer", diff --git a/crates/subc-core/src/bin/fake-aft-stub.rs b/crates/subc-core/src/bin/fake-aft-stub.rs index 9721235a..c9d6a058 100644 --- a/crates/subc-core/src/bin/fake-aft-stub.rs +++ b/crates/subc-core/src/bin/fake-aft-stub.rs @@ -136,6 +136,20 @@ const FAKE_AFT_ORPHAN_WRITER_MODE_ENV: &str = "FAKE_AFT_ORPHAN_WRITER_MODE"; /// cannot be asked what it did with a signal, and what it does with a signal is /// the other half of what these tests need to observe. const FAKE_AFT_NEVER_CONNECT_ENV: &str = "FAKE_AFT_NEVER_CONNECT"; +/// Presence marks a GRANDCHILD: park forever, touch nothing else. +/// +/// The child→grandchild shape a teardown test needs, because a supervision test +/// that only ever observes the direct child cannot tell a reaped tree from a +/// leaked helper. Checked before every other arm, since the grandchild must not +/// dial subc, exit on its own, or register a signal handler. +const FAKE_AFT_GRANDCHILD_MODE_ENV: &str = "FAKE_AFT_GRANDCHILD_MODE"; +/// Where a parent stub records the pid of the grandchild it spawned. +/// +/// The PARENT writes this, from the `Child` handle it already holds, rather than +/// the grandchild reporting itself: a self-report would need the grandchild to +/// reach a point where it can write, which is a race against the teardown the +/// test is about to perform. +const FAKE_AFT_GRANDCHILD_PID_FILE_ENV: &str = "FAKE_AFT_GRANDCHILD_PID_FILE"; /// Where to write a marker file when SIGTERM arrives, just before exiting 0. /// /// The marker is the witness that the supervisor ASKED before it forced: the @@ -242,6 +256,12 @@ async fn main() -> Result<(), StubError> { // never dialled subc either. Checking first also means a connect/HELLO // failure can never land its own noise in the very stderr ring this knob // is configured to control. + if env_flag(FAKE_AFT_GRANDCHILD_MODE_ENV) { + // A grandchild does nothing but exist: no subc dial, no exit, no signal + // handler. Its whole purpose is to be a process the teardown must reach. + std::future::pending::<()>().await; + unreachable!("a pending future never resolves"); + } if env_flag(FAKE_AFT_ORPHAN_WRITER_MODE_ENV) { return run_detached_orphan_writer().await; } @@ -297,6 +317,7 @@ async fn run_never_connect() -> Result<(), StubError> { } announce_never_connect_ready()?; + spawn_grandchild_if_requested()?; // Park. The supervisor's teardown -- signal, or the kill behind it -- is what // ends this process; nothing here decides to stop on its own, because a test @@ -306,6 +327,30 @@ async fn run_never_connect() -> Result<(), StubError> { unreachable!("a pending future never resolves"); } +/// Spawn a grandchild and record its pid, when the run asks for one. +/// +/// The pid is written from HERE, out of the `Child` handle, so the test can +/// address the grandchild without racing its startup. The `Child` is +/// deliberately dropped rather than kept: `std::process::Child` closes its handle +/// on drop and does not wait, so the grandchild keeps running while the parent +/// holds nothing that would keep the process object alive past its death. +fn spawn_grandchild_if_requested() -> Result<(), StubError> { + let Ok(pid_file) = env::var(FAKE_AFT_GRANDCHILD_PID_FILE_ENV) else { + return Ok(()); + }; + let exe = env::current_exe().map_err(StubError::Io)?; + let grandchild = Command::new(exe) + .env(FAKE_AFT_GRANDCHILD_MODE_ENV, "1") + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .map_err(StubError::Io)?; + fs::write(&pid_file, grandchild.id().to_string()).map_err(StubError::Io)?; + drop(grandchild); + Ok(()) +} + fn announce_never_connect_ready() -> Result<(), StubError> { let Ok(path) = env::var(FAKE_AFT_NEVER_CONNECT_READY_PATH_ENV) else { return Ok(()); diff --git a/crates/subc-daemon/Cargo.toml b/crates/subc-daemon/Cargo.toml index 1da33cab..18ee219b 100644 --- a/crates/subc-daemon/Cargo.toml +++ b/crates/subc-daemon/Cargo.toml @@ -33,6 +33,11 @@ tracing = "0.1" [target.'cfg(target_os = "linux")'.dependencies] subc-cgroup = { path = "../subc-cgroup", version = "0.1.1" } +[target.'cfg(windows)'.dependencies] +# Job-object containment for module grandchildren (issue #109). A leaf crate +# because this crate forbids unsafe code and the Win32 calls need it. +subc-jobobject = { path = "../subc-jobobject" } + [target.'cfg(unix)'.dependencies] # SIGTERM for `protocol: "none"` teardown. This crate forbids unsafe code, so # `libc::kill` is not reachable from here; rustix wraps the same syscall safely diff --git a/crates/subc-daemon/src/supervise.rs b/crates/subc-daemon/src/supervise.rs index 20750b7c..26a90697 100644 --- a/crates/subc-daemon/src/supervise.rs +++ b/crates/subc-daemon/src/supervise.rs @@ -92,6 +92,13 @@ struct SupervisedChild { module_id: String, #[cfg(target_os = "linux")] cgroup_placement: Option, + /// The job that contains this child and every process it spawns (issue #109). + /// + /// Dropping this handle is what reaps a surviving tree when no supervisor + /// code runs — a daemon crash — because the job carries + /// `JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE`. + #[cfg(windows)] + job: Option, stdout_pump: Option>, stderr_pump: Option>, stderr_ring: Arc>, @@ -131,7 +138,28 @@ impl SupervisedChild { result } + /// Kill the child, and on Windows the whole tree it spawned (issue #109). + /// + /// `Child::kill` is `TerminateProcess` scoped to one pid, so a module with a + /// helper process leaked the helper — the Synapse embedding module's CUDA + /// worker holds the GPU allocation, so the leak cost VRAM until the next + /// restart of something else. Terminating the job reaches grandchildren that + /// a tree walk cannot, including one whose parent has already exited and + /// been reparented away. + /// + /// Best-effort like `request_graceful_stop`: a job failure is logged and the + /// direct-child kill still decides the outcome, so containment can never + /// change whether a module is reported as stopped. fn start_kill(&mut self) -> io::Result<()> { + #[cfg(windows)] + if let Some(job) = &self.job { + if let Err(error) = job.terminate() { + debug!( + error = %error, + "job termination failed; the direct-child kill still owns the outcome" + ); + } + } self.child.start_kill() } @@ -5424,6 +5452,13 @@ fn spawn_child_in_slot( #[cfg(unix)] command.process_group(0); command.stdin(Stdio::null()); + + // Containment, step 1 of 3 (issue #109): create the child suspended so it + // cannot run a single instruction -- and therefore cannot spawn a + // grandchild -- before it is in the job. See `contain_spawned_child` for the + // other two steps and why the window matters. + #[cfg(windows)] + subc_jobobject::suspend_on_create_async(&mut command); let mut child = match command.spawn() { Ok(child) => child, Err(source) => { @@ -5438,6 +5473,10 @@ fn spawn_child_in_slot( }); } }; + + // Containment, steps 2 and 3: assign while suspended, then resume. + #[cfg(windows)] + let job = contain_spawned_child(&child, spec)?; let spawned_at_ms = unix_ms_now(); let spawned_from = spec.program.clone(); let spawned_file_identity = spawned_file_identity(&spawned_from); @@ -5497,6 +5536,8 @@ fn spawn_child_in_slot( module_id: cgroup_name, #[cfg(target_os = "linux")] cgroup_placement: cgroup_placement.cloned(), + #[cfg(windows)] + job, stdout_pump, stderr_pump, stderr_ring: Arc::clone(ring), @@ -5510,6 +5551,91 @@ fn spawn_child_in_slot( }) } +/// Contain a freshly spawned Windows child and start it. +/// +/// Steps 2 and 3 of the suspended-create contract: the job is created and the +/// child assigned **while it is still suspended** (step 1 is +/// `suspend_on_create_async` at the spawn site), then the child is resumed. +/// +/// A child that is never resumed hangs forever holding a pid, so a resume +/// failure kills the child and fails the spawn rather than returning a +/// `SupervisedChild` that can never run. +/// +/// An assignment failure is NOT fatal: an uncontained module behaves exactly as +/// it did before this existed, whereas refusing to start one would be a new +/// outage. It is logged at warn because it means a helper process could leak. +#[cfg(windows)] +fn contain_spawned_child( + child: &Child, + spec: &ModuleSpec, +) -> Result, SuperviseError> { + let module_id = spec.module_id.as_str(); + let Some(pid) = child.id() else { + // The child exited between spawn and here. Its tree, if it made one, + // needs no containment: nothing is left to contain. + warn!( + module_id, + "spawned child had already exited before containment; no job object attached" + ); + return Ok(None); + }; + + let job = match subc_jobobject::JobObject::new() { + Ok(job) => job, + Err(source) => { + warn!( + module_id, + error = %source, + "could not create a job object; this module's helper processes will not be \ + reaped on teardown" + ); + // Resume regardless: leaving the child suspended would turn a + // containment gap into a hung module. + resume_suspended_child(pid, spec)?; + return Ok(None); + } + }; + + if let Err(source) = job.assign(child) { + warn!( + module_id, + error = %source, + "could not assign the child to its job object; this module's helper processes \ + will not be reaped on teardown" + ); + resume_suspended_child(pid, spec)?; + return Ok(None); + } + + resume_suspended_child(pid, spec)?; + Ok(Some(job)) +} + +/// Resume a suspended child, killing it if it cannot be started. +/// +/// A suspended process holds a pid and does nothing, so there is no useful +/// state to return: the caller gets an error and the spawn fails. +#[cfg(windows)] +fn resume_suspended_child(pid: u32, spec: &ModuleSpec) -> Result<(), SuperviseError> { + if let Err(source) = subc_jobobject::resume_main_thread(pid) { + // Kill it here rather than leaving a suspended process for the caller + // to notice; `kill_on_drop` would eventually do this, but the module + // would have been reported as running in between. + let _ = std::process::Command::new("taskkill.exe") + .args(["/PID", &pid.to_string(), "/T", "/F"]) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status(); + return Err(SuperviseError::Spawn { + program: spec.program.clone(), + source, + cgroup_path: None, + }); + } + Ok(()) +} + #[cfg(target_os = "linux")] fn remove_module_cgroup(placement: &subc_cgroup::Placement, module_id: &str) { match placement.remove_module(module_id) { @@ -8768,3 +8894,251 @@ mod cgroup_placement_tests { ); } } + +/// Containment of a module's process tree (issue #109). +/// +/// The behaviour these defend against is a module helper surviving its module: +/// on a real machine the Synapse embedding module's CUDA worker holds ~2.2 GB of +/// VRAM, so a leaked grandchild is a leaked GPU allocation, and a day of restarts +/// compounds it. +/// +/// They run against the SUPERVISOR rather than the job-object crate because the +/// claim is about teardown: a crate-level test proves a job can reap a tree, not +/// that the daemon's drain path reaches it. +/// +/// Windows-only, like the mechanism. On Unix this arm compiles out; the cgroup +/// lane there is a separate containment path with its own tests. +#[cfg(all(test, windows))] +mod job_containment_tests { + use super::*; + use crate::test_support::TestTempDir; + use std::{ + path::{Path, PathBuf}, + sync::{Arc, Mutex}, + time::{Duration, Instant}, + }; + + /// The stub, expected beside this test executable. + /// + /// The existence check is here for the reason its twin at `fake_aft_stub_path` + /// documents: `--lib` does not build `[[bin]]` targets, and a bare spawn + /// failure then reads as a broken test rather than an unbuilt dependency. + fn stub_path() -> PathBuf { + let mut path = std::env::current_exe().expect("current_exe available in tests"); + path.pop(); + path.pop(); + path.push("fake-aft-stub.exe"); + assert!( + path.exists(), + "fake-aft-stub not built at {}: run `cargo test -p subc-core` (which builds \ + [[bin]] targets) rather than `cargo test -p subc-core --lib` (which does not)", + path.display() + ); + path + } + + /// Poll for the grandchild pid the stub records, and parse it. + fn read_grandchild_pid(path: &Path) -> u32 { + let deadline = Instant::now() + Duration::from_secs(10); + loop { + if let Ok(contents) = std::fs::read_to_string(path) { + if let Ok(pid) = contents.trim().parse() { + return pid; + } + } + assert!( + Instant::now() < deadline, + "the stub never recorded a grandchild pid at {}", + path.display() + ); + std::thread::sleep(Duration::from_millis(10)); + } + } + + /// Everything one fixture run needs, so the two tests below differ in exactly + /// one place: whether the child is contained. + struct Fixture { + _dir: TestTempDir, + module_id: String, + grandchild: u32, + child: Option, + registry: Arc, + snapshot: Arc>, + terminal_ring: Arc>, + spawn_events: SpawnEventFeed, + } + + fn fixture(label: &str, module_id: &str) -> Fixture { + let dir = TestTempDir::new(label); + let pid_file = dir.join("grandchild.pid"); + let supervisor = Supervisor::new( + Arc::new(Registry::default()), + RestartPolicy::new(3, Duration::ZERO), + ); + let runtime = supervisor.runtime_config(); + let snapshot = Arc::new(Mutex::new(SupervisorSnapshot::starting())); + let spec = ModuleSpec { + module_id: module_id.to_string(), + program: stub_path(), + // Zero args deliberately: a `--subc` argument would make the stub dial + // a daemon that is not there, and the failure would land in the same + // stderr ring this fixture exists to keep quiet. + args: Vec::new(), + env: vec![ + ("FAKE_AFT_NEVER_CONNECT".to_string(), "1".to_string()), + ( + "FAKE_AFT_GRANDCHILD_PID_FILE".to_string(), + pid_file.display().to_string(), + ), + ], + reserved: false, + reserved_prefixes: Vec::new(), + protocol: ModuleProtocol::Subc, + overlap: Default::default(), + }; + let child = spawn_and_mark_running(&spec, &runtime, &snapshot) + .expect("spawn the supervised fixture"); + let grandchild = read_grandchild_pid(&pid_file); + Fixture { + _dir: dir, + module_id: module_id.to_string(), + grandchild, + child: Some(child), + registry: Arc::new(Registry::default()), + snapshot, + terminal_ring: Arc::clone(&runtime.terminal_ring), + spawn_events: SpawnEventFeed::default(), + } + } + + impl Fixture { + /// Drain through the supervisor's own teardown path. + async fn drain(&mut self) { + let child = self + .child + .take() + .expect("the fixture child is still present"); + drain_child_to_state( + &self.module_id, + ModuleProtocol::Subc, + &self.registry, + &self.snapshot, + &self.terminal_ring, + &self.spawn_events, + child, + Duration::from_millis(500), + ModuleState::Stopped, + Some(false), + ) + .await + .expect("drain the supervised fixture"); + } + } + + /// Teardown reaps the grandchild, not merely the direct child. + /// + /// This is the assertion the change exists for. Before containment the + /// grandchild survived: it is a separate process, and `start_kill` is + /// `TerminateProcess` scoped to one pid. + #[tokio::test] + async fn teardown_reaps_the_grandchild() { + let mut fixture = fixture("teardown-grandchild", "tree-teardown"); + let grandchild = fixture.grandchild; + + assert!( + subc_jobobject::process_exists(grandchild), + "grandchild {grandchild} must be alive before teardown, or this proves nothing" + ); + + fixture.drain().await; + + assert!( + subc_jobobject::wait_for_process_exit(grandchild, Duration::from_secs(10)), + "grandchild {grandchild} outlived module teardown: the tree was not contained" + ); + } + + /// The mutation control: with containment withheld, the grandchild survives + /// the same kill. + /// + /// This is the defect reproduction from #109 — a direct-child kill reaches + /// one pid, and the grandchild is a different process. It spawns OUTSIDE the + /// supervisor because `spawn_and_mark_running` now always contains on + /// Windows, which is the point: there is no longer a path that spawns + /// uncontained, so the control has to construct one. + /// + /// Its job is to keep `teardown_reaps_the_grandchild` honest. If the + /// grandchild ever dies here, that test is passing for a reason unrelated to + /// the job object and the containment claim is unproven. + #[test] + fn an_uncontained_grandchild_survives_a_direct_child_kill() { + let dir = TestTempDir::new("teardown-uncontained"); + let pid_file = dir.join("grandchild.pid"); + let mut child = std::process::Command::new(stub_path()) + .env("FAKE_AFT_NEVER_CONNECT", "1") + .env( + "FAKE_AFT_GRANDCHILD_PID_FILE", + pid_file.display().to_string(), + ) + .stdin(std::process::Stdio::null()) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .spawn() + .expect("spawn the uncontained fixture"); + let grandchild = read_grandchild_pid(&pid_file); + + // Exactly what the pre-fix teardown did: kill the direct child. + child.kill().expect("kill the direct child"); + let _ = child.wait(); + + assert!( + subc_jobobject::process_exists(grandchild), + "grandchild {grandchild} died with the direct child, so this control no longer \ + distinguishes contained from uncontained teardown and the regression test is \ + passing vacuously" + ); + + // The orphan this control demonstrates is the leak the fix prevents, so + // the control must not leave one behind. + kill_tree(grandchild); + } + + /// Crash durability: closing the containment handle reaps the tree with no + /// teardown code running at all. + /// + /// This is the case `taskkill /T` cannot cover — a daemon that dies cannot + /// call anything — and it is why containment is a kernel property of the + /// handle rather than a step in the drain. Discovered by getting the + /// mutation control wrong: clearing `job` to "disable" containment instead + /// killed the tree, which is the guarantee, not a mistake. + #[tokio::test] + async fn dropping_containment_reaps_the_grandchild() { + let mut fixture = fixture("drop-containment", "tree-drop"); + let grandchild = fixture.grandchild; + + assert!(subc_jobobject::process_exists(grandchild)); + + // No `drain` call, no kill: dropping the handle is the entire mechanism. + fixture.child.as_mut().expect("child present").job = None; + + assert!( + subc_jobobject::wait_for_process_exit(grandchild, Duration::from_secs(10)), + "grandchild {grandchild} survived the containment handle closing, so a daemon \ + crash would leave the tree behind" + ); + } + + /// Kill a pid and its tree, then confirm it is gone. + fn kill_tree(pid: u32) { + let _ = std::process::Command::new("taskkill.exe") + .args(["/PID", &pid.to_string(), "/T", "/F"]) + .stdin(std::process::Stdio::null()) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .status(); + assert!( + subc_jobobject::wait_for_process_exit(pid, Duration::from_secs(10)), + "could not clean up grandchild {pid}" + ); + } +} diff --git a/crates/subc-jobobject/Cargo.toml b/crates/subc-jobobject/Cargo.toml new file mode 100644 index 00000000..aeb107d0 --- /dev/null +++ b/crates/subc-jobobject/Cargo.toml @@ -0,0 +1,22 @@ +[package] +name = "subc-jobobject" +version = "0.1.0" +edition = "2021" +license = "MIT" +publish = false +description = "Windows job-object containment for supervised subc module children." + +[dependencies] +tokio = { version = "1", features = ["process"] } + +[target.'cfg(windows)'.dependencies] +windows-sys = { version = "0.61", features = [ + "Win32_Foundation", + "Win32_Security", + "Win32_System_Diagnostics_ToolHelp", + "Win32_System_JobObjects", + "Win32_System_Threading", +] } + +[lints.rust] +unsafe_code = "deny" diff --git a/crates/subc-jobobject/src/bin/jobobject-fixture.rs b/crates/subc-jobobject/src/bin/jobobject-fixture.rs new file mode 100644 index 00000000..73d66be5 --- /dev/null +++ b/crates/subc-jobobject/src/bin/jobobject-fixture.rs @@ -0,0 +1,66 @@ +//! Test fixture: a process that spawns a grandchild and then parks. +//! +//! Both modes exist so containment can be measured against a real descendant +//! tree rather than a single pid. The grandchild is a genuine re-exec of this +//! binary, so it is a separate process with its own pid that the job must +//! contain by membership — a thread or a task would prove nothing. +//! +//! Run by the integration tests in `tests/containment.rs`, never by hand. + +use std::{env, fs, process::Command, thread, time::Duration}; + +/// Which role this invocation plays. +const MODE_ENV: &str = "SUBC_JOBOBJECT_FIXTURE_MODE"; +/// Where the grandchild writes its own pid, so the test can address it. +const GRANDCHILD_PID_ENV: &str = "SUBC_JOBOBJECT_GRANDCHILD_PID_FILE"; + +const PARENT: &str = "parent"; +const GRANDCHILD: &str = "grandchild"; + +fn main() { + match env::var(MODE_ENV).as_deref() { + Ok(PARENT) => run_parent(), + Ok(GRANDCHILD) => run_grandchild(), + _ => { + eprintln!("fixture: set {MODE_ENV} to '{PARENT}' or '{GRANDCHILD}'"); + std::process::exit(2); + } + } +} + +/// Spawn the grandchild, then park. +/// +/// Parking rather than exiting is what makes this a teardown test: a parent +/// that exited on its own would leave the grandchild reparented for reasons +/// unrelated to containment, and the two cases would be indistinguishable. +fn run_parent() { + let exe = env::current_exe().expect("current_exe"); + // The grandchild inherits this environment, so it sees the same pid-file + // path and publishes its own pid there. No rewriting in between. + let grandchild = Command::new(exe) + .env(MODE_ENV, GRANDCHILD) + .spawn() + .expect("spawn grandchild"); + // The handle is deliberately forgotten rather than dropped-and-waited: the + // parent parks below, the grandchild must outlive it independently, and on + // Windows dropping a Child neither kills it nor reaps it. + std::mem::forget(grandchild); + park_forever(); +} + +/// Publish this process's pid and park. +fn run_grandchild() { + let path = env::var(GRANDCHILD_PID_ENV).expect("grandchild pid file"); + fs::write(&path, std::process::id().to_string()).expect("write grandchild pid"); + park_forever(); +} + +/// Sleep until something kills this process. +/// +/// The supervisor's teardown is the only thing that ends these processes, which +/// is the property under test. +fn park_forever() -> ! { + loop { + thread::sleep(Duration::from_secs(3600)); + } +} diff --git a/crates/subc-jobobject/src/lib.rs b/crates/subc-jobobject/src/lib.rs new file mode 100644 index 00000000..4c16b393 --- /dev/null +++ b/crates/subc-jobobject/src/lib.rs @@ -0,0 +1,218 @@ +//! Windows job-object containment for supervised subc module children. +//! +//! A supervised module may spawn helpers of its own — the Synapse embedding +//! module spawns a CUDA worker that holds the GPU allocation — and until this +//! crate existed nothing in teardown took them. `Child::kill` is `TerminateProcess` +//! scoped to one pid, and Windows has no process group to signal, so a module +//! that was terminated rather than asked could not close its own pipes and its +//! grandchildren outlived it. A day of restarts accumulated orphans. +//! +//! A job object contains by **membership** rather than ancestry, which covers the +//! two cases a `taskkill /T` walk cannot: +//! +//! * a grandchild spawned between the walk and the kill; +//! * a grandchild whose parent has already exited and been reparented, so it is +//! no longer in the tree to be found. Those are precisely the orphans that +//! accumulate across restarts. +//! +//! # The suspended-create contract +//! +//! [`suspend_on_create`] and [`resume_main_thread`] exist because assignment must +//! happen before the child can run a single instruction. `CreateProcess` gives the +//! caller no way to place a process in a job atomically short of building the +//! process by hand with `PROC_THREAD_ATTRIBUTE_JOB_LIST`, so the child is created +//! suspended, assigned, and then resumed. A window between spawn and assignment +//! is a grandchild that escapes — the defect this crate fixes, in smaller form. +//! +//! The order is load-bearing and there are exactly three steps: +//! +//! ```no_run +//! # #[cfg(windows)] +//! # fn main() -> std::io::Result<()> { +//! let job = subc_jobobject::JobObject::new()?; +//! let mut command = std::process::Command::new("module.exe"); +//! subc_jobobject::suspend_on_create(&mut command); +//! let child = command.spawn()?; +//! job.assign(&child)?; +//! subc_jobobject::resume_main_thread(child.id())?; +//! # Ok(()) +//! # } +//! # #[cfg(not(windows))] fn main() {} +//! ``` +//! +//! [`spawn_contained`] performs all three for the synchronous case. A caller that +//! spawns asynchronously sets the flag, spawns, calls [`JobObject::assign`], and +//! resumes — which is what the daemon's `contain_spawned_child` does. +//! +//! If the caller never resumes, the child stays suspended forever, so +//! [`resume_main_thread`] reports the failure rather than logging it: the caller +//! must kill the child and fail the spawn. + +#![cfg(windows)] +#![deny(unsafe_code)] + +mod sys; + +use std::{ + io, + os::windows::io::AsRawHandle, + process::Child, + thread, + time::{Duration, Instant}, +}; + +pub use sys::{JobObject, CREATE_SUSPENDED_FLAG}; + +/// A process handle suitable for [`JobObject::assign`]. +/// +/// `std` and `tokio` children expose their handle differently — `AsRawHandle` +/// for the former, an inherent `raw_handle()` that returns `Option` for the +/// latter — so this normalizes both rather than making the caller reach for the +/// right accessor and get the `None` case wrong. +pub trait ProcessHandle { + /// The process handle, or `None` if the process has already been reaped. + fn handle(&self) -> Option<*mut std::ffi::c_void>; +} + +/// A raw handle wrapper, so `JobObject::assign` can take an already-resolved +/// handle without dereferencing one that a safe caller cannot validate. +#[derive(Clone, Copy)] +pub struct RawProcessHandle(pub *mut std::ffi::c_void); + +impl ProcessHandle for RawProcessHandle { + fn handle(&self) -> Option<*mut std::ffi::c_void> { + Some(self.0) + } +} + +impl ProcessHandle for Child { + fn handle(&self) -> Option<*mut std::ffi::c_void> { + Some(self.as_raw_handle().cast()) + } +} + +impl ProcessHandle for tokio::process::Child { + fn handle(&self) -> Option<*mut std::ffi::c_void> { + self.raw_handle().map(|handle| handle.cast()) + } +} + +/// [`ProcessHandle::handle`] by reference, for the call site's readability. +pub fn process_handle(child: &C) -> Option<*mut std::ffi::c_void> { + child.handle() +} + +/// How long [`resume_main_thread`] keeps looking for the new process's thread. +/// +/// A snapshot taken immediately after `CreateProcess` can miss the thread even +/// though the process exists, so this is a retry budget rather than an expected +/// wait: the ordinary case finds the thread on the first attempt. +const RESUME_DISCOVERY_TIMEOUT: Duration = Duration::from_millis(500); + +/// Mark `command` so the child is created suspended. +/// +/// The child must not run a single instruction before [`JobObject::assign`], or +/// it could spawn a grandchild that escapes containment. Call +/// [`resume_main_thread`] once the child is assigned. +pub fn suspend_on_create(command: &mut std::process::Command) { + sys::set_suspended_creation_flags(command); +} + +/// [`suspend_on_create`] for the async command the supervisor spawns. +pub fn suspend_on_create_async(command: &mut tokio::process::Command) { + sys::set_suspended_creation_flags_async(command); +} + +/// Start the child created by [`suspend_on_create`], now that it is contained. +/// +/// `std` keeps no handle to the primary thread that `CreateProcess` returns, so +/// it is found through a toolhelp snapshot. A failure here means the child cannot +/// be started and the caller must kill it rather than leave a suspended process +/// holding a pid. +pub fn resume_main_thread(pid: u32) -> io::Result<()> { + let started = Instant::now(); + let mut backoff = Duration::from_millis(1); + loop { + match sys::open_first_thread(pid)? { + Some(thread) => return sys::resume_and_close(thread), + None if started.elapsed() < RESUME_DISCOVERY_TIMEOUT => { + thread::sleep(backoff); + backoff = (backoff * 2).min(Duration::from_millis(20)); + } + None => { + return Err(io::Error::new( + io::ErrorKind::NotFound, + format!("no thread appeared for pid {pid} within {RESUME_DISCOVERY_TIMEOUT:?}"), + )); + } + } + } +} + +/// Wait for `pid` to leave the process table. +/// +/// Teardown is asynchronous: `TerminateJobObject` returns once the members are +/// signalled, not once they have been reaped, so a caller that asserted +/// immediately would be racing the kernel. `true` means gone, `false` means still +/// present at the deadline. +pub fn wait_for_process_exit(pid: u32, timeout: Duration) -> bool { + let deadline = Instant::now() + timeout; + loop { + if !process_exists(pid) { + return true; + } + if Instant::now() >= deadline { + return false; + } + thread::sleep(Duration::from_millis(10)); + } +} + +/// Whether a pid currently names a running process. +pub fn process_exists(pid: u32) -> bool { + sys::process_exists(pid) +} + +/// A child and the job that contains it, held together. +/// +/// The daemon spawns, assigns, and resumes in sequence; keeping the job beside +/// the child is what makes the association survive a refactor that reorders those +/// steps. +#[derive(Debug)] +pub struct ContainedChild { + /// The supervised process. + pub child: C, + /// The job that contains it and everything it spawns. + pub job: JobObject, +} + +impl ContainedChild { + /// Contain `child` and start it. + /// + /// The one call that gets the order right: assign while suspended, then + /// resume. A failure to resume kills the child, because a suspended process + /// holding a pid with no way to start is a leak rather than a failed spawn. + pub fn contain(child: C, job: JobObject, pid: u32) -> io::Result { + let handle = child + .handle() + .ok_or_else(|| io::Error::other("child was reaped before it could be contained"))?; + job.assign(&RawProcessHandle(handle))?; + if let Err(error) = resume_main_thread(pid) { + let _ = job.terminate(); + return Err(error); + } + Ok(Self { child, job }) + } +} + +/// Spawn a command already contained, with no window for an escape. +/// +/// Applies [`suspend_on_create`], spawns, assigns, and resumes in one place so +/// the ordering cannot be got wrong at a call site. +pub fn spawn_contained(command: &mut std::process::Command) -> io::Result> { + suspend_on_create(command); + let child = command.spawn()?; + let pid = child.id(); + let job = JobObject::new()?; + ContainedChild::contain(child, job, pid) +} diff --git a/crates/subc-jobobject/src/sys.rs b/crates/subc-jobobject/src/sys.rs new file mode 100644 index 00000000..5c14f55a --- /dev/null +++ b/crates/subc-jobobject/src/sys.rs @@ -0,0 +1,315 @@ +//! Every Win32 call in this crate, plus the handle-owning type, and the safety +//! argument for each. +//! +//! `unsafe_code` is denied at the crate root and allowed here by exception, which +//! keeps the FFI boundary auditable in one file rather than spread across the +//! public API. [`JobObject`] lives here rather than in `lib.rs` because it owns a +//! raw handle and therefore needs a manual `Send`: an `unsafe impl` belongs on +//! the side of the boundary that is allowed to write one. + +#![allow(unsafe_code)] + +use std::{io, mem::size_of, os::windows::process::CommandExt, process::Command}; + +use windows_sys::Win32::{ + Foundation::{CloseHandle, GetLastError, HANDLE, INVALID_HANDLE_VALUE, WAIT_TIMEOUT}, + System::{ + Diagnostics::ToolHelp::{ + CreateToolhelp32Snapshot, Thread32First, Thread32Next, TH32CS_SNAPTHREAD, THREADENTRY32, + }, + JobObjects::{ + AssignProcessToJobObject, CreateJobObjectW, JobObjectBasicProcessIdList, + JobObjectExtendedLimitInformation, QueryInformationJobObject, SetInformationJobObject, + TerminateJobObject, JOBOBJECT_EXTENDED_LIMIT_INFORMATION, + JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE, + }, + Threading::{ + OpenProcess, OpenThread, ResumeThread, WaitForSingleObject, CREATE_SUSPENDED, + PROCESS_SYNCHRONIZE, THREAD_SUSPEND_RESUME, + }, + }, +}; + +/// Preceded by `CREATE_SUSPENDED` so a child cannot run before it is contained. +pub const CREATE_SUSPENDED_FLAG: u32 = CREATE_SUSPENDED; + +/// `ERROR_NO_MORE_FILES`. The toolhelp thread walk reports ordinary exhaustion +/// through this code, so it is the one value that separates "finished +/// enumerating" from a real failure. +const ERROR_NO_MORE_FILES: u32 = 18; + +/// `ERROR_MORE_DATA`. A job holding more processes than the query's buffer can +/// enumerate reports this while still filling in the assigned count. +const ERROR_MORE_DATA: u32 = 234; + +/// How many pids one [`job_process_count`] query can enumerate. +/// +/// The count of *assigned* processes is reported by the kernel regardless of what +/// the buffer can hold; only the enumerated subset is bounded. Far past any +/// legitimate module tree, and small enough to keep the query a stack +/// allocation. +const PROCESS_ID_LIST_WIDTH: usize = 256; + +/// A job object that owns one supervised module's whole process tree. +/// +/// `JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE` is set at creation, so dropping the last +/// handle kills every member. That is what makes containment survive a daemon +/// crash, where no code runs to call [`JobObject::terminate`] — the handle closes +/// with the process and the kernel reaps the tree. +/// +/// Breakaway is deliberately NOT permitted: leaving `JOB_OBJECT_LIMIT_BREAKAWAY_OK` +/// unset is what stops a child from calling `CreateProcess` with +/// `CREATE_BREAKAWAY_FROM_JOB` and escaping. A module that could leave the job +/// could leave its grandchildren behind, which is the whole defect. +#[derive(Debug)] +pub struct JobObject { + handle: HANDLE, +} + +// SAFETY: a `HANDLE` is a process-wide kernel handle value, not a pointer into +// this process's address space. Every operation this crate performs on it +// (`AssignProcessToJobObject`, `TerminateJobObject`, `QueryInformationJobObject`, +// `CloseHandle`) is thread-agnostic, and the daemon holds the job across an await +// in a spawned task, which is why `Send` is required at all. `Sync` is NOT +// implemented: concurrent `TerminateJobObject` and `Drop` on the same handle +// would be a double-close, and nothing needs it. +unsafe impl Send for JobObject {} + +impl JobObject { + /// Create a job object whose members die with it. + pub fn new() -> io::Result { + // A null `SECURITY_ATTRIBUTES` gives the default descriptor and a null + // name means unnamed, so no other process can open it to interfere. + let handle = unsafe { CreateJobObjectW(std::ptr::null(), std::ptr::null()) }; + if handle.is_null() { + return Err(io::Error::last_os_error()); + } + + // Breakaway is left unset deliberately, per the type documentation. + let mut limits = JOBOBJECT_EXTENDED_LIMIT_INFORMATION::default(); + limits.BasicLimitInformation.LimitFlags = JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE; + let applied = unsafe { + SetInformationJobObject( + handle, + JobObjectExtendedLimitInformation, + std::ptr::from_ref(&limits).cast(), + size_of::() as u32, + ) + }; + if applied == 0 { + let error = io::Error::last_os_error(); + // The handle is already owned by us, so a failure here must close it + // or the job outlives the error that rejected it. + unsafe { CloseHandle(handle) }; + return Err(error); + } + + Ok(Self { handle }) + } + + /// Put a process in the job, and with it every process it goes on to create. + /// + /// Generic over [`crate::ProcessHandle`] rather than taking a raw `HANDLE`: + /// a raw handle cannot be validated by a safe caller, and the trait keeps the + /// `None` (already reaped) case expressed in the type. Must be called while + /// the process is still suspended: see the crate root. + pub fn assign(&self, child: &C) -> io::Result<()> { + let Some(process) = child.handle() else { + return Err(io::Error::other( + "child was reaped before it could be assigned to the job", + )); + }; + let assigned = unsafe { AssignProcessToJobObject(self.handle, process) }; + if assigned == 0 { + return Err(io::Error::last_os_error()); + } + Ok(()) + } + + /// Kill every member of the job. + /// + /// The explicit counterpart to the close-on-drop guarantee, for a teardown + /// that wants the tree gone now rather than when the handle is released. + /// Membership is what lets this reach grandchildren whose parent has already + /// exited and been reparented away from the tree. + pub fn terminate(&self) -> io::Result<()> { + let terminated = unsafe { TerminateJobObject(self.handle, 1) }; + if terminated == 0 { + return Err(io::Error::last_os_error()); + } + Ok(()) + } + + /// How many processes the job currently contains. + /// + /// Exists so a test can witness the grandchild rather than infer it: a + /// teardown that reaped the direct child and nothing else reports `1` here, + /// which is a different observation from the child being gone. + pub fn process_count(&self) -> io::Result { + let mut list = PidList::default(); + let queried = unsafe { + QueryInformationJobObject( + self.handle, + JobObjectBasicProcessIdList, + std::ptr::from_mut(&mut list).cast(), + size_of::() as u32, + std::ptr::null_mut(), + ) + }; + if queried == 0 { + // A job with more members than the buffer enumerates returns this + // while still filling in the true assigned count, so the count is + // readable rather than unavailable. + let code = unsafe { GetLastError() }; + if code != ERROR_MORE_DATA { + return Err(io::Error::from_raw_os_error(code as i32)); + } + } + Ok(list.number_of_assigned_processes) + } +} + +impl Drop for JobObject { + fn drop(&mut self) { + // Close-on-drop is the crash-durability guarantee, not merely cleanup: if + // the daemon is killed, this handle closes with it and the kernel reaps + // the tree that no supervisor code is left alive to kill. + unsafe { CloseHandle(self.handle) }; + } +} + +/// The kernel's variable-length pid list, sized for a single query. +/// +/// Mirrors `JOBOBJECT_BASIC_PROCESS_ID_LIST`, whose trailing array the SDK +/// declares as one element. Field order and `#[repr(C)]` are the SDK's, because +/// the kernel writes this layout directly. +#[repr(C)] +struct PidList { + number_of_assigned_processes: u32, + number_of_process_ids_in_list: u32, + process_ids: [usize; PROCESS_ID_LIST_WIDTH], +} + +impl Default for PidList { + fn default() -> Self { + Self { + number_of_assigned_processes: 0, + number_of_process_ids_in_list: 0, + process_ids: [0; PROCESS_ID_LIST_WIDTH], + } + } +} + +/// Mark `command` to create its child suspended. +pub fn set_suspended_creation_flags(command: &mut Command) { + command.creation_flags(CREATE_SUSPENDED_FLAG); +} + +/// The same, for the async command the supervisor's spawn path uses. +pub fn set_suspended_creation_flags_async(command: &mut tokio::process::Command) { + command.creation_flags(CREATE_SUSPENDED_FLAG); +} + +/// Find the primary thread of `pid` and open it with resume rights. +/// +/// A suspended process has exactly one thread, so the first match is the one to +/// resume. +/// +/// Safety: the snapshot handle is checked against `INVALID_HANDLE_VALUE` before +/// use and closed on every exit path. `entry` is a correctly-sized stack value +/// whose `dwSize` is set as the API requires. +pub fn open_first_thread(pid: u32) -> io::Result> { + let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPTHREAD, 0) }; + if snapshot == INVALID_HANDLE_VALUE { + return Err(io::Error::last_os_error()); + } + + let mut entry = THREADENTRY32 { + dwSize: size_of::() as u32, + ..Default::default() + }; + let mut found = None; + let mut ok = unsafe { Thread32First(snapshot, &mut entry) }; + while ok != 0 { + if entry.th32OwnerProcessID == pid { + let thread = unsafe { OpenThread(THREAD_SUSPEND_RESUME, 0, entry.th32ThreadID) }; + if thread.is_null() { + let error = io::Error::last_os_error(); + unsafe { CloseHandle(snapshot) }; + return Err(error); + } + found = Some(thread); + break; + } + ok = unsafe { Thread32Next(snapshot, &mut entry) }; + if ok == 0 { + if let Err(error) = exhaustion_or_error() { + unsafe { CloseHandle(snapshot) }; + return Err(error); + } + } + } + + // A first call that reported exhaustion still needs classifying: treating a + // broken snapshot as "this process has no threads" would make the caller kill + // a child it could have started. + if found.is_none() { + if let Err(error) = exhaustion_or_error() { + unsafe { CloseHandle(snapshot) }; + return Err(error); + } + } + + unsafe { CloseHandle(snapshot) }; + Ok(found) +} + +/// Classify a `0` from the toolhelp thread walk: `Ok(())` for ordinary +/// exhaustion, `Err` for a real failure. +fn exhaustion_or_error() -> Result<(), io::Error> { + let code = unsafe { GetLastError() }; + if code == ERROR_NO_MORE_FILES { + Ok(()) + } else { + Err(io::Error::from_raw_os_error(code as i32)) + } +} + +/// Resume a suspended thread and close the handle opened for it. +/// +/// Safety: `thread` was opened by `open_first_thread` with resume rights, is +/// resumed at most once, and is closed exactly once here. +pub fn resume_and_close(thread: HANDLE) -> io::Result<()> { + let previous = unsafe { ResumeThread(thread) }; + let result = if previous == u32::MAX { + Err(io::Error::last_os_error()) + } else { + Ok(()) + }; + unsafe { CloseHandle(thread) }; + result +} + +/// Whether `pid` names a process that is still running. +/// +/// `OpenProcess` succeeding is NOT the answer: a terminated process stays +/// openable while any handle to it survives, and a killed module's handles +/// outlive the kill. That version reported a dead parent as alive and made a +/// passing teardown look like a leak — so liveness is read from the process's +/// signal state, which is what `WaitForSingleObject` with a zero timeout +/// answers. +/// +/// Safety: the handle is opened with synchronize rights (the minimum +/// `WaitForSingleObject` accepts), checked for null (the ordinary "no such +/// process" answer, reported as absence rather than an error), used for exactly +/// one zero-timeout wait, and closed. +pub fn process_exists(pid: u32) -> bool { + let handle = unsafe { OpenProcess(PROCESS_SYNCHRONIZE, 0, pid) }; + if handle.is_null() { + return false; + } + // Signalled means terminated; a timeout means it is still running. + let alive = unsafe { WaitForSingleObject(handle, 0) } == WAIT_TIMEOUT; + unsafe { CloseHandle(handle) }; + alive +} diff --git a/crates/subc-jobobject/tests/containment.rs b/crates/subc-jobobject/tests/containment.rs new file mode 100644 index 00000000..7514c047 --- /dev/null +++ b/crates/subc-jobobject/tests/containment.rs @@ -0,0 +1,298 @@ +//! Containment tests: a grandchild must not outlive teardown. +//! +//! These are the evidence for the fix. The failure they defend against is a +//! module whose helper process survives the module: on this machine the Synapse +//! embedding module's CUDA worker holds ~2.2 GB of VRAM, so a leaked grandchild +//! is a leaked GPU allocation, and a day of restarts compounds it. +//! +//! Every assertion names a **pid**, not "the job is empty". A test that only +//! checked the job's own count would pass against a job that was never populated +//! and against a tree that leaked but was never assigned. + +#![cfg(windows)] + +use std::{ + path::PathBuf, + process::{Command, Stdio}, + time::Duration, +}; + +use subc_jobobject::{process_exists, wait_for_process_exit, JobObject}; + +/// How long a killed tree is given to leave the process table. +/// +/// `TerminateJobObject` returns once members are signalled, not once they are +/// reaped, so teardown is asynchronous and an immediate assertion would be +/// racing the kernel. +const EXIT_DEADLINE: Duration = Duration::from_secs(10); + +/// How long the grandchild's pid file is given to appear. +const GRANDCHILD_APPEARANCE_DEADLINE: Duration = Duration::from_secs(10); + +/// The fixture binary, expected beside this test executable. +/// +/// The existence check is here because a bare spawn failure reads as a broken +/// test rather than an unbuilt dependency: `cargo test -p subc-jobobject` builds +/// `[[bin]]` targets, `--lib` does not. +fn fixture_path() -> PathBuf { + let mut path = std::env::current_exe().expect("current_exe available in tests"); + path.pop(); + path.pop(); + path.push("jobobject-fixture.exe"); + assert!( + path.exists(), + "jobobject-fixture not built at {}: run `cargo test -p subc-jobobject` \ + (which builds [[bin]] targets) rather than `--lib`", + path.display() + ); + path +} + +/// A spawn whose parent has already produced its grandchild. +struct Fixture { + child: std::process::Child, + grandchild_pid: u32, + _dir: TempDir, +} + +/// Minimal RAII temp dir, matching the daemon's `TestTempDir` convention of +/// keeping orphans attributable rather than hand-assembled. +struct TempDir(PathBuf); + +impl TempDir { + fn new(label: &str) -> Self { + let path = std::env::temp_dir() + .join("subc-jobobject-tests") + .join(format!("{label}-{}", std::process::id())); + std::fs::create_dir_all(&path).expect("create test temp dir"); + Self(path) + } + + fn join(&self, name: &str) -> PathBuf { + self.0.join(name) + } +} + +impl Drop for TempDir { + fn drop(&mut self) { + // A panicking test's evidence must outlive it, so the tree is kept and + // its path printed rather than silently removed. + if std::thread::panicking() { + eprintln!("fixture temp dir kept for inspection: {}", self.0.display()); + return; + } + let _ = std::fs::remove_dir_all(&self.0); + } +} + +/// Spawn the fixture parent, contained, and wait for its grandchild to appear. +/// +/// The wait is on the grandchild's published pid, not a sleep: the test must not +/// assert against a tree that has not finished forming, or it would pass for the +/// wrong reason. +fn spawn_contained_fixture(label: &str) -> (Fixture, JobObject) { + let dir = TempDir::new(label); + let grandchild_pid_file = dir.join("grandchild.pid"); + + let mut command = Command::new(fixture_path()); + command + .env("SUBC_JOBOBJECT_FIXTURE_MODE", "parent") + .env("SUBC_JOBOBJECT_GRANDCHILD_PID_FILE", &grandchild_pid_file) + // The fixture parks; it must not hold this test's stdout/stderr open. + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()); + + let (child, job) = { + let contained = subc_jobobject::spawn_contained(&mut command).expect("spawn fixture"); + (contained.child, contained.job) + }; + + let grandchild_pid = wait_for_pid_file(&grandchild_pid_file) + .unwrap_or_else(|| panic!("grandchild pid file never appeared at {grandchild_pid_file:?}")); + + ( + Fixture { + child, + grandchild_pid, + _dir: dir, + }, + job, + ) +} + +/// Poll for the pid file and parse it. +fn wait_for_pid_file(path: &PathBuf) -> Option { + let deadline = std::time::Instant::now() + GRANDCHILD_APPEARANCE_DEADLINE; + loop { + if let Ok(contents) = std::fs::read_to_string(path) { + if let Ok(pid) = contents.trim().parse() { + return Some(pid); + } + } + if std::time::Instant::now() >= deadline { + return None; + } + std::thread::sleep(Duration::from_millis(10)); + } +} + +/// The fix: killing the job takes the grandchild with it. +/// +/// This is the assertion the whole change exists for. Before containment the +/// grandchild survived — it was a separate process whose parent had been +/// `TerminateProcess`'d, so it could not be reached by any walk of the tree. +#[test] +fn terminating_the_job_kills_the_grandchild() { + let (mut fixture, job) = spawn_contained_fixture("terminate"); + + let parent_pid = fixture.child.id(); + assert!( + process_exists(fixture.grandchild_pid), + "grandchild {} must be alive before teardown, or this test proves nothing", + fixture.grandchild_pid + ); + + job.terminate().expect("terminate job"); + + assert!( + wait_for_process_exit(fixture.grandchild_pid, EXIT_DEADLINE), + "grandchild {} survived job termination; the tree was not contained", + fixture.grandchild_pid + ); + assert!( + wait_for_process_exit(parent_pid, EXIT_DEADLINE), + "parent {parent_pid} survived job termination" + ); + + let _ = fixture.child.wait(); +} + +/// The crash-durability guarantee: the tree dies when the handle closes, with no +/// code running to ask it to. +/// +/// This is the case `taskkill /T` cannot cover at all — a daemon that dies +/// cannot call anything, so containment has to be a kernel property of the +/// handle rather than a teardown step. +#[test] +fn dropping_the_job_kills_the_grandchild() { + let (mut fixture, job) = spawn_contained_fixture("drop"); + + assert!(process_exists(fixture.grandchild_pid)); + + // No `terminate` call: closing the last handle is the whole mechanism. + drop(job); + + assert!( + wait_for_process_exit(fixture.grandchild_pid, EXIT_DEADLINE), + "grandchild {} survived the job handle closing", + fixture.grandchild_pid + ); + + let _ = fixture.child.wait(); +} + +/// The mutation control the review asked for: with nothing assigned, the +/// grandchild survives — and this test goes red if that ever stops being true. +/// +/// Without this, a passing `terminating_the_job_kills_the_grandchild` would not +/// distinguish "the job contained the tree" from "the fixture's grandchild died +/// for some unrelated reason". Here the child is spawned and resumed exactly as +/// the fix does, but never assigned, so the only difference is membership. +#[test] +fn an_unassigned_grandchild_survives_teardown() { + let dir = TempDir::new("unassigned"); + let grandchild_pid_file = dir.join("grandchild.pid"); + + let mut command = Command::new(fixture_path()); + command + .env("SUBC_JOBOBJECT_FIXTURE_MODE", "parent") + .env("SUBC_JOBOBJECT_GRANDCHILD_PID_FILE", &grandchild_pid_file) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()); + subc_jobobject::suspend_on_create(&mut command); + let mut child = command.spawn().expect("spawn unassigned fixture"); + // Resumed, so the fixture actually runs and forms its tree — but no job. + subc_jobobject::resume_main_thread(child.id()).expect("resume"); + + let grandchild_pid = + wait_for_pid_file(&grandchild_pid_file).expect("grandchild pid file never appeared"); + assert!(process_exists(grandchild_pid), "grandchild must be alive"); + + // Kill only the direct child, exactly as the pre-fix teardown did. + child.kill().expect("kill direct child"); + let _ = child.wait(); + + assert!( + process_exists(grandchild_pid), + "grandchild {grandchild_pid} died with its parent, which would mean this \ + control no longer distinguishes contained from uncontained — the \ + regression test above would then be passing vacuously" + ); + + // Leave nothing behind: this is the leak the fix prevents, so the test that + // demonstrates it must clean it up itself. + let mut cleanup = Command::new("taskkill.exe"); + cleanup + .args(["/PID", &grandchild_pid.to_string(), "/T", "/F"]) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()); + let _ = cleanup.status(); + assert!( + wait_for_process_exit(grandchild_pid, EXIT_DEADLINE), + "control test could not clean up grandchild {grandchild_pid}" + ); +} + +/// Assignment must happen while the child is suspended, so the job holds the +/// process before it can create anything. +/// +/// Proves the ordering contract rather than assuming it: a child assigned after +/// it ran would already have had the chance to spawn an escapee. +#[test] +fn a_suspended_child_is_assigned_before_it_runs() { + let dir = TempDir::new("suspended"); + let pid_file = dir.join("grandchild.pid"); + + let mut command = Command::new(fixture_path()); + command + .env("SUBC_JOBOBJECT_FIXTURE_MODE", "parent") + .env("SUBC_JOBOBJECT_GRANDCHILD_PID_FILE", &pid_file) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()); + subc_jobobject::suspend_on_create(&mut command); + let child = command.spawn().expect("spawn suspended fixture"); + + // While suspended the fixture has run nothing, so it cannot have spawned + // its grandchild yet. If this file exists here, `CREATE_SUSPENDED` was not + // applied and the containment window is open. + assert!( + !pid_file.exists(), + "the child ran before it was resumed; the suspended-create contract is broken" + ); + + let job = JobObject::new().expect("create job"); + job.assign(&child).expect("assign while suspended"); + assert_eq!( + job.process_count().expect("count"), + 1, + "only the child itself is contained at this point" + ); + + subc_jobobject::resume_main_thread(child.id()).expect("resume"); + + let grandchild_pid = wait_for_pid_file(&pid_file).expect("grandchild pid file never appeared"); + assert!(process_exists(grandchild_pid)); + + job.terminate().expect("terminate"); + assert!( + wait_for_process_exit(grandchild_pid, EXIT_DEADLINE), + "grandchild {grandchild_pid} survived: it escaped the suspension window" + ); + + let mut child = child; + let _ = child.wait(); +} From a4720c95f7d73fe95f1b393a176ecb9bf60c211b Mon Sep 17 00:00:00 2001 From: Qiiks Date: Thu, 24 Sep 2026 04:10:51 +0530 Subject: [PATCH 2/2] subc-daemon: document what KILL_ON_JOB_CLOSE costs on a Windows stop The job field's doc explained the limit only as crash containment. It also fires on every ordinary Windows daemon stop, because the daemon has no stop notice there: the SIGTERM handler is #[cfg(unix)], so a stop is taskkill or the scheduler's /End, the kernel closes the job handle, and every module is TerminateProcess'd at once rather than running its own teardown. The Windows stop path is still outstanding; the Unix arm now has one (process groups plus the child roster). Stating the trade and naming the fix, so the next reader does not take the crash rationale as the whole story. Requested in review of #111. --- crates/subc-daemon/src/supervise.rs | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/crates/subc-daemon/src/supervise.rs b/crates/subc-daemon/src/supervise.rs index 26a90697..91ec8ebf 100644 --- a/crates/subc-daemon/src/supervise.rs +++ b/crates/subc-daemon/src/supervise.rs @@ -97,6 +97,22 @@ struct SupervisedChild { /// Dropping this handle is what reaps a surviving tree when no supervisor /// code runs — a daemon crash — because the job carries /// `JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE`. + /// + /// That limit is not crash-only, and the difference is worth knowing: a + /// Windows daemon *stop* is `taskkill` or the scheduler's `/End` — the + /// SIGTERM handler is `#[cfg(unix)]` — so the daemon dies with no stop + /// notice and the kernel closes the job handle, `TerminateProcess`ing every + /// module at once. Before this change they survived that, saw EOF on the + /// control socket, and ran their own teardown; Unix keeps that path + /// deliberately, so a module can seal a WAL or close a capture rather than + /// be killed mid-write. So this trades graceful teardown on every Windows + /// daemon stop for containment on a crash, which is the right way round + /// today: orphaned GPU workers are a reported, recurring problem, and the + /// modules that write most heavily do not run on Windows. + /// + /// The fix is a real Windows stop path — the daemon draining before it + /// exits, the twin of the Unix SIGTERM handler. Once it exists, this limit + /// reaches only what the drain left behind, which is what it should reach. #[cfg(windows)] job: Option, stdout_pump: Option>,