//! Periodic housekeeping: point thinning, hard retention, token staleness. //! //! Runs every 15 minutes in bounded batches. Bounded matters: an unbounded //! `DELETE` on a large table holds the write lock for as long as it takes, which //! on the single-writer design means every device is stalled meanwhile. use std::sync::Arc; use std::time::Duration; use anyhow::{Context, Result}; use sqlx::SqlitePool; use tracing::{debug, info, warn}; use crate::db::now; use crate::ingest::Ingest; use otproto::RevokeReason; const INTERVAL: Duration = Duration::from_secs(15 * 60); /// Rows touched per statement, so the write lock is never held for long. const BATCH: i64 = 5_000; /// Points older than this are thinned to [`THIN_SPACING_S`] apart. const THIN_AFTER_S: i64 = 24 * 3_600; const THIN_SPACING_S: i64 = 5 * 60; pub struct Retention { pub pool: SqlitePool, pub retention_days: u32, pub token_stale_days: u32, pub ingest: Arc, } impl Retention { pub fn spawn(self) -> tokio::task::JoinHandle<()> { tokio::spawn(async move { // Sleep first: startup already has enough to do. let mut ticker = tokio::time::interval(INTERVAL); ticker.tick().await; loop { ticker.tick().await; if let Err(e) = self.run_once().await { warn!(error = %e, "retention sweep failed; will retry next tick"); } } }) } pub async fn run_once(&self) -> Result<()> { let now = now(); let dropped = self.hard_drop(now).await?; let thinned = self.thin(now).await?; let tokens = self.expire_tokens(now).await?; if dropped + thinned > 0 || !tokens.is_empty() { info!( dropped, thinned, tokens_expired = tokens.len(), "retention sweep" ); } self.incremental_vacuum().await?; Ok(()) } /// Points beyond the retention horizon. /// /// `user_latest` is a separate table precisely so this cannot delete a live /// marker: a user who has not moved in a fortnight still has a position on the /// map, even with no surviving `points` row. async fn hard_drop(&self, now: i64) -> Result { let cutoff = now - i64::from(self.retention_days) * 86_400; let mut total = 0; loop { let affected = sqlx::query( // `points` is WITHOUT ROWID, so there is no rowid to select on; // the row-value form addresses the primary key directly. "DELETE FROM points WHERE (user_id, ts) IN \ (SELECT user_id, ts FROM points WHERE ts < ? LIMIT ?)", ) .bind(cutoff) .bind(BATCH) .execute(&self.pool) .await .context("dropping expired points")? .rows_affected(); total += affected; if affected < BATCH as u64 { break; } } Ok(total) } /// Thin points older than a day down to roughly one per five minutes. /// /// Keeps the first point of each 5-minute bucket. A day-old trail does not /// need per-second resolution, and this is where most of the space goes. async fn thin(&self, now: i64) -> Result { let cutoff = now - THIN_AFTER_S; let mut total = 0; loop { let affected = sqlx::query( "DELETE FROM points WHERE (user_id, ts) IN ( \ SELECT p.user_id, p.ts FROM points p WHERE p.ts < ? AND EXISTS ( \ SELECT 1 FROM points q \ WHERE q.user_id = p.user_id \ AND q.ts / ? = p.ts / ? \ AND q.ts < p.ts \ ) LIMIT ? )", ) .bind(cutoff) .bind(THIN_SPACING_S) .bind(THIN_SPACING_S) .bind(BATCH) .execute(&self.pool) .await .context("thinning points")? .rows_affected(); total += affected; if affected < BATCH as u64 { break; } } Ok(total) } /// Delete tokens with no activity for `token_stale_days`. /// /// A phone genuinely idle for a month has to log in again — the same contract /// as an expiring browser session. Removing the slot from the ingest cache is /// the part that actually takes effect immediately; the database row is just /// bookkeeping. /// /// `created_at` stands in for `last_seen_at` when a token was minted and never /// used, so an abandoned login does not live forever. async fn expire_tokens(&self, now: i64) -> Result> { let cutoff = now - i64::from(self.token_stale_days) * 86_400; let stale: Vec = sqlx::query_scalar( "SELECT token_id FROM tokens \ WHERE COALESCE(last_seen_at, created_at) < ? AND revoked_at IS NULL LIMIT ?", ) .bind(cutoff) .bind(BATCH) .fetch_all(&self.pool) .await .context("finding stale tokens")?; for token_id in &stale { // Marked, not deleted. A phone that wakes up after a month must be // told to log in again, and the only message it can trust is one // sealed with its own key — which deleting the row would destroy. // A token row is about a hundred bytes; a device that cannot be told // why it stopped working is a support ticket. sqlx::query( "UPDATE tokens SET revoked_at = ? WHERE token_id = ? AND revoked_at IS NULL", ) .bind(now) .bind(token_id) .execute(&self.pool) .await .context("expiring stale token")?; // History survives regardless: points.src_token_id is deliberately // not a cascading foreign key, because positions belong to the // account rather than to the phone that reported them. self.ingest .mark_revoked(*token_id as u64, RevokeReason::Expired); } Ok(stale) } /// Return freed pages to the filesystem a little at a time. async fn incremental_vacuum(&self) -> Result<()> { sqlx::query("PRAGMA incremental_vacuum(1000)") .execute(&self.pool) .await .context("incremental vacuum")?; debug!("incremental vacuum done"); Ok(()) } } /// Also on a timer: sweep the state that attacker-chosen keys accumulate in — the /// UDP rate limiter and the failed-login throttle. Cheap, and the reason neither /// can become a memory-exhaustion vector itself. pub fn spawn_limits_gc( ingest: Arc, throttle: Arc, ) -> tokio::task::JoinHandle<()> { tokio::spawn(async move { let mut ticker = tokio::time::interval(Duration::from_secs(60)); loop { ticker.tick().await; ingest.limits().gc(); throttle.gc(); } }) } #[cfg(test)] mod tests { use super::*; use crate::db::Db; async fn fixture() -> (Db, Arc, tempfile::TempDir) { let dir = tempfile::tempdir().expect("temp dir"); let db = Db::open(&dir.path().join("t.db")).await.expect("open"); sqlx::query( "INSERT INTO users (id, username, pw_hash, display_name, created_at, pw_changed_at) \ VALUES (1, 'a', 'x', 'A', 0, 0)", ) .execute(&db.write) .await .expect("user"); let (writer, _task) = crate::writer::spawn(db.write.clone()); let ingest = Arc::new(Ingest::new(writer, 30 * 86_400, None)); (db, ingest, dir) } async fn insert_point(db: &Db, ts: i64) { sqlx::query("INSERT INTO points (user_id, ts, lat, lon, recv_at) VALUES (1, ?, 1, 2, ?)") .bind(ts) .bind(ts) .execute(&db.write) .await .expect("point"); } fn retention(db: &Db, ingest: Arc) -> Retention { Retention { pool: db.write.clone(), retention_days: 7, token_stale_days: 30, ingest, } } #[tokio::test] async fn points_beyond_the_horizon_are_dropped_and_recent_ones_kept() { let (db, ingest, _dir) = fixture().await; let now = now(); insert_point(&db, now - 8 * 86_400).await; // too old insert_point(&db, now - 60).await; // fresh retention(&db, ingest).run_once().await.expect("sweep"); let remaining: Vec = sqlx::query_scalar("SELECT ts FROM points ORDER BY ts") .fetch_all(&db.read) .await .expect("points"); assert_eq!(remaining, vec![now - 60]); } #[tokio::test] async fn the_live_marker_survives_the_retention_horizon() { let (db, ingest, _dir) = fixture().await; let ancient = now() - 30 * 86_400; insert_point(&db, ancient).await; sqlx::query( "INSERT INTO user_latest (user_id, ts, lat, lon, recv_at) VALUES (1, ?, 1, 2, ?)", ) .bind(ancient) .bind(ancient) .execute(&db.write) .await .expect("latest"); retention(&db, ingest).run_once().await.expect("sweep"); let points: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points") .fetch_one(&db.read) .await .expect("count"); let latest: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM user_latest") .fetch_one(&db.read) .await .expect("count"); assert_eq!(points, 0, "the old point should be gone"); assert_eq!( latest, 1, "a separate table for the live marker is the whole point: GC must not erase someone \ from the map for standing still" ); } #[tokio::test] async fn day_old_points_are_thinned_to_one_per_bucket() { let (db, ingest, _dir) = fixture().await; let base = now() - 2 * 86_400; // Ten points inside a single 5-minute bucket. for i in 0..10 { insert_point(&db, base + i * 10).await; } // And one in the next bucket, which must be kept too. insert_point(&db, base + THIN_SPACING_S).await; retention(&db, ingest).run_once().await.expect("sweep"); let remaining: Vec = sqlx::query_scalar("SELECT ts FROM points ORDER BY ts") .fetch_all(&db.read) .await .expect("points"); assert_eq!( remaining.len(), 2, "each 5-minute bucket should keep exactly one point, got {remaining:?}" ); assert_eq!( remaining[0], base, "the first point of a bucket is the one kept" ); } #[tokio::test] async fn recent_points_are_not_thinned() { let (db, ingest, _dir) = fixture().await; let base = now() - 600; for i in 0..5 { insert_point(&db, base + i * 10).await; } retention(&db, ingest).run_once().await.expect("sweep"); let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points") .fetch_one(&db.read) .await .expect("count"); assert_eq!(count, 5, "points inside the last 24h keep full resolution"); } #[tokio::test] async fn a_token_idle_for_a_month_is_revoked_but_keeps_its_key() { let (db, ingest, _dir) = fixture().await; let stale_id = 1234i64; let fresh_id = 5678i64; let now = now(); for (id, last_seen) in [(stale_id, now - 31 * 86_400), (fresh_id, now - 60)] { sqlx::query( "INSERT INTO tokens (token_id, user_id, key_wrapped, name, created_at, last_seen_at) \ VALUES (?, 1, x'00', 'phone', 0, ?)", ) .bind(id) .bind(last_seen) .execute(&db.write) .await .expect("token"); ingest.insert_token(crate::ingest::TokenSlot::new(id as u64, 1, &[1; 32], 1)); } assert_eq!(ingest.active_token_count(), 2); retention(&db, Arc::clone(&ingest)) .run_once() .await .expect("sweep"); // The row survives, marked. Deleting it would destroy the only key that // can seal a message this phone will believe — and a phone idle for a // month is exactly the one that needs telling. let expired: Vec = sqlx::query_scalar("SELECT token_id FROM tokens WHERE revoked_at IS NOT NULL") .fetch_all(&db.read) .await .expect("tokens"); assert_eq!(expired, vec![stale_id]); let live: Vec = sqlx::query_scalar("SELECT token_id FROM tokens WHERE revoked_at IS NULL") .fetch_all(&db.read) .await .expect("tokens"); assert_eq!(live, vec![fresh_id]); assert_eq!( ingest.active_token_count(), 1, "an expired token must stop authorising writes immediately" ); } /// The sweep must not keep re-expiring what it already expired, or every run /// would rewrite the same rows forever. #[tokio::test] async fn expiring_is_idempotent() { let (db, ingest, _dir) = fixture().await; sqlx::query( "INSERT INTO tokens (token_id, user_id, key_wrapped, name, created_at) \ VALUES (7, 1, x'00', 'phone', ?)", ) .bind(now() - 31 * 86_400) .execute(&db.write) .await .expect("token"); let sweep = retention(&db, Arc::clone(&ingest)); sweep.run_once().await.expect("first sweep"); let first: i64 = sqlx::query_scalar("SELECT revoked_at FROM tokens WHERE token_id = 7") .fetch_one(&db.read) .await .expect("revoked_at"); retention(&db, ingest) .run_once() .await .expect("second sweep"); let second: i64 = sqlx::query_scalar("SELECT revoked_at FROM tokens WHERE token_id = 7") .fetch_one(&db.read) .await .expect("revoked_at"); assert_eq!( first, second, "the sweep rewrote a token it had already expired" ); } #[tokio::test] async fn a_token_minted_and_never_used_still_expires() { let (db, ingest, _dir) = fixture().await; sqlx::query( "INSERT INTO tokens (token_id, user_id, key_wrapped, name, created_at) \ VALUES (99, 1, x'00', 'phone', ?)", ) .bind(now() - 31 * 86_400) .execute(&db.write) .await .expect("token"); retention(&db, ingest).run_once().await.expect("sweep"); let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM tokens WHERE revoked_at IS NOT NULL") .fetch_one(&db.read) .await .expect("count"); assert_eq!(count, 1, "created_at must stand in for a never-used token"); } #[tokio::test] async fn expiring_a_token_does_not_delete_history() { let (db, ingest, _dir) = fixture().await; let now = now(); sqlx::query( "INSERT INTO tokens (token_id, user_id, key_wrapped, name, created_at, last_seen_at) \ VALUES (42, 1, x'00', 'phone', 0, ?)", ) .bind(now - 31 * 86_400) .execute(&db.write) .await .expect("token"); sqlx::query( "INSERT INTO points (user_id, ts, lat, lon, recv_at, src_token_id) \ VALUES (1, ?, 1, 2, ?, 42)", ) .bind(now - 60) .bind(now) .execute(&db.write) .await .expect("point"); retention(&db, ingest).run_once().await.expect("sweep"); let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points") .fetch_one(&db.read) .await .expect("count"); assert_eq!( count, 1, "positions belong to the account, so losing the token must not lose the trail" ); } }