1use std::{
9 collections::BTreeMap,
10 path::PathBuf,
11 process::Stdio,
12 sync::{
13 Arc, Mutex,
14 atomic::{AtomicBool, Ordering},
15 },
16 time::Duration,
17};
18
19use async_trait::async_trait;
20use serde_json::Value;
21use tokio::{
22 io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
23 process::{Child, Command},
24 sync::mpsc,
25};
26
27use crate::{
28 AgentCapabilities, AgentEvent, Mode, PermissionAnswer, RosterSlot, ToolStatus, ToolUpdate,
29};
30
31use super::{
32 AdapterError, AdapterResult, AgentAdapter, CANCEL_SETTLE_TIMEOUT, drain_bounded,
33 isolate_process_group,
34 native::{NativeTurn, spawn_native_turn},
35 parse_command_line, terminate_child,
36};
37
38const MODE_AUTO: &str = "codeswarm:mode:full-access";
39const MODE_PLAN: &str = "codeswarm:mode:plan";
40const MODEL_CONFIG_ID: &str = "codex:model";
41
42#[derive(Debug, Default)]
43struct ParserState {
44 messages: BTreeMap<String, String>,
45 thoughts: BTreeMap<String, String>,
46 tools: BTreeMap<String, ToolUpdate>,
47}
48
49fn text_value(value: Option<&Value>) -> Option<String> {
50 value.and_then(|value| match value {
51 Value::String(text) => (!text.is_empty()).then(|| text.to_owned()),
52 Value::Null => None,
53 Value::Object(object) => object
54 .get("message")
55 .and_then(|message| text_value(Some(message)))
56 .or_else(|| {
57 object
58 .get("detail")
59 .and_then(|detail| text_value(Some(detail)))
60 })
61 .or_else(|| Some(value.to_string())),
62 value => Some(value.to_string()),
63 })
64}
65
66fn item_text(item: &Value) -> Option<String> {
67 if let Some(text) = item.get("text").and_then(Value::as_str) {
68 return (!text.is_empty()).then(|| text.to_owned());
69 }
70 if let Some(delta) = item.get("delta").and_then(Value::as_str) {
71 return (!delta.is_empty()).then(|| delta.to_owned());
72 }
73 let summary = item.get("summary").and_then(Value::as_array)?;
74 let text = summary
75 .iter()
76 .filter_map(|entry| {
77 entry
78 .get("text")
79 .and_then(Value::as_str)
80 .or_else(|| entry.as_str())
81 })
82 .collect::<Vec<_>>()
83 .join("\n");
84 (!text.is_empty()).then_some(text)
85}
86
87fn incremental_text(
88 previous: &mut BTreeMap<String, String>,
89 id: &str,
90 text: String,
91 is_delta: bool,
92) -> Option<String> {
93 if is_delta {
94 previous
95 .entry(id.to_owned())
96 .and_modify(|current| current.push_str(&text))
97 .or_insert_with(|| text.clone());
98 return Some(text);
99 }
100 let current = previous.entry(id.to_owned()).or_default();
101 if current == &text {
102 return None;
103 }
104 let visible = text
105 .strip_prefix(current.as_str())
106 .map_or_else(|| text.clone(), str::to_owned);
107 *current = text;
108 (!visible.is_empty()).then_some(visible)
109}
110
111fn tool_title(item: &Value, kind: &str) -> String {
112 item.get("command")
113 .and_then(Value::as_str)
114 .or_else(|| item.get("name").and_then(Value::as_str))
115 .or_else(|| item.get("tool").and_then(Value::as_str))
116 .filter(|title| !title.is_empty())
117 .map(str::to_owned)
118 .unwrap_or_else(|| kind.strip_suffix("_call").unwrap_or(kind).replace('_', " "))
119}
120
121fn tool_status(item: &Value, event_type: &str) -> ToolStatus {
122 if item
123 .get("exit_code")
124 .and_then(Value::as_i64)
125 .is_some_and(|exit_code| exit_code != 0)
126 {
127 return ToolStatus::Failed;
128 }
129 match item
130 .get("status")
131 .and_then(Value::as_str)
132 .or(Some(event_type))
133 {
134 Some("completed") | Some("success") | Some("item.completed") => ToolStatus::Completed,
135 Some("failed") | Some("error") | Some("errored") | Some("declined") | Some("cancelled")
136 | Some("interrupted") | Some("item.failed") => ToolStatus::Failed,
137 Some("in_progress") | Some("inProgress") | Some("running") | Some("item.started")
138 | Some("item.updated") => ToolStatus::Running,
139 _ => ToolStatus::Pending,
140 }
141}
142
143fn parse_tool(
144 slot: RosterSlot,
145 event_type: &str,
146 item: &Value,
147 state: &mut ParserState,
148) -> Option<AgentEvent> {
149 let kind = item.get("type").and_then(Value::as_str)?;
150 if matches!(kind, "agent_message" | "reasoning") {
151 return None;
152 }
153 let id = item
154 .get("id")
155 .and_then(Value::as_str)
156 .filter(|id| !id.is_empty())?;
157 let title = tool_title(item, kind);
158 let update = state
159 .tools
160 .entry(id.to_owned())
161 .or_insert_with(|| ToolUpdate {
162 id: id.to_owned(),
163 title: title.clone(),
164 status: ToolStatus::Pending,
165 detail: None,
166 });
167 update.title = title;
168 update.status = tool_status(item, event_type);
169 for key in ["aggregated_output", "output", "result", "error", "detail"] {
170 if let Some(detail) = text_value(item.get(key)) {
171 update.detail = Some(detail);
172 break;
173 }
174 }
175 Some(AgentEvent::Tool {
176 slot,
177 update: update.clone(),
178 })
179}
180
181fn parse_value(slot: RosterSlot, value: &Value, state: &mut ParserState) -> Option<AgentEvent> {
184 let event_type = value
185 .get("type")
186 .and_then(Value::as_str)
187 .unwrap_or_default();
188 if matches!(
189 event_type,
190 "item.started" | "item.updated" | "item.completed" | "item.failed"
191 ) {
192 let item = value.get("item")?;
193 let kind = item.get("type").and_then(Value::as_str).unwrap_or_default();
194 if kind == "agent_message" {
195 let id = item
196 .get("id")
197 .and_then(Value::as_str)
198 .unwrap_or("agent-message");
199 let text = item_text(item)?;
200 let is_delta = item.get("delta").is_some() || value.get("delta").is_some();
201 return incremental_text(&mut state.messages, id, text, is_delta)
202 .map(|text| AgentEvent::Text { slot, text });
203 }
204 if kind == "reasoning" {
205 let id = item
206 .get("id")
207 .and_then(Value::as_str)
208 .unwrap_or("reasoning");
209 let text = item_text(item)?;
210 let is_delta = item.get("delta").is_some() || value.get("delta").is_some();
211 return incremental_text(&mut state.thoughts, id, text, is_delta)
212 .map(|text| AgentEvent::Thought { slot, text });
213 }
214 return parse_tool(slot, event_type, item, state);
215 }
216 if event_type.ends_with(".delta")
218 && let Some(delta) = value.get("delta").and_then(Value::as_str)
219 && !delta.is_empty()
220 {
221 return Some(AgentEvent::Text {
222 slot,
223 text: delta.to_owned(),
224 });
225 }
226 None
227}
228
229fn failure_detail(value: &Value) -> Option<String> {
230 ["error", "message", "detail", "reason"]
231 .into_iter()
232 .find_map(|key| text_value(value.get(key)))
233 .or_else(|| {
234 value.get("item").and_then(|item| {
235 ["error", "message", "detail"]
236 .into_iter()
237 .find_map(|key| text_value(item.get(key)))
238 })
239 })
240}
241
242fn thread_id(value: &Value) -> Option<String> {
243 value
244 .get("thread_id")
245 .or_else(|| value.get("threadId"))
246 .and_then(Value::as_str)
247 .filter(|id| !id.is_empty())
248 .map(str::to_owned)
249}
250
251fn cached_models_at(path: &std::path::Path) -> Vec<Mode> {
252 let Ok(contents) = std::fs::read_to_string(path) else {
253 return Vec::new();
254 };
255 let Ok(cache) = serde_json::from_str::<Value>(&contents) else {
256 return Vec::new();
257 };
258 let Some(models) = cache.get("models").and_then(Value::as_array) else {
259 return Vec::new();
260 };
261 let mut catalog = Vec::new();
262 for model in models {
263 if model.get("visibility").and_then(Value::as_str) != Some("list") {
264 continue;
265 }
266 let Some(id) = model
267 .get("slug")
268 .and_then(Value::as_str)
269 .filter(|id| !id.is_empty())
270 else {
271 continue;
272 };
273 if catalog.iter().any(|candidate: &Mode| candidate.id == id) {
274 continue;
275 }
276 let label = model
277 .get("display_name")
278 .and_then(Value::as_str)
279 .filter(|label| !label.is_empty())
280 .unwrap_or(id);
281 catalog.push(Mode {
282 id: id.to_owned(),
283 label: label.to_owned(),
284 });
285 }
286 catalog
287}
288
289fn load_cached_codex_models() -> Vec<Mode> {
290 codex_home()
291 .map(|path| cached_models_at(&path.join("models_cache.json")))
292 .unwrap_or_default()
293}
294
295fn codex_home() -> Option<PathBuf> {
296 std::env::var_os("CODEX_HOME")
297 .filter(|path| !path.is_empty())
298 .map(PathBuf::from)
299 .or_else(|| {
300 std::env::var_os("HOME")
301 .filter(|path| !path.is_empty())
302 .map(|path| PathBuf::from(path).join(".codex"))
303 })
304}
305
306#[derive(Debug, Default)]
307struct CodexDiscovery {
308 models: Vec<Mode>,
309 current_model: Option<String>,
310 config_received: bool,
311 models_received: bool,
312 thread_received: bool,
313}
314
315fn apply_discovery_response(value: &Value, discovery: &mut CodexDiscovery) {
316 match value.get("id").and_then(Value::as_u64) {
317 Some(2) => {
318 discovery.config_received = true;
319 discovery.current_model = value
320 .pointer("/result/config/model")
321 .and_then(Value::as_str)
322 .filter(|model| !model.is_empty())
323 .map(str::to_owned);
324 }
325 Some(3) => {
326 discovery.models_received = true;
327 let Some(models) = value.pointer("/result/data").and_then(Value::as_array) else {
328 return;
329 };
330 discovery.models.clear();
331 for model in models {
332 if model.get("hidden").and_then(Value::as_bool) == Some(true) {
333 continue;
334 }
335 let Some(id) = model
336 .get("id")
337 .or_else(|| model.get("model"))
338 .and_then(Value::as_str)
339 .filter(|id| !id.is_empty())
340 else {
341 continue;
342 };
343 if discovery.models.iter().any(|candidate| candidate.id == id) {
344 continue;
345 }
346 let label = model
347 .get("displayName")
348 .and_then(Value::as_str)
349 .filter(|label| !label.is_empty())
350 .unwrap_or(id);
351 discovery.models.push(Mode {
352 id: id.to_owned(),
353 label: label.to_owned(),
354 });
355 }
356 }
357 Some(4) => {
358 discovery.thread_received = true;
359 if let Some(model) = value
360 .pointer("/result/thread/model")
361 .and_then(Value::as_str)
362 .filter(|model| !model.is_empty())
363 {
364 discovery.current_model = Some(model.to_owned());
365 }
366 }
367 _ => {}
368 }
369}
370
371async fn discover_codex(
372 command_line: &str,
373 cwd: &std::path::Path,
374 session_id: Option<&str>,
375) -> Option<CodexDiscovery> {
376 let (program, args) = parse_command_line(command_line).ok()?;
377 let executable = std::path::Path::new(&program)
378 .file_name()
379 .and_then(|name| name.to_str())?;
380 if !matches!(executable, "codex" | "codex.exe") {
381 return None;
382 }
383 let mut command = Command::new(program);
384 isolate_process_group(&mut command);
385 command
386 .args(args)
387 .arg("app-server")
388 .current_dir(cwd)
389 .stdin(Stdio::piped())
390 .stdout(Stdio::piped())
391 .stderr(Stdio::null());
392 let mut child = command.spawn().ok()?;
393 let mut stdin = child.stdin.take()?;
394 let stdout = child.stdout.take()?;
395 let initialize = serde_json::json!({
396 "method": "initialize",
397 "id": 1,
398 "params": {
399 "clientInfo": {"name": "codeswarm", "title": "CodeSwarm", "version": env!("CARGO_PKG_VERSION")},
400 "capabilities": null
401 }
402 });
403 let mut requests = vec![
404 initialize,
405 serde_json::json!({"method": "initialized", "params": {}}),
406 serde_json::json!({"method": "config/read", "id": 2, "params": {"includeLayers": false, "cwd": cwd}}),
407 serde_json::json!({"method": "model/list", "id": 3, "params": {"cursor": null, "limit": null}}),
408 ];
409 if let Some(session_id) = session_id {
410 requests.push(serde_json::json!({
411 "method": "thread/read",
412 "id": 4,
413 "params": {"threadId": session_id, "includeTurns": false}
414 }));
415 }
416 for request in requests {
417 if stdin
418 .write_all(request.to_string().as_bytes())
419 .await
420 .is_err()
421 || stdin.write_all(b"\n").await.is_err()
422 {
423 let _ = terminate_child(&mut child).await;
424 return None;
425 }
426 }
427 let mut lines = BufReader::new(stdout).lines();
428 let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
429 let mut discovery = CodexDiscovery::default();
430 loop {
431 let complete = discovery.config_received
432 && discovery.models_received
433 && (session_id.is_none() || discovery.thread_received);
434 if complete {
435 break;
436 }
437 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
438 if remaining.is_zero() {
439 break;
440 }
441 let Ok(Ok(Some(line))) = tokio::time::timeout(remaining, lines.next_line()).await else {
442 break;
443 };
444 if let Ok(value) = serde_json::from_str::<Value>(&line) {
445 apply_discovery_response(&value, &mut discovery);
446 }
447 }
448 drop(stdin);
449 let _ = terminate_child(&mut child).await;
450 discovery.models_received.then_some(discovery)
451}
452
453#[derive(Debug)]
455pub struct CodexAdapter {
456 slot: RosterSlot,
457 cwd: PathBuf,
458 command: String,
459 mode: String,
460 model: Option<String>,
461 model_overridden: bool,
462 models: Vec<Mode>,
463 session_id: Option<String>,
464 child: Option<Child>,
465 sender: mpsc::Sender<AdapterResult<AgentEvent>>,
466 receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
467 announced_session: Arc<Mutex<Option<String>>>,
468 cancel_requested: Arc<AtomicBool>,
469}
470
471impl CodexAdapter {
472 pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
473 let (sender, receiver) = mpsc::channel(256);
474 Self {
475 slot,
476 cwd,
477 command: command.into(),
478 mode: MODE_AUTO.into(),
479 model: None,
480 model_overridden: false,
481 models: load_cached_codex_models(),
482 session_id: None,
483 child: None,
484 sender,
485 receiver,
486 announced_session: Arc::new(Mutex::new(None)),
487 cancel_requested: Arc::new(AtomicBool::new(false)),
488 }
489 }
490
491 pub fn with_session_id(
492 slot: RosterSlot,
493 cwd: PathBuf,
494 command: impl Into<String>,
495 session_id: impl Into<String>,
496 ) -> Self {
497 let mut adapter = Self::new(slot, cwd, command);
498 adapter.session_id = Some(session_id.into());
499 adapter
500 }
501
502 fn modes() -> Vec<Mode> {
503 vec![
504 Mode {
505 id: MODE_AUTO.into(),
506 label: "Auto pilot".into(),
507 },
508 Mode {
509 id: MODE_PLAN.into(),
510 label: "Plan".into(),
511 },
512 ]
513 }
514
515 fn retain_selected_model(&mut self) {
516 let Some(model) = self.model.as_ref() else {
517 return;
518 };
519 if !self.models.iter().any(|candidate| candidate.id == *model) {
520 self.models.push(Mode {
521 id: model.clone(),
522 label: model.clone(),
523 });
524 }
525 }
526
527 async fn emit(&self, event: AdapterResult<AgentEvent>) {
528 let _ = self.sender.send(event).await;
529 }
530
531 fn append_mode_flags(&self, command: &mut Command, fresh: bool) {
532 if self.mode == MODE_AUTO {
536 command.arg("--dangerously-bypass-approvals-and-sandbox");
537 } else if self.mode == MODE_PLAN {
538 if fresh {
539 command.arg("--sandbox").arg("read-only");
540 } else {
541 command.arg("-c").arg("sandbox_mode=\"read-only\"");
544 }
545 }
546 }
547}
548
549#[async_trait]
550impl AgentAdapter for CodexAdapter {
551 fn slot(&self) -> RosterSlot {
552 self.slot
553 }
554
555 fn session_id(&self) -> Option<String> {
556 self.session_id.clone()
557 }
558
559 fn protocol(&self) -> &'static str {
560 "native"
561 }
562
563 fn capabilities(&self) -> AgentCapabilities {
564 AgentCapabilities {
565 supports_cancel: true,
566 supports_modes: true,
567 supports_permissions: false,
568 supports_terminals: false,
569 supports_session_load: true,
570 supports_models: !self.models.is_empty(),
571 }
572 }
573
574 async fn start(&mut self) -> AdapterResult<()> {
575 if self.child.is_some() {
576 self.stop().await?;
577 }
578 self.cancel_requested.store(false, Ordering::Release);
579 if let Some(discovery) =
580 discover_codex(&self.command, &self.cwd, self.session_id.as_deref()).await
581 {
582 self.models = discovery.models;
583 if !self.model_overridden {
584 self.model = discovery.current_model;
585 }
586 } else {
587 let refreshed_models = load_cached_codex_models();
588 if !refreshed_models.is_empty() {
589 self.models = refreshed_models;
590 }
591 if !self.model_overridden {
592 self.model = None;
593 }
594 }
595 self.retain_selected_model();
596 self.emit(Ok(AgentEvent::ModesReplaced {
597 slot: self.slot,
598 modes: Self::modes(),
599 current_mode: Some(self.mode.clone()),
600 }))
601 .await;
602 if !self.models.is_empty() {
603 self.emit(Ok(AgentEvent::ModelsReplaced {
604 slot: self.slot,
605 config_id: MODEL_CONFIG_ID.into(),
606 models: self.models.clone(),
607 current_model: self.model.clone(),
608 }))
609 .await;
610 }
611 self.emit(Ok(AgentEvent::Ready {
612 slot: self.slot,
613 capabilities: self.capabilities(),
614 }))
615 .await;
616 Ok(())
617 }
618
619 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
620 if self.child.is_some() {
621 return Err(AdapterError::Transport(
622 "agent is already handling a turn".into(),
623 ));
624 }
625 self.cancel_requested.store(false, Ordering::Release);
626 let fresh = self.session_id.is_none();
627 let (program, args) = parse_command_line(&self.command)
628 .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
629 let mut command = Command::new(program);
630 isolate_process_group(&mut command);
631 command.args(args).arg("exec");
632 if !fresh {
633 command.arg("resume");
634 }
635 command
636 .arg("--json")
637 .arg("-c")
642 .arg("show_raw_agent_reasoning=true")
643 .arg("-c")
644 .arg("model_reasoning_summary=\"detailed\"");
645 if self.model_overridden
646 && let Some(model) = &self.model
647 {
648 command.arg("--model").arg(model);
649 }
650 self.append_mode_flags(&mut command, fresh);
651 if !fresh && let Some(session_id) = &self.session_id {
652 command.arg(session_id);
653 }
654 command
655 .arg("-")
656 .current_dir(&self.cwd)
657 .env("CODESWARM_CWD", &self.cwd);
658 let NativeTurn {
659 child,
660 stdout,
661 stderr,
662 } = spawn_native_turn(command, prompt).await?;
663 let sender = self.sender.clone();
664 let slot = self.slot;
665 let announced_session = Arc::clone(&self.announced_session);
666 let cancel_requested = Arc::clone(&self.cancel_requested);
667 tokio::spawn(async move {
668 let stderr_task = tokio::spawn(drain_bounded(stderr, 32 * 1024));
669 let mut lines = BufReader::new(stdout).lines();
670 let mut state = ParserState::default();
671 let mut turn_completed = false;
672 let mut failure = None;
673 while let Ok(Some(line)) = lines.next_line().await {
674 let Ok(value) = serde_json::from_str::<Value>(&line) else {
675 continue;
676 };
677 let event_type = value
678 .get("type")
679 .and_then(Value::as_str)
680 .unwrap_or_default();
681 if event_type == "thread.started"
682 && let Some(id) = thread_id(&value)
683 && let Ok(mut announced) = announced_session.lock()
684 {
685 *announced = Some(id);
686 }
687 if event_type == "turn.completed" {
688 turn_completed = true;
689 }
690 if event_type == "turn.failed" || event_type == "error" {
691 failure = failure_detail(&value);
692 }
693 if let Some(event) = parse_value(slot, &value, &mut state)
694 && sender.send(Ok(event)).await.is_err()
695 {
696 break;
697 }
698 }
699 let stderr = stderr_task.await.ok().unwrap_or_default();
700 if turn_completed || cancel_requested.load(Ordering::Acquire) {
701 let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
702 } else {
703 let detail = failure
704 .or_else(|| (!stderr.is_empty()).then_some(stderr))
705 .unwrap_or_else(|| "Codex stream ended before a successful turn".into());
706 let _ = sender
707 .send(Ok(AgentEvent::Failed {
708 slot,
709 started: true,
710 detail,
711 }))
712 .await;
713 }
714 });
715 self.child = Some(child);
716 Ok(())
717 }
718
719 async fn cancel(&mut self) -> AdapterResult<bool> {
720 self.cancel_requested.store(true, Ordering::Release);
721 let Some(mut child) = self.child.take() else {
722 return Ok(false);
723 };
724 terminate_child(&mut child).await?;
725 let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
726 while let Some(event) = self.receiver.recv().await {
727 if matches!(
728 event,
729 Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
730 ) {
731 break;
732 }
733 }
734 })
735 .await;
736 Ok(true)
737 }
738
739 async fn answer_permission(
740 &mut self,
741 _request_id: String,
742 _answer: PermissionAnswer,
743 ) -> AdapterResult<()> {
744 Err(AdapterError::Unsupported("permission answer"))
745 }
746
747 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
748 let mode = match mode.as_str() {
749 "full-access" | "auto" | "autopilot" | MODE_AUTO => MODE_AUTO,
750 "plan" | "readonly" | MODE_PLAN => MODE_PLAN,
751 _ => return Err(AdapterError::Unsupported("requested Codex mode")),
752 };
753 self.mode = mode.into();
754 self.emit(Ok(AgentEvent::ModesReplaced {
755 slot: self.slot,
756 modes: Self::modes(),
757 current_mode: Some(self.mode.clone()),
758 }))
759 .await;
760 Ok(())
761 }
762
763 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
764 let model = model.trim();
765 if model.is_empty() {
766 return Err(AdapterError::Protocol("model must not be empty".into()));
767 }
768 self.model = Some(model.to_owned());
769 self.model_overridden = true;
770 self.retain_selected_model();
771 self.emit(Ok(AgentEvent::ModelsReplaced {
772 slot: self.slot,
773 config_id: MODEL_CONFIG_ID.into(),
774 models: self.models.clone(),
775 current_model: self.model.clone(),
776 }))
777 .await;
778 Ok(())
779 }
780
781 async fn reload(&mut self) -> AdapterResult<()> {
782 let session_id = self.session_id.clone();
783 self.stop().await?;
784 self.session_id = session_id;
785 self.start().await
786 }
787
788 async fn stop(&mut self) -> AdapterResult<()> {
789 let _ = self.cancel().await?;
790 while self.receiver.try_recv().is_ok() {}
791 Ok(())
792 }
793
794 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
795 let event = self.receiver.recv().await;
796 if matches!(
797 event.as_ref(),
798 Some(Ok(
799 AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. }
800 ))
801 ) {
802 if self.session_id.is_none()
803 && let Ok(session) = self.announced_session.lock()
804 {
805 self.session_id = session.clone();
806 }
807 if let Some(mut child) = self.child.take() {
808 let _ = child.wait().await;
809 }
810 }
811 event
812 }
813}
814
815#[cfg(test)]
816mod tests {
817 use super::{
818 CodexAdapter, CodexDiscovery, ParserState, apply_discovery_response, cached_models_at,
819 parse_value,
820 };
821 use crate::{AgentAdapter, AgentEvent, ToolStatus};
822 use serde_json::json;
823
824 fn unique_test_path(stem: &str) -> std::path::PathBuf {
825 let nonce = std::time::SystemTime::now()
826 .duration_since(std::time::UNIX_EPOCH)
827 .expect("clock")
828 .as_nanos();
829 std::env::temp_dir().join(format!("{stem}-{}-{nonce}", std::process::id()))
830 }
831
832 async fn start_adapter(adapter: &mut CodexAdapter) {
833 adapter.start().await.expect("start");
834 assert!(matches!(
835 adapter.next_event().await,
836 Some(Ok(AgentEvent::ModesReplaced { modes, current_mode, .. }))
837 if modes.len() == 2
838 && modes.iter().any(|mode| mode.label == "Auto pilot")
839 && modes.iter().any(|mode| mode.label == "Plan")
840 && current_mode.as_deref() == Some("codeswarm:mode:full-access")
841 ));
842 let event = adapter.next_event().await;
843 if matches!(event, Some(Ok(AgentEvent::ModelsReplaced { .. }))) {
844 assert!(matches!(
845 adapter.next_event().await,
846 Some(Ok(AgentEvent::Ready { .. }))
847 ));
848 } else {
849 assert!(matches!(event, Some(Ok(AgentEvent::Ready { .. }))));
850 }
851 }
852
853 #[test]
854 fn reads_visible_models_from_codex_cache() {
855 let cache_path = unique_test_path("codeswarm-codex-model-cache");
856 std::fs::write(
857 &cache_path,
858 r#"{"models":[
859 {"slug":"gpt-visible","display_name":"GPT Visible","visibility":"list"},
860 {"slug":"gpt-hidden","display_name":"GPT Hidden","visibility":"hide"},
861 {"slug":"gpt-unlisted"},
862 {"slug":"gpt-visible","display_name":"Duplicate"},
863 {"display_name":"Missing slug"}
864 ]}"#,
865 )
866 .expect("cache");
867 let models = cached_models_at(&cache_path);
868 assert_eq!(models.len(), 1);
869 assert_eq!(models[0].id, "gpt-visible");
870 assert_eq!(models[0].label, "GPT Visible");
871 std::fs::remove_file(cache_path).expect("cleanup");
872 }
873
874 #[test]
875 fn parses_authoritative_codex_model_discovery() {
876 let mut discovery = CodexDiscovery::default();
877 apply_discovery_response(
878 &json!({"id": 2, "result": {"config": {"model": "gpt-config"}}}),
879 &mut discovery,
880 );
881 apply_discovery_response(
882 &json!({"id": 3, "result": {"data": [
883 {"id": "gpt-config", "displayName": "GPT Config", "hidden": false},
884 {"id": "gpt-hidden", "displayName": "GPT Hidden", "hidden": true},
885 {"model": "gpt-fallback", "displayName": "", "hidden": false}
886 ], "nextCursor": null}}),
887 &mut discovery,
888 );
889 assert_eq!(discovery.current_model.as_deref(), Some("gpt-config"));
890 assert_eq!(discovery.models.len(), 2);
891 assert_eq!(discovery.models[1].label, "gpt-fallback");
892 apply_discovery_response(
893 &json!({"id": 4, "result": {"thread": {"model": "gpt-resumed"}}}),
894 &mut discovery,
895 );
896 assert_eq!(discovery.current_model.as_deref(), Some("gpt-resumed"));
897 }
898
899 #[cfg(unix)]
900 #[tokio::test]
901 async fn startup_uses_app_server_catalog_and_resumed_thread_model() {
902 use std::os::unix::fs::PermissionsExt;
903
904 let directory = unique_test_path("codeswarm-codex-discovery");
905 std::fs::create_dir_all(&directory).expect("directory");
906 let command_path = directory.join("codex");
907 std::fs::write(
908 &command_path,
909 r#"#!/bin/sh
910if [ "$1" = "app-server" ]; then
911 printf '%s\n' \
912 '{"id":1,"result":{}}' \
913 '{"id":2,"result":{"config":{"model":"gpt-config"}}}' \
914 '{"id":3,"result":{"data":[{"id":"gpt-config","displayName":"GPT Config","hidden":false},{"id":"gpt-hidden","displayName":"GPT Hidden","hidden":true}],"nextCursor":null}}' \
915 '{"id":4,"result":{"thread":{"model":"gpt-resumed"}}}'
916 sleep 10
917fi
918"#,
919 )
920 .expect("script");
921 let mut permissions = std::fs::metadata(&command_path)
922 .expect("metadata")
923 .permissions();
924 permissions.set_mode(0o700);
925 std::fs::set_permissions(&command_path, permissions).expect("permissions");
926
927 let mut adapter = CodexAdapter::with_session_id(
928 0,
929 std::env::current_dir().expect("cwd"),
930 command_path.to_string_lossy(),
931 "resume-thread",
932 );
933 adapter.start().await.expect("start");
934 assert!(matches!(
935 adapter.next_event().await,
936 Some(Ok(AgentEvent::ModesReplaced { .. }))
937 ));
938 assert!(matches!(
939 adapter.next_event().await,
940 Some(Ok(AgentEvent::ModelsReplaced { models, current_model, .. }))
941 if current_model.as_deref() == Some("gpt-resumed")
942 && models.len() == 2
943 && models[0].id == "gpt-config"
944 && models[1].id == "gpt-resumed"
945 ));
946 assert!(matches!(
947 adapter.next_event().await,
948 Some(Ok(AgentEvent::Ready { .. }))
949 ));
950 adapter.stop().await.expect("stop");
951 std::fs::remove_dir_all(directory).expect("cleanup");
952 }
953
954 #[tokio::test]
955 async fn exposes_only_noninteractive_codex_modes() {
956 let mut adapter = CodexAdapter::new(0, std::env::current_dir().expect("cwd"), "codex");
957 start_adapter(&mut adapter).await;
958 assert!(adapter.set_mode("manual".into()).await.is_err());
959 assert!(adapter.set_mode("accept-edits".into()).await.is_err());
960 adapter.set_mode("plan".into()).await.expect("plan mode");
961 assert!(matches!(
962 adapter.next_event().await,
963 Some(Ok(AgentEvent::ModesReplaced { modes, current_mode, .. }))
964 if modes.len() == 2
965 && current_mode.as_deref() == Some("codeswarm:mode:plan")
966 ));
967 }
968
969 #[test]
970 fn parses_codex_item_lifecycle_and_deduplicates_snapshots() {
971 let mut state = ParserState::default();
972 assert!(
973 parse_value(
974 1,
975 &json!({"type":"thread.started","thread_id":"t1"}),
976 &mut state
977 )
978 .is_none()
979 );
980 assert!(
981 parse_value(
982 1,
983 &json!({"type":"item.started","item":{"id":"m1","type":"agent_message"}}),
984 &mut state
985 )
986 .is_none()
987 );
988 assert!(matches!(
989 parse_value(1, &json!({"type":"item.updated","item":{"id":"m1","type":"agent_message","text":"Hello"}}), &mut state),
990 Some(AgentEvent::Text { slot: 1, text }) if text == "Hello"
991 ));
992 assert!(parse_value(1, &json!({"type":"item.completed","item":{"id":"m1","type":"agent_message","text":"Hello"}}), &mut state).is_none());
993 assert!(matches!(
994 parse_value(1, &json!({"type":"item.completed","item":{"id":"r1","type":"reasoning","summary":[{"type":"summary_text","text":"Checked the patch"}]}}), &mut state),
995 Some(AgentEvent::Thought { text, .. }) if text == "Checked the patch"
996 ));
997 assert!(matches!(
998 parse_value(1, &json!({"type":"item.started","item":{"id":"c1","type":"command_execution","command":"cargo test","status":"in_progress"}}), &mut state),
999 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Running && update.title == "cargo test" && update.detail.is_none()
1000 ));
1001 assert!(matches!(
1002 parse_value(1, &json!({"type":"item.completed","item":{"id":"c1","type":"command_execution","status":"completed","aggregated_output":"ok"}}), &mut state),
1003 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Completed && update.detail.as_deref() == Some("ok")
1004 ));
1005 assert!(matches!(
1006 parse_value(1, &json!({"type":"item.completed","item":{"id":"c2","type":"command_execution","exit_code":1,"aggregated_output":"command failed"}}), &mut state),
1007 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Failed && update.detail.as_deref() == Some("command failed")
1008 ));
1009 assert!(matches!(
1010 parse_value(1, &json!({"type":"item.failed","item":{"id":"c3","type":"command_execution","status":"error","error":"spawn failed"}}), &mut state),
1011 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Failed && update.detail.as_deref() == Some("spawn failed")
1012 ));
1013 assert!(matches!(
1014 parse_value(1, &json!({"type":"item.updated","item":{"id":"c4","type":"command_execution","status":"inProgress"}}), &mut state),
1015 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Running
1016 ));
1017 assert!(parse_value(1, &json!({"type":"turn.completed"}), &mut state).is_none());
1018 assert!(
1019 parse_value(
1020 1,
1021 &json!({"type":"turn.failed","error":{"message":"rate limit"}}),
1022 &mut state
1023 )
1024 .is_none()
1025 );
1026 }
1027
1028 #[tokio::test]
1029 async fn native_codex_process_captures_thread_and_resumes_it() {
1030 let args_path = unique_test_path("codeswarm-codex-args");
1031 let prompts_path = unique_test_path("codeswarm-codex-prompts");
1032 let script_path = unique_test_path("codeswarm-codex-script");
1033 let script = format!(
1034 r#"printf '%s\n' "$*" >> '{}'
1035cat >> '{}'
1036printf '%s\n' '{{"type":"thread.started","thread_id":"thread-native"}}' '{{"type":"item.completed","item":{{"id":"m1","type":"agent_message","text":"hello"}}}}' '{{"type":"turn.completed"}}'
1037"#,
1038 args_path.display(),
1039 prompts_path.display(),
1040 );
1041 std::fs::write(&script_path, script).expect("script");
1042 let cwd = std::env::current_dir().expect("cwd");
1043 let mut adapter = CodexAdapter::new(0, cwd, format!("sh {}", script_path.display()));
1044 start_adapter(&mut adapter).await;
1045 adapter
1046 .send_prompt("first".into())
1047 .await
1048 .expect("first prompt");
1049 assert!(
1050 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "hello")
1051 );
1052 assert!(matches!(
1053 adapter.next_event().await,
1054 Some(Ok(AgentEvent::TurnComplete { .. }))
1055 ));
1056 assert_eq!(adapter.session_id(), Some("thread-native".into()));
1057 adapter.set_mode("plan".into()).await.expect("plan mode");
1058 assert!(matches!(
1059 adapter.next_event().await,
1060 Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. }))
1061 if mode == "codeswarm:mode:plan"
1062 ));
1063 adapter
1064 .send_prompt("follow-up".into())
1065 .await
1066 .expect("resume prompt");
1067 assert!(matches!(
1068 adapter.next_event().await,
1069 Some(Ok(AgentEvent::Text { .. }))
1070 ));
1071 assert!(matches!(
1072 adapter.next_event().await,
1073 Some(Ok(AgentEvent::TurnComplete { .. }))
1074 ));
1075 let args = std::fs::read_to_string(&args_path).expect("captured arguments");
1076 assert!(
1077 args.lines()
1078 .any(|line| line.contains("exec --json") && line.ends_with(" -"))
1079 && args.lines().any(|line| line.contains("exec resume --json")
1080 && line.contains("thread-native")
1081 && line.ends_with(" -"))
1082 && args.lines().any(|line| {
1083 line.contains("-c sandbox_mode=\"read-only\"")
1084 && line.contains("exec resume --json")
1085 })
1086 && !args.contains("--model")
1087 && !args.contains("first")
1088 && !args.contains("follow-up"),
1089 "{args}"
1090 );
1091 assert_eq!(
1092 std::fs::read_to_string(&prompts_path).expect("captured prompts"),
1093 "firstfollow-up"
1094 );
1095 adapter.stop().await.expect("stop");
1096 std::fs::remove_file(args_path).expect("cleanup");
1097 std::fs::remove_file(prompts_path).expect("cleanup");
1098 std::fs::remove_file(script_path).expect("cleanup");
1099 }
1100
1101 #[tokio::test]
1102 async fn native_codex_forwards_model_and_auto_approval_flags() {
1103 let args_path = unique_test_path("codeswarm-codex-model");
1104 let prompt_path = unique_test_path("codeswarm-codex-model-prompt");
1105 let script_path = unique_test_path("codeswarm-codex-model-script");
1106 let script = format!(
1107 r#"printf '%s\n' "$*" > '{}'
1108cat > '{}'
1109printf '%s\n' '{{"type":"thread.started","thread_id":"thread-model"}}' '{{"type":"turn.completed"}}'
1110"#,
1111 args_path.display(),
1112 prompt_path.display(),
1113 );
1114 std::fs::write(&script_path, script).expect("script");
1115 let mut adapter = CodexAdapter::new(
1116 0,
1117 std::env::current_dir().expect("cwd"),
1118 format!("sh {}", script_path.display()),
1119 );
1120 start_adapter(&mut adapter).await;
1121 adapter.set_model("gpt-test".into()).await.expect("model");
1122 assert!(matches!(
1123 adapter.next_event().await,
1124 Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
1125 if config_id == "codex:model"
1126 && models.iter().any(|model| model.id == "gpt-test")
1127 && current_model.as_deref() == Some("gpt-test")
1128 ));
1129 let prompt = "task with\nmultiple lines\nand leading -flags";
1130 adapter.send_prompt(prompt.into()).await.expect("prompt");
1131 assert!(matches!(
1132 adapter.next_event().await,
1133 Some(Ok(AgentEvent::TurnComplete { .. }))
1134 ));
1135 let args = std::fs::read_to_string(&args_path).expect("captured arguments");
1136 assert!(args.contains("--model gpt-test"), "{args}");
1137 assert!(args.contains("show_raw_agent_reasoning=true"), "{args}");
1138 assert!(
1139 args.contains("model_reasoning_summary=\"detailed\""),
1140 "{args}"
1141 );
1142 assert!(args.ends_with(" -\n"), "{args}");
1143 assert!(!args.contains("task with"), "{args}");
1144 assert!(
1145 args.contains("--dangerously-bypass-approvals-and-sandbox"),
1146 "{args}"
1147 );
1148 assert_eq!(
1149 std::fs::read_to_string(&prompt_path).expect("captured prompt"),
1150 prompt
1151 );
1152 adapter.stop().await.expect("stop");
1153 std::fs::remove_file(args_path).expect("cleanup");
1154 std::fs::remove_file(prompt_path).expect("cleanup");
1155 std::fs::remove_file(script_path).expect("cleanup");
1156 }
1157
1158 #[tokio::test]
1159 async fn native_codex_surfaces_turn_failure_with_nested_message() {
1160 let script_path = unique_test_path("codeswarm-codex-failure-script");
1161 std::fs::write(
1162 &script_path,
1163 r#"printf '%s\n' '{"type":"turn.failed","error":{"message":"rate limit"}}'
1164"#,
1165 )
1166 .expect("script");
1167 let mut adapter = CodexAdapter::new(
1168 0,
1169 std::env::current_dir().expect("cwd"),
1170 format!("sh {}", script_path.display()),
1171 );
1172 start_adapter(&mut adapter).await;
1173 adapter.send_prompt("task".into()).await.expect("prompt");
1174 assert!(matches!(
1175 adapter.next_event().await,
1176 Some(Ok(AgentEvent::Failed { detail, started: true, .. })) if detail == "rate limit"
1177 ));
1178 adapter.stop().await.expect("stop");
1179 std::fs::remove_file(script_path).expect("cleanup");
1180 }
1181
1182 #[tokio::test]
1183 async fn native_codex_cancellation_reaps_the_turn_process() {
1184 let mut adapter =
1185 CodexAdapter::new(0, std::env::current_dir().expect("cwd"), "sh -c 'sleep 10'");
1186 start_adapter(&mut adapter).await;
1187 adapter
1188 .send_prompt("long task".into())
1189 .await
1190 .expect("prompt");
1191 assert!(adapter.cancel().await.expect("cancel"));
1192 assert!(adapter.child.is_none());
1193 }
1194}