From 2549cc51bac75fa19a63cca80667e130d9ce26f7 Mon Sep 17 00:00:00 2001 From: damocles Date: Wed, 1 Jul 2026 18:37:08 +0200 Subject: [PATCH 1/2] fix(#2106): persist forge-notify dedupe cursor to disk so rebuilds don't re-deliver the unread backlog --- docs/forge.md | 33 ++++--- hive-ag3nt/src/forge_notify.rs | 155 ++++++++++++++++++++++++++++++--- 2 files changed, 164 insertions(+), 24 deletions(-) diff --git a/docs/forge.md b/docs/forge.md index 9db84275..759d9074 100644 --- a/docs/forge.md +++ b/docs/forge.md @@ -92,17 +92,28 @@ that unread signal would be consumed before the agent acts and the guard could never fire. Because a delivered thread stays unread, it reappears in every -`?all=false` poll. An in-memory **delivery-dedupe cursor** (thread -id → last-delivered `updated_at`, held in the poll loop) stops the -same version from re-firing a wake; a new comment bumps `updated_at` -so genuinely new activity re-delivers. The cursor is pure anti-spam, -not a correctness oracle: lost on harness restart it just -re-delivers currently-unread threads once (harmless — `recv` -tolerates redelivery), so it carries none of the persisted-mirror -fragility that ruled out an on-disk seen-cursor. Each poll prunes -the cursor to the threads still in the unread set. A failed wake -delivery is left unread **and** out of the cursor, so it resurfaces -next tick. +`?all=false` poll. A **delivery-dedupe cursor** (thread id → +last-delivered `updated_at`) stops the same version from re-firing a +wake; a new comment bumps `updated_at` so genuinely new activity +re-delivers. The cursor is pure anti-spam, not a correctness oracle. +Each poll prunes it to the threads still in the unread set. A failed +wake delivery is left unread **and** out of the cursor, so it +resurfaces next tick. + +The cursor is **persisted** to +`$HYPERHIVE_STATE_DIR/forge-notify-cursor.json` (atomic tmp+rename, +flushed only when it changed) and reloaded on boot, so a container +rebuild/restart doesn't re-deliver the whole currently-unread backlog +(#2106 — previously the in-memory-only cursor was lost on restart and +every old still-unread thread re-fired a wake). This is safe because a +thread is recorded **after** a successful broker delivery, and the +broker inbox is durable sqlite — so a persisted "delivered" entry can +never swallow a wake the agent never received. Crucially the cursor +file is a private dedup mirror, **not** forge's read-state: it does +not reintroduce the read-before-comment coupling that ruled out the +old mark-read-on-delivery approach. A missing (first boot) or corrupt +cursor file degrades to empty — re-deliver the unread set once — never +an abort. Self-echo notifications (the agent's own writes, see below) are the one path still marked-read directly (no read-before-comment value). diff --git a/hive-ag3nt/src/forge_notify.rs b/hive-ag3nt/src/forge_notify.rs index fa039ca0..529d556f 100644 --- a/hive-ag3nt/src/forge_notify.rs +++ b/hive-ag3nt/src/forge_notify.rs @@ -4,10 +4,14 @@ //! delivers it to the agent's inbox. Delivered threads are deliberately //! left UNREAD in forge — the hive-forge read-before-comment guard keys //! off forge's own unread-state, and the agent reading the thread via the -//! CLI is what marks it read. An in-memory delivery-dedupe cursor -//! (thread id → last-delivered `updated_at`) stops the still-unread -//! notification from re-firing a wake every poll; self-echo and -//! drop-listed notifications are still marked read directly. +//! CLI is what marks it read. A delivery-dedupe cursor (thread id → +//! last-delivered `updated_at`) stops the still-unread notification from +//! re-firing a wake every poll; self-echo and drop-listed notifications +//! are still marked read directly. The cursor is persisted to +//! `$HYPERHIVE_STATE_DIR/forge-notify-cursor.json` (atomic tmp+rename) so +//! a container rebuild/restart doesn't re-deliver the whole currently- +//! unread backlog — the on-disk cursor is a private dedup mirror, NOT +//! forge's read-state, so the read-before-comment guard is untouched. //! //! Activation gates, self-notification filtering, body excerpt + //! truncation + heading escape, wrapper formats (comment / review / @@ -123,15 +127,35 @@ pub async fn run(socket: PathBuf) { // Delivery-dedupe cursor: notification thread id -> the `updated_at` // of the version we last woke the agent for. We no longer mark a // thread read on delivery (that would consume the unread signal the - // hive-forge read-before-comment guard relies on), so this in-memory - // map is what stops the same unread notification from re-firing a - // wake every poll. A new comment bumps `updated_at`, so the thread - // re-delivers. This is purely anti-spam, NOT a correctness oracle: - // lost on harness restart it just re-delivers currently-unread - // threads once (harmless — recv tolerates redelivery), so it carries - // none of the persisted-mirror fragility that sank the on-disk - // cursor approach. - let mut delivered: HashMap = HashMap::new(); + // hive-forge read-before-comment guard relies on), so this map is what + // stops the same unread notification from re-firing a wake every poll. + // A new comment bumps `updated_at`, so the thread re-delivers. This is + // purely anti-spam, NOT a correctness oracle. + // + // It is persisted to `$HYPERHIVE_STATE_DIR/forge-notify-cursor.json` + // and reloaded here on boot so a container rebuild/restart doesn't + // re-deliver the entire currently-unread backlog (#2106). Persisting + // is safe because we only record a thread AFTER a successful broker + // delivery, and the broker inbox is durable sqlite — so a persisted + // "delivered" entry can never swallow a wake the agent never received. + // The cursor file is a private dedup mirror, decoupled from forge's + // own read-state, so it doesn't reintroduce the read-before-comment + // coupling that the mark-read-on-delivery approach suffered. + let cursor_path: Option = if state_dir.is_empty() { + None + } else { + Some(PathBuf::from(format!( + "{state_dir}/forge-notify-cursor.json" + ))) + }; + let mut delivered: HashMap = + cursor_path.as_deref().map(load_cursor).unwrap_or_default(); + if !delivered.is_empty() { + info!( + entries = delivered.len(), + "forge_notify: restored delivery-dedupe cursor from disk" + ); + } loop { interval.tick().await; @@ -141,12 +165,53 @@ pub async fn run(socket: PathBuf) { &token, &socket, &mut delivered, + cursor_path.as_deref(), &own_login, ) .await; } } +/// Load the persisted delivery-dedupe cursor from disk. Returns an empty +/// map on any error — a missing file (first boot) or corrupt JSON both +/// degrade to the pre-persistence behaviour (re-deliver the currently- +/// unread set once) rather than aborting the poller. +fn load_cursor(path: &Path) -> HashMap { + match std::fs::read_to_string(path) { + Ok(s) => serde_json::from_str(&s).unwrap_or_else(|e| { + warn!(?path, error = %e, "forge_notify: cursor parse failed — starting empty"); + HashMap::new() + }), + Err(e) => { + debug!(?path, error = %e, "forge_notify: no cursor file — starting empty"); + HashMap::new() + } + } +} + +/// Persist the delivery-dedupe cursor atomically (write tmp + rename). +/// Best-effort: logs on failure and lets the poll loop continue. Safe to +/// call after a successful broker delivery because the wake is already +/// durably queued in the broker sqlite inbox — recording "delivered" here +/// can never swallow a notification the agent hasn't received. +async fn persist_cursor(path: &Path, delivered: &HashMap) { + let json = match serde_json::to_string(delivered) { + Ok(j) => j, + Err(e) => { + warn!(?path, error = %e, "forge_notify: cursor serialize failed"); + return; + } + }; + let tmp = path.with_extension("json.tmp"); + if let Err(e) = tokio::fs::write(&tmp, json.as_bytes()).await { + warn!(?tmp, error = %e, "forge_notify: cursor tmp write failed"); + return; + } + if let Err(e) = tokio::fs::rename(&tmp, path).await { + warn!(?path, error = %e, "forge_notify: cursor rename failed"); + } +} + /// Fetch a JSON value from a URL using the agent's forge token. Returns /// `None` on any HTTP or parse error (best-effort enrichment). async fn fetch_json(client: &reqwest::Client, url: &str, token: &str) -> Option { @@ -723,6 +788,7 @@ async fn poll_once( token: &str, socket: &Path, delivered: &mut HashMap, + cursor_path: Option<&Path>, own_login: &str, ) { let url = format!("{forge_url}/api/v1/notifications?all=false&limit=50"); @@ -761,6 +827,11 @@ async fn poll_once( "forge_notify: delivering notifications" ); + // Tracks whether the dedupe cursor changed this poll (a new delivery + // recorded, or the prune below dropped now-read threads) so we only + // rewrite the on-disk cursor when there's something to persist. + let mut cursor_dirty = false; + for notif in ¬ifications { let Some(id) = notif["id"].as_u64() else { continue; @@ -804,6 +875,7 @@ async fn poll_once( // here in the Ok arm — a failed delivery leaves the cursor // untouched, so it re-delivers next tick. delivered.insert(id, updated_at); + cursor_dirty = true; } Err(e) => { warn!(%id, error = ?e, "forge_notify: deliver failed — leaving unread"); @@ -821,7 +893,17 @@ async fn poll_once( .iter() .filter_map(|n| n["id"].as_u64()) .collect(); + let before_prune = delivered.len(); delivered.retain(|id, _| current_ids.contains(id)); + if delivered.len() != before_prune { + cursor_dirty = true; + } + + // Flush the cursor to disk only when it changed, so a rebuild/restart + // reloads it instead of re-delivering the whole unread backlog (#2106). + if cursor_dirty && let Some(path) = cursor_path { + persist_cursor(path, delivered).await; + } } /// Whether a notification should be delivered as a wake given the @@ -894,6 +976,53 @@ mod tests { assert!(should_deliver(&delivered, 99, "2026-06-22T16:00:00Z")); } + #[test] + fn cursor_serde_roundtrips_u64_keys() { + // serde_json stringifies integer map keys; make sure the + // round-trip preserves the u64 thread ids the cursor is keyed on. + let mut delivered = HashMap::new(); + delivered.insert(42u64, "2026-06-22T16:00:00Z".to_owned()); + delivered.insert(99u64, "2026-06-22T17:30:00Z".to_owned()); + let json = serde_json::to_string(&delivered).unwrap(); + let back: HashMap = serde_json::from_str(&json).unwrap(); + assert_eq!(delivered, back); + } + + #[test] + fn load_cursor_missing_file_is_empty() { + // First boot: no cursor file yet → empty map (re-deliver once). + let path = std::env::temp_dir().join("forge-notify-cursor-missing-xyzzy.json"); + let _ = std::fs::remove_file(&path); + assert!(load_cursor(&path).is_empty()); + } + + #[test] + fn load_cursor_corrupt_file_is_empty() { + // A truncated / garbage cursor degrades to empty rather than + // aborting the poller. + let path = std::env::temp_dir().join(format!( + "forge-notify-cursor-corrupt-{}.json", + std::process::id() + )); + std::fs::write(&path, "{not valid json").unwrap(); + assert!(load_cursor(&path).is_empty()); + let _ = std::fs::remove_file(&path); + } + + #[test] + fn load_cursor_reads_written_map() { + // The happy path: a persisted cursor reloads to the same map. + let path = std::env::temp_dir().join(format!( + "forge-notify-cursor-rw-{}.json", + std::process::id() + )); + let mut delivered = HashMap::new(); + delivered.insert(7u64, "2026-06-22T16:00:00Z".to_owned()); + std::fs::write(&path, serde_json::to_string(&delivered).unwrap()).unwrap(); + assert_eq!(load_cursor(&path), delivered); + let _ = std::fs::remove_file(&path); + } + #[test] fn escape_md_headings_escapes_top_level_atx() { // Argus reviews start with `## argus review`, which would From f6d9ca7f9984062f50f22688477ab3ea71ff413b Mon Sep 17 00:00:00 2001 From: damocles Date: Wed, 1 Jul 2026 18:48:23 +0200 Subject: [PATCH 2/2] fold forge_cursor into hyperhive-harness.json + prose comments (mara/argus review) --- docs/forge.md | 28 ++++--- docs/persistence.md | 13 ++- hive-ag3nt/src/events.rs | 118 +++++++++++++++++++++++---- hive-ag3nt/src/forge_notify.rs | 140 ++++++--------------------------- 4 files changed, 153 insertions(+), 146 deletions(-) diff --git a/docs/forge.md b/docs/forge.md index 759d9074..c74d07f3 100644 --- a/docs/forge.md +++ b/docs/forge.md @@ -100,20 +100,24 @@ Each poll prunes it to the threads still in the unread set. A failed wake delivery is left unread **and** out of the cursor, so it resurfaces next tick. -The cursor is **persisted** to -`$HYPERHIVE_STATE_DIR/forge-notify-cursor.json` (atomic tmp+rename, -flushed only when it changed) and reloaded on boot, so a container -rebuild/restart doesn't re-deliver the whole currently-unread backlog -(#2106 — previously the in-memory-only cursor was lost on restart and -every old still-unread thread re-fired a wake). This is safe because a +The cursor is **persisted** as the `forge_cursor` field of the +harness's consolidated `hyperhive-harness.json` state file (atomic +tmp+rename, flushed only when it changed) and reloaded on boot, so a +container rebuild/restart doesn't re-deliver the whole currently-unread +backlog (#2106 — previously the in-memory-only cursor was lost on +restart and every old still-unread thread re-fired a wake). The poller +runs in the same harness process that owns that file, so it's one +daemon → one state file rather than a second json; both writers +(turn-loop fields + this cursor) go read-modify-write under a shared +lock so neither clobbers the other's fields. This is safe because a thread is recorded **after** a successful broker delivery, and the broker inbox is durable sqlite — so a persisted "delivered" entry can -never swallow a wake the agent never received. Crucially the cursor -file is a private dedup mirror, **not** forge's read-state: it does -not reintroduce the read-before-comment coupling that ruled out the -old mark-read-on-delivery approach. A missing (first boot) or corrupt -cursor file degrades to empty — re-deliver the unread set once — never -an abort. +never swallow a wake the agent never received. Crucially the cursor is +a private dedup mirror, **not** forge's read-state: it does not +reintroduce the read-before-comment coupling that ruled out the old +mark-read-on-delivery approach. A missing (first boot) or malformed +cursor degrades to empty — re-deliver the unread set once — never an +abort. Self-echo notifications (the agent's own writes, see below) are the one path still marked-read directly (no read-before-comment value). diff --git a/docs/persistence.md b/docs/persistence.md index 6cae08b0..af900b95 100644 --- a/docs/persistence.md +++ b/docs/persistence.md @@ -129,7 +129,7 @@ Consolidated harness state file written atomically (`.tmp` + rename) by Shape: ```json -{ "rate_limited": false, "needs_login": false } +{ "rate_limited": false, "needs_login": false, "active_model": "…", "forge_cursor": { "42": "2026-07-01T18:00:00Z" } } ``` - `rate_limited` — set when the harness detects a 429 from the Claude @@ -138,6 +138,17 @@ Shape: - `needs_login` — set when a turn hits 401 (expired OAuth credentials); cleared by `"online"` status (re-auth completed). Drives the `needs_login` flag alongside the `claude_has_session` check. +- `active_model` — the resolved Claude model for the dashboard badge. +- `forge_cursor` — the `forge_notify` delivery-dedupe cursor + (notification thread id → last-delivered `updated_at`), so a + rebuild/restart doesn't re-deliver the whole currently-unread forge + backlog. See [`forge.md`](forge.md). + +Multiple harness tasks write this file (the turn loop for the first +three fields, the `forge_notify` poller for `forge_cursor`), so every +writer goes read-modify-write under a shared in-process lock — each +preserves the fields it doesn't own rather than reconstructing the file +from scratch. hive-c0re reads this file on each `build_all` sweep (~10s) via `container_view::read_harness_flags`. Falls back to the legacy individual diff --git a/hive-ag3nt/src/events.rs b/hive-ag3nt/src/events.rs index b255d3f3..e69a9cff 100644 --- a/hive-ag3nt/src/events.rs +++ b/hive-ag3nt/src/events.rs @@ -106,6 +106,34 @@ fn harness_json_path() -> PathBuf { crate::paths::state_dir().join(HARNESS_JSON) } +// Serialises the read-modify-write of `hyperhive-harness.json`. Two +// harness tasks touch it in the same process — the turn loop (rate-limit +// / needs-login / active-model) and the forge_notify poller (the +// delivery-dedupe cursor) — writing disjoint fields, so each writer must +// preserve the other's. The lock closes the lost-update window between a +// writer's read and its rename. +static HARNESS_JSON_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + +/// Read the consolidated state file as a JSON object, or an empty object +/// when it is missing / unparseable / not an object. +fn read_harness_json() -> serde_json::Value { + std::fs::read_to_string(harness_json_path()) + .ok() + .and_then(|raw| serde_json::from_str::(&raw).ok()) + .filter(serde_json::Value::is_object) + .unwrap_or_else(|| serde_json::json!({})) +} + +/// Atomically overwrite the state file (`.tmp` + rename) so hive-c0re +/// never reads a partial file. +fn write_harness_json(v: &serde_json::Value) { + let path = harness_json_path(); + let tmp = path.with_extension("json.tmp"); + if std::fs::write(&tmp, v.to_string()).is_ok() { + let _ = std::fs::rename(&tmp, &path); + } +} + fn read_harness_state() -> (bool, bool, Option) { // Try the new consolidated file first. if let Ok(raw) = std::fs::read_to_string(harness_json_path()) @@ -133,24 +161,55 @@ fn read_harness_state() -> (bool, bool, Option) { (rate_limited, needs_login, None) } -/// Write harness state atomically via a `.tmp` + `rename` pair so -/// hive-c0re never reads a partial file. Pass `active_model: Some(s)` -/// to include the resolved model (as surfaced in the dashboard badge); -/// `None` omits the field, which hive-c0re treats as "not yet known". +/// Write the turn-loop's harness state fields via a read-modify-write so +/// any other writer's fields (e.g. `forge_notify`'s `forge_cursor`) survive. +/// Pass `active_model: Some(s)` to update the resolved model (surfaced in +/// the dashboard badge); `None` leaves the stored value untouched. fn write_harness_state(rate_limited: bool, needs_login: bool, active_model: Option<&str>) { - let path = harness_json_path(); - let mut json = serde_json::json!({ - "rate_limited": rate_limited, - "needs_login": needs_login, - }); + let _guard = HARNESS_JSON_LOCK + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let mut v = read_harness_json(); + v["rate_limited"] = rate_limited.into(); + v["needs_login"] = needs_login.into(); if let Some(model) = active_model { - json["active_model"] = serde_json::Value::String(model.to_string()); - } - let body = json.to_string(); - let tmp = path.with_extension("json.tmp"); - if std::fs::write(&tmp, &body).is_ok() { - let _ = std::fs::rename(&tmp, &path); + v["active_model"] = model.into(); } + write_harness_json(&v); +} + +/// Parse the `forge_notify` delivery-dedupe cursor (notification thread id +/// -> last-delivered `updated_at`) out of a harness-state JSON value. +/// Empty when the field is absent (first boot) or malformed. +fn forge_cursor_from_json(v: &serde_json::Value) -> std::collections::HashMap { + v.get("forge_cursor") + .cloned() + .and_then(|c| serde_json::from_value(c).ok()) + .unwrap_or_default() +} + +/// Restore the `forge_notify` delivery-dedupe cursor from the consolidated +/// state file so a container rebuild/restart doesn't re-deliver the whole +/// currently-unread backlog. +pub fn read_forge_cursor() -> std::collections::HashMap { + forge_cursor_from_json(&read_harness_json()) +} + +/// Persist the `forge_notify` delivery-dedupe cursor into the consolidated +/// state file, read-modify-write under the shared lock so the turn-loop's +/// own fields survive. Best-effort: a serialize failure is a no-op. +pub fn write_forge_cursor( + cursor: &std::collections::HashMap, +) { + let Ok(value) = serde_json::to_value(cursor) else { + return; + }; + let _guard = HARNESS_JSON_LOCK + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let mut v = read_harness_json(); + v["forge_cursor"] = value; + write_harness_json(&v); } fn now_unix() -> i64 { @@ -1145,10 +1204,37 @@ impl Default for Bus { mod tests { use super::{ BusEvent, DEFAULT_EFFORT, EFFORT_LEVELS, LiveEvent, StoredEvent, TokenUsage, - is_valid_effort, + forge_cursor_from_json, is_valid_effort, }; use serde_json::json; + #[test] + fn forge_cursor_absent_field_is_empty() { + // First boot / a state file that only carries the turn-loop fields: + // no cursor yet, so we re-deliver the currently-unread set once. + assert!(forge_cursor_from_json(&json!({ "rate_limited": false })).is_empty()); + } + + #[test] + fn forge_cursor_malformed_field_is_empty() { + // A wrong-typed / corrupt cursor degrades to empty rather than + // aborting the poller. + assert!(forge_cursor_from_json(&json!({ "forge_cursor": "nonsense" })).is_empty()); + } + + #[test] + fn forge_cursor_roundtrips_u64_keys() { + // serde_json stringifies integer map keys; confirm the u64 thread + // ids the cursor is keyed on survive the JSON round-trip. + let v = json!({ + "forge_cursor": { "42": "2026-06-22T16:00:00Z", "99": "2026-06-22T17:30:00Z" } + }); + let cursor = forge_cursor_from_json(&v); + assert_eq!(cursor.len(), 2); + assert_eq!(cursor.get(&42), Some(&"2026-06-22T16:00:00Z".to_owned())); + assert_eq!(cursor.get(&99), Some(&"2026-06-22T17:30:00Z".to_owned())); + } + #[test] fn stored_event_serializes_ts_beside_kind() { // History-row wire shape: `ts` is a flattened sibling of `kind`, diff --git a/hive-ag3nt/src/forge_notify.rs b/hive-ag3nt/src/forge_notify.rs index 529d556f..7eb98cfe 100644 --- a/hive-ag3nt/src/forge_notify.rs +++ b/hive-ag3nt/src/forge_notify.rs @@ -7,11 +7,12 @@ //! CLI is what marks it read. A delivery-dedupe cursor (thread id → //! last-delivered `updated_at`) stops the still-unread notification from //! re-firing a wake every poll; self-echo and drop-listed notifications -//! are still marked read directly. The cursor is persisted to -//! `$HYPERHIVE_STATE_DIR/forge-notify-cursor.json` (atomic tmp+rename) so -//! a container rebuild/restart doesn't re-deliver the whole currently- -//! unread backlog — the on-disk cursor is a private dedup mirror, NOT -//! forge's read-state, so the read-before-comment guard is untouched. +//! are still marked read directly. The cursor is persisted as the +//! `forge_cursor` field of the harness's consolidated `hyperhive-harness.json` +//! (via [`crate::events`]) and reloaded on boot so a container +//! rebuild/restart doesn't re-deliver the whole currently-unread backlog — +//! it is a private dedup mirror, NOT forge's read-state, so the +//! read-before-comment guard is untouched. //! //! Activation gates, self-notification filtering, body excerpt + //! truncation + heading escape, wrapper formats (comment / review / @@ -132,28 +133,21 @@ pub async fn run(socket: PathBuf) { // A new comment bumps `updated_at`, so the thread re-delivers. This is // purely anti-spam, NOT a correctness oracle. // - // It is persisted to `$HYPERHIVE_STATE_DIR/forge-notify-cursor.json` - // and reloaded here on boot so a container rebuild/restart doesn't - // re-deliver the entire currently-unread backlog (#2106). Persisting - // is safe because we only record a thread AFTER a successful broker - // delivery, and the broker inbox is durable sqlite — so a persisted - // "delivered" entry can never swallow a wake the agent never received. - // The cursor file is a private dedup mirror, decoupled from forge's - // own read-state, so it doesn't reintroduce the read-before-comment - // coupling that the mark-read-on-delivery approach suffered. - let cursor_path: Option = if state_dir.is_empty() { - None - } else { - Some(PathBuf::from(format!( - "{state_dir}/forge-notify-cursor.json" - ))) - }; - let mut delivered: HashMap = - cursor_path.as_deref().map(load_cursor).unwrap_or_default(); + // It is persisted as the `forge_cursor` field of the harness's + // consolidated state file and reloaded here on boot so a container + // rebuild/restart doesn't re-deliver the entire currently-unread + // backlog. Persisting is safe because we only record a thread AFTER a + // successful broker delivery, and the broker inbox is durable sqlite — + // so a persisted "delivered" entry can never swallow a wake the agent + // never received. The cursor is a private dedup mirror, decoupled from + // forge's own read-state, so it doesn't reintroduce the + // read-before-comment coupling that the mark-read-on-delivery approach + // suffered. + let mut delivered: HashMap = crate::events::read_forge_cursor(); if !delivered.is_empty() { info!( entries = delivered.len(), - "forge_notify: restored delivery-dedupe cursor from disk" + "forge_notify: restored delivery-dedupe cursor" ); } @@ -165,53 +159,12 @@ pub async fn run(socket: PathBuf) { &token, &socket, &mut delivered, - cursor_path.as_deref(), &own_login, ) .await; } } -/// Load the persisted delivery-dedupe cursor from disk. Returns an empty -/// map on any error — a missing file (first boot) or corrupt JSON both -/// degrade to the pre-persistence behaviour (re-deliver the currently- -/// unread set once) rather than aborting the poller. -fn load_cursor(path: &Path) -> HashMap { - match std::fs::read_to_string(path) { - Ok(s) => serde_json::from_str(&s).unwrap_or_else(|e| { - warn!(?path, error = %e, "forge_notify: cursor parse failed — starting empty"); - HashMap::new() - }), - Err(e) => { - debug!(?path, error = %e, "forge_notify: no cursor file — starting empty"); - HashMap::new() - } - } -} - -/// Persist the delivery-dedupe cursor atomically (write tmp + rename). -/// Best-effort: logs on failure and lets the poll loop continue. Safe to -/// call after a successful broker delivery because the wake is already -/// durably queued in the broker sqlite inbox — recording "delivered" here -/// can never swallow a notification the agent hasn't received. -async fn persist_cursor(path: &Path, delivered: &HashMap) { - let json = match serde_json::to_string(delivered) { - Ok(j) => j, - Err(e) => { - warn!(?path, error = %e, "forge_notify: cursor serialize failed"); - return; - } - }; - let tmp = path.with_extension("json.tmp"); - if let Err(e) = tokio::fs::write(&tmp, json.as_bytes()).await { - warn!(?tmp, error = %e, "forge_notify: cursor tmp write failed"); - return; - } - if let Err(e) = tokio::fs::rename(&tmp, path).await { - warn!(?path, error = %e, "forge_notify: cursor rename failed"); - } -} - /// Fetch a JSON value from a URL using the agent's forge token. Returns /// `None` on any HTTP or parse error (best-effort enrichment). async fn fetch_json(client: &reqwest::Client, url: &str, token: &str) -> Option { @@ -788,7 +741,6 @@ async fn poll_once( token: &str, socket: &Path, delivered: &mut HashMap, - cursor_path: Option<&Path>, own_login: &str, ) { let url = format!("{forge_url}/api/v1/notifications?all=false&limit=50"); @@ -899,10 +851,11 @@ async fn poll_once( cursor_dirty = true; } - // Flush the cursor to disk only when it changed, so a rebuild/restart - // reloads it instead of re-delivering the whole unread backlog (#2106). - if cursor_dirty && let Some(path) = cursor_path { - persist_cursor(path, delivered).await; + // Flush the cursor to the consolidated state file only when it + // changed, so a rebuild/restart reloads it instead of re-delivering + // the whole unread backlog. + if cursor_dirty { + crate::events::write_forge_cursor(delivered); } } @@ -976,53 +929,6 @@ mod tests { assert!(should_deliver(&delivered, 99, "2026-06-22T16:00:00Z")); } - #[test] - fn cursor_serde_roundtrips_u64_keys() { - // serde_json stringifies integer map keys; make sure the - // round-trip preserves the u64 thread ids the cursor is keyed on. - let mut delivered = HashMap::new(); - delivered.insert(42u64, "2026-06-22T16:00:00Z".to_owned()); - delivered.insert(99u64, "2026-06-22T17:30:00Z".to_owned()); - let json = serde_json::to_string(&delivered).unwrap(); - let back: HashMap = serde_json::from_str(&json).unwrap(); - assert_eq!(delivered, back); - } - - #[test] - fn load_cursor_missing_file_is_empty() { - // First boot: no cursor file yet → empty map (re-deliver once). - let path = std::env::temp_dir().join("forge-notify-cursor-missing-xyzzy.json"); - let _ = std::fs::remove_file(&path); - assert!(load_cursor(&path).is_empty()); - } - - #[test] - fn load_cursor_corrupt_file_is_empty() { - // A truncated / garbage cursor degrades to empty rather than - // aborting the poller. - let path = std::env::temp_dir().join(format!( - "forge-notify-cursor-corrupt-{}.json", - std::process::id() - )); - std::fs::write(&path, "{not valid json").unwrap(); - assert!(load_cursor(&path).is_empty()); - let _ = std::fs::remove_file(&path); - } - - #[test] - fn load_cursor_reads_written_map() { - // The happy path: a persisted cursor reloads to the same map. - let path = std::env::temp_dir().join(format!( - "forge-notify-cursor-rw-{}.json", - std::process::id() - )); - let mut delivered = HashMap::new(); - delivered.insert(7u64, "2026-06-22T16:00:00Z".to_owned()); - std::fs::write(&path, serde_json::to_string(&delivered).unwrap()).unwrap(); - assert_eq!(load_cursor(&path), delivered); - let _ = std::fs::remove_file(&path); - } - #[test] fn escape_md_headings_escapes_top_level_atx() { // Argus reviews start with `## argus review`, which would