Watch
0
0
Fork
You've already forked hyperhive
0

hive-runtime: cancel and idle watchdog for ACP turns

`Runtime` gets a fourth operation, `canceller()`: a handle that stops the
turn in flight from outside `run`. The ACP backend returns one; claude
returns `None`, because the harness stops a claude turn by signalling the
`claude` process, and that path is unchanged.

Both stops send the agent `session/cancel`:

- `Canceller::cancel()`, when asked from outside. The turn then ends
  normally, reported with stop reason `cancelled` whatever reason the agent
  gives. opencode 1.15.10, for one, answers a cancelled prompt with
  `end_turn` (`acp/agent.ts` `prompt()` always returns `end_turn`).
- The idle watchdog, once no `session/update` has arrived for
  `Config::idle_timeout`, the same field claude's watchdog reads. The turn
  fails with `AcpError::IdleTimeout`.

An agent that has not answered the prompt 10s after `session/cancel` is
killed (`IdleKilled` / `CancelIgnored`), and the next turn respawns it.

The watchdog is also what ends a turn stuck on a provider HTTP 429.
opencode 1.15.10 retries a retryable provider error with no attempt limit
(`session/retry.ts` `policy`, `session/processor.ts` `Effect.retry`) and
forwards neither `session.status` nor `session.error` over ACP (its
`handleEvent` only handles `permission.asked` and `message.part.*`). So the
ACP client sees nothing at all until the provider recovers. The
`IdleTimeout` message says a silently retried provider error looks like
this.

Refs #4391
This commit is contained in:
atlas 2026-09-29 23:22:50 +02:00
commit bcbeb8ac9c
6 changed files with 333 additions and 28 deletions

View file

@ -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`.

View file

@ -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<CancelState>);
#[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<Option<Live>>,
cancel: Arc<CancelState>,
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<Progress> {
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<Canceller> {
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<String> {
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<bool> {
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<Vec<String>>);
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);
}
}

View file

@ -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(

View file

@ -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<P: CompactionPolicy> Runtime for ClaudeRuntime<P> {
fn archive(&self) -> Result<Option<PathBuf>> {
Ok(self.store.archive_by_title(&self.title)?)
}
fn canceller(&self) -> Option<Canceller> {
None
}
}

View file

@ -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<Option<PathBuf>>;
/// 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<Canceller>;
}
/// The runtime an agent was configured with, chosen at startup from a
@ -82,6 +87,13 @@ impl<P: CompactionPolicy> Runtime for AgentRuntime<P> {
Self::Acp(r) => r.archive(),
}
}
fn canceller(&self) -> Option<Canceller> {
match self {
Self::Claude(r) => r.canceller(),
Self::Acp(r) => r.canceller(),
}
}
}
/// Why a runtime operation did not complete.