diff --git a/docs/turn-loop/README.md b/docs/turn-loop/README.md index 03fa27ac..15338e7e 100644 --- a/docs/turn-loop/README.md +++ b/docs/turn-loop/README.md @@ -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(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(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] \`\` claude turn failed:\n` to `operator` via `send_to_operator` | `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 `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` reports the result to `serve_loop` via `TurnControl { auth_failed }` — on auth failure the loop parks in `wait_for_login`; otherwise it loops diff --git a/docs/web-ui/agent.md b/docs/web-ui/agent.md index e00ebf89..de4dc5d6 100644 --- a/docs/web-ui/agent.md +++ b/docs/web-ui/agent.md @@ -244,8 +244,10 @@ Slash commands today: - `/help` — list commands locally. - `/clear` — wipe the local terminal view (server history kept). - `/cancel` — `POST /api/cancel` → host shellouts `pkill -INT - claude`, emits a Note. Also surfaces as a `■ cancel turn` - button in the state row while state=thinking. + claude`, emits a Note. On an ACP agent it sends the agent + `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 `Bus::request_compact()` flag and returns immediately. The harness runs the `/compact` at the next turn boundary (end of the in-flight diff --git a/hive-agent/src/events.rs b/hive-agent/src/events.rs index 2ad82a8b..ef48be69 100644 --- a/hive-agent/src/events.rs +++ b/hive-agent/src/events.rs @@ -10,7 +10,7 @@ use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering}; -use std::sync::{Arc, Mutex}; +use std::sync::{Arc, Mutex, OnceLock}; use chrono::Utc; use hive_claude::TokenUsage; @@ -331,6 +331,10 @@ pub struct Bus { /// set, the compact has already run and that flag has already been /// cleared by `take_compact`. post_compact_wake: Arc>>, + /// 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>, /// 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 /// 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)), compact_pending: Arc::new(Mutex::new(None)), post_compact_wake: Arc::new(Mutex::new(None)), + turn_canceller: Arc::default(), session_id: Arc::new(Mutex::new(None)), fresh_session: Arc::new(AtomicBool::new(false)), 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 } + /// 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: /// 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 diff --git a/hive-agent/src/main.rs b/hive-agent/src/main.rs index 3ed7743f..f65f884a 100644 --- a/hive-agent/src/main.rs +++ b/hive-agent/src/main.rs @@ -662,6 +662,9 @@ async fn serve_loop( // The durable agent session, built once and reused for every turn + // idle compaction below (it's effectively stateless). 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 // once each instead of twelve lines a minute while parked. let mut was_paused = false; @@ -1040,9 +1043,13 @@ async fn handle_turn_error_recovery( S::requeue_inflight(socket).await; bus.emit_status("online"); } - if matches!(outcome, Err(turn::TurnError::ApiStall)) { - // Idle watchdog killed claude on a suspected API stall. Park briefly to - // let the API recover, then requeue — same shape as the rate-limit path. + if matches!( + outcome, + 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(); bus.emit_status("api_stall"); bus.emit(LiveEvent::Note { diff --git a/hive-agent/src/serve_common.rs b/hive-agent/src/serve_common.rs index 4beb839d..d56e3ccc 100644 --- a/hive-agent/src/serve_common.rs +++ b/hive-agent/src/serve_common.rs @@ -100,6 +100,7 @@ pub fn build_row(args: TurnRowArgs<'_>) -> TurnStatRow { Err(TurnError::AuthFailed) => ("auth_failed", None), Err(TurnError::SessionNotFound) => ("session_not_found", None), Err(TurnError::ApiStall) => ("api_stall", None), + Err(TurnError::AgentStall(note)) => ("api_stall", Some(note.clone())), Err(TurnError::Failed(e)) => ("failed", Some(format!("{e:#}"))), }; TurnStatRow { diff --git a/hive-agent/src/turn.rs b/hive-agent/src/turn.rs index 8d05dfc2..6b007318 100644 --- a/hive-agent/src/turn.rs +++ b/hive-agent/src/turn.rs @@ -195,6 +195,12 @@ pub enum TurnError { /// The serve loop parks for `stall_sleep_secs()` and requeues, like the /// rate-limit path — NOT a crash. 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 /// operator (`send_to_operator`). Failed(anyhow::Error), @@ -565,6 +571,13 @@ pub fn emit_turn_end(bus: &Bus, outcome: &TurnOutcome) { }); 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)) => { let note = format!("{e:#}"); 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 -/// 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 { + use hive_runtime::{AcpError, Error}; 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())), } } @@ -809,7 +827,10 @@ fn archive_session(bus: &Bus, session: &AgentSession) { #[cfg(test)] 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] 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"), 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(_)) + )); + } } diff --git a/hive-agent/src/web_ui/actions.rs b/hive-agent/src/web_ui/actions.rs index c5e7d115..f4ca0cd2 100644 --- a/hive-agent/src/web_ui/actions.rs +++ b/hive-agent/src/web_ui/actions.rs @@ -52,6 +52,22 @@ pub(super) async fn post_send( } pub(super) async fn post_cancel_turn(State(state): State) -> 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 note = match out { SigintOutcome::Signalled => {