1use std::collections::BTreeSet;
10use std::sync::Arc;
11
12use parking_lot::Mutex;
13use serde::{Deserialize, Serialize};
14
15use crate::catalog::{
16 CatalogFacade, CatalogIndexRow, ColumnStatsInput, ColumnStatsRow, EdgeRow, ForeignTableRow,
17 GraphSnapshot, RelationIdentity, RelationKind, SequenceRow, TableSchema, ViewRow,
18};
19use crate::{StorageBackendError, StorageBackendResult};
20
21use super::codec::{
22 decode_document_value, decode_string, decode_value, doc_length_key, doc_length_key_prefix,
23 document_key_prefix, encode_document_value, encode_value, field_stats_key,
24 field_stats_key_prefix, key_with_tag, posting_cluster_positions_field_prefix,
25 posting_cluster_positions_key_prefix, posting_cluster_score_field_prefix,
26 posting_cluster_score_key_prefix, posting_document_key, posting_document_key_prefix,
27 posting_field_prefix, posting_key_prefix, push_str, push_u64, read_str, read_u64,
28 reverse_posting_key, reverse_posting_key_prefix, single_str_key, string_value,
29 vector_field_prefix, vector_key_prefix,
30};
31use super::{
32 KeyValueBatch, KeyValueStore, TAG_ANALYZER, TAG_CATALOG_INDEX, TAG_COLUMN_STATS, TAG_EDGE,
33 TAG_FOREIGN_SERVER, TAG_FOREIGN_TABLE, TAG_GRAPH_MEMBERSHIP, TAG_METADATA, TAG_MODEL,
34 TAG_NAMED_GRAPH, TAG_PATH_INDEX, TAG_RELATION, TAG_SCHEMA, TAG_SCORING_PARAMS, TAG_SEQUENCE,
35 TAG_TABLE, TAG_TABLE_FIELD_ANALYZER, TAG_VERTEX, TAG_VIEW,
36};
37
38mod analyzers;
39mod foreign;
40mod graphs;
41mod indexes;
42mod keys;
43mod migration;
44mod models;
45mod physical_indexes;
46mod records;
47mod relations;
48mod schema_table;
49mod sequences;
50mod views;
51
52use keys::{
53 batch_put_or_keep_existing, batch_rekey_prefix, batch_rekey_prefix_or_keep_existing,
54 catalog_index_references_column, catalog_index_rename_column, column_stats_key,
55 column_stats_prefix, decode_catalog_relation_key, decode_relation_key, edge_key,
56 ensure_prefix_absent, graph_membership_graph_prefix, graph_membership_key,
57 graph_membership_prefix, load_single_keys, load_single_string_rows,
58 register_migration_relation, relation_key, table_field_analyzer_field_prefix,
59 table_field_analyzer_key, table_field_analyzer_prefix, vertex_key,
60};
61use migration::{
62 apply_relation_migrations, collect_relation_migrations, validate_relation_parents,
63};
64use records::{
65 StoredCatalogIndex, StoredColumnStats, StoredEdge, StoredForeignServer, StoredForeignTable,
66 StoredRelation, StoredSequence, StoredVertex, StoredView,
67};
68
69#[derive(Clone)]
70pub struct KeyValueCatalog {
71 store: Arc<dyn KeyValueStore>,
72 sequence_lock: Arc<Mutex<()>>,
73}
74
75impl KeyValueCatalog {
76 pub fn new(store: Arc<dyn KeyValueStore>) -> Self {
77 Self {
78 store,
79 sequence_lock: Arc::new(Mutex::new(())),
80 }
81 }
82
83 pub fn store(&self) -> Arc<dyn KeyValueStore> {
84 Arc::clone(&self.store)
85 }
86}
87
88impl CatalogFacade for KeyValueCatalog {
89 fn set_metadata(&self, key: &str, value: &str) -> StorageBackendResult<()> {
90 self.set_metadata_impl(key, value)
91 }
92
93 fn get_metadata(&self, key: &str) -> StorageBackendResult<Option<String>> {
94 self.get_metadata_impl(key)
95 }
96
97 fn migrate_relation_namespace(&self) -> StorageBackendResult<()> {
98 self.migrate_relation_namespace_impl()
99 }
100
101 fn save_schema(&self, name: &str) -> StorageBackendResult<()> {
102 self.save_schema_impl(name)
103 }
104
105 fn drop_schema(&self, name: &str) -> StorageBackendResult<()> {
106 self.drop_schema_impl(name)
107 }
108
109 fn load_schemas(&self) -> StorageBackendResult<Vec<String>> {
110 self.load_schemas_impl()
111 }
112
113 fn save_table(&self, schema: &TableSchema) -> StorageBackendResult<()> {
114 self.save_table_impl(schema)
115 }
116
117 fn load_tables(&self) -> StorageBackendResult<Vec<TableSchema>> {
118 self.load_tables_impl()
119 }
120
121 fn drop_table(&self, name: &str) -> StorageBackendResult<()> {
122 self.drop_table_impl(name)
123 }
124
125 fn drop_table_and_data(&self, name: &str) -> StorageBackendResult<()> {
126 self.drop_table_and_data_impl(name)
127 }
128
129 fn purge_table_data(&self, name: &str) -> StorageBackendResult<()> {
130 self.purge_table_data_impl(name)
131 }
132
133 fn rename_table_data(&self, from: &str, to: &str) -> StorageBackendResult<()> {
134 self.rename_table_data_impl(from, to)
135 }
136
137 fn drop_column_data(&self, table_name: &str, column_name: &str) -> StorageBackendResult<()> {
138 self.drop_column_data_impl(table_name, column_name)
139 }
140
141 fn rename_column_data(
142 &self,
143 table_name: &str,
144 from: &str,
145 to: &str,
146 ) -> StorageBackendResult<()> {
147 self.rename_column_data_impl(table_name, from, to)
148 }
149
150 fn save_model(&self, name: &str, json: &str) -> StorageBackendResult<()> {
151 self.save_model_impl(name, json)
152 }
153
154 fn load_models(&self) -> StorageBackendResult<Vec<(String, String)>> {
155 self.load_models_impl()
156 }
157
158 fn load_model(&self, name: &str) -> StorageBackendResult<Option<String>> {
159 self.load_model_impl(name)
160 }
161
162 fn drop_model(&self, name: &str) -> StorageBackendResult<()> {
163 self.drop_model_impl(name)
164 }
165
166 fn save_scoring_params(&self, name: &str, params_json: &str) -> StorageBackendResult<()> {
167 self.save_scoring_params_impl(name, params_json)
168 }
169
170 fn load_scoring_params(&self, name: &str) -> StorageBackendResult<Option<String>> {
171 self.load_scoring_params_impl(name)
172 }
173
174 fn load_all_scoring_params(&self) -> StorageBackendResult<Vec<(String, String)>> {
175 self.load_all_scoring_params_impl()
176 }
177
178 fn drop_scoring_params(&self, name: &str) -> StorageBackendResult<()> {
179 self.drop_scoring_params_impl(name)
180 }
181
182 fn create_sequence_row(&self, sequence: &SequenceRow) -> StorageBackendResult<bool> {
183 self.create_sequence_row_impl(sequence)
184 }
185
186 fn replace_sequence_row(&self, sequence: &SequenceRow) -> StorageBackendResult<bool> {
187 self.replace_sequence_row_impl(sequence)
188 }
189
190 fn drop_sequence_row(&self, name: &str) -> StorageBackendResult<bool> {
191 self.drop_sequence_row_impl(name)
192 }
193
194 fn load_sequence_rows(&self) -> StorageBackendResult<Vec<SequenceRow>> {
195 self.load_sequence_rows_impl()
196 }
197
198 fn next_sequence_value(&self, name: &str) -> StorageBackendResult<Option<i64>> {
199 self.next_sequence_value_impl(name)
200 }
201
202 fn set_sequence_value(&self, name: &str, value: i64) -> StorageBackendResult<Option<i64>> {
203 self.set_sequence_value_impl(name, value)
204 }
205
206 fn save_view(&self, view: &ViewRow) -> StorageBackendResult<()> {
207 self.save_view_impl(view)
208 }
209
210 fn drop_view(&self, relation: &RelationIdentity) -> StorageBackendResult<bool> {
211 self.drop_view_impl(relation)
212 }
213
214 fn load_views(&self) -> StorageBackendResult<Vec<ViewRow>> {
215 self.load_views_impl()
216 }
217
218 fn save_named_graph(&self, name: &str) -> StorageBackendResult<()> {
219 self.save_named_graph_impl(name)
220 }
221
222 fn drop_named_graph(&self, name: &str) -> StorageBackendResult<()> {
223 self.drop_named_graph_impl(name)
224 }
225
226 fn load_named_graphs(&self) -> StorageBackendResult<Vec<String>> {
227 self.load_named_graphs_impl()
228 }
229
230 fn save_vertex(
231 &self,
232 vertex_id: u64,
233 label: &str,
234 properties_json: &str,
235 ) -> StorageBackendResult<()> {
236 self.save_vertex_impl(vertex_id, label, properties_json)
237 }
238
239 fn delete_vertex(&self, vertex_id: u64) -> StorageBackendResult<()> {
240 self.delete_vertex_impl(vertex_id)
241 }
242
243 fn load_vertices(&self) -> StorageBackendResult<Vec<(u64, String, String)>> {
244 self.load_vertices_impl()
245 }
246
247 fn save_edge(
248 &self,
249 edge_id: u64,
250 source_id: u64,
251 target_id: u64,
252 label: &str,
253 properties_json: &str,
254 ) -> StorageBackendResult<()> {
255 self.save_edge_impl(edge_id, source_id, target_id, label, properties_json)
256 }
257
258 fn delete_edge(&self, edge_id: u64) -> StorageBackendResult<()> {
259 self.delete_edge_impl(edge_id)
260 }
261
262 fn load_edges(&self) -> StorageBackendResult<Vec<EdgeRow>> {
263 self.load_edges_impl()
264 }
265
266 fn save_graph_membership(
267 &self,
268 entity_type: &str,
269 entity_id: u64,
270 graph_name: &str,
271 ) -> StorageBackendResult<()> {
272 self.save_graph_membership_impl(entity_type, entity_id, graph_name)
273 }
274
275 fn delete_graph_membership(
276 &self,
277 entity_type: &str,
278 entity_id: u64,
279 graph_name: &str,
280 ) -> StorageBackendResult<()> {
281 self.delete_graph_membership_impl(entity_type, entity_id, graph_name)
282 }
283
284 fn delete_graph_membership_for_graph(&self, graph_name: &str) -> StorageBackendResult<()> {
285 self.delete_graph_membership_for_graph_impl(graph_name)
286 }
287
288 fn load_graph_memberships(&self) -> StorageBackendResult<Vec<(String, u64, String)>> {
289 self.load_graph_memberships_impl()
290 }
291
292 fn purge_orphan_graph_entities(&self) -> StorageBackendResult<()> {
293 self.purge_orphan_graph_entities_impl()
294 }
295
296 fn replace_named_graph(
297 &self,
298 graph_name: &str,
299 snapshot: &GraphSnapshot,
300 ) -> StorageBackendResult<()> {
301 self.replace_named_graph_impl(graph_name, snapshot)
302 }
303
304 fn drop_named_graph_data(&self, graph_name: &str) -> StorageBackendResult<()> {
305 self.drop_named_graph_data_impl(graph_name)
306 }
307
308 fn save_analyzer(&self, name: &str, config_json: &str) -> StorageBackendResult<()> {
309 self.save_analyzer_impl(name, config_json)
310 }
311
312 fn drop_analyzer(&self, name: &str) -> StorageBackendResult<()> {
313 self.drop_analyzer_impl(name)
314 }
315
316 fn load_analyzers(&self) -> StorageBackendResult<Vec<(String, String)>> {
317 self.load_analyzers_impl()
318 }
319
320 fn save_table_field_analyzer(
321 &self,
322 table_name: &str,
323 field: &str,
324 phase: &str,
325 analyzer_name: &str,
326 ) -> StorageBackendResult<()> {
327 self.save_table_field_analyzer_impl(table_name, field, phase, analyzer_name)
328 }
329
330 fn replace_table_field_analyzer(
331 &self,
332 table_name: &str,
333 field: &str,
334 phase: &str,
335 analyzer_name: &str,
336 ) -> StorageBackendResult<()> {
337 self.replace_table_field_analyzer_impl(table_name, field, phase, analyzer_name)
338 }
339
340 fn drop_table_field_analyzer_field(
341 &self,
342 table_name: &str,
343 field: &str,
344 ) -> StorageBackendResult<()> {
345 self.drop_table_field_analyzer_field_impl(table_name, field)
346 }
347
348 fn drop_table_field_analyzers(&self, table_name: &str) -> StorageBackendResult<()> {
349 self.drop_table_field_analyzers_impl(table_name)
350 }
351
352 fn load_table_field_analyzers(
353 &self,
354 ) -> StorageBackendResult<Vec<(String, String, String, String)>> {
355 self.load_table_field_analyzers_impl()
356 }
357
358 fn save_foreign_server(
359 &self,
360 name: &str,
361 fdw_type: &str,
362 options_json: &str,
363 ) -> StorageBackendResult<()> {
364 self.save_foreign_server_impl(name, fdw_type, options_json)
365 }
366
367 fn drop_foreign_server(&self, name: &str) -> StorageBackendResult<()> {
368 self.drop_foreign_server_impl(name)
369 }
370
371 fn load_foreign_servers(&self) -> StorageBackendResult<Vec<(String, String, String)>> {
372 self.load_foreign_servers_impl()
373 }
374
375 fn save_foreign_table(
376 &self,
377 relation: &RelationIdentity,
378 server_name: &str,
379 columns_json: &str,
380 options_json: &str,
381 ) -> StorageBackendResult<()> {
382 self.save_foreign_table_impl(relation, server_name, columns_json, options_json)
383 }
384
385 fn drop_foreign_table(&self, relation: &RelationIdentity) -> StorageBackendResult<()> {
386 self.drop_foreign_table_impl(relation)
387 }
388
389 fn load_foreign_tables(&self) -> StorageBackendResult<Vec<ForeignTableRow>> {
390 self.load_foreign_tables_impl()
391 }
392
393 fn save_catalog_index(
394 &self,
395 name: &str,
396 index_type: &str,
397 table_name: &str,
398 columns_json: &str,
399 parameters_json: &str,
400 ) -> StorageBackendResult<()> {
401 self.save_catalog_index_impl(name, index_type, table_name, columns_json, parameters_json)
402 }
403
404 fn drop_catalog_index(&self, name: &str) -> StorageBackendResult<()> {
405 self.drop_catalog_index_impl(name)
406 }
407
408 fn drop_catalog_indexes_for_table(&self, table_name: &str) -> StorageBackendResult<()> {
409 self.drop_catalog_indexes_for_table_impl(table_name)
410 }
411
412 fn load_catalog_indexes(&self) -> StorageBackendResult<Vec<CatalogIndexRow>> {
413 self.load_catalog_indexes_impl()
414 }
415
416 fn save_path_index(
417 &self,
418 graph_name: &str,
419 label_sequences_json: &str,
420 ) -> StorageBackendResult<()> {
421 self.save_path_index_impl(graph_name, label_sequences_json)
422 }
423
424 fn drop_path_index(&self, graph_name: &str) -> StorageBackendResult<()> {
425 self.drop_path_index_impl(graph_name)
426 }
427
428 fn load_path_indexes(&self) -> StorageBackendResult<Vec<(String, String)>> {
429 self.load_path_indexes_impl()
430 }
431
432 fn save_column_stats(&self, stats: ColumnStatsInput<'_>) -> StorageBackendResult<()> {
433 self.save_column_stats_impl(stats)
434 }
435
436 fn replace_column_stats(
437 &self,
438 table_name: &str,
439 stats: &[ColumnStatsInput<'_>],
440 ) -> StorageBackendResult<()> {
441 self.replace_column_stats_impl(table_name, stats)
442 }
443
444 fn load_column_stats(&self, table_name: &str) -> StorageBackendResult<Vec<ColumnStatsRow>> {
445 self.load_column_stats_impl(table_name)
446 }
447
448 fn delete_column_stats(&self, table_name: &str) -> StorageBackendResult<()> {
449 self.delete_column_stats_impl(table_name)
450 }
451}
452
453#[cfg(test)]
454mod tests;