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 let right_keys: std::collections::HashSet<String> = next_table
1980 .rows
1981 .iter()
1982 .map(|r| format!("{:?}", r.values))
1983 .collect();
1984 let mut seen = std::collections::HashSet::new();
1985 let retained: Vec<_> = combined_table
1986 .rows
1987 .iter()
1988 .filter(|r| {
1989 let key = format!("{:?}", r.values);
1990 right_keys.contains(&key) && seen.insert(key)
1991 })
1992 .cloned()
1993 .collect();
1994 combined_table.rows = retained;
1995 plan.add_detail(format!(
1996 "Result: {} rows (intersection, deduplicated)",
1997 combined_table.row_count()
1998 ));
1999 }
2000 SetOperation::Except => {
2001 let right_keys: std::collections::HashSet<String> = next_table
2004 .rows
2005 .iter()
2006 .map(|r| format!("{:?}", r.values))
2007 .collect();
2008 let mut seen = std::collections::HashSet::new();
2009 let retained: Vec<_> = combined_table
2010 .rows
2011 .iter()
2012 .filter(|r| {
2013 let key = format!("{:?}", r.values);
2014 !right_keys.contains(&key) && seen.insert(key)
2015 })
2016 .cloned()
2017 .collect();
2018 combined_table.rows = retained;
2019 plan.add_detail(format!(
2020 "Result: {} rows (difference, deduplicated)",
2021 combined_table.row_count()
2022 ));
2023 }
2024 }
2025
2026 plan.add_detail(format!(
2027 "Operation time: {:.3}ms",
2028 op_start.elapsed().as_secs_f64() * 1000.0
2029 ));
2030 plan.set_rows_out(combined_table.row_count());
2031 plan.end_step();
2032 }
2033
2034 plan.set_rows_out(combined_table.row_count());
2035 plan.add_detail(format!(
2036 "Combined result: {} rows after {} operations",
2037 combined_table.row_count(),
2038 statement.set_operations.len()
2039 ));
2040 plan.end_step();
2041
2042 view = DataView::new(Arc::new(combined_table));
2044
2045 if needs_deduplication {
2047 plan.begin_step(
2048 StepType::Distinct,
2049 "UNION deduplication - remove duplicate rows".to_string(),
2050 );
2051 plan.set_rows_in(view.row_count());
2052 plan.add_detail(format!("Input: {} rows", view.row_count()));
2053
2054 let distinct_start = Instant::now();
2055 view = self.apply_distinct(view)?;
2056
2057 plan.set_rows_out(view.row_count());
2058 plan.add_detail(format!("Output: {} unique rows", view.row_count()));
2059 plan.add_detail(format!(
2060 "Deduplication time: {:.3}ms",
2061 distinct_start.elapsed().as_secs_f64() * 1000.0
2062 ));
2063 plan.end_step();
2064 }
2065 }
2066
2067 Ok(view)
2068 }
2069
2070 fn resolve_column_indices(&self, table: &DataTable, columns: &[String]) -> Result<Vec<usize>> {
2072 let mut indices = Vec::new();
2073 let table_columns = table.column_names();
2074
2075 for col_name in columns {
2076 let index = table_columns
2077 .iter()
2078 .position(|c| c.eq_ignore_ascii_case(col_name))
2079 .ok_or_else(|| {
2080 let suggestion = self.find_similar_column(table, col_name);
2081 match suggestion {
2082 Some(similar) => anyhow::anyhow!(
2083 "Column '{}' not found. Did you mean '{}'?",
2084 col_name,
2085 similar
2086 ),
2087 None => anyhow::anyhow!("Column '{}' not found", col_name),
2088 }
2089 })?;
2090 indices.push(index);
2091 }
2092
2093 Ok(indices)
2094 }
2095
2096 fn apply_select_items(
2098 &self,
2099 view: DataView,
2100 select_items: &[SelectItem],
2101 _statement: &SelectStatement,
2102 exec_context: Option<&ExecutionContext>,
2103 plan: &mut ExecutionPlanBuilder,
2104 ) -> Result<DataView> {
2105 debug!(
2106 "QueryEngine::apply_select_items - items: {:?}",
2107 select_items
2108 );
2109 debug!(
2110 "QueryEngine::apply_select_items - input view has {} rows",
2111 view.row_count()
2112 );
2113
2114 let has_window_functions = select_items.iter().any(|item| match item {
2116 SelectItem::Expression { expr, .. } => Self::contains_window_function(expr),
2117 _ => false,
2118 });
2119
2120 let window_func_count: usize = select_items
2122 .iter()
2123 .filter(|item| match item {
2124 SelectItem::Expression { expr, .. } => Self::contains_window_function(expr),
2125 _ => false,
2126 })
2127 .count();
2128
2129 let window_start = if has_window_functions {
2131 debug!(
2132 "QueryEngine::apply_select_items - detected {} window functions",
2133 window_func_count
2134 );
2135
2136 let window_specs = Self::extract_window_specs(select_items);
2138 debug!("Extracted {} window function specs", window_specs.len());
2139
2140 Some(Instant::now())
2141 } else {
2142 None
2143 };
2144
2145 let has_unnest = select_items.iter().any(|item| match item {
2147 SelectItem::Expression { expr, .. } => Self::contains_unnest(expr),
2148 _ => false,
2149 });
2150
2151 if has_unnest {
2152 debug!("QueryEngine::apply_select_items - UNNEST detected, using row expansion");
2153 return self.apply_select_with_row_expansion(view, select_items);
2154 }
2155
2156 let has_aggregates = select_items.iter().any(|item| match item {
2160 SelectItem::Expression { expr, .. } => contains_aggregate(expr),
2161 SelectItem::Column { .. } => false,
2162 SelectItem::Star { .. } => false,
2163 SelectItem::StarExclude { .. } => false,
2164 });
2165
2166 let all_aggregate_compatible = select_items.iter().all(|item| match item {
2167 SelectItem::Expression { expr, .. } => is_aggregate_compatible(expr),
2168 SelectItem::Column { .. } => false, SelectItem::Star { .. } => false, SelectItem::StarExclude { .. } => false, });
2172
2173 if has_aggregates && all_aggregate_compatible && view.row_count() > 0 {
2174 debug!("QueryEngine::apply_select_items - detected aggregate query with constants");
2177 return self.apply_aggregate_select(view, select_items);
2178 }
2179
2180 let has_computed_expressions = select_items
2182 .iter()
2183 .any(|item| matches!(item, SelectItem::Expression { .. }));
2184
2185 debug!(
2186 "QueryEngine::apply_select_items - has_computed_expressions: {}",
2187 has_computed_expressions
2188 );
2189
2190 if !has_computed_expressions {
2191 let column_indices = self.resolve_select_columns(view.source(), select_items)?;
2193 return Ok(view.with_columns(column_indices));
2194 }
2195
2196 let source_table = view.source();
2201 let visible_rows = view.visible_row_indices();
2202
2203 let mut computed_table = DataTable::new("query_result");
2206
2207 let mut expanded_items = Vec::new();
2209 for item in select_items {
2210 match item {
2211 SelectItem::Star { table_prefix, .. } => {
2212 if let Some(prefix) = table_prefix {
2213 debug!("QueryEngine::apply_select_items - expanding {}.*", prefix);
2215 for col in &source_table.columns {
2216 if Self::column_matches_table(col, prefix) {
2217 expanded_items.push(SelectItem::Column {
2218 column: ColumnRef::unquoted(col.name.clone()),
2219 leading_comments: vec![],
2220 trailing_comment: None,
2221 });
2222 }
2223 }
2224 } else {
2225 debug!("QueryEngine::apply_select_items - expanding *");
2227 for col_name in source_table.column_names() {
2228 expanded_items.push(SelectItem::Column {
2229 column: ColumnRef::unquoted(col_name.to_string()),
2230 leading_comments: vec![],
2231 trailing_comment: None,
2232 });
2233 }
2234 }
2235 }
2236 _ => expanded_items.push(item.clone()),
2237 }
2238 }
2239
2240 let mut column_name_counts: std::collections::HashMap<String, usize> =
2242 std::collections::HashMap::new();
2243
2244 for item in &expanded_items {
2245 let base_name = match item {
2246 SelectItem::Column {
2247 column: col_ref, ..
2248 } => col_ref.name.clone(),
2249 SelectItem::Expression { alias, .. } => alias.clone(),
2250 SelectItem::Star { .. } => unreachable!("Star should have been expanded"),
2251 SelectItem::StarExclude { .. } => {
2252 unreachable!("StarExclude should have been expanded")
2253 }
2254 };
2255
2256 let count = column_name_counts.entry(base_name.clone()).or_insert(0);
2258 let column_name = if *count == 0 {
2259 base_name.clone()
2261 } else {
2262 format!("{base_name}_{count}")
2264 };
2265 *count += 1;
2266
2267 computed_table.add_column(DataColumn::new(&column_name));
2268 }
2269
2270 let can_use_batch = expanded_items.iter().all(|item| {
2274 match item {
2275 SelectItem::Expression { expr, .. } => {
2276 matches!(expr, SqlExpression::WindowFunction { .. })
2279 || !Self::contains_window_function(expr)
2280 }
2281 _ => true, }
2283 });
2284
2285 let use_batch_evaluation = can_use_batch
2288 && std::env::var("SQL_CLI_BATCH_WINDOW")
2289 .map(|v| v != "0" && v.to_lowercase() != "false")
2290 .unwrap_or(true);
2291
2292 let batch_window_specs = if use_batch_evaluation && has_window_functions {
2294 debug!("BATCH window function evaluation flag is enabled");
2295 let specs = Self::extract_window_specs(&expanded_items);
2297 debug!(
2298 "Extracted {} window function specs for batch evaluation",
2299 specs.len()
2300 );
2301 Some(specs)
2302 } else {
2303 None
2304 };
2305
2306 let mut evaluator =
2308 ArithmeticEvaluator::with_date_notation(source_table, self.date_notation.clone());
2309
2310 if let Some(exec_ctx) = exec_context {
2312 let aliases = exec_ctx.get_aliases();
2313 if !aliases.is_empty() {
2314 debug!(
2315 "Applying {} aliases to evaluator: {:?}",
2316 aliases.len(),
2317 aliases
2318 );
2319 evaluator = evaluator.with_table_aliases(aliases);
2320 }
2321 }
2322
2323 if has_window_functions {
2326 let preload_start = Instant::now();
2327
2328 let mut window_specs = Vec::new();
2330 for item in &expanded_items {
2331 if let SelectItem::Expression { expr, .. } = item {
2332 Self::collect_window_specs(expr, &mut window_specs);
2333 }
2334 }
2335
2336 for spec in &window_specs {
2338 let _ = evaluator.get_or_create_window_context(spec);
2339 }
2340
2341 debug!(
2342 "Pre-created {} WindowContext(s) in {:.2}ms",
2343 window_specs.len(),
2344 preload_start.elapsed().as_secs_f64() * 1000.0
2345 );
2346 }
2347
2348 if let Some(window_specs) = batch_window_specs {
2350 debug!("Starting batch window function evaluation");
2351 let batch_start = Instant::now();
2352
2353 let mut batch_results: Vec<Vec<DataValue>> =
2355 vec![vec![DataValue::Null; expanded_items.len()]; visible_rows.len()];
2356
2357 let detailed_window_specs = &window_specs;
2359
2360 let mut specs_by_window: HashMap<
2362 u64,
2363 Vec<&crate::data::batch_window_evaluator::WindowFunctionSpec>,
2364 > = HashMap::new();
2365 for spec in detailed_window_specs {
2366 let hash = spec.spec.compute_hash();
2367 specs_by_window
2368 .entry(hash)
2369 .or_insert_with(Vec::new)
2370 .push(spec);
2371 }
2372
2373 for (_window_hash, specs) in specs_by_window {
2375 let context = evaluator.get_or_create_window_context(&specs[0].spec)?;
2377
2378 for spec in specs {
2380 match spec.function_name.as_str() {
2381 "LAG" => {
2382 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2384 let column_name = col_ref.name.as_str();
2385 let offset = if let Some(SqlExpression::NumberLiteral(n)) =
2386 spec.args.get(1)
2387 {
2388 n.parse::<i64>().unwrap_or(1)
2389 } else {
2390 1 };
2392
2393 let values = context.evaluate_lag_batch(
2394 visible_rows,
2395 column_name,
2396 offset,
2397 )?;
2398
2399 for (row_idx, value) in values.into_iter().enumerate() {
2401 batch_results[row_idx][spec.output_column_index] = value;
2402 }
2403 }
2404 }
2405 "LEAD" => {
2406 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2408 let column_name = col_ref.name.as_str();
2409 let offset = if let Some(SqlExpression::NumberLiteral(n)) =
2410 spec.args.get(1)
2411 {
2412 n.parse::<i64>().unwrap_or(1)
2413 } else {
2414 1 };
2416
2417 let values = context.evaluate_lead_batch(
2418 visible_rows,
2419 column_name,
2420 offset,
2421 )?;
2422
2423 for (row_idx, value) in values.into_iter().enumerate() {
2425 batch_results[row_idx][spec.output_column_index] = value;
2426 }
2427 }
2428 }
2429 "ROW_NUMBER" => {
2430 let values = context.evaluate_row_number_batch(visible_rows)?;
2431
2432 for (row_idx, value) in values.into_iter().enumerate() {
2434 batch_results[row_idx][spec.output_column_index] = value;
2435 }
2436 }
2437 "RANK" => {
2438 let values = context.evaluate_rank_batch(visible_rows)?;
2439
2440 for (row_idx, value) in values.into_iter().enumerate() {
2442 batch_results[row_idx][spec.output_column_index] = value;
2443 }
2444 }
2445 "DENSE_RANK" => {
2446 let values = context.evaluate_dense_rank_batch(visible_rows)?;
2447
2448 for (row_idx, value) in values.into_iter().enumerate() {
2450 batch_results[row_idx][spec.output_column_index] = value;
2451 }
2452 }
2453 "SUM" => {
2454 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2455 let column_name = col_ref.name.as_str();
2456 let values =
2457 context.evaluate_sum_batch(visible_rows, column_name)?;
2458
2459 for (row_idx, value) in values.into_iter().enumerate() {
2460 batch_results[row_idx][spec.output_column_index] = value;
2461 }
2462 }
2463 }
2464 "AVG" => {
2465 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2466 let column_name = col_ref.name.as_str();
2467 let values =
2468 context.evaluate_avg_batch(visible_rows, column_name)?;
2469
2470 for (row_idx, value) in values.into_iter().enumerate() {
2471 batch_results[row_idx][spec.output_column_index] = value;
2472 }
2473 }
2474 }
2475 "MIN" => {
2476 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2477 let column_name = col_ref.name.as_str();
2478 let values =
2479 context.evaluate_min_batch(visible_rows, column_name)?;
2480
2481 for (row_idx, value) in values.into_iter().enumerate() {
2482 batch_results[row_idx][spec.output_column_index] = value;
2483 }
2484 }
2485 }
2486 "MAX" => {
2487 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2488 let column_name = col_ref.name.as_str();
2489 let values =
2490 context.evaluate_max_batch(visible_rows, column_name)?;
2491
2492 for (row_idx, value) in values.into_iter().enumerate() {
2493 batch_results[row_idx][spec.output_column_index] = value;
2494 }
2495 }
2496 }
2497 "COUNT" => {
2498 let column_name = match spec.args.get(0) {
2500 Some(SqlExpression::Column(col_ref)) => Some(col_ref.name.as_str()),
2501 Some(SqlExpression::StringLiteral(s)) if s == "*" => None,
2502 _ => None,
2503 };
2504
2505 let values = context.evaluate_count_batch(visible_rows, column_name)?;
2506
2507 for (row_idx, value) in values.into_iter().enumerate() {
2508 batch_results[row_idx][spec.output_column_index] = value;
2509 }
2510 }
2511 "FIRST_VALUE" => {
2512 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2513 let column_name = col_ref.name.as_str();
2514 let values = context
2515 .evaluate_first_value_batch(visible_rows, column_name)?;
2516
2517 for (row_idx, value) in values.into_iter().enumerate() {
2518 batch_results[row_idx][spec.output_column_index] = value;
2519 }
2520 }
2521 }
2522 "LAST_VALUE" => {
2523 if let Some(SqlExpression::Column(col_ref)) = spec.args.get(0) {
2524 let column_name = col_ref.name.as_str();
2525 let values =
2526 context.evaluate_last_value_batch(visible_rows, column_name)?;
2527
2528 for (row_idx, value) in values.into_iter().enumerate() {
2529 batch_results[row_idx][spec.output_column_index] = value;
2530 }
2531 }
2532 }
2533 _ => {
2534 debug!(
2536 "Window function {} not supported in batch mode, using per-row",
2537 spec.function_name
2538 );
2539 }
2540 }
2541 }
2542 }
2543
2544 for (result_row_idx, &source_row_idx) in visible_rows.iter().enumerate() {
2546 for (col_idx, item) in expanded_items.iter().enumerate() {
2547 if !matches!(batch_results[result_row_idx][col_idx], DataValue::Null) {
2549 continue;
2550 }
2551
2552 let value = match item {
2553 SelectItem::Column {
2554 column: col_ref, ..
2555 } => {
2556 match evaluator
2557 .evaluate(&SqlExpression::Column(col_ref.clone()), source_row_idx)
2558 {
2559 Ok(val) => val,
2560 Err(e) => {
2561 return Err(anyhow!(
2562 "Failed to evaluate column {}: {}",
2563 col_ref.to_sql(),
2564 e
2565 ));
2566 }
2567 }
2568 }
2569 SelectItem::Expression { expr, .. } => {
2570 if matches!(expr, SqlExpression::WindowFunction { .. }) {
2573 continue;
2575 }
2576 evaluator.evaluate(&expr, source_row_idx)?
2579 }
2580 SelectItem::Star { .. } => unreachable!("Star should have been expanded"),
2581 SelectItem::StarExclude { .. } => {
2582 unreachable!("StarExclude should have been expanded")
2583 }
2584 };
2585 batch_results[result_row_idx][col_idx] = value;
2586 }
2587 }
2588
2589 for row_values in batch_results {
2591 computed_table
2592 .add_row(DataRow::new(row_values))
2593 .map_err(|e| anyhow::anyhow!("Failed to add row: {}", e))?;
2594 }
2595
2596 debug!(
2597 "Batch window evaluation completed in {:.3}ms",
2598 batch_start.elapsed().as_secs_f64() * 1000.0
2599 );
2600 } else {
2601 for &row_idx in visible_rows {
2603 let mut row_values = Vec::new();
2604
2605 for item in &expanded_items {
2606 let value = match item {
2607 SelectItem::Column {
2608 column: col_ref, ..
2609 } => {
2610 match evaluator
2612 .evaluate(&SqlExpression::Column(col_ref.clone()), row_idx)
2613 {
2614 Ok(val) => val,
2615 Err(e) => {
2616 return Err(anyhow!(
2617 "Failed to evaluate column {}: {}",
2618 col_ref.to_sql(),
2619 e
2620 ));
2621 }
2622 }
2623 }
2624 SelectItem::Expression { expr, .. } => {
2625 evaluator.evaluate(&expr, row_idx)?
2627 }
2628 SelectItem::Star { .. } => unreachable!("Star should have been expanded"),
2629 SelectItem::StarExclude { .. } => {
2630 unreachable!("StarExclude should have been expanded")
2631 }
2632 };
2633 row_values.push(value);
2634 }
2635
2636 computed_table
2637 .add_row(DataRow::new(row_values))
2638 .map_err(|e| anyhow::anyhow!("Failed to add row: {}", e))?;
2639 }
2640 }
2641
2642 if let Some(start) = window_start {
2644 let window_duration = start.elapsed();
2645 info!(
2646 "Window function evaluation took {:.2}ms for {} rows ({} window functions)",
2647 window_duration.as_secs_f64() * 1000.0,
2648 visible_rows.len(),
2649 window_func_count
2650 );
2651
2652 plan.begin_step(
2654 StepType::WindowFunction,
2655 format!("Evaluate {} window function(s)", window_func_count),
2656 );
2657 plan.set_rows_in(visible_rows.len());
2658 plan.set_rows_out(visible_rows.len());
2659 plan.add_detail(format!("Input: {} rows", visible_rows.len()));
2660 plan.add_detail(format!("{} window functions evaluated", window_func_count));
2661 plan.add_detail(format!(
2662 "Evaluation time: {:.3}ms",
2663 window_duration.as_secs_f64() * 1000.0
2664 ));
2665 plan.end_step();
2666 }
2667
2668 Ok(DataView::new(Arc::new(computed_table)))
2671 }
2672
2673 fn apply_select_with_row_expansion(
2675 &self,
2676 view: DataView,
2677 select_items: &[SelectItem],
2678 ) -> Result<DataView> {
2679 debug!("QueryEngine::apply_select_with_row_expansion - expanding rows");
2680
2681 let source_table = view.source();
2682 let visible_rows = view.visible_row_indices();
2683 let expander_registry = RowExpanderRegistry::new();
2684
2685 let mut result_table = DataTable::new("unnest_result");
2687
2688 let mut expanded_items = Vec::new();
2690 for item in select_items {
2691 match item {
2692 SelectItem::Star { table_prefix, .. } => {
2693 if let Some(prefix) = table_prefix {
2694 debug!(
2696 "QueryEngine::apply_select_with_row_expansion - expanding {}.*",
2697 prefix
2698 );
2699 for col in &source_table.columns {
2700 if Self::column_matches_table(col, prefix) {
2701 expanded_items.push(SelectItem::Column {
2702 column: ColumnRef::unquoted(col.name.clone()),
2703 leading_comments: vec![],
2704 trailing_comment: None,
2705 });
2706 }
2707 }
2708 } else {
2709 debug!("QueryEngine::apply_select_with_row_expansion - expanding *");
2711 for col_name in source_table.column_names() {
2712 expanded_items.push(SelectItem::Column {
2713 column: ColumnRef::unquoted(col_name.to_string()),
2714 leading_comments: vec![],
2715 trailing_comment: None,
2716 });
2717 }
2718 }
2719 }
2720 _ => expanded_items.push(item.clone()),
2721 }
2722 }
2723
2724 for item in &expanded_items {
2726 let column_name = match item {
2727 SelectItem::Column {
2728 column: col_ref, ..
2729 } => col_ref.name.clone(),
2730 SelectItem::Expression { alias, .. } => alias.clone(),
2731 SelectItem::Star { .. } => unreachable!("Star should have been expanded"),
2732 SelectItem::StarExclude { .. } => {
2733 unreachable!("StarExclude should have been expanded")
2734 }
2735 };
2736 result_table.add_column(DataColumn::new(&column_name));
2737 }
2738
2739 let mut evaluator =
2741 ArithmeticEvaluator::with_date_notation(source_table, self.date_notation.clone());
2742
2743 for &row_idx in visible_rows {
2744 let mut unnest_expansions = Vec::new();
2746 let mut unnest_indices = Vec::new();
2747
2748 for (col_idx, item) in expanded_items.iter().enumerate() {
2749 if let SelectItem::Expression { expr, .. } = item {
2750 if let Some(expansion_result) = self.try_expand_unnest(
2751 &expr,
2752 source_table,
2753 row_idx,
2754 &mut evaluator,
2755 &expander_registry,
2756 )? {
2757 unnest_expansions.push(expansion_result);
2758 unnest_indices.push(col_idx);
2759 }
2760 }
2761 }
2762
2763 let expansion_count = if unnest_expansions.is_empty() {
2765 1 } else {
2767 unnest_expansions
2768 .iter()
2769 .map(|exp| exp.row_count())
2770 .max()
2771 .unwrap_or(1)
2772 };
2773
2774 for output_idx in 0..expansion_count {
2776 let mut row_values = Vec::new();
2777
2778 for (col_idx, item) in expanded_items.iter().enumerate() {
2779 let unnest_position = unnest_indices.iter().position(|&idx| idx == col_idx);
2781
2782 let value = if let Some(unnest_idx) = unnest_position {
2783 let expansion = &unnest_expansions[unnest_idx];
2785 expansion
2786 .values
2787 .get(output_idx)
2788 .cloned()
2789 .unwrap_or(DataValue::Null)
2790 } else {
2791 match item {
2793 SelectItem::Column {
2794 column: col_ref, ..
2795 } => {
2796 let col_idx =
2797 source_table.get_column_index(&col_ref.name).ok_or_else(
2798 || anyhow::anyhow!("Column '{}' not found", col_ref.name),
2799 )?;
2800 let row = source_table
2801 .get_row(row_idx)
2802 .ok_or_else(|| anyhow::anyhow!("Row {} not found", row_idx))?;
2803 row.get(col_idx)
2804 .ok_or_else(|| {
2805 anyhow::anyhow!("Column {} not found in row", col_idx)
2806 })?
2807 .clone()
2808 }
2809 SelectItem::Expression { expr, .. } => {
2810 evaluator.evaluate(&expr, row_idx)?
2812 }
2813 SelectItem::Star { .. } => unreachable!(),
2814 SelectItem::StarExclude { .. } => {
2815 unreachable!("StarExclude should have been expanded")
2816 }
2817 }
2818 };
2819
2820 row_values.push(value);
2821 }
2822
2823 result_table
2824 .add_row(DataRow::new(row_values))
2825 .map_err(|e| anyhow::anyhow!("Failed to add expanded row: {}", e))?;
2826 }
2827 }
2828
2829 debug!(
2830 "QueryEngine::apply_select_with_row_expansion - input rows: {}, output rows: {}",
2831 visible_rows.len(),
2832 result_table.row_count()
2833 );
2834
2835 Ok(DataView::new(Arc::new(result_table)))
2836 }
2837
2838 fn try_expand_unnest(
2841 &self,
2842 expr: &SqlExpression,
2843 _source_table: &DataTable,
2844 row_idx: usize,
2845 evaluator: &mut ArithmeticEvaluator,
2846 expander_registry: &RowExpanderRegistry,
2847 ) -> Result<Option<crate::data::row_expanders::ExpansionResult>> {
2848 if let SqlExpression::Unnest { column, delimiter } = expr {
2850 let column_value = evaluator.evaluate(column, row_idx)?;
2852
2853 let delimiter_value = DataValue::String(delimiter.clone());
2855
2856 let expander = expander_registry
2858 .get("UNNEST")
2859 .ok_or_else(|| anyhow::anyhow!("UNNEST expander not found"))?;
2860
2861 let expansion = expander.expand(&column_value, &[delimiter_value])?;
2863 return Ok(Some(expansion));
2864 }
2865
2866 if let SqlExpression::FunctionCall { name, args, .. } = expr {
2868 if name.to_uppercase() == "UNNEST" {
2869 if args.len() != 2 {
2871 return Err(anyhow::anyhow!(
2872 "UNNEST requires exactly 2 arguments: UNNEST(column, delimiter)"
2873 ));
2874 }
2875
2876 let column_value = evaluator.evaluate(&args[0], row_idx)?;
2878
2879 let delimiter_value = evaluator.evaluate(&args[1], row_idx)?;
2881
2882 let expander = expander_registry
2884 .get("UNNEST")
2885 .ok_or_else(|| anyhow::anyhow!("UNNEST expander not found"))?;
2886
2887 let expansion = expander.expand(&column_value, &[delimiter_value])?;
2889 return Ok(Some(expansion));
2890 }
2891 }
2892
2893 Ok(None)
2894 }
2895
2896 fn apply_aggregate_select(
2898 &self,
2899 view: DataView,
2900 select_items: &[SelectItem],
2901 ) -> Result<DataView> {
2902 debug!("QueryEngine::apply_aggregate_select - creating single row aggregate result");
2903
2904 let source_table = view.source();
2905 let mut result_table = DataTable::new("aggregate_result");
2906
2907 for item in select_items {
2909 let column_name = match item {
2910 SelectItem::Expression { alias, .. } => alias.clone(),
2911 _ => unreachable!("Should only have expressions in aggregate-only query"),
2912 };
2913 result_table.add_column(DataColumn::new(&column_name));
2914 }
2915
2916 let visible_rows = view.visible_row_indices().to_vec();
2918 let mut evaluator =
2919 ArithmeticEvaluator::with_date_notation(source_table, self.date_notation.clone())
2920 .with_visible_rows(visible_rows);
2921
2922 let mut row_values = Vec::new();
2924 for item in select_items {
2925 match item {
2926 SelectItem::Expression { expr, .. } => {
2927 let value = evaluator.evaluate(expr, 0)?;
2930 row_values.push(value);
2931 }
2932 _ => unreachable!("Should only have expressions in aggregate-only query"),
2933 }
2934 }
2935
2936 result_table
2938 .add_row(DataRow::new(row_values))
2939 .map_err(|e| anyhow::anyhow!("Failed to add aggregate result row: {}", e))?;
2940
2941 Ok(DataView::new(Arc::new(result_table)))
2942 }
2943
2944 fn column_matches_table(col: &DataColumn, table_name: &str) -> bool {
2956 if let Some(ref source) = col.source_table {
2958 if source == table_name || source.ends_with(&format!(".{}", table_name)) {
2960 return true;
2961 }
2962 }
2963
2964 if let Some(ref qualified) = col.qualified_name {
2966 if qualified.starts_with(&format!("{}.", table_name)) {
2968 return true;
2969 }
2970 }
2971
2972 false
2973 }
2974
2975 fn resolve_select_columns(
2977 &self,
2978 table: &DataTable,
2979 select_items: &[SelectItem],
2980 ) -> Result<Vec<usize>> {
2981 let mut indices = Vec::new();
2982 let table_columns = table.column_names();
2983
2984 for item in select_items {
2985 match item {
2986 SelectItem::Column {
2987 column: col_ref, ..
2988 } => {
2989 let index = if let Some(table_prefix) = &col_ref.table_prefix {
2991 let qualified_name = format!("{}.{}", table_prefix, col_ref.name);
3000 table.find_column_by_qualified_name(&qualified_name)
3001 .or_else(|| {
3002 table_columns
3003 .iter()
3004 .position(|c| c.eq_ignore_ascii_case(&col_ref.name))
3005 })
3006 .ok_or_else(|| {
3007 let has_qualified = table.columns.iter()
3009 .any(|c| c.qualified_name.is_some());
3010 if !has_qualified {
3011 anyhow::anyhow!(
3012 "Column '{}' not found. Note: Table '{}' may not support qualified column names",
3013 qualified_name, table_prefix
3014 )
3015 } else {
3016 anyhow::anyhow!("Column '{}' not found", qualified_name)
3017 }
3018 })?
3019 } else {
3020 table_columns
3022 .iter()
3023 .position(|c| c.eq_ignore_ascii_case(&col_ref.name))
3024 .ok_or_else(|| {
3025 let suggestion = self.find_similar_column(table, &col_ref.name);
3026 match suggestion {
3027 Some(similar) => anyhow::anyhow!(
3028 "Column '{}' not found. Did you mean '{}'?",
3029 col_ref.name,
3030 similar
3031 ),
3032 None => anyhow::anyhow!("Column '{}' not found", col_ref.name),
3033 }
3034 })?
3035 };
3036 indices.push(index);
3037 }
3038 SelectItem::Star { table_prefix, .. } => {
3039 if let Some(prefix) = table_prefix {
3040 for (i, col) in table.columns.iter().enumerate() {
3042 if Self::column_matches_table(col, prefix) {
3043 indices.push(i);
3044 }
3045 }
3046 } else {
3047 for i in 0..table_columns.len() {
3049 indices.push(i);
3050 }
3051 }
3052 }
3053 SelectItem::StarExclude {
3054 table_prefix,
3055 excluded_columns,
3056 ..
3057 } => {
3058 if let Some(prefix) = table_prefix {
3060 for (i, col) in table.columns.iter().enumerate() {
3062 if Self::column_matches_table(col, prefix)
3063 && !excluded_columns.contains(&col.name)
3064 {
3065 indices.push(i);
3066 }
3067 }
3068 } else {
3069 for (i, col_name) in table_columns.iter().enumerate() {
3071 if !excluded_columns
3072 .iter()
3073 .any(|exc| exc.eq_ignore_ascii_case(col_name))
3074 {
3075 indices.push(i);
3076 }
3077 }
3078 }
3079 }
3080 SelectItem::Expression { .. } => {
3081 return Err(anyhow::anyhow!(
3082 "Computed expressions require new table creation"
3083 ));
3084 }
3085 }
3086 }
3087
3088 Ok(indices)
3089 }
3090
3091 fn apply_distinct(&self, view: DataView) -> Result<DataView> {
3093 use std::collections::HashSet;
3094
3095 let source = view.source();
3096 let visible_cols = view.visible_column_indices();
3097 let visible_rows = view.visible_row_indices();
3098
3099 let mut seen_rows = HashSet::new();
3101 let mut unique_row_indices = Vec::new();
3102
3103 for &row_idx in visible_rows {
3104 let mut row_key = Vec::new();
3106 for &col_idx in visible_cols {
3107 let value = source
3108 .get_value(row_idx, col_idx)
3109 .ok_or_else(|| anyhow!("Invalid cell reference"))?;
3110 row_key.push(format!("{:?}", value));
3112 }
3113
3114 if seen_rows.insert(row_key) {
3116 unique_row_indices.push(row_idx);
3118 }
3119 }
3120
3121 Ok(view.with_rows(unique_row_indices))
3123 }
3124
3125 fn apply_multi_order_by(
3127 &self,
3128 view: DataView,
3129 order_by_columns: &[OrderByItem],
3130 ) -> Result<DataView> {
3131 self.apply_multi_order_by_with_context(view, order_by_columns, None)
3132 }
3133
3134 fn apply_multi_order_by_with_context(
3136 &self,
3137 mut view: DataView,
3138 order_by_columns: &[OrderByItem],
3139 _exec_context: Option<&ExecutionContext>,
3140 ) -> Result<DataView> {
3141 let mut sort_columns = Vec::new();
3143
3144 for order_col in order_by_columns {
3145 let column_name = match &order_col.expr {
3147 SqlExpression::Column(col_ref) => col_ref.name.clone(),
3148 _ => {
3149 return Err(anyhow!(
3151 "ORDER BY expressions not yet supported - only simple columns allowed"
3152 ));
3153 }
3154 };
3155
3156 let col_index = if column_name.contains('.') {
3158 if let Some(dot_pos) = column_name.rfind('.') {
3160 let col_name = &column_name[dot_pos + 1..];
3161
3162 debug!(
3165 "ORDER BY: Extracting unqualified column '{}' from '{}'",
3166 col_name, column_name
3167 );
3168 view.source().get_column_index(col_name)
3169 } else {
3170 view.source().get_column_index(&column_name)
3171 }
3172 } else {
3173 view.source().get_column_index(&column_name)
3175 }
3176 .ok_or_else(|| {
3177 let suggestion = self.find_similar_column(view.source(), &column_name);
3179 match suggestion {
3180 Some(similar) => anyhow::anyhow!(
3181 "Column '{}' not found. Did you mean '{}'?",
3182 column_name,
3183 similar
3184 ),
3185 None => {
3186 let available_cols = view.source().column_names().join(", ");
3188 anyhow::anyhow!(
3189 "Column '{}' not found. Available columns: {}",
3190 column_name,
3191 available_cols
3192 )
3193 }
3194 }
3195 })?;
3196
3197 let ascending = matches!(order_col.direction, SortDirection::Asc);
3198 sort_columns.push((col_index, ascending));
3199 }
3200
3201 view.apply_multi_sort(&sort_columns)?;
3203 Ok(view)
3204 }
3205
3206 fn apply_group_by(
3208 &self,
3209 view: DataView,
3210 group_by_exprs: &[SqlExpression],
3211 select_items: &[SelectItem],
3212 having: Option<&SqlExpression>,
3213 plan: &mut ExecutionPlanBuilder,
3214 ) -> Result<DataView> {
3215 let (result_view, phase_info) = self.apply_group_by_expressions(
3217 view,
3218 group_by_exprs,
3219 select_items,
3220 having,
3221 self.case_insensitive,
3222 self.date_notation.clone(),
3223 )?;
3224
3225 plan.add_detail(format!("=== GROUP BY Phase Breakdown ==="));
3227 plan.add_detail(format!(
3228 "Phase 1 - Group Building: {:.3}ms",
3229 phase_info.phase2_key_building.as_secs_f64() * 1000.0
3230 ));
3231 plan.add_detail(format!(
3232 " • Processing {} rows into {} groups",
3233 phase_info.total_rows, phase_info.num_groups
3234 ));
3235 plan.add_detail(format!(
3236 "Phase 2 - Aggregation: {:.3}ms",
3237 phase_info.phase4_aggregation.as_secs_f64() * 1000.0
3238 ));
3239 if phase_info.phase4_having_evaluation > Duration::ZERO {
3240 plan.add_detail(format!(
3241 "Phase 3 - HAVING Filter: {:.3}ms",
3242 phase_info.phase4_having_evaluation.as_secs_f64() * 1000.0
3243 ));
3244 plan.add_detail(format!(
3245 " • Filtered {} groups",
3246 phase_info.groups_filtered_by_having
3247 ));
3248 }
3249 plan.add_detail(format!(
3250 "Total GROUP BY time: {:.3}ms",
3251 phase_info.total_time.as_secs_f64() * 1000.0
3252 ));
3253
3254 Ok(result_view)
3255 }
3256
3257 pub fn estimate_group_cardinality(
3260 &self,
3261 view: &DataView,
3262 group_by_exprs: &[SqlExpression],
3263 ) -> usize {
3264 let row_count = view.get_visible_rows().len();
3266 if row_count <= 100 {
3267 return row_count;
3268 }
3269
3270 let sample_size = min(1000, row_count / 10).max(100);
3272 let mut seen = FxHashSet::default();
3273
3274 let visible_rows = view.get_visible_rows();
3275 for (i, &row_idx) in visible_rows.iter().enumerate() {
3276 if i >= sample_size {
3277 break;
3278 }
3279
3280 let mut key_values = Vec::new();
3282 for expr in group_by_exprs {
3283 let mut evaluator = ArithmeticEvaluator::new(view.source());
3284 let value = evaluator.evaluate(expr, row_idx).unwrap_or(DataValue::Null);
3285 key_values.push(value);
3286 }
3287
3288 seen.insert(key_values);
3289 }
3290
3291 let sample_cardinality = seen.len();
3293 let estimated = (sample_cardinality * row_count) / sample_size;
3294
3295 estimated.min(row_count).max(sample_cardinality)
3297 }
3298}
3299
3300#[cfg(test)]
3301mod tests {
3302 use super::*;
3303 use crate::data::datatable::{DataColumn, DataRow, DataValue};
3304
3305 fn create_test_table() -> Arc<DataTable> {
3306 let mut table = DataTable::new("test");
3307
3308 table.add_column(DataColumn::new("id"));
3310 table.add_column(DataColumn::new("name"));
3311 table.add_column(DataColumn::new("age"));
3312
3313 table
3315 .add_row(DataRow::new(vec![
3316 DataValue::Integer(1),
3317 DataValue::String("Alice".to_string()),
3318 DataValue::Integer(30),
3319 ]))
3320 .unwrap();
3321
3322 table
3323 .add_row(DataRow::new(vec![
3324 DataValue::Integer(2),
3325 DataValue::String("Bob".to_string()),
3326 DataValue::Integer(25),
3327 ]))
3328 .unwrap();
3329
3330 table
3331 .add_row(DataRow::new(vec![
3332 DataValue::Integer(3),
3333 DataValue::String("Charlie".to_string()),
3334 DataValue::Integer(35),
3335 ]))
3336 .unwrap();
3337
3338 Arc::new(table)
3339 }
3340
3341 #[test]
3342 fn test_select_all() {
3343 let table = create_test_table();
3344 let engine = QueryEngine::new();
3345
3346 let view = engine
3347 .execute(table.clone(), "SELECT * FROM users")
3348 .unwrap();
3349 assert_eq!(view.row_count(), 3);
3350 assert_eq!(view.column_count(), 3);
3351 }
3352
3353 #[test]
3354 fn test_select_columns() {
3355 let table = create_test_table();
3356 let engine = QueryEngine::new();
3357
3358 let view = engine
3359 .execute(table.clone(), "SELECT name, age FROM users")
3360 .unwrap();
3361 assert_eq!(view.row_count(), 3);
3362 assert_eq!(view.column_count(), 2);
3363 }
3364
3365 #[test]
3366 fn test_select_with_limit() {
3367 let table = create_test_table();
3368 let engine = QueryEngine::new();
3369
3370 let view = engine
3371 .execute(table.clone(), "SELECT * FROM users LIMIT 2")
3372 .unwrap();
3373 assert_eq!(view.row_count(), 2);
3374 }
3375
3376 #[test]
3377 fn test_type_coercion_contains() {
3378 let _ = tracing_subscriber::fmt()
3380 .with_max_level(tracing::Level::DEBUG)
3381 .try_init();
3382
3383 let mut table = DataTable::new("test");
3384 table.add_column(DataColumn::new("id"));
3385 table.add_column(DataColumn::new("status"));
3386 table.add_column(DataColumn::new("price"));
3387
3388 table
3390 .add_row(DataRow::new(vec![
3391 DataValue::Integer(1),
3392 DataValue::String("Pending".to_string()),
3393 DataValue::Float(99.99),
3394 ]))
3395 .unwrap();
3396
3397 table
3398 .add_row(DataRow::new(vec![
3399 DataValue::Integer(2),
3400 DataValue::String("Confirmed".to_string()),
3401 DataValue::Float(150.50),
3402 ]))
3403 .unwrap();
3404
3405 table
3406 .add_row(DataRow::new(vec![
3407 DataValue::Integer(3),
3408 DataValue::String("Pending".to_string()),
3409 DataValue::Float(75.00),
3410 ]))
3411 .unwrap();
3412
3413 let table = Arc::new(table);
3414 let engine = QueryEngine::new();
3415
3416 println!("\n=== Testing WHERE clause with Contains ===");
3417 println!("Table has {} rows", table.row_count());
3418 for i in 0..table.row_count() {
3419 let status = table.get_value(i, 1);
3420 println!("Row {i}: status = {status:?}");
3421 }
3422
3423 println!("\n--- Test 1: status.Contains('pend') ---");
3425 let result = engine.execute(
3426 table.clone(),
3427 "SELECT * FROM test WHERE status.Contains('pend')",
3428 );
3429 match result {
3430 Ok(view) => {
3431 println!("SUCCESS: Found {} matching rows", view.row_count());
3432 assert_eq!(view.row_count(), 2); }
3434 Err(e) => {
3435 panic!("Query failed: {e}");
3436 }
3437 }
3438
3439 println!("\n--- Test 2: price.Contains('9') ---");
3441 let result = engine.execute(
3442 table.clone(),
3443 "SELECT * FROM test WHERE price.Contains('9')",
3444 );
3445 match result {
3446 Ok(view) => {
3447 println!(
3448 "SUCCESS: Found {} matching rows with price containing '9'",
3449 view.row_count()
3450 );
3451 assert!(view.row_count() >= 1);
3453 }
3454 Err(e) => {
3455 panic!("Numeric coercion query failed: {e}");
3456 }
3457 }
3458
3459 println!("\n=== All tests passed! ===");
3460 }
3461
3462 #[test]
3463 fn test_not_in_clause() {
3464 let _ = tracing_subscriber::fmt()
3466 .with_max_level(tracing::Level::DEBUG)
3467 .try_init();
3468
3469 let mut table = DataTable::new("test");
3470 table.add_column(DataColumn::new("id"));
3471 table.add_column(DataColumn::new("country"));
3472
3473 table
3475 .add_row(DataRow::new(vec![
3476 DataValue::Integer(1),
3477 DataValue::String("CA".to_string()),
3478 ]))
3479 .unwrap();
3480
3481 table
3482 .add_row(DataRow::new(vec![
3483 DataValue::Integer(2),
3484 DataValue::String("US".to_string()),
3485 ]))
3486 .unwrap();
3487
3488 table
3489 .add_row(DataRow::new(vec![
3490 DataValue::Integer(3),
3491 DataValue::String("UK".to_string()),
3492 ]))
3493 .unwrap();
3494
3495 let table = Arc::new(table);
3496 let engine = QueryEngine::new();
3497
3498 println!("\n=== Testing NOT IN clause ===");
3499 println!("Table has {} rows", table.row_count());
3500 for i in 0..table.row_count() {
3501 let country = table.get_value(i, 1);
3502 println!("Row {i}: country = {country:?}");
3503 }
3504
3505 println!("\n--- Test: country NOT IN ('CA') ---");
3507 let result = engine.execute(
3508 table.clone(),
3509 "SELECT * FROM test WHERE country NOT IN ('CA')",
3510 );
3511 match result {
3512 Ok(view) => {
3513 println!("SUCCESS: Found {} rows not in ('CA')", view.row_count());
3514 assert_eq!(view.row_count(), 2); }
3516 Err(e) => {
3517 panic!("NOT IN query failed: {e}");
3518 }
3519 }
3520
3521 println!("\n=== NOT IN test complete! ===");
3522 }
3523
3524 #[test]
3525 fn test_case_insensitive_in_and_not_in() {
3526 let _ = tracing_subscriber::fmt()
3528 .with_max_level(tracing::Level::DEBUG)
3529 .try_init();
3530
3531 let mut table = DataTable::new("test");
3532 table.add_column(DataColumn::new("id"));
3533 table.add_column(DataColumn::new("country"));
3534
3535 table
3537 .add_row(DataRow::new(vec![
3538 DataValue::Integer(1),
3539 DataValue::String("CA".to_string()), ]))
3541 .unwrap();
3542
3543 table
3544 .add_row(DataRow::new(vec![
3545 DataValue::Integer(2),
3546 DataValue::String("us".to_string()), ]))
3548 .unwrap();
3549
3550 table
3551 .add_row(DataRow::new(vec![
3552 DataValue::Integer(3),
3553 DataValue::String("UK".to_string()), ]))
3555 .unwrap();
3556
3557 let table = Arc::new(table);
3558
3559 println!("\n=== Testing Case-Insensitive IN clause ===");
3560 println!("Table has {} rows", table.row_count());
3561 for i in 0..table.row_count() {
3562 let country = table.get_value(i, 1);
3563 println!("Row {i}: country = {country:?}");
3564 }
3565
3566 println!("\n--- Test: country IN ('ca') with case_insensitive=true ---");
3568 let engine = QueryEngine::with_case_insensitive(true);
3569 let result = engine.execute(table.clone(), "SELECT * FROM test WHERE country IN ('ca')");
3570 match result {
3571 Ok(view) => {
3572 println!(
3573 "SUCCESS: Found {} rows matching 'ca' (case-insensitive)",
3574 view.row_count()
3575 );
3576 assert_eq!(view.row_count(), 1); }
3578 Err(e) => {
3579 panic!("Case-insensitive IN query failed: {e}");
3580 }
3581 }
3582
3583 println!("\n--- Test: country NOT IN ('ca') with case_insensitive=true ---");
3585 let result = engine.execute(
3586 table.clone(),
3587 "SELECT * FROM test WHERE country NOT IN ('ca')",
3588 );
3589 match result {
3590 Ok(view) => {
3591 println!(
3592 "SUCCESS: Found {} rows not matching 'ca' (case-insensitive)",
3593 view.row_count()
3594 );
3595 assert_eq!(view.row_count(), 2); }
3597 Err(e) => {
3598 panic!("Case-insensitive NOT IN query failed: {e}");
3599 }
3600 }
3601
3602 println!("\n--- Test: country IN ('ca') with case_insensitive=false ---");
3604 let engine_case_sensitive = QueryEngine::new(); let result = engine_case_sensitive
3606 .execute(table.clone(), "SELECT * FROM test WHERE country IN ('ca')");
3607 match result {
3608 Ok(view) => {
3609 println!(
3610 "SUCCESS: Found {} rows matching 'ca' (case-sensitive)",
3611 view.row_count()
3612 );
3613 assert_eq!(view.row_count(), 0); }
3615 Err(e) => {
3616 panic!("Case-sensitive IN query failed: {e}");
3617 }
3618 }
3619
3620 println!("\n=== Case-insensitive IN/NOT IN test complete! ===");
3621 }
3622
3623 #[test]
3624 #[ignore = "Parentheses in WHERE clause not yet implemented"]
3625 fn test_parentheses_in_where_clause() {
3626 let _ = tracing_subscriber::fmt()
3628 .with_max_level(tracing::Level::DEBUG)
3629 .try_init();
3630
3631 let mut table = DataTable::new("test");
3632 table.add_column(DataColumn::new("id"));
3633 table.add_column(DataColumn::new("status"));
3634 table.add_column(DataColumn::new("priority"));
3635
3636 table
3638 .add_row(DataRow::new(vec![
3639 DataValue::Integer(1),
3640 DataValue::String("Pending".to_string()),
3641 DataValue::String("High".to_string()),
3642 ]))
3643 .unwrap();
3644
3645 table
3646 .add_row(DataRow::new(vec![
3647 DataValue::Integer(2),
3648 DataValue::String("Complete".to_string()),
3649 DataValue::String("High".to_string()),
3650 ]))
3651 .unwrap();
3652
3653 table
3654 .add_row(DataRow::new(vec![
3655 DataValue::Integer(3),
3656 DataValue::String("Pending".to_string()),
3657 DataValue::String("Low".to_string()),
3658 ]))
3659 .unwrap();
3660
3661 table
3662 .add_row(DataRow::new(vec![
3663 DataValue::Integer(4),
3664 DataValue::String("Complete".to_string()),
3665 DataValue::String("Low".to_string()),
3666 ]))
3667 .unwrap();
3668
3669 let table = Arc::new(table);
3670 let engine = QueryEngine::new();
3671
3672 println!("\n=== Testing Parentheses in WHERE clause ===");
3673 println!("Table has {} rows", table.row_count());
3674 for i in 0..table.row_count() {
3675 let status = table.get_value(i, 1);
3676 let priority = table.get_value(i, 2);
3677 println!("Row {i}: status = {status:?}, priority = {priority:?}");
3678 }
3679
3680 println!("\n--- Test: (status = 'Pending' AND priority = 'High') OR (status = 'Complete' AND priority = 'Low') ---");
3682 let result = engine.execute(
3683 table.clone(),
3684 "SELECT * FROM test WHERE (status = 'Pending' AND priority = 'High') OR (status = 'Complete' AND priority = 'Low')",
3685 );
3686 match result {
3687 Ok(view) => {
3688 println!(
3689 "SUCCESS: Found {} rows with parenthetical logic",
3690 view.row_count()
3691 );
3692 assert_eq!(view.row_count(), 2); }
3694 Err(e) => {
3695 panic!("Parentheses query failed: {e}");
3696 }
3697 }
3698
3699 println!("\n=== Parentheses test complete! ===");
3700 }
3701
3702 #[test]
3703 #[ignore = "Numeric type coercion needs fixing"]
3704 fn test_numeric_type_coercion() {
3705 let _ = tracing_subscriber::fmt()
3707 .with_max_level(tracing::Level::DEBUG)
3708 .try_init();
3709
3710 let mut table = DataTable::new("test");
3711 table.add_column(DataColumn::new("id"));
3712 table.add_column(DataColumn::new("price"));
3713 table.add_column(DataColumn::new("quantity"));
3714
3715 table
3717 .add_row(DataRow::new(vec![
3718 DataValue::Integer(1),
3719 DataValue::Float(99.50), DataValue::Integer(100),
3721 ]))
3722 .unwrap();
3723
3724 table
3725 .add_row(DataRow::new(vec![
3726 DataValue::Integer(2),
3727 DataValue::Float(150.0), DataValue::Integer(200),
3729 ]))
3730 .unwrap();
3731
3732 table
3733 .add_row(DataRow::new(vec![
3734 DataValue::Integer(3),
3735 DataValue::Integer(75), DataValue::Integer(50),
3737 ]))
3738 .unwrap();
3739
3740 let table = Arc::new(table);
3741 let engine = QueryEngine::new();
3742
3743 println!("\n=== Testing Numeric Type Coercion ===");
3744 println!("Table has {} rows", table.row_count());
3745 for i in 0..table.row_count() {
3746 let price = table.get_value(i, 1);
3747 let quantity = table.get_value(i, 2);
3748 println!("Row {i}: price = {price:?}, quantity = {quantity:?}");
3749 }
3750
3751 println!("\n--- Test: price.Contains('.') ---");
3753 let result = engine.execute(
3754 table.clone(),
3755 "SELECT * FROM test WHERE price.Contains('.')",
3756 );
3757 match result {
3758 Ok(view) => {
3759 println!(
3760 "SUCCESS: Found {} rows with decimal points in price",
3761 view.row_count()
3762 );
3763 assert_eq!(view.row_count(), 2); }
3765 Err(e) => {
3766 panic!("Numeric Contains query failed: {e}");
3767 }
3768 }
3769
3770 println!("\n--- Test: quantity.Contains('0') ---");
3772 let result = engine.execute(
3773 table.clone(),
3774 "SELECT * FROM test WHERE quantity.Contains('0')",
3775 );
3776 match result {
3777 Ok(view) => {
3778 println!(
3779 "SUCCESS: Found {} rows with '0' in quantity",
3780 view.row_count()
3781 );
3782 assert_eq!(view.row_count(), 2); }
3784 Err(e) => {
3785 panic!("Integer Contains query failed: {e}");
3786 }
3787 }
3788
3789 println!("\n=== Numeric type coercion test complete! ===");
3790 }
3791
3792 #[test]
3793 fn test_datetime_comparisons() {
3794 let _ = tracing_subscriber::fmt()
3796 .with_max_level(tracing::Level::DEBUG)
3797 .try_init();
3798
3799 let mut table = DataTable::new("test");
3800 table.add_column(DataColumn::new("id"));
3801 table.add_column(DataColumn::new("created_date"));
3802
3803 table
3805 .add_row(DataRow::new(vec![
3806 DataValue::Integer(1),
3807 DataValue::String("2024-12-15".to_string()),
3808 ]))
3809 .unwrap();
3810
3811 table
3812 .add_row(DataRow::new(vec![
3813 DataValue::Integer(2),
3814 DataValue::String("2025-01-15".to_string()),
3815 ]))
3816 .unwrap();
3817
3818 table
3819 .add_row(DataRow::new(vec![
3820 DataValue::Integer(3),
3821 DataValue::String("2025-02-15".to_string()),
3822 ]))
3823 .unwrap();
3824
3825 let table = Arc::new(table);
3826 let engine = QueryEngine::new();
3827
3828 println!("\n=== Testing DateTime Comparisons ===");
3829 println!("Table has {} rows", table.row_count());
3830 for i in 0..table.row_count() {
3831 let date = table.get_value(i, 1);
3832 println!("Row {i}: created_date = {date:?}");
3833 }
3834
3835 println!("\n--- Test: created_date > DateTime(2025,1,1) ---");
3837 let result = engine.execute(
3838 table.clone(),
3839 "SELECT * FROM test WHERE created_date > DateTime(2025,1,1)",
3840 );
3841 match result {
3842 Ok(view) => {
3843 println!("SUCCESS: Found {} rows after 2025-01-01", view.row_count());
3844 assert_eq!(view.row_count(), 2); }
3846 Err(e) => {
3847 panic!("DateTime comparison query failed: {e}");
3848 }
3849 }
3850
3851 println!("\n=== DateTime comparison test complete! ===");
3852 }
3853
3854 #[test]
3855 fn test_not_with_method_calls() {
3856 let _ = tracing_subscriber::fmt()
3858 .with_max_level(tracing::Level::DEBUG)
3859 .try_init();
3860
3861 let mut table = DataTable::new("test");
3862 table.add_column(DataColumn::new("id"));
3863 table.add_column(DataColumn::new("status"));
3864
3865 table
3867 .add_row(DataRow::new(vec![
3868 DataValue::Integer(1),
3869 DataValue::String("Pending Review".to_string()),
3870 ]))
3871 .unwrap();
3872
3873 table
3874 .add_row(DataRow::new(vec![
3875 DataValue::Integer(2),
3876 DataValue::String("Complete".to_string()),
3877 ]))
3878 .unwrap();
3879
3880 table
3881 .add_row(DataRow::new(vec![
3882 DataValue::Integer(3),
3883 DataValue::String("Pending Approval".to_string()),
3884 ]))
3885 .unwrap();
3886
3887 let table = Arc::new(table);
3888 let engine = QueryEngine::with_case_insensitive(true);
3889
3890 println!("\n=== Testing NOT with Method Calls ===");
3891 println!("Table has {} rows", table.row_count());
3892 for i in 0..table.row_count() {
3893 let status = table.get_value(i, 1);
3894 println!("Row {i}: status = {status:?}");
3895 }
3896
3897 println!("\n--- Test: NOT status.Contains('pend') ---");
3899 let result = engine.execute(
3900 table.clone(),
3901 "SELECT * FROM test WHERE NOT status.Contains('pend')",
3902 );
3903 match result {
3904 Ok(view) => {
3905 println!(
3906 "SUCCESS: Found {} rows NOT containing 'pend'",
3907 view.row_count()
3908 );
3909 assert_eq!(view.row_count(), 1); }
3911 Err(e) => {
3912 panic!("NOT Contains query failed: {e}");
3913 }
3914 }
3915
3916 println!("\n--- Test: NOT status.StartsWith('Pending') ---");
3918 let result = engine.execute(
3919 table.clone(),
3920 "SELECT * FROM test WHERE NOT status.StartsWith('Pending')",
3921 );
3922 match result {
3923 Ok(view) => {
3924 println!(
3925 "SUCCESS: Found {} rows NOT starting with 'Pending'",
3926 view.row_count()
3927 );
3928 assert_eq!(view.row_count(), 1); }
3930 Err(e) => {
3931 panic!("NOT StartsWith query failed: {e}");
3932 }
3933 }
3934
3935 println!("\n=== NOT with method calls test complete! ===");
3936 }
3937
3938 #[test]
3939 #[ignore = "Complex logical expressions with parentheses not yet implemented"]
3940 fn test_complex_logical_expressions() {
3941 let _ = tracing_subscriber::fmt()
3943 .with_max_level(tracing::Level::DEBUG)
3944 .try_init();
3945
3946 let mut table = DataTable::new("test");
3947 table.add_column(DataColumn::new("id"));
3948 table.add_column(DataColumn::new("status"));
3949 table.add_column(DataColumn::new("priority"));
3950 table.add_column(DataColumn::new("assigned"));
3951
3952 table
3954 .add_row(DataRow::new(vec![
3955 DataValue::Integer(1),
3956 DataValue::String("Pending".to_string()),
3957 DataValue::String("High".to_string()),
3958 DataValue::String("John".to_string()),
3959 ]))
3960 .unwrap();
3961
3962 table
3963 .add_row(DataRow::new(vec![
3964 DataValue::Integer(2),
3965 DataValue::String("Complete".to_string()),
3966 DataValue::String("High".to_string()),
3967 DataValue::String("Jane".to_string()),
3968 ]))
3969 .unwrap();
3970
3971 table
3972 .add_row(DataRow::new(vec![
3973 DataValue::Integer(3),
3974 DataValue::String("Pending".to_string()),
3975 DataValue::String("Low".to_string()),
3976 DataValue::String("John".to_string()),
3977 ]))
3978 .unwrap();
3979
3980 table
3981 .add_row(DataRow::new(vec![
3982 DataValue::Integer(4),
3983 DataValue::String("In Progress".to_string()),
3984 DataValue::String("Medium".to_string()),
3985 DataValue::String("Jane".to_string()),
3986 ]))
3987 .unwrap();
3988
3989 let table = Arc::new(table);
3990 let engine = QueryEngine::new();
3991
3992 println!("\n=== Testing Complex Logical Expressions ===");
3993 println!("Table has {} rows", table.row_count());
3994 for i in 0..table.row_count() {
3995 let status = table.get_value(i, 1);
3996 let priority = table.get_value(i, 2);
3997 let assigned = table.get_value(i, 3);
3998 println!(
3999 "Row {i}: status = {status:?}, priority = {priority:?}, assigned = {assigned:?}"
4000 );
4001 }
4002
4003 println!("\n--- Test: status = 'Pending' AND (priority = 'High' OR assigned = 'John') ---");
4005 let result = engine.execute(
4006 table.clone(),
4007 "SELECT * FROM test WHERE status = 'Pending' AND (priority = 'High' OR assigned = 'John')",
4008 );
4009 match result {
4010 Ok(view) => {
4011 println!(
4012 "SUCCESS: Found {} rows with complex logic",
4013 view.row_count()
4014 );
4015 assert_eq!(view.row_count(), 2); }
4017 Err(e) => {
4018 panic!("Complex logic query failed: {e}");
4019 }
4020 }
4021
4022 println!("\n--- Test: NOT (status.Contains('Complete') OR priority = 'Low') ---");
4024 let result = engine.execute(
4025 table.clone(),
4026 "SELECT * FROM test WHERE NOT (status.Contains('Complete') OR priority = 'Low')",
4027 );
4028 match result {
4029 Ok(view) => {
4030 println!(
4031 "SUCCESS: Found {} rows with NOT complex logic",
4032 view.row_count()
4033 );
4034 assert_eq!(view.row_count(), 2); }
4036 Err(e) => {
4037 panic!("NOT complex logic query failed: {e}");
4038 }
4039 }
4040
4041 println!("\n=== Complex logical expressions test complete! ===");
4042 }
4043
4044 #[test]
4045 fn test_mixed_data_types_and_edge_cases() {
4046 let _ = tracing_subscriber::fmt()
4048 .with_max_level(tracing::Level::DEBUG)
4049 .try_init();
4050
4051 let mut table = DataTable::new("test");
4052 table.add_column(DataColumn::new("id"));
4053 table.add_column(DataColumn::new("value"));
4054 table.add_column(DataColumn::new("nullable_field"));
4055
4056 table
4058 .add_row(DataRow::new(vec![
4059 DataValue::Integer(1),
4060 DataValue::String("123.45".to_string()),
4061 DataValue::String("present".to_string()),
4062 ]))
4063 .unwrap();
4064
4065 table
4066 .add_row(DataRow::new(vec![
4067 DataValue::Integer(2),
4068 DataValue::Float(678.90),
4069 DataValue::Null,
4070 ]))
4071 .unwrap();
4072
4073 table
4074 .add_row(DataRow::new(vec![
4075 DataValue::Integer(3),
4076 DataValue::Boolean(true),
4077 DataValue::String("also present".to_string()),
4078 ]))
4079 .unwrap();
4080
4081 table
4082 .add_row(DataRow::new(vec![
4083 DataValue::Integer(4),
4084 DataValue::String("false".to_string()),
4085 DataValue::Null,
4086 ]))
4087 .unwrap();
4088
4089 let table = Arc::new(table);
4090 let engine = QueryEngine::new();
4091
4092 println!("\n=== Testing Mixed Data Types and Edge Cases ===");
4093 println!("Table has {} rows", table.row_count());
4094 for i in 0..table.row_count() {
4095 let value = table.get_value(i, 1);
4096 let nullable = table.get_value(i, 2);
4097 println!("Row {i}: value = {value:?}, nullable_field = {nullable:?}");
4098 }
4099
4100 println!("\n--- Test: value.Contains('true') (boolean to string coercion) ---");
4102 let result = engine.execute(
4103 table.clone(),
4104 "SELECT * FROM test WHERE value.Contains('true')",
4105 );
4106 match result {
4107 Ok(view) => {
4108 println!(
4109 "SUCCESS: Found {} rows with boolean coercion",
4110 view.row_count()
4111 );
4112 assert_eq!(view.row_count(), 1); }
4114 Err(e) => {
4115 panic!("Boolean coercion query failed: {e}");
4116 }
4117 }
4118
4119 println!("\n--- Test: id IN (1, 3) ---");
4121 let result = engine.execute(table.clone(), "SELECT * FROM test WHERE id IN (1, 3)");
4122 match result {
4123 Ok(view) => {
4124 println!("SUCCESS: Found {} rows with IN clause", view.row_count());
4125 assert_eq!(view.row_count(), 2); }
4127 Err(e) => {
4128 panic!("Multiple IN values query failed: {e}");
4129 }
4130 }
4131
4132 println!("\n=== Mixed data types test complete! ===");
4133 }
4134
4135 #[test]
4137 fn test_aggregate_only_single_row() {
4138 let table = create_test_stock_data();
4139 let engine = QueryEngine::new();
4140
4141 let result = engine
4143 .execute(
4144 table.clone(),
4145 "SELECT COUNT(*), MIN(close), MAX(close), AVG(close) FROM stock",
4146 )
4147 .expect("Query should succeed");
4148
4149 assert_eq!(
4150 result.row_count(),
4151 1,
4152 "Aggregate-only query should return exactly 1 row"
4153 );
4154 assert_eq!(result.column_count(), 4, "Should have 4 aggregate columns");
4155
4156 let source = result.source();
4158 let row = source.get_row(0).expect("Should have first row");
4159
4160 assert_eq!(row.values[0], DataValue::Integer(5));
4162
4163 assert_eq!(row.values[1], DataValue::Float(99.5));
4165
4166 assert_eq!(row.values[2], DataValue::Float(105.0));
4168
4169 if let DataValue::Float(avg) = &row.values[3] {
4171 assert!(
4172 (avg - 102.4).abs() < 0.01,
4173 "Average should be approximately 102.4, got {}",
4174 avg
4175 );
4176 } else {
4177 panic!("AVG should return a Float value");
4178 }
4179 }
4180
4181 #[test]
4183 fn test_single_aggregate_single_row() {
4184 let table = create_test_stock_data();
4185 let engine = QueryEngine::new();
4186
4187 let result = engine
4188 .execute(table.clone(), "SELECT COUNT(*) FROM stock")
4189 .expect("Query should succeed");
4190
4191 assert_eq!(
4192 result.row_count(),
4193 1,
4194 "Single aggregate query should return exactly 1 row"
4195 );
4196 assert_eq!(result.column_count(), 1, "Should have 1 column");
4197
4198 let source = result.source();
4199 let row = source.get_row(0).expect("Should have first row");
4200 assert_eq!(row.values[0], DataValue::Integer(5));
4201 }
4202
4203 #[test]
4205 fn test_aggregate_with_where_single_row() {
4206 let table = create_test_stock_data();
4207 let engine = QueryEngine::new();
4208
4209 let result = engine
4211 .execute(
4212 table.clone(),
4213 "SELECT COUNT(*), MIN(close), MAX(close) FROM stock WHERE close >= 103.0",
4214 )
4215 .expect("Query should succeed");
4216
4217 assert_eq!(
4218 result.row_count(),
4219 1,
4220 "Filtered aggregate query should return exactly 1 row"
4221 );
4222 assert_eq!(result.column_count(), 3, "Should have 3 aggregate columns");
4223
4224 let source = result.source();
4225 let row = source.get_row(0).expect("Should have first row");
4226
4227 assert_eq!(row.values[0], DataValue::Integer(2));
4229 assert_eq!(row.values[1], DataValue::Float(103.5)); assert_eq!(row.values[2], DataValue::Float(105.0)); }
4232
4233 #[test]
4234 fn test_not_in_parsing() {
4235 use crate::sql::recursive_parser::Parser;
4236
4237 let query = "SELECT * FROM test WHERE country NOT IN ('CA')";
4238 println!("\n=== Testing NOT IN parsing ===");
4239 println!("Parsing query: {query}");
4240
4241 let mut parser = Parser::new(query);
4242 match parser.parse() {
4243 Ok(statement) => {
4244 println!("Parsed statement: {statement:#?}");
4245 if let Some(where_clause) = statement.where_clause {
4246 println!("WHERE conditions: {:#?}", where_clause.conditions);
4247 if let Some(first_condition) = where_clause.conditions.first() {
4248 println!("First condition expression: {:#?}", first_condition.expr);
4249 }
4250 }
4251 }
4252 Err(e) => {
4253 panic!("Parse error: {e}");
4254 }
4255 }
4256 }
4257
4258 fn create_test_stock_data() -> Arc<DataTable> {
4260 let mut table = DataTable::new("stock");
4261
4262 table.add_column(DataColumn::new("symbol"));
4263 table.add_column(DataColumn::new("close"));
4264 table.add_column(DataColumn::new("volume"));
4265
4266 let test_data = vec![
4268 ("AAPL", 99.5, 1000),
4269 ("AAPL", 101.2, 1500),
4270 ("AAPL", 103.5, 2000),
4271 ("AAPL", 105.0, 1200),
4272 ("AAPL", 102.8, 1800),
4273 ];
4274
4275 for (symbol, close, volume) in test_data {
4276 table
4277 .add_row(DataRow::new(vec![
4278 DataValue::String(symbol.to_string()),
4279 DataValue::Float(close),
4280 DataValue::Integer(volume),
4281 ]))
4282 .expect("Should add row successfully");
4283 }
4284
4285 Arc::new(table)
4286 }
4287}
4288
4289#[cfg(test)]
4290#[path = "query_engine_tests.rs"]
4291mod query_engine_tests;