lib.rsannotatedlib.rssource258 lines · 11.5 KB · raw
1//! Whiskers' storage on SQLite: one database file per device and one for the service, with the pictures
2//! as blobs in it. See README.md, and `SCHEMA.md` for every table and how it merges.
3//!
4//! The merge rules are not here. They are pure functions in `whiskers-core` (`MemoryDoc::merge`,
5//! `Household::merge`, `ChatState::merge`); this crate keeps the result, and keeps it so that the rules
6//! about what must be gone (a forgotten memory) are properties of the schema and not of the code that
7//! happens to run.
8
9#![forbid(unsafe_code)]
10
11mod chat;
12mod error;
13mod household;
14mod hubs;
15mod journal;
16mod log;
17mod memory;
18mod pictures;
19mod reflections;
20mod speech;
21mod spend;
22
23use std::path::{Path, PathBuf};
24use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
25
26use ::log::{debug, error, info, warn};
27use rusqlite::Connection;
28
29pub use chat::SqliteChat;
30pub use error::DbError;
31pub use household::SqliteHousehold;
32pub use hubs::{SqliteChatHub, SqliteHouseholdHub, SqliteMemoryHub};
33pub use journal::SqliteJournal;
34pub use log::{EXCHANGE_SLACK_MS, LogDay, MAX_DAYS_PAGE, Replica, SqliteLog};
35pub use memory::SqliteMemory;
36pub use pictures::{SqlitePictureShelf, SqlitePictures};
37pub use reflections::{SqliteReflections, Waiting};
38pub use speech::SqliteSpeechCache;
39pub use spend::SqliteAllowance;
40
41/// The numbered migrations, in order. Migration `n` takes a database at `user_version` `n - 1` to `n`. A
42/// migration is never edited once a build that has it is released: a change is the next number.
43const MIGRATIONS: &[&str] = &[include_str!("../migrations/0001_init.sql"), include_str!("../migrations/0002_forgetting_empties_the_log.sql"), include_str!("../migrations/0003_soft_actions.sql"), include_str!("../migrations/0004_speech_cache.sql"), include_str!("../migrations/0005_household_characters.sql"), include_str!("../migrations/0006_forgetting_the_turn.sql")];
44
45/// The schema version this build writes.
46pub const SCHEMA_VERSION: u32 = MIGRATIONS.len() as u32;
47
48/// One database, shared by everything that keeps state in it. Cloning gives another handle on the same
49/// connection: every operation runs under one lock, so each is atomic with respect to every other, which
50/// is the atomicity the hubs' contracts ask for.
51#[derive(Clone)]
52pub struct Store {
53    conn: Arc<Mutex<Connection>>,
54    path: Option<PathBuf>,
55}
56
57impl Store {
58    /// Opens (creating it if need be) the database at `path`, brings its schema up to this build's and
59    /// sets the pragmas the rest relies on. A database written by a newer build is refused, not read.
60    pub fn open(path: &Path) -> Result<Store, DbError> {
61        debug!("opening the database {}", path.display());
62        let conn = Connection::open(path).map_err(|e| {
63            error!("cannot open the database {}: {e}", path.display());
64            DbError::from(e)
65        })?;
66        Self::prepare(conn, Some(path.to_owned()))
67    }
68
69    /// A database that lives only as long as its handles: for tests.
70    pub fn open_in_memory() -> Result<Store, DbError> {
71        Self::prepare(Connection::open_in_memory()?, None)
72    }
73
74    fn prepare(conn: Connection, path: Option<PathBuf>) -> Result<Store, DbError> {
75        conn.busy_timeout(std::time::Duration::from_secs(5))?;
76        // `secure_delete` makes SQLite overwrite what it deletes (with zeros) in the database file and in
77        // the pages the write-ahead log hands back, which is what lets "forgotten" mean gone. It is a
78        // per-connection setting, so it is set on every open and never left to the caller.
79        conn.pragma_update(None, "secure_delete", "ON")?;
80        conn.pragma_update(None, "foreign_keys", "ON")?;
81        conn.pragma_update(None, "synchronous", "FULL")?;
82        if path.is_some() {
83            // The write-ahead log gives crash safety without rewriting the file on every change. Anything
84            // that must be gone from the log as well as the file asks for a checkpoint that truncates it
85            // (see `Store::scrub`), and so does opening a database, so a crash cannot leave old bytes behind.
86            let mode: String = conn.query_row("PRAGMA journal_mode = WAL", [], |r| r.get(0))?;
87            if mode != "wal" {
88                warn!("the database would not enter write-ahead mode (it says {mode})");
89            }
90        }
91        let store = Store { conn: Arc::new(Mutex::new(conn)), path };
92        // A log left over from a process that stopped after a forgetting and before it was emptied is emptied
93        // now, before anything else is written. (A database with no log has nothing to heal.)
94        if store.path.as_deref().is_some_and(|p| std::fs::metadata(wal_of(p)).is_ok_and(|m| m.len() > 0)) {
95            store.scrub()?;
96        }
97        let migrated = Self::migrate(&mut store.lock())?;
98        if migrated {
99            // A migration may have emptied words out of the log: the old bytes must not wait in the write-ahead log.
100            store.scrub()?;
101        }
102        info!("database ready at schema version {SCHEMA_VERSION}");
103        Ok(store)
104    }
105
106    /// Whether anything was migrated.
107    fn migrate(conn: &mut Connection) -> Result<bool, DbError> {
108        let found: u32 = conn.query_row("PRAGMA user_version", [], |r| r.get(0))?;
109        if found > SCHEMA_VERSION {
110            error!("the database is at schema version {found}, newer than this build's {SCHEMA_VERSION}");
111            return Err(DbError::Newer { found, known: SCHEMA_VERSION });
112        }
113        let mut ran = false;
114        for (i, sql) in MIGRATIONS.iter().enumerate().skip(found as usize) {
115            ran = true;
116            let version = i as u32 + 1;
117            info!("migrating the database to schema version {version}");
118            let tx = conn.transaction()?;
119            tx.execute_batch(sql)?;
120            // `user_version` cannot be bound, and `version` is a number made here.
121            tx.execute_batch(&format!("PRAGMA user_version = {version}"))?;
122            tx.commit()?;
123        }
124        Ok(ran)
125    }
126
127    /// The path of the database file, if it has one.
128    pub fn path(&self) -> Option<&Path> {
129        self.path.as_deref()
130    }
131
132    pub(crate) fn lock(&self) -> MutexGuard<'_, Connection> {
133        // A panic elsewhere must not stop the child's memory.
134        self.conn.lock().unwrap_or_else(|e| {
135            warn!("the database lock was poisoned by a panic elsewhere; carrying on");
136            PoisonError::into_inner(e)
137        })
138    }
139
140    /// Runs `f` in one transaction: all of it or none of it.
141    pub(crate) fn transaction<T>(&self, f: impl FnOnce(&rusqlite::Transaction<'_>) -> Result<T, DbError>) -> Result<T, DbError> {
142        let mut conn = self.lock();
143        let tx = conn.transaction()?;
144        let out = f(&tx)?;
145        tx.commit()?;
146        Ok(out)
147    }
148
149    /// Empties the write-ahead log into the database file and truncates it to nothing, so that what was
150    /// deleted is in neither. Run after every forgetting, and when a database is opened.
151    pub fn scrub(&self) -> Result<(), DbError> {
152        if self.path.is_none() {
153            return Ok(());
154        }
155        let conn = self.lock();
156        let (busy, _log, _moved): (i64, i64, i64) = conn.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?;
157        if busy != 0 {
158            error!("the write-ahead log could not be truncated (another connection is reading); deleted bytes may still be in it");
159            return Err(DbError::Busy);
160        }
161        Ok(())
162    }
163
164    /// A consistent copy of the whole database at `dest`, taken with SQLite's backup interface while the
165    /// database stays in use. The copy is a plain database file, with no write-ahead log beside it. It
166    /// carries every row, `device` rows of the time count included; the engine's backup leaves out the
167    /// device's own identity, which is a file and not a row.
168    pub fn backup_to(&self, dest: &Path) -> Result<Inspection, DbError> {
169        copy_into(&self.lock(), dest)
170    }
171}
172
173/// The write-ahead log of the database at `db`.
174fn wal_of(db: &Path) -> PathBuf {
175    let mut name = db.file_name().map(|n| n.to_os_string()).unwrap_or_default();
176    name.push("-wal");
177    db.with_file_name(name)
178}
179
180fn copy_into(conn: &Connection, dest: &Path) -> Result<Inspection, DbError> {
181    use rusqlite::backup::{Backup, StepResult};
182    let mut out = Connection::open(dest)?;
183    {
184        let backup = Backup::new(conn, &mut out)?;
185        loop {
186            match backup.step(256)? {
187                StepResult::Done => break,
188                StepResult::More => {}
189                StepResult::Busy | StepResult::Locked => std::thread::sleep(std::time::Duration::from_millis(20)),
190                _ => {}
191            }
192        }
193    }
194    // The copy is made into a database of its own: leave it in rollback-journal mode so one file is all of it.
195    out.pragma_update(None, "journal_mode", "DELETE")?;
196    let pictures: i64 = out.query_row("SELECT COUNT(*) FROM picture", [], |r| r.get(0))?;
197    Ok(Inspection { pictures: pictures as u32 })
198}
199
200/// A consistent copy of the database file at `src` into `dest`, for a database some `Store` may have open (the app
201/// is running): it reads through a connection of its own and takes nothing from the one in use. Says what the copy holds.
202pub fn copy_database(src: &Path, dest: &Path) -> Result<Inspection, DbError> {
203    let conn = Connection::open_with_flags(src, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX)?;
204    copy_into(&conn, dest)
205}
206
207/// What a database file holds that a backup is described by.
208#[derive(Debug, PartialEq, Eq)]
209pub struct Inspection {
210    pub pictures: u32,
211}
212
213/// Opens the file at `path` as a Whiskers database and reads it, to say whether this build can use it: it is
214/// brought to this build's schema if it is older, and refused (`DbError::Newer`) if it is newer or
215/// (`DbError::is_unreadable`) if it is not a database at all.
216pub fn inspect(path: &Path) -> Result<Inspection, DbError> {
217    let store = Store::open(path)?;
218    let pictures: i64 = store.lock().query_row("SELECT COUNT(*) FROM picture", [], |r| r.get(0))?;
219    Ok(Inspection { pictures: pictures as u32 })
220}
221
222#[cfg(test)]
223mod tests {
224    use super::*;
225
226    #[test]
227    fn a_new_database_is_at_the_current_version_and_a_newer_one_is_refused() {
228        let dir = tempdir("version");
229        let path = dir.join("w.db");
230        drop(Store::open(&path).unwrap());
231        let version: u32 = Connection::open(&path).unwrap().query_row("PRAGMA user_version", [], |r| r.get(0)).unwrap();
232        assert_eq!(version, SCHEMA_VERSION);
233        Connection::open(&path).unwrap().execute_batch(&format!("PRAGMA user_version = {}", SCHEMA_VERSION + 1)).unwrap();
234        assert!(matches!(Store::open(&path), Err(DbError::Newer { .. })), "a database from a newer build is not read");
235    }
236
237    #[test]
238    fn opening_again_changes_nothing_and_secure_delete_is_on() {
239        let dir = tempdir("again");
240        let path = dir.join("w.db");
241        let a = Store::open(&path).unwrap();
242        drop(a);
243        let b = Store::open(&path).unwrap();
244        let on: i64 = b.lock().query_row("PRAGMA secure_delete", [], |r| r.get(0)).unwrap();
245        assert_eq!(on, 1);
246        let fk: i64 = b.lock().query_row("PRAGMA foreign_keys", [], |r| r.get(0)).unwrap();
247        assert_eq!(fk, 1);
248    }
249
250    pub(crate) fn tempdir(what: &str) -> PathBuf {
251        use std::sync::atomic::{AtomicUsize, Ordering};
252        static N: AtomicUsize = AtomicUsize::new(0);
253        let d = std::env::temp_dir().join(format!("whiskers-store-{what}-{}-{}", std::process::id(), N.fetch_add(1, Ordering::SeqCst)));
254        let _ = std::fs::remove_dir_all(&d);
255        std::fs::create_dir_all(&d).unwrap();
256        d
257    }
258}