1use std::fs;
10use std::path::{Path, PathBuf};
11use std::sync::{Mutex, MutexGuard, PoisonError};
12
13use anyhow::{Context, Result};
14use recall_wire::{AdminTotals, File, ProjectStats};
15use rusqlite::{Connection, OptionalExtension};
16
17use crate::audit::merkle::Tree;
18use crate::now;
19
20mod audit;
21mod devices;
22mod jobs;
23mod passkeys;
24
25pub use audit::{AuditEntry, ConsistencyError, Outcome};
26pub use devices::{
27 plain_name, Created, Decision, Inserted, NewAuthkey, NewDevice, NewEnrollment, Poll, Waiting,
28};
29pub use jobs::{
30 clip, Failure, Queued, Retried, Settled, Settlement, MAX_ATTEMPTS, MAX_ERROR_BYTES, MAX_LINKS,
31 MAX_OPEN_JOBS,
32};
33pub use passkeys::{
34 AddedCredential, AdminCredential, AdminSession, BootstrapCode, FirstPasskey,
35 NewAdminCredential, RemovedCredential,
36};
37
38const SCHEMA: &str = "
40 CREATE TABLE IF NOT EXISTS memory_files (
41 project_key TEXT NOT NULL,
42 file_path TEXT NOT NULL,
43 content TEXT NOT NULL,
44 source_env TEXT,
45 updated_at TEXT NOT NULL,
46 deleted INTEGER NOT NULL DEFAULT 0,
47 PRIMARY KEY (project_key, file_path)
48 );
49";
50
51#[derive(Debug, Clone, PartialEq, Eq)]
54pub struct Existing {
55 pub content: String,
58 pub deleted: bool,
60 pub source_env: String,
62 pub updated_at: String,
64}
65
66struct StoreState {
83 conn: Connection,
84 audit: Tree,
85 audit_at: String,
86 file: Option<FileId>,
87 log_moved: bool,
88 moved_said: bool,
89}
90
91type FileId = (u64, u64);
110
111#[cfg(unix)]
115fn file_id(conn: &Connection) -> Option<FileId> {
116 use std::os::unix::fs::MetadataExt;
117 let path = conn.path().filter(|p| !p.is_empty())?;
118 fs::metadata(path).ok().map(|m| (m.dev(), m.ino()))
119}
120
121#[cfg(not(unix))]
122fn file_id(_conn: &Connection) -> Option<FileId> {
123 None
124}
125
126impl std::ops::Deref for StoreState {
127 type Target = Connection;
128 fn deref(&self) -> &Connection {
129 &self.conn
130 }
131}
132
133impl std::ops::DerefMut for StoreState {
134 fn deref_mut(&mut self) -> &mut Connection {
135 &mut self.conn
136 }
137}
138
139pub struct Store {
141 state: Mutex<StoreState>,
148}
149
150impl Store {
151 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
161 let path = path.as_ref();
162 if let Some(dir) = path.parent() {
163 if !dir.as_os_str().is_empty() {
164 fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
165 }
166 }
167 let conn = Connection::open(path).with_context(|| format!("opening {}", path.display()))?;
168 use_durable_wal(&conn).with_context(|| {
169 format!(
170 "switching {} to SQLite's WAL journal. It needs a local filesystem, and a \
171 moment with no other process holding the file (sqlite-web mid-read, an admin \
172 command): start the server again",
173 path.display()
174 )
175 })?;
176 Self::with_connection(conn)
177 }
178
179 pub fn open_in_memory() -> Result<Self> {
181 Self::with_connection(Connection::open_in_memory()?)
182 }
183
184 fn with_connection(conn: Connection) -> Result<Self> {
185 let file = file_id(&conn);
186 let store = Self {
187 state: Mutex::new(StoreState {
188 conn,
189 audit: Tree::new(),
190 audit_at: String::new(),
191 file,
192 log_moved: true,
193 moved_said: false,
194 }),
195 };
196 store.migrate()?;
197 Ok(store)
198 }
199
200 fn lock(&self) -> MutexGuard<'_, StoreState> {
204 self.state.lock().unwrap_or_else(PoisonError::into_inner)
205 }
206
207 fn migrate(&self) -> Result<()> {
208 let mut state = self.lock();
209 state.conn.execute_batch(SCHEMA)?;
210
211 let has_deleted = {
214 let mut stmt = state.conn.prepare("PRAGMA table_info(memory_files)")?;
215 let mut rows = stmt.query([])?;
216 let mut found = false;
217 while let Some(row) = rows.next()? {
218 if row.get::<_, String>(1)? == "deleted" {
219 found = true;
220 }
221 }
222 found
223 };
224 if !has_deleted {
225 state.conn.execute(
226 "ALTER TABLE memory_files ADD COLUMN deleted INTEGER NOT NULL DEFAULT 0",
227 [],
228 )?;
229 }
230
231 state.conn.execute_batch(devices::SCHEMA)?;
236 state.conn.execute_batch(passkeys::SCHEMA)?;
237 devices::allow_worker_scope(&state.conn)?;
242 state.conn.execute_batch(jobs::SCHEMA)?;
243 state.conn.execute_batch(audit::SCHEMA)?;
244
245 let loaded = audit::load(&state.conn)?;
246 state.audit = loaded.tree;
247 state.audit_at = loaded.last_at;
248 Ok(())
249 }
250
251 pub fn get(&self, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
253 read_file(&self.lock(), project_key, file_path)
254 }
255
256 pub fn upsert_audited(
261 &self,
262 project_key: &str,
263 file_path: &str,
264 content: &str,
265 source_env: &str,
266 build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
267 ) -> Result<String> {
268 self.audited(
269 |tx, at| {
270 write_file(tx, project_key, file_path, content, source_env, at)?;
271 Ok(Outcome::Commit(at.to_string()))
272 },
273 |seq, at, _| build_leaf(seq, at),
274 )
275 }
276
277 pub fn tombstone_audited(
287 &self,
288 project_key: &str,
289 file_path: &str,
290 source_env: &str,
291 build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
292 ) -> Result<String> {
293 self.audited(
294 |tx, at| {
295 tx.execute(
296 "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
297 VALUES (?1, ?2, '', ?3, ?4, 1)
298 ON CONFLICT(project_key, file_path) DO UPDATE SET
299 source_env = excluded.source_env,
300 updated_at = excluded.updated_at,
301 deleted = 1",
302 (project_key, file_path, nullable(source_env), at),
303 )?;
304 jobs::close_for_delete(tx, project_key, file_path, at)?;
305 Ok(Outcome::Commit(at.to_string()))
306 },
307 |seq, at, _| build_leaf(seq, at),
308 )
309 }
310
311 pub fn list(&self, project_key: &str) -> Result<Vec<File>> {
315 let conn = self.lock();
316 let mut stmt = conn.prepare(
317 "SELECT file_path, content, COALESCE(source_env, ''), updated_at, deleted
318 FROM memory_files WHERE project_key = ?1 ORDER BY file_path",
319 )?;
320 let rows = stmt.query_map((project_key,), |r| {
321 let content: String = r.get(1)?;
322 let deleted = r.get::<_, i64>(4)? != 0;
323 Ok(File {
324 file_path: r.get(0)?,
325 content: if deleted { None } else { Some(content) },
326 source_env: r.get(2)?,
327 updated_at: r.get(3)?,
328 deleted,
329 })
330 })?;
331 let mut files = Vec::new();
332 for row in rows {
333 files.push(row?);
334 }
335 Ok(files)
336 }
337
338 pub fn last_sync_at(&self) -> Result<String> {
341 let conn = self.lock();
342 let v: Option<String> =
343 conn.query_row("SELECT MAX(updated_at) FROM memory_files", [], |r| r.get(0))?;
344 Ok(v.unwrap_or_default())
345 }
346
347 pub fn admin_stats(&self) -> Result<(Vec<ProjectStats>, AdminTotals)> {
349 let conn = self.lock();
350
351 let mut projects = Vec::new();
352 let mut totals = AdminTotals::default();
353 {
354 let mut stmt = conn.prepare(
355 "SELECT project_key,
356 SUM(CASE WHEN deleted = 0 THEN 1 ELSE 0 END),
357 SUM(CASE WHEN deleted = 1 THEN 1 ELSE 0 END),
358 MAX(updated_at)
359 FROM memory_files GROUP BY project_key ORDER BY MAX(updated_at) DESC",
360 )?;
361 let rows = stmt.query_map([], |r| {
362 Ok(ProjectStats {
363 project_key: r.get(0)?,
364 file_count: r.get(1)?,
365 deleted_count: r.get(2)?,
366 sources: Vec::new(),
367 last_updated_at: r.get::<_, Option<String>>(3)?.unwrap_or_default(),
368 })
369 })?;
370 for row in rows {
371 let p = row?;
372 totals.file_count += p.file_count;
373 totals.deleted_count += p.deleted_count;
374 projects.push(p);
375 }
376 }
377 totals.project_count = projects.len() as i64;
378
379 {
384 let mut stmt = conn.prepare(
385 "SELECT DISTINCT project_key, source_env FROM memory_files WHERE source_env IS NOT NULL",
386 )?;
387 let rows =
388 stmt.query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?)))?;
389 for row in rows {
390 let (key, src) = row?;
391 if src.is_empty() {
392 continue;
393 }
394 if let Some(p) = projects.iter_mut().find(|p| p.project_key == key) {
395 p.sources.push(src);
396 }
397 }
398 }
399 for p in &mut projects {
400 p.sources.sort();
401 }
402 Ok((projects, totals))
403 }
404
405 pub fn checkpoint(&self) -> Result<bool> {
419 let (frames, copied): (i64, i64) =
420 self.lock()
421 .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |r| {
422 Ok((r.get(1)?, r.get(2)?))
423 })?;
424 Ok(copied >= frames)
425 }
426
427 pub fn checkpoint_all(&self) -> Result<bool> {
438 let busy: i64 = self
439 .lock()
440 .query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |r| r.get(0))?;
441 Ok(busy == 0)
442 }
443
444 pub fn backup(&self, dir: impl AsRef<Path>, keep: usize) -> Result<PathBuf> {
455 let dir = dir.as_ref();
456 fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
457
458 let stamp = now().replace([':', '.'], "-");
465 let dest = dir.join(format!("recall-{stamp}.db"));
466 let dest_str = dest
467 .to_str()
468 .context("backup path is not valid UTF-8")?
469 .to_owned();
470
471 let existed = dest.exists();
474 let vacuumed = {
475 let conn = self.lock();
476 conn.execute("VACUUM INTO ?1", (&dest_str,))
477 .with_context(|| format!("VACUUM INTO {dest_str}"))
478 };
479 if let Err(err) = vacuumed {
480 if !existed {
486 let _ = fs::remove_file(&dest);
487 }
488 return Err(err);
489 }
490
491 let mut snapshots: Vec<PathBuf> = fs::read_dir(dir)?
492 .filter_map(|e| e.ok())
493 .map(|e| e.path())
494 .filter(|p| {
495 p.file_name()
496 .and_then(|n| n.to_str())
497 .is_some_and(|n| n.starts_with("recall-") && n.ends_with(".db"))
498 })
499 .collect();
500 snapshots.sort();
501 for stale in snapshots.iter().take(snapshots.len().saturating_sub(keep)) {
502 let _ = fs::remove_file(stale);
503 }
504 Ok(dest)
505 }
506}
507
508#[cfg(test)]
509impl Store {
510 pub(crate) fn with_raw<T>(
514 &self,
515 f: impl FnOnce(&Connection) -> rusqlite::Result<T>,
516 ) -> rusqlite::Result<T> {
517 f(&self.lock())
518 }
519}
520
521#[cfg(test)]
524pub(crate) fn test_leaf(seq: u64, at: &str) -> Vec<u8> {
525 use crate::audit::leaf;
526 leaf::encode(
527 seq,
528 at,
529 leaf::action::START,
530 &leaf::Actor::Server,
531 leaf::subject_start("test"),
532 None,
533 )
534}
535
536fn use_durable_wal(conn: &Connection) -> Result<()> {
574 conn.busy_timeout(admin::BUSY_TIMEOUT)?;
575 let mode: String = conn.query_row("PRAGMA journal_mode = WAL", [], |r| r.get(0))?;
576 if !mode.eq_ignore_ascii_case("wal") {
577 anyhow::bail!("SQLite kept the {mode} journal");
578 }
579 conn.execute_batch(&format!(
580 "PRAGMA synchronous = FULL; PRAGMA journal_size_limit = {WAL_SIZE_LIMIT}"
581 ))?;
582 Ok(())
583}
584
585const WAL_SIZE_LIMIT: i64 = 16 * 1024 * 1024;
593
594const EXISTING_COLUMNS: &str = "content, deleted, COALESCE(source_env, ''), updated_at";
596
597fn existing_from(r: &rusqlite::Row<'_>) -> rusqlite::Result<Existing> {
598 Ok(Existing {
599 content: r.get(0)?,
600 deleted: r.get::<_, i64>(1)? != 0,
601 source_env: r.get(2)?,
602 updated_at: r.get(3)?,
603 })
604}
605
606fn read_file(conn: &Connection, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
608 Ok(conn
609 .query_row(
610 &format!("SELECT {EXISTING_COLUMNS} FROM memory_files WHERE project_key = ?1 AND file_path = ?2"),
611 (project_key, file_path),
612 existing_from,
613 )
614 .optional()?)
615}
616
617fn write_file(
622 conn: &Connection,
623 project_key: &str,
624 file_path: &str,
625 content: &str,
626 source_env: &str,
627 updated_at: &str,
628) -> Result<()> {
629 conn.execute(
630 "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
631 VALUES (?1, ?2, ?3, ?4, ?5, 0)
632 ON CONFLICT(project_key, file_path) DO UPDATE SET
633 content = excluded.content,
634 source_env = excluded.source_env,
635 updated_at = excluded.updated_at,
636 deleted = 0",
637 (project_key, file_path, content, nullable(source_env), updated_at),
638 )?;
639 Ok(())
640}
641
642fn nullable(s: &str) -> Option<&str> {
645 if s.is_empty() {
646 None
647 } else {
648 Some(s)
649 }
650}
651
652pub(crate) mod admin;
657
658#[cfg(test)]
659mod tests {
660 use super::*;
661
662 fn store() -> Store {
663 Store::open_in_memory().unwrap()
664 }
665
666 fn put(st: &Store, project_key: &str, file_path: &str, content: &str, source_env: &str) {
667 st.upsert_audited(project_key, file_path, content, source_env, test_leaf)
668 .unwrap();
669 }
670
671 fn del(st: &Store, project_key: &str, file_path: &str, source_env: &str) {
672 st.tombstone_audited(project_key, file_path, source_env, test_leaf)
673 .unwrap();
674 }
675
676 #[test]
679 fn a_write_is_stamped_with_its_leafs_at() {
680 let st = store();
681 let mut leaf_at = String::new();
682 let updated_at = st
683 .upsert_audited("acme/app", "a.md", "x", "laptop", |seq, at| {
684 leaf_at = at.to_string();
685 test_leaf(seq, at)
686 })
687 .unwrap();
688 assert_eq!(updated_at, leaf_at);
689 assert_eq!(st.list("acme/app").unwrap()[0].updated_at, updated_at);
690 assert_eq!(st.audit_checkpoint().0, 1);
691 }
692
693 #[test]
694 fn upsert_get_and_list_round_trip() {
695 let st = store();
696 put(&st, "acme/app", "MEMORY.md", "hello", "laptop");
697
698 let got = st.get("acme/app", "MEMORY.md").unwrap().unwrap();
699 assert_eq!(got.content, "hello");
700 assert!(!got.deleted);
701
702 let files = st.list("acme/app").unwrap();
703 assert_eq!(files.len(), 1);
704 assert_eq!(files[0].content.as_deref(), Some("hello"));
705 assert_eq!(files[0].source_env, "laptop");
706 assert!(st.get("acme/app", "missing.md").unwrap().is_none());
707 }
708
709 #[test]
712 fn tombstone_preserves_content_but_list_withholds_it() {
713 let st = store();
714 put(&st, "acme/app", "gone.md", "secret", "laptop");
715 del(&st, "acme/app", "gone.md", "laptop");
716
717 let row = st.get("acme/app", "gone.md").unwrap().unwrap();
718 assert_eq!(row.content, "secret", "content must stay recoverable");
719 assert!(row.deleted);
720
721 let files = st.list("acme/app").unwrap();
722 assert_eq!(
723 files.len(),
724 1,
725 "tombstones are listed so clients can delete locally"
726 );
727 assert!(files[0].deleted);
728 assert_eq!(files[0].content, None, "a pull must not resurrect it");
729 }
730
731 #[test]
733 fn upsert_clears_a_tombstone() {
734 let st = store();
735 del(&st, "acme/app", "f.md", "laptop");
736 put(&st, "acme/app", "f.md", "back", "laptop");
737 let row = st.get("acme/app", "f.md").unwrap().unwrap();
738 assert!(!row.deleted);
739 assert_eq!(row.content, "back");
740 }
741
742 #[test]
743 fn last_sync_at_is_empty_on_a_fresh_database() {
744 assert_eq!(store().last_sync_at().unwrap(), "");
745 }
746
747 #[test]
750 fn admin_stats_keeps_commas_inside_a_source_env() {
751 let st = store();
752 put(&st, "acme/app", "a.md", "x", "laptop,evil");
753 let (projects, _) = st.admin_stats().unwrap();
754 assert_eq!(projects[0].sources, vec!["laptop,evil".to_string()]);
755 }
756
757 #[test]
760 fn migrates_a_database_that_predates_tombstones() {
761 let dir = tempfile::tempdir().unwrap();
762 let path = dir.path().join("old.db");
763 {
764 let conn = Connection::open(&path).unwrap();
765 conn.execute_batch(
766 "CREATE TABLE memory_files (
767 project_key TEXT NOT NULL,
768 file_path TEXT NOT NULL,
769 content TEXT NOT NULL,
770 source_env TEXT,
771 updated_at TEXT NOT NULL,
772 PRIMARY KEY (project_key, file_path)
773 );
774 INSERT INTO memory_files VALUES ('acme/app','old.md','kept','node-era','2026-09-03T21:49:55.191Z');",
775 )
776 .unwrap();
777 }
778 let st = Store::open(&path).unwrap();
779 let files = st.list("acme/app").unwrap();
780 assert_eq!(files.len(), 1);
781 assert_eq!(files[0].content.as_deref(), Some("kept"));
782 assert!(!files[0].deleted);
783 }
784
785 #[test]
786 fn backup_names_carry_milliseconds() {
787 let dir = tempfile::tempdir().unwrap();
788 let st = store();
789 let dest = st.backup(dir.path(), 7).unwrap();
790 let name = dest.file_name().unwrap().to_str().unwrap();
791 assert!(
793 name.starts_with("recall-") && name.ends_with("Z.db"),
794 "got {name}"
795 );
796 let stamp = &name["recall-".len()..name.len() - ".db".len()];
797 assert_eq!(stamp.len(), 24, "got {stamp}");
798 assert!(
801 stamp[20..23].chars().all(|c| c.is_ascii_digit()),
802 "no millisecond field in {stamp}"
803 );
804 }
805
806 fn journal_mode(conn: &Connection) -> String {
809 conn.query_row("PRAGMA journal_mode", [], |r| r.get(0))
810 .unwrap()
811 }
812
813 fn header_mode(path: &Path) -> (u8, u8) {
816 let head = fs::read(path).unwrap();
817 (head[18], head[19])
818 }
819
820 fn wal_len(db: &Path) -> u64 {
821 let mut wal = db.as_os_str().to_owned();
822 wal.push("-wal");
823 fs::metadata(wal).map(|m| m.len()).unwrap_or(0)
824 }
825
826 #[test]
831 fn the_store_keeps_the_file_in_wal_and_syncs_every_commit() {
832 let dir = tempfile::tempdir().unwrap();
833 let path = dir.path().join("recall.db");
834 let st = Store::open(&path).unwrap();
835 put(&st, "acme/app", "a.md", "x", "laptop");
836 let (mode, sync, busy) = st
837 .with_raw(|c| {
838 Ok((
839 journal_mode(c),
840 c.query_row("PRAGMA synchronous", [], |r| r.get::<_, i64>(0))?,
841 c.query_row("PRAGMA busy_timeout", [], |r| r.get::<_, i64>(0))?,
842 ))
843 })
844 .unwrap();
845 assert_eq!(mode, "wal");
846 assert_eq!(sync, 2, "synchronous=FULL");
847 assert_eq!(busy, admin::BUSY_TIMEOUT.as_millis() as i64);
848 let limit: i64 = st
849 .with_raw(|c| c.query_row("PRAGMA journal_size_limit", [], |r| r.get(0)))
850 .unwrap();
851 assert_eq!(limit, WAL_SIZE_LIMIT);
852
853 assert_eq!(header_mode(&path), (2, 2));
854 assert_eq!(journal_mode(&Connection::open(&path).unwrap()), "wal");
855 assert!(wal_len(&path) > 0, "the commit went to the WAL");
856 }
857
858 #[test]
864 fn the_sqlite_compiled_in_has_the_wal_restart_fix() {
865 assert!(
866 rusqlite::version_number() >= 3_051_003,
867 "SQLite {} predates 3.51.3",
868 rusqlite::version()
869 );
870 }
871
872 #[test]
876 fn a_rollback_journal_database_from_an_older_server_is_converted_on_open() {
877 let dir = tempfile::tempdir().unwrap();
878 let path = dir.path().join("recall.db");
879 let before = {
880 let st = Store::open(&path).unwrap();
881 put(&st, "acme/app", "MEMORY.md", "kept\n", "laptop");
882 del(&st, "acme/app", "gone.md", "laptop");
883 st.audit_checkpoint()
884 };
885 Connection::open(&path)
888 .unwrap()
889 .query_row("PRAGMA journal_mode = DELETE", [], |_| Ok(()))
890 .unwrap();
891 assert_eq!(header_mode(&path), (1, 1));
892
893 let st = Store::open(&path).unwrap();
894 assert_eq!(header_mode(&path), (2, 2));
895 let files = st.list("acme/app").unwrap();
896 assert_eq!(files.len(), 2);
897 assert_eq!(files[0].content.as_deref(), Some("kept\n"));
898 assert!(files[1].deleted);
899 assert_eq!(st.audit_checkpoint(), before);
900 put(&st, "acme/app", "after.md", "y", "laptop");
901 drop(st);
902
903 let plain = Connection::open(&path).unwrap();
906 let n: i64 = plain
907 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
908 .unwrap();
909 assert_eq!(n, 3);
910 }
911
912 #[test]
919 fn a_switch_held_up_past_the_busy_timeout_refuses_to_start_and_changes_nothing() {
920 let dir = tempfile::tempdir().unwrap();
921 let path = dir.path().join("recall.db");
922 drop(Store::open(&path).unwrap());
923 Connection::open(&path)
924 .unwrap()
925 .query_row("PRAGMA journal_mode = DELETE", [], |_| Ok(()))
926 .unwrap();
927 let reader = Connection::open(&path).unwrap();
928 reader.execute_batch("BEGIN").unwrap();
929 let _: i64 = reader
930 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
931 .unwrap();
932
933 let err = match Store::open(&path) {
934 Ok(_) => panic!("switched to WAL under a reader holding the file"),
935 Err(e) => format!("{e:#}"),
936 };
937 assert!(err.contains("switching"), "{err}");
938 assert!(err.contains("start the server again"), "{err}");
939 assert!(err.contains("locked"), "{err}");
940 assert_eq!(header_mode(&path), (1, 1), "still the rollback journal");
941
942 reader.execute_batch("COMMIT").unwrap();
943 drop(reader);
944 drop(Store::open(&path).unwrap());
945 assert_eq!(header_mode(&path), (2, 2));
946 }
947
948 #[test]
953 fn the_node_servers_database_is_converted_with_every_row() {
954 let dir = tempfile::tempdir().unwrap();
955 let path = dir.path().join("recall.db");
956 fs::copy(
957 concat!(
958 env!("CARGO_MANIFEST_DIR"),
959 "/../../fixtures/node-written.db"
960 ),
961 &path,
962 )
963 .unwrap();
964 assert_eq!(header_mode(&path), (1, 1));
965 let dump = || -> Vec<(String, String, String, Option<String>, String, i64)> {
966 let conn = Connection::open(&path).unwrap();
967 let mut stmt = conn
968 .prepare(
969 "SELECT project_key, file_path, content, source_env, updated_at, deleted
970 FROM memory_files ORDER BY project_key, file_path",
971 )
972 .unwrap();
973 let rows = stmt
974 .query_map([], |r| {
975 Ok((
976 r.get(0)?,
977 r.get(1)?,
978 r.get(2)?,
979 r.get(3)?,
980 r.get(4)?,
981 r.get(5)?,
982 ))
983 })
984 .unwrap();
985 rows.map(Result::unwrap).collect()
986 };
987 let before = dump();
988 assert!(!before.is_empty());
989
990 let st = Store::open(&path).unwrap();
991 assert_eq!(header_mode(&path), (2, 2));
992 assert_eq!(dump(), before);
993 drop(st);
994 assert_eq!(dump(), before);
995 }
996
997 #[test]
1001 fn a_snapshot_is_one_self_contained_file_with_what_the_wal_holds() {
1002 let dir = tempfile::tempdir().unwrap();
1003 let path = dir.path().join("recall.db");
1004 let st = Store::open(&path).unwrap();
1005 for i in 0..20 {
1006 put(&st, "acme/app", &format!("f{i}.md"), "x", "laptop");
1007 }
1008 assert!(wal_len(&path) > 0, "the pushes are still in the WAL");
1009
1010 let snap = st.backup(dir.path().join("backups"), 7).unwrap();
1011 let names: Vec<_> = fs::read_dir(snap.parent().unwrap())
1012 .unwrap()
1013 .map(|e| e.unwrap().file_name().into_string().unwrap())
1014 .collect();
1015 assert_eq!(names.len(), 1, "no -wal or -shm beside it: {names:?}");
1016 assert_eq!(header_mode(&snap), (1, 1));
1017
1018 let read =
1019 Connection::open_with_flags(&snap, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY).unwrap();
1020 assert_eq!(journal_mode(&read), "delete");
1021 let n: i64 = read
1022 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1023 .unwrap();
1024 assert_eq!(n, 20);
1025 }
1026
1027 fn rows_in_the_file_alone(db: &Path) -> i64 {
1032 let alone = tempfile::tempdir().unwrap();
1033 let copy = alone.path().join("recall.db");
1034 fs::copy(db, ©).unwrap();
1035 Connection::open(©)
1036 .unwrap()
1037 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1038 .unwrap()
1039 }
1040
1041 #[test]
1045 fn a_checkpoint_brings_the_file_on_its_own_up_to_date() {
1046 let dir = tempfile::tempdir().unwrap();
1047 let path = dir.path().join("recall.db");
1048 let st = Store::open(&path).unwrap();
1049 put(&st, "acme/app", "a.md", "x", "laptop");
1050 assert!(st.checkpoint_all().unwrap());
1051 assert_eq!(wal_len(&path), 0, "emptied");
1052 assert_eq!(rows_in_the_file_alone(&path), 1);
1053
1054 put(&st, "acme/app", "b.md", "x", "laptop");
1055 assert_eq!(rows_in_the_file_alone(&path), 1, "b.md is only in the WAL");
1056 assert!(st.checkpoint().unwrap(), "nothing held it back");
1057 assert_eq!(rows_in_the_file_alone(&path), 2);
1058
1059 put(&st, "acme/app", "c.md", "x", "laptop");
1060 assert!(st.checkpoint_all().unwrap());
1061 assert_eq!(wal_len(&path), 0, "emptied");
1062 assert_eq!(rows_in_the_file_alone(&path), 3);
1063 }
1064
1065 #[test]
1069 fn a_checkpoint_says_when_a_reader_held_it_back() {
1070 let dir = tempfile::tempdir().unwrap();
1071 let path = dir.path().join("recall.db");
1072 let st = Store::open(&path).unwrap();
1073 put(&st, "acme/app", "a.md", "x", "laptop");
1074 assert!(st.checkpoint_all().unwrap());
1075
1076 let reader = Connection::open(&path).unwrap();
1077 reader.execute_batch("BEGIN").unwrap();
1078 let _: i64 = reader
1079 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1080 .unwrap();
1081 put(&st, "acme/app", "b.md", "x", "laptop");
1082 assert!(!st.checkpoint().unwrap(), "held back by the reader");
1083 assert_eq!(rows_in_the_file_alone(&path), 1);
1084
1085 reader.execute_batch("COMMIT").unwrap();
1086 assert!(st.checkpoint().unwrap());
1087 assert_eq!(rows_in_the_file_alone(&path), 2);
1088 }
1089
1090 #[test]
1094 fn a_wal_grown_large_is_cut_back_to_the_limit() {
1095 let dir = tempfile::tempdir().unwrap();
1096 let path = dir.path().join("recall.db");
1097 let st = Store::open(&path).unwrap();
1098 let big = "x".repeat(24 * 1024 * 1024);
1099 put(&st, "acme/app", "big.md", &big, "laptop");
1100 let grown = wal_len(&path);
1101 assert!(grown > WAL_SIZE_LIMIT as u64, "{grown}");
1102 assert!(st.checkpoint().unwrap());
1103
1104 put(&st, "acme/app", "small.md", "x", "laptop");
1105 let now = wal_len(&path);
1106 assert!(now <= WAL_SIZE_LIMIT as u64, "{now} after {grown}");
1107 assert_eq!(st.get("acme/app", "big.md").unwrap().unwrap().content, big);
1108 }
1109
1110 #[cfg(unix)]
1117 #[test]
1118 fn a_database_moved_aside_under_a_running_store_is_not_written_to() {
1119 let dir = tempfile::tempdir().unwrap();
1120 let path = dir.path().join("recall.db");
1121 let st = Store::open(&path).unwrap();
1122 put(&st, "acme/app", "a.md", "x", "laptop");
1123 let snapshot = st.backup(dir.path().join("backups"), 7).unwrap();
1124
1125 let aside = dir.path().join("aside");
1126 fs::create_dir(&aside).unwrap();
1127 for f in ["recall.db", "recall.db-wal", "recall.db-shm"] {
1128 if dir.path().join(f).exists() {
1129 fs::rename(dir.path().join(f), aside.join(f)).unwrap();
1130 }
1131 }
1132 let refused = |st: &Store| {
1135 let err = st
1136 .upsert_audited("acme/app", "b.md", "y", "laptop", test_leaf)
1137 .unwrap_err();
1138 let said = format!("{err:#}");
1139 assert!(said.contains("moved or replaced"), "{said}");
1140 assert!(!said.contains(" "), "a run of spaces: {said:?}");
1141 let err = st.audit_append(test_leaf).unwrap_err();
1142 assert!(format!("{err:#}").contains("moved or replaced"), "{err:#}");
1143 };
1144 refused(&st);
1145 fs::copy(&snapshot, &path).unwrap();
1146 refused(&st);
1147 drop(st);
1148
1149 for db in [aside.join("recall.db"), path] {
1150 let files = Store::open(&db).unwrap().list("acme/app").unwrap();
1151 let names: Vec<_> = files.iter().map(|f| f.file_path.as_str()).collect();
1152 assert_eq!(names, ["a.md"], "{}", db.display());
1153 }
1154 }
1155
1156 #[tokio::test]
1165 #[ignore = "a measurement, not a check: run it by hand"]
1166 async fn push_latency_by_journal_mode() {
1167 use axum::body::Body;
1168 use axum::http::Request;
1169 use std::time::{Duration, Instant};
1170 use tower::ServiceExt;
1171
1172 const PUSHES: usize = 400;
1173 const TOKEN: &str = "bench-token";
1174 let content = "- a remembered fact, about as long as one usually is\n".repeat(20);
1175 println!("| journal | push p50 | push p90 | push mean | pull p50 | pull p90 |");
1176 println!("|---|---|---|---|---|---|");
1177 for (label, journal, sync) in [
1178 ("rollback (DELETE), FULL", "DELETE", "FULL"),
1179 ("WAL, NORMAL", "WAL", "NORMAL"),
1180 ("WAL, FULL", "WAL", "FULL"),
1181 ] {
1182 let dir = tempfile::tempdir().unwrap();
1183 let store = std::sync::Arc::new(Store::open(dir.path().join("recall.db")).unwrap());
1184 store
1185 .with_raw(|c| {
1186 c.query_row(&format!("PRAGMA journal_mode = {journal}"), [], |_| Ok(()))?;
1187 c.execute_batch(&format!("PRAGMA synchronous = {sync}"))
1188 })
1189 .unwrap();
1190 let server = crate::Server::new(
1191 crate::Config {
1192 token: TOKEN.into(),
1193 merge_enabled: false,
1194 rate_limit_max: 1_000_000,
1195 ..crate::Config::default()
1196 },
1197 store,
1198 );
1199 let router = server.router();
1200 let push = |i: usize| {
1201 let body = serde_json::json!({
1202 "project_key": "bench/app",
1203 "file_path": format!("f{i}.md"),
1204 "content": content,
1205 "source_env": "bench",
1206 });
1207 Request::post("/sync")
1208 .header("authorization", format!("Bearer {TOKEN}"))
1209 .header("content-type", "application/json")
1210 .body(Body::from(body.to_string()))
1211 .unwrap()
1212 };
1213 let pull = || {
1214 Request::get("/sync?project_key=small/app")
1215 .header("authorization", format!("Bearer {TOKEN}"))
1216 .body(Body::empty())
1217 .unwrap()
1218 };
1219 let time = |mut samples: Vec<Duration>| {
1220 samples.sort();
1221 let ms = |d: Duration| d.as_secs_f64() * 1000.0;
1222 let mean = ms(samples.iter().sum::<Duration>()) / samples.len() as f64;
1223 (
1224 ms(samples[samples.len() / 2]),
1225 ms(samples[samples.len() * 9 / 10]),
1226 mean,
1227 )
1228 };
1229 for i in 0..20 {
1230 let resp = router.clone().oneshot(push(PUSHES + i)).await.unwrap();
1231 assert_eq!(resp.status(), 200);
1232 }
1233 let mut pushes = Vec::with_capacity(PUSHES);
1234 for i in 0..PUSHES {
1235 let started = Instant::now();
1236 let resp = router.clone().oneshot(push(i)).await.unwrap();
1237 pushes.push(started.elapsed());
1238 assert_eq!(resp.status(), 200);
1239 }
1240 let mut pulls = Vec::with_capacity(PUSHES);
1241 for _ in 0..PUSHES {
1242 let started = Instant::now();
1243 let resp = router.clone().oneshot(pull()).await.unwrap();
1244 pulls.push(started.elapsed());
1245 assert_eq!(resp.status(), 200);
1246 }
1247 let (p50, p90, mean) = time(pushes);
1248 let (l50, l90, _) = time(pulls);
1249 println!(
1250 "| {label} | {p50:.2} ms | {p90:.2} ms | {mean:.2} ms | {l50:.2} ms | {l90:.2} ms |"
1251 );
1252 }
1253 }
1254}