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