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};
12use std::time::Duration;
13
14use anyhow::{Context, Result};
15use recall_wire::{AdminTotals, File, ProjectStats};
16use rusqlite::{Connection, OptionalExtension};
17
18use crate::audit::merkle::Tree;
19use crate::now;
20
21mod audit;
22mod devices;
23mod evaluations;
24mod jobs;
25mod passkeys;
26
27pub use audit::{AuditEntry, ConsistencyError, Outcome};
28pub use devices::{
29    plain_name, Created, Decision, Inserted, NewAuthkey, NewDevice, NewEnrollment, Poll, Waiting,
30};
31pub use evaluations::Requested;
32pub use jobs::{
33    clip, Failure, Queued, Retried, Settled, Settlement, MAX_ATTEMPTS, MAX_ERROR_BYTES, MAX_LINKS,
34    MAX_OPEN_JOBS,
35};
36pub use passkeys::{
37    AddedCredential, AdminCredential, AdminSession, BootstrapCode, FirstPasskey,
38    NewAdminCredential, RemovedCredential,
39};
40
41/// Frozen: an already-deployed database was created with exactly this.
42const SCHEMA: &str = "
43    CREATE TABLE IF NOT EXISTS memory_files (
44        project_key TEXT NOT NULL,
45        file_path   TEXT NOT NULL,
46        content     TEXT NOT NULL,
47        source_env  TEXT,
48        updated_at  TEXT NOT NULL,
49        deleted     INTEGER NOT NULL DEFAULT 0,
50        PRIMARY KEY (project_key, file_path)
51    );
52";
53
54/// The stored state of one file, used to decide whether a push needs
55/// merging.
56#[derive(Debug, Clone, PartialEq, Eq)]
57pub struct Existing {
58    /// The stored bytes. Present even for a tombstone — the server keeps the
59    /// last known content, it just refuses to hand it back over the wire.
60    pub content: String,
61    /// Whether the row is a tombstone.
62    pub deleted: bool,
63    /// The machine that wrote it; empty when none was recorded.
64    pub source_env: String,
65    /// When it was written.
66    pub updated_at: String,
67}
68
69/// The connection, plus the in-memory state built from it: the audit log's
70/// [`Tree`], rebuilt at open from `audit_log`, so an append, a checkpoint
71/// and a consistency proof each cost a few hashes rather than a pass over
72/// the whole table; and the newest leaf's `at`, which the next may not go
73/// below.
74///
75/// [`std::ops::Deref`] and [`std::ops::DerefMut`] to [`Connection`] mean
76/// every existing call site — `conn.execute(...)`, `conn.transaction()` —
77/// keeps compiling unchanged; only the audit-specific code reaches `audit`
78/// directly.
79///
80/// `file` is the database file this connection opened, as (device, inode),
81/// which every audited write checks is still the file at its path: see
82/// [`FileId`]. When it is not, `log_moved` says whether to say so on stderr
83/// as well as in the refusal (the server does; an admin command, which
84/// prints the refusal itself, does not), and `moved_said` that it has.
85struct StoreState {
86    conn: Connection,
87    audit: Tree,
88    audit_at: String,
89    file: Option<FileId>,
90    log_moved: bool,
91    moved_said: bool,
92}
93
94/// Which file a path names: its device and inode.
95///
96/// A connection keeps the file it opened, whatever happens to the path. So
97/// a `recall.db` moved aside or replaced under a running server, which is
98/// what a restore done without stopping it does, leaves the server writing
99/// into the file it had, now under another name or no name at all, and
100/// answering 200 for writes nobody will read. Under the rollback journal
101/// the next read noticed the new file's change counter; under WAL nothing
102/// does, since a connection goes by its WAL index and its cache. So every
103/// audited write, which is every push and pull and every other change the
104/// log records, first checks that the path still names the file it
105/// opened, and refuses (a 500 to the client, who keeps its copy and pushes
106/// again later) if not.
107///
108/// A file overwritten in place keeps its inode and is not seen; copying
109/// over a live `recall.db` is what `deploy/README.md` forbids outright.
110/// Only on Unix, where a file has an inode to compare: elsewhere this is
111/// always `None` and nothing is checked.
112type FileId = (u64, u64);
113
114/// The [`FileId`] of the file `conn` has open, by its path; `None` for an
115/// in-memory database, a path that no longer names a file, or a platform
116/// without inodes.
117#[cfg(unix)]
118fn file_id(conn: &Connection) -> Option<FileId> {
119    use std::os::unix::fs::MetadataExt;
120    let path = conn.path().filter(|p| !p.is_empty())?;
121    fs::metadata(path).ok().map(|m| (m.dev(), m.ino()))
122}
123
124#[cfg(not(unix))]
125fn file_id(_conn: &Connection) -> Option<FileId> {
126    None
127}
128
129impl std::ops::Deref for StoreState {
130    type Target = Connection;
131    fn deref(&self) -> &Connection {
132        &self.conn
133    }
134}
135
136impl std::ops::DerefMut for StoreState {
137    fn deref_mut(&mut self) -> &mut Connection {
138        &mut self.conn
139    }
140}
141
142/// The SQLite database, and every query the server makes against it.
143pub struct Store {
144    // A single connection behind a mutex. This is a single-owner server
145    // against a local file; a pool would buy nothing and SQLite would
146    // serialize the writes anyway. The audit log's append-then-commit
147    // relies on this too: every write already goes through this one lock,
148    // so a leaf and the state change it records are never interleaved with
149    // another request's.
150    state: Mutex<StoreState>,
151}
152
153impl Store {
154    /// Opens (creating if needed) the database at `path`, switching it to
155    /// SQLite's WAL journal with every commit synced before it returns (see
156    /// `use_durable_wal` in this module for why both). A file an older
157    /// server wrote with the rollback journal is converted here, once; the
158    /// mode is stored in the file. Fails, rather than serving in another
159    /// mode, if the switch cannot be made: the file is held by another
160    /// process for longer than the busy timeout, or SQLite cannot keep the
161    /// WAL's index beside it. A network filesystem may well not fail here
162    /// and still not work; `deploy/README.md` says not to use one.
163    pub fn open(path: impl AsRef<Path>) -> Result<Self> {
164        Self::open_waiting(path.as_ref(), admin::BUSY_TIMEOUT)
165    }
166
167    /// [`Store::open`], its connection waiting `busy` rather than
168    /// [`admin::BUSY_TIMEOUT`] on a lock another process holds. The server
169    /// always opens through `open`; taking the wait apart is for the test
170    /// of what happens past it, which need not sit through five seconds.
171    fn open_waiting(path: &Path, busy: Duration) -> Result<Self> {
172        if let Some(dir) = path.parent() {
173            if !dir.as_os_str().is_empty() {
174                fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
175            }
176        }
177        let conn = Connection::open(path).with_context(|| format!("opening {}", path.display()))?;
178        use_durable_wal(&conn, busy).with_context(|| {
179            format!(
180                "switching {} to SQLite's WAL journal. It needs a local filesystem, and a \
181                 moment with no other process holding the file (sqlite-web mid-read, an admin \
182                 command): start the server again",
183                path.display()
184            )
185        })?;
186        Self::with_connection(conn)
187    }
188
189    /// An in-memory database, for tests.
190    pub fn open_in_memory() -> Result<Self> {
191        Self::with_connection(Connection::open_in_memory()?)
192    }
193
194    fn with_connection(conn: Connection) -> Result<Self> {
195        let file = file_id(&conn);
196        let store = Self {
197            state: Mutex::new(StoreState {
198                conn,
199                audit: Tree::new(),
200                audit_at: String::new(),
201                file,
202                log_moved: true,
203                moved_said: false,
204            }),
205        };
206        store.migrate()?;
207        Ok(store)
208    }
209
210    // A panic in one request must not render the whole store unusable, and
211    // nothing here leaves the database in a half-written state, so a
212    // poisoned mutex is recovered rather than propagated.
213    fn lock(&self) -> MutexGuard<'_, StoreState> {
214        self.state.lock().unwrap_or_else(PoisonError::into_inner)
215    }
216
217    fn migrate(&self) -> Result<()> {
218        let mut state = self.lock();
219        state.conn.execute_batch(SCHEMA)?;
220
221        // Databases created before tombstones existed have no `deleted`
222        // column. Adding it is safe and idempotent when guarded like this.
223        let has_deleted = {
224            let mut stmt = state.conn.prepare("PRAGMA table_info(memory_files)")?;
225            let mut rows = stmt.query([])?;
226            let mut found = false;
227            while let Some(row) = rows.next()? {
228                if row.get::<_, String>(1)? == "deleted" {
229                    found = true;
230                }
231            }
232            found
233        };
234        if !has_deleted {
235            state.conn.execute(
236                "ALTER TABLE memory_files ADD COLUMN deleted INTEGER NOT NULL DEFAULT 0",
237                [],
238            )?;
239        }
240
241        // Tables of their own, beside memory_files rather than in it, so a
242        // server from before devices existed still opens this file and
243        // simply never looks at them: rolling back stays a matter of
244        // starting the older image. Same for audit_log.
245        state.conn.execute_batch(devices::SCHEMA)?;
246        state.conn.execute_batch(passkeys::SCHEMA)?;
247        // The one change to an existing table the merge queue needs:
248        // devices may now have the worker scope, which SQLite can only
249        // allow by rebuilding the table. Then the queue's own table, beside
250        // the others for the same reason they are.
251        devices::allow_worker_scope(&state.conn)?;
252        state.conn.execute_batch(jobs::SCHEMA)?;
253        state.conn.execute_batch(evaluations::SCHEMA)?;
254        state.conn.execute_batch(audit::SCHEMA)?;
255
256        let loaded = audit::load(&state.conn)?;
257        state.audit = loaded.tree;
258        state.audit_at = loaded.last_at;
259        Ok(())
260    }
261
262    /// Reads one row.
263    pub fn get(&self, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
264        read_file(&self.lock(), project_key, file_path)
265    }
266
267    /// Writes content, clearing any tombstone, and appends the leaf
268    /// `build_leaf` makes from its `seq` and `at` in the same transaction.
269    /// Answers the time it was written, which is that `at`: the row's
270    /// `updated_at` and its leaf's `at` are one timestamp.
271    pub fn upsert_audited(
272        &self,
273        project_key: &str,
274        file_path: &str,
275        content: &str,
276        source_env: &str,
277        build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
278    ) -> Result<String> {
279        self.audited(
280            |tx, at| {
281                write_file(tx, project_key, file_path, content, source_env, at)?;
282                Ok(Outcome::Commit(at.to_string()))
283            },
284            |seq, at, _| build_leaf(seq, at),
285        )
286    }
287
288    /// Marks a file deleted while deliberately leaving its content in
289    /// place: a mistaken delete stays recoverable at the database level,
290    /// even though nothing in the app surfaces an undo yet. [`Store::list`]
291    /// withholds the content so a pull can't resurrect it.
292    ///
293    /// In the same transaction it closes the file's open merge jobs, so
294    /// nothing merges the deleted notes back into a file pushed after it,
295    /// and appends its leaf, answering the time, as
296    /// [`Store::upsert_audited`] does.
297    pub fn tombstone_audited(
298        &self,
299        project_key: &str,
300        file_path: &str,
301        source_env: &str,
302        build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
303    ) -> Result<String> {
304        self.audited(
305            |tx, at| {
306                tx.execute(
307                    "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
308                     VALUES (?1, ?2, '', ?3, ?4, 1)
309                     ON CONFLICT(project_key, file_path) DO UPDATE SET
310                         source_env = excluded.source_env,
311                         updated_at = excluded.updated_at,
312                         deleted = 1",
313                    (project_key, file_path, nullable(source_env), at),
314                )?;
315                jobs::close_for_delete(tx, project_key, file_path, at)?;
316                Ok(Outcome::Commit(at.to_string()))
317            },
318            |seq, at, _| build_leaf(seq, at),
319        )
320    }
321
322    /// Every file for a project, tombstones included so a pulling client
323    /// knows what to remove locally — but with the deleted content
324    /// withheld.
325    pub fn list(&self, project_key: &str) -> Result<Vec<File>> {
326        let conn = self.lock();
327        let mut stmt = conn.prepare(
328            "SELECT file_path, content, COALESCE(source_env, ''), updated_at, deleted
329             FROM memory_files WHERE project_key = ?1 ORDER BY file_path",
330        )?;
331        let rows = stmt.query_map((project_key,), |r| {
332            let content: String = r.get(1)?;
333            let deleted = r.get::<_, i64>(4)? != 0;
334            Ok(File {
335                file_path: r.get(0)?,
336                content: if deleted { None } else { Some(content) },
337                source_env: r.get(2)?,
338                updated_at: r.get(3)?,
339                deleted,
340            })
341        })?;
342        let mut files = Vec::new();
343        for row in rows {
344            files.push(row?);
345        }
346        Ok(files)
347    }
348
349    /// The most recent write across all projects, for `/health`. Empty when
350    /// nothing has ever been synced.
351    pub fn last_sync_at(&self) -> Result<String> {
352        let conn = self.lock();
353        let v: Option<String> =
354            conn.query_row("SELECT MAX(updated_at) FROM memory_files", [], |r| r.get(0))?;
355        Ok(v.unwrap_or_default())
356    }
357
358    /// Aggregates per project for the admin page.
359    pub fn admin_stats(&self) -> Result<(Vec<ProjectStats>, AdminTotals)> {
360        let conn = self.lock();
361
362        let mut projects = Vec::new();
363        let mut totals = AdminTotals::default();
364        {
365            let mut stmt = conn.prepare(
366                "SELECT project_key,
367                        SUM(CASE WHEN deleted = 0 THEN 1 ELSE 0 END),
368                        SUM(CASE WHEN deleted = 1 THEN 1 ELSE 0 END),
369                        MAX(updated_at)
370                 FROM memory_files GROUP BY project_key ORDER BY MAX(updated_at) DESC",
371            )?;
372            let rows = stmt.query_map([], |r| {
373                Ok(ProjectStats {
374                    project_key: r.get(0)?,
375                    file_count: r.get(1)?,
376                    deleted_count: r.get(2)?,
377                    sources: Vec::new(),
378                    last_updated_at: r.get::<_, Option<String>>(3)?.unwrap_or_default(),
379                })
380            })?;
381            for row in rows {
382                let p = row?;
383                totals.file_count += p.file_count;
384                totals.deleted_count += p.deleted_count;
385                projects.push(p);
386            }
387        }
388        totals.project_count = projects.len() as i64;
389
390        // Sources are collected separately rather than with GROUP_CONCAT:
391        // SQLite won't take a custom separator together with DISTINCT, and
392        // source_env is client-supplied, so a value containing a comma
393        // would silently split into bogus entries.
394        {
395            let mut stmt = conn.prepare(
396                "SELECT DISTINCT project_key, source_env FROM memory_files WHERE source_env IS NOT NULL",
397            )?;
398            let rows =
399                stmt.query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?)))?;
400            for row in rows {
401                let (key, src) = row?;
402                if src.is_empty() {
403                    continue;
404                }
405                if let Some(p) = projects.iter_mut().find(|p| p.project_key == key) {
406                    p.sources.push(src);
407                }
408            }
409        }
410        for p in &mut projects {
411            p.sources.sort();
412        }
413        Ok((projects, totals))
414    }
415
416    /// Copies the commits the WAL holds into the database file, as far as no
417    /// reader still needs them, waiting on nothing. Answers whether that
418    /// was all of them.
419    ///
420    /// SQLite already does this by itself once the WAL passes 1000 pages,
421    /// which on a personal server can be days of pushes. The server also
422    /// does it every sweep, so that `recall.db` on its own, all a reader
423    /// that cannot see `recall.db-wal` has (sqlite-web's single-file mount
424    /// in `deploy/docker-compose.direct.yml`), is never far behind. Nothing
425    /// relies on it for correctness: the WAL is part of the database. A
426    /// reader that keeps one transaction open holds back everything
427    /// committed after it began, and the WAL grows until it lets go; the
428    /// server says so when that lasts several sweeps.
429    pub fn checkpoint(&self) -> Result<bool> {
430        let (frames, copied): (i64, i64) =
431            self.lock()
432                .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |r| {
433                    Ok((r.get(1)?, r.get(2)?))
434                })?;
435        Ok(copied >= frames)
436    }
437
438    /// Copies every commit the WAL holds into the database file and empties
439    /// the WAL, waiting up to the busy timeout for a reader in the way.
440    /// Answers whether it got all the way.
441    ///
442    /// The last thing a server does as it stops. Closing the connection
443    /// does the same, but only when no other process has the file open, and
444    /// sqlite-web keeps it open. Emptied here, the WAL left beside the file
445    /// holds nothing, so a stopped server's `recall.db` is the whole
446    /// database. Not something to rely on after a crash, which is why every
447    /// procedure in `deploy/README.md` moves or copies the WAL with the file.
448    pub fn checkpoint_all(&self) -> Result<bool> {
449        let busy: i64 = self
450            .lock()
451            .query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |r| r.get(0))?;
452        Ok(busy == 0)
453    }
454
455    /// Writes a consistent snapshot via `VACUUM INTO`, then prunes the
456    /// oldest snapshots beyond `keep`.
457    ///
458    /// `VACUUM INTO` reads through this connection, so the snapshot holds
459    /// every commit, those still only in the WAL included, as of one
460    /// moment, whatever is writing meanwhile. Copying `recall.db` instead
461    /// would miss whatever the WAL has not yet handed back to it. And what
462    /// it writes is one self-contained file in the rollback journal's mode,
463    /// with no `-wal` of its own, so a snapshot can be copied, uploaded,
464    /// opened read-only, or copied over `recall.db` in a restore, as it is.
465    pub fn backup(&self, dir: impl AsRef<Path>, keep: usize) -> Result<PathBuf> {
466        let dir = dir.as_ref();
467        fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
468
469        // Mirrors the Node server's naming, which was
470        // toISOString().replace(/[:.]/g, "-"). Millisecond precision
471        // matters twice over: without it two snapshots in the same second
472        // collide (and VACUUM INTO refuses to overwrite an existing file),
473        // and the prune below sorts these names lexicographically alongside
474        // any snapshots the Node server already wrote into the directory.
475        let stamp = now().replace([':', '.'], "-");
476        let dest = dir.join(format!("recall-{stamp}.db"));
477        let dest_str = dest
478            .to_str()
479            .context("backup path is not valid UTF-8")?
480            .to_owned();
481
482        // VACUUM INTO refuses a file that is already there, so one that is
483        // there now is not this call's to delete if the vacuum fails.
484        let existed = dest.exists();
485        let vacuumed = {
486            let conn = self.lock();
487            conn.execute("VACUUM INTO ?1", (&dest_str,))
488                .with_context(|| format!("VACUUM INTO {dest_str}"))
489        };
490        if let Err(err) = vacuumed {
491            // A vacuum that fails part way (a full disk, a file size limit)
492            // leaves the part it wrote. Left there, it is named like every
493            // good snapshot, sorts among them, and counts toward `keep`,
494            // so the prune would delete a good one to make room for it,
495            // and anyone restoring would find a file that does not open.
496            if !existed {
497                let _ = fs::remove_file(&dest);
498            }
499            return Err(err);
500        }
501
502        let mut snapshots: Vec<PathBuf> = fs::read_dir(dir)?
503            .filter_map(|e| e.ok())
504            .map(|e| e.path())
505            .filter(|p| {
506                p.file_name()
507                    .and_then(|n| n.to_str())
508                    .is_some_and(|n| n.starts_with("recall-") && n.ends_with(".db"))
509            })
510            .collect();
511        snapshots.sort();
512        for stale in snapshots.iter().take(snapshots.len().saturating_sub(keep)) {
513            let _ = fs::remove_file(stale);
514        }
515        Ok(dest)
516    }
517}
518
519#[cfg(test)]
520impl Store {
521    /// Test-only: runs `f` against the raw connection. Used to assert things
522    /// no public method goes anywhere near on purpose, such as the audit
523    /// log's append-only triggers refusing a raw `UPDATE` or `DELETE`.
524    pub(crate) fn with_raw<T>(
525        &self,
526        f: impl FnOnce(&Connection) -> rusqlite::Result<T>,
527    ) -> rusqlite::Result<T> {
528        f(&self.lock())
529    }
530}
531
532/// Test-only: a leaf for the store's own tests, which are about rows rather
533/// than what a leaf says. The store has no way to write without one.
534#[cfg(test)]
535pub(crate) fn test_leaf(seq: u64, at: &str) -> Vec<u8> {
536    use crate::audit::leaf;
537    leaf::encode(
538        seq,
539        at,
540        leaf::action::START,
541        &leaf::Actor::Server,
542        leaf::subject_start("test"),
543        None,
544    )
545}
546
547/// Puts the server's connection in WAL mode with `synchronous=FULL`, and
548/// states its busy timeout.
549///
550/// WAL for two reasons. A commit appends its pages to `recall.db-wal` and
551/// syncs that one file, where the rollback journal wrote and synced a
552/// journal, then the database, then deleted the journal; since the audit
553/// log, a pull is a write transaction too, so every sync request pays for
554/// one. And a reader no longer holds up a writer: sqlite-web mid-read, or
555/// an admin command reading its backup, used to stall a push's commit for
556/// up to the busy timeout and fail it past that. Readers see the last
557/// commit before they began, and the one writer at a time is unchanged:
558/// every write here still goes through the store's lock, and an admin
559/// command's through `BEGIN IMMEDIATE`, so the audit log's tree and its
560/// triggers see exactly what they did.
561///
562/// FULL rather than NORMAL, WAL's usual partner. Under NORMAL a commit is
563/// not synced until the next checkpoint, so a power cut or a crash of the
564/// host (not of this process, whose writes the kernel already has) can
565/// take back the last pushes after their client was told 200, and that
566/// client may be a cloud session that no longer exists, whose memory this
567/// server was the only copy of. FULL syncs the WAL at every commit. The
568/// ignored `push_latency_by_journal_mode` test measures what that costs.
569/// On the VM this was written on (a release build, through the router,
570/// medians), a push took 1.5 ms under the rollback journal, 0.12 ms under
571/// WAL with NORMAL and 0.58 ms under WAL with FULL, and a pull 1.2, 0.07
572/// and 0.29 ms. The one sync is the difference between the last two, and
573/// grows with a slower disk, as the rollback journal's several did; a 200
574/// that means stored is worth it.
575///
576/// Both settings are the connection's, and the journal mode is also
577/// stored in the file, so every other connection (an admin command,
578/// sqlite-web, a `sqlite3` shell, an older recall-server) opens it in WAL
579/// too; an admin command states FULL for itself. What WAL asks of anything
580/// that copies or replaces the file is in `deploy/README.md`: `recall.db`
581/// alone is not the whole database while `recall.db-wal` holds commits not
582/// yet copied back into it, and a WAL left beside a replaced `recall.db`
583/// is replayed into it.
584fn use_durable_wal(conn: &Connection, busy: Duration) -> Result<()> {
585    conn.busy_timeout(busy)?;
586    let mode: String = conn.query_row("PRAGMA journal_mode = WAL", [], |r| r.get(0))?;
587    if !mode.eq_ignore_ascii_case("wal") {
588        anyhow::bail!("SQLite kept the {mode} journal");
589    }
590    conn.execute_batch(&format!(
591        "PRAGMA synchronous = FULL; PRAGMA journal_size_limit = {WAL_SIZE_LIMIT}"
592    ))?;
593    Ok(())
594}
595
596/// What `recall.db-wal` is cut back to, in bytes, when SQLite next starts
597/// it over. The file is reused rather than shrunk as it cycles, so without
598/// a limit it stays as large as it ever grew: after one large transaction
599/// (an admin restore of a big key), or after a reader held checkpoints
600/// back for a day. 16 MiB is four times what the WAL reaches between
601/// SQLite's own checkpoints (1000 pages of 4 KiB), so ordinary traffic
602/// never meets it.
603const WAL_SIZE_LIMIT: i64 = 16 * 1024 * 1024;
604
605/// What [`existing_from`] reads, in its order.
606const EXISTING_COLUMNS: &str = "content, deleted, COALESCE(source_env, ''), updated_at";
607
608fn existing_from(r: &rusqlite::Row<'_>) -> rusqlite::Result<Existing> {
609    Ok(Existing {
610        content: r.get(0)?,
611        deleted: r.get::<_, i64>(1)? != 0,
612        source_env: r.get(2)?,
613        updated_at: r.get(3)?,
614    })
615}
616
617/// One row, on a connection or transaction the caller already holds.
618fn read_file(conn: &Connection, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
619    Ok(conn
620        .query_row(
621            &format!("SELECT {EXISTING_COLUMNS} FROM memory_files WHERE project_key = ?1 AND file_path = ?2"),
622            (project_key, file_path),
623            existing_from,
624        )
625        .optional()?)
626}
627
628/// Writes content, clearing any tombstone, on a connection or transaction
629/// the caller already holds: the one statement behind
630/// [`Store::upsert_audited`] and a merged result being applied, each in a
631/// transaction that appends its leaf.
632fn write_file(
633    conn: &Connection,
634    project_key: &str,
635    file_path: &str,
636    content: &str,
637    source_env: &str,
638    updated_at: &str,
639) -> Result<()> {
640    conn.execute(
641        "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
642         VALUES (?1, ?2, ?3, ?4, ?5, 0)
643         ON CONFLICT(project_key, file_path) DO UPDATE SET
644             content = excluded.content,
645             source_env = excluded.source_env,
646             updated_at = excluded.updated_at,
647             deleted = 0",
648        (project_key, file_path, content, nullable(source_env), updated_at),
649    )?;
650    Ok(())
651}
652
653/// An absent `source_env` is stored as NULL, not `''` — `admin_stats`
654/// distinguishes the two.
655fn nullable(s: &str) -> Option<&str> {
656    if s.is_empty() {
657        None
658    } else {
659        Some(s)
660    }
661}
662
663// What `recall-server admin` does to a database. A child module so it can
664// share the one connection type and `Store::backup` without widening what
665// `Store` exposes; crate-private, because nothing but that subcommand, and
666// nothing reachable from the HTTP router, may call it.
667pub(crate) mod admin;
668
669#[cfg(test)]
670mod tests {
671    use super::*;
672
673    fn store() -> Store {
674        Store::open_in_memory().unwrap()
675    }
676
677    fn put(st: &Store, project_key: &str, file_path: &str, content: &str, source_env: &str) {
678        st.upsert_audited(project_key, file_path, content, source_env, test_leaf)
679            .unwrap();
680    }
681
682    fn del(st: &Store, project_key: &str, file_path: &str, source_env: &str) {
683        st.tombstone_audited(project_key, file_path, source_env, test_leaf)
684            .unwrap();
685    }
686
687    /// The row's `updated_at` is its leaf's `at`: one moment, taken under
688    /// the lock the write holds, answered to the caller.
689    #[test]
690    fn a_write_is_stamped_with_its_leafs_at() {
691        let st = store();
692        let mut leaf_at = String::new();
693        let updated_at = st
694            .upsert_audited("acme/app", "a.md", "x", "laptop", |seq, at| {
695                leaf_at = at.to_string();
696                test_leaf(seq, at)
697            })
698            .unwrap();
699        assert_eq!(updated_at, leaf_at);
700        assert_eq!(st.list("acme/app").unwrap()[0].updated_at, updated_at);
701        assert_eq!(st.audit_checkpoint().0, 1);
702    }
703
704    #[test]
705    fn upsert_get_and_list_round_trip() {
706        let st = store();
707        put(&st, "acme/app", "MEMORY.md", "hello", "laptop");
708
709        let got = st.get("acme/app", "MEMORY.md").unwrap().unwrap();
710        assert_eq!(got.content, "hello");
711        assert!(!got.deleted);
712
713        let files = st.list("acme/app").unwrap();
714        assert_eq!(files.len(), 1);
715        assert_eq!(files[0].content.as_deref(), Some("hello"));
716        assert_eq!(files[0].source_env, "laptop");
717        assert!(st.get("acme/app", "missing.md").unwrap().is_none());
718    }
719
720    /// Both halves of the tombstone contract in one place: the row keeps
721    /// its content, the listing does not hand it back.
722    #[test]
723    fn tombstone_preserves_content_but_list_withholds_it() {
724        let st = store();
725        put(&st, "acme/app", "gone.md", "secret", "laptop");
726        del(&st, "acme/app", "gone.md", "laptop");
727
728        let row = st.get("acme/app", "gone.md").unwrap().unwrap();
729        assert_eq!(row.content, "secret", "content must stay recoverable");
730        assert!(row.deleted);
731
732        let files = st.list("acme/app").unwrap();
733        assert_eq!(
734            files.len(),
735            1,
736            "tombstones are listed so clients can delete locally"
737        );
738        assert!(files[0].deleted);
739        assert_eq!(files[0].content, None, "a pull must not resurrect it");
740    }
741
742    /// A push after a delete revives the row and clears the tombstone.
743    #[test]
744    fn upsert_clears_a_tombstone() {
745        let st = store();
746        del(&st, "acme/app", "f.md", "laptop");
747        put(&st, "acme/app", "f.md", "back", "laptop");
748        let row = st.get("acme/app", "f.md").unwrap().unwrap();
749        assert!(!row.deleted);
750        assert_eq!(row.content, "back");
751    }
752
753    #[test]
754    fn last_sync_at_is_empty_on_a_fresh_database() {
755        assert_eq!(store().last_sync_at().unwrap(), "");
756    }
757
758    /// A `source_env` containing a comma must survive as one value — the
759    /// reason sources aren't gathered with GROUP_CONCAT.
760    #[test]
761    fn admin_stats_keeps_commas_inside_a_source_env() {
762        let st = store();
763        put(&st, "acme/app", "a.md", "x", "laptop,evil");
764        let (projects, _) = st.admin_stats().unwrap();
765        assert_eq!(projects[0].sources, vec!["laptop,evil".to_string()]);
766    }
767
768    /// The `deleted` column is added to databases that predate tombstones,
769    /// without touching their rows.
770    #[test]
771    fn migrates_a_database_that_predates_tombstones() {
772        let dir = tempfile::tempdir().unwrap();
773        let path = dir.path().join("old.db");
774        {
775            let conn = Connection::open(&path).unwrap();
776            conn.execute_batch(
777                "CREATE TABLE memory_files (
778                    project_key TEXT NOT NULL,
779                    file_path   TEXT NOT NULL,
780                    content     TEXT NOT NULL,
781                    source_env  TEXT,
782                    updated_at  TEXT NOT NULL,
783                    PRIMARY KEY (project_key, file_path)
784                );
785                INSERT INTO memory_files VALUES ('acme/app','old.md','kept','node-era','2026-09-03T21:49:55.191Z');",
786            )
787            .unwrap();
788        }
789        let st = Store::open(&path).unwrap();
790        let files = st.list("acme/app").unwrap();
791        assert_eq!(files.len(), 1);
792        assert_eq!(files[0].content.as_deref(), Some("kept"));
793        assert!(!files[0].deleted);
794    }
795
796    #[test]
797    fn backup_names_carry_milliseconds() {
798        let dir = tempfile::tempdir().unwrap();
799        let st = store();
800        let dest = st.backup(dir.path(), 7).unwrap();
801        let name = dest.file_name().unwrap().to_str().unwrap();
802        // recall-2026-09-03T21-49-55-191Z.db
803        assert!(
804            name.starts_with("recall-") && name.ends_with("Z.db"),
805            "got {name}"
806        );
807        let stamp = &name["recall-".len()..name.len() - ".db".len()];
808        assert_eq!(stamp.len(), 24, "got {stamp}");
809        // The three characters before the Z are the milliseconds, which
810        // keep two snapshots in the same second from colliding.
811        assert!(
812            stamp[20..23].chars().all(|c| c.is_ascii_digit()),
813            "no millisecond field in {stamp}"
814        );
815    }
816
817    // ------------------------------------------------------------ journal
818
819    fn journal_mode(conn: &Connection) -> String {
820        conn.query_row("PRAGMA journal_mode", [], |r| r.get(0))
821            .unwrap()
822    }
823
824    /// Bytes 18 and 19 of the header: 1 for the rollback journal, 2 for
825    /// WAL. What SQLite, any version, reads the mode from.
826    fn header_mode(path: &Path) -> (u8, u8) {
827        let head = fs::read(path).unwrap();
828        (head[18], head[19])
829    }
830
831    fn wal_len(db: &Path) -> u64 {
832        let mut wal = db.as_os_str().to_owned();
833        wal.push("-wal");
834        fs::metadata(wal).map(|m| m.len()).unwrap_or(0)
835    }
836
837    /// The server's connection is in WAL, syncs every commit (FULL, not
838    /// NORMAL: see `use_durable_wal`) and waits the admin commands' busy
839    /// timeout; and the mode is the file's, so any other connection opens
840    /// it in WAL too.
841    #[test]
842    fn the_store_keeps_the_file_in_wal_and_syncs_every_commit() {
843        let dir = tempfile::tempdir().unwrap();
844        let path = dir.path().join("recall.db");
845        let st = Store::open(&path).unwrap();
846        put(&st, "acme/app", "a.md", "x", "laptop");
847        let (mode, sync, busy) = st
848            .with_raw(|c| {
849                Ok((
850                    journal_mode(c),
851                    c.query_row("PRAGMA synchronous", [], |r| r.get::<_, i64>(0))?,
852                    c.query_row("PRAGMA busy_timeout", [], |r| r.get::<_, i64>(0))?,
853                ))
854            })
855            .unwrap();
856        assert_eq!(mode, "wal");
857        assert_eq!(sync, 2, "synchronous=FULL");
858        assert_eq!(busy, admin::BUSY_TIMEOUT.as_millis() as i64);
859        let limit: i64 = st
860            .with_raw(|c| c.query_row("PRAGMA journal_size_limit", [], |r| r.get(0)))
861            .unwrap();
862        assert_eq!(limit, WAL_SIZE_LIMIT);
863
864        assert_eq!(header_mode(&path), (2, 2));
865        assert_eq!(journal_mode(&Connection::open(&path).unwrap()), "wal");
866        assert!(wal_len(&path) > 0, "the commit went to the WAL");
867    }
868
869    /// The SQLite compiled in is 3.51.3 or later, the first with the fix
870    /// for a checkpoint on one connection racing a write that starts the
871    /// WAL over on another (the server's and an admin command's, under
872    /// WAL). A rusqlite pinned back, or a bundle that goes back, fails here
873    /// rather than in someone's database.
874    #[test]
875    fn the_sqlite_compiled_in_has_the_wal_restart_fix() {
876        assert!(
877            rusqlite::version_number() >= 3_051_003,
878            "SQLite {} predates 3.51.3",
879            rusqlite::version()
880        );
881    }
882
883    /// The upgrade: a file a server before WAL wrote, rows, audit log and
884    /// all, in the rollback journal it always had. The first open converts
885    /// it, loses nothing, and the log reads back to the same tree.
886    #[test]
887    fn a_rollback_journal_database_from_an_older_server_is_converted_on_open() {
888        let dir = tempfile::tempdir().unwrap();
889        let path = dir.path().join("recall.db");
890        let before = {
891            let st = Store::open(&path).unwrap();
892            put(&st, "acme/app", "MEMORY.md", "kept\n", "laptop");
893            del(&st, "acme/app", "gone.md", "laptop");
894            st.audit_checkpoint()
895        };
896        // As 0.4.2 left it: the rollback journal, which is what a server
897        // that never asked for WAL got.
898        Connection::open(&path)
899            .unwrap()
900            .query_row("PRAGMA journal_mode = DELETE", [], |_| Ok(()))
901            .unwrap();
902        assert_eq!(header_mode(&path), (1, 1));
903
904        let st = Store::open(&path).unwrap();
905        assert_eq!(header_mode(&path), (2, 2));
906        let files = st.list("acme/app").unwrap();
907        assert_eq!(files.len(), 2);
908        assert_eq!(files[0].content.as_deref(), Some("kept\n"));
909        assert!(files[1].deleted);
910        assert_eq!(st.audit_checkpoint(), before);
911        put(&st, "acme/app", "after.md", "y", "laptop");
912        drop(st);
913
914        // And back: a connection that asks for nothing, as an older server
915        // rolled back to, reads the converted file, every commit included.
916        let plain = Connection::open(&path).unwrap();
917        let n: i64 = plain
918            .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
919            .unwrap();
920        assert_eq!(n, 3);
921    }
922
923    /// The upgrade's one way to fail: another process reading the old
924    /// file (sqlite-web, say) through the whole busy timeout, which under
925    /// the rollback journal keeps the switch from committing. The server
926    /// refuses to start, saying why and what to do, rather than serving in
927    /// the old mode; the file is left exactly as it was, and the next
928    /// start, with the reader gone, converts it.
929    #[test]
930    fn a_switch_held_up_past_the_busy_timeout_refuses_to_start_and_changes_nothing() {
931        let dir = tempfile::tempdir().unwrap();
932        let path = dir.path().join("recall.db");
933        drop(Store::open(&path).unwrap());
934        Connection::open(&path)
935            .unwrap()
936            .query_row("PRAGMA journal_mode = DELETE", [], |_| Ok(()))
937            .unwrap();
938        let reader = Connection::open(&path).unwrap();
939        reader.execute_batch("BEGIN").unwrap();
940        let _: i64 = reader
941            .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
942            .unwrap();
943
944        // Waiting a quarter of a second rather than the server's five: the
945        // reader holds on until after the open has failed, so it is held
946        // up past whatever the wait is, and what is checked is what
947        // happens then, which does not depend on how long that was.
948        let err = match Store::open_waiting(&path, Duration::from_millis(250)) {
949            Ok(_) => panic!("switched to WAL under a reader holding the file"),
950            Err(e) => format!("{e:#}"),
951        };
952        assert!(err.contains("switching"), "{err}");
953        assert!(err.contains("start the server again"), "{err}");
954        assert!(err.contains("locked"), "{err}");
955        assert_eq!(header_mode(&path), (1, 1), "still the rollback journal");
956
957        reader.execute_batch("COMMIT").unwrap();
958        drop(reader);
959        drop(Store::open(&path).unwrap());
960        assert_eq!(header_mode(&path), (2, 2));
961    }
962
963    /// The same, from the oldest file there is: the one the Node server
964    /// wrote, which production's rows came from (`scripts/compat-check.sh`
965    /// drives a release binary against it). Every row reads back the same
966    /// after the switch.
967    #[test]
968    fn the_node_servers_database_is_converted_with_every_row() {
969        let dir = tempfile::tempdir().unwrap();
970        let path = dir.path().join("recall.db");
971        fs::copy(
972            concat!(
973                env!("CARGO_MANIFEST_DIR"),
974                "/../../fixtures/node-written.db"
975            ),
976            &path,
977        )
978        .unwrap();
979        assert_eq!(header_mode(&path), (1, 1));
980        let dump = || -> Vec<(String, String, String, Option<String>, String, i64)> {
981            let conn = Connection::open(&path).unwrap();
982            let mut stmt = conn
983                .prepare(
984                    "SELECT project_key, file_path, content, source_env, updated_at, deleted
985                     FROM memory_files ORDER BY project_key, file_path",
986                )
987                .unwrap();
988            let rows = stmt
989                .query_map([], |r| {
990                    Ok((
991                        r.get(0)?,
992                        r.get(1)?,
993                        r.get(2)?,
994                        r.get(3)?,
995                        r.get(4)?,
996                        r.get(5)?,
997                    ))
998                })
999                .unwrap();
1000            rows.map(Result::unwrap).collect()
1001        };
1002        let before = dump();
1003        assert!(!before.is_empty());
1004
1005        let st = Store::open(&path).unwrap();
1006        assert_eq!(header_mode(&path), (2, 2));
1007        assert_eq!(dump(), before);
1008        drop(st);
1009        assert_eq!(dump(), before);
1010    }
1011
1012    /// What a restore copies over `recall.db`: one file, in the rollback
1013    /// journal's mode, with no WAL of its own, holding the commits the live
1014    /// file does not have yet because they are still only in its WAL.
1015    #[test]
1016    fn a_snapshot_is_one_self_contained_file_with_what_the_wal_holds() {
1017        let dir = tempfile::tempdir().unwrap();
1018        let path = dir.path().join("recall.db");
1019        let st = Store::open(&path).unwrap();
1020        for i in 0..20 {
1021            put(&st, "acme/app", &format!("f{i}.md"), "x", "laptop");
1022        }
1023        assert!(wal_len(&path) > 0, "the pushes are still in the WAL");
1024
1025        let snap = st.backup(dir.path().join("backups"), 7).unwrap();
1026        let names: Vec<_> = fs::read_dir(snap.parent().unwrap())
1027            .unwrap()
1028            .map(|e| e.unwrap().file_name().into_string().unwrap())
1029            .collect();
1030        assert_eq!(names.len(), 1, "no -wal or -shm beside it: {names:?}");
1031        assert_eq!(header_mode(&snap), (1, 1));
1032
1033        let read =
1034            Connection::open_with_flags(&snap, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY).unwrap();
1035        assert_eq!(journal_mode(&read), "delete");
1036        let n: i64 = read
1037            .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1038            .unwrap();
1039        assert_eq!(n, 20);
1040    }
1041
1042    /// How many rows `recall.db` holds on its own: the main file alone,
1043    /// copied into a directory of its own with no WAL beside it and opened
1044    /// there, which is all a reader that cannot see `recall.db-wal` has
1045    /// (sqlite-web's single-file mount with direct TLS).
1046    fn rows_in_the_file_alone(db: &Path) -> i64 {
1047        let alone = tempfile::tempdir().unwrap();
1048        let copy = alone.path().join("recall.db");
1049        fs::copy(db, &copy).unwrap();
1050        Connection::open(&copy)
1051            .unwrap()
1052            .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1053            .unwrap()
1054    }
1055
1056    /// The file on its own is behind while the WAL holds commits, caught up
1057    /// by the sweep's checkpoint, and whole after the one a server takes as
1058    /// it stops, which also empties the WAL.
1059    #[test]
1060    fn a_checkpoint_brings_the_file_on_its_own_up_to_date() {
1061        let dir = tempfile::tempdir().unwrap();
1062        let path = dir.path().join("recall.db");
1063        let st = Store::open(&path).unwrap();
1064        put(&st, "acme/app", "a.md", "x", "laptop");
1065        assert!(st.checkpoint_all().unwrap());
1066        assert_eq!(wal_len(&path), 0, "emptied");
1067        assert_eq!(rows_in_the_file_alone(&path), 1);
1068
1069        put(&st, "acme/app", "b.md", "x", "laptop");
1070        assert_eq!(rows_in_the_file_alone(&path), 1, "b.md is only in the WAL");
1071        assert!(st.checkpoint().unwrap(), "nothing held it back");
1072        assert_eq!(rows_in_the_file_alone(&path), 2);
1073
1074        put(&st, "acme/app", "c.md", "x", "laptop");
1075        assert!(st.checkpoint_all().unwrap());
1076        assert_eq!(wal_len(&path), 0, "emptied");
1077        assert_eq!(rows_in_the_file_alone(&path), 3);
1078    }
1079
1080    /// A reader in the middle of a transaction holds back what was committed
1081    /// after it began, and the sweep's checkpoint says so rather than
1082    /// claiming the file caught up; once the reader is done, it does.
1083    #[test]
1084    fn a_checkpoint_says_when_a_reader_held_it_back() {
1085        let dir = tempfile::tempdir().unwrap();
1086        let path = dir.path().join("recall.db");
1087        let st = Store::open(&path).unwrap();
1088        put(&st, "acme/app", "a.md", "x", "laptop");
1089        assert!(st.checkpoint_all().unwrap());
1090
1091        let reader = Connection::open(&path).unwrap();
1092        reader.execute_batch("BEGIN").unwrap();
1093        let _: i64 = reader
1094            .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1095            .unwrap();
1096        put(&st, "acme/app", "b.md", "x", "laptop");
1097        assert!(!st.checkpoint().unwrap(), "held back by the reader");
1098        assert_eq!(rows_in_the_file_alone(&path), 1);
1099
1100        reader.execute_batch("COMMIT").unwrap();
1101        assert!(st.checkpoint().unwrap());
1102        assert_eq!(rows_in_the_file_alone(&path), 2);
1103    }
1104
1105    /// The WAL is cut back to `WAL_SIZE_LIMIT` when it next starts over
1106    /// after something made it large, instead of staying that large for
1107    /// good.
1108    #[test]
1109    fn a_wal_grown_large_is_cut_back_to_the_limit() {
1110        let dir = tempfile::tempdir().unwrap();
1111        let path = dir.path().join("recall.db");
1112        let st = Store::open(&path).unwrap();
1113        let big = "x".repeat(24 * 1024 * 1024);
1114        put(&st, "acme/app", "big.md", &big, "laptop");
1115        let grown = wal_len(&path);
1116        assert!(grown > WAL_SIZE_LIMIT as u64, "{grown}");
1117        assert!(st.checkpoint().unwrap());
1118
1119        put(&st, "acme/app", "small.md", "x", "laptop");
1120        let now = wal_len(&path);
1121        assert!(now <= WAL_SIZE_LIMIT as u64, "{now} after {grown}");
1122        assert_eq!(st.get("acme/app", "big.md").unwrap().unwrap().content, big);
1123    }
1124
1125    /// `recall.db` moved aside under a running store, the way a restore run
1126    /// without stopping the server does it, and a new file put where it
1127    /// was: every audited write (a push, and a pull's leaf) is refused, and
1128    /// neither file is written to. Without the check, the store writes on
1129    /// into the moved file, and the push is answered 200 for a write the
1130    /// database at the path never gets.
1131    #[cfg(unix)]
1132    #[test]
1133    fn a_database_moved_aside_under_a_running_store_is_not_written_to() {
1134        let dir = tempfile::tempdir().unwrap();
1135        let path = dir.path().join("recall.db");
1136        let st = Store::open(&path).unwrap();
1137        put(&st, "acme/app", "a.md", "x", "laptop");
1138        let snapshot = st.backup(dir.path().join("backups"), 7).unwrap();
1139
1140        let aside = dir.path().join("aside");
1141        fs::create_dir(&aside).unwrap();
1142        for f in ["recall.db", "recall.db-wal", "recall.db-shm"] {
1143            if dir.path().join(f).exists() {
1144                fs::rename(dir.path().join(f), aside.join(f)).unwrap();
1145            }
1146        }
1147        // The message is a 500's body and a log line: one sentence, with no
1148        // run of spaces where a line of the source was continued.
1149        let refused = |st: &Store| {
1150            let err = st
1151                .upsert_audited("acme/app", "b.md", "y", "laptop", test_leaf)
1152                .unwrap_err();
1153            let said = format!("{err:#}");
1154            assert!(said.contains("moved or replaced"), "{said}");
1155            assert!(!said.contains("  "), "a run of spaces: {said:?}");
1156            let err = st.audit_append(test_leaf).unwrap_err();
1157            assert!(format!("{err:#}").contains("moved or replaced"), "{err:#}");
1158        };
1159        refused(&st);
1160        fs::copy(&snapshot, &path).unwrap();
1161        refused(&st);
1162        drop(st);
1163
1164        for db in [aside.join("recall.db"), path] {
1165            let files = Store::open(&db).unwrap().list("acme/app").unwrap();
1166            let names: Vec<_> = files.iter().map(|f| f.file_path.as_str()).collect();
1167            assert_eq!(names, ["a.md"], "{}", db.display());
1168        }
1169    }
1170
1171    /// Push latency under the rollback journal, WAL with NORMAL and WAL
1172    /// with FULL, through the real router against a file in a temporary
1173    /// directory. The numbers quoted beside `use_durable_wal` come from
1174    /// here:
1175    ///
1176    /// ```text
1177    /// cargo test -p recall-server --release --lib -- --ignored --nocapture push_latency
1178    /// ```
1179    #[tokio::test]
1180    #[ignore = "a measurement, not a check: run it by hand"]
1181    async fn push_latency_by_journal_mode() {
1182        use axum::body::Body;
1183        use axum::http::Request;
1184        use std::time::{Duration, Instant};
1185        use tower::ServiceExt;
1186
1187        const PUSHES: usize = 400;
1188        const TOKEN: &str = "bench-token";
1189        let content = "- a remembered fact, about as long as one usually is\n".repeat(20);
1190        println!("| journal | push p50 | push p90 | push mean | pull p50 | pull p90 |");
1191        println!("|---|---|---|---|---|---|");
1192        for (label, journal, sync) in [
1193            ("rollback (DELETE), FULL", "DELETE", "FULL"),
1194            ("WAL, NORMAL", "WAL", "NORMAL"),
1195            ("WAL, FULL", "WAL", "FULL"),
1196        ] {
1197            let dir = tempfile::tempdir().unwrap();
1198            let store = std::sync::Arc::new(Store::open(dir.path().join("recall.db")).unwrap());
1199            store
1200                .with_raw(|c| {
1201                    c.query_row(&format!("PRAGMA journal_mode = {journal}"), [], |_| Ok(()))?;
1202                    c.execute_batch(&format!("PRAGMA synchronous = {sync}"))
1203                })
1204                .unwrap();
1205            let server = crate::Server::new(
1206                crate::Config {
1207                    token: TOKEN.into(),
1208                    merge_enabled: false,
1209                    rate_limit_max: 1_000_000,
1210                    ..crate::Config::default()
1211                },
1212                store,
1213            );
1214            let router = server.router();
1215            let push = |i: usize| {
1216                let body = serde_json::json!({
1217                    "project_key": "bench/app",
1218                    "file_path": format!("f{i}.md"),
1219                    "content": content,
1220                    "source_env": "bench",
1221                });
1222                Request::post("/sync")
1223                    .header("authorization", format!("Bearer {TOKEN}"))
1224                    .header("content-type", "application/json")
1225                    .body(Body::from(body.to_string()))
1226                    .unwrap()
1227            };
1228            let pull = || {
1229                Request::get("/sync?project_key=small/app")
1230                    .header("authorization", format!("Bearer {TOKEN}"))
1231                    .body(Body::empty())
1232                    .unwrap()
1233            };
1234            let time = |mut samples: Vec<Duration>| {
1235                samples.sort();
1236                let ms = |d: Duration| d.as_secs_f64() * 1000.0;
1237                let mean = ms(samples.iter().sum::<Duration>()) / samples.len() as f64;
1238                (
1239                    ms(samples[samples.len() / 2]),
1240                    ms(samples[samples.len() * 9 / 10]),
1241                    mean,
1242                )
1243            };
1244            for i in 0..20 {
1245                let resp = router.clone().oneshot(push(PUSHES + i)).await.unwrap();
1246                assert_eq!(resp.status(), 200);
1247            }
1248            let mut pushes = Vec::with_capacity(PUSHES);
1249            for i in 0..PUSHES {
1250                let started = Instant::now();
1251                let resp = router.clone().oneshot(push(i)).await.unwrap();
1252                pushes.push(started.elapsed());
1253                assert_eq!(resp.status(), 200);
1254            }
1255            let mut pulls = Vec::with_capacity(PUSHES);
1256            for _ in 0..PUSHES {
1257                let started = Instant::now();
1258                let resp = router.clone().oneshot(pull()).await.unwrap();
1259                pulls.push(started.elapsed());
1260                assert_eq!(resp.status(), 200);
1261            }
1262            let (p50, p90, mean) = time(pushes);
1263            let (l50, l90, _) = time(pulls);
1264            println!(
1265                "| {label} | {p50:.2} ms | {p90:.2} ms | {mean:.2} ms | {l50:.2} ms | {l90:.2} ms |"
1266            );
1267        }
1268    }
1269}