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