1use 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;