1use 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
51pub struct SeatState {
52 pub key: String,
54 pub agent: String,
56 pub turns: usize,
58 pub claude_session: Option<String>,
61 pub captured_session: Option<String>,
63}
64
65impl SeatState {
66 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
79pub 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#[derive(Debug)]
95pub struct Invocation<'a> {
96 pub cwd: &'a Path,
99 pub prompt: &'a str,
101 pub timeout: Duration,
103 pub allow_write: bool,
105 pub sessions: bool,
107 pub artifacts: &'a Path,
109 pub stem: &'a str,
111 pub run: &'a str,
115 pub node: &'a str,
119 pub cache_dir: Option<&'a Path>,
125 pub attachments: &'a [PathBuf],
132 pub writable: &'a [PathBuf],
137}
138
139#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
146pub struct Quota {
147 #[serde(default)]
149 pub reset: Option<String>,
150}
151
152#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
159pub struct Dropped {
160 pub why: String,
162 pub output_tokens: u64,
165}
166
167#[derive(Debug, Clone, Serialize, Deserialize)]
181pub struct CommandEvidence {
182 pub id: String,
184 pub description: String,
186 pub exit_code: Option<i32>,
188 pub result_summary: String,
190 pub source: String,
192}
193
194#[derive(Debug, Clone, Serialize, Deserialize)]
196pub struct AgentOutput {
197 pub text: String,
199 pub exit_code: Option<i32>,
201 pub timed_out: bool,
203 pub duration_ms: u64,
205 pub artifacts: Vec<String>,
207 #[serde(default)]
211 pub quota: Option<Quota>,
212 #[serde(default)]
216 pub dropped: Option<Dropped>,
217 #[serde(default)]
221 pub commands: Vec<CommandEvidence>,
222 #[serde(default)]
227 pub context_tokens: Option<u64>,
228}
229
230impl AgentOutput {
231 pub fn usable(&self) -> bool {
233 !self.timed_out && self.exit_code == Some(0) && !self.text.trim().is_empty()
234 }
235
236 pub fn quota_exhausted(&self) -> bool {
238 self.quota.is_some()
239 }
240
241 pub fn work_undelivered(&self) -> bool {
246 self.dropped.is_some()
247 }
248}
249
250const PIPE_GRACE: Duration = Duration::from_secs(3);
255
256type Captured = Arc<Mutex<Vec<u8>>>;
258
259fn 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
293async 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
311pub 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 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 .quiet();
364 if let Some(cache) = inv.cache_dir {
365 cmd.env("CARGO_TARGET_DIR", cache);
368 } else {
369 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 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 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 child.start_kill().ok();
420 (None, true)
421 }
422 };
423
424 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 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 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#[derive(Debug)]
503struct Plan {
504 argv: Vec<String>,
505 stdin: Option<String>,
506}
507
508fn 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 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 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 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 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 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 argv.push("codex".to_owned());
671 argv.push("exec".to_owned());
672 argv.push("--json".to_owned());
673 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 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 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 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 argv.push("omp".to_owned());
730 argv.push("-p".to_owned());
731 argv.push("--mode=json".to_owned());
732 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 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 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 if spec.kind == AgentKind::Antigravity {
791 argv.push("-p".to_owned());
792 }
793 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 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#[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
823fn 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
830fn uint(v: &serde_json::Value, key: &str) -> Option<u64> {
832 v.get(key).and_then(serde_json::Value::as_u64)
833}
834
835fn 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
846fn 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 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 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 dropped: None,
961 commands: Vec::new(),
962 context_tokens: None,
963 }
964 }
965 AgentKind::Opencode => {
966 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 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 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 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 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 let parsed = serde_json::from_str::<serde_json::Value>(stdout.trim()).ok();
1177 let quota = parsed.as_ref().and_then(claude_quota);
1178 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
1194fn 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
1231fn 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
1240fn 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 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
1266fn 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
1314fn 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 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 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
1374pub 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
1393pub fn artifacts_dir(run_dir: &Path) -> PathBuf {
1395 run_dir.join("artifacts")
1396}
1397
1398pub fn installed(spec: &AgentSpec) -> bool {
1400 spec.kind.program().is_none_or(crate::config::which)
1403}
1404
1405pub 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
1480pub 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
1518pub fn chain_advances(outcome: &Result<AgentOutput>) -> bool {
1524 match outcome {
1525 Err(_) => true,
1526 Ok(out) => output_advances(out),
1527 }
1528}
1529
1530pub 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 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 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 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 #[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 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 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 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 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 #[test]
1805 fn omp_reads_stdin_auto_approves_and_resumes_by_id() {
1806 let mut seat = SeatState::new("review-1", "a", 7);
1807
1808 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 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 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 assert!(!resumed.argv.iter().any(|a| a == "--continue"));
1861 assert_eq!(resumed.stdin.as_deref(), Some("do the thing"));
1863 }
1864
1865 #[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 #[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 #[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 #[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 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 #[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 #[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 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 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 let cmd_fail = extract(AgentKind::Command, "boom");
2186 assert!(cmd_fail.quota.is_none());
2187
2188 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 #[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 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 #[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 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 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 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 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 let bare = r#"{"conversation_id":"c1","status":"ERROR","response":"","error":"boom"}"#;
2424 assert!(extract(AgentKind::Antigravity, bare).dropped.is_none());
2425
2426 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 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 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 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 let previous = std::env::var("CARGO_TARGET_DIR").ok();
2551 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 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 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 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 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 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), ¬_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}