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