use std::collections::{HashMap, HashSet};
#[cfg(feature = "persistence")]
use std::path::Path;
#[cfg(feature = "persistence")]
use std::sync::Arc;
use serde_json::{Map, Value};
pub type Metadata = Map<String, Value>;
use crate::clock;
use crate::embedder::Embedder;
use crate::error::MemoryError;
use crate::extract::{ExtractedAttribute, ExtractedRelation, Extractor};
use crate::id;
use crate::model::{
ColumnFilter, EntityProfile, EntityRelation, Explanation, Link, MemoryEdge, MemoryNode,
Recollection, RememberedExtraction, UnrelateOutcome,
};
#[cfg(feature = "persistence")]
use crate::mutation::MutationObserver;
#[cfg(feature = "persistence")]
use crate::storage::NativeStore;
use crate::storage::{
is_reserved_key, strip_reserved_keys, ColumnStore, FactStore, GraphStore, RecallStore,
AUTO_DATE_FIELD,
};
#[path = "fused_recall.rs"]
mod fused_recall;
#[path = "service_graph.rs"]
mod graph;
#[cfg(feature = "persistence")]
#[allow(dead_code)]
#[path = "online_migration.rs"]
mod online_migration;
#[cfg(feature = "persistence")]
pub(crate) use online_migration::recover_startup;
#[cfg(feature = "mcp")]
pub(crate) use online_migration::{
JobPhase, JobTarget, LiveGenerationSlot, MigrationStartConfig, MigrationStatus,
OnlineMigrationManager,
};
#[cfg(feature = "persistence")]
#[path = "reinforce.rs"]
mod reinforce;
#[cfg(feature = "context")]
#[path = "context/memory_bridge.rs"]
mod memory_bridge;
use crate::storage::HUB_FIELD;
const HUB_ID_SALT: &str = "\u{0}_veles_entity_hub\u{0}";
const MENTIONS_RELATION: &str = "mentions";
const ABOUT_RELATION: &str = "about";
#[cfg(feature = "persistence")]
pub struct MemoryService<E: Embedder, S: FactStore = NativeStore> {
store: S,
embedder: E,
autograph: Option<crate::extract::DynExtractor>,
autograph_queue: AutographQueue,
generation_gate: parking_lot::RwLock<()>,
}
#[cfg(not(feature = "persistence"))]
pub struct MemoryService<E: Embedder, S: FactStore> {
store: S,
embedder: E,
autograph: Option<crate::extract::DynExtractor>,
autograph_queue: AutographQueue,
}
struct GenerationGuard<'a> {
#[cfg(feature = "persistence")]
_guard: parking_lot::RwLockReadGuard<'a, ()>,
#[cfg(not(feature = "persistence"))]
_lifetime: std::marker::PhantomData<&'a ()>,
}
#[cfg_attr(target_arch = "wasm32", allow(dead_code))]
struct AutographJob {
fact_id: u64,
fact: String,
}
#[derive(Default)]
#[cfg_attr(target_arch = "wasm32", allow(dead_code))]
struct AutographQueue {
tx: parking_lot::Mutex<Option<std::sync::mpsc::SyncSender<AutographJob>>>,
dropped: std::sync::atomic::AtomicU64,
closing: std::sync::atomic::AtomicBool,
}
pub struct AutographWorkerHandle {
close_queue: Option<Box<dyn FnOnce() + Send + Sync>>,
join: Option<std::thread::JoinHandle<()>>,
}
impl Drop for AutographWorkerHandle {
fn drop(&mut self) {
if let Some(close) = self.close_queue.take() {
close();
}
if let Some(join) = self.join.take() {
let _ = join.join();
}
}
}
#[cfg(feature = "persistence")]
impl<E: Embedder> MemoryService<E, NativeStore> {
pub fn open<P: AsRef<Path>>(path: P, embedder: E) -> Result<Self, MemoryError> {
let store = NativeStore::open(path, embedder.dimension())?;
Ok(Self {
store,
embedder,
autograph: None,
autograph_queue: AutographQueue::default(),
generation_gate: parking_lot::RwLock::new(()),
})
}
pub(crate) fn install_mutation_observer(
&self,
observer: Option<Arc<dyn MutationObserver>>,
) -> Result<(), MemoryError> {
let _generation = self.generation_gate.write();
self.store.set_mutation_observer(observer)
}
pub(crate) fn migration_capture_active(&self) -> bool {
self.store.mutation_capture_active()
}
}
impl<E: Embedder, S: FactStore> MemoryService<E, S> {
pub fn with_store(store: S, embedder: E) -> Self {
Self {
store,
embedder,
autograph: None,
autograph_queue: AutographQueue::default(),
#[cfg(feature = "persistence")]
generation_gate: parking_lot::RwLock::new(()),
}
}
#[cfg(feature = "persistence")]
fn enter_generation(&self) -> GenerationGuard<'_> {
GenerationGuard {
_guard: self.generation_gate.read(),
}
}
#[cfg(not(feature = "persistence"))]
fn enter_generation(&self) -> GenerationGuard<'_> {
let _ = self;
GenerationGuard {
_lifetime: std::marker::PhantomData,
}
}
#[must_use]
pub fn with_autograph(mut self, extractor: crate::extract::DynExtractor) -> Self {
self.autograph = Some(extractor);
self
}
pub fn remember(
&self,
fact: &str,
links: &[Link],
metadata: Option<&Metadata>,
) -> Result<u64, MemoryError>
where
S: GraphStore,
{
let _generation = self.enter_generation();
self.remember_inner(fact, links, metadata, None, true)
}
pub fn remember_with_ttl(
&self,
fact: &str,
links: &[Link],
metadata: Option<&Metadata>,
ttl_seconds: Option<u64>,
) -> Result<u64, MemoryError>
where
S: GraphStore,
{
let _generation = self.enter_generation();
self.remember_inner(fact, links, metadata, ttl_seconds, true)
}
fn remember_inner(
&self,
fact: &str,
links: &[Link],
metadata: Option<&Metadata>,
ttl_seconds: Option<u64>,
run_autograph: bool,
) -> Result<u64, MemoryError>
where
S: GraphStore,
{
let fact = fact.trim();
self.validate_write(fact, links, metadata, ttl_seconds)?;
let fact_id = id::stable_id(fact);
reject_self_links(fact_id, links)?;
let existed_before = !links.is_empty() && self.store.get(fact_id)?.is_some();
self.write_fact(fact_id, fact, metadata, ttl_seconds)?;
self.link_or_rollback(fact_id, links, existed_before)?;
self.autograph_if(run_autograph, fact_id, fact);
Ok(fact_id)
}
fn validate_write(
&self,
fact: &str,
links: &[Link],
metadata: Option<&Metadata>,
ttl_seconds: Option<u64>,
) -> Result<(), MemoryError> {
validate_fact(fact)?;
reject_zero_ttl(ttl_seconds)?;
reject_reserved_keys(metadata)?;
reject_oversized_metadata(metadata)?;
self.validate_links(links)
}
fn write_fact(
&self,
fact_id: u64,
fact: &str,
metadata: Option<&Metadata>,
ttl_seconds: Option<u64>,
) -> Result<(), MemoryError> {
let embedding = self.embedder.embed(fact)?;
let stamped = stamp_with_today(metadata);
self.store_fact(fact_id, fact, &embedding, stamped.as_ref(), ttl_seconds)
}
fn validate_links(&self, links: &[Link]) -> Result<(), MemoryError> {
for link in links {
validate_relation(&link.relation)?;
}
self.ensure_link_targets_exist(links)
}
fn link_or_rollback(
&self,
fact_id: u64,
links: &[Link],
existed_before: bool,
) -> Result<(), MemoryError>
where
S: GraphStore,
{
let Err(cause) = self.relate_links(fact_id, links) else {
return Ok(());
};
if existed_before {
return Err(cause);
}
match self.store.delete(fact_id) {
Ok(()) => Err(cause),
Err(rollback) => Err(MemoryError::RollbackFailed {
cause: Box::new(cause),
rollback: Box::new(rollback),
}),
}
}
fn autograph_if(&self, run: bool, fact_id: u64, fact: &str)
where
S: GraphStore,
{
if !run {
return;
}
let guard = self.autograph_queue.tx.lock();
if let Some(tx) = guard.as_ref() {
use std::sync::mpsc::TrySendError;
match tx.try_send(AutographJob {
fact_id,
fact: fact.to_owned(),
}) {
Ok(()) => return,
Err(TrySendError::Full(_)) => {
self.autograph_queue
.dropped
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
#[cfg(feature = "mcp")]
tracing::warn!(
fact_id,
"autograph queue full: enrichment dropped — the fact is \
stored, its graph structure is not; re-remembering \
rebuilds it"
);
return;
}
Err(TrySendError::Disconnected(_)) => {
}
}
}
drop(guard);
self.autograph(fact_id, fact);
}
#[must_use]
pub fn autograph_dropped(&self) -> u64 {
self.autograph_queue
.dropped
.load(std::sync::atomic::Ordering::Relaxed)
}
#[must_use]
pub fn autograph_queue_open(&self) -> bool {
self.autograph_queue.tx.lock().is_some()
}
#[must_use]
pub fn has_autograph(&self) -> bool {
self.autograph.is_some()
}
#[must_use]
pub fn fact_count(&self) -> usize {
let _generation = self.enter_generation();
self.store.count()
}
#[must_use]
pub fn edge_count(&self) -> Option<usize>
where
S: GraphStore,
{
let _generation = self.enter_generation();
self.store.edge_count()
}
pub fn list(
&self,
cursor: Option<u64>,
limit: usize,
filter: Option<&Metadata>,
include_internal: bool,
) -> Result<(Vec<crate::model::ListedMemory>, Option<u64>), MemoryError> {
let _generation = self.enter_generation();
let limit = crate::limits::clamp_recall_limit(limit.max(1));
let (page, next) = self.store.list(cursor, limit)?;
let memories = page
.into_iter()
.filter_map(|fact| audited(fact, filter, include_internal))
.collect();
Ok((memories, next))
}
}
#[cfg(not(target_arch = "wasm32"))]
impl<E, S> MemoryService<E, S>
where
E: Embedder + Send + Sync + 'static,
S: FactStore + Send + Sync + 'static,
{
fn autograph_worker_loop(&self, rx: &std::sync::mpsc::Receiver<AutographJob>)
where
S: GraphStore,
{
let mut skipped_on_close: u64 = 0;
for job in rx {
if self
.autograph_queue
.closing
.load(std::sync::atomic::Ordering::Acquire)
{
skipped_on_close += 1;
continue;
}
let _generation = self.enter_generation();
self.autograph(job.fact_id, &job.fact);
}
if skipped_on_close > 0 {
self.autograph_queue
.dropped
.fetch_add(skipped_on_close, std::sync::atomic::Ordering::Relaxed);
#[cfg(feature = "mcp")]
tracing::warn!(
skipped = skipped_on_close,
"autograph worker closing: queued enrichments skipped — \
the facts are stored, their graph structure is not; \
re-remembering rebuilds it"
);
}
}
pub fn spawn_autograph_worker(
self: &std::sync::Arc<Self>,
capacity: usize,
) -> Result<AutographWorkerHandle, MemoryError>
where
S: GraphStore,
{
let (tx, rx) = std::sync::mpsc::sync_channel::<AutographJob>(capacity);
{
let mut guard = self.autograph_queue.tx.lock();
if guard.is_some() {
return Err(MemoryError::Extract(crate::extract::ExtractError::Backend(
"autograph worker already spawned for this service".to_owned(),
)));
}
*guard = Some(tx);
self.autograph_queue
.closing
.store(false, std::sync::atomic::Ordering::Release);
}
let worker_service = std::sync::Arc::clone(self);
let join = std::thread::Builder::new()
.name("velesdb-autograph".to_owned())
.spawn(move || worker_service.autograph_worker_loop(&rx))
.map_err(|err| {
MemoryError::Extract(crate::extract::ExtractError::Backend(format!(
"spawn autograph worker: {err}"
)))
})?;
let closer_service = std::sync::Arc::clone(self);
Ok(AutographWorkerHandle {
close_queue: Some(Box::new(move || {
closer_service
.autograph_queue
.closing
.store(true, std::sync::atomic::Ordering::Release);
closer_service.autograph_queue.tx.lock().take();
})),
join: Some(join),
})
}
}
impl<E: Embedder, S: FactStore> MemoryService<E, S> {
fn autograph(&self, fact_id: u64, fact: &str)
where
S: GraphStore,
{
let Some(extractor) = self.autograph.as_ref() else {
return;
};
let Ok(mut extraction) = extractor.extract_graph(fact) else {
return;
};
if !self.fact_exists(fact_id) {
return;
}
crate::extract::orient_kinship(fact, &mut extraction.relations);
let mut entity_ids: HashMap<String, u64> = HashMap::new();
let mut edges: HashSet<(u64, u64, String)> = HashSet::new();
for extracted in &extraction.facts {
let _ = self.wire_entities(fact_id, &extracted.entities, &mut entity_ids, &mut edges);
}
if !self.fact_exists(fact_id) {
return;
}
let _ = self.wire_relations(&extraction.relations, &mut entity_ids, &mut edges);
let _ = self.wire_attributes(&extraction.attributes, &mut entity_ids);
}
fn fact_exists(&self, fact_id: u64) -> bool {
matches!(self.store.get(fact_id), Ok(Some(_)))
}
fn relate_links(&self, fact_id: u64, links: &[Link]) -> Result<(), MemoryError>
where
S: GraphStore,
{
for link in links {
self.store.relate(fact_id, link.target, &link.relation)?;
}
Ok(())
}
pub fn remember_extracted<X: Extractor>(
&self,
text: &str,
extractor: &X,
metadata: Option<&Metadata>,
) -> Result<RememberedExtraction, MemoryError>
where
S: GraphStore,
{
let _generation = self.enter_generation();
let extraction = Self::extract_passage(text, extractor)?;
self.store_extraction_inner(&extraction, metadata)
}
pub(crate) fn extract_passage<X: Extractor>(
text: &str,
extractor: &X,
) -> Result<crate::extract::Extraction, MemoryError> {
let text = text.trim();
if text.is_empty() {
return Err(MemoryError::EmptyFact);
}
let mut extraction = extractor.extract_graph(text)?;
crate::extract::orient_kinship(text, &mut extraction.relations);
Ok(extraction)
}
#[cfg(feature = "mcp")]
pub(crate) fn store_extraction(
&self,
extraction: &crate::extract::Extraction,
metadata: Option<&Metadata>,
) -> Result<RememberedExtraction, MemoryError>
where
S: GraphStore,
{
let _generation = self.enter_generation();
self.store_extraction_inner(extraction, metadata)
}
fn store_extraction_inner(
&self,
extraction: &crate::extract::Extraction,
metadata: Option<&Metadata>,
) -> Result<RememberedExtraction, MemoryError>
where
S: GraphStore,
{
let mut entity_ids: HashMap<String, u64> = HashMap::new();
let mut edges: HashSet<(u64, u64, String)> = HashSet::new();
let outcome =
self.store_extracted_facts(&extraction.facts, metadata, &mut entity_ids, &mut edges)?;
self.wire_relations(&extraction.relations, &mut entity_ids, &mut edges)?;
self.wire_attributes(&extraction.attributes, &mut entity_ids)?;
Ok(outcome)
}
pub fn entity_profile(&self, name: &str) -> Result<Option<EntityProfile>, MemoryError>
where
S: GraphStore,
{
let _generation = self.enter_generation();
let key = canonical_entity_name(name);
if key.is_empty() {
return Ok(None);
}
let id = id::stable_id(&format!("{HUB_ID_SALT}{key}"));
if self.store.get(id)?.is_none() {
return Ok(None);
}
let (relations, relations_truncated) = self.outgoing_entity_relations(id)?;
let (relations_in, relations_in_truncated) = self.incoming_entity_relations(id)?;
Ok(Some(EntityProfile {
id,
name: key,
attributes: strip_reserved_keys(self.store.get_metadata(id)?).unwrap_or_default(),
relations,
relations_in,
relations_truncated,
relations_in_truncated,
}))
}
fn outgoing_entity_relations(&self, id: u64) -> Result<(Vec<EntityRelation>, bool), MemoryError>
where
S: GraphStore,
{
let scanned = self
.store
.relations_bounded(id, crate::limits::MAX_ENTITY_SCAN_EDGES)?;
self.resolve_entity_relations(scanned, |edge| edge.to)
}
fn incoming_entity_relations(&self, id: u64) -> Result<(Vec<EntityRelation>, bool), MemoryError>
where
S: GraphStore,
{
let scanned = self
.store
.incoming_relations_bounded(id, crate::limits::MAX_ENTITY_SCAN_EDGES)?;
self.resolve_entity_relations(scanned, |edge| edge.from)
}
fn resolve_entity_relations(
&self,
scanned: crate::model::BoundedMemoryEdges,
far_end: impl Fn(&MemoryEdge) -> u64,
) -> Result<(Vec<EntityRelation>, bool), MemoryError> {
let mut relations = Vec::new();
let mut truncated = scanned.truncated;
for edge in scanned.edges {
if edge.relation == MENTIONS_RELATION || edge.relation == ABOUT_RELATION {
continue;
}
if relations.len() >= crate::limits::MAX_ENTITY_RELATIONS {
truncated = true;
break;
}
let far = far_end(&edge);
let content = self.store.get(far)?.map(|(content, _)| content);
relations.push(EntityRelation {
predicate: edge.relation,
target_id: far,
target: content.unwrap_or_default(),
});
}
Ok((relations, truncated))
}
fn wire_relations(
&self,
relations: &[ExtractedRelation],
entity_ids: &mut HashMap<String, u64>,
edges: &mut HashSet<(u64, u64, String)>,
) -> Result<(), MemoryError>
where
S: GraphStore,
{
for relation in relations {
if validate_relation(&relation.predicate).is_err() {
continue;
}
let subject_id = self.entity_hub(&relation.subject, entity_ids)?;
let object_id = self.entity_hub(&relation.object, entity_ids)?;
if subject_id == object_id {
continue;
}
self.add_edge(subject_id, object_id, &relation.predicate, edges)?;
}
Ok(())
}
fn wire_attributes(
&self,
attributes: &[ExtractedAttribute],
entity_ids: &mut HashMap<String, u64>,
) -> Result<(), MemoryError> {
let mut per_entity: HashMap<String, Metadata> = HashMap::new();
for attribute in attributes {
if is_reserved_key(&attribute.key) {
continue;
}
per_entity
.entry(attribute.entity.clone())
.or_default()
.insert(attribute.key.clone(), attribute.value.clone());
}
for (entity, meta) in per_entity {
if meta.is_empty() {
continue;
}
reject_oversized_metadata(Some(&meta))?;
let hub_id = self.entity_hub(&entity, entity_ids)?;
self.store.update_metadata(hub_id, &meta)?;
}
Ok(())
}
fn store_extracted_facts(
&self,
facts: &[crate::extract::ExtractedFact],
metadata: Option<&Metadata>,
entity_ids: &mut HashMap<String, u64>,
edges: &mut HashSet<(u64, u64, String)>,
) -> Result<RememberedExtraction, MemoryError>
where
S: GraphStore,
{
let mut ids = Vec::with_capacity(facts.len());
let mut skipped_over_cap = 0;
for fact in facts {
let content = fact.text.trim();
if content.is_empty() {
continue;
}
let fact_id = match self.remember_inner(content, &[], metadata, None, false) {
Ok(id) => id,
Err(MemoryError::FactTooLarge { .. }) => {
skipped_over_cap += 1;
continue;
}
Err(error) => return Err(error),
};
ids.push(fact_id);
self.wire_entities(fact_id, &fact.entities, entity_ids, edges)?;
}
Ok(RememberedExtraction {
ids,
skipped_over_cap,
})
}
fn wire_entities(
&self,
fact_id: u64,
entities: &[String],
entity_ids: &mut HashMap<String, u64>,
edges: &mut HashSet<(u64, u64, String)>,
) -> Result<(), MemoryError>
where
S: GraphStore,
{
for entity in entities {
if entity.chars().any(char::is_alphanumeric) {
self.wire_entity(fact_id, entity, entity_ids, edges)?;
}
}
Ok(())
}
fn wire_entity(
&self,
fact_id: u64,
entity: &str,
entity_ids: &mut HashMap<String, u64>,
edges: &mut HashSet<(u64, u64, String)>,
) -> Result<(), MemoryError>
where
S: GraphStore,
{
let entity_id = self.entity_hub(entity, entity_ids)?;
if entity_id == fact_id {
return Ok(());
}
self.add_edge(fact_id, entity_id, ABOUT_RELATION, edges)?;
self.add_edge(entity_id, fact_id, MENTIONS_RELATION, edges)?;
Ok(())
}
fn add_edge(
&self,
from: u64,
to: u64,
label: &str,
edges: &mut HashSet<(u64, u64, String)>,
) -> Result<(), MemoryError>
where
S: GraphStore,
{
if edges.insert((from, to, label.to_string())) {
self.relate_inner(from, to, label)?;
}
Ok(())
}
fn entity_hub(
&self,
entity: &str,
entity_ids: &mut HashMap<String, u64>,
) -> Result<u64, MemoryError> {
let key = entity.trim().to_lowercase();
if let Some(&id) = entity_ids.get(&key) {
return Ok(id);
}
let id = self.remember_hub(&key)?;
entity_ids.insert(key, id);
Ok(id)
}
fn remember_hub(&self, key: &str) -> Result<u64, MemoryError> {
let id = id::stable_id(&format!("{HUB_ID_SALT}{key}"));
if self.store.get(id)?.is_some() {
return Ok(id);
}
let content = format!("Entity: {key}");
let embedding = self.embedder.embed(&content)?;
let mut meta = Map::new();
meta.insert(HUB_FIELD.to_string(), Value::Bool(true));
self.store_fact(id, &content, &embedding, Some(&meta), None)?;
Ok(id)
}
fn ensure_exists(&self, id: u64) -> Result<(), MemoryError> {
if self.store.get(id)?.is_none() {
return Err(MemoryError::UnknownMemory(id));
}
Ok(())
}
fn ensure_link_targets_exist(&self, links: &[Link]) -> Result<(), MemoryError> {
for link in links {
self.ensure_exists(link.target)?;
}
Ok(())
}
fn store_fact(
&self,
id: u64,
fact: &str,
embedding: &[f32],
metadata: Option<&Metadata>,
ttl_seconds: Option<u64>,
) -> Result<(), MemoryError> {
match (metadata, ttl_seconds) {
(Some(meta), Some(ttl)) => {
self.store
.store_with_metadata_and_ttl(id, fact, embedding, meta, ttl)?;
}
(Some(meta), None) => self.store.store_with_metadata(id, fact, embedding, meta)?,
(None, Some(ttl)) => self.store.store_with_ttl(id, fact, embedding, ttl)?,
(None, None) => self.store.store(id, fact, embedding)?,
}
Ok(())
}
pub fn recall(
&self,
query: &str,
k: usize,
filter: Option<&Metadata>,
) -> Result<Vec<Recollection>, MemoryError>
where
S: RecallStore,
{
let _generation = self.enter_generation();
self.recall_inner(query, k, filter)
}
fn recall_inner(
&self,
query: &str,
k: usize,
filter: Option<&Metadata>,
) -> Result<Vec<Recollection>, MemoryError>
where
S: RecallStore,
{
let query = query.trim();
if query.is_empty() {
return Ok(Vec::new());
}
reject_reserved_keys(filter)?;
let embedding = self.embedder.embed(query)?;
let hits = self.search(&embedding, k, filter)?;
let ids: Vec<u64> = hits.iter().map(|(id, _, _)| *id).collect();
let payloads = self.store.get_metadata_batch(&ids)?;
#[cfg(feature = "persistence")]
let (hits, payloads) = Self::rl_rerank(hits, payloads);
Ok(hits
.into_iter()
.zip(payloads)
.map(|((id, score, content), payload)| Recollection {
id,
score,
content,
metadata: strip_reserved_keys(payload),
})
.collect())
}
fn search(
&self,
embedding: &[f32],
k: usize,
filter: Option<&Metadata>,
) -> Result<Vec<(u64, f32, String)>, MemoryError>
where
S: RecallStore,
{
match filter {
Some(meta) if !meta.is_empty() => self.store.query_filtered(embedding, k, meta, 0),
_ => self
.store
.query_excluding(embedding, k, &hub_exclude_filter()),
}
}
pub fn recall_where(
&self,
query: &str,
k: usize,
filters: &[ColumnFilter],
) -> Result<Vec<Recollection>, MemoryError>
where
S: ColumnStore + RecallStore,
{
let _generation = self.enter_generation();
let query = query.trim();
if query.is_empty() || k == 0 {
return Ok(Vec::new());
}
if filters.is_empty() {
return self.recall_inner(query, k, None);
}
let embedding = self.embedder.embed(query)?;
self.store.query_columnar(&embedding, k, filters)
}
}
fn hub_exclude_filter() -> Metadata {
let mut exclude = Map::new();
exclude.insert(HUB_FIELD.to_string(), Value::Bool(true));
exclude
}
fn reject_reserved_keys(metadata: Option<&Metadata>) -> Result<(), MemoryError> {
let Some(meta) = metadata else {
return Ok(());
};
for key in meta.keys() {
if is_reserved_key(key) {
return Err(MemoryError::ReservedKey(key.clone()));
}
}
Ok(())
}
fn reject_oversized_metadata(metadata: Option<&Metadata>) -> Result<(), MemoryError> {
let Some(meta) = metadata else {
return Ok(());
};
let bytes = crate::limits::metadata_bytes(meta);
if bytes > crate::limits::MAX_METADATA_BYTES {
return Err(MemoryError::MetadataTooLarge {
bytes,
max: crate::limits::MAX_METADATA_BYTES,
});
}
Ok(())
}
#[cfg(feature = "context")]
pub(crate) fn positive_ttl(ttl_seconds: Option<u64>) -> Option<u64> {
ttl_seconds.filter(|&seconds| seconds > 0)
}
fn audited(
fact: crate::storage::RawListedFact,
filter: Option<&Metadata>,
include_internal: bool,
) -> Option<crate::model::ListedMemory> {
if !include_internal && crate::storage::is_internal_scaffolding(&fact.payload) {
return None;
}
let matches = filter.is_none_or(|wanted| {
wanted
.iter()
.all(|(key, value)| fact.payload.get(key) == Some(value))
});
if !matches {
return None;
}
let metadata = if include_internal {
(!fact.payload.is_empty()).then_some(fact.payload)
} else {
strip_reserved_keys(Some(fact.payload))
};
Some(crate::model::ListedMemory {
id: fact.id,
content: fact.content,
metadata,
})
}
#[must_use]
pub fn canonical_entity_name(name: &str) -> String {
name.trim().to_lowercase()
}
fn validate_fact(fact: &str) -> Result<(), MemoryError> {
if fact.is_empty() {
return Err(MemoryError::EmptyFact);
}
validate_embeddable(fact)
}
pub(crate) fn validate_embeddable(text: &str) -> Result<(), MemoryError> {
if text.len() > crate::limits::MAX_EMBEDDABLE_TEXT_BYTES {
return Err(MemoryError::FactTooLarge {
bytes: text.len(),
max: crate::limits::MAX_EMBEDDABLE_TEXT_BYTES,
});
}
Ok(())
}
#[cfg(feature = "context")]
pub(crate) fn embeddable_prefix(text: &str) -> &str {
let cap = crate::limits::MAX_EMBEDDABLE_TEXT_BYTES;
if text.len() <= cap {
return text;
}
let mut end = cap;
while end > 0 && !text.is_char_boundary(end) {
end -= 1;
}
&text[..end]
}
fn reject_zero_ttl(ttl_seconds: Option<u64>) -> Result<(), MemoryError> {
if ttl_seconds == Some(0) {
return Err(MemoryError::ZeroTtl);
}
Ok(())
}
fn reject_self_links(fact_id: u64, links: &[Link]) -> Result<(), MemoryError> {
if links.iter().any(|link| link.target == fact_id) {
return Err(MemoryError::SelfRelation(fact_id));
}
Ok(())
}
fn stamp_with_today(metadata: Option<&Metadata>) -> Option<Metadata> {
if metadata.is_some_and(|meta| meta.contains_key(AUTO_DATE_FIELD)) {
return metadata.cloned();
}
let Some(today) = clock::today_ymd() else {
return metadata.cloned();
};
let mut stamped = metadata.cloned().unwrap_or_default();
stamped.insert(AUTO_DATE_FIELD.to_owned(), Value::from(today));
Some(stamped)
}
const MAX_RELATION_BYTES: usize = 512;
fn validate_relation(label: &str) -> Result<(), MemoryError> {
if label.is_empty() {
return Err(MemoryError::InvalidRelation(
"relation label must not be empty".to_owned(),
));
}
if label.len() > MAX_RELATION_BYTES {
return Err(MemoryError::InvalidRelation(format!(
"relation label exceeds maximum of {MAX_RELATION_BYTES} bytes ({} given)",
label.len()
)));
}
if label.chars().any(|c| c.is_ascii_control()) {
return Err(MemoryError::InvalidRelation(
"relation label must not contain ASCII control characters".to_owned(),
));
}
Ok(())
}