velesdb_memory/storage.rs
1//! Storage backend abstraction for [`crate::service::MemoryService`].
2//!
3//! The wedge orchestration (remember/recall/relate/forget/why/fusion) is
4//! written once, generic over [`MemoryStore`], so it runs unchanged over any
5//! backend: the native, file-backed [`NativeStore`] (the default — nothing
6//! changes for existing callers), or an in-memory backend such as the one
7//! `velesdb-wasm` provides for the browser (no filesystem, no `persistence`
8//! feature).
9
10#[cfg(feature = "persistence")]
11use std::collections::HashMap;
12#[cfg(feature = "persistence")]
13use std::path::Path;
14#[cfg(feature = "persistence")]
15use std::sync::Arc;
16
17#[cfg(feature = "persistence")]
18use serde_json::json;
19use serde_json::Value;
20#[cfg(feature = "persistence")]
21use velesdb_core::agent::AgentMemory;
22#[cfg(feature = "persistence")]
23use velesdb_core::{Database, SearchResult};
24
25use crate::error::MemoryError;
26use crate::model::{ColumnFilter, MemoryEdge, Recollection};
27use crate::service::Metadata;
28
29/// The storage primitives [`crate::service::MemoryService`] needs: write,
30/// vector search, graph edges, and by-id lookup. A backend that implements
31/// this trait can run the full wedge (`remember`/`recall`/`recall_fused`/
32/// `relate`/`forget`/`why`/`remember_extracted`) with no orchestration code
33/// duplicated.
34pub trait MemoryStore {
35 /// Store a fact with no metadata or expiry.
36 ///
37 /// # Errors
38 /// Returns [`MemoryError`] if persistence fails.
39 fn store(&self, id: u64, content: &str, embedding: &[f32]) -> Result<(), MemoryError>;
40
41 /// Store a fact tagged with `metadata`, no expiry.
42 ///
43 /// # Errors
44 /// Returns [`MemoryError`] if persistence fails.
45 fn store_with_metadata(
46 &self,
47 id: u64,
48 content: &str,
49 embedding: &[f32],
50 metadata: &Metadata,
51 ) -> Result<(), MemoryError>;
52
53 /// Store a fact that expires after `ttl_seconds`, no metadata.
54 ///
55 /// # Errors
56 /// Returns [`MemoryError`] if persistence fails.
57 fn store_with_ttl(
58 &self,
59 id: u64,
60 content: &str,
61 embedding: &[f32],
62 ttl_seconds: u64,
63 ) -> Result<(), MemoryError>;
64
65 /// Store a fact with BOTH metadata and a durable TTL, in ONE write.
66 ///
67 /// Default: the historical two-call sequence, so a backend written before
68 /// this method keeps compiling and behaving as it did. Backends that can
69 /// write both at once should override it — the two-call form leaves the
70 /// fact live and expiring between the calls, so a short TTL can lapse in
71 /// the gap and the metadata write then fails on a fact that was perfectly
72 /// valid when the caller asked for it.
73 ///
74 /// # Errors
75 /// Returns [`MemoryError`] if persistence fails.
76 fn store_with_metadata_and_ttl(
77 &self,
78 id: u64,
79 content: &str,
80 embedding: &[f32],
81 metadata: &Metadata,
82 ttl_seconds: u64,
83 ) -> Result<(), MemoryError> {
84 self.store_with_ttl(id, content, embedding, ttl_seconds)?;
85 self.update_metadata(id, metadata)
86 }
87
88 /// Merge `metadata` into an already-stored fact's payload, preserving any
89 /// durable TTL. Used to combine metadata with an expiry (store both in
90 /// two calls rather than needing every metadata×TTL combination as a
91 /// separate primitive).
92 ///
93 /// # Errors
94 /// Returns [`MemoryError`] if `id` is unknown or persistence fails.
95 fn update_metadata(&self, id: u64, metadata: &Metadata) -> Result<(), MemoryError>;
96
97 /// A fact's content and embedding, or `None` if unknown/expired.
98 ///
99 /// # Errors
100 /// Returns [`MemoryError`] if storage access fails.
101 fn get(&self, id: u64) -> Result<Option<(String, Vec<f32>)>, MemoryError>;
102
103 /// A fact's raw stored payload — reserved system keys (`_veles_*`)
104 /// included, so the service layer can check the hub flag before
105 /// stripping them for the caller — or `None` when the fact is
106 /// unknown/expired.
107 ///
108 /// # Errors
109 /// Returns [`MemoryError`] if storage access fails.
110 fn get_metadata(&self, id: u64) -> Result<Option<Metadata>, MemoryError>;
111
112 /// Batched [`Self::get_metadata`]: one storage round trip for every id
113 /// in `ids`, results in the same order and length (an unknown or expired
114 /// id maps to `None`). Same raw-payload semantics as the single-id form.
115 ///
116 /// # Errors
117 /// Returns [`MemoryError`] if storage access fails.
118 fn get_metadata_batch(&self, ids: &[u64]) -> Result<Vec<Option<Metadata>>, MemoryError>;
119
120 /// Delete a fact.
121 ///
122 /// # Errors
123 /// Returns [`MemoryError`] if deletion fails.
124 fn delete(&self, id: u64) -> Result<(), MemoryError>;
125
126 /// Vector search for up to `k` ids, narrowed to facts whose metadata
127 /// exactly matches every key in `filter`.
128 ///
129 /// # Errors
130 /// Returns [`MemoryError`] if the query fails.
131 fn query_filtered(
132 &self,
133 embedding: &[f32],
134 k: usize,
135 filter: &Metadata,
136 offset: usize,
137 ) -> Result<Vec<(u64, f32, String)>, MemoryError>;
138
139 /// Vector search for up to `k` ids, dropping facts whose metadata matches
140 /// every key in `exclude`.
141 ///
142 /// # Errors
143 /// Returns [`MemoryError`] if the query fails.
144 fn query_excluding(
145 &self,
146 embedding: &[f32],
147 k: usize,
148 exclude: &Metadata,
149 ) -> Result<Vec<(u64, f32, String)>, MemoryError>;
150
151 /// Vector search fused with structured `ColumnStore` predicates (ranges
152 /// and comparisons, not just equality) — the engine behind
153 /// [`crate::service::MemoryService::recall_where`].
154 ///
155 /// # Errors
156 /// Returns [`MemoryError::InvalidFilter`] if a filter field is not a
157 /// plain identifier or a filter value is non-scalar, or [`MemoryError`]
158 /// if the query fails.
159 fn query_columnar(
160 &self,
161 embedding: &[f32],
162 k: usize,
163 filters: &[ColumnFilter],
164 ) -> Result<Vec<Recollection>, MemoryError>;
165
166 /// Create a typed edge `from -> to`. Returns the edge id.
167 ///
168 /// # Errors
169 /// Returns [`MemoryError`] if either endpoint is missing or persistence fails.
170 fn relate(&self, from: u64, to: u64, relation: &str) -> Result<u64, MemoryError>;
171
172 /// The outgoing edges of `id`.
173 ///
174 /// # Errors
175 /// Returns [`MemoryError`] if storage access fails.
176 fn relations(&self, id: u64) -> Result<Vec<MemoryEdge>, MemoryError>;
177
178 /// The total number of live (non-expired) tracked facts, including
179 /// internal entity hubs — used as a corpus-size proxy for idf weighting.
180 fn count(&self) -> usize;
181}
182
183/// The default [`MemoryStore`]: the native, file-backed engine
184/// (`velesdb-core`'s `Database`/`AgentMemory`, requiring the `persistence`
185/// feature). Existing callers of `MemoryService::open` see no change — this
186/// is exactly what they already ran.
187#[cfg(feature = "persistence")]
188pub struct NativeStore {
189 memory: AgentMemory,
190}
191
192#[cfg(feature = "persistence")]
193impl NativeStore {
194 /// Open (or create) a native store at `path`, sized for `dimension`.
195 ///
196 /// # Errors
197 /// Returns [`MemoryError`] if the store cannot be opened.
198 pub fn open<P: AsRef<Path>>(path: P, dimension: usize) -> Result<Self, MemoryError> {
199 let db = Arc::new(Database::open(path)?);
200 let memory = AgentMemory::with_dimension(db, dimension)?;
201 Ok(Self { memory })
202 }
203}
204
205#[cfg(feature = "persistence")]
206impl MemoryStore for NativeStore {
207 fn store(&self, id: u64, content: &str, embedding: &[f32]) -> Result<(), MemoryError> {
208 self.memory
209 .semantic()
210 .store(id, content, embedding)
211 .map_err(MemoryError::from)
212 }
213
214 fn store_with_metadata(
215 &self,
216 id: u64,
217 content: &str,
218 embedding: &[f32],
219 metadata: &Metadata,
220 ) -> Result<(), MemoryError> {
221 self.memory
222 .semantic()
223 .store_with_metadata(id, content, embedding, metadata)
224 .map_err(MemoryError::from)
225 }
226
227 fn store_with_ttl(
228 &self,
229 id: u64,
230 content: &str,
231 embedding: &[f32],
232 ttl_seconds: u64,
233 ) -> Result<(), MemoryError> {
234 self.memory
235 .semantic()
236 .store_with_ttl(id, content, embedding, ttl_seconds)
237 .map_err(MemoryError::from)
238 }
239
240 fn update_metadata(&self, id: u64, metadata: &Metadata) -> Result<(), MemoryError> {
241 self.memory
242 .semantic()
243 .update_metadata(id, metadata)
244 .map_err(MemoryError::from)
245 }
246
247 fn store_with_metadata_and_ttl(
248 &self,
249 id: u64,
250 content: &str,
251 embedding: &[f32],
252 metadata: &Metadata,
253 ttl_seconds: u64,
254 ) -> Result<(), MemoryError> {
255 // Ordre delibere : le fait est ecrit avec sa metadata et SANS
256 // expiration, donc il ne peut pas expirer entre les deux appels.
257 // L'expiration est posee ensuite. C'est l'inverse de la sequence
258 // historique (store_with_ttl puis update_metadata), ou le fait etait
259 // deja vivant et deja en train d'expirer pendant la seconde ecriture.
260 self.memory
261 .semantic()
262 .store_with_metadata(id, content, embedding, metadata)
263 .map_err(MemoryError::from)?;
264 self.memory
265 .semantic()
266 .set_ttl_durable(id, ttl_seconds)
267 .map_err(MemoryError::from)
268 }
269
270 fn get(&self, id: u64) -> Result<Option<(String, Vec<f32>)>, MemoryError> {
271 self.memory.semantic().get(id).map_err(MemoryError::from)
272 }
273
274 fn get_metadata(&self, id: u64) -> Result<Option<Metadata>, MemoryError> {
275 self.memory
276 .semantic()
277 .get_metadata(id)
278 .map_err(MemoryError::from)
279 }
280
281 fn get_metadata_batch(&self, ids: &[u64]) -> Result<Vec<Option<Metadata>>, MemoryError> {
282 self.memory
283 .semantic()
284 .get_metadata_batch(ids)
285 .map_err(MemoryError::from)
286 }
287
288 fn delete(&self, id: u64) -> Result<(), MemoryError> {
289 self.memory.semantic().delete(id).map_err(MemoryError::from)
290 }
291
292 fn query_filtered(
293 &self,
294 embedding: &[f32],
295 k: usize,
296 filter: &Metadata,
297 offset: usize,
298 ) -> Result<Vec<(u64, f32, String)>, MemoryError> {
299 self.memory
300 .semantic()
301 .query_filtered(embedding, k, filter, offset)
302 .map_err(MemoryError::from)
303 }
304
305 fn query_excluding(
306 &self,
307 embedding: &[f32],
308 k: usize,
309 exclude: &Metadata,
310 ) -> Result<Vec<(u64, f32, String)>, MemoryError> {
311 self.memory
312 .semantic()
313 .query_excluding(embedding, k, exclude)
314 .map_err(MemoryError::from)
315 }
316
317 fn query_columnar(
318 &self,
319 embedding: &[f32],
320 k: usize,
321 filters: &[ColumnFilter],
322 ) -> Result<Vec<Recollection>, MemoryError> {
323 let (sql, params) = self.build_fused_query(embedding, k, filters)?;
324 // Field names are validated by `build_fused_query`; ensure each one is
325 // indexed so the planner uses a bitmap prefilter instead of an O(n)
326 // post-filter scan. Idempotent and incrementally maintained thereafter.
327 for filter in filters {
328 self.memory
329 .semantic()
330 .ensure_index(&filter.field)
331 .map_err(MemoryError::from)?;
332 }
333 let results = self
334 .memory
335 .query_semantic(&sql, ¶ms)
336 .map_err(MemoryError::from)?;
337 Ok(results.iter().map(to_recollection).collect())
338 }
339
340 fn relate(&self, from: u64, to: u64, relation: &str) -> Result<u64, MemoryError> {
341 self.memory
342 .semantic()
343 .relate(from, to, relation, None)
344 .map_err(MemoryError::from)
345 }
346
347 fn relations(&self, id: u64) -> Result<Vec<MemoryEdge>, MemoryError> {
348 Ok(self
349 .memory
350 .semantic()
351 .relations(id)?
352 .into_iter()
353 .map(|edge| MemoryEdge {
354 from: edge.source(),
355 to: edge.target(),
356 relation: edge.label().to_owned(),
357 })
358 .collect())
359 }
360
361 fn count(&self) -> usize {
362 self.memory.semantic().count()
363 }
364}
365
366#[cfg(feature = "persistence")]
367impl NativeStore {
368 /// Build the `VelesQL` for [`Self::query_columnar`]: a `NEAR` predicate
369 /// plus one bound parameter per filter, against the semantic collection.
370 /// Filter *values* are bound as query parameters (never interpolated);
371 /// filter *field names* are validated to be plain identifiers.
372 fn build_fused_query(
373 &self,
374 embedding: &[f32],
375 k: usize,
376 filters: &[ColumnFilter],
377 ) -> Result<(String, HashMap<String, Value>), MemoryError> {
378 use std::fmt::Write as _;
379 let mut params: HashMap<String, Value> = HashMap::new();
380 params.insert("q".to_string(), json!(embedding));
381 let mut predicate = String::from("vector NEAR $q");
382 for (index, filter) in filters.iter().enumerate() {
383 validate_column_filter(filter)?;
384 let key = format!("p{index}");
385 let _ = write!(
386 predicate,
387 " AND {} {} ${key}",
388 filter.field,
389 filter.op.as_sql()
390 );
391 params.insert(key, filter.value.clone());
392 }
393 let sql = format!(
394 "SELECT * FROM {} WHERE {predicate} LIMIT {k}",
395 self.memory.semantic().collection_name()
396 );
397 Ok((sql, params))
398 }
399}
400
401/// Reserved metadata key `remember`/`remember_with_ttl` auto-stamp with
402/// today's date (a `YYYYMMDD` integer, [`crate::clock::today_ymd`]) whenever
403/// the caller didn't already set it — see
404/// [`crate::service::MemoryService::remember_with_ttl`] for the full
405/// contract. A deliberate, documented **exception** to every other
406/// `_veles_`-namespaced key: [`is_reserved_key`] still names it (so it can
407/// never be confused with an arbitrary caller field), but unlike a true
408/// system key —
409/// - a caller MAY set it explicitly (to date a fact retroactively; never
410/// overwritten once present), and
411/// - it is NOT stripped from caller-facing results, so
412/// [`crate::dated_context::format_dated_context`]'s `date_field` (wired
413/// through `recall_fused`'s `date_field` parameter) can read it back with
414/// zero caller effort.
415///
416/// `pub` (re-exported at the crate root) so every caller of `date_field`
417/// names this one string in exactly one place, not a copy-pasted literal.
418pub const AUTO_DATE_FIELD: &str = "_veles_date";
419
420/// True for metadata keys the memory layer reserves: the engine's `content`
421/// payload, and any `_veles_`-namespaced system key (durable TTL, entity
422/// hubs) — [`AUTO_DATE_FIELD`] EXCEPTED, since (unlike every other reserved
423/// key) it is caller-settable and caller-visible by design. The single
424/// source of the reserved-key contract — the service layer (reject/strip)
425/// and every backend enforce it through this one predicate.
426pub(crate) fn is_reserved_key(key: &str) -> bool {
427 key != AUTO_DATE_FIELD && (key == "content" || key.starts_with("_veles_"))
428}
429
430/// Drop reserved system keys from a raw payload, and collapse an
431/// empty-after-stripping map to `None` — the caller-facing shape every
432/// [`Recollection::metadata`] is built from. `pub` because a [`MemoryStore`]
433/// backend that assembles `Recollection`s itself (`query_columnar`) must
434/// apply the same stripping the service layer applies on every other recall
435/// path, or reserved keys leak to callers on that one path only.
436#[must_use]
437pub fn strip_reserved_keys(payload: Option<Metadata>) -> Option<Metadata> {
438 payload.and_then(|payload| {
439 let metadata: Metadata = payload
440 .into_iter()
441 .filter(|(key, _)| !is_reserved_key(key))
442 .collect();
443 (!metadata.is_empty()).then_some(metadata)
444 })
445}
446
447/// [`strip_reserved_keys`] over a *borrowed* payload: clones only the
448/// surviving non-reserved entries. Use this when the payload isn't already
449/// owned — cloning the whole map first would deep-copy the reserved
450/// `content` value (the full fact text) per hit, only to discard it.
451#[must_use]
452pub fn strip_reserved_keys_ref(payload: Option<&Metadata>) -> Option<Metadata> {
453 payload.and_then(|payload| {
454 let metadata: Metadata = payload
455 .iter()
456 .filter(|(key, _)| !is_reserved_key(key))
457 .map(|(key, value)| (key.clone(), value.clone()))
458 .collect();
459 (!metadata.is_empty()).then_some(metadata)
460 })
461}
462
463/// Map a core search result to a [`Recollection`], lifting the fact text out
464/// of the reserved `content` payload key and surfacing any remaining
465/// caller-supplied metadata (reserved system keys excluded).
466#[cfg(feature = "persistence")]
467fn to_recollection(result: &SearchResult) -> Recollection {
468 let payload = result.point.payload.as_ref().and_then(Value::as_object);
469 let content = payload
470 .and_then(|payload| payload.get("content"))
471 .and_then(Value::as_str)
472 .unwrap_or_default()
473 .to_owned();
474 Recollection {
475 id: result.point.id,
476 score: result.score,
477 content,
478 metadata: strip_reserved_keys_ref(payload),
479 }
480}
481
482/// Validate one `recall_where` column filter: a plain, non-reserved
483/// identifier field name and a scalar (string/number/boolean) value. `pub`
484/// and shared so every [`MemoryStore`] backend enforces the *same* documented
485/// contract — the field-name rule keeps a filter safe to place into query
486/// text (`NativeStore` builds `VelesQL`; values are always bound parameters),
487/// and rejects the reserved system columns (`content`, `_veles_*`) regardless
488/// of backend; the scalar rule turns what would be an opaque engine error
489/// into a clear client-input error.
490///
491/// # Errors
492/// Returns [`MemoryError::InvalidFilter`] when either rule is violated.
493pub fn validate_column_filter(filter: &ColumnFilter) -> Result<(), MemoryError> {
494 let field = &filter.field;
495 let plain = !field.is_empty() && field.chars().all(|c| c.is_ascii_alphanumeric() || c == '_');
496 if !plain || is_reserved_key(field) {
497 return Err(MemoryError::InvalidFilter(field.clone()));
498 }
499 match &filter.value {
500 Value::String(_) | Value::Number(_) | Value::Bool(_) => Ok(()),
501 value => Err(MemoryError::InvalidFilter(format!(
502 "value must be a string, number, or boolean, got {value}"
503 ))),
504 }
505}
506
507#[cfg(all(test, feature = "persistence"))]
508#[path = "storage_tests.rs"]
509mod tests;