hive-subagent-mcp: run an agent's subagents on its runtime
The subagent daemon now reads the parent agent's runtime at startup (`hive_runtime::RuntimeSpec`, from the harness's `HIVE_RUNTIME` / `HIVE_ACP_*`, which `mcp.nix` forwards onto its unit). On claude nothing changes. On ACP, each run drives an `AcpRuntime` whose session id is kept per name under the harness dir: `start` archives the old one, `continue` loads it (and fails when none is recorded), `interrupt` sends `session/cancel`, a role goes in front of the first prompt, and permission requests get the answers a claude subagent's tool list gives. The unit loads `backendEnvironmentFile` on ACP only, so the agent can authenticate. The end-of-turn handling moves out of the claude loop into `after_turn` unchanged, so both loops share it. Refs #4391
This commit is contained in:
parent
a483b23ccf
commit
84d4808d54
10 changed files with 689 additions and 71 deletions
|
|
@ -114,8 +114,9 @@ hand-maintained per-file tree drifts out of sync with the code.
|
|||
directly over streamable-http (no stdio bridge), writes task files
|
||||
under `/harness/bash-tasks/`, and records the favorite-tools
|
||||
`bash_commands` stat into turn-stats.sqlite.
|
||||
- **`hive-subagent-mcp/`** — per-agent claude-subagent runner daemon
|
||||
(`hive-subagent-daemon`); spawns nested claude sessions on request and
|
||||
- **`hive-subagent-mcp/`** — per-agent subagent runner daemon
|
||||
(`hive-subagent-daemon`); spawns nested sessions on request, on the
|
||||
parent agent's runtime via `hive-runtime` (claude or ACP), and
|
||||
serves `start`/`continue`/`status`/`interrupt` directly over
|
||||
streamable-http (no stdio bridge), plus a second subagent-facing route
|
||||
carrying `goal_reached`/`need_help` — served per session under a minted
|
||||
|
|
|
|||
1
Cargo.lock
generated
1
Cargo.lock
generated
|
|
@ -2040,6 +2040,7 @@ dependencies = [
|
|||
"clap",
|
||||
"hive-agent-sock",
|
||||
"hive-claude",
|
||||
"hive-runtime",
|
||||
"hive-sh4re",
|
||||
"hive-sock-client",
|
||||
"hive-types",
|
||||
|
|
|
|||
|
|
@ -288,6 +288,31 @@ agent that needs a stable or non-default port.
|
|||
Own systemd unit, defined alongside the other per-agent MCP daemons in
|
||||
`nix/agent-modules/mcp.nix`.
|
||||
|
||||
## Runtime
|
||||
|
||||
A subagent runs on its parent agent's runtime
|
||||
(`services.hyperhive.agent.runtime`). The unit gets the harness's
|
||||
`HIVE_RUNTIME`/`HIVE_ACP_*` variables, and the daemon reads them once at
|
||||
startup (`hive_runtime::RuntimeSpec`). On `claude`, everything else on this
|
||||
page applies as written.
|
||||
|
||||
On `acp`:
|
||||
|
||||
- each run spawns its own copy of the parent's ACP agent, which exits
|
||||
when the run ends. The unit loads `backendEnvironmentFile` for it.
|
||||
- the session id lives in `subagent-acp-session-<name>` under the harness
|
||||
dir. `start` moves an old one aside, `continue` loads it into a fresh
|
||||
agent process, and a `continue` with none recorded fails straight away.
|
||||
- a role goes in front of the session's first prompt, since ACP has no
|
||||
system prompt.
|
||||
- `interrupt` sends `session/cancel`, and `force` changes nothing. The
|
||||
runtime kills an agent that ignores the cancel for 10s.
|
||||
- `model` and `effort` do nothing: the agent runs the model it's
|
||||
configured with.
|
||||
- the agent's permission requests get the same answers a claude subagent's
|
||||
tool list gives: MCP tools and file tools, `fetch` with `web_tools`,
|
||||
`execute` with `execution` (`session::acp_permits`).
|
||||
|
||||
## The tool surface a subagent gets
|
||||
|
||||
Two flags, each covering one half, and neither covering the other:
|
||||
|
|
|
|||
|
|
@ -13,6 +13,7 @@ axum.workspace = true
|
|||
clap.workspace = true
|
||||
hive-agent-sock.workspace = true
|
||||
hive-claude.workspace = true
|
||||
hive-runtime.workspace = true
|
||||
# `permissions::builtin_tools_arg` — the same `--tools` resolution the parent
|
||||
# harness spawns its own claude with, so a subagent's built-in surface is its
|
||||
# parent's rather than a second list that drifts. See `session::build_config`.
|
||||
|
|
|
|||
|
|
@ -51,10 +51,11 @@ async fn main() -> Result<()> {
|
|||
// for that session alone, which is what makes a signal's identity a
|
||||
// property of the endpoint rather than of the payload.
|
||||
let signal_base = format!("http://{}{}", cli.http, hive_subagent_mcp::mcp::SIGNAL_PATH);
|
||||
let state = Arc::new(hive_subagent_mcp::session::State::new(
|
||||
todo_socket,
|
||||
signal_base,
|
||||
));
|
||||
// The parent agent's runtime, from the same variables the harness reads.
|
||||
let runtime = hive_runtime::RuntimeSpec::from_env()?;
|
||||
let state = Arc::new(
|
||||
hive_subagent_mcp::session::State::new(todo_socket, signal_base).on_runtime(runtime),
|
||||
);
|
||||
|
||||
// Serve the MCP tools over streamable-http forever. No background poll
|
||||
// loop to start — unlike the bash daemon's task-file queue, `start`/
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@
|
|||
//! continuation loop records.
|
||||
//!
|
||||
//! **No task files, no restart recovery.** The daemon's only state is an
|
||||
//! in-memory `name -> Option<Cancel>` map (see `State`'s own doc for the
|
||||
//! in-memory `name -> Option<Stopper>` map (see `State`'s own doc for the
|
||||
//! rest of the maps and for what the `None`/`Some` split is for) — all of
|
||||
//! it living for exactly as long as the process is. A daemon restart means
|
||||
//! whatever was running gets killed
|
||||
|
|
@ -74,6 +74,17 @@
|
|||
//! long-lived enough to need in-place compaction — a real follow-up if that
|
||||
//! assumption stops holding, not shipped here.
|
||||
|
||||
//! **A subagent runs on its parent agent's runtime.** `State::runtime` is
|
||||
//! read once at startup from the same `HIVE_RUNTIME`/`HIVE_ACP_*` variables
|
||||
//! the harness reads (`hive_runtime::RuntimeSpec`), forwarded onto this unit
|
||||
//! by `nix/agent-modules/mcp.nix`. On claude, everything above holds as
|
||||
//! written. On ACP, each run drives one `hive_runtime::AcpRuntime` whose
|
||||
//! session id is kept in a per-name file under the harness dir
|
||||
//! (`acp_session_file`): the agent process lives for one run, `interrupt`
|
||||
//! sends `session/cancel` rather than a signal, a failed turn is `Failed`
|
||||
//! and never `Killed`, and `model`/`effort` are not applied — the ACP agent
|
||||
//! runs whatever model it is configured with.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::os::unix::process::ExitStatusExt as _;
|
||||
use std::path::{Path, PathBuf};
|
||||
|
|
@ -81,6 +92,10 @@ use std::sync::{Arc, Mutex, PoisonError};
|
|||
use std::time::{Duration, Instant};
|
||||
|
||||
use hive_claude::{Attach, Cancel, Claude, Config, SessionStore};
|
||||
use hive_runtime::{
|
||||
AcpCommand, AcpRuntime, Canceller, PermissionAsk, PermissionPolicy, Runtime as _, RuntimeSpec,
|
||||
};
|
||||
use hive_sh4re::permissions::ToolGroup;
|
||||
use tokio::sync::oneshot;
|
||||
|
||||
/// How long a `continue` holds its tool call open waiting to find out
|
||||
|
|
@ -196,6 +211,27 @@ enum TurnEnd {
|
|||
Failed(String),
|
||||
}
|
||||
|
||||
/// What a subagent's turns run on.
|
||||
enum SubagentRuntime {
|
||||
/// A `claude` process per turn (`hive_claude::Claude::spawn`).
|
||||
Claude,
|
||||
/// The parent agent's ACP agent, spawned from `command`, with each
|
||||
/// name's session id kept in `sessions` (see [`acp_session_file`]).
|
||||
Acp {
|
||||
command: AcpCommand,
|
||||
sessions: PathBuf,
|
||||
},
|
||||
}
|
||||
|
||||
/// How `interrupt` reaches a tracked turn, per runtime.
|
||||
enum Stopper {
|
||||
/// Signals the turn's `claude` process.
|
||||
Claude(Cancel),
|
||||
/// Sends the ACP agent `session/cancel`; one that ignores it is killed
|
||||
/// by the runtime after its own grace period.
|
||||
Acp(Canceller),
|
||||
}
|
||||
|
||||
/// What settled a resumed turn's bounded wait — the one thing `continue`
|
||||
/// blocks on before it answers.
|
||||
enum ResumeVerdict {
|
||||
|
|
@ -294,9 +330,9 @@ fn searched_location(config: &Config) -> Option<String> {
|
|||
/// and where to push the completion todo. `Arc`-wrapped so the background
|
||||
/// task driving a run outlives the tool call that started it.
|
||||
///
|
||||
/// The map value is `Option<Cancel>`: `None` means `name` is reserved for
|
||||
/// The map value is `Option<Stopper>`: `None` means `name` is reserved for
|
||||
/// an in-flight `start`/`continue` that hasn't reached a confirmed
|
||||
/// `Claude::spawn` yet; `Some(cancel)` means a real process is tracked and
|
||||
/// spawn yet; `Some(stopper)` means a real turn is tracked and
|
||||
/// interruptible. The `None` state exists to close a real TOCTOU window a
|
||||
/// reviewer caught in the original check-then-insert version: checking "is
|
||||
/// `name` free" and committing to it are two different lock acquisitions
|
||||
|
|
@ -320,7 +356,7 @@ fn searched_location(config: &Config) -> Option<String> {
|
|||
/// name across directories at your own risk, the tool doesn't disambiguate
|
||||
/// it (flagged in review when `dir` was added).
|
||||
pub struct State {
|
||||
running: Mutex<HashMap<String, Option<Cancel>>>,
|
||||
running: Mutex<HashMap<String, Option<Stopper>>>,
|
||||
dirs: Mutex<HashMap<String, String>>,
|
||||
killed: Mutex<HashMap<String, i32>>,
|
||||
last_event: Mutex<HashMap<String, Instant>>,
|
||||
|
|
@ -341,6 +377,9 @@ pub struct State {
|
|||
/// to make.
|
||||
signal_base: String,
|
||||
signal_tokens: Mutex<SignalTokens>,
|
||||
/// What a subagent's turns run on: claude unless [`State::on_runtime`]
|
||||
/// says otherwise.
|
||||
runtime: SubagentRuntime,
|
||||
}
|
||||
|
||||
/// Which opaque URL segment belongs to which session — the whole of a
|
||||
|
|
@ -392,9 +431,25 @@ impl State {
|
|||
socket,
|
||||
signal_base,
|
||||
signal_tokens: Mutex::new(SignalTokens::default()),
|
||||
runtime: SubagentRuntime::Claude,
|
||||
}
|
||||
}
|
||||
|
||||
/// Run every subagent on `runtime` — the daemon passes the parent
|
||||
/// agent's, read from its environment at startup. On ACP this resolves
|
||||
/// the harness dir, which only exists inside a container.
|
||||
#[must_use]
|
||||
pub fn on_runtime(mut self, runtime: RuntimeSpec) -> Self {
|
||||
self.runtime = match runtime {
|
||||
RuntimeSpec::Claude => SubagentRuntime::Claude,
|
||||
RuntimeSpec::Acp(command) => SubagentRuntime::Acp {
|
||||
command,
|
||||
sessions: crate::paths::harness_dir(),
|
||||
},
|
||||
};
|
||||
self
|
||||
}
|
||||
|
||||
/// Mint `name` a fresh signal URL: an unguessable token appended to
|
||||
/// [`State::signal_base`], resolvable back to this one session and to no
|
||||
/// other. Called once per spawned run, and the result goes into exactly
|
||||
|
|
@ -649,11 +704,11 @@ impl State {
|
|||
/// Upgrade `name`'s `None` reservation to a real, interruptible process.
|
||||
/// Same key as the reservation, so there is no window in which `name`
|
||||
/// reads as unoccupied between the two.
|
||||
fn track(&self, name: &str, cancel: Cancel) {
|
||||
fn track(&self, name: &str, stopper: Stopper) {
|
||||
self.running
|
||||
.lock()
|
||||
.unwrap_or_else(PoisonError::into_inner)
|
||||
.insert(name.to_owned(), Some(cancel));
|
||||
.insert(name.to_owned(), Some(stopper));
|
||||
}
|
||||
|
||||
/// Drop `name`'s cancel handle but keep the name claimed — the moment
|
||||
|
|
@ -1185,6 +1240,44 @@ fn start_reserved(
|
|||
dir,
|
||||
Some(&signal_url),
|
||||
);
|
||||
start_on_runtime(
|
||||
state,
|
||||
name,
|
||||
config,
|
||||
system_prompt.map(PathBuf::from),
|
||||
trigger,
|
||||
)
|
||||
}
|
||||
|
||||
/// Archive any prior session under `name` and spawn a fresh run on the
|
||||
/// daemon's runtime. `system_prompt` is the role file `config` already
|
||||
/// hands claude; ACP has no system prompt, so it goes to the ACP runtime,
|
||||
/// which puts it in front of the session's first prompt.
|
||||
fn start_on_runtime(
|
||||
state: &Arc<State>,
|
||||
name: &str,
|
||||
config: Config,
|
||||
system_prompt: Option<PathBuf>,
|
||||
trigger: String,
|
||||
) -> anyhow::Result<String> {
|
||||
if let SubagentRuntime::Acp { command, sessions } = &state.runtime {
|
||||
let runtime = acp_runtime(command.clone(), sessions, name);
|
||||
if let Some(archived) = runtime
|
||||
.archive()
|
||||
.map_err(|e| anyhow::anyhow!("archiving the prior `{name}` session failed: {e}"))?
|
||||
{
|
||||
tracing::info!(
|
||||
name,
|
||||
archived = %archived.display(),
|
||||
"start: archived a finished prior session for a fresh start"
|
||||
);
|
||||
}
|
||||
let config = Config {
|
||||
system_prompt_file: system_prompt,
|
||||
..config
|
||||
};
|
||||
return spawn_and_track_acp(state, name, runtime, config, trigger, None);
|
||||
}
|
||||
let store = build_store(&config)?;
|
||||
if store.find_by_title(name).is_some() {
|
||||
tracing::info!(
|
||||
|
|
@ -1345,6 +1438,41 @@ fn continue_reserved(
|
|||
// a token is per run, and this is a new one.
|
||||
let signal_url = state.mint_signal_url(name);
|
||||
let config = build_config(name, model, effort, None, dir, Some(&signal_url));
|
||||
continue_on_runtime(state, name, config, prompt, verdict)
|
||||
}
|
||||
|
||||
/// Resume `name`'s session on the daemon's runtime for one more run.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// On ACP, a name with no recorded session: the ACP runtime would start a
|
||||
/// new one rather than fail, so the check has to happen here.
|
||||
fn continue_on_runtime(
|
||||
state: &Arc<State>,
|
||||
name: &str,
|
||||
config: Config,
|
||||
prompt: String,
|
||||
verdict: &VerdictTx,
|
||||
) -> anyhow::Result<String> {
|
||||
if let SubagentRuntime::Acp { command, sessions } = &state.runtime {
|
||||
let session = acp_session_file(sessions, name);
|
||||
if !session.exists() {
|
||||
anyhow::bail!(
|
||||
"no ACP session is recorded for subagent `{name}` (looked for {}) — `start` \
|
||||
creates one",
|
||||
session.display()
|
||||
);
|
||||
}
|
||||
let runtime = acp_runtime(command.clone(), sessions, name);
|
||||
return spawn_and_track_acp(
|
||||
state,
|
||||
name,
|
||||
runtime,
|
||||
config,
|
||||
prompt,
|
||||
Some(Arc::clone(verdict)),
|
||||
);
|
||||
}
|
||||
// No existence pre-check: claude's own `--resume` is the authority on
|
||||
// whether the session is there, and it errors rather than quietly
|
||||
// starting a fresh one. `verdict` is how that answer gets back to the
|
||||
|
|
@ -1451,7 +1579,7 @@ fn spawn_and_track(
|
|||
) -> anyhow::Result<String> {
|
||||
let running = Claude::spawn(config, attach)
|
||||
.map_err(|e| anyhow::anyhow!("starting the subagent process failed: {e}"))?;
|
||||
state.track(name, running.cancel_handle());
|
||||
state.track(name, Stopper::Claude(running.cancel_handle()));
|
||||
// This turn supersedes whatever the previous one did, including having
|
||||
// been killed — the record is about the turn before this one.
|
||||
state.clear_kill(name);
|
||||
|
|
@ -1483,59 +1611,25 @@ fn spawn_and_track(
|
|||
verdict: verdict.clone(),
|
||||
};
|
||||
let end = classify_end(running.wait(&prompt, &sink).await, searched.as_deref());
|
||||
log_turn_end(&task_name, &end);
|
||||
// The todo is how a turn's end reaches an agent that is no
|
||||
// longer looking — so it is pushed for every end *except* the
|
||||
// one the caller is being handed as a tool-call error right now.
|
||||
// `settle` saying the verdict was delivered is what makes that
|
||||
// certain: a `continue` whose grace had already run out gets
|
||||
// `false` here and its todo, same as before. Only the first turn
|
||||
// can ever win this — the sender is consumed — which is right,
|
||||
// since only the first turn is one a caller is still waiting on.
|
||||
let reported = settle(verdict.as_ref(), ResumeVerdict::Ended(end.clone()));
|
||||
|
||||
if !matches!(end, TurnEnd::Complete) {
|
||||
// A killed or failed turn ends the run, goal or not: there is
|
||||
// nothing to re-prompt a child that isn't there any more, and
|
||||
// these two ends already have records of their own.
|
||||
state.finish_turn(&task_name, &end);
|
||||
if !(matches!(end, TurnEnd::Failed(_)) && reported) {
|
||||
push_turn_end_todo(&state.socket, &task_name, &end, None).await;
|
||||
}
|
||||
let Some(next) = after_turn(&state, &task_name, end, verdict.as_ref()).await else {
|
||||
return;
|
||||
}
|
||||
|
||||
match state.plan_after_turn(&task_name) {
|
||||
Continuation::Stop(stop) => {
|
||||
state.record_stop(&task_name, stop.clone());
|
||||
state.finish_turn(&task_name, &end);
|
||||
write_stop_to_report(state.report_file(&task_name), &task_name, &stop).await;
|
||||
push_turn_end_todo(&state.socket, &task_name, &end, Some(&stop)).await;
|
||||
return;
|
||||
};
|
||||
match Claude::spawn(&config, &Attach::Resume(task_name.clone())) {
|
||||
Ok(next_running) => {
|
||||
state.track(&task_name, Stopper::Claude(next_running.cancel_handle()));
|
||||
state.note_event(&task_name);
|
||||
running = next_running;
|
||||
prompt = next;
|
||||
searched = resume_searched.clone();
|
||||
}
|
||||
Continuation::Continue { prompt: next } => {
|
||||
// The name stays claimed across the gap — see
|
||||
// `between_turns` for what a concurrent `start` would
|
||||
// otherwise be able to do with it.
|
||||
state.between_turns(&task_name);
|
||||
match Claude::spawn(&config, &Attach::Resume(task_name.clone())) {
|
||||
Ok(next_running) => {
|
||||
state.track(&task_name, next_running.cancel_handle());
|
||||
state.note_event(&task_name);
|
||||
running = next_running;
|
||||
prompt = next;
|
||||
searched = resume_searched.clone();
|
||||
}
|
||||
Err(e) => {
|
||||
let end = TurnEnd::Failed(format!(
|
||||
"claude error: starting the next goal turn failed: {e}"
|
||||
));
|
||||
log_turn_end(&task_name, &end);
|
||||
state.finish_turn(&task_name, &end);
|
||||
push_turn_end_todo(&state.socket, &task_name, &end, None).await;
|
||||
return;
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
let end = TurnEnd::Failed(format!(
|
||||
"claude error: starting the next goal turn failed: {e}"
|
||||
));
|
||||
log_turn_end(&task_name, &end);
|
||||
state.finish_turn(&task_name, &end);
|
||||
push_turn_end_todo(&state.socket, &task_name, &end, None).await;
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1544,6 +1638,136 @@ fn spawn_and_track(
|
|||
Ok(format!("subagent `{name}` started"))
|
||||
}
|
||||
|
||||
/// Everything a run does once one of its turns has ended, on either
|
||||
/// runtime: report it, and either stop the run or claim the name for the
|
||||
/// next goal turn. Returns that turn's prompt, or `None` once the run is
|
||||
/// over.
|
||||
async fn after_turn(
|
||||
state: &State,
|
||||
name: &str,
|
||||
end: TurnEnd,
|
||||
verdict: Option<&VerdictTx>,
|
||||
) -> Option<String> {
|
||||
log_turn_end(name, &end);
|
||||
// The todo is how a turn's end reaches an agent that is no
|
||||
// longer looking — so it is pushed for every end *except* the
|
||||
// one the caller is being handed as a tool-call error right now.
|
||||
// `settle` saying the verdict was delivered is what makes that
|
||||
// certain: a `continue` whose grace had already run out gets
|
||||
// `false` here and its todo, same as before. Only the first turn
|
||||
// can ever win this — the sender is consumed — which is right,
|
||||
// since only the first turn is one a caller is still waiting on.
|
||||
let reported = settle(verdict, ResumeVerdict::Ended(end.clone()));
|
||||
|
||||
if !matches!(end, TurnEnd::Complete) {
|
||||
// A killed or failed turn ends the run, goal or not: there is
|
||||
// nothing to re-prompt a child that isn't there any more, and
|
||||
// these two ends already have records of their own.
|
||||
state.finish_turn(name, &end);
|
||||
if !(matches!(end, TurnEnd::Failed(_)) && reported) {
|
||||
push_turn_end_todo(&state.socket, name, &end, None).await;
|
||||
}
|
||||
return None;
|
||||
}
|
||||
|
||||
match state.plan_after_turn(name) {
|
||||
Continuation::Stop(stop) => {
|
||||
state.record_stop(name, stop.clone());
|
||||
state.finish_turn(name, &end);
|
||||
write_stop_to_report(state.report_file(name), name, &stop).await;
|
||||
push_turn_end_todo(&state.socket, name, &end, Some(&stop)).await;
|
||||
None
|
||||
}
|
||||
Continuation::Continue { prompt } => {
|
||||
// The name stays claimed across the gap — see
|
||||
// `between_turns` for what a concurrent `start` would
|
||||
// otherwise be able to do with it.
|
||||
state.between_turns(name);
|
||||
Some(prompt)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// [`spawn_and_track`] on ACP: every turn of the run goes through
|
||||
/// `runtime`, whose agent process is spawned by the first turn and killed
|
||||
/// when the run ends and the runtime is dropped. Returns before the agent
|
||||
/// exists, so a spawn failure is the first turn's `Failed` end rather than
|
||||
/// this call's error.
|
||||
fn spawn_and_track_acp(
|
||||
state: &Arc<State>,
|
||||
name: &str,
|
||||
runtime: AcpRuntime,
|
||||
config: Config,
|
||||
prompt: String,
|
||||
verdict: Option<VerdictTx>,
|
||||
) -> anyhow::Result<String> {
|
||||
let canceller = runtime
|
||||
.canceller()
|
||||
.ok_or_else(|| anyhow::anyhow!("the ACP runtime offers no cancel handle"))?;
|
||||
state.track(name, Stopper::Acp(canceller.clone()));
|
||||
state.clear_kill(name);
|
||||
state.note_event(name);
|
||||
|
||||
let state = Arc::clone(state);
|
||||
let task_name = name.to_owned();
|
||||
tokio::spawn(async move {
|
||||
let mut prompt = prompt;
|
||||
loop {
|
||||
let sink = LivenessSink {
|
||||
state: Arc::clone(&state),
|
||||
name: task_name.clone(),
|
||||
verdict: verdict.clone(),
|
||||
};
|
||||
let end = match runtime.run(&config, &prompt, &sink).await {
|
||||
Ok(_) => TurnEnd::Complete,
|
||||
Err(e) => TurnEnd::Failed(format!("ACP error: {e}")),
|
||||
};
|
||||
let Some(next) = after_turn(&state, &task_name, end, verdict.as_ref()).await else {
|
||||
return;
|
||||
};
|
||||
state.track(&task_name, Stopper::Acp(canceller.clone()));
|
||||
state.note_event(&task_name);
|
||||
prompt = next;
|
||||
}
|
||||
});
|
||||
|
||||
Ok(format!("subagent `{name}` started"))
|
||||
}
|
||||
|
||||
/// The file `name`'s ACP session id is kept in, under `sessions`. Named
|
||||
/// like the per-session MCP config and system-prompt files beside it.
|
||||
fn acp_session_file(sessions: &Path, name: &str) -> PathBuf {
|
||||
sessions.join(format!("subagent-acp-session-{name}"))
|
||||
}
|
||||
|
||||
/// The ACP runtime `name`'s run drives, keeping its session id in
|
||||
/// [`acp_session_file`] and answering permission requests per
|
||||
/// [`acp_permits`].
|
||||
fn acp_runtime(command: AcpCommand, sessions: &Path, name: &str) -> AcpRuntime {
|
||||
let groups = hive_sh4re::permissions::effective_tool_groups();
|
||||
let permit: PermissionPolicy =
|
||||
Arc::new(move |ask: &PermissionAsk<'_>| acp_permits(ask, &groups));
|
||||
AcpRuntime::new(command, acp_session_file(sessions, name), permit)
|
||||
}
|
||||
|
||||
/// Whether an ACP subagent may run the tool call it asks about.
|
||||
/// Default-deny, allowing what a claude subagent gets from `build_config`:
|
||||
///
|
||||
/// - `other`-kind calls to tools of the MCP servers the session was handed;
|
||||
/// - the file tools: `read`, `edit`, `search`;
|
||||
/// - `fetch` with the `web_tools` group;
|
||||
/// - `execute` with the `execution` group, as a claude subagent gets `Bash`
|
||||
/// ([`hive_sh4re::permissions::subagent_builtin_tools_for`]).
|
||||
fn acp_permits(ask: &PermissionAsk<'_>, groups: &[ToolGroup]) -> bool {
|
||||
match ask.kind {
|
||||
"other" => ask.mcp_server.is_some(),
|
||||
"read" | "edit" | "search" => true,
|
||||
"fetch" => groups.contains(&ToolGroup::WebTools),
|
||||
"execute" => groups.contains(&ToolGroup::Execution),
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Log a turn's end at the level its severity deserves: a completion is
|
||||
/// unremarkable, the other two are not.
|
||||
fn log_turn_end(name: &str, end: &TurnEnd) {
|
||||
|
|
@ -1740,11 +1964,15 @@ pub fn status(state: &State, name: &str, dir: Option<&str>) -> anyhow::Result<St
|
|||
// proof this daemon ran the session, which is what the lookup asks.
|
||||
let session_exists =
|
||||
if facts.occupancy.is_none() && facts.killed.is_none() && facts.stop.is_none() {
|
||||
// No signal surface in this config: it exists only to resolve the
|
||||
// session store, and rendering a subagent's MCP config for a
|
||||
// read-only status check would be writing a file for nobody.
|
||||
let config = build_config(name, None, None, None, dir.as_deref(), None);
|
||||
build_store(&config)?.find_by_title(name).is_some()
|
||||
if let SubagentRuntime::Acp { sessions, .. } = &state.runtime {
|
||||
acp_session_file(sessions, name).exists()
|
||||
} else {
|
||||
// No signal surface in this config: it exists only to resolve
|
||||
// the session store, and rendering a subagent's MCP config for
|
||||
// a read-only status check would be writing a file for nobody.
|
||||
let config = build_config(name, None, None, None, dir.as_deref(), None);
|
||||
build_store(&config)?.find_by_title(name).is_some()
|
||||
}
|
||||
} else {
|
||||
false
|
||||
};
|
||||
|
|
@ -1986,10 +2214,17 @@ pub fn interrupt(state: &State, name: &str, force: bool) -> anyhow::Result<Strin
|
|||
shortly"
|
||||
);
|
||||
}
|
||||
Some(Some(cancel)) => {
|
||||
Some(Some(stopper)) => {
|
||||
state.record_stop(name, StopReason::Cancelled);
|
||||
drop(running);
|
||||
cancel.cancel(force);
|
||||
match stopper {
|
||||
Stopper::Claude(cancel) => cancel.cancel(force),
|
||||
// `force` has no stronger form on ACP: the runtime already kills
|
||||
// an agent that ignores `session/cancel`.
|
||||
Stopper::Acp(canceller) => {
|
||||
let _ = canceller.cancel();
|
||||
}
|
||||
}
|
||||
Ok(format!(
|
||||
"interrupt sent to subagent `{name}` — the run is cancelled, not just the turn \
|
||||
that was in flight, so no further goal turn will start. `continue` is what \
|
||||
|
|
@ -3848,4 +4083,258 @@ mod tests {
|
|||
"a var naming nothing restricts nothing — default open"
|
||||
);
|
||||
}
|
||||
|
||||
/// An ACP agent in plain `sh`, after `hive-runtime`'s own test agent. It
|
||||
/// has one session, `s1`, appends every method it is sent to
|
||||
/// `<log>.methods` and every `session/prompt` line to `<log>`, and
|
||||
/// answers each prompt with `end_turn` — or, in `hold` mode, only once
|
||||
/// `session/cancel` arrives.
|
||||
const ACP_AGENT: &str = r#"
|
||||
prompt=
|
||||
while IFS= read -r line; do
|
||||
id=${line#*\"id\":}; id=${id%%[,\}]*}
|
||||
m=${line#*\"method\":\"}; m=${m%%\"*}
|
||||
printf '%s\n' "$m" >> "$1.methods"
|
||||
case $m in
|
||||
initialize)
|
||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"protocolVersion":1,"agentCapabilities":{"loadSession":true,"mcpCapabilities":{"http":true}}}}\n' "$id" ;;
|
||||
session/new)
|
||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"sessionId":"s1"}}\n' "$id" ;;
|
||||
session/load)
|
||||
printf '{"jsonrpc":"2.0","id":%s,"result":{}}\n' "$id" ;;
|
||||
session/prompt)
|
||||
printf '%s\n' "$line" >> "$1"
|
||||
if [ "$2" = hold ]; then
|
||||
prompt=$id
|
||||
else
|
||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$id"
|
||||
fi ;;
|
||||
session/cancel)
|
||||
if [ -n "$prompt" ]; then
|
||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"cancelled"}}\n' "$prompt"
|
||||
prompt=
|
||||
fi ;;
|
||||
esac
|
||||
done
|
||||
"#;
|
||||
|
||||
fn scratch_dir(tag: &str) -> PathBuf {
|
||||
let dir = std::env::temp_dir().join(format!(
|
||||
"hive-subagent-{tag}-{}-{:?}",
|
||||
std::process::id(),
|
||||
std::thread::current().id()
|
||||
));
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
std::fs::create_dir_all(&dir).expect("scratch dir");
|
||||
dir
|
||||
}
|
||||
|
||||
/// `state` running its subagents on [`ACP_AGENT`] in `mode`, keeping
|
||||
/// sessions and the agent's logs in `dir`.
|
||||
fn on_acp_agent(mut state: State, dir: &Path, mode: &str) -> State {
|
||||
let script = dir.join("agent.sh");
|
||||
std::fs::write(&script, ACP_AGENT).expect("write the ACP agent");
|
||||
state.runtime = SubagentRuntime::Acp {
|
||||
command: AcpCommand {
|
||||
command: "/bin/sh".into(),
|
||||
args: vec![
|
||||
script.display().to_string(),
|
||||
dir.join("prompts").display().to_string(),
|
||||
mode.to_owned(),
|
||||
],
|
||||
env: std::collections::BTreeMap::new(),
|
||||
},
|
||||
sessions: dir.to_path_buf(),
|
||||
};
|
||||
state
|
||||
}
|
||||
|
||||
/// A config that would run [`fake_claude`], so a test on the ACP runtime
|
||||
/// can show that claude was never spawned, with claude's session store
|
||||
/// in `dir` rather than the ambient `HOME`.
|
||||
fn claude_config(dir: &Path) -> Config {
|
||||
Config {
|
||||
program: Some(fake_claude(dir, &dir.join("claude-turns"))),
|
||||
cwd: Some(dir.to_path_buf()),
|
||||
claude_config_dir: Some(dir.join("claude-home")),
|
||||
..Config::default()
|
||||
}
|
||||
}
|
||||
|
||||
fn lines(path: &Path) -> Vec<String> {
|
||||
std::fs::read_to_string(path)
|
||||
.map(|body| body.lines().map(str::to_owned).collect())
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
async fn until(what: &str, done: impl Fn() -> bool) {
|
||||
for _ in 0..200 {
|
||||
if done() {
|
||||
return;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(50)).await;
|
||||
}
|
||||
panic!("timed out waiting until {what}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_claude_parent_runs_its_subagents_on_claude() {
|
||||
let dir = scratch_dir("runtime-claude");
|
||||
let state = Arc::new(State::new(PathBuf::from("/dev/null"), signal_url()));
|
||||
assert!(state.reserve("n"));
|
||||
|
||||
start_on_runtime(&state, "n", claude_config(&dir), None, "go".into())
|
||||
.expect("the fake claude spawns");
|
||||
until("claude runs the turn", || {
|
||||
turns_spawned(&dir.join("claude-turns")) == 1
|
||||
})
|
||||
.await;
|
||||
interrupt(&state, "n", false).expect("a running claude turn is interruptible");
|
||||
// `interrupt` drops the tracking entry itself; the run is over once
|
||||
// `finish_turn` has cleared its liveness clock too.
|
||||
until("the run ends", || state.last_event_age("n").is_none()).await;
|
||||
|
||||
assert!(
|
||||
!acp_session_file(&dir, "n").exists() && !dir.join("prompts.methods").exists(),
|
||||
"a claude agent's subagent never touches an ACP agent"
|
||||
);
|
||||
std::fs::remove_dir_all(&dir).ok();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_acp_parent_runs_its_subagents_on_its_acp_agent() {
|
||||
let dir = scratch_dir("runtime-acp");
|
||||
let state = State::new(PathBuf::from("/dev/null"), signal_url());
|
||||
let state = Arc::new(on_acp_agent(state, &dir, "answer"));
|
||||
let role = dir.join("role.md");
|
||||
std::fs::write(&role, "ROLE-TEXT").expect("write the role");
|
||||
|
||||
assert!(state.reserve("n"));
|
||||
start_on_runtime(
|
||||
&state,
|
||||
"n",
|
||||
claude_config(&dir),
|
||||
Some(role),
|
||||
"TRIGGER-ONE".into(),
|
||||
)
|
||||
.expect("the ACP run starts");
|
||||
until("the start's run ends", || state.occupancy("n").is_none()).await;
|
||||
|
||||
let prompts = lines(&dir.join("prompts"));
|
||||
assert_eq!(prompts.len(), 1, "{prompts:?}");
|
||||
assert!(
|
||||
prompts[0].contains("ROLE-TEXT") && prompts[0].contains("TRIGGER-ONE"),
|
||||
"the first prompt carries the role, since ACP has no system prompt: {prompts:?}"
|
||||
);
|
||||
assert_eq!(
|
||||
std::fs::read_to_string(acp_session_file(&dir, "n")).unwrap(),
|
||||
"s1"
|
||||
);
|
||||
|
||||
assert!(state.reserve("n"));
|
||||
let (tx, rx) = oneshot::channel();
|
||||
let verdict: VerdictTx = Arc::new(Mutex::new(Some(tx)));
|
||||
continue_on_runtime(
|
||||
&state,
|
||||
"n",
|
||||
claude_config(&dir),
|
||||
"TRIGGER-TWO".into(),
|
||||
&verdict,
|
||||
)
|
||||
.expect("the ACP continue starts");
|
||||
await_resume(rx).await.expect("the resumed turn lands");
|
||||
until("the continue's run ends", || state.occupancy("n").is_none()).await;
|
||||
|
||||
let prompts = lines(&dir.join("prompts"));
|
||||
assert_eq!(prompts.len(), 2, "{prompts:?}");
|
||||
assert!(
|
||||
prompts[1].contains("TRIGGER-TWO") && !prompts[1].contains("ROLE-TEXT"),
|
||||
"the continue resumes the session rather than starting one: {prompts:?}"
|
||||
);
|
||||
assert!(
|
||||
lines(&dir.join("prompts.methods")).contains(&"session/load".to_owned()),
|
||||
"the continue's agent process loads the recorded session"
|
||||
);
|
||||
assert_eq!(
|
||||
turns_spawned(&dir.join("claude-turns")),
|
||||
0,
|
||||
"an ACP agent's subagent never runs claude"
|
||||
);
|
||||
std::fs::remove_dir_all(&dir).ok();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_acp_continue_with_no_recorded_session_is_refused() {
|
||||
let dir = scratch_dir("runtime-acp-missing");
|
||||
let state = State::new(PathBuf::from("/dev/null"), signal_url());
|
||||
let state = Arc::new(on_acp_agent(state, &dir, "answer"));
|
||||
let (tx, _rx) = oneshot::channel();
|
||||
let verdict: VerdictTx = Arc::new(Mutex::new(Some(tx)));
|
||||
|
||||
let err = continue_on_runtime(&state, "n", claude_config(&dir), "go".into(), &verdict)
|
||||
.expect_err("nothing to resume");
|
||||
assert!(err.to_string().contains("no ACP session"), "{err}");
|
||||
assert!(
|
||||
!dir.join("prompts.methods").exists(),
|
||||
"the agent is never started"
|
||||
);
|
||||
std::fs::remove_dir_all(&dir).ok();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_interrupt_cancels_an_acp_goal_run() {
|
||||
let dir = scratch_dir("runtime-acp-cancel");
|
||||
let state = Arc::new(on_acp_agent(
|
||||
mid_run("n", "carry on indefinitely", 5),
|
||||
&dir,
|
||||
"hold",
|
||||
));
|
||||
start_on_runtime(&state, "n", claude_config(&dir), None, "go".into())
|
||||
.expect("the ACP run starts");
|
||||
until("the first prompt reaches the agent", || {
|
||||
dir.join("prompts").exists()
|
||||
})
|
||||
.await;
|
||||
|
||||
interrupt(&state, "n", false).expect("a running ACP turn is interruptible");
|
||||
until("the run ends", || state.last_event_age("n").is_none()).await;
|
||||
|
||||
assert_eq!(lines(&dir.join("prompts")).len(), 1, "no second goal turn");
|
||||
assert!(
|
||||
lines(&dir.join("prompts.methods")).contains(&"session/cancel".to_owned()),
|
||||
"the turn was stopped with `session/cancel`"
|
||||
);
|
||||
assert_eq!(state.stop_reason("n"), Some(StopReason::Cancelled));
|
||||
std::fs::remove_dir_all(&dir).ok();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_acp_subagent_gets_the_tools_a_claude_subagent_would() {
|
||||
let ask = |kind| PermissionAsk {
|
||||
kind,
|
||||
mcp_server: None,
|
||||
};
|
||||
let none: &[ToolGroup] = &[];
|
||||
let both = &[ToolGroup::WebTools, ToolGroup::Execution];
|
||||
for kind in ["read", "edit", "search"] {
|
||||
assert!(acp_permits(&ask(kind), none), "{kind}");
|
||||
}
|
||||
for kind in ["fetch", "execute"] {
|
||||
assert!(!acp_permits(&ask(kind), none), "{kind} needs its group");
|
||||
assert!(acp_permits(&ask(kind), both), "{kind} with its group");
|
||||
}
|
||||
assert!(
|
||||
!acp_permits(&ask("delete"), both),
|
||||
"an unnamed kind is refused"
|
||||
);
|
||||
assert!(
|
||||
!acp_permits(&ask("other"), both),
|
||||
"other outside the MCP servers"
|
||||
);
|
||||
let mcp = PermissionAsk {
|
||||
kind: "other",
|
||||
mcp_server: Some("subagent_control"),
|
||||
};
|
||||
assert!(acp_permits(&mcp, none));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -228,6 +228,10 @@ in
|
|||
On `"acp"`, `model`, `effortLevel` and `autoCompact` have no effect
|
||||
(the model is whatever the agent is configured with), and neither the
|
||||
web UI's cancel button nor `/compact` works yet.
|
||||
|
||||
Subagents (`hive-subagent-daemon`) run on the same runtime. On
|
||||
`"acp"`, each subagent run spawns its own copy of the ACP agent, which
|
||||
also gets `services.hyperhive.agent.backendEnvironmentFile`.
|
||||
'';
|
||||
};
|
||||
|
||||
|
|
|
|||
|
|
@ -391,6 +391,13 @@ in
|
|||
# `null` when the agent has no groups declared, which systemd drops
|
||||
# — the same "absent" the harness itself would see.
|
||||
HIVE_TOOL_GROUPS = config.systemd.services.hive-agent.environment.HIVE_TOOL_GROUPS or null;
|
||||
# The harness's runtime selection, forwarded the same way, so an ACP
|
||||
# agent's subagents run on its ACP agent. All four are absent on a
|
||||
# claude agent, which the daemon reads as claude.
|
||||
HIVE_RUNTIME = config.systemd.services.hive-agent.environment.HIVE_RUNTIME or null;
|
||||
HIVE_ACP_COMMAND = config.systemd.services.hive-agent.environment.HIVE_ACP_COMMAND or null;
|
||||
HIVE_ACP_ARGS = config.systemd.services.hive-agent.environment.HIVE_ACP_ARGS or null;
|
||||
HIVE_ACP_ENV = config.systemd.services.hive-agent.environment.HIVE_ACP_ENV or null;
|
||||
# Same `services.hyperhive.agent.availableModels` the harness's own assertions gate
|
||||
# the primary session's model against, so a subagent can't be spawned
|
||||
# on a model the operator didn't make available to this agent. The
|
||||
|
|
@ -451,7 +458,19 @@ in
|
|||
}
|
||||
// lib.optionalAttrs (containerMemoryMaxBytes != null) {
|
||||
MemoryHigh = toString subagentMemoryHigh;
|
||||
};
|
||||
}
|
||||
//
|
||||
lib.optionalAttrs
|
||||
(
|
||||
config.services.hyperhive.agent.runtime == "acp"
|
||||
&& config.services.hyperhive.agent.backendEnvironmentFile != null
|
||||
)
|
||||
{
|
||||
# An ACP subagent authenticates to its provider with the same
|
||||
# credentials as the harness's ACP agent. Claude subagents are
|
||||
# not handed this file.
|
||||
EnvironmentFile = "-${config.services.hyperhive.agent.backendEnvironmentFile}";
|
||||
};
|
||||
};
|
||||
|
||||
# Persistent streamable-http MCP daemon for the built-in hyperhive
|
||||
|
|
|
|||
|
|
@ -134,6 +134,10 @@ in
|
|||
inherit pkgs self nixosSystem;
|
||||
inherit (pkgs) lib;
|
||||
};
|
||||
module-eval-agent-runtime = import ./module-eval/agent-runtime.nix {
|
||||
inherit pkgs self nixosSystem;
|
||||
inherit (pkgs) lib;
|
||||
};
|
||||
module-eval-agent-icon = import ./module-eval/agent-icon.nix {
|
||||
inherit pkgs self nixosSystem;
|
||||
inherit (pkgs) lib;
|
||||
|
|
|
|||
73
nix/module-eval/agent-runtime.nix
Normal file
73
nix/module-eval/agent-runtime.nix
Normal file
|
|
@ -0,0 +1,73 @@
|
|||
# `checks.module-eval-agent-runtime` — see ./lib.nix for the shared
|
||||
# rationale (why this suite exists, naming convention, "evaluates
|
||||
# not executes").
|
||||
{
|
||||
pkgs,
|
||||
lib,
|
||||
self,
|
||||
nixosSystem,
|
||||
}:
|
||||
let
|
||||
inherit
|
||||
(import ./lib.nix {
|
||||
inherit
|
||||
pkgs
|
||||
lib
|
||||
self
|
||||
nixosSystem
|
||||
;
|
||||
})
|
||||
agentWith
|
||||
runGroup
|
||||
;
|
||||
|
||||
backendEnv = "/agents/a1/harness/backend.env";
|
||||
|
||||
# Both arms carry a backend file, so the claude arm's lack of one on the
|
||||
# subagent unit is the gate and not a missing input.
|
||||
agentOn =
|
||||
runtime:
|
||||
agentWith {
|
||||
services.hyperhive.agent = {
|
||||
inherit runtime;
|
||||
backendEnvironmentFile = backendEnv;
|
||||
acp.command = "/bin/agent";
|
||||
acp.args = [ "acp" ];
|
||||
};
|
||||
};
|
||||
|
||||
claude = agentOn "claude";
|
||||
acp = agentOn "acp";
|
||||
|
||||
subagent = machine: machine.systemd.services.hive-subagent-daemon;
|
||||
harness = machine: machine.systemd.services.hive-agent;
|
||||
runtimeVars = [
|
||||
"HIVE_RUNTIME"
|
||||
"HIVE_ACP_COMMAND"
|
||||
"HIVE_ACP_ARGS"
|
||||
"HIVE_ACP_ENV"
|
||||
];
|
||||
cases = [
|
||||
{
|
||||
# The daemon reads these with `hive_runtime::RuntimeSpec`, the same
|
||||
# parser as the harness, so equal values mean the same runtime.
|
||||
name = "an ACP agent's subagent daemon gets the harness's runtime selection";
|
||||
ok = lib.all (
|
||||
var: (subagent acp).environment.${var} or null == (harness acp).environment.${var}
|
||||
) runtimeVars;
|
||||
}
|
||||
{
|
||||
name = "an ACP agent's subagent daemon loads the backend credentials";
|
||||
ok = (subagent acp).serviceConfig.EnvironmentFile or null == "-${backendEnv}";
|
||||
}
|
||||
{
|
||||
# Unset is what `RuntimeSpec` reads as claude, and a claude subagent
|
||||
# is not handed the backend file.
|
||||
name = "a claude agent's subagent daemon has no runtime selection and no backend file";
|
||||
ok =
|
||||
lib.all (var: (subagent claude).environment.${var} or null == null) runtimeVars
|
||||
&& !((subagent claude).serviceConfig ? EnvironmentFile);
|
||||
}
|
||||
];
|
||||
in
|
||||
runGroup "agent-runtime" cases
|
||||
Loading…
Reference in a new issue