hive-forge-notify: a failed assigned-count poll leaves the rollup todo unchanged

update_assigned_rollup treated a failed count_assigned as zero and issued
ClearTodo, deleting an acked rollup todo on every transient forge error
(and, on a one-sided failure, corrupting the breakdown for one poll).
count_assigned now returns Result<u64, String> instead of Option<u64> so
the caller can log the cause, and decide_rollup only acts once both
counts are known — either failing leaves the existing todo state as-is
and logs a warn! naming which poll failed.

Closes #4720
This commit is contained in:
atlas 2026-09-27 00:02:18 +02:00 • committed by mara
commit ee25b7de20

View file

@ -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 /// Count of open issues or PRs assigned to this agent, via Forgejo's
/// global search API. `issue_type` is `"issues"` or `"pulls"`. Reads the /// global search API. `issue_type` is `"issues"` or `"pulls"`. Reads the
/// `X-Total-Count` header rather than deserialising the body — only page 1 /// `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( async fn count_assigned(
client: &reqwest::Client, client: &reqwest::Client,
forge_url: &str, forge_url: &str,
token: &str, token: &str,
issue_type: &str, issue_type: &str,
) -> Option<u64> { ) -> Result<u64, String> {
let url = format!( let url = format!(
"{forge_url}/api/v1/issues/search\ "{forge_url}/api/v1/issues/search\
?type={issue_type}&state=open&assigned=true&limit=1&page=1" ?type={issue_type}&state=open&assigned=true&limit=1&page=1"
); );
let resp = match client let resp = client
.get(&url) .get(&url)
.header("Authorization", format!("token {token}")) .header("Authorization", format!("token {token}"))
.send() .send()
.await .await
{ .map_err(|e| e.to_string())?;
Ok(r) if r.status().is_success() => r, if !resp.status().is_success() {
_ => return None, 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::<u64>().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<u64, String>, pulls: Result<u64, String>) -> RollupDecision {
let (Ok(issues), Ok(pulls)) = (issues, pulls) else {
return RollupDecision::Skip;
}; };
let count_str = resp.headers().get("x-total-count")?.to_str().ok()?; let total = issues + pulls;
count_str.parse::<u64>().ok() 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 /// Keep a keyed `"rollup"` todo in sync with the count of open assigned
/// issues + PRs: a positive count summarises the breakdown, zero clears /// issues + PRs: a positive count summarises the breakdown, zero clears
/// the todo. The rollup key is distinct from per-thread keys so clearing /// the todo, and a failed poll leaves it untouched (see [`decide_rollup`]).
/// it never touches notification todos. /// 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 /// Forge-only by nature — it asks the forge what is assigned to this
/// agent, which is not a notification-protocol concern and has no /// agent, which is not a notification-protocol concern and has no
@ -171,43 +215,79 @@ async fn update_assigned_rollup(
token: &str, token: &str,
socket: &Path, socket: &Path,
) { ) {
let issues = count_assigned(client, forge_url, token, "issues") let issues = count_assigned(client, forge_url, token, "issues").await;
.await let pulls = count_assigned(client, forge_url, token, "pulls").await;
.unwrap_or(0); if let Err(e) = &issues {
let pulls = count_assigned(client, forge_url, token, "pulls") warn!("forge_notify: assigned-issues poll failed, leaving rollup todo unchanged: {e}");
.await }
.unwrap_or(0); if let Err(e) = &pulls {
let total = issues + pulls; warn!("forge_notify: assigned-pulls poll failed, leaving rollup todo unchanged: {e}");
}
let req = if total == 0 { let req = match decide_rollup(issues, pulls) {
hive_agent_sock::Request::ClearTodo { RollupDecision::Skip => return,
RollupDecision::Clear => hive_agent_sock::Request::ClearTodo {
subsystem: "forge".to_owned(), subsystem: "forge".to_owned(),
key: Some("rollup".to_owned()), key: Some("rollup".to_owned()),
all: false, all: false,
} },
} else { RollupDecision::Upsert(summary) => hive_agent_sock::Request::UpsertTodo {
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 {
subsystem: "forge".to_owned(), subsystem: "forge".to_owned(),
key: Some("rollup".to_owned()), key: Some("rollup".to_owned()),
summary: format!("{total} open assigned: {breakdown}"), summary,
source: None, source: None,
reopen_if_acked: false, reopen_if_acked: false,
} },
}; };
match hive_sock_client::request::<_, hive_agent_sock::Response>(socket, &req, TODO_SOCKET_RETRY) match hive_sock_client::request::<_, hive_agent_sock::Response>(socket, &req, TODO_SOCKET_RETRY)
.await .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}"), 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:?}"),
}
}
}