diff --git a/CLAUDE.md b/CLAUDE.md index 104348f7..6683f878 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -45,13 +45,21 @@ hand-maintained per-file tree drifts out of sync with the code. single `hive-ag3nt/` dir — that's the runtime/binary-family nickname, not a directory). - **`hive-agent/`** — the serve-loop binary: turn-loop *policy* layer - (`turn.rs`) over the `hive-claude` driver, per-agent web UI (`web_ui/` + (`turn.rs`) over `hive-runtime`, per-agent web UI (`web_ui/` module dir), event + turn-stats sqlite sinks, login flow, system-prompt renderer. - **`hive-agent-mcp/`** — the embedded MCP server (long-lived streamable-http listener, `hive-mcp-http` systemd unit) + its claude launch-config layer (tool-group/capability → `--allowedTools`, `--mcp-config` render). +- **`hive-runtime/`** — the runtime an agent's turns run on: one + `Runtime` interface (`run`/`compact`/`archive`) with a `claude` backend + (a pass-through to `hive-claude`) and an `acp` backend (any Agent Client + Protocol agent, spawned from the command/args/env nix hands it — no agent + is named in Rust). The ACP backend translates its stream into claude's + `stream-json` shape, so the harness's consumers read either unchanged. + Depends on no hyperhive binary crate, so the subagent daemon can use it + too. - **`hive-jobq/`** — job-DAG scheduler, extracted from hive-c0re's in-tree `job_queue` as a domain-agnostic library. **Runtime-only — nothing writes the graph to disk**; hive-c0re starts empty each boot diff --git a/hive-agent/Cargo.toml b/hive-agent/Cargo.toml index a34b24df..98afdae2 100644 --- a/hive-agent/Cargo.toml +++ b/hive-agent/Cargo.toml @@ -22,6 +22,7 @@ http-body-util.workspace = true futures-util = "0.3" clap.workspace = true hive-claude.workspace = true +hive-runtime.workspace = true hive-agent-sock.workspace = true hive-core-agent-sock.workspace = true hive-log.workspace = true diff --git a/hive-agent/src/main.rs b/hive-agent/src/main.rs index e1614007..3ed7743f 100644 --- a/hive-agent/src/main.rs +++ b/hive-agent/src/main.rs @@ -659,9 +659,9 @@ async fn serve_loop( ) -> Result<()> { tracing::info!(socket = %socket.display(), "harness serve"); S::requeue_inflight(socket).await; - // The durable claude session, built once and reused for every turn + + // The durable agent session, built once and reused for every turn + // idle compaction below (it's effectively stateless). - let session = turn::make_session(&bus); + let session = turn::make_session(&bus)?; // Tracks the last observed pause state so the transitions get logged // once each instead of twelve lines a minute while parked. let mut was_paused = false; diff --git a/hive-agent/src/turn.rs b/hive-agent/src/turn.rs index 9254fc06..535a9693 100644 --- a/hive-agent/src/turn.rs +++ b/hive-agent/src/turn.rs @@ -1,6 +1,6 @@ -//! Per-turn claude policy layer. The generic subprocess mechanics — spawning -//! `claude --print`, streaming + classifying stream-json, session -//! lookup/archive — live in the `hive-claude` crate. This module owns the +//! Per-turn agent policy layer. The runtime mechanics — spawning the agent +//! (`claude --print`, or an ACP agent), streaming its output, session +//! lookup/archive — live in the `hive-runtime` crate. This module owns the //! hyperhive-specific policy on top: building the per-turn config from the //! bus, bridging the output stream onto the event bus (`BusSink`), and the //! compaction / auto-reset / retry state machine (`drive_turn`). @@ -8,7 +8,10 @@ use std::path::{Path, PathBuf}; use anyhow::Result; -use hive_claude::{Config, InfiniteSession, PercentPolicy, Sink}; +use hive_runtime::{ + AcpRuntime, AgentRuntime, ClaudeRuntime, Config, PercentPolicy, PermissionPolicy, Runtime, + RuntimeSpec, Sink, +}; use serde_json::Value; use crate::events::{Bus, LiveEvent}; @@ -24,9 +27,9 @@ use crate::mcp_config; // claude-code; if a key gets renamed we'll spot it because the corresponding // behavior will start firing mid-turn again. // -// The subprocess mechanics — spawning `claude --print`, streaming + -// classifying stream-json, session lookup/archive — live in the generic -// `hive-claude` crate. This module is the hyperhive *policy* layer on top: +// The runtime mechanics — spawning the agent, streaming its output, session +// lookup/archive — live in the generic `hive-runtime` crate. This module is +// the hyperhive *policy* layer on top: // it builds the per-turn [`Config`] from the bus, forwards the stream to the // event bus via [`BusSink`], and owns compaction / auto-reset / retry. @@ -165,7 +168,7 @@ pub type TurnOutcome = std::result::Result; #[derive(Debug)] pub enum TurnError { /// claude saw "Prompt is too long" and even a reactive compact + retry - /// (inside [`InfiniteSession::run`]) couldn't bring it back under the + /// (inside [`hive_claude::InfiniteSession::run`]) couldn't bring it back under the /// window. Rare. [`drive_turn`] archives the session (so the next turn /// starts fresh) and the serve loop requeues the in-flight message, which /// redelivers into that fresh session — the wake prompt itself is tiny, so @@ -182,7 +185,7 @@ pub enum TurnError { AuthFailed, /// `--resume ` missed AND the lib's create self-heal also failed to /// resolve the session — "shouldn't happen" (a resume-miss is normally - /// self-healed inside [`InfiniteSession::attempt`]). Rather than ack + drop + /// self-healed inside [`hive_claude::InfiniteSession::attempt`]). Rather than ack + drop /// the wake message, the serve loop requeues it so the next turn retries; /// no status park. SessionNotFound, @@ -292,32 +295,60 @@ fn compact_percent() -> u8 { u8::try_from(pct.min(100)).expect("value clamped to <= 100 fits in u8") } -/// The agent's durable session type: the constant-title [`InfiniteSession`] -/// with hyperhive's percent-of-window compaction policy. Built once by the -/// serve loop (see [`make_session`]) and threaded through the turns, rather -/// than rebuilt each time — it's effectively stateless, so one instance serves -/// the whole run. -pub type AgentSession = InfiniteSession<PercentPolicy>; +/// Where an ACP agent's session id is kept, under the harness dir. +const ACP_SESSION_FILE: &str = "hyperhive-acp-session"; -/// Construct the agent's durable session: constant title + on-disk store + a -/// percent-of-window compaction policy that checkpoints (`CHECKPOINT_PROMPT`) -/// before compacting. Called once at serve-loop start. `percent` comes from a -/// boot-time env var and `default_window` is only a fallback for turns where -/// the model didn't report a window, so a single build at startup is fine. -#[must_use] -pub fn make_session(bus: &Bus) -> AgentSession { - InfiniteSession::new( - session_title(), - session_store(), - PercentPolicy { - percent: compact_percent(), - default_window: Some(effective_context_window(bus)), - checkpoint_prompt: Some(CHECKPOINT_PROMPT.to_string()), - }, - ) +/// The agent's durable session, on the runtime its environment selects +/// (`hive_runtime::RuntimeSpec`). On claude it is the constant-title +/// [`hive_claude::InfiniteSession`] with hyperhive's percent-of-window +/// compaction policy. Built once by the serve loop (see [`make_session`]) and +/// threaded through the turns, rather than rebuilt each time. +pub type AgentSession = AgentRuntime<PercentPolicy>; + +/// Construct the agent's durable session. On claude: constant title + on-disk +/// store + a percent-of-window compaction policy that checkpoints +/// (`CHECKPOINT_PROMPT`) before compacting. Called once at serve-loop start. +/// `percent` comes from a boot-time env var and `default_window` is only a +/// fallback for turns where the model didn't report a window, so a single +/// build at startup is fine. +/// +/// # Errors +/// +/// Returns an error if the runtime selection in the environment is invalid. +pub fn make_session(bus: &Bus) -> Result<AgentSession> { + Ok(match RuntimeSpec::from_env()? { + RuntimeSpec::Claude => AgentRuntime::Claude(ClaudeRuntime::new( + session_title(), + session_store(), + PercentPolicy { + percent: compact_percent(), + default_window: Some(effective_context_window(bus)), + checkpoint_prompt: Some(CHECKPOINT_PROMPT.to_string()), + }, + )), + RuntimeSpec::Acp(command) => AgentRuntime::Acp(Box::new(AcpRuntime::new( + command, + crate::paths::harness_dir().join(ACP_SESSION_FILE), + acp_permission_policy(), + ))), + }) } -/// Drive one turn end-to-end. The durable [`InfiniteSession`] owns the +/// Which ACP tool-call kinds an ACP agent may run when it asks. Mirrors the +/// claude built-ins (`hive_sh4re::permissions`): never a built-in shell +/// (`execute`) — shell goes through `mcp__bash__run` — and web access +/// (`fetch`) only with the `web_tools` group. +fn acp_permission_policy() -> PermissionPolicy { + let web = + mcp_config::effective_tool_groups().contains(&hive_sh4re::permissions::ToolGroup::WebTools); + std::sync::Arc::new(move |kind: &str| match kind { + "execute" => false, + "fetch" => web, + _ => true, + }) +} + +/// Drive one turn end-to-end. The durable [`AgentSession`] owns the /// resume-or-create + compaction loop (reactive on overflow, and proactive per /// the percent policy — including the pre-compaction checkpoint turn). This /// layer wraps it with the two hyperhive-specific concerns: @@ -344,15 +375,18 @@ pub async fn drive_turn( bus.emit(LiveEvent::Note { text: "operator: resetting session — archiving before this turn".into(), }); - archive_session(bus); + archive_session(bus, session); } else { // Heuristic: context large AND prompt cache gone cold. - maybe_auto_reset(bus); + maybe_auto_reset(bus, session); } let config = claude_config(bus, files); let sink = BusSink::new(bus); let mut result = session.run(&config, prompt, &sink).await; - if matches!(result, Err(hive_claude::Error::AuthFailed)) { + if matches!( + result, + Err(hive_runtime::Error::Claude(hive_claude::Error::AuthFailed)) + ) { bus.emit(LiveEvent::Note { text: "got 401 — retrying once before parking for re-login".into(), }); @@ -373,7 +407,7 @@ pub async fn drive_turn( } Ok(progress.compacted) } - Err(e) => error_to_turn(e), + Err(e) => runtime_error_to_turn(e), }; if matches!(outcome, Err(TurnError::PromptTooLong)) { // The lib already compacted + retried and the session is still over the @@ -385,7 +419,7 @@ pub async fn drive_turn( retried message starts fresh" .into(), }); - archive_session(bus); + archive_session(bus, session); return Err(TurnError::PromptTooLong); } // Operator `/compact` (`POST /api/compact`) or an agent's own `compact` @@ -406,7 +440,15 @@ pub async fn drive_turn( // Reflect `Compacting` in the UI like the idle path (`run_pending_compact`) // does; the serve loop resets to `Idle` once this turn returns. bus.set_state(crate::events::TurnState::Compacting); - let _ = session.compact(&config, &sink).await; + let compacted = match session.compact(&config, &sink).await { + Err(e @ hive_runtime::Error::Unsupported(_)) => { + bus.emit(LiveEvent::Note { + text: format!("/compact skipped: {e}"), + }); + false + } + _ => true, + }; // If the compact call asked to be woken (the agent's own `compact` // tool with a `wake_prompt`), stash it — the serve loop reads it // back after this turn returns and drives a synthetic follow-up @@ -415,7 +457,7 @@ pub async fn drive_turn( if let Some(prompt) = request.wake_prompt { bus.set_post_compact_wake(prompt); } - return Ok(true); + return Ok(compacted); } outcome } @@ -426,7 +468,7 @@ pub async fn drive_turn( /// `--name <title>` session. No preceding checkpoint turn — running any turn /// before the reset would re-upload and re-warm the cache, which defeats the /// cost-optimisation purpose entirely. -fn maybe_auto_reset(bus: &Bus) { +fn maybe_auto_reset(bus: &Bus, session: &AgentSession) { let watermark = auto_reset_watermark_tokens(bus); if watermark == 0 { return; // auto-reset disabled @@ -457,7 +499,7 @@ fn maybe_auto_reset(bus: &Bus) { — dropping session (cache cold, fresh start is equally cheap)" ), }); - archive_session(bus); + archive_session(bus, session); } /// Emit the per-turn `TurnEnd` event + log line. Single owner so outcome @@ -524,7 +566,7 @@ pub fn emit_turn_end(bus: &Bus, outcome: &TurnOutcome) { /// agent is idle — the serve loop calls this when a `recv` returns no message, /// so a queued `/compact` runs even when no turn is driving. (The in-flight /// case is handled at the end of [`drive_turn`].) Resume-only via -/// [`InfiniteSession::compact`]: a missing session is a harmless no-op. Returns +/// [`hive_claude::InfiniteSession::compact`]: a missing session is a harmless no-op. Returns /// `true` if a compaction ran; the serve loop follows up with /// `Bus::take_post_compact_wake` to see whether a synthetic follow-up turn /// should run (set below when the compact request carried a `wake_prompt`). @@ -637,6 +679,15 @@ fn error_to_turn(err: hive_claude::Error) -> TurnOutcome { } } +/// [`error_to_turn`] for a runtime error: the claude backend's errors map as +/// above, every other runtime's as `Failed`. +fn runtime_error_to_turn(err: hive_runtime::Error) -> TurnOutcome { + match err { + hive_runtime::Error::Claude(err) => error_to_turn(err), + other => Err(TurnError::Failed(other.into())), + } +} + /// Bridges a claude run's raw output stream onto the hyperhive event bus: /// per-turn tool-call counting (`observe_stream`), the live SSE stream, and /// non-JSON stdout + stderr as Notes. Stateless — usage/model/context-window @@ -705,24 +756,25 @@ fn apply_telemetry(bus: &Bus, telemetry: &hive_claude::Telemetry) { } } -/// Archive (do NOT delete) the harness's own session so the next turn's -/// `--resume <title>` misses and self-heals into a fresh `--name <title>` -/// session. Delegates the rename to [`hive_claude::SessionStore::archive_by_title`] -/// (which touches only the file carrying OUR `customTitle`, leaving any `choom` -/// session sharing the cwd alone) and surfaces the result as a Note. Best- -/// effort: never fails a turn. Only ever called at a turn boundary (top of -/// `drive_turn` for an operator reset, or `maybe_auto_reset` pre-turn) so no -/// claude process holds the session file open when it's renamed. -fn archive_session(bus: &Bus) { +/// Archive (do NOT delete) the harness's own session so the next turn starts +/// a fresh one. On claude the rename is +/// [`hive_claude::SessionStore::archive_by_title`] (which touches only the file +/// carrying OUR `customTitle`, leaving any `choom` session sharing the cwd +/// alone), so the next `--resume <title>` misses and self-heals into a fresh +/// `--name <title>` session. Surfaces the result as a Note. Best-effort: never +/// fails a turn. Only ever called at a turn boundary (top of `drive_turn` for +/// an operator reset, or `maybe_auto_reset` pre-turn) so no turn holds the +/// session open when it's renamed. +fn archive_session(bus: &Bus, session: &AgentSession) { let title = session_title(); - match session_store().archive_by_title(&title) { + match session.archive() { Ok(Some(path)) => { let name = path .file_name() .and_then(|n| n.to_str()) .unwrap_or("?") .to_string(); - tracing::info!(path = %path.display(), "archived claude session"); + tracing::info!(path = %path.display(), "archived agent session"); bus.emit(LiveEvent::Note { text: format!("archived session \"{title}\" ({name}) — next turn starts fresh"), }); @@ -733,7 +785,7 @@ fn archive_session(bus: &Bus) { ), }), Err(e) => { - tracing::warn!(error = %e, "failed to archive claude session"); + tracing::warn!(error = %e, "failed to archive agent session"); bus.emit(LiveEvent::Note { text: format!("failed to archive session \"{title}\": {e}"), });