1use 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
46pub 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 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 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 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 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 Some(None) => crate::adapter::SourceWatermark::Opaque,
304 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
334fn is_legacy_session_row(row: &Value) -> bool {
339 row.get("type").is_none() && row.get("id").is_some()
340}
341
342fn 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 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
454fn 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 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 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
570fn capture_tool_call_name(row: &Value, map: &mut HashMap<String, Extracted<String>>) {
574 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 let provenance = message_provenance(role, content);
617 let mut parts = Vec::with_capacity(content.len());
618 for (ordinal, item) in content.iter().enumerate() {
619 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" | "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
701fn 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
723fn 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 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 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 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 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 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 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 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 #![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 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 #[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 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 #[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 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 #[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 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 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 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 assert_eq!(session.messages.len(), 11, "every legacy data row ingests");
1183 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}