Watch
0
0
Fork
You've already forked hyperhive
0

hive-agent: drive turns through hive-runtime

AgentSession becomes hive_runtime::AgentRuntime, picked at startup from
the environment. Unset HIVE_RUNTIME keeps the claude backend, built
from the same title, store and PercentPolicy as before; drive_turn,
the 401 retry and the error mapping see the claude errors unchanged.

On acp the harness keeps the session id in the harness dir, answers
permission requests like the claude built-in allow-list (no built-in
shell, web fetch only with web_tools), maps runtime errors to Failed,
and reports an operator /compact as skipped instead of done.

Refs #4391
This commit is contained in:
atlas 2026-09-29 21:22:41 +02:00 • committed by mara
commit f0110be76c
4 changed files with 119 additions and 58 deletions

View file

@ -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, single `hive-ag3nt/` dir — that's the runtime/binary-family nickname,
not a directory). not a directory).
- **`hive-agent/`** — the serve-loop binary: turn-loop *policy* layer - **`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 module dir), event + turn-stats sqlite sinks, login flow, system-prompt
renderer. renderer.
- **`hive-agent-mcp/`** — the embedded MCP server (long-lived - **`hive-agent-mcp/`** — the embedded MCP server (long-lived
streamable-http listener, `hive-mcp-http` systemd unit) + its claude streamable-http listener, `hive-mcp-http` systemd unit) + its claude
launch-config layer (tool-group/capability → `--allowedTools`, launch-config layer (tool-group/capability → `--allowedTools`,
`--mcp-config` render). `--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 - **`hive-jobq/`** — job-DAG scheduler, extracted from hive-c0re's
in-tree `job_queue` as a domain-agnostic library. **Runtime-only — in-tree `job_queue` as a domain-agnostic library. **Runtime-only —
nothing writes the graph to disk**; hive-c0re starts empty each boot nothing writes the graph to disk**; hive-c0re starts empty each boot

View file

@ -22,6 +22,7 @@ http-body-util.workspace = true
futures-util = "0.3" futures-util = "0.3"
clap.workspace = true clap.workspace = true
hive-claude.workspace = true hive-claude.workspace = true
hive-runtime.workspace = true
hive-agent-sock.workspace = true hive-agent-sock.workspace = true
hive-core-agent-sock.workspace = true hive-core-agent-sock.workspace = true
hive-log.workspace = true hive-log.workspace = true

View file

@ -659,9 +659,9 @@ async fn serve_loop<S: Surface>(
) -> Result<()> { ) -> Result<()> {
tracing::info!(socket = %socket.display(), "harness serve"); tracing::info!(socket = %socket.display(), "harness serve");
S::requeue_inflight(socket).await; 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). // 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 // Tracks the last observed pause state so the transitions get logged
// once each instead of twelve lines a minute while parked. // once each instead of twelve lines a minute while parked.
let mut was_paused = false; let mut was_paused = false;

View file

