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) => out.quota_exhausted() || !out.usable(),
1527 }
1528}
1529
1530fn ids(agents: &[AgentSpec]) -> String {
1531 if agents.is_empty() {
1532 return "no agents at all".to_owned();
1533 }
1534 agents
1535 .iter()
1536 .map(|a| a.id.clone())
1537 .collect::<Vec<_>>()
1538 .join(", ")
1539}
1540
1541#[cfg(test)]
1542mod tests {
1543 use super::*;
1544
1545 const COMMAND_HELPER_MODE: &str = "MAGI_TEST_COMMAND_HELPER_MODE";
1546
1547 fn command_helper(mode: &str) -> AgentSpec {
1550 AgentSpec {
1551 id: "helper".to_owned(),
1552 kind: AgentKind::Command,
1553 model: None,
1554 command: vec![
1555 std::env::current_exe()
1556 .expect("locate test helper")
1557 .to_string_lossy()
1558 .into_owned(),
1559 "--exact".to_owned(),
1560 "agent::tests::command_agent_test_helper".to_owned(),
1561 "--nocapture".to_owned(),
1562 ],
1563 extra_args: Vec::new(),
1564 env: BTreeMap::from([(COMMAND_HELPER_MODE.to_owned(), mode.to_owned())]),
1565 prompt_delivery: None,
1566 }
1567 }
1568
1569 #[test]
1570 fn command_agent_test_helper() {
1571 match std::env::var(COMMAND_HELPER_MODE).as_deref() {
1572 Ok("reply") => println!("hello {}", std::env::var("MAGI_SEAT").unwrap()),
1573 Ok("cache") => println!("{}", std::env::var("CARGO_TARGET_DIR").unwrap()),
1574 Ok("no-cache") => println!(
1575 "{}",
1576 std::env::var("CARGO_TARGET_DIR").unwrap_or_else(|_| "ABSENT".to_owned())
1577 ),
1578 Ok("ignore-stdin") => println!("done"),
1579 Ok("chatty-sleep") => {
1580 println!("i-said-something");
1581 std::thread::sleep(Duration::from_secs(30));
1582 }
1583 Ok("sleep") => std::thread::sleep(Duration::from_secs(30)),
1584 Ok(other) => panic!("unknown command helper mode {other}"),
1585 Err(_) => {}
1586 }
1587 }
1588
1589 fn spec(kind: AgentKind, model: Option<&str>) -> AgentSpec {
1590 AgentSpec {
1591 id: "a".to_owned(),
1592 kind,
1593 model: model.map(str::to_owned),
1594 command: vec!["echo".to_owned(), "{label}".to_owned()],
1595 extra_args: Vec::new(),
1596 env: BTreeMap::new(),
1597 prompt_delivery: None,
1598 }
1599 }
1600
1601 fn inv<'a>(cwd: &'a Path, art: &'a Path, allow_write: bool) -> Invocation<'a> {
1602 Invocation {
1603 cwd,
1604 prompt: "do the thing",
1605 timeout: Duration::from_secs(900),
1606 allow_write,
1607 sessions: true,
1608 artifacts: art,
1609 stem: "t",
1610 run: "test-run",
1611 node: "test",
1612 cache_dir: None,
1613 attachments: &[],
1614 writable: &[],
1615 }
1616 }
1617
1618 fn plan_for(kind: AgentKind, seat: &SeatState, allow_write: bool) -> Plan {
1619 build_command(
1620 &spec(kind, None),
1621 seat,
1622 &inv(Path::new("."), Path::new("/art"), allow_write),
1623 Path::new("/art/p.md"),
1624 )
1625 .unwrap()
1626 }
1627
1628 #[test]
1629 fn claude_mints_then_resumes_the_same_uuid() {
1630 let mut seat = SeatState::new("judge-1", "a", 7);
1631 let uuid = seat.claude_session.clone().unwrap();
1632 let first = plan_for(AgentKind::Claude, &seat, true);
1633 assert!(first.argv.windows(2).any(|w| w == ["--session-id", &uuid]));
1634 assert!(!first.argv.iter().any(|a| a == "--resume"));
1635
1636 seat.turns = 1;
1637 let second = plan_for(AgentKind::Claude, &seat, true);
1638 assert!(second.argv.windows(2).any(|w| w == ["--resume", &uuid]));
1639 assert!(!second.argv.iter().any(|a| a == "--session-id"));
1640 }
1641
1642 #[test]
1643 fn read_only_seats_cannot_edit() {
1644 let seat = SeatState::new("judge-1", "a", 7);
1645 let claude = plan_for(AgentKind::Claude, &seat, false);
1646 assert!(claude.argv.iter().any(|a| a == "--disallowed-tools"));
1647 assert!(
1648 !plan_for(AgentKind::Claude, &seat, true)
1649 .argv
1650 .iter()
1651 .any(|a| a == "--disallowed-tools")
1652 );
1653
1654 let agy = plan_for(AgentKind::Antigravity, &seat, false);
1655 assert!(agy.argv.windows(2).any(|w| w == ["--mode", "plan"]));
1656 assert!(
1657 !agy.argv
1658 .iter()
1659 .any(|a| a == "--dangerously-skip-permissions")
1660 );
1661 let agy_rw = plan_for(AgentKind::Antigravity, &seat, true);
1662 assert!(
1663 agy_rw
1664 .argv
1665 .windows(2)
1666 .any(|w| w == ["--mode", "accept-edits"])
1667 );
1668 assert!(
1669 agy_rw
1670 .argv
1671 .iter()
1672 .any(|a| a == "--dangerously-skip-permissions")
1673 );
1674 let agy_prompt = agy_rw
1679 .argv
1680 .iter()
1681 .position(|a| a == "-p")
1682 .map(|i| agy_rw.argv[i + 1].clone())
1683 .expect("agy takes its prompt with -p");
1684 assert!(
1685 agy_prompt.starts_with('@'),
1686 "agy must get a file reference, got {agy_prompt:?}"
1687 );
1688 assert!(
1689 !agy_prompt.contains("Read the file at"),
1690 "the prose pointer is for CLIs with no file syntax"
1691 );
1692
1693 for allow_write in [false, true] {
1698 assert!(
1699 plan_for(AgentKind::Opencode, &seat, allow_write)
1700 .argv
1701 .iter()
1702 .any(|a| a == "--auto"),
1703 "opencode needs --auto even to read (allow_write = {allow_write})"
1704 );
1705 }
1706 }
1707
1708 #[test]
1711 fn codex_gets_extra_writable_roots_only_when_it_may_write() {
1712 let seat = SeatState::new("deputy-x", "a", 7);
1713 let roots = [PathBuf::from("/data/questions")];
1714 let mk = |allow_write: bool| {
1715 let mut i = inv(Path::new("."), Path::new("/art"), allow_write);
1716 i.writable = &roots;
1717 build_command(
1718 &spec(AgentKind::Codex, None),
1719 &seat,
1720 &i,
1721 Path::new("/art/p.md"),
1722 )
1723 .unwrap()
1724 };
1725 let want = "sandbox_workspace_write.writable_roots=[\"/data/questions\"]";
1726 assert!(mk(true).argv.windows(2).any(|w| w == ["-c", want]));
1727 assert!(!mk(false).argv.iter().any(|a| a.contains("writable_roots")));
1728 }
1729
1730 #[test]
1731 fn codex_is_sandboxed_reads_stdin_and_puts_resume_last() {
1732 let mut seat = SeatState::new("judge-1", "a", 7);
1733
1734 let ro = plan_for(AgentKind::Codex, &seat, false);
1738 assert!(ro.argv.windows(2).any(|w| w == ["--sandbox", "read-only"]));
1739 let rw = plan_for(AgentKind::Codex, &seat, true);
1740 assert!(
1741 rw.argv
1742 .windows(2)
1743 .any(|w| w == ["--sandbox", "workspace-write"])
1744 );
1745 for p in [&ro, &rw] {
1746 assert!(
1747 !p.argv
1748 .iter()
1749 .any(|a| a == "--dangerously-bypass-approvals-and-sandbox"),
1750 "the bypass defeats the only enforced read-only mode we have"
1751 );
1752 assert!(
1754 p.argv
1755 .windows(2)
1756 .any(|w| w == ["-c", "approval_policy=\"never\""]),
1757 "an unattended seat that asks for approval blocks until timeout"
1758 );
1759 }
1760
1761 assert_eq!(ro.stdin.as_deref(), Some("do the thing"));
1763 assert_eq!(
1764 ro.argv.last().map(String::as_str),
1765 Some("-"),
1766 "without the `-` argument codex waits for a prompt it never gets"
1767 );
1768
1769 seat.turns = 1;
1773 assert!(!has_session(AgentKind::Codex, &seat, true));
1774 assert!(
1775 !plan_for(AgentKind::Codex, &seat, true)
1776 .argv
1777 .iter()
1778 .any(|a| a == "resume")
1779 );
1780 seat.captured_session = Some("01a07440-4545-7492-85c1-024e3259a90a".to_owned());
1781 let resumed = plan_for(AgentKind::Codex, &seat, true);
1782 let at = resumed
1783 .argv
1784 .iter()
1785 .position(|a| a == "resume")
1786 .expect("resumes by subcommand");
1787 assert_eq!(resumed.argv[at + 1], "01a07440-4545-7492-85c1-024e3259a90a");
1788 assert!(
1789 resumed.argv[..at].iter().any(|a| a == "--sandbox"),
1790 "every option precedes the subcommand"
1791 );
1792 assert_eq!(resumed.argv.last().map(String::as_str), Some("-"));
1793 }
1794
1795 #[test]
1798 fn omp_reads_stdin_auto_approves_and_resumes_by_id() {
1799 let mut seat = SeatState::new("review-1", "a", 7);
1800
1801 let first = plan_for(AgentKind::Omp, &seat, false);
1805 assert!(first.argv.iter().any(|a| a == "-p"));
1806 assert!(first.argv.iter().any(|a| a == "--mode=json"));
1807 assert_eq!(first.stdin.as_deref(), Some("do the thing"));
1808 assert!(
1809 !first.argv.iter().any(|a| a == "do the thing"),
1810 "the prompt reached argv, where Windows caps it"
1811 );
1812
1813 for allow_write in [false, true] {
1819 let p = plan_for(AgentKind::Omp, &seat, allow_write);
1820 assert!(
1821 p.argv.iter().any(|a| a == "--auto-approve"),
1822 "omp needs --auto-approve even to read (allow_write = {allow_write})"
1823 );
1824 assert!(
1825 !p.argv
1826 .iter()
1827 .any(|a| a == "--dangerously-bypass-approvals-and-sandbox"),
1828 "nothing ever asks for the bypass"
1829 );
1830 }
1831
1832 seat.turns = 1;
1835 assert!(!has_session(AgentKind::Omp, &seat, true));
1836 assert!(
1837 !plan_for(AgentKind::Omp, &seat, true)
1838 .argv
1839 .iter()
1840 .any(|a| a == "--resume")
1841 );
1842 seat.captured_session = Some("01a09fe9-4e31-7226-85b3-fda6f46689d5".to_owned());
1843 let resumed = plan_for(AgentKind::Omp, &seat, true);
1844 assert!(
1845 resumed
1846 .argv
1847 .windows(2)
1848 .any(|w| w == ["--resume", "01a09fe9-4e31-7226-85b3-fda6f46689d5"]),
1849 "a captured id is what makes the next turn a resume"
1850 );
1851 assert!(!resumed.argv.iter().any(|a| a == "--continue"));
1854 assert_eq!(resumed.stdin.as_deref(), Some("do the thing"));
1856 }
1857
1858 #[test]
1863 fn omp_takes_the_answer_without_an_agent_end_line() {
1864 let stream = concat!(
1865 r#"{"type":"session","version":3,"id":"01a09fe9-4e31-7226-85b3-fda6f46689d5","cwd":"C:\\w"}"#,
1866 "\n",
1867 r#"{"type":"agent_start"}"#,
1868 "\n",
1869 r#"{"type":"turn_start"}"#,
1870 "\n",
1871 r#"{"type":"message_update","assistantMessageEvent":{"type":"text_delta","contentIndex":1,"delta":"."}}"#,
1872 "\n",
1873 r#"{"type":"message_end","message":{"role":"assistant","content":[{"type":"thinking","thinking":"checking"},{"type":"text","text":"."}]}}"#,
1874 "\n",
1875 r#"{"type":"turn_end","message":{"role":"assistant","content":[{"type":"thinking","thinking":"done"},{"type":"text","text":"{\"vote\":\"approve\"}"}]}}"#,
1876 "\n",
1877 );
1878 let out = extract(AgentKind::Omp, stream);
1879 assert_eq!(
1880 out.text, "{\"vote\":\"approve\"}",
1881 "the last assistant text block is the answer even with no agent_end"
1882 );
1883 assert_eq!(
1884 out.session.as_deref(),
1885 Some("01a09fe9-4e31-7226-85b3-fda6f46689d5")
1886 );
1887 }
1888
1889 #[test]
1893 fn omp_walks_agent_end_and_ignores_tool_loop_narration() {
1894 let stream = concat!(
1895 r#"{"type":"session","version":3,"id":"s1"}"#,
1896 "\n",
1897 "{\"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問題ありません。\"}]}]}",
1898 "\n",
1899 );
1900 let out = extract(AgentKind::Omp, stream);
1901 assert_eq!(
1902 out.text, "## 判定\n\n問題ありません。",
1903 "the narration is not the answer, and non-ASCII survives intact"
1904 );
1905 assert_eq!(out.session.as_deref(), Some("s1"));
1906 }
1907
1908 #[test]
1911 fn omp_skips_non_json_lines() {
1912 let stream = concat!(
1913 "Warning: some omp notice\n",
1914 r#"{"type":"session","version":3,"id":"s2"}"#,
1915 "\n",
1916 r#"{"type":"message_end","message":{"role":"assistant","content":[{"type":"text","text":"the answer"}]}}"#,
1917 "\n",
1918 "trailing junk",
1919 "\n",
1920 );
1921 let out = extract(AgentKind::Omp, stream);
1922 assert_eq!(out.text, "the answer");
1923 assert_eq!(out.session.as_deref(), Some("s2"));
1924 }
1925
1926 #[test]
1928 fn codex_takes_the_last_agent_message_and_the_thread_id() {
1929 let stream = concat!(
1930 "2026-09-06T01:05:49.394445Z ERROR codex_models_manager: failed to load models cache\n",
1931 r#"{"type":"thread.started","thread_id":"01a07440-4545-7492-85c1-024e3259a90a"}"#,
1932 "\n",
1933 r#"{"type":"turn.started"}"#,
1934 "\n",
1935 r#"{"type":"item.completed","item":{"id":"item_0","type":"agent_message","text":"Looking into it."}}"#,
1936 "\n",
1937 r#"{"type":"item.completed","item":{"id":"item_1","type":"command_execution","text":"cargo test"}}"#,
1938 "\n",
1939 r#"{"type":"item.completed","item":{"id":"item_2","type":"agent_message","text":"{\"verdict\": \"ok\"}"}}"#,
1940 "\n",
1941 r#"{"type":"turn.completed","usage":{"input_tokens":17137}}"#,
1942 "\n",
1943 );
1944 let out = extract(AgentKind::Codex, stream);
1945 assert_eq!(
1946 out.text, "{\"verdict\": \"ok\"}",
1947 "the last agent message is the answer; earlier ones narrate"
1948 );
1949 assert_eq!(
1950 out.session.as_deref(),
1951 Some("01a07440-4545-7492-85c1-024e3259a90a")
1952 );
1953 assert_eq!(out.status.as_deref(), Some("success"));
1954
1955 let failed = concat!(
1956 r#"{"type":"thread.started","thread_id":"t1"}"#,
1957 "\n",
1958 r#"{"type":"turn.failed","error":{"message":"nope"}}"#,
1959 "\n",
1960 );
1961 assert_eq!(
1962 extract(AgentKind::Codex, failed).status.as_deref(),
1963 Some("error")
1964 );
1965 }
1966
1967 #[test]
1968 fn captured_sessions_resume_only_once_reported() {
1969 let mut seat = SeatState::new("impl-A", "a", 7);
1970 seat.turns = 1;
1971 for kind in [AgentKind::Opencode, AgentKind::Antigravity] {
1972 assert!(!has_session(kind, &seat, true));
1973 let p = plan_for(kind, &seat, true);
1974 assert!(!p.argv.iter().any(|a| a == "-s" || a == "--conversation"));
1975 }
1976
1977 seat.captured_session = Some("sid".to_owned());
1978 assert!(has_session(AgentKind::Opencode, &seat, true));
1979 assert!(
1980 plan_for(AgentKind::Opencode, &seat, true)
1981 .argv
1982 .windows(2)
1983 .any(|w| w == ["-s", "sid"])
1984 );
1985 assert!(
1986 plan_for(AgentKind::Antigravity, &seat, true)
1987 .argv
1988 .windows(2)
1989 .any(|w| w == ["--conversation", "sid"])
1990 );
1991 }
1992
1993 #[test]
1994 fn sessions_disabled_never_resumes() {
1995 let mut seat = SeatState::new("impl-A", "a", 7);
1996 seat.turns = 3;
1997 seat.captured_session = Some("sid".to_owned());
1998 for kind in [
1999 AgentKind::Claude,
2000 AgentKind::Opencode,
2001 AgentKind::Antigravity,
2002 ] {
2003 assert!(!has_session(kind, &seat, false));
2004 }
2005 }
2006
2007 #[test]
2008 fn long_prompts_never_reach_argv_for_file_delivery_clis() {
2009 let seat = SeatState::new("judge-1", "a", 7);
2010 for kind in [AgentKind::Opencode, AgentKind::Antigravity] {
2011 let p = plan_for(kind, &seat, false);
2012 assert!(
2013 p.argv.iter().all(|a| a != "do the thing"),
2014 "{kind:?} put the prompt on the command line"
2015 );
2016 assert!(p.argv.iter().any(|a| a.contains("/art/p.md")));
2017 }
2018 let p = plan_for(AgentKind::Antigravity, &seat, false);
2020 let at = p.argv.iter().position(|a| a == "-p").unwrap();
2021 assert!(p.argv.get(at + 1).is_some_and(|v| v.contains("p.md")));
2022 assert!(p.stdin.is_none());
2023 }
2024
2025 #[test]
2026 fn agy_print_timeout_tracks_the_node_budget() {
2027 let seat = SeatState::new("impl-A", "a", 7);
2028 let p = build_command(
2029 &spec(AgentKind::Antigravity, None),
2030 &seat,
2031 &Invocation {
2032 cwd: Path::new("."),
2033 prompt: "p",
2034 timeout: Duration::from_secs(3600),
2035 allow_write: true,
2036 sessions: true,
2037 artifacts: Path::new("/art"),
2038 stem: "t",
2039 run: "test-run",
2040 node: "test",
2041 cache_dir: None,
2042 attachments: &[],
2043 writable: &[],
2044 },
2045 Path::new("/art/p.md"),
2046 )
2047 .unwrap();
2048 assert!(p.argv.windows(2).any(|w| w == ["--print-timeout", "3600s"]));
2049 }
2050
2051 #[test]
2058 fn attachments_widen_antigravitys_add_dir_even_off_file_delivery() {
2059 let mut s = spec(AgentKind::Antigravity, None);
2060 s.prompt_delivery = Some(Delivery::Argv);
2061 let seat = SeatState::new("talk", "a", 7);
2062 let atts = [PathBuf::from("/art/attachments/abc.png")];
2063
2064 let without = build_command(
2065 &s,
2066 &seat,
2067 &Invocation {
2068 attachments: &[],
2069 writable: &[],
2070 ..inv(Path::new("."), Path::new("/art"), true)
2071 },
2072 Path::new("/art/p.md"),
2073 )
2074 .unwrap();
2075 assert!(
2076 !without.argv.iter().any(|a| a == "--add-dir"),
2077 "no attachment, no reason to widen the sandbox: {without:?}"
2078 );
2079
2080 let with = build_command(
2081 &s,
2082 &seat,
2083 &Invocation {
2084 attachments: &atts,
2085 writable: &[],
2086 ..inv(Path::new("."), Path::new("/art"), true)
2087 },
2088 Path::new("/art/p.md"),
2089 )
2090 .unwrap();
2091 assert!(
2092 with.argv.windows(2).any(|w| w == ["--add-dir", "/art"]),
2093 "an attachment outside cwd must widen the sandbox even off File delivery: {with:?}"
2094 );
2095 }
2096
2097 #[test]
2103 fn an_inherited_attachment_outside_this_conversations_artifacts_dir_gets_its_own_add_dir() {
2104 let seat = SeatState::new("plan", "a", 7);
2105 let atts = [
2106 PathBuf::from("/art/attachments/own.png"),
2107 PathBuf::from("/other-chat/attachments/inherited.png"),
2108 ];
2109
2110 let p = build_command(
2111 &spec(AgentKind::Antigravity, None),
2112 &seat,
2113 &Invocation {
2114 attachments: &atts,
2115 writable: &[],
2116 ..inv(Path::new("."), Path::new("/art"), true)
2117 },
2118 Path::new("/art/p.md"),
2119 )
2120 .unwrap();
2121
2122 assert!(
2123 p.argv.windows(2).any(|w| w == ["--add-dir", "/art"]),
2124 "this conversation's own artifacts dir must still be granted: {p:?}"
2125 );
2126 assert!(
2127 p.argv
2128 .windows(2)
2129 .any(|w| w == ["--add-dir", "/other-chat/attachments"]),
2130 "the inherited attachment's own directory must be granted too: {p:?}"
2131 );
2132 }
2133
2134 #[test]
2135 fn command_agents_get_placeholders_substituted() {
2136 let seat = SeatState::new("impl-A", "a", 7);
2137 let p = plan_for(AgentKind::Command, &seat, true);
2138 assert_eq!(p.argv[0], "echo");
2139 assert_eq!(p.argv[1], "impl-A");
2140 assert_eq!(p.stdin.as_deref(), Some("do the thing"));
2141 }
2142
2143 #[test]
2144 fn claude_rate_limit_is_detected_and_reset_read_when_present() {
2145 let stdout = r#"{"is_error": true, "terminal_reason": "api_error",
2147 "result": "You've hit your session limit · resets 4:50am (Asia/Tokyo)",
2148 "session_id": "b8e928f1-754e-4bd3-86c5-0567763654e3"}"#;
2149 let out = extract(AgentKind::Claude, stdout);
2150 let quota = out.quota.as_ref().expect("rate limit must be detected");
2151 assert_eq!(
2152 quota.reset.as_deref(),
2153 Some("4:50am (Asia/Tokyo)"),
2154 "reset time read from the body"
2155 );
2156 }
2157
2158 #[test]
2159 fn claude_rate_limit_without_a_readable_reset_is_still_detected() {
2160 let out = extract(
2161 AgentKind::Claude,
2162 r#"{"is_error":true,"result":"session limit reached"}"#,
2163 );
2164 let quota = out.quota.expect("rate limit detected without a reset");
2165 assert!(quota.reset.is_none(), "unknown reset is kept as unknown");
2166 }
2167
2168 #[test]
2169 fn ordinary_failures_are_never_quota() {
2170 let claude_fail = extract(
2172 AgentKind::Claude,
2173 r#"{"is_error":true,"result":"account does not exist"}"#,
2174 );
2175 assert!(claude_fail.quota.is_none());
2176
2177 let cmd_fail = extract(AgentKind::Command, "boom");
2179 assert!(cmd_fail.quota.is_none());
2180
2181 let success = extract(
2183 AgentKind::Command,
2184 r#"{"is_error":false,"result":"session limit is fine"}"#,
2185 );
2186 assert!(success.quota.is_none());
2187 }
2188
2189 #[test]
2197 fn codex_command_execution_events_are_captured_alongside_the_final_message() {
2198 let stream = concat!(
2199 r#"{"type":"thread.started","thread_id":"t1"}"#,
2200 "\n",
2201 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"}}"#,
2202 "\n",
2203 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"}}"#,
2204 "\n",
2205 r#"{"type":"item.completed","item":{"id":"item99","type":"agent_message","text":"Both tests in the target pass."}}"#,
2206 "\n",
2207 r#"{"type":"turn.completed"}"#,
2208 "\n",
2209 );
2210 let out = extract(AgentKind::Codex, stream);
2211 assert_eq!(out.text, "Both tests in the target pass.");
2212 assert_eq!(out.commands.len(), 2, "{:?}", out.commands);
2213
2214 let paired = &out.commands[0];
2215 assert_eq!(paired.id, "item49");
2216 assert_eq!(paired.exit_code, Some(1));
2217 assert!(paired.description.contains("graph_cached_gate"));
2218 assert!(paired.result_summary.contains("1 failed"));
2219
2220 let solo = &out.commands[1];
2221 assert_eq!(solo.exit_code, Some(0));
2222
2223 assert!(
2227 out.commands
2228 .iter()
2229 .any(|c| c.exit_code != Some(0) && c.description.contains("graph_cached_gate")),
2230 "a failed run of the actual target must still be visible: {:?}",
2231 out.commands
2232 );
2233 }
2234
2235 #[test]
2236 fn command_agent_can_carry_the_claude_quota_shape() {
2237 let out = extract(
2238 AgentKind::Command,
2239 r#"{"is_error":true,"result":"You've hit your session limit · resets 1:00am (UTC)"}"#,
2240 );
2241 assert!(
2242 out.quota.is_some(),
2243 "a wrapper emitting the claude shape counts as quota"
2244 );
2245 }
2246
2247 #[test]
2248 fn claude_json_result_is_extracted() {
2249 let out = extract(
2250 AgentKind::Claude,
2251 r#"{"result":"all done","session_id":"abc","is_error":false}"#,
2252 );
2253 assert_eq!(out.text, "all done");
2254 assert_eq!(out.session.as_deref(), Some("abc"));
2255 assert_eq!(out.status.as_deref(), Some("success"));
2256 }
2257
2258 #[test]
2273 fn a_clean_cli_turn_is_not_the_same_fact_as_the_nodes_own_work_being_done() {
2274 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"}"#;
2275 let out = extract(AgentKind::Claude, stdout);
2276 assert_eq!(out.status.as_deref(), Some("success"));
2277 assert!(out.quota.is_none());
2278 assert!(!out.text.trim().is_empty());
2279
2280 let agent_out = AgentOutput {
2281 text: out.text.clone(),
2282 exit_code: Some(0),
2283 timed_out: false,
2284 duration_ms: 500,
2285 artifacts: Vec::new(),
2286 quota: out.quota,
2287 dropped: out.dropped,
2288 commands: out.commands,
2289 context_tokens: None,
2290 };
2291 assert!(
2292 agent_out.usable(),
2293 "the CLI turn itself ended cleanly and must read as usable"
2294 );
2295 assert!(
2296 crate::verdict::extract_json::<crate::verdict::FixReport>(&agent_out.text).is_err(),
2297 "a clean CLI turn is not proof the node's own report ever arrived"
2298 );
2299 }
2300
2301 #[test]
2302 fn opencode_event_stream_is_concatenated() {
2303 let stream = concat!(
2304 r#"{"type":"step_start","sessionID":"ses_1","part":{"type":"step-start"}}"#,
2305 "\n",
2306 r#"{"type":"text","sessionID":"ses_1","part":{"type":"text","text":"first"}}"#,
2307 "\n",
2308 "garbage line\n",
2309 r#"{"type":"text","sessionID":"ses_1","part":{"type":"text","text":"second"}}"#,
2310 "\n"
2311 );
2312 let out = extract(AgentKind::Opencode, stream);
2313 assert_eq!(out.text, "first\nsecond");
2314 assert_eq!(out.session.as_deref(), Some("ses_1"));
2315 }
2316
2317 #[test]
2318 fn agy_json_survives_a_leading_warning_line() {
2319 let stdout = concat!(
2320 "warning: --mode plan has no effect while slash commands are disabled.\n",
2321 r#"{"conversation_id":"eaf2d00a","status":"SUCCESS","response":"persimmon\n"}"#,
2322 "\n"
2323 );
2324 let out = extract(AgentKind::Antigravity, stdout);
2325 assert_eq!(out.text, "persimmon");
2326 assert_eq!(out.session.as_deref(), Some("eaf2d00a"));
2327 assert_eq!(out.status.as_deref(), Some("SUCCESS"));
2328 }
2329
2330 const AGY_DROPPED: &str = concat!(
2337 r#"{"conversation_id":"36743d06-c0b3-4b79-9fa2-23869289d7b6","status":"ERROR","#,
2338 r#""response":"","error":"the connection to the agent was interrupted before "#,
2339 r#"the response finished: subscriber fell behind updates, stalled for 5s","#,
2340 r#""duration_seconds":431.1941803,"num_turns":1,"usage":{"input_tokens":260113,"#,
2341 r#""output_tokens":14267,"thinking_tokens":9695,"cache_read_tokens":2200925,"#,
2342 r#""total_tokens":274380}}"#
2343 );
2344
2345 const AGY_QUOTA_OUT: &str = concat!(
2347 r#"{"conversation_id":"323c3b5b-0000","status":"ERROR","response":"","#,
2348 r#""error":"Individual quota reached. Please upgrade your subscription to "#,
2349 r#"increase your limits. Resets in 1h2m49s.","duration_seconds":265.9,"#,
2350 r#""num_turns":2,"usage":{"input_tokens":1000,"output_tokens":50}}"#
2351 );
2352 const AGY_QUOTA_ERR: &str = concat!(
2353 "error: Individual quota reached. Resets in 1h2m49s.\n",
2354 r#"AGY_ERROR: {"short_error":"RESOURCE_EXHAUSTED (code 429): Individual quota "#,
2355 r#"reached.","status":"RESOURCE_EXHAUSTED","error_code":429,"code_kind":"http","#,
2356 r#""retryable":true}"#,
2357 "\n"
2358 );
2359
2360 #[test]
2361 fn agy_out_of_quota_is_a_quota_with_the_reset_hint() {
2362 let both = agy_quota(AGY_QUOTA_OUT, AGY_QUOTA_ERR).expect("both streams");
2363 assert_eq!(both.reset.as_deref(), Some("in 1h2m49s"));
2364 let stdout_only = agy_quota(AGY_QUOTA_OUT, "").expect("stdout alone");
2365 assert_eq!(stdout_only.reset.as_deref(), Some("in 1h2m49s"));
2366 let stderr_only = agy_quota("not json", AGY_QUOTA_ERR).expect("stderr alone");
2368 assert!(stderr_only.reset.is_none());
2369 }
2370
2371 #[test]
2372 fn ordinary_agy_failures_are_not_a_quota() {
2373 assert!(agy_quota(AGY_DROPPED, "").is_none());
2374 assert!(agy_quota(r#"{"status":"ERROR","error":"boom"}"#, "").is_none());
2375 assert!(
2376 agy_quota(
2377 "",
2378 r#"AGY_ERROR: {"status":"RESOURCE_EXHAUSTED","error_code":500}"#
2379 )
2380 .is_none()
2381 );
2382 assert!(
2383 agy_quota(
2384 "",
2385 r#"AGY_ERROR: {"status":"UNAVAILABLE","error_code":429}"#
2386 )
2387 .is_none()
2388 );
2389 assert!(agy_quota(r#"{"status":"SUCCESS","response":"ok"}"#, "").is_none());
2390 }
2391
2392 #[test]
2393 fn a_cli_that_hangs_up_on_billed_work_is_not_an_agent_that_produced_nothing() {
2394 let out = extract(AgentKind::Antigravity, AGY_DROPPED);
2395 let dropped = out.dropped.expect("recognised as undelivered work");
2396 assert_eq!(dropped.output_tokens, 14267);
2397 assert!(
2398 dropped.why.contains("subscriber fell behind"),
2399 "the CLI's own words are kept for the record: {}",
2400 dropped.why
2401 );
2402 assert_eq!(
2405 out.session.as_deref(),
2406 Some("36743d06-c0b3-4b79-9fa2-23869289d7b6")
2407 );
2408 assert!(out.quota.is_none(), "a dropped stream is not a rate limit");
2409 }
2410
2411 #[test]
2412 fn an_error_with_nothing_produced_stays_an_ordinary_failure() {
2413 let bare = r#"{"conversation_id":"c1","status":"ERROR","response":"","error":"boom"}"#;
2417 assert!(extract(AgentKind::Antigravity, bare).dropped.is_none());
2418
2419 let answered = concat!(
2422 r#"{"conversation_id":"c2","status":"ERROR","response":"here it is","#,
2423 r#""usage":{"output_tokens":10}}"#
2424 );
2425 assert!(extract(AgentKind::Antigravity, answered).dropped.is_none());
2426
2427 let ok = concat!(
2429 r#"{"conversation_id":"c3","status":"SUCCESS","response":"done","#,
2430 r#""usage":{"output_tokens":10}}"#
2431 );
2432 assert!(extract(AgentKind::Antigravity, ok).dropped.is_none());
2433 }
2434
2435 #[test]
2436 fn an_undelivered_output_is_not_usable_but_is_worth_asking_again() {
2437 let out = AgentOutput {
2438 text: String::new(),
2439 exit_code: Some(1),
2440 timed_out: false,
2441 duration_ms: 431_194,
2442 artifacts: Vec::new(),
2443 quota: None,
2444 dropped: Some(Dropped {
2445 why: "subscriber fell behind updates".to_owned(),
2446 output_tokens: 14267,
2447 }),
2448 commands: Vec::new(),
2449 context_tokens: None,
2450 };
2451 assert!(!out.usable());
2452 assert!(out.work_undelivered());
2453 assert!(!out.quota_exhausted());
2456 }
2457
2458 #[test]
2459 fn non_json_stdout_falls_back_to_raw_text() {
2460 let out = extract(AgentKind::Antigravity, "plain answer\n");
2461 assert_eq!(out.text, "plain answer");
2462 assert!(out.session.is_none());
2463 }
2464
2465 #[tokio::test]
2466 async fn command_agent_round_trip_writes_artifacts() {
2467 let dir = tempfile::tempdir().unwrap();
2468 let art = dir.path().join("artifacts");
2469 let mut seat = SeatState::new("impl-A", "a", 7);
2470 let s = command_helper("reply");
2471 let out = invoke(
2472 &s,
2473 &mut seat,
2474 &Invocation {
2475 cwd: dir.path(),
2476 prompt: "unused",
2477 timeout: Duration::from_secs(30),
2478 allow_write: true,
2479 sessions: true,
2480 artifacts: &art,
2481 stem: "impl-A",
2482 run: "test-run",
2483 node: "test",
2484 cache_dir: None,
2485 attachments: &[],
2486 writable: &[],
2487 },
2488 )
2489 .await
2490 .unwrap();
2491 assert!(out.usable(), "{out:?}");
2492 assert!(out.text.contains("hello impl-A"), "{}", out.text);
2493 assert_eq!(seat.turns, 1);
2494 assert!(art.join("impl-A.prompt.md").is_file());
2495 assert!(art.join("impl-A.out").is_file());
2496 }
2497
2498 #[tokio::test]
2499 async fn the_invocation_cache_dir_reaches_the_seat_as_cargo_target_dir() {
2500 let dir = tempfile::tempdir().unwrap();
2504 let cache = dir.path().join("magi-cache");
2505 let mut seat = SeatState::new("impl-A", "a", 7);
2506 let s = command_helper("cache");
2507 let out = invoke(
2508 &s,
2509 &mut seat,
2510 &Invocation {
2511 cwd: dir.path(),
2512 prompt: "unused",
2513 timeout: Duration::from_secs(30),
2514 allow_write: true,
2515 sessions: true,
2516 artifacts: &dir.path().join("artifacts"),
2517 stem: "cache",
2518 run: "test-run",
2519 node: "test",
2520 cache_dir: Some(&cache),
2521 attachments: &[],
2522 writable: &[],
2523 },
2524 )
2525 .await
2526 .unwrap();
2527 assert!(out.usable(), "{out:?}");
2528 assert!(
2529 out.text.contains(cache.to_string_lossy().as_ref()),
2530 "the seat must see CARGO_TARGET_DIR = the shared cache"
2531 );
2532 }
2533
2534 #[tokio::test]
2535 async fn cache_dir_none_strips_a_cargo_target_dir_inherited_from_this_process() {
2536 let previous = std::env::var("CARGO_TARGET_DIR").ok();
2544 unsafe {
2549 std::env::set_var("CARGO_TARGET_DIR", "/should/never/reach/a/read-only/seat");
2550 }
2551 let dir = tempfile::tempdir().unwrap();
2552 let mut seat = SeatState::new("review-1", "a", 7);
2553 let s = command_helper("no-cache");
2554 let result = invoke(
2555 &s,
2556 &mut seat,
2557 &Invocation {
2558 cwd: dir.path(),
2559 prompt: "unused",
2560 timeout: Duration::from_secs(30),
2561 allow_write: false,
2562 sessions: true,
2563 artifacts: &dir.path().join("artifacts"),
2564 stem: "no-cache",
2565 run: "test-run",
2566 node: "test",
2567 cache_dir: None,
2568 attachments: &[],
2569 writable: &[],
2570 },
2571 )
2572 .await;
2573 unsafe {
2578 match &previous {
2579 Some(v) => std::env::set_var("CARGO_TARGET_DIR", v),
2580 None => std::env::remove_var("CARGO_TARGET_DIR"),
2581 }
2582 }
2583 let out = result.unwrap();
2584 assert!(out.usable(), "{out:?}");
2585 assert!(
2586 out.text.contains("ABSENT"),
2587 "a read-only seat must never inherit the process's own CARGO_TARGET_DIR: {}",
2588 out.text
2589 );
2590 }
2591
2592 #[tokio::test]
2593 async fn a_prompt_larger_than_the_pipe_buffer_does_not_deadlock() {
2594 let dir = tempfile::tempdir().unwrap();
2595 let mut seat = SeatState::new("impl-A", "a", 7);
2596 let s = command_helper("ignore-stdin");
2599 let big = "x".repeat(1_000_000);
2600 let out = invoke(
2601 &s,
2602 &mut seat,
2603 &Invocation {
2604 cwd: dir.path(),
2605 prompt: &big,
2606 timeout: Duration::from_secs(60),
2607 allow_write: true,
2608 sessions: true,
2609 artifacts: &dir.path().join("artifacts"),
2610 stem: "big",
2611 run: "test-run",
2612 node: "test",
2613 cache_dir: None,
2614 attachments: &[],
2615 writable: &[],
2616 },
2617 )
2618 .await
2619 .unwrap();
2620 assert!(out.usable(), "{out:?}");
2621 assert!(out.text.contains("done"), "{}", out.text);
2622 }
2623
2624 #[tokio::test]
2625 async fn timeout_is_reported_not_hung() {
2626 let dir = tempfile::tempdir().unwrap();
2627 let mut seat = SeatState::new("impl-A", "a", 7);
2628 let s = command_helper("sleep");
2629 let out = invoke(
2630 &s,
2631 &mut seat,
2632 &Invocation {
2633 cwd: dir.path(),
2634 prompt: "unused",
2635 timeout: Duration::from_millis(300),
2636 allow_write: true,
2637 sessions: true,
2638 artifacts: &dir.path().join("artifacts"),
2639 stem: "slow",
2640 run: "test-run",
2641 node: "test",
2642 cache_dir: None,
2643 attachments: &[],
2644 writable: &[],
2645 },
2646 )
2647 .await
2648 .unwrap();
2649 assert!(out.timed_out);
2650 assert!(!out.usable());
2651 }
2652
2653 #[tokio::test]
2654 async fn a_timeout_keeps_what_the_agent_had_already_printed() {
2655 let dir = tempfile::tempdir().unwrap();
2661 let artifacts = dir.path().join("artifacts");
2662 let mut seat = SeatState::new("impl-A", "a", 7);
2663 let s = command_helper("chatty-sleep");
2664 let out = invoke(
2665 &s,
2666 &mut seat,
2667 &Invocation {
2668 cwd: dir.path(),
2669 prompt: "unused",
2670 timeout: Duration::from_secs(10),
2675 allow_write: true,
2676 sessions: true,
2677 artifacts: &artifacts,
2678 stem: "chatty",
2679 run: "test-run",
2680 node: "test",
2681 cache_dir: None,
2682 attachments: &[],
2683 writable: &[],
2684 },
2685 )
2686 .await
2687 .unwrap();
2688
2689 assert!(out.timed_out, "{out:?}");
2690 assert!(!out.usable(), "a cut-off answer is still not an answer");
2691 let recorded = std::fs::read_to_string(artifacts.join("chatty.out")).unwrap();
2692 assert!(
2693 recorded.contains("i-said-something"),
2694 "the artifact must keep what arrived before the kill, got {recorded:?}"
2695 );
2696 assert!(
2697 out.text.contains("i-said-something"),
2698 "and the graph must be able to see it too, got {:?}",
2699 out.text
2700 );
2701 }
2702
2703 #[test]
2704 fn missing_programs_reports_command_binaries() {
2705 let mut s = spec(AgentKind::Command, None);
2706 s.command = vec!["definitely-not-a-real-binary-xyz".to_owned()];
2707 assert_eq!(
2708 missing_programs(&[s]),
2709 ["definitely-not-a-real-binary-xyz".to_owned()]
2710 );
2711 }
2712
2713 fn pick_spec(id: &str, kind: AgentKind) -> AgentSpec {
2714 AgentSpec {
2715 id: id.to_owned(),
2716 kind,
2717 model: None,
2718 command: Vec::new(),
2719 extra_args: Vec::new(),
2720 env: BTreeMap::new(),
2721 prompt_delivery: None,
2722 }
2723 }
2724
2725 fn without<'a>(missing: &'a [&'a str]) -> impl Fn(&AgentSpec) -> bool + 'a {
2728 move |a: &AgentSpec| !missing.contains(&a.id.as_str())
2729 }
2730
2731 #[test]
2732 fn pick_prefers_the_claude_seat_even_when_it_is_not_first_in_the_roster() {
2733 let agents = [
2734 pick_spec("oc", AgentKind::Opencode),
2735 pick_spec("opus", AgentKind::Claude),
2736 pick_spec("agy", AgentKind::Antigravity),
2737 ];
2738 let got = pick(&agents, None, &without(&[])).expect("a pick");
2739 assert_eq!(got.id, "opus");
2740 }
2741
2742 #[test]
2743 fn pick_falls_back_to_the_first_installed_agent_in_roster_order() {
2744 let agents = [
2745 pick_spec("opus", AgentKind::Claude),
2746 pick_spec("oc", AgentKind::Opencode),
2747 pick_spec("agy", AgentKind::Antigravity),
2748 ];
2749 let got = pick(&agents, None, &without(&["opus", "oc"])).expect("a pick");
2750 assert_eq!(got.id, "agy");
2751 }
2752
2753 #[test]
2754 fn pick_on_an_empty_roster_says_what_to_install() {
2755 let msg = pick(&[], None, &without(&[]))
2756 .expect_err("nobody to ask")
2757 .to_string();
2758 assert!(msg.contains("roster is empty"), "{msg}");
2759 assert!(msg.contains("claude"), "{msg}");
2760 assert!(msg.contains("magi.toml"), "{msg}");
2761 }
2762
2763 #[test]
2764 fn pick_on_a_roster_with_nothing_installed_names_the_programs_that_are_missing() {
2765 let agents = [
2766 pick_spec("opus", AgentKind::Claude),
2767 pick_spec("oc", AgentKind::Opencode),
2768 ];
2769 let err = pick(&agents, None, &without(&["opus", "oc"])).expect_err("nothing runnable");
2770 let msg = format!("{err:#}");
2771 assert!(msg.contains("claude"), "{msg}");
2772 assert!(msg.contains("opencode"), "{msg}");
2773 }
2774
2775 #[test]
2776 fn an_explicitly_named_agent_wins_over_the_claude_preference() {
2777 let agents = [
2778 pick_spec("opus", AgentKind::Claude),
2779 pick_spec("oc", AgentKind::Opencode),
2780 ];
2781 let got = pick(&agents, Some("oc"), &without(&[])).expect("a pick");
2782 assert_eq!(got.id, "oc");
2783 }
2784
2785 #[test]
2786 fn an_unknown_agent_id_lists_the_ids_that_do_exist() {
2787 let agents = [
2788 pick_spec("opus", AgentKind::Claude),
2789 pick_spec("oc", AgentKind::Opencode),
2790 ];
2791 let msg = pick(&agents, Some("gemini"), &without(&[]))
2792 .expect_err("no such agent")
2793 .to_string();
2794 assert!(msg.contains("gemini"), "{msg}");
2795 assert!(msg.contains("opus, oc"), "{msg}");
2796 }
2797
2798 #[test]
2799 fn an_explicitly_named_agent_that_is_not_installed_is_an_error_not_a_fallback() {
2800 let agents = [
2801 pick_spec("opus", AgentKind::Claude),
2802 pick_spec("oc", AgentKind::Opencode),
2803 ];
2804 let msg = pick(&agents, Some("oc"), &without(&["oc"]))
2805 .expect_err("must not silently substitute another model")
2806 .to_string();
2807 assert!(msg.contains("opencode"), "{msg}");
2808 assert!(msg.contains("--agent"), "{msg}");
2809 }
2810
2811 fn named(id: &str) -> AgentSpec {
2812 AgentSpec {
2813 id: id.to_owned(),
2814 ..spec(AgentKind::Command, None)
2815 }
2816 }
2817
2818 fn output(text: &str, exit: i32, quota: bool) -> AgentOutput {
2819 AgentOutput {
2820 text: text.to_owned(),
2821 exit_code: Some(exit),
2822 timed_out: false,
2823 duration_ms: 0,
2824 artifacts: Vec::new(),
2825 quota: quota.then_some(Quota { reset: None }),
2826 dropped: None,
2827 commands: Vec::new(),
2828 context_tokens: None,
2829 }
2830 }
2831
2832 #[test]
2833 fn a_chain_keeps_the_written_order_and_a_string_is_a_chain_of_one() {
2834 let agents = [named("a"), named("b"), named("c")];
2835 let all = |_: &AgentSpec| true;
2836 let chain = AgentChoice::Chain(vec!["c".into(), "a".into()]);
2837 let got = pick_chain(&agents, Some(&chain), &all, "synthesizer").unwrap();
2838 assert_eq!(
2839 got.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(),
2840 ["c", "a"]
2841 );
2842
2843 let one = AgentChoice::from("b");
2844 let got = pick_chain(&agents, Some(&one), &all, "synthesizer").unwrap();
2845 assert_eq!(got.len(), 1);
2846 assert_eq!(got[0].id, "b");
2847
2848 let got = pick_chain(&agents, None, &all, "synthesizer").unwrap();
2849 assert_eq!(got.len(), 1, "unset keeps pick's default");
2850 assert_eq!(got[0].id, "a");
2851 let empty = AgentChoice::Chain(Vec::new());
2852 assert_eq!(
2853 pick_chain(&agents, Some(&empty), &all, "x").unwrap()[0].id,
2854 "a"
2855 );
2856 }
2857
2858 #[test]
2859 fn a_chain_skips_unknown_and_uninstalled_ids_and_tries_each_once() {
2860 let agents = [named("a"), named("b")];
2861 let not_a = |s: &AgentSpec| s.id != "a";
2862 let chain = AgentChoice::Chain(
2863 ["a", "ghost", "b", "b"]
2864 .iter()
2865 .map(|s| (*s).to_owned())
2866 .collect(),
2867 );
2868 let got = pick_chain(&agents, Some(&chain), ¬_a, "chatter").unwrap();
2869 assert_eq!(got.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(), ["b"]);
2870
2871 let dup = AgentChoice::Chain(vec!["b".into(), "a".into(), "b".into()]);
2872 let got = pick_chain(&agents, Some(&dup), &|_| true, "chatter").unwrap();
2873 assert_eq!(
2874 got.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(),
2875 ["b", "a"]
2876 );
2877 }
2878
2879 #[test]
2880 fn a_chain_with_nothing_runnable_names_the_role() {
2881 let agents = [named("a")];
2882 let chain = AgentChoice::Chain(vec!["a".into(), "ghost".into()]);
2883 let err = pick_chain(&agents, Some(&chain), &|_| false, "conductor")
2884 .unwrap_err()
2885 .to_string();
2886 assert!(err.contains("conductor"), "{err}");
2887 }
2888
2889 #[test]
2890 fn a_chain_advances_on_error_quota_or_an_unusable_answer_only() {
2891 assert!(chain_advances(&Err(anyhow::anyhow!("spawn failed"))));
2892 assert!(chain_advances(&Ok(output("limit", 0, true))));
2893 assert!(chain_advances(&Ok(output("", 0, false))));
2894 assert!(chain_advances(&Ok(output("x", 1, false))));
2895 assert!(!chain_advances(&Ok(output("answer", 0, false))));
2896 }
2897
2898 #[test]
2899 fn claude_context_tokens_are_unknown_because_usage_is_aggregated() {
2900 for out in [
2901 r#"{"result":"ok","num_turns":1,"usage":{"input_tokens":10,"cache_read_input_tokens":3000}}"#,
2902 r#"{"result":"ok","num_turns":3,"usage":{"input_tokens":10}}"#,
2903 r#"{"result":"ok"}"#,
2904 "not json",
2905 ] {
2906 assert_eq!(context_tokens(AgentKind::Claude, out), None, "{out}");
2907 }
2908 }
2909
2910 #[test]
2911 fn opencode_context_tokens_take_the_last_step_finish() {
2912 let out = concat!(
2913 r#"{"type":"step_finish","sessionID":"s","part":{"type":"step-finish","tokens":{"input":100,"output":5,"cache":{"read":1000,"write":50}}}}"#,
2914 "\n",
2915 r#"{"type":"text","part":{"type":"text","text":"hi"}}"#,
2916 "\n",
2917 r#"{"type":"step_finish","sessionID":"s","part":{"type":"step-finish","tokens":{"input":120,"output":9,"cache":{"read":1500}}}}"#,
2918 "\n"
2919 );
2920 assert_eq!(context_tokens(AgentKind::Opencode, out), Some(1620));
2921 let none = r#"{"type":"step_finish","part":{"type":"step-finish"}}"#;
2922 assert_eq!(context_tokens(AgentKind::Opencode, none), None);
2923 assert_eq!(context_tokens(AgentKind::Opencode, ""), None);
2924 }
2925
2926 #[test]
2927 fn agy_context_tokens_are_unknown_because_usage_is_aggregated() {
2928 let out = r#"{"conversation_id":"c","status":"OK","response":"x","usage":{"input_tokens":260113,"cache_read_tokens":2200925}}"#;
2929 assert_eq!(context_tokens(AgentKind::Antigravity, out), None);
2930 }
2931
2932 #[test]
2933 fn codex_context_tokens_are_unknown_because_usage_is_cumulative() {
2934 let out = concat!(
2935 r#"{"type":"turn.completed","usage":{"input_tokens":1000,"cached_input_tokens":900}}"#,
2936 "\n"
2937 );
2938 assert_eq!(context_tokens(AgentKind::Codex, out), None);
2939 }
2940
2941 #[test]
2942 fn omp_context_tokens_come_from_messages_without_an_agent_end() {
2943 let out = concat!(
2944 r#"{"type":"session","id":"s"}"#,
2945 "\n",
2946 r#"{"type":"message_end","message":{"role":"assistant","content":[],"usage":{"input":50,"cacheRead":400,"cacheWrite":10}}}"#,
2947 "\n",
2948 r#"{"type":"turn_end","message":{"role":"assistant","content":[],"usage":{"input":70,"cacheRead":500}}}"#,
2949 "\n"
2950 );
2951 assert_eq!(context_tokens(AgentKind::Omp, out), Some(570));
2952 let user_only = r#"{"type":"message_end","message":{"role":"user","usage":{"input":9}}}"#;
2953 assert_eq!(context_tokens(AgentKind::Omp, user_only), None);
2954 let no_usage = r#"{"type":"message_end","message":{"role":"assistant","content":[]}}"#;
2955 assert_eq!(context_tokens(AgentKind::Omp, no_usage), None);
2956 }
2957
2958 #[test]
2959 fn command_agents_report_no_context_tokens() {
2960 let out = r#"{"usage":{"input_tokens":5}}"#;
2961 assert_eq!(context_tokens(AgentKind::Command, out), None);
2962 }
2963}