Skip to main content

recall_server/
store.rs

1//! SQLite persistence.
2//!
3//! The schema here is deliberately identical to the one the Node server
4//! created, because this server opens the *existing production database
5//! file* rather than migrating to a new one. That is what makes the cutover
6//! reversible: roll back by starting the old container against the same
7//! untouched file.
8
9use std::fs;
10use std::path::{Path, PathBuf};
11use std::sync::{Mutex, MutexGuard, PoisonError};
12
13use anyhow::{Context, Result};
14use recall_wire::{AdminTotals, File, ProjectStats};
15use rusqlite::{Connection, OptionalExtension};
16
17use crate::audit::merkle::Tree;
18use crate::now;
19
20mod audit;
21mod devices;
22mod jobs;
23mod passkeys;
24
25pub use audit::{AuditEntry, ConsistencyError, Outcome};
26pub use devices::{
27    plain_name, Created, Decision, Inserted, NewAuthkey, NewDevice, NewEnrollment, Poll, Waiting,
28};
29pub use jobs::{
30    clip, Failure, Queued, Retried, Settled, Settlement, MAX_ATTEMPTS, MAX_ERROR_BYTES, MAX_LINKS,
31    MAX_OPEN_JOBS,
32};
33pub use passkeys::{
34    AddedCredential, AdminCredential, AdminSession, BootstrapCode, FirstPasskey,
35    NewAdminCredential, RemovedCredential,
36};
37
38/// Frozen: an already-deployed database was created with exactly this.
39const SCHEMA: &str = "
40    CREATE TABLE IF NOT EXISTS memory_files (
41        project_key TEXT NOT NULL,
42        file_path   TEXT NOT NULL,
43        content     TEXT NOT NULL,
44        source_env  TEXT,
45        updated_at  TEXT NOT NULL,
46        deleted     INTEGER NOT NULL DEFAULT 0,
47        PRIMARY KEY (project_key, file_path)
48    );
49";
50
51/// The stored state of one file, used to decide whether a push needs
52/// merging.
53#[derive(Debug, Clone, PartialEq, Eq)]
54pub struct Existing {
55    /// The stored bytes. Present even for a tombstone — the server keeps the
56    /// last known content, it just refuses to hand it back over the wire.
57    pub content: String,
58    /// Whether the row is a tombstone.
59    pub deleted: bool,
60    /// The machine that wrote it; empty when none was recorded.
61    pub source_env: String,
62    /// When it was written.
63    pub updated_at: String,
64}
65
66/// The connection, plus the in-memory state built from it: the audit log's
67/// [`Tree`], rebuilt at open from `audit_log`, so an append, a checkpoint
68/// and a consistency proof each cost a few hashes rather than a pass over
69/// the whole table; and the newest leaf's `at`, which the next may not go
70/// below.
71///
72/// [`std::ops::Deref`] and [`std::ops::DerefMut`] to [`Connection`] mean
73/// every existing call site — `conn.execute(...)`, `conn.transaction()` —
74/// keeps compiling unchanged; only the audit-specific code reaches `audit`
75/// directly.
76struct StoreState {
77    conn: Connection,
78    audit: Tree,
79    audit_at: String,
80}
81
82impl std::ops::Deref for StoreState {
83    type Target = Connection;
84    fn deref(&self) -> &Connection {
85        &self.conn
86    }
87}
88
89impl std::ops::DerefMut for StoreState {
90    fn deref_mut(&mut self) -> &mut Connection {
91        &mut self.conn
92    }
93}
94
95/// The SQLite database, and every query the server makes against it.
96pub struct Store {
97    // A single connection behind a mutex. This is a single-owner server
98    // against a local file; a pool would buy nothing and SQLite would
99    // serialize the writes anyway. The audit log's append-then-commit
100    // relies on this too: every write already goes through this one lock,
101    // so a leaf and the state change it records are never interleaved with
102    // another request's.
103    state: Mutex<StoreState>,
104}
105
106impl Store {
107    /// Opens (creating if needed) the database at `path`.
108    pub fn open(path: impl AsRef<Path>) -> Result<Self> {
109        let path = path.as_ref();
110        if let Some(dir) = path.parent() {
111            if !dir.as_os_str().is_empty() {
112                fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
113            }
114        }
115        let conn = Connection::open(path).with_context(|| format!("opening {}", path.display()))?;
116        Self::with_connection(conn)
117    }
118
119    /// An in-memory database, for tests.
120    pub fn open_in_memory() -> Result<Self> {
121        Self::with_connection(Connection::open_in_memory()?)
122    }
123
124    fn with_connection(conn: Connection) -> Result<Self> {
125        let store = Self {
126            state: Mutex::new(StoreState {
127                conn,
128                audit: Tree::new(),
129                audit_at: String::new(),
130            }),
131        };
132        store.migrate()?;
133        Ok(store)
134    }
135
136    // A panic in one request must not render the whole store unusable, and
137    // nothing here leaves the database in a half-written state, so a
138    // poisoned mutex is recovered rather than propagated.
139    fn lock(&self) -> MutexGuard<'_, StoreState> {
140        self.state.lock().unwrap_or_else(PoisonError::into_inner)
141    }
142
143    fn migrate(&self) -> Result<()> {
144        let mut state = self.lock();
145        state.conn.execute_batch(SCHEMA)?;
146
147        // Databases created before tombstones existed have no `deleted`
148        // column. Adding it is safe and idempotent when guarded like this.
149        let has_deleted = {
150            let mut stmt = state.conn.prepare("PRAGMA table_info(memory_files)")?;
151            let mut rows = stmt.query([])?;
152            let mut found = false;
153            while let Some(row) = rows.next()? {
154                if row.get::<_, String>(1)? == "deleted" {
155                    found = true;
156                }
157            }
158            found
159        };
160        if !has_deleted {
161            state.conn.execute(
162                "ALTER TABLE memory_files ADD COLUMN deleted INTEGER NOT NULL DEFAULT 0",
163                [],
164            )?;
165        }
166
167        // Tables of their own, beside memory_files rather than in it, so a
168        // server from before devices existed still opens this file and
169        // simply never looks at them: rolling back stays a matter of
170        // starting the older image. Same for audit_log.
171        state.conn.execute_batch(devices::SCHEMA)?;
172        state.conn.execute_batch(passkeys::SCHEMA)?;
173        // The one change to an existing table the merge queue needs:
174        // devices may now have the worker scope, which SQLite can only
175        // allow by rebuilding the table. Then the queue's own table, beside
176        // the others for the same reason they are.
177        devices::allow_worker_scope(&state.conn)?;
178        state.conn.execute_batch(jobs::SCHEMA)?;
179        state.conn.execute_batch(audit::SCHEMA)?;
180
181        let loaded = audit::load(&state.conn)?;
182        state.audit = loaded.tree;
183        state.audit_at = loaded.last_at;
184        Ok(())
185    }
186
187    /// Reads one row.
188    pub fn get(&self, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
189        read_file(&self.lock(), project_key, file_path)
190    }
191
192    /// Writes content, clearing any tombstone, and appends the leaf
193    /// `build_leaf` makes from its `seq` and `at` in the same transaction.
194    /// Answers the time it was written, which is that `at`: the row's
195    /// `updated_at` and its leaf's `at` are one timestamp.
196    pub fn upsert_audited(
197        &self,
198        project_key: &str,
199        file_path: &str,
200        content: &str,
201        source_env: &str,
202        build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
203    ) -> Result<String> {
204        self.audited(
205            |tx, at| {
206                write_file(tx, project_key, file_path, content, source_env, at)?;
207                Ok(Outcome::Commit(at.to_string()))
208            },
209            |seq, at, _| build_leaf(seq, at),
210        )
211    }
212
213    /// Marks a file deleted while deliberately leaving its content in
214    /// place: a mistaken delete stays recoverable at the database level,
215    /// even though nothing in the app surfaces an undo yet. [`Store::list`]
216    /// withholds the content so a pull can't resurrect it.
217    ///
218    /// In the same transaction it closes the file's open merge jobs, so
219    /// nothing merges the deleted notes back into a file pushed after it,
220    /// and appends its leaf, answering the time, as
221    /// [`Store::upsert_audited`] does.
222    pub fn tombstone_audited(
223        &self,
224        project_key: &str,
225        file_path: &str,
226        source_env: &str,
227        build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
228    ) -> Result<String> {
229        self.audited(
230            |tx, at| {
231                tx.execute(
232                    "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
233                     VALUES (?1, ?2, '', ?3, ?4, 1)
234                     ON CONFLICT(project_key, file_path) DO UPDATE SET
235                         source_env = excluded.source_env,
236                         updated_at = excluded.updated_at,
237                         deleted = 1",
238                    (project_key, file_path, nullable(source_env), at),
239                )?;
240                jobs::close_for_delete(tx, project_key, file_path, at)?;
241                Ok(Outcome::Commit(at.to_string()))
242            },
243            |seq, at, _| build_leaf(seq, at),
244        )
245    }
246
247    /// Every file for a project, tombstones included so a pulling client
248    /// knows what to remove locally — but with the deleted content
249    /// withheld.
250    pub fn list(&self, project_key: &str) -> Result<Vec<File>> {
251        let conn = self.lock();
252        let mut stmt = conn.prepare(
253            "SELECT file_path, content, COALESCE(source_env, ''), updated_at, deleted
254             FROM memory_files WHERE project_key = ?1 ORDER BY file_path",
255        )?;
256        let rows = stmt.query_map((project_key,), |r| {
257            let content: String = r.get(1)?;
258            let deleted = r.get::<_, i64>(4)? != 0;
259            Ok(File {
260                file_path: r.get(0)?,
261                content: if deleted { None } else { Some(content) },
262                source_env: r.get(2)?,
263                updated_at: r.get(3)?,
264                deleted,
265            })
266        })?;
267        let mut files = Vec::new();
268        for row in rows {
269            files.push(row?);
270        }
271        Ok(files)
272    }
273
274    /// The most recent write across all projects, for `/health`. Empty when
275    /// nothing has ever been synced.
276    pub fn last_sync_at(&self) -> Result<String> {
277        let conn = self.lock();
278        let v: Option<String> =
279            conn.query_row("SELECT MAX(updated_at) FROM memory_files", [], |r| r.get(0))?;
280        Ok(v.unwrap_or_default())
281    }
282
283    /// Aggregates per project for the admin page.
284    pub fn admin_stats(&self) -> Result<(Vec<ProjectStats>, AdminTotals)> {
285        let conn = self.lock();
286
287        let mut projects = Vec::new();
288        let mut totals = AdminTotals::default();
289        {
290            let mut stmt = conn.prepare(
291                "SELECT project_key,
292                        SUM(CASE WHEN deleted = 0 THEN 1 ELSE 0 END),
293                        SUM(CASE WHEN deleted = 1 THEN 1 ELSE 0 END),
294                        MAX(updated_at)
295                 FROM memory_files GROUP BY project_key ORDER BY MAX(updated_at) DESC",
296            )?;
297            let rows = stmt.query_map([], |r| {
298                Ok(ProjectStats {
299                    project_key: r.get(0)?,
300                    file_count: r.get(1)?,
301                    deleted_count: r.get(2)?,
302                    sources: Vec::new(),
303                    last_updated_at: r.get::<_, Option<String>>(3)?.unwrap_or_default(),
304                })
305            })?;
306            for row in rows {
307                let p = row?;
308                totals.file_count += p.file_count;
309                totals.deleted_count += p.deleted_count;
310                projects.push(p);
311            }
312        }
313        totals.project_count = projects.len() as i64;
314
315        // Sources are collected separately rather than with GROUP_CONCAT:
316        // SQLite won't take a custom separator together with DISTINCT, and
317        // source_env is client-supplied, so a value containing a comma
318        // would silently split into bogus entries.
319        {
320            let mut stmt = conn.prepare(
321                "SELECT DISTINCT project_key, source_env FROM memory_files WHERE source_env IS NOT NULL",
322            )?;
323            let rows =
324                stmt.query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?)))?;
325            for row in rows {
326                let (key, src) = row?;
327                if src.is_empty() {
328                    continue;
329                }
330                if let Some(p) = projects.iter_mut().find(|p| p.project_key == key) {
331                    p.sources.push(src);
332                }
333            }
334        }
335        for p in &mut projects {
336            p.sources.sort();
337        }
338        Ok((projects, totals))
339    }
340
341    /// Writes a consistent snapshot via `VACUUM INTO`, which is safe
342    /// against a live database — unlike copying the file — then prunes the
343    /// oldest snapshots beyond `keep`.
344    pub fn backup(&self, dir: impl AsRef<Path>, keep: usize) -> Result<PathBuf> {
345        let dir = dir.as_ref();
346        fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
347
348        // Mirrors the Node server's naming, which was
349        // toISOString().replace(/[:.]/g, "-"). Millisecond precision
350        // matters twice over: without it two snapshots in the same second
351        // collide (and VACUUM INTO refuses to overwrite an existing file),
352        // and the prune below sorts these names lexicographically alongside
353        // any snapshots the Node server already wrote into the directory.
354        let stamp = now().replace([':', '.'], "-");
355        let dest = dir.join(format!("recall-{stamp}.db"));
356        let dest_str = dest
357            .to_str()
358            .context("backup path is not valid UTF-8")?
359            .to_owned();
360
361        // VACUUM INTO refuses a file that is already there, so one that is
362        // there now is not this call's to delete if the vacuum fails.
363        let existed = dest.exists();
364        let vacuumed = {
365            let conn = self.lock();
366            conn.execute("VACUUM INTO ?1", (&dest_str,))
367                .with_context(|| format!("VACUUM INTO {dest_str}"))
368        };
369        if let Err(err) = vacuumed {
370            // A vacuum that fails part way (a full disk, a file size limit)
371            // leaves the part it wrote. Left there, it is named like every
372            // good snapshot, sorts among them, and counts toward `keep`,
373            // so the prune would delete a good one to make room for it,
374            // and anyone restoring would find a file that does not open.
375            if !existed {
376                let _ = fs::remove_file(&dest);
377            }
378            return Err(err);
379        }
380
381        let mut snapshots: Vec<PathBuf> = fs::read_dir(dir)?
382            .filter_map(|e| e.ok())
383            .map(|e| e.path())
384            .filter(|p| {
385                p.file_name()
386                    .and_then(|n| n.to_str())
387                    .is_some_and(|n| n.starts_with("recall-") && n.ends_with(".db"))
388            })
389            .collect();
390        snapshots.sort();
391        for stale in snapshots.iter().take(snapshots.len().saturating_sub(keep)) {
392            let _ = fs::remove_file(stale);
393        }
394        Ok(dest)
395    }
396}
397
398#[cfg(test)]
399impl Store {
400    /// Test-only: runs `f` against the raw connection. Used to assert things
401    /// no public method goes anywhere near on purpose, such as the audit
402    /// log's append-only triggers refusing a raw `UPDATE` or `DELETE`.
403    pub(crate) fn with_raw<T>(
404        &self,
405        f: impl FnOnce(&Connection) -> rusqlite::Result<T>,
406    ) -> rusqlite::Result<T> {
407        f(&self.lock())
408    }
409}
410
411/// Test-only: a leaf for the store's own tests, which are about rows rather
412/// than what a leaf says. The store has no way to write without one.
413#[cfg(test)]
414pub(crate) fn test_leaf(seq: u64, at: &str) -> Vec<u8> {
415    use crate::audit::leaf;
416    leaf::encode(
417        seq,
418        at,
419        leaf::action::START,
420        &leaf::Actor::Server,
421        leaf::subject_start("test"),
422        None,
423    )
424}
425
426/// What [`existing_from`] reads, in its order.
427const EXISTING_COLUMNS: &str = "content, deleted, COALESCE(source_env, ''), updated_at";
428
429fn existing_from(r: &rusqlite::Row<'_>) -> rusqlite::Result<Existing> {
430    Ok(Existing {
431        content: r.get(0)?,
432        deleted: r.get::<_, i64>(1)? != 0,
433        source_env: r.get(2)?,
434        updated_at: r.get(3)?,
435    })
436}
437
438/// One row, on a connection or transaction the caller already holds.
439fn read_file(conn: &Connection, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
440    Ok(conn
441        .query_row(
442            &format!("SELECT {EXISTING_COLUMNS} FROM memory_files WHERE project_key = ?1 AND file_path = ?2"),
443            (project_key, file_path),
444            existing_from,
445        )
446        .optional()?)
447}
448
449/// Writes content, clearing any tombstone, on a connection or transaction
450/// the caller already holds: the one statement behind
451/// [`Store::upsert_audited`] and a merged result being applied, each in a
452/// transaction that appends its leaf.
453fn write_file(
454    conn: &Connection,
455    project_key: &str,
456    file_path: &str,
457    content: &str,
458    source_env: &str,
459    updated_at: &str,
460) -> Result<()> {
461    conn.execute(
462        "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
463         VALUES (?1, ?2, ?3, ?4, ?5, 0)
464         ON CONFLICT(project_key, file_path) DO UPDATE SET
465             content = excluded.content,
466             source_env = excluded.source_env,
467             updated_at = excluded.updated_at,
468             deleted = 0",
469        (project_key, file_path, content, nullable(source_env), updated_at),
470    )?;
471    Ok(())
472}
473
474/// An absent `source_env` is stored as NULL, not `''` — `admin_stats`
475/// distinguishes the two.
476fn nullable(s: &str) -> Option<&str> {
477    if s.is_empty() {
478        None
479    } else {
480        Some(s)
481    }
482}
483
484// What `recall-server admin` does to a database. A child module so it can
485// share the one connection type and `Store::backup` without widening what
486// `Store` exposes; crate-private, because nothing but that subcommand, and
487// nothing reachable from the HTTP router, may call it.
488pub(crate) mod admin;
489
490#[cfg(test)]
491mod tests {
492    use super::*;
493
494    fn store() -> Store {
495        Store::open_in_memory().unwrap()
496    }
497
498    fn put(st: &Store, project_key: &str, file_path: &str, content: &str, source_env: &str) {
499        st.upsert_audited(project_key, file_path, content, source_env, test_leaf)
500            .unwrap();
501    }
502
503    fn del(st: &Store, project_key: &str, file_path: &str, source_env: &str) {
504        st.tombstone_audited(project_key, file_path, source_env, test_leaf)
505            .unwrap();
506    }
507
508    /// The row's `updated_at` is its leaf's `at`: one moment, taken under
509    /// the lock the write holds, answered to the caller.
510    #[test]
511    fn a_write_is_stamped_with_its_leafs_at() {
512        let st = store();
513        let mut leaf_at = String::new();
514        let updated_at = st
515            .upsert_audited("acme/app", "a.md", "x", "laptop", |seq, at| {
516                leaf_at = at.to_string();
517                test_leaf(seq, at)
518            })
519            .unwrap();
520        assert_eq!(updated_at, leaf_at);
521        assert_eq!(st.list("acme/app").unwrap()[0].updated_at, updated_at);
522        assert_eq!(st.audit_checkpoint().0, 1);
523    }
524
525    #[test]
526    fn upsert_get_and_list_round_trip() {
527        let st = store();
528        put(&st, "acme/app", "MEMORY.md", "hello", "laptop");
529
530        let got = st.get("acme/app", "MEMORY.md").unwrap().unwrap();
531        assert_eq!(got.content, "hello");
532        assert!(!got.deleted);
533
534        let files = st.list("acme/app").unwrap();
535        assert_eq!(files.len(), 1);
536        assert_eq!(files[0].content.as_deref(), Some("hello"));
537        assert_eq!(files[0].source_env, "laptop");
538        assert!(st.get("acme/app", "missing.md").unwrap().is_none());
539    }
540
541    /// Both halves of the tombstone contract in one place: the row keeps
542    /// its content, the listing does not hand it back.
543    #[test]
544    fn tombstone_preserves_content_but_list_withholds_it() {
545        let st = store();
546        put(&st, "acme/app", "gone.md", "secret", "laptop");
547        del(&st, "acme/app", "gone.md", "laptop");
548
549        let row = st.get("acme/app", "gone.md").unwrap().unwrap();
550        assert_eq!(row.content, "secret", "content must stay recoverable");
551        assert!(row.deleted);
552
553        let files = st.list("acme/app").unwrap();
554        assert_eq!(
555            files.len(),
556            1,
557            "tombstones are listed so clients can delete locally"
558        );
559        assert!(files[0].deleted);
560        assert_eq!(files[0].content, None, "a pull must not resurrect it");
561    }
562
563    /// A push after a delete revives the row and clears the tombstone.
564    #[test]
565    fn upsert_clears_a_tombstone() {
566        let st = store();
567        del(&st, "acme/app", "f.md", "laptop");
568        put(&st, "acme/app", "f.md", "back", "laptop");
569        let row = st.get("acme/app", "f.md").unwrap().unwrap();
570        assert!(!row.deleted);
571        assert_eq!(row.content, "back");
572    }
573
574    #[test]
575    fn last_sync_at_is_empty_on_a_fresh_database() {
576        assert_eq!(store().last_sync_at().unwrap(), "");
577    }
578
579    /// A `source_env` containing a comma must survive as one value — the
580    /// reason sources aren't gathered with GROUP_CONCAT.
581    #[test]
582    fn admin_stats_keeps_commas_inside_a_source_env() {
583        let st = store();
584        put(&st, "acme/app", "a.md", "x", "laptop,evil");
585        let (projects, _) = st.admin_stats().unwrap();
586        assert_eq!(projects[0].sources, vec!["laptop,evil".to_string()]);
587    }
588
589    /// The `deleted` column is added to databases that predate tombstones,
590    /// without touching their rows.
591    #[test]
592    fn migrates_a_database_that_predates_tombstones() {
593        let dir = tempfile::tempdir().unwrap();
594        let path = dir.path().join("old.db");
595        {
596            let conn = Connection::open(&path).unwrap();
597            conn.execute_batch(
598                "CREATE TABLE memory_files (
599                    project_key TEXT NOT NULL,
600                    file_path   TEXT NOT NULL,
601                    content     TEXT NOT NULL,
602                    source_env  TEXT,
603                    updated_at  TEXT NOT NULL,
604                    PRIMARY KEY (project_key, file_path)
605                );
606                INSERT INTO memory_files VALUES ('acme/app','old.md','kept','node-era','2026-09-03T21:49:55.191Z');",
607            )
608            .unwrap();
609        }
610        let st = Store::open(&path).unwrap();
611        let files = st.list("acme/app").unwrap();
612        assert_eq!(files.len(), 1);
613        assert_eq!(files[0].content.as_deref(), Some("kept"));
614        assert!(!files[0].deleted);
615    }
616
617    #[test]
618    fn backup_names_carry_milliseconds() {
619        let dir = tempfile::tempdir().unwrap();
620        let st = store();
621        let dest = st.backup(dir.path(), 7).unwrap();
622        let name = dest.file_name().unwrap().to_str().unwrap();
623        // recall-2026-09-03T21-49-55-191Z.db
624        assert!(
625            name.starts_with("recall-") && name.ends_with("Z.db"),
626            "got {name}"
627        );
628        let stamp = &name["recall-".len()..name.len() - ".db".len()];
629        assert_eq!(stamp.len(), 24, "got {stamp}");
630        // The three characters before the Z are the milliseconds, which
631        // keep two snapshots in the same second from colliding.
632        assert!(
633            stamp[20..23].chars().all(|c| c.is_ascii_digit()),
634            "no millisecond field in {stamp}"
635        );
636    }
637}