diff --git a/crates/statsai-adapters/src/archive/scan/codex.rs b/crates/statsai-adapters/src/archive/scan/codex.rs index 94ad504..f903535 100644 --- a/crates/statsai-adapters/src/archive/scan/codex.rs +++ b/crates/statsai-adapters/src/archive/scan/codex.rs @@ -41,7 +41,10 @@ pub(crate) fn collect_codex_quota_observations( let file = File::open(path).with_context(|| format!("read {}", path.display()))?; let mut reader = BufReader::new(file); let fallback_timestamp = file_modified_timestamp(path).unwrap_or_else(Utc::now); - let mut previous_totals: Option = None; + // Keyed like the usage parser: a line's own session id, else the file's + // latest `session_meta` id. + let mut previous_totals = HashMap::new(); + let mut file_session = String::new(); let mut observations = Vec::new(); let mut line_bytes = Vec::new(); let mut line_number = 0usize; @@ -60,26 +63,23 @@ pub(crate) fn collect_codex_quota_observations( let Ok(value) = serde_json::from_str::(line) else { continue; }; + if value.get("type").and_then(Value::as_str) == Some("session_meta") { + if let Some(id) = value.pointer("/payload/id").and_then(Value::as_str) { + file_session = id.to_owned(); + } + continue; + } if value.get("type").and_then(Value::as_str) != Some("event_msg") || value.pointer("/payload/type").and_then(Value::as_str) != Some("token_count") { continue; } - let info = value.pointer("/payload/info"); - let total_usage = info - .and_then(|info| info.get("total_token_usage")) - .map(codex_usage_counts_from_value); - let usage_sample = info - .and_then(|info| info.get("last_token_usage")) - .map(codex_usage_counts_from_value) - .or_else(|| { - total_usage - .as_ref() - .map(|total| subtract_usage_counts(total, previous_totals.as_ref())) - }); - if let Some(total_usage) = total_usage { - previous_totals = Some(total_usage); - } + let usage_sample = codex_token_count_usage( + value.pointer("/payload/info"), + previous_totals + .entry(session_raw_from_value(&value).unwrap_or_else(|| file_session.clone())) + .or_default(), + ); let observed_at = timestamp_from_nested_value(&value).unwrap_or(fallback_timestamp); if let Some(observation) = codex_quota_observation(source, path, line_number, observed_at, usage_sample, &value) diff --git a/crates/statsai-adapters/src/archive/scan/mod.rs b/crates/statsai-adapters/src/archive/scan/mod.rs index c0d632b..f0c2039 100644 --- a/crates/statsai-adapters/src/archive/scan/mod.rs +++ b/crates/statsai-adapters/src/archive/scan/mod.rs @@ -1,8 +1,8 @@ use super::super::{ canonical_display, codex_project_context_from_value, codex_quota_observation, - codex_usage_counts_from_value, codex_usage_roots, collect_jsonl_files, expand_home_path, + codex_token_count_usage, codex_usage_roots, collect_jsonl_files, expand_home_path, file_modified_timestamp, grok_sessions_root, model_from_nested_value, open_sqlite_readonly, - read_bounded_jsonl_line, resolve_project_context, source_root_path, subtract_usage_counts, + read_bounded_jsonl_line, resolve_project_context, session_raw_from_value, source_root_path, timestamp_from_nested_value, BoundedLineRead, ProjectContextCache, CLAUDE_CODE_PROVIDER, CODEX_PROVIDER, GROK_BUILD_PROVIDER, MAX_JSONL_RECORD_BYTES, OPENCODE_PROVIDER, }; @@ -17,7 +17,7 @@ use rusqlite::OptionalExtension; use serde_json::Value; use statsai_core::{ ArchiveItemKind, ArchiveRole, CoverageStatus, ModelInfo, ProjectInfo, QuotaObservationRecordV1, - SourceLocation, UsageCounts, + SourceLocation, }; use std::collections::{HashMap, HashSet}; use std::fs::File; diff --git a/crates/statsai-adapters/src/cache.rs b/crates/statsai-adapters/src/cache.rs index 98f4a99..370c00d 100644 --- a/crates/statsai-adapters/src/cache.rs +++ b/crates/statsai-adapters/src/cache.rs @@ -8,9 +8,10 @@ use std::time::UNIX_EPOCH; pub(crate) const SCAN_CACHE_SIGNATURE_VERSION: &str = "scan-cache.v1"; // Invalidate unchanged-file scan cache entries whenever provider parsing semantics change, // so historical sessions get rescanned for runtime, pricing, and project context updates. -// activity-invocations.v33: Codex command extraction from parsed_cmd, legacy -// function_call arguments/input, and item_started CommandExecution. -pub(crate) const CODEX_SCAN_CACHE_PARSER_REVISION: &str = "activity-invocations.v33"; +// activity-invocations.v34: Codex usage skips repeated token_count totals and +// counts token_usage_record lines, including compaction, instead of the +// token_count each one pairs with. +pub(crate) const CODEX_SCAN_CACHE_PARSER_REVISION: &str = "activity-invocations.v34"; pub(crate) const CLAUDE_SCAN_CACHE_PARSER_REVISION: &str = "activity-invocations.v28"; // activity-invocations.v20: session rows report their message count as requests and // are priced per message, so long-context tiers stop being decided session-wide. diff --git a/crates/statsai-adapters/src/codex/parse/line.rs b/crates/statsai-adapters/src/codex/parse/line.rs index bd2bff2..c19a3ff 100644 --- a/crates/statsai-adapters/src/codex/parse/line.rs +++ b/crates/statsai-adapters/src/codex/parse/line.rs @@ -17,6 +17,10 @@ pub(crate) fn is_codex_token_count(value: &Value) -> bool { && value.pointer("/payload/type").and_then(Value::as_str) == Some("token_count") } +pub(crate) fn is_codex_token_usage_record(value: &Value) -> bool { + value.get("type").and_then(Value::as_str) == Some("token_usage_record") +} + pub(crate) fn is_codex_task_started(value: &Value) -> bool { value.get("type").and_then(Value::as_str) == Some("event_msg") && value.pointer("/payload/type").and_then(Value::as_str) == Some("task_started") @@ -42,6 +46,8 @@ pub(crate) enum CodexLineKind { ResponseItemMessage, EventUserMessage, TokenCount, + TokenUsageRecord, + Compacted, TaskStarted, TaskComplete, HeadlessUsage, @@ -94,6 +100,9 @@ pub(crate) fn codex_line_kind(line: &str) -> CodexLineKind { if header.contains("\"type\":\"turn_context\"") { return CodexLineKind::TurnContext; } + if header.contains("\"type\":\"compacted\"") { + return CodexLineKind::Compacted; + } if header.contains("\"type\":\"response_item\"") { return if header.contains("\"payload\":{\"type\":\"message\"") { CodexLineKind::ResponseItemMessage @@ -130,6 +139,12 @@ pub(crate) fn codex_line_kind(line: &str) -> CodexLineKind { } return CodexLineKind::Irrelevant; } + // Match this before the headless `"usage":` fallback. On current rollout + // lines the usage object sits past the 256-byte header, so without this + // check the row is dropped as irrelevant. + if header.contains("\"type\":\"token_usage_record\"") { + return CodexLineKind::TokenUsageRecord; + } if header.contains("\"usage\":") || header.contains("\"token_count\":") || header.contains("\"message\":{\"usage\":") diff --git a/crates/statsai-adapters/src/codex/parse/mod.rs b/crates/statsai-adapters/src/codex/parse/mod.rs index 2798609..3b2edc2 100644 --- a/crates/statsai-adapters/src/codex/parse/mod.rs +++ b/crates/statsai-adapters/src/codex/parse/mod.rs @@ -23,7 +23,17 @@ pub(crate) fn parse_codex_file( let mut reader = BufReader::new(file); let fallback_timestamp = file_modified_timestamp(path).unwrap_or_else(Utc::now); let file_fallback_project = project_context_from_path_fallback(root, path); - let mut previous_totals: Option = None; + // Both are per session: one file can interleave several sessions. + let mut previous_totals: HashMap> = HashMap::new(); + // A `token_usage_record` precedes the `token_count` for the same response + // and carries identical usage; compaction inference has a record and no + // token_count. The record is counted and remembered here until the next + // advancing token_count, which is dropped only when it repeats that usage. + // A response whose record is missing or malformed keeps its token_count. + // `compacted` and turn boundaries end the pairing window. Keyed by record + // line so a paired token_count's quota sample can follow the record. + let mut unpaired_records: HashMap = HashMap::new(); + let mut paired_quota_lines: HashMap = HashMap::new(); let mut current_model: Option = None; let mut current_reasoning = ModelReasoningState::default(); let mut current_model_is_fallback = false; @@ -96,6 +106,18 @@ pub(crate) fn parse_codex_file( activity_started_at.elapsed().as_millis() as u64; continue; } + if line_kind == CodexLineKind::Compacted { + // The row carries the whole compacted conversation, so read the + // session from its header instead of parsing it. + let session = codex_json_string_prefix_after_marker( + codex_line_header(line), + "\"session_id\":\"", + 256, + ) + .unwrap_or_else(|| session_raw.clone()); + unpaired_records.remove(&session); + continue; + } if line_kind == CodexLineKind::Irrelevant && !is_codex_quota_line_structurally(line) { continue; } @@ -159,6 +181,7 @@ pub(crate) fn parse_codex_file( model_explicit: false, usage: None, is_token_count_event: false, + is_usage_record: false, is_task_started: false, is_task_complete: false, message_role: role, @@ -224,6 +247,7 @@ pub(crate) fn parse_codex_file( model_explicit: false, usage: None, is_token_count_event: false, + is_usage_record: false, is_task_started: false, is_task_complete: false, message_role, @@ -286,6 +310,7 @@ pub(crate) fn parse_codex_file( model_explicit: false, usage: None, is_token_count_event: false, + is_usage_record: false, is_task_started: false, is_task_complete: false, message_role: None, @@ -351,8 +376,27 @@ pub(crate) fn parse_codex_file( } let is_token_count_event = is_codex_token_count(&value); + let is_token_usage_record = is_codex_token_usage_record(&value); + // Records are keyed exactly like the token_count lines they pair with. + // `payload.thread_id` is not used: a fork copies its parent's history, + // records and all, and token_count carries no thread id to match it. + let event_session_raw = + session_raw_from_value(&value).unwrap_or_else(|| session_raw.clone()); + let record_usage = is_token_usage_record + .then(|| codex_token_usage_record_from_value(&value)) + .flatten() + .map(|record| record.usage) + .filter(|usage| usage.total_tokens.is_some()); + if let Some(usage) = &record_usage { + unpaired_records.insert(event_session_raw.clone(), (usage.clone(), index)); + } let is_task_started = is_codex_task_started(&value); let is_task_complete = is_codex_task_complete(&value); + // A record and its token_count belong to one turn. An unpaired record + // must not suppress a later turn's token_count that happens to match. + if is_task_started || is_task_complete { + unpaired_records.remove(&event_session_raw); + } let task_started_at = is_task_started .then(|| codex_task_timestamp(&value, &["/payload/started_at"])) .flatten(); @@ -377,29 +421,38 @@ pub(crate) fn parse_codex_file( let user_message_preview = collect_tasks .then(|| codex_user_message_preview(&value)) .flatten(); - let event_session_raw = - session_raw_from_value(&value).unwrap_or_else(|| session_raw.clone()); - let usage = if is_token_count_event { - let info = value.pointer("/payload/info"); - let total_usage = info - .and_then(|info| info.get("total_token_usage")) - .map(codex_usage_counts_from_value); - let usage = info - .and_then(|info| info.get("last_token_usage")) - .map(codex_usage_counts_from_value) - .or_else(|| { - total_usage - .as_ref() - .map(|total| subtract_usage_counts(total, previous_totals.as_ref())) - }); - if let Some(total) = total_usage { - previous_totals = Some(total); - } - usage + let token_count_usage = if is_token_count_event { + codex_token_count_usage( + value.pointer("/payload/info"), + previous_totals + .entry(event_session_raw.clone()) + .or_default(), + ) + } else { + None + }; + let usage = if is_token_usage_record { + record_usage + } else if is_token_count_event { + token_count_usage.as_ref().and_then(|usage| { + match unpaired_records.remove(&event_session_raw) { + Some((record_usage, record_line)) + if same_codex_response_usage(&record_usage, usage) => + { + paired_quota_lines.insert(record_line, index); + None + } + _ => Some(usage.clone()), + } + }) } else { codex_headless_usage_value(&value).map(codex_usage_counts_from_value) }; - let quota_usage_sample = usage.clone(); + let quota_usage_sample = if is_token_count_event { + token_count_usage + } else { + usage.clone() + }; let (timestamp, timestamp_inferred) = timestamp_from_nested_value(&value) .map(|timestamp| (timestamp, false)) @@ -474,6 +527,7 @@ pub(crate) fn parse_codex_file( model_explicit, usage, is_token_count_event, + is_usage_record: is_token_usage_record, is_task_started, is_task_complete, message_role, @@ -514,6 +568,7 @@ pub(crate) fn parse_codex_file( .as_ref() .map(|_| vec![record.line_number]) .unwrap_or_default(), + quota_lines: Vec::new(), project: record.project.clone(), }); if record.usage.is_some() { @@ -604,6 +659,8 @@ pub(crate) fn parse_codex_file( turn.last_usage = Some(usage); turn.usage_lines.push(record.line_number); } + } else if quota_observation_indices.contains_key(&record.line_number) { + turn.quota_lines.push(record.line_number); } } @@ -678,6 +735,7 @@ pub(crate) fn parse_codex_file( }, ); let mut linked_quota_lines = turn.usage_lines.clone(); + linked_quota_lines.extend_from_slice(&turn.quota_lines); if record.usage.is_some() { linked_quota_lines.push(record.line_number); } @@ -888,7 +946,9 @@ pub(crate) fn parse_codex_file( project: record .project .or_else(|| project_context_from_path_fallback(root, path)), - event_kind: if record.is_token_count_event { + event_kind: if record.is_usage_record { + "codex_usage_record" + } else if record.is_token_count_event { "codex_token_count" } else { "codex_headless_usage" @@ -898,7 +958,7 @@ pub(crate) fn parse_codex_file( source_type: "jsonl", model_inferred: record.model_inferred, timestamp_inferred: record.timestamp_inferred, - deduplication: if record.is_token_count_event { + deduplication: if record.is_token_count_event || record.is_usage_record { EventDeduplication::PathIndependent } else { EventDeduplication::SessionScoped @@ -906,10 +966,13 @@ pub(crate) fn parse_codex_file( dedupe_salt: None, }, ); + // A record's paired token_count carries the quota sample for it. + let mut linked_quota_lines = vec![record.line_number]; + linked_quota_lines.extend(paired_quota_lines.get(&record.line_number)); link_quota_observations( ctx.scan, "a_observation_indices, - &[record.line_number], + &linked_quota_lines, &event.event_id, QuotaUsageLinkKind::RecordEvent, ); diff --git a/crates/statsai-adapters/src/codex/parse/types.rs b/crates/statsai-adapters/src/codex/parse/types.rs index 1b1bbe4..2b7749e 100644 --- a/crates/statsai-adapters/src/codex/parse/types.rs +++ b/crates/statsai-adapters/src/codex/parse/types.rs @@ -11,6 +11,7 @@ pub(crate) struct CodexLineRecord { pub(crate) model_explicit: bool, pub(crate) usage: Option, pub(crate) is_token_count_event: bool, + pub(crate) is_usage_record: bool, pub(crate) is_task_started: bool, pub(crate) is_task_complete: bool, pub(crate) message_role: Option, @@ -108,5 +109,9 @@ pub(crate) struct ActiveCodexTurn { pub(crate) prompt_previews: Vec, pub(crate) last_activity_at: DateTime, pub(crate) usage_lines: Vec, + /// `token_count` lines inside the turn that carry no usage, because they + /// repeat a total or pair with a `token_usage_record`. Their quota + /// observations still belong to this turn's event. + pub(crate) quota_lines: Vec, pub(crate) project: Option, } diff --git a/crates/statsai-adapters/src/codex/parse/usage.rs b/crates/statsai-adapters/src/codex/parse/usage.rs index 7a3fcab..87a70ad 100644 --- a/crates/statsai-adapters/src/codex/parse/usage.rs +++ b/crates/statsai-adapters/src/codex/parse/usage.rs @@ -1,5 +1,140 @@ use super::*; +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct CodexTokenUsageRecord { + pub(crate) usage: UsageCounts, + pub(crate) turn_id: Option, + pub(crate) root_turn_id: Option, + pub(crate) response_id: Option, + pub(crate) thread_id: Option, + pub(crate) session_id: Option, +} + +impl CodexTokenUsageRecord { + /// Parent session for a sub-agent record. Main threads use the same id for + /// `thread_id` and `session_id`, so they have no parent to roll up to. + /// + /// Callers that roll sub-agent cost onto the parent turn should read this + /// with `root_turn_id`. Nothing in the store reads it yet. + #[cfg_attr(not(test), allow(dead_code))] + pub(crate) fn parent_session_id(&self) -> Option<&str> { + match (self.session_id.as_deref(), self.thread_id.as_deref()) { + (Some(session_id), Some(thread_id)) if session_id != thread_id => Some(session_id), + _ => None, + } + } +} + +pub(crate) fn codex_token_usage_record_from_value(value: &Value) -> Option { + if value.get("type").and_then(Value::as_str) != Some("token_usage_record") { + return None; + } + let payload = value.get("payload")?; + let usage = payload.get("usage")?; + let string_at = |key: &str| { + payload + .get(key) + .and_then(Value::as_str) + .map(ToOwned::to_owned) + }; + Some(CodexTokenUsageRecord { + usage: codex_usage_counts_from_value(usage), + turn_id: string_at("turn_id"), + root_turn_id: string_at("root_turn_id"), + response_id: string_at("response_id"), + thread_id: string_at("thread_id"), + session_id: string_at("session_id"), + }) +} + +pub(crate) struct CodexCumulativeTotal { + raw: Value, + counts: UsageCounts, +} + +/// Whether a `token_count` repeats the response a `token_usage_record` already +/// counted. Compares the normalized counts, so whether either side reported +/// `total_tokens` explicitly does not matter. +pub(crate) fn same_codex_response_usage(record: &UsageCounts, token_count: &UsageCounts) -> bool { + let key = |usage: &UsageCounts| { + ( + usage.input_tokens.unwrap_or(0), + usage.output_tokens.unwrap_or(0), + usage.cache_creation_tokens.unwrap_or(0), + usage.cache_read_tokens.unwrap_or(0), + usage.reasoning_tokens.unwrap_or(0), + ) + }; + key(record) == key(token_count) +} + +/// Usage to attribute for one `token_count` line. +/// +/// When `total_token_usage` repeats the previous cumulative total for this +/// session, the line is a duplicate snapshot (or the post-compaction context +/// size) and contributes nothing. A smaller total is a fork or resume reset +/// and still counts, from zero. Lines with no cumulative total keep +/// `last_token_usage`. +pub(crate) fn codex_token_count_usage( + info: Option<&Value>, + previous: &mut Option, +) -> Option { + let info = info?; + let last_usage = info + .get("last_token_usage") + .map(codex_usage_counts_from_value); + let Some(total_value) = info + .get("total_token_usage") + .filter(|value| value.is_object()) + else { + return last_usage; + }; + let total_counts = codex_usage_counts_from_value(total_value); + let unchanged = previous + .as_ref() + .is_some_and(|previous| cumulative_total_unchanged(previous, total_value)); + let reset = previous + .as_ref() + .is_some_and(|previous| cumulative_total_decreased(previous, total_value)); + let usage = if unchanged { + None + } else { + last_usage.or_else(|| { + // After a reset the counter restarted at zero, so the whole new + // total is usage; subtracting the larger old total would clamp + // every field to zero. + let baseline = (!reset) + .then(|| previous.as_ref().map(|previous| &previous.counts)) + .flatten(); + Some(crate::subtract_usage_counts(&total_counts, baseline)) + }) + }; + *previous = Some(CodexCumulativeTotal { + raw: total_value.clone(), + counts: total_counts, + }); + usage +} + +fn cumulative_total_unchanged(previous: &CodexCumulativeTotal, raw: &Value) -> bool { + match (previous.raw.get("total_tokens"), raw.get("total_tokens")) { + // Compare the provider's cumulative counter, not the normalized sum. + // An equal counter is a repeated snapshot. A smaller one is a reset. + (Some(before), Some(after)) => before == after, + _ => &previous.raw == raw, + } +} + +fn cumulative_total_decreased(previous: &CodexCumulativeTotal, raw: &Value) -> bool { + match ( + previous.raw.get("total_tokens").and_then(Value::as_u64), + raw.get("total_tokens").and_then(Value::as_u64), + ) { + (Some(before), Some(after)) => after < before, + _ => false, + } +} + pub(crate) fn codex_usage_counts_from_value(value: &Value) -> UsageCounts { let raw_input = number_at_any(value, &["input_tokens", "prompt_tokens", "input"]); let raw_output = number_at_any(value, &["output_tokens", "completion_tokens", "output"]); @@ -10,6 +145,8 @@ pub(crate) fn codex_usage_counts_from_value(value: &Value) -> UsageCounts { "cacheCreationInputTokens", "cache_creation_tokens", "cacheCreationTokens", + "cache_write_input_tokens", + "cacheWriteInputTokens", ], ); let raw_cache_read = number_at_any( @@ -33,9 +170,9 @@ pub(crate) fn codex_usage_counts_from_value(value: &Value) -> UsageCounts { ) } -// Codex reports cached input and reasoning output as subsets of the top-level -// input/output counters. Normalize that inclusive provider shape into the -// additive contract used everywhere else in statsai. +// Codex reports cached input, cache writes, and reasoning output as subsets of +// the top-level input/output counters. Normalize that inclusive provider shape +// into the additive contract used everywhere else in statsai. pub(crate) fn normalize_codex_usage_counts( raw_input: Option, raw_output: Option, diff --git a/crates/statsai-adapters/src/codex/tests/parse/session.rs b/crates/statsai-adapters/src/codex/tests/parse/session.rs index f0bc803..61e7073 100644 --- a/crates/statsai-adapters/src/codex/tests/parse/session.rs +++ b/crates/statsai-adapters/src/codex/tests/parse/session.rs @@ -105,6 +105,34 @@ fn codex_line_kind_uses_header_window_for_large_user_messages() { ); } +#[test] +fn codex_token_usage_record_is_never_classified_as_headless_usage() { + // `"usage":` on a real record sits past the 256-byte header, under payload. + // A compact line still carries that key inside the header. Neither shape + // may fall through to HeadlessUsage. + let compact = r#"{"timestamp":"2026-09-01T10:00:02.080Z","type":"token_usage_record","payload":{"usage":{"input_tokens":1,"output_tokens":1,"total_tokens":2}}}"#; + assert!(compact.len() < 256); + assert!(compact.contains("\"usage\":")); + assert_eq!(codex_line_kind(compact), CodexLineKind::TokenUsageRecord); + + let fixture = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join( + "tests/fixtures/codex/usage-record/mixed/sessions/2026/09/01/rollout-fixture-usage-record.jsonl", + ); + let realistic = std::fs::read_to_string(fixture) + .expect("fixture") + .lines() + .find(|line| line.contains("\"token_usage_record\"")) + .expect("record line") + .to_string(); + let usage_at = realistic.find("\"usage\":").expect("usage field"); + assert!( + usage_at > 256, + "the fixture must keep usage outside the header window, at {usage_at}" + ); + assert!(!codex_line_header(&realistic).contains("\"usage\":")); + assert_eq!(codex_line_kind(&realistic), CodexLineKind::TokenUsageRecord); +} + #[test] fn codex_task_spans_prefer_real_user_message_over_wrapper_response_item() { let dir = tempfile::tempdir().expect("tempdir"); diff --git a/crates/statsai-adapters/src/codex/tests/parse/usage.rs b/crates/statsai-adapters/src/codex/tests/parse/usage.rs index 523f64d..8d5df7d 100644 --- a/crates/statsai-adapters/src/codex/tests/parse/usage.rs +++ b/crates/statsai-adapters/src/codex/tests/parse/usage.rs @@ -254,3 +254,708 @@ fn codex_caps_cached_input_to_input() { assert_eq!(scan.events[0].usage.input_tokens, Some(0)); assert_eq!(scan.events[0].usage.cache_read_tokens, Some(10)); } + +fn scan_session_lines(lines: &[String]) -> (crate::AdapterScan, usize, usize) { + let dir = tempfile::tempdir().expect("tempdir"); + let sessions = dir.path().join("sessions"); + std::fs::create_dir_all(&sessions).expect("sessions"); + let path = sessions.join("session.jsonl"); + let mut file = File::create(&path).expect("fixture"); + for line in lines { + writeln!(file, "{line}").expect("write"); + } + drop(file); + let source = SourceLocation::local_adapter( + CODEX_PROVIDER, + "test", + "0", + dir.path(), + LocationOrigin::Configured, + ); + let scan = scan_codex_source(&CodexAdapter, &source, &options()).expect("scan"); + let quota = crate::archive::collect_codex_quota_observations(&source, &path).expect("quota"); + let positive_samples = quota + .iter() + .filter(|record| { + record + .observation + .usage_sample + .as_ref() + .is_some_and(|usage| usage.computed_total() > 0) + }) + .count(); + (scan, quota.len(), positive_samples) +} + +fn token_count_line(timestamp: &str, limit_id: &str, total: Option, last: Value) -> String { + // `type` has to lead `payload`. The line classifier matches a fixed header + // prefix, and serde_json's map orders keys alphabetically. + let total_field = total + .map(|total| format!(r#""total_token_usage":{total},"#)) + .unwrap_or_default(); + format!( + r#"{{"timestamp":"{timestamp}","type":"event_msg","payload":{{"type":"token_count","info":{{{total_field}"last_token_usage":{last}}},"rate_limits":{{"limit_id":"{limit_id}","primary":{{"used_percent":10,"window_minutes":300,"resets_at":1790000000}},"secondary":{{"used_percent":4,"window_minutes":10080,"resets_at":1790400000}},"plan_type":"plus"}}}}}}"# + ) +} + +fn turn_bounds(started_at: &str, completed_at: &str) -> (String, String) { + ( + format!( + r#"{{"timestamp":"{started_at}","type":"event_msg","payload":{{"type":"task_started","started_at":"{started_at}"}}}}"# + ), + format!( + r#"{{"timestamp":"{completed_at}","type":"event_msg","payload":{{"type":"task_complete","completed_at":"{completed_at}"}}}}"# + ), + ) +} + +#[test] +fn codex_counts_an_unchanged_token_total_once_for_every_limit_bucket() { + // Older builds repeat total_token_usage with an identical last_token_usage. + // Newer builds do it once more per rate-limit bucket. Neither copy is a + // new response, so the turn total has to stay on the first line. + let usage = serde_json::json!({ + "input_tokens": 100, + "cached_input_tokens": 25, + "output_tokens": 20, + "reasoning_output_tokens": 5, + "total_tokens": 120 + }); + let (started, completed) = turn_bounds("2026-09-01T10:00:00Z", "2026-09-01T10:00:04Z"); + let lines = vec![ + started, + token_count_line( + "2026-09-01T10:00:01Z", + "codex", + Some(usage.clone()), + usage.clone(), + ), + token_count_line( + "2026-09-01T10:00:02Z", + "codex", + Some(usage.clone()), + usage.clone(), + ), + token_count_line( + "2026-09-01T10:00:03Z", + "premium", + Some(usage.clone()), + usage, + ), + completed, + ]; + + let (scan, quota_len, positive_samples) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.events[0].usage.input_tokens, Some(75)); + assert_eq!(scan.events[0].usage.cache_read_tokens, Some(25)); + assert_eq!(scan.events[0].usage.output_tokens, Some(15)); + assert_eq!(scan.events[0].usage.reasoning_tokens, Some(5)); + assert_eq!(scan.events[0].usage.computed_total(), 120); + assert_eq!(scan.events[0].usage.requests, Some(1)); + assert_eq!( + quota_len, 3, + "each bucket still produces a quota observation" + ); + assert_eq!(positive_samples, 1, "repeated totals carry no usage sample"); +} + +#[test] +fn codex_ignores_the_phantom_token_count_after_compaction() { + // After `compacted`, Codex emits a token_count whose cumulative total is + // unchanged and whose last_token_usage is only the new context size. + let total = serde_json::json!({ + "input_tokens": 200, + "cached_input_tokens": 40, + "output_tokens": 30, + "reasoning_output_tokens": 10, + "total_tokens": 230 + }); + let last = serde_json::json!({ + "input_tokens": 80, + "cached_input_tokens": 20, + "output_tokens": 16, + "reasoning_output_tokens": 4, + "total_tokens": 96 + }); + let phantom_last = serde_json::json!({ + "input_tokens": 4, + "cached_input_tokens": 0, + "output_tokens": 1, + "reasoning_output_tokens": 0, + "total_tokens": 5 + }); + let (started, completed) = turn_bounds("2026-09-01T10:00:00Z", "2026-09-01T10:00:05Z"); + let lines = vec![ + started, + token_count_line("2026-09-01T10:00:01Z", "codex", Some(total.clone()), last), + r#"{"timestamp":"2026-09-01T10:00:02Z","type":"compacted","payload":{"message":"synthetic","window_number":1}}"#.to_string(), + token_count_line("2026-09-01T10:00:03Z", "codex", Some(total), phantom_last), + completed, + ]; + + let (scan, quota_len, positive_samples) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.events[0].usage.computed_total(), 96); + assert_eq!(scan.events[0].usage.requests, Some(1)); + assert_eq!(quota_len, 2); + assert_eq!(positive_samples, 1); +} + +#[test] +fn codex_interleaved_sessions_keep_their_own_usage_source_and_totals() { + let usage = r#"{"input_tokens":90,"output_tokens":10,"total_tokens":100}"#; + let record = format!( + r#"{{"timestamp":"2026-09-01T10:00:02Z","session_id":"session-a","type":"token_usage_record","payload":{{"thread_id":"session-a","usage":{usage}}}}}"# + ); + assert_interleaved_sessions_keep_their_own_usage(usage, record); +} + +#[test] +fn codex_forked_history_pairs_records_named_for_the_parent_thread() { + // A fork copies its parent's history: the copied record names the parent + // thread, and the token_count after it carries no thread id at all. Both + // belong to the forked file's session and describe one response. + let usage = serde_json::json!({"input_tokens": 90, "output_tokens": 10, "total_tokens": 100}); + let (started, completed) = turn_bounds("2026-09-01T10:00:00Z", "2026-09-01T10:00:05Z"); + let lines = vec![ + r#"{"timestamp":"2026-09-01T09:59:59Z","type":"session_meta","payload":{"id":"00000000-0000-7000-8000-00000000f0f0"}}"#.to_string(), + started, + format!( + r#"{{"timestamp":"2026-09-01T10:00:01Z","type":"token_usage_record","payload":{{"thread_id":"00000000-0000-7000-8000-000000000001","session_id":"00000000-0000-7000-8000-000000000001","usage":{usage}}}}}"# + ), + token_count_line("2026-09-01T10:00:02Z", "codex", Some(usage.clone()), usage), + completed, + ]; + + let (scan, _, _) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.events[0].usage.computed_total(), 100); + assert_eq!(scan.events[0].usage.requests, Some(1)); +} + +#[test] +fn codex_compaction_record_does_not_suppress_a_later_response_in_the_turn() { + // Compaction inference has a record and no token_count. A later response + // in the same turn has no usable record, and its token_count happens to + // repeat the compaction usage. + let usage = serde_json::json!({"input_tokens": 90, "output_tokens": 10, "total_tokens": 100}); + let total = serde_json::json!({"input_tokens": 180, "output_tokens": 20, "total_tokens": 200}); + let (started, completed) = turn_bounds("2026-09-01T10:00:00Z", "2026-09-01T10:00:09Z"); + let lines = vec![ + started, + format!( + r#"{{"timestamp":"2026-09-01T10:00:01Z","type":"token_usage_record","payload":{{"usage":{usage}}}}}"# + ), + r#"{"timestamp":"2026-09-01T10:00:02Z","type":"compacted","payload":{"message":"synthetic","window_number":1}}"#.to_string(), + token_count_line("2026-09-01T10:00:03Z", "codex", Some(total), usage), + completed, + ]; + + let (scan, _, _) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.events[0].usage.computed_total(), 200); + assert_eq!(scan.events[0].usage.requests, Some(2)); +} + +#[test] +fn codex_record_outside_a_turn_links_its_paired_quota_sample() { + let usage = serde_json::json!({"input_tokens": 90, "output_tokens": 10, "total_tokens": 100}); + let lines = vec![ + format!( + r#"{{"timestamp":"2026-09-01T10:00:01Z","type":"token_usage_record","payload":{{"usage":{usage}}}}}"# + ), + token_count_line("2026-09-01T10:00:02Z", "codex", Some(usage.clone()), usage), + ]; + + let (scan, _, _) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.quota_observations.len(), 1); + assert_eq!( + scan.quota_observations[0] + .observation + .usage_event_id + .as_ref(), + Some(&scan.events[0].event_id) + ); +} + +#[test] +fn codex_counter_reset_without_last_usage_counts_the_new_total() { + let earlier = + serde_json::json!({"input_tokens": 900, "output_tokens": 100, "total_tokens": 1000}); + let restarted = serde_json::json!({"input_tokens": 45, "output_tokens": 5, "total_tokens": 50}); + let token_count = |timestamp: &str, total: &serde_json::Value| { + format!( + r#"{{"timestamp":"{timestamp}","type":"event_msg","payload":{{"type":"token_count","info":{{"total_token_usage":{total}}}}}}}"# + ) + }; + let lines = vec![ + token_count("2026-09-01T10:00:01Z", &earlier), + token_count("2026-09-01T10:00:02Z", &restarted), + ]; + + let (scan, _, _) = scan_session_lines(&lines); + + let mut totals = scan + .events + .iter() + .map(|event| event.usage.computed_total()) + .collect::>(); + totals.sort_unstable(); + assert_eq!(totals, vec![50, 1000]); +} + +fn assert_interleaved_sessions_keep_their_own_usage(usage: &str, record: String) { + // Session A switches to records; session B stays on token_count. B's first + // cumulative total equals A's last one and is still a new response for B. + let token_count = |timestamp: &str, session: &str| { + format!( + r#"{{"timestamp":"{timestamp}","session_id":"{session}","type":"event_msg","payload":{{"type":"token_count","info":{{"total_token_usage":{usage},"last_token_usage":{usage}}},"rate_limits":{{"limit_id":"codex","primary":{{"used_percent":10,"window_minutes":300,"resets_at":1790000000}},"plan_type":"plus"}}}}}}"# + ) + }; + let task = |timestamp: &str, session: &str, kind: &str, at: &str| { + format!( + r#"{{"timestamp":"{timestamp}","session_id":"{session}","type":"event_msg","payload":{{"type":"{kind}","{at}":"{timestamp}"}}}}"# + ) + }; + let lines = vec![ + task( + "2026-09-01T10:00:00Z", + "session-a", + "task_started", + "started_at", + ), + task( + "2026-09-01T10:00:01Z", + "session-b", + "task_started", + "started_at", + ), + record, + token_count("2026-09-01T10:00:03Z", "session-a"), + token_count("2026-09-01T10:00:04Z", "session-b"), + task( + "2026-09-01T10:00:05Z", + "session-a", + "task_complete", + "completed_at", + ), + task( + "2026-09-01T10:00:06Z", + "session-b", + "task_complete", + "completed_at", + ), + ]; + + let (scan, quota_len, positive_samples) = scan_session_lines(&lines); + + // The archive quota scan keys cumulative totals the same way, so B's + // sample is not mistaken for a repeat of A's. + assert_eq!(quota_len, 2); + assert_eq!(positive_samples, 2); + let mut totals = scan + .events + .iter() + .map(|event| { + ( + event.session.local_session_id_hash.clone(), + event.usage.computed_total(), + ) + }) + .collect::>(); + totals.sort(); + let mut expected = vec![ + (Some(hash_text("session-a")), 100), + (Some(hash_text("session-b")), 100), + ]; + expected.sort(); + assert_eq!(totals, expected); +} + +#[test] +fn codex_record_without_usage_keeps_token_count_as_the_usage_source() { + // A record that carries no usable usage cannot stand in for the + // token_count lines after it, or those responses would vanish. + let total = serde_json::json!({ + "input_tokens": 100, + "output_tokens": 20, + "total_tokens": 120 + }); + let (started, completed) = turn_bounds("2026-09-01T10:00:00Z", "2026-09-01T10:00:05Z"); + let lines = vec![ + started, + r#"{"timestamp":"2026-09-01T10:00:01Z","type":"token_usage_record","payload":{"turn_id":"00000000-0000-7000-8000-0000000000a1"}}"#.to_string(), + r#"{"timestamp":"2026-09-01T10:00:02Z","type":"token_usage_record","payload":{"usage":{}}}"#.to_string(), + token_count_line("2026-09-01T10:00:03Z", "codex", Some(total.clone()), total), + completed, + ]; + + let (scan, _, _) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.events[0].usage.computed_total(), 120); + assert_eq!(scan.events[0].usage.requests, Some(1)); +} + +#[test] +fn codex_unpaired_record_does_not_suppress_a_later_turn() { + // Turn one ends on a record with no token_count (compaction). Turn two + // has no usable record, and its token_count happens to repeat that usage. + let usage = serde_json::json!({"input_tokens": 90, "output_tokens": 10, "total_tokens": 100}); + let (first_started, first_completed) = + turn_bounds("2026-09-01T10:00:00Z", "2026-09-01T10:00:02Z"); + let (second_started, second_completed) = + turn_bounds("2026-09-01T10:01:00Z", "2026-09-01T10:01:02Z"); + let lines = vec![ + first_started, + format!( + r#"{{"timestamp":"2026-09-01T10:00:01Z","type":"token_usage_record","payload":{{"usage":{usage}}}}}"# + ), + first_completed, + second_started, + token_count_line("2026-09-01T10:01:01Z", "codex", Some(usage.clone()), usage), + second_completed, + ]; + + let (scan, _, _) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 2); + assert!(scan + .events + .iter() + .all(|event| event.usage.computed_total() == 100)); +} + +#[test] +fn codex_malformed_record_after_records_began_keeps_its_token_count() { + // The first response is recorded normally; the second response's record + // is malformed, so its token_count is the only evidence of it. + let first = serde_json::json!({"input_tokens": 90, "output_tokens": 10, "total_tokens": 100}); + let first_total = first.clone(); + let second = serde_json::json!({"input_tokens": 45, "output_tokens": 5, "total_tokens": 50}); + let second_total = + serde_json::json!({"input_tokens": 135, "output_tokens": 15, "total_tokens": 150}); + let (started, completed) = turn_bounds("2026-09-01T10:00:00Z", "2026-09-01T10:00:09Z"); + let lines = vec![ + started, + format!( + r#"{{"timestamp":"2026-09-01T10:00:01Z","type":"token_usage_record","payload":{{"usage":{first}}}}}"# + ), + token_count_line("2026-09-01T10:00:02Z", "codex", Some(first_total), first), + r#"{"timestamp":"2026-09-01T10:00:03Z","type":"token_usage_record","payload":{"usage":{}}}"#.to_string(), + token_count_line("2026-09-01T10:00:04Z", "codex", Some(second_total), second), + completed, + ]; + + let (scan, _, _) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.events[0].usage.computed_total(), 150); + assert_eq!(scan.events[0].usage.requests, Some(2)); +} + +#[test] +fn codex_counts_a_decreased_token_total_as_a_reset() { + // A fork or resume can restart the cumulative total. That is new usage, + // unlike an exact repeat of the previous total. + let first = serde_json::json!({ + "input_tokens": 80, + "output_tokens": 20, + "total_tokens": 100 + }); + let reset = serde_json::json!({ + "input_tokens": 30, + "output_tokens": 10, + "total_tokens": 40 + }); + let (started, completed) = turn_bounds("2026-09-01T10:00:00Z", "2026-09-01T10:00:04Z"); + let lines = vec![ + started, + token_count_line("2026-09-01T10:00:01Z", "codex", Some(first.clone()), first), + token_count_line("2026-09-01T10:00:02Z", "codex", Some(reset.clone()), reset), + completed, + ]; + + let (scan, _, positive_samples) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.events[0].usage.input_tokens, Some(110)); + assert_eq!(scan.events[0].usage.output_tokens, Some(30)); + assert_eq!(scan.events[0].usage.computed_total(), 140); + assert_eq!(scan.events[0].usage.requests, Some(2)); + assert_eq!(positive_samples, 2); +} + +#[test] +fn codex_keeps_last_token_usage_when_the_cumulative_total_is_missing() { + let (started, completed) = turn_bounds("2026-09-01T10:00:00Z", "2026-09-01T10:00:04Z"); + let lines = vec![ + started, + token_count_line( + "2026-09-01T10:00:01Z", + "codex", + None, + serde_json::json!({"input_tokens": 40, "output_tokens": 10, "total_tokens": 50}), + ), + token_count_line( + "2026-09-01T10:00:02Z", + "codex", + None, + serde_json::json!({"input_tokens": 20, "output_tokens": 5, "total_tokens": 25}), + ), + completed, + ]; + + let (scan, _, _) = scan_session_lines(&lines); + + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.events[0].usage.input_tokens, Some(60)); + assert_eq!(scan.events[0].usage.output_tokens, Some(15)); + assert_eq!(scan.events[0].usage.computed_total(), 75); + assert_eq!(scan.events[0].usage.requests, Some(2)); +} + +#[test] +fn codex_compaction_fixture_skips_the_repeated_total_after_compacted() { + // The fixture's second token_count repeats total_token_usage + // (total_tokens 11081643) after `compacted`, while last_token_usage changes + // to the post-compaction context (total_tokens 17983). That line is not a + // billed response. Counting it added 17983 on top of the four advancing + // rows: 318182 + 23975 + 24069 + 23998 = 390224. + let root = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/codex/compaction"); + let source = SourceLocation::local_adapter( + CODEX_PROVIDER, + "test", + "0", + &root, + LocationOrigin::Configured, + ); + let scan = scan_codex_source(&CodexAdapter, &source, &options()).expect("scan"); + let mut totals: Vec = scan + .events + .iter() + .map(|event| event.usage.computed_total()) + .collect(); + totals.sort_unstable(); + + assert_eq!( + totals, + vec![23_975, 23_998, 24_069, 318_182], + "the phantom 17983 context size must not be added" + ); + assert_eq!(totals.iter().sum::(), 390_224); + assert!( + scan.quota_observations.len() >= 5, + "the phantom line still produces a quota observation" + ); + assert!(scan + .quota_observations + .iter() + .any(|record| { record.observation.usage_sample.is_none() })); + assert!(scan.quota_observations.iter().all(|record| { + record + .observation + .usage_sample + .as_ref() + .map(UsageCounts::computed_total) + != Some(17_983) + })); +} + +#[test] +fn codex_usage_counts_treat_cache_writes_as_an_inclusive_input_subset() { + let value = serde_json::json!({ + "input_tokens": 1000, + "cached_input_tokens": 600, + "cache_write_input_tokens": 100, + "output_tokens": 0, + "total_tokens": 1000 + }); + + let usage = codex_usage_counts_from_value(&value); + + assert_eq!(usage.input_tokens, Some(300)); + assert_eq!(usage.cache_read_tokens, Some(600)); + assert_eq!(usage.cache_creation_tokens, Some(100)); + assert_eq!(usage.computed_total(), 1000); + + let camel = serde_json::json!({ + "input_tokens": 1000, + "cached_input_tokens": 600, + "cacheWriteInputTokens": 100 + }); + let camel_usage = codex_usage_counts_from_value(&camel); + assert_eq!(camel_usage.input_tokens, Some(300)); + assert_eq!(camel_usage.cache_read_tokens, Some(600)); + assert_eq!(camel_usage.cache_creation_tokens, Some(100)); +} + +fn scan_fixture(relative_root: &str) -> crate::AdapterScan { + let root = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join(relative_root); + let source = SourceLocation::local_adapter( + CODEX_PROVIDER, + "test", + "0", + &root, + LocationOrigin::Configured, + ); + scan_codex_source(&CodexAdapter, &source, &options()).expect("scan") +} + +#[test] +fn codex_usage_record_fixture_counts_each_record_once_across_eras() { + // Mixed era: token_count rows before the file's first token_usage_record, + // then records only. The paired token_count and the post-compaction + // phantom must not be added, and thread_token_usage is not the response. + let scan = scan_fixture("tests/fixtures/codex/usage-record/mixed"); + + assert_eq!( + scan.events.len(), + 2, + "one legacy turn and one record-era turn" + ); + let legacy = &scan.events[0].usage; + assert_eq!(legacy.input_tokens, Some(100)); + assert_eq!(legacy.output_tokens, Some(10)); + assert_eq!(legacy.computed_total(), 110); + assert_eq!(legacy.requests, Some(1)); + assert_eq!( + scan.events[0] + .model + .as_ref() + .and_then(|model| model.provider_model_id.as_deref()), + Some("gpt-5") + ); + + let recorded = &scan.events[1].usage; + assert_eq!(recorded.input_tokens, Some(650)); + assert_eq!(recorded.cache_read_tokens, Some(350)); + assert_eq!(recorded.cache_creation_tokens, Some(0)); + assert_eq!(recorded.output_tokens, Some(45)); + assert_eq!(recorded.reasoning_tokens, Some(15)); + assert_eq!(recorded.computed_total(), 1_060); + assert_eq!(recorded.requests, Some(2)); + assert_eq!( + scan.events[1] + .model + .as_ref() + .and_then(|model| model.provider_model_id.as_deref()), + Some("gpt-5") + ); + + let session_total: u64 = scan + .events + .iter() + .map(|event| event.usage.computed_total()) + .sum(); + assert_eq!(session_total, 1_170); + assert_eq!( + scan.quota_observations.len(), + 5, + "token_count lines still supply quota after records begin" + ); + // A token_count carries no usage once records begin, but its quota + // observation still belongs to the turn that the records fed. + let linked_to = |event_index: usize| { + scan.quota_observations + .iter() + .filter(|record| { + record.observation.usage_event_id.as_ref() + == Some(&scan.events[event_index].event_id) + }) + .count() + }; + assert_eq!(linked_to(0), 1, "legacy advancing token_count"); + assert_eq!(linked_to(1), 1, "record-era paired token_count"); + assert!(scan.quota_observations.iter().all(|record| { + let positive = record + .observation + .usage_sample + .as_ref() + .is_some_and(|usage| usage.computed_total() > 0); + positive == record.observation.usage_event_id.is_some() + })); +} + +#[test] +fn codex_subagent_usage_record_fixture_sums_records_and_keeps_parent_ids() { + let root = PathBuf::from(env!("CARGO_MANIFEST_DIR")) + .join("tests/fixtures/codex/usage-record/subagent"); + let path = root.join("sessions/2026/09/01/rollout-fixture-usage-record-subagent.jsonl"); + let records: Vec = std::fs::read_to_string(&path) + .expect("fixture") + .lines() + .filter_map(|line| { + let value: Value = serde_json::from_str(line).ok()?; + codex_token_usage_record_from_value(&value) + }) + .collect(); + assert_eq!(records.len(), 2); + assert_eq!( + records[0].parent_session_id(), + Some("00000000-0000-7000-8000-000000000001") + ); + assert_eq!( + records[0].root_turn_id.as_deref(), + Some("00000000-0000-7000-8000-0000000000a1") + ); + assert_eq!( + records[0].turn_id.as_deref(), + Some("00000000-0000-7000-8000-0000000000b2") + ); + assert_ne!(records[0].turn_id, records[0].root_turn_id); + assert_eq!( + records[0].response_id.as_deref(), + Some("resp_000000000000000000000000000000000000000000000000003") + ); + assert_eq!( + records[1].parent_session_id(), + Some("00000000-0000-7000-8000-000000000001") + ); + let record_total: u64 = records + .iter() + .map(|record| record.usage.computed_total()) + .sum(); + + let scan = scan_fixture("tests/fixtures/codex/usage-record/subagent"); + assert_eq!(scan.events.len(), 1); + assert_eq!(scan.events[0].usage.computed_total(), record_total); + assert_eq!(scan.events[0].usage.computed_total(), 120); + assert_eq!(scan.events[0].usage.input_tokens, Some(80)); + assert_eq!(scan.events[0].usage.cache_read_tokens, Some(20)); + assert_eq!(scan.events[0].usage.output_tokens, Some(20)); + assert_eq!(scan.events[0].usage.requests, Some(2)); + assert_eq!( + scan.events[0] + .model + .as_ref() + .and_then(|model| model.provider_model_id.as_deref()), + Some("gpt-5") + ); + + let mixed = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join( + "tests/fixtures/codex/usage-record/mixed/sessions/2026/09/01/rollout-fixture-usage-record.jsonl", + ); + let main_records: Vec = std::fs::read_to_string(mixed) + .expect("fixture") + .lines() + .filter_map(|line| { + let value: Value = serde_json::from_str(line).ok()?; + codex_token_usage_record_from_value(&value) + }) + .collect(); + assert!(main_records.len() >= 2); + assert!(main_records + .iter() + .all(|record| record.parent_session_id().is_none())); + assert!(main_records + .iter() + .all(|record| record.turn_id == record.root_turn_id)); +} diff --git a/crates/statsai-adapters/src/codex/tests/quota.rs b/crates/statsai-adapters/src/codex/tests/quota.rs index e340887..397d303 100644 --- a/crates/statsai-adapters/src/codex/tests/quota.rs +++ b/crates/statsai-adapters/src/codex/tests/quota.rs @@ -138,7 +138,7 @@ fn codex_quota_parser_requires_integer_reset_epochs_and_leniently_reads_balances } #[test] -fn codex_quota_links_consumed_samples_to_turn_events_and_preserves_zero_samples() { +fn codex_quota_links_consumed_samples_and_drops_usage_on_a_repeated_total() { let directory = tempfile::tempdir().expect("tempdir"); let root = directory.path().join("codex"); let sessions = root.join("sessions"); @@ -209,14 +209,12 @@ fn codex_quota_links_consumed_samples_to_turn_events_and_preserves_zero_samples( scan.quota_observations[0].observation.usage_event_id, Some(scan.events[0].event_id.clone()) ); - assert_eq!( - scan.quota_observations[1] - .observation - .usage_sample - .as_ref() - .map(UsageCounts::computed_total), - Some(0) - ); + // The second token_count repeats total_token_usage, so it is not usage. + // The observation is still recorded and stays unlinked. + assert!(scan.quota_observations[1] + .observation + .usage_sample + .is_none()); assert_eq!( scan.quota_observations[1].observation.usage_link_kind, QuotaUsageLinkKind::None diff --git a/crates/statsai-adapters/src/lib.rs b/crates/statsai-adapters/src/lib.rs index a781691..5924efb 100644 --- a/crates/statsai-adapters/src/lib.rs +++ b/crates/statsai-adapters/src/lib.rs @@ -83,8 +83,8 @@ pub(crate) use claude::claude_usage_counts_from_value; pub use claude::{ClaudeCodeAdapter, ClaudeProjectPathMemo}; pub use codex::CodexAdapter; pub(crate) use codex::{ - codex_project_context_from_value, codex_quota_observation, codex_usage_counts_from_value, - codex_usage_roots, + codex_project_context_from_value, codex_quota_observation, codex_token_count_usage, + codex_usage_roots, session_raw_from_value, }; pub(crate) use grok::grok_sessions_root; pub use grok::GrokBuildAdapter; diff --git a/crates/statsai-adapters/tests/fixtures/MANIFEST.json b/crates/statsai-adapters/tests/fixtures/MANIFEST.json index 1c10ea4..394d043 100644 --- a/crates/statsai-adapters/tests/fixtures/MANIFEST.json +++ b/crates/statsai-adapters/tests/fixtures/MANIFEST.json @@ -105,6 +105,16 @@ "path": "codex/reasoning/sessions/2026/08/17/rollout-fixture-reasoning.jsonl", "sha256": "bb86399f4fbaeb92c45c6b6833a81e65afc3702bad4591dc71f5481b2667ebe1" }, + { + "bytes": 7067, + "path": "codex/usage-record/mixed/sessions/2026/09/01/rollout-fixture-usage-record.jsonl", + "sha256": "0475c7ff12c76c2bcb6ba47f7721792c83879f650974b83dbcbc5ad88a4ac52a" + }, + { + "bytes": 3294, + "path": "codex/usage-record/subagent/sessions/2026/09/01/rollout-fixture-usage-record-subagent.jsonl", + "sha256": "35a9d4d0f9a9f8ec914ec960d58e5303c9d4fc4b61756cd50f473c434f0457f2" + }, { "bytes": 899, "path": "cursor/basic/usage-events-basic.csv", @@ -250,7 +260,11 @@ "basic": 10, "compaction": 43, "malformed": 10, - "reasoning": 10 + "reasoning": 10, + "usage-record": { + "mixed": 16, + "subagent": 7 + } }, "cursor": { "derived_from_local_store": false, diff --git a/crates/statsai-adapters/tests/fixtures/codex/usage-record/mixed/sessions/2026/09/01/rollout-fixture-usage-record.jsonl b/crates/statsai-adapters/tests/fixtures/codex/usage-record/mixed/sessions/2026/09/01/rollout-fixture-usage-record.jsonl new file mode 100644 index 0000000..7b5867d --- /dev/null +++ b/crates/statsai-adapters/tests/fixtures/codex/usage-record/mixed/sessions/2026/09/01/rollout-fixture-usage-record.jsonl @@ -0,0 +1,16 @@ +{"timestamp":"2026-09-01T10:00:00.000Z","ordinal":1,"type":"session_meta","payload":{"id":"00000000-0000-7000-8000-000000000001","timestamp":"2026-09-01T10:00:00.000Z","cwd":"/workspace/sample-project","cli_version":"0.150.0","source":"cli","model_provider":"openai"}} +{"timestamp":"2026-09-01T10:00:00.100Z","ordinal":2,"type":"turn_context","payload":{"model":"gpt-5","turn_id":"00000000-0000-7000-8000-0000000000a1"}} +{"timestamp":"2026-09-01T10:00:00.200Z","ordinal":3,"type":"event_msg","payload":{"type":"task_started","started_at":"2026-09-01T10:00:00.200Z"}} +{"timestamp":"2026-09-01T10:00:01.000Z","ordinal":4,"type":"event_msg","payload":{"type":"token_count","info":{"total_token_usage":{"input_tokens":100,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":110},"last_token_usage":{"input_tokens":100,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":110},"model_context_window":258400},"rate_limits":{"limit_id":"codex","limit_name":null,"primary":{"used_percent":10,"window_minutes":300,"resets_at":1790000000},"secondary":{"used_percent":5,"window_minutes":10080,"resets_at":1790400000},"credits":{"has_credits":false,"unlimited":false,"balance":"0"},"plan_type":"plus","rate_limit_reached_type":null}}} +{"timestamp":"2026-09-01T10:00:01.200Z","ordinal":5,"type":"event_msg","payload":{"type":"token_count","info":{"total_token_usage":{"input_tokens":100,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":110},"last_token_usage":{"input_tokens":100,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":110},"model_context_window":258400},"rate_limits":{"limit_id":"codex","limit_name":null,"primary":{"used_percent":10,"window_minutes":300,"resets_at":1790000000},"secondary":{"used_percent":5,"window_minutes":10080,"resets_at":1790400000},"credits":{"has_credits":false,"unlimited":false,"balance":"0"},"plan_type":"plus","rate_limit_reached_type":null}}} +{"timestamp":"2026-09-01T10:00:01.400Z","ordinal":6,"type":"event_msg","payload":{"type":"token_count","info":{"total_token_usage":{"input_tokens":100,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":110},"last_token_usage":{"input_tokens":100,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":110},"model_context_window":258400},"rate_limits":{"limit_id":"premium","limit_name":null,"primary":{"used_percent":10,"window_minutes":300,"resets_at":1790000000},"secondary":{"used_percent":5,"window_minutes":10080,"resets_at":1790400000},"credits":{"has_credits":false,"unlimited":false,"balance":"0"},"plan_type":"plus","rate_limit_reached_type":null}}} +{"timestamp":"2026-09-01T10:00:02.000Z","ordinal":7,"type":"event_msg","payload":{"type":"task_complete","completed_at":"2026-09-01T10:00:02.000Z"}} +{"timestamp":"2026-09-01T10:01:00.000Z","ordinal":8,"type":"event_msg","payload":{"type":"task_started","started_at":"2026-09-01T10:01:00.000Z"}} +{"timestamp":"2026-09-01T10:01:01.000Z","ordinal":9,"type":"token_usage_record","payload":{"thread_id":"00000000-0000-7000-8000-000000000001","turn_id":"00000000-0000-7000-8000-0000000000a2","session_id":"00000000-0000-7000-8000-000000000001","root_turn_id":"00000000-0000-7000-8000-0000000000a2","response_id":"resp_000000000000000000000000000000000000000000000000001","usage":{"input_tokens":200,"cached_input_tokens":50,"cache_write_input_tokens":0,"output_tokens":20,"reasoning_output_tokens":5,"total_tokens":220},"turn_token_usage":{"input_tokens":5000,"cached_input_tokens":1000,"cache_write_input_tokens":0,"output_tokens":400,"reasoning_output_tokens":50,"total_tokens":5400},"thread_token_usage":{"input_tokens":9000,"cached_input_tokens":2000,"cache_write_input_tokens":0,"output_tokens":800,"reasoning_output_tokens":100,"total_tokens":9800}}} +{"timestamp":"2026-09-01T10:01:01.100Z","ordinal":10,"type":"event_msg","payload":{"type":"token_count","info":{"total_token_usage":{"input_tokens":300,"cached_input_tokens":50,"cache_write_input_tokens":0,"output_tokens":30,"reasoning_output_tokens":5,"total_tokens":330},"last_token_usage":{"input_tokens":200,"cached_input_tokens":50,"cache_write_input_tokens":0,"output_tokens":20,"reasoning_output_tokens":5,"total_tokens":220},"model_context_window":258400},"rate_limits":{"limit_id":"codex","limit_name":null,"primary":{"used_percent":10,"window_minutes":300,"resets_at":1790000000},"secondary":{"used_percent":5,"window_minutes":10080,"resets_at":1790400000},"credits":{"has_credits":false,"unlimited":false,"balance":"0"},"plan_type":"plus","rate_limit_reached_type":null}}} +{"timestamp":"2026-09-01T10:01:02.000Z","ordinal":11,"type":"token_usage_record","payload":{"thread_id":"00000000-0000-7000-8000-000000000001","turn_id":"00000000-0000-7000-8000-0000000000a2","session_id":"00000000-0000-7000-8000-000000000001","root_turn_id":"00000000-0000-7000-8000-0000000000a2","response_id":"resp_000000000000000000000000000000000000000000000000002","usage":{"input_tokens":800,"cached_input_tokens":300,"cache_write_input_tokens":0,"output_tokens":40,"reasoning_output_tokens":10,"total_tokens":840},"turn_token_usage":{"input_tokens":5000,"cached_input_tokens":1000,"cache_write_input_tokens":0,"output_tokens":400,"reasoning_output_tokens":50,"total_tokens":5400},"thread_token_usage":{"input_tokens":9000,"cached_input_tokens":2000,"cache_write_input_tokens":0,"output_tokens":800,"reasoning_output_tokens":100,"total_tokens":9800}}} +{"timestamp":"2026-09-01T10:01:02.100Z","ordinal":12,"type":"compacted","payload":{"message":"synthetic compaction marker","window_number":1}} +{"timestamp":"2026-09-01T10:01:02.200Z","ordinal":13,"type":"world_state","payload":{"full":true,"state":{"current_date":"2026-09-01"}}} +{"timestamp":"2026-09-01T10:01:02.300Z","ordinal":14,"type":"turn_context","payload":{"model":"gpt-5","turn_id":"00000000-0000-7000-8000-0000000000a2"}} +{"timestamp":"2026-09-01T10:01:02.400Z","ordinal":15,"type":"event_msg","payload":{"type":"token_count","info":{"total_token_usage":{"input_tokens":300,"cached_input_tokens":50,"cache_write_input_tokens":0,"output_tokens":30,"reasoning_output_tokens":5,"total_tokens":330},"last_token_usage":{"input_tokens":4,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":1,"reasoning_output_tokens":0,"total_tokens":5},"model_context_window":258400},"rate_limits":{"limit_id":"codex","limit_name":null,"primary":{"used_percent":10,"window_minutes":300,"resets_at":1790000000},"secondary":{"used_percent":5,"window_minutes":10080,"resets_at":1790400000},"credits":{"has_credits":false,"unlimited":false,"balance":"0"},"plan_type":"plus","rate_limit_reached_type":null}}} +{"timestamp":"2026-09-01T10:01:03.000Z","ordinal":16,"type":"event_msg","payload":{"type":"task_complete","completed_at":"2026-09-01T10:01:03.000Z"}} diff --git a/crates/statsai-adapters/tests/fixtures/codex/usage-record/subagent/sessions/2026/09/01/rollout-fixture-usage-record-subagent.jsonl b/crates/statsai-adapters/tests/fixtures/codex/usage-record/subagent/sessions/2026/09/01/rollout-fixture-usage-record-subagent.jsonl new file mode 100644 index 0000000..4411428 --- /dev/null +++ b/crates/statsai-adapters/tests/fixtures/codex/usage-record/subagent/sessions/2026/09/01/rollout-fixture-usage-record-subagent.jsonl @@ -0,0 +1,7 @@ +{"timestamp":"2026-09-01T11:00:00.000Z","ordinal":1,"type":"session_meta","payload":{"id":"00000000-0000-7000-8000-0000000000b1","timestamp":"2026-09-01T11:00:00.000Z","cwd":"/workspace/sample-project","cli_version":"0.150.0","source":{"subagent":{"thread_spawn":{"parent_thread_id":"00000000-0000-7000-8000-000000000001","depth":1}}},"model_provider":"openai"}} +{"timestamp":"2026-09-01T11:00:00.100Z","ordinal":2,"type":"turn_context","payload":{"model":"gpt-5","turn_id":"00000000-0000-7000-8000-0000000000b2"}} +{"timestamp":"2026-09-01T11:00:00.200Z","ordinal":3,"type":"event_msg","payload":{"type":"task_started","started_at":"2026-09-01T11:00:00.200Z"}} +{"timestamp":"2026-09-01T11:00:01.000Z","ordinal":4,"type":"token_usage_record","payload":{"thread_id":"00000000-0000-7000-8000-0000000000b1","turn_id":"00000000-0000-7000-8000-0000000000b2","session_id":"00000000-0000-7000-8000-000000000001","root_turn_id":"00000000-0000-7000-8000-0000000000a1","response_id":"resp_000000000000000000000000000000000000000000000000003","usage":{"input_tokens":40,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":50},"turn_token_usage":{"input_tokens":5000,"cached_input_tokens":1000,"cache_write_input_tokens":0,"output_tokens":400,"reasoning_output_tokens":50,"total_tokens":5400},"thread_token_usage":{"input_tokens":9000,"cached_input_tokens":2000,"cache_write_input_tokens":0,"output_tokens":800,"reasoning_output_tokens":100,"total_tokens":9800}}} +{"timestamp":"2026-09-01T11:00:01.100Z","ordinal":5,"type":"event_msg","payload":{"type":"token_count","info":{"total_token_usage":{"input_tokens":40,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":50},"last_token_usage":{"input_tokens":40,"cached_input_tokens":0,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":50},"model_context_window":258400},"rate_limits":{"limit_id":"codex","limit_name":null,"primary":{"used_percent":10,"window_minutes":300,"resets_at":1790000000},"secondary":{"used_percent":5,"window_minutes":10080,"resets_at":1790400000},"credits":{"has_credits":false,"unlimited":false,"balance":"0"},"plan_type":"plus","rate_limit_reached_type":null}}} +{"timestamp":"2026-09-01T11:00:02.000Z","ordinal":6,"type":"token_usage_record","payload":{"thread_id":"00000000-0000-7000-8000-0000000000b1","turn_id":"00000000-0000-7000-8000-0000000000b2","session_id":"00000000-0000-7000-8000-000000000001","root_turn_id":"00000000-0000-7000-8000-0000000000a1","response_id":"resp_000000000000000000000000000000000000000000000000004","usage":{"input_tokens":60,"cached_input_tokens":20,"cache_write_input_tokens":0,"output_tokens":10,"reasoning_output_tokens":0,"total_tokens":70},"turn_token_usage":{"input_tokens":5000,"cached_input_tokens":1000,"cache_write_input_tokens":0,"output_tokens":400,"reasoning_output_tokens":50,"total_tokens":5400},"thread_token_usage":{"input_tokens":9000,"cached_input_tokens":2000,"cache_write_input_tokens":0,"output_tokens":800,"reasoning_output_tokens":100,"total_tokens":9800}}} +{"timestamp":"2026-09-01T11:00:03.000Z","ordinal":7,"type":"event_msg","payload":{"type":"task_complete","completed_at":"2026-09-01T11:00:03.000Z"}}