Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions rust/src/cli/serve/dashboard/source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ use std::collections::{BTreeSet, HashMap};
use std::pin::Pin;
use std::time::Duration;

use chrono::{Local, Utc};
use chrono::Utc;

use crate::core::{CostScanOptions, FetchContext, ProviderId, SourceMode, instantiate_provider};
use crate::cost_scanner::{self, CostScanner};
Expand Down Expand Up @@ -218,7 +218,10 @@ async fn collect_costs(pi_selected: bool) -> HashMap<String, RawCostPayload> {
let pi = scanner.scan_pi_with_cancel(None);
let pi_contract =
build_local_spend_contract_from_summary("pi", 30, false, false, false, pi);
let today = Local::now().date_naive().format("%Y-%m-%d").to_string();
let today = crate::cost_reporting_period::cost_bucket_zone()
.date(Utc::now())
.format("%Y-%m-%d")
.to_string();
let today_of = |provider: &str| {
cost_scanner::get_daily_cost_history(provider, 30)
.into_iter()
Expand Down
9 changes: 5 additions & 4 deletions rust/src/codex_costs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,14 +15,15 @@ pub(crate) use summary_contract::{
decode_remote_codex_summary,
};

use chrono::{Duration, Local, NaiveDate, Utc};
use chrono::{Duration, NaiveDate, Utc};
use std::collections::HashSet;
use std::path::Path;

use crate::core::{
CodexUsageRecord, CostUsageCache, CostUsageDayRange, CostUsagePricing, JsonlScanner,
is_unpriced_codex_routing_model,
};
use crate::cost_reporting_period::cost_bucket_zone;
use crate::cost_scanner::{CostSummary, ModelPricingCompleteness, ModelTokenCounts};
use crate::spend_contract::CostCoverageCounts;

Expand All @@ -40,12 +41,12 @@ pub(crate) fn build_codex_cost_summary(
&today,
history_days,
Utc::now(),
crate::core::local_timezone_name(),
cost_bucket_zone().identifier(),
)
}

fn codex_today_summary(history: &CostSummary, cache: &CostUsageCache) -> CostSummary {
let today = Local::now().date_naive();
let today = cost_bucket_zone().date(Utc::now());
let range = CostUsageDayRange::new(today, today);
let mut summary = CostSummary {
period_start: Some(today),
Expand Down Expand Up @@ -228,7 +229,7 @@ pub(crate) fn scan_codex_file_cost_for_range(path: &Path, range: &CostUsageDayRa

#[cfg(test)]
pub(crate) fn scan_codex_file_cost(path: &Path) -> f64 {
let today = Local::now().date_naive();
let today = chrono::Local::now().date_naive();
let range = CostUsageDayRange::new(codex_period_start(today, 30), today);
scan_codex_file_cost_for_range(path, &range)
}
Expand Down
26 changes: 11 additions & 15 deletions rust/src/codex_costs/quota_windows.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,14 @@
//! window describes the account now; these rows describe local historical
//! evidence for a future display/transport surface.

use chrono::{DateTime, Duration, Local, NaiveDate, TimeZone, Utc};
use chrono::{DateTime, Duration, NaiveDate, Utc};
use serde::{Deserialize, Serialize};
use std::collections::HashSet;

use crate::core::{
CodexSourceRowCache, CodexSourceUsageRow, CostUsageCache, CostUsagePricing, RateWindow,
};
use crate::cost_reporting_period::cost_bucket_zone;

const NOMINAL_WEEK_MINUTES: i64 = 7 * 24 * 60;
const RESET_TOLERANCE_SECONDS: i64 = 2 * 60;
Expand Down Expand Up @@ -146,7 +147,7 @@ pub fn codex_quota_windows_from_cache(
let history_start = cache
.scan_since_key
.as_deref()
.and_then(local_day_start)
.and_then(bucket_day_start)
.unwrap_or_else(|| {
let count = i32::try_from(count).expect("quota window count is capped");
current_end - duration * (count + 1)
Expand Down Expand Up @@ -359,9 +360,9 @@ fn cache_slices(cache: &CostUsageCache) -> Vec<Slice> {
}

fn slice_from_row(row: &CodexSourceUsageRow) -> Slice {
let timestamp = row.timestamp.or_else(|| local_day_start(&row.day_key));
let timestamp = row.timestamp.or_else(|| bucket_day_start(&row.day_key));
let end = row.timestamp.map(|_| None).unwrap_or_else(|| {
local_day_start(&row.day_key).and_then(|start| start.checked_add_signed(Duration::days(1)))
bucket_day_start(&row.day_key).and_then(|start| start.checked_add_signed(Duration::days(1)))
});
let input = u64::try_from(row.input.max(0)).unwrap_or(0);
let output = u64::try_from(row.output.max(0)).unwrap_or(0);
Expand All @@ -374,7 +375,7 @@ fn slice_from_row(row: &CodexSourceUsageRow) -> Slice {
} else {
model.to_string()
};
let date = timestamp.map(|value| value.with_timezone(&Local).date_naive())?;
let date = timestamp.map(|value| cost_bucket_zone().date(value))?;
CostUsagePricing::codex_cost_usd_at_date(
&model,
input,
Expand All @@ -398,7 +399,7 @@ fn legacy_day_slices(cache: &CostUsageCache) -> Vec<Slice> {
let mut days: Vec<_> = cache.days.iter().collect();
days.sort_by_key(|(day, _)| *day);
for (day, models) in days {
let Some(start) = local_day_start(day) else {
let Some(start) = bucket_day_start(day) else {
continue;
};
let Some(end) = start.checked_add_signed(Duration::days(1)) else {
Expand All @@ -425,7 +426,7 @@ fn legacy_day_slices(cache: &CostUsageCache) -> Vec<Slice> {
output,
NaiveDate::parse_from_str(day, "%Y-%m-%d")
.ok()
.unwrap_or_else(|| start.with_timezone(&Local).date_naive()),
.unwrap_or_else(|| cost_bucket_zone().date(start)),
);
slices.push(Slice {
start,
Expand All @@ -440,15 +441,10 @@ fn legacy_day_slices(cache: &CostUsageCache) -> Vec<Slice> {
slices
}

fn local_day_start(day: &str) -> Option<DateTime<Utc>> {
/// The instant a `YYYY-MM-DD` day key begins in the pinned bucket zone.
fn bucket_day_start(day: &str) -> Option<DateTime<Utc>> {
let date = NaiveDate::parse_from_str(day, "%Y-%m-%d").ok()?;
let naive = date.and_hms_opt(0, 0, 0)?;
Local
.from_local_datetime(&naive)
.single()
.or_else(|| Local.from_local_datetime(&naive).earliest())
.or_else(|| Local.from_local_datetime(&naive).latest())
.map(|value| value.with_timezone(&Utc))
Some(cost_bucket_zone().start_of_day_utc(date))
}

#[cfg(test)]
Expand Down
5 changes: 3 additions & 2 deletions rust/src/codex_workspaces/indexer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ use std::collections::HashMap;
use std::fs;
use std::path::{Path, PathBuf};

use chrono::{DateTime, Local, NaiveDate, TimeZone, Utc};
use chrono::{DateTime, NaiveDate, TimeZone, Utc};
use rusqlite::{Connection, OpenFlags};

use crate::agent_sessions::CodexRolloutFirstLineParser;
Expand Down Expand Up @@ -156,7 +156,7 @@ impl CodexWorkspacesIndex {
}

progress(Progress::phase(ProgressPhase::ScanningLogs));
let today = Local::now().date_naive();
let today = crate::cost_reporting_period::cost_bucket_zone().date(Utc::now());
let since = codex_period_start(today, self.history_days);
let range = CostUsageDayRange::new(since, today);

Expand Down Expand Up @@ -745,6 +745,7 @@ fn meta_day_out_of_range(path: &Path, range: &CostUsageDayRange) -> bool {
mod tests {
use super::*;
use crate::codex_workspaces::types::SourceStatus;
use chrono::Local;
use rusqlite::Connection;
use std::fs::File;
use std::io::Write;
Expand Down
10 changes: 9 additions & 1 deletion rust/src/core/jsonl_scanner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,8 @@ struct CachedCostReadStatusProjection {
previous_report: Option<CachedCostReport>,
#[serde(default)]
codex_scan_pause_reason: Option<CodexScanPauseReason>,
#[serde(default)]
bucket_time_zone: Option<String>,
}

fn deserialize_nonempty_object<'de, D>(deserializer: D) -> Result<bool, D::Error>
Expand Down Expand Up @@ -239,6 +241,11 @@ pub struct CostUsageCache {
/// caches remain valid and can be upgraded lazily.
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub codex_source_rows: HashMap<String, CodexSourceRowCache>,
/// Zone the day keys were bucketed in (upstream `timeZoneIdentifier`).
/// A cache from another zone is rebuilt; caches written before the stamp
/// existed are kept.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub bucket_time_zone: Option<String>,
/// Content stamp of the decoded on-disk baseline. This is process-local
/// and omitted from JSON so a stale reader cannot replace a newer cache.
#[serde(skip)]
Expand Down Expand Up @@ -592,7 +599,8 @@ impl JsonlScanner {
return CachedCostReadStatus::default();
};
if provider == ProviderId::Codex
&& !codex::codex_cache_schema_is_current(projection.codex_cache_schema_version)
&& (!codex::codex_cache_schema_is_current(projection.codex_cache_schema_version)
|| !codex::codex_cache_zone_is_current(projection.bucket_time_zone.as_deref()))
{
return CachedCostReadStatus::default();
}
Expand Down
19 changes: 16 additions & 3 deletions rust/src/core/jsonl_scanner/codex.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ use helpers::{
};
use parser::CodexParserState;

use crate::cost_reporting_period::cost_bucket_zone;

/// Persisted Codex cache schema version. Version 0 predates 64-bit totals;
/// version 1 can retain a terminal pause after treating a paginated v2
/// subagent's independent counters as an inherited fork. Version 3 adds
Expand All @@ -23,7 +25,14 @@ pub(crate) fn codex_cache_schema_is_current(schema_version: u32) -> bool {
schema_version == CODEX_CACHE_SCHEMA_VERSION
}

/// Apply the Codex cache schema version policy to a freshly decoded artifact.
/// Whether a persisted Codex cache bucketed its days in the zone now in
/// effect. Artifacts written before the zone was stamped are kept.
pub(crate) fn codex_cache_zone_is_current(bucket_time_zone: Option<&str>) -> bool {
bucket_time_zone.is_none_or(|zone| zone == cost_bucket_zone().identifier())
}

/// Apply the Codex cache schema version and bucket zone policy to a freshly
/// decoded artifact.
///
/// A mismatched artifact is invalidated: a fresh, current-version cache is
/// returned with the decoded baseline stamp retained so the caller stays
Expand All @@ -33,7 +42,9 @@ pub(crate) fn codex_cache_apply_load_policy(
mut cache: CostUsageCache,
stamp: CacheStamp,
) -> CostUsageCache {
if !codex_cache_schema_is_current(cache.codex_cache_schema_version) {
if !codex_cache_schema_is_current(cache.codex_cache_schema_version)
|| !codex_cache_zone_is_current(cache.bucket_time_zone.as_deref())
{
return CostUsageCache {
codex_cache_schema_version: CODEX_CACHE_SCHEMA_VERSION,
loaded_stamp: Some(Some(stamp)),
Expand All @@ -44,9 +55,11 @@ pub(crate) fn codex_cache_apply_load_policy(
cache
}

/// Stamp the current schema version before a Codex cache is persisted.
/// Stamp the current schema version and bucket zone before a Codex cache is
/// persisted.
pub(crate) fn codex_cache_stamp_schema_version(cache: &mut CostUsageCache) {
cache.codex_cache_schema_version = CODEX_CACHE_SCHEMA_VERSION;
cache.bucket_time_zone = Some(cost_bucket_zone().identifier());
}

#[cfg(test)]
Expand Down
7 changes: 3 additions & 4 deletions rust/src/core/jsonl_scanner/codex/helpers.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
use super::CodexTotals;
use chrono::{DateTime, FixedOffset, Local, NaiveDate, TimeZone};
use chrono::{DateTime, FixedOffset, NaiveDate, TimeZone, Utc};
use serde::Deserialize;
use serde_json::Value;
use std::io::BufRead;
Expand Down Expand Up @@ -346,9 +346,8 @@ impl ParsedCodexTimestamp {
self.parsed
.as_ref()
.map(|timestamp| {
timestamp
.with_timezone(&Local)
.date_naive()
crate::cost_reporting_period::cost_bucket_zone()
.date(timestamp.with_timezone(&Utc))
.format("%Y-%m-%d")
.to_string()
})
Expand Down
13 changes: 10 additions & 3 deletions rust/src/core/jsonl_scanner/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1603,8 +1603,10 @@ fn save_cache_at_exact_limit_is_accepted() {
let root = tempfile::tempdir().unwrap();
let cache_root = root.path().to_path_buf();

let cache = CostUsageCache::default();
// Serialize to learn the actual encoded size for this exact struct.
let mut cache = CostUsageCache::default();
// Saving stamps the schema version and bucket zone first; serialize the
// stamped struct to learn the exact encoded size.
codex_cache_stamp_schema_version(&mut cache);
let json = serde_json::to_string(&cache).unwrap();
let exact_limit = json.len();

Expand All @@ -1630,7 +1632,8 @@ fn save_cache_one_over_limit_is_refused_and_removes_destination() {
let root = tempfile::tempdir().unwrap();
let cache_root = root.path().to_path_buf();

let cache = CostUsageCache::default();
let mut cache = CostUsageCache::default();
codex_cache_stamp_schema_version(&mut cache);
let json = serde_json::to_string(&cache).unwrap();
// One byte short of the encoded size forces refusal on the next attempt.
let under_by_one = json.len().saturating_sub(1);
Expand All @@ -1653,3 +1656,7 @@ fn save_cache_one_over_limit_is_refused_and_removes_destination() {
#[cfg(test)]
#[path = "tests/codex_metadata.rs"]
mod codex_metadata;

#[cfg(test)]
#[path = "tests/cache_zone.rs"]
mod cache_zone;
Loading