Skip to main content

uqa_storage/document_store/
memory.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! In-memory document storage and borrowed projection scans.
8
9use std::collections::BTreeMap;
10use std::sync::Arc;
11
12use uqa_core::{DocId, Value};
13
14use crate::backend::StorageBackendResult;
15
16use super::{Document, DocumentMetadata, DocumentStore, SharedDocumentRow, StoredDocument};
17
18mod projection;
19
20#[derive(Debug, Default)]
21pub struct MemoryDocumentStore {
22    /// Rows and their layout IDs advance together when a shared snapshot is changed.
23    state: Arc<MemoryDocumentState>,
24}
25
26/// Cloning retains independent row maps and layouts; snapshot APIs share them until a write.
27impl Clone for MemoryDocumentStore {
28    fn clone(&self) -> Self {
29        Self {
30            state: Arc::new((*self.state).clone()),
31        }
32    }
33}
34
35#[derive(Debug, Default, Clone)]
36struct MemoryDocumentState {
37    documents: BTreeMap<DocId, MemoryDocumentRow>,
38    layouts: Vec<Vec<String>>,
39}
40
41#[derive(Debug, Clone)]
42pub(super) struct MemoryDocumentRow {
43    layout_id: usize,
44    values: Arc<Vec<Value>>,
45    metadata: DocumentMetadata,
46}
47
48impl MemoryDocumentStore {
49    pub fn new() -> Self {
50        Self::default()
51    }
52
53    fn shared_snapshot(&self) -> Self {
54        Self {
55            state: Arc::clone(&self.state),
56        }
57    }
58
59    pub fn iter(&self) -> impl Iterator<Item = (DocId, Document)> + '_ {
60        self.state
61            .documents
62            .iter()
63            .map(|(doc_id, stored)| (*doc_id, self.materialize_document(stored)))
64    }
65
66    fn materialize_document(&self, stored: &MemoryDocumentRow) -> Document {
67        self.state.layouts[stored.layout_id]
68            .iter()
69            .cloned()
70            .zip(stored.values.iter().cloned())
71            .collect()
72    }
73
74    fn materialize_stored_document(&self, stored: &MemoryDocumentRow) -> StoredDocument {
75        StoredDocument::with_metadata(self.materialize_document(stored), stored.metadata)
76    }
77
78    fn field<'a>(&'a self, stored: &'a MemoryDocumentRow, field: &str) -> Option<&'a Value> {
79        let slot = self.state.layouts[stored.layout_id]
80            .binary_search_by(|stored| stored.as_str().cmp(field))
81            .ok()?;
82        stored.values.get(slot)
83    }
84
85    fn put_stored_inner(&mut self, doc_id: DocId, document: StoredDocument) {
86        let state = Arc::make_mut(&mut self.state);
87        let (document, metadata) = document.into_parts();
88        let (layout_id, values) =
89            if let Some(layout_id) = document_layout_id(&state.layouts, &document) {
90                (layout_id, document.into_values().collect())
91            } else {
92                let (layout, values): (Vec<_>, Vec<_>) = document.into_iter().unzip();
93                let layout_id = state.layouts.len();
94                state.layouts.push(layout);
95                (layout_id, values)
96            };
97        state.documents.insert(
98            doc_id,
99            MemoryDocumentRow {
100                layout_id,
101                values: Arc::new(values),
102                metadata,
103            },
104        );
105    }
106}
107
108fn document_layout_id(layouts: &[Vec<String>], document: &Document) -> Option<usize> {
109    layouts.iter().position(|layout| {
110        layout.len() == document.len()
111            && layout
112                .iter()
113                .map(String::as_str)
114                .eq(document.keys().map(String::as_str))
115    })
116}
117
118impl DocumentStore for MemoryDocumentStore {
119    fn put(&mut self, doc_id: DocId, document: Document) -> StorageBackendResult<()> {
120        let metadata = self
121            .state
122            .documents
123            .get(&doc_id)
124            .map_or_else(DocumentMetadata::default, |stored| stored.metadata);
125        self.put_stored_inner(doc_id, StoredDocument::with_metadata(document, metadata));
126        Ok(())
127    }
128
129    fn get(&self, doc_id: DocId) -> StorageBackendResult<Option<Document>> {
130        Ok(self
131            .state
132            .documents
133            .get(&doc_id)
134            .map(|stored| self.materialize_document(stored)))
135    }
136
137    fn put_stored(&mut self, doc_id: DocId, document: StoredDocument) -> StorageBackendResult<()> {
138        self.put_stored_inner(doc_id, document);
139        Ok(())
140    }
141
142    fn get_stored(&self, doc_id: DocId) -> StorageBackendResult<Option<StoredDocument>> {
143        Ok(self
144            .state
145            .documents
146            .get(&doc_id)
147            .map(|stored| self.materialize_stored_document(stored)))
148    }
149
150    fn get_stored_many(
151        &self,
152        doc_ids: &[DocId],
153    ) -> StorageBackendResult<BTreeMap<DocId, StoredDocument>> {
154        Ok(doc_ids
155            .iter()
156            .filter_map(|doc_id| {
157                self.state
158                    .documents
159                    .get(doc_id)
160                    .map(|stored| (*doc_id, self.materialize_stored_document(stored)))
161            })
162            .collect())
163    }
164
165    fn get_metadata(&self, doc_id: DocId) -> StorageBackendResult<Option<DocumentMetadata>> {
166        Ok(self
167            .state
168            .documents
169            .get(&doc_id)
170            .map(|stored| stored.metadata))
171    }
172
173    fn contains_doc_id(&self, doc_id: DocId) -> StorageBackendResult<bool> {
174        Ok(self.state.documents.contains_key(&doc_id))
175    }
176
177    fn get_field(&self, doc_id: DocId, field: &str) -> StorageBackendResult<Option<Value>> {
178        Ok(self
179            .state
180            .documents
181            .get(&doc_id)
182            .and_then(|stored| self.field(stored, field).cloned()))
183    }
184
185    fn get_fields_multi(
186        &self,
187        doc_ids: &[DocId],
188        fields: &[&str],
189    ) -> StorageBackendResult<BTreeMap<DocId, Vec<Value>>> {
190        Ok(doc_ids
191            .iter()
192            .filter_map(|doc_id| {
193                let stored = self.state.documents.get(doc_id)?;
194                let values = fields
195                    .iter()
196                    .map(|field| self.field(stored, field).cloned().unwrap_or(Value::Null))
197                    .collect();
198                Some((*doc_id, values))
199            })
200            .collect())
201    }
202
203    fn for_each_fields_multi(
204        &self,
205        doc_ids: &[DocId],
206        fields: &[&str],
207        visitor: &mut dyn FnMut(DocId, Vec<Value>) -> bool,
208    ) -> StorageBackendResult<()> {
209        for doc_id in doc_ids {
210            let values = self.state.documents.get(doc_id).map_or_else(
211                || vec![Value::Null; fields.len()],
212                |stored| {
213                    fields
214                        .iter()
215                        .map(|field| self.field(stored, field).cloned().unwrap_or(Value::Null))
216                        .collect()
217                },
218            );
219            if !visitor(*doc_id, values) {
220                break;
221            }
222        }
223        Ok(())
224    }
225
226    fn for_each_fields_multi_ref(
227        &self,
228        doc_ids: &[DocId],
229        fields: &[&str],
230        visitor: &mut dyn FnMut(DocId, &[&Value]) -> bool,
231    ) -> StorageBackendResult<()> {
232        self.visit_fields_multi_ref_with_presence(doc_ids, fields, &mut |doc_id, _, values| {
233            visitor(doc_id, values)
234        });
235        Ok(())
236    }
237
238    fn for_each_fields_multi_ref_with_presence(
239        &self,
240        doc_ids: &[DocId],
241        fields: &[&str],
242        visitor: &mut dyn FnMut(DocId, bool, &[&Value]) -> bool,
243    ) -> StorageBackendResult<()> {
244        self.visit_fields_multi_ref_with_presence(doc_ids, fields, visitor);
245        Ok(())
246    }
247
248    fn get_shared_fields(
249        &self,
250        doc_ids: &[DocId],
251        fields: &[&str],
252    ) -> StorageBackendResult<Option<Vec<Option<SharedDocumentRow>>>> {
253        Ok(Some(self.shared_fields(doc_ids, fields)))
254    }
255
256    fn find_doc_id_by_field(
257        &self,
258        field: &str,
259        value: &Value,
260    ) -> StorageBackendResult<Option<DocId>> {
261        Ok(self.state.documents.iter().find_map(|(doc_id, stored)| {
262            (self.field(stored, field) == Some(value)).then_some(*doc_id)
263        }))
264    }
265
266    fn find_doc_id_by_fields(
267        &self,
268        fields: &[String],
269        values: &[Value],
270    ) -> StorageBackendResult<Option<DocId>> {
271        if fields.is_empty() || fields.len() != values.len() {
272            return Ok(None);
273        }
274        Ok(self.state.documents.iter().find_map(|(doc_id, stored)| {
275            fields
276                .iter()
277                .zip(values.iter())
278                .all(|(field, value)| self.field(stored, field).unwrap_or(&Value::Null) == value)
279                .then_some(*doc_id)
280        }))
281    }
282
283    fn patch_fields(
284        &mut self,
285        doc_id: DocId,
286        updates: &BTreeMap<String, Value>,
287    ) -> StorageBackendResult<bool> {
288        let Some(mut document) = self
289            .state
290            .documents
291            .get(&doc_id)
292            .map(|stored| self.materialize_stored_document(stored))
293        else {
294            return Ok(false);
295        };
296        for (field, value) in updates {
297            if matches!(value, Value::Null) {
298                document.fields_mut().remove(field);
299            } else {
300                document.fields_mut().insert(field.clone(), value.clone());
301            }
302        }
303        self.put_stored(doc_id, document)?;
304        Ok(true)
305    }
306
307    fn delete(&mut self, doc_id: DocId) -> StorageBackendResult<()> {
308        if self.state.documents.contains_key(&doc_id) {
309            Arc::make_mut(&mut self.state).documents.remove(&doc_id);
310        }
311        Ok(())
312    }
313
314    fn clear(&mut self) -> StorageBackendResult<()> {
315        if let Some(state) = Arc::get_mut(&mut self.state) {
316            *state = MemoryDocumentState::default();
317        } else {
318            self.state = Arc::default();
319        }
320        Ok(())
321    }
322
323    fn doc_ids(&self) -> StorageBackendResult<Vec<DocId>> {
324        Ok(self.state.documents.keys().copied().collect())
325    }
326
327    fn next_doc_id(&self, after: Option<DocId>) -> StorageBackendResult<Option<DocId>> {
328        use std::ops::Bound::{Excluded, Unbounded};
329
330        Ok(match after {
331            Some(after) => self
332                .state
333                .documents
334                .range((Excluded(after), Unbounded))
335                .next()
336                .map(|(doc_id, _)| *doc_id),
337            None => self.state.documents.keys().next().copied(),
338        })
339    }
340
341    fn next_doc_ids(&self, after: Option<DocId>, limit: usize) -> StorageBackendResult<Vec<DocId>> {
342        use std::ops::Bound::{Excluded, Unbounded};
343
344        if limit == 0 {
345            return Ok(Vec::new());
346        }
347        Ok(match after {
348            Some(after) => self
349                .state
350                .documents
351                .range((Excluded(after), Unbounded))
352                .take(limit)
353                .map(|(doc_id, _)| *doc_id)
354                .collect(),
355            None => self.state.documents.keys().take(limit).copied().collect(),
356        })
357    }
358
359    fn next_shared_fields(
360        &self,
361        after: Option<DocId>,
362        limit: usize,
363        fields: &[&str],
364    ) -> StorageBackendResult<Option<Vec<(DocId, SharedDocumentRow)>>> {
365        Ok(Some(self.next_shared_rows(after, limit, fields)))
366    }
367
368    fn for_each_next_fields(
369        &self,
370        after: Option<DocId>,
371        limit: usize,
372        fields: &[&str],
373        visitor: &mut dyn FnMut(DocId, &[&Value]) -> bool,
374    ) -> StorageBackendResult<Option<usize>> {
375        Ok(Some(self.visit_next_rows(after, limit, fields, visitor)))
376    }
377
378    fn len(&self) -> StorageBackendResult<usize> {
379        Ok(self.state.documents.len())
380    }
381
382    fn snapshot(&self) -> StorageBackendResult<Arc<dyn DocumentStore>> {
383        Ok(Arc::new(self.shared_snapshot()))
384    }
385
386    fn writable_snapshot(&self) -> StorageBackendResult<Box<dyn DocumentStore>> {
387        Ok(Box::new(self.shared_snapshot()))
388    }
389}
390
391#[cfg(test)]
392mod tests;