use std::io::Write;
use std::os::unix::net::UnixStream;
use std::path::Path;
use std::time::{Duration, Instant};
use anyhow::Context;
use zygo_core::pool::AGENT_FD;
use zygo_core::protocol::{ErrorCode, Message, PROTOCOL_VERSION, frame};
use crate::cli::Cli;
use crate::output::{self, Style};
const REPLY_TIMEOUT: Duration = Duration::from_secs(30);
const STALLED: &str = "stopped answering";
const OUT_OF_STEP: &str = "not asked: the agent is no longer in step with this suite";
const GO_GRACE: Duration = Duration::from_millis(300);
enum Outcome {
Pass(String),
Skipped(String),
Partial(String),
}
struct Report {
style: Style,
passed: u32,
failed: u32,
json: bool,
stalled: bool,
results: Vec<(String, bool, String)>,
skips: Vec<(String, String)>,
partials: Vec<(String, String)>,
}
impl Report {
fn new(json: bool) -> Self {
Self {
style: Style::stdout(),
passed: 0,
failed: 0,
json,
stalled: false,
results: Vec::new(),
skips: Vec::new(),
partials: Vec::new(),
}
}
fn ok(&mut self, what: impl Into<String>, detail: impl Into<String>) {
let (what, detail) = (what.into(), detail.into());
self.passed += 1;
if !self.json {
if detail.is_empty() {
println!(" {} {what}", self.style.green("PASS"));
} else {
println!(
" {} {what} — {}",
self.style.green("PASS"),
self.style.dim(&detail)
);
}
}
self.results.push((what, true, detail));
}
fn bad(&mut self, what: impl Into<String>, detail: impl Into<String>) {
let (what, detail) = (what.into(), detail.into());
self.stalled |= detail.contains(STALLED);
self.failed += 1;
if !self.json {
println!(" {} {what} — {detail}", self.style.red("FAIL"));
}
self.results.push((what, false, detail));
}
fn skipped(&mut self, what: &str, why: &str) {
if !self.json {
println!(
" {} {what} — {}",
self.style.yellow("SKIP"),
self.style.dim(why)
);
}
self.skips.push((what.to_string(), why.to_string()));
}
fn check(&mut self, what: &str, outcome: anyhow::Result<String>) -> bool {
self.record(what, outcome.map(Outcome::Pass))
}
fn partial(&mut self, what: &str, detail: &str) {
if !self.json {
println!(" {} {what} — {detail}", self.style.yellow("PART"));
}
self.partials.push((what.to_string(), detail.to_string()));
}
fn record(&mut self, what: &str, outcome: anyhow::Result<Outcome>) -> bool {
match outcome {
Ok(Outcome::Pass(detail)) => {
self.ok(what, detail);
true
}
Ok(Outcome::Skipped(why)) => {
self.skipped(what, &why);
true
}
Ok(Outcome::Partial(detail)) => {
self.partial(what, &detail);
true
}
Err(e) => {
self.bad(what, format!("{e:#}"));
false
}
}
}
}
struct Agent {
stream: UnixStream,
reader: frame::FrameReader<UnixStream>,
child: std::process::Child,
announced_pid: u32,
pool: Option<zygo_core::protocol::Script>,
}
impl Agent {
fn start(binary: &Path, args: &[String]) -> anyhow::Result<Agent> {
Agent::start_with_env(binary, args, &[])
}
fn start_with_env(
binary: &Path,
args: &[String],
env: &[(&str, String)],
) -> anyhow::Result<Agent> {
use std::os::fd::AsRawFd;
use std::os::unix::process::CommandExt;
let (ours, theirs) = UnixStream::pair().context("could not create a socket pair")?;
let their_fd = theirs.as_raw_fd();
let mut command = std::process::Command::new(binary);
command.args(args);
for (key, value) in env {
command.env(key, value);
}
unsafe {
command.pre_exec(move || {
if libc::setpgid(0, 0) < 0 {
return Err(std::io::Error::last_os_error());
}
if their_fd != AGENT_FD {
if libc::dup2(their_fd, AGENT_FD) < 0 {
return Err(std::io::Error::last_os_error());
}
} else if libc::fcntl(AGENT_FD, libc::F_SETFD, 0) < 0 {
return Err(std::io::Error::last_os_error());
}
Ok(())
});
}
let child = command
.spawn()
.with_context(|| format!("could not start the agent `{}`", binary.display()))?;
drop(theirs);
ours.set_read_timeout(Some(REPLY_TIMEOUT))
.context("could not set a read timeout")?;
let reader =
frame::FrameReader::new(ours.try_clone().context("could not clone the socket")?);
Ok(Agent {
stream: ours,
reader,
child,
announced_pid: 0,
pool: None,
})
}
fn exec(&self, id: &str, event: serde_json::Value) -> Message {
Message::Exec {
script: self.pool.clone(),
id: id.to_string(),
event,
timeout_ms: 30_000,
env_overrides: Default::default(),
stream: false,
workspace: None,
}
}
fn send(&mut self, message: &Message) -> anyhow::Result<()> {
let bytes = frame::encode(message)?;
self.stream.write_all(&bytes)?;
self.stream.flush()?;
Ok(())
}
fn send_raw(&mut self, body: &[u8]) -> anyhow::Result<()> {
let mut out = (body.len() as u32).to_be_bytes().to_vec();
out.extend_from_slice(body);
self.stream.write_all(&out)?;
self.stream.flush()?;
Ok(())
}
fn recv(&mut self) -> anyhow::Result<Message> {
match self.reader.read() {
Ok(Some(message)) => Ok(message),
Ok(None) => anyhow::bail!("the agent closed the connection"),
Err(frame::FrameError::Io(e))
if matches!(
e.kind(),
std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
) =>
{
anyhow::bail!("the agent {STALLED}: nothing at all for {REPLY_TIMEOUT:?}")
}
Err(e) => Err(e.into()),
}
}
fn recv_matching(
&mut self,
what: &str,
mut want: impl FnMut(&Message) -> bool,
) -> anyhow::Result<Message> {
let deadline = Instant::now() + REPLY_TIMEOUT;
loop {
let message = self.recv()?;
if want(&message) {
return Ok(message);
}
anyhow::ensure!(
Instant::now() < deadline,
"the agent {STALLED}: no {what} within {REPLY_TIMEOUT:?}"
);
}
}
fn quiet_for(&mut self, window: Duration) -> anyhow::Result<Option<Message>> {
self.stream.set_read_timeout(Some(window))?;
let outcome = match self.reader.read() {
Ok(Some(m)) => Ok(Some(m)),
Ok(None) => Err(anyhow::anyhow!("the agent closed the connection")),
Err(frame::FrameError::Io(e))
if matches!(
e.kind(),
std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
) =>
{
Ok(None)
}
Err(e) => Err(e.into()),
};
self.stream.set_read_timeout(Some(REPLY_TIMEOUT))?;
outcome
}
}
impl Drop for Agent {
fn drop(&mut self) {
unsafe {
libc::kill(-(self.child.id() as libc::pid_t), libc::SIGKILL);
}
let _ = self.child.kill();
let _ = self.child.wait();
}
}
fn pool_handler(file: &Path) -> anyhow::Result<zygo_core::protocol::Script> {
let source = std::fs::read_to_string(file)
.with_context(|| format!("could not read the pool handler {}", file.display()))?;
let digest = zygo_core::scripts::ScriptDigest::of(&source).to_string();
Ok(zygo_core::protocol::Script {
digest: Some(digest),
..zygo_core::protocol::Script::inline(source)
})
}
pub fn test(
cli: &Cli,
binary: &Path,
script: Option<&Path>,
spawning: Option<&Path>,
pool: Option<&Path>,
args: &[String],
) -> anyhow::Result<u8> {
let mut report = Report::new(cli.json);
let pool = pool.map(pool_handler).transpose()?;
if !cli.json {
println!(
"protocol conformance: {} (proto {PROTOCOL_VERSION})",
binary.display()
);
println!(
" {}",
Report::new(false).style.dim(
"the handler must echo the event, write `stdout`/`stderr` when present, \
and start a program for `spawn`"
)
);
if pool.is_some() {
println!(
" {}",
Report::new(false)
.style
.dim("runtime pool: the agent holds no handler, so every request carries one")
);
}
println!();
}
let mut agent = Agent::start(binary, args)?;
agent.pool = pool;
let ready = report.check(
"the agent announces itself with READY",
ready_check(&mut agent),
);
if !ready {
return finish(cli, report);
}
macro_rules! step {
($what:expr, $check:expr) => {
if report.stalled {
report.skipped($what, OUT_OF_STEP);
} else {
report.check($what, $check);
}
};
}
macro_rules! step_or_skip {
($what:expr, $check:expr) => {
if report.stalled {
report.skipped($what, OUT_OF_STEP);
} else {
report.record($what, $check);
}
};
}
step!(
"PING is answered by PONG with the same seq",
ping_check(&mut agent)
);
step!(
"EXEC is answered by FORKED naming a process that is not the agent",
forked_check(&mut agent)
);
step!("the child does nothing until GO", go_check(&mut agent));
step!(
"the event reaches the handler and its result comes back",
result_check(&mut agent)
);
step!(
"stdout and stderr come back in separate fields",
streams_check(&mut agent)
);
step!(
"two requests in flight are both answered, with their own ids",
concurrency_check(&mut agent)
);
step!(
"a frame that is not a message is an ERROR, not a crash",
bad_message_check(&mut agent)
);
step_or_skip!(
"a script in EXEC is loaded by the child (proto 1.1)",
script_in_exec_check(&mut agent, script)
);
step_or_skip!(
"ZYGO_CHILD_SECCOMP is installed in the child, or the request is refused",
child_seccomp_check(&mut agent, binary, args)
);
step_or_skip!(
"the child filter is installed before the script's first line (proto 1.1)",
script_before_filter_check(&mut agent, binary, args, spawning)
);
step_or_skip!(
"a cancelled request comes back as DONE{cancelled} (proto 1.2)",
cancel_check(&mut agent)
);
step_or_skip!(
"output arrives in CHUNKs before the DONE (proto 1.3)",
stream_check(&mut agent)
);
step_or_skip!(
"a long request is reported alive with PING{id} (proto 1.4)",
heartbeat_check(&mut agent)
);
step!("SHUTDOWN makes the agent exit", shutdown_check(&mut agent));
finish(cli, report)
}
fn finish(cli: &Cli, report: Report) -> anyhow::Result<u8> {
if cli.json {
output::json(&serde_json::json!({
"passed": report.passed,
"failed": report.failed,
"skipped": report.skips.len(),
"checks": report
.results
.iter()
.map(|(what, ok, detail)| serde_json::json!({
"check": what,
"passed": ok,
"detail": detail,
}))
.collect::<Vec<_>>(),
"skips": report
.skips
.iter()
.map(|(what, why)| serde_json::json!({ "check": what, "reason": why }))
.collect::<Vec<_>>(),
"partial": report
.partials
.iter()
.map(|(what, detail)| serde_json::json!({ "check": what, "detail": detail }))
.collect::<Vec<_>>(),
}))?;
} else {
println!();
let skipped = match report.skips.len() {
0 => String::new(),
n => format!(", {n} not asked"),
};
let partial = match report.partials.len() {
0 => String::new(),
n => format!(", {n} partial"),
};
if report.failed == 0 && !report.partials.is_empty() {
println!(
"{} {} checks passed{skipped}{partial} — conforms to protocol \
{PROTOCOL_VERSION}, with gaps named above",
report.style.yellow("~"),
report.passed
);
} else if report.failed == 0 {
println!(
"{} {} checks passed{skipped} — this agent conforms to protocol {PROTOCOL_VERSION}",
report.style.green("✓"),
report.passed
);
} else {
println!(
"{} {} passed, {} failed{skipped}{partial}",
report.style.red("✗"),
report.passed,
report.failed
);
}
}
Ok(u8::from(report.failed > 0))
}
fn ready_check(agent: &mut Agent) -> anyhow::Result<String> {
match agent.recv()? {
Message::Ready {
proto,
pid,
runtime,
rss_kb,
imports_ms,
} => {
anyhow::ensure!(
proto == PROTOCOL_VERSION,
"announced proto {proto}, this supervisor speaks {PROTOCOL_VERSION}"
);
anyhow::ensure!(pid > 0, "READY carries no pid");
anyhow::ensure!(!runtime.trim().is_empty(), "READY carries no runtime name");
agent.announced_pid = pid;
Ok(format!(
"{runtime}, pid {pid}, {} MB resident after {imports_ms:.0} ms",
rss_kb / 1024
))
}
other => anyhow::bail!("the first message was {} rather than READY", other.kind()),
}
}
fn ping_check(agent: &mut Agent) -> anyhow::Result<String> {
agent.send(&Message::Ping { seq: 7, id: None })?;
let reply = agent.recv_matching("PONG", |m| matches!(m, Message::Pong { .. }))?;
match reply {
Message::Pong { seq } => {
anyhow::ensure!(seq == 7, "PONG carried seq {seq}, not the 7 that was sent");
Ok("seq 7".into())
}
other => anyhow::bail!("expected PONG, got {}", other.kind()),
}
}
fn forked_check(agent: &mut Agent) -> anyhow::Result<String> {
let request = agent.exec("c1", serde_json::json!({ "n": 1 }));
agent.send(&request)?;
let forked = agent.recv_matching("FORKED", |m| matches!(m, Message::Forked { .. }))?;
match forked {
Message::Forked { id, pid } => {
anyhow::ensure!(id == "c1", "FORKED carried id `{id}`, not `c1`");
anyhow::ensure!(pid > 0, "FORKED carries no pid");
anyhow::ensure!(
pid != agent.announced_pid,
"the request runs in the agent's own process ({pid}); \
a request must have a pid of its own to be put in a cgroup"
);
Ok(format!("child pid {pid}"))
}
other => anyhow::bail!("expected FORKED, got {}", other.kind()),
}
}
fn go_check(agent: &mut Agent) -> anyhow::Result<String> {
if let Some(early) = agent.quiet_for(GO_GRACE)? {
anyhow::bail!(
"{} arrived {GO_GRACE:?} after FORKED and before GO; the child \
started work outside its cgroup",
early.kind()
);
}
agent.send(&Message::Go { id: "c1".into() })?;
let done = agent.recv_matching("DONE", |m| matches!(m, Message::Done { .. }))?;
match done {
Message::Done { id, .. } => {
anyhow::ensure!(id == "c1", "DONE carried id `{id}`, not `c1`");
Ok(format!("silent for {GO_GRACE:?}, then answered"))
}
other => anyhow::bail!("expected DONE, got {}", other.kind()),
}
}
fn result_check(agent: &mut Agent) -> anyhow::Result<String> {
let event = serde_json::json!({ "n": 42, "s": "hello" });
let done = one_request(agent, "c2", event.clone())?;
match done {
Message::Done {
exit_code,
result,
error,
..
} => {
anyhow::ensure!(
exit_code == 0,
"the echo handler exited {exit_code}{}",
error.map(|e| format!(": {e}")).unwrap_or_default()
);
anyhow::ensure!(
result == event,
"the handler was given {event} and the result was {result}"
);
Ok("the event came back unchanged".into())
}
other => anyhow::bail!("expected DONE, got {}", other.kind()),
}
}
fn streams_check(agent: &mut Agent) -> anyhow::Result<String> {
let event = serde_json::json!({ "stdout": "to-stdout", "stderr": "to-stderr" });
let done = one_request(agent, "c3", event)?;
match done {
Message::Done { stdout, stderr, .. } => {
anyhow::ensure!(
stdout.contains("to-stdout"),
"what the handler printed is not in `stdout`: {stdout:?}"
);
anyhow::ensure!(
stderr.contains("to-stderr"),
"what the handler wrote to stderr is not in `stderr`: {stderr:?}"
);
anyhow::ensure!(
!stdout.contains("to-stderr"),
"the two streams are mixed: stdout holds {stdout:?}"
);
Ok("kept apart".into())
}
other => anyhow::bail!("expected DONE, got {}", other.kind()),
}
}
fn concurrency_check(agent: &mut Agent) -> anyhow::Result<String> {
let first = agent.exec("a", serde_json::json!({ "which": "a" }));
let second = agent.exec("b", serde_json::json!({ "which": "b" }));
agent.send(&first)?;
agent.send(&second)?;
let mut answered: std::collections::BTreeMap<String, Option<serde_json::Value>> =
Default::default();
let mut refused = 0;
let deadline = Instant::now() + REPLY_TIMEOUT;
while answered.len() < 2 {
anyhow::ensure!(
Instant::now() < deadline,
"the agent {STALLED}: {} of 2 requests answered within {REPLY_TIMEOUT:?}",
answered.len()
);
match agent.recv()? {
Message::Forked { id, .. } => agent.send(&Message::Go { id })?,
Message::Done { id, result, .. } => {
anyhow::ensure!(
answered.insert(id.clone(), Some(result)).is_none(),
"`{id}` was answered twice"
);
}
Message::Error {
id: Some(id),
code: ErrorCode::Overloaded,
..
} => {
refused += 1;
anyhow::ensure!(
answered.insert(id.clone(), None).is_none(),
"`{id}` was answered twice"
);
}
Message::Error { id, code, message } => {
anyhow::bail!("{id:?} failed with {code:?}: {message}")
}
_ => {}
}
}
for id in ["a", "b"] {
let entry = answered
.get(id)
.with_context(|| format!("`{id}` was never answered"))?;
if let Some(result) = entry {
anyhow::ensure!(
result.get("which").and_then(|v| v.as_str()) == Some(id),
"`{id}` came back with another request's event: {result}"
);
}
}
anyhow::ensure!(refused < 2, "both requests were refused as overloaded");
Ok(match refused {
0 => "both ran, neither crossed".into(),
_ => format!(
"{} ran, {refused} refused as overloaded (a serial agent)",
2 - refused
),
})
}
fn bad_message_check(agent: &mut Agent) -> anyhow::Result<String> {
agent.send_raw(b"{ this is not json")?;
let reply = agent.recv_matching("ERROR", |m| matches!(m, Message::Error { .. }))?;
match reply {
Message::Error { code, .. } => {
anyhow::ensure!(
code == ErrorCode::BadMessage,
"reported {code:?} rather than `bad_message`"
);
agent.send(&Message::Ping { seq: 9, id: None })?;
agent.recv_matching("PONG", |m| matches!(m, Message::Pong { seq: 9 }))?;
Ok("reported, and the agent is still serving".into())
}
other => anyhow::bail!("expected ERROR, got {}", other.kind()),
}
}
const SPAWN_MARK: &str = "zygo-child-spawned";
const SCRIPT_MARK: &str = "the-request";
fn script_in_exec_check(agent: &mut Agent, script: Option<&Path>) -> anyhow::Result<Outcome> {
let Some(script) = script else {
return Ok(Outcome::Skipped(
"not asked: pass --script <file> in this agent's own language to check \
protocol 1.1"
.into(),
));
};
let source = std::fs::read_to_string(script)
.with_context(|| format!("could not read the script {}", script.display()))?;
let digest = zygo_core::scripts::ScriptDigest::of(&source);
let inline = zygo_core::protocol::Script::inline(source);
let event = serde_json::json!({ "n": 11 });
match ran_the_script(agent, "c-script", &event, inline.clone())? {
ScriptRun::Ran => {}
ScriptRun::NotSupported(why) => return Ok(Outcome::Skipped(why)),
}
let staged = tempfile::Builder::new()
.prefix("zygo-conformance-")
.suffix(&script_suffix(script))
.tempfile()
.context("could not stage a script for the `path` shape")?;
std::fs::write(staged.path(), inline.source.as_deref().unwrap_or_default())
.context("could not write the staged script")?;
let by_path =
zygo_core::protocol::Script::at(staged.path().display().to_string(), digest.to_string());
let path_shape = ran_the_script(agent, "c-script-path", &event, by_path)?;
let mut tampered = inline.clone();
tampered.digest = Some(zygo_core::scripts::ScriptDigest::of("not this script").to_string());
let refused = matches!(
ran_the_script(agent, "c-script-digest", &event, tampered)?,
ScriptRun::NotSupported(_)
);
match (path_shape, refused) {
(ScriptRun::Ran, true) => Ok(Outcome::Pass(
"the request's own script ran, from `source` and from `path`, and a \
digest that did not match was refused"
.into(),
)),
(ScriptRun::Ran, false) => anyhow::bail!(
"the agent ran a script whose bytes do not hash to the `digest` in \
its own EXEC\n → §3.8: hash what you read and answer ERROR / \
handler_load instead. Without it a tenant can swap \
/run/script/<hash> for its own code and be served it"
),
(ScriptRun::NotSupported(why), _) => Ok(Outcome::Partial(format!(
"`source` works and `path` does not ({why})\n → the supervisor \
sends `path` whenever it can write into the sandbox, which is the \
usual case, so this agent would fail those requests"
))),
}
}
fn script_suffix(script: &Path) -> String {
match script.extension().and_then(|e| e.to_str()) {
Some(extension) => format!(".{extension}"),
None => String::new(),
}
}
enum ScriptRun {
Ran,
NotSupported(String),
}
fn ran_the_script(
agent: &mut Agent,
id: &str,
event: &serde_json::Value,
script: zygo_core::protocol::Script,
) -> anyhow::Result<ScriptRun> {
agent.send(&Message::Exec {
id: id.to_string(),
event: event.clone(),
timeout_ms: 30_000,
env_overrides: Default::default(),
script: Some(script),
stream: false,
workspace: None,
})?;
let done = loop {
match agent.recv()? {
Message::Forked { id: forked, .. } if forked == id => {
agent.send(&Message::Go { id: forked })?;
}
done @ Message::Done { .. } if done.request_id() == Some(id) => break done,
Message::Error { code, message, .. } => {
return Ok(ScriptRun::NotSupported(format!(
"refused with {code:?} ({})",
first_line(&message)
)));
}
_ => {}
}
};
match done {
Message::Done {
exit_code,
result,
error,
..
} => {
if result.get("from").and_then(|v| v.as_str()) == Some(SCRIPT_MARK) {
return Ok(ScriptRun::Ran);
}
if result == *event {
return Ok(ScriptRun::NotSupported(
"this agent serves functions only: it ignored the script and ran \
the handler it was warmed with, which §5 allows"
.into(),
));
}
anyhow::ensure!(
exit_code != 0 || error.is_some(),
"the agent answered neither the script's result nor its own \
handler's: {result}"
);
Ok(ScriptRun::NotSupported(format!(
"the script failed rather than running ({})",
error.map(|e| first_line(&e)).unwrap_or_default()
)))
}
other => anyhow::bail!("expected DONE, got {}", other.kind()),
}
}
fn child_seccomp_check(
warm: &mut Agent,
binary: &Path,
args: &[String],
) -> anyhow::Result<Outcome> {
let Some(filter) = strict_child_filter() else {
return Ok(Outcome::Skipped(
"not asked: seccomp is a Linux facility, and this is not Linux".into(),
));
};
let control = one_request(warm, "c-spawn", serde_json::json!({ "spawn": SPAWN_MARK }))?;
if !spawned(&control) {
return Ok(Outcome::Skipped(format!(
"not asked: the handler did not start a program even without a filter, \
so a refusal would prove nothing ({})",
summarise(&control)
)));
}
let mut tightened = Agent::start_with_env(
binary,
args,
&[(zygo_core::protocol::CHILD_SECCOMP_ENV, filter)],
)?;
tightened.pool = warm.pool.clone();
match tightened.recv()? {
Message::Ready { .. } => {}
Message::Error { code, message, .. } => {
return Ok(Outcome::Pass(format!(
"refused to start under the filter ({code:?}: {})",
first_line(&message)
)));
}
other => anyhow::bail!(
"the first message under {} was {} rather than READY",
zygo_core::protocol::CHILD_SECCOMP_ENV,
other.kind()
),
}
match one_request(
&mut tightened,
"c-filtered",
serde_json::json!({ "spawn": SPAWN_MARK }),
) {
Ok(done) => {
anyhow::ensure!(
!spawned(&done),
"the child started a program with {} set: the filter was ignored and \
the request ran with the sandbox's filter alone",
zygo_core::protocol::CHILD_SECCOMP_ENV
);
Ok(Outcome::Pass(format!(
"the child could not start a program ({})",
summarise(&done)
)))
}
Err(e) if e.to_string().contains("the agent reported") => Ok(Outcome::Pass(format!(
"the request was refused rather than run unfiltered ({})",
first_line(&e.to_string())
))),
Err(e) => Err(e),
}
}
fn script_before_filter_check(
warm: &mut Agent,
binary: &Path,
args: &[String],
spawning: Option<&Path>,
) -> anyhow::Result<Outcome> {
let Some(filter) = strict_child_filter() else {
return Ok(Outcome::Skipped(
"not asked: seccomp is a Linux facility, and this is not Linux".into(),
));
};
let Some(spawning) = spawning else {
return Ok(Outcome::Skipped(
"not asked: pass --script-spawn <file> — a script whose module body starts \
a program — to check when the child filter is installed"
.into(),
));
};
let source = std::fs::read_to_string(spawning)
.with_context(|| format!("could not read {}", spawning.display()))?;
let script = zygo_core::protocol::Script::inline(source);
let control = one_script_request(warm, "c-spawn-script", script.clone())?;
if !spawned(&control) {
return Ok(Outcome::Skipped(format!(
"not asked: the script did not start a program even without a filter, \
so a refusal would prove nothing ({})",
summarise(&control)
)));
}
let mut tightened = Agent::start_with_env(
binary,
args,
&[(zygo_core::protocol::CHILD_SECCOMP_ENV, filter)],
)?;
match tightened.recv()? {
Message::Ready { .. } => {}
Message::Error { code, message, .. } => {
return Ok(Outcome::Pass(format!(
"refused to start under the filter ({code:?}: {})",
first_line(&message)
)));
}
other => anyhow::bail!("expected READY, got {}", other.kind()),
}
match one_script_request(&mut tightened, "c-spawn-filtered", script) {
Ok(done) => {
anyhow::ensure!(
!spawned(&done),
"the script started a program from its module body with {} set: the \
filter goes on after the script loads, so a `strict` pool does not \
restrict the script's own import-time code (spec/protocol.md §2)",
zygo_core::protocol::CHILD_SECCOMP_ENV
);
Ok(Outcome::Pass(format!(
"the script could not start a program while loading ({})",
summarise(&done)
)))
}
Err(e) if e.to_string().contains("the agent reported") => Ok(Outcome::Pass(format!(
"the request was refused rather than loaded unfiltered ({})",
first_line(&e.to_string())
))),
Err(e) => Err(e),
}
}
fn spawned(done: &Message) -> bool {
match done {
Message::Done { stdout, .. } => stdout.contains(SPAWN_MARK),
_ => false,
}
}
fn summarise(done: &Message) -> String {
match done {
Message::Done {
exit_code, error, ..
} => match error {
Some(e) => format!("exit {exit_code}: {}", first_line(e)),
None => format!("exit {exit_code}"),
},
other => other.kind().to_string(),
}
}
fn first_line(text: &str) -> String {
let line = text.lines().next().unwrap_or("").trim();
if line.len() > 90 {
format!("{}…", &line[..89])
} else {
line.to_string()
}
}
#[cfg(target_os = "linux")]
fn strict_child_filter() -> Option<String> {
use base64::Engine as _;
use zygo_core::backend::ns::seccomp;
let program = seccomp::child_program(zygo_core::spec::SeccompProfile::Strict).ok()??;
Some(base64::engine::general_purpose::STANDARD.encode(seccomp::encode(&program)))
}
#[cfg(not(target_os = "linux"))]
fn strict_child_filter() -> Option<String> {
None
}
fn shutdown_check(agent: &mut Agent) -> anyhow::Result<String> {
agent.send(&Message::Shutdown { grace_ms: 2_000 })?;
let deadline = Instant::now() + Duration::from_secs(5);
loop {
match agent.child.try_wait()? {
Some(status) => return Ok(format!("exited with {status}")),
None if Instant::now() >= deadline => {
anyhow::bail!("still running 5 s after SHUTDOWN")
}
None => std::thread::sleep(Duration::from_millis(20)),
}
}
}
fn one_script_request(
agent: &mut Agent,
id: &str,
script: zygo_core::protocol::Script,
) -> anyhow::Result<Message> {
agent.send(&Message::Exec {
id: id.to_string(),
event: serde_json::json!({ "spawn": SPAWN_MARK }),
timeout_ms: 30_000,
env_overrides: Default::default(),
script: Some(script),
stream: false,
workspace: None,
})?;
await_done(agent, id)
}
const HEARTBEAT_WAIT: Duration = Duration::from_secs(6);
fn heartbeat_check(agent: &mut Agent) -> anyhow::Result<Outcome> {
let id = "c11";
let request = agent.exec(id, serde_json::json!({ "sleep_ms": 30_000 }));
agent.send(&request)?;
let forked = agent.recv_matching("FORKED", |m| matches!(m, Message::Forked { .. }))?;
let Message::Forked { pid, .. } = forked else {
anyhow::bail!("expected FORKED, got {}", forked.kind());
};
agent.send(&Message::Go { id: id.to_string() })?;
let started = Instant::now();
let mut beat = None;
while started.elapsed() < HEARTBEAT_WAIT {
let left = HEARTBEAT_WAIT.saturating_sub(started.elapsed());
match agent.quiet_for(left)? {
Some(Message::Ping {
id: Some(pinged), ..
}) if pinged == id => {
beat = Some(started.elapsed());
break;
}
Some(_) => continue,
None => break,
}
}
unsafe { libc::kill(pid as libc::pid_t, libc::SIGKILL) };
let _ = agent.recv_matching("DONE", |m| matches!(m, Message::Done { .. }));
Ok(match beat {
Some(at) => Outcome::Pass(format!("a heartbeat for this request after {at:?}")),
None => Outcome::Skipped(format!(
"nothing for {HEARTBEAT_WAIT:?} while a request ran: heartbeats are \
not implemented"
)),
})
}
fn stream_check(agent: &mut Agent) -> anyhow::Result<Outcome> {
let id = "c10";
let early = "first";
agent.send(&Message::Exec {
id: id.to_string(),
event: serde_json::json!({ "stdout": early, "sleep_ms": 1_500 }),
timeout_ms: 30_000,
env_overrides: Default::default(),
script: agent.pool.clone(),
stream: true,
workspace: None,
})?;
let forked = agent.recv_matching("FORKED", |m| matches!(m, Message::Forked { .. }))?;
anyhow::ensure!(
matches!(forked, Message::Forked { .. }),
"expected FORKED, got {}",
forked.kind()
);
agent.send(&Message::Go { id: id.to_string() })?;
let started = Instant::now();
let mut first: Option<Duration> = None;
let mut streamed = String::new();
loop {
match agent.recv()? {
Message::Chunk {
stream: zygo_core::protocol::Stream::Stdout,
data,
..
} => {
first.get_or_insert_with(|| started.elapsed());
streamed.push_str(&data);
}
done @ Message::Done { .. } if done.request_id() == Some(id) => {
let Message::Done { stdout, .. } = done else {
unreachable!("matched above")
};
let Some(at) = first else {
return Ok(Outcome::Skipped(
"the request was answered, but nothing arrived before \
the DONE: `stream` is not implemented"
.into(),
));
};
anyhow::ensure!(
streamed.contains(early),
"what the handler printed first is not in the chunks: {streamed:?}"
);
anyhow::ensure!(
at < Duration::from_millis(1_200),
"the first chunk arrived after {at:?}, so it waited for the handler to finish rather than being forwarded"
);
anyhow::ensure!(
stdout.contains(early),
"DONE no longer carries the output it always did: {stdout:?}"
);
return Ok(Outcome::Pass(format!(
"the first chunk arrived after {at:?}"
)));
}
Message::Error { code, message, .. } => {
anyhow::bail!("the agent reported {code:?}: {message}")
}
_ => {}
}
}
}
const CANCEL_SLEEP_MS: u64 = 5_000;
const CANCEL_GRACE: Duration = Duration::from_secs(5);
fn cancel_check(agent: &mut Agent) -> anyhow::Result<Outcome> {
let id = "c9";
let request = agent.exec(id, serde_json::json!({ "sleep_ms": CANCEL_SLEEP_MS }));
agent.send(&request)?;
let forked = agent.recv_matching("FORKED", |m| matches!(m, Message::Forked { .. }))?;
let Message::Forked { pid, .. } = forked else {
anyhow::bail!("expected FORKED, got {}", forked.kind());
};
agent.send(&Message::Go { id: id.to_string() })?;
std::thread::sleep(Duration::from_millis(200));
agent.send(&Message::Cancel { id: id.to_string() })?;
unsafe { libc::kill(pid as libc::pid_t, libc::SIGKILL) };
let started = Instant::now();
loop {
match agent.recv()? {
done @ Message::Done { .. } if done.request_id() == Some(id) => {
let Message::Done { cancelled, .. } = done else {
unreachable!("matched above")
};
anyhow::ensure!(
started.elapsed() < CANCEL_GRACE,
"the agent answered {:?} after the kill, past the {CANCEL_GRACE:?} bound",
started.elapsed()
);
return Ok(match cancelled {
true => Outcome::Pass("answered with `cancelled` after the kill".into()),
false => Outcome::Skipped(
"answered the killed request, but without `cancelled`: \
CANCEL is not implemented"
.into(),
),
});
}
Message::Error { code, message, .. }
if code != zygo_core::protocol::ErrorCode::BadMessage =>
{
anyhow::bail!("the agent reported {code:?}: {message}")
}
_ => {}
}
anyhow::ensure!(
started.elapsed() < CANCEL_GRACE,
"the agent {STALLED}: no DONE for a killed request within {CANCEL_GRACE:?}"
);
}
}
fn one_request(agent: &mut Agent, id: &str, event: serde_json::Value) -> anyhow::Result<Message> {
let request = agent.exec(id, event);
agent.send(&request)?;
await_done(agent, id)
}
fn await_done(agent: &mut Agent, id: &str) -> anyhow::Result<Message> {
loop {
match agent.recv()? {
Message::Forked { id: forked, .. } if forked == id => {
agent.send(&Message::Go { id: id.to_string() })?;
}
done @ Message::Done { .. } if done.request_id() == Some(id) => return Ok(done),
Message::Error { code, message, .. } => {
anyhow::bail!("the agent reported {code:?}: {message}")
}
_ => {}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_missed_deadline_stops_the_run_and_a_wrong_answer_does_not() {
let mut report = Report::new(true);
report.bad("a check", "the agent said Pong when Forked was due");
assert!(!report.stalled, "a wrong answer is still an answer");
report.bad(
"another",
format!("the agent {STALLED}: nothing at all for 30s"),
);
assert!(report.stalled);
}
}