Skip to main content

uqa_storage/key_value/
document_store.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Document-store adapter over an ordered key/value store.
8
9use super::codec::{
10    decode_document_value, document_key, document_key_prefix, encode_document_value, read_str,
11    read_u64,
12};
13use super::{Arc, DocId, Document, DocumentStore, KeyValueStore, StorageBackendResult, Value};
14
15/// Document store implemented over [`KeyValueStore`].
16#[derive(Clone)]
17pub struct KeyValueDocumentStore {
18    store: Arc<dyn KeyValueStore>,
19    table: String,
20}
21
22impl KeyValueDocumentStore {
23    pub fn new(store: Arc<dyn KeyValueStore>, table: impl Into<String>) -> Self {
24        Self {
25            store,
26            table: table.into(),
27        }
28    }
29}
30
31impl DocumentStore for KeyValueDocumentStore {
32    fn put(&mut self, doc_id: DocId, document: Document) -> StorageBackendResult<()> {
33        let document: Document = document
34            .into_iter()
35            .filter(|(_, value)| !matches!(value, Value::Null))
36            .collect();
37        let value = encode_document_value(&document)?;
38        self.store.put(&document_key(&self.table, doc_id)?, &value)
39    }
40
41    fn get(&self, doc_id: DocId) -> StorageBackendResult<Option<Document>> {
42        self.store
43            .get(&document_key(&self.table, doc_id)?)?
44            .map(|bytes| decode_document_value(&bytes))
45            .transpose()
46    }
47
48    fn contains_doc_id(&self, doc_id: DocId) -> StorageBackendResult<bool> {
49        self.store.contains_key(&document_key(&self.table, doc_id)?)
50    }
51
52    fn delete(&mut self, doc_id: DocId) -> StorageBackendResult<()> {
53        self.store.delete(&document_key(&self.table, doc_id)?)
54    }
55
56    fn clear(&mut self) -> StorageBackendResult<()> {
57        self.store
58            .delete_prefix(&document_key_prefix(&self.table)?)
59            .map(|_| ())
60    }
61
62    fn doc_ids(&self) -> StorageBackendResult<Vec<DocId>> {
63        let mut out = Vec::new();
64        for (key, _) in self.store.scan_prefix(&document_key_prefix(&self.table)?)? {
65            let mut offset = 1;
66            let _table = read_str(&key, &mut offset)?;
67            out.push(read_u64(&key, &mut offset)?);
68        }
69        Ok(out)
70    }
71
72    fn next_doc_id(&self, after: Option<DocId>) -> StorageBackendResult<Option<DocId>> {
73        Ok(self.next_doc_ids(after, 1)?.into_iter().next())
74    }
75
76    fn next_doc_ids(&self, after: Option<DocId>, limit: usize) -> StorageBackendResult<Vec<DocId>> {
77        if limit == 0 {
78            return Ok(Vec::new());
79        }
80        let prefix = document_key_prefix(&self.table)?;
81        let after_key = after
82            .map(|doc_id| document_key(&self.table, doc_id))
83            .transpose()?;
84        let mut out = Vec::with_capacity(limit);
85        for key in self
86            .store
87            .scan_prefix_keys_after(&prefix, after_key.as_deref(), limit)?
88        {
89            let mut offset = 1;
90            let _table = read_str(&key, &mut offset)?;
91            let doc_id = read_u64(&key, &mut offset)?;
92            out.push(doc_id);
93        }
94        Ok(out)
95    }
96
97    fn len(&self) -> StorageBackendResult<usize> {
98        Ok(self
99            .store
100            .scan_prefix(&document_key_prefix(&self.table)?)?
101            .len())
102    }
103
104    fn snapshot(&self) -> StorageBackendResult<Arc<dyn DocumentStore>> {
105        Ok(Arc::new(self.clone()))
106    }
107}