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
40 changes: 39 additions & 1 deletion rust/src/providers/chart.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ use crate::providers::claude::quota_history::{
ClaudeQuotaHistoryOptions, ClaudeQuotaResetObservation, aggregate_claude_quota_windows,
};
use crate::providers::claude::reset_observations;
use crate::providers::codex::reset_observations as codex_reset_observations;
use chrono::{DateTime, Utc};
use std::sync::atomic::AtomicBool;

Expand Down Expand Up @@ -156,13 +157,50 @@ fn load_and_persist_reset_observations(
}
}

/// Load the persisted Codex reset observations and record the live window's
/// reset on this refresh (upstream 0.62.0 #3358). Read failures never break
/// the chart: an error degrades to the previously stored observations.
fn load_and_persist_codex_reset_observations(
live_window: &RateWindow,
now: DateTime<Utc>,
) -> Vec<DateTime<Utc>> {
let Some(reset_at) = live_window.resets_at else {
return Vec::new();
};
let config_root = match dirs::config_dir().map(|root| root.join("CodexBar")) {
Some(root) => root,
None => return Vec::new(),
};
let scope = codex_reset_observations::CODEX_ACCOUNT_SCOPE;
match codex_reset_observations::merge_and_persist_reset_observation(
&config_root,
scope,
reset_at,
now,
) {
Ok(result) => result
.observations
.iter()
.map(|observation| observation.resets_at)
.collect(),
Err(_) => codex_reset_observations::load_reset_observations(&config_root, scope)
.unwrap_or_default()
.iter()
.map(|observation| observation.resets_at)
.collect(),
}
}

