Skip to content
32 changes: 16 additions & 16 deletions crates/statsai-adapters/src/archive/scan/codex.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<UsageCounts> = 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;
Expand All @@ -60,26 +63,23 @@ pub(crate) fn collect_codex_quota_observations(
let Ok(value) = serde_json::from_str::<Value>(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)
Expand Down
6 changes: 3 additions & 3 deletions crates/statsai-adapters/src/archive/scan/mod.rs
Original file line number Diff line number Diff line change
@@ -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,
};
Expand All @@ -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;
Expand Down
7 changes: 4 additions & 3 deletions crates/statsai-adapters/src/cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
15 changes: 15 additions & 0 deletions crates/statsai-adapters/src/codex/parse/line.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand All @@ -42,6 +46,8 @@ pub(crate) enum CodexLineKind {
ResponseItemMessage,
EventUserMessage,
TokenCount,
TokenUsageRecord,
Compacted,
TaskStarted,
TaskComplete,
HeadlessUsage,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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\":")
Expand Down
111 changes: 87 additions & 24 deletions crates/statsai-adapters/src/codex/parse/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<UsageCounts> = None;
// Both are per session: one file can interleave several sessions.
let mut previous_totals: HashMap<String, Option<CodexCumulativeTotal>> = 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<String, (UsageCounts, usize)> = HashMap::new();
let mut paired_quota_lines: HashMap<usize, usize> = HashMap::new();
let mut current_model: Option<String> = None;
let mut current_reasoning = ModelReasoningState::default();
let mut current_model_is_fallback = false;
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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();
Expand All @@ -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))
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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() {
Expand Down Expand Up @@ -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);
}
}

Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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"
Expand All @@ -898,18 +958,21 @@ 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
},
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,
&quota_observation_indices,
&[record.line_number],
&linked_quota_lines,
&event.event_id,
QuotaUsageLinkKind::RecordEvent,
);
Expand Down
5 changes: 5 additions & 0 deletions crates/statsai-adapters/src/codex/parse/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ pub(crate) struct CodexLineRecord {
pub(crate) model_explicit: bool,
pub(crate) usage: Option<UsageCounts>,
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<String>,
Expand Down Expand Up @@ -108,5 +109,9 @@ pub(crate) struct ActiveCodexTurn {
pub(crate) prompt_previews: Vec<CodexPromptPreviewCandidate>,
pub(crate) last_activity_at: DateTime<Utc>,
pub(crate) usage_lines: Vec<usize>,
/// `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<usize>,
pub(crate) project: Option<ProjectInfo>,
}
Loading
Loading