diff --git a/tests/cluster_two_node.rs b/tests/cluster_two_node.rs new file mode 100644 index 0000000..712b370 --- /dev/null +++ b/tests/cluster_two_node.rs @@ -0,0 +1,108 @@ +//! RFC 010 c4 — two-node harness smoke tests. +//! +//! Roadmap: "spawn two, handshake-less connect, both exit clean." The +//! listener node binds port 0 and announces its concrete address; the +//! dialer connects raw (no Hello — c5 doesn't exist yet), pushes one +//! Heartbeat through the real framed codec, and closes. Assertions are on +//! protocol-visible lines only. Flake budget: see tests/common/mod.rs. +#![cfg(feature = "cluster")] + +mod common; + +use common::{maybe_child, spawn_node}; +use smarm::cluster::envelope::Frame; +use smarm::cluster::transport::tcp::TcpTransport; +use smarm::cluster::transport::{FramedConn, Transport}; + +const ROLES: &[(&str, fn())] = &[ + ("listener", role_listener), + ("dialer", role_dialer), + ("hang", role_hang), + ("fail", role_fail), +]; + +fn role_listener() { + smarm::run(|| { + let mut l = TcpTransport.listen("127.0.0.1:0").unwrap(); + println!("LISTENING {}", l.local_addr()); + let mut fc = FramedConn::new(l.accept().unwrap()); + match fc.recv() { + Ok(Some(Frame::Heartbeat)) => println!("RECV heartbeat"), + other => { + println!("RECV unexpected: {other:?}"); + std::process::exit(3); + } + } + match fc.recv() { + Ok(None) => println!("PEER-CLOSED clean"), + other => { + println!("PEER-CLOSED unexpected: {other:?}"); + std::process::exit(3); + } + } + }); + println!("EXIT ok"); +} + +fn role_dialer() { + let addr = std::env::var("SMARM_PEER_ADDR").expect("SMARM_PEER_ADDR not set"); + smarm::run(move || { + let mut fc = FramedConn::new(TcpTransport.dial(&addr).unwrap()); + fc.send(&Frame::Heartbeat).unwrap(); + fc.close(); + println!("SENT heartbeat"); + }); + println!("EXIT ok"); +} + +fn role_hang() { + println!("HANGING"); + loop { + std::thread::sleep(std::time::Duration::from_secs(3600)); + } +} + +fn role_fail() { + std::process::exit(7); +} + +/// The roadmap smoke test: two real processes, raw transport connect, one +/// frame across, clean close observed on both sides, both exit 0. +#[test] +fn two_nodes_connect_and_exit_clean() { + maybe_child(ROLES); + let mut listener = spawn_node("listener", &[]); + let addr = listener.wait_listening(); + let mut dialer = spawn_node("dialer", &[("SMARM_PEER_ADDR", &addr)]); + dialer.wait_line("SENT heartbeat", |l| l == "SENT heartbeat"); + listener.wait_line("RECV heartbeat", |l| l == "RECV heartbeat"); + listener.wait_line("clean peer close", |l| l == "PEER-CLOSED clean"); + dialer.wait_exit_ok(); + listener.wait_exit_ok(); +} + +/// Reap guarantee: dropping a Node kills a hung child — no orphan survives +/// a panicking test. +#[test] +fn drop_reaps_hung_node() { + maybe_child(ROLES); + let mut node = spawn_node("hang", &[]); + node.wait_line("HANGING", |l| l == "HANGING"); + let pid = node.pid().expect("live child has a pid") as libc::pid_t; + drop(node); + // After Drop's kill+wait the pid is fully reaped: signalling it fails + // with ESRCH (pid-reuse in this instant is not a realistic race). + let rc = unsafe { libc::kill(pid, 0) }; + assert_eq!(rc, -1, "process still signallable after Drop"); + let errno = std::io::Error::last_os_error().raw_os_error(); + assert_eq!(errno, Some(libc::ESRCH), "expected ESRCH, got {errno:?}"); +} + +/// Nonzero child exits surface as statuses, not hangs or panics. +#[test] +fn nonzero_exit_is_reported() { + maybe_child(ROLES); + let mut node = spawn_node("fail", &[]); + let status = node.wait_exit(); + assert_eq!(status.code(), Some(7)); +} diff --git a/tests/common/mod.rs b/tests/common/mod.rs new file mode 100644 index 0000000..ba522e7 --- /dev/null +++ b/tests/common/mod.rs @@ -0,0 +1,239 @@ +//! 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") + } + } + + /// 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) => { + panic!( + "node {:?}: timed out waiting for {what} after {WAIT:?}; transcript:\n{}", + self.role, + self.dump() + ); + } + Err(RecvTimeoutError::Disconnected) => { + 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(); + } +}