use lmdb::{
Cursor, Database, DatabaseFlags, Environment, EnvironmentFlags, RwTransaction, Transaction,
WriteFlags,
};
use std::collections::HashMap;
use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, RwLock};
use wm_core::{CoreError, Galaxy, Result};
#[cfg(unix)]
use std::os::unix::fs::PermissionsExt;
use crate::episodic::EpisodicStore;
use crate::indexes::IndexDbs;
use crate::memory::{Memory, MemoryId, decode_embedding, encode_embedding};
use crate::semantic::SemanticEncoder;
#[derive(Debug, Clone, Default)]
pub struct MemoryQuery {
pub tags: Vec<String>,
pub min_importance: Option<f32>,
pub max_importance: Option<f32>,
pub created_after: Option<chrono::DateTime<chrono::Utc>>,
pub created_before: Option<chrono::DateTime<chrono::Utc>>,
pub content_substring: Option<String>,
pub limit: usize,
}
impl MemoryQuery {
#[must_use]
pub fn new() -> Self {
Self {
limit: 100,
..Default::default()
}
}
#[must_use]
pub fn with_tags(mut self, tags: Vec<String>) -> Self {
self.tags = tags;
self
}
#[must_use]
pub const fn with_importance_range(mut self, min: f32, max: f32) -> Self {
self.min_importance = Some(min);
self.max_importance = Some(max);
self
}
#[must_use]
pub const fn with_time_range(
mut self,
after: chrono::DateTime<chrono::Utc>,
before: chrono::DateTime<chrono::Utc>,
) -> Self {
self.created_after = Some(after);
self.created_before = Some(before);
self
}
#[must_use]
pub const fn with_created_after(mut self, after: chrono::DateTime<chrono::Utc>) -> Self {
self.created_after = Some(after);
self
}
#[must_use]
pub const fn with_created_before(mut self, before: chrono::DateTime<chrono::Utc>) -> Self {
self.created_before = Some(before);
self
}
#[must_use]
pub const fn with_limit(mut self, limit: usize) -> Self {
self.limit = limit;
self
}
#[must_use]
pub fn with_content_substring(mut self, substring: impl Into<String>) -> Self {
self.content_substring = Some(substring.into().to_lowercase());
self
}
#[must_use]
pub fn matches(&self, mem: &Memory) -> bool {
if !self.tags.is_empty() {
for tag in &self.tags {
if !mem.metadata.tags.iter().any(|t| t == tag) {
return false;
}
}
}
if let Some(min) = self.min_importance {
if mem.metadata.importance < min {
return false;
}
}
if let Some(max) = self.max_importance {
if mem.metadata.importance > max {
return false;
}
}
if let Some(after) = self.created_after {
if mem.metadata.created_at < after {
return false;
}
}
if let Some(before) = self.created_before {
if mem.metadata.created_at > before {
return false;
}
}
if let Some(sub) = &self.content_substring {
if !mem.content.to_lowercase().contains(sub) {
return false;
}
}
true
}
}
pub struct MemoryStore {
path: std::path::PathBuf,
env: Environment,
index_dbs: IndexDbs,
semantic_encoder: SemanticEncoder,
max_entries_per_galaxy: Option<usize>,
mutation_count: AtomicU64,
episodic_db: Database,
episodic_terms_v2_db: Database,
embedding_cache_db: Database,
revisions_db: Database,
attestations_db: Database,
pub(crate) cold_storage_db: Database,
session_sequences_db: Option<Database>,
index_pending_db: Option<Database>,
keyring_db: Option<Database>,
at_rest: Option<crate::at_rest::AtRestState>,
episodic_term_cache: std::sync::Arc<RwLock<HashMap<String, Vec<uuid::Uuid>>>>,
episodic_embedder:
std::sync::OnceLock<Option<Arc<dyn crate::embedder::Embedder + Send + Sync>>>,
episodic_sidecar_ensured: std::sync::OnceLock<()>,
episodic_aliases: std::sync::OnceLock<Option<crate::episodic_keys::AdaptiveAliases>>,
episodic_enrichment: std::sync::OnceLock<Option<crate::enrichment::VocabularyEnrichment>>,
}
impl MemoryStore {
#[cfg(unix)]
pub fn probe_write_lock(store_root: &Path) -> std::io::Result<()> {
use rustix::fs::{FlockOperation, fcntl_lock};
let lock_path = store_root.join("lmdb").join("lock.mdb");
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.open(&lock_path)?;
match fcntl_lock(&file, FlockOperation::NonBlockingLockExclusive) {
Ok(()) => Ok(()),
Err(rustix::io::Errno::AGAIN | rustix::io::Errno::ACCESS) => Err(std::io::Error::new(
std::io::ErrorKind::WouldBlock,
"LMDB writer lock held by another process",
)),
Err(e) => Err(e.into()),
}
}
#[cfg(not(unix))]
pub fn probe_write_lock(_store_root: &Path) -> std::io::Result<()> {
Ok(())
}
pub fn open(path: impl AsRef<Path>, map_size: usize) -> Result<Self> {
Self::open_with_at_rest(path, map_size, &crate::at_rest::AtRestConfig::from_env()?)
}
pub fn open_with_at_rest(
path: impl AsRef<Path>,
map_size: usize,
at_rest_config: &crate::at_rest::AtRestConfig,
) -> Result<Self> {
let path = path.as_ref().to_path_buf();
std::fs::create_dir_all(&path)
.map_err(|e| CoreError::Memory(format!("Cannot create store dir: {e}")))?;
#[cfg(unix)]
{
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o700))
.map_err(|e| CoreError::Memory(format!("Cannot set store dir permissions: {e}")))?;
}
let env = Environment::new()
.set_map_size(map_size)
.set_max_dbs(64)
.open(&path)
.map_err(|e| CoreError::Memory(format!("LMDB open failed: {e}")))?;
for galaxy in Galaxy::all() {
let db = env
.create_db(Some(galaxy.db_name()), DatabaseFlags::default())
.map_err(|e| {
CoreError::Memory(format!(
"LMDB create_db failed for {}: {e}",
galaxy.db_name()
))
})?;
let _ = db;
}
for (name, flags) in crate::indexes::INDEX_DBS {
let db = env
.create_db(Some(name), *flags)
.map_err(|e| CoreError::Memory(format!("LMDB create_db failed for {name}: {e}")))?;
let _ = db;
}
let index_dbs = IndexDbs::open(&env)?;
let episodic_db = env
.create_db(Some("episodic_records"), DatabaseFlags::default())
.map_err(|e| {
CoreError::Memory(format!("LMDB create_db failed for episodic_records: {e}"))
})?;
let episodic_terms_v2_db = env
.create_db(Some("episodic_terms_v2"), DatabaseFlags::DUP_SORT)
.map_err(|e| {
CoreError::Memory(format!("LMDB create_db failed for episodic_terms_v2: {e}"))
})?;
let embedding_cache_db = env
.create_db(Some("embedding_cache"), DatabaseFlags::default())
.map_err(|e| {
CoreError::Memory(format!("LMDB create_db failed for embedding_cache: {e}"))
})?;
let revisions_db = env
.create_db(Some("revisions"), DatabaseFlags::default())
.map_err(|e| CoreError::Memory(format!("LMDB create_db failed for revisions: {e}")))?;
let attestations_db = env
.create_db(
Some(crate::attestation::ATTESTATIONS_DB),
DatabaseFlags::default(),
)
.map_err(|e| {
CoreError::Memory(format!("LMDB create_db failed for attestations: {e}"))
})?;
let cold_storage_db = env
.create_db(Some("cold_storage"), DatabaseFlags::default())
.map_err(|e| {
CoreError::Memory(format!("LMDB create_db failed for cold_storage: {e}"))
})?;
let session_sequences_db = env
.create_db(Some("session_sequences"), DatabaseFlags::default())
.map_err(|e| {
CoreError::Memory(format!("LMDB create_db failed for session_sequences: {e}"))
})?;
let index_pending_db = env
.create_db(Some("index_pending"), DatabaseFlags::default())
.map_err(|e| {
CoreError::Memory(format!("LMDB create_db failed for index_pending: {e}"))
})?;
let (keyring_db, at_rest) = crate::at_rest::open_at_rest(&env, &path, at_rest_config)?;
Ok(Self {
path,
env,
index_dbs,
semantic_encoder: SemanticEncoder::new(),
max_entries_per_galaxy: None,
mutation_count: AtomicU64::new(0),
episodic_db,
episodic_terms_v2_db,
embedding_cache_db,
revisions_db,
attestations_db,
cold_storage_db,
session_sequences_db: Some(session_sequences_db),
index_pending_db: Some(index_pending_db),
keyring_db,
at_rest,
episodic_term_cache: std::sync::Arc::new(RwLock::new(HashMap::new())),
episodic_embedder: std::sync::OnceLock::new(),
episodic_sidecar_ensured: std::sync::OnceLock::new(),
episodic_aliases: std::sync::OnceLock::new(),
episodic_enrichment: std::sync::OnceLock::new(),
})
}
#[must_use]
pub fn at_rest_status(&self) -> crate::at_rest::AtRestStatus {
if let Some(state) = &self.at_rest {
return crate::at_rest::AtRestStatus::Present(state.status());
}
match &self.keyring_db {
None => crate::at_rest::AtRestStatus::Absent,
Some(db) => crate::at_rest::read_status(&self.env, *db, &self.path),
}
}
#[must_use]
pub const fn at_rest_state(&self) -> Option<&crate::at_rest::AtRestState> {
self.at_rest.as_ref()
}
#[must_use]
pub const fn keyring_db(&self) -> Option<Database> {
self.keyring_db
}
pub(crate) fn record_cipher(&self, galaxy: Galaxy) -> Option<&[u8; 32]> {
self.at_rest
.as_ref()
.and_then(|state| state.galaxy_dek(galaxy.db_name()))
}
pub(crate) fn encode_record_value(&self, galaxy: Galaxy, memory: &Memory) -> Result<Vec<u8>> {
let plaintext = rmp_serde::to_vec_named(memory)
.map_err(|e| CoreError::Memory(format!("serialize failed: {e}")))?;
let Some(key) = self.record_cipher(galaxy) else {
return Ok(plaintext);
};
crate::codec::seal_record(
&plaintext,
key,
galaxy.db_name(),
memory.metadata.id.as_bytes(),
memory.metadata.version,
)
.map_err(|e| CoreError::Memory(format!("at-rest seal failed: {e}")))
}
pub(crate) fn decode_record_value(
&self,
galaxy: Galaxy,
key_bytes: &[u8],
value: &[u8],
) -> Result<Memory> {
if crate::codec::is_sealed_record(value) {
let Some(dek) = self.record_cipher(galaxy) else {
return Err(CoreError::Memory(format!(
"sealed record in {} but no at-rest key is loaded (WM_AT_REST_MODE off?)",
galaxy.db_name()
)));
};
let record_id: [u8; 16] = key_bytes
.try_into()
.map_err(|_| CoreError::Memory("sealed record key is not a 16-byte id".into()))?;
let opened = crate::codec::open_record(value, dek, galaxy.db_name(), &record_id)
.map_err(|e| CoreError::Memory(format!("at-rest open failed: {e}")))?;
return crate::codec::decode(&opened)
.map_err(|e| CoreError::Memory(format!("deserialize failed: {e}")));
}
crate::codec::decode(value)
.map_err(|e| CoreError::Memory(format!("deserialize failed: {e}")))
}
pub fn open_readonly_bounded(
path: impl AsRef<Path>,
timeout: std::time::Duration,
) -> Result<Option<Self>> {
let path = path.as_ref().to_path_buf();
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
let result = Self::open_readonly(&path);
let _ = tx.send(result);
});
match rx.recv_timeout(timeout) {
Ok(result) => result.map(Some),
Err(_) => Ok(None),
}
}
pub fn open_default_bounded(
path: impl AsRef<Path>,
timeout: std::time::Duration,
) -> Result<Option<Self>> {
let path = path.as_ref().to_path_buf();
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
let result = Self::open_default(&path);
let _ = tx.send(result);
});
match rx.recv_timeout(timeout) {
Ok(result) => result.map(Some),
Err(_) => Ok(None),
}
}
#[must_use]
pub fn default_map_size() -> usize {
let platform_default = if cfg!(windows) {
256 * 1024 * 1024
} else {
4 * 1024 * 1024 * 1024
};
std::env::var("WM_DEFAULT_MAP_SIZE")
.ok()
.and_then(|v| v.parse::<usize>().ok())
.filter(|&v| v > 0)
.unwrap_or(platform_default)
}
pub fn open_default(path: impl AsRef<Path>) -> Result<Self> {
let size = Self::default_map_size();
Self::open(path, size)
}
pub fn open_readonly(path: impl AsRef<Path>) -> Result<Self> {
let path = path.as_ref().to_path_buf();
if !path.is_dir() {
return Err(CoreError::Memory(format!(
"Read-only LMDB store directory does not exist: {}",
path.display()
)));
}
if !path.join("data.mdb").is_file() {
return Err(CoreError::Memory(format!(
"Read-only LMDB store is missing data.mdb: {}",
path.display()
)));
}
let env = Environment::new()
.set_max_dbs(32)
.set_flags(EnvironmentFlags::READ_ONLY)
.open(&path)
.map_err(|e| CoreError::Memory(format!("Read-only LMDB open failed: {e}")))?;
let index_dbs = IndexDbs::open(&env)?;
let open_named = |name: &str| {
env.open_db(Some(name)).map_err(|e| {
CoreError::Memory(format!("Read-only LMDB missing database {name}: {e}"))
})
};
let episodic_db = open_named("episodic_records")?;
let episodic_terms_v2_db = open_named("episodic_terms_v2")?;
let embedding_cache_db = open_named("embedding_cache")?;
let revisions_db = open_named("revisions")?;
let attestations_db = open_named(crate::attestation::ATTESTATIONS_DB)?;
let cold_storage_db = open_named("cold_storage")?;
let session_sequences_db = env.open_db(Some("session_sequences")).ok();
let index_pending_db = env.open_db(Some("index_pending")).ok();
let keyring_db = crate::at_rest::open_keyring_optional(&env)?;
Ok(Self {
path,
env,
index_dbs,
semantic_encoder: SemanticEncoder::new(),
max_entries_per_galaxy: None,
mutation_count: AtomicU64::new(0),
episodic_db,
episodic_terms_v2_db,
embedding_cache_db,
revisions_db,
attestations_db,
cold_storage_db,
session_sequences_db,
index_pending_db,
keyring_db,
at_rest: None,
episodic_term_cache: std::sync::Arc::new(RwLock::new(HashMap::new())),
episodic_embedder: std::sync::OnceLock::new(),
episodic_sidecar_ensured: std::sync::OnceLock::new(),
episodic_aliases: std::sync::OnceLock::new(),
episodic_enrichment: std::sync::OnceLock::new(),
})
}
const NAMED_DBIS: [&'static str; 6] = [
"episodic_records",
"episodic_terms_v2",
"embedding_cache",
"revisions",
crate::attestation::ATTESTATIONS_DB,
"cold_storage",
];
pub fn open_inspection(path: impl AsRef<Path>) -> Result<Self> {
let path = path.as_ref().to_path_buf();
if !path.is_dir() {
return Err(CoreError::Memory(format!(
"Read-only LMDB store directory does not exist: {}",
path.display()
)));
}
if !path.join("data.mdb").is_file() {
return Err(CoreError::Memory(format!(
"Read-only LMDB store is missing data.mdb: {}",
path.display()
)));
}
let env = Environment::new()
.set_max_dbs(32)
.set_flags(EnvironmentFlags::READ_ONLY | EnvironmentFlags::NO_LOCK)
.open(&path)
.map_err(|e| CoreError::Memory(format!("Inspection LMDB open failed: {e}")))?;
let index_dbs = IndexDbs::open(&env)?;
let open_named = |name: &str| {
env.open_db(Some(name)).map_err(|e| {
CoreError::Memory(format!("Inspection LMDB missing database {name}: {e}"))
})
};
let episodic_db = open_named("episodic_records")?;
let episodic_terms_v2_db = open_named("episodic_terms_v2")?;
let embedding_cache_db = open_named("embedding_cache")?;
let revisions_db = open_named("revisions")?;
let attestations_db = open_named(crate::attestation::ATTESTATIONS_DB)?;
let cold_storage_db = open_named("cold_storage")?;
let session_sequences_db = env.open_db(Some("session_sequences")).ok();
let index_pending_db = env.open_db(Some("index_pending")).ok();
let keyring_db = crate::at_rest::open_keyring_optional(&env)?;
Ok(Self {
path,
env,
index_dbs,
semantic_encoder: SemanticEncoder::new(),
max_entries_per_galaxy: None,
mutation_count: AtomicU64::new(0),
episodic_db,
episodic_terms_v2_db,
embedding_cache_db,
revisions_db,
attestations_db,
cold_storage_db,
session_sequences_db,
index_pending_db,
keyring_db,
at_rest: None,
episodic_term_cache: std::sync::Arc::new(RwLock::new(HashMap::new())),
episodic_embedder: std::sync::OnceLock::new(),
episodic_sidecar_ensured: std::sync::OnceLock::new(),
episodic_aliases: std::sync::OnceLock::new(),
episodic_enrichment: std::sync::OnceLock::new(),
})
}
pub fn ensure_schema(path: impl AsRef<Path>) -> Result<Vec<String>> {
let path = path.as_ref().to_path_buf();
if !path.is_dir() {
return Err(CoreError::Memory(format!(
"Store directory does not exist: {}",
path.display()
)));
}
let expected = || {
Galaxy::all()
.into_iter()
.map(|galaxy| galaxy.db_name().to_string())
.chain(
crate::indexes::INDEX_DBS
.iter()
.map(|(name, _)| (*name).to_string()),
)
.chain(Self::NAMED_DBIS.iter().map(|name| (*name).to_string()))
};
let missing: Vec<String> = {
let env = Environment::new()
.set_max_dbs(64)
.set_flags(EnvironmentFlags::READ_ONLY)
.open(&path)
.map_err(|e| CoreError::Memory(format!("Read-only LMDB open failed: {e}")))?;
expected()
.filter(|name| env.open_db(Some(name.as_str())).is_err())
.collect()
};
if missing.is_empty() {
return Ok(missing);
}
let store = Self::open_default(&path)?;
drop(store);
Ok(missing)
}
#[must_use]
pub const fn with_entry_limit(mut self, limit: usize) -> Self {
self.max_entries_per_galaxy = Some(limit);
self
}
pub fn path(&self) -> &Path {
&self.path
}
pub const fn env(&self) -> &Environment {
&self.env
}
pub fn mutation_count(&self) -> u64 {
self.mutation_count.load(Ordering::Relaxed)
}
pub const fn index_dbs(&self) -> &IndexDbs {
&self.index_dbs
}
pub const fn semantic_encoder(&self) -> &SemanticEncoder {
&self.semantic_encoder
}
pub fn episodic_sidecar_health(&self) -> Result<(u64, bool)> {
let view = EpisodicStore::new(
&self.env,
self.episodic_db,
self.episodic_terms_v2_db,
self.episodic_term_cache.clone(),
&self.mutation_count,
);
let empty = view.sidecar_is_empty()?;
if !empty {
return Ok((view.record_count()?, false));
}
let indexable = view
.scan(None, usize::MAX)?
.iter()
.filter(|record| !record.is_private && !record.model_exclude)
.count() as u64;
Ok((indexable, true))
}
pub fn episodic_record_count(&self) -> Result<u64> {
let view = EpisodicStore::new(
&self.env,
self.episodic_db,
self.episodic_terms_v2_db,
self.episodic_term_cache.clone(),
&self.mutation_count,
);
view.record_count()
}
fn ensure_episodic_sidecar(&self) {
if self.episodic_sidecar_ensured.get().is_some() {
return;
}
let _ = self.episodic_sidecar_ensured.set(());
let view = EpisodicStore::new(
&self.env,
self.episodic_db,
self.episodic_terms_v2_db,
self.episodic_term_cache.clone(),
&self.mutation_count,
);
let needs_rebuild = matches!(
(view.sidecar_is_empty(), view.record_count()),
(Ok(true), Ok(n)) if n > 0
);
if needs_rebuild {
match view.rebuild_sidecar() {
Ok(n) => tracing::info!("episodic sidecar rebuilt from {n} records"),
Err(e) => {
tracing::warn!("episodic sidecar rebuild failed: {e}");
}
}
}
}
#[must_use]
pub fn episodic(&self) -> EpisodicStore<'_> {
self.ensure_episodic_sidecar();
let mut store = EpisodicStore::new(
&self.env,
self.episodic_db,
self.episodic_terms_v2_db,
self.episodic_term_cache.clone(),
&self.mutation_count,
);
if let Some(Some(embedder)) = self.episodic_embedder.get() {
store = store.with_embedder(embedder.clone());
}
if let Some(Some(aliases)) = self.episodic_aliases.get() {
store = store.with_adaptive_aliases(aliases.clone());
}
if let Some(Some(enrichment)) = self.episodic_enrichment.get() {
store = store.with_enrichment(enrichment.clone());
}
store
}
pub fn set_episodic_embedder(
&self,
embedder: Arc<dyn crate::embedder::Embedder + Send + Sync>,
) {
let _ = self.episodic_embedder.set(Some(embedder));
}
pub fn set_episodic_aliases(&self, aliases: crate::episodic_keys::AdaptiveAliases) {
let _ = self.episodic_aliases.set(Some(aliases));
}
pub fn set_episodic_enrichment(&self, enrichment: crate::enrichment::VocabularyEnrichment) {
let _ = self.episodic_enrichment.set(Some(enrichment));
}
pub fn galaxy_db(&self, galaxy: Galaxy) -> Result<Database> {
self.env.open_db(Some(galaxy.db_name())).map_err(|e| {
CoreError::Memory(format!("LMDB open_db failed for {}: {e}", galaxy.db_name()))
})
}
pub fn put(&self, galaxy: Galaxy, memory: &Memory) -> Result<()> {
if let Some(limit) = self.max_entries_per_galaxy {
let current = self.count(galaxy)?;
if current >= limit {
return Err(CoreError::Memory(format!(
"galaxy {} entry limit reached ({current}/{limit}), write rejected",
galaxy.db_name()
)));
}
}
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
self.put_in_txn(&mut tx, galaxy, memory)?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count.fetch_add(1, Ordering::Relaxed);
Ok(())
}
fn put_in_txn(&self, tx: &mut RwTransaction, galaxy: Galaxy, memory: &Memory) -> Result<()> {
let db = self.galaxy_db(galaxy)?;
let key = memory.metadata.id.as_bytes();
let val = self.encode_record_value(galaxy, memory)?;
let existing = tx
.get(db, key)
.ok()
.and_then(|bytes| self.decode_record_value(galaxy, key, bytes).ok());
match tx.put(db, key, &val, lmdb::WriteFlags::default()) {
Ok(()) => {}
Err(lmdb::Error::MapFull) => {
return Err(CoreError::Memory(format!(
"LMDB map full: galaxy {}, consider growing map size or pruning old memories",
galaxy.db_name()
)));
}
Err(e) => {
return Err(CoreError::Memory(format!("LMDB put failed: {e}")));
}
}
if let Some(existing) = existing {
self.index_dbs.remove(tx, galaxy, &existing)?;
}
self.index_dbs.add(tx, galaxy, memory)?;
Ok(())
}
pub fn put_session_turn<F>(&self, session_id: &str, build: F) -> Result<(u64, Memory)>
where
F: FnOnce(u64) -> Memory,
{
let seq_db = self.session_sequences_db.ok_or_else(|| {
CoreError::Memory(
"session_sequences DBI missing (legacy store opened read-only); \
a writable open repairs it via ensure_schema"
.into(),
)
})?;
if let Some(limit) = self.max_entries_per_galaxy {
let current = self.count(Galaxy::Sessions)?;
if current >= limit {
return Err(CoreError::Memory(format!(
"galaxy {} entry limit reached ({current}/{limit}), write rejected",
Galaxy::Sessions.db_name()
)));
}
}
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let key: &[u8] = session_id.as_bytes();
let current = match tx.get(seq_db, &key) {
Ok(bytes) => {
let arr: [u8; 8] = <[u8; 8]>::try_from(bytes).map_err(|_| {
CoreError::Memory(format!(
"session_sequences value for {session_id} is malformed \
({} bytes, want 8)",
bytes.len()
))
})?;
u64::from_be_bytes(arr)
}
Err(lmdb::Error::NotFound) => 0,
Err(e) => {
return Err(CoreError::Memory(format!(
"LMDB get failed (session_sequences): {e}"
)));
}
};
let next = current
.checked_add(1)
.ok_or_else(|| CoreError::Memory("session sequence overflow (u64)".into()))?;
tx.put(
seq_db,
&key,
&next.to_be_bytes(),
lmdb::WriteFlags::default(),
)
.map_err(|e| CoreError::Memory(format!("LMDB put failed (session_sequences): {e}")))?;
let memory = build(next);
self.put_in_txn(&mut tx, Galaxy::Sessions, &memory)?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count.fetch_add(1, Ordering::Relaxed);
Ok((next, memory))
}
pub fn last_session_sequence(&self, session_id: &str) -> Result<Option<u64>> {
let Some(seq_db) = self.session_sequences_db else {
return Ok(None);
};
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let key: &[u8] = session_id.as_bytes();
match tx.get(seq_db, &key) {
Ok(bytes) => {
let arr: [u8; 8] = <[u8; 8]>::try_from(bytes).map_err(|_| {
CoreError::Memory(format!(
"session_sequences value for {session_id} is malformed \
({} bytes, want 8)",
bytes.len()
))
})?;
Ok(Some(u64::from_be_bytes(arr)))
}
Err(lmdb::Error::NotFound) => Ok(None),
Err(e) => Err(CoreError::Memory(format!(
"LMDB get failed (session_sequences): {e}"
))),
}
}
pub fn mark_index_pending(&self, galaxy: &str, memory_id: &str, at_ms: i64) -> Result<()> {
let db = self.index_pending_db.ok_or_else(|| {
CoreError::Memory(
"index_pending DBI missing (store opened without the writable schema repair)"
.into(),
)
})?;
let value = serde_json::to_vec(&serde_json::json!({"galaxy": galaxy, "at_ms": at_ms}))
.map_err(|e| CoreError::Memory(format!("index_pending encode: {e}")))?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let key = memory_id.as_bytes();
tx.put(db, &key, &value, lmdb::WriteFlags::default())
.map_err(|e| CoreError::Memory(format!("LMDB put failed (index_pending): {e}")))?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(())
}
pub fn index_pending_entries(&self) -> Result<Vec<(String, String, i64)>> {
let Some(db) = self.index_pending_db else {
return Ok(Vec::new());
};
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let mut cursor = tx
.open_ro_cursor(db)
.map_err(|e| CoreError::Memory(format!("LMDB cursor failed (index_pending): {e}")))?;
let mut out: Vec<(String, String, i64)> = Vec::new();
for (key, value) in cursor.iter() {
let id = String::from_utf8_lossy(key).into_owned();
let parsed: serde_json::Value =
serde_json::from_slice(value).unwrap_or(serde_json::Value::Null);
let galaxy = parsed
.get("galaxy")
.and_then(serde_json::Value::as_str)
.unwrap_or_default()
.to_string();
let at_ms = parsed
.get("at_ms")
.and_then(serde_json::Value::as_i64)
.unwrap_or_default();
out.push((id, galaxy, at_ms));
}
drop(cursor);
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
out.sort_by_key(|(_, _, at)| *at);
Ok(out)
}
pub fn count_index_pending(&self) -> Result<usize> {
Ok(self.index_pending_entries()?.len())
}
pub fn clear_index_pending(&self, ids: &[String]) -> Result<usize> {
let Some(db) = self.index_pending_db else {
return Ok(0);
};
if ids.is_empty() {
return Ok(0);
}
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let mut cleared = 0usize;
for id in ids {
let key = id.as_bytes();
match tx.del(db, &key, None) {
Ok(()) => cleared += 1,
Err(lmdb::Error::NotFound) => {}
Err(e) => {
return Err(CoreError::Memory(format!(
"LMDB del failed (index_pending): {e}"
)));
}
}
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(cleared)
}
pub fn clear_all_index_pending(&self) -> Result<usize> {
let Some(db) = self.index_pending_db else {
return Ok(0);
};
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let count = tx
.open_ro_cursor(db)
.map_err(|e| CoreError::Memory(format!("LMDB cursor failed (index_pending): {e}")))?
.iter()
.count();
tx.clear_db(db)
.map_err(|e| CoreError::Memory(format!("LMDB clear failed (index_pending): {e}")))?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(count)
}
pub fn get(&self, galaxy: Galaxy, id: uuid::Uuid) -> Result<Option<Memory>> {
let db = self.galaxy_db(galaxy)?;
let key = id.as_bytes();
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let result = tx.get(db, key);
match result {
Ok(bytes) => {
let memory = self.decode_record_value(galaxy, key, bytes)?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(Some(memory))
}
Err(lmdb::Error::NotFound) => {
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(None)
}
Err(e) => Err(CoreError::Memory(format!("LMDB get failed: {e}"))),
}
}
pub fn find_across_galaxies(&self, id: uuid::Uuid) -> Result<Option<(Galaxy, Memory)>> {
for galaxy in Galaxy::memory_galaxies() {
if let Some(mem) = self.get(galaxy, id)? {
return Ok(Some((galaxy, mem)));
}
}
Ok(None)
}
pub fn delete(&self, galaxy: Galaxy, id: uuid::Uuid) -> Result<bool> {
let db = self.galaxy_db(galaxy)?;
let key = id.as_bytes();
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let exists = tx.get(db, key).is_ok();
if exists {
if let Ok(bytes) = tx.get(db, key) {
if let Ok(memory) = self.decode_record_value(galaxy, key, bytes) {
let _ = self.index_dbs.remove(&mut tx, galaxy, &memory);
}
}
tx.del(db, key, None)
.map_err(|e| CoreError::Memory(format!("LMDB del failed: {e}")))?;
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
if exists {
self.mutation_count.fetch_add(1, Ordering::Relaxed);
}
Ok(exists)
}
pub fn scan(&self, galaxy: Galaxy, limit: usize) -> Result<Vec<Memory>> {
let db = self.galaxy_db(galaxy)?;
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let mut cursor = tx
.open_ro_cursor(db)
.map_err(|e| CoreError::Memory(format!("LMDB cursor failed: {e}")))?;
let mut memories = Vec::with_capacity(limit.min(256));
for (i, (key, val)) in cursor.iter().enumerate() {
if memories.len() >= limit {
break;
}
match self.decode_record_value(galaxy, key, val) {
Ok(memory) => memories.push(memory),
Err(e) => {
tracing::warn!(
"Skipping corrupted entry at index {i} in galaxy {:?}: {e}",
galaxy
);
}
}
}
drop(cursor);
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(memories)
}
pub fn scan_all(&self, galaxy: Galaxy) -> Result<Vec<Memory>> {
self.scan_all_impl(galaxy, false)
}
pub fn scan_all_strict(&self, galaxy: Galaxy) -> Result<Vec<Memory>> {
self.scan_all_impl(galaxy, true)
}
fn scan_all_impl(&self, galaxy: Galaxy, strict: bool) -> Result<Vec<Memory>> {
let db = self.galaxy_db(galaxy)?;
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let mut cursor = tx
.open_ro_cursor(db)
.map_err(|e| CoreError::Memory(format!("LMDB cursor failed: {e}")))?;
let mut memories = Vec::new();
for (i, (key, val)) in cursor.iter().enumerate() {
match self.decode_record_value(galaxy, key, val) {
Ok(memory) => memories.push(memory),
Err(e) => {
if strict {
return Err(CoreError::Memory(format!(
"refusing incomplete scan of {}: record {i} cannot be decoded: {e}",
galaxy.db_name()
)));
}
tracing::warn!(
"Skipping corrupted entry at index {i} in galaxy {:?}: {e}",
galaxy
);
}
}
}
drop(cursor);
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(memories)
}
pub fn count(&self, galaxy: Galaxy) -> Result<usize> {
let db = self.galaxy_db(galaxy)?;
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let mut cursor = tx
.open_ro_cursor(db)
.map_err(|e| CoreError::Memory(format!("LMDB cursor failed: {e}")))?;
let count = cursor.iter().count();
drop(cursor);
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(count)
}
pub fn count_by_tag(&self, galaxy: Galaxy, tag: &str) -> Result<usize> {
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let ids = self.index_dbs.find_by_tag(&tx, galaxy, tag)?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(ids.len())
}
pub fn clear_galaxy(&self, galaxy: Galaxy) -> Result<usize> {
let db = self.galaxy_db(galaxy)?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let mut cursor = tx
.open_ro_cursor(db)
.map_err(|e| CoreError::Memory(format!("LMDB cursor failed: {e}")))?;
let mut count = 0usize;
let keys_to_delete: Vec<(Vec<u8>, Memory)> = cursor
.iter()
.filter_map(|(key, val)| {
if let Ok(memory) = self.decode_record_value(galaxy, key, val) {
Some((key.to_vec(), memory))
} else {
None
}
})
.collect();
drop(cursor);
for (key, memory) in &keys_to_delete {
let _ = self.index_dbs.remove(&mut tx, galaxy, memory);
tx.del(db, &key, None)
.map_err(|e| CoreError::Memory(format!("LMDB del failed: {e}")))?;
count += 1;
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count
.fetch_add(count as u64, Ordering::Relaxed);
Ok(count)
}
pub fn batch_put(&self, galaxy: Galaxy, memories: &[Memory]) -> Result<usize> {
if memories.is_empty() {
return Ok(0);
}
let db = self.galaxy_db(galaxy)?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let mut count = 0usize;
for memory in memories {
let key = memory.metadata.id.as_bytes();
let val = self.encode_record_value(galaxy, memory)?;
match tx.put(db, key, &val, WriteFlags::default()) {
Ok(()) => {}
Err(lmdb::Error::MapFull) => {
tx.abort();
return Err(CoreError::Memory(format!(
"LMDB map full: galaxy {}, consider growing map size",
galaxy.db_name()
)));
}
Err(e) => {
tx.abort();
return Err(CoreError::Memory(format!("LMDB put failed: {e}")));
}
}
self.index_dbs.add(&mut tx, galaxy, memory)?;
count += 1;
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count
.fetch_add(count as u64, Ordering::Relaxed);
Ok(count)
}
pub fn get_raw(&self, galaxy: Galaxy, key: &[u8]) -> Result<Option<Vec<u8>>> {
let db = self.galaxy_db(galaxy)?;
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
match tx.get(db, &key) {
Ok(bytes) => {
let data = bytes.to_vec();
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(Some(data))
}
Err(lmdb::Error::NotFound) => {
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(None)
}
Err(e) => Err(CoreError::Memory(format!("LMDB get_raw failed: {e}"))),
}
}
pub fn put_raw(&self, galaxy: Galaxy, key: &[u8], val: &[u8]) -> Result<()> {
let db = self.galaxy_db(galaxy)?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
tx.put(db, &key, &val, lmdb::WriteFlags::default())
.map_err(|e| CoreError::Memory(format!("LMDB put_raw failed: {e}")))?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count.fetch_add(1, Ordering::Relaxed);
Ok(())
}
pub fn delete_raw(&self, galaxy: Galaxy, key: &[u8]) -> Result<bool> {
let db = self.galaxy_db(galaxy)?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let deleted = tx.del(db, &key, None).is_ok();
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
if deleted {
self.mutation_count.fetch_add(1, Ordering::Relaxed);
}
Ok(deleted)
}
pub fn put_raw_batch(&self, galaxy: Galaxy, entries: &[(&[u8], &[u8])]) -> Result<()> {
self.put_raw_batch_impl(galaxy, entries)?;
self.mutation_count
.fetch_add(entries.len() as u64, Ordering::Relaxed);
Ok(())
}
pub fn put_raw_batch_untracked(
&self,
galaxy: Galaxy,
entries: &[(&[u8], &[u8])],
) -> Result<()> {
self.put_raw_batch_impl(galaxy, entries)
}
fn put_raw_batch_impl(&self, galaxy: Galaxy, entries: &[(&[u8], &[u8])]) -> Result<()> {
let db = self.galaxy_db(galaxy)?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
for (key, val) in entries {
tx.put(db, key, val, WriteFlags::default())
.map_err(|e| CoreError::Memory(format!("LMDB put_raw_batch failed: {e}")))?;
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(())
}
pub fn find_by_content_hash(&self, galaxy: Galaxy, hash: &str) -> Result<Option<uuid::Uuid>> {
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let result = self.index_dbs.find_by_content_hash(&tx, galaxy, hash)?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(result)
}
pub fn find_by_content_hash_scan(
&self,
galaxy: Galaxy,
hash: &str,
) -> Result<Option<uuid::Uuid>> {
let memories = self.scan(galaxy, 10_000)?;
for mem in memories {
if mem.metadata.content_hash == hash {
return Ok(Some(mem.metadata.id));
}
}
Ok(None)
}
pub fn put_dedup(&self, galaxy: Galaxy, memory: &Memory) -> Result<uuid::Uuid> {
if let Some(existing_id) =
self.find_by_content_hash(galaxy, &memory.metadata.content_hash)?
{
return Ok(existing_id);
}
let id = memory.metadata.id;
self.put(galaxy, memory)?;
Ok(id)
}
pub fn put_batch(&self, galaxy: Galaxy, memories: &[Memory]) -> Result<()> {
let db = self.galaxy_db(galaxy)?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
for memory in memories {
let key = memory.metadata.id.as_bytes();
let val = self.encode_record_value(galaxy, memory)?;
tx.put(db, key, &val, WriteFlags::default())
.map_err(|e| CoreError::Memory(format!("LMDB put_batch failed: {e}")))?;
self.index_dbs.add(&mut tx, galaxy, memory)?;
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count
.fetch_add(memories.len() as u64, Ordering::Relaxed);
Ok(())
}
pub fn query(&self, galaxy: Galaxy, query: &MemoryQuery) -> Result<Vec<Memory>> {
if query.content_substring.is_none()
&& query.tags.len() == 1
&& query.min_importance.is_none()
&& query.max_importance.is_none()
&& query.created_after.is_none()
&& query.created_before.is_none()
{
return self.query_by_tag_indexed(galaxy, &query.tags[0], query.limit);
}
if query.content_substring.is_none()
&& query.tags.is_empty()
&& let Some(min) = query.min_importance
&& let Some(max) = query.max_importance
&& query.created_after.is_none()
&& query.created_before.is_none()
{
return self.query_by_importance_indexed(galaxy, min, max, query.limit);
}
if query.content_substring.is_none()
&& query.tags.is_empty()
&& query.min_importance.is_none()
&& query.max_importance.is_none()
&& let Some(after) = query.created_after
&& let Some(before) = query.created_before
{
return self.query_by_time_indexed(galaxy, after, before, query.limit);
}
let memories = self.scan(galaxy, 10_000)?;
let mut results = Vec::new();
for mem in memories {
if query.matches(&mem) {
results.push(mem);
if results.len() >= query.limit {
break;
}
}
}
Ok(results)
}
fn query_by_tag_indexed(&self, galaxy: Galaxy, tag: &str, limit: usize) -> Result<Vec<Memory>> {
let db = self.galaxy_db(galaxy)?;
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let ids = self.index_dbs.find_by_tag(&tx, galaxy, tag)?;
let mut results = Vec::new();
for id in &ids {
if results.len() >= limit {
break;
}
if let Ok(bytes) = tx.get(db, id.as_bytes()) {
if let Ok(mem) = self.decode_record_value(galaxy, id.as_bytes(), bytes) {
results.push(mem);
}
}
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(results)
}
fn query_by_importance_indexed(
&self,
galaxy: Galaxy,
min: f32,
max: f32,
limit: usize,
) -> Result<Vec<Memory>> {
let db = self.galaxy_db(galaxy)?;
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let ids = self
.index_dbs
.find_by_importance_range(&tx, galaxy, min, max)?;
let mut results = Vec::new();
for id in &ids {
if results.len() >= limit {
break;
}
if let Ok(bytes) = tx.get(db, id.as_bytes()) {
if let Ok(mem) = self.decode_record_value(galaxy, id.as_bytes(), bytes) {
results.push(mem);
}
}
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(results)
}
fn query_by_time_indexed(
&self,
galaxy: Galaxy,
after: chrono::DateTime<chrono::Utc>,
before: chrono::DateTime<chrono::Utc>,
limit: usize,
) -> Result<Vec<Memory>> {
let db = self.galaxy_db(galaxy)?;
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let ids = self
.index_dbs
.find_by_time_range(&tx, galaxy, after, before)?;
let mut results = Vec::new();
for id in &ids {
if results.len() >= limit {
break;
}
if let Ok(bytes) = tx.get(db, id.as_bytes()) {
if let Ok(mem) = self.decode_record_value(galaxy, id.as_bytes(), bytes) {
results.push(mem);
}
}
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(results)
}
pub fn put_semantic(&self, galaxy: Galaxy, memory: &mut Memory) -> Result<()> {
let temporal_weight = memory.metadata.coord5d.w;
let importance = memory.metadata.importance;
memory.metadata.coord5d =
self.semantic_encoder
.encode_coordinate(&memory.content, temporal_weight, importance);
self.put(galaxy, memory)
}
pub fn find_similar(
&self,
galaxy: Galaxy,
query_text: &str,
limit: usize,
) -> Result<Vec<(Memory, f32)>> {
let query_coord = self
.semantic_encoder
.encode_coordinate(query_text, 0.5, 0.5);
let memories = self.scan(galaxy, 10_000)?;
let mut results: Vec<(Memory, f32)> = memories
.into_iter()
.map(|m| {
let dist = query_coord.semantic_distance_to(&m.metadata.coord5d);
(m, dist)
})
.collect();
results.sort_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal));
results.truncate(limit);
Ok(results)
}
pub fn put_embedding(&self, memory_id: uuid::Uuid, embedding: &[f32]) -> Result<()> {
let db = self.galaxy_db(Galaxy::Embeddings)?;
let key = memory_id.as_bytes();
let val = encode_embedding(embedding);
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
tx.put(db, key, &val, WriteFlags::default())
.map_err(|e| CoreError::Memory(format!("LMDB put_embedding failed: {e}")))?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count.fetch_add(1, Ordering::Relaxed);
Ok(())
}
pub fn get_embedding(&self, memory_id: uuid::Uuid) -> Result<Option<Vec<f32>>> {
let db = self.galaxy_db(Galaxy::Embeddings)?;
let key = memory_id.as_bytes();
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
match tx.get(db, key) {
Ok(bytes) => {
let embedding = decode_embedding(bytes);
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(Some(embedding))
}
Err(lmdb::Error::NotFound) => {
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(None)
}
Err(e) => Err(CoreError::Memory(format!("LMDB get_embedding failed: {e}"))),
}
}
pub fn delete_embedding(&self, memory_id: uuid::Uuid) -> Result<bool> {
let db = self.galaxy_db(Galaxy::Embeddings)?;
let key = memory_id.as_bytes();
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let exists = tx.get(db, key).is_ok();
if exists {
tx.del(db, key, None)
.map_err(|e| CoreError::Memory(format!("LMDB del_embedding failed: {e}")))?;
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
if exists {
self.mutation_count.fetch_add(1, Ordering::Relaxed);
}
Ok(exists)
}
pub fn put_embedding_cache(&self, cache_key: &str, embedding: &[f32]) -> Result<()> {
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
tx.put(
self.embedding_cache_db,
&cache_key.as_bytes().to_vec(),
&encode_embedding(embedding),
WriteFlags::default(),
)
.map_err(|e| CoreError::Memory(format!("LMDB put_embedding_cache failed: {e}")))?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count.fetch_add(1, Ordering::Relaxed);
Ok(())
}
pub fn put_embedding_cache_batch(&self, entries: &[(String, Vec<f32>)]) -> Result<()> {
if entries.is_empty() {
return Ok(());
}
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
for (key, embedding) in entries {
tx.put(
self.embedding_cache_db,
&key.as_bytes().to_vec(),
&encode_embedding(embedding),
WriteFlags::default(),
)
.map_err(|e| CoreError::Memory(format!("LMDB put_embedding_cache failed: {e}")))?;
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count
.fetch_add(entries.len() as u64, Ordering::Relaxed);
Ok(())
}
pub fn get_embedding_cache(&self, cache_key: &str) -> Result<Option<Vec<f32>>> {
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
match tx.get(self.embedding_cache_db, &cache_key.as_bytes().to_vec()) {
Ok(bytes) => {
let embedding = decode_embedding(bytes);
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(Some(embedding))
}
Err(lmdb::Error::NotFound) => {
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(None)
}
Err(e) => Err(CoreError::Memory(format!(
"LMDB get_embedding_cache failed: {e}"
))),
}
}
pub fn get_embedding_cache_batch(&self, keys: &[String]) -> Result<Vec<Option<Vec<f32>>>> {
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let mut out = Vec::with_capacity(keys.len());
for key in keys {
out.push(
tx.get(self.embedding_cache_db, &key.as_bytes().to_vec())
.ok()
.map(decode_embedding),
);
}
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(out)
}
pub fn embedding_cache_count(&self) -> Result<u64> {
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let mut cursor = tx
.open_ro_cursor(self.embedding_cache_db)
.map_err(|e| CoreError::Memory(format!("LMDB cursor embedding_cache failed: {e}")))?;
let mut count = 0u64;
for _ in cursor.iter() {
count += 1;
}
Ok(count)
}
pub fn record_revision(
&self,
galaxy: Galaxy,
id: MemoryId,
old_hash: &str,
new_hash: &str,
actor: crate::revision::RevisionActor,
) -> Result<crate::revision::MemoryRevision> {
let seq = self.revisions(galaxy, id)?.len() as u32;
let entry = crate::revision::MemoryRevision {
seq,
timestamp: wm_core::time::now_unix_secs(),
old_hash: old_hash.to_string(),
new_hash: new_hash.to_string(),
actor_session: actor.session,
actor_user: actor.user,
actor_compartment: actor.compartment,
};
let key = crate::revision::revision_key(galaxy, id, seq);
let val = serde_json::to_vec(&entry)
.map_err(|e| CoreError::Memory(format!("revision serialize failed: {e}")))?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
tx.put(self.revisions_db, &key, &val, WriteFlags::default())
.map_err(|e| CoreError::Memory(format!("LMDB put revision failed: {e}")))?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count.fetch_add(1, Ordering::Relaxed);
Ok(entry)
}
pub fn revisions(
&self,
galaxy: Galaxy,
id: MemoryId,
) -> Result<Vec<crate::revision::MemoryRevision>> {
const MDB_GET_CURRENT: u32 = 4;
const MDB_NEXT: u32 = 8;
const MDB_SET_RANGE: u32 = 17;
let prefix = crate::revision::revision_prefix(galaxy, id);
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let cursor = tx
.open_ro_cursor(self.revisions_db)
.map_err(|e| CoreError::Memory(format!("LMDB cursor revisions failed: {e}")))?;
let mut out = Vec::new();
if cursor.get(Some(&prefix), None, MDB_SET_RANGE).is_ok() {
while let Ok((key, val)) = cursor.get(None, None, MDB_GET_CURRENT) {
if !key.is_some_and(|k| k.starts_with(&prefix)) {
break;
}
let entry: crate::revision::MemoryRevision = serde_json::from_slice(val)
.map_err(|e| CoreError::Memory(format!("revision deserialize failed: {e}")))?;
out.push(entry);
if out.len() >= 10_000 || cursor.get(None, None, MDB_NEXT).is_err() {
break;
}
}
}
drop(cursor);
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(out)
}
pub fn verify_revision_chain(
&self,
galaxy: Galaxy,
id: MemoryId,
current_hash: &str,
) -> Result<crate::revision::RevisionChainReport> {
let entries = self.revisions(galaxy, id)?;
Ok(crate::revision::verify_chain(&entries, current_hash))
}
pub fn record_attestation(
&self,
galaxy: Galaxy,
id: MemoryId,
entry: &crate::attestation::RecordAttestation,
) -> Result<()> {
let key = crate::attestation::attestation_key(galaxy, id);
let val = serde_json::to_vec(entry)
.map_err(|e| CoreError::Memory(format!("attestation serialize failed: {e}")))?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
tx.put(self.attestations_db, &key, &val, WriteFlags::default())
.map_err(|e| CoreError::Memory(format!("LMDB put attestation failed: {e}")))?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count.fetch_add(1, Ordering::Relaxed);
Ok(())
}
pub fn attestation(
&self,
galaxy: Galaxy,
id: MemoryId,
) -> Result<Option<crate::attestation::RecordAttestation>> {
let key = crate::attestation::attestation_key(galaxy, id);
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let out =
match tx.get(self.attestations_db, &key) {
Ok(val) => Some(serde_json::from_slice(val).map_err(|e| {
CoreError::Memory(format!("attestation deserialize failed: {e}"))
})?),
Err(lmdb::Error::NotFound) => None,
Err(e) => {
return Err(CoreError::Memory(format!(
"LMDB get attestation failed: {e}"
)));
}
};
drop(tx);
Ok(out)
}
pub fn scan_attestations(&self) -> Result<Vec<crate::attestation::RecordAttestation>> {
const MDB_GET_CURRENT: u32 = 4;
const MDB_NEXT: u32 = 8;
const MDB_FIRST: u32 = 9;
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let cursor = tx
.open_ro_cursor(self.attestations_db)
.map_err(|e| CoreError::Memory(format!("LMDB cursor attestations failed: {e}")))?;
let mut out = Vec::new();
if cursor.get(None, None, MDB_FIRST).is_ok() {
while let Ok((_, val)) = cursor.get(None, None, MDB_GET_CURRENT) {
let entry: crate::attestation::RecordAttestation = serde_json::from_slice(val)
.map_err(|e| {
CoreError::Memory(format!("attestation deserialize failed: {e}"))
})?;
out.push(entry);
if out.len() >= 1_000_000 || cursor.get(None, None, MDB_NEXT).is_err() {
break;
}
}
}
drop(cursor);
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
Ok(out)
}
pub fn verify_attestation(
&self,
galaxy: Galaxy,
id: MemoryId,
) -> Result<crate::attestation::AttestationReport> {
use crate::attestation::AttestationReport;
let Some(att) = self.attestation(galaxy, id)? else {
return Ok(AttestationReport {
attested: false,
signature_valid: false,
matches_head: false,
memory_present: self.get(galaxy, id)?.is_some(),
breaks: vec!["no attestation recorded for this memory".to_string()],
});
};
let mut breaks = Vec::new();
let signature_valid = crate::attestation::verify_attestation(&att);
if !signature_valid {
breaks.push("signature does not verify against recorded pubkey".to_string());
}
let (matches_head, memory_present) = if let Some(memory) = self.get(galaxy, id)? {
let matches = memory.metadata.content_hash == att.record_hash;
if !matches {
breaks.push(
"attested record_hash != live content_hash (memory updated after attestation)"
.to_string(),
);
}
(matches, true)
} else {
breaks.push("attested memory id not present in galaxy".to_string());
(false, false)
};
Ok(AttestationReport {
attested: true,
signature_valid,
matches_head,
memory_present,
breaks,
})
}
pub fn attestation_sweep(
&self,
) -> Result<
Vec<(
crate::attestation::RecordAttestation,
crate::attestation::AttestationReport,
)>,
> {
use crate::attestation::AttestationReport;
let mut out = Vec::new();
for att in self.scan_attestations()? {
let parsed = match (
Galaxy::from_db_name(&att.galaxy),
uuid::Uuid::parse_str(&att.memory_id),
) {
(Some(galaxy), Ok(id)) => Some((galaxy, id)),
_ => None,
};
match parsed {
Some((galaxy, id)) => out.push((att, self.verify_attestation(galaxy, id)?)),
None => out.push((
att,
AttestationReport {
attested: true,
signature_valid: false,
matches_head: false,
memory_present: false,
breaks: vec![
"attestation row has unparseable galaxy or memory id".to_string(),
],
},
)),
}
}
Ok(out)
}
pub fn put_cold_record(&self, record: &crate::cold_storage::ColdRecord) -> Result<()> {
let key = record.id.as_bytes();
let val = rmp_serde::to_vec_named(record)
.map_err(|e| CoreError::Memory(format!("Cold record serialization failed: {e}")))?;
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
tx.put(self.cold_storage_db, key, &val, WriteFlags::default())
.map_err(|e| CoreError::Memory(format!("LMDB put cold_storage failed: {e}")))?;
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
self.mutation_count.fetch_add(1, Ordering::Relaxed);
Ok(())
}
pub fn get_cold_record(&self, id: MemoryId) -> Result<Option<crate::cold_storage::ColdRecord>> {
let key = id.as_bytes();
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
match tx.get(self.cold_storage_db, key) {
Ok(bytes) => {
let record: crate::cold_storage::ColdRecord = rmp_serde::from_slice(bytes)
.map_err(|e| {
CoreError::Memory(format!("Cold record deserialization failed: {e}"))
})?;
Ok(Some(record))
}
Err(lmdb::Error::NotFound) => Ok(None),
Err(e) => Err(CoreError::Memory(format!(
"LMDB get cold_storage failed: {e}"
))),
}
}
pub fn delete_cold_record(&self, id: MemoryId) -> Result<bool> {
let key = id.as_bytes();
let mut tx = self
.env
.begin_rw_txn()
.map_err(|e| CoreError::Memory(format!("LMDB rw_txn failed: {e}")))?;
let deleted = match tx.del(self.cold_storage_db, key, None) {
Ok(()) => true,
Err(lmdb::Error::NotFound) => false,
Err(e) => {
return Err(CoreError::Memory(format!(
"LMDB del cold_storage failed: {e}"
)));
}
};
tx.commit()
.map_err(|e| CoreError::Memory(format!("LMDB commit failed: {e}")))?;
if deleted {
self.mutation_count.fetch_add(1, Ordering::Relaxed);
}
Ok(deleted)
}
pub fn count_cold(&self, galaxy: Option<Galaxy>) -> Result<usize> {
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let mut cursor = tx
.open_ro_cursor(self.cold_storage_db)
.map_err(|e| CoreError::Memory(format!("LMDB open_ro_cursor failed: {e}")))?;
let mut count = 0;
for (_key, val) in cursor.iter() {
if let Some(target_g) = galaxy {
let record: crate::cold_storage::ColdRecord =
rmp_serde::from_slice(val).map_err(|e| {
CoreError::Memory(format!("Cold record deserialization failed: {e}"))
})?;
if record.galaxy == target_g {
count += 1;
}
} else {
count += 1;
}
}
Ok(count)
}
pub fn list_cold_records(
&self,
galaxy: Option<Galaxy>,
limit: usize,
) -> Result<Vec<crate::cold_storage::ColdRecordSummary>> {
let query = crate::cold_storage::ColdQuery {
galaxy,
limit: if limit == 0 { 100 } else { limit },
..Default::default()
};
self.query_cold_records(&query)
}
pub fn query_cold_records(
&self,
query: &crate::cold_storage::ColdQuery,
) -> Result<Vec<crate::cold_storage::ColdRecordSummary>> {
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let mut cursor = tx
.open_ro_cursor(self.cold_storage_db)
.map_err(|e| CoreError::Memory(format!("LMDB open_ro_cursor failed: {e}")))?;
let mut results = Vec::new();
let limit = if query.limit == 0 {
usize::MAX
} else {
query.limit
};
for (_key, val) in cursor.iter() {
let record: crate::cold_storage::ColdRecord =
rmp_serde::from_slice(val).map_err(|e| {
CoreError::Memory(format!("Cold record deserialization failed: {e}"))
})?;
let summary = record.summary();
if query.matches(&summary) {
results.push(summary);
if results.len() >= limit {
break;
}
}
}
Ok(results)
}
pub fn find_cold_matching(
&self,
terms: &[String],
galaxy: Option<Galaxy>,
limit: usize,
max_scan: usize,
) -> Result<crate::cold_storage::ColdDiscoveryOutcome> {
self.find_cold_matching_eligible(terms, galaxy, limit, max_scan, |_| true)
}
pub fn find_cold_matching_eligible(
&self,
terms: &[String],
galaxy: Option<Galaxy>,
limit: usize,
max_scan: usize,
eligible: impl Fn(&Memory) -> bool,
) -> Result<crate::cold_storage::ColdDiscoveryOutcome> {
use crate::cold_storage::ColdDiscoveryStop;
let mut out = crate::cold_storage::ColdDiscoveryOutcome::default();
if terms.is_empty() || limit == 0 || max_scan == 0 {
return Ok(out);
}
let tx = self
.env
.begin_ro_txn()
.map_err(|e| CoreError::Memory(format!("LMDB ro_txn failed: {e}")))?;
let mut cursor = tx
.open_ro_cursor(self.cold_storage_db)
.map_err(|e| CoreError::Memory(format!("LMDB open_ro_cursor failed: {e}")))?;
let mut iter = cursor.iter();
loop {
if out.scanned >= max_scan {
out.stop_reason = ColdDiscoveryStop::ScanLimit;
break;
}
if out.records.len() >= limit {
out.stop_reason = ColdDiscoveryStop::ResultLimit;
break;
}
let Some((key, val)) = iter.next() else {
out.stop_reason = ColdDiscoveryStop::Exhausted;
break;
};
out.scanned += 1;
let record: crate::cold_storage::ColdRecord = if let Ok(r) = rmp_serde::from_slice(val)
{
r
} else {
out.integrity_rejected += 1;
continue;
};
if let Some(g) = galaxy {
if record.galaxy != g {
continue;
}
}
out.candidates += 1;
let mem = if let Ok(m) = record.decompress() {
m
} else {
out.integrity_rejected += 1;
continue;
};
if mem.metadata.is_private {
out.private_skipped += 1;
continue;
}
if !mem.metadata.validity.is_current() {
out.non_current_skipped += 1;
continue;
}
let integrity_ok = key == record.id.as_bytes()
&& mem.metadata.id == record.id
&& mem.metadata.galaxy == record.galaxy
&& mem.metadata.content_hash == record.content_hash
&& crate::content_hash(&mem.content) == record.content_hash;
if !integrity_ok {
out.integrity_rejected += 1;
continue;
}
let haystack = format!(
"{} {}",
mem.content.to_lowercase(),
mem.metadata.tags.join(" ").to_lowercase()
);
if !terms.iter().all(|t| haystack.contains(t.as_str())) {
continue;
}
out.matched += 1;
if !eligible(&mem) {
out.eligibility_skipped += 1;
continue;
}
out.records.push(record);
}
Ok(out)
}
pub fn freeze_to_cold(
&self,
search: Option<&crate::SearchEngine>,
memory_id: MemoryId,
distance: f32,
factors: crate::cold_storage::OuterRimFactors,
digest_id: Option<MemoryId>,
notes: Option<String>,
codec: crate::cold_storage::CompressionCodec,
) -> Result<crate::cold_storage::ColdRecord> {
let (galaxy, mut mem) = self.find_across_galaxies(memory_id)?.ok_or_else(|| {
CoreError::NotFound(format!("Memory {memory_id} not found in hot store"))
})?;
if mem.metadata.tier != crate::memory::Tier::Archival {
let _ = mem.transition_tier(crate::memory::Tier::Archival);
}
let record =
crate::cold_storage::ColdRecord::new(&mem, distance, factors, digest_id, notes, codec)?;
self.put_cold_record(&record)?;
self.delete(galaxy, memory_id)?;
if let Some(engine) = search {
if let Ok(mut writer_guard) = engine.writer() {
let _ = engine.delete_document(&mut writer_guard, &memory_id.to_string());
let _ = engine.commit(&mut writer_guard);
}
}
Ok(record)
}
pub fn thaw_from_cold(
&self,
search: Option<&crate::SearchEngine>,
memory_id: MemoryId,
) -> Result<Memory> {
let record = self.get_cold_record(memory_id)?.ok_or_else(|| {
CoreError::NotFound(format!("Memory {memory_id} not found in cold storage"))
})?;
let mut mem = record.decompress()?;
let _ = mem.transition_tier(crate::memory::Tier::Episodic);
mem.metadata.accessed_at = chrono::Utc::now();
mem.metadata.access_count += 1;
mem.metadata.recall_count += 1;
if !mem.metadata.tags.iter().any(|t| t == "thawed:phagic") {
mem.metadata.tags.push("thawed:phagic".to_string());
}
self.put(record.galaxy, &mem)?;
if let Some(engine) = search {
if let Ok(mut writer_guard) = engine.writer() {
let _ = engine.index_memory(&mut writer_guard, &mem);
let _ = engine.commit(&mut writer_guard);
}
}
self.delete_cold_record(memory_id)?;
Ok(mem)
}
pub fn find_anywhere(&self, id: MemoryId) -> Result<Option<(Galaxy, Memory, bool)>> {
if let Some((galaxy, mem)) = self.find_across_galaxies(id)? {
return Ok(Some((galaxy, mem, false)));
}
if let Some(cold_record) = self.get_cold_record(id)? {
let galaxy = cold_record.galaxy;
let mem = cold_record.decompress()?;
return Ok(Some((galaxy, mem, true)));
}
Ok(None)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::content_hash;
#[test]
fn open_and_create_galaxies() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
for galaxy in Galaxy::all() {
let _db = store.galaxy_db(galaxy).unwrap();
}
}
#[test]
fn ensure_schema_completes_a_pre_cold_store() {
use lmdb::{DatabaseFlags as LmdbFlags, Environment as LmdbEnv};
use uuid::Uuid;
let tmp = tempfile::tempdir().unwrap();
let path = tmp.path().join("old-store");
std::fs::create_dir_all(&path).unwrap();
{
let env = LmdbEnv::new().set_max_dbs(64).open(&path).unwrap();
for galaxy in Galaxy::all() {
env.create_db(Some(galaxy.db_name()), LmdbFlags::default())
.unwrap();
}
for (name, flags) in crate::indexes::INDEX_DBS {
env.create_db(Some(name), *flags).unwrap();
}
for (name, flags) in [
("episodic_records", LmdbFlags::default()),
("episodic_terms_v2", LmdbFlags::DUP_SORT),
("embedding_cache", LmdbFlags::default()),
("revisions", LmdbFlags::default()),
(crate::attestation::ATTESTATIONS_DB, LmdbFlags::default()),
] {
env.create_db(Some(name), flags).unwrap();
}
}
let files_before: Vec<String> = {
let mut names: Vec<String> = std::fs::read_dir(&path)
.unwrap()
.map(|e| e.unwrap().file_name().to_string_lossy().to_string())
.collect();
names.sort();
names
};
let data_before = std::fs::read(path.join("data.mdb")).unwrap();
let error = match MemoryStore::open_readonly(&path) {
Ok(_) => panic!("strict read-only open must refuse an incomplete store"),
Err(e) => e.to_string(),
};
assert!(error.contains("cold_storage"), "{error}");
assert_eq!(
std::fs::read(path.join("data.mdb")).unwrap(),
data_before,
"readonly refusal must not mutate the pre-cold store"
);
let files_after: Vec<String> = {
let mut names: Vec<String> = std::fs::read_dir(&path)
.unwrap()
.map(|e| e.unwrap().file_name().to_string_lossy().to_string())
.collect();
names.sort();
names
};
assert_eq!(
files_after, files_before,
"readonly refusal changed the store directory"
);
let created = MemoryStore::ensure_schema(&path).unwrap();
assert_eq!(created, vec!["cold_storage".to_string()], "{created:?}");
let store = MemoryStore::open_readonly(&path).unwrap();
assert!(store.get_cold_record(Uuid::nil()).unwrap().is_none());
assert!(MemoryStore::ensure_schema(&path).unwrap().is_empty());
}
#[test]
fn cold_discovery_hydrates_verifies_and_respects_visibility() {
use crate::cold_storage::{ColdRecord, CompressionCodec, OuterRimFactors};
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let factors = OuterRimFactors {
age_factor: 0.5,
access_factor: 0.5,
resonance_factor: 0.5,
emotional_factor: 0.5,
importance_factor: 0.5,
distance: 0.5,
};
let pub_mem = Memory::new(
Galaxy::Codex,
"public needle zxquniquecoldfact741 buried".into(),
);
let rec = ColdRecord::new(
&pub_mem,
0.5,
factors.clone(),
None,
None,
CompressionCodec::Gzip,
)
.unwrap();
store.put_cold_record(&rec).unwrap();
let out = store
.find_cold_matching(&["zxquniquecoldfact741".to_string()], None, 10, 100)
.unwrap();
assert_eq!(out.matched, 1);
assert_eq!(out.integrity_rejected, 0);
assert_eq!(out.records.len(), 1);
let mut priv_mem = Memory::new(Galaxy::Codex, "private needle zxquniquecoldfact742".into());
priv_mem.metadata.is_private = true;
store
.put_cold_record(
&ColdRecord::new(&priv_mem, 0.5, factors, None, None, CompressionCodec::Gzip)
.unwrap(),
)
.unwrap();
let out_priv = store
.find_cold_matching(&["zxquniquecoldfact742".to_string()], None, 10, 100)
.unwrap();
assert_eq!(out_priv.matched, 0);
assert_eq!(out_priv.private_skipped, 1);
let mut tampered = rec.clone();
let mut bad = Memory::new(Galaxy::Codex, "tampered needle zxquniquecoldfact743".into());
bad.metadata.id = rec.id;
let (payload, size) =
crate::cold_storage::compress_memory(&bad, CompressionCodec::Gzip).unwrap();
tampered.compressed_payload = payload;
tampered.uncompressed_size = size;
store.put_cold_record(&tampered).unwrap();
let out_tamper = store
.find_cold_matching(&["zxquniquecoldfact743".to_string()], None, 10, 100)
.unwrap();
assert_eq!(out_tamper.matched, 0);
assert_eq!(out_tamper.integrity_rejected, 1);
store.delete_cold_record(rec.id).unwrap();
let wrong_key = uuid::Uuid::from_u128(741);
assert_ne!(wrong_key, rec.id);
let value = rmp_serde::to_vec_named(&rec).unwrap();
let mut tx = store.env.begin_rw_txn().unwrap();
tx.put(
store.cold_storage_db,
wrong_key.as_bytes(),
&value,
WriteFlags::default(),
)
.unwrap();
tx.commit().unwrap();
let wrong_key_out = store
.find_cold_matching(&["zxquniquecoldfact741".into()], None, 10, 100)
.unwrap();
assert!(wrong_key_out.records.is_empty());
assert_eq!(wrong_key_out.integrity_rejected, 1);
}
#[test]
fn query_substring_filters_galaxy_wide() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
for (i, content) in [
"the mesh joins at dawn",
"unrelated content entirely",
"MESH joins at dusk",
]
.iter()
.enumerate()
{
let mut m = Memory::new(Galaxy::Codex, content.to_string());
m.metadata.importance = 0.5 + i as f32 / 10.0;
store.put(Galaxy::Codex, &m).unwrap();
}
let hits = store
.query(
Galaxy::Codex,
&MemoryQuery::new().with_content_substring("mesh joins"),
)
.unwrap();
assert_eq!(hits.len(), 2, "CI substring must match both: {hits:?}");
assert!(
hits.iter()
.all(|m| m.content.to_lowercase().contains("mesh joins"))
);
let none = store
.query(
Galaxy::Codex,
&MemoryQuery::new().with_content_substring("quantum calendar"),
)
.unwrap();
assert!(none.is_empty(), "no match must be an honest empty set");
let mut tagged = Memory::new(Galaxy::Codex, "mesh joins again".to_string());
tagged.metadata.tags = vec!["mesh".into()];
store.put(Galaxy::Codex, &tagged).unwrap();
let combined = store
.query(
Galaxy::Codex,
&MemoryQuery::new()
.with_tags(vec!["mesh".into()])
.with_content_substring("again"),
)
.unwrap();
assert_eq!(combined.len(), 1);
assert_eq!(combined[0].content, "mesh joins again");
}
#[cfg(unix)]
#[test]
fn store_dir_has_restrictive_permissions() {
let tmp = tempfile::tempdir().unwrap();
let store_path = tmp.path().join("lmdb");
let _store = MemoryStore::open_default(&store_path).unwrap();
let perms = std::fs::metadata(&store_path).unwrap().permissions().mode();
assert_eq!(
perms & 0o777,
0o700,
"store directory should have 0700 permissions, got {:o}",
perms & 0o777
);
}
#[test]
fn put_get_delete_memory() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mem = Memory::new(Galaxy::Codex, "Hello world".to_string());
let id = mem.metadata.id;
store.put(Galaxy::Codex, &mem).unwrap();
let retrieved = store.get(Galaxy::Codex, id).unwrap();
assert!(retrieved.is_some());
assert_eq!(retrieved.unwrap().content, "Hello world");
let deleted = store.delete(Galaxy::Codex, id).unwrap();
assert!(deleted);
let gone = store.get(Galaxy::Codex, id).unwrap();
assert!(gone.is_none());
}
#[test]
fn scan_memories() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
for i in 0..5 {
let mem = Memory::new(Galaxy::Codex, format!("memory-{i}"));
store.put(Galaxy::Codex, &mem).unwrap();
}
let all = store.scan(Galaxy::Codex, 100).unwrap();
assert_eq!(all.len(), 5);
let limited = store.scan(Galaxy::Codex, 3).unwrap();
assert_eq!(limited.len(), 3);
}
#[test]
fn overwrite_removes_stale_index_entries() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mut mem = Memory::new(Galaxy::Codex, "overwrite target".to_string());
mem.metadata.tags = vec!["alpha".to_string()];
mem.metadata.importance = 0.9;
let id = mem.metadata.id;
store.put(Galaxy::Codex, &mem).unwrap();
let mut updated = Memory::new(Galaxy::Codex, "overwritten content".to_string());
updated.metadata.id = id;
updated.metadata.tags = vec!["beta".to_string()];
updated.metadata.importance = 0.1;
store.put(Galaxy::Codex, &updated).unwrap();
let tx = store.env().begin_ro_txn().unwrap();
let by_alpha = store
.index_dbs()
.find_by_tag(&tx, Galaxy::Codex, "alpha")
.unwrap();
let by_beta = store
.index_dbs()
.find_by_tag(&tx, Galaxy::Codex, "beta")
.unwrap();
assert!(
by_alpha.is_empty(),
"stale tag index entries must be removed on overwrite"
);
assert_eq!(by_beta, vec![id]);
let by_importance = store
.index_dbs()
.find_by_importance_range(&tx, Galaxy::Codex, 0.0, 0.2)
.unwrap();
assert!(
by_importance.contains(&id),
"new importance must be indexed"
);
let by_high = store
.index_dbs()
.find_by_importance_range(&tx, Galaxy::Codex, 0.8, 1.0)
.unwrap();
assert!(
!by_high.contains(&id),
"stale importance index entries must be removed on overwrite"
);
let old_hash = content_hash("overwrite target");
let new_hash = content_hash("overwritten content");
assert_eq!(
store
.index_dbs()
.find_by_content_hash(&tx, Galaxy::Codex, &old_hash)
.unwrap(),
None,
"stale content-hash index entry must be removed"
);
assert_eq!(
store
.index_dbs()
.find_by_content_hash(&tx, Galaxy::Codex, &new_hash)
.unwrap(),
Some(id)
);
}
#[test]
fn count_memories() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
assert_eq!(store.count(Galaxy::Codex).unwrap(), 0);
for i in 0..3 {
let mem = Memory::new(Galaxy::Codex, format!("count-{i}"));
store.put(Galaxy::Codex, &mem).unwrap();
}
assert_eq!(store.count(Galaxy::Codex).unwrap(), 3);
}
#[test]
fn count_by_tag_counts_indexed_records() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
assert_eq!(store.count_by_tag(Galaxy::Sessions, "start").unwrap(), 0);
let mut start = Memory::new(Galaxy::Sessions, "{\"type\":\"session_start\"}".into());
start.metadata.tags = vec!["session".into(), "start".into()];
store.put(Galaxy::Sessions, &start).unwrap();
for i in 0..2 {
let mut turn =
Memory::new(Galaxy::Sessions, format!("{{\"type\":\"turn\",\"i\":{i}}}"));
turn.metadata.tags = vec!["session".into(), "turn".into()];
store.put(Galaxy::Sessions, &turn).unwrap();
}
assert_eq!(store.count(Galaxy::Sessions).unwrap(), 3);
assert_eq!(store.count_by_tag(Galaxy::Sessions, "start").unwrap(), 1);
assert_eq!(store.count_by_tag(Galaxy::Sessions, "turn").unwrap(), 2);
assert_eq!(store.count_by_tag(Galaxy::Sessions, "absent").unwrap(), 0);
}
#[test]
fn get_nonexistent_returns_none() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let result = store.get(Galaxy::Codex, uuid::Uuid::new_v4()).unwrap();
assert!(result.is_none());
}
#[test]
fn raw_put_get() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
store
.put_raw(Galaxy::Substrate, b"config:key", b"value123")
.unwrap();
let val = store.get_raw(Galaxy::Substrate, b"config:key").unwrap();
assert_eq!(val, Some(b"value123".to_vec()));
}
#[test]
fn put_dedup_prevents_duplicates() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mem1 = Memory::new(Galaxy::Codex, "duplicate content".into());
let id1 = store.put_dedup(Galaxy::Codex, &mem1).unwrap();
let mem2 = Memory::new(Galaxy::Codex, "duplicate content".into());
let id2 = store.put_dedup(Galaxy::Codex, &mem2).unwrap();
assert_eq!(id1, id2, "dedup should return same ID for same content");
assert_eq!(store.count(Galaxy::Codex).unwrap(), 1);
}
#[test]
fn put_dedup_allows_different_content() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mem1 = Memory::new(Galaxy::Codex, "content A".into());
store.put_dedup(Galaxy::Codex, &mem1).unwrap();
let mem2 = Memory::new(Galaxy::Codex, "content B".into());
store.put_dedup(Galaxy::Codex, &mem2).unwrap();
assert_eq!(store.count(Galaxy::Codex).unwrap(), 2);
}
#[test]
fn put_batch_atomic_write() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let memories: Vec<Memory> = (0..10)
.map(|i| Memory::new(Galaxy::Codex, format!("batch-{i}")))
.collect();
store.put_batch(Galaxy::Codex, &memories).unwrap();
assert_eq!(store.count(Galaxy::Codex).unwrap(), 10);
}
#[test]
fn query_by_tags() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mem1 = Memory::new(Galaxy::Codex, "tagged memory".into())
.with_tags(vec!["rust".into(), "memory".into()]);
let mem2 =
Memory::new(Galaxy::Codex, "other memory".into()).with_tags(vec!["python".into()]);
store.put(Galaxy::Codex, &mem1).unwrap();
store.put(Galaxy::Codex, &mem2).unwrap();
let query = MemoryQuery::new().with_tags(vec!["rust".into()]);
let results = store.query(Galaxy::Codex, &query).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].content, "tagged memory");
}
#[test]
fn query_by_importance_range() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "low".into()).with_importance(0.1),
)
.unwrap();
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "mid".into()).with_importance(0.5),
)
.unwrap();
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "high".into()).with_importance(0.9),
)
.unwrap();
let query = MemoryQuery::new().with_importance_range(0.4, 0.6);
let results = store.query(Galaxy::Codex, &query).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].content, "mid");
}
#[test]
fn memory_query_one_sided_time_bounds() {
let old = Memory::new(Galaxy::Codex, "old".into());
let mut recent = Memory::new(Galaxy::Codex, "recent".into());
recent.metadata.created_at = old.metadata.created_at + chrono::Duration::days(30);
let cutoff = old.metadata.created_at + chrono::Duration::days(10);
let after = MemoryQuery::new().with_created_after(cutoff);
assert!(!after.matches(&old), "pre-cutoff memory must not match");
assert!(after.matches(&recent), "post-cutoff memory must match");
let before = MemoryQuery::new().with_created_before(cutoff);
assert!(before.matches(&old), "pre-cutoff memory must match");
assert!(
!before.matches(&recent),
"post-cutoff memory must not match"
);
let edge = MemoryQuery::new().with_created_after(cutoff);
let mut at = Memory::new(Galaxy::Codex, "at cutoff".into());
at.metadata.created_at = cutoff;
assert!(edge.matches(&at), "created_at == after bound is inclusive");
}
#[test]
fn embedding_put_get_delete() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let id = uuid::Uuid::new_v4();
let embedding = vec![0.1, 0.2, 0.3, 0.4, 0.5];
store.put_embedding(id, &embedding).unwrap();
let retrieved = store.get_embedding(id).unwrap();
assert!(retrieved.is_some());
let retrieved = retrieved.unwrap();
assert_eq!(retrieved.len(), 5);
assert!((retrieved[0] - 0.1).abs() < f32::EPSILON);
assert!(store.delete_embedding(id).unwrap());
assert!(store.get_embedding(id).unwrap().is_none());
}
#[test]
fn embedding_cache_roundtrip_batch_and_count() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let entries: Vec<(String, Vec<f32>)> = (0..5)
.map(|i| (format!("ns:model:{i:016x}"), vec![i as f32; 8]))
.collect();
store.put_embedding_cache_batch(&entries).unwrap();
assert_eq!(store.embedding_cache_count().unwrap(), 5);
let hit = store
.get_embedding_cache("ns:model:0000000000000003")
.unwrap();
assert_eq!(hit.unwrap(), vec![3.0f32; 8]);
assert!(
store
.get_embedding_cache("ns:model:absent")
.unwrap()
.is_none()
);
let keys: Vec<String> = (0..6).map(|i| format!("ns:model:{i:016x}")).collect();
let batch = store.get_embedding_cache_batch(&keys).unwrap();
assert_eq!(batch.len(), 6);
assert!(batch[0..5].iter().all(Option::is_some));
assert!(batch[5].is_none());
store
.put_embedding_cache("ns:model:0000000000000001", &[9.0; 8])
.unwrap();
assert_eq!(store.embedding_cache_count().unwrap(), 5);
assert_eq!(
store
.get_embedding_cache("ns:model:0000000000000001")
.unwrap()
.unwrap(),
vec![9.0f32; 8]
);
}
#[test]
fn embedding_cache_survives_store_reopen() {
let tmp = tempfile::tempdir().unwrap();
{
let store = MemoryStore::open_default(tmp.path()).unwrap();
store
.put_embedding_cache("onnx:bge-small:abc", &[0.5; 384])
.unwrap();
}
let reopened = MemoryStore::open_default(tmp.path()).unwrap();
let cached = reopened.get_embedding_cache("onnx:bge-small:abc").unwrap();
assert_eq!(cached.unwrap(), vec![0.5f32; 384]);
}
#[test]
fn content_hash_is_sha256() {
let hash1 = content_hash("test content");
let hash2 = content_hash("test content");
let hash3 = content_hash("different content");
assert_eq!(hash1, hash2, "same content should produce same hash");
assert_ne!(
hash1, hash3,
"different content should produce different hash"
);
assert_eq!(hash1.len(), 64, "SHA-256 hex should be 64 chars");
}
#[test]
fn query_by_tag_uses_index() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mem1 = Memory::new(Galaxy::Codex, "tagged".into())
.with_tags(vec!["rust".into(), "memory".into()]);
let mem2 = Memory::new(Galaxy::Codex, "other".into()).with_tags(vec!["python".into()]);
store.put(Galaxy::Codex, &mem1).unwrap();
store.put(Galaxy::Codex, &mem2).unwrap();
let query = MemoryQuery::new().with_tags(vec!["rust".into()]);
let results = store.query(Galaxy::Codex, &query).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].content, "tagged");
}
#[test]
fn query_by_importance_uses_index() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "low".into()).with_importance(0.1),
)
.unwrap();
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "mid".into()).with_importance(0.5),
)
.unwrap();
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "high".into()).with_importance(0.9),
)
.unwrap();
let query = MemoryQuery::new().with_importance_range(0.4, 0.6);
let results = store.query(Galaxy::Codex, &query).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].content, "mid");
}
#[test]
fn query_by_time_uses_index() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let t0 = chrono::Utc::now();
std::thread::sleep(std::time::Duration::from_millis(10));
let mem = Memory::new(Galaxy::Codex, "timed".into());
store.put(Galaxy::Codex, &mem).unwrap();
std::thread::sleep(std::time::Duration::from_millis(10));
let t2 = chrono::Utc::now();
let query = MemoryQuery::new().with_time_range(t0, t2);
let results = store.query(Galaxy::Codex, &query).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].content, "timed");
}
#[test]
fn delete_removes_index_entries() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mem = Memory::new(Galaxy::Codex, "test".into())
.with_tags(vec!["tag1".into()])
.with_importance(0.7);
let id = mem.metadata.id;
let hash = mem.metadata.content_hash.clone();
store.put(Galaxy::Codex, &mem).unwrap();
assert!(
store
.find_by_content_hash(Galaxy::Codex, &hash)
.unwrap()
.is_some()
);
store.delete(Galaxy::Codex, id).unwrap();
assert!(
store
.find_by_content_hash(Galaxy::Codex, &hash)
.unwrap()
.is_none()
);
let query = MemoryQuery::new().with_tags(vec!["tag1".into()]);
let results = store.query(Galaxy::Codex, &query).unwrap();
assert!(results.is_empty());
}
#[test]
fn put_batch_updates_indexes() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let memories: Vec<Memory> = (0..5)
.map(|i| {
Memory::new(Galaxy::Codex, format!("batch-{i}"))
.with_tags(vec![format!("tag{i}")])
.with_importance(i as f32 * 0.2)
})
.collect();
store.put_batch(Galaxy::Codex, &memories).unwrap();
for i in 0..5 {
let query = MemoryQuery::new().with_tags(vec![format!("tag{i}")]);
let results = store.query(Galaxy::Codex, &query).unwrap();
assert_eq!(results.len(), 1, "tag{i} should have 1 result");
}
}
#[test]
fn find_by_content_hash_indexed_matches_scan() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mem = Memory::new(Galaxy::Codex, "dedup test".into());
let id = mem.metadata.id;
let hash = mem.metadata.content_hash.clone();
store.put(Galaxy::Codex, &mem).unwrap();
let indexed = store.find_by_content_hash(Galaxy::Codex, &hash).unwrap();
let scanned = store
.find_by_content_hash_scan(Galaxy::Codex, &hash)
.unwrap();
assert_eq!(indexed, scanned);
assert_eq!(indexed, Some(id));
}
#[test]
fn put_dedup_uses_index() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mem1 = Memory::new(Galaxy::Codex, "duplicate content".into());
let id1 = store.put_dedup(Galaxy::Codex, &mem1).unwrap();
let mem2 = Memory::new(Galaxy::Codex, "duplicate content".into());
let id2 = store.put_dedup(Galaxy::Codex, &mem2).unwrap();
assert_eq!(id1, id2, "dedup should return same ID for same content");
assert_eq!(store.count(Galaxy::Codex).unwrap(), 1);
}
#[test]
fn put_semantic_updates_coord5d() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mut mem = Memory::new(
Galaxy::Codex,
"The algorithm computes data using a systematic method".to_string(),
);
let original_coord = mem.metadata.coord5d.clone();
store.put_semantic(Galaxy::Codex, &mut mem).unwrap();
assert_ne!(
mem.metadata.coord5d.x, original_coord.x,
"semantic encoding should change x"
);
assert_ne!(
mem.metadata.coord5d.y, original_coord.y,
"semantic encoding should change y"
);
let retrieved = store.get(Galaxy::Codex, mem.metadata.id).unwrap().unwrap();
assert_eq!(retrieved.metadata.coord5d.x, mem.metadata.coord5d.x);
}
#[test]
fn put_semantic_preserves_temporal_and_importance() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mut mem = Memory::new(Galaxy::Codex, "test content".into()).with_importance(0.8);
mem.metadata.coord5d.w = 0.6;
store.put_semantic(Galaxy::Codex, &mut mem).unwrap();
assert!((mem.metadata.coord5d.w - 0.6).abs() < f32::EPSILON);
assert!((mem.metadata.coord5d.v - 0.8).abs() < f32::EPSILON);
}
#[test]
fn find_similar_returns_nearest_first() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mut logic_mem = Memory::new(
Galaxy::Codex,
"The algorithm computes data using systematic logic and analysis".to_string(),
);
store.put_semantic(Galaxy::Codex, &mut logic_mem).unwrap();
let mut emotion_mem = Memory::new(
Galaxy::Codex,
"I feel love and joy with deep passion and empathy in my heart".to_string(),
);
store.put_semantic(Galaxy::Codex, &mut emotion_mem).unwrap();
let results = store
.find_similar(Galaxy::Codex, "algorithm data systematic method", 10)
.unwrap();
assert!(!results.is_empty());
assert_eq!(results[0].0.metadata.id, logic_mem.metadata.id);
let results = store
.find_similar(Galaxy::Codex, "love joy passion heart feeling", 10)
.unwrap();
assert!(!results.is_empty());
assert_eq!(results[0].0.metadata.id, emotion_mem.metadata.id);
}
#[test]
fn find_similar_empty_galaxy() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let results = store.find_similar(Galaxy::Codex, "anything", 10).unwrap();
assert!(results.is_empty());
}
#[test]
fn find_similar_respects_limit() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
for i in 0..5 {
let mut mem = Memory::new(Galaxy::Codex, format!("algorithm data method {i}"));
store.put_semantic(Galaxy::Codex, &mut mem).unwrap();
}
let results = store
.find_similar(Galaxy::Codex, "algorithm data", 3)
.unwrap();
assert_eq!(results.len(), 3);
}
#[test]
fn semantic_encoder_accessible() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let scores = store.semantic_encoder().encode("algorithm data logic");
assert!(scores.x < 0.5);
}
#[test]
fn put_raw_batch_writes_atomically() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let entries: &[(&[u8], &[u8])] =
&[(b"key1", b"val1"), (b"key2", b"val2"), (b"key3", b"val3")];
store.put_raw_batch(Galaxy::Karma, entries).unwrap();
assert_eq!(
store.get_raw(Galaxy::Karma, b"key1").unwrap().unwrap(),
b"val1"
);
assert_eq!(
store.get_raw(Galaxy::Karma, b"key2").unwrap().unwrap(),
b"val2"
);
assert_eq!(
store.get_raw(Galaxy::Karma, b"key3").unwrap().unwrap(),
b"val3"
);
}
#[test]
fn put_raw_batch_empty_is_noop() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
store.put_raw_batch(Galaxy::Karma, &[]).unwrap();
assert_eq!(store.count(Galaxy::Karma).unwrap(), 0);
}
#[test]
fn entry_limit_rejects_excess_writes() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path())
.unwrap()
.with_entry_limit(3);
for i in 0..3 {
let mem = Memory::new(Galaxy::Codex, format!("memory {i}"));
store.put(Galaxy::Codex, &mem).unwrap();
}
let mem = Memory::new(Galaxy::Codex, "overflow memory".to_string());
let result = store.put(Galaxy::Codex, &mem);
assert!(result.is_err(), "write beyond limit should be rejected");
let err_msg = result.unwrap_err().to_string();
assert!(
err_msg.contains("entry limit reached"),
"error should mention entry limit: {err_msg}"
);
assert_eq!(store.count(Galaxy::Codex).unwrap(), 3);
}
#[test]
fn entry_limit_per_galaxy_independent() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path())
.unwrap()
.with_entry_limit(2);
for i in 0..2 {
let mem = Memory::new(Galaxy::Codex, format!("codex {i}"));
store.put(Galaxy::Codex, &mem).unwrap();
}
let mem = Memory::new(Galaxy::Research, "science memory".to_string());
let result = store.put(Galaxy::Research, &mem);
assert!(
result.is_ok(),
"different galaxy should not be affected by limit"
);
}
#[test]
fn entry_limit_none_allows_unlimited() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
for i in 0..50 {
let mem = Memory::new(Galaxy::Codex, format!("memory {i}"));
store.put(Galaxy::Codex, &mem).unwrap();
}
assert_eq!(store.count(Galaxy::Codex).unwrap(), 50);
}
#[test]
fn map_full_error_is_graceful() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open(tmp.path(), 512 * 1024).unwrap();
let mut written = 0;
let mut got_map_full = false;
for i in 0..1000 {
let mem = Memory::new(
Galaxy::Codex,
format!("memory content {i} {}", "with padding ".repeat(50)),
);
match store.put(Galaxy::Codex, &mem) {
Ok(()) => written += 1,
Err(e) => {
let msg = e.to_string();
if msg.contains("map full") {
got_map_full = true;
break;
}
break;
}
}
}
assert!(
got_map_full || written < 1000,
"should eventually hit map full or error"
);
assert!(written > 0, "should have written at least some memories");
}
#[test]
fn test_find_across_galaxies() {
let tmp = tempfile::tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let mem = Memory::new(Galaxy::Research, "Cross-galaxy research memo".into());
let id = mem.metadata.id;
store.put(Galaxy::Research, &mem).unwrap();
let found = store.find_across_galaxies(id).unwrap();
assert!(found.is_some());
let (galaxy, retrieved) = found.unwrap();
assert_eq!(galaxy, Galaxy::Research);
assert_eq!(retrieved.metadata.id, id);
assert_eq!(retrieved.content, "Cross-galaxy research memo");
assert!(
store
.find_across_galaxies(uuid::Uuid::new_v4())
.unwrap()
.is_none()
);
}
}