1use 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#[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}