mod common;
use std::process::{Command, Stdio};
use std::time::{Duration, Instant};
const HARD_TIMEOUT: Duration = Duration::from_secs(30);
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 run_agentd_bounded(config: &str, inbox: &str) -> (Option<i32>, String) {
let err_path = common::unique_path("agentd-overflow", "err");
let err = std::fs::File::create(&err_path).expect("create stderr file");
let mut child = Command::new(env!("CARGO_BIN_EXE_agentd"))
.args(["--config", config])
.env("AGENTD_TEST_INBOX_FILE", inbox)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::from(err))
.spawn()
.expect("spawn agentd");
let deadline = Instant::now() + 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);
}
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:?}): the inbox drain is livelocked on a queued start event.\nstderr:\n{log}"
);
}
None => std::thread::sleep(Duration::from_millis(50)),
}
}
}
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 nth_event_index(stderr: &str, name: &str, n: usize) -> usize {
stderr
.lines()
.filter_map(|l| serde_json::from_str::<serde_json::Value>(l).ok())
.enumerate()
.filter(|(_, v)| v["event"] == name)
.map(|(i, _)| i)
.nth(n)
.unwrap_or_else(|| panic!("no {name}[{n}] line in:\n{stderr}"))
}
fn capped_config(concurrency: &str) -> String {
let steps = r#"{
"start": {"kind": "manual"},
"hold": {"kind": "sleep", "depends_on": ["start"], "duration": "400ms"},
"done": {"kind": "finish", "depends_on": ["hold"], "status": "completed", "output": {"n": "{{inputs.n}}"}}
}"#;
write_file(
"agentd-overflow",
"yaml",
&format!(
"config_version: \"1\"\nagent:\n name: overflow\nworkflows:\n - name: capped\n{concurrency} steps: {steps}\nlifecycle:\n run_until: idle\n idle_grace: 1s\nobservability:\n log_level: info\n log_content: true\n"
),
)
}
fn inbox_for(n: usize) -> String {
let events: Vec<serde_json::Value> = (1..=n)
.map(|i| {
serde_json::json!({"kind": "workflow_run", "payload": {"workflow": "capped", "node": "start", "inputs": {"n": i.to_string()}}})
})
.collect();
write_file(
"inbox-overflow",
"json",
&serde_json::Value::Array(events).to_string(),
)
}
fn assert_all_runs_completed(code: Option<i32>, stderr: &str, expected: usize) {
assert_eq!(code, Some(0), "stderr:\n{stderr}");
let started = events(stderr, "run.start");
assert_eq!(started.len(), expected, "every start event ran:\n{stderr}");
let done = events(stderr, "run.done");
assert_eq!(
done.len(),
expected,
"every run reached a terminal status:\n{stderr}"
);
assert!(
done.iter().all(|d| d["status"] == "completed"),
"no run failed or was dropped: {done:?}"
);
let ids: std::collections::BTreeSet<&str> =
done.iter().filter_map(|d| d["run"].as_str()).collect();
assert_eq!(ids.len(), expected, "distinct runs: {done:?}");
}
#[test]
fn a_start_event_over_an_explicit_max_runs_cap_runs_on_a_later_tick() {
let cfg = capped_config(" concurrency: {max_runs: 1, on_overflow: queue}\n");
let (code, stderr) = run_agentd_bounded(&cfg, &inbox_for(2));
assert_all_runs_completed(code, &stderr, 2);
assert!(
nth_event_index(&stderr, "run.start", 1) > nth_event_index(&stderr, "run.done", 0),
"the second run must wait for the first to finish:\n{stderr}"
);
}
#[test]
fn start_events_over_the_default_cap_run_on_later_ticks() {
let cfg = capped_config("");
let (code, stderr) = run_agentd_bounded(&cfg, &inbox_for(5));
assert_all_runs_completed(code, &stderr, 5);
assert!(
nth_event_index(&stderr, "run.start", 4) > nth_event_index(&stderr, "run.done", 0),
"the fifth run must wait for a slot:\n{stderr}"
);
}