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//! * `task_started` and `task_complete` bracket a turn span. An inference
15//!   span runs from a user `message` item or a `*_output` item to the next
16//!   thing the model produced: a call, a `reasoning` item, a
17//!   `web_search_call`, or an assistant `message`.
18//! * `response_item` `web_search_call` is one server-side web search.
19//! * `info.last_token_usage` beside the cumulative record is the one
20//!   response's usage. Repeated cumulative usage marks a snapshot; equal
21//!   per-response usage alone need not. Each `*_output` item is filed under
22//!   its call's name, re-filed under the MCP server when its
23//!   `mcp_tool_call_end` follows, and sized by the consuming response. Codex
24//!   writes no compaction marker that was seen, so the ledger's halving
25//!   rule stands in. See `ContextLedger`.
26//! * Usage ordering (0.130 fixture / 0.154.0 live and source, 2026-09-18):
27//!   newer versions emit `token_count` after draining tool outputs, older
28//!   ones before. The first model item freezes that response's input batch;
29//!   results produced during it wait for the next response.
30//!
31//! Codex model prices are not in the static table, so cost is reported as
32//! unpriced tokens.
33//!
34//! Process semantics (source-verified against Codex 0.154.0, 2026-09-17):
35//! the npm launcher spawns a native runtime, which owns the rollout handles.
36//! Native `spawn_agent` creates an in-process session, not an OS process.
37//! Its lineage is `payload.source.subagent.thread_spawn.parent_thread_id` in
38//! session metadata, or `payload.parent_thread_id` when `thread_source` is
39//! explicitly `subagent`, never a PID relationship or `forked_from_id`.
40//! Nickname and role come from that source record (or the top-level metadata);
41//! older records name the role `agent_type`. Rollouts retain separate usage
42//! totals while the TUI groups them by parent. Memory belongs to the process,
43//! not to individual logical subagents sharing that process.
44//! Live 0.154.0 rollouts (2026-09-17) can copy ancestor `session_meta` records
45//! after the child's own header when forking history. The first identified
46//! header owns this rollout's identity and lineage; later headers do not.
47//!
48//! Code mode (live/source verified on 0.154.0, 2026-09-18): `exec` wraps
49//! nested tools recorded as `event_msg/item_completed`, with `exec-<uuid>`
50//! ids and `started_at_ms`/`completed_at_ms`. Typed metadata supplies names;
51//! unknown types are ignored. Counts include wrappers and nested calls;
52//! each wrapper's context share is split between its contained children by
53//! output-text bytes (evenly if sizes are unavailable). Only sizes are retained;
54//! no prompts or inputs are inspected. Ambiguous wrappers keep their name.
55
56use super::{
57    AttributeContext, ContextWeights, HarnessAdapter, REFRESH_BUDGET_BYTES, SessionSummary, SessionTracker, SpanRetention,
58    parse_rfc3339_utc,
59};
60use crate::jsonl::TailReader;
61use crate::model::{Activity, Attribution, ContextOrigin, Harness, ProcNode, SpanKind, SubagentInfo, TokenUsage};
62use crate::pricing::{self, Table};
63use crate::process::RawProc;
64use serde_json::Value;
65use std::collections::{HashMap, HashSet};
66use std::path::{Path, PathBuf};
67use std::time::{Duration, SystemTime};
68
69pub fn codex_dir() -> Option<PathBuf> {
70    if let Some(d) = std::env::var_os("CODEX_HOME") {
71        return Some(PathBuf::from(d));
72    }
73    std::env::var_os("HOME").map(|h| PathBuf::from(h).join(".codex"))
74}
75
76pub fn sessions_dir() -> Option<PathBuf> {
77    codex_dir().map(|d| d.join("sessions"))
78}
79
80/// Rollout files modified after `since`. Walks `YYYY/MM/DD` and prunes by
81/// directory mtime so the walk stays cheap on a long history.
82/// Rollouts the process has open: the app-server's live threads, or the CLI's
83/// one conversation. `None` when the platform cannot say. Filtered to the
84/// sessions directory so an unrelated file the process holds (a log, a
85/// config) is never mistaken for a thread, and mapped back under the
86/// un-canonicalised sessions directory so the paths compare equal to those
87/// from `recent_rollouts`.
88pub fn rollouts_open_by(pid: u32) -> Option<Vec<PathBuf>> {
89    let root = sessions_dir()?;
90    let canonical = std::fs::canonicalize(&root).unwrap_or_else(|_| root.clone());
91    let open = crate::openfiles::open_files(pid)?;
92    Some(
93        open.into_iter()
94            .filter(|p| p.extension().and_then(|x| x.to_str()) == Some("jsonl"))
95            .filter_map(|p| p.strip_prefix(&canonical).ok().map(|rel| root.join(rel)))
96            .collect(),
97    )
98}
99
100/// An npm launcher holds no rollouts; its native runtime does. Match the
101/// forwarded argv, never just an `Agent` child: a nested agent invocation
102/// owns its own sessions even though it is beneath this process.
103fn rollout_owner<'a>(root: &'a ProcNode, by_pid: &HashMap<u32, &RawProc>) -> &'a ProcNode {
104    let Some(parent) = by_pid.get(&root.pid) else { return root };
105    root.children
106        .iter()
107        .find(|p| p.harness == Some(Harness::Codex) && by_pid.get(&p.pid).is_some_and(|child| child.is_codex_runtime_of(parent)))
108        .unwrap_or(root)
109}
110
111/// Every rollout written since `since`.
112pub fn recent_rollouts(since: SystemTime) -> Vec<PathBuf> {
113    let Some(root) = sessions_dir() else { return Vec::new() };
114    rollouts_under(&root, since)
115}
116
117/// The tree is `YYYY/MM/DD/*.jsonl` and is walked in full, three levels deep,
118/// with only the files filtered by mtime. Pruning directories by their mtime
119/// looked cheaper and was wrong: a directory's mtime moves only when an entry
120/// is created directly inside it, so the year directory is touched once a
121/// month and every rollout written after the first of the month was invisible.
122/// Pruning by name would be wrong too, since a directory's date says when a
123/// thread started, not whether it is still being written to; the app-server
124/// keeps a thread for days. A few hundred directories cost a few milliseconds.
125pub(crate) fn rollouts_under(root: &Path, since: SystemTime) -> Vec<PathBuf> {
126    let mut out = Vec::new();
127    walk(root, 0, since, &mut out);
128    out
129}
130
131fn walk(dir: &Path, depth: usize, since: SystemTime, out: &mut Vec<PathBuf>) {
132    let Ok(rd) = std::fs::read_dir(dir) else { return };
133    for e in rd.flatten() {
134        let p = e.path();
135        let Ok(md) = e.metadata() else { continue };
136        if md.is_dir() {
137            if depth < 3 {
138                walk(&p, depth + 1, since, out);
139            }
140        } else if p.extension().and_then(|x| x.to_str()) == Some("jsonl") && md.modified().map(|m| m >= since).unwrap_or(false) {
141            out.push(p);
142        }
143    }
144}
145
146/// Cheap header read: cwd and start time from the first line only.
147pub fn read_meta(path: &Path) -> Option<(PathBuf, SystemTime)> {
148    use std::io::{BufRead, BufReader};
149    let f = std::fs::File::open(path).ok()?;
150    let mut first = String::new();
151    BufReader::new(f).read_line(&mut first).ok()?;
152    let v: Value = serde_json::from_str(&first).ok()?;
153    if v.get("type").and_then(Value::as_str) != Some("session_meta") {
154        return None;
155    }
156    let cwd = v.pointer("/payload/cwd").and_then(Value::as_str).map(PathBuf::from)?;
157    let ts = v.get("timestamp").and_then(Value::as_str).and_then(parse_rfc3339_utc)?;
158    Some((cwd, ts))
159}
160
161/// The Codex adapter: a process is matched to the rollouts it holds open,
162/// and only where the platform cannot say to the cwd and activity heuristics.
163/// See DEC-006.
164#[derive(Default)]
165pub struct CodexAdapter {
166    /// Recent rollouts with the cwd and start time from their header.
167    recent: Vec<(PathBuf, PathBuf, SystemTime)>,
168    /// Which rollouts each Codex process has open, gathered before any
169    /// attribution so that no process's fallback can claim a thread another
170    /// process is demonstrably writing. `None` when the platform cannot say.
171    held: HashMap<u32, Option<Vec<PathBuf>>>,
172    all_held: HashSet<PathBuf>,
173}
174
175impl HarnessAdapter for CodexAdapter {
176    fn harness(&self) -> Harness {
177        Harness::Codex
178    }
179
180    fn rescan(&mut self, since: SystemTime) {
181        self.recent = recent_rollouts(since).into_iter().filter_map(|p| read_meta(&p).map(|(cwd, ts)| (p, cwd, ts))).collect();
182    }
183
184    fn prepare(&mut self, roots: &[&ProcNode], by_pid: &HashMap<u32, &RawProc>) {
185        self.held = roots.iter().map(|r| (r.pid, rollouts_open_by(rollout_owner(r, by_pid).pid))).collect();
186        self.all_held = self.held.values().flatten().flatten().cloned().collect();
187    }
188
189    fn attribute(&self, root: &ProcNode, _raw: Option<&RawProc>, ctx: &AttributeContext) -> (Vec<PathBuf>, Attribution) {
190        let mine: Option<Vec<PathBuf>> =
191            self.held.get(&root.pid).and_then(|h| h.as_ref()).map(|h| h.iter().filter(|p| !ctx.attached.contains(*p)).cloned().collect());
192        let taken: HashSet<PathBuf> = ctx.attached.union(&self.all_held).cloned().collect();
193        attribute(ctx.cwd, ctx.proc_start, mine.as_deref(), &self.recent, &taken, ctx.now, ctx.activity_timeout)
194    }
195
196    fn unowned(&self, attached: &HashSet<PathBuf>) -> Vec<PathBuf> {
197        self.recent.iter().map(|(p, _, _)| p).filter(|p| !attached.contains(*p)).cloned().collect()
198    }
199
200    fn open(&self, path: &Path, spans: SpanRetention) -> Box<dyn SessionTracker> {
201        Box::new(CodexTranscript::new(path).with_spans(spans))
202    }
203
204    /// Every rollout opens with a `session_meta` record.
205    fn detect(&self, path: &Path) -> bool {
206        super::head_lines(path).iter().any(|v| v.get("type").and_then(Value::as_str) == Some("session_meta"))
207    }
208
209    fn transcripts(&self) -> Vec<(String, PathBuf)> {
210        recent_rollouts(SystemTime::UNIX_EPOCH).into_iter().map(|p| (rollout_id(&p), p)).collect()
211    }
212}
213
214/// The id in `rollout-2026-05-14T21-37-50-<id>`: what follows the fixed-width
215/// timestamp. A file named some other way is matched on its whole stem.
216pub fn rollout_id(p: &Path) -> String {
217    const TS_LEN: usize = "2026-05-14T21-37-50-".len();
218    let stem = p.file_stem().map(|s| s.to_string_lossy().into_owned()).unwrap_or_default();
219    stem.strip_prefix("rollout-").and_then(|s| s.get(TS_LEN..)).map(str::to_string).unwrap_or(stem)
220}
221
222/// Codex conversations belonging to one process, newest activity first.
223///
224/// `held` are the rollouts the process has open, which is not a guess: Codex
225/// opens a thread's rollout when the thread starts and closes it when the
226/// thread ends. When the platform can say (`Some`), that list is the answer,
227/// an empty one included: a process holding no rollout is hosting no thread,
228/// and a rollout nobody holds is a finished conversation for the stopped
229/// list. The heuristics below are for when it cannot (`None`).
230///
231/// A `codex` CLI runs one conversation from the directory it was started in, so
232/// a cwd match finds it. The VS Code app-server is a different shape: one
233/// long-lived process, running from `/`, hosting any number of conversations
234/// over its life. Returning a single rollout for it collapses every one of
235/// those into one row and attributes whichever happened to be newest, so this
236/// returns all of them that are currently live and lets the caller give each
237/// its own row.
238///
239/// A rollout in `taken` is skipped: one already claimed by another process,
240/// or one some process has open, so that two Codex processes cannot both
241/// show the same conversation and an older app-server cannot collect the
242/// threads of a newer one.
243pub(crate) fn attribute(
244    cwd: Option<&Path>,
245    proc_start: SystemTime,
246    held: Option<&[PathBuf]>,
247    recent: &[(PathBuf, PathBuf, SystemTime)],
248    taken: &HashSet<PathBuf>,
249    now: SystemTime,
250    activity_timeout: Duration,
251) -> (Vec<PathBuf>, Attribution) {
252    if let Some(held) = held {
253        let mut mine = held.to_vec();
254        mine.sort_by_key(|p| std::cmp::Reverse(written_at(p)));
255        mine.truncate(MAX_THREADS);
256        let attribution = if mine.is_empty() { Attribution::None } else { Attribution::OpenFile };
257        return (mine, attribution);
258    }
259
260    let slack = Duration::from_secs(60);
261    let started_after = |ts: &SystemTime| *ts + slack >= proc_start;
262    let candidates = || recent.iter().filter(|(p, _, ts)| started_after(ts) && !taken.contains(p));
263
264    // The CLI case: the conversation runs where the process runs.
265    if let Some(cwd) = cwd {
266        let mut matched: Vec<&(PathBuf, PathBuf, SystemTime)> = candidates().filter(|(_, c, _)| c == cwd).collect();
267        if !matched.is_empty() {
268            matched.sort_by_key(|(p, _, _)| std::cmp::Reverse(written_at(p)));
269            return (matched.into_iter().map(|(p, _, _)| p.clone()).collect(), Attribution::CwdHeuristic);
270        }
271    }
272
273    // The app-server case: no cwd to match on, so take the conversations that
274    // are actually being written to. A rollout nobody has touched in a while is
275    // a finished conversation, not a thread of this process.
276    let mut live: Vec<&(PathBuf, PathBuf, SystemTime)> = candidates()
277        .filter(|(p, _, _)| written_at(p).map(|w| now.duration_since(w).unwrap_or_default() <= activity_timeout).unwrap_or(false))
278        .collect();
279    live.sort_by_key(|(p, _, _)| std::cmp::Reverse(written_at(p)));
280    live.truncate(MAX_THREADS);
281    let attribution = if live.is_empty() { Attribution::None } else { Attribution::CwdHeuristic };
282    (live.into_iter().map(|(p, _, _)| p.clone()).collect(), attribution)
283}
284
285/// One process is not plausibly running more conversations than this at once,
286/// and an unbounded fan-out would let a stale directory fill the table.
287const MAX_THREADS: usize = 12;
288
289fn written_at(p: &Path) -> Option<SystemTime> {
290    std::fs::metadata(p).and_then(|m| m.modified()).ok()
291}
292
293fn subagent_info(meta: &Value) -> Option<SubagentInfo> {
294    let source = meta.pointer("/source/subagent/thread_spawn");
295    let parent = source.and_then(|s| s.get("parent_thread_id")).and_then(Value::as_str).or_else(|| {
296        (meta.get("thread_source").and_then(Value::as_str) == Some("subagent"))
297            .then(|| meta.get("parent_thread_id").and_then(Value::as_str))
298            .flatten()
299    })?;
300    let id = meta.get("id").or_else(|| meta.get("session_id")).and_then(Value::as_str);
301    if parent.is_empty() || Some(parent) == id {
302        return None;
303    }
304    let field = |key| source.and_then(|s| s.get(key)).and_then(Value::as_str).or_else(|| meta.get(key).and_then(Value::as_str));
305    Some(SubagentInfo {
306        parent_session_id: parent.to_string(),
307        nickname: field("agent_nickname").map(str::to_string),
308        role: field("agent_role").or_else(|| field("agent_type")).map(str::to_string),
309    })
310}
311
312pub struct CodexTranscript {
313    reader: TailReader,
314    prices: &'static Table,
315    summary: SessionSummary,
316    /// Counters naming the turn and inference spans, and the ids of the ones
317    /// currently being extended.
318    turns: u64,
319    inferences: u64,
320    turn: Option<String>,
321    inference: Option<String>,
322    /// Tool calls awaiting their output, by call id, so the output can be
323    /// filed under the call's name.
324    pending_tools: HashMap<String, PendingTool>,
325    /// Completed results awaiting the next model response.
326    context_results: Vec<(String, ContextWeights)>,
327    response_in_progress: bool,
328    /// Per-response and cumulative usage, to skip repeated snapshots.
329    last_response: Option<(TokenUsage, Option<TokenUsage>)>,
330}
331
332struct PendingTool {
333    name: String,
334    started_at: SystemTime,
335    nested: Vec<NestedSource>,
336    ambiguous: bool,
337    active: bool,
338}
339
340struct NestedSource {
341    id: String,
342    origin: ContextOrigin,
343    name: String,
344    ended_at: SystemTime,
345    output_bytes: Option<u64>,
346}
347
348/// Measure decoded output text, never command arguments or patch contents.
349fn nested_output_bytes(item: &Value) -> Option<u64> {
350    let text = |key| item.get(key).and_then(Value::as_str).map(|s| s.len() as u64);
351    let streams = || match (text("stdout"), text("stderr")) {
352        (None, None) => None,
353        (out, err) => Some(out.unwrap_or(0) + err.unwrap_or(0)),
354    };
355    match item.get("type").and_then(Value::as_str)? {
356        "CommandExecution" => text("formatted_output").or_else(|| text("aggregated_output")).or_else(streams),
357        "FileChange" => streams(),
358        "McpToolCall" => match item.get("result").filter(|v| !v.is_null()) {
359            Some(result) if result.get("structuredContent").is_some_and(|v| !v.is_null()) => None,
360            Some(result) => text_content_bytes(result.get("content")?, "text"),
361            None => item.pointer("/error/message").and_then(Value::as_str).map(|s| s.len() as u64),
362        },
363        "DynamicToolCall" => match item.get("content_items").filter(|v| !v.is_null()) {
364            Some(content) => text_content_bytes(content, "inputText"),
365            None => text("error"),
366        },
367        _ => None,
368    }
369}
370
371fn text_content_bytes(content: &Value, kind: &str) -> Option<u64> {
372    content.as_array()?.iter().try_fold(0, |total, item| {
373        if item.get("type").and_then(Value::as_str) != Some(kind) {
374            return None;
375        }
376        Some(total + item.get("text")?.as_str()?.len() as u64)
377    })
378}
379
380impl CodexTranscript {
381    pub fn new(path: impl Into<PathBuf>) -> Self {
382        CodexTranscript {
383            reader: TailReader::new(path),
384            prices: pricing::table(),
385            summary: SessionSummary { harness: Some(Harness::Codex), ..Default::default() },
386            turns: 0,
387            inferences: 0,
388            turn: None,
389            inference: None,
390            pending_tools: HashMap::new(),
391            context_results: Vec::new(),
392            response_in_progress: false,
393            last_response: None,
394        }
395    }
396
397    /// Something was submitted to the model. One inference at a time: a
398    /// developer message followed by a user message is one submission.
399    fn begin_inference(&mut self, ts: SystemTime) {
400        if self.summary.spans.open_of_kind(SpanKind::Inference).is_some() {
401            return;
402        }
403        // A turn that ended without the model replying (aborted) leaves the
404        // previous inference open; it produced nothing, so it goes.
405        if let Some(id) = self.inference.take() {
406            self.summary.spans.discard_open(&id);
407        }
408        self.inferences += 1;
409        let id = format!("inference:{}", self.inferences);
410        self.summary.spans.open_kind(id.clone(), "inference".into(), ts, false, SpanKind::Inference);
411        self.inference = Some(id);
412    }
413
414    /// The model produced something: the inference in progress ends here.
415    fn end_inference(&mut self, ts: SystemTime) {
416        if let Some(id) = self.inference.take() {
417            self.summary.spans.end_at(&id, ts);
418        }
419    }
420
421    /// See `ClaudeTranscript::with_prices`.
422    pub fn with_prices(mut self, prices: &'static Table) -> Self {
423        self.prices = prices;
424        self
425    }
426
427    /// Keep every span instead of the newest `MAX_SPANS`. See `SpanRetention`.
428    pub fn with_spans(mut self, retention: SpanRetention) -> Self {
429        self.summary.spans = retention.log();
430        self
431    }
432
433    /// Nested `exec-<uuid>` calls lack response_item pairs; direct calls already
434    /// have spans and must not be counted again.
435    fn nested_tool_completed(&mut self, payload: &Value) {
436        let Some(item) = payload.get("item") else { return };
437        let Some(id) = item.get("id").and_then(Value::as_str).filter(|id| id.starts_with("exec-") && id.len() > 5) else {
438            return;
439        };
440        if self.summary.spans.iter().any(|s| s.id == id) || self.pending_tools.values().any(|p| p.nested.iter().any(|n| n.id == id)) {
441            return;
442        }
443        let name = match item.get("type").and_then(Value::as_str) {
444            Some("CommandExecution") => match item.get("source").and_then(Value::as_str) {
445                Some("unified_exec_startup") => Some("exec_command"),
446                Some("unified_exec_interaction") => Some("write_stdin"),
447                _ => None,
448            },
449            Some("FileChange") => Some("apply_patch"),
450            Some("McpToolCall" | "DynamicToolCall") => item.get("tool").and_then(Value::as_str).filter(|s| !s.is_empty()),
451            _ => None,
452        };
453        let time =
454            |key| payload.get(key).and_then(Value::as_u64).and_then(|ms| SystemTime::UNIX_EPOCH.checked_add(Duration::from_millis(ms)));
455        let (start, end) = (time("started_at_ms"), time("completed_at_ms"));
456        let source = match (name, item.get("type").and_then(Value::as_str)) {
457            (Some(_), Some("McpToolCall")) => {
458                item.get("server").and_then(Value::as_str).filter(|s| !s.is_empty()).map(|server| (ContextOrigin::Mcp, server))
459            }
460            (Some(name), _) => Some((ContextOrigin::Tool, name)),
461            _ => None,
462        };
463        // Timing is a heuristic; keep ambiguous results under the wrapper.
464        for pending in self.pending_tools.values_mut().filter(|p| p.active && matches!(p.name.as_str(), "exec" | "wait")) {
465            match (source, start, end) {
466                (Some((origin, name)), Some(start), Some(end)) if start >= pending.started_at && end >= start => {
467                    pending.nested.push(NestedSource {
468                        id: id.into(),
469                        origin,
470                        name: name.into(),
471                        ended_at: end,
472                        output_bytes: nested_output_bytes(item),
473                    });
474                }
475                _ => pending.ambiguous = true,
476            }
477        }
478        let (Some(name), Some(start), Some(end)) = (name, start, end) else { return };
479        if end < start {
480            return;
481        }
482        let error = matches!(item.get("status").and_then(Value::as_str), Some("failed" | "declined"))
483            || item.get("exit_code").and_then(Value::as_i64).is_some_and(|code| code != 0)
484            || item.get("success").and_then(Value::as_bool) == Some(false);
485        self.summary.tool_calls += 1;
486        self.summary.spans.open(id.to_string(), name.to_string(), start, false);
487        self.summary.spans.close(id, end, error);
488    }
489
490    fn ingest(&mut self, line: &str) {
491        let Ok(v) = serde_json::from_str::<Value>(line) else { return };
492        let ts = v.get("timestamp").and_then(Value::as_str).and_then(parse_rfc3339_utc);
493        if let Some(ts) = ts {
494            if self.summary.started_at.is_none() {
495                self.summary.started_at = Some(ts);
496            }
497            self.summary.last_activity = Some(ts);
498        }
499        let kind = v.get("type").and_then(Value::as_str).unwrap_or("");
500        let payload = v.get("payload");
501        let ptype = payload.and_then(|p| p.get("type")).and_then(Value::as_str).unwrap_or("");
502        if !self.response_in_progress
503            && kind == "response_item"
504            && (matches!(
505                ptype,
506                "function_call" | "custom_tool_call" | "local_shell_call" | "tool_search_call" | "web_search_call" | "reasoning"
507            ) || (ptype == "message" && payload.and_then(|p| p.get("role")).and_then(Value::as_str) == Some("assistant")))
508        {
509            for (id, sources) in self.context_results.drain(..) {
510                self.summary.context.result_weighted(&id, sources);
511            }
512            self.response_in_progress = true;
513        }
514        match kind {
515            // Forked history can contain parent headers, including across tail
516            // refreshes. They must not replace this rollout's own metadata.
517            "session_meta" if self.summary.session_id.is_none() => {
518                if let Some(p) = payload {
519                    self.summary.session_id = p.get("id").or(p.get("session_id")).and_then(Value::as_str).map(str::to_string);
520                    self.summary.subagent = subagent_info(p);
521                    self.summary.cwd = p.get("cwd").and_then(Value::as_str).map(PathBuf::from);
522                    self.summary.harness_version = p.get("cli_version").and_then(Value::as_str).map(str::to_string);
523                }
524            }
525            "turn_context" => {
526                if let Some(m) = payload.and_then(|p| p.get("model")).and_then(Value::as_str) {
527                    self.summary.model = Some(m.to_string());
528                }
529            }
530            "event_msg" => match ptype {
531                "token_count" => {
532                    let total = payload.and_then(|p| p.pointer("/info/total_token_usage")).filter(|v| v.is_object());
533                    if let Some(total) = total {
534                        let g = |k: &str| total.get(k).and_then(Value::as_u64).unwrap_or(0);
535                        self.summary.health.usage_records += 1;
536                        if g("input_tokens") + g("output_tokens") + g("cached_input_tokens") == 0 {
537                            self.summary.health.empty_usage_records += 1;
538                        }
539                        let cached = g("cached_input_tokens");
540                        let usage = TokenUsage {
541                            input: g("input_tokens").saturating_sub(cached),
542                            cache_read: cached,
543                            output: g("output_tokens"),
544                            ..Default::default()
545                        };
546                        self.summary.usage = usage;
547                        let price = self.summary.model.as_deref().and_then(|m| self.prices.lookup(m));
548                        match price {
549                            Some(p) => {
550                                self.summary.cost_breakdown = p.breakdown(&usage);
551                                self.summary.cost_usd = self.summary.cost_breakdown.total();
552                                self.summary.unpriced_tokens = 0;
553                            }
554                            None => {
555                                self.summary.cost_breakdown = Default::default();
556                                self.summary.cost_usd = 0.0;
557                                self.summary.unpriced_tokens = usage.total();
558                            }
559                        }
560                    }
561                    if let Some(last) = payload.and_then(|p| p.pointer("/info/last_token_usage")).filter(|v| v.is_object()) {
562                        let g = |k: &str| last.get(k).and_then(Value::as_u64).unwrap_or(0);
563                        let cached = g("cached_input_tokens");
564                        let usage = TokenUsage {
565                            input: g("input_tokens").saturating_sub(cached),
566                            cache_read: cached,
567                            output: g("output_tokens"),
568                            ..Default::default()
569                        };
570                        let response = (usage, total.map(|_| self.summary.usage));
571                        if usage.prompt() > 0 && self.last_response != Some(response) {
572                            let cost = self
573                                .summary
574                                .model
575                                .as_deref()
576                                .and_then(|m| self.prices.lookup(m))
577                                .map(|p| p.breakdown(&usage))
578                                .unwrap_or_default();
579                            self.summary.context.response(&usage, &cost);
580                            self.last_response = Some(response);
581                            self.response_in_progress = false;
582                        }
583                    }
584                    // The rate-limit snapshot rides on every token_count; the
585                    // latest one is the current state. Codex also reports other
586                    // quotas by `limit_id` (`premium`, seen since 0.131) that
587                    // can arrive with no windows at all; one of those says
588                    // nothing about the limit and must not blank the last one.
589                    if let Some(rl) = payload.and_then(|p| p.get("rate_limits")).filter(|v| v.is_object()) {
590                        let rl = parse_rate_limits(rl);
591                        if rl.primary.is_some() || rl.secondary.is_some() {
592                            self.summary.rate_limit = Some(rl);
593                        }
594                    }
595                }
596                "task_started" => {
597                    self.summary.activity = Activity::Working;
598                    if let Some(ts) = ts {
599                        self.turns += 1;
600                        let id = format!("turn:{}", self.turns);
601                        self.summary.spans.open_kind(id.clone(), "turn".into(), ts, false, SpanKind::Turn);
602                        self.turn = Some(id);
603                    }
604                }
605                "user_message" => self.summary.activity = Activity::Working,
606                "item_completed" => {
607                    if let Some(p) = payload {
608                        self.nested_tool_completed(p);
609                    }
610                }
611                // An MCP tool call. Codex records the call as a `response_item`
612                // `function_call` too, which the block below counts as a tool
613                // call and turns into a span; this line is the only one that
614                // names the server, so it feeds the per-server map and nothing
615                // else, to avoid double counting. `mcp_tool_call_begin` carries
616                // the same `invocation`; the pair brackets the call, but the
617                // `end` alone is enough for a count and is the one always
618                // present in the versions seen.
619                "mcp_tool_call_end" => {
620                    if let Some(inv) = payload.and_then(|p| p.get("invocation"))
621                        && let Some(server) = inv.get("server").and_then(Value::as_str).filter(|s| !s.is_empty())
622                    {
623                        let error = payload
624                            .and_then(|p| p.get("result"))
625                            .and_then(Value::as_object)
626                            .map(|r| !r.contains_key("Ok"))
627                            .unwrap_or(false);
628                        let u = self.summary.mcp.entry(server.to_string()).or_default();
629                        u.calls += 1;
630                        u.errors += u64::from(error);
631                        u.last_call = u.last_call.max(ts);
632                        let id = payload.map(call_id).unwrap_or_default();
633                        if let Some((_, sources)) = self.context_results.iter_mut().find(|(i, _)| *i == id) {
634                            for (origin, name, _) in sources {
635                                *origin = ContextOrigin::Mcp;
636                                *name = server.to_string();
637                            }
638                        }
639                        self.summary.context.retag(&id, ContextOrigin::Mcp, server);
640                    }
641                }
642                "task_complete" | "turn_aborted" | "error" => {
643                    self.summary.activity = Activity::Waiting;
644                    if let (Some(ts), Some(id)) = (ts, self.turn.take()) {
645                        self.summary.spans.end_at(&id, ts);
646                    }
647                    if let Some(id) = self.inference.take() {
648                        self.summary.spans.discard_open(&id);
649                    }
650                    // Keep names for late outputs, not overlap attribution.
651                    for pending in self.pending_tools.values_mut() {
652                        pending.active = false;
653                        pending.ambiguous = true;
654                    }
655                    self.response_in_progress = false;
656                }
657                _ => {}
658            },
659            "response_item" => match ptype {
660                "function_call" | "custom_tool_call" | "local_shell_call" => {
661                    self.summary.tool_calls += 1;
662                    if let (Some(ts), Some(p)) = (ts, payload) {
663                        self.end_inference(ts);
664                        let id = call_id(p);
665                        let name = p.get("name").and_then(Value::as_str).unwrap_or(ptype);
666                        let mut ambiguous = false;
667                        if matches!(name, "exec" | "wait") {
668                            for pending in
669                                self.pending_tools.values_mut().filter(|p| p.active && matches!(p.name.as_str(), "exec" | "wait"))
670                            {
671                                pending.ambiguous = true;
672                                ambiguous = true;
673                            }
674                        }
675                        self.pending_tools.insert(
676                            id.clone(),
677                            PendingTool { name: name.into(), started_at: ts, nested: Vec::new(), ambiguous, active: true },
678                        );
679                        self.summary.spans.open(id, name.to_string(), ts, false);
680                    }
681                }
682                "function_call_output" | "custom_tool_call_output" | "local_shell_call_output" => {
683                    if let (Some(ts), Some(p)) = (ts, payload) {
684                        // Wrapper output text is opaque; errors need typed metadata.
685                        let id = call_id(p);
686                        self.summary.spans.close(&id, ts, false);
687                        let sources = match self.pending_tools.remove(&id) {
688                            Some(p) if !p.ambiguous && !p.nested.is_empty() && p.nested.iter().all(|n| n.ended_at <= ts) => {
689                                let sized = p.nested.iter().all(|n| n.output_bytes.is_some());
690                                p.nested
691                                    .into_iter()
692                                    .map(|n| (n.origin, n.name, if sized { n.output_bytes.unwrap_or(0) } else { 1 }))
693                                    .collect()
694                            }
695                            Some(p) => vec![(ContextOrigin::Tool, p.name, 1)],
696                            None => vec![(ContextOrigin::Tool, "tool".into(), 1)],
697                        };
698                        self.context_results.push((id, sources));
699                        self.begin_inference(ts);
700                    }
701                }
702                // A server-side web search: billed per search by OpenAI, but
703                // at a rate this table does not carry, so counted only.
704                "web_search_call" => {
705                    self.summary.web_searches += 1;
706                    if let Some(ts) = ts {
707                        self.end_inference(ts);
708                    }
709                }
710                "reasoning" => {
711                    if let Some(ts) = ts {
712                        self.end_inference(ts);
713                    }
714                }
715                "message" => match payload.and_then(|p| p.get("role")).and_then(Value::as_str) {
716                    Some("assistant") => {
717                        self.summary.turns += 1;
718                        self.summary.health.billable_messages += 1;
719                        if let Some(ts) = ts {
720                            self.end_inference(ts);
721                        }
722                    }
723                    Some("user") => {
724                        if let Some(ts) = ts {
725                            self.begin_inference(ts);
726                        }
727                    }
728                    _ => {}
729                },
730                _ => {}
731            },
732            _ => {}
733        }
734    }
735}
736
737/// Codex's `rate_limits`: a short window (`primary`) and a long one
738/// (`secondary`), each a used-percent, a window length and a reset time in
739/// epoch seconds, plus the plan and whether the limit is currently hit.
740fn parse_rate_limits(v: &Value) -> crate::model::RateLimit {
741    use crate::model::{RateLimit, RateWindow};
742    let window = |w: Option<&Value>| -> Option<RateWindow> {
743        // `"primary": null` is a missing window, not an empty one.
744        let w = w.filter(|w| w.is_object())?;
745        Some(RateWindow {
746            used_percent: w.get("used_percent").and_then(Value::as_f64).unwrap_or(0.0),
747            window_minutes: w.get("window_minutes").and_then(Value::as_u64).unwrap_or(0),
748            resets_at: w
749                .get("resets_at")
750                .and_then(Value::as_i64)
751                .filter(|s| *s > 0)
752                .map(|s| std::time::UNIX_EPOCH + std::time::Duration::from_secs(s as u64)),
753        })
754    };
755    RateLimit {
756        primary: window(v.get("primary")),
757        secondary: window(v.get("secondary")),
758        plan: v.get("plan_type").and_then(Value::as_str).map(str::to_string),
759        reached: v.get("rate_limit_reached_type").map(|x| !x.is_null()).unwrap_or(false),
760    }
761}
762
763/// `call_id` on function calls, `id` on the shell-call variants.
764fn call_id(payload: &Value) -> String {
765    payload.get("call_id").or_else(|| payload.get("id")).and_then(Value::as_str).unwrap_or_default().to_string()
766}
767
768impl SessionTracker for CodexTranscript {
769    fn refresh(&mut self) -> anyhow::Result<bool> {
770        let (lines, more) = self.reader.read_new_lines(REFRESH_BUDGET_BYTES)?;
771        for l in &lines {
772            self.ingest(l);
773        }
774        Ok(more)
775    }
776
777    fn summary(&self) -> &SessionSummary {
778        &self.summary
779    }
780
781    fn path(&self) -> &Path {
782        self.reader.path()
783    }
784}
785
786#[cfg(test)]
787mod tests {
788    use super::*;
789    use std::io::Write;
790    use std::time::Duration;
791
792    fn record_response(t: &mut CodexTranscript, input: u64, cached: u64, output: u64) -> String {
793        let line = serde_json::json!({"type":"event_msg","payload":{"type":"token_count","info":{
794            "last_token_usage":{"input_tokens":input,"cached_input_tokens":cached,"output_tokens":output},
795            "total_token_usage":{
796                "input_tokens":t.summary.usage.prompt() + input,
797                "cached_input_tokens":t.summary.usage.cache_read + cached,
798                "output_tokens":t.summary.usage.output + output,
799            },
800        }}})
801        .to_string();
802        t.ingest(&line);
803        line
804    }
805
806    /// The bug this guards: the year and month directories were last touched
807    /// when a child directory was created, long before the rollout of
808    /// interest was written.
809    /// Write a rollout with an explicit modification time.
810    ///
811    /// Ordering must not be left to how finely the filesystem happens to
812    /// timestamp three writes microseconds apart: Linux gave all three the
813    /// same mtime, the stable sort preserved insertion order, and the test
814    /// failed there while passing on macOS.
815    fn rollout(dir: &Path, name: &str, written: SystemTime) -> PathBuf {
816        let p = dir.join(name);
817        std::fs::write(&p, b"x").unwrap();
818        let f = std::fs::File::options().write(true).open(&p).unwrap();
819        f.set_times(std::fs::FileTimes::new().set_accessed(written).set_modified(written)).unwrap();
820        p
821    }
822
823    const TIMEOUT: Duration = Duration::from_secs(15 * 60);
824
825    /// One app-server, several conversations. Every live one must get a row:
826    /// returning only the newest is what collapsed them into a single
827    /// mis-attributed row.
828    #[test]
829    fn every_live_codex_thread_is_returned_newest_first() {
830        let dir = std::env::temp_dir().join(format!("agent-top-threads-{}", std::process::id()));
831        let _ = std::fs::remove_dir_all(&dir);
832        std::fs::create_dir_all(&dir).unwrap();
833        let now = SystemTime::now();
834        let started = now - Duration::from_secs(600);
835
836        // Distinct write times, oldest first, so "newest first" has a single
837        // correct answer.
838        let a = rollout(&dir, "a.jsonl", now - Duration::from_secs(300));
839        let b = rollout(&dir, "b.jsonl", now - Duration::from_secs(200));
840        let c = rollout(&dir, "c.jsonl", now - Duration::from_secs(100));
841        let recent: Vec<(PathBuf, PathBuf, SystemTime)> =
842            [&a, &b, &c].iter().map(|p| ((*p).clone(), PathBuf::from("/Users/dev/code/one"), started)).collect();
843
844        // The app-server case: the process cwd matches no conversation.
845        let (paths, attribution) = attribute(Some(Path::new("/")), started, None, &recent, &HashSet::new(), now, TIMEOUT);
846        assert_eq!(paths.len(), 3, "all three conversations get a row");
847        assert_eq!(paths[0], c, "newest activity first");
848        assert_eq!(attribution, Attribution::CwdHeuristic, "still a heuristic, and still labelled one");
849
850        // A conversation already claimed by another process is not shown twice.
851        let taken: HashSet<PathBuf> = [c.clone()].into_iter().collect();
852        let (paths, _) = attribute(Some(Path::new("/")), started, None, &recent, &taken, now, TIMEOUT);
853        assert_eq!(paths.len(), 2);
854        assert!(!paths.contains(&c));
855
856        // A conversation nobody has written to for longer than the activity
857        // window has finished; it belongs in the stopped list, not on this
858        // process.
859        let stale = now + TIMEOUT + Duration::from_secs(60);
860        let (paths, attribution) = attribute(Some(Path::new("/")), started, None, &recent, &HashSet::new(), stale, TIMEOUT);
861        assert!(paths.is_empty());
862        assert_eq!(attribution, Attribution::None);
863
864        // The CLI case: one conversation, in the directory the process runs in.
865        let (paths, _) = attribute(Some(Path::new("/Users/dev/code/one")), started, None, &recent, &HashSet::new(), now, TIMEOUT);
866        assert_eq!(paths.len(), 3, "a cwd match takes every conversation in that directory");
867        assert_eq!(paths[0], c);
868
869        // A rollout that predates the process is not this process's.
870        let (paths, _) = attribute(Some(Path::new("/")), now + Duration::from_secs(3600), None, &recent, &HashSet::new(), now, TIMEOUT);
871        assert!(paths.is_empty());
872
873        let _ = std::fs::remove_dir_all(&dir);
874    }
875
876    /// Two app-servers at once, the VS Code one and a CLI-spawned one, both
877    /// running from `/`. Without the open-file signal the one asked first
878    /// took every live thread. The bug this guards was found live on
879    /// 2026-09-04: two threads of a fresh app-server were shown on the four
880    /// day old VS Code one.
881    #[test]
882    fn an_open_rollout_belongs_to_the_process_holding_it() {
883        let dir = std::env::temp_dir().join(format!("agent-top-held-{}", std::process::id()));
884        let _ = std::fs::remove_dir_all(&dir);
885        std::fs::create_dir_all(&dir).unwrap();
886        let now = SystemTime::now();
887        let started = now - Duration::from_secs(600);
888        let a = rollout(&dir, "a.jsonl", now - Duration::from_secs(200));
889        let b = rollout(&dir, "b.jsonl", now - Duration::from_secs(100));
890        let recent: Vec<(PathBuf, PathBuf, SystemTime)> =
891            [&a, &b].iter().map(|p| ((*p).clone(), PathBuf::from("/Users/dev/code/one"), started)).collect();
892
893        // The newer app-server holds both rollouts open. It started after the
894        // rollouts' recorded start, which the heuristic would reject; the open
895        // file settles it.
896        let held = vec![a.clone(), b.clone()];
897        let (paths, attribution) = attribute(Some(Path::new("/")), now, Some(&held), &recent, &HashSet::new(), now, TIMEOUT);
898        assert_eq!(paths, vec![b.clone(), a.clone()], "held rollouts, newest written first");
899        assert_eq!(attribution, Attribution::OpenFile);
900
901        // The older app-server holds nothing. Its fallback would have taken
902        // both live rollouts; with them marked taken it gets no row.
903        let taken: HashSet<PathBuf> = held.iter().cloned().collect();
904        let (paths, attribution) = attribute(Some(Path::new("/")), started, None, &recent, &taken, now, TIMEOUT);
905        assert!(paths.is_empty());
906        assert_eq!(attribution, Attribution::None);
907
908        let _ = std::fs::remove_dir_all(&dir);
909    }
910
911    #[test]
912    fn rollout_owner_follows_the_launcher_runtime_but_not_nested_agents() {
913        let proc = |pid, ppid, cmd: &[&str]| RawProc {
914            pid,
915            ppid,
916            name: cmd[0].into(),
917            exe: None,
918            cmd: cmd.iter().map(|s| (*s).into()).collect(),
919            cwd: None,
920            cpu_percent: 0.0,
921            rss_bytes: 0,
922            start_time: 0,
923            run_time: 1,
924        };
925        let procs = [
926            proc(10, None, &["node", "/usr/lib/node_modules/@openai/codex/bin/codex.js", "--yolo"]),
927            proc(11, Some(10), &["codex", "--yolo"]),
928            proc(12, Some(11), &["codex", "exec", "review"]),
929        ];
930        let (roots, _) = crate::process::build_forest(&procs);
931        let by_pid = procs.iter().map(|p| (p.pid, p)).collect();
932        let root = &roots[0];
933        assert_eq!(rollout_owner(root, &by_pid).pid, 11, "read the native runtime's descriptors, not the launcher's");
934        assert_eq!(rollout_owner(&root.children[0], &by_pid).pid, 11, "never claim a nested agent's rollouts");
935        assert_eq!(rollout_owner(&root.children[0].children[0], &by_pid).pid, 12);
936    }
937
938    #[test]
939    fn reads_explicit_subagent_lineage_and_identity() {
940        use serde_json::json;
941
942        let mut t = CodexTranscript::new("unused.jsonl");
943        t.ingest(
944            &json!({
945                "type": "session_meta",
946                "payload": {
947                    "id": "child", "session_id": "root", "cli_version": "0.154.0",
948                    "forked_from_id": "history-source",
949                    "source": {"subagent": {"thread_spawn": {
950                        "parent_thread_id": "parent", "depth": 2,
951                        "agent_nickname": "Scout", "agent_role": "explorer"
952                    }}}
953                }
954            })
955            .to_string(),
956        );
957        assert_eq!(t.summary.session_id.as_deref(), Some("child"), "thread id wins over root session_id");
958        assert_eq!(
959            t.summary.subagent,
960            Some(SubagentInfo { parent_session_id: "parent".into(), nickname: Some("Scout".into()), role: Some("explorer".into()) })
961        );
962
963        // New metadata also carries explicit lineage at the top level. A bare
964        // parent_thread_id without the subagent source must not be inferred.
965        let meta = json!({"id": "child", "thread_source": "subagent", "parent_thread_id": "parent", "agent_type": "worker"});
966        assert_eq!(
967            subagent_info(&meta),
968            Some(SubagentInfo { parent_session_id: "parent".into(), nickname: None, role: Some("worker".into()) })
969        );
970        let meta = json!({"id": "child", "agent_nickname": "Scout", "source": {"subagent": {"thread_spawn": {
971            "parent_thread_id": "parent", "agent_type": "worker"
972        }}}});
973        assert_eq!(subagent_info(&meta).unwrap().nickname.as_deref(), Some("Scout"));
974        assert_eq!(subagent_info(&meta).unwrap().role.as_deref(), Some("worker"));
975    }
976
977    #[test]
978    fn inherited_session_headers_do_not_replace_the_rollouts_identity() {
979        use serde_json::json;
980
981        let dir = std::env::temp_dir().join(format!("agent-top-codex-inherited-meta-{}", std::process::id()));
982        std::fs::create_dir_all(&dir).unwrap();
983        let path = dir.join("child.jsonl");
984        let mut file = std::fs::File::create(&path).unwrap();
985        let own = json!({"type": "session_meta", "payload": {
986            "id": "child", "session_id": "parent", "cwd": "/child", "cli_version": "0.154.0",
987            "source": {"subagent": {"thread_spawn": {
988                "parent_thread_id": "parent", "agent_nickname": "Scout", "agent_role": "explorer"
989            }}}
990        }});
991        writeln!(file, "{own}").unwrap();
992        let mut t = CodexTranscript::new(&path);
993        t.refresh().unwrap();
994
995        // Forked history follows the child's own header, possibly across refreshes.
996        for id in ["parent", "grandparent"] {
997            let inherited = json!({"type": "session_meta", "payload": {
998                "id": id, "cwd": "/ancestor", "cli_version": "0.149.0", "source": "cli"
999            }});
1000            writeln!(file, "{inherited}").unwrap();
1001            t.refresh().unwrap();
1002            assert_eq!(t.summary.session_id.as_deref(), Some("child"));
1003            assert_eq!(t.summary.cwd.as_deref(), Some(Path::new("/child")));
1004            assert_eq!(t.summary.harness_version.as_deref(), Some("0.154.0"));
1005            assert_eq!(
1006                t.summary.subagent,
1007                Some(SubagentInfo { parent_session_id: "parent".into(), nickname: Some("Scout".into()), role: Some("explorer".into()) })
1008            );
1009        }
1010        let usage = json!({"type": "event_msg", "payload": {"type": "token_count", "info": {
1011            "total_token_usage": {"input_tokens": 42}
1012        }}});
1013        writeln!(file, "{usage}").unwrap();
1014        t.refresh().unwrap();
1015        assert_eq!(t.summary.usage.input, 42, "continue reading events after inherited headers");
1016        std::fs::remove_dir_all(&dir).unwrap();
1017    }
1018
1019    #[test]
1020    fn never_infers_subagents_from_forks_or_incomplete_metadata() {
1021        use serde_json::json;
1022
1023        for meta in [
1024            json!({"id": "child", "source": "cli", "forked_from_id": "parent"}),
1025            json!({"id": "child", "parent_thread_id": "parent"}),
1026            json!({"id": "child", "source": {"subagent": "review"}}),
1027            json!({"id": "child", "source": {"subagent": {"thread_spawn": {"depth": 1}}}}),
1028            json!({"id": "child", "source": {"subagent": {"thread_spawn": {"parent_thread_id": ""}}}}),
1029            json!({"id": "child", "thread_source": "subagent", "parent_thread_id": "child"}),
1030            json!({"id": "child", "source": {"subagent": {"thread_spawn": {"parent_thread_id": 42}}}}),
1031        ] {
1032            assert_eq!(subagent_info(&meta), None, "{meta}");
1033        }
1034    }
1035
1036    #[test]
1037    fn the_rollout_id_follows_the_timestamp() {
1038        assert_eq!(
1039            rollout_id(Path::new("/x/2026/05/14/rollout-2026-05-14T21-37-50-01000000-0000-7000-0000-000000000000.jsonl")),
1040            "01000000-0000-7000-0000-000000000000"
1041        );
1042        assert_eq!(rollout_id(Path::new("/x/odd.jsonl")), "odd");
1043    }
1044
1045    #[test]
1046    fn finds_a_fresh_rollout_under_stale_directories() {
1047        let root = std::env::temp_dir().join(format!("agent-top-rollouts-{}", std::process::id()));
1048        let day = root.join("2026").join("09").join("04");
1049        std::fs::create_dir_all(&day).unwrap();
1050        let fresh = day.join("rollout-fresh.jsonl");
1051        let stale = day.join("rollout-stale.jsonl");
1052        std::fs::write(&fresh, "{}\n").unwrap();
1053        std::fs::write(&stale, "{}\n").unwrap();
1054        let now = SystemTime::now();
1055        let long_ago = now - Duration::from_secs(40 * 86_400);
1056        std::fs::File::open(&stale).unwrap().set_modified(long_ago).unwrap();
1057        for dir in [&root, &root.join("2026"), &root.join("2026").join("09"), &day] {
1058            std::fs::File::open(dir).unwrap().set_modified(long_ago).unwrap();
1059        }
1060        let found = rollouts_under(&root, now - Duration::from_secs(1800));
1061        assert_eq!(found, vec![fresh], "the fresh file is found through directories nobody has touched in weeks");
1062        std::fs::remove_dir_all(&root).unwrap();
1063    }
1064
1065    #[test]
1066    fn reads_the_latest_rate_limit_snapshot() {
1067        let dir = std::env::temp_dir().join(format!("agent-top-codex-rl-{}", std::process::id()));
1068        std::fs::create_dir_all(&dir).unwrap();
1069        let path = dir.join("rollout.jsonl");
1070        let mut f = std::fs::File::create(&path).unwrap();
1071        writeln!(f, r#"{{"timestamp":"2026-06-16T20:45:04.000Z","type":"event_msg","payload":{{"type":"token_count","info":{{"total_token_usage":{{"input_tokens":10,"output_tokens":1,"total_tokens":11}}}},"rate_limits":{{"limit_id":"codex","primary":{{"used_percent":1.0,"window_minutes":300,"resets_at":1781660699}},"secondary":{{"used_percent":27.0,"window_minutes":10080,"resets_at":1782080576}},"plan_type":"plus","rate_limit_reached_type":null}}}}}}"#).unwrap();
1072        // A later snapshot with higher usage; the latest wins.
1073        writeln!(f, r#"{{"timestamp":"2026-06-16T20:50:00.000Z","type":"event_msg","payload":{{"type":"token_count","info":{{"total_token_usage":{{"input_tokens":20,"output_tokens":2,"total_tokens":22}}}},"rate_limits":{{"limit_id":"codex","primary":{{"used_percent":42.0,"window_minutes":300,"resets_at":1781660999}},"secondary":{{"used_percent":28.0,"window_minutes":10080,"resets_at":1782080576}},"plan_type":"plus","rate_limit_reached_type":"primary"}}}}}}"#).unwrap();
1074        let mut t = CodexTranscript::new(&path);
1075        t.refresh().unwrap();
1076        let rl = t.summary().rate_limit.as_ref().expect("rate limit parsed");
1077        assert_eq!(rl.plan.as_deref(), Some("plus"));
1078        assert!(rl.reached, "the latest snapshot reports the primary window hit");
1079        let p = rl.primary.expect("primary window");
1080        assert_eq!(p.used_percent, 42.0, "the latest value, not the first");
1081        assert_eq!(p.window_minutes, 300);
1082        assert_eq!(rl.secondary.unwrap().used_percent, 28.0);
1083        assert_eq!(rl.tightest().map(|w| w.used_percent), Some(42.0));
1084        let _ = std::fs::remove_dir_all(&dir);
1085    }
1086
1087    #[test]
1088    fn a_snapshot_with_no_windows_keeps_the_last_one() {
1089        let dir = std::env::temp_dir().join(format!("agent-top-codex-rl-premium-{}", std::process::id()));
1090        std::fs::create_dir_all(&dir).unwrap();
1091        let path = dir.join("rollout.jsonl");
1092        let mut f = std::fs::File::create(&path).unwrap();
1093        writeln!(f, r#"{{"timestamp":"2026-09-24T10:00:00.000Z","type":"event_msg","payload":{{"type":"token_count","info":null,"rate_limits":{{"limit_id":"codex","primary":{{"used_percent":12.0,"window_minutes":300,"resets_at":1790000000}},"secondary":{{"used_percent":40.0,"window_minutes":10080,"resets_at":1790500000}},"plan_type":"plus","rate_limit_reached_type":null}}}}}}"#).unwrap();
1094        writeln!(f, r#"{{"timestamp":"2026-09-24T10:00:05.000Z","type":"event_msg","payload":{{"type":"token_count","info":null,"rate_limits":{{"limit_id":"premium","limit_name":null,"primary":null,"secondary":null,"credits":{{"has_credits":false,"unlimited":false,"balance":"0"}},"plan_type":"plus","rate_limit_reached_type":null}}}}}}"#).unwrap();
1095        drop(f);
1096        let mut t = CodexTranscript::new(&path);
1097        t.refresh().unwrap();
1098        let rl = t.summary().rate_limit.as_ref().expect("the codex snapshot survives");
1099        assert_eq!(rl.primary.map(|w| w.used_percent), Some(12.0));
1100        assert_eq!(rl.secondary.map(|w| w.used_percent), Some(40.0));
1101        let _ = std::fs::remove_dir_all(&dir);
1102    }
1103
1104    #[test]
1105    fn counts_mcp_calls_per_server_from_the_end_event() {
1106        let dir = std::env::temp_dir().join(format!("agent-top-codex-mcp-{}", std::process::id()));
1107        std::fs::create_dir_all(&dir).unwrap();
1108        let path = dir.join("rollout.jsonl");
1109        let mut f = std::fs::File::create(&path).unwrap();
1110        // Two MCP servers, one call failing. The matching response_item
1111        // function_call/output pair is what makes the tool-call count and span;
1112        // the mcp_tool_call_end is the only line naming the server.
1113        writeln!(f, r#"{{"timestamp":"2026-05-27T09:00:01.000Z","type":"response_item","payload":{{"type":"function_call","call_id":"c1","name":"github_fetch_file"}}}}"#).unwrap();
1114        writeln!(f, r#"{{"timestamp":"2026-05-27T09:00:02.000Z","type":"response_item","payload":{{"type":"function_call_output","call_id":"c1"}}}}"#).unwrap();
1115        writeln!(f, r#"{{"timestamp":"2026-05-27T09:00:02.100Z","type":"event_msg","payload":{{"type":"mcp_tool_call_end","call_id":"c1","invocation":{{"server":"codex_apps","tool":"github_fetch_file"}},"duration":{{"secs":1,"nanos":0}},"result":{{"Ok":{{}}}}}}}}"#).unwrap();
1116        writeln!(f, r#"{{"timestamp":"2026-05-27T09:00:05.000Z","type":"event_msg","payload":{{"type":"mcp_tool_call_end","call_id":"c2","invocation":{{"server":"codex_apps","tool":"github_search"}},"duration":{{"secs":0,"nanos":0}},"result":{{"Err":"boom"}}}}}}"#).unwrap();
1117        writeln!(f, r#"{{"timestamp":"2026-05-27T09:00:07.000Z","type":"event_msg","payload":{{"type":"mcp_tool_call_end","call_id":"c3","invocation":{{"server":"node_repl","tool":"js"}},"duration":{{"secs":0,"nanos":0}},"result":{{"Ok":{{}}}}}}}}"#).unwrap();
1118        let mut t = CodexTranscript::new(&path);
1119        t.refresh().unwrap();
1120        let s = t.summary();
1121        assert_eq!(s.mcp.len(), 2);
1122        let apps = &s.mcp["codex_apps"];
1123        assert_eq!((apps.calls, apps.errors), (2, 1));
1124        assert_eq!(apps.last_call, parse_rfc3339_utc("2026-05-27T09:00:05.000Z"));
1125        assert_eq!(s.mcp["node_repl"].calls, 1);
1126        // The one call with a response_item pair is one tool call and one span;
1127        // the mcp_tool_call_end lines do not add to that.
1128        assert_eq!(s.tool_calls, 1, "mcp_tool_call_end must not double-count tool calls");
1129        let _ = std::fs::remove_dir_all(&dir);
1130    }
1131
1132    #[test]
1133    fn sizes_context_per_tool_from_last_token_usage_and_skips_the_repeated_snapshot() {
1134        for delayed in [false, true] {
1135            let mut t = CodexTranscript::new("unused.jsonl").with_prices(pricing::builtin_table());
1136            t.ingest(r#"{"timestamp":"2026-09-18T09:00:01Z","type":"response_item","payload":{"type":"function_call","call_id":"c1","name":"exec_command"}}"#);
1137            t.ingest(r#"{"timestamp":"2026-09-18T09:00:01Z","type":"response_item","payload":{"type":"function_call","call_id":"c2","name":"github_fetch_file"}}"#);
1138            if !delayed {
1139                record_response(&mut t, 10_000, 0, 100);
1140            }
1141            t.ingest(
1142                r#"{"timestamp":"2026-09-18T09:00:02Z","type":"response_item","payload":{"type":"function_call_output","call_id":"c1"}}"#,
1143            );
1144            t.ingest(
1145                r#"{"timestamp":"2026-09-18T09:00:02Z","type":"response_item","payload":{"type":"function_call_output","call_id":"c2"}}"#,
1146            );
1147            t.ingest(r#"{"type":"event_msg","payload":{"type":"mcp_tool_call_end","call_id":"c2","invocation":{"server":"codex_apps","tool":"github_fetch_file"},"result":{"Ok":{}}}}"#);
1148            if delayed {
1149                record_response(&mut t, 10_000, 0, 100);
1150            }
1151            let initial = t.summary.context.sources();
1152            assert_eq!(initial.len(), 1);
1153            assert_eq!((initial[0].name.as_str(), initial[0].tokens), ("other", 10_000));
1154
1155            t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1156            let snapshot = record_response(&mut t, 13_100, 12_000, 40);
1157            t.ingest(r#"{"type":"event_msg","payload":{"type":"task_started"}}"#);
1158            t.ingest(&snapshot);
1159            let c: HashMap<_, _> = t.summary.context.sources().into_iter().map(|c| (c.name.clone(), c)).collect();
1160            assert_eq!((c["exec_command"].calls, c["exec_command"].tokens), (1, 1_500));
1161            assert_eq!((c["codex_apps"].calls, c["codex_apps"].tokens, c["codex_apps"].origin), (1, 1_500, ContextOrigin::Mcp));
1162            assert!(!c.contains_key("github_fetch_file"));
1163            assert_eq!(c["other"].tokens, 10_100);
1164            assert_eq!(c["other"].cost_usd, 0.0);
1165        }
1166    }
1167
1168    #[test]
1169    fn delayed_usage_charges_only_results_consumed_by_that_response() {
1170        let mut t = CodexTranscript::new("unused.jsonl").with_prices(pricing::builtin_table());
1171        t.ingest(r#"{"type":"turn_context","payload":{"model":"gpt-5.4-mini"}}"#);
1172        t.ingest(r#"{"timestamp":"2026-09-18T09:00:01Z","type":"response_item","payload":{"type":"function_call","call_id":"c1","name":"exec_command"}}"#);
1173        t.ingest(r#"{"timestamp":"2026-09-18T09:00:02Z","type":"response_item","payload":{"type":"function_call_output","call_id":"c1"}}"#);
1174        // A later item in the same streamed response must not consume c1.
1175        t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1176        let snapshot = record_response(&mut t, 1_000, 200, 10);
1177        let initial = t.summary.context.sources();
1178        assert_eq!(initial.len(), 1);
1179        assert_eq!((initial[0].name.as_str(), initial[0].tokens), ("other", 1_000));
1180
1181        t.ingest(r#"{"type":"response_item","payload":{"type":"reasoning"}}"#);
1182        t.ingest(r#"{"timestamp":"2026-09-18T09:00:03Z","type":"response_item","payload":{"type":"custom_tool_call","call_id":"c2","name":"apply_patch"}}"#);
1183        t.ingest(
1184            r#"{"timestamp":"2026-09-18T09:00:04Z","type":"response_item","payload":{"type":"custom_tool_call_output","call_id":"c2"}}"#,
1185        );
1186        for snapshot in [
1187            snapshot.as_str(),
1188            r#"{"type":"event_msg","payload":{"type":"token_count","info":null}}"#,
1189            r#"{"type":"event_msg","payload":{"type":"token_count","info":{"last_token_usage":null}}}"#,
1190            r#"{"type":"event_msg","payload":{"type":"token_count","info":{"last_token_usage":{}}}}"#,
1191        ] {
1192            t.ingest(snapshot);
1193        }
1194        t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1195        assert_eq!(t.summary.context.sources(), initial);
1196        record_response(&mut t, 1_300, 1_100, 20);
1197        let sources = t.summary.context.sources();
1198        let command = sources.iter().find(|s| s.name == "exec_command").unwrap();
1199        assert_eq!((command.calls, command.tokens), (1, 290));
1200        assert!(!sources.iter().any(|s| s.name == "apply_patch"));
1201
1202        t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1203        record_response(&mut t, 1_350, 1_200, 5);
1204        let sources = t.summary.context.sources();
1205        let patch = sources.iter().find(|s| s.name == "apply_patch").unwrap();
1206        assert_eq!((patch.calls, patch.tokens), (1, 30));
1207        assert_eq!(sources.iter().map(|s| s.tokens).sum::<u64>(), 1_350);
1208        let prompt_cost = t.summary.cost_breakdown.input + t.summary.cost_breakdown.cache_read;
1209        assert!(prompt_cost > 0.0);
1210        assert!((sources.iter().map(|s| s.cost_usd).sum::<f64>() - prompt_cost).abs() < 1e-9);
1211        assert_eq!((t.summary.usage.prompt(), t.summary.usage.output, t.summary.tool_calls), (3_650, 35, 2));
1212    }
1213
1214    #[test]
1215    fn tool_search_errors_are_not_part_of_the_requesting_responses_prompt() {
1216        let mut t = CodexTranscript::new("unused.jsonl");
1217        t.ingest(r#"{"type":"response_item","payload":{"type":"tool_search_call","call_id":"search"}}"#);
1218        t.ingest(
1219            r#"{"timestamp":"2026-09-18T09:00:01Z","type":"response_item","payload":{"type":"function_call_output","call_id":"search"}}"#,
1220        );
1221        t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1222        record_response(&mut t, 1_000, 0, 10);
1223        let sources = t.summary.context.sources();
1224        assert_eq!(sources.len(), 1);
1225        assert_eq!((sources[0].name.as_str(), sources[0].tokens), ("other", 1_000));
1226
1227        t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1228        record_response(&mut t, 1_060, 0, 10);
1229        let sources = t.summary.context.sources();
1230        let tool = sources.iter().find(|s| s.origin == ContextOrigin::Tool).unwrap();
1231        assert_eq!((tool.calls, tool.tokens), (1, 50));
1232    }
1233
1234    #[test]
1235    fn identical_response_usage_is_not_a_snapshot_when_cumulative_usage_advances() {
1236        let mut t = CodexTranscript::new("unused.jsonl").with_prices(pricing::builtin_table());
1237        t.ingest(r#"{"type":"turn_context","payload":{"model":"gpt-5.4-mini"}}"#);
1238        t.ingest(r#"{"timestamp":"2026-09-18T09:00:01Z","type":"response_item","payload":{"type":"function_call","call_id":"c1","name":"exec_command"}}"#);
1239        t.ingest(r#"{"timestamp":"2026-09-18T09:00:02Z","type":"response_item","payload":{"type":"function_call_output","call_id":"c1"}}"#);
1240        let snapshot = record_response(&mut t, 1_000, 0, 10);
1241        t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1242        t.ingest(&snapshot);
1243        assert_eq!(t.summary.context.sources().len(), 1);
1244        record_response(&mut t, 1_000, 0, 10);
1245        let sources = t.summary.context.sources();
1246        let command = sources.iter().find(|s| s.name == "exec_command").unwrap();
1247        assert_eq!((command.calls, command.tokens), (1, 0));
1248        assert_eq!(sources.iter().map(|s| s.tokens).sum::<u64>(), 1_000);
1249        assert!((sources.iter().map(|s| s.cost_usd).sum::<f64>() - t.summary.cost_breakdown.input).abs() < 1e-9);
1250    }
1251
1252    #[test]
1253    fn interrupted_responses_keep_results_for_the_next_consuming_response() {
1254        for ending in ["turn_aborted", "error"] {
1255            let mut t = CodexTranscript::new("unused.jsonl");
1256            let snapshot = record_response(&mut t, 1_000, 0, 10);
1257            t.ingest(r#"{"timestamp":"2026-09-18T09:00:01Z","type":"response_item","payload":{"type":"function_call","call_id":"c1","name":"exec_command"}}"#);
1258            t.ingest(
1259                r#"{"timestamp":"2026-09-18T09:00:02Z","type":"response_item","payload":{"type":"function_call_output","call_id":"c1"}}"#,
1260            );
1261            t.ingest(&serde_json::json!({"type":"event_msg","payload":{"type":ending}}).to_string());
1262            t.ingest(r#"{"type":"event_msg","payload":{"type":"task_started"}}"#);
1263            t.ingest(&snapshot);
1264            t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1265            record_response(&mut t, 1_110, 0, 10);
1266            let sources = t.summary.context.sources();
1267            let command = sources.iter().find(|s| s.name == "exec_command").unwrap();
1268            assert_eq!((command.calls, command.tokens), (1, 100));
1269        }
1270    }
1271
1272    #[test]
1273    fn reads_cumulative_usage_and_state() {
1274        let dir = std::env::temp_dir().join(format!("agent-top-codex-{}", std::process::id()));
1275        std::fs::create_dir_all(&dir).unwrap();
1276        let path = dir.join("rollout.jsonl");
1277        let mut f = std::fs::File::create(&path).unwrap();
1278        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();
1279        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:21.000Z","type":"turn_context","payload":{{"model":"gpt-5-codex"}}}}"#).unwrap();
1280        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:22.000Z","type":"event_msg","payload":{{"type":"task_started"}}}}"#).unwrap();
1281        writeln!(
1282            f,
1283            r#"{{"timestamp":"2026-08-28T08:53:23.000Z","type":"response_item","payload":{{"type":"function_call","call_id":"call_1","name":"shell"}}}}"#
1284        )
1285        .unwrap();
1286        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();
1287        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:25.000Z","type":"event_msg","payload":{{"type":"token_count","info":null}}}}"#)
1288            .unwrap();
1289        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:25.500Z","type":"response_item","payload":{{"type":"web_search_call","status":"completed"}}}}"#).unwrap();
1290        let mut t = CodexTranscript::new(&path).with_prices(pricing::builtin_table());
1291        t.refresh().unwrap();
1292        let s = t.summary();
1293        assert_eq!(s.session_id.as_deref(), Some("01a0"));
1294        assert_eq!(s.model.as_deref(), Some("gpt-5-codex"));
1295        assert_eq!(s.usage.input, 14778 - 12672);
1296        assert_eq!(s.usage.cache_read, 12672);
1297        assert_eq!(s.usage.total(), 15019);
1298        // gpt-5-codex has no entry of its own; it resolves to gpt-5 by the
1299        // longest-prefix rule (input 1.25, cache read 0.125, output 10).
1300        assert_eq!(s.unpriced_tokens, 0, "priced now that OpenAI's rows are in the table");
1301        assert!((s.cost_usd - (2106.0 * 1.25 + 12672.0 * 0.125 + 241.0 * 10.0) / 1_000_000.0).abs() < 1e-9, "{}", s.cost_usd);
1302        assert_eq!(s.tool_calls, 1);
1303        assert_eq!(s.activity, Activity::Working);
1304        assert_eq!(read_meta(&path).unwrap().0, PathBuf::from("/tmp/p"));
1305        assert_eq!(s.web_searches, 1);
1306        let tools: Vec<_> = s.spans.iter().filter(|sp| sp.kind == SpanKind::Tool).collect();
1307        assert_eq!(tools.len(), 1);
1308        assert_eq!(tools[0].name, "shell");
1309        assert!(tools[0].is_open(), "no output item yet");
1310        let turns: Vec<_> = s.spans.iter().filter(|sp| sp.kind == SpanKind::Turn).collect();
1311        assert_eq!(turns.len(), 1);
1312        assert!(turns[0].is_open(), "task_started with no task_complete");
1313        let _ = std::fs::remove_dir_all(&dir);
1314    }
1315
1316    #[test]
1317    fn pairs_calls_with_their_outputs() {
1318        let dir = std::env::temp_dir().join(format!("agent-top-codex-spans-{}", std::process::id()));
1319        std::fs::create_dir_all(&dir).unwrap();
1320        let path = dir.join("rollout.jsonl");
1321        let mut f = std::fs::File::create(&path).unwrap();
1322        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();
1323        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();
1324        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();
1325        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();
1326        // The model answers the outputs 1.5 s after the last one, then the turn completes.
1327        writeln!(
1328            f,
1329            r#"{{"timestamp":"2026-08-28T08:53:27.600Z","type":"response_item","payload":{{"type":"message","role":"assistant"}}}}"#
1330        )
1331        .unwrap();
1332        writeln!(f, r#"{{"timestamp":"2026-08-28T08:53:27.700Z","type":"event_msg","payload":{{"type":"task_complete"}}}}"#).unwrap();
1333        let mut t = CodexTranscript::new(&path);
1334        t.refresh().unwrap();
1335        let all = t.summary().spans.to_vec();
1336        let spans: Vec<_> = all.iter().filter(|sp| sp.kind == SpanKind::Tool).collect();
1337        assert_eq!(spans.len(), 2);
1338        assert_eq!(spans[0].name, "exec_command");
1339        assert_eq!(spans[0].duration_ms, Some(1_000));
1340        assert_eq!(spans[1].name, "apply_patch");
1341        assert_eq!(spans[1].duration_ms, Some(3_000));
1342        assert_eq!(t.summary().tool_calls, 2);
1343        // One inference: opened by the first output at :24, not re-opened by the
1344        // second at :26.1, ended by the assistant message at :27.6.
1345        let inf: Vec<_> = all.iter().filter(|sp| sp.kind == SpanKind::Inference).collect();
1346        assert_eq!(inf.len(), 1);
1347        assert_eq!(inf[0].duration_ms, Some(3_600));
1348        assert!(all.iter().all(|sp| sp.kind != SpanKind::Turn), "no task_started in this file");
1349        let _ = std::fs::remove_dir_all(&dir);
1350    }
1351
1352    #[test]
1353    fn code_mode_records_nested_commands_and_patches_alongside_exec() {
1354        // Sanitised 0.154.0 metadata; nested ids are independent of wrapper ids.
1355        let path = std::env::temp_dir().join(format!("agent-top-codex-nested-{}.jsonl", std::process::id()));
1356        let mut f = std::fs::File::create(&path).unwrap();
1357        let mut t = CodexTranscript::new(&path);
1358        for line in [
1359            r#"{"timestamp":"2026-09-18T07:22:23.296Z","type":"response_item","payload":{"type":"custom_tool_call","name":"exec","call_id":"call_wrapper_1"}}"#,
1360            r#"{"timestamp":"2026-09-18T07:22:23.334Z","type":"event_msg","payload":{"type":"item_completed","started_at_ms":1789716143334,"completed_at_ms":1789716143334,"item":{"type":"CommandExecution","id":"exec-command-1","source":"unified_exec_startup","status":"completed","exit_code":0}}}"#,
1361            r#"{"timestamp":"2026-09-18T07:22:23.348Z","type":"response_item","payload":{"type":"custom_tool_call_output","call_id":"call_wrapper_1"}}"#,
1362            r#"{"timestamp":"2026-09-18T07:22:30.281Z","type":"response_item","payload":{"type":"custom_tool_call","name":"exec","call_id":"call_wrapper_2"}}"#,
1363            r#"{"timestamp":"2026-09-18T07:22:30.288Z","type":"event_msg","payload":{"type":"item_completed","started_at_ms":1789716150288,"completed_at_ms":1789716150288,"item":{"type":"FileChange","id":"exec-patch-1","status":"completed"}}}"#,
1364            r#"{"timestamp":"2026-09-18T07:22:30.351Z","type":"response_item","payload":{"type":"custom_tool_call_output","call_id":"call_wrapper_2"}}"#,
1365        ] {
1366            writeln!(f, "{line}").unwrap();
1367            t.refresh().unwrap();
1368        }
1369        let _ = std::fs::remove_file(&path);
1370        let tools: Vec<_> = t.summary.spans.iter().filter(|s| s.kind == SpanKind::Tool).collect();
1371        assert_eq!(t.summary.tool_calls, 4, "two wrapper invocations and two actual nested tool invocations");
1372        assert_eq!(
1373            tools.iter().map(|s| (s.name.as_str(), s.duration_ms)).collect::<Vec<_>>(),
1374            [("exec", Some(52)), ("exec_command", Some(0)), ("exec", Some(70)), ("apply_patch", Some(0)),]
1375        );
1376        assert!(tools.iter().all(|s| !s.error));
1377        assert!(t.pending_tools.is_empty(), "completion items do not leave unmatched response calls");
1378        // Only wrapper outputs contribute to context accounting.
1379        record_response(&mut t, 1_000, 0, 0);
1380        t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1381        let usage = TokenUsage { input: 1_100, ..Default::default() };
1382        let cost = crate::model::CostBreakdown { input: 1.0, ..Default::default() };
1383        t.summary.context.response(&usage, &cost);
1384        let sources = t.summary.context.sources();
1385        for name in ["exec_command", "apply_patch"] {
1386            let source = sources.iter().find(|s| s.name == name).unwrap();
1387            assert_eq!((source.calls, source.tokens), (1, 50));
1388        }
1389        assert!(!sources.iter().any(|s| s.name == "exec"));
1390
1391        let mut baseline = crate::harness::ContextLedger::default();
1392        baseline.response(&TokenUsage { input: 1_000, ..Default::default() }, &Default::default());
1393        baseline.result("call_wrapper_1", ContextOrigin::Tool, "exec");
1394        baseline.result("call_wrapper_2", ContextOrigin::Tool, "exec");
1395        baseline.response(&usage, &cost);
1396        for input in [1_150, 1_180] {
1397            let usage = TokenUsage { input, ..Default::default() };
1398            t.summary.context.response(&usage, &cost);
1399            baseline.response(&usage, &cost);
1400            let totals = |ledger: &crate::harness::ContextLedger| {
1401                ledger
1402                    .sources()
1403                    .iter()
1404                    .fold((0, 0, 0.0), |(calls, tokens, cost), s| (calls + s.calls, tokens + s.tokens, cost + s.cost_usd))
1405            };
1406            let (calls, tokens, cost) = totals(&t.summary.context);
1407            let (old_calls, old_tokens, old_cost) = totals(&baseline);
1408            assert_eq!((calls, tokens), (old_calls, old_tokens));
1409            assert!((cost - old_cost).abs() < 1e-9);
1410        }
1411    }
1412
1413    #[test]
1414    fn nested_output_sizes_use_decoded_text_not_inputs_or_duplicate_fields() {
1415        use serde_json::json;
1416        for (item, expected) in [
1417            (
1418                json!({"type":"CommandExecution","formatted_output":"é\n", "aggregated_output":"longer output", "stdout":"duplicate"}),
1419                Some(3),
1420            ),
1421            (json!({"type":"CommandExecution","aggregated_output":"abcd","stdout":"duplicate","stderr":"duplicate"}), Some(4)),
1422            (json!({"type":"CommandExecution","stdout":"ab","stderr":"c"}), Some(3)),
1423            (json!({"type":"CommandExecution","command":["must not size this"]}), None),
1424            (json!({"type":"FileChange","stdout":"ok\n","stderr":"!","changes":{"file":"not output"}}), Some(4)),
1425            (json!({"type":"FileChange","stdout":"","stderr":""}), Some(0)),
1426            (json!({"type":"FileChange","changes":{"file":"not output"}}), None),
1427            (json!({"type":"McpToolCall","result":{"content":[{"type":"text","text":"abc"},{"type":"text","text":"d"}]}}), Some(4)),
1428            (json!({"type":"McpToolCall","result":{"content":[{"type":"image","data":"not text"}]}}), None),
1429            (json!({"type":"McpToolCall","result":{"content":[],"structuredContent":{"large":"result"}}}), None),
1430            (json!({"type":"McpToolCall","error":{"message":"failed"}}), Some(6)),
1431            (json!({"type":"DynamicToolCall","content_items":[{"type":"inputText","text":"abc"}],"arguments":"not output"}), Some(3)),
1432            (json!({"type":"DynamicToolCall","content_items":[{"type":"inputImage","imageUrl":"not text"}]}), None),
1433            (json!({"type":"DynamicToolCall","content_items":[{"type":"inputText","text":42}]}), None),
1434            (json!({"type":"DynamicToolCall","error":"failed"}), Some(6)),
1435        ] {
1436            assert_eq!(nested_output_bytes(&item), expected, "{item}");
1437        }
1438    }
1439
1440    #[test]
1441    fn code_mode_context_weights_children_without_changing_wrapper_totals() {
1442        use serde_json::json;
1443        for wrapper in ["exec", "wait"] {
1444            for (command, patch, expected) in [
1445                (Some("abc"), Some("d"), (30, 10)),
1446                (Some("abc"), None, (20, 20)),
1447                (None, Some("d"), (20, 20)),
1448                (Some(""), Some(""), (20, 20)),
1449                (Some(""), Some("d"), (0, 40)),
1450            ] {
1451                let mut t = CodexTranscript::new("unused.jsonl").with_prices(pricing::builtin_table());
1452                t.ingest(r#"{"type":"turn_context","payload":{"model":"gpt-5.4-mini"}}"#);
1453                t.ingest(
1454                    &json!({"timestamp":"1970-01-01T00:00:01Z","type":"response_item",
1455                    "payload":{"type":"custom_tool_call","name":wrapper,"call_id":"wrapper"}})
1456                    .to_string(),
1457                );
1458                for (item, start, end) in [
1459                    (
1460                        json!({"type":"CommandExecution","id":"exec-command","source":"unified_exec_startup","formatted_output":command}),
1461                        1100,
1462                        1200,
1463                    ),
1464                    (json!({"type":"FileChange","id":"exec-patch","stdout":patch}), 1300, 1400),
1465                ] {
1466                    let line = json!({"type":"event_msg","payload":{"type":"item_completed","item":item,
1467                        "started_at_ms":start,"completed_at_ms":end}})
1468                    .to_string();
1469                    t.ingest(&line);
1470                    t.ingest(&line);
1471                }
1472                t.ingest(r#"{"timestamp":"1970-01-01T00:00:02Z","type":"response_item","payload":{"type":"custom_tool_call_output","call_id":"wrapper"}}"#);
1473                record_response(&mut t, 1_000, 0, 10);
1474                assert_eq!(t.summary.context.sources().len(), 1);
1475                t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1476                let snapshot = record_response(&mut t, 1_050, 1_000, 10);
1477                t.ingest(&snapshot);
1478                let sources = t.summary.context.sources();
1479                let command = sources.iter().find(|s| s.name == "exec_command").unwrap();
1480                let patch = sources.iter().find(|s| s.name == "apply_patch").unwrap();
1481                assert_eq!((command.tokens, patch.tokens), expected);
1482                assert_eq!((command.calls, patch.calls, t.summary.tool_calls), (1, 1, 3));
1483                assert!(!sources.iter().any(|s| s.name == wrapper));
1484                assert_eq!(sources.iter().map(|s| s.tokens).sum::<u64>(), 1_050);
1485                let prompt_cost = t.summary.cost_breakdown.input + t.summary.cost_breakdown.cache_read;
1486                assert!((sources.iter().map(|s| s.cost_usd).sum::<f64>() - prompt_cost).abs() < 1e-9);
1487            }
1488        }
1489    }
1490
1491    #[test]
1492    fn code_mode_context_weights_mcp_and_dynamic_results_and_preserves_errors() {
1493        use serde_json::json;
1494        let mut t = CodexTranscript::new("unused.jsonl");
1495        t.ingest(r#"{"timestamp":"1970-01-01T00:00:01Z","type":"response_item","payload":{"type":"custom_tool_call","name":"exec","call_id":"wrapper"}}"#);
1496        for item in [
1497            json!({"id":"exec-mcp","type":"McpToolCall","tool":"lookup","server":"docs","result":{"content":[{"type":"text","text":"abc"}]}}),
1498            json!({"id":"exec-dynamic","type":"DynamicToolCall","tool":"check","success":false,"content_items":[{"type":"inputText","text":"d"}]}),
1499        ] {
1500            t.ingest(
1501                &json!({"type":"event_msg","payload":{"type":"item_completed","started_at_ms":1100,"completed_at_ms":1500,"item":item}})
1502                    .to_string(),
1503            );
1504        }
1505        t.ingest(r#"{"timestamp":"1970-01-01T00:00:02Z","type":"response_item","payload":{"type":"custom_tool_call_output","call_id":"wrapper"}}"#);
1506        record_response(&mut t, 1_000, 0, 0);
1507        t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1508        record_response(&mut t, 1_040, 0, 0);
1509        let sources = t.summary.context.sources();
1510        let mcp = sources.iter().find(|s| s.name == "docs").unwrap();
1511        let dynamic = sources.iter().find(|s| s.name == "check").unwrap();
1512        assert_eq!((mcp.origin, mcp.calls, mcp.tokens), (ContextOrigin::Mcp, 1, 30));
1513        assert_eq!((dynamic.origin, dynamic.calls, dynamic.tokens), (ContextOrigin::Tool, 1, 10));
1514        assert!(t.summary.spans.iter().any(|s| s.name == "check" && s.error));
1515    }
1516
1517    #[test]
1518    fn code_mode_context_keeps_ambiguous_wrappers_and_survives_span_eviction() {
1519        use serde_json::json;
1520        let child = |id, kind, start, end| {
1521            json!({
1522                "item":{"id":id,"type":kind}, "started_at_ms":start, "completed_at_ms":end,
1523            })
1524        };
1525        let patch = child("exec-patch", "FileChange", Some(1250), Some(1500));
1526        let unknown = child("exec-unknown", "FutureTool", Some(1250), Some(1500));
1527        for (children, expected, calls) in [
1528            (vec![patch.clone()], "apply_patch", 1),
1529            (vec![patch.clone(), patch.clone()], "apply_patch", 1),
1530            (vec![patch.clone(), child("exec-patch-2", "FileChange", Some(1500), Some(1750))], "apply_patch", 2),
1531            (vec![patch.clone(), unknown.clone()], "exec", 1),
1532            (vec![unknown, patch], "exec", 1),
1533            (vec![child("exec-patch", "FileChange", Some(500), Some(1500))], "exec", 1),
1534            (vec![child("exec-patch", "FileChange", Some(1250), Some(2500))], "exec", 1),
1535            (vec![child("exec-patch", "FileChange", None, Some(1500))], "exec", 1),
1536            (vec![child("exec-patch", "FileChange", Some(1250), None)], "exec", 1),
1537            (vec![], "exec", 1),
1538        ] {
1539            let mut t = CodexTranscript::new("unused.jsonl");
1540            t.ingest(r#"{"timestamp":"1970-01-01T00:00:01.000Z","type":"response_item","payload":{"type":"custom_tool_call","name":"exec","call_id":"wrapper"}}"#);
1541            for completion in children {
1542                t.nested_tool_completed(&completion);
1543            }
1544            for i in 0..t.summary.spans.cap() {
1545                t.summary.spans.open(format!("filler-{i}"), "tool".into(), SystemTime::UNIX_EPOCH, false);
1546            }
1547            if t.pending_tools["wrapper"].nested.iter().any(|n| n.id == "exec-patch") {
1548                let before = t.summary.tool_calls;
1549                t.nested_tool_completed(&child("exec-patch", "FileChange", Some(1250), Some(1500)));
1550                assert_eq!(t.summary.tool_calls, before);
1551            }
1552            t.ingest(r#"{"timestamp":"1970-01-01T00:00:02.000Z","type":"response_item","payload":{"type":"custom_tool_call_output","call_id":"wrapper"}}"#);
1553            record_response(&mut t, 1_000, 0, 0);
1554            t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1555            record_response(&mut t, 1_100, 0, 0);
1556            let sources: Vec<_> = t.summary.context.sources().into_iter().filter(|s| s.origin != ContextOrigin::Other).collect();
1557            assert_eq!(sources.len(), 1);
1558            assert_eq!((sources[0].name.as_str(), sources[0].calls, sources[0].tokens), (expected, calls, 100));
1559        }
1560    }
1561
1562    #[test]
1563    fn code_mode_context_rejects_overlapping_exec_and_wait_but_not_wait_agent() {
1564        use serde_json::json;
1565        for other in ["exec", "wait", "wait_agent"] {
1566            for exec_first in [true, false] {
1567                let mut t = CodexTranscript::new("unused.jsonl");
1568                for (id, name) in [("wrapper", "exec"), ("other", other)] {
1569                    t.ingest(
1570                        &json!({"timestamp":"1970-01-01T00:00:01.000Z", "type":"response_item",
1571                        "payload":{"type":"function_call","name":name,"call_id":id}})
1572                        .to_string(),
1573                    );
1574                }
1575                t.nested_tool_completed(&json!({"started_at_ms":1250,"completed_at_ms":1500,
1576                    "item":{"id":"exec-patch","type":"FileChange"}}));
1577                for id in if exec_first { ["wrapper", "other"] } else { ["other", "wrapper"] } {
1578                    t.ingest(
1579                        &json!({"timestamp":"1970-01-01T00:00:02.000Z", "type":"response_item",
1580                        "payload":{"type":"function_call_output","call_id":id}})
1581                        .to_string(),
1582                    );
1583                }
1584                record_response(&mut t, 1_000, 0, 0);
1585                t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1586                record_response(&mut t, 1_100, 0, 0);
1587                let sources = t.summary.context.sources();
1588                assert_eq!(sources.iter().any(|s| s.name == "apply_patch"), other == "wait_agent");
1589                assert_eq!(sources.iter().map(|s| s.calls).sum::<u64>(), 2);
1590                assert_eq!(sources.iter().map(|s| s.tokens).sum::<u64>(), 1_100);
1591            }
1592        }
1593    }
1594
1595    #[test]
1596    fn unfinished_wrappers_do_not_block_context_attribution_in_later_turns() {
1597        use serde_json::json;
1598        for ending in ["task_complete", "turn_aborted", "error"] {
1599            let mut t = CodexTranscript::new("unused.jsonl");
1600            for line in [
1601                json!({"timestamp":"1970-01-01T00:00:01.000Z","type":"response_item",
1602                    "payload":{"type":"custom_tool_call","name":"exec","call_id":"old"}}),
1603                json!({"timestamp":"1970-01-01T00:00:02.000Z","type":"event_msg","payload":{"type":ending}}),
1604                json!({"timestamp":"1970-01-01T00:00:03.000Z","type":"event_msg","payload":{"type":"task_started"}}),
1605                json!({"timestamp":"1970-01-01T00:00:03.100Z","type":"response_item",
1606                    "payload":{"type":"custom_tool_call","name":"exec","call_id":"new"}}),
1607                json!({"type":"event_msg","payload":{"type":"item_completed","started_at_ms":3200,"completed_at_ms":3500,
1608                    "item":{"type":"FileChange","id":"exec-patch"}}}),
1609                json!({"timestamp":"1970-01-01T00:00:04.000Z","type":"response_item",
1610                    "payload":{"type":"custom_tool_call_output","call_id":"new"}}),
1611                json!({"timestamp":"1970-01-01T00:00:05.000Z","type":"response_item",
1612                    "payload":{"type":"custom_tool_call_output","call_id":"old"}}),
1613            ] {
1614                t.ingest(&line.to_string());
1615            }
1616            record_response(&mut t, 1_000, 0, 0);
1617            t.ingest(r#"{"type":"response_item","payload":{"type":"message","role":"assistant"}}"#);
1618            record_response(&mut t, 1_100, 0, 0);
1619            let sources = t.summary.context.sources();
1620            assert_eq!(sources.len(), 3);
1621            for name in ["apply_patch", "exec"] {
1622                let source = sources.iter().find(|s| s.name == name).unwrap();
1623                assert_eq!((source.calls, source.tokens), (1, 50));
1624            }
1625        }
1626    }
1627
1628    #[test]
1629    fn code_mode_completion_metadata_preserves_timing_names_and_errors() {
1630        use serde_json::json;
1631        let mut t = CodexTranscript::new("unused.jsonl");
1632        for (item, name, error) in [
1633            (json!({"type":"CommandExecution","source":"unified_exec_startup","exit_code":1}), "exec_command", true),
1634            (json!({"type":"CommandExecution","source":"unified_exec_interaction","status":"completed"}), "write_stdin", false),
1635            (json!({"type":"FileChange","status":"declined"}), "apply_patch", true),
1636            (json!({"type":"McpToolCall","tool":"search","status":"failed"}), "search", true),
1637            (json!({"type":"DynamicToolCall","tool":"lookup","success":false}), "lookup", true),
1638        ] {
1639            let mut item = item;
1640            item["id"] = json!(format!("exec-{name}"));
1641            let line = json!({
1642                "type":"event_msg", "timestamp":"2026-09-18T07:22:25.000Z",
1643                "payload":{"type":"item_completed", "started_at_ms":1000, "completed_at_ms":1250, "item":item},
1644            })
1645            .to_string();
1646            let before = t.summary.tool_calls;
1647            t.ingest(&line);
1648            let span = t.summary.spans.iter().next_back().unwrap();
1649            assert_eq!((span.name.as_str(), span.duration_ms, span.error), (name, Some(250), error));
1650            assert_eq!(span.started_at, SystemTime::UNIX_EPOCH + Duration::from_secs(1));
1651            // Replayed completions must not count twice.
1652            t.ingest(&line);
1653            assert_eq!(t.summary.tool_calls, before + 1);
1654        }
1655        assert_eq!(t.summary.tool_calls, 5);
1656        assert_eq!(t.summary.spans.len(), 5);
1657        assert!(t.inference.is_none(), "nested results are not new model requests");
1658    }
1659
1660    #[test]
1661    fn completion_items_do_not_duplicate_direct_calls_or_guess_unknown_tools() {
1662        use serde_json::json;
1663        let mut t = CodexTranscript::new("unused.jsonl");
1664        t.ingest(r#"{"timestamp":"2026-09-18T07:22:23.000Z","type":"response_item","payload":{"type":"function_call","name":"exec_command","call_id":"call_direct"}}"#);
1665        for item in [
1666            json!({"type":"CommandExecution","id":"call_direct","source":"unified_exec_startup"}),
1667            json!({"type":"FileChange","id":"call_patch"}),
1668            json!({"type":"CommandExecution","id":"exec-user","source":"user_shell"}),
1669            json!({"type":"CommandExecution","id":"exec-unknown","source":"new_source"}),
1670            json!({"type":"UnknownTool","id":"exec-future"}),
1671            json!({"type":"DynamicToolCall","id":"exec-missing-name"}),
1672            json!({"type":"McpToolCall","id":"exec-empty-name","tool":""}),
1673            json!({"type":"FileChange","id":"exec-"}),
1674            json!({"type":"FileChange"}),
1675        ] {
1676            t.ingest(
1677                &json!({
1678                    "type":"event_msg", "payload":{"type":"item_completed", "started_at_ms":1000, "completed_at_ms":1250, "item":item},
1679                })
1680                .to_string(),
1681            );
1682        }
1683        for (start, end) in [(json!(null), json!(1250)), (json!(1000), json!(null)), (json!(1250), json!(1000))] {
1684            t.ingest(
1685                &json!({
1686                    "type":"event_msg", "payload":{"type":"item_completed", "started_at_ms":start, "completed_at_ms":end,
1687                        "item":{"type":"FileChange","id":"exec-no-timing"}},
1688                })
1689                .to_string(),
1690            );
1691        }
1692        t.ingest(r#"{"timestamp":"2026-09-18T07:22:24.000Z","type":"response_item","payload":{"type":"function_call_output","call_id":"call_direct"}}"#);
1693        assert_eq!(t.summary.tool_calls, 1);
1694        let tools: Vec<_> = t.summary.spans.iter().filter(|s| s.kind == SpanKind::Tool).collect();
1695        assert_eq!(tools.len(), 1);
1696        assert_eq!(tools[0].duration_ms, Some(1000));
1697    }
1698}