From 8b35b8c4af05078584eae9cd6c58029768bc0350 Mon Sep 17 00:00:00 2001 From: damocles Date: Wed, 22 Jul 2026 21:02:44 +0200 Subject: [PATCH] repoint per-agent stats reminder rollup off the broker (#2635 inc 1) --- hive-agent/src/web_ui/stats.rs | 61 ++++++++++++++++++++++------------ 1 file changed, 39 insertions(+), 22 deletions(-) diff --git a/hive-agent/src/web_ui/stats.rs b/hive-agent/src/web_ui/stats.rs index 8bf59747..1852e873 100644 --- a/hive-agent/src/web_ui/stats.rs +++ b/hive-agent/src/web_ui/stats.rs @@ -12,38 +12,55 @@ pub(super) struct StatsQuery { } pub(super) async fn api_stats( - State(state): State, + State(_state): State, axum::extract::Query(q): axum::extract::Query, ) -> axum::Json { let window = crate::stats::Window::parse(q.window.as_deref().unwrap_or("24h")); let mut snapshot = crate::stats::snapshot_default(window); - // Pass the window span to the reminder-stats RPC so the broker - // filters its counts to the same time range as the chart data. + // Pass the window span so the local reminder rollup filters its counts + // to the same time range as the chart data. let window_secs = window.span_secs(); let window_secs_u = u64::try_from(window_secs).unwrap_or(0); - snapshot.reminder_stats = fetch_reminder_stats(&state.socket, window_secs_u).await; + snapshot.reminder_stats = fetch_reminder_stats(window_secs_u).await; axum::Json(snapshot) } -/// Fetch reminder activity stats from the broker via the per-agent / manager -/// socket. Returns None on any transport / decode failure — the stats are -/// decorative, not authoritative. -async fn fetch_reminder_stats( - socket: &std::path::Path, - window_secs: u64, -) -> Option { - match super::broker_request( - socket, - &hive_core_agent_sock::Request::ReminderRollup { - since_secs: window_secs, - agent: None, - }, - ) - .await - { - Ok(hive_core_agent_sock::Response::ReminderRollup(stats)) => Some(stats), - _ => None, +/// Fetch reminder activity stats from the harness-local reminder store over +/// `HIVE_AGENT_SOCKET` (#2635 inc 1 — was a broker RPC before reminders +/// moved in-container). Returns `None` on any transport / decode failure or +/// when the socket is unset — the stats are decorative, not authoritative. +async fn fetch_reminder_stats(window_secs: u64) -> Option { + use tokio::io::{AsyncBufReadExt as _, AsyncWriteExt as _, BufReader}; + use tokio::net::UnixStream; + + let socket_path = std::env::var_os("HIVE_AGENT_SOCKET").map(std::path::PathBuf::from)?; + if !socket_path.exists() { + return None; } + tokio::time::timeout(std::time::Duration::from_secs(3), async move { + let mut stream = UnixStream::connect(&socket_path).await?; + let req = hive_agent_sock::Request::ReminderRollup { + since_secs: window_secs, + }; + let mut line = serde_json::to_string(&req)?; + line.push('\n'); + stream.write_all(line.as_bytes()).await?; + stream.flush().await?; + let mut lines = BufReader::new(stream).lines(); + let resp_line = lines + .next_line() + .await? + .ok_or_else(|| anyhow::anyhow!("agent socket closed without response"))?; + let resp: hive_agent_sock::Response = serde_json::from_str(&resp_line)?; + anyhow::Ok(match resp { + hive_agent_sock::Response::ReminderRollup { stats } => Some(stats), + _ => None, + }) + }) + .await + .ok()? + .ok() + .flatten() } /// `GET /api/todos` — snapshot of this agent's local todos (loose-ends v2).