use std::io::Write;
use std::process::{Command, Stdio};
use std::sync::Arc;
use futures_util::{SinkExt, StreamExt};
use super::{run_compile_session, serve_session, session_framed};
use crate::broker::protocol_v2::{session_frame, SessionFrame, SessionStart};
use crate::containment::ContainedProcessGroup;
struct TestEnvVarGuard {
key: &'static str,
previous: Option<std::ffi::OsString>,
}
impl TestEnvVarGuard {
fn set(key: &'static str, value: &str) -> Self {
let guard = Self {
key,
previous: std::env::var_os(key),
};
std::env::set_var(key, value);
guard
}
}
impl Drop for TestEnvVarGuard {
fn drop(&mut self) {
match self.previous.take() {
Some(value) => std::env::set_var(self.key, value),
None => std::env::remove_var(self.key),
}
}
}
fn fixture_program() -> String {
let exe = std::env::current_exe().expect("test executable path");
let dir = exe
.parent()
.and_then(std::path::Path::parent)
.expect("test binary should live in <profile>/deps/");
dir.join(format!(
"testbin-stdio-scripted{}",
std::env::consts::EXE_SUFFIX
))
.to_string_lossy()
.into_owned()
}
fn fixture() -> Command {
let exe = std::env::current_exe().expect("test executable path");
let dir = exe
.parent()
.and_then(std::path::Path::parent)
.expect("test binary should live in <profile>/deps/");
let path = dir.join(format!(
"testbin-stdio-scripted{}",
std::env::consts::EXE_SUFFIX
));
assert!(
path.is_file(),
"test fixture is missing at {} — run `soldr cargo build -p testbins` first",
path.display()
);
Command::new(path)
}
fn frame(kind: session_frame::Kind) -> SessionFrame {
SessionFrame { kind: Some(kind) }
}
#[derive(Debug, PartialEq, Eq)]
struct Reassembled {
stdout: Vec<u8>,
stderr: Vec<u8>,
code: i32,
}
fn run_direct(directives: &[&str], stdin: &[u8]) -> Reassembled {
let mut child = fixture()
.args(directives)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("spawn oracle");
child.stdin.take().unwrap().write_all(stdin).unwrap();
let out = child.wait_with_output().unwrap();
Reassembled {
stdout: out.stdout,
stderr: out.stderr,
code: out.status.code().unwrap_or(-1),
}
}
async fn run_over_daemon_transport(directives: &[&str], stdin: &[u8]) -> Reassembled {
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let server = session_framed(server_io);
let mut client = session_framed(client_io);
let mut cmd = fixture();
cmd.args(directives);
let group = Arc::new(ContainedProcessGroup::new().expect("contained group"));
let handler = tokio::spawn(run_compile_session(server, cmd, group));
if !stdin.is_empty() {
client
.send(frame(session_frame::Kind::Stdin(stdin.to_vec())))
.await
.unwrap();
}
client
.send(frame(session_frame::Kind::StdinEof(true)))
.await
.unwrap();
let mut stdout = Vec::new();
let mut stderr = Vec::new();
let mut code = None;
while let Some(Ok(sf)) = client.next().await {
match sf.kind {
Some(session_frame::Kind::Stdout(b)) => stdout.extend_from_slice(&b),
Some(session_frame::Kind::Stderr(b)) => stderr.extend_from_slice(&b),
Some(session_frame::Kind::Exit(e)) => {
code = Some(e.code);
break;
}
_ => panic!("unexpected inbound-only frame on the outbound lane"),
}
}
let exit = handler.await.unwrap().expect("handler reaps child");
assert_eq!(code, Some(exit.code), "client-visible exit matches handler");
Reassembled {
stdout,
stderr,
code: exit.code,
}
}
#[tokio::test]
async fn compile_session_matches_oracle_over_daemon_transport() {
let script = &["out:ARTIFACT", "err:DIAGNOSTIC", "out:MORE", "exit:5"];
assert_eq!(
run_over_daemon_transport(script, b"").await,
run_direct(script, b"")
);
}
#[tokio::test]
async fn compile_session_is_byte_transparent_for_raw_non_utf8() {
let script = &["outhex:00ff80", "errhex:8081ff"];
let got = run_over_daemon_transport(script, b"").await;
assert_eq!(got.stdout, vec![0x00, 0xff, 0x80]);
assert_eq!(got.stderr, vec![0x80, 0x81, 0xff]);
assert_eq!(got, run_direct(script, b""));
}
#[tokio::test]
async fn compile_session_delivers_full_byte_range_stdin() {
let payload: Vec<u8> = (0u8..=255).collect();
let got = run_over_daemon_transport(&["echo"], &payload).await;
assert_eq!(got.stdout, payload);
assert_eq!(got, run_direct(&["echo"], &payload));
}
#[cfg(unix)]
#[tokio::test]
async fn compile_session_matches_oracle_over_real_unix_socket() {
use tokio::net::{UnixListener, UnixStream};
let sock = std::env::temp_dir().join(format!("rp-sess-{}.sock", std::process::id()));
let _ = std::fs::remove_file(&sock);
let listener = UnixListener::bind(&sock).expect("bind unix listener");
let script: &[&str] = &["out:ARTIFACT", "err:DIAGNOSTIC", "out:MORE", "exit:5"];
let mut cmd = fixture();
cmd.args(script);
let group = Arc::new(ContainedProcessGroup::new().expect("contained group"));
let server = tokio::spawn(async move {
let (io, _addr) = listener.accept().await.expect("accept");
run_compile_session(session_framed(io), cmd, group).await
});
let client_io = UnixStream::connect(&sock)
.await
.expect("connect unix socket");
let mut client = session_framed(client_io);
client
.send(frame(session_frame::Kind::StdinEof(true)))
.await
.expect("send stdin eof");
let mut stdout = Vec::new();
let mut stderr = Vec::new();
let mut code = None;
while let Some(Ok(sf)) = client.next().await {
match sf.kind {
Some(session_frame::Kind::Stdout(b)) => stdout.extend_from_slice(&b),
Some(session_frame::Kind::Stderr(b)) => stderr.extend_from_slice(&b),
Some(session_frame::Kind::Exit(e)) => {
code = Some(e.code);
break;
}
_ => panic!("unexpected inbound-only frame on the outbound lane"),
}
}
let exit = server
.await
.expect("server task")
.expect("handler reaps child");
let _ = std::fs::remove_file(&sock);
let oracle = run_direct(script, b"");
assert_eq!(code, Some(exit.code), "client-visible exit matches handler");
assert_eq!(
stdout, oracle.stdout,
"stdout byte-exact across a real socket"
);
assert_eq!(
stderr, oracle.stderr,
"stderr byte-exact across a real socket"
);
assert_eq!(
exit.code, oracle.code,
"exit code fidelity across a real socket"
);
}
#[cfg(unix)]
#[tokio::test]
async fn compile_session_matches_oracle_across_broker_relay() {
use tokio::net::{UnixListener, UnixStream};
let pid = std::process::id();
let daemon_sock = std::env::temp_dir().join(format!("rp-relay-d-{pid}.sock"));
let broker_sock = std::env::temp_dir().join(format!("rp-relay-b-{pid}.sock"));
let _ = std::fs::remove_file(&daemon_sock);
let _ = std::fs::remove_file(&broker_sock);
let daemon_listener = UnixListener::bind(&daemon_sock).expect("bind daemon listener");
let broker_listener = UnixListener::bind(&broker_sock).expect("bind broker listener");
let script: &[&str] = &["out:ARTIFACT", "err:DIAGNOSTIC", "out:MORE", "exit:5"];
let mut cmd = fixture();
cmd.args(script);
let group = Arc::new(ContainedProcessGroup::new().expect("contained group"));
let daemon = tokio::spawn(async move {
let (io, _addr) = daemon_listener.accept().await.expect("daemon accept");
run_compile_session(session_framed(io), cmd, group).await
});
let daemon_sock_for_broker = daemon_sock.clone();
let broker = tokio::spawn(async move {
let (mut client_conn, _addr) = broker_listener.accept().await.expect("broker accept");
let mut daemon_conn = UnixStream::connect(&daemon_sock_for_broker)
.await
.expect("broker dials daemon");
let _ = tokio::io::copy_bidirectional(&mut client_conn, &mut daemon_conn).await;
});
let client_io = UnixStream::connect(&broker_sock)
.await
.expect("client dials broker");
let mut client = session_framed(client_io);
client
.send(frame(session_frame::Kind::StdinEof(true)))
.await
.expect("send stdin eof");
let mut stdout = Vec::new();
let mut stderr = Vec::new();
let mut code = None;
while let Some(Ok(sf)) = client.next().await {
match sf.kind {
Some(session_frame::Kind::Stdout(b)) => stdout.extend_from_slice(&b),
Some(session_frame::Kind::Stderr(b)) => stderr.extend_from_slice(&b),
Some(session_frame::Kind::Exit(e)) => {
code = Some(e.code);
break;
}
_ => panic!("unexpected inbound-only frame on the outbound lane"),
}
}
drop(client);
let exit = daemon
.await
.expect("daemon task")
.expect("handler reaps child");
let _ = broker.await;
let _ = std::fs::remove_file(&daemon_sock);
let _ = std::fs::remove_file(&broker_sock);
let oracle = run_direct(script, b"");
assert_eq!(code, Some(exit.code), "client-visible exit matches handler");
assert_eq!(stdout, oracle.stdout, "stdout byte-exact across the relay");
assert_eq!(stderr, oracle.stderr, "stderr byte-exact across the relay");
assert_eq!(
exit.code, oracle.code,
"exit code fidelity across the relay"
);
}
#[cfg(unix)]
#[tokio::test]
async fn broker_kill_mid_session_detected_by_client_within_bound() {
use std::time::{Duration, Instant};
use tokio::net::{UnixListener, UnixStream};
let pid = std::process::id();
let daemon_sock = std::env::temp_dir().join(format!("rp-kill-d-{pid}.sock"));
let broker_sock = std::env::temp_dir().join(format!("rp-kill-b-{pid}.sock"));
let _ = std::fs::remove_file(&daemon_sock);
let _ = std::fs::remove_file(&broker_sock);
let daemon_listener = UnixListener::bind(&daemon_sock).expect("bind daemon listener");
let broker_listener = UnixListener::bind(&broker_sock).expect("bind broker listener");
let mut cmd = fixture();
cmd.args(["out:ping", "echo"]);
let group = Arc::new(ContainedProcessGroup::new().expect("contained group"));
let daemon = tokio::spawn(async move {
let (io, _addr) = daemon_listener.accept().await.expect("daemon accept");
run_compile_session(session_framed(io), cmd, group).await
});
let daemon_sock_for_broker = daemon_sock.clone();
let broker = tokio::spawn(async move {
let (mut client_conn, _addr) = broker_listener.accept().await.expect("broker accept");
let mut daemon_conn = UnixStream::connect(&daemon_sock_for_broker)
.await
.expect("broker dials daemon");
let _ = tokio::io::copy_bidirectional(&mut client_conn, &mut daemon_conn).await;
});
let client_io = UnixStream::connect(&broker_sock)
.await
.expect("client dials broker");
let mut client = session_framed(client_io);
let live = tokio::time::timeout(Duration::from_secs(5), client.next())
.await
.expect("first output within 5s")
.expect("a frame")
.expect("decode ok");
assert_eq!(
live.kind,
Some(session_frame::Kind::Stdout(b"ping".to_vec())),
"session is live over the relay before the kill"
);
let t0 = Instant::now();
broker.abort();
let detected = tokio::time::timeout(Duration::from_secs(2), async {
loop {
match client.next().await {
None | Some(Err(_)) => break,
Some(Ok(_)) => continue,
}
}
})
.await;
let latency = t0.elapsed();
assert!(
detected.is_ok(),
"client hung after broker death instead of detecting it"
);
assert!(
latency < Duration::from_secs(1),
"broker-death detection latency {latency:?} exceeds the 1s bound"
);
let _ = tokio::time::timeout(Duration::from_secs(5), daemon)
.await
.expect("daemon session terminates after its connection drops");
let _ = std::fs::remove_file(&daemon_sock);
let _ = std::fs::remove_file(&broker_sock);
}
#[cfg(unix)]
#[tokio::test]
async fn client_kill_mid_session_torn_down_by_daemon_within_bound() {
use std::time::{Duration, Instant};
use tokio::net::{UnixListener, UnixStream};
let pid = std::process::id();
let daemon_sock = std::env::temp_dir().join(format!("rp-ckill-d-{pid}.sock"));
let broker_sock = std::env::temp_dir().join(format!("rp-ckill-b-{pid}.sock"));
let _ = std::fs::remove_file(&daemon_sock);
let _ = std::fs::remove_file(&broker_sock);
let daemon_listener = UnixListener::bind(&daemon_sock).expect("bind daemon listener");
let broker_listener = UnixListener::bind(&broker_sock).expect("bind broker listener");
let mut cmd = fixture();
cmd.args(["out:ping", "echo"]);
let group = Arc::new(ContainedProcessGroup::new().expect("contained group"));
let daemon = tokio::spawn(async move {
let (io, _addr) = daemon_listener.accept().await.expect("daemon accept");
run_compile_session(session_framed(io), cmd, group).await
});
let daemon_sock_for_broker = daemon_sock.clone();
let broker = tokio::spawn(async move {
let (mut client_conn, _addr) = broker_listener.accept().await.expect("broker accept");
let mut daemon_conn = UnixStream::connect(&daemon_sock_for_broker)
.await
.expect("broker dials daemon");
let _ = tokio::io::copy_bidirectional(&mut client_conn, &mut daemon_conn).await;
});
let client_io = UnixStream::connect(&broker_sock)
.await
.expect("client dials broker");
let mut client = session_framed(client_io);
let live = tokio::time::timeout(Duration::from_secs(5), client.next())
.await
.expect("first output within 5s")
.expect("a frame")
.expect("decode ok");
assert_eq!(
live.kind,
Some(session_frame::Kind::Stdout(b"ping".to_vec())),
"session is live over the relay before the kill"
);
let t0 = Instant::now();
drop(client);
let outcome = tokio::time::timeout(Duration::from_secs(2), daemon).await;
let latency = t0.elapsed();
assert!(
outcome.is_ok(),
"daemon leaked the session after client death instead of tearing it down"
);
outcome
.unwrap()
.expect("daemon task")
.expect("daemon reaps the child cleanly");
assert!(
latency < Duration::from_secs(1),
"client-death teardown latency {latency:?} exceeds the 1s bound"
);
let _ = tokio::time::timeout(Duration::from_secs(5), broker)
.await
.expect("broker relay completes after both legs close");
let _ = std::fs::remove_file(&daemon_sock);
let _ = std::fs::remove_file(&broker_sock);
}
#[cfg(unix)]
#[tokio::test]
async fn daemon_kill_mid_session_detected_by_client_within_bound() {
use std::time::{Duration, Instant};
use tokio::net::{UnixListener, UnixStream};
let pid = std::process::id();
let daemon_sock = std::env::temp_dir().join(format!("rp-dkill-d-{pid}.sock"));
let broker_sock = std::env::temp_dir().join(format!("rp-dkill-b-{pid}.sock"));
let _ = std::fs::remove_file(&daemon_sock);
let _ = std::fs::remove_file(&broker_sock);
let daemon_listener = UnixListener::bind(&daemon_sock).expect("bind daemon listener");
let broker_listener = UnixListener::bind(&broker_sock).expect("bind broker listener");
let mut cmd = fixture();
cmd.args(["out:ping", "echo"]);
let group = Arc::new(ContainedProcessGroup::new().expect("contained group"));
let daemon = tokio::spawn(async move {
let (io, _addr) = daemon_listener.accept().await.expect("daemon accept");
run_compile_session(session_framed(io), cmd, group).await
});
let daemon_sock_for_broker = daemon_sock.clone();
let broker = tokio::spawn(async move {
let (mut client_conn, _addr) = broker_listener.accept().await.expect("broker accept");
let mut daemon_conn = UnixStream::connect(&daemon_sock_for_broker)
.await
.expect("broker dials daemon");
let _ = tokio::io::copy_bidirectional(&mut client_conn, &mut daemon_conn).await;
});
let client_io = UnixStream::connect(&broker_sock)
.await
.expect("client dials broker");
let mut client = session_framed(client_io);
let live = tokio::time::timeout(Duration::from_secs(5), client.next())
.await
.expect("first output within 5s")
.expect("a frame")
.expect("decode ok");
assert_eq!(
live.kind,
Some(session_frame::Kind::Stdout(b"ping".to_vec())),
"session is live over the relay before the kill"
);
let t0 = Instant::now();
daemon.abort();
let detected = tokio::time::timeout(Duration::from_secs(2), async {
loop {
match client.next().await {
None | Some(Err(_)) => break,
Some(Ok(_)) => continue,
}
}
})
.await;
let latency = t0.elapsed();
assert!(
detected.is_ok(),
"client hung after daemon death instead of detecting it"
);
assert!(
latency < Duration::from_secs(1),
"daemon-death detection latency {latency:?} exceeds the 1s bound"
);
broker.abort();
let _ = std::fs::remove_file(&daemon_sock);
let _ = std::fs::remove_file(&broker_sock);
}
#[tokio::test]
async fn serve_session_runs_command_from_start_frame() {
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let server = session_framed(server_io);
let mut client = session_framed(client_io);
let group = Arc::new(ContainedProcessGroup::new().expect("contained group"));
let handler = tokio::spawn(serve_session(server, group));
let script = ["out:ARTIFACT", "err:DIAGNOSTIC", "out:MORE", "exit:5"];
client
.send(frame(session_frame::Kind::Start(SessionStart {
program: fixture_program(),
args: script.iter().map(|s| s.to_string()).collect(),
cwd: String::new(),
env: Vec::new(),
clear_inherited_env: false,
environment_policy: 0,
})))
.await
.expect("send start");
client
.send(frame(session_frame::Kind::StdinEof(true)))
.await
.expect("send stdin eof");
let mut stdout = Vec::new();
let mut stderr = Vec::new();
let mut code = None;
while let Some(Ok(sf)) = client.next().await {
match sf.kind {
Some(session_frame::Kind::Stdout(b)) => stdout.extend_from_slice(&b),
Some(session_frame::Kind::Stderr(b)) => stderr.extend_from_slice(&b),
Some(session_frame::Kind::Exit(e)) => {
code = Some(e.code);
break;
}
_ => panic!("unexpected inbound-only frame on the outbound lane"),
}
}
let exit = handler
.await
.expect("handler task")
.expect("serve reaps child");
let oracle = run_direct(&script, b"");
assert_eq!(code, Some(exit.code), "client-visible exit matches handler");
assert_eq!(stdout, oracle.stdout, "stdout byte-exact");
assert_eq!(stderr, oracle.stderr, "stderr byte-exact");
assert_eq!(exit.code, oracle.code, "exit code fidelity");
}
#[tokio::test]
async fn serve_session_enforces_all_environment_policies_and_metadata() {
const DAEMON_ONLY: &str = "RUNNING_PROCESS_TEST_DAEMON_ONLY_ENV";
const CLIENT_ONLY: &str = "RUNNING_PROCESS_TEST_CLIENT_ONLY_ENV";
#[cfg(windows)]
const BASELINE_KEY: &str = "USERNAME";
#[cfg(unix)]
const BASELINE_KEY: &str = "HOME";
for (name, wire_policy, legacy_clear, inherits, baseline) in [
("legacy-inherit", 0, false, true, false),
("legacy-clear", 0, true, false, false),
("explicit-inherit", 1, false, true, false),
("user-baseline", 2, true, false, true),
("client-snapshot", 3, true, false, false),
] {
let _daemon_guard = TestEnvVarGuard::set(DAEMON_ONLY, "daemon-only");
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let server = session_framed(server_io);
let mut client = session_framed(client_io);
let group = Arc::new(
ContainedProcessGroup::with_originator("SESSION-POLICY").expect("contained group"),
);
let handler = tokio::spawn(serve_session(server, group));
#[cfg(windows)]
let (program, args) = (
"cmd.exe".to_owned(),
vec!["/D".to_owned(), "/C".to_owned(), "set".to_owned()],
);
#[cfg(not(windows))]
let (program, args) = ("/usr/bin/env".to_owned(), Vec::new());
client
.send(frame(session_frame::Kind::Start(SessionStart {
program,
args,
cwd: String::new(),
env: vec![crate::broker::protocol_v2::SessionEnvVar {
key: CLIENT_ONLY.to_owned(),
value: "forwarded".to_owned(),
}],
clear_inherited_env: legacy_clear,
environment_policy: wire_policy,
})))
.await
.expect("send start");
client
.send(frame(session_frame::Kind::StdinEof(true)))
.await
.expect("send stdin eof");
let mut stdout = Vec::new();
while let Some(Ok(sf)) = client.next().await {
match sf.kind {
Some(session_frame::Kind::Stdout(bytes)) => stdout.extend_from_slice(&bytes),
Some(session_frame::Kind::Exit(_)) => break,
_ => {}
}
}
handler.await.expect("handler task").expect("serve session");
let output = String::from_utf8_lossy(&stdout);
assert!(
output.contains(&format!("{CLIENT_ONLY}=forwarded")),
"{name}: client environment did not reach child"
);
assert_eq!(
output.contains(DAEMON_ONLY),
inherits,
"{name}: daemon ambient environment visibility"
);
assert!(
output.contains("RUNNING_PROCESS_ORIGINATOR=SESSION-POLICY:"),
"{name}: allowlisted originator metadata missing"
);
if baseline {
assert!(
output.contains(&format!("{BASELINE_KEY}=")),
"{name}: baseline identity key missing"
);
}
}
}
#[tokio::test]
async fn serve_session_rejects_missing_start_frame() {
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let server = session_framed(server_io);
let mut client = session_framed(client_io);
let group = Arc::new(ContainedProcessGroup::new().expect("contained group"));
let handler = tokio::spawn(serve_session(server, group));
client
.send(frame(session_frame::Kind::Stdin(b"oops".to_vec())))
.await
.expect("send stdin");
drop(client);
let result = handler.await.expect("handler task");
assert!(
result.is_err(),
"serve_session must reject a session that skips SessionStart"
);
}
#[tokio::test]
async fn serve_session_rejects_unknown_environment_policy() {
let (server_io, client_io) = tokio::io::duplex(64 * 1024);
let server = session_framed(server_io);
let mut client = session_framed(client_io);
let group = Arc::new(ContainedProcessGroup::new().expect("contained group"));
let handler = tokio::spawn(serve_session(server, group));
client
.send(frame(session_frame::Kind::Start(SessionStart {
program: fixture_program(),
args: Vec::new(),
cwd: String::new(),
env: Vec::new(),
clear_inherited_env: false,
environment_policy: 99,
})))
.await
.expect("send start");
drop(client);
let error = handler
.await
.expect("handler task")
.expect_err("unknown policy must fail closed");
assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput);
}