db.rs
⎇
Raw
1//! SQLite setup: pragmas, migrations, and the two-pool split.
2//!
3//! There are two pools, and the split is the whole performance story:
4//!
5//! * a **read pool** (4–8 connections) for HTTP handlers, and
6//! * a **writer pool of exactly one connection**, owned by the writer task.
7//!
8//! SQLite allows one writer at a time. Rather than discovering that as
9//! `SQLITE_BUSY` under load, the design makes it structural: all writes funnel
10//! through one task that batches them. 100 devices reporting once a minute becomes
11//! ~4 transactions per second instead of 100 fsyncs, and `SQLITE_BUSY` cannot
12//! happen because there is never a second writer to contend with.
13
14use std::path::Path;
15use std::str::FromStr;
16use std::time::Duration;
17
18use anyhow::{Context, Result};
19use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions, SqliteSynchronous};
20use sqlx::{Executor, SqlitePool};
21
22pub struct Db {
23 /// For HTTP handlers. Concurrent readers are free under WAL.
24 pub read: SqlitePool,
25 /// Single connection, held by the writer task. Nothing else may write.
26 pub write: SqlitePool,
27}
28
29/// Pragmas that are not optional.
30///
31/// * `journal_mode=WAL` — readers never block the writer, and vice versa.
32/// * `synchronous=NORMAL` — with WAL this risks losing the last few
33/// *transactions* on an OS crash, not corruption. For location history that is
34/// the right trade against an fsync per commit.
35/// * `busy_timeout` — belt and braces; the single-writer design should make it
36/// unreachable.
37/// * `foreign_keys=ON` — off by default in SQLite, which surprises everyone once.
38/// * `auto_vacuum=INCREMENTAL` — lets the retention sweep return space in bounded
39/// chunks instead of a stop-the-world VACUUM.
40fn write_options(path: &Path) -> Result<SqliteConnectOptions> {
41 Ok(base_options(path)?
42 .create_if_missing(true)
43 .journal_mode(SqliteJournalMode::Wal)
44 .synchronous(SqliteSynchronous::Normal)
45 .pragma("auto_vacuum", "incremental")
46 // Keep the WAL from growing without bound between checkpoints.
47 .pragma("journal_size_limit", "67108864")) // 64 MiB
48}
49
50/// Options for the read pool.
51///
52/// Read-only is enforced with `PRAGMA query_only` rather than by opening the file
53/// `SQLITE_OPEN_READONLY`. The distinction matters: several of the pragmas above
54/// are themselves writes, and a genuinely read-only handle also cannot create the
55/// `-shm` file a WAL database needs, so it fails in ways that depend on whether a
56/// writer happens to be attached. `query_only` rejects writes at the statement
57/// level, which is the property actually wanted here.
58fn read_options(path: &Path) -> Result<SqliteConnectOptions> {
59 Ok(base_options(path)?
60 .create_if_missing(false)
61 .pragma("query_only", "ON"))
62}
63
64fn base_options(path: &Path) -> Result<SqliteConnectOptions> {
65 Ok(
66 SqliteConnectOptions::from_str(&format!("sqlite://{}", path.display()))
67 .with_context(|| format!("bad database path {}", path.display()))?
68 .foreign_keys(true)
69 .busy_timeout(Duration::from_secs(5))
70 .pragma("mmap_size", "268435456"), // 256 MiB
71 )
72}
73
74impl Db {
75 pub async fn open(path: &Path) -> Result<Self> {
76 if let Some(parent) = path.parent().filter(|p| !p.as_os_str().is_empty()) {
77 std::fs::create_dir_all(parent)
78 .with_context(|| format!("creating {}", parent.display()))?;
79 }
80 // Migrations run on the writer pool: they are writes, and running them
81 // here means the read pool never sees a half-migrated schema.
82 let write = SqlitePoolOptions::new()
83 .max_connections(1)
84 .min_connections(1)
85 .connect_with(write_options(path)?)
86 .await
87 .with_context(|| format!("opening {} for writing", path.display()))?;
88
89 sqlx::migrate!("./migrations")
90 .run(&write)
91 .await
92 .context("running migrations")?;
93
94 let read = SqlitePoolOptions::new()
95 .max_connections(8)
96 .min_connections(2)
97 .connect_with(read_options(path)?)
98 .await
99 .with_context(|| format!("opening {} for reading", path.display()))?;
100
101 Ok(Self { read, write })
102 }
103
104 /// Flush the WAL back into the main database file. Called on shutdown so the
105 /// on-disk file is self-contained.
106 pub async fn checkpoint(&self) -> Result<()> {
107 self.write
108 .execute("PRAGMA wal_checkpoint(TRUNCATE);")
109 .await
110 .context("WAL checkpoint")?;
111 Ok(())
112 }
113
114 pub async fn close(&self) {
115 self.read.close().await;
116 self.write.close().await;
117 }
118}
119
120/// Seconds since the Unix epoch.
121///
122/// Every timestamp in this system is an `i64` of Unix seconds. No `TEXT`
123/// datetimes, no local time, nowhere.
124pub fn now() -> i64 {
125 std::time::SystemTime::now()
126 .duration_since(std::time::UNIX_EPOCH)
127 .map_or(0, |d| d.as_secs() as i64)
128}
129
130#[cfg(test)]
131mod tests {
132 use super::*;
133
134 /// A throwaway database in a temp directory, migrated and ready.
135 pub async fn test_db() -> (Db, tempfile::TempDir) {
136 let dir = tempfile::tempdir().expect("temp dir");
137 let db = Db::open(&dir.path().join("test.db")).await.expect("open");
138 (db, dir)
139 }
140
141 #[tokio::test]
142 async fn migrations_apply_and_pragmas_take_effect() {
143 let (db, _dir) = test_db().await;
144
145 let journal: String = sqlx::query_scalar("PRAGMA journal_mode")
146 .fetch_one(&db.write)
147 .await
148 .expect("journal_mode");
149 assert_eq!(journal.to_lowercase(), "wal");
150
151 let fk: i64 = sqlx::query_scalar("PRAGMA foreign_keys")
152 .fetch_one(&db.write)
153 .await
154 .expect("foreign_keys");
155 assert_eq!(fk, 1, "foreign keys are off by default and must be enabled");
156 }
157
158 #[tokio::test]
159 async fn every_table_is_strict() {
160 let (db, _dir) = test_db().await;
161 let sql: Vec<String> = sqlx::query_scalar(
162 "SELECT sql FROM sqlite_master WHERE type = 'table' AND name NOT LIKE 'sqlite_%' \
163 AND name NOT LIKE '_sqlx%'",
164 )
165 .fetch_all(&db.read)
166 .await
167 .expect("schema");
168 assert!(!sql.is_empty());
169 for stmt in sql {
170 assert!(
171 stmt.to_uppercase().contains("STRICT"),
172 "table is not STRICT, so type affinity could store a string in an integer \
173 column:\n{stmt}"
174 );
175 }
176 }
177
178 #[tokio::test]
179 async fn the_read_pool_cannot_write() {
180 let (db, _dir) = test_db().await;
181 let err = sqlx::query("INSERT INTO settings (key, value) VALUES ('x', 'y')")
182 .execute(&db.read)
183 .await;
184 assert!(err.is_err(), "the read pool must be read-only");
185 }
186
187 #[tokio::test]
188 async fn the_points_primary_key_makes_replay_idempotent() {
189 let (db, _dir) = test_db().await;
190 sqlx::query(
191 "INSERT INTO users (id, username, pw_hash, display_name, created_at, pw_changed_at) \
192 VALUES (1, 'a', 'x', 'A', 0, 0)",
193 )
194 .execute(&db.write)
195 .await
196 .expect("user");
197
198 let insert = "INSERT INTO points (user_id, ts, lat, lon, acc_dm, recv_at) \
199 VALUES (1, 100, 5, 6, ?, 0) \
200 ON CONFLICT (user_id, ts) DO UPDATE SET \
201 acc_dm = excluded.acc_dm WHERE excluded.acc_dm < points.acc_dm";
202
203 sqlx::query(insert)
204 .bind(80)
205 .execute(&db.write)
206 .await
207 .expect("first");
208 // The same point again — a retry or a replay.
209 sqlx::query(insert)
210 .bind(80)
211 .execute(&db.write)
212 .await
213 .expect("replay");
214 // A second phone, same second, worse accuracy: must not win.
215 sqlx::query(insert)
216 .bind(200)
217 .execute(&db.write)
218 .await
219 .expect("worse");
220 // A second phone, same second, better accuracy: must win.
221 sqlx::query(insert)
222 .bind(30)
223 .execute(&db.write)
224 .await
225 .expect("better");
226
227 let (count, acc): (i64, i64) =
228 sqlx::query_as("SELECT COUNT(*), MIN(acc_dm) FROM points WHERE user_id = 1")
229 .fetch_one(&db.read)
230 .await
231 .expect("count");
232 assert_eq!(count, 1, "a replayed datagram must not create a second row");
233 assert_eq!(
234 acc, 30,
235 "the better-accuracy point must win a same-second collision"
236 );
237 }
238
239 #[tokio::test]
240 async fn a_share_must_target_exactly_one_of_user_or_group() {
241 let (db, _dir) = test_db().await;
242 for (id, username) in [(1, "a"), (2, "b")] {
243 sqlx::query(
244 "INSERT INTO users (id, username, pw_hash, display_name, created_at, pw_changed_at) \
245 VALUES (?, ?, 'x', 'X', 0, 0)",
246 )
247 .bind(id)
248 .bind(username)
249 .execute(&db.write)
250 .await
251 .expect("user");
252 }
253
254 // Neither target.
255 assert!(
256 sqlx::query("INSERT INTO shares (owner_user_id, created_at) VALUES (1, 0)")
257 .execute(&db.write)
258 .await
259 .is_err()
260 );
261 // Sharing with yourself.
262 assert!(
263 sqlx::query(
264 "INSERT INTO shares (owner_user_id, viewer_user_id, created_at) VALUES (1, 1, 0)"
265 )
266 .execute(&db.write)
267 .await
268 .is_err()
269 );
270 // A real share.
271 sqlx::query(
272 "INSERT INTO shares (owner_user_id, viewer_user_id, created_at) VALUES (1, 2, 0)",
273 )
274 .execute(&db.write)
275 .await
276 .expect("valid share");
277 }
278}
279