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 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 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 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 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 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 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, ®istry);
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
437fn 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
450fn 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}