Skip to main content

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 incoming edges of `id` — the mirror of [`Self::relations`], with
179    /// the same liveness rule applied to the far end (here the *source*).
180    ///
181    /// # Errors
182    /// Returns [`MemoryError`] if storage access fails.
183    fn incoming_relations(&self, id: u64) -> Result<Vec<MemoryEdge>, MemoryError>;
184
185    /// Remove the edge with `edge_id`. Returns `true` when it existed —
186    /// idempotent: removing an absent edge is `Ok(false)`, never an error.
187    ///
188    /// # Errors
189    /// Returns [`MemoryError`] if storage access fails.
190    fn unrelate(&self, edge_id: u64) -> Result<bool, MemoryError>;
191
192    /// The total number of live (non-expired) tracked facts, including
193    /// internal entity hubs — used as a corpus-size proxy for idf weighting.
194    fn count(&self) -> usize;
195}
196
197/// The default [`MemoryStore`]: the native, file-backed engine
198/// (`velesdb-core`'s `Database`/`AgentMemory`, requiring the `persistence`
199/// feature). Existing callers of `MemoryService::open` see no change — this
200/// is exactly what they already ran.
201#[cfg(feature = "persistence")]
202pub struct NativeStore {
203    memory: AgentMemory,
204}
205
206#[cfg(feature = "persistence")]
207impl NativeStore {
208    /// Open (or create) a native store at `path`, sized for `dimension`.
209    ///
210    /// # Errors
211    /// Returns [`MemoryError`] if the store cannot be opened.
212    pub fn open<P: AsRef<Path>>(path: P, dimension: usize) -> Result<Self, MemoryError> {
213        let db = Arc::new(Database::open(path)?);
214        let memory = AgentMemory::with_dimension(db, dimension)?;
215        Ok(Self { memory })
216    }
217}
218
219#[cfg(feature = "persistence")]
220impl MemoryStore for NativeStore {
221    fn store(&self, id: u64, content: &str, embedding: &[f32]) -> Result<(), MemoryError> {
222        self.memory
223            .semantic()
224            .store(id, content, embedding)
225            .map_err(MemoryError::from)
226    }
227
228    fn store_with_metadata(
229        &self,
230        id: u64,
231        content: &str,
232        embedding: &[f32],
233        metadata: &Metadata,
234    ) -> Result<(), MemoryError> {
235        self.memory
236            .semantic()
237            .store_with_metadata(id, content, embedding, metadata)
238            .map_err(MemoryError::from)
239    }
240
241    fn store_with_ttl(
242        &self,
243        id: u64,
244        content: &str,
245        embedding: &[f32],
246        ttl_seconds: u64,
247    ) -> Result<(), MemoryError> {
248        self.memory
249            .semantic()
250            .store_with_ttl(id, content, embedding, ttl_seconds)
251            .map_err(MemoryError::from)
252    }
253
254    fn update_metadata(&self, id: u64, metadata: &Metadata) -> Result<(), MemoryError> {
255        self.memory
256            .semantic()
257            .update_metadata(id, metadata)
258            .map_err(MemoryError::from)
259    }
260
261    fn store_with_metadata_and_ttl(
262        &self,
263        id: u64,
264        content: &str,
265        embedding: &[f32],
266        metadata: &Metadata,
267        ttl_seconds: u64,
268    ) -> Result<(), MemoryError> {
269        // Ordre delibere : le fait est ecrit avec sa metadata et SANS
270        // expiration, donc il ne peut pas expirer entre les deux appels.
271        // L'expiration est posee ensuite. C'est l'inverse de la sequence
272        // historique (store_with_ttl puis update_metadata), ou le fait etait
273        // deja vivant et deja en train d'expirer pendant la seconde ecriture.
274        self.memory
275            .semantic()
276            .store_with_metadata(id, content, embedding, metadata)
277            .map_err(MemoryError::from)?;
278        self.memory
279            .semantic()
280            .set_ttl_durable(id, ttl_seconds)
281            .map_err(MemoryError::from)
282    }
283
284    fn get(&self, id: u64) -> Result<Option<(String, Vec<f32>)>, MemoryError> {
285        self.memory.semantic().get(id).map_err(MemoryError::from)
286    }
287
288    fn get_metadata(&self, id: u64) -> Result<Option<Metadata>, MemoryError> {
289        self.memory
290            .semantic()
291            .get_metadata(id)
292            .map_err(MemoryError::from)
293    }
294
295    fn get_metadata_batch(&self, ids: &[u64]) -> Result<Vec<Option<Metadata>>, MemoryError> {
296        self.memory
297            .semantic()
298            .get_metadata_batch(ids)
299            .map_err(MemoryError::from)
300    }
301
302    fn delete(&self, id: u64) -> Result<(), MemoryError> {
303        self.memory.semantic().delete(id).map_err(MemoryError::from)
304    }
305
306    fn query_filtered(
307        &self,
308        embedding: &[f32],
309        k: usize,
310        filter: &Metadata,
311        offset: usize,
312    ) -> Result<Vec<(u64, f32, String)>, MemoryError> {
313        self.memory
314            .semantic()
315            .query_filtered(embedding, k, filter, offset)
316            .map_err(MemoryError::from)
317    }
318
319    fn query_excluding(
320        &self,
321        embedding: &[f32],
322        k: usize,
323        exclude: &Metadata,
324    ) -> Result<Vec<(u64, f32, String)>, MemoryError> {
325        self.memory
326            .semantic()
327            .query_excluding(embedding, k, exclude)
328            .map_err(MemoryError::from)
329    }
330
331    fn query_columnar(
332        &self,
333        embedding: &[f32],
334        k: usize,
335        filters: &[ColumnFilter],
336    ) -> Result<Vec<Recollection>, MemoryError> {
337        let (sql, params) = self.build_fused_query(embedding, k, filters)?;
338        // Field names are validated by `build_fused_query`; ensure each one is
339        // indexed so the planner uses a bitmap prefilter instead of an O(n)
340        // post-filter scan. Idempotent and incrementally maintained thereafter.
341        for filter in filters {
342            self.memory
343                .semantic()
344                .ensure_index(&filter.field)
345                .map_err(MemoryError::from)?;
346        }
347        let results = self
348            .memory
349            .query_semantic(&sql, &params)
350            .map_err(MemoryError::from)?;
351        Ok(results.iter().map(to_recollection).collect())
352    }
353
354    fn relate(&self, from: u64, to: u64, relation: &str) -> Result<u64, MemoryError> {
355        self.memory
356            .semantic()
357            .relate(from, to, relation, None)
358            .map_err(MemoryError::from)
359    }
360
361    fn relations(&self, id: u64) -> Result<Vec<MemoryEdge>, MemoryError> {
362        Ok(to_memory_edges(self.memory.semantic().relations(id)?))
363    }
364
365    fn incoming_relations(&self, id: u64) -> Result<Vec<MemoryEdge>, MemoryError> {
366        Ok(to_memory_edges(
367            self.memory.semantic().incoming_relations(id)?,
368        ))
369    }
370
371    fn unrelate(&self, edge_id: u64) -> Result<bool, MemoryError> {
372        self.memory
373            .semantic()
374            .unrelate(edge_id)
375            .map_err(MemoryError::from)
376    }
377
378    fn count(&self) -> usize {
379        self.memory.semantic().count()
380    }
381}
382
383/// Map core [`GraphEdge`](velesdb_core::collection::graph::GraphEdge)s to the
384/// wire-facing [`MemoryEdge`] shape — shared by both edge directions, so the
385/// two can never disagree on which endpoint or id they report.
386#[cfg(feature = "persistence")]
387fn to_memory_edges(edges: Vec<velesdb_core::collection::graph::GraphEdge>) -> Vec<MemoryEdge> {
388    edges
389        .into_iter()
390        .map(|edge| MemoryEdge {
391            id: edge.id(),
392            from: edge.source(),
393            to: edge.target(),
394            relation: edge.label().to_owned(),
395        })
396        .collect()
397}
398
399#[cfg(feature = "persistence")]
400impl NativeStore {
401    /// Build the `VelesQL` for [`Self::query_columnar`]: a `NEAR` predicate
402    /// plus one bound parameter per filter, against the semantic collection.
403    /// Filter *values* are bound as query parameters (never interpolated);
404    /// filter *field names* are validated to be plain identifiers.
405    fn build_fused_query(
406        &self,
407        embedding: &[f32],
408        k: usize,
409        filters: &[ColumnFilter],
410    ) -> Result<(String, HashMap<String, Value>), MemoryError> {
411        use std::fmt::Write as _;
412        let mut params: HashMap<String, Value> = HashMap::new();
413        params.insert("q".to_string(), json!(embedding));
414        let mut predicate = String::from("vector NEAR $q");
415        for (index, filter) in filters.iter().enumerate() {
416            validate_column_filter(filter)?;
417            let key = format!("p{index}");
418            let _ = write!(
419                predicate,
420                " AND {} {} ${key}",
421                filter.field,
422                filter.op.as_sql()
423            );
424            params.insert(key, filter.value.clone());
425        }
426        let sql = format!(
427            "SELECT * FROM {} WHERE {predicate} LIMIT {k}",
428            self.memory.semantic().collection_name()
429        );
430        Ok((sql, params))
431    }
432}
433
434/// Reserved metadata key `remember`/`remember_with_ttl` auto-stamp with
435/// today's date (a `YYYYMMDD` integer, [`crate::clock::today_ymd`]) whenever
436/// the caller didn't already set it — see
437/// [`crate::service::MemoryService::remember_with_ttl`] for the full
438/// contract. A deliberate, documented **exception** to every other
439/// `_veles_`-namespaced key: [`is_reserved_key`] still names it (so it can
440/// never be confused with an arbitrary caller field), but unlike a true
441/// system key —
442/// - a caller MAY set it explicitly (to date a fact retroactively; never
443///   overwritten once present), and
444/// - it is NOT stripped from caller-facing results, so
445///   [`crate::dated_context::format_dated_context`]'s `date_field` (wired
446///   through `recall_fused`'s `date_field` parameter) can read it back with
447///   zero caller effort.
448///
449/// `pub` (re-exported at the crate root) so every caller of `date_field`
450/// names this one string in exactly one place, not a copy-pasted literal.
451pub const AUTO_DATE_FIELD: &str = "_veles_date";
452
453/// True for metadata keys the memory layer reserves: the engine's `content`
454/// payload, and any `_veles_`-namespaced system key (durable TTL, entity
455/// hubs) — [`AUTO_DATE_FIELD`] EXCEPTED, since (unlike every other reserved
456/// key) it is caller-settable and caller-visible by design. The single
457/// source of the reserved-key contract — the service layer (reject/strip)
458/// and every backend enforce it through this one predicate.
459pub(crate) fn is_reserved_key(key: &str) -> bool {
460    key != AUTO_DATE_FIELD && (key == "content" || key.starts_with("_veles_"))
461}
462
463/// Drop reserved system keys from a raw payload, and collapse an
464/// empty-after-stripping map to `None` — the caller-facing shape every
465/// [`Recollection::metadata`] is built from. `pub` because a [`MemoryStore`]
466/// backend that assembles `Recollection`s itself (`query_columnar`) must
467/// apply the same stripping the service layer applies on every other recall
468/// path, or reserved keys leak to callers on that one path only.
469#[must_use]
470pub fn strip_reserved_keys(payload: Option<Metadata>) -> Option<Metadata> {
471    payload.and_then(|payload| {
472        let metadata: Metadata = payload
473            .into_iter()
474            .filter(|(key, _)| !is_reserved_key(key))
475            .collect();
476        (!metadata.is_empty()).then_some(metadata)
477    })
478}
479
480/// [`strip_reserved_keys`] over a *borrowed* payload: clones only the
481/// surviving non-reserved entries. Use this when the payload isn't already
482/// owned — cloning the whole map first would deep-copy the reserved
483/// `content` value (the full fact text) per hit, only to discard it.
484#[must_use]
485pub fn strip_reserved_keys_ref(payload: Option<&Metadata>) -> Option<Metadata> {
486    payload.and_then(|payload| {
487        let metadata: Metadata = payload
488            .iter()
489            .filter(|(key, _)| !is_reserved_key(key))
490            .map(|(key, value)| (key.clone(), value.clone()))
491            .collect();
492        (!metadata.is_empty()).then_some(metadata)
493    })
494}
495
496/// Map a core search result to a [`Recollection`], lifting the fact text out
497/// of the reserved `content` payload key and surfacing any remaining
498/// caller-supplied metadata (reserved system keys excluded).
499#[cfg(feature = "persistence")]
500fn to_recollection(result: &SearchResult) -> Recollection {
501    let payload = result.point.payload.as_ref().and_then(Value::as_object);
502    let content = payload
503        .and_then(|payload| payload.get("content"))
504        .and_then(Value::as_str)
505        .unwrap_or_default()
506        .to_owned();
507    Recollection {
508        id: result.point.id,
509        score: result.score,
510        content,
511        metadata: strip_reserved_keys_ref(payload),
512    }
513}
514
515/// Validate one `recall_where` column filter: a plain, non-reserved
516/// identifier field name and a scalar (string/number/boolean) value. `pub`
517/// and shared so every [`MemoryStore`] backend enforces the *same* documented
518/// contract — the field-name rule keeps a filter safe to place into query
519/// text (`NativeStore` builds `VelesQL`; values are always bound parameters),
520/// and rejects the reserved system columns (`content`, `_veles_*`) regardless
521/// of backend; the scalar rule turns what would be an opaque engine error
522/// into a clear client-input error.
523///
524/// # Errors
525/// Returns [`MemoryError::InvalidFilter`] when either rule is violated.
526pub fn validate_column_filter(filter: &ColumnFilter) -> Result<(), MemoryError> {
527    let field = &filter.field;
528    let plain = !field.is_empty() && field.chars().all(|c| c.is_ascii_alphanumeric() || c == '_');
529    if !plain || is_reserved_key(field) {
530        return Err(MemoryError::InvalidFilter(field.clone()));
531    }
532    match &filter.value {
533        Value::String(_) | Value::Number(_) | Value::Bool(_) => Ok(()),
534        value => Err(MemoryError::InvalidFilter(format!(
535            "value must be a string, number, or boolean, got {value}"
536        ))),
537    }
538}
539
540#[cfg(all(test, feature = "persistence"))]
541#[path = "storage_tests.rs"]
542mod tests;