1mod registry;
31
32use std::fmt;
33use std::ops::Deref;
34use std::path::{Path, PathBuf};
35use std::sync::atomic::{AtomicU64, Ordering};
36use std::sync::{Mutex, MutexGuard};
37
38use crate::{Database, HostError};
39
40pub use registry::{
41 ARCHIVED_TAG, DbEntry, Description, ENTRY_TAG, ReindexReport, SELF_ENTITY, WorkspaceIssue,
42};
43
44pub const MAX_DB_NAME: usize = 64;
51
52const DB_DIR: &str = "db";
54
55const REGISTRY_FILE: &str = "registry.plugmem";
57
58const DB_EXT: &str = "plugmem";
62
63const RESERVED_DEVICE_NAMES: &[&str] = &[
81 "con", "prn", "aux", "nul", "com1", "com2", "com3", "com4", "com5", "com6", "com7", "com8",
82 "com9", "lpt1", "lpt2", "lpt3", "lpt4", "lpt5", "lpt6", "lpt7", "lpt8", "lpt9",
83];
84
85#[derive(Clone, Copy, Debug, PartialEq, Eq)]
90#[non_exhaustive]
91pub enum NameProblem {
92 Empty,
94 TooLong,
96 LeadingChar,
100 Character,
102 ReservedDevice,
105}
106
107impl fmt::Display for NameProblem {
108 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
109 match self {
110 Self::Empty => f.write_str("it is empty"),
111 Self::TooLong => write!(f, "it is longer than {MAX_DB_NAME} bytes"),
112 Self::LeadingChar => f.write_str("it must start with a lowercase letter or a digit"),
113 Self::Character => {
114 f.write_str("it may hold only lowercase letters, digits, '-' and '_'")
115 }
116 Self::ReservedDevice => f.write_str(
117 "it is a Windows device name, which would open a device rather than a file there",
118 ),
119 }
120 }
121}
122
123#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
129pub struct DbName(String);
130
131impl DbName {
132 pub fn parse(s: &str) -> Result<Self, WorkspaceError> {
141 let bad = |why| {
142 Err(WorkspaceError::BadName {
143 name: s.to_string(),
144 why,
145 })
146 };
147 let Some(&first) = s.as_bytes().first() else {
148 return bad(NameProblem::Empty);
149 };
150 if !first.is_ascii_lowercase() && !first.is_ascii_digit() {
151 return bad(NameProblem::LeadingChar);
152 }
153 if !s.bytes().all(is_name_byte) {
154 return bad(NameProblem::Character);
155 }
156 if s.len() > MAX_DB_NAME {
157 return bad(NameProblem::TooLong);
158 }
159 if RESERVED_DEVICE_NAMES.contains(&s) {
160 return bad(NameProblem::ReservedDevice);
161 }
162 Ok(DbName(s.to_string()))
163 }
164
165 pub fn as_str(&self) -> &str {
167 &self.0
168 }
169}
170
171impl fmt::Display for DbName {
172 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
173 f.write_str(&self.0)
174 }
175}
176
177fn is_name_byte(b: u8) -> bool {
179 b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-' || b == b'_'
180}
181
182#[derive(Debug, thiserror::Error)]
187#[non_exhaustive]
188pub enum WorkspaceError {
189 #[error("{name:?} is not a usable database name: {why}")]
191 BadName {
192 name: String,
194 why: NameProblem,
196 },
197
198 #[error("no database named {name} in this workspace (looked for {})", path.display())]
200 NoSuchDatabase {
201 name: DbName,
203 path: PathBuf,
205 },
206
207 #[error(
214 "database {name} is in use by another process; it is released once that process closes it (a pooled handle does so after its idle timeout)"
215 )]
216 Busy {
217 name: DbName,
219 },
220
221 #[error(
224 "workspace has {max_open} active databases (the max_open limit); retry after one call finishes or raise the limit"
225 )]
226 AtCapacity {
227 max_open: usize,
229 },
230
231 #[error("database {name} is in use by an active workspace operation")]
234 InUse {
235 name: DbName,
237 },
238
239 #[error("i/o on {}: {source}", path.display())]
241 Io {
242 path: PathBuf,
244 #[source]
246 source: std::io::Error,
247 },
248
249 #[error(transparent)]
251 Host(#[from] HostError),
252}
253
254impl WorkspaceError {
255 pub(crate) fn io(path: &Path, source: std::io::Error) -> Self {
257 Self::Io {
258 path: path.to_path_buf(),
259 source,
260 }
261 }
262}
263
264#[derive(Clone, Debug, PartialEq, Eq)]
270pub struct WorkspaceLayout {
271 root: PathBuf,
272}
273
274impl WorkspaceLayout {
275 pub fn new(root: impl Into<PathBuf>) -> Self {
278 Self { root: root.into() }
279 }
280
281 pub fn root(&self) -> &Path {
283 &self.root
284 }
285
286 pub fn db_dir(&self) -> PathBuf {
288 self.root.join(DB_DIR)
289 }
290
291 pub fn path_of(&self, name: &DbName) -> PathBuf {
296 self.db_dir().join(format!("{}.{DB_EXT}", name.0))
297 }
298
299 pub fn registry_path(&self) -> PathBuf {
302 self.root.join(REGISTRY_FILE)
303 }
304
305 pub fn exists(&self, name: &DbName) -> bool {
311 crate::storage::database_exists(&self.path_of(name))
312 }
313
314 pub fn list(&self) -> Result<Vec<DbName>, WorkspaceError> {
331 let dir = self.db_dir();
332 let entries = match std::fs::read_dir(&dir) {
333 Ok(entries) => entries,
334 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
335 Err(e) => return Err(WorkspaceError::io(&dir, e)),
336 };
337 let mut names = Vec::new();
338 for entry in entries {
339 let entry = entry.map_err(|e| WorkspaceError::io(&dir, e))?;
340 let file_name = entry.file_name();
341 let Some(candidate) = file_name.to_str().and_then(|n| n.split('.').next()) else {
342 continue;
343 };
344 if let Ok(name) = DbName::parse(candidate)
345 && !names.contains(&name)
346 && self.exists(&name)
347 {
348 names.push(name);
349 }
350 }
351 names.sort_unstable();
352 Ok(names)
353 }
354}
355
356#[derive(Clone, Copy, Debug, PartialEq, Eq)]
358pub struct WorkspaceLimits {
359 pub max_open: usize,
363 pub idle_timeout_ms: u64,
371}
372
373pub const DEFAULT_MAX_OPEN: usize = 16;
377
378pub const DEFAULT_IDLE_TIMEOUT_MS: u64 = 60_000;
380
381const FILES_PER_OPEN_DATABASE: usize = 4;
384
385const ASSUMED_FD_LIMIT: usize = 1024;
388
389const RESERVED_FDS: usize = 64;
392
393pub const MAX_OPEN_CEILING: usize = (ASSUMED_FD_LIMIT - RESERVED_FDS) / FILES_PER_OPEN_DATABASE;
401
402const _: () = {
403 assert!(MAX_OPEN_CEILING > DEFAULT_MAX_OPEN);
404};
405
406impl Default for WorkspaceLimits {
407 fn default() -> Self {
408 Self {
409 max_open: DEFAULT_MAX_OPEN,
410 idle_timeout_ms: DEFAULT_IDLE_TIMEOUT_MS,
411 }
412 }
413}
414
415impl WorkspaceLimits {
416 pub fn ceiling(&self) -> usize {
420 self.max_open.clamp(1, MAX_OPEN_CEILING)
421 }
422}
423
424#[derive(Clone, Copy, Debug, PartialEq, Eq)]
426pub enum IfMissing {
427 Create,
430 Fail,
433}
434
435pub type Opener = Box<dyn Fn(&Path) -> Result<Database, HostError> + Send + Sync>;
439
440struct Pooled {
442 name: DbName,
443 db: Database,
444 last_used_ms: u64,
445 active: usize,
448 token: u64,
450}
451
452pub struct WorkspaceLease<'a> {
464 workspace: &'a Workspace,
465 db: Option<Database>,
466 token: u64,
467}
468
469impl Deref for WorkspaceLease<'_> {
470 type Target = Database;
471
472 fn deref(&self) -> &Self::Target {
473 self.db
474 .as_ref()
475 .expect("a live workspace lease always owns its database")
476 }
477}
478
479impl Drop for WorkspaceLease<'_> {
480 fn drop(&mut self) {
481 drop(self.db.take());
485 let mut pool = self.workspace.pooled();
486 if let Some(slot) = pool.iter_mut().find(|p| p.token == self.token) {
487 debug_assert!(slot.active > 0);
488 slot.active = slot.active.saturating_sub(1);
489 }
490 }
491}
492
493pub struct Workspace {
530 layout: WorkspaceLayout,
531 open: Opener,
532 limits: WorkspaceLimits,
533 pool: Mutex<Vec<Pooled>>,
534 registry: Mutex<Option<Database>>,
535 next_token: AtomicU64,
536}
537
538impl Workspace {
539 pub fn new(layout: WorkspaceLayout, open: Opener, limits: WorkspaceLimits) -> Self {
541 Self {
542 layout,
543 open,
544 limits,
545 pool: Mutex::new(Vec::new()),
546 registry: Mutex::new(None),
547 next_token: AtomicU64::new(1),
548 }
549 }
550
551 pub fn registry(&self) -> Result<Database, WorkspaceError> {
569 let mut slot = self.registry.lock().unwrap_or_else(|e| e.into_inner());
570 if let Some(db) = slot.as_ref() {
571 return Ok(db.clone());
572 }
573 let root = self.layout.root();
574 std::fs::create_dir_all(root).map_err(|e| WorkspaceError::io(root, e))?;
575 let db = (self.open)(&self.layout.registry_path())?;
576 *slot = Some(db.clone());
577 Ok(db)
578 }
579
580 pub fn close_registry(&self) -> bool {
584 self.registry
585 .lock()
586 .unwrap_or_else(|e| e.into_inner())
587 .take()
588 .is_some()
589 }
590
591 pub fn layout(&self) -> &WorkspaceLayout {
593 &self.layout
594 }
595
596 pub fn limits(&self) -> WorkspaceLimits {
598 self.limits
599 }
600
601 pub fn open_count(&self) -> usize {
604 self.pooled().len()
605 }
606
607 pub fn get(
620 &self,
621 name: &DbName,
622 now_ms: u64,
623 missing: IfMissing,
624 ) -> Result<Database, WorkspaceError> {
625 self.acquire(name, now_ms, missing, false).map(|(db, _)| db)
626 }
627
628 pub fn lease(
640 &self,
641 name: &DbName,
642 now_ms: u64,
643 missing: IfMissing,
644 ) -> Result<WorkspaceLease<'_>, WorkspaceError> {
645 let (db, token) = self.acquire(name, now_ms, missing, true)?;
646 Ok(WorkspaceLease {
647 workspace: self,
648 db: Some(db),
649 token,
650 })
651 }
652
653 fn acquire(
655 &self,
656 name: &DbName,
657 now_ms: u64,
658 missing: IfMissing,
659 pin: bool,
660 ) -> Result<(Database, u64), WorkspaceError> {
661 let mut pool = self.pooled();
666
667 if let Some(slot) = pool.iter_mut().find(|p| &p.name == name) {
668 slot.last_used_ms = now_ms;
669 if pin {
670 slot.active = slot.active.saturating_add(1);
671 }
672 return Ok((slot.db.clone(), slot.token));
673 }
674
675 let path = self.layout.path_of(name);
676 if !self.layout.exists(name) {
677 if missing == IfMissing::Fail {
678 return Err(WorkspaceError::NoSuchDatabase {
679 name: name.clone(),
680 path,
681 });
682 }
683 let dir = self.layout.db_dir();
684 std::fs::create_dir_all(&dir).map_err(|e| WorkspaceError::io(&dir, e))?;
685 }
686
687 let ceiling = self.limits.ceiling();
689 while pool.len() >= ceiling {
690 let Some(lru) = Self::least_recently_used_available(&pool) else {
691 return Err(WorkspaceError::AtCapacity { max_open: ceiling });
692 };
693 pool.remove(lru);
694 }
695
696 let db = (self.open)(&path).map_err(|e| match e {
697 HostError::Locked { .. } => WorkspaceError::Busy { name: name.clone() },
698 other => WorkspaceError::Host(other),
699 })?;
700 let token = self.next_token.fetch_add(1, Ordering::Relaxed);
701 pool.push(Pooled {
702 name: name.clone(),
703 db: db.clone(),
704 last_used_ms: now_ms,
705 active: usize::from(pin),
706 token,
707 });
708 Ok((db, token))
709 }
710
711 pub fn close_idle(&self, now_ms: u64) -> usize {
717 let timeout = self.limits.idle_timeout_ms;
718 if timeout == 0 {
719 return 0;
720 }
721 let closing: Vec<Pooled> = self
728 .pooled()
729 .extract_if(.., |p| {
730 p.active == 0 && now_ms.saturating_sub(p.last_used_ms) >= timeout
731 })
732 .collect();
733 closing.len()
734 }
735
736 pub fn release(&self, name: &DbName) -> Result<bool, WorkspaceError> {
742 let closing = {
743 let mut pool = self.pooled();
744 let Some(index) = pool.iter().position(|p| &p.name == name) else {
745 return Ok(false);
746 };
747 if pool[index].active > 0 {
748 return Err(WorkspaceError::InUse { name: name.clone() });
749 }
750 pool.remove(index)
751 };
752 drop(closing);
753 Ok(true)
754 }
755
756 pub fn close_all(&self) -> usize {
759 let closing = std::mem::take(&mut *self.pooled());
762 closing.len()
763 }
764
765 fn pooled(&self) -> MutexGuard<'_, Vec<Pooled>> {
769 self.pool.lock().unwrap_or_else(|e| e.into_inner())
770 }
771
772 fn least_recently_used_available(pool: &[Pooled]) -> Option<usize> {
774 let mut oldest: Option<usize> = None;
775 for (i, p) in pool.iter().enumerate() {
776 let is_older = match oldest {
777 Some(candidate) => p.last_used_ms < pool[candidate].last_used_ms,
778 None => true,
779 };
780 if p.active == 0 && is_older {
781 oldest = Some(i);
782 }
783 }
784 oldest
785 }
786}
787
788impl fmt::Debug for Workspace {
789 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
790 f.debug_struct("Workspace")
791 .field("root", &self.layout.root())
792 .field("limits", &self.limits)
793 .field("open", &self.open_count())
794 .finish()
795 }
796}
797
798#[cfg(test)]
800pub(crate) mod testkit {
801 use std::path::PathBuf;
802 use std::sync::Arc;
803 use std::sync::atomic::{AtomicUsize, Ordering};
804
805 use super::{DbName, Opener, Workspace, WorkspaceLayout, WorkspaceLimits};
806 use crate::Database;
807
808 pub(crate) struct TempDir(pub PathBuf);
810
811 impl TempDir {
812 pub(crate) fn new(tag: &str) -> Self {
813 let dir = std::env::temp_dir().join(format!(
814 "plugmem-workspace-{tag}-{}-{}",
815 std::process::id(),
816 std::time::SystemTime::now()
817 .duration_since(std::time::UNIX_EPOCH)
818 .unwrap()
819 .as_nanos()
820 ));
821 std::fs::create_dir_all(&dir).unwrap();
822 TempDir(dir)
823 }
824 }
825
826 impl Drop for TempDir {
827 fn drop(&mut self) {
828 let _ = std::fs::remove_dir_all(&self.0);
829 }
830 }
831
832 pub(crate) fn workspace(
835 tmp: &TempDir,
836 limits: WorkspaceLimits,
837 ) -> (Workspace, Arc<AtomicUsize>) {
838 let opens = Arc::new(AtomicUsize::new(0));
839 let counted = Arc::clone(&opens);
840 let open: Opener = Box::new(move |path: &std::path::Path| {
841 counted.fetch_add(1, Ordering::SeqCst);
842 Ok(Database::open(path, crate::Config::default())?.0)
843 });
844 (
845 Workspace::new(WorkspaceLayout::new(&tmp.0), open, limits),
846 opens,
847 )
848 }
849
850 pub(crate) fn name(s: &str) -> DbName {
852 DbName::parse(s).unwrap()
853 }
854}
855
856#[cfg(test)]
857mod tests {
858 use std::sync::atomic::Ordering;
859
860 use super::testkit::{TempDir, name, workspace};
861 use super::*;
862
863 fn problem(s: &str) -> NameProblem {
864 match DbName::parse(s) {
865 Err(WorkspaceError::BadName { why, .. }) => why,
866 other => panic!("expected {s:?} to be refused, got {other:?}"),
867 }
868 }
869
870 #[test]
871 fn a_name_admits_only_the_safe_alphabet() {
872 for ok in [
873 "a",
874 "0",
875 "chat-42",
876 "common",
877 "x_y-9",
878 &"a".repeat(MAX_DB_NAME),
879 ] {
880 assert_eq!(DbName::parse(ok).unwrap().as_str(), ok);
881 }
882
883 assert_eq!(problem(""), NameProblem::Empty);
884 assert_eq!(problem(&"a".repeat(MAX_DB_NAME + 1)), NameProblem::TooLong);
885
886 for bad in ["-x", "_x", ".x", "..", "/x", "Ab", "чат"] {
891 assert_eq!(problem(bad), NameProblem::LeadingChar, "{bad:?}");
892 }
893
894 for bad in ["a/b", "a\\b", "a.b", "a b", "aB", "a:b", "aчат", "a\0b"] {
897 assert_eq!(problem(bad), NameProblem::Character, "{bad:?}");
898 }
899
900 for bad in ["con", "nul", "prn", "aux", "com1", "lpt9"] {
905 assert_eq!(problem(bad), NameProblem::ReservedDevice, "{bad:?}");
906 }
907 for ok in ["console", "con1", "com0", "com10", "nula"] {
909 assert!(DbName::parse(ok).is_ok(), "{ok:?}");
910 }
911 }
912
913 #[test]
914 fn a_name_prints_as_itself() {
915 assert_eq!(DbName::parse("chat-42").unwrap().to_string(), "chat-42");
916 assert_eq!(NameProblem::Empty.to_string(), "it is empty");
917 assert_eq!(
918 NameProblem::TooLong.to_string(),
919 format!("it is longer than {MAX_DB_NAME} bytes")
920 );
921 assert!(NameProblem::LeadingChar.to_string().contains("start with"));
922 assert!(NameProblem::Character.to_string().contains("lowercase"));
923 assert!(
924 NameProblem::ReservedDevice
925 .to_string()
926 .contains("Windows device name")
927 );
928 }
929
930 #[test]
931 fn the_layout_puts_the_registry_out_of_reach_of_names() {
932 let layout = WorkspaceLayout::new("/ws");
933 let name = DbName::parse("chat-42").unwrap();
934
935 assert_eq!(layout.root(), Path::new("/ws"));
936 assert_eq!(layout.db_dir(), Path::new("/ws/db"));
937 assert_eq!(layout.path_of(&name), Path::new("/ws/db/chat-42.plugmem"));
938 assert_eq!(layout.registry_path(), Path::new("/ws/registry.plugmem"));
939
940 let lookalike = DbName::parse("registry").unwrap();
943 assert_ne!(layout.path_of(&lookalike), layout.registry_path());
944 }
945
946 #[test]
947 fn listing_reads_the_directory_and_ignores_what_is_not_a_database() {
948 let tmp = TempDir::new("list");
949 let layout = WorkspaceLayout::new(&tmp.0);
950
951 assert!(layout.list().unwrap().is_empty());
953
954 std::fs::create_dir_all(layout.db_dir()).unwrap();
955 for file in [
956 "chat-42.plugmem",
957 "common.plugmem",
958 "chat-42.plugmem.lock",
961 "chat-42.plugmem.journal",
962 "chat-42.plugmem.snap.3",
963 "notes.txt",
965 "Chat-43.plugmem",
966 ] {
967 std::fs::write(layout.db_dir().join(file), b"").unwrap();
968 }
969
970 let names: Vec<String> = layout
971 .list()
972 .unwrap()
973 .iter()
974 .map(DbName::to_string)
975 .collect();
976 assert_eq!(names, ["chat-42", "common"]);
977
978 assert!(layout.exists(&DbName::parse("chat-42").unwrap()));
979 assert!(!layout.exists(&DbName::parse("nope").unwrap()));
980 }
981
982 #[test]
983 fn a_database_exists_before_its_first_checkpoint() {
984 let tmp = TempDir::new("list-uncheckpointed");
985 let layout = WorkspaceLayout::new(&tmp.0);
986 let fresh = DbName::parse("fresh").unwrap();
987 std::fs::create_dir_all(layout.db_dir()).unwrap();
988
989 let (db, _) = Database::open(layout.path_of(&fresh), crate::Config::default()).unwrap();
993 db.remember(crate::RememberInput::text(1_000, "not yet checkpointed"))
994 .unwrap();
995 assert!(!layout.path_of(&fresh).exists());
996 assert!(layout.exists(&fresh));
997 assert_eq!(layout.list().unwrap(), [fresh]);
998 }
999
1000 #[test]
1001 fn an_unreadable_directory_is_an_error_not_an_empty_workspace() {
1002 let tmp = TempDir::new("list-io");
1003 let layout = WorkspaceLayout::new(&tmp.0);
1004 std::fs::write(layout.db_dir(), b"not a directory").unwrap();
1007 assert!(matches!(layout.list(), Err(WorkspaceError::Io { .. })));
1008 }
1009
1010 #[test]
1011 fn every_failure_names_what_the_caller_typed() {
1012 let busy = WorkspaceError::Busy {
1013 name: DbName::parse("chat-42").unwrap(),
1014 };
1015 assert!(busy.to_string().contains("chat-42"));
1016
1017 let host = WorkspaceError::from(HostError::Embed("no".into()));
1018 assert!(matches!(host, WorkspaceError::Host(HostError::Embed(_))));
1019
1020 let io = WorkspaceError::io(
1021 Path::new("/ws"),
1022 std::io::Error::new(std::io::ErrorKind::PermissionDenied, "denied"),
1023 );
1024 assert!(io.to_string().contains("/ws"));
1025
1026 let missing = WorkspaceError::NoSuchDatabase {
1027 name: DbName::parse("gone").unwrap(),
1028 path: PathBuf::from("/ws/db/gone.plugmem"),
1029 };
1030 assert!(missing.to_string().contains("gone"));
1031 }
1032
1033 #[test]
1034 fn a_pooled_database_is_reused_and_a_missing_one_is_created_only_on_request() {
1035 let tmp = TempDir::new("pool-reuse");
1036 let (ws, opens) = workspace(&tmp, WorkspaceLimits::default());
1037 let chat = name("chat-42");
1038
1039 let missed = ws.get(&chat, 1_000, IfMissing::Fail).unwrap_err();
1042 assert!(
1043 matches!(&missed, WorkspaceError::NoSuchDatabase { name, .. } if name == &chat),
1044 "{missed}"
1045 );
1046 assert_eq!(opens.load(Ordering::SeqCst), 0);
1047 assert!(!ws.layout().db_dir().exists());
1048
1049 let db = ws.get(&chat, 1_000, IfMissing::Create).unwrap();
1051 db.remember(crate::RememberInput::text(1_000, "prefers tokio"))
1052 .unwrap();
1053 assert!(ws.layout().exists(&chat));
1054 assert_eq!(ws.open_count(), 1);
1055
1056 let again = ws.get(&chat, 2_000, IfMissing::Fail).unwrap();
1058 assert_eq!(again.stats().facts, 1);
1059 assert_eq!(opens.load(Ordering::SeqCst), 1);
1060
1061 assert!(format!("{ws:?}").contains("open: 1"));
1062 }
1063
1064 #[test]
1065 fn databases_in_one_workspace_do_not_see_each_other() {
1066 let tmp = TempDir::new("isolation");
1067 let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
1068
1069 for (db, text) in [
1070 ("chat-42", "the sky is blue"),
1071 ("chat-43", "the sky is red"),
1072 ] {
1073 ws.get(&name(db), 1_000, IfMissing::Create)
1074 .unwrap()
1075 .remember(crate::RememberInput::text(1_000, text))
1076 .unwrap();
1077 }
1078
1079 for (db, expected) in [("chat-42", "blue"), ("chat-43", "red")] {
1082 let out = ws
1083 .get(&name(db), 2_000, IfMissing::Fail)
1084 .unwrap()
1085 .recall(crate::RecallQuery::text(2_000, "sky"))
1086 .unwrap();
1087 assert_eq!(out.facts.len(), 1, "{db}");
1088 assert!(out.rendered.contains(expected), "{db}: {}", out.rendered);
1089 }
1090 }
1091
1092 #[test]
1093 fn the_ceiling_evicts_the_least_recently_used() {
1094 let tmp = TempDir::new("pool-evict");
1095 let (ws, opens) = workspace(
1096 &tmp,
1097 WorkspaceLimits {
1098 max_open: 2,
1099 ..WorkspaceLimits::default()
1100 },
1101 );
1102
1103 ws.get(&name("a"), 1_000, IfMissing::Create).unwrap();
1106 ws.get(&name("b"), 2_000, IfMissing::Create).unwrap();
1107 ws.get(&name("a"), 3_000, IfMissing::Fail).unwrap();
1108 ws.get(&name("c"), 4_000, IfMissing::Create).unwrap();
1109 assert_eq!(ws.open_count(), 2);
1110 assert_eq!(opens.load(Ordering::SeqCst), 3);
1111
1112 ws.get(&name("a"), 5_000, IfMissing::Fail).unwrap();
1114 assert_eq!(opens.load(Ordering::SeqCst), 3);
1115 ws.get(&name("b"), 6_000, IfMissing::Fail).unwrap();
1116 assert_eq!(opens.load(Ordering::SeqCst), 4);
1117 }
1118
1119 #[test]
1120 fn a_ceiling_of_zero_still_serves_one_database() {
1121 let tmp = TempDir::new("pool-zero");
1122 let (ws, opens) = workspace(
1123 &tmp,
1124 WorkspaceLimits {
1125 max_open: 0,
1126 idle_timeout_ms: 0,
1127 },
1128 );
1129 ws.get(&name("a"), 1_000, IfMissing::Create).unwrap();
1130 ws.get(&name("b"), 2_000, IfMissing::Create).unwrap();
1131 assert_eq!(ws.open_count(), 1);
1132 assert_eq!(opens.load(Ordering::SeqCst), 2);
1133
1134 assert_eq!(ws.close_idle(u64::MAX), 0);
1136 assert_eq!(ws.open_count(), 1);
1137 }
1138
1139 #[test]
1140 fn an_idle_database_is_closed_and_its_lock_released() {
1141 let tmp = TempDir::new("pool-idle");
1142 let (ws, _) = workspace(
1143 &tmp,
1144 WorkspaceLimits {
1145 max_open: 8,
1146 idle_timeout_ms: 1_000,
1147 },
1148 );
1149 let chat = name("chat-42");
1150 let path = ws.layout().path_of(&chat);
1151 drop(ws.get(&chat, 1_000, IfMissing::Create).unwrap());
1152
1153 assert_eq!(ws.close_idle(1_500), 0);
1155 assert!(matches!(
1156 Database::open(&path, crate::Config::default()),
1157 Err(HostError::Locked { .. })
1158 ));
1159
1160 assert_eq!(ws.close_idle(500), 0);
1162
1163 assert_eq!(ws.close_idle(2_000), 1);
1166 assert_eq!(ws.open_count(), 0);
1167 assert!(Database::open(&path, crate::Config::default()).is_ok());
1168 }
1169
1170 #[test]
1171 fn a_scoped_lease_cannot_be_swept_released_or_evicted() {
1172 let tmp = TempDir::new("pool-lease-pin");
1173 let (ws, opens) = workspace(
1174 &tmp,
1175 WorkspaceLimits {
1176 max_open: 1,
1177 idle_timeout_ms: 1,
1178 },
1179 );
1180 let a = name("a");
1181 let b = name("b");
1182 let path = ws.layout().path_of(&a);
1183 let lease = ws.lease(&a, 1_000, IfMissing::Create).unwrap();
1184
1185 assert_eq!(ws.close_idle(u64::MAX), 0);
1188 assert!(matches!(
1189 ws.release(&a),
1190 Err(WorkspaceError::InUse { name }) if name == a
1191 ));
1192 assert!(matches!(
1193 ws.lease(&b, 2_000, IfMissing::Create),
1194 Err(WorkspaceError::AtCapacity { max_open: 1 })
1195 ));
1196 assert_eq!(opens.load(Ordering::SeqCst), 1);
1197 assert!(matches!(
1198 Database::open(&path, crate::Config::default()),
1199 Err(HostError::Locked { .. })
1200 ));
1201
1202 drop(lease);
1203 assert!(ws.release(&a).unwrap());
1204 assert!(!ws.release(&a).unwrap());
1205 assert!(Database::open(&path, crate::Config::default()).is_ok());
1206 }
1207
1208 #[test]
1209 fn lru_evicts_an_inactive_entry_instead_of_an_active_one() {
1210 let tmp = TempDir::new("pool-lease-lru");
1211 let (ws, opens) = workspace(
1212 &tmp,
1213 WorkspaceLimits {
1214 max_open: 2,
1215 idle_timeout_ms: 0,
1216 },
1217 );
1218 let a = name("a");
1219 let b = name("b");
1220 let c = name("c");
1221 let a_lease = ws.lease(&a, 1_000, IfMissing::Create).unwrap();
1222 drop(ws.get(&b, 2_000, IfMissing::Create).unwrap());
1223
1224 let c_lease = ws.lease(&c, 3_000, IfMissing::Create).unwrap();
1226 assert_eq!(ws.open_count(), 2);
1227 assert_eq!(opens.load(Ordering::SeqCst), 3);
1228 assert_eq!(a_lease.stats().facts, 0);
1229 drop(c_lease);
1230 drop(a_lease);
1231
1232 drop(ws.get(&a, 4_000, IfMissing::Fail).unwrap());
1234 assert_eq!(opens.load(Ordering::SeqCst), 3);
1235 drop(ws.get(&b, 5_000, IfMissing::Fail).unwrap());
1236 assert_eq!(opens.load(Ordering::SeqCst), 4);
1237 }
1238
1239 #[test]
1240 fn same_name_leases_share_one_slot_until_the_last_drop() {
1241 let tmp = TempDir::new("pool-lease-shared");
1242 let (ws, opens) = workspace(
1243 &tmp,
1244 WorkspaceLimits {
1245 max_open: 1,
1246 idle_timeout_ms: 0,
1247 },
1248 );
1249 let a = name("a");
1250 let first = ws.lease(&a, 1_000, IfMissing::Create).unwrap();
1251 let second = ws.lease(&a, 2_000, IfMissing::Fail).unwrap();
1252 assert_eq!(opens.load(Ordering::SeqCst), 1);
1253
1254 drop(first);
1255 assert!(matches!(ws.release(&a), Err(WorkspaceError::InUse { .. })));
1256 assert_eq!(second.stats().facts, 0);
1257
1258 drop(second);
1259 assert!(ws.release(&a).unwrap());
1260 assert_eq!(ws.open_count(), 0);
1261 }
1262
1263 #[test]
1264 fn auto_maintain_preserves_scoped_workspace_ownership_and_tag_catalogue() {
1265 let tmp = TempDir::new("pool-auto-maintain");
1266 let open: Opener = Box::new(|path| {
1267 Ok(Database::builder(crate::Config::default())
1268 .maintain_every_forgets(1)
1269 .open(path)?
1270 .0)
1271 });
1272 let ws = Workspace::new(
1273 WorkspaceLayout::new(&tmp.0),
1274 open,
1275 WorkspaceLimits {
1276 max_open: 1,
1277 idle_timeout_ms: 0,
1278 },
1279 );
1280 let chat = name("chat");
1281 let id = {
1282 let lease = ws.lease(&chat, 1, IfMissing::Create).unwrap();
1283 lease
1284 .remember(crate::RememberInput {
1285 tags: &["temporary"],
1286 ..crate::RememberInput::text(1, "short lived")
1287 })
1288 .unwrap()
1289 .id
1290 };
1291 {
1292 let lease = ws.lease(&chat, 2, IfMissing::Fail).unwrap();
1293 assert!(lease.forget(2, id).unwrap());
1294 assert_eq!(lease.stats().tombstones, 0, "auto maintain purged it");
1295 assert!(
1296 lease
1297 .list_tags(crate::TagQuery::default())
1298 .unwrap()
1299 .items
1300 .is_empty()
1301 );
1302 }
1303 assert!(ws.release(&chat).unwrap());
1304 assert!(Database::open(ws.layout().path_of(&chat), crate::Config::default()).is_ok());
1305 }
1306
1307 #[test]
1308 fn dropping_a_lease_after_unwind_makes_the_entry_available_again() {
1309 let tmp = TempDir::new("pool-lease-unwind");
1310 let (ws, _) = workspace(
1311 &tmp,
1312 WorkspaceLimits {
1313 max_open: 1,
1314 idle_timeout_ms: 0,
1315 },
1316 );
1317 let a = name("a");
1318 let b = name("b");
1319
1320 let panicked = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1321 let _lease = ws.lease(&a, 1_000, IfMissing::Create).unwrap();
1322 panic!("stand in for a binding conversion panic");
1323 }));
1324 assert!(panicked.is_err());
1325
1326 assert!(ws.lease(&b, 2_000, IfMissing::Create).is_ok());
1329 }
1330
1331 #[test]
1332 fn close_all_may_remove_a_leased_pool_entry_without_breaking_lease_drop() {
1333 let tmp = TempDir::new("pool-lease-close-all");
1334 let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
1335 let a = name("a");
1336 let path = ws.layout().path_of(&a);
1337 let lease = ws.lease(&a, 1_000, IfMissing::Create).unwrap();
1338
1339 assert_eq!(ws.close_all(), 1);
1340 assert_eq!(ws.open_count(), 0);
1341 assert!(matches!(
1342 Database::open(&path, crate::Config::default()),
1343 Err(HostError::Locked { .. })
1344 ));
1345 drop(lease);
1346 assert!(Database::open(&path, crate::Config::default()).is_ok());
1347 }
1348
1349 #[test]
1350 fn a_handle_held_by_a_caller_outlives_its_pool_entry() {
1351 let tmp = TempDir::new("pool-outlive");
1352 let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
1353 let chat = name("chat-42");
1354 let held = ws.get(&chat, 1_000, IfMissing::Create).unwrap();
1355 let path = ws.layout().path_of(&chat);
1356
1357 assert_eq!(ws.close_all(), 1);
1358 assert_eq!(ws.open_count(), 0);
1359
1360 held.remember(crate::RememberInput::text(2_000, "still mine"))
1363 .unwrap();
1364 assert!(matches!(
1365 Database::open(&path, crate::Config::default()),
1366 Err(HostError::Locked { .. })
1367 ));
1368 drop(held);
1369 assert!(Database::open(&path, crate::Config::default()).is_ok());
1370 }
1371
1372 #[test]
1373 fn a_database_held_by_another_process_is_reported_by_name() {
1374 let tmp = TempDir::new("pool-busy");
1375 let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
1376 let chat = name("chat-42");
1377 let path = ws.layout().path_of(&chat);
1378 std::fs::create_dir_all(ws.layout().db_dir()).unwrap();
1379
1380 let outsider = Database::open(&path, crate::Config::default()).unwrap().0;
1383 let e = ws.get(&chat, 1_000, IfMissing::Create).unwrap_err();
1384 assert!(matches!(&e, WorkspaceError::Busy { name } if name == &chat));
1385 assert!(e.to_string().contains("chat-42"), "{e}");
1386
1387 drop(outsider);
1388 assert!(ws.get(&chat, 2_000, IfMissing::Fail).is_ok());
1389 }
1390
1391 #[test]
1392 fn an_open_that_fails_for_another_reason_keeps_its_own_error() {
1393 let tmp = TempDir::new("pool-open-err");
1394 let open: Opener = Box::new(|_| Err(HostError::Embed("no provider".into())));
1395 let ws = Workspace::new(
1396 WorkspaceLayout::new(&tmp.0),
1397 open,
1398 WorkspaceLimits::default(),
1399 );
1400 let e = ws.get(&name("a"), 1_000, IfMissing::Create).unwrap_err();
1401 assert!(
1402 matches!(e, WorkspaceError::Host(HostError::Embed(_))),
1403 "{e}"
1404 );
1405 }
1406
1407 #[test]
1408 fn a_directory_that_cannot_be_created_is_an_error_not_a_panic() {
1409 let tmp = TempDir::new("pool-mkdir");
1410 std::fs::write(tmp.0.join(DB_DIR), b"in the way").unwrap();
1412 let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
1413 assert!(matches!(
1414 ws.get(&name("a"), 1_000, IfMissing::Create),
1415 Err(WorkspaceError::Io { .. })
1416 ));
1417 }
1418
1419 proptest::proptest! {
1420 #[test]
1426 fn a_name_that_parses_can_only_resolve_inside_the_workspace(s in ".*") {
1427 let Ok(name) = DbName::parse(&s) else { return Ok(()) };
1428 let layout = WorkspaceLayout::new("/ws");
1429 let path = layout.path_of(&name);
1430
1431 let rest: Vec<_> = path
1432 .strip_prefix(layout.db_dir())
1433 .expect("resolved outside the workspace")
1434 .components()
1435 .collect();
1436 let expected = format!("{s}.{DB_EXT}");
1437 proptest::prop_assert_eq!(rest.len(), 1);
1438 proptest::prop_assert_eq!(
1439 path.file_name().and_then(|n| n.to_str()),
1440 Some(expected.as_str())
1441 );
1442 }
1443 }
1444}