hive-runtime: exempt ACP compaction turns from the empty end_turn check
A checkpoint or `compact` command turn can end with a bare `end_turn` as a normal answer, since the agent does the work on its side. Treating it as `EmptyEndTurn` made `compact_session` count every such compaction as failed and archive the session. `prompt()` and `turn()` now take a `TurnKind`; only `run()`'s ordinary turns get the check. The test agent's `BLANK` env answers one prompt text with a bare `end_turn`, for the compact and checkpoint tests.
This commit is contained in:
parent
f9c6a56ab9
commit
af8e0fe681
2 changed files with 103 additions and 21 deletions
|
|
@ -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<P: CompactionPolicy> AcpRuntime<P> {
|
|||
live: &mut Live,
|
||||
config: &Config,
|
||||
prompt: &str,
|
||||
kind: TurnKind,
|
||||
sink: &impl Sink,
|
||||
mut cancelled: Pin<&mut Notified<'_>>,
|
||||
) -> Result<Progress> {
|
||||
|
|
@ -415,12 +430,10 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
|||
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<P: CompactionPolicy> AcpRuntime<P> {
|
|||
guard: &mut Option<Live>,
|
||||
config: &Config,
|
||||
prompt: &str,
|
||||
kind: TurnKind,
|
||||
sink: &impl Sink,
|
||||
) -> Result<Progress> {
|
||||
// Registered before the turn is marked in flight, so a cancel is
|
||||
|
|
@ -469,7 +483,7 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
|||
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<P: CompactionPolicy> AcpRuntime<P> {
|
|||
..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<P: CompactionPolicy> AcpRuntime<P> {
|
|||
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<P: CompactionPolicy> AcpRuntime<P> {
|
|||
impl<P: CompactionPolicy> Runtime for AcpRuntime<P> {
|
||||
async fn run(&self, config: &Config, prompt: &str, sink: &impl Sink) -> Result<Progress> {
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Reference in a new issue