diff --git a/src/daemon/telemetry_worker.rs b/src/daemon/telemetry_worker.rs index fa4bfe3850..45adbc7628 100644 --- a/src/daemon/telemetry_worker.rs +++ b/src/daemon/telemetry_worker.rs @@ -471,7 +471,7 @@ fn backfill_metrics_event_metadata() -> Result<(), GitAiError> { let mut db_lock = db .lock() .map_err(|_| GitAiError::Generic("metrics DB lock poisoned".to_string()))?; - db_lock.backfill_event_metadata_batch_after(after_id, METADATA_BACKFILL_BATCH_SIZE)? + db_lock.backfill_event_metadata_batch_once(after_id, METADATA_BACKFILL_BATCH_SIZE)? }; let Some(id) = last_id else { diff --git a/src/metrics/db.rs b/src/metrics/db.rs index 9fcac3884c..96ac603c89 100644 --- a/src/metrics/db.rs +++ b/src/metrics/db.rs @@ -24,6 +24,7 @@ const SCHEMA_VERSION: usize = 5; const MAX_METRIC_UPLOAD_ATTEMPTS: u32 = 6; const METRIC_PROCESSING_LOCK_TIMEOUT_SECS: u64 = 10 * 60; pub(crate) const METADATA_BACKFILL_BATCH_SIZE: usize = 1000; +const EVENT_METADATA_BACKFILL_COMPLETED_KEY: &str = "event_metadata_backfill_completed"; const NS_PER_SECOND: u128 = 1_000_000_000; const RETRYABLE_METRIC_IDS_SQL: &str = "SELECT id FROM metrics \ @@ -1197,6 +1198,39 @@ impl MetricsDatabase { .map(|(summary, _)| summary) } + pub(crate) fn event_metadata_backfill_completed(&self) -> Result { + let completed: Option = self + .conn + .query_row( + "SELECT value FROM schema_metadata WHERE key = ?1", + params![EVENT_METADATA_BACKFILL_COMPLETED_KEY], + |row| row.get(0), + ) + .optional()?; + Ok(completed.as_deref() == Some("1")) + } + + /// Backfill one bounded batch, permanently marking the one-time migration complete + /// after a successful scan reaches the end of the table. + pub(crate) fn backfill_event_metadata_batch_once( + &mut self, + after_id: i64, + limit: usize, + ) -> Result<(MetricMetadataBackfillSummary, Option), GitAiError> { + if limit == 0 || self.event_metadata_backfill_completed()? { + return Ok((MetricMetadataBackfillSummary::default(), None)); + } + + let result = self.backfill_event_metadata_batch_after(after_id, limit)?; + if result.0.scanned < limit { + self.conn.execute( + "INSERT OR REPLACE INTO schema_metadata (key, value) VALUES (?1, '1')", + params![EVENT_METADATA_BACKFILL_COMPLETED_KEY], + )?; + } + Ok(result) + } + /// Backfill cached event metadata for all currently eligible legacy rows. pub fn backfill_event_metadata(&mut self) -> Result { let mut total = MetricMetadataBackfillSummary::default(); @@ -2468,6 +2502,55 @@ mod tests { assert_eq!(empty_last_id, None); } + #[test] + fn test_backfill_event_metadata_batch_once_marks_completion() { + let (mut db, temp_dir) = create_test_db(); + db.conn + .execute( + "INSERT INTO metrics (event_json) VALUES (?1), (?2)", + params![event_json(days_ago(2)), event_json(days_ago(1))], + ) + .unwrap(); + + let (first_summary, last_id) = db.backfill_event_metadata_batch_once(0, 1).unwrap(); + assert_eq!(first_summary.scanned, 1); + assert!(!db.event_metadata_backfill_completed().unwrap()); + + let (second_summary, _) = db + .backfill_event_metadata_batch_once(last_id.unwrap(), 2) + .unwrap(); + assert_eq!(second_summary.scanned, 1); + assert!(db.event_metadata_backfill_completed().unwrap()); + + drop(db); + let conn = crate::sqlite::open_with_memory_limits(temp_dir.path().join("test-metrics.db")) + .unwrap(); + let mut db = MetricsDatabase { conn }; + db.initialize_schema().unwrap(); + assert!(db.event_metadata_backfill_completed().unwrap()); + + db.conn + .execute( + "INSERT INTO metrics (event_json) VALUES (?1)", + params![event_json(days_ago(0))], + ) + .unwrap(); + let inserted_after_completion = db.conn.last_insert_rowid(); + + let skipped = db.backfill_event_metadata_batch_once(0, 100).unwrap(); + assert_eq!(skipped, (MetricMetadataBackfillSummary::default(), None)); + assert_eq!( + db.conn + .query_row( + "SELECT event_ts FROM metrics WHERE id = ?1", + params![inserted_after_completion], + |row| row.get::<_, Option>(0), + ) + .unwrap(), + None + ); + } + #[test] fn test_dequeue_pending_batch_locks_rows() { let (mut db, _temp_dir) = create_test_db();