From 795dd882bbac9a618337644ae2f77cbad4b57fdc Mon Sep 17 00:00:00 2001 From: damocles Date: Mon, 20 Jul 2026 21:58:28 +0200 Subject: [PATCH] refactor(#2569): rename hive-agent-sock to hive-core-agent-sock --- Cargo.lock | 26 ++-- Cargo.toml | 4 +- hive-agent-mcp/Cargo.toml | 2 +- hive-agent-mcp/src/mcp/mod.rs | 75 +++++----- hive-agent-mcp/src/mcp/render.rs | 31 ++-- hive-agent-wake/Cargo.toml | 2 +- hive-agent-wake/src/main.rs | 2 +- hive-agent/Cargo.toml | 2 +- hive-agent/src/forge_notify.rs | 9 +- hive-agent/src/main.rs | 2 +- hive-agent/src/web_ui/actions.rs | 8 +- hive-agent/src/web_ui/mod.rs | 6 +- hive-agent/src/web_ui/state.rs | 9 +- hive-agent/src/web_ui/stats.rs | 10 +- hive-bash-mcp/Cargo.toml | 2 +- hive-bash-mcp/src/runner.rs | 2 +- hive-c0re/Cargo.toml | 2 +- .../src/socket_server/config_approvals.rs | 2 +- .../src/socket_server/lifecycle_handlers.rs | 2 +- hive-c0re/src/socket_server/mod.rs | 137 ++++++++++-------- hive-c0re/src/socket_server/reminders.rs | 2 +- hive-c0re/src/socket_server/schedules.rs | 2 +- .../Cargo.toml | 2 +- .../src/lib.rs | 0 hive-matrix-mcp/src/wake.rs | 2 +- hive-types/src/lib.rs | 2 +- 26 files changed, 184 insertions(+), 161 deletions(-) rename {hive-agent-sock => hive-core-agent-sock}/Cargo.toml (83%) rename {hive-agent-sock => hive-core-agent-sock}/src/lib.rs (100%) diff --git a/Cargo.lock b/Cargo.lock index b69af92e..7d573164 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1523,8 +1523,8 @@ dependencies = [ "clap", "forgejo-api", "futures-util", - "hive-agent-sock", "hive-claude", + "hive-core-agent-sock", "hive-sh4re", "http-body-util", "hyper", @@ -1552,7 +1552,7 @@ dependencies = [ "anyhow", "axum", "clap", - "hive-agent-sock", + "hive-core-agent-sock", "hive-sh4re", "rmcp", "serde", @@ -1562,21 +1562,13 @@ dependencies = [ "tracing-subscriber", ] -[[package]] -name = "hive-agent-sock" -version = "0.1.0" -dependencies = [ - "hive-sh4re", - "serde", -] - [[package]] name = "hive-agent-wake" version = "0.1.0" dependencies = [ "anyhow", "clap", - "hive-agent-sock", + "hive-core-agent-sock", "serde", "serde_json", "tokio", @@ -1589,7 +1581,7 @@ name = "hive-bash-mcp" version = "0.1.0" dependencies = [ "anyhow", - "hive-agent-sock", + "hive-core-agent-sock", "hive-sh4re", "libc", "rmcp", @@ -1615,7 +1607,7 @@ dependencies = [ "clap-markdown", "clap_complete", "forgejo-api", - "hive-agent-sock", + "hive-core-agent-sock", "hive-host-sock", "hive-priv-sock", "hive-sh4re", @@ -1653,6 +1645,14 @@ dependencies = [ "tokio", ] +[[package]] +name = "hive-core-agent-sock" +version = "0.1.0" +dependencies = [ + "hive-sh4re", + "serde", +] + [[package]] name = "hive-forge" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index cc38c6fb..549abeb5 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -3,7 +3,7 @@ resolver = "3" members = [ "hive-agent", "hive-agent-mcp", - "hive-agent-sock", + "hive-core-agent-sock", "hive-agent-wake", "hive-bash-mcp", "hive-c0re", @@ -51,7 +51,7 @@ clap = { version = "4", features = ["derive"] } clap_complete = "4" indicatif = "0.18" hive-sh4re = { path = "hive-sh4re" } -hive-agent-sock = { path = "hive-agent-sock" } +hive-core-agent-sock = { path = "hive-core-agent-sock" } hive-claude = { path = "hive-claude" } hive-host-sock = { path = "hive-host-sock" } hive-priv-sock = { path = "hive-priv-sock" } diff --git a/hive-agent-mcp/Cargo.toml b/hive-agent-mcp/Cargo.toml index 7c4ee9c2..6fad2eb9 100644 --- a/hive-agent-mcp/Cargo.toml +++ b/hive-agent-mcp/Cargo.toml @@ -14,7 +14,7 @@ workspace = true anyhow.workspace = true axum.workspace = true clap.workspace = true -hive-agent-sock.workspace = true +hive-core-agent-sock.workspace = true hive-sh4re.workspace = true rmcp.workspace = true serde.workspace = true diff --git a/hive-agent-mcp/src/mcp/mod.rs b/hive-agent-mcp/src/mcp/mod.rs index 6cf8cb1b..5ac7c35d 100644 --- a/hive-agent-mcp/src/mcp/mod.rs +++ b/hive-agent-mcp/src/mcp/mod.rs @@ -121,9 +121,10 @@ impl AgentServer { /// covers both. async fn dispatch( &self, - req: hive_agent_sock::Request, - ) -> (Result, u32) { - match client::request_retried::<_, hive_agent_sock::Response>(&self.socket, &req).await { + req: hive_core_agent_sock::Request, + ) -> (Result, u32) { + match client::request_retried::<_, hive_core_agent_sock::Response>(&self.socket, &req).await + { Ok((r, n)) => (Ok(r), n), Err(e) => (Err(e), 0), } @@ -150,7 +151,7 @@ impl AgentServer { } run_tool_envelope("send", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::Send { + .dispatch(hive_core_agent_sock::Request::Send { to: args.to, body: args.body, in_reply_to: args.in_reply_to, @@ -181,7 +182,7 @@ impl AgentServer { let log = format!("{args:?}"); run_tool_envelope("ask", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::Ask { + .dispatch(hive_core_agent_sock::Request::Ask { question: args.question, options: args.options, multi: args.multi, @@ -190,7 +191,7 @@ impl AgentServer { }) .await; let s = match resp { - Ok(hive_agent_sock::Response::QuestionQueued { id }) => format!( + Ok(hive_core_agent_sock::Response::QuestionQueued { id }) => format!( "question queued (id={id}); answer will arrive as a system \ `question_answered` event in your inbox" ), @@ -215,7 +216,7 @@ impl AgentServer { let id = args.id; run_tool_envelope("answer", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::Answer { + .dispatch(hive_core_agent_sock::Request::Answer { id, answer: args.answer, }) @@ -253,7 +254,7 @@ impl AgentServer { run_tool_envelope("recv", log, async move { let waited = args.wait_seconds.is_some_and(|w| w > 0); let (resp, retries) = self - .dispatch(hive_agent_sock::Request::Recv { + .dispatch(hive_core_agent_sock::Request::Recv { wait_seconds: args.wait_seconds, max: args.max, }) @@ -278,10 +279,10 @@ impl AgentServer { let log = format!("{args:?}"); run_tool_envelope("ack_until", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::AckUntil { up_to: args.up_to }) + .dispatch(hive_core_agent_sock::Request::AckUntil { up_to: args.up_to }) .await; let rendered = match resp { - Ok(hive_agent_sock::Response::Acked { count }) => { + Ok(hive_core_agent_sock::Response::Acked { count }) => { format!("acked {count} message(s) up to id {}", args.up_to) } other => reply_err(other, "ack_until"), @@ -310,11 +311,11 @@ impl AgentServer { run_tool_envelope("get_loose_ends", String::new(), async move { let is_self_query = args.agent.is_none(); let (resp, retries) = self - .dispatch(hive_agent_sock::Request::GetLooseEnds { agent: args.agent }) + .dispatch(hive_core_agent_sock::Request::GetLooseEnds { agent: args.agent }) .await; // Extract the vec so we can augment before rendering. let mut loose_ends = match resp { - Ok(hive_agent_sock::Response::LooseEnds { loose_ends }) => loose_ends, + Ok(hive_core_agent_sock::Response::LooseEnds { loose_ends }) => loose_ends, other => return annotate_retries(reply_err(other, "get_loose_ends"), retries), }; // Prepend matrix unread entry for self-queries only (can't @@ -363,7 +364,7 @@ impl AgentServer { return e; } let (resp, retries) = self - .dispatch(hive_agent_sock::Request::SetStatus { text: args.text }) + .dispatch(hive_core_agent_sock::Request::SetStatus { text: args.text }) .await; annotate_retries( format_ack(resp, "set_status", "status updated".to_owned()), @@ -395,7 +396,7 @@ impl AgentServer { let log = args.name.clone().unwrap_or_else(|| "".to_owned()); run_tool_envelope("get_agent_meta", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::GetAgentMeta { name: args.name }) + .dispatch(hive_core_agent_sock::Request::GetAgentMeta { name: args.name }) .await; annotate_retries(format_agent_meta(resp), retries) }) @@ -423,7 +424,7 @@ impl AgentServer { }; let kind_label = loose_end_kind_label(kind); let (resp, retries) = self - .dispatch(hive_agent_sock::Request::CancelLooseEnd { kind, id }) + .dispatch(hive_core_agent_sock::Request::CancelLooseEnd { kind, id }) .await; annotate_retries( format_ack( @@ -450,10 +451,10 @@ impl AgentServer { let log = format!("{args:?}"); run_tool_envelope("create_repo", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::CreateRepo { repo: args.repo }) + .dispatch(hive_core_agent_sock::Request::CreateRepo { repo: args.repo }) .await; let s = match resp { - Ok(hive_agent_sock::Response::RepoCreated { + Ok(hive_core_agent_sock::Response::RepoCreated { full_name, clone_url, }) => format!("created repo {full_name} — clone: {clone_url}"), @@ -493,7 +494,7 @@ impl AgentServer { (None, Some(t)) => hive_sh4re::ReminderTiming::At { unix_timestamp: t }, }; let (resp, retries) = self - .dispatch(hive_agent_sock::Request::Remind { + .dispatch(hive_core_agent_sock::Request::Remind { message: args.message, timing, file_path: args.file_path, @@ -547,7 +548,7 @@ impl AgentServer { let name = args.name.clone(); run_tool_envelope("restart", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::Restart { name: args.name }) + .dispatch(hive_core_agent_sock::Request::Restart { name: args.name }) .await; annotate_retries( format_ack(resp, "restart", format!("restarted {name}")), @@ -569,7 +570,7 @@ impl AgentServer { let name = args.name.clone(); run_tool_envelope("kill", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::Kill { name: args.name }) + .dispatch(hive_core_agent_sock::Request::Kill { name: args.name }) .await; annotate_retries(format_ack(resp, "kill", format!("killed {name}")), retries) }) @@ -590,7 +591,7 @@ impl AgentServer { let name = args.name.clone(); run_tool_envelope("update", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::Update { name: args.name }) + .dispatch(hive_core_agent_sock::Request::Update { name: args.name }) .await; annotate_retries( format_ack(resp, "update", format!("updated {name}")), @@ -613,10 +614,10 @@ impl AgentServer { async fn list_containers(&self) -> String { run_tool_envelope("list_containers", String::new(), async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::ListDescendants) + .dispatch(hive_core_agent_sock::Request::ListDescendants) .await; let body = match resp { - Ok(hive_agent_sock::Response::Containers { containers }) => { + Ok(hive_core_agent_sock::Response::Containers { containers }) => { if containers.is_empty() { "no descendant containers".to_owned() } else { @@ -659,7 +660,7 @@ impl AgentServer { let log = format!("{args:?}"); run_tool_envelope("get_host_journal", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::GetHostJournal { + .dispatch(hive_core_agent_sock::Request::GetHostJournal { unit: args.unit, container: args.container, lines: args.lines, @@ -670,7 +671,7 @@ impl AgentServer { }) .await; let result = match resp { - Ok(hive_agent_sock::Response::HostJournal { content }) => content, + Ok(hive_core_agent_sock::Response::HostJournal { content }) => content, other => reply_err(other, "get_host_journal"), }; annotate_retries(result, retries) @@ -699,7 +700,7 @@ impl AgentServer { let name = args.name.clone(); run_tool_envelope("request_init_config", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::RequestInitConfig { + .dispatch(hive_core_agent_sock::Request::RequestInitConfig { name: args.name, description: args.description, }) @@ -727,7 +728,7 @@ impl AgentServer { let name = args.name.clone(); run_tool_envelope("start", log, async move { let (resp, retries) = self - .dispatch(hive_agent_sock::Request::Start { name: args.name }) + .dispatch(hive_core_agent_sock::Request::Start { name: args.name }) .await; annotate_retries( format_ack(resp, "start", format!("started {name}")), @@ -750,13 +751,13 @@ impl AgentServer { run_tool_envelope("get_logs", log, async move { let lines = args.lines.map(|n| n.min(500)); let (resp, retries) = self - .dispatch(hive_agent_sock::Request::GetLogs { + .dispatch(hive_core_agent_sock::Request::GetLogs { agent: agent.clone(), lines, }) .await; let s = match resp { - Ok(hive_agent_sock::Response::Logs { content }) => { + Ok(hive_core_agent_sock::Response::Logs { content }) => { if content.is_empty() { format!("(no journal output for {agent})") } else { @@ -790,7 +791,7 @@ impl AgentServer { args.inputs.join(", ") }; let (resp, retries) = self - .dispatch(hive_agent_sock::Request::RequestUpdateMetaInputs { + .dispatch(hive_core_agent_sock::Request::RequestUpdateMetaInputs { inputs: args.inputs, description: args.description, }) @@ -829,7 +830,7 @@ impl AgentServer { run_tool_envelope("request_schedule_prompt", log, async move { let target_count = args.targets.len(); let (resp, retries) = self - .dispatch(hive_agent_sock::Request::RequestSchedulePrompt( + .dispatch(hive_core_agent_sock::Request::RequestSchedulePrompt( hive_sh4re::SchedulePromptPayload { targets: args.targets, body: args.body, @@ -865,7 +866,7 @@ impl AgentServer { run_tool_envelope("fire_schedule_now", log, async move { let id = args.id; let (resp, retries) = self - .dispatch(hive_agent_sock::Request::FireScheduleNow { id }) + .dispatch(hive_core_agent_sock::Request::FireScheduleNow { id }) .await; annotate_retries( format_ack(resp, "fire_schedule_now", format!("fired #{id} now")), @@ -888,7 +889,7 @@ impl AgentServer { run_tool_envelope("cancel_schedule", log, async move { let id = args.id; let (resp, retries) = self - .dispatch(hive_agent_sock::Request::CancelSchedule { + .dispatch(hive_core_agent_sock::Request::CancelSchedule { id: args.id, targets: args.targets, }) @@ -920,7 +921,7 @@ impl AgentServer { run_tool_envelope("edit_schedule", log, async move { let id = args.id; let (resp, retries) = self - .dispatch(hive_agent_sock::Request::EditSchedule { + .dispatch(hive_core_agent_sock::Request::EditSchedule { id: args.id, body: args.body, description: args.description.map(Some), @@ -947,9 +948,11 @@ impl AgentServer { )] async fn list_schedules(&self) -> String { run_tool_envelope("list_schedules", String::new(), async move { - let (resp, retries) = self.dispatch(hive_agent_sock::Request::ListSchedules).await; + let (resp, retries) = self + .dispatch(hive_core_agent_sock::Request::ListSchedules) + .await; let body = match resp { - Ok(hive_agent_sock::Response::Schedules { schedules }) => { + Ok(hive_core_agent_sock::Response::Schedules { schedules }) => { serde_json::to_string(&schedules) .unwrap_or_else(|e| format!("list_schedules: serialise: {e:#}")) } diff --git a/hive-agent-mcp/src/mcp/render.rs b/hive-agent-mcp/src/mcp/render.rs index ccce7934..dddc9eca 100644 --- a/hive-agent-mcp/src/mcp/render.rs +++ b/hive-agent-mcp/src/mcp/render.rs @@ -11,11 +11,11 @@ /// everything else here via a catch-all arm (`other => reply_err(other, tool)`), /// so the triplet lives in exactly one place. pub(super) fn reply_err( - resp: Result, + resp: Result, tool: &str, ) -> String { match resp { - Ok(hive_agent_sock::Response::Err { message }) => format!("{tool} failed: {message}"), + Ok(hive_core_agent_sock::Response::Err { message }) => format!("{tool} failed: {message}"), Ok(other) => format!("{tool} unexpected response: {other:?}"), Err(e) => format!("{tool} transport error: {e:#}"), } @@ -26,12 +26,12 @@ pub(super) fn reply_err( /// behavior. #[must_use] pub fn format_ack( - resp: Result, + resp: Result, tool: &str, ok_msg: String, ) -> String { match resp { - Ok(hive_agent_sock::Response::Ok) => ok_msg, + Ok(hive_core_agent_sock::Response::Ok) => ok_msg, other => reply_err(other, tool), } } @@ -47,9 +47,12 @@ pub fn format_ack( /// so the model can tell where one ends and the next begins; /// per-message redelivery banners included. #[must_use] -pub fn format_recv(resp: Result, waited: bool) -> String { +pub fn format_recv( + resp: Result, + waited: bool, +) -> String { match resp { - Ok(hive_agent_sock::Response::Messages { + Ok(hive_core_agent_sock::Response::Messages { messages, remaining, }) => render_recv_messages(&messages, remaining, waited), @@ -60,7 +63,7 @@ pub fn format_recv(resp: Result, waite // stop unmissably tells claude to flush + end. `remaining` is forced // to 0 — the inbox is fenced, so a "N more pending" hint would be // misleading. - Ok(hive_agent_sock::Response::GracefulStop) => { + Ok(hive_core_agent_sock::Response::GracefulStop) => { render_recv_messages(&[graceful_stop_message()], 0, waited) } other => reply_err(other, "recv"), @@ -357,9 +360,9 @@ pub(super) fn loose_end_kind_label(kind: hive_sh4re::CancelLooseEndKind) -> &'st /// `running: no` line tells the caller WHY. See /// `docs/turn-loop/mcp.md::Core tools` (`get_agent_meta`). #[must_use] -pub fn format_agent_meta(resp: Result) -> String { +pub fn format_agent_meta(resp: Result) -> String { match resp { - Ok(hive_agent_sock::Response::AgentMeta { + Ok(hive_core_agent_sock::Response::AgentMeta { name, running, hyperhive_rev, @@ -480,7 +483,7 @@ mod tests { #[test] fn empty_recv_after_wait_appends_idle_hint() { let out = format_recv( - Ok(hive_agent_sock::Response::Messages { + Ok(hive_core_agent_sock::Response::Messages { messages: vec![], remaining: 0, }), @@ -493,7 +496,7 @@ mod tests { #[test] fn empty_recv_without_wait_has_no_hint() { let out = format_recv( - Ok(hive_agent_sock::Response::Messages { + Ok(hive_core_agent_sock::Response::Messages { messages: vec![], remaining: 0, }), @@ -505,7 +508,7 @@ mod tests { #[test] fn single_recv_with_remaining_appends_pending_hint() { let out = format_recv( - Ok(hive_agent_sock::Response::Messages { + Ok(hive_core_agent_sock::Response::Messages { messages: vec![msg(7, "alice", "hi")], remaining: 3, }), @@ -519,7 +522,7 @@ mod tests { #[test] fn single_recv_no_remaining_has_no_pending_hint() { let out = format_recv( - Ok(hive_agent_sock::Response::Messages { + Ok(hive_core_agent_sock::Response::Messages { messages: vec![msg(7, "alice", "hi")], remaining: 0, }), @@ -531,7 +534,7 @@ mod tests { #[test] fn batch_recv_with_remaining_appends_pending_hint_once() { let out = format_recv( - Ok(hive_agent_sock::Response::Messages { + Ok(hive_core_agent_sock::Response::Messages { messages: vec![msg(7, "alice", "hi"), msg(8, "bob", "yo")], remaining: 9, }), diff --git a/hive-agent-wake/Cargo.toml b/hive-agent-wake/Cargo.toml index c672f09d..2fa7aa0b 100644 --- a/hive-agent-wake/Cargo.toml +++ b/hive-agent-wake/Cargo.toml @@ -13,7 +13,7 @@ workspace = true [dependencies] anyhow.workspace = true clap.workspace = true -hive-agent-sock.workspace = true +hive-core-agent-sock.workspace = true serde.workspace = true serde_json.workspace = true tokio.workspace = true diff --git a/hive-agent-wake/src/main.rs b/hive-agent-wake/src/main.rs index d21cfca9..6bc12cec 100644 --- a/hive-agent-wake/src/main.rs +++ b/hive-agent-wake/src/main.rs @@ -13,7 +13,7 @@ use std::path::PathBuf; use anyhow::Result; use clap::Parser; -use hive_agent_sock::{Request, Response}; +use hive_core_agent_sock::{Request, Response}; /// Per-agent MCP socket, bind-mounted from the host into every container. const DEFAULT_SOCKET: &str = "/run/hive/mcp.sock"; diff --git a/hive-agent/Cargo.toml b/hive-agent/Cargo.toml index f16a7e9a..2b0c74c1 100644 --- a/hive-agent/Cargo.toml +++ b/hive-agent/Cargo.toml @@ -19,7 +19,7 @@ time.workspace = true futures-util = "0.3" clap.workspace = true hive-claude.workspace = true -hive-agent-sock.workspace = true +hive-core-agent-sock.workspace = true hive-sh4re.workspace = true rmcp.workspace = true rusqlite.workspace = true diff --git a/hive-agent/src/forge_notify.rs b/hive-agent/src/forge_notify.rs index 2ecae583..9f650877 100644 --- a/hive-agent/src/forge_notify.rs +++ b/hive-agent/src/forge_notify.rs @@ -1064,13 +1064,14 @@ async fn poll_once( continue; }; - let req = hive_agent_sock::Request::Wake { + let req = hive_core_agent_sock::Request::Wake { from: "forge".to_owned(), body, }; - let deliver_result = crate::client::request::<_, hive_agent_sock::Response>(socket, &req) - .await - .map(|_| ()); + let deliver_result = + crate::client::request::<_, hive_core_agent_sock::Response>(socket, &req) + .await + .map(|_| ()); match deliver_result { Ok(()) => { debug!(%id, "forge_notify: delivered"); diff --git a/hive-agent/src/main.rs b/hive-agent/src/main.rs index 23ee0e7d..69b8b0b1 100644 --- a/hive-agent/src/main.rs +++ b/hive-agent/src/main.rs @@ -43,7 +43,7 @@ use crate::login::LoginState; use crate::turn_stats::TurnStats; use anyhow::Result; use clap::Parser; -use hive_agent_sock::{Request, Response}; +use hive_core_agent_sock::{Request, Response}; use hive_sh4re::{HelperEvent, SYSTEM_SENDER}; #[derive(Parser)] diff --git a/hive-agent/src/web_ui/actions.rs b/hive-agent/src/web_ui/actions.rs index 7e6de725..de060f92 100644 --- a/hive-agent/src/web_ui/actions.rs +++ b/hive-agent/src/web_ui/actions.rs @@ -25,7 +25,7 @@ pub(super) async fn post_send( } match super::broker_request( &state.socket, - &hive_agent_sock::Request::OperatorMsg { body }, + &hive_core_agent_sock::Request::OperatorMsg { body }, ) .await { @@ -35,8 +35,10 @@ pub(super) async fn post_send( // resulting `TurnStart` SSE event drives the terminal + the // inbox row gets consumed by the time `TurnEnd` fires the // existing turn-end refresh. - Ok(hive_agent_sock::Response::Ok) => (axum::http::StatusCode::OK, "ok").into_response(), - Ok(hive_agent_sock::Response::Err { message }) => error_response( + Ok(hive_core_agent_sock::Response::Ok) => { + (axum::http::StatusCode::OK, "ok").into_response() + } + Ok(hive_core_agent_sock::Response::Err { message }) => error_response( StatusCode::INTERNAL_SERVER_ERROR, &format!("send failed: {message}"), ), diff --git a/hive-agent/src/web_ui/mod.rs b/hive-agent/src/web_ui/mod.rs index 75c0768d..2908a212 100644 --- a/hive-agent/src/web_ui/mod.rs +++ b/hive-agent/src/web_ui/mod.rs @@ -302,11 +302,11 @@ enum BrokerError { /// web-UI handler that talks to the broker goes through it. async fn broker_request( socket: &Path, - req: &hive_agent_sock::Request, -) -> std::result::Result { + req: &hive_core_agent_sock::Request, +) -> std::result::Result { match tokio::time::timeout( SOCKET_FETCH_TIMEOUT, - crate::client::request::<_, hive_agent_sock::Response>(socket, req), + crate::client::request::<_, hive_core_agent_sock::Response>(socket, req), ) .await { diff --git a/hive-agent/src/web_ui/state.rs b/hive-agent/src/web_ui/state.rs index b475c619..9cb001be 100644 --- a/hive-agent/src/web_ui/state.rs +++ b/hive-agent/src/web_ui/state.rs @@ -366,8 +366,13 @@ async fn recent_inbox(socket: &std::path::Path) -> Vec { const LIMIT: u64 = 30; // Deadline-bounded (via `broker_request`): `/api/state` must render even // when hive-c0re is busy — an empty inbox section beats a hung snapshot. - match super::broker_request(socket, &hive_agent_sock::Request::Recent { limit: LIMIT }).await { - Ok(hive_agent_sock::Response::Recent { rows }) => rows, + match super::broker_request( + socket, + &hive_core_agent_sock::Request::Recent { limit: LIMIT }, + ) + .await + { + Ok(hive_core_agent_sock::Response::Recent { rows }) => rows, _ => Vec::new(), } } diff --git a/hive-agent/src/web_ui/stats.rs b/hive-agent/src/web_ui/stats.rs index bbb3eb8a..cc310603 100644 --- a/hive-agent/src/web_ui/stats.rs +++ b/hive-agent/src/web_ui/stats.rs @@ -35,14 +35,14 @@ async fn fetch_reminder_stats( ) -> Option { match super::broker_request( socket, - &hive_agent_sock::Request::ReminderRollup { + &hive_core_agent_sock::Request::ReminderRollup { since_secs: window_secs, agent: None, }, ) .await { - Ok(hive_agent_sock::Response::ReminderRollup(stats)) => Some(stats), + Ok(hive_core_agent_sock::Response::ReminderRollup(stats)) => Some(stats), _ => None, } } @@ -57,14 +57,14 @@ async fn fetch_reminder_stats( pub(super) async fn api_loose_ends(State(state): State) -> Response { match super::broker_request( &state.socket, - &hive_agent_sock::Request::GetLooseEnds { agent: None }, + &hive_core_agent_sock::Request::GetLooseEnds { agent: None }, ) .await { - Ok(hive_agent_sock::Response::LooseEnds { loose_ends }) => { + Ok(hive_core_agent_sock::Response::LooseEnds { loose_ends }) => { axum::Json(serde_json::json!({ "loose_ends": loose_ends })).into_response() } - Ok(hive_agent_sock::Response::Err { message }) => error_response( + Ok(hive_core_agent_sock::Response::Err { message }) => error_response( StatusCode::INTERNAL_SERVER_ERROR, &format!("get_loose_ends: {message}"), ), diff --git a/hive-bash-mcp/Cargo.toml b/hive-bash-mcp/Cargo.toml index 06269257..440e73a1 100644 --- a/hive-bash-mcp/Cargo.toml +++ b/hive-bash-mcp/Cargo.toml @@ -8,7 +8,7 @@ workspace = true [dependencies] anyhow.workspace = true -hive-agent-sock.workspace = true +hive-core-agent-sock.workspace = true hive-sh4re.workspace = true libc.workspace = true rmcp.workspace = true diff --git a/hive-bash-mcp/src/runner.rs b/hive-bash-mcp/src/runner.rs index 1db267f1..bba8dcb6 100644 --- a/hive-bash-mcp/src/runner.rs +++ b/hive-bash-mcp/src/runner.rs @@ -760,7 +760,7 @@ pub(crate) async fn send_wake( ); } } - let req = hive_agent_sock::Request::Wake { + let req = hive_core_agent_sock::Request::Wake { from: format!("bash-task-{id}"), body, }; diff --git a/hive-c0re/Cargo.toml b/hive-c0re/Cargo.toml index 9f304815..8809dc9a 100644 --- a/hive-c0re/Cargo.toml +++ b/hive-c0re/Cargo.toml @@ -30,7 +30,7 @@ opentelemetry-otlp = { version = "0.32", default-features = false, features = [ "reqwest-rustls", ] } indicatif.workspace = true -hive-agent-sock.workspace = true +hive-core-agent-sock.workspace = true hive-sh4re.workspace = true hive-host-sock.workspace = true hive-priv-sock.workspace = true diff --git a/hive-c0re/src/socket_server/config_approvals.rs b/hive-c0re/src/socket_server/config_approvals.rs index e852892c..95f1caaa 100644 --- a/hive-c0re/src/socket_server/config_approvals.rs +++ b/hive-c0re/src/socket_server/config_approvals.rs @@ -9,7 +9,7 @@ use std::sync::Arc; -use hive_agent_sock::Response; +use hive_core_agent_sock::Response; use super::require_new_child; use crate::coordinator::Coordinator; diff --git a/hive-c0re/src/socket_server/lifecycle_handlers.rs b/hive-c0re/src/socket_server/lifecycle_handlers.rs index ae57b3bd..29182a69 100644 --- a/hive-c0re/src/socket_server/lifecycle_handlers.rs +++ b/hive-c0re/src/socket_server/lifecycle_handlers.rs @@ -5,7 +5,7 @@ use std::sync::Arc; -use hive_agent_sock::Response; +use hive_core_agent_sock::Response; use super::require_descendant; use crate::coordinator::Coordinator; diff --git a/hive-c0re/src/socket_server/mod.rs b/hive-c0re/src/socket_server/mod.rs index 86b0a6a7..99be0eb2 100644 --- a/hive-c0re/src/socket_server/mod.rs +++ b/hive-c0re/src/socket_server/mod.rs @@ -13,7 +13,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use anyhow::{Context, Result}; -use hive_agent_sock::{Request, Response}; +use hive_core_agent_sock::{Request, Response}; use hive_sh4re::{MANAGER_AGENT, Message}; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::{UnixListener, UnixStream}; @@ -190,24 +190,26 @@ pub(crate) fn recv_timeout(wait_seconds: Option) -> std::time::Duration { /// The unified `dispatch` calls this first; the remaining arms (which gate /// on topology / capabilities / tool-groups) are handled there. pub(crate) async fn dispatch_shared( - req: &hive_agent_sock::Request, + req: &hive_core_agent_sock::Request, agent: &str, coord: &Arc, -) -> Option { +) -> Option { Some(match req { - hive_agent_sock::Request::Send { + hive_core_agent_sock::Request::Send { to, body, in_reply_to, } => handle_send(coord, agent, to, body, *in_reply_to), - hive_agent_sock::Request::Recv { wait_seconds, max } => { + hive_core_agent_sock::Request::Recv { wait_seconds, max } => { handle_recv(coord, agent, *wait_seconds, *max).await } - hive_agent_sock::Request::Status => handle_status(coord, agent), - hive_agent_sock::Request::OperatorMsg { body } => handle_operator_msg(coord, agent, body), - hive_agent_sock::Request::Wake { from, body } => handle_wake(coord, agent, from, body), - hive_agent_sock::Request::Recent { limit } => handle_recent(coord, agent, *limit), - hive_agent_sock::Request::Ask { + hive_core_agent_sock::Request::Status => handle_status(coord, agent), + hive_core_agent_sock::Request::OperatorMsg { body } => { + handle_operator_msg(coord, agent, body) + } + hive_core_agent_sock::Request::Wake { from, body } => handle_wake(coord, agent, from, body), + hive_core_agent_sock::Request::Recent { limit } => handle_recent(coord, agent, *limit), + hive_core_agent_sock::Request::Ask { question, options, multi, @@ -223,42 +225,42 @@ pub(crate) async fn dispatch_shared( to.as_deref(), ) .map_or_else( - |message| hive_agent_sock::Response::Err { message }, - |id| hive_agent_sock::Response::QuestionQueued { id }, + |message| hive_core_agent_sock::Response::Err { message }, + |id| hive_core_agent_sock::Response::QuestionQueued { id }, ), - hive_agent_sock::Request::Answer { id, answer } => { + hive_core_agent_sock::Request::Answer { id, answer } => { crate::questions::handle_answer(coord, agent, *id, answer).map_or_else( - |message| hive_agent_sock::Response::Err { message }, - |()| hive_agent_sock::Response::Ok, + |message| hive_core_agent_sock::Response::Err { message }, + |()| hive_core_agent_sock::Response::Ok, ) } - hive_agent_sock::Request::Remind { + hive_core_agent_sock::Request::Remind { message, timing, file_path, } => handle_remind(coord, agent, message, timing, file_path.as_deref()), - hive_agent_sock::Request::SetStatus { text } => handle_set_status(coord, text), - hive_agent_sock::Request::GetAgentMeta { name } => { + hive_core_agent_sock::Request::SetStatus { text } => handle_set_status(coord, text), + hive_core_agent_sock::Request::GetAgentMeta { name } => { handle_get_agent_meta(coord, agent, name.as_deref()).await } - hive_agent_sock::Request::CancelLooseEnd { kind, id } => { + hive_core_agent_sock::Request::CancelLooseEnd { kind, id } => { crate::questions::handle_cancel_loose_end(coord, agent, *kind, *id).map_or_else( - |message| hive_agent_sock::Response::Err { message }, - |()| hive_agent_sock::Response::Ok, + |message| hive_core_agent_sock::Response::Err { message }, + |()| hive_core_agent_sock::Response::Ok, ) } - hive_agent_sock::Request::CreateRepo { repo } => handle_create_repo(agent, repo).await, - hive_agent_sock::Request::AckTurn => handle_ack_turn(coord, agent), - hive_agent_sock::Request::AckUntil { up_to } => handle_ack_until(coord, agent, *up_to), - hive_agent_sock::Request::RequeueInflight => handle_requeue_inflight(coord, agent), - hive_agent_sock::Request::GracefulStopComplete => { + hive_core_agent_sock::Request::CreateRepo { repo } => handle_create_repo(agent, repo).await, + hive_core_agent_sock::Request::AckTurn => handle_ack_turn(coord, agent), + hive_core_agent_sock::Request::AckUntil { up_to } => handle_ack_until(coord, agent, *up_to), + hive_core_agent_sock::Request::RequeueInflight => handle_requeue_inflight(coord, agent), + hive_core_agent_sock::Request::GracefulStopComplete => { // Harness drained + is exiting: clear the fence so the // `GracefulStop` orchestration (which polls this flag) proceeds // to stop the container without waiting out its timeout. coord.clear_graceful_stop(agent); - hive_agent_sock::Response::Ok + hive_core_agent_sock::Response::Ok } - hive_agent_sock::Request::GetHostJournal { + hive_core_agent_sock::Request::GetHostJournal { unit, container, lines, @@ -293,7 +295,7 @@ async fn handle_recv( agent: &str, wait_seconds: Option, max: Option, -) -> hive_agent_sock::Response { +) -> hive_core_agent_sock::Response { // Graceful-stop fence: while a graceful stop is pending for this agent, // return `GracefulStop` instead of polling the broker. The harness runs // one stop-checkpoint turn then exits; new sends keep queueing in the @@ -301,7 +303,7 @@ async fn handle_recv( // so a flag set between polls is seen on the next Recv — the orchestration // also fires a transient wake to break an in-flight long-poll. if coord.is_graceful_stop_pending(agent) { - return hive_agent_sock::Response::GracefulStop; + return hive_core_agent_sock::Response::GracefulStop; } let cap = max.unwrap_or(1).min(RECV_BATCH_MAX) as usize; match coord @@ -315,7 +317,7 @@ async fn handle_recv( // exactly how many still-pending messages remain to drain. A count // error is non-fatal — fall back to 0 rather than fail the recv. let remaining = coord.broker.count_pending(agent).unwrap_or(0); - hive_agent_sock::Response::Messages { + hive_core_agent_sock::Response::Messages { messages: deliveries .into_iter() .map(|d| hive_sh4re::DeliveredMessage { @@ -329,7 +331,7 @@ async fn handle_recv( remaining, } } - Err(e) => hive_agent_sock::Response::Err { + Err(e) => hive_core_agent_sock::Response::Err { message: format!("{e:#}"), }, } @@ -343,15 +345,15 @@ fn handle_wake( agent: &str, from: &str, body: &str, -) -> hive_agent_sock::Response { +) -> hive_core_agent_sock::Response { match coord.broker.send(&Message { from: from.to_owned(), to: agent.to_owned(), body: body.to_owned(), in_reply_to: None, }) { - Ok(()) => hive_agent_sock::Response::Ok, - Err(e) => hive_agent_sock::Response::Err { + Ok(()) => hive_core_agent_sock::Response::Ok, + Err(e) => hive_core_agent_sock::Response::Err { message: format!("{e:#}"), }, } @@ -361,13 +363,13 @@ fn handle_wake( /// rescan. The harness has already written the status file to its own /// `state/` dir (it runs as the agent user), so this only refreshes the /// dashboard's view. -fn handle_set_status(coord: &Arc, text: &str) -> hive_agent_sock::Response { +fn handle_set_status(coord: &Arc, text: &str) -> hive_core_agent_sock::Response { if let Err(message) = crate::limits::check_status_text(text) { - return hive_agent_sock::Response::Err { message }; + return hive_core_agent_sock::Response::Err { message }; } let coord2 = Arc::clone(coord); tokio::spawn(async move { coord2.rescan_containers_and_emit().await }); - hive_agent_sock::Response::Ok + hive_core_agent_sock::Response::Ok } /// Validate an agent-supplied repo name: a single safe slug segment, no @@ -385,9 +387,9 @@ fn valid_repo_name(name: &str) -> bool { /// `CreateRepo` — create a repo for `agent` *through hive-c0re* in the /// c0re-owned `agents` org with operator-team branch protection. /// The sanctioned create path now that agents can't create repos directly. -async fn handle_create_repo(agent: &str, repo: &str) -> hive_agent_sock::Response { +async fn handle_create_repo(agent: &str, repo: &str) -> hive_core_agent_sock::Response { if !valid_repo_name(repo) { - return hive_agent_sock::Response::Err { + return hive_core_agent_sock::Response::Err { message: format!( "invalid repo name {repo:?} — single segment of letters, digits, '-', '_', '.' \ (no leading '-'/'.', max 100 chars)" @@ -395,16 +397,16 @@ async fn handle_create_repo(agent: &str, repo: &str) -> hive_agent_sock::Respons }; } let Some(core_token) = crate::forge::core_token() else { - return hive_agent_sock::Response::Err { + return hive_core_agent_sock::Response::Err { message: "forge unavailable (no core token) — cannot create repo".to_owned(), }; }; match crate::forge::create_agent_repo(agent, repo, &core_token).await { - Ok(full_name) => hive_agent_sock::Response::RepoCreated { + Ok(full_name) => hive_core_agent_sock::Response::RepoCreated { clone_url: format!("{}/{full_name}.git", crate::forge::forge_http_base()), full_name, }, - Err(e) => hive_agent_sock::Response::Err { + Err(e) => hive_core_agent_sock::Response::Err { message: format!("create repo {repo:?} failed: {e:#}"), }, } @@ -417,7 +419,7 @@ async fn handle_get_agent_meta( coord: &Arc, agent: &str, name: Option<&str>, -) -> hive_agent_sock::Response { +) -> hive_core_agent_sock::Response { let target = name.unwrap_or(agent); // `name` is agent-supplied and flows into filesystem reads below // (`read_agent_status_live`, `read_agent_matrix_identities` → @@ -428,7 +430,7 @@ async fn handle_get_agent_meta( let target_id = match hive_types::Ident::parse(target) { Ok(id) => id, Err(reason) => { - return hive_agent_sock::Response::Err { + return hive_core_agent_sock::Response::Err { message: format!("get_agent_meta: invalid agent name {target:?}: {reason}"), }; } @@ -436,7 +438,7 @@ async fn handle_get_agent_meta( let (status_text, status_set_at, running) = crate::container_view::read_agent_status_live(&target_id).await; let (hive_name, swarm_name) = crate::container_view::hive_swarm_names(); - hive_agent_sock::Response::AgentMeta { + hive_core_agent_sock::Response::AgentMeta { name: target.to_owned(), running, hyperhive_rev: crate::auto_update::current_flake_rev(&coord.hyperhive_flake), @@ -470,10 +472,10 @@ fn read_agent_matrix_identities(agent: &hive_types::Ident) -> Vec, agent: &str) -> hive_agent_sock::Response { +fn handle_status(coord: &Arc, agent: &str) -> hive_core_agent_sock::Response { match coord.broker.count_pending(agent) { - Ok(unread) => hive_agent_sock::Response::Status { unread }, - Err(e) => hive_agent_sock::Response::Err { + Ok(unread) => hive_core_agent_sock::Response::Status { unread }, + Err(e) => hive_core_agent_sock::Response::Err { message: format!("{e:#}"), }, } @@ -485,15 +487,15 @@ fn handle_operator_msg( coord: &Arc, agent: &str, body: &str, -) -> hive_agent_sock::Response { +) -> hive_core_agent_sock::Response { match coord.broker.send(&Message { from: hive_sh4re::OPERATOR_RECIPIENT.to_owned(), to: agent.to_owned(), body: body.to_owned(), in_reply_to: None, }) { - Ok(()) => hive_agent_sock::Response::Ok, - Err(e) => hive_agent_sock::Response::Err { + Ok(()) => hive_core_agent_sock::Response::Ok, + Err(e) => hive_core_agent_sock::Response::Err { message: format!("{e:#}"), }, } @@ -501,10 +503,14 @@ fn handle_operator_msg( /// `Recent` — the last `limit` inbox rows for `agent` (read-only, /// doesn't consume). -fn handle_recent(coord: &Arc, agent: &str, limit: u64) -> hive_agent_sock::Response { +fn handle_recent( + coord: &Arc, + agent: &str, + limit: u64, +) -> hive_core_agent_sock::Response { match coord.broker.recent_for(agent, limit) { - Ok(rows) => hive_agent_sock::Response::Recent { rows }, - Err(e) => hive_agent_sock::Response::Err { + Ok(rows) => hive_core_agent_sock::Response::Recent { rows }, + Err(e) => hive_core_agent_sock::Response::Err { message: format!("{e:#}"), }, } @@ -512,10 +518,10 @@ fn handle_recent(coord: &Arc, agent: &str, limit: u64) -> hive_agen /// `AckTurn` — mark `agent`'s in-flight delivered messages acked so /// they don't redeliver on the next turn. -fn handle_ack_turn(coord: &Arc, agent: &str) -> hive_agent_sock::Response { +fn handle_ack_turn(coord: &Arc, agent: &str) -> hive_core_agent_sock::Response { match coord.broker.ack_turn(agent) { - Ok(_n) => hive_agent_sock::Response::Ok, - Err(e) => hive_agent_sock::Response::Err { + Ok(_n) => hive_core_agent_sock::Response::Ok, + Err(e) => hive_core_agent_sock::Response::Err { message: format!("{e:#}"), }, } @@ -527,10 +533,10 @@ fn handle_ack_until( coord: &Arc, agent: &str, up_to: i64, -) -> hive_agent_sock::Response { +) -> hive_core_agent_sock::Response { match coord.broker.ack_until(agent, up_to) { - Ok(count) => hive_agent_sock::Response::Acked { count }, - Err(e) => hive_agent_sock::Response::Err { + Ok(count) => hive_core_agent_sock::Response::Acked { count }, + Err(e) => hive_core_agent_sock::Response::Err { message: format!("{e:#}"), }, } @@ -538,15 +544,18 @@ fn handle_ack_until( /// `RequeueInflight` — resurface `agent`'s unacked in-flight messages /// (crash recovery on harness boot). -fn handle_requeue_inflight(coord: &Arc, agent: &str) -> hive_agent_sock::Response { +fn handle_requeue_inflight( + coord: &Arc, + agent: &str, +) -> hive_core_agent_sock::Response { match coord.broker.requeue_inflight(agent) { Ok(n) => { if n > 0 { tracing::info!(%agent, requeued = %n, "requeued in-flight messages"); } - hive_agent_sock::Response::Ok + hive_core_agent_sock::Response::Ok } - Err(e) => hive_agent_sock::Response::Err { + Err(e) => hive_core_agent_sock::Response::Err { message: format!("{e:#}"), }, } diff --git a/hive-c0re/src/socket_server/reminders.rs b/hive-c0re/src/socket_server/reminders.rs index e2b7b734..ae879da1 100644 --- a/hive-c0re/src/socket_server/reminders.rs +++ b/hive-c0re/src/socket_server/reminders.rs @@ -5,7 +5,7 @@ use std::sync::Arc; -use hive_agent_sock::Response; +use hive_core_agent_sock::Response; use crate::coordinator::Coordinator; diff --git a/hive-c0re/src/socket_server/schedules.rs b/hive-c0re/src/socket_server/schedules.rs index 9ba19b1b..bebe855d 100644 --- a/hive-c0re/src/socket_server/schedules.rs +++ b/hive-c0re/src/socket_server/schedules.rs @@ -6,7 +6,7 @@ use std::sync::Arc; -use hive_agent_sock::Response; +use hive_core_agent_sock::Response; use crate::coordinator::Coordinator; diff --git a/hive-agent-sock/Cargo.toml b/hive-core-agent-sock/Cargo.toml similarity index 83% rename from hive-agent-sock/Cargo.toml rename to hive-core-agent-sock/Cargo.toml index 57a61004..f810785d 100644 --- a/hive-agent-sock/Cargo.toml +++ b/hive-core-agent-sock/Cargo.toml @@ -1,5 +1,5 @@ [package] -name = "hive-agent-sock" +name = "hive-core-agent-sock" edition.workspace = true version.workspace = true diff --git a/hive-agent-sock/src/lib.rs b/hive-core-agent-sock/src/lib.rs similarity index 100% rename from hive-agent-sock/src/lib.rs rename to hive-core-agent-sock/src/lib.rs diff --git a/hive-matrix-mcp/src/wake.rs b/hive-matrix-mcp/src/wake.rs index 2ce4e737..18c41d62 100644 --- a/hive-matrix-mcp/src/wake.rs +++ b/hive-matrix-mcp/src/wake.rs @@ -24,7 +24,7 @@ use tokio::net::UnixStream; /// failure; callers log + ignore so a wake delivery hiccup doesn't tear /// down the matrix sync loop. /// -/// Wire format matches `hive_agent_sock::Request` tagged with `"cmd"` per +/// Wire format matches `hive_core_agent_sock::Request` tagged with `"cmd"` per /// `#[serde(tag = "cmd", rename_all = "snake_case")]`. Must be `"cmd"`, /// not `"kind"` — the harness deserialises against the hive-sh4re type /// and silently discards requests that don't match. diff --git a/hive-types/src/lib.rs b/hive-types/src/lib.rs index c6b6a536..40d04857 100644 --- a/hive-types/src/lib.rs +++ b/hive-types/src/lib.rs @@ -1,7 +1,7 @@ //! Foundational shared newtypes for the hyperhive workspace. //! //! A zero-dependency (bar `serde`) leaf crate so every wire-type crate -//! (`hive-sh4re`, `hive-host-sock`, `hive-agent-sock`) and both binaries +//! (`hive-sh4re`, `hive-host-sock`, `hive-core-agent-sock`) and both binaries //! (`hive-c0re`, `hivectl`) can type their agent-name fields as [`Ident`] //! and get serde-validated parsing at the socket boundary for free — with //! no cross-crate coupling and without growing `hive-sh4re`.