1use anyhow::{anyhow, Result};
2use fxhash::FxHashSet;
3use std::cmp::min;
4use std::collections::HashMap;
5use std::sync::Arc;
6use std::time::{Duration, Instant};
7use tracing::{debug, info};
8
9use crate::config::config::BehaviorConfig;
10use crate::config::global::get_date_notation;
11use crate::data::arithmetic_evaluator::ArithmeticEvaluator;
12use crate::data::data_view::DataView;
13use crate::data::datatable::{DataColumn, DataRow, DataTable, DataValue};
14use crate::data::evaluation_context::EvaluationContext;
15use crate::data::group_by_expressions::GroupByExpressions;
16use crate::data::hash_join::HashJoinExecutor;
17use crate::data::recursive_where_evaluator::RecursiveWhereEvaluator;
18use crate::data::row_expanders::RowExpanderRegistry;
19use crate::data::subquery_executor::SubqueryExecutor;
20use crate::data::temp_table_registry::TempTableRegistry;
21use crate::execution_plan::{ExecutionPlan, ExecutionPlanBuilder, StepType};
22use crate::sql::aggregates::{contains_aggregate, is_aggregate_compatible};
23use crate::sql::parser::ast::ColumnRef;
24use crate::sql::parser::ast::SetOperation;
25use crate::sql::parser::ast::TableSource;
26use crate::sql::parser::ast::WindowSpec;
27use crate::sql::recursive_parser::{
28 CTEType, OrderByItem, Parser, SelectItem, SelectStatement, SortDirection, SqlExpression,
29 TableFunction,
30};
31
32fn resolve_cte<'a>(
37 context: &'a HashMap<String, Arc<DataView>>,
38 name: &str,
39) -> Option<&'a Arc<DataView>> {
40 if let Some(v) = context.get(name) {
41 return Some(v);
42 }
43 let lower = name.to_lowercase();
44 context
45 .iter()
46 .find(|(k, _)| k.to_lowercase() == lower)
47 .map(|(_, v)| v)
48}
49
50#[derive(Debug, Clone)]
52pub struct ExecutionContext {
53 alias_map: HashMap<String, String>,
56}
57
58impl ExecutionContext {
59 pub fn new() -> Self {
61 Self {
62 alias_map: HashMap::new(),
63 }
64 }
65
66 pub fn register_alias(&mut self, alias: String, table_name: String) {
68 debug!("Registering alias: {} -> {}", alias, table_name);
69 self.alias_map.insert(alias, table_name);
70 }
71
72 pub fn resolve_alias(&self, name: &str) -> String {
75 self.alias_map
76 .get(name)
77 .cloned()
78 .unwrap_or_else(|| name.to_string())
79 }
80
81 pub fn is_alias(&self, name: &str) -> bool {
83 self.alias_map.contains_key(name)
84 }
85
86 pub fn get_aliases(&self) -> HashMap<String, String> {
88 self.alias_map.clone()
89 }
90
91 pub fn resolve_column_index(&self, table: &DataTable, column_ref: &ColumnRef) -> Result<usize> {
106 if let Some(table_prefix) = &column_ref.table_prefix {
107 let actual_table = self.resolve_alias(table_prefix);
109
110 let qualified_name = format!("{}.{}", actual_table, column_ref.name);
112 if let Some(idx) = table.find_column_by_qualified_name(&qualified_name) {
113 debug!(
114 "Resolved {}.{} -> qualified column '{}' at index {}",
115 table_prefix, column_ref.name, qualified_name, idx
116 );
117 return Ok(idx);
118 }
119
120 if let Some(idx) = table.get_column_index(&column_ref.name) {
122 debug!(
123 "Resolved {}.{} -> unqualified column '{}' at index {}",
124 table_prefix, column_ref.name, column_ref.name, idx
125 );
126 return Ok(idx);
127 }
128
129 Err(anyhow!(
131 "Column '{}' not found. Table '{}' may not support qualified column names",
132 qualified_name,
133 actual_table
134 ))
135 } else {
136 if let Some(idx) = table.get_column_index(&column_ref.name) {
138 debug!(
139 "Resolved unqualified column '{}' at index {}",
140 column_ref.name, idx
141 );
142 return Ok(idx);
143 }
144
145 if column_ref.name.contains('.') {
147 if let Some(idx) = table.find_column_by_qualified_name(&column_ref.name) {
148 debug!(
149 "Resolved '{}' as qualified column at index {}",
150 column_ref.name, idx
151 );
152 return Ok(idx);
153 }
154 }
155
156 let suggestion = self.find_similar_column(table, &column_ref.name);
158 match suggestion {
159 Some(similar) => Err(anyhow!(
160 "Column '{}' not found. Did you mean '{}'?",
161 column_ref.name,
162 similar
163 )),
164 None => Err(anyhow!("Column '{}' not found", column_ref.name)),
165 }
166 }
167 }
168
169 fn find_similar_column(&self, table: &DataTable, name: &str) -> Option<String> {
171 let columns = table.column_names();
172 let mut best_match: Option<(String, usize)> = None;
173
174 for col in columns {
175 let distance = edit_distance(name, &col);
176 if distance <= 2 {
177 match best_match {
179 Some((_, best_dist)) if distance < best_dist => {
180 best_match = Some((col.clone(), distance));
181 }
182 None => {
183 best_match = Some((col.clone(), distance));
184 }
185 _ => {}
186 }
187 }
188 }
189
190 best_match.map(|(name, _)| name)
191 }
192}
193
194impl Default for ExecutionContext {
195 fn default() -> Self {
196 Self::new()
197 }
198}
199
200fn edit_distance(a: &str, b: &str) -> usize {
202 let len_a = a.chars().count();
203 let len_b = b.chars().count();
204
205 if len_a == 0 {
206 return len_b;
207 }
208 if len_b == 0 {
209 return len_a;
210 }
211
212 let mut matrix = vec![vec![0; len_b + 1]; len_a + 1];
213
214 for i in 0..=len_a {
215 matrix[i][0] = i;
216 }
217 for j in 0..=len_b {
218 matrix[0][j] = j;
219 }
220
221 let a_chars: Vec<char> = a.chars().collect();
222 let b_chars: Vec<char> = b.chars().collect();
223
224 for i in 1..=len_a {
225 for j in 1..=len_b {
226 let cost = if a_chars[i - 1] == b_chars[j - 1] {
227 0
228 } else {
229 1
230 };
231 matrix[i][j] = min(
232 min(matrix[i - 1][j] + 1, matrix[i][j - 1] + 1),
233 matrix[i - 1][j - 1] + cost,
234 );
235 }
236 }
237
238 matrix[len_a][len_b]
239}
240
241#[derive(Clone)]
243pub struct QueryEngine {
244 case_insensitive: bool,
245 date_notation: String,
246 _behavior_config: Option<BehaviorConfig>,
247}
248
249impl Default for QueryEngine {
250 fn default() -> Self {
251 Self::new()
252 }
253}
254
255impl QueryEngine {
256 #[must_use]
257 pub fn new() -> Self {
258 Self {
259 case_insensitive: false,
260 date_notation: get_date_notation(),
261 _behavior_config: None,
262 }
263 }
264
265 #[must_use]
266 pub fn with_behavior_config(config: BehaviorConfig) -> Self {
267 let case_insensitive = config.case_insensitive_default;
268 let date_notation = get_date_notation();
270 Self {
271 case_insensitive,
272 date_notation,
273 _behavior_config: Some(config),
274 }
275 }
276
277 #[must_use]
278 pub fn with_date_notation(_date_notation: String) -> Self {
279 Self {
280 case_insensitive: false,
281 date_notation: get_date_notation(), _behavior_config: None,
283 }
284 }
285
286 #[must_use]
287 pub fn with_case_insensitive(case_insensitive: bool) -> Self {
288 Self {
289 case_insensitive,
290 date_notation: get_date_notation(),
291 _behavior_config: None,
292 }
293 }
294
295 #[must_use]
296 pub fn with_case_insensitive_and_date_notation(
297 case_insensitive: bool,
298 _date_notation: String, ) -> Self {
300 Self {
301 case_insensitive,
302 date_notation: get_date_notation(), _behavior_config: None,
304 }
305 }
306
307 fn find_similar_column(&self, table: &DataTable, name: &str) -> Option<String> {
309 let columns = table.column_names();
310 let mut best_match: Option<(String, usize)> = None;
311
312 for col in columns {
313 let distance = self.edit_distance(&col.to_lowercase(), &name.to_lowercase());
314 let max_distance = if name.len() > 10 { 3 } else { 2 };
317 if distance <= max_distance {
318 match &best_match {
319 None => best_match = Some((col, distance)),
320 Some((_, best_dist)) if distance < *best_dist => {
321 best_match = Some((col, distance));
322 }
323 _ => {}
324 }
325 }
326 }
327
328 best_match.map(|(name, _)| name)
329 }
330
331 fn edit_distance(&self, s1: &str, s2: &str) -> usize {
333 let len1 = s1.len();
334 let len2 = s2.len();
335 let mut matrix = vec![vec![0; len2 + 1]; len1 + 1];
336
337 for i in 0..=len1 {
338 matrix[i][0] = i;
339 }
340 for j in 0..=len2 {
341 matrix[0][j] = j;
342 }
343
344 for (i, c1) in s1.chars().enumerate() {
345 for (j, c2) in s2.chars().enumerate() {
346 let cost = usize::from(c1 != c2);
347 matrix[i + 1][j + 1] = std::cmp::min(
348 matrix[i][j + 1] + 1, std::cmp::min(
350 matrix[i + 1][j] + 1, matrix[i][j] + cost, ),
353 );
354 }
355 }
356
357 matrix[len1][len2]
358 }
359
360 fn contains_unnest(expr: &SqlExpression) -> bool {
362 match expr {
363 SqlExpression::Unnest { .. } => true,
365 SqlExpression::FunctionCall { name, args, .. } => {
366 if name.to_uppercase() == "UNNEST" {
367 return true;
368 }
369 args.iter().any(Self::contains_unnest)
371 }
372 SqlExpression::BinaryOp { left, right, .. } => {
373 Self::contains_unnest(left) || Self::contains_unnest(right)
374 }
375 SqlExpression::Not { expr } => Self::contains_unnest(expr),
376 SqlExpression::CaseExpression {
377 when_branches,
378 else_branch,
379 } => {
380 when_branches.iter().any(|branch| {
381 Self::contains_unnest(&branch.condition)
382 || Self::contains_unnest(&branch.result)
383 }) || else_branch
384 .as_ref()
385 .map_or(false, |e| Self::contains_unnest(e))
386 }
387 SqlExpression::SimpleCaseExpression {
388 expr,
389 when_branches,
390 else_branch,
391 } => {
392 Self::contains_unnest(expr)
393 || when_branches.iter().any(|branch| {
394 Self::contains_unnest(&branch.value)
395 || Self::contains_unnest(&branch.result)
396 })
397 || else_branch
398 .as_ref()
399 .map_or(false, |e| Self::contains_unnest(e))
400 }
401 SqlExpression::InList { expr, values } => {
402 Self::contains_unnest(expr) || values.iter().any(Self::contains_unnest)
403 }
404 SqlExpression::NotInList { expr, values } => {
405 Self::contains_unnest(expr) || values.iter().any(Self::contains_unnest)
406 }
407 SqlExpression::Between { expr, lower, upper } => {
408 Self::contains_unnest(expr)
409 || Self::contains_unnest(lower)
410 || Self::contains_unnest(upper)
411 }
412 SqlExpression::InSubquery { expr, .. } => Self::contains_unnest(expr),
413 SqlExpression::NotInSubquery { expr, .. } => Self::contains_unnest(expr),
414 SqlExpression::ScalarSubquery { .. } => false, SqlExpression::WindowFunction { args, .. } => args.iter().any(Self::contains_unnest),
416 SqlExpression::MethodCall { args, .. } => args.iter().any(Self::contains_unnest),
417 SqlExpression::ChainedMethodCall { base, args, .. } => {
418 Self::contains_unnest(base) || args.iter().any(Self::contains_unnest)
419 }
420 _ => false,
421 }
422 }
423
424 fn collect_window_specs(expr: &SqlExpression, specs: &mut Vec<WindowSpec>) {
426 match expr {
427 SqlExpression::WindowFunction {
428 window_spec, args, ..
429 } => {
430 specs.push(window_spec.clone());
432 for arg in args {
434 Self::collect_window_specs(arg, specs);
435 }
436 }
437 SqlExpression::BinaryOp { left, right, .. } => {
438 Self::collect_window_specs(left, specs);
439 Self::collect_window_specs(right, specs);
440 }
441 SqlExpression::Not { expr } => {
442 Self::collect_window_specs(expr, specs);
443 }
444 SqlExpression::FunctionCall { args, .. } => {
445 for arg in args {
446 Self::collect_window_specs(arg, specs);
447 }
448 }
449 SqlExpression::CaseExpression {
450 when_branches,
451 else_branch,
452 } => {
453 for branch in when_branches {
454 Self::collect_window_specs(&branch.condition, specs);
455 Self::collect_window_specs(&branch.result, specs);
456 }
457 if let Some(else_expr) = else_branch {
458 Self::collect_window_specs(else_expr, specs);
459 }
460 }
461 SqlExpression::SimpleCaseExpression {
462 expr,
463 when_branches,
464 else_branch,
465 } => {
466 Self::collect_window_specs(expr, specs);
467 for branch in when_branches {
468 Self::collect_window_specs(&branch.value, specs);
469 Self::collect_window_specs(&branch.result, specs);
470 }
471 if let Some(else_expr) = else_branch {
472 Self::collect_window_specs(else_expr, specs);
473 }
474 }
475 SqlExpression::InList { expr, values, .. } => {
476 Self::collect_window_specs(expr, specs);
477 for item in values {
478 Self::collect_window_specs(item, specs);
479 }
480 }
481 SqlExpression::ChainedMethodCall { base, args, .. } => {
482 Self::collect_window_specs(base, specs);
483 for arg in args {
484 Self::collect_window_specs(arg, specs);
485 }
486 }
487 SqlExpression::Column(_)
489 | SqlExpression::NumberLiteral(_)
490 | SqlExpression::StringLiteral(_)
491 | SqlExpression::BooleanLiteral(_)
492 | SqlExpression::Null
493 | SqlExpression::DateTimeToday { .. }
494 | SqlExpression::DateTimeConstructor { .. }
495 | SqlExpression::MethodCall { .. } => {}
496 _ => {}
498 }
499 }
500
501 fn contains_window_function(expr: &SqlExpression) -> bool {
503 match expr {
504 SqlExpression::WindowFunction { .. } => true,
505 SqlExpression::BinaryOp { left, right, .. } => {
506 Self::contains_window_function(left) || Self::contains_window_function(right)
507 }
508 SqlExpression::Not { expr } => Self::contains_window_function(expr),
509 SqlExpression::FunctionCall { args, .. } => {
510 args.iter().any(Self::contains_window_function)
511 }
512 SqlExpression::CaseExpression {
513 when_branches,
514 else_branch,
515 } => {
516 when_branches.iter().any(|branch| {
517 Self::contains_window_function(&branch.condition)
518 || Self::contains_window_function(&branch.result)
519 }) || else_branch
520 .as_ref()
521 .map_or(false, |e| Self::contains_window_function(e))
522 }
523 SqlExpression::SimpleCaseExpression {
524 expr,
525 when_branches,
526 else_branch,
527 } => {
528 Self::contains_window_function(expr)
529 || when_branches.iter().any(|branch| {
530 Self::contains_window_function(&branch.value)
531 || Self::contains_window_function(&branch.result)
532 })
533 || else_branch
534 .as_ref()
535 .map_or(false, |e| Self::contains_window_function(e))
536 }
537 SqlExpression::InList { expr, values } => {
538 Self::contains_window_function(expr)
539 || values.iter().any(Self::contains_window_function)
540 }
541 SqlExpression::NotInList { expr, values } => {
542 Self::contains_window_function(expr)
543 || values.iter().any(Self::contains_window_function)
544 }
545 SqlExpression::Between { expr, lower, upper } => {
546 Self::contains_window_function(expr)
547 || Self::contains_window_function(lower)
548 || Self::contains_window_function(upper)
549 }
550 SqlExpression::InSubquery { expr, .. } => Self::contains_window_function(expr),
551 SqlExpression::NotInSubquery { expr, .. } => Self::contains_window_function(expr),
552 SqlExpression::MethodCall { args, .. } => {
553 args.iter().any(Self::contains_window_function)
554 }
555 SqlExpression::ChainedMethodCall { base, args, .. } => {
556 Self::contains_window_function(base)
557 || args.iter().any(Self::contains_window_function)
558 }
559 _ => false,
560 }
561 }
562
563 fn extract_window_specs(
565 items: &[SelectItem],
566 ) -> Vec<crate::data::batch_window_evaluator::WindowFunctionSpec> {
567 let mut specs = Vec::new();
568 for (idx, item) in items.iter().enumerate() {
569 if let SelectItem::Expression { expr, .. } = item {
570 Self::collect_window_function_specs(expr, idx, &mut specs);
571 }
572 }
573 specs
574 }
575
576 fn collect_window_function_specs(
578 expr: &SqlExpression,
579 output_column_index: usize,
580 specs: &mut Vec<crate::data::batch_window_evaluator::WindowFunctionSpec>,
581 ) {
582 match expr {
583 SqlExpression::WindowFunction {
584 name,
585 args,
586 window_spec,
587 } => {
588 specs.push(crate::data::batch_window_evaluator::WindowFunctionSpec {
589 spec: window_spec.clone(),
590 function_name: name.clone(),
591 args: args.clone(),
592 output_column_index,
593 });
594 }
595 SqlExpression::BinaryOp { left, right, .. } => {
596 Self::collect_window_function_specs(left, output_column_index, specs);
597 Self::collect_window_function_specs(right, output_column_index, specs);
598 }
599 SqlExpression::Not { expr } => {
600 Self::collect_window_function_specs(expr, output_column_index, specs);
601 }
602 SqlExpression::FunctionCall { args, .. } => {
603 for arg in args {
604 Self::collect_window_function_specs(arg, output_column_index, specs);
605 }
606 }
607 SqlExpression::CaseExpression {
608 when_branches,
609 else_branch,
610 } => {
611 for branch in when_branches {
612 Self::collect_window_function_specs(
613 &branch.condition,
614 output_column_index,
615 specs,
616 );
617 Self::collect_window_function_specs(&branch.result, output_column_index, specs);
618 }
619 if let Some(e) = else_branch {
620 Self::collect_window_function_specs(e, output_column_index, specs);
621 }
622 }
623 SqlExpression::SimpleCaseExpression {
624 expr,
625 when_branches,
626 else_branch,
627 } => {
628 Self::collect_window_function_specs(expr, output_column_index, specs);
629 for branch in when_branches {
630 Self::collect_window_function_specs(&branch.value, output_column_index, specs);
631 Self::collect_window_function_specs(&branch.result, output_column_index, specs);
632 }
633 if let Some(e) = else_branch {
634 Self::collect_window_function_specs(e, output_column_index, specs);
635 }
636 }
637 SqlExpression::InList { expr, values } => {
638 Self::collect_window_function_specs(expr, output_column_index, specs);
639 for val in values {
640 Self::collect_window_function_specs(val, output_column_index, specs);
641 }
642 }
643 SqlExpression::NotInList { expr, values } => {
644 Self::collect_window_function_specs(expr, output_column_index, specs);
645 for val in values {
646 Self::collect_window_function_specs(val, output_column_index, specs);
647 }
648 }
649 SqlExpression::Between { expr, lower, upper } => {
650 Self::collect_window_function_specs(expr, output_column_index, specs);
651 Self::collect_window_function_specs(lower, output_column_index, specs);
652 Self::collect_window_function_specs(upper, output_column_index, specs);
653 }
654 SqlExpression::InSubquery { expr, .. } => {
655 Self::collect_window_function_specs(expr, output_column_index, specs);
656 }
657 SqlExpression::NotInSubquery { expr, .. } => {
658 Self::collect_window_function_specs(expr, output_column_index, specs);
659 }
660 SqlExpression::MethodCall { args, .. } => {
661 for arg in args {
662 Self::collect_window_function_specs(arg, output_column_index, specs);
663 }
664 }
665 SqlExpression::ChainedMethodCall { base, args, .. } => {
666 Self::collect_window_function_specs(base, output_column_index, specs);
667 for arg in args {
668 Self::collect_window_function_specs(arg, output_column_index, specs);
669 }
670 }
671 _ => {} }
673 }
674
675 pub fn execute(&self, table: Arc<DataTable>, sql: &str) -> Result<DataView> {
677 let (view, _plan) = self.execute_with_plan(table, sql)?;
678 Ok(view)
679 }
680
681 pub fn execute_with_temp_tables(
683 &self,
684 table: Arc<DataTable>,
685 sql: &str,
686 temp_tables: Option<&TempTableRegistry>,
687 ) -> Result<DataView> {
688 let (view, _plan) = self.execute_with_plan_and_temp_tables(table, sql, temp_tables)?;
689 Ok(view)
690 }
691
692 pub fn execute_statement(
694 &self,
695 table: Arc<DataTable>,
696 statement: SelectStatement,
697 ) -> Result<DataView> {
698 self.execute_statement_with_temp_tables(table, statement, None)
699 }
700
701 pub fn execute_statement_with_temp_tables(
703 &self,
704 table: Arc<DataTable>,
705 statement: SelectStatement,
706 temp_tables: Option<&TempTableRegistry>,
707 ) -> Result<DataView> {
708 let mut cte_context = HashMap::new();
710
711 if let Some(temp_registry) = temp_tables {
713 for table_name in temp_registry.list_tables() {
714 if let Some(temp_table) = temp_registry.get(&table_name) {
715 debug!("Adding temp table {} to CTE context", table_name);
716 let view = DataView::new(temp_table);
717 cte_context.insert(table_name, Arc::new(view));
718 }
719 }
720 }
721
722 for cte in &statement.ctes {
723 debug!("QueryEngine: Pre-processing CTE '{}'...", cte.name);
724 let cte_result = match &cte.cte_type {
726 CTEType::Standard(query) => {
727 let view = self.build_view_with_context(
729 table.clone(),
730 query.clone(),
731 &mut cte_context,
732 )?;
733
734 let mut materialized = self.materialize_view(view)?;
736
737 for column in materialized.columns_mut() {
739 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
740 column.source_table = Some(cte.name.clone());
741 }
742
743 DataView::new(Arc::new(materialized))
744 }
745 CTEType::Web(web_spec) => {
746 use crate::web::http_fetcher::WebDataFetcher;
748
749 let fetcher = WebDataFetcher::new()?;
750 let mut data_table = fetcher.fetch(web_spec, &cte.name, None)?;
752
753 for column in data_table.columns_mut() {
755 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
756 column.source_table = Some(cte.name.clone());
757 }
758
759 DataView::new(Arc::new(data_table))
761 }
762 CTEType::File(file_spec) => {
763 let mut data_table =
764 crate::data::file_walker::walk_filesystem(file_spec, &cte.name)?;
765
766 for column in data_table.columns_mut() {
767 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
768 column.source_table = Some(cte.name.clone());
769 }
770
771 DataView::new(Arc::new(data_table))
772 }
773 };
774 cte_context.insert(cte.name.clone(), Arc::new(cte_result));
776 debug!(
777 "QueryEngine: CTE '{}' pre-processed, stored in context",
778 cte.name
779 );
780 }
781
782 let mut subquery_executor =
784 SubqueryExecutor::with_cte_context(self.clone(), table.clone(), cte_context.clone());
785 let processed_statement = subquery_executor.execute_subqueries(&statement)?;
786
787 self.build_view_with_context(table, processed_statement, &mut cte_context)
789 }
790
791 pub fn execute_statement_with_cte_context(
793 &self,
794 table: Arc<DataTable>,
795 statement: SelectStatement,
796 cte_context: &HashMap<String, Arc<DataView>>,
797 ) -> Result<DataView> {
798 let mut local_context = cte_context.clone();
800
801 for cte in &statement.ctes {
803 debug!("QueryEngine: Processing nested CTE '{}'...", cte.name);
804 let cte_result = match &cte.cte_type {
805 CTEType::Standard(query) => {
806 let view = self.build_view_with_context(
807 table.clone(),
808 query.clone(),
809 &mut local_context,
810 )?;
811
812 let mut materialized = self.materialize_view(view)?;
814
815 for column in materialized.columns_mut() {
817 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
818 column.source_table = Some(cte.name.clone());
819 }
820
821 DataView::new(Arc::new(materialized))
822 }
823 CTEType::Web(web_spec) => {
824 use crate::web::http_fetcher::WebDataFetcher;
826
827 let fetcher = WebDataFetcher::new()?;
828 let mut data_table = fetcher.fetch(web_spec, &cte.name, None)?;
830
831 for column in data_table.columns_mut() {
833 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
834 column.source_table = Some(cte.name.clone());
835 }
836
837 DataView::new(Arc::new(data_table))
839 }
840 CTEType::File(file_spec) => {
841 let mut data_table =
842 crate::data::file_walker::walk_filesystem(file_spec, &cte.name)?;
843
844 for column in data_table.columns_mut() {
845 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
846 column.source_table = Some(cte.name.clone());
847 }
848
849 DataView::new(Arc::new(data_table))
850 }
851 };
852 local_context.insert(cte.name.clone(), Arc::new(cte_result));
853 }
854
855 let mut subquery_executor =
857 SubqueryExecutor::with_cte_context(self.clone(), table.clone(), local_context.clone());
858 let processed_statement = subquery_executor.execute_subqueries(&statement)?;
859
860 self.build_view_with_context(table, processed_statement, &mut local_context)
862 }
863
864 pub fn execute_with_plan(
866 &self,
867 table: Arc<DataTable>,
868 sql: &str,
869 ) -> Result<(DataView, ExecutionPlan)> {
870 self.execute_with_plan_and_temp_tables(table, sql, None)
871 }
872
873 pub fn execute_with_plan_and_temp_tables(
875 &self,
876 table: Arc<DataTable>,
877 sql: &str,
878 temp_tables: Option<&TempTableRegistry>,
879 ) -> Result<(DataView, ExecutionPlan)> {
880 let mut plan_builder = ExecutionPlanBuilder::new();
881 let start_time = Instant::now();
882
883 plan_builder.begin_step(StepType::Parse, "Parse SQL query".to_string());
885 plan_builder.add_detail(format!("Query: {}", sql));
886 let mut parser = Parser::new(sql);
887 let statement = parser
888 .parse()
889 .map_err(|e| anyhow::anyhow!("Parse error: {}", e))?;
890 plan_builder.add_detail(format!("Parsed successfully"));
891 if let Some(ref from_source) = statement.from_source {
892 match from_source {
893 TableSource::Table(name) => {
894 plan_builder.add_detail(format!("FROM: {}", name));
895 }
896 TableSource::DerivedTable { alias, .. } => {
897 plan_builder.add_detail(format!("FROM: derived table (alias: {})", alias));
898 }
899 TableSource::Pivot { .. } => {
900 plan_builder.add_detail("FROM: PIVOT".to_string());
901 }
902 }
903 }
904 if statement.where_clause.is_some() {
905 plan_builder.add_detail("WHERE clause present".to_string());
906 }
907 plan_builder.end_step();
908
909 let mut cte_context = HashMap::new();
911
912 if let Some(temp_registry) = temp_tables {
914 for table_name in temp_registry.list_tables() {
915 if let Some(temp_table) = temp_registry.get(&table_name) {
916 debug!("Adding temp table {} to CTE context", table_name);
917 let view = DataView::new(temp_table);
918 cte_context.insert(table_name, Arc::new(view));
919 }
920 }
921 }
922
923 if !statement.ctes.is_empty() {
924 plan_builder.begin_step(
925 StepType::CTE,
926 format!("Process {} CTEs", statement.ctes.len()),
927 );
928
929 for cte in &statement.ctes {
930 let cte_start = Instant::now();
931 plan_builder.begin_step(StepType::CTE, format!("CTE '{}'", cte.name));
932
933 let cte_result = match &cte.cte_type {
934 CTEType::Standard(query) => {
935 if let Some(ref from_source) = query.from_source {
937 match from_source {
938 TableSource::Table(name) => {
939 plan_builder.add_detail(format!("Source: {}", name));
940 }
941 TableSource::DerivedTable { alias, .. } => {
942 plan_builder
943 .add_detail(format!("Source: derived table ({})", alias));
944 }
945 TableSource::Pivot { .. } => {
946 plan_builder.add_detail("Source: PIVOT".to_string());
947 }
948 }
949 }
950 if query.where_clause.is_some() {
951 plan_builder.add_detail("Has WHERE clause".to_string());
952 }
953 if query.group_by.is_some() {
954 plan_builder.add_detail("Has GROUP BY".to_string());
955 }
956
957 debug!(
958 "QueryEngine: Processing CTE '{}' with existing context: {:?}",
959 cte.name,
960 cte_context.keys().collect::<Vec<_>>()
961 );
962
963 let mut subquery_executor = SubqueryExecutor::with_cte_context(
966 self.clone(),
967 table.clone(),
968 cte_context.clone(),
969 );
970 let processed_query = subquery_executor.execute_subqueries(query)?;
971
972 let view = self.build_view_with_context(
973 table.clone(),
974 processed_query,
975 &mut cte_context,
976 )?;
977
978 let mut materialized = self.materialize_view(view)?;
980
981 for column in materialized.columns_mut() {
983 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
984 column.source_table = Some(cte.name.clone());
985 }
986
987 DataView::new(Arc::new(materialized))
988 }
989 CTEType::Web(web_spec) => {
990 plan_builder.add_detail(format!("URL: {}", web_spec.url));
991 if let Some(format) = &web_spec.format {
992 plan_builder.add_detail(format!("Format: {:?}", format));
993 }
994 if let Some(cache) = web_spec.cache_seconds {
995 plan_builder.add_detail(format!("Cache: {} seconds", cache));
996 }
997
998 use crate::web::http_fetcher::WebDataFetcher;
1000
1001 let fetcher = WebDataFetcher::new()?;
1002 let mut data_table = fetcher.fetch(web_spec, &cte.name, None)?;
1004
1005 for column in data_table.columns_mut() {
1007 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
1008 column.source_table = Some(cte.name.clone());
1009 }
1010
1011 DataView::new(Arc::new(data_table))
1013 }
1014 CTEType::File(file_spec) => {
1015 plan_builder.add_detail(format!("PATH: {}", file_spec.path));
1016 if file_spec.recursive {
1017 plan_builder.add_detail("RECURSIVE".to_string());
1018 }
1019 if let Some(ref g) = file_spec.glob {
1020 plan_builder.add_detail(format!("GLOB: {}", g));
1021 }
1022 if let Some(d) = file_spec.max_depth {
1023 plan_builder.add_detail(format!("MAX_DEPTH: {}", d));
1024 }
1025
1026 let mut data_table =
1027 crate::data::file_walker::walk_filesystem(file_spec, &cte.name)?;
1028
1029 for column in data_table.columns_mut() {
1030 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
1031 column.source_table = Some(cte.name.clone());
1032 }
1033
1034 DataView::new(Arc::new(data_table))
1035 }
1036 };
1037
1038 plan_builder.set_rows_out(cte_result.row_count());
1040 plan_builder.add_detail(format!(
1041 "Result: {} rows, {} columns",
1042 cte_result.row_count(),
1043 cte_result.column_count()
1044 ));
1045 plan_builder.add_detail(format!(
1046 "Execution time: {:.3}ms",
1047 cte_start.elapsed().as_secs_f64() * 1000.0
1048 ));
1049
1050 debug!(
1051 "QueryEngine: Storing CTE '{}' in context with {} rows",
1052 cte.name,
1053 cte_result.row_count()
1054 );
1055 cte_context.insert(cte.name.clone(), Arc::new(cte_result));
1056 plan_builder.end_step();
1057 }
1058
1059 plan_builder.add_detail(format!(
1060 "All {} CTEs cached in context",
1061 statement.ctes.len()
1062 ));
1063 plan_builder.end_step();
1064 }
1065
1066 plan_builder.begin_step(StepType::Subquery, "Process subqueries".to_string());
1068 let mut subquery_executor =
1069 SubqueryExecutor::with_cte_context(self.clone(), table.clone(), cte_context.clone());
1070
1071 let has_subqueries = statement.where_clause.as_ref().map_or(false, |w| {
1073 format!("{:?}", w).contains("Subquery")
1075 });
1076
1077 if has_subqueries {
1078 plan_builder.add_detail("Evaluating subqueries in WHERE clause".to_string());
1079 }
1080
1081 let processed_statement = subquery_executor.execute_subqueries(&statement)?;
1082
1083 if has_subqueries {
1084 plan_builder.add_detail("Subqueries replaced with materialized values".to_string());
1085 } else {
1086 plan_builder.add_detail("No subqueries to process".to_string());
1087 }
1088
1089 plan_builder.end_step();
1090 let result = self.build_view_with_context_and_plan(
1091 table,
1092 processed_statement,
1093 &mut cte_context,
1094 &mut plan_builder,
1095 )?;
1096
1097 let total_duration = start_time.elapsed();
1098 info!(
1099 "Query execution complete: total={:?}, rows={}",
1100 total_duration,
1101 result.row_count()
1102 );
1103
1104 let plan = plan_builder.build();
1105 Ok((result, plan))
1106 }
1107
1108 fn build_view(&self, table: Arc<DataTable>, statement: SelectStatement) -> Result<DataView> {
1110 let mut cte_context = HashMap::new();
1111 self.build_view_with_context(table, statement, &mut cte_context)
1112 }
1113
1114 fn build_view_with_context(
1116 &self,
1117 table: Arc<DataTable>,
1118 statement: SelectStatement,
1119 cte_context: &mut HashMap<String, Arc<DataView>>,
1120 ) -> Result<DataView> {
1121 let mut dummy_plan = ExecutionPlanBuilder::new();
1122 let mut exec_context = ExecutionContext::new();
1123 self.build_view_with_context_and_plan_and_exec(
1124 table,
1125 statement,
1126 cte_context,
1127 &mut dummy_plan,
1128 &mut exec_context,
1129 )
1130 }
1131
1132 fn build_view_with_context_and_plan(
1134 &self,
1135 table: Arc<DataTable>,
1136 statement: SelectStatement,
1137 cte_context: &mut HashMap<String, Arc<DataView>>,
1138 plan: &mut ExecutionPlanBuilder,
1139 ) -> Result<DataView> {
1140 let mut exec_context = ExecutionContext::new();
1141 self.build_view_with_context_and_plan_and_exec(
1142 table,
1143 statement,
1144 cte_context,
1145 plan,
1146 &mut exec_context,
1147 )
1148 }
1149
1150 fn build_view_with_context_and_plan_and_exec(
1152 &self,
1153 table: Arc<DataTable>,
1154 statement: SelectStatement,
1155 cte_context: &mut HashMap<String, Arc<DataView>>,
1156 plan: &mut ExecutionPlanBuilder,
1157 exec_context: &mut ExecutionContext,
1158 ) -> Result<DataView> {
1159 for cte in &statement.ctes {
1161 if cte_context.contains_key(&cte.name) {
1163 debug!(
1164 "QueryEngine: CTE '{}' already in context, skipping",
1165 cte.name
1166 );
1167 continue;
1168 }
1169
1170 debug!("QueryEngine: Processing CTE '{}'...", cte.name);
1171 debug!(
1172 "QueryEngine: Available CTEs for '{}': {:?}",
1173 cte.name,
1174 cte_context.keys().collect::<Vec<_>>()
1175 );
1176
1177 let cte_result = match &cte.cte_type {
1179 CTEType::Standard(query) => {
1180 let view =
1181 self.build_view_with_context(table.clone(), query.clone(), cte_context)?;
1182
1183 let mut materialized = self.materialize_view(view)?;
1185
1186 for column in materialized.columns_mut() {
1188 column.qualified_name = Some(format!("{}.{}", cte.name, column.name));
1189 column.source_table = Some(cte.name.clone());
1190 }
1191
1192 DataView::new(Arc::new(materialized))
1193 }
1194 CTEType::Web(_web_spec) => {
1195 return Err(anyhow!(
1197 "Web CTEs should be processed in execute_select method"
1198 ));
1199 }
1200 CTEType::File(_file_spec) => {
1201 return Err(anyhow!(
1203 "FILE CTEs should be processed in execute_select method"
1204 ));
1205 }
1206 };
1207
1208 cte_context.insert(cte.name.clone(), Arc::new(cte_result));
1210 debug!(
1211 "QueryEngine: CTE '{}' processed, stored in context",
1212 cte.name
1213 );
1214 }
1215
1216 let source_table = if let Some(ref from_source) = statement.from_source {
1218 match from_source {
1219 TableSource::Table(table_name) => {
1220 if let Some(cte_view) = resolve_cte(cte_context, table_name) {
1222 debug!("QueryEngine: Using CTE '{}' as source table", table_name);
1223 let mut materialized = self.materialize_view((**cte_view).clone())?;
1225
1226 #[allow(deprecated)]
1228 if let Some(ref alias) = statement.from_alias {
1229 debug!(
1230 "QueryEngine: Applying alias '{}' to CTE '{}' qualified column names",
1231 alias, table_name
1232 );
1233 for column in materialized.columns_mut() {
1234 if let Some(ref qualified_name) = column.qualified_name {
1236 if qualified_name.starts_with(&format!("{}.", table_name)) {
1237 column.qualified_name = Some(qualified_name.replace(
1238 &format!("{}.", table_name),
1239 &format!("{}.", alias),
1240 ));
1241 }
1242 }
1243 if column.source_table.as_ref() == Some(table_name) {
1245 column.source_table = Some(alias.clone());
1246 }
1247 }
1248 }
1249
1250 Arc::new(materialized)
1251 } else {
1252 table.clone()
1254 }
1255 }
1256 TableSource::DerivedTable { query, alias } => {
1257 debug!(
1259 "QueryEngine: Processing FROM derived table (alias: {})",
1260 alias
1261 );
1262 let subquery_result =
1263 self.build_view_with_context(table.clone(), *query.clone(), cte_context)?;
1264
1265 let mut materialized = self.materialize_view(subquery_result)?;
1268
1269 for column in materialized.columns_mut() {
1273 column.source_table = Some(alias.clone());
1274 }
1275
1276 Arc::new(materialized)
1277 }
1278 TableSource::Pivot { .. } => {
1279 return Err(anyhow!(
1281 "PIVOT in FROM clause should have been expanded by preprocessing pipeline"
1282 ));
1283 }
1284 }
1285 } else {
1286 #[allow(deprecated)]
1288 if let Some(ref table_func) = statement.from_function {
1289 debug!("QueryEngine: Processing table function (deprecated field)...");
1291 match table_func {
1292 TableFunction::Generator { name, args } => {
1293 use crate::sql::generators::GeneratorRegistry;
1295
1296 let registry = GeneratorRegistry::new();
1298
1299 if let Some(generator) = registry.get(name) {
1300 let mut evaluator = ArithmeticEvaluator::with_date_notation(
1302 &table,
1303 self.date_notation.clone(),
1304 );
1305 let dummy_row = 0;
1306
1307 let mut evaluated_args = Vec::new();
1308 for arg in args {
1309 evaluated_args.push(evaluator.evaluate(arg, dummy_row)?);
1310 }
1311
1312 generator.generate(evaluated_args)?
1314 } else {
1315 return Err(anyhow!("Unknown generator function: {}", name));
1316 }
1317 }
1318 }
1319 } else {
1320 #[allow(deprecated)]
1321 if let Some(ref subquery) = statement.from_subquery {
1322 debug!("QueryEngine: Processing FROM subquery (deprecated field)...");
1324 let subquery_result = self.build_view_with_context(
1325 table.clone(),
1326 *subquery.clone(),
1327 cte_context,
1328 )?;
1329
1330 let materialized = self.materialize_view(subquery_result)?;
1333 Arc::new(materialized)
1334 } else {
1335 #[allow(deprecated)]
1336 if let Some(ref table_name) = statement.from_table {
1337 if let Some(cte_view) = resolve_cte(cte_context, table_name) {
1339 debug!(
1340 "QueryEngine: Using CTE '{}' as source table (deprecated field)",
1341 table_name
1342 );
1343 let mut materialized = self.materialize_view((**cte_view).clone())?;
1345
1346 #[allow(deprecated)]
1348 if let Some(ref alias) = statement.from_alias {
1349 debug!(
1350 "QueryEngine: Applying alias '{}' to CTE '{}' qualified column names",
1351 alias, table_name
1352 );
1353 for column in materialized.columns_mut() {
1354 if let Some(ref qualified_name) = column.qualified_name {
1356 if qualified_name.starts_with(&format!("{}.", table_name)) {
1357 column.qualified_name = Some(qualified_name.replace(
1358 &format!("{}.", table_name),
1359 &format!("{}.", alias),
1360 ));
1361 }
1362 }
1363 if column.source_table.as_ref() == Some(table_name) {
1365 column.source_table = Some(alias.clone());
1366 }
1367 }
1368 }
1369
1370 Arc::new(materialized)
1371 } else {
1372 table.clone()
1374 }
1375 } else {
1376 Arc::new(DataTable::dual())
1381 }
1382 }
1383 }
1384 };
1385
1386 #[allow(deprecated)]
1388 if let Some(ref alias) = statement.from_alias {
1389 #[allow(deprecated)]
1390 if let Some(ref table_name) = statement.from_table {
1391 exec_context.register_alias(alias.clone(), table_name.clone());
1392 }
1393 }
1394
1395 let final_table = if !statement.joins.is_empty() {
1397 plan.begin_step(
1398 StepType::Join,
1399 format!("Process {} JOINs", statement.joins.len()),
1400 );
1401 plan.set_rows_in(source_table.row_count());
1402
1403 let join_executor = HashJoinExecutor::new(self.case_insensitive);
1404
1405 #[allow(deprecated)]
1409 let base_table_name = match statement.from_source {
1410 Some(TableSource::Table(ref n)) => Some(n.clone()),
1411 _ => statement.from_table.clone(),
1412 };
1413
1414 let mut current_table = source_table;
1415
1416 for (idx, join_clause) in statement.joins.iter().enumerate() {
1417 let join_start = Instant::now();
1418 plan.begin_step(StepType::Join, format!("JOIN #{}", idx + 1));
1419 plan.add_detail(format!("Type: {:?}", join_clause.join_type));
1420 plan.add_detail(format!("Left table: {} rows", current_table.row_count()));
1421 plan.add_detail(format!(
1422 "Executing {:?} JOIN on {} condition(s)",
1423 join_clause.join_type,
1424 join_clause.condition.conditions.len()
1425 ));
1426
1427 let right_table = match &join_clause.table {
1429 TableSource::Table(name) => {
1430 if let Some(cte_view) = resolve_cte(cte_context, name) {
1432 let mut materialized = self.materialize_view((**cte_view).clone())?;
1433
1434 if let Some(ref alias) = join_clause.alias {
1436 debug!("QueryEngine: Applying JOIN alias '{}' to CTE '{}' qualified column names", alias, name);
1437 for column in materialized.columns_mut() {
1438 if let Some(ref qualified_name) = column.qualified_name {
1440 if qualified_name.starts_with(&format!("{}.", name)) {
1441 column.qualified_name = Some(qualified_name.replace(
1442 &format!("{}.", name),
1443 &format!("{}.", alias),
1444 ));
1445 }
1446 }
1447 if column.source_table.as_ref() == Some(name) {
1449 column.source_table = Some(alias.clone());
1450 }
1451 }
1452 }
1453
1454 Arc::new(materialized)
1455 } else if base_table_name.as_deref().is_some_and(|base| {
1456 if self.case_insensitive {
1457 base.eq_ignore_ascii_case(name)
1458 } else {
1459 base == name
1460 }
1461 }) {
1462 let mut materialized = (*table).clone();
1468 if let Some(ref alias) = join_clause.alias {
1469 for column in materialized.columns_mut() {
1470 if let Some(ref qualified_name) = column.qualified_name {
1471 if qualified_name.starts_with(&format!("{}.", name)) {
1472 column.qualified_name = Some(qualified_name.replace(
1473 &format!("{}.", name),
1474 &format!("{}.", alias),
1475 ));
1476 }
1477 }
1478 if column.source_table.as_ref() == Some(name) {
1479 column.source_table = Some(alias.clone());
1480 }
1481 }
1482 }
1483 Arc::new(materialized)
1484 } else {
1485 return Err(anyhow!("Cannot resolve table '{}' for JOIN", name));
1488 }
1489 }
1490 TableSource::DerivedTable { query, alias: _ } => {
1491 let subquery_result = self.build_view_with_context(
1493 table.clone(),
1494 *query.clone(),
1495 cte_context,
1496 )?;
1497 let materialized = self.materialize_view(subquery_result)?;
1498 Arc::new(materialized)
1499 }
1500 TableSource::Pivot { .. } => {
1501 return Err(anyhow!("PIVOT in JOIN clause is not yet supported"));
1503 }
1504 };
1505
1506 let joined = join_executor.execute_join(
1508 current_table.clone(),
1509 join_clause,
1510 right_table.clone(),
1511 )?;
1512
1513 plan.add_detail(format!("Right table: {} rows", right_table.row_count()));
1514 plan.set_rows_out(joined.row_count());
1515 plan.add_detail(format!("Result: {} rows", joined.row_count()));
1516 plan.add_detail(format!(
1517 "Join time: {:.3}ms",
1518 join_start.elapsed().as_secs_f64() * 1000.0
1519 ));
1520 plan.end_step();
1521
1522 current_table = Arc::new(joined);
1523 }
1524
1525 plan.set_rows_out(current_table.row_count());
1526 plan.add_detail(format!(
1527 "Final result after all joins: {} rows",
1528 current_table.row_count()
1529 ));
1530 plan.end_step();
1531 current_table
1532 } else {
1533 source_table
1534 };
1535
1536 self.build_view_internal_with_plan_and_exec(
1538 final_table,
1539 statement,
1540 plan,
1541 Some(exec_context),
1542 )
1543 }
1544
1545 pub fn materialize_view(&self, view: DataView) -> Result<DataTable> {
1547 let source = view.source();
1548 let mut result_table = DataTable::new("derived");
1549
1550 let visible_cols = view.visible_column_indices().to_vec();
1552
1553 for col_idx in &visible_cols {
1555 let col = &source.columns[*col_idx];
1556 let new_col = DataColumn {
1557 name: col.name.clone(),
1558 data_type: col.data_type.clone(),
1559 nullable: col.nullable,
1560 unique_values: col.unique_values,
1561 null_count: col.null_count,
1562 metadata: col.metadata.clone(),
1563 qualified_name: col.qualified_name.clone(), source_table: col.source_table.clone(), };
1566 result_table.add_column(new_col);
1567 }
1568
1569 for row_idx in view.visible_row_indices() {
1571 let source_row = &source.rows[*row_idx];
1572 let mut new_row = DataRow { values: Vec::new() };
1573
1574 for col_idx in &visible_cols {
1575 new_row.values.push(source_row.values[*col_idx].clone());
1576 }
1577
1578 result_table.add_row(new_row);
1579 }
1580
1581 Ok(result_table)
1582 }
1583
1584 fn build_view_internal(
1585 &self,
1586 table: Arc<DataTable>,
1587 statement: SelectStatement,
1588 ) -> Result<DataView> {
1589 let mut dummy_plan = ExecutionPlanBuilder::new();
1590 self.build_view_internal_with_plan(table, statement, &mut dummy_plan)
1591 }
1592
1593 fn build_view_internal_with_plan(
1594 &self,
1595 table: Arc<DataTable>,
1596 statement: SelectStatement,
1597 plan: &mut ExecutionPlanBuilder,
1598 ) -> Result<DataView> {
1599 self.build_view_internal_with_plan_and_exec(table, statement, plan, None)
1600 }
1601
1602 fn build_view_internal_with_plan_and_exec(
1603 &self,
1604 table: Arc<DataTable>,
1605 statement: SelectStatement,
1606 plan: &mut ExecutionPlanBuilder,
1607 exec_context: Option<&ExecutionContext>,
1608 ) -> Result<DataView> {
1609 debug!(
1610 "QueryEngine::build_view - select_items: {:?}",
1611 statement.select_items
1612 );
1613 debug!(
1614 "QueryEngine::build_view - where_clause: {:?}",
1615 statement.where_clause
1616 );
1617
1618 let mut visible_rows: Vec<usize> = (0..table.row_count()).collect();
1620
1621 if let Some(where_clause) = &statement.where_clause {
1623 let total_rows = table.row_count();
1624 debug!("QueryEngine: Applying WHERE clause to {} rows", total_rows);
1625 debug!("QueryEngine: WHERE clause = {:?}", where_clause);
1626
1627 plan.begin_step(StepType::Filter, "WHERE clause filtering".to_string());
1628 plan.set_rows_in(total_rows);
1629 plan.add_detail(format!("Input: {} rows", total_rows));
1630
1631 for condition in &where_clause.conditions {
1633 plan.add_detail(format!("Condition: {:?}", condition.expr));
1634 }
1635
1636 let filter_start = Instant::now();
1637 let mut eval_context = EvaluationContext::new(self.case_insensitive);
1639
1640 let mut evaluator = if let Some(exec_ctx) = exec_context {
1642 RecursiveWhereEvaluator::with_both_contexts(&table, &mut eval_context, exec_ctx)
1644 } else {
1645 RecursiveWhereEvaluator::with_context(&table, &mut eval_context)
1646 };
1647
1648 let mut filtered_rows = Vec::new();
1650 for row_idx in visible_rows {
1651 if row_idx < 3 {
1653 debug!("QueryEngine: Evaluating WHERE clause for row {}", row_idx);
1654 }
1655
1656 match evaluator.evaluate(where_clause, row_idx) {
1657 Ok(result) => {
1658 if row_idx < 3 {
1659 debug!("QueryEngine: Row {} WHERE result: {}", row_idx, result);
1660 }
1661 if result {
1662 filtered_rows.push(row_idx);
1663 }
1664 }
1665 Err(e) => {
1666 if row_idx < 3 {
1667 debug!(
1668 "QueryEngine: WHERE evaluation error for row {}: {}",
1669 row_idx, e
1670 );
1671 }
1672 return Err(e);
1674 }
1675 }
1676 }
1677
1678 let (compilations, cache_hits) = eval_context.get_stats();
1680 if compilations > 0 || cache_hits > 0 {
1681 debug!(
1682 "LIKE pattern cache: {} compilations, {} cache hits",
1683 compilations, cache_hits
1684 );
1685 }
1686 visible_rows = filtered_rows;
1687 let filter_duration = filter_start.elapsed();
1688 info!(
1689 "WHERE clause filtering: {} rows -> {} rows in {:?}",
1690 total_rows,
1691 visible_rows.len(),
1692 filter_duration
1693 );
1694
1695 plan.set_rows_out(visible_rows.len());
1696 plan.add_detail(format!("Output: {} rows", visible_rows.len()));
1697 plan.add_detail(format!(
1698 "Filter time: {:.3}ms",
1699 filter_duration.as_secs_f64() * 1000.0
1700 ));
1701 plan.end_step();
1702 }
1703
1704 let mut view = DataView::new(table.clone());
1706 view = view.with_rows(visible_rows);
1707
1708 if let Some(group_by_exprs) = &statement.group_by {
1710 if !group_by_exprs.is_empty() {
1711 debug!("QueryEngine: Processing GROUP BY: {:?}", group_by_exprs);
1712
1713 plan.begin_step(
1714 StepType::GroupBy,
1715 format!("GROUP BY {} expressions", group_by_exprs.len()),
1716 );
1717 plan.set_rows_in(view.row_count());
1718 plan.add_detail(format!("Input: {} rows", view.row_count()));
1719 for expr in group_by_exprs {
1720 plan.add_detail(format!("Group by: {:?}", expr));
1721 }
1722
1723 let group_start = Instant::now();
1724 view = self.apply_group_by(
1725 view,
1726 group_by_exprs,
1727 &statement.select_items,
1728 statement.having.as_ref(),
1729 plan,
1730 )?;
1731
1732 use crate::query_plan::having_alias_transformer::HIDDEN_AGG_PREFIX;
1735 let hidden_indices: Vec<usize> = view
1736 .source()
1737 .columns
1738 .iter()
1739 .enumerate()
1740 .filter_map(|(i, c)| {
1741 if c.name.starts_with(HIDDEN_AGG_PREFIX) {
1742 Some(i)
1743 } else {
1744 None
1745 }
1746 })
1747 .collect();
1748 for &idx in hidden_indices.iter().rev() {
1749 view.hide_column(idx);
1750 }
1751
1752 plan.set_rows_out(view.row_count());
1753 plan.add_detail(format!("Output: {} groups", view.row_count()));
1754 plan.add_detail(format!(
1755 "Overall time: {:.3}ms",
1756 group_start.elapsed().as_secs_f64() * 1000.0
1757 ));
1758 plan.end_step();
1759 }
1760 } else {
1761 if !statement.select_items.is_empty() {
1763 let has_non_star_items = statement
1765 .select_items
1766 .iter()
1767 .any(|item| !matches!(item, SelectItem::Star { .. }));
1768
1769 if has_non_star_items || statement.select_items.len() > 1 {
1773 view = self.apply_select_items(
1774 view,
1775 &statement.select_items,
1776 &statement,
1777 exec_context,
1778 plan,
1779 )?;
1780 }
1781 } else if !statement.columns.is_empty() && statement.columns[0] != "*" {
1783 debug!("QueryEngine: Using legacy columns path");
1784 let source_table = view.source();
1787 let column_indices =
1788 self.resolve_column_indices(source_table, &statement.columns)?;
1789 view = view.with_columns(column_indices);
1790 }
1791 }
1792
1793 if statement.distinct {
1795 plan.begin_step(StepType::Distinct, "Remove duplicate rows".to_string());
1796 plan.set_rows_in(view.row_count());
1797 plan.add_detail(format!("Input: {} rows", view.row_count()));
1798
1799 let distinct_start = Instant::now();
1800 view = self.apply_distinct(view)?;
1801
1802 plan.set_rows_out(view.row_count());
1803 plan.add_detail(format!("Output: {} unique rows", view.row_count()));
1804 plan.add_detail(format!(
1805 "Distinct time: {:.3}ms",
1806 distinct_start.elapsed().as_secs_f64() * 1000.0
1807 ));
1808 plan.end_step();
1809 }
1810
1811 if let Some(order_by_columns) = &statement.order_by {
1813 if !order_by_columns.is_empty() {
1814 plan.begin_step(
1815 StepType::Sort,
1816 format!("ORDER BY {} columns", order_by_columns.len()),
1817 );
1818 plan.set_rows_in(view.row_count());
1819 for col in order_by_columns {
1820 let expr_str = match &col.expr {
1822 SqlExpression::Column(col_ref) => col_ref.name.clone(),
1823 _ => "expr".to_string(),
1824 };
1825 plan.add_detail(format!("{} {:?}", expr_str, col.direction));
1826 }
1827
1828 let sort_start = Instant::now();
1829 view =
1830 self.apply_multi_order_by_with_context(view, order_by_columns, exec_context)?;
1831
1832 plan.add_detail(format!(
1833 "Sort time: {:.3}ms",
1834 sort_start.elapsed().as_secs_f64() * 1000.0
1835 ));
1836 plan.end_step();
1837 }
1838 }
1839
1840 {
1845 use crate::query_plan::order_by_alias_transformer::HIDDEN_ORDERBY_PREFIX;
1846 let hidden_indices: Vec<usize> = view
1847 .source()
1848 .columns
1849 .iter()
1850 .enumerate()
1851 .filter_map(|(i, c)| {
1852 if c.name.starts_with(HIDDEN_ORDERBY_PREFIX) {
1853 Some(i)
1854 } else {
1855 None
1856 }
1857 })
1858 .collect();
1859 for &idx in hidden_indices.iter().rev() {
1860 view.hide_column(idx);
1861 }
1862 }
1863
1864 if let Some(limit) = statement.limit {
1866 let offset = statement.offset.unwrap_or(0);
1867 plan.begin_step(StepType::Limit, format!("LIMIT {}", limit));
1868 plan.set_rows_in(view.row_count());
1869 if offset > 0 {
1870 plan.add_detail(format!("OFFSET: {}", offset));
1871 }
1872 view = view.with_limit(limit, offset);
1873 plan.set_rows_out(view.row_count());
1874 plan.add_detail(format!("Output: {} rows", view.row_count()));
1875 plan.end_step();
1876 }
1877
1878 if !statement.set_operations.is_empty() {
1880 plan.begin_step(
1881 StepType::SetOperation,
1882 format!("Process {} set operations", statement.set_operations.len()),
1883 );
1884 plan.set_rows_in(view.row_count());
1885
1886 let mut combined_table = self.materialize_view(view)?;
1888 let first_columns = combined_table.column_names();
1889 let first_column_count = first_columns.len();
1890
1891 let mut needs_deduplication = false;
1893
1894 for (idx, (operation, next_statement)) in statement.set_operations.iter().enumerate() {
1896 let op_start = Instant::now();
1897 plan.begin_step(
1898 StepType::SetOperation,
1899 format!("{:?} operation #{}", operation, idx + 1),
1900 );
1901
1902 let next_view = if let Some(exec_ctx) = exec_context {
1905 self.build_view_internal_with_plan_and_exec(
1906 table.clone(),
1907 *next_statement.clone(),
1908 plan,
1909 Some(exec_ctx),
1910 )?
1911 } else {
1912 self.build_view_internal_with_plan(
1913 table.clone(),
1914 *next_statement.clone(),
1915 plan,
1916 )?
1917 };
1918
1919 let next_table = self.materialize_view(next_view)?;
1921 let next_columns = next_table.column_names();
1922 let next_column_count = next_columns.len();
1923
1924 if first_column_count != next_column_count {
1926 return Err(anyhow!(
1927 "UNION queries must have the same number of columns: first query has {} columns, but query #{} has {} columns",
1928 first_column_count,
1929 idx + 2,
1930 next_column_count
1931 ));
1932 }
1933
1934 for (col_idx, (first_col, next_col)) in
1936 first_columns.iter().zip(next_columns.iter()).enumerate()
1937 {
1938 if !first_col.eq_ignore_ascii_case(next_col) {
1939 debug!(
1940 "UNION column name mismatch at position {}: '{}' vs '{}' (using first query's name)",
1941 col_idx + 1,
1942 first_col,
1943 next_col
1944 );
1945 }
1946 }
1947
1948 plan.add_detail(format!("Left: {} rows", combined_table.row_count()));
1949 plan.add_detail(format!("Right: {} rows", next_table.row_count()));
1950
1951 match operation {
1953 SetOperation::UnionAll => {
1954 for row in next_table.rows.iter() {
1956 combined_table.add_row(row.clone());
1957 }
1958 plan.add_detail(format!(
1959 "Result: {} rows (no deduplication)",
1960 combined_table.row_count()
1961 ));
1962 }
1963 SetOperation::Union => {
1964 for row in next_table.rows.iter() {
1966 combined_table.add_row(row.clone());
1967 }
1968 needs_deduplication = true;
1969 plan.add_detail(format!(
1970 "Combined: {} rows (deduplication pending)",
1971 combined_table.row_count()
1972 ));
1973 }
1974 SetOperation::Intersect => {
1975 return Err(anyhow!("INTERSECT is not yet implemented"));
1978 }
1979 SetOperation::Except => {
1980 return Err(anyhow!("EXCEPT is not yet implemented"));
1983 }
1984 }
1985
1986 plan.add_detail(format!(
1987 "Operation time: {:.3}ms",
1988 op_start.elapsed().as_secs_f64() * 1000.0
1989 ));
1990 plan.set_rows_out(combined_table.row_count());
1991 plan.end_step();
1992 }
1993
1994 plan.set_rows_out(combined_table.row_count());
1995 plan.add_detail(format!(
1996 "Combined result: {} rows after {} operations",
1997 combined_table.row_count(),
1998 statement.set_operations.len()
1999 ));
2000 plan.end_step();
2001
2002 view = DataView::new(Arc::new(combined_table));
2004
2005 if needs_deduplication {
2007 plan.begin_step(
2008 StepType::Distinct,
2009 "UNION deduplication - remove duplicate rows".to_string(),
2010 );
2011 plan.set_rows_in(view.row_count());
2012 plan.add_detail(format!("Input: {} rows", view.row_count()));
2013
2014 let distinct_start = Instant::now();
2015 view = self.apply_distinct(view)?;
2016
2017 plan.set_rows_out(view.row_count());
2018 plan.add_detail(format!("Output: {} unique rows", view.row_count()));
2019 plan.add_detail(format!(
2020 "Deduplication time: {:.3}ms",
2021 distinct_start.elapsed().as_secs_f64() * 1000.0
2022 ));
2023 plan.end_step();
2024 }
2025 }
2026
2027 Ok(view)
2028 }
2029
2030 fn resolve_column_indices(&self, table: &DataTable, columns: &[String]) -> Result<Vec<usize>> {
2032 let mut indices = Vec::new();
2033 let table_columns = table.column_names();
2034
2035 for col_name in columns {
2036 let index = table_columns
2037 .iter()
2038 .position(|c| c.eq_ignore_ascii_case(col_name))
2039 .ok_or_else(|| {
2040 let suggestion = self.find_similar_column(table, col_name);
2041 match suggestion {
2042 Some(similar) => anyhow::anyhow!(
2043 "Column '{}' not found. Did you mean '{}'?",
2044 col_name,
2045 similar
2046 ),
2047 None => anyhow::anyhow!("Column '{}' not found", col_name),
2048 }
2049 })?;
2050 indices.push(index);
2051 }
2052
2053 Ok(indices)
2054 }
2055
2056 fn apply_select_items(
2058 &self,
2059 view: DataView,
2060 select_items: &[SelectItem],
2061 _statement: &SelectStatement,
2062 exec_context: Option<&ExecutionContext>,
2063 plan: &mut ExecutionPlanBuilder,
2064 ) -> Result<DataView> {
2065 debug!(
2066 "QueryEngine::apply_select_items - items: {:?}",
2067 select_items
2068 );
2069 debug!(
2070 "QueryEngine::apply_select_items - input view has {} rows",
2071 view.row_count()
2072 );
2073
2074 let has_window_functions = select_items.iter().any(|item| match item {
2076 SelectItem::Expression { expr, .. } => Self::contains_window_function(expr),
2077 _ => false,
2078 });
2079
2080 let window_func_count: usize = select_items
2082 .iter()
2083 .filter(|item| match item {
2084 SelectItem::Expression { expr, .. } => Self::contains_window_function(expr),
2085 _ => false,
2086 })
2087 .count();
2088
2089 let window_start = if has_window_functions {
2091 debug!(
2092 "QueryEngine::apply_select_items - detected {} window functions",
2093 window_func_count
2094 );
2095
2096 let window_specs = Self::extract_window_specs(select_items);
2098 debug!("Extracted {} window function specs", window_specs.len());
2099
2100 Some(Instant::now())
2101 } else {
2102 None
2103 };
2104
2105 let has_unnest = select_items.iter().any(|item| match item {
2107 SelectItem::Expression { expr, .. } => Self::contains_unnest(expr),
2108 _ => false,
2109 });
2110
2111 if has_unnest {
2112 debug!("QueryEngine::apply_select_items - UNNEST detected, using row expansion");
2113 return self.apply_select_with_row_expansion(view, select_items);
2114 }
2115
2116 let has_aggregates = select_items.iter().any(|item| match item {
2120 SelectItem::Expression { expr, .. } => contains_aggregate(expr),
2121 SelectItem::Column { .. } => false,
2122 SelectItem::Star { .. } => false,
2123 SelectItem::StarExclude { .. } => false,
2124 });
2125
2126 let all_aggregate_compatible = select_items.iter().all(|item| match item {
2127 SelectItem::Expression { expr, .. } => is_aggregate_compatible(expr),
2128 SelectItem::Column { .. } => false, SelectItem::Star { .. } => false, SelectItem::StarExclude { .. } => false, });
2132
2133 if has_aggregates && all_aggregate_compatible && view.row_count() > 0 {
2134 debug!("QueryEngine::apply_select_items - detected aggregate query with constants");
2137 return self.apply_aggregate_select(view, select_items);
2138 }
2139
2140 let has_computed_expressions = select_items
2142 .iter()
2143 .any(|item| matches!(item, SelectItem::Expression { .. }));
2144
2145 debug!(
2146 "QueryEngine::apply_select_items - has_computed_expressions: {}",
2147 has_computed_expressions
2148 );
2149
2150 if !has_computed_expressions {
2151 let column_indices = self.resolve_select_columns(view.source(), select_items)?;
2153 return Ok(view.with_columns(column_indices));
2154 }
2155
2156 let source_table = view.source();
2161 let visible_rows = view.visible_row_indices();
2162
2163 let mut computed_table = DataTable::new("query_result");
2166
2167 let mut expanded_items = Vec::new();
2169 for item in select_items {
2170 match item {
2171 SelectItem::Star { table_prefix, .. } => {
2172 if let Some(prefix) = table_prefix {
2173 debug!("QueryEngine::apply_select_items - expanding {}.*", prefix);
2175 for col in &source_table.columns {
2176 if Self::column_matches_table(col, prefix) {
2177 expanded_items.push(SelectItem::Column {
2178 column: ColumnRef::unquoted(col.name.clone()),
2179 leading_comments: vec![],
2180 trailing_comment: None,
2181 });
2182 }
2183 }
2184 } else {
2185 debug!("QueryEngine::apply_select_items - expanding *");
2187 for col_name in source_table.column_names() {
2188 expanded_items.push(SelectItem::Column {
2189 column: ColumnRef::unquoted(col_name.to_string()),
2190 leading_comments: vec![],
2191 trailing_comment: None,
2192 });
2193 }
2194 }
2195 }
2196 _ => expanded_items.push(item.clone()),
2197 }
2198 }
2199
2200 let mut column_name_counts: std::collections::HashMap<String, usize> =
2202 std::collections::HashMap::new();
2203
2204 for item in &expanded_items {
2205 let base_name = match item {
2206 SelectItem::Column {
2207 column: col_ref, ..
2208 } => col_ref.name.clone(),
2209 SelectItem::Expression { alias, .. } => alias.clone(),
2210 SelectItem::Star { .. } => unreachable!("Star should have been expanded"),
2211 SelectItem::StarExclude { .. } => {
2212 unreachable!("StarExclude should have been expanded")
2213 }
2214 };
2215
2216 let count = column_name_counts.entry(base_name.clone()).or_insert(0);
2218 let column_name = if *count == 0 {
2219 base_name.clone()
2221 } else {
2222 format!("{base_name}_{count}")
2224 };
2225 *count += 1;
2226
2227 computed_table.add_column(DataColumn::new(&column_name));
2228 }
2229
2230 let can_use_batch = expanded_items.iter().all(|item| {
2234 match item {
2235 SelectItem::Expression { expr, .. } => {
2236 matches!(expr, SqlExpression::WindowFunction { .. })
2239 || !Self::contains_window_function(expr)
2240 }
2241 _ => true, }
2243 });
2244
2245 let use_batch_evaluation = can_use_batch
2248 && std::env::var("SQL_CLI_BATCH_WINDOW")
2249 .map(|v| v != "0" && v.to_lowercase() != "false")
2250 .unwrap_or(true);
2251
2252 let batch_window_specs = if use_batch_evaluation && has_window_functions {
2254 debug!("BATCH window function evaluation flag is enabled");
2255 let specs = Self::extract_window_specs(&expanded_items);
2257 debug!(
2258 "Extracted {} window function specs for batch evaluation",
2259 specs.len()
2260 );
2261 Some(specs)
2262 } else {
2263 None
2264 };
2265
2266 let mut evaluator =
2268 ArithmeticEvaluator::with_date_notation(source_table, self.date_notation.clone());
2269
2270 if let Some(exec_ctx) = exec_context {
2272 let aliases = exec_ctx.get_aliases();
2273 if !aliases.is_empty() {
2274 debug!(
2275 "Applying {} aliases to evaluator: {:?}",
2276 aliases.len(),
2277 aliases
2278 );
2279 evaluator = evaluator.with_table_aliases(aliases);
2280 }
2281 }
2282
2283 if has_window_functions {
2286 let preload_start = Instant::now();
2287
2288 let mut window_specs = Vec::new();
2290 for item in &expanded_items {
2291 if let SelectItem::Expression { expr, .. } = item {
2292 Self::collect_window_specs(expr, &mut window_specs);
2293 }
2294 }
2295
2296 for spec in &window_specs {
2298 let _ = evaluator.get_or_create_window_context(spec);
2299 }
2300
2301 debug!(
2302 "Pre-created {} WindowContext(s) in {:.2}ms",
2303 window_specs.len(),
2304 preload_start.elapsed().as_secs_f64() * 1000.0
2305 );
2306 }
2307
2308 if let Some(window_specs) = batch_window_specs {
2310 debug!("Starting batch window function evaluation");
2311 let batch_start = Instant::now();
2312
2313 let mut batch_results: Vec<Vec<DataValue>> =
2315 vec![vec![DataValue::Null; expanded_items.len()]; visible_rows.len()];
2316
2317 let detailed_window_specs = &window_specs;
2319
2320 let mut specs_by_window: HashMap<
2322 u64,
2323 Vec<&crate::data::batch_window_evaluator::WindowFunctionSpec>,
2324 > = HashMap::new();
2325 for spec in detailed_window_specs {
2326 let hash = spec.spec.compute_hash();
2327 specs_by_window
2328 .entry(hash)
2329 .or_insert_with(Vec::new)
2330 .push(spec);
2331 }
2332
2333 for (_window_hash, specs) in specs_by_window {
2335 let context = evaluator.get_or_create_window_context(&specs[0].spec)?;
2337
2338 for spec in specs {
2340 match spec.function_name.as_str() {
2341 "LAG" => {
2342 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2344 let column_name = col_ref.name.as_str();
2345 let offset = if let Some(SqlExpression::NumberLiteral(n)) =
2346 spec.args.get(1)
2347 {
2348 n.parse::<i64>().unwrap_or(1)
2349 } else {
2350 1 };
2352
2353 let values = context.evaluate_lag_batch(
2354 visible_rows,
2355 column_name,
2356 offset,
2357 )?;
2358
2359 for (row_idx, value) in values.into_iter().enumerate() {
2361 batch_results[row_idx][spec.output_column_index] = value;
2362 }
2363 }
2364 }
2365 "LEAD" => {
2366 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2368 let column_name = col_ref.name.as_str();
2369 let offset = if let Some(SqlExpression::NumberLiteral(n)) =
2370 spec.args.get(1)
2371 {
2372 n.parse::<i64>().unwrap_or(1)
2373 } else {
2374 1 };
2376
2377 let values = context.evaluate_lead_batch(
2378 visible_rows,
2379 column_name,
2380 offset,
2381 )?;
2382
2383 for (row_idx, value) in values.into_iter().enumerate() {
2385 batch_results[row_idx][spec.output_column_index] = value;
2386 }
2387 }
2388 }
2389 "ROW_NUMBER" => {
2390 let values = context.evaluate_row_number_batch(visible_rows)?;
2391
2392 for (row_idx, value) in values.into_iter().enumerate() {
2394 batch_results[row_idx][spec.output_column_index] = value;
2395 }
2396 }
2397 "RANK" => {
2398 let values = context.evaluate_rank_batch(visible_rows)?;
2399
2400 for (row_idx, value) in values.into_iter().enumerate() {
2402 batch_results[row_idx][spec.output_column_index] = value;
2403 }
2404 }
2405 "DENSE_RANK" => {
2406 let values = context.evaluate_dense_rank_batch(visible_rows)?;
2407
2408 for (row_idx, value) in values.into_iter().enumerate() {
2410 batch_results[row_idx][spec.output_column_index] = value;
2411 }
2412 }
2413 "SUM" => {
2414 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2415 let column_name = col_ref.name.as_str();
2416 let values =
2417 context.evaluate_sum_batch(visible_rows, column_name)?;
2418
2419 for (row_idx, value) in values.into_iter().enumerate() {
2420 batch_results[row_idx][spec.output_column_index] = value;
2421 }
2422 }
2423 }
2424 "AVG" => {
2425 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2426 let column_name = col_ref.name.as_str();
2427 let values =
2428 context.evaluate_avg_batch(visible_rows, column_name)?;
2429
2430 for (row_idx, value) in values.into_iter().enumerate() {
2431 batch_results[row_idx][spec.output_column_index] = value;
2432 }
2433 }
2434 }
2435 "MIN" => {
2436 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2437 let column_name = col_ref.name.as_str();
2438 let values =
2439 context.evaluate_min_batch(visible_rows, column_name)?;
2440
2441 for (row_idx, value) in values.into_iter().enumerate() {
2442 batch_results[row_idx][spec.output_column_index] = value;
2443 }
2444 }
2445 }
2446 "MAX" => {
2447 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2448 let column_name = col_ref.name.as_str();
2449 let values =
2450 context.evaluate_max_batch(visible_rows, column_name)?;
2451
2452 for (row_idx, value) in values.into_iter().enumerate() {
2453 batch_results[row_idx][spec.output_column_index] = value;
2454 }
2455 }
2456 }
2457 "COUNT" => {
2458 let column_name = match spec.args.get(0) {
2460 Some(SqlExpression::Column(col_ref)) => Some(col_ref.name.as_str()),
2461 Some(SqlExpression::StringLiteral(s)) if s == "*" => None,
2462 _ => None,
2463 };
2464
2465 let values = context.evaluate_count_batch(visible_rows, column_name)?;
2466
2467 for (row_idx, value) in values.into_iter().enumerate() {
2468 batch_results[row_idx][spec.output_column_index] = value;
2469 }
2470 }
2471 "FIRST_VALUE" => {
2472 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2473 let column_name = col_ref.name.as_str();
2474 let values = context
2475 .evaluate_first_value_batch(visible_rows, column_name)?;
2476
2477 for (row_idx, value) in values.into_iter().enumerate() {
2478 batch_results[row_idx][spec.output_column_index] = value;
2479 }
2480 }
2481 }
2482 "LAST_VALUE" => {
2483 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2484 let column_name = col_ref.name.as_str();
2485 let values =
2486 context.evaluate_last_value_batch(visible_rows, column_name)?;
2487
2488 for (row_idx, value) in values.into_iter().enumerate() {
2489 batch_results[row_idx][spec.output_column_index] = value;
2490 }
2491 }
2492 }
2493 _ => {
2494 debug!(
2496 "Window function {} not supported in batch mode, using per-row",
2497 spec.function_name
2498 );
2499 }
2500 }
2501 }
2502 }
2503
2504 for (result_row_idx, &source_row_idx) in visible_rows.iter().enumerate() {
2506 for (col_idx, item) in expanded_items.iter().enumerate() {
2507 if !matches!(batch_results[result_row_idx][col_idx], DataValue::Null) {
2509 continue;
2510 }
2511
2512 let value = match item {
2513 SelectItem::Column {
2514 column: col_ref, ..
2515 } => {
2516 match evaluator
2517 .evaluate(&SqlExpression::Column(col_ref.clone()), source_row_idx)
2518 {
2519 Ok(val) => val,
2520 Err(e) => {
2521 return Err(anyhow!(
2522 "Failed to evaluate column {}: {}",
2523 col_ref.to_sql(),
2524 e
2525 ));
2526 }
2527 }
2528 }
2529 SelectItem::Expression { expr, .. } => {
2530 if matches!(expr, SqlExpression::WindowFunction { .. }) {
2533 continue;
2535 }
2536 evaluator.evaluate(&expr, source_row_idx)?
2539 }
2540 SelectItem::Star { .. } => unreachable!("Star should have been expanded"),
2541 SelectItem::StarExclude { .. } => {
2542 unreachable!("StarExclude should have been expanded")
2543 }
2544 };
2545 batch_results[result_row_idx][col_idx] = value;
2546 }
2547 }
2548
2549 for row_values in batch_results {
2551 computed_table
2552 .add_row(DataRow::new(row_values))
2553 .map_err(|e| anyhow::anyhow!("Failed to add row: {}", e))?;
2554 }
2555
2556 debug!(
2557 "Batch window evaluation completed in {:.3}ms",
2558 batch_start.elapsed().as_secs_f64() * 1000.0
2559 );
2560 } else {
2561 for &row_idx in visible_rows {
2563 let mut row_values = Vec::new();
2564
2565 for item in &expanded_items {
2566 let value = match item {
2567 SelectItem::Column {
2568 column: col_ref, ..
2569 } => {
2570 match evaluator
2572 .evaluate(&SqlExpression::Column(col_ref.clone()), row_idx)
2573 {
2574 Ok(val) => val,
2575 Err(e) => {
2576 return Err(anyhow!(
2577 "Failed to evaluate column {}: {}",
2578 col_ref.to_sql(),
2579 e
2580 ));
2581 }
2582 }
2583 }
2584 SelectItem::Expression { expr, .. } => {
2585 evaluator.evaluate(&expr, row_idx)?
2587 }
2588 SelectItem::Star { .. } => unreachable!("Star should have been expanded"),
2589 SelectItem::StarExclude { .. } => {
2590 unreachable!("StarExclude should have been expanded")
2591 }
2592 };
2593 row_values.push(value);
2594 }
2595
2596 computed_table
2597 .add_row(DataRow::new(row_values))
2598 .map_err(|e| anyhow::anyhow!("Failed to add row: {}", e))?;
2599 }
2600 }
2601
2602 if let Some(start) = window_start {
2604 let window_duration = start.elapsed();
2605 info!(
2606 "Window function evaluation took {:.2}ms for {} rows ({} window functions)",
2607 window_duration.as_secs_f64() * 1000.0,
2608 visible_rows.len(),
2609 window_func_count
2610 );
2611
2612 plan.begin_step(
2614 StepType::WindowFunction,
2615 format!("Evaluate {} window function(s)", window_func_count),
2616 );
2617 plan.set_rows_in(visible_rows.len());
2618 plan.set_rows_out(visible_rows.len());
2619 plan.add_detail(format!("Input: {} rows", visible_rows.len()));
2620 plan.add_detail(format!("{} window functions evaluated", window_func_count));
2621 plan.add_detail(format!(
2622 "Evaluation time: {:.3}ms",
2623 window_duration.as_secs_f64() * 1000.0
2624 ));
2625 plan.end_step();
2626 }
2627
2628 Ok(DataView::new(Arc::new(computed_table)))
2631 }
2632
2633 fn apply_select_with_row_expansion(
2635 &self,
2636 view: DataView,
2637 select_items: &[SelectItem],
2638 ) -> Result<DataView> {
2639 debug!("QueryEngine::apply_select_with_row_expansion - expanding rows");
2640
2641 let source_table = view.source();
2642 let visible_rows = view.visible_row_indices();
2643 let expander_registry = RowExpanderRegistry::new();
2644
2645 let mut result_table = DataTable::new("unnest_result");
2647
2648 let mut expanded_items = Vec::new();
2650 for item in select_items {
2651 match item {
2652 SelectItem::Star { table_prefix, .. } => {
2653 if let Some(prefix) = table_prefix {
2654 debug!(
2656 "QueryEngine::apply_select_with_row_expansion - expanding {}.*",
2657 prefix
2658 );
2659 for col in &source_table.columns {
2660 if Self::column_matches_table(col, prefix) {
2661 expanded_items.push(SelectItem::Column {
2662 column: ColumnRef::unquoted(col.name.clone()),
2663 leading_comments: vec![],
2664 trailing_comment: None,
2665 });
2666 }
2667 }
2668 } else {
2669 debug!("QueryEngine::apply_select_with_row_expansion - expanding *");
2671 for col_name in source_table.column_names() {
2672 expanded_items.push(SelectItem::Column {
2673 column: ColumnRef::unquoted(col_name.to_string()),
2674 leading_comments: vec![],
2675 trailing_comment: None,
2676 });
2677 }
2678 }
2679 }
2680 _ => expanded_items.push(item.clone()),
2681 }
2682 }
2683
2684 for item in &expanded_items {
2686 let column_name = match item {
2687 SelectItem::Column {
2688 column: col_ref, ..
2689 } => col_ref.name.clone(),
2690 SelectItem::Expression { alias, .. } => alias.clone(),
2691 SelectItem::Star { .. } => unreachable!("Star should have been expanded"),
2692 SelectItem::StarExclude { .. } => {
2693 unreachable!("StarExclude should have been expanded")
2694 }
2695 };
2696 result_table.add_column(DataColumn::new(&column_name));
2697 }
2698
2699 let mut evaluator =
2701 ArithmeticEvaluator::with_date_notation(source_table, self.date_notation.clone());
2702
2703 for &row_idx in visible_rows {
2704 let mut unnest_expansions = Vec::new();
2706 let mut unnest_indices = Vec::new();
2707
2708 for (col_idx, item) in expanded_items.iter().enumerate() {
2709 if let SelectItem::Expression { expr, .. } = item {
2710 if let Some(expansion_result) = self.try_expand_unnest(
2711 &expr,
2712 source_table,
2713 row_idx,
2714 &mut evaluator,
2715 &expander_registry,
2716 )? {
2717 unnest_expansions.push(expansion_result);
2718 unnest_indices.push(col_idx);
2719 }
2720 }
2721 }
2722
2723 let expansion_count = if unnest_expansions.is_empty() {
2725 1 } else {
2727 unnest_expansions
2728 .iter()
2729 .map(|exp| exp.row_count())
2730 .max()
2731 .unwrap_or(1)
2732 };
2733
2734 for output_idx in 0..expansion_count {
2736 let mut row_values = Vec::new();
2737
2738 for (col_idx, item) in expanded_items.iter().enumerate() {
2739 let unnest_position = unnest_indices.iter().position(|&idx| idx == col_idx);
2741
2742 let value = if let Some(unnest_idx) = unnest_position {
2743 let expansion = &unnest_expansions[unnest_idx];
2745 expansion
2746 .values
2747 .get(output_idx)
2748 .cloned()
2749 .unwrap_or(DataValue::Null)
2750 } else {
2751 match item {
2753 SelectItem::Column {
2754 column: col_ref, ..
2755 } => {
2756 let col_idx =
2757 source_table.get_column_index(&col_ref.name).ok_or_else(
2758 || anyhow::anyhow!("Column '{}' not found", col_ref.name),
2759 )?;
2760 let row = source_table
2761 .get_row(row_idx)
2762 .ok_or_else(|| anyhow::anyhow!("Row {} not found", row_idx))?;
2763 row.get(col_idx)
2764 .ok_or_else(|| {
2765 anyhow::anyhow!("Column {} not found in row", col_idx)
2766 })?
2767 .clone()
2768 }
2769 SelectItem::Expression { expr, .. } => {
2770 evaluator.evaluate(&expr, row_idx)?
2772 }
2773 SelectItem::Star { .. } => unreachable!(),
2774 SelectItem::StarExclude { .. } => {
2775 unreachable!("StarExclude should have been expanded")
2776 }
2777 }
2778 };
2779
2780 row_values.push(value);
2781 }
2782
2783 result_table
2784 .add_row(DataRow::new(row_values))
2785 .map_err(|e| anyhow::anyhow!("Failed to add expanded row: {}", e))?;
2786 }
2787 }
2788
2789 debug!(
2790 "QueryEngine::apply_select_with_row_expansion - input rows: {}, output rows: {}",
2791 visible_rows.len(),
2792 result_table.row_count()
2793 );
2794
2795 Ok(DataView::new(Arc::new(result_table)))
2796 }
2797
2798 fn try_expand_unnest(
2801 &self,
2802 expr: &SqlExpression,
2803 _source_table: &DataTable,
2804 row_idx: usize,
2805 evaluator: &mut ArithmeticEvaluator,
2806 expander_registry: &RowExpanderRegistry,
2807 ) -> Result<Option<crate::data::row_expanders::ExpansionResult>> {
2808 if let SqlExpression::Unnest { column, delimiter } = expr {
2810 let column_value = evaluator.evaluate(column, row_idx)?;
2812
2813 let delimiter_value = DataValue::String(delimiter.clone());
2815
2816 let expander = expander_registry
2818 .get("UNNEST")
2819 .ok_or_else(|| anyhow::anyhow!("UNNEST expander not found"))?;
2820
2821 let expansion = expander.expand(&column_value, &[delimiter_value])?;
2823 return Ok(Some(expansion));
2824 }
2825
2826 if let SqlExpression::FunctionCall { name, args, .. } = expr {
2828 if name.to_uppercase() == "UNNEST" {
2829 if args.len() != 2 {
2831 return Err(anyhow::anyhow!(
2832 "UNNEST requires exactly 2 arguments: UNNEST(column, delimiter)"
2833 ));
2834 }
2835
2836 let column_value = evaluator.evaluate(&args[0], row_idx)?;
2838
2839 let delimiter_value = evaluator.evaluate(&args[1], row_idx)?;
2841
2842 let expander = expander_registry
2844 .get("UNNEST")
2845 .ok_or_else(|| anyhow::anyhow!("UNNEST expander not found"))?;
2846
2847 let expansion = expander.expand(&column_value, &[delimiter_value])?;
2849 return Ok(Some(expansion));
2850 }
2851 }
2852
2853 Ok(None)
2854 }
2855
2856 fn apply_aggregate_select(
2858 &self,
2859 view: DataView,
2860 select_items: &[SelectItem],
2861 ) -> Result<DataView> {
2862 debug!("QueryEngine::apply_aggregate_select - creating single row aggregate result");
2863
2864 let source_table = view.source();
2865 let mut result_table = DataTable::new("aggregate_result");
2866
2867 for item in select_items {
2869 let column_name = match item {
2870 SelectItem::Expression { alias, .. } => alias.clone(),
2871 _ => unreachable!("Should only have expressions in aggregate-only query"),
2872 };
2873 result_table.add_column(DataColumn::new(&column_name));
2874 }
2875
2876 let visible_rows = view.visible_row_indices().to_vec();
2878 let mut evaluator =
2879 ArithmeticEvaluator::with_date_notation(source_table, self.date_notation.clone())
2880 .with_visible_rows(visible_rows);
2881
2882 let mut row_values = Vec::new();
2884 for item in select_items {
2885 match item {
2886 SelectItem::Expression { expr, .. } => {
2887 let value = evaluator.evaluate(expr, 0)?;
2890 row_values.push(value);
2891 }
2892 _ => unreachable!("Should only have expressions in aggregate-only query"),
2893 }
2894 }
2895
2896 result_table
2898 .add_row(DataRow::new(row_values))
2899 .map_err(|e| anyhow::anyhow!("Failed to add aggregate result row: {}", e))?;
2900
2901 Ok(DataView::new(Arc::new(result_table)))
2902 }
2903
2904 fn column_matches_table(col: &DataColumn, table_name: &str) -> bool {
2916 if let Some(ref source) = col.source_table {
2918 if source == table_name || source.ends_with(&format!(".{}", table_name)) {
2920 return true;
2921 }
2922 }
2923
2924 if let Some(ref qualified) = col.qualified_name {
2926 if qualified.starts_with(&format!("{}.", table_name)) {
2928 return true;
2929 }
2930 }
2931
2932 false
2933 }
2934
2935 fn resolve_select_columns(
2937 &self,
2938 table: &DataTable,
2939 select_items: &[SelectItem],
2940 ) -> Result<Vec<usize>> {
2941 let mut indices = Vec::new();
2942 let table_columns = table.column_names();
2943
2944 for item in select_items {
2945 match item {
2946 SelectItem::Column {
2947 column: col_ref, ..
2948 } => {
2949 let index = if let Some(table_prefix) = &col_ref.table_prefix {
2951 let qualified_name = format!("{}.{}", table_prefix, col_ref.name);
2960 table.find_column_by_qualified_name(&qualified_name)
2961 .or_else(|| {
2962 table_columns
2963 .iter()
2964 .position(|c| c.eq_ignore_ascii_case(&col_ref.name))
2965 })
2966 .ok_or_else(|| {
2967 let has_qualified = table.columns.iter()
2969 .any(|c| c.qualified_name.is_some());
2970 if !has_qualified {
2971 anyhow::anyhow!(
2972 "Column '{}' not found. Note: Table '{}' may not support qualified column names",
2973 qualified_name, table_prefix
2974 )
2975 } else {
2976 anyhow::anyhow!("Column '{}' not found", qualified_name)
2977 }
2978 })?
2979 } else {
2980 table_columns
2982 .iter()
2983 .position(|c| c.eq_ignore_ascii_case(&col_ref.name))
2984 .ok_or_else(|| {
2985 let suggestion = self.find_similar_column(table, &col_ref.name);
2986 match suggestion {
2987 Some(similar) => anyhow::anyhow!(
2988 "Column '{}' not found. Did you mean '{}'?",
2989 col_ref.name,
2990 similar
2991 ),
2992 None => anyhow::anyhow!("Column '{}' not found", col_ref.name),
2993 }
2994 })?
2995 };
2996 indices.push(index);
2997 }
2998 SelectItem::Star { table_prefix, .. } => {
2999 if let Some(prefix) = table_prefix {
3000 for (i, col) in table.columns.iter().enumerate() {
3002 if Self::column_matches_table(col, prefix) {
3003 indices.push(i);
3004 }
3005 }
3006 } else {
3007 for i in 0..table_columns.len() {
3009 indices.push(i);
3010 }
3011 }
3012 }
3013 SelectItem::StarExclude {
3014 table_prefix,
3015 excluded_columns,
3016 ..
3017 } => {
3018 if let Some(prefix) = table_prefix {
3020 for (i, col) in table.columns.iter().enumerate() {
3022 if Self::column_matches_table(col, prefix)
3023 && !excluded_columns.contains(&col.name)
3024 {
3025 indices.push(i);
3026 }
3027 }
3028 } else {
3029 for (i, col_name) in table_columns.iter().enumerate() {
3031 if !excluded_columns
3032 .iter()
3033 .any(|exc| exc.eq_ignore_ascii_case(col_name))
3034 {
3035 indices.push(i);
3036 }
3037 }
3038 }
3039 }
3040 SelectItem::Expression { .. } => {
3041 return Err(anyhow::anyhow!(
3042 "Computed expressions require new table creation"
3043 ));
3044 }
3045 }
3046 }
3047
3048 Ok(indices)
3049 }
3050
3051 fn apply_distinct(&self, view: DataView) -> Result<DataView> {
3053 use std::collections::HashSet;
3054
3055 let source = view.source();
3056 let visible_cols = view.visible_column_indices();
3057 let visible_rows = view.visible_row_indices();
3058
3059 let mut seen_rows = HashSet::new();
3061 let mut unique_row_indices = Vec::new();
3062
3063 for &row_idx in visible_rows {
3064 let mut row_key = Vec::new();
3066 for &col_idx in visible_cols {
3067 let value = source
3068 .get_value(row_idx, col_idx)
3069 .ok_or_else(|| anyhow!("Invalid cell reference"))?;
3070 row_key.push(format!("{:?}", value));
3072 }
3073
3074 if seen_rows.insert(row_key) {
3076 unique_row_indices.push(row_idx);
3078 }
3079 }
3080
3081 Ok(view.with_rows(unique_row_indices))
3083 }
3084
3085 fn apply_multi_order_by(
3087 &self,
3088 view: DataView,
3089 order_by_columns: &[OrderByItem],
3090 ) -> Result<DataView> {
3091 self.apply_multi_order_by_with_context(view, order_by_columns, None)
3092 }
3093
3094 fn apply_multi_order_by_with_context(
3096 &self,
3097 mut view: DataView,
3098 order_by_columns: &[OrderByItem],
3099 _exec_context: Option<&ExecutionContext>,
3100 ) -> Result<DataView> {
3101 let mut sort_columns = Vec::new();
3103
3104 for order_col in order_by_columns {
3105 let column_name = match &order_col.expr {
3107 SqlExpression::Column(col_ref) => col_ref.name.clone(),
3108 _ => {
3109 return Err(anyhow!(
3111 "ORDER BY expressions not yet supported - only simple columns allowed"
3112 ));
3113 }
3114 };
3115
3116 let col_index = if column_name.contains('.') {
3118 if let Some(dot_pos) = column_name.rfind('.') {
3120 let col_name = &column_name[dot_pos + 1..];
3121
3122 debug!(
3125 "ORDER BY: Extracting unqualified column '{}' from '{}'",
3126 col_name, column_name
3127 );
3128 view.source().get_column_index(col_name)
3129 } else {
3130 view.source().get_column_index(&column_name)
3131 }
3132 } else {
3133 view.source().get_column_index(&column_name)
3135 }
3136 .ok_or_else(|| {
3137 let suggestion = self.find_similar_column(view.source(), &column_name);
3139 match suggestion {
3140 Some(similar) => anyhow::anyhow!(
3141 "Column '{}' not found. Did you mean '{}'?",
3142 column_name,
3143 similar
3144 ),
3145 None => {
3146 let available_cols = view.source().column_names().join(", ");
3148 anyhow::anyhow!(
3149 "Column '{}' not found. Available columns: {}",
3150 column_name,
3151 available_cols
3152 )
3153 }
3154 }
3155 })?;
3156
3157 let ascending = matches!(order_col.direction, SortDirection::Asc);
3158 sort_columns.push((col_index, ascending));
3159 }
3160
3161 view.apply_multi_sort(&sort_columns)?;
3163 Ok(view)
3164 }
3165
3166 fn apply_group_by(
3168 &self,
3169 view: DataView,
3170 group_by_exprs: &[SqlExpression],
3171 select_items: &[SelectItem],
3172 having: Option<&SqlExpression>,
3173 plan: &mut ExecutionPlanBuilder,
3174 ) -> Result<DataView> {
3175 let (result_view, phase_info) = self.apply_group_by_expressions(
3177 view,
3178 group_by_exprs,
3179 select_items,
3180 having,
3181 self.case_insensitive,
3182 self.date_notation.clone(),
3183 )?;
3184
3185 plan.add_detail(format!("=== GROUP BY Phase Breakdown ==="));
3187 plan.add_detail(format!(
3188 "Phase 1 - Group Building: {:.3}ms",
3189 phase_info.phase2_key_building.as_secs_f64() * 1000.0
3190 ));
3191 plan.add_detail(format!(
3192 " • Processing {} rows into {} groups",
3193 phase_info.total_rows, phase_info.num_groups
3194 ));
3195 plan.add_detail(format!(
3196 "Phase 2 - Aggregation: {:.3}ms",
3197 phase_info.phase4_aggregation.as_secs_f64() * 1000.0
3198 ));
3199 if phase_info.phase4_having_evaluation > Duration::ZERO {
3200 plan.add_detail(format!(
3201 "Phase 3 - HAVING Filter: {:.3}ms",
3202 phase_info.phase4_having_evaluation.as_secs_f64() * 1000.0
3203 ));
3204 plan.add_detail(format!(
3205 " • Filtered {} groups",
3206 phase_info.groups_filtered_by_having
3207 ));
3208 }
3209 plan.add_detail(format!(
3210 "Total GROUP BY time: {:.3}ms",
3211 phase_info.total_time.as_secs_f64() * 1000.0
3212 ));
3213
3214 Ok(result_view)
3215 }
3216
3217 pub fn estimate_group_cardinality(
3220 &self,
3221 view: &DataView,
3222 group_by_exprs: &[SqlExpression],
3223 ) -> usize {
3224 let row_count = view.get_visible_rows().len();
3226 if row_count <= 100 {
3227 return row_count;
3228 }
3229
3230 let sample_size = min(1000, row_count / 10).max(100);
3232 let mut seen = FxHashSet::default();
3233
3234 let visible_rows = view.get_visible_rows();
3235 for (i, &row_idx) in visible_rows.iter().enumerate() {
3236 if i >= sample_size {
3237 break;
3238 }
3239
3240 let mut key_values = Vec::new();
3242 for expr in group_by_exprs {
3243 let mut evaluator = ArithmeticEvaluator::new(view.source());
3244 let value = evaluator.evaluate(expr, row_idx).unwrap_or(DataValue::Null);
3245 key_values.push(value);
3246 }
3247
3248 seen.insert(key_values);
3249 }
3250
3251 let sample_cardinality = seen.len();
3253 let estimated = (sample_cardinality * row_count) / sample_size;
3254
3255 estimated.min(row_count).max(sample_cardinality)
3257 }
3258}
3259
3260#[cfg(test)]
3261mod tests {
3262 use super::*;
3263 use crate::data::datatable::{DataColumn, DataRow, DataValue};
3264
3265 fn create_test_table() -> Arc<DataTable> {
3266 let mut table = DataTable::new("test");
3267
3268 table.add_column(DataColumn::new("id"));
3270 table.add_column(DataColumn::new("name"));
3271 table.add_column(DataColumn::new("age"));
3272
3273 table
3275 .add_row(DataRow::new(vec![
3276 DataValue::Integer(1),
3277 DataValue::String("Alice".to_string()),
3278 DataValue::Integer(30),
3279 ]))
3280 .unwrap();
3281
3282 table
3283 .add_row(DataRow::new(vec![
3284 DataValue::Integer(2),
3285 DataValue::String("Bob".to_string()),
3286 DataValue::Integer(25),
3287 ]))
3288 .unwrap();
3289
3290 table
3291 .add_row(DataRow::new(vec![
3292 DataValue::Integer(3),
3293 DataValue::String("Charlie".to_string()),
3294 DataValue::Integer(35),
3295 ]))
3296 .unwrap();
3297
3298 Arc::new(table)
3299 }
3300
3301 #[test]
3302 fn test_select_all() {
3303 let table = create_test_table();
3304 let engine = QueryEngine::new();
3305
3306 let view = engine
3307 .execute(table.clone(), "SELECT * FROM users")
3308 .unwrap();
3309 assert_eq!(view.row_count(), 3);
3310 assert_eq!(view.column_count(), 3);
3311 }
3312
3313 #[test]
3314 fn test_select_columns() {
3315 let table = create_test_table();
3316 let engine = QueryEngine::new();
3317
3318 let view = engine
3319 .execute(table.clone(), "SELECT name, age FROM users")
3320 .unwrap();
3321 assert_eq!(view.row_count(), 3);
3322 assert_eq!(view.column_count(), 2);
3323 }
3324
3325 #[test]
3326 fn test_select_with_limit() {
3327 let table = create_test_table();
3328 let engine = QueryEngine::new();
3329
3330 let view = engine
3331 .execute(table.clone(), "SELECT * FROM users LIMIT 2")
3332 .unwrap();
3333 assert_eq!(view.row_count(), 2);
3334 }
3335
3336 #[test]
3337 fn test_type_coercion_contains() {
3338 let _ = tracing_subscriber::fmt()
3340 .with_max_level(tracing::Level::DEBUG)
3341 .try_init();
3342
3343 let mut table = DataTable::new("test");
3344 table.add_column(DataColumn::new("id"));
3345 table.add_column(DataColumn::new("status"));
3346 table.add_column(DataColumn::new("price"));
3347
3348 table
3350 .add_row(DataRow::new(vec![
3351 DataValue::Integer(1),
3352 DataValue::String("Pending".to_string()),
3353 DataValue::Float(99.99),
3354 ]))
3355 .unwrap();
3356
3357 table
3358 .add_row(DataRow::new(vec![
3359 DataValue::Integer(2),
3360 DataValue::String("Confirmed".to_string()),
3361 DataValue::Float(150.50),
3362 ]))
3363 .unwrap();
3364
3365 table
3366 .add_row(DataRow::new(vec![
3367 DataValue::Integer(3),
3368 DataValue::String("Pending".to_string()),
3369 DataValue::Float(75.00),
3370 ]))
3371 .unwrap();
3372
3373 let table = Arc::new(table);
3374 let engine = QueryEngine::new();
3375
3376 println!("\n=== Testing WHERE clause with Contains ===");
3377 println!("Table has {} rows", table.row_count());
3378 for i in 0..table.row_count() {
3379 let status = table.get_value(i, 1);
3380 println!("Row {i}: status = {status:?}");
3381 }
3382
3383 println!("\n--- Test 1: status.Contains('pend') ---");
3385 let result = engine.execute(
3386 table.clone(),
3387 "SELECT * FROM test WHERE status.Contains('pend')",
3388 );
3389 match result {
3390 Ok(view) => {
3391 println!("SUCCESS: Found {} matching rows", view.row_count());
3392 assert_eq!(view.row_count(), 2); }
3394 Err(e) => {
3395 panic!("Query failed: {e}");
3396 }
3397 }
3398
3399 println!("\n--- Test 2: price.Contains('9') ---");
3401 let result = engine.execute(
3402 table.clone(),
3403 "SELECT * FROM test WHERE price.Contains('9')",
3404 );
3405 match result {
3406 Ok(view) => {
3407 println!(
3408 "SUCCESS: Found {} matching rows with price containing '9'",
3409 view.row_count()
3410 );
3411 assert!(view.row_count() >= 1);
3413 }
3414 Err(e) => {
3415 panic!("Numeric coercion query failed: {e}");
3416 }
3417 }
3418
3419 println!("\n=== All tests passed! ===");
3420 }
3421
3422 #[test]
3423 fn test_not_in_clause() {
3424 let _ = tracing_subscriber::fmt()
3426 .with_max_level(tracing::Level::DEBUG)
3427 .try_init();
3428
3429 let mut table = DataTable::new("test");
3430 table.add_column(DataColumn::new("id"));
3431 table.add_column(DataColumn::new("country"));
3432
3433 table
3435 .add_row(DataRow::new(vec![
3436 DataValue::Integer(1),
3437 DataValue::String("CA".to_string()),
3438 ]))
3439 .unwrap();
3440
3441 table
3442 .add_row(DataRow::new(vec![
3443 DataValue::Integer(2),
3444 DataValue::String("US".to_string()),
3445 ]))
3446 .unwrap();
3447
3448 table
3449 .add_row(DataRow::new(vec![
3450 DataValue::Integer(3),
3451 DataValue::String("UK".to_string()),
3452 ]))
3453 .unwrap();
3454
3455 let table = Arc::new(table);
3456 let engine = QueryEngine::new();
3457
3458 println!("\n=== Testing NOT IN clause ===");
3459 println!("Table has {} rows", table.row_count());
3460 for i in 0..table.row_count() {
3461 let country = table.get_value(i, 1);
3462 println!("Row {i}: country = {country:?}");
3463 }
3464
3465 println!("\n--- Test: country NOT IN ('CA') ---");
3467 let result = engine.execute(
3468 table.clone(),
3469 "SELECT * FROM test WHERE country NOT IN ('CA')",
3470 );
3471 match result {
3472 Ok(view) => {
3473 println!("SUCCESS: Found {} rows not in ('CA')", view.row_count());
3474 assert_eq!(view.row_count(), 2); }
3476 Err(e) => {
3477 panic!("NOT IN query failed: {e}");
3478 }
3479 }
3480
3481 println!("\n=== NOT IN test complete! ===");
3482 }
3483
3484 #[test]
3485 fn test_case_insensitive_in_and_not_in() {
3486 let _ = tracing_subscriber::fmt()
3488 .with_max_level(tracing::Level::DEBUG)
3489 .try_init();
3490
3491 let mut table = DataTable::new("test");
3492 table.add_column(DataColumn::new("id"));
3493 table.add_column(DataColumn::new("country"));
3494
3495 table
3497 .add_row(DataRow::new(vec![
3498 DataValue::Integer(1),
3499 DataValue::String("CA".to_string()), ]))
3501 .unwrap();
3502
3503 table
3504 .add_row(DataRow::new(vec![
3505 DataValue::Integer(2),
3506 DataValue::String("us".to_string()), ]))
3508 .unwrap();
3509
3510 table
3511 .add_row(DataRow::new(vec![
3512 DataValue::Integer(3),
3513 DataValue::String("UK".to_string()), ]))
3515 .unwrap();
3516
3517 let table = Arc::new(table);
3518
3519 println!("\n=== Testing Case-Insensitive IN clause ===");
3520 println!("Table has {} rows", table.row_count());
3521 for i in 0..table.row_count() {
3522 let country = table.get_value(i, 1);
3523 println!("Row {i}: country = {country:?}");
3524 }
3525
3526 println!("\n--- Test: country IN ('ca') with case_insensitive=true ---");
3528 let engine = QueryEngine::with_case_insensitive(true);
3529 let result = engine.execute(table.clone(), "SELECT * FROM test WHERE country IN ('ca')");
3530 match result {
3531 Ok(view) => {
3532 println!(
3533 "SUCCESS: Found {} rows matching 'ca' (case-insensitive)",
3534 view.row_count()
3535 );
3536 assert_eq!(view.row_count(), 1); }
3538 Err(e) => {
3539 panic!("Case-insensitive IN query failed: {e}");
3540 }
3541 }
3542
3543 println!("\n--- Test: country NOT IN ('ca') with case_insensitive=true ---");
3545 let result = engine.execute(
3546 table.clone(),
3547 "SELECT * FROM test WHERE country NOT IN ('ca')",
3548 );
3549 match result {
3550 Ok(view) => {
3551 println!(
3552 "SUCCESS: Found {} rows not matching 'ca' (case-insensitive)",
3553 view.row_count()
3554 );
3555 assert_eq!(view.row_count(), 2); }
3557 Err(e) => {
3558 panic!("Case-insensitive NOT IN query failed: {e}");
3559 }
3560 }
3561
3562 println!("\n--- Test: country IN ('ca') with case_insensitive=false ---");
3564 let engine_case_sensitive = QueryEngine::new(); let result = engine_case_sensitive
3566 .execute(table.clone(), "SELECT * FROM test WHERE country IN ('ca')");
3567 match result {
3568 Ok(view) => {
3569 println!(
3570 "SUCCESS: Found {} rows matching 'ca' (case-sensitive)",
3571 view.row_count()
3572 );
3573 assert_eq!(view.row_count(), 0); }
3575 Err(e) => {
3576 panic!("Case-sensitive IN query failed: {e}");
3577 }
3578 }
3579
3580 println!("\n=== Case-insensitive IN/NOT IN test complete! ===");
3581 }
3582
3583 #[test]
3584 #[ignore = "Parentheses in WHERE clause not yet implemented"]
3585 fn test_parentheses_in_where_clause() {
3586 let _ = tracing_subscriber::fmt()
3588 .with_max_level(tracing::Level::DEBUG)
3589 .try_init();
3590
3591 let mut table = DataTable::new("test");
3592 table.add_column(DataColumn::new("id"));
3593 table.add_column(DataColumn::new("status"));
3594 table.add_column(DataColumn::new("priority"));
3595
3596 table
3598 .add_row(DataRow::new(vec![
3599 DataValue::Integer(1),
3600 DataValue::String("Pending".to_string()),
3601 DataValue::String("High".to_string()),
3602 ]))
3603 .unwrap();
3604
3605 table
3606 .add_row(DataRow::new(vec![
3607 DataValue::Integer(2),
3608 DataValue::String("Complete".to_string()),
3609 DataValue::String("High".to_string()),
3610 ]))
3611 .unwrap();
3612
3613 table
3614 .add_row(DataRow::new(vec![
3615 DataValue::Integer(3),
3616 DataValue::String("Pending".to_string()),
3617 DataValue::String("Low".to_string()),
3618 ]))
3619 .unwrap();
3620
3621 table
3622 .add_row(DataRow::new(vec![
3623 DataValue::Integer(4),
3624 DataValue::String("Complete".to_string()),
3625 DataValue::String("Low".to_string()),
3626 ]))
3627 .unwrap();
3628
3629 let table = Arc::new(table);
3630 let engine = QueryEngine::new();
3631
3632 println!("\n=== Testing Parentheses in WHERE clause ===");
3633 println!("Table has {} rows", table.row_count());
3634 for i in 0..table.row_count() {
3635 let status = table.get_value(i, 1);
3636 let priority = table.get_value(i, 2);
3637 println!("Row {i}: status = {status:?}, priority = {priority:?}");
3638 }
3639
3640 println!("\n--- Test: (status = 'Pending' AND priority = 'High') OR (status = 'Complete' AND priority = 'Low') ---");
3642 let result = engine.execute(
3643 table.clone(),
3644 "SELECT * FROM test WHERE (status = 'Pending' AND priority = 'High') OR (status = 'Complete' AND priority = 'Low')",
3645 );
3646 match result {
3647 Ok(view) => {
3648 println!(
3649 "SUCCESS: Found {} rows with parenthetical logic",
3650 view.row_count()
3651 );
3652 assert_eq!(view.row_count(), 2); }
3654 Err(e) => {
3655 panic!("Parentheses query failed: {e}");
3656 }
3657 }
3658
3659 println!("\n=== Parentheses test complete! ===");
3660 }
3661
3662 #[test]
3663 #[ignore = "Numeric type coercion needs fixing"]
3664 fn test_numeric_type_coercion() {
3665 let _ = tracing_subscriber::fmt()
3667 .with_max_level(tracing::Level::DEBUG)
3668 .try_init();
3669
3670 let mut table = DataTable::new("test");
3671 table.add_column(DataColumn::new("id"));
3672 table.add_column(DataColumn::new("price"));
3673 table.add_column(DataColumn::new("quantity"));
3674
3675 table
3677 .add_row(DataRow::new(vec![
3678 DataValue::Integer(1),
3679 DataValue::Float(99.50), DataValue::Integer(100),
3681 ]))
3682 .unwrap();
3683
3684 table
3685 .add_row(DataRow::new(vec![
3686 DataValue::Integer(2),
3687 DataValue::Float(150.0), DataValue::Integer(200),
3689 ]))
3690 .unwrap();
3691
3692 table
3693 .add_row(DataRow::new(vec![
3694 DataValue::Integer(3),
3695 DataValue::Integer(75), DataValue::Integer(50),
3697 ]))
3698 .unwrap();
3699
3700 let table = Arc::new(table);
3701 let engine = QueryEngine::new();
3702
3703 println!("\n=== Testing Numeric Type Coercion ===");
3704 println!("Table has {} rows", table.row_count());
3705 for i in 0..table.row_count() {
3706 let price = table.get_value(i, 1);
3707 let quantity = table.get_value(i, 2);
3708 println!("Row {i}: price = {price:?}, quantity = {quantity:?}");
3709 }
3710
3711 println!("\n--- Test: price.Contains('.') ---");
3713 let result = engine.execute(
3714 table.clone(),
3715 "SELECT * FROM test WHERE price.Contains('.')",
3716 );
3717 match result {
3718 Ok(view) => {
3719 println!(
3720 "SUCCESS: Found {} rows with decimal points in price",
3721 view.row_count()
3722 );
3723 assert_eq!(view.row_count(), 2); }
3725 Err(e) => {
3726 panic!("Numeric Contains query failed: {e}");
3727 }
3728 }
3729
3730 println!("\n--- Test: quantity.Contains('0') ---");
3732 let result = engine.execute(
3733 table.clone(),
3734 "SELECT * FROM test WHERE quantity.Contains('0')",
3735 );
3736 match result {
3737 Ok(view) => {
3738 println!(
3739 "SUCCESS: Found {} rows with '0' in quantity",
3740 view.row_count()
3741 );
3742 assert_eq!(view.row_count(), 2); }
3744 Err(e) => {
3745 panic!("Integer Contains query failed: {e}");
3746 }
3747 }
3748
3749 println!("\n=== Numeric type coercion test complete! ===");
3750 }
3751
3752 #[test]
3753 fn test_datetime_comparisons() {
3754 let _ = tracing_subscriber::fmt()
3756 .with_max_level(tracing::Level::DEBUG)
3757 .try_init();
3758
3759 let mut table = DataTable::new("test");
3760 table.add_column(DataColumn::new("id"));
3761 table.add_column(DataColumn::new("created_date"));
3762
3763 table
3765 .add_row(DataRow::new(vec![
3766 DataValue::Integer(1),
3767 DataValue::String("2024-12-15".to_string()),
3768 ]))
3769 .unwrap();
3770
3771 table
3772 .add_row(DataRow::new(vec![
3773 DataValue::Integer(2),
3774 DataValue::String("2025-01-15".to_string()),
3775 ]))
3776 .unwrap();
3777
3778 table
3779 .add_row(DataRow::new(vec![
3780 DataValue::Integer(3),
3781 DataValue::String("2025-02-15".to_string()),
3782 ]))
3783 .unwrap();
3784
3785 let table = Arc::new(table);
3786 let engine = QueryEngine::new();
3787
3788 println!("\n=== Testing DateTime Comparisons ===");
3789 println!("Table has {} rows", table.row_count());
3790 for i in 0..table.row_count() {
3791 let date = table.get_value(i, 1);
3792 println!("Row {i}: created_date = {date:?}");
3793 }
3794
3795 println!("\n--- Test: created_date > DateTime(2025,1,1) ---");
3797 let result = engine.execute(
3798 table.clone(),
3799 "SELECT * FROM test WHERE created_date > DateTime(2025,1,1)",
3800 );
3801 match result {
3802 Ok(view) => {
3803 println!("SUCCESS: Found {} rows after 2025-01-01", view.row_count());
3804 assert_eq!(view.row_count(), 2); }
3806 Err(e) => {
3807 panic!("DateTime comparison query failed: {e}");
3808 }
3809 }
3810
3811 println!("\n=== DateTime comparison test complete! ===");
3812 }
3813
3814 #[test]
3815 fn test_not_with_method_calls() {
3816 let _ = tracing_subscriber::fmt()
3818 .with_max_level(tracing::Level::DEBUG)
3819 .try_init();
3820
3821 let mut table = DataTable::new("test");
3822 table.add_column(DataColumn::new("id"));
3823 table.add_column(DataColumn::new("status"));
3824
3825 table
3827 .add_row(DataRow::new(vec![
3828 DataValue::Integer(1),
3829 DataValue::String("Pending Review".to_string()),
3830 ]))
3831 .unwrap();
3832
3833 table
3834 .add_row(DataRow::new(vec![
3835 DataValue::Integer(2),
3836 DataValue::String("Complete".to_string()),
3837 ]))
3838 .unwrap();
3839
3840 table
3841 .add_row(DataRow::new(vec![
3842 DataValue::Integer(3),
3843 DataValue::String("Pending Approval".to_string()),
3844 ]))
3845 .unwrap();
3846
3847 let table = Arc::new(table);
3848 let engine = QueryEngine::with_case_insensitive(true);
3849
3850 println!("\n=== Testing NOT with Method Calls ===");
3851 println!("Table has {} rows", table.row_count());
3852 for i in 0..table.row_count() {
3853 let status = table.get_value(i, 1);
3854 println!("Row {i}: status = {status:?}");
3855 }
3856
3857 println!("\n--- Test: NOT status.Contains('pend') ---");
3859 let result = engine.execute(
3860 table.clone(),
3861 "SELECT * FROM test WHERE NOT status.Contains('pend')",
3862 );
3863 match result {
3864 Ok(view) => {
3865 println!(
3866 "SUCCESS: Found {} rows NOT containing 'pend'",
3867 view.row_count()
3868 );
3869 assert_eq!(view.row_count(), 1); }
3871 Err(e) => {
3872 panic!("NOT Contains query failed: {e}");
3873 }
3874 }
3875
3876 println!("\n--- Test: NOT status.StartsWith('Pending') ---");
3878 let result = engine.execute(
3879 table.clone(),
3880 "SELECT * FROM test WHERE NOT status.StartsWith('Pending')",
3881 );
3882 match result {
3883 Ok(view) => {
3884 println!(
3885 "SUCCESS: Found {} rows NOT starting with 'Pending'",
3886 view.row_count()
3887 );
3888 assert_eq!(view.row_count(), 1); }
3890 Err(e) => {
3891 panic!("NOT StartsWith query failed: {e}");
3892 }
3893 }
3894
3895 println!("\n=== NOT with method calls test complete! ===");
3896 }
3897
3898 #[test]
3899 #[ignore = "Complex logical expressions with parentheses not yet implemented"]
3900 fn test_complex_logical_expressions() {
3901 let _ = tracing_subscriber::fmt()
3903 .with_max_level(tracing::Level::DEBUG)
3904 .try_init();
3905
3906 let mut table = DataTable::new("test");
3907 table.add_column(DataColumn::new("id"));
3908 table.add_column(DataColumn::new("status"));
3909 table.add_column(DataColumn::new("priority"));
3910 table.add_column(DataColumn::new("assigned"));
3911
3912 table
3914 .add_row(DataRow::new(vec![
3915 DataValue::Integer(1),
3916 DataValue::String("Pending".to_string()),
3917 DataValue::String("High".to_string()),
3918 DataValue::String("John".to_string()),
3919 ]))
3920 .unwrap();
3921
3922 table
3923 .add_row(DataRow::new(vec![
3924 DataValue::Integer(2),
3925 DataValue::String("Complete".to_string()),
3926 DataValue::String("High".to_string()),
3927 DataValue::String("Jane".to_string()),
3928 ]))
3929 .unwrap();
3930
3931 table
3932 .add_row(DataRow::new(vec![
3933 DataValue::Integer(3),
3934 DataValue::String("Pending".to_string()),
3935 DataValue::String("Low".to_string()),
3936 DataValue::String("John".to_string()),
3937 ]))
3938 .unwrap();
3939
3940 table
3941 .add_row(DataRow::new(vec![
3942 DataValue::Integer(4),
3943 DataValue::String("In Progress".to_string()),
3944 DataValue::String("Medium".to_string()),
3945 DataValue::String("Jane".to_string()),
3946 ]))
3947 .unwrap();
3948
3949 let table = Arc::new(table);
3950 let engine = QueryEngine::new();
3951
3952 println!("\n=== Testing Complex Logical Expressions ===");
3953 println!("Table has {} rows", table.row_count());
3954 for i in 0..table.row_count() {
3955 let status = table.get_value(i, 1);
3956 let priority = table.get_value(i, 2);
3957 let assigned = table.get_value(i, 3);
3958 println!(
3959 "Row {i}: status = {status:?}, priority = {priority:?}, assigned = {assigned:?}"
3960 );
3961 }
3962
3963 println!("\n--- Test: status = 'Pending' AND (priority = 'High' OR assigned = 'John') ---");
3965 let result = engine.execute(
3966 table.clone(),
3967 "SELECT * FROM test WHERE status = 'Pending' AND (priority = 'High' OR assigned = 'John')",
3968 );
3969 match result {
3970 Ok(view) => {
3971 println!(
3972 "SUCCESS: Found {} rows with complex logic",
3973 view.row_count()
3974 );
3975 assert_eq!(view.row_count(), 2); }
3977 Err(e) => {
3978 panic!("Complex logic query failed: {e}");
3979 }
3980 }
3981
3982 println!("\n--- Test: NOT (status.Contains('Complete') OR priority = 'Low') ---");
3984 let result = engine.execute(
3985 table.clone(),
3986 "SELECT * FROM test WHERE NOT (status.Contains('Complete') OR priority = 'Low')",
3987 );
3988 match result {
3989 Ok(view) => {
3990 println!(
3991 "SUCCESS: Found {} rows with NOT complex logic",
3992 view.row_count()
3993 );
3994 assert_eq!(view.row_count(), 2); }
3996 Err(e) => {
3997 panic!("NOT complex logic query failed: {e}");
3998 }
3999 }
4000
4001 println!("\n=== Complex logical expressions test complete! ===");
4002 }
4003
4004 #[test]
4005 fn test_mixed_data_types_and_edge_cases() {
4006 let _ = tracing_subscriber::fmt()
4008 .with_max_level(tracing::Level::DEBUG)
4009 .try_init();
4010
4011 let mut table = DataTable::new("test");
4012 table.add_column(DataColumn::new("id"));
4013 table.add_column(DataColumn::new("value"));
4014 table.add_column(DataColumn::new("nullable_field"));
4015
4016 table
4018 .add_row(DataRow::new(vec![
4019 DataValue::Integer(1),
4020 DataValue::String("123.45".to_string()),
4021 DataValue::String("present".to_string()),
4022 ]))
4023 .unwrap();
4024
4025 table
4026 .add_row(DataRow::new(vec![
4027 DataValue::Integer(2),
4028 DataValue::Float(678.90),
4029 DataValue::Null,
4030 ]))
4031 .unwrap();
4032
4033 table
4034 .add_row(DataRow::new(vec![
4035 DataValue::Integer(3),
4036 DataValue::Boolean(true),
4037 DataValue::String("also present".to_string()),
4038 ]))
4039 .unwrap();
4040
4041 table
4042 .add_row(DataRow::new(vec![
4043 DataValue::Integer(4),
4044 DataValue::String("false".to_string()),
4045 DataValue::Null,
4046 ]))
4047 .unwrap();
4048
4049 let table = Arc::new(table);
4050 let engine = QueryEngine::new();
4051
4052 println!("\n=== Testing Mixed Data Types and Edge Cases ===");
4053 println!("Table has {} rows", table.row_count());
4054 for i in 0..table.row_count() {
4055 let value = table.get_value(i, 1);
4056 let nullable = table.get_value(i, 2);
4057 println!("Row {i}: value = {value:?}, nullable_field = {nullable:?}");
4058 }
4059
4060 println!("\n--- Test: value.Contains('true') (boolean to string coercion) ---");
4062 let result = engine.execute(
4063 table.clone(),
4064 "SELECT * FROM test WHERE value.Contains('true')",
4065 );
4066 match result {
4067 Ok(view) => {
4068 println!(
4069 "SUCCESS: Found {} rows with boolean coercion",
4070 view.row_count()
4071 );
4072 assert_eq!(view.row_count(), 1); }
4074 Err(e) => {
4075 panic!("Boolean coercion query failed: {e}");
4076 }
4077 }
4078
4079 println!("\n--- Test: id IN (1, 3) ---");
4081 let result = engine.execute(table.clone(), "SELECT * FROM test WHERE id IN (1, 3)");
4082 match result {
4083 Ok(view) => {
4084 println!("SUCCESS: Found {} rows with IN clause", view.row_count());
4085 assert_eq!(view.row_count(), 2); }
4087 Err(e) => {
4088 panic!("Multiple IN values query failed: {e}");
4089 }
4090 }
4091
4092 println!("\n=== Mixed data types test complete! ===");
4093 }
4094
4095 #[test]
4097 fn test_aggregate_only_single_row() {
4098 let table = create_test_stock_data();
4099 let engine = QueryEngine::new();
4100
4101 let result = engine
4103 .execute(
4104 table.clone(),
4105 "SELECT COUNT(*), MIN(close), MAX(close), AVG(close) FROM stock",
4106 )
4107 .expect("Query should succeed");
4108
4109 assert_eq!(
4110 result.row_count(),
4111 1,
4112 "Aggregate-only query should return exactly 1 row"
4113 );
4114 assert_eq!(result.column_count(), 4, "Should have 4 aggregate columns");
4115
4116 let source = result.source();
4118 let row = source.get_row(0).expect("Should have first row");
4119
4120 assert_eq!(row.values[0], DataValue::Integer(5));
4122
4123 assert_eq!(row.values[1], DataValue::Float(99.5));
4125
4126 assert_eq!(row.values[2], DataValue::Float(105.0));
4128
4129 if let DataValue::Float(avg) = &row.values[3] {
4131 assert!(
4132 (avg - 102.4).abs() < 0.01,
4133 "Average should be approximately 102.4, got {}",
4134 avg
4135 );
4136 } else {
4137 panic!("AVG should return a Float value");
4138 }
4139 }
4140
4141 #[test]
4143 fn test_single_aggregate_single_row() {
4144 let table = create_test_stock_data();
4145 let engine = QueryEngine::new();
4146
4147 let result = engine
4148 .execute(table.clone(), "SELECT COUNT(*) FROM stock")
4149 .expect("Query should succeed");
4150
4151 assert_eq!(
4152 result.row_count(),
4153 1,
4154 "Single aggregate query should return exactly 1 row"
4155 );
4156 assert_eq!(result.column_count(), 1, "Should have 1 column");
4157
4158 let source = result.source();
4159 let row = source.get_row(0).expect("Should have first row");
4160 assert_eq!(row.values[0], DataValue::Integer(5));
4161 }
4162
4163 #[test]
4165 fn test_aggregate_with_where_single_row() {
4166 let table = create_test_stock_data();
4167 let engine = QueryEngine::new();
4168
4169 let result = engine
4171 .execute(
4172 table.clone(),
4173 "SELECT COUNT(*), MIN(close), MAX(close) FROM stock WHERE close >= 103.0",
4174 )
4175 .expect("Query should succeed");
4176
4177 assert_eq!(
4178 result.row_count(),
4179 1,
4180 "Filtered aggregate query should return exactly 1 row"
4181 );
4182 assert_eq!(result.column_count(), 3, "Should have 3 aggregate columns");
4183
4184 let source = result.source();
4185 let row = source.get_row(0).expect("Should have first row");
4186
4187 assert_eq!(row.values[0], DataValue::Integer(2));
4189 assert_eq!(row.values[1], DataValue::Float(103.5)); assert_eq!(row.values[2], DataValue::Float(105.0)); }
4192
4193 #[test]
4194 fn test_not_in_parsing() {
4195 use crate::sql::recursive_parser::Parser;
4196
4197 let query = "SELECT * FROM test WHERE country NOT IN ('CA')";
4198 println!("\n=== Testing NOT IN parsing ===");
4199 println!("Parsing query: {query}");
4200
4201 let mut parser = Parser::new(query);
4202 match parser.parse() {
4203 Ok(statement) => {
4204 println!("Parsed statement: {statement:#?}");
4205 if let Some(where_clause) = statement.where_clause {
4206 println!("WHERE conditions: {:#?}", where_clause.conditions);
4207 if let Some(first_condition) = where_clause.conditions.first() {
4208 println!("First condition expression: {:#?}", first_condition.expr);
4209 }
4210 }
4211 }
4212 Err(e) => {
4213 panic!("Parse error: {e}");
4214 }
4215 }
4216 }
4217
4218 fn create_test_stock_data() -> Arc<DataTable> {
4220 let mut table = DataTable::new("stock");
4221
4222 table.add_column(DataColumn::new("symbol"));
4223 table.add_column(DataColumn::new("close"));
4224 table.add_column(DataColumn::new("volume"));
4225
4226 let test_data = vec![
4228 ("AAPL", 99.5, 1000),
4229 ("AAPL", 101.2, 1500),
4230 ("AAPL", 103.5, 2000),
4231 ("AAPL", 105.0, 1200),
4232 ("AAPL", 102.8, 1800),
4233 ];
4234
4235 for (symbol, close, volume) in test_data {
4236 table
4237 .add_row(DataRow::new(vec![
4238 DataValue::String(symbol.to_string()),
4239 DataValue::Float(close),
4240 DataValue::Integer(volume),
4241 ]))
4242 .expect("Should add row successfully");
4243 }
4244
4245 Arc::new(table)
4246 }
4247}
4248
4249#[cfg(test)]
4250#[path = "query_engine_tests.rs"]
4251mod query_engine_tests;