1use alopex_core::kv::KVStore;
2
3use crate::ast::LITERAL_TABLE;
4use crate::catalog::{Catalog, StorageType};
5use crate::executor::evaluator::EvalContext;
6use crate::executor::memory::MemoryPolicy;
7use crate::executor::{ExecutionResult, ExecutorError, QueryResult, QueryRowIterator, Result};
8use crate::planner::logical_plan::LogicalPlan;
9use crate::planner::typed_expr::{Projection, SortExpr};
10use crate::storage::{SqlTxn, SqlValue};
11
12use super::{ColumnInfo, Row};
13
14pub mod aggregate;
15pub mod columnar_scan;
16pub mod iterator;
17pub mod join;
18mod knn;
19mod project;
20mod scan;
21pub mod subquery;
22
23pub use columnar_scan::{ColumnarScanIterator, create_columnar_scan_iterator};
24pub use iterator::{FilterIterator, LimitIterator, RowIterator, ScanIterator, SortIterator};
25pub use project::{project_row_values, projected_columns};
26pub use scan::{
27 create_fenced_range_scan_iterator, create_scan_iterator, execute_fenced_range_scan,
28};
29
30pub fn execute_query<'txn, S: KVStore + 'txn, C: Catalog + ?Sized, T: SqlTxn<'txn, S>>(
42 txn: &mut T,
43 catalog: &C,
44 plan: LogicalPlan,
45) -> Result<ExecutionResult> {
46 execute_query_with_policy(txn, catalog, plan, None)
47}
48
49pub fn execute_query_with_policy<
50 'txn,
51 S: KVStore + 'txn,
52 C: Catalog + ?Sized,
53 T: SqlTxn<'txn, S>,
54>(
55 txn: &mut T,
56 catalog: &C,
57 plan: LogicalPlan,
58 memory: Option<&MemoryPolicy>,
59) -> Result<ExecutionResult> {
60 if let Some((pattern, projection, filter)) = knn::extract_knn_context(&plan) {
61 return knn::execute_knn_query(txn, catalog, &pattern, &projection, filter.as_ref());
62 }
63
64 let result = execute_query_result_with_outer_and_policy(txn, catalog, plan, None, memory)?;
65 Ok(ExecutionResult::Query(result))
66}
67
68pub(crate) fn execute_query_result_with_outer<
69 'txn,
70 S: KVStore + 'txn,
71 C: Catalog + ?Sized,
72 T: SqlTxn<'txn, S>,
73>(
74 txn: &mut T,
75 catalog: &C,
76 plan: LogicalPlan,
77 outer: Option<&Row>,
78) -> Result<QueryResult> {
79 execute_query_result_with_outer_and_policy(txn, catalog, plan, outer, None)
80}
81
82fn execute_query_result_with_outer_and_policy<
83 'txn,
84 S: KVStore + 'txn,
85 C: Catalog + ?Sized,
86 T: SqlTxn<'txn, S>,
87>(
88 txn: &mut T,
89 catalog: &C,
90 plan: LogicalPlan,
91 outer: Option<&Row>,
92 memory: Option<&MemoryPolicy>,
93) -> Result<QueryResult> {
94 let (mut iter, projection, schema) =
95 build_iterator_pipeline_with_outer(txn, catalog, plan, memory, outer)?;
96 let mut rows = Vec::new();
97 while let Some(result) = iter.next_row() {
98 rows.push(result?);
99 }
100 execute_project_with_subqueries(txn, catalog, rows, &projection, &schema, outer)
101}
102
103pub fn execute_query_streaming<'txn, S: KVStore + 'txn, C: Catalog + ?Sized, T: SqlTxn<'txn, S>>(
119 txn: &mut T,
120 catalog: &C,
121 plan: LogicalPlan,
122) -> Result<QueryRowIterator<'static>> {
123 execute_query_streaming_with_policy(txn, catalog, plan, None)
124}
125
126pub fn execute_query_streaming_with_policy<
127 'txn,
128 S: KVStore + 'txn,
129 C: Catalog + ?Sized,
130 T: SqlTxn<'txn, S>,
131>(
132 txn: &mut T,
133 catalog: &C,
134 plan: LogicalPlan,
135 memory: Option<&MemoryPolicy>,
136) -> Result<QueryRowIterator<'static>> {
137 if knn::extract_knn_context(&plan).is_some() {
139 let result = execute_query_with_policy(txn, catalog, plan, memory)?;
141 if let ExecutionResult::Query(qr) = result {
142 let (iter, projection, schema) = materialize_query_result(qr);
143 return Ok(QueryRowIterator::new(iter, projection, schema));
144 }
145 return Err(ExecutorError::InvalidOperation {
146 operation: "execute_query_streaming".into(),
147 reason: "KNN query did not return Query result".into(),
148 });
149 }
150
151 if subquery::plan_contains_subquery(&plan) {
156 let result = execute_query_result_with_outer_and_policy(txn, catalog, plan, None, memory)?;
157 let (iter, projection, schema) = materialize_query_result(result);
158 return Ok(QueryRowIterator::new(iter, projection, schema));
159 }
160
161 let (iter, projection, schema) = build_iterator_pipeline(txn, catalog, plan, memory)?;
162
163 Ok(QueryRowIterator::new(iter, projection, schema))
164}
165
166fn materialize_query_result(
171 result: QueryResult,
172) -> (
173 Box<dyn RowIterator>,
174 Projection,
175 Vec<crate::catalog::ColumnMetadata>,
176) {
177 let column_names: Vec<String> = result.columns.iter().map(|c| c.name.clone()).collect();
178 let schema: Vec<crate::catalog::ColumnMetadata> = result
179 .columns
180 .iter()
181 .map(|c| crate::catalog::ColumnMetadata::new(&c.name, c.data_type.clone()))
182 .collect();
183 let rows: Vec<Row> = result
184 .rows
185 .into_iter()
186 .enumerate()
187 .map(|(i, values)| Row::new(i as u64, values))
188 .collect();
189 let iter = iterator::VecIterator::new(rows, schema.clone());
190 (Box::new(iter), Projection::All(column_names), schema)
191}
192
193fn build_iterator_pipeline<'txn, S: KVStore + 'txn, C: Catalog + ?Sized, T: SqlTxn<'txn, S>>(
200 txn: &mut T,
201 catalog: &C,
202 plan: LogicalPlan,
203 memory: Option<&MemoryPolicy>,
204) -> Result<(
205 Box<dyn RowIterator>,
206 Projection,
207 Vec<crate::catalog::ColumnMetadata>,
208)> {
209 build_iterator_pipeline_with_outer(txn, catalog, plan, memory, None)
210}
211
212fn build_iterator_pipeline_with_outer<
213 'txn,
214 S: KVStore + 'txn,
215 C: Catalog + ?Sized,
216 T: SqlTxn<'txn, S>,
217>(
218 txn: &mut T,
219 catalog: &C,
220 plan: LogicalPlan,
221 memory: Option<&MemoryPolicy>,
222 outer: Option<&Row>,
223) -> Result<(
224 Box<dyn RowIterator>,
225 Projection,
226 Vec<crate::catalog::ColumnMetadata>,
227)> {
228 match plan {
229 LogicalPlan::Scan { table, projection } => {
230 if table == LITERAL_TABLE {
231 let schema = Vec::new();
232 let rows = vec![Row::new(0, Vec::new())];
233 let iter = iterator::VecIterator::new(rows, schema.clone());
234 return Ok((Box::new(iter), projection, schema));
235 }
236 let table_meta = catalog
237 .get_table(&table)
238 .cloned()
239 .ok_or_else(|| ExecutorError::TableNotFound(table.clone()))?;
240
241 if table_meta.storage_options.storage_type == StorageType::Columnar {
242 let columnar_scan = columnar_scan::build_columnar_scan(&table_meta, &projection);
243 let rows = columnar_scan::execute_columnar_scan(txn, &table_meta, &columnar_scan)?;
244 let schema = table_meta.columns.clone();
245 let iter = iterator::VecIterator::new(rows, schema.clone());
246 return Ok((Box::new(iter), projection, schema));
247 }
248
249 let rows = scan::execute_scan(txn, &table_meta)?;
253 let schema = table_meta.columns.clone();
254
255 let iter = iterator::VecIterator::new(rows, schema.clone());
257 Ok((Box::new(iter), projection, schema))
258 }
259 LogicalPlan::Filter { input, predicate } => {
260 if let LogicalPlan::Scan { table, projection } = input.as_ref()
261 && let Some(table_meta) = catalog.get_table(table)
262 && table_meta.storage_options.storage_type == StorageType::Columnar
263 {
264 let columnar_scan = columnar_scan::build_columnar_scan_for_filter(
265 table_meta,
266 projection.clone(),
267 &predicate,
268 );
269 let rows = columnar_scan::execute_columnar_scan(txn, table_meta, &columnar_scan)?;
270 let schema = table_meta.columns.clone();
271 let iter = iterator::VecIterator::new(rows, schema.clone());
272 return Ok((Box::new(iter), projection.clone(), schema));
273 }
274 let (mut input_iter, projection, schema) =
275 build_iterator_pipeline_with_outer(txn, catalog, *input, memory, outer)?;
276 if outer.is_some() || subquery::contains_subquery(&predicate) {
277 let mut rows = Vec::new();
278 while let Some(result) = input_iter.next_row() {
279 let row = result?;
280 let eval_row = combine_outer_for_eval(&row, outer);
281 if let SqlValue::Boolean(true) = subquery::evaluate_expr_with_subqueries(
282 txn, catalog, &predicate, &eval_row,
283 )? {
284 rows.push(row);
285 }
286 }
287 let iter = iterator::VecIterator::new(rows, schema.clone());
288 return Ok((Box::new(iter), projection, schema));
289 }
290 let filter_iter = FilterIterator::new(input_iter, predicate);
291 Ok((Box::new(filter_iter), projection, schema))
292 }
293 LogicalPlan::Project { input, projection } => {
294 let (mut input_iter, _input_projection, schema) =
295 build_iterator_pipeline_with_outer(txn, catalog, *input, memory, outer)?;
296 let mut rows = Vec::new();
297 while let Some(result) = input_iter.next_row() {
298 rows.push(result?);
299 }
300 let projected =
301 execute_project_with_subqueries(txn, catalog, rows, &projection, &schema, outer)?;
302 let output_schema = projected
303 .columns
304 .iter()
305 .map(|col| crate::catalog::ColumnMetadata::new(&col.name, col.data_type.clone()))
306 .collect::<Vec<_>>();
307 let rows = projected
308 .rows
309 .into_iter()
310 .enumerate()
311 .map(|(idx, values)| Row::new(idx as u64, values))
312 .collect::<Vec<_>>();
313 let output_projection =
314 Projection::All(output_schema.iter().map(|col| col.name.clone()).collect());
315 let iter = iterator::VecIterator::new(rows, output_schema.clone());
316 Ok((Box::new(iter), output_projection, output_schema))
317 }
318 LogicalPlan::Join {
319 left,
320 right,
321 join_type,
322 condition,
323 using: _,
324 } => {
325 let (mut left_iter, _left_projection, left_schema) =
326 build_iterator_pipeline_with_outer(txn, catalog, *left, memory, outer)?;
327 let (mut right_iter, _right_projection, right_schema) =
328 build_iterator_pipeline_with_outer(txn, catalog, *right, memory, outer)?;
329 let mut left_rows = Vec::new();
330 while let Some(result) = left_iter.next_row() {
331 left_rows.push(result?);
332 }
333 let mut right_rows = Vec::new();
334 while let Some(result) = right_iter.next_row() {
335 right_rows.push(result?);
336 }
337 let left_width = left_schema.len();
338 let right_width = right_schema.len();
339 let rows = join::execute_join_with_widths(
340 left_rows,
341 right_rows,
342 join_type,
343 condition.as_ref(),
344 left_width,
345 right_width,
346 )?;
347 let mut schema = left_schema;
348 schema.extend(right_schema);
349 let projection = Projection::All(schema.iter().map(|col| col.name.clone()).collect());
350 let iter = iterator::VecIterator::new(rows, schema.clone());
351 Ok((Box::new(iter), projection, schema))
352 }
353 LogicalPlan::Aggregate {
354 input,
355 group_keys,
356 aggregates,
357 having,
358 projection,
359 } => {
360 let (input_iter, _projection, _schema) =
361 build_iterator_pipeline_with_outer(txn, catalog, *input, memory, outer)?;
362 let schema = aggregate::build_aggregate_schema(&group_keys, &aggregates);
363 if let Some(policy) = memory
364 && policy.spill_directory().is_some()
365 {
366 if group_keys.is_empty() {
367 let iter = aggregate::StreamingAggregateIterator::new(
368 input_iter,
369 group_keys,
370 aggregates,
371 having,
372 schema.clone(),
373 );
374 return Ok((Box::new(iter), projection, schema));
375 }
376 let order_by = group_keys
377 .iter()
378 .cloned()
379 .map(|expr| SortExpr {
380 expr,
381 asc: true,
382 nulls_first: false,
383 })
384 .collect::<Vec<_>>();
385 let sort_iter =
386 SortIterator::new_with_policy(input_iter, &order_by, Some(policy.clone()))?;
387 let iter = aggregate::StreamingAggregateIterator::new(
388 Box::new(sort_iter),
389 group_keys,
390 aggregates,
391 having,
392 schema.clone(),
393 );
394 return Ok((Box::new(iter), projection, schema));
395 }
396
397 let parallelism = std::thread::available_parallelism()
398 .map(usize::from)
399 .unwrap_or(1);
400 if !aggregate::should_use_single_for_parallel(parallelism, &aggregates) {
401 let rows = aggregate::execute_parallel_aggregate_rows_with_policy(
402 input_iter,
403 group_keys,
404 aggregates,
405 having,
406 schema.clone(),
407 parallelism,
408 memory.cloned(),
409 1_000_000,
410 )?;
411 let iter = iterator::VecIterator::new(rows, schema.clone());
412 return Ok((Box::new(iter), projection, schema));
413 }
414
415 let mut iter = aggregate::AggregateIterator::new(
416 input_iter,
417 group_keys,
418 aggregates,
419 having,
420 schema.clone(),
421 );
422 if let Some(policy) = memory {
423 iter = iter.with_memory_policy(Some(policy.clone()));
424 }
425 Ok((Box::new(iter), projection, schema))
426 }
427 LogicalPlan::Sort { input, order_by } => {
428 let (input_iter, projection, schema) =
429 build_iterator_pipeline_with_outer(txn, catalog, *input, memory, outer)?;
430 let sort_iter = if let Some(policy) = memory {
431 SortIterator::new_with_policy(input_iter, &order_by, Some(policy.clone()))?
432 } else {
433 SortIterator::new(input_iter, &order_by)?
434 };
435 Ok((Box::new(sort_iter), projection, schema))
436 }
437 LogicalPlan::Limit {
438 input,
439 limit,
440 offset,
441 } => {
442 let (input_iter, projection, schema) =
443 build_iterator_pipeline_with_outer(txn, catalog, *input, memory, outer)?;
444 let limit_iter = LimitIterator::new(input_iter, limit, offset);
445 Ok((Box::new(limit_iter), projection, schema))
446 }
447 other => Err(ExecutorError::UnsupportedOperation(format!(
448 "unsupported query plan: {other:?}"
449 ))),
450 }
451}
452
453pub fn build_streaming_pipeline<
465 'a,
466 'txn: 'a,
467 S: KVStore + 'txn,
468 C: Catalog + ?Sized,
469 T: SqlTxn<'txn, S>,
470>(
471 txn: &'a mut T,
472 catalog: &C,
473 plan: LogicalPlan,
474) -> Result<(
475 Box<dyn RowIterator + 'a>,
476 Projection,
477 Vec<crate::catalog::ColumnMetadata>,
478)> {
479 build_streaming_pipeline_with_policy(txn, catalog, plan, None)
480}
481
482pub fn build_streaming_pipeline_with_policy<
483 'a,
484 'txn: 'a,
485 S: KVStore + 'txn,
486 C: Catalog + ?Sized,
487 T: SqlTxn<'txn, S>,
488>(
489 txn: &'a mut T,
490 catalog: &C,
491 plan: LogicalPlan,
492 memory: Option<&MemoryPolicy>,
493) -> Result<(
494 Box<dyn RowIterator + 'a>,
495 Projection,
496 Vec<crate::catalog::ColumnMetadata>,
497)> {
498 if subquery::plan_contains_subquery(&plan) {
504 let result = execute_query_result_with_outer_and_policy(txn, catalog, plan, None, memory)?;
505 return Ok(materialize_query_result(result));
506 }
507
508 build_streaming_pipeline_inner(txn, catalog, plan, memory)
509}
510
511fn build_streaming_pipeline_inner<
513 'a,
514 'txn: 'a,
515 S: KVStore + 'txn,
516 C: Catalog + ?Sized,
517 T: SqlTxn<'txn, S>,
518>(
519 txn: &'a mut T,
520 catalog: &C,
521 plan: LogicalPlan,
522 memory: Option<&MemoryPolicy>,
523) -> Result<(
524 Box<dyn RowIterator + 'a>,
525 Projection,
526 Vec<crate::catalog::ColumnMetadata>,
527)> {
528 match plan {
529 LogicalPlan::Scan { table, projection } => {
530 if table == LITERAL_TABLE {
531 let schema = Vec::new();
532 let rows = vec![Row::new(0, Vec::new())];
533 let iter = iterator::VecIterator::new(rows, schema.clone());
534 return Ok((Box::new(iter), projection, schema));
535 }
536 let table_meta = catalog
537 .get_table(&table)
538 .cloned()
539 .ok_or_else(|| ExecutorError::TableNotFound(table.clone()))?;
540
541 if table_meta.storage_options.storage_type == StorageType::Columnar {
542 let columnar_scan = columnar_scan::build_columnar_scan(&table_meta, &projection);
544 let schema = table_meta.columns.clone();
545 let iter =
546 columnar_scan::create_columnar_scan_iterator(txn, &table_meta, &columnar_scan)?;
547 return Ok((Box::new(iter), projection, schema));
548 }
549
550 let schema = table_meta.columns.clone();
552 let scan_iter = scan::create_scan_iterator(txn, &table_meta)?;
553 Ok((Box::new(scan_iter), projection, schema))
554 }
555 LogicalPlan::Filter { input, predicate } => {
556 if let LogicalPlan::Scan { table, projection } = input.as_ref()
557 && let Some(table_meta) = catalog.get_table(table)
558 && table_meta.storage_options.storage_type == StorageType::Columnar
559 {
560 let columnar_scan = columnar_scan::build_columnar_scan_for_filter(
562 table_meta,
563 projection.clone(),
564 &predicate,
565 );
566 let schema = table_meta.columns.clone();
567 let iter =
568 columnar_scan::create_columnar_scan_iterator(txn, table_meta, &columnar_scan)?;
569 return Ok((Box::new(iter), projection.clone(), schema));
570 }
571 let (input_iter, projection, schema) =
572 build_streaming_pipeline_inner(txn, catalog, *input, memory)?;
573 let filter_iter = FilterIterator::new(input_iter, predicate);
574 Ok((Box::new(filter_iter), projection, schema))
575 }
576 LogicalPlan::Project { input, projection } => {
577 let (mut input_iter, _input_projection, schema) =
578 build_streaming_pipeline_inner(txn, catalog, *input, memory)?;
579 let mut rows = Vec::new();
580 while let Some(result) = input_iter.next_row() {
581 rows.push(result?);
582 }
583 let projected = project::execute_project(rows, &projection, &schema)?;
584 let output_schema = projected
585 .columns
586 .iter()
587 .map(|col| crate::catalog::ColumnMetadata::new(&col.name, col.data_type.clone()))
588 .collect::<Vec<_>>();
589 let rows = projected
590 .rows
591 .into_iter()
592 .enumerate()
593 .map(|(idx, values)| Row::new(idx as u64, values))
594 .collect::<Vec<_>>();
595 let output_projection =
596 Projection::All(output_schema.iter().map(|col| col.name.clone()).collect());
597 let iter = iterator::VecIterator::new(rows, output_schema.clone());
598 Ok((Box::new(iter), output_projection, output_schema))
599 }
600 LogicalPlan::Join {
601 left,
602 right,
603 join_type,
604 condition,
605 using: _,
606 } => {
607 let (mut left_iter, _left_projection, left_schema) =
608 build_streaming_pipeline_inner(txn, catalog, *left, memory)?;
609 let mut left_rows = Vec::new();
610 while let Some(result) = left_iter.next_row() {
611 left_rows.push(result?);
612 }
613 drop(left_iter);
614 let (mut right_iter, _right_projection, right_schema) =
615 build_streaming_pipeline_inner(txn, catalog, *right, memory)?;
616 let mut right_rows = Vec::new();
617 while let Some(result) = right_iter.next_row() {
618 right_rows.push(result?);
619 }
620 let rows = join::execute_join_with_widths(
621 left_rows,
622 right_rows,
623 join_type,
624 condition.as_ref(),
625 left_schema.len(),
626 right_schema.len(),
627 )?;
628 let mut schema = left_schema;
629 schema.extend(right_schema);
630 let projection = Projection::All(schema.iter().map(|col| col.name.clone()).collect());
631 let iter = iterator::VecIterator::new(rows, schema.clone());
632 Ok((Box::new(iter), projection, schema))
633 }
634 LogicalPlan::Aggregate {
635 input,
636 group_keys,
637 aggregates,
638 having,
639 projection,
640 } => {
641 let (input_iter, _projection, _schema) =
642 build_streaming_pipeline_inner(txn, catalog, *input, memory)?;
643 let schema = aggregate::build_aggregate_schema(&group_keys, &aggregates);
644 if let Some(policy) = memory
645 && policy.spill_directory().is_some()
646 {
647 if group_keys.is_empty() {
648 let iter = aggregate::StreamingAggregateIterator::new(
649 input_iter,
650 group_keys,
651 aggregates,
652 having,
653 schema.clone(),
654 );
655 return Ok((Box::new(iter), projection, schema));
656 }
657 let order_by = group_keys
658 .iter()
659 .cloned()
660 .map(|expr| SortExpr {
661 expr,
662 asc: true,
663 nulls_first: false,
664 })
665 .collect::<Vec<_>>();
666 let sort_iter =
667 SortIterator::new_with_policy(input_iter, &order_by, Some(policy.clone()))?;
668 let iter = aggregate::StreamingAggregateIterator::new(
669 Box::new(sort_iter),
670 group_keys,
671 aggregates,
672 having,
673 schema.clone(),
674 );
675 return Ok((Box::new(iter), projection, schema));
676 }
677
678 let parallelism = std::thread::available_parallelism()
679 .map(usize::from)
680 .unwrap_or(1);
681 if !aggregate::should_use_single_for_parallel(parallelism, &aggregates) {
682 let rows = aggregate::execute_parallel_aggregate_rows_with_policy(
683 input_iter,
684 group_keys,
685 aggregates,
686 having,
687 schema.clone(),
688 parallelism,
689 memory.cloned(),
690 1_000_000,
691 )?;
692 let iter = iterator::VecIterator::new(rows, schema.clone());
693 return Ok((Box::new(iter), projection, schema));
694 }
695
696 let mut iter = aggregate::AggregateIterator::new(
697 input_iter,
698 group_keys,
699 aggregates,
700 having,
701 schema.clone(),
702 );
703 if let Some(policy) = memory {
704 iter = iter.with_memory_policy(Some(policy.clone()));
705 }
706 Ok((Box::new(iter), projection, schema))
707 }
708 LogicalPlan::Sort { input, order_by } => {
709 let (input_iter, projection, schema) =
710 build_streaming_pipeline_inner(txn, catalog, *input, memory)?;
711 let sort_iter = if let Some(policy) = memory {
712 SortIterator::new_with_policy(input_iter, &order_by, Some(policy.clone()))?
713 } else {
714 SortIterator::new(input_iter, &order_by)?
715 };
716 Ok((Box::new(sort_iter), projection, schema))
717 }
718 LogicalPlan::Limit {
719 input,
720 limit,
721 offset,
722 } => {
723 let (input_iter, projection, schema) =
724 build_streaming_pipeline_inner(txn, catalog, *input, memory)?;
725 let limit_iter = LimitIterator::new(input_iter, limit, offset);
726 Ok((Box::new(limit_iter), projection, schema))
727 }
728 other => Err(ExecutorError::UnsupportedOperation(format!(
729 "unsupported query plan: {other:?}"
730 ))),
731 }
732}
733
734fn eval_expr(expr: &crate::planner::typed_expr::TypedExpr, row: &Row) -> Result<SqlValue> {
736 let ctx = EvalContext::new(&row.values);
737 crate::executor::evaluator::evaluate(expr, &ctx)
738}
739
740fn combine_outer_for_eval(row: &Row, outer: Option<&Row>) -> Row {
741 let Some(outer) = outer else {
742 return row.clone();
743 };
744 let mut values = Vec::with_capacity(row.len() + outer.len());
745 values.extend(row.values.clone());
746 values.extend(outer.values.clone());
747 Row::new(row.row_id, values)
748}
749
750fn execute_project_with_subqueries<
751 'txn,
752 S: KVStore + 'txn,
753 C: Catalog + ?Sized,
754 T: SqlTxn<'txn, S>,
755>(
756 txn: &mut T,
757 catalog: &C,
758 rows: Vec<Row>,
759 projection: &Projection,
760 schema: &[crate::catalog::ColumnMetadata],
761 outer: Option<&Row>,
762) -> Result<QueryResult> {
763 match projection {
764 Projection::All(_) => project::execute_project(rows, projection, schema),
765 Projection::Columns(cols)
766 if outer.is_some() || cols.iter().any(|c| subquery::contains_subquery(&c.expr)) =>
767 {
768 let columns: Vec<_> = cols
769 .iter()
770 .enumerate()
771 .map(|(i, c)| column_info_from_projection(c, i))
772 .collect();
773 let mut projected_rows = Vec::with_capacity(rows.len());
774 for row in rows {
775 let eval_row = combine_outer_for_eval(&row, outer);
776 let mut values = Vec::with_capacity(cols.len());
777 for col in cols {
778 values.push(subquery::evaluate_expr_with_subqueries(
779 txn, catalog, &col.expr, &eval_row,
780 )?);
781 }
782 projected_rows.push(values);
783 }
784 Ok(QueryResult::new(columns, projected_rows))
785 }
786 Projection::Columns(_) => project::execute_project(rows, projection, schema),
787 }
788}
789
790fn column_name_from_projection(
792 projected: &crate::planner::typed_expr::ProjectedColumn,
793 idx: usize,
794) -> String {
795 projected
796 .alias
797 .clone()
798 .or_else(|| match &projected.expr.kind {
799 crate::planner::typed_expr::TypedExprKind::ColumnRef { column, .. } => {
800 Some(column.clone())
801 }
802 _ => None,
803 })
804 .unwrap_or_else(|| format!("col_{idx}"))
805}
806
807fn column_info_from_projection(
809 projected: &crate::planner::typed_expr::ProjectedColumn,
810 idx: usize,
811) -> ColumnInfo {
812 ColumnInfo::new(
813 column_name_from_projection(projected, idx),
814 projected.expr.resolved_type.clone(),
815 )
816}
817
818fn column_infos_from_all(
820 schema: &[crate::catalog::ColumnMetadata],
821 names: &[String],
822) -> Result<Vec<ColumnInfo>> {
823 names
824 .iter()
825 .map(|name| {
826 let col = schema
827 .iter()
828 .find(|c| &c.name == name)
829 .ok_or_else(|| ExecutorError::ColumnNotFound(name.clone()))?;
830 Ok(ColumnInfo::new(name.clone(), col.data_type.clone()))
831 })
832 .collect()
833}
834
835#[cfg(test)]
836mod tests {
837 use super::*;
838 use crate::catalog::{ColumnMetadata, MemoryCatalog, TableMetadata};
839 use crate::executor::ddl::create_table::execute_create_table;
840 use crate::planner::typed_expr::TypedExpr;
841 use crate::planner::types::ResolvedType;
842 use crate::storage::TxnBridge;
843 use alopex_core::kv::memory::MemoryKV;
844 use std::sync::Arc;
845
846 #[test]
847 fn execute_query_scan_only_returns_rows() {
848 let bridge = TxnBridge::new(Arc::new(MemoryKV::new()));
849 let mut catalog = MemoryCatalog::new();
850 let table = TableMetadata::new(
851 "users",
852 vec![
853 ColumnMetadata::new("id", ResolvedType::Integer),
854 ColumnMetadata::new("name", ResolvedType::Text),
855 ],
856 );
857 let mut ddl_txn = bridge.begin_write().unwrap();
858 execute_create_table(&mut ddl_txn, &mut catalog, table.clone(), vec![], false).unwrap();
859 ddl_txn.commit().unwrap();
860
861 let mut txn = bridge.begin_write().unwrap();
862 crate::executor::dml::execute_insert(
863 &mut txn,
864 &catalog,
865 "users",
866 vec!["id".into(), "name".into()],
867 vec![vec![
868 TypedExpr::literal(
869 crate::ast::expr::Literal::Number("1".into()),
870 ResolvedType::Integer,
871 crate::Span::default(),
872 ),
873 TypedExpr::literal(
874 crate::ast::expr::Literal::String("alice".into()),
875 ResolvedType::Text,
876 crate::Span::default(),
877 ),
878 ]],
879 )
880 .unwrap();
881
882 let result = execute_query(
883 &mut txn,
884 &catalog,
885 LogicalPlan::scan(
886 "users".into(),
887 Projection::All(vec!["id".into(), "name".into()]),
888 ),
889 )
890 .unwrap();
891
892 match result {
893 ExecutionResult::Query(q) => {
894 assert_eq!(q.rows.len(), 1);
895 assert_eq!(q.columns.len(), 2);
896 assert_eq!(
897 q.rows[0],
898 vec![SqlValue::Integer(1), SqlValue::Text("alice".into())]
899 );
900 }
901 other => panic!("unexpected result {other:?}"),
902 }
903 }
904}