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