Skip to main content

tokenburn_core/features/codex/
ops.rs

1use std::collections::HashSet;
2use std::path::{Path, PathBuf};
3use std::sync::LazyLock;
4
5use anyhow::Result;
6use chrono::{DateTime, Utc};
7use serde_json::Value;
8
9use crate::features::usage::{Row, Tool};
10use crate::utils::cache::{FileCache, Sig};
11use crate::utils::files::{dir_basename, env_path, modified_since, num, read_lines, walk_ext};
12use crate::utils::time::parse_ts;
13
14/// `$CODEX_HOME` (default `~/.codex`) → its `sessions/` and `archived_sessions/`.
15pub fn roots() -> Vec<PathBuf> {
16    let home = env_path("CODEX_HOME").or_else(|| dirs::home_dir().map(|h| h.join(".codex")));
17    home.map(|h| vec![h.join("sessions"), h.join("archived_sessions")])
18        .unwrap_or_default()
19}
20
21pub fn collect_codex(start: DateTime<Utc>) -> Result<Vec<Row>> {
22    let mut rows = Vec::new();
23    for root in roots().into_iter().filter(|r| r.is_dir()) {
24        rows.extend(collect_codex_from(&root, start)?);
25    }
26    Ok(rows)
27}
28
29/// Parsed rows per rollout file, reused while the file is unchanged.
30static CACHE: LazyLock<FileCache<Vec<Row>>> = LazyLock::new(FileCache::default);
31
32pub fn collect_codex_from(root: &Path, start: DateTime<Utc>) -> Result<Vec<Row>> {
33    let files = walk_ext(root, &["jsonl"]);
34    let live: HashSet<PathBuf> = files.iter().cloned().collect();
35    let mut rows = Vec::new();
36    for file in files {
37        if !modified_since(&file, start) {
38            continue;
39        }
40        let Some(sig) = Sig::of(&file, 0) else {
41            continue;
42        };
43        let parsed = CACHE.get_or_parse(&file, sig, || {
44            read_lines(&file)
45                .map(|lines| parse_session(lines, &file))
46                .unwrap_or_default()
47        });
48        rows.extend(parsed.iter().cloned());
49    }
50    CACHE.prune_under(root, &live);
51    Ok(rows)
52}
53
54/// Cumulative token counters of one Codex session.
55#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
56struct Totals {
57    input: u64,
58    cached: u64,
59    output: u64,
60}
61
62impl Totals {
63    fn of(v: &Value) -> Totals {
64        Totals {
65            input: num(v.get("input_tokens")),
66            cached: num(v.get("cached_input_tokens")),
67            output: num(v.get("output_tokens")),
68        }
69    }
70    fn is_zero(self) -> bool {
71        self == Totals::default()
72    }
73}
74
75/// Codex emits a `token_count` event after every model call carrying the
76/// *cumulative* `total_token_usage` (and often a repeated, unchanged copy), so
77/// per-call usage is the positive delta between successive totals. Sessions
78/// without totals fall back to `last_token_usage`.
79fn parse_session(lines: impl Iterator<Item = String>, file: &Path) -> Vec<Row> {
80    let id = file
81        .file_stem()
82        .map(|s| s.to_string_lossy().into_owned())
83        .unwrap_or_default();
84    let mut project: Option<String> = None;
85    let mut model: Option<String> = None;
86    let mut prev = Totals::default();
87    let mut rows = Vec::new();
88
89    for line in lines {
90        let Ok(rec) = serde_json::from_str::<Value>(&line) else {
91            continue;
92        };
93        let payload = rec.get("payload");
94        match rec.get("type").and_then(Value::as_str) {
95            Some("session_meta") => {
96                project = payload
97                    .and_then(|p| p.get("cwd"))
98                    .and_then(Value::as_str)
99                    .and_then(dir_basename);
100            }
101            Some("turn_context") => {
102                if let Some(m) = payload.and_then(|p| p.get("model")).and_then(Value::as_str) {
103                    model = Some(m.to_string());
104                }
105            }
106            Some("event_msg") => {}
107            _ => continue,
108        }
109        let Some(p) =
110            payload.filter(|p| p.get("type").and_then(Value::as_str) == Some("token_count"))
111        else {
112            continue;
113        };
114        let Some(info) = p.get("info").filter(|i| i.is_object()) else {
115            continue;
116        };
117        if let Some(m) = info.get("model").and_then(Value::as_str) {
118            model = Some(m.to_string());
119        }
120        let Some(ts) = rec
121            .get("timestamp")
122            .and_then(Value::as_str)
123            .and_then(parse_ts)
124        else {
125            continue;
126        };
127
128        let delta = if let Some(t) = info.get("total_token_usage").filter(|t| t.is_object()) {
129            let cur = Totals::of(t);
130            let d = Totals {
131                input: cur.input.saturating_sub(prev.input),
132                cached: cur.cached.saturating_sub(prev.cached),
133                output: cur.output.saturating_sub(prev.output),
134            };
135            prev = cur;
136            d
137        } else if let Some(l) = info.get("last_token_usage").filter(|l| l.is_object()) {
138            Totals::of(l)
139        } else {
140            continue;
141        };
142        if delta.is_zero() {
143            continue;
144        }
145        rows.push(Row {
146            tool: Tool::Codex,
147            project: project
148                .clone()
149                .or_else(|| model.clone())
150                .unwrap_or_else(|| "unknown".into()),
151            id: id.clone(),
152            ts,
153            // OpenAI counts cached tokens inside `input_tokens`: split them out.
154            input: delta.input.saturating_sub(delta.cached),
155            output: delta.output,
156            cache_read: delta.cached,
157            cache_write: 0,
158            cost: 0.0,
159        });
160    }
161    rows
162}
163
164#[cfg(test)]
165mod tests {
166    use super::*;
167    use serde_json::json;
168
169    fn meta(cwd: &str) -> String {
170        json!({"timestamp":"2026-05-01T10:00:00Z","type":"session_meta","payload":{"cwd":cwd}})
171            .to_string()
172    }
173
174    fn tc(ts: &str, total: (u64, u64, u64)) -> String {
175        json!({"timestamp":ts,"type":"event_msg","payload":{"type":"token_count","info":{
176            "total_token_usage":{"input_tokens":total.0,"cached_input_tokens":total.1,"output_tokens":total.2}}}})
177        .to_string()
178    }
179
180    fn rows(lines: Vec<String>) -> Vec<Row> {
181        parse_session(lines.into_iter(), Path::new("rollout-1.jsonl"))
182    }
183
184    #[test]
185    fn cumulative_totals_become_per_call_deltas() {
186        let r = rows(vec![
187            meta("/home/me/app"),
188            tc("2026-05-01T10:00:01Z", (100, 40, 10)),
189            tc("2026-05-01T10:00:02Z", (100, 40, 10)), // repeated, unchanged → ignored
190            tc("2026-05-01T10:00:03Z", (250, 90, 30)),
191        ]);
192        assert_eq!(r.len(), 2);
193        // first call: 100 input of which 40 cached
194        assert_eq!((r[0].input, r[0].cache_read, r[0].output), (60, 40, 10));
195        // second call: +150 input of which +50 cached, +20 output
196        assert_eq!((r[1].input, r[1].cache_read, r[1].output), (100, 50, 20));
197        assert_eq!(r[0].project, "app");
198        assert_eq!(r[0].tool, Tool::Codex);
199    }
200
201    #[test]
202    fn the_sum_of_deltas_equals_the_final_total() {
203        let r = rows(vec![
204            tc("2026-05-01T10:00:01Z", (10, 0, 1)),
205            tc("2026-05-01T10:00:02Z", (30, 5, 4)),
206            tc("2026-05-01T10:00:03Z", (31, 5, 9)),
207        ]);
208        let input: u64 = r.iter().map(|x| x.input + x.cache_read).sum();
209        let output: u64 = r.iter().map(|x| x.output).sum();
210        assert_eq!((input, output), (31, 9));
211    }
212
213    #[test]
214    fn falls_back_to_last_token_usage_without_totals() {
215        let l = json!({"timestamp":"2026-05-01T10:00:00Z","type":"event_msg","payload":{"type":"token_count","info":{
216            "model":"gpt-x","last_token_usage":{"input_tokens":7,"cached_input_tokens":2,"output_tokens":3}}}});
217        let r = rows(vec![l.to_string()]);
218        assert_eq!((r[0].input, r[0].cache_read, r[0].output), (5, 2, 3));
219        assert_eq!(
220            r[0].project, "gpt-x",
221            "falls back to the model when there is no cwd"
222        );
223    }
224
225    #[test]
226    fn ignores_other_events_and_null_info() {
227        let null_info = json!({"timestamp":"2026-05-01T10:00:00Z","type":"event_msg","payload":{"type":"token_count","info":null}});
228        let other = json!({"timestamp":"2026-05-01T10:00:00Z","type":"event_msg","payload":{"type":"agent_message"}});
229        assert!(rows(vec![
230            null_info.to_string(),
231            other.to_string(),
232            "garbage".into()
233        ])
234        .is_empty());
235    }
236
237    #[test]
238    fn model_from_turn_context_is_remembered() {
239        let ctx = json!({"timestamp":"2026-05-01T10:00:00Z","type":"turn_context","payload":{"model":"gpt-5"}});
240        let r = rows(vec![ctx.to_string(), tc("2026-05-01T10:00:01Z", (5, 0, 1))]);
241        assert_eq!(r[0].project, "gpt-5");
242    }
243
244    #[test]
245    fn collects_from_a_dated_directory_tree() {
246        let d = tempfile::tempdir().unwrap();
247        let day = d.path().join("2026/05/01");
248        std::fs::create_dir_all(&day).unwrap();
249        std::fs::write(
250            day.join("rollout-a.jsonl"),
251            tc("2026-05-01T10:00:01Z", (8, 0, 2)),
252        )
253        .unwrap();
254        let r = collect_codex_from(d.path(), DateTime::<Utc>::UNIX_EPOCH).unwrap();
255        assert_eq!(r.len(), 1);
256        assert_eq!(r[0].id, "rollout-a");
257    }
258}