use crate::physical_operator::*;
use crate::processor::QueryProcessor;
use crate::processor::merge_union_chunks;
use akar_binder::bound_statement::BoundExpression;
use akar_common::types::{LogicalTypeID, PhysicalTypeID, Value};
use akar_common::vector::{DataChunk, ValueVector};
use akar_function::registry::FunctionRegistry;
use akar_parser::ast::{Constant, Expression};
use akar_planner::logical_operator::LogicalOperator;
use akar_storage::table::{ColumnDefinition, TableCatalog};
use hashbrown::HashMap;
use std::sync::{Arc, Mutex};
#[cfg(test)]
mod tests {
use super::*;
fn make_scan_op() -> LogicalOperator {
LogicalOperator::ScanNode(akar_planner::logical_operator::LogicalScanNode {
predicate: None,
table_name: "Person".into(),
table_id: 0,
alias: Some("a".into()),
columns: vec![],
cardinality: 0,
fts_query: None,
})
}
fn make_filter_op() -> LogicalOperator {
LogicalOperator::Filter(akar_planner::logical_operator::LogicalFilter {
expression: Expression::Constant(Constant::Bool(true)),
children: vec![],
cardinality: 0,
})
}
fn make_proj_op() -> LogicalOperator {
LogicalOperator::Projection(akar_planner::logical_operator::LogicalProjection {
expressions: vec![BoundExpression {
expression: Expression::Variable("a".into()),
resolved_type: LogicalTypeID::Any,
is_constant: false,
}],
children: vec![],
cardinality: 0,
})
}
fn make_limit_op() -> LogicalOperator {
LogicalOperator::Limit(akar_planner::logical_operator::LogicalLimit {
limit: 10,
offset: 0,
children: vec![],
cardinality: 0,
})
}
fn make_processor_with_person_table() -> QueryProcessor {
let catalog = Arc::new(TableCatalog::new());
{
catalog.create_node_table(
"Person".into(),
vec![
ColumnDefinition {
compression: akar_common::enums::CompressionType::Uncompressed,
name: "name".into(),
logical_type: LogicalTypeID::String,
is_primary_key: true,
},
ColumnDefinition {
compression: akar_common::enums::CompressionType::Uncompressed,
name: "age".into(),
logical_type: LogicalTypeID::Int64,
is_primary_key: false,
},
],
);
let mut table = catalog.get_node_table_by_name_mut("Person").unwrap();
table
.insert_row(vec![Value::String("Alice".into()), Value::Int64(30)])
.unwrap();
table
.insert_row(vec![Value::String("Bob".into()), Value::Int64(25)])
.unwrap();
}
let registry = Arc::new(Mutex::new(FunctionRegistry::new()));
QueryProcessor::with_catalog(
registry,
catalog,
std::sync::Arc::new(akar_common::file_system::VirtualFileSystemRegistry::new()),
)
}
#[test]
fn test_empty_plan() {
let proc = QueryProcessor::new();
let result = proc.execute(&[]).unwrap();
assert_eq!(result.len(), 1);
}
#[test]
fn test_scan_only() {
let proc = make_processor_with_person_table();
let result = proc.execute(&[make_scan_op()]).unwrap();
assert!(!result.is_empty());
assert!(result[0].num_fields() > 0);
assert_eq!(result[0].size, 2); }
#[test]
fn test_scan_filter_projection() {
let proc = make_processor_with_person_table();
let plan = vec![make_scan_op(), make_filter_op(), make_proj_op()];
let result = proc.execute(&plan).unwrap();
assert!(!result.is_empty());
}
#[test]
fn test_scan_filter_limit() {
let proc = make_processor_with_person_table();
let plan = vec![make_scan_op(), make_filter_op(), make_limit_op()];
let result = proc.execute(&plan).unwrap();
assert!(!result.is_empty());
}
#[test]
fn test_filter_true_passthrough() {
let filter = PhysicalFilter::new(Expression::Constant(Constant::Bool(true)));
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
for i in 0..5 {
v.set_i64(i, i as i64);
}
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = filter.execute(input).unwrap();
assert!(!result.is_empty());
assert_eq!(result[0].size, 5); }
#[test]
fn test_filter_false_removes_all() {
let filter = PhysicalFilter::new(Expression::Constant(Constant::Bool(false)));
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
for i in 0..5 {
v.set_i64(i, i as i64);
}
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = filter.execute(input).unwrap();
assert!(result.is_empty());
}
#[test]
fn test_limit() {
let limit = PhysicalLimit { limit: 3, offset: 0 };
let mut v = ValueVector::new(PhysicalTypeID::Int64, 10);
for i in 0..10 {
v.set_i64(i, i as i64);
}
v.resize(10);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = limit.execute(input).unwrap();
assert_eq!(result[0].size, 3);
}
#[test]
fn test_limit_with_offset() {
let limit = PhysicalLimit { limit: 2, offset: 5 };
let mut v = ValueVector::new(PhysicalTypeID::Int64, 10);
for i in 0..10 {
v.set_i64(i, i as i64);
}
v.resize(10);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = limit.execute(input).unwrap();
assert!(!result.is_empty());
}
#[test]
fn test_projection() {
let proj = PhysicalProjection {
column_indices: vec![0],
};
let mut v1 = ValueVector::new(PhysicalTypeID::Int64, 5);
let mut v2 = ValueVector::new(PhysicalTypeID::Int64, 5);
for i in 0..5 {
v1.set_i64(i, i as i64);
v2.set_i64(i, (i * 10) as i64);
}
v1.resize(5);
v2.resize(5);
let input = vec![DataChunk::from_legacy(vec![v1, v2])];
let result = proj.execute(input).unwrap();
assert_eq!(result[0].num_fields(), 1); }
#[test]
fn test_projection_evaluates_function_call_no_input_source() {
let state = Arc::new(Mutex::new(HashMap::new()));
state.lock().unwrap().insert("s".to_string(), 1_i64);
let state_for_fn = state.clone();
let seq_fn: Arc<dyn Fn(&str, bool) -> Result<Value, akar_common::error::ProcessorError> + Send + Sync> =
Arc::new(move |seq_name: &str, is_nextval: bool| {
let mut m = state_for_fn.lock().map_err(|e| format!("Lock error: {e}"))?;
let v = m
.get_mut(seq_name)
.ok_or_else(|| format!("Sequence '{}' not found", seq_name))?;
if is_nextval {
let out = *v;
*v += 1;
Ok(Value::Int64(out))
} else {
Ok(Value::Int64(*v))
}
});
let proc =
QueryProcessor::with_registry(Arc::new(Mutex::new(FunctionRegistry::new()))).with_sequence_fn(seq_fn);
let plan = vec![LogicalOperator::Projection(
akar_planner::logical_operator::LogicalProjection {
expressions: vec![BoundExpression {
expression: Expression::FunctionCall(
"nextval".into(),
vec![Expression::Constant(Constant::String("s".into()))],
),
resolved_type: LogicalTypeID::Int64,
is_constant: false,
}],
children: vec![],
cardinality: 1,
},
)];
let result = proc.execute(&plan).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 1);
assert_eq!(result[0].get_value(0, 0), Some(Value::Int64(1)));
}
#[test]
fn test_projection_sequence_missing_callback_errors() {
let proc = QueryProcessor::with_registry(Arc::new(Mutex::new(FunctionRegistry::new())));
let plan = vec![LogicalOperator::Projection(
akar_planner::logical_operator::LogicalProjection {
expressions: vec![BoundExpression {
expression: Expression::FunctionCall(
"nextval".into(),
vec![Expression::Constant(Constant::String("s".into()))],
),
resolved_type: LogicalTypeID::Int64,
is_constant: false,
}],
children: vec![],
cardinality: 1,
},
)];
let err = proc.execute(&plan).unwrap_err();
assert!(
err.to_string().contains("No sequence callback configured"),
"Unexpected error: {err}"
);
}
#[test]
fn test_order_by_ascending() {
let order = PhysicalOrderBy {
sort_keys: vec![(0, true)],
};
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
let vals = [5, 3, 1, 4, 2];
for i in 0..5 {
v.set_i64(i, vals[i]);
}
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = order.execute(input).unwrap();
assert!(!result.is_empty());
let sorted = result[0].get_i64(0, 0).unwrap();
assert_eq!(sorted, 1); }
#[test]
fn test_order_by_descending() {
let order = PhysicalOrderBy {
sort_keys: vec![(0, false)],
};
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
let vals = [5, 3, 1, 4, 2];
for i in 0..5 {
v.set_i64(i, vals[i]);
}
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = order.execute(input).unwrap();
assert!(!result.is_empty());
let sorted = result[0].get_i64(0, 0).unwrap();
assert_eq!(sorted, 5); }
#[test]
fn test_order_by_empty_input() {
let order = PhysicalOrderBy {
sort_keys: vec![(0, true)],
};
let result = order.execute(vec![]).unwrap();
assert!(result.is_empty());
}
#[test]
fn test_aggregate_count() {
let agg = PhysicalAggregate {
group_by_cols: vec![],
aggregate_functions: vec!["COUNT".into()],
};
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
for i in 0..5 {
v.set_i64(i, i as i64);
}
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = agg.execute(input).unwrap();
assert_eq!(result[0].get_value(0, 0).unwrap(), Value::Int64(5)); }
#[test]
fn test_aggregate_sum() {
let agg = PhysicalAggregate {
group_by_cols: vec![],
aggregate_functions: vec!["SUM".into()],
};
let mut v = ValueVector::new(PhysicalTypeID::Int64, 4);
for i in 0..4 {
v.set_i64(i, (i + 1) as i64);
}
v.resize(4);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = agg.execute(input).unwrap();
assert_eq!(result[0].get_value(0, 0).unwrap(), Value::Int64(10)); }
#[test]
fn test_aggregate_min_max() {
let agg = PhysicalAggregate {
group_by_cols: vec![],
aggregate_functions: vec!["MIN".into(), "MAX".into()],
};
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
let vals = [42, 7, 99, 15, 3];
for i in 0..5 {
v.set_i64(i, vals[i]);
}
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = agg.execute(input).unwrap();
assert_eq!(result[0].get_value(0, 0).unwrap(), Value::Int64(3)); assert_eq!(result[0].get_value(1, 0).unwrap(), Value::Int64(99)); }
#[test]
fn test_aggregate_avg() {
let agg = PhysicalAggregate {
group_by_cols: vec![],
aggregate_functions: vec!["AVG".into()],
};
let mut v = ValueVector::new(PhysicalTypeID::Int64, 4);
for i in 0..4 {
v.set_i64(i, (i + 1) as i64);
}
v.resize(4);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = agg.execute(input).unwrap();
assert_eq!(result[0].get_value(0, 0).unwrap(), Value::Double(2.5)); }
#[test]
fn test_aggregate_empty_input() {
let agg = PhysicalAggregate {
group_by_cols: vec![],
aggregate_functions: vec!["COUNT".into()],
};
let result = agg.execute(vec![]).unwrap();
assert_eq!(result[0].get_value(0, 0).unwrap(), Value::Int64(0)); }
#[test]
fn test_hash_join_basic() {
let join = PhysicalHashJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let mut build = ValueVector::new(PhysicalTypeID::Int64, 3);
for i in 0..3 {
build.set_i64(i, (i + 1) as i64);
}
build.resize(3);
let mut probe = ValueVector::new(PhysicalTypeID::Int64, 3);
probe.set_i64(0, 2);
probe.set_i64(1, 3);
probe.set_i64(2, 4);
probe.resize(3);
let build_chunk = DataChunk::from_legacy(vec![build]);
let probe_chunk = DataChunk::from_legacy(vec![probe]);
let result = join.execute_binary(&[build_chunk], &[probe_chunk]).unwrap();
assert!(!result.is_empty());
}
#[test]
fn test_hash_join_no_match() {
let join = PhysicalHashJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let mut build = ValueVector::new(PhysicalTypeID::Int64, 2);
build.set_i64(0, 1);
build.set_i64(1, 2);
build.resize(2);
let mut probe = ValueVector::new(PhysicalTypeID::Int64, 2);
probe.set_i64(0, 3);
probe.set_i64(1, 4);
probe.resize(2);
let build_chunk = DataChunk::from_legacy(vec![build]);
let probe_chunk = DataChunk::from_legacy(vec![probe]);
let result = join.execute_binary(&[build_chunk], &[probe_chunk]).unwrap();
assert!(result.is_empty()); }
#[test]
fn test_hash_join_empty_build() {
let join = PhysicalHashJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let build = ValueVector::new(PhysicalTypeID::Int64, 0);
let mut probe = ValueVector::new(PhysicalTypeID::Int64, 3);
probe.set_i64(0, 1);
probe.set_i64(1, 2);
probe.set_i64(2, 3);
probe.resize(3);
let build_chunk = DataChunk::from_legacy(vec![build]);
let probe_chunk = DataChunk::from_legacy(vec![probe]);
let result = join.execute_binary(&[build_chunk], &[probe_chunk]).unwrap();
assert!(result.is_empty()); }
#[test]
fn test_hash_join_null_keys_no_match() {
let join = PhysicalHashJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let mut build = ValueVector::new(PhysicalTypeID::Int64, 3);
build.set_i64(0, 1);
build.set_i64(2, 3);
build.resize(3);
let mut probe = ValueVector::new(PhysicalTypeID::Int64, 3);
probe.set_i64(0, 1);
probe.set_i64(1, 3);
probe.resize(3);
let build_chunk = DataChunk::from_legacy(vec![build]);
let probe_chunk = DataChunk::from_legacy(vec![probe]);
let result = join.execute_binary(&[build_chunk], &[probe_chunk]).unwrap();
assert!(!result.is_empty(), "Expected at least one matching row");
}
#[test]
fn test_hash_join_all_null_keys() {
let join = PhysicalHashJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let mut build = ValueVector::new(PhysicalTypeID::Int64, 3);
build.resize(3);
build.set_null(0, true);
build.set_null(1, true);
build.set_null(2, true);
let mut probe = ValueVector::new(PhysicalTypeID::Int64, 3);
probe.resize(3);
probe.set_null(0, true);
probe.set_null(1, true);
probe.set_null(2, true);
let build_chunk = DataChunk::from_legacy(vec![build]);
let probe_chunk = DataChunk::from_legacy(vec![probe]);
let result = join.execute_binary(&[build_chunk], &[probe_chunk]).unwrap();
assert!(result.is_empty());
}
#[test]
fn test_order_by_with_nulls() {
let order = PhysicalOrderBy {
sort_keys: vec![(0, true)],
};
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
v.set_i64(0, 3);
v.set_null(1, true); v.set_i64(2, 1);
v.set_i64(3, 2);
v.set_null(4, true); v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = order.execute(input).unwrap();
assert!(!result.is_empty());
assert_eq!(result[0].get_i64(0, 0).unwrap(), 1);
assert_eq!(result[0].get_i64(0, 1).unwrap(), 2);
assert_eq!(result[0].get_i64(0, 2).unwrap(), 3);
assert!(result[0].is_null(0, 3));
assert!(result[0].is_null(0, 4));
}
#[test]
fn test_limit_zero() {
let limit = PhysicalLimit { limit: 0, offset: 0 };
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
for i in 0..5 {
v.set_i64(i, i as i64);
}
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = limit.execute(input).unwrap();
assert!(result.is_empty());
}
#[test]
fn test_limit_offset_exceeds_total() {
let limit = PhysicalLimit { limit: 5, offset: 100 };
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
for i in 0..5 {
v.set_i64(i, i as i64);
}
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = limit.execute(input).unwrap();
assert!(result.is_empty());
}
#[test]
fn test_aggregate_count_with_nulls() {
let agg = PhysicalAggregate {
group_by_cols: vec![],
aggregate_functions: vec!["COUNT".into()],
};
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
v.set_i64(0, 10);
v.set_null(1, true);
v.set_i64(2, 20);
v.set_null(3, true);
v.set_i64(4, 30);
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = agg.execute(input).unwrap();
assert_eq!(result[0].get_value(0, 0).unwrap(), Value::Int64(3));
}
#[test]
fn test_aggregate_sum_with_nulls() {
let agg = PhysicalAggregate {
group_by_cols: vec![],
aggregate_functions: vec!["SUM".into()],
};
let mut v = ValueVector::new(PhysicalTypeID::Int64, 5);
v.set_i64(0, 10);
v.set_null(1, true);
v.set_i64(2, 20);
v.set_null(3, true);
v.set_i64(4, 30);
v.resize(5);
let input = vec![DataChunk::from_legacy(vec![v])];
let result = agg.execute(input).unwrap();
assert_eq!(result[0].get_value(0, 0).unwrap(), Value::Int64(60));
}
#[test]
fn test_aggregate_group_by_with_nulls() {
let agg = PhysicalAggregate {
group_by_cols: vec![0],
aggregate_functions: vec!["COUNT".into()],
};
let n = 6;
let mut keys = ValueVector::new(PhysicalTypeID::Int64, n);
keys.set_i64(0, 1);
keys.set_i64(1, 1);
keys.set_null(2, true);
keys.set_null(3, true);
keys.set_i64(4, 2);
keys.set_i64(5, 2);
keys.resize(n);
let mut vals = ValueVector::new(PhysicalTypeID::Int64, n);
for i in 0..n {
vals.set_i64(i, i as i64);
}
vals.resize(n);
let input = vec![DataChunk::from_legacy(vec![keys, vals])];
let result = agg.execute(input).unwrap();
assert!(!result.is_empty());
assert_eq!(result[0].size, 3);
}
#[test]
fn test_filter_with_nulls() {
let mut v = ValueVector::new(PhysicalTypeID::Int64, 4);
v.set_i64(0, 1);
v.set_null(1, true);
v.set_i64(2, 3);
v.set_i64(3, 4);
v.resize(4);
let input = vec![DataChunk::from_legacy(vec![v])];
let filter = PhysicalFilter::new(Expression::Variable("a".into()));
let result = filter.execute(input.clone()).unwrap();
assert!(!result.is_empty());
assert_eq!(result[0].size, 3); }
#[test]
fn test_empty_table_scan() {
let scan = PhysicalScan::new("EmptyTable".into(), 0, 0);
let result = scan.execute(vec![]).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 0);
}
#[test]
fn test_empty_input_through_pipeline() {
let filter = PhysicalFilter::new(Expression::Constant(Constant::Bool(true)));
let result = filter.execute(vec![DataChunk::from_legacy(vec![])]).unwrap();
assert!(result.is_empty());
}
fn make_i64_chunk(values: &[i64]) -> DataChunk {
let mut v = ValueVector::new(PhysicalTypeID::Int64, values.len().max(1));
for (i, val) in values.iter().enumerate() {
v.set_i64(i, *val);
}
v.resize(values.len());
DataChunk::from_legacy(vec![v])
}
#[test]
fn test_union_all_basic() {
let mut left_v = ValueVector::new(PhysicalTypeID::Int64, 3);
left_v.set_i64(0, 1);
left_v.set_i64(1, 2);
left_v.set_i64(2, 3);
left_v.resize(3);
let left_data = vec![DataChunk::from_legacy(vec![left_v])];
let mut right_v = ValueVector::new(PhysicalTypeID::Int64, 2);
right_v.set_i64(0, 4);
right_v.set_i64(1, 5);
right_v.resize(2);
let right_data = vec![DataChunk::from_legacy(vec![right_v])];
let result = merge_union_chunks(left_data, right_data, true).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 5);
assert_eq!(result[0].get_i64(0, 0), Some(1));
assert_eq!(result[0].get_i64(0, 1), Some(2));
assert_eq!(result[0].get_i64(0, 2), Some(3));
assert_eq!(result[0].get_i64(0, 3), Some(4));
assert_eq!(result[0].get_i64(0, 4), Some(5));
}
#[test]
fn test_union_all_multiple_chunks() {
let mut v1 = ValueVector::new(PhysicalTypeID::Int64, 2);
v1.set_i64(0, 1);
v1.set_i64(1, 2);
v1.resize(2);
let mut v2 = ValueVector::new(PhysicalTypeID::Int64, 1);
v2.set_i64(0, 3);
v2.resize(1);
let left = vec![DataChunk::from_legacy(vec![v1]), DataChunk::from_legacy(vec![v2])];
let mut rv = ValueVector::new(PhysicalTypeID::Int64, 2);
rv.set_i64(0, 4);
rv.set_i64(1, 5);
rv.resize(2);
let right = vec![DataChunk::from_legacy(vec![rv])];
let result = merge_union_chunks(left, right, true).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 5);
assert_eq!(result[0].get_i64(0, 0), Some(1));
assert_eq!(result[0].get_i64(0, 4), Some(5));
}
#[test]
fn test_union_distinct_dedup() {
let mut lv = ValueVector::new(PhysicalTypeID::Int64, 3);
lv.set_i64(0, 1);
lv.set_i64(1, 2);
lv.set_i64(2, 3);
lv.resize(3);
let left = vec![DataChunk::from_legacy(vec![lv])];
let mut rv = ValueVector::new(PhysicalTypeID::Int64, 3);
rv.set_i64(0, 2);
rv.set_i64(1, 3);
rv.set_i64(2, 4);
rv.resize(3);
let right = vec![DataChunk::from_legacy(vec![rv])];
let result = merge_union_chunks(left, right, false).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 4);
}
#[test]
fn test_union_column_mismatch() {
let left = vec![DataChunk::from_legacy(vec![
ValueVector::new(PhysicalTypeID::Int64, 1),
ValueVector::new(PhysicalTypeID::Int64, 1),
])];
let right = vec![DataChunk::from_legacy(vec![ValueVector::new(PhysicalTypeID::Int64, 1)])];
let result = merge_union_chunks(left, right, true);
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("column count mismatch"));
}
#[test]
fn test_union_empty_left() {
let left = vec![];
let mut rv = ValueVector::new(PhysicalTypeID::Int64, 2);
rv.set_i64(0, 42);
rv.set_i64(1, 43);
rv.resize(2);
let right = vec![DataChunk::from_legacy(vec![rv])];
let result = merge_union_chunks(left, right, true).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 2);
assert_eq!(result[0].get_i64(0, 0), Some(42));
}
#[test]
fn test_union_empty_right() {
let mut lv = ValueVector::new(PhysicalTypeID::Int64, 2);
lv.set_i64(0, 99);
lv.set_i64(1, 100);
lv.resize(2);
let left = vec![DataChunk::from_legacy(vec![lv])];
let right = vec![];
let result = merge_union_chunks(left, right, true).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 2);
}
#[test]
fn test_union_all_multi_column() {
let mut left_v1 = ValueVector::new(PhysicalTypeID::Int64, 2);
left_v1.set_i64(0, 1);
left_v1.set_i64(1, 2);
left_v1.resize(2);
let mut left_v2 = ValueVector::new(PhysicalTypeID::String, 2);
left_v2.push_string("hello").unwrap();
left_v2.push_string("world").unwrap();
let left = vec![DataChunk::from_legacy(vec![left_v1, left_v2])];
let mut right_v1 = ValueVector::new(PhysicalTypeID::Int64, 1);
right_v1.set_i64(0, 3);
right_v1.resize(1);
let mut right_v2 = ValueVector::new(PhysicalTypeID::String, 1);
right_v2.push_string("foo").unwrap();
let right = vec![DataChunk::from_legacy(vec![right_v1, right_v2])];
let result = merge_union_chunks(left, right, true).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 3);
assert_eq!(result[0].get_i64(0, 0), Some(1));
assert_eq!(result[0].get_i64(0, 1), Some(2));
assert_eq!(result[0].get_i64(0, 2), Some(3));
}
#[test]
fn test_union_distinct_all_duplicates() {
let mut lv = ValueVector::new(PhysicalTypeID::Int64, 2);
lv.set_i64(0, 1);
lv.set_i64(1, 1);
lv.resize(2);
let left = vec![DataChunk::from_legacy(vec![lv])];
let mut rv = ValueVector::new(PhysicalTypeID::Int64, 2);
rv.set_i64(0, 1);
rv.set_i64(1, 1);
rv.resize(2);
let right = vec![DataChunk::from_legacy(vec![rv])];
let result = merge_union_chunks(left, right, false).unwrap();
assert_eq!(result[0].size, 1);
assert_eq!(result[0].get_i64(0, 0), Some(1));
}
#[test]
fn test_union_all_empty_chunks() {
let empty = ValueVector::new(PhysicalTypeID::Int64, 0);
let left = vec![DataChunk::from_legacy(vec![empty])];
let mut rv = ValueVector::new(PhysicalTypeID::Int64, 1);
rv.set_i64(0, 42);
rv.resize(1);
let right = vec![DataChunk::from_legacy(vec![rv])];
let result = merge_union_chunks(left, right, true).unwrap();
assert_eq!(result[0].size, 1);
assert_eq!(result[0].get_i64(0, 0), Some(42));
}
#[test]
fn test_cross_product_basic() {
let cross = PhysicalCrossProduct;
let left = vec![make_i64_chunk(&[1, 2, 3])];
let right = vec![make_i64_chunk(&[4, 5])];
let result = cross.execute_binary(&left, &right).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 6); assert_eq!(result[0].get_i64(0, 0), Some(1));
assert_eq!(result[0].get_i64(0, 1), Some(1));
assert_eq!(result[0].get_i64(0, 2), Some(2));
assert_eq!(result[0].get_i64(0, 3), Some(2));
assert_eq!(result[0].get_i64(0, 4), Some(3));
assert_eq!(result[0].get_i64(0, 5), Some(3));
}
#[test]
fn test_cross_product_multi_column() {
let cross = PhysicalCrossProduct;
let mut l1 = ValueVector::new(PhysicalTypeID::Int64, 2);
l1.set_i64(0, 1);
l1.set_i64(1, 2);
l1.resize(2);
let mut l2 = ValueVector::new(PhysicalTypeID::String, 2);
l2.push_string("a").unwrap();
l2.push_string("b").unwrap();
let left = DataChunk::from_legacy(vec![l1, l2]);
let mut r1 = ValueVector::new(PhysicalTypeID::Int64, 2);
r1.set_i64(0, 10);
r1.set_i64(1, 20);
r1.resize(2);
let right = DataChunk::from_legacy(vec![r1]);
let result = cross.execute_binary(&[left], &[right]).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].size, 4); assert_eq!(result[0].get_i64(0, 0), Some(1));
assert_eq!(result[0].get_i64(0, 1), Some(1));
assert_eq!(result[0].get_i64(0, 2), Some(2));
assert_eq!(result[0].get_i64(0, 3), Some(2));
assert_eq!(result[0].get_i64(2, 0), Some(10));
assert_eq!(result[0].get_i64(2, 1), Some(20));
assert_eq!(result[0].get_i64(2, 2), Some(10));
assert_eq!(result[0].get_i64(2, 3), Some(20));
}
#[test]
fn test_cross_product_empty_left() {
let cross = PhysicalCrossProduct;
let left = make_i64_chunk(&[]);
let right = make_i64_chunk(&[1, 2]);
let result = cross.execute_binary(&[left], &[right]).unwrap();
assert_eq!(result.len(), 0);
}
#[test]
fn test_cross_product_empty_right() {
let cross = PhysicalCrossProduct;
let left = make_i64_chunk(&[1, 2, 3]);
let right = make_i64_chunk(&[]);
let result = cross.execute_binary(&[left], &[right]).unwrap();
assert!(result.is_empty() || result[0].size == 0);
}
#[test]
fn test_cross_product_multi_chunk() {
let cross = PhysicalCrossProduct;
let left = vec![make_i64_chunk(&[1, 2]), make_i64_chunk(&[3])];
let right = vec![make_i64_chunk(&[4, 5])];
let result = cross.execute_binary(&left, &right).unwrap();
assert_eq!(result[0].size, 6); }
#[test]
fn test_semi_join_basic() {
let semi = PhysicalSemiJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let build = make_i64_chunk(&[2, 3]);
let probe = make_i64_chunk(&[1, 2, 3]);
let result = semi.execute_binary(&[build], &[probe]).unwrap();
assert_eq!(result[0].size, 2); }
#[test]
fn test_semi_join_no_match() {
let semi = PhysicalSemiJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let build = make_i64_chunk(&[4, 5]);
let probe = make_i64_chunk(&[1, 2, 3]);
let result = semi.execute_binary(&[build], &[probe]).unwrap();
assert!(result.is_empty() || result[0].size == 0);
}
#[test]
fn test_anti_join_basic() {
let anti = PhysicalAntiJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let build = make_i64_chunk(&[2, 3]);
let probe = make_i64_chunk(&[1, 2, 3]);
let result = anti.execute_binary(&[build], &[probe]).unwrap();
assert_eq!(result[0].size, 1); }
#[test]
fn test_anti_join_all_match() {
let anti = PhysicalAntiJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let build = make_i64_chunk(&[1, 2, 3]);
let probe = make_i64_chunk(&[1, 2, 3]);
let result = anti.execute_binary(&[build], &[probe]).unwrap();
assert!(result.is_empty() || result[0].size == 0);
}
#[test]
fn test_semi_join_empty_build() {
let semi = PhysicalSemiJoin {
build_columns: vec![0],
probe_columns: vec![0],
};
let build = make_i64_chunk(&[]);
let probe = make_i64_chunk(&[1, 2, 3]);
let result = semi.execute_binary(&[build], &[probe]).unwrap();
assert!(result.is_empty() || result[0].size == 0);
}
#[test]
fn test_intersect_basic() {
let intersect = PhysicalIntersect {
num_build_sides: 2,
probe_key_col: 0,
build_key_col: 0,
};
let build1 = make_i64_chunk(&[1, 2, 3]);
let build2 = make_i64_chunk(&[2, 3, 4]);
let probe = make_i64_chunk(&[2, 3]);
let build_chunks = vec![build1, build2];
let probe_chunks = vec![probe];
let result = intersect.execute_binary(&build_chunks, &probe_chunks).unwrap();
assert!(!result.is_empty(), "Expected non-empty result");
assert!(result[0].size > 0, "Expected at least one output row");
}
#[test]
fn test_intersect_no_common() {
let intersect = PhysicalIntersect {
num_build_sides: 2,
probe_key_col: 0,
build_key_col: 0,
};
let build1 = make_i64_chunk(&[1, 2, 3]);
let build2 = make_i64_chunk(&[4, 5, 6]);
let probe = make_i64_chunk(&[1, 2, 3, 4, 5, 6]);
let build_chunks = vec![build1, build2];
let probe_chunks = vec![probe];
let result = intersect.execute_binary(&build_chunks, &probe_chunks).unwrap();
assert!(result.is_empty() || result[0].size == 0);
}
#[test]
fn test_intersect_probe_key_missing() {
let intersect = PhysicalIntersect {
num_build_sides: 2,
probe_key_col: 0,
build_key_col: 0,
};
let build1 = make_i64_chunk(&[1, 3]);
let build2 = make_i64_chunk(&[3, 5]);
let probe = make_i64_chunk(&[1, 5]); let build_chunks = vec![build1, build2];
let probe_chunks = vec![probe];
let result = intersect.execute_binary(&build_chunks, &probe_chunks).unwrap();
assert!(result.is_empty() || result[0].size == 0);
}
#[test]
fn test_intersect_single_build_side() {
let intersect = PhysicalIntersect {
num_build_sides: 1,
probe_key_col: 0,
build_key_col: 0,
};
let build = make_i64_chunk(&[2, 3]);
let probe = make_i64_chunk(&[1, 2, 3, 4]);
let build_chunks = vec![build];
let probe_chunks = vec![probe];
let result = intersect.execute_binary(&build_chunks, &probe_chunks).unwrap();
assert!(!result.is_empty(), "Expected non-empty result for single build side");
assert!(result[0].size > 0, "Expected matching rows");
}
#[test]
fn test_intersect_no_probe() {
let intersect = PhysicalIntersect {
num_build_sides: 2,
probe_key_col: 0,
build_key_col: 0,
};
let build1 = make_i64_chunk(&[1, 2, 3]);
let build2 = make_i64_chunk(&[2, 3, 4]);
let probe = make_i64_chunk(&[]);
let build_chunks = vec![build1, build2];
let probe_chunks = vec![probe];
let result = intersect.execute_binary(&build_chunks, &probe_chunks).unwrap();
assert!(result.is_empty() || result[0].size == 0);
}
#[test]
fn test_intersect_cross_product_multi_match() {
let intersect = PhysicalIntersect {
num_build_sides: 2,
probe_key_col: 0,
build_key_col: 0,
};
let build1 = make_i64_chunk(&[1, 1, 5]);
let build2 = make_i64_chunk(&[1, 1, 1, 7]);
let probe = make_i64_chunk(&[1, 2]);
let build_chunks = vec![build1, build2];
let probe_chunks = vec![probe];
let result = intersect
.execute_sides(
&vec![vec![build_chunks[0].clone()], vec![build_chunks[1].clone()]],
&probe_chunks,
)
.unwrap();
assert!(!result.is_empty(), "Expected non-empty result");
assert_eq!(result[0].size, 6, "expected 2x3 cross product for probe key 1");
}
#[test]
fn test_intersect_key_col_resolution() {
let mut probe_v = ValueVector::new(PhysicalTypeID::Int64, 2);
probe_v.set_i64(0, 10);
probe_v.set_i64(1, 20);
let mut id_v = ValueVector::new(PhysicalTypeID::Int64, 2);
id_v.set_i64(0, 1);
id_v.set_i64(1, 2);
let mut probe = DataChunk::from_legacy(vec![probe_v, id_v]);
probe.field_names = vec!["a.other".into(), "a.id".into()];
let mut build_v = ValueVector::new(PhysicalTypeID::Int64, 2);
build_v.set_i64(0, 30);
build_v.set_i64(1, 40);
let mut build_id = ValueVector::new(PhysicalTypeID::Int64, 2);
build_id.set_i64(0, 1);
build_id.set_i64(1, 1);
let mut build = DataChunk::from_legacy(vec![build_v, build_id]);
build.field_names = vec!["a.other".into(), "a.id".into()];
let intersect = PhysicalIntersect {
num_build_sides: 1,
probe_key_col: 1,
build_key_col: 1,
};
let result = intersect.execute_sides(&vec![vec![build]], &vec![probe]).unwrap();
assert!(!result.is_empty(), "Expected non-empty result");
assert_eq!(
result[0].size, 2,
"probe id 1 matches build row 0, id 2 matches nothing"
);
}
}