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