//! Database access. //! //! §8 imposes two structural requirements that this module exists to satisfy: //! //! 1. **A single writer connection, serialized through one owner**, with a read //! pool alongside. SQLite permits only one writer at a time even in WAL mode; //! pointing a multi-connection pool at writes and relying on `busy_timeout` //! to sort it out is explicitly rejected by the spec. Here the writer lives //! behind a `Mutex`, so contention queues in Rust rather than surfacing as //! `SQLITE_BUSY`. //! 2. **All access behind a thin repository layer** rather than queries //! scattered through handlers — this is what keeps the Turso/Postgres options //! cheap and localises the serialization in one place. //! //! rusqlite is synchronous, so every call is wrapped in `spawn_blocking`: a //! write that waits on the mutex must never block a Tokio worker thread. pub mod repo; use std::sync::{Arc, Mutex}; use anyhow::Context; use rusqlite::Connection; /// The whole schema, as portable SQL — it runs unchanged on Postgres, so /// SQLite-specific forms (`INSERT OR REPLACE`) are avoided in favour of the /// standard `INSERT ... ON CONFLICT` (§8). Keeping it in one `.sql` file rather /// than scattered through the repository is what makes that reviewable. /// /// TRACES: DR-014 | PR-004 const SCHEMA: &str = include_str!("schema.sql"); /// Handle to the database: one serialized writer, plus read connections. /// /// Cloning is cheap and shares the same underlying connections. /// TRACES: DR-003 | PR-004 #[derive(Clone)] pub struct Db { writer: Arc>, readers: Arc, } struct ReadPool { conns: Mutex>, path: String, } impl ReadPool { fn acquire(&self) -> anyhow::Result { if let Some(c) = self.conns.lock().expect("read pool poisoned").pop() { return Ok(c); } open_conn(&self.path, false) } fn release(&self, conn: Connection) { let mut conns = self.conns.lock().expect("read pool poisoned"); // Bounded: excess connections are dropped rather than accumulating. if conns.len() < 8 { conns.push(conn); } } } fn open_conn(path: &str, writer: bool) -> anyhow::Result { let conn = Connection::open(path).with_context(|| format!("opening database {path}"))?; // WAL gives concurrent readers alongside the single writer, which suits a // read-dominated workload; `synchronous = NORMAL` is safe under WAL, and // `busy_timeout` makes contention wait rather than error (§8). conn.pragma_update(None, "journal_mode", "WAL")?; conn.pragma_update(None, "synchronous", "NORMAL")?; conn.pragma_update(None, "busy_timeout", 5_000)?; conn.pragma_update(None, "foreign_keys", true)?; if !writer { conn.pragma_update(None, "query_only", true)?; } Ok(conn) } impl Db { /// Opens the database, applying the schema. Idempotent — every statement in /// `schema.sql` is `IF NOT EXISTS`. pub fn open(path: &str) -> anyhow::Result { let writer = open_conn(path, true)?; writer.execute_batch(SCHEMA).context("applying schema")?; Ok(Self { writer: Arc::new(Mutex::new(writer)), readers: Arc::new(ReadPool { conns: Mutex::new(Vec::new()), path: path.to_string() }), }) } /// Runs `f` against the serialized writer connection on a blocking thread. /// /// `f` receives a `Transaction`, so a manifest's scene rows go in as one /// transaction rather than one per row (§8), and a failure rolls back. pub async fn write(&self, f: F) -> anyhow::Result where T: Send + 'static, F: FnOnce(&rusqlite::Transaction<'_>) -> anyhow::Result + Send + 'static, { let writer = self.writer.clone(); tokio::task::spawn_blocking(move || { let mut conn = writer.lock().expect("writer poisoned"); let tx = conn.transaction()?; let out = f(&tx)?; tx.commit()?; Ok(out) }) .await .context("writer task panicked")? } /// Runs `f` against a read connection on a blocking thread. pub async fn read(&self, f: F) -> anyhow::Result where T: Send + 'static, F: FnOnce(&Connection) -> anyhow::Result + Send + 'static, { let readers = self.readers.clone(); tokio::task::spawn_blocking(move || { let conn = readers.acquire()?; let out = f(&conn); readers.release(conn); out }) .await .context("reader task panicked")? } } #[cfg(test)] mod tests { use super::*; #[tokio::test] async fn schema_applies_and_roundtrips() { let db = Db::open(":memory:").unwrap(); // In-memory databases are per-connection, so only exercise the writer. let n = db .write(|tx| { tx.execute( "INSERT INTO contributors (id, token_hash, created_at) VALUES (?1, ?2, ?3)", rusqlite::params!["c1", "hash", "2026-01-01T00:00:00Z"], )?; Ok(tx.query_row("SELECT COUNT(*) FROM contributors", [], |r| r.get::<_, i64>(0))?) }) .await .unwrap(); assert_eq!(n, 1); } #[tokio::test] async fn write_rolls_back_on_error() { let db = Db::open(":memory:").unwrap(); let res: anyhow::Result<()> = db .write(|tx| { tx.execute( "INSERT INTO contributors (id, token_hash, created_at) VALUES ('c1','h','t')", [], )?; anyhow::bail!("deliberate failure") }) .await; assert!(res.is_err()); let n = db .write(|tx| { Ok(tx.query_row("SELECT COUNT(*) FROM contributors", [], |r| r.get::<_, i64>(0))?) }) .await .unwrap(); assert_eq!(n, 0, "failed transaction must not persist rows"); } }