feat(#2569): DB-backed per-agent todo store + mcp.sock upsert/clear/list/mark-done ops

This commit is contained in:
damocles 2026-07-19 01:15:17 +02:00 committed by mara
commit 6685b33c9d
8 changed files with 505 additions and 0 deletions

View file

@ -228,6 +228,27 @@ pub(super) fn render_loose_ends(loose_ends: &[hive_sh4re::LooseEnd]) -> String {
);
}
}
hive_sh4re::LooseEnd::Todo {
id,
subsystem,
subsystem_key,
summary,
source,
age_seconds,
} => {
let key = subsystem_key
.as_deref()
.map(|k| format!(" {k}"))
.unwrap_or_default();
let src = source
.as_deref()
.map(|s| format!("{s}"))
.unwrap_or_default();
let _ = writeln!(
out,
"- todo #{id} [{subsystem}{key}, {age_seconds}s old]: {summary}{src} (mark_todo_done to clear)"
);
}
}
}
out

View file

@ -16,6 +16,7 @@ use crate::container_view::{self, ContainerView};
use crate::dashboard_events::DashboardEvent;
use crate::operator_questions::OperatorQuestions;
use crate::socket_server::{self, AgentSocket};
use crate::todos::Todos;
/// Capacity of the dashboard event channel. Slow browser subscribers
/// (idle tab, throttled connection) drop frames past this — that's
@ -32,6 +33,11 @@ pub struct Coordinator {
pub broker: Arc<Broker>,
pub approvals: Arc<Approvals>,
pub questions: Arc<OperatorQuestions>,
/// Dynamic, subsystem-pushed todos (loose-ends v2, #2569). In-agent
/// subsystems (matrix, forge, bash) upsert/clear todos over mcp.sock
/// instead of firing wakes directly; `get_todos` merges these with
/// the computed static loose ends.
pub todos: Arc<Todos>,
/// Scheduled-prompts queue. One sqlite connection,
/// internal mutex; the worker drains due rows and the manager
/// handlers insert / cancel through the same handle.
@ -460,6 +466,7 @@ impl Coordinator {
let broker = Broker::open(db_path).context("open broker")?;
let approvals = Approvals::open(db_path).context("open approvals")?;
let questions = OperatorQuestions::open(db_path).context("open operator_questions")?;
let todos = Todos::open(db_path).context("open todos")?;
let scheduled_prompts = crate::scheduled_prompts::ScheduledPrompts::open(db_path)
.context("open scheduled_prompts")?;
// BuildLogs wants a directory (it picks its own `build_logs.sqlite`
@ -489,6 +496,7 @@ impl Coordinator {
broker: Arc::new(broker),
approvals: Arc::new(approvals),
questions: Arc::new(questions),
todos: Arc::new(todos),
scheduled_prompts: Arc::new(scheduled_prompts),
build_logs,
audit_log,

View file

@ -97,9 +97,36 @@ pub fn for_agent(coord: &Coordinator, agent: &str) -> Result<Vec<LooseEnd>> {
age_seconds: saturating_age(now, r.created_at.timestamp()),
});
}
// Dynamic, subsystem-pushed todos (loose-ends v2, #2569). Scoped to
// this agent; the producing subsystem or the agent itself clears them.
out.extend(todos_for(coord, agent, None)?);
Ok(out)
}
/// This agent's dynamic todos as `LooseEnd::Todo` rows, optionally
/// filtered to one `subsystem`. Shared by [`for_agent`] and the
/// `ListTodos` handler so the row-mapping lives in one place.
pub fn todos_for(
coord: &Coordinator,
agent: &str,
subsystem: Option<&str>,
) -> Result<Vec<LooseEnd>> {
let now = now_unix();
Ok(coord
.todos
.list(agent, subsystem)?
.into_iter()
.map(|t| LooseEnd::Todo {
id: t.id,
subsystem: t.subsystem,
subsystem_key: t.subsystem_key,
summary: t.summary,
source: t.source,
age_seconds: saturating_age(now, t.updated_at.timestamp()),
})
.collect())
}
/// Hive-wide loose-ends view: EVERY pending approval + EVERY
/// unanswered question + EVERY pending reminder. Manager surface
/// only; sub-agents can't see each other's threads via the agent

