diff --git a/CLAUDE.md b/CLAUDE.md index 6683f878..9d04c554 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -53,7 +53,7 @@ hand-maintained per-file tree drifts out of sync with the code. launch-config layer (tool-group/capability → `--allowedTools`, `--mcp-config` render). - **`hive-runtime/`** — the runtime an agent's turns run on: one - `Runtime` interface (`run`/`compact`/`archive`) with a `claude` backend + `Runtime` interface (`run`/`compact`/`archive`/`canceller`) with a `claude` backend (a pass-through to `hive-claude`) and an `acp` backend (any Agent Client Protocol agent, spawned from the command/args/env nix hands it — no agent is named in Rust). The ACP backend translates its stream into claude's diff --git a/hive-runtime/README.md b/hive-runtime/README.md index d6cbf8ff..581146a7 100644 --- a/hive-runtime/README.md +++ b/hive-runtime/README.md @@ -1,7 +1,7 @@ # hive-runtime The layer an agent's turns are driven through: one `Runtime` interface -(`run`, `compact`, `archive`) with a backend per runtime. +(`run`, `compact`, `archive`, `canceller`) with a backend per runtime. - **claude** — `claude --print` through the `hive-claude` crate's `InfiniteSession`. A pass-through: same spawn, same session handling, same @@ -30,8 +30,16 @@ The crate depends on no hyperhive binary crate, so `hive-agent` and - `loadSession`, to pick its session back up after a harness restart. Without it every restart starts a new session. +## ACP backend: stopping a turn + +A turn is stopped with `session/cancel`: by the `Canceller` handle, or by the +idle watchdog once no `session/update` has arrived for +`Config::idle_timeout`. An agent that has not answered the prompt 10s later +is killed, and the next turn respawns it. A watchdog stop fails the turn +with `AcpError::IdleTimeout` (`IdleKilled` if the agent was killed). This is +the only way a provider error the agent retries without reporting it, such as +an HTTP 429, ends the turn. + ## ACP backend: not yet -`compact` returns `Error::Unsupported`, there is no cancel, and no idle -watchdog (`Config::idle_timeout` is ignored). An agent that retries a failing -provider on its own keeps the turn open until it gives up. +`compact` returns `Error::Unsupported`. diff --git a/hive-runtime/src/acp/mod.rs b/hive-runtime/src/acp/mod.rs index 49cce1f8..7ebf3505 100644 --- a/hive-runtime/src/acp/mod.rs +++ b/hive-runtime/src/acp/mod.rs @@ -6,11 +6,16 @@ mod stream; use std::collections::VecDeque; use std::path::{Path, PathBuf}; +use std::pin::Pin; use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::time::Duration; use hive_claude::{Config, Progress, Sink}; use serde_json::{Value, json}; -use tokio::sync::Mutex; +use tokio::sync::futures::Notified; +use tokio::sync::{Mutex, Notify}; +use tokio::time::Instant; use self::rpc::{Connection, Incoming}; use self::stream::StreamMapper; @@ -25,8 +30,12 @@ const STDERR_TAIL: usize = 20; /// After a turn's reply, how long the agent must stay quiet before the turn /// is taken to be over, and the most that wait may take in all. -const SETTLE_QUIET: std::time::Duration = std::time::Duration::from_millis(500); -const SETTLE_MAX: std::time::Duration = std::time::Duration::from_secs(5); +const SETTLE_QUIET: Duration = Duration::from_millis(500); +const SETTLE_MAX: Duration = Duration::from_secs(5); + +/// How long the agent has to end its turn after `session/cancel` before it is +/// killed. +const CANCEL_GRACE: Duration = Duration::from_secs(10); /// One `session/request_permission` request, as a [`PermissionPolicy`] sees it. #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -80,6 +89,61 @@ pub enum AcpError { NoHttpMcp, #[error("reading the MCP config {path}: {detail}")] McpConfig { path: PathBuf, detail: String }, + #[error( + "ACP agent sent nothing for {silent_secs}s, so its turn was cancelled; a \ + provider error the agent retries without reporting it, such as HTTP 429, looks like this" + )] + IdleTimeout { silent_secs: u64 }, + #[error( + "ACP agent sent nothing for {silent_secs}s and did not stop within {grace_secs}s of \ + `session/cancel`, so it was killed" + )] + IdleKilled { silent_secs: u64, grace_secs: u64 }, + #[error("ACP agent did not stop within {grace_secs}s of `session/cancel`, so it was killed")] + CancelIgnored { grace_secs: u64 }, +} + +/// Stops the turn an [`AcpRuntime`] has in flight, from outside the call +/// driving it. Cheap to clone. +#[derive(Clone)] +pub struct Canceller(Arc); + +#[derive(Default)] +struct CancelState { + in_turn: AtomicBool, + requested: Notify, +} + +impl Canceller { + /// Ask the turn in flight to stop. The agent is sent `session/cancel`, + /// and the turn then ends normally with what it produced so far; an + /// agent that has not ended it within `CANCEL_GRACE` is killed, and the + /// turn fails with [`AcpError::CancelIgnored`]. Returns whether a turn + /// was in flight; with none, this does nothing. + #[must_use] + pub fn cancel(&self) -> bool { + let in_turn = self.0.in_turn.load(Ordering::SeqCst); + if in_turn { + self.0.requested.notify_waiters(); + } + in_turn + } +} + +/// Marks a turn in flight for [`Canceller::cancel`] while it lives. +struct InTurn<'a>(&'a AtomicBool); + +impl Drop for InTurn<'_> { + fn drop(&mut self) { + self.0.store(false, Ordering::SeqCst); + } +} + +/// Why a turn's prompt is being stopped. +#[derive(Clone, Copy)] +enum Stop { + Operator, + Idle, } /// Turns on an ACP agent's durable session. @@ -87,13 +151,17 @@ pub enum AcpError { /// From the [`Config`] it reads `cwd` (the agent's working directory), /// `mcp_config` (converted to the session's `mcpServers`) and /// `system_prompt_file` (prepended to the first prompt of each new session, -/// since ACP has no system prompt). Everything else in it is claude's and is -/// ignored, including `idle_timeout`. +/// since ACP has no system prompt) and `idle_timeout` (with no +/// `session/update` for that long the turn is cancelled, and fails with +/// [`AcpError::IdleTimeout`]). Everything else in it is claude's and is +/// ignored. pub struct AcpRuntime { command: AcpCommand, session_file: PathBuf, permit: PermissionPolicy, live: Mutex>, + cancel: Arc, + cancel_grace: Duration, } /// The running agent process and the session loaded into it. @@ -115,6 +183,8 @@ impl AcpRuntime { session_file, permit, live: Mutex::new(None), + cancel: Arc::default(), + cancel_grace: CANCEL_GRACE, } } @@ -220,38 +290,62 @@ impl AcpRuntime { config: &Config, prompt: &str, sink: &impl Sink, + mut cancelled: Pin<&mut Notified<'_>>, ) -> Result { let cwd = session_cwd(config); let (servers, names) = mcp_servers(config)?; let (session, created) = self.attach(live, &cwd, &servers).await?; - let text = match (created, &config.system_prompt_file) { - (true, Some(path)) => match std::fs::read_to_string(path) { - Ok(system) => format!("{system}\n\n{prompt}"), - Err(e) => { - tracing::warn!(path = %path.display(), error = %e, "system prompt unreadable"); - prompt.to_owned() - } - }, - _ => prompt.to_owned(), - }; + let text = prompt_text(config, prompt, created); let params = json!({ "sessionId": session, "prompt": [{ "type": "text", "text": text }] }); discard_stale(&mut live.conn)?; let mut reply = live.conn.send("session/prompt", params).await?; let mut mapper = StreamMapper::new(names); let mut stderr_tail = VecDeque::new(); + let idle = config.idle_timeout; + let mut quiet_until = idle.map(|window| Instant::now() + window); + let mut stop: Option<(Stop, Instant)> = None; let response = loop { + let now = Instant::now(); tokio::select! { // Notifications first: the ones sent before the reply belong // to this turn and must reach the sink before it ends. biased; incoming = live.conn.incoming.recv() => { let incoming = incoming.unwrap_or(Incoming::Closed); + if matches!(incoming, Incoming::Update(_)) { + quiet_until = idle.map(|window| Instant::now() + window); + } if !deliver(incoming, &session, &mut mapper, &mut stderr_tail, sink) { let stderr_tail = Vec::from(stderr_tail).join("\n"); return Err(AcpError::Exited { stderr_tail }.into()); } } reply = &mut reply => break reply.unwrap_or(Err(AcpError::Closed))?, + () = cancelled.as_mut(), if stop.is_none() => { + tracing::info!("cancelling the ACP turn"); + live.conn.notify("session/cancel", json!({ "sessionId": session })).await?; + stop = Some((Stop::Operator, Instant::now() + self.cancel_grace)); + } + () = tokio::time::sleep_until(quiet_until.unwrap_or(now)), + if stop.is_none() && quiet_until.is_some() => + { + tracing::warn!(?idle, "ACP agent silent; cancelling the turn"); + live.conn.notify("session/cancel", json!({ "sessionId": session })).await?; + stop = Some((Stop::Idle, Instant::now() + self.cancel_grace)); + } + () = tokio::time::sleep_until(stop.map_or(now, |(_, deadline)| deadline)), + if stop.is_some() => + { + let grace_secs = self.cancel_grace.as_secs(); + return Err(match stop { + Some((Stop::Idle, _)) => AcpError::IdleKilled { + silent_secs: idle.unwrap_or_default().as_secs(), + grace_secs, + }, + _ => AcpError::CancelIgnored { grace_secs }, + } + .into()); + } } }; if created { @@ -273,7 +367,18 @@ impl AcpRuntime { for event in mapper.finish() { sink.on_event(&event); } - match response["stopReason"].as_str() { + let stop_reason = match stop { + Some((Stop::Idle, _)) => { + return Err(AcpError::IdleTimeout { + silent_secs: idle.unwrap_or_default().as_secs(), + } + .into()); + } + // Agents differ in the reason they give a cancelled prompt. + Some((Stop::Operator, _)) => Some("cancelled"), + None => response["stopReason"].as_str(), + }; + match stop_reason { Some("end_turn") => {} reason => { let reason = reason.unwrap_or("none"); @@ -297,16 +402,26 @@ impl Runtime for AcpRuntime { tracing::warn!("ACP agent has exited; respawning"); *guard = None; } + // Registered before the turn is marked in flight, so a cancel is + // never missed; one that comes before the prompt is sent stops it + // right after. + let cancelled = self.cancel.requested.notified(); + tokio::pin!(cancelled); + cancelled.as_mut().enable(); + self.cancel.in_turn.store(true, Ordering::SeqCst); + let _in_turn = InTurn(&self.cancel.in_turn); let live = match guard.as_mut() { Some(live) => live, None => guard.insert(self.start(config).await?), }; - let result = self.turn(live, config, prompt, sink).await; + let result = self.turn(live, config, prompt, sink, cancelled).await; if let Err(Error::Acp( AcpError::Closed | AcpError::Exited { .. } | AcpError::Io(_) - | AcpError::Malformed { .. }, + | AcpError::Malformed { .. } + | AcpError::IdleKilled { .. } + | AcpError::CancelIgnored { .. }, )) = &result { *guard = None; @@ -318,6 +433,10 @@ impl Runtime for AcpRuntime { Err(Error::Unsupported("compact")) } + fn canceller(&self) -> Option { + Some(Canceller(self.cancel.clone())) + } + /// Moves the session file aside; the agent keeps its own copy of the /// session, so nothing is lost. A new session whose first prompt was never /// answered is dropped too. @@ -383,6 +502,20 @@ fn discard_stale(conn: &mut Connection) -> std::result::Result<(), AcpError> { Ok(()) } +/// The prompt as sent: the first one of a session carries the system prompt. +fn prompt_text(config: &Config, prompt: &str, first: bool) -> String { + match (first, &config.system_prompt_file) { + (true, Some(path)) => match std::fs::read_to_string(path) { + Ok(system) => format!("{system}\n\n{prompt}"), + Err(e) => { + tracing::warn!(path = %path.display(), error = %e, "system prompt unreadable"); + prompt.to_owned() + } + }, + _ => prompt.to_owned(), + } +} + fn read_id(file: &Path) -> Option { std::fs::read_to_string(file) .ok() @@ -431,8 +564,9 @@ mod tests { use std::collections::BTreeMap; use std::path::Path; use std::sync::Arc; + use std::time::Duration; - use hive_claude::{Config, NoopSink}; + use hive_claude::{Config, NoopSink, Sink}; use super::{AcpError, AcpRuntime}; use crate::{AcpCommand, Error, Runtime}; @@ -444,9 +578,14 @@ mod tests { /// `end_turn`, except the very first one it is sent across restarts, /// which it handles per its mode: /// - /// - `fail`: an error. + /// - `fail`: an error; + /// - `silent`: nothing until `session/cancel`, then `end_turn` — the + /// reply an agent gives a cancelled prompt it was retrying against a + /// rate-limited (HTTP 429) provider without reporting it; + /// - `deaf`: nothing, and `session/cancel` is ignored; + /// - `trickle`: five text chunks 100ms apart, then `end_turn`. const AGENT: &str = r#" -n=0 +n=0 prompt= printf 'start\n' >> "$1.methods" while IFS= read -r line; do id=${line#*\"id\":}; id=${id%%[,\}]*} @@ -469,8 +608,20 @@ while IFS= read -r line; do case $2 in fail) printf '{"jsonrpc":"2.0","id":%s,"error":{"code":-32603,"message":"provider down"}}\n' "$id" ;; + silent) prompt=$id ;; + trickle) + for _ in 1 2 3 4 5; do + sleep 0.1 + printf '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"s%s","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"x"}}}}\n' "$n" + done + printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$id" ;; esac fi ;; + session/cancel) + if [ -n "$prompt" ]; then + printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$prompt" + prompt= + fi ;; esac done "#; @@ -499,6 +650,13 @@ done } } + fn idle_config(dir: &Path, idle_ms: u64) -> Config { + Config { + idle_timeout: Some(Duration::from_millis(idle_ms)), + ..config(dir) + } + } + /// Which of the prompts the agent was sent carried the system prompt. fn carried_system_prompt(dir: &Path) -> Vec { std::fs::read_to_string(dir.join("prompts")) @@ -521,6 +679,22 @@ done std::fs::read_to_string(dir.join("session")).unwrap() } + /// Keeps the stderr lines a turn reports. + #[derive(Default)] + struct Stderr(std::sync::Mutex>); + + impl Sink for Stderr { + fn on_stderr_line(&self, line: &str) { + self.0.lock().unwrap().push(line.to_owned()); + } + } + + async fn prompt_sent(dir: &Path) { + while !dir.join("prompts").exists() { + tokio::time::sleep(Duration::from_millis(10)).await; + } + } + #[tokio::test] async fn a_session_whose_first_prompt_failed_is_retried_not_resumed() { let dir = tempfile::tempdir().unwrap(); @@ -563,4 +737,105 @@ done assert_eq!(count(dir.path(), "session/load"), 1); assert_eq!(recorded(dir.path()), "s1"); } + + #[tokio::test] + async fn a_turn_silent_for_the_idle_window_is_cancelled_and_fails() { + let dir = tempfile::tempdir().unwrap(); + let runtime = runtime(dir.path(), "silent"); + let config = idle_config(dir.path(), 300); + + let stalled = runtime.run(&config, "one", &NoopSink).await; + assert!( + matches!(stalled, Err(Error::Acp(AcpError::IdleTimeout { .. }))), + "{stalled:?}" + ); + assert_eq!(count(dir.path(), "session/cancel"), 1); + + // The agent stopped when asked, so it keeps running and the session + // it answered on is kept. + let next = runtime.run(&config, "two", &NoopSink).await.unwrap(); + assert!(!next.created); + assert_eq!(count(dir.path(), "start"), 1); + assert_eq!(recorded(dir.path()), "s1"); + } + + #[tokio::test] + async fn an_agent_ignoring_the_idle_cancel_is_killed_and_respawned() { + let dir = tempfile::tempdir().unwrap(); + let mut runtime = runtime(dir.path(), "deaf"); + runtime.cancel_grace = Duration::from_millis(200); + let config = idle_config(dir.path(), 200); + + let stalled = runtime.run(&config, "one", &NoopSink).await; + assert!( + matches!(stalled, Err(Error::Acp(AcpError::IdleKilled { .. }))), + "{stalled:?}" + ); + + let next = runtime.run(&config, "two", &NoopSink).await.unwrap(); + assert!(next.created); + assert_eq!(count(dir.path(), "start"), 2); + assert_eq!(count(dir.path(), "session/new"), 1); + assert_eq!(carried_system_prompt(dir.path()), [true, true]); + } + + #[tokio::test] + async fn updates_keep_the_idle_watchdog_from_firing() { + let dir = tempfile::tempdir().unwrap(); + let runtime = runtime(dir.path(), "trickle"); + + // Five chunks 100ms apart: longer than the window in all, never + // silent for as long as it. + let done = runtime + .run(&idle_config(dir.path(), 300), "one", &NoopSink) + .await; + assert!(done.is_ok(), "{done:?}"); + assert_eq!(count(dir.path(), "session/cancel"), 0); + } + + #[tokio::test] + async fn a_cancelled_turn_stops_and_ends_normally() { + let dir = tempfile::tempdir().unwrap(); + let runtime = runtime(dir.path(), "silent"); + let canceller = runtime.canceller().unwrap(); + assert!(!canceller.cancel(), "no turn is in flight yet"); + + let (config, stderr) = (config(dir.path()), Stderr::default()); + let (done, was_in_flight) = tokio::join!(runtime.run(&config, "one", &stderr), async { + prompt_sent(dir.path()).await; + canceller.cancel() + }); + assert!(was_in_flight); + assert!(done.is_ok(), "{done:?}"); + assert_eq!(count(dir.path(), "session/cancel"), 1); + assert!( + stderr + .0 + .lock() + .unwrap() + .contains(&"ACP turn stopped: cancelled".to_owned()), + "{stderr:?}", + stderr = stderr.0.lock().unwrap() + ); + assert!(!canceller.cancel(), "the turn is over"); + } + + #[tokio::test] + async fn an_agent_ignoring_a_cancel_is_killed() { + let dir = tempfile::tempdir().unwrap(); + let mut runtime = runtime(dir.path(), "deaf"); + runtime.cancel_grace = Duration::from_millis(200); + let (canceller, config) = (runtime.canceller().unwrap(), config(dir.path())); + + let (done, _) = tokio::join!(runtime.run(&config, "one", &NoopSink), async { + prompt_sent(dir.path()).await; + canceller.cancel() + }); + assert!( + matches!(done, Err(Error::Acp(AcpError::CancelIgnored { .. }))), + "{done:?}" + ); + runtime.run(&config, "two", &NoopSink).await.unwrap(); + assert_eq!(count(dir.path(), "start"), 2); + } } diff --git a/hive-runtime/src/acp/rpc.rs b/hive-runtime/src/acp/rpc.rs index d3db7cea..6cccfb25 100644 --- a/hive-runtime/src/acp/rpc.rs +++ b/hive-runtime/src/acp/rpc.rs @@ -135,6 +135,12 @@ impl Connection { Ok(rx) } + /// Send a notification: a message the agent does not reply to. + pub(super) async fn notify(&self, method: &'static str, params: Value) -> Result<(), AcpError> { + let message = json!({ "jsonrpc": "2.0", "method": method, "params": params }); + write(&self.stdin, &message).await + } + /// Send a request and wait for its reply, queueing anything else the agent /// sends meanwhile on [`Self::incoming`]. pub(super) async fn request( diff --git a/hive-runtime/src/claude.rs b/hive-runtime/src/claude.rs index c6690069..c2bcc1f8 100644 --- a/hive-runtime/src/claude.rs +++ b/hive-runtime/src/claude.rs @@ -4,7 +4,7 @@ use std::path::PathBuf; use hive_claude::{CompactionPolicy, Config, InfiniteSession, Progress, SessionStore, Sink}; -use crate::{Result, Runtime}; +use crate::{Canceller, Result, Runtime}; /// `claude --print` turns on a titled [`InfiniteSession`], which owns /// resume-or-create and compaction. @@ -39,4 +39,8 @@ impl Runtime for ClaudeRuntime

