Skip to main content

alopex_sql/executor/query/
mod.rs

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
30/// Execute a SELECT logical plan and return a query result.
31///
32/// This function uses an iterator-based execution model that processes rows
33/// through a pipeline of operators. This approach:
34/// - Enables early termination for LIMIT queries
35/// - Provides streaming execution after the initial scan
36/// - Allows composable query operators
37///
38/// Note: The Scan stage reads all matching rows into memory, but subsequent
39/// operators (Filter, Sort, Limit) process rows through an iterator pipeline.
40/// Sort operations additionally require materializing all input rows.
41pub 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
103/// Execute a SELECT logical plan and return a streaming query result.
104///
105/// This function returns a `QueryRowIterator` that yields rows one at a time,
106/// enabling true streaming output without materializing all rows upfront.
107///
108/// # FR-7 Streaming Output
109///
110/// This function implements the FR-7 requirement for streaming output.
111/// Rows are yielded through an iterator interface, and projection is applied
112/// on-the-fly as each row is consumed.
113///
114/// # Note
115///
116/// KNN queries currently fall back to the non-streaming path as they require
117/// specialized handling.
118pub 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    // KNN queries not yet supported for streaming - fall back would need different handling
138    if knn::extract_knn_context(&plan).is_some() {
139        // For KNN, we materialize and wrap in VecIterator
140        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    // Subqueries need transaction access during evaluation, which streaming
152    // iterators borrow exclusively. Execute through the materializing path
153    // (the same one used by `execute_query`) so results are identical to the
154    // non-streaming API instead of failing or silently dropping rows.
155    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
166/// Convert a materialized `QueryResult` into pipeline outputs.
167///
168/// The resulting rows are already fully projected, so the returned projection
169/// is `Projection::All` over the output column names.
170fn 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
193/// Build an iterator pipeline from a logical plan.
194///
195/// This recursively constructs a tree of iterators that mirrors the logical plan
196/// structure. The scan phase reads rows into memory, then subsequent operators
197/// process them through an iterator pipeline enabling streaming execution and
198/// early termination.
199fn 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            // TODO: 現状は Scan で一度全件をメモリに載せてから iterator に渡しています。
250            // 将来ストリーミングを徹底する場合は、ScanIterator を活用できるよう
251            // トランザクションのライフタイム設計を見直すとよいです。
252            let rows = scan::execute_scan(txn, &table_meta)?;
253            let schema = table_meta.columns.clone();
254
255            // Wrap in VecIterator for consistent iterator-based processing
256            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
453/// Build a streaming iterator pipeline from a logical plan (FR-7).
454///
455/// This version uses `ScanIterator` for row-based tables to enable true
456/// streaming without materializing all rows upfront. The returned iterator
457/// has lifetime `'a` tied to the transaction borrow.
458///
459/// # Limitations
460///
461/// - Columnar storage still materializes rows (uses VecIterator)
462/// - Sort operations materialize all input rows
463/// - KNN queries are not supported (use `build_iterator_pipeline` instead)
464pub 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    // Subqueries need transaction access during evaluation, which streaming
499    // iterators borrow exclusively. Execute through the materializing path
500    // (the same one used by `execute_query`) so results are identical to the
501    // non-streaming API instead of failing or silently dropping rows
502    // (GitHub issues #23 / #24).
503    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
511/// Inner implementation of streaming pipeline builder.
512fn 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                // Columnar storage: use ColumnarScanIterator for FR-7 streaming
543                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            // Row-based storage: use ScanIterator for true streaming (FR-7)
551            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                // Columnar storage with filter: use ColumnarScanIterator for FR-7 streaming
561                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
734/// Evaluate a typed expression against a row, returning SqlValue.
735fn 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
790/// Build column info name using alias fallback.
791fn 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
807/// Build ColumnInfo from projection.
808fn 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
818/// Build ColumnInfo for Projection::All using schema.
819fn 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}