lib.rsannotatedlib.rssource258 lines · 11.5 KB · raw

Whiskers' storage on SQLite: one database file per device and one for the service, with the pictures as blobs in it. See README.md, and SCHEMA.md for every table and how it merges.

The merge rules are not here. They are pure functions in whiskers-core (MemoryDoc::merge, Household::merge, ChatState::merge); this crate keeps the result, and keeps it so that the rules about what must be gone (a forgotten memory) are properties of the schema and not of the code that happens to run.

9#![forbid(unsafe_code)]
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;

The numbered migrations, in order. Migration n takes a database at user_version n - 1 to n. A 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")];

The schema version this build writes.

46pub const SCHEMA_VERSION: u32 = MIGRATIONS.len() as u32;

One database, shared by everything that keeps state in it. Cloning gives another handle on the same connection: every operation runs under one lock, so each is atomic with respect to every other, which is the atomicity the hubs' contracts ask for.

51#[derive(Clone)]
52pub struct Store {
53    conn: Arc<Mutex<Connection>>,
54    path: Option<PathBuf>,
55}
57impl Store {

Opens (creating it if need be) the database at path, brings its schema up to this build's and 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    }

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    }
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    }

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    }

The path of the database file, if it has one.

128    pub fn path(&self) -> Option<&Path> {
129        self.path.as_deref()
130    }
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    }

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    }

Empties the write-ahead log into the database file and truncates it to nothing, so that what was 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    }

A consistent copy of the whole database at dest, taken with SQLite's backup interface while the database stays in use. The copy is a plain database file, with no write-ahead log beside it. It carries every row, device rows of the time count included; the engine's backup leaves out the 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}

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}
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}

A consistent copy of the database file at src into dest, for a database some Store may have open (the app 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}

What a database file holds that a backup is described by.

208#[derive(Debug, PartialEq, Eq)]
209pub struct Inspection {
210    pub pictures: u32,
211}

Opens the file at path as a Whiskers database and reads it, to say whether this build can use it: it is brought to this build's schema if it is older, and refused (DbError::Newer) if it is newer or (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}
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}