diff --git a/src/db.rs b/src/db.rs index 249f89eb..f44a2229 100644 --- a/src/db.rs +++ b/src/db.rs @@ -362,7 +362,17 @@ fn default_optimizer_pipeline() -> HepOptimizerPipeline { .after_batch( "Parameterize Mark Apply".to_string(), HepBatchStrategy::once_topdown(), - vec![NormalizationRuleImpl::ParameterizeMarkApply], + vec![ + NormalizationRuleImpl::ParameterizeMarkApply, + NormalizationRuleImpl::ParameterizeInnerJoin, + ], + ) + // Only discard predicates after parameterization has chosen the final + // lookup. A Probe must not inherit residuals from a replaced static range. + .after_batch( + "Eliminate Index Filter".to_string(), + HepBatchStrategy::once_topdown(), + vec![NormalizationRuleImpl::EliminateIndexFilter], ) .after_batch( "Expression Remapper".to_string(), diff --git a/src/execution/dql/mark_apply.rs b/src/execution/dql/mark_apply.rs index cec3927b..5f1a812f 100644 --- a/src/execution/dql/mark_apply.rs +++ b/src/execution/dql/mark_apply.rs @@ -31,13 +31,15 @@ enum QuantifiedPredicateOutcome { Skip, } -pub struct MarkApply { +pub struct MarkApply<'a, T: Transaction + 'a> { op: MarkApplyOperator, right_input_plan: LogicalPlan, left_input: ExecId, + // Retain a streaming inner input across next_tuple calls, not its result rows. + join_input: Option<(Box>, ExecId, Tuple)>, } -impl<'a, T: Transaction + 'a> ReadExecutor<'a, T> for MarkApply { +impl<'a, T: Transaction + 'a> ReadExecutor<'a, T> for MarkApply<'a, T> { type Input = (MarkApplyOperator, LogicalPlan, LogicalPlan); fn into_executor( @@ -52,16 +54,20 @@ impl<'a, T: Transaction + 'a> ReadExecutor<'a, T> for MarkApply { op, right_input_plan: right_input, left_input, + join_input: None, })) } } -impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for MarkApply { +impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for MarkApply<'a, T> { fn next_tuple( &mut self, arena: &mut ExecArena<'a, T>, plan_arena: &mut crate::planner::PlanArena<'a>, ) -> Result<(), DatabaseError> { + if matches!(self.op.kind, MarkApplyKind::InnerJoin) { + return self.next_join_tuple(arena, plan_arena); + } if !arena.next_tuple(self.left_input, plan_arena)? { arena.finish(); return Ok(()); @@ -76,7 +82,44 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for MarkApply { } } -impl MarkApply { +impl<'a, T: Transaction + 'a> MarkApply<'a, T> { + fn next_join_tuple( + &mut self, + arena: &mut ExecArena<'a, T>, + plan_arena: &mut crate::planner::PlanArena<'a>, + ) -> Result<(), DatabaseError> { + loop { + if let Some((inner, root, left)) = &mut self.join_input { + while inner.next_tuple(*root, plan_arena)? { + let right = inner.result_tuple(); + if Self::predicates_matched(self.op.predicates(), left, right, plan_arena)? { + let mut output = left.clone(); + output.pk = output.pk.or_else(|| right.pk.clone()); + output.values.extend(right.values.iter().cloned()); + arena.produce_tuple(output); + return Ok(()); + } + } + } + if !arena.next_tuple(self.left_input, plan_arena)? { + self.join_input = None; + arena.finish(); + return Ok(()); + } + let left: Tuple = arena.result_tuple().clone(); + let value = self.parameterized_probe_value(&left, plan_arena)?; + let mut inner = self + .join_input + .take() + .map(|(inner, _, _)| inner) + .unwrap_or_else(|| Box::new(ExecArena::new())); + inner.reset_for_rebuild(); + inner.init_context(arena.context(), arena.transaction()); + let root = self.build_right_input(&mut inner, plan_arena, value); + self.join_input = Some((inner, root, left)); + } + } + fn runtime_probe_for(&self, param_value: Option) -> Option { self.op.parameterized_probe()?; @@ -96,7 +139,25 @@ impl MarkApply { } } - fn with_right_input<'a, T: Transaction + 'a, R>( + fn build_right_input( + &self, + arena: &mut ExecArena<'a, T>, + plan_arena: &mut crate::planner::PlanArena<'a>, + param_value: Option, + ) -> ExecId { + if let Some(probe) = self.runtime_probe_for(param_value) { + arena.push_runtime_probe(probe); + } + build_read( + arena, + plan_arena, + self.right_input_plan.clone(), + arena.context(), + arena.transaction(), + ) + } + + fn with_right_input( &self, arena: &mut ExecArena<'a, T>, plan_arena: &mut crate::planner::PlanArena<'a>, @@ -107,24 +168,9 @@ impl MarkApply { ExecId, ) -> Result, ) -> Result { - let runtime_probe = self.runtime_probe_for(param_value); let depth_before = arena.runtime_probe_depth(); - if let Some(runtime_probe) = runtime_probe { - arena.push_runtime_probe(runtime_probe); - } - - let cache = arena.context(); - let transaction = arena.transaction(); - let result = { - let right_input = build_read( - arena, - plan_arena, - self.right_input_plan.clone(), - cache, - transaction, - ); - f(arena, plan_arena, right_input) - }; + let right_input = self.build_right_input(arena, plan_arena, param_value); + let result = f(arena, plan_arena, right_input); let depth_after = arena.runtime_probe_depth(); debug_assert!( @@ -153,13 +199,14 @@ impl MarkApply { .transpose() } - fn mark_value<'a, T: Transaction + 'a>( + fn mark_value( &self, arena: &mut ExecArena<'a, T>, plan_arena: &mut crate::planner::PlanArena<'a>, left_tuple: &Tuple, ) -> Result { match self.op.kind { + MarkApplyKind::InnerJoin => unreachable!("inner join streams tuples"), MarkApplyKind::Exists => self.with_right_input( arena, plan_arena, @@ -167,7 +214,12 @@ impl MarkApply { |arena, plan_arena, right_input| { while arena.next_tuple(right_input, plan_arena)? { let right_tuple = arena.result_tuple(); - if self.exists_predicate_matched(left_tuple, right_tuple, plan_arena)? { + if Self::predicates_matched( + self.op.predicates(), + left_tuple, + right_tuple, + plan_arena, + )? { return Ok(DataValue::Boolean(true)); } } @@ -252,7 +304,7 @@ impl MarkApply { } } - fn scan_quantified_right_input<'a, T: Transaction + 'a>( + fn scan_quantified_right_input( &self, arena: &mut ExecArena<'a, T>, plan_arena: &mut crate::planner::PlanArena<'a>, @@ -290,15 +342,15 @@ impl MarkApply { } } - fn exists_predicate_matched( - &self, + fn predicates_matched( + predicates: &[crate::planner::ExprRef], left_tuple: &Tuple, right_tuple: &Tuple, plan_arena: &crate::planner::PlanArena<'_>, ) -> Result { let values = SplitTupleRef::new(left_tuple, right_tuple); - for predicate in self.op.predicates() { + for predicate in predicates { match plan_arena .expression(*predicate) .eval(plan_arena, Some(values))? @@ -459,6 +511,118 @@ mod tests { })) } + #[test] + fn inner_join_apply_emits_all_matches_and_reuses_inner_arena() -> Result<(), DatabaseError> { + let table_arena = crate::planner::TableArenaCell::default(); + let mut plan_arena = crate::planner::PlanArena::new(&table_arena); + let mut left = build_values( + &mut plan_arena, + "left_key", + vec![ + vec![DataValue::Int32(2)], + vec![DataValue::Null], + vec![DataValue::Int32(99)], + vec![DataValue::Int32(2)], + ], + ); + let mut right = build_values_with_schema( + &mut plan_arena, + vec![ + ("right_key", LogicalType::Integer), + ("flag", LogicalType::Boolean), + ], + vec![ + vec![DataValue::Int32(2), DataValue::Boolean(true)], + vec![DataValue::Int32(2), DataValue::Boolean(false)], + vec![DataValue::Int32(2), DataValue::Null], + vec![DataValue::Int32(2), DataValue::Boolean(true)], + vec![DataValue::Null, DataValue::Boolean(true)], + ], + ); + let left_column = left.output_schema(&mut plan_arena)[0]; + let right_schema = right.output_schema(&mut plan_arena).clone(); + let equality = + build_equality_predicate(&mut plan_arena, left_column, 0, right_schema[0], 1)?; + let residual = + plan_arena.alloc_expression(ScalarExpression::column_expr(right_schema[1], 2)); + let probe = plan_arena.alloc_expression(ScalarExpression::column_expr(left_column, 0)); + let mut op = MarkApplyOperator::new_inner_join(vec![equality, residual], probe); + // Values supplies the inner rows directly; there is no IndexScan to consume a probe. + op.set_parameterized_probe(None); + let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; + let transaction = storage.transaction()?; + let cache = crate::execution::empty_context(&table_cache, &view_cache, &meta_cache); + let mut arena = ExecArena::new(); + arena.init_context(cache, &transaction); + let left_input = build_read(&mut arena, &mut plan_arena, left, cache, &transaction); + let mut exec = MarkApply { + op, + right_input_plan: right, + left_input, + join_input: None, + }; + let mut inner_address = None; + for _ in 0..4 { + exec.next_tuple(&mut arena, &mut plan_arena)?; + assert_eq!( + arena.result_tuple().values, + vec![ + DataValue::Int32(2), + DataValue::Int32(2), + DataValue::Boolean(true) + ] + ); + let (inner, _, _) = exec.join_input.as_ref().expect("active inner scan"); + let address = &**inner as *const _; + assert_eq!(*inner_address.get_or_insert(address), address); + assert_eq!( + inner.nodes.len(), + 1, + "inner executors must not accumulate per outer row" + ); + assert_eq!(inner.runtime_probe_depth(), 0); + } + exec.next_tuple(&mut arena, &mut plan_arena)?; + assert!(exec.join_input.is_none(), "all outer rows exhausted"); + Ok(()) + } + + #[test] + fn inner_join_apply_does_not_read_past_the_returned_match() -> Result<(), DatabaseError> { + let table_arena = crate::planner::TableArenaCell::default(); + let mut plan_arena = crate::planner::PlanArena::new(&table_arena); + let mut left = build_values(&mut plan_arena, "left_key", vec![vec![DataValue::Int32(1)]]); + let mut right = build_values_with_schema( + &mut plan_arena, + vec![("flag", LogicalType::Boolean)], + // A later row that cannot be cast must not fail the first fetch. + vec![ + vec![DataValue::Boolean(true)], + vec![DataValue::from("not-a-boolean".to_string())], + ], + ); + let left_column = left.output_schema(&mut plan_arena)[0]; + let right_column = right.output_schema(&mut plan_arena)[0]; + let predicate = plan_arena.alloc_expression(ScalarExpression::column_expr(right_column, 1)); + let probe = plan_arena.alloc_expression(ScalarExpression::column_expr(left_column, 0)); + let mut op = MarkApplyOperator::new_inner_join(vec![predicate], probe); + op.set_parameterized_probe(None); + let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; + let transaction = storage.transaction()?; + let mut executor = execute_input::<_, MarkApply<_>>( + (op, left, right), + crate::execution::empty_context(&table_cache, &view_cache, &meta_cache), + plan_arena, + &transaction, + ); + assert_eq!( + executor.next_tuple()?.unwrap().values, + vec![DataValue::Int32(1), DataValue::Boolean(true)] + ); + assert!(executor.next_tuple().is_err()); + Ok(()) + } + #[test] fn mark_exists_apply_appends_boolean_match_column() -> Result<(), DatabaseError> { let table_arena = crate::planner::TableArenaCell::default(); @@ -480,7 +644,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; - let tuples = try_collect(execute_input::<_, MarkApply>( + let tuples = try_collect(execute_input::<_, MarkApply<_>>( ( MarkApplyOperator::new_exists( build_marker_column(&mut plan_arena), @@ -531,7 +695,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; - let tuples = try_collect(execute_input::<_, MarkApply>( + let tuples = try_collect(execute_input::<_, MarkApply<_>>( ( MarkApplyOperator::new_exists( build_marker_column(&mut plan_arena), @@ -612,10 +776,11 @@ mod tests { &transaction, ); - let exec = MarkApply { + let exec: MarkApply = MarkApply { op, right_input_plan: right, left_input: 0, + join_input: None, }; let left_tuple = Tuple::new(None, vec![DataValue::Int32(2), DataValue::Int32(1)]); @@ -663,10 +828,11 @@ mod tests { &transaction, ); - let exec = MarkApply { + let exec: MarkApply = MarkApply { op, right_input_plan: right, left_input: 0, + join_input: None, }; let left_tuple = Tuple::new(None, vec![DataValue::Int32(2)]); @@ -714,10 +880,11 @@ mod tests { &transaction, ); - let exec = MarkApply { + let exec: MarkApply = MarkApply { op, right_input_plan: right, left_input: 0, + join_input: None, }; let left_tuple = Tuple::new(None, vec![DataValue::Null]); @@ -757,7 +924,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; - let tuples = try_collect(execute_input::<_, MarkApply>( + let tuples = try_collect(execute_input::<_, MarkApply<_>>( ( MarkApplyOperator::new_in(build_marker_column(&mut plan_arena), vec![predicate]), left, @@ -805,7 +972,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; - let tuples = try_collect(execute_input::<_, MarkApply>( + let tuples = try_collect(execute_input::<_, MarkApply<_>>( ( MarkApplyOperator::new_in(build_marker_column(&mut plan_arena), vec![predicate]), left, @@ -873,7 +1040,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; - let tuples = try_collect(execute_input::<_, MarkApply>( + let tuples = try_collect(execute_input::<_, MarkApply<_>>( ( MarkApplyOperator::new_in( build_marker_column(&mut plan_arena), diff --git a/src/execution/mod.rs b/src/execution/mod.rs index b3cc8e2e..3bf60c08 100644 --- a/src/execution/mod.rs +++ b/src/execution/mod.rs @@ -195,7 +195,7 @@ pub(crate) enum ExecNode<'a, T: Transaction + 'a> { IndexScan(IndexScan<'a, T>), Insert(Insert), Limit(Limit), - MarkApply(MarkApply), + MarkApply(MarkApply<'a, T>), NestedLoopJoin(NestedLoopJoin), Projection(Projection), RecursiveCte(RecursiveCte<'a, T>), @@ -310,7 +310,7 @@ impl<'a, T: Transaction + 'a> ExecNode<'a, T> { >::next_tuple(exec, arena, plan_arena) } ExecNode::MarkApply(exec) => { - >::next_tuple(exec, arena, plan_arena) + as ExecutorNode<'a, T>>::next_tuple(exec, arena, plan_arena) } ExecNode::NestedLoopJoin(exec) => { >::next_tuple(exec, arena, plan_arena) @@ -746,7 +746,7 @@ where } Operator::MarkApply(op) => { let (left, right) = childrens.pop_twins(); - >::into_executor( + as ReadExecutor<'a, T>>::into_executor( (op, left, right), arena, plan_arena, diff --git a/src/expression/range_detacher.rs b/src/expression/range_detacher.rs index d1ccf29f..0a65d2ef 100644 --- a/src/expression/range_detacher.rs +++ b/src/expression/range_detacher.rs @@ -174,13 +174,20 @@ impl Range { ) -> Bound { match bound { Bound::Included(v) => Bound::Included(merge_value(tuple, is_upper, v)), - Bound::Excluded(v) => Bound::Excluded(merge_value(tuple, is_upper, v)), + Bound::Excluded(v) => Bound::Excluded(merge_value(tuple, !is_upper, v)), Bound::Unbounded => { if tuple.is_empty() { return Bound::Unbounded; } let values = tuple.iter().map(|v| (*v).clone()).collect_vec(); - Bound::Excluded(DataValue::Tuple(values, is_upper)) + // Excluding a lower equality prefix skips its entire key + // range when storage encodes the exclusive bound. Start at + // the prefix itself; the upper sentinel stays exclusive. + if is_upper { + Bound::Excluded(DataValue::Tuple(values, true)) + } else { + Bound::Included(DataValue::Tuple(values, false)) + } } } } @@ -2258,7 +2265,7 @@ mod test { range, Some(Range::SortedRanges(vec![ Range::Scope { - min: Bound::Excluded(DataValue::Tuple( + min: Bound::Included(DataValue::Tuple( vec![DataValue::Int32(1), DataValue::Null, DataValue::Int32(1),], false )), @@ -2273,7 +2280,7 @@ mod test { )), }, Range::Scope { - min: Bound::Excluded(DataValue::Tuple( + min: Bound::Included(DataValue::Tuple( vec![DataValue::Int32(1), DataValue::Null, DataValue::Int32(2),], false )), @@ -2288,7 +2295,7 @@ mod test { )), }, Range::Scope { - min: Bound::Excluded(DataValue::Tuple( + min: Bound::Included(DataValue::Tuple( vec![ DataValue::Int32(1), DataValue::Int32(1), @@ -2307,7 +2314,7 @@ mod test { )), }, Range::Scope { - min: Bound::Excluded(DataValue::Tuple( + min: Bound::Included(DataValue::Tuple( vec![ DataValue::Int32(1), DataValue::Int32(1), @@ -2326,7 +2333,7 @@ mod test { )), }, Range::Scope { - min: Bound::Excluded(DataValue::Tuple( + min: Bound::Included(DataValue::Tuple( vec![ DataValue::Int32(1), DataValue::Int32(2), @@ -2345,7 +2352,7 @@ mod test { )), }, Range::Scope { - min: Bound::Excluded(DataValue::Tuple( + min: Bound::Included(DataValue::Tuple( vec![ DataValue::Int32(1), DataValue::Int32(2), diff --git a/src/optimizer/heuristic/optimizer.rs b/src/optimizer/heuristic/optimizer.rs index 37189bd7..a95778e9 100644 --- a/src/optimizer/heuristic/optimizer.rs +++ b/src/optimizer/heuristic/optimizer.rs @@ -691,6 +691,11 @@ mod tests { NormalizationRuleImpl::PushPredicateIntoScan, ], ) + .after_batch( + "Eliminate Index Filter".to_string(), + HepBatchStrategy::once_topdown(), + vec![NormalizationRuleImpl::EliminateIndexFilter], + ) .implementations(vec![ ImplementationRuleImpl::Projection, ImplementationRuleImpl::Filter, diff --git a/src/optimizer/rule/normalization/column_pruning.rs b/src/optimizer/rule/normalization/column_pruning.rs index 91123f0c..4eae5a7d 100644 --- a/src/optimizer/rule/normalization/column_pruning.rs +++ b/src/optimizer/rule/normalization/column_pruning.rs @@ -141,8 +141,13 @@ impl ColumnPruning { &mut self, op: &'a crate::planner::operator::mark_apply::MarkApplyOperator, ) -> Result<(), DatabaseError> { - self.referenced_columns - .insert(*op.output_column(), self.arena); + if !matches!( + op.kind, + crate::planner::operator::mark_apply::MarkApplyKind::InnerJoin + ) { + self.referenced_columns + .insert(*op.output_column(), self.arena); + } Ok(()) } diff --git a/src/optimizer/rule/normalization/elimination.rs b/src/optimizer/rule/normalization/elimination.rs index 0517c588..8ff77852 100644 --- a/src/optimizer/rule/normalization/elimination.rs +++ b/src/optimizer/rule/normalization/elimination.rs @@ -339,9 +339,6 @@ pub(crate) fn apply_annotated_post_rules( if EliminateRedundantSort.apply(plan, arena)? { changed = true; } - if EliminateIndexFilter.apply(plan, arena)? { - changed = true; - } if UseStreamAggregate.apply(plan, arena)? { changed = true; } @@ -690,7 +687,17 @@ mod tests { let mut arena = crate::planner::PlanArena::new(&table_arena); let predicate = arena.alloc_expression(ScalarExpression::Constant(DataValue::Boolean(true))); - let mut plan = build_filter_with_selected_index(&mut arena, predicate, None); + let stale_residual = + arena.alloc_expression(ScalarExpression::Constant(DataValue::Boolean(false))); + let mut plan = + build_filter_with_selected_index(&mut arena, predicate, Some(stale_residual)); + // Physical annotation must leave the full predicate available for a later + // parameterization rule, even when a static index initially won. + super::apply_annotated_post_rules(&mut plan, &mut arena)?; + let Operator::Filter(filter) = &plan.operator else { + panic!("filter removed before parameterization") + }; + assert_eq!(filter.predicate, predicate); let Childrens::Only(child) = plan.childrens.as_mut() else { unreachable!("filter should have a scan child"); }; @@ -708,7 +715,10 @@ mod tests { let rule = EliminateIndexFilter; assert!(!rule.apply(&mut plan, &mut arena)?); - assert!(matches!(plan.operator, Operator::Filter(_))); + let Operator::Filter(filter) = &plan.operator else { + panic!("probe cannot discharge the static filter") + }; + assert_eq!(filter.predicate, predicate); Ok(()) } diff --git a/src/optimizer/rule/normalization/mod.rs b/src/optimizer/rule/normalization/mod.rs index cec699b1..b85d4118 100644 --- a/src/optimizer/rule/normalization/mod.rs +++ b/src/optimizer/rule/normalization/mod.rs @@ -47,9 +47,10 @@ mod simplification; mod top_k; pub(crate) use compilation_in_advance::evaluator_bind_current; pub(crate) use elimination::{ - apply_annotated_post_rules, apply_scan_order_hint, OrderHintKind, ScanOrderHint, + apply_annotated_post_rules, apply_scan_order_hint, EliminateIndexFilter, OrderHintKind, + ScanOrderHint, }; -pub(crate) use parameterized_index::ParameterizeMarkApply; +pub(crate) use parameterized_index::{ParameterizeInnerJoin, ParameterizeMarkApply}; pub(crate) use simplification::constant_calculation_current; #[derive(Debug, Copy, Clone)] @@ -76,6 +77,8 @@ pub enum NormalizationRuleImpl { MinMaxToTopK, TopK, ParameterizeMarkApply, + ParameterizeInnerJoin, + EliminateIndexFilter, } #[derive(Debug, Copy, Clone, Eq, PartialEq)] @@ -179,6 +182,8 @@ impl NormalizationRuleImpl { NormalizationRuleImpl::ConstantCalculation => NormalizationRuleRootTag::Any, NormalizationRuleImpl::EvaluatorBind => NormalizationRuleRootTag::Any, NormalizationRuleImpl::MinMaxToTopK => NormalizationRuleRootTag::Aggregate, + NormalizationRuleImpl::EliminateIndexFilter => NormalizationRuleRootTag::Filter, + NormalizationRuleImpl::ParameterizeInnerJoin => NormalizationRuleRootTag::Join, NormalizationRuleImpl::ParameterizeMarkApply => NormalizationRuleRootTag::MarkApply, } } @@ -214,6 +219,10 @@ impl NormalizationRule for NormalizationRuleImpl { NormalizationRuleImpl::EvaluatorBind => EvaluatorBind.apply(plan, arena), NormalizationRuleImpl::MinMaxToTopK => MinMaxToTopK.apply(plan, arena), NormalizationRuleImpl::TopK => TopK.apply(plan, arena), + NormalizationRuleImpl::EliminateIndexFilter => EliminateIndexFilter.apply(plan, arena), + NormalizationRuleImpl::ParameterizeInnerJoin => { + ParameterizeInnerJoin.apply(plan, arena) + } NormalizationRuleImpl::ParameterizeMarkApply => { ParameterizeMarkApply.apply(plan, arena) } diff --git a/src/optimizer/rule/normalization/parameterized_index.rs b/src/optimizer/rule/normalization/parameterized_index.rs index 5bd00392..4ea2eacb 100644 --- a/src/optimizer/rule/normalization/parameterized_index.rs +++ b/src/optimizer/rule/normalization/parameterized_index.rs @@ -14,14 +14,19 @@ use crate::catalog::ColumnRef; use crate::errors::DatabaseError; -use crate::expression::{BinaryOperator, ScalarExpression}; +use crate::expression::visitor_mut::{ExprVisitorMut, PositionShift}; +use crate::expression::{BinaryOperator, ScalarExpression, TypeCast}; use crate::optimizer::core::rule::NormalizationRule; -use crate::planner::operator::mark_apply::{MarkApplyKind, MarkApplyQuantifier}; +use crate::planner::operator::filter::FilterOperator; +use crate::planner::operator::join::{JoinCondition, JoinType}; +use crate::planner::operator::mark_apply::{MarkApplyKind, MarkApplyOperator, MarkApplyQuantifier}; +use crate::planner::operator::project::ProjectOperator; use crate::planner::operator::table_scan::TableScanOperator; -use crate::planner::operator::{Operator, PhysicalOption, PlanImpl}; -use crate::planner::{Childrens, ExprRef, LogicalPlan}; +use crate::planner::operator::{Operator, PhysicalOption, PlanImpl, SortOption}; +use crate::planner::{Childrens, ExprRef, LogicalPlan, PlanArena}; use crate::types::index::{IndexLookup, IndexType}; use crate::types::tuple::Schema; +use crate::types::LogicalType; pub(crate) struct ParameterizeMarkApply; @@ -32,7 +37,9 @@ impl NormalizationRule for ParameterizeMarkApply { arena: &mut crate::planner::PlanArena, ) -> Result { let (op, new_probe) = match (&mut plan.operator, plan.childrens.as_mut()) { - (Operator::MarkApply(op), Childrens::Twins { left, right }) => { + (Operator::MarkApply(op), Childrens::Twins { left, right }) + if !matches!(op.kind, MarkApplyKind::InnerJoin) => + { let probe = find_parameterized_probe( op.kind, op.predicates(), @@ -79,7 +86,7 @@ fn find_parameterized_probe( Ok(None) } } - MarkApplyKind::Quantified(MarkApplyQuantifier::All) => Ok(None), + MarkApplyKind::InnerJoin | MarkApplyKind::Quantified(MarkApplyQuantifier::All) => Ok(None), } } @@ -238,6 +245,192 @@ fn schema_contains_column( .any(|candidate| arena.same_column(*candidate, *column)) } +pub(crate) struct ParameterizeInnerJoin; + +// Collect constant equalities without consuming the filter: it must still check +// all conditions after a static index range is replaced by a runtime probe. +fn constant_keys(expr: ExprRef, keys: &mut Vec<(ExprRef, ExprRef)>, arena: &PlanArena<'_>) { + if let ScalarExpression::Binary { + op, + left_expr, + right_expr, + .. + } = arena.expression(expr) + { + match op { + BinaryOperator::And => { + constant_keys(*left_expr, keys, arena); + constant_keys(*right_expr, keys, arena); + } + BinaryOperator::Eq => { + if matches!(arena.expression(*right_expr), ScalarExpression::Constant(_)) { + keys.push((*right_expr, *left_expr)); + } else if matches!(arena.expression(*left_expr), ScalarExpression::Constant(_)) { + keys.push((*left_expr, *right_expr)); + } + } + _ => {} + } + } +} + +fn parameterize( + plan: &mut LogicalPlan, + mut keys: Vec<(ExprRef, ExprRef)>, + arena: &PlanArena<'_>, +) -> Option> { + match (&mut plan.operator, plan.childrens.as_mut()) { + (Operator::Filter(filter), Childrens::Only(child)) => { + constant_keys(filter.predicate, &mut keys, arena); + parameterize(child, keys, arena) + } + (Operator::TableScan(scan), _) if scan.limit == (None, None) => { + 'indexes: for info in &mut scan.index_infos { + let meta = arena.index(info.meta); + let mut probe = Vec::with_capacity(meta.column_ids.len()); + 'columns: for id in &meta.column_ids { + let Some(candidate) = scan + .columns + .iter() + .find(|column| arena.column(**column).id() == Some(*id)) + else { + continue 'indexes; + }; + for &(value, column) in &keys { + let ScalarExpression::ColumnRef { column, .. } = + arena.expression(column.unpack_alias(arena)) + else { + continue; + }; + if arena.same_column(*candidate, *column) + && value.return_type(arena).as_ref() + == arena.column(*candidate).datatype() + { + probe.push(value); + continue 'columns; + } + } + continue 'indexes; + } + info.lookup = Some(IndexLookup::Probe); + info.residual_predicate = None; + plan.physical_option = Some(PhysicalOption::new( + PlanImpl::IndexScan(Box::new(info.clone())), + info.sort_option.clone(), + )); + return Some(probe); + } + None + } + // LIMIT/aggregation/projection are not row-local filters. Moving a probe + // below them could change which rows the original inner input produces. + _ => None, + } +} + +impl NormalizationRule for ParameterizeInnerJoin { + fn apply( + &self, + plan: &mut LogicalPlan, + arena: &mut PlanArena<'_>, + ) -> Result { + let (Operator::Join(join), Childrens::Twins { left, right }) = + (&mut plan.operator, plan.childrens.as_mut()) + else { + return Ok(false); + }; + if join.join_type != JoinType::Inner || join.force_nested_loop { + return Ok(false); + } + let JoinCondition::On { on, filter } = &mut join.on else { + return Ok(false); + }; + if on.is_empty() { + return Ok(false); + } + let mut keys = on.clone(); + let (project, probe_keys) = if let Some(probe) = parameterize(right, keys.clone(), arena) { + (None, probe) + } else { + for (l, r) in &mut keys { + std::mem::swap(l, r); + } + let Some(probe) = parameterize(left, keys.clone(), arena) else { + return Ok(false); + }; + let left_schema = left.output_schema(arena); + let right_schema = right.output_schema(arena); + let old_left_len = left_schema.len(); + let left_len = right_schema.len(); + // Preserve the original output slots after changing the driving side. + let exprs = left_schema + .iter() + .chain(right_schema.iter()) + .copied() + .enumerate() + .map(|(i, column)| { + let position = if i < old_left_len { + left_len + i + } else { + i - old_left_len + }; + arena.alloc_expression(ScalarExpression::column_expr(column, position)) + }) + .collect(); + std::mem::swap(left, right); + let mut project = LogicalPlan::new( + Operator::Project(ProjectOperator { exprs }), + Childrens::None, + ); + project.physical_option = + Some(PhysicalOption::new(PlanImpl::Project, SortOption::Follow)); + (Some(project), probe) + }; + + let filter = filter.take(); + let (mut left, right) = plan.take().childrens.pop_twins(); + let left_len = left.output_schema(arena).len(); + let probe = if probe_keys.len() == 1 { + probe_keys[0] + } else { + arena.alloc_expression(ScalarExpression::Tuple(probe_keys)) + }; + let mut predicates = Vec::with_capacity(keys.len()); + for (left_expr, mut right_expr) in keys { + PositionShift { + delta: left_len as isize, + } + .visit(&mut right_expr, arena)?; + predicates.push(arena.alloc_expression(ScalarExpression::Binary { + op: BinaryOperator::Eq, + left_expr, + right_expr, + evaluator: None, + ty: LogicalType::Boolean, + })); + } + let apply = LogicalPlan::new( + Operator::MarkApply(MarkApplyOperator::new_inner_join(predicates, probe)), + Childrens::Twins { + left: Box::new(left), + right: Box::new(right), + }, + ); + *plan = if let Some(mut project) = project { + project.childrens = Box::new(Childrens::Only(Box::new(apply))); + project + } else { + apply + }; + if let Some(filter) = filter { + // The original join filter uses the original left/right output slots. + *plan = FilterOperator::build(filter, plan.take(), false); + plan.physical_option = Some(PhysicalOption::new(PlanImpl::Filter, SortOption::Follow)); + } + Ok(true) + } +} + // GRCOV_EXCL_START #[cfg(all(test, not(target_arch = "wasm32")))] mod tests { @@ -382,3 +575,315 @@ mod tests { } } // GRCOV_EXCL_STOP + +#[cfg(test)] +mod inner_join_tests { + use super::*; + use crate::catalog::{ColumnCatalog, ColumnDesc}; + use crate::expression::range_detacher::Range; + use crate::planner::operator::join::JoinOperator; + use crate::planner::operator::mark_apply::MarkApplyKind; + use crate::planner::operator::table_scan::TableScanOperator; + use crate::planner::TableArenaCell; + use crate::types::index::{IndexInfo, IndexMeta, IndexType}; + use crate::types::value::DataValue; + + fn scan( + arena: &mut PlanArena<'_>, + name: &str, + types: &[LogicalType], + indexes: &[&[usize]], + ) -> LogicalPlan { + let columns: Vec<_> = types + .iter() + .enumerate() + .map(|(id, ty)| { + let mut column = ColumnCatalog::new( + format!("c{id}"), + true, + ColumnDesc::new(ty.clone(), None, false, None).unwrap(), + ); + column.set_ref_table(name.into(), id as _, true); + arena.alloc_column(column) + }) + .collect(); + let index_infos = indexes + .iter() + .enumerate() + .map(|(id, keys)| IndexInfo { + meta: arena.alloc_index(IndexMeta { + id: id as _, + column_ids: keys.iter().map(|id| *id as _).collect(), + table_name: name.into(), + pk_ty: LogicalType::Integer, + value_ty: if keys.len() == 1 { + types[keys[0]].clone() + } else { + LogicalType::Tuple(keys.iter().map(|id| types[*id].clone()).collect()) + }, + name: format!("idx{id}"), + ty: if keys.len() == 1 { + IndexType::Normal + } else { + IndexType::Composite + }, + }), + lookup: None, + residual_predicate: None, + sort_option: SortOption::None, + covered_deserializers: None, + cover_mapping: None, + sort_elimination_hint: None, + stream_aggregate_hint: None, + }) + .collect(); + LogicalPlan::new( + Operator::TableScan(TableScanOperator { + table_name: name.into(), + columns, + limit: (None, None), + index_infos, + with_pk: false, + }), + Childrens::None, + ) + } + + fn column(plan: &LogicalPlan, position: usize, arena: &mut PlanArena<'_>) -> ExprRef { + let Operator::TableScan(scan) = &plan.operator else { + panic!("expected scan") + }; + arena.alloc_expression(ScalarExpression::column_expr( + scan.columns[position], + position, + )) + } + + fn binary( + arena: &mut PlanArena<'_>, + op: BinaryOperator, + left_expr: ExprRef, + right_expr: ExprRef, + ) -> ExprRef { + arena.alloc_expression(ScalarExpression::Binary { + op, + left_expr, + right_expr, + evaluator: None, + ty: LogicalType::Boolean, + }) + } + + fn position(expr: ExprRef, arena: &PlanArena<'_>) -> usize { + let ScalarExpression::ColumnRef { position, .. } = arena.expression(expr) else { + panic!("expected column") + }; + *position + } + + fn join( + left: LogicalPlan, + right: LogicalPlan, + on: Vec<(ExprRef, ExprRef)>, + filter: Option, + ) -> LogicalPlan { + JoinOperator::build( + left, + right, + JoinCondition::On { on, filter }, + JoinType::Inner, + false, + ) + } + + #[test] + fn composite_probe_uses_index_order_and_preserves_static_filter() -> Result<(), DatabaseError> { + let tables = TableArenaCell::default(); + let mut arena = PlanArena::new(&tables); + let left = scan( + &mut arena, + "outer", + &[const { LogicalType::Integer }; 2], + &[], + ); + // First index cannot be probed. Second requires constant + two join keys. + let mut right = scan( + &mut arena, + "inner", + &[const { LogicalType::Integer }; 4], + &[&[3], &[0, 2, 1]], + ); + let l0 = column(&left, 0, &mut arena); + let l1 = column(&left, 1, &mut arena); + let r0 = column(&right, 0, &mut arena); + let r1 = column(&right, 1, &mut arena); + let r2 = column(&right, 2, &mut arena); + let r3 = column(&right, 3, &mut arena); + let constant = arena.alloc_expression(ScalarExpression::Constant(DataValue::Int32(7))); + let equality = binary(&mut arena, BinaryOperator::Eq, constant, r0); + let residual = binary(&mut arena, BinaryOperator::Gt, r3, constant); + let predicate = binary(&mut arena, BinaryOperator::And, equality, residual); + let Operator::TableScan(scan) = &mut right.operator else { + unreachable!() + }; + let info = &mut scan.index_infos[1]; + info.lookup = Some(IndexLookup::Static(Range::Eq(DataValue::Int32(7)))); + info.residual_predicate = Some(residual); + let index = info.meta; + right.physical_option = Some(PhysicalOption::new( + PlanImpl::IndexScan(Box::new(info.clone())), + SortOption::None, + )); + let right = FilterOperator::build(predicate, right, false); + let mut plan = join(left, right, vec![(l0, r1), (l1, r2)], None); + + assert!(ParameterizeInnerJoin.apply(&mut plan, &mut arena)?); + let Operator::MarkApply(apply) = &plan.operator else { + panic!("expected apply without projection") + }; + assert_eq!(apply.kind, MarkApplyKind::InnerJoin); + assert_eq!( + arena.expression(*apply.parameterized_probe().unwrap()), + &ScalarExpression::Tuple(vec![constant, l1, l0]) + ); + assert_eq!((position(l0, &arena), position(l1, &arena)), (0, 1)); + assert_eq!((position(r1, &arena), position(r2, &arena)), (3, 4)); + let Childrens::Twins { right, .. } = plan.childrens.as_ref() else { + unreachable!() + }; + let Operator::Filter(filter) = &right.operator else { + panic!("full filter must survive") + }; + assert_eq!(filter.predicate, predicate); + let Childrens::Only(scan) = right.childrens.as_ref() else { + unreachable!() + }; + let PlanImpl::IndexScan(info) = &scan.physical_option.as_ref().unwrap().plan else { + panic!("expected index scan") + }; + assert_eq!(info.meta, index); + assert_eq!(info.lookup, Some(IndexLookup::Probe)); + assert_eq!(info.residual_predicate, None); + assert_eq!(position(r0, &arena), 0); + assert_eq!(position(r3, &arena), 3); + Ok(()) + } + + #[test] + fn swapped_join_restores_unequal_widths_and_filter_positions() -> Result<(), DatabaseError> { + let tables = TableArenaCell::default(); + let mut arena = PlanArena::new(&tables); + let mut left = scan( + &mut arena, + "left", + &[const { LogicalType::Integer }; 3], + &[&[1]], + ); + let mut right = scan(&mut arena, "right", &[LogicalType::Integer], &[]); + let mut schema = left.output_schema(&mut arena).clone(); + schema.extend_from_slice(right.output_schema(&mut arena)); + let l = column(&left, 1, &mut arena); + let r = column(&right, 0, &mut arena); + let filter_left = arena.alloc_expression(ScalarExpression::column_expr(schema[2], 2)); + let filter_right = arena.alloc_expression(ScalarExpression::column_expr(schema[3], 3)); + let predicate = binary(&mut arena, BinaryOperator::Gt, filter_left, filter_right); + let mut plan = join(left, right, vec![(l, r)], Some(predicate)); + + assert!(ParameterizeInnerJoin.apply(&mut plan, &mut arena)?); + assert_eq!(plan.output_schema(&mut arena), &schema); + let Operator::Filter(filter) = &plan.operator else { + panic!("join filter must remain above projection") + }; + assert_eq!(filter.predicate, predicate); + assert_eq!( + ( + position(filter_left, &arena), + position(filter_right, &arena) + ), + (2, 3) + ); + let Childrens::Only(project) = plan.childrens.as_ref() else { + unreachable!() + }; + let Operator::Project(project_op) = &project.operator else { + panic!("expected reorder projection") + }; + assert_eq!( + project_op + .exprs + .iter() + .map(|expr| position(*expr, &arena)) + .collect::>(), + vec![1, 2, 3, 0] + ); + let Childrens::Only(apply_plan) = project.childrens.as_ref() else { + unreachable!() + }; + let Operator::MarkApply(apply) = &apply_plan.operator else { + panic!("expected apply") + }; + assert_eq!(apply.parameterized_probe(), Some(&r)); + assert_eq!((position(r, &arena), position(l, &arena)), (0, 2)); + let ScalarExpression::Binary { + left_expr, + right_expr, + .. + } = arena.expression(apply.predicates()[0]) + else { + unreachable!() + }; + assert_eq!((*left_expr, *right_expr), (r, l)); + let Childrens::Twins { left, right } = apply_plan.childrens.as_ref() else { + unreachable!() + }; + let Operator::TableScan(outer) = &left.operator else { + unreachable!() + }; + let Operator::TableScan(inner) = &right.operator else { + unreachable!() + }; + assert_eq!(outer.table_name.as_ref(), "right"); + assert_eq!(inner.table_name.as_ref(), "left"); + Ok(()) + } + + #[test] + fn rejected_probe_leaves_join_and_expression_positions_unchanged() -> Result<(), DatabaseError> + { + for case in ["missing_key", "type_mismatch", "inner_limit", "outer_join"] { + let tables = TableArenaCell::default(); + let mut arena = PlanArena::new(&tables); + let left = scan(&mut arena, "left", &[LogicalType::Integer], &[]); + let ty = if case == "type_mismatch" { + LogicalType::Bigint + } else { + LogicalType::Integer + }; + let index: &[usize] = if case == "missing_key" { &[0, 1] } else { &[0] }; + let mut right = scan(&mut arena, "right", &[ty, LogicalType::Integer], &[index]); + if case == "inner_limit" { + let Operator::TableScan(scan) = &mut right.operator else { + unreachable!() + }; + scan.limit = (None, Some(1)); + } + let l = column(&left, 0, &mut arena); + let r = column(&right, 0, &mut arena); + let mut plan = join(left, right, vec![(l, r)], None); + if case == "outer_join" { + let Operator::Join(join) = &mut plan.operator else { + unreachable!() + }; + join.join_type = JoinType::LeftOuter; + } + let before = plan.clone(); + assert!( + !ParameterizeInnerJoin.apply(&mut plan, &mut arena)?, + "{case}" + ); + assert_eq!(plan, before, "{case}"); + assert_eq!((position(l, &arena), position(r, &arena)), (0, 0), "{case}"); + } + Ok(()) + } +} diff --git a/src/optimizer/rule/normalization/pushdown_predicates.rs b/src/optimizer/rule/normalization/pushdown_predicates.rs index 104ad37c..637c8267 100644 --- a/src/optimizer/rule/normalization/pushdown_predicates.rs +++ b/src/optimizer/rule/normalization/pushdown_predicates.rs @@ -762,7 +762,7 @@ mod tests { Range::Scope { min: Bound::Excluded(DataValue::Tuple( vec![DataValue::Int32(1), DataValue::Int32(2)], - false, + true, )), max: Bound::Excluded(DataValue::Tuple(vec![DataValue::Int32(1)], true)), } diff --git a/src/orm/mod.rs b/src/orm/mod.rs index 2d8958e5..8b1aae23 100644 --- a/src/orm/mod.rs +++ b/src/orm/mod.rs @@ -4074,9 +4074,9 @@ mod tests { inner_using_plan, concat!( "Projection [orm_unit_users.id] [Project => (Sort Option: Follow)] ", - "Inner Join On orm_unit_users.id = orm_unit_orders.id [HashJoin => (Sort Option: None)] ", + "InnerJoinApply ", "TableScan orm_unit_users -> [orm_unit_users.id] [SeqScan => (Sort Option: None)] ", - "TableScan orm_unit_orders -> [orm_unit_orders.id] [SeqScan => (Sort Option: None)]" + "TableScan orm_unit_orders -> [orm_unit_orders.id] [IndexScan By pk_index => Probe ? => (Sort Option: OrderBy: (orm_unit_orders.id Asc Nulls Last) ignore_prefix_len: 0)]" ), "{inner_using_plan}" ); diff --git a/src/planner/mod.rs b/src/planner/mod.rs index 93a2b979..bcadaafd 100644 --- a/src/planner/mod.rs +++ b/src/planner/mod.rs @@ -272,6 +272,16 @@ impl LogicalPlan { } _ => unreachable!(), }, + Operator::MarkApply(op) + if matches!(op.kind, operator::mark_apply::MarkApplyKind::InnerJoin) => + { + let Childrens::Twins { left, right } = childrens else { + unreachable!("inner apply requires two inputs") + }; + let mut schema = left.output_schema(arena).clone(); + schema.extend_from_slice(right.output_schema(arena)); + schema + } Operator::MarkApply(op) => { let mut schema = match childrens { Childrens::Only(left) => left.output_schema(arena).clone(), diff --git a/src/planner/operator/mark_apply.rs b/src/planner/operator/mark_apply.rs index 21e61ddf..32e056ac 100644 --- a/src/planner/operator/mark_apply.rs +++ b/src/planner/operator/mark_apply.rs @@ -29,13 +29,14 @@ pub enum MarkApplyQuantifier { pub enum MarkApplyKind { Exists, Quantified(MarkApplyQuantifier), + InnerJoin, } #[derive(Debug, PartialEq, Eq, Clone, Hash, ReferenceSerialization)] pub struct MarkApplyOperator { pub kind: MarkApplyKind, pub predicates: Vec, - output_column: ColumnRef, + output_column: Option, pub parameterized_probe: Option, } @@ -44,7 +45,7 @@ impl MarkApplyOperator { Self { kind: MarkApplyKind::Exists, predicates, - output_column, + output_column: Some(output_column), parameterized_probe: None, } } @@ -76,7 +77,7 @@ impl MarkApplyOperator { Self { kind: MarkApplyKind::Quantified(quantifier), predicates, - output_column, + output_column: Some(output_column), parameterized_probe: None, } } @@ -125,7 +126,18 @@ impl MarkApplyOperator { } pub fn output_column(&self) -> &ColumnRef { - &self.output_column + self.output_column + .as_ref() + .expect("marker apply output column") + } + + pub(crate) fn new_inner_join(predicates: Vec, probe: ExprRef) -> Self { + Self { + kind: MarkApplyKind::InnerJoin, + predicates, + output_column: None, + parameterized_probe: Some(probe), + } } pub fn parameterized_probe(&self) -> Option<&ExprRef> { @@ -140,6 +152,7 @@ impl MarkApplyOperator { impl fmt::Display for MarkApplyOperator { fn fmt(&self, f: &mut Formatter) -> fmt::Result { match self.kind { + MarkApplyKind::InnerJoin => write!(f, "InnerJoinApply"), MarkApplyKind::Exists => write!(f, "MarkExistsApply"), MarkApplyKind::Quantified(MarkApplyQuantifier::Any) => write!(f, "MarkAnyApply"), MarkApplyKind::Quantified(MarkApplyQuantifier::All) => write!(f, "MarkAllApply"), diff --git a/tests/slt/join.slt b/tests/slt/join.slt index b4aafd61..5496a260 100644 --- a/tests/slt/join.slt +++ b/tests/slt/join.slt @@ -205,3 +205,123 @@ select * from a as aliased(only_one) # Every USING column must exist on both sides. statement error select * from a join b using (missing) + + +# Parameterized inner join +statement ok +create table lookup_outer(id int primary key, k int); + +statement ok +create table lookup_inner(id int primary key, k int, v int); + +statement ok +create index lookup_inner_k on lookup_inner(k); + +statement ok +insert into lookup_outer values (1,10),(2,20),(3,null),(4,99),(5,10); + +statement ok +insert into lookup_inner values (1,10,7),(2,10,9),(3,20,3),(4,null,8); + +# A plain SeqScan outer input is eligible; no size/limit heuristic is required. +query T +explain select lookup_outer.id,lookup_inner.id from lookup_outer join lookup_inner on lookup_outer.k=lookup_inner.k; +---- +Projection [lookup_outer.id, lookup_inner.id] [Project => (Sort Option: Follow)] InnerJoinApply TableScan lookup_outer -> [lookup_outer.id, lookup_outer.k] [SeqScan => (Sort Option: None)] TableScan lookup_inner -> [lookup_inner.id, lookup_inner.k] [IndexScan By lookup_inner_k => Probe ? => (Sort Option: OrderBy: (lookup_inner.k Asc Nulls Last) ignore_prefix_len: 0)] + +# Resume the same non-unique probe across calls, then rebuild after NULL and +# missing keys. In particular, NULL = NULL must not produce the inner NULL row. +query II +select lookup_outer.id,lookup_inner.id from lookup_outer +join lookup_inner on lookup_outer.k=lookup_inner.k +order by lookup_outer.id,lookup_inner.id; +---- +1 1 +1 2 +2 3 +5 1 +5 2 + +# All matches are emitted, the non-key filter remains, and NULL never matches. +query III +select o.id,i.id,i.v from (select * from lookup_outer limit 5) o +join lookup_inner i on o.k=i.k where i.v>5 order by o.id,i.id; +---- +1 1 7 +1 2 9 +5 1 7 +5 2 9 + +query III +select o.id,i.id,i.v from (select * from lookup_outer limit 5) o +join lookup_inner i on o.k=i.k and o.ido.id order by i.id,o.id; +---- +2 9 1 +3 3 2 + +query III +select o.id,i.id,i.v from (select * from lookup_outer limit 5) o +join lookup_inner i on o.k=i.k where i.v>5 limit 1; +---- +1 1 7 + +statement ok +create table lookup_composite(w int, item int, qty int, primary key(w,item)); + +statement ok +insert into lookup_composite values (1,10,7),(1,20,9),(2,10,1); + +# A composite probe containing NULL also fails the retained join equality. +query II +select lookup_outer.id,lookup_composite.item from lookup_outer +join lookup_composite on lookup_outer.k=lookup_composite.item +where lookup_composite.w=1 and lookup_composite.qty<8 order by lookup_outer.id; +---- +1 10 +5 10 + +# Constant leading key + outer item key forms one composite point probe. +query III +select c.w,c.item,o.id from lookup_composite c +join (select * from lookup_outer limit 5) o on c.item=o.k +where c.w=1 and c.qty<8 order by o.id; +---- +1 10 1 +1 10 5 + +# Inner LIMIT is a semantic boundary and must not be moved inside the probe. +query II +select o.id,i.id from (select * from lookup_outer limit 5) o +join (select * from lookup_inner order by id limit 1) i on o.k=i.k order by o.id; +---- +1 1 +5 1 + +query II +select o.id,i.id from (select * from lookup_outer limit 5) o +left join lookup_inner i on o.k=i.k order by o.id,i.id; +---- +1 1 +1 2 +2 3 +3 null +4 null +5 1 +5 2 + +statement ok +drop table lookup_composite; + +statement ok +drop table lookup_inner; + +statement ok +drop table lookup_outer; diff --git a/tests/slt/where_by_index.slt b/tests/slt/where_by_index.slt index 860b15d2..902a2ae4 100644 --- a/tests/slt/where_by_index.slt +++ b/tests/slt/where_by_index.slt @@ -303,3 +303,147 @@ select c2, c3 from t_cover where c1 = 2; statement ok drop table t_cover; + +# Composite index range bounds +# Composite equality-prefix bounds must include the complete matching prefix +# and distinguish < / <= / > / >= when trailing index columns are present. +statement ok +create table stock_bounds(w smallint, item int, qty smallint, payload int, primary key(w,item)); + +statement ok +insert into stock_bounds values (1,1,1,10),(1,2,2,20),(1,3,3,30),(1,4,4,40),(1,5,5,50),(1,6,6,60),(1,7,7,70),(1,8,8,80),(1,9,9,90),(1,10,10,100),(1,11,11,110),(1,12,12,120),(1,13,13,130),(1,14,14,140),(1,15,15,150),(1,16,16,160),(1,17,17,170),(1,18,18,180),(1,19,19,190),(1,20,20,200),(1,21,21,210),(1,22,22,220),(1,23,23,230),(1,24,24,240),(1,25,25,250),(1,26,26,260),(1,27,27,270),(1,28,28,280),(1,29,29,290),(1,30,30,300),(1,31,31,310),(1,32,32,320),(1,33,33,330),(1,34,34,340),(1,35,35,350),(1,36,36,360),(1,37,37,370),(1,38,38,380),(1,39,39,390),(1,40,40,400),(1,41,41,410),(1,42,42,420),(1,43,43,430),(1,44,44,440),(1,45,45,450),(1,46,46,460),(1,47,47,470),(1,48,48,480),(1,49,49,490),(1,50,50,500),(1,51,51,510),(1,52,52,520),(1,53,53,530),(1,54,54,540),(1,55,55,550),(1,56,56,560),(1,57,57,570),(1,58,58,580),(1,59,59,590),(1,60,60,600),(1,61,61,610),(1,62,62,620),(1,63,63,630),(1,64,64,640),(1,65,65,650),(1,66,66,660),(1,67,67,670),(1,68,68,680),(1,69,69,690),(1,70,70,700),(1,71,71,710),(1,72,72,720),(1,73,73,730),(1,74,74,740),(1,75,75,750),(1,76,76,760),(1,77,77,770),(1,78,78,780),(1,79,79,790),(1,80,80,800),(1,81,81,810),(1,82,82,820),(1,83,83,830),(1,84,84,840),(1,85,85,850),(1,86,86,860),(1,87,87,870),(1,88,88,880),(1,89,89,890),(1,90,90,900),(1,91,91,910),(1,92,92,920),(1,93,93,930),(1,94,94,940),(1,95,95,950),(1,96,96,960),(1,97,97,970),(1,98,98,980),(1,99,99,990),(1,100,0,1000),(1,101,1,1010),(1,102,2,1020),(1,103,3,1030),(1,104,4,1040),(1,105,5,1050),(1,106,6,1060),(1,107,7,1070),(1,108,8,1080),(1,109,9,1090),(1,110,10,1100),(1,111,11,1110),(1,112,12,1120),(1,113,13,1130),(1,114,14,1140),(1,115,15,1150),(1,116,16,1160),(1,117,17,1170),(1,118,18,1180),(1,119,19,1190),(1,120,20,1200),(1,121,21,1210),(1,122,22,1220),(1,123,23,1230),(1,124,24,1240),(1,125,25,1250),(1,126,26,1260),(1,127,27,1270),(1,128,28,1280),(1,129,29,1290),(1,130,30,1300),(1,131,31,1310),(1,132,32,1320),(1,133,33,1330),(1,134,34,1340),(1,135,35,1350),(1,136,36,1360),(1,137,37,1370),(1,138,38,1380),(1,139,39,1390),(1,140,40,1400),(1,141,41,1410),(1,142,42,1420),(1,143,43,1430),(1,144,44,1440),(1,145,45,1450),(1,146,46,1460),(1,147,47,1470),(1,148,48,1480),(1,149,49,1490),(1,150,50,1500),(1,151,51,1510),(1,152,52,1520),(1,153,53,1530),(1,154,54,1540),(1,155,55,1550),(1,156,56,1560),(1,157,57,1570),(1,158,58,1580),(1,159,59,1590),(1,160,60,1600),(1,161,61,1610),(1,162,62,1620),(1,163,63,1630),(1,164,64,1640),(1,165,65,1650),(1,166,66,1660),(1,167,67,1670),(1,168,68,1680),(1,169,69,1690),(1,170,70,1700),(1,171,71,1710),(1,172,72,1720),(1,173,73,1730),(1,174,74,1740),(1,175,75,1750),(1,176,76,1760),(1,177,77,1770),(1,178,78,1780),(1,179,79,1790),(1,180,80,1800),(1,181,81,1810),(1,182,82,1820),(1,183,83,1830),(1,184,84,1840),(1,185,85,1850),(1,186,86,1860),(1,187,87,1870),(1,188,88,1880),(1,189,89,1890),(1,190,90,1900),(1,191,91,1910),(1,192,92,1920),(1,193,93,1930),(1,194,94,1940),(1,195,95,1950),(1,196,96,1960),(1,197,97,1970),(1,198,98,1980),(1,199,99,1990),(1,200,0,2000); + +statement ok +insert into stock_bounds values (0,1,0,1),(2,1,0,1),(2,2,99,2); + +statement ok +create index stock_bounds_qty on stock_bounds(w,qty,item); + +statement ok +analyze table stock_bounds; + +query III +select item,w,qty from stock_bounds where w=1 and qty < 2 order by item; +---- +1 1 1 +100 1 0 +101 1 1 +200 1 0 + +query IIII +select item,w,qty,payload from stock_bounds where w=1 and qty < 2 order by item; +---- +1 1 1 10 +100 1 0 1000 +101 1 1 1010 +200 1 0 2000 + +query III +select item,w,qty from stock_bounds where w=1 and qty <= 2 order by item; +---- +1 1 1 +2 1 2 +100 1 0 +101 1 1 +102 1 2 +200 1 0 + +query IIII +select item,w,qty,payload from stock_bounds where w=1 and qty <= 2 order by item; +---- +1 1 1 10 +2 1 2 20 +100 1 0 1000 +101 1 1 1010 +102 1 2 1020 +200 1 0 2000 + +query III +select item,w,qty from stock_bounds where w=1 and qty > 97 order by item; +---- +98 1 98 +99 1 99 +198 1 98 +199 1 99 + +query IIII +select item,w,qty,payload from stock_bounds where w=1 and qty > 97 order by item; +---- +98 1 98 980 +99 1 99 990 +198 1 98 1980 +199 1 99 1990 + +query III +select item,w,qty from stock_bounds where w=1 and qty >= 97 order by item; +---- +97 1 97 +98 1 98 +99 1 99 +197 1 97 +198 1 98 +199 1 99 + +query IIII +select item,w,qty,payload from stock_bounds where w=1 and qty >= 97 order by item; +---- +97 1 97 970 +98 1 98 980 +99 1 99 990 +197 1 97 1970 +198 1 98 1980 +199 1 99 1990 + +query III +select item,w,qty from stock_bounds where w=1 and qty = 2 order by item; +---- +2 1 2 +102 1 2 + +query IIII +select item,w,qty,payload from stock_bounds where w=1 and qty = 2 order by item; +---- +2 1 2 20 +102 1 2 1020 + +query III +select item,w,qty from stock_bounds where w=1 and qty < 0 order by item; +---- + +query IIII +select item,w,qty,payload from stock_bounds where w=1 and qty < 0 order by item; +---- + +query III +select item,w,qty from stock_bounds where w=1 and qty > 99 order by item; +---- + +query IIII +select item,w,qty,payload from stock_bounds where w=1 and qty > 99 order by item; +---- + +statement ok +update stock_bounds set qty=1 where w=1 and item=50; + +query III +select item,w,qty from stock_bounds where w=1 and qty<2 order by item; +---- +1 1 1 +50 1 1 +100 1 0 +101 1 1 +200 1 0 + +statement ok +delete from stock_bounds where w=1 and item=101; + +query III +select item,w,qty from stock_bounds where w=1 and qty<2 order by item; +---- +1 1 1 +50 1 1 +100 1 0 +200 1 0 + +statement ok +drop table stock_bounds;