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 {
77 conn: Connection,
78 audit: Tree,
79 audit_at: String,
80}
81
82impl std::ops::Deref for StoreState {
83 type Target = Connection;
84 fn deref(&self) -> &Connection {
85 &self.conn
86 }
87}
88
89impl std::ops::DerefMut for StoreState {
90 fn deref_mut(&mut self) -> &mut Connection {
91 &mut self.conn
92 }
93}
94
95pub struct Store {
97 state: Mutex<StoreState>,
104}
105
106impl Store {
107 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
109 let path = path.as_ref();
110 if let Some(dir) = path.parent() {
111 if !dir.as_os_str().is_empty() {
112 fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
113 }
114 }
115 let conn = Connection::open(path).with_context(|| format!("opening {}", path.display()))?;
116 Self::with_connection(conn)
117 }
118
119 pub fn open_in_memory() -> Result<Self> {
121 Self::with_connection(Connection::open_in_memory()?)
122 }
123
124 fn with_connection(conn: Connection) -> Result<Self> {
125 let store = Self {
126 state: Mutex::new(StoreState {
127 conn,
128 audit: Tree::new(),
129 audit_at: String::new(),
130 }),
131 };
132 store.migrate()?;
133 Ok(store)
134 }
135
136 fn lock(&self) -> MutexGuard<'_, StoreState> {
140 self.state.lock().unwrap_or_else(PoisonError::into_inner)
141 }
142
143 fn migrate(&self) -> Result<()> {
144 let mut state = self.lock();
145 state.conn.execute_batch(SCHEMA)?;
146
147 let has_deleted = {
150 let mut stmt = state.conn.prepare("PRAGMA table_info(memory_files)")?;
151 let mut rows = stmt.query([])?;
152 let mut found = false;
153 while let Some(row) = rows.next()? {
154 if row.get::<_, String>(1)? == "deleted" {
155 found = true;
156 }
157 }
158 found
159 };
160 if !has_deleted {
161 state.conn.execute(
162 "ALTER TABLE memory_files ADD COLUMN deleted INTEGER NOT NULL DEFAULT 0",
163 [],
164 )?;
165 }
166
167 state.conn.execute_batch(devices::SCHEMA)?;
172 state.conn.execute_batch(passkeys::SCHEMA)?;
173 devices::allow_worker_scope(&state.conn)?;
178 state.conn.execute_batch(jobs::SCHEMA)?;
179 state.conn.execute_batch(audit::SCHEMA)?;
180
181 let loaded = audit::load(&state.conn)?;
182 state.audit = loaded.tree;
183 state.audit_at = loaded.last_at;
184 Ok(())
185 }
186
187 pub fn get(&self, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
189 read_file(&self.lock(), project_key, file_path)
190 }
191
192 pub fn upsert_audited(
197 &self,
198 project_key: &str,
199 file_path: &str,
200 content: &str,
201 source_env: &str,
202 build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
203 ) -> Result<String> {
204 self.audited(
205 |tx, at| {
206 write_file(tx, project_key, file_path, content, source_env, at)?;
207 Ok(Outcome::Commit(at.to_string()))
208 },
209 |seq, at, _| build_leaf(seq, at),
210 )
211 }
212
213 pub fn tombstone_audited(
223 &self,
224 project_key: &str,
225 file_path: &str,
226 source_env: &str,
227 build_leaf: impl FnOnce(u64, &str) -> Vec<u8>,
228 ) -> Result<String> {
229 self.audited(
230 |tx, at| {
231 tx.execute(
232 "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
233 VALUES (?1, ?2, '', ?3, ?4, 1)
234 ON CONFLICT(project_key, file_path) DO UPDATE SET
235 source_env = excluded.source_env,
236 updated_at = excluded.updated_at,
237 deleted = 1",
238 (project_key, file_path, nullable(source_env), at),
239 )?;
240 jobs::close_for_delete(tx, project_key, file_path, at)?;
241 Ok(Outcome::Commit(at.to_string()))
242 },
243 |seq, at, _| build_leaf(seq, at),
244 )
245 }
246
247 pub fn list(&self, project_key: &str) -> Result<Vec<File>> {
251 let conn = self.lock();
252 let mut stmt = conn.prepare(
253 "SELECT file_path, content, COALESCE(source_env, ''), updated_at, deleted
254 FROM memory_files WHERE project_key = ?1 ORDER BY file_path",
255 )?;
256 let rows = stmt.query_map((project_key,), |r| {
257 let content: String = r.get(1)?;
258 let deleted = r.get::<_, i64>(4)? != 0;
259 Ok(File {
260 file_path: r.get(0)?,
261 content: if deleted { None } else { Some(content) },
262 source_env: r.get(2)?,
263 updated_at: r.get(3)?,
264 deleted,
265 })
266 })?;
267 let mut files = Vec::new();
268 for row in rows {
269 files.push(row?);
270 }
271 Ok(files)
272 }
273
274 pub fn last_sync_at(&self) -> Result<String> {
277 let conn = self.lock();
278 let v: Option<String> =
279 conn.query_row("SELECT MAX(updated_at) FROM memory_files", [], |r| r.get(0))?;
280 Ok(v.unwrap_or_default())
281 }
282
283 pub fn admin_stats(&self) -> Result<(Vec<ProjectStats>, AdminTotals)> {
285 let conn = self.lock();
286
287 let mut projects = Vec::new();
288 let mut totals = AdminTotals::default();
289 {
290 let mut stmt = conn.prepare(
291 "SELECT project_key,
292 SUM(CASE WHEN deleted = 0 THEN 1 ELSE 0 END),
293 SUM(CASE WHEN deleted = 1 THEN 1 ELSE 0 END),
294 MAX(updated_at)
295 FROM memory_files GROUP BY project_key ORDER BY MAX(updated_at) DESC",
296 )?;
297 let rows = stmt.query_map([], |r| {
298 Ok(ProjectStats {
299 project_key: r.get(0)?,
300 file_count: r.get(1)?,
301 deleted_count: r.get(2)?,
302 sources: Vec::new(),
303 last_updated_at: r.get::<_, Option<String>>(3)?.unwrap_or_default(),
304 })
305 })?;
306 for row in rows {
307 let p = row?;
308 totals.file_count += p.file_count;
309 totals.deleted_count += p.deleted_count;
310 projects.push(p);
311 }
312 }
313 totals.project_count = projects.len() as i64;
314
315 {
320 let mut stmt = conn.prepare(
321 "SELECT DISTINCT project_key, source_env FROM memory_files WHERE source_env IS NOT NULL",
322 )?;
323 let rows =
324 stmt.query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?)))?;
325 for row in rows {
326 let (key, src) = row?;
327 if src.is_empty() {
328 continue;
329 }
330 if let Some(p) = projects.iter_mut().find(|p| p.project_key == key) {
331 p.sources.push(src);
332 }
333 }
334 }
335 for p in &mut projects {
336 p.sources.sort();
337 }
338 Ok((projects, totals))
339 }
340
341 pub fn backup(&self, dir: impl AsRef<Path>, keep: usize) -> Result<PathBuf> {
345 let dir = dir.as_ref();
346 fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
347
348 let stamp = now().replace([':', '.'], "-");
355 let dest = dir.join(format!("recall-{stamp}.db"));
356 let dest_str = dest
357 .to_str()
358 .context("backup path is not valid UTF-8")?
359 .to_owned();
360
361 let existed = dest.exists();
364 let vacuumed = {
365 let conn = self.lock();
366 conn.execute("VACUUM INTO ?1", (&dest_str,))
367 .with_context(|| format!("VACUUM INTO {dest_str}"))
368 };
369 if let Err(err) = vacuumed {
370 if !existed {
376 let _ = fs::remove_file(&dest);
377 }
378 return Err(err);
379 }
380
381 let mut snapshots: Vec<PathBuf> = fs::read_dir(dir)?
382 .filter_map(|e| e.ok())
383 .map(|e| e.path())
384 .filter(|p| {
385 p.file_name()
386 .and_then(|n| n.to_str())
387 .is_some_and(|n| n.starts_with("recall-") && n.ends_with(".db"))
388 })
389 .collect();
390 snapshots.sort();
391 for stale in snapshots.iter().take(snapshots.len().saturating_sub(keep)) {
392 let _ = fs::remove_file(stale);
393 }
394 Ok(dest)
395 }
396}
397
398#[cfg(test)]
399impl Store {
400 pub(crate) fn with_raw<T>(
404 &self,
405 f: impl FnOnce(&Connection) -> rusqlite::Result<T>,
406 ) -> rusqlite::Result<T> {
407 f(&self.lock())
408 }
409}
410
411#[cfg(test)]
414pub(crate) fn test_leaf(seq: u64, at: &str) -> Vec<u8> {
415 use crate::audit::leaf;
416 leaf::encode(
417 seq,
418 at,
419 leaf::action::START,
420 &leaf::Actor::Server,
421 leaf::subject_start("test"),
422 None,
423 )
424}
425
426const EXISTING_COLUMNS: &str = "content, deleted, COALESCE(source_env, ''), updated_at";
428
429fn existing_from(r: &rusqlite::Row<'_>) -> rusqlite::Result<Existing> {
430 Ok(Existing {
431 content: r.get(0)?,
432 deleted: r.get::<_, i64>(1)? != 0,
433 source_env: r.get(2)?,
434 updated_at: r.get(3)?,
435 })
436}
437
438fn read_file(conn: &Connection, project_key: &str, file_path: &str) -> Result<Option<Existing>> {
440 Ok(conn
441 .query_row(
442 &format!("SELECT {EXISTING_COLUMNS} FROM memory_files WHERE project_key = ?1 AND file_path = ?2"),
443 (project_key, file_path),
444 existing_from,
445 )
446 .optional()?)
447}
448
449fn write_file(
454 conn: &Connection,
455 project_key: &str,
456 file_path: &str,
457 content: &str,
458 source_env: &str,
459 updated_at: &str,
460) -> Result<()> {
461 conn.execute(
462 "INSERT INTO memory_files (project_key, file_path, content, source_env, updated_at, deleted)
463 VALUES (?1, ?2, ?3, ?4, ?5, 0)
464 ON CONFLICT(project_key, file_path) DO UPDATE SET
465 content = excluded.content,
466 source_env = excluded.source_env,
467 updated_at = excluded.updated_at,
468 deleted = 0",
469 (project_key, file_path, content, nullable(source_env), updated_at),
470 )?;
471 Ok(())
472}
473
474fn nullable(s: &str) -> Option<&str> {
477 if s.is_empty() {
478 None
479 } else {
480 Some(s)
481 }
482}
483
484pub(crate) mod admin;
489
490#[cfg(test)]
491mod tests {
492 use super::*;
493
494 fn store() -> Store {
495 Store::open_in_memory().unwrap()
496 }
497
498 fn put(st: &Store, project_key: &str, file_path: &str, content: &str, source_env: &str) {
499 st.upsert_audited(project_key, file_path, content, source_env, test_leaf)
500 .unwrap();
501 }
502
503 fn del(st: &Store, project_key: &str, file_path: &str, source_env: &str) {
504 st.tombstone_audited(project_key, file_path, source_env, test_leaf)
505 .unwrap();
506 }
507
508 #[test]
511 fn a_write_is_stamped_with_its_leafs_at() {
512 let st = store();
513 let mut leaf_at = String::new();
514 let updated_at = st
515 .upsert_audited("acme/app", "a.md", "x", "laptop", |seq, at| {
516 leaf_at = at.to_string();
517 test_leaf(seq, at)
518 })
519 .unwrap();
520 assert_eq!(updated_at, leaf_at);
521 assert_eq!(st.list("acme/app").unwrap()[0].updated_at, updated_at);
522 assert_eq!(st.audit_checkpoint().0, 1);
523 }
524
525 #[test]
526 fn upsert_get_and_list_round_trip() {
527 let st = store();
528 put(&st, "acme/app", "MEMORY.md", "hello", "laptop");
529
530 let got = st.get("acme/app", "MEMORY.md").unwrap().unwrap();
531 assert_eq!(got.content, "hello");
532 assert!(!got.deleted);
533
534 let files = st.list("acme/app").unwrap();
535 assert_eq!(files.len(), 1);
536 assert_eq!(files[0].content.as_deref(), Some("hello"));
537 assert_eq!(files[0].source_env, "laptop");
538 assert!(st.get("acme/app", "missing.md").unwrap().is_none());
539 }
540
541 #[test]
544 fn tombstone_preserves_content_but_list_withholds_it() {
545 let st = store();
546 put(&st, "acme/app", "gone.md", "secret", "laptop");
547 del(&st, "acme/app", "gone.md", "laptop");
548
549 let row = st.get("acme/app", "gone.md").unwrap().unwrap();
550 assert_eq!(row.content, "secret", "content must stay recoverable");
551 assert!(row.deleted);
552
553 let files = st.list("acme/app").unwrap();
554 assert_eq!(
555 files.len(),
556 1,
557 "tombstones are listed so clients can delete locally"
558 );
559 assert!(files[0].deleted);
560 assert_eq!(files[0].content, None, "a pull must not resurrect it");
561 }
562
563 #[test]
565 fn upsert_clears_a_tombstone() {
566 let st = store();
567 del(&st, "acme/app", "f.md", "laptop");
568 put(&st, "acme/app", "f.md", "back", "laptop");
569 let row = st.get("acme/app", "f.md").unwrap().unwrap();
570 assert!(!row.deleted);
571 assert_eq!(row.content, "back");
572 }
573
574 #[test]
575 fn last_sync_at_is_empty_on_a_fresh_database() {
576 assert_eq!(store().last_sync_at().unwrap(), "");
577 }
578
579 #[test]
582 fn admin_stats_keeps_commas_inside_a_source_env() {
583 let st = store();
584 put(&st, "acme/app", "a.md", "x", "laptop,evil");
585 let (projects, _) = st.admin_stats().unwrap();
586 assert_eq!(projects[0].sources, vec!["laptop,evil".to_string()]);
587 }
588
589 #[test]
592 fn migrates_a_database_that_predates_tombstones() {
593 let dir = tempfile::tempdir().unwrap();
594 let path = dir.path().join("old.db");
595 {
596 let conn = Connection::open(&path).unwrap();
597 conn.execute_batch(
598 "CREATE TABLE memory_files (
599 project_key TEXT NOT NULL,
600 file_path TEXT NOT NULL,
601 content TEXT NOT NULL,
602 source_env TEXT,
603 updated_at TEXT NOT NULL,
604 PRIMARY KEY (project_key, file_path)
605 );
606 INSERT INTO memory_files VALUES ('acme/app','old.md','kept','node-era','2026-09-03T21:49:55.191Z');",
607 )
608 .unwrap();
609 }
610 let st = Store::open(&path).unwrap();
611 let files = st.list("acme/app").unwrap();
612 assert_eq!(files.len(), 1);
613 assert_eq!(files[0].content.as_deref(), Some("kept"));
614 assert!(!files[0].deleted);
615 }
616
617 #[test]
618 fn backup_names_carry_milliseconds() {
619 let dir = tempfile::tempdir().unwrap();
620 let st = store();
621 let dest = st.backup(dir.path(), 7).unwrap();
622 let name = dest.file_name().unwrap().to_str().unwrap();
623 assert!(
625 name.starts_with("recall-") && name.ends_with("Z.db"),
626 "got {name}"
627 );
628 let stamp = &name["recall-".len()..name.len() - ".db".len()];
629 assert_eq!(stamp.len(), 24, "got {stamp}");
630 assert!(
633 stamp[20..23].chars().all(|c| c.is_ascii_digit()),
634 "no millisecond field in {stamp}"
635 );
636 }
637}