llkv_table/catalog/
manager.rs

1//! High-level service for creating tables.
2//!
3//! Use `CatalogManager` to create tables. It coordinates metadata persistence,
4//! catalog registration, and storage initialization.
5
6#![forbid(unsafe_code)]
7
8use std::convert::TryFrom;
9use std::sync::Arc;
10use std::time::{SystemTime, UNIX_EPOCH};
11
12use arrow::array::ArrayRef;
13use arrow::datatypes::{DataType, Field, Schema};
14use arrow::record_batch::RecordBatch;
15use llkv_column_map::ColumnStore;
16use llkv_column_map::store::IndexKind;
17use llkv_plan::{DropIndexPlan, ForeignKeySpec, PlanColumnSpec};
18use llkv_result::{Error, Result as LlkvResult};
19use llkv_storage::pager::Pager;
20use rustc_hash::{FxHashMap, FxHashSet};
21use simd_r_drive_entry_handle::EntryHandle;
22
23use super::table_catalog::{FieldDefinition, TableCatalog};
24use crate::constraints::{ConstraintId, ConstraintKind};
25use crate::metadata::{MetadataManager, MultiColumnUniqueRegistration, SingleColumnIndexEntry};
26use crate::sys_catalog::{
27    ColMeta, MultiColumnIndexEntryMeta, SysCatalog, TriggerEntryMeta, TriggerEventMeta,
28    TriggerTimingMeta,
29};
30use crate::table::Table;
31use crate::types::{FieldId, RowId, TableColumn, TableId};
32use crate::{
33    ForeignKeyColumn, ForeignKeyTableInfo, ForeignKeyView, TableConstraintSummaryView, TableView,
34    ValidatedForeignKey,
35};
36
37/// Result of creating a table. The caller is responsible for wiring executor
38/// caches and any higher-level state that depends on the table schema.
39pub struct CreateTableResult<P>
40where
41    P: Pager<Blob = EntryHandle> + Send + Sync,
42{
43    pub table_id: TableId,
44    pub table: Arc<Table<P>>,
45    pub table_columns: Vec<TableColumn>,
46    pub column_lookup: FxHashMap<String, usize>,
47}
48
49/// Result of attempting to register a single-column index definition.
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub enum SingleColumnIndexRegistration {
52    Created { index_name: String },
53    AlreadyExists { index_name: String },
54}
55
56/// Descriptor for a single-column index resolved from catalog metadata.
57#[derive(Debug, Clone, PartialEq, Eq)]
58pub struct SingleColumnIndexDescriptor {
59    pub index_name: String,
60    pub table_id: TableId,
61    pub canonical_table_name: String,
62    pub display_table_name: String,
63    pub field_id: FieldId,
64    pub column_name: String,
65    pub was_unique: bool,
66}
67
68/// Trait for constructing MVCC columns and Arrow metadata during batch ingestion.
69///
70/// Implementations can delegate to `llkv_transaction::mvcc` or provide custom logic
71/// for synthetic data sources. This indirection avoids coupling the table crate to the
72/// transaction crate while still centralizing MVCC helpers.
73pub trait MvccColumnBuilder: Send + Sync {
74    /// Build MVCC columns (row_id, created_by, deleted_by) for INSERT/CTAS operations.
75    fn build_insert_columns(
76        &self,
77        row_count: usize,
78        start_row_id: RowId,
79        creator_txn_id: u64,
80        deleted_marker: u64,
81    ) -> (ArrayRef, ArrayRef, ArrayRef);
82
83    /// Return the Arrow field definitions for MVCC columns.
84    fn mvcc_fields(&self) -> Vec<Field>;
85
86    /// Construct a user column field with the required metadata assigned.
87    fn field_with_metadata(
88        &self,
89        name: &str,
90        data_type: DataType,
91        nullable: bool,
92        field_id: FieldId,
93    ) -> Field;
94}
95
96/// Service for creating tables.
97///
98/// Coordinates metadata persistence (`MetadataManager`), catalog registration
99/// (`TableCatalog`), and storage initialization (`ColumnStore`).
100#[derive(Clone)]
101pub struct CatalogManager<P>
102where
103    P: Pager<Blob = EntryHandle> + Send + Sync,
104{
105    metadata: Arc<MetadataManager<P>>,
106    catalog: Arc<TableCatalog>,
107    store: Arc<ColumnStore<P>>,
108    type_registry: Arc<std::sync::RwLock<FxHashMap<String, sqlparser::ast::DataType>>>,
109}
110
111impl<P> CatalogManager<P>
112where
113    P: Pager<Blob = EntryHandle> + Send + Sync,
114{
115    /// Creates a new CatalogManager coordinating metadata, catalog, and storage layers.
116    pub fn new(
117        metadata: Arc<MetadataManager<P>>,
118        catalog: Arc<TableCatalog>,
119        store: Arc<ColumnStore<P>>,
120    ) -> Self {
121        Self {
122            metadata,
123            catalog,
124            store,
125            type_registry: Arc::new(std::sync::RwLock::new(FxHashMap::default())),
126        }
127    }
128
129    // ============================================================================
130    // Type Registry Management
131    // ============================================================================
132
133    /// Load custom types from system catalog.
134    /// Should be called during initialization to restore persisted types.
135    pub fn load_types_from_catalog(&self) -> LlkvResult<()> {
136        use crate::sys_catalog::SysCatalog;
137
138        let sys_catalog = SysCatalog::new(&self.store);
139        match sys_catalog.all_custom_type_metas() {
140            Ok(type_metas) => {
141                tracing::debug!(
142                    "[CATALOG] Loaded {} custom type(s) from catalog",
143                    type_metas.len()
144                );
145
146                let mut registry = self.type_registry.write().unwrap();
147                for type_meta in type_metas {
148                    // Parse the base_type_sql back to a DataType
149                    if let Ok(parsed_type) = parse_data_type_from_sql(&type_meta.base_type_sql) {
150                        registry.insert(type_meta.name.to_lowercase(), parsed_type);
151                    } else {
152                        tracing::warn!(
153                            "[CATALOG] Failed to parse base type SQL for type '{}': {}",
154                            type_meta.name,
155                            type_meta.base_type_sql
156                        );
157                    }
158                }
159
160                tracing::debug!(
161                    "[CATALOG] Type registry initialized with {} type(s)",
162                    registry.len()
163                );
164                Ok(())
165            }
166            Err(e) => {
167                tracing::warn!(
168                    "[CATALOG] Failed to load custom types: {}, starting with empty type registry",
169                    e
170                );
171                Ok(()) // Non-fatal, start with empty registry
172            }
173        }
174    }
175
176    /// Register a custom type alias (CREATE TYPE/DOMAIN).
177    pub fn register_type(&self, name: String, data_type: sqlparser::ast::DataType) {
178        let mut registry = self.type_registry.write().unwrap();
179        registry.insert(name.to_lowercase(), data_type);
180    }
181
182    /// Drop a custom type alias (DROP TYPE/DOMAIN).
183    pub fn drop_type(&self, name: &str) -> LlkvResult<()> {
184        let mut registry = self.type_registry.write().unwrap();
185        if registry.remove(&name.to_lowercase()).is_none() {
186            return Err(Error::InvalidArgumentError(format!(
187                "Type '{}' does not exist",
188                name
189            )));
190        }
191        Ok(())
192    }
193
194    /// Resolve a type name to its base DataType, recursively following aliases.
195    pub fn resolve_type(&self, data_type: &sqlparser::ast::DataType) -> sqlparser::ast::DataType {
196        use sqlparser::ast::DataType;
197
198        match data_type {
199            DataType::Custom(obj_name, _) => {
200                let name = obj_name.to_string().to_lowercase();
201                let registry = self.type_registry.read().unwrap();
202                if let Some(base_type) = registry.get(&name) {
203                    // Recursively resolve in case the base type is also an alias
204                    self.resolve_type(base_type)
205                } else {
206                    // Not a custom type, return as-is
207                    data_type.clone()
208                }
209            }
210            // For non-custom types, return as-is
211            _ => data_type.clone(),
212        }
213    }
214
215    // ============================================================================
216    // View Management
217    // ============================================================================
218
219    /// Create a view by storing its SQL definition in the catalog.
220    /// The view will be registered as a table with a view_definition.
221    pub fn create_view(
222        &self,
223        display_name: &str,
224        view_definition: String,
225        column_specs: Vec<PlanColumnSpec>,
226    ) -> LlkvResult<crate::types::TableId> {
227        if column_specs.is_empty() {
228            return Err(Error::InvalidArgumentError(
229                "CREATE VIEW requires at least one column".into(),
230            ));
231        }
232
233        use crate::sys_catalog::TableMeta;
234
235        // Reserve a new table ID for the view
236        let table_id = self.metadata.reserve_table_id()?;
237
238        let created_at_micros = SystemTime::now()
239            .duration_since(UNIX_EPOCH)
240            .unwrap_or_default()
241            .as_micros() as u64;
242
243        // Create the table metadata with view_definition set
244        let table_meta = TableMeta {
245            table_id,
246            name: Some(display_name.to_string()),
247            created_at_micros,
248            flags: 0,
249            epoch: 0,
250            view_definition: Some(view_definition),
251        };
252
253        let mut table_columns = Vec::with_capacity(column_specs.len());
254        for (idx, spec) in column_specs.iter().enumerate() {
255            let field_id = field_id_for_index(idx)?;
256            table_columns.push(TableColumn {
257                field_id,
258                name: spec.name.clone(),
259                data_type: spec.data_type.clone(),
260                nullable: spec.nullable,
261                primary_key: spec.primary_key,
262                unique: spec.unique,
263                check_expr: spec.check_expr.clone(),
264            });
265        }
266
267        // Store the metadata and flush to disk
268        self.metadata.set_table_meta(table_id, table_meta)?;
269        self.metadata
270            .apply_column_definitions(table_id, &table_columns, created_at_micros)?;
271        self.metadata.flush_table(table_id)?;
272
273        // Register the view in the catalog (no namespace prefix - namespacing handled at runtime session layer)
274        self.catalog.register_table(display_name, table_id)?;
275
276        if let Some(field_resolver) = self.catalog.field_resolver(table_id) {
277            for column in &table_columns {
278                let definition = FieldDefinition::new(&column.name)
279                    .with_primary_key(column.primary_key)
280                    .with_unique(column.unique)
281                    .with_check_expr(column.check_expr.clone());
282                if let Err(err) = field_resolver.register_field(definition) {
283                    self.catalog.unregister_table(table_id);
284                    self.metadata.remove_table_state(table_id);
285                    return Err(err);
286                }
287            }
288        }
289
290        tracing::debug!("Created view '{}' with table_id={}", display_name, table_id);
291        Ok(table_id)
292    }
293
294    /// Check if a table is actually a view by looking at its metadata.
295    /// Returns true if the table exists and has a view_definition.
296    pub fn is_view(&self, table_id: crate::types::TableId) -> LlkvResult<bool> {
297        match self.metadata.table_meta(table_id)? {
298            Some(meta) => Ok(meta.view_definition.is_some()),
299            None => Ok(false),
300        }
301    }
302
303    /// Drop a view by removing its metadata and catalog entry.
304    pub fn drop_view(&self, canonical_name: &str, table_id: TableId) -> LlkvResult<()> {
305        let (_, field_ids) = self.sorted_user_fields(table_id);
306        self.metadata.prepare_table_drop(table_id, &field_ids)?;
307        self.metadata.flush_table(table_id)?;
308        self.metadata.remove_table_state(table_id);
309
310        if let Some(table_id_from_catalog) = self.catalog.table_id(canonical_name) {
311            let _ = self.catalog.unregister_table(table_id_from_catalog);
312        } else {
313            let _ = self.catalog.unregister_table(table_id);
314        }
315
316        Ok(())
317    }
318
319    // ============================================================================
320    // Table Creation
321    // ============================================================================
322
323    /// Create a new table using column specifications.
324    ///
325    /// Reserves table ID from metadata, validates columns, persists schema,
326    /// registers in catalog, and returns a Table handle for data operations.
327    pub(crate) fn create_table_from_columns(
328        &self,
329        display_name: &str,
330        canonical_name: &str,
331        columns: &[PlanColumnSpec],
332    ) -> LlkvResult<CreateTableResult<P>> {
333        if columns.is_empty() {
334            return Err(Error::InvalidArgumentError(
335                "CREATE TABLE requires at least one column".into(),
336            ));
337        }
338
339        let mut lookup: FxHashMap<String, usize> =
340            FxHashMap::with_capacity_and_hasher(columns.len(), Default::default());
341        let mut table_columns: Vec<TableColumn> = Vec::with_capacity(columns.len());
342
343        for (idx, column) in columns.iter().enumerate() {
344            let normalized = column.name.to_ascii_lowercase();
345            if lookup.insert(normalized.clone(), idx).is_some() {
346                return Err(Error::InvalidArgumentError(format!(
347                    "duplicate column name '{}' in table '{}'",
348                    column.name, display_name
349                )));
350            }
351
352            table_columns.push(TableColumn {
353                field_id: field_id_for_index(idx)?,
354                name: column.name.clone(),
355                data_type: column.data_type.clone(),
356                nullable: column.nullable,
357                primary_key: column.primary_key,
358                unique: column.unique,
359                check_expr: column.check_expr.clone(),
360            });
361        }
362
363        self.create_table_inner(display_name, canonical_name, table_columns, lookup)
364    }
365
366    /// Create a new table using an Arrow schema (used by CTAS flows).
367    pub fn create_table_from_schema(
368        &self,
369        display_name: &str,
370        canonical_name: &str,
371        schema: &Schema,
372    ) -> LlkvResult<CreateTableResult<P>> {
373        if schema.fields().is_empty() {
374            return Err(Error::InvalidArgumentError(
375                "CREATE TABLE AS SELECT requires at least one column".into(),
376            ));
377        }
378
379        let mut lookup: FxHashMap<String, usize> =
380            FxHashMap::with_capacity_and_hasher(schema.fields().len(), Default::default());
381        let mut table_columns: Vec<TableColumn> = Vec::with_capacity(schema.fields().len());
382
383        for (idx, field) in schema.fields().iter().enumerate() {
384            let data_type = match field.data_type() {
385                DataType::Int64
386                | DataType::Float64
387                | DataType::Utf8
388                | DataType::Date32
389                | DataType::Struct(_) => field.data_type().clone(),
390                other => {
391                    return Err(Error::InvalidArgumentError(format!(
392                        "unsupported column type in CTAS result: {other:?}"
393                    )));
394                }
395            };
396
397            let normalized = field.name().to_ascii_lowercase();
398            if lookup.insert(normalized.clone(), idx).is_some() {
399                return Err(Error::InvalidArgumentError(format!(
400                    "duplicate column name '{}' in CTAS result",
401                    field.name()
402                )));
403            }
404
405            table_columns.push(TableColumn {
406                field_id: field_id_for_index(idx)?,
407                name: field.name().to_string(),
408                data_type,
409                nullable: field.is_nullable(),
410                primary_key: false,
411                unique: false,
412                check_expr: None,
413            });
414        }
415
416        self.create_table_inner(display_name, canonical_name, table_columns, lookup)
417    }
418
419    fn create_table_inner(
420        &self,
421        display_name: &str,
422        _canonical_name: &str,
423        table_columns: Vec<TableColumn>,
424        column_lookup: FxHashMap<String, usize>,
425    ) -> LlkvResult<CreateTableResult<P>> {
426        let table_id = self.metadata.reserve_table_id()?;
427        let timestamp = current_time_micros();
428        let table_meta = crate::sys_catalog::TableMeta {
429            table_id,
430            name: Some(display_name.to_string()),
431            created_at_micros: timestamp,
432            flags: 0,
433            epoch: 0,
434            view_definition: None, // Regular table, not a view
435        };
436
437        self.metadata.set_table_meta(table_id, table_meta)?;
438        self.metadata
439            .apply_column_definitions(table_id, &table_columns, timestamp)?;
440        self.metadata.flush_table(table_id)?;
441
442        let table = Table::from_id_and_store(table_id, Arc::clone(&self.store))?;
443
444        // Register table in catalog using the table_id from metadata
445        tracing::debug!(
446            "[CATALOG_REGISTER] Registering table '{}' (id={}) in catalog @ {:p}",
447            display_name,
448            table_id,
449            &*self.catalog
450        );
451        if let Err(err) = self.catalog.register_table(display_name, table_id) {
452            self.metadata.remove_table_state(table_id);
453            return Err(err);
454        }
455
456        if let Some(field_resolver) = self.catalog.field_resolver(table_id) {
457            for column in &table_columns {
458                let definition = FieldDefinition::new(&column.name)
459                    .with_primary_key(column.primary_key)
460                    .with_unique(column.unique)
461                    .with_check_expr(column.check_expr.clone());
462                if let Err(err) = field_resolver.register_field(definition) {
463                    self.catalog.unregister_table(table_id);
464                    self.metadata.remove_table_state(table_id);
465                    return Err(err);
466                }
467            }
468        }
469
470        Ok(CreateTableResult {
471            table_id,
472            table: Arc::new(table),
473            table_columns,
474            column_lookup,
475        })
476    }
477
478    /// Prepare metadata state and unregister catalog entries for a dropped table.
479    pub fn drop_table(
480        &self,
481        canonical_name: &str,
482        table_id: TableId,
483        column_field_ids: &[FieldId],
484    ) -> LlkvResult<()> {
485        self.metadata
486            .prepare_table_drop(table_id, column_field_ids)?;
487        self.metadata.flush_table(table_id)?;
488        self.metadata.remove_table_state(table_id);
489        if let Some(table_id_from_catalog) = self.catalog.table_id(canonical_name) {
490            let _ = self.catalog.unregister_table(table_id_from_catalog);
491        } else {
492            let _ = self.catalog.unregister_table(table_id);
493        }
494        Ok(())
495    }
496
497    /// Rename a table across metadata and catalog layers.
498    pub fn rename_table(
499        &self,
500        table_id: TableId,
501        current_name: &str,
502        new_name: &str,
503    ) -> LlkvResult<()> {
504        if !current_name.eq_ignore_ascii_case(new_name) && self.catalog.table_id(new_name).is_some()
505        {
506            return Err(Error::CatalogError(format!(
507                "Table '{}' already exists",
508                new_name
509            )));
510        }
511
512        let previous_meta = self.metadata.table_meta(table_id)?;
513        let mut prior_snapshot = None;
514        if let Some(mut meta) = previous_meta.clone() {
515            prior_snapshot = Some(meta.clone());
516            meta.name = Some(new_name.to_string());
517            self.metadata.set_table_meta(table_id, meta)?;
518        }
519
520        if let Err(err) = self.catalog.rename_registered_table(current_name, new_name) {
521            if let Some(prior) = prior_snapshot {
522                let _ = self.metadata.set_table_meta(table_id, prior);
523            }
524            return Err(err);
525        }
526
527        if let Some(prior) = prior_snapshot.clone()
528            && let Err(err) = self.metadata.flush_table(table_id)
529        {
530            let _ = self.metadata.set_table_meta(table_id, prior);
531            let _ = self.catalog.rename_registered_table(new_name, current_name);
532            let _ = self.metadata.flush_table(table_id);
533            return Err(err);
534        }
535
536        Ok(())
537    }
538
539    /// Rename a column in a table by updating its metadata.
540    pub fn rename_column(
541        &self,
542        table_id: TableId,
543        old_column_name: &str,
544        new_column_name: &str,
545    ) -> LlkvResult<()> {
546        // Get all column metas for this table
547        let (_, field_ids) = self.sorted_user_fields(table_id);
548        let column_metas = self.metadata.column_metas(table_id, &field_ids)?;
549
550        // Find the column by old name
551        let mut found_col: Option<(u32, ColMeta)> = None;
552        for (idx, meta_opt) in column_metas.iter().enumerate() {
553            if let Some(meta) = meta_opt
554                && let Some(name) = &meta.name
555                && name.eq_ignore_ascii_case(old_column_name)
556            {
557                found_col = Some((field_ids[idx], meta.clone()));
558                break;
559            }
560        }
561
562        let (_field_id, mut col_meta) = found_col.ok_or_else(|| {
563            Error::InvalidArgumentError(format!("column '{}' not found in table", old_column_name))
564        })?;
565
566        // Update the column name
567        col_meta.name = Some(new_column_name.to_string());
568
569        // Save to catalog
570        let catalog = SysCatalog::new(&self.store);
571        catalog.put_col_meta(table_id, &col_meta);
572
573        // Update metadata manager cache
574        self.metadata.set_column_meta(table_id, col_meta)?;
575
576        // Update field resolver mapping for this column name change.
577        if let Some(resolver) = self.catalog.field_resolver(table_id) {
578            resolver.rename_field(old_column_name, new_column_name)?;
579        }
580
581        self.metadata.flush_table(table_id)?;
582
583        Ok(())
584    }
585
586    /// Alter the data type of a column.
587    ///
588    /// This updates both the column metadata and the storage layer's data type fingerprint.
589    /// Note that actual data conversion is NOT performed - the caller must ensure that
590    /// existing data is compatible with the new type (or that no data exists).
591    ///
592    /// # Arguments
593    /// * `table_id` - The table containing the column
594    /// * `column_name` - Name of the column to alter
595    /// * `new_data_type` - Arrow DataType to set for this column
596    pub fn alter_column_type(
597        &self,
598        table_id: TableId,
599        column_name: &str,
600        new_data_type: &DataType,
601    ) -> LlkvResult<()> {
602        // Get all column metas for this table
603        let (logical_fields, field_ids) = self.sorted_user_fields(table_id);
604        let column_metas = self.metadata.column_metas(table_id, &field_ids)?;
605
606        // Find the column by name
607        let mut found_col: Option<(usize, u32, ColMeta)> = None;
608        for (idx, meta_opt) in column_metas.iter().enumerate() {
609            if let Some(meta) = meta_opt
610                && let Some(name) = &meta.name
611                && name.eq_ignore_ascii_case(column_name)
612            {
613                found_col = Some((idx, field_ids[idx], meta.clone()));
614                break;
615            }
616        }
617
618        let (col_idx, _field_id, col_meta) = found_col.ok_or_else(|| {
619            Error::InvalidArgumentError(format!("column '{}' not found in table", column_name))
620        })?;
621
622        // Update the data type in the storage layer
623        let lfid = logical_fields[col_idx];
624        self.store.update_data_type(lfid, new_data_type)?;
625
626        // Save metadata to catalog
627        let catalog = SysCatalog::new(&self.store);
628        catalog.put_col_meta(table_id, &col_meta);
629
630        // Update metadata manager cache
631        self.metadata.set_column_meta(table_id, col_meta)?;
632
633        Ok(())
634    }
635
636    /// Drop a column from a table by removing its metadata.
637    pub fn drop_column(&self, table_id: TableId, column_name: &str) -> LlkvResult<()> {
638        // Get all column metas for this table
639        let (_, field_ids) = self.sorted_user_fields(table_id);
640        let column_metas = self.metadata.column_metas(table_id, &field_ids)?;
641
642        // Find the column by name
643        let mut found_col_id: Option<u32> = None;
644        for (idx, meta_opt) in column_metas.iter().enumerate() {
645            if let Some(meta) = meta_opt
646                && let Some(name) = &meta.name
647                && name.eq_ignore_ascii_case(column_name)
648            {
649                found_col_id = Some(field_ids[idx]);
650                break;
651            }
652        }
653
654        let col_id = found_col_id.ok_or_else(|| {
655            Error::InvalidArgumentError(format!("column '{}' not found in table", column_name))
656        })?;
657
658        // Delete from catalog
659        let catalog = SysCatalog::new(&self.store);
660        catalog.delete_col_meta(table_id, &[col_id])?;
661
662        Ok(())
663    }
664
665    /// Register a single-column sort (B-tree) index. Optionally marks the field unique.
666    /// Returns `true` if the index was newly created, `false` if it already existed and `if_not_exists` was true.
667    #[allow(clippy::too_many_arguments)]
668    pub fn register_single_column_index(
669        &self,
670        display_name: &str,
671        canonical_name: &str,
672        table: &Table<P>,
673        field_id: FieldId,
674        column_name: &str,
675        index_name: Option<String>,
676        mark_unique: bool,
677        ascending: bool,
678        nulls_first: bool,
679        if_not_exists: bool,
680    ) -> LlkvResult<SingleColumnIndexRegistration> {
681        let table_id = table.table_id();
682        let existing_indexes = table.list_registered_indexes(field_id)?;
683        if existing_indexes.contains(&IndexKind::Sort) {
684            let existing_name = self
685                .metadata
686                .single_column_indexes(table_id)?
687                .into_iter()
688                .find(|entry| entry.column_id == field_id)
689                .map(|entry| entry.index_name)
690                .unwrap_or_else(|| column_name.to_string());
691
692            if if_not_exists {
693                return Ok(SingleColumnIndexRegistration::AlreadyExists {
694                    index_name: existing_name,
695                });
696            }
697
698            return Err(Error::CatalogError(format!(
699                "Index already exists on column '{}' in table '{}'",
700                column_name, display_name
701            )));
702        }
703
704        let index_display_name = match index_name {
705            Some(name) => name,
706            None => {
707                self.generate_single_column_index_name(table_id, canonical_name, column_name)?
708            }
709        };
710        if index_display_name.is_empty() {
711            return Err(Error::InvalidArgumentError(
712                "Index name must not be empty".into(),
713            ));
714        }
715        let canonical_index_name = index_display_name.to_ascii_lowercase();
716
717        if let Some(existing) = self
718            .metadata
719            .single_column_index(table_id, &canonical_index_name)?
720        {
721            if if_not_exists {
722                return Ok(SingleColumnIndexRegistration::AlreadyExists {
723                    index_name: existing.index_name,
724                });
725            }
726
727            return Err(Error::CatalogError(format!(
728                "Index '{}' already exists on table '{}'",
729                existing.index_name, display_name
730            )));
731        }
732
733        let entry = SingleColumnIndexEntry {
734            index_name: index_display_name.clone(),
735            canonical_name: canonical_index_name,
736            column_id: field_id,
737            column_name: column_name.to_string(),
738            unique: mark_unique,
739            ascending,
740            nulls_first,
741        };
742
743        self.metadata.put_single_column_index(table_id, entry)?;
744        self.metadata.register_sort_index(table_id, field_id)?;
745
746        if mark_unique {
747            let catalog_table_id = self.catalog.table_id(canonical_name).unwrap_or(table_id);
748            if let Some(resolver) = self.catalog.field_resolver(catalog_table_id) {
749                resolver.set_field_unique(column_name, true)?;
750            }
751        }
752
753        self.metadata.flush_table(table_id)?;
754
755        Ok(SingleColumnIndexRegistration::Created {
756            index_name: index_display_name,
757        })
758    }
759
760    pub fn drop_single_column_index(
761        &self,
762        plan: DropIndexPlan,
763    ) -> LlkvResult<Option<SingleColumnIndexDescriptor>> {
764        let canonical_index = plan.canonical_name.to_ascii_lowercase();
765        let snapshot = self.catalog.snapshot();
766
767        for canonical_table_name in snapshot.table_names() {
768            let Some(table_id) = snapshot.table_id(&canonical_table_name) else {
769                continue;
770            };
771
772            if let Some(entry) = self
773                .metadata
774                .single_column_index(table_id, &canonical_index)?
775            {
776                self.metadata
777                    .remove_single_column_index(table_id, &canonical_index)?;
778
779                if entry.unique
780                    && let Some(resolver) = self.catalog.field_resolver(table_id)
781                {
782                    resolver.set_field_unique(&entry.column_name, false)?;
783                }
784
785                self.metadata.flush_table(table_id)?;
786
787                let display_table_name = self
788                    .metadata
789                    .table_meta(table_id)?
790                    .and_then(|meta| meta.name)
791                    .unwrap_or_else(|| canonical_table_name.clone());
792
793                return Ok(Some(SingleColumnIndexDescriptor {
794                    index_name: entry.index_name,
795                    table_id,
796                    canonical_table_name,
797                    display_table_name,
798                    field_id: entry.column_id,
799                    column_name: entry.column_name,
800                    was_unique: entry.unique,
801                }));
802            }
803        }
804
805        if plan.if_exists {
806            Ok(None)
807        } else {
808            Err(Error::CatalogError(format!(
809                "Index '{}' does not exist",
810                plan.name
811            )))
812        }
813    }
814
815    /// Register a multi-column UNIQUE index.
816    pub fn register_multi_column_unique_index(
817        &self,
818        table_id: TableId,
819        field_ids: &[FieldId],
820        index_name: Option<String>,
821    ) -> LlkvResult<MultiColumnUniqueRegistration> {
822        let registration = self
823            .metadata
824            .register_multi_column_unique(table_id, field_ids, index_name)?;
825
826        if matches!(registration, MultiColumnUniqueRegistration::Created) {
827            self.metadata.flush_table(table_id)?;
828        }
829
830        Ok(registration)
831    }
832
833    #[allow(clippy::too_many_arguments)]
834    pub fn create_trigger(
835        &self,
836        trigger_display_name: &str,
837        canonical_trigger_name: &str,
838        table_display_name: &str,
839        canonical_table_name: &str,
840        timing: TriggerTimingMeta,
841        event: TriggerEventMeta,
842        for_each_row: bool,
843        condition: Option<String>,
844        body_sql: String,
845        if_not_exists: bool,
846    ) -> LlkvResult<bool> {
847        let Some(table_id) = self.catalog.table_id(canonical_table_name) else {
848            return Err(Error::CatalogError(format!(
849                "Table '{}' does not exist",
850                table_display_name
851            )));
852        };
853
854        let table_meta = self.metadata.table_meta(table_id)?;
855        let is_view = table_meta
856            .as_ref()
857            .and_then(|meta| meta.view_definition.as_ref())
858            .is_some();
859
860        match timing {
861            TriggerTimingMeta::InsteadOf => {
862                if !is_view {
863                    return Err(Error::InvalidArgumentError(format!(
864                        "INSTEAD OF trigger '{}' requires a view target",
865                        trigger_display_name
866                    )));
867                }
868            }
869            _ => {
870                if is_view {
871                    return Err(Error::InvalidArgumentError(format!(
872                        "Trigger '{}' must use INSTEAD OF when targeting a view",
873                        trigger_display_name
874                    )));
875                }
876            }
877        }
878
879        if self
880            .metadata
881            .trigger(table_id, canonical_trigger_name)?
882            .is_some()
883        {
884            if if_not_exists {
885                return Ok(false);
886            }
887            return Err(Error::CatalogError(format!(
888                "Trigger '{}' already exists",
889                trigger_display_name
890            )));
891        }
892
893        let entry = TriggerEntryMeta {
894            name: trigger_display_name.to_string(),
895            canonical_name: canonical_trigger_name.to_string(),
896            timing,
897            event,
898            for_each_row,
899            condition,
900            body_sql,
901        };
902
903        self.metadata.insert_trigger(table_id, entry)?;
904        self.metadata.flush_table(table_id)?;
905        Ok(true)
906    }
907
908    pub fn drop_trigger(
909        &self,
910        trigger_display_name: &str,
911        canonical_trigger_name: &str,
912        table_hint_display: Option<&str>,
913        table_hint_canonical: Option<&str>,
914        if_exists: bool,
915    ) -> LlkvResult<bool> {
916        let mut candidate_tables: Vec<(TableId, String)> = Vec::new();
917
918        if let Some(canonical_table) = table_hint_canonical {
919            match self.catalog.table_id(canonical_table) {
920                Some(table_id) => candidate_tables.push((table_id, canonical_table.to_string())),
921                None => {
922                    if if_exists {
923                        return Ok(false);
924                    }
925                    let display = table_hint_display.unwrap_or(canonical_table);
926                    return Err(Error::CatalogError(format!(
927                        "Table '{}' does not exist",
928                        display
929                    )));
930                }
931            }
932        } else {
933            let snapshot = self.catalog.snapshot();
934            for canonical_table in snapshot.table_names() {
935                if let Some(table_id) = snapshot.table_id(&canonical_table) {
936                    candidate_tables.push((table_id, canonical_table));
937                }
938            }
939        }
940
941        for (table_id, canonical_table) in candidate_tables {
942            if self
943                .metadata
944                .remove_trigger(table_id, canonical_trigger_name)?
945            {
946                self.metadata.flush_table(table_id)?;
947                return Ok(true);
948            } else if table_hint_canonical.is_some()
949                && table_hint_canonical
950                    .unwrap()
951                    .eq_ignore_ascii_case(&canonical_table)
952            {
953                break;
954            }
955        }
956
957        if if_exists {
958            Ok(false)
959        } else {
960            Err(Error::CatalogError(format!(
961                "Trigger '{}' does not exist",
962                trigger_display_name
963            )))
964        }
965    }
966
967    /// Register a multi-column index (unique or non-unique).
968    ///
969    /// Returns true if the index was created, false if it already exists.
970    pub fn register_multi_column_index(
971        &self,
972        table_id: TableId,
973        field_ids: &[FieldId],
974        index_name: String,
975        unique: bool,
976    ) -> LlkvResult<bool> {
977        let canonical_name = index_name.to_lowercase();
978
979        // Check if index already exists
980        if let Some(_existing) = self
981            .metadata
982            .get_multi_column_index(table_id, &canonical_name)?
983        {
984            return Ok(false);
985        }
986
987        // Create new index entry
988        let entry = MultiColumnIndexEntryMeta {
989            index_name: Some(index_name),
990            canonical_name,
991            column_ids: field_ids.to_vec(),
992            unique,
993        };
994
995        self.metadata.put_multi_column_index(table_id, entry)?;
996        self.metadata.flush_table(table_id)?;
997
998        Ok(true)
999    }
1000
1001    fn generate_single_column_index_name(
1002        &self,
1003        table_id: TableId,
1004        canonical_table_name: &str,
1005        column_name: &str,
1006    ) -> LlkvResult<String> {
1007        let table_token = if canonical_table_name.is_empty() {
1008            "table".to_string()
1009        } else {
1010            canonical_table_name.replace('.', "_")
1011        };
1012        let column_token = column_name.to_ascii_lowercase();
1013
1014        let mut candidate = format!("{}_{}_idx", table_token, column_token);
1015        let mut suffix: u32 = 1;
1016        loop {
1017            let canonical = candidate.to_ascii_lowercase();
1018            if self
1019                .metadata
1020                .single_column_index(table_id, &canonical)?
1021                .is_none()
1022            {
1023                return Ok(candidate);
1024            }
1025
1026            candidate = format!("{}_{}_idx{}", table_token, column_token, suffix);
1027            suffix = suffix.checked_add(1).ok_or_else(|| {
1028                Error::InvalidArgumentError("exhausted unique index name generation space".into())
1029            })?;
1030        }
1031    }
1032
1033    /// Append RecordBatches to a freshly created table, injecting MVCC columns.
1034    #[allow(clippy::too_many_arguments)]
1035    pub fn append_batches_with_mvcc(
1036        &self,
1037        table: &Table<P>,
1038        table_columns: &[TableColumn],
1039        batches: &[RecordBatch],
1040        creator_txn_id: u64,
1041        deleted_marker: u64,
1042        starting_row_id: RowId,
1043        mvcc_builder: &dyn MvccColumnBuilder,
1044    ) -> LlkvResult<(RowId, u64)> {
1045        let mut next_row_id = starting_row_id;
1046        let mut total_rows: u64 = 0;
1047
1048        for batch in batches {
1049            if batch.num_rows() == 0 {
1050                continue;
1051            }
1052
1053            if batch.num_columns() != table_columns.len() {
1054                return Err(Error::InvalidArgumentError(format!(
1055                    "CTAS query returned unexpected column count (expected {}, found {})",
1056                    table_columns.len(),
1057                    batch.num_columns()
1058                )));
1059            }
1060
1061            let row_count = batch.num_rows();
1062
1063            let (row_id_array, created_by_array, deleted_by_array) = mvcc_builder
1064                .build_insert_columns(row_count, next_row_id, creator_txn_id, deleted_marker);
1065
1066            let mut arrays: Vec<ArrayRef> = Vec::with_capacity(table_columns.len() + 3);
1067            arrays.push(row_id_array);
1068            arrays.push(created_by_array);
1069            arrays.push(deleted_by_array);
1070
1071            let mut fields = mvcc_builder.mvcc_fields();
1072
1073            for (idx, column) in table_columns.iter().enumerate() {
1074                let array = batch.column(idx).clone();
1075                let field = mvcc_builder.field_with_metadata(
1076                    &column.name,
1077                    column.data_type.clone(),
1078                    column.nullable,
1079                    column.field_id,
1080                );
1081                arrays.push(array);
1082                fields.push(field);
1083            }
1084
1085            let append_schema = Arc::new(Schema::new(fields));
1086            let append_batch = RecordBatch::try_new(append_schema, arrays).map_err(Error::Arrow)?;
1087            table.append(&append_batch)?;
1088
1089            next_row_id = next_row_id.saturating_add(row_count as u64);
1090            total_rows = total_rows.saturating_add(row_count as u64);
1091        }
1092
1093        Ok((next_row_id, total_rows))
1094    }
1095
1096    /// Validate and register foreign keys for a newly created table.
1097    #[allow(clippy::too_many_arguments)]
1098    pub fn register_foreign_keys_for_new_table<F>(
1099        &self,
1100        table_id: TableId,
1101        display_name: &str,
1102        canonical_name: &str,
1103        table_columns: &[TableColumn],
1104        specs: &[ForeignKeySpec],
1105        lookup_table: F,
1106        timestamp_micros: u64,
1107    ) -> LlkvResult<Vec<ValidatedForeignKey>>
1108    where
1109        F: FnMut(&str) -> LlkvResult<ForeignKeyTableInfo>,
1110    {
1111        if specs.is_empty() {
1112            return Ok(Vec::new());
1113        }
1114
1115        let referencing_columns: Vec<ForeignKeyColumn> = table_columns
1116            .iter()
1117            .map(|column| ForeignKeyColumn {
1118                name: column.name.clone(),
1119                data_type: column.data_type.clone(),
1120                nullable: column.nullable,
1121                primary_key: column.primary_key,
1122                unique: column.unique,
1123                field_id: column.field_id,
1124            })
1125            .collect();
1126
1127        let multi_column_uniques = {
1128            let catalog = SysCatalog::new(&self.store);
1129            let all_indexes = catalog.get_multi_column_indexes(table_id)?;
1130            all_indexes.into_iter().filter(|idx| idx.unique).collect()
1131        };
1132
1133        let referencing_table = ForeignKeyTableInfo {
1134            display_name: display_name.to_string(),
1135            canonical_name: canonical_name.to_string(),
1136            table_id,
1137            columns: referencing_columns,
1138            multi_column_uniques,
1139        };
1140
1141        self.metadata.validate_and_register_foreign_keys(
1142            &referencing_table,
1143            specs,
1144            lookup_table,
1145            timestamp_micros,
1146        )
1147    }
1148
1149    /// Resolve referenced tables for the provided foreign key view definitions.
1150    pub fn referenced_table_info(
1151        &self,
1152        views: &[ForeignKeyView],
1153    ) -> LlkvResult<Vec<ForeignKeyTableInfo>> {
1154        let mut results = Vec::with_capacity(views.len());
1155        for view in views {
1156            let Some(table_id) = self.catalog.table_id(&view.referenced_table_canonical) else {
1157                return Err(Error::CatalogError(format!(
1158                    "Catalog Error: referenced table '{}' does not exist",
1159                    view.referenced_table_display
1160                )));
1161            };
1162
1163            let Some(resolver) = self.catalog.field_resolver(table_id) else {
1164                return Err(Error::Internal(format!(
1165                    "catalog resolver missing for table '{}'",
1166                    view.referenced_table_display
1167                )));
1168            };
1169
1170            let mut columns = Vec::with_capacity(view.referenced_field_ids.len());
1171            for field_id in &view.referenced_field_ids {
1172                let info = resolver.field_info(*field_id).ok_or_else(|| {
1173                    Error::Internal(format!(
1174                        "field metadata missing for id {} in table '{}'",
1175                        field_id, view.referenced_table_display
1176                    ))
1177                })?;
1178
1179                let data_type = self.metadata.column_data_type(table_id, *field_id)?;
1180
1181                columns.push(ForeignKeyColumn {
1182                    name: info.display_name.to_string(),
1183                    data_type,
1184                    nullable: !info.constraints.primary_key,
1185                    primary_key: info.constraints.primary_key,
1186                    unique: info.constraints.unique,
1187                    field_id: *field_id,
1188                });
1189            }
1190
1191            let multi_column_uniques = {
1192                let catalog = SysCatalog::new(&self.store);
1193                let all_indexes = catalog.get_multi_column_indexes(table_id)?;
1194                all_indexes.into_iter().filter(|idx| idx.unique).collect()
1195            };
1196
1197            results.push(ForeignKeyTableInfo {
1198                display_name: view.referenced_table_display.clone(),
1199                canonical_name: view.referenced_table_canonical.clone(),
1200                table_id,
1201                columns,
1202                multi_column_uniques,
1203            });
1204        }
1205
1206        Ok(results)
1207    }
1208
1209    /// Return the current metadata snapshot for a table, including column metadata and constraints.
1210    pub fn table_view(&self, canonical_name: &str) -> LlkvResult<TableView> {
1211        let table_id = self.catalog.table_id(canonical_name).ok_or_else(|| {
1212            Error::InvalidArgumentError(format!("unknown table '{}'", canonical_name))
1213        })?;
1214
1215        let (_, field_ids) = self.sorted_user_fields(table_id);
1216        self.table_view_with_field_ids(table_id, &field_ids)
1217    }
1218
1219    /// Produce a read-only view of a table's catalog, including column metadata and constraints.
1220    pub fn table_column_specs(&self, canonical_name: &str) -> LlkvResult<Vec<PlanColumnSpec>> {
1221        let table_id = self.catalog.table_id(canonical_name).ok_or_else(|| {
1222            Error::InvalidArgumentError(format!("unknown table '{}'", canonical_name))
1223        })?;
1224
1225        let resolver = self
1226            .catalog
1227            .field_resolver(table_id)
1228            .ok_or_else(|| Error::Internal("missing field resolver for table".into()))?;
1229
1230        let (logical_fields, field_ids) = self.sorted_user_fields(table_id);
1231
1232        let table_view = self.table_view_with_field_ids(table_id, &field_ids)?;
1233        let column_metas = table_view.column_metas;
1234        let constraint_records = table_view.constraint_records;
1235
1236        let mut metadata_primary_keys: FxHashSet<FieldId> = FxHashSet::default();
1237        let mut metadata_unique_fields: FxHashSet<FieldId> = FxHashSet::default();
1238        let mut has_primary_key_records = false;
1239        let mut has_single_unique_records = false;
1240
1241        for record in constraint_records
1242            .iter()
1243            .filter(|record| record.is_active())
1244        {
1245            match &record.kind {
1246                ConstraintKind::PrimaryKey(pk) => {
1247                    has_primary_key_records = true;
1248                    for field_id in &pk.field_ids {
1249                        metadata_primary_keys.insert(*field_id);
1250                        metadata_unique_fields.insert(*field_id);
1251                    }
1252                }
1253                ConstraintKind::Unique(unique) => {
1254                    if unique.field_ids.len() == 1 {
1255                        has_single_unique_records = true;
1256                        if let Some(field_id) = unique.field_ids.first() {
1257                            metadata_unique_fields.insert(*field_id);
1258                        }
1259                    }
1260                }
1261                _ => {}
1262            }
1263        }
1264
1265        let mut specs = Vec::with_capacity(field_ids.len());
1266
1267        for (idx, lfid) in logical_fields.iter().enumerate() {
1268            let field_id = lfid.field_id();
1269
1270            let column_name = column_metas
1271                .get(idx)
1272                .and_then(|meta| meta.as_ref())
1273                .and_then(|meta| meta.name.clone())
1274                .unwrap_or_else(|| format!("col_{}", field_id));
1275
1276            let fallback_constraints = resolver
1277                .field_constraints_by_name(&column_name)
1278                .unwrap_or_default();
1279
1280            let metadata_primary = metadata_primary_keys.contains(&field_id);
1281            let primary_key = if has_primary_key_records {
1282                metadata_primary
1283            } else {
1284                fallback_constraints.primary_key
1285            };
1286
1287            let metadata_unique = metadata_primary || metadata_unique_fields.contains(&field_id);
1288            let unique = if has_primary_key_records || has_single_unique_records {
1289                metadata_unique
1290            } else {
1291                fallback_constraints.primary_key || fallback_constraints.unique
1292            };
1293
1294            let data_type = self.store.data_type(*lfid)?;
1295            let nullable = !primary_key;
1296
1297            let mut spec = PlanColumnSpec::new(column_name.clone(), data_type, nullable)
1298                .with_primary_key(primary_key)
1299                .with_unique(unique);
1300
1301            if let Some(check_expr) = fallback_constraints.check_expr.clone() {
1302                spec = spec.with_check(Some(check_expr));
1303            }
1304
1305            specs.push(spec);
1306        }
1307
1308        Ok(specs)
1309    }
1310
1311    /// Return the foreign key metadata for the specified table.
1312    pub fn foreign_key_views(&self, canonical_name: &str) -> LlkvResult<Vec<ForeignKeyView>> {
1313        let table_id = self.catalog.table_id(canonical_name).ok_or_else(|| {
1314            Error::InvalidArgumentError(format!("unknown table '{}'", canonical_name))
1315        })?;
1316
1317        self.metadata.foreign_key_views(&self.catalog, table_id)
1318    }
1319
1320    /// Return constraint-related catalog metadata for the specified table.
1321    pub fn table_constraint_summary(
1322        &self,
1323        canonical_name: &str,
1324    ) -> LlkvResult<TableConstraintSummaryView> {
1325        tracing::trace!(
1326            "[TABLE_CONSTRAINT_SUMMARY] Looking up table '{}' in catalog @ {:p}",
1327            canonical_name,
1328            &*self.catalog
1329        );
1330        let table_id = self.catalog.table_id(canonical_name).ok_or_else(|| {
1331            tracing::error!(
1332                "[TABLE_CONSTRAINT_SUMMARY] Table '{}' NOT FOUND in catalog @ {:p}",
1333                canonical_name,
1334                &*self.catalog
1335            );
1336            Error::InvalidArgumentError(format!("unknown table '{}'", canonical_name))
1337        })?;
1338        tracing::trace!(
1339            "[TABLE_CONSTRAINT_SUMMARY] Found table '{}' with id={} in catalog",
1340            canonical_name,
1341            table_id
1342        );
1343
1344        let (_, field_ids) = self.sorted_user_fields(table_id);
1345        let table_meta = self.metadata.table_meta(table_id)?;
1346        let column_metas = self.metadata.column_metas(table_id, &field_ids)?;
1347        let constraint_records = self.metadata.constraint_records(table_id)?;
1348        let multi_column_uniques = self.metadata.multi_column_uniques(table_id)?;
1349
1350        Ok(TableConstraintSummaryView {
1351            table_meta,
1352            column_metas,
1353            constraint_records,
1354            multi_column_uniques,
1355        })
1356    }
1357
1358    fn sorted_user_fields(
1359        &self,
1360        table_id: TableId,
1361    ) -> (Vec<llkv_column_map::types::LogicalFieldId>, Vec<FieldId>) {
1362        let mut logical_fields = self.store.user_field_ids_for_table(table_id);
1363        logical_fields.sort_by_key(|lfid| lfid.field_id());
1364        let field_ids = logical_fields
1365            .iter()
1366            .map(|lfid| lfid.field_id())
1367            .collect::<Vec<_>>();
1368
1369        (logical_fields, field_ids)
1370    }
1371
1372    fn table_view_with_field_ids(
1373        &self,
1374        table_id: TableId,
1375        field_ids: &[FieldId],
1376    ) -> LlkvResult<TableView> {
1377        self.metadata.table_view(&self.catalog, table_id, field_ids)
1378    }
1379
1380    // -------------------------------------------------------------------------
1381    // Catalog read-only views
1382    // -------------------------------------------------------------------------
1383
1384    /// Returns all table names in the catalog.
1385    pub fn table_names(&self) -> Vec<String> {
1386        self.catalog.table_names()
1387    }
1388
1389    /// Returns the TableId for a canonical table name.
1390    pub fn table_id(&self, canonical_name: &str) -> Option<TableId> {
1391        self.catalog.table_id(canonical_name)
1392    }
1393
1394    /// Returns a field resolver for the given table.
1395    pub fn field_resolver(&self, table_id: TableId) -> Option<crate::catalog::FieldResolver> {
1396        self.catalog.field_resolver(table_id)
1397    }
1398
1399    /// Returns a snapshot of the catalog for read-only access.
1400    pub fn catalog_snapshot(&self) -> crate::catalog::TableCatalogSnapshot {
1401        self.catalog.snapshot()
1402    }
1403
1404    /// Returns a reference to the internal catalog for services that need it.
1405    /// Note: This is primarily for internal use by services like ConstraintService.
1406    pub fn catalog(&self) -> &Arc<TableCatalog> {
1407        &self.catalog
1408    }
1409
1410    /// Returns all foreign keys that reference the specified table.
1411    /// Returns a vector of (referencing_table_id, constraint_id) pairs.
1412    pub fn foreign_keys_referencing(
1413        &self,
1414        referenced_table_id: TableId,
1415    ) -> LlkvResult<Vec<(TableId, ConstraintId)>> {
1416        self.metadata.foreign_keys_referencing(referenced_table_id)
1417    }
1418
1419    /// Returns detailed foreign key views for a specific table.
1420    /// This includes foreign keys where the specified table is the referencing table.
1421    pub fn foreign_key_views_for_table(
1422        &self,
1423        table_id: TableId,
1424    ) -> LlkvResult<Vec<ForeignKeyView>> {
1425        self.metadata.foreign_key_views(&self.catalog, table_id)
1426    }
1427}
1428
1429fn field_id_for_index(idx: usize) -> LlkvResult<FieldId> {
1430    FieldId::try_from(idx + 1).map_err(|_| {
1431        Error::Internal(format!(
1432            "column index {} exceeded supported field id range",
1433            idx + 1
1434        ))
1435    })
1436}
1437
1438// TODO: Dedupe (another instance exists in llkv-executor)
1439#[allow(clippy::unnecessary_wraps)]
1440fn current_time_micros() -> u64 {
1441    SystemTime::now()
1442        .duration_since(UNIX_EPOCH)
1443        .map(|duration| duration.as_micros() as u64)
1444        .unwrap_or(0)
1445}
1446
1447/// Parse a SQL type string (e.g., "INTEGER") back into a DataType.
1448fn parse_data_type_from_sql(sql: &str) -> LlkvResult<sqlparser::ast::DataType> {
1449    use sqlparser::dialect::GenericDialect;
1450    use sqlparser::parser::Parser;
1451
1452    // Try to parse as a simple CREATE DOMAIN statement
1453    let create_sql = format!("CREATE DOMAIN dummy AS {}", sql);
1454    let dialect = GenericDialect {};
1455
1456    match Parser::parse_sql(&dialect, &create_sql) {
1457        Ok(stmts) if !stmts.is_empty() => {
1458            if let sqlparser::ast::Statement::CreateDomain(create_domain) = &stmts[0] {
1459                Ok(create_domain.data_type.clone())
1460            } else {
1461                Err(Error::InvalidArgumentError(format!(
1462                    "Failed to parse type from SQL: {}",
1463                    sql
1464                )))
1465            }
1466        }
1467        _ => Err(Error::InvalidArgumentError(format!(
1468            "Failed to parse type from SQL: {}",
1469            sql
1470        ))),
1471    }
1472}