Skip to main content

alopex_sql/catalog/
persistent.rs

1//! 永続化対応カタログ実装。
2//!
3//! 既存の `TableMetadata` / `IndexMetadata` は `Expr` を含むため、そのままシリアライズして
4//! 永続化することができない。そこで、本モジュールでは永続化用 DTO を定義し、KV ストアへ
5//! bincode で保存する。
6//!
7//! 注意: 現状は `ColumnMetadata.default`(DEFAULT 式)を永続化しない。復元時は `None` となる。
8
9use 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
25/// カタログ用キープレフィックス。
26pub 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
847/// トランザクション内(オーバーレイ込み)で参照するための Catalog ビュー。
848///
849/// DML/SELECT の実行や Planner の参照用途に使う。書き込み系 API は利用しない前提のため、
850/// `Catalog` trait の書き込みメソッドは `unreachable!()` とする。
851pub 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/// 永続カタログ実装。
953#[derive(Debug)]
954/// 永続化対応のカタログ実装。
955///
956/// # Examples
957///
958/// ```
959/// use std::sync::Arc;
960/// use alopex_core::kv::memory::MemoryKV;
961/// use alopex_sql::Catalog;
962/// use alopex_sql::catalog::PersistentCatalog;
963///
964/// let store = Arc::new(MemoryKV::new());
965/// let catalog = PersistentCatalog::new(store);
966/// assert!(catalog.table_exists("users") == false);
967/// ```
968pub 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        // テーブルをロード(まずテーブルを入れてからインデックスを入れる)
1058        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            // 参照先テーブルがない場合はスキップ(破損対策)
1095            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        // テーブルに紐づくインデックスも削除する。
1538        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        // inner のカウンタを更新して meta が書き込まれることを担保する
2372        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}