use std::fs;
use std::path::{Path, PathBuf};
use std::sync::{Mutex, MutexGuard, PoisonError};
use std::time::Duration;
use anyhow::{Context, Result};
use recall_wire::{AdminTotals, File, ProjectStats};
use rusqlite::{Connection, OptionalExtension};
use crate::audit::merkle::Tree;
use crate::now;
mod audit;
mod devices;
mod evaluations;
mod jobs;
mod passkeys;
pub use audit::{AuditEntry, ConsistencyError, Outcome};
pub use devices::{
plain_name, Created, Decision, Inserted, NewAuthkey, NewDevice, NewEnrollment, Poll, Waiting,
};
pub use evaluations::Requested;
pub use jobs::{
clip, Failure, Queued, Retried, Settled, Settlement, MAX_ATTEMPTS, MAX_ERROR_BYTES, MAX_LINKS,
MAX_OPEN_JOBS,
};
pub use passkeys::{
AddedCredential, AdminCredential, AdminSession, BootstrapCode, FirstPasskey,
NewAdminCredential, RemovedCredential,
};
const SCHEMA: &str = "
CREATE TABLE IF NOT EXISTS memory_files (
project_key TEXT NOT NULL,
file_path TEXT NOT NULL,
content TEXT NOT NULL,
source_env TEXT,
updated_at TEXT NOT NULL,
deleted INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (project_key, file_path)
);
";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Existing {
pub content: String,
pub deleted: bool,
pub source_env: String,
pub updated_at: String,
}
struct StoreState {
conn: Connection,
audit: Tree,
audit_at: String,
file: Option<FileId>,
log_moved: bool,
moved_said: bool,
}
type FileId = (u64, u64);
#[cfg(unix)]
fn file_id(conn: &Connection) -> Option<FileId> {
use std::os::unix::fs::MetadataExt;
let path = conn.path().filter(|p| !p.is_empty())?;
fs::metadata(path).ok().map(|m| (m.dev(), m.ino()))
}
#[cfg(not(unix))]
fn file_id(_conn: &Connection) -> Option<FileId> {
None
}
impl std::ops::Deref for StoreState {
type Target = Connection;
fn deref(&self) -> &Connection {
&self.conn
}
}
impl std::ops::DerefMut for StoreState {
fn deref_mut(&mut self) -> &mut Connection {
&mut self.conn
}
}
pub struct Store {
state: Mutex<StoreState>,
}
impl Store {
pub fn open(path: impl AsRef<Path>) -> Result<Self> {
Self::open_waiting(path.as_ref(), admin::BUSY_TIMEOUT)
}
fn open_waiting(path: &Path, busy: Duration) -> Result<Self> {
if let Some(dir) = path.parent() {
if !dir.as_os_str().is_empty() {
fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
}
}
let conn = Connection::open(path).with_context(|| format!("opening {}", path.display()))?;
use_durable_wal(&conn, busy).with_context(|| {
format!(
"switching {} to SQLite's WAL journal. It needs a local filesystem, and a \
moment with no other process holding the file (sqlite-web mid-read, an admin \
command): start the server again",
path.display()
)
})?;
Self::with_connection(conn)
}
pub fn open_in_memory() -> Result<Self> {
Self::with_connection(Connection::open_in_memory()?)
}
fn with_connection(conn: Connection) -> Result<Self> {
let file = file_id(&conn);
let store = Self {
state: Mutex::new(StoreState {
conn,
audit: Tree::new(),
audit_at: String::new(),
file,
log_moved: true,
moved_said: false,
}),
};
store.migrate()?;
Ok(store)
}
fn lock(&self) -> MutexGuard<'_, StoreState> {
self.state.lock().unwrap_or_else(PoisonError::into_inner)
}
fn migrate(&self) -> Result<()> {
let mut state = self.lock();
state.conn.execute_batch(SCHEMA)?;
let has_deleted = {
let mut stmt = state.conn.prepare("PRAGMA table_info(memory_files)")?;
let mut rows = stmt.query([])?;
let mut found = false;
while let Some(row) = rows.next()? {
if row.get::<_, String>(1)? == "deleted" {
found = true;
}
}
found
};
if !has_deleted {
state.conn.execute(
"ALTER TABLE memory_files ADD COLUMN deleted INTEGER NOT NULL DEFAULT 0",
[],
)?;
}
state.conn.execute_batch(devices::SCHEMA)?;
state.conn.execute_batch(passkeys::SCHEMA)?;
devices::allow_worker_scope(&state.conn)?;
state.conn.execute_batch(jobs::SCHEMA)?;
state.conn.execute_batch(evaluations::SCHEMA)?;
state.conn.execute_batch(audit::SCHEMA)?;
let loaded = audit::load(&state.conn)?;
state.audit = loaded.tree;
state.audit_at = loaded.last_at;
Ok(())
}
pub fn get(&self, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
read_file(&self.lock(), project_key, file_path)
}
pub fn upsert_audited(
&self,
project_key: &str,
file_path: &str,
content: &str,
source_env: &str,
build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
) -> Result<String> {
self.audited(
|tx, at| {
write_file(tx, project_key, file_path, content, source_env, at)?;
Ok(Outcome::Commit(at.to_string()))
},
|seq, at, _| build_leaf(seq, at),
)
}
pub fn tombstone_audited(
&self,
project_key: &str,
file_path: &str,
source_env: &str,
build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
) -> Result<String> {
self.audited(
|tx, at| {
tx.execute(
"INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
VALUES (?1, ?2, '', ?3, ?4, 1)
ON CONFLICT(project_key, file_path) DO UPDATE SET
source_env = excluded.source_env,
updated_at = excluded.updated_at,
deleted = 1",
(project_key, file_path, nullable(source_env), at),
)?;
jobs::close_for_delete(tx, project_key, file_path, at)?;
Ok(Outcome::Commit(at.to_string()))
},
|seq, at, _| build_leaf(seq, at),
)
}
pub fn list(&self, project_key: &str) -> Result<Vec<File>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT file_path, content, COALESCE(source_env, ''), updated_at, deleted
FROM memory_files WHERE project_key = ?1 ORDER BY file_path",
)?;
let rows = stmt.query_map((project_key,), |r| {
let content: String = r.get(1)?;
let deleted = r.get::<_, i64>(4)? != 0;
Ok(File {
file_path: r.get(0)?,
content: if deleted { None } else { Some(content) },
source_env: r.get(2)?,
updated_at: r.get(3)?,
deleted,
})
})?;
let mut files = Vec::new();
for row in rows {
files.push(row?);
}
Ok(files)
}
pub fn last_sync_at(&self) -> Result<String> {
let conn = self.lock();
let v: Option<String> =
conn.query_row("SELECT MAX(updated_at) FROM memory_files", [], |r| r.get(0))?;
Ok(v.unwrap_or_default())
}
pub fn admin_stats(&self) -> Result<(Vec<ProjectStats>, AdminTotals)> {
let conn = self.lock();
let mut projects = Vec::new();
let mut totals = AdminTotals::default();
{
let mut stmt = conn.prepare(
"SELECT project_key,
SUM(CASE WHEN deleted = 0 THEN 1 ELSE 0 END),
SUM(CASE WHEN deleted = 1 THEN 1 ELSE 0 END),
MAX(updated_at)
FROM memory_files GROUP BY project_key ORDER BY MAX(updated_at) DESC",
)?;
let rows = stmt.query_map([], |r| {
Ok(ProjectStats {
project_key: r.get(0)?,
file_count: r.get(1)?,
deleted_count: r.get(2)?,
sources: Vec::new(),
last_updated_at: r.get::<_, Option<String>>(3)?.unwrap_or_default(),
})
})?;
for row in rows {
let p = row?;
totals.file_count += p.file_count;
totals.deleted_count += p.deleted_count;
projects.push(p);
}
}
totals.project_count = projects.len() as i64;
{
let mut stmt = conn.prepare(
"SELECT DISTINCT project_key, source_env FROM memory_files WHERE source_env IS NOT NULL",
)?;
let rows =
stmt.query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?)))?;
for row in rows {
let (key, src) = row?;
if src.is_empty() {
continue;
}
if let Some(p) = projects.iter_mut().find(|p| p.project_key == key) {
p.sources.push(src);
}
}
}
for p in &mut projects {
p.sources.sort();
}
Ok((projects, totals))
}
pub fn checkpoint(&self) -> Result<bool> {
let (frames, copied): (i64, i64) =
self.lock()
.query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |r| {
Ok((r.get(1)?, r.get(2)?))
})?;
Ok(copied >= frames)
}
pub fn checkpoint_all(&self) -> Result<bool> {
let busy: i64 = self
.lock()
.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |r| r.get(0))?;
Ok(busy == 0)
}
pub fn backup(&self, dir: impl AsRef<Path>, keep: usize) -> Result<PathBuf> {
let dir = dir.as_ref();
fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
let stamp = now().replace([':', '.'], "-");
let dest = dir.join(format!("recall-{stamp}.db"));
let dest_str = dest
.to_str()
.context("backup path is not valid UTF-8")?
.to_owned();
let existed = dest.exists();
let vacuumed = {
let conn = self.lock();
conn.execute("VACUUM INTO ?1", (&dest_str,))
.with_context(|| format!("VACUUM INTO {dest_str}"))
};
if let Err(err) = vacuumed {
if !existed {
let _ = fs::remove_file(&dest);
}
return Err(err);
}
let mut snapshots: Vec<PathBuf> = fs::read_dir(dir)?
.filter_map(|e| e.ok())
.map(|e| e.path())
.filter(|p| {
p.file_name()
.and_then(|n| n.to_str())
.is_some_and(|n| n.starts_with("recall-") && n.ends_with(".db"))
})
.collect();
snapshots.sort();
for stale in snapshots.iter().take(snapshots.len().saturating_sub(keep)) {
let _ = fs::remove_file(stale);
}
Ok(dest)
}
}
#[cfg(test)]
impl Store {
pub(crate) fn with_raw<T>(
&self,
f: impl FnOnce(&Connection) -> rusqlite::Result<T>,
) -> rusqlite::Result<T> {
f(&self.lock())
}
}
#[cfg(test)]
pub(crate) fn test_leaf(seq: u64, at: &str) -> Vec<u8> {
use crate::audit::leaf;
leaf::encode(
seq,
at,
leaf::action::START,
&leaf::Actor::Server,
leaf::subject_start("test"),
None,
)
}
fn use_durable_wal(conn: &Connection, busy: Duration) -> Result<()> {
conn.busy_timeout(busy)?;
let mode: String = conn.query_row("PRAGMA journal_mode = WAL", [], |r| r.get(0))?;
if !mode.eq_ignore_ascii_case("wal") {
anyhow::bail!("SQLite kept the {mode} journal");
}
conn.execute_batch(&format!(
"PRAGMA synchronous = FULL; PRAGMA journal_size_limit = {WAL_SIZE_LIMIT}"
))?;
Ok(())
}
const WAL_SIZE_LIMIT: i64 = 16 * 1024 * 1024;
const EXISTING_COLUMNS: &str = "content, deleted, COALESCE(source_env, ''), updated_at";
fn existing_from(r: &rusqlite::Row<'_>) -> rusqlite::Result<Existing> {
Ok(Existing {
content: r.get(0)?,
deleted: r.get::<_, i64>(1)? != 0,
source_env: r.get(2)?,
updated_at: r.get(3)?,
})
}
fn read_file(conn: &Connection, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
Ok(conn
.query_row(
&format!("SELECT {EXISTING_COLUMNS} FROM memory_files WHERE project_key = ?1 AND file_path = ?2"),
(project_key, file_path),
existing_from,
)
.optional()?)
}
fn write_file(
conn: &Connection,
project_key: &str,
file_path: &str,
content: &str,
source_env: &str,
updated_at: &str,
) -> Result<()> {
conn.execute(
"INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
VALUES (?1, ?2, ?3, ?4, ?5, 0)
ON CONFLICT(project_key, file_path) DO UPDATE SET
content = excluded.content,
source_env = excluded.source_env,
updated_at = excluded.updated_at,
deleted = 0",
(project_key, file_path, content, nullable(source_env), updated_at),
)?;
Ok(())
}
fn nullable(s: &str) -> Option<&str> {
if s.is_empty() {
None
} else {
Some(s)
}
}
pub(crate) mod admin;
#[cfg(test)]
mod tests {
use super::*;
fn store() -> Store {
Store::open_in_memory().unwrap()
}
fn put(st: &Store, project_key: &str, file_path: &str, content: &str, source_env: &str) {
st.upsert_audited(project_key, file_path, content, source_env, test_leaf)
.unwrap();
}
fn del(st: &Store, project_key: &str, file_path: &str, source_env: &str) {
st.tombstone_audited(project_key, file_path, source_env, test_leaf)
.unwrap();
}
#[test]
fn a_write_is_stamped_with_its_leafs_at() {
let st = store();
let mut leaf_at = String::new();
let updated_at = st
.upsert_audited("acme/app", "a.md", "x", "laptop", |seq, at| {
leaf_at = at.to_string();
test_leaf(seq, at)
})
.unwrap();
assert_eq!(updated_at, leaf_at);
assert_eq!(st.list("acme/app").unwrap()[0].updated_at, updated_at);
assert_eq!(st.audit_checkpoint().0, 1);
}
#[test]
fn upsert_get_and_list_round_trip() {
let st = store();
put(&st, "acme/app", "MEMORY.md", "hello", "laptop");
let got = st.get("acme/app", "MEMORY.md").unwrap().unwrap();
assert_eq!(got.content, "hello");
assert!(!got.deleted);
let files = st.list("acme/app").unwrap();
assert_eq!(files.len(), 1);
assert_eq!(files[0].content.as_deref(), Some("hello"));
assert_eq!(files[0].source_env, "laptop");
assert!(st.get("acme/app", "missing.md").unwrap().is_none());
}
#[test]
fn tombstone_preserves_content_but_list_withholds_it() {
let st = store();
put(&st, "acme/app", "gone.md", "secret", "laptop");
del(&st, "acme/app", "gone.md", "laptop");
let row = st.get("acme/app", "gone.md").unwrap().unwrap();
assert_eq!(row.content, "secret", "content must stay recoverable");
assert!(row.deleted);
let files = st.list("acme/app").unwrap();
assert_eq!(
files.len(),
1,
"tombstones are listed so clients can delete locally"
);
assert!(files[0].deleted);
assert_eq!(files[0].content, None, "a pull must not resurrect it");
}
#[test]
fn upsert_clears_a_tombstone() {
let st = store();
del(&st, "acme/app", "f.md", "laptop");
put(&st, "acme/app", "f.md", "back", "laptop");
let row = st.get("acme/app", "f.md").unwrap().unwrap();
assert!(!row.deleted);
assert_eq!(row.content, "back");
}
#[test]
fn last_sync_at_is_empty_on_a_fresh_database() {
assert_eq!(store().last_sync_at().unwrap(), "");
}
#[test]
fn admin_stats_keeps_commas_inside_a_source_env() {
let st = store();
put(&st, "acme/app", "a.md", "x", "laptop,evil");
let (projects, _) = st.admin_stats().unwrap();
assert_eq!(projects[0].sources, vec!["laptop,evil".to_string()]);
}
#[test]
fn migrates_a_database_that_predates_tombstones() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("old.db");
{
let conn = Connection::open(&path).unwrap();
conn.execute_batch(
"CREATE TABLE memory_files (
project_key TEXT NOT NULL,
file_path TEXT NOT NULL,
content TEXT NOT NULL,
source_env TEXT,
updated_at TEXT NOT NULL,
PRIMARY KEY (project_key, file_path)
);
INSERT INTO memory_files VALUES ('acme/app','old.md','kept','node-era','2026-09-03T21:49:55.191Z');",
)
.unwrap();
}
let st = Store::open(&path).unwrap();
let files = st.list("acme/app").unwrap();
assert_eq!(files.len(), 1);
assert_eq!(files[0].content.as_deref(), Some("kept"));
assert!(!files[0].deleted);
}
#[test]
fn backup_names_carry_milliseconds() {
let dir = tempfile::tempdir().unwrap();
let st = store();
let dest = st.backup(dir.path(), 7).unwrap();
let name = dest.file_name().unwrap().to_str().unwrap();
assert!(
name.starts_with("recall-") && name.ends_with("Z.db"),
"got {name}"
);
let stamp = &name["recall-".len()..name.len() - ".db".len()];
assert_eq!(stamp.len(), 24, "got {stamp}");
assert!(
stamp[20..23].chars().all(|c| c.is_ascii_digit()),
"no millisecond field in {stamp}"
);
}
fn journal_mode(conn: &Connection) -> String {
conn.query_row("PRAGMA journal_mode", [], |r| r.get(0))
.unwrap()
}
fn header_mode(path: &Path) -> (u8, u8) {
let head = fs::read(path).unwrap();
(head[18], head[19])
}
fn wal_len(db: &Path) -> u64 {
let mut wal = db.as_os_str().to_owned();
wal.push("-wal");
fs::metadata(wal).map(|m| m.len()).unwrap_or(0)
}
#[test]
fn the_store_keeps_the_file_in_wal_and_syncs_every_commit() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("recall.db");
let st = Store::open(&path).unwrap();
put(&st, "acme/app", "a.md", "x", "laptop");
let (mode, sync, busy) = st
.with_raw(|c| {
Ok((
journal_mode(c),
c.query_row("PRAGMA synchronous", [], |r| r.get::<_, i64>(0))?,
c.query_row("PRAGMA busy_timeout", [], |r| r.get::<_, i64>(0))?,
))
})
.unwrap();
assert_eq!(mode, "wal");
assert_eq!(sync, 2, "synchronous=FULL");
assert_eq!(busy, admin::BUSY_TIMEOUT.as_millis() as i64);
let limit: i64 = st
.with_raw(|c| c.query_row("PRAGMA journal_size_limit", [], |r| r.get(0)))
.unwrap();
assert_eq!(limit, WAL_SIZE_LIMIT);
assert_eq!(header_mode(&path), (2, 2));
assert_eq!(journal_mode(&Connection::open(&path).unwrap()), "wal");
assert!(wal_len(&path) > 0, "the commit went to the WAL");
}
#[test]
fn the_sqlite_compiled_in_has_the_wal_restart_fix() {
assert!(
rusqlite::version_number() >= 3_051_003,
"SQLite {} predates 3.51.3",
rusqlite::version()
);
}
#[test]
fn a_rollback_journal_database_from_an_older_server_is_converted_on_open() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("recall.db");
let before = {
let st = Store::open(&path).unwrap();
put(&st, "acme/app", "MEMORY.md", "kept\n", "laptop");
del(&st, "acme/app", "gone.md", "laptop");
st.audit_checkpoint()
};
Connection::open(&path)
.unwrap()
.query_row("PRAGMA journal_mode = DELETE", [], |_| Ok(()))
.unwrap();
assert_eq!(header_mode(&path), (1, 1));
let st = Store::open(&path).unwrap();
assert_eq!(header_mode(&path), (2, 2));
let files = st.list("acme/app").unwrap();
assert_eq!(files.len(), 2);
assert_eq!(files[0].content.as_deref(), Some("kept\n"));
assert!(files[1].deleted);
assert_eq!(st.audit_checkpoint(), before);
put(&st, "acme/app", "after.md", "y", "laptop");
drop(st);
let plain = Connection::open(&path).unwrap();
let n: i64 = plain
.query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
.unwrap();
assert_eq!(n, 3);
}
#[test]
fn a_switch_held_up_past_the_busy_timeout_refuses_to_start_and_changes_nothing() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("recall.db");
drop(Store::open(&path).unwrap());
Connection::open(&path)
.unwrap()
.query_row("PRAGMA journal_mode = DELETE", [], |_| Ok(()))
.unwrap();
let reader = Connection::open(&path).unwrap();
reader.execute_batch("BEGIN").unwrap();
let _: i64 = reader
.query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
.unwrap();
let err = match Store::open_waiting(&path, Duration::from_millis(250)) {
Ok(_) => panic!("switched to WAL under a reader holding the file"),
Err(e) => format!("{e:#}"),
};
assert!(err.contains("switching"), "{err}");
assert!(err.contains("start the server again"), "{err}");
assert!(err.contains("locked"), "{err}");
assert_eq!(header_mode(&path), (1, 1), "still the rollback journal");
reader.execute_batch("COMMIT").unwrap();
drop(reader);
drop(Store::open(&path).unwrap());
assert_eq!(header_mode(&path), (2, 2));
}
#[test]
fn the_node_servers_database_is_converted_with_every_row() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("recall.db");
fs::copy(
concat!(
env!("CARGO_MANIFEST_DIR"),
"/../../fixtures/node-written.db"
),
&path,
)
.unwrap();
assert_eq!(header_mode(&path), (1, 1));
let dump = || -> Vec<(String, String, String, Option<String>, String, i64)> {
let conn = Connection::open(&path).unwrap();
let mut stmt = conn
.prepare(
"SELECT project_key, file_path, content, source_env, updated_at, deleted
FROM memory_files ORDER BY project_key, file_path",
)
.unwrap();
let rows = stmt
.query_map([], |r| {
Ok((
r.get(0)?,
r.get(1)?,
r.get(2)?,
r.get(3)?,
r.get(4)?,
r.get(5)?,
))
})
.unwrap();
rows.map(Result::unwrap).collect()
};
let before = dump();
assert!(!before.is_empty());
let st = Store::open(&path).unwrap();
assert_eq!(header_mode(&path), (2, 2));
assert_eq!(dump(), before);
drop(st);
assert_eq!(dump(), before);
}
#[test]
fn a_snapshot_is_one_self_contained_file_with_what_the_wal_holds() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("recall.db");
let st = Store::open(&path).unwrap();
for i in 0..20 {
put(&st, "acme/app", &format!("f{i}.md"), "x", "laptop");
}
assert!(wal_len(&path) > 0, "the pushes are still in the WAL");
let snap = st.backup(dir.path().join("backups"), 7).unwrap();
let names: Vec<_> = fs::read_dir(snap.parent().unwrap())
.unwrap()
.map(|e| e.unwrap().file_name().into_string().unwrap())
.collect();
assert_eq!(names.len(), 1, "no -wal or -shm beside it: {names:?}");
assert_eq!(header_mode(&snap), (1, 1));
let read =
Connection::open_with_flags(&snap, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY).unwrap();
assert_eq!(journal_mode(&read), "delete");
let n: i64 = read
.query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
.unwrap();
assert_eq!(n, 20);
}
fn rows_in_the_file_alone(db: &Path) -> i64 {
let alone = tempfile::tempdir().unwrap();
let copy = alone.path().join("recall.db");
fs::copy(db, ©).unwrap();
Connection::open(©)
.unwrap()
.query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
.unwrap()
}
#[test]
fn a_checkpoint_brings_the_file_on_its_own_up_to_date() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("recall.db");
let st = Store::open(&path).unwrap();
put(&st, "acme/app", "a.md", "x", "laptop");
assert!(st.checkpoint_all().unwrap());
assert_eq!(wal_len(&path), 0, "emptied");
assert_eq!(rows_in_the_file_alone(&path), 1);
put(&st, "acme/app", "b.md", "x", "laptop");
assert_eq!(rows_in_the_file_alone(&path), 1, "b.md is only in the WAL");
assert!(st.checkpoint().unwrap(), "nothing held it back");
assert_eq!(rows_in_the_file_alone(&path), 2);
put(&st, "acme/app", "c.md", "x", "laptop");
assert!(st.checkpoint_all().unwrap());
assert_eq!(wal_len(&path), 0, "emptied");
assert_eq!(rows_in_the_file_alone(&path), 3);
}
#[test]
fn a_checkpoint_says_when_a_reader_held_it_back() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("recall.db");
let st = Store::open(&path).unwrap();
put(&st, "acme/app", "a.md", "x", "laptop");
assert!(st.checkpoint_all().unwrap());
let reader = Connection::open(&path).unwrap();
reader.execute_batch("BEGIN").unwrap();
let _: i64 = reader
.query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
.unwrap();
put(&st, "acme/app", "b.md", "x", "laptop");
assert!(!st.checkpoint().unwrap(), "held back by the reader");
assert_eq!(rows_in_the_file_alone(&path), 1);
reader.execute_batch("COMMIT").unwrap();
assert!(st.checkpoint().unwrap());
assert_eq!(rows_in_the_file_alone(&path), 2);
}
#[test]
fn a_wal_grown_large_is_cut_back_to_the_limit() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("recall.db");
let st = Store::open(&path).unwrap();
let big = "x".repeat(24 * 1024 * 1024);
put(&st, "acme/app", "big.md", &big, "laptop");
let grown = wal_len(&path);
assert!(grown > WAL_SIZE_LIMIT as u64, "{grown}");
assert!(st.checkpoint().unwrap());
put(&st, "acme/app", "small.md", "x", "laptop");
let now = wal_len(&path);
assert!(now <= WAL_SIZE_LIMIT as u64, "{now} after {grown}");
assert_eq!(st.get("acme/app", "big.md").unwrap().unwrap().content, big);
}
#[cfg(unix)]
#[test]
fn a_database_moved_aside_under_a_running_store_is_not_written_to() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("recall.db");
let st = Store::open(&path).unwrap();
put(&st, "acme/app", "a.md", "x", "laptop");
let snapshot = st.backup(dir.path().join("backups"), 7).unwrap();
let aside = dir.path().join("aside");
fs::create_dir(&aside).unwrap();
for f in ["recall.db", "recall.db-wal", "recall.db-shm"] {
if dir.path().join(f).exists() {
fs::rename(dir.path().join(f), aside.join(f)).unwrap();
}
}
let refused = |st: &Store| {
let err = st
.upsert_audited("acme/app", "b.md", "y", "laptop", test_leaf)
.unwrap_err();
let said = format!("{err:#}");
assert!(said.contains("moved or replaced"), "{said}");
assert!(!said.contains(" "), "a run of spaces: {said:?}");
let err = st.audit_append(test_leaf).unwrap_err();
assert!(format!("{err:#}").contains("moved or replaced"), "{err:#}");
};
refused(&st);
fs::copy(&snapshot, &path).unwrap();
refused(&st);
drop(st);
for db in [aside.join("recall.db"), path] {
let files = Store::open(&db).unwrap().list("acme/app").unwrap();
let names: Vec<_> = files.iter().map(|f| f.file_path.as_str()).collect();
assert_eq!(names, ["a.md"], "{}", db.display());
}
}
#[tokio::test]
#[ignore = "a measurement, not a check: run it by hand"]
async fn push_latency_by_journal_mode() {
use axum::body::Body;
use axum::http::Request;
use std::time::{Duration, Instant};
use tower::ServiceExt;
const PUSHES: usize = 400;
const TOKEN: &str = "bench-token";
let content = "- a remembered fact, about as long as one usually is\n".repeat(20);
println!("| journal | push p50 | push p90 | push mean | pull p50 | pull p90 |");
println!("|---|---|---|---|---|---|");
for (label, journal, sync) in [
("rollback (DELETE), FULL", "DELETE", "FULL"),
("WAL, NORMAL", "WAL", "NORMAL"),
("WAL, FULL", "WAL", "FULL"),
] {
let dir = tempfile::tempdir().unwrap();
let store = std::sync::Arc::new(Store::open(dir.path().join("recall.db")).unwrap());
store
.with_raw(|c| {
c.query_row(&format!("PRAGMA journal_mode = {journal}"), [], |_| Ok(()))?;
c.execute_batch(&format!("PRAGMA synchronous = {sync}"))
})
.unwrap();
let server = crate::Server::new(
crate::Config {
token: TOKEN.into(),
merge_enabled: false,
rate_limit_max: 1_000_000,
..crate::Config::default()
},
store,
);
let router = server.router();
let push = |i: usize| {
let body = serde_json::json!({
"project_key": "bench/app",
"file_path": format!("f{i}.md"),
"content": content,
"source_env": "bench",
});
Request::post("/sync")
.header("authorization", format!("Bearer {TOKEN}"))
.header("content-type", "application/json")
.body(Body::from(body.to_string()))
.unwrap()
};
let pull = || {
Request::get("/sync?project_key=small/app")
.header("authorization", format!("Bearer {TOKEN}"))
.body(Body::empty())
.unwrap()
};
let time = |mut samples: Vec<Duration>| {
samples.sort();
let ms = |d: Duration| d.as_secs_f64() * 1000.0;
let mean = ms(samples.iter().sum::<Duration>()) / samples.len() as f64;
(
ms(samples[samples.len() / 2]),
ms(samples[samples.len() * 9 / 10]),
mean,
)
};
for i in 0..20 {
let resp = router.clone().oneshot(push(PUSHES + i)).await.unwrap();
assert_eq!(resp.status(), 200);
}
let mut pushes = Vec::with_capacity(PUSHES);
for i in 0..PUSHES {
let started = Instant::now();
let resp = router.clone().oneshot(push(i)).await.unwrap();
pushes.push(started.elapsed());
assert_eq!(resp.status(), 200);
}
let mut pulls = Vec::with_capacity(PUSHES);
for _ in 0..PUSHES {
let started = Instant::now();
let resp = router.clone().oneshot(pull()).await.unwrap();
pulls.push(started.elapsed());
assert_eq!(resp.status(), 200);
}
let (p50, p90, mean) = time(pushes);
let (l50, l90, _) = time(pulls);
println!(
"| {label} | {p50:.2} ms | {p90:.2} ms | {mean:.2} ms | {l50:.2} ms | {l90:.2} ms |"
);
}
}
}