//! SQLite setup: pragmas, migrations, and the two-pool split. //! //! There are two pools, and the split is the whole performance story: //! //! * a **read pool** (4–8 connections) for HTTP handlers, and //! * a **writer pool of exactly one connection**, owned by the writer task. //! //! SQLite allows one writer at a time. Rather than discovering that as //! `SQLITE_BUSY` under load, the design makes it structural: all writes funnel //! through one task that batches them. 100 devices reporting once a minute becomes //! ~4 transactions per second instead of 100 fsyncs, and `SQLITE_BUSY` cannot //! happen because there is never a second writer to contend with. use std::path::Path; use std::str::FromStr; use std::time::Duration; use anyhow::{Context, Result}; use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions, SqliteSynchronous}; use sqlx::{Executor, SqlitePool}; pub struct Db { /// For HTTP handlers. Concurrent readers are free under WAL. pub read: SqlitePool, /// Single connection, held by the writer task. Nothing else may write. pub write: SqlitePool, } /// Pragmas that are not optional. /// /// * `journal_mode=WAL` — readers never block the writer, and vice versa. /// * `synchronous=NORMAL` — with WAL this risks losing the last few /// *transactions* on an OS crash, not corruption. For location history that is /// the right trade against an fsync per commit. /// * `busy_timeout` — belt and braces; the single-writer design should make it /// unreachable. /// * `foreign_keys=ON` — off by default in SQLite, which surprises everyone once. /// * `auto_vacuum=INCREMENTAL` — lets the retention sweep return space in bounded /// chunks instead of a stop-the-world VACUUM. fn write_options(path: &Path) -> Result { Ok(base_options(path)? .create_if_missing(true) .journal_mode(SqliteJournalMode::Wal) .synchronous(SqliteSynchronous::Normal) .pragma("auto_vacuum", "incremental") // Keep the WAL from growing without bound between checkpoints. .pragma("journal_size_limit", "67108864")) // 64 MiB } /// Options for the read pool. /// /// Read-only is enforced with `PRAGMA query_only` rather than by opening the file /// `SQLITE_OPEN_READONLY`. The distinction matters: several of the pragmas above /// are themselves writes, and a genuinely read-only handle also cannot create the /// `-shm` file a WAL database needs, so it fails in ways that depend on whether a /// writer happens to be attached. `query_only` rejects writes at the statement /// level, which is the property actually wanted here. fn read_options(path: &Path) -> Result { Ok(base_options(path)? .create_if_missing(false) .pragma("query_only", "ON")) } fn base_options(path: &Path) -> Result { Ok( SqliteConnectOptions::from_str(&format!("sqlite://{}", path.display())) .with_context(|| format!("bad database path {}", path.display()))? .foreign_keys(true) .busy_timeout(Duration::from_secs(5)) .pragma("mmap_size", "268435456"), // 256 MiB ) } impl Db { pub async fn open(path: &Path) -> Result { if let Some(parent) = path.parent().filter(|p| !p.as_os_str().is_empty()) { std::fs::create_dir_all(parent) .with_context(|| format!("creating {}", parent.display()))?; } // Migrations run on the writer pool: they are writes, and running them // here means the read pool never sees a half-migrated schema. let write = SqlitePoolOptions::new() .max_connections(1) .min_connections(1) .connect_with(write_options(path)?) .await .with_context(|| format!("opening {} for writing", path.display()))?; sqlx::migrate!("./migrations") .run(&write) .await .context("running migrations")?; let read = SqlitePoolOptions::new() .max_connections(8) .min_connections(2) .connect_with(read_options(path)?) .await .with_context(|| format!("opening {} for reading", path.display()))?; Ok(Self { read, write }) } /// Flush the WAL back into the main database file. Called on shutdown so the /// on-disk file is self-contained. pub async fn checkpoint(&self) -> Result<()> { self.write .execute("PRAGMA wal_checkpoint(TRUNCATE);") .await .context("WAL checkpoint")?; Ok(()) } pub async fn close(&self) { self.read.close().await; self.write.close().await; } } /// Seconds since the Unix epoch. /// /// Every timestamp in this system is an `i64` of Unix seconds. No `TEXT` /// datetimes, no local time, nowhere. pub fn now() -> i64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map_or(0, |d| d.as_secs() as i64) } #[cfg(test)] mod tests { use super::*; /// A throwaway database in a temp directory, migrated and ready. pub async fn test_db() -> (Db, tempfile::TempDir) { let dir = tempfile::tempdir().expect("temp dir"); let db = Db::open(&dir.path().join("test.db")).await.expect("open"); (db, dir) } #[tokio::test] async fn migrations_apply_and_pragmas_take_effect() { let (db, _dir) = test_db().await; let journal: String = sqlx::query_scalar("PRAGMA journal_mode") .fetch_one(&db.write) .await .expect("journal_mode"); assert_eq!(journal.to_lowercase(), "wal"); let fk: i64 = sqlx::query_scalar("PRAGMA foreign_keys") .fetch_one(&db.write) .await .expect("foreign_keys"); assert_eq!(fk, 1, "foreign keys are off by default and must be enabled"); } #[tokio::test] async fn every_table_is_strict() { let (db, _dir) = test_db().await; let sql: Vec = sqlx::query_scalar( "SELECT sql FROM sqlite_master WHERE type = 'table' AND name NOT LIKE 'sqlite_%' \ AND name NOT LIKE '_sqlx%'", ) .fetch_all(&db.read) .await .expect("schema"); assert!(!sql.is_empty()); for stmt in sql { assert!( stmt.to_uppercase().contains("STRICT"), "table is not STRICT, so type affinity could store a string in an integer \ column:\n{stmt}" ); } } #[tokio::test] async fn the_read_pool_cannot_write() { let (db, _dir) = test_db().await; let err = sqlx::query("INSERT INTO settings (key, value) VALUES ('x', 'y')") .execute(&db.read) .await; assert!(err.is_err(), "the read pool must be read-only"); } #[tokio::test] async fn the_points_primary_key_makes_replay_idempotent() { let (db, _dir) = test_db().await; 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 insert = "INSERT INTO points (user_id, ts, lat, lon, acc_dm, recv_at) \ VALUES (1, 100, 5, 6, ?, 0) \ ON CONFLICT (user_id, ts) DO UPDATE SET \ acc_dm = excluded.acc_dm WHERE excluded.acc_dm < points.acc_dm"; sqlx::query(insert) .bind(80) .execute(&db.write) .await .expect("first"); // The same point again — a retry or a replay. sqlx::query(insert) .bind(80) .execute(&db.write) .await .expect("replay"); // A second phone, same second, worse accuracy: must not win. sqlx::query(insert) .bind(200) .execute(&db.write) .await .expect("worse"); // A second phone, same second, better accuracy: must win. sqlx::query(insert) .bind(30) .execute(&db.write) .await .expect("better"); let (count, acc): (i64, i64) = sqlx::query_as("SELECT COUNT(*), MIN(acc_dm) FROM points WHERE user_id = 1") .fetch_one(&db.read) .await .expect("count"); assert_eq!(count, 1, "a replayed datagram must not create a second row"); assert_eq!( acc, 30, "the better-accuracy point must win a same-second collision" ); } #[tokio::test] async fn a_share_must_target_one_other_account() { let (db, _dir) = test_db().await; for (id, username) in [(1, "a"), (2, "b")] { sqlx::query( "INSERT INTO users (id, username, pw_hash, display_name, created_at, pw_changed_at) \ VALUES (?, ?, 'x', 'X', 0, 0)", ) .bind(id) .bind(username) .execute(&db.write) .await .expect("user"); } // No viewer at all. assert!( sqlx::query("INSERT INTO shares (owner_user_id, created_at) VALUES (1, 0)") .execute(&db.write) .await .is_err() ); // Sharing with yourself. assert!( sqlx::query( "INSERT INTO shares (owner_user_id, viewer_user_id, created_at) VALUES (1, 1, 0)" ) .execute(&db.write) .await .is_err() ); // A real share. sqlx::query( "INSERT INTO shares (owner_user_id, viewer_user_id, created_at) VALUES (1, 2, 0)", ) .execute(&db.write) .await .expect("valid share"); } }