projectatlas_db/
telemetry.rs

1//! Persist bounded token telemetry in the authoritative project database.
2
3use crate::{DbError, DbResult};
4use projectatlas_core::graph::ProjectInstanceId;
5use projectatlas_core::telemetry::{
6    TokenAccountingTotals, TokenBucketOverview, TokenOverview, TokenTrendPeriod, TokenTrendReport,
7    TokenTrendWindow, UsageDetailAvailability, UsageEvent, UsageInstanceId, UsageInstanceOwner,
8};
9use rusqlite::{Connection, OptionalExtension, params};
10use serde::Serialize;
11use std::collections::BTreeMap;
12use std::time::{Duration, SystemTime, UNIX_EPOCH};
13
14/// Version of the persisted telemetry retention contract.
15const POLICY_VERSION: u32 = 1;
16/// Version of the deterministic logical-byte accounting contract.
17const LOGICAL_BYTE_VERSION: u32 = 1;
18/// Reserved dimension value that owns aggregated overflow detail.
19const OVERFLOW_DIMENSION: &str = "<overflow>";
20/// Persisted state for a runtime that may still accept events.
21const INSTANCE_ACTIVE: &str = "active";
22/// Persisted state for a runtime that completed cleanly.
23const INSTANCE_SEALED: &str = "sealed";
24/// Persisted state for a runtime retired by bounded maintenance.
25const INSTANCE_EXPIRED: &str = "expired";
26/// Dedupe scope whose contribution belongs only to one event.
27const DEDUPE_SCOPE_EVENT: &str = "event";
28/// Number of seconds in a UTC reporting day.
29const SECONDS_PER_DAY: i64 = 86_400;
30/// Stable owner reported when `SQLite` rejects an aggregate addition overflow.
31const AGGREGATE_COUNTER_FIELD: &str = "aggregate_counter";
32/// Domain separator for bounded representations of predecessor telemetry text.
33const LEGACY_TEXT_HASH_DOMAIN: &[u8] = b"projectatlas:legacy-telemetry-text:v1\0";
34
35/// Capacity policy used while a modeled baseline is active.
36#[derive(Clone, Copy, Debug, Eq, PartialEq)]
37enum BaselineAdmission {
38    /// Enforce every runtime baseline row and witness-byte bound.
39    BoundedRuntime,
40    /// Preserve exact predecessor totals before sealing upgrade-owned baselines.
41    SupportedUpgrade,
42}
43
44/// Select the reporting dimension admitted for one event.
45#[derive(Clone, Copy, Debug, Eq, PartialEq)]
46enum DimensionAdmission {
47    /// Derive the normalized dimension from the validated event.
48    Event,
49    /// Route predecessor detail that cannot be represented exactly to overflow.
50    Overflow,
51}
52
53/// Detail that became unavailable while bounding one predecessor event.
54#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
55struct LegacyDetailLoss {
56    /// Raw event fields were replaced by bounded opaque representations.
57    raw: bool,
58    /// At least one reporting dimension could not be retained exactly.
59    dimension: bool,
60    /// The predecessor caller label could not be retained exactly.
61    label: bool,
62}
63
64/// Bind one aggregate value to the common aggregate table parameter order.
65macro_rules! aggregate_params {
66    ($first:expr, $second:expr, $value:expr) => {
67        params![
68            $first,
69            $second,
70            $value.calls,
71            $value.estimated_without,
72            $value.estimated_with,
73            $value.observed_without,
74            $value.observed_with,
75            $value.modeled_without,
76            $value.modeled_with,
77            $value.deduped_modeled_without,
78            $value.deduped_modeled_with,
79            $value.repeated_baselines,
80            $value.observed_file_read_replacements,
81            $value.modeled_file_reads_avoided,
82        ]
83    };
84}
85
86/// Bind one aggregate value to the common daily table parameter order.
87macro_rules! daily_aggregate_params {
88    ($first:expr, $day:expr, $dimension:expr, $value:expr) => {
89        params![
90            $first,
91            $day,
92            $dimension,
93            $value.calls,
94            $value.estimated_without,
95            $value.estimated_with,
96            $value.observed_without,
97            $value.observed_with,
98            $value.modeled_without,
99            $value.modeled_with,
100            $value.deduped_modeled_without,
101            $value.deduped_modeled_with,
102            $value.repeated_baselines,
103            $value.observed_file_read_replacements,
104            $value.modeled_file_reads_avoided,
105        ]
106    };
107}
108
109/// Generate one opaque compatibility identity without panicking on entropy failure.
110pub(crate) fn generate_usage_instance_id() -> DbResult<UsageInstanceId> {
111    let mut bytes = [0_u8; 16];
112    getrandom::fill(&mut bytes).map_err(|_source| DbError::TelemetryIdentityUnavailable)?;
113    UsageInstanceId::from_bytes(bytes).map_err(Into::into)
114}
115
116/// Retention limits for telemetry inside one authoritative database.
117#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
118pub struct TelemetryRetentionPolicy {
119    /// Maximum retained raw events.
120    pub max_raw_rows: usize,
121    /// Maximum logical bytes represented by retained raw events.
122    pub max_raw_logical_bytes: usize,
123    /// Maximum raw-event age in seconds.
124    pub max_raw_age_seconds: u64,
125    /// Maximum normalized reporting dimensions, including overflow.
126    pub max_dimensions: usize,
127    /// Maximum active runtime instances.
128    pub max_active_instances: usize,
129    /// Maximum retained runtime instances.
130    pub max_retained_instances: usize,
131    /// Maximum retained label tombstones.
132    pub max_label_tombstones: usize,
133    /// Maximum retained instance tombstones.
134    pub max_instance_tombstones: usize,
135    /// Maximum caller labels with retained detail state in the database.
136    pub max_retained_labels: usize,
137    /// Maximum baselines admitted for one active instance.
138    pub max_baselines_per_instance: usize,
139    /// Maximum active baseline rows across the database.
140    pub max_active_baseline_rows: usize,
141    /// Maximum logical witness bytes across active baselines.
142    pub max_baseline_logical_bytes: usize,
143    /// Maximum retained daily aggregate rows.
144    pub max_daily_rows: usize,
145    /// Maximum rows touched by one maintenance pass.
146    pub prune_batch_rows: usize,
147    /// Maximum caller-label bytes.
148    pub max_label_bytes: usize,
149    /// Maximum command bytes.
150    pub max_command_bytes: usize,
151    /// Maximum path bytes.
152    pub max_path_bytes: usize,
153    /// Maximum query bytes.
154    pub max_query_bytes: usize,
155    /// Maximum normalized dimension bytes.
156    pub max_dimension_bytes: usize,
157    /// Maximum modeled-baseline witness bytes.
158    pub max_baseline_witness_bytes: usize,
159    /// Active-instance idle timeout in seconds.
160    pub max_active_idle_seconds: u64,
161    /// Permitted future clock skew in seconds.
162    pub future_clock_tolerance_seconds: u64,
163    /// Writes between passive checkpoint attempts.
164    pub checkpoint_write_interval: usize,
165    /// Retained daily trend history.
166    pub retained_trend_days: u64,
167    /// Retained sealed-instance history.
168    pub retained_instance_seconds: u64,
169    /// Retained caller-label history.
170    pub retained_label_seconds: u64,
171    /// Retained tombstone history.
172    pub retained_tombstone_seconds: u64,
173}
174
175impl Default for TelemetryRetentionPolicy {
176    fn default() -> Self {
177        Self {
178            max_raw_rows: 50_000,
179            max_raw_logical_bytes: 64 * 1_024 * 1_024,
180            max_raw_age_seconds: 30 * 24 * 60 * 60,
181            max_dimensions: 128,
182            max_active_instances: 64,
183            max_retained_instances: 4_096,
184            max_label_tombstones: 1_024,
185            max_instance_tombstones: 4_096,
186            max_retained_labels: 256,
187            max_baselines_per_instance: 1_024,
188            max_active_baseline_rows: 16_384,
189            max_baseline_logical_bytes: 16 * 1_024 * 1_024,
190            max_daily_rows: 100_000,
191            prune_batch_rows: 512,
192            max_label_bytes: 128,
193            max_command_bytes: 96,
194            max_path_bytes: 4_096,
195            max_query_bytes: 4_096,
196            max_dimension_bytes: 96,
197            max_baseline_witness_bytes: 1_024,
198            max_active_idle_seconds: 24 * 60 * 60,
199            future_clock_tolerance_seconds: 5 * 60,
200            checkpoint_write_interval: 1_024,
201            retained_trend_days: 400,
202            retained_instance_seconds: 90 * 24 * 60 * 60,
203            retained_label_seconds: 365 * 24 * 60 * 60,
204            retained_tombstone_seconds: 730 * 24 * 60 * 60,
205        }
206    }
207}
208
209impl TelemetryRetentionPolicy {
210    /// Validate that every hard limit can make forward progress.
211    ///
212    /// # Errors
213    ///
214    /// Returns an error for zero or contradictory limits.
215    pub fn validate(self) -> DbResult<Self> {
216        for (field, value) in [
217            ("max_raw_rows", self.max_raw_rows),
218            ("max_raw_logical_bytes", self.max_raw_logical_bytes),
219            ("max_dimensions", self.max_dimensions),
220            ("max_active_instances", self.max_active_instances),
221            ("max_retained_instances", self.max_retained_instances),
222            ("max_label_tombstones", self.max_label_tombstones),
223            ("max_instance_tombstones", self.max_instance_tombstones),
224            ("max_retained_labels", self.max_retained_labels),
225            (
226                "max_baselines_per_instance",
227                self.max_baselines_per_instance,
228            ),
229            ("max_active_baseline_rows", self.max_active_baseline_rows),
230            (
231                "max_baseline_logical_bytes",
232                self.max_baseline_logical_bytes,
233            ),
234            ("max_daily_rows", self.max_daily_rows),
235            ("prune_batch_rows", self.prune_batch_rows),
236            ("max_label_bytes", self.max_label_bytes),
237            ("max_command_bytes", self.max_command_bytes),
238            ("max_path_bytes", self.max_path_bytes),
239            ("max_query_bytes", self.max_query_bytes),
240            ("max_dimension_bytes", self.max_dimension_bytes),
241            (
242                "max_baseline_witness_bytes",
243                self.max_baseline_witness_bytes,
244            ),
245            ("checkpoint_write_interval", self.checkpoint_write_interval),
246        ] {
247            if value == 0 {
248                return Err(DbError::TelemetryLimitInvalid { field, value });
249            }
250        }
251        if self.max_raw_age_seconds == 0
252            || self.max_active_idle_seconds == 0
253            || self.retained_trend_days == 0
254            || self.retained_instance_seconds == 0
255            || self.retained_label_seconds == 0
256            || self.retained_tombstone_seconds == 0
257            || self.max_retained_instances < self.max_active_instances
258            || self.max_daily_rows < 2
259        {
260            return Err(DbError::TelemetryLimitInvalid {
261                field: "telemetry_retention_policy",
262                value: 0,
263            });
264        }
265        Ok(self)
266    }
267}
268
269/// Spill cleanup state for the selected storage design.
270#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
271#[serde(rename_all = "snake_case")]
272pub enum SpillCleanupState {
273    /// Telemetry has no secondary spill database or file owner.
274    NotApplicable,
275}
276
277/// Persisted passive-checkpoint lifecycle state.
278#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
279#[serde(rename_all = "snake_case")]
280pub enum TelemetryCheckpointState {
281    /// The telemetry write threshold has not yet requested a checkpoint.
282    NotDue,
283    /// Every WAL frame visible to the passive attempt was checkpointed.
284    Completed,
285    /// A reader or writer prevented the passive attempt from completing.
286    Busy,
287    /// `SQLite` rejected the passive attempt.
288    Error,
289}
290
291impl TelemetryCheckpointState {
292    /// Return the stable `SQLite` representation.
293    const fn as_str(self) -> &'static str {
294        match self {
295            Self::NotDue => "not_due",
296            Self::Completed => "completed",
297            Self::Busy => "busy",
298            Self::Error => "error",
299        }
300    }
301
302    /// Decode one checked `SQLite` value.
303    fn from_str(value: &str) -> DbResult<Self> {
304        match value {
305            "not_due" => Ok(Self::NotDue),
306            "completed" => Ok(Self::Completed),
307            "busy" => Ok(Self::Busy),
308            "error" => Ok(Self::Error),
309            _ => Err(DbError::InvalidEnum {
310                field: "usage_retention_state.checkpoint_state",
311                value: value.to_string(),
312            }),
313        }
314    }
315}
316
317/// Planner-statistics maintenance policy currently owned by `ProjectAtlas`.
318#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
319#[serde(rename_all = "snake_case")]
320pub enum PlannerStatisticsPolicy {
321    /// No `ProjectAtlas` lifecycle currently runs `ANALYZE` or `PRAGMA optimize`.
322    NotConfigured,
323}
324
325/// Availability of `SQLite` planner statistics in the selected database.
326#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
327#[serde(rename_all = "snake_case")]
328pub enum PlannerStatisticsState {
329    /// No `sqlite_stat1` table has been initialized.
330    NotInitialized,
331    /// `SQLite` planner statistics are present.
332    Available,
333}
334
335/// Content-free retention and page-lifecycle state.
336#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
337pub struct TelemetryRetentionState {
338    /// Retention policy format version.
339    pub policy_version: u32,
340    /// Logical-byte accounting version.
341    pub logical_byte_version: u32,
342    /// Retained raw rows.
343    pub raw_rows: usize,
344    /// Maximum retained raw rows.
345    pub max_raw_rows: usize,
346    /// Maximum retained raw age in seconds.
347    pub max_raw_age_seconds: u64,
348    /// Retained raw logical bytes.
349    pub raw_logical_bytes: usize,
350    /// Maximum retained raw logical bytes.
351    pub max_raw_logical_bytes: usize,
352    /// Active baseline rows.
353    pub baseline_rows: usize,
354    /// Maximum baselines for one active instance.
355    pub max_baselines_per_instance: usize,
356    /// Maximum active baseline rows.
357    pub max_active_baseline_rows: usize,
358    /// Active baseline witness bytes.
359    pub baseline_logical_bytes: usize,
360    /// Maximum active baseline witness bytes.
361    pub max_baseline_logical_bytes: usize,
362    /// Normalized dimension rows.
363    pub dimension_rows: usize,
364    /// Maximum normalized dimensions including overflow.
365    pub max_dimensions: usize,
366    /// Retained runtime-instance rows.
367    pub instance_rows: usize,
368    /// Active runtime-instance rows for the selected project.
369    pub active_instance_rows: usize,
370    /// Maximum active runtime instances per project.
371    pub max_active_instances: usize,
372    /// Maximum retained runtime-instance rows.
373    pub max_retained_instances: usize,
374    /// Retained caller-label state rows in the authoritative database.
375    pub retained_label_rows: usize,
376    /// Maximum retained caller-label state rows in the authoritative database.
377    pub max_retained_labels: usize,
378    /// Retained daily aggregate rows.
379    pub daily_rows: usize,
380    /// Maximum retained daily aggregate rows.
381    pub max_daily_rows: usize,
382    /// Daily trend retention in days.
383    pub retained_trend_days: u64,
384    /// Retained label tombstones.
385    pub label_tombstone_rows: usize,
386    /// Maximum retained label tombstones.
387    pub max_label_tombstones: usize,
388    /// Retained runtime-instance tombstones.
389    pub instance_tombstone_rows: usize,
390    /// Maximum retained runtime-instance tombstones.
391    pub max_instance_tombstones: usize,
392    /// Lifetime pruned raw rows.
393    pub pruned_raw_rows: usize,
394    /// Lifetime pruned runtime instances.
395    pub pruned_instance_rows: usize,
396    /// Lifetime evicted tombstones.
397    pub evicted_tombstones: usize,
398    /// Whether more bounded maintenance is pending.
399    pub maintenance_pending: bool,
400    /// Fixed maximum rows touched by one maintenance category pass.
401    pub prune_batch_rows: usize,
402    /// Writes accumulated since the last passive checkpoint attempt.
403    pub writes_since_checkpoint: usize,
404    /// Writes between passive checkpoint attempts.
405    pub checkpoint_write_interval: usize,
406    /// Epoch of the most recent checkpoint attempt.
407    pub last_checkpoint_epoch: u64,
408    /// Oldest retained raw-event epoch, when detail exists.
409    pub oldest_retained_epoch: Option<u64>,
410    /// Whether anomalous wall-clock movement was observed.
411    pub clock_anomaly: bool,
412    /// Spill cleanup state; always not applicable for the one-database design.
413    pub spill_cleanup: SpillCleanupState,
414    /// Most recent `ProjectAtlas` passive-checkpoint state.
415    pub checkpoint_state: TelemetryCheckpointState,
416    /// Connection-local `SQLite` automatic-checkpoint threshold in WAL pages.
417    pub wal_autocheckpoint_pages: usize,
418    /// Reusable database pages observed live by the reporting connection.
419    pub freelist_pages: usize,
420    /// Total database pages observed live by the reporting connection.
421    pub page_count: usize,
422    /// `SQLite` database page size in bytes.
423    pub page_size: usize,
424    /// Active journal mode reported by `SQLite`.
425    pub journal_mode: String,
426    /// Active synchronous mode reported by `SQLite`.
427    pub synchronous_mode: String,
428    /// Busy timeout observed on the reporting connection in milliseconds.
429    pub connection_busy_timeout_ms: u64,
430    /// Busy timeout required for ordinary read/write connections in milliseconds.
431    pub normal_busy_timeout_ms: u64,
432    /// Busy timeout required for best-effort telemetry writers in milliseconds.
433    pub telemetry_busy_timeout_ms: u64,
434    /// Planner-statistics lifecycle policy.
435    pub statistics_policy: PlannerStatisticsPolicy,
436    /// Current planner-statistics availability.
437    pub statistics_state: PlannerStatisticsState,
438}
439
440/// Normalized aggregate dimension persisted once and referenced by identifier.
441#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
442struct DimensionValues {
443    /// Token-savings evidence bucket.
444    token_savings_bucket: String,
445    /// Token-count provider.
446    provider: String,
447    /// Token-count model.
448    model: String,
449    /// Tokenizer backend.
450    tokenizer_backend: String,
451    /// Accuracy classification.
452    accuracy: String,
453    /// Baseline scenario kind.
454    baseline_kind: String,
455    /// Baseline confidence.
456    confidence: String,
457    /// Accounting layer.
458    accounting_layer: String,
459    /// Estimate method.
460    estimate_method: String,
461    /// Baseline denominator kind.
462    denominator_kind: String,
463    /// Dedupe scope.
464    dedupe_scope: String,
465    /// Whether this is the reserved overflow dimension.
466    overflow: bool,
467}
468
469impl DimensionValues {
470    /// Normalize the reporting dimensions carried by one event.
471    fn from_event(event: &UsageEvent) -> Self {
472        Self {
473            token_savings_bucket: event.token_savings_bucket.clone(),
474            provider: event.provider.clone(),
475            model: event.model.clone(),
476            tokenizer_backend: event.tokenizer_backend.clone(),
477            accuracy: event.accuracy.clone(),
478            baseline_kind: event.baseline_kind.clone(),
479            confidence: event.confidence.clone(),
480            accounting_layer: event.report_accounting_layer().to_string(),
481            estimate_method: event.estimate_method.clone(),
482            denominator_kind: event.report_denominator_kind().to_string(),
483            dedupe_scope: event.report_dedupe_scope().to_string(),
484            overflow: false,
485        }
486    }
487
488    /// Construct the unique reserved overflow dimension.
489    fn overflow() -> Self {
490        Self {
491            token_savings_bucket: OVERFLOW_DIMENSION.to_string(),
492            provider: OVERFLOW_DIMENSION.to_string(),
493            model: OVERFLOW_DIMENSION.to_string(),
494            tokenizer_backend: OVERFLOW_DIMENSION.to_string(),
495            accuracy: OVERFLOW_DIMENSION.to_string(),
496            baseline_kind: OVERFLOW_DIMENSION.to_string(),
497            confidence: OVERFLOW_DIMENSION.to_string(),
498            accounting_layer: OVERFLOW_DIMENSION.to_string(),
499            estimate_method: OVERFLOW_DIMENSION.to_string(),
500            denominator_kind: OVERFLOW_DIMENSION.to_string(),
501            dedupe_scope: OVERFLOW_DIMENSION.to_string(),
502            overflow: true,
503        }
504    }
505}
506
507/// Exact nonnegative aggregate components stored in `SQLite` integers.
508#[derive(Clone, Copy, Debug, Default)]
509struct AggregateCounters {
510    /// Number of represented events.
511    calls: i64,
512    /// Total estimated tokens without `ProjectAtlas`.
513    estimated_without: i64,
514    /// Total estimated tokens emitted with `ProjectAtlas`.
515    estimated_with: i64,
516    /// Observed baseline tokens.
517    observed_without: i64,
518    /// Observed emitted tokens.
519    observed_with: i64,
520    /// Modeled baseline tokens before deduplication.
521    modeled_without: i64,
522    /// Modeled emitted tokens before deduplication.
523    modeled_with: i64,
524    /// Positive component of signed deduped modeled savings.
525    deduped_modeled_without: i64,
526    /// Negative component of signed deduped modeled savings.
527    deduped_modeled_with: i64,
528    /// Repeated modeled-baseline observations.
529    repeated_baselines: i64,
530    /// Observed full-file reads replaced by bounded output.
531    observed_file_read_replacements: i64,
532    /// Modeled file reads avoided by narrowing.
533    modeled_file_reads_avoided: i64,
534}
535
536/// Persisted singleton counter selected for a bounded exact update.
537#[derive(Clone, Copy, Debug)]
538enum RetentionCounter {
539    /// Retained raw event rows.
540    RawRows,
541    /// Logical bytes represented by retained raw events.
542    RawLogicalBytes,
543    /// Active modeled-baseline rows.
544    BaselineRows,
545    /// Witness bytes represented by active modeled baselines.
546    BaselineLogicalBytes,
547    /// Normalized aggregate dimension rows.
548    DimensionRows,
549    /// Retained runtime-instance rows.
550    InstanceRows,
551    /// Retained caller-label rows.
552    LabelRows,
553    /// Retained global and instance daily aggregate rows.
554    DailyRows,
555    /// Retained caller-label tombstones.
556    LabelTombstoneRows,
557    /// Retained runtime-instance tombstones.
558    InstanceTombstoneRows,
559}
560
561impl RetentionCounter {
562    /// Return the fixed query that reads this counter.
563    const fn select_sql(self) -> &'static str {
564        match self {
565            Self::RawRows => "SELECT raw_rows FROM usage_retention_state WHERE singleton = 1",
566            Self::RawLogicalBytes => {
567                "SELECT raw_logical_bytes FROM usage_retention_state WHERE singleton = 1"
568            }
569            Self::BaselineRows => {
570                "SELECT baseline_rows FROM usage_retention_state WHERE singleton = 1"
571            }
572            Self::BaselineLogicalBytes => {
573                "SELECT baseline_logical_bytes FROM usage_retention_state WHERE singleton = 1"
574            }
575            Self::DimensionRows => {
576                "SELECT dimension_rows FROM usage_retention_state WHERE singleton = 1"
577            }
578            Self::InstanceRows => {
579                "SELECT instance_rows FROM usage_retention_state WHERE singleton = 1"
580            }
581            Self::LabelRows => "SELECT label_rows FROM usage_retention_state WHERE singleton = 1",
582            Self::DailyRows => "SELECT daily_rows FROM usage_retention_state WHERE singleton = 1",
583            Self::LabelTombstoneRows => {
584                "SELECT label_tombstone_rows FROM usage_retention_state WHERE singleton = 1"
585            }
586            Self::InstanceTombstoneRows => {
587                "SELECT instance_tombstone_rows FROM usage_retention_state WHERE singleton = 1"
588            }
589        }
590    }
591
592    /// Return the fixed statement that replaces this counter.
593    const fn update_sql(self) -> &'static str {
594        match self {
595            Self::RawRows => "UPDATE usage_retention_state SET raw_rows = ?1 WHERE singleton = 1",
596            Self::RawLogicalBytes => {
597                "UPDATE usage_retention_state SET raw_logical_bytes = ?1 WHERE singleton = 1"
598            }
599            Self::BaselineRows => {
600                "UPDATE usage_retention_state SET baseline_rows = ?1 WHERE singleton = 1"
601            }
602            Self::BaselineLogicalBytes => {
603                "UPDATE usage_retention_state SET baseline_logical_bytes = ?1 WHERE singleton = 1"
604            }
605            Self::DimensionRows => {
606                "UPDATE usage_retention_state SET dimension_rows = ?1 WHERE singleton = 1"
607            }
608            Self::InstanceRows => {
609                "UPDATE usage_retention_state SET instance_rows = ?1 WHERE singleton = 1"
610            }
611            Self::LabelRows => {
612                "UPDATE usage_retention_state SET label_rows = ?1 WHERE singleton = 1"
613            }
614            Self::DailyRows => {
615                "UPDATE usage_retention_state SET daily_rows = ?1 WHERE singleton = 1"
616            }
617            Self::LabelTombstoneRows => {
618                "UPDATE usage_retention_state SET label_tombstone_rows = ?1 WHERE singleton = 1"
619            }
620            Self::InstanceTombstoneRows => {
621                "UPDATE usage_retention_state SET instance_tombstone_rows = ?1 WHERE singleton = 1"
622            }
623        }
624    }
625
626    /// Return the stable field name used by typed diagnostics.
627    const fn field(self) -> &'static str {
628        match self {
629            Self::RawRows => "raw_rows",
630            Self::RawLogicalBytes => "raw_logical_bytes",
631            Self::BaselineRows => "baseline_rows",
632            Self::BaselineLogicalBytes => "baseline_logical_bytes",
633            Self::DimensionRows => "dimension_rows",
634            Self::InstanceRows => "instance_rows",
635            Self::LabelRows => "label_rows",
636            Self::DailyRows => "daily_rows",
637            Self::LabelTombstoneRows => "label_tombstone_rows",
638            Self::InstanceTombstoneRows => "instance_tombstone_rows",
639        }
640    }
641}
642
643impl AggregateCounters {
644    /// Add every component while rejecting integer overflow.
645    fn checked_add(self, other: Self) -> DbResult<Self> {
646        macro_rules! add {
647            ($field:ident) => {
648                self.$field
649                    .checked_add(other.$field)
650                    .ok_or(DbError::TelemetryIntegerOverflow {
651                        field: stringify!($field),
652                    })?
653            };
654        }
655        Ok(Self {
656            calls: add!(calls),
657            estimated_without: add!(estimated_without),
658            estimated_with: add!(estimated_with),
659            observed_without: add!(observed_without),
660            observed_with: add!(observed_with),
661            modeled_without: add!(modeled_without),
662            modeled_with: add!(modeled_with),
663            deduped_modeled_without: add!(deduped_modeled_without),
664            deduped_modeled_with: add!(deduped_modeled_with),
665            repeated_baselines: add!(repeated_baselines),
666            observed_file_read_replacements: add!(observed_file_read_replacements),
667            modeled_file_reads_avoided: add!(modeled_file_reads_avoided),
668        })
669    }
670}
671
672/// Seed telemetry state after fresh-schema creation.
673pub(crate) fn initialize_empty_storage(connection: &Connection) -> DbResult<()> {
674    let policy = TelemetryRetentionPolicy::default().validate()?;
675    ensure_overflow_dimension(connection)?;
676    refresh_retention_state(connection, policy, now_epoch_seconds()?, 0, 0, 0, 0)
677}
678
679/// Convert schema-10 raw usage inside the outer schema transaction.
680pub(crate) fn migrate_legacy_usage(connection: &Connection) -> DbResult<()> {
681    let policy = TelemetryRetentionPolicy::default().validate()?;
682    let (project, _) = crate::project_identity::ensure_project_identity(connection)?;
683    ensure_overflow_dimension(connection)?;
684    let mut statement = connection.prepare(
685        "SELECT session_id, command, path, query,
686                estimated_tokens_without_projectatlas,
687                estimated_tokens_with_projectatlas, estimated_tokens_saved,
688                token_savings_bucket, provider, model, tokenizer_backend,
689                accuracy, baseline_kind, confidence, calculation_trace,
690                accounting_layer, estimate_method, denominator_kind,
691                baseline_identity, baseline_fingerprint, dedupe_scope,
692                created_at, unixepoch(created_at)
693         FROM usage_events_legacy
694         ORDER BY session_id, id",
695    )?;
696    let mut rows = statement.query([])?;
697    let mut previous: Option<UsageInstanceId> = None;
698    while let Some(row) = rows.next()? {
699        let label = row.get::<_, String>(0)?;
700        let created_text = row.get::<_, String>(21)?;
701        let created_at = row
702            .get::<_, Option<i64>>(22)?
703            .ok_or_else(|| DbError::InvalidEnum {
704                field: "usage_events_legacy.created_at",
705                value: created_text,
706            })?;
707        let event = UsageEvent {
708            session_id: label.clone(),
709            command: row.get(1)?,
710            path: row.get(2)?,
711            query: row.get(3)?,
712            estimated_tokens_without_projectatlas: row.get(4)?,
713            estimated_tokens_with_projectatlas: row.get(5)?,
714            estimated_tokens_saved: row.get(6)?,
715            token_savings_bucket: row.get(7)?,
716            provider: row.get(8)?,
717            model: row.get(9)?,
718            tokenizer_backend: row.get(10)?,
719            accuracy: row.get(11)?,
720            baseline_kind: row.get(12)?,
721            confidence: row.get(13)?,
722            calculation_trace: row.get(14)?,
723            accounting_layer: row.get(15)?,
724            estimate_method: row.get(16)?,
725            denominator_kind: row.get(17)?,
726            baseline_identity: row.get(18)?,
727            baseline_fingerprint: row.get(19)?,
728            dedupe_scope: row.get(20)?,
729        };
730        let instance = migrated_instance_id(project, &label)?;
731        if previous != Some(instance) {
732            if let Some(previous) = previous {
733                seal_usage_instance_for_project(connection, project, previous, created_at)?;
734            }
735            previous = Some(instance);
736        }
737        let (event, detail_loss) = normalize_legacy_event(event, policy)?;
738        validate_event(&event, policy)?;
739        record_usage_at(
740            connection,
741            project,
742            instance,
743            UsageInstanceOwner::MigratedLegacy,
744            &event,
745            policy,
746            created_at,
747            false,
748            BaselineAdmission::SupportedUpgrade,
749            if detail_loss.dimension {
750                DimensionAdmission::Overflow
751            } else {
752                DimensionAdmission::Event
753            },
754        )?;
755        mark_legacy_detail_loss(connection, project, instance, &event, detail_loss)?;
756    }
757    drop(rows);
758    drop(statement);
759    if let Some(instance) = previous {
760        seal_usage_instance_for_project(connection, project, instance, now_epoch_seconds()?)?;
761    }
762    reconcile_retention_counters(connection)?;
763    converge_retention(connection, project, policy, now_epoch_seconds()?)
764}
765
766/// Persist one event against the project identity captured by the adapter.
767pub(crate) fn record_usage_for_project(
768    connection: &Connection,
769    project: ProjectInstanceId,
770    instance_id: UsageInstanceId,
771    owner: UsageInstanceOwner,
772    event: &UsageEvent,
773    policy: TelemetryRetentionPolicy,
774    seal_after_record: bool,
775) -> DbResult<()> {
776    crate::project_identity::require_bound_project_identity(connection, project)?;
777    let policy = policy.validate()?;
778    validate_event(event, policy)?;
779    record_usage_at(
780        connection,
781        project,
782        instance_id,
783        owner,
784        event,
785        policy,
786        now_epoch_seconds()?,
787        seal_after_record,
788        BaselineAdmission::BoundedRuntime,
789        DimensionAdmission::Event,
790    )
791}
792
793#[allow(clippy::too_many_arguments)]
794/// Apply one validated event inside its caller-owned write transaction.
795fn record_usage_at(
796    connection: &Connection,
797    project: ProjectInstanceId,
798    instance_id: UsageInstanceId,
799    owner: UsageInstanceOwner,
800    event: &UsageEvent,
801    policy: TelemetryRetentionPolicy,
802    now: i64,
803    seal_after_record: bool,
804    baseline_admission: BaselineAdmission,
805    dimension_admission: DimensionAdmission,
806) -> DbResult<()> {
807    expire_idle_instances(connection, project, instance_id, policy, now)?;
808    let instance_exists = connection.query_row(
809        "SELECT EXISTS(
810             SELECT 1 FROM usage_instances
811             WHERE project_instance_id = ?1 AND runtime_instance_id = ?2
812         )",
813        params![
814            project.as_bytes().as_slice(),
815            instance_id.as_bytes().as_slice()
816        ],
817        |row| row.get::<_, i64>(0),
818    )?;
819    let reserve_instances = usize::from(instance_exists == 0);
820    let (pruned_instances, instance_raw_rows) =
821        prune_instances_once(connection, policy, now, reserve_instances)?;
822    let instance_row_id = ensure_active_instance(
823        connection,
824        project,
825        instance_id,
826        owner,
827        event_label(event),
828        policy,
829        now,
830    )?;
831    ensure_label(connection, project, event_label(event), policy, now)?;
832    let dimension = match dimension_admission {
833        DimensionAdmission::Event => DimensionValues::from_event(event),
834        DimensionAdmission::Overflow => DimensionValues::overflow(),
835    };
836    let dimension_id = ensure_dimension(connection, &dimension, policy)?;
837    let logical_bytes = logical_event_bytes(event, event_label(event))?;
838    let delta = aggregate_delta(
839        connection,
840        instance_row_id,
841        event,
842        policy,
843        baseline_admission,
844    )?;
845    insert_raw_event(
846        connection,
847        instance_row_id,
848        dimension_id,
849        event,
850        now,
851        logical_bytes,
852    )?;
853    prune_daily_once(connection, policy, now)?;
854    apply_aggregates(
855        connection,
856        project,
857        instance_row_id,
858        dimension_id,
859        now,
860        delta,
861        policy,
862    )?;
863    touch_instance(connection, instance_row_id, now, policy)?;
864    if seal_after_record {
865        seal_usage_instance_for_project(connection, project, instance_id, now)?;
866    }
867    let pruned_raw = prune_raw_once(connection, policy, now)?
868        .checked_add(instance_raw_rows)
869        .ok_or(DbError::TelemetryIntegerOverflow {
870            field: "pruned_raw_rows",
871        })?;
872    let evicted_tombstones = prune_tombstones_once(connection, policy, now)?;
873    prune_labels_once(connection, policy, now)?;
874    refresh_retention_state(
875        connection,
876        policy,
877        now,
878        pruned_raw,
879        pruned_instances,
880        evicted_tombstones,
881        1,
882    )?;
883    Ok(())
884}
885
886/// Seal one cleanly completed runtime instance.
887pub(crate) fn seal_usage_instance(
888    connection: &Connection,
889    instance_id: UsageInstanceId,
890) -> DbResult<()> {
891    let project = current_project(connection)?;
892    seal_usage_instance_for_project(connection, project, instance_id, now_epoch_seconds()?)
893}
894
895/// Seal one active instance under an explicitly captured project identity.
896fn seal_usage_instance_for_project(
897    connection: &Connection,
898    project: ProjectInstanceId,
899    instance_id: UsageInstanceId,
900    now: i64,
901) -> DbResult<()> {
902    let changed = connection.execute(
903        "UPDATE usage_instances
904         SET state = ?3, sealed_at_epoch = CASE
905                 WHEN last_seen_at_epoch > ?4 THEN last_seen_at_epoch ELSE ?4 END
906         WHERE project_instance_id = ?1 AND runtime_instance_id = ?2 AND state = ?5",
907        params![
908            project.as_bytes().as_slice(),
909            instance_id.as_bytes().as_slice(),
910            INSTANCE_SEALED,
911            now,
912            INSTANCE_ACTIVE,
913        ],
914    )?;
915    if changed == 0 {
916        return Err(DbError::TelemetryInstanceInactive);
917    }
918    let row_id = connection.query_row(
919        "SELECT instance_row_id FROM usage_instances
920         WHERE project_instance_id = ?1 AND runtime_instance_id = ?2",
921        params![
922            project.as_bytes().as_slice(),
923            instance_id.as_bytes().as_slice()
924        ],
925        |row| row.get::<_, i64>(0),
926    )?;
927    delete_instance_baselines(connection, row_id)?;
928    Ok(())
929}
930
931/// Seal active instances before a project-identity rotation.
932pub(crate) fn seal_project_usage_instances(
933    connection: &Connection,
934    project: ProjectInstanceId,
935) -> DbResult<usize> {
936    crate::project_identity::require_bound_project_identity(connection, project)?;
937    let now = now_epoch_seconds()?;
938    let project_bytes = project.as_bytes();
939    let (baseline_rows, baseline_bytes) = connection.query_row(
940        "SELECT COUNT(*), COALESCE(SUM(b.witness_logical_bytes), 0)
941         FROM usage_instance_baselines AS b
942         JOIN usage_instances AS i USING(instance_row_id)
943         WHERE i.project_instance_id = ?1 AND i.state = ?2",
944        params![project_bytes.as_slice(), INSTANCE_ACTIVE],
945        |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
946    )?;
947    connection.execute(
948        "DELETE FROM usage_instance_baselines
949         WHERE instance_row_id IN (
950             SELECT instance_row_id FROM usage_instances
951             WHERE project_instance_id = ?1 AND state = ?2
952         )",
953        params![project_bytes.as_slice(), INSTANCE_ACTIVE],
954    )?;
955    decrement_retention_counter(
956        connection,
957        RetentionCounter::BaselineRows,
958        count_usize("project_baseline_rows", baseline_rows)?,
959    )?;
960    decrement_retention_counter(
961        connection,
962        RetentionCounter::BaselineLogicalBytes,
963        count_usize("project_baseline_logical_bytes", baseline_bytes)?,
964    )?;
965    connection
966        .execute(
967            "UPDATE usage_instances
968             SET state = ?2,
969                 sealed_at_epoch = CASE
970                     WHEN last_seen_at_epoch > ?3 THEN last_seen_at_epoch ELSE ?3 END
971             WHERE project_instance_id = ?1 AND state = ?4",
972            params![
973                project_bytes.as_slice(),
974                INSTANCE_SEALED,
975                now,
976                INSTANCE_ACTIVE,
977            ],
978        )
979        .map_err(Into::into)
980}
981
982/// Run due post-commit maintenance only for the adapter's captured project identity.
983pub(crate) fn maintain_after_commit_for_project(
984    connection: &Connection,
985    project: ProjectInstanceId,
986    policy: TelemetryRetentionPolicy,
987) -> DbResult<()> {
988    crate::project_identity::require_bound_project_identity(connection, project)?;
989    let policy = policy.validate()?;
990    let writes = connection.query_row(
991        "SELECT writes_since_checkpoint FROM usage_retention_state WHERE singleton = 1",
992        [],
993        |row| row.get::<_, i64>(0),
994    )?;
995    if count_usize("writes_since_checkpoint", writes)? < policy.checkpoint_write_interval {
996        return Ok(());
997    }
998    let result = connection.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
999        Ok((
1000            row.get::<_, i64>(0)?,
1001            row.get::<_, i64>(1)?,
1002            row.get::<_, i64>(2)?,
1003        ))
1004    });
1005    let state = match result {
1006        Ok((0, log_frames, checkpointed_frames)) if checkpointed_frames == log_frames => {
1007            TelemetryCheckpointState::Completed
1008        }
1009        Ok(_) => TelemetryCheckpointState::Busy,
1010        Err(_) => TelemetryCheckpointState::Error,
1011    };
1012    crate::project_identity::require_bound_project_identity(connection, project)?;
1013    let now = now_epoch_seconds()?;
1014    let (raw_rows, raw_bytes, old_raw) = raw_pressure(connection, policy, now)?;
1015    let retention_pending = raw_rows > policy.max_raw_rows
1016        || raw_bytes > policy.max_raw_logical_bytes
1017        || old_raw > 0
1018        || retention_counter(connection, RetentionCounter::InstanceRows)?
1019            > policy.max_retained_instances
1020        || retention_counter(connection, RetentionCounter::LabelRows)? > policy.max_retained_labels
1021        || retention_counter(connection, RetentionCounter::DailyRows)? > policy.max_daily_rows
1022        || retention_counter(connection, RetentionCounter::LabelTombstoneRows)?
1023            > policy.max_label_tombstones
1024        || retention_counter(connection, RetentionCounter::InstanceTombstoneRows)?
1025            > policy.max_instance_tombstones
1026        || aged_maintenance_pending(connection, policy, now)?;
1027    let completed = state == TelemetryCheckpointState::Completed;
1028    connection.execute(
1029        "UPDATE usage_retention_state
1030         SET writes_since_checkpoint = ?1,
1031             last_checkpoint_epoch = ?2,
1032             checkpoint_state = ?3,
1033             maintenance_pending = ?4
1034         WHERE singleton = 1",
1035        params![
1036            if completed { 0 } else { writes },
1037            now,
1038            state.as_str(),
1039            i64::from(!completed || retention_pending),
1040        ],
1041    )?;
1042    Ok(())
1043}
1044
1045/// Return content-free bounded telemetry state.
1046pub(crate) fn retention_state(connection: &Connection) -> DbResult<TelemetryRetentionState> {
1047    let project = current_project(connection)?;
1048    retention_state_for_project(connection, project)
1049}
1050
1051/// Read retention state scoped to the selected project identity.
1052fn retention_state_for_project(
1053    connection: &Connection,
1054    project: ProjectInstanceId,
1055) -> DbResult<TelemetryRetentionState> {
1056    crate::project_identity::require_bound_project_identity(connection, project)?;
1057    let row = connection.query_row(
1058        "SELECT policy_version, logical_byte_version, raw_rows, raw_logical_bytes,
1059                baseline_rows, baseline_logical_bytes, dimension_rows, instance_rows,
1060                daily_rows, label_tombstone_rows, instance_tombstone_rows,
1061                pruned_raw_rows, pruned_instance_rows, evicted_tombstones,
1062                maintenance_pending, clock_anomaly, spill_state, checkpoint_state
1063         FROM usage_retention_state WHERE singleton = 1",
1064        [],
1065        |row| {
1066            Ok((
1067                row.get::<_, i64>(0)?,
1068                row.get::<_, i64>(1)?,
1069                row.get::<_, i64>(2)?,
1070                row.get::<_, i64>(3)?,
1071                row.get::<_, i64>(4)?,
1072                row.get::<_, i64>(5)?,
1073                row.get::<_, i64>(6)?,
1074                row.get::<_, i64>(7)?,
1075                row.get::<_, i64>(8)?,
1076                row.get::<_, i64>(9)?,
1077                row.get::<_, i64>(10)?,
1078                row.get::<_, i64>(11)?,
1079                row.get::<_, i64>(12)?,
1080                row.get::<_, i64>(13)?,
1081                row.get::<_, i64>(14)?,
1082                row.get::<_, i64>(15)?,
1083                row.get::<_, String>(16)?,
1084                row.get::<_, String>(17)?,
1085            ))
1086        },
1087    )?;
1088    if row.16 != "not_applicable" {
1089        return Err(DbError::InvalidEnum {
1090            field: "usage_retention_state.spill_state",
1091            value: row.16,
1092        });
1093    }
1094    let policy = TelemetryRetentionPolicy::default().validate()?;
1095    let lifecycle = connection.query_row(
1096        "SELECT writes_since_checkpoint, last_checkpoint_epoch, oldest_retained_epoch
1097         FROM usage_retention_state WHERE singleton = 1",
1098        [],
1099        |row| {
1100            Ok((
1101                row.get::<_, i64>(0)?,
1102                row.get::<_, i64>(1)?,
1103                row.get::<_, Option<i64>>(2)?,
1104            ))
1105        },
1106    )?;
1107    let active_instances = connection.query_row(
1108        "SELECT COUNT(*) FROM usage_instances
1109         WHERE project_instance_id = ?1 AND state = ?2",
1110        params![project.as_bytes().as_slice(), INSTANCE_ACTIVE],
1111        |row| row.get::<_, i64>(0),
1112    )?;
1113    let retained_labels = retention_counter(connection, RetentionCounter::LabelRows)?;
1114    let journal_mode =
1115        connection.query_row("PRAGMA journal_mode", [], |row| row.get::<_, String>(0))?;
1116    let synchronous = connection.query_row("PRAGMA synchronous", [], |row| row.get::<_, i64>(0))?;
1117    let busy_timeout =
1118        connection.query_row("PRAGMA busy_timeout", [], |row| row.get::<_, i64>(0))?;
1119    let statistics = connection.query_row(
1120        "SELECT EXISTS(
1121             SELECT 1 FROM sqlite_schema WHERE type = 'table' AND name = 'sqlite_stat1'
1122         )",
1123        [],
1124        |row| row.get::<_, i64>(0),
1125    )?;
1126    let checkpoint_state = TelemetryCheckpointState::from_str(&row.17)?;
1127    let wal_autocheckpoint =
1128        connection.query_row("PRAGMA wal_autocheckpoint", [], |row| row.get::<_, i64>(0))?;
1129    let normal_busy_timeout_ms =
1130        duration_millis("normal_busy_timeout_ms", crate::SQLITE_BUSY_TIMEOUT)?;
1131    let telemetry_busy_timeout_ms = duration_millis(
1132        "telemetry_busy_timeout_ms",
1133        crate::SQLITE_TELEMETRY_BUSY_TIMEOUT,
1134    )?;
1135    Ok(TelemetryRetentionState {
1136        policy_version: count_u32("policy_version", row.0)?,
1137        logical_byte_version: count_u32("logical_byte_version", row.1)?,
1138        raw_rows: count_usize("raw_rows", row.2)?,
1139        max_raw_rows: policy.max_raw_rows,
1140        max_raw_age_seconds: policy.max_raw_age_seconds,
1141        raw_logical_bytes: count_usize("raw_logical_bytes", row.3)?,
1142        max_raw_logical_bytes: policy.max_raw_logical_bytes,
1143        baseline_rows: count_usize("baseline_rows", row.4)?,
1144        max_baselines_per_instance: policy.max_baselines_per_instance,
1145        max_active_baseline_rows: policy.max_active_baseline_rows,
1146        baseline_logical_bytes: count_usize("baseline_logical_bytes", row.5)?,
1147        max_baseline_logical_bytes: policy.max_baseline_logical_bytes,
1148        dimension_rows: count_usize("dimension_rows", row.6)?,
1149        max_dimensions: policy.max_dimensions,
1150        instance_rows: count_usize("instance_rows", row.7)?,
1151        active_instance_rows: count_usize("active_instance_rows", active_instances)?,
1152        max_active_instances: policy.max_active_instances,
1153        max_retained_instances: policy.max_retained_instances,
1154        retained_label_rows: retained_labels,
1155        max_retained_labels: policy.max_retained_labels,
1156        daily_rows: count_usize("daily_rows", row.8)?,
1157        max_daily_rows: policy.max_daily_rows,
1158        retained_trend_days: policy.retained_trend_days,
1159        label_tombstone_rows: count_usize("label_tombstone_rows", row.9)?,
1160        max_label_tombstones: policy.max_label_tombstones,
1161        instance_tombstone_rows: count_usize("instance_tombstone_rows", row.10)?,
1162        max_instance_tombstones: policy.max_instance_tombstones,
1163        pruned_raw_rows: count_usize("pruned_raw_rows", row.11)?,
1164        pruned_instance_rows: count_usize("pruned_instance_rows", row.12)?,
1165        evicted_tombstones: count_usize("evicted_tombstones", row.13)?,
1166        maintenance_pending: bool_from_sql("maintenance_pending", row.14)?,
1167        prune_batch_rows: policy.prune_batch_rows,
1168        writes_since_checkpoint: count_usize("writes_since_checkpoint", lifecycle.0)?,
1169        checkpoint_write_interval: policy.checkpoint_write_interval,
1170        last_checkpoint_epoch: count_u64("last_checkpoint_epoch", lifecycle.1)?,
1171        oldest_retained_epoch: lifecycle
1172            .2
1173            .map(|value| count_u64("oldest_retained_epoch", value))
1174            .transpose()?,
1175        clock_anomaly: bool_from_sql("clock_anomaly", row.15)?,
1176        spill_cleanup: SpillCleanupState::NotApplicable,
1177        checkpoint_state,
1178        wal_autocheckpoint_pages: count_usize("wal_autocheckpoint", wal_autocheckpoint)?,
1179        freelist_pages: count_usize(
1180            "freelist_pages",
1181            pragma_count(connection, "freelist_count")?,
1182        )?,
1183        page_count: count_usize("page_count", pragma_count(connection, "page_count")?)?,
1184        page_size: count_usize("page_size", pragma_count(connection, "page_size")?)?,
1185        journal_mode,
1186        synchronous_mode: synchronous_mode(synchronous)?.to_string(),
1187        connection_busy_timeout_ms: count_u64("busy_timeout", busy_timeout)?,
1188        normal_busy_timeout_ms,
1189        telemetry_busy_timeout_ms,
1190        statistics_policy: PlannerStatisticsPolicy::NotConfigured,
1191        statistics_state: if statistics == 0 {
1192            PlannerStatisticsState::NotInitialized
1193        } else {
1194            PlannerStatisticsState::Available
1195        },
1196    })
1197}
1198
1199/// Load retained raw usage events.
1200pub(crate) fn usage_events(
1201    connection: &Connection,
1202    caller_label: Option<&str>,
1203) -> DbResult<Vec<UsageEvent>> {
1204    let project = current_project(connection)?;
1205    usage_events_for_project(connection, project, caller_label)
1206}
1207
1208/// Load retained raw events for one captured project and optional label.
1209fn usage_events_for_project(
1210    connection: &Connection,
1211    project: ProjectInstanceId,
1212    caller_label: Option<&str>,
1213) -> DbResult<Vec<UsageEvent>> {
1214    crate::project_identity::require_bound_project_identity(connection, project)?;
1215    let sql = if caller_label.is_some() {
1216        raw_event_select("AND i.caller_label = ?2")
1217    } else {
1218        raw_event_select("")
1219    };
1220    let mut statement = connection.prepare(&sql)?;
1221    let mut rows = if let Some(label) = caller_label {
1222        statement.query(params![project.as_bytes().as_slice(), label])?
1223    } else {
1224        statement.query([project.as_bytes().as_slice()])?
1225    };
1226    let mut events = Vec::new();
1227    while let Some(row) = rows.next()? {
1228        events.push(map_usage_event(row)?);
1229    }
1230    Ok(events)
1231}
1232
1233/// Build an exact all-time overview from bounded component aggregates.
1234pub(crate) fn token_overview(
1235    connection: &Connection,
1236    caller_label: Option<&str>,
1237) -> DbResult<TokenOverview> {
1238    let project = current_project(connection)?;
1239    token_overview_for_project(connection, project, caller_label)
1240}
1241
1242/// Aggregate all-time token totals for one project and optional label.
1243fn token_overview_for_project(
1244    connection: &Connection,
1245    project: ProjectInstanceId,
1246    caller_label: Option<&str>,
1247) -> DbResult<TokenOverview> {
1248    crate::project_identity::require_bound_project_identity(connection, project)?;
1249    let aggregates = load_overview_aggregates(connection, project, caller_label)?;
1250    let (buckets, totals) = aggregate_report_rows(aggregates)?;
1251    let mut overview = TokenOverview::from_buckets(buckets);
1252    overview.apply_accounting_totals(totals);
1253    overview.set_detail_availability(detail_availability(connection, project, caller_label)?);
1254    Ok(overview)
1255}
1256
1257/// Build trends from bounded daily component aggregates.
1258pub(crate) fn token_trends(
1259    connection: &Connection,
1260    caller_label: Option<&str>,
1261    window: TokenTrendWindow,
1262) -> DbResult<TokenTrendReport> {
1263    let project = current_project(connection)?;
1264    token_trends_for_project(connection, project, caller_label, window)
1265}
1266
1267/// Aggregate bounded daily token trends for one project and optional label.
1268fn token_trends_for_project(
1269    connection: &Connection,
1270    project: ProjectInstanceId,
1271    caller_label: Option<&str>,
1272    window: TokenTrendWindow,
1273) -> DbResult<TokenTrendReport> {
1274    crate::project_identity::require_bound_project_identity(connection, project)?;
1275    let rows = load_daily_aggregates(connection, project, caller_label, window)?;
1276    let mut by_period = BTreeMap::<String, BTreeMap<DimensionValues, AggregateCounters>>::new();
1277    for (period, dimension, counters) in rows {
1278        let entry = by_period
1279            .entry(period)
1280            .or_default()
1281            .entry(dimension)
1282            .or_default();
1283        *entry = entry.checked_add(counters)?;
1284    }
1285    let periods = by_period
1286        .into_iter()
1287        .map(|(period, rows)| {
1288            let buckets = rows
1289                .into_iter()
1290                .map(|(dimension, counters)| bucket_from_counters(dimension, counters))
1291                .collect::<DbResult<Vec<_>>>()?;
1292            Ok(TokenTrendPeriod::from_buckets(period, buckets))
1293        })
1294        .collect::<DbResult<Vec<_>>>()?;
1295    let mut report = TokenTrendReport::new(caller_label.map(str::to_owned), window, periods);
1296    report.set_detail_availability(detail_availability(connection, project, caller_label)?);
1297    Ok(report)
1298}
1299
1300/// Load the required current project identity from the authoritative database.
1301fn current_project(connection: &Connection) -> DbResult<ProjectInstanceId> {
1302    crate::project_identity::load_project_identity(connection)?
1303        .ok_or(DbError::ProjectInstanceIdentityMissing)
1304}
1305
1306/// Derive one deterministic nonzero runtime identity for a legacy caller label.
1307fn migrated_instance_id(
1308    project: ProjectInstanceId,
1309    caller_label: &str,
1310) -> DbResult<UsageInstanceId> {
1311    let mut hasher = blake3::Hasher::new();
1312    hasher.update(b"projectatlas:migrated-usage-instance:v2\0");
1313    hasher.update(&project.as_bytes());
1314    hasher.update(&(caller_label.len() as u64).to_le_bytes());
1315    hasher.update(caller_label.as_bytes());
1316    let mut bytes = [0_u8; 16];
1317    bytes.copy_from_slice(&hasher.finalize().as_bytes()[..16]);
1318    UsageInstanceId::from_bytes(bytes).map_err(Into::into)
1319}
1320
1321/// Bound predecessor text without weakening current-event admission.
1322fn normalize_legacy_event(
1323    mut event: UsageEvent,
1324    policy: TelemetryRetentionPolicy,
1325) -> DbResult<(UsageEvent, LegacyDetailLoss)> {
1326    let baseline_identity = event.effective_baseline_identity().into_owned();
1327    let baseline_fingerprint = event.effective_baseline_fingerprint().into_owned();
1328    let mut loss = LegacyDetailLoss::default();
1329
1330    if !event.session_id.is_empty() {
1331        loss.label = normalize_legacy_required_text(
1332            "session_id",
1333            &mut event.session_id,
1334            policy.max_label_bytes,
1335        )?;
1336    }
1337    loss.raw |=
1338        normalize_legacy_required_text("command", &mut event.command, policy.max_command_bytes)?;
1339    loss.raw |= normalize_legacy_optional_text("path", &mut event.path, policy.max_path_bytes)?;
1340    loss.raw |= normalize_legacy_optional_text("query", &mut event.query, policy.max_query_bytes)?;
1341
1342    for (field, value) in [
1343        ("token_savings_bucket", &mut event.token_savings_bucket),
1344        ("provider", &mut event.provider),
1345        ("model", &mut event.model),
1346        ("tokenizer_backend", &mut event.tokenizer_backend),
1347        ("accuracy", &mut event.accuracy),
1348        ("baseline_kind", &mut event.baseline_kind),
1349        ("confidence", &mut event.confidence),
1350        ("accounting_layer", &mut event.accounting_layer),
1351        ("estimate_method", &mut event.estimate_method),
1352        ("denominator_kind", &mut event.denominator_kind),
1353        ("dedupe_scope", &mut event.dedupe_scope),
1354    ] {
1355        loss.dimension |= normalize_legacy_required_text(field, value, policy.max_dimension_bytes)?;
1356    }
1357
1358    loss.raw |= normalize_legacy_required_text(
1359        "calculation_trace",
1360        &mut event.calculation_trace,
1361        256.min(policy.max_baseline_witness_bytes),
1362    )?;
1363    event.baseline_identity = baseline_identity;
1364    loss.raw |= normalize_legacy_required_text(
1365        "baseline_identity",
1366        &mut event.baseline_identity,
1367        policy.max_baseline_witness_bytes,
1368    )?;
1369    event.baseline_fingerprint = baseline_fingerprint;
1370    loss.raw |= normalize_legacy_required_text(
1371        "baseline_fingerprint",
1372        &mut event.baseline_fingerprint,
1373        256.min(policy.max_baseline_witness_bytes),
1374    )?;
1375    Ok((event, loss))
1376}
1377
1378/// Replace an empty or oversized required predecessor value deterministically.
1379fn normalize_legacy_required_text(
1380    field: &'static str,
1381    value: &mut String,
1382    limit: usize,
1383) -> DbResult<bool> {
1384    if !value.is_empty() && value.len() <= limit {
1385        return Ok(false);
1386    }
1387    *value = legacy_text_token(field, value, limit)?;
1388    Ok(true)
1389}
1390
1391/// Replace an oversized optional predecessor value deterministically.
1392fn normalize_legacy_optional_text(
1393    field: &'static str,
1394    value: &mut Option<String>,
1395    limit: usize,
1396) -> DbResult<bool> {
1397    let Some(text) = value else {
1398        return Ok(false);
1399    };
1400    if text.len() <= limit {
1401        return Ok(false);
1402    }
1403    *text = legacy_text_token(field, text, limit)?;
1404    Ok(true)
1405}
1406
1407/// Produce one bounded opaque value while preserving predecessor equality.
1408fn legacy_text_token(field: &'static str, value: &str, limit: usize) -> DbResult<String> {
1409    let mut hasher = blake3::Hasher::new();
1410    hasher.update(LEGACY_TEXT_HASH_DOMAIN);
1411    for bytes in [field.as_bytes(), value.as_bytes()] {
1412        hasher.update(&(bytes.len() as u64).to_le_bytes());
1413        hasher.update(bytes);
1414    }
1415    let token = format!("legacy:{field}:{}", hasher.finalize().to_hex());
1416    validate_required_text(field, &token, limit)?;
1417    Ok(token)
1418}
1419
1420/// Record which predecessor detail is no longer available exactly.
1421fn mark_legacy_detail_loss(
1422    connection: &Connection,
1423    project: ProjectInstanceId,
1424    instance: UsageInstanceId,
1425    event: &UsageEvent,
1426    loss: LegacyDetailLoss,
1427) -> DbResult<()> {
1428    if loss == LegacyDetailLoss::default() {
1429        return Ok(());
1430    }
1431    if loss.raw {
1432        connection.execute(
1433            "UPDATE usage_instances SET raw_detail_complete = 0
1434             WHERE project_instance_id = ?1 AND runtime_instance_id = ?2
1435               AND raw_detail_complete <> 0",
1436            params![
1437                project.as_bytes().as_slice(),
1438                instance.as_bytes().as_slice()
1439            ],
1440        )?;
1441    }
1442    if (loss.raw || loss.label)
1443        && let Some(label) = event_label(event)
1444    {
1445        connection.execute(
1446            "UPDATE usage_labels SET detail_complete = 0
1447             WHERE project_instance_id = ?1 AND caller_label = ?2
1448               AND detail_complete <> 0",
1449            params![project.as_bytes().as_slice(), label],
1450        )?;
1451    }
1452    connection.execute(
1453        "UPDATE usage_retention_state
1454         SET raw_detail_complete = CASE WHEN ?1 <> 0 THEN 0 ELSE raw_detail_complete END,
1455             dimension_detail_complete =
1456                 CASE WHEN ?2 <> 0 THEN 0 ELSE dimension_detail_complete END,
1457             label_history_complete =
1458                 CASE WHEN ?3 <> 0 THEN 0 ELSE label_history_complete END
1459         WHERE singleton = 1
1460           AND ((?1 <> 0 AND raw_detail_complete <> 0)
1461             OR (?2 <> 0 AND dimension_detail_complete <> 0)
1462             OR (?3 <> 0 AND label_history_complete <> 0))",
1463        params![
1464            i64::from(loss.raw),
1465            i64::from(loss.dimension),
1466            i64::from(loss.label),
1467        ],
1468    )?;
1469    Ok(())
1470}
1471
1472/// Normalize an empty compatibility session label to an absent caller label.
1473fn event_label(event: &UsageEvent) -> Option<&str> {
1474    (!event.session_id.is_empty()).then_some(event.session_id.as_str())
1475}
1476
1477/// Validate every bounded event field before any telemetry mutation.
1478pub(crate) fn validate_event(event: &UsageEvent, policy: TelemetryRetentionPolicy) -> DbResult<()> {
1479    validate_optional_text("session_id", event_label(event), policy.max_label_bytes)?;
1480    validate_required_text("command", &event.command, policy.max_command_bytes)?;
1481    validate_optional_text("path", event.path.as_deref(), policy.max_path_bytes)?;
1482    validate_optional_text("query", event.query.as_deref(), policy.max_query_bytes)?;
1483    for (field, value) in [
1484        ("token_savings_bucket", event.token_savings_bucket.as_str()),
1485        ("provider", event.provider.as_str()),
1486        ("model", event.model.as_str()),
1487        ("tokenizer_backend", event.tokenizer_backend.as_str()),
1488        ("accuracy", event.accuracy.as_str()),
1489        ("baseline_kind", event.baseline_kind.as_str()),
1490        ("confidence", event.confidence.as_str()),
1491        ("accounting_layer", event.report_accounting_layer()),
1492        ("estimate_method", event.estimate_method.as_str()),
1493        ("denominator_kind", event.report_denominator_kind()),
1494        ("dedupe_scope", event.report_dedupe_scope()),
1495    ] {
1496        validate_required_text(field, value, policy.max_dimension_bytes)?;
1497    }
1498    validate_required_text(
1499        "calculation_trace",
1500        &event.calculation_trace,
1501        256.min(policy.max_baseline_witness_bytes),
1502    )?;
1503    let baseline_identity = event.effective_baseline_identity();
1504    let baseline_fingerprint = event.effective_baseline_fingerprint();
1505    validate_required_text(
1506        "baseline_identity",
1507        baseline_identity.as_ref(),
1508        policy.max_baseline_witness_bytes,
1509    )?;
1510    validate_required_text(
1511        "baseline_fingerprint",
1512        baseline_fingerprint.as_ref(),
1513        256.min(policy.max_baseline_witness_bytes),
1514    )?;
1515    let _ = option_usize_to_i64(
1516        "estimated_tokens_without_projectatlas",
1517        event.estimated_tokens_without_projectatlas,
1518    )?;
1519    let _ = option_usize_to_i64(
1520        "estimated_tokens_with_projectatlas",
1521        event.estimated_tokens_with_projectatlas,
1522    )?;
1523    let _ = option_isize_to_i64("estimated_tokens_saved", event.estimated_tokens_saved)?;
1524    Ok(())
1525}
1526
1527/// Require one nonempty UTF-8 field within its byte limit.
1528fn validate_required_text(field: &'static str, value: &str, limit: usize) -> DbResult<()> {
1529    if value.is_empty() || value.len() > limit {
1530        return Err(DbError::TelemetryFieldTooLarge {
1531            field,
1532            bytes: value.len(),
1533            limit,
1534        });
1535    }
1536    Ok(())
1537}
1538
1539/// Validate an optional UTF-8 field against its byte limit.
1540fn validate_optional_text(field: &'static str, value: Option<&str>, limit: usize) -> DbResult<()> {
1541    if let Some(value) = value
1542        && value.len() > limit
1543    {
1544        return Err(DbError::TelemetryFieldTooLarge {
1545            field,
1546            bytes: value.len(),
1547            limit,
1548        });
1549    }
1550    Ok(())
1551}
1552
1553/// Return the identifier of the reserved overflow dimension, inserting it once.
1554fn ensure_overflow_dimension(connection: &Connection) -> DbResult<i64> {
1555    ensure_dimension_unbounded(connection, &DimensionValues::overflow())
1556}
1557
1558/// Resolve a normalized dimension or route it to bounded overflow.
1559fn ensure_dimension(
1560    connection: &Connection,
1561    dimension: &DimensionValues,
1562    policy: TelemetryRetentionPolicy,
1563) -> DbResult<i64> {
1564    if let Some(id) = find_dimension(connection, dimension)? {
1565        return Ok(id);
1566    }
1567    if retention_counter(connection, RetentionCounter::DimensionRows)? >= policy.max_dimensions {
1568        connection.execute(
1569            "UPDATE usage_retention_state SET dimension_detail_complete = 0 WHERE singleton = 1",
1570            [],
1571        )?;
1572        return ensure_overflow_dimension(connection);
1573    }
1574    ensure_dimension_unbounded(connection, dimension)
1575}
1576
1577/// Insert a known-admissible dimension and return its identifier.
1578fn ensure_dimension_unbounded(
1579    connection: &Connection,
1580    dimension: &DimensionValues,
1581) -> DbResult<i64> {
1582    let inserted = connection.execute(
1583        "INSERT OR IGNORE INTO usage_bucket_dimensions(
1584             token_savings_bucket, provider, model, tokenizer_backend,
1585             accuracy, baseline_kind, confidence, accounting_layer,
1586             estimate_method, denominator_kind, dedupe_scope, overflow
1587         ) VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
1588        params![
1589            dimension.token_savings_bucket,
1590            dimension.provider,
1591            dimension.model,
1592            dimension.tokenizer_backend,
1593            dimension.accuracy,
1594            dimension.baseline_kind,
1595            dimension.confidence,
1596            dimension.accounting_layer,
1597            dimension.estimate_method,
1598            dimension.denominator_kind,
1599            dimension.dedupe_scope,
1600            i64::from(dimension.overflow),
1601        ],
1602    )?;
1603    if inserted != 0 {
1604        increment_retention_counter(connection, RetentionCounter::DimensionRows, 1)?;
1605    }
1606    find_dimension(connection, dimension)?.ok_or(DbError::SchemaPostcondition { expected: 11 })
1607}
1608
1609/// Find one exact normalized dimension through its unique key.
1610fn find_dimension(connection: &Connection, dimension: &DimensionValues) -> DbResult<Option<i64>> {
1611    connection
1612        .query_row(
1613            "SELECT dimension_id FROM usage_bucket_dimensions
1614             WHERE token_savings_bucket = ?1 AND provider = ?2 AND model = ?3
1615               AND tokenizer_backend = ?4 AND accuracy = ?5
1616               AND baseline_kind = ?6 AND confidence = ?7
1617               AND accounting_layer = ?8 AND estimate_method = ?9
1618               AND denominator_kind = ?10 AND dedupe_scope = ?11
1619               AND overflow = ?12",
1620            params![
1621                dimension.token_savings_bucket,
1622                dimension.provider,
1623                dimension.model,
1624                dimension.tokenizer_backend,
1625                dimension.accuracy,
1626                dimension.baseline_kind,
1627                dimension.confidence,
1628                dimension.accounting_layer,
1629                dimension.estimate_method,
1630                dimension.denominator_kind,
1631                dimension.dedupe_scope,
1632                i64::from(dimension.overflow),
1633            ],
1634            |row| row.get(0),
1635        )
1636        .optional()
1637        .map_err(Into::into)
1638}
1639
1640/// Resolve or create one active runtime instance without reopening retired state.
1641fn ensure_active_instance(
1642    connection: &Connection,
1643    project: ProjectInstanceId,
1644    runtime: UsageInstanceId,
1645    owner: UsageInstanceOwner,
1646    caller_label: Option<&str>,
1647    policy: TelemetryRetentionPolicy,
1648    now: i64,
1649) -> DbResult<i64> {
1650    let existing = connection
1651        .query_row(
1652            "SELECT instance_row_id, owner, caller_label, state
1653             FROM usage_instances
1654             WHERE project_instance_id = ?1 AND runtime_instance_id = ?2",
1655            params![project.as_bytes().as_slice(), runtime.as_bytes().as_slice()],
1656            |row| {
1657                Ok((
1658                    row.get::<_, i64>(0)?,
1659                    row.get::<_, String>(1)?,
1660                    row.get::<_, Option<String>>(2)?,
1661                    row.get::<_, String>(3)?,
1662                ))
1663            },
1664        )
1665        .optional()?;
1666    if let Some((row_id, stored_owner, stored_label, state)) = existing {
1667        if stored_owner != owner.as_str() || stored_label.as_deref() != caller_label {
1668            return Err(DbError::TelemetryInstanceMismatch);
1669        }
1670        if state != INSTANCE_ACTIVE {
1671            return Err(DbError::TelemetryInstanceInactive);
1672        }
1673        return Ok(row_id);
1674    }
1675    let retired = connection.query_row(
1676        "SELECT EXISTS(
1677             SELECT 1 FROM usage_instance_tombstones
1678             WHERE project_instance_id = ?1 AND runtime_instance_id = ?2
1679         )",
1680        params![project.as_bytes().as_slice(), runtime.as_bytes().as_slice()],
1681        |row| row.get::<_, i64>(0),
1682    )?;
1683    if retired != 0 {
1684        return Err(DbError::TelemetryInstanceInactive);
1685    }
1686    let active = connection.query_row(
1687        "SELECT COUNT(*) FROM usage_instances
1688         WHERE project_instance_id = ?1 AND state = ?2",
1689        params![project.as_bytes().as_slice(), INSTANCE_ACTIVE],
1690        |row| row.get::<_, i64>(0),
1691    )?;
1692    if count_usize("active_usage_instances", active)? >= policy.max_active_instances {
1693        return Err(DbError::TelemetryInstanceCapacity);
1694    }
1695    if retention_counter(connection, RetentionCounter::InstanceRows)?
1696        >= policy.max_retained_instances
1697    {
1698        return Err(DbError::TelemetryInstanceCapacity);
1699    }
1700    connection.execute(
1701        "INSERT INTO usage_instances(
1702             project_instance_id, runtime_instance_id, owner, caller_label, state,
1703             started_at_epoch, last_seen_at_epoch
1704         ) VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?6)",
1705        params![
1706            project.as_bytes().as_slice(),
1707            runtime.as_bytes().as_slice(),
1708            owner.as_str(),
1709            caller_label,
1710            INSTANCE_ACTIVE,
1711            now,
1712        ],
1713    )?;
1714    increment_retention_counter(connection, RetentionCounter::InstanceRows, 1)?;
1715    Ok(connection.last_insert_rowid())
1716}
1717
1718/// Retain or refresh one optional caller label under the database-wide cap.
1719fn ensure_label(
1720    connection: &Connection,
1721    project: ProjectInstanceId,
1722    caller_label: Option<&str>,
1723    policy: TelemetryRetentionPolicy,
1724    now: i64,
1725) -> DbResult<()> {
1726    let Some(label) = caller_label else {
1727        return Ok(());
1728    };
1729    let exists = connection.query_row(
1730        "SELECT EXISTS(
1731             SELECT 1 FROM usage_labels
1732             WHERE project_instance_id = ?1 AND caller_label = ?2
1733         )",
1734        params![project.as_bytes().as_slice(), label],
1735        |row| row.get::<_, i64>(0),
1736    )?;
1737    let prior_history = connection.query_row(
1738        "SELECT EXISTS(
1739             SELECT 1 FROM usage_label_tombstones
1740             WHERE project_instance_id = ?1 AND caller_label = ?2
1741         )",
1742        params![project.as_bytes().as_slice(), label],
1743        |row| row.get::<_, i64>(0),
1744    )?;
1745    if exists == 0
1746        && retention_counter(connection, RetentionCounter::LabelRows)? >= policy.max_retained_labels
1747        && !evict_oldest_label(connection, policy, now)?
1748    {
1749        return Err(DbError::TelemetryInstanceCapacity);
1750    }
1751    let inserted = connection.execute(
1752        "INSERT INTO usage_labels(
1753             project_instance_id, caller_label, last_seen_at_epoch, detail_complete
1754         ) VALUES(?1, ?2, ?3, ?4)
1755         ON CONFLICT(project_instance_id, caller_label) DO UPDATE SET
1756             last_seen_at_epoch = excluded.last_seen_at_epoch",
1757        params![
1758            project.as_bytes().as_slice(),
1759            label,
1760            now,
1761            i64::from(prior_history == 0),
1762        ],
1763    )?;
1764    if inserted != 0 && exists == 0 {
1765        increment_retention_counter(connection, RetentionCounter::LabelRows, 1)?;
1766    }
1767    Ok(())
1768}
1769
1770/// Record that caller-label detail became incomplete or expired.
1771fn upsert_label_tombstone(
1772    connection: &Connection,
1773    project: &[u8],
1774    label: &str,
1775    now: i64,
1776    runtime: Option<&[u8]>,
1777) -> DbResult<()> {
1778    let exists = connection.query_row(
1779        "SELECT EXISTS(
1780             SELECT 1 FROM usage_label_tombstones
1781             WHERE project_instance_id = ?1 AND caller_label = ?2
1782         )",
1783        params![project, label],
1784        |row| row.get::<_, i64>(0),
1785    )?;
1786    connection.execute(
1787        "INSERT INTO usage_label_tombstones(
1788             project_instance_id, caller_label, expired_at_epoch, last_instance_id
1789         ) VALUES(?1, ?2, ?3, ?4)
1790         ON CONFLICT(project_instance_id, caller_label) DO UPDATE SET
1791             expired_at_epoch = excluded.expired_at_epoch,
1792             last_instance_id = excluded.last_instance_id",
1793        params![project, label, now, runtime],
1794    )?;
1795    if exists == 0 {
1796        increment_retention_counter(connection, RetentionCounter::LabelTombstoneRows, 1)?;
1797    }
1798    Ok(())
1799}
1800
1801/// Evict the oldest globally eligible inactive caller-label state.
1802fn evict_oldest_label(
1803    connection: &Connection,
1804    policy: TelemetryRetentionPolicy,
1805    now: i64,
1806) -> DbResult<bool> {
1807    let candidate = connection
1808        .query_row(
1809            "SELECT project_instance_id, caller_label FROM usage_labels
1810             WHERE NOT EXISTS(
1811                   SELECT 1 FROM usage_instances AS i
1812                   WHERE i.project_instance_id = usage_labels.project_instance_id
1813                     AND i.caller_label = usage_labels.caller_label
1814                     AND i.state = ?1
1815               )
1816               AND (
1817                   SELECT COUNT(*) FROM usage_instances AS i
1818                   WHERE i.project_instance_id = usage_labels.project_instance_id
1819                     AND i.caller_label = usage_labels.caller_label
1820               ) <= ?2
1821             ORDER BY last_seen_at_epoch, project_instance_id, caller_label LIMIT 1",
1822            params![
1823                INSTANCE_ACTIVE,
1824                to_i64("prune_batch_rows", policy.prune_batch_rows)?,
1825            ],
1826            |row| Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, String>(1)?)),
1827        )
1828        .optional()?;
1829    let Some((project, label)) = candidate else {
1830        return Ok(false);
1831    };
1832    let runtime = connection
1833        .query_row(
1834            "SELECT runtime_instance_id FROM usage_instances
1835             WHERE project_instance_id = ?1 AND caller_label = ?2
1836             ORDER BY last_seen_at_epoch DESC, instance_row_id DESC LIMIT 1",
1837            params![project, label],
1838            |row| row.get::<_, Vec<u8>>(0),
1839        )
1840        .optional()?;
1841    upsert_label_tombstone(
1842        connection,
1843        project.as_slice(),
1844        &label,
1845        now,
1846        runtime.as_deref(),
1847    )?;
1848    connection.execute(
1849        "UPDATE usage_instances SET caller_label = NULL
1850         WHERE project_instance_id = ?1 AND caller_label = ?2 AND state <> ?3",
1851        params![project, label, INSTANCE_ACTIVE],
1852    )?;
1853    let deleted = connection.execute(
1854        "DELETE FROM usage_labels
1855         WHERE project_instance_id = ?1 AND caller_label = ?2",
1856        params![project, label],
1857    )?;
1858    decrement_retention_counter(connection, RetentionCounter::LabelRows, deleted)?;
1859    connection.execute(
1860        "UPDATE usage_retention_state SET label_history_complete = 0 WHERE singleton = 1",
1861        [],
1862    )?;
1863    Ok(true)
1864}
1865
1866/// Calculate the exact aggregate delta contributed by one event.
1867fn aggregate_delta(
1868    connection: &Connection,
1869    instance_row_id: i64,
1870    event: &UsageEvent,
1871    policy: TelemetryRetentionPolicy,
1872    baseline_admission: BaselineAdmission,
1873) -> DbResult<AggregateCounters> {
1874    let (Some(without_source), Some(with_source)) = (
1875        event.estimated_tokens_without_projectatlas,
1876        event.estimated_tokens_with_projectatlas,
1877    ) else {
1878        return Ok(AggregateCounters::default());
1879    };
1880    let without = to_i64("estimated_tokens_without_projectatlas", without_source)?;
1881    let with = to_i64("estimated_tokens_with_projectatlas", with_source)?;
1882    let observed = event.is_observed();
1883    let modeled = event.is_modeled();
1884    let mut counters = AggregateCounters {
1885        calls: 1,
1886        estimated_without: without,
1887        estimated_with: with,
1888        observed_without: if observed { without } else { 0 },
1889        observed_with: if observed { with } else { 0 },
1890        modeled_without: if modeled { without } else { 0 },
1891        modeled_with: if modeled { with } else { 0 },
1892        observed_file_read_replacements: i64::from(
1893            observed && event.is_observed_file_read_replacement(without_source),
1894        ),
1895        modeled_file_reads_avoided: i64::from(
1896            modeled && event.is_modeled_file_read_avoidance(without_source),
1897        ),
1898        ..AggregateCounters::default()
1899    };
1900    if modeled {
1901        let contribution = if event.report_dedupe_scope() == DEDUPE_SCOPE_EVENT {
1902            without
1903                .checked_sub(with)
1904                .ok_or(DbError::TelemetryIntegerOverflow {
1905                    field: "event_modeled_contribution",
1906                })?
1907        } else {
1908            let (adjustment, repeated) = update_modeled_baseline(
1909                connection,
1910                instance_row_id,
1911                event,
1912                without,
1913                with,
1914                policy,
1915                baseline_admission,
1916            )?;
1917            counters.repeated_baselines = repeated;
1918            adjustment
1919        };
1920        let (positive, negative) = signed_components(contribution)?;
1921        counters.deduped_modeled_without = positive;
1922        counters.deduped_modeled_with = negative;
1923    }
1924    Ok(counters)
1925}
1926
1927/// Update one active modeled baseline and return its signed adjustment.
1928fn update_modeled_baseline(
1929    connection: &Connection,
1930    instance_row_id: i64,
1931    event: &UsageEvent,
1932    without: i64,
1933    with: i64,
1934    policy: TelemetryRetentionPolicy,
1935    baseline_admission: BaselineAdmission,
1936) -> DbResult<(i64, i64)> {
1937    let key = event.modeled_baseline_key();
1938    let identity = event.effective_baseline_identity();
1939    let fingerprint = event.effective_baseline_fingerprint();
1940    let existing = connection
1941        .query_row(
1942            "SELECT baseline_identity, baseline_fingerprint, denominator_kind,
1943                    maximum_without, emitted_with, calls
1944             FROM usage_instance_baselines
1945             WHERE instance_row_id = ?1 AND baseline_key = ?2",
1946            params![instance_row_id, key.as_slice()],
1947            |row| {
1948                Ok((
1949                    row.get::<_, String>(0)?,
1950                    row.get::<_, String>(1)?,
1951                    row.get::<_, String>(2)?,
1952                    row.get::<_, i64>(3)?,
1953                    row.get::<_, i64>(4)?,
1954                    row.get::<_, i64>(5)?,
1955                ))
1956            },
1957        )
1958        .optional()?;
1959    if let Some((stored_identity, stored_fingerprint, denominator, old_max, old_with, calls)) =
1960        existing
1961    {
1962        if stored_identity != identity
1963            || stored_fingerprint != fingerprint
1964            || denominator != event.report_denominator_kind()
1965        {
1966            return Err(DbError::TelemetryBaselineCollision);
1967        }
1968        let previous = old_max
1969            .checked_sub(old_with)
1970            .ok_or(DbError::TelemetryIntegerOverflow {
1971                field: "previous_modeled_baseline",
1972            })?;
1973        let maximum = old_max.max(without);
1974        let emitted = old_with
1975            .checked_add(with)
1976            .ok_or(DbError::TelemetryIntegerOverflow {
1977                field: "modeled_baseline_emitted_with",
1978            })?;
1979        let current = maximum
1980            .checked_sub(emitted)
1981            .ok_or(DbError::TelemetryIntegerOverflow {
1982                field: "current_modeled_baseline",
1983            })?;
1984        let adjustment =
1985            current
1986                .checked_sub(previous)
1987                .ok_or(DbError::TelemetryIntegerOverflow {
1988                    field: "modeled_baseline_adjustment",
1989                })?;
1990        let calls = calls
1991            .checked_add(1)
1992            .ok_or(DbError::TelemetryIntegerOverflow {
1993                field: "modeled_baseline_calls",
1994            })?;
1995        connection.execute(
1996            "UPDATE usage_instance_baselines
1997             SET maximum_without = ?3, emitted_with = ?4, calls = ?5
1998             WHERE instance_row_id = ?1 AND baseline_key = ?2",
1999            params![instance_row_id, key.as_slice(), maximum, emitted, calls],
2000        )?;
2001        return Ok((adjustment, 1));
2002    }
2003    let instance_count = connection.query_row(
2004        "SELECT COUNT(*) FROM usage_instance_baselines WHERE instance_row_id = ?1",
2005        [instance_row_id],
2006        |row| row.get::<_, i64>(0),
2007    )?;
2008    let total_count = retention_counter(connection, RetentionCounter::BaselineRows)?;
2009    let total_bytes = retention_counter(connection, RetentionCounter::BaselineLogicalBytes)?;
2010    let witness_bytes = identity
2011        .len()
2012        .checked_add(fingerprint.len())
2013        .and_then(|value| value.checked_add(event.report_denominator_kind().len()))
2014        .ok_or(DbError::TelemetryIntegerOverflow {
2015            field: "baseline_witness_logical_bytes",
2016        })?;
2017    let projected_bytes =
2018        total_bytes
2019            .checked_add(witness_bytes)
2020            .ok_or(DbError::TelemetryIntegerOverflow {
2021                field: "baseline_logical_bytes",
2022            })?;
2023    let exceeds_runtime_limit = count_usize("instance_baseline_rows", instance_count)?
2024        >= policy.max_baselines_per_instance
2025        || total_count >= policy.max_active_baseline_rows
2026        || projected_bytes > policy.max_baseline_logical_bytes;
2027    if baseline_admission == BaselineAdmission::BoundedRuntime && exceeds_runtime_limit {
2028        return Err(DbError::TelemetryBaselineCapacity);
2029    }
2030    connection.execute(
2031        "INSERT INTO usage_instance_baselines(
2032             instance_row_id, baseline_key, baseline_identity, baseline_fingerprint,
2033             denominator_kind, maximum_without, emitted_with, calls, witness_logical_bytes
2034         ) VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, 1, ?8)",
2035        params![
2036            instance_row_id,
2037            key.as_slice(),
2038            identity,
2039            fingerprint,
2040            event.report_denominator_kind(),
2041            without,
2042            with,
2043            to_i64("baseline_witness_logical_bytes", witness_bytes)?,
2044        ],
2045    )?;
2046    increment_retention_counter(connection, RetentionCounter::BaselineRows, 1)?;
2047    increment_retention_counter(
2048        connection,
2049        RetentionCounter::BaselineLogicalBytes,
2050        witness_bytes,
2051    )?;
2052    Ok((
2053        without
2054            .checked_sub(with)
2055            .ok_or(DbError::TelemetryIntegerOverflow {
2056                field: "initial_modeled_baseline",
2057            })?,
2058        0,
2059    ))
2060}
2061
2062/// Delete one instance's baseline witnesses and decrement exact counters.
2063fn delete_instance_baselines(connection: &Connection, instance_row_id: i64) -> DbResult<usize> {
2064    let (rows, bytes) = connection.query_row(
2065        "SELECT COUNT(*), COALESCE(SUM(witness_logical_bytes), 0)
2066         FROM usage_instance_baselines WHERE instance_row_id = ?1",
2067        [instance_row_id],
2068        |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
2069    )?;
2070    let rows = count_usize("instance_baseline_rows", rows)?;
2071    let bytes = count_usize("instance_baseline_logical_bytes", bytes)?;
2072    if rows != 0 {
2073        connection.execute(
2074            "DELETE FROM usage_instance_baselines WHERE instance_row_id = ?1",
2075            [instance_row_id],
2076        )?;
2077        decrement_retention_counter(connection, RetentionCounter::BaselineRows, rows)?;
2078        decrement_retention_counter(connection, RetentionCounter::BaselineLogicalBytes, bytes)?;
2079    }
2080    Ok(rows)
2081}
2082
2083/// Persist one bounded raw event row.
2084fn insert_raw_event(
2085    connection: &Connection,
2086    instance_row_id: i64,
2087    dimension_id: i64,
2088    event: &UsageEvent,
2089    created_at: i64,
2090    logical_bytes: usize,
2091) -> DbResult<()> {
2092    connection.execute(
2093        "INSERT INTO usage_events(
2094             instance_row_id, dimension_id, command, path, query,
2095             estimated_tokens_without_projectatlas,
2096             estimated_tokens_with_projectatlas, estimated_tokens_saved,
2097             calculation_trace, baseline_identity, baseline_fingerprint,
2098             created_at_epoch, logical_bytes
2099         ) VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)",
2100        params![
2101            instance_row_id,
2102            dimension_id,
2103            event.command,
2104            event.path,
2105            event.query,
2106            option_usize_to_i64(
2107                "estimated_tokens_without_projectatlas",
2108                event.estimated_tokens_without_projectatlas,
2109            )?,
2110            option_usize_to_i64(
2111                "estimated_tokens_with_projectatlas",
2112                event.estimated_tokens_with_projectatlas,
2113            )?,
2114            option_isize_to_i64("estimated_tokens_saved", event.estimated_tokens_saved)?,
2115            event.calculation_trace,
2116            event.effective_baseline_identity(),
2117            event.effective_baseline_fingerprint(),
2118            created_at,
2119            to_i64("raw_logical_bytes", logical_bytes)?,
2120        ],
2121    )?;
2122    increment_retention_counter(connection, RetentionCounter::RawRows, 1)?;
2123    increment_retention_counter(connection, RetentionCounter::RawLogicalBytes, logical_bytes)?;
2124    Ok(())
2125}
2126
2127/// Apply one event delta to all retained aggregate scopes.
2128fn apply_aggregates(
2129    connection: &Connection,
2130    project: ProjectInstanceId,
2131    instance_row_id: i64,
2132    dimension_id: i64,
2133    created_at: i64,
2134    delta: AggregateCounters,
2135    policy: TelemetryRetentionPolicy,
2136) -> DbResult<()> {
2137    if delta.calls == 0 {
2138        return Ok(());
2139    }
2140    upsert_global(connection, project, dimension_id, delta)?;
2141    upsert_instance(connection, instance_row_id, dimension_id, delta)?;
2142    let day = created_at - created_at.rem_euclid(SECONDS_PER_DAY);
2143    prepare_daily_capacity(
2144        connection,
2145        project,
2146        instance_row_id,
2147        day,
2148        dimension_id,
2149        policy,
2150    )?;
2151    upsert_global_daily(connection, project, day, dimension_id, delta)?;
2152    upsert_instance_daily(connection, instance_row_id, day, dimension_id, delta)?;
2153    Ok(())
2154}
2155
2156/// Upsert exact all-time aggregates for one project dimension.
2157fn upsert_global(
2158    connection: &Connection,
2159    project: ProjectInstanceId,
2160    dimension_id: i64,
2161    delta: AggregateCounters,
2162) -> DbResult<()> {
2163    connection.execute(
2164        "INSERT INTO usage_global_aggregates VALUES(
2165             ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)
2166         ON CONFLICT(project_instance_id, dimension_id) DO UPDATE SET
2167             calls = usage_global_aggregates.calls + excluded.calls,
2168             estimated_without = usage_global_aggregates.estimated_without + excluded.estimated_without,
2169             estimated_with = usage_global_aggregates.estimated_with + excluded.estimated_with,
2170             observed_without = usage_global_aggregates.observed_without + excluded.observed_without,
2171             observed_with = usage_global_aggregates.observed_with + excluded.observed_with,
2172             modeled_without = usage_global_aggregates.modeled_without + excluded.modeled_without,
2173             modeled_with = usage_global_aggregates.modeled_with + excluded.modeled_with,
2174             deduped_modeled_without = usage_global_aggregates.deduped_modeled_without + excluded.deduped_modeled_without,
2175             deduped_modeled_with = usage_global_aggregates.deduped_modeled_with + excluded.deduped_modeled_with,
2176             repeated_baselines = usage_global_aggregates.repeated_baselines + excluded.repeated_baselines,
2177             observed_file_read_replacements = usage_global_aggregates.observed_file_read_replacements + excluded.observed_file_read_replacements,
2178             modeled_file_reads_avoided = usage_global_aggregates.modeled_file_reads_avoided + excluded.modeled_file_reads_avoided",
2179        aggregate_params!(project.as_bytes().as_slice(), dimension_id, delta),
2180    )
2181    .map_err(aggregate_write_error)?;
2182    Ok(())
2183}
2184
2185/// Upsert exact all-time aggregates for one runtime dimension.
2186fn upsert_instance(
2187    connection: &Connection,
2188    instance_row_id: i64,
2189    dimension_id: i64,
2190    delta: AggregateCounters,
2191) -> DbResult<()> {
2192    connection.execute(
2193        "INSERT INTO usage_instance_aggregates VALUES(
2194             ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)
2195         ON CONFLICT(instance_row_id, dimension_id) DO UPDATE SET
2196             calls = usage_instance_aggregates.calls + excluded.calls,
2197             estimated_without = usage_instance_aggregates.estimated_without + excluded.estimated_without,
2198             estimated_with = usage_instance_aggregates.estimated_with + excluded.estimated_with,
2199             observed_without = usage_instance_aggregates.observed_without + excluded.observed_without,
2200             observed_with = usage_instance_aggregates.observed_with + excluded.observed_with,
2201             modeled_without = usage_instance_aggregates.modeled_without + excluded.modeled_without,
2202             modeled_with = usage_instance_aggregates.modeled_with + excluded.modeled_with,
2203             deduped_modeled_without = usage_instance_aggregates.deduped_modeled_without + excluded.deduped_modeled_without,
2204             deduped_modeled_with = usage_instance_aggregates.deduped_modeled_with + excluded.deduped_modeled_with,
2205             repeated_baselines = usage_instance_aggregates.repeated_baselines + excluded.repeated_baselines,
2206             observed_file_read_replacements = usage_instance_aggregates.observed_file_read_replacements + excluded.observed_file_read_replacements,
2207             modeled_file_reads_avoided = usage_instance_aggregates.modeled_file_reads_avoided + excluded.modeled_file_reads_avoided",
2208        aggregate_params!(instance_row_id, dimension_id, delta),
2209    )
2210    .map_err(aggregate_write_error)?;
2211    Ok(())
2212}
2213
2214/// Upsert one bounded project-wide daily aggregate row.
2215fn upsert_global_daily(
2216    connection: &Connection,
2217    project: ProjectInstanceId,
2218    day: i64,
2219    dimension_id: i64,
2220    delta: AggregateCounters,
2221) -> DbResult<()> {
2222    connection.execute(
2223        "INSERT INTO usage_daily_aggregates VALUES(
2224             ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15)
2225         ON CONFLICT(project_instance_id, day_epoch, dimension_id) DO UPDATE SET
2226             calls = usage_daily_aggregates.calls + excluded.calls,
2227             estimated_without = usage_daily_aggregates.estimated_without + excluded.estimated_without,
2228             estimated_with = usage_daily_aggregates.estimated_with + excluded.estimated_with,
2229             observed_without = usage_daily_aggregates.observed_without + excluded.observed_without,
2230             observed_with = usage_daily_aggregates.observed_with + excluded.observed_with,
2231             modeled_without = usage_daily_aggregates.modeled_without + excluded.modeled_without,
2232             modeled_with = usage_daily_aggregates.modeled_with + excluded.modeled_with,
2233             deduped_modeled_without = usage_daily_aggregates.deduped_modeled_without + excluded.deduped_modeled_without,
2234             deduped_modeled_with = usage_daily_aggregates.deduped_modeled_with + excluded.deduped_modeled_with,
2235             repeated_baselines = usage_daily_aggregates.repeated_baselines + excluded.repeated_baselines,
2236             observed_file_read_replacements = usage_daily_aggregates.observed_file_read_replacements + excluded.observed_file_read_replacements,
2237             modeled_file_reads_avoided = usage_daily_aggregates.modeled_file_reads_avoided + excluded.modeled_file_reads_avoided",
2238        daily_aggregate_params!(project.as_bytes().as_slice(), day, dimension_id, delta),
2239    )
2240    .map_err(aggregate_write_error)?;
2241    Ok(())
2242}
2243
2244/// Upsert one bounded instance-specific daily aggregate row.
2245fn upsert_instance_daily(
2246    connection: &Connection,
2247    instance_row_id: i64,
2248    day: i64,
2249    dimension_id: i64,
2250    delta: AggregateCounters,
2251) -> DbResult<()> {
2252    connection.execute(
2253        "INSERT INTO usage_instance_daily_aggregates VALUES(
2254             ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15)
2255         ON CONFLICT(instance_row_id, day_epoch, dimension_id) DO UPDATE SET
2256             calls = usage_instance_daily_aggregates.calls + excluded.calls,
2257             estimated_without = usage_instance_daily_aggregates.estimated_without + excluded.estimated_without,
2258             estimated_with = usage_instance_daily_aggregates.estimated_with + excluded.estimated_with,
2259             observed_without = usage_instance_daily_aggregates.observed_without + excluded.observed_without,
2260             observed_with = usage_instance_daily_aggregates.observed_with + excluded.observed_with,
2261             modeled_without = usage_instance_daily_aggregates.modeled_without + excluded.modeled_without,
2262             modeled_with = usage_instance_daily_aggregates.modeled_with + excluded.modeled_with,
2263             deduped_modeled_without = usage_instance_daily_aggregates.deduped_modeled_without + excluded.deduped_modeled_without,
2264             deduped_modeled_with = usage_instance_daily_aggregates.deduped_modeled_with + excluded.deduped_modeled_with,
2265             repeated_baselines = usage_instance_daily_aggregates.repeated_baselines + excluded.repeated_baselines,
2266             observed_file_read_replacements = usage_instance_daily_aggregates.observed_file_read_replacements + excluded.observed_file_read_replacements,
2267             modeled_file_reads_avoided = usage_instance_daily_aggregates.modeled_file_reads_avoided + excluded.modeled_file_reads_avoided",
2268        daily_aggregate_params!(instance_row_id, day, dimension_id, delta),
2269    )
2270    .map_err(aggregate_write_error)?;
2271    Ok(())
2272}
2273
2274/// Preserve the typed telemetry overflow contract for native `SQLite` additions.
2275fn aggregate_write_error(error: rusqlite::Error) -> DbError {
2276    if matches!(
2277        &error,
2278        rusqlite::Error::SqliteFailure(sqlite_error, _)
2279            if sqlite_error.extended_code == rusqlite::ffi::SQLITE_CONSTRAINT_DATATYPE
2280    ) {
2281        DbError::TelemetryIntegerOverflow {
2282            field: AGGREGATE_COUNTER_FIELD,
2283        }
2284    } else {
2285        error.into()
2286    }
2287}
2288
2289/// Identify one daily row that can be evicted without removing the current targets.
2290#[derive(Debug, Eq, PartialEq)]
2291enum DailyEvictionCandidate {
2292    /// One project-wide daily aggregate.
2293    Global {
2294        /// Owning project identity bytes.
2295        project: Vec<u8>,
2296        /// UTC day epoch.
2297        day: i64,
2298        /// Normalized dimension identifier.
2299        dimension_id: i64,
2300    },
2301    /// One runtime-specific daily aggregate.
2302    Instance {
2303        /// Owning runtime row.
2304        instance_row_id: i64,
2305        /// UTC day epoch.
2306        day: i64,
2307        /// Normalized dimension identifier.
2308        dimension_id: i64,
2309    },
2310}
2311
2312impl DailyEvictionCandidate {
2313    /// Return the stable ordering key used to merge both indexed retention heads.
2314    fn ordering_key(&self) -> (i64, u8, Vec<u8>, i64) {
2315        match self {
2316            Self::Global {
2317                project,
2318                day,
2319                dimension_id,
2320            } => (*day, 0, project.clone(), *dimension_id),
2321            Self::Instance {
2322                instance_row_id,
2323                day,
2324                dimension_id,
2325            } => (
2326                *day,
2327                1,
2328                instance_row_id.to_be_bytes().to_vec(),
2329                *dimension_id,
2330            ),
2331        }
2332    }
2333}
2334
2335/// Reserve both current daily rows together, evicting only the exact oldest pressure.
2336fn prepare_daily_capacity(
2337    connection: &Connection,
2338    project: ProjectInstanceId,
2339    instance_row_id: i64,
2340    day: i64,
2341    dimension_id: i64,
2342    policy: TelemetryRetentionPolicy,
2343) -> DbResult<()> {
2344    let global_exists = connection.query_row(
2345        "SELECT EXISTS(
2346             SELECT 1 FROM usage_daily_aggregates
2347             WHERE project_instance_id = ?1 AND day_epoch = ?2 AND dimension_id = ?3
2348         )",
2349        params![project.as_bytes().as_slice(), day, dimension_id],
2350        |row| row.get::<_, i64>(0),
2351    )?;
2352    let instance_exists = connection.query_row(
2353        "SELECT EXISTS(
2354             SELECT 1 FROM usage_instance_daily_aggregates
2355             WHERE instance_row_id = ?1 AND day_epoch = ?2 AND dimension_id = ?3
2356         )",
2357        params![instance_row_id, day, dimension_id],
2358        |row| row.get::<_, i64>(0),
2359    )?;
2360    let required = usize::from(global_exists == 0)
2361        .checked_add(usize::from(instance_exists == 0))
2362        .ok_or(DbError::TelemetryIntegerOverflow {
2363            field: "daily_rows",
2364        })?;
2365    if required == 0 {
2366        return Ok(());
2367    }
2368    let current = retention_counter(connection, RetentionCounter::DailyRows)?;
2369    if current > policy.max_daily_rows {
2370        return Err(DbError::TelemetryLimitInvalid {
2371            field: "daily_rows",
2372            value: current,
2373        });
2374    }
2375    let projected = current
2376        .checked_add(required)
2377        .ok_or(DbError::TelemetryIntegerOverflow {
2378            field: "daily_rows",
2379        })?;
2380    let evictions = projected.saturating_sub(policy.max_daily_rows);
2381    for _ in 0..evictions {
2382        let candidate = oldest_daily_eviction_candidate(
2383            connection,
2384            project,
2385            instance_row_id,
2386            day,
2387            dimension_id,
2388        )?
2389        .ok_or(DbError::TelemetryLimitInvalid {
2390            field: "max_daily_rows",
2391            value: policy.max_daily_rows,
2392        })?;
2393        delete_daily_eviction_candidate(connection, &candidate)?;
2394    }
2395    if evictions != 0 {
2396        decrement_retention_counter(connection, RetentionCounter::DailyRows, evictions)?;
2397        connection.execute(
2398            "UPDATE usage_retention_state SET label_history_complete = 0 WHERE singleton = 1",
2399            [],
2400        )?;
2401    }
2402    increment_retention_counter(connection, RetentionCounter::DailyRows, required)
2403}
2404
2405/// Load the exact oldest evictable row from the two indexed daily tables.
2406fn oldest_daily_eviction_candidate(
2407    connection: &Connection,
2408    project: ProjectInstanceId,
2409    instance_row_id: i64,
2410    day: i64,
2411    dimension_id: i64,
2412) -> DbResult<Option<DailyEvictionCandidate>> {
2413    let global = connection
2414        .query_row(
2415            "SELECT project_instance_id, day_epoch, dimension_id
2416             FROM usage_daily_aggregates
2417             WHERE NOT (
2418                 project_instance_id = ?1 AND day_epoch = ?2 AND dimension_id = ?3
2419             )
2420             ORDER BY day_epoch, project_instance_id, dimension_id LIMIT 1",
2421            params![project.as_bytes().as_slice(), day, dimension_id],
2422            |row| {
2423                Ok(DailyEvictionCandidate::Global {
2424                    project: row.get(0)?,
2425                    day: row.get(1)?,
2426                    dimension_id: row.get(2)?,
2427                })
2428            },
2429        )
2430        .optional()?;
2431    let instance = connection
2432        .query_row(
2433            "SELECT instance_row_id, day_epoch, dimension_id
2434             FROM usage_instance_daily_aggregates
2435             WHERE NOT (
2436                 instance_row_id = ?1 AND day_epoch = ?2 AND dimension_id = ?3
2437             )
2438             ORDER BY day_epoch, instance_row_id, dimension_id LIMIT 1",
2439            params![instance_row_id, day, dimension_id],
2440            |row| {
2441                Ok(DailyEvictionCandidate::Instance {
2442                    instance_row_id: row.get(0)?,
2443                    day: row.get(1)?,
2444                    dimension_id: row.get(2)?,
2445                })
2446            },
2447        )
2448        .optional()?;
2449    Ok(match (global, instance) {
2450        (Some(global), Some(instance)) => {
2451            if global.ordering_key() <= instance.ordering_key() {
2452                Some(global)
2453            } else {
2454                Some(instance)
2455            }
2456        }
2457        (Some(candidate), None) | (None, Some(candidate)) => Some(candidate),
2458        (None, None) => None,
2459    })
2460}
2461
2462/// Delete one selected daily aggregate row.
2463fn delete_daily_eviction_candidate(
2464    connection: &Connection,
2465    candidate: &DailyEvictionCandidate,
2466) -> DbResult<()> {
2467    let deleted = match candidate {
2468        DailyEvictionCandidate::Global {
2469            project,
2470            day,
2471            dimension_id,
2472        } => connection.execute(
2473            "DELETE FROM usage_daily_aggregates
2474             WHERE project_instance_id = ?1 AND day_epoch = ?2 AND dimension_id = ?3",
2475            params![project, day, dimension_id],
2476        )?,
2477        DailyEvictionCandidate::Instance {
2478            instance_row_id,
2479            day,
2480            dimension_id,
2481        } => connection.execute(
2482            "DELETE FROM usage_instance_daily_aggregates
2483             WHERE instance_row_id = ?1 AND day_epoch = ?2 AND dimension_id = ?3",
2484            params![instance_row_id, day, dimension_id],
2485        )?,
2486    };
2487    if deleted != 1 {
2488        return Err(DbError::TelemetryIntegerOverflow {
2489            field: "daily_rows",
2490        });
2491    }
2492    Ok(())
2493}
2494
2495/// Advance one instance's monotonic last-seen state and record clock anomalies.
2496fn touch_instance(
2497    connection: &Connection,
2498    instance_row_id: i64,
2499    now: i64,
2500    policy: TelemetryRetentionPolicy,
2501) -> DbResult<()> {
2502    let previous = connection.query_row(
2503        "SELECT last_seen_at_epoch FROM usage_instances WHERE instance_row_id = ?1",
2504        [instance_row_id],
2505        |row| row.get::<_, i64>(0),
2506    )?;
2507    let tolerance = to_i64(
2508        "future_clock_tolerance_seconds",
2509        policy.future_clock_tolerance_seconds,
2510    )?;
2511    let anomaly = previous
2512        > now
2513            .checked_add(tolerance)
2514            .ok_or(DbError::TelemetryIntegerOverflow {
2515                field: "future_clock_tolerance",
2516            })?;
2517    let observed = previous.max(now);
2518    connection.execute(
2519        "UPDATE usage_instances
2520         SET last_seen_at_epoch = ?2,
2521             clock_anomaly = CASE WHEN ?3 = 1 THEN 1 ELSE clock_anomaly END
2522         WHERE instance_row_id = ?1",
2523        params![instance_row_id, observed, i64::from(anomaly)],
2524    )?;
2525    if anomaly {
2526        connection.execute(
2527            "UPDATE usage_retention_state SET clock_anomaly = 1 WHERE singleton = 1",
2528            [],
2529        )?;
2530    }
2531    Ok(())
2532}
2533
2534/// Expire a bounded page of idle active instances except the current runtime.
2535fn expire_idle_instances(
2536    connection: &Connection,
2537    project: ProjectInstanceId,
2538    current: UsageInstanceId,
2539    policy: TelemetryRetentionPolicy,
2540    now: i64,
2541) -> DbResult<usize> {
2542    let cutoff = epoch_cutoff(
2543        now,
2544        policy.max_active_idle_seconds,
2545        "max_active_idle_seconds",
2546    )?;
2547    let mut statement = connection.prepare_cached(
2548        "SELECT instance_row_id FROM usage_instances
2549         WHERE project_instance_id = ?1 AND state = ?2
2550           AND runtime_instance_id <> ?3 AND last_seen_at_epoch < ?4
2551         ORDER BY last_seen_at_epoch, instance_row_id LIMIT ?5",
2552    )?;
2553    let rows = statement.query_map(
2554        params![
2555            project.as_bytes().as_slice(),
2556            INSTANCE_ACTIVE,
2557            current.as_bytes().as_slice(),
2558            cutoff,
2559            to_i64("prune_batch_rows", policy.prune_batch_rows)?,
2560        ],
2561        |row| row.get::<_, i64>(0),
2562    )?;
2563    let ids = rows.collect::<Result<Vec<_>, _>>()?;
2564    drop(statement);
2565    for row_id in &ids {
2566        connection.execute(
2567            "UPDATE usage_instances
2568             SET state = ?2, sealed_at_epoch = last_seen_at_epoch
2569             WHERE instance_row_id = ?1 AND state = ?3",
2570            params![row_id, INSTANCE_EXPIRED, INSTANCE_ACTIVE],
2571        )?;
2572        delete_instance_baselines(connection, *row_id)?;
2573    }
2574    Ok(ids.len())
2575}
2576
2577/// Remove one bounded page of expired, sealed, or capacity-reserved instances.
2578fn prune_instances_once(
2579    connection: &Connection,
2580    policy: TelemetryRetentionPolicy,
2581    now: i64,
2582    reserve_instances: usize,
2583) -> DbResult<(usize, usize)> {
2584    let active_project = current_project(connection)?.as_bytes();
2585    let projected = retention_counter(connection, RetentionCounter::InstanceRows)?
2586        .checked_add(reserve_instances)
2587        .ok_or(DbError::TelemetryIntegerOverflow {
2588            field: "instance_rows",
2589        })?;
2590    let excess = projected.saturating_sub(policy.max_retained_instances);
2591    let cutoff = epoch_cutoff(
2592        now,
2593        policy.retained_instance_seconds,
2594        "retained_instance_seconds",
2595    )?;
2596    let limit = if excess == 0 {
2597        policy.prune_batch_rows
2598    } else {
2599        excess.min(policy.prune_batch_rows)
2600    };
2601    let mut statement = connection.prepare_cached(
2602        "SELECT instance_row_id, project_instance_id, runtime_instance_id, caller_label
2603         FROM usage_instances
2604         WHERE state IN (?1, ?2) AND (?3 > 0 OR last_seen_at_epoch < ?4)
2605         ORDER BY last_seen_at_epoch, instance_row_id LIMIT ?5",
2606    )?;
2607    let rows = statement.query_map(
2608        params![
2609            INSTANCE_SEALED,
2610            INSTANCE_EXPIRED,
2611            to_i64("excess_instance_rows", excess)?,
2612            cutoff,
2613            to_i64("prune_batch_rows", limit)?,
2614        ],
2615        |row| {
2616            Ok((
2617                row.get::<_, i64>(0)?,
2618                row.get::<_, Vec<u8>>(1)?,
2619                row.get::<_, Vec<u8>>(2)?,
2620                row.get::<_, Option<String>>(3)?,
2621            ))
2622        },
2623    )?;
2624    let candidates = rows.collect::<Result<Vec<_>, _>>()?;
2625    drop(statement);
2626    let mut raw_rows = 0usize;
2627    for (row_id, project, runtime, caller_label) in &candidates {
2628        let (raw, raw_bytes) = connection.query_row(
2629            "SELECT COUNT(*), COALESCE(SUM(logical_bytes), 0)
2630             FROM usage_events WHERE instance_row_id = ?1",
2631            [row_id],
2632            |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
2633        )?;
2634        raw_rows = raw_rows
2635            .checked_add(count_usize("pruned_instance_raw_rows", raw)?)
2636            .ok_or(DbError::TelemetryIntegerOverflow {
2637                field: "pruned_instance_raw_rows",
2638            })?;
2639        let daily_rows = connection.query_row(
2640            "SELECT COUNT(*) FROM usage_instance_daily_aggregates WHERE instance_row_id = ?1",
2641            [row_id],
2642            |row| row.get::<_, i64>(0),
2643        )?;
2644        let tombstone_exists = connection.query_row(
2645            "SELECT EXISTS(
2646                 SELECT 1 FROM usage_instance_tombstones
2647                 WHERE project_instance_id = ?1 AND runtime_instance_id = ?2
2648             )",
2649            params![project, runtime],
2650            |row| row.get::<_, i64>(0),
2651        )?;
2652        connection.execute(
2653            "INSERT INTO usage_instance_tombstones(
2654                 project_instance_id, runtime_instance_id, retired_at_epoch
2655             ) VALUES(?1, ?2, ?3)
2656             ON CONFLICT(project_instance_id, runtime_instance_id) DO UPDATE SET
2657                 retired_at_epoch = excluded.retired_at_epoch",
2658            params![project, runtime, now],
2659        )?;
2660        if tombstone_exists == 0 {
2661            increment_retention_counter(connection, RetentionCounter::InstanceTombstoneRows, 1)?;
2662        }
2663        if let Some(label) = caller_label {
2664            connection.execute(
2665                "UPDATE usage_labels SET detail_complete = 0
2666                 WHERE project_instance_id = ?1 AND caller_label = ?2",
2667                params![project, label],
2668            )?;
2669            upsert_label_tombstone(connection, project, label, now, Some(runtime.as_slice()))?;
2670        }
2671        delete_instance_baselines(connection, *row_id)?;
2672        connection.execute(
2673            "DELETE FROM usage_instances WHERE instance_row_id = ?1",
2674            [row_id],
2675        )?;
2676        decrement_retention_counter(connection, RetentionCounter::InstanceRows, 1)?;
2677        decrement_retention_counter(
2678            connection,
2679            RetentionCounter::RawRows,
2680            count_usize("instance_raw_rows", raw)?,
2681        )?;
2682        decrement_retention_counter(
2683            connection,
2684            RetentionCounter::RawLogicalBytes,
2685            count_usize("instance_raw_logical_bytes", raw_bytes)?,
2686        )?;
2687        decrement_retention_counter(
2688            connection,
2689            RetentionCounter::DailyRows,
2690            count_usize("instance_daily_rows", daily_rows)?,
2691        )?;
2692        if project.as_slice() != active_project.as_slice() {
2693            let project_still_retained = connection.query_row(
2694                "SELECT EXISTS(
2695                     SELECT 1 FROM usage_instances WHERE project_instance_id = ?1 LIMIT 1
2696                 )",
2697                [project],
2698                |row| row.get::<_, i64>(0),
2699            )?;
2700            if project_still_retained == 0 {
2701                connection.execute(
2702                    "DELETE FROM usage_global_aggregates WHERE project_instance_id = ?1",
2703                    [project],
2704                )?;
2705            }
2706        }
2707    }
2708    Ok((candidates.len(), raw_rows))
2709}
2710
2711/// Read raw row, logical-byte, and age pressure through bounded counters and probes.
2712fn raw_pressure(
2713    connection: &Connection,
2714    policy: TelemetryRetentionPolicy,
2715    now: i64,
2716) -> DbResult<(usize, usize, usize)> {
2717    let cutoff = epoch_cutoff(now, policy.max_raw_age_seconds, "max_raw_age_seconds")?;
2718    let old_rows = connection.query_row(
2719        "SELECT EXISTS(
2720             SELECT 1 FROM usage_events
2721             WHERE created_at_epoch < ?1
2722             ORDER BY created_at_epoch, id LIMIT 1
2723         )",
2724        [cutoff],
2725        |row| row.get::<_, i64>(0),
2726    )?;
2727    Ok((
2728        retention_counter(connection, RetentionCounter::RawRows)?,
2729        retention_counter(connection, RetentionCounter::RawLogicalBytes)?,
2730        usize::from(old_rows != 0),
2731    ))
2732}
2733
2734/// Remove one oldest bounded raw-event page when a raw budget is exceeded.
2735fn prune_raw_once(
2736    connection: &Connection,
2737    policy: TelemetryRetentionPolicy,
2738    now: i64,
2739) -> DbResult<usize> {
2740    let (rows, bytes, old_rows) = raw_pressure(connection, policy, now)?;
2741    if rows <= policy.max_raw_rows && bytes <= policy.max_raw_logical_bytes && old_rows == 0 {
2742        return Ok(0);
2743    }
2744    let cutoff = epoch_cutoff(now, policy.max_raw_age_seconds, "max_raw_age_seconds")?;
2745    let mut statement = connection.prepare_cached(
2746        "SELECT created_at_epoch, logical_bytes FROM usage_events
2747         ORDER BY created_at_epoch, id LIMIT ?1",
2748    )?;
2749    let selected = statement.query_map(
2750        [to_i64("prune_batch_rows", policy.prune_batch_rows)?],
2751        |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
2752    )?;
2753    let selected = selected.collect::<Result<Vec<_>, _>>()?;
2754    drop(statement);
2755    let mut remaining_rows = rows;
2756    let mut remaining_bytes = bytes;
2757    let mut delete_count = 0usize;
2758    let mut logical_bytes = 0usize;
2759    for (created_at, event_bytes) in selected {
2760        if created_at >= cutoff
2761            && remaining_rows <= policy.max_raw_rows
2762            && remaining_bytes <= policy.max_raw_logical_bytes
2763        {
2764            break;
2765        }
2766        let event_bytes = count_usize("raw_logical_bytes", event_bytes)?;
2767        remaining_rows = remaining_rows
2768            .checked_sub(1)
2769            .ok_or(DbError::TelemetryIntegerOverflow { field: "raw_rows" })?;
2770        remaining_bytes =
2771            remaining_bytes
2772                .checked_sub(event_bytes)
2773                .ok_or(DbError::TelemetryIntegerOverflow {
2774                    field: "raw_logical_bytes",
2775                })?;
2776        delete_count = delete_count
2777            .checked_add(1)
2778            .ok_or(DbError::TelemetryIntegerOverflow { field: "raw_rows" })?;
2779        logical_bytes =
2780            logical_bytes
2781                .checked_add(event_bytes)
2782                .ok_or(DbError::TelemetryIntegerOverflow {
2783                    field: "raw_logical_bytes",
2784                })?;
2785    }
2786    if delete_count != 0 {
2787        let limit = to_i64("pruned_raw_rows", delete_count)?;
2788        connection.execute(
2789            "UPDATE usage_labels SET detail_complete = 0
2790             WHERE (project_instance_id, caller_label) IN (
2791                 SELECT DISTINCT i.project_instance_id, i.caller_label
2792                 FROM usage_events AS e
2793                 JOIN usage_instances AS i USING(instance_row_id)
2794                 WHERE i.caller_label IS NOT NULL
2795                   AND e.id IN (
2796                       SELECT id FROM usage_events
2797                       ORDER BY created_at_epoch, id LIMIT ?1
2798                   )
2799             )",
2800            [limit],
2801        )?;
2802        connection.execute(
2803            "UPDATE usage_instances SET raw_detail_complete = 0
2804             WHERE instance_row_id IN (
2805                 SELECT instance_row_id FROM usage_events
2806                 ORDER BY created_at_epoch, id LIMIT ?1
2807             )",
2808            [limit],
2809        )?;
2810        let deleted = connection.execute(
2811            "DELETE FROM usage_events WHERE id IN (
2812                 SELECT id FROM usage_events
2813                 ORDER BY created_at_epoch, id LIMIT ?1
2814             )",
2815            [limit],
2816        )?;
2817        if deleted != delete_count {
2818            return Err(DbError::TelemetryIntegerOverflow { field: "raw_rows" });
2819        }
2820        decrement_retention_counter(connection, RetentionCounter::RawRows, delete_count)?;
2821        decrement_retention_counter(connection, RetentionCounter::RawLogicalBytes, logical_bytes)?;
2822        connection.execute(
2823            "UPDATE usage_retention_state SET raw_detail_complete = 0 WHERE singleton = 1",
2824            [],
2825        )?;
2826    }
2827    Ok(delete_count)
2828}
2829
2830/// Remove one bounded page of old or excess daily aggregate rows.
2831fn prune_daily_once(
2832    connection: &Connection,
2833    policy: TelemetryRetentionPolicy,
2834    now: i64,
2835) -> DbResult<usize> {
2836    let days = policy
2837        .retained_trend_days
2838        .checked_mul(SECONDS_PER_DAY as u64)
2839        .ok_or(DbError::TelemetryIntegerOverflow {
2840            field: "retained_trend_seconds",
2841        })?;
2842    let cutoff = epoch_cutoff(now, days, "retained_trend_seconds")?;
2843    let limit = to_i64("prune_batch_rows", policy.prune_batch_rows)?;
2844    let global = connection.execute(
2845        "DELETE FROM usage_daily_aggregates
2846         WHERE (project_instance_id, day_epoch, dimension_id) IN (
2847             SELECT project_instance_id, day_epoch, dimension_id
2848             FROM usage_daily_aggregates WHERE day_epoch < ?1
2849             ORDER BY day_epoch, project_instance_id, dimension_id LIMIT ?2
2850         )",
2851        params![cutoff, limit],
2852    )?;
2853    let remaining = policy.prune_batch_rows.saturating_sub(global);
2854    let instance = if remaining == 0 {
2855        0
2856    } else {
2857        connection.execute(
2858            "DELETE FROM usage_instance_daily_aggregates
2859             WHERE (instance_row_id, day_epoch, dimension_id) IN (
2860                 SELECT instance_row_id, day_epoch, dimension_id
2861                 FROM usage_instance_daily_aggregates WHERE day_epoch < ?1
2862                 ORDER BY day_epoch, instance_row_id, dimension_id LIMIT ?2
2863             )",
2864            params![cutoff, to_i64("daily_prune_remaining", remaining)?],
2865        )?
2866    };
2867    let deleted = global
2868        .checked_add(instance)
2869        .ok_or(DbError::TelemetryIntegerOverflow {
2870            field: "pruned_daily_rows",
2871        })?;
2872    if deleted != 0 {
2873        decrement_retention_counter(connection, RetentionCounter::DailyRows, deleted)?;
2874    }
2875    Ok(deleted)
2876}
2877
2878/// Remove one bounded page of old inactive caller-label state.
2879fn prune_labels_once(
2880    connection: &Connection,
2881    policy: TelemetryRetentionPolicy,
2882    now: i64,
2883) -> DbResult<usize> {
2884    let cutoff = epoch_cutoff(now, policy.retained_label_seconds, "retained_label_seconds")?;
2885    let label = connection
2886        .query_row(
2887            "SELECT project_instance_id, caller_label FROM usage_labels AS labels
2888             WHERE last_seen_at_epoch < ?1
2889               AND NOT EXISTS(
2890                   SELECT 1 FROM usage_instances AS i
2891                   WHERE i.project_instance_id = labels.project_instance_id
2892                     AND i.caller_label = labels.caller_label
2893                     AND i.state = ?2
2894               )
2895               AND (
2896                   SELECT COUNT(*) FROM usage_instances AS i
2897                   WHERE i.project_instance_id = labels.project_instance_id
2898                     AND i.caller_label = labels.caller_label
2899               ) <= ?3
2900             ORDER BY last_seen_at_epoch, project_instance_id, caller_label LIMIT 1",
2901            params![
2902                cutoff,
2903                INSTANCE_ACTIVE,
2904                to_i64("prune_batch_rows", policy.prune_batch_rows)?,
2905            ],
2906            |row| Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, String>(1)?)),
2907        )
2908        .optional()?;
2909    let Some((project, label)) = label else {
2910        return Ok(0);
2911    };
2912    let runtime = connection
2913        .query_row(
2914            "SELECT runtime_instance_id FROM usage_instances
2915             WHERE project_instance_id = ?1 AND caller_label = ?2
2916             ORDER BY last_seen_at_epoch DESC, instance_row_id DESC LIMIT 1",
2917            params![project, label],
2918            |row| row.get::<_, Vec<u8>>(0),
2919        )
2920        .optional()?;
2921    upsert_label_tombstone(connection, &project, &label, now, runtime.as_deref())?;
2922    connection.execute(
2923        "UPDATE usage_instances SET caller_label = NULL
2924         WHERE project_instance_id = ?1 AND caller_label = ?2 AND state <> ?3",
2925        params![project, label, INSTANCE_ACTIVE],
2926    )?;
2927    let deleted = connection.execute(
2928        "DELETE FROM usage_labels WHERE project_instance_id = ?1 AND caller_label = ?2",
2929        params![project, label],
2930    )?;
2931    decrement_retention_counter(connection, RetentionCounter::LabelRows, deleted)?;
2932    connection.execute(
2933        "UPDATE usage_retention_state SET label_history_complete = 0 WHERE singleton = 1",
2934        [],
2935    )?;
2936    Ok(deleted)
2937}
2938
2939/// Remove one bounded page of old or excess label and runtime tombstones.
2940fn prune_tombstones_once(
2941    connection: &Connection,
2942    policy: TelemetryRetentionPolicy,
2943    now: i64,
2944) -> DbResult<usize> {
2945    let cutoff = epoch_cutoff(
2946        now,
2947        policy.retained_tombstone_seconds,
2948        "retained_tombstone_seconds",
2949    )?;
2950    let label_rows = retention_counter(connection, RetentionCounter::LabelTombstoneRows)?;
2951    let mut statement = connection.prepare_cached(
2952        "SELECT expired_at_epoch FROM usage_label_tombstones
2953         ORDER BY expired_at_epoch, project_instance_id, caller_label LIMIT ?1",
2954    )?;
2955    let rows = statement.query_map(
2956        [to_i64("prune_batch_rows", policy.prune_batch_rows)?],
2957        |row| row.get::<_, i64>(0),
2958    )?;
2959    let mut remaining_label_rows = label_rows;
2960    let mut label_limit = 0usize;
2961    for expired_at in rows {
2962        let expired_at = expired_at?;
2963        if expired_at >= cutoff && remaining_label_rows <= policy.max_label_tombstones {
2964            break;
2965        }
2966        remaining_label_rows =
2967            remaining_label_rows
2968                .checked_sub(1)
2969                .ok_or(DbError::TelemetryIntegerOverflow {
2970                    field: "label_tombstone_rows",
2971                })?;
2972        label_limit = label_limit
2973            .checked_add(1)
2974            .ok_or(DbError::TelemetryIntegerOverflow {
2975                field: "label_tombstone_rows",
2976            })?;
2977    }
2978    drop(statement);
2979    let label = connection.execute(
2980        "DELETE FROM usage_label_tombstones
2981         WHERE (project_instance_id, caller_label) IN (
2982             SELECT project_instance_id, caller_label FROM usage_label_tombstones
2983             ORDER BY expired_at_epoch, project_instance_id, caller_label LIMIT ?1
2984         )",
2985        [to_i64("label_tombstone_prune_rows", label_limit)?],
2986    )?;
2987    if label != label_limit {
2988        return Err(DbError::TelemetryIntegerOverflow {
2989            field: "label_tombstone_rows",
2990        });
2991    }
2992    let remaining = policy.prune_batch_rows.saturating_sub(label);
2993    let instance = if remaining == 0 {
2994        0
2995    } else {
2996        let instance_rows = retention_counter(connection, RetentionCounter::InstanceTombstoneRows)?;
2997        let mut statement = connection.prepare_cached(
2998            "SELECT retired_at_epoch FROM usage_instance_tombstones
2999             ORDER BY retired_at_epoch, project_instance_id, runtime_instance_id LIMIT ?1",
3000        )?;
3001        let rows = statement
3002            .query_map([to_i64("tombstone_prune_remaining", remaining)?], |row| {
3003                row.get::<_, i64>(0)
3004            })?;
3005        let mut remaining_instance_rows = instance_rows;
3006        let mut instance_limit = 0usize;
3007        for retired_at in rows {
3008            let retired_at = retired_at?;
3009            if retired_at >= cutoff && remaining_instance_rows <= policy.max_instance_tombstones {
3010                break;
3011            }
3012            remaining_instance_rows = remaining_instance_rows.checked_sub(1).ok_or(
3013                DbError::TelemetryIntegerOverflow {
3014                    field: "instance_tombstone_rows",
3015                },
3016            )?;
3017            instance_limit =
3018                instance_limit
3019                    .checked_add(1)
3020                    .ok_or(DbError::TelemetryIntegerOverflow {
3021                        field: "instance_tombstone_rows",
3022                    })?;
3023        }
3024        drop(statement);
3025        let deleted = connection.execute(
3026            "DELETE FROM usage_instance_tombstones
3027             WHERE (project_instance_id, runtime_instance_id) IN (
3028                 SELECT project_instance_id, runtime_instance_id
3029                 FROM usage_instance_tombstones
3030                 ORDER BY retired_at_epoch, project_instance_id, runtime_instance_id LIMIT ?1
3031             )",
3032            [to_i64("instance_tombstone_prune_rows", instance_limit)?],
3033        )?;
3034        if deleted != instance_limit {
3035            return Err(DbError::TelemetryIntegerOverflow {
3036                field: "instance_tombstone_rows",
3037            });
3038        }
3039        deleted
3040    };
3041    if label != 0 {
3042        decrement_retention_counter(connection, RetentionCounter::LabelTombstoneRows, label)?;
3043    }
3044    if instance != 0 {
3045        decrement_retention_counter(
3046            connection,
3047            RetentionCounter::InstanceTombstoneRows,
3048            instance,
3049        )?;
3050    }
3051    label
3052        .checked_add(instance)
3053        .ok_or(DbError::TelemetryIntegerOverflow {
3054            field: "evicted_tombstones",
3055        })
3056}
3057
3058/// Rebuild persisted counter state once during the supported telemetry migration.
3059/// Rebuild persisted retention counters during a supported schema upgrade.
3060fn reconcile_retention_counters(connection: &Connection) -> DbResult<()> {
3061    let counts = connection.query_row(
3062        "SELECT
3063             (SELECT COUNT(*) FROM usage_events),
3064             (SELECT COALESCE(SUM(logical_bytes), 0) FROM usage_events),
3065             (SELECT COUNT(*) FROM usage_instance_baselines),
3066             (SELECT COALESCE(SUM(witness_logical_bytes), 0)
3067                FROM usage_instance_baselines),
3068             (SELECT COUNT(*) FROM usage_bucket_dimensions),
3069             (SELECT COUNT(*) FROM usage_instances),
3070             (SELECT COUNT(*) FROM usage_labels),
3071             (SELECT COUNT(*) FROM usage_daily_aggregates)
3072                + (SELECT COUNT(*) FROM usage_instance_daily_aggregates),
3073             (SELECT COUNT(*) FROM usage_label_tombstones),
3074             (SELECT COUNT(*) FROM usage_instance_tombstones)",
3075        [],
3076        |row| {
3077            Ok((
3078                row.get::<_, i64>(0)?,
3079                row.get::<_, i64>(1)?,
3080                row.get::<_, i64>(2)?,
3081                row.get::<_, i64>(3)?,
3082                row.get::<_, i64>(4)?,
3083                row.get::<_, i64>(5)?,
3084                row.get::<_, i64>(6)?,
3085                row.get::<_, i64>(7)?,
3086                row.get::<_, i64>(8)?,
3087                row.get::<_, i64>(9)?,
3088            ))
3089        },
3090    )?;
3091    for (field, value) in [
3092        ("raw_rows", counts.0),
3093        ("raw_logical_bytes", counts.1),
3094        ("baseline_rows", counts.2),
3095        ("baseline_logical_bytes", counts.3),
3096        ("dimension_rows", counts.4),
3097        ("instance_rows", counts.5),
3098        ("label_rows", counts.6),
3099        ("daily_rows", counts.7),
3100        ("label_tombstone_rows", counts.8),
3101        ("instance_tombstone_rows", counts.9),
3102    ] {
3103        let _ = count_usize(field, value)?;
3104    }
3105    connection.execute(
3106        "UPDATE usage_retention_state
3107         SET raw_rows = ?1,
3108             raw_logical_bytes = ?2,
3109             baseline_rows = ?3,
3110             baseline_logical_bytes = ?4,
3111             dimension_rows = ?5,
3112             instance_rows = ?6,
3113             label_rows = ?7,
3114             daily_rows = ?8,
3115             label_tombstone_rows = ?9,
3116             instance_tombstone_rows = ?10
3117         WHERE singleton = 1",
3118        params![
3119            counts.0, counts.1, counts.2, counts.3, counts.4, counts.5, counts.6, counts.7,
3120            counts.8, counts.9,
3121        ],
3122    )?;
3123    Ok(())
3124}
3125
3126/// Probe indexed retention paths for additional age-based maintenance work.
3127fn aged_maintenance_pending(
3128    connection: &Connection,
3129    policy: TelemetryRetentionPolicy,
3130    now: i64,
3131) -> DbResult<bool> {
3132    let instance_cutoff = epoch_cutoff(
3133        now,
3134        policy.retained_instance_seconds,
3135        "retained_instance_seconds",
3136    )?;
3137    let trend_seconds = policy
3138        .retained_trend_days
3139        .checked_mul(SECONDS_PER_DAY as u64)
3140        .ok_or(DbError::TelemetryIntegerOverflow {
3141            field: "retained_trend_seconds",
3142        })?;
3143    let daily_cutoff = epoch_cutoff(now, trend_seconds, "retained_trend_seconds")?;
3144    let label_cutoff = epoch_cutoff(now, policy.retained_label_seconds, "retained_label_seconds")?;
3145    let tombstone_cutoff = epoch_cutoff(
3146        now,
3147        policy.retained_tombstone_seconds,
3148        "retained_tombstone_seconds",
3149    )?;
3150    let pending = connection.query_row(
3151        "SELECT
3152             EXISTS(
3153                 SELECT 1 FROM usage_instances
3154                 WHERE state IN (?1, ?2) AND last_seen_at_epoch < ?3 LIMIT 1
3155             ) OR EXISTS(
3156                 SELECT 1 FROM usage_daily_aggregates
3157                 WHERE day_epoch < ?4 LIMIT 1
3158             ) OR EXISTS(
3159                 SELECT 1 FROM usage_instance_daily_aggregates
3160                 WHERE day_epoch < ?4 LIMIT 1
3161             ) OR EXISTS(
3162                 SELECT 1 FROM usage_labels AS labels
3163                 WHERE last_seen_at_epoch < ?5
3164                   AND NOT EXISTS(
3165                       SELECT 1 FROM usage_instances AS i
3166                       WHERE i.project_instance_id = labels.project_instance_id
3167                         AND i.caller_label = labels.caller_label
3168                         AND i.state = ?6
3169                   )
3170                   AND (
3171                       SELECT COUNT(*) FROM usage_instances AS i
3172                       WHERE i.project_instance_id = labels.project_instance_id
3173                         AND i.caller_label = labels.caller_label
3174                   ) <= ?7
3175                 LIMIT 1
3176             ) OR EXISTS(
3177                 SELECT 1 FROM usage_label_tombstones
3178                 WHERE expired_at_epoch < ?8 LIMIT 1
3179             ) OR EXISTS(
3180                 SELECT 1 FROM usage_instance_tombstones
3181                 WHERE retired_at_epoch < ?8 LIMIT 1
3182             )",
3183        params![
3184            INSTANCE_SEALED,
3185            INSTANCE_EXPIRED,
3186            instance_cutoff,
3187            daily_cutoff,
3188            label_cutoff,
3189            INSTANCE_ACTIVE,
3190            to_i64("prune_batch_rows", policy.prune_batch_rows)?,
3191            tombstone_cutoff,
3192        ],
3193        |row| row.get::<_, i64>(0),
3194    )?;
3195    Ok(pending != 0)
3196}
3197
3198#[allow(clippy::too_many_arguments)]
3199/// Refresh content-free lifecycle state from exact counters and bounded probes.
3200fn refresh_retention_state(
3201    connection: &Connection,
3202    policy: TelemetryRetentionPolicy,
3203    now: i64,
3204    pruned_raw: usize,
3205    pruned_instances: usize,
3206    evicted_tombstones: usize,
3207    writes_added: usize,
3208) -> DbResult<()> {
3209    let counters = connection.query_row(
3210        "SELECT raw_rows, raw_logical_bytes, instance_rows, label_rows, daily_rows,
3211                label_tombstone_rows, instance_tombstone_rows,
3212                pruned_raw_rows, pruned_instance_rows, evicted_tombstones,
3213                writes_since_checkpoint
3214         FROM usage_retention_state WHERE singleton = 1",
3215        [],
3216        |row| {
3217            Ok((
3218                row.get::<_, i64>(0)?,
3219                row.get::<_, i64>(1)?,
3220                row.get::<_, i64>(2)?,
3221                row.get::<_, i64>(3)?,
3222                row.get::<_, i64>(4)?,
3223                row.get::<_, i64>(5)?,
3224                row.get::<_, i64>(6)?,
3225                row.get::<_, i64>(7)?,
3226                row.get::<_, i64>(8)?,
3227                row.get::<_, i64>(9)?,
3228                row.get::<_, i64>(10)?,
3229            ))
3230        },
3231    )?;
3232    let raw_rows = count_usize("raw_rows", counters.0)?;
3233    let raw_bytes = count_usize("raw_logical_bytes", counters.1)?;
3234    let instance_rows = count_usize("instance_rows", counters.2)?;
3235    let label_rows = count_usize("label_rows", counters.3)?;
3236    let daily_rows = count_usize("daily_rows", counters.4)?;
3237    let label_tombstones = count_usize("label_tombstone_rows", counters.5)?;
3238    let instance_tombstones = count_usize("instance_tombstone_rows", counters.6)?;
3239    let cutoff = epoch_cutoff(now, policy.max_raw_age_seconds, "max_raw_age_seconds")?;
3240    let old_raw = connection.query_row(
3241        "SELECT EXISTS(
3242             SELECT 1 FROM usage_events
3243             WHERE created_at_epoch < ?1 LIMIT 1
3244         )",
3245        [cutoff],
3246        |row| row.get::<_, i64>(0),
3247    )?;
3248    let oldest = connection
3249        .query_row(
3250            "SELECT created_at_epoch FROM usage_events
3251         ORDER BY created_at_epoch, id LIMIT 1",
3252            [],
3253            |row| row.get::<_, i64>(0),
3254        )
3255        .optional()?;
3256    let pruned_raw_total = checked_count_add("pruned_raw_rows", counters.7, pruned_raw)?;
3257    let pruned_instance_total =
3258        checked_count_add("pruned_instance_rows", counters.8, pruned_instances)?;
3259    let evicted_total = checked_count_add("evicted_tombstones", counters.9, evicted_tombstones)?;
3260    let writes = checked_count_add("writes_since_checkpoint", counters.10, writes_added)?;
3261    let maintenance_pending = raw_rows > policy.max_raw_rows
3262        || raw_bytes > policy.max_raw_logical_bytes
3263        || old_raw > 0
3264        || instance_rows > policy.max_retained_instances
3265        || label_rows > policy.max_retained_labels
3266        || daily_rows > policy.max_daily_rows
3267        || label_tombstones > policy.max_label_tombstones
3268        || instance_tombstones > policy.max_instance_tombstones
3269        || count_usize("writes_since_checkpoint", writes)? >= policy.checkpoint_write_interval
3270        || aged_maintenance_pending(connection, policy, now)?;
3271    connection.execute(
3272        "UPDATE usage_retention_state
3273         SET policy_version = ?1,
3274             logical_byte_version = ?2,
3275             pruned_raw_rows = ?3,
3276             pruned_instance_rows = ?4,
3277             evicted_tombstones = ?5,
3278             writes_since_checkpoint = ?6,
3279             last_maintenance_epoch = ?7,
3280             oldest_retained_epoch = ?8,
3281             maintenance_pending = ?9
3282         WHERE singleton = 1",
3283        params![
3284            i64::from(POLICY_VERSION),
3285            i64::from(LOGICAL_BYTE_VERSION),
3286            pruned_raw_total,
3287            pruned_instance_total,
3288            evicted_total,
3289            writes,
3290            now,
3291            oldest,
3292            i64::from(maintenance_pending),
3293        ],
3294    )?;
3295    Ok(())
3296}
3297
3298/// Converge all bounded retention categories during a supported upgrade.
3299fn converge_retention(
3300    connection: &Connection,
3301    project: ProjectInstanceId,
3302    policy: TelemetryRetentionPolicy,
3303    now: i64,
3304) -> DbResult<()> {
3305    let mut pruned_raw = 0usize;
3306    loop {
3307        let (rows, bytes, old_rows) = raw_pressure(connection, policy, now)?;
3308        if rows <= policy.max_raw_rows && bytes <= policy.max_raw_logical_bytes && old_rows == 0 {
3309            break;
3310        }
3311        let deleted = prune_raw_once(connection, policy, now)?;
3312        if deleted == 0 {
3313            break;
3314        }
3315        pruned_raw = pruned_raw
3316            .checked_add(deleted)
3317            .ok_or(DbError::TelemetryIntegerOverflow {
3318                field: "pruned_raw_rows",
3319            })?;
3320    }
3321    let mut pruned_instances = 0usize;
3322    loop {
3323        if retention_counter(connection, RetentionCounter::InstanceRows)?
3324            <= policy.max_retained_instances
3325        {
3326            break;
3327        }
3328        let (deleted, raw) = prune_instances_once(connection, policy, now, 0)?;
3329        if deleted == 0 {
3330            break;
3331        }
3332        pruned_instances =
3333            pruned_instances
3334                .checked_add(deleted)
3335                .ok_or(DbError::TelemetryIntegerOverflow {
3336                    field: "pruned_instance_rows",
3337                })?;
3338        pruned_raw = pruned_raw
3339            .checked_add(raw)
3340            .ok_or(DbError::TelemetryIntegerOverflow {
3341                field: "pruned_raw_rows",
3342            })?;
3343    }
3344    while prune_daily_once(connection, policy, now)? == policy.prune_batch_rows {}
3345    let mut evicted = 0usize;
3346    loop {
3347        let deleted = prune_tombstones_once(connection, policy, now)?;
3348        evicted = evicted
3349            .checked_add(deleted)
3350            .ok_or(DbError::TelemetryIntegerOverflow {
3351                field: "evicted_tombstones",
3352            })?;
3353        if deleted < policy.prune_batch_rows {
3354            break;
3355        }
3356    }
3357    refresh_retention_state(
3358        connection,
3359        policy,
3360        now,
3361        pruned_raw,
3362        pruned_instances,
3363        evicted,
3364        0,
3365    )?;
3366    retention_state_for_project(connection, project).map(|_| ())
3367}
3368
3369/// Build the shared raw-event projection with a fixed internal predicate suffix.
3370fn raw_event_select(extra_predicate: &str) -> String {
3371    format!(
3372        "SELECT COALESCE(i.caller_label, ''), e.command, e.path, e.query,
3373                e.estimated_tokens_without_projectatlas,
3374                e.estimated_tokens_with_projectatlas, e.estimated_tokens_saved,
3375                d.token_savings_bucket, d.provider, d.model, d.tokenizer_backend,
3376                d.accuracy, d.baseline_kind, d.confidence, e.calculation_trace,
3377                d.accounting_layer, d.estimate_method, d.denominator_kind,
3378                e.baseline_identity, e.baseline_fingerprint, d.dedupe_scope
3379         FROM usage_events AS e
3380         JOIN usage_instances AS i USING(instance_row_id)
3381         JOIN usage_bucket_dimensions AS d USING(dimension_id)
3382         WHERE i.project_instance_id = ?1 {extra_predicate}
3383         ORDER BY e.id"
3384    )
3385}
3386
3387/// Decode one retained raw event and its normalized dimensions.
3388fn map_usage_event(row: &rusqlite::Row<'_>) -> DbResult<UsageEvent> {
3389    Ok(UsageEvent {
3390        session_id: row.get(0)?,
3391        command: row.get(1)?,
3392        path: row.get(2)?,
3393        query: row.get(3)?,
3394        estimated_tokens_without_projectatlas: row.get(4)?,
3395        estimated_tokens_with_projectatlas: row.get(5)?,
3396        estimated_tokens_saved: row.get(6)?,
3397        token_savings_bucket: row.get(7)?,
3398        provider: row.get(8)?,
3399        model: row.get(9)?,
3400        tokenizer_backend: row.get(10)?,
3401        accuracy: row.get(11)?,
3402        baseline_kind: row.get(12)?,
3403        confidence: row.get(13)?,
3404        calculation_trace: row.get(14)?,
3405        accounting_layer: row.get(15)?,
3406        estimate_method: row.get(16)?,
3407        denominator_kind: row.get(17)?,
3408        baseline_identity: row.get(18)?,
3409        baseline_fingerprint: row.get(19)?,
3410        dedupe_scope: row.get(20)?,
3411    })
3412}
3413
3414/// Load exact all-time aggregate rows for a project or caller label.
3415fn load_overview_aggregates(
3416    connection: &Connection,
3417    project: ProjectInstanceId,
3418    caller_label: Option<&str>,
3419) -> DbResult<Vec<(DimensionValues, AggregateCounters)>> {
3420    let (sql, label) = if let Some(label) = caller_label {
3421        (
3422            "SELECT d.token_savings_bucket, d.provider, d.model, d.tokenizer_backend,
3423                    d.accuracy, d.baseline_kind, d.confidence, d.accounting_layer,
3424                    d.estimate_method, d.denominator_kind, d.dedupe_scope, d.overflow,
3425                    a.calls, a.estimated_without, a.estimated_with,
3426                    a.observed_without, a.observed_with, a.modeled_without, a.modeled_with,
3427                    a.deduped_modeled_without, a.deduped_modeled_with,
3428                    a.repeated_baselines, a.observed_file_read_replacements,
3429                    a.modeled_file_reads_avoided
3430             FROM (
3431                 SELECT aggregate.dimension_id,
3432                        SUM(aggregate.calls) AS calls,
3433                        SUM(aggregate.estimated_without) AS estimated_without,
3434                        SUM(aggregate.estimated_with) AS estimated_with,
3435                        SUM(aggregate.observed_without) AS observed_without,
3436                        SUM(aggregate.observed_with) AS observed_with,
3437                        SUM(aggregate.modeled_without) AS modeled_without,
3438                        SUM(aggregate.modeled_with) AS modeled_with,
3439                        SUM(aggregate.deduped_modeled_without) AS deduped_modeled_without,
3440                        SUM(aggregate.deduped_modeled_with) AS deduped_modeled_with,
3441                        SUM(aggregate.repeated_baselines) AS repeated_baselines,
3442                        SUM(aggregate.observed_file_read_replacements)
3443                            AS observed_file_read_replacements,
3444                        SUM(aggregate.modeled_file_reads_avoided)
3445                            AS modeled_file_reads_avoided
3446                 FROM usage_instance_aggregates AS aggregate
3447                 JOIN usage_instances AS instance USING(instance_row_id)
3448                 WHERE instance.project_instance_id = ?1 AND instance.caller_label = ?2
3449                 GROUP BY aggregate.dimension_id
3450             ) AS a
3451             JOIN usage_bucket_dimensions AS d USING(dimension_id)
3452             ORDER BY d.dimension_id",
3453            Some(label),
3454        )
3455    } else {
3456        (
3457            "SELECT d.token_savings_bucket, d.provider, d.model, d.tokenizer_backend,
3458                    d.accuracy, d.baseline_kind, d.confidence, d.accounting_layer,
3459                    d.estimate_method, d.denominator_kind, d.dedupe_scope, d.overflow,
3460                    a.calls, a.estimated_without, a.estimated_with,
3461                    a.observed_without, a.observed_with, a.modeled_without, a.modeled_with,
3462                    a.deduped_modeled_without, a.deduped_modeled_with,
3463                    a.repeated_baselines, a.observed_file_read_replacements,
3464                    a.modeled_file_reads_avoided
3465             FROM usage_global_aggregates AS a
3466             JOIN usage_bucket_dimensions AS d USING(dimension_id)
3467             WHERE a.project_instance_id = ?1
3468             ORDER BY d.dimension_id",
3469            None,
3470        )
3471    };
3472    let mut statement = connection.prepare_cached(sql)?;
3473    let mut rows = if let Some(label) = label {
3474        statement.query(params![project.as_bytes().as_slice(), label])?
3475    } else {
3476        statement.query([project.as_bytes().as_slice()])?
3477    };
3478    let mut result = Vec::new();
3479    while let Some(row) = rows.next()? {
3480        result.push((read_dimension(row, 0)?, read_counters_offset(row, 12)?));
3481    }
3482    Ok(result)
3483}
3484
3485/// Combine normalized aggregate rows into one overview and bounded buckets.
3486fn aggregate_report_rows(
3487    rows: Vec<(DimensionValues, AggregateCounters)>,
3488) -> DbResult<(Vec<TokenBucketOverview>, TokenAccountingTotals)> {
3489    let mut by_dimension = BTreeMap::<DimensionValues, AggregateCounters>::new();
3490    for (dimension, counters) in rows {
3491        let entry = by_dimension.entry(dimension).or_default();
3492        *entry = entry.checked_add(counters)?;
3493    }
3494    let mut totals = TokenAccountingTotals::default();
3495    let mut buckets = Vec::with_capacity(by_dimension.len());
3496    for (dimension, counters) in by_dimension {
3497        totals.measured_tokens_saved = totals
3498            .measured_tokens_saved
3499            .checked_add(component_difference(
3500                counters.observed_without,
3501                counters.observed_with,
3502            ))
3503            .ok_or(DbError::TelemetryIntegerOverflow {
3504                field: "measured_tokens_saved",
3505            })?;
3506        totals.gross_modeled_tokens_avoided = totals
3507            .gross_modeled_tokens_avoided
3508            .checked_add(component_difference(
3509                counters.modeled_without,
3510                counters.modeled_with,
3511            ))
3512            .ok_or(DbError::TelemetryIntegerOverflow {
3513                field: "gross_modeled_tokens_avoided",
3514            })?;
3515        totals.deduped_modeled_tokens_avoided = totals
3516            .deduped_modeled_tokens_avoided
3517            .checked_add(component_difference(
3518                counters.deduped_modeled_without,
3519                counters.deduped_modeled_with,
3520            ))
3521            .ok_or(DbError::TelemetryIntegerOverflow {
3522                field: "deduped_modeled_tokens_avoided",
3523            })?;
3524        totals.repeated_baselines_deduped = totals
3525            .repeated_baselines_deduped
3526            .checked_add(count_u128(
3527                "repeated_baselines",
3528                counters.repeated_baselines,
3529            )?)
3530            .ok_or(DbError::TelemetryIntegerOverflow {
3531                field: "repeated_baselines",
3532            })?;
3533        totals.observed_file_read_replacements = totals
3534            .observed_file_read_replacements
3535            .checked_add(count_u128(
3536                "observed_file_read_replacements",
3537                counters.observed_file_read_replacements,
3538            )?)
3539            .ok_or(DbError::TelemetryIntegerOverflow {
3540                field: "observed_file_read_replacements",
3541            })?;
3542        totals.modeled_file_reads_avoided = totals
3543            .modeled_file_reads_avoided
3544            .checked_add(count_u128(
3545                "modeled_file_reads_avoided",
3546                counters.modeled_file_reads_avoided,
3547            )?)
3548            .ok_or(DbError::TelemetryIntegerOverflow {
3549                field: "modeled_file_reads_avoided",
3550            })?;
3551        buckets.push(bucket_from_counters(dimension, counters)?);
3552    }
3553    Ok((buckets, totals))
3554}
3555
3556/// Convert one normalized counter row into the public bucket contract.
3557fn bucket_from_counters(
3558    dimension: DimensionValues,
3559    counters: AggregateCounters,
3560) -> DbResult<TokenBucketOverview> {
3561    Ok(TokenBucketOverview::from_totals(
3562        dimension.token_savings_bucket,
3563        dimension.provider,
3564        dimension.model,
3565        dimension.tokenizer_backend,
3566        dimension.accuracy,
3567        dimension.baseline_kind,
3568        dimension.confidence,
3569        dimension.accounting_layer,
3570        dimension.estimate_method,
3571        dimension.denominator_kind,
3572        dimension.dedupe_scope,
3573        count_u128("calls", counters.calls)?,
3574        count_u128("estimated_without", counters.estimated_without)?,
3575        count_u128("estimated_with", counters.estimated_with)?,
3576    ))
3577}
3578
3579/// Load retained daily aggregates for one project or caller label.
3580fn load_daily_aggregates(
3581    connection: &Connection,
3582    project: ProjectInstanceId,
3583    caller_label: Option<&str>,
3584    window: TokenTrendWindow,
3585) -> DbResult<Vec<(String, DimensionValues, AggregateCounters)>> {
3586    let period_expression = match window {
3587        TokenTrendWindow::Day => "strftime('%Y-%m-%d', a.day_epoch, 'unixepoch')",
3588        TokenTrendWindow::Week => "strftime('%Y-W%W', a.day_epoch, 'unixepoch')",
3589        TokenTrendWindow::Month => "strftime('%Y-%m', a.day_epoch, 'unixepoch')",
3590        TokenTrendWindow::Year => "strftime('%Y', a.day_epoch, 'unixepoch')",
3591    };
3592    let (table, join, predicate) = if caller_label.is_some() {
3593        (
3594            "usage_instance_daily_aggregates",
3595            "JOIN usage_instances AS i USING(instance_row_id)",
3596            "i.project_instance_id = ?1 AND i.caller_label = ?2",
3597        )
3598    } else {
3599        ("usage_daily_aggregates", "", "a.project_instance_id = ?1")
3600    };
3601    let sql = format!(
3602        "SELECT grouped.period,
3603                d.token_savings_bucket, d.provider, d.model, d.tokenizer_backend,
3604                d.accuracy, d.baseline_kind, d.confidence, d.accounting_layer,
3605                d.estimate_method, d.denominator_kind, d.dedupe_scope, d.overflow,
3606                grouped.calls, grouped.estimated_without, grouped.estimated_with,
3607                grouped.observed_without, grouped.observed_with,
3608                grouped.modeled_without, grouped.modeled_with,
3609                grouped.deduped_modeled_without, grouped.deduped_modeled_with,
3610                grouped.repeated_baselines, grouped.observed_file_read_replacements,
3611                grouped.modeled_file_reads_avoided
3612         FROM (
3613             SELECT {period_expression} AS period, a.dimension_id,
3614                    SUM(a.calls) AS calls,
3615                    SUM(a.estimated_without) AS estimated_without,
3616                    SUM(a.estimated_with) AS estimated_with,
3617                    SUM(a.observed_without) AS observed_without,
3618                    SUM(a.observed_with) AS observed_with,
3619                    SUM(a.modeled_without) AS modeled_without,
3620                    SUM(a.modeled_with) AS modeled_with,
3621                    SUM(a.deduped_modeled_without) AS deduped_modeled_without,
3622                    SUM(a.deduped_modeled_with) AS deduped_modeled_with,
3623                    SUM(a.repeated_baselines) AS repeated_baselines,
3624                    SUM(a.observed_file_read_replacements)
3625                        AS observed_file_read_replacements,
3626                    SUM(a.modeled_file_reads_avoided) AS modeled_file_reads_avoided
3627             FROM {table} AS a {join}
3628             WHERE {predicate}
3629             GROUP BY period, a.dimension_id
3630         ) AS grouped
3631         JOIN usage_bucket_dimensions AS d USING(dimension_id)
3632         ORDER BY grouped.period, grouped.dimension_id"
3633    );
3634    let mut statement = connection.prepare(&sql)?;
3635    let mut rows = if let Some(label) = caller_label {
3636        statement.query(params![project.as_bytes().as_slice(), label])?
3637    } else {
3638        statement.query([project.as_bytes().as_slice()])?
3639    };
3640    let mut result = Vec::new();
3641    while let Some(row) = rows.next()? {
3642        result.push((
3643            row.get::<_, String>(0)?,
3644            read_dimension(row, 1)?,
3645            read_counters_offset(row, 13)?,
3646        ));
3647    }
3648    Ok(result)
3649}
3650
3651/// Classify retained, partial, expired, or unavailable caller detail.
3652fn detail_availability(
3653    connection: &Connection,
3654    project: ProjectInstanceId,
3655    caller_label: Option<&str>,
3656) -> DbResult<UsageDetailAvailability> {
3657    let state = connection.query_row(
3658        "SELECT raw_detail_complete, dimension_detail_complete, label_history_complete
3659         FROM usage_retention_state WHERE singleton = 1",
3660        [],
3661        |row| {
3662            Ok((
3663                row.get::<_, i64>(0)?,
3664                row.get::<_, i64>(1)?,
3665                row.get::<_, i64>(2)?,
3666            ))
3667        },
3668    )?;
3669    let global_complete = bool_from_sql("raw_detail_complete", state.0)?
3670        && bool_from_sql("dimension_detail_complete", state.1)?
3671        && bool_from_sql("label_history_complete", state.2)?;
3672    let Some(label) = caller_label else {
3673        return Ok(if global_complete {
3674            UsageDetailAvailability::Retained
3675        } else {
3676            UsageDetailAvailability::Partial
3677        });
3678    };
3679    let instances = connection.query_row(
3680        "SELECT COUNT(*), MIN(raw_detail_complete)
3681         FROM usage_instances
3682         WHERE project_instance_id = ?1 AND caller_label = ?2",
3683        params![project.as_bytes().as_slice(), label],
3684        |row| Ok((row.get::<_, i64>(0)?, row.get::<_, Option<i64>>(1)?)),
3685    )?;
3686    let label_detail = connection
3687        .query_row(
3688            "SELECT detail_complete FROM usage_labels
3689             WHERE project_instance_id = ?1 AND caller_label = ?2",
3690            params![project.as_bytes().as_slice(), label],
3691            |row| row.get::<_, i64>(0),
3692        )
3693        .optional()?;
3694    if instances.0 > 0 {
3695        return Ok(
3696            if global_complete && instances.1 == Some(1) && label_detail == Some(1) {
3697                UsageDetailAvailability::Retained
3698            } else {
3699                UsageDetailAvailability::Partial
3700            },
3701        );
3702    }
3703    let tombstone = connection.query_row(
3704        "SELECT EXISTS(
3705             SELECT 1 FROM usage_label_tombstones
3706             WHERE project_instance_id = ?1 AND caller_label = ?2
3707         )",
3708        params![project.as_bytes().as_slice(), label],
3709        |row| row.get::<_, i64>(0),
3710    )?;
3711    Ok(if tombstone != 0 || label_detail == Some(0) {
3712        UsageDetailAvailability::Expired
3713    } else {
3714        UsageDetailAvailability::Unavailable
3715    })
3716}
3717
3718/// Decode one normalized dimension beginning at the selected column offset.
3719fn read_dimension(row: &rusqlite::Row<'_>, offset: usize) -> DbResult<DimensionValues> {
3720    Ok(DimensionValues {
3721        token_savings_bucket: row.get(offset)?,
3722        provider: row.get(offset + 1)?,
3723        model: row.get(offset + 2)?,
3724        tokenizer_backend: row.get(offset + 3)?,
3725        accuracy: row.get(offset + 4)?,
3726        baseline_kind: row.get(offset + 5)?,
3727        confidence: row.get(offset + 6)?,
3728        accounting_layer: row.get(offset + 7)?,
3729        estimate_method: row.get(offset + 8)?,
3730        denominator_kind: row.get(offset + 9)?,
3731        dedupe_scope: row.get(offset + 10)?,
3732        overflow: bool_from_sql("usage_bucket_dimensions.overflow", row.get(offset + 11)?)?,
3733    })
3734}
3735
3736/// Decode aggregate counters beginning at the selected column offset.
3737fn read_counters_offset(row: &rusqlite::Row<'_>, offset: usize) -> DbResult<AggregateCounters> {
3738    Ok(AggregateCounters {
3739        calls: row.get(offset)?,
3740        estimated_without: row.get(offset + 1)?,
3741        estimated_with: row.get(offset + 2)?,
3742        observed_without: row.get(offset + 3)?,
3743        observed_with: row.get(offset + 4)?,
3744        modeled_without: row.get(offset + 5)?,
3745        modeled_with: row.get(offset + 6)?,
3746        deduped_modeled_without: row.get(offset + 7)?,
3747        deduped_modeled_with: row.get(offset + 8)?,
3748        repeated_baselines: row.get(offset + 9)?,
3749        observed_file_read_replacements: row.get(offset + 10)?,
3750        modeled_file_reads_avoided: row.get(offset + 11)?,
3751    })
3752}
3753
3754/// Calculate deterministic retained logical bytes for one raw event.
3755fn logical_event_bytes(event: &UsageEvent, label: Option<&str>) -> DbResult<usize> {
3756    let identity = event.effective_baseline_identity();
3757    let fingerprint = event.effective_baseline_fingerprint();
3758    [
3759        label.unwrap_or_default(),
3760        event.command.as_str(),
3761        event.path.as_deref().unwrap_or_default(),
3762        event.query.as_deref().unwrap_or_default(),
3763        event.token_savings_bucket.as_str(),
3764        event.provider.as_str(),
3765        event.model.as_str(),
3766        event.tokenizer_backend.as_str(),
3767        event.accuracy.as_str(),
3768        event.baseline_kind.as_str(),
3769        event.confidence.as_str(),
3770        event.calculation_trace.as_str(),
3771        event.report_accounting_layer(),
3772        event.estimate_method.as_str(),
3773        event.report_denominator_kind(),
3774        identity.as_ref(),
3775        fingerprint.as_ref(),
3776        event.report_dedupe_scope(),
3777    ]
3778    .into_iter()
3779    .try_fold(96_usize, |total, value| {
3780        total
3781            .checked_add(value.len())
3782            .ok_or(DbError::TelemetryIntegerOverflow {
3783                field: "raw_logical_bytes",
3784            })
3785    })
3786}
3787
3788/// Encode one signed value as separate nonnegative positive and negative components.
3789fn signed_components(value: i64) -> DbResult<(i64, i64)> {
3790    if value >= 0 {
3791        Ok((value, 0))
3792    } else {
3793        Ok((
3794            0,
3795            value
3796                .checked_neg()
3797                .ok_or(DbError::TelemetryIntegerOverflow {
3798                    field: "signed_component",
3799                })?,
3800        ))
3801    }
3802}
3803
3804/// Reconstruct a signed aggregate from its nonnegative components.
3805fn component_difference(without: i64, with: i64) -> i128 {
3806    i128::from(without) - i128::from(with)
3807}
3808
3809/// Add an unsigned count to a persisted `SQLite` integer exactly.
3810fn checked_count_add(field: &'static str, previous: i64, value: usize) -> DbResult<i64> {
3811    previous
3812        .checked_add(to_i64(field, value)?)
3813        .ok_or(DbError::TelemetryIntegerOverflow { field })
3814}
3815
3816/// Convert one bounded duration to milliseconds without truncation.
3817fn duration_millis(field: &'static str, duration: Duration) -> DbResult<u64> {
3818    duration
3819        .as_secs()
3820        .checked_mul(1_000)
3821        .and_then(|milliseconds| milliseconds.checked_add(u64::from(duration.subsec_millis())))
3822        .ok_or(DbError::TelemetryIntegerOverflow { field })
3823}
3824
3825/// Read one persisted retention counter as an unsigned Rust count.
3826fn retention_counter(connection: &Connection, counter: RetentionCounter) -> DbResult<usize> {
3827    let value = connection.query_row(counter.select_sql(), [], |row| row.get::<_, i64>(0))?;
3828    count_usize(counter.field(), value)
3829}
3830
3831/// Increase one persisted retention counter with checked arithmetic.
3832fn increment_retention_counter(
3833    connection: &Connection,
3834    counter: RetentionCounter,
3835    amount: usize,
3836) -> DbResult<()> {
3837    let value = retention_counter(connection, counter)?
3838        .checked_add(amount)
3839        .ok_or(DbError::TelemetryIntegerOverflow {
3840            field: counter.field(),
3841        })?;
3842    connection.execute(counter.update_sql(), [to_i64(counter.field(), value)?])?;
3843    Ok(())
3844}
3845
3846/// Decrease one persisted retention counter with checked arithmetic.
3847fn decrement_retention_counter(
3848    connection: &Connection,
3849    counter: RetentionCounter,
3850    amount: usize,
3851) -> DbResult<()> {
3852    let current = retention_counter(connection, counter)?;
3853    let value = current
3854        .checked_sub(amount)
3855        .ok_or(DbError::TelemetryIntegerOverflow {
3856            field: counter.field(),
3857        })?;
3858    connection.execute(counter.update_sql(), [to_i64(counter.field(), value)?])?;
3859    Ok(())
3860}
3861
3862/// Decode a strict zero-or-one `SQLite` boolean.
3863fn bool_from_sql(field: &'static str, value: i64) -> DbResult<bool> {
3864    match value {
3865        0 => Ok(false),
3866        1 => Ok(true),
3867        _ => Err(DbError::InvalidEnum {
3868            field,
3869            value: value.to_string(),
3870        }),
3871    }
3872}
3873
3874/// Convert an optional unsigned count to a `SQLite` integer.
3875fn option_usize_to_i64(field: &'static str, value: Option<usize>) -> DbResult<Option<i64>> {
3876    value.map(|value| to_i64(field, value)).transpose()
3877}
3878
3879/// Convert an optional signed count to a `SQLite` integer.
3880fn option_isize_to_i64(field: &'static str, value: Option<isize>) -> DbResult<Option<i64>> {
3881    value
3882        .map(|value| {
3883            i64::try_from(value).map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3884        })
3885        .transpose()
3886}
3887
3888/// Convert one integer-like value into the exact `SQLite` integer range.
3889fn to_i64(field: &'static str, value: impl TryInto<i64>) -> DbResult<i64> {
3890    value
3891        .try_into()
3892        .map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3893}
3894
3895/// Decode one nonnegative persisted count as `usize`.
3896fn count_usize(field: &'static str, value: i64) -> DbResult<usize> {
3897    usize::try_from(value).map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3898}
3899
3900/// Decode one nonnegative persisted count as `u32`.
3901fn count_u32(field: &'static str, value: i64) -> DbResult<u32> {
3902    u32::try_from(value).map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3903}
3904
3905/// Decode one nonnegative persisted count as `u64`.
3906fn count_u64(field: &'static str, value: i64) -> DbResult<u64> {
3907    u64::try_from(value).map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3908}
3909
3910/// Decode one nonnegative persisted count as `u128`.
3911fn count_u128(field: &'static str, value: i64) -> DbResult<u128> {
3912    u128::try_from(value).map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3913}
3914
3915/// Calculate a nonnegative retention cutoff from an epoch and duration.
3916fn epoch_cutoff(now: i64, seconds: u64, field: &'static str) -> DbResult<i64> {
3917    let seconds = to_i64(field, seconds)?;
3918    Ok(now.checked_sub(seconds).unwrap_or(0))
3919}
3920
3921/// Read the current Unix epoch in the persisted integer range.
3922fn now_epoch_seconds() -> DbResult<i64> {
3923    let duration = SystemTime::now()
3924        .duration_since(UNIX_EPOCH)
3925        .map_err(|_source| DbError::TelemetryIntegerOverflow {
3926            field: "system_time",
3927        })?;
3928    to_i64("system_time", duration.as_secs())
3929}
3930
3931/// Read one allowlisted numeric `SQLite` pragma.
3932fn pragma_count(connection: &Connection, pragma: &str) -> DbResult<i64> {
3933    let sql = match pragma {
3934        "freelist_count" => "PRAGMA freelist_count",
3935        "page_count" => "PRAGMA page_count",
3936        "page_size" => "PRAGMA page_size",
3937        _ => {
3938            return Err(DbError::InvalidEnum {
3939                field: "telemetry_pragma",
3940                value: pragma.to_string(),
3941            });
3942        }
3943    };
3944    connection
3945        .query_row(sql, [], |row| row.get(0))
3946        .map_err(Into::into)
3947}
3948
3949/// Decode the numeric `SQLite` synchronous pragma into its stable name.
3950fn synchronous_mode(value: i64) -> DbResult<&'static str> {
3951    match value {
3952        0 => Ok("off"),
3953        1 => Ok("normal"),
3954        2 => Ok("full"),
3955        3 => Ok("extra"),
3956        _ => Err(DbError::InvalidEnum {
3957            field: "pragma.synchronous",
3958            value: value.to_string(),
3959        }),
3960    }
3961}
3962
3963#[cfg(test)]
3964mod tests {
3965    use super::*;
3966    use crate::AtlasStore;
3967    use projectatlas_core::telemetry::usage_from_estimates;
3968    use projectatlas_core::{Node, NodeKind, normalized_parent};
3969    use rusqlite::{Transaction, TransactionBehavior};
3970    use std::error::Error;
3971    use std::fs;
3972
3973    struct TestDatabase {
3974        temp: tempfile::TempDir,
3975        connection: Connection,
3976        project: ProjectInstanceId,
3977    }
3978
3979    struct ProductionTestDatabase {
3980        _temp: tempfile::TempDir,
3981        root: std::path::PathBuf,
3982        database_path: std::path::PathBuf,
3983        store: AtlasStore,
3984        project: ProjectInstanceId,
3985    }
3986
3987    fn test_database() -> Result<TestDatabase, Box<dyn Error>> {
3988        let temp = tempfile::tempdir()?;
3989        let database_path = temp.path().join("projectatlas.db");
3990        let connection = Connection::open(&database_path)?;
3991        crate::schema::initialize(&connection, None)?;
3992        let project = crate::project_identity::ensure_project_identity(&connection)?.0;
3993        Ok(TestDatabase {
3994            temp,
3995            connection,
3996            project,
3997        })
3998    }
3999
4000    fn production_database() -> Result<ProductionTestDatabase, Box<dyn Error>> {
4001        let temp = tempfile::tempdir()?;
4002        let root = temp.path().join("repository");
4003        let atlas = root.join(".projectatlas");
4004        fs::create_dir_all(&atlas)?;
4005        let database_path = atlas.join("projectatlas.db");
4006        let store = AtlasStore::open_for_project(&database_path, &root)?;
4007        let project = store
4008            .validated_project_instance_id
4009            .ok_or(DbError::ProjectInstanceIdentityMissing)?;
4010        Ok(ProductionTestDatabase {
4011            _temp: temp,
4012            root,
4013            database_path,
4014            store,
4015            project,
4016        })
4017    }
4018
4019    fn instance(byte: u8) -> Result<UsageInstanceId, Box<dyn Error>> {
4020        Ok(UsageInstanceId::from_bytes([byte; 16])?)
4021    }
4022
4023    fn event(label: &str, without: usize, with: usize) -> UsageEvent {
4024        let mut event = usage_from_estimates(
4025            label,
4026            "summary",
4027            Some("src/lib.rs".to_string()),
4028            None,
4029            without,
4030            with,
4031        );
4032        event.baseline_identity = "source:src/lib.rs".to_string();
4033        event.baseline_fingerprint = "source:src/lib.rs:v1".to_string();
4034        event
4035    }
4036
4037    fn record_transaction(
4038        connection: &Connection,
4039        project: ProjectInstanceId,
4040        instance: UsageInstanceId,
4041        owner: UsageInstanceOwner,
4042        event: &UsageEvent,
4043        policy: TelemetryRetentionPolicy,
4044        seal_after_record: bool,
4045    ) -> DbResult<()> {
4046        record_transaction_at(
4047            connection,
4048            project,
4049            instance,
4050            owner,
4051            event,
4052            policy,
4053            seal_after_record,
4054            now_epoch_seconds()?,
4055        )
4056    }
4057
4058    #[allow(clippy::too_many_arguments)]
4059    fn record_transaction_at(
4060        connection: &Connection,
4061        project: ProjectInstanceId,
4062        instance: UsageInstanceId,
4063        owner: UsageInstanceOwner,
4064        event: &UsageEvent,
4065        policy: TelemetryRetentionPolicy,
4066        seal_after_record: bool,
4067        now: i64,
4068    ) -> DbResult<()> {
4069        let transaction = Transaction::new_unchecked(connection, TransactionBehavior::Immediate)?;
4070        let result = (|| {
4071            crate::project_identity::require_bound_project_identity(&transaction, project)?;
4072            let policy = policy.validate()?;
4073            validate_event(event, policy)?;
4074            record_usage_at(
4075                &transaction,
4076                project,
4077                instance,
4078                owner,
4079                event,
4080                policy,
4081                now,
4082                seal_after_record,
4083                BaselineAdmission::BoundedRuntime,
4084                DimensionAdmission::Event,
4085            )
4086        })();
4087        match result {
4088            Ok(()) => transaction.commit().map_err(Into::into),
4089            Err(error) => {
4090                transaction.rollback()?;
4091                Err(error)
4092            }
4093        }
4094    }
4095
4096    fn scalar_count(connection: &Connection, sql: &str) -> Result<usize, Box<dyn Error>> {
4097        let value = connection.query_row(sql, [], |row| row.get::<_, i64>(0))?;
4098        Ok(usize::try_from(value)?)
4099    }
4100
4101    fn query_plan(connection: &Connection, sql: &str) -> Result<Vec<String>, Box<dyn Error>> {
4102        let mut statement = connection.prepare(&format!("EXPLAIN QUERY PLAN {sql}"))?;
4103        let rows = statement.query_map([], |row| row.get::<_, String>(3))?;
4104        Ok(rows.collect::<Result<Vec<_>, _>>()?)
4105    }
4106
4107    fn assert_plan_uses(plan: &[String], owner: &str) {
4108        assert!(
4109            plan.iter().any(|detail| detail.contains(owner)),
4110            "expected query plan to use {owner}; plan was {plan:?}"
4111        );
4112    }
4113
4114    #[test]
4115    fn retention_policy_rejects_zero_and_inverted_instance_limits() {
4116        let policy = TelemetryRetentionPolicy {
4117            prune_batch_rows: 0,
4118            ..TelemetryRetentionPolicy::default()
4119        };
4120        assert!(policy.validate().is_err());
4121
4122        let default_policy = TelemetryRetentionPolicy::default();
4123        let policy = TelemetryRetentionPolicy {
4124            max_retained_instances: default_policy.max_active_instances - 1,
4125            ..default_policy
4126        };
4127        assert!(policy.validate().is_err());
4128    }
4129
4130    #[test]
4131    fn oversized_event_fields_are_rejected_at_the_typed_storage_boundary() {
4132        type Mutator = fn(&mut UsageEvent, String);
4133
4134        let policy = TelemetryRetentionPolicy::default();
4135        let dimension_limit = policy.max_dimension_bytes;
4136        let short_witness_limit = 256.min(policy.max_baseline_witness_bytes);
4137        let cases: [(&str, usize, Mutator); 18] = [
4138            ("session_id", policy.max_label_bytes, |event, value| {
4139                event.session_id = value;
4140            }),
4141            ("command", policy.max_command_bytes, |event, value| {
4142                event.command = value;
4143            }),
4144            ("path", policy.max_path_bytes, |event, value| {
4145                event.path = Some(value);
4146            }),
4147            ("query", policy.max_query_bytes, |event, value| {
4148                event.query = Some(value);
4149            }),
4150            ("token_savings_bucket", dimension_limit, |event, value| {
4151                event.token_savings_bucket = value;
4152            }),
4153            ("provider", dimension_limit, |event, value| {
4154                event.provider = value;
4155            }),
4156            ("model", dimension_limit, |event, value| {
4157                event.model = value;
4158            }),
4159            ("tokenizer_backend", dimension_limit, |event, value| {
4160                event.tokenizer_backend = value;
4161            }),
4162            ("accuracy", dimension_limit, |event, value| {
4163                event.accuracy = value;
4164            }),
4165            ("baseline_kind", dimension_limit, |event, value| {
4166                event.baseline_kind = value;
4167            }),
4168            ("confidence", dimension_limit, |event, value| {
4169                event.confidence = value;
4170            }),
4171            ("accounting_layer", dimension_limit, |event, value| {
4172                event.accounting_layer = value;
4173            }),
4174            ("estimate_method", dimension_limit, |event, value| {
4175                event.estimate_method = value;
4176            }),
4177            ("denominator_kind", dimension_limit, |event, value| {
4178                event.denominator_kind = value;
4179            }),
4180            ("dedupe_scope", dimension_limit, |event, value| {
4181                event.dedupe_scope = value;
4182            }),
4183            ("calculation_trace", short_witness_limit, |event, value| {
4184                event.calculation_trace = value;
4185            }),
4186            (
4187                "baseline_identity",
4188                policy.max_baseline_witness_bytes,
4189                |event, value| {
4190                    event.baseline_identity = value;
4191                },
4192            ),
4193            (
4194                "baseline_fingerprint",
4195                short_witness_limit,
4196                |event, value| {
4197                    event.baseline_fingerprint = value;
4198                },
4199            ),
4200        ];
4201        for (field, limit, mutate) in cases {
4202            let mut candidate = event("bounded", 10, 1);
4203            mutate(&mut candidate, "x".repeat(limit + 1));
4204            assert!(matches!(
4205                validate_event(&candidate, policy),
4206                Err(DbError::TelemetryFieldTooLarge {
4207                    field: found,
4208                    bytes,
4209                    limit: found_limit,
4210                }) if found == field && bytes == limit + 1 && found_limit == limit
4211            ));
4212        }
4213    }
4214
4215    #[test]
4216    fn spill_cleanup_is_not_applicable() {
4217        assert_eq!(
4218            SpillCleanupState::NotApplicable,
4219            SpillCleanupState::NotApplicable
4220        );
4221        assert_eq!(POLICY_VERSION, 1);
4222        assert_eq!(LOGICAL_BYTE_VERSION, 1);
4223    }
4224
4225    #[test]
4226    fn on_disk_instances_preserve_exact_aggregates_and_atomic_sealing() {
4227        let result = (|| -> Result<(), Box<dyn Error>> {
4228            let database = test_database()?;
4229            let policy = TelemetryRetentionPolicy::default();
4230            let first = instance(1)?;
4231            let second = instance(2)?;
4232            record_transaction(
4233                &database.connection,
4234                database.project,
4235                first,
4236                UsageInstanceOwner::McpProcess,
4237                &event("agent", 100, 10),
4238                policy,
4239                false,
4240            )?;
4241            record_transaction(
4242                &database.connection,
4243                database.project,
4244                first,
4245                UsageInstanceOwner::McpProcess,
4246                &event("agent", 100, 95),
4247                policy,
4248                true,
4249            )?;
4250            record_transaction(
4251                &database.connection,
4252                database.project,
4253                second,
4254                UsageInstanceOwner::McpProcess,
4255                &event("agent", 40, 10),
4256                policy,
4257                false,
4258            )?;
4259
4260            let overview =
4261                token_overview_for_project(&database.connection, database.project, Some("agent"))?;
4262            assert_eq!(overview.calls, 3);
4263            assert_eq!(overview.deduped_modeled_tokens_avoided, 25);
4264            assert_eq!(overview.repeated_baselines_deduped, 1);
4265            assert_eq!(
4266                overview.detail_availability,
4267                UsageDetailAvailability::Retained
4268            );
4269            assert_eq!(
4270                scalar_count(
4271                    &database.connection,
4272                    "SELECT COUNT(*) FROM usage_instances WHERE caller_label = 'agent'",
4273                )?,
4274                2
4275            );
4276            assert_eq!(
4277                scalar_count(
4278                    &database.connection,
4279                    "SELECT COUNT(*) FROM usage_instance_baselines",
4280                )?,
4281                1
4282            );
4283
4284            let inactive = record_transaction(
4285                &database.connection,
4286                database.project,
4287                first,
4288                UsageInstanceOwner::McpProcess,
4289                &event("agent", 20, 5),
4290                policy,
4291                false,
4292            );
4293            assert!(matches!(inactive, Err(DbError::TelemetryInstanceInactive)));
4294            let mismatched = record_transaction(
4295                &database.connection,
4296                database.project,
4297                second,
4298                UsageInstanceOwner::CliInvocation,
4299                &event("agent", 20, 5),
4300                policy,
4301                false,
4302            );
4303            assert!(matches!(
4304                mismatched,
4305                Err(DbError::TelemetryInstanceMismatch)
4306            ));
4307
4308            let database_path = database.temp.path().join("projectatlas.db");
4309            drop(database.connection);
4310            let reopened = Connection::open(database_path)?;
4311            crate::schema::initialize(&reopened, None)?;
4312            let reopened_overview =
4313                token_overview_for_project(&reopened, database.project, Some("agent"))?;
4314            assert_eq!(reopened_overview.calls, overview.calls);
4315            assert_eq!(
4316                reopened_overview.deduped_modeled_tokens_avoided,
4317                overview.deduped_modeled_tokens_avoided
4318            );
4319            Ok(())
4320        })();
4321        assert!(result.is_ok(), "on-disk telemetry test failed: {result:?}");
4322    }
4323
4324    #[test]
4325    fn bounded_raw_and_label_retention_preserve_global_truth() {
4326        let result = (|| -> Result<(), Box<dyn Error>> {
4327            let database = test_database()?;
4328            let mut policy = TelemetryRetentionPolicy {
4329                max_raw_rows: 2,
4330                prune_batch_rows: 1,
4331                max_retained_labels: 1,
4332                ..TelemetryRetentionPolicy::default()
4333            };
4334            policy.checkpoint_write_interval = usize::MAX;
4335            let first = instance(3)?;
4336            record_transaction(
4337                &database.connection,
4338                database.project,
4339                first,
4340                UsageInstanceOwner::McpProcess,
4341                &event("first", 100, 10),
4342                policy,
4343                true,
4344            )?;
4345            let second = instance(4)?;
4346            record_transaction(
4347                &database.connection,
4348                database.project,
4349                second,
4350                UsageInstanceOwner::McpProcess,
4351                &event("second", 80, 20),
4352                policy,
4353                false,
4354            )?;
4355            record_transaction(
4356                &database.connection,
4357                database.project,
4358                second,
4359                UsageInstanceOwner::McpProcess,
4360                &event("second", 70, 20),
4361                policy,
4362                false,
4363            )?;
4364
4365            let state = retention_state_for_project(&database.connection, database.project)?;
4366            assert_eq!(state.raw_rows, 2);
4367            assert_eq!(state.retained_label_rows, 1);
4368            assert_eq!(state.label_tombstone_rows, 1);
4369            assert_eq!(state.spill_cleanup, SpillCleanupState::NotApplicable);
4370            let global = token_overview_for_project(&database.connection, database.project, None)?;
4371            assert_eq!(global.calls, 3);
4372            assert_eq!(global.deduped_modeled_tokens_avoided, 130);
4373            let expired =
4374                token_overview_for_project(&database.connection, database.project, Some("first"))?;
4375            assert_eq!(
4376                expired.detail_availability,
4377                UsageDetailAvailability::Expired
4378            );
4379            let unavailable = token_overview_for_project(
4380                &database.connection,
4381                database.project,
4382                Some("never-recorded"),
4383            )?;
4384            assert_eq!(
4385                unavailable.detail_availability,
4386                UsageDetailAvailability::Unavailable
4387            );
4388            let retained =
4389                token_overview_for_project(&database.connection, database.project, Some("second"))?;
4390            assert_eq!(retained.calls, 2);
4391            assert_eq!(
4392                retained.detail_availability,
4393                UsageDetailAvailability::Partial
4394            );
4395
4396            record_transaction(
4397                &database.connection,
4398                database.project,
4399                second,
4400                UsageInstanceOwner::McpProcess,
4401                &event("second", 60, 20),
4402                policy,
4403                false,
4404            )?;
4405            assert_eq!(
4406                token_overview_for_project(&database.connection, database.project, Some("second"))?
4407                    .calls,
4408                3
4409            );
4410            Ok(())
4411        })();
4412        assert!(
4413            result.is_ok(),
4414            "bounded telemetry retention test failed: {result:?}"
4415        );
4416    }
4417
4418    #[test]
4419    fn raw_age_and_logical_byte_retention_preserve_exact_aggregates() {
4420        let result = (|| -> Result<(), Box<dyn Error>> {
4421            let aged = test_database()?;
4422            let aged_policy = TelemetryRetentionPolicy {
4423                max_raw_age_seconds: 10,
4424                prune_batch_rows: 1,
4425                checkpoint_write_interval: usize::MAX,
4426                ..TelemetryRetentionPolicy::default()
4427            };
4428            let aged_runtime = instance(20)?;
4429            record_transaction_at(
4430                &aged.connection,
4431                aged.project,
4432                aged_runtime,
4433                UsageInstanceOwner::McpProcess,
4434                &event("aged", 100, 10),
4435                aged_policy,
4436                false,
4437                1_000,
4438            )?;
4439            record_transaction_at(
4440                &aged.connection,
4441                aged.project,
4442                aged_runtime,
4443                UsageInstanceOwner::McpProcess,
4444                &event("aged", 100, 10),
4445                aged_policy,
4446                false,
4447                1_020,
4448            )?;
4449            let aged_state = retention_state_for_project(&aged.connection, aged.project)?;
4450            assert_eq!(aged_state.raw_rows, 1);
4451            assert_eq!(aged_state.pruned_raw_rows, 1);
4452            let aged_overview =
4453                token_overview_for_project(&aged.connection, aged.project, Some("aged"))?;
4454            assert_eq!(aged_overview.calls, 2);
4455            assert_eq!(aged_overview.deduped_modeled_tokens_avoided, 80);
4456            assert_eq!(
4457                aged_overview.detail_availability,
4458                UsageDetailAvailability::Partial
4459            );
4460
4461            let byte_bounded = test_database()?;
4462            let sample = event("bytes", 100, 10);
4463            let one_event_bytes = logical_event_bytes(&sample, Some("bytes"))?;
4464            let byte_policy = TelemetryRetentionPolicy {
4465                max_raw_logical_bytes: one_event_bytes + 1,
4466                prune_batch_rows: 1,
4467                checkpoint_write_interval: usize::MAX,
4468                ..TelemetryRetentionPolicy::default()
4469            };
4470            let byte_runtime = instance(21)?;
4471            record_transaction_at(
4472                &byte_bounded.connection,
4473                byte_bounded.project,
4474                byte_runtime,
4475                UsageInstanceOwner::McpProcess,
4476                &sample,
4477                byte_policy,
4478                false,
4479                2_000,
4480            )?;
4481            record_transaction_at(
4482                &byte_bounded.connection,
4483                byte_bounded.project,
4484                byte_runtime,
4485                UsageInstanceOwner::McpProcess,
4486                &sample,
4487                byte_policy,
4488                false,
4489                2_001,
4490            )?;
4491            let byte_state =
4492                retention_state_for_project(&byte_bounded.connection, byte_bounded.project)?;
4493            assert_eq!(byte_state.raw_rows, 1);
4494            assert!(byte_state.raw_logical_bytes <= byte_policy.max_raw_logical_bytes);
4495            assert_eq!(byte_state.pruned_raw_rows, 1);
4496            let byte_overview = token_overview_for_project(
4497                &byte_bounded.connection,
4498                byte_bounded.project,
4499                Some("bytes"),
4500            )?;
4501            assert_eq!(byte_overview.calls, 2);
4502            assert_eq!(byte_overview.deduped_modeled_tokens_avoided, 80);
4503            Ok(())
4504        })();
4505        assert!(
4506            result.is_ok(),
4507            "raw telemetry budget test failed: {result:?}"
4508        );
4509    }
4510
4511    #[test]
4512    fn production_batch_raw_retention_removes_only_the_required_oldest_prefix() {
4513        let result = (|| -> Result<(), Box<dyn Error>> {
4514            let aged = test_database()?;
4515            let aged_policy = TelemetryRetentionPolicy {
4516                max_raw_rows: 10,
4517                max_raw_age_seconds: 10,
4518                prune_batch_rows: TelemetryRetentionPolicy::default().prune_batch_rows,
4519                checkpoint_write_interval: usize::MAX,
4520                ..TelemetryRetentionPolicy::default()
4521            };
4522            let runtime = instance(22)?;
4523            for now in [1_000, 1_020, 1_021] {
4524                record_transaction_at(
4525                    &aged.connection,
4526                    aged.project,
4527                    runtime,
4528                    UsageInstanceOwner::McpProcess,
4529                    &event("aged-prefix", 100, 10),
4530                    aged_policy,
4531                    false,
4532                    now,
4533                )?;
4534            }
4535            let epochs = aged
4536                .connection
4537                .prepare("SELECT created_at_epoch FROM usage_events ORDER BY created_at_epoch")?
4538                .query_map([], |row| row.get::<_, i64>(0))?
4539                .collect::<Result<Vec<_>, _>>()?;
4540            assert_eq!(epochs, vec![1_020, 1_021]);
4541
4542            let capped = test_database()?;
4543            let capped_policy = TelemetryRetentionPolicy {
4544                max_raw_rows: 2,
4545                prune_batch_rows: TelemetryRetentionPolicy::default().prune_batch_rows,
4546                checkpoint_write_interval: usize::MAX,
4547                ..TelemetryRetentionPolicy::default()
4548            };
4549            let runtime = instance(23)?;
4550            for now in [2_000, 2_001, 2_002] {
4551                record_transaction_at(
4552                    &capped.connection,
4553                    capped.project,
4554                    runtime,
4555                    UsageInstanceOwner::McpProcess,
4556                    &event("row-prefix", 100, 10),
4557                    capped_policy,
4558                    false,
4559                    now,
4560                )?;
4561            }
4562            let epochs = capped
4563                .connection
4564                .prepare("SELECT created_at_epoch FROM usage_events ORDER BY created_at_epoch")?
4565                .query_map([], |row| row.get::<_, i64>(0))?
4566                .collect::<Result<Vec<_>, _>>()?;
4567            assert_eq!(epochs, vec![2_001, 2_002]);
4568            assert_eq!(
4569                retention_state_for_project(&capped.connection, capped.project)?.pruned_raw_rows,
4570                1
4571            );
4572            Ok(())
4573        })();
4574        assert!(
4575            result.is_ok(),
4576            "exact raw-prefix retention test failed: {result:?}"
4577        );
4578    }
4579
4580    #[test]
4581    fn production_batch_tombstone_retention_removes_only_exact_excess() {
4582        let result = (|| -> Result<(), Box<dyn Error>> {
4583            let database = test_database()?;
4584            let policy = TelemetryRetentionPolicy {
4585                max_label_tombstones: 2,
4586                max_instance_tombstones: 2,
4587                prune_batch_rows: TelemetryRetentionPolicy::default().prune_batch_rows,
4588                checkpoint_write_interval: usize::MAX,
4589                ..TelemetryRetentionPolicy::default()
4590            };
4591            for (offset, label) in ["first", "second", "third"].into_iter().enumerate() {
4592                upsert_label_tombstone(
4593                    &database.connection,
4594                    database.project.as_bytes().as_slice(),
4595                    label,
4596                    1_000 + i64::try_from(offset)?,
4597                    None,
4598                )?;
4599            }
4600            for (offset, runtime) in [31_u8, 32, 33].into_iter().enumerate() {
4601                database.connection.execute(
4602                    "INSERT INTO usage_instance_tombstones(
4603                         project_instance_id, runtime_instance_id, retired_at_epoch
4604                     ) VALUES(?1, ?2, ?3)",
4605                    params![
4606                        database.project.as_bytes().as_slice(),
4607                        instance(runtime)?.as_bytes().as_slice(),
4608                        1_000 + i64::try_from(offset)?,
4609                    ],
4610                )?;
4611            }
4612            increment_retention_counter(
4613                &database.connection,
4614                RetentionCounter::InstanceTombstoneRows,
4615                3,
4616            )?;
4617            let deleted = prune_tombstones_once(&database.connection, policy, 1_100)?;
4618            assert_eq!(deleted, 2);
4619            assert_eq!(
4620                scalar_count(
4621                    &database.connection,
4622                    "SELECT COUNT(*) FROM usage_label_tombstones",
4623                )?,
4624                2
4625            );
4626            assert_eq!(
4627                scalar_count(
4628                    &database.connection,
4629                    "SELECT COUNT(*) FROM usage_instance_tombstones",
4630                )?,
4631                2
4632            );
4633            assert_eq!(
4634                database.connection.query_row(
4635                    "SELECT caller_label FROM usage_label_tombstones
4636                     ORDER BY expired_at_epoch LIMIT 1",
4637                    [],
4638                    |row| row.get::<_, String>(0),
4639                )?,
4640                "second"
4641            );
4642            Ok(())
4643        })();
4644        assert!(
4645            result.is_ok(),
4646            "exact tombstone retention test failed: {result:?}"
4647        );
4648    }
4649
4650    #[test]
4651    fn daily_capacity_prunes_old_history_before_reserving_both_current_rows() {
4652        let result = (|| -> Result<(), Box<dyn Error>> {
4653            let database = test_database()?;
4654            let policy = TelemetryRetentionPolicy {
4655                max_daily_rows: 4,
4656                retained_trend_days: 1,
4657                prune_batch_rows: TelemetryRetentionPolicy::default().prune_batch_rows,
4658                checkpoint_write_interval: usize::MAX,
4659                ..TelemetryRetentionPolicy::default()
4660            };
4661            let runtime = instance(24)?;
4662            let mut first = event("daily", 100, 10);
4663            first.provider = "first-provider".to_string();
4664            let mut second = event("daily", 90, 10);
4665            second.provider = "second-provider".to_string();
4666            for value in [&first, &second] {
4667                record_transaction_at(
4668                    &database.connection,
4669                    database.project,
4670                    runtime,
4671                    UsageInstanceOwner::McpProcess,
4672                    value,
4673                    policy,
4674                    false,
4675                    SECONDS_PER_DAY,
4676                )?;
4677            }
4678            assert_eq!(
4679                retention_counter(&database.connection, RetentionCounter::DailyRows)?,
4680                4
4681            );
4682
4683            let mut current = event("daily", 80, 10);
4684            current.provider = "current-provider".to_string();
4685            let current_epoch = 3 * SECONDS_PER_DAY;
4686            record_transaction_at(
4687                &database.connection,
4688                database.project,
4689                runtime,
4690                UsageInstanceOwner::McpProcess,
4691                &current,
4692                policy,
4693                false,
4694                current_epoch,
4695            )?;
4696            assert_eq!(
4697                retention_counter(&database.connection, RetentionCounter::DailyRows)?,
4698                2
4699            );
4700            let current_rows = database.connection.query_row(
4701                "SELECT
4702                     (SELECT COUNT(*) FROM usage_daily_aggregates AS aggregate
4703                      JOIN usage_bucket_dimensions AS dimension USING(dimension_id)
4704                      WHERE aggregate.day_epoch = ?1 AND dimension.provider = ?2),
4705                     (SELECT COUNT(*) FROM usage_instance_daily_aggregates AS aggregate
4706                      JOIN usage_bucket_dimensions AS dimension USING(dimension_id)
4707                      WHERE aggregate.day_epoch = ?1 AND dimension.provider = ?2)",
4708                params![current_epoch, "current-provider"],
4709                |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
4710            )?;
4711            assert_eq!(current_rows, (1, 1));
4712            let trends = token_trends_for_project(
4713                &database.connection,
4714                database.project,
4715                Some("daily"),
4716                TokenTrendWindow::Day,
4717            )?;
4718            assert_eq!(trends.periods.len(), 1);
4719            assert_eq!(trends.periods[0].calls, 1);
4720            Ok(())
4721        })();
4722        assert!(
4723            result.is_ok(),
4724            "daily capacity reservation test failed: {result:?}"
4725        );
4726    }
4727
4728    #[test]
4729    fn label_reports_group_many_instances_inside_sqlite() {
4730        let result = (|| -> Result<(), Box<dyn Error>> {
4731            let database = test_database()?;
4732            let policy = TelemetryRetentionPolicy {
4733                checkpoint_write_interval: usize::MAX,
4734                ..TelemetryRetentionPolicy::default()
4735            };
4736            for runtime in 1_u8..=48 {
4737                record_transaction_at(
4738                    &database.connection,
4739                    database.project,
4740                    instance(runtime)?,
4741                    UsageInstanceOwner::McpProcess,
4742                    &event("grouped", 100, 10),
4743                    policy,
4744                    true,
4745                    10 * SECONDS_PER_DAY,
4746                )?;
4747            }
4748            let overview = token_overview_for_project(
4749                &database.connection,
4750                database.project,
4751                Some("grouped"),
4752            )?;
4753            assert_eq!(overview.calls, 48);
4754            assert_eq!(overview.buckets.len(), 1);
4755            let trends = token_trends_for_project(
4756                &database.connection,
4757                database.project,
4758                Some("grouped"),
4759                TokenTrendWindow::Day,
4760            )?;
4761            assert_eq!(trends.periods.len(), 1);
4762            assert_eq!(trends.periods[0].calls, 48);
4763            assert_eq!(trends.periods[0].buckets.len(), 1);
4764            Ok(())
4765        })();
4766        assert!(
4767            result.is_ok(),
4768            "SQLite-side telemetry grouping test failed: {result:?}"
4769        );
4770    }
4771
4772    #[test]
4773    fn new_project_runtime_reclaims_inactive_instance_and_label_capacity() {
4774        let result = (|| -> Result<(), Box<dyn Error>> {
4775            let database = test_database()?;
4776            let policy = TelemetryRetentionPolicy {
4777                max_active_instances: 1,
4778                max_retained_instances: 1,
4779                max_retained_labels: 1,
4780                ..TelemetryRetentionPolicy::default()
4781            };
4782            let old_project = database.project;
4783            let old_runtime = instance(30)?;
4784            record_transaction(
4785                &database.connection,
4786                old_project,
4787                old_runtime,
4788                UsageInstanceOwner::McpProcess,
4789                &event("old-project", 100, 10),
4790                policy,
4791                false,
4792            )?;
4793            assert_eq!(
4794                seal_project_usage_instances(&database.connection, old_project)?,
4795                1
4796            );
4797            let sealed_state = retention_state_for_project(&database.connection, old_project)?;
4798            assert_eq!(sealed_state.baseline_rows, 0);
4799
4800            let new_project = ProjectInstanceId::from_bytes([31; 16])?;
4801            database.connection.execute(
4802                "UPDATE project_identity SET project_instance_id = ?1 WHERE singleton = 1",
4803                [new_project.as_bytes().as_slice()],
4804            )?;
4805            let new_runtime = old_runtime;
4806            record_transaction(
4807                &database.connection,
4808                new_project,
4809                new_runtime,
4810                UsageInstanceOwner::McpProcess,
4811                &event("new-project", 80, 20),
4812                policy,
4813                false,
4814            )?;
4815
4816            assert_eq!(
4817                scalar_count(&database.connection, "SELECT COUNT(*) FROM usage_instances")?,
4818                1
4819            );
4820            assert_eq!(
4821                scalar_count(&database.connection, "SELECT COUNT(*) FROM usage_labels")?,
4822                1
4823            );
4824            assert_eq!(
4825                scalar_count(
4826                    &database.connection,
4827                    "SELECT COUNT(*) FROM usage_global_aggregates",
4828                )?,
4829                1
4830            );
4831            assert_eq!(
4832                scalar_count(
4833                    &database.connection,
4834                    "SELECT COUNT(*) FROM usage_labels WHERE caller_label = 'new-project'",
4835                )?,
4836                1
4837            );
4838            let old_instance_tombstones = database.connection.query_row(
4839                "SELECT COUNT(*) FROM usage_instance_tombstones
4840             WHERE project_instance_id = ?1",
4841                [old_project.as_bytes().as_slice()],
4842                |row| row.get::<_, i64>(0),
4843            )?;
4844            assert_eq!(old_instance_tombstones, 1);
4845            let old_label_tombstones = database.connection.query_row(
4846                "SELECT COUNT(*) FROM usage_label_tombstones
4847             WHERE project_instance_id = ?1",
4848                [old_project.as_bytes().as_slice()],
4849                |row| row.get::<_, i64>(0),
4850            )?;
4851            assert_eq!(old_label_tombstones, 1);
4852            let overview = token_overview_for_project(&database.connection, new_project, None)?;
4853            assert_eq!(overview.calls, 1);
4854            Ok(())
4855        })();
4856        assert!(
4857            result.is_ok(),
4858            "project-rotation telemetry capacity test failed: {result:?}"
4859        );
4860    }
4861
4862    #[test]
4863    fn active_labels_and_reserved_overflow_dimension_remain_disjoint() {
4864        let result = (|| -> Result<(), Box<dyn Error>> {
4865            let database = test_database()?;
4866            let policy = TelemetryRetentionPolicy {
4867                max_retained_labels: 1,
4868                max_dimensions: 2,
4869                ..TelemetryRetentionPolicy::default()
4870            };
4871            let active = instance(5)?;
4872            let mut reserved = event("active", 100, 10);
4873            reserved.token_savings_bucket = OVERFLOW_DIMENSION.to_string();
4874            reserved.provider = OVERFLOW_DIMENSION.to_string();
4875            reserved.model = OVERFLOW_DIMENSION.to_string();
4876            reserved.tokenizer_backend = OVERFLOW_DIMENSION.to_string();
4877            reserved.accuracy = OVERFLOW_DIMENSION.to_string();
4878            reserved.baseline_kind = OVERFLOW_DIMENSION.to_string();
4879            reserved.confidence = OVERFLOW_DIMENSION.to_string();
4880            reserved.accounting_layer = OVERFLOW_DIMENSION.to_string();
4881            reserved.estimate_method = OVERFLOW_DIMENSION.to_string();
4882            reserved.denominator_kind = OVERFLOW_DIMENSION.to_string();
4883            reserved.dedupe_scope = OVERFLOW_DIMENSION.to_string();
4884            record_transaction(
4885                &database.connection,
4886                database.project,
4887                active,
4888                UsageInstanceOwner::McpProcess,
4889                &reserved,
4890                policy,
4891                false,
4892            )?;
4893            assert_eq!(
4894                scalar_count(
4895                    &database.connection,
4896                    "SELECT COUNT(*) FROM usage_bucket_dimensions
4897                 WHERE token_savings_bucket = '<overflow>'",
4898                )?,
4899                2
4900            );
4901            assert_eq!(
4902                scalar_count(
4903                    &database.connection,
4904                    "SELECT COUNT(DISTINCT overflow) FROM usage_bucket_dimensions
4905                 WHERE token_savings_bucket = '<overflow>'",
4906                )?,
4907                2
4908            );
4909
4910            let rejected = record_transaction(
4911                &database.connection,
4912                database.project,
4913                instance(6)?,
4914                UsageInstanceOwner::McpProcess,
4915                &event("replacement", 50, 10),
4916                policy,
4917                false,
4918            );
4919            assert!(matches!(rejected, Err(DbError::TelemetryInstanceCapacity)));
4920            assert_eq!(
4921                scalar_count(
4922                    &database.connection,
4923                    "SELECT COUNT(*) FROM usage_labels WHERE caller_label = 'active'",
4924                )?,
4925                1
4926            );
4927            assert_eq!(
4928                scalar_count(
4929                    &database.connection,
4930                    "SELECT COUNT(*) FROM usage_label_tombstones",
4931                )?,
4932                0
4933            );
4934            Ok(())
4935        })();
4936        assert!(
4937            result.is_ok(),
4938            "active-label overflow sentinel test failed: {result:?}"
4939        );
4940    }
4941
4942    #[test]
4943    fn dimension_capacity_uses_reserved_overflow_and_reports_partial_detail() {
4944        let result = (|| -> Result<(), Box<dyn Error>> {
4945            let database = test_database()?;
4946            let policy = TelemetryRetentionPolicy {
4947                max_dimensions: 2,
4948                ..TelemetryRetentionPolicy::default()
4949            };
4950            let runtime = instance(35)?;
4951            record_transaction(
4952                &database.connection,
4953                database.project,
4954                runtime,
4955                UsageInstanceOwner::McpProcess,
4956                &event("dimensions", 100, 10),
4957                policy,
4958                false,
4959            )?;
4960            let mut overflowed = event("dimensions", 50, 10);
4961            overflowed.provider = "another-provider".to_string();
4962            overflowed.baseline_identity = "source:src/other.rs".to_string();
4963            overflowed.baseline_fingerprint = "source:src/other.rs:v1".to_string();
4964            record_transaction(
4965                &database.connection,
4966                database.project,
4967                runtime,
4968                UsageInstanceOwner::McpProcess,
4969                &overflowed,
4970                policy,
4971                false,
4972            )?;
4973
4974            assert_eq!(
4975                scalar_count(
4976                    &database.connection,
4977                    "SELECT COUNT(*) FROM usage_bucket_dimensions",
4978                )?,
4979                2
4980            );
4981            let overflow_calls = database.connection.query_row(
4982                "SELECT aggregate.calls
4983             FROM usage_global_aggregates AS aggregate
4984             JOIN usage_bucket_dimensions AS dimension USING(dimension_id)
4985             WHERE aggregate.project_instance_id = ?1 AND dimension.overflow = 1",
4986                [database.project.as_bytes().as_slice()],
4987                |row| row.get::<_, i64>(0),
4988            )?;
4989            assert_eq!(overflow_calls, 1);
4990            let overview = token_overview_for_project(
4991                &database.connection,
4992                database.project,
4993                Some("dimensions"),
4994            )?;
4995            assert_eq!(overview.calls, 2);
4996            assert_eq!(
4997                overview.detail_availability,
4998                UsageDetailAvailability::Partial
4999            );
5000            Ok(())
5001        })();
5002        assert!(
5003            result.is_ok(),
5004            "dimension overflow retention test failed: {result:?}"
5005        );
5006    }
5007
5008    #[test]
5009    fn rejected_events_and_read_only_reports_leave_storage_unchanged() {
5010        let result = (|| -> Result<(), Box<dyn Error>> {
5011            let database = test_database()?;
5012            let policy = TelemetryRetentionPolicy::default();
5013            let runtime = instance(7)?;
5014            record_transaction(
5015                &database.connection,
5016                database.project,
5017                runtime,
5018                UsageInstanceOwner::McpProcess,
5019                &event("reader", 100, 10),
5020                policy,
5021                false,
5022            )?;
5023            let counters_before = database.connection.query_row(
5024                "SELECT raw_rows, baseline_rows, instance_rows, daily_rows,
5025                    writes_since_checkpoint
5026             FROM usage_retention_state WHERE singleton = 1",
5027                [],
5028                |row| {
5029                    Ok((
5030                        row.get::<_, i64>(0)?,
5031                        row.get::<_, i64>(1)?,
5032                        row.get::<_, i64>(2)?,
5033                        row.get::<_, i64>(3)?,
5034                        row.get::<_, i64>(4)?,
5035                    ))
5036                },
5037            )?;
5038            let mut oversized = event("reader", 10, 5);
5039            oversized.query = Some("x".repeat(policy.max_query_bytes + 1));
5040            let rejected = record_transaction(
5041                &database.connection,
5042                database.project,
5043                runtime,
5044                UsageInstanceOwner::McpProcess,
5045                &oversized,
5046                policy,
5047                false,
5048            );
5049            assert!(matches!(
5050                rejected,
5051                Err(DbError::TelemetryFieldTooLarge { .. })
5052            ));
5053            let _ =
5054                usage_events_for_project(&database.connection, database.project, Some("reader"))?;
5055            let _ =
5056                token_overview_for_project(&database.connection, database.project, Some("reader"))?;
5057            let _ = token_trends_for_project(
5058                &database.connection,
5059                database.project,
5060                Some("reader"),
5061                TokenTrendWindow::Day,
5062            )?;
5063            let state = retention_state_for_project(&database.connection, database.project)?;
5064            assert_eq!(state.spill_cleanup, SpillCleanupState::NotApplicable);
5065            let counters_after = database.connection.query_row(
5066                "SELECT raw_rows, baseline_rows, instance_rows, daily_rows,
5067                    writes_since_checkpoint
5068             FROM usage_retention_state WHERE singleton = 1",
5069                [],
5070                |row| {
5071                    Ok((
5072                        row.get::<_, i64>(0)?,
5073                        row.get::<_, i64>(1)?,
5074                        row.get::<_, i64>(2)?,
5075                        row.get::<_, i64>(3)?,
5076                        row.get::<_, i64>(4)?,
5077                    ))
5078                },
5079            )?;
5080            assert_eq!(counters_after, counters_before);
5081
5082            let storage_before = database.connection.query_row(
5083                "SELECT
5084                 (SELECT COUNT(*) FROM usage_events),
5085                 (SELECT COUNT(*) FROM usage_instance_baselines),
5086                 (SELECT calls FROM usage_global_aggregates),
5087                 raw_rows,
5088                 baseline_rows
5089             FROM usage_retention_state WHERE singleton = 1",
5090                [],
5091                |row| {
5092                    Ok((
5093                        row.get::<_, i64>(0)?,
5094                        row.get::<_, i64>(1)?,
5095                        row.get::<_, i64>(2)?,
5096                        row.get::<_, i64>(3)?,
5097                        row.get::<_, i64>(4)?,
5098                    ))
5099                },
5100            )?;
5101            database.connection.execute_batch(
5102                "CREATE TEMP TRIGGER reject_telemetry_aggregate_update
5103             BEFORE UPDATE ON usage_global_aggregates
5104             BEGIN
5105                 SELECT RAISE(ABORT, 'injected telemetry aggregate failure');
5106             END;",
5107            )?;
5108            let mut late_failure = event("reader", 200, 20);
5109            late_failure.baseline_identity = "source:src/lib.rs:second".to_string();
5110            late_failure.baseline_fingerprint = "source:src/lib.rs:v2".to_string();
5111            assert!(
5112                record_transaction(
5113                    &database.connection,
5114                    database.project,
5115                    runtime,
5116                    UsageInstanceOwner::McpProcess,
5117                    &late_failure,
5118                    policy,
5119                    false,
5120                )
5121                .is_err()
5122            );
5123            database
5124                .connection
5125                .execute_batch("DROP TRIGGER reject_telemetry_aggregate_update")?;
5126            let storage_after = database.connection.query_row(
5127                "SELECT
5128                 (SELECT COUNT(*) FROM usage_events),
5129                 (SELECT COUNT(*) FROM usage_instance_baselines),
5130                 (SELECT calls FROM usage_global_aggregates),
5131                 raw_rows,
5132                 baseline_rows
5133             FROM usage_retention_state WHERE singleton = 1",
5134                [],
5135                |row| {
5136                    Ok((
5137                        row.get::<_, i64>(0)?,
5138                        row.get::<_, i64>(1)?,
5139                        row.get::<_, i64>(2)?,
5140                        row.get::<_, i64>(3)?,
5141                        row.get::<_, i64>(4)?,
5142                    ))
5143                },
5144            )?;
5145            assert_eq!(storage_after, storage_before);
5146
5147            let lookalike = database.temp.path().join("usage-telemetry-spill.db");
5148            std::fs::write(&lookalike, b"not owned by ProjectAtlas")?;
5149            maintain_after_commit_for_project(&database.connection, database.project, policy)?;
5150            assert_eq!(std::fs::read(lookalike)?, b"not owned by ProjectAtlas");
5151            Ok(())
5152        })();
5153        assert!(
5154            result.is_ok(),
5155            "telemetry rejection rollback test failed: {result:?}"
5156        );
5157    }
5158
5159    #[test]
5160    fn baseline_capacity_collision_and_integer_overflow_roll_back_completely() {
5161        let result = (|| -> Result<(), Box<dyn Error>> {
5162            let database = test_database()?;
5163            let policy = TelemetryRetentionPolicy {
5164                max_baselines_per_instance: 1,
5165                ..TelemetryRetentionPolicy::default()
5166            };
5167
5168            let capacity_runtime = instance(40)?;
5169            record_transaction(
5170                &database.connection,
5171                database.project,
5172                capacity_runtime,
5173                UsageInstanceOwner::McpProcess,
5174                &event("capacity", 100, 10),
5175                policy,
5176                false,
5177            )?;
5178            let capacity_before = database.connection.query_row(
5179                "SELECT
5180                 (SELECT COUNT(*) FROM usage_events),
5181                 (SELECT COUNT(*) FROM usage_instance_baselines),
5182                 (SELECT SUM(calls) FROM usage_global_aggregates),
5183                 raw_rows,
5184                 baseline_rows
5185             FROM usage_retention_state WHERE singleton = 1",
5186                [],
5187                |row| {
5188                    Ok((
5189                        row.get::<_, i64>(0)?,
5190                        row.get::<_, i64>(1)?,
5191                        row.get::<_, i64>(2)?,
5192                        row.get::<_, i64>(3)?,
5193                        row.get::<_, i64>(4)?,
5194                    ))
5195                },
5196            )?;
5197            let mut second_baseline = event("capacity", 90, 10);
5198            second_baseline.baseline_identity = "source:src/other.rs".to_string();
5199            second_baseline.baseline_fingerprint = "source:src/other.rs:v1".to_string();
5200            assert!(matches!(
5201                record_transaction(
5202                    &database.connection,
5203                    database.project,
5204                    capacity_runtime,
5205                    UsageInstanceOwner::McpProcess,
5206                    &second_baseline,
5207                    policy,
5208                    false,
5209                ),
5210                Err(DbError::TelemetryBaselineCapacity)
5211            ));
5212            let capacity_after = database.connection.query_row(
5213                "SELECT
5214                 (SELECT COUNT(*) FROM usage_events),
5215                 (SELECT COUNT(*) FROM usage_instance_baselines),
5216                 (SELECT SUM(calls) FROM usage_global_aggregates),
5217                 raw_rows,
5218                 baseline_rows
5219             FROM usage_retention_state WHERE singleton = 1",
5220                [],
5221                |row| {
5222                    Ok((
5223                        row.get::<_, i64>(0)?,
5224                        row.get::<_, i64>(1)?,
5225                        row.get::<_, i64>(2)?,
5226                        row.get::<_, i64>(3)?,
5227                        row.get::<_, i64>(4)?,
5228                    ))
5229                },
5230            )?;
5231            assert_eq!(capacity_after, capacity_before);
5232
5233            let collision_runtime = instance(41)?;
5234            let collision_event = event("collision", 60, 10);
5235            record_transaction(
5236                &database.connection,
5237                database.project,
5238                collision_runtime,
5239                UsageInstanceOwner::McpProcess,
5240                &collision_event,
5241                TelemetryRetentionPolicy::default(),
5242                false,
5243            )?;
5244            database.connection.execute(
5245                "UPDATE usage_instance_baselines
5246             SET baseline_identity = 'source:src/foo.rs'
5247             WHERE instance_row_id = (
5248                 SELECT instance_row_id FROM usage_instances
5249                 WHERE project_instance_id = ?1 AND runtime_instance_id = ?2
5250             )",
5251                params![
5252                    database.project.as_bytes().as_slice(),
5253                    collision_runtime.as_bytes().as_slice()
5254                ],
5255            )?;
5256            let collision_raw_before =
5257                scalar_count(&database.connection, "SELECT COUNT(*) FROM usage_events")?;
5258            assert!(matches!(
5259                record_transaction(
5260                    &database.connection,
5261                    database.project,
5262                    collision_runtime,
5263                    UsageInstanceOwner::McpProcess,
5264                    &collision_event,
5265                    TelemetryRetentionPolicy::default(),
5266                    false,
5267                ),
5268                Err(DbError::TelemetryBaselineCollision)
5269            ));
5270            assert_eq!(
5271                scalar_count(&database.connection, "SELECT COUNT(*) FROM usage_events")?,
5272                collision_raw_before
5273            );
5274
5275            let overflow_runtime = instance(42)?;
5276            record_transaction(
5277                &database.connection,
5278                database.project,
5279                overflow_runtime,
5280                UsageInstanceOwner::McpProcess,
5281                &event("overflow", 50, 10),
5282                TelemetryRetentionPolicy::default(),
5283                false,
5284            )?;
5285            database
5286                .connection
5287                .execute("UPDATE usage_global_aggregates SET calls = ?1", [i64::MAX])?;
5288            let overflow_raw_before =
5289                scalar_count(&database.connection, "SELECT COUNT(*) FROM usage_events")?;
5290            let overflow_baselines_before = scalar_count(
5291                &database.connection,
5292                "SELECT COUNT(*) FROM usage_instance_baselines",
5293            )?;
5294            assert!(matches!(
5295                record_transaction(
5296                    &database.connection,
5297                    database.project,
5298                    overflow_runtime,
5299                    UsageInstanceOwner::McpProcess,
5300                    &event("overflow", 50, 10),
5301                    TelemetryRetentionPolicy::default(),
5302                    false,
5303                ),
5304                Err(DbError::TelemetryIntegerOverflow {
5305                    field: AGGREGATE_COUNTER_FIELD
5306                })
5307            ));
5308            assert_eq!(
5309                scalar_count(&database.connection, "SELECT COUNT(*) FROM usage_events")?,
5310                overflow_raw_before
5311            );
5312            assert_eq!(
5313                scalar_count(
5314                    &database.connection,
5315                    "SELECT COUNT(*) FROM usage_instance_baselines",
5316                )?,
5317                overflow_baselines_before
5318            );
5319            Ok(())
5320        })();
5321        assert!(
5322            result.is_ok(),
5323            "telemetry baseline rollback test failed: {result:?}"
5324        );
5325    }
5326
5327    #[test]
5328    fn production_store_reports_live_page_and_connection_policy_state() {
5329        let result = (|| -> Result<(), Box<dyn Error>> {
5330            let mut database = production_database()?;
5331            let initial = database.store.telemetry_retention_state()?;
5332            assert!(initial.page_count > 0);
5333            assert_eq!(initial.checkpoint_state, TelemetryCheckpointState::NotDue);
5334            assert_eq!(
5335                initial.statistics_policy,
5336                PlannerStatisticsPolicy::NotConfigured
5337            );
5338            assert_eq!(
5339                initial.statistics_state,
5340                PlannerStatisticsState::NotInitialized
5341            );
5342            assert_eq!(
5343                initial.connection_busy_timeout_ms,
5344                initial.normal_busy_timeout_ms
5345            );
5346            assert_eq!(initial.normal_busy_timeout_ms, 5_000);
5347            assert_eq!(initial.telemetry_busy_timeout_ms, 25);
5348            assert!(initial.wal_autocheckpoint_pages > 0);
5349
5350            let nodes = (0..1_024)
5351                .map(|index| {
5352                    let path = format!("src/generated/{index:04}.rs");
5353                    Node {
5354                        parent_path: normalized_parent(&path),
5355                        path,
5356                        kind: NodeKind::File,
5357                        extension: Some(".rs".to_string()),
5358                        language: Some("rust".to_string()),
5359                        size_bytes: Some(4_096),
5360                        mtime_ns: Some(i64::from(index)),
5361                        content_hash: Some(format!("hash-{index:04}")),
5362                    }
5363                })
5364                .collect::<Vec<_>>();
5365            database.store.replace_scan(&nodes)?;
5366            let grown = database.store.telemetry_retention_state()?;
5367            let live_grown_pages = pragma_count(&database.store.connection, "page_count")?;
5368            assert_eq!(
5369                grown.page_count,
5370                count_usize("page_count", live_grown_pages)?
5371            );
5372            assert!(grown.page_count >= initial.page_count);
5373
5374            database.store.replace_scan(&[])?;
5375            let deleted = database.store.telemetry_retention_state()?;
5376            assert_eq!(
5377                deleted.freelist_pages,
5378                count_usize(
5379                    "freelist_pages",
5380                    pragma_count(&database.store.connection, "freelist_count")?,
5381                )?
5382            );
5383            assert_eq!(
5384                deleted.page_count,
5385                count_usize(
5386                    "page_count",
5387                    pragma_count(&database.store.connection, "page_count")?,
5388                )?
5389            );
5390            Ok(())
5391        })();
5392        assert!(
5393            result.is_ok(),
5394            "production page-policy report test failed: {result:?}"
5395        );
5396    }
5397
5398    #[test]
5399    fn production_store_checkpoint_retries_after_a_held_reader_and_rejects_stale_binding() {
5400        let result = (|| -> Result<(), Box<dyn Error>> {
5401            let database = production_database()?;
5402            let runtime = instance(90)?;
5403            database.store.record_usage_for_instance(
5404                runtime,
5405                UsageInstanceOwner::McpProcess,
5406                &event("checkpoint", 100, 10),
5407                false,
5408            )?;
5409            let reader =
5410                AtlasStore::open_read_only_for_project(&database.database_path, &database.root)?;
5411            let _: i64 =
5412                reader
5413                    .connection
5414                    .query_row("SELECT COUNT(*) FROM usage_events", [], |row| row.get(0))?;
5415            database.store.record_usage_for_instance(
5416                runtime,
5417                UsageInstanceOwner::McpProcess,
5418                &event("checkpoint", 90, 10),
5419                false,
5420            )?;
5421            let policy = TelemetryRetentionPolicy::default();
5422            database.store.connection.execute(
5423                "UPDATE usage_retention_state
5424                 SET writes_since_checkpoint = ?1 WHERE singleton = 1",
5425                [to_i64(
5426                    "checkpoint_write_interval",
5427                    policy.checkpoint_write_interval,
5428                )?],
5429            )?;
5430            maintain_after_commit_for_project(
5431                &database.store.connection,
5432                database.project,
5433                policy,
5434            )?;
5435            let busy = database.store.telemetry_retention_state()?;
5436            assert_eq!(busy.checkpoint_state, TelemetryCheckpointState::Busy);
5437            assert_eq!(
5438                busy.writes_since_checkpoint,
5439                policy.checkpoint_write_interval
5440            );
5441
5442            reader.finish_index_read_snapshot()?;
5443            maintain_after_commit_for_project(
5444                &database.store.connection,
5445                database.project,
5446                policy,
5447            )?;
5448            let completed = database.store.telemetry_retention_state()?;
5449            assert_eq!(
5450                completed.checkpoint_state,
5451                TelemetryCheckpointState::Completed
5452            );
5453            assert_eq!(completed.writes_since_checkpoint, 0);
5454
5455            let old_project = database.project;
5456            let captured = database.store.captured_project_binding()?;
5457            assert_eq!(captured.project_instance_id, old_project);
5458            let detached = AtlasStore::transition_project_root(
5459                &database.database_path,
5460                &database.root,
5461                crate::ProjectRootTransition::Detach,
5462            )?;
5463            assert_ne!(detached.project_instance_id, old_project);
5464            assert!(
5465                database
5466                    .store
5467                    .revalidate_captured_project_binding()
5468                    .is_err()
5469            );
5470            let stale =
5471                maintain_after_commit_for_project(&database.store.connection, old_project, policy);
5472            assert!(stale.is_err());
5473            Ok(())
5474        })();
5475        assert!(
5476            result.is_ok(),
5477            "production checkpoint lifecycle test failed: {result:?}"
5478        );
5479    }
5480
5481    #[test]
5482    fn production_store_reuses_freed_raw_pages_without_request_path_vacuum() {
5483        let result = (|| -> Result<(), Box<dyn Error>> {
5484            let database = production_database()?;
5485            let runtime = instance(91)?;
5486            let generous = TelemetryRetentionPolicy {
5487                max_raw_rows: 200,
5488                max_raw_logical_bytes: 2 * 1_024 * 1_024,
5489                checkpoint_write_interval: usize::MAX,
5490                ..TelemetryRetentionPolicy::default()
5491            };
5492            let mut large = event("page-reuse", 100, 10);
5493            large.query = Some("x".repeat(3_500));
5494            for now in 10_000..10_128 {
5495                record_transaction_at(
5496                    &database.store.connection,
5497                    database.project,
5498                    runtime,
5499                    UsageInstanceOwner::McpProcess,
5500                    &large,
5501                    generous,
5502                    false,
5503                    now,
5504                )?;
5505            }
5506            let allocated = database.store.telemetry_retention_state()?;
5507            let compact = TelemetryRetentionPolicy {
5508                max_raw_rows: 16,
5509                checkpoint_write_interval: usize::MAX,
5510                ..generous
5511            };
5512            let transaction = Transaction::new_unchecked(
5513                &database.store.connection,
5514                TransactionBehavior::Immediate,
5515            )?;
5516            let mut pruned = 0usize;
5517            loop {
5518                let deleted = prune_raw_once(&transaction, compact, 10_128)?;
5519                pruned = pruned
5520                    .checked_add(deleted)
5521                    .ok_or(DbError::TelemetryIntegerOverflow {
5522                        field: "pruned_raw_rows",
5523                    })?;
5524                if deleted == 0 {
5525                    break;
5526                }
5527            }
5528            refresh_retention_state(&transaction, compact, 10_128, pruned, 0, 0, 0)?;
5529            transaction.commit()?;
5530            let after_prune = database.store.telemetry_retention_state()?;
5531            assert_eq!(after_prune.raw_rows, 16);
5532            assert!(after_prune.freelist_pages > 0);
5533            assert_eq!(after_prune.page_count, allocated.page_count);
5534
5535            let refill = TelemetryRetentionPolicy {
5536                max_raw_rows: 80,
5537                ..compact
5538            };
5539            for now in 20_000..20_064 {
5540                record_transaction_at(
5541                    &database.store.connection,
5542                    database.project,
5543                    runtime,
5544                    UsageInstanceOwner::McpProcess,
5545                    &large,
5546                    refill,
5547                    false,
5548                    now,
5549                )?;
5550            }
5551            let after_refill = database.store.telemetry_retention_state()?;
5552            assert_eq!(after_refill.raw_rows, 80);
5553            assert!(after_refill.page_count <= allocated.page_count);
5554            assert!(after_refill.freelist_pages < after_prune.freelist_pages);
5555            Ok(())
5556        })();
5557        assert!(
5558            result.is_ok(),
5559            "production page-reuse test failed: {result:?}"
5560        );
5561    }
5562
5563    #[test]
5564    fn hot_telemetry_queries_use_owned_indexes_without_duplicate_primary_key_indexes() {
5565        let result = (|| -> Result<(), Box<dyn Error>> {
5566            let database = test_database()?;
5567            assert_plan_uses(
5568                &query_plan(
5569                    &database.connection,
5570                    "SELECT instance_row_id FROM usage_instances
5571                 WHERE project_instance_id = X'01010101010101010101010101010101'
5572                   AND runtime_instance_id = X'02020202020202020202020202020202'",
5573                )?,
5574                "sqlite_autoindex_usage_instances_1",
5575            );
5576            assert_plan_uses(
5577                &query_plan(
5578                    &database.connection,
5579                    "SELECT maximum_without FROM usage_instance_baselines
5580                 WHERE instance_row_id = 1
5581                   AND baseline_key = zeroblob(32)",
5582                )?,
5583                "PRIMARY KEY",
5584            );
5585            assert_plan_uses(
5586                &query_plan(
5587                    &database.connection,
5588                    "SELECT id FROM usage_events
5589                 WHERE created_at_epoch < 1
5590                 ORDER BY created_at_epoch, id LIMIT 1",
5591                )?,
5592                "idx_usage_created_at",
5593            );
5594            assert_plan_uses(
5595                &query_plan(
5596                    &database.connection,
5597                    "SELECT calls FROM usage_global_aggregates
5598                 WHERE project_instance_id = X'01010101010101010101010101010101'
5599                   AND dimension_id = 1",
5600                )?,
5601                "PRIMARY KEY",
5602            );
5603            assert_plan_uses(
5604                &query_plan(
5605                    &database.connection,
5606                    "SELECT instance_row_id FROM usage_instances
5607                 WHERE project_instance_id = X'01010101010101010101010101010101'
5608                   AND caller_label = 'agent' AND state = 'active'
5609                 ORDER BY started_at_epoch, instance_row_id LIMIT 1",
5610                )?,
5611                "idx_usage_instances_label_state",
5612            );
5613            assert_plan_uses(
5614                &query_plan(
5615                    &database.connection,
5616                    "SELECT project_instance_id FROM usage_daily_aggregates
5617                 WHERE day_epoch < 1 ORDER BY day_epoch LIMIT 1",
5618                )?,
5619                "idx_usage_daily_retention",
5620            );
5621            assert_plan_uses(
5622                &query_plan(
5623                    &database.connection,
5624                    "SELECT project_instance_id FROM usage_label_tombstones
5625                 WHERE expired_at_epoch < 1 ORDER BY expired_at_epoch LIMIT 1",
5626                )?,
5627                "idx_usage_label_tombstones_retention",
5628            );
5629            for redundant in [
5630                "idx_usage_baselines_active",
5631                "idx_usage_daily_range",
5632                "idx_usage_instance_daily_range",
5633            ] {
5634                let exists = database.connection.query_row(
5635                    "SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE type = 'index' AND name = ?1)",
5636                    [redundant],
5637                    |row| row.get::<_, i64>(0),
5638                )?;
5639                assert_eq!(exists, 0, "redundant index {redundant} must stay absent");
5640            }
5641            Ok(())
5642        })();
5643        assert!(
5644            result.is_ok(),
5645            "telemetry query-plan test failed: {result:?}"
5646        );
5647    }
5648}