//! RFC 010 c4 — subprocess multi-node test harness. //! //! The runtime is a process singleton, so two real nodes means two //! processes. This harness re-execs the *current test binary* as node //! processes (precedent: tests/stack_diag.rs), tails their output live, //! waits on protocol-visible lines, and reaps reliably no matter how the //! test dies. //! //! Usage, per test file: //! //! - Declare roles as plain `fn()`s. A role prints protocol-visible facts //! as single lines (Rust's piped stdout is line-buffered, so `println!` //! is enough) and exits. //! - **Every** `#[test]` in the file starts with //! [`maybe_child`]`(ROLES)` — in the child re-exec, whichever test //! libtest runs first performs the role and exits before the rest of the //! suite runs (children are spawned with `--test-threads=1 --quiet`). //! - The parent side spawns nodes with [`spawn_node`], waits on lines with //! [`Node::wait_line`], and on exits with [`Node::wait_exit`]. //! //! Port assignment: children bind port 0 and *report* the concrete address //! (e.g. `LISTENING 127.0.0.1:41733`) rather than the parent pre-picking a //! port — no bind/steal race by construction. //! //! Reaping: [`Node`]'s `Drop` SIGKILLs and `wait(2)`s the child, so a //! panicking test (including a `wait_line` timeout) leaves no orphan and //! no zombie. Tail threads exit on pipe EOF. //! //! Flake budget (explicit, per roadmap): every wait is bounded by //! [`WAIT`] (10 s) against a typical cost of well under 1 s; the smoke //! suite ran 10/10 clean at authoring time. Treat >1 failure in 100 runs //! as a harness or runtime regression, not weather. On timeout the panic //! message carries the node's full transcript so far. #![allow(dead_code)] // Reusable surface: later phases use more of it than any one file. use std::env; use std::io::{BufRead, BufReader}; use std::process::{Child, Command, ExitStatus, Stdio}; use std::sync::mpsc::{Receiver, RecvTimeoutError}; use std::time::{Duration, Instant}; /// Env var selecting the child role in a re-exec. const ROLE_ENV: &str = "SMARM_TWO_NODE_ROLE"; /// Upper bound for every wait in the harness. See the flake budget above. pub const WAIT: Duration = Duration::from_secs(10); /// In the child re-exec: run the matching role and exit. In the parent (no /// role env set): return immediately. Call this first in every `#[test]` of /// any file using the harness, passing the file's full role table. pub fn maybe_child(roles: &[(&str, fn())]) { let role = match env::var(ROLE_ENV) { Ok(r) => r, Err(_) => return, }; for (name, f) in roles { if *name == role { f(); std::process::exit(0); } } eprintln!("two_node harness: unknown role {role:?}"); std::process::exit(2); } /// One spawned node process with live-tailed output. pub struct Node { /// Role name, for panic messages. pub role: String, child: Option, stdout_rx: Receiver, stderr_rx: Receiver, /// Every line consumed from stdout/stderr so far, for failure dumps. transcript: Vec, } fn tail(stream: impl std::io::Read + Send + 'static, prefix: &'static str) -> Receiver { let (tx, rx) = std::sync::mpsc::channel(); std::thread::spawn(move || { for line in BufReader::new(stream).lines() { let line = match line { Ok(l) => l, Err(_) => break, }; // Receiver gone (Node dropped): stop tailing. if tx.send(format!("{prefix}{line}")).is_err() { break; } } }); rx } /// Re-exec the current test binary as `role`, with any extra env vars. pub fn spawn_node(role: &str, extra_env: &[(&str, &str)]) -> Node { let exe = env::current_exe().expect("current_exe"); let mut cmd = Command::new(exe); cmd.env(ROLE_ENV, role) // --test-threads=1: exactly one test fn starts, hits maybe_child, // and becomes the role. --nocapture: libtest must not swallow the // role's println! lines — the parent tails them live. .args(["--test-threads=1", "--quiet", "--nocapture"]) .stdout(Stdio::piped()) .stderr(Stdio::piped()); for (k, v) in extra_env { cmd.env(k, v); } let mut child = cmd.spawn().expect("failed to spawn node process"); let stdout_rx = tail(child.stdout.take().expect("piped stdout"), ""); let stderr_rx = tail(child.stderr.take().expect("piped stderr"), "[stderr] "); Node { role: role.to_string(), child: Some(child), stdout_rx, stderr_rx, transcript: Vec::new(), } } impl Node { fn drain_stderr(&mut self) { while let Ok(l) = self.stderr_rx.try_recv() { self.transcript.push(l); } } fn dump(&self) -> String { if self.transcript.is_empty() { "".to_string() } else { self.transcript.join("\n") } } /// Every stdout/stderr line seen so far, in arrival order. For /// ordering-proof assertions ("by the time X arrived, Y had not"). #[allow(dead_code)] pub fn transcript(&self) -> &[String] { &self.transcript } /// Wait until a stdout line satisfies `pred`; return it. Panics with the /// full transcript after [`WAIT`]. `what` names the expectation in the /// panic message. pub fn wait_line(&mut self, what: &str, pred: impl Fn(&str) -> bool) -> String { let deadline = Instant::now() + WAIT; loop { self.drain_stderr(); let left = deadline.saturating_duration_since(Instant::now()); match self.stdout_rx.recv_timeout(left) { Ok(line) => { self.transcript.push(line.clone()); if pred(&line) { return line; } } Err(RecvTimeoutError::Timeout) => { // Pull in whatever stderr arrived since the last drain, // so a role's eprintln! diagnostics survive into the dump. self.drain_stderr(); panic!( "node {:?}: timed out waiting for {what} after {WAIT:?}; transcript:\n{}", self.role, self.dump() ); } Err(RecvTimeoutError::Disconnected) => { self.drain_stderr(); panic!( "node {:?}: output closed while waiting for {what}; transcript:\n{}", self.role, self.dump() ); } } } } /// Shorthand: wait for a `LISTENING ` announcement, return the addr. pub fn wait_listening(&mut self) -> String { let line = self.wait_line("LISTENING announcement", |l| l.starts_with("LISTENING ")); line["LISTENING ".len()..].to_string() } /// Wait for the process to exit; panics with the transcript on timeout. pub fn wait_exit(&mut self) -> ExitStatus { let deadline = Instant::now() + WAIT; loop { let polled = match self.child.as_mut() { Some(c) => c.try_wait(), None => panic!("node {:?}: already reaped", self.role), }; match polled { Ok(Some(status)) => { // Drain remaining output into the transcript for dumps. self.drain_stderr(); while let Ok(l) = self.stdout_rx.try_recv() { self.transcript.push(l); } self.child = None; return status; } Ok(None) => { if Instant::now() >= deadline { self.drain_stderr(); self.kill(); panic!( "node {:?}: did not exit within {WAIT:?}; transcript:\n{}", self.role, self.dump() ); } std::thread::sleep(Duration::from_millis(10)); } Err(e) => panic!("node {:?}: try_wait failed: {e}", self.role), } } } /// Wait for exit and require success, dumping the transcript otherwise. pub fn wait_exit_ok(&mut self) { let status = self.wait_exit(); assert!( status.success(), "node {:?}: exited with {status}; transcript:\n{}", self.role, self.dump() ); } /// The child's OS pid, if not yet reaped. pub fn pid(&self) -> Option { self.child.as_ref().map(Child::id) } /// SIGKILL + reap now (idempotent). pub fn kill(&mut self) { if let Some(mut child) = self.child.take() { let _ = child.kill(); let _ = child.wait(); } } } impl Drop for Node { fn drop(&mut self) { self.kill(); } }