`topology.json` was a map of `name -> parent | null`, and that value fed the whole agent hierarchy: `<parent>` / `<children>` recipient sentinels, the reparenting API (CLI verb, wire verb, dashboard endpoints, DAG node), the dashboard tree, the rebuild depth sort, and an unconditional bind-mount grant giving every agent RW on its direct children's state. Per the operator's ruling the field goes, and with it all of the above. The file survives as what remains once the value is gone: the roster of agent names, which is the set `ManageRootAgent` grants mounts over. It is now a JSON array; `read` still accepts the old map shape and keeps its keys, so a hive that upgrades across this does not blank its roster (and so no capability holder loses its mounts for the length of that window). Two sites kept their behaviour under a different recipient rather than losing it. Both addressed `<parent>`, which the broker already resolved to `operator` for a root agent, and every agent is now what that fallback called a root: - the harness's turn-failure / plugin-failure notification (`Surface::send_to_parent` -> `send_to_operator`), and - the send allow-list's always-permitted escape hatch, so an agent with a restrictive allow-list still has a way to say it is stuck. What is NOT preserved, deliberately: an agent with no capability no longer sees any other agent's dirs. `ManageRootAgent`'s own grant is unchanged -- still every agent in the roster, still state RW + config RO, still no `harness`. The dashboard's reparenting control (the M0V3 picker) is deleted with its CSS. The tree rendering that reads `ContainerView.parent` is left for the frontend owner -- it degrades to a flat list with the field gone.
791 lines
34 KiB
Rust
791 lines
34 KiB
Rust
//! 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
|
|
//! 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_claude::{Config, InfiniteSession, PercentPolicy, 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 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:
|
|
// 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 <title>`); 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 [`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 [`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")
|
|
}
|
|
|
|
/// 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>;
|
|
|
|
/// 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()),
|
|
},
|
|
)
|
|
}
|
|
|
|
/// Drive one turn end-to-end. The durable [`InfiniteSession`] 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);
|
|
} else {
|
|
// Heuristic: context large AND prompt cache gone cold.
|
|
maybe_auto_reset(bus);
|
|
}
|
|
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)) {
|
|
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) => 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);
|
|
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 _ = session.compact(&config, &sink).await;
|
|
// 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(true);
|
|
}
|
|
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) {
|
|
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);
|
|
}
|
|
|
|
/// 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
|
|
/// [`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())),
|
|
}
|
|
}
|
|
|
|
/// 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'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) {
|
|
let title = session_title();
|
|
match session_store().archive_by_title(&title) {
|
|
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");
|
|
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 claude session");
|
|
bus.emit(LiveEvent::Note {
|
|
text: format!("failed to archive session \"{title}\": {e}"),
|
|
});
|
|
}
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::{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"
|
|
);
|
|
}
|
|
}
|