retention.rs
⎇
Raw
1//! Periodic housekeeping: point thinning, hard retention, token staleness.
2//!
3//! Runs every 15 minutes in bounded batches. Bounded matters: an unbounded
4//! `DELETE` on a large table holds the write lock for as long as it takes, which
5//! on the single-writer design means every device is stalled meanwhile.
6
7use std::sync::Arc;
8use std::time::Duration;
9
10use anyhow::{Context, Result};
11use sqlx::SqlitePool;
12use tracing::{debug, info, warn};
13
14use crate::db::now;
15use crate::ingest::Ingest;
16use otproto::RevokeReason;
17
18const INTERVAL: Duration = Duration::from_secs(15 * 60);
19
20/// Rows touched per statement, so the write lock is never held for long.
21const BATCH: i64 = 5_000;
22
23/// Points older than this are thinned to [`THIN_SPACING_S`] apart.
24const THIN_AFTER_S: i64 = 24 * 3_600;
25const THIN_SPACING_S: i64 = 5 * 60;
26
27pub struct Retention {
28 pub pool: SqlitePool,
29 pub retention_days: u32,
30 pub token_stale_days: u32,
31 pub ingest: Arc<Ingest>,
32}
33
34impl Retention {
35 pub fn spawn(self) -> tokio::task::JoinHandle<()> {
36 tokio::spawn(async move {
37 // Sleep first: startup already has enough to do.
38 let mut ticker = tokio::time::interval(INTERVAL);
39 ticker.tick().await;
40 loop {
41 ticker.tick().await;
42 if let Err(e) = self.run_once().await {
43 warn!(error = %e, "retention sweep failed; will retry next tick");
44 }
45 }
46 })
47 }
48
49 pub async fn run_once(&self) -> Result<()> {
50 let now = now();
51 let dropped = self.hard_drop(now).await?;
52 let thinned = self.thin(now).await?;
53 let tokens = self.expire_tokens(now).await?;
54 if dropped + thinned > 0 || !tokens.is_empty() {
55 info!(
56 dropped,
57 thinned,
58 tokens_expired = tokens.len(),
59 "retention sweep"
60 );
61 }
62 self.incremental_vacuum().await?;
63 Ok(())
64 }
65
66 /// Points beyond the retention horizon.
67 ///
68 /// `user_latest` is a separate table precisely so this cannot delete a live
69 /// marker: a user who has not moved in a fortnight still has a position on the
70 /// map, even with no surviving `points` row.
71 async fn hard_drop(&self, now: i64) -> Result<u64> {
72 let cutoff = now - i64::from(self.retention_days) * 86_400;
73 let mut total = 0;
74 loop {
75 let affected = sqlx::query(
76 // `points` is WITHOUT ROWID, so there is no rowid to select on;
77 // the row-value form addresses the primary key directly.
78 "DELETE FROM points WHERE (user_id, ts) IN \
79 (SELECT user_id, ts FROM points WHERE ts < ? LIMIT ?)",
80 )
81 .bind(cutoff)
82 .bind(BATCH)
83 .execute(&self.pool)
84 .await
85 .context("dropping expired points")?
86 .rows_affected();
87 total += affected;
88 if affected < BATCH as u64 {
89 break;
90 }
91 }
92 Ok(total)
93 }
94
95 /// Thin points older than a day down to roughly one per five minutes.
96 ///
97 /// Keeps the first point of each 5-minute bucket. A day-old trail does not
98 /// need per-second resolution, and this is where most of the space goes.
99 async fn thin(&self, now: i64) -> Result<u64> {
100 let cutoff = now - THIN_AFTER_S;
101 let mut total = 0;
102 loop {
103 let affected = sqlx::query(
104 "DELETE FROM points WHERE (user_id, ts) IN ( \
105 SELECT p.user_id, p.ts FROM points p WHERE p.ts < ? AND EXISTS ( \
106 SELECT 1 FROM points q \
107 WHERE q.user_id = p.user_id \
108 AND q.ts / ? = p.ts / ? \
109 AND q.ts < p.ts \
110 ) LIMIT ? )",
111 )
112 .bind(cutoff)
113 .bind(THIN_SPACING_S)
114 .bind(THIN_SPACING_S)
115 .bind(BATCH)
116 .execute(&self.pool)
117 .await
118 .context("thinning points")?
119 .rows_affected();
120 total += affected;
121 if affected < BATCH as u64 {
122 break;
123 }
124 }
125 Ok(total)
126 }
127
128 /// Delete tokens with no activity for `token_stale_days`.
129 ///
130 /// A phone genuinely idle for a month has to log in again — the same contract
131 /// as an expiring browser session. Removing the slot from the ingest cache is
132 /// the part that actually takes effect immediately; the database row is just
133 /// bookkeeping.
134 ///
135 /// `created_at` stands in for `last_seen_at` when a token was minted and never
136 /// used, so an abandoned login does not live forever.
137 async fn expire_tokens(&self, now: i64) -> Result<Vec<i64>> {
138 let cutoff = now - i64::from(self.token_stale_days) * 86_400;
139 let stale: Vec<i64> = sqlx::query_scalar(
140 "SELECT token_id FROM tokens \
141 WHERE COALESCE(last_seen_at, created_at) < ? AND revoked_at IS NULL LIMIT ?",
142 )
143 .bind(cutoff)
144 .bind(BATCH)
145 .fetch_all(&self.pool)
146 .await
147 .context("finding stale tokens")?;
148
149 for token_id in &stale {
150 // Marked, not deleted. A phone that wakes up after a month must be
151 // told to log in again, and the only message it can trust is one
152 // sealed with its own key — which deleting the row would destroy.
153 // A token row is about a hundred bytes; a device that cannot be told
154 // why it stopped working is a support ticket.
155 sqlx::query(
156 "UPDATE tokens SET revoked_at = ? WHERE token_id = ? AND revoked_at IS NULL",
157 )
158 .bind(now)
159 .bind(token_id)
160 .execute(&self.pool)
161 .await
162 .context("expiring stale token")?;
163 // History survives regardless: points.src_token_id is deliberately
164 // not a cascading foreign key, because positions belong to the
165 // account rather than to the phone that reported them.
166 self.ingest
167 .mark_revoked(*token_id as u64, RevokeReason::Expired);
168 }
169 Ok(stale)
170 }
171
172 /// Return freed pages to the filesystem a little at a time.
173 async fn incremental_vacuum(&self) -> Result<()> {
174 sqlx::query("PRAGMA incremental_vacuum(1000)")
175 .execute(&self.pool)
176 .await
177 .context("incremental vacuum")?;
178 debug!("incremental vacuum done");
179 Ok(())
180 }
181}
182
183/// Also on a timer: sweep the state that attacker-chosen keys accumulate in — the
184/// UDP rate limiter and the failed-login throttle. Cheap, and the reason neither
185/// can become a memory-exhaustion vector itself.
186pub fn spawn_limits_gc(
187 ingest: Arc<Ingest>,
188 throttle: Arc<crate::auth::LoginThrottle>,
189) -> tokio::task::JoinHandle<()> {
190 tokio::spawn(async move {
191 let mut ticker = tokio::time::interval(Duration::from_secs(60));
192 loop {
193 ticker.tick().await;
194 ingest.limits().gc();
195 throttle.gc();
196 }
197 })
198}
199
200#[cfg(test)]
201mod tests {
202 use super::*;
203 use crate::db::Db;
204
205 async fn fixture() -> (Db, Arc<Ingest>, tempfile::TempDir) {
206 let dir = tempfile::tempdir().expect("temp dir");
207 let db = Db::open(&dir.path().join("t.db")).await.expect("open");
208 sqlx::query(
209 "INSERT INTO users (id, username, pw_hash, display_name, created_at, pw_changed_at) \
210 VALUES (1, 'a', 'x', 'A', 0, 0)",
211 )
212 .execute(&db.write)
213 .await
214 .expect("user");
215 let (writer, _task) = crate::writer::spawn(db.write.clone());
216 let ingest = Arc::new(Ingest::new(writer, 30 * 86_400, None));
217 (db, ingest, dir)
218 }
219
220 async fn insert_point(db: &Db, ts: i64) {
221 sqlx::query("INSERT INTO points (user_id, ts, lat, lon, recv_at) VALUES (1, ?, 1, 2, ?)")
222 .bind(ts)
223 .bind(ts)
224 .execute(&db.write)
225 .await
226 .expect("point");
227 }
228
229 fn retention(db: &Db, ingest: Arc<Ingest>) -> Retention {
230 Retention {
231 pool: db.write.clone(),
232 retention_days: 7,
233 token_stale_days: 30,
234 ingest,
235 }
236 }
237
238 #[tokio::test]
239 async fn points_beyond_the_horizon_are_dropped_and_recent_ones_kept() {
240 let (db, ingest, _dir) = fixture().await;
241 let now = now();
242 insert_point(&db, now - 8 * 86_400).await; // too old
243 insert_point(&db, now - 60).await; // fresh
244
245 retention(&db, ingest).run_once().await.expect("sweep");
246
247 let remaining: Vec<i64> = sqlx::query_scalar("SELECT ts FROM points ORDER BY ts")
248 .fetch_all(&db.read)
249 .await
250 .expect("points");
251 assert_eq!(remaining, vec![now - 60]);
252 }
253
254 #[tokio::test]
255 async fn the_live_marker_survives_the_retention_horizon() {
256 let (db, ingest, _dir) = fixture().await;
257 let ancient = now() - 30 * 86_400;
258 insert_point(&db, ancient).await;
259 sqlx::query(
260 "INSERT INTO user_latest (user_id, ts, lat, lon, recv_at) VALUES (1, ?, 1, 2, ?)",
261 )
262 .bind(ancient)
263 .bind(ancient)
264 .execute(&db.write)
265 .await
266 .expect("latest");
267
268 retention(&db, ingest).run_once().await.expect("sweep");
269
270 let points: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points")
271 .fetch_one(&db.read)
272 .await
273 .expect("count");
274 let latest: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM user_latest")
275 .fetch_one(&db.read)
276 .await
277 .expect("count");
278 assert_eq!(points, 0, "the old point should be gone");
279 assert_eq!(
280 latest, 1,
281 "a separate table for the live marker is the whole point: GC must not erase someone \
282 from the map for standing still"
283 );
284 }
285
286 #[tokio::test]
287 async fn day_old_points_are_thinned_to_one_per_bucket() {
288 let (db, ingest, _dir) = fixture().await;
289 let base = now() - 2 * 86_400;
290 // Ten points inside a single 5-minute bucket.
291 for i in 0..10 {
292 insert_point(&db, base + i * 10).await;
293 }
294 // And one in the next bucket, which must be kept too.
295 insert_point(&db, base + THIN_SPACING_S).await;
296
297 retention(&db, ingest).run_once().await.expect("sweep");
298
299 let remaining: Vec<i64> = sqlx::query_scalar("SELECT ts FROM points ORDER BY ts")
300 .fetch_all(&db.read)
301 .await
302 .expect("points");
303 assert_eq!(
304 remaining.len(),
305 2,
306 "each 5-minute bucket should keep exactly one point, got {remaining:?}"
307 );
308 assert_eq!(
309 remaining[0], base,
310 "the first point of a bucket is the one kept"
311 );
312 }
313
314 #[tokio::test]
315 async fn recent_points_are_not_thinned() {
316 let (db, ingest, _dir) = fixture().await;
317 let base = now() - 600;
318 for i in 0..5 {
319 insert_point(&db, base + i * 10).await;
320 }
321 retention(&db, ingest).run_once().await.expect("sweep");
322 let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points")
323 .fetch_one(&db.read)
324 .await
325 .expect("count");
326 assert_eq!(count, 5, "points inside the last 24h keep full resolution");
327 }
328
329 #[tokio::test]
330 async fn a_token_idle_for_a_month_is_revoked_but_keeps_its_key() {
331 let (db, ingest, _dir) = fixture().await;
332 let stale_id = 1234i64;
333 let fresh_id = 5678i64;
334 let now = now();
335
336 for (id, last_seen) in [(stale_id, now - 31 * 86_400), (fresh_id, now - 60)] {
337 sqlx::query(
338 "INSERT INTO tokens (token_id, user_id, key_wrapped, name, created_at, last_seen_at) \
339 VALUES (?, 1, x'00', 'phone', 0, ?)",
340 )
341 .bind(id)
342 .bind(last_seen)
343 .execute(&db.write)
344 .await
345 .expect("token");
346 ingest.insert_token(crate::ingest::TokenSlot::new(id as u64, 1, &[1; 32], 1));
347 }
348 assert_eq!(ingest.active_token_count(), 2);
349
350 retention(&db, Arc::clone(&ingest))
351 .run_once()
352 .await
353 .expect("sweep");
354
355 // The row survives, marked. Deleting it would destroy the only key that
356 // can seal a message this phone will believe — and a phone idle for a
357 // month is exactly the one that needs telling.
358 let expired: Vec<i64> =
359 sqlx::query_scalar("SELECT token_id FROM tokens WHERE revoked_at IS NOT NULL")
360 .fetch_all(&db.read)
361 .await
362 .expect("tokens");
363 assert_eq!(expired, vec![stale_id]);
364 let live: Vec<i64> =
365 sqlx::query_scalar("SELECT token_id FROM tokens WHERE revoked_at IS NULL")
366 .fetch_all(&db.read)
367 .await
368 .expect("tokens");
369 assert_eq!(live, vec![fresh_id]);
370 assert_eq!(
371 ingest.active_token_count(),
372 1,
373 "an expired token must stop authorising writes immediately"
374 );
375 }
376
377 /// The sweep must not keep re-expiring what it already expired, or every run
378 /// would rewrite the same rows forever.
379 #[tokio::test]
380 async fn expiring_is_idempotent() {
381 let (db, ingest, _dir) = fixture().await;
382 sqlx::query(
383 "INSERT INTO tokens (token_id, user_id, key_wrapped, name, created_at) \
384 VALUES (7, 1, x'00', 'phone', ?)",
385 )
386 .bind(now() - 31 * 86_400)
387 .execute(&db.write)
388 .await
389 .expect("token");
390
391 let sweep = retention(&db, Arc::clone(&ingest));
392 sweep.run_once().await.expect("first sweep");
393 let first: i64 = sqlx::query_scalar("SELECT revoked_at FROM tokens WHERE token_id = 7")
394 .fetch_one(&db.read)
395 .await
396 .expect("revoked_at");
397
398 retention(&db, ingest)
399 .run_once()
400 .await
401 .expect("second sweep");
402 let second: i64 = sqlx::query_scalar("SELECT revoked_at FROM tokens WHERE token_id = 7")
403 .fetch_one(&db.read)
404 .await
405 .expect("revoked_at");
406 assert_eq!(
407 first, second,
408 "the sweep rewrote a token it had already expired"
409 );
410 }
411
412 #[tokio::test]
413 async fn a_token_minted_and_never_used_still_expires() {
414 let (db, ingest, _dir) = fixture().await;
415 sqlx::query(
416 "INSERT INTO tokens (token_id, user_id, key_wrapped, name, created_at) \
417 VALUES (99, 1, x'00', 'phone', ?)",
418 )
419 .bind(now() - 31 * 86_400)
420 .execute(&db.write)
421 .await
422 .expect("token");
423
424 retention(&db, ingest).run_once().await.expect("sweep");
425 let count: i64 =
426 sqlx::query_scalar("SELECT COUNT(*) FROM tokens WHERE revoked_at IS NOT NULL")
427 .fetch_one(&db.read)
428 .await
429 .expect("count");
430 assert_eq!(count, 1, "created_at must stand in for a never-used token");
431 }
432
433 #[tokio::test]
434 async fn expiring_a_token_does_not_delete_history() {
435 let (db, ingest, _dir) = fixture().await;
436 let now = now();
437 sqlx::query(
438 "INSERT INTO tokens (token_id, user_id, key_wrapped, name, created_at, last_seen_at) \
439 VALUES (42, 1, x'00', 'phone', 0, ?)",
440 )
441 .bind(now - 31 * 86_400)
442 .execute(&db.write)
443 .await
444 .expect("token");
445 sqlx::query(
446 "INSERT INTO points (user_id, ts, lat, lon, recv_at, src_token_id) \
447 VALUES (1, ?, 1, 2, ?, 42)",
448 )
449 .bind(now - 60)
450 .bind(now)
451 .execute(&db.write)
452 .await
453 .expect("point");
454
455 retention(&db, ingest).run_once().await.expect("sweep");
456
457 let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points")
458 .fetch_one(&db.read)
459 .await
460 .expect("count");
461 assert_eq!(
462 count, 1,
463 "positions belong to the account, so losing the token must not lose the trail"
464 );
465 }
466}
467