remove Role::Manager + ManagerSurface + Flavor::Manager — there is only one role: agent
This commit is contained in:
parent
186ee430b5
commit
f56b272a23
8 changed files with 100 additions and 489 deletions
|
|
@ -1,10 +1,7 @@
|
|||
//! Unified hyperhive harness binary. Picks role from `HIVE_ROLE`
|
||||
//! (`"agent"` | `"manager"`), dispatches one of three subcommands
|
||||
//! (`serve` / `mcp` / `wake`), and runs the turn loop through a
|
||||
//! generic `Surface` trait so both wire surfaces stay in lockstep.
|
||||
//!
|
||||
//! Architecture (single-binary rationale, Surface-trait + zero-sized
|
||||
//! type tags, boot wiring, turn-outcome branch) lives in
|
||||
//! Unified hyperhive harness binary. Dispatches one of three subcommands
|
||||
//! (`serve` / `mcp` / `wake`). There is one role: agent. The `Surface`
|
||||
//! trait + `AgentSurface` zero-sized type tag keeps the turn loop
|
||||
//! generic and testable. Architecture lives in
|
||||
//! [`docs/turn-loop.md::Harness binary shape`](../../../docs/turn-loop.md).
|
||||
|
||||
use std::path::{Path, PathBuf};
|
||||
|
|
@ -13,7 +10,7 @@ use std::time::Duration;
|
|||
|
||||
use hive_ag3nt::web_ui::TurnLock;
|
||||
|
||||
use anyhow::{Result, bail};
|
||||
use anyhow::Result;
|
||||
use clap::{Parser, Subcommand};
|
||||
use hive_ag3nt::events::{Bus, LiveEvent, TurnState};
|
||||
use hive_ag3nt::login::{self, LoginState};
|
||||
|
|
@ -21,15 +18,10 @@ use hive_ag3nt::turn_stats::TurnStats;
|
|||
use hive_ag3nt::{
|
||||
DEFAULT_SOCKET, DEFAULT_WEB_PORT, client, mcp, plugins, serve_common, turn, web_ui,
|
||||
};
|
||||
use hive_sh4re::{
|
||||
AgentRequest, AgentResponse, HelperEvent, ManagerRequest, ManagerResponse, SYSTEM_SENDER,
|
||||
};
|
||||
use hive_sh4re::{AgentRequest, AgentResponse, HelperEvent, SYSTEM_SENDER};
|
||||
|
||||
#[derive(Parser)]
|
||||
#[command(
|
||||
name = "hive",
|
||||
about = "hyperhive harness — role from $HIVE_ROLE (agent|manager)"
|
||||
)]
|
||||
#[command(name = "hive", about = "hyperhive harness")]
|
||||
struct Cli {
|
||||
/// Path to the per-agent MCP socket (bind-mounted from the host).
|
||||
#[arg(long, global = true, default_value = DEFAULT_SOCKET)]
|
||||
|
|
@ -48,16 +40,14 @@ enum Cmd {
|
|||
#[arg(long, default_value_t = 1000)]
|
||||
poll_ms: u64,
|
||||
},
|
||||
/// Run this role's MCP server on stdio. Spawned by `claude` via
|
||||
/// Run the MCP server on stdio. Spawned by `claude` via
|
||||
/// `--mcp-config`; tools dispatch through `/run/hive/mcp.sock` back
|
||||
/// into the hyperhive broker.
|
||||
Mcp,
|
||||
/// Inject a wake-up event into this harness's inbox so the next
|
||||
/// turn fires with the given body. Intended for extra MCP servers
|
||||
/// / helpers (matrix bridge, scraper, webhook listener, etc.) that
|
||||
/// need to nudge claude on external events. Available on both
|
||||
/// agent and manager roles; mirrors the `AgentRequest::Wake` /
|
||||
/// `ManagerRequest::Wake` pair already on the wire.
|
||||
/// need to nudge claude on external events.
|
||||
Wake {
|
||||
#[arg(long)]
|
||||
from: String,
|
||||
|
|
@ -67,20 +57,6 @@ enum Cmd {
|
|||
},
|
||||
}
|
||||
|
||||
#[derive(Copy, Clone)]
|
||||
enum Role {
|
||||
Agent,
|
||||
Manager,
|
||||
}
|
||||
|
||||
fn resolve_role() -> Result<Role> {
|
||||
match std::env::var("HIVE_ROLE").as_deref() {
|
||||
Ok("agent") | Err(_) => Ok(Role::Agent),
|
||||
Ok("manager") => Ok(Role::Manager),
|
||||
Ok(other) => bail!("unknown HIVE_ROLE={other:?}; expected 'agent' or 'manager'"),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<()> {
|
||||
tracing_subscriber::fmt()
|
||||
|
|
@ -91,26 +67,11 @@ async fn main() -> Result<()> {
|
|||
.init();
|
||||
|
||||
let cli = Cli::parse();
|
||||
let role = resolve_role()?;
|
||||
|
||||
// Generic dispatch: one `serve_main` / `wake` body, two
|
||||
// monomorphisations driven by the `Surface` type parameter. See
|
||||
// `docs/turn-loop.md::Surface trait + zero-sized type tags`.
|
||||
match (role, cli.cmd) {
|
||||
(Role::Agent, Cmd::Serve { poll_ms }) => {
|
||||
serve_main::<AgentSurface>(&cli.socket, poll_ms).await
|
||||
}
|
||||
(Role::Manager, Cmd::Serve { poll_ms }) => {
|
||||
serve_main::<ManagerSurface>(&cli.socket, poll_ms).await
|
||||
}
|
||||
(Role::Agent, Cmd::Mcp) => mcp::serve_agent_stdio(cli.socket).await,
|
||||
(Role::Manager, Cmd::Mcp) => mcp::serve_agent_stdio(cli.socket).await,
|
||||
(Role::Agent, Cmd::Wake { from, body }) => {
|
||||
wake::<AgentSurface>(&cli.socket, from, body).await
|
||||
}
|
||||
(Role::Manager, Cmd::Wake { from, body }) => {
|
||||
wake::<ManagerSurface>(&cli.socket, from, body).await
|
||||
}
|
||||
match cli.cmd {
|
||||
Cmd::Serve { poll_ms } => serve_main::<AgentSurface>(&cli.socket, poll_ms).await,
|
||||
Cmd::Mcp => mcp::serve_agent_stdio(cli.socket).await,
|
||||
Cmd::Wake { from, body } => wake::<AgentSurface>(&cli.socket, from, body).await,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -186,16 +147,11 @@ enum RecvOutcome {
|
|||
TransportError,
|
||||
}
|
||||
|
||||
/// Per-role wire surface. Two impls — `AgentSurface`, `ManagerSurface`
|
||||
/// — wrap the disjoint `Request`/`Response` enums plus a handful of
|
||||
/// boot-time constants that vary by role. Every other function in this
|
||||
/// binary that talks to the broker goes through this trait so the turn
|
||||
/// loop itself has zero per-role branches.
|
||||
/// Wire surface abstraction. `AgentSurface` is the only impl — the trait
|
||||
/// exists to keep the turn loop generic and testable. Every function that
|
||||
/// talks to the broker goes through this so there are zero hard-coded
|
||||
/// `AgentRequest` / `AgentResponse` references in the turn loop itself.
|
||||
trait Surface {
|
||||
/// MCP flavor passed to `TurnFiles::prepare`. Picks which static
|
||||
/// system-prompt block + tool registration goes into the spawned
|
||||
/// `claude` process.
|
||||
const FLAVOR: mcp::Flavor;
|
||||
/// Ack the in-flight turn. Logs warnings on transport/broker
|
||||
/// errors but never propagates — turn loop continues either way.
|
||||
fn ack_turn(socket: &Path) -> impl Future<Output = ()>;
|
||||
|
|
@ -237,12 +193,11 @@ trait Surface {
|
|||
|
||||
// ---------- AgentSurface ----------
|
||||
|
||||
/// Zero-sized type tag for the sub-agent wire surface.
|
||||
/// Zero-sized type tag for the agent wire surface.
|
||||
/// Talks `AgentRequest` / `AgentResponse`.
|
||||
struct AgentSurface;
|
||||
|
||||
impl Surface for AgentSurface {
|
||||
const FLAVOR: mcp::Flavor = mcp::Flavor::Agent;
|
||||
|
||||
async fn ack_turn(socket: &Path) {
|
||||
match client::request::<_, AgentResponse>(socket, &AgentRequest::AckTurn).await {
|
||||
|
|
@ -382,159 +337,11 @@ impl Surface for AgentSurface {
|
|||
}
|
||||
}
|
||||
|
||||
// ---------- ManagerSurface ----------
|
||||
|
||||
/// Zero-sized type tag for the manager wire surface.
|
||||
/// Talks `ManagerRequest` / `ManagerResponse`.
|
||||
struct ManagerSurface;
|
||||
|
||||
impl Surface for ManagerSurface {
|
||||
const FLAVOR: mcp::Flavor = mcp::Flavor::Manager;
|
||||
|
||||
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"),
|
||||
}
|
||||
}
|
||||
|
||||
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"),
|
||||
}
|
||||
}
|
||||
|
||||
async fn inbox_unread(socket: &Path) -> u64 {
|
||||
match client::request::<_, ManagerResponse>(socket, &ManagerRequest::Status).await {
|
||||
Ok(ManagerResponse::Status { unread }) => unread,
|
||||
_ => 0,
|
||||
}
|
||||
}
|
||||
|
||||
async fn post_turn_counts(socket: &Path) -> (Option<u64>, Option<u64>) {
|
||||
let threads = match client::request::<_, ManagerResponse>(
|
||||
socket,
|
||||
&ManagerRequest::GetLooseEnds { agent: None },
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(ManagerResponse::LooseEnds { loose_ends }) => u64::try_from(loose_ends.len()).ok(),
|
||||
_ => None,
|
||||
};
|
||||
let reminders = match client::request::<_, ManagerResponse>(
|
||||
socket,
|
||||
&ManagerRequest::CountPendingReminders { agent: None },
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(ManagerResponse::PendingRemindersCount { count }) => Some(count),
|
||||
_ => None,
|
||||
};
|
||||
(threads, reminders)
|
||||
}
|
||||
|
||||
async fn send_to_parent(socket: &Path, body: String) {
|
||||
let res = client::request::<_, ManagerResponse>(
|
||||
socket,
|
||||
&ManagerRequest::Send {
|
||||
to: hive_sh4re::PARENT_RECIPIENT.into(),
|
||||
body,
|
||||
in_reply_to: None,
|
||||
},
|
||||
)
|
||||
.await;
|
||||
if let Err(e) = res {
|
||||
tracing::warn!(error = ?e, "failed to notify parent of turn failure");
|
||||
}
|
||||
}
|
||||
|
||||
async fn self_wake(socket: &Path) {
|
||||
let res = client::request::<_, ManagerResponse>(
|
||||
socket,
|
||||
&ManagerRequest::Wake {
|
||||
from: "self".into(),
|
||||
body: "continue".into(),
|
||||
transient: false,
|
||||
},
|
||||
)
|
||||
.await;
|
||||
match res {
|
||||
Ok(ManagerResponse::Ok) => {
|
||||
tracing::info!("request_next_turn: injected self-continue wake");
|
||||
}
|
||||
Ok(ManagerResponse::Err { message }) => {
|
||||
tracing::warn!(%message, "check_and_inject_continue: wake rejected");
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(error = ?e, "check_and_inject_continue: wake transport error");
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
async fn recv_next(socket: &Path) -> RecvOutcome {
|
||||
let recv: Result<ManagerResponse> = client::request(
|
||||
socket,
|
||||
&ManagerRequest::Recv {
|
||||
wait_seconds: Some(180),
|
||||
max: None,
|
||||
},
|
||||
)
|
||||
.await;
|
||||
match recv {
|
||||
Ok(ManagerResponse::Messages { messages }) if !messages.is_empty() => {
|
||||
let first = messages.into_iter().next().expect("checked non-empty");
|
||||
RecvOutcome::Message(first)
|
||||
}
|
||||
Ok(ManagerResponse::Messages { .. }) => RecvOutcome::Empty,
|
||||
Ok(ManagerResponse::Err { message }) => {
|
||||
tracing::warn!(%message, "recv error");
|
||||
RecvOutcome::TransportError
|
||||
}
|
||||
Ok(other) => {
|
||||
tracing::warn!(?other, "recv produced unexpected response kind");
|
||||
RecvOutcome::TransportError
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(error = ?e, "recv failed; retrying");
|
||||
RecvOutcome::TransportError
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn wake_external(socket: &Path, from: String, body: String) -> Result<()> {
|
||||
let resp: ManagerResponse = client::request(
|
||||
socket,
|
||||
&ManagerRequest::Wake {
|
||||
from,
|
||||
body,
|
||||
transient: false,
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
match resp {
|
||||
ManagerResponse::Ok => Ok(()),
|
||||
ManagerResponse::Err { message } => anyhow::bail!("wake: {message}"),
|
||||
other => anyhow::bail!("wake: unexpected response {other:?}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ---------- generic turn loop ----------
|
||||
|
||||
/// Per-role boot — wires up the web UI, login state, stats, plugins,
|
||||
/// forge notifier, and either drops into `serve_loop` directly
|
||||
/// (`Online`) or parks on the login flow first (`NeedsLogin`). See
|
||||
/// Boot — wires up the web UI, login state, stats, plugins, forge
|
||||
/// notifier, and either drops into `serve_loop` directly (`Online`) or
|
||||
/// parks on the login flow first (`NeedsLogin`). See
|
||||
/// `docs/turn-loop.md::Boot wiring`.
|
||||
async fn serve_main<S: Surface>(socket: &Path, poll_ms: u64) -> Result<()> {
|
||||
let port = std::env::var("HIVE_PORT")
|
||||
|
|
@ -559,13 +366,12 @@ async fn serve_main<S: Surface>(socket: &Path, poll_ms: u64) -> Result<()> {
|
|||
bus.seed_usage(ctx, cost);
|
||||
}
|
||||
}
|
||||
let files = turn::TurnFiles::prepare(socket, &label, S::FLAVOR).await?;
|
||||
let files = turn::TurnFiles::prepare(socket, &label).await?;
|
||||
let turn_lock: TurnLock = Arc::new(tokio::sync::Mutex::new(()));
|
||||
// Plugin install runs role-agnostic: failures come back as a
|
||||
// Vec<String> and we route each through `<parent>` via the same
|
||||
// `send_to_parent` failure-notify path the turn loop uses. The
|
||||
// broker resolves `<parent>` per `topology::parent_of`; root
|
||||
// agents and the manager fall through to operator.
|
||||
// Plugin install failures come back as a Vec<String> — route each
|
||||
// through `<parent>` via the `send_to_parent` failure-notify path.
|
||||
// The broker resolves `<parent>` per `topology::parent_of`;
|
||||
// root agents fall through to operator.
|
||||
for failure in plugins::install_configured(socket).await {
|
||||
S::send_to_parent(socket, failure).await;
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue