#![allow(dead_code)]
use std::io::Read;
use std::net::SocketAddr;
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Output, Stdio};
use std::time::{Duration, Instant};
pub fn target_dir() -> PathBuf {
std::env::var_os("CARGO_TARGET_DIR").map_or_else(
|| PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("target"),
PathBuf::from,
)
}
pub struct ChildGuard(Option<Child>);
impl ChildGuard {
pub fn new(child: Child) -> Self {
Self(Some(child))
}
fn take_status(&mut self) -> Option<std::process::ExitStatus> {
self.0.as_mut().and_then(|child| child.try_wait().ok())?
}
}
impl Drop for ChildGuard {
fn drop(&mut self) {
if let Some(mut child) = self.0.take() {
let _ = child.kill();
let _ = child.wait();
}
}
}
fn resolve_example_binary(
rel_path: &str,
build_hint: impl FnOnce(&str) -> String,
) -> (PathBuf, &str) {
let binary = target_dir().join(rel_path);
let example_name = Path::new(rel_path)
.file_name()
.and_then(|name| name.to_str())
.unwrap_or(rel_path);
assert!(
binary.is_file(),
"{} is missing. This leg FAILS rather than skipping, by design: a skip would \
restore the unenforced 'the example demonstrates the fix' criterion it exists \
to close. Build it first with {}.",
binary.display(),
build_hint(example_name)
);
assert_binary_is_not_stale(&binary, example_name);
(binary, example_name)
}
pub fn spawn_example(rel_path: &str, bind_addr: &str) -> (SocketAddr, ChildGuard) {
let (binary, _example_name) = resolve_example_binary(rel_path, |name| {
format!("`cargo build --features full --example {name}`")
});
let addr: SocketAddr = bind_addr
.parse()
.unwrap_or_else(|error| panic!("`{bind_addr}` is not a socket address: {error}"));
let child = Command::new(&binary)
.arg(bind_addr)
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.unwrap_or_else(|error| panic!("could not spawn {}: {error}", binary.display()));
(addr, ChildGuard::new(child))
}
const DRAIN_GRACE: Duration = Duration::from_millis(500);
#[derive(Clone)]
struct Drain {
captured: std::sync::Arc<std::sync::Mutex<Vec<u8>>>,
}
impl Drain {
fn new() -> Self {
Self {
captured: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
}
}
fn captured(&self) -> Vec<u8> {
self.captured
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
}
fn drain_into(mut source: impl Read, sink: &Drain) {
let mut chunk = [0_u8; 8192];
loop {
match source.read(&mut chunk) {
Ok(0) => break,
Err(error) if error.kind() == std::io::ErrorKind::Interrupted => {},
Err(_) => break,
Ok(count) => {
sink.captured
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.extend_from_slice(&chunk[..count]);
},
}
}
}
fn settle(readers: &[&std::thread::JoinHandle<()>]) {
let deadline = Instant::now() + DRAIN_GRACE;
while Instant::now() < deadline {
if readers.iter().all(|reader| reader.is_finished()) {
return;
}
std::thread::sleep(Duration::from_millis(10));
}
}
pub struct ExampleLeg<'a> {
pub rel_path: &'a str,
pub banner: &'a str,
pub rebuild: &'a str,
pub claim: &'a str,
pub banner_means: &'a str,
}
pub fn assert_ran_and_printed_banner(leg: &ExampleLeg<'_>, output: &Output) {
let rel_path = leg.rel_path;
let stdout = String::from_utf8_lossy(&output.stdout);
assert!(
output.status.success(),
"`{rel_path}` exited with {} rather than succeeding. {}\n\
Rebuild with {} — note that `cargo test` does NOT rebuild examples.\n\
--- stdout ---\n{}\n--- stderr ---\n{}",
output.status,
leg.claim,
leg.rebuild,
stdout,
String::from_utf8_lossy(&output.stderr)
);
assert!(
!stdout.trim().is_empty(),
"`{rel_path}` exited 0 but printed nothing on stdout. The evidence this leg asserts on \
IS the printed transcript, so an empty stdout means the run proved nothing.\n\
--- stderr ---\n{}",
String::from_utf8_lossy(&output.stderr)
);
assert!(
stdout.contains(leg.banner),
"`{rel_path}` exited 0 but never printed its completion banner {:?}. {}\n\
--- stdout ---\n{stdout}",
leg.banner,
leg.banner_means
);
}
pub fn run_example_to_completion(rel_path: &str, args: &[&str], timeout: Duration) -> Output {
let (binary, example_name) = resolve_example_binary(rel_path, |name| {
format!(
"`cargo build --example {name}` (add `-p <crate>` when the example lives under \
`crates/*/examples/`)"
)
});
let mut child = Command::new(&binary)
.args(args)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap_or_else(|error| panic!("could not spawn {}: {error}", binary.display()));
let stdout_pipe = child
.stdout
.take()
.expect("stdout was piped, so the handle is present");
let stderr_pipe = child
.stderr
.take()
.expect("stderr was piped, so the handle is present");
let stdout_drain = Drain::new();
let stderr_drain = Drain::new();
let stdout_reader = {
let sink = stdout_drain.clone();
std::thread::spawn(move || drain_into(stdout_pipe, &sink))
};
let stderr_reader = {
let sink = stderr_drain.clone();
std::thread::spawn(move || drain_into(stderr_pipe, &sink))
};
let deadline = Instant::now() + timeout;
loop {
match child.try_wait() {
Ok(Some(status)) => {
settle(&[&stdout_reader, &stderr_reader]);
return Output {
status,
stdout: stdout_drain.captured(),
stderr: stderr_drain.captured(),
};
},
Ok(None) => {},
Err(error) => panic!("cannot poll {}: {error}", binary.display()),
}
if Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
settle(&[&stdout_reader, &stderr_reader]);
let stdout = stdout_drain.captured();
let stderr = stderr_drain.captured();
panic!(
"{} did not exit within {timeout:?}: this leg converts a hang into a red rather \
than blocking the integration suite forever. The child was killed and reaped.\n\
If the example is simply slower than its budget, raise the budget constant in \
the OWNING test file (budgets are per-leg by design); if it is wedged, rebuild \
it with `cargo build --example {example_name}` (add `-p <crate>` for a \
`crates/*/examples/` binary) and run it by hand.\n\
--- partial stdout ---\n{}\n--- partial stderr ---\n{}",
binary.display(),
String::from_utf8_lossy(&stdout),
String::from_utf8_lossy(&stderr)
);
}
std::thread::sleep(Duration::from_millis(50));
}
}
fn assert_binary_is_not_stale(binary: &Path, example_name: &str) {
let Some(binary_mtime) = modified_at(binary) else {
return;
};
let manifest = PathBuf::from(env!("CARGO_MANIFEST_DIR"));
let mut newest: Option<(PathBuf, std::time::SystemTime)> = None;
let mut consider = |path: PathBuf| {
if let Some(mtime) = modified_at(&path) {
if newest.as_ref().is_none_or(|(_, best)| mtime > *best) {
newest = Some((path, mtime));
}
}
};
consider(manifest.join("examples").join(format!("{example_name}.rs")));
for source in rust_sources_under(&manifest.join("src")) {
consider(source);
}
let Some((newest_path, newest_mtime)) = newest else {
return;
};
assert!(
newest_mtime <= binary_mtime,
"{} is STALE: it was built BEFORE {}, so this leg would exercise OLD CODE \
and report a result about a source tree that no longer exists.\n\
`cargo test --test <name>` does NOT rebuild examples — target selection \
excludes them — so the binary in `target/` is whatever was left there last.\n\
REBUILD IT: `cargo build --features full --example {example_name}`\n\
This assertion exists because the merged-tree run of a Phase-118.1 leg once \
disagreed with the worktree run for exactly this reason.",
binary.display(),
newest_path.display()
);
}
fn modified_at(path: &Path) -> Option<std::time::SystemTime> {
std::fs::metadata(path).ok()?.modified().ok()
}
fn rust_sources_under(root: &Path) -> Vec<PathBuf> {
let mut found = Vec::new();
let mut stack = vec![root.to_path_buf()];
while let Some(dir) = stack.pop() {
let Ok(entries) = std::fs::read_dir(&dir) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
stack.push(path);
} else if path.extension().is_some_and(|ext| ext == "rs") {
found.push(path);
}
}
}
found
}
pub async fn wait_until_listening(
addr: SocketAddr,
guard: &mut ChildGuard,
ready_timeout: Duration,
) {
let deadline = Instant::now() + ready_timeout;
while Instant::now() < deadline {
if let Some(status) = guard.take_status() {
panic!(
"the example exited before binding {addr} (status {status}). The most likely \
cause is that {addr} is already held — check with \
`lsof -nP -iTCP:{} -sTCP:LISTEN`.",
addr.port()
);
}
if tokio::net::TcpStream::connect(addr).await.is_ok() {
return;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
panic!("the example never accepted a connection on {addr} within {ready_timeout:?}");
}
pub async fn wait_until_released(addr: SocketAddr, release_timeout: Duration) {
let deadline = Instant::now() + release_timeout;
while Instant::now() < deadline {
if tokio::net::TcpStream::connect(addr).await.is_err() {
return;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
panic!(
"{addr} still accepts connections {release_timeout:?} after the child was killed: the \
teardown did not tear down, and the next run would talk to a stale listener"
);
}