diff --git a/crates/statsai-store/src/lib.rs b/crates/statsai-store/src/lib.rs index 3bc5eb7..1ca46a8 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, @@ -532,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/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..e998bfa --- /dev/null +++ b/crates/statsai-store/src/quota/tests/weekly.rs @@ -0,0 +1,676 @@ +use super::super::weekly::boundary_day_events_sql; +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}" + ); + 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!( + sourced.contains("started_at"), + "source-filtered boundary lookup dropped its timestamp bounds: {sourced}" + ); +} + +#[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)); +} + +#[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"); + 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..bb69ebf --- /dev/null +++ b/crates/statsai-store/src/quota/weekly.rs @@ -0,0 +1,517 @@ +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()); + } + // 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, query, 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, + query: &QuotaQuery, + now: DateTime, + ) -> Result> { + 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) + .into_iter() + .filter(|end_epoch| cycle_overlaps_query(*end_epoch, query)) + .collect::>(); + 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, + query.source_id.as_ref(), + &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, + source_id: Option<&SourceId>, + ) -> Result> { + 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 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?; + 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, + source_id: Option<&SourceId>, + 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, 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() { + 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, + source_id: Option<&SourceId>, + days: &[NaiveDate], + ) -> Result> { + if days.is_empty() { + return Ok(Vec::new()); + } + // 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 { + bindings.push(Box::new(source_id.0.clone())); + } + 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) + } + + 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 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()) + && 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); + 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-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 22b46bc..80520eb 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,43 +181,79 @@ pub(crate) fn merge_provider_accounts( .count(); let evidence_to_move = store.account_evidence_reference_counts(provider, &from.provider_account_id)?; - - if !dry_run { - 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)?; + 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) } - move_direct_account_records( - store, - provider, - &from.provider_account_id, - &to.provider_account_id, - )?; - } + (Some(_), None) => (1, 0), + _ => (0, 0), + }; - 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)? + let (remaining_references, deleted_source_account) = if dry_run { + ( + account_reference_counts(store, &from.provider_account_id, Some(provider))?, + false, + ) } else { - false + 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, + )?; + 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(), @@ -225,6 +268,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 +288,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 +297,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 +632,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 +643,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..2c1b229 --- /dev/null +++ b/crates/statsai/src/cli/tests/account/weekly_reset.rs @@ -0,0 +1,355 @@ +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 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"); + 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.