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::{BTreeMap, 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, SchemaRow, SequenceOptions,
18    SequenceReservationResult, SequenceRow, TableAclEntry, TableSchema, ViewRow,
19};
20use crate::{StorageBackendError, StorageBackendResult};
21
22use super::codec::{
23    decode_stored_document_value, decode_string, decode_value, doc_length_key,
24    doc_length_key_prefix, document_key_prefix, encode_stored_document_value, encode_value,
25    field_stats_key, field_stats_key_prefix, key_with_tag, posting_cluster_positions_field_prefix,
26    posting_cluster_positions_key_prefix, posting_cluster_score_field_prefix,
27    posting_cluster_score_key_prefix, posting_document_key, posting_document_key_prefix,
28    posting_field_prefix, posting_key_prefix, push_str, push_u64, read_str, read_u64,
29    reverse_posting_key, reverse_posting_key_prefix, single_str_key, string_value,
30    vector_field_prefix, vector_key_prefix,
31};
32use super::{
33    KeyValueBatch, KeyValueStore, TAG_ANALYZER, TAG_ANALYZER_DESCRIPTOR, TAG_CATALOG_INDEX,
34    TAG_COLUMN_STATS, TAG_EDGE, TAG_FIELD_ANALYZER_BINDING, TAG_FOREIGN_SERVER, TAG_FOREIGN_TABLE,
35    TAG_GRAPH_MEMBERSHIP, TAG_METADATA, TAG_MODEL, TAG_NAMED_GRAPH, TAG_PATH_INDEX, TAG_RELATION,
36    TAG_SCHEMA, TAG_SCORING_PARAMS, TAG_SEQUENCE, TAG_TABLE, TAG_TABLE_FIELD_ANALYZER, TAG_VERTEX,
37    TAG_VIEW,
38};
39
40mod analyzers;
41mod foreign;
42mod graph_access;
43mod graphs;
44mod indexes;
45mod keys;
46mod migration;
47mod models;
48mod occurrence_lifecycle;
49mod path_index_data;
50mod physical_indexes;
51mod records;
52mod relations;
53mod schema_table;
54mod sequences;
55mod views;
56
57use keys::{
58    batch_put_or_keep_existing, batch_rekey_prefix, batch_rekey_prefix_or_keep_existing,
59    catalog_index_references_column, catalog_index_rename_column, column_stats_key,
60    column_stats_prefix, decode_catalog_relation_key, decode_relation_key, edge_key,
61    ensure_prefix_absent, graph_membership_graph_prefix, graph_membership_key,
62    graph_membership_prefix, load_single_keys, load_single_string_rows,
63    register_migration_relation, relation_key, table_field_analyzer_field_prefix,
64    table_field_analyzer_key, table_field_analyzer_prefix, vertex_key,
65};
66use migration::{
67    apply_relation_migrations, collect_relation_migrations, validate_relation_parents,
68};
69use records::{
70    StoredCatalogIndex, StoredColumnStats, StoredEdge, StoredForeignServer, StoredForeignTable,
71    StoredRelation, StoredSequence, StoredVertex, StoredView,
72    STORED_FOREIGN_TABLE_SECURITY_VERSION,
73};
74
75#[derive(Clone)]
76pub struct KeyValueCatalog {
77    store: Arc<dyn KeyValueStore>,
78    sequence_lock: Arc<Mutex<()>>,
79    graph_indexes_lock: Arc<Mutex<()>>,
80}
81
82impl KeyValueCatalog {
83    pub fn new(store: Arc<dyn KeyValueStore>) -> Self {
84        Self {
85            store,
86            sequence_lock: Arc::new(Mutex::new(())),
87            graph_indexes_lock: Arc::new(Mutex::new(())),
88        }
89    }
90
91    pub fn store(&self) -> Arc<dyn KeyValueStore> {
92        Arc::clone(&self.store)
93    }
94}
95
96impl CatalogFacade for KeyValueCatalog {
97    fn clear_path_index_data(&self, index: &str) -> StorageBackendResult<()> {
98        self.clear_path_index_data_impl(index)
99    }
100
101    fn save_path_index_pairs(
102        &self,
103        index: &str,
104        sequence: &str,
105        pairs: &[(u64, u64)],
106    ) -> StorageBackendResult<()> {
107        self.save_path_index_pairs_impl(index, sequence, pairs)
108    }
109
110    fn finish_path_index_data(
111        &self,
112        index: &str,
113        graph: &str,
114        definition: &str,
115    ) -> StorageBackendResult<()> {
116        self.finish_path_index_data_impl(index, graph, definition)
117    }
118
119    fn path_index_data_is_current(
120        &self,
121        index: &str,
122        definition: &str,
123    ) -> StorageBackendResult<bool> {
124        self.path_index_data_is_current_impl(index, definition)
125    }
126
127    fn path_index_pairs(
128        &self,
129        index: &str,
130        sequence: &str,
131        after: Option<(u64, u64)>,
132        limit: usize,
133    ) -> StorageBackendResult<Vec<(u64, u64)>> {
134        self.path_index_pairs_impl(index, sequence, after, limit)
135    }
136
137    fn graph_vertex(&self, id: u64) -> StorageBackendResult<Option<crate::GraphVertexRow>> {
138        self.graph_vertex_impl(id)
139    }
140
141    fn graph_edge(&self, id: u64) -> StorageBackendResult<Option<EdgeRow>> {
142        self.graph_edge_impl(id)
143    }
144
145    fn graph_entity_ids(
146        &self,
147        filter: crate::GraphEntityFilter<'_>,
148        after: Option<u64>,
149        limit: usize,
150    ) -> StorageBackendResult<Vec<u64>> {
151        self.graph_entity_ids_impl(filter, after, limit)
152    }
153
154    fn graph_entity_count(
155        &self,
156        filter: crate::GraphEntityFilter<'_>,
157    ) -> StorageBackendResult<u64> {
158        let mut after = None;
159        let mut count = 0_u64;
160        loop {
161            let ids = self.graph_entity_ids_impl(filter, after, 256)?;
162            if ids.is_empty() {
163                return Ok(count);
164            }
165            after = ids.last().copied();
166            count = count
167                .checked_add(
168                    u64::try_from(ids.len())
169                        .map_err(|error| StorageBackendError::Other(error.to_string()))?,
170                )
171                .ok_or_else(|| StorageBackendError::Other("graph entity count overflow".into()))?;
172        }
173    }
174
175    fn graph_entity_max_id(
176        &self,
177        kind: crate::GraphEntityKind,
178    ) -> StorageBackendResult<Option<u64>> {
179        let filter = crate::GraphEntityFilter::new(kind, None);
180        let mut after = None;
181        loop {
182            let ids = self.graph_entity_ids_impl(filter, after, 256)?;
183            if ids.is_empty() {
184                return Ok(after);
185            }
186            after = ids.last().copied();
187        }
188    }
189
190    fn graph_entity_memberships(
191        &self,
192        kind: crate::GraphEntityKind,
193        id: u64,
194    ) -> StorageBackendResult<Vec<String>> {
195        self.graph_entity_memberships_impl(kind, id)
196    }
197
198    fn graph_has_membership(
199        &self,
200        kind: crate::GraphEntityKind,
201        id: u64,
202        graph: &str,
203    ) -> StorageBackendResult<bool> {
204        self.store
205            .contains_key(&graph_membership_key(kind.as_str(), id, graph)?)
206    }
207
208    fn set_metadata(&self, key: &str, value: &str) -> StorageBackendResult<()> {
209        self.set_metadata_impl(key, value)
210    }
211
212    fn get_metadata(&self, key: &str) -> StorageBackendResult<Option<String>> {
213        self.get_metadata_impl(key)
214    }
215
216    fn migrate_relation_namespace(&self) -> StorageBackendResult<()> {
217        self.migrate_relation_namespace_impl()?;
218        self.ensure_graph_lookup_indexes()
219    }
220
221    fn save_schema_row(&self, schema: &SchemaRow) -> StorageBackendResult<()> {
222        self.save_schema_row_impl(schema)
223    }
224
225    fn drop_schema(&self, name: &str) -> StorageBackendResult<()> {
226        self.drop_schema_impl(name)
227    }
228
229    fn load_schema_rows(&self) -> StorageBackendResult<Vec<SchemaRow>> {
230        self.load_schema_rows_impl()
231    }
232
233    fn save_table(&self, schema: &TableSchema) -> StorageBackendResult<()> {
234        self.save_table_impl(schema)
235    }
236
237    fn load_tables(&self) -> StorageBackendResult<Vec<TableSchema>> {
238        self.load_tables_impl()
239    }
240
241    fn drop_table(&self, name: &str) -> StorageBackendResult<()> {
242        self.drop_table_impl(name)
243    }
244
245    fn drop_table_and_data(&self, name: &str) -> StorageBackendResult<()> {
246        self.drop_table_and_data_impl(name)
247    }
248
249    fn purge_table_data(&self, name: &str) -> StorageBackendResult<()> {
250        self.purge_table_data_impl(name)
251    }
252
253    fn rename_table_data(&self, from: &str, to: &str) -> StorageBackendResult<()> {
254        self.rename_table_data_impl(from, to)
255    }
256
257    fn drop_column_data(&self, table_name: &str, column_name: &str) -> StorageBackendResult<()> {
258        self.drop_column_data_impl(table_name, column_name)
259    }
260
261    fn rename_column_data(
262        &self,
263        table_name: &str,
264        from: &str,
265        to: &str,
266    ) -> StorageBackendResult<()> {
267        self.rename_column_data_impl(table_name, from, to)
268    }
269
270    fn save_model(&self, name: &str, json: &str) -> StorageBackendResult<()> {
271        self.save_model_impl(name, json)
272    }
273
274    fn load_models(&self) -> StorageBackendResult<Vec<(String, String)>> {
275        self.load_models_impl()
276    }
277
278    fn load_model(&self, name: &str) -> StorageBackendResult<Option<String>> {
279        self.load_model_impl(name)
280    }
281
282    fn drop_model(&self, name: &str) -> StorageBackendResult<()> {
283        self.drop_model_impl(name)
284    }
285
286    fn save_scoring_params(&self, name: &str, params_json: &str) -> StorageBackendResult<()> {
287        self.save_scoring_params_impl(name, params_json)
288    }
289
290    fn load_scoring_params(&self, name: &str) -> StorageBackendResult<Option<String>> {
291        self.load_scoring_params_impl(name)
292    }
293
294    fn load_all_scoring_params(&self) -> StorageBackendResult<Vec<(String, String)>> {
295        self.load_all_scoring_params_impl()
296    }
297
298    fn drop_scoring_params(&self, name: &str) -> StorageBackendResult<()> {
299        self.drop_scoring_params_impl(name)
300    }
301
302    fn create_sequence_row(&self, sequence: &SequenceRow) -> StorageBackendResult<bool> {
303        self.create_sequence_row_impl(sequence)
304    }
305
306    fn replace_sequence_row(&self, sequence: &SequenceRow) -> StorageBackendResult<bool> {
307        self.replace_sequence_row_impl(sequence)
308    }
309
310    fn rename_sequence_row(&self, from: &str, to: &str) -> StorageBackendResult<bool> {
311        self.rename_sequence_row_impl(from, to)
312    }
313
314    fn drop_sequence_row(&self, name: &str) -> StorageBackendResult<bool> {
315        self.drop_sequence_row_impl(name)
316    }
317
318    fn load_sequence_rows(&self) -> StorageBackendResult<Vec<SequenceRow>> {
319        self.load_sequence_rows_impl()
320    }
321
322    fn reserve_sequence_values(
323        &self,
324        name: &str,
325        object_id: [u8; 16],
326        definition_generation: [u8; 16],
327    ) -> StorageBackendResult<SequenceReservationResult> {
328        self.reserve_sequence_values_impl(name, object_id, definition_generation)
329    }
330
331    fn set_sequence_value(
332        &self,
333        name: &str,
334        object_id: [u8; 16],
335        value: i64,
336        called: bool,
337        log_count: i64,
338    ) -> StorageBackendResult<Option<i64>> {
339        self.set_sequence_value_impl(name, object_id, value, called, log_count)
340    }
341
342    fn save_view(&self, view: &ViewRow) -> StorageBackendResult<()> {
343        self.save_view_impl(view)
344    }
345
346    fn rename_view(
347        &self,
348        from: &RelationIdentity,
349        to: &RelationIdentity,
350    ) -> StorageBackendResult<bool> {
351        self.rename_view_impl(from, to)
352    }
353
354    fn drop_view(&self, relation: &RelationIdentity) -> StorageBackendResult<bool> {
355        self.drop_view_impl(relation)
356    }
357
358    fn load_views(&self) -> StorageBackendResult<Vec<ViewRow>> {
359        self.load_views_impl()
360    }
361
362    fn save_named_graph(&self, name: &str) -> StorageBackendResult<()> {
363        self.save_named_graph_impl(name)
364    }
365
366    fn drop_named_graph(&self, name: &str) -> StorageBackendResult<()> {
367        self.drop_named_graph_impl(name)
368    }
369
370    fn load_named_graphs(&self) -> StorageBackendResult<Vec<String>> {
371        self.load_named_graphs_impl()
372    }
373    fn named_graph_exists(&self, name: &str) -> StorageBackendResult<bool> {
374        self.store
375            .contains_key(&single_str_key(TAG_NAMED_GRAPH, name)?)
376    }
377
378    fn save_vertex(
379        &self,
380        vertex_id: u64,
381        label: &str,
382        properties_json: &str,
383    ) -> StorageBackendResult<()> {
384        self.save_vertex_impl(vertex_id, label, properties_json)
385    }
386
387    fn delete_vertex(&self, vertex_id: u64) -> StorageBackendResult<()> {
388        self.delete_vertex_impl(vertex_id)
389    }
390
391    fn load_vertices(&self) -> StorageBackendResult<Vec<(u64, String, String)>> {
392        self.load_vertices_impl()
393    }
394
395    fn save_edge(
396        &self,
397        edge_id: u64,
398        source_id: u64,
399        target_id: u64,
400        label: &str,
401        properties_json: &str,
402    ) -> StorageBackendResult<()> {
403        self.save_edge_impl(edge_id, source_id, target_id, label, properties_json)
404    }
405
406    fn delete_edge(&self, edge_id: u64) -> StorageBackendResult<()> {
407        self.delete_edge_impl(edge_id)
408    }
409
410    fn load_edges(&self) -> StorageBackendResult<Vec<EdgeRow>> {
411        self.load_edges_impl()
412    }
413
414    fn save_graph_membership(
415        &self,
416        entity_type: &str,
417        entity_id: u64,
418        graph_name: &str,
419    ) -> StorageBackendResult<()> {
420        self.save_graph_membership_impl(entity_type, entity_id, graph_name)
421    }
422
423    fn delete_graph_membership(
424        &self,
425        entity_type: &str,
426        entity_id: u64,
427        graph_name: &str,
428    ) -> StorageBackendResult<()> {
429        self.delete_graph_membership_impl(entity_type, entity_id, graph_name)
430    }
431
432    fn delete_graph_membership_for_graph(&self, graph_name: &str) -> StorageBackendResult<()> {
433        self.delete_graph_membership_for_graph_impl(graph_name)
434    }
435
436    fn load_graph_memberships(&self) -> StorageBackendResult<Vec<(String, u64, String)>> {
437        self.load_graph_memberships_impl()
438    }
439
440    fn purge_orphan_graph_entities(&self) -> StorageBackendResult<()> {
441        self.purge_orphan_graph_entities_impl()
442    }
443
444    fn replace_named_graph(
445        &self,
446        graph_name: &str,
447        snapshot: &GraphSnapshot,
448    ) -> StorageBackendResult<()> {
449        self.replace_named_graph_impl(graph_name, snapshot)
450    }
451
452    fn drop_named_graph_data(&self, graph_name: &str) -> StorageBackendResult<()> {
453        self.drop_named_graph_data_impl(graph_name)
454    }
455
456    fn save_analyzer(&self, name: &str, config_json: &str) -> StorageBackendResult<()> {
457        self.save_analyzer_impl(name, config_json)
458    }
459
460    fn drop_analyzer(&self, name: &str) -> StorageBackendResult<()> {
461        self.drop_analyzer_impl(name)
462    }
463
464    fn load_analyzers(&self) -> StorageBackendResult<Vec<(String, String)>> {
465        self.load_analyzers_impl()
466    }
467
468    fn save_analyzer_revision(
469        &self,
470        name: &str,
471        config_json: &str,
472        descriptor_json: &str,
473    ) -> StorageBackendResult<()> {
474        self.save_analyzer_revision_impl(name, config_json, descriptor_json)
475    }
476
477    fn load_analyzer_descriptors(&self) -> StorageBackendResult<Vec<(String, String)>> {
478        load_single_string_rows(self.store.as_ref(), TAG_ANALYZER_DESCRIPTOR)
479    }
480
481    fn replace_table_field_analyzer_binding(
482        &self,
483        table: &str,
484        field: &str,
485        phase: &str,
486        name: &str,
487        binding_json: &str,
488    ) -> StorageBackendResult<()> {
489        self.replace_table_field_analyzer_binding_impl(table, field, phase, name, binding_json)
490    }
491
492    fn load_table_field_analyzer_bindings(
493        &self,
494    ) -> StorageBackendResult<Vec<(String, String, String)>> {
495        self.load_table_field_analyzer_bindings_impl()
496    }
497
498    fn save_table_field_analyzer(
499        &self,
500        table_name: &str,
501        field: &str,
502        phase: &str,
503        analyzer_name: &str,
504    ) -> StorageBackendResult<()> {
505        self.save_table_field_analyzer_impl(table_name, field, phase, analyzer_name)
506    }
507
508    fn replace_table_field_analyzer(
509        &self,
510        table_name: &str,
511        field: &str,
512        phase: &str,
513        analyzer_name: &str,
514    ) -> StorageBackendResult<()> {
515        self.replace_table_field_analyzer_impl(table_name, field, phase, analyzer_name)
516    }
517
518    fn drop_table_field_analyzer_field(
519        &self,
520        table_name: &str,
521        field: &str,
522    ) -> StorageBackendResult<()> {
523        self.drop_table_field_analyzer_field_impl(table_name, field)
524    }
525
526    fn drop_table_field_analyzers(&self, table_name: &str) -> StorageBackendResult<()> {
527        self.drop_table_field_analyzers_impl(table_name)
528    }
529
530    fn load_table_field_analyzers(
531        &self,
532    ) -> StorageBackendResult<Vec<(String, String, String, String)>> {
533        self.load_table_field_analyzers_impl()
534    }
535
536    fn save_foreign_server(
537        &self,
538        name: &str,
539        fdw_type: &str,
540        options_json: &str,
541    ) -> StorageBackendResult<()> {
542        self.save_foreign_server_impl(name, fdw_type, options_json)
543    }
544
545    fn drop_foreign_server(&self, name: &str) -> StorageBackendResult<()> {
546        self.drop_foreign_server_impl(name)
547    }
548
549    fn load_foreign_servers(&self) -> StorageBackendResult<Vec<(String, String, String)>> {
550        self.load_foreign_servers_impl()
551    }
552
553    fn save_foreign_table(&self, row: &ForeignTableRow) -> StorageBackendResult<()> {
554        self.save_foreign_table_impl(row)
555    }
556
557    fn rename_foreign_table(
558        &self,
559        from: &RelationIdentity,
560        to: &RelationIdentity,
561    ) -> StorageBackendResult<bool> {
562        self.rename_foreign_table_impl(from, to)
563    }
564
565    fn update_foreign_table_security(
566        &self,
567        relation: &RelationIdentity,
568        role_owner: &str,
569        acl: Option<&[TableAclEntry]>,
570        column_acls: &BTreeMap<String, Vec<TableAclEntry>>,
571    ) -> StorageBackendResult<bool> {
572        self.update_foreign_table_security_impl(relation, role_owner, acl, column_acls)
573    }
574
575    fn drop_foreign_table(&self, relation: &RelationIdentity) -> StorageBackendResult<()> {
576        self.drop_foreign_table_impl(relation)
577    }
578
579    fn load_foreign_tables(&self) -> StorageBackendResult<Vec<ForeignTableRow>> {
580        self.load_foreign_tables_impl()
581    }
582
583    fn save_catalog_index_row(&self, index: &CatalogIndexRow) -> StorageBackendResult<()> {
584        self.save_catalog_index_impl(index)
585    }
586
587    fn drop_catalog_index(&self, relation: &RelationIdentity) -> StorageBackendResult<()> {
588        self.drop_catalog_index_impl(relation)
589    }
590
591    fn drop_catalog_indexes_for_table(&self, table_name: &str) -> StorageBackendResult<()> {
592        self.drop_catalog_indexes_for_table_impl(table_name)
593    }
594
595    fn load_catalog_indexes(&self) -> StorageBackendResult<Vec<CatalogIndexRow>> {
596        self.load_catalog_indexes_impl()
597    }
598
599    fn save_path_index(
600        &self,
601        graph_name: &str,
602        label_sequences_json: &str,
603    ) -> StorageBackendResult<()> {
604        self.save_path_index_impl(graph_name, label_sequences_json)
605    }
606
607    fn drop_path_index(&self, graph_name: &str) -> StorageBackendResult<()> {
608        self.drop_path_index_impl(graph_name)
609    }
610
611    fn load_path_indexes(&self) -> StorageBackendResult<Vec<(String, String)>> {
612        self.load_path_indexes_impl()
613    }
614
615    fn save_column_stats(&self, stats: ColumnStatsInput<'_>) -> StorageBackendResult<()> {
616        self.save_column_stats_impl(stats)
617    }
618
619    fn replace_column_stats(
620        &self,
621        table_name: &str,
622        stats: &[ColumnStatsInput<'_>],
623    ) -> StorageBackendResult<()> {
624        self.replace_column_stats_impl(table_name, stats)
625    }
626
627    fn load_column_stats(&self, table_name: &str) -> StorageBackendResult<Vec<ColumnStatsRow>> {
628        self.load_column_stats_impl(table_name)
629    }
630
631    fn delete_column_stats(&self, table_name: &str) -> StorageBackendResult<()> {
632        self.delete_column_stats_impl(table_name)
633    }
634}
635
636#[cfg(test)]
637mod tests;