Skip to main content

uqa_storage/key_value/
catalog.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Catalog facade implementation for key/value-backed persistence.
8
9use 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;