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, decode_stored_document_value,
11    decode_stored_document_value_for_migration, decode_value, document_key, document_key_prefix,
12    document_value_is_current, encode_stored_document_value, key_with_tag, other_error, read_str,
13    read_u64, single_str_key, string_value,
14};
15use super::{
16    Arc, DocId, Document, DocumentMetadata, DocumentStore, KeyValueStore, StorageBackendResult,
17    StoredDocument, Value, TAG_DOCUMENT, TAG_METADATA, TAG_TABLE,
18};
19use crate::TableSchema;
20
21const DOCUMENT_FORMAT_METADATA_KEY: &str = "document_storage_format";
22const DOCUMENT_FORMAT_NAME: &str = "record-v2";
23const MIGRATION_PAGE_SIZE: usize = 512;
24
25/// Document store implemented over [`KeyValueStore`].
26#[derive(Clone)]
27pub struct KeyValueDocumentStore {
28    store: Arc<dyn KeyValueStore>,
29    table: String,
30}
31
32impl KeyValueDocumentStore {
33    pub fn new(store: Arc<dyn KeyValueStore>, table: impl Into<String>) -> Self {
34        Self {
35            store,
36            table: table.into(),
37        }
38    }
39
40    pub(crate) fn migrate_legacy_storage(store: &dyn KeyValueStore) -> StorageBackendResult<()> {
41        let marker = single_str_key(TAG_METADATA, DOCUMENT_FORMAT_METADATA_KEY)?;
42        if let Some(format) = store.get(&marker)? {
43            if format == DOCUMENT_FORMAT_NAME.as_bytes() {
44                return Ok(());
45            }
46            return Err(other_error(format!(
47                "unsupported KeyValue document format `{}`",
48                String::from_utf8_lossy(&format)
49            )));
50        }
51        if store.in_transaction() {
52            return Err(other_error(
53                "cannot migrate KeyValue documents inside an active transaction",
54            ));
55        }
56        store.begin_transaction()?;
57        let migration = Self::migrate_legacy_storage_in_transaction(store, &marker);
58        match migration {
59            Ok(()) => store.commit_transaction(),
60            Err(error) => match store.rollback_transaction() {
61                Ok(()) => Err(error),
62                Err(rollback) => Err(other_error(format!(
63                    "{error}; KeyValue document migration rollback also failed: {rollback}"
64                ))),
65            },
66        }
67    }
68
69    fn migrate_legacy_storage_in_transaction(
70        store: &dyn KeyValueStore,
71        marker: &[u8],
72    ) -> StorageBackendResult<()> {
73        let (known_tables, declared_xmin_tables) = catalog_xmin_tables(store)?;
74        let prefix = key_with_tag(TAG_DOCUMENT);
75        let mut after = None::<Vec<u8>>;
76        loop {
77            let page = store.scan_prefix_after(&prefix, after.as_deref(), MIGRATION_PAGE_SIZE)?;
78            if page.is_empty() {
79                break;
80            }
81            for (key, value) in page {
82                after = Some(key.clone());
83                if document_value_is_current(&value) {
84                    continue;
85                }
86                let mut offset = 1;
87                let table = read_str(&key, &mut offset)?;
88                let preserve_public_xmin =
89                    !known_tables.contains(&table) || declared_xmin_tables.contains(&table);
90                let document =
91                    decode_stored_document_value_for_migration(&value, preserve_public_xmin)?;
92                store.put(&key, &encode_stored_document_value(&document)?)?;
93            }
94        }
95        store.put(marker, &string_value(DOCUMENT_FORMAT_NAME))
96    }
97}
98
99fn catalog_xmin_tables(
100    store: &dyn KeyValueStore,
101) -> StorageBackendResult<(
102    std::collections::BTreeSet<String>,
103    std::collections::BTreeSet<String>,
104)> {
105    let mut known = std::collections::BTreeSet::new();
106    let mut declared_xmin = std::collections::BTreeSet::new();
107    for (_, value) in store.scan_prefix(&key_with_tag(TAG_TABLE))? {
108        let schema = decode_value::<TableSchema>(&value)?;
109        let definitions = serde_json::from_str::<Vec<serde_json::Value>>(&schema.columns_json)?;
110        let has_declared_xmin = definitions.iter().any(|definition| {
111            definition
112                .as_object()
113                .and_then(|definition| definition.get("name"))
114                .and_then(serde_json::Value::as_str)
115                == Some("xmin")
116        });
117        let aliases = schema.relation.canonical_and_legacy_public_names();
118        known.extend(aliases.iter().cloned());
119        if has_declared_xmin {
120            declared_xmin.extend(aliases);
121        }
122    }
123    Ok((known, declared_xmin))
124}
125
126impl DocumentStore for KeyValueDocumentStore {
127    fn put(&mut self, doc_id: DocId, document: Document) -> StorageBackendResult<()> {
128        let metadata = self.get_metadata(doc_id)?.unwrap_or_default();
129        self.put_stored(doc_id, StoredDocument::with_metadata(document, metadata))
130    }
131
132    fn get(&self, doc_id: DocId) -> StorageBackendResult<Option<Document>> {
133        self.store
134            .get(&document_key(&self.table, doc_id)?)?
135            .map(|bytes| decode_document_value(&bytes))
136            .transpose()
137    }
138
139    fn put_stored(&mut self, doc_id: DocId, document: StoredDocument) -> StorageBackendResult<()> {
140        let (fields, metadata) = document.into_parts();
141        let fields = fields
142            .into_iter()
143            .filter(|(_, value)| !matches!(value, Value::Null))
144            .collect();
145        let document = StoredDocument::with_metadata(fields, metadata);
146        let value = encode_stored_document_value(&document)?;
147        self.store.put(&document_key(&self.table, doc_id)?, &value)
148    }
149
150    fn get_stored(&self, doc_id: DocId) -> StorageBackendResult<Option<StoredDocument>> {
151        self.store
152            .get(&document_key(&self.table, doc_id)?)?
153            .map(|bytes| decode_stored_document_value(&bytes))
154            .transpose()
155    }
156
157    fn get_metadata(&self, doc_id: DocId) -> StorageBackendResult<Option<DocumentMetadata>> {
158        self.get_stored(doc_id)
159            .map(|document| document.map(|document| document.metadata()))
160    }
161
162    fn contains_doc_id(&self, doc_id: DocId) -> StorageBackendResult<bool> {
163        self.store.contains_key(&document_key(&self.table, doc_id)?)
164    }
165
166    fn delete(&mut self, doc_id: DocId) -> StorageBackendResult<()> {
167        self.store.delete(&document_key(&self.table, doc_id)?)
168    }
169
170    fn clear(&mut self) -> StorageBackendResult<()> {
171        self.store
172            .delete_prefix(&document_key_prefix(&self.table)?)
173            .map(|_| ())
174    }
175
176    fn doc_ids(&self) -> StorageBackendResult<Vec<DocId>> {
177        let mut out = Vec::new();
178        for (key, _) in self.store.scan_prefix(&document_key_prefix(&self.table)?)? {
179            let mut offset = 1;
180            let _table = read_str(&key, &mut offset)?;
181            out.push(read_u64(&key, &mut offset)?);
182        }
183        Ok(out)
184    }
185
186    fn next_doc_id(&self, after: Option<DocId>) -> StorageBackendResult<Option<DocId>> {
187        Ok(self.next_doc_ids(after, 1)?.into_iter().next())
188    }
189
190    fn next_doc_ids(&self, after: Option<DocId>, limit: usize) -> StorageBackendResult<Vec<DocId>> {
191        if limit == 0 {
192            return Ok(Vec::new());
193        }
194        let prefix = document_key_prefix(&self.table)?;
195        let after_key = after
196            .map(|doc_id| document_key(&self.table, doc_id))
197            .transpose()?;
198        let mut out = Vec::with_capacity(limit);
199        for key in self
200            .store
201            .scan_prefix_keys_after(&prefix, after_key.as_deref(), limit)?
202        {
203            let mut offset = 1;
204            let _table = read_str(&key, &mut offset)?;
205            let doc_id = read_u64(&key, &mut offset)?;
206            out.push(doc_id);
207        }
208        Ok(out)
209    }
210
211    fn len(&self) -> StorageBackendResult<usize> {
212        Ok(self
213            .store
214            .scan_prefix(&document_key_prefix(&self.table)?)?
215            .len())
216    }
217
218    fn snapshot(&self) -> StorageBackendResult<Arc<dyn DocumentStore>> {
219        Ok(Arc::new(self.clone()))
220    }
221}