1use std::collections::{HashMap, HashSet};
10use std::sync::Arc;
11
12use alopex_core::kv::{KVStore, KVTransaction};
13use alopex_core::types::TxnMode;
14use serde::{Deserialize, Serialize};
15use thiserror::Error;
16
17use crate::ast::ddl::{IndexMethod, VectorMetric};
18use crate::catalog::{
19 Catalog, ColumnMetadata, Compression, IndexMetadata, MemoryCatalog, RowIdMode,
20};
21use crate::catalog::{StorageOptions, StorageType, TableMetadata};
22use crate::planner::PlannerError;
23use crate::planner::types::ResolvedType;
24
25pub const CATALOG_PREFIX: &[u8] = b"__catalog__/";
27pub const CATALOGS_PREFIX: &[u8] = b"__catalog__/catalogs/";
28pub const NAMESPACES_PREFIX: &[u8] = b"__catalog__/namespaces/";
29pub const TABLES_PREFIX: &[u8] = b"__catalog__/tables/";
30pub const INDEXES_PREFIX: &[u8] = b"__catalog__/indexes/";
31pub const META_KEY: &[u8] = b"__catalog__/meta";
32
33const CATALOG_VERSION: u32 = 2;
34
35#[derive(Debug, Error)]
36pub enum CatalogError {
37 #[error("kv error: {0}")]
38 Kv(#[from] alopex_core::Error),
39
40 #[error("serialize error: {0}")]
41 Serialize(#[from] bincode::Error),
42
43 #[error("invalid catalog key: {0}")]
44 InvalidKey(String),
45}
46
47#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
48struct CatalogState {
49 version: u32,
50 table_id_counter: u32,
51 index_id_counter: u32,
52}
53
54#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
55pub struct PersistedCatalogMeta {
56 pub name: String,
57 pub comment: Option<String>,
58 pub storage_root: Option<String>,
59}
60
61#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
62pub struct PersistedNamespaceMeta {
63 pub name: String,
64 pub catalog_name: String,
65 pub comment: Option<String>,
66 pub storage_root: Option<String>,
67}
68
69pub type CatalogMeta = PersistedCatalogMeta;
70pub type NamespaceMeta = PersistedNamespaceMeta;
71
72#[derive(Debug, Clone, PartialEq, Eq, Hash)]
73pub struct TableFqn {
74 pub catalog: String,
75 pub namespace: String,
76 pub table: String,
77}
78
79#[derive(Debug, Clone, PartialEq, Eq, Hash)]
80pub struct IndexFqn {
81 pub catalog: String,
82 pub namespace: String,
83 pub table: String,
84 pub index: String,
85}
86
87impl TableFqn {
88 pub fn new(catalog: &str, namespace: &str, table: &str) -> Self {
89 Self {
90 catalog: catalog.to_string(),
91 namespace: namespace.to_string(),
92 table: table.to_string(),
93 }
94 }
95}
96
97impl IndexFqn {
98 pub fn new(catalog: &str, namespace: &str, table: &str, index: &str) -> Self {
99 Self {
100 catalog: catalog.to_string(),
101 namespace: namespace.to_string(),
102 table: table.to_string(),
103 index: index.to_string(),
104 }
105 }
106}
107
108impl From<&TableMetadata> for TableFqn {
109 fn from(value: &TableMetadata) -> Self {
110 Self::new(&value.catalog_name, &value.namespace_name, &value.name)
111 }
112}
113
114impl From<&IndexMetadata> for IndexFqn {
115 fn from(value: &IndexMetadata) -> Self {
116 Self::new(
117 &value.catalog_name,
118 &value.namespace_name,
119 &value.table,
120 &value.name,
121 )
122 }
123}
124
125#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
126#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
127pub enum TableType {
128 Managed,
129 External,
130}
131
132#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
133#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
134pub enum DataSourceFormat {
135 #[default]
136 Alopex,
137 Parquet,
138 Delta,
139}
140
141#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
142pub enum PersistedVectorMetric {
143 Cosine,
144 L2,
145 Inner,
146}
147
148impl From<VectorMetric> for PersistedVectorMetric {
149 fn from(value: VectorMetric) -> Self {
150 match value {
151 VectorMetric::Cosine => Self::Cosine,
152 VectorMetric::L2 => Self::L2,
153 VectorMetric::Inner => Self::Inner,
154 }
155 }
156}
157
158impl From<PersistedVectorMetric> for VectorMetric {
159 fn from(value: PersistedVectorMetric) -> Self {
160 match value {
161 PersistedVectorMetric::Cosine => Self::Cosine,
162 PersistedVectorMetric::L2 => Self::L2,
163 PersistedVectorMetric::Inner => Self::Inner,
164 }
165 }
166}
167
168#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
169pub enum PersistedType {
170 Integer,
171 BigInt,
172 Float,
173 Double,
174 Text,
175 Blob,
176 Boolean,
177 Timestamp,
178 Vector {
179 dimension: u32,
180 metric: PersistedVectorMetric,
181 },
182 Null,
183 Date,
184 Time,
185 Interval,
186 Decimal {
187 precision: u8,
188 scale: u8,
189 },
190 Json,
191 Array(Box<PersistedType>),
192 Map {
193 key: Box<PersistedType>,
194 value: Box<PersistedType>,
195 },
196 Struct(Vec<(String, PersistedType)>),
197}
198
199impl From<ResolvedType> for PersistedType {
200 fn from(value: ResolvedType) -> Self {
201 match value {
202 ResolvedType::Integer => Self::Integer,
203 ResolvedType::BigInt => Self::BigInt,
204 ResolvedType::Float => Self::Float,
205 ResolvedType::Double => Self::Double,
206 ResolvedType::Text => Self::Text,
207 ResolvedType::Blob => Self::Blob,
208 ResolvedType::Boolean => Self::Boolean,
209 ResolvedType::Timestamp => Self::Timestamp,
210 ResolvedType::Date => Self::Date,
211 ResolvedType::Time => Self::Time,
212 ResolvedType::Interval => Self::Interval,
213 ResolvedType::Decimal { precision, scale } => Self::Decimal { precision, scale },
214 ResolvedType::Json => Self::Json,
215 ResolvedType::Array(element) => Self::Array(Box::new((*element).into())),
216 ResolvedType::Map { key, value } => Self::Map {
217 key: Box::new((*key).into()),
218 value: Box::new((*value).into()),
219 },
220 ResolvedType::Struct(fields) => Self::Struct(
221 fields
222 .into_iter()
223 .map(|(name, data_type)| (name, data_type.into()))
224 .collect(),
225 ),
226 ResolvedType::Vector { dimension, metric } => Self::Vector {
227 dimension,
228 metric: metric.into(),
229 },
230 ResolvedType::Null => Self::Null,
231 }
232 }
233}
234
235impl From<PersistedType> for ResolvedType {
236 fn from(value: PersistedType) -> Self {
237 match value {
238 PersistedType::Integer => Self::Integer,
239 PersistedType::BigInt => Self::BigInt,
240 PersistedType::Float => Self::Float,
241 PersistedType::Double => Self::Double,
242 PersistedType::Text => Self::Text,
243 PersistedType::Blob => Self::Blob,
244 PersistedType::Boolean => Self::Boolean,
245 PersistedType::Timestamp => Self::Timestamp,
246 PersistedType::Date => Self::Date,
247 PersistedType::Time => Self::Time,
248 PersistedType::Interval => Self::Interval,
249 PersistedType::Decimal { precision, scale } => Self::Decimal { precision, scale },
250 PersistedType::Json => Self::Json,
251 PersistedType::Array(element) => Self::Array(Box::new((*element).into())),
252 PersistedType::Map { key, value } => Self::Map {
253 key: Box::new((*key).into()),
254 value: Box::new((*value).into()),
255 },
256 PersistedType::Struct(fields) => Self::Struct(
257 fields
258 .into_iter()
259 .map(|(name, data_type)| (name, data_type.into()))
260 .collect(),
261 ),
262 PersistedType::Vector { dimension, metric } => Self::Vector {
263 dimension,
264 metric: metric.into(),
265 },
266 PersistedType::Null => Self::Null,
267 }
268 }
269}
270
271#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
272pub enum PersistedIndexType {
273 BTree,
274 Hnsw,
275 Fts,
276}
277
278impl From<PersistedIndexType> for IndexMethod {
279 fn from(value: PersistedIndexType) -> Self {
280 match value {
281 PersistedIndexType::BTree => IndexMethod::BTree,
282 PersistedIndexType::Hnsw => IndexMethod::Hnsw,
283 PersistedIndexType::Fts => IndexMethod::Fts,
284 }
285 }
286}
287
288impl TryFrom<IndexMethod> for PersistedIndexType {
289 type Error = ();
290
291 fn try_from(value: IndexMethod) -> Result<Self, Self::Error> {
292 match value {
293 IndexMethod::BTree => Ok(Self::BTree),
294 IndexMethod::Hnsw => Ok(Self::Hnsw),
295 IndexMethod::Fts => Ok(Self::Fts),
296 }
297 }
298}
299
300#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
301pub enum PersistedStorageType {
302 Row,
303 Columnar,
304}
305
306impl From<PersistedStorageType> for StorageType {
307 fn from(value: PersistedStorageType) -> Self {
308 match value {
309 PersistedStorageType::Row => Self::Row,
310 PersistedStorageType::Columnar => Self::Columnar,
311 }
312 }
313}
314
315impl From<StorageType> for PersistedStorageType {
316 fn from(value: StorageType) -> Self {
317 match value {
318 StorageType::Row => Self::Row,
319 StorageType::Columnar => Self::Columnar,
320 }
321 }
322}
323
324#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
325pub enum PersistedCompression {
326 None,
327 Lz4,
328 Zstd,
329}
330
331impl From<PersistedCompression> for Compression {
332 fn from(value: PersistedCompression) -> Self {
333 match value {
334 PersistedCompression::None => Self::None,
335 PersistedCompression::Lz4 => Self::Lz4,
336 PersistedCompression::Zstd => Self::Zstd,
337 }
338 }
339}
340
341impl From<Compression> for PersistedCompression {
342 fn from(value: Compression) -> Self {
343 match value {
344 Compression::None => Self::None,
345 Compression::Lz4 => Self::Lz4,
346 Compression::Zstd => Self::Zstd,
347 }
348 }
349}
350
351#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
352pub enum PersistedRowIdMode {
353 None,
354 Direct,
355}
356
357impl From<PersistedRowIdMode> for RowIdMode {
358 fn from(value: PersistedRowIdMode) -> Self {
359 match value {
360 PersistedRowIdMode::None => Self::None,
361 PersistedRowIdMode::Direct => Self::Direct,
362 }
363 }
364}
365
366impl From<RowIdMode> for PersistedRowIdMode {
367 fn from(value: RowIdMode) -> Self {
368 match value {
369 RowIdMode::None => Self::None,
370 RowIdMode::Direct => Self::Direct,
371 }
372 }
373}
374
375#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
376pub struct PersistedStorageOptions {
377 pub storage_type: PersistedStorageType,
378 pub compression: PersistedCompression,
379 pub row_group_size: u32,
380 pub row_id_mode: PersistedRowIdMode,
381}
382
383impl From<StorageOptions> for PersistedStorageOptions {
384 fn from(value: StorageOptions) -> Self {
385 Self {
386 storage_type: value.storage_type.into(),
387 compression: value.compression.into(),
388 row_group_size: value.row_group_size,
389 row_id_mode: value.row_id_mode.into(),
390 }
391 }
392}
393
394impl From<PersistedStorageOptions> for StorageOptions {
395 fn from(value: PersistedStorageOptions) -> Self {
396 Self {
397 storage_type: value.storage_type.into(),
398 compression: value.compression.into(),
399 row_group_size: value.row_group_size,
400 row_id_mode: value.row_id_mode.into(),
401 }
402 }
403}
404
405#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
406pub struct PersistedColumnMeta {
407 pub name: String,
408 pub data_type: PersistedType,
409 pub not_null: bool,
410 pub primary_key: bool,
411 pub unique: bool,
412}
413
414impl From<&ColumnMetadata> for PersistedColumnMeta {
415 fn from(value: &ColumnMetadata) -> Self {
416 Self {
417 name: value.name.clone(),
418 data_type: value.data_type.clone().into(),
419 not_null: value.not_null,
420 primary_key: value.primary_key,
421 unique: value.unique,
422 }
423 }
424}
425
426impl From<PersistedColumnMeta> for ColumnMetadata {
427 fn from(value: PersistedColumnMeta) -> Self {
428 ColumnMetadata::new(value.name, value.data_type.into())
429 .with_not_null(value.not_null)
430 .with_primary_key(value.primary_key)
431 .with_unique(value.unique)
432 }
433}
434
435#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
436pub struct PersistedTableMeta {
437 pub table_id: u32,
438 pub name: String,
439 pub catalog_name: String,
440 pub namespace_name: String,
441 pub table_type: TableType,
442 pub data_source_format: DataSourceFormat,
443 pub columns: Vec<PersistedColumnMeta>,
444 pub primary_key: Option<Vec<String>>,
445 pub storage_options: PersistedStorageOptions,
446 pub storage_location: Option<String>,
447 pub comment: Option<String>,
448 pub properties: HashMap<String, String>,
449}
450
451#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
452struct PersistedTableMetaV1 {
453 table_id: u32,
454 name: String,
455 columns: Vec<PersistedColumnMeta>,
456 primary_key: Option<Vec<String>>,
457 storage_options: PersistedStorageOptions,
458}
459
460impl From<&TableMetadata> for PersistedTableMeta {
461 fn from(value: &TableMetadata) -> Self {
462 Self {
463 table_id: value.table_id,
464 name: value.name.clone(),
465 catalog_name: value.catalog_name.clone(),
466 namespace_name: value.namespace_name.clone(),
467 table_type: value.table_type,
468 data_source_format: value.data_source_format,
469 columns: value
470 .columns
471 .iter()
472 .map(PersistedColumnMeta::from)
473 .collect(),
474 primary_key: value.primary_key.clone(),
475 storage_options: value.storage_options.clone().into(),
476 storage_location: value.storage_location.clone(),
477 comment: value.comment.clone(),
478 properties: value.properties.clone(),
479 }
480 }
481}
482
483impl From<PersistedTableMeta> for TableMetadata {
484 fn from(value: PersistedTableMeta) -> Self {
485 let mut table = TableMetadata::new(
486 value.name,
487 value
488 .columns
489 .into_iter()
490 .map(ColumnMetadata::from)
491 .collect(),
492 )
493 .with_table_id(value.table_id);
494 table.primary_key = value.primary_key;
495 table.storage_options = value.storage_options.into();
496 table.catalog_name = value.catalog_name;
497 table.namespace_name = value.namespace_name;
498 table.table_type = value.table_type;
499 table.data_source_format = value.data_source_format;
500 table.storage_location = value.storage_location;
501 table.comment = value.comment;
502 table.properties = value.properties;
503 table
504 }
505}
506
507#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
508pub struct PersistedIndexMeta {
509 pub index_id: u32,
510 pub name: String,
511 pub table: String,
512 pub columns: Vec<String>,
513 pub column_indices: Vec<usize>,
514 pub unique: bool,
515 pub method: Option<PersistedIndexType>,
516 pub options: Vec<(String, String)>,
517 pub catalog_name: String,
518 pub namespace_name: String,
519}
520
521#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
522struct PersistedIndexMetaV1 {
523 index_id: u32,
524 name: String,
525 table: String,
526 columns: Vec<String>,
527 column_indices: Vec<usize>,
528 unique: bool,
529 method: Option<PersistedIndexType>,
530 options: Vec<(String, String)>,
531}
532
533impl From<&IndexMetadata> for PersistedIndexMeta {
534 fn from(value: &IndexMetadata) -> Self {
535 Self {
536 index_id: value.index_id,
537 name: value.name.clone(),
538 table: value.table.clone(),
539 columns: value.columns.clone(),
540 column_indices: value.column_indices.clone(),
541 unique: value.unique,
542 method: value
543 .method
544 .and_then(|m| PersistedIndexType::try_from(m).ok()),
545 options: value.options.clone(),
546 catalog_name: value.catalog_name.clone(),
547 namespace_name: value.namespace_name.clone(),
548 }
549 }
550}
551
552impl From<PersistedIndexMeta> for IndexMetadata {
553 fn from(value: PersistedIndexMeta) -> Self {
554 let mut index = IndexMetadata::new(value.index_id, value.name, value.table, value.columns)
555 .with_column_indices(value.column_indices)
556 .with_unique(value.unique)
557 .with_options(value.options);
558 index.catalog_name = value.catalog_name;
559 index.namespace_name = value.namespace_name;
560 if let Some(method) = value.method {
561 index = index.with_method(method.into());
562 }
563 index
564 }
565}
566
567fn deserialize_table_meta(bytes: &[u8]) -> Result<PersistedTableMeta, CatalogError> {
568 match bincode::deserialize::<PersistedTableMeta>(bytes) {
569 Ok(meta) => Ok(meta),
570 Err(err) => {
571 let is_legacy = matches!(
572 err.as_ref(),
573 bincode::ErrorKind::Io(io)
574 if io.kind() == std::io::ErrorKind::UnexpectedEof
575 );
576 if !is_legacy {
577 return Err(err.into());
578 }
579 let legacy: PersistedTableMetaV1 = bincode::deserialize(bytes)?;
580 Ok(PersistedTableMeta {
581 table_id: legacy.table_id,
582 name: legacy.name,
583 catalog_name: "default".to_string(),
584 namespace_name: "default".to_string(),
585 table_type: TableType::Managed,
586 data_source_format: DataSourceFormat::Alopex,
587 columns: legacy.columns,
588 primary_key: legacy.primary_key,
589 storage_options: legacy.storage_options,
590 storage_location: None,
591 comment: None,
592 properties: HashMap::new(),
593 })
594 }
595 }
596}
597
598fn deserialize_index_meta(bytes: &[u8]) -> Result<PersistedIndexMeta, CatalogError> {
599 match bincode::deserialize::<PersistedIndexMeta>(bytes) {
600 Ok(meta) => Ok(meta),
601 Err(err) => {
602 let is_legacy = matches!(
603 err.as_ref(),
604 bincode::ErrorKind::Io(io)
605 if io.kind() == std::io::ErrorKind::UnexpectedEof
606 );
607 if !is_legacy {
608 return Err(err.into());
609 }
610 let legacy: PersistedIndexMetaV1 = bincode::deserialize(bytes)?;
611 Ok(PersistedIndexMeta {
612 index_id: legacy.index_id,
613 name: legacy.name,
614 table: legacy.table,
615 columns: legacy.columns,
616 column_indices: legacy.column_indices,
617 unique: legacy.unique,
618 method: legacy.method,
619 options: legacy.options,
620 catalog_name: "default".to_string(),
621 namespace_name: "default".to_string(),
622 })
623 }
624 }
625}
626
627fn table_key(catalog_name: &str, namespace_name: &str, table_name: &str) -> Vec<u8> {
628 let mut key = TABLES_PREFIX.to_vec();
629 key.extend_from_slice(catalog_name.as_bytes());
630 key.push(b'/');
631 key.extend_from_slice(namespace_name.as_bytes());
632 key.push(b'/');
633 key.extend_from_slice(table_name.as_bytes());
634 key
635}
636
637fn catalog_key(name: &str) -> Vec<u8> {
638 let mut key = CATALOGS_PREFIX.to_vec();
639 key.extend_from_slice(name.as_bytes());
640 key
641}
642
643fn namespace_key(catalog_name: &str, namespace_name: &str) -> Vec<u8> {
644 let mut key = NAMESPACES_PREFIX.to_vec();
645 key.extend_from_slice(catalog_name.as_bytes());
646 key.push(b'/');
647 key.extend_from_slice(namespace_name.as_bytes());
648 key
649}
650
651fn index_key(
652 catalog_name: &str,
653 namespace_name: &str,
654 table_name: &str,
655 index_name: &str,
656) -> Vec<u8> {
657 let mut key = INDEXES_PREFIX.to_vec();
658 key.extend_from_slice(catalog_name.as_bytes());
659 key.push(b'/');
660 key.extend_from_slice(namespace_name.as_bytes());
661 key.push(b'/');
662 key.extend_from_slice(table_name.as_bytes());
663 key.push(b'/');
664 key.extend_from_slice(index_name.as_bytes());
665 key
666}
667
668fn index_prefix(catalog_name: &str, namespace_name: &str, table_name: &str) -> Vec<u8> {
669 let mut key = INDEXES_PREFIX.to_vec();
670 key.extend_from_slice(catalog_name.as_bytes());
671 key.push(b'/');
672 key.extend_from_slice(namespace_name.as_bytes());
673 key.push(b'/');
674 key.extend_from_slice(table_name.as_bytes());
675 key.push(b'/');
676 key
677}
678
679fn key_suffix(prefix: &[u8], key: &[u8]) -> Result<String, CatalogError> {
680 let suffix = key
681 .strip_prefix(prefix)
682 .ok_or_else(|| CatalogError::InvalidKey(format!("{key:?}")))?;
683 std::str::from_utf8(suffix)
684 .map(|s| s.to_string())
685 .map_err(|_| CatalogError::InvalidKey(format!("{key:?}")))
686}
687
688fn parse_table_key_suffix(suffix: &str) -> Result<TableFqn, CatalogError> {
689 let mut parts = suffix.splitn(3, '/');
690 let catalog = parts
691 .next()
692 .filter(|part| !part.is_empty())
693 .ok_or_else(|| CatalogError::InvalidKey(suffix.to_string()))?;
694 let namespace = parts
695 .next()
696 .filter(|part| !part.is_empty())
697 .ok_or_else(|| CatalogError::InvalidKey(suffix.to_string()))?;
698 let table = parts
699 .next()
700 .filter(|part| !part.is_empty())
701 .ok_or_else(|| CatalogError::InvalidKey(suffix.to_string()))?;
702 Ok(TableFqn::new(catalog, namespace, table))
703}
704
705fn parse_index_key_suffix(suffix: &str) -> Result<IndexFqn, CatalogError> {
706 let mut parts = suffix.splitn(4, '/');
707 let catalog = parts
708 .next()
709 .filter(|part| !part.is_empty())
710 .ok_or_else(|| CatalogError::InvalidKey(suffix.to_string()))?;
711 let namespace = parts
712 .next()
713 .filter(|part| !part.is_empty())
714 .ok_or_else(|| CatalogError::InvalidKey(suffix.to_string()))?;
715 let table = parts
716 .next()
717 .filter(|part| !part.is_empty())
718 .ok_or_else(|| CatalogError::InvalidKey(suffix.to_string()))?;
719 let index = parts
720 .next()
721 .filter(|part| !part.is_empty())
722 .ok_or_else(|| CatalogError::InvalidKey(suffix.to_string()))?;
723 Ok(IndexFqn::new(catalog, namespace, table, index))
724}
725
726#[derive(Debug, Clone, Default)]
727pub struct CatalogOverlay {
728 added_catalogs: HashMap<String, CatalogMeta>,
729 dropped_catalogs: HashSet<String>,
730 added_namespaces: HashMap<(String, String), NamespaceMeta>,
731 dropped_namespaces: HashSet<(String, String)>,
732 added_tables: HashMap<TableFqn, TableMetadata>,
733 dropped_tables: HashSet<TableFqn>,
734 added_indexes: HashMap<IndexFqn, IndexMetadata>,
735 dropped_indexes: HashSet<IndexFqn>,
736}
737
738impl CatalogOverlay {
739 pub fn new() -> Self {
740 Self::default()
741 }
742
743 pub fn add_catalog(&mut self, meta: CatalogMeta) {
744 self.dropped_catalogs.remove(&meta.name);
745 self.added_catalogs.insert(meta.name.clone(), meta);
746 }
747
748 pub fn drop_catalog(&mut self, name: &str) {
749 self.added_catalogs.remove(name);
750 self.dropped_catalogs.insert(name.to_string());
751 }
752
753 pub fn add_namespace(&mut self, meta: NamespaceMeta) {
754 let key = (meta.catalog_name.clone(), meta.name.clone());
755 self.dropped_namespaces.remove(&key);
756 self.added_namespaces.insert(key, meta);
757 }
758
759 pub fn drop_namespace(&mut self, catalog_name: &str, namespace_name: &str) {
760 let key = (catalog_name.to_string(), namespace_name.to_string());
761 self.added_namespaces.remove(&key);
762 self.dropped_namespaces.insert(key);
763 }
764
765 pub fn add_table(&mut self, fqn: TableFqn, table: TableMetadata) {
766 self.dropped_tables.remove(&fqn);
767 self.added_tables.insert(fqn, table);
768 }
769
770 pub fn drop_table(&mut self, fqn: &TableFqn) {
771 self.added_tables.remove(fqn);
772 self.dropped_tables.insert(fqn.clone());
773 self.added_indexes.retain(|key, _| {
774 key.catalog != fqn.catalog || key.namespace != fqn.namespace || key.table != fqn.table
775 });
776 }
777
778 pub fn add_index(&mut self, fqn: IndexFqn, index: IndexMetadata) {
779 self.dropped_indexes.remove(&fqn);
780 self.added_indexes.insert(fqn, index);
781 }
782
783 pub fn drop_index(&mut self, fqn: &IndexFqn) {
784 self.added_indexes.remove(fqn);
785 self.dropped_indexes.insert(fqn.clone());
786 }
787
788 pub fn drop_cascade_catalog(&mut self, catalog: &str) {
789 self.drop_catalog(catalog);
790
791 let namespace_keys: Vec<(String, String)> = self
792 .added_namespaces
793 .keys()
794 .filter(|(cat, _)| cat == catalog)
795 .cloned()
796 .collect();
797 for (catalog_name, namespace_name) in namespace_keys {
798 self.drop_namespace(&catalog_name, &namespace_name);
799 }
800
801 let index_keys: Vec<IndexFqn> = self
802 .added_indexes
803 .keys()
804 .filter(|fqn| fqn.catalog == catalog)
805 .cloned()
806 .collect();
807 for fqn in index_keys {
808 self.drop_index(&fqn);
809 }
810
811 let table_keys: Vec<TableFqn> = self
812 .added_tables
813 .keys()
814 .filter(|fqn| fqn.catalog == catalog)
815 .cloned()
816 .collect();
817 for fqn in table_keys {
818 self.drop_table(&fqn);
819 }
820 }
821
822 pub fn drop_cascade_namespace(&mut self, catalog: &str, namespace: &str) {
823 self.drop_namespace(catalog, namespace);
824
825 let index_keys: Vec<IndexFqn> = self
826 .added_indexes
827 .keys()
828 .filter(|fqn| fqn.catalog == catalog && fqn.namespace == namespace)
829 .cloned()
830 .collect();
831 for fqn in index_keys {
832 self.drop_index(&fqn);
833 }
834
835 let table_keys: Vec<TableFqn> = self
836 .added_tables
837 .keys()
838 .filter(|fqn| fqn.catalog == catalog && fqn.namespace == namespace)
839 .cloned()
840 .collect();
841 for fqn in table_keys {
842 self.drop_table(&fqn);
843 }
844 }
845}
846
847pub struct TxnCatalogView<'a, S: KVStore> {
852 catalog: &'a PersistentCatalog<S>,
853 overlay: &'a CatalogOverlay,
854}
855
856impl<'a, S: KVStore> TxnCatalogView<'a, S> {
857 pub fn new(catalog: &'a PersistentCatalog<S>, overlay: &'a CatalogOverlay) -> Self {
858 Self { catalog, overlay }
859 }
860}
861
862impl<'a, S: KVStore> Catalog for TxnCatalogView<'a, S> {
863 fn create_table(&mut self, _table: TableMetadata) -> Result<(), PlannerError> {
864 unreachable!("TxnCatalogView は参照専用です")
865 }
866
867 fn get_table(&self, name: &str) -> Option<&TableMetadata> {
868 self.catalog.get_table_in_txn(name, self.overlay)
869 }
870
871 fn drop_table(&mut self, _name: &str) -> Result<(), PlannerError> {
872 unreachable!("TxnCatalogView は参照専用です")
873 }
874
875 fn create_index(&mut self, _index: IndexMetadata) -> Result<(), PlannerError> {
876 unreachable!("TxnCatalogView は参照専用です")
877 }
878
879 fn get_index(&self, name: &str) -> Option<&IndexMetadata> {
880 self.catalog.get_index_in_txn(name, self.overlay)
881 }
882
883 fn get_indexes_for_table(&self, table: &str) -> Vec<&IndexMetadata> {
884 let Some(table_meta) = self.catalog.get_table_in_txn(table, self.overlay) else {
885 return Vec::new();
886 };
887
888 let mut indexes: Vec<&IndexMetadata> = self
889 .catalog
890 .inner
891 .get_indexes_for_table(table)
892 .into_iter()
893 .filter(|idx| {
894 idx.catalog_name == table_meta.catalog_name
895 && idx.namespace_name == table_meta.namespace_name
896 && !self.catalog.index_hidden_by_overlay(idx, self.overlay)
897 })
898 .collect();
899
900 for idx in self.overlay.added_indexes.values() {
901 if idx.table == table
902 && idx.catalog_name == table_meta.catalog_name
903 && idx.namespace_name == table_meta.namespace_name
904 && !self.catalog.index_hidden_by_overlay(idx, self.overlay)
905 {
906 indexes.push(idx);
907 }
908 }
909
910 indexes
911 }
912
913 fn drop_index(&mut self, _name: &str) -> Result<(), PlannerError> {
914 unreachable!("TxnCatalogView は参照専用です")
915 }
916
917 fn table_exists(&self, name: &str) -> bool {
918 self.catalog.table_exists_in_txn(name, self.overlay)
919 }
920
921 fn index_exists(&self, name: &str) -> bool {
922 self.catalog.index_exists_in_txn(name, self.overlay)
923 }
924
925 fn next_table_id(&mut self) -> u32 {
926 unreachable!("TxnCatalogView は参照専用です")
927 }
928
929 fn next_index_id(&mut self) -> u32 {
930 unreachable!("TxnCatalogView は参照専用です")
931 }
932
933 fn list_tables(&self) -> Vec<TableMetadata> {
934 let mut names = HashSet::new();
935 for name in self.catalog.inner.table_names() {
936 names.insert(name.to_string());
937 }
938 for fqn in self.overlay.added_tables.keys() {
939 names.insert(fqn.table.clone());
940 }
941
942 let mut tables = Vec::new();
943 for name in names {
944 if let Some(table) = self.catalog.get_table_in_txn(&name, self.overlay) {
945 tables.push(table.clone());
946 }
947 }
948 tables
949 }
950}
951
952#[derive(Debug)]
954pub struct PersistentCatalog<S: KVStore> {
969 inner: MemoryCatalog,
970 store: Arc<S>,
971 catalogs: HashMap<String, CatalogMeta>,
972 namespaces: HashMap<(String, String), NamespaceMeta>,
973}
974
975impl<S: KVStore> PersistentCatalog<S> {
976 pub fn load(store: Arc<S>) -> Result<Self, CatalogError> {
977 let mut txn = store.begin(TxnMode::ReadOnly)?;
978 let meta_key = META_KEY.to_vec();
979 let mut meta_state: Option<CatalogState> = None;
980
981 if let Some(meta_bytes) = txn.get(&meta_key)? {
982 let meta: CatalogState = bincode::deserialize(&meta_bytes)?;
983 if meta.version > CATALOG_VERSION {
984 return Err(CatalogError::InvalidKey(format!(
985 "unsupported catalog version: {}",
986 meta.version
987 )));
988 }
989 meta_state = Some(meta);
990 }
991
992 let mut needs_migration = meta_state
993 .as_ref()
994 .is_some_and(|meta| meta.version < CATALOG_VERSION);
995 if !needs_migration && meta_state.is_none() {
996 for (key, _) in txn.scan_prefix(TABLES_PREFIX)? {
997 let suffix = key_suffix(TABLES_PREFIX, &key)?;
998 if !suffix.contains('/') {
999 needs_migration = true;
1000 break;
1001 }
1002 }
1003 if !needs_migration {
1004 for (key, _) in txn.scan_prefix(INDEXES_PREFIX)? {
1005 let suffix = key_suffix(INDEXES_PREFIX, &key)?;
1006 if !suffix.contains('/') {
1007 needs_migration = true;
1008 break;
1009 }
1010 }
1011 }
1012 }
1013
1014 if needs_migration {
1015 txn.rollback_self()?;
1016 Self::migrate_v1_to_v2(&store)?;
1017 return Self::load(store);
1018 }
1019
1020 let mut inner = MemoryCatalog::new();
1021 let mut catalogs = HashMap::new();
1022 let mut namespaces = HashMap::new();
1023
1024 let mut max_table_id = 0u32;
1025 let mut max_index_id = 0u32;
1026
1027 for (key, value) in txn.scan_prefix(CATALOGS_PREFIX)? {
1028 let catalog_name = key_suffix(CATALOGS_PREFIX, &key)?;
1029 let mut meta: CatalogMeta = bincode::deserialize(&value)?;
1030 if meta.name != catalog_name {
1031 meta.name = catalog_name.clone();
1032 }
1033 catalogs.insert(catalog_name, meta);
1034 }
1035
1036 for (key, value) in txn.scan_prefix(NAMESPACES_PREFIX)? {
1037 let suffix = key_suffix(NAMESPACES_PREFIX, &key)?;
1038 let mut parts = suffix.splitn(2, '/');
1039 let catalog_name = parts
1040 .next()
1041 .filter(|part| !part.is_empty())
1042 .ok_or_else(|| CatalogError::InvalidKey(suffix.clone()))?;
1043 let namespace_name = parts
1044 .next()
1045 .filter(|part| !part.is_empty())
1046 .ok_or_else(|| CatalogError::InvalidKey(suffix.clone()))?;
1047 let mut meta: NamespaceMeta = bincode::deserialize(&value)?;
1048 if meta.catalog_name != catalog_name {
1049 meta.catalog_name = catalog_name.to_string();
1050 }
1051 if meta.name != namespace_name {
1052 meta.name = namespace_name.to_string();
1053 }
1054 namespaces.insert((meta.catalog_name.clone(), meta.name.clone()), meta);
1055 }
1056
1057 for (key, value) in txn.scan_prefix(TABLES_PREFIX)? {
1059 let suffix = key_suffix(TABLES_PREFIX, &key)?;
1060 let fqn = parse_table_key_suffix(&suffix)?;
1061 let mut persisted = deserialize_table_meta(&value)?;
1062 if persisted.catalog_name != fqn.catalog {
1063 persisted.catalog_name = fqn.catalog.clone();
1064 }
1065 if persisted.namespace_name != fqn.namespace {
1066 persisted.namespace_name = fqn.namespace.clone();
1067 }
1068 if persisted.name != fqn.table {
1069 persisted.name = fqn.table.clone();
1070 }
1071 max_table_id = max_table_id.max(persisted.table_id);
1072 let table: TableMetadata = persisted.into();
1073 inner.insert_table_unchecked(table);
1074 }
1075
1076 for (key, value) in txn.scan_prefix(INDEXES_PREFIX)? {
1077 let suffix = key_suffix(INDEXES_PREFIX, &key)?;
1078 let fqn = parse_index_key_suffix(&suffix)?;
1079 let mut persisted = deserialize_index_meta(&value)?;
1080 if persisted.catalog_name != fqn.catalog {
1081 persisted.catalog_name = fqn.catalog.clone();
1082 }
1083 if persisted.namespace_name != fqn.namespace {
1084 persisted.namespace_name = fqn.namespace.clone();
1085 }
1086 if persisted.table != fqn.table {
1087 persisted.table = fqn.table.clone();
1088 }
1089 if persisted.name != fqn.index {
1090 persisted.name = fqn.index.clone();
1091 }
1092 max_index_id = max_index_id.max(persisted.index_id);
1093 let mut index: IndexMetadata = persisted.into();
1094 if let Some(table) = inner.get_table(&index.table) {
1096 if index.catalog_name != table.catalog_name
1097 || index.namespace_name != table.namespace_name
1098 {
1099 index.catalog_name = table.catalog_name.clone();
1100 index.namespace_name = table.namespace_name.clone();
1101 }
1102 inner.insert_index_unchecked(index);
1103 }
1104 }
1105
1106 let (mut table_id_counter, mut index_id_counter) = (max_table_id, max_index_id);
1107 if let Some(meta) = meta_state
1108 .as_ref()
1109 .filter(|meta| meta.version == CATALOG_VERSION)
1110 {
1111 table_id_counter = table_id_counter.max(meta.table_id_counter);
1112 index_id_counter = index_id_counter.max(meta.index_id_counter);
1113 }
1114 inner.set_counters(table_id_counter, index_id_counter);
1115
1116 txn.rollback_self()?;
1117
1118 Ok(Self {
1119 inner,
1120 store,
1121 catalogs,
1122 namespaces,
1123 })
1124 }
1125
1126 fn migrate_v1_to_v2(store: &Arc<S>) -> Result<(), CatalogError> {
1127 let mut txn = store.begin(TxnMode::ReadWrite)?;
1128
1129 if txn.get(&catalog_key("default"))?.is_none() {
1130 let meta = CatalogMeta {
1131 name: "default".to_string(),
1132 comment: None,
1133 storage_root: None,
1134 };
1135 let value = bincode::serialize(&meta)?;
1136 txn.put(catalog_key("default"), value)?;
1137 }
1138
1139 if txn.get(&namespace_key("default", "default"))?.is_none() {
1140 let meta = NamespaceMeta {
1141 name: "default".to_string(),
1142 catalog_name: "default".to_string(),
1143 comment: None,
1144 storage_root: None,
1145 };
1146 let value = bincode::serialize(&meta)?;
1147 txn.put(namespace_key("default", "default"), value)?;
1148 }
1149
1150 let mut table_updates = Vec::new();
1151 let mut table_keys_to_delete = Vec::new();
1152 let mut max_table_id = 0u32;
1153 for (key, value) in txn.scan_prefix(TABLES_PREFIX)? {
1154 let suffix = key_suffix(TABLES_PREFIX, &key)?;
1155 if suffix.contains('/') {
1156 continue;
1157 }
1158 let mut persisted = deserialize_table_meta(&value)?;
1159 if persisted.catalog_name.is_empty() {
1160 persisted.catalog_name = "default".to_string();
1161 }
1162 if persisted.namespace_name.is_empty() {
1163 persisted.namespace_name = "default".to_string();
1164 }
1165 persisted.table_type = TableType::Managed;
1166 persisted.data_source_format = DataSourceFormat::Alopex;
1167 max_table_id = max_table_id.max(persisted.table_id);
1168
1169 let new_key = table_key(
1170 &persisted.catalog_name,
1171 &persisted.namespace_name,
1172 &persisted.name,
1173 );
1174 let bytes = bincode::serialize(&persisted)?;
1175 table_updates.push((new_key, bytes));
1176 table_keys_to_delete.push(key);
1177 }
1178
1179 for (key, value) in table_updates {
1180 txn.put(key, value)?;
1181 }
1182 for key in table_keys_to_delete {
1183 txn.delete(key)?;
1184 }
1185
1186 let mut index_updates = Vec::new();
1187 let mut index_keys_to_delete = Vec::new();
1188 let mut max_index_id = 0u32;
1189 for (key, value) in txn.scan_prefix(INDEXES_PREFIX)? {
1190 let suffix = key_suffix(INDEXES_PREFIX, &key)?;
1191 if suffix.contains('/') {
1192 continue;
1193 }
1194 let mut persisted = deserialize_index_meta(&value)?;
1195 if persisted.catalog_name.is_empty() {
1196 persisted.catalog_name = "default".to_string();
1197 }
1198 if persisted.namespace_name.is_empty() {
1199 persisted.namespace_name = "default".to_string();
1200 }
1201 max_index_id = max_index_id.max(persisted.index_id);
1202
1203 let new_key = index_key(
1204 &persisted.catalog_name,
1205 &persisted.namespace_name,
1206 &persisted.table,
1207 &persisted.name,
1208 );
1209 let bytes = bincode::serialize(&persisted)?;
1210 index_updates.push((new_key, bytes));
1211 index_keys_to_delete.push(key);
1212 }
1213
1214 for (key, value) in index_updates {
1215 txn.put(key, value)?;
1216 }
1217 for key in index_keys_to_delete {
1218 txn.delete(key)?;
1219 }
1220
1221 let mut table_id_counter = max_table_id;
1222 let mut index_id_counter = max_index_id;
1223 if let Some(meta_bytes) = txn.get(&META_KEY.to_vec())? {
1224 let meta: CatalogState = bincode::deserialize(&meta_bytes)?;
1225 table_id_counter = table_id_counter.max(meta.table_id_counter);
1226 index_id_counter = index_id_counter.max(meta.index_id_counter);
1227 }
1228 let meta = CatalogState {
1229 version: CATALOG_VERSION,
1230 table_id_counter,
1231 index_id_counter,
1232 };
1233 let meta_bytes = bincode::serialize(&meta)?;
1234 txn.put(META_KEY.to_vec(), meta_bytes)?;
1235 txn.commit_self()?;
1236
1237 Ok(())
1238 }
1239
1240 pub fn new(store: Arc<S>) -> Self {
1241 Self {
1242 inner: MemoryCatalog::new(),
1243 store,
1244 catalogs: HashMap::new(),
1245 namespaces: HashMap::new(),
1246 }
1247 }
1248
1249 pub fn store(&self) -> &Arc<S> {
1250 &self.store
1251 }
1252
1253 pub fn list_catalogs(&self) -> Vec<CatalogMeta> {
1254 let mut catalogs: Vec<CatalogMeta> = self.catalogs.values().cloned().collect();
1255 catalogs.sort_by(|a, b| a.name.cmp(&b.name));
1256 catalogs
1257 }
1258
1259 pub fn get_catalog(&self, name: &str) -> Option<CatalogMeta> {
1260 self.catalogs.get(name).cloned()
1261 }
1262
1263 pub fn create_catalog(&mut self, meta: CatalogMeta) -> Result<(), CatalogError> {
1264 let mut txn = self.store.begin(TxnMode::ReadWrite)?;
1265 let value = bincode::serialize(&meta)?;
1266 txn.put(catalog_key(&meta.name), value)?;
1267 txn.commit_self()?;
1268 self.catalogs.insert(meta.name.clone(), meta);
1269 Ok(())
1270 }
1271
1272 pub fn delete_catalog(&mut self, name: &str) -> Result<(), CatalogError> {
1273 let mut txn = self.store.begin(TxnMode::ReadWrite)?;
1274 txn.delete(catalog_key(name))?;
1275 let mut namespace_prefix = NAMESPACES_PREFIX.to_vec();
1276 namespace_prefix.extend_from_slice(name.as_bytes());
1277 namespace_prefix.push(b'/');
1278 let mut namespace_keys = Vec::new();
1279 for (key, _) in txn.scan_prefix(&namespace_prefix)? {
1280 namespace_keys.push(key);
1281 }
1282 for key in namespace_keys {
1283 txn.delete(key)?;
1284 }
1285 let mut table_keys = Vec::new();
1286 let mut table_fqns = Vec::new();
1287 for (key, value) in txn.scan_prefix(TABLES_PREFIX)? {
1288 let persisted = deserialize_table_meta(&value)?;
1289 if persisted.catalog_name == name {
1290 table_fqns.push(TableFqn::new(
1291 &persisted.catalog_name,
1292 &persisted.namespace_name,
1293 &persisted.name,
1294 ));
1295 table_keys.push(key);
1296 }
1297 }
1298 let table_set: HashSet<TableFqn> = table_fqns.iter().cloned().collect();
1299 for key in table_keys {
1300 txn.delete(key)?;
1301 }
1302 if !table_set.is_empty() {
1303 let mut index_keys = Vec::new();
1304 for (key, value) in txn.scan_prefix(INDEXES_PREFIX)? {
1305 let persisted = deserialize_index_meta(&value)?;
1306 let fqn = TableFqn::new(
1307 &persisted.catalog_name,
1308 &persisted.namespace_name,
1309 &persisted.table,
1310 );
1311 if table_set.contains(&fqn) {
1312 index_keys.push(key);
1313 }
1314 }
1315 for key in index_keys {
1316 txn.delete(key)?;
1317 }
1318 }
1319 txn.commit_self()?;
1320 self.catalogs.remove(name);
1321 self.namespaces.retain(|(catalog, _), _| catalog != name);
1322 for fqn in table_fqns {
1323 self.inner.remove_table_unchecked(&fqn.table);
1324 }
1325 Ok(())
1326 }
1327
1328 pub fn list_namespaces(&self, catalog_name: &str) -> Vec<NamespaceMeta> {
1329 let mut namespaces: Vec<NamespaceMeta> = self
1330 .namespaces
1331 .values()
1332 .filter(|meta| meta.catalog_name == catalog_name)
1333 .cloned()
1334 .collect();
1335 namespaces.sort_by(|a, b| a.name.cmp(&b.name));
1336 namespaces
1337 }
1338
1339 pub fn get_namespace(&self, catalog_name: &str, namespace_name: &str) -> Option<NamespaceMeta> {
1340 self.namespaces
1341 .get(&(catalog_name.to_string(), namespace_name.to_string()))
1342 .cloned()
1343 }
1344
1345 pub fn create_namespace(&mut self, meta: NamespaceMeta) -> Result<(), CatalogError> {
1346 if !self.catalogs.contains_key(&meta.catalog_name) {
1347 return Err(CatalogError::InvalidKey(format!(
1348 "catalog not found: {}",
1349 meta.catalog_name
1350 )));
1351 }
1352
1353 let mut txn = self.store.begin(TxnMode::ReadWrite)?;
1354 let value = bincode::serialize(&meta)?;
1355 txn.put(namespace_key(&meta.catalog_name, &meta.name), value)?;
1356 txn.commit_self()?;
1357 self.namespaces
1358 .insert((meta.catalog_name.clone(), meta.name.clone()), meta);
1359 Ok(())
1360 }
1361
1362 pub fn delete_namespace(
1363 &mut self,
1364 catalog_name: &str,
1365 namespace_name: &str,
1366 ) -> Result<(), CatalogError> {
1367 if !self.catalogs.contains_key(catalog_name) {
1368 return Err(CatalogError::InvalidKey(format!(
1369 "catalog not found: {}",
1370 catalog_name
1371 )));
1372 }
1373
1374 let mut txn = self.store.begin(TxnMode::ReadWrite)?;
1375 txn.delete(namespace_key(catalog_name, namespace_name))?;
1376 txn.commit_self()?;
1377 self.namespaces
1378 .remove(&(catalog_name.to_string(), namespace_name.to_string()));
1379 Ok(())
1380 }
1381
1382 fn persist_create_catalog(
1383 &mut self,
1384 txn: &mut S::Transaction<'_>,
1385 meta: &CatalogMeta,
1386 ) -> Result<(), CatalogError> {
1387 let value = bincode::serialize(meta)?;
1388 txn.put(catalog_key(&meta.name), value)?;
1389 Ok(())
1390 }
1391
1392 fn persist_drop_catalog(
1393 &mut self,
1394 txn: &mut S::Transaction<'_>,
1395 name: &str,
1396 ) -> Result<(), CatalogError> {
1397 txn.delete(catalog_key(name))?;
1398
1399 let mut namespace_prefix = NAMESPACES_PREFIX.to_vec();
1400 namespace_prefix.extend_from_slice(name.as_bytes());
1401 namespace_prefix.push(b'/');
1402 let mut namespace_keys = Vec::new();
1403 for (key, _) in txn.scan_prefix(&namespace_prefix)? {
1404 namespace_keys.push(key);
1405 }
1406 for key in namespace_keys {
1407 txn.delete(key)?;
1408 }
1409
1410 let mut table_keys = Vec::new();
1411 let mut table_fqns = Vec::new();
1412 for (key, value) in txn.scan_prefix(TABLES_PREFIX)? {
1413 let persisted = deserialize_table_meta(&value)?;
1414 if persisted.catalog_name == name {
1415 table_fqns.push(TableFqn::new(
1416 &persisted.catalog_name,
1417 &persisted.namespace_name,
1418 &persisted.name,
1419 ));
1420 table_keys.push(key);
1421 }
1422 }
1423 let table_set: HashSet<TableFqn> = table_fqns.iter().cloned().collect();
1424 for key in table_keys {
1425 txn.delete(key)?;
1426 }
1427 if !table_set.is_empty() {
1428 let mut index_keys = Vec::new();
1429 for (key, value) in txn.scan_prefix(INDEXES_PREFIX)? {
1430 let persisted = deserialize_index_meta(&value)?;
1431 let fqn = TableFqn::new(
1432 &persisted.catalog_name,
1433 &persisted.namespace_name,
1434 &persisted.table,
1435 );
1436 if table_set.contains(&fqn) {
1437 index_keys.push(key);
1438 }
1439 }
1440 for key in index_keys {
1441 txn.delete(key)?;
1442 }
1443 }
1444 Ok(())
1445 }
1446
1447 fn persist_create_namespace(
1448 &mut self,
1449 txn: &mut S::Transaction<'_>,
1450 meta: &NamespaceMeta,
1451 ) -> Result<(), CatalogError> {
1452 let value = bincode::serialize(meta)?;
1453 txn.put(namespace_key(&meta.catalog_name, &meta.name), value)?;
1454 Ok(())
1455 }
1456
1457 fn persist_drop_namespace(
1458 &mut self,
1459 txn: &mut S::Transaction<'_>,
1460 catalog_name: &str,
1461 namespace_name: &str,
1462 ) -> Result<(), CatalogError> {
1463 txn.delete(namespace_key(catalog_name, namespace_name))?;
1464
1465 let mut table_keys = Vec::new();
1466 let mut table_fqns = Vec::new();
1467 for (key, value) in txn.scan_prefix(TABLES_PREFIX)? {
1468 let persisted = deserialize_table_meta(&value)?;
1469 if persisted.catalog_name == catalog_name && persisted.namespace_name == namespace_name
1470 {
1471 table_fqns.push(TableFqn::new(
1472 &persisted.catalog_name,
1473 &persisted.namespace_name,
1474 &persisted.name,
1475 ));
1476 table_keys.push(key);
1477 }
1478 }
1479 let table_set: HashSet<TableFqn> = table_fqns.iter().cloned().collect();
1480 for key in table_keys {
1481 txn.delete(key)?;
1482 }
1483 if !table_set.is_empty() {
1484 let mut index_keys = Vec::new();
1485 for (key, value) in txn.scan_prefix(INDEXES_PREFIX)? {
1486 let persisted = deserialize_index_meta(&value)?;
1487 let fqn = TableFqn::new(
1488 &persisted.catalog_name,
1489 &persisted.namespace_name,
1490 &persisted.table,
1491 );
1492 if table_set.contains(&fqn) {
1493 index_keys.push(key);
1494 }
1495 }
1496 for key in index_keys {
1497 txn.delete(key)?;
1498 }
1499 }
1500 Ok(())
1501 }
1502
1503 fn write_meta(&self, txn: &mut S::Transaction<'_>) -> Result<(), CatalogError> {
1504 let (table_id_counter, index_id_counter) = self.inner.counters();
1505 let meta = CatalogState {
1506 version: CATALOG_VERSION,
1507 table_id_counter,
1508 index_id_counter,
1509 };
1510 let meta_bytes = bincode::serialize(&meta)?;
1511 txn.put(META_KEY.to_vec(), meta_bytes)?;
1512 Ok(())
1513 }
1514
1515 pub fn persist_create_table(
1516 &mut self,
1517 txn: &mut S::Transaction<'_>,
1518 table: &TableMetadata,
1519 ) -> Result<(), CatalogError> {
1520 let persisted = PersistedTableMeta::from(table);
1521 let value = bincode::serialize(&persisted)?;
1522 txn.put(
1523 table_key(&table.catalog_name, &table.namespace_name, &table.name),
1524 value,
1525 )?;
1526 self.write_meta(txn)?;
1527 Ok(())
1528 }
1529
1530 pub fn persist_drop_table(
1531 &mut self,
1532 txn: &mut S::Transaction<'_>,
1533 fqn: &TableFqn,
1534 ) -> Result<(), CatalogError> {
1535 txn.delete(table_key(&fqn.catalog, &fqn.namespace, &fqn.table))?;
1536
1537 let mut to_delete: Vec<String> = Vec::new();
1539 let prefix = index_prefix(&fqn.catalog, &fqn.namespace, &fqn.table);
1540 for (key, _) in txn.scan_prefix(&prefix)? {
1541 let index_name = key_suffix(&prefix, &key)?;
1542 to_delete.push(index_name);
1543 }
1544 for index_name in to_delete {
1545 txn.delete(index_key(
1546 &fqn.catalog,
1547 &fqn.namespace,
1548 &fqn.table,
1549 &index_name,
1550 ))?;
1551 }
1552
1553 Ok(())
1554 }
1555
1556 pub fn persist_create_index(
1557 &mut self,
1558 txn: &mut S::Transaction<'_>,
1559 index: &IndexMetadata,
1560 ) -> Result<(), CatalogError> {
1561 let persisted = PersistedIndexMeta::from(index);
1562 let value = bincode::serialize(&persisted)?;
1563 txn.put(
1564 index_key(
1565 &index.catalog_name,
1566 &index.namespace_name,
1567 &index.table,
1568 &index.name,
1569 ),
1570 value,
1571 )?;
1572 self.write_meta(txn)?;
1573 Ok(())
1574 }
1575
1576 pub fn persist_drop_index(
1577 &mut self,
1578 txn: &mut S::Transaction<'_>,
1579 fqn: &IndexFqn,
1580 ) -> Result<(), CatalogError> {
1581 txn.delete(index_key(
1582 &fqn.catalog,
1583 &fqn.namespace,
1584 &fqn.table,
1585 &fqn.index,
1586 ))?;
1587 Ok(())
1588 }
1589
1590 pub fn persist_overlay(
1591 &mut self,
1592 txn: &mut S::Transaction<'_>,
1593 overlay: &CatalogOverlay,
1594 ) -> Result<(), CatalogError> {
1595 self.ensure_overlay_name_uniqueness(overlay)?;
1596
1597 for catalog in overlay.dropped_catalogs.iter() {
1598 self.persist_drop_catalog(txn, catalog)?;
1599 }
1600
1601 for (catalog, namespace) in overlay.dropped_namespaces.iter() {
1602 if overlay.dropped_catalogs.contains(catalog) {
1603 continue;
1604 }
1605 self.persist_drop_namespace(txn, catalog, namespace)?;
1606 }
1607
1608 for fqn in overlay.dropped_tables.iter() {
1609 if overlay.dropped_catalogs.contains(&fqn.catalog)
1610 || overlay
1611 .dropped_namespaces
1612 .contains(&(fqn.catalog.clone(), fqn.namespace.clone()))
1613 {
1614 continue;
1615 }
1616 self.persist_drop_table(txn, fqn)?;
1617 }
1618
1619 for fqn in overlay.dropped_indexes.iter() {
1620 if overlay.dropped_catalogs.contains(&fqn.catalog)
1621 || overlay
1622 .dropped_namespaces
1623 .contains(&(fqn.catalog.clone(), fqn.namespace.clone()))
1624 {
1625 continue;
1626 }
1627 self.persist_drop_index(txn, fqn)?;
1628 }
1629
1630 for meta in overlay.added_catalogs.values() {
1631 if overlay.dropped_catalogs.contains(&meta.name) {
1632 continue;
1633 }
1634 self.persist_create_catalog(txn, meta)?;
1635 }
1636
1637 for meta in overlay.added_namespaces.values() {
1638 if overlay.dropped_catalogs.contains(&meta.catalog_name)
1639 || overlay
1640 .dropped_namespaces
1641 .contains(&(meta.catalog_name.clone(), meta.name.clone()))
1642 {
1643 continue;
1644 }
1645 self.persist_create_namespace(txn, meta)?;
1646 }
1647
1648 for (fqn, table) in overlay.added_tables.iter() {
1649 if overlay.dropped_catalogs.contains(&fqn.catalog)
1650 || overlay
1651 .dropped_namespaces
1652 .contains(&(fqn.catalog.clone(), fqn.namespace.clone()))
1653 || overlay.dropped_tables.contains(fqn)
1654 {
1655 continue;
1656 }
1657 self.persist_create_table(txn, table)?;
1658 }
1659
1660 for (fqn, index) in overlay.added_indexes.iter() {
1661 if overlay.dropped_catalogs.contains(&fqn.catalog)
1662 || overlay
1663 .dropped_namespaces
1664 .contains(&(fqn.catalog.clone(), fqn.namespace.clone()))
1665 || overlay.dropped_indexes.contains(fqn)
1666 {
1667 continue;
1668 }
1669 self.persist_create_index(txn, index)?;
1670 }
1671
1672 Ok(())
1673 }
1674
1675 fn ensure_overlay_name_uniqueness(&self, overlay: &CatalogOverlay) -> Result<(), CatalogError> {
1676 let mut table_names: HashMap<String, TableFqn> = HashMap::new();
1677 for name in self.inner.table_names() {
1678 let Some(table) = self.inner.get_table(name) else {
1679 continue;
1680 };
1681 let fqn = TableFqn::from(table);
1682 if overlay.dropped_catalogs.contains(&fqn.catalog)
1683 || overlay
1684 .dropped_namespaces
1685 .contains(&(fqn.catalog.clone(), fqn.namespace.clone()))
1686 || overlay.dropped_tables.contains(&fqn)
1687 {
1688 continue;
1689 }
1690 table_names.insert(table.name.clone(), fqn);
1691 }
1692
1693 for (fqn, table) in overlay.added_tables.iter() {
1694 if overlay.dropped_catalogs.contains(&fqn.catalog)
1695 || overlay
1696 .dropped_namespaces
1697 .contains(&(fqn.catalog.clone(), fqn.namespace.clone()))
1698 || overlay.dropped_tables.contains(fqn)
1699 {
1700 continue;
1701 }
1702 if let Some(existing) = table_names.get(&table.name)
1703 && existing != fqn
1704 {
1705 return Err(CatalogError::InvalidKey(format!(
1706 "table name '{}' conflicts across namespaces: {}.{} vs {}.{}",
1707 table.name, existing.catalog, existing.namespace, fqn.catalog, fqn.namespace
1708 )));
1709 }
1710 table_names.insert(table.name.clone(), fqn.clone());
1711 }
1712
1713 let mut index_names: HashMap<String, IndexFqn> = HashMap::new();
1714 for name in self.inner.index_names() {
1715 let Some(index) = self.inner.get_index(name) else {
1716 continue;
1717 };
1718 let fqn = IndexFqn::from(index);
1719 if overlay.dropped_catalogs.contains(&fqn.catalog)
1720 || overlay
1721 .dropped_namespaces
1722 .contains(&(fqn.catalog.clone(), fqn.namespace.clone()))
1723 || overlay.dropped_indexes.contains(&fqn)
1724 || overlay.dropped_tables.contains(&TableFqn::new(
1725 &fqn.catalog,
1726 &fqn.namespace,
1727 &fqn.table,
1728 ))
1729 {
1730 continue;
1731 }
1732 index_names.insert(index.name.clone(), fqn);
1733 }
1734
1735 for (fqn, index) in overlay.added_indexes.iter() {
1736 if overlay.dropped_catalogs.contains(&fqn.catalog)
1737 || overlay
1738 .dropped_namespaces
1739 .contains(&(fqn.catalog.clone(), fqn.namespace.clone()))
1740 || overlay.dropped_indexes.contains(fqn)
1741 || overlay.dropped_tables.contains(&TableFqn::new(
1742 &fqn.catalog,
1743 &fqn.namespace,
1744 &fqn.table,
1745 ))
1746 {
1747 continue;
1748 }
1749 if let Some(existing) = index_names.get(&index.name)
1750 && existing != fqn
1751 {
1752 return Err(CatalogError::InvalidKey(format!(
1753 "index name '{}' conflicts across namespaces: {}.{} vs {}.{}",
1754 index.name, existing.catalog, existing.namespace, fqn.catalog, fqn.namespace
1755 )));
1756 }
1757 index_names.insert(index.name.clone(), fqn.clone());
1758 }
1759
1760 Ok(())
1761 }
1762
1763 fn namespace_dropped(overlay: &CatalogOverlay, catalog: &str, namespace: &str) -> bool {
1764 overlay.dropped_catalogs.contains(catalog)
1765 || overlay
1766 .dropped_namespaces
1767 .contains(&(catalog.to_string(), namespace.to_string()))
1768 }
1769
1770 fn overlay_added_table_by_name<'a>(
1771 overlay: &'a CatalogOverlay,
1772 name: &str,
1773 ) -> Option<&'a TableMetadata> {
1774 let mut iter = overlay
1775 .added_tables
1776 .values()
1777 .filter(|table| table.name == name);
1778 let first = iter.next()?;
1779 if iter.next().is_some() {
1780 return None;
1781 }
1782 Some(first)
1783 }
1784
1785 fn overlay_added_index_by_name<'a>(
1786 overlay: &'a CatalogOverlay,
1787 name: &str,
1788 ) -> Option<&'a IndexMetadata> {
1789 let mut iter = overlay
1790 .added_indexes
1791 .values()
1792 .filter(|index| index.name == name);
1793 let first = iter.next()?;
1794 if iter.next().is_some() {
1795 return None;
1796 }
1797 Some(first)
1798 }
1799
1800 fn base_table_conflicts_with_overlay(
1801 &self,
1802 overlay: &CatalogOverlay,
1803 table: &TableMetadata,
1804 ) -> bool {
1805 let Some(base) = self.inner.get_table(&table.name) else {
1806 return false;
1807 };
1808 if self.table_hidden_by_overlay(base, overlay) {
1809 return false;
1810 }
1811 if overlay.dropped_tables.contains(&TableFqn::from(base)) {
1812 return false;
1813 }
1814 TableFqn::from(base) != TableFqn::from(table)
1815 }
1816
1817 fn base_index_conflicts_with_overlay(
1818 &self,
1819 overlay: &CatalogOverlay,
1820 index: &IndexMetadata,
1821 ) -> bool {
1822 let Some(base) = self.inner.get_index(&index.name) else {
1823 return false;
1824 };
1825 if Self::namespace_dropped(overlay, &base.catalog_name, &base.namespace_name) {
1826 return false;
1827 }
1828 if overlay.dropped_indexes.contains(&IndexFqn::from(base)) {
1829 return false;
1830 }
1831 if self.dropped_table_matches_fqn(
1832 &base.table,
1833 &base.catalog_name,
1834 &base.namespace_name,
1835 overlay,
1836 ) {
1837 return false;
1838 }
1839 IndexFqn::from(base) != IndexFqn::from(index)
1840 }
1841
1842 fn table_hidden_by_overlay(&self, table: &TableMetadata, overlay: &CatalogOverlay) -> bool {
1843 Self::namespace_dropped(overlay, &table.catalog_name, &table.namespace_name)
1844 }
1845
1846 fn dropped_table_matches_fqn(
1847 &self,
1848 table_name: &str,
1849 catalog: &str,
1850 namespace: &str,
1851 overlay: &CatalogOverlay,
1852 ) -> bool {
1853 let fqn = TableFqn::new(catalog, namespace, table_name);
1854 overlay.dropped_tables.contains(&fqn)
1855 }
1856
1857 fn index_hidden_by_overlay(&self, index: &IndexMetadata, overlay: &CatalogOverlay) -> bool {
1858 let index_fqn = IndexFqn::from(index);
1859 if overlay.dropped_indexes.contains(&index_fqn) {
1860 return true;
1861 }
1862 if Self::namespace_dropped(overlay, &index.catalog_name, &index.namespace_name) {
1863 return true;
1864 }
1865 if self.dropped_table_matches_fqn(
1866 &index.table,
1867 &index.catalog_name,
1868 &index.namespace_name,
1869 overlay,
1870 ) {
1871 return true;
1872 }
1873 match self.get_table_in_txn(&index.table, overlay) {
1874 Some(table) => {
1875 table.catalog_name != index.catalog_name
1876 || table.namespace_name != index.namespace_name
1877 }
1878 None => true,
1879 }
1880 }
1881
1882 pub fn get_catalog_in_txn<'a>(
1883 &'a self,
1884 name: &str,
1885 overlay: &'a CatalogOverlay,
1886 ) -> Option<&'a CatalogMeta> {
1887 if overlay.dropped_catalogs.contains(name) {
1888 return None;
1889 }
1890 if let Some(catalog) = overlay.added_catalogs.get(name) {
1891 return Some(catalog);
1892 }
1893 self.catalogs.get(name)
1894 }
1895
1896 pub fn get_namespace_in_txn<'a>(
1897 &'a self,
1898 catalog_name: &str,
1899 namespace_name: &str,
1900 overlay: &'a CatalogOverlay,
1901 ) -> Option<&'a NamespaceMeta> {
1902 if overlay.dropped_catalogs.contains(catalog_name) {
1903 return None;
1904 }
1905 let key = (catalog_name.to_string(), namespace_name.to_string());
1906 if overlay.dropped_namespaces.contains(&key) {
1907 return None;
1908 }
1909 if let Some(namespace) = overlay.added_namespaces.get(&key) {
1910 return Some(namespace);
1911 }
1912 self.namespaces.get(&key)
1913 }
1914
1915 pub fn list_catalogs_in_txn(&self, overlay: &CatalogOverlay) -> Vec<CatalogMeta> {
1916 let mut catalogs: HashMap<String, CatalogMeta> = HashMap::new();
1917 for (name, meta) in &self.catalogs {
1918 if !overlay.dropped_catalogs.contains(name) {
1919 catalogs.insert(name.clone(), meta.clone());
1920 }
1921 }
1922 for (name, meta) in &overlay.added_catalogs {
1923 if !overlay.dropped_catalogs.contains(name) {
1924 catalogs.insert(name.clone(), meta.clone());
1925 }
1926 }
1927 let mut values: Vec<CatalogMeta> = catalogs.into_values().collect();
1928 values.sort_by(|a, b| a.name.cmp(&b.name));
1929 values
1930 }
1931
1932 pub fn list_namespaces_in_txn(
1933 &self,
1934 catalog_name: &str,
1935 overlay: &CatalogOverlay,
1936 ) -> Vec<NamespaceMeta> {
1937 if overlay.dropped_catalogs.contains(catalog_name) {
1938 return Vec::new();
1939 }
1940 let mut namespaces: HashMap<(String, String), NamespaceMeta> = HashMap::new();
1941 for ((catalog, namespace), meta) in &self.namespaces {
1942 if catalog != catalog_name {
1943 continue;
1944 }
1945 let key = (catalog.clone(), namespace.clone());
1946 if overlay.dropped_namespaces.contains(&key) {
1947 continue;
1948 }
1949 namespaces.insert(key, meta.clone());
1950 }
1951 for ((catalog, namespace), meta) in &overlay.added_namespaces {
1952 if catalog != catalog_name {
1953 continue;
1954 }
1955 let key = (catalog.clone(), namespace.clone());
1956 if overlay.dropped_namespaces.contains(&key) {
1957 continue;
1958 }
1959 namespaces.insert(key, meta.clone());
1960 }
1961 let mut values: Vec<NamespaceMeta> = namespaces.into_values().collect();
1962 values.sort_by(|a, b| a.name.cmp(&b.name));
1963 values
1964 }
1965
1966 pub fn table_exists_in_txn(&self, name: &str, overlay: &CatalogOverlay) -> bool {
1967 if let Some(table) = Self::overlay_added_table_by_name(overlay, name) {
1968 if self.table_hidden_by_overlay(table, overlay) {
1969 return false;
1970 }
1971 if self.base_table_conflicts_with_overlay(overlay, table) {
1972 return false;
1973 }
1974 return true;
1975 }
1976 match self.inner.get_table(name) {
1977 Some(table) => {
1978 !self.table_hidden_by_overlay(table, overlay)
1979 && !overlay.dropped_tables.contains(&TableFqn::from(table))
1980 }
1981 None => false,
1982 }
1983 }
1984
1985 pub fn get_table_in_txn<'a>(
1986 &'a self,
1987 name: &str,
1988 overlay: &'a CatalogOverlay,
1989 ) -> Option<&'a TableMetadata> {
1990 if let Some(table) = Self::overlay_added_table_by_name(overlay, name) {
1991 if self.table_hidden_by_overlay(table, overlay) {
1992 return None;
1993 }
1994 if self.base_table_conflicts_with_overlay(overlay, table) {
1995 return None;
1996 }
1997 return Some(table);
1998 }
1999 self.inner.get_table(name).filter(|table| {
2000 !self.table_hidden_by_overlay(table, overlay)
2001 && !overlay.dropped_tables.contains(&TableFqn::from(*table))
2002 })
2003 }
2004
2005 pub fn index_exists_in_txn(&self, name: &str, overlay: &CatalogOverlay) -> bool {
2006 if let Some(index) = Self::overlay_added_index_by_name(overlay, name) {
2007 if self.index_hidden_by_overlay(index, overlay) {
2008 return false;
2009 }
2010 if self.base_index_conflicts_with_overlay(overlay, index) {
2011 return false;
2012 }
2013 return true;
2014 }
2015 match self.inner.get_index(name) {
2016 Some(index) => !self.index_hidden_by_overlay(index, overlay),
2017 None => false,
2018 }
2019 }
2020
2021 pub fn get_index_in_txn<'a>(
2022 &'a self,
2023 name: &str,
2024 overlay: &'a CatalogOverlay,
2025 ) -> Option<&'a IndexMetadata> {
2026 if let Some(index) = Self::overlay_added_index_by_name(overlay, name) {
2027 if self.index_hidden_by_overlay(index, overlay) {
2028 return None;
2029 }
2030 if self.base_index_conflicts_with_overlay(overlay, index) {
2031 return None;
2032 }
2033 return Some(index);
2034 }
2035 match self.inner.get_index(name) {
2036 Some(index) if self.index_hidden_by_overlay(index, overlay) => None,
2037 other => other,
2038 }
2039 }
2040
2041 pub fn list_tables_in_txn(
2042 &self,
2043 catalog_name: &str,
2044 namespace_name: &str,
2045 overlay: &CatalogOverlay,
2046 ) -> Vec<TableMetadata> {
2047 if overlay.dropped_catalogs.contains(catalog_name) {
2048 return Vec::new();
2049 }
2050 if overlay
2051 .dropped_namespaces
2052 .contains(&(catalog_name.to_string(), namespace_name.to_string()))
2053 {
2054 return Vec::new();
2055 }
2056 let mut tables: HashMap<TableFqn, TableMetadata> = HashMap::new();
2057 for name in self.inner.table_names() {
2058 if let Some(table) = self.inner.get_table(name)
2059 && table.catalog_name == catalog_name
2060 && table.namespace_name == namespace_name
2061 {
2062 let fqn = TableFqn::from(table);
2063 if !overlay.dropped_tables.contains(&fqn) {
2064 tables.insert(fqn, table.clone());
2065 }
2066 }
2067 }
2068 for table in overlay.added_tables.values() {
2069 if table.catalog_name == catalog_name && table.namespace_name == namespace_name {
2070 tables.insert(TableFqn::from(table), table.clone());
2071 }
2072 }
2073 let mut values: Vec<TableMetadata> = tables.into_values().collect();
2074 values.sort_by(|a, b| a.name.cmp(&b.name));
2075 values
2076 }
2077
2078 pub fn list_indexes_in_txn(
2079 &self,
2080 fqn: &TableFqn,
2081 overlay: &CatalogOverlay,
2082 ) -> Vec<IndexMetadata> {
2083 if overlay.dropped_catalogs.contains(&fqn.catalog) {
2084 return Vec::new();
2085 }
2086 if overlay
2087 .dropped_namespaces
2088 .contains(&(fqn.catalog.clone(), fqn.namespace.clone()))
2089 {
2090 return Vec::new();
2091 }
2092 if overlay
2093 .dropped_tables
2094 .contains(&TableFqn::new(&fqn.catalog, &fqn.namespace, &fqn.table))
2095 {
2096 return Vec::new();
2097 }
2098
2099 let mut indexes: HashMap<IndexFqn, IndexMetadata> = HashMap::new();
2100 for index in self.inner.get_indexes_for_table(&fqn.table) {
2101 if index.catalog_name == fqn.catalog && index.namespace_name == fqn.namespace {
2102 let index_fqn = IndexFqn::from(index);
2103 if !overlay.dropped_indexes.contains(&index_fqn) {
2104 indexes.insert(index_fqn, index.clone());
2105 }
2106 }
2107 }
2108 for index in overlay.added_indexes.values() {
2109 if index.table == fqn.table
2110 && index.catalog_name == fqn.catalog
2111 && index.namespace_name == fqn.namespace
2112 {
2113 indexes.insert(IndexFqn::from(index), index.clone());
2114 }
2115 }
2116 let mut values: Vec<IndexMetadata> = indexes.into_values().collect();
2117 values.sort_by(|a, b| a.name.cmp(&b.name));
2118 values
2119 }
2120
2121 pub fn apply_overlay(&mut self, overlay: CatalogOverlay) {
2122 let CatalogOverlay {
2123 added_catalogs,
2124 dropped_catalogs,
2125 added_namespaces,
2126 dropped_namespaces,
2127 added_tables,
2128 dropped_tables,
2129 added_indexes,
2130 dropped_indexes,
2131 } = overlay;
2132
2133 for (name, meta) in added_catalogs {
2134 self.catalogs.insert(name, meta);
2135 }
2136 for (catalog_name, namespace_name) in dropped_namespaces.iter() {
2137 self.namespaces
2138 .remove(&(catalog_name.clone(), namespace_name.clone()));
2139 }
2140 for ((catalog_name, namespace_name), meta) in added_namespaces {
2141 self.namespaces.insert((catalog_name, namespace_name), meta);
2142 }
2143 for name in dropped_catalogs {
2144 self.catalogs.remove(&name);
2145 self.namespaces.retain(|(catalog, _), _| catalog != &name);
2146 }
2147 for (_, table) in added_tables {
2148 self.inner.insert_table_unchecked(table);
2149 }
2150 for fqn in dropped_tables {
2151 self.inner.remove_table_unchecked(&fqn.table);
2152 }
2153 for (_, index) in added_indexes {
2154 self.inner.insert_index_unchecked(index);
2155 }
2156 for fqn in dropped_indexes {
2157 self.inner.remove_index_unchecked(&fqn.index);
2158 }
2159 }
2160
2161 pub fn discard_overlay(_overlay: CatalogOverlay) {}
2162}
2163
2164impl<S: KVStore> Catalog for PersistentCatalog<S> {
2165 fn create_table(&mut self, table: TableMetadata) -> Result<(), PlannerError> {
2166 self.inner.create_table(table)
2167 }
2168
2169 fn get_table(&self, name: &str) -> Option<&TableMetadata> {
2170 self.inner.get_table(name)
2171 }
2172
2173 fn drop_table(&mut self, name: &str) -> Result<(), PlannerError> {
2174 self.inner.drop_table(name)
2175 }
2176
2177 fn create_index(&mut self, index: IndexMetadata) -> Result<(), PlannerError> {
2178 self.inner.create_index(index)
2179 }
2180
2181 fn get_index(&self, name: &str) -> Option<&IndexMetadata> {
2182 self.inner.get_index(name)
2183 }
2184
2185 fn get_indexes_for_table(&self, table: &str) -> Vec<&IndexMetadata> {
2186 self.inner.get_indexes_for_table(table)
2187 }
2188
2189 fn drop_index(&mut self, name: &str) -> Result<(), PlannerError> {
2190 self.inner.drop_index(name)
2191 }
2192
2193 fn table_exists(&self, name: &str) -> bool {
2194 self.inner.table_exists(name)
2195 }
2196
2197 fn index_exists(&self, name: &str) -> bool {
2198 self.inner.index_exists(name)
2199 }
2200
2201 fn next_table_id(&mut self) -> u32 {
2202 self.inner.next_table_id()
2203 }
2204
2205 fn next_index_id(&mut self) -> u32 {
2206 self.inner.next_index_id()
2207 }
2208
2209 fn list_tables(&self) -> Vec<TableMetadata> {
2210 let mut tables = Vec::new();
2211 for name in self.inner.table_names() {
2212 if let Some(table) = self.inner.get_table(name) {
2213 tables.push(table.clone());
2214 }
2215 }
2216 tables
2217 }
2218
2219 fn persistence_enabled(&self) -> bool {
2220 true
2221 }
2222}
2223
2224#[cfg(test)]
2225mod tests {
2226 use super::*;
2227 use crate::planner::types::ResolvedType;
2228 use std::collections::HashSet;
2229
2230 fn test_table(name: &str, id: u32) -> TableMetadata {
2231 TableMetadata::new(
2232 name,
2233 vec![ColumnMetadata::new("id", ResolvedType::Integer).with_primary_key(true)],
2234 )
2235 .with_table_id(id)
2236 .with_primary_key(vec!["id".to_string()])
2237 }
2238
2239 fn legacy_table_key(table_name: &str) -> Vec<u8> {
2240 let mut key = TABLES_PREFIX.to_vec();
2241 key.extend_from_slice(table_name.as_bytes());
2242 key
2243 }
2244
2245 fn legacy_index_key(index_name: &str) -> Vec<u8> {
2246 let mut key = INDEXES_PREFIX.to_vec();
2247 key.extend_from_slice(index_name.as_bytes());
2248 key
2249 }
2250
2251 fn seed_legacy_store(store: &Arc<alopex_core::kv::memory::MemoryKV>) {
2252 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
2253 let table = test_table("users", 7);
2254 let legacy_table = PersistedTableMetaV1 {
2255 table_id: table.table_id,
2256 name: table.name.clone(),
2257 columns: table
2258 .columns
2259 .iter()
2260 .map(PersistedColumnMeta::from)
2261 .collect(),
2262 primary_key: table.primary_key.clone(),
2263 storage_options: table.storage_options.clone().into(),
2264 };
2265 let table_bytes = bincode::serialize(&legacy_table).unwrap();
2266 txn.put(legacy_table_key("users"), table_bytes).unwrap();
2267
2268 let legacy_index = PersistedIndexMetaV1 {
2269 index_id: 3,
2270 name: "idx_users_id".to_string(),
2271 table: "users".to_string(),
2272 columns: vec!["id".to_string()],
2273 column_indices: vec![0],
2274 unique: false,
2275 method: Some(PersistedIndexType::BTree),
2276 options: Vec::new(),
2277 };
2278 let index_bytes = bincode::serialize(&legacy_index).unwrap();
2279 txn.put(legacy_index_key("idx_users_id"), index_bytes)
2280 .unwrap();
2281
2282 let meta = CatalogState {
2283 version: 1,
2284 table_id_counter: 7,
2285 index_id_counter: 3,
2286 };
2287 let meta_bytes = bincode::serialize(&meta).unwrap();
2288 txn.put(META_KEY.to_vec(), meta_bytes).unwrap();
2289 txn.commit_self().unwrap();
2290 }
2291
2292 #[test]
2293 fn load_empty_store() {
2294 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2295 let catalog = PersistentCatalog::load(store).unwrap();
2296 assert_eq!(catalog.inner.table_count(), 0);
2297 assert_eq!(catalog.inner.index_count(), 0);
2298 }
2299
2300 #[test]
2301 fn load_migrates_v1_keys_and_meta() {
2302 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2303 seed_legacy_store(&store);
2304
2305 let reloaded = PersistentCatalog::load(store.clone()).unwrap();
2306 assert!(reloaded.get_catalog("default").is_some());
2307 assert!(reloaded.get_namespace("default", "default").is_some());
2308
2309 let table = reloaded.get_table("users").unwrap();
2310 assert_eq!(table.catalog_name, "default");
2311 assert_eq!(table.namespace_name, "default");
2312
2313 let index = reloaded.get_index("idx_users_id").unwrap();
2314 assert_eq!(index.catalog_name, "default");
2315 assert_eq!(index.namespace_name, "default");
2316 assert_eq!(index.table, "users");
2317
2318 let mut txn = store.begin(TxnMode::ReadOnly).unwrap();
2319 assert!(txn.get(&legacy_table_key("users")).unwrap().is_none());
2320 assert!(
2321 txn.get(&table_key("default", "default", "users"))
2322 .unwrap()
2323 .is_some()
2324 );
2325 assert!(
2326 txn.get(&legacy_index_key("idx_users_id"))
2327 .unwrap()
2328 .is_none()
2329 );
2330 assert!(
2331 txn.get(&index_key("default", "default", "users", "idx_users_id"))
2332 .unwrap()
2333 .is_some()
2334 );
2335 let meta_bytes = txn.get(&META_KEY.to_vec()).unwrap().unwrap();
2336 let meta: CatalogState = bincode::deserialize(&meta_bytes).unwrap();
2337 assert_eq!(meta.version, CATALOG_VERSION);
2338 txn.rollback_self().unwrap();
2339 }
2340
2341 #[test]
2342 fn load_after_migration_keeps_v2_keys() {
2343 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2344 seed_legacy_store(&store);
2345
2346 let _ = PersistentCatalog::load(store.clone()).unwrap();
2347 let reloaded = PersistentCatalog::load(store.clone()).unwrap();
2348
2349 let table = reloaded.get_table("users").unwrap();
2350 assert_eq!(table.catalog_name, "default");
2351 assert_eq!(table.namespace_name, "default");
2352
2353 let mut txn = store.begin(TxnMode::ReadOnly).unwrap();
2354 assert!(txn.get(&legacy_table_key("users")).unwrap().is_none());
2355 assert!(
2356 txn.get(&table_key("default", "default", "users"))
2357 .unwrap()
2358 .is_some()
2359 );
2360 let meta_bytes = txn.get(&META_KEY.to_vec()).unwrap().unwrap();
2361 let meta: CatalogState = bincode::deserialize(&meta_bytes).unwrap();
2362 assert_eq!(meta.version, CATALOG_VERSION);
2363 txn.rollback_self().unwrap();
2364 }
2365
2366 #[test]
2367 fn create_table_persists() {
2368 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2369 let mut catalog = PersistentCatalog::new(store.clone());
2370
2371 catalog.inner.set_counters(1, 0);
2373
2374 let table = test_table("users", 1);
2375 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
2376 catalog.persist_create_table(&mut txn, &table).unwrap();
2377 txn.commit_self().unwrap();
2378
2379 let reloaded = PersistentCatalog::load(store).unwrap();
2380 assert!(reloaded.table_exists("users"));
2381 assert_eq!(reloaded.get_table("users").unwrap().table_id, 1);
2382 }
2383
2384 #[test]
2385 fn drop_table_removes() {
2386 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2387 let mut catalog = PersistentCatalog::new(store.clone());
2388 catalog.inner.set_counters(1, 0);
2389
2390 let table = test_table("users", 1);
2391 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
2392 catalog.persist_create_table(&mut txn, &table).unwrap();
2393 txn.commit_self().unwrap();
2394
2395 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
2396 let fqn = TableFqn::new("default", "default", "users");
2397 catalog.persist_drop_table(&mut txn, &fqn).unwrap();
2398 txn.commit_self().unwrap();
2399
2400 let reloaded = PersistentCatalog::load(store).unwrap();
2401 assert!(!reloaded.table_exists("users"));
2402 }
2403
2404 #[test]
2405 fn reload_preserves_state() {
2406 let temp_dir = tempfile::tempdir().unwrap();
2407 let wal_path = temp_dir.path().join("catalog.wal");
2408 let store = Arc::new(alopex_core::kv::memory::MemoryKV::open(&wal_path).unwrap());
2409 let mut catalog = PersistentCatalog::new(store.clone());
2410 catalog.inner.set_counters(1, 0);
2411
2412 let table = test_table("users", 1);
2413 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
2414 catalog.persist_create_table(&mut txn, &table).unwrap();
2415 txn.commit_self().unwrap();
2416 store.flush().unwrap();
2417
2418 drop(catalog);
2419 drop(store);
2420
2421 let store = Arc::new(alopex_core::kv::memory::MemoryKV::open(&wal_path).unwrap());
2422 let reloaded = PersistentCatalog::load(store).unwrap();
2423 assert!(reloaded.table_exists("users"));
2424 }
2425
2426 #[test]
2427 fn overlay_applied_on_commit() {
2428 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2429 let mut catalog = PersistentCatalog::new(store);
2430 let users = test_table("users", 1);
2431 catalog.inner.insert_table_unchecked(users.clone());
2432
2433 let mut overlay = CatalogOverlay::new();
2434 overlay.drop_table(&TableFqn::from(&users));
2435 let orders = test_table("orders", 2);
2436 overlay.add_table(TableFqn::from(&orders), orders);
2437
2438 assert!(!catalog.table_exists_in_txn("users", &overlay));
2439 assert!(catalog.table_exists_in_txn("orders", &overlay));
2440
2441 catalog.apply_overlay(overlay);
2442
2443 assert!(!catalog.table_exists("users"));
2444 assert!(catalog.table_exists("orders"));
2445 }
2446
2447 #[test]
2448 fn overlay_discarded_on_rollback() {
2449 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2450 let mut catalog = PersistentCatalog::new(store);
2451 let users = test_table("users", 1);
2452 catalog.inner.insert_table_unchecked(users.clone());
2453
2454 let mut overlay = CatalogOverlay::new();
2455 overlay.drop_table(&TableFqn::from(&users));
2456
2457 PersistentCatalog::<alopex_core::kv::memory::MemoryKV>::discard_overlay(overlay);
2458
2459 assert!(catalog.table_exists("users"));
2460 }
2461
2462 #[test]
2463 fn catalog_crud_persists() {
2464 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2465 let mut catalog = PersistentCatalog::new(store.clone());
2466
2467 let meta = CatalogMeta {
2468 name: "main".to_string(),
2469 comment: Some("primary".to_string()),
2470 storage_root: Some("/tmp/alopex".to_string()),
2471 };
2472
2473 catalog.create_catalog(meta.clone()).unwrap();
2474 assert!(catalog.get_catalog("main").is_some());
2475
2476 let mut txn = store.begin(TxnMode::ReadOnly).unwrap();
2477 let stored = txn.get(&catalog_key("main")).unwrap().unwrap();
2478 let decoded: CatalogMeta = bincode::deserialize(&stored).unwrap();
2479 txn.rollback_self().unwrap();
2480 assert_eq!(decoded, meta);
2481
2482 catalog.delete_catalog("main").unwrap();
2483 assert!(catalog.get_catalog("main").is_none());
2484 }
2485
2486 #[test]
2487 fn namespace_crud_persists_and_validates_catalog() {
2488 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2489 let mut catalog = PersistentCatalog::new(store.clone());
2490
2491 let missing_catalog = NamespaceMeta {
2492 name: "analytics".to_string(),
2493 catalog_name: "missing".to_string(),
2494 comment: None,
2495 storage_root: None,
2496 };
2497 let err = catalog.create_namespace(missing_catalog).unwrap_err();
2498 assert!(matches!(err, CatalogError::InvalidKey(_)));
2499
2500 catalog
2501 .create_catalog(CatalogMeta {
2502 name: "main".to_string(),
2503 comment: None,
2504 storage_root: None,
2505 })
2506 .unwrap();
2507
2508 let namespace = NamespaceMeta {
2509 name: "analytics".to_string(),
2510 catalog_name: "main".to_string(),
2511 comment: Some("warehouse".to_string()),
2512 storage_root: None,
2513 };
2514
2515 catalog.create_namespace(namespace.clone()).unwrap();
2516 assert!(catalog.get_namespace("main", "analytics").is_some());
2517
2518 let mut txn = store.begin(TxnMode::ReadOnly).unwrap();
2519 let stored = txn
2520 .get(&namespace_key("main", "analytics"))
2521 .unwrap()
2522 .unwrap();
2523 let decoded: NamespaceMeta = bincode::deserialize(&stored).unwrap();
2524 txn.rollback_self().unwrap();
2525 assert_eq!(decoded, namespace);
2526
2527 catalog.delete_namespace("main", "analytics").unwrap();
2528 assert!(catalog.get_namespace("main", "analytics").is_none());
2529 }
2530
2531 #[test]
2532 fn delete_catalog_removes_namespaces_from_store() {
2533 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2534 let mut catalog = PersistentCatalog::new(store.clone());
2535
2536 catalog
2537 .create_catalog(CatalogMeta {
2538 name: "main".to_string(),
2539 comment: None,
2540 storage_root: None,
2541 })
2542 .unwrap();
2543 catalog
2544 .create_namespace(NamespaceMeta {
2545 name: "analytics".to_string(),
2546 catalog_name: "main".to_string(),
2547 comment: None,
2548 storage_root: None,
2549 })
2550 .unwrap();
2551
2552 catalog.delete_catalog("main").unwrap();
2553
2554 let mut txn = store.begin(TxnMode::ReadOnly).unwrap();
2555 let mut prefix = NAMESPACES_PREFIX.to_vec();
2556 prefix.extend_from_slice(b"main");
2557 prefix.push(b'/');
2558 let remaining: Vec<_> = txn.scan_prefix(&prefix).unwrap().collect();
2559 txn.rollback_self().unwrap();
2560
2561 assert!(remaining.is_empty());
2562 assert!(catalog.list_namespaces("main").is_empty());
2563 }
2564
2565 #[test]
2566 fn delete_catalog_removes_tables_and_indexes_from_store() {
2567 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2568 let mut catalog = PersistentCatalog::new(store.clone());
2569
2570 catalog
2571 .create_catalog(CatalogMeta {
2572 name: "main".to_string(),
2573 comment: None,
2574 storage_root: None,
2575 })
2576 .unwrap();
2577
2578 let mut table = test_table("users", 1);
2579 table.catalog_name = "main".to_string();
2580 table.namespace_name = "default".to_string();
2581
2582 let mut index = IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
2583 .with_column_indices(vec![0]);
2584 index.catalog_name = "main".to_string();
2585 index.namespace_name = "default".to_string();
2586
2587 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
2588 catalog.persist_create_table(&mut txn, &table).unwrap();
2589 catalog.persist_create_index(&mut txn, &index).unwrap();
2590 txn.commit_self().unwrap();
2591
2592 catalog.inner.insert_table_unchecked(table);
2593 catalog.inner.insert_index_unchecked(index);
2594
2595 catalog.delete_catalog("main").unwrap();
2596
2597 assert!(catalog.inner.get_table("users").is_none());
2598 assert!(catalog.inner.get_index("idx_users_id").is_none());
2599
2600 let mut txn = store.begin(TxnMode::ReadOnly).unwrap();
2601 assert!(
2602 txn.get(&table_key("main", "default", "users"))
2603 .unwrap()
2604 .is_none()
2605 );
2606 assert!(
2607 txn.get(&index_key("main", "default", "users", "idx_users_id"))
2608 .unwrap()
2609 .is_none()
2610 );
2611 txn.rollback_self().unwrap();
2612 }
2613
2614 #[test]
2615 fn index_meta_loads_catalog_and_namespace_from_table() {
2616 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2617 let mut catalog = PersistentCatalog::new(store.clone());
2618 catalog.inner.set_counters(1, 1);
2619
2620 let mut table = test_table("users", 1);
2621 table.catalog_name = "main".to_string();
2622 table.namespace_name = "analytics".to_string();
2623
2624 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
2625 catalog.persist_create_table(&mut txn, &table).unwrap();
2626 let mut index = IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
2627 .with_column_indices(vec![0])
2628 .with_method(IndexMethod::BTree);
2629 index.catalog_name = "main".to_string();
2630 index.namespace_name = "analytics".to_string();
2631 catalog.persist_create_index(&mut txn, &index).unwrap();
2632 txn.commit_self().unwrap();
2633
2634 let reloaded = PersistentCatalog::load(store).unwrap();
2635 let index = reloaded.get_index("idx_users_id").unwrap();
2636 assert_eq!(index.catalog_name, "main");
2637 assert_eq!(index.namespace_name, "analytics");
2638 }
2639
2640 #[test]
2641 fn legacy_index_meta_loads_catalog_and_namespace_from_table() {
2642 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2643 let mut catalog = PersistentCatalog::new(store.clone());
2644 catalog.inner.set_counters(1, 1);
2645
2646 let mut table = test_table("users", 1);
2647 table.catalog_name = "main".to_string();
2648 table.namespace_name = "analytics".to_string();
2649
2650 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
2651 catalog.persist_create_table(&mut txn, &table).unwrap();
2652 let legacy = PersistedIndexMetaV1 {
2653 index_id: 1,
2654 name: "idx_users_id".to_string(),
2655 table: "users".to_string(),
2656 columns: vec!["id".to_string()],
2657 column_indices: vec![0],
2658 unique: false,
2659 method: Some(PersistedIndexType::BTree),
2660 options: Vec::new(),
2661 };
2662 let bytes = bincode::serialize(&legacy).unwrap();
2663 txn.put(
2664 index_key("main", "analytics", "users", "idx_users_id"),
2665 bytes,
2666 )
2667 .unwrap();
2668 txn.commit_self().unwrap();
2669
2670 let reloaded = PersistentCatalog::load(store).unwrap();
2671 let index = reloaded.get_index("idx_users_id").unwrap();
2672 assert_eq!(index.catalog_name, "main");
2673 assert_eq!(index.namespace_name, "analytics");
2674 }
2675
2676 #[test]
2677 fn overlay_catalog_get_and_list() {
2678 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2679 let mut catalog = PersistentCatalog::new(store);
2680 catalog
2681 .create_catalog(CatalogMeta {
2682 name: "main".to_string(),
2683 comment: None,
2684 storage_root: None,
2685 })
2686 .unwrap();
2687
2688 let mut overlay = CatalogOverlay::new();
2689 overlay.add_catalog(CatalogMeta {
2690 name: "temp".to_string(),
2691 comment: None,
2692 storage_root: None,
2693 });
2694 overlay.drop_catalog("main");
2695
2696 assert!(catalog.get_catalog_in_txn("main", &overlay).is_none());
2697 assert!(catalog.get_catalog_in_txn("temp", &overlay).is_some());
2698
2699 let names: Vec<String> = catalog
2700 .list_catalogs_in_txn(&overlay)
2701 .into_iter()
2702 .map(|meta| meta.name)
2703 .collect();
2704 assert_eq!(names, vec!["temp".to_string()]);
2705 }
2706
2707 #[test]
2708 fn overlay_namespace_get_and_list() {
2709 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2710 let mut catalog = PersistentCatalog::new(store);
2711 catalog
2712 .create_catalog(CatalogMeta {
2713 name: "main".to_string(),
2714 comment: None,
2715 storage_root: None,
2716 })
2717 .unwrap();
2718 catalog
2719 .create_namespace(NamespaceMeta {
2720 name: "default".to_string(),
2721 catalog_name: "main".to_string(),
2722 comment: None,
2723 storage_root: None,
2724 })
2725 .unwrap();
2726
2727 let mut overlay = CatalogOverlay::new();
2728 overlay.add_namespace(NamespaceMeta {
2729 name: "analytics".to_string(),
2730 catalog_name: "main".to_string(),
2731 comment: None,
2732 storage_root: None,
2733 });
2734 overlay.drop_namespace("main", "default");
2735
2736 assert!(
2737 catalog
2738 .get_namespace_in_txn("main", "default", &overlay)
2739 .is_none()
2740 );
2741 assert!(
2742 catalog
2743 .get_namespace_in_txn("main", "analytics", &overlay)
2744 .is_some()
2745 );
2746
2747 let names: Vec<String> = catalog
2748 .list_namespaces_in_txn("main", &overlay)
2749 .into_iter()
2750 .map(|meta| meta.name)
2751 .collect();
2752 assert_eq!(names, vec!["analytics".to_string()]);
2753 }
2754
2755 #[test]
2756 fn overlay_table_and_index_list() {
2757 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2758 let mut catalog = PersistentCatalog::new(store);
2759
2760 let mut users = test_table("users", 1);
2761 users.catalog_name = "main".to_string();
2762 users.namespace_name = "default".to_string();
2763 catalog.inner.insert_table_unchecked(users.clone());
2764
2765 let mut users_index =
2766 IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
2767 .with_column_indices(vec![0]);
2768 users_index.catalog_name = "main".to_string();
2769 users_index.namespace_name = "default".to_string();
2770 catalog.inner.insert_index_unchecked(users_index);
2771
2772 let mut overlay = CatalogOverlay::new();
2773
2774 let mut orders = test_table("orders", 2);
2775 orders.catalog_name = "main".to_string();
2776 orders.namespace_name = "default".to_string();
2777 overlay.add_table(TableFqn::from(&orders), orders.clone());
2778
2779 let mut orders_index =
2780 IndexMetadata::new(2, "idx_orders_id", "orders", vec!["id".to_string()])
2781 .with_column_indices(vec![0]);
2782 orders_index.catalog_name = "main".to_string();
2783 orders_index.namespace_name = "default".to_string();
2784 overlay.add_index(IndexFqn::from(&orders_index), orders_index);
2785
2786 overlay.drop_table(&TableFqn::from(&users));
2787
2788 let table_names: Vec<String> = catalog
2789 .list_tables_in_txn("main", "default", &overlay)
2790 .into_iter()
2791 .map(|table| table.name)
2792 .collect();
2793 assert_eq!(table_names, vec!["orders".to_string()]);
2794
2795 let users_fqn = TableFqn::new("main", "default", "users");
2796 assert!(catalog.list_indexes_in_txn(&users_fqn, &overlay).is_empty());
2797
2798 let orders_fqn = TableFqn::new("main", "default", "orders");
2799 let index_names: Vec<String> = catalog
2800 .list_indexes_in_txn(&orders_fqn, &overlay)
2801 .into_iter()
2802 .map(|index| index.name)
2803 .collect();
2804 assert_eq!(index_names, vec!["idx_orders_id".to_string()]);
2805 }
2806
2807 #[test]
2808 fn overlay_name_lookup_ambiguous_returns_none() {
2809 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2810 let catalog = PersistentCatalog::new(store);
2811
2812 let mut overlay = CatalogOverlay::new();
2813
2814 let mut users_default = test_table("users", 1);
2815 users_default.catalog_name = "main".to_string();
2816 users_default.namespace_name = "default".to_string();
2817 overlay.add_table(TableFqn::from(&users_default), users_default);
2818
2819 let mut users_analytics = test_table("users", 2);
2820 users_analytics.catalog_name = "main".to_string();
2821 users_analytics.namespace_name = "analytics".to_string();
2822 overlay.add_table(TableFqn::from(&users_analytics), users_analytics);
2823
2824 assert!(catalog.get_table_in_txn("users", &overlay).is_none());
2825 assert!(!catalog.table_exists_in_txn("users", &overlay));
2826
2827 let mut idx_default =
2828 IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
2829 .with_column_indices(vec![0]);
2830 idx_default.catalog_name = "main".to_string();
2831 idx_default.namespace_name = "default".to_string();
2832 overlay.add_index(IndexFqn::from(&idx_default), idx_default);
2833
2834 let mut idx_analytics =
2835 IndexMetadata::new(2, "idx_users_id", "users", vec!["id".to_string()])
2836 .with_column_indices(vec![0]);
2837 idx_analytics.catalog_name = "main".to_string();
2838 idx_analytics.namespace_name = "analytics".to_string();
2839 overlay.add_index(IndexFqn::from(&idx_analytics), idx_analytics);
2840
2841 assert!(catalog.get_index_in_txn("idx_users_id", &overlay).is_none());
2842 assert!(!catalog.index_exists_in_txn("idx_users_id", &overlay));
2843 }
2844
2845 #[test]
2846 fn overlay_name_lookup_ambiguous_with_base_returns_none() {
2847 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2848 let mut catalog = PersistentCatalog::new(store);
2849
2850 let mut overlay = CatalogOverlay::new();
2851
2852 let mut base_users = test_table("users", 1);
2853 base_users.catalog_name = "main".to_string();
2854 base_users.namespace_name = "default".to_string();
2855 catalog.inner.insert_table_unchecked(base_users);
2856
2857 let mut overlay_users = test_table("users", 2);
2858 overlay_users.catalog_name = "main".to_string();
2859 overlay_users.namespace_name = "analytics".to_string();
2860 overlay.add_table(TableFqn::from(&overlay_users), overlay_users);
2861
2862 assert!(catalog.get_table_in_txn("users", &overlay).is_none());
2863 assert!(!catalog.table_exists_in_txn("users", &overlay));
2864
2865 let mut base_index = IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
2866 .with_column_indices(vec![0]);
2867 base_index.catalog_name = "main".to_string();
2868 base_index.namespace_name = "default".to_string();
2869 catalog.inner.insert_index_unchecked(base_index);
2870
2871 let mut overlay_index =
2872 IndexMetadata::new(2, "idx_users_id", "users", vec!["id".to_string()])
2873 .with_column_indices(vec![0]);
2874 overlay_index.catalog_name = "main".to_string();
2875 overlay_index.namespace_name = "analytics".to_string();
2876 overlay.add_index(IndexFqn::from(&overlay_index), overlay_index);
2877
2878 assert!(catalog.get_index_in_txn("idx_users_id", &overlay).is_none());
2879 assert!(!catalog.index_exists_in_txn("idx_users_id", &overlay));
2880 }
2881
2882 #[test]
2883 fn overlay_fqn_tables_separate_namespaces() {
2884 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2885 let catalog = PersistentCatalog::new(store);
2886
2887 let mut overlay = CatalogOverlay::new();
2888
2889 let mut users_default = test_table("users", 1);
2890 users_default.catalog_name = "main".to_string();
2891 users_default.namespace_name = "default".to_string();
2892 overlay.add_table(TableFqn::from(&users_default), users_default.clone());
2893
2894 let mut users_analytics = test_table("users", 2);
2895 users_analytics.catalog_name = "main".to_string();
2896 users_analytics.namespace_name = "analytics".to_string();
2897 overlay.add_table(TableFqn::from(&users_analytics), users_analytics.clone());
2898
2899 let default_tables: Vec<String> = catalog
2900 .list_tables_in_txn("main", "default", &overlay)
2901 .into_iter()
2902 .map(|table| table.name)
2903 .collect();
2904 assert_eq!(default_tables, vec!["users".to_string()]);
2905
2906 let analytics_tables: Vec<String> = catalog
2907 .list_tables_in_txn("main", "analytics", &overlay)
2908 .into_iter()
2909 .map(|table| table.name)
2910 .collect();
2911 assert_eq!(analytics_tables, vec!["users".to_string()]);
2912
2913 overlay.drop_table(&TableFqn::from(&users_default));
2914
2915 let default_tables_after: Vec<String> = catalog
2916 .list_tables_in_txn("main", "default", &overlay)
2917 .into_iter()
2918 .map(|table| table.name)
2919 .collect();
2920 assert!(default_tables_after.is_empty());
2921
2922 let analytics_tables_after: Vec<String> = catalog
2923 .list_tables_in_txn("main", "analytics", &overlay)
2924 .into_iter()
2925 .map(|table| table.name)
2926 .collect();
2927 assert_eq!(analytics_tables_after, vec!["users".to_string()]);
2928 }
2929
2930 #[test]
2931 fn overlay_fqn_indexes_separate_namespaces() {
2932 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2933 let catalog = PersistentCatalog::new(store);
2934
2935 let mut overlay = CatalogOverlay::new();
2936
2937 let mut users_default = test_table("users", 1);
2938 users_default.catalog_name = "main".to_string();
2939 users_default.namespace_name = "default".to_string();
2940 overlay.add_table(TableFqn::from(&users_default), users_default.clone());
2941
2942 let mut users_analytics = test_table("users", 2);
2943 users_analytics.catalog_name = "main".to_string();
2944 users_analytics.namespace_name = "analytics".to_string();
2945 overlay.add_table(TableFqn::from(&users_analytics), users_analytics.clone());
2946
2947 let mut idx_default =
2948 IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
2949 .with_column_indices(vec![0]);
2950 idx_default.catalog_name = "main".to_string();
2951 idx_default.namespace_name = "default".to_string();
2952 overlay.add_index(IndexFqn::from(&idx_default), idx_default);
2953
2954 let mut idx_analytics =
2955 IndexMetadata::new(2, "idx_users_id", "users", vec!["id".to_string()])
2956 .with_column_indices(vec![0]);
2957 idx_analytics.catalog_name = "main".to_string();
2958 idx_analytics.namespace_name = "analytics".to_string();
2959 overlay.add_index(IndexFqn::from(&idx_analytics), idx_analytics);
2960
2961 let default_fqn = TableFqn::new("main", "default", "users");
2962 let analytics_fqn = TableFqn::new("main", "analytics", "users");
2963
2964 let default_indexes: Vec<String> = catalog
2965 .list_indexes_in_txn(&default_fqn, &overlay)
2966 .into_iter()
2967 .map(|index| index.name)
2968 .collect();
2969 assert_eq!(default_indexes, vec!["idx_users_id".to_string()]);
2970
2971 let analytics_indexes: Vec<String> = catalog
2972 .list_indexes_in_txn(&analytics_fqn, &overlay)
2973 .into_iter()
2974 .map(|index| index.name)
2975 .collect();
2976 assert_eq!(analytics_indexes, vec!["idx_users_id".to_string()]);
2977 }
2978
2979 #[test]
2980 fn persist_overlay_rejects_duplicate_table_names() {
2981 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
2982 let mut catalog = PersistentCatalog::new(store.clone());
2983
2984 let mut overlay = CatalogOverlay::new();
2985
2986 let mut users_default = test_table("users", 1);
2987 users_default.catalog_name = "main".to_string();
2988 users_default.namespace_name = "default".to_string();
2989 overlay.add_table(TableFqn::from(&users_default), users_default);
2990
2991 let mut users_analytics = test_table("users", 2);
2992 users_analytics.catalog_name = "main".to_string();
2993 users_analytics.namespace_name = "analytics".to_string();
2994 overlay.add_table(TableFqn::from(&users_analytics), users_analytics);
2995
2996 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
2997 let err = catalog.persist_overlay(&mut txn, &overlay).unwrap_err();
2998 assert!(matches!(err, CatalogError::InvalidKey(_)));
2999 txn.rollback_self().unwrap();
3000 }
3001
3002 #[test]
3003 fn persist_overlay_rejects_duplicate_index_names() {
3004 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
3005 let mut catalog = PersistentCatalog::new(store.clone());
3006
3007 let mut overlay = CatalogOverlay::new();
3008
3009 let mut idx_default =
3010 IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
3011 .with_column_indices(vec![0]);
3012 idx_default.catalog_name = "main".to_string();
3013 idx_default.namespace_name = "default".to_string();
3014 overlay.add_index(IndexFqn::from(&idx_default), idx_default);
3015
3016 let mut idx_analytics =
3017 IndexMetadata::new(2, "idx_users_id", "users", vec!["id".to_string()])
3018 .with_column_indices(vec![0]);
3019 idx_analytics.catalog_name = "main".to_string();
3020 idx_analytics.namespace_name = "analytics".to_string();
3021 overlay.add_index(IndexFqn::from(&idx_analytics), idx_analytics);
3022
3023 let mut txn = store.begin(TxnMode::ReadWrite).unwrap();
3024 let err = catalog.persist_overlay(&mut txn, &overlay).unwrap_err();
3025 assert!(matches!(err, CatalogError::InvalidKey(_)));
3026 txn.rollback_self().unwrap();
3027 }
3028
3029 #[test]
3030 fn overlay_drop_cascade_namespace_removes_children() {
3031 let mut overlay = CatalogOverlay::new();
3032
3033 overlay.add_namespace(NamespaceMeta {
3034 name: "default".to_string(),
3035 catalog_name: "main".to_string(),
3036 comment: None,
3037 storage_root: None,
3038 });
3039
3040 let mut users = test_table("users", 1);
3041 users.catalog_name = "main".to_string();
3042 users.namespace_name = "default".to_string();
3043 overlay.add_table(TableFqn::from(&users), users.clone());
3044
3045 let mut users_index =
3046 IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
3047 .with_column_indices(vec![0]);
3048 users_index.catalog_name = "main".to_string();
3049 users_index.namespace_name = "default".to_string();
3050 let users_index_fqn = IndexFqn::from(&users_index);
3051 overlay.add_index(users_index_fqn.clone(), users_index);
3052
3053 overlay.drop_cascade_namespace("main", "default");
3054
3055 assert!(
3056 overlay
3057 .dropped_namespaces
3058 .contains(&("main".to_string(), "default".to_string()))
3059 );
3060 assert!(overlay.dropped_tables.contains(&TableFqn::from(&users)));
3061 assert!(overlay.dropped_indexes.contains(&users_index_fqn));
3062 assert!(!overlay.added_tables.contains_key(&TableFqn::from(&users)));
3063 assert!(!overlay.added_indexes.contains_key(&users_index_fqn));
3064 }
3065
3066 #[test]
3067 fn overlay_drop_cascade_catalog_removes_children() {
3068 let mut overlay = CatalogOverlay::new();
3069
3070 overlay.add_catalog(CatalogMeta {
3071 name: "main".to_string(),
3072 comment: None,
3073 storage_root: None,
3074 });
3075 overlay.add_namespace(NamespaceMeta {
3076 name: "default".to_string(),
3077 catalog_name: "main".to_string(),
3078 comment: None,
3079 storage_root: None,
3080 });
3081
3082 let mut users = test_table("users", 1);
3083 users.catalog_name = "main".to_string();
3084 users.namespace_name = "default".to_string();
3085 overlay.add_table(TableFqn::from(&users), users.clone());
3086
3087 let mut users_index =
3088 IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
3089 .with_column_indices(vec![0]);
3090 users_index.catalog_name = "main".to_string();
3091 users_index.namespace_name = "default".to_string();
3092 let users_index_fqn = IndexFqn::from(&users_index);
3093 overlay.add_index(users_index_fqn.clone(), users_index);
3094
3095 overlay.drop_cascade_catalog("main");
3096
3097 assert!(overlay.dropped_catalogs.contains("main"));
3098 assert!(overlay.dropped_tables.contains(&TableFqn::from(&users)));
3099 assert!(overlay.dropped_indexes.contains(&users_index_fqn));
3100 assert!(!overlay.added_tables.contains_key(&TableFqn::from(&users)));
3101 assert!(!overlay.added_indexes.contains_key(&users_index_fqn));
3102 }
3103
3104 #[test]
3105 fn dropped_namespace_hides_tables_and_indexes() {
3106 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
3107 let mut catalog = PersistentCatalog::new(store);
3108
3109 let mut users = test_table("users", 1);
3110 users.catalog_name = "main".to_string();
3111 users.namespace_name = "default".to_string();
3112 catalog.inner.insert_table_unchecked(users);
3113
3114 let mut users_index =
3115 IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
3116 .with_column_indices(vec![0]);
3117 users_index.catalog_name = "main".to_string();
3118 users_index.namespace_name = "default".to_string();
3119 catalog.inner.insert_index_unchecked(users_index);
3120
3121 let mut overlay = CatalogOverlay::new();
3122 overlay.drop_namespace("main", "default");
3123
3124 let table_names: Vec<String> = catalog
3125 .list_tables_in_txn("main", "default", &overlay)
3126 .into_iter()
3127 .map(|table| table.name)
3128 .collect();
3129 assert!(table_names.is_empty());
3130
3131 let fqn = TableFqn {
3132 catalog: "main".to_string(),
3133 namespace: "default".to_string(),
3134 table: "users".to_string(),
3135 };
3136 let index_names: Vec<String> = catalog
3137 .list_indexes_in_txn(&fqn, &overlay)
3138 .into_iter()
3139 .map(|index| index.name)
3140 .collect();
3141 assert!(index_names.is_empty());
3142 }
3143
3144 #[test]
3145 fn dropped_namespace_hides_get_and_exists() {
3146 let store = Arc::new(alopex_core::kv::memory::MemoryKV::new());
3147 let mut catalog = PersistentCatalog::new(store);
3148
3149 let mut users = test_table("users", 1);
3150 users.catalog_name = "main".to_string();
3151 users.namespace_name = "default".to_string();
3152 catalog.inner.insert_table_unchecked(users);
3153
3154 let mut users_index =
3155 IndexMetadata::new(1, "idx_users_id", "users", vec!["id".to_string()])
3156 .with_column_indices(vec![0]);
3157 users_index.catalog_name = "main".to_string();
3158 users_index.namespace_name = "default".to_string();
3159 catalog.inner.insert_index_unchecked(users_index);
3160
3161 let mut overlay = CatalogOverlay::new();
3162 overlay.drop_namespace("main", "default");
3163
3164 assert!(!catalog.table_exists_in_txn("users", &overlay));
3165 assert!(catalog.get_table_in_txn("users", &overlay).is_none());
3166 assert!(!catalog.index_exists_in_txn("idx_users_id", &overlay));
3167 assert!(catalog.get_index_in_txn("idx_users_id", &overlay).is_none());
3168
3169 let view = TxnCatalogView::new(&catalog, &overlay);
3170 assert!(view.get_indexes_for_table("users").is_empty());
3171 }
3172
3173 #[test]
3174 fn persisted_catalog_meta_roundtrip() {
3175 let meta = PersistedCatalogMeta {
3176 name: "main".to_string(),
3177 comment: Some("primary catalog".to_string()),
3178 storage_root: Some("/tmp/alopex".to_string()),
3179 };
3180 let bytes = bincode::serialize(&meta).unwrap();
3181 let decoded: PersistedCatalogMeta = bincode::deserialize(&bytes).unwrap();
3182 assert_eq!(meta, decoded);
3183 }
3184
3185 #[test]
3186 fn persisted_namespace_meta_roundtrip() {
3187 let meta = PersistedNamespaceMeta {
3188 name: "analytics".to_string(),
3189 catalog_name: "main".to_string(),
3190 comment: Some("warehouse".to_string()),
3191 storage_root: Some("s3://bucket/ns".to_string()),
3192 };
3193 let bytes = bincode::serialize(&meta).unwrap();
3194 let decoded: PersistedNamespaceMeta = bincode::deserialize(&bytes).unwrap();
3195 assert_eq!(meta, decoded);
3196 }
3197
3198 #[test]
3199 fn table_fqn_hash_and_eq() {
3200 let first = TableFqn {
3201 catalog: "main".to_string(),
3202 namespace: "default".to_string(),
3203 table: "users".to_string(),
3204 };
3205 let same = TableFqn {
3206 catalog: "main".to_string(),
3207 namespace: "default".to_string(),
3208 table: "users".to_string(),
3209 };
3210 let different = TableFqn {
3211 catalog: "main".to_string(),
3212 namespace: "default".to_string(),
3213 table: "orders".to_string(),
3214 };
3215
3216 let mut set = HashSet::new();
3217 set.insert(first);
3218 assert!(set.contains(&same));
3219 assert!(!set.contains(&different));
3220 }
3221
3222 #[test]
3223 fn index_fqn_hash_and_eq() {
3224 let first = IndexFqn {
3225 catalog: "main".to_string(),
3226 namespace: "default".to_string(),
3227 table: "users".to_string(),
3228 index: "idx_users_id".to_string(),
3229 };
3230 let same = IndexFqn {
3231 catalog: "main".to_string(),
3232 namespace: "default".to_string(),
3233 table: "users".to_string(),
3234 index: "idx_users_id".to_string(),
3235 };
3236 let different = IndexFqn {
3237 catalog: "main".to_string(),
3238 namespace: "default".to_string(),
3239 table: "users".to_string(),
3240 index: "idx_users_email".to_string(),
3241 };
3242
3243 let mut set = HashSet::new();
3244 set.insert(first);
3245 assert!(set.contains(&same));
3246 assert!(!set.contains(&different));
3247 }
3248
3249 #[test]
3250 fn table_type_and_data_source_format_serde() {
3251 let managed = serde_json::to_string(&TableType::Managed).unwrap();
3252 let external = serde_json::to_string(&TableType::External).unwrap();
3253 let alopex = serde_json::to_string(&DataSourceFormat::Alopex).unwrap();
3254 let parquet = serde_json::to_string(&DataSourceFormat::Parquet).unwrap();
3255 let delta = serde_json::to_string(&DataSourceFormat::Delta).unwrap();
3256
3257 assert_eq!(managed, "\"MANAGED\"");
3258 assert_eq!(external, "\"EXTERNAL\"");
3259 assert_eq!(alopex, "\"ALOPEX\"");
3260 assert_eq!(parquet, "\"PARQUET\"");
3261 assert_eq!(delta, "\"DELTA\"");
3262 }
3263
3264 #[test]
3265 fn nested_catalog_types_roundtrip_recursively() {
3266 let data_type = ResolvedType::Struct(vec![
3267 (
3268 "items".into(),
3269 ResolvedType::Array(Box::new(ResolvedType::Integer)),
3270 ),
3271 (
3272 "attrs".into(),
3273 ResolvedType::Map {
3274 key: Box::new(ResolvedType::Text),
3275 value: Box::new(ResolvedType::BigInt),
3276 },
3277 ),
3278 ]);
3279 let persisted = PersistedType::from(data_type.clone());
3280 let bytes = bincode::serialize(&persisted).unwrap();
3281 let decoded: PersistedType = bincode::deserialize(&bytes).unwrap();
3282 assert_eq!(ResolvedType::from(decoded), data_type);
3283 }
3284}