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    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
172/// Parse one Codex JSONL event into the common adapter event vocabulary.
173/// Unknown events are ignored so newer Codex event kinds do not break turns.
174fn 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    // Keep compatibility with a possible Responses-style top-level delta.
208    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/// Native process-per-turn Codex adapter.
295#[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        // `exec resume --help` does not expose --sandbox or --approve-for-me;
372        // a resumed thread inherits its Codex policy. The bypass flag is
373        // accepted by both forms and is the only deterministic Auto setting.
374        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                // `exec resume` does not expose --sandbox, but its config
381                // override remains available and applies to this turn.
382                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}