1use std::fs;
10use std::path::{Path, PathBuf};
11use std::sync::{Mutex, MutexGuard, PoisonError};
12use std::time::Duration;
13
14use anyhow::{Context, Result};
15use recall_wire::{AdminTotals, File, ProjectStats};
16use rusqlite::{Connection, OptionalExtension};
17
18use crate::audit::merkle::Tree;
19use crate::now;
20
21mod audit;
22mod devices;
23mod evaluations;
24mod jobs;
25mod passkeys;
26
27pub use audit::{AuditEntry, ConsistencyError, Outcome};
28pub use devices::{
29 plain_name, Created, Decision, Inserted, NewAuthkey, NewDevice, NewEnrollment, Poll, Waiting,
30};
31pub use evaluations::Requested;
32pub use jobs::{
33 clip, Failure, Queued, Retried, Settled, Settlement, MAX_ATTEMPTS, MAX_ERROR_BYTES, MAX_LINKS,
34 MAX_OPEN_JOBS,
35};
36pub use passkeys::{
37 AddedCredential, AdminCredential, AdminSession, BootstrapCode, FirstPasskey,
38 NewAdminCredential, RemovedCredential,
39};
40
41const SCHEMA: &str = "
43 CREATE TABLE IF NOT EXISTS memory_files (
44 project_key TEXT NOT NULL,
45 file_path TEXT NOT NULL,
46 content TEXT NOT NULL,
47 source_env TEXT,
48 updated_at TEXT NOT NULL,
49 deleted INTEGER NOT NULL DEFAULT 0,
50 PRIMARY KEY (project_key, file_path)
51 );
52";
53
54#[derive(Debug, Clone, PartialEq, Eq)]
57pub struct Existing {
58 pub content: String,
61 pub deleted: bool,
63 pub source_env: String,
65 pub updated_at: String,
67}
68
69struct StoreState {
86 conn: Connection,
87 audit: Tree,
88 audit_at: String,
89 file: Option<FileId>,
90 log_moved: bool,
91 moved_said: bool,
92}
93
94type FileId = (u64, u64);
113
114#[cfg(unix)]
118fn file_id(conn: &Connection) -> Option<FileId> {
119 use std::os::unix::fs::MetadataExt;
120 let path = conn.path().filter(|p| !p.is_empty())?;
121 fs::metadata(path).ok().map(|m| (m.dev(), m.ino()))
122}
123
124#[cfg(not(unix))]
125fn file_id(_conn: &Connection) -> Option<FileId> {
126 None
127}
128
129impl std::ops::Deref for StoreState {
130 type Target = Connection;
131 fn deref(&self) -> &Connection {
132 &self.conn
133 }
134}
135
136impl std::ops::DerefMut for StoreState {
137 fn deref_mut(&mut self) -> &mut Connection {
138 &mut self.conn
139 }
140}
141
142pub struct Store {
144 state: Mutex<StoreState>,
151}
152
153impl Store {
154 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
164 Self::open_waiting(path.as_ref(), admin::BUSY_TIMEOUT)
165 }
166
167 fn open_waiting(path: &Path, busy: Duration) -> Result<Self> {
172 if let Some(dir) = path.parent() {
173 if !dir.as_os_str().is_empty() {
174 fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
175 }
176 }
177 let conn = Connection::open(path).with_context(|| format!("opening {}", path.display()))?;
178 use_durable_wal(&conn, busy).with_context(|| {
179 format!(
180 "switching {} to SQLite's WAL journal. It needs a local filesystem, and a \
181 moment with no other process holding the file (sqlite-web mid-read, an admin \
182 command): start the server again",
183 path.display()
184 )
185 })?;
186 Self::with_connection(conn)
187 }
188
189 pub fn open_in_memory() -> Result<Self> {
191 Self::with_connection(Connection::open_in_memory()?)
192 }
193
194 fn with_connection(conn: Connection) -> Result<Self> {
195 let file = file_id(&conn);
196 let store = Self {
197 state: Mutex::new(StoreState {
198 conn,
199 audit: Tree::new(),
200 audit_at: String::new(),
201 file,
202 log_moved: true,
203 moved_said: false,
204 }),
205 };
206 store.migrate()?;
207 Ok(store)
208 }
209
210 fn lock(&self) -> MutexGuard<'_, StoreState> {
214 self.state.lock().unwrap_or_else(PoisonError::into_inner)
215 }
216
217 fn migrate(&self) -> Result<()> {
218 let mut state = self.lock();
219 state.conn.execute_batch(SCHEMA)?;
220
221 let has_deleted = {
224 let mut stmt = state.conn.prepare("PRAGMA table_info(memory_files)")?;
225 let mut rows = stmt.query([])?;
226 let mut found = false;
227 while let Some(row) = rows.next()? {
228 if row.get::<_, String>(1)? == "deleted" {
229 found = true;
230 }
231 }
232 found
233 };
234 if !has_deleted {
235 state.conn.execute(
236 "ALTER TABLE memory_files ADD COLUMN deleted INTEGER NOT NULL DEFAULT 0",
237 [],
238 )?;
239 }
240
241 state.conn.execute_batch(devices::SCHEMA)?;
246 state.conn.execute_batch(passkeys::SCHEMA)?;
247 devices::allow_worker_scope(&state.conn)?;
252 state.conn.execute_batch(jobs::SCHEMA)?;
253 state.conn.execute_batch(evaluations::SCHEMA)?;
254 state.conn.execute_batch(audit::SCHEMA)?;
255
256 let loaded = audit::load(&state.conn)?;
257 state.audit = loaded.tree;
258 state.audit_at = loaded.last_at;
259 Ok(())
260 }
261
262 pub fn get(&self, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
264 read_file(&self.lock(), project_key, file_path)
265 }
266
267 pub fn upsert_audited(
272 &self,
273 project_key: &str,
274 file_path: &str,
275 content: &str,
276 source_env: &str,
277 build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
278 ) -> Result<String> {
279 self.audited(
280 |tx, at| {
281 write_file(tx, project_key, file_path, content, source_env, at)?;
282 Ok(Outcome::Commit(at.to_string()))
283 },
284 |seq, at, _| build_leaf(seq, at),
285 )
286 }
287
288 pub fn tombstone_audited(
298 &self,
299 project_key: &str,
300 file_path: &str,
301 source_env: &str,
302 build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
303 ) -> Result<String> {
304 self.audited(
305 |tx, at| {
306 tx.execute(
307 "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
308 VALUES (?1, ?2, '', ?3, ?4, 1)
309 ON CONFLICT(project_key, file_path) DO UPDATE SET
310 source_env = excluded.source_env,
311 updated_at = excluded.updated_at,
312 deleted = 1",
313 (project_key, file_path, nullable(source_env), at),
314 )?;
315 jobs::close_for_delete(tx, project_key, file_path, at)?;
316 Ok(Outcome::Commit(at.to_string()))
317 },
318 |seq, at, _| build_leaf(seq, at),
319 )
320 }
321
322 pub fn list(&self, project_key: &str) -> Result<Vec<File>> {
326 let conn = self.lock();
327 let mut stmt = conn.prepare(
328 "SELECT file_path, content, COALESCE(source_env, ''), updated_at, deleted
329 FROM memory_files WHERE project_key = ?1 ORDER BY file_path",
330 )?;
331 let rows = stmt.query_map((project_key,), |r| {
332 let content: String = r.get(1)?;
333 let deleted = r.get::<_, i64>(4)? != 0;
334 Ok(File {
335 file_path: r.get(0)?,
336 content: if deleted { None } else { Some(content) },
337 source_env: r.get(2)?,
338 updated_at: r.get(3)?,
339 deleted,
340 })
341 })?;
342 let mut files = Vec::new();
343 for row in rows {
344 files.push(row?);
345 }
346 Ok(files)
347 }
348
349 pub fn last_sync_at(&self) -> Result<String> {
352 let conn = self.lock();
353 let v: Option<String> =
354 conn.query_row("SELECT MAX(updated_at) FROM memory_files", [], |r| r.get(0))?;
355 Ok(v.unwrap_or_default())
356 }
357
358 pub fn admin_stats(&self) -> Result<(Vec<ProjectStats>, AdminTotals)> {
360 let conn = self.lock();
361
362 let mut projects = Vec::new();
363 let mut totals = AdminTotals::default();
364 {
365 let mut stmt = conn.prepare(
366 "SELECT project_key,
367 SUM(CASE WHEN deleted = 0 THEN 1 ELSE 0 END),
368 SUM(CASE WHEN deleted = 1 THEN 1 ELSE 0 END),
369 MAX(updated_at)
370 FROM memory_files GROUP BY project_key ORDER BY MAX(updated_at) DESC",
371 )?;
372 let rows = stmt.query_map([], |r| {
373 Ok(ProjectStats {
374 project_key: r.get(0)?,
375 file_count: r.get(1)?,
376 deleted_count: r.get(2)?,
377 sources: Vec::new(),
378 last_updated_at: r.get::<_, Option<String>>(3)?.unwrap_or_default(),
379 })
380 })?;
381 for row in rows {
382 let p = row?;
383 totals.file_count += p.file_count;
384 totals.deleted_count += p.deleted_count;
385 projects.push(p);
386 }
387 }
388 totals.project_count = projects.len() as i64;
389
390 {
395 let mut stmt = conn.prepare(
396 "SELECT DISTINCT project_key, source_env FROM memory_files WHERE source_env IS NOT NULL",
397 )?;
398 let rows =
399 stmt.query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?)))?;
400 for row in rows {
401 let (key, src) = row?;
402 if src.is_empty() {
403 continue;
404 }
405 if let Some(p) = projects.iter_mut().find(|p| p.project_key == key) {
406 p.sources.push(src);
407 }
408 }
409 }
410 for p in &mut projects {
411 p.sources.sort();
412 }
413 Ok((projects, totals))
414 }
415
416 pub fn checkpoint(&self) -> Result<bool> {
430 let (frames, copied): (i64, i64) =
431 self.lock()
432 .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |r| {
433 Ok((r.get(1)?, r.get(2)?))
434 })?;
435 Ok(copied >= frames)
436 }
437
438 pub fn checkpoint_all(&self) -> Result<bool> {
449 let busy: i64 = self
450 .lock()
451 .query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |r| r.get(0))?;
452 Ok(busy == 0)
453 }
454
455 pub fn backup(&self, dir: impl AsRef<Path>, keep: usize) -> Result<PathBuf> {
466 let dir = dir.as_ref();
467 fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
468
469 let stamp = now().replace([':', '.'], "-");
476 let dest = dir.join(format!("recall-{stamp}.db"));
477 let dest_str = dest
478 .to_str()
479 .context("backup path is not valid UTF-8")?
480 .to_owned();
481
482 let existed = dest.exists();
485 let vacuumed = {
486 let conn = self.lock();
487 conn.execute("VACUUM INTO ?1", (&dest_str,))
488 .with_context(|| format!("VACUUM INTO {dest_str}"))
489 };
490 if let Err(err) = vacuumed {
491 if !existed {
497 let _ = fs::remove_file(&dest);
498 }
499 return Err(err);
500 }
501
502 let mut snapshots: Vec<PathBuf> = fs::read_dir(dir)?
503 .filter_map(|e| e.ok())
504 .map(|e| e.path())
505 .filter(|p| {
506 p.file_name()
507 .and_then(|n| n.to_str())
508 .is_some_and(|n| n.starts_with("recall-") && n.ends_with(".db"))
509 })
510 .collect();
511 snapshots.sort();
512 for stale in snapshots.iter().take(snapshots.len().saturating_sub(keep)) {
513 let _ = fs::remove_file(stale);
514 }
515 Ok(dest)
516 }
517}
518
519#[cfg(test)]
520impl Store {
521 pub(crate) fn with_raw<T>(
525 &self,
526 f: impl FnOnce(&Connection) -> rusqlite::Result<T>,
527 ) -> rusqlite::Result<T> {
528 f(&self.lock())
529 }
530}
531
532#[cfg(test)]
535pub(crate) fn test_leaf(seq: u64, at: &str) -> Vec<u8> {
536 use crate::audit::leaf;
537 leaf::encode(
538 seq,
539 at,
540 leaf::action::START,
541 &leaf::Actor::Server,
542 leaf::subject_start("test"),
543 None,
544 )
545}
546
547fn use_durable_wal(conn: &Connection, busy: Duration) -> Result<()> {
585 conn.busy_timeout(busy)?;
586 let mode: String = conn.query_row("PRAGMA journal_mode = WAL", [], |r| r.get(0))?;
587 if !mode.eq_ignore_ascii_case("wal") {
588 anyhow::bail!("SQLite kept the {mode} journal");
589 }
590 conn.execute_batch(&format!(
591 "PRAGMA synchronous = FULL; PRAGMA journal_size_limit = {WAL_SIZE_LIMIT}"
592 ))?;
593 Ok(())
594}
595
596const WAL_SIZE_LIMIT: i64 = 16 * 1024 * 1024;
604
605const EXISTING_COLUMNS: &str = "content, deleted, COALESCE(source_env, ''), updated_at";
607
608fn existing_from(r: &rusqlite::Row<'_>) -> rusqlite::Result<Existing> {
609 Ok(Existing {
610 content: r.get(0)?,
611 deleted: r.get::<_, i64>(1)? != 0,
612 source_env: r.get(2)?,
613 updated_at: r.get(3)?,
614 })
615}
616
617fn read_file(conn: &Connection, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
619 Ok(conn
620 .query_row(
621 &format!("SELECT {EXISTING_COLUMNS} FROM memory_files WHERE project_key = ?1 AND file_path = ?2"),
622 (project_key, file_path),
623 existing_from,
624 )
625 .optional()?)
626}
627
628fn write_file(
633 conn: &Connection,
634 project_key: &str,
635 file_path: &str,
636 content: &str,
637 source_env: &str,
638 updated_at: &str,
639) -> Result<()> {
640 conn.execute(
641 "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
642 VALUES (?1, ?2, ?3, ?4, ?5, 0)
643 ON CONFLICT(project_key, file_path) DO UPDATE SET
644 content = excluded.content,
645 source_env = excluded.source_env,
646 updated_at = excluded.updated_at,
647 deleted = 0",
648 (project_key, file_path, content, nullable(source_env), updated_at),
649 )?;
650 Ok(())
651}
652
653fn nullable(s: &str) -> Option<&str> {
656 if s.is_empty() {
657 None
658 } else {
659 Some(s)
660 }
661}
662
663pub(crate) mod admin;
668
669#[cfg(test)]
670mod tests {
671 use super::*;
672
673 fn store() -> Store {
674 Store::open_in_memory().unwrap()
675 }
676
677 fn put(st: &Store, project_key: &str, file_path: &str, content: &str, source_env: &str) {
678 st.upsert_audited(project_key, file_path, content, source_env, test_leaf)
679 .unwrap();
680 }
681
682 fn del(st: &Store, project_key: &str, file_path: &str, source_env: &str) {
683 st.tombstone_audited(project_key, file_path, source_env, test_leaf)
684 .unwrap();
685 }
686
687 #[test]
690 fn a_write_is_stamped_with_its_leafs_at() {
691 let st = store();
692 let mut leaf_at = String::new();
693 let updated_at = st
694 .upsert_audited("acme/app", "a.md", "x", "laptop", |seq, at| {
695 leaf_at = at.to_string();
696 test_leaf(seq, at)
697 })
698 .unwrap();
699 assert_eq!(updated_at, leaf_at);
700 assert_eq!(st.list("acme/app").unwrap()[0].updated_at, updated_at);
701 assert_eq!(st.audit_checkpoint().0, 1);
702 }
703
704 #[test]
705 fn upsert_get_and_list_round_trip() {
706 let st = store();
707 put(&st, "acme/app", "MEMORY.md", "hello", "laptop");
708
709 let got = st.get("acme/app", "MEMORY.md").unwrap().unwrap();
710 assert_eq!(got.content, "hello");
711 assert!(!got.deleted);
712
713 let files = st.list("acme/app").unwrap();
714 assert_eq!(files.len(), 1);
715 assert_eq!(files[0].content.as_deref(), Some("hello"));
716 assert_eq!(files[0].source_env, "laptop");
717 assert!(st.get("acme/app", "missing.md").unwrap().is_none());
718 }
719
720 #[test]
723 fn tombstone_preserves_content_but_list_withholds_it() {
724 let st = store();
725 put(&st, "acme/app", "gone.md", "secret", "laptop");
726 del(&st, "acme/app", "gone.md", "laptop");
727
728 let row = st.get("acme/app", "gone.md").unwrap().unwrap();
729 assert_eq!(row.content, "secret", "content must stay recoverable");
730 assert!(row.deleted);
731
732 let files = st.list("acme/app").unwrap();
733 assert_eq!(
734 files.len(),
735 1,
736 "tombstones are listed so clients can delete locally"
737 );
738 assert!(files[0].deleted);
739 assert_eq!(files[0].content, None, "a pull must not resurrect it");
740 }
741
742 #[test]
744 fn upsert_clears_a_tombstone() {
745 let st = store();
746 del(&st, "acme/app", "f.md", "laptop");
747 put(&st, "acme/app", "f.md", "back", "laptop");
748 let row = st.get("acme/app", "f.md").unwrap().unwrap();
749 assert!(!row.deleted);
750 assert_eq!(row.content, "back");
751 }
752
753 #[test]
754 fn last_sync_at_is_empty_on_a_fresh_database() {
755 assert_eq!(store().last_sync_at().unwrap(), "");
756 }
757
758 #[test]
761 fn admin_stats_keeps_commas_inside_a_source_env() {
762 let st = store();
763 put(&st, "acme/app", "a.md", "x", "laptop,evil");
764 let (projects, _) = st.admin_stats().unwrap();
765 assert_eq!(projects[0].sources, vec!["laptop,evil".to_string()]);
766 }
767
768 #[test]
771 fn migrates_a_database_that_predates_tombstones() {
772 let dir = tempfile::tempdir().unwrap();
773 let path = dir.path().join("old.db");
774 {
775 let conn = Connection::open(&path).unwrap();
776 conn.execute_batch(
777 "CREATE TABLE memory_files (
778 project_key TEXT NOT NULL,
779 file_path TEXT NOT NULL,
780 content TEXT NOT NULL,
781 source_env TEXT,
782 updated_at TEXT NOT NULL,
783 PRIMARY KEY (project_key, file_path)
784 );
785 INSERT INTO memory_files VALUES ('acme/app','old.md','kept','node-era','2026-09-03T21:49:55.191Z');",
786 )
787 .unwrap();
788 }
789 let st = Store::open(&path).unwrap();
790 let files = st.list("acme/app").unwrap();
791 assert_eq!(files.len(), 1);
792 assert_eq!(files[0].content.as_deref(), Some("kept"));
793 assert!(!files[0].deleted);
794 }
795
796 #[test]
797 fn backup_names_carry_milliseconds() {
798 let dir = tempfile::tempdir().unwrap();
799 let st = store();
800 let dest = st.backup(dir.path(), 7).unwrap();
801 let name = dest.file_name().unwrap().to_str().unwrap();
802 assert!(
804 name.starts_with("recall-") && name.ends_with("Z.db"),
805 "got {name}"
806 );
807 let stamp = &name["recall-".len()..name.len() - ".db".len()];
808 assert_eq!(stamp.len(), 24, "got {stamp}");
809 assert!(
812 stamp[20..23].chars().all(|c| c.is_ascii_digit()),
813 "no millisecond field in {stamp}"
814 );
815 }
816
817 fn journal_mode(conn: &Connection) -> String {
820 conn.query_row("PRAGMA journal_mode", [], |r| r.get(0))
821 .unwrap()
822 }
823
824 fn header_mode(path: &Path) -> (u8, u8) {
827 let head = fs::read(path).unwrap();
828 (head[18], head[19])
829 }
830
831 fn wal_len(db: &Path) -> u64 {
832 let mut wal = db.as_os_str().to_owned();
833 wal.push("-wal");
834 fs::metadata(wal).map(|m| m.len()).unwrap_or(0)
835 }
836
837 #[test]
842 fn the_store_keeps_the_file_in_wal_and_syncs_every_commit() {
843 let dir = tempfile::tempdir().unwrap();
844 let path = dir.path().join("recall.db");
845 let st = Store::open(&path).unwrap();
846 put(&st, "acme/app", "a.md", "x", "laptop");
847 let (mode, sync, busy) = st
848 .with_raw(|c| {
849 Ok((
850 journal_mode(c),
851 c.query_row("PRAGMA synchronous", [], |r| r.get::<_, i64>(0))?,
852 c.query_row("PRAGMA busy_timeout", [], |r| r.get::<_, i64>(0))?,
853 ))
854 })
855 .unwrap();
856 assert_eq!(mode, "wal");
857 assert_eq!(sync, 2, "synchronous=FULL");
858 assert_eq!(busy, admin::BUSY_TIMEOUT.as_millis() as i64);
859 let limit: i64 = st
860 .with_raw(|c| c.query_row("PRAGMA journal_size_limit", [], |r| r.get(0)))
861 .unwrap();
862 assert_eq!(limit, WAL_SIZE_LIMIT);
863
864 assert_eq!(header_mode(&path), (2, 2));
865 assert_eq!(journal_mode(&Connection::open(&path).unwrap()), "wal");
866 assert!(wal_len(&path) > 0, "the commit went to the WAL");
867 }
868
869 #[test]
875 fn the_sqlite_compiled_in_has_the_wal_restart_fix() {
876 assert!(
877 rusqlite::version_number() >= 3_051_003,
878 "SQLite {} predates 3.51.3",
879 rusqlite::version()
880 );
881 }
882
883 #[test]
887 fn a_rollback_journal_database_from_an_older_server_is_converted_on_open() {
888 let dir = tempfile::tempdir().unwrap();
889 let path = dir.path().join("recall.db");
890 let before = {
891 let st = Store::open(&path).unwrap();
892 put(&st, "acme/app", "MEMORY.md", "kept\n", "laptop");
893 del(&st, "acme/app", "gone.md", "laptop");
894 st.audit_checkpoint()
895 };
896 Connection::open(&path)
899 .unwrap()
900 .query_row("PRAGMA journal_mode = DELETE", [], |_| Ok(()))
901 .unwrap();
902 assert_eq!(header_mode(&path), (1, 1));
903
904 let st = Store::open(&path).unwrap();
905 assert_eq!(header_mode(&path), (2, 2));
906 let files = st.list("acme/app").unwrap();
907 assert_eq!(files.len(), 2);
908 assert_eq!(files[0].content.as_deref(), Some("kept\n"));
909 assert!(files[1].deleted);
910 assert_eq!(st.audit_checkpoint(), before);
911 put(&st, "acme/app", "after.md", "y", "laptop");
912 drop(st);
913
914 let plain = Connection::open(&path).unwrap();
917 let n: i64 = plain
918 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
919 .unwrap();
920 assert_eq!(n, 3);
921 }
922
923 #[test]
930 fn a_switch_held_up_past_the_busy_timeout_refuses_to_start_and_changes_nothing() {
931 let dir = tempfile::tempdir().unwrap();
932 let path = dir.path().join("recall.db");
933 drop(Store::open(&path).unwrap());
934 Connection::open(&path)
935 .unwrap()
936 .query_row("PRAGMA journal_mode = DELETE", [], |_| Ok(()))
937 .unwrap();
938 let reader = Connection::open(&path).unwrap();
939 reader.execute_batch("BEGIN").unwrap();
940 let _: i64 = reader
941 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
942 .unwrap();
943
944 let err = match Store::open_waiting(&path, Duration::from_millis(250)) {
949 Ok(_) => panic!("switched to WAL under a reader holding the file"),
950 Err(e) => format!("{e:#}"),
951 };
952 assert!(err.contains("switching"), "{err}");
953 assert!(err.contains("start the server again"), "{err}");
954 assert!(err.contains("locked"), "{err}");
955 assert_eq!(header_mode(&path), (1, 1), "still the rollback journal");
956
957 reader.execute_batch("COMMIT").unwrap();
958 drop(reader);
959 drop(Store::open(&path).unwrap());
960 assert_eq!(header_mode(&path), (2, 2));
961 }
962
963 #[test]
968 fn the_node_servers_database_is_converted_with_every_row() {
969 let dir = tempfile::tempdir().unwrap();
970 let path = dir.path().join("recall.db");
971 fs::copy(
972 concat!(
973 env!("CARGO_MANIFEST_DIR"),
974 "/../../fixtures/node-written.db"
975 ),
976 &path,
977 )
978 .unwrap();
979 assert_eq!(header_mode(&path), (1, 1));
980 let dump = || -> Vec<(String, String, String, Option<String>, String, i64)> {
981 let conn = Connection::open(&path).unwrap();
982 let mut stmt = conn
983 .prepare(
984 "SELECT project_key, file_path, content, source_env, updated_at, deleted
985 FROM memory_files ORDER BY project_key, file_path",
986 )
987 .unwrap();
988 let rows = stmt
989 .query_map([], |r| {
990 Ok((
991 r.get(0)?,
992 r.get(1)?,
993 r.get(2)?,
994 r.get(3)?,
995 r.get(4)?,
996 r.get(5)?,
997 ))
998 })
999 .unwrap();
1000 rows.map(Result::unwrap).collect()
1001 };
1002 let before = dump();
1003 assert!(!before.is_empty());
1004
1005 let st = Store::open(&path).unwrap();
1006 assert_eq!(header_mode(&path), (2, 2));
1007 assert_eq!(dump(), before);
1008 drop(st);
1009 assert_eq!(dump(), before);
1010 }
1011
1012 #[test]
1016 fn a_snapshot_is_one_self_contained_file_with_what_the_wal_holds() {
1017 let dir = tempfile::tempdir().unwrap();
1018 let path = dir.path().join("recall.db");
1019 let st = Store::open(&path).unwrap();
1020 for i in 0..20 {
1021 put(&st, "acme/app", &format!("f{i}.md"), "x", "laptop");
1022 }
1023 assert!(wal_len(&path) > 0, "the pushes are still in the WAL");
1024
1025 let snap = st.backup(dir.path().join("backups"), 7).unwrap();
1026 let names: Vec<_> = fs::read_dir(snap.parent().unwrap())
1027 .unwrap()
1028 .map(|e| e.unwrap().file_name().into_string().unwrap())
1029 .collect();
1030 assert_eq!(names.len(), 1, "no -wal or -shm beside it: {names:?}");
1031 assert_eq!(header_mode(&snap), (1, 1));
1032
1033 let read =
1034 Connection::open_with_flags(&snap, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY).unwrap();
1035 assert_eq!(journal_mode(&read), "delete");
1036 let n: i64 = read
1037 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1038 .unwrap();
1039 assert_eq!(n, 20);
1040 }
1041
1042 fn rows_in_the_file_alone(db: &Path) -> i64 {
1047 let alone = tempfile::tempdir().unwrap();
1048 let copy = alone.path().join("recall.db");
1049 fs::copy(db, ©).unwrap();
1050 Connection::open(©)
1051 .unwrap()
1052 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1053 .unwrap()
1054 }
1055
1056 #[test]
1060 fn a_checkpoint_brings_the_file_on_its_own_up_to_date() {
1061 let dir = tempfile::tempdir().unwrap();
1062 let path = dir.path().join("recall.db");
1063 let st = Store::open(&path).unwrap();
1064 put(&st, "acme/app", "a.md", "x", "laptop");
1065 assert!(st.checkpoint_all().unwrap());
1066 assert_eq!(wal_len(&path), 0, "emptied");
1067 assert_eq!(rows_in_the_file_alone(&path), 1);
1068
1069 put(&st, "acme/app", "b.md", "x", "laptop");
1070 assert_eq!(rows_in_the_file_alone(&path), 1, "b.md is only in the WAL");
1071 assert!(st.checkpoint().unwrap(), "nothing held it back");
1072 assert_eq!(rows_in_the_file_alone(&path), 2);
1073
1074 put(&st, "acme/app", "c.md", "x", "laptop");
1075 assert!(st.checkpoint_all().unwrap());
1076 assert_eq!(wal_len(&path), 0, "emptied");
1077 assert_eq!(rows_in_the_file_alone(&path), 3);
1078 }
1079
1080 #[test]
1084 fn a_checkpoint_says_when_a_reader_held_it_back() {
1085 let dir = tempfile::tempdir().unwrap();
1086 let path = dir.path().join("recall.db");
1087 let st = Store::open(&path).unwrap();
1088 put(&st, "acme/app", "a.md", "x", "laptop");
1089 assert!(st.checkpoint_all().unwrap());
1090
1091 let reader = Connection::open(&path).unwrap();
1092 reader.execute_batch("BEGIN").unwrap();
1093 let _: i64 = reader
1094 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1095 .unwrap();
1096 put(&st, "acme/app", "b.md", "x", "laptop");
1097 assert!(!st.checkpoint().unwrap(), "held back by the reader");
1098 assert_eq!(rows_in_the_file_alone(&path), 1);
1099
1100 reader.execute_batch("COMMIT").unwrap();
1101 assert!(st.checkpoint().unwrap());
1102 assert_eq!(rows_in_the_file_alone(&path), 2);
1103 }
1104
1105 #[test]
1109 fn a_wal_grown_large_is_cut_back_to_the_limit() {
1110 let dir = tempfile::tempdir().unwrap();
1111 let path = dir.path().join("recall.db");
1112 let st = Store::open(&path).unwrap();
1113 let big = "x".repeat(24 * 1024 * 1024);
1114 put(&st, "acme/app", "big.md", &big, "laptop");
1115 let grown = wal_len(&path);
1116 assert!(grown > WAL_SIZE_LIMIT as u64, "{grown}");
1117 assert!(st.checkpoint().unwrap());
1118
1119 put(&st, "acme/app", "small.md", "x", "laptop");
1120 let now = wal_len(&path);
1121 assert!(now <= WAL_SIZE_LIMIT as u64, "{now} after {grown}");
1122 assert_eq!(st.get("acme/app", "big.md").unwrap().unwrap().content, big);
1123 }
1124
1125 #[cfg(unix)]
1132 #[test]
1133 fn a_database_moved_aside_under_a_running_store_is_not_written_to() {
1134 let dir = tempfile::tempdir().unwrap();
1135 let path = dir.path().join("recall.db");
1136 let st = Store::open(&path).unwrap();
1137 put(&st, "acme/app", "a.md", "x", "laptop");
1138 let snapshot = st.backup(dir.path().join("backups"), 7).unwrap();
1139
1140 let aside = dir.path().join("aside");
1141 fs::create_dir(&aside).unwrap();
1142 for f in ["recall.db", "recall.db-wal", "recall.db-shm"] {
1143 if dir.path().join(f).exists() {
1144 fs::rename(dir.path().join(f), aside.join(f)).unwrap();
1145 }
1146 }
1147 let refused = |st: &Store| {
1150 let err = st
1151 .upsert_audited("acme/app", "b.md", "y", "laptop", test_leaf)
1152 .unwrap_err();
1153 let said = format!("{err:#}");
1154 assert!(said.contains("moved or replaced"), "{said}");
1155 assert!(!said.contains(" "), "a run of spaces: {said:?}");
1156 let err = st.audit_append(test_leaf).unwrap_err();
1157 assert!(format!("{err:#}").contains("moved or replaced"), "{err:#}");
1158 };
1159 refused(&st);
1160 fs::copy(&snapshot, &path).unwrap();
1161 refused(&st);
1162 drop(st);
1163
1164 for db in [aside.join("recall.db"), path] {
1165 let files = Store::open(&db).unwrap().list("acme/app").unwrap();
1166 let names: Vec<_> = files.iter().map(|f| f.file_path.as_str()).collect();
1167 assert_eq!(names, ["a.md"], "{}", db.display());
1168 }
1169 }
1170
1171 #[tokio::test]
1180 #[ignore = "a measurement, not a check: run it by hand"]
1181 async fn push_latency_by_journal_mode() {
1182 use axum::body::Body;
1183 use axum::http::Request;
1184 use std::time::{Duration, Instant};
1185 use tower::ServiceExt;
1186
1187 const PUSHES: usize = 400;
1188 const TOKEN: &str = "bench-token";
1189 let content = "- a remembered fact, about as long as one usually is\n".repeat(20);
1190 println!("| journal | push p50 | push p90 | push mean | pull p50 | pull p90 |");
1191 println!("|---|---|---|---|---|---|");
1192 for (label, journal, sync) in [
1193 ("rollback (DELETE), FULL", "DELETE", "FULL"),
1194 ("WAL, NORMAL", "WAL", "NORMAL"),
1195 ("WAL, FULL", "WAL", "FULL"),
1196 ] {
1197 let dir = tempfile::tempdir().unwrap();
1198 let store = std::sync::Arc::new(Store::open(dir.path().join("recall.db")).unwrap());
1199 store
1200 .with_raw(|c| {
1201 c.query_row(&format!("PRAGMA journal_mode = {journal}"), [], |_| Ok(()))?;
1202 c.execute_batch(&format!("PRAGMA synchronous = {sync}"))
1203 })
1204 .unwrap();
1205 let server = crate::Server::new(
1206 crate::Config {
1207 token: TOKEN.into(),
1208 merge_enabled: false,
1209 rate_limit_max: 1_000_000,
1210 ..crate::Config::default()
1211 },
1212 store,
1213 );
1214 let router = server.router();
1215 let push = |i: usize| {
1216 let body = serde_json::json!({
1217 "project_key": "bench/app",
1218 "file_path": format!("f{i}.md"),
1219 "content": content,
1220 "source_env": "bench",
1221 });
1222 Request::post("/sync")
1223 .header("authorization", format!("Bearer {TOKEN}"))
1224 .header("content-type", "application/json")
1225 .body(Body::from(body.to_string()))
1226 .unwrap()
1227 };
1228 let pull = || {
1229 Request::get("/sync?project_key=small/app")
1230 .header("authorization", format!("Bearer {TOKEN}"))
1231 .body(Body::empty())
1232 .unwrap()
1233 };
1234 let time = |mut samples: Vec<Duration>| {
1235 samples.sort();
1236 let ms = |d: Duration| d.as_secs_f64() * 1000.0;
1237 let mean = ms(samples.iter().sum::<Duration>()) / samples.len() as f64;
1238 (
1239 ms(samples[samples.len() / 2]),
1240 ms(samples[samples.len() * 9 / 10]),
1241 mean,
1242 )
1243 };
1244 for i in 0..20 {
1245 let resp = router.clone().oneshot(push(PUSHES + i)).await.unwrap();
1246 assert_eq!(resp.status(), 200);
1247 }
1248 let mut pushes = Vec::with_capacity(PUSHES);
1249 for i in 0..PUSHES {
1250 let started = Instant::now();
1251 let resp = router.clone().oneshot(push(i)).await.unwrap();
1252 pushes.push(started.elapsed());
1253 assert_eq!(resp.status(), 200);
1254 }
1255 let mut pulls = Vec::with_capacity(PUSHES);
1256 for _ in 0..PUSHES {
1257 let started = Instant::now();
1258 let resp = router.clone().oneshot(pull()).await.unwrap();
1259 pulls.push(started.elapsed());
1260 assert_eq!(resp.status(), 200);
1261 }
1262 let (p50, p90, mean) = time(pushes);
1263 let (l50, l90, _) = time(pulls);
1264 println!(
1265 "| {label} | {p50:.2} ms | {p90:.2} ms | {mean:.2} ms | {l50:.2} ms | {l90:.2} ms |"
1266 );
1267 }
1268 }
1269}