From c44c7f04be7ca7319c311b6eb864aa7fe082724e Mon Sep 17 00:00:00 2001 From: john Date: Wed, 23 Sep 2026 14:25:12 +0800 Subject: [PATCH] fix(tkmind): pause harness remember/recall after repeated failures Three consecutive harness failures open a 10 minute breaker so a broken harness no longer adds a failing subprocess to every turn; stderr is logged as a bounded tail. Co-authored-by: Cursor --- crates/goose/src/acp/tkmind/memory.rs | 140 ++++++++++++++++++++++++-- 1 file changed, 133 insertions(+), 7 deletions(-) diff --git a/crates/goose/src/acp/tkmind/memory.rs b/crates/goose/src/acp/tkmind/memory.rs index 39874b64c..fdea710f7 100644 --- a/crates/goose/src/acp/tkmind/memory.rs +++ b/crates/goose/src/acp/tkmind/memory.rs @@ -1,7 +1,7 @@ use std::{ path::{Path, PathBuf}, - sync::Arc, - time::Duration, + sync::{Arc, Mutex}, + time::{Duration, Instant}, }; use axum::{extract::State, routing::post, Json, Router}; @@ -15,6 +15,82 @@ use super::{errors::ErrorResponse, state::AppState}; const DEFAULT_QUERY: &str = "project architecture conventions current goals decisions risks and recent work"; const DEFAULT_MAX_CHARS: usize = 6_000; +const BREAKER_FAILURE_THRESHOLD: u32 = 3; +const BREAKER_COOLDOWN: Duration = Duration::from_secs(600); +const STDERR_LOG_TAIL_CHARS: usize = 600; + +struct HarnessBreaker { + consecutive_failures: u32, + open_until: Option, +} + +impl HarnessBreaker { + const fn new() -> Self { + Self { + consecutive_failures: 0, + open_until: None, + } + } + + fn allows(&mut self, now: Instant) -> bool { + match self.open_until { + Some(until) if now < until => false, + Some(_) => { + self.open_until = None; + true + } + None => true, + } + } + + /// Returns true when this failure trips the breaker open. + fn record(&mut self, success: bool, now: Instant) -> bool { + if success { + self.consecutive_failures = 0; + self.open_until = None; + return false; + } + self.consecutive_failures += 1; + if self.consecutive_failures >= BREAKER_FAILURE_THRESHOLD { + self.consecutive_failures = 0; + self.open_until = Some(now + BREAKER_COOLDOWN); + return true; + } + false + } +} + +static HARNESS_BREAKER: Mutex = Mutex::new(HarnessBreaker::new()); + +fn harness_allowed() -> bool { + HARNESS_BREAKER + .lock() + .map(|mut breaker| breaker.allows(Instant::now())) + .unwrap_or(true) +} + +fn record_harness_result(success: bool) { + let tripped = HARNESS_BREAKER + .lock() + .map(|mut breaker| breaker.record(success, Instant::now())) + .unwrap_or(false); + if tripped { + tracing::warn!( + cooldown_secs = BREAKER_COOLDOWN.as_secs(), + "Harness failing repeatedly; pausing harness calls" + ); + } +} + +fn stderr_tail(stderr: &[u8]) -> String { + let text = String::from_utf8_lossy(stderr); + let text = text.trim(); + let count = text.chars().count(); + if count <= STDERR_LOG_TAIL_CHARS { + return text.to_string(); + } + text.chars().skip(count - STDERR_LOG_TAIL_CHARS).collect() +} #[derive(serde::Deserialize)] #[serde(rename_all = "camelCase")] @@ -69,16 +145,25 @@ fn harness_binary() -> PathBuf { } async fn recall(repo: &str, query: &str, working_dir: &Path) -> Option { - let output = timeout( + if !harness_allowed() { + return None; + } + let result = timeout( Duration::from_secs(8), Command::new(harness_binary()) .args(["recall", "--repo", repo, "--query", query]) .current_dir(working_dir) .output(), ) - .await - .ok()? - .ok()?; + .await; + let output = match result { + Ok(Ok(output)) => output, + _ => { + record_harness_result(false); + return None; + } + }; + record_harness_result(output.status.success()); if !output.status.success() { return None; } @@ -209,6 +294,10 @@ async fn harness_remember( .map(str::trim) .filter(|value| !value.is_empty()) .unwrap_or("H5 conversation memory"); + if !harness_allowed() { + tracing::debug!(session_id = %request.session_id, "Harness paused; skipping remember"); + return Ok(Json(HarnessRememberResponse { remembered: false })); + } let remembered = match timeout( Duration::from_secs(8), Command::new(harness_binary()) @@ -233,7 +322,7 @@ async fn harness_remember( tracing::warn!( session_id = %request.session_id, status = ?output.status.code(), - stderr = %String::from_utf8_lossy(&output.stderr).trim(), + stderr = %stderr_tail(&output.stderr), "Harness remember failed" ); false @@ -247,9 +336,46 @@ async fn harness_remember( false } }; + record_harness_result(remembered); Ok(Json(HarnessRememberResponse { remembered })) } +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn breaker_opens_after_threshold_and_recovers_after_cooldown() { + let mut breaker = HarnessBreaker::new(); + let t0 = Instant::now(); + assert!(breaker.allows(t0)); + assert!(!breaker.record(false, t0)); + assert!(!breaker.record(false, t0)); + assert!(breaker.record(false, t0)); + assert!(!breaker.allows(t0 + Duration::from_secs(1))); + assert!(breaker.allows(t0 + BREAKER_COOLDOWN + Duration::from_secs(1))); + } + + #[test] + fn breaker_success_resets_failure_count() { + let mut breaker = HarnessBreaker::new(); + let t0 = Instant::now(); + breaker.record(false, t0); + breaker.record(false, t0); + breaker.record(true, t0); + assert!(!breaker.record(false, t0)); + assert!(breaker.allows(t0)); + } + + #[test] + fn stderr_tail_keeps_last_chars() { + let long = format!("{}错误结尾", "x".repeat(2_000)); + let tail = stderr_tail(long.as_bytes()); + assert_eq!(tail.chars().count(), STDERR_LOG_TAIL_CHARS); + assert!(tail.ends_with("错误结尾")); + } +} + pub fn routes(state: Arc) -> Router { Router::new() .route("/agent/harness_bootstrap", post(harness_bootstrap))