mod common;
use std::process::{Command, Stdio};
use std::time::{Duration, Instant};
const SIBLING_SLEEP: &str = "90s";
const HARD_TIMEOUT: Duration = Duration::from_secs(20);
fn write_file(tag: &str, ext: &str, body: &str) -> String {
let path = common::unique_path(tag, ext);
std::fs::write(&path, body).expect("write test file");
path
}
fn job(steps: &str) -> String {
write_file(
"agentd-nested-cancel",
"yaml",
&format!(
"config_version: \"2\"\nagent:\n name: nested-cancel\nworkflows:\n - name: pipe\n steps: {steps}\nlifecycle:\n run_until: idle\n idle_grace: 1s\nobservability:\n log_level: info\n log_content: true\n"
),
)
}
fn run_agentd_bounded(config: &str) -> (Option<i32>, String, Duration) {
let err_path = common::unique_path("agentd-nested-cancel", "err");
let err = std::fs::File::create(&err_path).expect("create stderr file");
let started = Instant::now();
let mut child = Command::new(env!("CARGO_BIN_EXE_agentd"))
.args(["--config", config])
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::from(err))
.spawn()
.expect("spawn agentd");
let deadline = started + HARD_TIMEOUT;
loop {
match child.try_wait().expect("wait for agentd") {
Some(status) => {
let log = std::fs::read_to_string(&err_path).unwrap_or_default();
let _ = std::fs::remove_file(&err_path);
return (status.code(), log, started.elapsed());
}
None if Instant::now() >= deadline => {
let _ = child.kill();
let _ = child.wait();
let log = std::fs::read_to_string(&err_path).unwrap_or_default();
let _ = std::fs::remove_file(&err_path);
panic!(
"agentd never exited (waited {HARD_TIMEOUT:?}): a nested body's siblings were never cancelled, so their {SIBLING_SLEEP} sleep timers keep the instance busy.\nstderr:\n{log}"
);
}
None => std::thread::sleep(Duration::from_millis(25)),
}
}
}
fn events(stderr: &str, name: &str) -> Vec<serde_json::Value> {
stderr
.lines()
.filter_map(|l| serde_json::from_str::<serde_json::Value>(l).ok())
.filter(|v| v["event"] == name)
.collect()
}
fn saw_step(stderr: &str, name: &str, step: &str) -> bool {
events(stderr, name).iter().any(|e| e["step"] == step)
}
#[test]
fn a_failed_foreach_element_cancels_the_siblings_still_sleeping() {
let steps = format!(
r#"{{
"start": {{"kind": "once"}},
"each": {{"kind": "foreach", "depends_on": ["start"], "over": ["hold", "boom"], "batch": {{"parallel": 2}},
"body": {{"steps": {{
"route": {{"kind": "switch", "on": "{{{{item}}}}", "cases": {{"hold": "sleeper", "boom": "explode"}}}},
"sleeper": {{"kind": "sleep", "depends_on": ["route"], "duration": "{SIBLING_SLEEP}"}},
"explode": {{"kind": "fail", "depends_on": ["route"], "message": "element {{{{index}}}} exploded"}}}}}}}},
"done": {{"kind": "finish", "depends_on": ["each"], "status": "completed"}}
}}"#
);
let (code, stderr, took) = run_agentd_bounded(&job(&steps));
assert!(code.is_some(), "the daemon exited on its own:\n{stderr}");
assert!(
saw_step(&stderr, "step.start", "each[0].sleeper"),
"the sibling element started sleeping:\n{stderr}"
);
assert!(
events(&stderr, "step.done")
.iter()
.any(|e| e["step"] == "each[1].explode" && e["status"] == "failed"),
"the second element failed:\n{stderr}"
);
assert!(
!saw_step(&stderr, "step.done", "each[0].sleeper"),
"the cancelled sibling must never finish its sleep:\n{stderr}"
);
assert!(
took < Duration::from_secs(15),
"the instance idled out promptly (took {took:?}) — a still-armed sibling timer would hold it for {SIBLING_SLEEP}:\n{stderr}"
);
let done = events(&stderr, "run.done");
assert_eq!(
done.len(),
1,
"the run reached a terminal status:\n{stderr}"
);
assert_eq!(done[0]["status"], "failed", "{stderr}");
}
#[test]
fn a_race_timeout_fires_and_cancels_every_branch() {
let steps = format!(
r#"{{
"start": {{"kind": "once"}},
"pick": {{"kind": "race", "depends_on": ["start"], "timeout": "300ms",
"branches": {{"slow": {{"steps": {{"s": {{"kind": "sleep", "duration": "{SIBLING_SLEEP}"}}}}}},
"slower": {{"steps": {{"s": {{"kind": "sleep", "duration": "{SIBLING_SLEEP}"}}}}}}}}}},
"done": {{"kind": "finish", "depends_on": ["pick"], "status": "completed"}}
}}"#
);
let (code, stderr, took) = run_agentd_bounded(&job(&steps));
assert!(code.is_some(), "the daemon exited on its own:\n{stderr}");
for branch in ["pick{slow}.s", "pick{slower}.s"] {
assert!(
saw_step(&stderr, "step.start", branch),
"branch step {branch} started:\n{stderr}"
);
assert!(
!saw_step(&stderr, "step.done", branch),
"branch step {branch} was cancelled, not slept out:\n{stderr}"
);
}
assert!(
events(&stderr, "step.done")
.iter()
.any(|e| e["step"] == "pick" && e["status"] == "timeout"),
"the race timed out:\n{stderr}"
);
assert!(
took < Duration::from_secs(15),
"the instance idled out promptly (took {took:?}):\n{stderr}"
);
let done = events(&stderr, "run.done");
assert_eq!(
done.len(),
1,
"the run reached a terminal status:\n{stderr}"
);
assert_eq!(done[0]["status"], "failed", "{stderr}");
}