Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
565b1b90fc | ||
|
|
8149dc7633 | ||
|
|
711e0ece2a | ||
|
|
6685b33c9d |
12 changed files with 686 additions and 71 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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). 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,
|
||||
|
|
|
|||
|
|
@ -97,9 +97,40 @@ 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). 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.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Propagates the todo-store query failure.
|
||||
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
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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): 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_loose_ends to see them".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(
|
||||
|
|
|
|||
|
|
@ -12,3 +12,4 @@ pub mod db;
|
|||
pub mod operator_questions;
|
||||
pub mod power;
|
||||
pub mod scheduled_prompts;
|
||||
pub mod todos;
|
||||
|
|
|
|||
341
hive-c0re/src/stores/todos.rs
Normal file
341
hive-c0re/src/stores/todos.rs
Normal file
|
|
@ -0,0 +1,341 @@
|
|||
//! Todo store — the persistent, DB-backed half of the "todos"
|
||||
//! (loose-ends v2) system.
|
||||
//!
|
||||
//! 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): 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 {
|
||||
/// Open (creating if needed) the todo store at `path`.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Propagates sqlite open / schema-apply / migration failures.
|
||||
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).
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Propagates sqlite query / execute failures.
|
||||
///
|
||||
/// # Panics
|
||||
///
|
||||
/// Panics if the connection mutex is poisoned.
|
||||
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 producer-resolved todo(s) by `(agent, 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 (keyless rows have no distinguishing key — clear a
|
||||
/// specific one via [`Todos::mark_done`] by id instead). Returns the
|
||||
/// number of rows deleted (0 when nothing matched).
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Propagates the sqlite delete failure.
|
||||
///
|
||||
/// # Panics
|
||||
///
|
||||
/// Panics if the connection mutex is poisoned.
|
||||
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.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Propagates the sqlite delete failure.
|
||||
///
|
||||
/// # Panics
|
||||
///
|
||||
/// Panics if the connection mutex is poisoned.
|
||||
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).
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Propagates the sqlite delete failure.
|
||||
///
|
||||
/// # Panics
|
||||
///
|
||||
/// Panics if the connection mutex is poisoned.
|
||||
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.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Propagates the sqlite prepare / query failures.
|
||||
///
|
||||
/// # Panics
|
||||
///
|
||||
/// Panics if the connection mutex is poisoned.
|
||||
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");
|
||||
}
|
||||
}
|
||||
|
|
@ -901,6 +901,21 @@ pub fn unread_count(client: &Client) -> DaemonResponse {
|
|||
/// mind if latency becomes a concern.
|
||||
#[must_use]
|
||||
pub async fn collect_unread(client: &Client) -> Vec<crate::protocol::RoomUnread> {
|
||||
collect_unread_with_ids(client)
|
||||
.await
|
||||
.into_iter()
|
||||
.map(|(_, ru)| ru)
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Like [`collect_unread`] but pairs each entry with its `OwnedRoomId`.
|
||||
/// The todo producer (loose-ends v2) needs the room id as the
|
||||
/// per-room upsert/dedup key, which the claude-facing `RoomUnread`
|
||||
/// payload intentionally doesn't carry.
|
||||
#[must_use]
|
||||
pub async fn collect_unread_with_ids(
|
||||
client: &Client,
|
||||
) -> Vec<(matrix_sdk::ruma::OwnedRoomId, crate::protocol::RoomUnread)> {
|
||||
use crate::protocol::RoomUnread;
|
||||
let mut result = Vec::new();
|
||||
for room in client.joined_rooms() {
|
||||
|
|
@ -915,12 +930,15 @@ pub async fn collect_unread(client: &Client) -> Vec<crate::protocol::RoomUnread>
|
|||
} else {
|
||||
(None, None)
|
||||
};
|
||||
result.push(RoomUnread {
|
||||
label,
|
||||
count,
|
||||
last_body,
|
||||
last_sender,
|
||||
});
|
||||
result.push((
|
||||
room.room_id().to_owned(),
|
||||
RoomUnread {
|
||||
label,
|
||||
count,
|
||||
last_body,
|
||||
last_sender,
|
||||
},
|
||||
));
|
||||
}
|
||||
result
|
||||
}
|
||||
|
|
|
|||
|
|
@ -299,6 +299,11 @@ async fn bring_up_account(
|
|||
// each sweep (see the sweep fns) so re-invites / new messages re-wake.
|
||||
let invite_notified = Arc::new(tokio::sync::Mutex::new(std::collections::HashSet::new()));
|
||||
let unread_notified = Arc::new(tokio::sync::Mutex::new(std::collections::HashSet::new()));
|
||||
// Startup cancel-and-recreate (loose-ends v2): wipe this agent's
|
||||
// matrix todos so stale ones (rooms read while the daemon was down) don't
|
||||
// linger, then let the first sweep rebuild the set to match current
|
||||
// unread reality. Best-effort; the sweep converges regardless.
|
||||
let _ = crate::wake::send_todo_clear(hyperhive_socket, None, true).await;
|
||||
let sync_loop: SyncLoop = Box::pin(async move {
|
||||
sync_client
|
||||
.sync_with_callback(SyncSettings::default(), move |_response| {
|
||||
|
|
|
|||
|
|
@ -20,23 +20,19 @@ use tokio::sync::Mutex;
|
|||
|
||||
use crate::{handlers, wake};
|
||||
|
||||
/// Wake the agent for any joined room carrying unread notifications it has
|
||||
/// not yet been successfully woken about. Driven from the post-sync
|
||||
/// callback — the same reliable path `sweep_invites` uses — rather than a
|
||||
/// one-shot `m.room.message` event handler: a message-wake whose
|
||||
/// `send_wake` raced a hive-c0re / socket-down window (a container rebuild)
|
||||
/// was dropped with no retry, so the agent went deaf until manually
|
||||
/// prompted. The sync callback fires on every sync response (so promptly on
|
||||
/// new activity), so this re-checks unread each tick and retries a dropped
|
||||
/// wake on the next one — self-healing once the socket is back.
|
||||
/// Push a *todo* (loose-ends v2) for each joined room carrying
|
||||
/// unread notifications, and clear the todo for rooms that have been read.
|
||||
/// Replaces the old direct-wake path: instead of firing an all-rooms wake,
|
||||
/// each unread room becomes a per-room `upsert_todo` keyed by its room id,
|
||||
/// and hive-c0re coalesces exactly one wake when the todo set changes.
|
||||
///
|
||||
/// `notified` dedups so a standing unread wakes the agent once, not on
|
||||
/// every sync tick; it is pruned to the current unread set each pass so a
|
||||
/// room that has been read and then receives a new message wakes again. On
|
||||
/// a wake-send failure the freshly-seen rooms are rolled back out of
|
||||
/// `notified` so the next pass retries them. `account_tag` is `Some(name)`
|
||||
/// only in multi-account mode, prepended so the agent knows which account
|
||||
/// to `read_room` on.
|
||||
/// Driven from the post-sync callback — the same reliable path
|
||||
/// `sweep_invites` uses — so it re-checks unread each tick; a dropped
|
||||
/// upsert is retried on the next one (self-healing once the socket is
|
||||
/// back). `notified` tracks which rooms currently have an active todo so
|
||||
/// a room that becomes read gets its todo cleared. `account_tag` is
|
||||
/// `Some(name)` only in multi-account mode, prepended so the agent knows
|
||||
/// which account to `read_room` on.
|
||||
///
|
||||
/// Self-sent messages don't raise an unread notification server-side, so
|
||||
/// no explicit self-filter is needed here.
|
||||
|
|
@ -46,55 +42,42 @@ pub async fn sweep_unread(
|
|||
notified: &Mutex<HashSet<OwnedRoomId>>,
|
||||
account_tag: Option<&str>,
|
||||
) {
|
||||
// Current set of joined rooms carrying unread notifications.
|
||||
let unread_ids: HashSet<OwnedRoomId> = client
|
||||
.joined_rooms()
|
||||
.into_iter()
|
||||
.filter(|r| r.unread_notification_counts().notification_count > 0)
|
||||
.map(|r| r.room_id().to_owned())
|
||||
.collect();
|
||||
let unread = handlers::collect_unread_with_ids(client).await;
|
||||
let unread_ids: HashSet<OwnedRoomId> = unread.iter().map(|(id, _)| id.clone()).collect();
|
||||
|
||||
// Which unread rooms are newly-seen (not already successfully woken)?
|
||||
let mut fresh = Vec::new();
|
||||
{
|
||||
let mut seen = notified.lock().await;
|
||||
// Drop rooms no longer unread (read / left) so a future new message
|
||||
// in them wakes the agent again.
|
||||
seen.retain(|id| unread_ids.contains(id));
|
||||
for id in &unread_ids {
|
||||
if seen.insert(id.clone()) {
|
||||
fresh.push(id.clone());
|
||||
// Rooms we previously pushed a todo for that are no longer unread → the
|
||||
// agent read them; clear their todo. Snapshot under the lock, send
|
||||
// outside it, drop from `notified` on success (retry next tick on fail).
|
||||
let stale: Vec<OwnedRoomId> = {
|
||||
let active = notified.lock().await;
|
||||
active.difference(&unread_ids).cloned().collect()
|
||||
};
|
||||
for id in stale {
|
||||
if wake::send_todo_clear(socket, Some(id.as_str()), false)
|
||||
.await
|
||||
.is_ok()
|
||||
{
|
||||
notified.lock().await.remove(&id);
|
||||
}
|
||||
}
|
||||
|
||||
// Upsert a todo per currently-unread room. hive-c0re coalesces the wake
|
||||
// iff the summary is new or changed, so re-upserting an unchanged room
|
||||
// every sync tick is a cheap server-side no-op (no re-wake).
|
||||
for (id, ru) in &unread {
|
||||
let summary = wake::tag_account(
|
||||
account_tag,
|
||||
wake::format_unread_summary(std::slice::from_ref(ru)),
|
||||
);
|
||||
match wake::send_todo_upsert(socket, id.as_str(), &summary).await {
|
||||
Ok(()) => {
|
||||
notified.lock().await.insert(id.clone());
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(error = %e, room = %id, "matrix: todo upsert failed; will retry next sync");
|
||||
}
|
||||
}
|
||||
}
|
||||
if fresh.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
// Roll the freshly-seen rooms back out of `notified` so the next sweep
|
||||
// retries them. Used on both the sync-race (empty collect) and the
|
||||
// wake-send-failure paths.
|
||||
let rollback = || async {
|
||||
let mut seen = notified.lock().await;
|
||||
for id in &fresh {
|
||||
seen.remove(id);
|
||||
}
|
||||
};
|
||||
|
||||
// Same all-rooms unread summary the wake body has always used. Empty
|
||||
// only under a sync race (counts changed between the scan and collect).
|
||||
let unread = handlers::collect_unread(client).await;
|
||||
if unread.is_empty() {
|
||||
rollback().await;
|
||||
return;
|
||||
}
|
||||
let body = wake::tag_account(account_tag, wake::format_unread_summary(&unread));
|
||||
if let Err(e) = wake::send_wake(socket, &body).await {
|
||||
rollback().await;
|
||||
tracing::warn!(error = %e, "matrix: unread-sweep wake failed; will retry next sync");
|
||||
} else {
|
||||
tracing::info!(rooms = fresh.len(), "matrix: unread-sweep wake delivered");
|
||||
}
|
||||
}
|
||||
|
||||
/// Wake the agent for any pending room invite it has not yet been
|
||||
|
|
|
|||
|
|
@ -34,15 +34,65 @@ use tokio::net::UnixStream;
|
|||
/// Returns an error on socket connect failure, serialisation failure,
|
||||
/// or I/O error writing to or reading from the socket.
|
||||
pub async fn send_wake(socket: &Path, body: impl AsRef<str>) -> Result<()> {
|
||||
use tokio::io::AsyncBufReadExt;
|
||||
|
||||
let payload = serde_json::json!({
|
||||
"cmd": "wake",
|
||||
"from": "matrix",
|
||||
"body": body.as_ref(),
|
||||
"transient": true,
|
||||
});
|
||||
let line = format!("{}\n", serde_json::to_string(&payload)?);
|
||||
send_line(socket, &payload).await
|
||||
}
|
||||
|
||||
/// Upsert a matrix-subsystem *todo* (loose-ends v2) on the
|
||||
/// hyperhive control socket — the replacement for a direct wake. `key` is
|
||||
/// the room id (the dedup key); hive-c0re coalesces a wake iff the todo is
|
||||
/// new or its `summary` changed. Best-effort like [`send_wake`].
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Returns an error on socket connect failure, serialisation failure,
|
||||
/// or I/O error writing to or reading from the socket.
|
||||
pub async fn send_todo_upsert(socket: &Path, key: &str, summary: impl AsRef<str>) -> Result<()> {
|
||||
let payload = serde_json::json!({
|
||||
"cmd": "upsert_todo",
|
||||
"subsystem": "matrix",
|
||||
"key": key,
|
||||
"summary": summary.as_ref(),
|
||||
});
|
||||
send_line(socket, &payload).await
|
||||
}
|
||||
|
||||
/// Clear matrix-subsystem todos. `key = Some(room)` clears one room's
|
||||
/// todo (it was read); `all = true` wipes the whole matrix set
|
||||
/// (cancel-and-recreate on daemon restart). Best-effort.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Returns an error on socket connect failure, serialisation failure,
|
||||
/// or I/O error writing to or reading from the socket.
|
||||
pub async fn send_todo_clear(socket: &Path, key: Option<&str>, all: bool) -> Result<()> {
|
||||
let payload = serde_json::json!({
|
||||
"cmd": "clear_todo",
|
||||
"subsystem": "matrix",
|
||||
"key": key,
|
||||
"all": all,
|
||||
});
|
||||
send_line(socket, &payload).await
|
||||
}
|
||||
|
||||
/// Write one JSON request line to the hyperhive control socket and drain
|
||||
/// the response line (best-effort — the reply is not acted on, we just
|
||||
/// read it so the server doesn't get ECONNRESET on its write-back).
|
||||
/// Shared by [`send_wake`] and the todo senders.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Returns an error on socket connect failure, serialisation failure,
|
||||
/// or I/O error writing to or reading from the socket.
|
||||
async fn send_line(socket: &Path, payload: &serde_json::Value) -> Result<()> {
|
||||
use tokio::io::AsyncBufReadExt;
|
||||
|
||||
let line = format!("{}\n", serde_json::to_string(payload)?);
|
||||
let stream = UnixStream::connect(socket)
|
||||
.await
|
||||
.with_context(|| format!("connect hyperhive socket {}", socket.display()))?;
|
||||
|
|
@ -50,13 +100,11 @@ pub async fn send_wake(socket: &Path, body: impl AsRef<str>) -> Result<()> {
|
|||
write
|
||||
.write_all(line.as_bytes())
|
||||
.await
|
||||
.with_context(|| format!("write wake to {}", socket.display()))?;
|
||||
.with_context(|| format!("write to {}", socket.display()))?;
|
||||
write
|
||||
.shutdown()
|
||||
.await
|
||||
.with_context(|| format!("shutdown write to {}", socket.display()))?;
|
||||
// Drain the response line so the server doesn't get ECONNRESET on
|
||||
// its write-back. We don't act on the response — best-effort wake.
|
||||
let mut reader = tokio::io::BufReader::new(read);
|
||||
let mut resp = String::new();
|
||||
let _ = reader.read_line(&mut resp).await;
|
||||
|
|
|
|||
|
|
@ -269,6 +269,23 @@ pub enum LooseEnd {
|
|||
#[serde(default)]
|
||||
summary: String,
|
||||
},
|
||||
/// A dynamic, subsystem-pushed todo (loose-ends v2). 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,41 @@ pub enum Request {
|
|||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
agent: Option<String>,
|
||||
},
|
||||
/// Upsert a *todo* (loose-ends v2) 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 producer-resolved todo(s) by `(subsystem, key)`. `key =
|
||||
/// Some(k)` clears the one keyed row; `key = None` clears **all** of
|
||||
/// the subsystem's keyless todos (rows with no key can't be told
|
||||
/// apart — clear a specific one via `MarkTodoDone` by id). `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).
|
||||
|
|
|
|||
Loading…
Reference in a new issue