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 <cursoragent@cursor.com>
This commit is contained in:
john
2026-09-23 14:25:12 +08:00
parent 6af43341c8
commit c44c7f04be
+133 -7
View File
@@ -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<Instant>,
}
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<HarnessBreaker> = 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<String> {
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<AppState>) -> Router {
Router::new()
.route("/agent/harness_bootstrap", post(harness_bootstrap))