1#![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
37pub 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#[derive(Debug, Clone, PartialEq, Eq)]
51pub enum SingleColumnIndexRegistration {
52 Created { index_name: String },
53 AlreadyExists { index_name: String },
54}
55
56#[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
68pub trait MvccColumnBuilder: Send + Sync {
74 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 fn mvcc_fields(&self) -> Vec<Field>;
85
86 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#[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 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 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 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(()) }
173 }
174 }
175
176 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 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 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 self.resolve_type(base_type)
205 } else {
206 data_type.clone()
208 }
209 }
210 _ => data_type.clone(),
212 }
213 }
214
215 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 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 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 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 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 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 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 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 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, };
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 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 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 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 pub fn rename_column(
541 &self,
542 table_id: TableId,
543 old_column_name: &str,
544 new_column_name: &str,
545 ) -> LlkvResult<()> {
546 let (_, field_ids) = self.sorted_user_fields(table_id);
548 let column_metas = self.metadata.column_metas(table_id, &field_ids)?;
549
550 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 col_meta.name = Some(new_column_name.to_string());
568
569 let catalog = SysCatalog::new(&self.store);
571 catalog.put_col_meta(table_id, &col_meta);
572
573 self.metadata.set_column_meta(table_id, col_meta)?;
575
576 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 pub fn alter_column_type(
597 &self,
598 table_id: TableId,
599 column_name: &str,
600 new_data_type: &DataType,
601 ) -> LlkvResult<()> {
602 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 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 let lfid = logical_fields[col_idx];
624 self.store.update_data_type(lfid, new_data_type)?;
625
626 let catalog = SysCatalog::new(&self.store);
628 catalog.put_col_meta(table_id, &col_meta);
629
630 self.metadata.set_column_meta(table_id, col_meta)?;
632
633 Ok(())
634 }
635
636 pub fn drop_column(&self, table_id: TableId, column_name: &str) -> LlkvResult<()> {
638 let (_, field_ids) = self.sorted_user_fields(table_id);
640 let column_metas = self.metadata.column_metas(table_id, &field_ids)?;
641
642 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 let catalog = SysCatalog::new(&self.store);
660 catalog.delete_col_meta(table_id, &[col_id])?;
661
662 Ok(())
663 }
664
665 #[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 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 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 if let Some(_existing) = self
981 .metadata
982 .get_multi_column_index(table_id, &canonical_name)?
983 {
984 return Ok(false);
985 }
986
987 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 #[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 #[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 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 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 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 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 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 pub fn table_names(&self) -> Vec<String> {
1386 self.catalog.table_names()
1387 }
1388
1389 pub fn table_id(&self, canonical_name: &str) -> Option<TableId> {
1391 self.catalog.table_id(canonical_name)
1392 }
1393
1394 pub fn field_resolver(&self, table_id: TableId) -> Option<crate::catalog::FieldResolver> {
1396 self.catalog.field_resolver(table_id)
1397 }
1398
1399 pub fn catalog_snapshot(&self) -> crate::catalog::TableCatalogSnapshot {
1401 self.catalog.snapshot()
1402 }
1403
1404 pub fn catalog(&self) -> &Arc<TableCatalog> {
1407 &self.catalog
1408 }
1409
1410 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 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#[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
1447fn parse_data_type_from_sql(sql: &str) -> LlkvResult<sqlparser::ast::DataType> {
1449 use sqlparser::dialect::GenericDialect;
1450 use sqlparser::parser::Parser;
1451
1452 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}