@ -1,6 +1,6 @@
//! Per-turn claude policy layer. The generic subprocess mechanics — spawning //! Per-turn agent policy layer. The runtime mechanics — spawning the agent
//! `claude --print`, streaming + classifying stream-json, session //! (`claude --print`, or an ACP agent), streaming its output, session
//! lookup/archive — live in the `hive-claude` crate. This module owns the //! lookup/archive — live in the `hive-runtime` crate. This module owns the
//! hyperhive-specific policy on top: building the per-turn config from 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 //! bus, bridging the output stream onto the event bus (`BusSink`), and the
//! compaction / auto-reset / retry state machine (`drive_turn`). //! compaction / auto-reset / retry state machine (`drive_turn`).
@ -8,7 +8,10 @@
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use anyhow::Result; 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 serde_json::Value;
use crate::events::{Bus, LiveEvent}; 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 // claude-code; if a key gets renamed we'll spot it because the corresponding
// behavior will start firing mid-turn again. // behavior will start firing mid-turn again.
// //
// The subprocess mechanics — spawning `claude --print`, streaming + // The runtime mechanics — spawning the agent, streaming its output, session
// classifying stream-json, session lookup/archive — live in the generic // lookup/archive — live in the generic `hive-runtime` crate. This module is
// `hive-claude` crate. This module is the hyperhive *policy* layer on top: // the hyperhive *policy* layer on top:
// it builds the per-turn [`Config`] from the bus, forwards the stream to the // it builds the per-turn [`Config`] from the bus, forwards the stream to the
// event bus via [`BusSink`], and owns compaction / auto-reset / retry. // event bus via [`BusSink`], and owns compaction / auto-reset / retry.
@ -165,7 +168,7 @@ pub type TurnOutcome = std::result::Result<bool, TurnError>;
#[derive(Debug)] #[derive(Debug)]
pub enum TurnError { pub enum TurnError {
/// claude saw "Prompt is too long" and even a reactive compact + retry /// 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 /// window. Rare. [`drive_turn`] archives the session (so the next turn
/// starts fresh) and the serve loop requeues the in-flight message, which /// starts fresh) and the serve loop requeues the in-flight message, which
/// redelivers into that fresh session — the wake prompt itself is tiny, so /// redelivers into that fresh session — the wake prompt itself is tiny, so
@ -182,7 +185,7 @@ pub enum TurnError {
AuthFailed, AuthFailed,
/// `--resume <title>` missed AND the lib's create self-heal also failed to /// `--resume <title>` missed AND the lib's create self-heal also failed to
/// resolve the session — "shouldn't happen" (a resume-miss is normally /// 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; /// the wake message, the serve loop requeues it so the next turn retries;
/// no status park. /// no status park.
SessionNotFound, SessionNotFound,
@ -292,21 +295,29 @@ fn compact_percent() -> u8 {
u8::try_from(pct.min(100)).expect("value clamped to <= 100 fits in 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`] /// Where an ACP agent's session id is kept, under the harness dir.
/// with hyperhive's percent-of-window compaction policy. Built once by the const ACP_SESSION_FILE: &str = "hyperhive-acp-session";
/// 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>;
/// Construct the agent's durable session: constant title + on-disk store + a /// The agent's durable session, on the runtime its environment selects
/// percent-of-window compaction policy that checkpoints (`CHECKPOINT_PROMPT`) /// (`hive_runtime::RuntimeSpec`). On claude it is the constant-title
/// before compacting. Called once at serve-loop start. `percent` comes from a /// [`hive_claude::InfiniteSession`] with hyperhive's percent-of-window
/// boot-time env var and `default_window` is only a fallback for turns where /// compaction policy. Built once by the serve loop (see [`make_session`]) and
/// the model didn't report a window, so a single build at startup is fine. /// threaded through the turns, rather than rebuilt each time.
#[must_use] pub type AgentSession = AgentRuntime<PercentPolicy>;
pub fn make_session(bus: &Bus) -> AgentSession {
InfiniteSession::new( /// 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_title(),
session_store(), session_store(),
PercentPolicy { PercentPolicy {
@ -314,10 +325,30 @@ pub fn make_session(bus: &Bus) -> AgentSession {
default_window: Some(effective_context_window(bus)), default_window: Some(effective_context_window(bus)),
checkpoint_prompt: Some(CHECKPOINT_PROMPT.to_string()), 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 /// resume-or-create + compaction loop (reactive on overflow, and proactive per
/// the percent policy — including the pre-compaction checkpoint turn). This /// the percent policy — including the pre-compaction checkpoint turn). This
/// layer wraps it with the two hyperhive-specific concerns: /// layer wraps it with the two hyperhive-specific concerns:
@ -344,15 +375,18 @@ pub async fn drive_turn(
bus.emit(LiveEvent::Note { bus.emit(LiveEvent::Note {
text: "operator: resetting session — archiving before this turn".into(), text: "operator: resetting session — archiving before this turn".into(),
}); });
archive_session(bus); archive_session(bus, session);
} else { } else {
// Heuristic: context large AND prompt cache gone cold. // Heuristic: context large AND prompt cache gone cold.
maybe_auto_reset(bus); maybe_auto_reset(bus, session);
} }
let config = claude_config(bus, files); let config = claude_config(bus, files);
let sink = BusSink::new(bus); let sink = BusSink::new(bus);
let mut result = session.run(&config, prompt, &sink).await; 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 { bus.emit(LiveEvent::Note {
text: "got 401 — retrying once before parking for re-login".into(), text: "got 401 — retrying once before parking for re-login".into(),
}); });
@ -373,7 +407,7 @@ pub async fn drive_turn(
} }
Ok(progress.compacted) Ok(progress.compacted)
} }
Err(e) => error_to_turn(e), Err(e) => runtime_error_to_turn(e),
}; };
if matches!(outcome, Err(TurnError::PromptTooLong)) { if matches!(outcome, Err(TurnError::PromptTooLong)) {
// The lib already compacted + retried and the session is still over the // 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" retried message starts fresh"
.into(), .into(),
}); });
archive_session(bus); archive_session(bus, session);
return Err(TurnError::PromptTooLong); return Err(TurnError::PromptTooLong);
} }
// Operator `/compact` (`POST /api/compact`) or an agent's own `compact` // 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`) // Reflect `Compacting` in the UI like the idle path (`run_pending_compact`)
// does; the serve loop resets to `Idle` once this turn returns. // does; the serve loop resets to `Idle` once this turn returns.
bus.set_state(crate::events::TurnState::Compacting); 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` // 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 // tool with a `wake_prompt`), stash it — the serve loop reads it
// back after this turn returns and drives a synthetic follow-up // 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 { if let Some(prompt) = request.wake_prompt {
bus.set_post_compact_wake(prompt); bus.set_post_compact_wake(prompt);
} }
return Ok(true); return Ok(compacted);
} }
outcome outcome
} }
@ -426,7 +468,7 @@ pub async fn drive_turn(
/// `--name <title>` session. No preceding checkpoint turn — running any 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 /// before the reset would re-upload and re-warm the cache, which defeats the
/// cost-optimisation purpose entirely. /// 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); let watermark = auto_reset_watermark_tokens(bus);
if watermark == 0 { if watermark == 0 {
return; // auto-reset disabled return; // auto-reset disabled
@ -457,7 +499,7 @@ fn maybe_auto_reset(bus: &Bus) {
— dropping session (cache cold, fresh start is equally cheap)" — 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 /// 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, /// 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 /// 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 /// 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 /// `true` if a compaction ran; the serve loop follows up with
/// `Bus::take_post_compact_wake` to see whether a synthetic follow-up turn /// `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`). /// 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: /// 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 /// per-turn tool-call counting (`observe_stream`), the live SSE stream, and
/// non-JSON stdout + stderr as Notes. Stateless — usage/model/context-window /// 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 /// Archive (do NOT delete) the harness's own session so the next turn starts
/// `--resume <title>` misses and self-heals into a fresh `--name <title>` /// a fresh one. On claude the rename is
/// session. Delegates the rename to [`hive_claude::SessionStore::archive_by_title`] /// [`hive_claude::SessionStore::archive_by_title`] (which touches only the file
/// (which touches only the file carrying OUR `customTitle`, leaving any `choom` /// carrying OUR `customTitle`, leaving any `choom` session sharing the cwd
/// session sharing the cwd alone) and surfaces the result as a Note. Best- /// alone), so the next `--resume <title>` misses and self-heals into a fresh
/// effort: never fails a turn. Only ever called at a turn boundary (top of /// `--name <title>` session. Surfaces the result as a Note. Best-effort: never
/// `drive_turn` for an operator reset, or `maybe_auto_reset` pre-turn) so no /// fails a turn. Only ever called at a turn boundary (top of `drive_turn` for
/// claude process holds the session file open when it's renamed. /// an operator reset, or `maybe_auto_reset` pre-turn) so no turn holds the
fn archive_session(bus: &Bus) { /// session open when it's renamed.
fn archive_session(bus: &Bus, session: &AgentSession) {
let title = session_title(); let title = session_title();
match session_store().archive_by_title(&title) { match session.archive() {
Ok(Some(path)) => { Ok(Some(path)) => {
let name = path let name = path
.file_name() .file_name()
.and_then(|n| n.to_str()) .and_then(|n| n.to_str())
.unwrap_or("?") .unwrap_or("?")
.to_string(); .to_string();
tracing::info!(path = %path.display(), "archived claude session"); tracing::info!(path = %path.display(), "archived agent session");
bus.emit(LiveEvent::Note { bus.emit(LiveEvent::Note {
text: format!("archived session \"{title}\" ({name}) — next turn starts fresh"), text: format!("archived session \"{title}\" ({name}) — next turn starts fresh"),
}); });
@ -733,7 +785,7 @@ fn archive_session(bus: &Bus) {
), ),
}), }),
Err(e) => { Err(e) => {
tracing::warn!(error = %e, "failed to archive claude session"); tracing::warn!(error = %e, "failed to archive agent session");
bus.emit(LiveEvent::Note { bus.emit(LiveEvent::Note {
text: format!("failed to archive session \"{title}\": {e}"), text: format!("failed to archive session \"{title}\": {e}"),
}); });