hive-agent: dashboard Cancel and idle stalls for ACP agents
`/api/cancel` stopped a turn only by SIGINTing a child process whose argv0 is `claude`, so on an ACP agent it reported "no claude process to interrupt" and the turn ran on. The serve loop now stores the session's canceller on the `Bus` when the runtime has one, and `/api/cancel` uses it: the agent is sent `session/cancel`, and the next wake prompt carries the interrupted hint, as for a signalled claude turn. Without a canceller (claude), the SIGINT path is untouched. An ACP turn stopped by the idle watchdog (`HIVE_TURN_IDLE_SECS`) becomes `TurnError::AgentStall(note)`, handled like `ApiStall`: park for `HIVE_STALL_SLEEP_SECS`, requeue, and record `api_stall`. The TurnEnd note is the runtime's message (how long the agent was silent, whether it had to be killed, and that a silently retried provider error such as HTTP 429 looks like this), not the claude-specific `ApiStall` text. A cancel the agent ignores is a `Failed` turn. Refs #4391
This commit is contained in:
parent
bcbeb8ac9c
commit
91e47732a6
7 changed files with 115 additions and 9 deletions
|
|
@ -126,6 +126,7 @@ else a `TurnError`) drives the post-claude branch:
|
||||||
| `Err(AuthFailed)` | emit `needs_login_idle` sentinel, requeue inflight, park in `wait_for_login` |
|
| `Err(AuthFailed)` | emit `needs_login_idle` sentinel, requeue inflight, park in `wait_for_login` |
|
||||||
| `Err(SessionNotFound)` | resume + create self-heal both missed ("shouldn't happen"); requeue inflight so the next turn creates fresh — no status park, message not dropped |
|
| `Err(SessionNotFound)` | resume + create self-heal both missed ("shouldn't happen"); requeue inflight so the next turn creates fresh — no status park, message not dropped |
|
||||||
| `Err(ApiStall)` | idle watchdog killed claude after `HIVE_TURN_IDLE_SECS` (default 600) of output silence; sleep `HIVE_STALL_SLEEP_SECS` (default 60), requeue inflight, status back to `online` |
|
| `Err(ApiStall)` | idle watchdog killed claude after `HIVE_TURN_IDLE_SECS` (default 600) of output silence; sleep `HIVE_STALL_SLEEP_SECS` (default 60), requeue inflight, status back to `online` |
|
||||||
|
| `Err(AgentStall(note))` | the same, for an ACP agent: after `HIVE_TURN_IDLE_SECS` with no `session/update` the harness sent the agent `session/cancel`, and killed it if it ignored that; `note` says which |
|
||||||
| `Err(Failed(err))` | route `[system] \`<qualified-label>\` claude turn failed:\n<err>` to `operator` via `send_to_operator` |
|
| `Err(Failed(err))` | route `[system] \`<qualified-label>\` claude turn failed:\n<err>` to `operator` via `send_to_operator` |
|
||||||
|
|
||||||
`ApiStall` catches an Anthropic API stall — a multi-retry connection storm where
|
`ApiStall` catches an Anthropic API stall — a multi-retry connection storm where
|
||||||
|
|
@ -136,6 +137,10 @@ only complete output silence for the window trips it. The harness sets the
|
||||||
window from `HIVE_TURN_IDLE_SECS` (`0` disables) and maps the driver's
|
window from `HIVE_TURN_IDLE_SECS` (`0` disables) and maps the driver's
|
||||||
`Error::IdleTimeout` onto `TurnError::ApiStall`.
|
`Error::IdleTimeout` onto `TurnError::ApiStall`.
|
||||||
|
|
||||||
|
`hive-runtime` applies the same window to an ACP agent's turns. That's
|
||||||
|
also how a provider HTTP 429 ends an ACP turn when the agent retries it without
|
||||||
|
reporting it: the turn goes silent until the watchdog stops it.
|
||||||
|
|
||||||
After the outcome handler, the stats sink records a row. `handle_turn`
|
After the outcome handler, the stats sink records a row. `handle_turn`
|
||||||
reports the result to `serve_loop` via `TurnControl { auth_failed }` —
|
reports the result to `serve_loop` via `TurnControl { auth_failed }` —
|
||||||
on auth failure the loop parks in `wait_for_login`; otherwise it loops
|
on auth failure the loop parks in `wait_for_login`; otherwise it loops
|
||||||
|
|
|
||||||
|
|
@ -244,8 +244,10 @@ Slash commands today:
|
||||||
- `/help` — list commands locally.
|
- `/help` — list commands locally.
|
||||||
- `/clear` — wipe the local terminal view (server history kept).
|
- `/clear` — wipe the local terminal view (server history kept).
|
||||||
- `/cancel` — `POST /api/cancel` → host shellouts `pkill -INT
|
- `/cancel` — `POST /api/cancel` → host shellouts `pkill -INT
|
||||||
claude`, emits a Note. Also surfaces as a `■ cancel turn`
|
claude`, emits a Note. On an ACP agent it sends the agent
|
||||||
button in the state row while state=thinking.
|
`session/cancel` instead, and kills it if the turn hasn't ended 10s
|
||||||
|
later. Also surfaces as a `■ cancel turn` button in the state row
|
||||||
|
while state=thinking.
|
||||||
- `/compact` — `POST /api/compact` → sets the deferred
|
- `/compact` — `POST /api/compact` → sets the deferred
|
||||||
`Bus::request_compact()` flag and returns immediately. The harness
|
`Bus::request_compact()` flag and returns immediately. The harness
|
||||||
runs the `/compact` at the next turn boundary (end of the in-flight
|
runs the `/compact` at the next turn boundary (end of the in-flight
|
||||||
|
|
|
||||||
|
|
@ -10,7 +10,7 @@
|
||||||
|
|
||||||
use std::path::{Path, PathBuf};
|
use std::path::{Path, PathBuf};
|
||||||
use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering};
|
use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering};
|
||||||
use std::sync::{Arc, Mutex};
|
use std::sync::{Arc, Mutex, OnceLock};
|
||||||
|
|
||||||
use chrono::Utc;
|
use chrono::Utc;
|
||||||
use hive_claude::TokenUsage;
|
use hive_claude::TokenUsage;
|
||||||
|
|
@ -331,6 +331,10 @@ pub struct Bus {
|
||||||
/// set, the compact has already run and that flag has already been
|
/// set, the compact has already run and that flag has already been
|
||||||
/// cleared by `take_compact`.
|
/// cleared by `take_compact`.
|
||||||
post_compact_wake: Arc<Mutex<Option<String>>>,
|
post_compact_wake: Arc<Mutex<Option<String>>>,
|
||||||
|
/// Stops the turn in flight, for a runtime that offers it (ACP). Set once
|
||||||
|
/// by the serve loop when it builds the session; unset on claude, whose
|
||||||
|
/// turn `/api/cancel` stops by signalling the `claude` process.
|
||||||
|
turn_canceller: Arc<OnceLock<hive_runtime::Canceller>>,
|
||||||
/// Current fresh-claude-session id (FK to `sessions.id`). Set by the
|
/// Current fresh-claude-session id (FK to `sessions.id`). Set by the
|
||||||
/// bin loop after minting a session row on a fresh start; stamped onto
|
/// bin loop after minting a session row on a fresh start; stamped onto
|
||||||
/// every `turn_stats` row until the next fresh session. `None` before
|
/// every `turn_stats` row until the next fresh session. `None` before
|
||||||
|
|
@ -422,6 +426,7 @@ impl Bus {
|
||||||
session_reset_pending: Arc::new(AtomicBool::new(false)),
|
session_reset_pending: Arc::new(AtomicBool::new(false)),
|
||||||
compact_pending: Arc::new(Mutex::new(None)),
|
compact_pending: Arc::new(Mutex::new(None)),
|
||||||
post_compact_wake: Arc::new(Mutex::new(None)),
|
post_compact_wake: Arc::new(Mutex::new(None)),
|
||||||
|
turn_canceller: Arc::default(),
|
||||||
session_id: Arc::new(Mutex::new(None)),
|
session_id: Arc::new(Mutex::new(None)),
|
||||||
fresh_session: Arc::new(AtomicBool::new(false)),
|
fresh_session: Arc::new(AtomicBool::new(false)),
|
||||||
tool_calls: Arc::new(Mutex::new(std::collections::HashMap::new())),
|
tool_calls: Arc::new(Mutex::new(std::collections::HashMap::new())),
|
||||||
|
|
@ -445,6 +450,19 @@ impl Bus {
|
||||||
self.event_seq.fetch_add(1, Ordering::SeqCst) + 1
|
self.event_seq.fetch_add(1, Ordering::SeqCst) + 1
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Record how the session's turn in flight is stopped. Only the first call
|
||||||
|
/// takes effect; there is one session per harness.
|
||||||
|
pub fn set_turn_canceller(&self, canceller: hive_runtime::Canceller) {
|
||||||
|
let _ = self.turn_canceller.set(canceller);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// How to stop the turn in flight, if the runtime offers a way (see
|
||||||
|
/// [`Self::set_turn_canceller`]).
|
||||||
|
#[must_use]
|
||||||
|
pub fn turn_canceller(&self) -> Option<&hive_runtime::Canceller> {
|
||||||
|
self.turn_canceller.get()
|
||||||
|
}
|
||||||
|
|
||||||
/// Request a session reset (operator `POST /api/new-session`). Deferred:
|
/// Request a session reset (operator `POST /api/new-session`). Deferred:
|
||||||
/// the flag is consumed at the next turn boundary by `drive_turn`, which
|
/// the flag is consumed at the next turn boundary by `drive_turn`, which
|
||||||
/// archives the current session so no claude process is mid-write when the
|
/// archives the current session so no claude process is mid-write when the
|
||||||
|
|
|
||||||
|
|
@ -662,6 +662,9 @@ async fn serve_loop<S: Surface>(
|
||||||
// The durable agent session, built once and reused for every turn +
|
// The durable agent session, built once and reused for every turn +
|
||||||
// idle compaction below (it's effectively stateless).
|
// idle compaction below (it's effectively stateless).
|
||||||
let session = turn::make_session(&bus)?;
|
let session = turn::make_session(&bus)?;
|
||||||
|
if let Some(canceller) = hive_runtime::Runtime::canceller(&session) {
|
||||||
|
bus.set_turn_canceller(canceller);
|
||||||
|
}
|
||||||
// Tracks the last observed pause state so the transitions get logged
|
// Tracks the last observed pause state so the transitions get logged
|
||||||
// once each instead of twelve lines a minute while parked.
|
// once each instead of twelve lines a minute while parked.
|
||||||
let mut was_paused = false;
|
let mut was_paused = false;
|
||||||
|
|
@ -1040,9 +1043,13 @@ async fn handle_turn_error_recovery<S: Surface>(
|
||||||
S::requeue_inflight(socket).await;
|
S::requeue_inflight(socket).await;
|
||||||
bus.emit_status("online");
|
bus.emit_status("online");
|
||||||
}
|
}
|
||||||
if matches!(outcome, Err(turn::TurnError::ApiStall)) {
|
if matches!(
|
||||||
// Idle watchdog killed claude on a suspected API stall. Park briefly to
|
outcome,
|
||||||
// let the API recover, then requeue — same shape as the rate-limit path.
|
Err(turn::TurnError::ApiStall | turn::TurnError::AgentStall(_))
|
||||||
|
) {
|
||||||
|
// Idle watchdog stopped the turn on a suspected API stall. Park briefly
|
||||||
|
// to let the API recover, then requeue — same shape as the rate-limit
|
||||||
|
// path.
|
||||||
let secs = turn::stall_sleep_secs();
|
let secs = turn::stall_sleep_secs();
|
||||||
bus.emit_status("api_stall");
|
bus.emit_status("api_stall");
|
||||||
bus.emit(LiveEvent::Note {
|
bus.emit(LiveEvent::Note {
|
||||||
|
|
|
||||||
|
|
@ -100,6 +100,7 @@ pub fn build_row(args: TurnRowArgs<'_>) -> TurnStatRow {
|
||||||
Err(TurnError::AuthFailed) => ("auth_failed", None),
|
Err(TurnError::AuthFailed) => ("auth_failed", None),
|
||||||
Err(TurnError::SessionNotFound) => ("session_not_found", None),
|
Err(TurnError::SessionNotFound) => ("session_not_found", None),
|
||||||
Err(TurnError::ApiStall) => ("api_stall", None),
|
Err(TurnError::ApiStall) => ("api_stall", None),
|
||||||
|
Err(TurnError::AgentStall(note)) => ("api_stall", Some(note.clone())),
|
||||||
Err(TurnError::Failed(e)) => ("failed", Some(format!("{e:#}"))),
|
Err(TurnError::Failed(e)) => ("failed", Some(format!("{e:#}"))),
|
||||||
};
|
};
|
||||||
TurnStatRow {
|
TurnStatRow {
|
||||||
|
|
|
||||||
|
|
@ -195,6 +195,12 @@ pub enum TurnError {
|
||||||
/// The serve loop parks for `stall_sleep_secs()` and requeues, like the
|
/// The serve loop parks for `stall_sleep_secs()` and requeues, like the
|
||||||
/// rate-limit path — NOT a crash.
|
/// rate-limit path — NOT a crash.
|
||||||
ApiStall,
|
ApiStall,
|
||||||
|
/// The ACP backend's idle watchdog stopped a turn that sent nothing for
|
||||||
|
/// `turn_idle_secs()`: the agent was sent `session/cancel`, and killed if
|
||||||
|
/// it ignored that. Carries what happened. An ACP agent retrying a failing
|
||||||
|
/// provider (e.g. HTTP 429) without reporting it ends up here. Handled
|
||||||
|
/// like [`Self::ApiStall`].
|
||||||
|
AgentStall(String),
|
||||||
/// A hard failure with no recovery — the serve loop escalates it to the
|
/// A hard failure with no recovery — the serve loop escalates it to the
|
||||||
/// operator (`send_to_operator`).
|
/// operator (`send_to_operator`).
|
||||||
Failed(anyhow::Error),
|
Failed(anyhow::Error),
|
||||||
|
|
@ -565,6 +571,13 @@ pub fn emit_turn_end(bus: &Bus, outcome: &TurnOutcome) {
|
||||||
});
|
});
|
||||||
tracing::warn!("turn killed: API stall (idle watchdog)");
|
tracing::warn!("turn killed: API stall (idle watchdog)");
|
||||||
}
|
}
|
||||||
|
Err(TurnError::AgentStall(note)) => {
|
||||||
|
bus.emit(LiveEvent::TurnEnd {
|
||||||
|
ok: false,
|
||||||
|
note: Some(format!("{note} — parking + requeueing")),
|
||||||
|
});
|
||||||
|
tracing::warn!(note = %note, "ACP turn stopped by the idle watchdog");
|
||||||
|
}
|
||||||
Err(TurnError::Failed(e)) => {
|
Err(TurnError::Failed(e)) => {
|
||||||
let note = format!("{e:#}");
|
let note = format!("{e:#}");
|
||||||
bus.emit(LiveEvent::TurnEnd {
|
bus.emit(LiveEvent::TurnEnd {
|
||||||
|
|
@ -694,10 +707,15 @@ fn error_to_turn(err: hive_claude::Error) -> TurnOutcome {
|
||||||
}
|
}
|
||||||
|
|
||||||
/// [`error_to_turn`] for a runtime error: the claude backend's errors map as
|
/// [`error_to_turn`] for a runtime error: the claude backend's errors map as
|
||||||
/// above, every other runtime's as `Failed`.
|
/// above, the ACP idle watchdog's as `AgentStall`, everything else as
|
||||||
|
/// `Failed`.
|
||||||
fn runtime_error_to_turn(err: hive_runtime::Error) -> TurnOutcome {
|
fn runtime_error_to_turn(err: hive_runtime::Error) -> TurnOutcome {
|
||||||
|
use hive_runtime::{AcpError, Error};
|
||||||
match err {
|
match err {
|
||||||
hive_runtime::Error::Claude(err) => error_to_turn(err),
|
Error::Claude(err) => error_to_turn(err),
|
||||||
|
Error::Acp(err @ (AcpError::IdleTimeout { .. } | AcpError::IdleKilled { .. })) => {
|
||||||
|
Err(TurnError::AgentStall(err.to_string()))
|
||||||
|
}
|
||||||
other => Err(TurnError::Failed(other.into())),
|
other => Err(TurnError::Failed(other.into())),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -809,7 +827,10 @@ fn archive_session(bus: &Bus, session: &AgentSession) {
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{PermissionAsk, acp_permits, nonzero_or, parse_u64};
|
use super::{
|
||||||
|
PermissionAsk, TurnError, acp_permits, nonzero_or, parse_u64, runtime_error_to_turn,
|
||||||
|
};
|
||||||
|
use hive_runtime::{AcpError, Error};
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn an_absent_or_blank_value_is_not_a_number() {
|
fn an_absent_or_blank_value_is_not_a_number() {
|
||||||
|
|
@ -909,4 +930,40 @@ mod tests {
|
||||||
assert!(acp_permits(&ask("fetch"), true));
|
assert!(acp_permits(&ask("fetch"), true));
|
||||||
assert!(!acp_permits(&ask("fetch"), false));
|
assert!(!acp_permits(&ask("fetch"), false));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn a_stalled_acp_turn_parks_and_requeues_with_what_happened() {
|
||||||
|
for err in [
|
||||||
|
AcpError::IdleTimeout { silent_secs: 600 },
|
||||||
|
AcpError::IdleKilled {
|
||||||
|
silent_secs: 600,
|
||||||
|
grace_secs: 10,
|
||||||
|
},
|
||||||
|
] {
|
||||||
|
let expected = err.to_string();
|
||||||
|
match runtime_error_to_turn(Error::Acp(err)) {
|
||||||
|
Err(TurnError::AgentStall(note)) => assert_eq!(note, expected),
|
||||||
|
other => panic!("{expected}: {other:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn a_rate_limited_acp_turn_says_so_when_it_stalls() {
|
||||||
|
let Err(TurnError::AgentStall(note)) =
|
||||||
|
runtime_error_to_turn(Error::Acp(AcpError::IdleTimeout { silent_secs: 600 }))
|
||||||
|
else {
|
||||||
|
panic!("not a stall");
|
||||||
|
};
|
||||||
|
assert!(note.contains("600s"), "{note}");
|
||||||
|
assert!(note.contains("HTTP 429"), "{note}");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn a_cancel_the_agent_ignored_is_a_failure() {
|
||||||
|
assert!(matches!(
|
||||||
|
runtime_error_to_turn(Error::Acp(AcpError::CancelIgnored { grace_secs: 10 })),
|
||||||
|
Err(TurnError::Failed(_))
|
||||||
|
));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -52,6 +52,22 @@ pub(super) async fn post_send(
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(super) async fn post_cancel_turn(State(state): State<AppState>) -> Response {
|
pub(super) async fn post_cancel_turn(State(state): State<AppState>) -> Response {
|
||||||
|
if let Some(canceller) = state.bus.turn_canceller() {
|
||||||
|
let note = if canceller.cancel() {
|
||||||
|
// Same as a signalled claude turn: the next wake prompt says it
|
||||||
|
// was cut off.
|
||||||
|
state
|
||||||
|
.interrupted
|
||||||
|
.store(true, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
"operator: /cancel — sent session/cancel to the ACP agent"
|
||||||
|
} else {
|
||||||
|
"operator: /cancel — no ACP turn to cancel"
|
||||||
|
};
|
||||||
|
state
|
||||||
|
.bus
|
||||||
|
.emit(crate::events::LiveEvent::Note { text: note.into() });
|
||||||
|
return (axum::http::StatusCode::OK, "ok").into_response();
|
||||||
|
}
|
||||||
let out = super::sigint_claude().await;
|
let out = super::sigint_claude().await;
|
||||||
let note = match out {
|
let note = match out {
|
||||||
SigintOutcome::Signalled => {
|
SigintOutcome::Signalled => {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue