From 993e0bbd4af52659385bd2f779bb2b8c1fcbbf23 Mon Sep 17 00:00:00 2001 From: damocles Date: Mon, 20 Jul 2026 22:19:23 +0200 Subject: [PATCH] feat(#2569): add the harness-local todo store --- hive-agent/src/main.rs | 1 + hive-agent/src/paths.rs | 8 ++ hive-agent/src/todos.rs | 293 ++++++++++++++++++++++++++++++++++++++++ 3 files changed, 302 insertions(+) create mode 100644 hive-agent/src/todos.rs diff --git a/hive-agent/src/main.rs b/hive-agent/src/main.rs index 69b8b0b1..effadb98 100644 --- a/hive-agent/src/main.rs +++ b/hive-agent/src/main.rs @@ -23,6 +23,7 @@ mod prompt; mod serve_common; mod stats; mod stream_enrich; +mod todos; mod turn; mod turn_stats; mod vacuum; diff --git a/hive-agent/src/paths.rs b/hive-agent/src/paths.rs index 9535c7a5..2d7c7b99 100644 --- a/hive-agent/src/paths.rs +++ b/hive-agent/src/paths.rs @@ -40,6 +40,14 @@ pub fn harness_dir() -> PathBuf { hive_sh4re::paths::harness_dir() } +/// Harness-local todo store (loose-ends v2). A dedicated sqlite db under +/// the harness dir — the todos are mutable per-agent state the harness +/// owns, kept out of the append-only `hyperhive-events.sqlite` sink. +#[must_use] +pub fn todos_db() -> PathBuf { + harness_dir().join("hyperhive-todos.sqlite") +} + /// Per-turn config dir for the regenerated claude-{mcp-config,settings, /// system-prompt} files the harness drops before each turn. Set by /// systemd via `RuntimeDirectory = "hive-config"`: a per-service runtime diff --git a/hive-agent/src/todos.rs b/hive-agent/src/todos.rs new file mode 100644 index 00000000..d95a6cb3 --- /dev/null +++ b/hive-agent/src/todos.rs @@ -0,0 +1,293 @@ +//! Harness-local todo store — the persistent, DB-backed half of the +//! "todos" (loose-ends v2) system, owned by the in-container harness. +//! +//! In-container subsystems (matrix, forge-notify, bash) push *todos* to +//! the harness over the in-agent socket instead of firing wakes directly. +//! The harness owns this store locally (one sqlite db under the harness +//! dir) and signals its own turn loop on a new/changed row — hive-c0re is +//! not involved (no broker round-trip, no marker files). +//! +//! Because the store lives inside a single agent's container, todos are +//! **not** agent-scoped here (unlike the old c0re store): every row +//! belongs to this agent. A todo is tagged with a `subsystem` marker plus +//! an optional `subsystem_key` (a matrix room id, a bash task id, …), and +//! `(subsystem, subsystem_key)` is the upsert/dedup key — re-pushing the +//! same item is idempotent, and a producer can list / clear / rebuild only +//! its own set (e.g. matrix wipes + recreates its todos on daemon restart). +//! +//! Removal has two paths (mara's call): the producing subsystem `clear`s a +//! todo it has resolved (keyed by subsystem + key), or the agent itself +//! `mark_done`s one by id. Both delete the row. + +use std::path::Path; +use std::sync::Mutex; + +use anyhow::{Context, Result}; +use hive_sh4re::wire_time::now_unix; +use rusqlite::{Connection, params}; + +const SCHEMA: &str = r" +CREATE TABLE IF NOT EXISTS todos ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + subsystem TEXT NOT NULL, + subsystem_key TEXT, + summary TEXT NOT NULL, + source TEXT, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL +); +-- (subsystem, subsystem_key) is the upsert/dedup key. A NULL key never +-- conflicts (SQLite treats NULLs as distinct), so keyless todos always +-- insert as one-offs; keyed todos update in place. +CREATE UNIQUE INDEX IF NOT EXISTS idx_todos_dedup + ON todos (subsystem, subsystem_key); +"; + +/// One dynamic, subsystem-pushed todo. Timestamps are unix seconds; the +/// consumer derives `age_seconds` from `updated_at`. +#[derive(Debug, Clone)] +pub struct Todo { + pub id: i64, + /// Producing subsystem marker (`"matrix"`, `"forge"`, `"bash"`, …). + pub subsystem: String, + /// Optional subsystem-specific dedup key (matrix room id, bash task + /// id, …). `None` = a keyless one-off todo. + pub subsystem_key: Option, + /// Human-readable one-line summary shown to the agent. + pub summary: String, + /// Optional free-text provenance (e.g. the room name / task label). + pub source: Option, + pub created_at: i64, + pub updated_at: i64, +} + +/// The harness-local todo store. Cheap to share behind an `Arc`; the inner +/// connection is guarded by a `Mutex` (todo ops are short sqlite writes). +pub struct Todos { + conn: Mutex, +} + +impl Todos { + /// Open (creating if needed) the todo store at `path`. + /// + /// # Errors + /// + /// Propagates sqlite open / schema-apply failures. + pub fn open(path: &Path) -> Result { + let conn = Connection::open(path) + .with_context(|| format!("open todos db {}", path.display()))?; + conn.execute_batch(SCHEMA).context("apply todos schema")?; + Ok(Self { + conn: Mutex::new(conn), + }) + } + + /// Insert a todo, or update the existing one for `(subsystem, key)` + /// when `key` is `Some` and already present. A `None` key never + /// conflicts, so it always inserts a fresh row. + /// + /// Returns `(id, changed)` where `changed` is `true` when the row is + /// new OR its `summary`/`source` actually differed — the caller uses + /// this to decide whether to signal the turn loop (re-pushing an + /// identical keyed todo is a no-op and must not re-wake). + /// + /// # Errors + /// + /// Propagates sqlite query / execute failures. + /// + /// # Panics + /// + /// Panics if the connection mutex is poisoned. + pub fn upsert( + &self, + subsystem: &str, + key: Option<&str>, + summary: &str, + source: Option<&str>, + ) -> Result<(i64, bool)> { + let conn = self.conn.lock().unwrap(); + let now = now_unix(); + let existing: Option<(i64, String, Option)> = if key.is_some() { + conn.query_row( + "SELECT id, summary, source FROM todos \ + WHERE subsystem = ?1 AND subsystem_key IS ?2", + params![subsystem, key], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + ) + .ok() + } else { + None + }; + if let Some((id, cur_summary, cur_source)) = existing { + let unchanged = cur_summary == summary && cur_source.as_deref() == source; + if unchanged { + return Ok((id, false)); + } + conn.execute( + "UPDATE todos SET summary = ?1, source = ?2, updated_at = ?3 WHERE id = ?4", + params![summary, source, now, id], + )?; + return Ok((id, true)); + } + conn.execute( + "INSERT INTO todos \ + (subsystem, subsystem_key, summary, source, created_at, updated_at) \ + VALUES (?1, ?2, ?3, ?4, ?5, ?5)", + params![subsystem, key, summary, source, now], + )?; + Ok((conn.last_insert_rowid(), true)) + } + + /// Clear producer-resolved todo(s) by `(subsystem, key)`. `key = + /// Some(k)` targets the one keyed row; `key = None` matches + /// `subsystem_key IS NULL`, i.e. **all** keyless todos for that + /// subsystem (clear a specific keyless one via [`Todos::mark_done`] by + /// id instead). Returns the number of rows deleted. + /// + /// # Errors + /// + /// Propagates the sqlite delete failure. + /// + /// # Panics + /// + /// Panics if the connection mutex is poisoned. + pub fn clear(&self, subsystem: &str, key: Option<&str>) -> Result { + let conn = self.conn.lock().unwrap(); + let n = conn.execute( + "DELETE FROM todos WHERE subsystem = ?1 AND subsystem_key IS ?2", + params![subsystem, key], + )?; + Ok(n) + } + + /// Clear every todo `subsystem` owns — used by a producer that rebuilds + /// its whole set on restart (cancel-and-recreate). Returns the number + /// of rows deleted. + /// + /// # Errors + /// + /// Propagates the sqlite delete failure. + /// + /// # Panics + /// + /// Panics if the connection mutex is poisoned. + pub fn clear_subsystem(&self, subsystem: &str) -> Result { + let conn = self.conn.lock().unwrap(); + let n = conn.execute( + "DELETE FROM todos WHERE subsystem = ?1", + params![subsystem], + )?; + Ok(n) + } + + /// The agent marks one of its todos done, by id. Returns the number of + /// rows deleted (0 when the id was unknown / already gone). + /// + /// # Errors + /// + /// Propagates the sqlite delete failure. + /// + /// # Panics + /// + /// Panics if the connection mutex is poisoned. + pub fn mark_done(&self, id: i64) -> Result { + let conn = self.conn.lock().unwrap(); + let n = conn.execute("DELETE FROM todos WHERE id = ?1", params![id])?; + Ok(n) + } + + /// List todos, newest-updated first. `subsystem = Some(..)` filters to + /// one producer's set; `None` returns all. + /// + /// # Errors + /// + /// Propagates the sqlite prepare / query failures. + /// + /// # Panics + /// + /// Panics if the connection mutex is poisoned. + pub fn list(&self, subsystem: Option<&str>) -> Result> { + let conn = self.conn.lock().unwrap(); + let mut stmt = conn.prepare( + "SELECT id, subsystem, subsystem_key, summary, source, created_at, updated_at \ + FROM todos \ + WHERE (?1 IS NULL OR subsystem = ?1) \ + ORDER BY updated_at DESC, id DESC", + )?; + let rows = stmt + .query_map(params![subsystem], |row| { + Ok(Todo { + id: row.get(0)?, + subsystem: row.get(1)?, + subsystem_key: row.get(2)?, + summary: row.get(3)?, + source: row.get(4)?, + created_at: row.get(5)?, + updated_at: row.get(6)?, + }) + })? + .collect::>>()?; + Ok(rows) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + // Return the `TempDir` alongside the store so it outlives the test — + // dropping it early deletes the dir and SQLite fails with + // `SQLITE_READONLY_DBMOVED`. + fn store() -> (tempfile::TempDir, Todos) { + let dir = tempfile::tempdir().unwrap(); + let db = Todos::open(&dir.path().join("todos.sqlite")).unwrap(); + (dir, db) + } + + #[test] + fn keyed_upsert_dedups_and_reports_changed() { + let (_dir, s) = store(); + let (id1, changed1) = s.upsert("matrix", Some("!room:x"), "1 unread", None).unwrap(); + assert!(changed1, "first push is new → changed"); + // Same key + same summary → no-op, not changed (must not re-wake). + let (id2, changed2) = s.upsert("matrix", Some("!room:x"), "1 unread", None).unwrap(); + assert_eq!(id1, id2, "keyed upsert updates in place, same row"); + assert!(!changed2, "identical re-push is a no-op"); + // Same key, new summary → updates, changed. + let (id3, changed3) = s.upsert("matrix", Some("!room:x"), "3 unread", None).unwrap(); + assert_eq!(id1, id3); + assert!(changed3); + assert_eq!(s.list(Some("matrix")).unwrap().len(), 1); + } + + #[test] + fn keyless_todos_always_insert() { + let (_dir, s) = store(); + let (a, _) = s.upsert("bash", None, "task done", None).unwrap(); + let (b, _) = s.upsert("bash", None, "task done", None).unwrap(); + assert_ne!(a, b, "keyless pushes are distinct one-offs"); + assert_eq!(s.list(Some("bash")).unwrap().len(), 2); + } + + #[test] + fn clear_and_mark_done_remove_rows() { + let (_dir, s) = store(); + s.upsert("matrix", Some("!a:x"), "unread", None).unwrap(); + let (id, _) = s.upsert("forge", Some("pr-1"), "review", None).unwrap(); + assert_eq!(s.clear("matrix", Some("!a:x")).unwrap(), 1); + assert_eq!(s.mark_done(id).unwrap(), 1); + assert!(s.list(None).unwrap().is_empty()); + } + + #[test] + fn clear_subsystem_wipes_only_its_own() { + let (_dir, s) = store(); + s.upsert("matrix", Some("!a:x"), "u", None).unwrap(); + s.upsert("matrix", Some("!b:x"), "u", None).unwrap(); + s.upsert("forge", Some("pr-1"), "r", None).unwrap(); + assert_eq!(s.clear_subsystem("matrix").unwrap(), 2); + let left = s.list(None).unwrap(); + assert_eq!(left.len(), 1); + assert_eq!(left[0].subsystem, "forge"); + } +}