#![cfg(feature = "test-hooks")]
use std::io::{BufRead, BufReader, Read, Write};
use std::path::{Path, PathBuf};
use std::process::{Child, Command, ExitStatus, Stdio};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use serde_json::Value;
const CODE_WORD: &str = "gigantic-element-tango";
const BUILD_PROFILE: &str = env!("FILAMENT_BUILD_PROFILE");
#[derive(Debug)]
enum ChildOutcome {
ExitedSuccess(ExitStatus),
ExitedFailure(ExitStatus),
TimedOut,
SpawnFailed(String),
}
#[derive(Debug)]
struct CapturedChild {
outcome: ChildOutcome,
stdout: String,
stderr: String,
events: Vec<CapturedChunk>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ChildStream {
Stdout,
Stderr,
}
#[derive(Debug, Clone)]
struct CapturedChunk {
at: Instant,
stream: ChildStream,
bytes: Vec<u8>,
}
struct LiveChild {
child: Child,
stdout: Arc<Mutex<Vec<u8>>>,
stderr: Arc<Mutex<Vec<u8>>>,
events: Arc<Mutex<Vec<CapturedChunk>>>,
readers: Vec<std::thread::JoinHandle<()>>,
}
fn drain_pipe<R: Read + Send + 'static>(
pipe: R,
buffer: Arc<Mutex<Vec<u8>>>,
events: Arc<Mutex<Vec<CapturedChunk>>>,
stream: ChildStream,
) -> std::thread::JoinHandle<()> {
std::thread::spawn(move || {
let mut reader = pipe;
let mut chunk = [0u8; 4096];
loop {
match reader.read(&mut chunk) {
Ok(0) | Err(_) => break,
Ok(n) => {
let bytes = chunk[..n].to_vec();
buffer.lock().unwrap().extend_from_slice(&bytes);
events.lock().unwrap().push(CapturedChunk {
at: Instant::now(),
stream,
bytes,
});
}
}
}
})
}
fn spawn_captured(mut command: Command) -> Result<LiveChild, String> {
let mut child = command
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.map_err(|e| e.to_string())?;
let stdout = Arc::new(Mutex::new(Vec::new()));
let stderr = Arc::new(Mutex::new(Vec::new()));
let events = Arc::new(Mutex::new(Vec::new()));
let readers = vec![
drain_pipe(child.stdout.take().ok_or("child stdout was not piped")?, stdout.clone(), events.clone(), ChildStream::Stdout),
drain_pipe(child.stderr.take().ok_or("child stderr was not piped")?, stderr.clone(), events.clone(), ChildStream::Stderr),
];
Ok(LiveChild { child, stdout, stderr, events, readers })
}
fn run_captured(command: Command, deadline: Duration) -> CapturedChild {
match spawn_captured(command) {
Ok(child) => child.wait_until(deadline),
Err(error) => CapturedChild {
outcome: ChildOutcome::SpawnFailed(error),
stdout: String::new(),
stderr: String::new(),
events: Vec::new(),
},
}
}
fn first_marker_at(events: &[CapturedChunk], marker: &[u8]) -> Option<Instant> {
let mut ordered = events.to_vec();
ordered.sort_by_key(|event| event.at);
let mut captured = Vec::new();
for event in ordered {
if event.stream != ChildStream::Stderr {
continue;
}
captured.extend_from_slice(&event.bytes);
if captured.windows(marker.len()).any(|window| window == marker) {
return Some(event.at);
}
}
None
}
fn marker_delta(events: &[CapturedChunk]) -> Option<Duration> {
let blocked_at = first_marker_at(events, b"DIRECT-BLOCKED")?;
let fallback_at = first_marker_at(events, b"DIRECT-FALLBACK")?;
fallback_at.checked_duration_since(blocked_at)
}
impl LiveChild {
fn snapshot(&self) -> (String, String) {
let stdout = String::from_utf8_lossy(&self.stdout.lock().unwrap()).into_owned();
let stderr = String::from_utf8_lossy(&self.stderr.lock().unwrap()).into_owned();
(stdout, stderr)
}
fn events(&self) -> Vec<CapturedChunk> {
self.events.lock().unwrap().clone()
}
fn wait_until(mut self, deadline: Duration) -> CapturedChild {
let started = std::time::Instant::now();
let outcome = loop {
match self.child.try_wait() {
Ok(Some(status)) => {
break if status.success() {
ChildOutcome::ExitedSuccess(status)
} else {
ChildOutcome::ExitedFailure(status)
};
}
Ok(None) if started.elapsed() < deadline => {
std::thread::sleep(Duration::from_millis(20));
}
Ok(None) => {
let _ = self.child.kill();
let _ = self.child.wait();
break ChildOutcome::TimedOut;
}
Err(_) => break ChildOutcome::TimedOut,
}
};
if !matches!(outcome, ChildOutcome::TimedOut) {
for reader in std::mem::take(&mut self.readers) {
let _ = reader.join();
}
}
let (stdout, stderr) = self.snapshot();
let events = self.events();
CapturedChild { outcome, stdout, stderr, events }
}
}
fn binary() -> PathBuf {
let cand = PathBuf::from(env!("CARGO_BIN_EXE_filament"));
let profile = BUILD_PROFILE;
use sha2::{Digest, Sha256};
let bytes = std::fs::read(&cand)
.unwrap_or_else(|error| panic!("cannot hash harness binary {}: {error}", cand.display()));
let sha256 = Sha256::digest(bytes)
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();
eprintln!(
"HARNESS-BINARY profile={} path={} sha256={}",
profile,
cand.display(),
sha256,
);
cand
}
fn find_backend_app() -> PathBuf {
let from_manifest = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.parent()
.unwrap()
.join("backend")
.join("app.py");
if from_manifest.exists() {
return from_manifest;
}
PathBuf::from("../backend/app.py")
}
struct Harness {
backend: Child,
backend_port: u16,
daemon_a: Option<Child>,
daemon_b: Option<Child>,
daemon_a_log: Arc<Mutex<Vec<String>>>,
daemon_b_log: Arc<Mutex<Vec<String>>>,
work_dir: PathBuf,
a_dir: PathBuf,
b_dir: PathBuf,
}
impl Drop for Harness {
fn drop(&mut self) {
if let Some(ref mut c) = self.daemon_a {
let _ = c.kill();
let _ = c.wait();
}
if let Some(ref mut c) = self.daemon_b {
let _ = c.kill();
let _ = c.wait();
}
let _ = self.backend.kill();
let _ = self.backend.wait();
let _ = std::fs::remove_dir_all(&self.work_dir);
}
}
fn python_cmd() -> &'static str {
if cfg!(windows) { "python" } else { "python3" }
}
fn find_free_port() -> u16 {
use std::net::TcpListener;
TcpListener::bind("127.0.0.1:0")
.map(|l| l.local_addr().unwrap().port())
.unwrap_or(19079)
}
impl Harness {
fn new() -> Self {
let work = std::env::temp_dir()
.join(format!("filament-harness-{}", std::process::id()));
std::fs::create_dir_all(&work).expect("create work dir");
let a_dir = work.join("a");
let b_dir = work.join("b");
let app_path = find_backend_app();
let port = find_free_port();
let mut backend = {
let mut b = Command::new(python_cmd())
.arg(&app_path)
.env("PORT", port.to_string())
.env("FIL_ASYNC_MODE", "eventlet")
.env("FIL_SELF_MONKEYPATCH", "1")
.env("FIL_CLAIM_LIMIT", "1000000")
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.current_dir(app_path.parent().unwrap())
.spawn()
.expect("start backend");
let stderr = b.stderr.take().unwrap();
std::thread::spawn(move || {
let reader = BufReader::new(stderr);
for line in reader.lines() {
if let Ok(l) = line {
eprintln!("[backend] {l}");
}
}
});
let stdout = b.stdout.take().unwrap();
std::thread::spawn(move || {
let reader = BufReader::new(stdout);
for line in reader.lines() {
if let Ok(l) = line {
eprintln!("[backend] {l}");
}
}
});
b
};
let server_url = format!("http://127.0.0.1:{port}");
let mut backend_ok = false;
for _ in 0..90 {
if reqwest_blocking_head(&format!("{server_url}/api/health")) {
backend_ok = true;
break;
}
std::thread::sleep(Duration::from_millis(500));
}
if !backend_ok {
panic!("backend did not start within 45s at {server_url}");
}
let bin = binary();
for (dir, name) in [(&a_dir, "test-a"), (&b_dir, "test-b")] {
std::fs::create_dir_all(dir).expect("create peer dir");
let rec = dir.join(format!("{name}-recovery.txt"));
let out = Command::new(&bin)
.env("FILAMENT_CONFIG_DIR", dir)
.arg("init")
.arg("--name")
.arg(name)
.arg("--recovery-file")
.arg(&rec)
.arg("--yes")
.output()
.expect("init peer");
assert!(
out.status.success(),
"init {name} failed: {}",
String::from_utf8_lossy(&out.stderr)
);
}
Harness {
backend,
backend_port: port,
daemon_a: None,
daemon_b: None,
daemon_a_log: Arc::new(Mutex::new(Vec::new())),
daemon_b_log: Arc::new(Mutex::new(Vec::new())),
work_dir: work,
a_dir,
b_dir,
}
}
fn server_url(&self) -> String {
format!("http://127.0.0.1:{}", self.backend_port)
}
fn filament_bin(&self) -> &Path {
static BIN: std::sync::OnceLock<PathBuf> = std::sync::OnceLock::new();
BIN.get_or_init(binary)
}
fn spawn_daemon(&mut self, name: &str, config_dir: &Path) -> (Child, Arc<Mutex<Vec<String>>>) {
let bin = self.filament_bin().to_path_buf();
let server = self.server_url();
spawn_daemon_inner(&bin, &server, name, config_dir)
}
fn pair_daemons(&mut self) {
let bin = self.filament_bin().to_path_buf();
let server = self.server_url();
let (child_a, log_a) = spawn_daemon_inner(&bin, &server, "test-a", &self.a_dir);
let (child_b, log_b) = spawn_daemon_inner(&bin, &server, "test-b", &self.b_dir);
self.daemon_a = Some(child_a);
self.daemon_b = Some(child_b);
self.daemon_a_log = log_a;
self.daemon_b_log = log_b;
std::thread::sleep(Duration::from_secs(8));
}
}
fn spawn_daemon_inner(
bin: &Path,
server: &str,
name: &str,
config_dir: &Path,
) -> (Child, Arc<Mutex<Vec<String>>>) {
std::fs::create_dir_all(config_dir).expect("create config dir");
let stderr_log: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
let log_clone = stderr_log.clone();
let mut child = Command::new(bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_CONFIG_DIR", config_dir)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_LOG", "trace")
.env(
"FILAMENT_DIRECT_LOOPBACK_ONLY",
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY")
.unwrap_or_else(|_| "1".into()),
)
.env("FILAMENT_NAME", name)
.arg("up")
.arg("--userspace")
.arg("--shell")
.arg("--i-know")
.arg("--server")
.arg(server)
.arg("--relay")
.arg("--dir")
.arg(config_dir.join("drops").to_str().unwrap())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap_or_else(|e| panic!("spawn daemon {name}: {e}"));
let label1 = name.to_string();
let label2 = name.to_string();
if let Some(stderr) = child.stderr.take() {
std::thread::spawn(move || {
for line in BufReader::new(stderr).lines() {
if let Ok(l) = line {
eprintln!("[{label1} stderr] {l}");
if let Ok(mut log) = log_clone.lock() {
log.push(l);
}
}
}
});
}
if let Some(stdout) = child.stdout.take() {
std::thread::spawn(move || {
for line in BufReader::new(stdout).lines() {
if let Ok(l) = line { eprintln!("[{label2}] {l}"); }
}
});
}
(child, stderr_log)
}
fn wait_for_line(log: &Arc<Mutex<Vec<String>>>, needle: &str, timeout: Duration) -> bool {
let deadline = std::time::Instant::now() + timeout;
loop {
if let Ok(lines) = log.lock() {
if lines.iter().any(|l| l.contains(needle)) {
return true;
}
}
if std::time::Instant::now() >= deadline {
return false;
}
std::thread::sleep(Duration::from_millis(500));
}
}
fn reqwest_blocking_head(url: &str) -> bool {
use std::io::{Read, Write};
use std::net::TcpStream;
let host_port = url
.trim_start_matches("http://")
.trim_end_matches("/api/health");
match TcpStream::connect_timeout(&host_port.parse().unwrap(), Duration::from_secs(3)) {
Ok(mut stream) => {
let _ = stream.set_read_timeout(Some(Duration::from_secs(2)));
let _ = stream.write_all(
format!("GET /api/health HTTP/1.0\r\nHost: {host_port}\r\n\r\n").as_bytes(),
);
let mut buf = [0u8; 256];
if let Ok(n) = stream.read(&mut buf) {
let resp = std::str::from_utf8(&buf[..n]).unwrap_or("");
return resp.contains("200") || resp.contains("ok");
}
false
}
Err(_) => false,
}
}
#[test]
fn pair_and_transfer_smoke() {
let h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let test_file = h.work_dir.join("test_payload.bin");
let mut payload = vec![
0x0D, 0x0A, ];
for i in 0u8..=255 {
payload.push(i);
}
payload.extend_from_slice(b"\x00\x01\x02\x03\xff\xfe\xfd\xfc");
std::fs::write(&test_file, &payload).expect("write test file");
use std::hash::Hasher;
let expected_hash = {
use sha2::{Digest, Sha256};
let digest = Sha256::digest(&payload);
digest.iter().map(|b| format!("{b:02x}")).collect::<String>()
};
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let direct_flag = std::env::var("FILAMENT_DIRECT_PER_OS").unwrap_or_else(|_| "1".into());
let mut send_proc = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env(
"FILAMENT_DIRECT_LOOPBACK_ONLY",
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY")
.unwrap_or_else(|_| "1".into()),
)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.arg("send")
.arg(&test_file)
.arg("--word")
.arg(CODE_WORD)
.arg("--server")
.arg(&server)
.stderr(Stdio::piped())
.spawn()
.expect("send");
let (code_tx, code_rx) = std::sync::mpsc::channel::<String>();
let stderr = send_proc.stderr.take().unwrap();
std::thread::spawn(move || {
let reader = BufReader::new(stderr);
for line in reader.lines() {
let line = line.unwrap_or_default();
eprintln!("[send] {line}");
let lower = line.to_lowercase();
if let Some(start) = lower.find(&CODE_WORD.to_lowercase()) {
let rest = &line[start..];
let end = rest.find(|c: char| c.is_whitespace()).unwrap_or(rest.len());
let _ = code_tx.send(line[start..start + end].to_lowercase().to_string());
}
}
});
let full_code = code_rx.recv_timeout(Duration::from_secs(30))
.expect("send did not mint a code within 30s");
let recv_dir = h.b_dir.join("received");
std::fs::create_dir_all(&recv_dir).expect("create recv dir");
let mut recv_proc = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env(
"FILAMENT_DIRECT_LOOPBACK_ONLY",
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY")
.unwrap_or_else(|_| "1".into()),
)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.arg("receive")
.arg(&full_code)
.arg("--yes")
.arg("--dir")
.arg(&recv_dir)
.arg("--server")
.arg(&server)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("recv");
let recv_out = recv_proc.wait_with_output().expect("recv result");
let recv_stdout = String::from_utf8_lossy(&recv_out.stdout);
let recv_stderr = String::from_utf8_lossy(&recv_out.stderr);
eprintln!("recv stdout:\n{recv_stdout}");
eprintln!("recv stderr:\n{recv_stderr}");
let received_file = recv_dir.join("test_payload.bin");
assert!(
received_file.exists(),
"received file missing: {}",
received_file.display()
);
let received_data = std::fs::read(&received_file).expect("read received file");
let recv_hash = {
use sha2::{Digest, Sha256};
let digest = Sha256::digest(&received_data);
digest.iter().map(|b| format!("{b:02x}")).collect::<String>()
};
assert_eq!(
expected_hash, recv_hash,
"sha256 mismatch: sent {expected_hash} != received {recv_hash}"
);
assert_eq!(payload.len(), received_data.len(), "size mismatch");
assert_eq!(
received_data[0], 0x0D,
"byte 0 = 0x0D lost (CR)"
);
assert_eq!(
received_data[1], 0x0A,
"byte 1 = 0x0A lost (LF)"
);
}
#[test]
fn revoked_device_first_transfer_is_denied() {
let h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let direct_flag =
std::env::var("FILAMENT_DIRECT_PER_OS").unwrap_or_else(|_| "1".into());
let loopback_only =
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY").unwrap_or_else(|_| "1".into());
let base_env = [
("FILAMENT_CAP_AUTHORITATIVE", "0"),
("FILAMENT_DIRECT", &direct_flag),
("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only),
("FILAMENT_L3_USERSPACE", "1"),
];
let a_cert_path = h.a_dir.join("identity").join("device-cert.json");
let a_cert_raw = std::fs::read_to_string(&a_cert_path).expect("a device cert");
let a_cert_val: Value =
serde_json::from_str(&a_cert_raw).expect("parse a device cert");
let a_cert = &a_cert_val["cert"];
assert!(
a_cert["devicePub"].as_str().is_some(),
"A's device cert has a devicePub"
);
let a_record = serde_json::json!([{
"name": "test-a",
"secret": "abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789",
"deviceCert": a_cert,
"caps": ["transfer"],
}]);
std::fs::write(
h.b_dir.join("devices.json"),
serde_json::to_string_pretty(&a_record).unwrap(),
)
.expect("write b devices.json");
let revoke = Command::new(&bin)
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.args(["devices", "revoke", "test-a", "--yes"])
.output()
.expect("revoke");
assert!(
revoke.status.success(),
"devices revoke failed: {}",
String::from_utf8_lossy(&revoke.stderr)
);
let test_file = h.work_dir.join("revoked-payload.bin");
std::fs::write(&test_file, b"secret bytes that must never land").unwrap();
let mut send_proc = Command::new(&bin)
.envs(base_env.iter().map(|(k, v)| (*k, *v)))
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.arg("send")
.arg(&test_file)
.arg("--word")
.arg("revoked-transfer-code")
.arg("--server")
.arg(&server)
.stderr(Stdio::piped())
.spawn()
.expect("send");
let (code_tx, code_rx) = std::sync::mpsc::channel::<String>();
let stderr = send_proc.stderr.take().unwrap();
let code_word = "revoked-transfer-code".to_string();
std::thread::spawn(move || {
let reader = BufReader::new(stderr);
for line in reader.lines() {
let line = line.unwrap_or_default();
let lower = line.to_lowercase();
if let Some(start) = lower.find(&code_word.to_lowercase()) {
let rest = &line[start..];
let end = rest
.find(|c: char| c.is_whitespace())
.unwrap_or(rest.len());
let _ = code_tx.send(line[start..start + end].to_lowercase().to_string());
}
}
});
let full_code = code_rx
.recv_timeout(Duration::from_secs(30))
.expect("send did not mint a code within 30s");
let recv_dir = h.b_dir.join("received");
std::fs::create_dir_all(&recv_dir).unwrap();
let mut recv_proc = Command::new(&bin)
.envs(base_env.iter().map(|(k, v)| (*k, *v)))
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.arg("receive")
.arg(&full_code)
.arg("--yes")
.arg("--dir")
.arg(&recv_dir)
.arg("--server")
.arg(&server)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("receive");
let recv_out = recv_proc.wait_with_output().expect("recv result");
let recv_stderr = String::from_utf8_lossy(&recv_out.stderr);
let _ = send_proc.wait_with_output().expect("send result");
assert!(
!recv_dir.join("revoked-payload.bin").exists(),
"a revoked device's first transfer must be denied; file landed"
);
assert!(
recv_stderr.contains("declined") || recv_stderr.contains("revoked"),
"expected a revocation decline in receiver stderr, got: {recv_stderr}"
);
}
#[test]
fn direct_blocked_falls_back_to_webrtc_promptly() {
let h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let test_file = h.work_dir.join("fallback_payload.bin");
let payload: Vec<u8> = (0u16..4096).map(|i| (i % 251) as u8).collect();
std::fs::write(&test_file, &payload).expect("write test file");
let expected_hash = {
use sha2::{Digest, Sha256};
Sha256::digest(&payload).iter().map(|b| format!("{b:02x}")).collect::<String>()
};
let loopback = std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY").unwrap_or_else(|_| "1".into());
let mut send = Command::new(&bin);
send.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", "1")
.env("FILAMENT_DIRECT_TEST_BLOCK", "1")
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.arg("send").arg(&test_file)
.arg("--word").arg(CODE_WORD)
.arg("--server").arg(&server);
let send_proc = spawn_captured(send).expect("send");
let code_started = Instant::now();
let full_code = loop {
let (stdout, stderr) = send_proc.snapshot();
let text = format!("{stdout}\n{stderr}");
let lower = text.to_lowercase();
if let Some(start) = lower.find(&CODE_WORD.to_lowercase()) {
let rest = &text[start..];
let end = rest.find(char::is_whitespace).unwrap_or(rest.len());
break rest[..end].to_lowercase();
}
assert!(code_started.elapsed() < Duration::from_secs(30), "send did not mint a code within 30s");
std::thread::sleep(Duration::from_millis(20));
};
let recv_dir = h.b_dir.join("received_fallback");
std::fs::create_dir_all(&recv_dir).expect("create recv dir");
let mut recv = Command::new(&bin);
recv.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", "1")
.env("FILAMENT_DIRECT_TEST_BLOCK", "1")
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.arg("receive").arg(&full_code)
.arg("--yes")
.arg("--dir").arg(&recv_dir)
.arg("--server").arg(&server);
let recv_out = spawn_captured(recv)
.expect("recv spawn")
.wait_until(Duration::from_secs(60));
let send_out = send_proc.wait_until(Duration::from_secs(10));
let recv_all = format!(
"{}{}",
recv_out.stdout,
recv_out.stderr
);
eprintln!("recv:\n{recv_all}");
let send_all = format!("{}{}", send_out.stdout, send_out.stderr);
let both = format!("{send_all}\n{recv_all}");
eprintln!("send:\n{send_all}");
let received_file = recv_dir.join("fallback_payload.bin");
assert!(received_file.exists(), "received file missing: {}", received_file.display());
let got = std::fs::read(&received_file).expect("read received file");
let got_hash = {
use sha2::{Digest, Sha256};
Sha256::digest(&got).iter().map(|b| format!("{b:02x}")).collect::<String>()
};
assert_eq!(expected_hash, got_hash, "sha256 mismatch after fallback");
assert!(
both.contains("DIRECT-BLOCKED"),
"DIRECT-BLOCKED marker absent: the block never engaged, so this run is \
UNCLASSIFIED rather than a fallback finding"
);
assert!(
!both.contains("DIRECT-CONNECT ok"),
"direct connected despite the block; this did not test the fallback"
);
if BUILD_PROFILE != "release" {
panic!(
"UNCLASSIFIED: fallback completed functionally, but latency coverage requires build profile release; current profile is {}",
BUILD_PROFILE,
);
}
if both.contains("DIRECT-FALLBACK") {
let elapsed = [
marker_delta(&send_out.events),
marker_delta(&recv_out.events),
]
.into_iter()
.flatten()
.max()
.expect("DIRECT-FALLBACK marker was present without a complete marker transition");
assert!(
elapsed < Duration::from_secs(10),
"fallback marker delta was {elapsed:?}; investigate instead of raising \
this bound: the designed path is ~6s, while roster reconciliation is ~30s"
);
eprintln!(
"PASS via designed DIRECT-FALLBACK in {elapsed:?} (build profile {BUILD_PROFILE})"
);
} else {
eprintln!("PASS via link re-establishment before DIRECT-FALLBACK expiry");
}
}
#[test]
fn two_nodes_pair_each_other() {
let h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let out = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.args(["--server", &server, "--help"])
.output()
.expect("help");
assert!(out.status.success(), "binary help failed");
eprintln!("two_nodes_pair_each_other: filament binary and backend OK");
}
#[test]
fn pty_one_shot_exec_smoke() {
#[cfg(any(windows, target_os = "macos"))]
{
eprintln!("pty_one_shot_exec_smoke: skipped on {os} (cold establish not yet verified on this platform)",
os = if cfg!(windows) { "Windows" } else { "macOS" });
return;
}
let mut h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let direct_flag =
std::env::var("FILAMENT_DIRECT_PER_OS").unwrap_or_else(|_| "1".into());
let loopback_only =
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY").unwrap_or_else(|_| "1".into());
h.pair_daemons();
eprintln!("pty_one_shot_exec_smoke: daemons started");
let pair_word = format!("pairtest-mesh-p{:x}", std::process::id());
let mut create = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-a")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "add", "--word", &pair_word])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair create");
let stderr = create.stderr.take().unwrap();
let pair_word_lower = pair_word.to_lowercase();
let (code_tx, code_rx) = std::sync::mpsc::channel::<String>();
std::thread::spawn(move || {
let reader = BufReader::new(stderr);
for line in reader.lines() {
let line = line.unwrap_or_default();
eprintln!("[pair-create] {line}");
if line.to_lowercase().contains(&pair_word_lower) {
if let Some(code) = line
.split_whitespace()
.find(|w| {
w.to_lowercase().contains(&pair_word_lower)
&& w.split('-').count() >= 4
})
{
let _ = code_tx.send(code.to_string());
}
}
}
});
let pair_code = code_rx.recv_timeout(Duration::from_secs(60))
.expect("pair create did not mint a code within 60s");
eprintln!("pair code: {pair_code}");
let mut claim = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-b")
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.args(["--server", &server, "add", &pair_code])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair claim");
let claim_out = claim.wait_with_output().expect("pair claim result");
eprintln!("pair-claim stdout: {}", String::from_utf8_lossy(&claim_out.stdout));
eprintln!("pair-claim stderr: {}", String::from_utf8_lossy(&claim_out.stderr));
let create_out = create.wait_with_output().expect("pair create result");
eprintln!("pair-create exit: {}", create_out.status);
#[cfg(target_os = "macos")]
{
eprintln!("pty_one_shot_exec_smoke: restarting daemons (macOS)");
if let Some(ref mut c) = h.daemon_a {
let _ = c.kill();
let _ = c.wait();
}
if let Some(ref mut c) = h.daemon_b {
let _ = c.kill();
let _ = c.wait();
}
h.daemon_a = None;
h.daemon_b = None;
std::thread::sleep(Duration::from_secs(3));
let (child_a, log_a) = spawn_daemon_inner(&bin, &server, "test-a", &h.a_dir);
let (child_b, log_b) = spawn_daemon_inner(&bin, &server, "test-b", &h.b_dir);
h.daemon_a = Some(child_a);
h.daemon_b = Some(child_b);
h.daemon_a_log = log_a;
h.daemon_b_log = log_b;
std::thread::sleep(Duration::from_secs(12));
}
#[cfg(not(target_os = "macos"))]
{
eprintln!("pty_one_shot_exec_smoke: waiting for live-pairing discovery...");
let discovered = wait_for_line(
&h.daemon_a_log,
"known device 'test-b' appeared",
Duration::from_secs(30),
) || wait_for_line(
&h.daemon_b_log,
"known device 'test-a' appeared",
Duration::from_secs(0),
);
assert!(
discovered,
"live-pairing discovery did not occur within 30s"
);
}
let nonce = format!("PTY-OK-{}", std::process::id());
let mut pty_proc = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "shell", "test-b", "--", "echo", &nonce])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pty");
let pty_out = pty_proc.wait_with_output().expect("pty result");
let pty_stdout = String::from_utf8_lossy(&pty_out.stdout);
let pty_stderr = String::from_utf8_lossy(&pty_out.stderr);
eprintln!("pty stdout: {pty_stdout}");
eprintln!("pty stderr: {pty_stderr}");
assert!(
pty_stdout.contains(&nonce) || pty_stderr.contains(&nonce),
"pty output does not contain nonce '{nonce}'\nstdout: {pty_stdout}\nstderr: {pty_stderr}"
);
}
#[test]
fn shell_owner_gate_refuses_real_spawn() {
let root = std::env::temp_dir().join(format!("filament-shell-gate-{}", std::process::id()));
let drops = root.join("drops");
let out = Command::new(env!("CARGO_BIN_EXE_filament"))
.env("FILAMENT_CONFIG_DIR", &root)
.args([
"up", "--userspace", "--shell", "--server", "http://127.0.0.1:1", "--dir",
drops.to_str().unwrap(),
])
.output()
.expect("spawn shell gate probe");
let stderr = String::from_utf8_lossy(&out.stderr);
assert!(!out.status.success(), "ungated shell spawn unexpectedly succeeded: {stderr}");
assert!(stderr.contains("owner's authority"), "missing owner-equivalence refusal: {stderr}");
assert!(!stderr.contains("socket.io connect"), "shell gate did not fire before signaling: {stderr}");
let _ = std::fs::remove_dir_all(root);
}
#[test]
fn shell_daemon_live_pairing_no_restart() {
#[cfg(any(windows, target_os = "macos"))]
{
eprintln!("shell_daemon_live_pairing_no_restart: skipped on {os} (cold establish not yet verified on this platform)",
os = if cfg!(windows) { "Windows" } else { "macOS" });
return;
}
let mut h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let direct_flag =
std::env::var("FILAMENT_DIRECT_PER_OS").unwrap_or_else(|_| "1".into());
let loopback_only =
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY").unwrap_or_else(|_| "1".into());
let (child_a, log_a) = spawn_daemon_inner(&bin, &server, "test-a", &h.a_dir);
let (child_b, log_b) = spawn_daemon_inner(&bin, &server, "test-b", &h.b_dir);
h.daemon_a = Some(child_a);
h.daemon_b = Some(child_b);
h.daemon_a_log = log_a;
h.daemon_b_log = log_b;
std::thread::sleep(Duration::from_secs(4));
assert!(
h.daemon_a.as_mut().unwrap().try_wait().unwrap().is_none(),
"shell daemon A exited on empty devices (bail bug)"
);
assert!(
h.daemon_b.as_mut().unwrap().try_wait().unwrap().is_none(),
"shell daemon B exited on empty devices (bail bug)"
);
let pair_word = format!("livescan-pair-p{:x}", std::process::id());
let mut create = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-a")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "add", "--word", &pair_word])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair create");
let stderr = create.stderr.take().unwrap();
let pair_word_lower = pair_word.to_lowercase();
let (code_tx, code_rx) = std::sync::mpsc::channel::<String>();
std::thread::spawn(move || {
let reader = BufReader::new(stderr);
for line in reader.lines() {
let line = line.unwrap_or_default();
eprintln!("[livescan-pair-create] {line}");
if line.to_lowercase().contains(&pair_word_lower) {
if let Some(code) = line
.split_whitespace()
.find(|w| {
w.to_lowercase().contains(&pair_word_lower)
&& w.split('-').count() >= 4
})
{
let _ = code_tx.send(code.to_string());
}
}
}
});
let pair_code = code_rx.recv_timeout(Duration::from_secs(60))
.expect("live-pairing create did not mint a code within 60s");
eprintln!("live-pairing code: {pair_code}");
let mut claim = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-b")
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.args(["--server", &server, "add", &pair_code])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair claim");
let claim_out = claim.wait_with_output().expect("pair claim result");
eprintln!("pair-claim stdout: {}", String::from_utf8_lossy(&claim_out.stdout));
eprintln!("pair-claim stderr: {}", String::from_utf8_lossy(&claim_out.stderr));
let create_out = create.wait_with_output().expect("pair create result");
eprintln!("pair-create exit: {}", create_out.status);
#[cfg(target_os = "macos")]
{
if let Some(ref mut c) = h.daemon_a {
let _ = c.kill();
let _ = c.wait();
}
if let Some(ref mut c) = h.daemon_b {
let _ = c.kill();
let _ = c.wait();
}
h.daemon_a = None;
h.daemon_b = None;
std::thread::sleep(Duration::from_secs(3));
let (child_a, log_a) = spawn_daemon_inner(&bin, &server, "test-a", &h.a_dir);
let (child_b, log_b) = spawn_daemon_inner(&bin, &server, "test-b", &h.b_dir);
h.daemon_a = Some(child_a);
h.daemon_b = Some(child_b);
h.daemon_a_log = log_a;
h.daemon_b_log = log_b;
std::thread::sleep(Duration::from_secs(12));
}
#[cfg(not(target_os = "macos"))]
{
eprintln!("waiting for live-pairing scan to discover the new device...");
std::thread::sleep(Duration::from_secs(15));
}
let nonce = format!("LIVE-PTY-OK-{}", std::process::id());
let mut pty_proc = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "shell", "test-b", "--", "echo", &nonce])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pty");
let pty_out = pty_proc.wait_with_output().expect("pty result");
let pty_stdout = String::from_utf8_lossy(&pty_out.stdout);
let pty_stderr = String::from_utf8_lossy(&pty_out.stderr);
eprintln!("live-pty stdout: {pty_stdout}");
eprintln!("live-pty stderr: {pty_stderr}");
assert!(
pty_stdout.contains(&nonce) || pty_stderr.contains(&nonce),
"live-pairing pty failed — daemon did not discover the newly paired device\n\
nonce: {nonce}\nstdout: {pty_stdout}\nstderr: {pty_stderr}"
);
}
#[cfg(not(target_os = "macos"))]
#[test]
fn warm_all_makes_first_contact_warm() {
#[cfg(windows)]
{
eprintln!("warm_all_makes_first_contact_warm: skipped on Windows");
return;
}
let mut h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let direct_flag =
std::env::var("FILAMENT_DIRECT_PER_OS").unwrap_or_else(|_| "1".into());
let loopback_only =
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY").unwrap_or_else(|_| "1".into());
h.pair_daemons();
eprintln!("warm_all: daemons started");
let pair_word = format!("warmtest-mesh-p{:x}", std::process::id());
let mut create = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-a")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "add", "--word", &pair_word])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair create");
let stderr = create.stderr.take().unwrap();
let pair_word_lower = pair_word.to_lowercase();
let (code_tx, code_rx) = std::sync::mpsc::channel::<String>();
std::thread::spawn(move || {
let reader = BufReader::new(stderr);
for line in reader.lines() {
let line = line.unwrap_or_default();
eprintln!("[pair-create] {line}");
if line.to_lowercase().contains(&pair_word_lower) {
if let Some(code) = line
.split_whitespace()
.find(|w| {
w.to_lowercase().contains(&pair_word_lower)
&& w.split('-').count() >= 4
})
{
let _ = code_tx.send(code.to_string());
}
}
}
});
let pair_code = code_rx.recv_timeout(Duration::from_secs(60))
.expect("pair create did not mint a code within 60s");
eprintln!("pair code: {pair_code}");
let mut claim = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-b")
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.args(["--server", &server, "add", &pair_code])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair claim");
let claim_out = claim.wait_with_output().expect("pair claim result");
eprintln!("pair-claim stdout: {}", String::from_utf8_lossy(&claim_out.stdout));
eprintln!("pair-claim stderr: {}", String::from_utf8_lossy(&claim_out.stderr));
let create_out = create.wait_with_output().expect("pair create result");
eprintln!("pair-create exit: {}", create_out.status);
eprintln!("warm_all: restarting daemons with stock config");
if let Some(ref mut c) = h.daemon_a { let _ = c.kill(); let _ = c.wait(); }
if let Some(ref mut c) = h.daemon_b { let _ = c.kill(); let _ = c.wait(); }
h.daemon_a = None;
h.daemon_b = None;
std::thread::sleep(Duration::from_secs(3));
let (child_a, log_a) = spawn_daemon_inner(&bin, &server, "test-a", &h.a_dir);
let (child_b, log_b) = spawn_daemon_inner(&bin, &server, "test-b", &h.b_dir);
h.daemon_a = Some(child_a);
h.daemon_b = Some(child_b);
h.daemon_a_log = log_a;
h.daemon_b_log = log_b;
eprintln!("warm_all: waiting for warm-hold to establish the link...");
let link_held = wait_for_line(
&h.daemon_a_log,
"warm-hold: established connection to 'test-b'",
Duration::from_secs(60),
) || wait_for_line(
&h.daemon_a_log,
"warm-hold: skip 'test-b'",
Duration::from_secs(0),
);
assert!(
link_held,
"warm-hold did not establish a link to 'test-b' within 60s — \
the auto-warm setting may not be defaulting to ON"
);
wait_for_line(
&h.daemon_b_log,
"known device 'test-a' appeared",
Duration::from_secs(15),
);
std::thread::sleep(Duration::from_secs(10));
eprintln!("warm_all: running filament ping --json...");
let mut ping_proc = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "reach", "test-b", "--json"])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("ping");
let ping_out = ping_proc.wait_with_output().expect("ping result");
let ping_stdout = String::from_utf8_lossy(&ping_out.stdout);
let ping_stderr = String::from_utf8_lossy(&ping_out.stderr);
eprintln!("warm_all ping stdout: {ping_stdout}");
eprintln!("warm_all ping stderr: {ping_stderr}");
let ping_json: serde_json::Value = serde_json::from_str(&ping_stdout)
.expect("ping --json output is not valid JSON");
assert_eq!(ping_json["warm"], true, "first ping should be warm, got: {ping_json}");
assert!(
!ping_stderr.contains("still reaching"),
"ping output contains cold-establish banner 'still reaching'"
);
eprintln!("warm_all: waiting 20s then re-pinging...");
std::thread::sleep(Duration::from_secs(20));
let mut ping_proc2 = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "reach", "test-b", "--json"])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("ping 2");
let ping_out2 = ping_proc2.wait_with_output().expect("ping 2 result");
let ping_stdout2 = String::from_utf8_lossy(&ping_out2.stdout);
eprintln!("warm_all ping2 stdout: {ping_stdout2}");
let ping_json2: serde_json::Value = serde_json::from_str(&ping_stdout2)
.expect("ping 2 --json output is not valid JSON");
assert_eq!(ping_json2["warm"], true, "second ping should still be warm, got: {ping_json2}");
eprintln!("warm_all: testing negative control (auto-warm off)");
if let Some(ref mut c) = h.daemon_a { let _ = c.kill(); let _ = c.wait(); }
h.daemon_a = None;
std::thread::sleep(Duration::from_secs(3));
let mut daemon_a_off = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_LOG", "trace")
.env("FILAMENT_AUTO_WARM", "0")
.env(
"FILAMENT_DIRECT_LOOPBACK_ONLY",
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY")
.unwrap_or_else(|_| "1".into()),
)
.env("FILAMENT_NAME", "test-a")
.arg("up")
.arg("--userspace")
.arg("--shell")
.arg("--i-know")
.arg("--server")
.arg(&server)
.arg("--relay")
.arg("--dir")
.arg(h.a_dir.join("drops").to_str().unwrap())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("spawn daemon a (auto-warm off)");
let log_a_off: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
let log_a_off_clone = log_a_off.clone();
if let Some(stderr) = daemon_a_off.stderr.take() {
std::thread::spawn(move || {
for line in BufReader::new(stderr).lines() {
if let Ok(l) = line {
eprintln!("[daemon-a-off stderr] {l}");
if let Ok(mut log) = log_a_off_clone.lock() {
log.push(l);
}
}
}
});
}
if let Some(stdout) = daemon_a_off.stdout.take() {
std::thread::spawn(move || {
for line in BufReader::new(stdout).lines() {
if let Ok(l) = line { eprintln!("[daemon-a-off] {l}"); }
}
});
}
h.daemon_a = Some(daemon_a_off);
std::thread::sleep(Duration::from_secs(35));
let auto_warm_line = wait_for_line(
&log_a_off,
"warm-hold: auto-warming",
Duration::from_secs(0),
);
assert!(
!auto_warm_line,
"auto-warm OFF: daemon should NOT emit 'auto-warming' line"
);
eprintln!("warm_all: all assertions passed ✓");
}
#[test]
fn warm_one_shot_pty_reuse() {
#[cfg(any(windows, target_os = "macos"))]
{
eprintln!("warm_one_shot_pty_reuse: skipped on {os} (warm-reuse pty not yet verified)",
os = if cfg!(windows) { "Windows" } else { "macOS" });
return;
}
let mut h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let direct_flag =
std::env::var("FILAMENT_DIRECT_PER_OS").unwrap_or_else(|_| "1".into());
let loopback_only =
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY").unwrap_or_else(|_| "1".into());
h.pair_daemons();
eprintln!("warm_one_shot_pty_reuse: daemons started");
let pair_word = format!("warm-reuse-p{:x}", std::process::id());
let mut create = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-a")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "add", "--word", &pair_word])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair create");
let stderr = create.stderr.take().unwrap();
let pair_word_lower = pair_word.to_lowercase();
let (code_tx, code_rx) = std::sync::mpsc::channel::<String>();
std::thread::spawn(move || {
let reader = BufReader::new(stderr);
for line in reader.lines() {
let line = line.unwrap_or_default();
eprintln!("[warm-pair-create] {line}");
if line.to_lowercase().contains(&pair_word_lower) {
if let Some(code) = line
.split_whitespace()
.find(|w| {
w.to_lowercase().contains(&pair_word_lower)
&& w.split('-').count() >= 4
})
{
let _ = code_tx.send(code.to_string());
}
}
}
});
let pair_code = code_rx.recv_timeout(Duration::from_secs(60))
.expect("warm pair create did not mint a code within 60s");
eprintln!("warm pair code: {pair_code}");
let mut claim = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-b")
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.args(["--server", &server, "add", &pair_code])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair claim");
let claim_out = claim.wait_with_output().expect("pair claim result");
eprintln!("warm-claim stdout: {}", String::from_utf8_lossy(&claim_out.stdout));
eprintln!("warm-claim stderr: {}", String::from_utf8_lossy(&claim_out.stderr));
let create_out = create.wait_with_output().expect("pair create result");
eprintln!("warm-pair-create exit: {}", create_out.status);
let set = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "set", "warm-peers", "test-b"])
.output()
.expect("set warm-peers");
assert!(set.status.success(), "set warm-peers failed: {:?}", String::from_utf8_lossy(&set.stderr));
if let Some(ref mut c) = h.daemon_a {
let _ = c.kill();
let _ = c.wait();
}
if let Some(ref mut c) = h.daemon_b {
let _ = c.kill();
let _ = c.wait();
}
h.daemon_a = None;
h.daemon_b = None;
std::thread::sleep(Duration::from_secs(3));
let (child_a, log_a) = spawn_daemon_inner(&bin, &server, "test-a", &h.a_dir);
let (child_b, log_b) = spawn_daemon_inner(&bin, &server, "test-b", &h.b_dir);
h.daemon_a = Some(child_a);
h.daemon_b = Some(child_b);
h.daemon_a_log = log_a;
h.daemon_b_log = log_b;
let warm_ready = wait_for_line(
&h.daemon_a_log,
"warm-hold: established connection to 'test-b'",
Duration::from_secs(40),
) || wait_for_line(
&h.daemon_a_log,
"warm-hold: skip 'test-b'",
Duration::from_secs(0),
);
assert!(
warm_ready,
"warm link to 'test-b' was not established within 40s"
);
std::thread::sleep(Duration::from_secs(5));
let nonce = format!("WARM-OK-{}", std::process::id());
let out = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "shell", "test-b", "--", "echo", &nonce])
.output()
.expect("pty");
let out_stdout = String::from_utf8_lossy(&out.stdout);
let out_stderr = String::from_utf8_lossy(&out.stderr);
eprintln!("warm pty stdout: {out_stdout}");
eprintln!("warm pty stderr: {out_stderr}");
assert!(
out_stderr.contains("reusing warm link") && out_stderr.contains("one-shot pty"),
"pty did NOT use warm link - expected trace 'reusing warm link ... for one-shot pty'\n\
stdout: {out_stdout}\nstderr: {out_stderr}"
);
assert!(
out_stdout.contains(&nonce),
"pty output does not contain nonce '{nonce}'\nstdout: {out_stdout}\nstderr: {out_stderr}"
);
}
#[cfg(test)]
mod captured_child_tests {
use super::*;
fn command_that_prints_then_waits() -> Command {
if cfg!(windows) {
let mut command = Command::new("cmd");
command.args(["/C", "echo live && ping -n 4 127.0.0.1 >NUL"]);
command
} else {
let mut command = Command::new("sh");
command.args(["-c", "printf 'live\\n'; sleep 3"]);
command
}
}
fn hanging_command() -> Command {
if cfg!(windows) {
let mut command = Command::new("ping");
command.args(["-n", "20", "127.0.0.1"]);
command
} else {
let mut command = Command::new("sleep");
command.arg("20");
command
}
}
fn failing_command() -> Command {
if cfg!(windows) {
let mut command = Command::new("cmd");
command.args(["/C", "exit 7"]);
command
} else {
let mut command = Command::new("sh");
command.args(["-c", "exit 7"]);
command
}
}
#[test]
fn captured_child_exposes_output_before_exit() {
let child = spawn_captured(command_that_prints_then_waits()).expect("spawn");
let started = std::time::Instant::now();
let mut observed = false;
while started.elapsed() < Duration::from_secs(2) {
let (stdout, stderr) = child.snapshot();
if stdout.contains("live") || stderr.contains("live") {
observed = true;
break;
}
std::thread::sleep(Duration::from_millis(20));
}
assert!(observed, "child output was not readable before child exit");
let result = child.wait_until(Duration::from_secs(5));
assert!(matches!(result.outcome, ChildOutcome::ExitedSuccess(_)));
}
#[test]
fn captured_child_timeout_is_distinct_from_failure() {
let started = std::time::Instant::now();
let result = run_captured(hanging_command(), Duration::from_millis(150));
assert!(matches!(result.outcome, ChildOutcome::TimedOut));
assert!(started.elapsed() < Duration::from_secs(2), "timeout did not fire promptly");
}
#[test]
fn captured_child_reports_nonzero_exit_separately() {
let result = run_captured(failing_command(), Duration::from_secs(2));
assert!(matches!(result.outcome, ChildOutcome::ExitedFailure(_)));
}
#[test]
fn captured_child_reports_spawn_failure() {
let mut command = Command::new("filament-test-command-that-does-not-exist");
command.arg("--version");
let result = run_captured(command, Duration::from_millis(100));
assert!(matches!(result.outcome, ChildOutcome::SpawnFailed(_)));
}
}
#[test]
#[ignore = "reports #31 ladder exhaustion after recovery attempts; enable with cargo test --features test-hooks -- --ignored after recovery is fixed"]
fn freeze_stall_detector_classification() {
let h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let payload_path = h.work_dir.join("stall_payload.bin");
let payload: Vec<u8> = (0..4_000_000).map(|i| (i % 251) as u8).collect();
std::fs::write(&payload_path, &payload).expect("write stall payload");
let expected_hash = {
use sha2::{Digest, Sha256};
Sha256::digest(&payload).iter().map(|b| format!("{b:02x}")).collect::<String>()
};
let loopback = std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY").unwrap_or_else(|_| "1".into());
let mut send = Command::new(&bin);
send.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", "1")
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback)
.env("FILAMENT_LOG", "debug")
.env("FILAMENT_STALL_MS", "2500")
.env("FILAMENT_WARM_STANDBY", "0")
.env("FILAMENT_TEST_FREEZE_AFTER_BYTES", "700000")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.arg("send").arg(&payload_path).arg("--word").arg(CODE_WORD)
.arg("--server").arg(&server);
let send = spawn_captured(send).expect("spawn measured sender");
let started = std::time::Instant::now();
let code = loop {
let (stdout, stderr) = send.snapshot();
let text = format!("{stdout}\n{stderr}");
if let Some(start) = text.to_lowercase().find(&CODE_WORD.to_lowercase()) {
let rest = &text[start..];
let end = rest.find(char::is_whitespace).unwrap_or(rest.len());
break rest[..end].to_string();
}
assert!(started.elapsed() < Duration::from_secs(30), "sender did not mint a code");
std::thread::sleep(Duration::from_millis(20));
};
let recv_dir = h.b_dir.join("stall_received");
std::fs::create_dir_all(&recv_dir).expect("create receive directory");
let mut recv = Command::new(&bin);
recv.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", "1")
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback)
.env("FILAMENT_LOG", "debug")
.env("FILAMENT_STALL_MS", "2500")
.env("FILAMENT_WARM_STANDBY", "0")
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.arg("receive").arg(&code).arg("--yes").arg("--dir").arg(&recv_dir)
.arg("--server").arg(&server);
let recv = spawn_captured(recv).expect("spawn measured receiver");
let deadline_secs = std::env::var("FILAMENT_STALL_MEASUREMENT_DEADLINE_SECS")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.filter(|&v| v > 0)
.unwrap_or(90);
let recovery_deadline = Duration::from_secs(deadline_secs);
let recv_result = recv.wait_until(recovery_deadline);
let send_result = send.wait_until(Duration::from_secs(10));
let logs = format!("{}\n{}\n{}\n{}", send_result.stdout, send_result.stderr,
recv_result.stdout, recv_result.stderr);
let armed = logs.contains("STALL_WATCHDOG_ARMED idle_ms_tracked");
let froze = logs.contains("data-path FREEZE engaged");
let detected = logs.contains("stall detected") || logs.contains("inbound stall");
let recovery_started = logs.contains("repairing the link")
|| logs.contains("re-establish")
|| logs.contains("re-dial")
|| logs.contains("stall correction");
let received = recv_dir.join("stall_payload.bin");
let recovered = matches!(recv_result.outcome, ChildOutcome::ExitedSuccess(_))
&& received.exists()
&& std::fs::read(&received).ok().map(|data| {
use sha2::{Digest, Sha256};
Sha256::digest(data).iter().map(|b| format!("{b:02x}")).collect::<String>()
}).as_deref() == Some(expected_hash.as_str());
let dump_logs = || {
eprintln!("--- measured sender stdout ---\n{}", send_result.stdout);
eprintln!("--- measured sender stderr ---\n{}", send_result.stderr);
eprintln!("--- measured receiver stdout ---\n{}", recv_result.stdout);
eprintln!("--- measured receiver stderr ---\n{}", recv_result.stderr);
};
if !armed {
dump_logs();
panic!("UNCLASSIFIED: stall watchdog armed marker absent; instrument presence was not proven");
}
if !froze {
dump_logs();
panic!("UNCLASSIFIED: freeze hook did not engage; no stall was injected");
}
if !detected {
dump_logs();
panic!("FAIL: freeze engaged but no stall detector event was observed");
}
if !recovered {
dump_logs();
if recovery_started {
if !matches!(recv_result.outcome, ChildOutcome::TimedOut) {
panic!("DETECTED_NOT_RECOVERED: recovery started, receiver exited before the {recovery_deadline:?} deadline without byte-exact completion");
}
panic!("DETECTED_NOT_RECOVERED_WITHIN_DEADLINE: recovery started but did not complete within {recovery_deadline:?}");
}
panic!("DETECTED_RECOVERY_UNOBSERVED: detector fired but no recovery marker was observed");
}
eprintln!("PASS: freeze injected, detector fired, and transfer recovered byte-exact");
}
#[test]
fn warm_one_shot_pty_instant_eof() {
#[cfg(any(windows, target_os = "macos"))]
{
eprintln!("warm_one_shot_pty_instant_eof: skipped on {os}",
os = if cfg!(windows) { "Windows" } else { "macOS" });
return;
}
let mut h = Harness::new();
let bin = h.filament_bin().to_path_buf();
let server = h.server_url();
let direct_flag =
std::env::var("FILAMENT_DIRECT_PER_OS").unwrap_or_else(|_| "1".into());
let loopback_only =
std::env::var("FILAMENT_DIRECT_LOOPBACK_ONLY").unwrap_or_else(|_| "1".into());
h.pair_daemons();
eprintln!("warm_one_shot_pty_instant_eof: daemons started");
let pair_word = format!("warm-eof-p{:x}", std::process::id());
let mut create = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-a")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "add", "--word", &pair_word])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair create");
let stderr = create.stderr.take().unwrap();
let pair_word_lower = pair_word.to_lowercase();
let (code_tx, code_rx) = std::sync::mpsc::channel::<String>();
std::thread::spawn(move || {
let reader = BufReader::new(stderr);
for line in reader.lines() {
let line = line.unwrap_or_default();
eprintln!("[warm-eof-pair-create] {line}");
if line.to_lowercase().contains(&pair_word_lower) {
if let Some(code) = line
.split_whitespace()
.find(|w| {
w.to_lowercase().contains(&pair_word_lower)
&& w.split('-').count() >= 4
})
{
let _ = code_tx.send(code.to_string());
}
}
}
});
let pair_code = code_rx.recv_timeout(Duration::from_secs(60))
.expect("warm-eof pair create did not mint a code within 60s");
eprintln!("warm-eof pair code: {pair_code}");
let mut claim = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_NAME", "test-b")
.env("FILAMENT_CONFIG_DIR", &h.b_dir)
.args(["--server", &server, "add", &pair_code])
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("pair claim");
let claim_out = claim.wait_with_output().expect("pair claim result");
eprintln!("warm-eof-claim stderr: {}", String::from_utf8_lossy(&claim_out.stderr));
let create_out = create.wait_with_output().expect("pair create result");
eprintln!("warm-eof-pair-create exit: {}", create_out.status);
let set = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "set", "warm-peers", "test-b"])
.output()
.expect("set warm-peers");
assert!(set.status.success(), "set warm-peers failed: {:?}", String::from_utf8_lossy(&set.stderr));
if let Some(ref mut c) = h.daemon_a {
let _ = c.kill();
let _ = c.wait();
}
if let Some(ref mut c) = h.daemon_b {
let _ = c.kill();
let _ = c.wait();
}
h.daemon_a = None;
h.daemon_b = None;
std::thread::sleep(Duration::from_secs(3));
let (child_a, log_a) = spawn_daemon_inner(&bin, &server, "test-a", &h.a_dir);
let (child_b, log_b) = spawn_daemon_inner(&bin, &server, "test-b", &h.b_dir);
h.daemon_a = Some(child_a);
h.daemon_b = Some(child_b);
h.daemon_a_log = log_a;
h.daemon_b_log = log_b;
let warm_ready = wait_for_line(
&h.daemon_a_log,
"warm-hold: established connection to 'test-b'",
Duration::from_secs(40),
) || wait_for_line(
&h.daemon_a_log,
"warm-hold: skip 'test-b'",
Duration::from_secs(0),
);
assert!(
warm_ready,
"warm link to 'test-b' was not established within 40s after daemon restart"
);
std::thread::sleep(Duration::from_secs(5));
let nonce = format!("WARM-EOF-OK-{}", std::process::id());
let out = Command::new(&bin)
.env("FILAMENT_CAP_AUTHORITATIVE", "0")
.env("FILAMENT_DIRECT", &direct_flag)
.env("FILAMENT_DIRECT_LOOPBACK_ONLY", &loopback_only)
.env("FILAMENT_L3_USERSPACE", "1")
.env("FILAMENT_CONFIG_DIR", &h.a_dir)
.args(["--server", &server, "shell", "test-b", "--", "printf", &nonce])
.stdin(Stdio::null()) .output()
.expect("pty instant-eof");
let out_stdout = String::from_utf8_lossy(&out.stdout);
let out_stderr = String::from_utf8_lossy(&out.stderr);
eprintln!("warm pty instant-eof stdout: {out_stdout}");
eprintln!("warm pty instant-eof stderr: {out_stderr}");
assert!(
out.status.success(),
"warm one-shot pty with instant stdin-EOF failed: rc={:?}\n\
stdout: {out_stdout}\nstderr: {out_stderr}",
out.status.code()
);
assert!(
out_stderr.contains("reusing warm link") && out_stderr.contains("one-shot pty"),
"pty did NOT use warm link - expected 'reusing warm link ... for one-shot pty'\n\
stdout: {out_stdout}\nstderr: {out_stderr}"
);
assert!(
out_stdout.contains(&nonce),
"pty output does not contain nonce '{nonce}'\nstdout: {out_stdout}\nstderr: {out_stderr}"
);
}