repoint per-agent stats reminder rollup off the broker (#2635 inc 1)

This commit is contained in:
damocles 2026-07-22 21:02:44 +02:00 committed by mara
commit 8b35b8c4af

View file

@ -12,38 +12,55 @@ pub(super) struct StatsQuery {
}
pub(super) async fn api_stats(
State(state): State<AppState>,
State(_state): State<AppState>,
axum::extract::Query(q): axum::extract::Query<StatsQuery>,
) -> axum::Json<crate::stats::Snapshot> {
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<hive_sh4re::ReminderStats> {
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<hive_sh4re::ReminderStats> {
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).