Skip to main content

agent_top_core/harness/
codex.rs

1//! OpenAI Codex CLI: `~/.codex/sessions/YYYY/MM/DD/rollout-<ts>-<id>.jsonl`.
2//!
3//! Format notes (verified on Codex CLI 0.149, 2026-09-03):
4//! * The first line is `session_meta` with `payload.cwd`, `payload.id`,
5//!   `payload.cli_version` and `payload.originator`.
6//! * `event_msg` / `token_count` carries `info.total_token_usage`, which is
7//!   cumulative for the session; `info` is null on rate-limit-only events.
8//!   `input_tokens` includes `cached_input_tokens`.
9//! * `task_started` / `task_complete` / `turn_aborted` bracket a turn.
10//! * `response_item` with `payload.type` `function_call` or
11//!   `custom_tool_call` is one tool call; the matching `*_output` item
12//!   carries the same `payload.call_id`, and the two lines' timestamps
13//!   bracket the call. That pairing is the trace.
14//!
15//! Codex model prices are not in the static table, so cost is reported as
16//! unpriced tokens.
17
18use super::{REFRESH_BUDGET_BYTES, SessionSummary, SessionTracker, SpanRetention, parse_rfc3339_utc};
19use crate::jsonl::TailReader;
20use crate::model::{Activity, Harness, TokenUsage};
21use crate::pricing::{self, Table};
22use serde_json::Value;
23use std::path::{Path, PathBuf};
24use std::time::SystemTime;
25
26pub fn codex_dir() -> Option<PathBuf> {
27    if let Some(d) = std::env::var_os("CODEX_HOME") {
28        return Some(PathBuf::from(d));
29    }
30    std::env::var_os("HOME").map(|h| PathBuf::from(h).join(".codex"))
31}
32
33pub fn sessions_dir() -> Option<PathBuf> {
34    codex_dir().map(|d| d.join("sessions"))
35}
36
37/// Rollout files modified after `since`. Walks `YYYY/MM/DD` and prunes by
38/// directory mtime so the walk stays cheap on a long history.
39pub fn recent_rollouts(since: SystemTime) -> Vec<PathBuf> {
40    let Some(root) = sessions_dir() else { return Vec::new() };
41    let mut out = Vec::new();
42    walk(&root, 0, since, &mut out);
43    out
44}
45
46fn walk(dir: &Path, depth: usize, since: SystemTime, out: &mut Vec<PathBuf>) {
47    let Ok(rd) = std::fs::read_dir(dir) else { return };
48    for e in rd.flatten() {
49        let p = e.path();
50        let Ok(md) = e.metadata() else { continue };
51        if md.is_dir() {
52            if depth < 3 && md.modified().map(|m| m >= since).unwrap_or(true) {
53                walk(&p, depth + 1, since, out);
54            }
55        } else if p.extension().and_then(|x| x.to_str()) == Some("jsonl") && md.modified().map(|m| m >= since).unwrap_or(false) {
56            out.push(p);
57        }
58    }
59}
60
61/// Cheap header read: cwd and start time from the first line only.
62pub fn read_meta(path: &Path) -> Option<(PathBuf, SystemTime)> {
63    use std::io::{BufRead, BufReader};
64    let f = std::fs::File::open(path).ok()?;
65    let mut first = String::new();
66    BufReader::new(f).read_line(&mut first).ok()?;
67    let v: Value = serde_json::from_str(&first).ok()?;
68    if v.get("type").and_then(Value::as_str) != Some("session_meta") {
69        return None;
70    }
71    let cwd = v.pointer("/payload/cwd").and_then(Value::as_str).map(PathBuf::from)?;
72    let ts = v.get("timestamp").and_then(Value::as_str).and_then(parse_rfc3339_utc)?;
73    Some((cwd, ts))
74}
75
76pub struct CodexTranscript {
77    reader: TailReader,
78    prices: &'static Table,
79    summary: SessionSummary,
80}
81
82impl CodexTranscript {
83    pub fn new(path: impl Into<PathBuf>) -> Self {
84        CodexTranscript {
85            reader: TailReader::new(path),
86            prices: pricing::table(),
87            summary: SessionSummary { harness: Some(Harness::Codex), ..Default::default() },
88        }
89    }
90
91    /// See `ClaudeTranscript::with_prices`.
92    pub fn with_prices(mut self, prices: &'static Table) -> Self {
93        self.prices = prices;
94        self
95    }
96
97    /// Keep every span instead of the newest `MAX_SPANS`. See `SpanRetention`.
98    pub fn with_spans(mut self, retention: SpanRetention) -> Self {
99        self.summary.spans = retention.log();
100        self
101    }
102
103    fn ingest(&mut self, line: &str) {
104        let Ok(v) = serde_json::from_str::<Value>(line) else { return };
105        let ts = v.get("timestamp").and_then(Value::as_str).and_then(parse_rfc3339_utc);
106        if let Some(ts) = ts {
107            if self.summary.started_at.is_none() {
108                self.summary.started_at = Some(ts);
109            }
110            self.summary.last_activity = Some(ts);
111        }
112        let kind = v.get("type").and_then(Value::as_str).unwrap_or("");
113        let payload = v.get("payload");
114        let ptype = payload.and_then(|p| p.get("type")).and_then(Value::as_str).unwrap_or("");
115        match kind {
116            "session_meta" => {
117                if let Some(p) = payload {
118                    self.summary.session_id = p.get("id").or(p.get("session_id")).and_then(Value::as_str).map(str::to_string);
119                    self.summary.cwd = p.get("cwd").and_then(Value::as_str).map(PathBuf::from);
120                    self.summary.harness_version = p.get("cli_version").and_then(Value::as_str).map(str::to_string);
121                }
122            }
123            "turn_context" => {
124                if let Some(m) = payload.and_then(|p| p.get("model")).and_then(Value::as_str) {
125                    self.summary.model = Some(m.to_string());
126                }
127            }
128            "event_msg" => match ptype {
129                "token_count" => {
130                    if let Some(total) = payload.and_then(|p| p.pointer("/info/total_token_usage")) {
131                        let g = |k: &str| total.get(k).and_then(Value::as_u64).unwrap_or(0);
132                        self.summary.health.usage_records += 1;
133                        if g("input_tokens") + g("output_tokens") + g("cached_input_tokens") == 0 {
134                            self.summary.health.empty_usage_records += 1;
135                        }
136                        let cached = g("cached_input_tokens");
137                        let usage = TokenUsage {
138                            input: g("input_tokens").saturating_sub(cached),
139                            cache_read: cached,
140                            output: g("output_tokens"),
141                            ..Default::default()
142                        };
143                        self.summary.usage = usage;
144                        let price = self.summary.model.as_deref().and_then(|m| self.prices.lookup(m));
145                        match price {
146                            Some(p) => {
147                                self.summary.cost_usd = p.cost(&usage);
148                                self.summary.unpriced_tokens = 0;
149                            }
150                            None => {
151                                self.summary.cost_usd = 0.0;
152                                self.summary.unpriced_tokens = usage.total();
153                            }
154                        }
155                    }
156                }
157                "task_started" | "user_message" => self.summary.activity = Activity::Working,
158                "task_complete" | "turn_aborted" | "error" => self.summary.activity = Activity::Waiting,
159                _ => {}
160            },
161            "response_item" => match ptype {
162                "function_call" | "custom_tool_call" | "local_shell_call" => {
163                    self.summary.tool_calls += 1;
164                    if let (Some(ts), Some(p)) = (ts, payload) {
165                        let id = call_id(p);
166                        let name = p.get("name").and_then(Value::as_str).unwrap_or(ptype);
167                        self.summary.spans.open(id, name.to_string(), ts, false);
168                    }
169                }
170                "function_call_output" | "custom_tool_call_output" | "local_shell_call_output" => {
171                    if let (Some(ts), Some(p)) = (ts, payload) {
172                        // Codex reports the result as an opaque string, and
173                        // agent-top does not read tool output, so a failed call
174                        // is not distinguishable from a successful one here.
175                        self.summary.spans.close(&call_id(p), ts, false);
176                    }
177                }
178                "message" if payload.and_then(|p| p.get("role")).and_then(Value::as_str) == Some("assistant") => {
179                    self.summary.turns += 1;
180                    self.summary.health.billable_messages += 1;
181                }
182                _ => {}
183            },
184            _ => {}
185        }
186    }
187}
188
189/// `call_id` on function calls, `id` on the shell-call variants.
190fn call_id(payload: &Value) -> String {
191    payload.get("call_id").or_else(|| payload.get("id")).and_then(Value::as_str).unwrap_or_default().to_string()
192}
193
194impl SessionTracker for CodexTranscript {
195    fn refresh(&mut self) -> anyhow::Result<bool> {
196        let (lines, more) = self.reader.read_new_lines(REFRESH_BUDGET_BYTES)?;
197        for l in &lines {
198            self.ingest(l);
199        }
200        Ok(more)
201    }
202
203    fn summary(&self) -> &SessionSummary {
204        &self.summary
205    }
206
207    fn path(&self) -> &Path {
208        self.reader.path()
209    }
210}
211
212#[cfg(test)]
213mod tests {
214    use super::*;
215    use std::io::Write;
216
217    #[test]
218    fn reads_cumulative_usage_and_state() {
219        let dir = std::env::temp_dir().join(format!("agent-top-codex-{}", std::process::id()));
220        std::fs::create_dir_all(&dir).unwrap();
221        let path = dir.join("rollout.jsonl");
222        let mut f = std::fs::File::create(&path).unwrap();
223        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:20.787Z","type":"session_meta","payload":{{"id":"01a0","cwd":"/tmp/p","cli_version":"0.149.1"}}}}"#).unwrap();
224        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:21.000Z","type":"turn_context","payload":{{"model":"gpt-5-codex"}}}}"#).unwrap();
225        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:22.000Z","type":"event_msg","payload":{{"type":"task_started"}}}}"#).unwrap();
226        writeln!(
227            f,
228            r#"{{"timestamp":"2026-08-28T08:53:23.000Z","type":"response_item","payload":{{"type":"function_call","call_id":"call_1","name":"shell"}}}}"#
229        )
230        .unwrap();
231        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:24.000Z","type":"event_msg","payload":{{"type":"token_count","info":{{"total_token_usage":{{"input_tokens":14778,"cached_input_tokens":12672,"output_tokens":241,"total_tokens":15019}}}}}}}}"#).unwrap();
232        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:25.000Z","type":"event_msg","payload":{{"type":"token_count","info":null}}}}"#)
233            .unwrap();
234        let mut t = CodexTranscript::new(&path);
235        t.refresh().unwrap();
236        let s = t.summary();
237        assert_eq!(s.session_id.as_deref(), Some("01a0"));
238        assert_eq!(s.model.as_deref(), Some("gpt-5-codex"));
239        assert_eq!(s.usage.input, 14778 - 12672);
240        assert_eq!(s.usage.cache_read, 12672);
241        assert_eq!(s.usage.total(), 15019);
242        assert_eq!(s.unpriced_tokens, 15019);
243        assert_eq!(s.tool_calls, 1);
244        assert_eq!(s.activity, Activity::Working);
245        assert_eq!(read_meta(&path).unwrap().0, PathBuf::from("/tmp/p"));
246        let spans = s.spans.to_vec();
247        assert_eq!(spans.len(), 1);
248        assert_eq!(spans[0].name, "shell");
249        assert!(spans[0].is_open(), "no output item yet");
250        let _ = std::fs::remove_dir_all(&dir);
251    }
252
253    #[test]
254    fn pairs_calls_with_their_outputs() {
255        let dir = std::env::temp_dir().join(format!("agent-top-codex-spans-{}", std::process::id()));
256        std::fs::create_dir_all(&dir).unwrap();
257        let path = dir.join("rollout.jsonl");
258        let mut f = std::fs::File::create(&path).unwrap();
259        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:23.000Z","type":"response_item","payload":{{"type":"function_call","call_id":"call_1","name":"exec_command"}}}}"#).unwrap();
260        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:23.100Z","type":"response_item","payload":{{"type":"custom_tool_call","call_id":"call_2","name":"apply_patch"}}}}"#).unwrap();
261        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:24.000Z","type":"response_item","payload":{{"type":"function_call_output","call_id":"call_1","output":"ok"}}}}"#).unwrap();
262        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:26.100Z","type":"response_item","payload":{{"type":"custom_tool_call_output","call_id":"call_2","output":"ok"}}}}"#).unwrap();
263        let mut t = CodexTranscript::new(&path);
264        t.refresh().unwrap();
265        let spans = t.summary().spans.to_vec();
266        assert_eq!(spans.len(), 2);
267        assert_eq!(spans[0].name, "exec_command");
268        assert_eq!(spans[0].duration_ms, Some(1_000));
269        assert_eq!(spans[1].name, "apply_patch");
270        assert_eq!(spans[1].duration_ms, Some(3_000));
271        assert_eq!(t.summary().tool_calls, 2);
272        let _ = std::fs::remove_dir_all(&dir);
273    }
274}