diff --git a/hive-forge-notify/src/bin/hive-forge-notify/main.rs b/hive-forge-notify/src/bin/hive-forge-notify/main.rs index f9717ed6..eff5a5a7 100644 --- a/hive-forge-notify/src/bin/hive-forge-notify/main.rs +++ b/hive-forge-notify/src/bin/hive-forge-notify/main.rs @@ -133,34 +133,78 @@ async fn forgejo_loop(state_dir: String, socket: std::path::PathBuf) { /// Count of open issues or PRs assigned to this agent, via Forgejo's /// global search API. `issue_type` is `"issues"` or `"pulls"`. Reads the /// `X-Total-Count` header rather than deserialising the body — only page 1 -/// with limit 1 is fetched, keeping the request cheap. `None` on any error. +/// with limit 1 is fetched, keeping the request cheap. `Err` carries a +/// human-readable cause for the caller to log — a poll failure must not +/// look like a genuine zero count (see [`decide_rollup`]). async fn count_assigned( client: &reqwest::Client, forge_url: &str, token: &str, issue_type: &str, -) -> Option { +) -> Result { let url = format!( "{forge_url}/api/v1/issues/search\ ?type={issue_type}&state=open&assigned=true&limit=1&page=1" ); - let resp = match client + let resp = client .get(&url) .header("Authorization", format!("token {token}")) .send() .await - { - Ok(r) if r.status().is_success() => r, - _ => return None, + .map_err(|e| e.to_string())?; + if !resp.status().is_success() { + return Err(format!("HTTP {}", resp.status())); + } + let count_str = resp + .headers() + .get("x-total-count") + .ok_or("response carried no X-Total-Count header")? + .to_str() + .map_err(|e| e.to_string())?; + count_str.parse::().map_err(|e| e.to_string()) +} + +/// What [`update_assigned_rollup`] should do with the harness's todo, +/// given this poll's counts. Kept separate from the socket call so the +/// decision — the part a transient forge outage must not corrupt — is +/// testable without a live server. +#[derive(Debug)] +enum RollupDecision { + /// Either count fetch failed; the existing todo (if any) stands as-is + /// until a poll where both counts are known succeeds. + Skip, + Clear, + Upsert(String), +} + +/// Either count being an `Err` means this poll cannot tell zero assigned +/// apart from "forge didn't answer" — treating it as zero would clear (or +/// wrongly resize) a todo that a working poll would have kept. +fn decide_rollup(issues: Result, pulls: Result) -> RollupDecision { + let (Ok(issues), Ok(pulls)) = (issues, pulls) else { + return RollupDecision::Skip; }; - let count_str = resp.headers().get("x-total-count")?.to_str().ok()?; - count_str.parse::().ok() + let total = issues + pulls; + if total == 0 { + return RollupDecision::Clear; + } + let breakdown = match (issues, pulls) { + (i, 0) => format!("{i} issue{}", if i == 1 { "" } else { "s" }), + (0, p) => format!("{p} PR{}", if p == 1 { "" } else { "s" }), + (i, p) => format!( + "{i} issue{}, {p} PR{}", + if i == 1 { "" } else { "s" }, + if p == 1 { "" } else { "s" } + ), + }; + RollupDecision::Upsert(format!("{total} open assigned: {breakdown}")) } /// Keep a keyed `"rollup"` todo in sync with the count of open assigned /// issues + PRs: a positive count summarises the breakdown, zero clears -/// the todo. The rollup key is distinct from per-thread keys so clearing -/// it never touches notification todos. +/// the todo, and a failed poll leaves it untouched (see [`decide_rollup`]). +/// The rollup key is distinct from per-thread keys so clearing it never +/// touches notification todos. /// /// Forge-only by nature — it asks the forge what is assigned to this /// agent, which is not a notification-protocol concern and has no @@ -171,43 +215,79 @@ async fn update_assigned_rollup( token: &str, socket: &Path, ) { - let issues = count_assigned(client, forge_url, token, "issues") - .await - .unwrap_or(0); - let pulls = count_assigned(client, forge_url, token, "pulls") - .await - .unwrap_or(0); - let total = issues + pulls; + let issues = count_assigned(client, forge_url, token, "issues").await; + let pulls = count_assigned(client, forge_url, token, "pulls").await; + if let Err(e) = &issues { + warn!("forge_notify: assigned-issues poll failed, leaving rollup todo unchanged: {e}"); + } + if let Err(e) = &pulls { + warn!("forge_notify: assigned-pulls poll failed, leaving rollup todo unchanged: {e}"); + } - let req = if total == 0 { - hive_agent_sock::Request::ClearTodo { + let req = match decide_rollup(issues, pulls) { + RollupDecision::Skip => return, + RollupDecision::Clear => hive_agent_sock::Request::ClearTodo { subsystem: "forge".to_owned(), key: Some("rollup".to_owned()), all: false, - } - } else { - let breakdown = match (issues, pulls) { - (i, 0) => format!("{i} issue{}", if i == 1 { "" } else { "s" }), - (0, p) => format!("{p} PR{}", if p == 1 { "" } else { "s" }), - (i, p) => format!( - "{i} issue{}, {p} PR{}", - if i == 1 { "" } else { "s" }, - if p == 1 { "" } else { "s" } - ), - }; - hive_agent_sock::Request::UpsertTodo { + }, + RollupDecision::Upsert(summary) => hive_agent_sock::Request::UpsertTodo { subsystem: "forge".to_owned(), key: Some("rollup".to_owned()), - summary: format!("{total} open assigned: {breakdown}"), + summary, source: None, reopen_if_acked: false, - } + }, }; match hive_sock_client::request::<_, hive_agent_sock::Response>(socket, &req, TODO_SOCKET_RETRY) .await { - Ok(_) => debug!(total, "forge_notify: assigned rollup todo updated"), + Ok(_) => debug!("forge_notify: assigned rollup todo updated"), Err(e) => debug!("forge_notify: assigned rollup todo update failed: {e}"), } } + +#[cfg(test)] +mod tests { + use super::{RollupDecision, decide_rollup}; + + #[test] + fn failed_issues_poll_leaves_todo_untouched() { + assert!(matches!( + decide_rollup(Err("timeout".to_owned()), Ok(0)), + RollupDecision::Skip + )); + } + + #[test] + fn failed_pulls_poll_leaves_todo_untouched() { + assert!(matches!( + decide_rollup(Ok(3), Err("HTTP 500".to_owned())), + RollupDecision::Skip + )); + } + + #[test] + fn both_polls_failing_leaves_todo_untouched() { + assert!(matches!( + decide_rollup(Err("timeout".to_owned()), Err("timeout".to_owned())), + RollupDecision::Skip + )); + } + + #[test] + fn successful_zero_clears_todo() { + assert!(matches!(decide_rollup(Ok(0), Ok(0)), RollupDecision::Clear)); + } + + #[test] + fn successful_nonzero_sets_todo() { + match decide_rollup(Ok(2), Ok(1)) { + RollupDecision::Upsert(summary) => { + assert_eq!(summary, "3 open assigned: 2 issues, 1 PR"); + } + other => panic!("expected Upsert, got {other:?}"), + } + } +}