diff --git a/CLAUDE.md b/CLAUDE.md index 4a7b47d4..d60dfe1b 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -183,43 +183,6 @@ read them à la carte. In-flight or recent context that hasn't earned a section yet. Prune freely. -- **Just landed:** lease-style message delivery / no-drop - on turn fail. The `messages` table gained an `acked_at` - column (idempotent ALTER + backfill = `delivered_at` so - pre-migration delivered rows count as already-acked). - `Broker::recv` now returns `Delivery { id, redelivered, - message }` — the harness gets the row id back so - `AckTurn` can sweep every popped id at turn-end-OK. Two - new wire arms on both agent + manager surfaces: - `AckTurn` (drains the broker's per-recipient in-memory - `unacked_ids` list and stamps the rows `acked_at = NOW`) - and `RequeueInflight` (one-shot at harness boot: resets - `delivered_at = NULL` on every still-inflight row + - remembers each id so the next `Recv` carries - `redelivered: true`). Both bin loops call - `requeue_inflight` once before entering serve, and - `ack_turn` after every `TurnOutcome::Ok` (Failed + - PromptTooLong intentionally skip the ack so the popped - rows stay in-flight for the next boot's requeue). - `format_recv` + `format_wake_prompt` on both bins - surface a `[redelivered after harness restart — may - already be handled]` banner so claude knows the - side-effects of any previous handling may already have - happened. Lock order: `inflight` mutex first then - `conn` mutex in all three methods (`recv` / `ack_turn` - / `requeue_inflight`) so a concurrent pop can't race - the requeue's DB update vs in-memory populate and - miss the redelivered tag. `vacuum_delivered` filter - flipped from `delivered_at < cutoff` to `acked_at IS - NOT NULL AND acked_at < cutoff` so unacked-but- - delivered rows survive vacuum (they're recoverable via - `requeue_inflight`). 7 new tests in `broker::tests` - cover happy path, crash recovery, idempotency, per- - recipient isolation, batch ack, vacuum preservation, - and FIFO ordering on requeue. Closes the "post-rebuild - system-message missed wake" bug class entirely (any - turn that wakes from a `delivered_at NOT NULL, - acked_at NULL` row resurfaces on next boot). - **Just landed:** ctx + cost badges split. The per-agent page now shows TWO chips — `ctx · N` (last inference's prompt size = actual context window utilisation, parsed diff --git a/hive-ag3nt/src/bin/hive-ag3nt.rs b/hive-ag3nt/src/bin/hive-ag3nt.rs index 84dfeac9..a9b9b44e 100644 --- a/hive-ag3nt/src/bin/hive-ag3nt.rs +++ b/hive-ag3nt/src/bin/hive-ag3nt.rs @@ -148,7 +148,6 @@ async fn main() -> Result<()> { } } -#[allow(clippy::too_many_arguments, clippy::similar_names)] async fn serve( socket: &Path, interval: Duration, @@ -161,15 +160,6 @@ async fn serve( ) -> Result<()> { tracing::info!(socket = %socket.display(), "hive-ag3nt serve"); let _ = state; // reserved for future state transitions (turn-loop -> needs-login) - // Boot-time recovery: ask the broker to resurface anything we - // popped in a previous harness session but never acked - // (crashed mid-turn / OOM / container restart). The broker - // resets `delivered_at = NULL` on those rows and remembers - // their ids so the next `Recv` tags them `redelivered: true`; - // we then prepend a "may already be handled" hint to the wake - // prompt. Single shot before entering the serve loop; idempotent - // when there's nothing inflight. - requeue_inflight(socket).await; loop { let recv: Result = // Explicit long-poll: the new agent_server semantics treat @@ -184,13 +174,8 @@ async fn serve( ) .await; match recv { - Ok(AgentResponse::Message { - from, - body, - id: _, - redelivered, - }) => { - tracing::info!(%from, %body, %redelivered, "inbox"); + Ok(AgentResponse::Message { from, body }) => { + tracing::info!(%from, %body, "inbox"); let unread = inbox_unread(socket).await; bus.emit(LiveEvent::TurnStart { from: from.clone(), @@ -201,24 +186,13 @@ async fn serve( let started_at = now_unix(); let started_instant = std::time::Instant::now(); let model_at_start = bus.model(); - let prompt = format_wake_prompt(&from, &body, unread, redelivered); + let prompt = format_wake_prompt(&from, &body, unread); let outcome = { let _guard = turn_lock.lock().await; turn::drive_turn(&prompt, files, &bus).await }; turn::emit_turn_end(&bus, &outcome); bus.set_state(TurnState::Idle); - // Ack only on a clean turn-end. `Failed` leaves every - // message popped during the turn in the unacked list; - // next harness boot's `RequeueInflight` will reset - // `delivered_at = NULL` and tag them `redelivered`. - // `PromptTooLong` is absorbed inside `drive_turn` via - // compaction so it shouldn't reach here, but if it - // does we also skip the ack (safer to redeliver than - // to lose the message). - if matches!(outcome, turn::TurnOutcome::Ok) { - ack_turn(socket).await; - } // Failures are unhandled by definition — PromptTooLong is // absorbed inside drive_turn via compaction, so anything // that reaches Failed here is a real crash. Notify the @@ -246,13 +220,13 @@ async fn serve( s.record(&row); } - // After turn completes, log whether messages arrived during - // the turn — the outer loop will iterate back to recv() on - // its own (the Empty-arm sleep only fires when recv - // actually returned Empty), so no explicit continue needed. + // After turn completes, check if there are pending messages waiting. + // If so, immediately process them instead of blocking on recv(). + // This ensures messages queued during the turn are processed ASAP. let pending = inbox_unread(socket).await; if pending > 0 { tracing::info!(%pending, "pending messages after turn; fetching next"); + continue; // Loop back to recv() immediately instead of sleeping } } Ok(AgentResponse::Empty) => { @@ -287,16 +261,7 @@ async fn serve( /// (`prompts/agent.md` → `claude --system-prompt-file`); this is just the /// wake signal claude reacts to. `unread` is the count of *other* /// messages in the inbox right after this one was popped. -/// `redelivered` flags messages that were popped in a prior harness -/// session, never acked, and resurfaced after a restart — a banner -/// at the top of the wake prompt warns that any side-effects of -/// previous handling may already have happened. -fn format_wake_prompt(from: &str, body: &str, unread: u64, redelivered: bool) -> String { - let banner = if redelivered { - hive_ag3nt::mcp::REDELIVERY_HINT - } else { - "" - }; +fn format_wake_prompt(from: &str, body: &str, unread: u64) -> String { let pending = if unread == 0 { String::new() } else { @@ -304,42 +269,7 @@ fn format_wake_prompt(from: &str, body: &str, unread: u64, redelivered: bool) -> "\n\n({unread} more message(s) pending in your inbox — drain via `mcp__hyperhive__recv` if relevant.)" ) }; - format!("{banner}Incoming message from `{from}`:\n---\n{body}\n---{pending}") -} - -/// Best-effort: tell the broker every message we popped during the -/// turn is now fully handled (turn-end-OK). Swallows transport -/// errors — the worst case is a redundant requeue on next boot. -async fn ack_turn(socket: &Path) { - match client::request::<_, AgentResponse>(socket, &AgentRequest::AckTurn).await { - Ok(AgentResponse::Ok) => {} - Ok(AgentResponse::Err { message }) => { - tracing::warn!(%message, "ack_turn rejected by broker"); - } - Ok(other) => { - tracing::warn!(?other, "ack_turn unexpected response"); - } - Err(e) => tracing::warn!(error = ?e, "ack_turn transport error"), - } -} - -/// Boot-time recovery: ask the broker to resurface anything we -/// popped in a previous harness session but never acked. The broker -/// resets `delivered_at = NULL` on those rows and remembers their -/// ids so the next `Recv` carries `redelivered: true`. Swallows -/// transport errors — they degrade to "no recovery this boot", -/// which is no worse than the pre-feature behaviour (silent drop). -async fn requeue_inflight(socket: &Path) { - match client::request::<_, AgentResponse>(socket, &AgentRequest::RequeueInflight).await { - Ok(AgentResponse::Ok) => {} - Ok(AgentResponse::Err { message }) => { - tracing::warn!(%message, "requeue_inflight rejected by broker"); - } - Ok(other) => { - tracing::warn!(?other, "requeue_inflight unexpected response"); - } - Err(e) => tracing::warn!(error = ?e, "requeue_inflight transport error"), - } + format!("Incoming message from `{from}`:\n---\n{body}\n---{pending}") } /// Best-effort: tell the manager that this agent's last turn crashed @@ -410,11 +340,10 @@ async fn fetch_agent_post_turn_counts(socket: &Path) -> (Option, Option Result<()> { tracing::info!(socket = %socket.display(), "hive-m1nd serve"); - // Same boot-time recovery as hive-ag3nt — see that loop for the - // rationale. Manager-flavour socket so we requeue only manager - // inflight rows. - requeue_inflight(socket).await; loop { let recv: Result = // Explicit long-poll: see hive-ag3nt's serve loop for the @@ -138,12 +134,7 @@ async fn serve( ) .await; match recv { - Ok(ManagerResponse::Message { - from, - body, - id: _, - redelivered, - }) => { + Ok(ManagerResponse::Message { from, body }) => { if from == SYSTEM_SENDER { // Helper events (ApprovalResolved / Spawned / Rebuilt / // Killed / Destroyed) — these are FYI for the manager; @@ -163,14 +154,14 @@ async fn serve( // prompt body so claude sees it. Sender stays "system" // so the wake prompt can label it as such. } - tracing::info!(%from, %body, %redelivered, "manager inbox"); + tracing::info!(%from, %body, "manager inbox"); let unread = inbox_unread(socket).await; bus.emit(LiveEvent::TurnStart { from: from.clone(), body: body.clone(), unread, }); - let prompt = format_wake_prompt(&from, &body, unread, redelivered); + let prompt = format_wake_prompt(&from, &body, unread); bus.set_state(TurnState::Thinking); let started_at = now_unix(); let started_instant = std::time::Instant::now(); @@ -181,12 +172,6 @@ async fn serve( }; turn::emit_turn_end(&bus, &outcome); bus.set_state(TurnState::Idle); - // Ack only on a clean turn-end; Failed leaves the - // popped ids in-flight for the next boot's requeue. - // Mirrors hive-ag3nt; see that loop for full rationale. - if matches!(outcome, turn::TurnOutcome::Ok) { - ack_turn(socket).await; - } if let Some(s) = &stats { let ended_at = now_unix(); let duration_ms = @@ -206,12 +191,12 @@ async fn serve( ); s.record(&row); } - // Check for messages that arrived during the turn so we - // surface "draining" in the logs. The loop will already - // re-iterate from here — no explicit continue needed. + // Check for messages that arrived during the turn and loop + // immediately if any are waiting — mirrors hive-ag3nt behaviour. let pending = inbox_unread(socket).await; if pending > 0 { tracing::info!(%pending, "pending messages after turn; fetching next"); + continue; } } Ok(ManagerResponse::Empty) => { @@ -243,15 +228,8 @@ async fn serve( /// Per-turn user prompt. The role/tools/etc. is in the system prompt /// (`prompts/manager.md` → `claude --system-prompt-file`); this is just /// the wake signal. `unread` is the inbox depth after this message was -/// popped. `redelivered` adds a "may already be handled" banner above -/// the wake body when the broker resurfaced this row (see hive-ag3nt's -/// `format_wake_prompt` for the full story). -fn format_wake_prompt(from: &str, body: &str, unread: u64, redelivered: bool) -> String { - let banner = if redelivered { - hive_ag3nt::mcp::REDELIVERY_HINT - } else { - "" - }; +/// popped. +fn format_wake_prompt(from: &str, body: &str, unread: u64) -> String { let pending = if unread == 0 { String::new() } else { @@ -259,39 +237,7 @@ fn format_wake_prompt(from: &str, body: &str, unread: u64, redelivered: bool) -> "\n\n({unread} more message(s) pending in your inbox — drain via `mcp__hyperhive__recv` if relevant.)" ) }; - format!("{banner}Incoming message from `{from}`:\n---\n{body}\n---{pending}") -} - -/// Best-effort: tell the broker every message popped during the turn -/// is now handled. Mirror of `hive-ag3nt::ack_turn` on the manager -/// surface. -async fn ack_turn(socket: &Path) { - match client::request::<_, ManagerResponse>(socket, &ManagerRequest::AckTurn).await { - Ok(ManagerResponse::Ok) => {} - Ok(ManagerResponse::Err { message }) => { - tracing::warn!(%message, "ack_turn rejected by broker"); - } - Ok(other) => { - tracing::warn!(?other, "ack_turn unexpected response"); - } - Err(e) => tracing::warn!(error = ?e, "ack_turn transport error"), - } -} - -/// Boot-time recovery: ask the broker to resurface any inflight (popped -/// but not acked) messages so the next `Recv` re-delivers them with -/// the redelivery banner. Mirror of `hive-ag3nt::requeue_inflight`. -async fn requeue_inflight(socket: &Path) { - match client::request::<_, ManagerResponse>(socket, &ManagerRequest::RequeueInflight).await { - Ok(ManagerResponse::Ok) => {} - Ok(ManagerResponse::Err { message }) => { - tracing::warn!(%message, "requeue_inflight rejected by broker"); - } - Ok(other) => { - tracing::warn!(?other, "requeue_inflight unexpected response"); - } - Err(e) => tracing::warn!(error = ?e, "requeue_inflight transport error"), - } + format!("Incoming message from `{from}`:\n---\n{body}\n---{pending}") } async fn inbox_unread(socket: &Path) -> u64 { @@ -333,9 +279,8 @@ async fn fetch_manager_post_turn_counts(socket: &Path) -> (Option, Option Self { + fn from_usage_obj(u: &serde_json::Value) -> Option { let field = |k: &str| u.get(k).and_then(serde_json::Value::as_u64).unwrap_or(0); - Self { + Some(Self { input_tokens: field("input_tokens"), output_tokens: field("output_tokens"), cache_read_input_tokens: field("cache_read_input_tokens"), cache_creation_input_tokens: field("cache_creation_input_tokens"), - } + }) } } @@ -494,7 +494,7 @@ impl Bus { }); } - /// Broadcast a status flip (online / `needs_login_*`). Called by + /// Broadcast a status flip (online / needs_login_*). Called by /// the bin entry points + `turn::wait_for_login` + the /// `post_login_*` handlers — every site that mutates the /// `Arc>` should also call this so the web UI diff --git a/hive-ag3nt/src/mcp.rs b/hive-ag3nt/src/mcp.rs index d6c4ea28..ace175aa 100644 --- a/hive-ag3nt/src/mcp.rs +++ b/hive-ag3nt/src/mcp.rs @@ -34,18 +34,7 @@ use crate::client; pub enum SocketReply { Ok, Err(String), - /// `id` is the broker's row id — not surfaced to claude but - /// useful for harness-side bookkeeping (not used in this module - /// today; the bin loops drive ack via `AckTurn` instead of - /// per-id). `redelivered` triggers the "may already be handled" - /// hint in `format_recv` so claude sees it when draining the - /// inbox in-turn. - Message { - from: String, - body: String, - id: i64, - redelivered: bool, - }, + Message { from: String, body: String }, Empty, Status(u64), QuestionQueued(i64), @@ -65,17 +54,7 @@ impl From for SocketReply { match r { hive_sh4re::AgentResponse::Ok => Self::Ok, hive_sh4re::AgentResponse::Err { message } => Self::Err(message), - hive_sh4re::AgentResponse::Message { - from, - body, - id, - redelivered, - } => Self::Message { - from, - body, - id, - redelivered, - }, + hive_sh4re::AgentResponse::Message { from, body } => Self::Message { from, body }, hive_sh4re::AgentResponse::Empty => Self::Empty, hive_sh4re::AgentResponse::Status { unread } => Self::Status(unread), hive_sh4re::AgentResponse::Recent { rows } => Self::Recent(rows), @@ -102,17 +81,7 @@ impl From for SocketReply { match r { hive_sh4re::ManagerResponse::Ok => Self::Ok, hive_sh4re::ManagerResponse::Err { message } => Self::Err(message), - hive_sh4re::ManagerResponse::Message { - from, - body, - id, - redelivered, - } => Self::Message { - from, - body, - id, - redelivered, - }, + hive_sh4re::ManagerResponse::Message { from, body } => Self::Message { from, body }, hive_sh4re::ManagerResponse::Empty => Self::Empty, hive_sh4re::ManagerResponse::Status { unread } => Self::Status(unread), hive_sh4re::ManagerResponse::QuestionQueued { id } => Self::QuestionQueued(id), @@ -148,22 +117,10 @@ pub fn format_ack(resp: Result, tool: &str, ok_msg: } /// Format helper for `recv` tools: `Message` → from + body block; -/// `Empty` → marker; anything else surfaces as an error. When the -/// broker tags the row as `redelivered` (popped before, never acked, -/// resurfaced after a harness restart) a short banner is prepended -/// so claude knows the side-effects of any previous handling may -/// already have happened. +/// `Empty` → marker; anything else surfaces as an error. pub fn format_recv(resp: Result) -> String { match resp { - Ok(SocketReply::Message { - from, - body, - redelivered, - .. - }) => { - let banner = if redelivered { REDELIVERY_HINT } else { "" }; - format!("{banner}from: {from}\n\n{body}") - } + Ok(SocketReply::Message { from, body }) => format!("from: {from}\n\n{body}"), Ok(SocketReply::Empty) => "(empty)".into(), Ok(SocketReply::Err(m)) => format!("recv failed: {m}"), Ok(other) => format!("recv unexpected response: {other:?}"), @@ -171,14 +128,6 @@ pub fn format_recv(resp: Result) -> String { } } -/// Header prepended to message bodies that were popped by a prior -/// harness session, never acked (turn crash / OOM / restart), and -/// resurfaced by `RequeueInflight` on this session's boot. Same -/// string surfaces in the wake prompt (see the bin loops) and the -/// in-turn `recv` tool result so claude sees the warning either way. -pub const REDELIVERY_HINT: &str = - "[redelivered after harness restart — may already be handled]\n"; - /// Format helper for `get_loose_ends`: renders a short bulleted list /// of pending approvals + questions + reminders. Empty list collapses /// to a clear marker so claude doesn't go hunting for a payload that @@ -308,13 +257,11 @@ where /// content failure and the model would burn a turn retrying it. pub fn annotate_retries(mut s: String, retries: u32) -> String { if retries > 0 { - use std::fmt::Write as _; let suffix = if retries == 1 { "retry" } else { "retries" }; - let _ = write!( - s, + s.push_str(&format!( "\n\n(note: hive socket connect needed {retries} {suffix} — c0re likely \ restarted. Your request did succeed on the final attempt; no action needed.)" - ); + )); } s } diff --git a/hive-ag3nt/src/plugins.rs b/hive-ag3nt/src/plugins.rs index 0fb0b982..baea3630 100644 --- a/hive-ag3nt/src/plugins.rs +++ b/hive-ag3nt/src/plugins.rs @@ -27,8 +27,9 @@ const AUTO_UPDATE_PATH: &str = "/etc/hyperhive/claude-plugins-auto-update.json"; /// "already exists" message and exits non-zero on some versions). /// Required before any `@` install can resolve. async fn add_marketplaces() { - let Ok(raw) = tokio::fs::read_to_string(MARKETPLACES_PATH).await else { - return; + let raw = match tokio::fs::read_to_string(MARKETPLACES_PATH).await { + Ok(s) => s, + Err(_) => return, }; let sources: Vec = match serde_json::from_str(&raw) { Ok(v) => v, @@ -106,8 +107,9 @@ async fn update_marketplaces() { /// journald. The manager itself passes `None` — there's nobody above /// it to notify. pub async fn install_configured(socket: &Path, notify_recipient: Option<&str>) { - let Ok(raw) = tokio::fs::read_to_string(PLUGINS_PATH).await else { - return; + let raw = match tokio::fs::read_to_string(PLUGINS_PATH).await { + Ok(s) => s, + Err(_) => return, }; let specs: Vec = match serde_json::from_str(&raw) { Ok(v) => v, diff --git a/hive-ag3nt/src/turn.rs b/hive-ag3nt/src/turn.rs index 99ca80e1..b4ad7945 100644 --- a/hive-ag3nt/src/turn.rs +++ b/hive-ag3nt/src/turn.rs @@ -222,12 +222,7 @@ pub async fn compact_session(files: &TurnFiles, bus: &Bus) -> Result<()> { Ok(()) } -#[allow(clippy::too_many_lines)] async fn run_claude(prompt: &str, files: &TurnFiles, bus: &Bus) -> Result { - // Keep the last STDERR_TAIL_LINES of stderr so a non-zero exit can - // include real context in the bail message (and downstream in the - // failure notification to the manager) instead of just "exit 1". - const STDERR_TAIL_LINES: usize = 20; let model = bus.model(); let resume = !bus.take_skip_continue(); if !resume { @@ -316,6 +311,10 @@ async fn run_claude(prompt: &str, files: &TurnFiles, bus: &Bus) -> Result } } }); + // Keep the last STDERR_TAIL_LINES of stderr so a non-zero exit can + // include real context in the bail message (and downstream in the + // failure notification to the manager) instead of just "exit 1". + const STDERR_TAIL_LINES: usize = 20; let stderr_tail: Arc>> = Arc::new(Mutex::new(VecDeque::with_capacity(STDERR_TAIL_LINES))); let tail_clone = stderr_tail.clone(); diff --git a/hive-ag3nt/src/turn_stats.rs b/hive-ag3nt/src/turn_stats.rs index ffae2807..899a6e66 100644 --- a/hive-ag3nt/src/turn_stats.rs +++ b/hive-ag3nt/src/turn_stats.rs @@ -1,8 +1,8 @@ //! Per-turn analytics sink. One sqlite row per claude turn captures: -//! identity (`model`, `wake_from`, `result_kind`), timing (`started_at`, -//! `ended_at`, `duration_ms`), cost (token counts), behaviour (tool-call +//! identity (model, wake_from, result_kind), timing (started_at, +//! ended_at, duration_ms), cost (token counts), behaviour (tool-call //! count + per-tool breakdown), and post-turn snapshot metrics -//! (`open_threads_count`, `open_reminders_count`). +//! (open_threads_count, open_reminders_count). //! //! Lives next to `hyperhive-events.sqlite` in the agent's state dir //! so the host-side state vacuum sweep can reach both. Schema is @@ -65,7 +65,7 @@ const MIGRATIONS: &[&str] = &[ /// One row to be inserted. `Option`-wrapped fields default to NULL /// when the harness couldn't gather them (e.g. socket roundtrip for -/// `open_threads` failed) so a partial row beats no row. +/// open_threads failed) so a partial row beats no row. #[derive(Debug, Clone)] pub struct TurnStatRow { pub started_at: i64, diff --git a/hive-ag3nt/src/web_ui.rs b/hive-ag3nt/src/web_ui.rs index caca25ee..2d69153f 100644 --- a/hive-ag3nt/src/web_ui.rs +++ b/hive-ag3nt/src/web_ui.rs @@ -527,8 +527,11 @@ async fn post_compact(State(state): State) -> Response { let lock = state.turn_lock.clone(); // Reject immediately if a normal turn is in flight — concurrent access // to the claude session is unsafe and produces garbled output. - let Ok(guard) = lock.try_lock_owned() else { - return error_response("turn in flight — wait for it to finish before compacting"); + let guard = match lock.try_lock_owned() { + Ok(g) => g, + Err(_) => { + return error_response("turn in flight — wait for it to finish before compacting"); + } }; let bus = state.bus.clone(); let files = state.files.clone(); diff --git a/hive-c0re/src/actions.rs b/hive-c0re/src/actions.rs index 74c37b2b..edc837e8 100644 --- a/hive-c0re/src/actions.rs +++ b/hive-c0re/src/actions.rs @@ -324,7 +324,7 @@ pub async fn destroy(coord: &Arc, name: &str, purge: bool) -> Resul tracing::info!(%name, purge, "destroy"); // Guard auto-clears on the success path's final scope exit and on // every early-return / cancellation along the way. - let guard = coord.transient_guard(name, TransientKind::Destroying); + let _guard = coord.transient_guard(name, TransientKind::Destroying); lifecycle::destroy(name).await?; coord.unregister_agent(name); let runtime = Coordinator::agent_dir(name); @@ -359,7 +359,7 @@ pub async fn destroy(coord: &Arc, name: &str, purge: bool) -> Resul "agent destroyed" }, ); - drop(guard); + drop(_guard); coord.notify_manager(&HelperEvent::Destroyed { agent: name.to_owned(), }); diff --git a/hive-c0re/src/agent_server.rs b/hive-c0re/src/agent_server.rs index 5da185a9..8615f165 100644 --- a/hive-c0re/src/agent_server.rs +++ b/hive-c0re/src/agent_server.rs @@ -94,7 +94,6 @@ fn recv_timeout(wait_seconds: Option) -> std::time::Duration { } } -#[allow(clippy::too_many_lines)] async fn dispatch(req: &AgentRequest, agent: &str, coord: &Arc) -> AgentResponse { let broker = &coord.broker; match req { @@ -103,11 +102,9 @@ async fn dispatch(req: &AgentRequest, agent: &str, coord: &Arc) -> .recv_blocking(agent, recv_timeout(*wait_seconds)) .await { - Ok(Some(d)) => AgentResponse::Message { - from: d.message.from, - body: d.message.body, - id: d.id, - redelivered: d.redelivered, + Ok(Some(msg)) => AgentResponse::Message { + from: msg.from, + body: msg.body, }, Ok(None) => AgentResponse::Empty, Err(e) => AgentResponse::Err { @@ -203,23 +200,6 @@ async fn dispatch(req: &AgentRequest, agent: &str, coord: &Arc) -> |message| AgentResponse::Err { message }, |()| AgentResponse::Ok, ), - AgentRequest::AckTurn => match broker.ack_turn(agent) { - Ok(_n) => AgentResponse::Ok, - Err(e) => AgentResponse::Err { - message: format!("{e:#}"), - }, - }, - AgentRequest::RequeueInflight => match broker.requeue_inflight(agent) { - Ok(n) => { - if n > 0 { - tracing::info!(%agent, requeued = %n, "requeued in-flight messages"); - } - AgentResponse::Ok - } - Err(e) => AgentResponse::Err { - message: format!("{e:#}"), - }, - }, } } @@ -376,11 +356,7 @@ mod tests { fn auto_reminder_path_format() { let p = auto_reminder_path("damocles"); assert!(p.starts_with("/agents/damocles/state/reminders/auto-")); - assert!( - std::path::Path::new(&p) - .extension() - .is_some_and(|ext| ext.eq_ignore_ascii_case("md")) - ); + assert!(p.ends_with(".md")); } #[test] diff --git a/hive-c0re/src/approvals.rs b/hive-c0re/src/approvals.rs index 2deddf5b..2e9dad90 100644 --- a/hive-c0re/src/approvals.rs +++ b/hive-c0re/src/approvals.rs @@ -168,11 +168,8 @@ impl Approvals { /// Mark pending -> approved (or fail if not pending). Returns the (now-updated) /// approval so the caller can run the action and pass the agent name. - #[allow(clippy::type_complexity)] pub fn mark_approved(&self, id: i64) -> Result { let conn = self.conn.lock().unwrap(); - // Row shape: (agent, kind, commit_ref, requested_at, status, - // fetched_sha, description). let current: Option<( String, String, diff --git a/hive-c0re/src/auto_update.rs b/hive-c0re/src/auto_update.rs index 9a18bd76..36dc72cd 100644 --- a/hive-c0re/src/auto_update.rs +++ b/hive-c0re/src/auto_update.rs @@ -63,7 +63,7 @@ pub async fn rebuild_agent(coord: &Arc, name: &str, current_rev: &s // lifecycle::rebuild. Dashboard rebuilds already do this via // lifecycle_action; this catches the auto-update scan + any // other direct caller. - let guard = coord.transient_guard(name, crate::coordinator::TransientKind::Rebuilding); + let _guard = coord.transient_guard(name, crate::coordinator::TransientKind::Rebuilding); let result = lifecycle::rebuild( name, &coord.hyperhive_flake, @@ -75,7 +75,7 @@ pub async fn rebuild_agent(coord: &Arc, name: &str, current_rev: &s &coord.operator_pronouns, ) .await; - drop(guard); + drop(_guard); match &result { Ok(()) => { if let Err(e) = std::fs::write(rev_marker_path(name), current_rev) { diff --git a/hive-c0re/src/broker.rs b/hive-c0re/src/broker.rs index 4cf66c34..e29c21fc 100644 --- a/hive-c0re/src/broker.rs +++ b/hive-c0re/src/broker.rs @@ -1,7 +1,6 @@ //! Sqlite-backed message broker. Survives `hive-c0re` restart, and taps every //! send/recv onto a broadcast channel so the dashboard can stream it. -use std::collections::{HashMap, HashSet}; use std::path::Path; use std::sync::Mutex; use std::time::{SystemTime, UNIX_EPOCH}; @@ -47,18 +46,6 @@ const EVENT_CHANNEL: usize = 256; /// self-documenting. pub type DueReminder = (String, i64, String, Option); -/// A single message hand-off from broker to recipient. Carries the -/// broker's row id (so the harness can drive `ack_turn` later) and -/// the redelivery flag (so the harness can prepend the -/// "may already be handled" hint to the wake prompt). The -/// `Message` itself is identical to a pristine `Send` payload. -#[derive(Debug, Clone)] -pub struct Delivery { - pub id: i64, - pub redelivered: bool, - pub message: Message, -} - /// Row shape for [`Broker::list_pending_reminders`], shipped on the /// dashboard `/api/reminders` response. #[derive(Debug, Clone, Serialize)] @@ -112,33 +99,9 @@ pub enum MessageEvent { }, } -/// Per-recipient in-memory bookkeeping for the deliver-then-ack -/// flow. Source of truth is the DB columns `delivered_at` + -/// `acked_at`; the in-memory state here is purely an optimisation -/// (avoids scanning the messages table on `AckTurn`) plus the -/// redelivery-hint marker. -#[derive(Default)] -struct RecipientInflight { - /// Message ids the broker has handed to this recipient since the - /// last `AckTurn`. Drained on `ack_turn`, which then runs a - /// single `UPDATE … WHERE id IN (…)` to set `acked_at`. - unacked_ids: Vec, - /// Message ids resurfaced by the most recent `requeue_inflight` - /// call. The next `recv` pop of any id in this set tags the - /// response with `redelivered: true` so the harness can prepend - /// the "may already be handled" hint to the wake prompt; - /// successful pops drain the id from the set. - requeued_ids: HashSet, -} - pub struct Broker { conn: Mutex, events: broadcast::Sender, - /// Per-recipient deliver/ack tracking. Lost on hive-c0re restart - /// (harmless — the harness fires `RequeueInflight` on its own - /// boot, which rebuilds the `requeued_ids` set from the DB and - /// clears any stale `unacked_ids`). - inflight: Mutex>, } impl Broker { @@ -150,13 +113,11 @@ impl Broker { let conn = Connection::open(path).with_context(|| format!("open broker db {}", path.display()))?; conn.execute_batch(SCHEMA).context("apply broker schema")?; - ensure_message_columns(&conn).context("migrate messages columns")?; ensure_reminder_columns(&conn).context("migrate reminders columns")?; let (events, _) = broadcast::channel(EVENT_CHANNEL); Ok(Self { conn: Mutex::new(conn), events, - inflight: Mutex::new(HashMap::new()), }) } @@ -268,10 +229,10 @@ impl Broker { &self, recipient: &str, timeout: std::time::Duration, - ) -> Result> { + ) -> Result> { let mut rx = self.subscribe(); - if let Some(d) = self.recv(recipient)? { - return Ok(Some(d)); + if let Some(m) = self.recv(recipient)? { + return Ok(Some(m)); } let deadline = tokio::time::Instant::now() + timeout; loop { @@ -285,8 +246,8 @@ impl Broker { // pop (in case we missed our notification while behind). Ok(Err(_)) => return self.recv(recipient), Ok(Ok(MessageEvent::Sent { to, .. })) if to == recipient => { - if let Some(d) = self.recv(recipient)? { - return Ok(Some(d)); + if let Some(m) = self.recv(recipient)? { + return Ok(Some(m)); } // Lost a race (concurrent recv elsewhere). Keep waiting. } @@ -295,31 +256,22 @@ impl Broker { } } - /// Delete fully-acked messages older than `older_than_secs`. - /// Unacked rows (delivered but not yet acknowledged by a clean - /// turn-end, plus undelivered rows) are always kept regardless of - /// age — the former because they're recoverable via - /// `requeue_inflight`, the latter because they're still in flight + /// Delete delivered messages older than `older_than_secs`. Undelivered + /// rows are always kept regardless of age — those are still in flight /// from the broker's POV. Returns the number of rows removed. pub fn vacuum_delivered(&self, older_than_secs: i64) -> Result { let cutoff = now_unix() - older_than_secs; let conn = self.conn.lock().unwrap(); let n = conn.execute( "DELETE FROM messages - WHERE acked_at IS NOT NULL - AND acked_at < ?1", + WHERE delivered_at IS NOT NULL + AND delivered_at < ?1", params![cutoff], )?; Ok(u64::try_from(n).unwrap_or(0)) } - pub fn recv(&self, recipient: &str) -> Result> { - // Lock order: inflight FIRST, then conn. `requeue_inflight` + - // `ack_turn` follow the same order so we never deadlock; the - // requeue path also needs both locks held together so a pop - // can't sneak in between its DB update + in-memory populate - // and miss the `redelivered` flag. - let mut inflight = self.inflight.lock().unwrap(); + pub fn recv(&self, recipient: &str) -> Result> { let conn = self.conn.lock().unwrap(); let row: Option<(i64, String, String, String)> = conn .query_row( @@ -339,113 +291,14 @@ impl Broker { "UPDATE messages SET delivered_at = ?1 WHERE id = ?2", params![now_unix(), id], )?; - // Track the id so the next `ack_turn(recipient)` can sweep it, - // and check whether it was resurfaced by a recent - // `requeue_inflight` (in which case the wake prompt gets the - // "may already be handled" hint). Both ops are O(1) per pop; - // the hash-set lookup runs at most once per delivery. - let slot = inflight.entry(recipient.to_owned()).or_default(); - slot.unacked_ids.push(id); - let redelivered = slot.requeued_ids.remove(&id); drop(conn); - drop(inflight); let _ = self.events.send(MessageEvent::Delivered { from: from.clone(), to: to.clone(), body: body.clone(), at: now_unix(), }); - Ok(Some(Delivery { - id, - redelivered, - message: Message { from, to, body }, - })) - } - - /// Drain the per-recipient unacked-id list and mark every row - /// `acked_at = NOW`. Fired by the harness after `TurnOutcome::Ok`. - /// Returns the number of rows acked (zero is normal — claude - /// may have not called recv during the turn). Tolerant of ids - /// that no longer exist in the DB (vacuumed, manually deleted) - /// — `UPDATE … WHERE id IN (…)` simply matches zero rows. - pub fn ack_turn(&self, recipient: &str) -> Result { - // Same lock order as `recv` and `requeue_inflight`. - let mut inflight = self.inflight.lock().unwrap(); - let ids: Vec = inflight - .get_mut(recipient) - .map(|s| std::mem::take(&mut s.unacked_ids)) - .unwrap_or_default(); - if ids.is_empty() { - return Ok(0); - } - let now = now_unix(); - let conn = self.conn.lock().unwrap(); - // Bind every id explicitly. Caps in the hundreds in the worst - // case (a single very chatty turn); well under sqlite's 999 - // default param limit and we're already serialising on the - // broker mutex. - let placeholders = std::iter::repeat_n("?", ids.len()) - .collect::>() - .join(","); - let sql = format!("UPDATE messages SET acked_at = ? WHERE id IN ({placeholders})"); - let mut params_vec: Vec<&dyn rusqlite::ToSql> = Vec::with_capacity(ids.len() + 1); - params_vec.push(&now); - for id in &ids { - params_vec.push(id); - } - let n = conn.execute(&sql, params_vec.as_slice())?; - Ok(u64::try_from(n).unwrap_or(0)) - } - - /// Resurface every message the broker previously handed to this - /// recipient that never got `acked_at` set. Used by the harness at - /// boot to recover from the crashed-mid-turn / OOM-killed / - /// container-restarted cases. Three steps: - /// - /// 1. Clear any stale in-memory state for this recipient (the - /// previous harness session's `unacked_ids` are irrelevant — - /// the new session will repopulate from fresh pops). - /// 2. Find every row where `recipient = me`, `delivered_at IS NOT - /// NULL`, `acked_at IS NULL`. Reset `delivered_at = NULL` so - /// the next `Recv` pops them again. - /// 3. Remember each id in the per-recipient `requeued_ids` set so - /// the next pop tags the response with `redelivered: true`. - /// - /// Returns the number of rows requeued. Safe to call when there's - /// nothing in flight (returns 0). Safe to call multiple times - /// (idempotent — the second call finds nothing because the rows - /// are now back in the pending state). - pub fn requeue_inflight(&self, recipient: &str) -> Result { - // Hold inflight + conn together so a concurrent `recv` can't - // pop a just-requeued row between our DB update and our - // in-memory populate and miss the redelivered tag. - let mut inflight = self.inflight.lock().unwrap(); - let conn = self.conn.lock().unwrap(); - let mut stmt = conn.prepare( - "SELECT id FROM messages - WHERE recipient = ?1 - AND delivered_at IS NOT NULL - AND acked_at IS NULL", - )?; - let ids: Vec = stmt - .query_map(params![recipient], |row| row.get(0))? - .collect::>()?; - drop(stmt); - if !ids.is_empty() { - let placeholders = std::iter::repeat_n("?", ids.len()) - .collect::>() - .join(","); - let sql = - format!("UPDATE messages SET delivered_at = NULL WHERE id IN ({placeholders})"); - let params_vec: Vec<&dyn rusqlite::ToSql> = - ids.iter().map(|id| id as &dyn rusqlite::ToSql).collect(); - conn.execute(&sql, params_vec.as_slice())?; - } - let slot = inflight.entry(recipient.to_owned()).or_default(); - slot.unacked_ids.clear(); - slot.requeued_ids.clear(); - slot.requeued_ids.extend(ids.iter().copied()); - Ok(u64::try_from(ids.len()).unwrap_or(0)) + Ok(Some(Message { from, to, body })) } /// Store a new reminder. Returns the reminder id. @@ -512,7 +365,7 @@ impl Broker { /// Clear the failure state on a pending reminder so the /// scheduler picks it up again. No-op when the row is already - /// fresh (`attempt_count == 0`). Returns the number of rows + /// fresh (attempt_count == 0). Returns the number of rows /// affected so callers can distinguish "retried" from "no /// such pending reminder" (already delivered, or wrong id). pub fn reset_reminder_failure(&self, id: i64) -> Result { @@ -649,30 +502,6 @@ impl Broker { } } -/// Idempotent messages-table migrations. Adds `acked_at` and -/// back-fills it for every already-delivered row, so the -/// pre-migration sessions count as "fully handled" and won't be -/// resurfaced by the first `requeue_inflight` after upgrade. -fn ensure_message_columns(conn: &Connection) -> Result<()> { - let has: bool = conn - .prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = 'acked_at'")? - .exists([])?; - if !has { - conn.execute_batch("ALTER TABLE messages ADD COLUMN acked_at INTEGER;") - .context("add messages.acked_at column")?; - // Backfill: treat every existing delivered row as acked. The - // session it was delivered to is gone, so requeue would just - // surface phantom traffic to whatever harness reads next. - conn.execute( - "UPDATE messages SET acked_at = delivered_at \ - WHERE delivered_at IS NOT NULL AND acked_at IS NULL", - [], - ) - .context("backfill messages.acked_at from delivered_at")?; - } - Ok(()) -} - /// Idempotent reminder-table migrations. `ALTER TABLE ADD COLUMN` /// has no `IF NOT EXISTS` form in sqlite, so we probe /// `pragma_table_info` per column. New deploys (table created by @@ -709,179 +538,3 @@ fn now_unix() -> i64 { .and_then(|d| i64::try_from(d.as_secs()).ok()) .unwrap_or(0) } - -#[cfg(test)] -mod tests { - use super::*; - use std::sync::atomic::{AtomicU64, Ordering}; - - /// Per-process counter so each test gets a unique sqlite path even - /// when threads run concurrently. Avoids pulling in a `tempfile` - /// dep just for this one module. - static TEST_COUNTER: AtomicU64 = AtomicU64::new(0); - - struct TmpBroker { - path: std::path::PathBuf, - pub broker: Broker, - } - - impl Drop for TmpBroker { - fn drop(&mut self) { - let _ = std::fs::remove_file(&self.path); - } - } - - fn open_broker() -> TmpBroker { - let n = TEST_COUNTER.fetch_add(1, Ordering::Relaxed); - let pid = std::process::id(); - let path = std::env::temp_dir().join(format!("hive-broker-test-{pid}-{n}.sqlite")); - let _ = std::fs::remove_file(&path); - let broker = Broker::open(&path).expect("open broker"); - TmpBroker { path, broker } - } - - fn msg(from: &str, to: &str, body: &str) -> Message { - Message { - from: from.to_owned(), - to: to.to_owned(), - body: body.to_owned(), - } - } - - /// Happy path: send → recv → `ack_turn` drains the in-memory list - /// and marks the row `acked_at IS NOT NULL`. A second recv finds - /// nothing pending (the row stays in the table for vacuum). - #[test] - fn ack_turn_marks_delivered_rows_acked() { - let h = open_broker(); - let broker = &h.broker; - broker.send(&msg("a", "b", "hi")).unwrap(); - let d = broker.recv("b").unwrap().expect("popped"); - assert_eq!(d.message.body, "hi"); - assert!(!d.redelivered); - assert_eq!(broker.ack_turn("b").unwrap(), 1); - // ack_turn drained the unacked list; calling again is a no-op. - assert_eq!(broker.ack_turn("b").unwrap(), 0); - // Recv finds nothing — the row is now delivered + acked. - assert!(broker.recv("b").unwrap().is_none()); - } - - /// Crash-recovery: send → recv → (no ack) → `requeue_inflight` - /// resets `delivered_at` + tags the next pop as redelivered. After - /// that `ack_turn` closes it out cleanly. - #[test] - fn requeue_inflight_resurfaces_unacked_with_redelivered_flag() { - let h = open_broker(); - let broker = &h.broker; - broker.send(&msg("a", "b", "hi")).unwrap(); - let d1 = broker.recv("b").unwrap().expect("popped"); - assert!(!d1.redelivered); - // Simulate harness crash: never call ack_turn. Now boot the - // new harness — requeue_inflight resurfaces the row. - assert_eq!(broker.requeue_inflight("b").unwrap(), 1); - let d2 = broker.recv("b").unwrap().expect("popped again"); - assert_eq!(d2.message.body, "hi"); - assert!( - d2.redelivered, - "second pop should be tagged redelivered" - ); - assert_eq!(broker.ack_turn("b").unwrap(), 1); - } - - /// Idempotency: a second `requeue_inflight` on the same recipient - /// finds nothing because the prior call already reset - /// `delivered_at` (the row is back in the pending state, not - /// inflight). - #[test] - fn requeue_inflight_is_idempotent() { - let h = open_broker(); - let broker = &h.broker; - broker.send(&msg("a", "b", "hi")).unwrap(); - broker.recv("b").unwrap().expect("popped"); - assert_eq!(broker.requeue_inflight("b").unwrap(), 1); - // Second call: the row is pending (delivered_at IS NULL) so - // nothing matches the inflight filter. - assert_eq!(broker.requeue_inflight("b").unwrap(), 0); - } - - /// Multiple messages, partial drain: pop two, `ack_turn` covers - /// both even though one was popped before the other. - #[test] - fn ack_turn_handles_batch() { - let h = open_broker(); - let broker = &h.broker; - broker.send(&msg("a", "b", "one")).unwrap(); - broker.send(&msg("a", "b", "two")).unwrap(); - broker.send(&msg("a", "b", "three")).unwrap(); - broker.recv("b").unwrap().expect("popped 1"); - broker.recv("b").unwrap().expect("popped 2"); - broker.recv("b").unwrap().expect("popped 3"); - assert_eq!(broker.ack_turn("b").unwrap(), 3); - assert!(broker.recv("b").unwrap().is_none()); - } - - /// Vacuum filter respects the new `acked_at` semantics — a - /// delivered-but-not-acked row is NOT vacuumed regardless of - /// age (the requeue path needs it). - #[test] - fn vacuum_preserves_unacked_inflight_rows() { - let h = open_broker(); - let broker = &h.broker; - broker.send(&msg("a", "b", "stuck")).unwrap(); - broker.recv("b").unwrap().expect("popped"); - // Wide window — should still skip unacked rows. - let removed = broker.vacuum_delivered(-i64::from(u8::MAX)).unwrap(); - assert_eq!(removed, 0, "unacked inflight row must survive vacuum"); - // After ack_turn the row is fair game. - broker.ack_turn("b").unwrap(); - let removed = broker.vacuum_delivered(-i64::from(u8::MAX)).unwrap(); - assert_eq!(removed, 1, "acked row is now vacuumable"); - } - - /// Recv ordering: requeued rows go back into FIFO position - /// (they keep their original id). New sends added after the - /// requeue arrive after them. - #[test] - fn requeued_rows_come_back_in_original_order() { - let h = open_broker(); - let broker = &h.broker; - broker.send(&msg("a", "b", "first")).unwrap(); - broker.send(&msg("a", "b", "second")).unwrap(); - // Pop both, ack neither. - broker.recv("b").unwrap().expect("popped 1"); - broker.recv("b").unwrap().expect("popped 2"); - broker.requeue_inflight("b").unwrap(); - // Now add a brand new message AFTER the requeue. - broker.send(&msg("a", "b", "third")).unwrap(); - let d1 = broker.recv("b").unwrap().expect("re-pop 1"); - assert_eq!(d1.message.body, "first"); - assert!(d1.redelivered); - let d2 = broker.recv("b").unwrap().expect("re-pop 2"); - assert_eq!(d2.message.body, "second"); - assert!(d2.redelivered); - let d3 = broker.recv("b").unwrap().expect("re-pop 3"); - assert_eq!(d3.message.body, "third"); - assert!( - !d3.redelivered, - "fresh-send-after-requeue must NOT be tagged redelivered" - ); - } - - /// Per-recipient isolation: `requeue_inflight("a")` doesn't touch - /// b's inflight rows. - #[test] - fn requeue_inflight_is_per_recipient() { - let h = open_broker(); - let broker = &h.broker; - broker.send(&msg("x", "alice", "for alice")).unwrap(); - broker.send(&msg("x", "bob", "for bob")).unwrap(); - broker.recv("alice").unwrap().expect("popped alice"); - broker.recv("bob").unwrap().expect("popped bob"); - // Requeue only alice. Bob's row stays inflight. - assert_eq!(broker.requeue_inflight("alice").unwrap(), 1); - let d = broker.recv("alice").unwrap().expect("re-pop alice"); - assert!(d.redelivered); - // Bob has nothing pending (his row is still delivered, not requeued). - assert!(broker.recv("bob").unwrap().is_none()); - } -} diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index 11adf9d9..1128c6b6 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -215,7 +215,6 @@ impl Coordinator { /// already have an authoritative timestamp from the db update, /// the tiny skew between "row updated" and "event emitted" is /// presentation-only and doesn't matter to clients. - #[allow(clippy::too_many_arguments)] pub fn emit_approval_resolved( &self, id: i64, @@ -248,7 +247,6 @@ impl Coordinator { /// both operator-targeted (`target = None`) and peer-to-peer /// (`target = Some(agent)`) threads — the dashboard surfaces /// both, distinguishing visually + offering operator override. - #[allow(clippy::too_many_arguments)] pub fn emit_question_added( &self, id: i64, @@ -320,7 +318,7 @@ impl Coordinator { /// resolves to — lifecycle ops, destroy, approve (post-spawn), /// rebuild, meta-update, and the crash-watcher's periodic poll. /// Cheap when nothing changed (one `nixos-container list` + a - /// `HashMap` diff + zero emits). + /// HashMap diff + zero emits). pub async fn rescan_containers_and_emit(self: &Arc) { let fresh = container_view::build_all(self).await; let mut last = self.last_containers.lock().await; diff --git a/hive-c0re/src/dashboard.rs b/hive-c0re/src/dashboard.rs index c5efe5ba..2790cbdf 100644 --- a/hive-c0re/src/dashboard.rs +++ b/hive-c0re/src/dashboard.rs @@ -187,7 +187,7 @@ struct StateSnapshot { meta_inputs: Vec, } -/// `OpQuestion` + computed `question_refs` / `answer_refs`. Built +/// OpQuestion + computed `question_refs` / `answer_refs`. Built /// from the snapshot read; the live channel attaches the same /// fields directly on `QuestionAdded` / `QuestionResolved`. #[derive(Serialize)] @@ -1207,7 +1207,7 @@ pub(crate) fn emit_meta_inputs_snapshot(coord: &Coordinator) { /// natural-language text but aren't part of the path itself. Any /// token starting with `/agents/`, `/shared/`, or /// `/var/lib/hyperhive/{agents,shared}/` is a candidate. The -/// allow-list + `is_file` check happens via the same +/// allow-list + is_file check happens via the same /// `resolve_state_path` helper the read endpoint uses, so the /// security rules can't drift. pub(crate) fn scan_validated_paths(body: &str) -> Vec { @@ -1222,7 +1222,7 @@ pub(crate) fn scan_validated_paths(body: &str) -> Vec { // Trim trailing natural-language punctuation that wouldn't // be part of any real path. Inline rather than via a regex // dep — the set is small and the call is hot. - let token = raw.trim_end_matches([',', ';', ':', ')', ']', '}', '.', '\'', '"']); + let token = raw.trim_end_matches(|c: char| matches!(c, ',' | ';' | ':' | ')' | ']' | '}' | '.' | '\'' | '"')); if token.is_empty() { continue; } @@ -1265,8 +1265,9 @@ async fn get_state_file( let body_bytes = if truncated { &bytes[..MAX_BYTES] } else { &bytes[..] }; let mut body = String::from_utf8_lossy(body_bytes).into_owned(); if truncated { - use std::fmt::Write as _; - let _ = write!(body, "\n\n--- truncated at {MAX_BYTES} of {size} bytes ---\n"); + body.push_str(&format!( + "\n\n--- truncated at {MAX_BYTES} of {size} bytes ---\n" + )); } ([("content-type", "text/plain; charset=utf-8")], body).into_response() } @@ -1574,9 +1575,9 @@ where Fut: std::future::Future>, { let logical = strip_container_prefix(name); - let guard = state.coord.transient_guard(&logical, kind); + let _guard = state.coord.transient_guard(&logical, kind); let result = body(logical.clone()).await; - drop(guard); + drop(_guard); match result { Ok(()) => { extra(state, &logical); diff --git a/hive-c0re/src/dashboard_events.rs b/hive-c0re/src/dashboard_events.rs index 4d5346e6..fe9f57c2 100644 --- a/hive-c0re/src/dashboard_events.rs +++ b/hive-c0re/src/dashboard_events.rs @@ -146,15 +146,15 @@ pub enum DashboardEvent { /// Clients drop the spinner row. TransientCleared { seq: u64, name: String }, /// One container row changed — new container appeared (post-spawn - /// finalise), an existing one flipped `running` / `needs_update` / - /// `sha`, etc. Clients upsert by `container.name`. Payload carries - /// the full row so cold-loaded clients and event-driven clients + /// finalise), an existing one flipped running/needs_update/sha, + /// etc. Clients upsert by `container.name`. Payload carries the + /// full row so cold-loaded clients and event-driven clients /// converge on the same render. /// /// Fired by `Coordinator::rescan_containers_and_emit`, which diffs /// a fresh `nixos-container list`–derived snapshot against the /// last one cached on the coordinator. Mutation sites (lifecycle - /// endpoints, `actions::destroy` / approve, `crash_watch`'s poll loop) + /// endpoints, actions::destroy / approve, crash_watch's poll loop) /// call the rescan after their work lands. ContainerStateChanged { seq: u64, diff --git a/hive-c0re/src/forge.rs b/hive-c0re/src/forge.rs index b82a720e..7018f327 100644 --- a/hive-c0re/src/forge.rs +++ b/hive-c0re/src/forge.rs @@ -165,7 +165,6 @@ async fn ensure_user_exists(name: &str, admin: bool) -> Result<()> { /// monotonic clock so re-issuing doesn't collide with an existing /// token of the same name in the DB. async fn mint_and_persist_token(name: &str, path: &Path) -> Result<()> { - use std::os::unix::fs::PermissionsExt; let token_name = format!( "{TOKEN_NAME_PREFIX}-{}", std::time::SystemTime::now() @@ -191,6 +190,7 @@ async fn mint_and_persist_token(name: &str, path: &Path) -> Result<()> { } std::fs::write(path, format!("{token}\n")) .with_context(|| format!("write token to {}", path.display()))?; + use std::os::unix::fs::PermissionsExt; let _ = std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)); tracing::info!(%name, path = %path.display(), %token_name, "forge: persisted access token"); Ok(()) diff --git a/hive-c0re/src/manager_server.rs b/hive-c0re/src/manager_server.rs index 3cf7d30c..c6e7aa59 100644 --- a/hive-c0re/src/manager_server.rs +++ b/hive-c0re/src/manager_server.rs @@ -138,11 +138,9 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc) -> ManagerResp .recv_blocking(MANAGER_AGENT, manager_recv_timeout(*wait_seconds)) .await { - Ok(Some(d)) => ManagerResponse::Message { - from: d.message.from, - body: d.message.body, - id: d.id, - redelivered: d.redelivered, + Ok(Some(msg)) => ManagerResponse::Message { + from: msg.from, + body: msg.body, }, Ok(None) => ManagerResponse::Empty, Err(e) => ManagerResponse::Err { @@ -234,9 +232,9 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc) -> ManagerResp message: "update: hyperhive_flake has no canonical path".into(), }; }; - let guard = coord.transient_guard(name, crate::coordinator::TransientKind::Rebuilding); + let _guard = coord.transient_guard(name, crate::coordinator::TransientKind::Rebuilding); let result = crate::auto_update::rebuild_agent(coord, name, ¤t_rev).await; - drop(guard); + drop(_guard); match result { Ok(()) => { coord.kick_agent(name, "container rebuilt"); @@ -360,23 +358,6 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc) -> ManagerResp |message| ManagerResponse::Err { message }, |()| ManagerResponse::Ok, ), - ManagerRequest::AckTurn => match coord.broker.ack_turn(MANAGER_AGENT) { - Ok(_n) => ManagerResponse::Ok, - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - }, - ManagerRequest::RequeueInflight => match coord.broker.requeue_inflight(MANAGER_AGENT) { - Ok(n) => { - if n > 0 { - tracing::info!(agent = %MANAGER_AGENT, requeued = %n, "requeued in-flight messages"); - } - ManagerResponse::Ok - } - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - }, } } diff --git a/hive-c0re/src/migrate.rs b/hive-c0re/src/migrate.rs index 43fe95d4..9076aa52 100644 --- a/hive-c0re/src/migrate.rs +++ b/hive-c0re/src/migrate.rs @@ -102,9 +102,9 @@ pub async fn run(coord: &Arc) -> Result<()> { // update activation triggers. Without this, crash_watch // would fire ContainerCrash for every agent here and the // manager would spuriously try to recover them. - let guard = coord.transient_guard(name, crate::coordinator::TransientKind::Rebuilding); + let _guard = coord.transient_guard(name, crate::coordinator::TransientKind::Rebuilding); let result = repoint_container(name).await; - drop(guard); + drop(_guard); if let Err(e) = result { tracing::warn!(%name, error = ?e, "migration: container repoint failed"); all_ok = false; diff --git a/hive-c0re/src/operator_questions.rs b/hive-c0re/src/operator_questions.rs index 5440c115..ce9a2094 100644 --- a/hive-c0re/src/operator_questions.rs +++ b/hive-c0re/src/operator_questions.rs @@ -284,8 +284,7 @@ impl OperatorQuestions { ORDER BY answered_at DESC LIMIT ?1", )?; - let limit_i = i64::try_from(limit).unwrap_or(i64::MAX); - let rows = stmt.query_map(params![limit_i], row_to_question)?; + let rows = stmt.query_map(params![limit as i64], row_to_question)?; rows.collect::>>() .map_err(Into::into) } diff --git a/hive-sh4re/src/lib.rs b/hive-sh4re/src/lib.rs index 4cd729ce..47c84fb3 100644 --- a/hive-sh4re/src/lib.rs +++ b/hive-sh4re/src/lib.rs @@ -358,28 +358,6 @@ pub enum AgentRequest { /// row. The manager surface uses the same wire variant but /// accepts any id. CancelLooseEnd { kind: CancelLooseEndKind, id: i64 }, - /// Mark every message popped by this agent since the last `AckTurn` - /// as fully handled. Fired by the harness after `TurnOutcome::Ok` - /// — claude doesn't see this surface, it's harness↔broker only. - /// On `TurnOutcome::Failed` the harness intentionally skips this - /// call, so the unacked rows stay in-flight in the DB and get - /// requeued by the next `RequeueInflight` on harness boot. Tracks - /// the popped-id list in-memory on the broker side; no payload - /// needed (the broker knows which ids it handed to this - /// recipient). - AckTurn, - /// Requeue every message the broker handed to this agent that - /// never got acked. Fired by the harness exactly once at boot, - /// before entering the serve loop — catches the - /// crashed-mid-turn / OOM-killed / container-restarted cases - /// where a previous harness session popped messages but never - /// drove them to a clean turn-end. Resets `delivered_at` on each - /// row back to NULL (so the next `Recv` pops it) and remembers - /// the id in a per-recipient in-memory set so the next `Recv` - /// can tag the message with `redelivered: true` (the harness - /// then prepends a "may already be handled" hint to the wake - /// prompt). Idempotent + cheap when there's nothing in flight. - RequeueInflight, } /// Responses on a per-agent socket. @@ -390,22 +368,8 @@ pub enum AgentResponse { Ok, /// Either `Send` failed or `Recv` errored. Err { message: String }, - /// `Recv` produced a message. `id` is the broker's row id — opaque - /// to claude (the MCP surface strips it before handing the body - /// to the model) but tracked by the harness so the broker's - /// in-memory unacked list can be drained on `AckTurn`. When - /// `redelivered = true` this row was popped earlier, never - /// acked (turn crash / OOM / restart), and resurfaced by - /// `RequeueInflight` — the harness prepends a "may already be - /// handled" hint to the wake prompt so claude can DTRT. - Message { - from: String, - body: String, - #[serde(default)] - id: i64, - #[serde(default)] - redelivered: bool, - }, + /// `Recv` produced a message. + Message { from: String, body: String }, /// `Recv` found nothing pending. Empty, /// `Status` result: how many pending messages are in this agent's inbox. @@ -704,13 +668,6 @@ pub enum ManagerRequest { /// can cancel any row (no owner check) — same dispatch as /// `AgentRequest::CancelLooseEnd` but with privileged auth. CancelLooseEnd { kind: CancelLooseEndKind, id: i64 }, - /// Mirror of `AgentRequest::AckTurn` on the manager surface — fired - /// by the manager harness after `TurnOutcome::Ok` to close out - /// every message popped during the turn. - AckTurn, - /// Mirror of `AgentRequest::RequeueInflight` on the manager - /// surface — fired exactly once on manager harness boot. - RequeueInflight, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -720,18 +677,9 @@ pub enum ManagerResponse { Err { message: String, }, - /// Same delivery shape as `AgentResponse::Message` — `id` + - /// `redelivered` carry the broker's row id and the - /// "previously popped, not acked" flag through the manager - /// surface so the manager harness drives the same - /// requeue-with-hint flow as a sub-agent. Message { from: String, body: String, - #[serde(default)] - id: i64, - #[serde(default)] - redelivered: bool, }, Empty, Status {