Skip to main content

magi/
agent.rs

1//! Driving agent CLIs.
2//!
3//! Every agent in magi is a subscription CLI (`claude`, `opencode`, `agy`) or
4//! an arbitrary command, invoked headless in a working directory. There is no
5//! API-key path on purpose: the CLIs carry the operator's own plan, and they
6//! are the only interface that exposes an agent's whole tool loop rather than a
7//! single completion.
8//!
9//! # Seats, not agents
10//!
11//! Conversations are keyed by *seat* ([`SeatState::key`]), never by agent id. A
12//! model that implements candidate B and also sits as judge 3 gets two
13//! unrelated conversations, so the judge cannot recognise its own work from
14//! having written it. Sessions are what make deliberation affordable — a judge
15//! remembers its own argument instead of being re-fed the entire candidate set
16//! — and seat scoping is what keeps that from destroying blindness.
17//!
18//! # Session mechanics per CLI
19//!
20//! | CLI | open | resume |
21//! |-----|------|--------|
22//! | `claude` | `--session-id <uuid>` (magi mints it) | `--resume <uuid>` |
23//! | `opencode` | `--format json` reports `sessionID` | `-s <id>` |
24//! | `agy` | `--output-format json` reports `conversation_id` | `--conversation <id>` |
25//! | `codex` | `exec --json` reports `thread.started.thread_id` | `exec … resume <id>` |
26//! | `omp` | `-p --mode=json` reports `id` on its `"type":"session"` line | `--resume <id>` |
27//!
28//! Claude is the only one magi can address before the first turn; the others
29//! report an id back, so [`SeatState::captured_session`] stays `None` until a
30//! turn has completed and [`has_session`] answers honestly instead of
31//! optimistically.
32use std::collections::BTreeMap;
33use std::path::{Path, PathBuf};
34use std::process::Stdio;
35use std::time::{Duration, Instant};
36
37use anyhow::{Context as _, Result, bail};
38use serde::{Deserialize, Serialize};
39use std::sync::{Arc, Mutex};
40
41use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
42use tokio::process::Command;
43
44use crate::config::{AgentChoice, AgentKind, AgentSpec, Delivery};
45use crate::proc::Quiet as _;
46use crate::rng::SplitMix64;
47
48/// Conversation state for one seat, persisted with the run so `magi run
49/// --resume` continues the same CLI conversations.
50#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
51pub struct SeatState {
52    /// Stable seat name, e.g. `impl-A`, `judge-2`, `review-1`, `fix`.
53    pub key: String,
54    /// Agent id occupying the seat.
55    pub agent: String,
56    /// Turns already taken in this seat.
57    pub turns: usize,
58    /// Claude session uuid, minted up front so the first turn and every resume
59    /// agree on it without parsing anything back.
60    pub claude_session: Option<String>,
61    /// Session id reported by a CLI that mints its own (`opencode`, `agy`).
62    pub captured_session: Option<String>,
63}
64
65impl SeatState {
66    /// New seat. `run_seed` scopes the minted Claude uuid to this run.
67    pub fn new(key: &str, agent: &str, run_seed: u64) -> Self {
68        let mut rng = SplitMix64::new(run_seed ^ crate::rng::fnv1a(key));
69        Self {
70            key: key.to_owned(),
71            agent: agent.to_owned(),
72            turns: 0,
73            claude_session: Some(rng.uuid_v4()),
74            captured_session: None,
75        }
76    }
77}
78
79/// Can a follow-up prompt rely on this seat remembering the conversation?
80pub fn has_session(kind: AgentKind, seat: &SeatState, sessions_enabled: bool) -> bool {
81    if !sessions_enabled || seat.turns == 0 {
82        return false;
83    }
84    match kind {
85        AgentKind::Claude => seat.claude_session.is_some(),
86        AgentKind::Opencode | AgentKind::Antigravity | AgentKind::Codex | AgentKind::Omp => {
87            seat.captured_session.is_some()
88        }
89        AgentKind::Command => true,
90    }
91}
92
93/// One agent invocation.
94#[derive(Debug)]
95pub struct Invocation<'a> {
96    /// Working directory. Always a real checkout, so every CLI can read the
97    /// repository without per-vendor "extra directory" flags.
98    pub cwd: &'a Path,
99    /// The full prompt.
100    pub prompt: &'a str,
101    /// Wall-clock limit; the process tree is killed when it elapses.
102    pub timeout: Duration,
103    /// May the agent modify files? Judges and reviewers may not.
104    pub allow_write: bool,
105    /// Drop the CLI's own sandbox and approval gate. Conversational talk
106    /// turns only, and only together with `allow_write`: it is the codex
107    /// counterpart of claude's `bypassPermissions`. Unattended seats
108    /// (implementer, judge, reviewer, deputy, ...) never set it, and only
109    /// `talk.rs` may.
110    pub unsandboxed: bool,
111    /// Continue this seat's conversation when the CLI supports it.
112    pub sessions: bool,
113    /// Directory for prompt / stdout / stderr artifacts.
114    pub artifacts: &'a Path,
115    /// Artifact filename stem.
116    pub stem: &'a str,
117    /// Run this invocation belongs to. Exported as `MAGI_RUN` so an agent that
118    /// files a task with `magi task add` is attributed to the run that was
119    /// paying for it, rather than looking like a human wandered by.
120    pub run: &'a str,
121    /// Graph node being executed, e.g. `implement` or `review`. Exported as
122    /// `MAGI_NODE` for the same reason: "who asked for this" is the first
123    /// question about an autonomously created task.
124    pub node: &'a str,
125    /// Shared build cache the seat should build into, from the rendered
126    /// `CARGO_TARGET_DIR=` in the verify commands. Exported as
127    /// `CARGO_TARGET_DIR` so the implementer's compile lands inside the same
128    /// directory `verify` reads back from - one cache, one prune, and the
129    /// build the agent just paid for is the build the gate reuses.
130    pub cache_dir: Option<&'a Path>,
131    /// Absolute paths of images the operator attached to this conversation,
132    /// outside `cwd` - see `chat`/`talk`'s `attachments_dir`. Empty for every
133    /// invocation that is not a chat or talk turn. [`build_command`] uses
134    /// this only to decide whether a CLI's sandbox needs widening to read
135    /// them; the prompt text naming each path and its mime is built by the
136    /// caller, not here.
137    pub attachments: &'a [PathBuf],
138    /// Directories outside `cwd` the seat may write, when it may write at all.
139    /// Only codex's `workspace-write` sandbox is confined to the workspace, so
140    /// only it needs this: a deputy's `magi ask` writes the question store,
141    /// which lives under magi's data directory, not in the repository.
142    pub writable: &'a [PathBuf],
143}
144
145/// Evidence that a CLI ran out of its rate limit / quota, distinct from an
146/// ordinary failure.
147///
148/// `reset` is free text: CLIs render the reset time in their own locale, and
149/// parsing it exactly would be a bug factory. When it is not readable we say
150/// nothing rather than invent a format.
151#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
152pub struct Quota {
153    /// Human-readable reset time, when the CLI printed one.
154    #[serde(default)]
155    pub reset: Option<String>,
156}
157
158/// A CLI hung up on its own stream while the agent was working.
159///
160/// Separate from a failure because the work was done and billed, and separate
161/// from a [`Quota`] because it is worth asking again: the answer is in the
162/// conversation, not lost to a limit that has to reset first. See
163/// [`dropped_stream`].
164#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
165pub struct Dropped {
166    /// What the CLI said as it hung up, verbatim.
167    pub why: String,
168    /// Output tokens the CLI reported before it did - the evidence that this
169    /// was a delivery failure and not an agent that produced nothing.
170    pub output_tokens: u64,
171}
172
173/// One command a CLI's own structured event stream reported running, with
174/// the result it reported for it.
175///
176/// This is evidence the CLI chose to report about its own tool loop — never
177/// something magi polled, supervised, or inferred from a process list. Only
178/// the Codex arm of [`extract`] currently populates it (its `item.completed`
179/// / `command_execution` events name `id`, `command`, `exit_code` and
180/// `aggregated_output` directly); every other backend's CLI does not expose
181/// this in what magi currently captures, so its seats simply never produce
182/// any. A command a CLI never reported finishing (still running when the
183/// turn ended, or the event stream never named it) has no entry here either
184/// — there is no event to build one from, and this type must never be used
185/// to *guess* that a command is still in flight.
186#[derive(Debug, Clone, Serialize, Deserialize)]
187pub struct CommandEvidence {
188    /// The CLI's own id for this command.
189    pub id: String,
190    /// The command itself, as the CLI reported it.
191    pub description: String,
192    /// Exit code the CLI reported for it.
193    pub exit_code: Option<i32>,
194    /// Tail of the command's own output, when the CLI reported one.
195    pub result_summary: String,
196    /// Which CLI/event stream this came from, e.g. `"codex"`.
197    pub source: String,
198}
199
200/// Result of an agent invocation.
201#[derive(Debug, Clone, Serialize, Deserialize)]
202pub struct AgentOutput {
203    /// The agent's final message, extracted from whatever the CLI printed.
204    pub text: String,
205    /// Exit status code.
206    pub exit_code: Option<i32>,
207    /// Did the invocation hit its timeout?
208    pub timed_out: bool,
209    /// Wall-clock duration.
210    pub duration_ms: u64,
211    /// Artifact file names, relative to the run's `artifacts/` directory.
212    pub artifacts: Vec<String>,
213    /// Rate-limit / quota exhaustion, when it can be told apart from a normal
214    /// failure. `None` for a normal failure, a timeout, or a CLI we cannot
215    /// read — the conservative default.
216    #[serde(default)]
217    pub quota: Option<Quota>,
218    /// The CLI hung up on its own stream after the agent had done billed
219    /// work. `None` unless that exact shape was recognised — see
220    /// [`dropped_stream`].
221    #[serde(default)]
222    pub dropped: Option<Dropped>,
223    /// Commands the CLI's own event stream reported running, see
224    /// [`CommandEvidence`]. Always empty for a backend this crate does not
225    /// currently read structured job events from.
226    #[serde(default)]
227    pub commands: Vec<CommandEvidence>,
228    /// Input-side tokens the CLI reported for this turn - approximately the
229    /// size of the context the model just read. `None` when the CLI reports
230    /// no usage (or in a shape this build does not read): unknown, never 0.
231    /// See [`context_tokens`].
232    #[serde(default)]
233    pub context_tokens: Option<u64>,
234}
235
236impl AgentOutput {
237    /// Did the CLI exit cleanly with something to say?
238    pub fn usable(&self) -> bool {
239        !self.timed_out && self.exit_code == Some(0) && !self.text.trim().is_empty()
240    }
241
242    /// Did this invocation run out of the CLI's rate limit / quota?
243    pub fn quota_exhausted(&self) -> bool {
244        self.quota.is_some()
245    }
246
247    /// Did the agent work and the CLI fail to deliver it?
248    ///
249    /// Worth re-asking, unlike [`AgentOutput::quota_exhausted`]: the answer is
250    /// in a conversation this process can resume.
251    pub fn work_undelivered(&self) -> bool {
252        self.dropped.is_some()
253    }
254}
255
256/// How long to keep reading a pipe after the child is gone.
257///
258/// Bounded on purpose: a surviving grandchild can hold the write end open
259/// forever, and the graph must not hang on a process it has already killed.
260const PIPE_GRACE: Duration = Duration::from_secs(3);
261
262/// Bytes a pipe reader has accumulated so far, shared with whoever spawned it.
263type Captured = Arc<Mutex<Vec<u8>>>;
264
265/// Read `pipe` to end in its own task, appending into a buffer the caller can
266/// inspect at any time.
267///
268/// The buffer is shared rather than returned because the interesting moment is
269/// exactly the one where the reader has *not* finished: a killed agent's pipe
270/// may still be held open by a surviving grandchild, and the bytes that did
271/// arrive are the only evidence of what it was doing. An earlier version
272/// returned the buffer from the task and dropped it on timeout, which is how
273/// `<stem>.out` came to be empty on every timeout.
274fn drain<R>(pipe: Option<R>) -> (Captured, Option<tokio::task::JoinHandle<()>>)
275where
276    R: tokio::io::AsyncRead + Unpin + Send + 'static,
277{
278    let buf: Captured = Arc::new(Mutex::new(Vec::new()));
279    let Some(mut pipe) = pipe else {
280        return (buf, None);
281    };
282    let sink = Arc::clone(&buf);
283    let handle = tokio::spawn(async move {
284        let mut chunk = [0u8; 8192];
285        loop {
286            match pipe.read(&mut chunk).await {
287                Ok(0) | Err(_) => break,
288                Ok(n) => {
289                    if let Ok(mut guard) = sink.lock() {
290                        guard.extend_from_slice(&chunk[..n]);
291                    }
292                }
293            }
294        }
295    });
296    (buf, Some(handle))
297}
298
299/// Take whatever a reader has captured, giving it at most `grace` to finish.
300///
301/// A reader still blocked after that is abandoned, not awaited — but its bytes
302/// come back either way, which is the whole point.
303async fn collect(
304    buf: &Captured,
305    handle: Option<tokio::task::JoinHandle<()>>,
306    grace: Duration,
307) -> String {
308    if let Some(handle) = handle {
309        if tokio::time::timeout(grace, handle).await.is_err() {
310            tracing::debug!("a pipe is still held open after the child exited");
311        }
312    }
313    let bytes = buf.lock().map(|g| g.clone()).unwrap_or_default();
314    String::from_utf8_lossy(&bytes).into_owned()
315}
316
317/// Invoke `spec` for `seat`, updating the seat's conversation state.
318pub async fn invoke(
319    spec: &AgentSpec,
320    seat: &mut SeatState,
321    inv: &Invocation<'_>,
322) -> Result<AgentOutput> {
323    tokio::fs::create_dir_all(inv.artifacts)
324        .await
325        .with_context(|| format!("create {}", inv.artifacts.display()))?;
326    let prompt_path = inv.artifacts.join(format!("{}.prompt.md", inv.stem));
327    tokio::fs::write(&prompt_path, inv.prompt)
328        .await
329        .with_context(|| format!("write {}", prompt_path.display()))?;
330
331    let plan = build_command(spec, seat, inv, &prompt_path)?;
332    tracing::debug!(seat = %seat.key, agent = %spec.id, argv = ?plan.argv, "spawning agent");
333
334    let started = Instant::now();
335    // Resolve against the PATH the child will actually see: `spec.env` may
336    // override it, and the resolved absolute path bypasses any later lookup.
337    let child_path = spec
338        .env
339        .iter()
340        .find(|(k, _)| k.eq_ignore_ascii_case("PATH"))
341        .map(|(_, v)| std::ffi::OsString::from(v))
342        .or_else(|| std::env::var_os("PATH"))
343        .unwrap_or_default();
344    let program = crate::config::find_program_on(&plan.argv[0], &child_path).map_or_else(
345        || plan.argv[0].clone().into(),
346        std::path::PathBuf::into_os_string,
347    );
348    let mut cmd = Command::new(program);
349    cmd.args(&plan.argv[1..])
350        .current_dir(inv.cwd)
351        .envs(&spec.env)
352        .env("MAGI_SEAT", &seat.key)
353        .env("MAGI_TURN", seat.turns.to_string())
354        .env("MAGI_RUN", inv.run)
355        .env("MAGI_NODE", inv.node)
356        .env("MAGI_PROMPT_FILE", &prompt_path)
357        .env("MAGI_ALLOW_WRITE", if inv.allow_write { "1" } else { "0" })
358        .env("GIT_TERMINAL_PROMPT", "0")
359        .stdin(if plan.stdin.is_some() {
360            Stdio::piped()
361        } else {
362            Stdio::null()
363        })
364        .stdout(Stdio::piped())
365        .stderr(Stdio::piped())
366        .kill_on_drop(true)
367        // No console window. `magi web` has no console of its own, so Windows
368        // would give each agent a fresh one - and draw it. See `crate::proc`.
369        .quiet();
370    if let Some(cache) = inv.cache_dir {
371        // Same directory the verify commands build into: one cache to prune,
372        // and the compile the seat pays for is the compile the gate reuses.
373        cmd.env("CARGO_TARGET_DIR", cache);
374    } else {
375        // `Command` inherits this process's environment by default, so
376        // simply not setting the variable here is not the same as the seat
377        // not seeing it: if the magi process itself is running under a
378        // shared `CARGO_TARGET_DIR` (the ordinary case), a read-only seat
379        // would otherwise inherit that exact path and try to build there
380        // anyway - the write refusal this is meant to prevent in the first
381        // place. Strip it explicitly.
382        cmd.env_remove("CARGO_TARGET_DIR");
383    }
384
385    let mut child = cmd
386        .spawn()
387        .with_context(|| format!("spawn `{}` for seat {}", plan.argv[0], seat.key))?;
388    // Feed stdin from a task rather than inline: a `command` agent that never
389    // reads its stdin, or a prompt larger than the pipe buffer, would
390    // otherwise deadlock here before the process is ever waited on.
391    if let (Some(body), Some(mut sink)) = (plan.stdin.clone(), child.stdin.take()) {
392        tokio::spawn(async move {
393            sink.write_all(body.as_bytes()).await.ok();
394            sink.shutdown().await.ok();
395        });
396    }
397
398    // Drain the pipes in their own tasks, and wait on the *process*, not on
399    // end-of-file. Two failures come out of conflating those:
400    //
401    // 1. `wait_with_output` returns when both pipes reach EOF, which is not
402    //    when the child exits. A CLI that leaves a helper process holding the
403    //    inherited stdout handle - normal on Windows, where a `.cmd` shim and
404    //    its grandchildren share handles - never closes the pipe, so a seat
405    //    that answered in five minutes was billed the full hour and then
406    //    recorded as a timeout. The answer was thrown away with it.
407    // 2. Cancelling `wait_with_output` at the timeout drops the buffers it
408    //    owned, so `<stem>.out` and `<stem>.err` were written empty exactly
409    //    when an operator needs them most. "It printed nothing" and "we
410    //    discarded what it printed" looked identical on disk.
411    //
412    // Now the readers own the bytes, so a timeout keeps whatever arrived, and
413    // the wait ends at exit even if a stray handle stays open.
414    let (out_buf, out_reader) = drain(child.stdout.take());
415    let (err_buf, err_reader) = drain(child.stderr.take());
416
417    let (code, timed_out) = match tokio::time::timeout(inv.timeout, child.wait()).await {
418        Ok(res) => {
419            let status = res.with_context(|| format!("wait for seat {}", seat.key))?;
420            (status.code(), false)
421        }
422        Err(_) => {
423            tracing::warn!(seat = %seat.key, secs = inv.timeout.as_secs(), "agent timed out");
424            // Kill the tree so the readers see EOF instead of hanging with it.
425            child.start_kill().ok();
426            (None, true)
427        }
428    };
429
430    // The child is gone either way, so the readers are bounded now. A grace
431    // window rather than an unbounded await: a surviving grandchild can still
432    // hold the write end open, and losing a few trailing bytes beats hanging
433    // the graph on a process we no longer control.
434    let stdout = collect(&out_buf, out_reader, PIPE_GRACE).await;
435    let stderr = collect(&err_buf, err_reader, PIPE_GRACE).await;
436
437    let out_path = inv.artifacts.join(format!("{}.out", inv.stem));
438    let err_path = inv.artifacts.join(format!("{}.err", inv.stem));
439    tokio::fs::write(&out_path, &stdout).await.ok();
440    tokio::fs::write(&err_path, &stderr).await.ok();
441
442    let mut extracted = extract(spec.kind, &stdout);
443    if spec.kind == AgentKind::Antigravity && extracted.quota.is_none() {
444        extracted.quota = agy_quota(&stdout, &stderr);
445    }
446    if let Some(quota) = &extracted.quota {
447        // A quota response carries usage too, so it can look like a dropped
448        // stream; a rate limit is never worth re-asking, so it wins.
449        extracted.dropped = None;
450        tracing::warn!(
451            seat = %seat.key,
452            agent = %spec.id,
453            reset = ?quota.reset,
454            "agent is out of quota"
455        );
456    }
457    if let Some(session) = extracted.session {
458        match spec.kind {
459            AgentKind::Claude => seat.claude_session = Some(session),
460            AgentKind::Opencode | AgentKind::Antigravity | AgentKind::Codex | AgentKind::Omp => {
461                seat.captured_session = Some(session);
462            }
463            AgentKind::Command => {}
464        }
465    }
466    if let Some(status) = &extracted.status
467        && !status.eq_ignore_ascii_case("success")
468    {
469        tracing::warn!(seat = %seat.key, status = %status, "agent reported a non-success status");
470    }
471    let text = if extracted.text.trim().is_empty() {
472        // A CLI that printed only to stderr still told us something.
473        if stdout.trim().is_empty() {
474            stderr.trim().to_owned()
475        } else {
476            stdout.trim().to_owned()
477        }
478    } else {
479        extracted.text
480    };
481    seat.turns += 1;
482
483    Ok(AgentOutput {
484        text,
485        exit_code: code,
486        timed_out,
487        duration_ms: started.elapsed().as_millis() as u64,
488        artifacts: vec![
489            file_name(&prompt_path),
490            file_name(&out_path),
491            file_name(&err_path),
492        ],
493        quota: extracted.quota,
494        dropped: extracted.dropped,
495        commands: extracted.commands,
496        context_tokens: extracted.context_tokens,
497    })
498}
499
500fn file_name(p: &Path) -> String {
501    p.file_name()
502        .unwrap_or_default()
503        .to_string_lossy()
504        .into_owned()
505}
506
507/// The argv plus optional stdin body for one invocation.
508#[derive(Debug)]
509struct Plan {
510    argv: Vec<String>,
511    stdin: Option<String>,
512}
513
514/// How a file-delivered prompt is pointed at, per CLI.
515///
516/// `agy` has a native file-context syntax, `@<path>`, and it is measurably the
517/// better contract: on the same trivial task it finished in 17s against 73s for
518/// the prose form, because prose makes the model spend a tool round-trip
519/// deciding to read the file. It is also the form yukimemi/rvpm proved out.
520///
521/// opencode has no equivalent, so it gets the prose. That is not a fallback
522/// worth apologising for — it works, and it is what the winning opencode
523/// candidates on this repository have been driven by all along.
524fn pointer(kind: AgentKind, prompt_path: &Path) -> String {
525    if matches!(kind, AgentKind::Antigravity) {
526        return format!("@{}", prompt_path.display());
527    }
528    format!(
529        "Read the file at {} and follow every instruction in it exactly. That \
530         file is your complete task description; this message contains nothing \
531         else.",
532        prompt_path.display()
533    )
534}
535
536fn build_command(
537    spec: &AgentSpec,
538    seat: &SeatState,
539    inv: &Invocation<'_>,
540    prompt_path: &Path,
541) -> Result<Plan> {
542    let mut argv: Vec<String> = Vec::new();
543    let mut stdin: Option<String> = None;
544    let delivery = spec.delivery();
545    let resuming = has_session(spec.kind, seat, inv.sessions);
546
547    match spec.kind {
548        AgentKind::Claude => {
549            // Claude's own tools have no cwd-confined sandbox - the CLI can
550            // already `Read` any absolute path magi hands it, an attachment
551            // outside the repository included - so no extra flag is needed
552            // here.
553            argv.push("claude".to_owned());
554            argv.push("-p".to_owned());
555            argv.push("--output-format".to_owned());
556            argv.push("json".to_owned());
557            if let Some(m) = &spec.model {
558                argv.push("--model".to_owned());
559                argv.push(m.clone());
560            }
561            if inv.sessions {
562                let uuid = seat
563                    .claude_session
564                    .as_deref()
565                    .context("claude seat is missing its session uuid")?;
566                argv.push(if resuming { "--resume" } else { "--session-id" }.to_owned());
567                argv.push(uuid.to_owned());
568            }
569            argv.push("--permission-mode".to_owned());
570            argv.push("bypassPermissions".to_owned());
571            if !inv.allow_write {
572                argv.push("--disallowed-tools".to_owned());
573                argv.push("Edit,Write,MultiEdit,NotebookEdit".to_owned());
574            }
575        }
576        AgentKind::Opencode => {
577            // `--auto` below already bypasses every permission, reads of a
578            // path outside `--dir` included, so an attachment elsewhere
579            // needs no extra flag.
580            argv.push("opencode".to_owned());
581            argv.push("run".to_owned());
582            argv.push("--format".to_owned());
583            argv.push("json".to_owned());
584            argv.push("--dir".to_owned());
585            argv.push(inv.cwd.to_string_lossy().into_owned());
586            // `--auto` gates *every* permission, reads included: without it a
587            // non-interactive opencode cannot even open the prompt file, and
588            // the seat drops out of the panel with "the user rejected
589            // permission to use this specific tool call". opencode has no
590            // read-only mode, so read-only seats rely on the prompt plus the
591            // fact that judge and reviewer worktrees are disposable — judges'
592            // are deleted after the tally, reviewers' are reset to the commit
593            // under review every round.
594            argv.push("--auto".to_owned());
595            if let Some(m) = &spec.model {
596                argv.push("-m".to_owned());
597                argv.push(m.clone());
598            }
599            if resuming {
600                argv.push("-s".to_owned());
601                argv.push(
602                    seat.captured_session
603                        .clone()
604                        .expect("has_session checked the id is present"),
605                );
606            }
607        }
608        AgentKind::Antigravity => {
609            argv.push("agy".to_owned());
610            argv.push("--output-format".to_owned());
611            argv.push("json".to_owned());
612            // agy's print mode gives up after 5 minutes by default, which is
613            // far below an implementation node's budget.
614            argv.push("--print-timeout".to_owned());
615            argv.push(format!("{}s", inv.timeout.as_secs()));
616            argv.push("--mode".to_owned());
617            argv.push(
618                if inv.allow_write {
619                    "accept-edits"
620                } else {
621                    "plan"
622                }
623                .to_owned(),
624            );
625            if inv.allow_write {
626                argv.push("--dangerously-skip-permissions".to_owned());
627            }
628            if let Some(m) = &spec.model {
629                argv.push("--model".to_owned());
630                argv.push(m.clone());
631            }
632            if resuming {
633                argv.push("--conversation".to_owned());
634                argv.push(
635                    seat.captured_session
636                        .clone()
637                        .expect("has_session checked the id is present"),
638                );
639            }
640            // The prompt file lives outside the worktree, so the workspace has
641            // to be widened to reach it - and so does an attachment's own
642            // directory, which usually lives right beside it under the
643            // conversation's `artifacts_dir` (see `chat`/`talk`). "Usually":
644            // a chat derived from another one (`chat::derived_background`)
645            // can carry attachment paths that live under the *source*
646            // conversation's own artifacts dir instead, so each attachment
647            // outside `inv.artifacts` gets its own `--add-dir` rather than
648            // assuming one directory covers all of `inv.attachments`.
649            let mut add_dirs: Vec<String> = Vec::new();
650            if delivery == Delivery::File || !inv.attachments.is_empty() {
651                add_dirs.push(inv.artifacts.to_string_lossy().into_owned());
652            }
653            for path in inv.attachments {
654                let Some(parent) = path.parent() else {
655                    continue;
656                };
657                if parent.starts_with(inv.artifacts) {
658                    continue;
659                }
660                let dir = parent.to_string_lossy().into_owned();
661                if !add_dirs.contains(&dir) {
662                    add_dirs.push(dir);
663                }
664            }
665            for dir in add_dirs {
666                argv.push("--add-dir".to_owned());
667                argv.push(dir);
668            }
669        }
670        AgentKind::Codex => {
671            // `--sandbox` below governs writes, not reads (see the module
672            // doc: it is what makes codex the one kind whose *read-only*
673            // mode is enforced, by refusing edits) - both presets can read
674            // anywhere the OS lets the process, so an attachment outside
675            // `cwd` is already reachable without an extra flag.
676            argv.push("codex".to_owned());
677            argv.push("exec".to_owned());
678            argv.push("--json".to_owned());
679            // The worktrees magi hands out are real checkouts, but a judge's
680            // is detached and a fixture's may be no repository at all.
681            argv.push("--skip-git-repo-check".to_owned());
682            argv.push("-C".to_owned());
683            argv.push(inv.cwd.to_string_lossy().into_owned());
684            // Codex is the only kind whose read-only-ness is enforced by the
685            // CLI rather than by the prompt: a judge or reviewer seat cannot
686            // write even if it decides to try. Implementers get the workspace.
687            // The one exception is a conversational talk turn that the
688            // repository opted into (`unsandboxed` and `allow_write`): it
689            // drops sandbox and approvals altogether, as claude's
690            // `bypassPermissions` does, so `magi task add`, `git fetch` and
691            // the like are not refused. Nothing else ever gets the bypass.
692            if inv.allow_write && inv.unsandboxed {
693                argv.push("--dangerously-bypass-approvals-and-sandbox".to_owned());
694            } else {
695                argv.push("--sandbox".to_owned());
696                argv.push(
697                    if inv.allow_write {
698                        "workspace-write"
699                    } else {
700                        "read-only"
701                    }
702                    .to_owned(),
703                );
704                // Nothing is watching to approve anything: an unattended seat
705                // that asks blocks until its timeout kills it.
706                argv.push("-c".to_owned());
707                argv.push("approval_policy=\"never\"".to_owned());
708                if inv.allow_write && !inv.writable.is_empty() {
709                    let roots: Vec<String> = inv
710                        .writable
711                        .iter()
712                        .map(|p| p.to_string_lossy().into_owned())
713                        .collect();
714                    argv.push("-c".to_owned());
715                    argv.push(format!(
716                        "sandbox_workspace_write.writable_roots={}",
717                        serde_json::to_string(&roots).unwrap_or_else(|_| "[]".to_owned())
718                    ));
719                }
720            }
721            if let Some(m) = &spec.model {
722                argv.push("-m".to_owned());
723                argv.push(m.clone());
724            }
725            // `resume` is a subcommand of `exec`, and it rejects the flags
726            // above when they follow it - so every option is emitted first and
727            // the subcommand last. Established by hand against codex-cli
728            // 0.153.4: with the order reversed the CLI exits on
729            // `unexpected argument '--sandbox'`.
730            if resuming {
731                argv.push("resume".to_owned());
732                argv.push(
733                    seat.captured_session
734                        .clone()
735                        .expect("has_session checked the id is present"),
736                );
737            }
738        }
739        AgentKind::Omp => {
740            // `omp` reads the prompt from stdin in print mode (see
741            // `AgentSpec::delivery`), so the whole instruction arrives without
742            // an argv length limit - the same reason codex gets stdin.
743            argv.push("omp".to_owned());
744            argv.push("-p".to_owned());
745            argv.push("--mode=json".to_owned());
746            // `--auto-approve` is required, and is the same trade opencode's
747            // `--auto` makes: it gates *every* permission, reads included, so
748            // without it a non-interactive seat cannot even open the prompt
749            // file magi wrote and drops out of the panel on a permission
750            // rejection. `omp` has no read-only mode of its own, so a judge or
751            // reviewer seat rests on the prompt plus the worktree discipline
752            // (judge worktrees are deleted after the tally, reviewer worktrees
753            // are reset to the commit under review every round) - never on this
754            // flag, and never on a bypass flag.
755            argv.push("--auto-approve".to_owned());
756            if let Some(m) = &spec.model {
757                argv.push("--model".to_owned());
758                argv.push(m.clone());
759            }
760            // Established by hand against omp 18.1.19: `-p --mode=json` reports
761            // the session id on its `"type":"session"` line, and
762            // `--resume <id>` continues that conversation. `--continue` is
763            // deliberately not used - it opens a *new* session rather than the
764            // stored one, which silently loses the seat's memory.
765            if resuming {
766                argv.push("--resume".to_owned());
767                argv.push(
768                    seat.captured_session
769                        .clone()
770                        .expect("has_session checked the id is present"),
771                );
772            }
773        }
774        AgentKind::Command => {
775            // The operator's own command line, not one of the roster CLIs -
776            // there is no flag this function could add on its behalf, so an
777            // attachment's path has to reach it the same way the prompt
778            // does, through the substitutions below.
779            if spec.command.is_empty() {
780                bail!("agent `{}` has kind = \"command\" but no command", spec.id);
781            }
782            let vars: BTreeMap<&str, String> = BTreeMap::from([
783                ("{prompt_file}", prompt_path.to_string_lossy().into_owned()),
784                ("{cwd}", inv.cwd.to_string_lossy().into_owned()),
785                ("{label}", seat.key.clone()),
786                ("{session}", seat.claude_session.clone().unwrap_or_default()),
787            ]);
788            for raw in &spec.command {
789                let mut arg = raw.clone();
790                for (k, v) in &vars {
791                    if arg.contains(k) {
792                        arg = arg.replace(k, v);
793                    }
794                }
795                argv.push(arg);
796            }
797        }
798    }
799
800    argv.extend(spec.extra_args.iter().cloned());
801
802    // `agy` takes the prompt as the value of `-p`, so the flag has to be
803    // emitted right before whatever the delivery mode produces.
804    if spec.kind == AgentKind::Antigravity {
805        argv.push("-p".to_owned());
806    }
807    // `codex exec` reads stdin only when its prompt argument is `-`; without
808    // it the CLI waits on a prompt it will never be given.
809    if spec.kind == AgentKind::Codex && delivery == Delivery::Stdin {
810        argv.push("-".to_owned());
811    }
812    match delivery {
813        Delivery::Stdin if spec.kind == AgentKind::Antigravity => {
814            // agy has no text stdin path; fall back to the pointer file.
815            argv.push(pointer(spec.kind, prompt_path));
816        }
817        Delivery::Stdin => stdin = Some(inv.prompt.to_owned()),
818        Delivery::Argv => argv.push(inv.prompt.to_owned()),
819        Delivery::File => argv.push(pointer(spec.kind, prompt_path)),
820    }
821
822    Ok(Plan { argv, stdin })
823}
824
825/// What a CLI's stdout yielded.
826#[derive(Debug, Default)]
827struct Extracted {
828    text: String,
829    session: Option<String>,
830    status: Option<String>,
831    quota: Option<Quota>,
832    dropped: Option<Dropped>,
833    commands: Vec<CommandEvidence>,
834    context_tokens: Option<u64>,
835}
836
837/// Pull the agent's message (and any session id) out of a CLI's stdout.
838fn extract(kind: AgentKind, stdout: &str) -> Extracted {
839    let mut extracted = extract_answer(kind, stdout);
840    extracted.context_tokens = context_tokens(kind, stdout);
841    extracted
842}
843
844/// `v[key]` as an unsigned integer, or `None` when absent or not a number.
845fn uint(v: &serde_json::Value, key: &str) -> Option<u64> {
846    v.get(key).and_then(serde_json::Value::as_u64)
847}
848
849/// The input side of a usage object: `input` plus whatever the cache served
850/// or stored, each key read under the names the CLIs use. A missing cache
851/// counter is 0 (the CLI simply had no cache traffic to report) but a missing
852/// *input* counter makes the whole thing unknown - `None`, never a made-up 0.
853fn input_side(usage: &serde_json::Value, input: &[&str], cache: &[&str]) -> Option<u64> {
854    let first = |keys: &[&str]| keys.iter().find_map(|k| uint(usage, k));
855    let base = first(input)?;
856    let cached: u64 = cache.iter().filter_map(|k| uint(usage, k)).sum();
857    Some(base.saturating_add(cached))
858}
859
860/// The context size a CLI reported for its last turn: its input-side token
861/// count, which approximates how much of the model's window the conversation
862/// occupies. A pure function of stdout, so each CLI's shape is unit-testable.
863///
864/// Rules, all deliberate:
865/// - **Last report wins, never a sum.** A multi-event stream (opencode's
866///   `step_finish`, codex's `turn.completed`, omp's `message_end`) reports the
867///   input of each model call; the final one is the largest context the
868///   conversation reached this turn, and adding them would count the same
869///   prefix once per tool round-trip.
870/// - **Missing or mistyped usage is `None`**, which the UI shows as "unknown".
871/// - **A cumulative total is not a context size, so it is `None`.** claude's
872///   result `usage` sums every model call of the turn and, on a resumed
873///   session, carries the session's earlier usage too, so even `num_turns == 1`
874///   does not make it one call's context. codex's `turn.completed.usage` is the
875///   session's running input total, and agy's `usage` aggregates the whole
876///   print-mode call (the observed sample exceeds any window). None of the
877///   three carries a per-call figure, so each is unknown rather than made up.
878///   Only the streams that report each model call (opencode `step_finish`, omp
879///   message usage) yield a number.
880fn context_tokens(kind: AgentKind, stdout: &str) -> Option<u64> {
881    use serde_json::Value;
882    let lines = || {
883        stdout
884            .lines()
885            .filter_map(|l| serde_json::from_str::<Value>(l.trim()).ok())
886    };
887    match kind {
888        AgentKind::Claude => None,
889        AgentKind::Opencode => lines()
890            .filter(|v| {
891                // The event is `step_finish`; its part is `step-finish`.
892                v.get("type").and_then(|t| t.as_str()) == Some("step_finish")
893                    || v.get("part")
894                        .and_then(|p| p.get("type"))
895                        .and_then(|t| t.as_str())
896                        == Some("step-finish")
897            })
898            .filter_map(|v| {
899                let tokens = v.get("part")?.get("tokens")?;
900                let base = uint(tokens, "input")?;
901                let cache = tokens.get("cache");
902                let cached = ["read", "write"]
903                    .iter()
904                    .filter_map(|k| cache.and_then(|c| uint(c, k)))
905                    .sum::<u64>();
906                Some(base.saturating_add(cached))
907            })
908            .next_back(),
909        AgentKind::Antigravity => None,
910        AgentKind::Codex => None,
911        AgentKind::Omp => lines()
912            .flat_map(|v| {
913                // Same walk as the answer: `agent_end` carries the thread,
914                // `turn_end` / `message_end` one message each.
915                match v.get("type").and_then(|t| t.as_str()) {
916                    Some("agent_end") => v
917                        .get("messages")
918                        .and_then(|m| m.as_array())
919                        .cloned()
920                        .unwrap_or_default(),
921                    Some("turn_end") | Some("message_end") => {
922                        v.get("message").cloned().into_iter().collect()
923                    }
924                    _ => Vec::new(),
925                }
926            })
927            .filter(|m| m.get("role").and_then(|r| r.as_str()) == Some("assistant"))
928            .filter_map(|m| {
929                input_side(
930                    m.get("usage")?,
931                    &["input", "input_tokens"],
932                    &[
933                        "cacheRead",
934                        "cacheWrite",
935                        "cache_read_tokens",
936                        "cache_write_tokens",
937                    ],
938                )
939            })
940            .next_back(),
941        AgentKind::Command => None,
942    }
943}
944
945fn extract_answer(kind: AgentKind, stdout: &str) -> Extracted {
946    match kind {
947        AgentKind::Claude => {
948            let Ok(v) = serde_json::from_str::<serde_json::Value>(stdout.trim()) else {
949                return Extracted {
950                    text: stdout.trim().to_owned(),
951                    ..Extracted::default()
952                };
953            };
954            Extracted {
955                text: v
956                    .get("result")
957                    .and_then(|r| r.as_str())
958                    .unwrap_or_default()
959                    .to_owned(),
960                session: v
961                    .get("session_id")
962                    .and_then(|s| s.as_str())
963                    .map(str::to_owned),
964                status: v.get("is_error").and_then(|e| e.as_bool()).map(|e| {
965                    if e {
966                        "error".to_owned()
967                    } else {
968                        "success".to_owned()
969                    }
970                }),
971                quota: claude_quota(&v),
972                // Claude reports a truncated stream as an ordinary error; the
973                // shape `dropped_stream` keys on is agy's.
974                dropped: None,
975                commands: Vec::new(),
976                context_tokens: None,
977            }
978        }
979        AgentKind::Opencode => {
980            // A JSONL event stream: text parts concatenated in arrival order.
981            let mut text = String::new();
982            let mut session = None;
983            for line in stdout.lines() {
984                let Ok(v) = serde_json::from_str::<serde_json::Value>(line.trim()) else {
985                    continue;
986                };
987                if session.is_none() {
988                    session = v
989                        .get("sessionID")
990                        .and_then(|s| s.as_str())
991                        .map(str::to_owned);
992                }
993                let part = v.get("part").unwrap_or(&serde_json::Value::Null);
994                if part.get("type").and_then(|t| t.as_str()) == Some("text")
995                    && let Some(t) = part.get("text").and_then(|t| t.as_str())
996                {
997                    if !text.is_empty() {
998                        text.push('\n');
999                    }
1000                    text.push_str(t);
1001                }
1002            }
1003            Extracted {
1004                text,
1005                session,
1006                status: None,
1007                quota: None,
1008                dropped: None,
1009                commands: Vec::new(),
1010                context_tokens: None,
1011            }
1012        }
1013        AgentKind::Antigravity => {
1014            // agy prints warnings before the JSON object, so parse the last
1015            // line that is one rather than the whole stream.
1016            let obj = stdout
1017                .lines()
1018                .rev()
1019                .find_map(|l| serde_json::from_str::<serde_json::Value>(l.trim()).ok());
1020            let Some(v) = obj else {
1021                return Extracted {
1022                    text: stdout.trim().to_owned(),
1023                    ..Extracted::default()
1024                };
1025            };
1026            Extracted {
1027                text: v
1028                    .get("response")
1029                    .and_then(|r| r.as_str())
1030                    .unwrap_or_default()
1031                    .trim()
1032                    .to_owned(),
1033                session: v
1034                    .get("conversation_id")
1035                    .and_then(|s| s.as_str())
1036                    .map(str::to_owned),
1037                status: v.get("status").and_then(|s| s.as_str()).map(str::to_owned),
1038                quota: None,
1039                dropped: dropped_stream(&v),
1040                commands: Vec::new(),
1041                context_tokens: None,
1042            }
1043        }
1044        AgentKind::Codex => {
1045            // A JSONL event stream, prefixed on a real machine by tracing
1046            // lines the CLI writes about its own config and skills - so
1047            // non-JSON lines are skipped rather than treated as the answer.
1048            //
1049            // The thread id arrives once, in `thread.started`, and a resumed
1050            // turn reports the same one. The answer is the last
1051            // `item.completed` carrying an `agent_message`: earlier ones are
1052            // the model narrating its way through the tool loop, and taking
1053            // the first would hand the caller a progress note instead of a
1054            // verdict.
1055            //
1056            // A `command_execution` item is different evidence entirely: the
1057            // CLI's own record that it ran a command and what that command
1058            // reported back, kept as `CommandEvidence` — see run
1059            // 20260912-214939-b3bb's artifacts, where these very fields
1060            // (`id`/`command`/`exit_code`/`aggregated_output`) were what
1061            // caught a paired test result magi's own agent prose had missed.
1062            // Only what `item.completed` actually reports: a command still
1063            // running when the turn ended emits no such event at all, and is
1064            // not something this can detect — see [`CommandEvidence`]'s own
1065            // doc for why that must not be guessed at instead.
1066            let mut text = String::new();
1067            let mut session = None;
1068            let mut status = None;
1069            let mut commands = Vec::new();
1070            for line in stdout.lines() {
1071                let Ok(v) = serde_json::from_str::<serde_json::Value>(line.trim()) else {
1072                    continue;
1073                };
1074                match v.get("type").and_then(|t| t.as_str()) {
1075                    Some("thread.started") => {
1076                        session = v
1077                            .get("thread_id")
1078                            .and_then(|s| s.as_str())
1079                            .map(str::to_owned);
1080                    }
1081                    Some("item.completed") => {
1082                        let item = v.get("item").unwrap_or(&serde_json::Value::Null);
1083                        match item.get("type").and_then(|t| t.as_str()) {
1084                            Some("agent_message") => {
1085                                if let Some(t) = item.get("text").and_then(|t| t.as_str()) {
1086                                    text = t.trim().to_owned();
1087                                }
1088                            }
1089                            Some("command_execution") => {
1090                                commands.push(command_evidence(item));
1091                            }
1092                            _ => {}
1093                        }
1094                    }
1095                    Some("turn.completed") => status = Some("success".to_owned()),
1096                    Some("turn.failed") => status = Some("error".to_owned()),
1097                    _ => {}
1098                }
1099            }
1100            Extracted {
1101                text,
1102                session,
1103                status,
1104                quota: None,
1105                dropped: None,
1106                commands,
1107                context_tokens: None,
1108            }
1109        }
1110        AgentKind::Omp => {
1111            // A JSONL event stream. The session id arrives once, on the
1112            // `"type":"session"` line that opens the run.
1113            //
1114            // The answer is the *last* non-empty assistant text block anywhere
1115            // in the stream, and neither of the two obvious shortcuts works:
1116            //
1117            // 1. Do not key on `agent_end`. `omp` emits it only for a run that
1118            //    quiesces on a message turn; a turn that ends on a tool call
1119            //    (`stopReason: "toolUse"`) ends the run with **no `agent_end`
1120            //    line at all**, and the answer is in `message_end` / `turn_end`
1121            //    instead. Reading only `agent_end` silently discards a complete
1122            //    review - which is exactly what the first hand-written wrapper
1123            //    did, three times, before this arm existed.
1124            // 2. Do not take the first assistant text. Earlier ones narrate the
1125            //    tool loop (sometimes with a single `.`), so the last non-empty
1126            //    block is the answer and the one before it is a progress note.
1127            //
1128            // Every line is parsed independently: a non-JSON line (a CLI
1129            // warning, a truncated write) is skipped rather than treated as the
1130            // answer, the same way the codex arm treats its tracing prefix.
1131            let mut text = String::new();
1132            let mut session = None;
1133            for line in stdout.lines() {
1134                let Ok(v) = serde_json::from_str::<serde_json::Value>(line.trim()) else {
1135                    continue;
1136                };
1137                if v.get("type").and_then(|t| t.as_str()) == Some("session") {
1138                    session = v.get("id").and_then(|s| s.as_str()).map(str::to_owned);
1139                    continue;
1140                }
1141                // `agent_end` carries the whole thread; `turn_end` and
1142                // `message_end` each carry one message. Whichever appears, the
1143                // messages are walked the same way.
1144                let messages: Vec<&serde_json::Value> = match v.get("type").and_then(|t| t.as_str())
1145                {
1146                    Some("agent_end") => v
1147                        .get("messages")
1148                        .and_then(|m| m.as_array())
1149                        .map(|m| m.iter().collect())
1150                        .unwrap_or_default(),
1151                    Some("turn_end") | Some("message_end") => {
1152                        v.get("message").into_iter().collect()
1153                    }
1154                    _ => continue,
1155                };
1156                for message in messages {
1157                    if message.get("role").and_then(|r| r.as_str()) != Some("assistant") {
1158                        continue;
1159                    }
1160                    let Some(parts) = message.get("content").and_then(|c| c.as_array()) else {
1161                        continue;
1162                    };
1163                    for part in parts {
1164                        if part.get("type").and_then(|t| t.as_str()) != Some("text") {
1165                            continue;
1166                        }
1167                        if let Some(t) = part.get("text").and_then(|t| t.as_str())
1168                            && !t.trim().is_empty()
1169                        {
1170                            text = t.trim().to_owned();
1171                        }
1172                    }
1173                }
1174            }
1175            Extracted {
1176                text,
1177                session,
1178                status: None,
1179                quota: None,
1180                dropped: None,
1181                commands: Vec::new(),
1182                context_tokens: None,
1183            }
1184        }
1185        AgentKind::Command => {
1186            // A `command` agent may wrap a subscription CLI (a fixture, or a
1187            // thin shim around `claude`). If its output is the claude error
1188            // shape we recognise the quota the same way, so tests and wrappers
1189            // do not need their own detection; anything else is just text.
1190            let parsed = serde_json::from_str::<serde_json::Value>(stdout.trim()).ok();
1191            let quota = parsed.as_ref().and_then(claude_quota);
1192            // A `command` fixture may also stand in for a CLI that hangs up on
1193            // its own stream, which is how that path is tested.
1194            let dropped = parsed.as_ref().and_then(dropped_stream);
1195            Extracted {
1196                text: stdout.trim().to_owned(),
1197                session: None,
1198                status: None,
1199                quota,
1200                dropped,
1201                commands: Vec::new(),
1202                context_tokens: None,
1203            }
1204        }
1205    }
1206}
1207
1208/// Build one [`CommandEvidence`] from a Codex `command_execution` item.
1209///
1210/// `command` arrives either as a single string or as an argv array,
1211/// depending on how the CLI shaped the call; both are read rather than
1212/// assuming one. Missing fields are left at their honest defaults (an empty
1213/// id/description, `exit_code: None`) rather than guessed at.
1214fn command_evidence(item: &serde_json::Value) -> CommandEvidence {
1215    let description = match item.get("command") {
1216        Some(serde_json::Value::String(s)) => s.clone(),
1217        Some(serde_json::Value::Array(parts)) => parts
1218            .iter()
1219            .filter_map(|p| p.as_str())
1220            .collect::<Vec<_>>()
1221            .join(" "),
1222        _ => String::new(),
1223    };
1224    let result_summary = item
1225        .get("aggregated_output")
1226        .and_then(|o| o.as_str())
1227        .map(|s| tail_chars(s.trim(), 400))
1228        .unwrap_or_default();
1229    CommandEvidence {
1230        id: item
1231            .get("id")
1232            .and_then(|s| s.as_str())
1233            .unwrap_or_default()
1234            .to_owned(),
1235        description,
1236        exit_code: item
1237            .get("exit_code")
1238            .and_then(serde_json::Value::as_i64)
1239            .map(|e| e as i32),
1240        result_summary,
1241        source: "codex".to_owned(),
1242    }
1243}
1244
1245/// The last `max` characters of `s`, cut on a char boundary.
1246fn tail_chars(s: &str, max: usize) -> String {
1247    let count = s.chars().count();
1248    if count <= max {
1249        return s.to_owned();
1250    }
1251    s.chars().skip(count - max).collect()
1252}
1253
1254/// Recognise claude's rate-limit error shape, when it is present.
1255///
1256/// The only output we have observed is the JSON object carrying `is_error:
1257/// true` and a `result` mentioning the session limit. We key on exactly that;
1258/// every other CLI (and any future shape) returns `None` and is treated as an
1259/// ordinary failure — the conservative side.
1260fn claude_quota(v: &serde_json::Value) -> Option<Quota> {
1261    let is_err = v.get("is_error").and_then(|e| e.as_bool()).unwrap_or(false);
1262    if !is_err {
1263        return None;
1264    }
1265    let result = v.get("result").and_then(|r| r.as_str()).unwrap_or("");
1266    if !result.to_lowercase().contains("session limit") {
1267        return None;
1268    }
1269    // "…session limit · resets 4:50am (Asia/Tokyo)". The timezone read is not
1270    // worth parsing exactly; keep the whole phrase after "resets" as free text.
1271    let reset = result
1272        .split("resets ")
1273        .nth(1)
1274        .map(str::trim)
1275        .filter(|s| !s.is_empty())
1276        .map(str::to_owned);
1277    Some(Quota { reset })
1278}
1279
1280/// Recognise `agy` running out of quota, from either stream.
1281///
1282/// Observed shape while out of quota: stdout carries the JSON object
1283/// `{"status":"ERROR","response":"","error":"Individual quota reached. … Resets
1284/// in 1h2m49s."}` and stderr carries `AGY_ERROR: {"status":"RESOURCE_EXHAUSTED",
1285/// "error_code":429,…}`. Either alone is enough (a mangled stdout must not hide
1286/// a quota that stderr states plainly), but each is keyed on a *pair* of
1287/// structured fields: `status: ERROR` with the quota text, or
1288/// `RESOURCE_EXHAUSTED` together with 429. An error status alone, or a 429
1289/// alone, stays an ordinary failure - the conservative side, as for
1290/// [`claude_quota`].
1291///
1292/// The reset hint is what follows `Resets ` (e.g. `in 1h2m49s`), trailing full
1293/// stop dropped.
1294fn agy_quota(stdout: &str, stderr: &str) -> Option<Quota> {
1295    let from_stdout = stdout
1296        .lines()
1297        .rev()
1298        .find_map(|l| serde_json::from_str::<serde_json::Value>(l.trim()).ok())
1299        .and_then(|v| {
1300            let status = v.get("status").and_then(|s| s.as_str()).unwrap_or("");
1301            let error = v.get("error").and_then(|e| e.as_str()).unwrap_or("");
1302            (status.eq_ignore_ascii_case("error") && error.to_lowercase().contains("quota reached"))
1303                .then(|| error.to_owned())
1304        });
1305    let from_stderr = || {
1306        stderr.lines().find_map(|l| {
1307            let v: serde_json::Value =
1308                serde_json::from_str(l.trim().strip_prefix("AGY_ERROR:")?.trim()).ok()?;
1309            let exhausted = v.get("status").and_then(|s| s.as_str()) == Some("RESOURCE_EXHAUSTED");
1310            let code = v.get("error_code").and_then(serde_json::Value::as_u64) == Some(429);
1311            (exhausted && code).then(|| {
1312                v.get("short_error")
1313                    .and_then(|e| e.as_str())
1314                    .unwrap_or_default()
1315                    .to_owned()
1316            })
1317        })
1318    };
1319    let text = from_stdout.or_else(from_stderr)?;
1320    let reset = text
1321        .split_once("Resets ")
1322        .map(|(_, rest)| rest.trim().trim_end_matches('.').trim())
1323        .filter(|s| !s.is_empty())
1324        .map(str::to_owned);
1325    Some(Quota { reset })
1326}
1327
1328/// Recognise a CLI that gave up on its own stream while the agent was working.
1329///
1330/// Observed once, verbatim, from `agy` on a candidate that produced nothing:
1331///
1332/// ```text
1333/// {"conversation_id":"36743d06-…","status":"ERROR","response":"",
1334///  "error":"the connection to the agent was interrupted before the response
1335///           finished: subscriber fell behind updates, stalled for 5s",
1336///  "duration_seconds":431.19,"num_turns":1,
1337///  "usage":{"input_tokens":260113,"output_tokens":14267,
1338///           "thinking_tokens":9695,"cache_read_tokens":2200925}}
1339/// ```
1340///
1341/// Seven minutes of work and fourteen thousand output tokens, billed, with an
1342/// empty `response`: the agent did the job and the CLI's own subscriber fell
1343/// behind and hung up. That is **not** an agent that failed to implement, and
1344/// counting it as one is how `agy` came to read as 0 wins in 4 entries with
1345/// five empty candidates - a number that has twice been used to argue the seat
1346/// out of the roster, and twice been wrong (see `cb6b830`, which reverted the
1347/// first removal: *"agy does not fail to implement, it fails to report"*).
1348///
1349/// The distinction that matters is **billed work with nothing delivered**, so
1350/// that is what this keys on: an error status, an empty response, and a usage
1351/// report showing output tokens. Everything else - including an error with no
1352/// usage at all - returns `None` and stays an ordinary failure, the
1353/// conservative side, exactly as [`claude_quota`] treats shapes it does not
1354/// recognise.
1355///
1356/// Unlike a quota, this **is** worth re-asking: the work exists in the
1357/// conversation the CLI just abandoned, and `conversation_id` is right there.
1358fn dropped_stream(v: &serde_json::Value) -> Option<Dropped> {
1359    let status = v.get("status").and_then(|s| s.as_str()).unwrap_or("");
1360    if !status.eq_ignore_ascii_case("error") {
1361        return None;
1362    }
1363    let response = v.get("response").and_then(|r| r.as_str()).unwrap_or("");
1364    if !response.trim().is_empty() {
1365        // It answered. Whatever the status says, there is something to read.
1366        return None;
1367    }
1368    let produced = v
1369        .get("usage")
1370        .and_then(|u| u.get("output_tokens"))
1371        .and_then(serde_json::Value::as_u64)
1372        .unwrap_or(0);
1373    if produced == 0 {
1374        // An error with nothing produced is just an error.
1375        return None;
1376    }
1377    Some(Dropped {
1378        why: v
1379            .get("error")
1380            .and_then(|e| e.as_str())
1381            .unwrap_or("the CLI ended the stream without delivering its answer")
1382            .trim()
1383            .to_owned(),
1384        output_tokens: produced,
1385    })
1386}
1387
1388/// Preflight: which configured agents are not runnable here?
1389pub fn missing_programs(specs: &[AgentSpec]) -> Vec<String> {
1390    let mut missing = Vec::new();
1391    for s in specs {
1392        let program = match s.kind {
1393            AgentKind::Command => s.command.first().map(String::as_str),
1394            other => other.program(),
1395        };
1396        if let Some(p) = program
1397            && !crate::config::which(p)
1398            && !Path::new(p).is_file()
1399            && !missing.iter().any(|m: &String| m == p)
1400        {
1401            missing.push(p.to_owned());
1402        }
1403    }
1404    missing
1405}
1406
1407/// Absolute path of a run's artifact directory.
1408pub fn artifacts_dir(run_dir: &Path) -> PathBuf {
1409    run_dir.join("artifacts")
1410}
1411
1412/// Can this agent's CLI actually be run on this machine?
1413pub fn installed(spec: &AgentSpec) -> bool {
1414    // A `command` agent has no program of its own to look for - its argv is the
1415    // operator's, and they are the authority on whether it runs.
1416    spec.kind.program().is_none_or(crate::config::which)
1417}
1418
1419/// Choose the agent for a seat that stands alone rather than rotating through
1420/// the roster: [`crate::talk`]'s standing conversation, [`crate::bump`]'s
1421/// release-bump decision, or anything else that needs one agent picked once
1422/// rather than a panel filled in.
1423///
1424/// `available` is a parameter rather than a call to [`installed`] so the order
1425/// below is assertable on a machine with none of these CLIs installed, which is
1426/// every CI runner.
1427///
1428/// The order, and why:
1429///
1430/// 1. An explicit id always wins, and is an error rather than a fallback when
1431///    it is unusable. Naming a seat has a reason, and silently substituting a
1432///    different model would waste whatever that reason was.
1433/// 2. Otherwise a [`AgentKind::Claude`] seat, ahead of the roster order: it is
1434///    the only one of the three CLIs magi can address before the first turn
1435///    (see this module's own doc on session mechanics), which matters most for
1436///    a conversation that opens with nothing typed yet.
1437/// 3. Otherwise the first runnable agent in roster order, because the roster
1438///    order is the operator's own stated preference and magi has nothing
1439///    better to go on.
1440pub fn pick(
1441    agents: &[AgentSpec],
1442    want: Option<&str>,
1443    available: &dyn Fn(&AgentSpec) -> bool,
1444) -> Result<AgentSpec> {
1445    if let Some(id) = want {
1446        let spec = agents
1447            .iter()
1448            .find(|a| a.id == id)
1449            .with_context(|| format!("no agent `{id}` in the roster; it has {}", ids(agents)))?;
1450        if !available(spec) {
1451            bail!(
1452                "agent `{}` needs `{}` on PATH; install it or pass a different \
1453                 --agent",
1454                spec.id,
1455                spec.kind.program().unwrap_or("its command")
1456            );
1457        }
1458        return Ok(spec.clone());
1459    }
1460
1461    if agents.is_empty() {
1462        bail!(
1463            "the agent roster is empty, so there is nobody to ask: install one \
1464             of claude, opencode or agy - magi derives a roster from what is on \
1465             PATH - or add an [[agents]] entry to magi.toml."
1466        );
1467    }
1468
1469    if let Some(spec) = agents
1470        .iter()
1471        .find(|a| a.kind == AgentKind::Claude && available(a))
1472    {
1473        return Ok(spec.clone());
1474    }
1475
1476    agents
1477        .iter()
1478        .find(|a| available(a))
1479        .cloned()
1480        .with_context(|| {
1481            let missing = agents
1482                .iter()
1483                .filter_map(|a| a.kind.program())
1484                .collect::<Vec<_>>()
1485                .join(", ");
1486            format!(
1487                "no agent in the roster can be run here: install one of \
1488                 {missing}, or add an [[agents]] entry to magi.toml for a CLI \
1489                 you do have"
1490            )
1491        })
1492}
1493
1494/// Resolve a standalone seat's agent choice into the ordered chain to try.
1495///
1496/// Unset (or empty) is [`pick`]'s default order, one entry. Otherwise each
1497/// named id is kept at its first appearance only - that is what bounds a
1498/// chain to one attempt per agent - and an id that is not in the roster or
1499/// not installed is skipped with a warning rather than failing the chain.
1500/// Errors, naming `role`, only when nothing resolves.
1501pub fn pick_chain(
1502    agents: &[AgentSpec],
1503    choice: Option<&AgentChoice>,
1504    available: &dyn Fn(&AgentSpec) -> bool,
1505    role: &str,
1506) -> Result<Vec<AgentSpec>> {
1507    let wanted = choice.map(AgentChoice::ids).unwrap_or_default();
1508    if wanted.is_empty() {
1509        return pick(agents, None, available)
1510            .map(|s| vec![s])
1511            .with_context(|| format!("choose an agent for the {role} role"));
1512    }
1513    let mut chain: Vec<AgentSpec> = Vec::new();
1514    for id in wanted {
1515        if chain.iter().any(|s| s.id == id) {
1516            continue;
1517        }
1518        match pick(agents, Some(id), available) {
1519            Ok(spec) => chain.push(spec),
1520            Err(e) => tracing::warn!("[roles] {role}: skipping `{id}`: {e:#}"),
1521        }
1522    }
1523    if chain.is_empty() {
1524        bail!(
1525            "[roles] {role} names no agent that can run here; the roster has {}",
1526            ids(agents)
1527        );
1528    }
1529    Ok(chain)
1530}
1531
1532/// Should a chain move on to its next agent after this call?
1533///
1534/// The one place that decides it: an error, a quota hit (judged apart from
1535/// `usable`, since a CLI can exit 0 with a quota message), or an unusable
1536/// answer such as a timeout or empty reply.
1537pub fn chain_advances(outcome: &Result<AgentOutput>) -> bool {
1538    match outcome {
1539        Err(_) => true,
1540        Ok(out) => output_advances(out),
1541    }
1542}
1543
1544/// [`chain_advances`] for a call that did return an output: the `Ok` half,
1545/// shared with callers that classify the call into their own outcome type
1546/// first (the graph's fixer chain).
1547pub fn output_advances(out: &AgentOutput) -> bool {
1548    out.quota_exhausted() || !out.usable()
1549}
1550
1551fn ids(agents: &[AgentSpec]) -> String {
1552    if agents.is_empty() {
1553        return "no agents at all".to_owned();
1554    }
1555    agents
1556        .iter()
1557        .map(|a| a.id.clone())
1558        .collect::<Vec<_>>()
1559        .join(", ")
1560}
1561
1562#[cfg(test)]
1563mod tests {
1564    use super::*;
1565
1566    const COMMAND_HELPER_MODE: &str = "MAGI_TEST_COMMAND_HELPER_MODE";
1567
1568    /// Test-only command agent implemented by this test binary itself. Unlike
1569    /// `echo` and `sleep`, it is available wherever the Rust tests run.
1570    fn command_helper(mode: &str) -> AgentSpec {
1571        AgentSpec {
1572            id: "helper".to_owned(),
1573            kind: AgentKind::Command,
1574            model: None,
1575            command: vec![
1576                std::env::current_exe()
1577                    .expect("locate test helper")
1578                    .to_string_lossy()
1579                    .into_owned(),
1580                "--exact".to_owned(),
1581                "agent::tests::command_agent_test_helper".to_owned(),
1582                "--nocapture".to_owned(),
1583            ],
1584            extra_args: Vec::new(),
1585            env: BTreeMap::from([(COMMAND_HELPER_MODE.to_owned(), mode.to_owned())]),
1586            prompt_delivery: None,
1587        }
1588    }
1589
1590    #[test]
1591    fn command_agent_test_helper() {
1592        match std::env::var(COMMAND_HELPER_MODE).as_deref() {
1593            Ok("reply") => println!("hello {}", std::env::var("MAGI_SEAT").unwrap()),
1594            Ok("cache") => println!("{}", std::env::var("CARGO_TARGET_DIR").unwrap()),
1595            Ok("no-cache") => println!(
1596                "{}",
1597                std::env::var("CARGO_TARGET_DIR").unwrap_or_else(|_| "ABSENT".to_owned())
1598            ),
1599            Ok("ignore-stdin") => println!("done"),
1600            Ok("chatty-sleep") => {
1601                println!("i-said-something");
1602                std::thread::sleep(Duration::from_secs(30));
1603            }
1604            Ok("sleep") => std::thread::sleep(Duration::from_secs(30)),
1605            Ok(other) => panic!("unknown command helper mode {other}"),
1606            Err(_) => {}
1607        }
1608    }
1609
1610    fn spec(kind: AgentKind, model: Option<&str>) -> AgentSpec {
1611        AgentSpec {
1612            id: "a".to_owned(),
1613            kind,
1614            model: model.map(str::to_owned),
1615            command: vec!["echo".to_owned(), "{label}".to_owned()],
1616            extra_args: Vec::new(),
1617            env: BTreeMap::new(),
1618            prompt_delivery: None,
1619        }
1620    }
1621
1622    fn inv<'a>(cwd: &'a Path, art: &'a Path, allow_write: bool) -> Invocation<'a> {
1623        Invocation {
1624            cwd,
1625            prompt: "do the thing",
1626            timeout: Duration::from_secs(900),
1627            allow_write,
1628            unsandboxed: false,
1629            sessions: true,
1630            artifacts: art,
1631            stem: "t",
1632            run: "test-run",
1633            node: "test",
1634            cache_dir: None,
1635            attachments: &[],
1636            writable: &[],
1637        }
1638    }
1639
1640    fn plan_for(kind: AgentKind, seat: &SeatState, allow_write: bool) -> Plan {
1641        build_command(
1642            &spec(kind, None),
1643            seat,
1644            &inv(Path::new("."), Path::new("/art"), allow_write),
1645            Path::new("/art/p.md"),
1646        )
1647        .unwrap()
1648    }
1649
1650    #[test]
1651    fn claude_mints_then_resumes_the_same_uuid() {
1652        let mut seat = SeatState::new("judge-1", "a", 7);
1653        let uuid = seat.claude_session.clone().unwrap();
1654        let first = plan_for(AgentKind::Claude, &seat, true);
1655        assert!(first.argv.windows(2).any(|w| w == ["--session-id", &uuid]));
1656        assert!(!first.argv.iter().any(|a| a == "--resume"));
1657
1658        seat.turns = 1;
1659        let second = plan_for(AgentKind::Claude, &seat, true);
1660        assert!(second.argv.windows(2).any(|w| w == ["--resume", &uuid]));
1661        assert!(!second.argv.iter().any(|a| a == "--session-id"));
1662    }
1663
1664    #[test]
1665    fn read_only_seats_cannot_edit() {
1666        let seat = SeatState::new("judge-1", "a", 7);
1667        let claude = plan_for(AgentKind::Claude, &seat, false);
1668        assert!(claude.argv.iter().any(|a| a == "--disallowed-tools"));
1669        assert!(
1670            !plan_for(AgentKind::Claude, &seat, true)
1671                .argv
1672                .iter()
1673                .any(|a| a == "--disallowed-tools")
1674        );
1675
1676        let agy = plan_for(AgentKind::Antigravity, &seat, false);
1677        assert!(agy.argv.windows(2).any(|w| w == ["--mode", "plan"]));
1678        assert!(
1679            !agy.argv
1680                .iter()
1681                .any(|a| a == "--dangerously-skip-permissions")
1682        );
1683        let agy_rw = plan_for(AgentKind::Antigravity, &seat, true);
1684        assert!(
1685            agy_rw
1686                .argv
1687                .windows(2)
1688                .any(|w| w == ["--mode", "accept-edits"])
1689        );
1690        assert!(
1691            agy_rw
1692                .argv
1693                .iter()
1694                .any(|a| a == "--dangerously-skip-permissions")
1695        );
1696        // agy is pointed at its prompt with its own `@<path>` syntax, not with
1697        // prose asking it to read a file. Measured on one trivial task: 17s
1698        // against 73s, because prose costs a tool round-trip before the model
1699        // has even seen its instructions. It is also the form rvpm proved.
1700        let agy_prompt = agy_rw
1701            .argv
1702            .iter()
1703            .position(|a| a == "-p")
1704            .map(|i| agy_rw.argv[i + 1].clone())
1705            .expect("agy takes its prompt with -p");
1706        assert!(
1707            agy_prompt.starts_with('@'),
1708            "agy must get a file reference, got {agy_prompt:?}"
1709        );
1710        assert!(
1711            !agy_prompt.contains("Read the file at"),
1712            "the prose pointer is for CLIs with no file syntax"
1713        );
1714
1715        // opencode is the exception: `--auto` also gates reads, so withholding
1716        // it silently drops the seat out of the panel. Verified against the CLI
1717        // — a read-only judge failed with "the user rejected permission to use
1718        // this specific tool call" while trying to open its own prompt.
1719        for allow_write in [false, true] {
1720            assert!(
1721                plan_for(AgentKind::Opencode, &seat, allow_write)
1722                    .argv
1723                    .iter()
1724                    .any(|a| a == "--auto"),
1725                "opencode needs --auto even to read (allow_write = {allow_write})"
1726            );
1727        }
1728    }
1729
1730    /// The three things about `codex exec` that were established by hand and
1731    /// that a rewrite would silently get wrong.
1732    #[test]
1733    fn codex_gets_extra_writable_roots_only_when_it_may_write() {
1734        let seat = SeatState::new("deputy-x", "a", 7);
1735        let roots = [PathBuf::from("/data/questions")];
1736        let mk = |allow_write: bool| {
1737            let mut i = inv(Path::new("."), Path::new("/art"), allow_write);
1738            i.writable = &roots;
1739            build_command(
1740                &spec(AgentKind::Codex, None),
1741                &seat,
1742                &i,
1743                Path::new("/art/p.md"),
1744            )
1745            .unwrap()
1746        };
1747        let want = "sandbox_workspace_write.writable_roots=[\"/data/questions\"]";
1748        assert!(mk(true).argv.windows(2).any(|w| w == ["-c", want]));
1749        assert!(!mk(false).argv.iter().any(|a| a.contains("writable_roots")));
1750    }
1751
1752    #[test]
1753    fn codex_talk_turn_bypasses_sandbox_only_when_opted_in() {
1754        const BYPASS: &str = "--dangerously-bypass-approvals-and-sandbox";
1755        let mut seat = SeatState::new("talk", "a", 7);
1756        let build = |seat: &SeatState, allow_write: bool, unsandboxed: bool| {
1757            let mut i = inv(Path::new("."), Path::new("/art"), allow_write);
1758            i.unsandboxed = unsandboxed;
1759            build_command(
1760                &spec(AgentKind::Codex, None),
1761                seat,
1762                &i,
1763                Path::new("/art/p.md"),
1764            )
1765            .unwrap()
1766        };
1767
1768        // (a) both set: bypass, and none of the sandbox flags.
1769        let first = build(&seat, true, true);
1770        assert!(first.argv.iter().any(|a| a == BYPASS));
1771        assert!(!first.argv.iter().any(|a| a == "--sandbox"));
1772        assert!(!first.argv.iter().any(|a| a.starts_with("approval_policy")));
1773
1774        // (b) the flag precedes the `resume` subcommand.
1775        seat.captured_session = Some("thread-1".to_owned());
1776        seat.turns = 1;
1777        let resumed = build(&seat, true, true);
1778        let at = |p: &Plan, x: &str| p.argv.iter().position(|a| a == x).unwrap();
1779        assert!(at(&resumed, BYPASS) < at(&resumed, "resume"));
1780
1781        // (c) without write access the read-only sandbox stays.
1782        let ro = build(&seat, false, true);
1783        assert!(!ro.argv.iter().any(|a| a == BYPASS));
1784        assert!(ro.argv.windows(2).any(|w| w == ["--sandbox", "read-only"]));
1785    }
1786
1787    #[test]
1788    fn codex_is_sandboxed_reads_stdin_and_puts_resume_last() {
1789        let mut seat = SeatState::new("judge-1", "a", 7);
1790
1791        // 1. Read-only is enforced by the CLI, not by the prompt - the only
1792        //    roster member for which that is true - and nothing ever asks for
1793        //    the bypass.
1794        let ro = plan_for(AgentKind::Codex, &seat, false);
1795        assert!(ro.argv.windows(2).any(|w| w == ["--sandbox", "read-only"]));
1796        let rw = plan_for(AgentKind::Codex, &seat, true);
1797        assert!(
1798            rw.argv
1799                .windows(2)
1800                .any(|w| w == ["--sandbox", "workspace-write"])
1801        );
1802        // `unsandboxed` without a talk turn's opt-in never reaches an unattended
1803        // seat, and on its own it never lifts a read-only seat either.
1804        let mut ro_flagged = inv(Path::new("."), Path::new("/art"), false);
1805        ro_flagged.unsandboxed = true;
1806        let ro_flagged = build_command(
1807            &spec(AgentKind::Codex, None),
1808            &seat,
1809            &ro_flagged,
1810            Path::new("/art/p.md"),
1811        )
1812        .unwrap();
1813        assert!(
1814            ro_flagged
1815                .argv
1816                .windows(2)
1817                .any(|w| w == ["--sandbox", "read-only"])
1818        );
1819        for p in [&ro, &rw, &ro_flagged] {
1820            assert!(
1821                !p.argv
1822                    .iter()
1823                    .any(|a| a == "--dangerously-bypass-approvals-and-sandbox"),
1824                "the bypass defeats the only enforced read-only mode we have"
1825            );
1826            // Nobody is watching to approve anything.
1827            assert!(
1828                p.argv
1829                    .windows(2)
1830                    .any(|w| w == ["-c", "approval_policy=\"never\""]),
1831                "an unattended seat that asks for approval blocks until timeout"
1832            );
1833        }
1834
1835        // 2. The prompt arrives on stdin, and `-` is what makes codex read it.
1836        assert_eq!(ro.stdin.as_deref(), Some("do the thing"));
1837        assert_eq!(
1838            ro.argv.last().map(String::as_str),
1839            Some("-"),
1840            "without the `-` argument codex waits for a prompt it never gets"
1841        );
1842
1843        // 3. `resume` is a subcommand and rejects the options above when they
1844        //    follow it, so it has to be emitted after all of them - and only
1845        //    once the CLI has reported a thread id.
1846        seat.turns = 1;
1847        assert!(!has_session(AgentKind::Codex, &seat, true));
1848        assert!(
1849            !plan_for(AgentKind::Codex, &seat, true)
1850                .argv
1851                .iter()
1852                .any(|a| a == "resume")
1853        );
1854        seat.captured_session = Some("01a07440-4545-7492-85c1-024e3259a90a".to_owned());
1855        let resumed = plan_for(AgentKind::Codex, &seat, true);
1856        let at = resumed
1857            .argv
1858            .iter()
1859            .position(|a| a == "resume")
1860            .expect("resumes by subcommand");
1861        assert_eq!(resumed.argv[at + 1], "01a07440-4545-7492-85c1-024e3259a90a");
1862        assert!(
1863            resumed.argv[..at].iter().any(|a| a == "--sandbox"),
1864            "every option precedes the subcommand"
1865        );
1866        assert_eq!(resumed.argv.last().map(String::as_str), Some("-"));
1867    }
1868
1869    /// The three things about `omp -p --mode=json` that were established by
1870    /// hand against omp 18.1.19 and that a rewrite would silently get wrong.
1871    #[test]
1872    fn omp_reads_stdin_auto_approves_and_resumes_by_id() {
1873        let mut seat = SeatState::new("review-1", "a", 7);
1874
1875        // 1. Print mode plus JSON, and the prompt on stdin: a judging prompt
1876        //    carrying three patches is past the Windows argv cap, so argv
1877        //    delivery is not an option for every node.
1878        let first = plan_for(AgentKind::Omp, &seat, false);
1879        assert!(first.argv.iter().any(|a| a == "-p"));
1880        assert!(first.argv.iter().any(|a| a == "--mode=json"));
1881        assert_eq!(first.stdin.as_deref(), Some("do the thing"));
1882        assert!(
1883            !first.argv.iter().any(|a| a == "do the thing"),
1884            "the prompt reached argv, where Windows caps it"
1885        );
1886
1887        // 2. `--auto-approve` is required (an unattended seat that stops to ask
1888        //    blocks until its node timeout kills it), and it is the *only*
1889        //    permission flag: omp has no read-only mode, so the bypass flag
1890        //    that would throw away codex's one enforced guarantee must never
1891        //    appear here either.
1892        for allow_write in [false, true] {
1893            let p = plan_for(AgentKind::Omp, &seat, allow_write);
1894            assert!(
1895                p.argv.iter().any(|a| a == "--auto-approve"),
1896                "omp needs --auto-approve even to read (allow_write = {allow_write})"
1897            );
1898            assert!(
1899                !p.argv
1900                    .iter()
1901                    .any(|a| a == "--dangerously-bypass-approvals-and-sandbox"),
1902                "nothing ever asks for the bypass"
1903            );
1904        }
1905
1906        // 3. The id omp reports is the only resume token - magi cannot mint it
1907        //    up front, so a seat resumes only once a turn has reported one.
1908        seat.turns = 1;
1909        assert!(!has_session(AgentKind::Omp, &seat, true));
1910        assert!(
1911            !plan_for(AgentKind::Omp, &seat, true)
1912                .argv
1913                .iter()
1914                .any(|a| a == "--resume")
1915        );
1916        seat.captured_session = Some("01a09fe9-4e31-7226-85b3-fda6f46689d5".to_owned());
1917        let resumed = plan_for(AgentKind::Omp, &seat, true);
1918        assert!(
1919            resumed
1920                .argv
1921                .windows(2)
1922                .any(|w| w == ["--resume", "01a09fe9-4e31-7226-85b3-fda6f46689d5"]),
1923            "a captured id is what makes the next turn a resume"
1924        );
1925        // `--continue` opens a *new* session instead of the stored one, which
1926        // would silently drop the seat's memory.
1927        assert!(!resumed.argv.iter().any(|a| a == "--continue"));
1928        // stdin still carries the prompt on a resumed turn.
1929        assert_eq!(resumed.stdin.as_deref(), Some("do the thing"));
1930    }
1931
1932    /// The extraction trap that cost three complete reviews when it was done by
1933    /// hand: a turn that ends on a tool call emits **no** `agent_end` line, so
1934    /// keying on `agent_end` finds nothing and the seat reads as one that
1935    /// produced no answer at all.
1936    #[test]
1937    fn omp_takes_the_answer_without_an_agent_end_line() {
1938        let stream = concat!(
1939            r#"{"type":"session","version":3,"id":"01a09fe9-4e31-7226-85b3-fda6f46689d5","cwd":"C:\\w"}"#,
1940            "\n",
1941            r#"{"type":"agent_start"}"#,
1942            "\n",
1943            r#"{"type":"turn_start"}"#,
1944            "\n",
1945            r#"{"type":"message_update","assistantMessageEvent":{"type":"text_delta","contentIndex":1,"delta":"."}}"#,
1946            "\n",
1947            r#"{"type":"message_end","message":{"role":"assistant","content":[{"type":"thinking","thinking":"checking"},{"type":"text","text":"."}]}}"#,
1948            "\n",
1949            r#"{"type":"turn_end","message":{"role":"assistant","content":[{"type":"thinking","thinking":"done"},{"type":"text","text":"{\"vote\":\"approve\"}"}]}}"#,
1950            "\n",
1951        );
1952        let out = extract(AgentKind::Omp, stream);
1953        assert_eq!(
1954            out.text, "{\"vote\":\"approve\"}",
1955            "the last assistant text block is the answer even with no agent_end"
1956        );
1957        assert_eq!(
1958            out.session.as_deref(),
1959            Some("01a09fe9-4e31-7226-85b3-fda6f46689d5")
1960        );
1961    }
1962
1963    /// A stream that *does* carry `agent_end` walks the whole thread, and the
1964    /// last non-empty assistant text still wins over the tool-loop narration
1965    /// that came before it.
1966    #[test]
1967    fn omp_walks_agent_end_and_ignores_tool_loop_narration() {
1968        let stream = concat!(
1969            r#"{"type":"session","version":3,"id":"s1"}"#,
1970            "\n",
1971            "{\"type\":\"agent_end\",\"messages\":[{\"role\":\"user\",\"content\":[{\"type\":\"text\",\"text\":\"review this\"}]},{\"role\":\"assistant\",\"content\":[{\"type\":\"text\",\"text\":\"Looking at the diff…\"}]},{\"role\":\"assistant\",\"content\":[{\"type\":\"thinking\",\"thinking\":\"…\"},{\"type\":\"text\",\"text\":\"## 判定\\n\\n問題ありません。\"}]}]}",
1972            "\n",
1973        );
1974        let out = extract(AgentKind::Omp, stream);
1975        assert_eq!(
1976            out.text, "## 判定\n\n問題ありません。",
1977            "the narration is not the answer, and non-ASCII survives intact"
1978        );
1979        assert_eq!(out.session.as_deref(), Some("s1"));
1980    }
1981
1982    /// A line that is not JSON - a CLI warning, a half-written line - is
1983    /// skipped rather than becoming the answer.
1984    #[test]
1985    fn omp_skips_non_json_lines() {
1986        let stream = concat!(
1987            "Warning: some omp notice\n",
1988            r#"{"type":"session","version":3,"id":"s2"}"#,
1989            "\n",
1990            r#"{"type":"message_end","message":{"role":"assistant","content":[{"type":"text","text":"the answer"}]}}"#,
1991            "\n",
1992            "trailing junk",
1993            "\n",
1994        );
1995        let out = extract(AgentKind::Omp, stream);
1996        assert_eq!(out.text, "the answer");
1997        assert_eq!(out.session.as_deref(), Some("s2"));
1998    }
1999
2000    /// A real `codex exec --json` stream, tracing prefix included.
2001    #[test]
2002    fn codex_takes_the_last_agent_message_and_the_thread_id() {
2003        let stream = concat!(
2004            "2026-09-06T01:05:49.394445Z ERROR codex_models_manager: failed to load models cache\n",
2005            r#"{"type":"thread.started","thread_id":"01a07440-4545-7492-85c1-024e3259a90a"}"#,
2006            "\n",
2007            r#"{"type":"turn.started"}"#,
2008            "\n",
2009            r#"{"type":"item.completed","item":{"id":"item_0","type":"agent_message","text":"Looking into it."}}"#,
2010            "\n",
2011            r#"{"type":"item.completed","item":{"id":"item_1","type":"command_execution","text":"cargo test"}}"#,
2012            "\n",
2013            r#"{"type":"item.completed","item":{"id":"item_2","type":"agent_message","text":"{\"verdict\": \"ok\"}"}}"#,
2014            "\n",
2015            r#"{"type":"turn.completed","usage":{"input_tokens":17137}}"#,
2016            "\n",
2017        );
2018        let out = extract(AgentKind::Codex, stream);
2019        assert_eq!(
2020            out.text, "{\"verdict\": \"ok\"}",
2021            "the last agent message is the answer; earlier ones narrate"
2022        );
2023        assert_eq!(
2024            out.session.as_deref(),
2025            Some("01a07440-4545-7492-85c1-024e3259a90a")
2026        );
2027        assert_eq!(out.status.as_deref(), Some("success"));
2028
2029        let failed = concat!(
2030            r#"{"type":"thread.started","thread_id":"t1"}"#,
2031            "\n",
2032            r#"{"type":"turn.failed","error":{"message":"nope"}}"#,
2033            "\n",
2034        );
2035        assert_eq!(
2036            extract(AgentKind::Codex, failed).status.as_deref(),
2037            Some("error")
2038        );
2039    }
2040
2041    #[test]
2042    fn captured_sessions_resume_only_once_reported() {
2043        let mut seat = SeatState::new("impl-A", "a", 7);
2044        seat.turns = 1;
2045        for kind in [AgentKind::Opencode, AgentKind::Antigravity] {
2046            assert!(!has_session(kind, &seat, true));
2047            let p = plan_for(kind, &seat, true);
2048            assert!(!p.argv.iter().any(|a| a == "-s" || a == "--conversation"));
2049        }
2050
2051        seat.captured_session = Some("sid".to_owned());
2052        assert!(has_session(AgentKind::Opencode, &seat, true));
2053        assert!(
2054            plan_for(AgentKind::Opencode, &seat, true)
2055                .argv
2056                .windows(2)
2057                .any(|w| w == ["-s", "sid"])
2058        );
2059        assert!(
2060            plan_for(AgentKind::Antigravity, &seat, true)
2061                .argv
2062                .windows(2)
2063                .any(|w| w == ["--conversation", "sid"])
2064        );
2065    }
2066
2067    #[test]
2068    fn sessions_disabled_never_resumes() {
2069        let mut seat = SeatState::new("impl-A", "a", 7);
2070        seat.turns = 3;
2071        seat.captured_session = Some("sid".to_owned());
2072        for kind in [
2073            AgentKind::Claude,
2074            AgentKind::Opencode,
2075            AgentKind::Antigravity,
2076        ] {
2077            assert!(!has_session(kind, &seat, false));
2078        }
2079    }
2080
2081    #[test]
2082    fn long_prompts_never_reach_argv_for_file_delivery_clis() {
2083        let seat = SeatState::new("judge-1", "a", 7);
2084        for kind in [AgentKind::Opencode, AgentKind::Antigravity] {
2085            let p = plan_for(kind, &seat, false);
2086            assert!(
2087                p.argv.iter().all(|a| a != "do the thing"),
2088                "{kind:?} put the prompt on the command line"
2089            );
2090            assert!(p.argv.iter().any(|a| a.contains("/art/p.md")));
2091        }
2092        // agy has no text stdin, so its `-p` must always carry something.
2093        let p = plan_for(AgentKind::Antigravity, &seat, false);
2094        let at = p.argv.iter().position(|a| a == "-p").unwrap();
2095        assert!(p.argv.get(at + 1).is_some_and(|v| v.contains("p.md")));
2096        assert!(p.stdin.is_none());
2097    }
2098
2099    #[test]
2100    fn agy_print_timeout_tracks_the_node_budget() {
2101        let seat = SeatState::new("impl-A", "a", 7);
2102        let p = build_command(
2103            &spec(AgentKind::Antigravity, None),
2104            &seat,
2105            &Invocation {
2106                cwd: Path::new("."),
2107                prompt: "p",
2108                timeout: Duration::from_secs(3600),
2109                allow_write: true,
2110                unsandboxed: false,
2111                sessions: true,
2112                artifacts: Path::new("/art"),
2113                stem: "t",
2114                run: "test-run",
2115                node: "test",
2116                cache_dir: None,
2117                attachments: &[],
2118                writable: &[],
2119            },
2120            Path::new("/art/p.md"),
2121        )
2122        .unwrap();
2123        assert!(p.argv.windows(2).any(|w| w == ["--print-timeout", "3600s"]));
2124    }
2125
2126    /// `--add-dir` is what lets antigravity open a file outside the
2127    /// worktree at all. Today that only happens when the delivery mode is
2128    /// already `File`, but an attachment can arrive on a seat whose delivery
2129    /// is `Stdin` or `Argv` (an explicit `prompt_delivery` override), and the
2130    /// image still lives outside `cwd` - so the flag has to widen for that
2131    /// reason too, independent of how the prompt itself is delivered.
2132    #[test]
2133    fn attachments_widen_antigravitys_add_dir_even_off_file_delivery() {
2134        let mut s = spec(AgentKind::Antigravity, None);
2135        s.prompt_delivery = Some(Delivery::Argv);
2136        let seat = SeatState::new("talk", "a", 7);
2137        let atts = [PathBuf::from("/art/attachments/abc.png")];
2138
2139        let without = build_command(
2140            &s,
2141            &seat,
2142            &Invocation {
2143                attachments: &[],
2144                writable: &[],
2145                ..inv(Path::new("."), Path::new("/art"), true)
2146            },
2147            Path::new("/art/p.md"),
2148        )
2149        .unwrap();
2150        assert!(
2151            !without.argv.iter().any(|a| a == "--add-dir"),
2152            "no attachment, no reason to widen the sandbox: {without:?}"
2153        );
2154
2155        let with = build_command(
2156            &s,
2157            &seat,
2158            &Invocation {
2159                attachments: &atts,
2160                writable: &[],
2161                ..inv(Path::new("."), Path::new("/art"), true)
2162            },
2163            Path::new("/art/p.md"),
2164        )
2165        .unwrap();
2166        assert!(
2167            with.argv.windows(2).any(|w| w == ["--add-dir", "/art"]),
2168            "an attachment outside cwd must widen the sandbox even off File delivery: {with:?}"
2169        );
2170    }
2171
2172    /// A chat derived from another one (`chat::derived_background`) can pass
2173    /// `turn` attachment paths that live under the *source* conversation's
2174    /// own artifacts dir, not this invocation's `artifacts`. A single
2175    /// `--add-dir` for `inv.artifacts` alone would leave those unreadable, so
2176    /// each attachment directory outside it must get its own grant.
2177    #[test]
2178    fn an_inherited_attachment_outside_this_conversations_artifacts_dir_gets_its_own_add_dir() {
2179        let seat = SeatState::new("plan", "a", 7);
2180        let atts = [
2181            PathBuf::from("/art/attachments/own.png"),
2182            PathBuf::from("/other-chat/attachments/inherited.png"),
2183        ];
2184
2185        let p = build_command(
2186            &spec(AgentKind::Antigravity, None),
2187            &seat,
2188            &Invocation {
2189                attachments: &atts,
2190                writable: &[],
2191                ..inv(Path::new("."), Path::new("/art"), true)
2192            },
2193            Path::new("/art/p.md"),
2194        )
2195        .unwrap();
2196
2197        assert!(
2198            p.argv.windows(2).any(|w| w == ["--add-dir", "/art"]),
2199            "this conversation's own artifacts dir must still be granted: {p:?}"
2200        );
2201        assert!(
2202            p.argv
2203                .windows(2)
2204                .any(|w| w == ["--add-dir", "/other-chat/attachments"]),
2205            "the inherited attachment's own directory must be granted too: {p:?}"
2206        );
2207    }
2208
2209    #[test]
2210    fn command_agents_get_placeholders_substituted() {
2211        let seat = SeatState::new("impl-A", "a", 7);
2212        let p = plan_for(AgentKind::Command, &seat, true);
2213        assert_eq!(p.argv[0], "echo");
2214        assert_eq!(p.argv[1], "impl-A");
2215        assert_eq!(p.stdin.as_deref(), Some("do the thing"));
2216    }
2217
2218    #[test]
2219    fn claude_rate_limit_is_detected_and_reset_read_when_present() {
2220        // The exact shape observed in the wild (run 20260831-031005-ae94).
2221        let stdout = r#"{"is_error": true, "terminal_reason": "api_error",
2222                        "result": "You've hit your session limit · resets 4:50am (Asia/Tokyo)",
2223                        "session_id": "b8e928f1-754e-4bd3-86c5-0567763654e3"}"#;
2224        let out = extract(AgentKind::Claude, stdout);
2225        let quota = out.quota.as_ref().expect("rate limit must be detected");
2226        assert_eq!(
2227            quota.reset.as_deref(),
2228            Some("4:50am (Asia/Tokyo)"),
2229            "reset time read from the body"
2230        );
2231    }
2232
2233    #[test]
2234    fn claude_rate_limit_without_a_readable_reset_is_still_detected() {
2235        let out = extract(
2236            AgentKind::Claude,
2237            r#"{"is_error":true,"result":"session limit reached"}"#,
2238        );
2239        let quota = out.quota.expect("rate limit detected without a reset");
2240        assert!(quota.reset.is_none(), "unknown reset is kept as unknown");
2241    }
2242
2243    #[test]
2244    fn ordinary_failures_are_never_quota() {
2245        // A normal failed claude call (is_error with a different message).
2246        let claude_fail = extract(
2247            AgentKind::Claude,
2248            r#"{"is_error":true,"result":"account does not exist"}"#,
2249        );
2250        assert!(claude_fail.quota.is_none());
2251
2252        // A command agent that exits 1 with plain text.
2253        let cmd_fail = extract(AgentKind::Command, "boom");
2254        assert!(cmd_fail.quota.is_none());
2255
2256        // A successful call is not quota even if it mentions the phrase.
2257        let success = extract(
2258            AgentKind::Command,
2259            r#"{"is_error":false,"result":"session limit is fine"}"#,
2260        );
2261        assert!(success.quota.is_none());
2262    }
2263
2264    /// Minimal, anonymised shape of run 20260912-214939-b3bb's
2265    /// artifacts/impl-A.out: `command_execution` items reporting a paired
2266    /// test run as `1 passed, 1 failed` twice, while the final
2267    /// `agent_message` nevertheless claimed the target passed. The point of
2268    /// `CommandEvidence` is that this claim and the CLI's own structured
2269    /// record of what actually ran are now two separate things a caller can
2270    /// compare, rather than the prose being the only account available.
2271    #[test]
2272    fn codex_command_execution_events_are_captured_alongside_the_final_message() {
2273        let stream = concat!(
2274            r#"{"type":"thread.started","thread_id":"t1"}"#,
2275            "\n",
2276            r#"{"type":"item.completed","item":{"id":"item49","type":"command_execution","command":["bash","-lc","cargo test --test graph_cached_gate"],"exit_code":1,"aggregated_output":"test result: 1 passed; 1 failed"}}"#,
2277            "\n",
2278            r#"{"type":"item.completed","item":{"id":"item52","type":"command_execution","command":["bash","-lc","cargo test --test graph_cached_gate a_single_test"],"exit_code":0,"aggregated_output":"test result: 1 passed; 0 failed"}}"#,
2279            "\n",
2280            r#"{"type":"item.completed","item":{"id":"item99","type":"agent_message","text":"Both tests in the target pass."}}"#,
2281            "\n",
2282            r#"{"type":"turn.completed"}"#,
2283            "\n",
2284        );
2285        let out = extract(AgentKind::Codex, stream);
2286        assert_eq!(out.text, "Both tests in the target pass.");
2287        assert_eq!(out.commands.len(), 2, "{:?}", out.commands);
2288
2289        let paired = &out.commands[0];
2290        assert_eq!(paired.id, "item49");
2291        assert_eq!(paired.exit_code, Some(1));
2292        assert!(paired.description.contains("graph_cached_gate"));
2293        assert!(paired.result_summary.contains("1 failed"));
2294
2295        let solo = &out.commands[1];
2296        assert_eq!(solo.exit_code, Some(0));
2297
2298        // The structured evidence disagrees with the final prose - exactly
2299        // what a caller must be able to see instead of trusting the message
2300        // alone: the full target never passed in one command.
2301        assert!(
2302            out.commands
2303                .iter()
2304                .any(|c| c.exit_code != Some(0) && c.description.contains("graph_cached_gate")),
2305            "a failed run of the actual target must still be visible: {:?}",
2306            out.commands
2307        );
2308    }
2309
2310    #[test]
2311    fn command_agent_can_carry_the_claude_quota_shape() {
2312        let out = extract(
2313            AgentKind::Command,
2314            r#"{"is_error":true,"result":"You've hit your session limit · resets 1:00am (UTC)"}"#,
2315        );
2316        assert!(
2317            out.quota.is_some(),
2318            "a wrapper emitting the claude shape counts as quota"
2319        );
2320    }
2321
2322    #[test]
2323    fn claude_json_result_is_extracted() {
2324        let out = extract(
2325            AgentKind::Claude,
2326            r#"{"result":"all done","session_id":"abc","is_error":false}"#,
2327        );
2328        assert_eq!(out.text, "all done");
2329        assert_eq!(out.session.as_deref(), Some("abc"));
2330        assert_eq!(out.status.as_deref(), Some("success"));
2331    }
2332
2333    /// The shape behind fix-2 in run 20260912-114326-d3b8, minimal and
2334    /// anonymised: a Claude CLI turn that ended `subtype: success`,
2335    /// `is_error: false`, `terminal_reason: completed`, `stop_reason:
2336    /// end_turn` — every signal this crate reads as a clean CLI turn — while
2337    /// `result` is a progress update, not the report the fixer node needed,
2338    /// and no `FixReport` JSON is anywhere in it.
2339    ///
2340    /// `AgentOutput::usable()` (a CLI fact: exit 0, not timed out, non-empty
2341    /// text) must stay true here — that is the honest reading of what the
2342    /// CLI reported — while `verdict::extract_json` on the same text must
2343    /// fail. Conflating the two is exactly the bug this fixture reproduces:
2344    /// `run.json` recorded `fix.failed = "unparsable fix report: the reply
2345    /// contained no JSON object"` and moved straight to the next review round
2346    /// with no report ever recovered from that seat.
2347    #[test]
2348    fn a_clean_cli_turn_is_not_the_same_fact_as_the_nodes_own_work_being_done() {
2349        let stdout = r#"{"type":"result","subtype":"success","is_error":false,"terminal_reason":"completed","stop_reason":"end_turn","result":"I'll pause here until the `cargo make check` background run reports back.","session_id":"11111111-1111-1111-1111-111111111111"}"#;
2350        let out = extract(AgentKind::Claude, stdout);
2351        assert_eq!(out.status.as_deref(), Some("success"));
2352        assert!(out.quota.is_none());
2353        assert!(!out.text.trim().is_empty());
2354
2355        let agent_out = AgentOutput {
2356            text: out.text.clone(),
2357            exit_code: Some(0),
2358            timed_out: false,
2359            duration_ms: 500,
2360            artifacts: Vec::new(),
2361            quota: out.quota,
2362            dropped: out.dropped,
2363            commands: out.commands,
2364            context_tokens: None,
2365        };
2366        assert!(
2367            agent_out.usable(),
2368            "the CLI turn itself ended cleanly and must read as usable"
2369        );
2370        assert!(
2371            crate::verdict::extract_json::<crate::verdict::FixReport>(&agent_out.text).is_err(),
2372            "a clean CLI turn is not proof the node's own report ever arrived"
2373        );
2374    }
2375
2376    #[test]
2377    fn opencode_event_stream_is_concatenated() {
2378        let stream = concat!(
2379            r#"{"type":"step_start","sessionID":"ses_1","part":{"type":"step-start"}}"#,
2380            "\n",
2381            r#"{"type":"text","sessionID":"ses_1","part":{"type":"text","text":"first"}}"#,
2382            "\n",
2383            "garbage line\n",
2384            r#"{"type":"text","sessionID":"ses_1","part":{"type":"text","text":"second"}}"#,
2385            "\n"
2386        );
2387        let out = extract(AgentKind::Opencode, stream);
2388        assert_eq!(out.text, "first\nsecond");
2389        assert_eq!(out.session.as_deref(), Some("ses_1"));
2390    }
2391
2392    #[test]
2393    fn agy_json_survives_a_leading_warning_line() {
2394        let stdout = concat!(
2395            "warning: --mode plan has no effect while slash commands are disabled.\n",
2396            r#"{"conversation_id":"eaf2d00a","status":"SUCCESS","response":"persimmon\n"}"#,
2397            "\n"
2398        );
2399        let out = extract(AgentKind::Antigravity, stdout);
2400        assert_eq!(out.text, "persimmon");
2401        assert_eq!(out.session.as_deref(), Some("eaf2d00a"));
2402        assert_eq!(out.status.as_deref(), Some("SUCCESS"));
2403    }
2404
2405    /// Run 26c7's candidate B, verbatim from `artifacts/impl-B.out`.
2406    ///
2407    /// The seat read as an empty candidate. It was seven minutes of work and
2408    /// 14,267 output tokens, billed, that the CLI then declined to hand over.
2409    /// Five such candidates are why `agy` reads as 0 wins in 4 entries, and
2410    /// that number has twice been used to argue the seat out of the roster.
2411    const AGY_DROPPED: &str = concat!(
2412        r#"{"conversation_id":"36743d06-c0b3-4b79-9fa2-23869289d7b6","status":"ERROR","#,
2413        r#""response":"","error":"the connection to the agent was interrupted before "#,
2414        r#"the response finished: subscriber fell behind updates, stalled for 5s","#,
2415        r#""duration_seconds":431.1941803,"num_turns":1,"usage":{"input_tokens":260113,"#,
2416        r#""output_tokens":14267,"thinking_tokens":9695,"cache_read_tokens":2200925,"#,
2417        r#""total_tokens":274380}}"#
2418    );
2419
2420    /// Run 1798's review seat, from `review-1-1.out` / `.err`.
2421    const AGY_QUOTA_OUT: &str = concat!(
2422        r#"{"conversation_id":"323c3b5b-0000","status":"ERROR","response":"","#,
2423        r#""error":"Individual quota reached. Please upgrade your subscription to "#,
2424        r#"increase your limits. Resets in 1h2m49s.","duration_seconds":265.9,"#,
2425        r#""num_turns":2,"usage":{"input_tokens":1000,"output_tokens":50}}"#
2426    );
2427    const AGY_QUOTA_ERR: &str = concat!(
2428        "error: Individual quota reached. Resets in 1h2m49s.\n",
2429        r#"AGY_ERROR: {"short_error":"RESOURCE_EXHAUSTED (code 429): Individual quota "#,
2430        r#"reached.","status":"RESOURCE_EXHAUSTED","error_code":429,"code_kind":"http","#,
2431        r#""retryable":true}"#,
2432        "\n"
2433    );
2434
2435    #[test]
2436    fn agy_out_of_quota_is_a_quota_with_the_reset_hint() {
2437        let both = agy_quota(AGY_QUOTA_OUT, AGY_QUOTA_ERR).expect("both streams");
2438        assert_eq!(both.reset.as_deref(), Some("in 1h2m49s"));
2439        let stdout_only = agy_quota(AGY_QUOTA_OUT, "").expect("stdout alone");
2440        assert_eq!(stdout_only.reset.as_deref(), Some("in 1h2m49s"));
2441        // stderr alone: the hint is not in `short_error`, so none is carried.
2442        let stderr_only = agy_quota("not json", AGY_QUOTA_ERR).expect("stderr alone");
2443        assert!(stderr_only.reset.is_none());
2444    }
2445
2446    #[test]
2447    fn ordinary_agy_failures_are_not_a_quota() {
2448        assert!(agy_quota(AGY_DROPPED, "").is_none());
2449        assert!(agy_quota(r#"{"status":"ERROR","error":"boom"}"#, "").is_none());
2450        assert!(
2451            agy_quota(
2452                "",
2453                r#"AGY_ERROR: {"status":"RESOURCE_EXHAUSTED","error_code":500}"#
2454            )
2455            .is_none()
2456        );
2457        assert!(
2458            agy_quota(
2459                "",
2460                r#"AGY_ERROR: {"status":"UNAVAILABLE","error_code":429}"#
2461            )
2462            .is_none()
2463        );
2464        assert!(agy_quota(r#"{"status":"SUCCESS","response":"ok"}"#, "").is_none());
2465    }
2466
2467    #[test]
2468    fn a_cli_that_hangs_up_on_billed_work_is_not_an_agent_that_produced_nothing() {
2469        let out = extract(AgentKind::Antigravity, AGY_DROPPED);
2470        let dropped = out.dropped.expect("recognised as undelivered work");
2471        assert_eq!(dropped.output_tokens, 14267);
2472        assert!(
2473            dropped.why.contains("subscriber fell behind"),
2474            "the CLI's own words are kept for the record: {}",
2475            dropped.why
2476        );
2477        // And the conversation is still there to resume, which is the whole
2478        // reason this is worth re-asking where a quota is not.
2479        assert_eq!(
2480            out.session.as_deref(),
2481            Some("36743d06-c0b3-4b79-9fa2-23869289d7b6")
2482        );
2483        assert!(out.quota.is_none(), "a dropped stream is not a rate limit");
2484    }
2485
2486    #[test]
2487    fn an_error_with_nothing_produced_stays_an_ordinary_failure() {
2488        // No usage at all: the agent never got going, so there is nothing in
2489        // the conversation to resume and nothing was billed. Treating this as
2490        // undelivered work would buy a second call for no reason.
2491        let bare = r#"{"conversation_id":"c1","status":"ERROR","response":"","error":"boom"}"#;
2492        assert!(extract(AgentKind::Antigravity, bare).dropped.is_none());
2493
2494        // Produced tokens, but it did answer - so there is something to read
2495        // and the status is not our business.
2496        let answered = concat!(
2497            r#"{"conversation_id":"c2","status":"ERROR","response":"here it is","#,
2498            r#""usage":{"output_tokens":10}}"#
2499        );
2500        assert!(extract(AgentKind::Antigravity, answered).dropped.is_none());
2501
2502        // A success is a success.
2503        let ok = concat!(
2504            r#"{"conversation_id":"c3","status":"SUCCESS","response":"done","#,
2505            r#""usage":{"output_tokens":10}}"#
2506        );
2507        assert!(extract(AgentKind::Antigravity, ok).dropped.is_none());
2508    }
2509
2510    #[test]
2511    fn an_undelivered_output_is_not_usable_but_is_worth_asking_again() {
2512        let out = AgentOutput {
2513            text: String::new(),
2514            exit_code: Some(1),
2515            timed_out: false,
2516            duration_ms: 431_194,
2517            artifacts: Vec::new(),
2518            quota: None,
2519            dropped: Some(Dropped {
2520                why: "subscriber fell behind updates".to_owned(),
2521                output_tokens: 14267,
2522            }),
2523            commands: Vec::new(),
2524            context_tokens: None,
2525        };
2526        assert!(!out.usable());
2527        assert!(out.work_undelivered());
2528        // The distinction the retry policy rests on: a quota fails the same way
2529        // until it resets, an abandoned conversation can be picked up.
2530        assert!(!out.quota_exhausted());
2531    }
2532
2533    #[test]
2534    fn non_json_stdout_falls_back_to_raw_text() {
2535        let out = extract(AgentKind::Antigravity, "plain answer\n");
2536        assert_eq!(out.text, "plain answer");
2537        assert!(out.session.is_none());
2538    }
2539
2540    #[tokio::test]
2541    async fn command_agent_round_trip_writes_artifacts() {
2542        let dir = tempfile::tempdir().unwrap();
2543        let art = dir.path().join("artifacts");
2544        let mut seat = SeatState::new("impl-A", "a", 7);
2545        let s = command_helper("reply");
2546        let out = invoke(
2547            &s,
2548            &mut seat,
2549            &Invocation {
2550                cwd: dir.path(),
2551                prompt: "unused",
2552                timeout: Duration::from_secs(30),
2553                allow_write: true,
2554                unsandboxed: false,
2555                sessions: true,
2556                artifacts: &art,
2557                stem: "impl-A",
2558                run: "test-run",
2559                node: "test",
2560                cache_dir: None,
2561                attachments: &[],
2562                writable: &[],
2563            },
2564        )
2565        .await
2566        .unwrap();
2567        assert!(out.usable(), "{out:?}");
2568        assert!(out.text.contains("hello impl-A"), "{}", out.text);
2569        assert_eq!(seat.turns, 1);
2570        assert!(art.join("impl-A.prompt.md").is_file());
2571        assert!(art.join("impl-A.out").is_file());
2572    }
2573
2574    #[tokio::test]
2575    async fn the_invocation_cache_dir_reaches_the_seat_as_cargo_target_dir() {
2576        // The whole point of threading the cache path through `Invocation`:
2577        // the compile the agent pays for lands in the directory `verify` reads
2578        // back out of its rendered commands, so one cache has one prune.
2579        let dir = tempfile::tempdir().unwrap();
2580        let cache = dir.path().join("magi-cache");
2581        let mut seat = SeatState::new("impl-A", "a", 7);
2582        let s = command_helper("cache");
2583        let out = invoke(
2584            &s,
2585            &mut seat,
2586            &Invocation {
2587                cwd: dir.path(),
2588                prompt: "unused",
2589                timeout: Duration::from_secs(30),
2590                allow_write: true,
2591                unsandboxed: false,
2592                sessions: true,
2593                artifacts: &dir.path().join("artifacts"),
2594                stem: "cache",
2595                run: "test-run",
2596                node: "test",
2597                cache_dir: Some(&cache),
2598                attachments: &[],
2599                writable: &[],
2600            },
2601        )
2602        .await
2603        .unwrap();
2604        assert!(out.usable(), "{out:?}");
2605        assert!(
2606            out.text.contains(cache.to_string_lossy().as_ref()),
2607            "the seat must see CARGO_TARGET_DIR = the shared cache"
2608        );
2609    }
2610
2611    #[tokio::test]
2612    async fn cache_dir_none_strips_a_cargo_target_dir_inherited_from_this_process() {
2613        // `Command` inherits the parent's environment by default, so
2614        // `cache_dir: None` alone is not the same as a seat never seeing
2615        // `CARGO_TARGET_DIR` - it also has to be true when *this* process
2616        // (standing in for the real magi process, which normally does have
2617        // one set, from its own `[verify]` config) already has the variable
2618        // set. Simulating that is the only way to exercise the inheritance
2619        // path at all.
2620        let previous = std::env::var("CARGO_TARGET_DIR").ok();
2621        // SAFETY: this crate's tests run single-threaded
2622        // (`RUST_TEST_THREADS=1`); see
2623        // `updater::tests::env_kill_switch_semantics` for the same
2624        // reasoning applied to another process-global env var.
2625        unsafe {
2626            std::env::set_var("CARGO_TARGET_DIR", "/should/never/reach/a/read-only/seat");
2627        }
2628        let dir = tempfile::tempdir().unwrap();
2629        let mut seat = SeatState::new("review-1", "a", 7);
2630        let s = command_helper("no-cache");
2631        let result = invoke(
2632            &s,
2633            &mut seat,
2634            &Invocation {
2635                cwd: dir.path(),
2636                prompt: "unused",
2637                timeout: Duration::from_secs(30),
2638                allow_write: false,
2639                unsandboxed: false,
2640                sessions: true,
2641                artifacts: &dir.path().join("artifacts"),
2642                stem: "no-cache",
2643                run: "test-run",
2644                node: "test",
2645                cache_dir: None,
2646                attachments: &[],
2647                writable: &[],
2648            },
2649        )
2650        .await;
2651        // Restored before any assertion that could panic, so a failure here
2652        // never leaks a bogus `CARGO_TARGET_DIR` into whichever test runs
2653        // next in this same process.
2654        // SAFETY: see above.
2655        unsafe {
2656            match &previous {
2657                Some(v) => std::env::set_var("CARGO_TARGET_DIR", v),
2658                None => std::env::remove_var("CARGO_TARGET_DIR"),
2659            }
2660        }
2661        let out = result.unwrap();
2662        assert!(out.usable(), "{out:?}");
2663        assert!(
2664            out.text.contains("ABSENT"),
2665            "a read-only seat must never inherit the process's own CARGO_TARGET_DIR: {}",
2666            out.text
2667        );
2668    }
2669
2670    #[tokio::test]
2671    async fn a_prompt_larger_than_the_pipe_buffer_does_not_deadlock() {
2672        let dir = tempfile::tempdir().unwrap();
2673        let mut seat = SeatState::new("impl-A", "a", 7);
2674        // The helper never reads stdin, so an inline write_all would block once
2675        // the OS pipe buffer filled — long before the process could be waited on.
2676        let s = command_helper("ignore-stdin");
2677        let big = "x".repeat(1_000_000);
2678        let out = invoke(
2679            &s,
2680            &mut seat,
2681            &Invocation {
2682                cwd: dir.path(),
2683                prompt: &big,
2684                timeout: Duration::from_secs(60),
2685                allow_write: true,
2686                unsandboxed: false,
2687                sessions: true,
2688                artifacts: &dir.path().join("artifacts"),
2689                stem: "big",
2690                run: "test-run",
2691                node: "test",
2692                cache_dir: None,
2693                attachments: &[],
2694                writable: &[],
2695            },
2696        )
2697        .await
2698        .unwrap();
2699        assert!(out.usable(), "{out:?}");
2700        assert!(out.text.contains("done"), "{}", out.text);
2701    }
2702
2703    #[tokio::test]
2704    async fn timeout_is_reported_not_hung() {
2705        let dir = tempfile::tempdir().unwrap();
2706        let mut seat = SeatState::new("impl-A", "a", 7);
2707        let s = command_helper("sleep");
2708        let out = invoke(
2709            &s,
2710            &mut seat,
2711            &Invocation {
2712                cwd: dir.path(),
2713                prompt: "unused",
2714                timeout: Duration::from_millis(300),
2715                allow_write: true,
2716                unsandboxed: false,
2717                sessions: true,
2718                artifacts: &dir.path().join("artifacts"),
2719                stem: "slow",
2720                run: "test-run",
2721                node: "test",
2722                cache_dir: None,
2723                attachments: &[],
2724                writable: &[],
2725            },
2726        )
2727        .await
2728        .unwrap();
2729        assert!(out.timed_out);
2730        assert!(!out.usable());
2731    }
2732
2733    #[tokio::test]
2734    async fn a_timeout_keeps_what_the_agent_had_already_printed() {
2735        // The old implementation cancelled `wait_with_output`, which dropped
2736        // the buffers it owned, so `<stem>.out` was written empty on every
2737        // timeout. "It printed nothing" and "we discarded what it printed"
2738        // looked identical on disk — and one real hour-long stall was
2739        // diagnosed wrongly twice because of it.
2740        let dir = tempfile::tempdir().unwrap();
2741        let artifacts = dir.path().join("artifacts");
2742        let mut seat = SeatState::new("impl-A", "a", 7);
2743        let s = command_helper("chatty-sleep");
2744        let out = invoke(
2745            &s,
2746            &mut seat,
2747            &Invocation {
2748                cwd: dir.path(),
2749                prompt: "unused",
2750                // Wide enough to cover process-spawn latency inside a loaded
2751                // parallel test run, not merely the helper's first write. At
2752                // two seconds this passed alone and failed in the full suite,
2753                // which is a dice roll rather than a test.
2754                timeout: Duration::from_secs(10),
2755                allow_write: true,
2756                unsandboxed: false,
2757                sessions: true,
2758                artifacts: &artifacts,
2759                stem: "chatty",
2760                run: "test-run",
2761                node: "test",
2762                cache_dir: None,
2763                attachments: &[],
2764                writable: &[],
2765            },
2766        )
2767        .await
2768        .unwrap();
2769
2770        assert!(out.timed_out, "{out:?}");
2771        assert!(!out.usable(), "a cut-off answer is still not an answer");
2772        let recorded = std::fs::read_to_string(artifacts.join("chatty.out")).unwrap();
2773        assert!(
2774            recorded.contains("i-said-something"),
2775            "the artifact must keep what arrived before the kill, got {recorded:?}"
2776        );
2777        assert!(
2778            out.text.contains("i-said-something"),
2779            "and the graph must be able to see it too, got {:?}",
2780            out.text
2781        );
2782    }
2783
2784    #[test]
2785    fn missing_programs_reports_command_binaries() {
2786        let mut s = spec(AgentKind::Command, None);
2787        s.command = vec!["definitely-not-a-real-binary-xyz".to_owned()];
2788        assert_eq!(
2789            missing_programs(&[s]),
2790            ["definitely-not-a-real-binary-xyz".to_owned()]
2791        );
2792    }
2793
2794    fn pick_spec(id: &str, kind: AgentKind) -> AgentSpec {
2795        AgentSpec {
2796            id: id.to_owned(),
2797            kind,
2798            model: None,
2799            command: Vec::new(),
2800            extra_args: Vec::new(),
2801            env: BTreeMap::new(),
2802            prompt_delivery: None,
2803        }
2804    }
2805
2806    /// Availability stub: an agent is runnable unless its id was listed as
2807    /// missing. Keeps the selection tests off `PATH` entirely.
2808    fn without<'a>(missing: &'a [&'a str]) -> impl Fn(&AgentSpec) -> bool + 'a {
2809        move |a: &AgentSpec| !missing.contains(&a.id.as_str())
2810    }
2811
2812    #[test]
2813    fn pick_prefers_the_claude_seat_even_when_it_is_not_first_in_the_roster() {
2814        let agents = [
2815            pick_spec("oc", AgentKind::Opencode),
2816            pick_spec("opus", AgentKind::Claude),
2817            pick_spec("agy", AgentKind::Antigravity),
2818        ];
2819        let got = pick(&agents, None, &without(&[])).expect("a pick");
2820        assert_eq!(got.id, "opus");
2821    }
2822
2823    #[test]
2824    fn pick_falls_back_to_the_first_installed_agent_in_roster_order() {
2825        let agents = [
2826            pick_spec("opus", AgentKind::Claude),
2827            pick_spec("oc", AgentKind::Opencode),
2828            pick_spec("agy", AgentKind::Antigravity),
2829        ];
2830        let got = pick(&agents, None, &without(&["opus", "oc"])).expect("a pick");
2831        assert_eq!(got.id, "agy");
2832    }
2833
2834    #[test]
2835    fn pick_on_an_empty_roster_says_what_to_install() {
2836        let msg = pick(&[], None, &without(&[]))
2837            .expect_err("nobody to ask")
2838            .to_string();
2839        assert!(msg.contains("roster is empty"), "{msg}");
2840        assert!(msg.contains("claude"), "{msg}");
2841        assert!(msg.contains("magi.toml"), "{msg}");
2842    }
2843
2844    #[test]
2845    fn pick_on_a_roster_with_nothing_installed_names_the_programs_that_are_missing() {
2846        let agents = [
2847            pick_spec("opus", AgentKind::Claude),
2848            pick_spec("oc", AgentKind::Opencode),
2849        ];
2850        let err = pick(&agents, None, &without(&["opus", "oc"])).expect_err("nothing runnable");
2851        let msg = format!("{err:#}");
2852        assert!(msg.contains("claude"), "{msg}");
2853        assert!(msg.contains("opencode"), "{msg}");
2854    }
2855
2856    #[test]
2857    fn an_explicitly_named_agent_wins_over_the_claude_preference() {
2858        let agents = [
2859            pick_spec("opus", AgentKind::Claude),
2860            pick_spec("oc", AgentKind::Opencode),
2861        ];
2862        let got = pick(&agents, Some("oc"), &without(&[])).expect("a pick");
2863        assert_eq!(got.id, "oc");
2864    }
2865
2866    #[test]
2867    fn an_unknown_agent_id_lists_the_ids_that_do_exist() {
2868        let agents = [
2869            pick_spec("opus", AgentKind::Claude),
2870            pick_spec("oc", AgentKind::Opencode),
2871        ];
2872        let msg = pick(&agents, Some("gemini"), &without(&[]))
2873            .expect_err("no such agent")
2874            .to_string();
2875        assert!(msg.contains("gemini"), "{msg}");
2876        assert!(msg.contains("opus, oc"), "{msg}");
2877    }
2878
2879    #[test]
2880    fn an_explicitly_named_agent_that_is_not_installed_is_an_error_not_a_fallback() {
2881        let agents = [
2882            pick_spec("opus", AgentKind::Claude),
2883            pick_spec("oc", AgentKind::Opencode),
2884        ];
2885        let msg = pick(&agents, Some("oc"), &without(&["oc"]))
2886            .expect_err("must not silently substitute another model")
2887            .to_string();
2888        assert!(msg.contains("opencode"), "{msg}");
2889        assert!(msg.contains("--agent"), "{msg}");
2890    }
2891
2892    fn named(id: &str) -> AgentSpec {
2893        AgentSpec {
2894            id: id.to_owned(),
2895            ..spec(AgentKind::Command, None)
2896        }
2897    }
2898
2899    fn output(text: &str, exit: i32, quota: bool) -> AgentOutput {
2900        AgentOutput {
2901            text: text.to_owned(),
2902            exit_code: Some(exit),
2903            timed_out: false,
2904            duration_ms: 0,
2905            artifacts: Vec::new(),
2906            quota: quota.then_some(Quota { reset: None }),
2907            dropped: None,
2908            commands: Vec::new(),
2909            context_tokens: None,
2910        }
2911    }
2912
2913    #[test]
2914    fn a_chain_keeps_the_written_order_and_a_string_is_a_chain_of_one() {
2915        let agents = [named("a"), named("b"), named("c")];
2916        let all = |_: &AgentSpec| true;
2917        let chain = AgentChoice::Chain(vec!["c".into(), "a".into()]);
2918        let got = pick_chain(&agents, Some(&chain), &all, "synthesizer").unwrap();
2919        assert_eq!(
2920            got.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(),
2921            ["c", "a"]
2922        );
2923
2924        let one = AgentChoice::from("b");
2925        let got = pick_chain(&agents, Some(&one), &all, "synthesizer").unwrap();
2926        assert_eq!(got.len(), 1);
2927        assert_eq!(got[0].id, "b");
2928
2929        let got = pick_chain(&agents, None, &all, "synthesizer").unwrap();
2930        assert_eq!(got.len(), 1, "unset keeps pick's default");
2931        assert_eq!(got[0].id, "a");
2932        let empty = AgentChoice::Chain(Vec::new());
2933        assert_eq!(
2934            pick_chain(&agents, Some(&empty), &all, "x").unwrap()[0].id,
2935            "a"
2936        );
2937    }
2938
2939    #[test]
2940    fn a_chain_skips_unknown_and_uninstalled_ids_and_tries_each_once() {
2941        let agents = [named("a"), named("b")];
2942        let not_a = |s: &AgentSpec| s.id != "a";
2943        let chain = AgentChoice::Chain(
2944            ["a", "ghost", "b", "b"]
2945                .iter()
2946                .map(|s| (*s).to_owned())
2947                .collect(),
2948        );
2949        let got = pick_chain(&agents, Some(&chain), &not_a, "chatter").unwrap();
2950        assert_eq!(got.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(), ["b"]);
2951
2952        let dup = AgentChoice::Chain(vec!["b".into(), "a".into(), "b".into()]);
2953        let got = pick_chain(&agents, Some(&dup), &|_| true, "chatter").unwrap();
2954        assert_eq!(
2955            got.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(),
2956            ["b", "a"]
2957        );
2958    }
2959
2960    #[test]
2961    fn a_chain_with_nothing_runnable_names_the_role() {
2962        let agents = [named("a")];
2963        let chain = AgentChoice::Chain(vec!["a".into(), "ghost".into()]);
2964        let err = pick_chain(&agents, Some(&chain), &|_| false, "conductor")
2965            .unwrap_err()
2966            .to_string();
2967        assert!(err.contains("conductor"), "{err}");
2968    }
2969
2970    #[test]
2971    fn a_chain_advances_on_error_quota_or_an_unusable_answer_only() {
2972        assert!(chain_advances(&Err(anyhow::anyhow!("spawn failed"))));
2973        assert!(chain_advances(&Ok(output("limit", 0, true))));
2974        assert!(chain_advances(&Ok(output("", 0, false))));
2975        assert!(chain_advances(&Ok(output("x", 1, false))));
2976        assert!(!chain_advances(&Ok(output("answer", 0, false))));
2977    }
2978
2979    #[test]
2980    fn claude_context_tokens_are_unknown_because_usage_is_aggregated() {
2981        for out in [
2982            r#"{"result":"ok","num_turns":1,"usage":{"input_tokens":10,"cache_read_input_tokens":3000}}"#,
2983            r#"{"result":"ok","num_turns":3,"usage":{"input_tokens":10}}"#,
2984            r#"{"result":"ok"}"#,
2985            "not json",
2986        ] {
2987            assert_eq!(context_tokens(AgentKind::Claude, out), None, "{out}");
2988        }
2989    }
2990
2991    #[test]
2992    fn opencode_context_tokens_take_the_last_step_finish() {
2993        let out = concat!(
2994            r#"{"type":"step_finish","sessionID":"s","part":{"type":"step-finish","tokens":{"input":100,"output":5,"cache":{"read":1000,"write":50}}}}"#,
2995            "\n",
2996            r#"{"type":"text","part":{"type":"text","text":"hi"}}"#,
2997            "\n",
2998            r#"{"type":"step_finish","sessionID":"s","part":{"type":"step-finish","tokens":{"input":120,"output":9,"cache":{"read":1500}}}}"#,
2999            "\n"
3000        );
3001        assert_eq!(context_tokens(AgentKind::Opencode, out), Some(1620));
3002        let none = r#"{"type":"step_finish","part":{"type":"step-finish"}}"#;
3003        assert_eq!(context_tokens(AgentKind::Opencode, none), None);
3004        assert_eq!(context_tokens(AgentKind::Opencode, ""), None);
3005    }
3006
3007    #[test]
3008    fn agy_context_tokens_are_unknown_because_usage_is_aggregated() {
3009        let out = r#"{"conversation_id":"c","status":"OK","response":"x","usage":{"input_tokens":260113,"cache_read_tokens":2200925}}"#;
3010        assert_eq!(context_tokens(AgentKind::Antigravity, out), None);
3011    }
3012
3013    #[test]
3014    fn codex_context_tokens_are_unknown_because_usage_is_cumulative() {
3015        let out = concat!(
3016            r#"{"type":"turn.completed","usage":{"input_tokens":1000,"cached_input_tokens":900}}"#,
3017            "\n"
3018        );
3019        assert_eq!(context_tokens(AgentKind::Codex, out), None);
3020    }
3021
3022    #[test]
3023    fn omp_context_tokens_come_from_messages_without_an_agent_end() {
3024        let out = concat!(
3025            r#"{"type":"session","id":"s"}"#,
3026            "\n",
3027            r#"{"type":"message_end","message":{"role":"assistant","content":[],"usage":{"input":50,"cacheRead":400,"cacheWrite":10}}}"#,
3028            "\n",
3029            r#"{"type":"turn_end","message":{"role":"assistant","content":[],"usage":{"input":70,"cacheRead":500}}}"#,
3030            "\n"
3031        );
3032        assert_eq!(context_tokens(AgentKind::Omp, out), Some(570));
3033        let user_only = r#"{"type":"message_end","message":{"role":"user","usage":{"input":9}}}"#;
3034        assert_eq!(context_tokens(AgentKind::Omp, user_only), None);
3035        let no_usage = r#"{"type":"message_end","message":{"role":"assistant","content":[]}}"#;
3036        assert_eq!(context_tokens(AgentKind::Omp, no_usage), None);
3037    }
3038
3039    #[test]
3040    fn command_agents_report_no_context_tokens() {
3041        let out = r#"{"usage":{"input_tokens":5}}"#;
3042        assert_eq!(context_tokens(AgentKind::Command, out), None);
3043    }
3044}