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