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..91ec8ebf 100644 --- a/crates/subc-daemon/src/supervise.rs +++ b/crates/subc-daemon/src/supervise.rs @@ -92,6 +92,29 @@ 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`. + /// + /// 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>, stderr_pump: Option>, stderr_ring: Arc>, @@ -131,7 +154,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 +5468,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 +5489,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 +5552,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 +5567,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 +8910,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(); +}