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