fn build_codex_quota_history(
account_scope: Option<&str>,
live_window: Option<&RateWindow>,
) -> Option<QuotaWindowHistorySnapshot> {
let live_window = live_window?;
let cache = JsonlScanner::load_cache(ProviderId::Codex, None);
let windows = codex_quota_windows_from_cache(&cache, Some(live_window), &[], Utc::now(), 4);
let now = Utc::now();
let observed_next_resets = load_and_persist_codex_reset_observations(live_window, now);
let windows =
codex_quota_windows_from_cache(&cache, Some(live_window), &observed_next_resets, now, 4);
if windows.is_empty() {
return None;
}
Expand Down
1 change: 1 addition & 0 deletions rust/src/providers/codex/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

mod api;
mod pat;
pub mod reset_observations;
mod weekly_reset;

use async_trait::async_trait;
Expand Down
324 changes: 324 additions & 0 deletions rust/src/providers/codex/reset_observations.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,324 @@
//! Durable Codex quota reset observations (upstream 0.62.0 #3358).
//!
//! Each refresh of the Codex chart observes the live weekly window's
//! `resets_at`. Persisting those observations lets
//! [`crate::codex_costs::quota_windows::codex_quota_windows_from_cache`]
//! recover real window boundaries across restarts instead of estimating every
//! boundary from the live window alone. Modeled on the Claude
//! `reset_observations` store; Codex history is not account-partitioned, so
//! observations live under a single well-known scope.

use crate::{atomic_file, secure_file};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use thiserror::Error;
use uuid::Uuid;

const STORE_VERSION: u32 = 1;
const STORE_RELATIVE_PATH: &str = "codex/quota-reset-observations-v1.json";
const MAX_OBSERVATIONS_PER_ACCOUNT: usize = 512;

/// Every Codex chart observation shares this scope (no account partitioning).
pub const CODEX_ACCOUNT_SCOPE: &str = "codex";

#[derive(Debug, Error)]
pub enum CodexResetObservationError {
#[error("Codex reset observation account scope is empty")]
EmptyAccountScope,
#[error("Codex reset observation account scope does not match the requested partition")]
AccountScopeMismatch,
#[error("failed to read Codex reset observations: {0}")]
Read(#[source] std::io::Error),
#[error("failed to decode Codex reset observations: {0}")]
Deserialize(#[source] serde_json::Error),
#[error("unsupported Codex reset observation store version {0}")]
UnsupportedVersion(u32),
#[error("failed to encode Codex reset observations: {0}")]
Serialize(#[source] serde_json::Error),
#[error("failed to persist Codex reset observations: {0}")]
Persist(#[source] std::io::Error),
}

#[derive(Debug, Clone, Serialize, Deserialize)]
struct CodexResetObservationStore {
version: u32,
#[serde(default)]
accounts: BTreeMap<String, Vec<CodexResetObservation>>,
}

impl Default for CodexResetObservationStore {
fn default() -> Self {
Self {
version: STORE_VERSION,
accounts: BTreeMap::new(),
}
}
}

/// One observed weekly reset: when it was seen, and the boundary it pointed at.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CodexResetObservation {
pub captured_at: chrono::DateTime<chrono::Utc>,
pub resets_at: chrono::DateTime<chrono::Utc>,
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CodexResetObservationMergeResult {
pub observations: Vec<CodexResetObservation>,
pub changed: bool,
}

/// Path of the store under the shared configuration root.
pub fn default_store_path() -> Result<PathBuf, CodexResetObservationError> {
let root = dirs::config_dir().ok_or_else(|| {
CodexResetObservationError::Read(std::io::Error::new(
std::io::ErrorKind::NotFound,
"configuration directory not found",
))
})?;
Ok(root.join("CodexBar").join(STORE_RELATIVE_PATH))
}

/// Path of the store relative to an explicit config root (tests, proof homes).
pub fn store_path(config_root: &Path) -> PathBuf {
config_root.join(STORE_RELATIVE_PATH)
}

fn validate_scope(account_scope: &str) -> Result<(), CodexResetObservationError> {
if account_scope.trim().is_empty() {
Err(CodexResetObservationError::EmptyAccountScope)
} else {
Ok(())
}
}

/// Read the observations recorded for `account_scope`.
pub fn load_reset_observations(
config_root: &Path,
account_scope: &str,
) -> Result<Vec<CodexResetObservation>, CodexResetObservationError> {
validate_scope(account_scope)?;
let path = store_path(config_root);
let raw = match secure_file::read_string(&path) {
Ok(raw) => raw,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(error) => return Err(CodexResetObservationError::Read(error)),
};
let store: CodexResetObservationStore =
serde_json::from_str(&raw).map_err(CodexResetObservationError::Deserialize)?;
if store.version != STORE_VERSION {
return Err(CodexResetObservationError::UnsupportedVersion(
store.version,
));
}
Ok(store
.accounts
.get(account_scope)
.cloned()
.unwrap_or_default())
}

/// Deduplicate and order observations: unique by `resets_at` (newest
/// `captured_at` wins), then sorted by reset time. Keeps the log bounded.
pub fn merge_reset_observations(
existing: &mut Vec<CodexResetObservation>,
incoming: impl IntoIterator<Item = CodexResetObservation>,
) -> bool {
let before = existing.clone();
for observation in incoming {
match existing
.iter_mut()
.find(|row| row.resets_at == observation.resets_at)
{
Some(row) if row.captured_at < observation.captured_at => {
row.captured_at = observation.captured_at;
}
Some(_) => {}
None => existing.push(observation),
}
}
existing.sort_by_key(|row| row.resets_at);
if existing.len() > MAX_OBSERVATIONS_PER_ACCOUNT {
let excess = existing.len() - MAX_OBSERVATIONS_PER_ACCOUNT;
existing.drain(..excess);
}
*existing != before
}

/// Merge one freshly observed reset and persist the store.
///
/// The write is atomic: bytes are staged with `secure_file::write_string` and
/// published with `atomic_file::replace_staged`, so a failure leaves the
/// previous store intact and never exposes a partly written secret.
pub fn merge_and_persist_reset_observation(
config_root: &Path,
account_scope: &str,
observed_resets_at: chrono::DateTime<chrono::Utc>,
captured_at: chrono::DateTime<chrono::Utc>,
) -> Result<CodexResetObservationMergeResult, CodexResetObservationError> {
validate_scope(account_scope)?;
if account_scope != CODEX_ACCOUNT_SCOPE {
return Err(CodexResetObservationError::AccountScopeMismatch);
}

let path = store_path(config_root);
let mut store = match secure_file::read_string(&path) {
Ok(raw) if raw.trim().is_empty() => CodexResetObservationStore::default(),
Ok(raw) => serde_json::from_str(&raw).map_err(CodexResetObservationError::Deserialize)?,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
CodexResetObservationStore::default()
}
Err(error) => return Err(CodexResetObservationError::Read(error)),
};
if store.version != STORE_VERSION {
return Err(CodexResetObservationError::UnsupportedVersion(
store.version,
));
}
let (changed, observations) = {
let rows = store.accounts.entry(account_scope.to_string()).or_default();
let changed = merge_reset_observations(
rows,
[CodexResetObservation {
captured_at,
resets_at: observed_resets_at,
}],
);
(changed, rows.clone())
};
if changed || !path.exists() {
let raw =
serde_json::to_string_pretty(&store).map_err(CodexResetObservationError::Serialize)?;
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).map_err(CodexResetObservationError::Read)?;
}
let file_name = path
.file_name()
.and_then(|name| name.to_str())
.unwrap_or("quota-reset-observations-v1.json");
let temp = path.with_file_name(format!(".{file_name}.tmp-{}", Uuid::new_v4()));
let result = (|| {
secure_file::write_string(&temp, &raw)?;
atomic_file::replace_staged(&temp, &path)
})();
if let Err(error) = result {
let cleanup_result = std::fs::remove_file(&temp);
if let Err(cleanup_error) = cleanup_result {
tracing::warn!(
"failed to clean staged codex reset observations file: {cleanup_error}"
);
}
tracing::debug!("publishing reset observations failed: {error}");
return Err(CodexResetObservationError::Read(error));
}
}
Ok(CodexResetObservationMergeResult {
observations,
changed,
})
}

#[cfg(test)]
mod tests {
use super::*;
use chrono::TimeZone;

fn temp_root(tag: &str) -> tempfile::TempDir {
tempfile::Builder::new()
.prefix(&format!("codex-resets-{tag}-"))
.tempdir()
.expect("temp dir")
}

fn at(secs: i64) -> chrono::DateTime<chrono::Utc> {
chrono::Utc.timestamp_opt(secs, 0).unwrap()
}

#[test]
fn rejects_empty_scope() {
let root = temp_root("empty-scope");
assert!(matches!(
load_reset_observations(root.path(), " "),
Err(CodexResetObservationError::EmptyAccountScope)
));
}

#[test]
fn rejects_foreign_scope() {
let root = temp_root("foreign-scope");
assert!(matches!(
merge_and_persist_reset_observation(root.path(), "claude", at(1_000), at(2_000)),
Err(CodexResetObservationError::AccountScopeMismatch)
));
}

#[test]
fn missing_store_loads_empty() {
let root = temp_root("missing");
assert!(
load_reset_observations(root.path(), CODEX_ACCOUNT_SCOPE)
.expect("load")
.is_empty()
);
}

#[test]
fn observation_survives_reload_and_dedupes_by_reset() {
let root = temp_root("dedupe");
merge_and_persist_reset_observation(root.path(), CODEX_ACCOUNT_SCOPE, at(5_000), at(1_000))
.expect("first merge");
// Same reset observed again later: keep the newer capture, still one row.
merge_and_persist_reset_observation(root.path(), CODEX_ACCOUNT_SCOPE, at(5_000), at(2_000))
.expect("second merge");
let rows = load_reset_observations(root.path(), CODEX_ACCOUNT_SCOPE).expect("load");
assert_eq!(
rows,
vec![CodexResetObservation {
captured_at: at(2_000),
resets_at: at(5_000)
}]
);
// An older capture for the same reset must not regress it.
merge_and_persist_reset_observation(root.path(), CODEX_ACCOUNT_SCOPE, at(5_000), at(500))
.expect("third merge");
let rows = load_reset_observations(root.path(), CODEX_ACCOUNT_SCOPE).expect("load");
assert_eq!(rows[0].captured_at, at(2_000));
}

#[test]
fn observations_are_sorted_by_reset_and_bounded() {
let root = temp_root("sort");
merge_and_persist_reset_observation(root.path(), CODEX_ACCOUNT_SCOPE, at(9_000), at(1_000))
.expect("late");
merge_and_persist_reset_observation(root.path(), CODEX_ACCOUNT_SCOPE, at(3_000), at(1_100))
.expect("early");
let rows = load_reset_observations(root.path(), CODEX_ACCOUNT_SCOPE).expect("load");
assert_eq!(
rows.iter().map(|row| row.resets_at).collect::<Vec<_>>(),
vec![at(3_000), at(9_000)]
);
}

#[test]
fn unsupported_version_is_rejected() {
let root = temp_root("version");
let path = store_path(root.path());
std::fs::create_dir_all(path.parent().unwrap()).expect("mkdir");
secure_file::write_string(&path, r#"{"version": 99, "accounts": {}}"#).expect("write");
assert!(matches!(
load_reset_observations(root.path(), CODEX_ACCOUNT_SCOPE),
Err(CodexResetObservationError::UnsupportedVersion(99))
));
}

#[test]
fn persisted_file_is_secure_wrapped() {
let root = temp_root("secure");
merge_and_persist_reset_observation(root.path(), CODEX_ACCOUNT_SCOPE, at(7_000), at(1_000))
.expect("merge");
let raw = std::fs::read_to_string(store_path(root.path())).expect("read");
assert!(raw.contains("codexbar.secure-file"), "not DPAPI wrapped");
assert!(!raw.contains("resets_at"), "plaintext payload leaked");
}
}