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