use crate::paths::AgentsHome;
use serde::Serialize;
use serde_json::Value;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
const DEFAULT_WINDOW_SECS: u64 = 24 * 60 * 60;
const DEFAULT_FIRES_FLOOR: u64 = 2;
#[derive(Debug, Clone, Serialize, PartialEq)]
pub struct NeedItem {
pub kind: String,
pub session_id: String,
pub node: Option<String>,
pub name: Option<String>,
pub title: Option<String>,
pub ts: String,
pub evidence: String,
pub live: bool,
}
fn field<'a>(v: &'a Value, key: &str) -> Option<&'a Value> {
v.get("data")
.and_then(|d| d.get(key))
.or_else(|| v.get(key))
}
fn event_kind(v: &Value) -> Option<&str> {
v.get("type")
.and_then(|t| t.as_str())
.or_else(|| v.get("kind").and_then(|k| k.as_str()))
}
fn event_ts(v: &Value) -> &str {
v.get("ts").and_then(|t| t.as_str()).unwrap_or("")
}
fn str_field<'a>(v: &'a Value, key: &str) -> Option<&'a str> {
field(v, key).and_then(|f| f.as_str())
}
fn basename(path: &str) -> &str {
path.rsplit('/').next().unwrap_or(path)
}
fn to_epoch_lenient(ts: &str) -> Option<u64> {
crate::state::rfc3339_like_to_secs(ts).or_else(|| {
let secs = ts.get(..19)?;
crate::state::rfc3339_like_to_secs(&format!("{secs}Z"))
})
}
fn in_window(ts: &str, since: u64) -> bool {
to_epoch_lenient(ts).is_none_or(|secs| secs >= since)
}
#[derive(Default, Clone)]
struct LoopState {
decision: String,
intent: String,
ci: String,
pr_state: String,
reviewed: bool,
fires: u64,
ts: String,
epoch: u64,
seq: usize,
}
#[derive(Default, Clone)]
struct TermState {
reason: String,
ts: String,
epoch: u64,
seq: usize,
}
#[derive(Default)]
struct SessionAcc {
latest_loop: Option<LoopState>,
latest_term: Option<TermState>,
}
pub fn fold(events_raw: &str, ledger_raw: &str, since: u64, fires_floor: u64) -> Vec<NeedItem> {
let mut sessions: HashMap<String, SessionAcc> = HashMap::new();
let mut seq = 0usize;
let mut mail_escalations: HashMap<String, (u64, usize, NeedItem)> = HashMap::new();
for line in events_raw.lines() {
if line.trim().is_empty() {
continue;
}
let Ok(v) = serde_json::from_str::<Value>(line) else {
continue; };
let ts = event_ts(&v);
if !in_window(ts, since) {
continue;
}
let kind = event_kind(&v);
if kind == Some("mail_escalation") {
let Some(recipient) = str_field(&v, "recipient") else {
continue;
};
let reason = str_field(&v, "reason").unwrap_or("");
let sender = str_field(&v, "sender").unwrap_or("");
let summary = str_field(&v, "summary").unwrap_or("");
let epoch = to_epoch_lenient(ts).unwrap_or(0);
seq += 1;
if mail_escalations
.get(recipient)
.is_none_or(|(e, s, _)| (epoch, seq) >= (*e, *s))
{
mail_escalations.insert(
recipient.to_string(),
(
epoch,
seq,
NeedItem {
kind: "mail_question".to_string(),
session_id: recipient.to_string(),
node: None,
name: Some(recipient.to_string()),
title: None,
ts: ts.to_string(),
evidence: format!("{reason}: {sender} -> {recipient}: {summary}"),
live: false,
},
),
);
}
continue;
}
let Some(sid) = str_field(&v, "session_id") else {
continue; };
if !matches!(
kind,
Some("loop_check") | Some("termination") | Some("loop_terminated")
) {
continue;
}
let epoch = to_epoch_lenient(ts).unwrap_or(0);
seq += 1;
let acc = sessions.entry(sid.to_string()).or_default();
match kind {
Some("loop_check") => {
if acc
.latest_loop
.as_ref()
.is_none_or(|c| (epoch, seq) >= (c.epoch, c.seq))
{
acc.latest_loop = Some(LoopState {
decision: str_field(&v, "decision").unwrap_or("").to_string(),
intent: str_field(&v, "intent").unwrap_or("").to_string(),
ci: str_field(&v, "ci").unwrap_or("").to_string(),
pr_state: str_field(&v, "pr_state").unwrap_or("").to_string(),
reviewed: field(&v, "reviewed")
.and_then(|r| r.as_bool())
.unwrap_or(false),
fires: field(&v, "fires").and_then(|f| f.as_u64()).unwrap_or(0),
ts: ts.to_string(),
epoch,
seq,
});
}
}
_ => {
if acc
.latest_term
.as_ref()
.is_none_or(|c| (epoch, seq) >= (c.epoch, c.seq))
{
acc.latest_term = Some(TermState {
reason: str_field(&v, "reason").unwrap_or("").to_string(),
ts: ts.to_string(),
epoch,
seq,
});
}
}
}
}
let ledger = LedgerIndex::parse(ledger_raw);
let mut items: Vec<NeedItem> = Vec::new();
for (sid, acc) in &sessions {
if let Some((kind, ts, evidence)) = classify(acc, fires_floor) {
let (node, name, title) = ledger.resolve(sid);
items.push(NeedItem {
kind: kind.to_string(),
session_id: sid.clone(),
node,
name,
title,
ts,
evidence,
live: false, });
}
}
for (_, (_, _, item)) in mail_escalations {
items.push(item);
}
items.sort_by(|a, b| {
a.ts.cmp(&b.ts)
.then_with(|| a.session_id.cmp(&b.session_id))
});
items
}
fn classify(acc: &SessionAcc, fires_floor: u64) -> Option<(&'static str, String, String)> {
let terminated = match (&acc.latest_term, &acc.latest_loop) {
(Some(t), Some(l)) => (t.epoch, t.seq) >= (l.epoch, l.seq),
(Some(_), None) => true,
(None, _) => false,
};
if terminated {
let t = acc.latest_term.as_ref()?;
return match t.reason.as_str() {
"Budget" | "NoProgress" => Some((
"budget_stop",
t.ts.clone(),
format!("loop stopped: {}", t.reason),
)),
_ => None,
};
}
let l = acc.latest_loop.as_ref()?;
let wedged = l.decision == "block"
&& l.intent == "promise"
&& l.ci == "SUCCESS"
&& l.pr_state == "OPEN"
&& !l.reviewed
&& l.fires >= fires_floor;
if wedged {
return Some((
"review_wedged",
l.ts.clone(),
format!("green PR wedged on review ({} checks)", l.fires),
));
}
None
}
struct LedgerIndex {
entries: Vec<Value>,
}
impl LedgerIndex {
fn parse(ledger_raw: &str) -> Self {
let entries = serde_json::from_str::<Value>(ledger_raw)
.ok()
.and_then(|root| {
root.get("entries")
.and_then(|e| e.as_array())
.or_else(|| root.as_array())
.cloned()
})
.unwrap_or_default();
LedgerIndex { entries }
}
fn entry_has_session(entry: &Value, sid: &str) -> bool {
if entry.get("session_id").and_then(|s| s.as_str()) == Some(sid) {
return true;
}
entry
.get("sessions")
.and_then(|s| s.as_array())
.is_some_and(|arr| arr.iter().any(|s| s.as_str() == Some(sid)))
}
fn resolve(&self, sid: &str) -> (Option<String>, Option<String>, Option<String>) {
let Some(entry) = self
.entries
.iter()
.find(|e| Self::entry_has_session(e, sid))
else {
return (None, None, None);
};
let node = entry
.get("graph_node_id")
.and_then(|v| v.as_str())
.map(str::to_string);
let title = entry
.get("title")
.and_then(|v| v.as_str())
.map(str::to_string);
let name = ["worktree", "root_path"]
.iter()
.find_map(|k| entry.get(*k).and_then(|v| v.as_str()))
.map(|p| basename(p).to_string())
.or_else(|| node.clone());
(node, name, title)
}
}
struct NeedsArgs {
since_epoch: Option<u64>,
fires_floor: u64,
json: bool,
events_override: Vec<PathBuf>,
ledger_override: Option<PathBuf>,
}
fn parse_args(rest: &[String]) -> Result<NeedsArgs, String> {
let mut since_epoch: Option<u64> = None;
let mut fires_floor = DEFAULT_FIRES_FLOOR;
let mut json = false;
let mut events_override: Vec<PathBuf> = Vec::new();
let mut ledger_override: Option<PathBuf> = None;
let mut it = expand_eq(rest).into_iter();
while let Some(a) = it.next() {
match a.as_str() {
"--since-epoch" => {
since_epoch = Some(
it.next()
.and_then(|v| v.parse::<u64>().ok())
.ok_or("--since-epoch needs a non-negative integer")?,
)
}
"--fires-floor" => {
fires_floor = it
.next()
.and_then(|v| v.parse::<u64>().ok())
.ok_or("--fires-floor needs a non-negative integer")?
}
"--json" | "-J" => json = true,
"--events" => {
events_override.push(PathBuf::from(it.next().ok_or("--events needs a path")?))
}
"--ledger" => {
ledger_override = Some(PathBuf::from(it.next().ok_or("--ledger needs a path")?))
}
other => return Err(format!("unknown needs flag: {other}")),
}
}
Ok(NeedsArgs {
since_epoch,
fires_floor,
json,
events_override,
ledger_override,
})
}
fn expand_eq(rest: &[String]) -> Vec<String> {
let mut out = Vec::with_capacity(rest.len());
for a in rest {
if let Some(eq) = a.find('=') {
if a.starts_with("--") && eq > 2 {
out.push(a[..eq].to_string());
out.push(a[eq + 1..].to_string());
continue;
}
}
out.push(a.clone());
}
out
}
fn default_sources(home: &AgentsHome) -> (Vec<PathBuf>, PathBuf) {
let fno_dir = home
.root()
.parent()
.map(Path::to_path_buf)
.unwrap_or_else(|| PathBuf::from(".fno"));
let global_events = fno_dir.join("events.jsonl");
let project_events = PathBuf::from(".fno").join("events.jsonl");
let ledger = fno_dir.join("ledger.json");
(vec![project_events, global_events], ledger)
}
fn stamp_liveness(mut items: Vec<NeedItem>) -> Vec<NeedItem> {
for item in &mut items {
if item.kind == "mail_question" {
item.live = true;
continue;
}
item.live = item.node.as_deref().is_some_and(|n| {
let (state, _) = crate::claims::status(&format!("node:{n}"), None);
matches!(
state,
crate::claims::ClaimState::Live | crate::claims::ClaimState::Suspect
)
});
}
items
}
fn now_secs() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
pub async fn run_needs(rest: &[String], home: &AgentsHome) -> i32 {
let args = match parse_args(rest) {
Ok(a) => a,
Err(msg) => {
eprintln!("fno-agents: {msg}");
return 2;
}
};
let (default_events, default_ledger) = default_sources(home);
let event_paths = if args.events_override.is_empty() {
default_events
} else {
args.events_override
};
let ledger_path = args.ledger_override.unwrap_or(default_ledger);
let mut events_raw = String::new();
for p in &event_paths {
if let Ok(content) = std::fs::read_to_string(p) {
events_raw.push_str(&content);
if !content.ends_with('\n') {
events_raw.push('\n');
}
}
}
let ledger_raw = std::fs::read_to_string(&ledger_path).unwrap_or_default();
let since = args
.since_epoch
.unwrap_or_else(|| now_secs().saturating_sub(DEFAULT_WINDOW_SECS));
let items = stamp_liveness(fold(&events_raw, &ledger_raw, since, args.fires_floor));
if args.json {
println!(
"{}",
serde_json::to_string(&items).expect("serializing an owned value never fails")
);
} else {
for item in &items {
let name = item.name.as_deref().unwrap_or(&item.session_id);
println!("{} {} - {}", item.kind, name, item.evidence);
}
}
0
}
#[cfg(test)]
mod tests {
use super::*;
fn loop_check(
ts: &str,
session: &str,
decision: &str,
ci: &str,
pr_state: &str,
reviewed: bool,
fires: u64,
) -> String {
loop_check_i(
ts, session, decision, "promise", ci, pr_state, reviewed, fires,
)
}
#[allow(clippy::too_many_arguments)]
fn loop_check_i(
ts: &str,
session: &str,
decision: &str,
intent: &str,
ci: &str,
pr_state: &str,
reviewed: bool,
fires: u64,
) -> String {
format!(
r#"{{"ts":"{ts}","type":"loop_check","source":"hook","data":{{"session_id":"{session}","decision":"{decision}","intent":"{intent}","ci":"{ci}","pr_state":"{pr_state}","reviewed":{reviewed},"fires":{fires}}}}}"#
)
}
fn termination(ts: &str, session: &str, reason: &str) -> String {
format!(
r#"{{"ts":"{ts}","type":"termination","source":"hook","data":{{"session_id":"{session}","reason":"{reason}"}}}}"#
)
}
const ALL: u64 = 0;
#[test]
fn green_open_unreviewed_block_is_review_wedged() {
let events = loop_check(
"2026-07-03T02:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
5,
);
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items.len(), 1);
assert_eq!(items[0].kind, "review_wedged");
assert_eq!(items[0].session_id, "s");
assert!(items[0].evidence.contains("5 checks"));
}
#[test]
fn budget_termination_is_budget_stop() {
let events = termination("2026-07-03T02:00:00Z", "s", "Budget");
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items.len(), 1);
assert_eq!(items[0].kind, "budget_stop");
assert!(items[0].evidence.contains("Budget"));
}
#[test]
fn noprogress_termination_is_budget_stop() {
let events = termination("2026-07-03T02:00:00Z", "s", "NoProgress");
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items[0].kind, "budget_stop");
}
#[test]
fn done_pr_green_termination_yields_nothing() {
let events = termination("2026-07-03T02:00:00Z", "s", "DonePRGreen");
assert!(fold(&events, "", ALL, DEFAULT_FIRES_FLOOR).is_empty());
}
fn mail_escalation(
ts: &str,
reason: &str,
sender: &str,
recipient: &str,
summary: &str,
) -> String {
format!(
r#"{{"ts":"{ts}","type":"mail_escalation","source":"target","data":{{"reason":"{reason}","sender":"{sender}","recipient":"{recipient}","summary":"{summary}"}}}}"#
)
}
#[test]
fn mail_escalation_folds_to_mail_question_without_a_session() {
let events = mail_escalation(
"2026-07-03T02:00:00Z",
"question",
"etl",
"web",
"which schema?",
);
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items.len(), 1);
assert_eq!(items[0].kind, "mail_question");
assert_eq!(items[0].name.as_deref(), Some("web"));
assert_eq!(items[0].session_id, "web");
assert!(items[0].evidence.contains("question"));
assert!(items[0].evidence.contains("etl -> web"));
}
#[test]
fn mail_escalation_is_stamped_always_live_with_no_node() {
let events = mail_escalation(
"2026-07-03T02:00:00Z",
"attended-miss",
"ops",
"claude-9a06",
"need you",
);
let items = stamp_liveness(fold(&events, "", ALL, DEFAULT_FIRES_FLOOR));
assert_eq!(items.len(), 1);
assert_eq!(items[0].node, None);
assert!(
items[0].live,
"mail_question is always-live even with no node"
);
}
#[test]
fn mail_escalation_latest_per_recipient_wins() {
let events = format!(
"{}\n{}\n",
mail_escalation("2026-07-03T02:00:00Z", "question", "etl", "web", "old"),
mail_escalation("2026-07-03T03:00:00Z", "attended-miss", "ops", "web", "new"),
);
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items.len(), 1, "one row per recipient");
assert!(
items[0].evidence.contains("new"),
"latest (epoch, seq) wins"
);
}
#[test]
fn same_second_rearm_after_termination_is_not_terminated() {
let events = [
termination("2026-07-03T02:00:00Z", "s", "Budget"),
loop_check(
"2026-07-03T02:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
9,
),
]
.join("\n");
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items.len(), 1);
assert_eq!(items[0].kind, "review_wedged");
}
#[test]
fn newer_state_survives_older_line_from_a_later_source() {
let events = [
loop_check(
"2026-07-03T05:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
9,
),
loop_check(
"2026-07-03T01:00:00Z",
"s",
"allow",
"SUCCESS",
"OPEN",
true,
3,
),
]
.join("\n");
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(
items.len(),
1,
"the newer block state survives the older line"
);
assert_eq!(items[0].kind, "review_wedged");
}
#[test]
fn intent_none_block_is_not_wedged() {
let events = loop_check_i(
"2026-07-03T02:00:00Z",
"s",
"block",
"none",
"SUCCESS",
"OPEN",
false,
9,
);
assert!(fold(&events, "", ALL, DEFAULT_FIRES_FLOOR).is_empty());
}
#[test]
fn merged_pr_block_is_not_wedged() {
let events = loop_check(
"2026-07-03T02:00:00Z",
"s",
"block",
"SUCCESS",
"MERGED",
false,
144,
);
assert!(fold(&events, "", ALL, DEFAULT_FIRES_FLOOR).is_empty());
}
#[test]
fn later_allow_clears_the_wedge() {
let events = [
loop_check(
"2026-07-03T02:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
5,
),
loop_check(
"2026-07-03T03:00:00Z",
"s",
"allow",
"SUCCESS",
"OPEN",
true,
5,
),
]
.join("\n");
assert!(fold(&events, "", ALL, DEFAULT_FIRES_FLOOR).is_empty());
}
#[test]
fn termination_after_wedge_wins() {
let events = [
loop_check(
"2026-07-03T02:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
5,
),
termination("2026-07-03T03:00:00Z", "s", "DonePRGreen"),
]
.join("\n");
assert!(fold(&events, "", ALL, DEFAULT_FIRES_FLOOR).is_empty());
}
#[test]
fn wedge_after_a_stale_budget_stop_reads_as_wedge() {
let events = [
termination("2026-07-03T02:00:00Z", "s", "Budget"),
loop_check(
"2026-07-03T03:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
9,
),
]
.join("\n");
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items[0].kind, "review_wedged");
}
#[test]
fn termination_with_fractional_ts_still_wins_over_z_loop_check() {
let events = [
loop_check(
"2026-07-03T02:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
5,
),
termination("2026-07-03T02:00:00.5", "s", "Budget"),
]
.join("\n");
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items.len(), 1);
assert_eq!(
items[0].kind, "budget_stop",
"the termination wins despite its fractional ts"
);
}
#[test]
fn fires_below_floor_is_not_wedged() {
let events = loop_check(
"2026-07-03T02:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
1,
);
assert!(fold(&events, "", ALL, 2).is_empty());
}
#[test]
fn since_window_excludes_old_events() {
let events = loop_check(
"2026-07-03T02:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
5,
);
let future = crate::state::rfc3339_like_to_secs("2099-01-01T00:00:00Z").unwrap();
assert!(fold(&events, "", future, DEFAULT_FIRES_FLOOR).is_empty());
}
#[test]
fn malformed_line_is_skipped_not_aborted() {
let events = [
"{ this is not valid json".to_string(),
loop_check(
"2026-07-03T02:00:00Z",
"s",
"block",
"SUCCESS",
"OPEN",
false,
5,
),
]
.join("\n");
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items.len(), 1, "the good line still folds");
}
#[test]
fn one_item_per_session_latest_wins() {
let events = [
loop_check(
"2026-07-03T02:00:00Z",
"a",
"block",
"SUCCESS",
"OPEN",
false,
5,
),
termination("2026-07-03T02:30:00Z", "b", "Budget"),
]
.join("\n");
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items.len(), 2);
assert_eq!(items[0].kind, "review_wedged");
assert_eq!(items[1].kind, "budget_stop");
}
#[test]
fn ledger_resolves_node_name_title() {
let events = loop_check(
"2026-07-03T02:00:00Z",
"sess-x",
"block",
"SUCCESS",
"OPEN",
false,
5,
);
let ledger = r#"{"entries":[{"session_id":"sess-x","graph_node_id":"x-feec","title":"needs queue","worktree":"/w/footnote/x-feec"}]}"#;
let items = fold(&events, ledger, ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items[0].node.as_deref(), Some("x-feec"));
assert_eq!(items[0].name.as_deref(), Some("x-feec"));
assert_eq!(items[0].title.as_deref(), Some("needs queue"));
}
#[test]
fn ledger_resolves_via_sessions_array() {
let events = termination("2026-07-03T02:00:00Z", "fno-sess", "Budget");
let ledger = r#"[{"sessions":["uuid-1","fno-sess"],"graph_node_id":"x-1","worktree":"/w/footnote/x-1"}]"#;
let items = fold(&events, ledger, ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items[0].node.as_deref(), Some("x-1"));
}
#[test]
fn unresolved_session_renders_id_only() {
let events = termination("2026-07-03T02:00:00Z", "ghost", "Budget");
let items = fold(&events, "", ALL, DEFAULT_FIRES_FLOOR);
assert_eq!(items[0].node, None);
assert_eq!(items[0].name, None);
assert_eq!(items[0].session_id, "ghost");
}
}