Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
40 changes: 40 additions & 0 deletions src/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -255,6 +255,46 @@ 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
/// 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 {
Expand Down
45 changes: 45 additions & 0 deletions src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,42 @@ 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
/// 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
Expand Down Expand Up @@ -195,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 { .. })
}
}
170 changes: 169 additions & 1 deletion src/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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::<Value>(line) else {
self.term.unparsed += 1;
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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();
Expand Down
4 changes: 3 additions & 1 deletion src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -89,15 +89,17 @@ mod agent;
mod error;
mod event;
mod outcome;
mod probe;
mod proc;
mod request;
mod run;
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;
pub use run::{Run, run, stream};
pub use session::{Phase, SessionRecord, SessionStore};
Loading
Loading