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 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
40const 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#[derive(Debug, Clone, PartialEq, Eq)]
56pub struct Existing {
57 pub content: String,
60 pub deleted: bool,
62 pub source_env: String,
64 pub updated_at: String,
66}
67
68struct 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
93type FileId = (u64, u64);
112
113#[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
141pub struct Store {
143 state: Mutex<StoreState>,
150}
151
152impl Store {
153 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 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 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 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 state.conn.execute_batch(devices::SCHEMA)?;
238 state.conn.execute_batch(passkeys::SCHEMA)?;
239 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 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 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 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 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 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 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 {
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 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 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 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 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 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 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 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#[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
539fn 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
588const WAL_SIZE_LIMIT: i64 = 16 * 1024 * 1024;
596
597const 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
609fn 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
620fn 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
645fn nullable(s: &str) -> Option<&str> {
648 if s.is_empty() {
649 None
650 } else {
651 Some(s)
652 }
653}
654
655pub(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 #[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 #[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 #[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 #[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 #[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 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 assert!(
804 stamp[20..23].chars().all(|c| c.is_ascii_digit()),
805 "no millisecond field in {stamp}"
806 );
807 }
808
809 fn journal_mode(conn: &Connection) -> String {
812 conn.query_row("PRAGMA journal_mode", [], |r| r.get(0))
813 .unwrap()
814 }
815
816 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 #[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 #[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 #[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 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 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 #[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 #[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 #[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 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, ©).unwrap();
1038 Connection::open(©)
1039 .unwrap()
1040 .query_row("SELECT count(*) FROM memory_files", [], |r| r.get(0))
1041 .unwrap()
1042 }
1043
1044 #[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 #[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 #[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 #[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 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 #[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}