From 1f8aa6791aced7a6e83e96472dd4e2c56dd1665a Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 29 Sep 2026 11:58:26 +0000 Subject: [PATCH 1/4] Add manual weekly reset anchors for Claude Code A Claude Code account can store one manual UTC weekly reset. Completed weeks are emitted as the existing quota cycle contributions when this device has attributed usage in the week or on a boundary UTC day. Weeks with no such usage are omitted, and the cycle that contains now is never emitted. Anchors that differ by a whole number of weeks are the same schedule. Merging accounts transfers a schedule, keeps one copy when the phases match, and refuses the merge when they do not. Removing an account treats a stored schedule as a reference. Usage for these contributions is read only for the Claude Code account and only on the boundary days. EXPLAIN QUERY PLAN showed the existing provider/source index cannot bound that lookup, so this adds an index on provider, account, and start time. Co-authored-by: Dmitry Starkov <21260939+starkdmi@users.noreply.github.com> --- crates/statsai-store/src/lib.rs | 5 +- crates/statsai-store/src/migrations/mod.rs | 8 +- crates/statsai-store/src/migrations/v2.rs | 29 + .../src/pricing/tests/repricing.rs | 2 +- crates/statsai-store/src/quota/mod.rs | 2 + crates/statsai-store/src/quota/sync.rs | 1 + crates/statsai-store/src/quota/tests/mod.rs | 1 + .../statsai-store/src/quota/tests/weekly.rs | 547 ++++++++++++++++++ crates/statsai-store/src/quota/weekly.rs | 453 +++++++++++++++ crates/statsai/src/cli/account.rs | 47 +- crates/statsai/src/cli/args/account.rs | 44 ++ crates/statsai/src/cli/mod.rs | 1 + crates/statsai/src/cli/tests/account/mod.rs | 1 + .../src/cli/tests/account/weekly_reset.rs | 277 +++++++++ crates/statsai/src/cli/weekly_reset.rs | 121 ++++ crates/statsai/tests/cli_surface/help.txt | 60 +- docs/quota-sync.md | 44 ++ 17 files changed, 1633 insertions(+), 10 deletions(-) create mode 100644 crates/statsai-store/src/quota/tests/weekly.rs create mode 100644 crates/statsai-store/src/quota/weekly.rs create mode 100644 crates/statsai/src/cli/tests/account/weekly_reset.rs create mode 100644 crates/statsai/src/cli/weekly_reset.rs diff --git a/crates/statsai-store/src/lib.rs b/crates/statsai-store/src/lib.rs index 3bc5eb7..15d34b7 100644 --- a/crates/statsai-store/src/lib.rs +++ b/crates/statsai-store/src/lib.rs @@ -101,7 +101,10 @@ pub use privacy::{ FilteredConversationMetadata, FilteredConversationRecord, PrivacyDatasetStatus, PrivacyFailureRecord, PrivacyFindingRecord, }; -pub use quota::{QuotaDateRange, QuotaQuery, QuotaStatus}; +pub use quota::{ + weekly_cycle_containing, QuotaDateRange, QuotaQuery, QuotaStatus, WeeklyResetAnchor, + WEEKLY_RESET_PERIOD_SECONDS, +}; pub use tasks::{ derive_task_work_items, NamedTaskBenchmark, TaskBenchmarkMetrics, TaskBenchmarkReport, TaskDeletionImpact, TaskRebuildReport, TaskRebuildTimings, TaskStats, diff --git a/crates/statsai-store/src/migrations/mod.rs b/crates/statsai-store/src/migrations/mod.rs index 6d6b9a1..e944d3c 100644 --- a/crates/statsai-store/src/migrations/mod.rs +++ b/crates/statsai-store/src/migrations/mod.rs @@ -8,7 +8,7 @@ mod v2; pub(crate) use v1::*; pub(crate) use v2::*; -pub const CURRENT_SCHEMA_VERSION: i64 = 28; +pub const CURRENT_SCHEMA_VERSION: i64 = 29; pub fn migrate(conn: &Connection) -> Result<()> { if let Some(current) = existing_schema_version(conn)? { @@ -153,6 +153,7 @@ fn apply_migration(conn: &Connection, version: i64) -> Result<()> { 26 => apply_migration_026(conn), 27 => apply_migration_027(conn), 28 => apply_migration_028(conn), + 29 => apply_migration_029(conn), _ => bail!("unsupported schema migration version {version}"), } } @@ -350,6 +351,11 @@ mod tests { assert!(table_exists(&conn, "activity_rollups")); assert!(table_exists(&conn, "activity_coverage")); assert!(table_exists(&conn, "activity_scan_cursors")); + assert!(table_exists(&conn, "weekly_reset_anchors")); + assert!(index_exists( + &conn, + "usage_events_provider_account_started_idx" + )); } /// The scan filters are only fast while SQLite can match them to the indexes diff --git a/crates/statsai-store/src/migrations/v2.rs b/crates/statsai-store/src/migrations/v2.rs index 1b64771..021e46f 100644 --- a/crates/statsai-store/src/migrations/v2.rs +++ b/crates/statsai-store/src/migrations/v2.rs @@ -655,3 +655,32 @@ pub(crate) fn apply_migration_028(conn: &Connection) -> Result<()> { )?; Ok(()) } + +/// Manual weekly reset anchors, and the event index those anchors read through. +/// +/// A Claude weekly contribution asks two questions of `usage_events`: which UTC +/// days this account has any usage on, and which events fall on the partial +/// days at each reset. `usage_events_semantic_lookup_idx` leads with provider +/// and then source, so it cannot bound a single provider account, and both +/// questions would scan every provider's events. The new index serves +/// `(provider, provider_account_id, started_at)` as a range. +pub(crate) fn apply_migration_029(conn: &Connection) -> Result<()> { + conn.execute_batch( + r#" + CREATE TABLE IF NOT EXISTS weekly_reset_anchors ( + provider TEXT NOT NULL, + provider_account_id TEXT NOT NULL, + anchor_epoch_seconds INTEGER NOT NULL, + source TEXT NOT NULL, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (provider, provider_account_id), + CHECK (provider = 'claude_code'), + CHECK (source = 'manual') + ); + CREATE INDEX IF NOT EXISTS usage_events_provider_account_started_idx + ON usage_events (provider, provider_account_id, started_at); + "#, + )?; + Ok(()) +} diff --git a/crates/statsai-store/src/pricing/tests/repricing.rs b/crates/statsai-store/src/pricing/tests/repricing.rs index 266b670..20b9f8b 100644 --- a/crates/statsai-store/src/pricing/tests/repricing.rs +++ b/crates/statsai-store/src/pricing/tests/repricing.rs @@ -384,7 +384,7 @@ fn daily_rollups_are_unused_by_report_and_sync_paths() { 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, 28); + assert_eq!(CURRENT_SCHEMA_VERSION, 29); } #[test] diff --git a/crates/statsai-store/src/quota/mod.rs b/crates/statsai-store/src/quota/mod.rs index 9146c0a..9f5821b 100644 --- a/crates/statsai-store/src/quota/mod.rs +++ b/crates/statsai-store/src/quota/mod.rs @@ -18,10 +18,12 @@ mod cycles; mod observations; mod reconstruct; mod sync; +mod weekly; mod windows; pub(crate) use cycles::*; pub(crate) use reconstruct::*; +pub use weekly::{weekly_cycle_containing, WeeklyResetAnchor, WEEKLY_RESET_PERIOD_SECONDS}; pub(crate) const RESET_CLUSTER_TOLERANCE_SECONDS: i64 = 5 * 60; /// A window cannot be observed once it has reset: the provider issues a fresh diff --git a/crates/statsai-store/src/quota/sync.rs b/crates/statsai-store/src/quota/sync.rs index 1f80793..8f438c6 100644 --- a/crates/statsai-store/src/quota/sync.rs +++ b/crates/statsai-store/src/quota/sync.rs @@ -189,6 +189,7 @@ impl Store { boundary_slices, }); } + contributions.extend(self.manual_weekly_contributions(query, device_id, Utc::now())?); Ok(contributions) } diff --git a/crates/statsai-store/src/quota/tests/mod.rs b/crates/statsai-store/src/quota/tests/mod.rs index 74d396c..a657a08 100644 --- a/crates/statsai-store/src/quota/tests/mod.rs +++ b/crates/statsai-store/src/quota/tests/mod.rs @@ -5,4 +5,5 @@ mod observations; mod reconstruct; mod support; mod sync; +mod weekly; mod windows; diff --git a/crates/statsai-store/src/quota/tests/weekly.rs b/crates/statsai-store/src/quota/tests/weekly.rs new file mode 100644 index 0000000..52c29b9 --- /dev/null +++ b/crates/statsai-store/src/quota/tests/weekly.rs @@ -0,0 +1,547 @@ +use super::support::*; +use super::*; +use chrono::TimeZone; + +fn at(year: i32, month: u32, day: u32, hour: u32, min: u32, sec: u32) -> DateTime { + Utc.with_ymd_and_hms(year, month, day, hour, min, sec) + .single() + .expect("timestamp") +} + +fn anchor_account(store: &Store, anchor: DateTime) -> ProviderAccountId { + let account_id = ProviderAccountId("account-claude".to_string()); + store + .upsert_weekly_reset_anchor("claude_code", &account_id, anchor) + .expect("anchor"); + account_id +} + +fn claude_event( + store: &Store, + account_id: &ProviderAccountId, + started_at: DateTime, + record_id: &str, + usage: UsageCounts, + estimated_cost_micro_usd: i64, +) { + let (source_id, _) = assigned_source(store, at(2024, 1, 1, 0, 0, 0)); + let mut event = sample_usage_event( + &source_id, + account_id, + started_at, + record_id, + usage.input_tokens.unwrap_or(0), + usage.cache_read_tokens.unwrap_or(0), + usage.output_tokens.unwrap_or(0), + usage.reasoning_tokens.unwrap_or(0), + estimated_cost_micro_usd, + ); + event.provider = "claude_code".to_string(); + event.usage = usage; + event.event_id = event_id("claude_code", &source_id, record_id, None, started_at); + store.insert_event(&event).expect("event"); +} + +fn simple_usage(total: u64) -> UsageCounts { + UsageCounts { + input_tokens: Some(total), + total_tokens: Some(total), + ..UsageCounts::default() + } +} + +fn contributions_at( + store: &Store, + now: DateTime, +) -> Vec { + store + .manual_weekly_contributions(&QuotaQuery::default(), "device-weekly", now) + .expect("contributions") +} + +#[test] +fn weekly_cycle_containing_is_half_open_and_week_equivalent() { + let anchor = at(2026, 1, 8, 18, 0, 0); + let later = anchor + Duration::seconds(WEEKLY_RESET_PERIOD_SECONDS); + let now = at(2026, 1, 20, 12, 0, 0); + assert_eq!( + weekly_cycle_containing(anchor, now), + weekly_cycle_containing(later, now) + ); + let (start, end) = weekly_cycle_containing(anchor, now); + assert_eq!(end - start, Duration::seconds(WEEKLY_RESET_PERIOD_SECONDS)); + assert!(start <= now && now < end); + assert_eq!( + end.timestamp().rem_euclid(WEEKLY_RESET_PERIOD_SECONDS), + anchor.timestamp().rem_euclid(WEEKLY_RESET_PERIOD_SECONDS) + ); + let (reset_start, reset_end) = weekly_cycle_containing(anchor, anchor); + assert_eq!(reset_start, anchor); + assert_eq!( + reset_end, + anchor + Duration::seconds(WEEKLY_RESET_PERIOD_SECONDS) + ); +} + +#[test] +fn manual_weekly_anchor_rejects_other_providers_and_stores_whole_seconds() { + let store = Store::in_memory().expect("store"); + let account_id = ProviderAccountId("account-claude".to_string()); + let error = store + .upsert_weekly_reset_anchor("codex", &account_id, at(2026, 1, 8, 18, 0, 0)) + .expect_err("codex anchor"); + assert!(error.to_string().contains("only supported for claude_code")); + let fractional = at(2026, 1, 8, 18, 0, 0) + Duration::nanoseconds(900_000_000); + store + .upsert_weekly_reset_anchor("claude_code", &account_id, fractional) + .expect("anchor"); + let stored = store + .weekly_reset_anchor("claude_code", &account_id) + .expect("read") + .expect("stored"); + assert_eq!(stored.anchor, at(2026, 1, 8, 18, 0, 0)); + assert_eq!(stored.source, "manual"); + let created_at = stored.created_at; + store + .upsert_weekly_reset_anchor("claude_code", &account_id, at(2026, 3, 1, 4, 0, 0)) + .expect("replace"); + let replaced = store + .weekly_reset_anchor("claude_code", &account_id) + .expect("read") + .expect("stored"); + assert_eq!(replaced.anchor, at(2026, 3, 1, 4, 0, 0)); + assert_eq!(replaced.created_at, created_at); + assert!(replaced.updated_at >= created_at); +} + +#[test] +fn equivalent_anchors_emit_the_same_completed_weeks() { + let now = at(2026, 2, 1, 0, 0, 0); + let event_at = at(2026, 1, 10, 12, 0, 0); + let mut ids = Vec::new(); + for anchor in [ + at(2026, 1, 8, 18, 0, 0), + at(2026, 1, 8, 18, 0, 0) + Duration::seconds(WEEKLY_RESET_PERIOD_SECONDS), + ] { + let store = Store::in_memory().expect("store"); + let account_id = anchor_account(&store, anchor); + claude_event( + &store, + &account_id, + event_at, + "inside", + simple_usage(10), + 10, + ); + let contributions = contributions_at(&store, now); + assert_eq!(contributions.len(), 1); + ids.push(contributions[0].contribution_id.clone()); + assert_eq!( + contributions[0].representative_reset, + at(2026, 1, 15, 18, 0, 0) + ); + } + assert_eq!(ids[0], ids[1]); +} + +#[test] +fn emission_rules_skip_empty_weeks_and_the_current_cycle() { + let store = Store::in_memory().expect("store"); + let account_id = anchor_account(&store, at(2026, 1, 8, 18, 0, 0)); + claude_event( + &store, + &account_id, + at(2026, 1, 7, 12, 0, 0), + "early", + simple_usage(10), + 10, + ); + claude_event( + &store, + &account_id, + at(2026, 1, 20, 12, 0, 0), + "later", + simple_usage(10), + 10, + ); + let now = at(2026, 1, 30, 0, 0, 0); + let ends = contributions_at(&store, now) + .into_iter() + .map(|contribution| contribution.representative_reset) + .collect::>(); + assert_eq!( + ends, + vec![at(2026, 1, 8, 18, 0, 0), at(2026, 1, 22, 18, 0, 0)] + ); + + let during_current = contributions_at(&store, at(2026, 1, 21, 0, 0, 0)); + assert_eq!(during_current.len(), 1); + assert_eq!( + during_current[0].representative_reset, + at(2026, 1, 8, 18, 0, 0) + ); +} + +#[test] +fn a_reset_day_before_the_reset_emits_the_following_week_with_a_zero_slice() { + let store = Store::in_memory().expect("store"); + let account_id = anchor_account(&store, at(2026, 1, 8, 18, 0, 0)); + claude_event( + &store, + &account_id, + at(2026, 1, 15, 10, 0, 0), + "before-reset", + simple_usage(70), + 700, + ); + let contributions = contributions_at(&store, at(2026, 1, 23, 0, 0, 0)); + assert_eq!(contributions.len(), 2); + let previous = &contributions[0]; + let following = &contributions[1]; + assert_eq!(previous.representative_reset, at(2026, 1, 15, 18, 0, 0)); + assert_eq!(following.representative_reset, at(2026, 1, 22, 18, 0, 0)); + assert_eq!(previous.boundary_slices.len(), 2); + assert_eq!( + previous.boundary_slices[1].period_start, + at(2026, 1, 15, 0, 0, 0) + ); + assert_eq!( + previous.boundary_slices[1].period_end, + at(2026, 1, 15, 18, 0, 0) + ); + assert_eq!(previous.boundary_slices[1].input_tokens, 70); + assert_eq!( + following.boundary_slices[0].period_start, + at(2026, 1, 15, 18, 0, 0) + ); + assert_eq!( + following.boundary_slices[0].period_end, + at(2026, 1, 16, 0, 0, 0) + ); + assert_eq!(following.boundary_slices[0].total_tokens, 0); + assert_eq!(following.boundary_slices[0].estimated_cost_micro_usd, 0); +} + +#[test] +fn an_event_exactly_at_a_reset_belongs_to_the_next_week() { + let store = Store::in_memory().expect("store"); + let account_id = anchor_account(&store, at(2026, 1, 8, 18, 0, 0)); + claude_event( + &store, + &account_id, + at(2026, 1, 15, 18, 0, 0), + "on-reset", + simple_usage(9), + 90, + ); + let contributions = contributions_at(&store, at(2026, 1, 23, 0, 0, 0)); + // The instant is on the previous cycle's end day, so that week is emitted + // too, but the event itself starts the next week. + assert_eq!(contributions.len(), 2); + assert_eq!( + contributions[0].representative_reset, + at(2026, 1, 15, 18, 0, 0) + ); + assert_eq!(contributions[0].boundary_slices[1].input_tokens, 0); + assert_eq!( + contributions[1].representative_reset, + at(2026, 1, 22, 18, 0, 0) + ); + assert_eq!(contributions[1].boundary_slices[0].input_tokens, 9); + assert_eq!( + contributions[1].boundary_slices[0].period_start, + at(2026, 1, 15, 18, 0, 0) + ); +} + +#[test] +fn midnight_resets_have_no_boundary_slices_and_midday_resets_do_not_prorate() { + let midnight = Store::in_memory().expect("store"); + let midnight_account = anchor_account(&midnight, at(2026, 1, 12, 0, 0, 0)); + claude_event( + &midnight, + &midnight_account, + at(2026, 1, 7, 12, 0, 0), + "interior", + simple_usage(1000), + 1000, + ); + let midnight_contributions = contributions_at(&midnight, at(2026, 1, 20, 0, 0, 0)); + assert_eq!(midnight_contributions.len(), 1); + assert!(midnight_contributions[0].boundary_slices.is_empty()); + let containing = &midnight_contributions[0]; + assert_eq!(containing.representative_reset, at(2026, 1, 12, 0, 0, 0)); + assert!(containing.daily_envelopes.is_empty()); + assert_eq!(containing.limit_id, None); + assert!(!containing.has_schedule_overlap); + assert_eq!(containing.window_minutes, 10_080); + + let midday = Store::in_memory().expect("store"); + let account_id = anchor_account(&midday, at(2026, 1, 8, 18, 0, 0)); + claude_event( + &midday, + &account_id, + at(2026, 1, 10, 12, 0, 0), + "interior", + simple_usage(1000), + 1000, + ); + claude_event( + &midday, + &account_id, + at(2026, 1, 8, 20, 0, 0), + "boundary", + simple_usage(40), + 40, + ); + let contributions = contributions_at(&midday, at(2026, 1, 20, 0, 0, 0)); + let week = contributions + .iter() + .find(|contribution| contribution.representative_reset == at(2026, 1, 15, 18, 0, 0)) + .expect("week"); + assert_eq!(week.boundary_slices.len(), 2); + assert_eq!(week.boundary_slices[0].input_tokens, 40); + assert_eq!(week.boundary_slices[1].input_tokens, 0); + assert_eq!(week.boundary_slices[0].total_tokens, 40); + assert!(week + .boundary_slices + .iter() + .all(|slice| slice.total_tokens < 1000)); +} + +#[test] +fn boundary_slices_keep_token_categories_and_micro_usd() { + let store = Store::in_memory().expect("store"); + let account_id = anchor_account(&store, at(2026, 1, 8, 18, 0, 0)); + claude_event( + &store, + &account_id, + at(2026, 1, 8, 20, 0, 0), + "categories", + UsageCounts { + input_tokens: Some(11), + output_tokens: Some(44), + cache_creation_tokens: Some(22), + cache_creation_5m_tokens: Some(22), + cache_read_tokens: Some(33), + reasoning_tokens: Some(55), + total_tokens: Some(165), + ..UsageCounts::default() + }, + 123_456, + ); + let contributions = contributions_at(&store, at(2026, 1, 20, 0, 0, 0)); + let week = contributions + .iter() + .find(|contribution| contribution.representative_reset == at(2026, 1, 15, 18, 0, 0)) + .expect("week"); + let slice = &week.boundary_slices[0]; + assert_eq!(slice.input_tokens, 11); + assert_eq!(slice.cache_creation_tokens, 22); + assert_eq!(slice.cache_read_tokens, 33); + assert_eq!(slice.output_tokens, 44); + assert_eq!(slice.reasoning_tokens, 55); + assert_eq!(slice.total_tokens, 165); + assert_eq!(slice.estimated_cost_micro_usd, 123_456); + assert!(week.contribution_id.starts_with("quota_cycle_")); + assert_eq!(week.contribution_id.len(), "quota_cycle_".len() + 32); + assert!(week + .contribution_id + .chars() + .skip("quota_cycle_".len()) + .all(|ch| ch.is_ascii_hexdigit() && !ch.is_ascii_uppercase())); +} + +#[test] +fn late_imports_and_corrections_change_only_the_affected_contribution() { + let store = Store::in_memory().expect("store"); + let account_id = anchor_account(&store, at(2026, 1, 8, 18, 0, 0)); + let now = at(2026, 2, 1, 0, 0, 0); + claude_event( + &store, + &account_id, + at(2026, 1, 10, 12, 0, 0), + "first", + simple_usage(10), + 10, + ); + let before = contributions_at(&store, now); + assert_eq!(before.len(), 1); + let before_id = before[0].contribution_id.clone(); + let before_json = serde_json::to_string(&before[0]).expect("json"); + + claude_event( + &store, + &account_id, + at(2026, 1, 3, 12, 0, 0), + "imported", + simple_usage(5), + 5, + ); + let after_import = contributions_at(&store, now); + assert_eq!(after_import.len(), 2); + let unchanged = after_import + .iter() + .find(|contribution| contribution.contribution_id == before_id) + .expect("original week"); + assert_eq!(serde_json::to_string(unchanged).expect("json"), before_json); + + claude_event( + &store, + &account_id, + at(2026, 1, 8, 20, 0, 0), + "boundary-late", + simple_usage(8), + 8, + ); + let corrected = contributions_at(&store, now); + let changed = corrected + .iter() + .find(|contribution| contribution.contribution_id == before_id) + .expect("original week"); + assert_ne!(serde_json::to_string(changed).expect("json"), before_json); + assert_eq!(changed.boundary_slices[0].input_tokens, 8); + let other = corrected + .iter() + .find(|contribution| contribution.representative_reset == at(2026, 1, 8, 18, 0, 0)) + .expect("imported week"); + assert_eq!(other.boundary_slices.len(), 2); + assert_eq!(other.boundary_slices[0].input_tokens, 0); + assert_eq!(other.boundary_slices[1].input_tokens, 0); +} + +#[test] +fn changing_or_clearing_the_anchor_drops_the_old_contribution_ids() { + let store = Store::in_memory().expect("store"); + let account_id = ProviderAccountId("account-claude".to_string()); + store + .upsert_weekly_reset_anchor("claude_code", &account_id, at(2026, 1, 8, 18, 0, 0)) + .expect("anchor"); + claude_event( + &store, + &account_id, + at(2026, 1, 10, 12, 0, 0), + "used", + simple_usage(10), + 10, + ); + let now = at(2026, 2, 1, 0, 0, 0); + let original = contributions_at(&store, now); + assert_eq!(original.len(), 1); + store + .upsert_weekly_reset_anchor("claude_code", &account_id, at(2026, 1, 8, 21, 0, 0)) + .expect("shift"); + let shifted = contributions_at(&store, now); + assert!(!shifted.is_empty()); + assert!(shifted + .iter() + .all(|contribution| contribution.contribution_id != original[0].contribution_id)); + assert!(store + .delete_weekly_reset_anchor("claude_code", &account_id) + .expect("clear")); + assert!(contributions_at(&store, now).is_empty()); + assert_eq!( + store + .weekly_reset_anchor_count(Some("claude_code"), &account_id) + .expect("count"), + 0 + ); +} + +#[test] +fn usage_day_reads_use_the_provider_account_index() { + let store = Store::in_memory().expect("store"); + let plan = query_plan( + &store, + "SELECT DISTINCT substr(started_at, 1, 10) + FROM usage_events + WHERE provider = 'claude_code' AND provider_account_id = ?1", + ); + assert!( + plan.contains("usage_events_provider_account_started_idx"), + "usage day lookup does not use the account index: {plan}" + ); + let bounded = query_plan( + &store, + "SELECT payload FROM usage_events + WHERE provider = 'claude_code' + AND provider_account_id = ?1 + AND (started_at >= ?2 AND started_at < ?3)", + ); + assert!( + bounded.contains("usage_events_provider_account_started_idx"), + "boundary day lookup does not use the account index: {bounded}" + ); +} + +#[test] +fn other_providers_and_unattributed_events_do_not_emit_a_week() { + let store = Store::in_memory().expect("store"); + let account_id = anchor_account(&store, at(2026, 1, 8, 18, 0, 0)); + let (source_id, _) = assigned_source(&store, at(2024, 1, 1, 0, 0, 0)); + let mut codex = sample_usage_event( + &source_id, + &account_id, + at(2026, 1, 10, 12, 0, 0), + "codex-event", + 50, + 0, + 0, + 0, + 50, + ); + codex.provider_account_id = Some(account_id.clone()); + store.insert_event(&codex).expect("codex event"); + let mut unattributed = sample_usage_event( + &source_id, + &account_id, + at(2026, 1, 10, 13, 0, 0), + "unattributed", + 50, + 0, + 0, + 0, + 50, + ); + unattributed.provider = "claude_code".to_string(); + unattributed.provider_account_id = None; + store.insert_event(&unattributed).expect("unattributed"); + assert!(contributions_at(&store, at(2026, 2, 1, 0, 0, 0)).is_empty()); +} + +#[test] +fn quota_cycle_contributions_appends_completed_manual_weeks() { + let store = Store::in_memory().expect("store"); + let account_id = anchor_account(&store, at(2024, 6, 6, 18, 0, 0)); + claude_event( + &store, + &account_id, + at(2024, 6, 10, 12, 0, 0), + "historical", + simple_usage(4), + 4, + ); + let contributions = store + .quota_cycle_contributions(&QuotaQuery::default(), "device-weekly") + .expect("contributions"); + assert!(contributions + .iter() + .any(|contribution| contribution.provider == "claude_code" + && contribution.provider_account_id == account_id)); +} + +fn query_plan(store: &Store, sql: &str) -> String { + let explain = format!("EXPLAIN QUERY PLAN {sql}"); + let mut statement = store.conn.prepare(&explain).expect("explain"); + let bindings = (0..statement.parameter_count()) + .map(|_| rusqlite::types::Null) + .collect::>(); + let rows = statement + .query_map(rusqlite::params_from_iter(bindings), |row| { + row.get::<_, String>(3) + }) + .expect("plan") + .collect::>>() + .expect("plan rows"); + rows.join(" | ") +} diff --git a/crates/statsai-store/src/quota/weekly.rs b/crates/statsai-store/src/quota/weekly.rs new file mode 100644 index 0000000..867e357 --- /dev/null +++ b/crates/statsai-store/src/quota/weekly.rs @@ -0,0 +1,453 @@ +use super::*; +use crate::parse_rfc3339_for_row; +use anyhow::{bail, Context}; +use chrono::NaiveDate; + +/// Nominal length of a weekly reset cycle. Pure UTC; a device's timezone never +/// enters the arithmetic. +pub const WEEKLY_RESET_PERIOD_SECONDS: i64 = 604_800; + +const MANUAL_WEEKLY_PROVIDER: &str = "claude_code"; +const MANUAL_SCHEDULE_SOURCE: &str = "manual"; + +/// One manually entered weekly reset anchor. +/// +/// Anchors that differ by whole weeks are the same schedule. The stored instant +/// is whichever whole-second UTC value was entered; readers normalize it. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct WeeklyResetAnchor { + pub provider: String, + pub provider_account_id: ProviderAccountId, + pub anchor: DateTime, + pub source: String, + pub created_at: DateTime, + pub updated_at: DateTime, +} + +/// The half-open cycle `[end - 7d, end)` that contains `instant`. +/// +/// `end` is the unique instant strictly after `instant` with +/// `end ≡ anchor (mod 7d)`. An event timestamped exactly on a reset belongs to +/// the cycle that starts there. +pub fn weekly_cycle_containing( + anchor: DateTime, + instant: DateTime, +) -> (DateTime, DateTime) { + let end_epoch = cycle_end_after(anchor.timestamp(), instant.timestamp()); + let end = DateTime::from_timestamp(end_epoch, 0).expect("weekly cycle end"); + let start = DateTime::from_timestamp(end_epoch - WEEKLY_RESET_PERIOD_SECONDS, 0) + .expect("weekly cycle start"); + (start, end) +} + +impl Store { + pub fn upsert_weekly_reset_anchor( + &self, + provider: &str, + provider_account_id: &ProviderAccountId, + anchor: DateTime, + ) -> Result { + self.ensure_manual_weekly_provider(provider)?; + let anchor = anchor + .with_nanosecond(0) + .expect("clearing subseconds stays in range"); + let now = Utc::now().to_rfc3339(); + self.conn.execute( + r#" + INSERT INTO weekly_reset_anchors ( + provider, provider_account_id, anchor_epoch_seconds, source, created_at, updated_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?5) + ON CONFLICT(provider, provider_account_id) DO UPDATE SET + anchor_epoch_seconds = excluded.anchor_epoch_seconds, + source = excluded.source, + updated_at = excluded.updated_at + "#, + params![ + MANUAL_WEEKLY_PROVIDER, + &provider_account_id.0, + anchor.timestamp(), + MANUAL_SCHEDULE_SOURCE, + now, + ], + )?; + self.weekly_reset_anchor(provider, provider_account_id)? + .context("weekly reset anchor missing after upsert") + } + + pub fn weekly_reset_anchor( + &self, + provider: &str, + provider_account_id: &ProviderAccountId, + ) -> Result> { + if provider != MANUAL_WEEKLY_PROVIDER { + return Ok(None); + } + self.conn + .query_row( + r#" + SELECT provider, provider_account_id, anchor_epoch_seconds, source, created_at, updated_at + FROM weekly_reset_anchors + WHERE provider = ?1 AND provider_account_id = ?2 + "#, + params![provider, &provider_account_id.0], + |row| { + Ok(WeeklyResetAnchor { + provider: row.get(0)?, + provider_account_id: ProviderAccountId(row.get(1)?), + anchor: DateTime::from_timestamp(row.get(2)?, 0) + .expect("stored weekly reset anchor"), + source: row.get(3)?, + created_at: parse_rfc3339_for_row(&row.get::<_, String>(4)?, 4)?, + updated_at: parse_rfc3339_for_row(&row.get::<_, String>(5)?, 5)?, + }) + }, + ) + .optional() + .map_err(Into::into) + } + + pub fn delete_weekly_reset_anchor( + &self, + provider: &str, + provider_account_id: &ProviderAccountId, + ) -> Result { + let deleted = self.conn.execute( + r#" + DELETE FROM weekly_reset_anchors + WHERE provider = ?1 AND provider_account_id = ?2 + "#, + params![provider, &provider_account_id.0], + )?; + Ok(deleted > 0) + } + + pub fn move_weekly_reset_anchor( + &self, + provider: &str, + from_provider_account_id: &ProviderAccountId, + to_provider_account_id: &ProviderAccountId, + ) -> Result { + let updated_at = Utc::now().to_rfc3339(); + let moved = self.conn.execute( + r#" + UPDATE weekly_reset_anchors + SET provider_account_id = ?3, updated_at = ?4 + WHERE provider = ?1 AND provider_account_id = ?2 + "#, + params![ + provider, + &from_provider_account_id.0, + &to_provider_account_id.0, + updated_at, + ], + )?; + Ok(moved > 0) + } + + pub fn weekly_reset_anchor_count( + &self, + provider: Option<&str>, + provider_account_id: &ProviderAccountId, + ) -> Result { + let count: i64 = if let Some(provider) = provider { + self.conn.query_row( + r#" + SELECT COUNT(*) FROM weekly_reset_anchors + WHERE provider = ?1 AND provider_account_id = ?2 + "#, + params![provider, &provider_account_id.0], + |row| row.get(0), + )? + } else { + self.conn.query_row( + r#" + SELECT COUNT(*) FROM weekly_reset_anchors + WHERE provider_account_id = ?1 + "#, + params![&provider_account_id.0], + |row| row.get(0), + )? + }; + Ok(usize::try_from(count).unwrap_or(0)) + } + + /// Completed manual weekly contributions for the accounts that have an anchor. + /// + /// `now` decides which cycles have ended. The cycle that contains `now` is + /// never emitted. + pub(crate) fn manual_weekly_contributions( + &self, + query: &QuotaQuery, + device_id: &str, + now: DateTime, + ) -> Result> { + if query + .provider + .as_deref() + .is_some_and(|provider| provider != MANUAL_WEEKLY_PROVIDER) + { + return Ok(Vec::new()); + } + let anchors = self.weekly_reset_anchors(query.provider_account_id.as_ref())?; + let mut contributions = Vec::new(); + for anchor in anchors { + contributions + .extend(self.manual_weekly_contributions_for_anchor(device_id, &anchor, now)?); + } + Ok(contributions) + } + + fn weekly_reset_anchors( + &self, + provider_account_id: Option<&ProviderAccountId>, + ) -> Result> { + let mut statement = if provider_account_id.is_some() { + self.conn.prepare( + r#" + SELECT provider, provider_account_id, anchor_epoch_seconds, source, created_at, updated_at + FROM weekly_reset_anchors + WHERE provider = ?1 AND provider_account_id = ?2 + ORDER BY provider_account_id + "#, + )? + } else { + self.conn.prepare( + r#" + SELECT provider, provider_account_id, anchor_epoch_seconds, source, created_at, updated_at + FROM weekly_reset_anchors + WHERE provider = ?1 + ORDER BY provider_account_id + "#, + )? + }; + let map_row = |row: &rusqlite::Row<'_>| { + Ok(WeeklyResetAnchor { + provider: row.get(0)?, + provider_account_id: ProviderAccountId(row.get(1)?), + anchor: DateTime::from_timestamp(row.get(2)?, 0) + .expect("stored weekly reset anchor"), + source: row.get(3)?, + created_at: parse_rfc3339_for_row(&row.get::<_, String>(4)?, 4)?, + updated_at: parse_rfc3339_for_row(&row.get::<_, String>(5)?, 5)?, + }) + }; + let rows = if let Some(account_id) = provider_account_id { + statement.query_map(params![MANUAL_WEEKLY_PROVIDER, &account_id.0], map_row)? + } else { + statement.query_map(params![MANUAL_WEEKLY_PROVIDER], map_row)? + }; + rows.collect::, _>>().map_err(Into::into) + } + + fn manual_weekly_contributions_for_anchor( + &self, + device_id: &str, + anchor: &WeeklyResetAnchor, + now: DateTime, + ) -> Result> { + let usage_days = self.claude_code_usage_days(&anchor.provider_account_id)?; + let ends = + emitted_manual_cycle_ends(anchor.anchor.timestamp(), now.timestamp(), &usage_days); + if ends.is_empty() { + return Ok(Vec::new()); + } + let mut slices_by_end = BTreeMap::>::new(); + for end_epoch in &ends { + let start = DateTime::from_timestamp(end_epoch - WEEKLY_RESET_PERIOD_SECONDS, 0) + .expect("weekly cycle start"); + let end = DateTime::from_timestamp(*end_epoch, 0).expect("weekly cycle end"); + let builders = boundary_slice_builders( + start, + end, + MANUAL_WEEKLY_PROVIDER, + &anchor.provider_account_id, + ); + slices_by_end.insert(*end_epoch, builders.slices); + } + self.fill_manual_boundary_slices(&anchor.provider_account_id, &mut slices_by_end)?; + + let mut contributions = Vec::with_capacity(ends.len()); + for end_epoch in ends { + let end = DateTime::from_timestamp(end_epoch, 0).expect("weekly cycle end"); + contributions.push(QuotaCycleContributionV1 { + schema_version: QUOTA_CYCLE_CONTRIBUTION_SCHEMA_VERSION.to_string(), + contribution_id: manual_contribution_id( + device_id, + &anchor.provider_account_id.0, + end_epoch, + ), + provider: MANUAL_WEEKLY_PROVIDER.to_string(), + provider_account_id: anchor.provider_account_id.clone(), + limit_id: None, + window_minutes: QUOTA_WEEKLY_WINDOW_MINUTES, + representative_reset: end, + representative_reset_epoch_seconds: end_epoch, + has_schedule_overlap: false, + daily_envelopes: Vec::new(), + boundary_slices: slices_by_end.remove(&end_epoch).unwrap_or_default(), + }); + } + Ok(contributions) + } + + fn claude_code_usage_days( + &self, + provider_account_id: &ProviderAccountId, + ) -> Result> { + let mut statement = self.conn.prepare( + r#" + SELECT DISTINCT substr(started_at, 1, 10) + FROM usage_events + WHERE provider = 'claude_code' AND provider_account_id = ?1 + "#, + )?; + let rows = statement.query_map(params![&provider_account_id.0], |row| { + row.get::<_, String>(0) + })?; + let mut days = BTreeSet::new(); + for row in rows { + let day = row?; + if let Ok(parsed) = NaiveDate::parse_from_str(&day, "%Y-%m-%d") { + days.insert(parsed); + } + } + Ok(days) + } + + fn fill_manual_boundary_slices( + &self, + provider_account_id: &ProviderAccountId, + slices_by_end: &mut BTreeMap>, + ) -> Result<()> { + let mut days = BTreeSet::new(); + for slices in slices_by_end.values() { + for slice in slices { + days.insert(slice.period_start.date_naive()); + } + } + if days.is_empty() { + return Ok(()); + } + let day_list = days.into_iter().collect::>(); + for event in self.claude_code_events_on_utc_days(provider_account_id, &day_list)? { + let estimated_cost = event.cost.estimated_micro_usd(); + for slices in slices_by_end.values_mut() { + for slice in slices.iter_mut() { + if event.session.started_at < slice.period_start + || event.session.started_at >= slice.period_end + { + continue; + } + slice.add_usage(&event.usage, estimated_cost); + } + } + } + Ok(()) + } + + fn claude_code_events_on_utc_days( + &self, + provider_account_id: &ProviderAccountId, + days: &[NaiveDate], + ) -> Result> { + if days.is_empty() { + return Ok(Vec::new()); + } + const CHUNK_SIZE: usize = 64; + let mut events = Vec::new(); + for chunk in days.chunks(CHUNK_SIZE) { + let mut sql = String::from( + "SELECT payload FROM usage_events + WHERE provider = 'claude_code' + AND provider_account_id = ?1 + AND (", + ); + let mut bindings: Vec> = + vec![Box::new(provider_account_id.0.clone())]; + for (index, day) in chunk.iter().enumerate() { + if index > 0 { + sql.push_str(" OR "); + } + let start_index = bindings.len() + 1; + let end_index = start_index + 1; + sql.push_str(&format!( + "(started_at >= ?{start_index} AND started_at < ?{end_index})" + )); + let start = day.and_hms_opt(0, 0, 0).expect("utc midnight").and_utc(); + let end = start + Duration::days(1); + bindings.push(Box::new(start.to_rfc3339())); + bindings.push(Box::new(end.to_rfc3339())); + } + sql.push_str(") ORDER BY started_at, event_id"); + let mut statement = self.conn.prepare(&sql)?; + let params = bindings + .iter() + .map(|value| value.as_ref() as &dyn rusqlite::types::ToSql) + .collect::>(); + let rows = statement.query_map(params.as_slice(), |row| row.get::<_, String>(0))?; + for row in rows { + events.push(serde_json::from_str(&row?)?); + } + } + Ok(events) + } + + fn ensure_manual_weekly_provider(&self, provider: &str) -> Result<()> { + if provider != MANUAL_WEEKLY_PROVIDER { + bail!("weekly reset anchors are only supported for claude_code"); + } + Ok(()) + } +} + +fn cycle_end_after(anchor_epoch: i64, instant_epoch: i64) -> i64 { + let phase = anchor_epoch.rem_euclid(WEEKLY_RESET_PERIOD_SECONDS); + let instant_mod = instant_epoch.rem_euclid(WEEKLY_RESET_PERIOD_SECONDS); + if instant_mod < phase { + instant_epoch - instant_mod + phase + } else { + instant_epoch - instant_mod + phase + WEEKLY_RESET_PERIOD_SECONDS + } +} + +fn first_weekly_end_on_or_after(anchor_epoch: i64, instant_epoch: i64) -> i64 { + let phase = anchor_epoch.rem_euclid(WEEKLY_RESET_PERIOD_SECONDS); + let instant_mod = instant_epoch.rem_euclid(WEEKLY_RESET_PERIOD_SECONDS); + let delta = (phase - instant_mod).rem_euclid(WEEKLY_RESET_PERIOD_SECONDS); + instant_epoch + delta +} + +/// Completed cycle ends a usage day can force this device to emit. +/// +/// A day triggers every cycle whose start day or end day is that UTC day, which +/// is the half-open epoch window `[day_start, day_start + 8d)`. +pub(crate) fn emitted_manual_cycle_ends( + anchor_epoch: i64, + now_epoch: i64, + usage_days: &BTreeSet, +) -> Vec { + let mut ends = BTreeSet::new(); + for day in usage_days { + let day_start = day + .and_hms_opt(0, 0, 0) + .expect("utc midnight") + .and_utc() + .timestamp(); + let window_end = day_start + 86_400 + WEEKLY_RESET_PERIOD_SECONDS; + let mut end = first_weekly_end_on_or_after(anchor_epoch, day_start); + while end < window_end { + if end <= now_epoch { + ends.insert(end); + } + end += WEEKLY_RESET_PERIOD_SECONDS; + } + } + ends.into_iter().collect() +} + +fn manual_contribution_id(device_id: &str, account_id: &str, end_epoch: i64) -> String { + let digest = hash_text(&format!( + "quota_cycle_contribution.v1:{device_id}:claude_code:{account_id}:default:10080:manual:{end_epoch}" + )); + format!("quota_cycle_{}", &digest[..32]) +} diff --git a/crates/statsai/src/cli/account.rs b/crates/statsai/src/cli/account.rs index 22b46bc..cfd37e4 100644 --- a/crates/statsai/src/cli/account.rs +++ b/crates/statsai/src/cli/account.rs @@ -6,6 +6,8 @@ use statsai_core::{ IdentitySource, ProviderAccount, ProviderAccountId, Subscription, }; use statsai_store::{QuotaQuery, Store}; + +use super::weekly_reset::{weekly_reset, weekly_reset_phases_match}; use std::collections::{BTreeMap, HashMap}; use super::args::{AccountCommand, AccountSubcommand}; @@ -62,6 +64,7 @@ pub(crate) fn account(command: AccountCommand, store: &Store) -> Result<()> { )?; println!("{}", serde_json::to_string_pretty(&report)?); } + AccountSubcommand::WeeklyReset { command } => weekly_reset(command, store)?, AccountSubcommand::Remove { provider, account, @@ -89,6 +92,7 @@ pub(crate) struct AccountReferenceCounts { pub(crate) identity_observations: usize, pub(crate) plan_observations: usize, pub(crate) conversation_bindings: usize, + pub(crate) weekly_reset_schedules: usize, } impl AccountReferenceCounts { @@ -101,6 +105,7 @@ impl AccountReferenceCounts { + self.identity_observations + self.plan_observations + self.conversation_bindings + + self.weekly_reset_schedules } } @@ -118,6 +123,8 @@ pub(crate) struct AccountMergeReport { pub(crate) moved_identity_observations: usize, pub(crate) moved_plan_observations: usize, pub(crate) moved_conversation_bindings: usize, + pub(crate) moved_weekly_reset_schedules: usize, + pub(crate) coalesced_weekly_reset_schedules: usize, pub(crate) deleted_source_account: bool, pub(crate) remaining_references: AccountReferenceCounts, pub(crate) reset_local_sync_tracking: bool, @@ -174,8 +181,39 @@ pub(crate) fn merge_provider_accounts( .count(); let evidence_to_move = store.account_evidence_reference_counts(provider, &from.provider_account_id)?; + let from_schedule = store.weekly_reset_anchor(provider, &from.provider_account_id)?; + let to_schedule = store.weekly_reset_anchor(provider, &to.provider_account_id)?; + let (moved_weekly_reset_schedules, coalesced_weekly_reset_schedules) = match ( + &from_schedule, + &to_schedule, + ) { + (Some(from_anchor), Some(to_anchor)) => { + if !weekly_reset_phases_match( + from_anchor.anchor.timestamp(), + to_anchor.anchor.timestamp(), + ) { + bail!( + "weekly reset schedules disagree: {} and {} use different weekly reset times. Align or clear one schedule before merging these accounts", + display_account_identity(&from), + display_account_identity(&to) + ); + } + (0, 1) + } + (Some(_), None) => (1, 0), + _ => (0, 0), + }; if !dry_run { + if coalesced_weekly_reset_schedules == 1 { + store.delete_weekly_reset_anchor(provider, &from.provider_account_id)?; + } else if moved_weekly_reset_schedules == 1 { + store.move_weekly_reset_anchor( + provider, + &from.provider_account_id, + &to.provider_account_id, + )?; + } for assignment in &assignments_to_move { connect_source_to_account( store, @@ -225,6 +263,8 @@ pub(crate) fn merge_provider_accounts( moved_identity_observations: evidence_to_move.identity_observations, moved_plan_observations: evidence_to_move.plan_observations, moved_conversation_bindings: evidence_to_move.conversation_bindings, + moved_weekly_reset_schedules, + coalesced_weekly_reset_schedules, deleted_source_account, remaining_references, reset_local_sync_tracking: !dry_run, @@ -243,7 +283,7 @@ pub(crate) fn remove_orphan_provider_account( account_reference_counts(store, &account.provider_account_id, Some(provider))?; if remaining_references.total() > 0 { bail!( - "account {} still has references: {} source assignments, {} subscriptions, {} events, {} summaries, {} quota observations, {} identity observations, {} plan observations, {} conversation bindings", + "account {} still has references: {} source assignments, {} subscriptions, {} events, {} summaries, {} quota observations, {} identity observations, {} plan observations, {} conversation bindings, {} weekly reset schedules", display_account_identity(&account), remaining_references.source_account_assignments, remaining_references.subscriptions, @@ -252,7 +292,8 @@ pub(crate) fn remove_orphan_provider_account( remaining_references.quota_observations, remaining_references.identity_observations, remaining_references.plan_observations, - remaining_references.conversation_bindings + remaining_references.conversation_bindings, + remaining_references.weekly_reset_schedules ); } let deleted = if dry_run { @@ -586,6 +627,7 @@ pub(crate) fn account_reference_counts( conversation_bindings, } }; + let weekly_reset_schedules = store.weekly_reset_anchor_count(provider, provider_account_id)?; Ok(AccountReferenceCounts { source_account_assignments, @@ -596,5 +638,6 @@ pub(crate) fn account_reference_counts( identity_observations: evidence.identity_observations, plan_observations: evidence.plan_observations, conversation_bindings: evidence.conversation_bindings, + weekly_reset_schedules, }) } diff --git a/crates/statsai/src/cli/args/account.rs b/crates/statsai/src/cli/args/account.rs index da8b289..eab2fe3 100644 --- a/crates/statsai/src/cli/args/account.rs +++ b/crates/statsai/src/cli/args/account.rs @@ -39,6 +39,11 @@ pub(crate) enum AccountSubcommand { #[arg(long, help = "Preview the cleanup without writing")] dry_run: bool, }, + #[command(about = "Manage the manual weekly reset anchor for a Claude Code account")] + WeeklyReset { + #[command(subcommand)] + command: WeeklyResetSubcommand, + }, #[command(about = "Remove an unreferenced account row")] Remove { #[arg(long, help = "Provider name (claude_code, codex)")] @@ -52,3 +57,42 @@ pub(crate) enum AccountSubcommand { dry_run: bool, }, } + +#[derive(Debug, Subcommand)] +pub(crate) enum WeeklyResetSubcommand { + #[command(about = "Store a manual weekly reset anchor")] + Set { + #[arg(long, help = "Provider name. Only claude_code is accepted")] + provider: String, + #[arg( + long, + help = "Account identity (label, email, provider user id, or provider account id)" + )] + account: String, + #[arg( + long, + help = "Weekly reset instant as RFC3339 with an explicit offset or Z. Past and future instants are both accepted and stored as whole-second UTC" + )] + at: String, + }, + #[command(about = "Show the manual weekly reset anchor and the cycle that contains now")] + Show { + #[arg(long, help = "Provider name. Only claude_code is accepted")] + provider: String, + #[arg( + long, + help = "Account identity (label, email, provider user id, or provider account id)" + )] + account: String, + }, + #[command(about = "Remove the manual weekly reset anchor")] + Clear { + #[arg(long, help = "Provider name. Only claude_code is accepted")] + provider: String, + #[arg( + long, + help = "Account identity (label, email, provider user id, or provider account id)" + )] + account: String, + }, +} diff --git a/crates/statsai/src/cli/mod.rs b/crates/statsai/src/cli/mod.rs index e3e365a..cecbb2e 100644 --- a/crates/statsai/src/cli/mod.rs +++ b/crates/statsai/src/cli/mod.rs @@ -19,6 +19,7 @@ pub(super) mod store_admin; pub(super) mod subscription; pub(super) mod sync; pub(super) mod task; +pub(super) mod weekly_reset; pub(crate) use account::*; pub(crate) use activity::*; diff --git a/crates/statsai/src/cli/tests/account/mod.rs b/crates/statsai/src/cli/tests/account/mod.rs index 16ada8f..432135f 100644 --- a/crates/statsai/src/cli/tests/account/mod.rs +++ b/crates/statsai/src/cli/tests/account/mod.rs @@ -3,6 +3,7 @@ pub(crate) use super::*; mod claude_plans; mod merge; +mod weekly_reset; #[test] fn canonicalization_skips_accounts_without_surviving_evidence() { diff --git a/crates/statsai/src/cli/tests/account/weekly_reset.rs b/crates/statsai/src/cli/tests/account/weekly_reset.rs new file mode 100644 index 0000000..9b0a4c8 --- /dev/null +++ b/crates/statsai/src/cli/tests/account/weekly_reset.rs @@ -0,0 +1,277 @@ +use super::*; +use crate::cli::weekly_reset::{ + parse_weekly_reset_timestamp, resolve_manual_weekly_account, weekly_reset_view, +}; + +fn at(year: i32, month: u32, day: u32, hour: u32, min: u32, sec: u32) -> DateTime { + Utc.with_ymd_and_hms(year, month, day, hour, min, sec) + .single() + .expect("timestamp") +} + +fn claude_account(store: &Store, label: &str) -> statsai_core::ProviderAccount { + let account = test_account( + "claude_code", + Some(label), + None, + None, + None, + at(2026, 1, 1, 0, 0, 0), + ); + store.upsert_account(&account).expect("account"); + account +} + +#[test] +fn weekly_reset_timestamps_require_an_explicit_offset() { + let z = parse_weekly_reset_timestamp("2026-01-08T18:00:00Z").expect("z"); + let offset = parse_weekly_reset_timestamp("2026-01-08T18:00:00+00:00").expect("offset"); + let shifted = parse_weekly_reset_timestamp("2026-01-08T13:00:00-05:00").expect("shifted"); + assert_eq!(z, offset); + assert_eq!(z, shifted); + assert_eq!(z, at(2026, 1, 8, 18, 0, 0)); + let fractional = parse_weekly_reset_timestamp("2026-01-08T18:00:00.900Z").expect("fraction"); + assert_eq!(fractional, at(2026, 1, 8, 18, 0, 0)); + let past = parse_weekly_reset_timestamp("2020-01-08T18:00:00Z").expect("past"); + let future = parse_weekly_reset_timestamp("2030-01-08T18:00:00Z").expect("future"); + assert!(past < future); + let missing = parse_weekly_reset_timestamp("2026-01-08T18:00:00").expect_err("missing offset"); + assert!(missing.to_string().contains("explicit offset or Z")); +} + +#[test] +fn show_prints_the_same_cycle_for_anchors_a_week_apart() { + let store = Store::in_memory().expect("store"); + let account = claude_account(&store, "work"); + let earlier = store + .upsert_weekly_reset_anchor( + "claude_code", + &account.provider_account_id, + at(2026, 1, 8, 18, 0, 0), + ) + .expect("earlier"); + let now = at(2026, 1, 20, 12, 0, 0); + let first = weekly_reset_view(&account, &earlier, now); + store + .upsert_weekly_reset_anchor( + "claude_code", + &account.provider_account_id, + at(2026, 1, 15, 18, 0, 0), + ) + .expect("later"); + let later = store + .weekly_reset_anchor("claude_code", &account.provider_account_id) + .expect("read") + .expect("anchor"); + assert_ne!(earlier.anchor, later.anchor); + assert_eq!(weekly_reset_view(&account, &later, now), first); + assert_eq!(first.anchor, first.current_cycle_end); + assert_eq!(first.current_cycle_start, "2026-01-15T18:00:00Z"); + assert_eq!(first.current_cycle_end, "2026-01-22T18:00:00Z"); +} + +#[test] +fn weekly_reset_commands_reject_ambiguous_accounts_and_other_providers() { + let store = Store::in_memory().expect("store"); + claude_account(&store, "work"); + let codex = test_account( + "codex", + Some("work"), + None, + None, + None, + at(2026, 1, 1, 0, 0, 0), + ); + store.upsert_account(&codex).expect("codex"); + let wrong_provider = resolve_manual_weekly_account(&store, "codex", "work").expect_err("codex"); + assert!(wrong_provider + .to_string() + .contains("only supported for claude_code")); + let missing = + resolve_manual_weekly_account(&store, "claude_code", "missing").expect_err("missing"); + assert!(missing + .to_string() + .contains("no claude_code account matched")); + claude_account(&store, "ada"); + let by_email = test_account( + "claude_code", + Some("other"), + Some("ada"), + None, + None, + at(2026, 1, 1, 0, 0, 0), + ); + store.upsert_account(&by_email).expect("email account"); + let ambiguous = + resolve_manual_weekly_account(&store, "claude_code", "ada").expect_err("ambiguous"); + assert!(ambiguous + .to_string() + .contains("multiple claude_code accounts matched")); +} + +#[test] +fn merge_transfers_coalesces_and_rejects_weekly_reset_schedules() { + let store = Store::in_memory().expect("store"); + let from = claude_account(&store, "legacy"); + let to = claude_account(&store, "canonical"); + store + .upsert_weekly_reset_anchor( + "claude_code", + &from.provider_account_id, + at(2026, 1, 8, 18, 0, 0), + ) + .expect("from anchor"); + let transferred = merge_provider_accounts(&store, "claude_code", "legacy", "canonical", false) + .expect("transfer"); + assert_eq!(transferred.moved_weekly_reset_schedules, 1); + assert!(store + .weekly_reset_anchor("claude_code", &from.provider_account_id) + .expect("from") + .is_none()); + assert_eq!( + store + .weekly_reset_anchor("claude_code", &to.provider_account_id) + .expect("to") + .expect("moved") + .anchor, + at(2026, 1, 8, 18, 0, 0) + ); + + let other = claude_account(&store, "other"); + store + .upsert_weekly_reset_anchor( + "claude_code", + &other.provider_account_id, + at(2026, 1, 15, 18, 0, 0), + ) + .expect("equivalent"); + let coalesced = merge_provider_accounts(&store, "claude_code", "other", "canonical", false) + .expect("coalesce"); + assert_eq!(coalesced.coalesced_weekly_reset_schedules, 1); + assert_eq!(coalesced.moved_weekly_reset_schedules, 0); + assert!(store + .weekly_reset_anchor("claude_code", &to.provider_account_id) + .expect("kept") + .is_some()); + + let mismatched = claude_account(&store, "mismatched"); + store + .upsert_weekly_reset_anchor( + "claude_code", + &mismatched.provider_account_id, + at(2026, 1, 8, 19, 0, 0), + ) + .expect("mismatch"); + let rejected = merge_provider_accounts(&store, "claude_code", "mismatched", "canonical", false) + .expect_err("reject"); + assert!(rejected + .to_string() + .contains("different weekly reset times")); + assert!(store + .weekly_reset_anchor("claude_code", &mismatched.provider_account_id) + .expect("still stored") + .is_some()); + assert!(store + .weekly_reset_anchor("claude_code", &to.provider_account_id) + .expect("untouched") + .is_some()); +} + +#[test] +fn remove_treats_a_weekly_reset_schedule_as_a_reference() { + let store = Store::in_memory().expect("store"); + let account = claude_account(&store, "work"); + store + .upsert_weekly_reset_anchor( + "claude_code", + &account.provider_account_id, + at(2026, 1, 8, 18, 0, 0), + ) + .expect("anchor"); + let blocked = + remove_orphan_provider_account(&store, "claude_code", "work", false).expect_err("blocked"); + assert!(blocked.to_string().contains("weekly reset schedules")); + assert!(store + .delete_weekly_reset_anchor("claude_code", &account.provider_account_id) + .expect("clear")); + let removed = + remove_orphan_provider_account(&store, "claude_code", "work", false).expect("remove"); + assert!(removed.deleted); +} + +#[test] +fn anchor_changes_retire_manual_contribution_ids_on_the_next_sync() { + let store = Store::in_memory().expect("store"); + let account = claude_account(&store, "work"); + let source = SourceLocation::local_adapter( + "claude_code", + "test", + "0", + Path::new("/tmp/statsai-weekly-reset"), + LocationOrigin::Configured, + ); + store.upsert_source(&source).expect("source"); + store + .upsert_weekly_reset_anchor( + "claude_code", + &account.provider_account_id, + at(2026, 1, 8, 18, 0, 0), + ) + .expect("anchor"); + let mut event = test_event( + "claude_code", + &source, + at(2026, 1, 10, 12, 0, 0), + Some(account.provider_account_id.clone()), + TokenParts::total(12), + ); + event.provider = "claude_code".to_string(); + store.insert_event(&event).expect("event"); + + let command = test_sync_command("http"); + let target = sync_target(&command).expect("target"); + let (first, _) = build_sync_batch(&command, &store, "device-weekly", &target).expect("first"); + let original_ids = first + .quota_cycle_contributions + .iter() + .map(|contribution| contribution.contribution_id.clone()) + .collect::>(); + assert_eq!(original_ids.len(), 1); + assert!(first + .authoritative_snapshot + .as_ref() + .expect("first snapshot") + .quota_cycle_contribution_ids + .iter() + .any(|id| id == &original_ids[0])); + record_rollup_sync_success(&store, "http", &target, &first).expect("record"); + + store + .upsert_weekly_reset_anchor( + "claude_code", + &account.provider_account_id, + at(2026, 1, 8, 21, 0, 0), + ) + .expect("shift"); + let (shifted, _) = + build_sync_batch(&command, &store, "device-weekly", &target).expect("shifted"); + let snapshot = shifted + .authoritative_snapshot + .as_ref() + .expect("retirement snapshot"); + assert!(snapshot + .quota_cycle_contribution_ids + .iter() + .all(|id| id != &original_ids[0])); + assert!(!shifted.quota_cycle_contributions.is_empty()); + record_rollup_sync_success(&store, "http", &target, &shifted).expect("record shift"); + + assert!(store + .delete_weekly_reset_anchor("claude_code", &account.provider_account_id) + .expect("clear")); + let (cleared, _) = + build_sync_batch(&command, &store, "device-weekly", &target).expect("cleared"); + let cleared_snapshot = cleared.authoritative_snapshot.expect("clear snapshot"); + assert!(cleared_snapshot.quota_cycle_contribution_ids.is_empty()); + assert!(cleared.quota_cycle_contributions.is_empty()); +} diff --git a/crates/statsai/src/cli/weekly_reset.rs b/crates/statsai/src/cli/weekly_reset.rs new file mode 100644 index 0000000..fad5918 --- /dev/null +++ b/crates/statsai/src/cli/weekly_reset.rs @@ -0,0 +1,121 @@ +use anyhow::{bail, Context, Result}; +use chrono::{DateTime, SecondsFormat, Timelike, Utc}; +use serde::Serialize; +use statsai_core::{display_account_identity, ProviderAccount}; +use statsai_store::{ + weekly_cycle_containing, Store, WeeklyResetAnchor, WEEKLY_RESET_PERIOD_SECONDS, +}; + +use super::account::resolve_existing_provider_account_selector; +use super::args::WeeklyResetSubcommand; +use super::source::canonical_provider; + +const MANUAL_WEEKLY_PROVIDER: &str = "claude_code"; + +pub(crate) fn weekly_reset(command: WeeklyResetSubcommand, store: &Store) -> Result<()> { + let view = match command { + WeeklyResetSubcommand::Set { + provider, + account, + at, + } => { + let anchor = parse_weekly_reset_timestamp(&at)?; + let account = resolve_manual_weekly_account(store, &provider, &account)?; + let stored = store.upsert_weekly_reset_anchor( + MANUAL_WEEKLY_PROVIDER, + &account.provider_account_id, + anchor, + )?; + weekly_reset_view(&account, &stored, Utc::now()) + } + WeeklyResetSubcommand::Show { provider, account } => { + let account = resolve_manual_weekly_account(store, &provider, &account)?; + let stored = store + .weekly_reset_anchor(MANUAL_WEEKLY_PROVIDER, &account.provider_account_id)? + .context("no weekly reset anchor for this account")?; + weekly_reset_view(&account, &stored, Utc::now()) + } + WeeklyResetSubcommand::Clear { provider, account } => { + let account = resolve_manual_weekly_account(store, &provider, &account)?; + let cleared = store + .delete_weekly_reset_anchor(MANUAL_WEEKLY_PROVIDER, &account.provider_account_id)?; + println!( + "{}", + serde_json::to_string_pretty(&WeeklyResetClear { + provider: MANUAL_WEEKLY_PROVIDER.to_string(), + provider_account_id: account.provider_account_id.0, + cleared, + })? + ); + return Ok(()); + } + }; + println!("{}", serde_json::to_string_pretty(&view)?); + Ok(()) +} + +#[derive(Debug, Serialize, PartialEq, Eq)] +pub(crate) struct WeeklyResetView { + pub(crate) provider: String, + pub(crate) provider_account_id: String, + pub(crate) account: String, + /// The end of the cycle that contains `now`. Equivalent anchors, including + /// ones that differ by whole weeks, print this same instant. + pub(crate) anchor: String, + pub(crate) source: String, + pub(crate) current_cycle_start: String, + pub(crate) current_cycle_end: String, +} + +#[derive(Debug, Serialize)] +struct WeeklyResetClear { + provider: String, + provider_account_id: String, + cleared: bool, +} + +pub(crate) fn weekly_reset_view( + account: &ProviderAccount, + stored: &WeeklyResetAnchor, + now: DateTime, +) -> WeeklyResetView { + let (start, end) = weekly_cycle_containing(stored.anchor, now); + WeeklyResetView { + provider: stored.provider.clone(), + provider_account_id: account.provider_account_id.0.clone(), + account: display_account_identity(account), + anchor: format_utc(end), + source: stored.source.clone(), + current_cycle_start: format_utc(start), + current_cycle_end: format_utc(end), + } +} + +pub(crate) fn parse_weekly_reset_timestamp(value: &str) -> Result> { + let parsed = DateTime::parse_from_rfc3339(value.trim()) + .context("--at must be an RFC3339 timestamp with an explicit offset or Z")?; + let utc = parsed.with_timezone(&Utc); + Ok(utc + .with_nanosecond(0) + .expect("clearing subseconds stays in range")) +} + +pub(crate) fn resolve_manual_weekly_account( + store: &Store, + provider: &str, + selector: &str, +) -> Result { + let provider = canonical_provider(provider)?; + if provider != MANUAL_WEEKLY_PROVIDER { + bail!("weekly reset anchors are only supported for claude_code"); + } + resolve_existing_provider_account_selector(store, &provider, selector) +} + +fn format_utc(value: DateTime) -> String { + value.to_rfc3339_opts(SecondsFormat::Secs, true) +} + +pub(crate) fn weekly_reset_phases_match(left: i64, right: i64) -> bool { + left.rem_euclid(WEEKLY_RESET_PERIOD_SECONDS) == right.rem_euclid(WEEKLY_RESET_PERIOD_SECONDS) +} diff --git a/crates/statsai/tests/cli_surface/help.txt b/crates/statsai/tests/cli_surface/help.txt index d240415..7d72338 100644 --- a/crates/statsai/tests/cli_surface/help.txt +++ b/crates/statsai/tests/cli_surface/help.txt @@ -309,11 +309,12 @@ List canonical provider accounts Usage: statsai account [OPTIONS] Commands: - list List canonical provider accounts - plans Show detected plan evidence per account - merge Merge a legacy/manual account into an existing canonical account - remove Remove an unreferenced account row - help Print this message or the help of the given subcommand(s) + list List canonical provider accounts + plans Show detected plan evidence per account + merge Merge a legacy/manual account into an existing canonical account + weekly-reset Manage the manual weekly reset anchor for a Claude Code account + remove Remove an unreferenced account row + help Print this message or the help of the given subcommand(s) Options: --store Path to SQLite store @@ -353,6 +354,55 @@ Options: --to Destination account identity (label, email, provider user id, or provider account id) --dry-run Preview the cleanup without writing -h, --help Print help +===== statsai account weekly-reset --help ===== +Manage the manual weekly reset anchor for a Claude Code account + +Usage: statsai account weekly-reset [OPTIONS] + +Commands: + set Store a manual weekly reset anchor + show Show the manual weekly reset anchor and the cycle that contains now + clear Remove the manual weekly reset anchor + help Print this message or the help of the given subcommand(s) + +Options: + --store Path to SQLite store + --device-id Device identifier for multi-device sync + -h, --help Print help +===== statsai account weekly-reset set --help ===== +Store a manual weekly reset anchor + +Usage: statsai account weekly-reset set [OPTIONS] --provider --account --at + +Options: + --provider Provider name. Only claude_code is accepted + --store Path to SQLite store + --account Account identity (label, email, provider user id, or provider account id) + --device-id Device identifier for multi-device sync + --at Weekly reset instant as RFC3339 with an explicit offset or Z. Past and future instants are both accepted and stored as whole-second UTC + -h, --help Print help +===== statsai account weekly-reset show --help ===== +Show the manual weekly reset anchor and the cycle that contains now + +Usage: statsai account weekly-reset show [OPTIONS] --provider --account + +Options: + --provider Provider name. Only claude_code is accepted + --store Path to SQLite store + --account Account identity (label, email, provider user id, or provider account id) + --device-id Device identifier for multi-device sync + -h, --help Print help +===== statsai account weekly-reset clear --help ===== +Remove the manual weekly reset anchor + +Usage: statsai account weekly-reset clear [OPTIONS] --provider --account + +Options: + --provider Provider name. Only claude_code is accepted + --store Path to SQLite store + --account Account identity (label, email, provider user id, or provider account id) + --device-id Device identifier for multi-device sync + -h, --help Print help ===== statsai account remove --help ===== Remove an unreferenced account row diff --git a/docs/quota-sync.md b/docs/quota-sync.md index 1abcbe8..b3c27fe 100644 --- a/docs/quota-sync.md +++ b/docs/quota-sync.md @@ -102,3 +102,47 @@ installations cannot be blended. A logical cycle remains while any device still contributes it. Stale device contributions are retired through the same authoritative-snapshot ownership used by summaries and code-change metrics. + +## Manual weekly schedules + +Claude Code weekly value uses the same `quota_cycle_contribution.v1` record. +There is no separate wire schema. A device stores one manual weekly reset +anchor per Claude Code account (`statsai account weekly-reset`). The anchor is +a single UTC instant. Every device of the account is expected to store the same +phase; anchors that differ by a whole number of weeks are the same schedule. +Automatic capture from the status line `resets_at` is future work. + +Cycles are half-open `[end - 604800s, end)` with `end ≡ anchor (mod 604800)`. +An event exactly at a reset belongs to the following week. Arithmetic is UTC. +The device timezone does not move the boundary. Only completed cycles are +emitted (`end <= now`). The cycle that contains the current instant is omitted. + +A device emits a completed cycle when it has attributed Claude Code usage for +that account inside the cycle, or anywhere on the UTC day that contains the +cycle start or the UTC day that contains the cycle end. The second case covers +a reset that falls during a UTC day: a device whose last use that day was +before the reset still owes the following week a zero slice for the remainder +of the day, because the backend requires a boundary slice from every device +with usage on a boundary day. Weeks with no such usage are not emitted. The +backend fills those gaps from other devices and from daily summaries. + +A manual contribution sets `provider` to `claude_code`, `limit_id` to null, +`window_minutes` to `10080`, `representative_reset` to the cycle end, +`has_schedule_overlap` to false, and `daily_envelopes` to empty. Boundary +slices are the partial UTC days of `[start, end)`, including an explicit zero +slice when the device has no events in that partial day. A boundary that falls +on UTC midnight has no slice. Usage is not prorated. The contribution id is + +```text +quota_cycle_ + hash("quota_cycle_contribution.v1:{device}:claude_code:{account}:default:10080:manual:{end_epoch}")[:32] +``` + +Unchanged completed weeks keep their payload hash, so an incremental sync does +not resend them. Replacing the anchor with a different phase recomputes every +id. Clearing the anchor drops this device's manual contributions. The next +authoritative snapshot retires the ids that are no longer present. + +`account merge` transfers a schedule when only one account has one, keeps a +single copy when the two phases match, and refuses the merge when the phases +differ. `account remove` treats a stored schedule as a reference, so an account +that still has one is not deleted. From 83d32a18c83a260a37b899a8f110d9be6af648ac Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 29 Sep 2026 13:28:08 +0000 Subject: [PATCH 2/4] Honor quota query filters for manual weekly contributions Manual cycles now follow the same source, limit, and time filters as reconstructed cycles. A query for another source, a provider limit, or a date range that misses the week no longer returns that week's contribution. Co-authored-by: Dmitry Starkov <21260939+starkdmi@users.noreply.github.com> --- .../statsai-store/src/quota/tests/weekly.rs | 123 ++++++++++++++++++ crates/statsai-store/src/quota/weekly.rs | 74 +++++++++-- 2 files changed, 183 insertions(+), 14 deletions(-) diff --git a/crates/statsai-store/src/quota/tests/weekly.rs b/crates/statsai-store/src/quota/tests/weekly.rs index 52c29b9..fb84d34 100644 --- a/crates/statsai-store/src/quota/tests/weekly.rs +++ b/crates/statsai-store/src/quota/tests/weekly.rs @@ -530,6 +530,129 @@ fn quota_cycle_contributions_appends_completed_manual_weeks() { && contribution.provider_account_id == account_id)); } +#[test] +fn manual_contributions_honor_source_limit_and_time_filters() { + let store = Store::in_memory().expect("store"); + let account_id = anchor_account(&store, at(2026, 1, 8, 18, 0, 0)); + claude_event( + &store, + &account_id, + at(2026, 1, 10, 12, 0, 0), + "first-week", + simple_usage(4), + 4, + ); + claude_event( + &store, + &account_id, + at(2026, 1, 20, 12, 0, 0), + "second-week", + simple_usage(6), + 6, + ); + let source_id = store + .events() + .expect("events") + .into_iter() + .find(|event| event.provider == "claude_code") + .expect("claude event") + .source_id; + let all = store + .quota_cycle_contributions(&QuotaQuery::default(), "device-weekly") + .expect("all"); + let claude = all + .iter() + .filter(|contribution| contribution.provider == "claude_code") + .count(); + assert_eq!(claude, 2); + + let other_source = store + .quota_cycle_contributions( + &QuotaQuery { + source_id: Some(SourceId("other-source".to_string())), + ..QuotaQuery::default() + }, + "device-weekly", + ) + .expect("other source"); + assert!(other_source + .iter() + .all(|contribution| contribution.provider != "claude_code")); + + let matching_source = store + .quota_cycle_contributions( + &QuotaQuery { + source_id: Some(source_id), + ..QuotaQuery::default() + }, + "device-weekly", + ) + .expect("matching source"); + assert_eq!( + matching_source + .iter() + .filter(|contribution| contribution.provider == "claude_code") + .count(), + 2 + ); + + let limited = store + .quota_cycle_contributions( + &QuotaQuery { + limit_id: Some("five_hour".to_string()), + ..QuotaQuery::default() + }, + "device-weekly", + ) + .expect("limit"); + assert!(limited + .iter() + .all(|contribution| contribution.provider != "claude_code")); + + let future = store + .quota_cycle_contributions( + &QuotaQuery { + from: Some(at(2030, 1, 1, 0, 0, 0)), + ..QuotaQuery::default() + }, + "device-weekly", + ) + .expect("future"); + assert!(future + .iter() + .all(|contribution| contribution.provider != "claude_code")); + + let before = store + .quota_cycle_contributions( + &QuotaQuery { + to: Some(at(2020, 1, 1, 0, 0, 0)), + ..QuotaQuery::default() + }, + "device-weekly", + ) + .expect("before"); + assert!(before + .iter() + .all(|contribution| contribution.provider != "claude_code")); + + let second_only = store + .quota_cycle_contributions( + &QuotaQuery { + from: Some(at(2026, 1, 16, 0, 0, 0)), + to: Some(at(2026, 1, 21, 0, 0, 0)), + ..QuotaQuery::default() + }, + "device-weekly", + ) + .expect("second week"); + let resets = second_only + .iter() + .filter(|contribution| contribution.provider == "claude_code") + .map(|contribution| contribution.representative_reset) + .collect::>(); + assert_eq!(resets, vec![at(2026, 1, 22, 18, 0, 0)]); +} + fn query_plan(store: &Store, sql: &str) -> String { let explain = format!("EXPLAIN QUERY PLAN {sql}"); let mut statement = store.conn.prepare(&explain).expect("explain"); diff --git a/crates/statsai-store/src/quota/weekly.rs b/crates/statsai-store/src/quota/weekly.rs index 867e357..e2e0a85 100644 --- a/crates/statsai-store/src/quota/weekly.rs +++ b/crates/statsai-store/src/quota/weekly.rs @@ -188,11 +188,17 @@ impl Store { { return Ok(Vec::new()); } + // Manual cycles have no provider limit. A limit filter matches the + // reconstructed path, which drops windows whose limit id differs. + if query.limit_id.is_some() { + return Ok(Vec::new()); + } let anchors = self.weekly_reset_anchors(query.provider_account_id.as_ref())?; let mut contributions = Vec::new(); for anchor in anchors { - contributions - .extend(self.manual_weekly_contributions_for_anchor(device_id, &anchor, now)?); + contributions.extend( + self.manual_weekly_contributions_for_anchor(device_id, &anchor, query, now)?, + ); } Ok(contributions) } @@ -243,11 +249,16 @@ impl Store { &self, device_id: &str, anchor: &WeeklyResetAnchor, + query: &QuotaQuery, now: DateTime, ) -> Result> { - let usage_days = self.claude_code_usage_days(&anchor.provider_account_id)?; + let usage_days = + self.claude_code_usage_days(&anchor.provider_account_id, query.source_id.as_ref())?; let ends = - emitted_manual_cycle_ends(anchor.anchor.timestamp(), now.timestamp(), &usage_days); + emitted_manual_cycle_ends(anchor.anchor.timestamp(), now.timestamp(), &usage_days) + .into_iter() + .filter(|end_epoch| cycle_overlaps_query(*end_epoch, query)) + .collect::>(); if ends.is_empty() { return Ok(Vec::new()); } @@ -264,7 +275,11 @@ impl Store { ); slices_by_end.insert(*end_epoch, builders.slices); } - self.fill_manual_boundary_slices(&anchor.provider_account_id, &mut slices_by_end)?; + self.fill_manual_boundary_slices( + &anchor.provider_account_id, + query.source_id.as_ref(), + &mut slices_by_end, + )?; let mut contributions = Vec::with_capacity(ends.len()); for end_epoch in ends { @@ -293,17 +308,30 @@ impl Store { fn claude_code_usage_days( &self, provider_account_id: &ProviderAccountId, + source_id: Option<&SourceId>, ) -> Result> { - let mut statement = self.conn.prepare( + let sql = if source_id.is_some() { + r#" + SELECT DISTINCT substr(started_at, 1, 10) + FROM usage_events + WHERE provider = 'claude_code' AND provider_account_id = ?1 AND source_id = ?2 + "# + } else { r#" SELECT DISTINCT substr(started_at, 1, 10) FROM usage_events WHERE provider = 'claude_code' AND provider_account_id = ?1 - "#, - )?; - let rows = statement.query_map(params![&provider_account_id.0], |row| { - row.get::<_, String>(0) - })?; + "# + }; + let mut statement = self.conn.prepare(sql)?; + let rows = if let Some(source_id) = source_id { + statement.query_map( + params![&provider_account_id.0, &source_id.0], + usage_day_cell, + )? + } else { + statement.query_map(params![&provider_account_id.0], usage_day_cell)? + }; let mut days = BTreeSet::new(); for row in rows { let day = row?; @@ -317,6 +345,7 @@ impl Store { fn fill_manual_boundary_slices( &self, provider_account_id: &ProviderAccountId, + source_id: Option<&SourceId>, slices_by_end: &mut BTreeMap>, ) -> Result<()> { let mut days = BTreeSet::new(); @@ -329,7 +358,9 @@ impl Store { return Ok(()); } let day_list = days.into_iter().collect::>(); - for event in self.claude_code_events_on_utc_days(provider_account_id, &day_list)? { + for event in + self.claude_code_events_on_utc_days(provider_account_id, source_id, &day_list)? + { let estimated_cost = event.cost.estimated_micro_usd(); for slices in slices_by_end.values_mut() { for slice in slices.iter_mut() { @@ -348,6 +379,7 @@ impl Store { fn claude_code_events_on_utc_days( &self, provider_account_id: &ProviderAccountId, + source_id: Option<&SourceId>, days: &[NaiveDate], ) -> Result> { if days.is_empty() { @@ -359,11 +391,15 @@ impl Store { let mut sql = String::from( "SELECT payload FROM usage_events WHERE provider = 'claude_code' - AND provider_account_id = ?1 - AND (", + AND provider_account_id = ?1 ", ); let mut bindings: Vec> = vec![Box::new(provider_account_id.0.clone())]; + if let Some(source_id) = source_id { + sql.push_str("AND source_id = ?2 "); + bindings.push(Box::new(source_id.0.clone())); + } + sql.push_str("AND ("); for (index, day) in chunk.iter().enumerate() { if index > 0 { sql.push_str(" OR "); @@ -410,6 +446,16 @@ fn cycle_end_after(anchor_epoch: i64, instant_epoch: i64) -> i64 { } } +fn usage_day_cell(row: &rusqlite::Row<'_>) -> rusqlite::Result { + row.get(0) +} + +fn cycle_overlaps_query(end_epoch: i64, query: &QuotaQuery) -> bool { + let start_epoch = end_epoch - WEEKLY_RESET_PERIOD_SECONDS; + query.from.is_none_or(|from| end_epoch > from.timestamp()) + && query.to.is_none_or(|to| start_epoch <= to.timestamp()) +} + fn first_weekly_end_on_or_after(anchor_epoch: i64, instant_epoch: i64) -> i64 { let phase = anchor_epoch.rem_euclid(WEEKLY_RESET_PERIOD_SECONDS); let instant_mod = instant_epoch.rem_euclid(WEEKLY_RESET_PERIOD_SECONDS); From 06b3107f362841021b77bf149804edc2a644f2c5 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 29 Sep 2026 13:51:25 +0000 Subject: [PATCH 3/4] Keep boundary-day reads inside the timestamp index An OR of many UTC day ranges made bundled SQLite use only the provider and account prefix, so each sync scanned the account's full history. Each day is now its own indexed range, combined with UNION ALL. Co-authored-by: Dmitry Starkov <21260939+starkdmi@users.noreply.github.com> --- .../statsai-store/src/quota/tests/weekly.rs | 24 +++-- crates/statsai-store/src/quota/weekly.rs | 90 +++++++++++-------- 2 files changed, 69 insertions(+), 45 deletions(-) diff --git a/crates/statsai-store/src/quota/tests/weekly.rs b/crates/statsai-store/src/quota/tests/weekly.rs index fb84d34..e998bfa 100644 --- a/crates/statsai-store/src/quota/tests/weekly.rs +++ b/crates/statsai-store/src/quota/tests/weekly.rs @@ -1,3 +1,4 @@ +use super::super::weekly::boundary_day_events_sql; use super::support::*; use super::*; use chrono::TimeZone; @@ -461,16 +462,21 @@ fn usage_day_reads_use_the_provider_account_index() { plan.contains("usage_events_provider_account_started_idx"), "usage day lookup does not use the account index: {plan}" ); - let bounded = query_plan( - &store, - "SELECT payload FROM usage_events - WHERE provider = 'claude_code' - AND provider_account_id = ?1 - AND (started_at >= ?2 AND started_at < ?3)", - ); + for day_count in [1usize, 2, 64] { + let bounded = query_plan(&store, &boundary_day_events_sql(day_count, false)); + assert!( + bounded.contains("usage_events_provider_account_started_idx"), + "{day_count}-day boundary lookup misses the account index: {bounded}" + ); + assert!( + bounded.contains("started_at"), + "{day_count}-day boundary lookup dropped its timestamp bounds: {bounded}" + ); + } + let sourced = query_plan(&store, &boundary_day_events_sql(2, true)); assert!( - bounded.contains("usage_events_provider_account_started_idx"), - "boundary day lookup does not use the account index: {bounded}" + sourced.contains("started_at"), + "source-filtered boundary lookup dropped its timestamp bounds: {sourced}" ); } diff --git a/crates/statsai-store/src/quota/weekly.rs b/crates/statsai-store/src/quota/weekly.rs index e2e0a85..bb69ebf 100644 --- a/crates/statsai-store/src/quota/weekly.rs +++ b/crates/statsai-store/src/quota/weekly.rs @@ -385,45 +385,30 @@ impl Store { if days.is_empty() { return Ok(Vec::new()); } - const CHUNK_SIZE: usize = 64; - let mut events = Vec::new(); - for chunk in days.chunks(CHUNK_SIZE) { - let mut sql = String::from( - "SELECT payload FROM usage_events - WHERE provider = 'claude_code' - AND provider_account_id = ?1 ", - ); - let mut bindings: Vec> = - vec![Box::new(provider_account_id.0.clone())]; + // One indexed range per UTC day. Bundled SQLite drops `started_at` from + // the provider/account index when those ranges are combined with OR, and + // then each batch scans the account's whole history. + let sql = boundary_day_events_sql(days.len(), source_id.is_some()); + let mut statement = self.conn.prepare(&sql)?; + let mut bindings: Vec> = Vec::new(); + for day in days { + let start = day.and_hms_opt(0, 0, 0).expect("utc midnight").and_utc(); + let end = start + Duration::days(1); + bindings.push(Box::new(provider_account_id.0.clone())); if let Some(source_id) = source_id { - sql.push_str("AND source_id = ?2 "); bindings.push(Box::new(source_id.0.clone())); } - sql.push_str("AND ("); - for (index, day) in chunk.iter().enumerate() { - if index > 0 { - sql.push_str(" OR "); - } - let start_index = bindings.len() + 1; - let end_index = start_index + 1; - sql.push_str(&format!( - "(started_at >= ?{start_index} AND started_at < ?{end_index})" - )); - let start = day.and_hms_opt(0, 0, 0).expect("utc midnight").and_utc(); - let end = start + Duration::days(1); - bindings.push(Box::new(start.to_rfc3339())); - bindings.push(Box::new(end.to_rfc3339())); - } - sql.push_str(") ORDER BY started_at, event_id"); - let mut statement = self.conn.prepare(&sql)?; - let params = bindings - .iter() - .map(|value| value.as_ref() as &dyn rusqlite::types::ToSql) - .collect::>(); - let rows = statement.query_map(params.as_slice(), |row| row.get::<_, String>(0))?; - for row in rows { - events.push(serde_json::from_str(&row?)?); - } + bindings.push(Box::new(start.to_rfc3339())); + bindings.push(Box::new(end.to_rfc3339())); + } + let params = bindings + .iter() + .map(|value| value.as_ref() as &dyn rusqlite::types::ToSql) + .collect::>(); + let rows = statement.query_map(params.as_slice(), |row| row.get::<_, String>(0))?; + let mut events = Vec::new(); + for row in rows { + events.push(serde_json::from_str(&row?)?); } Ok(events) } @@ -450,6 +435,39 @@ fn usage_day_cell(row: &rusqlite::Row<'_>) -> rusqlite::Result { row.get(0) } +/// `UNION ALL` of one provider/account/start-time range per UTC day. +/// +/// Each arm keeps `started_at` inside the index constraint. An `OR` of the same +/// ranges does not, on the SQLite build this crate bundles. +pub(crate) fn boundary_day_events_sql(day_count: usize, with_source: bool) -> String { + let mut arms = Vec::with_capacity(day_count); + for index in 0..day_count { + let base = if with_source { + index * 4 + 1 + } else { + index * 3 + 1 + }; + let source_clause = if with_source { + format!("AND source_id = ?{} ", base + 1) + } else { + String::new() + }; + let start_param = if with_source { base + 2 } else { base + 1 }; + let end_param = start_param + 1; + arms.push(format!( + "SELECT payload, started_at, event_id FROM usage_events \ + WHERE provider = 'claude_code' \ + AND provider_account_id = ?{base} \ + {source_clause}AND started_at >= ?{start_param} \ + AND started_at < ?{end_param}" + )); + } + format!( + "SELECT payload FROM ({}) ORDER BY started_at, event_id", + arms.join(" UNION ALL ") + ) +} + fn cycle_overlaps_query(end_epoch: i64, query: &QuotaQuery) -> bool { let start_epoch = end_epoch - WEEKLY_RESET_PERIOD_SECONDS; query.from.is_none_or(|from| end_epoch > from.timestamp()) From 8d86a072b820723848f72f942cdede91431d7dfd Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 29 Sep 2026 14:12:11 +0000 Subject: [PATCH 4/4] Roll back weekly anchors when an account merge fails Anchor moves now share the merge transaction. An overlapping subscription rejects the merge and leaves both accounts' schedules where they were. Summary rewrites and sync-tracking clears join that transaction instead of starting a second one. Co-authored-by: Dmitry Starkov <21260939+starkdmi@users.noreply.github.com> --- crates/statsai-store/src/lib.rs | 8 ++ crates/statsai-store/src/summaries.rs | 16 +--- crates/statsai-store/src/sync_state/mod.rs | 16 +--- crates/statsai/src/cli/account.rs | 85 ++++++++++--------- .../src/cli/tests/account/weekly_reset.rs | 78 +++++++++++++++++ 5 files changed, 135 insertions(+), 68 deletions(-) diff --git a/crates/statsai-store/src/lib.rs b/crates/statsai-store/src/lib.rs index 15d34b7..1ca46a8 100644 --- a/crates/statsai-store/src/lib.rs +++ b/crates/statsai-store/src/lib.rs @@ -535,6 +535,14 @@ impl Store { self.with_immediate_transaction(|| operation(self)) } + /// Runs one account merge's writes in a single transaction. + /// + /// Nested store operations join this transaction. A failure rolls every + /// write back, including a weekly reset anchor moved before a later check. + pub fn apply_account_merge(&self, operation: impl FnOnce(&Self) -> Result) -> Result { + self.with_immediate_transaction(|| operation(self)) + } + /// Applies a complete incoming sync batch in one transaction. /// /// # Errors diff --git a/crates/statsai-store/src/summaries.rs b/crates/statsai-store/src/summaries.rs index 60799ea..93ea3b9 100644 --- a/crates/statsai-store/src/summaries.rs +++ b/crates/statsai-store/src/summaries.rs @@ -112,8 +112,7 @@ impl Store { if summaries.is_empty() { return Ok(0); } - begin_immediate_transaction_with_retry(&self.conn)?; - let result = (|| { + self.with_immediate_transaction(|| { let mut changed = 0u64; for summary in summaries { if self.upsert_summary(summary)? { @@ -121,18 +120,7 @@ impl Store { } } Ok(changed) - })(); - - match result { - Ok(changed) => { - commit_transaction(&self.conn)?; - Ok(changed) - } - Err(error) => { - rollback(&self.conn); - Err(error) - } - } + }) } pub fn summaries_after( diff --git a/crates/statsai-store/src/sync_state/mod.rs b/crates/statsai-store/src/sync_state/mod.rs index 9008f22..a19be32 100644 --- a/crates/statsai-store/src/sync_state/mod.rs +++ b/crates/statsai-store/src/sync_state/mod.rs @@ -186,25 +186,13 @@ impl Store { } pub fn clear_sync_tracking(&self) -> Result<()> { - begin_immediate_transaction_with_retry(&self.conn)?; - let result = (|| { + self.with_immediate_transaction(|| { self.conn.execute("DELETE FROM entity_sync_state", [])?; self.conn .execute("DELETE FROM task_bucket_sync_state", [])?; self.conn.execute("DELETE FROM sync_state", [])?; Ok(()) - })(); - - match result { - Ok(()) => { - commit_transaction(&self.conn)?; - Ok(()) - } - Err(error) => { - rollback(&self.conn); - Err(error) - } - } + }) } /// The cursor upsert without a transaction of its own, so a caller that must diff --git a/crates/statsai/src/cli/account.rs b/crates/statsai/src/cli/account.rs index cfd37e4..80520eb 100644 --- a/crates/statsai/src/cli/account.rs +++ b/crates/statsai/src/cli/account.rs @@ -204,51 +204,56 @@ pub(crate) fn merge_provider_accounts( _ => (0, 0), }; - if !dry_run { - if coalesced_weekly_reset_schedules == 1 { - store.delete_weekly_reset_anchor(provider, &from.provider_account_id)?; - } else if moved_weekly_reset_schedules == 1 { - store.move_weekly_reset_anchor( + let (remaining_references, deleted_source_account) = if dry_run { + ( + account_reference_counts(store, &from.provider_account_id, Some(provider))?, + false, + ) + } else { + store.apply_account_merge(|store| { + if coalesced_weekly_reset_schedules == 1 { + store.delete_weekly_reset_anchor(provider, &from.provider_account_id)?; + } else if moved_weekly_reset_schedules == 1 { + store.move_weekly_reset_anchor( + provider, + &from.provider_account_id, + &to.provider_account_id, + )?; + } + for assignment in &assignments_to_move { + connect_source_to_account( + store, + ConnectSourceToAccountInput { + source_id: &assignment.source_id, + provider_account_id_value: Some(&to.provider_account_id.0), + provider_user_id: None, + email: None, + label: None, + started_at: assignment.started_at, + ended_at: assignment.ended_at, + }, + )?; + } + for subscription in &subscriptions_to_move { + move_subscription_to_account(store, subscription, &to.provider_account_id)?; + } + move_direct_account_records( + store, provider, &from.provider_account_id, &to.provider_account_id, )?; - } - for assignment in &assignments_to_move { - connect_source_to_account( - store, - ConnectSourceToAccountInput { - source_id: &assignment.source_id, - provider_account_id_value: Some(&to.provider_account_id.0), - provider_user_id: None, - email: None, - label: None, - started_at: assignment.started_at, - ended_at: assignment.ended_at, - }, - )?; - } - for subscription in &subscriptions_to_move { - move_subscription_to_account(store, subscription, &to.provider_account_id)?; - } - move_direct_account_records( - store, - provider, - &from.provider_account_id, - &to.provider_account_id, - )?; - } - - let remaining_references = - account_reference_counts(store, &from.provider_account_id, Some(provider))?; - let deleted_source_account = if !dry_run && remaining_references.total() == 0 { - store.delete_account(&from.provider_account_id)? - } else { - false + let remaining_references = + account_reference_counts(store, &from.provider_account_id, Some(provider))?; + let deleted_source_account = if remaining_references.total() == 0 { + store.delete_account(&from.provider_account_id)? + } else { + false + }; + store.clear_sync_tracking()?; + Ok((remaining_references, deleted_source_account)) + })? }; - if !dry_run { - store.clear_sync_tracking()?; - } Ok(AccountMergeReport { provider: provider.to_string(), diff --git a/crates/statsai/src/cli/tests/account/weekly_reset.rs b/crates/statsai/src/cli/tests/account/weekly_reset.rs index 9b0a4c8..2c1b229 100644 --- a/crates/statsai/src/cli/tests/account/weekly_reset.rs +++ b/crates/statsai/src/cli/tests/account/weekly_reset.rs @@ -177,6 +177,84 @@ fn merge_transfers_coalesces_and_rejects_weekly_reset_schedules() { .is_some()); } +#[test] +fn failed_merge_keeps_the_source_weekly_reset_anchor() { + let store = Store::in_memory().expect("store"); + let from = claude_account(&store, "legacy"); + let to = claude_account(&store, "canonical"); + let anchor = at(2026, 1, 8, 18, 0, 0); + let started_at = at(2026, 1, 1, 0, 0, 0); + store + .upsert_weekly_reset_anchor("claude_code", &from.provider_account_id, anchor) + .expect("anchor"); + for (account_id, plan) in [ + (&from.provider_account_id, "Pro"), + (&to.provider_account_id, "Max"), + ] { + store + .upsert_subscription(&Subscription { + schema_version: SUBSCRIPTION_SCHEMA_VERSION.to_string(), + subscription_id: subscription_id("claude_code", account_id, plan, started_at), + provider: "claude_code".to_string(), + provider_account_id: account_id.clone(), + plan_name: plan.to_string(), + price: 2000, + currency: "USD".to_string(), + billing_period: BillingPeriod::Monthly, + paid_at: Some(started_at), + renewal_day: Some(1), + started_at, + ended_at: None, + current_period_ends_at: None, + status: SubscriptionStatus::Active, + record_source: IdentitySource::LocalAuth, + verified_at: Some(started_at), + notes: None, + }) + .expect("subscription"); + } + + let rejected = merge_provider_accounts(&store, "claude_code", "legacy", "canonical", false) + .expect_err("overlap"); + assert!(rejected.to_string().contains("subscription overlaps")); + assert_eq!( + store + .weekly_reset_anchor("claude_code", &from.provider_account_id) + .expect("source schedule") + .expect("anchor restored") + .anchor, + anchor + ); + assert!(store + .weekly_reset_anchor("claude_code", &to.provider_account_id) + .expect("destination schedule") + .is_none()); + + let equivalent = at(2026, 1, 15, 18, 0, 0); + store + .upsert_weekly_reset_anchor("claude_code", &to.provider_account_id, equivalent) + .expect("destination anchor"); + let coalesced = merge_provider_accounts(&store, "claude_code", "legacy", "canonical", false) + .expect_err("coalesce overlap"); + assert!(coalesced.to_string().contains("subscription overlaps")); + assert_eq!( + store + .weekly_reset_anchor("claude_code", &from.provider_account_id) + .expect("source schedule") + .expect("source anchor kept") + .anchor, + anchor + ); + assert_eq!( + store + .weekly_reset_anchor("claude_code", &to.provider_account_id) + .expect("destination schedule") + .expect("destination anchor kept") + .anchor, + equivalent + ); +} + #[test] fn remove_treats_a_weekly_reset_schedule_as_a_reference() { let store = Store::in_memory().expect("store");