1use lmdb::{
7 Cursor, Database, DatabaseFlags, Environment, EnvironmentFlags, Transaction, WriteFlags,
8};
9use std::collections::HashMap;
10use std::path::Path;
11use std::sync::atomic::{AtomicU64, Ordering};
12use std::sync::{Arc, RwLock};
13use wm_core::{CoreError, Galaxy, Result};
14
15#[cfg(unix)]
16use std::os::unix::fs::PermissionsExt;
17
18use crate::episodic::EpisodicStore;
19use crate::indexes::IndexDbs;
20use crate::memory::{Memory, MemoryId, decode_embedding, encode_embedding};
21use crate::semantic::SemanticEncoder;
22
23#[derive(Debug, Clone, Default)]
25pub struct MemoryQuery {
26 pub tags: Vec<String>,
28 pub min_importance: Option<f32>,
30 pub max_importance: Option<f32>,
32 pub created_after: Option<chrono::DateTime<chrono::Utc>>,
34 pub created_before: Option<chrono::DateTime<chrono::Utc>>,
36 pub content_substring: Option<String>,
39 pub limit: usize,
41}
42
43impl MemoryQuery {
44 #[must_use]
46 pub fn new() -> Self {
47 Self {
48 limit: 100,
49 ..Default::default()
50 }
51 }
52
53 #[must_use]
55 pub fn with_tags(mut self, tags: Vec<String>) -> Self {
56 self.tags = tags;
57 self
58 }
59
60 #[must_use]
62 pub const fn with_importance_range(mut self, min: f32, max: f32) -> Self {
63 self.min_importance = Some(min);
64 self.max_importance = Some(max);
65 self
66 }
67
68 #[must_use]
70 pub const fn with_time_range(
71 mut self,
72 after: chrono::DateTime<chrono::Utc>,
73 before: chrono::DateTime<chrono::Utc>,
74 ) -> Self {
75 self.created_after = Some(after);
76 self.created_before = Some(before);
77 self
78 }
79
80 #[must_use]
83 pub const fn with_created_after(mut self, after: chrono::DateTime<chrono::Utc>) -> Self {
84 self.created_after = Some(after);
85 self
86 }
87
88 #[must_use]
91 pub const fn with_created_before(mut self, before: chrono::DateTime<chrono::Utc>) -> Self {
92 self.created_before = Some(before);
93 self
94 }
95
96 #[must_use]
98 pub const fn with_limit(mut self, limit: usize) -> Self {
99 self.limit = limit;
100 self
101 }
102
103 #[must_use]
105 pub fn with_content_substring(mut self, substring: impl Into<String>) -> Self {
106 self.content_substring = Some(substring.into().to_lowercase());
107 self
108 }
109
110 #[must_use]
112 pub fn matches(&self, mem: &Memory) -> bool {
113 if !self.tags.is_empty() {
115 for tag in &self.tags {
116 if !mem.metadata.tags.iter().any(|t| t == tag) {
117 return false;
118 }
119 }
120 }
121
122 if let Some(min) = self.min_importance {
124 if mem.metadata.importance < min {
125 return false;
126 }
127 }
128 if let Some(max) = self.max_importance {
129 if mem.metadata.importance > max {
130 return false;
131 }
132 }
133
134 if let Some(after) = self.created_after {
136 if mem.metadata.created_at < after {
137 return false;
138 }
139 }
140 if let Some(before) = self.created_before {
141 if mem.metadata.created_at > before {
142 return false;
143 }
144 }
145
146 if let Some(sub) = &self.content_substring {
148 if !mem.content.to_lowercase().contains(sub) {
149 return false;
150 }
151 }
152
153 true
154 }
155}
156
157pub struct MemoryStore {
159 path: std::path::PathBuf,
161 env: Environment,
163 index_dbs: IndexDbs,
165 semantic_encoder: SemanticEncoder,
167 max_entries_per_galaxy: Option<usize>,
169 mutation_count: AtomicU64,
173 episodic_db: Database,
175 episodic_terms_v2_db: Database,
177 embedding_cache_db: Database,
182 revisions_db: Database,
186 attestations_db: Database,
190 pub(crate) cold_storage_db: Database,
192 episodic_term_cache: std::sync::Arc<RwLock<HashMap<String, Vec<uuid::Uuid>>>>,
194 episodic_embedder:
196 std::sync::OnceLock<Option<Arc<dyn crate::embedder::Embedder + Send + Sync>>>,
197 episodic_sidecar_ensured: std::sync::OnceLock<()>,
199 episodic_aliases: std::sync::OnceLock<Option<crate::episodic_keys::AdaptiveAliases>>,
201 episodic_enrichment: std::sync::OnceLock<Option<crate::enrichment::VocabularyEnrichment>>,
203}
204
205impl MemoryStore {
206 pub fn open(path: impl AsRef<Path>, map_size: usize) -> Result<Self> {
212 let path = path.as_ref().to_path_buf();
213
214 std::fs::create_dir_all(&path)
216 .map_err(|e| CoreError::Memory(format!("Cannot create store dir: {e}")))?;
217 #[cfg(unix)]
218 {
219 std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o700))
220 .map_err(|e| CoreError::Memory(format!("Cannot set store dir permissions: {e}")))?;
221 }
222
223 let env = Environment::new()
224 .set_map_size(map_size)
225 .set_max_dbs(64)
226 .open(&path)
227 .map_err(|e| CoreError::Memory(format!("LMDB open failed: {e}")))?;
228
229 for galaxy in Galaxy::all() {
231 let db = env
232 .create_db(Some(galaxy.db_name()), DatabaseFlags::default())
233 .map_err(|e| {
234 CoreError::Memory(format!(
235 "LMDB create_db failed for {}: {e}",
236 galaxy.db_name()
237 ))
238 })?;
239 let _ = db;
240 }
241
242 for (name, flags) in crate::indexes::INDEX_DBS {
244 let db = env
245 .create_db(Some(name), *flags)
246 .map_err(|e| CoreError::Memory(format!("LMDB create_db failed for {name}: {e}")))?;
247 let _ = db;
248 }
249
250 let index_dbs = IndexDbs::open(&env)?;
251 let episodic_db = env
252 .create_db(Some("episodic_records"), DatabaseFlags::default())
253 .map_err(|e| {
254 CoreError::Memory(format!("LMDB create_db failed for episodic_records: {e}"))
255 })?;
256 let episodic_terms_v2_db = env
263 .create_db(Some("episodic_terms_v2"), DatabaseFlags::DUP_SORT)
264 .map_err(|e| {
265 CoreError::Memory(format!("LMDB create_db failed for episodic_terms_v2: {e}"))
266 })?;
267 let embedding_cache_db = env
268 .create_db(Some("embedding_cache"), DatabaseFlags::default())
269 .map_err(|e| {
270 CoreError::Memory(format!("LMDB create_db failed for embedding_cache: {e}"))
271 })?;
272 let revisions_db = env
273 .create_db(Some("revisions"), DatabaseFlags::default())
274 .map_err(|e| CoreError::Memory(format!("LMDB create_db failed for revisions: {e}")))?;
275 let attestations_db = env
276 .create_db(
277 Some(crate::attestation::ATTESTATIONS_DB),
278 DatabaseFlags::default(),
279 )
280 .map_err(|e| {
281 CoreError::Memory(format!("LMDB create_db failed for attestations: {e}"))
282 })?;
283 let cold_storage_db = env
284 .create_db(Some("cold_storage"), DatabaseFlags::default())
285 .map_err(|e| {
286 CoreError::Memory(format!("LMDB create_db failed for cold_storage: {e}"))
287 })?;
288 Ok(Self {
289 path,
290 env,
291 index_dbs,
292 semantic_encoder: SemanticEncoder::new(),
293 max_entries_per_galaxy: None,
294 mutation_count: AtomicU64::new(0),
295 episodic_db,
296 episodic_terms_v2_db,
297 embedding_cache_db,
298 revisions_db,
299 attestations_db,
300 cold_storage_db,
301 episodic_term_cache: std::sync::Arc::new(RwLock::new(HashMap::new())),
302 episodic_embedder: std::sync::OnceLock::new(),
303 episodic_sidecar_ensured: std::sync::OnceLock::new(),
304 episodic_aliases: std::sync::OnceLock::new(),
305 episodic_enrichment: std::sync::OnceLock::new(),
306 })
307 }
308
309 pub fn open_default(path: impl AsRef<Path>) -> Result<Self> {
318 let platform_default = if cfg!(windows) {
323 256 * 1024 * 1024
324 } else {
325 4 * 1024 * 1024 * 1024
326 };
327 let size = std::env::var("WM_DEFAULT_MAP_SIZE")
328 .ok()
329 .and_then(|v| v.parse::<usize>().ok())
330 .filter(|&v| v > 0)
331 .unwrap_or(platform_default);
332 Self::open(path, size)
333 }
334
335 pub fn open_readonly(path: impl AsRef<Path>) -> Result<Self> {
341 let path = path.as_ref().to_path_buf();
342 if !path.is_dir() {
343 return Err(CoreError::Memory(format!(
344 "Read-only LMDB store directory does not exist: {}",
345 path.display()
346 )));
347 }
348 if !path.join("data.mdb").is_file() {
349 return Err(CoreError::Memory(format!(
350 "Read-only LMDB store is missing data.mdb: {}",
351 path.display()
352 )));
353 }
354
355 let env = Environment::new()
356 .set_max_dbs(32)
357 .set_flags(EnvironmentFlags::READ_ONLY)
358 .open(&path)
359 .map_err(|e| CoreError::Memory(format!("Read-only LMDB open failed: {e}")))?;
360
361 let index_dbs = IndexDbs::open(&env)?;
362 let open_named = |name: &str| {
363 env.open_db(Some(name)).map_err(|e| {
364 CoreError::Memory(format!("Read-only LMDB missing database {name}: {e}"))
365 })
366 };
367 let episodic_db = open_named("episodic_records")?;
368 let episodic_terms_v2_db = open_named("episodic_terms_v2")?;
369 let embedding_cache_db = open_named("embedding_cache")?;
370 let revisions_db = open_named("revisions")?;
371 let attestations_db = open_named(crate::attestation::ATTESTATIONS_DB)?;
372 let cold_storage_db = open_named("cold_storage")?;
373
374 Ok(Self {
375 path,
376 env,
377 index_dbs,
378 semantic_encoder: SemanticEncoder::new(),
379 max_entries_per_galaxy: None,
380 mutation_count: AtomicU64::new(0),
381 episodic_db,
382 episodic_terms_v2_db,
383 embedding_cache_db,
384 revisions_db,
385 attestations_db,
386 cold_storage_db,
387 episodic_term_cache: std::sync::Arc::new(RwLock::new(HashMap::new())),
388 episodic_embedder: std::sync::OnceLock::new(),
389 episodic_sidecar_ensured: std::sync::OnceLock::new(),
390 episodic_aliases: std::sync::OnceLock::new(),
391 episodic_enrichment: std::sync::OnceLock::new(),
392 })
393 }
394
395 #[must_use]
400 pub const fn with_entry_limit(mut self, limit: usize) -> Self {
401 self.max_entries_per_galaxy = Some(limit);
402 self
403 }
404
405 pub fn path(&self) -> &Path {
407 &self.path
408 }
409
410 pub const fn env(&self) -> &Environment {
412 &self.env
413 }
414
415 pub fn mutation_count(&self) -> u64 {
419 self.mutation_count.load(Ordering::Relaxed)
420 }
421
422 pub const fn index_dbs(&self) -> &IndexDbs {
424 &self.index_dbs
425 }
426
427 pub const fn semantic_encoder(&self) -> &SemanticEncoder {
429 &self.semantic_encoder
430 }
431
432 fn ensure_episodic_sidecar(&self) {
437 if self.episodic_sidecar_ensured.get().is_some() {
438 return;
439 }
440 let _ = self.episodic_sidecar_ensured.set(());
441 let view = EpisodicStore::new(
442 &self.env,
443 self.episodic_db,
444 self.episodic_terms_v2_db,
445 self.episodic_term_cache.clone(),
446 &self.mutation_count,
447 );
448 let needs_rebuild = matches!(
449 (view.sidecar_is_empty(), view.record_count()),
450 (Ok(true), Ok(n)) if n > 0
451 );
452 if needs_rebuild {
453 match view.rebuild_sidecar() {
454 Ok(n) => tracing::info!("episodic sidecar rebuilt from {n} records"),
455 Err(e) => {
456 tracing::warn!("episodic sidecar rebuild failed: {e}");
457 }
458 }
459 }
460 }
461
462 #[must_use]
464 pub fn episodic(&self) -> EpisodicStore<'_> {
465 self.ensure_episodic_sidecar();
466 let mut store = EpisodicStore::new(
467 &self.env,
468 self.episodic_db,
469 self.episodic_terms_v2_db,
470 self.episodic_term_cache.clone(),
471 &self.mutation_count,
472 );
473 if let Some(Some(embedder)) = self.episodic_embedder.get() {
474 store = store.with_embedder(embedder.clone());
475 }
476 if let Some(Some(aliases)) = self.episodic_aliases.get() {
477 store = store.with_adaptive_aliases(aliases.clone());
478 }
479 if let Some(Some(enrichment)) = self.episodic_enrichment.get() {
480 store = store.with_enrichment(enrichment.clone());
481 }
482 store
483 }
484
485 pub fn set_episodic_embedder(
487 &self,
488 embedder: Arc<dyn crate::embedder::Embedder + Send + Sync>,
489 ) {
490 let _ = self.episodic_embedder.set(Some(embedder));
491 }
492
493 pub fn set_episodic_aliases(&self, aliases: crate::episodic_keys::AdaptiveAliases) {
495 let _ = self.episodic_aliases.set(Some(aliases));
496 }
497
498 pub fn set_episodic_enrichment(&self, enrichment: crate::enrichment::VocabularyEnrichment) {
500 let _ = self.episodic_enrichment.set(Some(enrichment));
501 }
502
503 pub fn galaxy_db(&self, galaxy: Galaxy) -> Result<Database> {
505 self.env.open_db(Some(galaxy.db_name())).map_err(|e| {
506 CoreError::Memory(format!("LMDB open_db failed for {}: {e}", galaxy.db_name()))
507 })
508 }
509
510 pub fn put(&self, galaxy: Galaxy, memory: &Memory) -> Result<()> {
518 if let Some(limit) = self.max_entries_per_galaxy {
520 let current = self.count(galaxy)?;
521 if current >= limit {
522 return Err(CoreError::Memory(format!(
523 "galaxy {} entry limit reached ({current}/{limit}), write rejected",
524 galaxy.db_name()
525 )));
526 }
527 }
528
529 let db = self.galaxy_db(galaxy)?;
530 let key = memory.metadata.id.as_bytes();
531 let val = rmp_serde::to_vec_named(memory)
532 .map_err(|e| CoreError::Memory(format!("serialize failed: {e}")))?;
533
534 let mut tx = self
535 .env
536 .begin_rw_txn()
537 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
538
539 let existing = tx
544 .get(db, key)
545 .ok()
546 .and_then(|bytes| crate::codec::decode(bytes).ok());
547
548 match tx.put(db, key, &val, lmdb::WriteFlags::default()) {
549 Ok(()) => {}
550 Err(lmdb::Error::MapFull) => {
551 tx.abort();
552 return Err(CoreError::Memory(format!(
553 "LMDB map full: galaxy {}, consider growing map size or pruning old memories",
554 galaxy.db_name()
555 )));
556 }
557 Err(e) => {
558 tx.abort();
559 return Err(CoreError::Memory(format!("LMDB put failed: {e}")));
560 }
561 }
562 if let Some(existing) = existing {
563 self.index_dbs.remove(&mut tx, galaxy, &existing)?;
564 }
565 self.index_dbs.add(&mut tx, galaxy, memory)?;
566 tx.commit()
567 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
568 self.mutation_count.fetch_add(1, Ordering::Relaxed);
569 Ok(())
570 }
571
572 pub fn get(&self, galaxy: Galaxy, id: uuid::Uuid) -> Result<Option<Memory>> {
574 let db = self.galaxy_db(galaxy)?;
575 let key = id.as_bytes();
576
577 let tx = self
578 .env
579 .begin_ro_txn()
580 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
581 let result = tx.get(db, key);
582 match result {
583 Ok(bytes) => {
584 let memory: Memory = crate::codec::decode(bytes)
585 .map_err(|e| CoreError::Memory(format!("deserialize failed: {e}")))?;
586 tx.commit()
587 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
588 Ok(Some(memory))
589 }
590 Err(lmdb::Error::NotFound) => {
591 tx.commit()
593 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
594 Ok(None)
595 }
596 Err(e) => Err(CoreError::Memory(format!("LMDB get failed: {e}"))),
597 }
598 }
599
600 pub fn find_across_galaxies(&self, id: uuid::Uuid) -> Result<Option<(Galaxy, Memory)>> {
606 for galaxy in Galaxy::memory_galaxies() {
607 if let Some(mem) = self.get(galaxy, id)? {
608 return Ok(Some((galaxy, mem)));
609 }
610 }
611 Ok(None)
612 }
613
614 pub fn delete(&self, galaxy: Galaxy, id: uuid::Uuid) -> Result<bool> {
617 let db = self.galaxy_db(galaxy)?;
618 let key = id.as_bytes();
619
620 let mut tx = self
621 .env
622 .begin_rw_txn()
623 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
624
625 let exists = tx.get(db, key).is_ok();
627 if exists {
628 if let Ok(bytes) = tx.get(db, key) {
630 if let Ok(memory) = crate::codec::decode(bytes) {
631 let _ = self.index_dbs.remove(&mut tx, galaxy, &memory);
632 }
633 }
634 tx.del(db, key, None)
635 .map_err(|e| CoreError::Memory(format!("LMDB del failed: {e}")))?;
636 }
637 tx.commit()
638 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
639 if exists {
640 self.mutation_count.fetch_add(1, Ordering::Relaxed);
641 }
642 Ok(exists)
643 }
644
645 pub fn scan(&self, galaxy: Galaxy, limit: usize) -> Result<Vec<Memory>> {
647 let db = self.galaxy_db(galaxy)?;
648 let tx = self
649 .env
650 .begin_ro_txn()
651 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
652
653 let mut cursor = tx
654 .open_ro_cursor(db)
655 .map_err(|e| CoreError::Memory(format!("LMDB cursor failed: {e}")))?;
656
657 let mut memories = Vec::with_capacity(limit.min(256));
658 for (i, (_key, val)) in cursor.iter().enumerate() {
659 if memories.len() >= limit {
660 break;
661 }
662 match crate::codec::decode(val) {
663 Ok(memory) => memories.push(memory),
664 Err(e) => {
665 tracing::warn!(
666 "Skipping corrupted entry at index {i} in galaxy {:?}: {e}",
667 galaxy
668 );
669 }
670 }
671 }
672
673 drop(cursor);
674 tx.commit()
675 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
676 Ok(memories)
677 }
678
679 pub fn scan_all(&self, galaxy: Galaxy) -> Result<Vec<Memory>> {
684 self.scan_all_impl(galaxy, false)
685 }
686
687 pub fn scan_all_strict(&self, galaxy: Galaxy) -> Result<Vec<Memory>> {
691 self.scan_all_impl(galaxy, true)
692 }
693
694 fn scan_all_impl(&self, galaxy: Galaxy, strict: bool) -> Result<Vec<Memory>> {
695 let db = self.galaxy_db(galaxy)?;
696 let tx = self
697 .env
698 .begin_ro_txn()
699 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
700
701 let mut cursor = tx
702 .open_ro_cursor(db)
703 .map_err(|e| CoreError::Memory(format!("LMDB cursor failed: {e}")))?;
704
705 let mut memories = Vec::new();
706 for (i, (_key, val)) in cursor.iter().enumerate() {
707 match crate::codec::decode(val) {
708 Ok(memory) => memories.push(memory),
709 Err(e) => {
710 if strict {
711 return Err(CoreError::Memory(format!(
712 "refusing incomplete scan of {}: record {i} cannot be decoded: {e}",
713 galaxy.db_name()
714 )));
715 }
716 tracing::warn!(
717 "Skipping corrupted entry at index {i} in galaxy {:?}: {e}",
718 galaxy
719 );
720 }
721 }
722 }
723
724 drop(cursor);
725 tx.commit()
726 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
727 Ok(memories)
728 }
729
730 pub fn count(&self, galaxy: Galaxy) -> Result<usize> {
732 let db = self.galaxy_db(galaxy)?;
733 let tx = self
734 .env
735 .begin_ro_txn()
736 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
737 let mut cursor = tx
738 .open_ro_cursor(db)
739 .map_err(|e| CoreError::Memory(format!("LMDB cursor failed: {e}")))?;
740 let count = cursor.iter().count();
741 drop(cursor);
742 tx.commit()
743 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
744 Ok(count)
745 }
746
747 pub fn clear_galaxy(&self, galaxy: Galaxy) -> Result<usize> {
751 let db = self.galaxy_db(galaxy)?;
752
753 let mut tx = self
754 .env
755 .begin_rw_txn()
756 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
757
758 let mut cursor = tx
759 .open_ro_cursor(db)
760 .map_err(|e| CoreError::Memory(format!("LMDB cursor failed: {e}")))?;
761
762 let mut count = 0usize;
763 let keys_to_delete: Vec<(Vec<u8>, Memory)> = cursor
764 .iter()
765 .filter_map(|(key, val)| {
766 if let Ok(memory) = crate::codec::decode(val) {
767 Some((key.to_vec(), memory))
768 } else {
769 None
770 }
771 })
772 .collect();
773
774 drop(cursor);
775
776 for (key, memory) in &keys_to_delete {
777 let _ = self.index_dbs.remove(&mut tx, galaxy, memory);
778 tx.del(db, &key, None)
779 .map_err(|e| CoreError::Memory(format!("LMDB del failed: {e}")))?;
780 count += 1;
781 }
782
783 tx.commit()
784 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
785 self.mutation_count
786 .fetch_add(count as u64, Ordering::Relaxed);
787 Ok(count)
788 }
789
790 pub fn batch_put(&self, galaxy: Galaxy, memories: &[Memory]) -> Result<usize> {
793 if memories.is_empty() {
794 return Ok(0);
795 }
796
797 let db = self.galaxy_db(galaxy)?;
798
799 let mut tx = self
800 .env
801 .begin_rw_txn()
802 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
803
804 let mut count = 0usize;
805 for memory in memories {
806 let key = memory.metadata.id.as_bytes();
807 let val = rmp_serde::to_vec_named(memory)
808 .map_err(|e| CoreError::Memory(format!("serialize failed: {e}")))?;
809 match tx.put(db, key, &val, WriteFlags::default()) {
810 Ok(()) => {}
811 Err(lmdb::Error::MapFull) => {
812 tx.abort();
813 return Err(CoreError::Memory(format!(
814 "LMDB map full: galaxy {}, consider growing map size",
815 galaxy.db_name()
816 )));
817 }
818 Err(e) => {
819 tx.abort();
820 return Err(CoreError::Memory(format!("LMDB put failed: {e}")));
821 }
822 }
823 self.index_dbs.add(&mut tx, galaxy, memory)?;
824 count += 1;
825 }
826
827 tx.commit()
828 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
829 self.mutation_count
830 .fetch_add(count as u64, Ordering::Relaxed);
831 Ok(count)
832 }
833
834 pub fn get_raw(&self, galaxy: Galaxy, key: &[u8]) -> Result<Option<Vec<u8>>> {
836 let db = self.galaxy_db(galaxy)?;
837 let tx = self
838 .env
839 .begin_ro_txn()
840 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
841 match tx.get(db, &key) {
842 Ok(bytes) => {
843 let data = bytes.to_vec();
844 tx.commit()
845 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
846 Ok(Some(data))
847 }
848 Err(lmdb::Error::NotFound) => {
849 tx.commit()
850 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
851 Ok(None)
852 }
853 Err(e) => Err(CoreError::Memory(format!("LMDB get_raw failed: {e}"))),
854 }
855 }
856
857 pub fn put_raw(&self, galaxy: Galaxy, key: &[u8], val: &[u8]) -> Result<()> {
859 let db = self.galaxy_db(galaxy)?;
860 let mut tx = self
861 .env
862 .begin_rw_txn()
863 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
864 tx.put(db, &key, &val, lmdb::WriteFlags::default())
865 .map_err(|e| CoreError::Memory(format!("LMDB put_raw failed: {e}")))?;
866 tx.commit()
867 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
868 self.mutation_count.fetch_add(1, Ordering::Relaxed);
869 Ok(())
870 }
871
872 pub fn delete_raw(&self, galaxy: Galaxy, key: &[u8]) -> Result<bool> {
875 let db = self.galaxy_db(galaxy)?;
876 let mut tx = self
877 .env
878 .begin_rw_txn()
879 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
880 let deleted = tx.del(db, &key, None).is_ok();
881 tx.commit()
882 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
883 if deleted {
884 self.mutation_count.fetch_add(1, Ordering::Relaxed);
885 }
886 Ok(deleted)
887 }
888
889 pub fn put_raw_batch(&self, galaxy: Galaxy, entries: &[(&[u8], &[u8])]) -> Result<()> {
892 self.put_raw_batch_impl(galaxy, entries)?;
893 self.mutation_count
894 .fetch_add(entries.len() as u64, Ordering::Relaxed);
895 Ok(())
896 }
897
898 pub fn put_raw_batch_untracked(
908 &self,
909 galaxy: Galaxy,
910 entries: &[(&[u8], &[u8])],
911 ) -> Result<()> {
912 self.put_raw_batch_impl(galaxy, entries)
913 }
914
915 fn put_raw_batch_impl(&self, galaxy: Galaxy, entries: &[(&[u8], &[u8])]) -> Result<()> {
916 let db = self.galaxy_db(galaxy)?;
917 let mut tx = self
918 .env
919 .begin_rw_txn()
920 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
921 for (key, val) in entries {
922 tx.put(db, key, val, WriteFlags::default())
923 .map_err(|e| CoreError::Memory(format!("LMDB put_raw_batch failed: {e}")))?;
924 }
925 tx.commit()
926 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
927 Ok(())
928 }
929
930 pub fn find_by_content_hash(&self, galaxy: Galaxy, hash: &str) -> Result<Option<uuid::Uuid>> {
936 let tx = self
937 .env
938 .begin_ro_txn()
939 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
940 let result = self.index_dbs.find_by_content_hash(&tx, galaxy, hash)?;
941 tx.commit()
942 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
943 Ok(result)
944 }
945
946 pub fn find_by_content_hash_scan(
948 &self,
949 galaxy: Galaxy,
950 hash: &str,
951 ) -> Result<Option<uuid::Uuid>> {
952 let memories = self.scan(galaxy, 10_000)?;
953 for mem in memories {
954 if mem.metadata.content_hash == hash {
955 return Ok(Some(mem.metadata.id));
956 }
957 }
958 Ok(None)
959 }
960
961 pub fn put_dedup(&self, galaxy: Galaxy, memory: &Memory) -> Result<uuid::Uuid> {
965 if let Some(existing_id) =
966 self.find_by_content_hash(galaxy, &memory.metadata.content_hash)?
967 {
968 return Ok(existing_id);
969 }
970 let id = memory.metadata.id;
971 self.put(galaxy, memory)?;
972 Ok(id)
973 }
974
975 pub fn put_batch(&self, galaxy: Galaxy, memories: &[Memory]) -> Result<()> {
980 let db = self.galaxy_db(galaxy)?;
981 let mut tx = self
982 .env
983 .begin_rw_txn()
984 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
985
986 for memory in memories {
987 let key = memory.metadata.id.as_bytes();
988 let val = rmp_serde::to_vec_named(memory)
989 .map_err(|e| CoreError::Memory(format!("serialize failed: {e}")))?;
990 tx.put(db, key, &val, WriteFlags::default())
991 .map_err(|e| CoreError::Memory(format!("LMDB put_batch failed: {e}")))?;
992 self.index_dbs.add(&mut tx, galaxy, memory)?;
993 }
994
995 tx.commit()
996 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
997 self.mutation_count
998 .fetch_add(memories.len() as u64, Ordering::Relaxed);
999 Ok(())
1000 }
1001
1002 pub fn query(&self, galaxy: Galaxy, query: &MemoryQuery) -> Result<Vec<Memory>> {
1009 if query.content_substring.is_none()
1012 && query.tags.len() == 1
1013 && query.min_importance.is_none()
1014 && query.max_importance.is_none()
1015 && query.created_after.is_none()
1016 && query.created_before.is_none()
1017 {
1018 return self.query_by_tag_indexed(galaxy, &query.tags[0], query.limit);
1019 }
1020
1021 if query.content_substring.is_none()
1022 && query.tags.is_empty()
1023 && let Some(min) = query.min_importance
1024 && let Some(max) = query.max_importance
1025 && query.created_after.is_none()
1026 && query.created_before.is_none()
1027 {
1028 return self.query_by_importance_indexed(galaxy, min, max, query.limit);
1029 }
1030
1031 if query.content_substring.is_none()
1032 && query.tags.is_empty()
1033 && query.min_importance.is_none()
1034 && query.max_importance.is_none()
1035 && let Some(after) = query.created_after
1036 && let Some(before) = query.created_before
1037 {
1038 return self.query_by_time_indexed(galaxy, after, before, query.limit);
1039 }
1040
1041 let memories = self.scan(galaxy, 10_000)?;
1043 let mut results = Vec::new();
1044 for mem in memories {
1045 if query.matches(&mem) {
1046 results.push(mem);
1047 if results.len() >= query.limit {
1048 break;
1049 }
1050 }
1051 }
1052 Ok(results)
1053 }
1054
1055 fn query_by_tag_indexed(&self, galaxy: Galaxy, tag: &str, limit: usize) -> Result<Vec<Memory>> {
1057 let db = self.galaxy_db(galaxy)?;
1058 let tx = self
1059 .env
1060 .begin_ro_txn()
1061 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1062 let ids = self.index_dbs.find_by_tag(&tx, galaxy, tag)?;
1063 let mut results = Vec::new();
1064 for id in &ids {
1065 if results.len() >= limit {
1066 break;
1067 }
1068 if let Ok(bytes) = tx.get(db, id.as_bytes()) {
1069 if let Ok(mem) = crate::codec::decode(bytes) {
1070 results.push(mem);
1071 }
1072 }
1073 }
1074 tx.commit()
1075 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1076 Ok(results)
1077 }
1078
1079 fn query_by_importance_indexed(
1081 &self,
1082 galaxy: Galaxy,
1083 min: f32,
1084 max: f32,
1085 limit: usize,
1086 ) -> Result<Vec<Memory>> {
1087 let db = self.galaxy_db(galaxy)?;
1088 let tx = self
1089 .env
1090 .begin_ro_txn()
1091 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1092 let ids = self
1093 .index_dbs
1094 .find_by_importance_range(&tx, galaxy, min, max)?;
1095 let mut results = Vec::new();
1096 for id in &ids {
1097 if results.len() >= limit {
1098 break;
1099 }
1100 if let Ok(bytes) = tx.get(db, id.as_bytes()) {
1101 if let Ok(mem) = crate::codec::decode(bytes) {
1102 results.push(mem);
1103 }
1104 }
1105 }
1106 tx.commit()
1107 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1108 Ok(results)
1109 }
1110
1111 fn query_by_time_indexed(
1113 &self,
1114 galaxy: Galaxy,
1115 after: chrono::DateTime<chrono::Utc>,
1116 before: chrono::DateTime<chrono::Utc>,
1117 limit: usize,
1118 ) -> Result<Vec<Memory>> {
1119 let db = self.galaxy_db(galaxy)?;
1120 let tx = self
1121 .env
1122 .begin_ro_txn()
1123 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1124 let ids = self
1125 .index_dbs
1126 .find_by_time_range(&tx, galaxy, after, before)?;
1127 let mut results = Vec::new();
1128 for id in &ids {
1129 if results.len() >= limit {
1130 break;
1131 }
1132 if let Ok(bytes) = tx.get(db, id.as_bytes()) {
1133 if let Ok(mem) = crate::codec::decode(bytes) {
1134 results.push(mem);
1135 }
1136 }
1137 }
1138 tx.commit()
1139 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1140 Ok(results)
1141 }
1142
1143 pub fn put_semantic(&self, galaxy: Galaxy, memory: &mut Memory) -> Result<()> {
1151 let temporal_weight = memory.metadata.coord5d.w;
1152 let importance = memory.metadata.importance;
1153 memory.metadata.coord5d =
1154 self.semantic_encoder
1155 .encode_coordinate(&memory.content, temporal_weight, importance);
1156 self.put(galaxy, memory)
1157 }
1158
1159 pub fn find_similar(
1164 &self,
1165 galaxy: Galaxy,
1166 query_text: &str,
1167 limit: usize,
1168 ) -> Result<Vec<(Memory, f32)>> {
1169 let query_coord = self
1170 .semantic_encoder
1171 .encode_coordinate(query_text, 0.5, 0.5);
1172 let memories = self.scan(galaxy, 10_000)?;
1173 let mut results: Vec<(Memory, f32)> = memories
1174 .into_iter()
1175 .map(|m| {
1176 let dist = query_coord.semantic_distance_to(&m.metadata.coord5d);
1177 (m, dist)
1178 })
1179 .collect();
1180 results.sort_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal));
1181 results.truncate(limit);
1182 Ok(results)
1183 }
1184
1185 pub fn put_embedding(&self, memory_id: uuid::Uuid, embedding: &[f32]) -> Result<()> {
1190 let db = self.galaxy_db(Galaxy::Embeddings)?;
1191 let key = memory_id.as_bytes();
1192 let val = encode_embedding(embedding);
1193
1194 let mut tx = self
1195 .env
1196 .begin_rw_txn()
1197 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
1198 tx.put(db, key, &val, WriteFlags::default())
1199 .map_err(|e| CoreError::Memory(format!("LMDB put_embedding failed: {e}")))?;
1200 tx.commit()
1201 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1202 self.mutation_count.fetch_add(1, Ordering::Relaxed);
1203 Ok(())
1204 }
1205
1206 pub fn get_embedding(&self, memory_id: uuid::Uuid) -> Result<Option<Vec<f32>>> {
1208 let db = self.galaxy_db(Galaxy::Embeddings)?;
1209 let key = memory_id.as_bytes();
1210
1211 let tx = self
1212 .env
1213 .begin_ro_txn()
1214 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1215 match tx.get(db, key) {
1216 Ok(bytes) => {
1217 let embedding = decode_embedding(bytes);
1218 tx.commit()
1219 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1220 Ok(Some(embedding))
1221 }
1222 Err(lmdb::Error::NotFound) => {
1223 tx.commit()
1224 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1225 Ok(None)
1226 }
1227 Err(e) => Err(CoreError::Memory(format!("LMDB get_embedding failed: {e}"))),
1228 }
1229 }
1230
1231 pub fn delete_embedding(&self, memory_id: uuid::Uuid) -> Result<bool> {
1233 let db = self.galaxy_db(Galaxy::Embeddings)?;
1234 let key = memory_id.as_bytes();
1235
1236 let mut tx = self
1237 .env
1238 .begin_rw_txn()
1239 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
1240 let exists = tx.get(db, key).is_ok();
1241 if exists {
1242 tx.del(db, key, None)
1243 .map_err(|e| CoreError::Memory(format!("LMDB del_embedding failed: {e}")))?;
1244 }
1245 tx.commit()
1246 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1247 if exists {
1248 self.mutation_count.fetch_add(1, Ordering::Relaxed);
1249 }
1250 Ok(exists)
1251 }
1252
1253 pub fn put_embedding_cache(&self, cache_key: &str, embedding: &[f32]) -> Result<()> {
1259 let mut tx = self
1260 .env
1261 .begin_rw_txn()
1262 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
1263 tx.put(
1264 self.embedding_cache_db,
1265 &cache_key.as_bytes().to_vec(),
1266 &encode_embedding(embedding),
1267 WriteFlags::default(),
1268 )
1269 .map_err(|e| CoreError::Memory(format!("LMDB put_embedding_cache failed: {e}")))?;
1270 tx.commit()
1271 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1272 self.mutation_count.fetch_add(1, Ordering::Relaxed);
1273 Ok(())
1274 }
1275
1276 pub fn put_embedding_cache_batch(&self, entries: &[(String, Vec<f32>)]) -> Result<()> {
1278 if entries.is_empty() {
1279 return Ok(());
1280 }
1281 let mut tx = self
1282 .env
1283 .begin_rw_txn()
1284 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
1285 for (key, embedding) in entries {
1286 tx.put(
1287 self.embedding_cache_db,
1288 &key.as_bytes().to_vec(),
1289 &encode_embedding(embedding),
1290 WriteFlags::default(),
1291 )
1292 .map_err(|e| CoreError::Memory(format!("LMDB put_embedding_cache failed: {e}")))?;
1293 }
1294 tx.commit()
1295 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1296 self.mutation_count
1297 .fetch_add(entries.len() as u64, Ordering::Relaxed);
1298 Ok(())
1299 }
1300
1301 pub fn get_embedding_cache(&self, cache_key: &str) -> Result<Option<Vec<f32>>> {
1303 let tx = self
1304 .env
1305 .begin_ro_txn()
1306 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1307 match tx.get(self.embedding_cache_db, &cache_key.as_bytes().to_vec()) {
1308 Ok(bytes) => {
1309 let embedding = decode_embedding(bytes);
1310 tx.commit()
1311 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1312 Ok(Some(embedding))
1313 }
1314 Err(lmdb::Error::NotFound) => {
1315 tx.commit()
1316 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1317 Ok(None)
1318 }
1319 Err(e) => Err(CoreError::Memory(format!(
1320 "LMDB get_embedding_cache failed: {e}"
1321 ))),
1322 }
1323 }
1324
1325 pub fn get_embedding_cache_batch(&self, keys: &[String]) -> Result<Vec<Option<Vec<f32>>>> {
1328 let tx = self
1329 .env
1330 .begin_ro_txn()
1331 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1332 let mut out = Vec::with_capacity(keys.len());
1333 for key in keys {
1334 out.push(
1335 tx.get(self.embedding_cache_db, &key.as_bytes().to_vec())
1336 .ok()
1337 .map(decode_embedding),
1338 );
1339 }
1340 tx.commit()
1341 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1342 Ok(out)
1343 }
1344
1345 pub fn embedding_cache_count(&self) -> Result<u64> {
1347 let tx = self
1348 .env
1349 .begin_ro_txn()
1350 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1351 let mut cursor = tx
1352 .open_ro_cursor(self.embedding_cache_db)
1353 .map_err(|e| CoreError::Memory(format!("LMDB cursor embedding_cache failed: {e}")))?;
1354 let mut count = 0u64;
1355 for _ in cursor.iter() {
1356 count += 1;
1357 }
1358 Ok(count)
1359 }
1360
1361 pub fn record_revision(
1367 &self,
1368 galaxy: Galaxy,
1369 id: MemoryId,
1370 old_hash: &str,
1371 new_hash: &str,
1372 actor: crate::revision::RevisionActor,
1373 ) -> Result<crate::revision::MemoryRevision> {
1374 let seq = self.revisions(galaxy, id)?.len() as u32;
1375 let entry = crate::revision::MemoryRevision {
1376 seq,
1377 timestamp: wm_core::time::now_unix_secs(),
1378 old_hash: old_hash.to_string(),
1379 new_hash: new_hash.to_string(),
1380 actor_session: actor.session,
1381 actor_user: actor.user,
1382 actor_compartment: actor.compartment,
1383 };
1384 let key = crate::revision::revision_key(galaxy, id, seq);
1385 let val = serde_json::to_vec(&entry)
1386 .map_err(|e| CoreError::Memory(format!("revision serialize failed: {e}")))?;
1387 let mut tx = self
1388 .env
1389 .begin_rw_txn()
1390 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
1391 tx.put(self.revisions_db, &key, &val, WriteFlags::default())
1392 .map_err(|e| CoreError::Memory(format!("LMDB put revision failed: {e}")))?;
1393 tx.commit()
1394 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1395 self.mutation_count.fetch_add(1, Ordering::Relaxed);
1396 Ok(entry)
1397 }
1398
1399 pub fn revisions(
1402 &self,
1403 galaxy: Galaxy,
1404 id: MemoryId,
1405 ) -> Result<Vec<crate::revision::MemoryRevision>> {
1406 const MDB_GET_CURRENT: u32 = 4;
1411 const MDB_NEXT: u32 = 8;
1412 const MDB_SET_RANGE: u32 = 17;
1413 let prefix = crate::revision::revision_prefix(galaxy, id);
1414 let tx = self
1415 .env
1416 .begin_ro_txn()
1417 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1418 let cursor = tx
1419 .open_ro_cursor(self.revisions_db)
1420 .map_err(|e| CoreError::Memory(format!("LMDB cursor revisions failed: {e}")))?;
1421 let mut out = Vec::new();
1422 if cursor.get(Some(&prefix), None, MDB_SET_RANGE).is_ok() {
1423 while let Ok((key, val)) = cursor.get(None, None, MDB_GET_CURRENT) {
1424 if !key.is_some_and(|k| k.starts_with(&prefix)) {
1427 break;
1428 }
1429 let entry: crate::revision::MemoryRevision = serde_json::from_slice(val)
1430 .map_err(|e| CoreError::Memory(format!("revision deserialize failed: {e}")))?;
1431 out.push(entry);
1432 if out.len() >= 10_000 || cursor.get(None, None, MDB_NEXT).is_err() {
1434 break;
1435 }
1436 }
1437 }
1438 drop(cursor);
1439 tx.commit()
1440 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1441 Ok(out)
1442 }
1443
1444 pub fn verify_revision_chain(
1447 &self,
1448 galaxy: Galaxy,
1449 id: MemoryId,
1450 current_hash: &str,
1451 ) -> Result<crate::revision::RevisionChainReport> {
1452 let entries = self.revisions(galaxy, id)?;
1453 Ok(crate::revision::verify_chain(&entries, current_hash))
1454 }
1455
1456 pub fn record_attestation(
1464 &self,
1465 galaxy: Galaxy,
1466 id: MemoryId,
1467 entry: &crate::attestation::RecordAttestation,
1468 ) -> Result<()> {
1469 let key = crate::attestation::attestation_key(galaxy, id);
1470 let val = serde_json::to_vec(entry)
1471 .map_err(|e| CoreError::Memory(format!("attestation serialize failed: {e}")))?;
1472 let mut tx = self
1473 .env
1474 .begin_rw_txn()
1475 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
1476 tx.put(self.attestations_db, &key, &val, WriteFlags::default())
1477 .map_err(|e| CoreError::Memory(format!("LMDB put attestation failed: {e}")))?;
1478 tx.commit()
1479 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1480 self.mutation_count.fetch_add(1, Ordering::Relaxed);
1481 Ok(())
1482 }
1483
1484 pub fn attestation(
1487 &self,
1488 galaxy: Galaxy,
1489 id: MemoryId,
1490 ) -> Result<Option<crate::attestation::RecordAttestation>> {
1491 let key = crate::attestation::attestation_key(galaxy, id);
1492 let tx = self
1493 .env
1494 .begin_ro_txn()
1495 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1496 let out =
1497 match tx.get(self.attestations_db, &key) {
1498 Ok(val) => Some(serde_json::from_slice(val).map_err(|e| {
1499 CoreError::Memory(format!("attestation deserialize failed: {e}"))
1500 })?),
1501 Err(lmdb::Error::NotFound) => None,
1502 Err(e) => {
1503 return Err(CoreError::Memory(format!(
1504 "LMDB get attestation failed: {e}"
1505 )));
1506 }
1507 };
1508 drop(tx);
1509 Ok(out)
1510 }
1511
1512 pub fn scan_attestations(&self) -> Result<Vec<crate::attestation::RecordAttestation>> {
1515 const MDB_GET_CURRENT: u32 = 4;
1516 const MDB_NEXT: u32 = 8;
1517 const MDB_FIRST: u32 = 9;
1518 let tx = self
1519 .env
1520 .begin_ro_txn()
1521 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1522 let cursor = tx
1523 .open_ro_cursor(self.attestations_db)
1524 .map_err(|e| CoreError::Memory(format!("LMDB cursor attestations failed: {e}")))?;
1525 let mut out = Vec::new();
1526 if cursor.get(None, None, MDB_FIRST).is_ok() {
1527 while let Ok((_, val)) = cursor.get(None, None, MDB_GET_CURRENT) {
1528 let entry: crate::attestation::RecordAttestation = serde_json::from_slice(val)
1529 .map_err(|e| {
1530 CoreError::Memory(format!("attestation deserialize failed: {e}"))
1531 })?;
1532 out.push(entry);
1533 if out.len() >= 1_000_000 || cursor.get(None, None, MDB_NEXT).is_err() {
1535 break;
1536 }
1537 }
1538 }
1539 drop(cursor);
1540 tx.commit()
1541 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1542 Ok(out)
1543 }
1544
1545 pub fn verify_attestation(
1551 &self,
1552 galaxy: Galaxy,
1553 id: MemoryId,
1554 ) -> Result<crate::attestation::AttestationReport> {
1555 use crate::attestation::AttestationReport;
1556 let Some(att) = self.attestation(galaxy, id)? else {
1557 return Ok(AttestationReport {
1558 attested: false,
1559 signature_valid: false,
1560 matches_head: false,
1561 memory_present: self.get(galaxy, id)?.is_some(),
1562 breaks: vec!["no attestation recorded for this memory".to_string()],
1563 });
1564 };
1565 let mut breaks = Vec::new();
1566 let signature_valid = crate::attestation::verify_attestation(&att);
1567 if !signature_valid {
1568 breaks.push("signature does not verify against recorded pubkey".to_string());
1569 }
1570 let (matches_head, memory_present) = if let Some(memory) = self.get(galaxy, id)? {
1571 let matches = memory.metadata.content_hash == att.record_hash;
1572 if !matches {
1573 breaks.push(
1574 "attested record_hash != live content_hash (memory updated after attestation)"
1575 .to_string(),
1576 );
1577 }
1578 (matches, true)
1579 } else {
1580 breaks.push("attested memory id not present in galaxy".to_string());
1581 (false, false)
1582 };
1583 Ok(AttestationReport {
1584 attested: true,
1585 signature_valid,
1586 matches_head,
1587 memory_present,
1588 breaks,
1589 })
1590 }
1591
1592 pub fn attestation_sweep(
1596 &self,
1597 ) -> Result<
1598 Vec<(
1599 crate::attestation::RecordAttestation,
1600 crate::attestation::AttestationReport,
1601 )>,
1602 > {
1603 use crate::attestation::AttestationReport;
1604 let mut out = Vec::new();
1605 for att in self.scan_attestations()? {
1606 let parsed = match (
1607 Galaxy::from_db_name(&att.galaxy),
1608 uuid::Uuid::parse_str(&att.memory_id),
1609 ) {
1610 (Some(galaxy), Ok(id)) => Some((galaxy, id)),
1611 _ => None,
1612 };
1613 match parsed {
1614 Some((galaxy, id)) => out.push((att, self.verify_attestation(galaxy, id)?)),
1615 None => out.push((
1616 att,
1617 AttestationReport {
1618 attested: true,
1619 signature_valid: false,
1620 matches_head: false,
1621 memory_present: false,
1622 breaks: vec![
1623 "attestation row has unparseable galaxy or memory id".to_string(),
1624 ],
1625 },
1626 )),
1627 }
1628 }
1629 Ok(out)
1630 }
1631
1632 pub fn put_cold_record(&self, record: &crate::cold_storage::ColdRecord) -> Result<()> {
1636 let key = record.id.as_bytes();
1637 let val = rmp_serde::to_vec_named(record)
1638 .map_err(|e| CoreError::Memory(format!("Cold record serialization failed: {e}")))?;
1639 let mut tx = self
1640 .env
1641 .begin_rw_txn()
1642 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
1643 tx.put(self.cold_storage_db, key, &val, WriteFlags::default())
1644 .map_err(|e| CoreError::Memory(format!("LMDB put cold_storage failed: {e}")))?;
1645 tx.commit()
1646 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1647 self.mutation_count.fetch_add(1, Ordering::Relaxed);
1648 Ok(())
1649 }
1650
1651 pub fn get_cold_record(&self, id: MemoryId) -> Result<Option<crate::cold_storage::ColdRecord>> {
1653 let key = id.as_bytes();
1654 let tx = self
1655 .env
1656 .begin_ro_txn()
1657 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1658 match tx.get(self.cold_storage_db, key) {
1659 Ok(bytes) => {
1660 let record: crate::cold_storage::ColdRecord = rmp_serde::from_slice(bytes)
1661 .map_err(|e| {
1662 CoreError::Memory(format!("Cold record deserialization failed: {e}"))
1663 })?;
1664 Ok(Some(record))
1665 }
1666 Err(lmdb::Error::NotFound) => Ok(None),
1667 Err(e) => Err(CoreError::Memory(format!(
1668 "LMDB get cold_storage failed: {e}"
1669 ))),
1670 }
1671 }
1672
1673 pub fn delete_cold_record(&self, id: MemoryId) -> Result<bool> {
1675 let key = id.as_bytes();
1676 let mut tx = self
1677 .env
1678 .begin_rw_txn()
1679 .map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
1680 let deleted = match tx.del(self.cold_storage_db, key, None) {
1681 Ok(()) => true,
1682 Err(lmdb::Error::NotFound) => false,
1683 Err(e) => {
1684 return Err(CoreError::Memory(format!(
1685 "LMDB del cold_storage failed: {e}"
1686 )));
1687 }
1688 };
1689 tx.commit()
1690 .map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
1691 if deleted {
1692 self.mutation_count.fetch_add(1, Ordering::Relaxed);
1693 }
1694 Ok(deleted)
1695 }
1696
1697 pub fn count_cold(&self, galaxy: Option<Galaxy>) -> Result<usize> {
1699 let tx = self
1700 .env
1701 .begin_ro_txn()
1702 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1703 let mut cursor = tx
1704 .open_ro_cursor(self.cold_storage_db)
1705 .map_err(|e| CoreError::Memory(format!("LMDB open_ro_cursor failed: {e}")))?;
1706 let mut count = 0;
1707 for (_key, val) in cursor.iter() {
1708 if let Some(target_g) = galaxy {
1709 let record: crate::cold_storage::ColdRecord =
1710 rmp_serde::from_slice(val).map_err(|e| {
1711 CoreError::Memory(format!("Cold record deserialization failed: {e}"))
1712 })?;
1713 if record.galaxy == target_g {
1714 count += 1;
1715 }
1716 } else {
1717 count += 1;
1718 }
1719 }
1720 Ok(count)
1721 }
1722
1723 pub fn list_cold_records(
1725 &self,
1726 galaxy: Option<Galaxy>,
1727 limit: usize,
1728 ) -> Result<Vec<crate::cold_storage::ColdRecordSummary>> {
1729 let query = crate::cold_storage::ColdQuery {
1730 galaxy,
1731 limit: if limit == 0 { 100 } else { limit },
1732 ..Default::default()
1733 };
1734 self.query_cold_records(&query)
1735 }
1736
1737 pub fn query_cold_records(
1739 &self,
1740 query: &crate::cold_storage::ColdQuery,
1741 ) -> Result<Vec<crate::cold_storage::ColdRecordSummary>> {
1742 let tx = self
1743 .env
1744 .begin_ro_txn()
1745 .map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
1746 let mut cursor = tx
1747 .open_ro_cursor(self.cold_storage_db)
1748 .map_err(|e| CoreError::Memory(format!("LMDB open_ro_cursor failed: {e}")))?;
1749 let mut results = Vec::new();
1750 let limit = if query.limit == 0 {
1751 usize::MAX
1752 } else {
1753 query.limit
1754 };
1755
1756 for (_key, val) in cursor.iter() {
1757 let record: crate::cold_storage::ColdRecord =
1758 rmp_serde::from_slice(val).map_err(|e| {
1759 CoreError::Memory(format!("Cold record deserialization failed: {e}"))
1760 })?;
1761 let summary = record.summary();
1762 if query.matches(&summary) {
1763 results.push(summary);
1764 if results.len() >= limit {
1765 break;
1766 }
1767 }
1768 }
1769 Ok(results)
1770 }
1771
1772 pub fn freeze_to_cold(
1778 &self,
1779 search: Option<&crate::SearchEngine>,
1780 memory_id: MemoryId,
1781 distance: f32,
1782 factors: crate::cold_storage::OuterRimFactors,
1783 digest_id: Option<MemoryId>,
1784 notes: Option<String>,
1785 codec: crate::cold_storage::CompressionCodec,
1786 ) -> Result<crate::cold_storage::ColdRecord> {
1787 let (galaxy, mut mem) = self.find_across_galaxies(memory_id)?.ok_or_else(|| {
1788 CoreError::NotFound(format!("Memory {memory_id} not found in hot store"))
1789 })?;
1790
1791 if mem.metadata.tier != crate::memory::Tier::Archival {
1793 let _ = mem.transition_tier(crate::memory::Tier::Archival);
1794 }
1795
1796 let record =
1797 crate::cold_storage::ColdRecord::new(&mem, distance, factors, digest_id, notes, codec)?;
1798
1799 self.put_cold_record(&record)?;
1801
1802 self.delete(galaxy, memory_id)?;
1804
1805 if let Some(engine) = search {
1807 if let Ok(mut writer_guard) = engine.writer() {
1808 let _ = engine.delete_document(&mut writer_guard, &memory_id.to_string());
1809 let _ = engine.commit(&mut writer_guard);
1810 }
1811 }
1812
1813 Ok(record)
1814 }
1815
1816 pub fn thaw_from_cold(
1822 &self,
1823 search: Option<&crate::SearchEngine>,
1824 memory_id: MemoryId,
1825 ) -> Result<Memory> {
1826 let record = self.get_cold_record(memory_id)?.ok_or_else(|| {
1827 CoreError::NotFound(format!("Memory {memory_id} not found in cold storage"))
1828 })?;
1829
1830 let mut mem = record.decompress()?;
1831
1832 let _ = mem.transition_tier(crate::memory::Tier::Episodic);
1834 mem.metadata.accessed_at = chrono::Utc::now();
1835 mem.metadata.access_count += 1;
1836 mem.metadata.recall_count += 1;
1837 if !mem.metadata.tags.iter().any(|t| t == "thawed:phagic") {
1838 mem.metadata.tags.push("thawed:phagic".to_string());
1839 }
1840
1841 self.put(record.galaxy, &mem)?;
1843
1844 if let Some(engine) = search {
1846 if let Ok(mut writer_guard) = engine.writer() {
1847 let _ = engine.index_memory(&mut writer_guard, &mem);
1848 let _ = engine.commit(&mut writer_guard);
1849 }
1850 }
1851
1852 self.delete_cold_record(memory_id)?;
1854
1855 Ok(mem)
1856 }
1857
1858 pub fn find_anywhere(&self, id: MemoryId) -> Result<Option<(Galaxy, Memory, bool)>> {
1862 if let Some((galaxy, mem)) = self.find_across_galaxies(id)? {
1863 return Ok(Some((galaxy, mem, false)));
1864 }
1865 if let Some(cold_record) = self.get_cold_record(id)? {
1866 let galaxy = cold_record.galaxy;
1867 let mem = cold_record.decompress()?;
1868 return Ok(Some((galaxy, mem, true)));
1869 }
1870 Ok(None)
1871 }
1872}
1873
1874#[cfg(test)]
1875mod tests {
1876 use super::*;
1877 use crate::content_hash;
1878
1879 #[test]
1880 fn open_and_create_galaxies() {
1881 let tmp = tempfile::tempdir().unwrap();
1882 let store = MemoryStore::open_default(tmp.path()).unwrap();
1883 for galaxy in Galaxy::all() {
1884 let _db = store.galaxy_db(galaxy).unwrap();
1885 }
1886 }
1887
1888 #[test]
1892 fn query_substring_filters_galaxy_wide() {
1893 let tmp = tempfile::tempdir().unwrap();
1894 let store = MemoryStore::open_default(tmp.path()).unwrap();
1895
1896 for (i, content) in [
1897 "the mesh joins at dawn",
1898 "unrelated content entirely",
1899 "MESH joins at dusk",
1900 ]
1901 .iter()
1902 .enumerate()
1903 {
1904 let mut m = Memory::new(Galaxy::Codex, content.to_string());
1905 m.metadata.importance = 0.5 + i as f32 / 10.0;
1906 store.put(Galaxy::Codex, &m).unwrap();
1907 }
1908
1909 let hits = store
1910 .query(
1911 Galaxy::Codex,
1912 &MemoryQuery::new().with_content_substring("mesh joins"),
1913 )
1914 .unwrap();
1915 assert_eq!(hits.len(), 2, "CI substring must match both: {hits:?}");
1916 assert!(
1917 hits.iter()
1918 .all(|m| m.content.to_lowercase().contains("mesh joins"))
1919 );
1920
1921 let none = store
1922 .query(
1923 Galaxy::Codex,
1924 &MemoryQuery::new().with_content_substring("quantum calendar"),
1925 )
1926 .unwrap();
1927 assert!(none.is_empty(), "no match must be an honest empty set");
1928
1929 let mut tagged = Memory::new(Galaxy::Codex, "mesh joins again".to_string());
1931 tagged.metadata.tags = vec!["mesh".into()];
1932 store.put(Galaxy::Codex, &tagged).unwrap();
1933 let combined = store
1934 .query(
1935 Galaxy::Codex,
1936 &MemoryQuery::new()
1937 .with_tags(vec!["mesh".into()])
1938 .with_content_substring("again"),
1939 )
1940 .unwrap();
1941 assert_eq!(combined.len(), 1);
1942 assert_eq!(combined[0].content, "mesh joins again");
1943 }
1944
1945 #[cfg(unix)]
1946 #[test]
1947 fn store_dir_has_restrictive_permissions() {
1948 let tmp = tempfile::tempdir().unwrap();
1949 let store_path = tmp.path().join("lmdb");
1950 let _store = MemoryStore::open_default(&store_path).unwrap();
1951 let perms = std::fs::metadata(&store_path).unwrap().permissions().mode();
1952 assert_eq!(
1953 perms & 0o777,
1954 0o700,
1955 "store directory should have 0700 permissions, got {:o}",
1956 perms & 0o777
1957 );
1958 }
1959
1960 #[test]
1961 fn put_get_delete_memory() {
1962 let tmp = tempfile::tempdir().unwrap();
1963 let store = MemoryStore::open_default(tmp.path()).unwrap();
1964
1965 let mem = Memory::new(Galaxy::Codex, "Hello world".to_string());
1966 let id = mem.metadata.id;
1967
1968 store.put(Galaxy::Codex, &mem).unwrap();
1969 let retrieved = store.get(Galaxy::Codex, id).unwrap();
1970 assert!(retrieved.is_some());
1971 assert_eq!(retrieved.unwrap().content, "Hello world");
1972
1973 let deleted = store.delete(Galaxy::Codex, id).unwrap();
1974 assert!(deleted);
1975
1976 let gone = store.get(Galaxy::Codex, id).unwrap();
1977 assert!(gone.is_none());
1978 }
1979
1980 #[test]
1981 fn scan_memories() {
1982 let tmp = tempfile::tempdir().unwrap();
1983 let store = MemoryStore::open_default(tmp.path()).unwrap();
1984
1985 for i in 0..5 {
1986 let mem = Memory::new(Galaxy::Codex, format!("memory-{i}"));
1987 store.put(Galaxy::Codex, &mem).unwrap();
1988 }
1989
1990 let all = store.scan(Galaxy::Codex, 100).unwrap();
1991 assert_eq!(all.len(), 5);
1992
1993 let limited = store.scan(Galaxy::Codex, 3).unwrap();
1994 assert_eq!(limited.len(), 3);
1995 }
1996
1997 #[test]
1998 fn overwrite_removes_stale_index_entries() {
1999 let tmp = tempfile::tempdir().unwrap();
2000 let store = MemoryStore::open_default(tmp.path()).unwrap();
2001
2002 let mut mem = Memory::new(Galaxy::Codex, "overwrite target".to_string());
2004 mem.metadata.tags = vec!["alpha".to_string()];
2005 mem.metadata.importance = 0.9;
2006 let id = mem.metadata.id;
2007 store.put(Galaxy::Codex, &mem).unwrap();
2008
2009 let mut updated = Memory::new(Galaxy::Codex, "overwritten content".to_string());
2011 updated.metadata.id = id;
2012 updated.metadata.tags = vec!["beta".to_string()];
2013 updated.metadata.importance = 0.1;
2014 store.put(Galaxy::Codex, &updated).unwrap();
2015
2016 let tx = store.env().begin_ro_txn().unwrap();
2018 let by_alpha = store
2019 .index_dbs()
2020 .find_by_tag(&tx, Galaxy::Codex, "alpha")
2021 .unwrap();
2022 let by_beta = store
2023 .index_dbs()
2024 .find_by_tag(&tx, Galaxy::Codex, "beta")
2025 .unwrap();
2026 assert!(
2027 by_alpha.is_empty(),
2028 "stale tag index entries must be removed on overwrite"
2029 );
2030 assert_eq!(by_beta, vec![id]);
2031
2032 let by_importance = store
2033 .index_dbs()
2034 .find_by_importance_range(&tx, Galaxy::Codex, 0.0, 0.2)
2035 .unwrap();
2036 assert!(
2037 by_importance.contains(&id),
2038 "new importance must be indexed"
2039 );
2040 let by_high = store
2041 .index_dbs()
2042 .find_by_importance_range(&tx, Galaxy::Codex, 0.8, 1.0)
2043 .unwrap();
2044 assert!(
2045 !by_high.contains(&id),
2046 "stale importance index entries must be removed on overwrite"
2047 );
2048
2049 let old_hash = content_hash("overwrite target");
2050 let new_hash = content_hash("overwritten content");
2051 assert_eq!(
2052 store
2053 .index_dbs()
2054 .find_by_content_hash(&tx, Galaxy::Codex, &old_hash)
2055 .unwrap(),
2056 None,
2057 "stale content-hash index entry must be removed"
2058 );
2059 assert_eq!(
2060 store
2061 .index_dbs()
2062 .find_by_content_hash(&tx, Galaxy::Codex, &new_hash)
2063 .unwrap(),
2064 Some(id)
2065 );
2066 }
2067
2068 #[test]
2069 fn count_memories() {
2070 let tmp = tempfile::tempdir().unwrap();
2071 let store = MemoryStore::open_default(tmp.path()).unwrap();
2072
2073 assert_eq!(store.count(Galaxy::Codex).unwrap(), 0);
2074
2075 for i in 0..3 {
2076 let mem = Memory::new(Galaxy::Codex, format!("count-{i}"));
2077 store.put(Galaxy::Codex, &mem).unwrap();
2078 }
2079
2080 assert_eq!(store.count(Galaxy::Codex).unwrap(), 3);
2081 }
2082
2083 #[test]
2084 fn get_nonexistent_returns_none() {
2085 let tmp = tempfile::tempdir().unwrap();
2086 let store = MemoryStore::open_default(tmp.path()).unwrap();
2087 let result = store.get(Galaxy::Codex, uuid::Uuid::new_v4()).unwrap();
2088 assert!(result.is_none());
2089 }
2090
2091 #[test]
2092 fn raw_put_get() {
2093 let tmp = tempfile::tempdir().unwrap();
2094 let store = MemoryStore::open_default(tmp.path()).unwrap();
2095
2096 store
2097 .put_raw(Galaxy::Substrate, b"config:key", b"value123")
2098 .unwrap();
2099 let val = store.get_raw(Galaxy::Substrate, b"config:key").unwrap();
2100 assert_eq!(val, Some(b"value123".to_vec()));
2101 }
2102
2103 #[test]
2104 fn put_dedup_prevents_duplicates() {
2105 let tmp = tempfile::tempdir().unwrap();
2106 let store = MemoryStore::open_default(tmp.path()).unwrap();
2107
2108 let mem1 = Memory::new(Galaxy::Codex, "duplicate content".into());
2109 let id1 = store.put_dedup(Galaxy::Codex, &mem1).unwrap();
2110
2111 let mem2 = Memory::new(Galaxy::Codex, "duplicate content".into());
2112 let id2 = store.put_dedup(Galaxy::Codex, &mem2).unwrap();
2113
2114 assert_eq!(id1, id2, "dedup should return same ID for same content");
2115 assert_eq!(store.count(Galaxy::Codex).unwrap(), 1);
2116 }
2117
2118 #[test]
2119 fn put_dedup_allows_different_content() {
2120 let tmp = tempfile::tempdir().unwrap();
2121 let store = MemoryStore::open_default(tmp.path()).unwrap();
2122
2123 let mem1 = Memory::new(Galaxy::Codex, "content A".into());
2124 store.put_dedup(Galaxy::Codex, &mem1).unwrap();
2125
2126 let mem2 = Memory::new(Galaxy::Codex, "content B".into());
2127 store.put_dedup(Galaxy::Codex, &mem2).unwrap();
2128
2129 assert_eq!(store.count(Galaxy::Codex).unwrap(), 2);
2130 }
2131
2132 #[test]
2133 fn put_batch_atomic_write() {
2134 let tmp = tempfile::tempdir().unwrap();
2135 let store = MemoryStore::open_default(tmp.path()).unwrap();
2136
2137 let memories: Vec<Memory> = (0..10)
2138 .map(|i| Memory::new(Galaxy::Codex, format!("batch-{i}")))
2139 .collect();
2140
2141 store.put_batch(Galaxy::Codex, &memories).unwrap();
2142 assert_eq!(store.count(Galaxy::Codex).unwrap(), 10);
2143 }
2144
2145 #[test]
2146 fn query_by_tags() {
2147 let tmp = tempfile::tempdir().unwrap();
2148 let store = MemoryStore::open_default(tmp.path()).unwrap();
2149
2150 let mem1 = Memory::new(Galaxy::Codex, "tagged memory".into())
2151 .with_tags(vec!["rust".into(), "memory".into()]);
2152 let mem2 =
2153 Memory::new(Galaxy::Codex, "other memory".into()).with_tags(vec!["python".into()]);
2154 store.put(Galaxy::Codex, &mem1).unwrap();
2155 store.put(Galaxy::Codex, &mem2).unwrap();
2156
2157 let query = MemoryQuery::new().with_tags(vec!["rust".into()]);
2158 let results = store.query(Galaxy::Codex, &query).unwrap();
2159 assert_eq!(results.len(), 1);
2160 assert_eq!(results[0].content, "tagged memory");
2161 }
2162
2163 #[test]
2164 fn query_by_importance_range() {
2165 let tmp = tempfile::tempdir().unwrap();
2166 let store = MemoryStore::open_default(tmp.path()).unwrap();
2167
2168 store
2169 .put(
2170 Galaxy::Codex,
2171 &Memory::new(Galaxy::Codex, "low".into()).with_importance(0.1),
2172 )
2173 .unwrap();
2174 store
2175 .put(
2176 Galaxy::Codex,
2177 &Memory::new(Galaxy::Codex, "mid".into()).with_importance(0.5),
2178 )
2179 .unwrap();
2180 store
2181 .put(
2182 Galaxy::Codex,
2183 &Memory::new(Galaxy::Codex, "high".into()).with_importance(0.9),
2184 )
2185 .unwrap();
2186
2187 let query = MemoryQuery::new().with_importance_range(0.4, 0.6);
2188 let results = store.query(Galaxy::Codex, &query).unwrap();
2189 assert_eq!(results.len(), 1);
2190 assert_eq!(results[0].content, "mid");
2191 }
2192
2193 #[test]
2194 fn memory_query_one_sided_time_bounds() {
2195 let old = Memory::new(Galaxy::Codex, "old".into());
2198 let mut recent = Memory::new(Galaxy::Codex, "recent".into());
2199 recent.metadata.created_at = old.metadata.created_at + chrono::Duration::days(30);
2200
2201 let cutoff = old.metadata.created_at + chrono::Duration::days(10);
2202 let after = MemoryQuery::new().with_created_after(cutoff);
2203 assert!(!after.matches(&old), "pre-cutoff memory must not match");
2204 assert!(after.matches(&recent), "post-cutoff memory must match");
2205
2206 let before = MemoryQuery::new().with_created_before(cutoff);
2207 assert!(before.matches(&old), "pre-cutoff memory must match");
2208 assert!(
2209 !before.matches(&recent),
2210 "post-cutoff memory must not match"
2211 );
2212
2213 let edge = MemoryQuery::new().with_created_after(cutoff);
2215 let mut at = Memory::new(Galaxy::Codex, "at cutoff".into());
2216 at.metadata.created_at = cutoff;
2217 assert!(edge.matches(&at), "created_at == after bound is inclusive");
2218 }
2219
2220 #[test]
2221 fn embedding_put_get_delete() {
2222 let tmp = tempfile::tempdir().unwrap();
2223 let store = MemoryStore::open_default(tmp.path()).unwrap();
2224
2225 let id = uuid::Uuid::new_v4();
2226 let embedding = vec![0.1, 0.2, 0.3, 0.4, 0.5];
2227
2228 store.put_embedding(id, &embedding).unwrap();
2229 let retrieved = store.get_embedding(id).unwrap();
2230 assert!(retrieved.is_some());
2231 let retrieved = retrieved.unwrap();
2232 assert_eq!(retrieved.len(), 5);
2233 assert!((retrieved[0] - 0.1).abs() < f32::EPSILON);
2234
2235 assert!(store.delete_embedding(id).unwrap());
2236 assert!(store.get_embedding(id).unwrap().is_none());
2237 }
2238
2239 #[test]
2240 fn embedding_cache_roundtrip_batch_and_count() {
2241 let tmp = tempfile::tempdir().unwrap();
2242 let store = MemoryStore::open_default(tmp.path()).unwrap();
2243
2244 let entries: Vec<(String, Vec<f32>)> = (0..5)
2245 .map(|i| (format!("ns:model:{i:016x}"), vec![i as f32; 8]))
2246 .collect();
2247 store.put_embedding_cache_batch(&entries).unwrap();
2248 assert_eq!(store.embedding_cache_count().unwrap(), 5);
2249
2250 let hit = store
2252 .get_embedding_cache("ns:model:0000000000000003")
2253 .unwrap();
2254 assert_eq!(hit.unwrap(), vec![3.0f32; 8]);
2255 assert!(
2256 store
2257 .get_embedding_cache("ns:model:absent")
2258 .unwrap()
2259 .is_none()
2260 );
2261
2262 let keys: Vec<String> = (0..6).map(|i| format!("ns:model:{i:016x}")).collect();
2264 let batch = store.get_embedding_cache_batch(&keys).unwrap();
2265 assert_eq!(batch.len(), 6);
2266 assert!(batch[0..5].iter().all(Option::is_some));
2267 assert!(batch[5].is_none());
2268
2269 store
2271 .put_embedding_cache("ns:model:0000000000000001", &[9.0; 8])
2272 .unwrap();
2273 assert_eq!(store.embedding_cache_count().unwrap(), 5);
2274 assert_eq!(
2275 store
2276 .get_embedding_cache("ns:model:0000000000000001")
2277 .unwrap()
2278 .unwrap(),
2279 vec![9.0f32; 8]
2280 );
2281 }
2282
2283 #[test]
2284 fn embedding_cache_survives_store_reopen() {
2285 let tmp = tempfile::tempdir().unwrap();
2287 {
2288 let store = MemoryStore::open_default(tmp.path()).unwrap();
2289 store
2290 .put_embedding_cache("onnx:bge-small:abc", &[0.5; 384])
2291 .unwrap();
2292 }
2293 let reopened = MemoryStore::open_default(tmp.path()).unwrap();
2294 let cached = reopened.get_embedding_cache("onnx:bge-small:abc").unwrap();
2295 assert_eq!(cached.unwrap(), vec![0.5f32; 384]);
2296 }
2297
2298 #[test]
2299 fn content_hash_is_sha256() {
2300 let hash1 = content_hash("test content");
2301 let hash2 = content_hash("test content");
2302 let hash3 = content_hash("different content");
2303
2304 assert_eq!(hash1, hash2, "same content should produce same hash");
2305 assert_ne!(
2306 hash1, hash3,
2307 "different content should produce different hash"
2308 );
2309 assert_eq!(hash1.len(), 64, "SHA-256 hex should be 64 chars");
2310 }
2311
2312 #[test]
2313 fn query_by_tag_uses_index() {
2314 let tmp = tempfile::tempdir().unwrap();
2315 let store = MemoryStore::open_default(tmp.path()).unwrap();
2316
2317 let mem1 = Memory::new(Galaxy::Codex, "tagged".into())
2318 .with_tags(vec!["rust".into(), "memory".into()]);
2319 let mem2 = Memory::new(Galaxy::Codex, "other".into()).with_tags(vec!["python".into()]);
2320 store.put(Galaxy::Codex, &mem1).unwrap();
2321 store.put(Galaxy::Codex, &mem2).unwrap();
2322
2323 let query = MemoryQuery::new().with_tags(vec!["rust".into()]);
2324 let results = store.query(Galaxy::Codex, &query).unwrap();
2325 assert_eq!(results.len(), 1);
2326 assert_eq!(results[0].content, "tagged");
2327 }
2328
2329 #[test]
2330 fn query_by_importance_uses_index() {
2331 let tmp = tempfile::tempdir().unwrap();
2332 let store = MemoryStore::open_default(tmp.path()).unwrap();
2333
2334 store
2335 .put(
2336 Galaxy::Codex,
2337 &Memory::new(Galaxy::Codex, "low".into()).with_importance(0.1),
2338 )
2339 .unwrap();
2340 store
2341 .put(
2342 Galaxy::Codex,
2343 &Memory::new(Galaxy::Codex, "mid".into()).with_importance(0.5),
2344 )
2345 .unwrap();
2346 store
2347 .put(
2348 Galaxy::Codex,
2349 &Memory::new(Galaxy::Codex, "high".into()).with_importance(0.9),
2350 )
2351 .unwrap();
2352
2353 let query = MemoryQuery::new().with_importance_range(0.4, 0.6);
2354 let results = store.query(Galaxy::Codex, &query).unwrap();
2355 assert_eq!(results.len(), 1);
2356 assert_eq!(results[0].content, "mid");
2357 }
2358
2359 #[test]
2360 fn query_by_time_uses_index() {
2361 let tmp = tempfile::tempdir().unwrap();
2362 let store = MemoryStore::open_default(tmp.path()).unwrap();
2363
2364 let t0 = chrono::Utc::now();
2365 std::thread::sleep(std::time::Duration::from_millis(10));
2366 let mem = Memory::new(Galaxy::Codex, "timed".into());
2367 store.put(Galaxy::Codex, &mem).unwrap();
2368 std::thread::sleep(std::time::Duration::from_millis(10));
2369 let t2 = chrono::Utc::now();
2370
2371 let query = MemoryQuery::new().with_time_range(t0, t2);
2372 let results = store.query(Galaxy::Codex, &query).unwrap();
2373 assert_eq!(results.len(), 1);
2374 assert_eq!(results[0].content, "timed");
2375 }
2376
2377 #[test]
2378 fn delete_removes_index_entries() {
2379 let tmp = tempfile::tempdir().unwrap();
2380 let store = MemoryStore::open_default(tmp.path()).unwrap();
2381
2382 let mem = Memory::new(Galaxy::Codex, "test".into())
2383 .with_tags(vec!["tag1".into()])
2384 .with_importance(0.7);
2385 let id = mem.metadata.id;
2386 let hash = mem.metadata.content_hash.clone();
2387 store.put(Galaxy::Codex, &mem).unwrap();
2388
2389 assert!(
2391 store
2392 .find_by_content_hash(Galaxy::Codex, &hash)
2393 .unwrap()
2394 .is_some()
2395 );
2396
2397 store.delete(Galaxy::Codex, id).unwrap();
2399
2400 assert!(
2402 store
2403 .find_by_content_hash(Galaxy::Codex, &hash)
2404 .unwrap()
2405 .is_none()
2406 );
2407
2408 let query = MemoryQuery::new().with_tags(vec!["tag1".into()]);
2410 let results = store.query(Galaxy::Codex, &query).unwrap();
2411 assert!(results.is_empty());
2412 }
2413
2414 #[test]
2415 fn put_batch_updates_indexes() {
2416 let tmp = tempfile::tempdir().unwrap();
2417 let store = MemoryStore::open_default(tmp.path()).unwrap();
2418
2419 let memories: Vec<Memory> = (0..5)
2420 .map(|i| {
2421 Memory::new(Galaxy::Codex, format!("batch-{i}"))
2422 .with_tags(vec![format!("tag{i}")])
2423 .with_importance(i as f32 * 0.2)
2424 })
2425 .collect();
2426 store.put_batch(Galaxy::Codex, &memories).unwrap();
2427
2428 for i in 0..5 {
2429 let query = MemoryQuery::new().with_tags(vec![format!("tag{i}")]);
2430 let results = store.query(Galaxy::Codex, &query).unwrap();
2431 assert_eq!(results.len(), 1, "tag{i} should have 1 result");
2432 }
2433 }
2434
2435 #[test]
2436 fn find_by_content_hash_indexed_matches_scan() {
2437 let tmp = tempfile::tempdir().unwrap();
2438 let store = MemoryStore::open_default(tmp.path()).unwrap();
2439
2440 let mem = Memory::new(Galaxy::Codex, "dedup test".into());
2441 let id = mem.metadata.id;
2442 let hash = mem.metadata.content_hash.clone();
2443 store.put(Galaxy::Codex, &mem).unwrap();
2444
2445 let indexed = store.find_by_content_hash(Galaxy::Codex, &hash).unwrap();
2446 let scanned = store
2447 .find_by_content_hash_scan(Galaxy::Codex, &hash)
2448 .unwrap();
2449
2450 assert_eq!(indexed, scanned);
2451 assert_eq!(indexed, Some(id));
2452 }
2453
2454 #[test]
2455 fn put_dedup_uses_index() {
2456 let tmp = tempfile::tempdir().unwrap();
2457 let store = MemoryStore::open_default(tmp.path()).unwrap();
2458
2459 let mem1 = Memory::new(Galaxy::Codex, "duplicate content".into());
2460 let id1 = store.put_dedup(Galaxy::Codex, &mem1).unwrap();
2461
2462 let mem2 = Memory::new(Galaxy::Codex, "duplicate content".into());
2463 let id2 = store.put_dedup(Galaxy::Codex, &mem2).unwrap();
2464
2465 assert_eq!(id1, id2, "dedup should return same ID for same content");
2466 assert_eq!(store.count(Galaxy::Codex).unwrap(), 1);
2467 }
2468
2469 #[test]
2470 fn put_semantic_updates_coord5d() {
2471 let tmp = tempfile::tempdir().unwrap();
2472 let store = MemoryStore::open_default(tmp.path()).unwrap();
2473
2474 let mut mem = Memory::new(
2475 Galaxy::Codex,
2476 "The algorithm computes data using a systematic method".to_string(),
2477 );
2478 let original_coord = mem.metadata.coord5d.clone();
2479 store.put_semantic(Galaxy::Codex, &mut mem).unwrap();
2480
2481 assert_ne!(
2483 mem.metadata.coord5d.x, original_coord.x,
2484 "semantic encoding should change x"
2485 );
2486 assert_ne!(
2487 mem.metadata.coord5d.y, original_coord.y,
2488 "semantic encoding should change y"
2489 );
2490
2491 let retrieved = store.get(Galaxy::Codex, mem.metadata.id).unwrap().unwrap();
2493 assert_eq!(retrieved.metadata.coord5d.x, mem.metadata.coord5d.x);
2494 }
2495
2496 #[test]
2497 fn put_semantic_preserves_temporal_and_importance() {
2498 let tmp = tempfile::tempdir().unwrap();
2499 let store = MemoryStore::open_default(tmp.path()).unwrap();
2500
2501 let mut mem = Memory::new(Galaxy::Codex, "test content".into()).with_importance(0.8);
2502 mem.metadata.coord5d.w = 0.6;
2503 store.put_semantic(Galaxy::Codex, &mut mem).unwrap();
2504
2505 assert!((mem.metadata.coord5d.w - 0.6).abs() < f32::EPSILON);
2506 assert!((mem.metadata.coord5d.v - 0.8).abs() < f32::EPSILON);
2507 }
2508
2509 #[test]
2510 fn find_similar_returns_nearest_first() {
2511 let tmp = tempfile::tempdir().unwrap();
2512 let store = MemoryStore::open_default(tmp.path()).unwrap();
2513
2514 let mut logic_mem = Memory::new(
2515 Galaxy::Codex,
2516 "The algorithm computes data using systematic logic and analysis".to_string(),
2517 );
2518 store.put_semantic(Galaxy::Codex, &mut logic_mem).unwrap();
2519
2520 let mut emotion_mem = Memory::new(
2521 Galaxy::Codex,
2522 "I feel love and joy with deep passion and empathy in my heart".to_string(),
2523 );
2524 store.put_semantic(Galaxy::Codex, &mut emotion_mem).unwrap();
2525
2526 let results = store
2528 .find_similar(Galaxy::Codex, "algorithm data systematic method", 10)
2529 .unwrap();
2530 assert!(!results.is_empty());
2531 assert_eq!(results[0].0.metadata.id, logic_mem.metadata.id);
2532
2533 let results = store
2535 .find_similar(Galaxy::Codex, "love joy passion heart feeling", 10)
2536 .unwrap();
2537 assert!(!results.is_empty());
2538 assert_eq!(results[0].0.metadata.id, emotion_mem.metadata.id);
2539 }
2540
2541 #[test]
2542 fn find_similar_empty_galaxy() {
2543 let tmp = tempfile::tempdir().unwrap();
2544 let store = MemoryStore::open_default(tmp.path()).unwrap();
2545
2546 let results = store.find_similar(Galaxy::Codex, "anything", 10).unwrap();
2547 assert!(results.is_empty());
2548 }
2549
2550 #[test]
2551 fn find_similar_respects_limit() {
2552 let tmp = tempfile::tempdir().unwrap();
2553 let store = MemoryStore::open_default(tmp.path()).unwrap();
2554
2555 for i in 0..5 {
2556 let mut mem = Memory::new(Galaxy::Codex, format!("algorithm data method {i}"));
2557 store.put_semantic(Galaxy::Codex, &mut mem).unwrap();
2558 }
2559
2560 let results = store
2561 .find_similar(Galaxy::Codex, "algorithm data", 3)
2562 .unwrap();
2563 assert_eq!(results.len(), 3);
2564 }
2565
2566 #[test]
2567 fn semantic_encoder_accessible() {
2568 let tmp = tempfile::tempdir().unwrap();
2569 let store = MemoryStore::open_default(tmp.path()).unwrap();
2570
2571 let scores = store.semantic_encoder().encode("algorithm data logic");
2572 assert!(scores.x < 0.5);
2574 }
2575
2576 #[test]
2577 fn put_raw_batch_writes_atomically() {
2578 let tmp = tempfile::tempdir().unwrap();
2579 let store = MemoryStore::open_default(tmp.path()).unwrap();
2580
2581 let entries: &[(&[u8], &[u8])] =
2582 &[(b"key1", b"val1"), (b"key2", b"val2"), (b"key3", b"val3")];
2583 store.put_raw_batch(Galaxy::Karma, entries).unwrap();
2584
2585 assert_eq!(
2586 store.get_raw(Galaxy::Karma, b"key1").unwrap().unwrap(),
2587 b"val1"
2588 );
2589 assert_eq!(
2590 store.get_raw(Galaxy::Karma, b"key2").unwrap().unwrap(),
2591 b"val2"
2592 );
2593 assert_eq!(
2594 store.get_raw(Galaxy::Karma, b"key3").unwrap().unwrap(),
2595 b"val3"
2596 );
2597 }
2598
2599 #[test]
2600 fn put_raw_batch_empty_is_noop() {
2601 let tmp = tempfile::tempdir().unwrap();
2602 let store = MemoryStore::open_default(tmp.path()).unwrap();
2603
2604 store.put_raw_batch(Galaxy::Karma, &[]).unwrap();
2605 assert_eq!(store.count(Galaxy::Karma).unwrap(), 0);
2606 }
2607
2608 #[test]
2609 fn entry_limit_rejects_excess_writes() {
2610 let tmp = tempfile::tempdir().unwrap();
2611 let store = MemoryStore::open_default(tmp.path())
2612 .unwrap()
2613 .with_entry_limit(3);
2614
2615 for i in 0..3 {
2616 let mem = Memory::new(Galaxy::Codex, format!("memory {i}"));
2617 store.put(Galaxy::Codex, &mem).unwrap();
2618 }
2619
2620 let mem = Memory::new(Galaxy::Codex, "overflow memory".to_string());
2622 let result = store.put(Galaxy::Codex, &mem);
2623 assert!(result.is_err(), "write beyond limit should be rejected");
2624 let err_msg = result.unwrap_err().to_string();
2625 assert!(
2626 err_msg.contains("entry limit reached"),
2627 "error should mention entry limit: {err_msg}"
2628 );
2629 assert_eq!(store.count(Galaxy::Codex).unwrap(), 3);
2630 }
2631
2632 #[test]
2633 fn entry_limit_per_galaxy_independent() {
2634 let tmp = tempfile::tempdir().unwrap();
2635 let store = MemoryStore::open_default(tmp.path())
2636 .unwrap()
2637 .with_entry_limit(2);
2638
2639 for i in 0..2 {
2641 let mem = Memory::new(Galaxy::Codex, format!("codex {i}"));
2642 store.put(Galaxy::Codex, &mem).unwrap();
2643 }
2644
2645 let mem = Memory::new(Galaxy::Research, "science memory".to_string());
2647 let result = store.put(Galaxy::Research, &mem);
2648 assert!(
2649 result.is_ok(),
2650 "different galaxy should not be affected by limit"
2651 );
2652 }
2653
2654 #[test]
2655 fn entry_limit_none_allows_unlimited() {
2656 let tmp = tempfile::tempdir().unwrap();
2657 let store = MemoryStore::open_default(tmp.path()).unwrap();
2658
2659 for i in 0..50 {
2661 let mem = Memory::new(Galaxy::Codex, format!("memory {i}"));
2662 store.put(Galaxy::Codex, &mem).unwrap();
2663 }
2664 assert_eq!(store.count(Galaxy::Codex).unwrap(), 50);
2665 }
2666
2667 #[test]
2668 fn map_full_error_is_graceful() {
2669 let tmp = tempfile::tempdir().unwrap();
2673 let store = MemoryStore::open(tmp.path(), 512 * 1024).unwrap();
2674
2675 let mut written = 0;
2677 let mut got_map_full = false;
2678 for i in 0..1000 {
2679 let mem = Memory::new(
2680 Galaxy::Codex,
2681 format!("memory content {i} {}", "with padding ".repeat(50)),
2682 );
2683 match store.put(Galaxy::Codex, &mem) {
2684 Ok(()) => written += 1,
2685 Err(e) => {
2686 let msg = e.to_string();
2687 if msg.contains("map full") {
2688 got_map_full = true;
2689 break;
2690 }
2691 break;
2693 }
2694 }
2695 }
2696
2697 assert!(
2698 got_map_full || written < 1000,
2699 "should eventually hit map full or error"
2700 );
2701 assert!(written > 0, "should have written at least some memories");
2702 }
2703
2704 #[test]
2705 fn test_find_across_galaxies() {
2706 let tmp = tempfile::tempdir().unwrap();
2707 let store = MemoryStore::open_default(tmp.path()).unwrap();
2708
2709 let mem = Memory::new(Galaxy::Research, "Cross-galaxy research memo".into());
2710 let id = mem.metadata.id;
2711 store.put(Galaxy::Research, &mem).unwrap();
2712
2713 let found = store.find_across_galaxies(id).unwrap();
2714 assert!(found.is_some());
2715 let (galaxy, retrieved) = found.unwrap();
2716 assert_eq!(galaxy, Galaxy::Research);
2717 assert_eq!(retrieved.metadata.id, id);
2718 assert_eq!(retrieved.content, "Cross-galaxy research memo");
2719
2720 assert!(
2722 store
2723 .find_across_galaxies(uuid::Uuid::new_v4())
2724 .unwrap()
2725 .is_none()
2726 );
2727 }
2728}