View file

@ -43,6 +43,7 @@ pub(crate) use stats::{
};
pub(crate) use stores::{
approvals, audit_log, broker, build_logs, db, operator_questions, power, scheduled_prompts,
todos,
};
pub(crate) use workers::{
agent_sockets, auto_update, crash_watch, knowledge, mcp_sockets, reminder_scheduler,

View file

@ -575,6 +575,31 @@ async fn dispatch(req: &AgentRequest, agent: &str, coord: &Arc<Coordinator>) ->
since_secs,
agent: target,
} => handle_reminder_rollup(coord, agent, target.as_deref(), *since_secs),
// Todos (loose-ends v2, #2569): in-container subsystems push/clear
// their own; the agent lists / marks its own done. Scoped to the
// calling agent (the socket identity) — no cross-agent access.
AgentRequest::UpsertTodo {
subsystem,
key,
summary,
source,
} => handle_upsert_todo(
coord,
agent,
subsystem,
key.as_deref(),
summary,
source.as_deref(),
),
AgentRequest::ClearTodo {
subsystem,
key,
all,
} => handle_clear_todo(coord, agent, subsystem, key.as_deref(), *all),
AgentRequest::ListTodos { subsystem } => {
handle_list_todos(coord, agent, subsystem.as_deref())
}
AgentRequest::MarkTodoDone { id } => handle_mark_todo_done(coord, agent, *id),
// Orchestration / diagnostics verbs — gated per-verb on tool-group
// membership or topology (see `dispatch_orchestration`).
_ => dispatch_orchestration(req, agent, coord).await,
@ -777,6 +802,87 @@ fn handle_get_loose_ends(
}
}
/// `UpsertTodo` — a subsystem pushes/updates one of this agent's todos.
/// Coalesces a wake ONLY when the row is new or actually changed, so
/// re-pushing an identical keyed todo is a silent no-op.
fn handle_upsert_todo(
coord: &Arc<Coordinator>,
agent: &str,
subsystem: &str,
key: Option<&str>,
summary: &str,
source: Option<&str>,
) -> AgentResponse {
match coord.todos.upsert(agent, subsystem, key, summary, source) {
Ok((_, changed)) => {
if changed {
let _ = coord.broker.send(&Message {
from: "todo".to_owned(),
to: agent.to_owned(),
body: "you have todos — call get_todos".to_owned(),
in_reply_to: None,
});
}
AgentResponse::Ok
}
Err(e) => AgentResponse::Err {
message: format!("{e:#}"),
},
}
}
/// `ClearTodo` — a producer clears a resolved todo by `(subsystem, key)`,
/// or wipes its whole set when `all` (cancel-and-recreate on restart).
fn handle_clear_todo(
coord: &Arc<Coordinator>,
agent: &str,
subsystem: &str,
key: Option<&str>,
all: bool,
) -> AgentResponse {
let result = if all {
coord.todos.clear_subsystem(agent, subsystem)
} else {
coord.todos.clear(agent, subsystem, key)
};
match result {
Ok(count) => AgentResponse::Acked {
count: u64::try_from(count).unwrap_or(0),
},
Err(e) => AgentResponse::Err {
message: format!("{e:#}"),
},
}
}
/// `ListTodos` — enumerate this agent's todos (optionally one subsystem's)
/// as `LooseEnd::Todo` rows, so a producer can reconcile its own set.
fn handle_list_todos(
coord: &Arc<Coordinator>,
agent: &str,
subsystem: Option<&str>,
) -> AgentResponse {
match crate::loose_ends::todos_for(coord, agent, subsystem) {
Ok(loose_ends) => AgentResponse::LooseEnds { loose_ends },
Err(e) => AgentResponse::Err {
message: format!("{e:#}"),
},
}
}
/// `MarkTodoDone` — the agent clears one of its own todos by id (scoped to
/// the agent, so it can't touch another agent's).
fn handle_mark_todo_done(coord: &Arc<Coordinator>, agent: &str, id: i64) -> AgentResponse {
match coord.todos.mark_done(agent, id) {
Ok(count) => AgentResponse::Acked {
count: u64::try_from(count).unwrap_or(0),
},
Err(e) => AgentResponse::Err {
message: format!("{e:#}"),
},
}
}
/// `CountPendingReminders` — resolve the target (own / subtree free, else
/// `QueryAgentState`) then count its pending reminders.
fn handle_count_pending_reminders(

View file

@ -12,3 +12,4 @@ pub mod db;
pub mod operator_questions;
pub mod power;
pub mod scheduled_prompts;
pub mod todos;

View file

@ -0,0 +1,292 @@
//! Todo store — the persistent, DB-backed half of the "todos"
//! (loose-ends v2) system (issue #2569).
//!
//! Subsystems inside an agent's container (matrix, forge-notify, bash,
//! …) push *todos* to the agent over the mcp.sock protocol instead of
//! firing wakes directly. Todos are scoped to the owning `agent` (the
//! socket identity of the pushing container) and tagged with a
//! `subsystem` marker plus an optional `subsystem_key` (a matrix room
//! id, a bash task id, …); together `(agent, subsystem, subsystem_key)`
//! is the upsert/dedup key, so re-pushing the same item is idempotent
//! (no duplicate) 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 on #2569): 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.
//!
//! This table holds only the *dynamic* subsystem-pushed todos. The
//! static ones (pending approvals / questions / reminders / undelivered
//! messages) are still computed on demand in `loose_ends.rs`; `get_todos`
//! merges the two. Folding the static kinds into this table is a later
//! increment.
use std::path::Path;
use std::sync::Mutex;
use anyhow::{Context, Result};
use chrono::{DateTime, Utc};
use hive_sh4re::wire_time::now_unix;
use rusqlite::{Connection, params};
use serde::Serialize;
use crate::db::Migration;
const SCHEMA: &str = r"
CREATE TABLE IF NOT EXISTS todos (
id INTEGER PRIMARY KEY AUTOINCREMENT,
agent TEXT NOT NULL,
subsystem TEXT NOT NULL,
subsystem_key TEXT,
summary TEXT NOT NULL,
source TEXT,
created_at INTEGER NOT NULL,
updated_at INTEGER NOT NULL
);
-- (agent, 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 (agent, subsystem, subsystem_key);
";
// New table — no legacy rows to migrate past the initial schema.
const MIGRATIONS: &[Migration] = &[];
/// One dynamic, subsystem-pushed todo.
#[derive(Debug, Clone, Serialize)]
pub struct Todo {
pub id: i64,
/// Owning agent (the container whose subsystem pushed it).
pub agent: String,
/// 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.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub subsystem_key: Option<String>,
/// Human-readable one-line summary shown to the agent.
pub summary: String,
/// Optional free-text provenance (e.g. the room name / task label).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source: Option<String>,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
}
pub struct Todos {
conn: Mutex<Connection>,
}
impl Todos {
pub fn open(path: &Path) -> Result<Self> {
let conn = crate::db::open(path, "todos")?;
conn.execute_batch(SCHEMA).context("apply todos schema")?;
crate::db::apply_versioned_migrations(&conn, "todos", MIGRATIONS)?;
Ok(Self {
conn: Mutex::new(conn),
})
}
/// Insert a todo for `agent`, or update the existing one for
/// `(agent, 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 coalesce a wake (re-pushing an identical
/// keyed todo is a no-op and must not re-wake).
pub fn upsert(
&self,
agent: &str,
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<String>)> = if key.is_some() {
conn.query_row(
"SELECT id, summary, source FROM todos \
WHERE agent = ?1 AND subsystem = ?2 AND subsystem_key IS ?3",
params![agent, 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 \
(agent, subsystem, subsystem_key, summary, source, created_at, updated_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?6)",
params![agent, subsystem, key, summary, source, now],
)?;
Ok((conn.last_insert_rowid(), true))
}
/// Clear a producer-resolved todo, keyed by `(agent, subsystem, key)`.
/// Returns the number of rows deleted (0 when nothing matched).
pub fn clear(&self, agent: &str, subsystem: &str, key: Option<&str>) -> Result<usize> {
let conn = self.conn.lock().unwrap();
let n = conn.execute(
"DELETE FROM todos WHERE agent = ?1 AND subsystem = ?2 AND subsystem_key IS ?3",
params![agent, subsystem, key],
)?;
Ok(n)
}
/// Clear every todo `agent`'s `subsystem` owns — used by a producer
/// that rebuilds its whole set on restart (cancel-and-recreate).
/// Returns the number of rows deleted.
pub fn clear_subsystem(&self, agent: &str, subsystem: &str) -> Result<usize> {
let conn = self.conn.lock().unwrap();
let n = conn.execute(
"DELETE FROM todos WHERE agent = ?1 AND subsystem = ?2",
params![agent, subsystem],
)?;
Ok(n)
}
/// The agent marks one of *its own* todos done, by id. Scoped to
/// `agent` so one agent can't clear another's. Returns the number of
/// rows deleted (0 when the id was unknown / not owned / already gone).
pub fn mark_done(&self, agent: &str, id: i64) -> Result<usize> {
let conn = self.conn.lock().unwrap();
let n = conn.execute(
"DELETE FROM todos WHERE id = ?1 AND agent = ?2",
params![id, agent],
)?;
Ok(n)
}
/// List `agent`'s todos, newest-updated first. `subsystem = Some(..)`
/// filters to one producer's set (so a producer can enumerate +
/// reconcile only its own); `None` returns all of the agent's.
pub fn list(&self, agent: &str, subsystem: Option<&str>) -> Result<Vec<Todo>> {
let conn = self.conn.lock().unwrap();
let mut stmt = conn.prepare(
"SELECT id, agent, subsystem, subsystem_key, summary, source, created_at, updated_at \
FROM todos \
WHERE agent = ?1 AND (?2 IS NULL OR subsystem = ?2) \
ORDER BY updated_at DESC, id DESC",
)?;
let rows = stmt
.query_map(params![agent, subsystem], |row| {
let created: i64 = row.get(6)?;
let updated: i64 = row.get(7)?;
Ok(Todo {
id: row.get(0)?,
agent: row.get(1)?,
subsystem: row.get(2)?,
subsystem_key: row.get(3)?,
summary: row.get(4)?,
source: row.get(5)?,
created_at: DateTime::from_timestamp(created, 0).unwrap_or_default(),
updated_at: DateTime::from_timestamp(updated, 0).unwrap_or_default(),
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
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("alice", "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("alice", "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("alice", "matrix", Some("!room:x"), "3 unread", None)
.unwrap();
assert_eq!(id1, id3);
assert!(changed3);
assert_eq!(s.list("alice", Some("matrix")).unwrap().len(), 1);
}
#[test]
fn todos_are_scoped_per_agent() {
let (_dir, s) = store();
// Same subsystem+key for two agents → distinct rows.
s.upsert("alice", "matrix", Some("!r:x"), "u", None)
.unwrap();
s.upsert("bob", "matrix", Some("!r:x"), "u", None).unwrap();
assert_eq!(s.list("alice", None).unwrap().len(), 1);
assert_eq!(s.list("bob", None).unwrap().len(), 1);
// bob can't mark alice's todo done.
let alice_id = s.list("alice", None).unwrap()[0].id;
assert_eq!(s.mark_done("bob", alice_id).unwrap(), 0);
assert_eq!(s.mark_done("alice", alice_id).unwrap(), 1);
}
#[test]
fn keyless_todos_always_insert() {
let (_dir, s) = store();
let (a, _) = s.upsert("alice", "bash", None, "task done", None).unwrap();
let (b, _) = s.upsert("alice", "bash", None, "task done", None).unwrap();
assert_ne!(a, b, "keyless pushes are distinct one-offs");
assert_eq!(s.list("alice", Some("bash")).unwrap().len(), 2);
}
#[test]
fn clear_and_mark_done_remove_rows() {
let (_dir, s) = store();
s.upsert("alice", "matrix", Some("!a:x"), "unread", None)
.unwrap();
let (id, _) = s
.upsert("alice", "forge", Some("pr-1"), "review", None)
.unwrap();
assert_eq!(s.clear("alice", "matrix", Some("!a:x")).unwrap(), 1);
assert_eq!(s.mark_done("alice", id).unwrap(), 1);
assert!(s.list("alice", None).unwrap().is_empty());
}
#[test]
fn clear_subsystem_wipes_only_its_own() {
let (_dir, s) = store();
s.upsert("alice", "matrix", Some("!a:x"), "u", None)
.unwrap();
s.upsert("alice", "matrix", Some("!b:x"), "u", None)
.unwrap();
s.upsert("alice", "forge", Some("pr-1"), "r", None).unwrap();
assert_eq!(s.clear_subsystem("alice", "matrix").unwrap(), 2);
let left = s.list("alice", None).unwrap();
assert_eq!(left.len(), 1);
assert_eq!(left[0].subsystem, "forge");
}
}

View file

@ -269,6 +269,23 @@ pub enum LooseEnd {
#[serde(default)]
summary: String,
},
/// A dynamic, subsystem-pushed todo (loose-ends v2, #2569). Produced
/// by an in-container subsystem (matrix / forge / bash) via
/// `UpsertTodo`. Cleared by that subsystem (`ClearTodo`) or by the
/// agent itself (`MarkTodoDone`, by `id`).
Todo {
id: i64,
/// Producing subsystem marker (`"matrix"`, `"forge"`, `"bash"`, …).
subsystem: String,
/// Optional subsystem-specific key (matrix room id, bash task id).
#[serde(default, skip_serializing_if = "Option::is_none")]
subsystem_key: Option<String>,
summary: String,
/// Optional free-text provenance (room name / task label).
#[serde(default, skip_serializing_if = "Option::is_none")]
source: Option<String>,
age_seconds: u64,
},
}
/// Kind discriminator for `CancelLooseEnd`. Per-kind store +
@ -422,6 +439,38 @@ pub enum Request {
#[serde(default, skip_serializing_if = "Option::is_none")]
agent: Option<String>,
},
/// Upsert a *todo* (loose-ends v2, #2569) from an in-container
/// subsystem (matrix / forge / bash). `subsystem` is the producer
/// marker; `key` is the optional subsystem-specific dedup key (a
/// matrix room id, a bash task id). Re-pushing an identical keyed
/// todo is a no-op; a new-or-changed one coalesces a wake to the
/// agent. Keyless todos always insert as one-offs.
UpsertTodo {
subsystem: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
key: Option<String>,
summary: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
source: Option<String>,
},
/// Clear a producer-resolved todo by `(subsystem, key)`. `key = None`
/// targets the keyless one-off; `all = true` wipes the producer's
/// whole set (cancel-and-recreate on daemon restart).
ClearTodo {
subsystem: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
key: Option<String>,
#[serde(default)]
all: bool,
},
/// List todos, optionally filtered to one `subsystem` (a producer
/// enumerating its own set). `None` = all.
ListTodos {
#[serde(default, skip_serializing_if = "Option::is_none")]
subsystem: Option<String>,
},
/// The agent marks one of its own todos done, by id.
MarkTodoDone { id: i64 },
/// Count of pending (un-delivered) reminders. On the agent socket:
/// same target rules as `GetLooseEnds` (self/children free;
/// non-children require `query_agent_state`; `"*"` rejected).