Skip to main content

akar_processor/processor/mapper/
map_ddl.rs

1use super::ExecutionContext;
2use crate::physical_operator::*;
3use crate::processor::SchemaDdlOp;
4use crate::processor::plan_serializer::serialize_plan_tree;
5use akar_common::error::ProcessorError;
6use akar_common::vector::DataChunk;
7use akar_planner::logical_operator::LogicalOperator;
8
9pub fn map_and_execute_ddl(
10    op: &LogicalOperator,
11    current_input: Vec<DataChunk>,
12    ctx: &mut ExecutionContext,
13) -> Result<Vec<DataChunk>, ProcessorError> {
14    match op {
15        LogicalOperator::Explain(ex) => {
16            // Serialize the inner plan tree to a string
17            let plan_str = serialize_plan_tree(&ex.inner, 0);
18            let explain = PhysicalExplain { inner_plan: plan_str };
19            let result = explain.execute(vec![])?;
20            Ok(result)
21        }
22        LogicalOperator::StandaloneCall(c) => {
23            if let Some(ref handler) = ctx.standalone_call_handler {
24                let result = handler.execute_call(&c.function_name, &c.args)?;
25                Ok(result)
26            } else {
27                Err(format!("No standalone call handler available to execute '{}'", c.function_name).into())
28            }
29        }
30        LogicalOperator::TableFunctionCall(tf) => {
31            let result = ctx.processor.execute_table_function(tf)?;
32            Ok(result)
33        }
34        LogicalOperator::Foreach(fc) => {
35            let foreach_op = PhysicalForeach {
36                variable: fc.variable.clone(),
37                expression: fc.expression.clone(),
38                sub_plans: fc.sub_plans.clone(),
39                function_registry: ctx.function_registry.clone(),
40                table_catalog: ctx.table_catalog.clone(),
41                vfs: ctx.vfs.clone(),
42            };
43            let result = foreach_op.execute(current_input)?;
44            Ok(result)
45        }
46        LogicalOperator::CreateNodeTable(c) => {
47            let tc = ctx
48                .table_catalog
49                .as_ref()
50                .ok_or("CREATE NODE TABLE requires a table catalog")?;
51            let columns: Vec<akar_storage::table::ColumnDefinition> = c
52                .columns
53                .iter()
54                .map(|col| akar_storage::table::ColumnDefinition {
55                    name: col.name.clone(),
56                    logical_type: col.logical_type,
57                    is_primary_key: col.is_primary_key,
58                    compression: col.compression,
59                })
60                .collect();
61            tc.create_node_table(c.name.clone(), columns);
62
63            // Auto-create ART index for primary key (matches connection/ddl.rs behavior)
64            if c.columns.iter().any(|col| col.is_primary_key) {
65                let index_name = format!("{}_pk_idx", c.name);
66                tc.create_art_index(&c.name, &index_name)
67                    .map_err(|e| format!("Failed to auto-create ART PK index for table '{}': {e}", c.name))?;
68            }
69
70            tracing::info!("Pipeline: Created node table '{}'", c.name);
71            Ok(ddl_success_chunk(&format!("Node table '{}' created", c.name)))
72        }
73        LogicalOperator::CreateRelTable(c) => {
74            let tc = ctx
75                .table_catalog
76                .as_ref()
77                .ok_or("CREATE REL TABLE requires a table catalog")?;
78            let from_id = tc
79                .get_node_table_by_name(&c.from)
80                .map(|t| t.table_id)
81                .ok_or_else(|| format!("From table '{}' not found", c.from))?;
82            let to_id = tc
83                .get_node_table_by_name(&c.to)
84                .map(|t| t.table_id)
85                .ok_or_else(|| format!("To table '{}' not found", c.to))?;
86            let columns: Vec<akar_storage::table::ColumnDefinition> = c
87                .columns
88                .iter()
89                .map(|col| akar_storage::table::ColumnDefinition {
90                    name: col.name.clone(),
91                    logical_type: col.logical_type,
92                    is_primary_key: col.is_primary_key,
93                    compression: col.compression,
94                })
95                .collect();
96            tc.create_rel_table(c.name.clone(), from_id, to_id, columns);
97            tracing::info!("Pipeline: Created rel table '{}' ({} -> {})", c.name, c.from, c.to);
98            Ok(ddl_success_chunk(&format!("Rel table '{}' created", c.name)))
99        }
100        LogicalOperator::DropTable(d) => {
101            let tc = ctx
102                .table_catalog
103                .as_ref()
104                .ok_or("DROP TABLE requires a table catalog")?;
105            let dropped = tc.drop_node_table(&d.name) || tc.drop_rel_table(&d.name);
106            if dropped {
107                tracing::info!("Pipeline: Dropped table '{}'", d.name);
108                Ok(ddl_success_chunk(&format!("Table '{}' dropped", d.name)))
109            } else {
110                Err(format!("Table '{}' not found", d.name).into())
111            }
112        }
113        LogicalOperator::AlterTable(a) => {
114            let tc = ctx
115                .table_catalog
116                .as_ref()
117                .ok_or("ALTER TABLE requires a table catalog")?;
118            match &a.action {
119                akar_parser::ast::AlterAction::AddColumn { name, type_name } => {
120                    let logical_type = parse_type_simple(type_name)?;
121                    let mut table = tc
122                        .get_node_table_by_name_mut(&a.table_name)
123                        .ok_or_else(|| format!("Table '{}' not found", a.table_name))?;
124                    if table.columns.iter().any(|c| c.name.eq_ignore_ascii_case(name)) {
125                        return Err(format!("Column '{}' already exists in '{}'", name, a.table_name).into());
126                    }
127                    table.columns.push(akar_storage::table::ColumnDefinition {
128                        name: name.clone(),
129                        logical_type,
130                        is_primary_key: false,
131                        compression: akar_common::enums::CompressionType::Uncompressed,
132                    });
133                    tracing::info!("Pipeline: Added column '{}' to '{}'", name, a.table_name);
134                    Ok(ddl_success_chunk(&format!(
135                        "Column '{}' added to table '{}'",
136                        name, a.table_name
137                    )))
138                }
139                akar_parser::ast::AlterAction::DropColumn { name } => {
140                    let mut table = tc
141                        .get_node_table_by_name_mut(&a.table_name)
142                        .ok_or_else(|| format!("Table '{}' not found", a.table_name))?;
143                    let pos = table
144                        .columns
145                        .iter()
146                        .position(|c| c.name == *name)
147                        .ok_or_else(|| format!("Column '{}' not found in '{}'", name, a.table_name))?;
148                    if table.columns[pos].is_primary_key {
149                        return Err(format!("Cannot drop primary key column '{}'", name).into());
150                    }
151                    table.columns.remove(pos);
152                    tracing::info!("Pipeline: Dropped column '{}' from '{}'", name, a.table_name);
153                    Ok(ddl_success_chunk(&format!(
154                        "Column '{}' dropped from table '{}'",
155                        name, a.table_name
156                    )))
157                }
158                akar_parser::ast::AlterAction::RenameColumn { old_name, new_name } => {
159                    {
160                        let table = tc
161                            .get_node_table_by_name(&a.table_name)
162                            .ok_or_else(|| format!("Table '{}' not found", a.table_name))?;
163                        if !table.columns.iter().any(|c| c.name == *old_name) {
164                            return Err(format!("Column '{}' not found in '{}'", old_name, a.table_name).into());
165                        }
166                        if table.columns.iter().any(|c| c.name == *new_name) {
167                            return Err(format!("Column '{}' already exists in '{}'", new_name, a.table_name).into());
168                        }
169                    }
170                    let mut table = tc.get_node_table_by_name_mut(&a.table_name).unwrap();
171                    let col = table.columns.iter_mut().find(|c| c.name == *old_name).unwrap();
172                    col.name = new_name.clone();
173                    tracing::info!(
174                        "Pipeline: Renamed column '{}' to '{}' in '{}'",
175                        old_name,
176                        new_name,
177                        a.table_name
178                    );
179                    Ok(ddl_success_chunk(&format!(
180                        "Column '{}' renamed to '{}' in table '{}'",
181                        old_name, new_name, a.table_name
182                    )))
183                }
184                akar_parser::ast::AlterAction::RenameTable { new_name } => {
185                    if tc.get_node_table_by_name(new_name).is_some() || tc.get_rel_table_by_name(new_name).is_some() {
186                        return Err(format!("Table '{}' already exists", new_name).into());
187                    }
188                    if let Some(mut table) = tc.get_node_table_by_name_mut(&a.table_name) {
189                        table.name = new_name.clone();
190                    } else if let Some(mut table) = tc.get_rel_table_by_name_mut(&a.table_name) {
191                        table.name = new_name.clone();
192                    } else {
193                        return Err(format!("Table '{}' not found", a.table_name).into());
194                    }
195                    tracing::info!("Pipeline: Renamed table '{}' to '{}'", a.table_name, new_name);
196                    Ok(ddl_success_chunk(&format!(
197                        "Table '{}' renamed to '{}'",
198                        a.table_name, new_name
199                    )))
200                }
201            }
202        }
203        LogicalOperator::CreateIndex(idx) => {
204            let tc = ctx
205                .table_catalog
206                .as_ref()
207                .ok_or("CREATE INDEX requires a table catalog")?;
208            tc.create_art_index(&idx.table_name, &idx.index_name)?;
209            tracing::info!(
210                "Pipeline: Created ART index '{}' on '{}'",
211                idx.index_name,
212                idx.table_name
213            );
214            Ok(ddl_success_chunk(&format!(
215                "ART index '{}' created on table '{}'",
216                idx.index_name, idx.table_name
217            )))
218        }
219        LogicalOperator::DropIndex(idx) => {
220            let tc = ctx
221                .table_catalog
222                .as_ref()
223                .ok_or("DROP INDEX requires a table catalog")?;
224            tc.drop_art_index(&idx.table_name)?;
225            tracing::info!("Pipeline: Dropped index '{}' from '{}'", idx.index_name, idx.table_name);
226            Ok(ddl_success_chunk(&format!(
227                "Index '{}' dropped from table '{}'",
228                idx.index_name, idx.table_name
229            )))
230        }
231        LogicalOperator::CreateVectorIndex(vi) => {
232            let tc = ctx
233                .table_catalog
234                .as_ref()
235                .ok_or("CREATE VECTOR INDEX requires a table catalog")?;
236
237            // Map string metric to typed enum
238            let metric = match vi.metric.to_lowercase().as_str() {
239                "cosine" => akar_vector::hnsw::DistanceMetric::Cosine,
240                "euclidean" | "l2" => akar_vector::hnsw::DistanceMetric::L2Squared,
241                "dot" => akar_vector::hnsw::DistanceMetric::DotProduct,
242                other => return Err(format!("Unknown vector metric '{other}'").into()),
243            };
244
245            // Create the vector index in storage
246            tc.create_vector_index(
247                vi.index_name.clone(),
248                vi.table_name.clone(),
249                vi.column_name.clone(),
250                metric,
251                vi.dimensions as u32,
252            );
253
254            // Auto-populate from existing table data
255            if let Some(table) = tc.get_node_table_by_name(&vi.table_name) {
256                let col_idx = table.columns.iter().position(|c| c.name == vi.column_name);
257                if let Some(col_idx) = col_idx {
258                    for row_id in 0..table.num_rows as usize {
259                        if let Some(val) = table.get_value(row_id, col_idx) {
260                            if let Ok(vec) = akar_storage::extract_f64_list_from_value(val) {
261                                if let Some(mut vib) = tc.get_vector_index_by_name_mut(&vi.index_name) {
262                                    vib.hnsw_mut().insert(vec, row_id);
263                                }
264                            }
265                        }
266                    }
267                }
268            }
269
270            tracing::info!(
271                "Pipeline: Created vector index '{}' on '{}.{}'",
272                vi.index_name,
273                vi.table_name,
274                vi.column_name
275            );
276            Ok(ddl_success_chunk(&format!(
277                "Vector index '{}' created on '{}.{}'",
278                vi.index_name, vi.table_name, vi.column_name
279            )))
280        }
281        LogicalOperator::CreateSequence(s) => {
282            if let Some(ref ddl_fn) = ctx.schema_ddl_fn {
283                let result = ddl_fn(SchemaDdlOp::CreateSequence {
284                    name: s.name.clone(),
285                    if_not_exists: s.if_not_exists,
286                    start_value: s.start_with,
287                    increment: s.increment,
288                    min_value: s.min_value,
289                    max_value: s.max_value,
290                    cycle: s.cycle,
291                })?;
292                Ok(ddl_success_chunk(&result))
293            } else {
294                Err("CREATE SEQUENCE requires schema catalog access".into())
295            }
296        }
297        LogicalOperator::DropSequence(s) => {
298            if let Some(ref ddl_fn) = ctx.schema_ddl_fn {
299                let result = ddl_fn(SchemaDdlOp::DropSequence {
300                    name: s.name.clone(),
301                    if_exists: s.if_exists,
302                })?;
303                Ok(ddl_success_chunk(&result))
304            } else {
305                Err("DROP SEQUENCE requires schema catalog access".into())
306            }
307        }
308        LogicalOperator::CreateDml(c) => {
309            let tc = ctx
310                .table_catalog
311                .as_ref()
312                .ok_or("CREATE DML requires a table catalog")?;
313            let mut table = tc
314                .get_node_table_by_name_mut(&c.table_name)
315                .ok_or_else(|| format!("Table '{}' not found", c.table_name))?;
316
317            // Build values from pattern properties, defaulting to Null
318            let mut values: Vec<akar_common::types::Value> =
319                table.columns.iter().map(|_| akar_common::types::Value::Null).collect();
320            {
321                let registry = ctx
322                    .function_registry
323                    .clone()
324                    .ok_or("CREATE DML requires a function registry")?;
325                let registry = registry.lock().map_err(|e| format!("Lock poisoned: {e}"))?;
326                for (prop_name, expr) in &c.properties {
327                    if let Some(col_idx) = table.columns.iter().position(|col| col.name == *prop_name) {
328                        values[col_idx] = crate::physical::write_ops::set::evaluate_constant_expr(expr, &registry);
329                    }
330                }
331            }
332
333            table.insert_row(values)?;
334            tracing::info!("Pipeline: Created node in '{}'", c.table_name);
335            Ok(ddl_success_chunk(&format!("Created node in '{}'", c.table_name)))
336        }
337        LogicalOperator::ExportDatabase(e) => {
338            if let Some(ref ddl_fn) = ctx.schema_ddl_fn {
339                let result = ddl_fn(SchemaDdlOp::ExportDatabase {
340                    file_path: e.file_path.clone(),
341                    file_type: e.file_type.clone(),
342                    schema_only: e.schema_only,
343                })?;
344                Ok(ddl_success_chunk(&result))
345            } else {
346                Err("EXPORT DATABASE requires schema catalog access".into())
347            }
348        }
349        LogicalOperator::ImportDatabase(i) => {
350            if let Some(ref ddl_fn) = ctx.schema_ddl_fn {
351                let result = ddl_fn(SchemaDdlOp::ImportDatabase {
352                    file_path: i.file_path.clone(),
353                    query: i.query.clone(),
354                    index_query: i.index_query.clone(),
355                })?;
356                Ok(ddl_success_chunk(&result))
357            } else {
358                Err("IMPORT DATABASE requires schema catalog access".into())
359            }
360        }
361        LogicalOperator::CreateFtsIndex(c) => {
362            if let Some(ref tc) = ctx.table_catalog {
363                let fts_index = PhysicalCreateFtsIndex {
364                    index_name: c.index_name.clone(),
365                    table_name: c.table_name.clone(),
366                    column_name: c.column_name.clone(),
367                    docs_table: c.docs_table.clone(),
368                    terms_table: c.terms_table.clone(),
369                    posting_table: c.posting_table.clone(),
370                    table_catalog: tc.clone(),
371                };
372                let result = fts_index.execute(current_input)?;
373                Ok(result)
374            } else {
375                Err("CREATE FTS INDEX requires a table catalog".into())
376            }
377        }
378        LogicalOperator::FtsScan(s) => {
379            if let Some(ref tc) = ctx.table_catalog {
380                let fts_scan = PhysicalFtsScan {
381                    index_name: s.index_name.clone(),
382                    query_string: s.query_string.clone(),
383                    docs_table: s.docs_table.clone(),
384                    terms_table: s.terms_table.clone(),
385                    posting_table: s.posting_table.clone(),
386                    table_name: s.table_name.clone(),
387                    column_name: s.column_name.clone(),
388                    table_catalog: tc.clone(),
389                };
390                let result = fts_scan.execute(current_input)?;
391                Ok(result)
392            } else {
393                Err("FTS scan requires a table catalog".into())
394            }
395        }
396        LogicalOperator::EmptyResult(_) => {
397            let exec = crate::physical::misc::PhysicalEmptyResult;
398            let result = exec.execute(current_input)?;
399            Ok(result)
400        }
401        LogicalOperator::MultiplicityReducer(m) => {
402            let exec = crate::physical::misc::PhysicalMultiplicityReducer {
403                key_columns: m.key_columns.clone(),
404            };
405            let input = if !m.children.is_empty() {
406                ctx.execute_children(&m.children)?
407            } else {
408                current_input
409            };
410            let result = exec.execute(input)?;
411            Ok(result)
412        }
413        LogicalOperator::Skip(s) => {
414            let exec = crate::physical::misc::PhysicalSkip {
415                skip_count: s.offset as usize,
416            };
417            let input = if !s.children.is_empty() {
418                ctx.execute_children(&s.children)?
419            } else {
420                current_input
421            };
422            let result = exec.execute(input)?;
423            Ok(result)
424        }
425        LogicalOperator::ExtensionClause(e) => {
426            let exec = crate::physical::misc::PhysicalExtensionClause {
427                action: e.action.clone(),
428                extension_name: e.extension_name.clone(),
429            };
430            let result = exec.execute(current_input)?;
431            Ok(result)
432        }
433        _ => Err(format!("DDL operator not implemented in mapper: {:?}", op).into()),
434    }
435}
436
437/// Build a success DataChunk with a single message column.
438fn ddl_success_chunk(message: &str) -> Vec<DataChunk> {
439    let mut v = akar_common::vector::ValueVector::new(akar_common::types::PhysicalTypeID::String, 1);
440    v.resize(1);
441    v.set_value(0, &akar_common::types::Value::String(message.to_string()))
442        .unwrap();
443    let arr = akar_common::arrow_vector::ArrowVector::from_legacy(&v).array;
444    let mut chunk = DataChunk::new(vec![arr], vec![akar_common::types::PhysicalTypeID::String]);
445    chunk.size = 1;
446    chunk.field_names = vec!["result".to_string()];
447    vec![chunk]
448}
449
450/// Minimal type parser for ALTER TABLE ADD COLUMN (avoids Akar-binder dependency).
451fn parse_type_simple(type_name: &str) -> Result<akar_common::types::LogicalTypeID, ProcessorError> {
452    let upper = type_name.trim().to_uppercase();
453    match upper.as_str() {
454        "BOOL" | "BOOLEAN" => Ok(akar_common::types::LogicalTypeID::Bool),
455        "INT64" => Ok(akar_common::types::LogicalTypeID::Int64),
456        "INT32" => Ok(akar_common::types::LogicalTypeID::Int32),
457        "INT16" => Ok(akar_common::types::LogicalTypeID::Int16),
458        "INT8" => Ok(akar_common::types::LogicalTypeID::Int8),
459        "UINT64" => Ok(akar_common::types::LogicalTypeID::UInt64),
460        "UINT32" => Ok(akar_common::types::LogicalTypeID::UInt32),
461        "UINT16" => Ok(akar_common::types::LogicalTypeID::UInt16),
462        "UINT8" => Ok(akar_common::types::LogicalTypeID::UInt8),
463        "DOUBLE" => Ok(akar_common::types::LogicalTypeID::Double),
464        "FLOAT" => Ok(akar_common::types::LogicalTypeID::Float),
465        "STRING" => Ok(akar_common::types::LogicalTypeID::String),
466        "BLOB" => Ok(akar_common::types::LogicalTypeID::Blob),
467        "DATE" => Ok(akar_common::types::LogicalTypeID::Date),
468        "TIMESTAMP" => Ok(akar_common::types::LogicalTypeID::Timestamp),
469        "INTERVAL" => Ok(akar_common::types::LogicalTypeID::Interval),
470        "UUID" => Ok(akar_common::types::LogicalTypeID::Uuid),
471        _ => Err(format!("Unknown type '{type_name}'").into()),
472    }
473}