From e11e8294a38ed496a3f578ab60fa0fdcdc692833 Mon Sep 17 00:00:00 2001 From: damocles Date: Wed, 22 Jul 2026 21:14:31 +0200 Subject: [PATCH] repoint post_turn_counts + dedup web_ui stats onto todo_server::dial (#2635 inc 1) --- hive-agent/src/main.rs | 16 ++++---- hive-agent/src/todo_server.rs | 26 ++++++++++++ hive-agent/src/web_ui/stats.rs | 73 ++++++---------------------------- 3 files changed, 46 insertions(+), 69 deletions(-) diff --git a/hive-agent/src/main.rs b/hive-agent/src/main.rs index a6556c05..b18307cb 100644 --- a/hive-agent/src/main.rs +++ b/hive-agent/src/main.rs @@ -332,15 +332,13 @@ impl Surface for AgentSurface { Ok(Response::LooseEnds { loose_ends }) => u64::try_from(loose_ends.len()).ok(), _ => None, }; - let reminders = match client::request::<_, Response>( - socket, - &Request::CountPendingReminders { agent: None }, - ) - .await - { - Ok(Response::PendingRemindersCount { count }) => Some(count), - _ => None, - }; + // Reminders are harness-local (#2635 inc 1) — dial the in-agent + // socket directly instead of the broker. + let reminders = + match todo_server::dial(&hive_agent_sock::Request::CountPendingReminders).await { + Some(hive_agent_sock::Response::PendingRemindersCount { count }) => Some(count), + _ => None, + }; (threads, reminders) } diff --git a/hive-agent/src/todo_server.rs b/hive-agent/src/todo_server.rs index 5ce462e4..5d3a0ccf 100644 --- a/hive-agent/src/todo_server.rs +++ b/hive-agent/src/todo_server.rs @@ -35,6 +35,32 @@ fn socket_path() -> Option { .map(PathBuf::from) } +/// Dial this same in-agent socket from elsewhere IN THIS PROCESS — the +/// harness's own `web_ui` endpoints (`api_todos`, `/api/stats` reminder +/// rollup) and `Surface::post_turn_counts` all need read access to the +/// `Todos`/`Reminders` stores this module owns behind an `Arc` on a +/// different spawned task, and a loopback dial is simpler than threading +/// those `Arc`s through every caller. One-shot, best-effort (no retry, +/// unlike the broker client): a connect failure means the socket server +/// task itself isn't up, which a retry within one call wouldn't fix. +pub(crate) async fn dial(req: &Request) -> Option { + let path = socket_path()?; + if !path.exists() { + return None; + } + tokio::time::timeout(std::time::Duration::from_secs(3), async move { + let mut stream = UnixStream::connect(&path).await.ok()?; + let mut line = serde_json::to_string(req).ok()?; + line.push('\n'); + stream.write_all(line.as_bytes()).await.ok()?; + let mut lines = BufReader::new(stream).lines(); + let resp_line = lines.next_line().await.ok()??; + serde_json::from_str(&resp_line).ok() + }) + .await + .ok()? +} + /// Run the in-agent socket server: bind + accept loop, one request/response /// line per connection. A no-op (returns `Ok`) when `HIVE_AGENT_SOCKET` is /// unset, so a standalone harness without producers just skips it. diff --git a/hive-agent/src/web_ui/stats.rs b/hive-agent/src/web_ui/stats.rs index 1852e873..db3271f3 100644 --- a/hive-agent/src/web_ui/stats.rs +++ b/hive-agent/src/web_ui/stats.rs @@ -30,37 +30,14 @@ pub(super) async fn api_stats( /// 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, - }) + match crate::todo_server::dial(&hive_agent_sock::Request::ReminderRollup { + since_secs: window_secs, }) - .await - .ok()? - .ok() - .flatten() + .await? + { + hive_agent_sock::Response::ReminderRollup { stats } => Some(stats), + _ => None, + } } /// `GET /api/todos` — snapshot of this agent's local todos (loose-ends v2). @@ -70,36 +47,12 @@ async fn fetch_reminder_stats(window_secs: u64) -> Option Response { - use tokio::io::{AsyncBufReadExt as _, AsyncWriteExt as _, BufReader}; - use tokio::net::UnixStream; - - let socket_path = match std::env::var_os("HIVE_AGENT_SOCKET") { - Some(p) => std::path::PathBuf::from(p), - None => return axum::Json(serde_json::json!({ "todos": [] })).into_response(), - }; - if !socket_path.exists() { - return axum::Json(serde_json::json!({ "todos": [] })).into_response(); - } - let todos = 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::ListTodos { subsystem: None }; - 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::LooseEnds { loose_ends } => loose_ends, + let todos = + match crate::todo_server::dial(&hive_agent_sock::Request::ListTodos { subsystem: None }) + .await + { + Some(hive_agent_sock::Response::LooseEnds { loose_ends }) => loose_ends, _ => Vec::new(), - }) - }) - .await - .unwrap_or_else(|_| Err(anyhow::anyhow!("timeout"))) - .unwrap_or_default(); + }; axum::Json(serde_json::json!({ "todos": todos })).into_response() }