Compare commits
11 changed files with 63 additions and 262 deletions
|
|
@ -113,7 +113,6 @@ else a `TurnError`) drives the post-claude branch:
|
||||||
| `Err(PromptTooLong)` | `drive_turn` archived the session (the lib already compacted + retried and it still overflowed); requeue inflight so the message redelivers into a fresh session that fits — no status park |
|
| `Err(PromptTooLong)` | `drive_turn` archived the session (the lib already compacted + retried and it still overflowed); requeue inflight so the message redelivers into a fresh session that fits — no status park |
|
||||||
| `Err(RateLimited)` | sleep `HIVE_RATE_LIMIT_SLEEP_SECS` (default 300), requeue inflight, status back to `online` |
|
| `Err(RateLimited)` | sleep `HIVE_RATE_LIMIT_SLEEP_SECS` (default 300), requeue inflight, status back to `online` |
|
||||||
| `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(Failed(err))` | route `[system] \`<qualified-label>\` claude turn failed:\n<err>` to `<parent>` via `send_to_parent` |
|
| `Err(Failed(err))` | route `[system] \`<qualified-label>\` claude turn failed:\n<err>` to `<parent>` via `send_to_parent` |
|
||||||
|
|
||||||
After the outcome handler, the stats sink records a row and the
|
After the outcome handler, the stats sink records a row and the
|
||||||
|
|
|
||||||
|
|
@ -627,13 +627,6 @@ async fn handle_turn<S: Surface>(
|
||||||
tracing::warn!("prompt-too-long; session archived, requeueing message for a fresh turn");
|
tracing::warn!("prompt-too-long; session archived, requeueing message for a fresh turn");
|
||||||
S::requeue_inflight(socket).await;
|
S::requeue_inflight(socket).await;
|
||||||
}
|
}
|
||||||
if matches!(outcome, Err(turn::TurnError::SessionNotFound)) {
|
|
||||||
// "Shouldn't happen": resume missed and the lib's create self-heal
|
|
||||||
// didn't resolve it. Requeue rather than ack-and-drop so the wake
|
|
||||||
// message isn't silently lost; the next turn creates the session fresh.
|
|
||||||
tracing::warn!("session-not-found; requeueing message for a fresh turn");
|
|
||||||
S::requeue_inflight(socket).await;
|
|
||||||
}
|
|
||||||
if let Err(turn::TurnError::Failed(e)) = &outcome {
|
if let Err(turn::TurnError::Failed(e)) = &outcome {
|
||||||
S::send_to_parent(socket, format_turn_failure(e)).await;
|
S::send_to_parent(socket, format_turn_failure(e)).await;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -392,8 +392,8 @@ async fn format_notification(
|
||||||
fetch_json(client, subject_api_url, token).await
|
fetch_json(client, subject_api_url, token).await
|
||||||
};
|
};
|
||||||
|
|
||||||
// Forgejo's notification `subject.type` is "Pull" / "Issue", never
|
// Forgejo's notification `subject.type` is "Pull" / "Issue" (never
|
||||||
// "Pull Request".
|
// "Pull Request") — see the comment in `build_meta_suffix` below.
|
||||||
let is_pr = notif_type == "Pull";
|
let is_pr = notif_type == "Pull";
|
||||||
let meta_suffix = build_meta_suffix(subject.as_ref(), is_pr);
|
let meta_suffix = build_meta_suffix(subject.as_ref(), is_pr);
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -33,13 +33,6 @@ pub fn default_dir() -> PathBuf {
|
||||||
/// Rationale + the previous wholesale-wipe shape we replaced live in
|
/// Rationale + the previous wholesale-wipe shape we replaced live in
|
||||||
/// [`docs/web-ui/agent.md::Per-agent endpoints`](../../docs/web-ui/agent.md)
|
/// [`docs/web-ui/agent.md::Per-agent endpoints`](../../docs/web-ui/agent.md)
|
||||||
/// (the `/api/logout` bullet).
|
/// (the `/api/logout` bullet).
|
||||||
///
|
|
||||||
/// `.credentials.json` is the actual OAuth session; `mcp-needs-auth-cache.json`
|
|
||||||
/// is claude-code's MCP-auth cache and a weaker signal. Both are kept in the
|
|
||||||
/// set only because `/logout` deletes both, so the "either present ⇒ logged
|
|
||||||
/// in" check can never disagree with a logout. (If a future edit ever removes
|
|
||||||
/// `.credentials.json` without the cache — a state `/logout` doesn't produce —
|
|
||||||
/// keying purely on `.credentials.json` would be the stronger boot signal.)
|
|
||||||
pub const CRED_FILE_NAMES: &[&str] = &[".credentials.json", "mcp-needs-auth-cache.json"];
|
pub const CRED_FILE_NAMES: &[&str] = &[".credentials.json", "mcp-needs-auth-cache.json"];
|
||||||
|
|
||||||
/// Is `entry` a regular file whose name is one of [`CRED_FILE_NAMES`]?
|
/// Is `entry` a regular file whose name is one of [`CRED_FILE_NAMES`]?
|
||||||
|
|
|
||||||
|
|
@ -105,7 +105,6 @@ pub fn build_row(args: TurnRowArgs<'_>) -> TurnStatRow {
|
||||||
Err(TurnError::PromptTooLong) => ("prompt_too_long", None),
|
Err(TurnError::PromptTooLong) => ("prompt_too_long", None),
|
||||||
Err(TurnError::RateLimited) => ("rate_limited", None),
|
Err(TurnError::RateLimited) => ("rate_limited", None),
|
||||||
Err(TurnError::AuthFailed) => ("auth_failed", None),
|
Err(TurnError::AuthFailed) => ("auth_failed", None),
|
||||||
Err(TurnError::SessionNotFound) => ("session_not_found", None),
|
|
||||||
Err(TurnError::Failed(e)) => ("failed", Some(format!("{e:#}"))),
|
Err(TurnError::Failed(e)) => ("failed", Some(format!("{e:#}"))),
|
||||||
};
|
};
|
||||||
let wake_from = if wake_from.starts_with("bash-task-") {
|
let wake_from = if wake_from.starts_with("bash-task-") {
|
||||||
|
|
|
||||||
|
|
@ -14,15 +14,15 @@ use serde_json::Value;
|
||||||
use crate::events::{Bus, LiveEvent};
|
use crate::events::{Bus, LiveEvent};
|
||||||
use crate::mcp_config;
|
use crate::mcp_config;
|
||||||
|
|
||||||
// Hive-enforced claude settings ship at `/etc/claude-code/managed-settings.json`,
|
// Hive-enforced claude settings ship at `/etc/claude-code/managed-settings.json`
|
||||||
// which claude-code auto-discovers (precedence #1, read-only, un-overridable) —
|
// (wired in `nix/templates/harness-base.nix` from the `prompts/claude-settings.json`
|
||||||
// so the harness no longer passes `--settings`. We turn off claude's in-session
|
// asset). claude-code auto-discovers that managed path — precedence #1,
|
||||||
// auto-compaction and its cross-session auto-memory because hyperhive owns those
|
// read-only, un-overridable — so the harness no longer passes `--settings`.
|
||||||
// concerns (`/compact` on overflow, notes persistence under `/state`). How the
|
// We turn off claude's in-session auto-compaction and its cross-session
|
||||||
// file is wired (the nix asset) + the full rationale live in
|
// auto-memory because hyperhive owns those concerns (`/compact` on overflow,
|
||||||
// `docs/turn-loop/claude-invocation.md`. Unknown keys are silently ignored by
|
// notes persistence under `/state`). Unknown keys are silently ignored by
|
||||||
// claude-code; if a key gets renamed we'll spot it because the corresponding
|
// claude-code; if a key gets renamed we'll spot it because the
|
||||||
// behavior will start firing mid-turn again.
|
// corresponding behavior will start firing mid-turn again.
|
||||||
//
|
//
|
||||||
// The subprocess mechanics — spawning `claude --print`, streaming +
|
// The subprocess mechanics — spawning `claude --print`, streaming +
|
||||||
// classifying stream-json, session lookup/archive — live in the generic
|
// classifying stream-json, session lookup/archive — live in the generic
|
||||||
|
|
@ -170,12 +170,6 @@ pub enum TurnError {
|
||||||
/// into `needs_login_idle` and stop driving turns until the
|
/// into `needs_login_idle` and stop driving turns until the
|
||||||
/// operator re-auths via the per-agent web UI.
|
/// operator re-auths via the per-agent web UI.
|
||||||
AuthFailed,
|
AuthFailed,
|
||||||
/// `--resume <title>` missed AND the lib's create self-heal also failed to
|
|
||||||
/// resolve the session — "shouldn't happen" (a resume-miss is normally
|
|
||||||
/// self-healed inside [`InfiniteSession::attempt`]). Rather than ack + drop
|
|
||||||
/// the wake message, the serve loop requeues it so the next turn retries;
|
|
||||||
/// no status park.
|
|
||||||
SessionNotFound,
|
|
||||||
/// 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
|
||||||
/// parent (`send_to_parent`).
|
/// parent (`send_to_parent`).
|
||||||
Failed(anyhow::Error),
|
Failed(anyhow::Error),
|
||||||
|
|
@ -348,11 +342,7 @@ pub async fn drive_turn(
|
||||||
// run it now that the turn is done, so it works mid-turn rather than only
|
// run it now that the turn is done, so it works mid-turn rather than only
|
||||||
// when the agent is idle. Only on a healthy turn — no point spawning a
|
// when the agent is idle. Only on a healthy turn — no point spawning a
|
||||||
// compaction after a rate-limited / auth-failed / crashed one.
|
// compaction after a rate-limited / auth-failed / crashed one.
|
||||||
// `is_ok()` first: `take_compact()` clears the flag, so it must only fire
|
if bus.take_compact() && outcome.is_ok() {
|
||||||
// when the compaction will actually run. On an unhealthy turn
|
|
||||||
// (rate-limited / auth-failed / failed) the flag is left set for the next
|
|
||||||
// turn or the idle `run_pending_compact` to service — not silently eaten.
|
|
||||||
if outcome.is_ok() && bus.take_compact() {
|
|
||||||
bus.emit(LiveEvent::Note {
|
bus.emit(LiveEvent::Note {
|
||||||
text: "operator: /compact — running at turn end".into(),
|
text: "operator: /compact — running at turn end".into(),
|
||||||
});
|
});
|
||||||
|
|
@ -436,13 +426,6 @@ pub fn emit_turn_end(bus: &Bus, outcome: &TurnOutcome) {
|
||||||
});
|
});
|
||||||
tracing::warn!("turn auth-failed (401)");
|
tracing::warn!("turn auth-failed (401)");
|
||||||
}
|
}
|
||||||
Err(TurnError::SessionNotFound) => {
|
|
||||||
bus.emit(LiveEvent::TurnEnd {
|
|
||||||
ok: false,
|
|
||||||
note: Some("session resume + create both missed — requeueing".into()),
|
|
||||||
});
|
|
||||||
tracing::warn!("turn session-not-found; requeueing message");
|
|
||||||
}
|
|
||||||
Err(TurnError::Failed(e)) => {
|
Err(TurnError::Failed(e)) => {
|
||||||
let note = format!("{e:#}");
|
let note = format!("{e:#}");
|
||||||
bus.emit(LiveEvent::TurnEnd {
|
bus.emit(LiveEvent::TurnEnd {
|
||||||
|
|
@ -554,7 +537,9 @@ fn error_to_turn(err: hive_claude::Error) -> TurnOutcome {
|
||||||
Error::PromptTooLong => Err(TurnError::PromptTooLong),
|
Error::PromptTooLong => Err(TurnError::PromptTooLong),
|
||||||
Error::RateLimited => Err(TurnError::RateLimited),
|
Error::RateLimited => Err(TurnError::RateLimited),
|
||||||
Error::AuthFailed => Err(TurnError::AuthFailed),
|
Error::AuthFailed => Err(TurnError::AuthFailed),
|
||||||
Error::SessionNotFound => Err(TurnError::SessionNotFound),
|
// A resume-miss the lib couldn't self-heal is benign — treat it as a
|
||||||
|
// clean (non-compacted) turn; the next turn creates the session fresh.
|
||||||
|
Error::SessionNotFound => Ok(false),
|
||||||
other => Err(TurnError::Failed(other.into())),
|
other => Err(TurnError::Failed(other.into())),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -606,20 +591,10 @@ impl Sink for BusSink<'_> {
|
||||||
/// watermark). Skips a degenerate turn that parsed nothing so it doesn't reset
|
/// watermark). Skips a degenerate turn that parsed nothing so it doesn't reset
|
||||||
/// the badges to zero.
|
/// the badges to zero.
|
||||||
fn apply_telemetry(bus: &Bus, telemetry: &hive_claude::Telemetry) {
|
fn apply_telemetry(bus: &Bus, telemetry: &hive_claude::Telemetry) {
|
||||||
// On a degenerate turn that emitted a `result` but no `assistant` event,
|
if telemetry.context.context_tokens() == 0 && telemetry.cost.context_tokens() == 0 {
|
||||||
// the per-inference `context` stays zero while `cost` (cumulative) is not.
|
|
||||||
// Fall back to `cost` as the ctx proxy so the ctx badge + auto-reset
|
|
||||||
// watermark don't go stale-to-zero. Only a turn that parsed nothing at all
|
|
||||||
// (both zero) is skipped.
|
|
||||||
let ctx = if telemetry.context.context_tokens() == 0 {
|
|
||||||
telemetry.cost
|
|
||||||
} else {
|
|
||||||
telemetry.context
|
|
||||||
};
|
|
||||||
if ctx.context_tokens() == 0 {
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
bus.record_turn_usage(ctx, telemetry.cost);
|
bus.record_turn_usage(telemetry.context, telemetry.cost);
|
||||||
bus.set_resolved_model(telemetry.model.clone());
|
bus.set_resolved_model(telemetry.model.clone());
|
||||||
if let Some(window) = telemetry.context_window {
|
if let Some(window) = telemetry.context_window {
|
||||||
bus.set_api_context_window(window);
|
bus.set_api_context_window(window);
|
||||||
|
|
|
||||||
|
|
@ -11,6 +11,3 @@ serde = { workspace = true }
|
||||||
serde_json.workspace = true
|
serde_json.workspace = true
|
||||||
thiserror.workspace = true
|
thiserror.workspace = true
|
||||||
tokio.workspace = true
|
tokio.workspace = true
|
||||||
|
|
||||||
[dev-dependencies]
|
|
||||||
tempfile = "3"
|
|
||||||
|
|
|
||||||
|
|
@ -50,42 +50,10 @@ pub(crate) struct Sentinels {
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Sentinels {
|
impl Sentinels {
|
||||||
/// Scan a CLI-authored line — stderr, or a non-JSON stdout line — for every
|
/// Scan a raw line (stdout or stderr) for the always-on markers:
|
||||||
/// marker. These bytes are always claude-code's own output, never model
|
/// prompt-too-long, auth-failed, session-not-found. Rate-limit is handled
|
||||||
/// conversation, so all markers are trusted.
|
/// separately because on stdout it must only fire on JSON `error` events.
|
||||||
pub(crate) fn scan_cli_line(&self, line: &str) {
|
pub(crate) fn scan_line(&self, line: &str) {
|
||||||
self.scan_failure_markers(line);
|
|
||||||
if RATE_LIMIT_MARKERS.iter().any(|m| line.contains(m)) {
|
|
||||||
self.rate_limited.store(true, Ordering::Relaxed);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Scan a parsed stdout JSON event. **Skips model-authored `assistant` /
|
|
||||||
/// `user` message events**, whose serialized content can quote any marker
|
|
||||||
/// verbatim (an agent discussing this very code, say) — a false positive
|
|
||||||
/// that would otherwise trip a needless compact/retry or a spurious
|
|
||||||
/// auth/session error. Every real signal here is emitted *instead of* a
|
|
||||||
/// model turn (the API rejected the prompt, the auth failed, or `--resume`
|
|
||||||
/// missed before any inference), so it can only appear on a control event
|
|
||||||
/// (`error` / `result` / `system`) or as raw non-JSON text — never inside
|
|
||||||
/// an assistant/user message. This holds whatever exact shape claude-code
|
|
||||||
/// uses for the message, so the gate can't suppress a genuine signal.
|
|
||||||
pub(crate) fn scan_stdout_json(&self, event: &serde_json::Value, raw: &str) {
|
|
||||||
let ty = event.get("type").and_then(|t| t.as_str());
|
|
||||||
if matches!(ty, Some("assistant" | "user")) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
self.scan_failure_markers(raw);
|
|
||||||
// Rate-limit stays scoped to `error` events (unchanged): the only
|
|
||||||
// control event that carries a rate-limit marker.
|
|
||||||
if ty == Some("error") && RATE_LIMIT_MARKERS.iter().any(|m| raw.contains(m)) {
|
|
||||||
self.rate_limited.store(true, Ordering::Relaxed);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The prompt-too-long / auth-failed / session-not-found markers. Callers
|
|
||||||
/// gate *where* this runs (see `scan_cli_line` / `scan_stdout_json`).
|
|
||||||
fn scan_failure_markers(&self, line: &str) {
|
|
||||||
if line.contains(PROMPT_TOO_LONG_MARKER) {
|
if line.contains(PROMPT_TOO_LONG_MARKER) {
|
||||||
self.prompt_too_long.store(true, Ordering::Relaxed);
|
self.prompt_too_long.store(true, Ordering::Relaxed);
|
||||||
}
|
}
|
||||||
|
|
@ -97,6 +65,26 @@ impl Sentinels {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Trust a rate-limit hit only on a JSON `error` event (so a model
|
||||||
|
/// *discussing* a rate limit in prose can't trigger it). The `type` gate
|
||||||
|
/// needs the parsed `event`; the marker match runs on `raw`, the original
|
||||||
|
/// line — the same bytes, so we don't re-serialize the value.
|
||||||
|
pub(crate) fn scan_stdout_json(&self, event: &serde_json::Value, raw: &str) {
|
||||||
|
if event.get("type").and_then(|t| t.as_str()) == Some("error")
|
||||||
|
&& RATE_LIMIT_MARKERS.iter().any(|m| raw.contains(m))
|
||||||
|
{
|
||||||
|
self.rate_limited.store(true, Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Trust a rate-limit hit on raw text (non-JSON stdout, or any stderr) —
|
||||||
|
/// these are CLI messages, not conversation content.
|
||||||
|
pub(crate) fn scan_rate_limit_text(&self, line: &str) {
|
||||||
|
if RATE_LIMIT_MARKERS.iter().any(|m| line.contains(m)) {
|
||||||
|
self.rate_limited.store(true, Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// The recognized-sentinel error, if any fired — `None` means no sentinel
|
/// The recognized-sentinel error, if any fired — `None` means no sentinel
|
||||||
/// (so the run either completed or failed hard on its exit code). The
|
/// (so the run either completed or failed hard on its exit code). The
|
||||||
/// sentinels keep a fixed priority (too-long > rate > auth); a
|
/// sentinels keep a fixed priority (too-long > rate > auth); a
|
||||||
|
|
@ -117,51 +105,3 @@ impl Sentinels {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
use crate::Error;
|
|
||||||
|
|
||||||
fn json(raw: &str) -> serde_json::Value {
|
|
||||||
serde_json::from_str(raw).unwrap()
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn assistant_content_quoting_marker_is_ignored() {
|
|
||||||
// An agent discussing this code emits the marker verbatim in an
|
|
||||||
// assistant message — must NOT trip a sentinel.
|
|
||||||
let s = Sentinels::default();
|
|
||||||
let raw = r#"{"type":"assistant","message":{"content":[{"type":"text","text":"the CLI prints Prompt is too long on overflow"}]}}"#;
|
|
||||||
s.scan_stdout_json(&json(raw), raw);
|
|
||||||
assert!(s.soft_error().is_none());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn control_event_marker_is_detected() {
|
|
||||||
let s = Sentinels::default();
|
|
||||||
let raw = r#"{"type":"result","subtype":"error","error":"Prompt is too long"}"#;
|
|
||||||
s.scan_stdout_json(&json(raw), raw);
|
|
||||||
assert!(matches!(s.soft_error(), Some(Error::PromptTooLong)));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn raw_non_json_marker_is_detected() {
|
|
||||||
let s = Sentinels::default();
|
|
||||||
s.scan_cli_line("API Error: Prompt is too long");
|
|
||||||
assert!(matches!(s.soft_error(), Some(Error::PromptTooLong)));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn rate_limit_still_only_on_error_event() {
|
|
||||||
// A non-error control event mentioning the marker must not trip it.
|
|
||||||
let s = Sentinels::default();
|
|
||||||
let raw = r#"{"type":"result","summary":"we hit a rate_limit_error earlier"}"#;
|
|
||||||
s.scan_stdout_json(&json(raw), raw);
|
|
||||||
assert!(s.soft_error().is_none());
|
|
||||||
// A genuine error event does.
|
|
||||||
let raw2 = r#"{"type":"error","error":{"type":"rate_limit_error"}}"#;
|
|
||||||
s.scan_stdout_json(&json(raw2), raw2);
|
|
||||||
assert!(matches!(s.soft_error(), Some(Error::RateLimited)));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
|
||||||
|
|
@ -152,14 +152,12 @@ fn build_command(program: &str, config: &Config, attach: &Attach) -> Command {
|
||||||
async fn pump_stdout(stdout: ChildStdout, sink: &impl Sink, sentinels: &Sentinels) {
|
async fn pump_stdout(stdout: ChildStdout, sink: &impl Sink, sentinels: &Sentinels) {
|
||||||
let mut lines = BufReader::new(stdout).lines();
|
let mut lines = BufReader::new(stdout).lines();
|
||||||
while let Ok(Some(line)) = lines.next_line().await {
|
while let Ok(Some(line)) = lines.next_line().await {
|
||||||
|
sentinels.scan_line(&line);
|
||||||
if let Ok(event) = serde_json::from_str::<serde_json::Value>(&line) {
|
if let Ok(event) = serde_json::from_str::<serde_json::Value>(&line) {
|
||||||
// JSON stdout: classify with the model-content gate so an
|
|
||||||
// `assistant`/`user` message quoting a marker can't trip it.
|
|
||||||
sentinels.scan_stdout_json(&event, &line);
|
sentinels.scan_stdout_json(&event, &line);
|
||||||
sink.on_event(&event);
|
sink.on_event(&event);
|
||||||
} else {
|
} else {
|
||||||
// Non-JSON stdout is CLI text, not conversation — trust all markers.
|
sentinels.scan_rate_limit_text(&line);
|
||||||
sentinels.scan_cli_line(&line);
|
|
||||||
sink.on_stdout_line(&line);
|
sink.on_stdout_line(&line);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -172,8 +170,8 @@ async fn pump_stderr(stderr: ChildStderr, sink: &impl Sink, sentinels: &Sentinel
|
||||||
let mut lines = BufReader::new(stderr).lines();
|
let mut lines = BufReader::new(stderr).lines();
|
||||||
let mut tail: VecDeque<String> = VecDeque::with_capacity(STDERR_TAIL_LINES);
|
let mut tail: VecDeque<String> = VecDeque::with_capacity(STDERR_TAIL_LINES);
|
||||||
while let Ok(Some(line)) = lines.next_line().await {
|
while let Ok(Some(line)) = lines.next_line().await {
|
||||||
// stderr is always CLI output — trust all markers.
|
sentinels.scan_line(&line);
|
||||||
sentinels.scan_cli_line(&line);
|
sentinels.scan_rate_limit_text(&line);
|
||||||
sink.on_stderr_line(&line);
|
sink.on_stderr_line(&line);
|
||||||
if tail.len() >= STDERR_TAIL_LINES {
|
if tail.len() >= STDERR_TAIL_LINES {
|
||||||
tail.pop_front();
|
tail.pop_front();
|
||||||
|
|
|
||||||
|
|
@ -66,25 +66,19 @@ impl<P: CompactionPolicy> InfiniteSession<P> {
|
||||||
/// [`Error::PromptTooLong`]; rate-limit / auth / hard failures propagate
|
/// [`Error::PromptTooLong`]; rate-limit / auth / hard failures propagate
|
||||||
/// unchanged for the caller to handle.
|
/// unchanged for the caller to handle.
|
||||||
pub async fn run(&self, config: &Config, prompt: &str, sink: &impl Sink) -> Result<Progress> {
|
pub async fn run(&self, config: &Config, prompt: &str, sink: &impl Sink) -> Result<Progress> {
|
||||||
// Resolve resume-vs-create once, up front, so `created` reflects the
|
|
||||||
// session state at the START of the run: the reactive-retry path can't
|
|
||||||
// read it from the retry (the session exists by then).
|
|
||||||
let existed = self.store.find_by_title(&self.name).is_some();
|
|
||||||
let meter = TelemetrySink::new(sink);
|
let meter = TelemetrySink::new(sink);
|
||||||
let created = match self.attempt(config, prompt, &meter, existed).await {
|
let created = match self.attempt(config, prompt, &meter).await {
|
||||||
Ok(created) => created,
|
Ok(created) => created,
|
||||||
Err(Error::PromptTooLong) => {
|
Err(Error::PromptTooLong) => {
|
||||||
// The session is already past the window — no turn can run on
|
// The session is already past the window — no turn can run on
|
||||||
// it and the detail is gone (no checkpoint possible). Compact,
|
// it and the detail is gone (no checkpoint possible). Compact,
|
||||||
// then retry the same prompt once; the retry is the answering
|
// then retry the same prompt once; the retry is the answering
|
||||||
// turn, so its telemetry is what we report. The failed attempt
|
// turn, so its telemetry is what we report.
|
||||||
// created/resumed the session, so the retry resumes it, and
|
|
||||||
// `created` still reflects the pre-run state.
|
|
||||||
self.compact(config, sink).await?;
|
self.compact(config, sink).await?;
|
||||||
let retry = TelemetrySink::new(sink);
|
let retry = TelemetrySink::new(sink);
|
||||||
self.attempt(config, prompt, &retry, true).await?;
|
let created = self.attempt(config, prompt, &retry).await?;
|
||||||
return Ok(Progress {
|
return Ok(Progress {
|
||||||
created: !existed,
|
created,
|
||||||
compacted: true,
|
compacted: true,
|
||||||
telemetry: retry.snapshot(),
|
telemetry: retry.snapshot(),
|
||||||
});
|
});
|
||||||
|
|
@ -97,15 +91,12 @@ impl<P: CompactionPolicy> InfiniteSession<P> {
|
||||||
// says it's due, checkpoint (best-effort) then compact.
|
// says it's due, checkpoint (best-effort) then compact.
|
||||||
if self.policy.should_compact(telemetry.usage()) {
|
if self.policy.should_compact(telemetry.usage()) {
|
||||||
if let Some(checkpoint) = self.policy.checkpoint_prompt() {
|
if let Some(checkpoint) = self.policy.checkpoint_prompt() {
|
||||||
let _ = self.attempt(config, checkpoint, sink, true).await;
|
let _ = self.attempt(config, checkpoint, sink).await;
|
||||||
}
|
}
|
||||||
// Only claim a compaction if it actually succeeded — the flag feeds
|
let _ = self.compact(config, sink).await;
|
||||||
// stats + the auto-reset watermark, so a failed best-effort
|
|
||||||
// `/compact` must not report that the context shrank.
|
|
||||||
let compacted = self.compact(config, sink).await.is_ok();
|
|
||||||
return Ok(Progress {
|
return Ok(Progress {
|
||||||
created,
|
created,
|
||||||
compacted,
|
compacted: true,
|
||||||
telemetry,
|
telemetry,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
@ -130,29 +121,22 @@ impl<P: CompactionPolicy> InfiniteSession<P> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// One resume-or-create turn. `existed` is the caller's up-front
|
/// One resume-or-create turn. Uses the store to pick resume vs. create up
|
||||||
/// resume-vs-create decision (whether the titled session was on disk before
|
/// front (avoiding a wasted resume-miss spawn), and still self-heals if the
|
||||||
/// the run) — passing it in rather than re-checking keeps `created`
|
/// backing file vanished between the check and the run. Returns whether a
|
||||||
/// reporting consistent across the reactive-retry path. Still self-heals if
|
/// fresh session was created.
|
||||||
/// the backing file vanished between the check and the run. Returns whether
|
async fn attempt(&self, config: &Config, prompt: &str, sink: &impl Sink) -> Result<bool> {
|
||||||
/// a fresh session was created.
|
let exists = self.store.find_by_title(&self.name).is_some();
|
||||||
async fn attempt(
|
let attach = if exists {
|
||||||
&self,
|
|
||||||
config: &Config,
|
|
||||||
prompt: &str,
|
|
||||||
sink: &impl Sink,
|
|
||||||
existed: bool,
|
|
||||||
) -> Result<bool> {
|
|
||||||
let attach = if existed {
|
|
||||||
Attach::Resume(self.name.clone())
|
Attach::Resume(self.name.clone())
|
||||||
} else {
|
} else {
|
||||||
Attach::Create(self.name.clone())
|
Attach::Create(self.name.clone())
|
||||||
};
|
};
|
||||||
match Claude::run(config, &attach, prompt, sink).await {
|
match Claude::run(config, &attach, prompt, sink).await {
|
||||||
Ok(()) => Ok(!existed),
|
Ok(()) => Ok(!exists),
|
||||||
// We thought it existed but the resume missed (raced an archive) —
|
// We thought it existed but the resume missed (raced an archive) —
|
||||||
// self-heal by creating.
|
// self-heal by creating.
|
||||||
Err(Error::SessionNotFound) if existed => {
|
Err(Error::SessionNotFound) if exists => {
|
||||||
Claude::run(config, &Attach::Create(self.name.clone()), prompt, sink).await?;
|
Claude::run(config, &Attach::Create(self.name.clone()), prompt, sink).await?;
|
||||||
Ok(true)
|
Ok(true)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -50,6 +50,7 @@ impl SessionStore {
|
||||||
/// (including anything already archived to `*.jsonl.archived`) are skipped.
|
/// (including anything already archived to `*.jsonl.archived`) are skipped.
|
||||||
#[must_use]
|
#[must_use]
|
||||||
pub fn find_by_title(&self, title: &str) -> Option<PathBuf> {
|
pub fn find_by_title(&self, title: &str) -> Option<PathBuf> {
|
||||||
|
let marker = format!("\"customTitle\":\"{title}\"");
|
||||||
for entry in std::fs::read_dir(self.project_dir()).ok()?.flatten() {
|
for entry in std::fs::read_dir(self.project_dir()).ok()?.flatten() {
|
||||||
let path = entry.path();
|
let path = entry.path();
|
||||||
if path.extension().and_then(|e| e.to_str()) != Some("jsonl") {
|
if path.extension().and_then(|e| e.to_str()) != Some("jsonl") {
|
||||||
|
|
@ -61,7 +62,7 @@ impl SessionStore {
|
||||||
if std::io::BufReader::new(file)
|
if std::io::BufReader::new(file)
|
||||||
.lines()
|
.lines()
|
||||||
.map_while(std::result::Result::ok)
|
.map_while(std::result::Result::ok)
|
||||||
.any(|line| line_sets_title(&line, title))
|
.any(|line| line.contains(&marker))
|
||||||
{
|
{
|
||||||
return Some(path);
|
return Some(path);
|
||||||
}
|
}
|
||||||
|
|
@ -92,81 +93,3 @@ impl SessionStore {
|
||||||
Ok(Some(target))
|
Ok(Some(target))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// True if `line` is the session's `custom-title` event whose **top-level**
|
|
||||||
/// `customTitle` field equals `title`.
|
|
||||||
///
|
|
||||||
/// Parsing the line (rather than substring-matching the whole transcript) is
|
|
||||||
/// what makes this robust: a message that merely *quotes* the marker in its
|
|
||||||
/// content has no top-level `customTitle` key, so it can't cause a false match
|
|
||||||
/// on the wrong session file; the compact-vs-spaced JSON form is irrelevant;
|
|
||||||
/// and titles containing `"` / `\` are handled by the parser. Verified shape
|
|
||||||
/// (claude 2.1.x): `{"type":"custom-title","customTitle":"<title>",…}`.
|
|
||||||
///
|
|
||||||
/// The `contains` pre-check keeps the common case cheap — only the rare line
|
|
||||||
/// mentioning `customTitle` is parsed as JSON, not every transcript line.
|
|
||||||
fn line_sets_title(line: &str, title: &str) -> bool {
|
|
||||||
if !line.contains("customTitle") {
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
let Ok(value) = serde_json::from_str::<serde_json::Value>(line) else {
|
|
||||||
return false;
|
|
||||||
};
|
|
||||||
value.get("customTitle").and_then(|t| t.as_str()) == Some(title)
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
fn write(dir: &std::path::Path, name: &str, lines: &[&str]) {
|
|
||||||
std::fs::write(dir.join(name), lines.join("\n")).unwrap();
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn find_by_title_matches_only_the_custom_title_event() {
|
|
||||||
let home = tempfile::tempdir().unwrap();
|
|
||||||
let cwd = std::path::Path::new("/agents/iris/state");
|
|
||||||
let store = SessionStore::new(home.path(), cwd);
|
|
||||||
let proj = store.project_dir();
|
|
||||||
std::fs::create_dir_all(&proj).unwrap();
|
|
||||||
|
|
||||||
// The real titled session.
|
|
||||||
write(
|
|
||||||
&proj,
|
|
||||||
"real.jsonl",
|
|
||||||
&[
|
|
||||||
r#"{"type":"summary","summary":"x"}"#,
|
|
||||||
r#"{"type":"custom-title","customTitle":"iris","sessionId":"real"}"#,
|
|
||||||
],
|
|
||||||
);
|
|
||||||
// A DIFFERENT session whose transcript merely quotes the marker string
|
|
||||||
// in message content — must NOT match.
|
|
||||||
write(
|
|
||||||
&proj,
|
|
||||||
"other.jsonl",
|
|
||||||
&[
|
|
||||||
r#"{"type":"custom-title","customTitle":"someone-else","sessionId":"other"}"#,
|
|
||||||
r#"{"type":"assistant","message":{"content":[{"type":"text","text":"the file had \"customTitle\":\"iris\" in it"}]}}"#,
|
|
||||||
],
|
|
||||||
);
|
|
||||||
|
|
||||||
assert_eq!(store.find_by_title("iris"), Some(proj.join("real.jsonl")));
|
|
||||||
assert_eq!(store.find_by_title("nobody"), None);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn find_by_title_tolerates_spaced_json_and_escapes() {
|
|
||||||
let home = tempfile::tempdir().unwrap();
|
|
||||||
let store = SessionStore::new(home.path(), std::path::Path::new("/x"));
|
|
||||||
let proj = store.project_dir();
|
|
||||||
std::fs::create_dir_all(&proj).unwrap();
|
|
||||||
// Spaced JSON form + a title needing JSON escaping.
|
|
||||||
write(
|
|
||||||
&proj,
|
|
||||||
"s.jsonl",
|
|
||||||
&[r#"{ "type": "custom-title", "customTitle": "a\"b" }"#],
|
|
||||||
);
|
|
||||||
assert_eq!(store.find_by_title("a\"b"), Some(proj.join("s.jsonl")));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue