diff --git a/apps/desktop-tauri/src-tauri/src/commands/spend_contract.rs b/apps/desktop-tauri/src-tauri/src/commands/spend_contract.rs index 4403ba7ad0..65544c559b 100644 --- a/apps/desktop-tauri/src-tauri/src/commands/spend_contract.rs +++ b/apps/desktop-tauri/src-tauri/src/commands/spend_contract.rs @@ -18,6 +18,10 @@ pub async fn get_spend_contract( } let days = history_days.unwrap_or(30); let include_import = include_open_codex.unwrap_or(false) && provider == "codex"; + if include_import { + // Upstream 0.60.4: price imported OpenCodex rows from a fresh catalog. + codexbar::spend_contract::refresh_opencodex_pricing_if_needed().await; + } tauri::async_runtime::spawn_blocking(move || { let history_days = if days == 0 { 365 } else { days.clamp(1, 365) }; let scanner = CostScanner::new(history_days); diff --git a/apps/desktop-tauri/src-tauri/src/commands/usage_spend.rs b/apps/desktop-tauri/src-tauri/src/commands/usage_spend.rs index 0c6df8fb4b..dfd86294f5 100644 --- a/apps/desktop-tauri/src-tauri/src/commands/usage_spend.rs +++ b/apps/desktop-tauri/src-tauri/src/commands/usage_spend.rs @@ -126,6 +126,13 @@ impl UsageSpendCoordinator { self.current = None; true } + + /// The cached summary for `key`, unless the caller forces a rebuild. + fn reusable(&self, key: &str, force_refresh: bool) -> Option<&CachedUsageSpendSummary> { + self.cache + .as_ref() + .filter(|existing| !force_refresh && existing.key == key) + } } static USAGE_SPEND_COORDINATOR: OnceLock> = OnceLock::new(); @@ -191,6 +198,7 @@ pub async fn get_usage_spend_summary( let selected_days = history_days.unwrap_or(30); let force_refresh = force_refresh.unwrap_or(false); + refresh_opencodex_pricing_before_build(&cached, selected_days, force_refresh).await; let built = tauri::async_runtime::spawn_blocking(move || { build_usage_spend_summary_cached(&cached, selected_days, force_refresh) }) @@ -216,6 +224,27 @@ pub async fn get_usage_spend_summary( Ok(built.summary) } +/// Refreshes models.dev prices for the OpenCodex ledger when this request +/// will rebuild the summary (upstream 0.60.4 fresh-load refresh). A cached +/// summary read never starts network work. +async fn refresh_opencodex_pricing_before_build( + cached: &[ProviderUsageSnapshot], + selected_days: u32, + force_refresh: bool, +) { + let settings = codexbar::settings::Settings::load(); + if !settings.open_codex_usage_logs_enabled { + return; + } + let key = usage_spend_cache_key(cached, selected_days, &settings); + let rebuilds = usage_spend_coordinator() + .lock() + .is_ok_and(|coordinator| coordinator.reusable(&key, force_refresh).is_none()); + if rebuilds { + codexbar::spend_contract::refresh_opencodex_pricing_if_needed().await; + } +} + #[tauri::command] pub fn write_usage_spend_export(path: String, payload: String) -> Result<(), String> { const MAX_EXPORT_BYTES: usize = 8 * 1024 * 1024; @@ -240,10 +269,7 @@ fn build_usage_spend_summary_cached( let guard = usage_spend_coordinator() .lock() .map_err(|error| error.to_string())?; - if !force_refresh - && let Some(existing) = guard.cache.as_ref() - && existing.key == key - { + if let Some(existing) = guard.reusable(&key, force_refresh) { return Ok(BuiltUsageSpendSummary { key: existing.key.clone(), summary: existing.summary.clone(), @@ -820,6 +846,49 @@ mod cache_key_tests { ); } + fn empty_summary() -> UsageSpendSummary { + UsageSpendSummary { + rows: Vec::new(), + contract: SpendContract { + provider_id: "codex".to_string(), + history_days: 30, + known_cost_usd: None, + known_zero: false, + provenance: codexbar::spend_contract::CostProvenance::Unknown, + price_coverage: Default::default(), + price_coverage_ratio: None, + history_coverage_established: false, + token_mix: Default::default(), + conversation_count: 0, + models: Vec::new(), + projects: Vec::new(), + conversations: Vec::new(), + daily: Vec::new(), + hourly_activity: Vec::new(), + project_source_status: None, + custom_pricing_active: false, + imports: Vec::new(), + }, + reporting_day: "2026-10-01".to_string(), + dashboard_timezone: "UTC".to_string(), + } + } + + #[test] + fn only_a_summary_rebuild_may_refresh_opencodex_pricing() { + let mut coordinator = UsageSpendCoordinator::default(); + assert!(coordinator.reusable("day|30", false).is_none()); + coordinator.cache = Some(CachedUsageSpendSummary { + key: "day|30".to_string(), + summary: empty_summary(), + refresh_owner: None, + }); + // A cached read starts no network work; a forced or new key rebuilds. + assert!(coordinator.reusable("day|30", false).is_some()); + assert!(coordinator.reusable("day|30", true).is_none()); + assert!(coordinator.reusable("day|7", false).is_none()); + } + #[test] fn privacy_mode_is_part_of_usage_spend_cache_identity() { let public = usage_spend_cache_key_with_privacy(&[], 30, false, false, false); diff --git a/rust/src/cli/cost.rs b/rust/src/cli/cost.rs index 8e9be16dc2..ab2442b7e1 100755 --- a/rust/src/cli/cost.rs +++ b/rust/src/cli/cost.rs @@ -202,6 +202,13 @@ pub async fn run(args: CostArgs) -> anyhow::Result<()> { print_text_output(&results, use_color, args.days, group_by); } OutputFormat::Json => { + // Upstream 0.60.4 refreshes OpenCodex prices before the JSON + // payload; on Windows the import lives in the Codex contract. + if results.iter().any(|result| result.provider == "codex") + && Settings::load().open_codex_usage_logs_enabled + { + crate::spend_contract::refresh_opencodex_pricing_if_needed().await; + } print_json_output(&results, args.pretty, args.days)?; } } diff --git a/rust/src/core/cost_pricing.rs b/rust/src/core/cost_pricing.rs index 461b6c2ff2..c4a8effebd 100755 --- a/rust/src/core/cost_pricing.rs +++ b/rust/src/core/cost_pricing.rs @@ -747,38 +747,73 @@ impl CostUsagePricing { output_tokens: u64, pricing_date: NaiveDate, pricing_snapshot: Option<&models_dev_pricing::ModelsDevPricingSnapshot>, + ) -> Option { + Self::codex_cost_usd_at_date_with_cache_write_and_pricing_snapshot( + model, + input_tokens, + cached_input_tokens, + 0, + output_tokens, + pricing_date, + pricing_snapshot, + ) + } + + /// Codex cost on a historical usage day when the prompt also wrote cache + /// tokens. `input_tokens` is the inclusive prompt size: cache reads and + /// writes are subsets of it. The pre-cutoff GPT-5.6 Terra/Luna rates carry + /// their own 1.25x cache-write rate (upstream `codexHistoricalPricing`). + pub fn codex_cost_usd_at_date_with_cache_write_and_pricing_snapshot( + model: &str, + input_tokens: u64, + cached_input_tokens: u64, + cache_write_input_tokens: u64, + output_tokens: u64, + pricing_date: NaiveDate, + pricing_snapshot: Option<&models_dev_pricing::ModelsDevPricingSnapshot>, ) -> Option { let key = Self::normalize_codex_model(model); let cutoff = NaiveDate::from_ymd_opt(2026, 7, 30).expect("valid pricing cutoff"); if pricing_date < cutoff { let long = input_tokens > codex_pricing::CODEX_LONG_CONTEXT_THRESHOLD; + // (input, cache read, cache write, output) per token. let rates = match (key.as_str(), long) { - ("gpt-5.6-terra", false) => Some((2.5e-6, 2.5e-7, 1.5e-5)), - ("gpt-5.6-terra", true) => Some((5e-6, 5e-7, 2.25e-5)), - ("gpt-5.6-luna", false) => Some((1e-6, 1e-7, 6e-6)), - ("gpt-5.6-luna", true) => Some((2e-6, 2e-7, 9e-6)), + ("gpt-5.6-terra", false) => Some((2.5e-6, 2.5e-7, 3.125e-6, 1.5e-5)), + ("gpt-5.6-terra", true) => Some((5e-6, 5e-7, 6.25e-6, 2.25e-5)), + ("gpt-5.6-luna", false) => Some((1e-6, 1e-7, 1.25e-6, 6e-6)), + ("gpt-5.6-luna", true) => Some((2e-6, 2e-7, 2.5e-6, 9e-6)), _ => None, }; - if let Some((input_rate, cache_rate, output_rate)) = rates { - return Some(codex_pricing::codex_cost_from_rates( + if let Some((input_rate, cache_read_rate, cache_write_rate, output_rate)) = rates { + return Some(codex_pricing::codex_cost_from_rates_with_cache_write( input_tokens, cached_input_tokens, + cache_write_input_tokens, output_tokens, input_rate, - cache_rate, + cache_read_rate, + cache_write_rate, output_rate, )); } } - Self::codex_cost_usd_with_pricing_snapshot( + Self::codex_cost_usd_with_cache_write_and_pricing_snapshot( model, input_tokens, cached_input_tokens, + cache_write_input_tokens, output_tokens, pricing_snapshot, ) } + /// True when the bundled Codex table prices `model`. Windows resolves the + /// bundled rates before any models.dev entry, so such a model never needs + /// a catalog refresh. + pub fn has_bundled_codex_pricing(model: &str) -> bool { + CODEX_PRICING.contains_key(Self::normalize_codex_model(model).as_str()) + } + pub fn codex_fast_cost_usd_at_date( model: &str, input: u64, diff --git a/rust/src/core/cost_pricing/codex.rs b/rust/src/core/cost_pricing/codex.rs index 97ae798f66..09116be6f4 100644 --- a/rust/src/core/cost_pricing/codex.rs +++ b/rust/src/core/cost_pricing/codex.rs @@ -24,7 +24,7 @@ pub(super) fn codex_cost_from_rates( clippy::too_many_arguments, reason = "Arguments mirror independent token classes and their corresponding pricing rates." )] -fn codex_cost_from_rates_with_cache_write( +pub(super) fn codex_cost_from_rates_with_cache_write( input_tokens: u64, cached_input_tokens: u64, cache_write_input_tokens: u64, @@ -95,7 +95,7 @@ impl CostUsagePricing { ) } - fn codex_cost_usd_with_cache_write_and_pricing_snapshot( + pub(super) fn codex_cost_usd_with_cache_write_and_pricing_snapshot( model: &str, input_tokens: u64, cached_input_tokens: u64, diff --git a/rust/src/core/mod.rs b/rust/src/core/mod.rs index 436c7d5e0a..a6b63044f5 100755 --- a/rust/src/core/mod.rs +++ b/rust/src/core/mod.rs @@ -15,6 +15,7 @@ mod http; mod http_proxy; mod jsonl_scanner; mod models_dev_pricing; +mod models_dev_targets; mod openai_dashboard; mod provider; mod provider_factory; @@ -43,6 +44,7 @@ pub use http::*; pub use http_proxy::*; pub use jsonl_scanner::*; pub use models_dev_pricing::*; +pub use models_dev_targets::*; pub use openai_dashboard::*; pub use provider::*; pub use provider_factory::instantiate as instantiate_provider; diff --git a/rust/src/core/models_dev_pricing.rs b/rust/src/core/models_dev_pricing.rs index e35a8d6e84..6e12e94823 100644 --- a/rust/src/core/models_dev_pricing.rs +++ b/rust/src/core/models_dev_pricing.rs @@ -1,214 +1,8 @@ #[cfg(test)] -mod tests { - use super::{ - ModelsDevCache, ModelsDevCacheArtifact, ModelsDevCatalog, ModelsDevRefreshCoordinator, - }; - use std::path::PathBuf; - use std::sync::Arc; - use std::sync::atomic::{AtomicUsize, Ordering}; - use std::time::{Duration, UNIX_EPOCH}; - - #[test] - fn decodes_top_level_provider_map_and_converts_million_token_rates() { - let catalog = ModelsDevCatalog::decode( - r#"{ - "openai": { - "id": "openai", - "models": { - "openai/gpt-fresh": { - "id": "openai/gpt-fresh", - "cost": { - "input": 2.5, - "output": 10, - "cache_read": 0.25, - "cache_write": 3.75, - "context_over_200k": { - "input": 5, - "output": 15, - "cache_read": 0.5, - "cache_write": 7.5 - } - } - } - } - }, - "anthropic": { - "models": { - "claude-fresh": { - "id": "claude-fresh", - "cost": { "input": 3, "output": 15 } - } - } - } - }"#, - ) - .expect("top-level catalog"); - - let pricing = catalog.lookup("openai", "gpt-fresh").expect("pricing"); - assert_eq!(pricing.input_cost_per_token, 2.5e-6); - assert_eq!(pricing.output_cost_per_token, 10e-6); - assert_eq!(pricing.cache_read_input_cost_per_token, Some(0.25e-6)); - assert_eq!(pricing.cache_write_input_cost_per_token, Some(3.75e-6)); - assert_eq!(pricing.threshold_tokens, Some(200_000)); - assert_eq!(pricing.input_cost_per_token_above_threshold, Some(5e-6)); - } - - #[test] - fn decodes_providers_envelope() { - let catalog = ModelsDevCatalog::decode( - r#"{ - "providers": { - "anthropic": { - "id": "anthropic", - "models": { - "claude-fresh": { - "id": "claude-fresh", - "cost": { "input": 3, "output": 15 } - } - } - }, - "openai": { - "models": { - "gpt-fresh": { - "id": "gpt-fresh", - "cost": { "input": 2.5, "output": 10 } - } - } - } - } - }"#, - ) - .expect("enveloped catalog"); - - assert_eq!( - catalog - .lookup("anthropic", "claude-fresh") - .expect("pricing") - .output_cost_per_token, - 15e-6 - ); - } - - #[test] - fn exact_lookup_matches_only_the_trimmed_key_or_model_id() { - let catalog = ModelsDevCatalog::decode( - r#"{ - "Nous": { - "models": { - "z-ai/glm-5": {"id": "z-ai/glm-5", "cost": {"input": 1, "output": 2}}, - "catalog-key": {"id": "deepseek/deepseek-v4", "cost": {"input": 3, "output": 4}}, - "gpt-5": {"id": "gpt-5", "cost": {"input": 5, "output": 6}} - } - } - }"#, - ) - .expect("catalog"); - let input_rate = |model: &str| { - catalog - .lookup_exact(" NOUS ", model) - .map(|pricing| pricing.input_cost_per_token) - }; - - assert_eq!(input_rate(" z-ai/glm-5 "), Some(1e-6)); - assert_eq!(input_rate("deepseek/deepseek-v4"), Some(3e-6)); - assert_eq!(input_rate("Z-AI/GLM-5"), None); - // Aliases the fuzzy lookup resolves never supply an exact price. - for alias in ["z-ai/glm-5@20260101", "z-ai/glm-5-20260101", "openai/gpt-5"] { - assert!(catalog.lookup("nous", alias).is_some(), "{alias}"); - assert_eq!(input_rate(alias), None, "{alias}"); - } - } - - #[test] - fn cache_artifact_is_versioned_and_expires_after_one_day() { - let catalog = ModelsDevCatalog::decode( - r#"{ - "openai": { - "models": { - "gpt-fresh": { - "id": "gpt-fresh", - "cost": { "input": 2.5, "output": 10 } - } - } - }, - "anthropic": { - "models": { - "claude-fresh": { - "id": "claude-fresh", - "cost": { "input": 3, "output": 15 } - } - } - } - }"#, - ) - .expect("catalog"); - let fetched_at = UNIX_EPOCH + Duration::from_secs(1_000_000); - let artifact = ModelsDevCacheArtifact::new(catalog, fetched_at); - - assert_eq!(artifact.version, ModelsDevCache::ARTIFACT_VERSION); - assert!(!artifact.is_stale(fetched_at + Duration::from_secs(86_400))); - assert!(artifact.is_stale(fetched_at + Duration::from_secs(86_401))); - assert_eq!( - ModelsDevCache::cache_path(Some(PathBuf::from("cache-root").as_path())), - PathBuf::from("cache-root") - .join("model-pricing") - .join("models-dev-v1.json") - ); - } - - #[tokio::test] - async fn concurrent_refreshes_for_one_cache_path_share_one_operation() { - let coordinator = ModelsDevRefreshCoordinator::default(); - let calls = Arc::new(AtomicUsize::new(0)); - let first_calls = Arc::clone(&calls); - let path = PathBuf::from("pricing.json"); - let now = UNIX_EPOCH + Duration::from_secs(1_000_000); - - let first = coordinator.refresh(path.clone(), now, async move { - first_calls.fetch_add(1, Ordering::SeqCst); - tokio::time::sleep(Duration::from_millis(10)).await; - true - }); - let second = coordinator.refresh(path, now, async { - panic!("the second caller must await the first operation"); - }); - - assert!(tokio::join!(first, second).0); - assert_eq!(calls.load(Ordering::SeqCst), 1); - } - - #[tokio::test] - async fn failed_refresh_is_not_retried_within_the_attempt_window() { - let coordinator = ModelsDevRefreshCoordinator::default(); - let path = PathBuf::from("pricing.json"); - let now = UNIX_EPOCH + Duration::from_secs(1_000_000); - - assert!( - !coordinator - .refresh(path.clone(), now, async { false }) - .await - ); - assert!( - !coordinator - .refresh(path, now + Duration::from_secs(60), async { - panic!("the 15-minute bound must suppress this attempt"); - }) - .await - ); - } - - #[test] - fn cache_path_uses_the_existing_per_user_cache_root() { - let cache_root = ModelsDevCache::default_cache_root().expect("per-user cache root"); - assert_eq!( - ModelsDevCache::cache_path(None), - cache_root.join("model-pricing").join("models-dev-v1.json") - ); - } -} +mod tests; use serde::{Deserialize, Deserializer, Serialize}; -use std::collections::{HashMap, HashSet}; +use std::collections::{BTreeMap, HashMap, HashSet}; use std::fs; use std::future::Future; use std::path::{Path, PathBuf}; @@ -216,6 +10,8 @@ use std::sync::{Arc, LazyLock, Mutex}; use std::time::{Duration, SystemTime, UNIX_EPOCH}; use tokio::sync::{Mutex as AsyncMutex, watch}; +use super::ModelsDevPricingTarget; + const MODELS_DEV_URL: &str = "https://models.dev/api.json"; const CACHE_TTL: Duration = Duration::from_secs(24 * 60 * 60); const REFRESH_ATTEMPT_WINDOW: Duration = Duration::from_secs(15 * 60); @@ -239,6 +35,15 @@ struct ModelsDevCatalog { providers: HashMap, } +/// How a refresh decides that a model already has a catalog price. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum ModelIdMatch { + /// Dated, `@` and vendor-prefixed aliases may supply the price. + Fuzzy, + /// Only the trimmed catalog key or model id (upstream `exactModelID: true`). + Exact, +} + /// Immutable models.dev view for callers that price many rows in one pass. /// Loading this once avoids repeating cache metadata checks for every row. #[derive(Debug, Clone)] @@ -348,6 +153,18 @@ impl ModelsDevCatalog { }) } + fn lookup_matching( + &self, + provider_id: &str, + model_id: &str, + matching: ModelIdMatch, + ) -> Option { + match matching { + ModelIdMatch::Fuzzy => self.lookup(provider_id, model_id), + ModelIdMatch::Exact => self.lookup_exact(provider_id, model_id), + } + } + fn is_plausible_refresh(&self) -> bool { ["openai", "anthropic"].into_iter().all(|provider_id| { self.providers @@ -818,7 +635,12 @@ static REFRESH_COORDINATOR: LazyLock = /// Loads the cached models.dev catalog once for bulk-pricing callers. pub fn pricing_snapshot() -> ModelsDevPricingSnapshot { - let load = ModelsDevCache::load(SystemTime::now(), None); + pricing_snapshot_at(SystemTime::now(), None) +} + +/// A stale catalog prices nothing: its rows stay unpriced until a refresh. +fn pricing_snapshot_at(now: SystemTime, cache_root: Option<&Path>) -> ModelsDevPricingSnapshot { + let load = ModelsDevCache::load(now, cache_root); let artifact = (!load.is_stale).then_some(load.artifact).flatten(); ModelsDevPricingSnapshot { artifact } } @@ -832,6 +654,15 @@ pub fn lookup(provider_id: &str, model_id: &str) -> Option .and_then(|artifact| artifact.catalog.lookup(provider_id, model_id)) } +/// Why a coordinated refresh runs. A stale-catalog refresh re-checks the +/// cache inside the coordinator, so callers that queued behind a refresh that +/// already landed do not fetch again (upstream `refreshStaleCache`). +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum RefreshReason { + StaleCatalog, + UnknownModels, +} + /// Refreshes the models.dev cache once when supplied models lack cached pricing. /// /// Returns true only if at least one supplied model has pricing after the coordinated refresh. @@ -842,30 +673,89 @@ pub async fn refresh_unknown_models_if_needed( if model_ids.is_empty() { return false; } - refresh_unknown_models_at(provider_id, model_ids, SystemTime::now(), None).await + refresh_unknown_models_at( + provider_id, + model_ids, + ModelIdMatch::Fuzzy, + SystemTime::now(), + None, + fetch_catalog, + ) + .await +} + +/// Refreshes the models.dev cache for exact pricing identities (upstream +/// OpenCodex `refreshPricingIfNeeded`). A stale catalog refreshes first; then +/// each provider whose exact model ids still lack a price may trigger one +/// coordinated refresh, outside the 15-minute attempt window. Cached pricing +/// reads never start this network work; callers run it before a fresh build. +pub async fn refresh_exact_pricing_targets_if_needed(targets: &[ModelsDevPricingTarget]) { + refresh_exact_pricing_targets_with(targets, SystemTime::now(), None, fetch_catalog).await; +} + +async fn refresh_exact_pricing_targets_with( + targets: &[ModelsDevPricingTarget], + now: SystemTime, + cache_root: Option<&Path>, + fetch: F, +) where + F: FnOnce() -> Fut + Clone + Send + 'static, + Fut: Future> + Send + 'static, +{ + if targets.is_empty() { + return; + } + if ModelsDevCache::load(now, cache_root).is_stale { + let _refreshed = + coordinated_refresh(now, cache_root, RefreshReason::StaleCatalog, fetch.clone()).await; + } + let mut grouped: BTreeMap<&str, HashSet> = BTreeMap::new(); + for target in targets { + grouped + .entry(target.provider_id.as_str()) + .or_default() + .insert(target.model_id.clone()); + } + for (provider_id, model_ids) in grouped { + let _priced = refresh_unknown_models_at( + provider_id, + &model_ids, + ModelIdMatch::Exact, + now, + cache_root, + fetch.clone(), + ) + .await; + } } -async fn refresh_unknown_models_at( +async fn refresh_unknown_models_at( provider_id: &str, model_ids: &HashSet, + matching: ModelIdMatch, now: SystemTime, cache_root: Option<&Path>, -) -> bool { - let load = ModelsDevCache::load(now, cache_root); - let unknown_models: Vec = if load.is_stale { - model_ids.iter().cloned().collect() - } else { - model_ids - .iter() - .filter(|model_id| { - load.artifact - .as_ref() - .and_then(|artifact| artifact.catalog.lookup(provider_id, model_id)) - .is_none() + fetch: F, +) -> bool +where + F: FnOnce() -> Fut + Send + 'static, + Fut: Future> + Send + 'static, +{ + // A stale catalog prices nothing (see `pricing_snapshot_at`). + let is_priced = |load: &ModelsDevCacheLoad, model_id: &str| { + !load.is_stale + && load.artifact.as_ref().is_some_and(|artifact| { + artifact + .catalog + .lookup_matching(provider_id, model_id, matching) + .is_some() }) - .cloned() - .collect() }; + let load = ModelsDevCache::load(now, cache_root); + let unknown_models: Vec<&String> = model_ids + .iter() + .filter(|model_id| !is_priced(&load, model_id)) + .collect(); if unknown_models.is_empty() { return true; } @@ -876,46 +766,57 @@ async fn refresh_unknown_models_at( }) { return false; } + drop(load); + + let _refreshed = + coordinated_refresh(now, cache_root, RefreshReason::UnknownModels, fetch).await; + let refreshed = ModelsDevCache::load(now, cache_root); + unknown_models + .iter() + .any(|model_id| is_priced(&refreshed, model_id)) +} + +/// Runs one fetch through the per-cache-path coordinator: concurrent callers +/// share it, and a path that attempted a refresh in the last 15 minutes waits. +async fn coordinated_refresh( + now: SystemTime, + cache_root: Option<&Path>, + reason: RefreshReason, + fetch: F, +) -> bool +where + F: FnOnce() -> Fut + Send + 'static, + Fut: Future> + Send + 'static, +{ let cache_path = ModelsDevCache::cache_path(cache_root); if cache_path.as_os_str().is_empty() { return false; } let cache_root = cache_root.map(Path::to_path_buf); - let refresh_cache_root = cache_root.clone(); - let _ = REFRESH_COORDINATOR + REFRESH_COORDINATOR .refresh(cache_path, now, async move { - refresh_catalog(now, refresh_cache_root.as_deref()).await - }) - .await; - - let refreshed = ModelsDevCache::load(now, cache_root.as_deref()); - !refreshed.is_stale - && unknown_models.iter().any(|model_id| { - refreshed - .artifact - .as_ref() - .and_then(|artifact| artifact.catalog.lookup(provider_id, model_id)) - .is_some() + let cache_root = cache_root.as_deref(); + if reason == RefreshReason::StaleCatalog + && !ModelsDevCache::load(now, cache_root).is_stale + { + return true; + } + match fetch().await { + Some(catalog) => store_refreshed_catalog(catalog, now, cache_root), + None => false, + } }) + .await } -async fn refresh_catalog(now: SystemTime, cache_root: Option<&Path>) -> bool { - let Ok(client) = crate::core::apply_app_proxy(reqwest::Client::builder()) - .timeout(Duration::from_secs(20)) - .build() - else { - return false; - }; - let Ok(response) = client.get(MODELS_DEV_URL).send().await else { - return false; - }; - if !response.status().is_success() { - return false; - } - let Ok(mut catalog) = response.json::().await else { - return false; - }; +/// Saves a plausible fetched catalog, keeping priceable entries that the +/// previous catalog had and the new one dropped. +fn store_refreshed_catalog( + mut catalog: ModelsDevCatalog, + now: SystemTime, + cache_root: Option<&Path>, +) -> bool { if !catalog.is_plausible_refresh() { return false; } @@ -924,3 +825,73 @@ async fn refresh_catalog(now: SystemTime, cache_root: Option<&Path>) -> bool { } ModelsDevCache::save(catalog, now, cache_root) } + +async fn fetch_catalog() -> Option { + let client = crate::core::apply_app_proxy(reqwest::Client::builder()) + .timeout(Duration::from_secs(20)) + .build() + .ok()?; + let response = client.get(MODELS_DEV_URL).send().await.ok()?; + if !response.status().is_success() { + return None; + } + response.json::().await.ok() +} + +/// Saves `json` as the models.dev cache under `cache_root`, as a refresh at +/// `fetched_at` would. +#[cfg(test)] +pub(crate) fn save_catalog_json_for_tests( + json: &str, + fetched_at: SystemTime, + cache_root: &Path, +) -> bool { + ModelsDevCatalog::decode(json) + .is_some_and(|catalog| ModelsDevCache::save(catalog, fetched_at, Some(cache_root))) +} + +#[cfg(test)] +pub(crate) fn models_dev_cache_path_for_tests(cache_root: &Path) -> PathBuf { + ModelsDevCache::cache_path(Some(cache_root)) +} + +#[cfg(test)] +pub(crate) fn pricing_snapshot_for_tests( + now: SystemTime, + cache_root: &Path, +) -> ModelsDevPricingSnapshot { + pricing_snapshot_at(now, Some(cache_root)) +} + +/// Runs the exact-target refresh against `cache_root`; each fetch counts one +/// call and answers `response_json` (`None` is a failed download). +#[cfg(test)] +pub(crate) async fn refresh_exact_pricing_targets_for_tests( + targets: &[ModelsDevPricingTarget], + now: SystemTime, + cache_root: &Path, + response_json: Option, + calls: Arc, +) { + refresh_exact_pricing_targets_with( + targets, + now, + Some(cache_root), + counting_fetch(response_json, calls), + ) + .await; +} + +#[cfg(test)] +fn counting_fetch( + response_json: Option, + calls: Arc, +) -> impl FnOnce() -> std::future::Ready> + Clone + Send + 'static { + move || { + calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + std::future::ready(response_json.as_deref().and_then(ModelsDevCatalog::decode)) + } +} + +#[cfg(test)] +mod exact_refresh_tests; diff --git a/rust/src/core/models_dev_pricing/exact_refresh_tests.rs b/rust/src/core/models_dev_pricing/exact_refresh_tests.rs new file mode 100644 index 0000000000..9237285e36 --- /dev/null +++ b/rust/src/core/models_dev_pricing/exact_refresh_tests.rs @@ -0,0 +1,187 @@ +use super::{ + ModelIdMatch, ModelsDevCache, ModelsDevPricingTarget, counting_fetch, pricing_snapshot_at, + refresh_exact_pricing_targets_with, refresh_unknown_models_at, save_catalog_json_for_tests, +}; +use std::collections::HashSet; +use std::path::Path; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; +use tempfile::tempdir; + +fn now() -> SystemTime { + UNIX_EPOCH + Duration::from_secs(2_000_000_000) +} + +/// A plausible catalog (OpenAI and Anthropic are priceable) plus one +/// OpenAI model under `fixture_key` priced at `fixture_input` per million. +fn catalog(fixture_key: &str, fixture_input: f64) -> String { + format!( + r#"{{ + "openai": {{"models": {{ + "gpt-base": {{"id": "gpt-base", "cost": {{"input": 1, "output": 2}}}}, + "{fixture_key}": {{"id": "{fixture_key}", "cost": {{"input": {fixture_input}, "output": 8}}}} + }}}}, + "anthropic": {{"models": {{ + "claude-base": {{"id": "claude-base", "cost": {{"input": 3, "output": 15}}}} + }}}} + }}"# + ) +} + +fn fixture_input_rate(root: &Path, model_id: &str) -> Option { + pricing_snapshot_at(now(), Some(root)) + .lookup_exact("openai", model_id) + .map(|pricing| pricing.input_cost_per_token) +} + +fn targets(model_id: &str) -> Vec { + vec![ModelsDevPricingTarget { + provider_id: "openai".to_string(), + model_id: model_id.to_string(), + }] +} + +#[tokio::test] +async fn exact_miss_outside_the_attempt_window_refreshes_once() { + let root = tempdir().unwrap(); + let before = now() - Duration::from_secs(901); + assert!(save_catalog_json_for_tests( + &catalog("gpt-other", 1.0), + before, + root.path() + )); + let calls = Arc::new(AtomicUsize::new(0)); + let response = Some(catalog("gpt-fixture", 2.0)); + + let fetch = counting_fetch(response.clone(), Arc::clone(&calls)); + refresh_exact_pricing_targets_with(&targets("gpt-fixture"), now(), Some(root.path()), fetch) + .await; + assert_eq!(calls.load(Ordering::SeqCst), 1); + assert_eq!(fixture_input_rate(root.path(), "gpt-fixture"), Some(2e-6)); + // The replaced catalog keeps the earlier priceable model. + assert_eq!(fixture_input_rate(root.path(), "gpt-other"), Some(1e-6)); + + let fetch = counting_fetch(response, Arc::clone(&calls)); + refresh_exact_pricing_targets_with(&targets("gpt-fixture"), now(), Some(root.path()), fetch) + .await; + assert_eq!(calls.load(Ordering::SeqCst), 1); +} + +#[tokio::test] +async fn a_catalog_fetched_within_fifteen_minutes_is_not_refetched() { + let root = tempdir().unwrap(); + let recent = now() - Duration::from_secs(60); + assert!(save_catalog_json_for_tests( + &catalog("gpt-other", 1.0), + recent, + root.path() + )); + let calls = Arc::new(AtomicUsize::new(0)); + + let fetch = counting_fetch(Some(catalog("gpt-fixture", 2.0)), Arc::clone(&calls)); + refresh_exact_pricing_targets_with(&targets("gpt-fixture"), now(), Some(root.path()), fetch) + .await; + + assert_eq!(calls.load(Ordering::SeqCst), 0); + assert_eq!(fixture_input_rate(root.path(), "gpt-fixture"), None); +} + +#[tokio::test] +async fn a_dated_alias_satisfies_only_a_fuzzy_refresh() { + let root = tempdir().unwrap(); + let before = now() - Duration::from_secs(901); + let cached = catalog("gpt-fixture-20260101", 1.0); + let calls = Arc::new(AtomicUsize::new(0)); + let models = HashSet::from(["gpt-fixture".to_string()]); + + assert!(save_catalog_json_for_tests(&cached, before, root.path())); + let fetch = counting_fetch(Some(catalog("gpt-fixture", 2.0)), Arc::clone(&calls)); + assert!( + refresh_unknown_models_at( + "openai", + &models, + ModelIdMatch::Fuzzy, + now(), + Some(root.path()), + fetch, + ) + .await + ); + assert_eq!(calls.load(Ordering::SeqCst), 0); + + let fetch = counting_fetch(Some(catalog("gpt-fixture", 2.0)), Arc::clone(&calls)); + assert!( + refresh_unknown_models_at( + "openai", + &models, + ModelIdMatch::Exact, + now(), + Some(root.path()), + fetch, + ) + .await + ); + assert_eq!(calls.load(Ordering::SeqCst), 1); + assert_eq!(fixture_input_rate(root.path(), "gpt-fixture"), Some(2e-6)); +} + +#[tokio::test] +async fn a_stale_catalog_is_replaced_before_exact_lookups() { + let root = tempdir().unwrap(); + let stale = now() - Duration::from_secs(90_000); + assert!(save_catalog_json_for_tests( + &catalog("gpt-fixture", 1.0), + stale, + root.path() + )); + assert_eq!(fixture_input_rate(root.path(), "gpt-fixture"), None); + let calls = Arc::new(AtomicUsize::new(0)); + + let fetch = counting_fetch(Some(catalog("gpt-fixture", 2.0)), Arc::clone(&calls)); + refresh_exact_pricing_targets_with(&targets("gpt-fixture"), now(), Some(root.path()), fetch) + .await; + + assert_eq!(calls.load(Ordering::SeqCst), 1); + assert_eq!(fixture_input_rate(root.path(), "gpt-fixture"), Some(2e-6)); +} + +#[tokio::test] +async fn failed_or_implausible_refreshes_keep_the_previous_cache() { + for response in [None, Some(r#"{"openai": {"models": {}}}"#.to_string())] { + let root = tempdir().unwrap(); + let stale = now() - Duration::from_secs(90_000); + assert!(save_catalog_json_for_tests( + &catalog("gpt-other", 1.0), + stale, + root.path() + )); + let cache_path = ModelsDevCache::cache_path(Some(root.path())); + let before = std::fs::read(&cache_path).unwrap(); + let calls = Arc::new(AtomicUsize::new(0)); + + let fetch = counting_fetch(response.clone(), Arc::clone(&calls)); + refresh_exact_pricing_targets_with( + &targets("gpt-fixture"), + now(), + Some(root.path()), + fetch, + ) + .await; + + assert_eq!(calls.load(Ordering::SeqCst), 1, "{response:?}"); + assert_eq!(std::fs::read(&cache_path).unwrap(), before, "{response:?}"); + } +} + +#[tokio::test] +async fn no_targets_never_fetch() { + let root = tempdir().unwrap(); + let calls = Arc::new(AtomicUsize::new(0)); + let fetch = counting_fetch(Some(catalog("gpt-fixture", 2.0)), Arc::clone(&calls)); + + refresh_exact_pricing_targets_with(&[], now(), Some(root.path()), fetch).await; + + assert_eq!(calls.load(Ordering::SeqCst), 0); + assert!(!ModelsDevCache::cache_path(Some(root.path())).exists()); +} diff --git a/rust/src/core/models_dev_pricing/tests.rs b/rust/src/core/models_dev_pricing/tests.rs new file mode 100644 index 0000000000..b8a1f23583 --- /dev/null +++ b/rust/src/core/models_dev_pricing/tests.rs @@ -0,0 +1,205 @@ +use super::{ + ModelsDevCache, ModelsDevCacheArtifact, ModelsDevCatalog, ModelsDevRefreshCoordinator, +}; +use std::path::PathBuf; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::{Duration, UNIX_EPOCH}; + +#[test] +fn decodes_top_level_provider_map_and_converts_million_token_rates() { + let catalog = ModelsDevCatalog::decode( + r#"{ + "openai": { + "id": "openai", + "models": { + "openai/gpt-fresh": { + "id": "openai/gpt-fresh", + "cost": { + "input": 2.5, + "output": 10, + "cache_read": 0.25, + "cache_write": 3.75, + "context_over_200k": { + "input": 5, + "output": 15, + "cache_read": 0.5, + "cache_write": 7.5 + } + } + } + } + }, + "anthropic": { + "models": { + "claude-fresh": { + "id": "claude-fresh", + "cost": { "input": 3, "output": 15 } + } + } + } + }"#, + ) + .expect("top-level catalog"); + + let pricing = catalog.lookup("openai", "gpt-fresh").expect("pricing"); + assert_eq!(pricing.input_cost_per_token, 2.5e-6); + assert_eq!(pricing.output_cost_per_token, 10e-6); + assert_eq!(pricing.cache_read_input_cost_per_token, Some(0.25e-6)); + assert_eq!(pricing.cache_write_input_cost_per_token, Some(3.75e-6)); + assert_eq!(pricing.threshold_tokens, Some(200_000)); + assert_eq!(pricing.input_cost_per_token_above_threshold, Some(5e-6)); +} + +#[test] +fn decodes_providers_envelope() { + let catalog = ModelsDevCatalog::decode( + r#"{ + "providers": { + "anthropic": { + "id": "anthropic", + "models": { + "claude-fresh": { + "id": "claude-fresh", + "cost": { "input": 3, "output": 15 } + } + } + }, + "openai": { + "models": { + "gpt-fresh": { + "id": "gpt-fresh", + "cost": { "input": 2.5, "output": 10 } + } + } + } + } + }"#, + ) + .expect("enveloped catalog"); + + assert_eq!( + catalog + .lookup("anthropic", "claude-fresh") + .expect("pricing") + .output_cost_per_token, + 15e-6 + ); +} + +#[test] +fn exact_lookup_matches_only_the_trimmed_key_or_model_id() { + let catalog = ModelsDevCatalog::decode( + r#"{ + "Nous": { + "models": { + "z-ai/glm-5": {"id": "z-ai/glm-5", "cost": {"input": 1, "output": 2}}, + "catalog-key": {"id": "deepseek/deepseek-v4", "cost": {"input": 3, "output": 4}}, + "gpt-5": {"id": "gpt-5", "cost": {"input": 5, "output": 6}} + } + } + }"#, + ) + .expect("catalog"); + let input_rate = |model: &str| { + catalog + .lookup_exact(" NOUS ", model) + .map(|pricing| pricing.input_cost_per_token) + }; + + assert_eq!(input_rate(" z-ai/glm-5 "), Some(1e-6)); + assert_eq!(input_rate("deepseek/deepseek-v4"), Some(3e-6)); + assert_eq!(input_rate("Z-AI/GLM-5"), None); + // Aliases the fuzzy lookup resolves never supply an exact price. + for alias in ["z-ai/glm-5@20260101", "z-ai/glm-5-20260101", "openai/gpt-5"] { + assert!(catalog.lookup("nous", alias).is_some(), "{alias}"); + assert_eq!(input_rate(alias), None, "{alias}"); + } +} + +#[test] +fn cache_artifact_is_versioned_and_expires_after_one_day() { + let catalog = ModelsDevCatalog::decode( + r#"{ + "openai": { + "models": { + "gpt-fresh": { + "id": "gpt-fresh", + "cost": { "input": 2.5, "output": 10 } + } + } + }, + "anthropic": { + "models": { + "claude-fresh": { + "id": "claude-fresh", + "cost": { "input": 3, "output": 15 } + } + } + } + }"#, + ) + .expect("catalog"); + let fetched_at = UNIX_EPOCH + Duration::from_secs(1_000_000); + let artifact = ModelsDevCacheArtifact::new(catalog, fetched_at); + + assert_eq!(artifact.version, ModelsDevCache::ARTIFACT_VERSION); + assert!(!artifact.is_stale(fetched_at + Duration::from_secs(86_400))); + assert!(artifact.is_stale(fetched_at + Duration::from_secs(86_401))); + assert_eq!( + ModelsDevCache::cache_path(Some(PathBuf::from("cache-root").as_path())), + PathBuf::from("cache-root") + .join("model-pricing") + .join("models-dev-v1.json") + ); +} + +#[tokio::test] +async fn concurrent_refreshes_for_one_cache_path_share_one_operation() { + let coordinator = ModelsDevRefreshCoordinator::default(); + let calls = Arc::new(AtomicUsize::new(0)); + let first_calls = Arc::clone(&calls); + let path = PathBuf::from("pricing.json"); + let now = UNIX_EPOCH + Duration::from_secs(1_000_000); + + let first = coordinator.refresh(path.clone(), now, async move { + first_calls.fetch_add(1, Ordering::SeqCst); + tokio::time::sleep(Duration::from_millis(10)).await; + true + }); + let second = coordinator.refresh(path, now, async { + panic!("the second caller must await the first operation"); + }); + + assert!(tokio::join!(first, second).0); + assert_eq!(calls.load(Ordering::SeqCst), 1); +} + +#[tokio::test] +async fn failed_refresh_is_not_retried_within_the_attempt_window() { + let coordinator = ModelsDevRefreshCoordinator::default(); + let path = PathBuf::from("pricing.json"); + let now = UNIX_EPOCH + Duration::from_secs(1_000_000); + + assert!( + !coordinator + .refresh(path.clone(), now, async { false }) + .await + ); + assert!( + !coordinator + .refresh(path, now + Duration::from_secs(60), async { + panic!("the 15-minute bound must suppress this attempt"); + }) + .await + ); +} + +#[test] +fn cache_path_uses_the_existing_per_user_cache_root() { + let cache_root = ModelsDevCache::default_cache_root().expect("per-user cache root"); + assert_eq!( + ModelsDevCache::cache_path(None), + cache_root.join("model-pricing").join("models-dev-v1.json") + ); +} diff --git a/rust/src/core/models_dev_targets.rs b/rust/src/core/models_dev_targets.rs new file mode 100644 index 0000000000..4ba3a8f2d0 --- /dev/null +++ b/rust/src/core/models_dev_targets.rs @@ -0,0 +1,170 @@ +//! models.dev pricing identities for a recorded billing route (upstream +//! 0.60.4 `ModelsDevPricingTargetResolver`). +//! +//! The recorded provider is the billing provider. A model namespace that +//! names another vendor (`openrouter` + `openai/gpt-5`) is part of the model +//! id on that provider's catalog, never a route to the vendor's own rates. + +/// One provider/model identity that may price a recorded usage row. +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct ModelsDevPricingTarget { + pub provider_id: String, + pub model_id: String, +} + +/// Resolves the models.dev identities for a recorded provider and model. +/// +/// The provider id is trimmed and lowercased (`x-ai` is models.dev `xai`); the +/// model id is trimmed. One leading segment is stripped only when it repeats +/// the provider itself (`openai` + `openai/gpt-5` is `gpt-5`). Two providers +/// also check their canonical models.dev id: `kimi-coding` adds +/// `kimi-for-coding` and `opencode-free` adds `opencode`. Empty or malformed +/// ids (an empty provider, or a model that is empty or starts or ends with +/// `/`) resolve to nothing. +pub fn models_dev_pricing_targets( + provider_id: &str, + model_id: &str, +) -> Vec { + let provider_id = normalized_provider_id(provider_id); + let model_id = model_id.trim(); + if provider_id.is_empty() || !is_valid_model_id(model_id) { + return Vec::new(); + } + let model_id = strip_self_prefix(model_id, &provider_id); + if !is_valid_model_id(model_id) { + return Vec::new(); + } + let alias = match provider_id.as_str() { + "kimi-coding" => Some("kimi-for-coding"), + "opencode-free" => Some("opencode"), + _ => None, + }; + let mut targets = vec![ModelsDevPricingTarget { + provider_id, + model_id: model_id.to_string(), + }]; + if let Some(alias) = alias { + targets.push(ModelsDevPricingTarget { + provider_id: alias.to_string(), + model_id: model_id.to_string(), + }); + } + targets +} + +fn normalized_provider_id(provider_id: &str) -> String { + let normalized = provider_id.trim().to_ascii_lowercase(); + if normalized == "x-ai" { + "xai".to_string() + } else { + normalized + } +} + +fn strip_self_prefix<'a>(model_id: &'a str, provider_id: &str) -> &'a str { + match model_id.split_once('/') { + Some((prefix, remainder)) + if !remainder.is_empty() && normalized_provider_id(prefix) == provider_id => + { + remainder + } + _ => model_id, + } +} + +fn is_valid_model_id(model_id: &str) -> bool { + !model_id.is_empty() && !model_id.starts_with('/') && !model_id.ends_with('/') +} + +#[cfg(test)] +mod tests { + use super::{ModelsDevPricingTarget, models_dev_pricing_targets}; + + fn target(provider_id: &str, model_id: &str) -> ModelsDevPricingTarget { + ModelsDevPricingTarget { + provider_id: provider_id.to_string(), + model_id: model_id.to_string(), + } + } + + #[test] + fn normalizes_the_provider_and_trims_the_model() { + assert_eq!( + models_dev_pricing_targets(" x-ai ", " Grok-4.6 "), + vec![target("xai", "Grok-4.6")] + ); + } + + #[test] + fn strips_only_one_prefix_that_repeats_the_provider() { + assert_eq!( + models_dev_pricing_targets("OpenAI", "openai/GPT-5"), + vec![target("openai", "GPT-5")] + ); + assert_eq!( + models_dev_pricing_targets("x-ai", "x-ai/grok-4.6"), + vec![target("xai", "grok-4.6")] + ); + assert_eq!( + models_dev_pricing_targets("xai", "xai/xai/grok-4.6"), + vec![target("xai", "xai/grok-4.6")] + ); + assert_eq!( + models_dev_pricing_targets("google-vertex", "google-vertex/gemini-2.5-pro"), + vec![target("google-vertex", "gemini-2.5-pro")] + ); + } + + #[test] + fn router_namespace_stays_part_of_the_model_id() { + assert_eq!( + models_dev_pricing_targets("openrouter", "openai/gpt-5"), + vec![target("openrouter", "openai/gpt-5")] + ); + assert_eq!( + models_dev_pricing_targets("OpenRouter", "openrouter/openai/gpt-5"), + vec![target("openrouter", "openai/gpt-5")] + ); + assert_eq!( + models_dev_pricing_targets("private-proxy", "openai/gpt-5"), + vec![target("private-proxy", "openai/gpt-5")] + ); + } + + #[test] + fn canonical_aliases_follow_the_recorded_provider() { + assert_eq!( + models_dev_pricing_targets("kimi-coding", "kimi-coding/k3"), + vec![target("kimi-coding", "k3"), target("kimi-for-coding", "k3")] + ); + assert_eq!( + models_dev_pricing_targets("opencode-free", "opencode-free/opencode"), + vec![ + target("opencode-free", "opencode"), + target("opencode", "opencode") + ] + ); + assert_eq!( + models_dev_pricing_targets("kimi-coding", "k3"), + vec![target("kimi-coding", "k3"), target("kimi-for-coding", "k3")] + ); + } + + #[test] + fn malformed_identities_resolve_to_nothing() { + for (provider_id, model_id) in [ + ("", "gpt-5"), + (" ", "gpt-5"), + ("openai", ""), + ("openai", " / "), + ("openai", "/gpt-5"), + ("openai", "gpt-5/"), + ("openai", "openai/gpt-5/"), + ] { + assert!( + models_dev_pricing_targets(provider_id, model_id).is_empty(), + "{provider_id:?} {model_id:?}" + ); + } + } +} diff --git a/rust/src/spend_contract.rs b/rust/src/spend_contract.rs index af9cdfef34..9c4be32e12 100644 --- a/rust/src/spend_contract.rs +++ b/rust/src/spend_contract.rs @@ -246,18 +246,13 @@ struct CustomPricing { entries: HashMap, } -#[derive(Debug, Clone, Default, Deserialize)] +/// One `custom-pricing.json` entry, in USD per million tokens. `Some(0.0)` is +/// free; `None` is unknown and is never filled from another source. +#[derive(Debug, Clone, Default)] struct CustomRates { input: Option, output: Option, - #[serde(rename = "cacheRead", alias = "cache_read")] cache_read: Option, - #[serde( - rename = "cacheWrite", - alias = "cache_write", - alias = "cacheCreation", - alias = "cache_creation" - )] cache_write: Option, } @@ -269,19 +264,32 @@ impl CustomPricing { fn load() -> Self { Self::default_path() .and_then(|path| fs::read(path).ok()) - .and_then(|bytes| serde_json::from_slice::>(&bytes).ok()) - .map(|entries| Self { - entries: entries - .into_iter() - .filter_map(|(key, rates)| { - let key = key.trim().to_ascii_lowercase(); - (!key.is_empty() && rates.is_valid()).then_some((key, rates)) - }) - .collect(), - }) + .map(|bytes| Self::parse(&bytes)) .unwrap_or_default() } + /// Upstream `CostUsageCustomPricing.parse`: keys are trimmed and + /// lowercased, and every entry is read on its own, so one malformed entry + /// never discards the others. An entry without a single usable rate is + /// dropped. A document that is not a JSON object is empty. + fn parse(bytes: &[u8]) -> Self { + let Ok(serde_json::Value::Object(object)) = serde_json::from_slice(bytes) else { + return Self::default(); + }; + Self { + entries: object + .into_iter() + .filter_map(|(key, value)| { + let key = key.trim().to_ascii_lowercase(); + if key.is_empty() { + return None; + } + CustomRates::from_value(&value).map(|rates| (key, rates)) + }) + .collect(), + } + } + fn rates(&self, provider_id: &str, model: &str) -> Option<&CustomRates> { let model_key = model.trim().to_ascii_lowercase(); let provider_key = format!("{}/{}", provider_id.trim().to_ascii_lowercase(), model_key); @@ -289,14 +297,52 @@ impl CustomPricing { .get(&provider_key) .or_else(|| self.entries.get(&model_key)) } + + /// Upstream `CostUsageCustomPricing.rates(providerID:model:)` for imported + /// ledgers: the bare model key first, then `provider/model`. An empty model + /// has no override. + fn overlay_rates(&self, provider_id: &str, model: &str) -> Option<&CustomRates> { + let model_key = model.trim().to_ascii_lowercase(); + if model_key.is_empty() { + return None; + } + self.entries.get(&model_key).or_else(|| { + let provider_key = format!("{}/{}", provider_id.trim(), model.trim()); + self.entries.get(&provider_key.to_ascii_lowercase()) + }) + } } impl CustomRates { - fn is_valid(&self) -> bool { - [self.input, self.output, self.cache_read, self.cache_write] - .into_iter() - .flatten() - .all(|value| value.is_finite() && value >= 0.0) + /// Upstream `rates(from:)`: a rate that is missing, not a number, + /// negative or non-finite is unknown, and the camelCase spelling wins over + /// the snake_case one. `None` when the entry has no usable rate at all. + fn from_value(value: &serde_json::Value) -> Option { + let object = value.as_object()?; + let rate = |key: &str| { + object + .get(key) + .and_then(serde_json::Value::as_f64) + .filter(|rate| rate.is_finite() && *rate >= 0.0) + }; + let rates = Self { + input: rate("input"), + output: rate("output"), + cache_read: rate("cacheRead").or_else(|| rate("cache_read")), + cache_write: rate("cacheWrite") + .or_else(|| rate("cache_write")) + .or_else(|| rate("cacheCreation")) + .or_else(|| rate("cache_creation")), + }; + [ + rates.input, + rates.output, + rates.cache_read, + rates.cache_write, + ] + .iter() + .any(Option::is_some) + .then_some(rates) } fn cost(&self, counts: &ModelTokenCounts) -> Option { @@ -316,16 +362,22 @@ impl CustomRates { cache_write: u64, ) -> Option { let cached = cache_read.min(input); - let uncached = input.saturating_sub(cached); + self.lane_cost(input.saturating_sub(cached), output, cached, cache_write) + } + + /// Upstream `CostUsageCustomPricing.costUSD(rates:...)`: every token lane is + /// billed on its own, and a lane with tokens but no rate leaves the cost + /// unknown. A missing rate is never filled from another source. + fn lane_cost(&self, input: u64, output: u64, cache_read: u64, cache_write: u64) -> Option { let mut total = 0.0; - if uncached > 0 { - total += uncached as f64 * self.input? / 1_000_000.0; + if input > 0 { + total += input as f64 * self.input? / 1_000_000.0; } if output > 0 { total += output as f64 * self.output? / 1_000_000.0; } - if cached > 0 { - total += cached as f64 * self.cache_read? / 1_000_000.0; + if cache_read > 0 { + total += cache_read as f64 * self.cache_read? / 1_000_000.0; } if cache_write > 0 { total += cache_write as f64 * self.cache_write? / 1_000_000.0; @@ -334,6 +386,17 @@ impl CustomRates { } } +/// Refreshes models.dev prices for the OpenCodex ledger before a fresh +/// Usage & Spend build (upstream 0.60.4 `refreshPricingIfNeeded`). +/// +/// Call it only when a summary will be rebuilt: a cached read must never +/// start network work. It fetches at most once per models.dev cache path per +/// 15 minutes, and only when the catalog is stale or a priced row's exact +/// identity is missing from it. +pub async fn refresh_opencodex_pricing_if_needed() { + opencodex::refresh_pricing_if_needed().await; +} + /// Build a stable accounting contract for a local-log provider. /// means the upstream All-time UI window, bounded to 365 days locally. pub fn build_local_spend_contract( diff --git a/rust/src/spend_contract/opencodex.rs b/rust/src/spend_contract/opencodex.rs index 481956a652..94fd6bba0e 100644 --- a/rust/src/spend_contract/opencodex.rs +++ b/rust/src/spend_contract/opencodex.rs @@ -1,15 +1,17 @@ -use std::collections::{BTreeMap, HashMap, HashSet}; +use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet}; use std::path::PathBuf; -use chrono::{DateTime, Datelike, Duration, Local, TimeZone, Timelike, Utc}; +use chrono::{DateTime, Datelike, Duration, Local, NaiveDate, TimeZone, Timelike, Utc}; use serde::{Deserialize, Serialize}; use serde_json::Value; -use crate::core::CostUsagePricing; +use crate::core::{ + CostUsagePricing, ModelsDevPricingSnapshot, ModelsDevPricingTarget, models_dev_pricing_targets, +}; use super::{ - CostCoverageCounts, CostProvenance, CustomPricing, ImportedSpendSource, SpendActivityCell, - SpendDailyPoint, SpendModelRow, SpendTokenMix, + CostCoverageCounts, CostProvenance, CustomPricing, CustomRates, ImportedSpendSource, + SpendActivityCell, SpendDailyPoint, SpendModelRow, SpendTokenMix, }; #[derive(Debug, Clone, Serialize, Deserialize)] @@ -32,6 +34,8 @@ mod cache; mod nous; #[cfg(test)] mod nous_tests; +#[cfg(test)] +mod pricing_tests; #[derive(Default)] struct ModelAccumulator { @@ -111,11 +115,65 @@ pub(super) fn load_for_subscription( aggregate(entries, Utc::now(), history_days.clamp(1, 365), custom) } +/// Refreshes the models.dev catalog when a ledger row that a Usage & Spend +/// build may import needs a price it lacks (upstream 0.60.4 +/// `OpenCodexUsageStore.refreshPricingIfNeeded`). +pub(super) async fn refresh_pricing_if_needed() { + let targets = tokio::task::spawn_blocking(|| { + let entries = cache::load_entries(&usage_path()?)?; + Some(pricing_targets(&entries, Utc::now())) + }) + .await + .ok() + .flatten() + .unwrap_or_default(); + crate::core::refresh_exact_pricing_targets_if_needed(&targets).await; +} + +/// The exact models.dev identities the cached catalog must price for the +/// rows a build may show. Windows imports only rows routed to a +/// subscription. A direct OpenAI model with bundled Codex rates never reads +/// the catalog, and a model-less row is never priced, so neither refreshes. +fn pricing_targets(entries: &[OpenCodexEntry], now: DateTime) -> Vec { + let mut targets = BTreeSet::new(); + for entry in entries { + if entry.timestamp > now + || !is_priceable_status(entry) + || !matches!(route_entry(entry), RouteTarget::Subscription(_)) + { + continue; + } + let resolved = models_dev_pricing_targets(&pricing_provider(entry), &entry.model); + match resolved.first() { + Some(first) if is_codex_target(first) => { + if !CostUsagePricing::has_bundled_codex_pricing(&first.model_id) + && !CostUsagePricing::is_codex_unattributed_model(&first.model_id) + { + targets.insert(first.clone()); + } + } + _ => targets.extend(resolved), + } + } + targets.into_iter().collect() +} + fn aggregate( entries: Vec, now: DateTime, history_days: u32, custom: &CustomPricing, +) -> Option { + aggregate_with_pricing(entries, now, history_days, custom, None) +} + +/// `pricing_snapshot` pins the models.dev catalog; `None` reads the cached one. +fn aggregate_with_pricing( + entries: Vec, + now: DateTime, + history_days: u32, + custom: &CustomPricing, + pricing_snapshot: Option<&ModelsDevPricingSnapshot>, ) -> Option { let first_day = now.with_timezone(&Local).date_naive() - Duration::days(i64::from(history_days.saturating_sub(1))); @@ -154,7 +212,9 @@ fn aggregate( let mut saw_metered_cost = false; // Upstream 0.55.0 #3136: resolve the dynamic pricing catalog once per // aggregate instead of re-checking its cache metadata for every usage row. - let pricing_snapshot = crate::core::pricing_snapshot(); + let pricing_snapshot = pricing_snapshot + .cloned() + .unwrap_or_else(crate::core::pricing_snapshot); for entry in &entries { // A row without a conversationId is its own session (upstream 0.68.0). @@ -169,7 +229,8 @@ fn aggregate( token_mix.reasoning_tokens = add_optional(token_mix.reasoning_tokens, entry.reasoning_tokens); - let cost = entry_cost(entry, custom, &pricing_snapshot); + let pricing = RowPricing::resolve(entry, custom); + let cost = pricing.cost(entry, &pricing_snapshot); match entry.usage_status.as_str() { "reported" => saw_vendor_provenance = true, "estimated" => saw_list_provenance = true, @@ -235,7 +296,7 @@ fn aggregate( if let Some(cost) = cost { model.cost = Some(model.cost.unwrap_or(0.0) + cost); } - model.custom_pricing |= has_custom_rates(entry, custom); + model.custom_pricing |= pricing.custom.is_some(); } let mut model_rows: Vec<_> = models @@ -310,87 +371,176 @@ fn aggregate( }) } -fn entry_cost( - entry: &OpenCodexEntry, - custom: &CustomPricing, - pricing_snapshot: &crate::core::ModelsDevPricingSnapshot, -) -> Option { - if !matches!(entry.usage_status.as_str(), "reported" | "estimated") { - return None; - } - let has_usage = entry.total_tokens.is_some() - || entry.input_tokens.is_some() - || entry.output_tokens.is_some() - || entry.cache_read_tokens.is_some() - || entry.cache_creation_tokens.is_some(); - if !has_usage { - return None; - } - if route_entry(entry) == RouteTarget::Subscription(nous::SUBSCRIPTION_ID) { - return nous::cost(entry, custom, pricing_snapshot); - } - let input = entry.input_tokens.unwrap_or(0); - let output = entry.output_tokens.unwrap_or(0); - let cache_read = entry.cache_read_tokens.unwrap_or(0); - let cache_write = entry.cache_creation_tokens.unwrap_or(0); - if let Some(rates) = custom.rates(&entry.provider, &entry.model) { - return rates.cost_parts(input, output, cache_read, cache_write); +/// Providers a legacy OpenAI-transport row may name in its model prefix as +/// the billing route (upstream `CostUsagePricing.codexModelsDevProviderIDs`). +const ROUTED_MODEL_PROVIDER_IDS: [&str; 7] = [ + "deepseek", + "kimi-coding", + "kimi-for-coding", + "openai", + "opencode", + "opencode-free", + "opencode-go", +]; + +/// The provider whose catalog prices `entry` (upstream +/// `OpenCodexUsagePricing.providerID(for:)`): the recorded provider, except +/// that a legacy OpenAI-transport row names a known route in its model prefix. +fn pricing_provider(entry: &OpenCodexEntry) -> String { + let provider = entry.provider.trim().to_ascii_lowercase(); + if provider == "openai" + && let Some((prefix, _)) = entry.model.trim().split_once('/') + { + let prefix = prefix.to_ascii_lowercase(); + if ROUTED_MODEL_PROVIDER_IDS.contains(&prefix.as_str()) { + return prefix; + } } - let pricing_model = pricing_model(entry)?; - CostUsagePricing::codex_cost_usd_at_date_with_pricing_snapshot( - &pricing_model, - input, - cache_read, - output, - entry.timestamp.date_naive(), - Some(pricing_snapshot), - ) + provider } -fn pricing_model(entry: &OpenCodexEntry) -> Option { - let target = route_entry(entry); - let model = entry.model.trim(); - match target { - RouteTarget::Subscription("codex") => Some(model.to_string()), - RouteTarget::Subscription("opencodego") => { - Some(format!("opencode/{}", provider_model_id(entry, target))) - } - RouteTarget::Subscription("kimi") => { - Some(format!("kimi/{}", provider_model_id(entry, target))) - } - RouteTarget::Subscription("deepseek") => { - Some(format!("deepseek/{}", provider_model_id(entry, target))) +fn is_priceable_status(entry: &OpenCodexEntry) -> bool { + matches!(entry.usage_status.as_str(), "reported" | "estimated") +} + +/// A direct OpenAI model keeps the Codex pricing convention. +fn is_codex_target(target: &ModelsDevPricingTarget) -> bool { + target.provider_id == "openai" && !target.model_id.contains('/') +} + +/// The token lanes of a priceable row. A row without both input and output +/// is unpriced (upstream `listPriceUSD`); a missing cache lane is zero. +#[derive(Debug, Clone, Copy)] +struct RowTokens { + input: u64, + output: u64, + cache_read: u64, + cache_write: u64, +} + +impl RowTokens { + fn of(entry: &OpenCodexEntry) -> Option { + if !is_priceable_status(entry) { + return None; } - RouteTarget::Subscription(_) | RouteTarget::TokenOnly | RouteTarget::Unknown => None, + Some(Self { + input: entry.input_tokens?, + output: entry.output_tokens?, + cache_read: entry.cache_read_tokens.unwrap_or(0), + cache_write: entry.cache_creation_tokens.unwrap_or(0), + }) } } -/// Whether a custom override covers `entry`, resolved as its cost resolves it. -fn has_custom_rates(entry: &OpenCodexEntry, custom: &CustomPricing) -> bool { - if route_entry(entry) == RouteTarget::Subscription(nous::SUBSCRIPTION_ID) { - return nous::custom_rates(entry, custom).is_some(); - } - custom.rates(&entry.provider, &entry.model).is_some() +/// How one ledger row is priced (upstream 0.60.4 `OpenCodexUsagePricing`): +/// by a custom override, or by the models.dev catalog of its recorded +/// billing route. Another vendor's namespace in the model id never borrows +/// that vendor's rates. +struct RowPricing<'a> { + /// The override that prices the row. It also marks the model row as + /// custom-priced, even when its rates leave the cost unknown. + custom: Option<&'a CustomRates>, + targets: Vec, } -fn provider_model_id(entry: &OpenCodexEntry, target: RouteTarget) -> String { - let model = entry.model.trim(); - let Some((model_prefix, model_tail)) = model.split_once('/') else { - return model.to_string(); - }; - let recorded_provider = entry.provider.trim(); - let prefix_matches_recorded_provider = model_prefix.eq_ignore_ascii_case(recorded_provider) - || (recorded_provider.eq_ignore_ascii_case("kimi-for-coding") - && model_prefix.eq_ignore_ascii_case("kimi-coding")); - let is_legacy_openai_route = - recorded_provider.eq_ignore_ascii_case("openai") && route_provider(model_prefix) == target; - if prefix_matches_recorded_provider || is_legacy_openai_route { - model_tail.to_string() - } else { - model.to_string() +impl<'a> RowPricing<'a> { + /// Overrides resolve in upstream order: the recorded provider and model, + /// then (for a resolvable row) the billing route, then each catalog + /// identity. The first match wins whole; a rate it lacks is never filled + /// from a later override or from the catalog. + fn resolve(entry: &OpenCodexEntry, custom: &'a CustomPricing) -> Self { + let provider = pricing_provider(entry); + let targets = models_dev_pricing_targets(&provider, &entry.model); + let custom = custom + .overlay_rates(&entry.provider, &entry.model) + .or_else(|| { + targets + .first() + .and_then(|_| custom.overlay_rates(&provider, &entry.model)) + }) + .or_else(|| { + targets + .iter() + .find_map(|target| custom.overlay_rates(&target.provider_id, &target.model_id)) + }); + Self { custom, targets } + } + + fn cost(&self, entry: &OpenCodexEntry, snapshot: &ModelsDevPricingSnapshot) -> Option { + let tokens = RowTokens::of(entry)?; + let Some(rates) = self.custom else { + return self.catalog_cost(tokens, entry.timestamp.date_naive(), snapshot); + }; + // Custom rates keep the historical convention that input includes + // cache reads and writes. Nous rows record input without them and + // bill each lane on its own (upstream 0.68.0 Nous fixtures). + let input = if route_entry(entry) == RouteTarget::Subscription(nous::SUBSCRIPTION_ID) { + tokens.input + } else { + tokens + .input + .saturating_sub(tokens.cache_read) + .saturating_sub(tokens.cache_write) + }; + rates.lane_cost(input, tokens.output, tokens.cache_read, tokens.cache_write) + } + + /// Upstream `providerCostUSD`. A direct OpenAI model keeps the Codex + /// convention: inclusive input, request-day rates, bundled rates before + /// the catalog. Every other identity needs an exact catalog entry and + /// bills independent token lanes; a consumed cache lane without its own + /// catalog rate leaves the row unpriced instead of borrowing the input + /// rate. + fn catalog_cost( + &self, + tokens: RowTokens, + day: NaiveDate, + snapshot: &ModelsDevPricingSnapshot, + ) -> Option { + let first = self.targets.first()?; + if is_codex_target(first) { + return CostUsagePricing::codex_cost_usd_at_date_with_cache_write_and_pricing_snapshot( + &first.model_id, + tokens.input, + tokens.cache_read, + tokens.cache_write, + tokens.output, + day, + Some(snapshot), + ); + } + let pricing = self + .targets + .iter() + .find_map(|target| snapshot.lookup_exact(&target.provider_id, &target.model_id))?; + if (tokens.cache_read > 0 && pricing.cache_read_input_cost_per_token.is_none()) + || (tokens.cache_write > 0 && pricing.cache_write_input_cost_per_token.is_none()) + { + return None; + } + let inclusive_input = tokens + .input + .checked_add(tokens.cache_read)? + .checked_add(tokens.cache_write)?; + Some(CostUsagePricing::models_dev_cost_usd( + &pricing, + inclusive_input, + tokens.cache_read, + tokens.cache_write, + tokens.output, + )) } } +#[cfg(test)] +fn entry_cost( + entry: &OpenCodexEntry, + custom: &CustomPricing, + snapshot: &ModelsDevPricingSnapshot, +) -> Option { + RowPricing::resolve(entry, custom).cost(entry, snapshot) +} + fn usage_path() -> Option { if let Ok(home) = std::env::var("OPENCODEX_HOME") { let trimmed = home.trim(); @@ -533,460 +683,4 @@ fn add_optional(left: Option, right: Option) -> Option { } #[cfg(test)] -mod tests { - use super::super::CustomRates; - use super::cache::{load_entries_with_cache, read_cache}; - use super::*; - use std::fs; - - #[test] - fn aggregate_deduplicates_requests_and_applies_history_window() { - let now = DateTime::parse_from_rfc3339("2026-08-19T12:00:00Z") - .unwrap() - .with_timezone(&Utc); - let make = |request_id: &str, timestamp: &str, input: u64| OpenCodexEntry { - request_id: request_id.into(), - timestamp: DateTime::parse_from_rfc3339(timestamp) - .unwrap() - .with_timezone(&Utc), - provider: "openai".into(), - model: "gpt-5".into(), - usage_status: "reported".into(), - conversation_id: Some(request_id.into()), - input_tokens: Some(input), - output_tokens: Some(1), - cache_read_tokens: Some(0), - cache_creation_tokens: None, - reasoning_tokens: None, - total_tokens: Some(input + 1), - }; - let source = aggregate( - vec![ - make("same", "2026-08-18T10:00:00Z", 10), - make("same", "2026-08-18T11:00:00Z", 20), - make("old", "2026-08-01T10:00:00Z", 30), - ], - now, - 7, - &CustomPricing::default(), - ) - .expect("source"); - assert_eq!(source.request_count, 1); - assert_eq!(source.conversation_count, 1); - assert_eq!(source.token_mix.input_tokens, Some(20)); - assert_eq!(source.coverage.priced, 1); - assert!(source.known_cost_usd.is_some()); - assert_eq!(source.provenance, CostProvenance::VendorMetered); - } - - #[test] - fn aggregate_preserves_list_and_mixed_provenance() { - let now = DateTime::parse_from_rfc3339("2026-08-19T12:00:00Z") - .unwrap() - .with_timezone(&Utc); - let mut estimated = entry("openai", "gpt-5"); - estimated.request_id = "estimated".to_string(); - estimated.usage_status = "estimated".to_string(); - let list_only = aggregate(vec![estimated.clone()], now, 30, &CustomPricing::default()) - .expect("list-price source"); - assert_eq!(list_only.provenance, CostProvenance::ListPriceEstimate); - - let mut reported = entry("openai", "gpt-5"); - reported.request_id = "reported".to_string(); - let mixed = aggregate( - vec![reported, estimated], - now, - 30, - &CustomPricing::default(), - ) - .expect("mixed source"); - assert_eq!(mixed.provenance, CostProvenance::Mixed); - } - - #[test] - fn aggregate_preserves_zero_cost_authoritative_provenance() { - let now = DateTime::parse_from_rfc3339("2026-08-19T12:00:00Z") - .unwrap() - .with_timezone(&Utc); - let custom = CustomPricing { - entries: std::collections::HashMap::from([( - "openai/gpt-5".to_string(), - CustomRates { - input: Some(0.0), - output: Some(0.0), - cache_read: Some(0.0), - cache_write: Some(0.0), - }, - )]), - }; - - let reported = aggregate(vec![entry("openai", "gpt-5")], now, 30, &custom) - .expect("zero-cost vendor source"); - assert_eq!(reported.known_cost_usd, Some(0.0)); - assert_eq!(reported.provenance, CostProvenance::VendorMetered); - - let mut estimated_entry = entry("openai", "gpt-5"); - estimated_entry.usage_status = "estimated".to_string(); - let estimated = - aggregate(vec![estimated_entry], now, 30, &custom).expect("zero-cost list source"); - assert_eq!(estimated.known_cost_usd, Some(0.0)); - assert_eq!(estimated.provenance, CostProvenance::ListPriceEstimate); - } - - fn entry(provider: &str, model: &str) -> OpenCodexEntry { - OpenCodexEntry { - request_id: format!("{provider}:{model}"), - timestamp: DateTime::parse_from_rfc3339("2026-07-29T12:00:00Z") - .unwrap() - .with_timezone(&Utc), - provider: provider.to_string(), - model: model.to_string(), - usage_status: "reported".to_string(), - conversation_id: None, - input_tokens: Some(100), - output_tokens: Some(5), - cache_read_tokens: Some(10), - cache_creation_tokens: None, - reasoning_tokens: None, - total_tokens: Some(105), - } - } - - #[test] - fn routes_opencodex_entries_into_subscription_rows() { - assert_eq!( - route_entry(&entry("openai", "gpt-5.6-sol")), - RouteTarget::Subscription("codex") - ); - assert_eq!( - route_entry(&entry("opencode-go", "gpt-5.6-sol")), - RouteTarget::Subscription("opencodego") - ); - assert_eq!( - route_entry(&entry("kimi-coding", "k2p5")), - RouteTarget::Subscription("kimi") - ); - assert_eq!( - route_entry(&entry("deepseek", "deepseek-chat")), - RouteTarget::Subscription("deepseek") - ); - assert_eq!( - route_entry(&entry("opencode-free", "free-model")), - RouteTarget::TokenOnly - ); - } - - #[test] - fn recorded_provider_wins_over_mismatched_model_namespace() { - assert_eq!( - route_entry(&entry("opencode-go", "openai/gpt-5.6-sol")), - RouteTarget::Subscription("opencodego") - ); - assert_eq!( - route_entry(&entry("deepseek", "openai/gpt-5.6-sol")), - RouteTarget::Subscription("deepseek") - ); - } - - #[test] - fn legacy_openai_transport_still_uses_explicit_route() { - assert_eq!( - route_entry(&entry("openai", "opencode-go/deepseek-v4-flash")), - RouteTarget::Subscription("opencodego") - ); - assert_eq!( - pricing_model(&entry("openai", "opencode-go/gpt-5")), - Some("opencode/gpt-5".to_string()) - ); - } - - #[test] - fn pricing_model_uses_routed_vendor_catalog() { - assert_eq!( - pricing_model(&entry("opencode-go", "gpt-5")).as_deref(), - Some("opencode/gpt-5") - ); - assert_eq!( - pricing_model(&entry("kimi-coding", "k2p5")).as_deref(), - Some("kimi/k2p5") - ); - assert_eq!( - pricing_model(&entry("deepseek", "deepseek-chat")).as_deref(), - Some("deepseek/deepseek-chat") - ); - assert_eq!( - pricing_model(&entry("opencode-go", "openai/gpt-5")), - Some("opencode/openai/gpt-5".to_string()) - ); - } - - #[test] - fn unknown_provider_or_namespace_fails_closed_for_routing_and_pricing() { - assert_eq!( - route_entry(&entry("private-proxy", "openai/gpt-5")), - RouteTarget::Unknown - ); - assert_eq!(pricing_model(&entry("private-proxy", "openai/gpt-5")), None); - assert_eq!( - route_entry(&entry("openai", "/gpt-5")), - RouteTarget::Subscription("codex") - ); - assert_eq!( - pricing_model(&entry("openai", "/gpt-5")), - Some("/gpt-5".to_string()) - ); - } - - #[test] - fn opencodex_uses_request_day_for_historical_gpt56_pricing() { - let entry = entry("openai", "gpt-5.6-terra"); - let pricing_snapshot = crate::core::pricing_snapshot(); - let cost = entry_cost(&entry, &CustomPricing::default(), &pricing_snapshot).unwrap(); - let expected = 90.0 * 2.5e-6 + 10.0 * 2.5e-7 + 5.0 * 1.5e-5; - assert!((cost - expected).abs() < 1e-12); - } - - #[test] - fn parser_keeps_reported_token_classes() { - let value = serde_json::json!({ - "requestId": "r1", "timestamp": "2026-08-18T10:00:00Z", "provider": "openai", - "model": "gpt-test", "usageStatus": "reported", "conversationId": "c1", - "usage": {"inputTokens": 10, "outputTokens": 4, "cachedInputTokens": 3, "reasoningOutputTokens": 2} - }); - let entry = parse_line(&value.to_string()).expect("entry"); - assert_eq!(entry.model, "gpt-test"); - assert_eq!(entry.input_tokens, Some(10)); - assert_eq!(entry.output_tokens, Some(4)); - assert_eq!(entry.cache_read_tokens, Some(3)); - assert_eq!(entry.reasoning_tokens, Some(2)); - } - - #[test] - fn parser_normalizes_defaults_and_rejects_malformed_lines() { - let minimal = serde_json::json!({ - "requestId": " r1 ", "model": "gpt-test", "timestamp": "2026-08-18T10:00:00Z", - "usageStatus": " REPORTED ", "usage": {"cacheCreationInputTokens": 7} - }); - let entry = parse_line(&minimal.to_string()).expect("entry"); - assert_eq!(entry.request_id, "r1", "ids are trimmed"); - assert_eq!( - entry.provider, "openai", - "missing provider defaults to openai" - ); - assert_eq!( - entry.usage_status, "reported", - "status is lowercased and trimmed" - ); - assert_eq!(entry.conversation_id, None); - assert_eq!(entry.cache_creation_tokens, Some(7)); - - for malformed in [ - "{}", - r#"{"requestId": "", "model": "m", "timestamp": "2026-08-18T10:00:00Z"}"#, - r#"{"requestId": "r1", "model": " ", "timestamp": "2026-08-18T10:00:00Z"}"#, - r#"{"requestId": "r1", "model": "m"}"#, - "not json at all", - ] { - assert!(parse_line(malformed).is_none(), "rejected: {malformed}"); - } - } - - #[test] - fn incremental_cache_appends_only_newline_terminated_tail() { - let dir = tempfile::tempdir().unwrap(); - let log = dir.path().join("usage.jsonl"); - let cache = dir.path().join("cache.sqlite"); - let row = |id: &str, input: u64| { - format!( - r#"{{"requestId":"{id}","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z","usageStatus":"reported","usage":{{"inputTokens":{input}}}}}"# - ) - }; - - fs::write(&log, format!("{}\n{}\n", row("a", 1), row("b", 2))).unwrap(); - let first = load_entries_with_cache(&log, &cache).unwrap(); - assert_eq!( - first - .iter() - .map(|entry| entry.request_id.as_str()) - .collect::>(), - vec!["a", "b"] - ); - let first_cursor = read_cache(&cache).unwrap().cursor; - - let mut file = fs::OpenOptions::new().append(true).open(&log).unwrap(); - use std::io::Write as _; - writeln!(file, "{}", row("c", 3)).unwrap(); - drop(file); - - let second = load_entries_with_cache(&log, &cache).unwrap(); - assert_eq!( - second - .iter() - .map(|entry| entry.request_id.as_str()) - .collect::>(), - vec!["a", "b", "c"] - ); - let second_cursor = read_cache(&cache).unwrap().cursor; - assert!(second_cursor.parsed_offset > first_cursor.parsed_offset); - } - - #[test] - fn incomplete_trailing_opencodex_record_waits_for_newline() { - let dir = tempfile::tempdir().unwrap(); - let log = dir.path().join("usage.jsonl"); - let cache = dir.path().join("cache.sqlite"); - let complete = r#"{"requestId":"a","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; - let pending = r#"{"requestId":"b","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; - let split = pending.len() / 2; - fs::write(&log, format!("{complete}\n{}", &pending[..split])).unwrap(); - - let first = load_entries_with_cache(&log, &cache).unwrap(); - assert_eq!( - first - .iter() - .map(|entry| entry.request_id.as_str()) - .collect::>(), - vec!["a"] - ); - let cursor = read_cache(&cache).unwrap().cursor; - assert_eq!( - cursor.parsed_offset, - u64::try_from(complete.len() + 1).unwrap() - ); - - let mut file = fs::OpenOptions::new().append(true).open(&log).unwrap(); - use std::io::Write as _; - writeln!(file, "{}", &pending[split..]).unwrap(); - drop(file); - let second = load_entries_with_cache(&log, &cache).unwrap(); - assert_eq!( - second - .iter() - .map(|entry| entry.request_id.as_str()) - .collect::>(), - vec!["a", "b"] - ); - } - - #[test] - fn complete_trailing_opencodex_record_waits_for_newline() { - let dir = tempfile::tempdir().unwrap(); - let log = dir.path().join("usage.jsonl"); - let cache = dir.path().join("cache.sqlite"); - let first = r#"{"requestId":"a","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; - let trailing = r#"{"requestId":"b","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; - fs::write(&log, format!("{first}\n{trailing}")).unwrap(); - - let before_newline = load_entries_with_cache(&log, &cache).unwrap(); - assert_eq!( - before_newline - .iter() - .map(|entry| entry.request_id.as_str()) - .collect::>(), - vec!["a"] - ); - - let mut file = fs::OpenOptions::new().append(true).open(&log).unwrap(); - use std::io::Write as _; - writeln!(file).unwrap(); - drop(file); - - let after_newline = load_entries_with_cache(&log, &cache).unwrap(); - assert_eq!( - after_newline - .iter() - .map(|entry| entry.request_id.as_str()) - .collect::>(), - vec!["a", "b"] - ); - } - - #[test] - fn later_request_id_replaces_cached_entry_without_full_cache_loss() { - let dir = tempfile::tempdir().unwrap(); - let log = dir.path().join("usage.jsonl"); - let cache = dir.path().join("cache.sqlite"); - let row = |id: &str, input: u64| { - format!( - r#"{{"requestId":"{id}","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z","usageStatus":"reported","usage":{{"inputTokens":{input}}}}}"# - ) - }; - fs::write(&log, format!("{}\n{}\n", row("dup", 1), row("keep", 2))).unwrap(); - let _ = load_entries_with_cache(&log, &cache).unwrap(); - let mut file = fs::OpenOptions::new().append(true).open(&log).unwrap(); - use std::io::Write as _; - writeln!(file, "{}", row("dup", 9)).unwrap(); - drop(file); - - let entries = load_entries_with_cache(&log, &cache).unwrap(); - assert_eq!(entries.len(), 2); - assert_eq!( - entries - .iter() - .find(|entry| entry.request_id == "dup") - .unwrap() - .input_tokens, - Some(9) - ); - assert!(entries.iter().any(|entry| entry.request_id == "keep")); - } - - #[test] - fn truncation_invalidates_opencodex_cursor_and_rebuilds() { - let dir = tempfile::tempdir().unwrap(); - let log = dir.path().join("usage.jsonl"); - let cache = dir.path().join("cache.sqlite"); - let old = r#"{"requestId":"old","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; - let replacement = - r#"{"requestId":"new","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; - fs::write(&log, format!("{old}\n{old}\n")).unwrap(); - let _ = load_entries_with_cache(&log, &cache).unwrap(); - fs::write(&log, format!("{replacement}\n")).unwrap(); - - let rebuilt = load_entries_with_cache(&log, &cache).unwrap(); - assert_eq!( - rebuilt - .iter() - .map(|entry| entry.request_id.as_str()) - .collect::>(), - vec!["new"] - ); - } - - #[test] - fn timestamps_parse_rfc3339_epoch_seconds_and_millis() { - let expected = Utc.with_ymd_and_hms(2026, 8, 18, 10, 0, 0).unwrap(); - let rfc3339 = serde_json::json!("2026-08-18T10:00:00Z"); - assert_eq!(parse_timestamp(&rfc3339), Some(expected)); - let epoch_seconds = serde_json::json!(1_787_047_200i64); - assert_eq!(parse_timestamp(&epoch_seconds), Some(expected)); - let epoch_millis = serde_json::json!(1_787_047_200_000f64); - assert_eq!(parse_timestamp(&epoch_millis), Some(expected)); - let numeric_string = serde_json::json!("1787047200.0"); - assert_eq!(parse_timestamp(&numeric_string), Some(expected)); - for invalid in [ - serde_json::json!("not a date"), - serde_json::json!(0), - serde_json::json!(-5.0), - serde_json::Value::Null, - serde_json::json!(true), - ] { - assert!(parse_timestamp(&invalid).is_none(), "rejected: {invalid}"); - } - } - - #[test] - fn nonnegative_u64_accepts_json_numbers_and_bounded_floats() { - assert_eq!(nonnegative_u64(Some(&serde_json::json!(42))), Some(42)); - assert_eq!(nonnegative_u64(Some(&serde_json::json!(12.0))), Some(12)); - // Fractional floats are accepted via `as u64` truncation. - assert_eq!(nonnegative_u64(Some(&serde_json::json!(1.5))), Some(1)); - // `u64::MAX as f64` rounds up to 2^64; f64 spacing there is 4096, so - // +2048.0 rounds back into range. First out-of-range step is +4096.0. - assert_eq!( - nonnegative_u64(Some(&serde_json::json!(u64::MAX as f64 + 4096.0))), - None - ); - assert_eq!(nonnegative_u64(None), None); - } -} +mod tests; diff --git a/rust/src/spend_contract/opencodex/nous.rs b/rust/src/spend_contract/opencodex/nous.rs index dea1a2f16a..f8e4cea706 100644 --- a/rust/src/spend_contract/opencodex/nous.rs +++ b/rust/src/spend_contract/opencodex/nous.rs @@ -1,9 +1,10 @@ //! Nous Portal rows in the OpenCodex ledger (upstream 0.68.0 #4008). //! //! The Hermes extractor records `provider: "nous"` with the exact inference -//! model id (`anthropic/claude-sonnet-4.6`), whatever vendor it names. Pricing -//! never falls back to that vendor's rates: only a custom pricing override or -//! an exact Nous models.dev entry prices a row. +//! model id (`anthropic/claude-sonnet-4.6`), whatever vendor it names. Rows +//! are priced by the shared OpenCodex resolver: a custom override, or the +//! exact Nous models.dev entry. The vendor named in the model id never lends +//! its own rates. //! //! Ledger convention: `inputTokens` excludes cache reads and writes (a row may //! carry more cache-read than input tokens). A custom override bills input, @@ -11,75 +12,4 @@ //! expect; the catalog path adds both back to form the inclusive prompt size //! that upstream `providerCostUSD` prices. -use crate::core::{CostUsagePricing, ModelsDevPricingSnapshot}; - -use super::super::CustomRates; -use super::{CustomPricing, OpenCodexEntry}; - pub(super) const SUBSCRIPTION_ID: &str = "nous"; - -pub(super) fn cost( - entry: &OpenCodexEntry, - custom: &CustomPricing, - pricing_snapshot: &ModelsDevPricingSnapshot, -) -> Option { - // Upstream `listPriceUSD`: a row without both input and output is unpriced. - let input = entry.input_tokens?; - let output = entry.output_tokens?; - let cache_read = entry.cache_read_tokens.unwrap_or(0); - let cache_write = entry.cache_creation_tokens.unwrap_or(0); - if let Some(rates) = custom_rates(entry, custom) { - return rates.cost_parts( - input.checked_add(cache_read)?, - output, - cache_read, - cache_write, - ); - } - let pricing = pricing_snapshot.lookup_exact(SUBSCRIPTION_ID, catalog_model_id(entry)?)?; - // A consumed cache lane the catalog does not price keeps the row unknown; - // it never borrows the input rate. - if (cache_read > 0 && pricing.cache_read_input_cost_per_token.is_none()) - || (cache_write > 0 && pricing.cache_write_input_cost_per_token.is_none()) - { - return None; - } - let inclusive_input = input.checked_add(cache_read)?.checked_add(cache_write)?; - Some(CostUsagePricing::models_dev_cost_usd( - &pricing, - inclusive_input, - cache_read, - cache_write, - output, - )) -} - -/// The custom override for a Nous row, shared by its cost and the model's -/// custom-pricing flag: the recorded identity (`nous/`, then the bare -/// model key), then the catalog id when the model repeats the `nous/` prefix. -pub(super) fn custom_rates<'a>( - entry: &OpenCodexEntry, - custom: &'a CustomPricing, -) -> Option<&'a CustomRates> { - custom - .rates(&entry.provider, &entry.model) - .or_else(|| catalog_model_id(entry).and_then(|model| custom.rates(SUBSCRIPTION_ID, model))) -} - -/// models.dev id of a row recorded under `provider: "nous"`: the model without -/// a repeated `nous/` prefix (upstream `ModelsDevPricingTargetResolver`). -/// Legacy OpenAI-transport rows (`provider: "openai"`, model `nous/`) keep -/// upstream's OpenAI pricing route, which has no Nous catalog: only a custom -/// override prices them. -fn catalog_model_id(entry: &OpenCodexEntry) -> Option<&str> { - if !entry.provider.trim().eq_ignore_ascii_case(SUBSCRIPTION_ID) { - return None; - } - let model = entry.model.trim(); - let model = match model.split_once('/') { - Some((prefix, rest)) if prefix.trim().eq_ignore_ascii_case(SUBSCRIPTION_ID) => rest, - _ => model, - }; - // Every id the resolver rejects before stripping is also rejected here. - (!model.is_empty() && !model.starts_with('/') && !model.ends_with('/')).then_some(model) -} diff --git a/rust/src/spend_contract/opencodex/pricing_tests.rs b/rust/src/spend_contract/opencodex/pricing_tests.rs new file mode 100644 index 0000000000..2122d612ed --- /dev/null +++ b/rust/src/spend_contract/opencodex/pricing_tests.rs @@ -0,0 +1,547 @@ +//! OpenCodex rows priced by their recorded billing route (upstream 0.60.4 +//! `OpenCodexProviderPricingTests`). +//! +//! Windows reads one custom-pricing file, the upstream application overlay. +//! It keeps the historical convention that input includes cache reads and +//! writes, so overlay cases expect the upstream application figures (for +//! example 0.00122, where an upstream caller override bills 0.00142). + +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::{Duration as StdDuration, SystemTime, UNIX_EPOCH}; + +use super::*; +use crate::core::{ + models_dev_cache_path_for_tests, pricing_snapshot_for_tests, + refresh_exact_pricing_targets_for_tests, save_catalog_json_for_tests, +}; + +const NOW_SECONDS: u64 = 2_000_000_000; + +fn now() -> DateTime { + Utc.timestamp_opt(2_000_000_000, 0) + .single() + .expect("fixture time") +} + +fn system_now() -> SystemTime { + UNIX_EPOCH + StdDuration::from_secs(NOW_SECONDS) +} + +fn row(provider: &str, model: &str) -> OpenCodexEntry { + OpenCodexEntry { + request_id: "request".to_string(), + timestamp: now(), + provider: provider.to_string(), + model: model.to_string(), + usage_status: "reported".to_string(), + conversation_id: None, + input_tokens: Some(100), + output_tokens: Some(10), + cache_read_tokens: Some(20), + cache_creation_tokens: None, + reasoning_tokens: None, + total_tokens: Some(110), + } +} + +fn rate_json(rate: Option) -> String { + rate.map_or_else(|| "null".to_string(), |rate| rate.to_string()) +} + +/// The upstream fixture catalog. OpenRouter prices `router_model`. +fn catalog_json(router_model: &str, router_input: f64, cache_write: Option) -> String { + let cache_write = rate_json(cache_write); + format!( + r#"{{ + "openai":{{"models":{{ + "gpt-5.4":{{"id":"gpt-5.4","cost":{{"input":2,"output":8,"cache_read":0.2}}}}, + "gpt-fixture":{{"id":"gpt-fixture","cost":{{"input":2,"output":8,"cache_read":0.2}}}} + }}}}, + "anthropic":{{"models":{{"fixture":{{"id":"fixture","cost":{{"input":1,"output":2}}}}}}}}, + "opencode-go":{{"models":{{"gpt-5.4":{{"id":"gpt-5.4","cost":{{"input":20,"output":80,"cache_read":2}}}}}}}}, + "xai":{{"models":{{"grok-fixture":{{"id":"grok-fixture","cost":{{"input":30,"output":60,"cache_read":3}}}}}}}}, + "openrouter":{{"models":{{"{router_model}":{{"id":"{router_model}","cost":{{"input":{router_input},"output":40,"cache_read":1,"cache_write":{cache_write}}}}}}}}} + }}"# + ) +} + +fn snapshot_of(json: &str) -> ModelsDevPricingSnapshot { + ModelsDevPricingSnapshot::from_catalog_json_for_tests(json).expect("catalog") +} + +fn catalog() -> ModelsDevPricingSnapshot { + snapshot_of(&catalog_json("openai/gpt-5.4", 10.0, None)) +} + +fn custom(json: &str) -> CustomPricing { + CustomPricing::parse(json.as_bytes()) +} + +fn priced( + entries: Vec, + custom: &CustomPricing, + snapshot: &ModelsDevPricingSnapshot, +) -> ImportedSpendSource { + aggregate_with_pricing(entries, now(), 7, custom, Some(snapshot)).expect("source") +} + +fn assert_cost(source: &ImportedSpendSource, expected: f64) { + let cost = source.known_cost_usd.expect("priced"); + assert!((cost - expected).abs() < 1e-10, "{cost} vs {expected}"); + assert_eq!(source.daily[0].cost_usd, source.known_cost_usd); + assert_eq!(source.models[0].cost_usd, source.known_cost_usd); + assert_eq!(source.coverage.unpriced, 0); +} + +fn assert_unpriced(source: &ImportedSpendSource) { + assert_eq!(source.known_cost_usd, None); + assert_eq!(source.daily[0].cost_usd, None); + assert_eq!(source.models[0].cost_usd, None); + assert_eq!(source.coverage.unpriced, 1); +} + +fn pair(provider: &str, model: &str) -> (String, String) { + (provider.to_string(), model.to_string()) +} + +#[test] +fn same_model_is_priced_by_its_recorded_provider() { + let none = CustomPricing::default(); + // Upstream prices a direct `gpt-5.4` from the catalog; Windows keeps its + // bundled Codex rate, so an unbundled OpenAI model stands in. + let direct = priced(vec![row("openai", "gpt-fixture")], &none, &catalog()); + assert_cost(&direct, 0.000_244); + let router = priced(vec![row("openrouter", "openai/gpt-5.4")], &none, &catalog()); + assert_cost(&router, 0.001_42); +} + +#[test] +fn unqualified_subscription_model_uses_its_own_catalog_and_legacy_route_still_works() { + let none = CustomPricing::default(); + let bare = priced(vec![row("opencode-go", "gpt-5.4")], &none, &catalog()); + let legacy = priced( + vec![row("openai", "opencode-go/gpt-5.4")], + &none, + &catalog(), + ); + assert_cost(&bare, 0.002_84); + assert_eq!(bare.known_cost_usd, legacy.known_cost_usd); +} + +#[test] +fn unknown_route_cannot_borrow_openai_prices_or_subscription_attribution() { + let none = CustomPricing::default(); + // Neither the bundled Codex rates (empty catalog) nor the catalog's own + // OpenAI entry may price another route's `openai/` namespace. + let openai_only = snapshot_of(&catalog_json("other-model", 10.0, None)); + for provider in ["openrouter", "private-proxy", "xai", "google"] { + let entry = row(provider, "openai/gpt-5.4"); + for snapshot in [snapshot_of("{}"), openai_only.clone()] { + assert_unpriced(&priced(vec![entry.clone()], &none, &snapshot)); + } + assert_eq!(route_entry(&entry), RouteTarget::Unknown, "{provider}"); + } +} + +#[test] +fn router_namespace_does_not_fall_through_to_a_bare_model_in_the_same_catalog() { + let snapshot = snapshot_of(&catalog_json("gpt-5.4", 10.0, None)); + assert_unpriced(&priced( + vec![row("openrouter", "openai/gpt-5.4")], + &CustomPricing::default(), + &snapshot, + )); +} + +#[test] +fn partial_custom_price_remains_unknown_instead_of_silently_falling_through() { + let partial = custom(r#"{"openrouter/openai/gpt-5.4":{"input":1}}"#); + let source = priced( + vec![row("openrouter", "openai/gpt-5.4")], + &partial, + &catalog(), + ); + assert_unpriced(&source); + assert!(source.models[0].custom_pricing); +} + +#[test] +fn override_entries_without_a_usable_rate_never_block_the_catalog() { + // Upstream drops an entry with no usable rate, so the row keeps its + // catalog price and is not marked custom-priced. + for json in [ + r#"{"openai/gpt-5.4":{}}"#, + r#"{"openai/gpt-5.4":{"input":-1,"output":"free"}}"#, + ] { + let source = priced( + vec![row("openrouter", "openai/gpt-5.4")], + &custom(json), + &catalog(), + ); + assert_cost(&source, 0.001_42); + assert!(!source.models[0].custom_pricing, "{json}"); + } +} + +#[test] +fn custom_pricing_counts_cached_tokens_once_and_partial_usage_remains_unknown() { + let overlay = custom(r#"{"openrouter/openai/gpt-5.4":{"input":10,"output":40,"cacheRead":1}}"#); + let source = priced( + vec![row("openrouter", "openai/gpt-5.4")], + &overlay, + &catalog(), + ); + assert_cost(&source, 0.001_22); + assert!(source.models[0].custom_pricing); + + let mut partial = row("openrouter", "openai/gpt-5.4"); + partial.request_id = "partial".to_string(); + partial.input_tokens = None; + partial.output_tokens = None; + partial.cache_read_tokens = None; + partial.total_tokens = Some(500); + let unpriced = priced(vec![partial], &CustomPricing::default(), &catalog()); + assert_eq!(unpriced.daily[0].total_tokens, Some(500)); + assert_eq!(unpriced.models[0].total_tokens, 500); + assert_unpriced(&unpriced); +} + +#[test] +fn custom_pricing_follows_canonical_provider_alias_after_explicit_observed_override() { + let entry = row("kimi-coding", "kimi-coding/k3"); + let canonical = custom(r#"{"kimi-for-coding/k3":{"input":3,"output":6,"cacheRead":0.3}}"#); + assert_cost( + &priced(vec![entry.clone()], &canonical, &catalog()), + 0.000_306, + ); + + let explicit = custom( + r#"{ + "kimi-coding/k3":{"input":1,"output":2,"cacheRead":0.1}, + "kimi-for-coding/k3":{"input":3,"output":6,"cacheRead":0.3} + }"#, + ); + assert_cost(&priced(vec![entry], &explicit, &catalog()), 0.000_102); +} + +#[test] +fn raw_provider_custom_pricing_wins_before_normalized_targets_without_filling_missing_rates() { + let entry = row("x-ai", "x-ai/grok-fixture"); + let explicit = custom( + r#"{ + "x-ai/grok-fixture":{"input":1,"output":2,"cacheRead":0.1}, + "xai/grok-fixture":{"input":3,"output":6,"cacheRead":0.3} + }"#, + ); + assert_cost( + &priced(vec![entry.clone()], &explicit, &catalog()), + 0.000_102, + ); + + let free = custom(r#"{"x-ai/grok-fixture":{"input":0,"output":0,"cacheRead":0}}"#); + assert_cost(&priced(vec![entry.clone()], &free, &catalog()), 0.0); + + let incomplete = custom( + r#"{ + "x-ai/grok-fixture":{"input":1}, + "xai/grok-fixture":{"input":3,"output":6,"cacheRead":0.3} + }"#, + ); + assert_unpriced(&priced(vec![entry], &incomplete, &catalog())); +} + +#[test] +fn legacy_openai_transport_keeps_recorded_application_price_before_routed_price() { + let entry = row("openai", "opencode-go/gpt-5.4"); + let application = custom( + r#"{ + "openai/opencode-go/gpt-5.4":{"input":1,"output":2,"cacheRead":0.1}, + "gpt-5.4":{"input":3,"output":6,"cacheRead":0.3} + }"#, + ); + assert_cost( + &priced(vec![entry.clone()], &application, &catalog()), + 0.000_102, + ); + + // Without the recorded key, the routed catalog identity's override applies. + let routed = custom(r#"{"gpt-5.4":{"input":3,"output":6,"cacheRead":0.3}}"#); + assert_cost(&priced(vec![entry], &routed, &catalog()), 0.000_306); +} + +#[test] +fn legacy_recorded_application_zero_and_incomplete_rates_block_routed_fallback() { + let entry = row("openai", "opencode-go/gpt-5.4"); + let free = custom( + r#"{ + "openai/opencode-go/gpt-5.4":{"input":0,"output":0,"cacheRead":0}, + "gpt-5.4":{"input":3,"output":6,"cacheRead":0.3} + }"#, + ); + assert_cost(&priced(vec![entry.clone()], &free, &catalog()), 0.0); + + let incomplete = custom( + r#"{ + "openai/opencode-go/gpt-5.4":{"input":1}, + "gpt-5.4":{"input":3,"output":6,"cacheRead":0.3} + }"#, + ); + assert_unpriced(&priced(vec![entry], &incomplete, &catalog())); +} + +#[test] +fn bare_application_overrides_retain_precedence_over_provider_qualified_overrides() { + for (provider, model) in [("openai", "gpt-5.4"), ("openai", "opencode-go/gpt-5.4")] { + let application = custom(&format!( + r#"{{ + "{model}":{{"input":0,"output":0,"cacheRead":0}}, + "{provider}/{model}":{{"input":3,"output":6,"cacheRead":0.3}} + }}"# + )); + assert_cost( + &priced(vec![row(provider, model)], &application, &catalog()), + 0.0, + ); + } +} + +#[test] +fn generic_catalog_row_with_consumed_cache_tokens_and_no_cache_rate_stays_unpriced() { + let snapshot = snapshot_of( + r#"{ + "anthropic":{"models":{"fixture":{"id":"fixture","cost":{"input":1,"output":2}}}}, + "openai":{"models":{"fixture":{"id":"fixture","cost":{"input":1,"output":2}}}}, + "xai":{"models":{"grok-fixture":{"id":"grok-fixture","cost":{"input":2,"output":8}}}} + }"#, + ); + assert_unpriced(&priced( + vec![row("xai", "grok-fixture")], + &CustomPricing::default(), + &snapshot, + )); +} + +#[test] +fn cache_creation_is_priced_once_and_requires_an_explicit_rate() { + let mut entry = row("openrouter", "openai/gpt-5.4"); + entry.request_id = "cache-creation".to_string(); + entry.cache_creation_tokens = Some(30); + entry.total_tokens = Some(160); + let snapshot = snapshot_of( + r#"{"openrouter":{"models":{"openai/gpt-5.4":{"id":"openai/gpt-5.4", + "cost":{"input":10,"output":40,"cache_read":1,"cache_write":5}}}}}"#, + ); + let none = CustomPricing::default(); + assert_cost(&priced(vec![entry.clone()], &none, &snapshot), 0.001_57); + + // The overlay takes cache reads and writes out of the inclusive input. + let overlay = custom( + r#"{"openrouter/openai/gpt-5.4":{"input":10,"output":40,"cacheRead":1,"cacheWrite":5}}"#, + ); + assert_cost(&priced(vec![entry.clone()], &overlay, &snapshot), 0.001_07); + + let unpriced = priced(vec![entry], &none, &catalog()); + assert_eq!(unpriced.daily[0].total_tokens, Some(160)); + assert_unpriced(&unpriced); +} + +#[test] +fn cache_only_rows_preserve_complete_free_and_missing_cache_prices() { + let mut entry = row("openrouter", "openai/gpt-5.4"); + entry.request_id = "cache-only".to_string(); + entry.input_tokens = Some(0); + entry.output_tokens = Some(0); + entry.cache_read_tokens = None; + entry.cache_creation_tokens = Some(30); + entry.total_tokens = None; + for (rate, expected) in [ + (Some(5.0), Some(0.000_15)), + (Some(0.0), Some(0.0)), + (None, None), + ] { + let snapshot = snapshot_of(&catalog_json("openai/gpt-5.4", 10.0, rate)); + let overlay = custom(&format!( + r#"{{"openrouter/openai/gpt-5.4":{{"input":10,"output":40,"cacheRead":1,"cacheWrite":{}}}}}"#, + rate_json(rate) + )); + for custom in [CustomPricing::default(), overlay] { + let source = priced(vec![entry.clone()], &custom, &snapshot); + assert_eq!(source.daily[0].total_tokens, Some(30), "{rate:?}"); + match expected { + Some(expected) => assert_cost(&source, expected), + None => assert_unpriced(&source), + } + } + } +} + +#[test] +fn independent_cache_conversion_overflow_stays_unpriced() { + let mut entry = row("openrouter", "openai/gpt-5.4"); + entry.request_id = "cache-overflow".to_string(); + entry.input_tokens = Some(u64::MAX); + entry.output_tokens = Some(0); + entry.cache_read_tokens = None; + entry.cache_creation_tokens = Some(1); + entry.total_tokens = None; + let snapshot = snapshot_of(&catalog_json("openai/gpt-5.4", 10.0, Some(5.0))); + assert_unpriced(&priced(vec![entry], &CustomPricing::default(), &snapshot)); +} + +#[test] +fn direct_and_legacy_routed_rows_keep_the_historical_application_convention() { + for model in ["gpt-5.4", "opencode-go/gpt-5.4"] { + let overlay = custom(&format!( + r#"{{"{model}":{{"input":10,"output":40,"cacheRead":1}}}}"# + )); + assert_cost( + &priced(vec![row("openai", model)], &overlay, &catalog()), + 0.001_22, + ); + } +} + +#[test] +fn refresh_targets_cover_only_imported_rows_that_read_the_catalog() { + let mut future = row("opencode-go", "future-model"); + future.timestamp = now() + Duration::seconds(1); + let mut unsupported = row("opencode-go", "unsupported-model"); + unsupported.usage_status = "unsupported".to_string(); + let entries = [ + // Not imported on Windows: unknown and token-only routes. + row("openrouter", "openai/gpt-5.4"), + row("opencode-free", "free-model"), + // Never read the catalog: bundled Codex rates and model-less rows. + row("openai", "gpt-5.4"), + row("openai", "unknown"), + // Never priced: unsupported status, or not yet recorded. + unsupported, + future, + row("openai", "gpt-fixture"), + row("opencode-go", "gpt-5.4"), + row("openai", "opencode-go/gpt-5.4"), + row("kimi-coding", "kimi-coding/k3"), + ]; + let targets: Vec<_> = pricing_targets(&entries, now()) + .into_iter() + .map(|target| (target.provider_id, target.model_id)) + .collect(); + assert_eq!( + targets, + vec![ + pair("kimi-coding", "k3"), + pair("kimi-for-coding", "k3"), + pair("openai", "gpt-fixture"), + pair("opencode-go", "gpt-5.4"), + ] + ); +} + +/// A plausible catalog whose OpenCode Go entry is `go_model` at `go_input`. +fn go_catalog_json(go_model: &str, go_input: f64) -> String { + format!( + r#"{{ + "openai":{{"models":{{"gpt-5.4":{{"id":"gpt-5.4","cost":{{"input":2,"output":8,"cache_read":0.2}}}}}}}}, + "anthropic":{{"models":{{"fixture":{{"id":"fixture","cost":{{"input":1,"output":2}}}}}}}}, + "opencode-go":{{"models":{{"{go_model}":{{"id":"{go_model}","cost":{{"input":{go_input},"output":80,"cache_read":2}}}}}}}} + }}"# + ) +} + +fn cached_cost(entries: &[OpenCodexEntry], root: &std::path::Path) -> Option { + let snapshot = pricing_snapshot_for_tests(system_now(), root); + aggregate_with_pricing( + entries.to_vec(), + now(), + 7, + &CustomPricing::default(), + Some(&snapshot), + ) + .expect("source") + .known_cost_usd +} + +#[tokio::test] +async fn fresh_catalog_miss_refreshes_the_route_and_reprices() { + let root = tempfile::tempdir().unwrap(); + let before = system_now() - StdDuration::from_secs(901); + assert!(save_catalog_json_for_tests( + &go_catalog_json("gpt-5.4-mini", 20.0), + before, + root.path() + )); + let entries = [row("opencode-go", "gpt-5.4")]; + assert_eq!(cached_cost(&entries, root.path()), None); + + let targets = pricing_targets(&entries, now()); + let calls = Arc::new(AtomicUsize::new(0)); + let response = Some(go_catalog_json("gpt-5.4", 20.0)); + refresh_exact_pricing_targets_for_tests( + &targets, + system_now(), + root.path(), + response.clone(), + Arc::clone(&calls), + ) + .await; + assert_eq!(calls.load(Ordering::SeqCst), 1); + let cost = cached_cost(&entries, root.path()).expect("refreshed price"); + assert!((cost - 0.002_84).abs() < 1e-10, "{cost}"); + + refresh_exact_pricing_targets_for_tests( + &targets, + system_now(), + root.path(), + response, + Arc::clone(&calls), + ) + .await; + assert_eq!(calls.load(Ordering::SeqCst), 1); +} + +#[tokio::test] +async fn stale_price_refresh_changes_costs_while_failure_preserves_the_last_good_rates() { + let root = tempfile::tempdir().unwrap(); + let old = go_catalog_json("gpt-5.4", 5.0); + let stale = system_now() - StdDuration::from_secs(90_000); + assert!(save_catalog_json_for_tests(&old, stale, root.path())); + let entries = [row("opencode-go", "gpt-5.4")]; + // A stale catalog prices nothing until it is refreshed. + assert_eq!(cached_cost(&entries, root.path()), None); + + let targets = pricing_targets(&entries, now()); + let calls = Arc::new(AtomicUsize::new(0)); + refresh_exact_pricing_targets_for_tests( + &targets, + system_now(), + root.path(), + Some(go_catalog_json("gpt-5.4", 20.0)), + Arc::clone(&calls), + ) + .await; + assert_eq!(calls.load(Ordering::SeqCst), 1); + let old_cost = priced( + entries.to_vec(), + &CustomPricing::default(), + &snapshot_of(&old), + ) + .known_cost_usd; + let refreshed_cost = cached_cost(&entries, root.path()); + assert!(refreshed_cost.is_some()); + assert_ne!(old_cost, refreshed_cost); + + let cache_path = models_dev_cache_path_for_tests(root.path()); + let refreshed_bytes = std::fs::read(&cache_path).unwrap(); + let failures = Arc::new(AtomicUsize::new(0)); + refresh_exact_pricing_targets_for_tests( + &targets, + system_now() + StdDuration::from_secs(90_000), + root.path(), + None, + Arc::clone(&failures), + ) + .await; + assert_eq!(failures.load(Ordering::SeqCst), 1); + assert_eq!(std::fs::read(&cache_path).unwrap(), refreshed_bytes); +} diff --git a/rust/src/spend_contract/opencodex/tests.rs b/rust/src/spend_contract/opencodex/tests.rs new file mode 100644 index 0000000000..99cbc58b94 --- /dev/null +++ b/rust/src/spend_contract/opencodex/tests.rs @@ -0,0 +1,478 @@ +use super::super::CustomRates; +use super::cache::{load_entries_with_cache, read_cache}; +use super::*; +use std::fs; + +#[test] +fn aggregate_deduplicates_requests_and_applies_history_window() { + let now = DateTime::parse_from_rfc3339("2026-08-19T12:00:00Z") + .unwrap() + .with_timezone(&Utc); + let make = |request_id: &str, timestamp: &str, input: u64| OpenCodexEntry { + request_id: request_id.into(), + timestamp: DateTime::parse_from_rfc3339(timestamp) + .unwrap() + .with_timezone(&Utc), + provider: "openai".into(), + model: "gpt-5".into(), + usage_status: "reported".into(), + conversation_id: Some(request_id.into()), + input_tokens: Some(input), + output_tokens: Some(1), + cache_read_tokens: Some(0), + cache_creation_tokens: None, + reasoning_tokens: None, + total_tokens: Some(input + 1), + }; + let source = aggregate( + vec![ + make("same", "2026-08-18T10:00:00Z", 10), + make("same", "2026-08-18T11:00:00Z", 20), + make("old", "2026-08-01T10:00:00Z", 30), + ], + now, + 7, + &CustomPricing::default(), + ) + .expect("source"); + assert_eq!(source.request_count, 1); + assert_eq!(source.conversation_count, 1); + assert_eq!(source.token_mix.input_tokens, Some(20)); + assert_eq!(source.coverage.priced, 1); + assert!(source.known_cost_usd.is_some()); + assert_eq!(source.provenance, CostProvenance::VendorMetered); +} + +#[test] +fn aggregate_preserves_list_and_mixed_provenance() { + let now = DateTime::parse_from_rfc3339("2026-08-19T12:00:00Z") + .unwrap() + .with_timezone(&Utc); + let mut estimated = entry("openai", "gpt-5"); + estimated.request_id = "estimated".to_string(); + estimated.usage_status = "estimated".to_string(); + let list_only = aggregate(vec![estimated.clone()], now, 30, &CustomPricing::default()) + .expect("list-price source"); + assert_eq!(list_only.provenance, CostProvenance::ListPriceEstimate); + + let mut reported = entry("openai", "gpt-5"); + reported.request_id = "reported".to_string(); + let mixed = aggregate( + vec![reported, estimated], + now, + 30, + &CustomPricing::default(), + ) + .expect("mixed source"); + assert_eq!(mixed.provenance, CostProvenance::Mixed); +} + +#[test] +fn aggregate_preserves_zero_cost_authoritative_provenance() { + let now = DateTime::parse_from_rfc3339("2026-08-19T12:00:00Z") + .unwrap() + .with_timezone(&Utc); + let custom = CustomPricing { + entries: std::collections::HashMap::from([( + "openai/gpt-5".to_string(), + CustomRates { + input: Some(0.0), + output: Some(0.0), + cache_read: Some(0.0), + cache_write: Some(0.0), + }, + )]), + }; + + let reported = aggregate(vec![entry("openai", "gpt-5")], now, 30, &custom) + .expect("zero-cost vendor source"); + assert_eq!(reported.known_cost_usd, Some(0.0)); + assert_eq!(reported.provenance, CostProvenance::VendorMetered); + + let mut estimated_entry = entry("openai", "gpt-5"); + estimated_entry.usage_status = "estimated".to_string(); + let estimated = + aggregate(vec![estimated_entry], now, 30, &custom).expect("zero-cost list source"); + assert_eq!(estimated.known_cost_usd, Some(0.0)); + assert_eq!(estimated.provenance, CostProvenance::ListPriceEstimate); +} + +fn entry(provider: &str, model: &str) -> OpenCodexEntry { + OpenCodexEntry { + request_id: format!("{provider}:{model}"), + timestamp: DateTime::parse_from_rfc3339("2026-07-29T12:00:00Z") + .unwrap() + .with_timezone(&Utc), + provider: provider.to_string(), + model: model.to_string(), + usage_status: "reported".to_string(), + conversation_id: None, + input_tokens: Some(100), + output_tokens: Some(5), + cache_read_tokens: Some(10), + cache_creation_tokens: None, + reasoning_tokens: None, + total_tokens: Some(105), + } +} + +#[test] +fn routes_opencodex_entries_into_subscription_rows() { + assert_eq!( + route_entry(&entry("openai", "gpt-5.6-sol")), + RouteTarget::Subscription("codex") + ); + assert_eq!( + route_entry(&entry("opencode-go", "gpt-5.6-sol")), + RouteTarget::Subscription("opencodego") + ); + assert_eq!( + route_entry(&entry("kimi-coding", "k2p5")), + RouteTarget::Subscription("kimi") + ); + assert_eq!( + route_entry(&entry("deepseek", "deepseek-chat")), + RouteTarget::Subscription("deepseek") + ); + assert_eq!( + route_entry(&entry("opencode-free", "free-model")), + RouteTarget::TokenOnly + ); +} + +#[test] +fn recorded_provider_wins_over_mismatched_model_namespace() { + assert_eq!( + route_entry(&entry("opencode-go", "openai/gpt-5.6-sol")), + RouteTarget::Subscription("opencodego") + ); + assert_eq!( + route_entry(&entry("deepseek", "openai/gpt-5.6-sol")), + RouteTarget::Subscription("deepseek") + ); +} + +/// The models.dev identities that price a row, as (provider, model). +fn targets_of(provider: &str, model: &str) -> Vec<(String, String)> { + models_dev_pricing_targets(&pricing_provider(&entry(provider, model)), model) + .into_iter() + .map(|target| (target.provider_id, target.model_id)) + .collect() +} + +fn pair(provider: &str, model: &str) -> (String, String) { + (provider.to_string(), model.to_string()) +} + +#[test] +fn legacy_openai_transport_still_uses_explicit_route() { + assert_eq!( + route_entry(&entry("openai", "opencode-go/deepseek-v4-flash")), + RouteTarget::Subscription("opencodego") + ); + assert_eq!( + targets_of("openai", "opencode-go/gpt-5"), + vec![pair("opencode-go", "gpt-5")] + ); +} + +#[test] +fn pricing_targets_follow_the_recorded_provider() { + assert_eq!( + targets_of("opencode-go", "gpt-5"), + vec![pair("opencode-go", "gpt-5")] + ); + assert_eq!( + targets_of("kimi-coding", "k2p5"), + vec![pair("kimi-coding", "k2p5"), pair("kimi-for-coding", "k2p5")] + ); + assert_eq!( + targets_of("deepseek", "deepseek-chat"), + vec![pair("deepseek", "deepseek-chat")] + ); + // Another vendor's namespace is part of the model id on the recorded + // provider's catalog, never a route to that vendor's own rates. + assert_eq!( + targets_of("opencode-go", "openai/gpt-5"), + vec![pair("opencode-go", "openai/gpt-5")] + ); +} + +#[test] +fn unknown_provider_or_namespace_fails_closed_for_routing_and_pricing() { + let snapshot = ModelsDevPricingSnapshot::from_catalog_json_for_tests( + r#"{"openai":{"models":{"gpt-5":{"id":"gpt-5","cost":{"input":2,"output":8,"cache_read":0.2}}}}}"#, + ) + .expect("catalog"); + let none = CustomPricing::default(); + let proxy = entry("private-proxy", "openai/gpt-5"); + assert_eq!(route_entry(&proxy), RouteTarget::Unknown); + assert_eq!(entry_cost(&proxy, &none, &snapshot), None); + let malformed = entry("openai", "/gpt-5"); + assert_eq!(route_entry(&malformed), RouteTarget::Subscription("codex")); + assert!(targets_of("openai", "/gpt-5").is_empty()); + assert_eq!(entry_cost(&malformed, &none, &snapshot), None); +} + +#[test] +fn opencodex_uses_request_day_for_historical_gpt56_pricing() { + let entry = entry("openai", "gpt-5.6-terra"); + let empty = ModelsDevPricingSnapshot::from_catalog_json_for_tests("{}").expect("catalog"); + let cost = entry_cost(&entry, &CustomPricing::default(), &empty).unwrap(); + let expected = 90.0 * 2.5e-6 + 10.0 * 2.5e-7 + 5.0 * 1.5e-5; + assert!((cost - expected).abs() < 1e-12); +} + +#[test] +fn opencodex_historical_gpt56_bills_cache_writes_at_their_own_rate() { + let mut entry = entry("openai", "gpt-5.6-terra"); + entry.cache_creation_tokens = Some(20); + let empty = ModelsDevPricingSnapshot::from_catalog_json_for_tests("{}").expect("catalog"); + let cost = entry_cost(&entry, &CustomPricing::default(), &empty).unwrap(); + // Input includes cache reads and writes; each lane has its own rate. + let expected = 70.0 * 2.5e-6 + 10.0 * 2.5e-7 + 20.0 * 3.125e-6 + 5.0 * 1.5e-5; + assert!((cost - expected).abs() < 1e-12); +} + +#[test] +fn parser_keeps_reported_token_classes() { + let value = serde_json::json!({ + "requestId": "r1", "timestamp": "2026-08-18T10:00:00Z", "provider": "openai", + "model": "gpt-test", "usageStatus": "reported", "conversationId": "c1", + "usage": {"inputTokens": 10, "outputTokens": 4, "cachedInputTokens": 3, "reasoningOutputTokens": 2} + }); + let entry = parse_line(&value.to_string()).expect("entry"); + assert_eq!(entry.model, "gpt-test"); + assert_eq!(entry.input_tokens, Some(10)); + assert_eq!(entry.output_tokens, Some(4)); + assert_eq!(entry.cache_read_tokens, Some(3)); + assert_eq!(entry.reasoning_tokens, Some(2)); +} + +#[test] +fn parser_normalizes_defaults_and_rejects_malformed_lines() { + let minimal = serde_json::json!({ + "requestId": " r1 ", "model": "gpt-test", "timestamp": "2026-08-18T10:00:00Z", + "usageStatus": " REPORTED ", "usage": {"cacheCreationInputTokens": 7} + }); + let entry = parse_line(&minimal.to_string()).expect("entry"); + assert_eq!(entry.request_id, "r1", "ids are trimmed"); + assert_eq!( + entry.provider, "openai", + "missing provider defaults to openai" + ); + assert_eq!( + entry.usage_status, "reported", + "status is lowercased and trimmed" + ); + assert_eq!(entry.conversation_id, None); + assert_eq!(entry.cache_creation_tokens, Some(7)); + + for malformed in [ + "{}", + r#"{"requestId": "", "model": "m", "timestamp": "2026-08-18T10:00:00Z"}"#, + r#"{"requestId": "r1", "model": " ", "timestamp": "2026-08-18T10:00:00Z"}"#, + r#"{"requestId": "r1", "model": "m"}"#, + "not json at all", + ] { + assert!(parse_line(malformed).is_none(), "rejected: {malformed}"); + } +} + +#[test] +fn incremental_cache_appends_only_newline_terminated_tail() { + let dir = tempfile::tempdir().unwrap(); + let log = dir.path().join("usage.jsonl"); + let cache = dir.path().join("cache.sqlite"); + let row = |id: &str, input: u64| { + format!( + r#"{{"requestId":"{id}","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z","usageStatus":"reported","usage":{{"inputTokens":{input}}}}}"# + ) + }; + + fs::write(&log, format!("{}\n{}\n", row("a", 1), row("b", 2))).unwrap(); + let first = load_entries_with_cache(&log, &cache).unwrap(); + assert_eq!( + first + .iter() + .map(|entry| entry.request_id.as_str()) + .collect::>(), + vec!["a", "b"] + ); + let first_cursor = read_cache(&cache).unwrap().cursor; + + let mut file = fs::OpenOptions::new().append(true).open(&log).unwrap(); + use std::io::Write as _; + writeln!(file, "{}", row("c", 3)).unwrap(); + drop(file); + + let second = load_entries_with_cache(&log, &cache).unwrap(); + assert_eq!( + second + .iter() + .map(|entry| entry.request_id.as_str()) + .collect::>(), + vec!["a", "b", "c"] + ); + let second_cursor = read_cache(&cache).unwrap().cursor; + assert!(second_cursor.parsed_offset > first_cursor.parsed_offset); +} + +#[test] +fn incomplete_trailing_opencodex_record_waits_for_newline() { + let dir = tempfile::tempdir().unwrap(); + let log = dir.path().join("usage.jsonl"); + let cache = dir.path().join("cache.sqlite"); + let complete = r#"{"requestId":"a","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; + let pending = r#"{"requestId":"b","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; + let split = pending.len() / 2; + fs::write(&log, format!("{complete}\n{}", &pending[..split])).unwrap(); + + let first = load_entries_with_cache(&log, &cache).unwrap(); + assert_eq!( + first + .iter() + .map(|entry| entry.request_id.as_str()) + .collect::>(), + vec!["a"] + ); + let cursor = read_cache(&cache).unwrap().cursor; + assert_eq!( + cursor.parsed_offset, + u64::try_from(complete.len() + 1).unwrap() + ); + + let mut file = fs::OpenOptions::new().append(true).open(&log).unwrap(); + use std::io::Write as _; + writeln!(file, "{}", &pending[split..]).unwrap(); + drop(file); + let second = load_entries_with_cache(&log, &cache).unwrap(); + assert_eq!( + second + .iter() + .map(|entry| entry.request_id.as_str()) + .collect::>(), + vec!["a", "b"] + ); +} + +#[test] +fn complete_trailing_opencodex_record_waits_for_newline() { + let dir = tempfile::tempdir().unwrap(); + let log = dir.path().join("usage.jsonl"); + let cache = dir.path().join("cache.sqlite"); + let first = r#"{"requestId":"a","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; + let trailing = r#"{"requestId":"b","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; + fs::write(&log, format!("{first}\n{trailing}")).unwrap(); + + let before_newline = load_entries_with_cache(&log, &cache).unwrap(); + assert_eq!( + before_newline + .iter() + .map(|entry| entry.request_id.as_str()) + .collect::>(), + vec!["a"] + ); + + let mut file = fs::OpenOptions::new().append(true).open(&log).unwrap(); + use std::io::Write as _; + writeln!(file).unwrap(); + drop(file); + + let after_newline = load_entries_with_cache(&log, &cache).unwrap(); + assert_eq!( + after_newline + .iter() + .map(|entry| entry.request_id.as_str()) + .collect::>(), + vec!["a", "b"] + ); +} + +#[test] +fn later_request_id_replaces_cached_entry_without_full_cache_loss() { + let dir = tempfile::tempdir().unwrap(); + let log = dir.path().join("usage.jsonl"); + let cache = dir.path().join("cache.sqlite"); + let row = |id: &str, input: u64| { + format!( + r#"{{"requestId":"{id}","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z","usageStatus":"reported","usage":{{"inputTokens":{input}}}}}"# + ) + }; + fs::write(&log, format!("{}\n{}\n", row("dup", 1), row("keep", 2))).unwrap(); + let _ = load_entries_with_cache(&log, &cache).unwrap(); + let mut file = fs::OpenOptions::new().append(true).open(&log).unwrap(); + use std::io::Write as _; + writeln!(file, "{}", row("dup", 9)).unwrap(); + drop(file); + + let entries = load_entries_with_cache(&log, &cache).unwrap(); + assert_eq!(entries.len(), 2); + assert_eq!( + entries + .iter() + .find(|entry| entry.request_id == "dup") + .unwrap() + .input_tokens, + Some(9) + ); + assert!(entries.iter().any(|entry| entry.request_id == "keep")); +} + +#[test] +fn truncation_invalidates_opencodex_cursor_and_rebuilds() { + let dir = tempfile::tempdir().unwrap(); + let log = dir.path().join("usage.jsonl"); + let cache = dir.path().join("cache.sqlite"); + let old = r#"{"requestId":"old","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; + let replacement = r#"{"requestId":"new","model":"gpt-5","timestamp":"2026-08-18T10:00:00Z"}"#; + fs::write(&log, format!("{old}\n{old}\n")).unwrap(); + let _ = load_entries_with_cache(&log, &cache).unwrap(); + fs::write(&log, format!("{replacement}\n")).unwrap(); + + let rebuilt = load_entries_with_cache(&log, &cache).unwrap(); + assert_eq!( + rebuilt + .iter() + .map(|entry| entry.request_id.as_str()) + .collect::>(), + vec!["new"] + ); +} + +#[test] +fn timestamps_parse_rfc3339_epoch_seconds_and_millis() { + let expected = Utc.with_ymd_and_hms(2026, 8, 18, 10, 0, 0).unwrap(); + let rfc3339 = serde_json::json!("2026-08-18T10:00:00Z"); + assert_eq!(parse_timestamp(&rfc3339), Some(expected)); + let epoch_seconds = serde_json::json!(1_787_047_200i64); + assert_eq!(parse_timestamp(&epoch_seconds), Some(expected)); + let epoch_millis = serde_json::json!(1_787_047_200_000f64); + assert_eq!(parse_timestamp(&epoch_millis), Some(expected)); + let numeric_string = serde_json::json!("1787047200.0"); + assert_eq!(parse_timestamp(&numeric_string), Some(expected)); + for invalid in [ + serde_json::json!("not a date"), + serde_json::json!(0), + serde_json::json!(-5.0), + serde_json::Value::Null, + serde_json::json!(true), + ] { + assert!(parse_timestamp(&invalid).is_none(), "rejected: {invalid}"); + } +} + +#[test] +fn nonnegative_u64_accepts_json_numbers_and_bounded_floats() { + assert_eq!(nonnegative_u64(Some(&serde_json::json!(42))), Some(42)); + assert_eq!(nonnegative_u64(Some(&serde_json::json!(12.0))), Some(12)); + // Fractional floats are accepted via `as u64` truncation. + assert_eq!(nonnegative_u64(Some(&serde_json::json!(1.5))), Some(1)); + // `u64::MAX as f64` rounds up to 2^64; f64 spacing there is 4096, so + // +2048.0 rounds back into range. First out-of-range step is +4096.0. + assert_eq!( + nonnegative_u64(Some(&serde_json::json!(u64::MAX as f64 + 4096.0))), + None + ); + assert_eq!(nonnegative_u64(None), None); +} diff --git a/rust/src/spend_contract/tests.rs b/rust/src/spend_contract/tests.rs index fb74ff09de..0d89c428af 100644 --- a/rust/src/spend_contract/tests.rs +++ b/rust/src/spend_contract/tests.rs @@ -54,6 +54,42 @@ fn explicit_zero_custom_rate_is_known_free_but_missing_rate_is_unknown() { assert_eq!(missing.cost(&counts), None); } +#[test] +fn custom_pricing_reads_each_entry_on_its_own_like_upstream() { + let custom = CustomPricing::parse( + br#"{ + " GPT-5 ": {"input": 1.25, "output": 10, "cacheRead": 0.125, "cache_read": 9}, + "negative-input": {"input": -1, "output": 2}, + "string-rate": {"input": "3", "output": 4}, + "snake": {"cache_write": 0, "cacheCreation": 7}, + "empty": {}, + "only-unusable": {"input": -1, "output": "free"}, + "not-an-object": 5, + " ": {"input": 1} + }"#, + ); + let rates = |key: &str| { + let rates = custom.entries.get(key).expect(key); + ( + rates.input, + rates.output, + rates.cache_read, + rates.cache_write, + ) + }; + assert_eq!(rates("gpt-5"), (Some(1.25), Some(10.0), Some(0.125), None)); + assert_eq!(rates("negative-input"), (None, Some(2.0), None, None)); + assert_eq!(rates("string-rate"), (None, Some(4.0), None, None)); + assert_eq!(rates("snake"), (None, None, None, Some(0.0))); + assert_eq!( + custom.entries.len(), + 4, + "entries without a usable rate are dropped" + ); + assert!(CustomPricing::parse(b"[1, 2]").entries.is_empty()); + assert!(CustomPricing::parse(b"not json").entries.is_empty()); +} + #[test] fn local_spend_contract_exposes_reasoning_tokens_and_preserves_unknown() { let known_summary = CostSummary {