//! 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`). use std::path::{Path, PathBuf}; use anyhow::Result; use hive_runtime::{ AcpRuntime, AgentRuntime, ClaudeRuntime, Config, PercentPolicy, PermissionAsk, PermissionPolicy, Runtime, RuntimeSpec, Sink, }; use serde_json::Value; use crate::events::{Bus, LiveEvent}; use crate::mcp_config; // Hive-enforced claude settings ship at `/etc/claude-code/managed-settings.json`, // which claude-code auto-discovers (precedence #1, read-only, un-overridable) — // so the harness no longer passes `--settings`. We turn off claude's in-session // auto-compaction and its cross-session auto-memory because hyperhive owns those // concerns (`/compact` on overflow, notes persistence under `/state`). How the // file is wired (the nix asset) + the full rationale live in // `docs/turn-loop/claude-invocation.md`. Unknown keys are silently ignored by // claude-code; if a key gets renamed we'll spot it because the corresponding // behavior will start firing mid-turn again. // // 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. /// Fixed, harness-owned claude session title. Every turn / compact / /// checkpoint resumes THIS title (`--resume `); the create path /// names it (`--name <title>`). One constant identity per agent means /// compaction and the post-compact retry provably target the same session — /// there is no scraped UUID to go stale, empty, or diverge. Each agent runs /// in its own container (own `~/.claude` + own `/state` cwd), so even the /// shared default never collides across agents. Override via /// `HIVE_SESSION_TITLE`. const DEFAULT_SESSION_TITLE: &str = "hive-session"; /// How long to sleep after detecting a rate-limit before re-entering the /// serve loop. Overridable via `HIVE_RATE_LIMIT_SLEEP_SECS`. Default is /// 5 minutes — enough for most short-lived throttles; the operator can /// tune down for tight retry scenarios or up if they're hitting sustained /// capacity limits. const DEFAULT_RATE_LIMIT_SLEEP_SECS: u64 = 300; /// Idle watchdog window: kill claude and surface `ApiStall` if it produces no /// stdout for this long. The timer resets on every stdout line, so a large but /// still-streaming turn is never cut — only a complete silence trips it. /// Default is 10 minutes: long enough for legitimate big-context turns, short /// enough to cut a multi-retry Anthropic connection storm (atlas telemetry: /// attempt:11 took ~6.7min on a 371k-token context). Override via /// `HIVE_TURN_IDLE_SECS`; `0` disables the watchdog (wait indefinitely). const DEFAULT_TURN_IDLE_SECS: u64 = 600; /// How long to park after an `ApiStall` before requeueing the message, giving /// the API a chance to recover. Overridable via `HIVE_STALL_SLEEP_SECS`. const DEFAULT_STALL_SLEEP_SECS: u64 = 60; /// Assumed prompt-cache TTL. Claude caches prompt prefixes — ~5 minutes on /// the API (pay-per-token), ~1 hour on Claude Max (subscription). When the /// idle gap exceeds this, the cache prefix has likely expired and the next /// turn re-uploads the full transcript regardless of whether we resume or /// start fresh. A fresh session with a small context is therefore equally /// cheap but gives the model a clean slate. Default is 3600s (1h) matching /// the subscription TTL; API (pay-per-token) users should set /// `HIVE_CACHE_TTL_SECS=300`. Override via `HIVE_CACHE_TTL_SECS`; set to /// `0` to disable (always resume). const DEFAULT_CACHE_TTL_SECS: u64 = 3600; /// Default proactive-compaction watermark, as a percent of the effective /// context window. Overridable via `HIVE_COMPACT_WATERMARK_PERCENT`. const DEFAULT_COMPACT_PERCENT: u8 = 75; /// Synthetic wake prompt for the proactive notes-checkpoint turn. Not an /// inbox message — the harness injects it directly so the agent gets one /// turn to persist durable state before `/compact` collapses the /// turn-by-turn history into a summary. const CHECKPOINT_PROMPT: &str = "[system] Context checkpoint — no inbox message to handle.\n\n\ Your conversation context has grown large and the harness is about to run `/compact`, \ which collapses the detailed turn-by-turn history into a short summary. Anything you \ do not persist now is effectively lost after the next turn.\n\n\ Use THIS turn to flush anything worth keeping into your durable `/state` files: update \ your notes / CLAUDE.md / TODO.md with in-flight task state, decisions made, important \ file paths, and whatever you would need to resume cleanly with only a summary of this \ conversation to go on. Do not start new work or reply to anyone — just write your notes \ and end the turn."; /// The set of files claude reads on every invocation: the MCP server /// config (`--mcp-config`) and the pre-rendered role/tools system /// prompt (`--system-prompt-file`). Static settings are no longer /// passed here — they live at `/etc/claude-code/managed-settings.json` /// and claude auto-discovers them. /// Materialised once at harness startup; shared between the turn loop /// and the operator-driven `/compact` path so both invocations look /// identical to claude (same MCP surface, same allowed tools, same /// role prompt — only the stdin payload differs). #[derive(Clone)] pub struct TurnFiles { pub mcp_config: PathBuf, pub system_prompt: PathBuf, } impl TurnFiles { /// Write the two per-turn files (MCP config + system prompt) into the /// agent's config dir. Idempotent — overwrites whatever was there. /// /// # Errors /// /// Returns an error if any of the config files cannot be written to disk. pub async fn prepare(socket: &Path, label: &str) -> Result<Self> { Ok(Self { mcp_config: write_mcp_config().await?, system_prompt: write_system_prompt(socket, label).await?, }) } } /// Drop the MCP config blob claude reads from `--mcp-config <path>`. /// The built-in hyperhive surface is served over HTTP by the persistent /// `hive-mcp-http` daemon, so no per-turn stdio child is spawned; extra /// servers declared via `services.hyperhive.agent.extraMcpServers` are still stdio bridges. /// /// # Errors /// /// Returns an error if the config file cannot be written. pub async fn write_mcp_config() -> Result<PathBuf> { let parent = crate::paths::config_dir(); tokio::fs::create_dir_all(&parent).await.ok(); let path = parent.join("claude-mcp-config.json"); let body = mcp_config::render_claude_config(); tokio::fs::write(&path, body).await?; tracing::info!(path = %path.display(), "wrote claude MCP config"); Ok(path) } /// Thin re-export of [`crate::prompt::write_system_prompt`] for /// callers that already import this module. The actual rendering + /// marker-block logic lives in `prompt.rs`; this is just the public /// entry point the binaries call. /// /// # Errors /// /// Returns an error if the system prompt file cannot be written. pub async fn write_system_prompt(socket: &Path, label: &str) -> Result<PathBuf> { crate::prompt::write_system_prompt(socket, label).await } /// One claude turn's outcome: `Ok(compacted)` on success, or a [`TurnError`] /// the serve loop must act on. The `compacted` bool is `true` when a /// compaction ran this turn (reactively on overflow, or proactively per the /// policy — or an operator `/compact` at turn end); it's recorded as /// `result_kind = "compacted"` in turn stats so the stats page can distinguish /// those turns. Both `Ok(true)` and `Ok(false)` are ack'd; the error cases /// each map to a distinct serve-loop action (see [`emit_turn_end`] and the /// `hive-agent` serve loop). pub type TurnOutcome = std::result::Result<bool, TurnError>; /// The ways a turn can end without a usable result. Each is deliberately *not* /// a generic failure — the serve loop reacts to each differently (requeue, /// park, escalate). #[derive(Debug)] pub enum TurnError { /// claude saw "Prompt is too long" and even a reactive compact + retry /// (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 /// the overflow was the accumulated context, which the archive clears. PromptTooLong, /// The Anthropic API refused the request due to a rate limit, per-account /// usage cap, or exhausted credit balance. The serve loop should park for /// `rate_limit_sleep_secs()` and requeue — NOT bubble up as a crash. RateLimited, /// The Anthropic API rejected the request with 401 (OAuth session /// expired or revoked). The serve loop should flip the container /// into `needs_login_idle` and stop driving turns until the /// operator re-auths via the per-agent web UI. AuthFailed, /// `--resume <title>` 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 [`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, /// The harness-side idle watchdog killed claude after `turn_idle_secs()` of /// output silence — indicative of an Anthropic API stall (e.g. a multi-retry /// connection storm burning minutes of wall-clock with no stream progress). /// The serve loop parks for `stall_sleep_secs()` and requeues, like the /// rate-limit path — NOT a crash. ApiStall, /// A hard failure with no recovery — the serve loop escalates it to the /// operator (`send_to_operator`). Failed(anyhow::Error), } /// Parse an env var as `u64`, ignoring absent / blank / unparseable values. /// Returns the raw value including `0` (several knobs use `0` as "disable"). fn env_u64(name: &str) -> Option<u64> { parse_u64(std::env::var(name).ok().as_deref()) } /// The parsing half of [`env_u64`], split from the lookup so it can be /// exercised without setting a process-global variable. fn parse_u64(raw: Option<&str>) -> Option<u64> { raw.and_then(|s| s.trim().parse::<u64>().ok()) } /// Like [`env_u64`] but also rejects `0`, falling back to `default` — for /// knobs where `0` is meaningless rather than a "disable" sentinel. fn env_u64_positive(name: &str, default: u64) -> u64 { nonzero_or(env_u64(name), default) } /// The `0`-rejecting half of [`env_u64_positive`]. /// /// Which of the two a knob uses is the whole distinction: a knob reading /// [`env_u64`] treats `0` as "disable" and keeps it, one reading /// [`env_u64_positive`] treats `0` as nonsense and falls back. Picking the /// wrong one silently changes when the loop sleeps, compacts or resets. fn nonzero_or(parsed: Option<u64>, default: u64) -> u64 { parsed.filter(|&v| v > 0).unwrap_or(default) } /// How long to sleep after a rate-limit before re-entering the serve loop. /// Reads `HIVE_RATE_LIMIT_SLEEP_SECS` if set to a valid positive integer. #[must_use] pub fn rate_limit_sleep_secs() -> u64 { env_u64_positive("HIVE_RATE_LIMIT_SLEEP_SECS", DEFAULT_RATE_LIMIT_SLEEP_SECS) } /// Idle-watchdog window in seconds. Reads `HIVE_TURN_IDLE_SECS`; `0` disables /// the watchdog. Absent / unparseable falls back to [`DEFAULT_TURN_IDLE_SECS`]. #[must_use] pub fn turn_idle_secs() -> u64 { env_u64("HIVE_TURN_IDLE_SECS").unwrap_or(DEFAULT_TURN_IDLE_SECS) } /// How long to park after an `ApiStall` before requeueing. Reads /// `HIVE_STALL_SLEEP_SECS` if set to a valid positive integer. #[must_use] pub fn stall_sleep_secs() -> u64 { env_u64_positive("HIVE_STALL_SLEEP_SECS", DEFAULT_STALL_SLEEP_SECS) } /// Resolve the effective context-window size for watermark calculations. /// Priority order (first wins): /// 1. API-reported window from the last `result` event's `modelUsage.*.contextWindow`. /// 2. `HIVE_CONTEXT_WINDOW_TOKENS_*` env vars (Nix-configured per-model defaults). /// 3. Hard fallback: 200 000. /// /// The API-reported window is the authoritative per-inference active /// context limit. It reflects what the model actually enforces — which /// for models with large prompt caches (e.g. 1 M total cache) may be /// significantly smaller than the cache capacity (e.g. 200 k active window /// for `claude-sonnet-4-6`). fn effective_context_window(bus: &Bus) -> u64 { bus.api_context_window() .unwrap_or_else(|| crate::harness_state::context_window_tokens(&bus.model())) } /// Resolve the auto-reset watermark. Priority order: /// 1. `HIVE_AUTO_RESET_WATERMARK_TOKENS` env var (explicit override). /// 2. 50% of `effective_context_window(bus)`. /// /// `0` disables auto-reset entirely. fn auto_reset_watermark_tokens(bus: &Bus) -> u64 { env_u64("HIVE_AUTO_RESET_WATERMARK_TOKENS").unwrap_or_else(|| effective_context_window(bus) / 2) } /// Resolve the assumed cache TTL: `HIVE_CACHE_TTL_SECS` if set, else /// `DEFAULT_CACHE_TTL_SECS`. fn cache_ttl_secs() -> u64 { env_u64_positive("HIVE_CACHE_TTL_SECS", DEFAULT_CACHE_TTL_SECS) } /// Proactive-compaction watermark as a percent of the effective context /// window (default [`DEFAULT_COMPACT_PERCENT`]). `0` disables proactive /// compaction — the reactive on-overflow path still applies. Reads /// `HIVE_COMPACT_WATERMARK_PERCENT`; the legacy `autoCompact = false` switch /// (which sets `HIVE_COMPACT_WATERMARK_TOKENS=0`) is still honoured as disable. fn compact_percent() -> u8 { if env_u64("HIVE_COMPACT_WATERMARK_TOKENS") == Some(0) { return 0; } let pct = env_u64("HIVE_COMPACT_WATERMARK_PERCENT").unwrap_or(u64::from(DEFAULT_COMPACT_PERCENT)); // `min(100)` is ≤ 100, so this `try_from` is infallible. u8::try_from(pct.min(100)).expect("value clamped to <= 100 fits in u8") } /// Where an ACP agent's session id is kept, under the harness dir. const ACP_SESSION_FILE: &str = "hyperhive-acp-session"; /// 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(), ))), }) } /// The answer to an ACP agent's permission requests, see [`acp_permits`]. fn acp_permission_policy() -> PermissionPolicy { let web = mcp_config::effective_tool_groups().contains(&hive_sh4re::permissions::ToolGroup::WebTools); std::sync::Arc::new(move |ask: &PermissionAsk<'_>| acp_permits(ask, web)) } /// Whether an ACP agent may run the tool call it asks about. Default-deny, /// allowing what a claude agent has: /// /// - `other`-kind calls to tools of the MCP servers the session was handed /// — which of those an agent gets is already decided by its tool groups /// (`mcp_config`); /// - the file tools, mirroring the claude built-ins /// (`hive_sh4re::permissions`): `read`, `edit`, `search`; /// - `fetch` with the `web_tools` group (`web`). /// /// Everything else is refused, including a built-in shell (`execute`) — /// shell goes through `mcp__bash__run` — and any kind this list doesn't name. fn acp_permits(ask: &PermissionAsk<'_>, web: bool) -> bool { match ask.kind { "other" => ask.mcp_server.is_some(), "read" | "edit" | "search" => true, "fetch" => web, _ => false, } } /// 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: /// /// - **Session reset (pre-turn)** — an operator reset (`/api/new-session`) or /// the auto-reset heuristic (context large AND prompt cache gone cold) /// archives the current session at this turn boundary so the run starts /// fresh. The two are mutually exclusive. This is deliberately *not* part of /// the infinite-session abstraction — it's the hive escape hatch. /// - **401 retry** — a transient token-refresh race can 401 once and clear, so /// the whole turn is retried a single time before bubbling `AuthFailed` to /// the serve loop (which parks for re-login). /// /// Called once per turn by the `hive-agent` serve loop, which owns the shared /// `session` ([`make_session`]) and threads it in. pub async fn drive_turn( prompt: &str, files: &TurnFiles, bus: &Bus, session: &AgentSession, ) -> TurnOutcome { if bus.take_session_reset() { // Operator-requested (deferred from `POST /api/new-session`). bus.emit(LiveEvent::Note { text: "operator: resetting session — archiving before this turn".into(), }); archive_session(bus, session); } else { // Heuristic: context large AND prompt cache gone cold. 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_runtime::Error::Claude(hive_claude::Error::AuthFailed)) ) { bus.emit(LiveEvent::Note { text: "got 401 — retrying once before parking for re-login".into(), }); result = session.run(&config, prompt, &sink).await; } let outcome = match result { Ok(progress) => { // Apply the turn's parsed usage / model / context-window to the bus // (badges, stats, auto-reset watermark input). apply_telemetry(bus, &progress.telemetry); if progress.created { // Fresh session minted this turn → flag it so the bin loop // mints a `sessions` row + stamps its id onto this turn's stats. bus.mark_fresh_session(); bus.emit(LiveEvent::Note { text: format!("created fresh session titled \"{}\"", session_title()), }); } Ok(progress.compacted) } 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 // window. Archive it here (session lifecycle stays hive-side) so the // requeued message — handled by the serve loop — redelivers into a // fresh session that fits. bus.emit(LiveEvent::Note { text: "context still over the window after compaction — archiving session so the \ retried message starts fresh" .into(), }); archive_session(bus, session); return Err(TurnError::PromptTooLong); } // Operator `/compact` (`POST /api/compact`) or an agent's own `compact` // MCP tool call, deferred to the turn boundary: run it now that the turn // is done, so it works mid-turn rather than only when the agent is idle. // Only on a healthy turn — no point spawning a compaction after a // rate-limited / auth-failed / crashed one. `is_ok()` first: `take_compact()` // clears the flag, so it must only fire when the compaction will actually // run. On an unhealthy turn (rate-limited / auth-failed / failed) the flag // is left set for the next turn or the idle `run_pending_compact` to // service — not silently eaten. if outcome.is_ok() && let Some(request) = bus.take_compact() { bus.emit(LiveEvent::Note { text: "operator: /compact — running at turn end".into(), }); // 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 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 // turn, so a self-requested compact provably doesn't strand the // agent idle waiting for the next external event. if let Some(prompt) = request.wake_prompt { bus.set_post_compact_wake(prompt); } return Ok(compacted); } outcome } /// Pre-turn auto-reset check. If context is large AND the prompt cache has /// gone cold (idle time >= cache TTL), archive the current session so the /// next wake-up turn's `--resume <title>` misses and self-heals into a fresh /// `--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, session: &AgentSession) { let watermark = auto_reset_watermark_tokens(bus); if watermark == 0 { return; // auto-reset disabled } let Some(ctx_tokens) = bus.last_ctx_usage().map(|u| u.context_tokens()) else { return; // no usage reading yet — first turn, nothing to reset }; if ctx_tokens < watermark { return; } let last_ended = bus.last_turn_ended_unix(); if last_ended == 0 { return; // no completed turn yet } // Compute idle seconds using the same clock semantics as elsewhere // in this crate (unix epoch, i64 seconds via `Utc::now().timestamp()`). let now = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map_or(0, |d| d.as_secs()); let idle_secs = now.saturating_sub(u64::try_from(last_ended).unwrap_or(0)); let ttl = cache_ttl_secs(); if idle_secs < ttl { return; } bus.emit(LiveEvent::Note { text: format!( "context {ctx_tokens} tokens, idle {idle_secs}s >= cache TTL {ttl}s \ — dropping session (cache cold, fresh start is equally cheap)" ), }); archive_session(bus, session); } /// Emit the per-turn `TurnEnd` event + log line. Single owner so outcome /// semantics stay consistent across every agent role. pub fn emit_turn_end(bus: &Bus, outcome: &TurnOutcome) { match outcome { Ok(_) => { bus.emit(LiveEvent::TurnEnd { ok: true, note: None, }); tracing::info!("turn finished"); } Err(TurnError::PromptTooLong) => { bus.emit(LiveEvent::TurnEnd { ok: false, note: Some("context too long after compaction — session archived, retrying".into()), }); tracing::warn!("turn prompt-too-long; archived session and requeueing"); } Err(TurnError::RateLimited) => { bus.emit(LiveEvent::TurnEnd { ok: false, note: Some("rate limited — parking until quota resets".into()), }); tracing::warn!("turn rate-limited"); } Err(TurnError::AuthFailed) => { bus.emit(LiveEvent::TurnEnd { ok: false, note: Some("authentication failed (401) — waiting for re-login".into()), }); tracing::warn!("turn auth-failed (401)"); } Err(TurnError::SessionNotFound) => { bus.emit(LiveEvent::TurnEnd { ok: false, note: Some("session resume + create both missed — requeueing".into()), }); tracing::warn!("turn session-not-found; requeueing message"); } Err(TurnError::ApiStall) => { bus.emit(LiveEvent::TurnEnd { ok: false, note: Some(format!( "claude killed after {}s of output silence — API stall suspected, parking + requeueing", turn_idle_secs() )), }); tracing::warn!("turn killed: API stall (idle watchdog)"); } Err(TurnError::Failed(e)) => { let note = format!("{e:#}"); bus.emit(LiveEvent::TurnEnd { ok: false, note: Some(note.clone()), }); tracing::warn!(error = %note, "turn failed"); } } } /// Service a pending operator `/compact` (`Bus::request_compact`) while the /// 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 /// [`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`). pub async fn run_pending_compact(files: &TurnFiles, bus: &Bus, session: &AgentSession) -> bool { let Some(request) = bus.take_compact() else { return false; }; bus.emit(LiveEvent::Note { text: "operator: /compact — running on idle session".into(), }); bus.set_state(crate::events::TurnState::Compacting); let config = claude_config(bus, files); let sink = BusSink::new(bus); match session.compact(&config, &sink).await { Ok(()) => bus.emit(LiveEvent::Note { text: "/compact done".into(), }), Err(e) => bus.emit(LiveEvent::Note { text: format!("/compact failed: {e}"), }), } bus.set_state(crate::events::TurnState::Idle); if let Some(prompt) = request.wake_prompt { bus.set_post_compact_wake(prompt); } true } /// The constant session title for this agent. `HIVE_SESSION_TITLE` overrides /// the compiled-in [`DEFAULT_SESSION_TITLE`]; each agent runs in its own /// container (own `~/.claude` + own `/state` cwd), so even the shared default /// never collides across agents. #[must_use] pub fn session_title() -> String { std::env::var("HIVE_SESSION_TITLE") .ok() .map(|s| s.trim().to_string()) .filter(|s| !s.is_empty()) .unwrap_or_else(|| DEFAULT_SESSION_TITLE.to_string()) } /// The cwd claude is spawned in: the agent's durable `/state` dir when it /// exists, else the harness process cwd. Claude derives its per-project /// session dir from this path, so the same value feeds both the [`Config`] and /// the [`hive_claude::SessionStore`]. fn session_cwd() -> PathBuf { let state_dir = crate::paths::state_dir(); if state_dir.is_dir() { state_dir } else { std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")) } } /// The on-disk session store for this agent (claude home + spawn cwd), used to /// locate + archive the harness session by title. fn session_store() -> hive_claude::SessionStore { hive_claude::SessionStore::new(crate::paths::claude_dir(), session_cwd()) } /// Build the per-turn `hive_claude::Config` from the bus (model / effort) and /// the materialised `TurnFiles` (system prompt + MCP config), plus the fixed /// tool allow-lists and the optional docs `--add-dir`. fn claude_config(bus: &Bus, files: &TurnFiles) -> Config { let mut add_dirs = Vec::new(); // services.hyperhive.agent.docs.enable wires HIVE_DOCS_DIR to the in-container reference // docs; expose it as an additional readable directory when set. if let Some(docs_dir) = std::env::var_os("HIVE_DOCS_DIR") && !docs_dir.is_empty() { add_dirs.push(PathBuf::from(docs_dir)); } let cwd = { let state_dir = crate::paths::state_dir(); state_dir.is_dir().then_some(state_dir) }; Config { model: Some(bus.model()), effort: Some(bus.effort()), cwd, system_prompt_file: Some(files.system_prompt.clone()), mcp_config: Some(files.mcp_config.clone()), strict_mcp_config: true, tools: Some(mcp_config::builtin_tools_arg()), allowed_tools: Some(mcp_config::allowed_tools_arg()), add_dirs, // Idle watchdog: `0` disables (wait indefinitely), any positive value // caps output silence. Policy lives here; the driver just enforces it. idle_timeout: match turn_idle_secs() { 0 => None, secs => Some(std::time::Duration::from_secs(secs)), }, ..Config::default() } } /// Map a `hive_claude::Error` onto the harness's `TurnOutcome`. The recognized /// sentinels become their matching outcomes; a residual `SessionNotFound` /// (create path itself missed — shouldn't happen) settles as `Ok`; genuine /// failures become `Failed` (converting the typed lib error into `anyhow`). fn error_to_turn(err: hive_claude::Error) -> TurnOutcome { use hive_claude::Error; match err { Error::PromptTooLong => Err(TurnError::PromptTooLong), Error::RateLimited => Err(TurnError::RateLimited), Error::AuthFailed => Err(TurnError::AuthFailed), Error::SessionNotFound => Err(TurnError::SessionNotFound), Error::IdleTimeout => Err(TurnError::ApiStall), other => Err(TurnError::Failed(other.into())), } } /// [`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 /// parsing lives in `hive-claude` and is applied from the run's returned /// `Telemetry` (see `apply_telemetry`). struct BusSink<'a> { bus: &'a Bus, } impl<'a> BusSink<'a> { fn new(bus: &'a Bus) -> Self { Self { bus } } } impl Sink for BusSink<'_> { fn on_event(&self, event: &Value) { // Raw-event concerns only: per-turn tool-call counting + the live SSE // stream. Usage / model / context-window parsing lives in the lib now // and is applied from the run's returned `Telemetry` (see `drive_turn` // → `apply_telemetry`). self.bus.observe_stream(event); self.bus.observe_mcp_health(event); self.bus.emit(LiveEvent::Stream(event.clone())); } fn on_stdout_line(&self, line: &str) { self.bus.emit(LiveEvent::Note { text: format!("(non-json) {line}"), }); } fn on_stderr_line(&self, line: &str) { // Mirror to journald so post-mortems work without the web UI / events // sqlite; the bus Note is what the dashboard renders. tracing::warn!(line = %line, "claude stderr"); self.bus.emit(LiveEvent::Note { text: format!("stderr: {line}"), }); } } /// Apply a completed turn's parsed [`hive_claude::Telemetry`] to the bus: /// per-inference context usage + cumulative cost, the resolved model id, and /// the API-reported context window (the authoritative window for the auto-reset /// watermark). Skips a degenerate turn that parsed nothing so it doesn't reset /// the badges to zero. fn apply_telemetry(bus: &Bus, telemetry: &hive_claude::Telemetry) { // On a degenerate turn that emitted a `result` but no `assistant` event, // the per-inference `context` stays zero while `cost` (cumulative) is not. // Fall back to `cost` as the ctx proxy so the ctx badge + auto-reset // watermark don't go stale-to-zero. Only a turn that parsed nothing at all // (both zero) is skipped. let ctx = if telemetry.context.context_tokens() == 0 { telemetry.cost } else { telemetry.context }; if ctx.context_tokens() == 0 { return; } bus.record_turn_usage(ctx, telemetry.cost); bus.set_resolved_model(telemetry.model.clone()); if let Some(window) = telemetry.context_window { bus.set_api_context_window(window); } } /// 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.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 agent session"); bus.emit(LiveEvent::Note { text: format!("archived session \"{title}\" ({name}) — next turn starts fresh"), }); } Ok(None) => bus.emit(LiveEvent::Note { text: format!( "no existing session titled \"{title}\" to archive — next turn starts fresh" ), }), Err(e) => { tracing::warn!(error = %e, "failed to archive agent session"); bus.emit(LiveEvent::Note { text: format!("failed to archive session \"{title}\": {e}"), }); } } } #[cfg(test)] mod tests { use super::{PermissionAsk, acp_permits, nonzero_or, parse_u64}; #[test] fn an_absent_or_blank_value_is_not_a_number() { assert_eq!(parse_u64(None), None); assert_eq!(parse_u64(Some("")), None); assert_eq!(parse_u64(Some(" ")), None); } #[test] fn surrounding_whitespace_is_trimmed() { assert_eq!(parse_u64(Some("12")), Some(12)); assert_eq!(parse_u64(Some(" 12\n")), Some(12)); } #[test] fn a_value_that_is_not_a_u64_is_rejected_rather_than_coerced() { for raw in ["abc", "1.5", "-1", "12s", "0x10"] { assert_eq!(parse_u64(Some(raw)), None, "for {raw}"); } // One past u64::MAX: rejected, not wrapped or saturated. assert_eq!(parse_u64(Some("18446744073709551616")), None); assert_eq!(parse_u64(Some("18446744073709551615")), Some(u64::MAX)); } #[test] fn nonzero_or_falls_back_when_there_is_no_value() { assert_eq!(nonzero_or(None, 9), 9); assert_eq!(nonzero_or(Some(5), 9), 5); } /// The distinction the two helpers exist for. A knob reading the raw /// parse keeps `0` and treats it as "disable"; one reading /// `nonzero_or` discards it and uses its default. Swapping which /// helper a knob calls is a one-word edit that silently changes /// whether `0` turns the feature off or does nothing at all. #[test] fn zero_survives_parsing_but_is_rejected_as_a_positive_knob() { assert_eq!(parse_u64(Some("0")), Some(0), "kept: 0 can mean disable"); assert_eq!( nonzero_or(parse_u64(Some("0")), 7), 7, "rejected: 0 is not a duration" ); } fn ask(kind: &str) -> PermissionAsk<'_> { PermissionAsk { kind, mcp_server: None, } } #[test] fn acp_permissions_default_to_deny() { for kind in [ "execute", "delete", "move", "switch_mode", "other", "think", "", "new_kind", ] { assert!(!acp_permits(&ask(kind), true), "{kind:?} was allowed"); } } #[test] fn acp_file_tools_and_mcp_servers_are_allowed() { for kind in ["read", "edit", "search"] { assert!(acp_permits(&ask(kind), false), "{kind:?} was refused"); } let mcp = PermissionAsk { kind: "other", mcp_server: Some("hyperhive"), }; assert!(acp_permits(&mcp, false)); } #[test] fn an_mcp_looking_title_does_not_lift_execute_or_fetch() { let execute = PermissionAsk { kind: "execute", mcp_server: Some("hyperhive"), }; assert!(!acp_permits(&execute, true)); let fetch = PermissionAsk { kind: "fetch", mcp_server: Some("hyperhive"), }; assert!(!acp_permits(&fetch, false)); } #[test] fn acp_fetch_follows_the_web_tools_group() { assert!(acp_permits(&ask("fetch"), true)); assert!(!acp_permits(&ask("fetch"), false)); } }