#![cfg(feature = "workflow")]
mod common;
use std::process::{Child, Command, Stdio};
use std::time::{Duration, Instant};
const HARD_TIMEOUT: Duration = Duration::from_secs(20);
fn write_config(yaml: &str) -> String {
let path = common::unique_path("schedule-at", "yaml");
std::fs::write(&path, yaml).expect("write config");
path
}
fn state_root(tag: &str) -> String {
let p = common::unique_path(tag, "state");
let _ = std::fs::remove_dir_all(&p);
p
}
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 spawn_daemon(config: &str, state_dir: &str) -> (Child, String) {
let err_path = common::unique_path("schedule-at", "err");
let err = std::fs::File::create(&err_path).expect("create stderr file");
let child = Command::new(env!("CARGO_BIN_EXE_agentd"))
.args(["--config", config])
.env("AGENTD_STATE_DIR", state_dir)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::from(err))
.spawn()
.expect("spawn agentd");
(child, err_path)
}
fn wait_bounded(mut child: Child, err_path: &str, what: &str) -> (Option<i32>, String) {
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();
panic!("{what}: agentd never exited (waited {HARD_TIMEOUT:?}).\nstderr:\n{log}");
}
None => std::thread::sleep(Duration::from_millis(25)),
}
}
}
fn run_daemon_until(
config: &str,
state_dir: &str,
marker: &str,
hold_ms: u64,
) -> (Option<i32>, String) {
let (child, err_path) = spawn_daemon(config, state_dir);
let pid = child.id() as i32;
let deadline = Instant::now() + HARD_TIMEOUT;
loop {
let log = std::fs::read_to_string(&err_path).unwrap_or_default();
if !events(&log, marker).is_empty() {
break;
}
assert!(
Instant::now() < deadline,
"no {marker:?} within {HARD_TIMEOUT:?}:
{log}"
);
std::thread::sleep(Duration::from_millis(25));
}
std::thread::sleep(Duration::from_millis(hold_ms));
unsafe { libc::kill(pid, libc::SIGTERM) };
wait_bounded(child, &err_path, "SIGTERMed life")
}
const AT_STEPS: &str = r#"{
"at_three": {"kind": "schedule", "at": "0s"},
"note": {"kind": "memory.set", "depends_on": ["at_three"], "key": "ran", "value": "1"},
"done": {"kind": "finish", "depends_on": ["note"], "status": "completed"}
}"#;
fn at_config(dir: &str) -> String {
write_config(&format!(
"config_version: \"1\"\nagent:\n name: oneshot\nstore:\n kind: file\n file:\n path: {dir}\n checkpoint:\n debounce_ms: 0\nworkflows:\n - name: nightly\n steps: {AT_STEPS}\nlifecycle:\n run_until: drained\n drain_timeout: 3s\nobservability:\n log_level: info\n"
))
}
#[test]
fn a_schedule_with_at_fires_exactly_once_and_not_again_after_a_restart() {
let dir = state_root("at-once");
let cfg = at_config(&dir);
let (code, first) = run_daemon_until(&cfg, &dir, "run.done", 700);
assert_eq!(code, Some(0), "life 1 drains; stderr:\n{first}");
let fired = events(&first, "start.fired");
let sched: Vec<&serde_json::Value> = fired.iter().filter(|e| e["kind"] == "schedule").collect();
assert_eq!(
sched.len(),
1,
"a one-shot `at` fires ONCE, not once per tick ({} firings):\n{first}",
sched.len()
);
assert_eq!(
events(&first, "run.start").len(),
1,
"and starts exactly one run:\n{first}"
);
let done = events(&first, "run.done");
assert_eq!(done.len(), 1, "which completes:\n{first}");
assert_eq!(done[0]["status"], "completed");
let (code, second) = run_daemon_until(&cfg, &dir, "start.schedule.done", 500);
assert_eq!(code, Some(0), "life 2 drains; stderr:\n{second}");
assert!(
events(&second, "restore.fresh").is_empty(),
"life 2 restored the first life's state:\n{second}"
);
assert!(
events(&second, "start.fired").is_empty(),
"the one-shot `at` must NOT fire again after a restart:\n{second}"
);
assert!(
events(&second, "run.start").is_empty(),
"and no run starts:\n{second}"
);
assert_eq!(
events(&second, "start.schedule.done").len(),
1,
"the consumed one-shot is reported as done, not as an invalid schedule:\n{second}"
);
let _ = std::fs::remove_dir_all(&dir);
}
const SLEEP_STEPS: &str = r#"{
"start": {"kind": "once"},
"nap": {"kind": "sleep", "depends_on": ["start"], "duration": "1s"},
"done": {"kind": "finish", "depends_on": ["nap"], "status": "completed"}
}"#;
fn sleep_config(dir: &str) -> String {
write_config(&format!(
"config_version: \"1\"\nagent:\n name: napper\nstore:\n kind: file\n file:\n path: {dir}\n checkpoint:\n debounce_ms: 0\nworkflows:\n - name: nap\n steps: {SLEEP_STEPS}\nlifecycle:\n run_until: idle\n idle_grace: 500ms\nobservability:\n log_level: info\n"
))
}
fn timer_dir(dir: &str) -> std::path::PathBuf {
std::path::Path::new(dir)
.join("agentd")
.join("napper")
.join("timer")
}
fn wait_for_timer_rows(dir: &str) -> Vec<std::path::PathBuf> {
let deadline = Instant::now() + Duration::from_secs(10);
loop {
let rows: Vec<std::path::PathBuf> = std::fs::read_dir(timer_dir(dir))
.map(|rd| rd.filter_map(Result::ok).map(|e| e.path()).collect())
.unwrap_or_default();
if !rows.is_empty() {
return rows;
}
assert!(
Instant::now() < deadline,
"the sleep step never armed a durable timer under {}",
timer_dir(dir).display()
);
std::thread::sleep(Duration::from_millis(20));
}
}
#[test]
fn a_suspended_step_whose_timer_is_gone_is_repaired_at_restore() {
let dir = state_root("orphan-timer");
let cfg = sleep_config(&dir);
let (child, err_path) = spawn_daemon(&cfg, &dir);
let pid = child.id() as i32;
let rows = wait_for_timer_rows(&dir);
unsafe { libc::kill(pid, libc::SIGKILL) };
let mut child = child;
let _ = child.wait();
let first = std::fs::read_to_string(&err_path).unwrap_or_default();
let _ = std::fs::remove_file(&err_path);
assert!(
events(&first, "run.start").len() == 1,
"life 1 started the run:\n{first}"
);
for row in &rows {
std::fs::remove_file(row).expect("delete the timer row");
}
assert!(
std::fs::read_dir(timer_dir(&dir))
.map(|rd| rd.filter_map(Result::ok).count())
.unwrap_or(0)
== 0,
"the timers are gone"
);
let (child2, err2) = spawn_daemon(&cfg, &dir);
let (code, second) = wait_bounded(child2, &err2, "life 2 (a step with no timer)");
assert_eq!(code, Some(0), "life 2 exits cleanly; stderr:\n{second}");
let repaired = events(&second, "restore.timer.repaired").len();
let replayed = events(&second, "restore.step.replay").len();
assert!(
repaired + replayed >= 1,
"restore rescued the step with no timer by one path or the other \
(repaired={repaired}, replayed={replayed}):\n{second}"
);
let done = events(&second, "run.done");
assert_eq!(
done.len(),
1,
"the run finished rather than wedging:\n{second}"
);
assert_eq!(done[0]["status"], "completed", "{second}");
let _ = std::fs::remove_dir_all(&dir);
}