use fno_agents::loop_runtime::{
run_loop, CloseOutcome, DispatchCtx, Dispatcher, Evidence, Journal, LoopBudget, LoopError,
Queue, Session, Unit,
};
use fno_agents::loopcheck::TerminationReason;
use std::fs;
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use tempfile::TempDir;
fn make_unit(id: &str, session_key: &str) -> Unit {
Unit {
id: id.to_string(),
title: format!("Test unit {id}"),
session_key: session_key.to_string(),
plan_path: None,
extra_env: vec![],
}
}
fn make_evidence(reason: TerminationReason) -> Evidence {
Evidence {
reason,
message: "test done".to_string(),
}
}
fn seed_termination_event(journal_path: &Path, session_key: &str, reason: &str) {
let line = format!(
"{{\"ts\":\"2026-06-06T00:00:00Z\",\"type\":\"termination\",\"source\":\"hook\",\
\"data\":{{\"session_id\":\"{session_key}\",\"reason\":\"{reason}\",\"message\":\"pre-seeded\"}}}}\n"
);
let mut f = fs::OpenOptions::new()
.create(true)
.append(true)
.open(journal_path)
.expect("seed termination event");
f.write_all(line.as_bytes()).expect("write seed");
}
fn read_jsonl(path: &Path) -> Vec<serde_json::Value> {
if !path.exists() {
return vec![];
}
let content = fs::read_to_string(path).expect("read jsonl");
content
.lines()
.filter(|l| !l.trim().is_empty())
.filter_map(|l| serde_json::from_str(l).ok())
.collect()
}
fn count_events(path: &Path, event_type: &str) -> usize {
read_jsonl(path)
.into_iter()
.filter(|v| v["type"].as_str() == Some(event_type))
.count()
}
struct FixedQueue {
units: Mutex<Vec<Unit>>,
closed: Mutex<Vec<String>>,
close_outcome: CloseOutcome,
}
impl FixedQueue {
fn new(units: Vec<Unit>) -> Self {
Self {
units: Mutex::new(units),
closed: Mutex::new(vec![]),
close_outcome: CloseOutcome::Closed,
}
}
fn with_close_outcome(units: Vec<Unit>, outcome: CloseOutcome) -> Self {
Self {
units: Mutex::new(units),
closed: Mutex::new(vec![]),
close_outcome: outcome,
}
}
fn close_count(&self) -> usize {
self.closed.lock().unwrap().len()
}
}
impl Queue for FixedQueue {
fn next(&mut self) -> Result<Option<Unit>, LoopError> {
let mut units = self.units.lock().unwrap();
if units.is_empty() {
Ok(None)
} else {
Ok(Some(units.remove(0)))
}
}
fn close(&mut self, unit: &Unit, _evidence: &Evidence) -> Result<CloseOutcome, LoopError> {
self.closed.lock().unwrap().push(unit.id.clone());
Ok(match &self.close_outcome {
CloseOutcome::Closed => CloseOutcome::Closed,
CloseOutcome::Refused(s) => CloseOutcome::Refused(s.clone()),
CloseOutcome::Parked(s) => CloseOutcome::Parked(s.clone()),
})
}
}
struct PanicDispatcher;
impl Dispatcher for PanicDispatcher {
fn run(&self, _unit: &Unit, _ctx: &DispatchCtx) -> Result<Box<dyn Session>, LoopError> {
panic!("PanicDispatcher: dispatch must not be called in this test");
}
}
struct MockSession {
journal_path: PathBuf,
session_key: String,
termination_reason: Option<String>, exit_code: i32,
}
impl Session for MockSession {
fn wait(&mut self) -> Result<i32, LoopError> {
if let Some(reason) = &self.termination_reason {
seed_termination_event(&self.journal_path, &self.session_key, reason);
}
Ok(self.exit_code)
}
}
struct MockDispatcher {
journal_path: PathBuf,
responses: Mutex<Vec<(Option<String>, i32)>>,
dispatch_count: AtomicU64,
}
impl MockDispatcher {
fn new(journal_path: PathBuf, responses: Vec<(Option<String>, i32)>) -> Self {
Self {
journal_path,
responses: Mutex::new(responses),
dispatch_count: AtomicU64::new(0),
}
}
fn count(&self) -> u64 {
self.dispatch_count.load(Ordering::SeqCst)
}
}
impl Dispatcher for MockDispatcher {
fn run(&self, unit: &Unit, _ctx: &DispatchCtx) -> Result<Box<dyn Session>, LoopError> {
self.dispatch_count.fetch_add(1, Ordering::SeqCst);
let mut responses = self.responses.lock().unwrap();
let (reason, exit_code) = if responses.len() > 1 {
responses.remove(0)
} else {
responses[0].clone()
};
Ok(Box::new(MockSession {
journal_path: self.journal_path.clone(),
session_key: unit.session_key.clone(),
termination_reason: reason,
exit_code,
}))
}
}
#[test]
fn empty_queue_terminates_nowork_without_dispatch() {
let dir = TempDir::new().unwrap();
let project_events = dir.path().join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
let mut queue = FixedQueue::new(vec![]);
let dispatcher = PanicDispatcher;
let budget = LoopBudget::new(10).unwrap();
let journal = Journal::new_raw(project_events.clone(), global_events);
let outcome = run_loop(&mut queue, &dispatcher, &budget, &journal, &|| false, None).unwrap();
assert_eq!(outcome.reason, TerminationReason::NoWork);
assert_eq!(outcome.iterations_used, 0);
assert!(outcome.units.is_empty());
let events = read_jsonl(&project_events);
let terminated: Vec<_> = events
.iter()
.filter(|v| v["type"].as_str() == Some("loop_terminated"))
.collect();
assert_eq!(
terminated.len(),
1,
"expected exactly one loop_terminated event"
);
assert_eq!(
terminated[0]["data"]["reason"].as_str(),
Some("NoWork"),
"loop_terminated reason must be NoWork"
);
}
#[test]
fn unit_closes_on_termination_event() {
let dir = TempDir::new().unwrap();
let project_events = dir.path().join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
let unit = make_unit("ab-001", "sess-001");
let mut queue = FixedQueue::new(vec![unit]);
let dispatcher = MockDispatcher::new(
project_events.clone(),
vec![(Some("DonePRGreen".to_string()), 0)],
);
let budget = LoopBudget::new(10).unwrap();
let journal = Journal::new_raw(project_events.clone(), global_events);
let outcome = run_loop(&mut queue, &dispatcher, &budget, &journal, &|| false, None).unwrap();
assert_eq!(outcome.reason, TerminationReason::NoWork);
assert_eq!(outcome.iterations_used, 1);
assert_eq!(outcome.units.len(), 1);
assert_eq!(
outcome.units[0].evidence.reason,
TerminationReason::DonePRGreen
);
assert_eq!(outcome.units[0].close, CloseOutcome::Closed);
assert_eq!(queue.close_count(), 1);
assert_eq!(
count_events(&project_events, "loop_unit_dispatched"),
1,
"expected one loop_unit_dispatched"
);
assert_eq!(
count_events(&project_events, "node_failed"),
0,
"expected no node_failed"
);
}
#[test]
fn no_event_exit_emits_node_failed_then_redispatches() {
let dir = TempDir::new().unwrap();
let project_events = dir.path().join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
let unit = make_unit("ab-002", "sess-002");
let mut queue = FixedQueue::new(vec![unit]);
let dispatcher = MockDispatcher::new(
project_events.clone(),
vec![(None, 1), (Some("DonePRGreen".to_string()), 0)],
);
let budget = LoopBudget::new(10).unwrap();
let journal = Journal::new_raw(project_events.clone(), global_events);
let outcome = run_loop(&mut queue, &dispatcher, &budget, &journal, &|| false, None).unwrap();
assert_eq!(outcome.reason, TerminationReason::NoWork);
assert_eq!(outcome.iterations_used, 2);
assert_eq!(
outcome.units[0].evidence.reason,
TerminationReason::DonePRGreen
);
assert_eq!(dispatcher.count(), 2, "expected exactly 2 dispatches");
assert_eq!(
count_events(&project_events, "node_failed"),
1,
"expected exactly one node_failed"
);
assert_eq!(
count_events(&project_events, "loop_unit_dispatched"),
2,
"expected two loop_unit_dispatched events"
);
}
#[test]
fn iteration_ceiling_terminates_budget() {
let dir = TempDir::new().unwrap();
let project_events = dir.path().join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
let unit = make_unit("ab-003", "sess-003");
let mut queue = FixedQueue::new(vec![unit]);
let dispatcher = MockDispatcher::new(project_events.clone(), vec![(None, 1)]);
let budget = LoopBudget::new(3).unwrap();
let journal = Journal::new_raw(project_events.clone(), global_events);
let outcome = run_loop(&mut queue, &dispatcher, &budget, &journal, &|| false, None).unwrap();
assert_eq!(outcome.reason, TerminationReason::Budget);
assert_eq!(
outcome.iterations_used, 3,
"should use exactly budget iterations"
);
assert_eq!(dispatcher.count(), 3, "expected exactly 3 dispatches");
assert_eq!(
count_events(&project_events, "node_failed"),
3,
"expected 3 node_failed"
);
let events = read_jsonl(&project_events);
let terminated: Vec<_> = events
.iter()
.filter(|v| v["type"].as_str() == Some("loop_terminated"))
.collect();
assert_eq!(terminated.len(), 1);
assert_eq!(terminated[0]["data"]["reason"].as_str(), Some("Budget"));
assert_eq!(
terminated[0]["data"]["axis"].as_str(),
Some("iterations"),
"Budget termination must carry axis=iterations"
);
}
#[test]
fn preexisting_termination_skips_dispatch() {
let dir = TempDir::new().unwrap();
let project_events = dir.path().join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
seed_termination_event(&project_events, "sess-004", "DoneAdvisory");
let unit = make_unit("ab-004", "sess-004");
let mut queue = FixedQueue::new(vec![unit]);
let dispatcher = PanicDispatcher;
let budget = LoopBudget::new(5).unwrap();
let journal = Journal::new_raw(project_events.clone(), global_events);
let outcome = run_loop(&mut queue, &dispatcher, &budget, &journal, &|| false, None).unwrap();
assert_eq!(
outcome.iterations_used, 0,
"resume guard: no dispatch iterations"
);
assert_eq!(outcome.units.len(), 1);
assert_eq!(
outcome.units[0].evidence.reason,
TerminationReason::DoneAdvisory
);
assert_eq!(outcome.units[0].close, CloseOutcome::Closed);
assert_eq!(queue.close_count(), 1);
assert_eq!(
count_events(&project_events, "loop_unit_dispatched"),
0,
"resume guard: no dispatch events"
);
}
#[test]
fn cancel_returns_interrupted() {
let dir = TempDir::new().unwrap();
let project_events = dir.path().join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
let mut queue = FixedQueue::new(vec![make_unit("ab-005", "sess-005")]);
let dispatcher = PanicDispatcher;
let budget = LoopBudget::new(10).unwrap();
let journal = Journal::new_raw(project_events.clone(), global_events);
let cancel_flag = Arc::new(AtomicBool::new(true)); let flag = cancel_flag.clone();
let outcome = run_loop(
&mut queue,
&dispatcher,
&budget,
&journal,
&move || flag.load(Ordering::SeqCst),
None,
)
.unwrap();
assert_eq!(outcome.reason, TerminationReason::Interrupted);
let events = read_jsonl(&project_events);
let terminated: Vec<_> = events
.iter()
.filter(|v| v["type"].as_str() == Some("loop_terminated"))
.collect();
assert_eq!(terminated.len(), 1);
assert_eq!(
terminated[0]["data"]["reason"].as_str(),
Some("Interrupted")
);
}
#[test]
fn zero_max_iterations_rejected() {
let result = LoopBudget::new(0);
assert!(result.is_err(), "LoopBudget::new(0) must return Err");
}
#[test]
fn journal_write_failure_is_fatal() {
let dir = TempDir::new().unwrap();
let blocker = dir.path().join("blocker");
fs::write(&blocker, b"not a dir").unwrap();
let project_events = blocker.join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
let mut queue = FixedQueue::new(vec![make_unit("ab-006", "sess-006")]);
let dispatcher = PanicDispatcher; let budget = LoopBudget::new(5).unwrap();
let journal = Journal::new_raw(project_events, global_events);
let result = run_loop(&mut queue, &dispatcher, &budget, &journal, &|| false, None);
assert!(
result.is_err(),
"journal write failure to project path must be fatal (Err)"
);
}
#[test]
fn envelope_shape() {
let dir = TempDir::new().unwrap();
let project_events = dir.path().join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
{
let mut f = fs::OpenOptions::new()
.create(true)
.append(true)
.open(&project_events)
.unwrap();
f.write_all(b"not json at all\n").unwrap();
f.write_all(b"{\"type\":\"other\",\"data\":{}}\n").unwrap();
}
let unit = make_unit("ab-007", "sess-007");
let mut queue = FixedQueue::new(vec![unit]);
let dispatcher = MockDispatcher::new(
project_events.clone(),
vec![(Some("DonePRGreen".to_string()), 0)],
);
let budget = LoopBudget::new(5).unwrap();
let journal = Journal::new_raw(project_events.clone(), global_events.clone());
run_loop(&mut queue, &dispatcher, &budget, &journal, &|| false, None).unwrap();
let content = fs::read_to_string(&project_events).unwrap();
let runtime_lines: Vec<serde_json::Value> = content
.lines()
.filter(|l| !l.trim().is_empty())
.filter_map(|l| serde_json::from_str::<serde_json::Value>(l).ok())
.filter(|v| v["source"].as_str() == Some("loop"))
.collect();
assert!(
!runtime_lines.is_empty(),
"runtime must have written at least one loop-source event"
);
for line in &runtime_lines {
let ts = line["ts"].as_str().expect("ts must be a string");
assert!(ts.ends_with('Z'), "ts must end with Z: {ts}");
ts.parse::<chrono::DateTime<chrono::Utc>>()
.expect("ts must be valid RFC3339");
assert!(line["type"].is_string(), "type must be a string: {line}");
assert_eq!(
line["source"].as_str(),
Some("loop"),
"source must be 'loop': {line}"
);
assert!(line["data"].is_object(), "data must be an object: {line}");
}
let global_content = fs::read_to_string(&global_events).unwrap_or_default();
let global_runtime_lines: Vec<serde_json::Value> = global_content
.lines()
.filter(|l| !l.trim().is_empty())
.filter_map(|l| serde_json::from_str(l).ok())
.filter(|v: &serde_json::Value| v["source"].as_str() == Some("loop"))
.collect();
assert!(
!global_runtime_lines.is_empty(),
"global mirror must also have loop-source events"
);
}
#[test]
fn termination_event_missing_reason_skips_no_panic() {
let dir = TempDir::new().unwrap();
let project_events = dir.path().join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
{
let mut f = fs::OpenOptions::new()
.create(true)
.append(true)
.open(&project_events)
.unwrap();
f.write_all(
b"{\"ts\":\"2026-06-06T00:00:00Z\",\"type\":\"termination\",\"source\":\"hook\",\
\"data\":{\"session_id\":\"sess-f1\",\"message\":\"corrupt\"}}\n",
)
.unwrap();
}
let unit = make_unit("ab-f1", "sess-f1");
let mut queue = FixedQueue::new(vec![unit]);
let dispatcher = MockDispatcher::new(project_events.clone(), vec![(None, 1)]);
let budget = LoopBudget::new(1).unwrap();
let journal = Journal::new_raw(project_events.clone(), global_events);
let outcome = run_loop(&mut queue, &dispatcher, &budget, &journal, &|| false, None).unwrap();
assert_eq!(
outcome.reason,
fno_agents::loopcheck::TerminationReason::Budget,
"corrupt termination event must not match; walk terminates with Budget"
);
}
#[test]
fn parse_termination_reason_serde_roundtrip() {
let dir = TempDir::new().unwrap();
let project_events = dir.path().join("events.jsonl");
let global_events = dir.path().join("global-events.jsonl");
for (reason_str, expected) in &[
(
"DonePRGreen",
fno_agents::loopcheck::TerminationReason::DonePRGreen,
),
(
"DoneAdvisory",
fno_agents::loopcheck::TerminationReason::DoneAdvisory,
),
("NoWork", fno_agents::loopcheck::TerminationReason::NoWork),
("Budget", fno_agents::loopcheck::TerminationReason::Budget),
(
"NoProgress",
fno_agents::loopcheck::TerminationReason::NoProgress,
),
(
"Interrupted",
fno_agents::loopcheck::TerminationReason::Interrupted,
),
("Aborted", fno_agents::loopcheck::TerminationReason::Aborted),
] {
let _ = fs::remove_file(&project_events);
seed_termination_event(&project_events, "sess-serde", reason_str);
let journal = Journal::new_raw(project_events.clone(), global_events.clone());
let evidence = journal
.find_termination("sess-serde")
.expect("find_termination must not error")
.expect(&format!("must find termination for reason '{reason_str}'"));
assert_eq!(
&evidence.reason, expected,
"serde parse must round-trip '{reason_str}'"
);
}
let _ = fs::remove_file(&project_events);
{
let mut f = fs::OpenOptions::new()
.create(true)
.append(true)
.open(&project_events)
.unwrap();
f.write_all(
b"{\"ts\":\"2026-06-06T00:00:00Z\",\"type\":\"termination\",\"source\":\"hook\",\
\"data\":{\"session_id\":\"sess-serde\",\"reason\":\"UnknownFuture\",\"message\":\"\"}}\n",
)
.unwrap();
}
let journal = Journal::new_raw(project_events.clone(), global_events.clone());
let result = journal
.find_termination("sess-serde")
.expect("find_termination must not error");
assert!(
result.is_none(),
"unknown reason string must not match (returns None)"
);
}