{ fn archive(&self) -> Result> { Ok(self.store.archive_by_title(&self.title)?) } + + fn canceller(&self) -> Option { + None + } } diff --git a/hive-runtime/src/lib.rs b/hive-runtime/src/lib.rs index 50064625..63a3ba1f 100644 --- a/hive-runtime/src/lib.rs +++ b/hive-runtime/src/lib.rs @@ -21,7 +21,7 @@ mod acp; mod claude; mod spec; -pub use acp::{AcpError, AcpRuntime, PermissionAsk, PermissionPolicy}; +pub use acp::{AcpError, AcpRuntime, Canceller, PermissionAsk, PermissionPolicy}; pub use claude::ClaudeRuntime; pub use hive_claude::{ CompactionPolicy, Config, PercentPolicy, Progress, SessionStore, Sink, Telemetry, TokenUsage, @@ -52,6 +52,11 @@ pub trait Runtime { /// Returns the file the session was moved to, or `None` if there was no /// session to archive. Only call it between turns. fn archive(&self) -> Result>; + + /// A handle that stops the turn in flight from outside [`Self::run`]. + /// `None` for a runtime without one: a claude turn is stopped by + /// signalling its `claude` process. + fn canceller(&self) -> Option; } /// The runtime an agent was configured with, chosen at startup from a @@ -82,6 +87,13 @@ impl Runtime for AgentRuntime

{ Self::Acp(r) => r.archive(), } } + + fn canceller(&self) -> Option { + match self { + Self::Claude(r) => r.canceller(), + Self::Acp(r) => r.canceller(), + } + } } /// Why a runtime operation did not complete.