1use super::{REFRESH_BUDGET_BYTES, SessionSummary, SessionTracker, 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
37pub 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
61pub 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 pub fn with_prices(mut self, prices: &'static Table) -> Self {
93 self.prices = prices;
94 self
95 }
96
97 fn ingest(&mut self, line: &str) {
98 let Ok(v) = serde_json::from_str::<Value>(line) else { return };
99 let ts = v.get("timestamp").and_then(Value::as_str).and_then(parse_rfc3339_utc);
100 if let Some(ts) = ts {
101 if self.summary.started_at.is_none() {
102 self.summary.started_at = Some(ts);
103 }
104 self.summary.last_activity = Some(ts);
105 }
106 let kind = v.get("type").and_then(Value::as_str).unwrap_or("");
107 let payload = v.get("payload");
108 let ptype = payload.and_then(|p| p.get("type")).and_then(Value::as_str).unwrap_or("");
109 match kind {
110 "session_meta" => {
111 if let Some(p) = payload {
112 self.summary.session_id = p.get("id").or(p.get("session_id")).and_then(Value::as_str).map(str::to_string);
113 self.summary.cwd = p.get("cwd").and_then(Value::as_str).map(PathBuf::from);
114 self.summary.harness_version = p.get("cli_version").and_then(Value::as_str).map(str::to_string);
115 }
116 }
117 "turn_context" => {
118 if let Some(m) = payload.and_then(|p| p.get("model")).and_then(Value::as_str) {
119 self.summary.model = Some(m.to_string());
120 }
121 }
122 "event_msg" => match ptype {
123 "token_count" => {
124 if let Some(total) = payload.and_then(|p| p.pointer("/info/total_token_usage")) {
125 let g = |k: &str| total.get(k).and_then(Value::as_u64).unwrap_or(0);
126 self.summary.health.usage_records += 1;
127 if g("input_tokens") + g("output_tokens") + g("cached_input_tokens") == 0 {
128 self.summary.health.empty_usage_records += 1;
129 }
130 let cached = g("cached_input_tokens");
131 let usage = TokenUsage {
132 input: g("input_tokens").saturating_sub(cached),
133 cache_read: cached,
134 output: g("output_tokens"),
135 ..Default::default()
136 };
137 self.summary.usage = usage;
138 let price = self.summary.model.as_deref().and_then(|m| self.prices.lookup(m));
139 match price {
140 Some(p) => {
141 self.summary.cost_usd = p.cost(&usage);
142 self.summary.unpriced_tokens = 0;
143 }
144 None => {
145 self.summary.cost_usd = 0.0;
146 self.summary.unpriced_tokens = usage.total();
147 }
148 }
149 }
150 }
151 "task_started" | "user_message" => self.summary.activity = Activity::Working,
152 "task_complete" | "turn_aborted" | "error" => self.summary.activity = Activity::Waiting,
153 _ => {}
154 },
155 "response_item" => match ptype {
156 "function_call" | "custom_tool_call" | "local_shell_call" => {
157 self.summary.tool_calls += 1;
158 if let (Some(ts), Some(p)) = (ts, payload) {
159 let id = call_id(p);
160 let name = p.get("name").and_then(Value::as_str).unwrap_or(ptype);
161 self.summary.spans.open(id, name.to_string(), ts, false);
162 }
163 }
164 "function_call_output" | "custom_tool_call_output" | "local_shell_call_output" => {
165 if let (Some(ts), Some(p)) = (ts, payload) {
166 self.summary.spans.close(&call_id(p), ts, false);
170 }
171 }
172 "message" if payload.and_then(|p| p.get("role")).and_then(Value::as_str) == Some("assistant") => {
173 self.summary.turns += 1;
174 self.summary.health.billable_messages += 1;
175 }
176 _ => {}
177 },
178 _ => {}
179 }
180 }
181}
182
183fn call_id(payload: &Value) -> String {
185 payload.get("call_id").or_else(|| payload.get("id")).and_then(Value::as_str).unwrap_or_default().to_string()
186}
187
188impl SessionTracker for CodexTranscript {
189 fn refresh(&mut self) -> anyhow::Result<bool> {
190 let (lines, more) = self.reader.read_new_lines(REFRESH_BUDGET_BYTES)?;
191 for l in &lines {
192 self.ingest(l);
193 }
194 Ok(more)
195 }
196
197 fn summary(&self) -> &SessionSummary {
198 &self.summary
199 }
200
201 fn path(&self) -> &Path {
202 self.reader.path()
203 }
204}
205
206#[cfg(test)]
207mod tests {
208 use super::*;
209 use std::io::Write;
210
211 #[test]
212 fn reads_cumulative_usage_and_state() {
213 let dir = std::env::temp_dir().join(format!("agent-top-codex-{}", std::process::id()));
214 std::fs::create_dir_all(&dir).unwrap();
215 let path = dir.join("rollout.jsonl");
216 let mut f = std::fs::File::create(&path).unwrap();
217 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();
218 writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:21.000Z","type":"turn_context","payload":{{"model":"gpt-5-codex"}}}}"#).unwrap();
219 writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:22.000Z","type":"event_msg","payload":{{"type":"task_started"}}}}"#).unwrap();
220 writeln!(
221 f,
222 r#"{{"timestamp":"2026-08-28T08:53:23.000Z","type":"response_item","payload":{{"type":"function_call","call_id":"call_1","name":"shell"}}}}"#
223 )
224 .unwrap();
225 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();
226 writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:25.000Z","type":"event_msg","payload":{{"type":"token_count","info":null}}}}"#)
227 .unwrap();
228 let mut t = CodexTranscript::new(&path);
229 t.refresh().unwrap();
230 let s = t.summary();
231 assert_eq!(s.session_id.as_deref(), Some("01a0"));
232 assert_eq!(s.model.as_deref(), Some("gpt-5-codex"));
233 assert_eq!(s.usage.input, 14778 - 12672);
234 assert_eq!(s.usage.cache_read, 12672);
235 assert_eq!(s.usage.total(), 15019);
236 assert_eq!(s.unpriced_tokens, 15019);
237 assert_eq!(s.tool_calls, 1);
238 assert_eq!(s.activity, Activity::Working);
239 assert_eq!(read_meta(&path).unwrap().0, PathBuf::from("/tmp/p"));
240 let spans = s.spans.to_vec();
241 assert_eq!(spans.len(), 1);
242 assert_eq!(spans[0].name, "shell");
243 assert!(spans[0].is_open(), "no output item yet");
244 let _ = std::fs::remove_dir_all(&dir);
245 }
246
247 #[test]
248 fn pairs_calls_with_their_outputs() {
249 let dir = std::env::temp_dir().join(format!("agent-top-codex-spans-{}", std::process::id()));
250 std::fs::create_dir_all(&dir).unwrap();
251 let path = dir.join("rollout.jsonl");
252 let mut f = std::fs::File::create(&path).unwrap();
253 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();
254 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();
255 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();
256 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();
257 let mut t = CodexTranscript::new(&path);
258 t.refresh().unwrap();
259 let spans = t.summary().spans.to_vec();
260 assert_eq!(spans.len(), 2);
261 assert_eq!(spans[0].name, "exec_command");
262 assert_eq!(spans[0].duration_ms, Some(1_000));
263 assert_eq!(spans[1].name, "apply_patch");
264 assert_eq!(spans[1].duration_ms, Some(3_000));
265 assert_eq!(t.summary().tool_calls, 2);
266 let _ = std::fs::remove_dir_all(&dir);
267 }
268}