diff --git a/.github/workflows/dev-build.yml b/.github/workflows/dev-build.yml index 4d2a172..6291cfc 100644 --- a/.github/workflows/dev-build.yml +++ b/.github/workflows/dev-build.yml @@ -69,6 +69,7 @@ jobs: mkdir -p "$artifact_dir" cp "target/$BUILD_TARGET/release/statsai" "$artifact_dir/statsai" store_schema_version=$("$artifact_dir/statsai" store supported-schema-version) + pricing_ruleset_version=$("$artifact_dir/statsai" store supported-pricing-ruleset-version) if [[ "${{ github.event_name }}" == "pull_request" ]]; then source_json=$(jq -n --argjson number "$PR_NUMBER" '{kind: "pull_request", number: $number}') @@ -84,21 +85,24 @@ jobs: --argjson workflow_run_id "$WORKFLOW_RUN_ID" \ --argjson workflow_attempt "$WORKFLOW_ATTEMPT" \ --argjson store_schema_version "$store_schema_version" \ + --argjson pricing_ruleset_version "$pricing_ruleset_version" \ '{ - schema: 1, + schema: 2, repository: $repository, sha: $sha, target: $target, source: $source, workflow_run_id: $workflow_run_id, workflow_attempt: $workflow_attempt, - store_schema_version: $store_schema_version + store_schema_version: $store_schema_version, + pricing_ruleset_version: $pricing_ruleset_version }' > "$artifact_dir/build.json" (cd "$artifact_dir" && shasum -a 256 statsai > SHA256SUMS) test "$(find "$artifact_dir" -type f | wc -l | tr -d ' ')" = 3 test "$(jq -r .sha "$artifact_dir/build.json")" = "$BUILD_SHA" test "$(jq -r .store_schema_version "$artifact_dir/build.json")" = "$store_schema_version" + test "$(jq -r .pricing_ruleset_version "$artifact_dir/build.json")" = "$pricing_ruleset_version" (cd "$artifact_dir" && shasum -a 256 -c SHA256SUMS) echo "sha=$BUILD_SHA" >> "$GITHUB_OUTPUT" echo "path=$artifact_dir" >> "$GITHUB_OUTPUT" diff --git a/Cargo.lock b/Cargo.lock index ffdd553..2e2f087 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2556,6 +2556,7 @@ dependencies = [ "serde_json", "sha2", "statsai-core", + "statsai-pricing", "tempfile", ] diff --git a/README.md b/README.md index 809b920..4ddf3ae 100644 --- a/README.md +++ b/README.md @@ -436,6 +436,17 @@ Cost figures are API-equivalent estimates, not subscription invoices. StatsAI selects published pricing by usage timestamp, preserves integer micro-USD values, and leaves cost unknown when a source does not prove the billable model. +Pricing updates are version-driven and automatic. The selected `statsai` binary +reprices persisted normalized events when its compiled ruleset is newer than +the store's last applied ruleset, before scan, report, sync, snapshot, task, +export, import, or daemon use price-derived data. Diagnostic commands such as +`status`, `doctor`, `quota`, `conversation`, and `account` do not trigger that +pass. A raw provider rescan (`--no-cache`) is not required. Ordinary incremental +`statsai sync --sink http` then uploads the corrected dirty rollups; +`--rebuild-rollups` and `--full` are unnecessary just because prices changed. A +store already processed by a newer ruleset is refused rather than silently +repriced backward. + ## Develop locally This is a Rust workspace. Run the same checks used in CI: diff --git a/crates/statsai-dev/src/artifact.rs b/crates/statsai-dev/src/artifact.rs index 81e1ac2..9455694 100644 --- a/crates/statsai-dev/src/artifact.rs +++ b/crates/statsai-dev/src/artifact.rs @@ -10,7 +10,7 @@ use std::path::{Path, PathBuf}; use zip::ZipArchive; pub(crate) const TARGET: &str = "aarch64-apple-darwin"; -const MANIFEST_SCHEMA: u32 = 1; +const MANIFEST_SCHEMA: u32 = 2; const MAX_BINARY_BYTES: u64 = 192 * 1024 * 1024; const MAX_METADATA_BYTES: u64 = 64 * 1024; @@ -24,6 +24,7 @@ pub(crate) struct BuildManifest { pub(crate) workflow_run_id: u64, pub(crate) workflow_attempt: u64, pub(crate) store_schema_version: i64, + pub(crate) pricing_ruleset_version: u64, } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] @@ -103,8 +104,7 @@ pub(crate) fn verify_download( let binary = binary.context("artifact is missing statsai")?; let manifest_bytes = manifest_bytes.context("artifact is missing build.json")?; let checksums = checksums.context("artifact is missing SHA256SUMS")?; - let manifest: BuildManifest = - serde_json::from_slice(&manifest_bytes).context("parse artifact build.json")?; + let manifest = parse_manifest(&manifest_bytes, "parse artifact build.json")?; validate_manifest(&manifest, expected_sha, Some(run))?; validate_macho_arm64(&binary)?; let expected_checksum = parse_statsai_checksum(&checksums)?; @@ -171,10 +171,12 @@ pub(crate) fn load_cached(paths: &Paths, sha: &str) -> Result Result<()> { .with_context(|| format!("remove cached build {}", directory.display())) } +fn parse_manifest(bytes: &[u8], context: &str) -> Result { + let value: serde_json::Value = + serde_json::from_slice(bytes).with_context(|| context.to_string())?; + match value.get("schema").and_then(serde_json::Value::as_u64) { + Some(schema) if schema == u64::from(MANIFEST_SCHEMA) => {} + Some(schema) => { + bail!("unsupported build manifest schema {schema} (expected {MANIFEST_SCHEMA})") + } + None => bail!("build manifest is missing schema (expected {MANIFEST_SCHEMA})"), + } + serde_json::from_value(value).with_context(|| context.to_string()) +} + fn validate_manifest( manifest: &BuildManifest, expected_sha: &str, @@ -403,7 +418,7 @@ mod tests { fn manifest(sha: &str) -> BuildManifest { BuildManifest { - schema: 1, + schema: MANIFEST_SCHEMA, repository: REPOSITORY.to_string(), sha: sha.to_string(), target: TARGET.to_string(), @@ -411,9 +426,30 @@ mod tests { workflow_run_id: 104, workflow_attempt: 1, store_schema_version: 17, + pricing_ruleset_version: 1, } } + fn archive_raw_manifest(manifest: &[u8], binary: &[u8], checksum: &str) -> Vec { + let mut writer = ZipWriter::new(Cursor::new(Vec::new())); + let options = SimpleFileOptions::default(); + writer + .start_file("statsai", options) + .expect("start binary entry"); + writer.write_all(binary).expect("write binary entry"); + writer + .start_file("build.json", options) + .expect("start manifest entry"); + writer.write_all(manifest).expect("write manifest entry"); + writer + .start_file("SHA256SUMS", options) + .expect("start checksum entry"); + writeln!(writer, "{checksum} statsai").expect("write checksum entry"); + let mut cursor = writer.finish().expect("finish artifact ZIP"); + cursor.rewind().expect("rewind artifact ZIP"); + cursor.into_inner() + } + fn run(sha: &str) -> WorkflowRun { WorkflowRun { id: 104, @@ -459,9 +495,75 @@ mod tests { .expect("verify artifact"); assert_eq!(verified.manifest.sha, sha); + assert_eq!(verified.manifest.schema, MANIFEST_SCHEMA); + assert_eq!(verified.manifest.pricing_ruleset_version, 1); assert_eq!(verified.binary_sha256, checksum); } + #[test] + fn shipped_schema_1_manifest_is_rejected() { + let sha = "a".repeat(40); + let binary = fake_binary(); + let checksum = sha256_bytes(&binary); + let schema_1 = serde_json::json!({ + "schema": 1, + "repository": REPOSITORY, + "sha": sha, + "target": TARGET, + "source": { "kind": "pull_request", "number": 12 }, + "workflow_run_id": 104, + "workflow_attempt": 1, + "store_schema_version": 17 + }); + let error = verify_download( + &archive_raw_manifest( + &serde_json::to_vec(&schema_1).expect("serialize schema 1"), + &binary, + &checksum, + ), + &sha, + &run(&sha), + ) + .expect_err("shipped schema 1 must fail"); + + assert!( + error + .to_string() + .contains("unsupported build manifest schema 1 (expected 2)"), + "{error}" + ); + } + + #[test] + fn schema_2_manifest_requires_pricing_ruleset_version() { + let sha = "a".repeat(40); + let binary = fake_binary(); + let checksum = sha256_bytes(&binary); + let missing_pricing = serde_json::json!({ + "schema": 2, + "repository": REPOSITORY, + "sha": sha, + "target": TARGET, + "source": { "kind": "main" }, + "workflow_run_id": 104, + "workflow_attempt": 1, + "store_schema_version": 17 + }); + let error = verify_download( + &archive_raw_manifest( + &serde_json::to_vec(&missing_pricing).expect("serialize incomplete schema 2"), + &binary, + &checksum, + ), + &sha, + &run(&sha), + ) + .expect_err("schema 2 without pricing_ruleset_version must fail"); + + let message = format!("{error:#}"); + assert!(message.contains("pricing_ruleset_version"), "{message}"); + } + #[test] fn checksum_mismatch_is_rejected() { let sha = "a".repeat(40); diff --git a/crates/statsai-dev/src/launcher.rs b/crates/statsai-dev/src/launcher.rs index 37e1f24..48971f1 100644 --- a/crates/statsai-dev/src/launcher.rs +++ b/crates/statsai-dev/src/launcher.rs @@ -1,7 +1,7 @@ use crate::artifact::load_cached; use crate::state::{Environment, Paths, State}; use anyhow::{bail, Context, Result}; -use statsai_store::database_schema_version; +use statsai_store::{database_applied_pricing_ruleset_version, database_schema_version}; use std::ffi::{OsStr, OsString}; use std::process::{Command, ExitCode}; @@ -62,11 +62,17 @@ pub(crate) fn forward( production_schema, installed.manifest.store_schema_version, )?; + let production_pricing = database_applied_pricing_ruleset_version(store)?; + ensure_production_pricing_compatible( + production_pricing, + installed.manifest.pricing_ruleset_version, + )?; eprintln!( - "WARNING: running development build {} against the production database {} (schema {})", + "WARNING: running development build {} against the production database {} (schema {}, pricing ruleset {})", &selected.sha[..8], paths.display(&paths.prod_store), - production_schema + production_schema, + installed.manifest.pricing_ruleset_version ); } @@ -101,6 +107,21 @@ fn ensure_production_schema_compatible(production_schema: i64, build_schema: i64 Ok(()) } +fn ensure_production_pricing_compatible( + production_ruleset: Option, + build_ruleset: u64, +) -> Result<()> { + match production_ruleset { + Some(applied) if applied == build_ruleset => Ok(()), + Some(applied) => bail!( + "refusing `--prod-data`: selected build supports pricing ruleset {build_ruleset}, but production uses pricing ruleset {applied}; development builds may not reprice an incompatible production database" + ), + None => bail!( + "refusing `--prod-data`: production database has no applied pricing ruleset, but the selected build requires pricing ruleset {build_ruleset}; development builds may not reprice an incompatible production database" + ), + } +} + fn apply_environment(command: &mut Command, environment: Environment) { match environment { Environment::Local | Environment::Dev => { @@ -182,11 +203,58 @@ mod tests { use crate::github::REPOSITORY; use crate::state::{BuildSource, SelectedBuild}; use sha2::{Digest, Sha256}; + use statsai_store::{ + APPLIED_PRICING_RULESET_VERSION_KEY, CURRENT_SCHEMA_VERSION, PRICING_RULESET_VERSION, + }; fn arguments(values: &[&str]) -> Vec { values.iter().map(OsString::from).collect() } + fn install_selected_build( + paths: &Paths, + store_schema_version: i64, + pricing_ruleset_version: u64, + ) -> (State, artifact::InstalledBuild) { + let sha = "a".repeat(40); + let mut binary = vec![0u8; 64]; + binary[..4].copy_from_slice(&[0xcf, 0xfa, 0xed, 0xfe]); + binary[4..8].copy_from_slice(&[0x0c, 0x00, 0x00, 0x01]); + let binary_sha256 = hex::encode(Sha256::digest(&binary)); + let installed = artifact::install( + paths, + &VerifiedArtifact { + manifest: BuildManifest { + schema: 2, + repository: REPOSITORY.to_string(), + sha: sha.clone(), + target: TARGET.to_string(), + source: ManifestSource::PullRequest { number: 12 }, + workflow_run_id: 104, + workflow_attempt: 1, + store_schema_version, + pricing_ruleset_version, + }, + binary, + binary_sha256: binary_sha256.clone(), + }, + ) + .expect("install fixture"); + let state = State { + build: Some(SelectedBuild { + source: BuildSource::Pr, + pr: Some(12), + sha, + workflow_run_id: installed.manifest.workflow_run_id, + workflow_attempt: installed.manifest.workflow_attempt, + target: installed.manifest.target.clone(), + binary_sha256, + }), + ..State::default() + }; + (state, installed) + } + #[test] fn daemon_and_mutating_service_commands_are_blocked() { assert!(validate_forward_arguments(&arguments(&["daemon", "--watch"])).is_err()); @@ -233,6 +301,14 @@ mod tests { assert!(ensure_production_schema_compatible(18, 17).is_err()); } + #[test] + fn production_data_requires_an_exact_pricing_match() { + assert!(ensure_production_pricing_compatible(Some(1), 1).is_ok()); + assert!(ensure_production_pricing_compatible(None, 1).is_err()); + assert!(ensure_production_pricing_compatible(Some(1), 2).is_err()); + assert!(ensure_production_pricing_compatible(Some(2), 1).is_err()); + } + #[test] fn production_data_refuses_a_build_that_would_migrate_the_store() { let directory = tempfile::tempdir().expect("temporary directory"); @@ -241,41 +317,7 @@ mod tests { .expect("create production store at current schema"); drop(production); - let sha = "a".repeat(40); - let mut binary = vec![0u8; 64]; - binary[..4].copy_from_slice(&[0xcf, 0xfa, 0xed, 0xfe]); - binary[4..8].copy_from_slice(&[0x0c, 0x00, 0x00, 0x01]); - let binary_sha256 = hex::encode(Sha256::digest(&binary)); - let installed = artifact::install( - &paths, - &VerifiedArtifact { - manifest: BuildManifest { - schema: 1, - repository: REPOSITORY.to_string(), - sha: sha.clone(), - target: TARGET.to_string(), - source: ManifestSource::PullRequest { number: 12 }, - workflow_run_id: 104, - workflow_attempt: 1, - store_schema_version: statsai_store::CURRENT_SCHEMA_VERSION + 1, - }, - binary, - binary_sha256: binary_sha256.clone(), - }, - ) - .expect("install future-schema fixture"); - let state = State { - build: Some(SelectedBuild { - source: BuildSource::Pr, - pr: Some(12), - sha, - workflow_run_id: installed.manifest.workflow_run_id, - workflow_attempt: installed.manifest.workflow_attempt, - target: installed.manifest.target.clone(), - binary_sha256, - }), - ..State::default() - }; + let (state, _) = install_selected_build(&paths, CURRENT_SCHEMA_VERSION + 1, 1); let error = forward(&paths, &state, &arguments(&["status"]), true) .expect_err("future-schema build must not run against production"); @@ -283,7 +325,63 @@ mod tests { assert!(error.to_string().contains("may not migrate")); assert_eq!( database_schema_version(&paths.prod_store).expect("read unchanged production schema"), - Some(statsai_store::CURRENT_SCHEMA_VERSION) + Some(CURRENT_SCHEMA_VERSION) + ); + } + + #[test] + fn production_data_refuses_a_missing_applied_pricing_ruleset() { + let directory = tempfile::tempdir().expect("temporary directory"); + let paths = Paths::for_test(directory.path()); + drop( + statsai_store::Store::open(&paths.prod_store) + .expect("create production store at current schema"), + ); + let (state, _) = + install_selected_build(&paths, CURRENT_SCHEMA_VERSION, PRICING_RULESET_VERSION); + + let error = forward(&paths, &state, &arguments(&["status"]), true) + .expect_err("unpriced production store must not run against development builds"); + + assert!(error.to_string().contains("no applied pricing ruleset")); + assert_eq!( + database_applied_pricing_ruleset_version(&paths.prod_store) + .expect("read unchanged production pricing"), + None + ); + } + + #[test] + fn production_data_refuses_older_and_newer_pricing_rulesets() { + let directory = tempfile::tempdir().expect("temporary directory"); + let paths = Paths::for_test(directory.path()); + let production = statsai_store::Store::open(&paths.prod_store) + .expect("create production store at current schema"); + production + .set_metadata_value(APPLIED_PRICING_RULESET_VERSION_KEY, "1") + .expect("record older applied ruleset"); + drop(production); + let (older_state, _) = install_selected_build(&paths, CURRENT_SCHEMA_VERSION, 2); + let older = forward(&paths, &older_state, &arguments(&["status"]), true) + .expect_err("older production ruleset must refuse"); + assert!(older.to_string().contains("pricing ruleset 2")); + assert!(older.to_string().contains("pricing ruleset 1")); + + let production = + statsai_store::Store::open(&paths.prod_store).expect("reopen production store"); + production + .set_metadata_value(APPLIED_PRICING_RULESET_VERSION_KEY, "99") + .expect("record newer applied ruleset"); + drop(production); + let (newer_state, _) = + install_selected_build(&paths, CURRENT_SCHEMA_VERSION, PRICING_RULESET_VERSION); + let newer = forward(&paths, &newer_state, &arguments(&["status"]), true) + .expect_err("newer production ruleset must refuse"); + assert!(newer.to_string().contains("pricing ruleset 99")); + assert_eq!( + database_applied_pricing_ruleset_version(&paths.prod_store) + .expect("read unchanged production pricing"), + Some(99) ); } } diff --git a/crates/statsai-dev/src/lib.rs b/crates/statsai-dev/src/lib.rs index 0ab2c53..14533cf 100644 --- a/crates/statsai-dev/src/lib.rs +++ b/crates/statsai-dev/src/lib.rs @@ -617,7 +617,7 @@ mod tests { fn installed_build(paths: &Paths, sha: &str) -> InstalledBuild { InstalledBuild { manifest: BuildManifest { - schema: 1, + schema: 2, repository: github::REPOSITORY.to_string(), sha: sha.to_string(), target: TARGET.to_string(), @@ -625,6 +625,7 @@ mod tests { workflow_run_id: 104, workflow_attempt: 1, store_schema_version: 17, + pricing_ruleset_version: 1, }, binary_sha256: "b".repeat(64), binary_path: paths.builds_dir.join(sha).join("statsai"), diff --git a/crates/statsai-pricing/src/lib.rs b/crates/statsai-pricing/src/lib.rs index 8bbc303..082f8b0 100644 --- a/crates/statsai-pricing/src/lib.rs +++ b/crates/statsai-pricing/src/lib.rs @@ -2,11 +2,25 @@ //! //! Provides static model pricing lookup and cost estimation //! decoupled from any specific adapter. +//! +//! [`PRICING_RULESET_VERSION`] is the monotonic numeric identity of the compiled +//! pricing rules. Increment it on every semantic pricing-rule change: new +//! models, rate changes, date-boundary mappings, multiplier logic, or anything +//! else that can change an estimated cost. Do not order +//! [`PRICING_CATALOG_VERSION`] strings lexicographically. use chrono::{DateTime, Datelike, Utc}; use statsai_core::{micro_usd_to_cents_rounded, Confidence, CostInfo, ModelInfo, UsageCounts}; -const PRICING_CATALOG_VERSION: &str = "official:2026-08-19"; +/// Monotonic numeric identity of the compiled pricing ruleset. +/// +/// Increment this constant on every semantic pricing-rule change so persisted +/// stores can reprice automatically. The descriptive catalog string is not an +/// ordering key. +pub const PRICING_RULESET_VERSION: u64 = 1; + +/// Descriptive identifier for the compiled price list. Not an ordering key. +pub const PRICING_CATALOG_VERSION: &str = "official:2026-08-19"; const MICRO_USD_PER_USD: i128 = 1_000_000; const TOKENS_PER_MILLION: i128 = 1_000_000; const MULTIPLIER_SCALE: i128 = 10_000; @@ -622,6 +636,40 @@ pub fn unknown_cost() -> CostInfo { } } +/// Overlays a freshly estimated cost onto a persisted [`CostInfo`]. +/// +/// Provider-reported amounts are never replaced. When a provider-reported +/// value is present, its provenance and confidence are preserved and only the +/// estimated fields and catalog version are updated. +#[must_use] +pub fn overlay_estimated_cost(existing: &CostInfo, estimated: CostInfo) -> CostInfo { + let has_provider_reported = + existing.provider_reported_usd.is_some() || existing.provider_reported_micro_usd.is_some(); + if has_provider_reported { + CostInfo { + currency: existing.currency.clone(), + estimated_api_equivalent_usd: estimated.estimated_api_equivalent_usd, + provider_reported_usd: existing.provider_reported_usd, + estimated_api_equivalent_micro_usd: estimated.estimated_api_equivalent_micro_usd, + provider_reported_micro_usd: existing.provider_reported_micro_usd, + pricing_source: existing.pricing_source.clone(), + pricing_version: estimated.pricing_version, + confidence: existing.confidence.clone(), + } + } else { + CostInfo { + currency: estimated.currency, + estimated_api_equivalent_usd: estimated.estimated_api_equivalent_usd, + provider_reported_usd: existing.provider_reported_usd, + estimated_api_equivalent_micro_usd: estimated.estimated_api_equivalent_micro_usd, + provider_reported_micro_usd: existing.provider_reported_micro_usd, + pricing_source: estimated.pricing_source, + pricing_version: estimated.pricing_version, + confidence: estimated.confidence, + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -1642,6 +1690,76 @@ mod tests { assert_eq!(terra_cut.estimated_api_equivalent_usd, Some(1_670)); } + #[test] + fn overlay_preserves_provider_reported_provenance() { + let existing = CostInfo { + currency: "USD".to_string(), + estimated_api_equivalent_usd: Some(10), + provider_reported_usd: Some(42), + estimated_api_equivalent_micro_usd: Some(100_000), + provider_reported_micro_usd: Some(420_000), + pricing_source: Some("claude_stats_cache:costUSD".to_string()), + pricing_version: Some("legacy".to_string()), + confidence: Confidence::High, + }; + let estimated = estimate_cost( + "codex", + Some(&test_model("gpt-5")), + &UsageCounts { + input_tokens: Some(1_000_000), + output_tokens: Some(500_000), + ..UsageCounts::default() + }, + ); + + let overlaid = overlay_estimated_cost(&existing, estimated.clone()); + + assert_eq!(overlaid.provider_reported_usd, Some(42)); + assert_eq!(overlaid.provider_reported_micro_usd, Some(420_000)); + assert_eq!( + overlaid.estimated_api_equivalent_usd, + estimated.estimated_api_equivalent_usd + ); + assert_eq!( + overlaid.estimated_api_equivalent_micro_usd, + estimated.estimated_api_equivalent_micro_usd + ); + assert_eq!( + overlaid.pricing_source.as_deref(), + Some("claude_stats_cache:costUSD") + ); + assert_eq!( + overlaid.pricing_version.as_deref(), + Some(PRICING_CATALOG_VERSION) + ); + assert_eq!(overlaid.confidence, Confidence::High); + } + + #[test] + fn overlay_replaces_estimated_only_cost() { + let existing = unknown_cost(); + let estimated = estimate_cost( + "codex", + Some(&test_model("gpt-5")), + &UsageCounts { + input_tokens: Some(1_000_000), + ..UsageCounts::default() + }, + ); + + let overlaid = overlay_estimated_cost(&existing, estimated.clone()); + + assert_eq!(overlaid, estimated); + assert!(overlaid.provider_reported_usd.is_none()); + } + + #[test] + fn ruleset_version_is_numeric_and_catalog_is_descriptive() { + const { assert!(PRICING_RULESET_VERSION >= 1) }; + assert!(PRICING_CATALOG_VERSION.contains(':')); + assert_ne!(PRICING_RULESET_VERSION.to_string(), PRICING_CATALOG_VERSION); + } + #[test] fn gpt_5_6_luna_and_terra_report_aggregate_periods_that_cross_the_july_30_cut() { let before = chrono::NaiveDate::from_ymd_opt(2026, 7, 29).expect("before boundary"); diff --git a/crates/statsai-store/Cargo.toml b/crates/statsai-store/Cargo.toml index 712573c..80f6c92 100644 --- a/crates/statsai-store/Cargo.toml +++ b/crates/statsai-store/Cargo.toml @@ -9,6 +9,7 @@ authors.workspace = true [dependencies] statsai-core.workspace = true +statsai-pricing.workspace = true anyhow.workspace = true base64.workspace = true chrono.workspace = true diff --git a/crates/statsai-store/src/lib.rs b/crates/statsai-store/src/lib.rs index 1a9e7bd..97eccea 100644 --- a/crates/statsai-store/src/lib.rs +++ b/crates/statsai-store/src/lib.rs @@ -4,6 +4,7 @@ mod account_plan; mod archive; mod code_changes; mod migrations; +mod pricing; mod privacy; mod quota; mod snapshot; @@ -36,7 +37,15 @@ use std::path::Path; pub use account_plan::AccountEvidenceReferenceCounts; pub use migrations::CURRENT_SCHEMA_VERSION; -pub use snapshot::{clone_database_to, database_schema_version, DatabaseClone}; +pub use pricing::{ + apply_current_estimated_pricing, RepricingReport, APPLIED_PRICING_CATALOG_VERSION_KEY, + APPLIED_PRICING_RULESET_VERSION_KEY, +}; +pub use snapshot::{ + clone_database_to, database_applied_pricing_ruleset_version, database_schema_version, + DatabaseClone, +}; +pub use statsai_pricing::{PRICING_CATALOG_VERSION, PRICING_RULESET_VERSION}; use std::time::Duration; #[cfg(test)] @@ -337,7 +346,7 @@ fn sanitize_summary_for_default_http_sync(summary: UsageSummary) -> UsageSummary sanitize_summary_for_http_sync(summary, false) } -fn is_daily_rollup_summary(summary: &UsageSummary) -> bool { +pub(crate) fn is_daily_rollup_summary(summary: &UsageSummary) -> bool { summary.metadata.summary_format == "daily_rollup.v1" } @@ -364,7 +373,7 @@ fn summary_sync_day(summary: &UsageSummary) -> NaiveDate { .unwrap_or_else(|| summary.observed_at.date_naive()) } -fn summary_period_bounds(summary: &UsageSummary) -> (DateTime, DateTime) { +pub(crate) fn summary_period_bounds(summary: &UsageSummary) -> (DateTime, DateTime) { let start = summary .period_start .or(summary.period_end) @@ -1703,10 +1712,30 @@ impl Store { }) } - fn update_event_payload(&self, event: &UsageEvent) -> Result> { + pub(crate) fn update_event_payload( + &self, + event: &UsageEvent, + ) -> Result> { let existing_bucket = self .event_by_id(&event.event_id.0)? .map(|existing| sync_rollup_bucket_key(&existing)); + let bucket = self.update_event_cost_payload(event)?; + let mut dirty_keys = BTreeSet::new(); + if let Some(existing_bucket) = existing_bucket { + dirty_keys.insert(existing_bucket); + } + dirty_keys.insert(bucket); + Ok(dirty_keys) + } + + /// Updates a persisted event's payload without re-reading it. + /// + /// Repricing only changes estimated cost, so the sync-rollup bucket is the + /// in-memory event's bucket. + pub(crate) fn update_event_cost_payload( + &self, + event: &UsageEvent, + ) -> Result { let payload = serde_json::to_string(event)?; let fingerprint = event_fingerprint(event); self.conn.execute( @@ -1732,12 +1761,7 @@ impl Store { &payload ], )?; - let mut dirty_keys = BTreeSet::new(); - if let Some(existing_bucket) = existing_bucket { - dirty_keys.insert(existing_bucket); - } - dirty_keys.insert(sync_rollup_bucket_key(event)); - Ok(dirty_keys) + Ok(sync_rollup_bucket_key(event)) } fn insert_event_in_batch( @@ -3908,7 +3932,7 @@ impl Store { Ok(deleted) } - fn event_by_id(&self, event_id: &str) -> Result> { + pub(crate) fn event_by_id(&self, event_id: &str) -> Result> { self.conn .query_row( "SELECT payload FROM usage_events WHERE event_id = ?1", @@ -3931,20 +3955,30 @@ impl Store { } fn refresh_sync_rollups_for_keys(&self, keys: &BTreeSet) -> Result<()> { + self.refresh_sync_rollups_for_keys_counted(keys).map(|_| ()) + } + + pub(crate) fn refresh_sync_rollups_for_keys_counted( + &self, + keys: &BTreeSet, + ) -> Result { + let mut refreshed = 0u64; for key in keys { - self.refresh_sync_rollup_for_key(key)?; + if self.refresh_sync_rollup_for_key(key)? { + refreshed += 1; + } } - Ok(()) + Ok(refreshed) } - fn refresh_sync_rollup_for_key(&self, key: &SyncRollupBucketKey) -> Result<()> { + fn refresh_sync_rollup_for_key(&self, key: &SyncRollupBucketKey) -> Result { let events = self.sync_rollup_events(key)?; if events.is_empty() { - self.conn.execute( + let deleted = self.conn.execute( "DELETE FROM sync_rollups WHERE summary_id = ?1", params![sync_rollup_summary_id(key).0], )?; - return Ok(()); + return Ok(deleted > 0); } let summary = build_sync_rollup_summary(&events); @@ -3963,7 +3997,7 @@ impl Store { .as_ref() .is_some_and(|(existing_hash, _)| existing_hash == &payload_hash) { - return Ok(()); + return Ok(false); } let dirty = existing.as_ref().map_or(1, |(_, dirty)| (*dirty).max(1)); @@ -3998,7 +4032,7 @@ impl Store { &payload, ], )?; - Ok(()) + Ok(true) } fn sync_rollup_events(&self, key: &SyncRollupBucketKey) -> Result> { @@ -4046,7 +4080,9 @@ impl Store { |row| row.get::<_, String>(0), )?; for row in rows { - events.push(serde_json::from_str(&row?)?); + if let Ok(event) = serde_json::from_str(&row?) { + events.push(event); + } } } else { let rows = stmt.query_map( @@ -4054,7 +4090,9 @@ impl Store { |row| row.get::<_, String>(0), )?; for row in rows { - events.push(serde_json::from_str(&row?)?); + if let Ok(event) = serde_json::from_str(&row?) { + events.push(event); + } } } events.retain(|event| sync_rollup_project_key(event.project.as_ref()) == key.project_key); @@ -5248,7 +5286,7 @@ fn assignment_for_timestamp( } #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] -struct SyncRollupBucketKey { +pub(crate) struct SyncRollupBucketKey { provider: String, source_id: String, provider_account_id: Option, @@ -5263,7 +5301,7 @@ struct EventInsertOutcome { dirty_keys: BTreeSet, } -fn sync_rollup_bucket_key(event: &UsageEvent) -> SyncRollupBucketKey { +pub(crate) fn sync_rollup_bucket_key(event: &UsageEvent) -> SyncRollupBucketKey { SyncRollupBucketKey { provider: event.provider.clone(), source_id: event.source_id.0.clone(), diff --git a/crates/statsai-store/src/pricing.rs b/crates/statsai-store/src/pricing.rs new file mode 100644 index 0000000..4ce4743 --- /dev/null +++ b/crates/statsai-store/src/pricing.rs @@ -0,0 +1,1750 @@ +//! Versioned automatic repricing of persisted normalized usage. + +use super::{is_daily_rollup_summary, summary_period_bounds, Store, SyncRollupBucketKey}; +use anyhow::Result; +use rusqlite::params; +use serde::de::DeserializeOwned; +use statsai_core::{CostAccumulator, CostInfo, ModelInfo, UsageEvent, UsageSummary}; +use statsai_pricing::{ + estimate_cost_at, overlay_estimated_cost, pricing_changes_between, unknown_cost, + PRICING_CATALOG_VERSION, PRICING_RULESET_VERSION, +}; +use std::collections::{BTreeSet, HashMap}; +use std::fmt; + +pub const APPLIED_PRICING_RULESET_VERSION_KEY: &str = "pricing.applied_ruleset_version"; +pub const APPLIED_PRICING_CATALOG_VERSION_KEY: &str = "pricing.applied_catalog_version"; + +const EVENT_PAGE_SIZE: usize = 256; +const SUMMARY_PAGE_SIZE: usize = 128; +const TASK_SPAN_PAGE_SIZE: usize = 128; + +/// Counts produced by one automatic pricing-ruleset application. +#[derive(Debug, Clone, PartialEq, Eq, Default)] +pub struct RepricingReport { + pub examined_events: u64, + pub changed_events: u64, + pub skipped_unreadable_events: u64, + pub examined_summaries: u64, + pub changed_summaries: u64, + pub skipped_unreadable_summaries: u64, + pub refreshed_rollups: u64, + pub changed_task_spans: u64, + pub skipped_unreadable_spans: u64, + pub rebuilt_work_items: u64, + pub already_current: bool, +} + +impl RepricingReport { + #[must_use] + pub fn did_work(&self) -> bool { + self.changed_events > 0 + || self.changed_summaries > 0 + || self.refreshed_rollups > 0 + || self.changed_task_spans > 0 + || self.rebuilt_work_items > 0 + } +} + +impl fmt::Display for RepricingReport { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + if self.already_current { + return write!( + formatter, + "pricing already at ruleset {PRICING_RULESET_VERSION} ({PRICING_CATALOG_VERSION})" + ); + } + write!( + formatter, + "repriced store to ruleset {PRICING_RULESET_VERSION} ({PRICING_CATALOG_VERSION}): examined_events={} changed_events={} skipped_unreadable_events={} changed_summaries={} refreshed_rollups={}", + self.examined_events, + self.changed_events, + self.skipped_unreadable_events, + self.changed_summaries, + self.refreshed_rollups + )?; + let skipped_other = self + .skipped_unreadable_summaries + .saturating_add(self.skipped_unreadable_spans); + if skipped_other > 0 { + write!( + formatter, + " skipped_unreadable_summaries={} skipped_unreadable_spans={}", + self.skipped_unreadable_summaries, self.skipped_unreadable_spans + )?; + } + Ok(()) + } +} + +impl Store { + /// Ensures persisted estimated prices match this binary's pricing ruleset. + /// + /// This is not called from [`Store::open`]. The selected StatsAI binary + /// owns the decision so development launchers compiled against a different + /// catalog cannot reprice a store they did not select. + pub fn ensure_current_pricing(&self) -> Result { + self.ensure_pricing_ruleset(PRICING_RULESET_VERSION, PRICING_CATALOG_VERSION) + } + + /// Returns the last successfully applied pricing ruleset, if recorded. + pub fn applied_pricing_ruleset_version(&self) -> Result> { + super::snapshot::parse_applied_pricing_ruleset_value( + self.metadata_value(APPLIED_PRICING_RULESET_VERSION_KEY)? + .as_deref(), + ) + } + + pub(crate) fn ensure_pricing_ruleset( + &self, + current_ruleset: u64, + catalog_version: &str, + ) -> Result { + match self.applied_pricing_ruleset_version()? { + Some(applied) if applied > current_ruleset => { + return Err(forward_pricing_version_error(applied, current_ruleset)); + } + Some(applied) if applied == current_ruleset => { + return Ok(RepricingReport { + already_current: true, + ..RepricingReport::default() + }); + } + _ => {} + } + + self.with_immediate_transaction(|| { + match self.applied_pricing_ruleset_version()? { + Some(applied) if applied > current_ruleset => { + return Err(forward_pricing_version_error(applied, current_ruleset)); + } + Some(applied) if applied == current_ruleset => { + return Ok(RepricingReport { + already_current: true, + ..RepricingReport::default() + }); + } + _ => {} + } + + let mut report = RepricingReport::default(); + let mut dirty_keys = BTreeSet::new(); + + self.reprice_events_in_tx(&mut report, &mut dirty_keys)?; + self.reprice_summaries_in_tx(&mut report)?; + report.refreshed_rollups = self.refresh_sync_rollups_for_keys_counted(&dirty_keys)?; + self.reprice_task_spans_in_tx(&mut report)?; + + self.set_metadata_value( + APPLIED_PRICING_RULESET_VERSION_KEY, + ¤t_ruleset.to_string(), + )?; + self.set_metadata_value(APPLIED_PRICING_CATALOG_VERSION_KEY, catalog_version)?; + Ok(report) + }) + } +} + +fn forward_pricing_version_error(applied: u64, supported: u64) -> anyhow::Error { + anyhow::anyhow!( + "database pricing ruleset version {applied} is newer than this StatsAI binary supports ({supported}); upgrade StatsAI or use a compatible database" + ) +} + +struct RepricePage { + items: Vec, + last_id: Option, + fetched: usize, + skipped: u64, +} + +fn decode_id_payload_page( + rows: impl Iterator>, +) -> Result> { + let mut page = RepricePage { + items: Vec::new(), + last_id: None, + fetched: 0, + skipped: 0, + }; + for row in rows { + let (id, payload) = row?; + page.fetched += 1; + page.last_id = Some(id); + match serde_json::from_str(&payload) { + Ok(item) => page.items.push(item), + Err(_) => page.skipped += 1, + } + } + Ok(page) +} + +impl Store { + fn reprice_events_in_tx( + &self, + report: &mut RepricingReport, + dirty_keys: &mut BTreeSet, + ) -> Result<()> { + let mut after: Option = None; + loop { + let page = self.event_page_after(after.as_deref(), EVENT_PAGE_SIZE)?; + if page.fetched == 0 { + break; + } + after = page.last_id; + report.skipped_unreadable_events += page.skipped; + for event in page.items { + report.examined_events += 1; + maybe_fail_after_event_writes(report.changed_events)?; + if let Some(updated) = reprice_event(&event) { + dirty_keys.insert(self.update_event_cost_payload(&updated)?); + report.changed_events += 1; + } + } + if page.fetched < EVENT_PAGE_SIZE { + break; + } + } + Ok(()) + } + + fn event_page_after( + &self, + after: Option<&str>, + limit: usize, + ) -> Result> { + if let Some(after) = after { + let mut statement = self.conn.prepare( + "SELECT event_id, payload FROM usage_events WHERE event_id > ?1 ORDER BY event_id LIMIT ?2", + )?; + let rows = statement.query_map(params![after, limit as i64], |row| { + Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)) + })?; + decode_id_payload_page(rows) + } else { + let mut statement = self + .conn + .prepare("SELECT event_id, payload FROM usage_events ORDER BY event_id LIMIT ?1")?; + let rows = statement.query_map(params![limit as i64], |row| { + Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)) + })?; + decode_id_payload_page(rows) + } + } + + fn reprice_summaries_in_tx(&self, report: &mut RepricingReport) -> Result<()> { + let mut after: Option = None; + loop { + let page = self.summary_page_after(after.as_deref(), SUMMARY_PAGE_SIZE)?; + if page.fetched == 0 { + break; + } + after = page.last_id; + report.skipped_unreadable_summaries += page.skipped; + for summary in page.items { + report.examined_summaries += 1; + if is_daily_rollup_summary(&summary) { + continue; + } + if let Some(updated) = reprice_summary(&summary) { + self.upsert_summary(&updated)?; + report.changed_summaries += 1; + } + } + if page.fetched < SUMMARY_PAGE_SIZE { + break; + } + } + Ok(()) + } + + fn summary_page_after( + &self, + after: Option<&str>, + limit: usize, + ) -> Result> { + if let Some(after) = after { + let mut statement = self.conn.prepare( + "SELECT summary_id, payload FROM usage_summaries WHERE summary_id > ?1 ORDER BY summary_id LIMIT ?2", + )?; + let rows = statement.query_map(params![after, limit as i64], |row| { + Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)) + })?; + decode_id_payload_page(rows) + } else { + let mut statement = self.conn.prepare( + "SELECT summary_id, payload FROM usage_summaries ORDER BY summary_id LIMIT ?1", + )?; + let rows = statement.query_map(params![limit as i64], |row| { + Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)) + })?; + decode_id_payload_page(rows) + } + } + + fn reprice_task_spans_in_tx(&self, report: &mut RepricingReport) -> Result<()> { + let mut after: Option = None; + let mut changed_buckets = BTreeSet::new(); + loop { + let page = self.linked_task_span_page_after(after.as_deref(), TASK_SPAN_PAGE_SIZE)?; + if page.fetched == 0 { + break; + } + after = page.last_id; + report.skipped_unreadable_spans += page.skipped; + let event_ids = page + .items + .iter() + .flat_map(|span| { + span.linked_event_ids + .iter() + .map(|event_id| event_id.0.clone()) + }) + .collect::>(); + let events = self.events_by_ids(&event_ids)?; + let mut updated_spans = Vec::new(); + for mut span in page.items { + let Some((cents, micro)) = + estimated_cost_for_loaded_events(&span.linked_event_ids, &events) + else { + continue; + }; + if span.estimated_cost_usd != cents || span.estimated_cost_micro_usd != micro { + span.estimated_cost_usd = cents; + span.estimated_cost_micro_usd = micro; + changed_buckets.insert(span.project_bucket.clone()); + updated_spans.push(span); + report.changed_task_spans += 1; + } + } + if !updated_spans.is_empty() { + self.upsert_task_spans_in_tx(&updated_spans)?; + } + if page.fetched < TASK_SPAN_PAGE_SIZE { + break; + } + } + if !changed_buckets.is_empty() { + report.rebuilt_work_items = + self.rebuild_task_work_items_for_project_buckets(&changed_buckets)?; + } + Ok(()) + } + + fn linked_task_span_page_after( + &self, + after: Option<&str>, + limit: usize, + ) -> Result> { + let sql = if after.is_some() { + "SELECT span_id, payload FROM task_spans + WHERE span_id > ?1 + AND EXISTS ( + SELECT 1 FROM task_span_event_links WHERE span_id = task_spans.span_id + ) + ORDER BY span_id LIMIT ?2" + } else { + "SELECT span_id, payload FROM task_spans + WHERE EXISTS ( + SELECT 1 FROM task_span_event_links WHERE span_id = task_spans.span_id + ) + ORDER BY span_id LIMIT ?1" + }; + if let Some(after) = after { + let mut statement = self.conn.prepare(sql)?; + let rows = statement.query_map(params![after, limit as i64], |row| { + Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)) + })?; + decode_id_payload_page(rows) + } else { + let mut statement = self.conn.prepare(sql)?; + let rows = statement.query_map(params![limit as i64], |row| { + Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)) + })?; + decode_id_payload_page(rows) + } + } + + fn events_by_ids(&self, event_ids: &BTreeSet) -> Result> { + let mut events = HashMap::new(); + let ids = event_ids.iter().cloned().collect::>(); + for chunk in ids.chunks(128) { + if chunk.is_empty() { + continue; + } + let placeholders = (0..chunk.len()).map(|_| "?").collect::>().join(","); + let sql = format!( + "SELECT event_id, payload FROM usage_events WHERE event_id IN ({placeholders})" + ); + let mut statement = self.conn.prepare(&sql)?; + let params: Vec<&dyn rusqlite::types::ToSql> = chunk + .iter() + .map(|id| id as &dyn rusqlite::types::ToSql) + .collect(); + let rows = statement.query_map(params.as_slice(), |row| { + Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)) + })?; + for row in rows { + let (event_id, payload) = row?; + if let Ok(event) = serde_json::from_str(&payload) { + events.insert(event_id, event); + } + } + } + Ok(events) + } +} + +fn estimated_cost_for_loaded_events( + event_ids: &[statsai_core::EventId], + events: &HashMap, +) -> Option<(Option, Option)> { + if event_ids.is_empty() { + return None; + } + let mut estimated = CostAccumulator::default(); + let mut found = 0usize; + for event_id in event_ids { + if let Some(event) = events.get(&event_id.0) { + estimated.add_estimated(&event.cost); + found += 1; + } + } + if found == 0 { + return Some((None, None)); + } + Some((estimated.cents_rounded(), estimated.micro_usd())) +} + +fn reprice_event(event: &UsageEvent) -> Option { + let estimated = estimate_cost_at( + &event.provider, + event.model.as_ref(), + &event.usage, + &event.session.started_at, + ); + let cost = overlay_estimated_cost(&event.cost, estimated); + if cost == event.cost { + return None; + } + let mut updated = event.clone(); + updated.cost = cost; + Some(updated) +} + +/// Overlays this binary's estimated pricing onto a summary. +/// +/// Daily rollups are left unchanged. Returns the original summary when +/// estimated fields already match. +#[must_use] +pub fn apply_current_estimated_pricing(summary: UsageSummary) -> UsageSummary { + if is_daily_rollup_summary(&summary) { + return summary; + } + reprice_summary(&summary).unwrap_or(summary) +} + +fn reprice_summary(summary: &UsageSummary) -> Option { + let (period_start, period_end) = summary_period_bounds(summary); + let pricing_at = summary.period_end.unwrap_or(summary.observed_at); + let mut updated = summary.clone(); + let mut changed = false; + + if !updated.models.is_empty() { + for model_usage in &mut updated.models { + let next = estimated_summary_cost( + &updated.provider, + Some(&model_usage.model), + &model_usage.usage, + period_start.date_naive(), + period_end.date_naive(), + pricing_at, + &model_usage.cost, + ); + if next != model_usage.cost { + model_usage.cost = next; + changed = true; + } + } + } + + let next = estimated_summary_cost( + &updated.provider, + updated.model.as_ref(), + &updated.usage, + period_start.date_naive(), + period_end.date_naive(), + pricing_at, + &updated.cost, + ); + if next != updated.cost { + updated.cost = next; + changed = true; + } + + changed.then_some(updated) +} + +fn estimated_summary_cost( + provider: &str, + model: Option<&ModelInfo>, + usage: &statsai_core::UsageCounts, + period_start: chrono::NaiveDate, + period_end: chrono::NaiveDate, + pricing_at: chrono::DateTime, + existing: &CostInfo, +) -> CostInfo { + let estimated = if summary_crosses_pricing_boundary(model, period_start, period_end) { + unknown_cost() + } else { + estimate_cost_at(provider, model, usage, &pricing_at) + }; + overlay_estimated_cost(existing, estimated) +} + +fn summary_crosses_pricing_boundary( + model: Option<&ModelInfo>, + period_start: chrono::NaiveDate, + period_end: chrono::NaiveDate, +) -> bool { + model_pricing_names(model) + .into_iter() + .any(|name| pricing_changes_between(&name, period_start, period_end)) +} + +fn model_pricing_names(model: Option<&ModelInfo>) -> Vec { + let Some(model) = model else { + return Vec::new(); + }; + [ + model.normalized_name.as_deref(), + model.name.as_deref(), + model.provider_model_id.as_deref(), + ] + .into_iter() + .flatten() + .map(ToOwned::to_owned) + .collect() +} + +fn maybe_fail_after_event_writes(changed_events: u64) -> Result<()> { + #[cfg(test)] + { + FAIL_AFTER_EVENT_WRITES.with(|cell| { + if cell.get().is_some_and(|limit| changed_events >= limit) { + anyhow::bail!("injected repricing failure after {changed_events} event writes") + } else { + Ok(()) + } + })?; + } + let _ = changed_events; + Ok(()) +} + +#[cfg(test)] +thread_local! { + static FAIL_AFTER_EVENT_WRITES: std::cell::Cell> = const { std::cell::Cell::new(None) }; +} + +#[cfg(test)] +pub(crate) fn fail_repricing_after_event_writes(limit: Option) { + FAIL_AFTER_EVENT_WRITES.with(|cell| cell.set(limit)); +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::CURRENT_SCHEMA_VERSION; + use chrono::Utc; + use statsai_core::{ + event_id, summary_id, task_span_id, Confidence, CostInfo, EventSource, LocationOrigin, + ModelInfo, PrivacyInfo, PrivacyMode, SessionInfo, SourceKind, SummaryMetadata, TaskSpan, + UsageCounts, UsageEvent, UsageSummary, TASK_SPAN_SCHEMA_VERSION, + USAGE_EVENT_SCHEMA_VERSION, USAGE_SUMMARY_SCHEMA_VERSION, + }; + use std::path::Path; + use std::sync::{Arc, Barrier}; + + fn test_model(name: &str) -> ModelInfo { + ModelInfo { + name: Some(name.to_string()), + normalized_name: Some(name.to_string()), + provider_model_id: Some(name.to_string()), + speed: None, + reasoning_level: None, + reasoning_level_raw: None, + } + } + + fn parse_utc(value: &str) -> chrono::DateTime { + chrono::DateTime::parse_from_rfc3339(value) + .expect("valid timestamp") + .with_timezone(&chrono::Utc) + } + + fn test_source(path: &str) -> statsai_core::SourceLocation { + statsai_core::SourceLocation::local_adapter( + "codex", + "test", + "0", + Path::new(path), + LocationOrigin::Configured, + ) + } + + fn test_event( + source: &statsai_core::SourceLocation, + started_at: chrono::DateTime, + record_id: &str, + model: &str, + usage: UsageCounts, + cost: CostInfo, + ) -> UsageEvent { + UsageEvent { + schema_version: USAGE_EVENT_SCHEMA_VERSION.to_string(), + event_id: event_id("codex", &source.source_id, record_id, None, started_at), + device_id: "device".to_string(), + provider: "codex".to_string(), + source_id: source.source_id.clone(), + provider_account_id: None, + subscription_id: None, + source: EventSource { + adapter_id: "test".to_string(), + adapter_version: "0".to_string(), + source_kind: SourceKind::LocalAdapter, + location_origin: Some(LocationOrigin::Configured), + source_type: "jsonl".to_string(), + source_path_hash: source.path_hash.clone(), + source_record_id: Some(record_id.to_string()), + parse_confidence: Confidence::High, + }, + session: SessionInfo { + session_id: "session".to_string(), + local_session_id_hash: Some("same-session".to_string()), + title: None, + started_at, + ended_at: None, + duration_seconds: None, + }, + model: Some(test_model(model)), + usage, + runtime: None, + cost, + parse_evidence: None, + project: None, + git: None, + privacy: PrivacyInfo { + mode: PrivacyMode::MetadataOnly, + contains_prompt_text: false, + contains_response_text: false, + contains_file_paths: false, + }, + created_at: started_at, + imported_at: started_at, + } + } + + fn missing_cost() -> CostInfo { + CostInfo { + currency: "USD".to_string(), + estimated_api_equivalent_usd: None, + provider_reported_usd: None, + estimated_api_equivalent_micro_usd: None, + provider_reported_micro_usd: None, + pricing_source: Some("unknown".to_string()), + pricing_version: None, + confidence: Confidence::Low, + } + } + + fn million_token_usage() -> UsageCounts { + UsageCounts { + input_tokens: Some(1_000_000), + cache_creation_tokens: Some(1_000_000), + cache_read_tokens: Some(1_000_000), + output_tokens: Some(1_000_000), + total_tokens: Some(4_000_000), + ..UsageCounts::default() + } + } + + fn expected_review_cost(started_at: chrono::DateTime) -> CostInfo { + estimate_cost_at( + "codex", + Some(&test_model("codex-auto-review")), + &million_token_usage(), + &started_at, + ) + } + + fn store_with_source(path: &str) -> (Store, statsai_core::SourceLocation) { + let store = Store::in_memory().expect("store"); + let source = test_source(path); + store.upsert_source(&source).expect("source"); + (store, source) + } + + fn stored_event(store: &Store, event_id: &str) -> UsageEvent { + store + .event_by_id(event_id) + .expect("load event") + .expect("event present") + } + + fn test_span( + source: &statsai_core::SourceLocation, + started_at: chrono::DateTime, + record_id: &str, + linked_event_ids: Vec, + estimated_cost_usd: Option, + estimated_cost_micro_usd: Option, + ) -> TaskSpan { + TaskSpan { + schema_version: TASK_SPAN_SCHEMA_VERSION.to_string(), + span_id: task_span_id("codex", &source.source_id, record_id), + provider: "codex".to_string(), + source_id: source.source_id.clone(), + span_kind: "codex_session".to_string(), + source_record_id: Some(record_id.to_string()), + source_file_path_hash: None, + summary_id: None, + session_id: Some("session".to_string()), + thread_id: None, + title: "Review".to_string(), + normalized_title: "review".to_string(), + title_source: Some("summary".to_string()), + summary_preview: None, + todo_excerpt: None, + issue_keys: Vec::new(), + branch_family: None, + project_bucket: "none".to_string(), + project: None, + git: None, + usage: million_token_usage(), + estimated_cost_usd, + estimated_cost_micro_usd, + event_count: linked_event_ids.len() as u64, + has_usage_evidence: !linked_event_ids.is_empty(), + total_messages: 0, + user_messages: 0, + assistant_messages: 0, + developer_messages: 0, + linked_event_ids, + confidence: Confidence::Medium, + is_meta: false, + started_at, + ended_at: Some(started_at), + duration_seconds: Some(0), + } + } + + fn test_summary( + source: &statsai_core::SourceLocation, + model: &str, + start: chrono::DateTime, + end: chrono::DateTime, + cost: CostInfo, + ) -> UsageSummary { + UsageSummary { + schema_version: USAGE_SUMMARY_SCHEMA_VERSION.to_string(), + summary_id: summary_id("codex", &source.source_id, "period"), + device_id: "device".to_string(), + provider: "codex".to_string(), + source_id: source.source_id.clone(), + provider_account_id: None, + source: EventSource { + adapter_id: "test".to_string(), + adapter_version: "0".to_string(), + source_kind: SourceKind::LocalAdapter, + location_origin: Some(LocationOrigin::Configured), + source_type: "stats-cache.json".to_string(), + source_path_hash: source.path_hash.clone(), + source_record_id: Some("period".to_string()), + parse_confidence: Confidence::Medium, + }, + model: Some(test_model(model)), + models: Vec::new(), + usage: million_token_usage(), + cost, + parse_evidence: None, + project: None, + privacy: PrivacyInfo { + mode: PrivacyMode::MetadataOnly, + contains_prompt_text: false, + contains_response_text: false, + contains_file_paths: false, + }, + metrics: None, + period_start: Some(start), + period_end: Some(end), + observed_at: end, + metadata: SummaryMetadata { + summary_format: "grok_build_session_summary".to_string(), + summary_version: Some("1".to_string()), + total_sessions: Some(1), + total_messages: Some(2), + last_computed_at: Some(end), + }, + imported_at: end, + } + } + + #[test] + fn store_open_does_not_reprice() { + let (store, source) = store_with_source("/tmp/codex-open-no-reprice"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store.insert_event(&event).expect("insert"); + + assert_eq!(store.applied_pricing_ruleset_version().expect("meta"), None); + assert!(stored_event(&store, &event.event_id.0) + .cost + .estimated_api_equivalent_usd + .is_none()); + } + + #[test] + fn legacy_codex_auto_review_event_is_repriced_without_source_files() { + let (store, source) = store_with_source("/tmp/codex-legacy-review"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store.insert_event(&event).expect("insert"); + + let report = store.ensure_current_pricing().expect("reprice"); + + assert_eq!(report.examined_events, 1); + assert_eq!(report.changed_events, 1); + assert!(!report.already_current); + let stored = stored_event(&store, &event.event_id.0); + assert_eq!(stored.cost, expected_review_cost(started_at)); + assert_eq!(stored.event_id, event.event_id); + assert_eq!(stored.usage, event.usage); + assert_eq!(stored.session.started_at, event.session.started_at); + assert_eq!( + store.applied_pricing_ruleset_version().expect("applied"), + Some(PRICING_RULESET_VERSION) + ); + assert_eq!( + store + .metadata_value(APPLIED_PRICING_CATALOG_VERSION_KEY) + .expect("catalog"), + Some(PRICING_CATALOG_VERSION.to_string()) + ); + } + + #[test] + fn date_aware_codex_auto_review_boundary_is_preserved() { + let (store, source) = store_with_source("/tmp/codex-review-boundary"); + let before = parse_utc("2026-07-29T23:59:59Z"); + let after = parse_utc("2026-07-30T00:00:00Z"); + let before_event = test_event( + &source, + before, + "before-boundary", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + let after_event = test_event( + &source, + after, + "after-boundary", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store + .insert_events(&[before_event.clone(), after_event.clone()]) + .expect("insert"); + + store.ensure_current_pricing().expect("reprice"); + + let before_cost = stored_event(&store, &before_event.event_id.0).cost; + let after_cost = stored_event(&store, &after_event.event_id.0).cost; + assert_eq!(before_cost, expected_review_cost(before)); + assert_eq!(after_cost, expected_review_cost(after)); + assert_ne!( + before_cost.estimated_api_equivalent_usd, + after_cost.estimated_api_equivalent_usd + ); + assert_eq!(after_cost.estimated_api_equivalent_usd, Some(167)); + } + + #[test] + fn applied_metadata_advances_only_after_success() { + let (store, source) = store_with_source("/tmp/codex-reprice-success-meta"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + store + .insert_event(&test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + )) + .expect("insert"); + + assert_eq!( + store.applied_pricing_ruleset_version().expect("before"), + None + ); + store.ensure_current_pricing().expect("reprice"); + assert_eq!( + store.applied_pricing_ruleset_version().expect("after"), + Some(PRICING_RULESET_VERSION) + ); + } + + #[test] + fn second_invocation_at_the_same_version_is_a_noop() { + let (store, source) = store_with_source("/tmp/codex-reprice-noop"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + store + .insert_event(&test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + )) + .expect("insert"); + + let first = store.ensure_current_pricing().expect("first"); + let payload_before = store.events().expect("events"); + let second = store.ensure_current_pricing().expect("second"); + + assert_eq!(first.changed_events, 1); + assert!(second.already_current); + assert_eq!(second.changed_events, 0); + assert_eq!(second.examined_events, 0); + assert_eq!(store.events().expect("events after noop"), payload_before); + } + + #[test] + fn provider_reported_cost_and_provenance_survive_repricing() { + let (store, source) = store_with_source("/tmp/codex-provider-reported"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let mut cost = missing_cost(); + cost.provider_reported_usd = Some(99); + cost.provider_reported_micro_usd = Some(990_000); + cost.pricing_source = Some("provider_invoice".to_string()); + cost.confidence = Confidence::High; + let event = test_event( + &source, + started_at, + "reported", + "codex-auto-review", + million_token_usage(), + cost, + ); + store.insert_event(&event).expect("insert"); + + store.ensure_current_pricing().expect("reprice"); + let stored = stored_event(&store, &event.event_id.0); + assert_eq!(stored.cost.provider_reported_usd, Some(99)); + assert_eq!(stored.cost.provider_reported_micro_usd, Some(990_000)); + assert_eq!( + stored.cost.pricing_source.as_deref(), + Some("provider_invoice") + ); + assert_eq!(stored.cost.confidence, Confidence::High); + assert_eq!( + stored.cost.estimated_api_equivalent_usd, + expected_review_cost(started_at).estimated_api_equivalent_usd + ); + assert_eq!( + stored.cost.pricing_version.as_deref(), + Some(PRICING_CATALOG_VERSION) + ); + } + + #[test] + fn apply_current_estimated_pricing_overlays_stale_estimated_only_summary() { + let source = test_source("/tmp/codex-import-overlay"); + let start = parse_utc("2026-07-29T00:00:00Z"); + let end = parse_utc("2026-07-29T23:59:59Z"); + let mut summary = test_summary(&source, "codex-auto-review", start, end, missing_cost()); + summary.cost.estimated_api_equivalent_usd = Some(1); + summary.cost.estimated_api_equivalent_micro_usd = Some(10_000); + summary.cost.pricing_source = Some("official:stale".to_string()); + summary.cost.pricing_version = Some("official:stale".to_string()); + + let priced = apply_current_estimated_pricing(summary); + assert_eq!(priced.cost, expected_review_cost(end)); + } + + #[test] + fn apply_current_estimated_pricing_keeps_provider_reported_amount() { + let source = test_source("/tmp/codex-import-provider-reported"); + let start = parse_utc("2026-07-29T00:00:00Z"); + let end = parse_utc("2026-07-29T23:59:59Z"); + let mut cost = missing_cost(); + cost.provider_reported_usd = Some(99); + cost.provider_reported_micro_usd = Some(990_000); + cost.pricing_source = Some("provider_invoice".to_string()); + cost.confidence = Confidence::High; + let summary = test_summary(&source, "codex-auto-review", start, end, cost); + + let priced = apply_current_estimated_pricing(summary); + assert_eq!(priced.cost.provider_reported_usd, Some(99)); + assert_eq!(priced.cost.provider_reported_micro_usd, Some(990_000)); + assert_eq!( + priced.cost.pricing_source.as_deref(), + Some("provider_invoice") + ); + assert_eq!( + priced.cost.estimated_api_equivalent_usd, + expected_review_cost(end).estimated_api_equivalent_usd + ); + } + + #[test] + fn unknown_model_stays_unknown_while_ruleset_is_marked_applied() { + let (store, source) = store_with_source("/tmp/codex-unknown-model"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "unknown", + "not-a-real-model", + million_token_usage(), + missing_cost(), + ); + store.insert_event(&event).expect("insert"); + + let report = store.ensure_current_pricing().expect("reprice"); + let stored = stored_event(&store, &event.event_id.0); + assert_eq!(report.changed_events, 0); + assert!(stored.cost.estimated_api_equivalent_usd.is_none()); + assert_eq!(stored.cost.pricing_source.as_deref(), Some("unknown")); + assert_eq!( + store.applied_pricing_ruleset_version().expect("applied"), + Some(PRICING_RULESET_VERSION) + ); + } + + #[test] + fn summary_spanning_a_pricing_boundary_remains_unknown() { + let (store, source) = store_with_source("/tmp/codex-boundary-summary"); + let start = parse_utc("2026-07-29T00:00:00Z"); + let end = parse_utc("2026-07-31T00:00:00Z"); + let mut summary = UsageSummary { + schema_version: USAGE_SUMMARY_SCHEMA_VERSION.to_string(), + summary_id: summary_id("codex", &source.source_id, "crossing"), + device_id: "device".to_string(), + provider: "codex".to_string(), + source_id: source.source_id.clone(), + provider_account_id: None, + source: EventSource { + adapter_id: "test".to_string(), + adapter_version: "0".to_string(), + source_kind: SourceKind::LocalSummary, + location_origin: Some(LocationOrigin::Configured), + source_type: "stats-cache.json".to_string(), + source_path_hash: source.path_hash.clone(), + source_record_id: Some("crossing".to_string()), + parse_confidence: Confidence::Medium, + }, + model: Some(test_model("codex-auto-review")), + models: Vec::new(), + usage: million_token_usage(), + cost: missing_cost(), + parse_evidence: None, + project: None, + privacy: PrivacyInfo { + mode: PrivacyMode::MetadataOnly, + contains_prompt_text: false, + contains_response_text: false, + contains_file_paths: false, + }, + metrics: None, + period_start: Some(start), + period_end: Some(end), + observed_at: end, + metadata: SummaryMetadata { + summary_format: "claude_stats_cache".to_string(), + summary_version: Some("1".to_string()), + total_sessions: Some(1), + total_messages: Some(2), + last_computed_at: Some(end), + }, + imported_at: end, + }; + summary.cost.estimated_api_equivalent_usd = Some(999); + store.upsert_summary(&summary).expect("summary"); + + let report = store.ensure_current_pricing().expect("reprice"); + let stored = store + .summaries() + .expect("summaries") + .into_iter() + .next() + .expect("one summary"); + assert_eq!(report.changed_summaries, 1); + assert!(stored.cost.estimated_api_equivalent_usd.is_none()); + assert_eq!(stored.cost.pricing_source.as_deref(), Some("unknown")); + } + + #[test] + fn changed_sync_rollups_are_refreshed_and_marked_dirty() { + let (store, source) = store_with_source("/tmp/codex-reprice-rollups"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store.insert_event(&event).expect("insert"); + store.rebuild_sync_rollups().expect("rebuild"); + let rollups = store.all_sync_rollup_summaries().expect("rollups"); + store + .mark_sync_rollups_synced( + &rollups + .iter() + .map(|summary| summary.summary_id.clone()) + .collect::>(), + ) + .expect("mark synced"); + assert!(store + .dirty_sync_rollup_summaries() + .expect("clean") + .is_empty()); + + let report = store.ensure_current_pricing().expect("reprice"); + let dirty = store.dirty_sync_rollup_summaries().expect("dirty"); + assert_eq!(report.refreshed_rollups, 1); + assert_eq!(dirty.len(), 1); + assert_eq!( + dirty[0].cost.estimated_api_equivalent_usd, + expected_review_cost(started_at).estimated_api_equivalent_usd + ); + } + + #[test] + fn injected_mid_operation_error_rolls_back_payloads_rollups_and_metadata() { + let (store, source) = store_with_source("/tmp/codex-reprice-rollback"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let first = test_event( + &source, + started_at, + "first", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + let second = test_event( + &source, + started_at + chrono::Duration::seconds(1), + "second", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store + .insert_events(&[first.clone(), second.clone()]) + .expect("insert"); + store.rebuild_sync_rollups().expect("rebuild"); + let rollups = store.all_sync_rollup_summaries().expect("rollups"); + store + .mark_sync_rollups_synced( + &rollups + .iter() + .map(|summary| summary.summary_id.clone()) + .collect::>(), + ) + .expect("mark synced"); + let payloads_before = store.events().expect("events before"); + + fail_repricing_after_event_writes(Some(1)); + let error = store + .ensure_current_pricing() + .expect_err("injected failure"); + fail_repricing_after_event_writes(None); + + assert!(error.to_string().contains("injected repricing failure")); + assert_eq!(store.events().expect("events after"), payloads_before); + assert!(store + .dirty_sync_rollup_summaries() + .expect("dirty after rollback") + .is_empty()); + assert_eq!(store.applied_pricing_ruleset_version().expect("meta"), None); + } + + #[test] + fn older_ruleset_refuses_a_newer_store_and_does_not_mutate_it() { + let (store, source) = store_with_source("/tmp/codex-forward-pricing"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store.insert_event(&event).expect("insert"); + store + .set_metadata_value(APPLIED_PRICING_RULESET_VERSION_KEY, "99") + .expect("future ruleset"); + store + .set_metadata_value(APPLIED_PRICING_CATALOG_VERSION_KEY, "future") + .expect("future catalog"); + let payload_before = + serde_json::to_string(&stored_event(&store, &event.event_id.0)).expect("serialize"); + + let error = store + .ensure_pricing_ruleset(1, PRICING_CATALOG_VERSION) + .expect_err("forward pricing must refuse"); + + assert!(error + .to_string() + .contains("pricing ruleset version 99 is newer than this StatsAI binary supports (1)")); + assert_eq!( + serde_json::to_string(&stored_event(&store, &event.event_id.0)).expect("serialize"), + payload_before + ); + assert_eq!( + store.applied_pricing_ruleset_version().expect("unchanged"), + Some(99) + ); + } + + #[test] + fn concurrent_callers_do_not_publish_partial_or_duplicate_repricing() { + let directory = tempfile::tempdir().expect("tempdir"); + let path = directory.path().join("statsai.sqlite"); + let setup = Store::open(&path).expect("create store"); + let source = test_source("/tmp/codex-concurrent-reprice"); + setup.upsert_source(&source).expect("source"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + setup + .insert_event(&test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + )) + .expect("insert"); + drop(setup); + + let barrier = Arc::new(Barrier::new(2)); + let reports = std::thread::scope(|scope| { + let first_barrier = Arc::clone(&barrier); + let second_barrier = Arc::clone(&barrier); + let first_path = path.clone(); + let second_path = path.clone(); + let first = scope.spawn(move || { + let store = Store::open(&first_path).expect("open first"); + first_barrier.wait(); + store.ensure_current_pricing() + }); + let second = scope.spawn(move || { + let store = Store::open(&second_path).expect("open second"); + second_barrier.wait(); + store.ensure_current_pricing() + }); + vec![ + first.join().expect("first thread"), + second.join().expect("second thread"), + ] + }); + + let reports = reports + .into_iter() + .map(|result| result.expect("repricing")) + .collect::>(); + let workers = reports + .iter() + .filter(|report| !report.already_current) + .count(); + let noops = reports + .iter() + .filter(|report| report.already_current) + .count(); + assert_eq!(workers, 1); + assert_eq!(noops, 1); + assert_eq!( + reports + .iter() + .map(|report| report.changed_events) + .sum::(), + 1 + ); + + let store = Store::open(&path).expect("reopen"); + assert_eq!( + store.applied_pricing_ruleset_version().expect("applied"), + Some(PRICING_RULESET_VERSION) + ); + let events = store.events().expect("events"); + assert_eq!(events.len(), 1); + assert_eq!(events[0].cost, expected_review_cost(started_at)); + } + + #[test] + fn daily_rollups_are_unused_by_report_and_sync_paths() { + let (store, source) = store_with_source("/tmp/codex-daily-rollup-exclusion"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + store + .insert_event(&test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + )) + .expect("insert"); + let stale = store + .compute_daily_rollup("2026-07-29", "device") + .expect("compute"); + store + .upsert_daily_rollup(&stale) + .expect("seed unused table"); + let before = store + .daily_rollups_between("2026-07-29", "2026-07-29") + .expect("before"); + + store.ensure_current_pricing().expect("reprice"); + + let after = store + .daily_rollups_between("2026-07-29", "2026-07-29") + .expect("after"); + assert_eq!(before, after); + let events = store.events().expect("events"); + assert_eq!(events[0].cost, expected_review_cost(started_at)); + // Pin: bump this when CURRENT_SCHEMA_VERSION changes, after confirming + // the daily_rollups table is still unused by report/sync/snapshot. + assert_eq!(CURRENT_SCHEMA_VERSION, 22); + } + + #[test] + fn task_spans_with_linked_events_are_repriced_from_persisted_events() { + let (store, source) = store_with_source("/tmp/codex-task-span-reprice"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store.insert_event(&event).expect("insert"); + store + .upsert_task_spans(&[test_span( + &source, + started_at, + "span", + vec![event.event_id.clone()], + Some(0), + Some(0), + )]) + .expect("span"); + + let report = store.ensure_current_pricing().expect("reprice"); + let stored = store.task_spans().expect("spans"); + assert_eq!(report.changed_task_spans, 1); + assert_eq!( + stored[0].estimated_cost_usd, + expected_review_cost(started_at).estimated_api_equivalent_usd + ); + } + + #[test] + fn stale_task_spans_are_repriced_when_events_already_match() { + let (store, source) = store_with_source("/tmp/codex-stale-span-current-event"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let expected = expected_review_cost(started_at); + let event = test_event( + &source, + started_at, + "already-priced", + "codex-auto-review", + million_token_usage(), + expected.clone(), + ); + store.insert_event(&event).expect("insert"); + store + .upsert_task_spans(&[test_span( + &source, + started_at, + "stale-span", + vec![event.event_id.clone()], + Some(0), + Some(0), + )]) + .expect("span"); + + let report = store.ensure_current_pricing().expect("reprice"); + let stored = store.task_spans().expect("spans"); + assert_eq!(report.changed_events, 0); + assert_eq!(report.changed_task_spans, 1); + assert!(!report.already_current); + assert_eq!( + stored[0].estimated_cost_usd, + expected.estimated_api_equivalent_usd + ); + assert_eq!( + store.applied_pricing_ruleset_version().expect("applied"), + Some(PRICING_RULESET_VERSION) + ); + } + + #[test] + fn unlinked_task_spans_are_left_unchanged() { + let (store, source) = store_with_source("/tmp/codex-unlinked-span"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + store + .insert_event(&test_event( + &source, + started_at, + "priced", + "codex-auto-review", + million_token_usage(), + expected_review_cost(started_at), + )) + .expect("insert"); + store + .upsert_task_spans(&[test_span( + &source, + started_at, + "unlinked", + Vec::new(), + Some(0), + Some(0), + )]) + .expect("span"); + + let report = store.ensure_current_pricing().expect("reprice"); + let stored = store.task_spans().expect("spans"); + assert_eq!(report.changed_task_spans, 0); + assert_eq!(stored[0].estimated_cost_usd, Some(0)); + assert_eq!(stored[0].estimated_cost_micro_usd, Some(0)); + } + + #[test] + fn unreadable_usage_payloads_are_skipped_without_blocking_repricing() { + let (store, source) = store_with_source("/tmp/codex-corrupt-usage-payload"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store.insert_event(&event).expect("insert"); + store + .conn + .execute( + "INSERT INTO usage_events ( + event_id, provider, source_id, started_at, total_tokens, payload + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", + rusqlite::params![ + "aaa-corrupt-event", + "codex", + source.source_id.0.as_str(), + started_at.to_rfc3339(), + 0, + "{ this is not json", + ], + ) + .expect("corrupt event"); + store + .conn + .execute( + "INSERT INTO usage_summaries ( + summary_id, provider, source_id, observed_at, total_tokens, payload + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", + rusqlite::params![ + "aaa-corrupt-summary", + "codex", + source.source_id.0.as_str(), + started_at.to_rfc3339(), + 0, + "{ also not json", + ], + ) + .expect("corrupt summary"); + + let report = store.ensure_current_pricing().expect("reprice"); + assert_eq!(report.skipped_unreadable_events, 1); + assert_eq!(report.skipped_unreadable_summaries, 1); + assert_eq!(report.changed_events, 1); + assert_eq!(report.refreshed_rollups, 1); + assert_eq!( + stored_event(&store, &event.event_id.0).cost, + expected_review_cost(started_at) + ); + assert_eq!( + store.applied_pricing_ruleset_version().expect("applied"), + Some(PRICING_RULESET_VERSION) + ); + let corrupt = store + .conn + .query_row( + "SELECT payload FROM usage_events WHERE event_id = ?1", + ["aaa-corrupt-event"], + |row| row.get::<_, String>(0), + ) + .expect("corrupt row remains"); + assert_eq!(corrupt, "{ this is not json"); + } + + #[test] + fn dangling_task_span_links_clear_stale_cost() { + let (store, source) = store_with_source("/tmp/codex-dangling-span-link"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "missing-later", + "codex-auto-review", + million_token_usage(), + expected_review_cost(started_at), + ); + store.insert_event(&event).expect("insert"); + store + .upsert_task_spans(&[test_span( + &source, + started_at, + "dangling", + vec![event.event_id.clone()], + Some(99), + Some(990_000), + )]) + .expect("span"); + store + .conn + .execute( + "DELETE FROM usage_events WHERE event_id = ?1", + [&event.event_id.0], + ) + .expect("drop linked event"); + + let report = store.ensure_current_pricing().expect("reprice"); + let stored = store.task_spans().expect("spans"); + assert_eq!(report.changed_task_spans, 1); + assert!(stored[0].estimated_cost_usd.is_none()); + assert!(stored[0].estimated_cost_micro_usd.is_none()); + } + + #[test] + fn summary_inside_one_pricing_window_is_repriced() { + let (store, source) = store_with_source("/tmp/codex-window-summary"); + let start = parse_utc("2026-07-28T00:00:00Z"); + let end = parse_utc("2026-07-29T23:59:59Z"); + store + .upsert_summary(&test_summary( + &source, + "codex-auto-review", + start, + end, + missing_cost(), + )) + .expect("summary"); + + let report = store.ensure_current_pricing().expect("reprice"); + let stored = store + .summaries() + .expect("summaries") + .into_iter() + .next() + .expect("one summary"); + assert_eq!(report.changed_summaries, 1); + assert_eq!(stored.cost, expected_review_cost(end)); + } + + #[test] + fn invalid_applied_ruleset_metadata_fails_closed() { + let (store, source) = store_with_source("/tmp/codex-invalid-ruleset-meta"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store.insert_event(&event).expect("insert"); + store + .set_metadata_value(APPLIED_PRICING_RULESET_VERSION_KEY, "not-a-number") + .expect("write invalid metadata"); + let payload_before = + serde_json::to_string(&stored_event(&store, &event.event_id.0)).expect("serialize"); + + let error = store + .ensure_current_pricing() + .expect_err("invalid metadata must fail"); + + assert!(error + .to_string() + .contains("invalid pricing.applied_ruleset_version")); + assert_eq!( + serde_json::to_string(&stored_event(&store, &event.event_id.0)).expect("serialize"), + payload_before + ); + } + + #[test] + fn repriced_task_buckets_are_marked_dirty_for_incremental_sync() { + let (store, source) = store_with_source("/tmp/codex-task-bucket-dirty"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let event = test_event( + &source, + started_at, + "legacy-review", + "codex-auto-review", + million_token_usage(), + missing_cost(), + ); + store.insert_event(&event).expect("insert"); + store + .upsert_task_spans(&[test_span( + &source, + started_at, + "span", + vec![event.event_id.clone()], + Some(0), + Some(0), + )]) + .expect("span"); + store + .rebuild_task_work_items_for_project_buckets(&BTreeSet::from(["none".to_string()])) + .expect("seed work items"); + let snapshots = store + .pending_task_bucket_snapshots_for_sync( + "http", + "https://example.invalid/api/sync/batches", + "device", + true, + None, + ) + .expect("initial snapshots"); + assert!(!snapshots.is_empty()); + store + .record_task_bucket_snapshots_synced( + "http", + "https://example.invalid/api/sync/batches", + "device", + &snapshots, + ) + .expect("mark synced"); + assert!(store + .pending_task_bucket_snapshots_for_sync( + "http", + "https://example.invalid/api/sync/batches", + "device", + false, + None, + ) + .expect("clean pending") + .is_empty()); + + let report = store.ensure_current_pricing().expect("reprice"); + assert_eq!(report.changed_task_spans, 1); + assert!(report.rebuilt_work_items > 0); + let pending = store + .pending_task_bucket_snapshots_for_sync( + "http", + "https://example.invalid/api/sync/batches", + "device", + false, + None, + ) + .expect("dirty pending"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].project_bucket, "none"); + let work_items = store.work_items().expect("work items"); + assert_eq!( + work_items[0].estimated_cost_usd, + expected_review_cost(started_at).estimated_api_equivalent_usd + ); + } + + #[test] + fn representative_fixture_streams_events_instead_of_loading_the_whole_store() { + let directory = tempfile::tempdir().expect("tempdir"); + let path = directory.path().join("large.sqlite"); + let store = Store::open(&path).expect("open"); + let source = test_source("/tmp/codex-large-reprice"); + store.upsert_source(&source).expect("source"); + let started_at = parse_utc("2026-07-29T12:00:00Z"); + let events = (0..512) + .map(|index| { + test_event( + &source, + started_at + chrono::Duration::seconds(index), + &format!("event-{index}"), + "codex-auto-review", + million_token_usage(), + missing_cost(), + ) + }) + .collect::>(); + store.insert_events(&events).expect("insert"); + + let started = std::time::Instant::now(); + let report = store.ensure_current_pricing().expect("reprice"); + let elapsed = started.elapsed(); + + assert_eq!(report.examined_events, 512); + assert_eq!(report.changed_events, 512); + assert!( + elapsed.as_secs() < 30, + "repricing 512 events should stay well under 30s, took {elapsed:?}" + ); + eprintln!( + "representative fixture: examined={} changed={} elapsed={elapsed:?}", + report.examined_events, report.changed_events + ); + } + + #[test] + fn read_only_applied_version_probe_does_not_open_or_reprice() { + let directory = tempfile::tempdir().expect("tempdir"); + let path = directory.path().join("probe.sqlite"); + assert_eq!( + crate::database_applied_pricing_ruleset_version(&path).expect("missing"), + None + ); + let store = Store::open(&path).expect("create"); + assert_eq!( + crate::database_applied_pricing_ruleset_version(&path).expect("legacy"), + None + ); + store + .set_metadata_value(APPLIED_PRICING_RULESET_VERSION_KEY, "7") + .expect("write"); + drop(store); + assert_eq!( + crate::database_applied_pricing_ruleset_version(&path).expect("present"), + Some(7) + ); + } +} diff --git a/crates/statsai-store/src/snapshot.rs b/crates/statsai-store/src/snapshot.rs index 5974690..c3b162b 100644 --- a/crates/statsai-store/src/snapshot.rs +++ b/crates/statsai-store/src/snapshot.rs @@ -3,7 +3,7 @@ use anyhow::{bail, Context, Result}; #[cfg(target_os = "macos")] use rusqlite::TransactionBehavior; -use rusqlite::{Connection, OpenFlags}; +use rusqlite::{Connection, OpenFlags, OptionalExtension}; #[cfg(target_os = "macos")] use std::fs; #[cfg(target_os = "macos")] @@ -68,6 +68,58 @@ pub fn database_schema_version(path: impl AsRef) -> Result> { .with_context(|| format!("read database schema version at {}", path.display())) } +/// Reads the last applied pricing ruleset without creating, migrating, or +/// repricing the database. +/// +/// `Ok(None)` means the path does not exist, the metadata table is missing, or +/// no applied ruleset has been recorded yet. +/// +/// # Errors +/// +/// Returns an error when the file exists but cannot be opened as a readable +/// SQLite database, or when the stored value is not a valid unsigned integer. +pub fn database_applied_pricing_ruleset_version(path: impl AsRef) -> Result> { + let path = path.as_ref(); + if !path + .try_exists() + .with_context(|| format!("check whether database {} exists", path.display()))? + { + return Ok(None); + } + + let connection = open_read_only(path)?; + let has_metadata_table = connection + .query_row( + "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'local_metadata')", + [], + |row| row.get::<_, bool>(0), + ) + .with_context(|| format!("inspect local metadata at {}", path.display()))?; + if !has_metadata_table { + return Ok(None); + } + + let value: Option = connection + .query_row( + "SELECT value FROM local_metadata WHERE key = ?1", + [super::APPLIED_PRICING_RULESET_VERSION_KEY], + |row| row.get(0), + ) + .optional() + .with_context(|| format!("read applied pricing ruleset version at {}", path.display()))?; + parse_applied_pricing_ruleset_value(value.as_deref()) +} + +pub(crate) fn parse_applied_pricing_ruleset_value(value: Option<&str>) -> Result> { + let Some(value) = value.map(str::trim).filter(|value| !value.is_empty()) else { + return Ok(None); + }; + value + .parse::() + .map(Some) + .with_context(|| format!("invalid pricing.applied_ruleset_version value {value:?}")) +} + /// Replaces `destination` with a SQLite-consistent APFS copy-on-write clone. /// /// The source is opened without running StatsAI migrations. In WAL mode, this diff --git a/crates/statsai-store/src/tasks.rs b/crates/statsai-store/src/tasks.rs index 59e4dcc..ba4aca9 100644 --- a/crates/statsai-store/src/tasks.rs +++ b/crates/statsai-store/src/tasks.rs @@ -958,7 +958,7 @@ impl Store { }) } - fn task_spans_by_sql( + pub(crate) fn task_spans_by_sql( &self, sql: &str, params: &[&dyn rusqlite::types::ToSql], @@ -972,7 +972,7 @@ impl Store { Ok(spans) } - fn upsert_task_spans_in_tx(&self, spans: &[TaskSpan]) -> Result { + pub(crate) fn upsert_task_spans_in_tx(&self, spans: &[TaskSpan]) -> Result { let mut changed = 0u64; let mut span_stmt = self.conn.prepare( r#" diff --git a/crates/statsai/src/lib.rs b/crates/statsai/src/lib.rs index b0be68b..f31a968 100644 --- a/crates/statsai/src/lib.rs +++ b/crates/statsai/src/lib.rs @@ -4,12 +4,29 @@ pub mod privacy_cli; pub mod service; pub mod snapshot; +use anyhow::Result; use chrono::Utc; use getrandom::getrandom; use statsai_core::{hash_text, home_dir}; +use statsai_store::{RepricingReport, Store}; use std::fs::OpenOptions; use std::io::{ErrorKind, Write}; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; + +/// Opens a store for a command that reads or publishes price-derived data and +/// applies the compiled pricing ruleset first. +pub fn open_operational_store(path: &Path) -> Result { + let store = Store::open(path)?; + let report = store.ensure_current_pricing()?; + log_repricing_report(&report); + Ok(store) +} + +pub fn log_repricing_report(report: &RepricingReport) { + if !report.already_current { + eprintln!("{report}"); + } +} pub fn default_store_path() -> PathBuf { home_dir() @@ -257,6 +274,142 @@ mod tests { } } + #[test] + fn open_operational_store_reprices_a_legacy_store() { + use chrono::TimeZone; + use statsai_core::{ + event_id, Confidence, CostInfo, EventSource, LocationOrigin, ModelInfo, PrivacyInfo, + PrivacyMode, SessionInfo, SourceKind, SourceLocation, UsageCounts, UsageEvent, + USAGE_EVENT_SCHEMA_VERSION, + }; + use statsai_store::PRICING_RULESET_VERSION; + use std::path::Path; + + let directory = tempfile::tempdir().expect("tempdir"); + let path = directory.path().join("statsai.sqlite"); + let started_at = Utc + .with_ymd_and_hms(2026, 7, 29, 12, 0, 0) + .single() + .expect("started_at"); + let source = SourceLocation::local_adapter( + "codex", + "test", + "0", + Path::new("/tmp/codex-open-operational-reprice"), + LocationOrigin::Configured, + ); + { + let store = Store::open(&path).expect("create store"); + store.upsert_source(&source).expect("source"); + let event = UsageEvent { + schema_version: USAGE_EVENT_SCHEMA_VERSION.to_string(), + event_id: event_id( + "codex", + &source.source_id, + "legacy-review", + None, + started_at, + ), + device_id: "device".to_string(), + provider: "codex".to_string(), + source_id: source.source_id.clone(), + provider_account_id: None, + subscription_id: None, + source: EventSource { + adapter_id: "test".to_string(), + adapter_version: "0".to_string(), + source_kind: SourceKind::LocalAdapter, + location_origin: Some(LocationOrigin::Configured), + source_type: "jsonl".to_string(), + source_path_hash: source.path_hash.clone(), + source_record_id: Some("legacy-review".to_string()), + parse_confidence: Confidence::High, + }, + session: SessionInfo { + session_id: "session".to_string(), + local_session_id_hash: Some("same-session".to_string()), + title: None, + started_at, + ended_at: None, + duration_seconds: None, + }, + model: Some(ModelInfo { + name: Some("codex-auto-review".to_string()), + normalized_name: Some("codex-auto-review".to_string()), + provider_model_id: Some("codex-auto-review".to_string()), + speed: None, + reasoning_level: None, + reasoning_level_raw: None, + }), + usage: UsageCounts { + input_tokens: Some(1_000_000), + cache_creation_tokens: Some(1_000_000), + cache_read_tokens: Some(1_000_000), + output_tokens: Some(1_000_000), + total_tokens: Some(4_000_000), + ..UsageCounts::default() + }, + runtime: None, + cost: CostInfo { + currency: "USD".to_string(), + estimated_api_equivalent_usd: None, + provider_reported_usd: None, + estimated_api_equivalent_micro_usd: None, + provider_reported_micro_usd: None, + pricing_source: Some("unknown".to_string()), + pricing_version: None, + confidence: Confidence::Low, + }, + parse_evidence: None, + project: None, + git: None, + privacy: PrivacyInfo { + mode: PrivacyMode::MetadataOnly, + contains_prompt_text: false, + contains_response_text: false, + contains_file_paths: false, + }, + created_at: started_at, + imported_at: started_at, + }; + store.insert_event(&event).expect("legacy event"); + assert_eq!(store.applied_pricing_ruleset_version().expect("meta"), None); + drop(store); + } + + let store = open_operational_store(&path).expect("operational open"); + assert_eq!( + store.applied_pricing_ruleset_version().expect("applied"), + Some(PRICING_RULESET_VERSION) + ); + let stored = store.events().expect("events"); + assert_eq!(stored.len(), 1); + assert!(stored[0].cost.estimated_api_equivalent_usd.is_some()); + } + + #[test] + fn open_operational_store_refuses_a_newer_pricing_ruleset() { + let directory = tempfile::tempdir().expect("tempdir"); + let path = directory.path().join("statsai.sqlite"); + let store = Store::open(&path).expect("create store"); + store + .set_metadata_value(statsai_store::APPLIED_PRICING_RULESET_VERSION_KEY, "99") + .expect("future ruleset"); + drop(store); + + let error = match open_operational_store(&path) { + Ok(_) => panic!("forward pricing must refuse"), + Err(error) => error, + }; + assert!(error + .to_string() + .contains("pricing ruleset version 99 is newer than this StatsAI binary supports")); + assert_eq!( + statsai_store::database_applied_pricing_ruleset_version(&path).expect("unchanged"), + Some(99) + ); + } + #[test] fn daemon_auth_token_is_random_persistent_and_private() { let temp = tempfile::tempdir().expect("temp dir"); diff --git a/crates/statsai/src/main.rs b/crates/statsai/src/main.rs index 5c67088..6827555 100644 --- a/crates/statsai/src/main.rs +++ b/crates/statsai/src/main.rs @@ -38,14 +38,15 @@ use statsai_sdk::{ build_reported_usage_summary, ReportedUsageSummaryInput, ReportedUsageSummaryRecord, REPORTED_USAGE_IMPORT_ADAPTER_ID, }; -#[cfg(test)] -use statsai_store::{apply_verified_source_state, verified_source_state_hash}; use statsai_store::{ - close_active_verified_source_linkages, derive_task_work_items, find_existing_provider_account, - reconcile_verified_source_state, upsert_provider_account, verified_source_observation_hash, - QuotaQuery, ScanFileStateEntry, Store, SyncPreferences, SyncState, TaskRebuildReport, - UpsertProviderAccountInput, CURRENT_SCHEMA_VERSION, + apply_current_estimated_pricing, close_active_verified_source_linkages, derive_task_work_items, + find_existing_provider_account, reconcile_verified_source_state, upsert_provider_account, + verified_source_observation_hash, QuotaQuery, ScanFileStateEntry, Store, SyncPreferences, + SyncState, TaskRebuildReport, UpsertProviderAccountInput, CURRENT_SCHEMA_VERSION, + PRICING_RULESET_VERSION, }; +#[cfg(test)] +use statsai_store::{apply_verified_source_state, verified_source_state_hash}; use statsai_sync::{ validate_authenticated_http_endpoint, FileSink, HttpSink, StdoutSink, SyncSink, }; @@ -932,6 +933,8 @@ enum StoreAdminSubcommand { }, #[command(about = "Print the store schema version supported by this binary")] SupportedSchemaVersion, + #[command(about = "Print the pricing ruleset version supported by this binary")] + SupportedPricingRulesetVersion, } #[derive(Debug, Subcommand)] @@ -973,7 +976,11 @@ fn main() -> Result<()> { Command::Service(command) => service(command), Command::Snapshot(command) => snapshot::run(command, &store_path, &device_id), command => { - let store = Store::open(&store_path)?; + let store = if command_reprices_persisted_usage(&command) { + statsai::open_operational_store(&store_path)? + } else { + Store::open(&store_path)? + }; match command { Command::Scan(command) => scan(command, &store, &device_id), Command::Report(command) => report(command, &store), @@ -1004,6 +1011,19 @@ fn main() -> Result<()> { } } +fn command_reprices_persisted_usage(command: &Command) -> bool { + matches!( + command, + Command::Scan(_) + | Command::Report(_) + | Command::Import(_) + | Command::Export(_) + | Command::Task(_) + | Command::Sync(_) + | Command::Daemon(_) + ) +} + fn store_admin(command: StoreAdminCommand, source: &Path) -> Result<()> { match command.command { StoreAdminSubcommand::CloneTo { destination } => { @@ -1021,6 +1041,10 @@ fn store_admin(command: StoreAdminCommand, source: &Path) -> Result<()> { println!("{CURRENT_SCHEMA_VERSION}"); Ok(()) } + StoreAdminSubcommand::SupportedPricingRulesetVersion => { + println!("{PRICING_RULESET_VERSION}"); + Ok(()) + } } } @@ -3808,7 +3832,11 @@ fn build_reported_import_record( device_id: &str, ) -> Result { let legacy_replacement_source_ids = legacy_alias_replacement_source_ids(&input); - let record = build_reported_usage_summary(input, device_id)?; + let mut record = build_reported_usage_summary(input, device_id)?; + // Overlay before upsert so a store already at this ruleset cannot keep an + // imported estimated-only figure from an older catalog. A later + // ensure_current_pricing pass would no-op. + record.summary = apply_current_estimated_pricing(record.summary); Ok(ReportedImportRecord { record, legacy_replacement_source_ids, @@ -11813,6 +11841,115 @@ mod tests { assert_eq!(summaries[0].metadata.summary_format, "manual_daily"); } + #[test] + fn import_overlays_estimated_cost_when_store_ruleset_is_already_current() { + let store = Store::in_memory().expect("store"); + store + .ensure_current_pricing() + .expect("apply current ruleset"); + let noop = store.ensure_current_pricing().expect("already current"); + assert!(noop.already_current); + + let period_end = Utc + .with_ymd_and_hms(2026, 7, 29, 23, 59, 59) + .single() + .expect("end"); + let input = ReportedUsageSummaryInput { + schema_version: "reported_usage_summary_input.v1".to_string(), + provider: "codex".to_string(), + provider_account_id: None, + provider_user_id: None, + email: None, + account_label: None, + source_kind: SourceKind::ExternalReport, + source_name: "legacy_catalog_export".to_string(), + evidence_id: Some("legacy-catalog:2026-07-29".to_string()), + evidence_path: None, + report_format: "manual_period_summary".to_string(), + report_version: Some("manual.v1".to_string()), + period_start: Some( + Utc.with_ymd_and_hms(2026, 7, 29, 0, 0, 0) + .single() + .expect("start"), + ), + period_end: Some(period_end), + observed_at: None, + model: Some(ModelInfo { + name: Some("codex-auto-review".to_string()), + normalized_name: Some("codex-auto-review".to_string()), + provider_model_id: Some("codex-auto-review".to_string()), + speed: None, + reasoning_level: None, + reasoning_level_raw: None, + }), + usage: UsageCounts { + input_tokens: Some(1_000_000), + cache_creation_tokens: Some(1_000_000), + cache_read_tokens: Some(1_000_000), + output_tokens: Some(1_000_000), + total_tokens: Some(4_000_000), + ..UsageCounts::default() + }, + cost: Some(CostInfo { + currency: "USD".to_string(), + estimated_api_equivalent_usd: Some(1), + provider_reported_usd: None, + estimated_api_equivalent_micro_usd: Some(10_000), + provider_reported_micro_usd: None, + pricing_source: Some("official:stale".to_string()), + pricing_version: Some("official:stale".to_string()), + confidence: Confidence::Low, + }), + confidence: Some(Confidence::Medium), + }; + + let incoming = build_reported_import_record(input, "device").expect("incoming"); + assert_ne!( + incoming.record.summary.cost.estimated_api_equivalent_usd, + Some(1) + ); + assert_eq!( + incoming.record.summary.cost.pricing_version.as_deref(), + Some(statsai_store::PRICING_CATALOG_VERSION) + ); + assert!(incoming + .record + .summary + .cost + .estimated_api_equivalent_usd + .is_some()); + let expected_usd = incoming.record.summary.cost.estimated_api_equivalent_usd; + + import_reported_summary_records( + &store, + &[ReportedImportReport { + path: PathBuf::from("legacy-catalog.json"), + records: vec![incoming], + warnings: Vec::new(), + }], + false, + false, + false, + ) + .expect("import"); + + let stored = store.summaries().expect("summaries"); + assert_eq!(stored.len(), 1); + assert_eq!(stored[0].cost.estimated_api_equivalent_usd, expected_usd); + assert_eq!( + stored[0].cost.pricing_version.as_deref(), + Some(statsai_store::PRICING_CATALOG_VERSION) + ); + let still_current = store.ensure_current_pricing().expect("still current"); + assert!(still_current.already_current); + assert_eq!( + store.summaries().expect("unchanged")[0] + .cost + .estimated_api_equivalent_usd, + expected_usd + ); + } + #[test] fn configured_claude_projects_path_normalizes_to_config_root() { let dir = tempfile::tempdir().expect("tempdir"); @@ -15162,6 +15299,179 @@ mod tests { assert!(is_daily_rollup_summary(&batch.summaries[0])); } + #[test] + fn incremental_http_sync_includes_repriced_rollups_without_full() { + let store = Store::in_memory().expect("store"); + let source = SourceLocation::local_adapter( + "codex", + "test", + "0", + Path::new("/tmp/codex-http-reprice-sync"), + LocationOrigin::Configured, + ); + store.upsert_source(&source).expect("source"); + + let started_at = Utc + .with_ymd_and_hms(2026, 7, 29, 12, 0, 0) + .single() + .expect("started_at"); + let mut event = test_event( + "codex", + &source, + started_at, + None, + TokenParts { + input: 1_000_000, + cached_input: 1_000_000, + output: 1_000_000, + reasoning: 0, + total: 4_000_000, + cost: None, + }, + ); + event.model = Some(ModelInfo { + name: Some("codex-auto-review".to_string()), + normalized_name: Some("codex-auto-review".to_string()), + provider_model_id: Some("codex-auto-review".to_string()), + speed: None, + reasoning_level: None, + reasoning_level_raw: None, + }); + store.insert_event(&event).expect("legacy unpriced event"); + store.rebuild_sync_rollups().expect("rebuild"); + + let command = SyncCommand { + endpoint: Some("https://api.example.com/api/sync/batches".to_string()), + ..test_sync_command("http") + }; + assert!(!command.full); + assert!(!command.rebuild_rollups); + let target = sync_target(&command).expect("target"); + let (initial_batch, initial_mode) = + build_sync_batch(&command, &store, "device", &target).expect("initial batch"); + assert_eq!(initial_mode, SyncPayloadMode::Rollups); + assert_eq!(initial_batch.summaries.len(), 1); + assert!(initial_batch.summaries[0] + .cost + .estimated_api_equivalent_usd + .is_none()); + record_rollup_sync_success(&store, "http", &target, &initial_batch) + .expect("record initial sync"); + + let (repeat_batch, _) = + build_sync_batch(&command, &store, "device", &target).expect("repeat batch"); + assert!( + repeat_batch.summaries.is_empty(), + "synced rollups must stay unpublished until pricing changes them" + ); + + let report = store.ensure_current_pricing().expect("automatic reprice"); + assert_eq!(report.changed_events, 1); + assert_eq!(report.refreshed_rollups, 1); + + let (incremental_batch, incremental_mode) = + build_sync_batch(&command, &store, "device", &target).expect("incremental batch"); + assert_eq!(incremental_mode, SyncPayloadMode::Rollups); + assert_eq!(incremental_batch.summaries.len(), 1); + assert!(is_daily_rollup_summary(&incremental_batch.summaries[0])); + assert_eq!( + incremental_batch.summaries[0] + .cost + .estimated_api_equivalent_usd, + store + .events() + .expect("repriced events") + .into_iter() + .next() + .expect("one event") + .cost + .estimated_api_equivalent_usd + ); + assert!(incremental_batch.summaries[0] + .cost + .estimated_api_equivalent_usd + .is_some()); + } + + #[test] + fn incremental_http_sync_includes_repriced_passthrough_summaries_without_full() { + let store = Store::in_memory().expect("store"); + let source = SourceLocation::local_adapter( + "codex", + "test", + "0", + Path::new("/tmp/codex-http-reprice-passthrough"), + LocationOrigin::Configured, + ); + store.upsert_source(&source).expect("source"); + + let start = Utc + .with_ymd_and_hms(2026, 7, 28, 0, 0, 0) + .single() + .expect("start"); + let end = Utc + .with_ymd_and_hms(2026, 7, 29, 23, 59, 59) + .single() + .expect("end"); + let mut summary = test_summary("codex", &source, end, 4_000_000, None); + summary.source.source_kind = SourceKind::LocalAdapter; + summary.source.source_type = "build-session.json".to_string(); + summary.metadata.summary_format = "grok_build_session_summary".to_string(); + summary.period_start = Some(start); + summary.period_end = Some(end); + summary.model = Some(ModelInfo { + name: Some("codex-auto-review".to_string()), + normalized_name: Some("codex-auto-review".to_string()), + provider_model_id: Some("codex-auto-review".to_string()), + speed: None, + reasoning_level: None, + reasoning_level_raw: None, + }); + summary.usage = UsageCounts { + input_tokens: Some(1_000_000), + cache_creation_tokens: Some(1_000_000), + cache_read_tokens: Some(1_000_000), + output_tokens: Some(1_000_000), + total_tokens: Some(4_000_000), + ..UsageCounts::default() + }; + store.upsert_summary(&summary).expect("passthrough summary"); + + let command = SyncCommand { + endpoint: Some("https://api.example.com/api/sync/batches".to_string()), + ..test_sync_command("http") + }; + assert!(!command.full); + let target = sync_target(&command).expect("target"); + let (initial_batch, _) = + build_sync_batch(&command, &store, "device", &target).expect("initial batch"); + assert_eq!(initial_batch.summaries.len(), 1); + assert!(!is_daily_rollup_summary(&initial_batch.summaries[0])); + assert!(initial_batch.summaries[0] + .cost + .estimated_api_equivalent_usd + .is_none()); + record_rollup_sync_success(&store, "http", &target, &initial_batch) + .expect("record initial sync"); + + let (repeat_batch, _) = + build_sync_batch(&command, &store, "device", &target).expect("repeat batch"); + assert!(repeat_batch.summaries.is_empty()); + + let report = store.ensure_current_pricing().expect("automatic reprice"); + assert_eq!(report.changed_summaries, 1); + + let (incremental_batch, incremental_mode) = + build_sync_batch(&command, &store, "device", &target).expect("incremental batch"); + assert_eq!(incremental_mode, SyncPayloadMode::Rollups); + assert_eq!(incremental_batch.summaries.len(), 1); + assert!(!is_daily_rollup_summary(&incremental_batch.summaries[0])); + assert!(incremental_batch.summaries[0] + .cost + .estimated_api_equivalent_usd + .is_some()); + } + #[test] fn http_incremental_rollups_are_tracked_per_target() { let store = Store::in_memory().expect("store"); @@ -21601,6 +21911,50 @@ mod tests { assert!(!store_path.exists()); } + #[test] + fn supported_pricing_ruleset_version_is_available_without_opening_a_store() { + let cli = Cli::try_parse_from(["statsai", "store", "supported-pricing-ruleset-version"]) + .expect("parse supported pricing query"); + + assert!(matches!( + cli.command, + Command::Store(StoreAdminCommand { + command: StoreAdminSubcommand::SupportedPricingRulesetVersion, + }) + )); + + let directory = tempfile::tempdir().expect("temporary directory"); + let store_path = directory.path().join("must-not-be-created.sqlite"); + store_admin( + StoreAdminCommand { + command: StoreAdminSubcommand::SupportedPricingRulesetVersion, + }, + &store_path, + ) + .expect("print supported pricing ruleset"); + assert!(!store_path.exists()); + } + + #[test] + fn price_derived_commands_reprice_and_diagnostic_commands_do_not() { + let reprice = |args: &[&str]| { + let cli = Cli::try_parse_from(args).expect("parse"); + command_reprices_persisted_usage(&cli.command) + }; + assert!(reprice(&["statsai", "scan"])); + assert!(reprice(&["statsai", "report", "monthly"])); + assert!(reprice(&["statsai", "sync"])); + assert!(reprice(&["statsai", "export", "--json"])); + assert!(reprice(&["statsai", "task", "list"])); + assert!(!reprice(&["statsai", "status"])); + assert!(!reprice(&["statsai", "quota", "status"])); + assert!(!reprice(&["statsai", "conversation", "list"])); + assert!(!reprice(&["statsai", "account", "list"])); + assert!(!reprice(&["statsai", "source", "list"])); + let doctor = Cli::try_parse_from(["statsai", "doctor"]).expect("parse doctor"); + assert!(!command_reprices_persisted_usage(&doctor.command)); + } + #[test] fn report_range_cli_parses_from_and_to() { let cli = Cli::try_parse_from([ diff --git a/crates/statsai/src/snapshot.rs b/crates/statsai/src/snapshot.rs index c909110..de60982 100644 --- a/crates/statsai/src/snapshot.rs +++ b/crates/statsai/src/snapshot.rs @@ -88,7 +88,7 @@ pub struct SnapshotSourceStatus { } pub fn run(command: SnapshotCommand, store_path: &Path, device_id: &str) -> Result<()> { - let store = Store::open(store_path)?; + let store = crate::open_operational_store(store_path)?; let snapshot = collect_for_device(&store, device_id)?; if command.json { diff --git a/docs/statsai-dev.md b/docs/statsai-dev.md index eb683dd..215470c 100644 --- a/docs/statsai-dev.md +++ b/docs/statsai-dev.md @@ -117,7 +117,8 @@ Each download is checked before selection: 4. `build.json.repository` is `starkdmi/statsai`; 5. `build.json.sha` is the exact resolved SHA; 6. workflow run ID and attempt match the downloaded run; -7. the supported store schema version is recorded in the manifest; +7. the supported store schema version and pricing ruleset version are recorded + in the manifest (`build.json` schema 2); 8. the target and Mach-O header are ARM64 macOS. Selection is an atomic state-file replacement. The cache retains the current @@ -185,10 +186,19 @@ statsai-dev --prod-data report monthly This choice is never persisted. -The escape hatch is allowed only when the production database schema exactly -matches the schema supported by the selected build. A development build is -never allowed to migrate production data; use the isolated development clone -to test a schema-changing PR. +The escape hatch is allowed only when the production database schema **and** +applied pricing ruleset exactly match the versions supported by the selected +build. Missing, older, or newer production pricing metadata is refused. A +development build is never allowed to migrate or reprice production data; use +the isolated development clone to test a schema-changing or pricing-changing +PR. + +Ordinary isolated `statsai-dev` stores are opened by the selected exact-SHA +`statsai` binary. That binary applies its own pricing ruleset automatically +before scan, report, sync, snapshot, and other price-derived commands. Status, +doctor, quota, and conversation do not trigger a reprice. A pricing catalog +change does not require a raw rescan; a later incremental sync publishes +corrected dirty rollups. Development daemon commands and mutating service commands are blocked: diff --git a/docs/sync-contract.md b/docs/sync-contract.md index d57c74f..e18ea23 100644 --- a/docs/sync-contract.md +++ b/docs/sync-contract.md @@ -228,6 +228,12 @@ Cost payloads may include `provider_reported_micro_usd`, `estimated_cost_micro_usd`. Receivers prefer these integer micro-USD values and fall back to legacy rounded-cent fields when they are absent. +When a local pricing ruleset advances, StatsAI reprices persisted normalized +events and refreshes the affected `sync_rollups` in place. Changed rollups are +marked dirty so a later incremental HTTP sync publishes the corrected values +without `--rebuild-rollups` or `--full`. Provider-reported cost fields are +never replaced by estimates. + User-defined aliases are still retained in `ProviderAccount.account_label` for display, but they are no longer the primary account key.