diff --git a/hive-runtime/src/acp/mod.rs b/hive-runtime/src/acp/mod.rs index 36b8761f..94bbef47 100644 --- a/hive-runtime/src/acp/mod.rs +++ b/hive-runtime/src/acp/mod.rs @@ -1,5 +1,9 @@ //! The ACP backend: one long-lived agent process per runtime, spawned on the //! first turn, holding one durable session whose id is kept in a file. +//! +//! An ordinary turn that ends with `end_turn` after sending no event and +//! reporting no usage fails with [`AcpError::EmptyEndTurn`]. An agent that +//! never reports usage therefore surfaces every such turn as that error. mod rpc; mod stream; @@ -161,6 +165,16 @@ enum Stop { Idle, } +/// What a turn is for. +#[derive(Clone, Copy, PartialEq, Eq)] +enum TurnKind { + /// A prompt the runtime was asked to run. + Ordinary, + /// A checkpoint or `compact` command turn. A bare `end_turn` is a normal + /// answer here, so it is not an [`AcpError::EmptyEndTurn`]. + Compaction, +} + /// Turns on an ACP agent's durable session. /// /// From the [`Config`] it reads `cwd` (the agent's working directory), @@ -326,6 +340,7 @@ impl AcpRuntime

{ live: &mut Live, config: &Config, prompt: &str, + kind: TurnKind, sink: &impl Sink, mut cancelled: Pin<&mut Notified<'_>>, ) -> Result { @@ -415,12 +430,10 @@ impl AcpRuntime

{ Some((Stop::Operator, _)) => Some("cancelled"), None => response["stopReason"].as_str(), }; - let empty = stop_reason == Some("end_turn") - && !mapper.has_content() - && !mapper.has_usage_update() - && response.get("usage").is_none(); match stop_reason { - Some("end_turn") if empty => return Err(AcpError::EmptyEndTurn.into()), + Some("end_turn") if kind == TurnKind::Ordinary && mapper.sent_nothing(&response) => { + return Err(AcpError::EmptyEndTurn.into()); + } Some("end_turn") => {} reason => { let reason = reason.unwrap_or("none"); @@ -458,6 +471,7 @@ impl AcpRuntime

{ guard: &mut Option, config: &Config, prompt: &str, + kind: TurnKind, sink: &impl Sink, ) -> Result { // Registered before the turn is marked in flight, so a cancel is @@ -469,7 +483,7 @@ impl AcpRuntime

{ self.cancel.in_turn.store(true, Ordering::SeqCst); let _in_turn = InTurn(&self.cancel.in_turn); let live = self.running(guard, config).await?; - let result = self.turn(live, config, prompt, sink, cancelled).await; + let result = self.turn(live, config, prompt, kind, sink, cancelled).await; if let Err(Error::Acp( AcpError::Closed | AcpError::Exited { .. } @@ -520,7 +534,10 @@ impl AcpRuntime

{ ..config.clone() }; let command = format!("/{COMPACT_COMMAND}"); - let Err(e) = self.prompt(guard, &bounded, &command, sink).await else { + let Err(e) = self + .prompt(guard, &bounded, &command, TurnKind::Compaction, sink) + .await + else { return Ok(()); }; tracing::warn!(error = %e, "ACP compact command failed; starting a new session"); @@ -543,7 +560,10 @@ impl AcpRuntime

{ let Some(prompt) = self.policy.checkpoint_prompt() else { return; }; - if let Err(e) = self.prompt(guard, config, prompt, sink).await { + if let Err(e) = self + .prompt(guard, config, prompt, TurnKind::Compaction, sink) + .await + { tracing::warn!(error = %e, "ACP checkpoint turn failed"); } } @@ -552,7 +572,9 @@ impl AcpRuntime

{ impl Runtime for AcpRuntime

{ async fn run(&self, config: &Config, prompt: &str, sink: &impl Sink) -> Result { let mut guard = self.live.lock().await; - let mut progress = self.prompt(&mut guard, config, prompt, sink).await?; + let mut progress = self + .prompt(&mut guard, config, prompt, TurnKind::Ordinary, sink) + .await?; if self.policy.should_compact(progress.telemetry.usage()) { match self.compact_session(&mut guard, config, sink, true).await { Ok(()) => progress.compacted = true, @@ -736,7 +758,7 @@ mod tests { use hive_claude::{Config, NoopSink, PercentPolicy, Sink}; - use super::{AcpError, AcpRuntime}; + use super::{AcpError, AcpRuntime, TurnKind}; use crate::{AcpCommand, Error, Runtime}; /// An ACP agent in plain `sh`. It answers `initialize`, numbers its @@ -760,7 +782,8 @@ mod tests { /// With `COMMANDS` set, it advertises that one command on each session it /// creates or loads. With `USED` set, it reports that many of 1000 context /// tokens used on each prompt to session `s1`. With `STALL_COMPACT` set, it - /// answers `/compact` only when cancelled, as `silent` does. + /// answers `/compact` only when cancelled, as `silent` does. With `BLANK` + /// set, it answers a prompt whose text is exactly that as `blank` does. const AGENT: &str = r#" n=0 prompt= printf 'start\n' >> "$1.methods" @@ -773,6 +796,9 @@ advertise() { ended() { printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn","usage":{"inputTokens":1,"outputTokens":1}}}\n' "$1" } +bare() { + printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$1" +} while IFS= read -r line; do id=${line#*\"id\":}; id=${id%%[,\}]*} m=${line#*\"method\":\"}; m=${m%%\"*} @@ -794,6 +820,9 @@ while IFS= read -r line; do case $line in *'"text":"/compact"'*) [ -n "$STALL_COMPACT" ] && prompt=$id && continue ;; esac + case $line in + *"\"text\":\"$BLANK\""*) [ -n "$BLANK" ] && bare "$id" && continue ;; + esac if [ -e "$1.once" ]; then ended "$id" else @@ -802,8 +831,7 @@ while IFS= read -r line; do fail) printf '{"jsonrpc":"2.0","id":%s,"error":{"code":-32603,"message":"provider down"}}\n' "$id" ;; ok) ended "$id" ;; - blank) - printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$id" ;; + blank) bare "$id" ;; silent) prompt=$id ;; trickle) for _ in 1 2 3 4 5; do @@ -1126,6 +1154,64 @@ done assert_eq!(recorded(dir.path()), "s1"); } + #[tokio::test] + async fn a_compact_command_answered_with_a_bare_end_turn_keeps_the_session() { + let dir = tempfile::tempdir().unwrap(); + let env = [("COMMANDS", "compact"), ("BLANK", "/compact")]; + let runtime = agent(dir.path(), "ok", &env, policy(0)); + let config = config(dir.path()); + + runtime.run(&config, "one", &NoopSink).await.unwrap(); + runtime.compact(&config, &NoopSink).await.unwrap(); + let next = runtime.run(&config, "two", &NoopSink).await.unwrap(); + + assert!(!next.created); + assert_eq!( + prompts(dir.path()), + [ + sent("s1", "SYSTEM PROMPT\n\none"), + sent("s1", "/compact"), + sent("s1", "two"), + ] + ); + assert_eq!(count(dir.path(), "session/new"), 1); + assert_eq!(recorded(dir.path()), "s1"); + } + + #[tokio::test] + async fn a_checkpoint_turn_answered_with_a_bare_end_turn_succeeds() { + let dir = tempfile::tempdir().unwrap(); + let runtime = agent(dir.path(), "ok", &[("BLANK", "CHECKPOINT")], policy(0)); + let config = config(dir.path()); + runtime.run(&config, "one", &NoopSink).await.unwrap(); + let mut guard = runtime.live.lock().await; + + let checkpoint = runtime + .prompt( + &mut guard, + &config, + "CHECKPOINT", + TurnKind::Compaction, + &NoopSink, + ) + .await; + assert!(checkpoint.is_ok(), "{checkpoint:?}"); + // The same reply to an ordinary turn is the empty-turn error. + let ordinary = runtime + .prompt( + &mut guard, + &config, + "CHECKPOINT", + TurnKind::Ordinary, + &NoopSink, + ) + .await; + assert!( + matches!(ordinary, Err(Error::Acp(AcpError::EmptyEndTurn))), + "{ordinary:?}" + ); + } + #[tokio::test] async fn with_no_compact_command_notes_are_written_then_a_new_session_starts() { let dir = tempfile::tempdir().unwrap(); diff --git a/hive-runtime/src/acp/stream.rs b/hive-runtime/src/acp/stream.rs index 50876041..c053e838 100644 --- a/hive-runtime/src/acp/stream.rs +++ b/hive-runtime/src/acp/stream.rs @@ -89,14 +89,10 @@ impl StreamMapper { out } - /// Whether the turn has emitted any event. - pub(super) fn has_content(&self) -> bool { - self.emitted - } - - /// Whether the turn has had a `usage_update`. - pub(super) fn has_usage_update(&self) -> bool { - self.usage.is_some() + /// Whether the turn has emitted no event and reported no usage, neither + /// in a `usage_update` nor in the `session/prompt` `response`. + pub(super) fn sent_nothing(&self, response: &Value) -> bool { + !self.emitted && self.usage.is_none() && response.get("usage").is_none() } /// The turn's telemetry: context from the last `usage_update`, cost from