Skip to main content

pond/adapter/
codex_cli.rs

1//! OpenAI Codex CLI adapter.
2//!
3//! Source path: `~/.codex/sessions/<year>/<month>/<day>/rollout-<ts>-<uuid>.jsonl`.
4//! Each line is an envelope `{timestamp, type, payload}`. Top-level types:
5//! `session_meta` (consumed up front for Session), `event_msg` /
6//! `turn_context` (transport noise, skipped), `response_item` (the per-turn
7//! model interaction: subtypes `message`, `reasoning`, `function_call`,
8//! `function_call_output`, `custom_tool_call`).
9//!
10//! Pre-Oct-2025 legacy rollouts (spec.md#adapters) predate the envelope: the
11//! first row is a bare metadata object and each data row is an un-enveloped
12//! payload, interleaved with `{record_type:"state"}` noise. The adapter
13//! accepts both shapes.
14
15use std::{
16    collections::HashMap,
17    path::{Path, PathBuf},
18};
19
20use chrono::{DateTime, Datelike, SecondsFormat, Utc};
21use serde_json::{Value, json};
22
23use crate::{
24    sessions::IngestEvent,
25    wire::{Message, Part, PartKind, Provenance, ProviderOptions, Session},
26};
27
28use super::{
29    Adapter, AdapterError, AdapterFactory, AdapterYieldStream, DiscoverFuture, Env,
30    RestoreFidelity, RestoredFile, SkipOracle, by_timestamp_then_id, compact_json, config_path,
31    empty_options,
32    extract::{
33        Extracted, extract_compact_repr, extract_raw_record, extract_self_str, extract_str,
34        json_or_string,
35    },
36    extracted_text,
37    jsonl::{
38        BoundedRow, JsonlTree, jsonl_tree_discover, jsonl_tree_events, peek_first_line,
39        peek_last_mapped,
40    },
41    jsonl_bytes, part_id, part_ordinal, raw_record,
42};
43
44const NAME: &str = "codex-cli";
45
46/// Stateless factory: opens [`CodexCliAdapter`] instances and probes for the
47/// canonical install location under `~/.codex/sessions`.
48pub struct CodexCliFactory;
49
50impl AdapterFactory for CodexCliFactory {
51    fn name(&self) -> &'static str {
52        NAME
53    }
54
55    fn open(&self, config: Value) -> Result<Box<dyn Adapter>, AdapterError> {
56        Ok(Box::new(CodexCliAdapter::new(config_path(NAME, config)?)))
57    }
58
59    fn probe_default(&self, env: &Env) -> Option<Value> {
60        let path = env.home.join(".codex").join("sessions");
61        path.exists().then(|| json!({ "path": path }))
62    }
63
64    fn serialize(
65        &self,
66        session: &crate::sessions::SessionWithMessages,
67        fidelity: RestoreFidelity,
68    ) -> Result<Vec<RestoredFile>, AdapterError> {
69        serialize_session(session, fidelity)
70    }
71}
72
73fn serialize_session(
74    session: &crate::sessions::SessionWithMessages,
75    fidelity: RestoreFidelity,
76) -> Result<Vec<RestoredFile>, AdapterError> {
77    // Native replays verbatim `options.source.raw_record` rows (session_meta,
78    // then one per message); `codex_session_meta` / `codex_response_item` below
79    // are foreign-only. Replay echoes a frozen snapshot - safe only while
80    // canonical is append-only (spec.md#adapter-integrity-additive-sync).
81    let mut records = Vec::new();
82    if fidelity == RestoreFidelity::Native
83        && let Some(raw) = raw_record(&session.session.options)
84    {
85        records.push(raw);
86    } else {
87        records.push(codex_session_meta(session));
88    }
89    let mut messages = session.messages.clone();
90    messages.sort_by(by_timestamp_then_id);
91    for message in &messages {
92        if fidelity == RestoreFidelity::Native
93            && let Some(raw) = raw_record(message.message.options())
94        {
95            records.push(raw);
96            continue;
97        }
98        // Foreign restore: a System message (a rule-3 carrier, or a source's
99        // own system/developer turn) has no idiomatic home in another
100        // client's transcript - drop it; the content stays in canonical
101        // (spec.md#adapter-native-restore-lossless, foreign clause).
102        if matches!(message.message, Message::System { .. }) {
103            continue;
104        }
105        records.push(codex_response_item(message));
106    }
107    Ok(vec![RestoredFile::new(
108        codex_relative_path(session),
109        jsonl_bytes(NAME, &records)?,
110        fidelity,
111    )])
112}
113
114fn codex_relative_path(session: &crate::sessions::SessionWithMessages) -> PathBuf {
115    let ts = session.session.created_at;
116    let filename_ts = ts.format("%Y-%m-%dT%H-%M-%S");
117    PathBuf::from("sessions")
118        .join(format!("{:04}", ts.year()))
119        .join(format!("{:02}", ts.month()))
120        .join(format!("{:02}", ts.day()))
121        .join(format!(
122            "rollout-{filename_ts}-{}.jsonl",
123            session.session.id
124        ))
125}
126
127fn codex_session_meta(session: &crate::sessions::SessionWithMessages) -> Value {
128    json!({
129        "timestamp": session.session.created_at.to_rfc3339_opts(SecondsFormat::Millis, true),
130        "type": "session_meta",
131        "payload": {
132            "id": session.session.id,
133            "timestamp": session.session.created_at.to_rfc3339_opts(SecondsFormat::Millis, true),
134            "cwd": &*session.session.project,
135        }
136    })
137}
138
139fn codex_response_item(message: &crate::sessions::MessageWithParts) -> Value {
140    json!({
141        "timestamp": message.message.timestamp().to_rfc3339_opts(SecondsFormat::Millis, true),
142        "type": "response_item",
143        "payload": codex_payload(message),
144    })
145}
146
147fn codex_payload(message: &crate::sessions::MessageWithParts) -> Value {
148    if let Some(part) = message.parts.first() {
149        match &part.kind {
150            PartKind::ToolCall {
151                call_id,
152                name,
153                params,
154                ..
155            } if matches!(message.message, Message::Assistant { .. }) => {
156                return json!({
157                    "type": "function_call",
158                    "call_id": extracted_text(call_id),
159                    "name": extracted_text(name),
160                    "arguments": compact_json(params),
161                });
162            }
163            PartKind::ToolResult {
164                call_id, result, ..
165            } if matches!(message.message, Message::Tool { .. }) => {
166                return json!({
167                    "type": "function_call_output",
168                    "call_id": extracted_text(call_id),
169                    "output": result,
170                });
171            }
172            PartKind::Reasoning { text }
173                if matches!(message.message, Message::Assistant { .. }) =>
174            {
175                if let Some(text) = text
176                    && let Ok(value) = serde_json::from_str::<Value>(text.as_ref())
177                {
178                    return value;
179                }
180                return json!({
181                    "type": "reasoning",
182                    "summary": [{"type": "summary_text", "text": extracted_text(text)}],
183                });
184            }
185            _ => {}
186        }
187    }
188    let is_assistant = matches!(message.message, Message::Assistant { .. });
189    json!({
190        "type": "message",
191        "role": match message.message.role() {
192            crate::wire::Role::System => "developer",
193            crate::wire::Role::User => "user",
194            crate::wire::Role::Assistant => "assistant",
195            crate::wire::Role::Tool => "tool",
196        },
197        "content": message
198            .parts
199            .iter()
200            .map(|part| codex_content_part(part, is_assistant))
201            .collect::<Vec<_>>(),
202    })
203}
204
205fn codex_content_part(part: &Part, is_assistant: bool) -> Value {
206    // Codex tags an assistant turn's content `output_text` and a user or
207    // developer turn's content `input_text` - the discriminator is the
208    // owning message's role, not the part.
209    let text_type = if is_assistant {
210        "output_text"
211    } else {
212        "input_text"
213    };
214    match &part.kind {
215        PartKind::Text { text } => json!({
216            "type": text_type,
217            "text": extracted_text(text),
218        }),
219        PartKind::File { data, .. } => json!({
220            "type": text_type,
221            "text": match data {
222                crate::wire::FileData::String(value) => value.clone(),
223                crate::wire::FileData::Bytes(value) => format!("<{} bytes>", value.len()),
224                crate::wire::FileData::Url(value) => value.clone(),
225            },
226        }),
227        other => json!({
228            "type": text_type,
229            "text": compact_json(&serde_json::to_value(other).unwrap_or(Value::Null)),
230        }),
231    }
232}
233
234#[derive(Debug, Clone)]
235pub struct CodexCliAdapter {
236    root: PathBuf,
237}
238
239impl CodexCliAdapter {
240    pub fn new(root: impl Into<PathBuf>) -> Self {
241        Self { root: root.into() }
242    }
243}
244
245impl Adapter for CodexCliAdapter {
246    fn discover(&self) -> DiscoverFuture<'_> {
247        jsonl_tree_discover(self)
248    }
249
250    fn events_with<'a>(&'a self, oracle: &'a dyn SkipOracle) -> AdapterYieldStream<'a> {
251        jsonl_tree_events(self, oracle)
252    }
253
254    fn plan<'a>(&'a self, oracle: &'a dyn SkipOracle) -> crate::adapter::PlanFuture<'a> {
255        crate::adapter::jsonl::jsonl_tree_plan(self, oracle)
256    }
257}
258
259impl JsonlTree for CodexCliAdapter {
260    type State = HashMap<String, Extracted<String>>;
261
262    fn name(&self) -> &'static str {
263        NAME
264    }
265
266    fn root(&self) -> &Path {
267        &self.root
268    }
269
270    fn peek_session_id(&self, _path: &Path, first_line: &str) -> Option<String> {
271        let row: Value = serde_json::from_str(first_line).ok()?;
272        if row.get("type").and_then(Value::as_str) == Some("session_meta") {
273            row.get("payload")?
274                .get("id")?
275                .as_str()
276                .map(ToOwned::to_owned)
277        } else if is_legacy_session_row(&row) {
278            row.get("id")?.as_str().map(ToOwned::to_owned)
279        } else {
280            None
281        }
282    }
283
284    fn peek_watermark(&self, path: &Path) -> crate::adapter::SourceWatermark {
285        // Codex message ids are physical line numbers, not tail-recoverable on a
286        // multi-GB rollout, so freshness keys on the watermark timestamp instead.
287        // The session's max stored timestamp is its last `response_item`'s. Scan
288        // a bounded tail backward for it - never the whole file. The nested
289        // `Option` keeps the walk stopping at the newest response_item even
290        // when its timestamp fails to parse, exactly like the pre-seam walk.
291        match peek_last_mapped(path, |line| {
292            let row: Value = serde_json::from_str(line).ok()?;
293            (row.get("type").and_then(Value::as_str) == Some("response_item")).then(|| {
294                let text = row.get("timestamp").and_then(Value::as_str)?;
295                DateTime::parse_from_rfc3339(text)
296                    .ok()
297                    .map(|ts| ts.with_timezone(&Utc).timestamp_micros())
298            })
299        }) {
300            Some(Some(ts)) => crate::adapter::SourceWatermark::At(ts),
301            // The newest response_item exists but its timestamp is unreadable:
302            // re-read, exactly as before.
303            Some(None) => crate::adapter::SourceWatermark::Opaque,
304            // No response_item in the scan (a legacy rollout, or a session with
305            // no completed turn yet). Every message such a rollout stores
306            // carries a timestamp at or after the session-start header (legacy
307            // payload rows and noise carriers default to it; carriers with own
308            // timestamps are later), so the header is a valid source watermark:
309            // the gate skips only when the store already holds this session at
310            // or past session start - i.e. it was ingested. A never-ingested
311            // rollout has no stored watermark and always re-reads.
312            None => match session_start_ts(path) {
313                Some(ts) => crate::adapter::SourceWatermark::At(ts),
314                None => crate::adapter::SourceWatermark::Opaque,
315            },
316        }
317    }
318
319    fn session(&self, path: &Path, rows: &[BoundedRow]) -> Result<Session, AdapterError> {
320        session_from_rows(path, rows)
321    }
322
323    fn events_from_row(
324        &self,
325        session: &Session,
326        row: &BoundedRow,
327        state: &mut Self::State,
328    ) -> Result<Vec<IngestEvent>, String> {
329        capture_tool_call_name(&row.value, state);
330        events_from_row(&session.id, row.line, &row.value, session.created_at, state)
331    }
332}
333
334/// True for a pre-Oct-2025 legacy rollout's bare first row: session metadata
335/// (`id`/`timestamp`/`git`/`instructions`) at the top level with no `type`
336/// envelope. spec.md#adapters: legacy rollouts predate the `session_meta`
337/// wrapper, so the first row IS the payload.
338fn is_legacy_session_row(row: &Value) -> bool {
339    row.get("type").is_none() && row.get("id").is_some()
340}
341
342/// Session-start timestamp (micros) from the rollout's first line, mirroring
343/// `session_from_rows`' anchor: the `session_meta` payload timestamp (envelope
344/// timestamp as fallback) or a legacy bare header's own. `None` for anything
345/// that is not a recognizable session header.
346fn session_start_ts(path: &Path) -> Option<i64> {
347    let line = peek_first_line(path)?;
348    let row: Value = serde_json::from_str(&line).ok()?;
349    let payload = if row.get("type").and_then(Value::as_str) == Some("session_meta") {
350        row.get("payload").unwrap_or(&Value::Null)
351    } else if is_legacy_session_row(&row) {
352        &row
353    } else {
354        return None;
355    };
356    let text = payload
357        .get("timestamp")
358        .and_then(Value::as_str)
359        .or_else(|| row.get("timestamp").and_then(Value::as_str))?;
360    Some(
361        DateTime::parse_from_rfc3339(text)
362            .ok()?
363            .with_timezone(&Utc)
364            .timestamp_micros(),
365    )
366}
367
368fn session_from_rows(path: &Path, rows: &[BoundedRow]) -> Result<Session, AdapterError> {
369    let path_display = path.display().to_string();
370    let first = rows
371        .first()
372        .ok_or_else(|| AdapterError::schema(NAME, path_display.clone(), "empty jsonl session"))?;
373    let row = &first.value;
374    let at_first = format!("{path_display}:{}", first.line);
375    // The current rollout wraps session metadata in a `session_meta` envelope;
376    // a legacy rollout (spec.md#adapters) has none - the first row is a bare
377    // metadata object. Either way, read fields from `payload`.
378    let payload = if row.get("type").and_then(Value::as_str) == Some("session_meta") {
379        row.get("payload").cloned().unwrap_or(Value::Null)
380    } else if is_legacy_session_row(row) {
381        row.clone()
382    } else {
383        return Err(AdapterError::schema(
384            NAME,
385            at_first,
386            "first row must be session_meta",
387        ));
388    };
389    let id = payload
390        .get("id")
391        .and_then(Value::as_str)
392        .ok_or_else(|| {
393            AdapterError::schema(NAME, at_first.clone(), "session_meta missing payload.id")
394        })?
395        .to_owned();
396    let created_at = payload
397        .get("timestamp")
398        .and_then(Value::as_str)
399        .and_then(|text| DateTime::parse_from_rfc3339(text).ok())
400        .map(|dt| dt.with_timezone(&Utc))
401        .or_else(|| {
402            row.get("timestamp")
403                .and_then(Value::as_str)
404                .and_then(|text| DateTime::parse_from_rfc3339(text).ok())
405                .map(|dt| dt.with_timezone(&Utc))
406        })
407        .ok_or_else(|| {
408            AdapterError::schema(NAME, at_first, "session_meta has no parseable timestamp")
409        })?;
410    let project = match extract_str(&payload, "cwd") {
411        Some(value) => value,
412        None => {
413            let path_str = path
414                .file_name()
415                .and_then(|n| n.to_str())
416                .unwrap_or(path_display.as_str())
417                .to_owned();
418            extract_self_str(&Value::String(path_str)).ok_or_else(|| {
419                AdapterError::schema(
420                    NAME,
421                    path_display.clone(),
422                    "internal: Value::String produced None from Source::as_str",
423                )
424            })?
425        }
426    };
427    let mut options = ProviderOptions::new();
428    options.insert(
429        "source".to_owned(),
430        json!({
431            "adapter": "codex-cli",
432            "originator": payload.get("originator"),
433            "cli_version": payload.get("cli_version"),
434            "model_provider": payload.get("model_provider"),
435            "git": payload.get("git"),
436            "base_instructions": payload.get("base_instructions"),
437            "instructions": payload.get("instructions"),
438            "source": payload.get("source"),
439            "raw_record": extract_raw_record(row),
440        }),
441    );
442
443    Ok(Session {
444        id,
445        parent_session_id: None,
446        parent_message_id: None,
447        source_agent: "codex-cli".to_owned(),
448        created_at,
449        project,
450        options,
451    })
452}
453
454/// Map one codex-cli JSONL record into zero-or-more `IngestEvent`s. Records pond
455/// keeps: `response_item` with `payload.type = "message"` (User/Assistant/
456/// System message + text Parts), `function_call` (Assistant + ToolCall),
457/// `function_call_output` (Tool + ToolResult), `reasoning` (Assistant +
458/// Reasoning Part). `session_meta` is consumed up front; `event_msg` and
459/// `turn_context` are transport noise. Legacy rows (spec.md#adapters) carry
460/// the same payload shapes un-enveloped; `{record_type:"state"}` markers and
461/// the bare first row are eventless.
462fn events_from_row(
463    session_id: &str,
464    line: usize,
465    row: &Value,
466    default_timestamp: DateTime<Utc>,
467    tool_call_names: &HashMap<String, Extracted<String>>,
468) -> Result<Vec<IngestEvent>, String> {
469    let kind = row.get("type").and_then(Value::as_str);
470    // Eventless rows: `session_meta` (current) and the legacy bare first row
471    // are both consumed up front by session_meta(); legacy
472    // `{record_type:"state"}` markers are transport noise (spec.md#adapters).
473    if kind == Some("session_meta")
474        || is_legacy_session_row(row)
475        || (kind.is_none() && row.get("record_type").is_some())
476    {
477        return Ok(Vec::new());
478    }
479    // Normalize to the per-turn payload. A current row wraps it in a
480    // `response_item` envelope carrying its own timestamp; a legacy data row
481    // IS the payload (spec.md#adapters) and inherits the session timestamp.
482    let (payload, timestamp) = if kind == Some("response_item") {
483        let timestamp = row
484            .get("timestamp")
485            .and_then(Value::as_str)
486            .and_then(|text| DateTime::parse_from_rfc3339(text).ok())
487            .map(|dt| dt.with_timezone(&Utc))
488            .unwrap_or(default_timestamp);
489        (row.get("payload").unwrap_or(&Value::Null), timestamp)
490    } else {
491        (row, default_timestamp)
492    };
493    let payload_type = payload.get("type").and_then(Value::as_str).unwrap_or("");
494    let message_id = format!("{session_id}:{line:06}");
495
496    match payload_type {
497        "message" => message_events(session_id, &message_id, timestamp, payload, row),
498        "function_call" => Ok(tool_call_events(
499            session_id,
500            &message_id,
501            timestamp,
502            payload,
503            row,
504        )),
505        "function_call_output" => Ok(tool_result_events(
506            session_id,
507            &message_id,
508            timestamp,
509            payload,
510            row,
511            tool_call_names,
512        )),
513        "reasoning" => Ok(reasoning_events(
514            session_id,
515            &message_id,
516            timestamp,
517            payload,
518            row,
519        )),
520        "custom_tool_call" => Ok(custom_tool_call_events(
521            session_id,
522            &message_id,
523            timestamp,
524            payload,
525            row,
526        )),
527        "custom_tool_call_output" => Ok(custom_tool_result_events(
528            session_id,
529            &message_id,
530            timestamp,
531            payload,
532            row,
533        )),
534        _ => Ok(vec![raw_carrier_event(session_id, line, row, timestamp)]),
535    }
536}
537
538fn row_options(row: &Value) -> ProviderOptions {
539    let mut options = ProviderOptions::new();
540    options.insert(
541        "source".to_owned(),
542        json!({ "raw_record": extract_raw_record(row) }),
543    );
544    options
545}
546
547fn raw_carrier_event(
548    session_id: &str,
549    line: usize,
550    row: &Value,
551    timestamp: DateTime<Utc>,
552) -> IngestEvent {
553    IngestEvent::Message(Message::System {
554        id: row
555            .get("id")
556            .and_then(Value::as_str)
557            .map_or_else(|| format!("{session_id}:{line:06}:raw"), ToOwned::to_owned),
558        session_id: session_id.to_owned(),
559        timestamp: row
560            .get("timestamp")
561            .and_then(Value::as_str)
562            .and_then(|text| DateTime::parse_from_rfc3339(text).ok())
563            .map(|dt| dt.with_timezone(&Utc))
564            .unwrap_or(timestamp),
565        content: None,
566        options: row_options(row),
567    })
568}
569
570/// Stash one row's `function_call` (call_id -> name) into the per-file
571/// map so the matching `function_call_output` row downstream can resolve
572/// the tool name rather than fall back to a sentinel.
573fn capture_tool_call_name(row: &Value, map: &mut HashMap<String, Extracted<String>>) {
574    // Mirror events_from_row's payload normalization: a current row wraps the
575    // payload under `response_item`, a legacy row IS the payload.
576    let payload = match row.get("type").and_then(Value::as_str) {
577        Some("response_item") => row.get("payload"),
578        Some(_) => Some(row),
579        None => None,
580    };
581    let Some(payload) = payload else {
582        return;
583    };
584    if payload.get("type").and_then(Value::as_str) != Some("function_call") {
585        return;
586    }
587    let Some(call_id) = payload.get("call_id").and_then(Value::as_str) else {
588        return;
589    };
590    let Some(name) = extract_str(payload, "name") else {
591        return;
592    };
593    map.insert(call_id.to_owned(), name);
594}
595
596fn message_events(
597    session_id: &str,
598    message_id: &str,
599    timestamp: DateTime<Utc>,
600    payload: &Value,
601    row: &Value,
602) -> Result<Vec<IngestEvent>, String> {
603    let role = payload
604        .get("role")
605        .and_then(Value::as_str)
606        .ok_or_else(|| "message missing role".to_owned())?;
607    let Some(content) = payload.get("content").and_then(Value::as_array) else {
608        return Ok(vec![message_raw_carrier_event(
609            session_id, message_id, row, timestamp,
610        )]);
611    };
612    // spec.md#model-part-provenance: a `developer` record is a harness instruction
613    // block; a `user`-slot record whose body is `<environment_context>` or an
614    // `# AGENTS.md instructions` blob is injected context, not a genuine
615    // prompt. Everything else in a message record is conversation.
616    let provenance = message_provenance(role, content);
617    let mut parts = Vec::with_capacity(content.len());
618    for (ordinal, item) in content.iter().enumerate() {
619        // Faithful encoding of one content item: prefer the raw `text`
620        // field when present; otherwise compact-encode the structured
621        // body as a JSON string. The fallback is lossless (preserves the
622        // item bytes) and explicit (not a synthesised "unknown" or "").
623        let text = extract_str(item, "text").or_else(|| Some(extract_compact_repr(item)));
624        parts.push(Part {
625            session_id: session_id.to_owned(),
626            id: part_id(message_id, ordinal),
627            message_id: message_id.to_owned(),
628            ordinal: part_ordinal(ordinal),
629            provenance,
630            options: empty_options(),
631            kind: PartKind::Text { text },
632        });
633    }
634
635    let (message, keep_parts) = match role {
636        "user" => (
637            Message::User {
638                id: message_id.to_owned(),
639                session_id: session_id.to_owned(),
640                timestamp,
641                options: row_options(row),
642            },
643            true,
644        ),
645        "assistant" => (
646            Message::Assistant {
647                id: message_id.to_owned(),
648                session_id: session_id.to_owned(),
649                timestamp,
650                options: row_options(row),
651            },
652            true,
653        ),
654        // `developer` rows are codex-cli's system-prompt frames; map to System
655        // with `content: None` and let the inner Text Parts carry the body.
656        "developer" | "system" => (
657            Message::System {
658                id: message_id.to_owned(),
659                session_id: session_id.to_owned(),
660                timestamp,
661                content: None,
662                options: row_options(row),
663            },
664            true,
665        ),
666        _ => {
667            return Ok(vec![message_raw_carrier_event(
668                session_id, message_id, row, timestamp,
669            )]);
670        }
671    };
672
673    let mut events = Vec::with_capacity(parts.len() + 1);
674    events.push(IngestEvent::Message(message));
675    if keep_parts {
676        events.extend(parts.into_iter().map(IngestEvent::Part));
677    }
678    Ok(events)
679}
680
681fn message_raw_carrier_event(
682    session_id: &str,
683    message_id: &str,
684    row: &Value,
685    timestamp: DateTime<Utc>,
686) -> IngestEvent {
687    IngestEvent::Message(Message::System {
688        id: message_id.to_owned(),
689        session_id: session_id.to_owned(),
690        timestamp,
691        content: row
692            .get("payload")
693            .and_then(|payload| payload.get("role"))
694            .or_else(|| row.get("role"))
695            .and_then(Value::as_str)
696            .and_then(|role| extract_self_str(&Value::String(role.to_owned()))),
697        options: row_options(row),
698    })
699}
700
701/// Provenance of a codex `message` record (spec.md#model-part-provenance). A
702/// `developer` record is a harness instruction block; a `user`-slot record
703/// whose only content is `<environment_context>` or `# AGENTS.md instructions`
704/// is injected context rather than a typed prompt. v1 codex never interleaves
705/// authored and injected content within one record.
706fn message_provenance(role: &str, content: &[Value]) -> Provenance {
707    if role == "developer" || role == "system" {
708        return Provenance::Injected;
709    }
710    if role == "user" {
711        let injected = content.iter().any(|item| {
712            item.get("text")
713                .and_then(Value::as_str)
714                .is_some_and(is_injected_user_text)
715        });
716        if injected {
717            return Provenance::Injected;
718        }
719    }
720    Provenance::Conversational
721}
722
723/// Harness-injected user-slot content codex emits as a non-prompt record.
724fn is_injected_user_text(text: &str) -> bool {
725    let trimmed = text.trim_start();
726    trimmed.starts_with("<environment_context>")
727        || trimmed.starts_with("<user_instructions>")
728        || trimmed.starts_with("# AGENTS.md")
729}
730
731fn tool_call_events(
732    session_id: &str,
733    message_id: &str,
734    timestamp: DateTime<Utc>,
735    payload: &Value,
736    row: &Value,
737) -> Vec<IngestEvent> {
738    let call_id = extract_str(payload, "call_id");
739    let name = extract_str(payload, "name");
740    let params = match payload.get("arguments") {
741        Some(Value::String(text)) => json_or_string(text),
742        Some(other) => other.clone(),
743        None => Value::Null,
744    };
745    let part = Part {
746        session_id: session_id.to_owned(),
747        id: part_id(message_id, 0),
748        message_id: message_id.to_owned(),
749        ordinal: 0,
750        // spec.md#model-part-provenance: the model authored the tool call.
751        provenance: Provenance::Conversational,
752        options: empty_options(),
753        kind: PartKind::ToolCall {
754            call_id,
755            name,
756            params,
757            provider_executed: false,
758        },
759    };
760    vec![
761        IngestEvent::Message(Message::Assistant {
762            id: message_id.to_owned(),
763            session_id: session_id.to_owned(),
764            timestamp,
765            options: row_options(row),
766        }),
767        IngestEvent::Part(part),
768    ]
769}
770
771fn custom_tool_call_events(
772    session_id: &str,
773    message_id: &str,
774    timestamp: DateTime<Utc>,
775    payload: &Value,
776    row: &Value,
777) -> Vec<IngestEvent> {
778    let part = Part {
779        session_id: session_id.to_owned(),
780        id: part_id(message_id, 0),
781        message_id: message_id.to_owned(),
782        ordinal: 0,
783        // spec.md#model-part-provenance: the model authored the tool call.
784        provenance: Provenance::Conversational,
785        options: empty_options(),
786        kind: PartKind::ToolCall {
787            call_id: extract_str(payload, "call_id"),
788            name: extract_str(payload, "name"),
789            params: payload.get("input").cloned().unwrap_or(Value::Null),
790            provider_executed: true,
791        },
792    };
793    vec![
794        IngestEvent::Message(Message::Assistant {
795            id: message_id.to_owned(),
796            session_id: session_id.to_owned(),
797            timestamp,
798            options: row_options(row),
799        }),
800        IngestEvent::Part(part),
801    ]
802}
803
804fn custom_tool_result_events(
805    session_id: &str,
806    message_id: &str,
807    timestamp: DateTime<Utc>,
808    payload: &Value,
809    row: &Value,
810) -> Vec<IngestEvent> {
811    let part = Part {
812        session_id: session_id.to_owned(),
813        id: part_id(message_id, 0),
814        message_id: message_id.to_owned(),
815        ordinal: 0,
816        // spec.md#model-part-provenance: tool output is runtime-produced.
817        provenance: Provenance::Injected,
818        options: empty_options(),
819        kind: PartKind::ToolResult {
820            call_id: extract_str(payload, "call_id"),
821            name: extract_str(payload, "name"),
822            is_failure: false,
823            result: payload.get("output").cloned().unwrap_or(Value::Null),
824        },
825    };
826    vec![
827        IngestEvent::Message(Message::Tool {
828            id: message_id.to_owned(),
829            session_id: session_id.to_owned(),
830            timestamp,
831            options: row_options(row),
832        }),
833        IngestEvent::Part(part),
834    ]
835}
836
837fn tool_result_events(
838    session_id: &str,
839    message_id: &str,
840    timestamp: DateTime<Utc>,
841    payload: &Value,
842    row: &Value,
843    tool_call_names: &HashMap<String, Extracted<String>>,
844) -> Vec<IngestEvent> {
845    let call_id = extract_str(payload, "call_id");
846    // Resolve tool name from the earlier `function_call` row via the
847    // per-file `call_id -> name` map. Misses (e.g. compaction pruned the
848    // originating call) yield `None`, a faithful "unresolved" value.
849    let name = call_id
850        .as_ref()
851        .and_then(|id| tool_call_names.get(id.as_str()))
852        .cloned();
853    let result = payload.get("output").cloned().unwrap_or(Value::Null);
854    let part = Part {
855        session_id: session_id.to_owned(),
856        id: part_id(message_id, 0),
857        message_id: message_id.to_owned(),
858        ordinal: 0,
859        // spec.md#model-part-provenance: tool output is runtime-produced.
860        provenance: Provenance::Injected,
861        options: empty_options(),
862        kind: PartKind::ToolResult {
863            call_id,
864            name,
865            is_failure: false,
866            result,
867        },
868    };
869    vec![
870        IngestEvent::Message(Message::Tool {
871            id: message_id.to_owned(),
872            session_id: session_id.to_owned(),
873            timestamp,
874            options: row_options(row),
875        }),
876        IngestEvent::Part(part),
877    ]
878}
879
880fn reasoning_events(
881    session_id: &str,
882    message_id: &str,
883    timestamp: DateTime<Utc>,
884    payload: &Value,
885    row: &Value,
886) -> Vec<IngestEvent> {
887    // The source `summary` array is the only place reasoning text lives in
888    // codex-cli's format. Empty array (or missing field) -> `None`. Joined
889    // text -> `Some(...)`. Don't synthesize an empty string.
890    let summary = payload
891        .get("summary")
892        .and_then(Value::as_array)
893        .and_then(|items| {
894            let joined = items
895                .iter()
896                .filter_map(|item| extract_str(item, "text"))
897                .map(|e| (*e).clone())
898                .collect::<Vec<_>>()
899                .join("\n");
900            if joined.is_empty() {
901                None
902            } else {
903                Some(extract_compact_repr(payload))
904            }
905        });
906    let part = Part {
907        session_id: session_id.to_owned(),
908        id: part_id(message_id, 0),
909        message_id: message_id.to_owned(),
910        ordinal: 0,
911        // spec.md#model-part-provenance: model-authored reasoning.
912        provenance: Provenance::Conversational,
913        options: empty_options(),
914        kind: PartKind::Reasoning { text: summary },
915    };
916    vec![
917        IngestEvent::Message(Message::Assistant {
918            id: message_id.to_owned(),
919            session_id: session_id.to_owned(),
920            timestamp,
921            options: row_options(row),
922        }),
923        IngestEvent::Part(part),
924    ]
925}
926
927#[cfg(test)]
928mod tests {
929    //! End-to-end test for the codex-cli adapter: ingest the committed fixture
930    //! corpus and assert pond's canonical Session/Message/Part shape comes out
931    //! the other side. The fixture lives under
932    //! `tests/fixtures/adapter/codex_cli/`.
933    #![allow(clippy::expect_used, clippy::unwrap_used)]
934
935    use super::*;
936    use crate::{handlers::ingest_adapter, sessions::Store, wire::PartKind};
937    use tempfile::TempDir;
938
939    // Manifest-dir anchored: unit tests must not depend on the process cwd
940    // (figment::Jail chdirs the whole test process while config tests run).
941    const FIXTURES: &str = concat!(
942        env!("CARGO_MANIFEST_DIR"),
943        "/tests/fixtures/adapter/codex_cli/sessions"
944    );
945
946    #[test]
947    fn probe_default_finds_codex_sessions_under_home() -> anyhow::Result<()> {
948        crate::adapter::test_support::assert_probe_default(
949            &CodexCliFactory,
950            &[".codex", "sessions"],
951        )
952    }
953
954    /// `peek_watermark` is the freshness watermark for multi-GB rollouts where
955    /// the line-numbered message id is not tail-recoverable. It must read the
956    /// last `response_item`'s envelope timestamp - pond's stored max - and
957    /// ignore the trailing `event_msg` noise (whose stored timestamp is the
958    /// session default), scanning only the file tail.
959    #[test]
960    fn peek_watermark_targets_last_response_item_ignoring_trailing_event_msg() {
961        let dir = TempDir::new().unwrap();
962        let path = dir.path().join("rollout.jsonl");
963        let lines = [
964            r#"{"type":"session_meta","timestamp":"2026-03-20T03:00:00.000Z","payload":{"id":"sess-x","timestamp":"2026-03-20T03:00:00.000Z"}}"#,
965            r#"{"type":"response_item","timestamp":"2026-03-20T03:10:00.000Z","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"hi"}]}}"#,
966            r#"{"type":"response_item","timestamp":"2026-03-20T03:20:30.500Z","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"yo"}]}}"#,
967            // Trailing token-count noise, later wall-clock but stored at the
968            // session default - must NOT be picked as the watermark.
969            r#"{"type":"event_msg","timestamp":"2026-03-20T03:59:59.000Z","payload":{"type":"token_count","info":{}}}"#,
970        ];
971        std::fs::write(&path, lines.join("\n") + "\n").unwrap();
972
973        let adapter = CodexCliAdapter::new(dir.path());
974        let expected = DateTime::parse_from_rfc3339("2026-03-20T03:20:30.500Z")
975            .unwrap()
976            .with_timezone(&Utc)
977            .timestamp_micros();
978        assert_eq!(
979            adapter.peek_watermark(&path),
980            crate::adapter::SourceWatermark::At(expected)
981        );
982    }
983
984    /// A rollout with no `response_item` at all (legacy data rows only, or a
985    /// session with no completed turn) keys freshness on the session-start
986    /// header: pond stores all such rows at or after that timestamp, so a
987    /// stored watermark at/past it proves the file was ingested. A
988    /// never-ingested file has no stored watermark and still re-reads.
989    #[test]
990    fn peek_watermark_falls_back_to_session_start_without_response_items() {
991        let dir = TempDir::new().unwrap();
992        let path = dir.path().join("rollout-legacy.jsonl");
993        let lines = [
994            r#"{"id":"sess-legacy","timestamp":"2025-09-10T12:48:29.371Z","git":{},"instructions":null}"#,
995            r#"{"record_type":"state"}"#,
996            r#"{"type":"message","role":"user","content":[{"type":"input_text","text":"hi"}]}"#,
997        ];
998        std::fs::write(&path, lines.join("\n") + "\n").unwrap();
999
1000        let adapter = CodexCliAdapter::new(dir.path());
1001        let expected = DateTime::parse_from_rfc3339("2025-09-10T12:48:29.371Z")
1002            .unwrap()
1003            .with_timezone(&Utc)
1004            .timestamp_micros();
1005        assert_eq!(
1006            adapter.peek_watermark(&path),
1007            crate::adapter::SourceWatermark::At(expected)
1008        );
1009    }
1010
1011    #[tokio::test(flavor = "multi_thread")]
1012    async fn native_restore_is_value_equal_to_fixture_corpus() -> anyhow::Result<()> {
1013        let adapter = CodexCliAdapter::new(FIXTURES);
1014        crate::adapter::test_support::assert_native_restore(
1015            &CodexCliFactory,
1016            &adapter,
1017            // Codex rollout paths embed the `sessions/` segment, so the corpus
1018            // root is FIXTURES' parent, not FIXTURES itself.
1019            std::path::Path::new(FIXTURES)
1020                .parent()
1021                .expect("FIXTURES is nested under a corpus root"),
1022        )
1023        .await
1024    }
1025
1026    #[tokio::test(flavor = "multi_thread")]
1027    async fn codex_cli_adapter_ingests_fixture_corpus_into_canonical_shape() -> anyhow::Result<()> {
1028        let temp = TempDir::new()?;
1029        let store = Store::open_local(temp.path()).await?;
1030        let adapter = CodexCliAdapter::new(FIXTURES);
1031
1032        let summary = ingest_adapter(&store, &adapter, &crate::adapter::NoopOracle, |_| {}).await?;
1033        assert!(summary.accepted() > 0, "ingest must accept rows");
1034        assert_eq!(summary.dropped_events, 0, "no per-event drops expected");
1035        assert_eq!(
1036            summary.dropped_sessions, 0,
1037            "no session-level rejections expected"
1038        );
1039        assert_eq!(summary.skipped_files, 0, "no whole-file skips expected");
1040
1041        let (sessions, messages, parts) = store.row_counts().await?;
1042        assert!(sessions > 0, "at least one codex-cli session");
1043        assert!(messages > 0, "at least one codex-cli message");
1044        assert!(parts > 0, "at least one codex-cli Part");
1045
1046        let mut saw_text_part = false;
1047        for session_id in store.session_ids().await? {
1048            let session = store
1049                .get_session(&session_id)
1050                .await?
1051                .expect("session round-trips");
1052            assert_eq!(session.session.source_agent, "codex-cli");
1053            assert!(
1054                !session.messages.is_empty(),
1055                "session {session_id} must carry messages",
1056            );
1057            for stored in &session.messages {
1058                for part in &stored.parts {
1059                    if matches!(part.kind, PartKind::Text { .. }) {
1060                        saw_text_part = true;
1061                    }
1062                }
1063            }
1064        }
1065        assert!(
1066            saw_text_part,
1067            "codex-cli corpus must contain at least one Text Part",
1068        );
1069        Ok(())
1070    }
1071
1072    /// spec.md#model-part-provenance: a `developer` record and a `user`-slot record
1073    /// whose body is `<environment_context>` are harness-injected; a genuine
1074    /// user prompt and an assistant message are conversation.
1075    #[test]
1076    fn message_provenance_separates_prompts_from_harness_records() {
1077        let prompt = vec![json!({"type": "input_text", "text": "refactor this"})];
1078        assert_eq!(
1079            message_provenance("user", &prompt),
1080            Provenance::Conversational,
1081        );
1082        assert_eq!(
1083            message_provenance("assistant", &[]),
1084            Provenance::Conversational,
1085        );
1086
1087        let developer = vec![json!({"type": "input_text", "text": "you are an agent"})];
1088        assert_eq!(
1089            message_provenance("developer", &developer),
1090            Provenance::Injected,
1091        );
1092
1093        let env = vec![json!({
1094            "type": "input_text",
1095            "text": "<environment_context>cwd=/tmp</environment_context>",
1096        })];
1097        assert_eq!(message_provenance("user", &env), Provenance::Injected);
1098    }
1099
1100    #[test]
1101    fn legacy_rows_normalize_to_payloads() {
1102        let ts = Utc::now();
1103        let map: HashMap<String, Extracted<String>> = HashMap::new();
1104
1105        // The bare first row and `{record_type:"state"}` markers are eventless.
1106        let first = json!({"id": "s1", "timestamp": "2025-09-13T04:30:17.447Z"});
1107        let state = json!({"record_type": "state"});
1108        assert!(
1109            events_from_row("s1", 1, &first, ts, &map)
1110                .expect("legacy first row parses")
1111                .is_empty(),
1112        );
1113        assert!(
1114            events_from_row("s1", 2, &state, ts, &map)
1115                .expect("state marker parses")
1116                .is_empty(),
1117        );
1118
1119        // An un-enveloped legacy `message` row yields a Message + Text Part.
1120        let message = json!({
1121            "type": "message",
1122            "role": "user",
1123            "content": [{"type": "input_text", "text": "hi"}],
1124        });
1125        let events = events_from_row("s1", 3, &message, ts, &map).expect("legacy message parses");
1126        assert_eq!(events.len(), 2, "message + one Text Part");
1127        assert!(matches!(
1128            events[0],
1129            IngestEvent::Message(Message::User { .. })
1130        ));
1131        assert!(matches!(
1132            &events[1],
1133            IngestEvent::Part(part) if matches!(part.kind, PartKind::Text { .. }),
1134        ));
1135    }
1136
1137    #[test]
1138    fn unknown_message_role_becomes_lossless_carrier() {
1139        let ts = Utc::now();
1140        let map: HashMap<String, Extracted<String>> = HashMap::new();
1141        let row = json!({
1142            "type": "response_item",
1143            "timestamp": "2026-06-01T00:00:00Z",
1144            "payload": {
1145                "type": "message",
1146                "role": "future_role",
1147                "content": [{"type": "input_text", "text": "keep me"}],
1148            },
1149        });
1150
1151        let events = events_from_row("s1", 4, &row, ts, &map).expect("carrier is valid");
1152        assert_eq!(events.len(), 1);
1153        assert!(matches!(
1154            &events[0],
1155            IngestEvent::Message(Message::System { id, content, .. })
1156                if id == "s1:000004" && content.as_deref().map(String::as_str) == Some("future_role")
1157        ));
1158    }
1159
1160    #[tokio::test(flavor = "multi_thread")]
1161    async fn legacy_rollout_ingests_into_canonical_shape() -> anyhow::Result<()> {
1162        let temp = TempDir::new()?;
1163        let store = Store::open_local(temp.path()).await?;
1164        let adapter = CodexCliAdapter::new(FIXTURES);
1165        ingest_adapter(&store, &adapter, &crate::adapter::NoopOracle, |_| {}).await?;
1166
1167        // The legacy fixture's bare first row -> Session: id and timestamp
1168        // read from the top level, with no `session_meta` envelope.
1169        let session = store
1170            .get_session("67c52f3f-d25e-4194-a006-93de58f28d7c")
1171            .await?
1172            .expect("legacy rollout ingests as a session");
1173        assert_eq!(session.session.source_agent, "codex-cli");
1174        assert_eq!(
1175            session
1176                .session
1177                .created_at
1178                .to_rfc3339_opts(SecondsFormat::Millis, true),
1179            "2025-09-13T04:30:17.447Z",
1180        );
1181        // Eleven un-enveloped data rows -> eleven messages.
1182        assert_eq!(session.messages.len(), 11, "every legacy data row ingests");
1183        // The legacy `function_call_output` resolves its tool name from the
1184        // prior legacy `function_call` row via the per-file call_id map.
1185        let resolved = session.messages.iter().any(|message| {
1186            message
1187                .parts
1188                .iter()
1189                .any(|part| matches!(&part.kind, PartKind::ToolResult { name: Some(_), .. }))
1190        });
1191        assert!(
1192            resolved,
1193            "legacy function_call_output resolves its tool name"
1194        );
1195        Ok(())
1196    }
1197}