Skip to main content

codeswarm_adapters/
codex.rs

1//! Native adapter for the Codex command-line exec protocol.
2//!
3//! Codex's `exec --json` command is a JSONL stream for one turn. A process is
4//! intentionally created per prompt: Codex persists the thread and exposes a
5//! stable thread ID, while `exec resume` restores that thread for the next
6//! prompt. Keeping the process boundary here avoids an ACP/Node bridge.
7
8use 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
181/// Parse one Codex JSONL event into the common adapter event vocabulary.
182/// Unknown events are ignored so newer Codex event kinds do not break turns.
183fn 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    // Keep compatibility with a possible Responses-style top-level delta.
217    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/// Native process-per-turn Codex adapter.
454#[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        // `exec resume --help` does not expose --sandbox or --approve-for-me;
533        // a resumed thread inherits its Codex policy. The bypass flag is
534        // accepted by both forms and is the only deterministic Auto setting.
535        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                // `exec resume` does not expose --sandbox, but its config
542                // override remains available and applies to this turn.
543                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            // `exec --json` reports reasoning token usage but omits reasoning
638            // items under Codex's default `none` summary policy. These
639            // invocation-local overrides make the provider's own reasoning
640            // summaries available to CodeSwarm's Thought event parser.
641            .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}