From 2e57093721930de43cba4320fb22854b3bee51ae Mon Sep 17 00:00:00 2001 From: meh Date: Wed, 29 Jul 2026 02:37:51 +0700 Subject: [PATCH 1/4] Add runtime capability probing, and name a rejected flag as version drift Every flag mapping in this crate was verified against one specific release, and those releases move. There was no way to ask whether the installed CLI is the one the mappings were checked against, so drift surfaced as a run failing halfway through with "unexpected argument", which names a flag rather than the cause. That is exactly how `codex exec resume` rejecting `--sandbox` was found. `Probe::run(agent)` reads `--version` and compares it against `Agent::verified_version`, reporting Verified / Newer / Older / Unrecognized with an `advisory()` sentence written for a person to read. Probing is not automatic: it spawns a process, and paying that per request to guard against an occasional upstream change is the wrong trade. A host probes at startup. Two changes on top of the idea as borrowed from unified-agent-api, which keys an embedded table by target triple and refuses unvalidated combinations: - Also classify the failure. `Error::FlagRejected` is separated from `Error::Failed` when a CLI refuses an argument, because the remedy differs: nothing about the request is wrong, the wrapper and the CLI disagree about what it accepts. Refusing up front only helps someone who probed; this helps whoever is reading the error. - No `semver` dependency. These CLIs print a dotted triple inside prose, none publishes pre-release metadata, and a crate to compare three integers is not worth the graph. `Version::find` scans for the triple, which handles "2.1.205 (Claude Code)", "codex-cli 0.145.0" and the trailing period in "GitHub Copilot CLI 1.0.75." alike. The live test that checks installed agents against their verified versions runs by default rather than behind --ignored, since `--version` costs no quota. It turns flag drift into a red build that names the version, instead of a mysterious failure later. --- src/agent.rs | 23 ++++ src/error.rs | 17 +++ src/lib.rs | 2 + src/probe.rs | 308 ++++++++++++++++++++++++++++++++++++++++++++++++++ src/run.rs | 66 +++++++++++ tests/live.rs | 32 +++++- 6 files changed, 447 insertions(+), 1 deletion(-) create mode 100644 src/probe.rs diff --git a/src/agent.rs b/src/agent.rs index ef3e6d1..f323486 100644 --- a/src/agent.rs +++ b/src/agent.rs @@ -255,6 +255,29 @@ impl Agent { } } + /// The release this crate's flag mappings were verified against. + /// + /// Every mapping in this module was checked by running these exact + /// versions, not by reading their documentation. [`crate::Probe`] compares + /// an installed CLI against this so drift is a question a host can ask up + /// front rather than something a failing run reveals. + #[must_use] + pub fn verified_version(self) -> crate::Version { + let (major, minor, patch) = match self { + // `claude --version` -> "2.1.205 (Claude Code)" + Agent::Claude => (2, 1, 205), + // `codex --version` -> "codex-cli 0.145.0" + Agent::Codex => (0, 145, 0), + // `copilot --version` -> "GitHub Copilot CLI 1.0.75." + Agent::Copilot => (1, 0, 75), + }; + crate::Version { + major, + minor, + patch, + } + } + /// The documented install command, surfaced by [`Error::NotInstalled`]. #[must_use] pub fn install_hint(self) -> &'static str { diff --git a/src/error.rs b/src/error.rs index b622de6..e93671f 100644 --- a/src/error.rs +++ b/src/error.rs @@ -124,6 +124,23 @@ pub enum Error { detail: String, }, + /// The CLI rejected an argument this crate passed it. + /// + /// Almost always a version mismatch: the flag was verified against the + /// release named in [`crate::Agent::verified_version`] and the installed + /// one differs. Separated from [`Error::Failed`] because the remedy is + /// different: nothing about the request is wrong, the wrapper and the CLI + /// disagree. Run [`crate::Probe`] to confirm. + #[error( + "`{bin}` rejected an argument, which usually means its version differs from the one these flags were verified against: {detail}" + )] + FlagRejected { + /// The binary that refused. + bin: String, + /// Its own complaint, unedited. + detail: String, + }, + /// The run was stopped by [`crate::Run::cancel`] or by dropping its handle. /// /// Not a fault: the caller asked for this. Distinguished from diff --git a/src/lib.rs b/src/lib.rs index 246620a..ccaa3a0 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -89,6 +89,7 @@ mod agent; mod error; mod event; mod outcome; +mod probe; mod proc; mod request; mod run; @@ -98,6 +99,7 @@ pub use agent::{Agent, Caps, EnvPolicy, Format, NETWORK_ENV, Permission, Session pub use error::{Error, Result}; pub use event::{Event, MAX_CAPTURE}; pub use outcome::{Outcome, RateLimit, Stop, Usage}; +pub use probe::{Probe, Version, VersionStatus}; pub use request::Request; pub use run::{Run, run, stream}; pub use session::{Phase, SessionRecord, SessionStore}; diff --git a/src/probe.rs b/src/probe.rs new file mode 100644 index 0000000..8336357 --- /dev/null +++ b/src/probe.rs @@ -0,0 +1,308 @@ +//! Asking a CLI what version it is, and whether this crate was built for it. +//! +//! Every flag mapping here was verified against a specific release. Those +//! releases move: `codex exec resume` accepts `--sandbox` in some versions and +//! rejects it outright in 0.145.0, and Copilot gained a headless session id and +//! an event stream that an older wrapper still models as absent. Without a +//! version check, that drift is discovered by a run failing halfway through with +//! "unexpected argument", which names a flag rather than the cause. +//! +//! This module makes the check available up front. It is **not** run +//! automatically: probing spawns a process, and paying that on every request to +//! guard against an occasional upstream change is the wrong trade. A host +//! should probe once at startup, or when a run fails with +//! [`crate::Error::FlagRejected`], and surface the result to whoever can act on +//! it. + +use std::fmt; + +use crate::agent::Agent; +use crate::error::{Error, Result}; + +/// A three-part version, compared numerically. +/// +/// Deliberately not a `semver` dependency: these CLIs report a plain dotted +/// triple inside prose ("codex-cli 0.145.0", "2.1.205 (Claude Code)"), none of +/// them publishes pre-release or build metadata here, and a whole crate to +/// compare three integers is not worth the graph. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +pub struct Version { + /// Breaking-change component. + pub major: u32, + /// Feature component. + pub minor: u32, + /// Fix component. + pub patch: u32, +} + +impl Version { + /// Find the first dotted triple anywhere in `text`. + /// + /// Each CLI wraps its version in different prose, so this scans rather than + /// parsing a fixed shape: `2.1.205 (Claude Code)`, `codex-cli 0.145.0`, and + /// `GitHub Copilot CLI 1.0.75.` all yield their number. + #[must_use] + pub fn find(text: &str) -> Option { + let bytes = text.as_bytes(); + let mut start = 0; + while start < bytes.len() { + if bytes[start].is_ascii_digit() + // Only start at a boundary, so `1.0.75` inside `abc1.0.75` is + // not read from the middle of a token. + && (start == 0 || !bytes[start - 1].is_ascii_digit() && bytes[start - 1] != b'.') + && let Some(version) = Version::parse_at(&text[start..]) + { + return Some(version); + } + start += 1; + } + None + } + + /// Parse a triple anchored at the start of `text`, ignoring any trailing + /// prose. + fn parse_at(text: &str) -> Option { + let mut parts = [0u32; 3]; + let mut rest = text; + for (index, part) in parts.iter_mut().enumerate() { + if index > 0 { + rest = rest.strip_prefix('.')?; + } + let digits: String = rest.chars().take_while(char::is_ascii_digit).collect(); + if digits.is_empty() { + return None; + } + *part = digits.parse().ok()?; + rest = &rest[digits.len()..]; + } + Some(Version { + major: parts[0], + minor: parts[1], + patch: parts[2], + }) + } +} + +impl fmt::Display for Version { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{}.{}.{}", self.major, self.minor, self.patch) + } +} + +/// How an installed CLI relates to the release this crate was verified against. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum VersionStatus { + /// Exactly the verified release. Mappings are known good. + Verified, + /// Newer. The likely direction of drift: flags may have been renamed, + /// replaced, or moved between subcommands. + Newer, + /// Older. Flags this crate relies on may not exist yet. + Older, + /// The CLI answered, but no version could be read out of it. + Unrecognized, +} + +impl VersionStatus { + /// Whether mappings are known good for this version. + #[must_use] + pub fn is_verified(self) -> bool { + self == VersionStatus::Verified + } +} + +/// What an installed agent CLI reported about itself. +#[derive(Debug, Clone)] +#[non_exhaustive] +pub struct Probe { + /// The agent probed. + pub agent: Agent, + /// The binary that answered. + pub bin: String, + /// Its `--version` output, trimmed. + pub reported: String, + /// The version read out of that output, if one could be. + pub version: Option, + /// The release this crate's mappings were verified against. + pub verified: Version, + /// How the two relate. + pub status: VersionStatus, +} + +impl Probe { + /// Ask `agent`'s default binary for its version. + /// + /// # Errors + /// [`Error::NotInstalled`] if the binary is missing, [`Error::Spawn`] if it + /// cannot be run. + pub async fn run(agent: Agent) -> Result { + Probe::run_bin(agent, agent.bin()).await + } + + /// Ask a specific binary for its version, for a caller that overrides the + /// path with [`crate::Request::bin`]. + /// + /// # Errors + /// [`Error::NotInstalled`] if the binary is missing, [`Error::Spawn`] if it + /// cannot be run. + pub async fn run_bin(agent: Agent, bin: &str) -> Result { + let output = tokio::process::Command::new(bin) + .arg("--version") + .output() + .await + .map_err(|source| { + if source.kind() == std::io::ErrorKind::NotFound { + Error::NotInstalled { + agent, + bin: bin.to_string(), + hint: agent.install_hint(), + } + } else { + Error::Spawn { + bin: bin.to_string(), + source, + } + } + })?; + + // Some CLIs print their version to stderr; take whichever answered. + let mut reported = String::from_utf8_lossy(&output.stdout).trim().to_string(); + if reported.is_empty() { + reported = String::from_utf8_lossy(&output.stderr).trim().to_string(); + } + + let verified = agent.verified_version(); + let version = Version::find(&reported); + let status = match version { + None => VersionStatus::Unrecognized, + Some(found) if found == verified => VersionStatus::Verified, + Some(found) if found > verified => VersionStatus::Newer, + Some(_) => VersionStatus::Older, + }; + + Ok(Probe { + agent, + bin: bin.to_string(), + reported, + version, + verified, + status, + }) + } + + /// A sentence explaining a non-verified version, or `None` when it matches. + /// + /// Written to be shown to a person: it says what was found, what was + /// expected, and what that means for them. + #[must_use] + pub fn advisory(&self) -> Option { + let agent = self.agent; + let verified = self.verified; + match self.status { + VersionStatus::Verified => None, + VersionStatus::Newer => Some(format!( + "{agent} is newer than the {verified} this crate's flags were verified \ + against ({}). Flags may have been renamed or moved between subcommands; \ + a run failing with an unexpected-argument error is the likely symptom.", + self.version + .map_or_else(|| "unknown".into(), |v| v.to_string()), + )), + VersionStatus::Older => Some(format!( + "{agent} is older than the {verified} this crate's flags were verified \ + against ({}). Flags this crate relies on may not exist in it yet.", + self.version + .map_or_else(|| "unknown".into(), |v| v.to_string()), + )), + VersionStatus::Unrecognized => Some(format!( + "could not read a version out of what {agent} reported ({:?}), so its flags \ + cannot be checked against the verified {verified}.", + self.reported, + )), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + /// The real strings, as each CLI actually prints them. + #[test] + fn versions_are_found_inside_each_cli_s_own_prose() { + assert_eq!( + Version::find("2.1.205 (Claude Code)"), + Some(Version { + major: 2, + minor: 1, + patch: 205 + }) + ); + assert_eq!( + Version::find("codex-cli 0.145.0"), + Some(Version { + major: 0, + minor: 145, + patch: 0 + }) + ); + // Note the trailing period, which must not be read as another part. + assert_eq!( + Version::find("GitHub Copilot CLI 1.0.75."), + Some(Version { + major: 1, + minor: 0, + patch: 75 + }) + ); + } + + #[test] + fn text_without_a_triple_yields_nothing() { + for text in ["", "no version here", "1.2", "v1", "beta"] { + assert_eq!(Version::find(text), None, "{text:?}"); + } + } + + /// A version embedded in a longer token must not be read from its middle. + #[test] + fn a_triple_is_only_read_from_a_token_boundary() { + assert_eq!( + Version::find("build20251.0.75"), + Some(Version { + major: 20251, + minor: 0, + patch: 75 + }), + "the whole leading number belongs to the version" + ); + } + + #[test] + fn versions_order_numerically_not_lexically() { + let v = |major, minor, patch| Version { + major, + minor, + patch, + }; + // The case a string comparison gets wrong: 205 > 99. + assert!(v(2, 1, 205) > v(2, 1, 99)); + assert!(v(0, 145, 0) > v(0, 99, 9)); + assert!(v(1, 0, 0) > v(0, 999, 999)); + } + + #[test] + fn every_agent_declares_a_parseable_verified_version() { + for agent in Agent::ALL { + let verified = agent.verified_version(); + assert!(verified.major > 0 || verified.minor > 0, "{agent}"); + } + } + + #[tokio::test] + async fn probing_a_missing_binary_says_how_to_install_it() { + let err = Probe::run_bin(Agent::Claude, "agent-abstraction-no-such-binary") + .await + .unwrap_err(); + assert!(matches!(err, Error::NotInstalled { .. }), "{err:?}"); + } +} diff --git a/src/run.rs b/src/run.rs index f34fea7..033666a 100644 --- a/src/run.rs +++ b/src/run.rs @@ -617,6 +617,15 @@ fn classify(bin: &str, code: i32, stderr: &str, stdout: &str, terminal: &Termina .unwrap_or_else(|| "usage limit reached".to_string()), }; } + // A rejected flag is not a failed request, it is this crate and the CLI + // disagreeing about what the CLI accepts. Naming that is the difference + // between "the run failed" and "your codex is a different version". + if let Some(detail) = rejected_flag(stderr).or_else(|| rejected_flag(stdout)) { + return Error::FlagRejected { + bin: bin.to_string(), + detail, + }; + } Error::Failed { bin: bin.to_string(), code, @@ -624,6 +633,27 @@ fn classify(bin: &str, code: i32, stderr: &str, stdout: &str, terminal: &Termina } } +/// The CLI's complaint, if it refused an argument. +/// +/// The phrasings are clap's and commander's, which is what all three CLIs are +/// built on. Matched narrowly: a false positive would relabel a genuine failure +/// as a version problem and send someone chasing the wrong thing. +fn rejected_flag(text: &str) -> Option { + const REJECTIONS: &[&str] = &[ + "unexpected argument", + "unknown option", + "unrecognized option", + "unknown flag", + "invalid option", + "unexpected option", + ]; + let lower = text.to_ascii_lowercase(); + REJECTIONS + .iter() + .any(|needle| lower.contains(needle)) + .then(|| first_meaningful_line(text).unwrap_or_default()) +} + /// Whether text carries a provider quota refusal. /// /// Deliberately a small set of unambiguous phrases: a false positive here would @@ -734,6 +764,42 @@ mod tests { )); } + /// The exact failure that cost a round of debugging: `codex exec resume` + /// rejects `--sandbox`, which `Error::Failed` reported as a generic + /// non-zero exit naming a flag rather than a version mismatch. + #[test] + fn a_rejected_flag_is_named_as_a_version_mismatch() { + let err = classify( + "codex", + 2, + "error: unexpected argument '--sandbox' found", + "", + &Terminal::default(), + ); + let Error::FlagRejected { bin, detail } = err else { + panic!("expected FlagRejected, got {err:?}") + }; + assert_eq!(bin, "codex"); + assert!(detail.contains("--sandbox"), "{detail}"); + } + + #[test] + fn ordinary_failures_are_not_mistaken_for_version_drift() { + for stderr in [ + "error: no such file or directory", + "model not found", + "permission denied", + ] { + assert!( + matches!( + classify("codex", 1, stderr, "", &Terminal::default()), + Error::Failed { .. } + ), + "{stderr:?} should stay a plain failure" + ); + } + } + #[test] fn failures_report_the_first_useful_line() { let err = classify( diff --git a/tests/live.rs b/tests/live.rs index 156afb6..0e42a64 100644 --- a/tests/live.rs +++ b/tests/live.rs @@ -14,7 +14,8 @@ use std::time::Duration; use agent_abstraction::{ - Agent, EnvPolicy, Event, Format, Permission, Request, SessionStore, run, stream, + Agent, EnvPolicy, Event, Format, Permission, Probe, Request, SessionStore, VersionStatus, run, + stream, }; /// A prompt with exactly one correct answer, so the assertion is about the @@ -447,3 +448,32 @@ async fn a_missing_agent_reports_how_to_install_it() { "a missing agent must say how to install it" ); } + +/// Probing costs no quota, just `--version`, so unlike the rest of this file it +/// runs by default. It is also the test that catches flag drift *before* a run +/// fails on it: if an agent updates underneath us, this goes red and names the +/// version rather than leaving someone to decode an unexpected-argument error. +#[tokio::test] +async fn installed_agents_match_the_versions_the_flags_were_verified_against() { + let mut checked = 0; + for agent in Agent::ALL { + if !available(agent) { + continue; + } + let probe = Probe::run(agent).await.expect("probe failed"); + checked += 1; + + assert!( + probe.version.is_some(), + "{agent} reported {:?}, which carries no readable version", + probe.reported + ); + assert_eq!( + probe.status, + VersionStatus::Verified, + "{}", + probe.advisory().unwrap_or_default() + ); + } + eprintln!("probed {checked} installed agents"); +} From 9d23281fd79f689194eb2315c87aefc4307b884a Mon Sep 17 00:00:00 2001 From: meh Date: Wed, 29 Jul 2026 02:39:28 +0700 Subject: [PATCH 2/4] Classify an unauthenticated agent as its own failure, with the fix MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A missing login was reported as a generic failure, or worse, not reported at all: an unauthenticated Claude run exits **0** and puts "Not logged in · Please run /login" in its result text, so checking the exit code alone handed back a successful Outcome whose answer was a login prompt. Auth is now checked regardless of exit code, the same reasoning that already applied to quota refusals. `Error::NotAuthenticated` carries the agent, the provider's own wording, and the command that fixes it. The hints are per agent because the routes genuinely differ, verified against each CLI's help: Codex and Copilot expose `login` subcommands, Claude authenticates interactively or via `setup-token`. This matters more since EnvPolicy::Minimal became the default. A credential variable this crate does not know to pass through now presents as a login failure, and "not authenticated, run codex login" sends someone to a login screen when the real fix is the environment policy. Naming the category is what makes that distinguishable at all. Deliberately narrow phrase matching, and it loses to nothing: an unrelated failure misread as an auth problem sends someone to re-authenticate over something else entirely. Tested against the failures it must not claim, including rate limits and rejected flags. --- src/agent.rs | 17 +++++++ src/error.rs | 28 ++++++++++++ src/run.rs | 125 ++++++++++++++++++++++++++++++++++++++++++++++++++- 3 files changed, 168 insertions(+), 2 deletions(-) diff --git a/src/agent.rs b/src/agent.rs index f323486..5baa050 100644 --- a/src/agent.rs +++ b/src/agent.rs @@ -255,6 +255,23 @@ impl Agent { } } + /// The command that resolves a missing login for this agent. + /// + /// Verified against each CLI's own help: Codex and Copilot expose a `login` + /// subcommand, while Claude authenticates interactively or through a + /// long-lived token. + #[must_use] + pub fn login_hint(self) -> &'static str { + match self { + Agent::Claude => { + "run `claude` and use /login, or `claude setup-token` for a \ + long-lived token" + } + Agent::Codex => "run `codex login`", + Agent::Copilot => "run `copilot login`", + } + } + /// The release this crate's flag mappings were verified against. /// /// Every mapping in this module was checked by running these exact diff --git a/src/error.rs b/src/error.rs index e93671f..3225d6a 100644 --- a/src/error.rs +++ b/src/error.rs @@ -124,6 +124,25 @@ pub enum Error { detail: String, }, + /// The agent has no usable credentials. + /// + /// Its own category because the remedy is a specific human action rather + /// than anything about the request, and because it is easy to reach by + /// accident: [`crate::EnvPolicy::Minimal`] withholds the environment by + /// default, so a credential this crate does not know to pass through + /// presents as a login failure rather than a configuration one. + #[error("`{bin}` is not authenticated: {message}. To fix: {hint}")] + NotAuthenticated { + /// The agent that refused. + agent: Agent, + /// The binary that refused. + bin: String, + /// The provider's own wording, unedited. + message: String, + /// The command that resolves it. + hint: &'static str, + }, + /// The CLI rejected an argument this crate passed it. /// /// Almost always a version mismatch: the flag was verified against the @@ -212,4 +231,13 @@ impl Error { pub fn is_cancelled(&self) -> bool { matches!(self, Error::Cancelled { .. }) } + + /// Whether this failed because the agent has no usable credentials. + /// + /// Worth branching on in a UI: unlike most failures, the user can fix it, + /// and [`Error::NotAuthenticated`] carries the command that does. + #[must_use] + pub fn is_auth_failure(&self) -> bool { + matches!(self, Error::NotAuthenticated { .. }) + } } diff --git a/src/run.rs b/src/run.rs index 033666a..cd04485 100644 --- a/src/run.rs +++ b/src/run.rs @@ -563,8 +563,19 @@ async fn drive( .rate_limit .as_ref() .is_some_and(crate::outcome::RateLimit::is_blocking); - if exit_code != 0 || quota_blocked { - return Err(classify(&bin, exit_code, &stderr, &raw, &terminal)); + // An unauthenticated Claude run exits 0 and reports the problem in its + // result text, so checking only the exit code would hand back a successful + // Outcome whose answer is "Please run /login". + let unauthenticated = looks_unauthenticated(&terminal.text); + if exit_code != 0 || quota_blocked || unauthenticated { + return Err(classify_run( + request.agent, + &bin, + exit_code, + &stderr, + &raw, + &terminal, + )); } // A fork lands on a *new* id the agent only reveals at the end, so the name @@ -603,6 +614,55 @@ async fn shut_down(child: &mut ChildGuard, stderr_task: tokio::task::JoinHandle< stderr_task.await.unwrap_or_default() } +/// Turn a failure into the most specific error available, agent included so an +/// auth failure can carry the right login command. +fn classify_run( + agent: crate::Agent, + bin: &str, + code: i32, + stderr: &str, + stdout: &str, + terminal: &Terminal, +) -> Error { + // Checked before quota and before a plain failure: a login problem is the + // most specific reading of the output, and the only one a user can act on + // directly. + for source in [terminal.text.as_str(), stderr, stdout] { + if looks_unauthenticated(source) { + return Error::NotAuthenticated { + agent, + bin: bin.to_string(), + message: first_meaningful_line(source).unwrap_or_default(), + hint: agent.login_hint(), + }; + } + } + classify(bin, code, stderr, stdout, terminal) +} + +/// Whether text is an agent saying it has no usable credentials. +/// +/// Narrow on purpose. Mislabelling an ordinary failure as an auth problem sends +/// someone to re-login over something unrelated, so these are phrases the CLIs +/// actually emit rather than every string containing "auth". +fn looks_unauthenticated(text: &str) -> bool { + const PHRASES: &[&str] = &[ + // Claude, verified: an unauthenticated run answers exactly this. + "not logged in", + "please run /login", + "invalid api key", + "authentication_error", + "unauthorized", + "not authenticated", + "no credentials", + "credentials not found", + "please log in", + "401", + ]; + let lower = text.to_ascii_lowercase(); + PHRASES.iter().any(|needle| lower.contains(needle)) +} + /// Turn a non-zero exit into the most specific error available. fn classify(bin: &str, code: i32, stderr: &str, stdout: &str, terminal: &Terminal) -> Error { let quota_signalled = terminal @@ -764,6 +824,67 @@ mod tests { )); } + /// Verified against the real CLI: with `USER` withheld, claude answers + /// "Not logged in · Please run /login" and exits **0**. Checking only the + /// exit code hands back a successful Outcome whose answer is a login + /// prompt. + #[test] + fn an_unauthenticated_run_is_named_even_though_it_exits_zero() { + let terminal = Terminal { + text: "Not logged in · Please run /login".into(), + ..Terminal::default() + }; + let err = classify_run(Agent::Claude, "claude", 0, "", "", &terminal); + let Error::NotAuthenticated { agent, hint, .. } = &err else { + panic!("expected NotAuthenticated, got {err:?}") + }; + assert_eq!(*agent, Agent::Claude); + assert!(hint.contains("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/login"), "{hint}"); + assert!(err.is_auth_failure()); + } + + /// Each agent's hint has to name its own login route, since they differ: + /// Codex and Copilot have `login` subcommands, Claude does not. + #[test] + fn every_agent_offers_its_own_login_route() { + for (agent, expected) in [ + (Agent::Claude, "setup-token"), + (Agent::Codex, "codex login"), + (Agent::Copilot, "copilot login"), + ] { + let err = classify_run( + agent, + agent.bin(), + 1, + "error: unauthorized", + "", + &Terminal::default(), + ); + let Error::NotAuthenticated { hint, .. } = &err else { + panic!("{agent}: expected NotAuthenticated, got {err:?}") + }; + assert!(hint.contains(expected), "{agent}: {hint}"); + } + } + + /// Auth is the most specific reading, so it wins over a generic failure, + /// but must not swallow unrelated errors. + #[test] + fn ordinary_failures_are_not_mistaken_for_auth_problems() { + for stderr in [ + "error: no such file or directory", + "model not found", + "rate limit exceeded", + "error: unexpected argument '--sandbox' found", + ] { + let err = classify_run(Agent::Codex, "codex", 1, stderr, "", &Terminal::default()); + assert!( + !err.is_auth_failure(), + "{stderr:?} was misread as an auth failure: {err:?}" + ); + } + } + /// The exact failure that cost a round of debugging: `codex exec resume` /// rejects `--sandbox`, which `Error::Failed` reported as a generic /// non-zero exit naming a flag rather than a version mismatch. From bc734280093b00e9ae3e424681621507a059b3d8 Mon Sep 17 00:00:00 2001 From: meh Date: Wed, 29 Jul 2026 02:41:08 +0700 Subject: [PATCH 3/4] Bound individual event payloads, not just the total kept MAX_CAPTURE bounds what is kept and the channel bounds how many events queue, but nothing bounded how large a single event is. With a 512 KiB line limit and a 256-deep channel, a consumer that stopped reading could hold roughly 130 MiB of events. Payloads are now capped at MAX_EVENT_BYTES, bringing that to about 16 MiB, which is at least a number worth being able to state. Borrowed from unified-agent-api's per-field bounds, with two changes: - Identifiers are exempt. The session id and tool-call ids pass through whole however long they are, because shortening one breaks the thing it exists for: a truncated session id resumes nothing and a truncated tool id matches no call. A large identifier is a nuisance; a shortened one is a bug. - Tool arguments are replaced rather than truncated. They are JSON, and cutting JSON in half yields something that no longer parses, which is worse for a consumer than an honest placeholder recording what was dropped. Truncated text carries TRUNCATION_MARK so a shortened payload is never mistaken for what the agent actually produced. Enforced at the single point where events leave the parser, so all three agents are covered without each one remembering. --- README.md | 6 ++ src/event.rs | 170 ++++++++++++++++++++++++++++++++++++++++++++++++++- src/lib.rs | 2 +- 3 files changed, 176 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 4e95bab..20d35e0 100644 --- a/README.md +++ b/README.md @@ -262,6 +262,12 @@ run into an OOM instead of an answer. Individual lines are bounded separately at because a reader that accumulates until a newline can exhaust memory on one line that never ends, long before any total cap applies. +Individual event payloads are bounded at `MAX_EVENT_BYTES` (64 KiB) and marked with +`TRUNCATION_MARK` when shortened. The channel bounds how many events queue, not how large +they are, so without this a stalled consumer could hold roughly 130 MiB; with it, about +16 MiB. Identifiers are exempt: a shortened session id cannot resume anything and a +shortened tool id cannot be matched to its call. + Under a structured format there is no silent fallback to raw stdout: a run that produced no recognizable records, or never reached its terminal record, returns `Error::Parse` rather than a plausible-looking answer assembled from whatever was printed. diff --git a/src/event.rs b/src/event.rs index b11ffa4..7dcc35a 100644 --- a/src/event.rs +++ b/src/event.rs @@ -82,6 +82,22 @@ pub const MAX_CAPTURE: usize = 1024 * 1024; /// one emitting an endless line is the case this exists for. pub const MAX_LINE: usize = 512 * 1024; +/// The ceiling on any single event's payload. +/// +/// [`MAX_CAPTURE`] bounds what is *kept*, and the channel bounds how many events +/// are queued, but neither bounds how large one event is. With a 512 KiB line +/// limit and a 256-deep channel, a stalled consumer could hold roughly 130 MiB +/// of events. Bounding the payload brings that to about 16 MiB, which is a +/// number worth being able to state. +/// +/// 64 KiB is far more than a UI renders of a single tool result and is generous +/// for a model turn. +pub const MAX_EVENT_BYTES: usize = 64 * 1024; + +/// Marks a payload this crate shortened, so a truncated value is never mistaken +/// for what the agent actually produced. +pub const TRUNCATION_MARK: &str = "…(truncated)"; + /// The ceiling on how many tool calls may be tracked at once. /// /// Entries are removed as results arrive, so this only bites when an agent @@ -115,6 +131,63 @@ pub(crate) fn append_capped(buf: &mut String, line: &str) -> bool { true } +/// Shorten `text` to [`MAX_EVENT_BYTES`], marking it if anything was dropped. +fn bound_text(text: String) -> String { + if text.len() <= MAX_EVENT_BYTES { + return text; + } + let mut cut = MAX_EVENT_BYTES - TRUNCATION_MARK.len(); + while cut > 0 && !text.is_char_boundary(cut) { + cut -= 1; + } + let mut out = text[..cut].to_string(); + out.push_str(TRUNCATION_MARK); + out +} + +/// Shorten a tool call's arguments, which are structured rather than text. +/// +/// A truncated JSON value would no longer parse, so an oversized one is +/// replaced wholesale by an object recording what was dropped. That keeps the +/// value valid JSON, which is what a consumer expects of this field. +fn bound_value(value: Value) -> Value { + let size = value.to_string().len(); + if size <= MAX_EVENT_BYTES { + return value; + } + serde_json::json!({ + "truncated": true, + "original_bytes": size, + "note": "arguments exceeded MAX_EVENT_BYTES and were dropped rather than \ + truncated, which would have produced invalid JSON", + }) +} + +/// Apply [`MAX_EVENT_BYTES`] to an event's payload. +/// +/// Payloads only. Identifiers, the session id and tool-call ids, are left +/// whole however long they are: they are short in practice, and shortening one +/// would break the thing it exists for, resuming a conversation or matching a +/// result to its call. A truncated identifier is worse than a large one. +fn enforce_bounds(event: Event) -> Event { + match event { + Event::Text(text) => Event::Text(bound_text(text)), + Event::Thinking(text) => Event::Thinking(bound_text(text)), + Event::ToolCall { id, name, input } => Event::ToolCall { + id, + name: bound_text(name), + input: bound_value(input), + }, + Event::ToolResult { id, ok, output } => Event::ToolResult { + id, + ok, + output: bound_text(output), + }, + // Identifiers and quota fields are bounded by their own nature. + other @ (Event::Started { .. } | Event::RateLimit(_)) => other, + } +} + /// Facts that are only known once the stream ends. #[derive(Debug, Clone, Default, PartialEq)] pub struct Terminal { @@ -185,7 +258,7 @@ impl Parser { // is the answer. if self.format == Format::Text { append_capped(&mut self.term.text, line); - return vec![Event::Text(line.to_string())]; + return vec![enforce_bounds(Event::Text(line.to_string()))]; } let Ok(value) = serde_json::from_str::(line) else { self.term.unparsed += 1; @@ -213,6 +286,10 @@ impl Parser { Agent::Codex => self.codex(&value), Agent::Copilot => self.copilot(&value), }; + // Every event leaves through here, so bounding once at the exit covers + // all three agents rather than each parser remembering. + out = out.into_iter().map(enforce_bounds).collect(); + // Fire `Started` exactly once, from whichever record first revealed the // id, and put it ahead of that record's own events. if !self.started { @@ -908,6 +985,97 @@ mod tests { assert_eq!(term.text, "DONE"); } + /// The channel bounds how many events queue, not how large they are. With + /// a 512 KiB line limit that left ~130 MiB reachable in flight. + #[test] + fn an_enormous_tool_result_is_bounded_and_marked() { + let huge = "x".repeat(MAX_EVENT_BYTES * 4); + let line = serde_json::json!({ + "type": "user", + "session_id": "s", + "message": {"content": [{ + "type": "tool_result", "tool_use_id": "t1", "content": huge + }]} + }) + .to_string(); + + let (events, _) = run(Agent::Claude, &[&line]); + let Some(Event::ToolResult { output, id, .. }) = events + .iter() + .find(|e| matches!(e, Event::ToolResult { .. })) + .cloned() + else { + panic!("expected a tool result, got {events:?}") + }; + assert!( + output.len() <= MAX_EVENT_BYTES, + "kept {} bytes", + output.len() + ); + assert!( + output.ends_with(TRUNCATION_MARK), + "truncation must be visible" + ); + assert_eq!(id.as_deref(), Some("t1"), "the id must survive whole"); + } + + /// Identifiers are exempt: a shortened session id cannot resume anything, + /// and a shortened tool id cannot be matched to its call. + #[test] + fn identifiers_are_never_truncated() { + let long_id = "s".repeat(MAX_EVENT_BYTES * 2); + let line = serde_json::json!({"type": "system", "subtype": "init", "session_id": long_id}) + .to_string(); + let (events, term) = run(Agent::Claude, &[&line]); + + let Some(Event::Started { session, .. }) = events.first().cloned() else { + panic!("expected Started, got {events:?}") + }; + assert_eq!(session.len(), long_id.len(), "the session id was shortened"); + assert_eq!(term.session.as_deref(), Some(long_id.as_str())); + } + + /// Truncating JSON would produce something that no longer parses, so an + /// oversized argument object is replaced rather than cut. + #[test] + fn oversized_tool_arguments_stay_valid_json() { + let line = serde_json::json!({ + "type": "assistant", + "session_id": "s", + "message": {"content": [{ + "type": "tool_use", "id": "t1", "name": "Bash", + "input": {"command": "y".repeat(MAX_EVENT_BYTES * 3)} + }]} + }) + .to_string(); + + let (events, _) = run(Agent::Claude, &[&line]); + let Some(Event::ToolCall { input, .. }) = events + .iter() + .find(|e| matches!(e, Event::ToolCall { .. })) + .cloned() + else { + panic!("expected a tool call, got {events:?}") + }; + assert_eq!(input["truncated"], true, "got {input}"); + assert!( + input.is_object(), + "the replacement must still be valid JSON" + ); + assert!(input.to_string().len() <= MAX_EVENT_BYTES); + } + + #[test] + fn ordinary_payloads_pass_through_untouched() { + let (events, _) = run( + Agent::Claude, + &[ + r#"{"type":"assistant","session_id":"s","message":{"content":[{"type":"text","text":"pong"}]}}"#, + ], + ); + assert!(events.contains(&Event::Text("pong".into())), "{events:?}"); + } + #[test] fn capture_is_bounded_and_keeps_the_earliest_output() { let mut buf = String::new(); diff --git a/src/lib.rs b/src/lib.rs index ccaa3a0..3e7db53 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -97,7 +97,7 @@ mod session; pub use agent::{Agent, Caps, EnvPolicy, Format, NETWORK_ENV, Permission, SessionSupport}; pub use error::{Error, Result}; -pub use event::{Event, MAX_CAPTURE}; +pub use event::{Event, MAX_CAPTURE, MAX_EVENT_BYTES, MAX_LINE, TRUNCATION_MARK}; pub use outcome::{Outcome, RateLimit, Stop, Usage}; pub use probe::{Probe, Version, VersionStatus}; pub use request::Request; From c3658cc081dd8e03bf7a5e5d0ea68e3fc682ba65 Mon Sep 17 00:00:00 2001 From: meh Date: Wed, 29 Jul 2026 02:48:10 +0700 Subject: [PATCH 4/4] Kill the process group from Drop directly, not via task abort CI caught this: a dropped Run left grandchildren alive on Linux, while the identical teardown worked from cancel and from a timeout. Drop was signalling the driver and aborting it, which leaves the actual kill waiting on the runtime to poll the aborted task so its ChildGuard runs. That is a dependency on scheduling, and it did not reliably happen. cancel and timeout were unaffected because both call kill_process_group directly. Run now holds the child's pid and kills the group itself in Drop. It is the one thing Drop can do synchronously, and it depends on nothing being polled. The abort stays as the way to stop the driver. A shared `reaped` flag stops Drop signalling a pid the OS may have since handed to someone else, and the three consuming methods disarm it: finish and cancel hand teardown to the driver, and detach would otherwise kill the run it exists to keep alive, which is how the detach test caught that omission immediately. Worth recording how this was found, since I got it wrong twice. The first failure I called flaky and added polling; it failed again identically. The second I diagnosed as a zombie being counted as alive, and fixed the liveness check to read process state. That was also wrong, but it made the third failure say `state Some("S")`: sleeping, not a zombie, genuinely alive. A test that reports what it observed rather than only that it failed is what turned a third round of guessing into a diagnosis. --- src/proc.rs | 19 +++++++++++++++++++ src/run.rs | 53 +++++++++++++++++++++++++++++++++++++++++++++-------- 2 files changed, 64 insertions(+), 8 deletions(-) diff --git a/src/proc.rs b/src/proc.rs index 164782e..576a25f 100644 --- a/src/proc.rs +++ b/src/proc.rs @@ -42,6 +42,21 @@ pub(crate) fn kill_process_group(child: &tokio::process::Child) { // would risk hitting a pid the OS has since recycled. return; }; + kill_group_by_pid(pid); +} + +/// Signal a group by its leader pid, for a caller holding the pid rather than +/// the [`tokio::process::Child`]. +/// +/// [`crate::Run`]'s `Drop` needs this. `Drop` cannot await, so its only other +/// option is to abort the driver task and rely on the runtime polling that task +/// so its guard runs. That makes teardown depend on scheduling, and it does not +/// reliably happen: a dropped `Run` left grandchildren alive and *sleeping* on +/// Linux, while the identical teardown worked from `cancel` and from a timeout, +/// both of which call this directly. Killing here is synchronous and depends on +/// nothing being polled. +#[cfg(unix)] +pub(crate) fn kill_group_by_pid(pid: u32) { let Ok(pid) = i32::try_from(pid) else { // Unreachable in practice: a pid always fits in an i32. return; @@ -66,3 +81,7 @@ pub(crate) fn kill_process_group(child: &tokio::process::Child) { /// No-op: see the module docs. Only the direct child is killed on Windows. #[cfg(not(unix))] pub(crate) fn kill_process_group(_child: &tokio::process::Child) {} + +/// No-op counterpart for non-unix. +#[cfg(not(unix))] +pub(crate) fn kill_group_by_pid(_pid: u32) {} diff --git a/src/run.rs b/src/run.rs index cd04485..c391322 100644 --- a/src/run.rs +++ b/src/run.rs @@ -19,7 +19,7 @@ use crate::agent::{Continue, EnvPolicy}; use crate::error::{Error, Result}; use crate::event::{Event, MAX_LINE, Parser, Terminal, append_capped}; use crate::outcome::{Outcome, Stop}; -use crate::proc::kill_process_group; +use crate::proc::{kill_group_by_pid, kill_process_group}; use crate::request::Request; /// Read one line, giving up on a line that never ends. @@ -90,6 +90,12 @@ pub struct Run { /// The typed command line, kept so both the plain and redacted views come /// from the same source. typed: Vec, + /// The child's pid, so `Drop` can tear the group down itself rather than + /// depending on an aborted task being polled. + pid: Option, + /// Set by the driver once the child has been reaped, so `Drop` never + /// signals a pid the OS may since have handed to someone else. + reaped: std::sync::Arc, /// Dropping or firing this asks the driver to tear down in order. Held as /// an `Option` so `detach` can discard it without signalling. cancel: Option>, @@ -136,6 +142,8 @@ impl Run { /// # Errors /// Whatever the run failed with. See [`Error`]. pub async fn finish(mut self) -> Result { + // The driver owns teardown from here; `Drop` must not also fire. + self.pid = None; while self.events.recv().await.is_some() {} // Taking the handle disarms the `Drop` guard: this run is settling // normally, not being abandoned. @@ -172,6 +180,9 @@ impl Run { /// [`Error::Cancelled`] in the normal case, or whatever the run failed with /// if it failed before the request arrived. pub async fn cancel(mut self) -> Result { + // The driver tears down cooperatively and this awaits it, so `Drop` + // must not race that with a kill of its own. + self.pid = None; // Dropping the sender is itself the signal, so this cannot fail in a // way that leaves the driver waiting. drop(self.cancel.take()); @@ -197,6 +208,9 @@ impl Run { /// afterwards, so reach for this only when an unsupervised background run /// is genuinely intended. pub fn detach(mut self) { + // Disarm `Drop` before it runs, or detaching would immediately kill the + // run it exists to keep alive. + self.pid = None; // Leak the cancel signal rather than dropping it: a dropped sender is // read by the driver as "stop", which is the opposite of detaching. if let Some(cancel) = self.cancel.take() { @@ -209,11 +223,19 @@ impl Run { impl Drop for Run { fn drop(&mut self) { - // Abandoned rather than finished, cancelled or detached. Signal the - // driver so it tears down in order if it gets the chance, then abort so - // the teardown happens even if nothing polls it again. `Drop` cannot - // await, so abort remains the backstop: it drops the driver's - // `ChildGuard`, which kills the process group synchronously. + // Abandoned rather than finished, cancelled or detached. + // + // Kill the group here, directly. Signalling the driver and aborting it + // is not enough on its own: that leaves teardown waiting on the runtime + // to poll the aborted task so its guard runs, and a dropped `Run` was + // observed leaving grandchildren alive and sleeping on Linux while the + // same teardown worked from `cancel`. `Drop` cannot await, so it does + // the one thing it can do synchronously. + if let Some(pid) = self.pid + && !self.reaped.load(std::sync::atomic::Ordering::SeqCst) + { + kill_group_by_pid(pid); + } drop(self.cancel.take()); if let Some(task) = self.task.take() { task.abort(); @@ -334,13 +356,23 @@ pub fn stream(request: &Request) -> Result { } })?; + let pid = child.id(); let (tx, rx) = mpsc::channel(EVENT_BUFFER); let (cancel_tx, cancel_rx) = tokio::sync::oneshot::channel(); + let reaped = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); let request = request.clone(); - let task = runtime.spawn(drive(child, request, tx, cancel_rx)); + let task = runtime.spawn(drive( + child, + request, + tx, + cancel_rx, + std::sync::Arc::clone(&reaped), + )); Ok(Run { events: rx, typed, + pid, + reaped, cancel: Some(cancel_tx), task: Some(task), argv, @@ -391,6 +423,7 @@ async fn drive( request: Request, events: mpsc::Sender, cancel: tokio::sync::oneshot::Receiver<()>, + reaped: std::sync::Arc, ) -> Result { // From here on the child is owned by a guard, so every exit path from this // task, including an abort, takes the process group with it. @@ -495,6 +528,7 @@ async fn drive( // the child's pid, and the group kill needs that pid to target the // group, so the other order silently leaves grandchildren running. let partial = shut_down(&mut child, stderr_task).await; + reaped.store(true, std::sync::atomic::Ordering::SeqCst); return Err(Error::Timeout { bin, timeout: request.timeout.unwrap_or_default(), @@ -506,6 +540,7 @@ async fn drive( // Cooperative teardown: the caller is waiting on this, so the tree // is signalled, reaped and joined before returning. shut_down(&mut child, stderr_task).await; + reaped.store(true, std::sync::atomic::Ordering::SeqCst); return Err(Error::Cancelled { bin }); } } @@ -514,8 +549,10 @@ async fn drive( source, })?; - // The child has been reaped, so its pid must not be signalled again. + // The child has been reaped, so its pid must not be signalled again, by the + // guard here or by `Run::drop` racing this. child.armed = false; + reaped.store(true, std::sync::atomic::Ordering::SeqCst); drop(events); let stderr = stderr_task.await.unwrap_or_default();