use crate::loopcheck::TerminationReason;
use chrono::Utc;
use serde_json::{json, Value};
use std::fs;
use std::io::{BufRead, BufReader, Write};
use std::path::{Path, PathBuf};
pub struct ProjectJournalPath(pub PathBuf);
pub struct GlobalJournalPath(pub PathBuf);
#[derive(Debug, thiserror::Error)]
pub enum LoopError {
#[error("I/O error: {0}")]
Io(#[from] std::io::Error),
#[error("journal write failure (project): {0}")]
Journal(String),
#[error("queue error: {0}")]
Queue(String),
#[error("dispatch error: {0}")]
Dispatch(String),
#[error("configuration error: {0}")]
Config(String),
}
pub struct Unit {
pub id: String,
pub title: String,
pub session_key: String,
pub plan_path: Option<String>,
pub extra_env: Vec<(String, String)>,
}
pub struct Evidence {
pub reason: TerminationReason,
pub message: String,
}
#[derive(Debug, PartialEq)]
pub enum CloseOutcome {
Closed,
Refused(String),
Parked(String),
AwaitingMerge,
}
pub struct DispatchCtx {
pub iteration: u64,
}
pub trait Queue {
fn next(&mut self) -> Result<Option<Unit>, LoopError>;
fn close(&mut self, unit: &Unit, evidence: &Evidence) -> Result<CloseOutcome, LoopError>;
}
pub trait Session {
fn wait(&mut self) -> Result<i32, LoopError>;
fn output_tail(&self) -> Option<String> {
None
}
}
pub trait Dispatcher {
fn run(&self, unit: &Unit, ctx: &DispatchCtx) -> Result<Box<dyn Session>, LoopError>;
}
pub struct LoopBudget {
max_iterations: u64,
}
impl LoopBudget {
pub fn new(max_iterations: u64) -> Result<Self, LoopError> {
if max_iterations == 0 {
return Err(LoopError::Config("max_iterations must be > 0".to_string()));
}
Ok(Self { max_iterations })
}
}
pub struct Journal {
project_path: PathBuf,
global_path: PathBuf,
}
impl Journal {
pub fn new(project_path: ProjectJournalPath, global_path: GlobalJournalPath) -> Self {
Self {
project_path: project_path.0,
global_path: global_path.0,
}
}
pub fn new_raw(project_path: PathBuf, global_path: PathBuf) -> Self {
Self {
project_path,
global_path,
}
}
pub fn append(&self, event_type: &str, data: Value) -> Result<(), LoopError> {
let ts = Utc::now().format("%Y-%m-%dT%H:%M:%SZ").to_string();
let env = json!({
"ts": ts,
"type": event_type,
"source": "loop",
"data": data,
});
let mut line = serde_json::to_string(&env)
.map_err(|e| LoopError::Journal(format!("serialize {event_type}: {e}")))?;
line.push('\n');
self.append_to_file(&self.project_path, &line, true)?;
if self.project_path != self.global_path {
if let Err(e) = self.append_to_file(&self.global_path, &line, false) {
eprintln!("loop-runtime: global mirror write failed (non-fatal): {e}");
}
}
Ok(())
}
fn append_to_file(&self, path: &Path, line: &str, fatal: bool) -> Result<(), LoopError> {
if let Some(parent) = path.parent() {
if let Err(e) = fs::create_dir_all(parent) {
let msg = format!("create_dir_all {}: {e}", parent.display());
if fatal {
return Err(LoopError::Journal(msg));
} else {
return Err(LoopError::Io(e));
}
}
}
match fs::OpenOptions::new().create(true).append(true).open(path) {
Ok(mut f) => {
if let Err(e) = f.write_all(line.as_bytes()) {
let msg = format!("write to {}: {e}", path.display());
if fatal {
return Err(LoopError::Journal(msg));
} else {
return Err(LoopError::Io(e));
}
}
Ok(())
}
Err(e) => {
let msg = format!("open {}: {e}", path.display());
if fatal {
Err(LoopError::Journal(msg))
} else {
Err(LoopError::Io(e))
}
}
}
}
pub fn find_termination(&self, session_key: &str) -> Result<Option<Evidence>, LoopError> {
if let Some(ev) = Self::scan_journal_with_rotation(&self.project_path, session_key) {
return Ok(Some(ev));
}
if self.project_path != self.global_path {
if let Some(ev) = Self::scan_journal_with_rotation(&self.global_path, session_key) {
return Ok(Some(ev));
}
}
Ok(None)
}
pub fn find_termination_strict(
&self,
session_key: &str,
) -> Result<Option<Evidence>, LoopError> {
let mut first_error: Option<LoopError> = None;
let mut paths = vec![self.project_path.clone()];
if self.global_path != self.project_path {
paths.push(self.global_path.clone());
}
for path in paths {
for candidate in [path.clone(), rotated_journal_path(&path)] {
match Self::scan_journal_strict(&candidate, session_key) {
Ok(Some(ev)) => return Ok(Some(ev)),
Ok(None) => {}
Err(err) => {
if first_error.is_none() {
first_error = Some(err);
}
}
}
}
}
match first_error {
Some(err) => Err(err),
None => Ok(None),
}
}
fn scan_journal_with_rotation(path: &Path, session_key: &str) -> Option<Evidence> {
Self::scan_journal(path, session_key)
.or_else(|| Self::scan_journal(&rotated_journal_path(path), session_key))
}
fn scan_journal_strict(path: &Path, session_key: &str) -> Result<Option<Evidence>, LoopError> {
let file = match fs::File::open(path) {
Ok(file) => file,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(err) => {
return Err(LoopError::Journal(format!(
"could not read journal {}: {err}",
path.display()
)))
}
};
let mut last_match = None;
for line_result in BufReader::new(file).lines() {
let raw = line_result.map_err(|err| {
LoopError::Journal(format!("read journal {}: {err}", path.display()))
})?;
let line = raw.trim();
if line.is_empty() {
continue;
}
let value: Value = match serde_json::from_str(line) {
Ok(value) => value,
Err(_) => continue,
};
if value["type"].as_str() != Some("termination")
|| value["data"]["session_id"].as_str() != Some(session_key)
{
continue;
}
let reason_raw = value["data"]["reason"].as_str().ok_or_else(|| {
LoopError::Journal(format!(
"termination event for {session_key} in {} has no reason",
path.display()
))
})?;
let reason = parse_termination_reason(reason_raw).ok_or_else(|| {
LoopError::Journal(format!(
"termination event for {session_key} in {} has unknown reason {reason_raw}",
path.display()
))
})?;
let message = value["data"]["message"].as_str().unwrap_or("").to_string();
last_match = Some(Evidence { reason, message });
}
Ok(last_match)
}
fn scan_journal(path: &Path, session_key: &str) -> Option<Evidence> {
if !path.exists() {
return None;
}
let file = match fs::File::open(path) {
Ok(f) => f,
Err(e) => {
eprintln!(
"loop-runtime: could not read journal {}: {e}",
path.display()
);
return None;
}
};
let mut last_match: Option<Evidence> = None;
for line_result in BufReader::new(file).lines() {
let raw = match line_result {
Ok(l) => l,
Err(_) => continue,
};
let line = raw.trim();
if line.is_empty() {
continue;
}
let v: Value = match serde_json::from_str(line) {
Ok(v) => v,
Err(_) => continue,
};
if v["type"].as_str() != Some("termination") {
continue;
}
if v["data"]["session_id"].as_str() != Some(session_key) {
continue;
}
let reason_str = match v["data"]["reason"].as_str() {
Some(s) => s,
None => {
eprintln!(
"loop-runtime: termination event for {session_key} missing 'reason' field, skipping"
);
continue;
}
};
let reason = match parse_termination_reason(reason_str) {
Some(r) => r,
None => {
eprintln!(
"loop-runtime: unknown TerminationReason '{reason_str}' in journal, skipping"
);
continue;
}
};
let message = v["data"]["message"].as_str().unwrap_or("").to_string();
last_match = Some(Evidence { reason, message });
}
last_match
}
}
fn rotated_journal_path(path: &Path) -> PathBuf {
let mut name = path.as_os_str().to_os_string();
name.push(".1");
PathBuf::from(name)
}
fn parse_termination_reason(s: &str) -> Option<TerminationReason> {
serde_json::from_value(Value::String(s.to_string())).ok()
}
pub struct UnitResult {
pub unit_id: String,
pub evidence: Evidence,
pub close: CloseOutcome,
}
pub struct LoopOutcome {
pub reason: TerminationReason,
pub iterations_used: u64,
pub units: Vec<UnitResult>,
}
const BG_GUARD_MARKER: &str = "running as a background agent";
fn is_bg_guard_refusal(exit_code: i32, output_tail: Option<&str>) -> bool {
exit_code != 0
&& output_tail
.map(|t| t.to_ascii_lowercase().contains(BG_GUARD_MARKER))
.unwrap_or(false)
}
pub fn run_loop(
queue: &mut dyn Queue,
dispatcher: &dyn Dispatcher,
budget: &LoopBudget,
journal: &Journal,
cancel: &dyn Fn() -> bool,
per_unit_max_dispatches: Option<u64>,
) -> Result<LoopOutcome, LoopError> {
let mut iterations_used: u64 = 0;
let mut units: Vec<UnitResult> = Vec::new();
loop {
if cancel() {
journal.append(
"loop_terminated",
json!({
"reason": "Interrupted",
"iterations_used": iterations_used,
"units_closed": units.len(),
}),
)?;
return Ok(LoopOutcome {
reason: TerminationReason::Interrupted,
iterations_used,
units,
});
}
let unit = match queue.next() {
Ok(None) => {
journal.append(
"loop_terminated",
json!({
"reason": "NoWork",
"iterations_used": iterations_used,
"units_closed": units.len(),
}),
)?;
return Ok(LoopOutcome {
reason: TerminationReason::NoWork,
iterations_used,
units,
});
}
Ok(Some(u)) => u,
Err(e) => return Err(e),
};
if let Some(evidence) = journal.find_termination(&unit.session_key)? {
let close = queue.close(&unit, &evidence)?;
journal_node_closed(journal, &unit, &evidence, &close, iterations_used)?;
units.push(UnitResult {
unit_id: unit.id.clone(),
evidence,
close,
});
continue;
}
let mut unit_dispatches: u64 = 0;
loop {
if iterations_used >= budget.max_iterations {
journal.append(
"loop_terminated",
json!({
"reason": "Budget",
"iterations_used": iterations_used,
"units_closed": units.len(),
"axis": "iterations",
}),
)?;
return Ok(LoopOutcome {
reason: TerminationReason::Budget,
iterations_used,
units,
});
}
if cancel() {
journal.append(
"loop_terminated",
json!({
"reason": "Interrupted",
"iterations_used": iterations_used,
"units_closed": units.len(),
}),
)?;
return Ok(LoopOutcome {
reason: TerminationReason::Interrupted,
iterations_used,
units,
});
}
if let Some(cap) = per_unit_max_dispatches {
if unit_dispatches >= cap {
let evidence = Evidence {
reason: TerminationReason::NoProgress,
message: format!(
"no termination event after {cap} dispatch(es); unit parked"
),
};
let close = queue.close(&unit, &evidence)?;
journal_node_closed(journal, &unit, &evidence, &close, iterations_used)?;
units.push(UnitResult {
unit_id: unit.id.clone(),
evidence,
close,
});
break; }
}
iterations_used += 1;
unit_dispatches += 1;
journal.append(
"loop_unit_dispatched",
json!({
"unit_id": unit.id,
"session_id": unit.session_key,
"iteration": iterations_used,
"title": unit.title,
}),
)?;
let mut session = dispatcher
.run(
&unit,
&DispatchCtx {
iteration: iterations_used,
},
)
.map_err(|e| LoopError::Dispatch(e.to_string()))?;
let exit_code = session.wait()?;
if let Some(evidence) = journal.find_termination(&unit.session_key)? {
let close = queue.close(&unit, &evidence)?;
journal_node_closed(journal, &unit, &evidence, &close, iterations_used)?;
units.push(UnitResult {
unit_id: unit.id.clone(),
evidence,
close,
});
break; }
if is_bg_guard_refusal(exit_code, session.output_tail().as_deref()) {
let evidence = Evidence {
reason: TerminationReason::NoProgress,
message: "claude bg-guard refusal (session running as a background agent); re-dispatch halted".to_string(),
};
let close = queue.close(&unit, &evidence)?;
journal_node_closed(journal, &unit, &evidence, &close, iterations_used)?;
units.push(UnitResult {
unit_id: unit.id.clone(),
evidence,
close,
});
break; }
journal.append(
"node_failed",
json!({
"unit_id": unit.id,
"session_id": unit.session_key,
"iteration": iterations_used,
"exit_code": exit_code,
}),
)?;
}
}
}
fn journal_node_closed(
journal: &Journal,
unit: &Unit,
evidence: &Evidence,
close: &CloseOutcome,
iterations_used: u64,
) -> Result<(), LoopError> {
let (close_str, detail) = match close {
CloseOutcome::Closed => ("closed", String::new()),
CloseOutcome::Parked(s) => ("parked", s.clone()),
CloseOutcome::Refused(s) => ("refused", s.clone()),
CloseOutcome::AwaitingMerge => (
"awaiting-merge",
"PR not merged; node stays in_review, reconcile/advance close it at merge".to_string(),
),
};
let reason_str = format!("{:?}", evidence.reason);
journal.append(
"node_closed",
json!({
"unit_id": unit.id,
"session_id": unit.session_key,
"reason": reason_str,
"close": close_str,
"detail": detail,
"iterations_used": iterations_used,
}),
)
}
#[cfg(test)]
mod bg_guard_tests {
use super::is_bg_guard_refusal;
#[test]
fn refusal_marker_with_nonzero_exit_is_terminal() {
let out = "abc123 is currently running as a background agent (bg). \
Use 'claude agents' to view it, or add --fork-session.";
assert!(is_bg_guard_refusal(1, Some(out)));
assert!(is_bg_guard_refusal(
1,
Some("RUNNING AS A BACKGROUND AGENT")
));
}
#[test]
fn bare_nonzero_exit_without_marker_is_not_terminal() {
assert!(!is_bg_guard_refusal(1, Some("panic: index out of bounds")));
assert!(!is_bg_guard_refusal(1, None));
assert!(!is_bg_guard_refusal(137, Some("killed"))); }
#[test]
fn clean_exit_is_never_a_refusal_even_with_marker() {
assert!(!is_bg_guard_refusal(
0,
Some("running as a background agent")
));
}
}