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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 12 additions & 1 deletion crates/statsai-store/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<T>(&self, operation: impl FnOnce(&Self) -> Result<T>) -> Result<T> {
self.with_immediate_transaction(|| operation(self))
}

/// Applies a complete incoming sync batch in one transaction.
///
/// # Errors
Expand Down
8 changes: 7 additions & 1 deletion crates/statsai-store/src/migrations/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)? {
Expand Down Expand Up @@ -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}"),
}
}
Expand Down Expand Up @@ -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
Expand Down
29 changes: 29 additions & 0 deletions crates/statsai-store/src/migrations/v2.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}
2 changes: 1 addition & 1 deletion crates/statsai-store/src/pricing/tests/repricing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
2 changes: 2 additions & 0 deletions crates/statsai-store/src/quota/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions crates/statsai-store/src/quota/sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,7 @@ impl Store {
boundary_slices,
});
}
contributions.extend(self.manual_weekly_contributions(query, device_id, Utc::now())?);
Ok(contributions)
}

Expand Down
1 change: 1 addition & 0 deletions crates/statsai-store/src/quota/tests/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,4 +5,5 @@ mod observations;
mod reconstruct;
mod support;
mod sync;
mod weekly;
mod windows;
Loading
Loading