1use 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
14const POLICY_VERSION: u32 = 1;
16const LOGICAL_BYTE_VERSION: u32 = 1;
18const OVERFLOW_DIMENSION: &str = "<overflow>";
20const INSTANCE_ACTIVE: &str = "active";
22const INSTANCE_SEALED: &str = "sealed";
24const INSTANCE_EXPIRED: &str = "expired";
26const DEDUPE_SCOPE_EVENT: &str = "event";
28const SECONDS_PER_DAY: i64 = 86_400;
30const AGGREGATE_COUNTER_FIELD: &str = "aggregate_counter";
32const LEGACY_TEXT_HASH_DOMAIN: &[u8] = b"projectatlas:legacy-telemetry-text:v1\0";
34
35#[derive(Clone, Copy, Debug, Eq, PartialEq)]
37enum BaselineAdmission {
38 BoundedRuntime,
40 SupportedUpgrade,
42}
43
44#[derive(Clone, Copy, Debug, Eq, PartialEq)]
46enum DimensionAdmission {
47 Event,
49 Overflow,
51}
52
53#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
55struct LegacyDetailLoss {
56 raw: bool,
58 dimension: bool,
60 label: bool,
62}
63
64macro_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
86macro_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
109pub(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#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
118pub struct TelemetryRetentionPolicy {
119 pub max_raw_rows: usize,
121 pub max_raw_logical_bytes: usize,
123 pub max_raw_age_seconds: u64,
125 pub max_dimensions: usize,
127 pub max_active_instances: usize,
129 pub max_retained_instances: usize,
131 pub max_label_tombstones: usize,
133 pub max_instance_tombstones: usize,
135 pub max_retained_labels: usize,
137 pub max_baselines_per_instance: usize,
139 pub max_active_baseline_rows: usize,
141 pub max_baseline_logical_bytes: usize,
143 pub max_daily_rows: usize,
145 pub prune_batch_rows: usize,
147 pub max_label_bytes: usize,
149 pub max_command_bytes: usize,
151 pub max_path_bytes: usize,
153 pub max_query_bytes: usize,
155 pub max_dimension_bytes: usize,
157 pub max_baseline_witness_bytes: usize,
159 pub max_active_idle_seconds: u64,
161 pub future_clock_tolerance_seconds: u64,
163 pub checkpoint_write_interval: usize,
165 pub retained_trend_days: u64,
167 pub retained_instance_seconds: u64,
169 pub retained_label_seconds: u64,
171 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 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#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
271#[serde(rename_all = "snake_case")]
272pub enum SpillCleanupState {
273 NotApplicable,
275}
276
277#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
279#[serde(rename_all = "snake_case")]
280pub enum TelemetryCheckpointState {
281 NotDue,
283 Completed,
285 Busy,
287 Error,
289}
290
291impl TelemetryCheckpointState {
292 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 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#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
319#[serde(rename_all = "snake_case")]
320pub enum PlannerStatisticsPolicy {
321 NotConfigured,
323}
324
325#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
327#[serde(rename_all = "snake_case")]
328pub enum PlannerStatisticsState {
329 NotInitialized,
331 Available,
333}
334
335#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
337pub struct TelemetryRetentionState {
338 pub policy_version: u32,
340 pub logical_byte_version: u32,
342 pub raw_rows: usize,
344 pub max_raw_rows: usize,
346 pub max_raw_age_seconds: u64,
348 pub raw_logical_bytes: usize,
350 pub max_raw_logical_bytes: usize,
352 pub baseline_rows: usize,
354 pub max_baselines_per_instance: usize,
356 pub max_active_baseline_rows: usize,
358 pub baseline_logical_bytes: usize,
360 pub max_baseline_logical_bytes: usize,
362 pub dimension_rows: usize,
364 pub max_dimensions: usize,
366 pub instance_rows: usize,
368 pub active_instance_rows: usize,
370 pub max_active_instances: usize,
372 pub max_retained_instances: usize,
374 pub retained_label_rows: usize,
376 pub max_retained_labels: usize,
378 pub daily_rows: usize,
380 pub max_daily_rows: usize,
382 pub retained_trend_days: u64,
384 pub label_tombstone_rows: usize,
386 pub max_label_tombstones: usize,
388 pub instance_tombstone_rows: usize,
390 pub max_instance_tombstones: usize,
392 pub pruned_raw_rows: usize,
394 pub pruned_instance_rows: usize,
396 pub evicted_tombstones: usize,
398 pub maintenance_pending: bool,
400 pub prune_batch_rows: usize,
402 pub writes_since_checkpoint: usize,
404 pub checkpoint_write_interval: usize,
406 pub last_checkpoint_epoch: u64,
408 pub oldest_retained_epoch: Option<u64>,
410 pub clock_anomaly: bool,
412 pub spill_cleanup: SpillCleanupState,
414 pub checkpoint_state: TelemetryCheckpointState,
416 pub wal_autocheckpoint_pages: usize,
418 pub freelist_pages: usize,
420 pub page_count: usize,
422 pub page_size: usize,
424 pub journal_mode: String,
426 pub synchronous_mode: String,
428 pub connection_busy_timeout_ms: u64,
430 pub normal_busy_timeout_ms: u64,
432 pub telemetry_busy_timeout_ms: u64,
434 pub statistics_policy: PlannerStatisticsPolicy,
436 pub statistics_state: PlannerStatisticsState,
438}
439
440#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
442struct DimensionValues {
443 token_savings_bucket: String,
445 provider: String,
447 model: String,
449 tokenizer_backend: String,
451 accuracy: String,
453 baseline_kind: String,
455 confidence: String,
457 accounting_layer: String,
459 estimate_method: String,
461 denominator_kind: String,
463 dedupe_scope: String,
465 overflow: bool,
467}
468
469impl DimensionValues {
470 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 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#[derive(Clone, Copy, Debug, Default)]
509struct AggregateCounters {
510 calls: i64,
512 estimated_without: i64,
514 estimated_with: i64,
516 observed_without: i64,
518 observed_with: i64,
520 modeled_without: i64,
522 modeled_with: i64,
524 deduped_modeled_without: i64,
526 deduped_modeled_with: i64,
528 repeated_baselines: i64,
530 observed_file_read_replacements: i64,
532 modeled_file_reads_avoided: i64,
534}
535
536#[derive(Clone, Copy, Debug)]
538enum RetentionCounter {
539 RawRows,
541 RawLogicalBytes,
543 BaselineRows,
545 BaselineLogicalBytes,
547 DimensionRows,
549 InstanceRows,
551 LabelRows,
553 DailyRows,
555 LabelTombstoneRows,
557 InstanceTombstoneRows,
559}
560
561impl RetentionCounter {
562 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 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 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 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
672pub(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
679pub(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
766pub(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)]
794fn 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
886pub(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
895fn 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
931pub(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
982pub(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
1045pub(crate) fn retention_state(connection: &Connection) -> DbResult<TelemetryRetentionState> {
1047 let project = current_project(connection)?;
1048 retention_state_for_project(connection, project)
1049}
1050
1051fn 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
1199pub(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
1208fn 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
1233pub(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
1242fn 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
1257pub(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
1267fn 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
1300fn current_project(connection: &Connection) -> DbResult<ProjectInstanceId> {
1302 crate::project_identity::load_project_identity(connection)?
1303 .ok_or(DbError::ProjectInstanceIdentityMissing)
1304}
1305
1306fn 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
1321fn 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
1378fn 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
1391fn 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
1407fn 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
1420fn 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
1472fn event_label(event: &UsageEvent) -> Option<&str> {
1474 (!event.session_id.is_empty()).then_some(event.session_id.as_str())
1475}
1476
1477pub(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
1527fn 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
1539fn 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
1553fn ensure_overflow_dimension(connection: &Connection) -> DbResult<i64> {
1555 ensure_dimension_unbounded(connection, &DimensionValues::overflow())
1556}
1557
1558fn 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
1577fn 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
1609fn 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
1640fn 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
1718fn 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
1770fn 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
1801fn 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
1866fn 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
1927fn 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
2062fn 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
2083fn 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
2127fn 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
2156fn 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
2185fn 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
2214fn 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
2244fn 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
2274fn 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#[derive(Debug, Eq, PartialEq)]
2291enum DailyEvictionCandidate {
2292 Global {
2294 project: Vec<u8>,
2296 day: i64,
2298 dimension_id: i64,
2300 },
2301 Instance {
2303 instance_row_id: i64,
2305 day: i64,
2307 dimension_id: i64,
2309 },
2310}
2311
2312impl DailyEvictionCandidate {
2313 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
2335fn 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
2405fn 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
2462fn 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
2495fn 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
2534fn 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
2577fn 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
2711fn 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
2734fn 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
2830fn 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
2878fn 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
2939fn 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
3058fn 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
3126fn 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)]
3199fn 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
3298fn 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
3369fn 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
3387fn 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
3414fn 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
3485fn 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
3556fn 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
3579fn 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
3651fn 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
3718fn 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
3736fn 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
3754fn 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
3788fn 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
3804fn component_difference(without: i64, with: i64) -> i128 {
3806 i128::from(without) - i128::from(with)
3807}
3808
3809fn 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
3816fn 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
3825fn 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
3831fn 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
3846fn 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
3862fn 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
3874fn option_usize_to_i64(field: &'static str, value: Option<usize>) -> DbResult<Option<i64>> {
3876 value.map(|value| to_i64(field, value)).transpose()
3877}
3878
3879fn 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
3888fn 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
3895fn count_usize(field: &'static str, value: i64) -> DbResult<usize> {
3897 usize::try_from(value).map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3898}
3899
3900fn count_u32(field: &'static str, value: i64) -> DbResult<u32> {
3902 u32::try_from(value).map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3903}
3904
3905fn count_u64(field: &'static str, value: i64) -> DbResult<u64> {
3907 u64::try_from(value).map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3908}
3909
3910fn count_u128(field: &'static str, value: i64) -> DbResult<u128> {
3912 u128::try_from(value).map_err(|_source| DbError::TelemetryIntegerOverflow { field })
3913}
3914
3915fn 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
3921fn 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
3931fn 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
3949fn 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 ¤t,
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}