1use std::{
9 collections::BTreeMap,
10 path::PathBuf,
11 sync::{
12 Arc, Mutex,
13 atomic::{AtomicBool, Ordering},
14 },
15};
16
17use async_trait::async_trait;
18use serde_json::Value;
19use tokio::{
20 io::{AsyncBufReadExt, BufReader},
21 process::{Child, Command},
22 sync::mpsc,
23};
24
25use crate::{
26 AgentCapabilities, AgentEvent, Mode, PermissionAnswer, RosterSlot, ToolStatus, ToolUpdate,
27};
28
29use super::{
30 AdapterError, AdapterResult, AgentAdapter, CANCEL_SETTLE_TIMEOUT, drain_bounded,
31 isolate_process_group,
32 native::{NativeTurn, spawn_native_turn},
33 parse_command_line, terminate_child,
34};
35
36const MODE_AUTO: &str = "codeswarm:mode:full-access";
37const MODE_PLAN: &str = "codeswarm:mode:plan";
38const MODEL_CONFIG_ID: &str = "codex:model";
39
40#[derive(Debug, Default)]
41struct ParserState {
42 messages: BTreeMap<String, String>,
43 thoughts: BTreeMap<String, String>,
44 tools: BTreeMap<String, ToolUpdate>,
45}
46
47fn text_value(value: Option<&Value>) -> Option<String> {
48 value.and_then(|value| match value {
49 Value::String(text) if !text.is_empty() => Some(text.to_owned()),
50 Value::Null => None,
51 Value::Object(object) => object
52 .get("message")
53 .and_then(|message| text_value(Some(message)))
54 .or_else(|| {
55 object
56 .get("detail")
57 .and_then(|detail| text_value(Some(detail)))
58 })
59 .or_else(|| Some(value.to_string())),
60 value => Some(value.to_string()),
61 })
62}
63
64fn item_text(item: &Value) -> Option<String> {
65 if let Some(text) = item.get("text").and_then(Value::as_str) {
66 return (!text.is_empty()).then(|| text.to_owned());
67 }
68 if let Some(delta) = item.get("delta").and_then(Value::as_str) {
69 return (!delta.is_empty()).then(|| delta.to_owned());
70 }
71 let summary = item.get("summary").and_then(Value::as_array)?;
72 let text = summary
73 .iter()
74 .filter_map(|entry| {
75 entry
76 .get("text")
77 .and_then(Value::as_str)
78 .or_else(|| entry.as_str())
79 })
80 .collect::<Vec<_>>()
81 .join("\n");
82 (!text.is_empty()).then_some(text)
83}
84
85fn incremental_text(
86 previous: &mut BTreeMap<String, String>,
87 id: &str,
88 text: String,
89 is_delta: bool,
90) -> Option<String> {
91 if is_delta {
92 previous
93 .entry(id.to_owned())
94 .and_modify(|current| current.push_str(&text))
95 .or_insert_with(|| text.clone());
96 return Some(text);
97 }
98 let current = previous.entry(id.to_owned()).or_default();
99 if current == &text {
100 return None;
101 }
102 let visible = text
103 .strip_prefix(current.as_str())
104 .map_or_else(|| text.clone(), str::to_owned);
105 *current = text;
106 (!visible.is_empty()).then_some(visible)
107}
108
109fn tool_title(item: &Value, kind: &str) -> String {
110 item.get("command")
111 .and_then(Value::as_str)
112 .or_else(|| item.get("name").and_then(Value::as_str))
113 .or_else(|| item.get("tool").and_then(Value::as_str))
114 .filter(|title| !title.is_empty())
115 .map(str::to_owned)
116 .unwrap_or_else(|| kind.strip_suffix("_call").unwrap_or(kind).replace('_', " "))
117}
118
119fn tool_status(item: &Value, event_type: &str) -> ToolStatus {
120 if item
121 .get("exit_code")
122 .and_then(Value::as_i64)
123 .is_some_and(|exit_code| exit_code != 0)
124 {
125 return ToolStatus::Failed;
126 }
127 match item
128 .get("status")
129 .and_then(Value::as_str)
130 .or(Some(event_type))
131 {
132 Some("completed") | Some("success") | Some("item.completed") => ToolStatus::Completed,
133 Some("failed") | Some("error") | Some("errored") | Some("declined") | Some("cancelled")
134 | Some("interrupted") | Some("item.failed") => ToolStatus::Failed,
135 Some("in_progress") | Some("inProgress") | Some("running") | Some("item.started")
136 | Some("item.updated") => ToolStatus::Running,
137 _ => ToolStatus::Pending,
138 }
139}
140
141fn parse_tool(
142 slot: RosterSlot,
143 event_type: &str,
144 item: &Value,
145 state: &mut ParserState,
146) -> Option<AgentEvent> {
147 let kind = item.get("type").and_then(Value::as_str)?;
148 if matches!(kind, "agent_message" | "reasoning") {
149 return None;
150 }
151 let id = item
152 .get("id")
153 .and_then(Value::as_str)
154 .filter(|id| !id.is_empty())?;
155 let title = tool_title(item, kind);
156 let update = state
157 .tools
158 .entry(id.to_owned())
159 .or_insert_with(|| ToolUpdate {
160 id: id.to_owned(),
161 title: title.clone(),
162 status: ToolStatus::Pending,
163 detail: None,
164 });
165 update.title = title;
166 update.status = tool_status(item, event_type);
167 for key in ["aggregated_output", "output", "result", "error", "detail"] {
168 if let Some(detail) = text_value(item.get(key)) {
169 update.detail = Some(detail);
170 break;
171 }
172 }
173 Some(AgentEvent::Tool {
174 slot,
175 update: update.clone(),
176 })
177}
178
179fn parse_value(slot: RosterSlot, value: &Value, state: &mut ParserState) -> Option<AgentEvent> {
182 let event_type = value
183 .get("type")
184 .and_then(Value::as_str)
185 .unwrap_or_default();
186 if matches!(
187 event_type,
188 "item.started" | "item.updated" | "item.completed" | "item.failed"
189 ) {
190 let item = value.get("item")?;
191 let kind = item.get("type").and_then(Value::as_str).unwrap_or_default();
192 if kind == "agent_message" {
193 let id = item
194 .get("id")
195 .and_then(Value::as_str)
196 .unwrap_or("agent-message");
197 let text = item_text(item)?;
198 let is_delta = item.get("delta").is_some() || value.get("delta").is_some();
199 return incremental_text(&mut state.messages, id, text, is_delta)
200 .map(|text| AgentEvent::Text { slot, text });
201 }
202 if kind == "reasoning" {
203 let id = item
204 .get("id")
205 .and_then(Value::as_str)
206 .unwrap_or("reasoning");
207 let text = item_text(item)?;
208 let is_delta = item.get("delta").is_some() || value.get("delta").is_some();
209 return incremental_text(&mut state.thoughts, id, text, is_delta)
210 .map(|text| AgentEvent::Thought { slot, text });
211 }
212 return parse_tool(slot, event_type, item, state);
213 }
214 if event_type.ends_with(".delta")
216 && let Some(delta) = value.get("delta").and_then(Value::as_str)
217 && !delta.is_empty()
218 {
219 return Some(AgentEvent::Text {
220 slot,
221 text: delta.to_owned(),
222 });
223 }
224 None
225}
226
227fn failure_detail(value: &Value) -> Option<String> {
228 ["error", "message", "detail", "reason"]
229 .into_iter()
230 .find_map(|key| text_value(value.get(key)))
231 .or_else(|| {
232 value.get("item").and_then(|item| {
233 ["error", "message", "detail"]
234 .into_iter()
235 .find_map(|key| text_value(item.get(key)))
236 })
237 })
238}
239
240fn thread_id(value: &Value) -> Option<String> {
241 value
242 .get("thread_id")
243 .or_else(|| value.get("threadId"))
244 .and_then(Value::as_str)
245 .filter(|id| !id.is_empty())
246 .map(str::to_owned)
247}
248
249fn cached_models_at(path: &std::path::Path) -> Vec<Mode> {
250 let Ok(contents) = std::fs::read_to_string(path) else {
251 return Vec::new();
252 };
253 let Ok(cache) = serde_json::from_str::<Value>(&contents) else {
254 return Vec::new();
255 };
256 let Some(models) = cache.get("models").and_then(Value::as_array) else {
257 return Vec::new();
258 };
259 let mut catalog = Vec::new();
260 for model in models {
261 if model.get("visibility").and_then(Value::as_str) == Some("hide") {
262 continue;
263 }
264 let Some(id) = model
265 .get("slug")
266 .and_then(Value::as_str)
267 .filter(|id| !id.is_empty())
268 else {
269 continue;
270 };
271 if catalog.iter().any(|candidate: &Mode| candidate.id == id) {
272 continue;
273 }
274 let label = model
275 .get("display_name")
276 .and_then(Value::as_str)
277 .filter(|label| !label.is_empty())
278 .unwrap_or(id);
279 catalog.push(Mode {
280 id: id.to_owned(),
281 label: label.to_owned(),
282 });
283 }
284 catalog
285}
286
287fn load_codex_models() -> Vec<Mode> {
288 let codex_home = std::env::var_os("CODEX_HOME")
289 .filter(|path| !path.is_empty())
290 .map(PathBuf::from)
291 .or_else(|| {
292 std::env::var_os("HOME")
293 .filter(|path| !path.is_empty())
294 .map(|path| PathBuf::from(path).join(".codex"))
295 });
296 codex_home
297 .map(|path| cached_models_at(&path.join("models_cache.json")))
298 .unwrap_or_default()
299}
300
301#[derive(Debug)]
303pub struct CodexAdapter {
304 slot: RosterSlot,
305 cwd: PathBuf,
306 command: String,
307 mode: String,
308 model: Option<String>,
309 models: Vec<Mode>,
310 session_id: Option<String>,
311 child: Option<Child>,
312 sender: mpsc::Sender<AdapterResult<AgentEvent>>,
313 receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
314 announced_session: Arc<Mutex<Option<String>>>,
315 cancel_requested: Arc<AtomicBool>,
316}
317
318impl CodexAdapter {
319 pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
320 let (sender, receiver) = mpsc::channel(256);
321 Self {
322 slot,
323 cwd,
324 command: command.into(),
325 mode: MODE_AUTO.into(),
326 model: None,
327 models: load_codex_models(),
328 session_id: None,
329 child: None,
330 sender,
331 receiver,
332 announced_session: Arc::new(Mutex::new(None)),
333 cancel_requested: Arc::new(AtomicBool::new(false)),
334 }
335 }
336
337 pub fn with_session_id(
338 slot: RosterSlot,
339 cwd: PathBuf,
340 command: impl Into<String>,
341 session_id: impl Into<String>,
342 ) -> Self {
343 let mut adapter = Self::new(slot, cwd, command);
344 adapter.session_id = Some(session_id.into());
345 adapter
346 }
347
348 fn modes() -> Vec<Mode> {
349 vec![
350 Mode {
351 id: MODE_AUTO.into(),
352 label: "Auto pilot".into(),
353 },
354 Mode {
355 id: MODE_PLAN.into(),
356 label: "Plan".into(),
357 },
358 ]
359 }
360
361 fn retain_selected_model(&mut self) {
362 let Some(model) = self.model.as_ref() else {
363 return;
364 };
365 if !self.models.iter().any(|candidate| candidate.id == *model) {
366 self.models.push(Mode {
367 id: model.clone(),
368 label: model.clone(),
369 });
370 }
371 }
372
373 async fn emit(&self, event: AdapterResult<AgentEvent>) {
374 let _ = self.sender.send(event).await;
375 }
376
377 fn append_mode_flags(&self, command: &mut Command, fresh: bool) {
378 if self.mode == MODE_AUTO {
382 command.arg("--dangerously-bypass-approvals-and-sandbox");
383 } else if self.mode == MODE_PLAN {
384 if fresh {
385 command.arg("--sandbox").arg("read-only");
386 } else {
387 command.arg("-c").arg("sandbox_mode=\"read-only\"");
390 }
391 }
392 }
393}
394
395#[async_trait]
396impl AgentAdapter for CodexAdapter {
397 fn slot(&self) -> RosterSlot {
398 self.slot
399 }
400
401 fn session_id(&self) -> Option<String> {
402 self.session_id.clone()
403 }
404
405 fn protocol(&self) -> &'static str {
406 "native"
407 }
408
409 fn capabilities(&self) -> AgentCapabilities {
410 AgentCapabilities {
411 supports_cancel: true,
412 supports_modes: true,
413 supports_permissions: false,
414 supports_terminals: false,
415 supports_session_load: true,
416 supports_models: !self.models.is_empty(),
417 }
418 }
419
420 async fn start(&mut self) -> AdapterResult<()> {
421 if self.child.is_some() {
422 self.stop().await?;
423 }
424 self.cancel_requested.store(false, Ordering::Release);
425 let refreshed_models = load_codex_models();
426 if !refreshed_models.is_empty() {
427 self.models = refreshed_models;
428 }
429 self.retain_selected_model();
430 self.emit(Ok(AgentEvent::ModesReplaced {
431 slot: self.slot,
432 modes: Self::modes(),
433 current_mode: Some(self.mode.clone()),
434 }))
435 .await;
436 if !self.models.is_empty() {
437 self.emit(Ok(AgentEvent::ModelsReplaced {
438 slot: self.slot,
439 config_id: MODEL_CONFIG_ID.into(),
440 models: self.models.clone(),
441 current_model: self.model.clone(),
442 }))
443 .await;
444 }
445 self.emit(Ok(AgentEvent::Ready {
446 slot: self.slot,
447 capabilities: self.capabilities(),
448 }))
449 .await;
450 Ok(())
451 }
452
453 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
454 if self.child.is_some() {
455 return Err(AdapterError::Transport(
456 "agent is already handling a turn".into(),
457 ));
458 }
459 self.cancel_requested.store(false, Ordering::Release);
460 let fresh = self.session_id.is_none();
461 let (program, args) = parse_command_line(&self.command)
462 .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
463 let mut command = Command::new(program);
464 isolate_process_group(&mut command);
465 command.args(args).arg("exec");
466 if !fresh {
467 command.arg("resume");
468 }
469 command.arg("--json");
470 if let Some(model) = &self.model {
471 command.arg("--model").arg(model);
472 }
473 self.append_mode_flags(&mut command, fresh);
474 if !fresh && let Some(session_id) = &self.session_id {
475 command.arg(session_id);
476 }
477 command
478 .arg("-")
479 .current_dir(&self.cwd)
480 .env("CODESWARM_CWD", &self.cwd);
481 let NativeTurn {
482 child,
483 stdout,
484 stderr,
485 } = spawn_native_turn(command, prompt).await?;
486 let sender = self.sender.clone();
487 let slot = self.slot;
488 let announced_session = Arc::clone(&self.announced_session);
489 let cancel_requested = Arc::clone(&self.cancel_requested);
490 tokio::spawn(async move {
491 let stderr_task = tokio::spawn(drain_bounded(stderr, 32 * 1024));
492 let mut lines = BufReader::new(stdout).lines();
493 let mut state = ParserState::default();
494 let mut turn_completed = false;
495 let mut failure = None;
496 while let Ok(Some(line)) = lines.next_line().await {
497 let Ok(value) = serde_json::from_str::<Value>(&line) else {
498 continue;
499 };
500 let event_type = value
501 .get("type")
502 .and_then(Value::as_str)
503 .unwrap_or_default();
504 if event_type == "thread.started"
505 && let Some(id) = thread_id(&value)
506 && let Ok(mut announced) = announced_session.lock()
507 {
508 *announced = Some(id);
509 }
510 if event_type == "turn.completed" {
511 turn_completed = true;
512 }
513 if event_type == "turn.failed" || event_type == "error" {
514 failure = failure_detail(&value);
515 }
516 if let Some(event) = parse_value(slot, &value, &mut state)
517 && sender.send(Ok(event)).await.is_err()
518 {
519 break;
520 }
521 }
522 let stderr = stderr_task.await.ok().unwrap_or_default();
523 if turn_completed || cancel_requested.load(Ordering::Acquire) {
524 let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
525 } else {
526 let detail = failure
527 .or_else(|| (!stderr.is_empty()).then_some(stderr))
528 .unwrap_or_else(|| "Codex stream ended before a successful turn".into());
529 let _ = sender
530 .send(Ok(AgentEvent::Failed {
531 slot,
532 started: true,
533 detail,
534 }))
535 .await;
536 }
537 });
538 self.child = Some(child);
539 Ok(())
540 }
541
542 async fn cancel(&mut self) -> AdapterResult<bool> {
543 self.cancel_requested.store(true, Ordering::Release);
544 let Some(mut child) = self.child.take() else {
545 return Ok(false);
546 };
547 terminate_child(&mut child).await?;
548 let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
549 while let Some(event) = self.receiver.recv().await {
550 if matches!(
551 event,
552 Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
553 ) {
554 break;
555 }
556 }
557 })
558 .await;
559 Ok(true)
560 }
561
562 async fn answer_permission(
563 &mut self,
564 _request_id: String,
565 _answer: PermissionAnswer,
566 ) -> AdapterResult<()> {
567 Err(AdapterError::Unsupported("permission answer"))
568 }
569
570 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
571 let mode = match mode.as_str() {
572 "full-access" | "auto" | "autopilot" | MODE_AUTO => MODE_AUTO,
573 "plan" | "readonly" | MODE_PLAN => MODE_PLAN,
574 _ => return Err(AdapterError::Unsupported("requested Codex mode")),
575 };
576 self.mode = mode.into();
577 self.emit(Ok(AgentEvent::ModesReplaced {
578 slot: self.slot,
579 modes: Self::modes(),
580 current_mode: Some(self.mode.clone()),
581 }))
582 .await;
583 Ok(())
584 }
585
586 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
587 let model = model.trim();
588 if model.is_empty() {
589 return Err(AdapterError::Protocol("model must not be empty".into()));
590 }
591 self.model = Some(model.to_owned());
592 self.retain_selected_model();
593 self.emit(Ok(AgentEvent::ModelsReplaced {
594 slot: self.slot,
595 config_id: MODEL_CONFIG_ID.into(),
596 models: self.models.clone(),
597 current_model: self.model.clone(),
598 }))
599 .await;
600 Ok(())
601 }
602
603 async fn reload(&mut self) -> AdapterResult<()> {
604 let session_id = self.session_id.clone();
605 self.stop().await?;
606 self.session_id = session_id;
607 self.start().await
608 }
609
610 async fn stop(&mut self) -> AdapterResult<()> {
611 let _ = self.cancel().await?;
612 while self.receiver.try_recv().is_ok() {}
613 Ok(())
614 }
615
616 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
617 let event = self.receiver.recv().await;
618 if matches!(
619 event.as_ref(),
620 Some(Ok(
621 AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. }
622 ))
623 ) {
624 if self.session_id.is_none()
625 && let Ok(session) = self.announced_session.lock()
626 {
627 self.session_id = session.clone();
628 }
629 if let Some(mut child) = self.child.take() {
630 let _ = child.wait().await;
631 }
632 }
633 event
634 }
635}
636
637#[cfg(test)]
638mod tests {
639 use super::{CodexAdapter, ParserState, cached_models_at, parse_value};
640 use crate::{AgentAdapter, AgentEvent, ToolStatus};
641 use serde_json::json;
642
643 fn unique_test_path(stem: &str) -> std::path::PathBuf {
644 let nonce = std::time::SystemTime::now()
645 .duration_since(std::time::UNIX_EPOCH)
646 .expect("clock")
647 .as_nanos();
648 std::env::temp_dir().join(format!("{stem}-{}-{nonce}", std::process::id()))
649 }
650
651 async fn start_adapter(adapter: &mut CodexAdapter) {
652 adapter.start().await.expect("start");
653 assert!(matches!(
654 adapter.next_event().await,
655 Some(Ok(AgentEvent::ModesReplaced { modes, current_mode, .. }))
656 if modes.len() == 2
657 && modes.iter().any(|mode| mode.label == "Auto pilot")
658 && modes.iter().any(|mode| mode.label == "Plan")
659 && current_mode.as_deref() == Some("codeswarm:mode:full-access")
660 ));
661 let event = adapter.next_event().await;
662 if matches!(event, Some(Ok(AgentEvent::ModelsReplaced { .. }))) {
663 assert!(matches!(
664 adapter.next_event().await,
665 Some(Ok(AgentEvent::Ready { .. }))
666 ));
667 } else {
668 assert!(matches!(event, Some(Ok(AgentEvent::Ready { .. }))));
669 }
670 }
671
672 #[test]
673 fn reads_visible_models_from_codex_cache() {
674 let cache_path = unique_test_path("codeswarm-codex-model-cache");
675 std::fs::write(
676 &cache_path,
677 r#"{"models":[
678 {"slug":"gpt-visible","display_name":"GPT Visible","visibility":"list"},
679 {"slug":"gpt-hidden","display_name":"GPT Hidden","visibility":"hide"},
680 {"slug":"gpt-fallback"},
681 {"slug":"gpt-visible","display_name":"Duplicate"},
682 {"display_name":"Missing slug"}
683 ]}"#,
684 )
685 .expect("cache");
686 let models = cached_models_at(&cache_path);
687 assert_eq!(models.len(), 2);
688 assert_eq!(models[0].id, "gpt-visible");
689 assert_eq!(models[0].label, "GPT Visible");
690 assert_eq!(models[1].id, "gpt-fallback");
691 assert_eq!(models[1].label, "gpt-fallback");
692 std::fs::remove_file(cache_path).expect("cleanup");
693 }
694
695 #[tokio::test]
696 async fn exposes_only_noninteractive_codex_modes() {
697 let mut adapter = CodexAdapter::new(0, std::env::current_dir().expect("cwd"), "codex");
698 start_adapter(&mut adapter).await;
699 assert!(adapter.set_mode("manual".into()).await.is_err());
700 assert!(adapter.set_mode("accept-edits".into()).await.is_err());
701 adapter.set_mode("plan".into()).await.expect("plan mode");
702 assert!(matches!(
703 adapter.next_event().await,
704 Some(Ok(AgentEvent::ModesReplaced { modes, current_mode, .. }))
705 if modes.len() == 2
706 && current_mode.as_deref() == Some("codeswarm:mode:plan")
707 ));
708 }
709
710 #[test]
711 fn parses_codex_item_lifecycle_and_deduplicates_snapshots() {
712 let mut state = ParserState::default();
713 assert!(
714 parse_value(
715 1,
716 &json!({"type":"thread.started","thread_id":"t1"}),
717 &mut state
718 )
719 .is_none()
720 );
721 assert!(
722 parse_value(
723 1,
724 &json!({"type":"item.started","item":{"id":"m1","type":"agent_message"}}),
725 &mut state
726 )
727 .is_none()
728 );
729 assert!(matches!(
730 parse_value(1, &json!({"type":"item.updated","item":{"id":"m1","type":"agent_message","text":"Hello"}}), &mut state),
731 Some(AgentEvent::Text { slot: 1, text }) if text == "Hello"
732 ));
733 assert!(parse_value(1, &json!({"type":"item.completed","item":{"id":"m1","type":"agent_message","text":"Hello"}}), &mut state).is_none());
734 assert!(matches!(
735 parse_value(1, &json!({"type":"item.completed","item":{"id":"r1","type":"reasoning","summary":[{"type":"summary_text","text":"Checked the patch"}]}}), &mut state),
736 Some(AgentEvent::Thought { text, .. }) if text == "Checked the patch"
737 ));
738 assert!(matches!(
739 parse_value(1, &json!({"type":"item.started","item":{"id":"c1","type":"command_execution","command":"cargo test","status":"in_progress"}}), &mut state),
740 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Running && update.title == "cargo test"
741 ));
742 assert!(matches!(
743 parse_value(1, &json!({"type":"item.completed","item":{"id":"c1","type":"command_execution","status":"completed","aggregated_output":"ok"}}), &mut state),
744 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Completed && update.detail.as_deref() == Some("ok")
745 ));
746 assert!(matches!(
747 parse_value(1, &json!({"type":"item.completed","item":{"id":"c2","type":"command_execution","exit_code":1,"aggregated_output":"command failed"}}), &mut state),
748 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Failed && update.detail.as_deref() == Some("command failed")
749 ));
750 assert!(matches!(
751 parse_value(1, &json!({"type":"item.failed","item":{"id":"c3","type":"command_execution","status":"error","error":"spawn failed"}}), &mut state),
752 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Failed && update.detail.as_deref() == Some("spawn failed")
753 ));
754 assert!(matches!(
755 parse_value(1, &json!({"type":"item.updated","item":{"id":"c4","type":"command_execution","status":"inProgress"}}), &mut state),
756 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Running
757 ));
758 assert!(parse_value(1, &json!({"type":"turn.completed"}), &mut state).is_none());
759 assert!(
760 parse_value(
761 1,
762 &json!({"type":"turn.failed","error":{"message":"rate limit"}}),
763 &mut state
764 )
765 .is_none()
766 );
767 }
768
769 #[tokio::test]
770 async fn native_codex_process_captures_thread_and_resumes_it() {
771 let args_path = unique_test_path("codeswarm-codex-args");
772 let prompts_path = unique_test_path("codeswarm-codex-prompts");
773 let script_path = unique_test_path("codeswarm-codex-script");
774 let script = format!(
775 r#"printf '%s\n' "$*" >> '{}'
776cat >> '{}'
777printf '%s\n' '{{"type":"thread.started","thread_id":"thread-native"}}' '{{"type":"item.completed","item":{{"id":"m1","type":"agent_message","text":"hello"}}}}' '{{"type":"turn.completed"}}'
778"#,
779 args_path.display(),
780 prompts_path.display(),
781 );
782 std::fs::write(&script_path, script).expect("script");
783 let cwd = std::env::current_dir().expect("cwd");
784 let mut adapter = CodexAdapter::new(0, cwd, format!("sh {}", script_path.display()));
785 start_adapter(&mut adapter).await;
786 adapter
787 .send_prompt("first".into())
788 .await
789 .expect("first prompt");
790 assert!(
791 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "hello")
792 );
793 assert!(matches!(
794 adapter.next_event().await,
795 Some(Ok(AgentEvent::TurnComplete { .. }))
796 ));
797 assert_eq!(adapter.session_id(), Some("thread-native".into()));
798 adapter.set_mode("plan".into()).await.expect("plan mode");
799 assert!(matches!(
800 adapter.next_event().await,
801 Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. }))
802 if mode == "codeswarm:mode:plan"
803 ));
804 adapter
805 .send_prompt("follow-up".into())
806 .await
807 .expect("resume prompt");
808 assert!(matches!(
809 adapter.next_event().await,
810 Some(Ok(AgentEvent::Text { .. }))
811 ));
812 assert!(matches!(
813 adapter.next_event().await,
814 Some(Ok(AgentEvent::TurnComplete { .. }))
815 ));
816 let args = std::fs::read_to_string(&args_path).expect("captured arguments");
817 assert!(
818 args.lines()
819 .any(|line| line.contains("exec --json") && line.ends_with(" -"))
820 && args.lines().any(|line| line.contains("exec resume --json")
821 && line.contains("thread-native")
822 && line.ends_with(" -"))
823 && args.lines().any(|line| {
824 line.contains("-c sandbox_mode=\"read-only\"")
825 && line.contains("exec resume --json")
826 })
827 && !args.contains("first")
828 && !args.contains("follow-up"),
829 "{args}"
830 );
831 assert_eq!(
832 std::fs::read_to_string(&prompts_path).expect("captured prompts"),
833 "firstfollow-up"
834 );
835 adapter.stop().await.expect("stop");
836 std::fs::remove_file(args_path).expect("cleanup");
837 std::fs::remove_file(prompts_path).expect("cleanup");
838 std::fs::remove_file(script_path).expect("cleanup");
839 }
840
841 #[tokio::test]
842 async fn native_codex_forwards_model_and_auto_approval_flags() {
843 let args_path = unique_test_path("codeswarm-codex-model");
844 let prompt_path = unique_test_path("codeswarm-codex-model-prompt");
845 let script_path = unique_test_path("codeswarm-codex-model-script");
846 let script = format!(
847 r#"printf '%s\n' "$*" > '{}'
848cat > '{}'
849printf '%s\n' '{{"type":"thread.started","thread_id":"thread-model"}}' '{{"type":"turn.completed"}}'
850"#,
851 args_path.display(),
852 prompt_path.display(),
853 );
854 std::fs::write(&script_path, script).expect("script");
855 let mut adapter = CodexAdapter::new(
856 0,
857 std::env::current_dir().expect("cwd"),
858 format!("sh {}", script_path.display()),
859 );
860 start_adapter(&mut adapter).await;
861 adapter.set_model("gpt-test".into()).await.expect("model");
862 assert!(matches!(
863 adapter.next_event().await,
864 Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
865 if config_id == "codex:model"
866 && models.iter().any(|model| model.id == "gpt-test")
867 && current_model.as_deref() == Some("gpt-test")
868 ));
869 let prompt = "task with\nmultiple lines\nand leading -flags";
870 adapter.send_prompt(prompt.into()).await.expect("prompt");
871 assert!(matches!(
872 adapter.next_event().await,
873 Some(Ok(AgentEvent::TurnComplete { .. }))
874 ));
875 let args = std::fs::read_to_string(&args_path).expect("captured arguments");
876 assert!(args.contains("--model gpt-test"), "{args}");
877 assert!(args.ends_with(" -\n"), "{args}");
878 assert!(!args.contains("task with"), "{args}");
879 assert!(
880 args.contains("--dangerously-bypass-approvals-and-sandbox"),
881 "{args}"
882 );
883 assert_eq!(
884 std::fs::read_to_string(&prompt_path).expect("captured prompt"),
885 prompt
886 );
887 adapter.stop().await.expect("stop");
888 std::fs::remove_file(args_path).expect("cleanup");
889 std::fs::remove_file(prompt_path).expect("cleanup");
890 std::fs::remove_file(script_path).expect("cleanup");
891 }
892
893 #[tokio::test]
894 async fn native_codex_surfaces_turn_failure_with_nested_message() {
895 let script_path = unique_test_path("codeswarm-codex-failure-script");
896 std::fs::write(
897 &script_path,
898 r#"printf '%s\n' '{"type":"turn.failed","error":{"message":"rate limit"}}'
899"#,
900 )
901 .expect("script");
902 let mut adapter = CodexAdapter::new(
903 0,
904 std::env::current_dir().expect("cwd"),
905 format!("sh {}", script_path.display()),
906 );
907 start_adapter(&mut adapter).await;
908 adapter.send_prompt("task".into()).await.expect("prompt");
909 assert!(matches!(
910 adapter.next_event().await,
911 Some(Ok(AgentEvent::Failed { detail, started: true, .. })) if detail == "rate limit"
912 ));
913 adapter.stop().await.expect("stop");
914 std::fs::remove_file(script_path).expect("cleanup");
915 }
916
917 #[tokio::test]
918 async fn native_codex_cancellation_reaps_the_turn_process() {
919 let mut adapter =
920 CodexAdapter::new(0, std::env::current_dir().expect("cwd"), "sh -c 'sleep 10'");
921 start_adapter(&mut adapter).await;
922 adapter
923 .send_prompt("long task".into())
924 .await
925 .expect("prompt");
926 assert!(adapter.cancel().await.expect("cancel"));
927 assert!(adapter.child.is_none());
928 }
929}