diff --git a/docs/PROVIDERS.md b/docs/PROVIDERS.md index 237f7dab3d..38cd3327c4 100644 --- a/docs/PROVIDERS.md +++ b/docs/PROVIDERS.md @@ -87,6 +87,8 @@ Custom pricing overlays are exact-match overrides used only where the local spen OpenCode-held OpenAI/Codex OAuth can be reused for **remote Codex account quota** only when the Codex provider's `External OAuth sources` setting is explicitly enabled. Native Codex credentials still take precedence, an explicit `CODEX_HOME` stays isolated, and external credentials remain read-only. This does **not** import ordinary OpenCode sessions into Codex token or spend totals. OpenCode Go's local SQLite reader remains scoped to its own `opencode-go` assistant records; OpenAI API-platform usage is a separate provider. +Codex local cost prices **Priority (Fast) turns** at the Fast rate. A turn counts as Priority when `/logs_2.sqlite` (Codex's trace database) holds a `response.create` websocket request with `service_tier == "priority"` for that turn's id. The database is opened read-only, scanned incrementally with a persisted cursor in the cost cache, and only turn ids, model names, and timestamps are kept; row bodies contain prompts and are never stored or logged. A missing or unreadable database keeps Standard pricing, and turn evidence never crosses `CODEX_HOME` scopes. Models without a Fast lane stay Standard. Older cost caches rebuild once (Codex cache schema v4 records each row's turn id). + ### z.ai Coding Plan quotas z.ai Coding Plans accept both `TOKENS_LIMIT` and `CREDIT_LIMIT` rows. The shortest known Coding Plan window becomes primary and the longest becomes secondary; `TIME_LIMIT` is the separate MCP lane. When absolute usage/remaining counts are available they determine the used percentage, otherwise the provider percentage is used, always clamped to 0–100%. This behavior is shared by the tray, provider detail, CLI, and other Windows surfaces. diff --git a/rust/src/codex_costs.rs b/rust/src/codex_costs.rs index cf80594c02..cabbd12435 100644 --- a/rust/src/codex_costs.rs +++ b/rust/src/codex_costs.rs @@ -361,27 +361,17 @@ fn add_codex_tokens_to_summary( return Some(0.0); } - let priced = pricing_day - .and_then(|day| { - CostUsagePricing::codex_cost_usd_at_date( - &model_key, - tokens.input, - tokens.cached, - tokens.output, - day, - ) - }) - .or_else(|| { - CostUsagePricing::codex_cost_usd(&model_key, tokens.input, tokens.cached, tokens.output) - }); - let uses_fallback_pricing = priced.is_none(); - let cost = codex_cost_usd_for_day( + let priced = CostUsagePricing::codex_day_aggregate_cost_usd( &model_key, tokens.input, tokens.cached, tokens.output, pricing_day, ); + let uses_fallback_pricing = priced.is_none(); + let cost = priced.unwrap_or_else(|| { + codex_cost_usd_fallback(&model_key, tokens.input, tokens.cached, tokens.output) + }); if uses_fallback_pricing { summary.unknown_models.insert(model_key.clone()); match &mut summary.model_pricing_completeness { @@ -484,26 +474,8 @@ fn codex_cost_usd_for_day( if CostUsagePricing::is_codex_unattributed_model(model) { return 0.0; } - let priced = pricing_day - .and_then(|day| CostUsagePricing::codex_cost_usd_at_date(model, input, cached, output, day)) - .or_else(|| CostUsagePricing::codex_cost_usd(model, input, cached, output)); - if let Some(cost) = priced { - return cost; - } - - let normalized = CostUsagePricing::normalize_codex_model(model); - if normalized.contains("fast") || normalized.contains("priority") { - let fast = pricing_day - .and_then(|day| { - CostUsagePricing::codex_fast_cost_usd_at_date(model, input, cached, output, day) - }) - .or_else(|| CostUsagePricing::codex_fast_cost_usd(model, input, cached, output)); - if let Some(cost) = fast { - return cost; - } - } - - codex_cost_usd_fallback(model, input, cached, output) + CostUsagePricing::codex_day_aggregate_cost_usd(model, input, cached, output, pricing_day) + .unwrap_or_else(|| codex_cost_usd_fallback(model, input, cached, output)) } fn codex_cost_usd_fallback(model: &str, input: u64, cached: u64, output: u64) -> f64 { @@ -581,6 +553,7 @@ mod tests { cached: 0, output: 20, reasoning, + turn_id: None, }; let mut known_summary = CostSummary::default(); @@ -612,6 +585,7 @@ mod tests { cached: 0, output: 20, reasoning, + turn_id: None, }; let records = vec![ (make_record(Some(7)), 0), @@ -675,6 +649,7 @@ mod tests { cached: 0, output: 0, reasoning: None, + turn_id: None, }, 0, ), @@ -687,6 +662,7 @@ mod tests { cached: 0, output: 0, reasoning: None, + turn_id: None, }, 0, ), @@ -699,6 +675,7 @@ mod tests { cached: 0, output: 0, reasoning: None, + turn_id: None, }, 0, ), @@ -750,6 +727,7 @@ mod tests { cached: 0, output: 5, reasoning: None, + turn_id: None, }, 0, ), @@ -762,6 +740,7 @@ mod tests { cached: 0, output: 1_000_000, reasoning: None, + turn_id: None, }, 0, ), @@ -789,6 +768,7 @@ mod tests { cached: 0, output: 1, reasoning: None, + turn_id: None, }, 0, )]; @@ -810,6 +790,7 @@ mod tests { cached: 0, output: 0, reasoning: None, + turn_id: None, }, 0, )]; @@ -843,6 +824,7 @@ mod tests { cached: 0, output: 1_000_000, reasoning: None, + turn_id: None, }, 0, )]; diff --git a/rust/src/codex_costs/quota_windows.rs b/rust/src/codex_costs/quota_windows.rs index 7119676488..65c8990c65 100644 --- a/rust/src/codex_costs/quota_windows.rs +++ b/rust/src/codex_costs/quota_windows.rs @@ -7,9 +7,11 @@ use chrono::{DateTime, Duration, Local, NaiveDate, TimeZone, Utc}; use serde::{Deserialize, Serialize}; use std::collections::HashSet; +use std::path::Path; use crate::core::{ - CodexSourceRowCache, CodexSourceUsageRow, CostUsageCache, CostUsagePricing, RateWindow, + CodexPriorityOverlay, CodexSourceRowCache, CodexSourceUsageRow, CostUsageCache, + CostUsagePricing, RateWindow, row_priced_model, }; const NOMINAL_WEEK_MINUTES: i64 = 7 * 24 * 60; @@ -341,7 +343,9 @@ fn cache_slices(cache: &CostUsageCache) -> Vec { }); let mut identities = HashSet::new(); let mut slices = Vec::new(); + let cursor = cache.codex_priority_turns_cursor.as_ref(); for (path, source) in sources { + let overlay = cursor.and_then(|cursor| cursor.overlay_for_file(Path::new(path))); let identity = if source.file_identity.is_empty() { path.as_str() } else { @@ -351,14 +355,14 @@ fn cache_slices(cache: &CostUsageCache) -> Vec { continue; } for row in &source.rows { - slices.push(slice_from_row(row)); + slices.push(slice_from_row(row, overlay.as_ref())); } } slices.sort_by_key(|slice| (slice.start, slice.end)); slices } -fn slice_from_row(row: &CodexSourceUsageRow) -> Slice { +fn slice_from_row(row: &CodexSourceUsageRow, overlay: Option<&CodexPriorityOverlay<'_>>) -> Slice { let timestamp = row.timestamp.or_else(|| local_day_start(&row.day_key)); let end = row.timestamp.map(|_| None).unwrap_or_else(|| { local_day_start(&row.day_key).and_then(|start| start.checked_add_signed(Duration::days(1))) @@ -366,22 +370,14 @@ fn slice_from_row(row: &CodexSourceUsageRow) -> Slice { let input = u64::try_from(row.input.max(0)).unwrap_or(0); let output = u64::try_from(row.output.max(0)).unwrap_or(0); let tokens = Some(input.saturating_add(output)); - let cost_usd = row.pricing.pricing_model.as_deref().and_then(|model| { - let model = if row.pricing.pricing_mode.as_deref() == Some("priority") - && !model.ends_with("-priority") - { - format!("{model}-priority") - } else { - model.to_string() - }; + let cost_usd = row_priced_model(row, overlay).and_then(|model| { let date = timestamp.map(|value| value.with_timezone(&Local).date_naive())?; - CostUsagePricing::codex_cost_usd_at_date( - &model, - input, - u64::try_from(row.cached.max(0)).unwrap_or(0).min(input), - output, - date, - ) + let cached = u64::try_from(row.cached.max(0)).unwrap_or(0).min(input); + if model.ends_with("-priority") { + CostUsagePricing::codex_fast_cost_usd_at_date(&model, input, cached, output, date) + } else { + CostUsagePricing::codex_cost_usd_at_date(&model, input, cached, output, date) + } }); Slice { start: timestamp.unwrap_or(DateTime::::UNIX_EPOCH), @@ -418,14 +414,16 @@ fn legacy_day_slices(cache: &CostUsageCache) -> Vec { let cached = u64::try_from(packed.get(1).copied().unwrap_or(0).max(0)) .unwrap_or(0) .min(input); - let cost_usd = CostUsagePricing::codex_cost_usd_at_date( + let cost_usd = CostUsagePricing::codex_day_aggregate_cost_usd( model, input, cached, output, - NaiveDate::parse_from_str(day, "%Y-%m-%d") - .ok() - .unwrap_or_else(|| start.with_timezone(&Local).date_naive()), + Some( + NaiveDate::parse_from_str(day, "%Y-%m-%d") + .ok() + .unwrap_or_else(|| start.with_timezone(&Local).date_naive()), + ), ); slices.push(Slice { start, diff --git a/rust/src/codex_costs/quota_windows/tests.rs b/rust/src/codex_costs/quota_windows/tests.rs index f4e0b5e91b..004368ac00 100644 --- a/rust/src/codex_costs/quota_windows/tests.rs +++ b/rust/src/codex_costs/quota_windows/tests.rs @@ -34,6 +34,7 @@ fn row( output, reasoning: None, source_end_offset: 1, + turn_id: None, pricing: CodexSourcePricingEvidence { pricing_model: pricing_model.map(str::to_string), pricing_mode: None, diff --git a/rust/src/core/cost_pricing/codex.rs b/rust/src/core/cost_pricing/codex.rs index 8b9a20baaa..fe0512ba89 100644 --- a/rust/src/core/cost_pricing/codex.rs +++ b/rust/src/core/cost_pricing/codex.rs @@ -1,3 +1,5 @@ +use chrono::NaiveDate; + use super::super::{codex_routed_pricing, models_dev_pricing}; use super::{CODEX_PRICING, CostUsagePricing}; @@ -60,6 +62,107 @@ pub(super) fn codex_fast_allows_long_context(model: &str) -> bool { } impl CostUsagePricing { + /// Whether one request of `input_tokens` can run in the Fast lane of + /// `model`. Older models offer no Fast lane above the long-context + /// threshold, so upstream charges such a Priority request the Standard + /// cost; Astra publishes long-context Fast rates. + pub fn codex_fast_lane_covers(model: &str, input_tokens: u64) -> bool { + Self::codex_api_fast_multiplier(model).is_some() + && (input_tokens <= CODEX_LONG_CONTEXT_THRESHOLD + || codex_fast_allows_long_context(model)) + } + + /// Fast cost in USD of a day aggregate under a Fast key (`-priority` or + /// `-fast`), or `None` when `model` names no Fast lane. + /// + /// Upstream prices every request on its own. A day aggregate sums + /// requests that each ran in the Fast lane, so the summed input must + /// neither refuse the surcharge nor switch to long-context rates: older + /// models price at the base model's short-context rates times the + /// multiplier. Astra's Fast lane has long-context rates, so its + /// aggregate keeps the whole-aggregate rule that Standard aggregates use. + pub fn codex_fast_aggregate_cost_usd( + model: &str, + input_tokens: u64, + cached_input_tokens: u64, + output_tokens: u64, + pricing_date: Option, + ) -> Option { + let base = Self::codex_fast_base_model(model); + if base == Self::normalize_codex_model(model) { + return None; + } + if codex_fast_allows_long_context(model) { + return pricing_date + .and_then(|date| { + Self::codex_fast_cost_usd_at_date( + model, + input_tokens, + cached_input_tokens, + output_tokens, + date, + ) + }) + .or_else(|| { + Self::codex_fast_cost_usd( + model, + input_tokens, + cached_input_tokens, + output_tokens, + ) + }); + } + let multiplier = Self::codex_api_fast_multiplier(model)?; + let (input_rate, cache_read_rate, output_rate) = + Self::codex_short_context_rates(&base, pricing_date)?; + Some( + codex_cost_from_rates( + input_tokens, + cached_input_tokens, + output_tokens, + input_rate, + cache_read_rate, + output_rate, + ) * multiplier, + ) + } + + /// Known cost in USD of one Codex day aggregate: a Fast key prices + /// through its base model's Fast lane + /// ([`Self::codex_fast_aggregate_cost_usd`]), any other model at the + /// rates in effect on `pricing_date`. `None` means no rate is known. + pub fn codex_day_aggregate_cost_usd( + model: &str, + input_tokens: u64, + cached_input_tokens: u64, + output_tokens: u64, + pricing_date: Option, + ) -> Option { + let (input, cached, output) = (input_tokens, cached_input_tokens, output_tokens); + Self::codex_fast_aggregate_cost_usd(model, input, cached, output, pricing_date) + .or_else(|| { + pricing_date.and_then(|date| { + Self::codex_cost_usd_at_date(model, input, cached, output, date) + }) + }) + .or_else(|| Self::codex_cost_usd(model, input, cached, output)) + } + + /// Short-context `(input, cache read, output)` rates of `model` on + /// `pricing_date`. Short-context pricing is linear per token, so + /// one-token probes read the exact dated rates back. + fn codex_short_context_rates( + model: &str, + pricing_date: Option, + ) -> Option<(f64, f64, f64)> { + let cost = |input, cached, output| { + pricing_date + .and_then(|date| Self::codex_cost_usd_at_date(model, input, cached, output, date)) + .or_else(|| Self::codex_cost_usd(model, input, cached, output)) + }; + Some((cost(1, 0, 0)?, cost(1, 1, 0)?, cost(0, 0, 1)?)) + } + /// Calculate Codex cost in USD when input includes cache-write tokens. pub fn codex_cost_usd_with_cache_write( model: &str, diff --git a/rust/src/core/cost_pricing_tests.rs b/rust/src/core/cost_pricing_tests.rs index 9984061644..c506ac67a1 100644 --- a/rust/src/core/cost_pricing_tests.rs +++ b/rust/src/core/cost_pricing_tests.rs @@ -513,3 +513,99 @@ fn gpt6_astra_unknown_models_fail_closed() { assert!(CostUsagePricing::codex_cost_usd(model, 1000, 0, 100).is_none()); } } + +// Upstream 0.65.0 #3820: Priority rows and day aggregates. + +#[test] +fn codex_fast_lane_covers_requests_up_to_the_long_context_threshold() { + assert!(CostUsagePricing::codex_fast_lane_covers("gpt-5.5", 272_000)); + assert!(!CostUsagePricing::codex_fast_lane_covers( + "gpt-5.5", 272_001 + )); + assert!(CostUsagePricing::codex_fast_lane_covers( + "gpt-5.5-priority", + 1_000 + )); + // Astra publishes long-context Fast rates. + assert!(CostUsagePricing::codex_fast_lane_covers( + "gpt-6-astra-priority", + 272_001 + )); + // gpt-5 has no Fast lane at any size. + assert!(!CostUsagePricing::codex_fast_lane_covers( + "gpt-5-priority", + 1_000 + )); +} + +#[test] +fn codex_priority_day_aggregate_prices_short_rates_times_multiplier() { + // Two 200k-input Fast requests sum to 400k. Each ran in the Fast lane, + // so the aggregate neither drops the surcharge nor switches tiers. + let one = CostUsagePricing::codex_cost_usd("gpt-5.5", 200_000, 50_000, 1_000).unwrap(); + let aggregate = CostUsagePricing::codex_day_aggregate_cost_usd( + "gpt-5.5-priority", + 400_000, + 100_000, + 2_000, + None, + ) + .unwrap(); + assert!((aggregate - 2.0 * one * 2.5).abs() < 1e-10); + assert_eq!( + CostUsagePricing::codex_fast_cost_usd("gpt-5.5-priority", 400_000, 100_000, 2_000), + None + ); +} + +#[test] +fn codex_astra_priority_day_aggregate_keeps_long_context_fast_rates() { + let aggregate = CostUsagePricing::codex_day_aggregate_cost_usd( + "gpt-6-astra-priority", + 300_000, + 0, + 1_000, + None, + ) + .unwrap(); + let fast = + CostUsagePricing::codex_fast_cost_usd("gpt-6-astra-priority", 300_000, 0, 1_000).unwrap(); + assert!((aggregate - fast).abs() < 1e-12); + let expected = (300_000.0 * 2e-5 + 1_000.0 * 7.5e-5) * 2.0; + assert!((aggregate - expected).abs() < 1e-10); +} + +#[test] +fn codex_standard_day_aggregate_uses_standard_pricing() { + let aggregate = + CostUsagePricing::codex_day_aggregate_cost_usd("gpt-5.5", 1_000, 200, 100, None).unwrap(); + let standard = CostUsagePricing::codex_cost_usd("gpt-5.5", 1_000, 200, 100).unwrap(); + assert!((aggregate - standard).abs() < 1e-15); +} + +#[test] +fn codex_priority_day_aggregate_uses_the_rates_of_its_day() { + use chrono::NaiveDate; + + let before_cut = NaiveDate::from_ymd_opt(2026, 7, 1).unwrap(); + let aggregate = CostUsagePricing::codex_day_aggregate_cost_usd( + "gpt-5.6-terra-priority", + 300_000, + 0, + 1_000, + Some(before_cut), + ) + .unwrap(); + // Pre-cut short-context Terra rates, times the 2.0 Fast multiplier. + let expected = (300_000.0 * 2.5e-6 + 1_000.0 * 1.5e-5) * 2.0; + assert!((aggregate - expected).abs() < 1e-10); + assert!((aggregate - 1.53).abs() < 1e-9); +} + +#[test] +fn codex_priority_day_aggregate_without_a_fast_lane_is_unpriced() { + assert_eq!( + CostUsagePricing::codex_day_aggregate_cost_usd("gpt-5-priority", 1_000, 0, 100, None), + None + ); +} diff --git a/rust/src/core/jsonl_scanner.rs b/rust/src/core/jsonl_scanner.rs index 883a31d89f..f0c63e0729 100755 --- a/rust/src/core/jsonl_scanner.rs +++ b/rust/src/core/jsonl_scanner.rs @@ -239,6 +239,22 @@ pub struct CostUsageCache { /// caches remain valid and can be upgraded lazily. #[serde(default, skip_serializing_if = "HashMap::is_empty")] pub codex_source_rows: HashMap, + /// Request rows of fork-shaped Codex files (a `forked_from_id` or a + /// parent-baseline lineage), built from the same parsed records as the + /// file's day totals. Source-row evidence skips these files, so this map + /// is what lets Priority trace evidence reach forked sessions. + #[serde(default, skip_serializing_if = "HashMap::is_empty")] + pub codex_fork_rows: HashMap>, + /// Priority (Fast) turn evidence read from the Codex trace database. + /// Applied as a pricing overlay whenever day totals are rebuilt. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub codex_priority_turns_cursor: Option, + /// Trace-database path and presence seen by the last full scan + /// (`sqlite:` or `missing:`, upstream + /// `codexPriorityMetadataKey`). A database that appears later bypasses + /// the scan debounce once so its evidence is applied promptly. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub codex_priority_metadata_key: Option, /// Content stamp of the decoded on-disk baseline. This is process-local /// and omitted from JSON so a stale reader cannot replace a newer cache. #[serde(skip)] @@ -273,6 +289,10 @@ pub struct CodexSourceUsageRow { pub source_end_offset: i64, #[serde(default)] pub pricing: CodexSourcePricingEvidence, + /// Codex turn (`task_started`) the request belongs to; matches Priority + /// trace evidence. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub turn_id: Option, } /// Source identity and rows retained for a cached Codex file. @@ -468,6 +488,7 @@ pub struct CodexUsageRecord { pub cached: i64, pub output: i64, pub reasoning: Option, + pub turn_id: Option, } /// Day range for scanning @@ -507,8 +528,14 @@ impl CostUsageDayRange { /// JSONL Scanner for cost/usage logs pub struct JsonlScanner; pub(crate) mod codex; +pub(crate) use codex::priority::CodexPriorityOverlay; +pub use codex::priority::{ + CODEX_PRIORITY_COMPLETED_MODEL_RETENTION_LIMIT, CodexPriorityCursorAnchor, + CodexPriorityTurnMetadata, CodexPriorityTurnsCursor, +}; pub(crate) use codex::source_rows::{ read_source_rows, recover_rows, row_cache, row_cache_matches, row_cache_needs_recovery, + row_priced_model, rows_from_records, }; impl JsonlScanner { @@ -668,24 +695,13 @@ impl JsonlScanner { if !CostUsagePricing::counts_toward_codex_subscription(model) { continue; } - let priced = pricing_day - .and_then(|day| { - CostUsagePricing::codex_cost_usd_at_date( - model, - u64::try_from(input).unwrap_or(0), - u64::try_from(cached).unwrap_or(0), - u64::try_from(output).unwrap_or(0), - day, - ) - }) - .or_else(|| { - CostUsagePricing::codex_cost_usd( - model, - u64::try_from(input).unwrap_or(0), - u64::try_from(cached).unwrap_or(0), - u64::try_from(output).unwrap_or(0), - ) - }); + let priced = CostUsagePricing::codex_day_aggregate_cost_usd( + model, + u64::try_from(input).unwrap_or(0), + u64::try_from(cached).unwrap_or(0), + u64::try_from(output).unwrap_or(0), + pricing_day, + ); if let Some(cost) = priced { total_cost_usd += cost; } else { diff --git a/rust/src/core/jsonl_scanner/codex.rs b/rust/src/core/jsonl_scanner/codex.rs index 30f159c0ca..ff66c711bc 100644 --- a/rust/src/core/jsonl_scanner/codex.rs +++ b/rust/src/core/jsonl_scanner/codex.rs @@ -2,6 +2,7 @@ use super::*; mod helpers; mod parser; +pub(crate) mod priority; pub(crate) mod source_rows; use helpers::{ @@ -13,8 +14,10 @@ use parser::CodexParserState; /// Persisted Codex cache schema version. Version 0 predates 64-bit totals; /// version 1 can retain a terminal pause after treating a paginated v2 /// subagent's independent counters as an inherited fork. Version 3 adds -/// persisted paginated-fork accounting state. Rebuild older artifacts. -pub(crate) const CODEX_CACHE_SCHEMA_VERSION: u32 = 3; +/// persisted paginated-fork accounting state. Version 4 records the Codex +/// turn id on source rows so Priority trace evidence can be matched. +/// Rebuild older artifacts. +pub(crate) const CODEX_CACHE_SCHEMA_VERSION: u32 = 4; /// Whether a persisted Codex cache artifact matches the current schema. /// A mismatched artifact (e.g. a pre-64-bit cache from an older release) is diff --git a/rust/src/core/jsonl_scanner/codex/helpers.rs b/rust/src/core/jsonl_scanner/codex/helpers.rs index e559700cbb..81734ce954 100644 --- a/rust/src/core/jsonl_scanner/codex/helpers.rs +++ b/rust/src/core/jsonl_scanner/codex/helpers.rs @@ -28,6 +28,10 @@ pub(super) struct CodexFastPayload<'a> { pub(super) model: Option<&'a str>, #[serde(default, borrow)] pub(super) model_name: Option<&'a str>, + #[serde(default, borrow, alias = "turnId")] + pub(super) turn_id: Option<&'a str>, + #[serde(default, borrow)] + pub(super) id: Option<&'a str>, #[serde(default, borrow)] pub(super) info: Option>, #[serde(default)] @@ -48,6 +52,8 @@ pub(super) struct CodexFastInfo<'a> { pub(super) model: Option<&'a str>, #[serde(default, borrow)] pub(super) model_name: Option<&'a str>, + #[serde(default, borrow, alias = "turnId")] + pub(super) turn_id: Option<&'a str>, #[serde(default)] pub(super) total_token_usage: Option, #[serde(default)] @@ -80,6 +86,10 @@ pub(super) enum CodexFastEvent<'a> { timestamp: &'a str, payload: CodexFastPayload<'a>, }, + /// A `task_started` event opens a turn; later token rows belong to it. + TaskStarted { + turn_id: Option<&'a str>, + }, } pub(super) fn model_evidence(raw: &str) -> Option<&str> { @@ -311,10 +321,16 @@ pub(super) fn parse_codex_fast_event(line: &str) -> Option> { } "event_msg" => { let payload = parsed.payload.or(parsed.event_msg)?; - (payload.payload_type == Some("token_count")).then_some(CodexFastEvent::TokenCount { - timestamp: parsed.timestamp?, - payload, - }) + match payload.payload_type { + Some("token_count") => Some(CodexFastEvent::TokenCount { + timestamp: parsed.timestamp?, + payload, + }), + Some("task_started") => Some(CodexFastEvent::TaskStarted { + turn_id: payload.turn_id.or(payload.id).and_then(model_evidence), + }), + _ => None, + } } _ => None, } @@ -328,7 +344,29 @@ pub(super) fn is_candidate_codex_line(line: &str) -> bool { return false; } - !line.contains("\"type\":\"event_msg\"") || line.contains("\"token_count\"") + !line.contains("\"type\":\"event_msg\"") + || line.contains("\"token_count\"") + || line.contains("\"task_started\"") +} + +/// Turn id carried by a `task_started` or `token_count` payload. +pub(super) fn codex_turn_id(payload: &Value) -> Option<&str> { + fn direct(value: &Value) -> Option<&str> { + value + .get("turn_id") + .or_else(|| value.get("turnId")) + .and_then(Value::as_str) + .and_then(model_evidence) + } + direct(payload).or_else(|| payload.get("info").and_then(direct)) +} + +/// The `event_msg` payload type of a parsed line, when it has one. +pub(super) fn event_payload_type(obj: &Value) -> Option<&str> { + obj.get("payload") + .or_else(|| obj.get("event_msg")) + .and_then(|payload| payload.get("type")) + .and_then(Value::as_str) } pub(super) fn codex_timestamp_day_key(timestamp: &str) -> Option { diff --git a/rust/src/core/jsonl_scanner/codex/parser.rs b/rust/src/core/jsonl_scanner/codex/parser.rs index 2ec70f9cf1..c22200989e 100644 --- a/rust/src/core/jsonl_scanner/codex/parser.rs +++ b/rust/src/core/jsonl_scanner/codex/parser.rs @@ -6,6 +6,8 @@ use serde_json::Value; pub(super) struct CodexParserState { pub(super) current_model: Option, + /// Turn opened by the latest `task_started` event; token rows inherit it. + current_turn_id: Option, pub(super) previous_totals: Option, /// High watermark of observed cumulative totals (never lowered). Used for /// Ultra interleaved-lineage containment (issue #2037 Phase 1). @@ -80,6 +82,7 @@ impl CodexParserState { .and_then(|baseline| remaining_inherited_totals.or_else(|| Some(baseline.clone()))); Self { current_model: initial_model, + current_turn_id: None, previous_totals: initial_totals.clone(), totals_watermark: initial_totals, saw_interleaved_totals: false, @@ -169,12 +172,30 @@ impl CodexParserState { totals.cached, totals.output, totals.reasoning, + self.current_turn_id.clone(), source_end_offset, ); } return; } + if event_payload_type(&obj) == Some("task_started") { + // A turn boundary is not a usage row, so it applies regardless of + // the requested day window. + let payload = obj.get("payload").or_else(|| obj.get("event_msg")); + self.current_turn_id = payload + .and_then(|payload| { + codex_turn_id(payload).or_else(|| { + payload + .get("id") + .and_then(Value::as_str) + .and_then(model_evidence) + }) + }) + .map(str::to_string); + return; + } + let is_token_count = token_count_payload(&obj).is_some(); let parsed_timestamp = obj .get("timestamp") @@ -225,6 +246,9 @@ impl CodexParserState { self.current_model = model_evidence(raw).map(str::to_string); } } + CodexFastEvent::TaskStarted { turn_id } => { + self.current_turn_id = turn_id.map(str::to_string); + } CodexFastEvent::TokenCount { timestamp, payload } => { let parsed_timestamp = parse_codex_timestamp(timestamp); self.observe_token_timestamp(timestamp, parsed_timestamp.as_ref()); @@ -300,6 +324,9 @@ impl CodexParserState { let info = payload.get("info"); let model = self.resolve_token_model(info, payload, obj); + let turn_id = codex_turn_id(payload) + .map(str::to_string) + .or_else(|| self.current_turn_id.clone()); self.record_usage( range, day_key, @@ -309,6 +336,7 @@ impl CodexParserState { delta_cached, delta_output, reasoning, + turn_id, source_end_offset, ); } @@ -345,6 +373,12 @@ impl CodexParserState { .or(event_model) .unwrap_or(CostUsagePricing::CODEX_UNATTRIBUTED_MODEL) .to_string(); + let turn_id = payload + .turn_id + .or_else(|| payload.info.as_ref().and_then(|info| info.turn_id)) + .and_then(model_evidence) + .map(str::to_string) + .or_else(|| self.current_turn_id.clone()); self.record_usage( range, day_key, @@ -354,6 +388,7 @@ impl CodexParserState { delta_cached, delta_output, reasoning, + turn_id, source_end_offset, ); } @@ -372,6 +407,7 @@ impl CodexParserState { cached: i64, output: i64, reasoning: Option, + turn_id: Option, source_end_offset: i64, ) { if !CostUsageDayRange::is_in_range(&day_key, &range.since_key, &range.until_key) { @@ -386,6 +422,7 @@ impl CodexParserState { cached: cached.min(input), output, reasoning: clamp_reasoning(reasoning, output), + turn_id, }, source_end_offset, )); diff --git a/rust/src/core/jsonl_scanner/codex/priority.rs b/rust/src/core/jsonl_scanner/codex/priority.rs new file mode 100644 index 0000000000..49fc4a0e3b --- /dev/null +++ b/rust/src/core/jsonl_scanner/codex/priority.rs @@ -0,0 +1,136 @@ +//! Persisted state for Codex Priority (Fast) trace evidence. +//! +//! Codex records each websocket request in `/logs_2.sqlite`. A +//! request carrying `service_tier == "priority"` marks its whole turn as a +//! Priority turn, but plain `gpt-5.x` session rows never say so. The scanner +//! reads the trace database incrementally (see `cost_scanner::codex::priority_trace`) +//! and stores only ids, model names and timestamps here; row bodies contain +//! prompts and are never persisted. +//! +//! Priority is applied as an overlay when day totals are rebuilt. It is never +//! baked into the persisted source-row pricing evidence, so cached-price +//! recovery keeps comparing like with like. + +use std::collections::{BTreeMap, VecDeque}; +use std::path::Path; + +use serde::{Deserialize, Serialize}; + +use crate::core::CostUsagePricing; + +/// Completions seen before their request may belong to non-priority turns, so +/// that pending map is bounded to keep memory and cache size constant. +pub const CODEX_PRIORITY_COMPLETED_MODEL_RETENTION_LIMIT: usize = 4096; + +/// Evidence for one Priority turn read from a trace row. +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct CodexPriorityTurnMetadata { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub thread_id: Option, + pub turn_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub model: Option, + /// Trace row timestamp in Unix seconds. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub timestamp: Option, +} + +/// Digest of one trace row, used to detect a replaced or rewritten database. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct CodexPriorityCursorAnchor { + pub row_id: i64, + pub digest: String, +} + +/// Durable scan cursor for one trace database. +/// +/// Every map is ordered so the serialized cursor is byte-stable across +/// processes; an unchanged cursor must not look like a cache change. +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct CodexPriorityTurnsCursor { + pub database_path: String, + /// Earliest `ts` (Unix seconds) the accumulated evidence covers. + pub coverage_since_epoch: i64, + /// Highest `logs.rowid` already examined. Rowids are monotonic. + pub last_row_id: i64, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub file_identity: Option, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub anchors: Vec, + /// Source trace rows behind each Priority turn, keyed by rowid. + #[serde(default)] + pub request_sources: BTreeMap>, + /// `response.completed` models for known Priority turns, keyed by rowid. + #[serde(default)] + pub priority_completed_models: BTreeMap>, + /// Completions seen before their request (bounded, oldest evicted first). + #[serde(default)] + pub completed_models: BTreeMap>, + #[serde(default)] + pub completed_order: VecDeque, +} + +/// A borrowed view that prices session rows of one file as Priority. +pub(crate) struct CodexPriorityOverlay<'a> { + cursor: &'a CodexPriorityTurnsCursor, +} + +impl CodexPriorityTurnsCursor { + /// The overlay for a session file, or `None` when the file lives outside + /// the Codex home this trace database belongs to. Turn sets are never + /// shared across `CODEX_HOME` scopes. + pub(crate) fn overlay_for_file(&self, file_path: &Path) -> Option> { + if self.request_sources.is_empty() { + return None; + } + let home = Path::new(&self.database_path).parent()?; + path_is_under(file_path, home).then_some(CodexPriorityOverlay { cursor: self }) + } + + /// Latest retained Priority request for `turn_id`. + pub(crate) fn turn(&self, turn_id: &str) -> Option<&CodexPriorityTurnMetadata> { + self.request_sources.get(turn_id)?.values().next_back() + } + + /// Model evidence for a Priority turn: the latest completion, else the + /// request model. + pub(crate) fn turn_model(&self, turn_id: &str) -> Option<&str> { + self.priority_completed_models + .get(turn_id) + .and_then(|models| models.values().next_back()) + .map(String::as_str) + .or_else(|| self.turn(turn_id)?.model.as_deref()) + } +} + +impl CodexPriorityOverlay<'_> { + /// The `-priority` name a row prices under, when its turn is a + /// Priority turn and the model has a Fast lane. A model without a Fast + /// lane keeps Standard pricing, as upstream does. + pub(crate) fn priority_model(&self, turn_id: Option<&str>, row_model: &str) -> Option { + let turn_id = turn_id?; + self.cursor.turn(turn_id)?; + if row_model.ends_with("-priority") { + return None; + } + let priced = self + .cursor + .turn_model(turn_id) + .filter(|model| CostUsagePricing::codex_api_fast_multiplier(model).is_some()) + .unwrap_or(row_model); + CostUsagePricing::codex_api_fast_multiplier(priced)?; + Some(format!( + "{}-priority", + CostUsagePricing::codex_fast_base_model(priced) + )) + } +} + +fn path_is_under(path: &Path, root: &Path) -> bool { + let normalize = |value: &Path| value.to_string_lossy().replace('\\', "/").to_lowercase(); + let root = normalize(root); + let root = root.trim_end_matches('/'); + let path = normalize(path); + path.strip_prefix(root) + .is_some_and(|rest| rest.starts_with('/')) +} diff --git a/rust/src/core/jsonl_scanner/codex/source_rows.rs b/rust/src/core/jsonl_scanner/codex/source_rows.rs index 9e0296663f..17a5a2aedb 100644 --- a/rust/src/core/jsonl_scanner/codex/source_rows.rs +++ b/rust/src/core/jsonl_scanner/codex/source_rows.rs @@ -21,6 +21,7 @@ pub(crate) fn rows_from_records(records: &[(CodexUsageRecord, i64)]) -> Vec String { } } +/// Apply the row's stored pricing mode and optional trace evidence. +pub(crate) fn priced_model( + pricing: &CodexSourcePricingEvidence, + turn_id: Option<&str>, + overlay: Option<&CodexPriorityOverlay<'_>>, +) -> Option { + let model = pricing + .pricing_model + .as_deref() + .filter(|model| !model.is_empty())?; + let model = + if pricing.pricing_mode.as_deref() == Some("priority") && !model.ends_with("-priority") { + format!("{model}-priority") + } else { + model.to_string() + }; + Some( + overlay + .and_then(|overlay| overlay.priority_model(turn_id, &model)) + .unwrap_or(model), + ) +} + +/// The model one source row prices under: [`priced_model`], except that a +/// Priority row the Fast lane cannot serve (too long for its model, or a +/// model without a Fast lane) prices at its Standard base model. Upstream +/// charges such a row the base cost because no Priority rate applies. +pub(crate) fn row_priced_model( + row: &CodexSourceUsageRow, + overlay: Option<&CodexPriorityOverlay<'_>>, +) -> Option { + let model = priced_model(&row.pricing, row.turn_id.as_deref(), overlay)?; + let input = u64::try_from(row.input.max(0)).unwrap_or(0); + if pricing_mode_of_model(&model) == "priority" + && !CostUsagePricing::codex_fast_lane_covers(&model, input) + { + return Some(model_of_pricing_mode(&model)); + } + Some(model) +} + /// Re-read the bounded reporting partition to obtain request-row order. /// The normal scanner still owns aggregate parsing and its byte budget; /// this path is used only after a complete file pass has established that diff --git a/rust/src/core/jsonl_scanner/codex/source_rows_tests.rs b/rust/src/core/jsonl_scanner/codex/source_rows_tests.rs index a82aa4416a..1fba1b18f6 100644 --- a/rust/src/core/jsonl_scanner/codex/source_rows_tests.rs +++ b/rust/src/core/jsonl_scanner/codex/source_rows_tests.rs @@ -13,6 +13,7 @@ fn source_row_recovery_preserves_unanimous_pricing_only() { output: 0, reasoning: None, source_end_offset: 100, + turn_id: None, pricing: CodexSourcePricingEvidence { pricing_model: Some("gpt-5.5".to_string()), pricing_mode: Some("priority".to_string()), @@ -27,6 +28,7 @@ fn source_row_recovery_preserves_unanimous_pricing_only() { output: 0, reasoning: None, source_end_offset: 200, + turn_id: None, pricing: CodexSourcePricingEvidence { pricing_model: Some("gpt-5.5".to_string()), pricing_mode: Some("standard".to_string()), @@ -44,6 +46,7 @@ fn source_row_recovery_preserves_unanimous_pricing_only() { output: 0, reasoning: None, source_end_offset: 300, + turn_id: None, pricing: CodexSourcePricingEvidence { pricing_model: Some("gpt-5.4".to_string()), pricing_mode: Some("standard".to_string()), @@ -68,6 +71,7 @@ fn source_row_recovery_does_not_price_an_appended_duplicate() { output: 0, reasoning: None, source_end_offset: 100, + turn_id: None, pricing: CodexSourcePricingEvidence { pricing_model: Some("gpt-5.5".to_string()), pricing_mode: Some("priority".to_string()), @@ -103,6 +107,7 @@ fn source_rows_from_records_use_source_model_as_initial_evidence() { cached: 20, output: -4, reasoning: Some(-2), + turn_id: None, }, 42, )]; @@ -185,3 +190,90 @@ fn pricing_mode_suffix_bijection_round_trips() { assert_eq!(model_of_pricing_mode(model), base); } } + +fn evidence_row(model: &str, mode: &str, input: i64, turn_id: Option<&str>) -> CodexSourceUsageRow { + CodexSourceUsageRow { + day_key: "2026-09-19".to_string(), + timestamp: None, + model: model.to_string(), + input, + cached: 0, + output: 10, + reasoning: None, + source_end_offset: 100, + turn_id: turn_id.map(str::to_string), + pricing: CodexSourcePricingEvidence { + pricing_model: Some(model.to_string()), + pricing_mode: Some(mode.to_string()), + }, + } +} + +#[test] +fn row_priced_model_keeps_priority_only_inside_the_fast_lane() { + let price = + |model, mode, input| row_priced_model(&evidence_row(model, mode, input, None), None); + assert_eq!( + price("gpt-5.5", "priority", 272_000).as_deref(), + Some("gpt-5.5-priority") + ); + // Too long for the gpt-5.5 Fast lane: upstream charges the Standard cost. + assert_eq!( + price("gpt-5.5", "priority", 272_001).as_deref(), + Some("gpt-5.5") + ); + // Astra publishes long-context Fast rates. + assert_eq!( + price("gpt-6-astra", "priority", 300_000).as_deref(), + Some("gpt-6-astra-priority") + ); + // gpt-5 has no Fast lane at all. + assert_eq!(price("gpt-5", "priority", 1_000).as_deref(), Some("gpt-5")); + assert_eq!( + price("gpt-5.5", "standard", 300_000).as_deref(), + Some("gpt-5.5") + ); +} + +#[test] +fn row_priced_model_without_pricing_evidence_stays_unattributed() { + let mut row = evidence_row("gpt-5.5", "priority", 1_000, None); + row.pricing = CodexSourcePricingEvidence::default(); + assert_eq!(row_priced_model(&row, None), None); +} + +#[test] +fn row_priced_model_applies_trace_evidence_only_inside_the_fast_lane() { + let home = std::env::temp_dir().join("codexbar-row-priced-model"); + let mut cursor = CodexPriorityTurnsCursor { + database_path: home.join("logs_2.sqlite").to_string_lossy().into_owned(), + ..Default::default() + }; + cursor.request_sources.insert( + "turn-fast".to_string(), + std::collections::BTreeMap::from([( + 1, + CodexPriorityTurnMetadata { + turn_id: "turn-fast".to_string(), + model: Some("gpt-5.5".to_string()), + ..Default::default() + }, + )]), + ); + let file = home.join("sessions").join("rollout.jsonl"); + let overlay = cursor + .overlay_for_file(&file) + .expect("a session file under the trace database home"); + let price = |input, turn| { + row_priced_model( + &evidence_row("gpt-5.5", "standard", input, Some(turn)), + Some(&overlay), + ) + }; + assert_eq!( + price(1_000, "turn-fast").as_deref(), + Some("gpt-5.5-priority") + ); + assert_eq!(price(300_000, "turn-fast").as_deref(), Some("gpt-5.5")); + assert_eq!(price(1_000, "turn-other").as_deref(), Some("gpt-5.5")); +} diff --git a/rust/src/core/jsonl_scanner/tests.rs b/rust/src/core/jsonl_scanner/tests.rs index bcb1129cb1..4190cd5068 100644 --- a/rust/src/core/jsonl_scanner/tests.rs +++ b/rust/src/core/jsonl_scanner/tests.rs @@ -142,6 +142,7 @@ fn codex_token_pipeline_preserves_counts_above_i32_max() { cached: 1_400_000_000, output: 100, reasoning: None, + turn_id: None, }, ); } @@ -223,6 +224,7 @@ fn legacy_packed_rows_remain_three_slots_and_report_reasoning_is_unknown() { cached: 1, output: 3, reasoning: Some(2), + turn_id: None, }; let mut packed = vec![10, 2, 4]; JsonlScanner::merge_codex_record_into_packed(&mut packed, &record); diff --git a/rust/src/cost_scanner.rs b/rust/src/cost_scanner.rs index 619f9ba7b0..4a77c72629 100755 --- a/rust/src/cost_scanner.rs +++ b/rust/src/cost_scanner.rs @@ -492,6 +492,9 @@ pub struct CostScanner { cache_root: Option, /// When set, bypass normal sessions-dir discovery (tests / inject roots). sessions_dirs_override: Option>, + /// Explicit Codex trace database (tests / inject roots). Without it the + /// ambient `CODEX_HOME` database is used, unless sessions dirs are injected. + codex_trace_database_override: Option, } impl CostScanner { @@ -502,6 +505,7 @@ impl CostScanner { options: CostScanOptions::default(), cache_root: None, sessions_dirs_override: None, + codex_trace_database_override: None, } } @@ -523,6 +527,12 @@ impl CostScanner { self } + /// Override the Codex Priority trace database (primarily for tests). + pub fn with_codex_trace_database(mut self, path: impl Into) -> Self { + self.codex_trace_database_override = Some(path.into()); + self + } + /// Scan standalone Pi and OMP local history. /// /// Pi session rows can represent either Codex or Claude models. They are diff --git a/rust/src/cost_scanner/codex.rs b/rust/src/cost_scanner/codex.rs index f054c6eb8f..bd76953e73 100644 --- a/rust/src/cost_scanner/codex.rs +++ b/rust/src/cost_scanner/codex.rs @@ -4,6 +4,7 @@ use crate::core::{CodexForkAccountingState, CodexSessionLineage}; mod cache_days; mod logical_target; mod pending_range; +mod priority_trace; mod reconciliation; mod scan; use cache_days::rebuild_cache_days; @@ -102,6 +103,25 @@ fn codex_parent_baseline( baseline } +/// `/logs_2.sqlite`, else `~/.codex/logs_2.sqlite`. +/// +/// Deliberate Windows deviation: upstream pins the trace database to +/// `~/.codex` and ignores `CODEX_HOME`. Here the scanned session roots come +/// from `CODEX_HOME`, and the Priority overlay only prices session files under +/// the database's own home, so both must resolve the same home or a +/// `CODEX_HOME` user would never see Priority pricing. +fn ambient_codex_trace_database_path( + codex_home: Option, + home_dir: Option, +) -> Option { + let home = codex_home + .map(|value| value.trim().to_string()) + .filter(|value| !value.is_empty()) + .map(PathBuf::from) + .or_else(|| home_dir.map(|home| home.join(".codex")))?; + Some(home.join(priority_trace::CODEX_TRACE_DATABASE_FILE)) +} + fn is_codex_path_in_scan_window( path: &Path, sessions_dirs: &[PathBuf], @@ -180,6 +200,21 @@ impl CostScanner { ) } + /// The Codex trace database that supplies Priority evidence. + /// + /// An explicit fixture is authoritative. Like upstream's test isolation, + /// unit tests and injected sessions roots never fall back to the ambient + /// user database, so no test reads (or caches a cursor for) real traces. + pub(super) fn codex_trace_database_path(&self) -> Option { + if let Some(path) = &self.codex_trace_database_override { + return Some(path.clone()); + } + if self.sessions_dirs_override.is_some() || cfg!(test) { + return None; + } + ambient_codex_trace_database_path(std::env::var("CODEX_HOME").ok(), dirs::home_dir()) + } + fn collect_codex_candidates( &self, sessions_dirs: &[PathBuf], @@ -486,6 +521,7 @@ impl CostScanner { .is_some_and(|history_base| Some(history_base) != codex_forked_from_id.as_deref()); if is_fork && fork_baseline.is_none() { + cache.codex_fork_rows.remove(&path_key); cache.files.insert( path_key, CostUsageFileUsage { @@ -579,6 +615,15 @@ impl CostScanner { .saturating_add(parse_result.token_timestamp_comparisons); let mut days = entry.days.clone(); merge_codex_records_into_days(&mut days, &parse_result.records); + if codex_forked_from_id.is_some() { + // Only extend rows that already cover the parsed prefix; + // a partial row set could never match the day totals. + if let Some(rows) = cache.codex_fork_rows.get_mut(&path_key) { + rows.extend(codex_fork_rows_from_records(&parse_result.records)); + } + } else { + cache.codex_fork_rows.remove(&path_key); + } let (session_cost, has_tokens) = add_codex_days_map_to_summary(summary, &days, range); if has_tokens { @@ -657,6 +702,7 @@ impl CostScanner { .token_timestamp_comparisons .saturating_add(parse_result.token_timestamp_comparisons); if parse_result.fork_baseline_ambiguous { + cache.codex_fork_rows.remove(&path_key); cache.files.insert( path_key, CostUsageFileUsage { @@ -686,6 +732,14 @@ impl CostScanner { } let mut days = HashMap::new(); merge_codex_records_into_days(&mut days, &parse_result.records); + if is_fork || codex_forked_from_id.is_some() { + cache.codex_fork_rows.insert( + path_key.clone(), + codex_fork_rows_from_records(&parse_result.records), + ); + } else { + cache.codex_fork_rows.remove(&path_key); + } let (session_cost, has_tokens) = add_codex_records_to_summary(summary, &parse_result.records, range); if has_tokens { @@ -737,6 +791,17 @@ impl CostScanner { } } +/// Request rows for a fork-shaped file, keeping exactly the records that +/// [`merge_codex_records_into_days`] folds into its day totals. +fn codex_fork_rows_from_records( + records: &[(crate::core::CodexUsageRecord, i64)], +) -> Vec { + crate::core::rows_from_records(records) + .into_iter() + .filter(|row| crate::core::CostUsagePricing::counts_toward_codex_subscription(&row.model)) + .collect() +} + #[cfg(test)] #[path = "tests.rs"] mod tests; diff --git a/rust/src/cost_scanner/codex/cache_days.rs b/rust/src/cost_scanner/codex/cache_days.rs index 6d63a40dc1..e1502a9874 100644 --- a/rust/src/cost_scanner/codex/cache_days.rs +++ b/rust/src/cost_scanner/codex/cache_days.rs @@ -1,10 +1,33 @@ use super::*; -use crate::core::{CodexSourceUsageRow, CodexUsageRecord, CostUsagePricing}; +use crate::core::{ + CodexPriorityOverlay, CodexSourceUsageRow, CodexUsageRecord, CostUsagePricing, row_priced_model, +}; +use std::collections::BTreeMap; + +type DayModels = HashMap>>; pub(super) fn rebuild_cache_days(cache: &mut CostUsageCache) { cache.days.clear(); - for usage in cache.files.values() { - for (day, models) in &usage.days { + let files = &cache.files; + cache.codex_fork_rows.retain(|path, _| { + files.get(path).is_some_and(|usage| { + (usage.codex_forked_from_id.is_some() || usage.codex_lineage.uses_parent_baseline()) + && !usage.codex_unresolved_fork_parent + }) + }); + let cursor = cache.codex_priority_turns_cursor.as_ref(); + for (path, usage) in &cache.files { + let rows = cache + .codex_source_rows + .get(path) + .map(|source| source.rows.as_slice()) + .or_else(|| cache.codex_fork_rows.get(path).map(Vec::as_slice)); + let overlaid = cursor + .and_then(|cursor| cursor.overlay_for_file(Path::new(path))) + .zip(rows) + .and_then(|(overlay, rows)| priority_days(rows, &usage.days, &overlay)); + let file_days = overlaid.as_ref().unwrap_or(&usage.days); + for (day, models) in file_days { let day_entry = cache.days.entry(day.clone()).or_default(); for (model, packed) in models { let dest = day_entry @@ -59,23 +82,47 @@ pub(super) fn rebuild_cache_days(cache: &mut CostUsageCache) { } } -pub(super) fn days_from_codex_source_rows( +/// Day totals for one file with Priority trace evidence applied, or `None` +/// when no row is affected or the retained source rows no longer match the +/// file's parsed day totals (a stale row cache must never replace newer data). +fn priority_days( + rows: &[CodexSourceUsageRow], + parsed_days: &DayModels, + overlay: &CodexPriorityOverlay<'_>, +) -> Option { + let plain = days_from_codex_source_rows_with_priority(rows, None); + if day_token_totals(&plain) != day_token_totals(parsed_days) { + return None; + } + let overlaid = days_from_codex_source_rows_with_priority(rows, Some(overlay)); + (overlaid != plain).then_some(overlaid) +} + +fn day_token_totals(days: &DayModels) -> BTreeMap<&str, (i64, i64, i64)> { + let mut totals = BTreeMap::new(); + for (day, models) in days { + let entry: &mut (i64, i64, i64) = totals.entry(day.as_str()).or_default(); + for packed in models.values() { + entry.0 = entry.0.saturating_add(packed.first().copied().unwrap_or(0)); + entry.1 = entry.1.saturating_add(packed.get(1).copied().unwrap_or(0)); + entry.2 = entry.2.saturating_add(packed.get(2).copied().unwrap_or(0)); + } + } + totals +} + +pub(super) fn days_from_codex_source_rows(rows: &[CodexSourceUsageRow]) -> DayModels { + days_from_codex_source_rows_with_priority(rows, None) +} + +fn days_from_codex_source_rows_with_priority( rows: &[CodexSourceUsageRow], -) -> HashMap>> { - let mut days: HashMap>> = HashMap::new(); + overlay: Option<&CodexPriorityOverlay<'_>>, +) -> DayModels { + let mut days: DayModels = HashMap::new(); for row in rows { - let model = match row.pricing.pricing_model.as_deref() { - Some(model) if !model.is_empty() => { - if row.pricing.pricing_mode.as_deref() == Some("priority") - && !model.ends_with("-priority") - { - format!("{model}-priority") - } else { - model.to_string() - } - } - _ => CostUsagePricing::CODEX_UNATTRIBUTED_MODEL.to_string(), - }; + let model = row_priced_model(row, overlay) + .unwrap_or_else(|| CostUsagePricing::CODEX_UNATTRIBUTED_MODEL.to_string()); let record = CodexUsageRecord { day_key: row.day_key.clone(), timestamp: row.timestamp, @@ -84,6 +131,7 @@ pub(super) fn days_from_codex_source_rows( cached: row.cached, output: row.output, reasoning: row.reasoning, + turn_id: row.turn_id.clone(), }; let packed = days .entry(record.day_key.clone()) diff --git a/rust/src/cost_scanner/codex/priority_trace.rs b/rust/src/cost_scanner/codex/priority_trace.rs new file mode 100644 index 0000000000..459e8cf9fa --- /dev/null +++ b/rust/src/cost_scanner/codex/priority_trace.rs @@ -0,0 +1,701 @@ +//! Incremental reader for Codex Priority (Fast) trace evidence. +//! +//! Upstream CodexBar #3820 (v0.65.0). Codex logs every websocket request to +//! `/logs_2.sqlite`; a `response.create` request with +//! `service_tier == "priority"` marks its turn as Priority. The database is +//! opened read-only (WAL safe, no sidecars created) and scanned incrementally: +//! +//! - the durable [`CodexPriorityTurnsCursor`] records the last examined rowid; +//! - SHA-256 anchors over sampled rows detect a replaced or rewritten +//! database, and retained source rows are revalidated so retention pruning +//! inside Codex removes stale turns; +//! - a cold scan runs in rowid chunks so it honors cancellation and keeps its +//! partial progress. +//! +//! Only ids, model names and timestamps are retained; row bodies contain +//! prompts and are neither stored nor logged. + +mod parse; + +use std::collections::{BTreeMap, HashMap, HashSet}; +use std::path::Path; +use std::sync::atomic::{AtomicBool, Ordering}; + +use chrono::{Local, NaiveDate}; +use rusqlite::types::ValueRef; +use rusqlite::{Connection, params}; +use sha2::{Digest, Sha256}; + +use crate::core::{ + CODEX_PRIORITY_COMPLETED_MODEL_RETENTION_LIMIT, CodexPriorityCursorAnchor, + CodexPriorityTurnMetadata, CodexPriorityTurnsCursor, DEFAULT_SQLITE_BUSY_TIMEOUT, JsonlScanner, + open_readonly_sqlite_connection, +}; +use parse::{parse_completed_trace_row, parse_priority_trace_row}; + +#[cfg(test)] +mod tests; + +/// Trace-database file name inside `CODEX_HOME`. +pub(super) const CODEX_TRACE_DATABASE_FILE: &str = "logs_2.sqlite"; + +/// Rowid span examined per chunk; cancellation is checked between chunks. +const ACCUMULATE_CHUNK_ROWS: i64 = 20_000; +/// Rows per `rowid in (...)` lookup when revalidating retained sources. +const SOURCE_LOOKUP_CHUNK: usize = 500; + +const BODY_FILTER: &str = "(feedback_log_body like '%websocket request:%' \ + or feedback_log_body like '%response.completed%' \ + or feedback_log_body like '%service_tier: Some(Some(\"priority\"))%')"; + +/// Result of one trace resolution. `cursor` is always the best-known state, +/// so a failed or interrupted scan never drops evidence already collected. +#[derive(Debug, Clone)] +pub(super) struct PriorityTraceResolution { + pub(super) cursor: Option, + /// True when the database could not be fully validated this pass; the + /// previous evidence stays in force and the next scan retries. + pub(super) validation_pending: bool, +} + +impl PriorityTraceResolution { + fn keep(cursor: Option, validation_pending: bool) -> Self { + Self { + cursor, + validation_pending, + } + } +} + +/// Resolve Priority turns from `database_path`, continuing `previous` when it +/// is still valid for this database. `expect_existing_database` is true when +/// this same path supplied evidence on an earlier scan (upstream +/// `expectExistingDatabase`). +pub(super) fn resolve_priority_turns( + database_path: &Path, + previous: Option, + coverage_since_epoch: i64, + expect_existing_database: bool, + cancel: Option<&AtomicBool>, +) -> PriorityTraceResolution { + let path_key = database_path.to_string_lossy().to_string(); + // A cursor for a different database is never reused. + let previous = previous.filter(|cursor| cursor.database_path == path_key); + + let Ok(metadata) = std::fs::metadata(database_path) else { + // A missing optional source is normal until this same path has + // supplied evidence; after that its absence is a validation failure + // that keeps the previous evidence and retries on the next scan. + return PriorityTraceResolution::keep(previous, expect_existing_database); + }; + let Some(identity) = JsonlScanner::codex_file_identity(database_path, &metadata) else { + return PriorityTraceResolution::keep(previous, true); + }; + let conn = match open_readonly_sqlite_connection(database_path, DEFAULT_SQLITE_BUSY_TIMEOUT) { + Ok(conn) => conn, + Err(_) => return PriorityTraceResolution::keep(previous, true), + }; + // The file may have been replaced between the identity read and the open. + let identity_after = std::fs::metadata(database_path) + .ok() + .and_then(|metadata| JsonlScanner::codex_file_identity(database_path, &metadata)); + if identity_after.as_deref() != Some(identity.as_str()) { + return PriorityTraceResolution::keep(previous, true); + } + let Some(max_row_id) = max_logs_row_id(&conn) else { + return PriorityTraceResolution::keep(previous, true); + }; + + let fresh = || CodexPriorityTurnsCursor { + database_path: path_key.clone(), + coverage_since_epoch, + file_identity: Some(identity.clone()), + ..CodexPriorityTurnsCursor::default() + }; + + let mut state = previous.clone().filter(|cursor| { + max_row_id >= cursor.last_row_id + && coverage_since_epoch >= cursor.coverage_since_epoch + && cursor.file_identity.as_deref() == Some(identity.as_str()) + }); + let mut anchors_need_refresh = false; + if let Some(cursor) = &state { + match validate_anchors(&conn, cursor) { + AnchorValidation::Valid { needs_refresh } => anchors_need_refresh = needs_refresh, + AnchorValidation::Invalid => state = None, + } + } + if let Some(cursor) = &mut state { + advance_coverage(cursor, coverage_since_epoch); + } + + let had_state = state.is_some(); + let mut resolved = state.unwrap_or_else(fresh); + let mut failed_replacement_fallback = None; + if had_state { + let mut pruned = resolved.clone(); + if prune_deleted_sources(&conn, &mut pruned).is_some() { + resolved = pruned; + } else { + // A retained source changed or could not be checked. Rebuild + // rather than publish evidence a cold scan might not reproduce, + // but keep the prior state as the failure result. + failed_replacement_fallback = Some(resolved); + resolved = fresh(); + anchors_need_refresh = false; + } + } + + if max_row_id > resolved.last_row_id { + let mut updated = resolved.clone(); + match accumulate(&conn, &mut updated, max_row_id, cancel) { + Accumulation::Complete => {} + Accumulation::Failed => { + return PriorityTraceResolution::keep( + failed_replacement_fallback.or(Some(resolved)), + true, + ); + } + Accumulation::Cancelled => { + // Keep partial progress; the next scan continues from it. + if updated.last_row_id == 0 { + return PriorityTraceResolution::keep( + failed_replacement_fallback.or(Some(resolved)), + true, + ); + } + if !capture_anchors(&conn, &mut updated) { + updated.anchors.clear(); + } + return PriorityTraceResolution::keep(Some(updated), true); + } + } + updated.last_row_id = max_row_id; + let anchored = capture_anchors(&conn, &mut updated); + return PriorityTraceResolution::keep(Some(updated), !anchored); + } + if anchors_need_refresh { + let mut updated = resolved.clone(); + let anchored = capture_anchors(&conn, &mut updated); + return PriorityTraceResolution::keep( + Some(if anchored { updated } else { resolved }), + !anchored, + ); + } + PriorityTraceResolution::keep(Some(resolved), false) +} + +/// Move a retained cursor with the active scan window and discard turn +/// evidence that can no longer match any scanned session rows. +fn advance_coverage(state: &mut CodexPriorityTurnsCursor, coverage_since_epoch: i64) { + if coverage_since_epoch <= state.coverage_since_epoch { + return; + } + state.coverage_since_epoch = coverage_since_epoch; + let mut expired_turns = Vec::new(); + state.request_sources.retain(|turn_id, sources| { + sources.retain(|_, metadata| { + metadata + .timestamp + .is_none_or(|timestamp| timestamp >= coverage_since_epoch) + }); + if sources.is_empty() { + expired_turns.push(turn_id.clone()); + false + } else { + true + } + }); + for turn_id in expired_turns { + state.priority_completed_models.remove(&turn_id); + } +} + +/// Best-effort rollback; the read-only connection has nothing to lose. +fn rollback(conn: &Connection) { + if let Err(error) = conn.execute_batch("rollback") { + tracing::debug!(%error, "Codex trace database rollback failed"); + } +} + +fn max_logs_row_id(conn: &Connection) -> Option { + conn.query_row("select max(rowid) from logs", [], |row| { + row.get::<_, Option>(0) + }) + .ok() + .map(|value| value.unwrap_or(0)) +} + +fn has_timestamp_index(conn: &Connection) -> bool { + conn.query_row( + "select 1 from sqlite_master where type = 'index' and tbl_name = 'logs' \ + and name = 'idx_logs_ts' limit 1", + [], + |_| Ok(()), + ) + .is_ok() +} + +fn column_text(row: &rusqlite::Row<'_>, index: usize) -> Option { + match row.get_ref(index).ok()? { + ValueRef::Text(bytes) | ValueRef::Blob(bytes) => { + Some(String::from_utf8_lossy(bytes).into_owned()) + } + _ => None, + } +} + +fn column_timestamp(row: &rusqlite::Row<'_>, index: usize) -> Option { + match row.get_ref(index).ok()? { + ValueRef::Integer(value) => Some(value), + // Truncation to whole seconds is intended for a fractional timestamp. + #[allow( + clippy::cast_possible_truncation, + reason = "whole-second trace timestamp" + )] + ValueRef::Real(value) => Some(value as i64), + ValueRef::Text(bytes) => text_timestamp(std::str::from_utf8(bytes).ok()?), + _ => None, + } +} + +/// A text `ts` value in Unix seconds. Upstream keys a non-integer text +/// timestamp by its leading `YYYY-MM-DD`, so such a row counts from the local +/// midnight of that day; coverage expiry then treats it like upstream's +/// day-key window check instead of retaining it forever. +fn text_timestamp(text: &str) -> Option { + let text = text.trim(); + if let Ok(seconds) = text.parse::() { + return Some(seconds); + } + let day = NaiveDate::parse_from_str(text.get(..10)?, "%Y-%m-%d").ok()?; + day.and_hms_opt(0, 0, 0)? + .and_local_timezone(Local) + .earliest() + .map(|midnight| midnight.timestamp()) +} + +enum Accumulation { + Complete, + Cancelled, + Failed, +} + +/// Examine rows in `(state.last_row_id, max_row_id]`, one rowid chunk at a +/// time. `state.last_row_id` advances only after a chunk fully succeeds. +fn accumulate( + conn: &Connection, + state: &mut CodexPriorityTurnsCursor, + max_row_id: i64, + cancel: Option<&AtomicBool>, +) -> Accumulation { + let mut from = state.last_row_id; + if from == 0 && state.coverage_since_epoch > 0 && has_timestamp_index(conn) { + // No matching row can precede the first row inside the window, so a + // cold scan skips the older history without reading its bodies. + let first = conn + .query_row( + "select min(rowid) from logs indexed by idx_logs_ts where ts >= ?1", + [state.coverage_since_epoch], + |row| row.get::<_, Option>(0), + ) + .ok(); + match first { + Some(Some(first)) => from = first.saturating_sub(1).max(0), + Some(None) => { + state.last_row_id = max_row_id; + return Accumulation::Complete; + } + None => return Accumulation::Failed, + } + } + let query = format!( + "select rowid, ts, feedback_log_body from logs \ + where rowid > ?1 and rowid <= ?2 and ts >= ?3 and {BODY_FILTER} order by rowid" + ); + let Ok(mut statement) = conn.prepare(&query) else { + return Accumulation::Failed; + }; + while from < max_row_id { + if cancel.is_some_and(|flag| flag.load(Ordering::Relaxed)) { + return Accumulation::Cancelled; + } + let to = from.saturating_add(ACCUMULATE_CHUNK_ROWS).min(max_row_id); + let Ok(mut rows) = statement.query(params![from, to, state.coverage_since_epoch]) else { + return Accumulation::Failed; + }; + loop { + match rows.next() { + Ok(Some(row)) => { + let row_id = row.get::<_, i64>(0).unwrap_or(0); + let timestamp = column_timestamp(row, 1); + if let Some(body) = column_text(row, 2) { + absorb_row(state, row_id, timestamp, &body); + } + } + Ok(None) => break, + Err(_) => return Accumulation::Failed, + } + } + state.last_row_id = to; + from = to; + } + Accumulation::Complete +} + +/// Fold one trace row into the cursor state. +fn absorb_row( + state: &mut CodexPriorityTurnsCursor, + row_id: i64, + timestamp: Option, + body: &str, +) { + if let Some(completed) = parse_completed_trace_row(body) { + if state.turn(&completed.turn_id).is_some() { + state + .priority_completed_models + .entry(completed.turn_id) + .or_default() + .insert(row_id, completed.model); + } else { + store_pending_completed_models( + state, + &completed.turn_id, + BTreeMap::from([(row_id, completed.model)]), + ); + } + return; + } + let Some(parsed) = parse_priority_trace_row(timestamp, body) else { + return; + }; + let turn_id = parsed.turn_id.clone(); + state + .request_sources + .entry(turn_id.clone()) + .or_default() + .insert(row_id, parsed); + if let Some(models) = state.completed_models.remove(&turn_id) { + state.completed_order.retain(|id| id != &turn_id); + state.priority_completed_models.insert(turn_id, models); + } +} + +/// Retain completions that arrived before their request. The map is bounded: +/// once over the limit the oldest turn's models are evicted. +fn store_pending_completed_models( + state: &mut CodexPriorityTurnsCursor, + turn_id: &str, + models: BTreeMap, +) { + if !state.completed_models.contains_key(turn_id) { + state.completed_order.push_back(turn_id.to_string()); + if state.completed_order.len() > CODEX_PRIORITY_COMPLETED_MODEL_RETENTION_LIMIT { + let evicted = state + .completed_order + .pop_front() + .expect("queue length checked"); + state.completed_models.remove(&evicted); + } + } + state + .completed_models + .entry(turn_id.to_string()) + .or_default() + .extend(models); +} + +enum AnchorValidation { + Valid { needs_refresh: bool }, + Invalid, +} + +enum AnchorLookup { + Found(String), + Missing, + Failed, +} + +fn hash_anchor_value(hasher: &mut Sha256, value: ValueRef<'_>) { + match value { + ValueRef::Null => hasher.update([0]), + ValueRef::Integer(value) => { + hasher.update([1]); + hasher.update(value.to_le_bytes()); + } + ValueRef::Real(value) => { + hasher.update([2]); + hasher.update(value.to_bits().to_le_bytes()); + } + ValueRef::Text(value) => { + hasher.update([3]); + hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_le_bytes()); + hasher.update(value); + } + ValueRef::Blob(value) => { + hasher.update([4]); + hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_le_bytes()); + hasher.update(value); + } + } +} + +fn row_anchor_digest(row: &rusqlite::Row<'_>, timestamp_index: usize, body_index: usize) -> String { + let mut hasher = Sha256::new(); + if let Ok(timestamp) = row.get_ref(timestamp_index) { + hash_anchor_value(&mut hasher, timestamp); + } + if let Ok(body) = row.get_ref(body_index) { + hash_anchor_value(&mut hasher, body); + } + hasher + .finalize() + .iter() + .map(|byte| format!("{byte:02x}")) + .collect() +} + +/// Sample four rows across the scanned range (a quarter, half, three +/// quarters and the last) so a rewritten database cannot look unchanged. +fn capture_anchors(conn: &Connection, state: &mut CodexPriorityTurnsCursor) -> bool { + if state.last_row_id <= 0 { + state.anchors.clear(); + return true; + } + if conn.execute_batch("begin deferred transaction").is_err() { + return false; + } + let anchors = sample_anchors(conn, state.last_row_id); + let committed = conn.execute_batch("commit").is_ok(); + if !committed { + rollback(conn); + } + match anchors { + Some(anchors) if committed && !anchors.is_empty() => { + state.anchors = anchors; + true + } + _ => false, + } +} + +fn sample_anchors(conn: &Connection, last_row_id: i64) -> Option> { + let minimum = conn + .query_row( + "select min(rowid) from logs where rowid <= ?1", + [last_row_id], + |row| row.get::<_, Option>(0), + ) + .ok()??; + let span = last_row_id - minimum; + let targets = [ + minimum + span / 4, + minimum + span / 2, + minimum + (span / 4) * 3 + (span % 4) * 3 / 4, + last_row_id, + ]; + let mut anchors: Vec = Vec::new(); + for target in targets { + let anchor = conn + .query_row( + "select rowid, ts, feedback_log_body from logs where rowid <= ?1 \ + order by rowid desc limit 1", + [target], + |row| { + Ok(CodexPriorityCursorAnchor { + row_id: row.get(0)?, + digest: row_anchor_digest(row, 1, 2), + }) + }, + ) + .ok()?; + if !anchors.iter().any(|known| known.row_id == anchor.row_id) { + anchors.push(anchor); + } + } + Some(anchors) +} + +fn lookup_anchor(conn: &Connection, row_id: i64) -> AnchorLookup { + let result = conn.query_row( + "select ts, feedback_log_body from logs where rowid = ?1", + [row_id], + |row| Ok(row_anchor_digest(row, 0, 1)), + ); + match result { + Ok(digest) => AnchorLookup::Found(digest), + Err(rusqlite::Error::QueryReturnedNoRows) => AnchorLookup::Missing, + Err(_) => AnchorLookup::Failed, + } +} + +/// Check that the sampled rows still hold what was recorded. Codex prunes old +/// rows in place, so a missing anchor only asks for a refresh; a changed one +/// means the database was rewritten. +fn validate_anchors(conn: &Connection, state: &CodexPriorityTurnsCursor) -> AnchorValidation { + if state.anchors.is_empty() { + return if state.last_row_id == 0 { + AnchorValidation::Valid { + needs_refresh: true, + } + } else { + AnchorValidation::Invalid + }; + } + if conn.execute_batch("begin deferred transaction").is_err() { + return AnchorValidation::Invalid; + } + let (mut matching, mut missing) = (0_usize, 0_usize); + let (mut failed, mut mismatched) = (false, false); + for anchor in &state.anchors { + match lookup_anchor(conn, anchor.row_id) { + AnchorLookup::Found(digest) if digest == anchor.digest => matching += 1, + AnchorLookup::Found(_) => mismatched = true, + AnchorLookup::Missing => missing += 1, + AnchorLookup::Failed => failed = true, + } + } + if conn.execute_batch("commit").is_err() { + rollback(conn); + return AnchorValidation::Invalid; + } + if failed || mismatched { + return AnchorValidation::Invalid; + } + // With fewer than four distinct rows there is no distributed quorum. + let required = if state.anchors.len() < 4 { + state.anchors.len() + } else { + 2 + }; + if matching < required { + return AnchorValidation::Invalid; + } + AnchorValidation::Valid { + needs_refresh: missing > 0, + } +} + +/// Drop turns whose source rows Codex has since deleted. Returns `None` when +/// a retained source no longer matches (rebuild required) or cannot be read. +fn prune_deleted_sources(conn: &Connection, state: &mut CodexPriorityTurnsCursor) -> Option { + let retained = retained_source_row_ids(conn, state)?; + let mut pruned = false; + + let turn_ids: Vec = state.request_sources.keys().cloned().collect(); + for turn_id in turn_ids { + let Some(sources) = state.request_sources.get(&turn_id) else { + continue; + }; + let kept: BTreeMap = sources + .iter() + .filter(|(row_id, _)| retained.contains(row_id)) + .map(|(row_id, meta)| (*row_id, meta.clone())) + .collect(); + if kept.len() == sources.len() { + continue; + } + pruned = true; + if kept.is_empty() { + state.request_sources.remove(&turn_id); + if let Some(models) = state.priority_completed_models.remove(&turn_id) { + store_pending_completed_models(state, &turn_id, models); + } + } else { + state.request_sources.insert(turn_id, kept); + } + } + pruned |= prune_completed(&retained, &mut state.priority_completed_models); + pruned |= prune_completed(&retained, &mut state.completed_models); + let completed = &state.completed_models; + state + .completed_order + .retain(|id| completed.contains_key(id)); + Some(pruned) +} + +fn prune_completed( + retained: &HashSet, + models: &mut BTreeMap>, +) -> bool { + let mut pruned = false; + models.retain(|_, by_row| { + let before = by_row.len(); + by_row.retain(|row_id, _| retained.contains(row_id)); + pruned |= by_row.len() != before; + !by_row.is_empty() + }); + pruned +} + +/// Row ids of retained sources that still parse to the recorded evidence. +fn retained_source_row_ids( + conn: &Connection, + state: &CodexPriorityTurnsCursor, +) -> Option> { + let mut requests: HashMap = HashMap::new(); + for sources in state.request_sources.values() { + for (row_id, meta) in sources { + if requests.get(row_id).is_some_and(|known| *known != meta) { + return None; + } + requests.insert(*row_id, meta); + } + } + let mut completions: HashMap = HashMap::new(); + for by_turn in [&state.priority_completed_models, &state.completed_models] { + for (turn_id, by_row) in by_turn { + for (row_id, model) in by_row { + let entry = (turn_id.as_str(), model.as_str()); + if completions.get(row_id).is_some_and(|known| *known != entry) { + return None; + } + completions.insert(*row_id, entry); + } + } + } + if requests + .keys() + .any(|row_id| completions.contains_key(row_id)) + { + return None; + } + let row_ids: Vec = requests + .keys() + .chain(completions.keys()) + .copied() + .collect::>() + .into_iter() + .collect(); + + let mut retained = HashSet::new(); + for chunk in row_ids.chunks(SOURCE_LOOKUP_CHUNK) { + let placeholders = vec!["?"; chunk.len()].join(","); + let query = format!( + "select rowid, ts, feedback_log_body from logs where rowid in ({placeholders})" + ); + let mut statement = conn.prepare(&query).ok()?; + let mut rows = statement + .query(rusqlite::params_from_iter(chunk.iter())) + .ok()?; + while let Some(row) = rows.next().ok()? { + let row_id: i64 = row.get(0).ok()?; + let body = column_text(row, 2)?; + if let Some(expected) = requests.get(&row_id) { + let parsed = parse_priority_trace_row(column_timestamp(row, 1), &body)?; + if &parsed != *expected { + return None; + } + } else if let Some((turn_id, model)) = completions.get(&row_id) { + let parsed = parse_completed_trace_row(&body)?; + if parsed.turn_id != *turn_id || parsed.model != *model { + return None; + } + } else { + return None; + } + retained.insert(row_id); + } + } + Some(retained) +} diff --git a/rust/src/cost_scanner/codex/priority_trace/parse.rs b/rust/src/cost_scanner/codex/priority_trace/parse.rs new file mode 100644 index 0000000000..6908eaf14f --- /dev/null +++ b/rust/src/cost_scanner/codex/priority_trace/parse.rs @@ -0,0 +1,259 @@ +//! Parsers for Codex trace-log rows. Bodies contain prompts, so every parser +//! works in memory and returns only ids, model names and timestamps. + +use std::borrow::Cow; +use std::fmt; +use std::marker::PhantomData; + +use serde::Deserialize; +use serde::de::{self, Deserializer, IgnoredAny, MapAccess, SeqAccess, Visitor}; + +use crate::core::CodexPriorityTurnMetadata; + +const REQUEST_MARKER: &str = "websocket request:"; +const EVENT_MARKER: &str = "websocket event:"; +const SUBMISSION_MARKER: &str = "Submission sub=Submission {"; +const PRIORITY_SUBMISSION_TIER: &str = "service_tier: Some(Some(\"priority\"))"; + +// Upstream reads these bodies as `[String: Any]` and casts each field with +// `as? String`: only a JSON object is a body, and a field of another type is +// absent rather than an error. The lenient readers below keep that contract, +// so a numeric `model` no longer drops a Priority turn and an array is never +// read positionally as a struct. +#[derive(Deserialize)] +struct RequestBody<'a> { + #[serde(default, rename = "type", borrow, deserialize_with = "lenient_str")] + kind: Option>, + #[serde(default, borrow, deserialize_with = "lenient_str")] + service_tier: Option>, + #[serde(default, borrow, deserialize_with = "lenient_str")] + turn_id: Option>, + #[serde(default, borrow, deserialize_with = "lenient_str")] + model: Option>, +} + +#[derive(Deserialize)] +struct EventBody<'a> { + #[serde(default, rename = "type", borrow, deserialize_with = "lenient_str")] + kind: Option>, + #[serde(default, borrow, deserialize_with = "object_only")] + response: Option>, +} + +#[derive(Deserialize)] +struct EventResponse<'a> { + #[serde(default, borrow, deserialize_with = "lenient_str")] + model: Option>, +} + +/// A JSON string, or `None` for any other value (upstream `as? String`). +fn lenient_str<'de, D: Deserializer<'de>>( + deserializer: D, +) -> Result>, D::Error> { + deserializer.deserialize_any(LenientStr) +} + +/// A JSON object read as `T`, or `None` for any other value (upstream +/// `as? [String: Any]`). +fn object_only<'de, D: Deserializer<'de>, T: Deserialize<'de>>( + deserializer: D, +) -> Result, D::Error> { + deserializer.deserialize_any(ObjectOnly(PhantomData)) +} + +struct LenientStr; + +impl<'de> Visitor<'de> for LenientStr { + type Value = Option>; + + fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("any JSON value") + } + + fn visit_borrowed_str(self, value: &'de str) -> Result { + Ok(Some(Cow::Borrowed(value))) + } + + fn visit_str(self, value: &str) -> Result { + Ok(Some(Cow::Owned(value.to_owned()))) + } + + fn visit_string(self, value: String) -> Result { + Ok(Some(Cow::Owned(value))) + } + + fn visit_bool(self, _: bool) -> Result { + Ok(None) + } + + fn visit_i64(self, _: i64) -> Result { + Ok(None) + } + + fn visit_u64(self, _: u64) -> Result { + Ok(None) + } + + fn visit_f64(self, _: f64) -> Result { + Ok(None) + } + + fn visit_unit(self) -> Result { + Ok(None) + } + + fn visit_seq>(self, seq: A) -> Result { + skip_seq(seq) + } + + fn visit_map>(self, mut map: A) -> Result { + while map.next_entry::()?.is_some() {} + Ok(None) + } +} + +struct ObjectOnly(PhantomData); + +impl<'de, T: Deserialize<'de>> Visitor<'de> for ObjectOnly { + type Value = Option; + + fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("any JSON value") + } + + fn visit_map>(self, map: A) -> Result { + T::deserialize(de::value::MapAccessDeserializer::new(map)).map(Some) + } + + fn visit_str(self, _: &str) -> Result { + Ok(None) + } + + fn visit_bool(self, _: bool) -> Result { + Ok(None) + } + + fn visit_i64(self, _: i64) -> Result { + Ok(None) + } + + fn visit_u64(self, _: u64) -> Result { + Ok(None) + } + + fn visit_f64(self, _: f64) -> Result { + Ok(None) + } + + fn visit_unit(self) -> Result { + Ok(None) + } + + fn visit_seq>(self, seq: A) -> Result { + skip_seq(seq) + } +} + +fn skip_seq<'de, A: SeqAccess<'de>, T>(mut seq: A) -> Result, A::Error> { + while seq.next_element::()?.is_some() {} + Ok(None) +} + +/// The JSON object after a log marker; any other JSON value is not a body. +fn object_after<'a, T: Deserialize<'a>>(body: &'a str, marker_end: usize) -> Option { + let json = body[marker_end..].trim(); + if !json.starts_with('{') { + return None; + } + serde_json::from_str(json).ok() +} + +/// A `response.completed` event: the turn and the model that served it. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(super) struct CompletedTrace { + pub(super) turn_id: String, + pub(super) model: String, +} + +/// Priority evidence from one trace row: either a `response.create` request +/// with `service_tier == "priority"` or a Priority `Submission` row. +pub(super) fn parse_priority_trace_row( + timestamp: Option, + body: &str, +) -> Option { + let Some(marker) = body.find(REQUEST_MARKER) else { + return parse_priority_submission_row(timestamp, body); + }; + let prefix = &body[..marker]; + let request: RequestBody<'_> = object_after(body, marker + REQUEST_MARKER.len())?; + if request.kind.as_deref() != Some("response.create") + || request.service_tier.as_deref() != Some("priority") + { + return None; + } + let turn_id = value_named("turn.id", prefix) + .or_else(|| value_named("turn_id", prefix)) + .or(request.turn_id.as_deref()) + .filter(|id| !id.is_empty())?; + Some(CodexPriorityTurnMetadata { + thread_id: value_named("thread_id", prefix).map(str::to_string), + turn_id: turn_id.to_string(), + model: request.model.map(Cow::into_owned), + timestamp, + }) +} + +fn parse_priority_submission_row( + timestamp: Option, + body: &str, +) -> Option { + if !body.contains(PRIORITY_SUBMISSION_TIER) { + return None; + } + let submission = body.find(SUBMISSION_MARKER)?; + let tail = &body[submission + SUBMISSION_MARKER.len()..]; + let turn_id = quoted_value_named("id", tail)?; + Some(CodexPriorityTurnMetadata { + thread_id: value_named("thread_id", &body[..submission]).map(str::to_string), + turn_id: turn_id.to_string(), + model: None, + timestamp, + }) +} + +/// A `response.completed` websocket event with a turn id and response model. +pub(super) fn parse_completed_trace_row(body: &str) -> Option { + let marker = body.find(EVENT_MARKER)?; + let prefix = &body[..marker]; + let event: EventBody<'_> = object_after(body, marker + EVENT_MARKER.len())?; + if event.kind.as_deref() != Some("response.completed") { + return None; + } + let model = event.response?.model.filter(|model| !model.is_empty())?; + let turn_id = value_named("turn.id", prefix) + .or_else(|| value_named("turn_id", prefix)) + .filter(|id| !id.is_empty())?; + Some(CompletedTrace { + turn_id: turn_id.to_string(), + model: model.into_owned(), + }) +} + +/// The token after `name=` in a log prefix, up to whitespace or punctuation. +fn value_named<'a>(name: &str, text: &'a str) -> Option<&'a str> { + let start = text.find(&format!("{name}="))? + name.len() + 1; + let tail = &text[start..]; + let end = tail + .find(|ch: char| ch.is_whitespace() || matches!(ch, ',' | ']' | ')' | '}' | ':')) + .unwrap_or(tail.len()); + let value = &tail[..end]; + (!value.is_empty()).then_some(value) +} + +/// The string after `name: "` up to the closing quote. +fn quoted_value_named<'a>(name: &str, text: &'a str) -> Option<&'a str> { + let start = text.find(&format!("{name}: \""))? + name.len() + 3; + let tail = &text[start..]; + let value = &tail[..tail.find('"')?]; + (!value.is_empty()).then_some(value) +} diff --git a/rust/src/cost_scanner/codex/priority_trace/tests.rs b/rust/src/cost_scanner/codex/priority_trace/tests.rs new file mode 100644 index 0000000000..3db5eee6cf --- /dev/null +++ b/rust/src/cost_scanner/codex/priority_trace/tests.rs @@ -0,0 +1,544 @@ +//! Tests for the Priority trace resolver and its scanner integration. + +use super::*; +use crate::core::CostUsagePricing; +use crate::cost_scanner::{CostScanOptions, CostScanner}; +use chrono::Local; +use std::path::PathBuf; + +mod scanner; + +fn request_body(turn: &str, tier: &str, model: &str) -> String { + format!( + "session_loop{{thread_id=thread-1}}:turn{{turn.id={turn}}}: \ + codex_api::endpoint::responses_websocket: websocket request: \ + {{\"type\":\"response.create\",\"model\":\"{model}\",\"service_tier\":\"{tier}\"}}" + ) +} + +fn completed_body(turn: &str, model: &str) -> String { + format!( + "session_loop{{thread_id=thread-1}}:turn{{turn.id={turn}}}: \ + codex_api::endpoint::responses_websocket: websocket event: \ + {{\"type\":\"response.completed\",\"response\":{{\"model\":\"{model}\"}}}}" + ) +} + +struct TraceDb { + _dir: Option, + path: PathBuf, +} + +impl TraceDb { + fn new(rows: &[(i64, String)]) -> Self { + let dir = tempfile::tempdir().unwrap(); + let mut db = Self::in_dir(dir.path(), rows); + db._dir = Some(dir); + db + } + + /// A trace database inside an existing Codex home directory. + fn in_dir(home: &Path, rows: &[(i64, String)]) -> Self { + let path = home.join(CODEX_TRACE_DATABASE_FILE); + let conn = Connection::open(&path).unwrap(); + conn.execute_batch( + "create table logs (id integer primary key autoincrement, ts integer not null, \ + feedback_log_body text); create index idx_logs_ts on logs(ts);", + ) + .unwrap(); + let db = Self { _dir: None, path }; + db.insert(rows); + db + } + + /// Insert rows in one transaction, like upstream `insertTestLogs`. + fn insert(&self, rows: &[(i64, String)]) { + let mut conn = Connection::open(&self.path).unwrap(); + let transaction = conn.transaction().unwrap(); + for (ts, body) in rows { + transaction + .execute( + "insert into logs (ts, feedback_log_body) values (?1, ?2)", + params![ts, body], + ) + .unwrap(); + } + transaction.commit().unwrap(); + } + + fn execute(&self, sql: &str) { + Connection::open(&self.path) + .unwrap() + .execute_batch(sql) + .unwrap(); + } +} + +fn resolve(db: &TraceDb, previous: Option) -> PriorityTraceResolution { + resolve_priority_turns(&db.path, previous, 0, false, None) +} + +fn noise(count: usize) -> Vec<(i64, String)> { + (0..count) + .map(|index| { + ( + 1_000 + i64::try_from(index).unwrap(), + format!("unrelated log line {index}"), + ) + }) + .collect() +} + +#[test] +fn parses_priority_request_with_turn_thread_and_model() { + let parsed = + parse_priority_trace_row(Some(42), &request_body("turn-1", "priority", "gpt-5.5")).unwrap(); + assert_eq!(parsed.turn_id, "turn-1"); + assert_eq!(parsed.thread_id.as_deref(), Some("thread-1")); + assert_eq!(parsed.model.as_deref(), Some("gpt-5.5")); + assert_eq!(parsed.timestamp, Some(42)); +} + +#[test] +fn ignores_standard_requests_and_other_event_types() { + assert!(parse_priority_trace_row(None, &request_body("t", "default", "gpt-5.5")).is_none()); + let create = request_body("t", "priority", "gpt-5.5").replace("response.create", "response.x"); + assert!(parse_priority_trace_row(None, &create).is_none()); + assert!(parse_priority_trace_row(None, "plain log").is_none()); +} + +#[test] +fn request_turn_id_falls_back_to_the_json_body() { + let body = "websocket request: {\"type\":\"response.create\",\"service_tier\":\"priority\",\ + \"turn_id\":\"json-turn\"}"; + let parsed = parse_priority_trace_row(None, body).unwrap(); + assert_eq!(parsed.turn_id, "json-turn"); + let without = "websocket request: {\"type\":\"response.create\",\"service_tier\":\"priority\"}"; + assert!(parse_priority_trace_row(None, without).is_none()); +} + +#[test] +fn parses_priority_submission_rows() { + let body = "thread_id=thread-9: Submission sub=Submission { id: \"sub-turn\", \ + op: UserTurn { service_tier: Some(Some(\"priority\")) } }"; + let parsed = parse_priority_trace_row(Some(7), body).unwrap(); + assert_eq!(parsed.turn_id, "sub-turn"); + assert_eq!(parsed.thread_id.as_deref(), Some("thread-9")); + assert_eq!(parsed.model, None); + let standard = body.replace("priority", "flex"); + assert!(parse_priority_trace_row(None, &standard).is_none()); +} + +#[test] +fn parses_completed_events() { + let parsed = parse_completed_trace_row(&completed_body("turn-1", "gpt-5.5")).unwrap(); + assert_eq!(parsed.turn_id, "turn-1"); + assert_eq!(parsed.model, "gpt-5.5"); + let other = completed_body("turn-1", "gpt-5.5").replace("response.completed", "response.x"); + assert!(parse_completed_trace_row(&other).is_none()); +} + +#[test] +fn request_fields_of_another_type_are_absent_not_fatal() { + // Upstream casts each field with `as? String`: a numeric model is absent, + // but the turn is still a Priority turn. + let numeric_model = "turn.id=t1 websocket request: \ + {\"type\":\"response.create\",\"model\":5,\"service_tier\":\"priority\"}"; + let parsed = parse_priority_trace_row(None, numeric_model).unwrap(); + assert_eq!(parsed.turn_id, "t1"); + assert_eq!(parsed.model, None); + + let numeric_tier = "turn.id=t1 websocket request: \ + {\"type\":\"response.create\",\"model\":\"gpt-5.5\",\"service_tier\":1}"; + assert!(parse_priority_trace_row(None, numeric_tier).is_none()); + + // Without the object guard serde would read an array positionally. + let array = "turn.id=t1 websocket request: \ + [\"response.create\",\"priority\",\"t1\",\"gpt-5.5\"]"; + assert!(parse_priority_trace_row(None, array).is_none()); +} + +#[test] +fn completed_event_without_an_object_response_is_ignored() { + let string_response = "turn.id=t1 websocket event: \ + {\"type\":\"response.completed\",\"response\":\"gpt-5.5\"}"; + assert!(parse_completed_trace_row(string_response).is_none()); + + let numeric_model = "turn.id=t1 websocket event: \ + {\"type\":\"response.completed\",\"response\":{\"model\":5}}"; + assert!(parse_completed_trace_row(numeric_model).is_none()); + + let array = "turn.id=t1 websocket event: \ + [\"response.completed\",{\"model\":\"gpt-5.5\"}]"; + assert!(parse_completed_trace_row(array).is_none()); +} + +fn local_midnight(year: i32, month: u32, day: u32) -> i64 { + NaiveDate::from_ymd_opt(year, month, day) + .and_then(|date| date.and_hms_opt(0, 0, 0)) + .and_then(|midnight| midnight.and_local_timezone(Local).earliest()) + .map(|midnight| midnight.timestamp()) + .unwrap() +} + +#[test] +fn text_timestamps_parse_seconds_or_day_keys() { + assert_eq!(text_timestamp("42"), Some(42)); + assert_eq!(text_timestamp(" 42 "), Some(42)); + let midnight = local_midnight(2026, 9, 4); + assert_eq!(text_timestamp("2026-09-04 12:34:56.789"), Some(midnight)); + assert_eq!(text_timestamp("2026-09-04T12:34:56Z"), Some(midnight)); + assert_eq!(text_timestamp("2026-09-04"), Some(midnight)); + assert_eq!(text_timestamp("garbage"), None); + assert_eq!(text_timestamp("2026-13-40 00:00:00"), None); + assert_eq!(text_timestamp(""), None); +} + +#[test] +fn text_timestamp_rows_count_from_their_local_day() { + let db = TraceDb::new(&[]); + db.execute(&format!( + "insert into logs (ts, feedback_log_body) values ('2026-09-04 12:00:00', '{}')", + request_body("turn-1", "priority", "gpt-5.5") + )); + let cursor = resolve(&db, None).cursor.unwrap(); + assert_eq!( + cursor.turn("turn-1").and_then(|turn| turn.timestamp), + Some(local_midnight(2026, 9, 4)) + ); +} + +#[test] +fn cold_scan_collects_priority_turns_and_completed_model() { + let mut rows = noise(3); + rows.push((2_000, request_body("turn-1", "priority", "gpt-5.5"))); + rows.push((2_001, request_body("turn-2", "default", "gpt-5.5"))); + rows.push((2_002, completed_body("turn-1", "gpt-5.4"))); + let db = TraceDb::new(&rows); + + let resolution = resolve(&db, None); + let cursor = resolution.cursor.unwrap(); + assert!(!resolution.validation_pending); + assert_eq!(cursor.request_sources.len(), 1); + assert!(cursor.turn("turn-1").is_some()); + assert_eq!(cursor.last_row_id, 6); + assert_eq!(cursor.anchors.len(), 4); + assert_eq!(cursor.turn_model("turn-1"), Some("gpt-5.4")); +} + +#[test] +fn incremental_scan_appends_only_new_rows() { + let db = TraceDb::new(&[(2_000, request_body("turn-1", "priority", "gpt-5.5"))]); + let first = resolve(&db, None).cursor.unwrap(); + assert_eq!(first.last_row_id, 1); + + db.insert(&[(2_001, request_body("turn-2", "priority", "gpt-5.4"))]); + let second = resolve(&db, Some(first)).cursor.unwrap(); + assert_eq!(second.last_row_id, 2); + assert_eq!(second.request_sources.len(), 2); +} + +#[test] +fn completion_before_request_is_matched_when_the_request_arrives() { + let db = TraceDb::new(&[(2_000, completed_body("turn-1", "gpt-5.4"))]); + let first = resolve(&db, None).cursor.unwrap(); + assert!(first.request_sources.is_empty()); + assert!(first.completed_models.contains_key("turn-1")); + + db.insert(&[(2_001, request_body("turn-1", "priority", "gpt-5.5"))]); + let second = resolve(&db, Some(first)).cursor.unwrap(); + assert!(second.completed_models.is_empty()); + assert_eq!(second.turn_model("turn-1"), Some("gpt-5.4")); +} + +#[test] +fn rewritten_database_rebuilds_instead_of_reusing_evidence() { + let mut rows = noise(6); + rows.push((2_000, request_body("old-turn", "priority", "gpt-5.5"))); + let db = TraceDb::new(&rows); + let first = resolve(&db, None).cursor.unwrap(); + assert!(first.turn("old-turn").is_some()); + + db.execute("update logs set feedback_log_body = 'rewritten ' || id"); + db.insert(&[(3_000, request_body("new-turn", "priority", "gpt-5.5"))]); + let second = resolve(&db, Some(first)).cursor.unwrap(); + assert!(second.turn("old-turn").is_none()); + assert!(second.turn("new-turn").is_some()); +} + +#[test] +fn deleted_source_rows_drop_their_turns() { + let db = TraceDb::new(&[ + (2_000, request_body("turn-1", "priority", "gpt-5.5")), + (2_001, request_body("turn-2", "priority", "gpt-5.5")), + (2_002, "tail row".to_string()), + ]); + let first = resolve(&db, None).cursor.unwrap(); + assert_eq!(first.request_sources.len(), 2); + + db.execute("delete from logs where id = 1"); + let second = resolve(&db, Some(first)).cursor.unwrap(); + assert!(second.turn("turn-1").is_none()); + assert!(second.turn("turn-2").is_some()); +} + +#[test] +fn cancelled_cold_scan_reports_pending_without_evidence() { + let db = TraceDb::new(&[(2_000, request_body("turn-1", "priority", "gpt-5.5"))]); + let cancel = AtomicBool::new(true); + let resolution = resolve_priority_turns(&db.path, None, 0, false, Some(&cancel)); + assert!(resolution.validation_pending); + assert!( + resolution + .cursor + .is_none_or(|cursor| cursor.request_sources.is_empty()) + ); +} + +#[test] +fn missing_database_is_pending_only_after_it_supplied_evidence() { + let db = TraceDb::new(&[(2_000, request_body("turn-1", "priority", "gpt-5.5"))]); + let cursor = resolve(&db, None).cursor.unwrap(); + + // A path that never held a database is a normal optional source. + let missing = db.path.with_file_name("absent.sqlite"); + let never = resolve_priority_turns(&missing, None, 0, false, None); + assert!(never.cursor.is_none()); + assert!(!never.validation_pending); + + // Once the path supplied evidence, its absence keeps that evidence and + // retries (upstream `expectExistingDatabase`). + let mut previous = cursor; + previous.database_path = missing.to_string_lossy().to_string(); + let kept = resolve_priority_turns(&missing, Some(previous.clone()), 0, true, None); + assert!(kept.validation_pending); + assert_eq!(kept.cursor, Some(previous)); +} + +#[test] +fn unreadable_database_is_not_an_error() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join(CODEX_TRACE_DATABASE_FILE); + std::fs::write(&path, b"not a sqlite database").unwrap(); + let resolution = resolve_priority_turns(&path, None, 0, false, None); + assert!(resolution.cursor.is_none()); + assert!(resolution.validation_pending); +} + +#[test] +fn cursor_for_another_database_is_discarded() { + let db = TraceDb::new(&[(2_000, request_body("turn-1", "priority", "gpt-5.5"))]); + let mut stale = resolve(&db, None).cursor.unwrap(); + stale.database_path = "elsewhere".to_string(); + stale.request_sources.insert( + "ghost".to_string(), + BTreeMap::from([( + 2, + CodexPriorityTurnMetadata { + turn_id: "ghost".to_string(), + ..CodexPriorityTurnMetadata::default() + }, + )]), + ); + let cursor = resolve(&db, Some(stale)).cursor.unwrap(); + assert!(cursor.turn("ghost").is_none()); + assert!(cursor.turn("turn-1").is_some()); +} + +#[test] +fn coverage_window_skips_older_history() { + let db = TraceDb::new(&[ + (100, request_body("ancient", "priority", "gpt-5.5")), + (5_000, request_body("recent", "priority", "gpt-5.5")), + ]); + let cursor = resolve_priority_turns(&db.path, None, 1_000, false, None) + .cursor + .unwrap(); + assert!(cursor.turn("ancient").is_none()); + assert!(cursor.turn("recent").is_some()); +} + +#[test] +fn advancing_coverage_prunes_expired_turns_without_restarting_the_cursor() { + let db = TraceDb::new(&[ + (1_000, request_body("expired", "priority", "gpt-5.5")), + (2_000, request_body("current", "priority", "gpt-5.5")), + ]); + let first = resolve(&db, None).cursor.unwrap(); + let last_row_id = first.last_row_id; + assert!(first.turn("expired").is_some()); + + let updated = resolve_priority_turns(&db.path, Some(first), 1_500, false, None) + .cursor + .unwrap(); + + assert_eq!(updated.last_row_id, last_row_id); + assert_eq!(updated.coverage_since_epoch, 1_500); + assert!(updated.turn("expired").is_none()); + assert!(updated.turn("current").is_some()); +} + +#[test] +fn anchor_digest_tracks_fractional_sqlite_timestamps() { + let db = TraceDb::new(&noise(6)); + let first = resolve(&db, None).cursor.unwrap(); + let anchor = first.anchors[0].clone(); + + db.execute(&format!( + "update logs set ts = ts + 0.25 where id = {}", + anchor.row_id + )); + + let updated = resolve(&db, Some(first)).cursor.unwrap(); + let updated_anchor = updated + .anchors + .iter() + .find(|candidate| candidate.row_id == anchor.row_id) + .unwrap(); + assert_ne!(updated_anchor.digest, anchor.digest); +} + +#[test] +fn pending_completions_are_bounded() { + let mut state = CodexPriorityTurnsCursor::default(); + for index in 0..(CODEX_PRIORITY_COMPLETED_MODEL_RETENTION_LIMIT + 10) { + let row_id = i64::try_from(index).unwrap() + 1; + absorb_row( + &mut state, + row_id, + Some(1), + &completed_body(&format!("turn-{index}"), "gpt-5.5"), + ); + } + assert_eq!( + state.completed_models.len(), + CODEX_PRIORITY_COMPLETED_MODEL_RETENTION_LIMIT + ); + assert!(!state.completed_models.contains_key("turn-0")); +} + +#[test] +fn overlay_is_scoped_to_the_codex_home_of_the_database() { + let db = TraceDb::new(&[(2_000, request_body("turn-1", "priority", "gpt-5.5"))]); + let cursor = resolve(&db, None).cursor.unwrap(); + let home = db.path.parent().unwrap(); + let inside = home.join("sessions").join("2026").join("a.jsonl"); + let outside = home + .parent() + .unwrap() + .join("other-home") + .join("sessions") + .join("a.jsonl"); + let overlay = cursor.overlay_for_file(&inside).unwrap(); + assert!(cursor.overlay_for_file(&outside).is_none()); + + assert_eq!( + overlay.priority_model(Some("turn-1"), "gpt-5.5").as_deref(), + Some("gpt-5.5-priority") + ); + assert_eq!(overlay.priority_model(Some("turn-2"), "gpt-5.5"), None); + assert_eq!(overlay.priority_model(None, "gpt-5.5"), None); + assert_eq!( + overlay.priority_model(Some("turn-1"), "gpt-5.5-priority"), + None + ); +} + +#[test] +fn model_without_a_fast_lane_stays_standard() { + let db = TraceDb::new(&[(2_000, request_body("turn-1", "priority", "gpt-5-mini"))]); + let cursor = resolve(&db, None).cursor.unwrap(); + let inside = db.path.parent().unwrap().join("sessions").join("a.jsonl"); + let overlay = cursor.overlay_for_file(&inside).unwrap(); + assert_eq!(overlay.priority_model(Some("turn-1"), "gpt-5-mini"), None); +} + +fn write_session(sessions: &Path, lines: &[serde_json::Value]) { + let today = Local::now().date_naive(); + let day_dir = sessions + .join(today.format("%Y").to_string()) + .join(today.format("%m").to_string()) + .join(today.format("%d").to_string()); + std::fs::create_dir_all(&day_dir).unwrap(); + let body: String = lines.iter().map(|line| format!("{line}\n")).collect(); + std::fs::write(day_dir.join("priority.jsonl"), body).unwrap(); +} + +fn turn_lines() -> Vec { + let now = Local::now().to_rfc3339(); + let turn = |turn_id: &str, input: i64, output: i64| { + [ + serde_json::json!({ + "timestamp": now, "type": "event_msg", + "payload": {"type": "task_started", "turn_id": turn_id} + }), + serde_json::json!({ + "timestamp": now, "type": "event_msg", + "payload": {"type": "token_count", "info": { + "model": "gpt-5.5", + "total_token_usage": { + "input_tokens": input, "cached_input_tokens": 0, + "output_tokens": output, "reasoning_output_tokens": 0 + } + }} + }), + ] + }; + // Cumulative totals: turn-fast adds 1000/100, turn-std adds 400/40. + let mut lines = turn("turn-fast", 1000, 100).to_vec(); + lines.extend(turn("turn-std", 1400, 140)); + lines +} + +#[test] +fn scan_prices_priority_turns_as_priority_and_others_as_standard() { + let root = tempfile::tempdir().unwrap(); + let sessions = root.path().join("sessions"); + write_session(&sessions, &turn_lines()); + let day = Local::now().date_naive().format("%Y-%m-%d").to_string(); + let now = Local::now().timestamp(); + let trace = TraceDb::in_dir( + root.path(), + &[(now, request_body("turn-fast", "priority", "gpt-5.5"))], + ); + + let scan = |cache: &str, trace_path: Option<&Path>| { + let mut scanner = CostScanner::new(7) + .with_options(CostScanOptions::app_driven()) + .with_cache_root(root.path().join(cache)) + .with_sessions_dirs(vec![sessions.clone()]); + if let Some(path) = trace_path { + scanner = scanner.with_codex_trace_database(path); + } + scanner.scan_codex_detailed_with_cache(None) + }; + + let (baseline, _, _) = scan("cache-plain", None); + let (summary, _, cache) = scan("cache-priority", Some(&trace.path)); + + let models = &cache.days[&day]; + assert_eq!(models["gpt-5.5-priority"][0], 1000); + assert_eq!(models["gpt-5.5"][0], 400); + assert_eq!(summary.input_tokens, baseline.input_tokens); + + let standard = CostUsagePricing::codex_cost_usd("gpt-5.5", 1000, 0, 100).unwrap(); + let fast = CostUsagePricing::codex_fast_cost_usd("gpt-5.5", 1000, 0, 100).unwrap(); + assert!(fast > standard); + assert!((summary.total_cost_usd - baseline.total_cost_usd - (fast - standard)).abs() < 1e-9); +} + +#[test] +fn scan_without_trace_database_keeps_standard_pricing() { + let root = tempfile::tempdir().unwrap(); + let sessions = root.path().join("sessions"); + write_session(&sessions, &turn_lines()); + let day = Local::now().date_naive().format("%Y-%m-%d").to_string(); + let scanner = CostScanner::new(7) + .with_options(CostScanOptions::app_driven()) + .with_cache_root(root.path().join("cache")) + .with_sessions_dirs(vec![sessions]) + .with_codex_trace_database(root.path().join("missing").join(CODEX_TRACE_DATABASE_FILE)); + let (_, _, cache) = scanner.scan_codex_detailed_with_cache(None); + assert!(cache.days[&day].contains_key("gpt-5.5")); + assert!(!cache.days[&day].contains_key("gpt-5.5-priority")); +} diff --git a/rust/src/cost_scanner/codex/priority_trace/tests/scanner.rs b/rust/src/cost_scanner/codex/priority_trace/tests/scanner.rs new file mode 100644 index 0000000000..7214bad13d --- /dev/null +++ b/rust/src/cost_scanner/codex/priority_trace/tests/scanner.rs @@ -0,0 +1,859 @@ +//! Scanner-level Priority pricing, translated from upstream +//! `CostUsageScannerPriorityTests` (v0.65.0), plus the Windows cache edges +//! (fork rows, a vanished database, budget pruning, the trace path). +//! +//! Upstream pins each fixture to a fixed local noon and passes `now`; the +//! Windows scanner always reads a window ending today, so these fixtures sit +//! just before the current time and "time passing" ages the persisted cache. + +use super::*; +use crate::core::{ + CodexPriorityTurnMetadata, CodexPriorityTurnsCursor, CodexSessionLineage, + CodexSourcePricingEvidence, CodexSourceRowCache, CodexSourceUsageRow, CostUsageCache, + CostUsageCacheBudget, CostUsageFileUsage, JsonlScanner, ProviderId, +}; +use crate::cost_scanner::codex::ambient_codex_trace_database_path; +use crate::cost_scanner::codex::scan::{codex_priority_metadata_appeared, save_codex_cache}; +use crate::cost_scanner::{ + CostScanOptions, CostScanStats, CostScanner, CostSummary, ModelPricingCompleteness, +}; +use chrono::{DateTime, Duration, Local}; +use serde_json::{Value, json}; +use std::collections::{BTreeMap, HashMap}; +use std::path::{Path, PathBuf}; + +type DayModels = HashMap>>; + +/// gpt-5.5 Standard and Priority cost of upstream's 100/20/10 request. +const GPT55_STANDARD: f64 = 80.0 * 5e-6 + 20.0 * 5e-7 + 10.0 * 3e-5; +const GPT55_PRIORITY: f64 = 80.0 * 1.25e-5 + 20.0 * 1.25e-6 + 10.0 * 7.5e-5; +/// gpt-5.4 Standard and Priority (2x) cost of the same request. +const GPT54_STANDARD: f64 = 80.0 * 2.5e-6 + 20.0 * 2.5e-7 + 10.0 * 1.5e-5; +const GPT54_PRIORITY: f64 = 80.0 * 5e-6 + 20.0 * 5e-7 + 10.0 * 3e-5; + +/// An isolated Codex home: `sessions/`, `logs_2.sqlite` and a cache root. +struct Env { + root: tempfile::TempDir, + base: DateTime, +} + +impl Env { + fn new() -> Self { + let now = Local::now(); + // Keep every fixture event on today's date and not after `now`. + let day_start = now + .date_naive() + .and_hms_opt(0, 0, 1) + .and_then(|start| start.and_local_timezone(Local).earliest()) + .unwrap_or(now); + Self { + root: tempfile::tempdir().unwrap(), + base: (now - Duration::seconds(30)).max(day_start), + } + } + + fn at(&self, offset: i64) -> DateTime { + self.base + Duration::seconds(offset) + } + + fn iso(&self, offset: i64) -> String { + self.at(offset).to_rfc3339() + } + + fn epoch(&self, offset: i64) -> i64 { + self.at(offset).timestamp() + } + + fn day(&self) -> String { + self.base.format("%Y-%m-%d").to_string() + } + + fn sessions(&self) -> PathBuf { + self.root.path().join("sessions") + } + + fn db_path(&self) -> PathBuf { + self.root.path().join(CODEX_TRACE_DATABASE_FILE) + } + + fn cache_root(&self) -> PathBuf { + self.root.path().join("cache") + } + + fn write_session(&self, name: &str, lines: &[Value]) -> PathBuf { + let day_dir = self + .sessions() + .join(self.base.format("%Y").to_string()) + .join(self.base.format("%m").to_string()) + .join(self.base.format("%d").to_string()); + std::fs::create_dir_all(&day_dir).unwrap(); + let body: String = lines.iter().map(|line| format!("{line}\n")).collect(); + let path = day_dir.join(name); + std::fs::write(&path, body).unwrap(); + path + } + + fn create_db(&self) -> TraceDb { + TraceDb::in_dir(self.root.path(), &[]) + } + + fn scan( + &self, + options: CostScanOptions, + database: &Path, + ) -> (CostSummary, CostScanStats, CostUsageCache) { + CostScanner::new(7) + .with_options(options) + .with_cache_root(self.cache_root()) + .with_sessions_dirs(vec![self.sessions()]) + .with_codex_trace_database(database) + .scan_codex_detailed_with_cache(None) + } + + /// Move the last full scan `secs` into the past, like upstream's later + /// `now` argument, so the next debounced pass may rescan. + fn age_cache(&self, secs: i64) { + let root = self.cache_root(); + let mut cache = JsonlScanner::load_cache(ProviderId::Codex, Some(&root)); + cache.last_scan_unix_ms -= secs * 1000; + JsonlScanner::save_cache(ProviderId::Codex, &mut cache, Some(&root)); + } +} + +fn turn_context(timestamp: &str, model: &str) -> Value { + json!({"type": "turn_context", "timestamp": timestamp, "payload": {"model": model}}) +} + +fn task_started(timestamp: &str, turn: &str) -> Value { + json!({ + "type": "event_msg", "timestamp": timestamp, + "payload": {"type": "task_started", "turn_id": turn} + }) +} + +fn usage(input: i64, cached: i64, output: i64) -> Value { + json!({"input_tokens": input, "cached_input_tokens": cached, "output_tokens": output}) +} + +fn token_count(timestamp: &str, input: i64, cached: i64, output: i64) -> Value { + json!({ + "type": "event_msg", "timestamp": timestamp, + "payload": {"type": "token_count", "info": {"last_token_usage": usage(input, cached, output)}} + }) +} + +fn total_token_count(timestamp: &str, input: i64, cached: i64, output: i64) -> Value { + json!({ + "type": "event_msg", "timestamp": timestamp, + "payload": {"type": "token_count", "info": {"total_token_usage": usage(input, cached, output)}} + }) +} + +/// Upstream `insertPriorityTrace` body. +fn trace_request(turn: &str, model: &str) -> String { + format!( + "thread_id=thread turn.id={turn} websocket request: \ + {{\"type\":\"response.create\",\"model\":\"{model}\",\"service_tier\":\"priority\"}}" + ) +} + +/// Upstream completed-response body. +fn trace_completed(turn: &str, model: &str) -> String { + format!( + "thread_id=thread turn.id={turn} websocket event: \ + {{\"type\":\"response.completed\",\"response\":{{\"model\":\"{model}\"}}}}" + ) +} + +/// One Standard turn, then one Priority turn, each a 100/20/10 request. +fn two_turn_session(env: &Env, model: &str) { + env.write_session( + "session.jsonl", + &[ + turn_context(&env.iso(0), model), + task_started(&env.iso(1), "standard-turn"), + token_count(&env.iso(2), 100, 20, 10), + task_started(&env.iso(3), "priority-turn"), + token_count(&env.iso(3), 100, 20, 10), + ], + ); +} + +/// A single request inside `priority-turn`. +fn priority_session(env: &Env, model: &str, input: i64, cached: i64, output: i64) { + env.write_session( + "session.jsonl", + &[ + turn_context(&env.iso(0), model), + task_started(&env.iso(1), "priority-turn"), + token_count(&env.iso(1), input, cached, output), + ], + ); +} + +/// A long Standard request, then a Priority turn whose first request is +/// over the Fast-lane limit and whose second one is under it. +fn long_context_session(env: &Env, model: &str) { + env.write_session( + "session.jsonl", + &[ + turn_context(&env.iso(0), model), + task_started(&env.iso(1), "standard-turn"), + token_count(&env.iso(1), 272_001, 0, 10), + task_started(&env.iso(2), "priority-turn"), + token_count(&env.iso(2), 300_000, 0, 5), + token_count(&env.iso(3), 100_001, 0, 5), + ], + ); +} + +#[track_caller] +fn assert_cost(actual: f64, expected: f64) { + assert!( + (actual - expected).abs() < 1e-9, + "cost {actual} != expected {expected}" + ); +} + +#[track_caller] +fn assert_exact_cost(actual: f64, expected: f64) { + assert!( + (actual - expected).abs() < 1e-12, + "cost {actual} != expected {expected}" + ); +} + +#[test] +fn astra_priority_traces_price_short_and_long_requests() { + for (input, cached, output, expected) in [ + (100_000, 20_000, 10_000, 2.64), + (300_000, 100_000, 20_000, 11.4), + ] { + let env = Env::new(); + priority_session(&env, "gpt-6-astra", input, cached, output); + let db = env.create_db(); + db.insert(&[(env.epoch(1), trace_request("priority-turn", "gpt-6-astra"))]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + assert_exact_cost(summary.total_cost_usd, expected); + assert_exact_cost(summary.by_speed["fast"], expected); + assert_exact_cost(summary.by_model["gpt-6-astra-priority"], expected); + assert_eq!( + summary.by_speed_tokens["fast"].total(), + u64::try_from(input + output).unwrap() + ); + } +} + +#[test] +fn gpt55_priority_turn_uses_priority_rates() { + let env = Env::new(); + two_turn_session(&env, "gpt-5.5"); + let db = env.create_db(); + db.insert(&[(env.epoch(3), trace_request("priority-turn", "gpt-5.5"))]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + assert_cost(summary.total_cost_usd, GPT55_STANDARD + GPT55_PRIORITY); + assert_cost(summary.by_model["gpt-5.5"], GPT55_STANDARD); + assert_cost(summary.by_model["gpt-5.5-priority"], GPT55_PRIORITY); + assert_cost(summary.by_speed["standard"], GPT55_STANDARD); + assert_cost(summary.by_speed["fast"], GPT55_PRIORITY); + assert_eq!(summary.by_speed_tokens["standard"].total(), 110); + assert_eq!(summary.by_speed_tokens["fast"].total(), 110); +} + +#[test] +fn gpt56_priority_turn_doubles_the_standard_brief_rate() { + // Upstream seeds models.dev with the brief's gpt-5.6-sol rates; the + // built-in table already carries the same 5/0.5/30 per-million rates. + let env = Env::new(); + priority_session(&env, "gpt-5.6-sol", 100_000, 20_000, 20_000); + let db = env.create_db(); + db.insert(&[(env.epoch(1), trace_request("priority-turn", "gpt-5.6-sol"))]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + // The brief's Standard total is $1.01; API Fast is 2x for GPT-5.6. + assert_exact_cost(summary.by_speed["fast"], 2.02); + assert_exact_cost(summary.by_model["gpt-5.6-sol-priority"], 2.02); +} + +#[test] +fn cached_priority_surcharge_survives_a_missing_database() { + let env = Env::new(); + priority_session(&env, "gpt-5.5", 100, 20, 10); + let db = env.create_db(); + db.insert(&[(env.epoch(1), trace_request("priority-turn", "gpt-5.5"))]); + let (first, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + assert_cost(first.total_cost_usd, GPT55_PRIORITY); + + let missing = env.root.path().join("missing.sqlite"); + let (cached, stats, _) = env.scan(CostScanOptions::default(), &missing); + + assert!(stats.used_cache_debounce); + assert_cost(cached.total_cost_usd, GPT55_PRIORITY); + assert_cost(cached.by_speed["fast"], GPT55_PRIORITY); + assert_eq!(cached.by_speed_tokens["fast"].total(), 110); +} + +#[test] +fn appearing_database_bypasses_the_scan_debounce() { + let env = Env::new(); + priority_session(&env, "gpt-5.5", 100, 20, 10); + let db_path = env.db_path(); + let (first, _, cache) = env.scan(CostScanOptions::app_driven(), &db_path); + assert_cost(first.total_cost_usd, GPT55_STANDARD); + assert_eq!( + cache.codex_priority_metadata_key, + Some(format!("missing:{}", db_path.to_string_lossy())) + ); + + let db = env.create_db(); + db.insert(&[(env.epoch(1), trace_request("priority-turn", "gpt-5.5"))]); + let (rescanned, stats, _) = env.scan(CostScanOptions::default(), &db.path); + + assert!(!stats.used_cache_debounce); + assert_cost(rescanned.total_cost_usd, GPT55_PRIORITY); +} + +#[test] +fn unrelated_wal_changes_keep_the_scan_debounce() { + let env = Env::new(); + priority_session(&env, "gpt-5.5", 100, 20, 10); + let db = env.create_db(); + let (first, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + assert_cost(first.total_cost_usd, GPT55_STANDARD); + + let mut wal = db.path.clone().into_os_string(); + wal.push("-wal"); + std::fs::write(PathBuf::from(wal), b"wal-changed").unwrap(); + let (cached, stats, _) = env.scan(CostScanOptions::default(), &db.path); + + assert!(stats.used_cache_debounce); + assert_cost(cached.total_cost_usd, GPT55_STANDARD); +} + +#[test] +fn new_priority_turn_reprices_the_cached_file() { + let env = Env::new(); + priority_session(&env, "gpt-5.5", 100, 20, 10); + let db = env.create_db(); + let (first, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + assert_cost(first.total_cost_usd, GPT55_STANDARD); + + db.insert(&[(env.epoch(1), trace_request("priority-turn", "gpt-5.5"))]); + env.age_cache(61); + let (repriced, stats, _) = env.scan(CostScanOptions::default(), &db.path); + + assert!(!stats.used_cache_debounce); + assert_cost(repriced.total_cost_usd, GPT55_PRIORITY); +} + +#[test] +fn gpt54_priority_turn_uses_priority_rates() { + let env = Env::new(); + two_turn_session(&env, "gpt-5.4"); + let db = env.create_db(); + db.insert(&[(env.epoch(3), trace_request("priority-turn", "gpt-5.4"))]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + assert_cost(summary.total_cost_usd, GPT54_STANDARD + GPT54_PRIORITY); + assert_cost(summary.by_model["gpt-5.4"], GPT54_STANDARD); + assert_cost(summary.by_model["gpt-5.4-priority"], GPT54_PRIORITY); +} + +#[test] +fn priority_alias_is_priced_with_the_completed_response_model() { + let env = Env::new(); + priority_session(&env, "codex-auto-review", 100, 20, 10); + let db = env.create_db(); + db.insert(&[ + ( + env.epoch(1), + trace_request("priority-turn", "codex-auto-review"), + ), + (env.epoch(1), trace_completed("priority-turn", "gpt-5.4")), + ]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + assert_cost(summary.total_cost_usd, GPT54_PRIORITY); + assert_cost(summary.by_model["gpt-5.4-priority"], GPT54_PRIORITY); + assert_cost(summary.by_speed["fast"], GPT54_PRIORITY); + assert_eq!(summary.by_speed_tokens["fast"].total(), 110); + assert_eq!( + summary.model_pricing_completeness, + ModelPricingCompleteness::Complete + ); +} + +#[test] +fn completed_response_model_prices_the_priority_turn() { + let env = Env::new(); + priority_session(&env, "gpt-5.4", 100, 20, 10); + let db = env.create_db(); + db.insert(&[ + (env.epoch(1), trace_request("priority-turn", "gpt-5.4")), + (env.epoch(1), trace_completed("priority-turn", "gpt-5.5")), + ]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + assert_cost(summary.total_cost_usd, GPT55_PRIORITY); + assert_cost(summary.by_model["gpt-5.5-priority"], GPT55_PRIORITY); + assert!(!summary.by_model.contains_key("gpt-5.4-priority")); +} + +#[test] +fn cached_priority_alias_is_repriced_when_the_completion_arrives() { + let env = Env::new(); + priority_session(&env, "codex-auto-review", 100, 20, 10); + let db = env.create_db(); + db.insert(&[( + env.epoch(1), + trace_request("priority-turn", "codex-auto-review"), + )]); + let (first, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + // Upstream reports no total; the Windows summary labels it partial. + assert!(first.total_cost_usd.abs() < 1e-12); + match &first.model_pricing_completeness { + ModelPricingCompleteness::Partial { unpriced_models } => { + assert!( + unpriced_models + .iter() + .any(|model| model == "codex-auto-review"), + "{unpriced_models:?}" + ); + } + ModelPricingCompleteness::Complete => panic!("an unpriced alias must be partial"), + } + + db.insert(&[(env.epoch(1), trace_completed("priority-turn", "gpt-5.4"))]); + env.age_cache(61); + let (repriced, stats, _) = env.scan(CostScanOptions::default(), &db.path); + + assert!(!stats.used_cache_debounce); + assert_cost(repriced.total_cost_usd, GPT54_PRIORITY); + assert_cost(repriced.by_speed["fast"], GPT54_PRIORITY); + assert_eq!(repriced.by_speed_tokens["fast"].total(), 110); + assert_eq!( + repriced.model_pricing_completeness, + ModelPricingCompleteness::Complete + ); +} + +#[test] +fn unpriced_priority_alias_falls_back_to_the_session_model() { + let env = Env::new(); + priority_session(&env, "gpt-5.4", 100, 20, 10); + let db = env.create_db(); + db.insert(&[( + env.epoch(1), + trace_request("priority-turn", "codex-auto-review"), + )]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + assert_cost(summary.total_cost_usd, GPT54_PRIORITY); + assert_cost(summary.by_model["gpt-5.4-priority"], GPT54_PRIORITY); + assert_eq!(summary.by_speed_tokens["fast"].total(), 110); +} + +#[test] +fn missing_trace_database_keeps_the_base_cost() { + let env = Env::new(); + priority_session(&env, "gpt-5.5", 100, 20, 10); + let missing = env.root.path().join("missing.sqlite"); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &missing); + + assert_cost(summary.total_cost_usd, GPT55_STANDARD); + assert_cost(summary.by_model["gpt-5.5"], GPT55_STANDARD); + assert!(!summary.by_model.contains_key("gpt-5.5-priority")); + assert!(!summary.by_speed.contains_key("fast")); + // Windows deviation: the speed split is always derived from the priced + // model, so the base cost shows as Standard; upstream leaves both + // buckets empty when no trace metadata exists. + assert_cost(summary.by_speed["standard"], GPT55_STANDARD); +} + +#[test] +fn priority_row_without_fast_rates_keeps_the_base_cost() { + let env = Env::new(); + priority_session(&env, "gpt-5.4-nano", 100, 20, 10); + let db = env.create_db(); + db.insert(&[(env.epoch(1), trace_request("priority-turn", "gpt-5.4-nano"))]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + let expected = 80.0 * 2e-7 + 20.0 * 2e-8 + 10.0 * 1.25e-6; + assert_cost(summary.total_cost_usd, expected); + assert_cost(summary.by_model["gpt-5.4-nano"], expected); + // Windows deviation: the day cache keys a row by its priced model, so a + // Priority row priced at the base rate lands in the Standard bucket; + // upstream attributes it to the priority bucket. Totals match. + assert_eq!(summary.by_speed_tokens["standard"].total(), 110); +} + +#[test] +fn long_context_priority_rows_skip_the_surcharge() { + let env = Env::new(); + long_context_session(&env, "gpt-5.5"); + let db = env.create_db(); + db.insert(&[(env.epoch(2), trace_request("priority-turn", "gpt-5.5"))]); + + let (_, _, cache) = env.scan(CostScanOptions::app_driven(), &db.path); + + // Only the second Priority request fits the Fast lane. (The Windows + // gpt-5.5 table has no long-context tier, so its cost is asserted on + // gpt-5.6-sol below, which prices both tiers like upstream gpt-5.5.) + let models = &cache.days[&env.day()]; + assert_eq!(models["gpt-5.5"][..3], [572_001, 0, 15]); + assert_eq!(models["gpt-5.5-priority"][..3], [100_001, 0, 5]); +} + +#[test] +fn long_context_rows_use_long_rates_and_short_priority_rows_use_fast_rates() { + let env = Env::new(); + long_context_session(&env, "gpt-5.6-sol"); + let db = env.create_db(); + db.insert(&[(env.epoch(2), trace_request("priority-turn", "gpt-5.6-sol"))]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + let standard_turn = 272_001.0 * 1e-5 + 10.0 * 4.5e-5; + let standard_first_row = 300_000.0 * 1e-5 + 5.0 * 4.5e-5; + let priority_second_row = (100_001.0 * 5e-6 + 5.0 * 3e-5) * 2.0; + assert_cost( + summary.total_cost_usd, + standard_turn + standard_first_row + priority_second_row, + ); + assert_cost(summary.by_speed["fast"], priority_second_row); +} + +#[test] +fn gpt56_long_context_priority_row_keeps_the_long_base_cost() { + let env = Env::new(); + priority_session(&env, "gpt-5.6-sol", 272_001, 100_000, 5); + let db = env.create_db(); + db.insert(&[(env.epoch(1), trace_request("priority-turn", "gpt-5.6-sol"))]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + let expected = 172_001.0 * 1e-5 + 100_000.0 * 1e-6 + 5.0 * 4.5e-5; + assert_cost(summary.total_cost_usd, expected); + assert_cost(summary.by_model["gpt-5.6-sol"], expected); + assert!(!summary.by_speed.contains_key("fast")); + // Same bucket deviation as the gpt-5.4-nano case above. + assert_eq!(summary.by_speed_tokens["standard"].total(), 272_006); +} + +#[test] +fn cached_reads_do_not_count_toward_the_fast_lane_limit() { + let env = Env::new(); + priority_session(&env, "gpt-5.5", 200_000, 100_000, 5); + let db = env.create_db(); + db.insert(&[(env.epoch(1), trace_request("priority-turn", "gpt-5.5"))]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + // Cached input is a subset of input, so the 272K limit applies to the + // 200K input alone and only the 100K uncached input pays the input rate. + let expected = 100_000.0 * 1.25e-5 + 100_000.0 * 1.25e-6 + 5.0 * 7.5e-5; + assert_cost(summary.total_cost_usd, expected); + assert_cost(summary.by_speed["fast"], expected); + assert_eq!(summary.by_speed_tokens["fast"].total(), 200_005); +} + +#[test] +fn cumulative_totals_do_not_trigger_long_context_pricing() { + let env = Env::new(); + env.write_session( + "session.jsonl", + &[ + turn_context(&env.iso(0), "gpt-5.5"), + task_started(&env.iso(1), "standard-turn"), + total_token_count(&env.iso(1), 120_000, 60_000, 100), + total_token_count(&env.iso(2), 240_000, 120_000, 200), + task_started(&env.iso(3), "priority-turn"), + total_token_count(&env.iso(3), 360_000, 180_000, 300), + ], + ); + let db = env.create_db(); + db.insert(&[(env.epoch(3), trace_request("priority-turn", "gpt-5.5"))]); + + let (summary, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + + let standard_row = 60_000.0 * 5e-6 + 60_000.0 * 5e-7 + 100.0 * 3e-5; + let priority_row = 60_000.0 * 1.25e-5 + 60_000.0 * 1.25e-6 + 100.0 * 7.5e-5; + assert_cost(summary.total_cost_usd, 2.0 * standard_row + priority_row); +} + +fn session_meta(timestamp: &str, session_id: &str, forked_from: Option<&str>) -> Value { + let mut payload = json!({"session_id": session_id}); + if let Some(parent) = forked_from { + payload["forked_from_id"] = json!(parent); + } + json!({"type": "session_meta", "timestamp": timestamp, "payload": payload}) +} + +fn model_total(timestamp: &str, input: i64, output: i64) -> Value { + json!({ + "type": "event_msg", "timestamp": timestamp, + "payload": {"type": "token_count", "info": { + "model": "gpt-5.5", + "total_token_usage": { + "input_tokens": input, "cached_input_tokens": 0, "output_tokens": output + } + }} + }) +} + +fn set_modified_ago(path: &Path, secs: u64) { + std::fs::File::options() + .write(true) + .open(path) + .unwrap() + .set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(secs)) + .unwrap(); +} + +#[test] +fn fork_child_growth_is_priced_from_trace_evidence() { + let env = Env::new(); + let parent = env.write_session( + "parent.jsonl", + &[ + session_meta(&env.iso(0), "parent-id", None), + model_total(&env.iso(0), 1_000, 5), + ], + ); + let child = env.write_session( + "child.jsonl", + &[ + session_meta(&env.iso(2), "child-id", Some("parent-id")), + model_total(&env.iso(3), 1_000, 5), + task_started(&env.iso(4), "child-fast"), + model_total(&env.iso(5), 1_400, 45), + ], + ); + set_modified_ago(&parent, 10); + set_modified_ago(&child, 5); + let db = env.create_db(); + db.insert(&[(env.epoch(4), trace_request("child-fast", "gpt-5.5"))]); + let mut options = CostScanOptions::app_driven(); + options.prefer_newest_codex_sessions_first = false; + + let (_, _, cache) = env.scan(options, &db.path); + + // The child inherits the parent's 1000/5 baseline; only its growth + // belongs to the Priority turn. + let models = &cache.days[&env.day()]; + assert_eq!(models["gpt-5.5-priority"][..3], [400, 0, 40]); + assert_eq!(models["gpt-5.5"][..3], [1_000, 0, 5]); +} + +#[test] +fn vanished_database_keeps_its_evidence_and_metadata_key() { + let env = Env::new(); + priority_session(&env, "gpt-5.5", 100, 20, 10); + let db = env.create_db(); + db.insert(&[(env.epoch(1), trace_request("priority-turn", "gpt-5.5"))]); + let (first, _, _) = env.scan(CostScanOptions::app_driven(), &db.path); + assert_cost(first.total_cost_usd, GPT55_PRIORITY); + + std::fs::rename(&db.path, env.root.path().join("moved.sqlite")).unwrap(); + let (rescanned, stats, cache) = env.scan(CostScanOptions::app_driven(), &db.path); + + assert!(!stats.used_cache_debounce); + assert_cost(rescanned.total_cost_usd, GPT55_PRIORITY); + // Validation stays pending, so the key still names the observed database. + assert_eq!( + cache.codex_priority_metadata_key, + Some(format!("sqlite:{}", db.path.to_string_lossy())) + ); +} + +fn file_usage(days: DayModels) -> CostUsageFileUsage { + CostUsageFileUsage { + mtime_unix_ms: 0, + size: 10, + codex_file_identity: None, + days, + parsed_bytes: Some(10), + codex_scan_target_size: None, + last_model: None, + last_totals: None, + codex_token_timestamps_monotonic: None, + codex_last_token_timestamp: None, + codex_session_id: None, + codex_forked_from_id: None, + codex_fork_accounting_state: None, + codex_lineage: CodexSessionLineage::Root, + codex_fork_timestamp: None, + codex_unresolved_fork_parent: false, + } +} + +fn one_day(day: &str, model: &str, tokens: [i64; 3]) -> DayModels { + HashMap::from([( + day.to_string(), + HashMap::from([(model.to_string(), tokens.to_vec())]), + )]) +} + +fn standard_row(day: &str, input: i64, turn: &str) -> CodexSourceUsageRow { + CodexSourceUsageRow { + day_key: day.to_string(), + timestamp: None, + model: "gpt-5.5".to_string(), + input, + cached: 0, + output: 5, + reasoning: None, + source_end_offset: 10, + turn_id: Some(turn.to_string()), + pricing: CodexSourcePricingEvidence { + pricing_model: Some("gpt-5.5".to_string()), + pricing_mode: Some("standard".to_string()), + }, + } +} + +#[track_caller] +fn assert_only_day(days: &DayModels, day: &str, model: &str, tokens: [i64; 3]) { + assert_eq!(days.len(), 1, "{days:?}"); + let models = &days[day]; + assert_eq!(models.len(), 1, "{models:?}"); + assert_eq!(models[model][..3], tokens); +} + +#[test] +fn budget_pruning_rebuilds_priority_days_from_retained_files() { + let root = tempfile::tempdir().unwrap(); + let home = root.path(); + let sessions = home.join("sessions"); + let key = |name: &str| sessions.join(name).to_string_lossy().into_owned(); + let mut cache = CostUsageCache::default(); + for index in 0..CostUsageCacheBudget::MAX_FILE_ENTRIES { + cache.files.insert( + key(&format!("filler-{index}.jsonl")), + file_usage(HashMap::new()), + ); + } + let old = key("old.jsonl"); + let new = key("new.jsonl"); + let mut cursor = CodexPriorityTurnsCursor { + database_path: home + .join(CODEX_TRACE_DATABASE_FILE) + .to_string_lossy() + .into_owned(), + ..CodexPriorityTurnsCursor::default() + }; + for (row_id, path, day, input, turn) in [ + (1, &old, "2026-08-01", 100, "turn-old"), + (2, &new, "2026-09-15", 50, "turn-new"), + ] { + cache.files.insert( + path.clone(), + file_usage(one_day(day, "gpt-5.5", [input, 0, 5])), + ); + cache.codex_source_rows.insert( + path.clone(), + CodexSourceRowCache { + file_identity: String::new(), + size: 10, + mtime_unix_ms: 0, + prefix_hash: 0, + rows: vec![standard_row(day, input, turn)], + }, + ); + cursor.request_sources.insert( + turn.to_string(), + BTreeMap::from([( + row_id, + CodexPriorityTurnMetadata { + turn_id: turn.to_string(), + model: Some("gpt-5.5".to_string()), + ..CodexPriorityTurnMetadata::default() + }, + )]), + ); + cache + .days + .extend(one_day(day, "gpt-5.5-priority", [input, 0, 5])); + } + cache.codex_priority_turns_cursor = Some(cursor); + cache.scan_since_key = Some("2026-09-01".to_string()); + cache.scan_until_key = Some("2026-09-30".to_string()); + let cache_root = home.join("cache"); + + save_codex_cache(&mut cache, Some(&cache_root)); + + // Pruning drops the fillers and the out-of-window file; subtracting the + // old file's plain totals cannot undo its overlay, so the aggregate is + // rebuilt from the retained file instead. + assert_eq!(cache.files.len(), 1); + assert!(cache.files.contains_key(&new)); + assert_only_day(&cache.days, "2026-09-15", "gpt-5.5-priority", [50, 0, 5]); + let reloaded = JsonlScanner::load_cache(ProviderId::Codex, Some(&cache_root)); + assert_only_day(&reloaded.days, "2026-09-15", "gpt-5.5-priority", [50, 0, 5]); +} + +#[test] +fn priority_metadata_appears_only_for_a_new_database() { + let appeared = codex_priority_metadata_appeared; + assert!(!appeared(None, None)); + // The first scan has nothing to compare against. + assert!(!appeared(None, Some("sqlite:/a"))); + assert!(appeared(Some("missing:/a"), Some("sqlite:/a"))); + assert!(appeared(Some("sqlite:/a"), Some("sqlite:/b"))); + assert!(!appeared(Some("sqlite:/a"), Some("sqlite:/a"))); + // A database that went away keeps the cached evidence. + assert!(!appeared(Some("sqlite:/a"), Some("missing:/a"))); + assert!(!appeared(Some("sqlite:/a"), None)); +} + +#[test] +fn trace_database_path_follows_codex_home() { + let root = tempfile::tempdir().unwrap(); + let home = root.path().join("home"); + let codex_home = root.path().join("codex-home"); + let default_path = home.join(".codex").join(CODEX_TRACE_DATABASE_FILE); + + assert_eq!( + ambient_codex_trace_database_path( + Some(format!(" {} ", codex_home.display())), + Some(home.clone()) + ), + Some(codex_home.join(CODEX_TRACE_DATABASE_FILE)) + ); + assert_eq!( + ambient_codex_trace_database_path(Some(" ".to_string()), Some(home.clone())), + Some(default_path.clone()) + ); + assert_eq!( + ambient_codex_trace_database_path(None, Some(home)), + Some(default_path) + ); + assert_eq!(ambient_codex_trace_database_path(None, None), None); + + // Tests and injected session roots never read the ambient database. + assert_eq!(CostScanner::new(7).codex_trace_database_path(), None); + assert_eq!( + CostScanner::new(7) + .with_sessions_dirs(vec![root.path().join("sessions")]) + .codex_trace_database_path(), + None + ); + let configured = root.path().join("configured.sqlite"); + assert_eq!( + CostScanner::new(7) + .with_codex_trace_database(&configured) + .codex_trace_database_path(), + Some(configured) + ); +} diff --git a/rust/src/cost_scanner/codex/scan.rs b/rust/src/cost_scanner/codex/scan.rs index e983adfd6a..04e4507ecb 100644 --- a/rust/src/cost_scanner/codex/scan.rs +++ b/rust/src/cost_scanner/codex/scan.rs @@ -83,6 +83,85 @@ fn key_path(key: &str) -> PathBuf { PathBuf::from(key) } +/// Persist the Codex cache. Budget pruning subtracts each dropped file's +/// plain day totals from `days`, which cannot undo a Priority overlay and +/// would leave stale `-priority` totals beside negative base totals. When +/// pruning dropped a file, rebuild the aggregate from the retained files and +/// persist that instead. +pub(super) fn save_codex_cache(cache: &mut CostUsageCache, cache_root: Option<&Path>) { + let files_before = cache.files.len(); + JsonlScanner::save_cache(ProviderId::Codex, cache, cache_root); + if cache.files.len() != files_before { + rebuild_cache_days(cache); + JsonlScanner::save_cache(ProviderId::Codex, cache, cache_root); + } +} + +/// Path and presence of the configured trace database (upstream +/// `codexPriorityMetadataKey`); `None` when no database is configured. +fn codex_priority_metadata_key(scanner: &CostScanner) -> Option { + let path = scanner.codex_trace_database_path()?; + let state = if path.exists() { "sqlite" } else { "missing" }; + Some(format!("{state}:{}", path.to_string_lossy())) +} + +/// True when the persisted metadata key says `database_path` existed on the +/// last validated scan (upstream `previouslyObservedDatabase`). +fn codex_priority_database_previously_observed( + persisted_key: Option<&str>, + database_path: &Path, +) -> bool { + persisted_key + .and_then(|key| key.strip_prefix("sqlite:")) + .is_some_and(|path| path == database_path.to_string_lossy()) +} + +/// True when a trace database appeared (or moved) since the last full scan. +/// Its evidence must reprice history even inside the debounce window, while +/// a database that went missing keeps the cached evidence (upstream +/// `codexPriorityMetadataChanged`). +pub(super) fn codex_priority_metadata_appeared(old: Option<&str>, new: Option<&str>) -> bool { + matches!((old, new), (Some(old), Some(new)) if old != new && new.starts_with("sqlite:")) +} + +/// Refresh the durable Priority-trace cursor before day totals are rebuilt. +/// A missing or unreadable trace database keeps the previous evidence, so a +/// transient failure never reprices history. Returns true when validation is +/// pending; the caller then keeps the persisted metadata key so the next scan +/// retries (upstream persists nothing for a pending pass). +fn resolve_codex_priority_evidence( + scanner: &CostScanner, + cache: &mut CostUsageCache, + start_date: NaiveDate, + cancel: Option<&AtomicBool>, +) -> bool { + let Some(database_path) = scanner.codex_trace_database_path() else { + return false; + }; + // One day of slack covers the local/UTC offset at the window edge. + let coverage_since_epoch = start_date + .and_hms_opt(0, 0, 0) + .and_then(|midnight| midnight.and_local_timezone(Local).earliest()) + .map_or(0, |start| start.timestamp() - 86_400) + .max(0); + let previously_observed = codex_priority_database_previously_observed( + cache.codex_priority_metadata_key.as_deref(), + &database_path, + ); + let resolution = super::priority_trace::resolve_priority_turns( + &database_path, + cache.codex_priority_turns_cursor.take(), + coverage_since_epoch, + previously_observed, + cancel, + ); + cache.codex_priority_turns_cursor = resolution.cursor; + if resolution.validation_pending { + tracing::debug!("Codex priority trace validation is pending; retrying next scan"); + } + resolution.validation_pending +} + pub(super) fn scan_codex_detailed_with_cache( scanner: &CostScanner, cancel: Option<&AtomicBool>, @@ -100,6 +179,7 @@ pub(super) fn scan_codex_detailed_with_cache( let cache_root = scanner.cache_root.as_deref(); let mut cache = JsonlScanner::load_cache(ProviderId::Codex, cache_root); let sessions_dirs = scanner.get_codex_sessions_dirs(); + let priority_metadata_key = codex_priority_metadata_key(scanner); let pending_scan = CodexPendingScanContext::new( &cache, &range, @@ -132,6 +212,10 @@ pub(super) fn scan_codex_detailed_with_cache( && !cache.codex_scan_incomplete && JsonlScanner::cache_covers_range(&cache, &range) && (!cache.days.is_empty() || !cache.files.is_empty()) + && !codex_priority_metadata_appeared( + cache.codex_priority_metadata_key.as_deref(), + priority_metadata_key.as_deref(), + ) { stats.used_cache_debounce = true; // A16 (upstream 0.48.0): cache hit within debounce = coverage established @@ -331,6 +415,11 @@ pub(super) fn scan_codex_detailed_with_cache( !pruned_paths_pending.is_empty(), bytes_read_this_refresh, ); + let priority_validation_pending = + resolve_codex_priority_evidence(scanner, &mut cache, start_date, cancel); + if !is_cancelled(cancel) && !priority_validation_pending { + cache.codex_priority_metadata_key = priority_metadata_key; + } rebuild_cache_days(&mut cache); cache.last_scan_unix_ms = now_ms; if cache.codex_scan_incomplete { @@ -366,7 +455,7 @@ pub(super) fn scan_codex_detailed_with_cache( cache.previous_report = None; cache.codex_scan_pause_reason = None; } - JsonlScanner::save_cache(ProviderId::Codex, &mut cache, cache_root); + save_codex_cache(&mut cache, cache_root); // Build the current native summary from the complete decoded cache view, // including prior cached files that were not reread in this bounded pass. diff --git a/rust/src/cost_scanner/tests.rs b/rust/src/cost_scanner/tests.rs index 3b16c92a1c..0e19a1940f 100644 --- a/rust/src/cost_scanner/tests.rs +++ b/rust/src/cost_scanner/tests.rs @@ -1561,6 +1561,10 @@ fn codex_source_recovery_keeps_appended_duplicate_unpriced_after_cache_reload() .get_mut(&path_key) .expect("source rows persisted"); assert_eq!(source_rows.rows.len(), 1); + // Cached evidence that disagrees with the source: a Priority row on a + // model with a Fast lane, so the recovered row keeps a `-priority` key + // (a Priority row without a Fast lane would price at its Standard base). + source_rows.rows[0].pricing.pricing_model = Some("gpt-5.5".to_string()); source_rows.rows[0].pricing.pricing_mode = Some("priority".to_string()); first_cache.last_scan_unix_ms = 1; JsonlScanner::save_cache(ProviderId::Codex, &mut first_cache, Some(&cache_root)); @@ -1581,7 +1585,7 @@ fn codex_source_recovery_keeps_appended_duplicate_unpriced_after_cache_reload() let (_, _, second_cache) = scanner.scan_codex_detailed_with_cache(None); let usage = second_cache.files.get(&path_key).expect("file cache"); let day = Local::now().format("%Y-%m-%d").to_string(); - assert_eq!(usage.days[&day]["gpt-5-priority"], vec![100, 0, 5]); + assert_eq!(usage.days[&day]["gpt-5.5-priority"], vec![100, 0, 5]); assert_eq!( usage.days[&day][CostUsagePricing::CODEX_UNATTRIBUTED_MODEL], vec![100, 0, 5]