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 unsandboxed: bool,
111 pub sessions: bool,
113 pub artifacts: &'a Path,
115 pub stem: &'a str,
117 pub run: &'a str,
121 pub node: &'a str,
125 pub cache_dir: Option<&'a Path>,
131 pub attachments: &'a [PathBuf],
138 pub writable: &'a [PathBuf],
143}
144
145#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
152pub struct Quota {
153 #[serde(default)]
155 pub reset: Option<String>,
156}
157
158#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
165pub struct Dropped {
166 pub why: String,
168 pub output_tokens: u64,
171}
172
173#[derive(Debug, Clone, Serialize, Deserialize)]
187pub struct CommandEvidence {
188 pub id: String,
190 pub description: String,
192 pub exit_code: Option<i32>,
194 pub result_summary: String,
196 pub source: String,
198}
199
200#[derive(Debug, Clone, Serialize, Deserialize)]
202pub struct AgentOutput {
203 pub text: String,
205 pub exit_code: Option<i32>,
207 pub timed_out: bool,
209 pub duration_ms: u64,
211 pub artifacts: Vec<String>,
213 #[serde(default)]
217 pub quota: Option<Quota>,
218 #[serde(default)]
222 pub dropped: Option<Dropped>,
223 #[serde(default)]
227 pub commands: Vec<CommandEvidence>,
228 #[serde(default)]
233 pub context_tokens: Option<u64>,
234}
235
236impl AgentOutput {
237 pub fn usable(&self) -> bool {
239 !self.timed_out && self.exit_code == Some(0) && !self.text.trim().is_empty()
240 }
241
242 pub fn quota_exhausted(&self) -> bool {
244 self.quota.is_some()
245 }
246
247 pub fn work_undelivered(&self) -> bool {
252 self.dropped.is_some()
253 }
254}
255
256const PIPE_GRACE: Duration = Duration::from_secs(3);
261
262type Captured = Arc<Mutex<Vec<u8>>>;
264
265fn drain<R>(pipe: Option<R>) -> (Captured, Option<tokio::task::JoinHandle<()>>)
275where
276 R: tokio::io::AsyncRead + Unpin + Send + 'static,
277{
278 let buf: Captured = Arc::new(Mutex::new(Vec::new()));
279 let Some(mut pipe) = pipe else {
280 return (buf, None);
281 };
282 let sink = Arc::clone(&buf);
283 let handle = tokio::spawn(async move {
284 let mut chunk = [0u8; 8192];
285 loop {
286 match pipe.read(&mut chunk).await {
287 Ok(0) | Err(_) => break,
288 Ok(n) => {
289 if let Ok(mut guard) = sink.lock() {
290 guard.extend_from_slice(&chunk[..n]);
291 }
292 }
293 }
294 }
295 });
296 (buf, Some(handle))
297}
298
299async fn collect(
304 buf: &Captured,
305 handle: Option<tokio::task::JoinHandle<()>>,
306 grace: Duration,
307) -> String {
308 if let Some(handle) = handle {
309 if tokio::time::timeout(grace, handle).await.is_err() {
310 tracing::debug!("a pipe is still held open after the child exited");
311 }
312 }
313 let bytes = buf.lock().map(|g| g.clone()).unwrap_or_default();
314 String::from_utf8_lossy(&bytes).into_owned()
315}
316
317pub async fn invoke(
319 spec: &AgentSpec,
320 seat: &mut SeatState,
321 inv: &Invocation<'_>,
322) -> Result<AgentOutput> {
323 tokio::fs::create_dir_all(inv.artifacts)
324 .await
325 .with_context(|| format!("create {}", inv.artifacts.display()))?;
326 let prompt_path = inv.artifacts.join(format!("{}.prompt.md", inv.stem));
327 tokio::fs::write(&prompt_path, inv.prompt)
328 .await
329 .with_context(|| format!("write {}", prompt_path.display()))?;
330
331 let plan = build_command(spec, seat, inv, &prompt_path)?;
332 tracing::debug!(seat = %seat.key, agent = %spec.id, argv = ?plan.argv, "spawning agent");
333
334 let started = Instant::now();
335 let child_path = spec
338 .env
339 .iter()
340 .find(|(k, _)| k.eq_ignore_ascii_case("PATH"))
341 .map(|(_, v)| std::ffi::OsString::from(v))
342 .or_else(|| std::env::var_os("PATH"))
343 .unwrap_or_default();
344 let program = crate::config::find_program_on(&plan.argv[0], &child_path).map_or_else(
345 || plan.argv[0].clone().into(),
346 std::path::PathBuf::into_os_string,
347 );
348 let mut cmd = Command::new(program);
349 cmd.args(&plan.argv[1..])
350 .current_dir(inv.cwd)
351 .envs(&spec.env)
352 .env("MAGI_SEAT", &seat.key)
353 .env("MAGI_TURN", seat.turns.to_string())
354 .env("MAGI_RUN", inv.run)
355 .env("MAGI_NODE", inv.node)
356 .env("MAGI_PROMPT_FILE", &prompt_path)
357 .env("MAGI_ALLOW_WRITE", if inv.allow_write { "1" } else { "0" })
358 .env("GIT_TERMINAL_PROMPT", "0")
359 .stdin(if plan.stdin.is_some() {
360 Stdio::piped()
361 } else {
362 Stdio::null()
363 })
364 .stdout(Stdio::piped())
365 .stderr(Stdio::piped())
366 .kill_on_drop(true)
367 .quiet();
370 if let Some(cache) = inv.cache_dir {
371 cmd.env("CARGO_TARGET_DIR", cache);
374 } else {
375 cmd.env_remove("CARGO_TARGET_DIR");
383 }
384
385 let mut child = cmd
386 .spawn()
387 .with_context(|| format!("spawn `{}` for seat {}", plan.argv[0], seat.key))?;
388 if let (Some(body), Some(mut sink)) = (plan.stdin.clone(), child.stdin.take()) {
392 tokio::spawn(async move {
393 sink.write_all(body.as_bytes()).await.ok();
394 sink.shutdown().await.ok();
395 });
396 }
397
398 let (out_buf, out_reader) = drain(child.stdout.take());
415 let (err_buf, err_reader) = drain(child.stderr.take());
416
417 let (code, timed_out) = match tokio::time::timeout(inv.timeout, child.wait()).await {
418 Ok(res) => {
419 let status = res.with_context(|| format!("wait for seat {}", seat.key))?;
420 (status.code(), false)
421 }
422 Err(_) => {
423 tracing::warn!(seat = %seat.key, secs = inv.timeout.as_secs(), "agent timed out");
424 child.start_kill().ok();
426 (None, true)
427 }
428 };
429
430 let stdout = collect(&out_buf, out_reader, PIPE_GRACE).await;
435 let stderr = collect(&err_buf, err_reader, PIPE_GRACE).await;
436
437 let out_path = inv.artifacts.join(format!("{}.out", inv.stem));
438 let err_path = inv.artifacts.join(format!("{}.err", inv.stem));
439 tokio::fs::write(&out_path, &stdout).await.ok();
440 tokio::fs::write(&err_path, &stderr).await.ok();
441
442 let mut extracted = extract(spec.kind, &stdout);
443 if spec.kind == AgentKind::Antigravity && extracted.quota.is_none() {
444 extracted.quota = agy_quota(&stdout, &stderr);
445 }
446 if let Some(quota) = &extracted.quota {
447 extracted.dropped = None;
450 tracing::warn!(
451 seat = %seat.key,
452 agent = %spec.id,
453 reset = ?quota.reset,
454 "agent is out of quota"
455 );
456 }
457 if let Some(session) = extracted.session {
458 match spec.kind {
459 AgentKind::Claude => seat.claude_session = Some(session),
460 AgentKind::Opencode | AgentKind::Antigravity | AgentKind::Codex | AgentKind::Omp => {
461 seat.captured_session = Some(session);
462 }
463 AgentKind::Command => {}
464 }
465 }
466 if let Some(status) = &extracted.status
467 && !status.eq_ignore_ascii_case("success")
468 {
469 tracing::warn!(seat = %seat.key, status = %status, "agent reported a non-success status");
470 }
471 let text = if extracted.text.trim().is_empty() {
472 if stdout.trim().is_empty() {
474 stderr.trim().to_owned()
475 } else {
476 stdout.trim().to_owned()
477 }
478 } else {
479 extracted.text
480 };
481 seat.turns += 1;
482
483 Ok(AgentOutput {
484 text,
485 exit_code: code,
486 timed_out,
487 duration_ms: started.elapsed().as_millis() as u64,
488 artifacts: vec![
489 file_name(&prompt_path),
490 file_name(&out_path),
491 file_name(&err_path),
492 ],
493 quota: extracted.quota,
494 dropped: extracted.dropped,
495 commands: extracted.commands,
496 context_tokens: extracted.context_tokens,
497 })
498}
499
500fn file_name(p: &Path) -> String {
501 p.file_name()
502 .unwrap_or_default()
503 .to_string_lossy()
504 .into_owned()
505}
506
507#[derive(Debug)]
509struct Plan {
510 argv: Vec<String>,
511 stdin: Option<String>,
512}
513
514fn pointer(kind: AgentKind, prompt_path: &Path) -> String {
525 if matches!(kind, AgentKind::Antigravity) {
526 return format!("@{}", prompt_path.display());
527 }
528 format!(
529 "Read the file at {} and follow every instruction in it exactly. That \
530 file is your complete task description; this message contains nothing \
531 else.",
532 prompt_path.display()
533 )
534}
535
536fn build_command(
537 spec: &AgentSpec,
538 seat: &SeatState,
539 inv: &Invocation<'_>,
540 prompt_path: &Path,
541) -> Result<Plan> {
542 let mut argv: Vec<String> = Vec::new();
543 let mut stdin: Option<String> = None;
544 let delivery = spec.delivery();
545 let resuming = has_session(spec.kind, seat, inv.sessions);
546
547 match spec.kind {
548 AgentKind::Claude => {
549 argv.push("claude".to_owned());
554 argv.push("-p".to_owned());
555 argv.push("--output-format".to_owned());
556 argv.push("json".to_owned());
557 if let Some(m) = &spec.model {
558 argv.push("--model".to_owned());
559 argv.push(m.clone());
560 }
561 if inv.sessions {
562 let uuid = seat
563 .claude_session
564 .as_deref()
565 .context("claude seat is missing its session uuid")?;
566 argv.push(if resuming { "--resume" } else { "--session-id" }.to_owned());
567 argv.push(uuid.to_owned());
568 }
569 argv.push("--permission-mode".to_owned());
570 argv.push("bypassPermissions".to_owned());
571 if !inv.allow_write {
572 argv.push("--disallowed-tools".to_owned());
573 argv.push("Edit,Write,MultiEdit,NotebookEdit".to_owned());
574 }
575 }
576 AgentKind::Opencode => {
577 argv.push("opencode".to_owned());
581 argv.push("run".to_owned());
582 argv.push("--format".to_owned());
583 argv.push("json".to_owned());
584 argv.push("--dir".to_owned());
585 argv.push(inv.cwd.to_string_lossy().into_owned());
586 argv.push("--auto".to_owned());
595 if let Some(m) = &spec.model {
596 argv.push("-m".to_owned());
597 argv.push(m.clone());
598 }
599 if resuming {
600 argv.push("-s".to_owned());
601 argv.push(
602 seat.captured_session
603 .clone()
604 .expect("has_session checked the id is present"),
605 );
606 }
607 }
608 AgentKind::Antigravity => {
609 argv.push("agy".to_owned());
610 argv.push("--output-format".to_owned());
611 argv.push("json".to_owned());
612 argv.push("--print-timeout".to_owned());
615 argv.push(format!("{}s", inv.timeout.as_secs()));
616 argv.push("--mode".to_owned());
617 argv.push(
618 if inv.allow_write {
619 "accept-edits"
620 } else {
621 "plan"
622 }
623 .to_owned(),
624 );
625 if inv.allow_write {
626 argv.push("--dangerously-skip-permissions".to_owned());
627 }
628 if let Some(m) = &spec.model {
629 argv.push("--model".to_owned());
630 argv.push(m.clone());
631 }
632 if resuming {
633 argv.push("--conversation".to_owned());
634 argv.push(
635 seat.captured_session
636 .clone()
637 .expect("has_session checked the id is present"),
638 );
639 }
640 let mut add_dirs: Vec<String> = Vec::new();
650 if delivery == Delivery::File || !inv.attachments.is_empty() {
651 add_dirs.push(inv.artifacts.to_string_lossy().into_owned());
652 }
653 for path in inv.attachments {
654 let Some(parent) = path.parent() else {
655 continue;
656 };
657 if parent.starts_with(inv.artifacts) {
658 continue;
659 }
660 let dir = parent.to_string_lossy().into_owned();
661 if !add_dirs.contains(&dir) {
662 add_dirs.push(dir);
663 }
664 }
665 for dir in add_dirs {
666 argv.push("--add-dir".to_owned());
667 argv.push(dir);
668 }
669 }
670 AgentKind::Codex => {
671 argv.push("codex".to_owned());
677 argv.push("exec".to_owned());
678 argv.push("--json".to_owned());
679 argv.push("--skip-git-repo-check".to_owned());
682 argv.push("-C".to_owned());
683 argv.push(inv.cwd.to_string_lossy().into_owned());
684 if inv.allow_write && inv.unsandboxed {
693 argv.push("--dangerously-bypass-approvals-and-sandbox".to_owned());
694 } else {
695 argv.push("--sandbox".to_owned());
696 argv.push(
697 if inv.allow_write {
698 "workspace-write"
699 } else {
700 "read-only"
701 }
702 .to_owned(),
703 );
704 argv.push("-c".to_owned());
707 argv.push("approval_policy=\"never\"".to_owned());
708 if inv.allow_write && !inv.writable.is_empty() {
709 let roots: Vec<String> = inv
710 .writable
711 .iter()
712 .map(|p| p.to_string_lossy().into_owned())
713 .collect();
714 argv.push("-c".to_owned());
715 argv.push(format!(
716 "sandbox_workspace_write.writable_roots={}",
717 serde_json::to_string(&roots).unwrap_or_else(|_| "[]".to_owned())
718 ));
719 }
720 }
721 if let Some(m) = &spec.model {
722 argv.push("-m".to_owned());
723 argv.push(m.clone());
724 }
725 if resuming {
731 argv.push("resume".to_owned());
732 argv.push(
733 seat.captured_session
734 .clone()
735 .expect("has_session checked the id is present"),
736 );
737 }
738 }
739 AgentKind::Omp => {
740 argv.push("omp".to_owned());
744 argv.push("-p".to_owned());
745 argv.push("--mode=json".to_owned());
746 argv.push("--auto-approve".to_owned());
756 if let Some(m) = &spec.model {
757 argv.push("--model".to_owned());
758 argv.push(m.clone());
759 }
760 if resuming {
766 argv.push("--resume".to_owned());
767 argv.push(
768 seat.captured_session
769 .clone()
770 .expect("has_session checked the id is present"),
771 );
772 }
773 }
774 AgentKind::Command => {
775 if spec.command.is_empty() {
780 bail!("agent `{}` has kind = \"command\" but no command", spec.id);
781 }
782 let vars: BTreeMap<&str, String> = BTreeMap::from([
783 ("{prompt_file}", prompt_path.to_string_lossy().into_owned()),
784 ("{cwd}", inv.cwd.to_string_lossy().into_owned()),
785 ("{label}", seat.key.clone()),
786 ("{session}", seat.claude_session.clone().unwrap_or_default()),
787 ]);
788 for raw in &spec.command {
789 let mut arg = raw.clone();
790 for (k, v) in &vars {
791 if arg.contains(k) {
792 arg = arg.replace(k, v);
793 }
794 }
795 argv.push(arg);
796 }
797 }
798 }
799
800 argv.extend(spec.extra_args.iter().cloned());
801
802 if spec.kind == AgentKind::Antigravity {
805 argv.push("-p".to_owned());
806 }
807 if spec.kind == AgentKind::Codex && delivery == Delivery::Stdin {
810 argv.push("-".to_owned());
811 }
812 match delivery {
813 Delivery::Stdin if spec.kind == AgentKind::Antigravity => {
814 argv.push(pointer(spec.kind, prompt_path));
816 }
817 Delivery::Stdin => stdin = Some(inv.prompt.to_owned()),
818 Delivery::Argv => argv.push(inv.prompt.to_owned()),
819 Delivery::File => argv.push(pointer(spec.kind, prompt_path)),
820 }
821
822 Ok(Plan { argv, stdin })
823}
824
825#[derive(Debug, Default)]
827struct Extracted {
828 text: String,
829 session: Option<String>,
830 status: Option<String>,
831 quota: Option<Quota>,
832 dropped: Option<Dropped>,
833 commands: Vec<CommandEvidence>,
834 context_tokens: Option<u64>,
835}
836
837fn extract(kind: AgentKind, stdout: &str) -> Extracted {
839 let mut extracted = extract_answer(kind, stdout);
840 extracted.context_tokens = context_tokens(kind, stdout);
841 extracted
842}
843
844fn uint(v: &serde_json::Value, key: &str) -> Option<u64> {
846 v.get(key).and_then(serde_json::Value::as_u64)
847}
848
849fn input_side(usage: &serde_json::Value, input: &[&str], cache: &[&str]) -> Option<u64> {
854 let first = |keys: &[&str]| keys.iter().find_map(|k| uint(usage, k));
855 let base = first(input)?;
856 let cached: u64 = cache.iter().filter_map(|k| uint(usage, k)).sum();
857 Some(base.saturating_add(cached))
858}
859
860fn context_tokens(kind: AgentKind, stdout: &str) -> Option<u64> {
881 use serde_json::Value;
882 let lines = || {
883 stdout
884 .lines()
885 .filter_map(|l| serde_json::from_str::<Value>(l.trim()).ok())
886 };
887 match kind {
888 AgentKind::Claude => None,
889 AgentKind::Opencode => lines()
890 .filter(|v| {
891 v.get("type").and_then(|t| t.as_str()) == Some("step_finish")
893 || v.get("part")
894 .and_then(|p| p.get("type"))
895 .and_then(|t| t.as_str())
896 == Some("step-finish")
897 })
898 .filter_map(|v| {
899 let tokens = v.get("part")?.get("tokens")?;
900 let base = uint(tokens, "input")?;
901 let cache = tokens.get("cache");
902 let cached = ["read", "write"]
903 .iter()
904 .filter_map(|k| cache.and_then(|c| uint(c, k)))
905 .sum::<u64>();
906 Some(base.saturating_add(cached))
907 })
908 .next_back(),
909 AgentKind::Antigravity => None,
910 AgentKind::Codex => None,
911 AgentKind::Omp => lines()
912 .flat_map(|v| {
913 match v.get("type").and_then(|t| t.as_str()) {
916 Some("agent_end") => v
917 .get("messages")
918 .and_then(|m| m.as_array())
919 .cloned()
920 .unwrap_or_default(),
921 Some("turn_end") | Some("message_end") => {
922 v.get("message").cloned().into_iter().collect()
923 }
924 _ => Vec::new(),
925 }
926 })
927 .filter(|m| m.get("role").and_then(|r| r.as_str()) == Some("assistant"))
928 .filter_map(|m| {
929 input_side(
930 m.get("usage")?,
931 &["input", "input_tokens"],
932 &[
933 "cacheRead",
934 "cacheWrite",
935 "cache_read_tokens",
936 "cache_write_tokens",
937 ],
938 )
939 })
940 .next_back(),
941 AgentKind::Command => None,
942 }
943}
944
945fn extract_answer(kind: AgentKind, stdout: &str) -> Extracted {
946 match kind {
947 AgentKind::Claude => {
948 let Ok(v) = serde_json::from_str::<serde_json::Value>(stdout.trim()) else {
949 return Extracted {
950 text: stdout.trim().to_owned(),
951 ..Extracted::default()
952 };
953 };
954 Extracted {
955 text: v
956 .get("result")
957 .and_then(|r| r.as_str())
958 .unwrap_or_default()
959 .to_owned(),
960 session: v
961 .get("session_id")
962 .and_then(|s| s.as_str())
963 .map(str::to_owned),
964 status: v.get("is_error").and_then(|e| e.as_bool()).map(|e| {
965 if e {
966 "error".to_owned()
967 } else {
968 "success".to_owned()
969 }
970 }),
971 quota: claude_quota(&v),
972 dropped: None,
975 commands: Vec::new(),
976 context_tokens: None,
977 }
978 }
979 AgentKind::Opencode => {
980 let mut text = String::new();
982 let mut session = None;
983 for line in stdout.lines() {
984 let Ok(v) = serde_json::from_str::<serde_json::Value>(line.trim()) else {
985 continue;
986 };
987 if session.is_none() {
988 session = v
989 .get("sessionID")
990 .and_then(|s| s.as_str())
991 .map(str::to_owned);
992 }
993 let part = v.get("part").unwrap_or(&serde_json::Value::Null);
994 if part.get("type").and_then(|t| t.as_str()) == Some("text")
995 && let Some(t) = part.get("text").and_then(|t| t.as_str())
996 {
997 if !text.is_empty() {
998 text.push('\n');
999 }
1000 text.push_str(t);
1001 }
1002 }
1003 Extracted {
1004 text,
1005 session,
1006 status: None,
1007 quota: None,
1008 dropped: None,
1009 commands: Vec::new(),
1010 context_tokens: None,
1011 }
1012 }
1013 AgentKind::Antigravity => {
1014 let obj = stdout
1017 .lines()
1018 .rev()
1019 .find_map(|l| serde_json::from_str::<serde_json::Value>(l.trim()).ok());
1020 let Some(v) = obj else {
1021 return Extracted {
1022 text: stdout.trim().to_owned(),
1023 ..Extracted::default()
1024 };
1025 };
1026 Extracted {
1027 text: v
1028 .get("response")
1029 .and_then(|r| r.as_str())
1030 .unwrap_or_default()
1031 .trim()
1032 .to_owned(),
1033 session: v
1034 .get("conversation_id")
1035 .and_then(|s| s.as_str())
1036 .map(str::to_owned),
1037 status: v.get("status").and_then(|s| s.as_str()).map(str::to_owned),
1038 quota: None,
1039 dropped: dropped_stream(&v),
1040 commands: Vec::new(),
1041 context_tokens: None,
1042 }
1043 }
1044 AgentKind::Codex => {
1045 let mut text = String::new();
1067 let mut session = None;
1068 let mut status = None;
1069 let mut commands = Vec::new();
1070 for line in stdout.lines() {
1071 let Ok(v) = serde_json::from_str::<serde_json::Value>(line.trim()) else {
1072 continue;
1073 };
1074 match v.get("type").and_then(|t| t.as_str()) {
1075 Some("thread.started") => {
1076 session = v
1077 .get("thread_id")
1078 .and_then(|s| s.as_str())
1079 .map(str::to_owned);
1080 }
1081 Some("item.completed") => {
1082 let item = v.get("item").unwrap_or(&serde_json::Value::Null);
1083 match item.get("type").and_then(|t| t.as_str()) {
1084 Some("agent_message") => {
1085 if let Some(t) = item.get("text").and_then(|t| t.as_str()) {
1086 text = t.trim().to_owned();
1087 }
1088 }
1089 Some("command_execution") => {
1090 commands.push(command_evidence(item));
1091 }
1092 _ => {}
1093 }
1094 }
1095 Some("turn.completed") => status = Some("success".to_owned()),
1096 Some("turn.failed") => status = Some("error".to_owned()),
1097 _ => {}
1098 }
1099 }
1100 Extracted {
1101 text,
1102 session,
1103 status,
1104 quota: None,
1105 dropped: None,
1106 commands,
1107 context_tokens: None,
1108 }
1109 }
1110 AgentKind::Omp => {
1111 let mut text = String::new();
1132 let mut session = None;
1133 for line in stdout.lines() {
1134 let Ok(v) = serde_json::from_str::<serde_json::Value>(line.trim()) else {
1135 continue;
1136 };
1137 if v.get("type").and_then(|t| t.as_str()) == Some("session") {
1138 session = v.get("id").and_then(|s| s.as_str()).map(str::to_owned);
1139 continue;
1140 }
1141 let messages: Vec<&serde_json::Value> = match v.get("type").and_then(|t| t.as_str())
1145 {
1146 Some("agent_end") => v
1147 .get("messages")
1148 .and_then(|m| m.as_array())
1149 .map(|m| m.iter().collect())
1150 .unwrap_or_default(),
1151 Some("turn_end") | Some("message_end") => {
1152 v.get("message").into_iter().collect()
1153 }
1154 _ => continue,
1155 };
1156 for message in messages {
1157 if message.get("role").and_then(|r| r.as_str()) != Some("assistant") {
1158 continue;
1159 }
1160 let Some(parts) = message.get("content").and_then(|c| c.as_array()) else {
1161 continue;
1162 };
1163 for part in parts {
1164 if part.get("type").and_then(|t| t.as_str()) != Some("text") {
1165 continue;
1166 }
1167 if let Some(t) = part.get("text").and_then(|t| t.as_str())
1168 && !t.trim().is_empty()
1169 {
1170 text = t.trim().to_owned();
1171 }
1172 }
1173 }
1174 }
1175 Extracted {
1176 text,
1177 session,
1178 status: None,
1179 quota: None,
1180 dropped: None,
1181 commands: Vec::new(),
1182 context_tokens: None,
1183 }
1184 }
1185 AgentKind::Command => {
1186 let parsed = serde_json::from_str::<serde_json::Value>(stdout.trim()).ok();
1191 let quota = parsed.as_ref().and_then(claude_quota);
1192 let dropped = parsed.as_ref().and_then(dropped_stream);
1195 Extracted {
1196 text: stdout.trim().to_owned(),
1197 session: None,
1198 status: None,
1199 quota,
1200 dropped,
1201 commands: Vec::new(),
1202 context_tokens: None,
1203 }
1204 }
1205 }
1206}
1207
1208fn command_evidence(item: &serde_json::Value) -> CommandEvidence {
1215 let description = match item.get("command") {
1216 Some(serde_json::Value::String(s)) => s.clone(),
1217 Some(serde_json::Value::Array(parts)) => parts
1218 .iter()
1219 .filter_map(|p| p.as_str())
1220 .collect::<Vec<_>>()
1221 .join(" "),
1222 _ => String::new(),
1223 };
1224 let result_summary = item
1225 .get("aggregated_output")
1226 .and_then(|o| o.as_str())
1227 .map(|s| tail_chars(s.trim(), 400))
1228 .unwrap_or_default();
1229 CommandEvidence {
1230 id: item
1231 .get("id")
1232 .and_then(|s| s.as_str())
1233 .unwrap_or_default()
1234 .to_owned(),
1235 description,
1236 exit_code: item
1237 .get("exit_code")
1238 .and_then(serde_json::Value::as_i64)
1239 .map(|e| e as i32),
1240 result_summary,
1241 source: "codex".to_owned(),
1242 }
1243}
1244
1245fn tail_chars(s: &str, max: usize) -> String {
1247 let count = s.chars().count();
1248 if count <= max {
1249 return s.to_owned();
1250 }
1251 s.chars().skip(count - max).collect()
1252}
1253
1254fn claude_quota(v: &serde_json::Value) -> Option<Quota> {
1261 let is_err = v.get("is_error").and_then(|e| e.as_bool()).unwrap_or(false);
1262 if !is_err {
1263 return None;
1264 }
1265 let result = v.get("result").and_then(|r| r.as_str()).unwrap_or("");
1266 if !result.to_lowercase().contains("session limit") {
1267 return None;
1268 }
1269 let reset = result
1272 .split("resets ")
1273 .nth(1)
1274 .map(str::trim)
1275 .filter(|s| !s.is_empty())
1276 .map(str::to_owned);
1277 Some(Quota { reset })
1278}
1279
1280fn agy_quota(stdout: &str, stderr: &str) -> Option<Quota> {
1295 let from_stdout = stdout
1296 .lines()
1297 .rev()
1298 .find_map(|l| serde_json::from_str::<serde_json::Value>(l.trim()).ok())
1299 .and_then(|v| {
1300 let status = v.get("status").and_then(|s| s.as_str()).unwrap_or("");
1301 let error = v.get("error").and_then(|e| e.as_str()).unwrap_or("");
1302 (status.eq_ignore_ascii_case("error") && error.to_lowercase().contains("quota reached"))
1303 .then(|| error.to_owned())
1304 });
1305 let from_stderr = || {
1306 stderr.lines().find_map(|l| {
1307 let v: serde_json::Value =
1308 serde_json::from_str(l.trim().strip_prefix("AGY_ERROR:")?.trim()).ok()?;
1309 let exhausted = v.get("status").and_then(|s| s.as_str()) == Some("RESOURCE_EXHAUSTED");
1310 let code = v.get("error_code").and_then(serde_json::Value::as_u64) == Some(429);
1311 (exhausted && code).then(|| {
1312 v.get("short_error")
1313 .and_then(|e| e.as_str())
1314 .unwrap_or_default()
1315 .to_owned()
1316 })
1317 })
1318 };
1319 let text = from_stdout.or_else(from_stderr)?;
1320 let reset = text
1321 .split_once("Resets ")
1322 .map(|(_, rest)| rest.trim().trim_end_matches('.').trim())
1323 .filter(|s| !s.is_empty())
1324 .map(str::to_owned);
1325 Some(Quota { reset })
1326}
1327
1328fn dropped_stream(v: &serde_json::Value) -> Option<Dropped> {
1359 let status = v.get("status").and_then(|s| s.as_str()).unwrap_or("");
1360 if !status.eq_ignore_ascii_case("error") {
1361 return None;
1362 }
1363 let response = v.get("response").and_then(|r| r.as_str()).unwrap_or("");
1364 if !response.trim().is_empty() {
1365 return None;
1367 }
1368 let produced = v
1369 .get("usage")
1370 .and_then(|u| u.get("output_tokens"))
1371 .and_then(serde_json::Value::as_u64)
1372 .unwrap_or(0);
1373 if produced == 0 {
1374 return None;
1376 }
1377 Some(Dropped {
1378 why: v
1379 .get("error")
1380 .and_then(|e| e.as_str())
1381 .unwrap_or("the CLI ended the stream without delivering its answer")
1382 .trim()
1383 .to_owned(),
1384 output_tokens: produced,
1385 })
1386}
1387
1388pub fn missing_programs(specs: &[AgentSpec]) -> Vec<String> {
1390 let mut missing = Vec::new();
1391 for s in specs {
1392 let program = match s.kind {
1393 AgentKind::Command => s.command.first().map(String::as_str),
1394 other => other.program(),
1395 };
1396 if let Some(p) = program
1397 && !crate::config::which(p)
1398 && !Path::new(p).is_file()
1399 && !missing.iter().any(|m: &String| m == p)
1400 {
1401 missing.push(p.to_owned());
1402 }
1403 }
1404 missing
1405}
1406
1407pub fn artifacts_dir(run_dir: &Path) -> PathBuf {
1409 run_dir.join("artifacts")
1410}
1411
1412pub fn installed(spec: &AgentSpec) -> bool {
1414 spec.kind.program().is_none_or(crate::config::which)
1417}
1418
1419pub fn pick(
1441 agents: &[AgentSpec],
1442 want: Option<&str>,
1443 available: &dyn Fn(&AgentSpec) -> bool,
1444) -> Result<AgentSpec> {
1445 if let Some(id) = want {
1446 let spec = agents
1447 .iter()
1448 .find(|a| a.id == id)
1449 .with_context(|| format!("no agent `{id}` in the roster; it has {}", ids(agents)))?;
1450 if !available(spec) {
1451 bail!(
1452 "agent `{}` needs `{}` on PATH; install it or pass a different \
1453 --agent",
1454 spec.id,
1455 spec.kind.program().unwrap_or("its command")
1456 );
1457 }
1458 return Ok(spec.clone());
1459 }
1460
1461 if agents.is_empty() {
1462 bail!(
1463 "the agent roster is empty, so there is nobody to ask: install one \
1464 of claude, opencode or agy - magi derives a roster from what is on \
1465 PATH - or add an [[agents]] entry to magi.toml."
1466 );
1467 }
1468
1469 if let Some(spec) = agents
1470 .iter()
1471 .find(|a| a.kind == AgentKind::Claude && available(a))
1472 {
1473 return Ok(spec.clone());
1474 }
1475
1476 agents
1477 .iter()
1478 .find(|a| available(a))
1479 .cloned()
1480 .with_context(|| {
1481 let missing = agents
1482 .iter()
1483 .filter_map(|a| a.kind.program())
1484 .collect::<Vec<_>>()
1485 .join(", ");
1486 format!(
1487 "no agent in the roster can be run here: install one of \
1488 {missing}, or add an [[agents]] entry to magi.toml for a CLI \
1489 you do have"
1490 )
1491 })
1492}
1493
1494pub fn pick_chain(
1502 agents: &[AgentSpec],
1503 choice: Option<&AgentChoice>,
1504 available: &dyn Fn(&AgentSpec) -> bool,
1505 role: &str,
1506) -> Result<Vec<AgentSpec>> {
1507 let wanted = choice.map(AgentChoice::ids).unwrap_or_default();
1508 if wanted.is_empty() {
1509 return pick(agents, None, available)
1510 .map(|s| vec![s])
1511 .with_context(|| format!("choose an agent for the {role} role"));
1512 }
1513 let mut chain: Vec<AgentSpec> = Vec::new();
1514 for id in wanted {
1515 if chain.iter().any(|s| s.id == id) {
1516 continue;
1517 }
1518 match pick(agents, Some(id), available) {
1519 Ok(spec) => chain.push(spec),
1520 Err(e) => tracing::warn!("[roles] {role}: skipping `{id}`: {e:#}"),
1521 }
1522 }
1523 if chain.is_empty() {
1524 bail!(
1525 "[roles] {role} names no agent that can run here; the roster has {}",
1526 ids(agents)
1527 );
1528 }
1529 Ok(chain)
1530}
1531
1532pub fn chain_advances(outcome: &Result<AgentOutput>) -> bool {
1538 match outcome {
1539 Err(_) => true,
1540 Ok(out) => output_advances(out),
1541 }
1542}
1543
1544pub fn output_advances(out: &AgentOutput) -> bool {
1548 out.quota_exhausted() || !out.usable()
1549}
1550
1551fn ids(agents: &[AgentSpec]) -> String {
1552 if agents.is_empty() {
1553 return "no agents at all".to_owned();
1554 }
1555 agents
1556 .iter()
1557 .map(|a| a.id.clone())
1558 .collect::<Vec<_>>()
1559 .join(", ")
1560}
1561
1562#[cfg(test)]
1563mod tests {
1564 use super::*;
1565
1566 const COMMAND_HELPER_MODE: &str = "MAGI_TEST_COMMAND_HELPER_MODE";
1567
1568 fn command_helper(mode: &str) -> AgentSpec {
1571 AgentSpec {
1572 id: "helper".to_owned(),
1573 kind: AgentKind::Command,
1574 model: None,
1575 command: vec![
1576 std::env::current_exe()
1577 .expect("locate test helper")
1578 .to_string_lossy()
1579 .into_owned(),
1580 "--exact".to_owned(),
1581 "agent::tests::command_agent_test_helper".to_owned(),
1582 "--nocapture".to_owned(),
1583 ],
1584 extra_args: Vec::new(),
1585 env: BTreeMap::from([(COMMAND_HELPER_MODE.to_owned(), mode.to_owned())]),
1586 prompt_delivery: None,
1587 }
1588 }
1589
1590 #[test]
1591 fn command_agent_test_helper() {
1592 match std::env::var(COMMAND_HELPER_MODE).as_deref() {
1593 Ok("reply") => println!("hello {}", std::env::var("MAGI_SEAT").unwrap()),
1594 Ok("cache") => println!("{}", std::env::var("CARGO_TARGET_DIR").unwrap()),
1595 Ok("no-cache") => println!(
1596 "{}",
1597 std::env::var("CARGO_TARGET_DIR").unwrap_or_else(|_| "ABSENT".to_owned())
1598 ),
1599 Ok("ignore-stdin") => println!("done"),
1600 Ok("chatty-sleep") => {
1601 println!("i-said-something");
1602 std::thread::sleep(Duration::from_secs(30));
1603 }
1604 Ok("sleep") => std::thread::sleep(Duration::from_secs(30)),
1605 Ok(other) => panic!("unknown command helper mode {other}"),
1606 Err(_) => {}
1607 }
1608 }
1609
1610 fn spec(kind: AgentKind, model: Option<&str>) -> AgentSpec {
1611 AgentSpec {
1612 id: "a".to_owned(),
1613 kind,
1614 model: model.map(str::to_owned),
1615 command: vec!["echo".to_owned(), "{label}".to_owned()],
1616 extra_args: Vec::new(),
1617 env: BTreeMap::new(),
1618 prompt_delivery: None,
1619 }
1620 }
1621
1622 fn inv<'a>(cwd: &'a Path, art: &'a Path, allow_write: bool) -> Invocation<'a> {
1623 Invocation {
1624 cwd,
1625 prompt: "do the thing",
1626 timeout: Duration::from_secs(900),
1627 allow_write,
1628 unsandboxed: false,
1629 sessions: true,
1630 artifacts: art,
1631 stem: "t",
1632 run: "test-run",
1633 node: "test",
1634 cache_dir: None,
1635 attachments: &[],
1636 writable: &[],
1637 }
1638 }
1639
1640 fn plan_for(kind: AgentKind, seat: &SeatState, allow_write: bool) -> Plan {
1641 build_command(
1642 &spec(kind, None),
1643 seat,
1644 &inv(Path::new("."), Path::new("/art"), allow_write),
1645 Path::new("/art/p.md"),
1646 )
1647 .unwrap()
1648 }
1649
1650 #[test]
1651 fn claude_mints_then_resumes_the_same_uuid() {
1652 let mut seat = SeatState::new("judge-1", "a", 7);
1653 let uuid = seat.claude_session.clone().unwrap();
1654 let first = plan_for(AgentKind::Claude, &seat, true);
1655 assert!(first.argv.windows(2).any(|w| w == ["--session-id", &uuid]));
1656 assert!(!first.argv.iter().any(|a| a == "--resume"));
1657
1658 seat.turns = 1;
1659 let second = plan_for(AgentKind::Claude, &seat, true);
1660 assert!(second.argv.windows(2).any(|w| w == ["--resume", &uuid]));
1661 assert!(!second.argv.iter().any(|a| a == "--session-id"));
1662 }
1663
1664 #[test]
1665 fn read_only_seats_cannot_edit() {
1666 let seat = SeatState::new("judge-1", "a", 7);
1667 let claude = plan_for(AgentKind::Claude, &seat, false);
1668 assert!(claude.argv.iter().any(|a| a == "--disallowed-tools"));
1669 assert!(
1670 !plan_for(AgentKind::Claude, &seat, true)
1671 .argv
1672 .iter()
1673 .any(|a| a == "--disallowed-tools")
1674 );
1675
1676 let agy = plan_for(AgentKind::Antigravity, &seat, false);
1677 assert!(agy.argv.windows(2).any(|w| w == ["--mode", "plan"]));
1678 assert!(
1679 !agy.argv
1680 .iter()
1681 .any(|a| a == "--dangerously-skip-permissions")
1682 );
1683 let agy_rw = plan_for(AgentKind::Antigravity, &seat, true);
1684 assert!(
1685 agy_rw
1686 .argv
1687 .windows(2)
1688 .any(|w| w == ["--mode", "accept-edits"])
1689 );
1690 assert!(
1691 agy_rw
1692 .argv
1693 .iter()
1694 .any(|a| a == "--dangerously-skip-permissions")
1695 );
1696 let agy_prompt = agy_rw
1701 .argv
1702 .iter()
1703 .position(|a| a == "-p")
1704 .map(|i| agy_rw.argv[i + 1].clone())
1705 .expect("agy takes its prompt with -p");
1706 assert!(
1707 agy_prompt.starts_with('@'),
1708 "agy must get a file reference, got {agy_prompt:?}"
1709 );
1710 assert!(
1711 !agy_prompt.contains("Read the file at"),
1712 "the prose pointer is for CLIs with no file syntax"
1713 );
1714
1715 for allow_write in [false, true] {
1720 assert!(
1721 plan_for(AgentKind::Opencode, &seat, allow_write)
1722 .argv
1723 .iter()
1724 .any(|a| a == "--auto"),
1725 "opencode needs --auto even to read (allow_write = {allow_write})"
1726 );
1727 }
1728 }
1729
1730 #[test]
1733 fn codex_gets_extra_writable_roots_only_when_it_may_write() {
1734 let seat = SeatState::new("deputy-x", "a", 7);
1735 let roots = [PathBuf::from("/data/questions")];
1736 let mk = |allow_write: bool| {
1737 let mut i = inv(Path::new("."), Path::new("/art"), allow_write);
1738 i.writable = &roots;
1739 build_command(
1740 &spec(AgentKind::Codex, None),
1741 &seat,
1742 &i,
1743 Path::new("/art/p.md"),
1744 )
1745 .unwrap()
1746 };
1747 let want = "sandbox_workspace_write.writable_roots=[\"/data/questions\"]";
1748 assert!(mk(true).argv.windows(2).any(|w| w == ["-c", want]));
1749 assert!(!mk(false).argv.iter().any(|a| a.contains("writable_roots")));
1750 }
1751
1752 #[test]
1753 fn codex_talk_turn_bypasses_sandbox_only_when_opted_in() {
1754 const BYPASS: &str = "--dangerously-bypass-approvals-and-sandbox";
1755 let mut seat = SeatState::new("talk", "a", 7);
1756 let build = |seat: &SeatState, allow_write: bool, unsandboxed: bool| {
1757 let mut i = inv(Path::new("."), Path::new("/art"), allow_write);
1758 i.unsandboxed = unsandboxed;
1759 build_command(
1760 &spec(AgentKind::Codex, None),
1761 seat,
1762 &i,
1763 Path::new("/art/p.md"),
1764 )
1765 .unwrap()
1766 };
1767
1768 let first = build(&seat, true, true);
1770 assert!(first.argv.iter().any(|a| a == BYPASS));
1771 assert!(!first.argv.iter().any(|a| a == "--sandbox"));
1772 assert!(!first.argv.iter().any(|a| a.starts_with("approval_policy")));
1773
1774 seat.captured_session = Some("thread-1".to_owned());
1776 seat.turns = 1;
1777 let resumed = build(&seat, true, true);
1778 let at = |p: &Plan, x: &str| p.argv.iter().position(|a| a == x).unwrap();
1779 assert!(at(&resumed, BYPASS) < at(&resumed, "resume"));
1780
1781 let ro = build(&seat, false, true);
1783 assert!(!ro.argv.iter().any(|a| a == BYPASS));
1784 assert!(ro.argv.windows(2).any(|w| w == ["--sandbox", "read-only"]));
1785 }
1786
1787 #[test]
1788 fn codex_is_sandboxed_reads_stdin_and_puts_resume_last() {
1789 let mut seat = SeatState::new("judge-1", "a", 7);
1790
1791 let ro = plan_for(AgentKind::Codex, &seat, false);
1795 assert!(ro.argv.windows(2).any(|w| w == ["--sandbox", "read-only"]));
1796 let rw = plan_for(AgentKind::Codex, &seat, true);
1797 assert!(
1798 rw.argv
1799 .windows(2)
1800 .any(|w| w == ["--sandbox", "workspace-write"])
1801 );
1802 let mut ro_flagged = inv(Path::new("."), Path::new("/art"), false);
1805 ro_flagged.unsandboxed = true;
1806 let ro_flagged = build_command(
1807 &spec(AgentKind::Codex, None),
1808 &seat,
1809 &ro_flagged,
1810 Path::new("/art/p.md"),
1811 )
1812 .unwrap();
1813 assert!(
1814 ro_flagged
1815 .argv
1816 .windows(2)
1817 .any(|w| w == ["--sandbox", "read-only"])
1818 );
1819 for p in [&ro, &rw, &ro_flagged] {
1820 assert!(
1821 !p.argv
1822 .iter()
1823 .any(|a| a == "--dangerously-bypass-approvals-and-sandbox"),
1824 "the bypass defeats the only enforced read-only mode we have"
1825 );
1826 assert!(
1828 p.argv
1829 .windows(2)
1830 .any(|w| w == ["-c", "approval_policy=\"never\""]),
1831 "an unattended seat that asks for approval blocks until timeout"
1832 );
1833 }
1834
1835 assert_eq!(ro.stdin.as_deref(), Some("do the thing"));
1837 assert_eq!(
1838 ro.argv.last().map(String::as_str),
1839 Some("-"),
1840 "without the `-` argument codex waits for a prompt it never gets"
1841 );
1842
1843 seat.turns = 1;
1847 assert!(!has_session(AgentKind::Codex, &seat, true));
1848 assert!(
1849 !plan_for(AgentKind::Codex, &seat, true)
1850 .argv
1851 .iter()
1852 .any(|a| a == "resume")
1853 );
1854 seat.captured_session = Some("01a07440-4545-7492-85c1-024e3259a90a".to_owned());
1855 let resumed = plan_for(AgentKind::Codex, &seat, true);
1856 let at = resumed
1857 .argv
1858 .iter()
1859 .position(|a| a == "resume")
1860 .expect("resumes by subcommand");
1861 assert_eq!(resumed.argv[at + 1], "01a07440-4545-7492-85c1-024e3259a90a");
1862 assert!(
1863 resumed.argv[..at].iter().any(|a| a == "--sandbox"),
1864 "every option precedes the subcommand"
1865 );
1866 assert_eq!(resumed.argv.last().map(String::as_str), Some("-"));
1867 }
1868
1869 #[test]
1872 fn omp_reads_stdin_auto_approves_and_resumes_by_id() {
1873 let mut seat = SeatState::new("review-1", "a", 7);
1874
1875 let first = plan_for(AgentKind::Omp, &seat, false);
1879 assert!(first.argv.iter().any(|a| a == "-p"));
1880 assert!(first.argv.iter().any(|a| a == "--mode=json"));
1881 assert_eq!(first.stdin.as_deref(), Some("do the thing"));
1882 assert!(
1883 !first.argv.iter().any(|a| a == "do the thing"),
1884 "the prompt reached argv, where Windows caps it"
1885 );
1886
1887 for allow_write in [false, true] {
1893 let p = plan_for(AgentKind::Omp, &seat, allow_write);
1894 assert!(
1895 p.argv.iter().any(|a| a == "--auto-approve"),
1896 "omp needs --auto-approve even to read (allow_write = {allow_write})"
1897 );
1898 assert!(
1899 !p.argv
1900 .iter()
1901 .any(|a| a == "--dangerously-bypass-approvals-and-sandbox"),
1902 "nothing ever asks for the bypass"
1903 );
1904 }
1905
1906 seat.turns = 1;
1909 assert!(!has_session(AgentKind::Omp, &seat, true));
1910 assert!(
1911 !plan_for(AgentKind::Omp, &seat, true)
1912 .argv
1913 .iter()
1914 .any(|a| a == "--resume")
1915 );
1916 seat.captured_session = Some("01a09fe9-4e31-7226-85b3-fda6f46689d5".to_owned());
1917 let resumed = plan_for(AgentKind::Omp, &seat, true);
1918 assert!(
1919 resumed
1920 .argv
1921 .windows(2)
1922 .any(|w| w == ["--resume", "01a09fe9-4e31-7226-85b3-fda6f46689d5"]),
1923 "a captured id is what makes the next turn a resume"
1924 );
1925 assert!(!resumed.argv.iter().any(|a| a == "--continue"));
1928 assert_eq!(resumed.stdin.as_deref(), Some("do the thing"));
1930 }
1931
1932 #[test]
1937 fn omp_takes_the_answer_without_an_agent_end_line() {
1938 let stream = concat!(
1939 r#"{"type":"session","version":3,"id":"01a09fe9-4e31-7226-85b3-fda6f46689d5","cwd":"C:\\w"}"#,
1940 "\n",
1941 r#"{"type":"agent_start"}"#,
1942 "\n",
1943 r#"{"type":"turn_start"}"#,
1944 "\n",
1945 r#"{"type":"message_update","assistantMessageEvent":{"type":"text_delta","contentIndex":1,"delta":"."}}"#,
1946 "\n",
1947 r#"{"type":"message_end","message":{"role":"assistant","content":[{"type":"thinking","thinking":"checking"},{"type":"text","text":"."}]}}"#,
1948 "\n",
1949 r#"{"type":"turn_end","message":{"role":"assistant","content":[{"type":"thinking","thinking":"done"},{"type":"text","text":"{\"vote\":\"approve\"}"}]}}"#,
1950 "\n",
1951 );
1952 let out = extract(AgentKind::Omp, stream);
1953 assert_eq!(
1954 out.text, "{\"vote\":\"approve\"}",
1955 "the last assistant text block is the answer even with no agent_end"
1956 );
1957 assert_eq!(
1958 out.session.as_deref(),
1959 Some("01a09fe9-4e31-7226-85b3-fda6f46689d5")
1960 );
1961 }
1962
1963 #[test]
1967 fn omp_walks_agent_end_and_ignores_tool_loop_narration() {
1968 let stream = concat!(
1969 r#"{"type":"session","version":3,"id":"s1"}"#,
1970 "\n",
1971 "{\"type\":\"agent_end\",\"messages\":[{\"role\":\"user\",\"content\":[{\"type\":\"text\",\"text\":\"review this\"}]},{\"role\":\"assistant\",\"content\":[{\"type\":\"text\",\"text\":\"Looking at the diff…\"}]},{\"role\":\"assistant\",\"content\":[{\"type\":\"thinking\",\"thinking\":\"…\"},{\"type\":\"text\",\"text\":\"## 判定\\n\\n問題ありません。\"}]}]}",
1972 "\n",
1973 );
1974 let out = extract(AgentKind::Omp, stream);
1975 assert_eq!(
1976 out.text, "## 判定\n\n問題ありません。",
1977 "the narration is not the answer, and non-ASCII survives intact"
1978 );
1979 assert_eq!(out.session.as_deref(), Some("s1"));
1980 }
1981
1982 #[test]
1985 fn omp_skips_non_json_lines() {
1986 let stream = concat!(
1987 "Warning: some omp notice\n",
1988 r#"{"type":"session","version":3,"id":"s2"}"#,
1989 "\n",
1990 r#"{"type":"message_end","message":{"role":"assistant","content":[{"type":"text","text":"the answer"}]}}"#,
1991 "\n",
1992 "trailing junk",
1993 "\n",
1994 );
1995 let out = extract(AgentKind::Omp, stream);
1996 assert_eq!(out.text, "the answer");
1997 assert_eq!(out.session.as_deref(), Some("s2"));
1998 }
1999
2000 #[test]
2002 fn codex_takes_the_last_agent_message_and_the_thread_id() {
2003 let stream = concat!(
2004 "2026-09-06T01:05:49.394445Z ERROR codex_models_manager: failed to load models cache\n",
2005 r#"{"type":"thread.started","thread_id":"01a07440-4545-7492-85c1-024e3259a90a"}"#,
2006 "\n",
2007 r#"{"type":"turn.started"}"#,
2008 "\n",
2009 r#"{"type":"item.completed","item":{"id":"item_0","type":"agent_message","text":"Looking into it."}}"#,
2010 "\n",
2011 r#"{"type":"item.completed","item":{"id":"item_1","type":"command_execution","text":"cargo test"}}"#,
2012 "\n",
2013 r#"{"type":"item.completed","item":{"id":"item_2","type":"agent_message","text":"{\"verdict\": \"ok\"}"}}"#,
2014 "\n",
2015 r#"{"type":"turn.completed","usage":{"input_tokens":17137}}"#,
2016 "\n",
2017 );
2018 let out = extract(AgentKind::Codex, stream);
2019 assert_eq!(
2020 out.text, "{\"verdict\": \"ok\"}",
2021 "the last agent message is the answer; earlier ones narrate"
2022 );
2023 assert_eq!(
2024 out.session.as_deref(),
2025 Some("01a07440-4545-7492-85c1-024e3259a90a")
2026 );
2027 assert_eq!(out.status.as_deref(), Some("success"));
2028
2029 let failed = concat!(
2030 r#"{"type":"thread.started","thread_id":"t1"}"#,
2031 "\n",
2032 r#"{"type":"turn.failed","error":{"message":"nope"}}"#,
2033 "\n",
2034 );
2035 assert_eq!(
2036 extract(AgentKind::Codex, failed).status.as_deref(),
2037 Some("error")
2038 );
2039 }
2040
2041 #[test]
2042 fn captured_sessions_resume_only_once_reported() {
2043 let mut seat = SeatState::new("impl-A", "a", 7);
2044 seat.turns = 1;
2045 for kind in [AgentKind::Opencode, AgentKind::Antigravity] {
2046 assert!(!has_session(kind, &seat, true));
2047 let p = plan_for(kind, &seat, true);
2048 assert!(!p.argv.iter().any(|a| a == "-s" || a == "--conversation"));
2049 }
2050
2051 seat.captured_session = Some("sid".to_owned());
2052 assert!(has_session(AgentKind::Opencode, &seat, true));
2053 assert!(
2054 plan_for(AgentKind::Opencode, &seat, true)
2055 .argv
2056 .windows(2)
2057 .any(|w| w == ["-s", "sid"])
2058 );
2059 assert!(
2060 plan_for(AgentKind::Antigravity, &seat, true)
2061 .argv
2062 .windows(2)
2063 .any(|w| w == ["--conversation", "sid"])
2064 );
2065 }
2066
2067 #[test]
2068 fn sessions_disabled_never_resumes() {
2069 let mut seat = SeatState::new("impl-A", "a", 7);
2070 seat.turns = 3;
2071 seat.captured_session = Some("sid".to_owned());
2072 for kind in [
2073 AgentKind::Claude,
2074 AgentKind::Opencode,
2075 AgentKind::Antigravity,
2076 ] {
2077 assert!(!has_session(kind, &seat, false));
2078 }
2079 }
2080
2081 #[test]
2082 fn long_prompts_never_reach_argv_for_file_delivery_clis() {
2083 let seat = SeatState::new("judge-1", "a", 7);
2084 for kind in [AgentKind::Opencode, AgentKind::Antigravity] {
2085 let p = plan_for(kind, &seat, false);
2086 assert!(
2087 p.argv.iter().all(|a| a != "do the thing"),
2088 "{kind:?} put the prompt on the command line"
2089 );
2090 assert!(p.argv.iter().any(|a| a.contains("/art/p.md")));
2091 }
2092 let p = plan_for(AgentKind::Antigravity, &seat, false);
2094 let at = p.argv.iter().position(|a| a == "-p").unwrap();
2095 assert!(p.argv.get(at + 1).is_some_and(|v| v.contains("p.md")));
2096 assert!(p.stdin.is_none());
2097 }
2098
2099 #[test]
2100 fn agy_print_timeout_tracks_the_node_budget() {
2101 let seat = SeatState::new("impl-A", "a", 7);
2102 let p = build_command(
2103 &spec(AgentKind::Antigravity, None),
2104 &seat,
2105 &Invocation {
2106 cwd: Path::new("."),
2107 prompt: "p",
2108 timeout: Duration::from_secs(3600),
2109 allow_write: true,
2110 unsandboxed: false,
2111 sessions: true,
2112 artifacts: Path::new("/art"),
2113 stem: "t",
2114 run: "test-run",
2115 node: "test",
2116 cache_dir: None,
2117 attachments: &[],
2118 writable: &[],
2119 },
2120 Path::new("/art/p.md"),
2121 )
2122 .unwrap();
2123 assert!(p.argv.windows(2).any(|w| w == ["--print-timeout", "3600s"]));
2124 }
2125
2126 #[test]
2133 fn attachments_widen_antigravitys_add_dir_even_off_file_delivery() {
2134 let mut s = spec(AgentKind::Antigravity, None);
2135 s.prompt_delivery = Some(Delivery::Argv);
2136 let seat = SeatState::new("talk", "a", 7);
2137 let atts = [PathBuf::from("/art/attachments/abc.png")];
2138
2139 let without = build_command(
2140 &s,
2141 &seat,
2142 &Invocation {
2143 attachments: &[],
2144 writable: &[],
2145 ..inv(Path::new("."), Path::new("/art"), true)
2146 },
2147 Path::new("/art/p.md"),
2148 )
2149 .unwrap();
2150 assert!(
2151 !without.argv.iter().any(|a| a == "--add-dir"),
2152 "no attachment, no reason to widen the sandbox: {without:?}"
2153 );
2154
2155 let with = build_command(
2156 &s,
2157 &seat,
2158 &Invocation {
2159 attachments: &atts,
2160 writable: &[],
2161 ..inv(Path::new("."), Path::new("/art"), true)
2162 },
2163 Path::new("/art/p.md"),
2164 )
2165 .unwrap();
2166 assert!(
2167 with.argv.windows(2).any(|w| w == ["--add-dir", "/art"]),
2168 "an attachment outside cwd must widen the sandbox even off File delivery: {with:?}"
2169 );
2170 }
2171
2172 #[test]
2178 fn an_inherited_attachment_outside_this_conversations_artifacts_dir_gets_its_own_add_dir() {
2179 let seat = SeatState::new("plan", "a", 7);
2180 let atts = [
2181 PathBuf::from("/art/attachments/own.png"),
2182 PathBuf::from("/other-chat/attachments/inherited.png"),
2183 ];
2184
2185 let p = build_command(
2186 &spec(AgentKind::Antigravity, None),
2187 &seat,
2188 &Invocation {
2189 attachments: &atts,
2190 writable: &[],
2191 ..inv(Path::new("."), Path::new("/art"), true)
2192 },
2193 Path::new("/art/p.md"),
2194 )
2195 .unwrap();
2196
2197 assert!(
2198 p.argv.windows(2).any(|w| w == ["--add-dir", "/art"]),
2199 "this conversation's own artifacts dir must still be granted: {p:?}"
2200 );
2201 assert!(
2202 p.argv
2203 .windows(2)
2204 .any(|w| w == ["--add-dir", "/other-chat/attachments"]),
2205 "the inherited attachment's own directory must be granted too: {p:?}"
2206 );
2207 }
2208
2209 #[test]
2210 fn command_agents_get_placeholders_substituted() {
2211 let seat = SeatState::new("impl-A", "a", 7);
2212 let p = plan_for(AgentKind::Command, &seat, true);
2213 assert_eq!(p.argv[0], "echo");
2214 assert_eq!(p.argv[1], "impl-A");
2215 assert_eq!(p.stdin.as_deref(), Some("do the thing"));
2216 }
2217
2218 #[test]
2219 fn claude_rate_limit_is_detected_and_reset_read_when_present() {
2220 let stdout = r#"{"is_error": true, "terminal_reason": "api_error",
2222 "result": "You've hit your session limit · resets 4:50am (Asia/Tokyo)",
2223 "session_id": "b8e928f1-754e-4bd3-86c5-0567763654e3"}"#;
2224 let out = extract(AgentKind::Claude, stdout);
2225 let quota = out.quota.as_ref().expect("rate limit must be detected");
2226 assert_eq!(
2227 quota.reset.as_deref(),
2228 Some("4:50am (Asia/Tokyo)"),
2229 "reset time read from the body"
2230 );
2231 }
2232
2233 #[test]
2234 fn claude_rate_limit_without_a_readable_reset_is_still_detected() {
2235 let out = extract(
2236 AgentKind::Claude,
2237 r#"{"is_error":true,"result":"session limit reached"}"#,
2238 );
2239 let quota = out.quota.expect("rate limit detected without a reset");
2240 assert!(quota.reset.is_none(), "unknown reset is kept as unknown");
2241 }
2242
2243 #[test]
2244 fn ordinary_failures_are_never_quota() {
2245 let claude_fail = extract(
2247 AgentKind::Claude,
2248 r#"{"is_error":true,"result":"account does not exist"}"#,
2249 );
2250 assert!(claude_fail.quota.is_none());
2251
2252 let cmd_fail = extract(AgentKind::Command, "boom");
2254 assert!(cmd_fail.quota.is_none());
2255
2256 let success = extract(
2258 AgentKind::Command,
2259 r#"{"is_error":false,"result":"session limit is fine"}"#,
2260 );
2261 assert!(success.quota.is_none());
2262 }
2263
2264 #[test]
2272 fn codex_command_execution_events_are_captured_alongside_the_final_message() {
2273 let stream = concat!(
2274 r#"{"type":"thread.started","thread_id":"t1"}"#,
2275 "\n",
2276 r#"{"type":"item.completed","item":{"id":"item49","type":"command_execution","command":["bash","-lc","cargo test --test graph_cached_gate"],"exit_code":1,"aggregated_output":"test result: 1 passed; 1 failed"}}"#,
2277 "\n",
2278 r#"{"type":"item.completed","item":{"id":"item52","type":"command_execution","command":["bash","-lc","cargo test --test graph_cached_gate a_single_test"],"exit_code":0,"aggregated_output":"test result: 1 passed; 0 failed"}}"#,
2279 "\n",
2280 r#"{"type":"item.completed","item":{"id":"item99","type":"agent_message","text":"Both tests in the target pass."}}"#,
2281 "\n",
2282 r#"{"type":"turn.completed"}"#,
2283 "\n",
2284 );
2285 let out = extract(AgentKind::Codex, stream);
2286 assert_eq!(out.text, "Both tests in the target pass.");
2287 assert_eq!(out.commands.len(), 2, "{:?}", out.commands);
2288
2289 let paired = &out.commands[0];
2290 assert_eq!(paired.id, "item49");
2291 assert_eq!(paired.exit_code, Some(1));
2292 assert!(paired.description.contains("graph_cached_gate"));
2293 assert!(paired.result_summary.contains("1 failed"));
2294
2295 let solo = &out.commands[1];
2296 assert_eq!(solo.exit_code, Some(0));
2297
2298 assert!(
2302 out.commands
2303 .iter()
2304 .any(|c| c.exit_code != Some(0) && c.description.contains("graph_cached_gate")),
2305 "a failed run of the actual target must still be visible: {:?}",
2306 out.commands
2307 );
2308 }
2309
2310 #[test]
2311 fn command_agent_can_carry_the_claude_quota_shape() {
2312 let out = extract(
2313 AgentKind::Command,
2314 r#"{"is_error":true,"result":"You've hit your session limit · resets 1:00am (UTC)"}"#,
2315 );
2316 assert!(
2317 out.quota.is_some(),
2318 "a wrapper emitting the claude shape counts as quota"
2319 );
2320 }
2321
2322 #[test]
2323 fn claude_json_result_is_extracted() {
2324 let out = extract(
2325 AgentKind::Claude,
2326 r#"{"result":"all done","session_id":"abc","is_error":false}"#,
2327 );
2328 assert_eq!(out.text, "all done");
2329 assert_eq!(out.session.as_deref(), Some("abc"));
2330 assert_eq!(out.status.as_deref(), Some("success"));
2331 }
2332
2333 #[test]
2348 fn a_clean_cli_turn_is_not_the_same_fact_as_the_nodes_own_work_being_done() {
2349 let stdout = r#"{"type":"result","subtype":"success","is_error":false,"terminal_reason":"completed","stop_reason":"end_turn","result":"I'll pause here until the `cargo make check` background run reports back.","session_id":"11111111-1111-1111-1111-111111111111"}"#;
2350 let out = extract(AgentKind::Claude, stdout);
2351 assert_eq!(out.status.as_deref(), Some("success"));
2352 assert!(out.quota.is_none());
2353 assert!(!out.text.trim().is_empty());
2354
2355 let agent_out = AgentOutput {
2356 text: out.text.clone(),
2357 exit_code: Some(0),
2358 timed_out: false,
2359 duration_ms: 500,
2360 artifacts: Vec::new(),
2361 quota: out.quota,
2362 dropped: out.dropped,
2363 commands: out.commands,
2364 context_tokens: None,
2365 };
2366 assert!(
2367 agent_out.usable(),
2368 "the CLI turn itself ended cleanly and must read as usable"
2369 );
2370 assert!(
2371 crate::verdict::extract_json::<crate::verdict::FixReport>(&agent_out.text).is_err(),
2372 "a clean CLI turn is not proof the node's own report ever arrived"
2373 );
2374 }
2375
2376 #[test]
2377 fn opencode_event_stream_is_concatenated() {
2378 let stream = concat!(
2379 r#"{"type":"step_start","sessionID":"ses_1","part":{"type":"step-start"}}"#,
2380 "\n",
2381 r#"{"type":"text","sessionID":"ses_1","part":{"type":"text","text":"first"}}"#,
2382 "\n",
2383 "garbage line\n",
2384 r#"{"type":"text","sessionID":"ses_1","part":{"type":"text","text":"second"}}"#,
2385 "\n"
2386 );
2387 let out = extract(AgentKind::Opencode, stream);
2388 assert_eq!(out.text, "first\nsecond");
2389 assert_eq!(out.session.as_deref(), Some("ses_1"));
2390 }
2391
2392 #[test]
2393 fn agy_json_survives_a_leading_warning_line() {
2394 let stdout = concat!(
2395 "warning: --mode plan has no effect while slash commands are disabled.\n",
2396 r#"{"conversation_id":"eaf2d00a","status":"SUCCESS","response":"persimmon\n"}"#,
2397 "\n"
2398 );
2399 let out = extract(AgentKind::Antigravity, stdout);
2400 assert_eq!(out.text, "persimmon");
2401 assert_eq!(out.session.as_deref(), Some("eaf2d00a"));
2402 assert_eq!(out.status.as_deref(), Some("SUCCESS"));
2403 }
2404
2405 const AGY_DROPPED: &str = concat!(
2412 r#"{"conversation_id":"36743d06-c0b3-4b79-9fa2-23869289d7b6","status":"ERROR","#,
2413 r#""response":"","error":"the connection to the agent was interrupted before "#,
2414 r#"the response finished: subscriber fell behind updates, stalled for 5s","#,
2415 r#""duration_seconds":431.1941803,"num_turns":1,"usage":{"input_tokens":260113,"#,
2416 r#""output_tokens":14267,"thinking_tokens":9695,"cache_read_tokens":2200925,"#,
2417 r#""total_tokens":274380}}"#
2418 );
2419
2420 const AGY_QUOTA_OUT: &str = concat!(
2422 r#"{"conversation_id":"323c3b5b-0000","status":"ERROR","response":"","#,
2423 r#""error":"Individual quota reached. Please upgrade your subscription to "#,
2424 r#"increase your limits. Resets in 1h2m49s.","duration_seconds":265.9,"#,
2425 r#""num_turns":2,"usage":{"input_tokens":1000,"output_tokens":50}}"#
2426 );
2427 const AGY_QUOTA_ERR: &str = concat!(
2428 "error: Individual quota reached. Resets in 1h2m49s.\n",
2429 r#"AGY_ERROR: {"short_error":"RESOURCE_EXHAUSTED (code 429): Individual quota "#,
2430 r#"reached.","status":"RESOURCE_EXHAUSTED","error_code":429,"code_kind":"http","#,
2431 r#""retryable":true}"#,
2432 "\n"
2433 );
2434
2435 #[test]
2436 fn agy_out_of_quota_is_a_quota_with_the_reset_hint() {
2437 let both = agy_quota(AGY_QUOTA_OUT, AGY_QUOTA_ERR).expect("both streams");
2438 assert_eq!(both.reset.as_deref(), Some("in 1h2m49s"));
2439 let stdout_only = agy_quota(AGY_QUOTA_OUT, "").expect("stdout alone");
2440 assert_eq!(stdout_only.reset.as_deref(), Some("in 1h2m49s"));
2441 let stderr_only = agy_quota("not json", AGY_QUOTA_ERR).expect("stderr alone");
2443 assert!(stderr_only.reset.is_none());
2444 }
2445
2446 #[test]
2447 fn ordinary_agy_failures_are_not_a_quota() {
2448 assert!(agy_quota(AGY_DROPPED, "").is_none());
2449 assert!(agy_quota(r#"{"status":"ERROR","error":"boom"}"#, "").is_none());
2450 assert!(
2451 agy_quota(
2452 "",
2453 r#"AGY_ERROR: {"status":"RESOURCE_EXHAUSTED","error_code":500}"#
2454 )
2455 .is_none()
2456 );
2457 assert!(
2458 agy_quota(
2459 "",
2460 r#"AGY_ERROR: {"status":"UNAVAILABLE","error_code":429}"#
2461 )
2462 .is_none()
2463 );
2464 assert!(agy_quota(r#"{"status":"SUCCESS","response":"ok"}"#, "").is_none());
2465 }
2466
2467 #[test]
2468 fn a_cli_that_hangs_up_on_billed_work_is_not_an_agent_that_produced_nothing() {
2469 let out = extract(AgentKind::Antigravity, AGY_DROPPED);
2470 let dropped = out.dropped.expect("recognised as undelivered work");
2471 assert_eq!(dropped.output_tokens, 14267);
2472 assert!(
2473 dropped.why.contains("subscriber fell behind"),
2474 "the CLI's own words are kept for the record: {}",
2475 dropped.why
2476 );
2477 assert_eq!(
2480 out.session.as_deref(),
2481 Some("36743d06-c0b3-4b79-9fa2-23869289d7b6")
2482 );
2483 assert!(out.quota.is_none(), "a dropped stream is not a rate limit");
2484 }
2485
2486 #[test]
2487 fn an_error_with_nothing_produced_stays_an_ordinary_failure() {
2488 let bare = r#"{"conversation_id":"c1","status":"ERROR","response":"","error":"boom"}"#;
2492 assert!(extract(AgentKind::Antigravity, bare).dropped.is_none());
2493
2494 let answered = concat!(
2497 r#"{"conversation_id":"c2","status":"ERROR","response":"here it is","#,
2498 r#""usage":{"output_tokens":10}}"#
2499 );
2500 assert!(extract(AgentKind::Antigravity, answered).dropped.is_none());
2501
2502 let ok = concat!(
2504 r#"{"conversation_id":"c3","status":"SUCCESS","response":"done","#,
2505 r#""usage":{"output_tokens":10}}"#
2506 );
2507 assert!(extract(AgentKind::Antigravity, ok).dropped.is_none());
2508 }
2509
2510 #[test]
2511 fn an_undelivered_output_is_not_usable_but_is_worth_asking_again() {
2512 let out = AgentOutput {
2513 text: String::new(),
2514 exit_code: Some(1),
2515 timed_out: false,
2516 duration_ms: 431_194,
2517 artifacts: Vec::new(),
2518 quota: None,
2519 dropped: Some(Dropped {
2520 why: "subscriber fell behind updates".to_owned(),
2521 output_tokens: 14267,
2522 }),
2523 commands: Vec::new(),
2524 context_tokens: None,
2525 };
2526 assert!(!out.usable());
2527 assert!(out.work_undelivered());
2528 assert!(!out.quota_exhausted());
2531 }
2532
2533 #[test]
2534 fn non_json_stdout_falls_back_to_raw_text() {
2535 let out = extract(AgentKind::Antigravity, "plain answer\n");
2536 assert_eq!(out.text, "plain answer");
2537 assert!(out.session.is_none());
2538 }
2539
2540 #[tokio::test]
2541 async fn command_agent_round_trip_writes_artifacts() {
2542 let dir = tempfile::tempdir().unwrap();
2543 let art = dir.path().join("artifacts");
2544 let mut seat = SeatState::new("impl-A", "a", 7);
2545 let s = command_helper("reply");
2546 let out = invoke(
2547 &s,
2548 &mut seat,
2549 &Invocation {
2550 cwd: dir.path(),
2551 prompt: "unused",
2552 timeout: Duration::from_secs(30),
2553 allow_write: true,
2554 unsandboxed: false,
2555 sessions: true,
2556 artifacts: &art,
2557 stem: "impl-A",
2558 run: "test-run",
2559 node: "test",
2560 cache_dir: None,
2561 attachments: &[],
2562 writable: &[],
2563 },
2564 )
2565 .await
2566 .unwrap();
2567 assert!(out.usable(), "{out:?}");
2568 assert!(out.text.contains("hello impl-A"), "{}", out.text);
2569 assert_eq!(seat.turns, 1);
2570 assert!(art.join("impl-A.prompt.md").is_file());
2571 assert!(art.join("impl-A.out").is_file());
2572 }
2573
2574 #[tokio::test]
2575 async fn the_invocation_cache_dir_reaches_the_seat_as_cargo_target_dir() {
2576 let dir = tempfile::tempdir().unwrap();
2580 let cache = dir.path().join("magi-cache");
2581 let mut seat = SeatState::new("impl-A", "a", 7);
2582 let s = command_helper("cache");
2583 let out = invoke(
2584 &s,
2585 &mut seat,
2586 &Invocation {
2587 cwd: dir.path(),
2588 prompt: "unused",
2589 timeout: Duration::from_secs(30),
2590 allow_write: true,
2591 unsandboxed: false,
2592 sessions: true,
2593 artifacts: &dir.path().join("artifacts"),
2594 stem: "cache",
2595 run: "test-run",
2596 node: "test",
2597 cache_dir: Some(&cache),
2598 attachments: &[],
2599 writable: &[],
2600 },
2601 )
2602 .await
2603 .unwrap();
2604 assert!(out.usable(), "{out:?}");
2605 assert!(
2606 out.text.contains(cache.to_string_lossy().as_ref()),
2607 "the seat must see CARGO_TARGET_DIR = the shared cache"
2608 );
2609 }
2610
2611 #[tokio::test]
2612 async fn cache_dir_none_strips_a_cargo_target_dir_inherited_from_this_process() {
2613 let previous = std::env::var("CARGO_TARGET_DIR").ok();
2621 unsafe {
2626 std::env::set_var("CARGO_TARGET_DIR", "/should/never/reach/a/read-only/seat");
2627 }
2628 let dir = tempfile::tempdir().unwrap();
2629 let mut seat = SeatState::new("review-1", "a", 7);
2630 let s = command_helper("no-cache");
2631 let result = invoke(
2632 &s,
2633 &mut seat,
2634 &Invocation {
2635 cwd: dir.path(),
2636 prompt: "unused",
2637 timeout: Duration::from_secs(30),
2638 allow_write: false,
2639 unsandboxed: false,
2640 sessions: true,
2641 artifacts: &dir.path().join("artifacts"),
2642 stem: "no-cache",
2643 run: "test-run",
2644 node: "test",
2645 cache_dir: None,
2646 attachments: &[],
2647 writable: &[],
2648 },
2649 )
2650 .await;
2651 unsafe {
2656 match &previous {
2657 Some(v) => std::env::set_var("CARGO_TARGET_DIR", v),
2658 None => std::env::remove_var("CARGO_TARGET_DIR"),
2659 }
2660 }
2661 let out = result.unwrap();
2662 assert!(out.usable(), "{out:?}");
2663 assert!(
2664 out.text.contains("ABSENT"),
2665 "a read-only seat must never inherit the process's own CARGO_TARGET_DIR: {}",
2666 out.text
2667 );
2668 }
2669
2670 #[tokio::test]
2671 async fn a_prompt_larger_than_the_pipe_buffer_does_not_deadlock() {
2672 let dir = tempfile::tempdir().unwrap();
2673 let mut seat = SeatState::new("impl-A", "a", 7);
2674 let s = command_helper("ignore-stdin");
2677 let big = "x".repeat(1_000_000);
2678 let out = invoke(
2679 &s,
2680 &mut seat,
2681 &Invocation {
2682 cwd: dir.path(),
2683 prompt: &big,
2684 timeout: Duration::from_secs(60),
2685 allow_write: true,
2686 unsandboxed: false,
2687 sessions: true,
2688 artifacts: &dir.path().join("artifacts"),
2689 stem: "big",
2690 run: "test-run",
2691 node: "test",
2692 cache_dir: None,
2693 attachments: &[],
2694 writable: &[],
2695 },
2696 )
2697 .await
2698 .unwrap();
2699 assert!(out.usable(), "{out:?}");
2700 assert!(out.text.contains("done"), "{}", out.text);
2701 }
2702
2703 #[tokio::test]
2704 async fn timeout_is_reported_not_hung() {
2705 let dir = tempfile::tempdir().unwrap();
2706 let mut seat = SeatState::new("impl-A", "a", 7);
2707 let s = command_helper("sleep");
2708 let out = invoke(
2709 &s,
2710 &mut seat,
2711 &Invocation {
2712 cwd: dir.path(),
2713 prompt: "unused",
2714 timeout: Duration::from_millis(300),
2715 allow_write: true,
2716 unsandboxed: false,
2717 sessions: true,
2718 artifacts: &dir.path().join("artifacts"),
2719 stem: "slow",
2720 run: "test-run",
2721 node: "test",
2722 cache_dir: None,
2723 attachments: &[],
2724 writable: &[],
2725 },
2726 )
2727 .await
2728 .unwrap();
2729 assert!(out.timed_out);
2730 assert!(!out.usable());
2731 }
2732
2733 #[tokio::test]
2734 async fn a_timeout_keeps_what_the_agent_had_already_printed() {
2735 let dir = tempfile::tempdir().unwrap();
2741 let artifacts = dir.path().join("artifacts");
2742 let mut seat = SeatState::new("impl-A", "a", 7);
2743 let s = command_helper("chatty-sleep");
2744 let out = invoke(
2745 &s,
2746 &mut seat,
2747 &Invocation {
2748 cwd: dir.path(),
2749 prompt: "unused",
2750 timeout: Duration::from_secs(10),
2755 allow_write: true,
2756 unsandboxed: false,
2757 sessions: true,
2758 artifacts: &artifacts,
2759 stem: "chatty",
2760 run: "test-run",
2761 node: "test",
2762 cache_dir: None,
2763 attachments: &[],
2764 writable: &[],
2765 },
2766 )
2767 .await
2768 .unwrap();
2769
2770 assert!(out.timed_out, "{out:?}");
2771 assert!(!out.usable(), "a cut-off answer is still not an answer");
2772 let recorded = std::fs::read_to_string(artifacts.join("chatty.out")).unwrap();
2773 assert!(
2774 recorded.contains("i-said-something"),
2775 "the artifact must keep what arrived before the kill, got {recorded:?}"
2776 );
2777 assert!(
2778 out.text.contains("i-said-something"),
2779 "and the graph must be able to see it too, got {:?}",
2780 out.text
2781 );
2782 }
2783
2784 #[test]
2785 fn missing_programs_reports_command_binaries() {
2786 let mut s = spec(AgentKind::Command, None);
2787 s.command = vec!["definitely-not-a-real-binary-xyz".to_owned()];
2788 assert_eq!(
2789 missing_programs(&[s]),
2790 ["definitely-not-a-real-binary-xyz".to_owned()]
2791 );
2792 }
2793
2794 fn pick_spec(id: &str, kind: AgentKind) -> AgentSpec {
2795 AgentSpec {
2796 id: id.to_owned(),
2797 kind,
2798 model: None,
2799 command: Vec::new(),
2800 extra_args: Vec::new(),
2801 env: BTreeMap::new(),
2802 prompt_delivery: None,
2803 }
2804 }
2805
2806 fn without<'a>(missing: &'a [&'a str]) -> impl Fn(&AgentSpec) -> bool + 'a {
2809 move |a: &AgentSpec| !missing.contains(&a.id.as_str())
2810 }
2811
2812 #[test]
2813 fn pick_prefers_the_claude_seat_even_when_it_is_not_first_in_the_roster() {
2814 let agents = [
2815 pick_spec("oc", AgentKind::Opencode),
2816 pick_spec("opus", AgentKind::Claude),
2817 pick_spec("agy", AgentKind::Antigravity),
2818 ];
2819 let got = pick(&agents, None, &without(&[])).expect("a pick");
2820 assert_eq!(got.id, "opus");
2821 }
2822
2823 #[test]
2824 fn pick_falls_back_to_the_first_installed_agent_in_roster_order() {
2825 let agents = [
2826 pick_spec("opus", AgentKind::Claude),
2827 pick_spec("oc", AgentKind::Opencode),
2828 pick_spec("agy", AgentKind::Antigravity),
2829 ];
2830 let got = pick(&agents, None, &without(&["opus", "oc"])).expect("a pick");
2831 assert_eq!(got.id, "agy");
2832 }
2833
2834 #[test]
2835 fn pick_on_an_empty_roster_says_what_to_install() {
2836 let msg = pick(&[], None, &without(&[]))
2837 .expect_err("nobody to ask")
2838 .to_string();
2839 assert!(msg.contains("roster is empty"), "{msg}");
2840 assert!(msg.contains("claude"), "{msg}");
2841 assert!(msg.contains("magi.toml"), "{msg}");
2842 }
2843
2844 #[test]
2845 fn pick_on_a_roster_with_nothing_installed_names_the_programs_that_are_missing() {
2846 let agents = [
2847 pick_spec("opus", AgentKind::Claude),
2848 pick_spec("oc", AgentKind::Opencode),
2849 ];
2850 let err = pick(&agents, None, &without(&["opus", "oc"])).expect_err("nothing runnable");
2851 let msg = format!("{err:#}");
2852 assert!(msg.contains("claude"), "{msg}");
2853 assert!(msg.contains("opencode"), "{msg}");
2854 }
2855
2856 #[test]
2857 fn an_explicitly_named_agent_wins_over_the_claude_preference() {
2858 let agents = [
2859 pick_spec("opus", AgentKind::Claude),
2860 pick_spec("oc", AgentKind::Opencode),
2861 ];
2862 let got = pick(&agents, Some("oc"), &without(&[])).expect("a pick");
2863 assert_eq!(got.id, "oc");
2864 }
2865
2866 #[test]
2867 fn an_unknown_agent_id_lists_the_ids_that_do_exist() {
2868 let agents = [
2869 pick_spec("opus", AgentKind::Claude),
2870 pick_spec("oc", AgentKind::Opencode),
2871 ];
2872 let msg = pick(&agents, Some("gemini"), &without(&[]))
2873 .expect_err("no such agent")
2874 .to_string();
2875 assert!(msg.contains("gemini"), "{msg}");
2876 assert!(msg.contains("opus, oc"), "{msg}");
2877 }
2878
2879 #[test]
2880 fn an_explicitly_named_agent_that_is_not_installed_is_an_error_not_a_fallback() {
2881 let agents = [
2882 pick_spec("opus", AgentKind::Claude),
2883 pick_spec("oc", AgentKind::Opencode),
2884 ];
2885 let msg = pick(&agents, Some("oc"), &without(&["oc"]))
2886 .expect_err("must not silently substitute another model")
2887 .to_string();
2888 assert!(msg.contains("opencode"), "{msg}");
2889 assert!(msg.contains("--agent"), "{msg}");
2890 }
2891
2892 fn named(id: &str) -> AgentSpec {
2893 AgentSpec {
2894 id: id.to_owned(),
2895 ..spec(AgentKind::Command, None)
2896 }
2897 }
2898
2899 fn output(text: &str, exit: i32, quota: bool) -> AgentOutput {
2900 AgentOutput {
2901 text: text.to_owned(),
2902 exit_code: Some(exit),
2903 timed_out: false,
2904 duration_ms: 0,
2905 artifacts: Vec::new(),
2906 quota: quota.then_some(Quota { reset: None }),
2907 dropped: None,
2908 commands: Vec::new(),
2909 context_tokens: None,
2910 }
2911 }
2912
2913 #[test]
2914 fn a_chain_keeps_the_written_order_and_a_string_is_a_chain_of_one() {
2915 let agents = [named("a"), named("b"), named("c")];
2916 let all = |_: &AgentSpec| true;
2917 let chain = AgentChoice::Chain(vec!["c".into(), "a".into()]);
2918 let got = pick_chain(&agents, Some(&chain), &all, "synthesizer").unwrap();
2919 assert_eq!(
2920 got.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(),
2921 ["c", "a"]
2922 );
2923
2924 let one = AgentChoice::from("b");
2925 let got = pick_chain(&agents, Some(&one), &all, "synthesizer").unwrap();
2926 assert_eq!(got.len(), 1);
2927 assert_eq!(got[0].id, "b");
2928
2929 let got = pick_chain(&agents, None, &all, "synthesizer").unwrap();
2930 assert_eq!(got.len(), 1, "unset keeps pick's default");
2931 assert_eq!(got[0].id, "a");
2932 let empty = AgentChoice::Chain(Vec::new());
2933 assert_eq!(
2934 pick_chain(&agents, Some(&empty), &all, "x").unwrap()[0].id,
2935 "a"
2936 );
2937 }
2938
2939 #[test]
2940 fn a_chain_skips_unknown_and_uninstalled_ids_and_tries_each_once() {
2941 let agents = [named("a"), named("b")];
2942 let not_a = |s: &AgentSpec| s.id != "a";
2943 let chain = AgentChoice::Chain(
2944 ["a", "ghost", "b", "b"]
2945 .iter()
2946 .map(|s| (*s).to_owned())
2947 .collect(),
2948 );
2949 let got = pick_chain(&agents, Some(&chain), ¬_a, "chatter").unwrap();
2950 assert_eq!(got.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(), ["b"]);
2951
2952 let dup = AgentChoice::Chain(vec!["b".into(), "a".into(), "b".into()]);
2953 let got = pick_chain(&agents, Some(&dup), &|_| true, "chatter").unwrap();
2954 assert_eq!(
2955 got.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(),
2956 ["b", "a"]
2957 );
2958 }
2959
2960 #[test]
2961 fn a_chain_with_nothing_runnable_names_the_role() {
2962 let agents = [named("a")];
2963 let chain = AgentChoice::Chain(vec!["a".into(), "ghost".into()]);
2964 let err = pick_chain(&agents, Some(&chain), &|_| false, "conductor")
2965 .unwrap_err()
2966 .to_string();
2967 assert!(err.contains("conductor"), "{err}");
2968 }
2969
2970 #[test]
2971 fn a_chain_advances_on_error_quota_or_an_unusable_answer_only() {
2972 assert!(chain_advances(&Err(anyhow::anyhow!("spawn failed"))));
2973 assert!(chain_advances(&Ok(output("limit", 0, true))));
2974 assert!(chain_advances(&Ok(output("", 0, false))));
2975 assert!(chain_advances(&Ok(output("x", 1, false))));
2976 assert!(!chain_advances(&Ok(output("answer", 0, false))));
2977 }
2978
2979 #[test]
2980 fn claude_context_tokens_are_unknown_because_usage_is_aggregated() {
2981 for out in [
2982 r#"{"result":"ok","num_turns":1,"usage":{"input_tokens":10,"cache_read_input_tokens":3000}}"#,
2983 r#"{"result":"ok","num_turns":3,"usage":{"input_tokens":10}}"#,
2984 r#"{"result":"ok"}"#,
2985 "not json",
2986 ] {
2987 assert_eq!(context_tokens(AgentKind::Claude, out), None, "{out}");
2988 }
2989 }
2990
2991 #[test]
2992 fn opencode_context_tokens_take_the_last_step_finish() {
2993 let out = concat!(
2994 r#"{"type":"step_finish","sessionID":"s","part":{"type":"step-finish","tokens":{"input":100,"output":5,"cache":{"read":1000,"write":50}}}}"#,
2995 "\n",
2996 r#"{"type":"text","part":{"type":"text","text":"hi"}}"#,
2997 "\n",
2998 r#"{"type":"step_finish","sessionID":"s","part":{"type":"step-finish","tokens":{"input":120,"output":9,"cache":{"read":1500}}}}"#,
2999 "\n"
3000 );
3001 assert_eq!(context_tokens(AgentKind::Opencode, out), Some(1620));
3002 let none = r#"{"type":"step_finish","part":{"type":"step-finish"}}"#;
3003 assert_eq!(context_tokens(AgentKind::Opencode, none), None);
3004 assert_eq!(context_tokens(AgentKind::Opencode, ""), None);
3005 }
3006
3007 #[test]
3008 fn agy_context_tokens_are_unknown_because_usage_is_aggregated() {
3009 let out = r#"{"conversation_id":"c","status":"OK","response":"x","usage":{"input_tokens":260113,"cache_read_tokens":2200925}}"#;
3010 assert_eq!(context_tokens(AgentKind::Antigravity, out), None);
3011 }
3012
3013 #[test]
3014 fn codex_context_tokens_are_unknown_because_usage_is_cumulative() {
3015 let out = concat!(
3016 r#"{"type":"turn.completed","usage":{"input_tokens":1000,"cached_input_tokens":900}}"#,
3017 "\n"
3018 );
3019 assert_eq!(context_tokens(AgentKind::Codex, out), None);
3020 }
3021
3022 #[test]
3023 fn omp_context_tokens_come_from_messages_without_an_agent_end() {
3024 let out = concat!(
3025 r#"{"type":"session","id":"s"}"#,
3026 "\n",
3027 r#"{"type":"message_end","message":{"role":"assistant","content":[],"usage":{"input":50,"cacheRead":400,"cacheWrite":10}}}"#,
3028 "\n",
3029 r#"{"type":"turn_end","message":{"role":"assistant","content":[],"usage":{"input":70,"cacheRead":500}}}"#,
3030 "\n"
3031 );
3032 assert_eq!(context_tokens(AgentKind::Omp, out), Some(570));
3033 let user_only = r#"{"type":"message_end","message":{"role":"user","usage":{"input":9}}}"#;
3034 assert_eq!(context_tokens(AgentKind::Omp, user_only), None);
3035 let no_usage = r#"{"type":"message_end","message":{"role":"assistant","content":[]}}"#;
3036 assert_eq!(context_tokens(AgentKind::Omp, no_usage), None);
3037 }
3038
3039 #[test]
3040 fn command_agents_report_no_context_tokens() {
3041 let out = r#"{"usage":{"input_tokens":5}}"#;
3042 assert_eq!(context_tokens(AgentKind::Command, out), None);
3043 }
3044}