use std::sync::{Arc, OnceLock};
use rudb_common::bounds::{Bound, Op};
use rudb_common::{Result, SessionTimeZone};
use rudb_plan::{ColumnBinding, Expr, ExprRef, Node, NodeRef, Plan};
use rudb_storage::{Blocked, Range};
use rudb_vector::{Chunk, Vector};
use crate::expr::evaluate_all_in_time_zone;
use crate::lookup::has_nulls;
use crate::schema::Schema;
use crate::table::{Across, hash};
const BUDGET: usize = 32 << 20;
#[derive(Debug, Default)]
pub(crate) struct Sideways<'a> {
keyed: OnceLock<Keyed<'a>>,
binding: OnceLock<ColumnBinding>,
found: OnceLock<Found>,
}
#[derive(Debug, Default)]
pub(crate) struct Found {
range: Option<(Bound, Bound)>,
filter: Option<Blocked>,
}
#[derive(Debug)]
pub(crate) struct Keyed<'a> {
plan: &'a Plan,
expr: ExprRef,
schema: Schema,
time_zone: SessionTimeZone,
}
impl<'a> Keyed<'a> {
pub(crate) fn new(
plan: &'a Plan,
expr: ExprRef,
schema: Schema,
time_zone: SessionTimeZone,
) -> Self {
Self { plan, expr, schema, time_zone }
}
pub(crate) fn parts(&self) -> (&'a Plan, [ExprRef; 1], &Schema, SessionTimeZone) {
(self.plan, [self.expr], &self.schema, self.time_zone)
}
}
impl<'a> Sideways<'a> {
pub(crate) fn new() -> Arc<Self> {
Arc::new(Self::default())
}
pub(crate) fn keying(&self, keyed: Keyed<'a>) {
let _ = self.keyed.set(keyed);
}
pub(crate) fn about(&self, binding: ColumnBinding) {
let _ = self.binding.set(binding);
}
pub(crate) fn keyed(&self) -> Option<&Keyed<'a>> {
self.keyed.get()
}
pub(crate) fn found(&self, found: Found) {
let _ = self.found.set(found);
}
pub(crate) fn tests(&self, index: u32) -> Vec<(usize, Op, Bound)> {
let (Some(binding), Some(Some((low, high)))) =
(self.binding.get(), self.found.get().map(|found| &found.range))
else {
return Vec::new();
};
if binding.table != index {
return Vec::new();
}
let column = binding.column as usize;
vec![(column, Op::GreaterOrEqual, low.clone()), (column, Op::LessOrEqual, high.clone())]
}
pub(crate) fn sifting(&self, index: u32) -> Option<(usize, &Blocked)> {
let binding = self.binding.get()?;
if binding.table != index {
return None;
}
Some((binding.column as usize, self.found.get()?.filter.as_ref()?))
}
}
pub(crate) fn beneath(plan: &Plan, node: NodeRef, binding: ColumnBinding) -> Option<ColumnBinding> {
let mut at = node;
let mut binding = binding;
loop {
match *plan.node(at) {
Node::Get { index, .. } | Node::TableFunction { index, .. } => {
return (binding.table == index).then_some(binding);
}
Node::Filter { input, .. } => at = input,
Node::Project { input, index, exprs, .. } => {
if binding.table == index {
let exprs = plan.expr_list(exprs);
let at = exprs.get(binding.column as usize)?;
let Expr::Column(inner) = *plan.expr(*at) else { return None };
binding = inner;
}
at = input;
}
_ => return None,
}
}
}
pub(crate) fn found(keyed: &Keyed<'_>, chunks: &[Chunk]) -> Result<Found> {
let (plan, exprs, schema, time_zone) = keyed.parts();
let rows: usize = chunks.iter().map(Chunk::len).sum();
let mut extremes = Extremes::default();
let mut filter = Blocked::sized(rows, BUDGET);
let mut hashes = Vec::new();
for chunk in chunks {
let keys = evaluate_all_in_time_zone(plan, &exprs, schema, chunk, time_zone)?;
let Some(keys) = keys.first() else { continue };
extremes.widen(keys);
let Some(filter) = filter.as_mut() else { continue };
hash(std::slice::from_ref(keys), chunk.len(), &mut hashes, Across::TwoInputs);
let nullable = has_nulls(keys, chunk.len());
for (row, &word) in hashes.iter().enumerate() {
if nullable && keys.is_null_at(row) {
continue;
}
filter.add(word);
}
}
Ok(Found { range: extremes.into_range(), filter })
}
#[derive(Debug, Default, Clone)]
pub(crate) struct Extremes {
low: Option<Bound>,
high: Option<Bound>,
}
impl Extremes {
pub(crate) fn widen(&mut self, keys: &Vector) {
let range = Range::of(keys);
if let Some(low) = range.low {
self.low = Some(match self.low.take() {
Some(held) => held.smaller(low),
None => low,
});
}
if let Some(high) = range.high {
self.high = Some(match self.high.take() {
Some(held) => held.larger(high),
None => high,
});
}
}
pub(crate) fn into_range(self) -> Option<(Bound, Bound)> {
Some((self.low?, self.high?))
}
}
impl Found {
#[cfg(test)]
pub(crate) fn of(range: Option<(Bound, Bound)>, filter: Option<Blocked>) -> Self {
Self { range, filter }
}
}
#[cfg(test)]
mod tests {
use rudb_common::bounds::{Bound, Op};
use rudb_common::{Field, LogicalType, SessionTimeZone, Value};
use rudb_plan::{ColumnBinding, Expr, ExprRef, Plan};
use rudb_storage::Blocked;
use rudb_vector::{Chunk, Vector};
use super::{Across, Extremes, Found, Keyed, Schema, Sideways, beneath, found, hash};
fn column(values: &[Option<i32>]) -> Vector {
let values: Vec<Value> =
values.iter().map(|value| value.map_or(Value::Null, Value::Integer)).collect();
Vector::from_values(LogicalType::Integer, &values).expect("a column of integers")
}
#[test]
fn the_range_of_several_chunks_covers_every_one_of_them() {
let mut extremes = Extremes::default();
extremes.widen(&column(&[Some(5), Some(9)]));
extremes.widen(&column(&[Some(2), Some(7)]));
assert_eq!(extremes.into_range(), Some((Bound::Int(2), Bound::Int(9))));
}
#[test]
fn a_column_of_nulls_widens_nothing() {
let mut extremes = Extremes::default();
extremes.widen(&column(&[Some(4)]));
extremes.widen(&column(&[None, None]));
assert_eq!(extremes.into_range(), Some((Bound::Int(4), Bound::Int(4))));
}
#[test]
fn nothing_seen_is_no_range() {
assert_eq!(Extremes::default().into_range(), None);
}
#[test]
fn a_scan_is_told_only_about_its_own_column() {
let sideways = Sideways::new();
sideways.found(Found::of(Some((Bound::Int(1), Bound::Int(4))), None));
sideways.about(ColumnBinding::new(7, 2));
assert!(sideways.tests(8).is_empty(), "another table's scan");
assert_eq!(
sideways.tests(7),
vec![(2, Op::GreaterOrEqual, Bound::Int(1)), (2, Op::LessOrEqual, Bound::Int(4)),]
);
}
#[test]
fn an_unarmed_handoff_and_an_empty_build_side_both_say_nothing() {
let unarmed = Sideways::new();
assert!(unarmed.tests(1).is_empty());
assert!(unarmed.sifting(1).is_none());
let empty = Sideways::new();
empty.about(ColumnBinding::new(1, 0));
empty.found(Found::of(None, None));
assert!(empty.tests(1).is_empty());
assert!(empty.sifting(1).is_none());
}
fn chunk(values: &[Option<i32>]) -> Chunk {
Chunk::new(vec![column(values)]).expect("one column is one length")
}
fn chunks(values: &[Option<i32>]) -> Vec<Chunk> {
values.chunks(512).map(chunk).collect()
}
fn key(plan: &mut Plan) -> (ExprRef, Schema) {
let expr = plan.add_expr(Expr::Column(ColumnBinding::new(1, 0)), LogicalType::Integer);
let schema = Schema::numbered(vec![Field::new("k", LogicalType::Integer)], 1);
(expr, schema)
}
fn through(filter: &Blocked, values: &[Option<i32>]) -> Vec<bool> {
let probe = column(values);
let mut hashes = Vec::new();
hash(std::slice::from_ref(&probe), values.len(), &mut hashes, Across::TwoInputs);
hashes.iter().map(|&word| filter.holds(word)).collect()
}
#[test]
fn a_build_side_is_read_for_both_its_range_and_its_keys() {
let mut plan = Plan::new();
let (expr, schema) = key(&mut plan);
let keyed = Keyed::new(&plan, expr, schema, SessionTimeZone::default());
let found = found(&keyed, &[chunk(&[Some(5), Some(9)]), chunk(&[Some(2)])])
.expect("a column of integers");
assert_eq!(found.range, Some((Bound::Int(2), Bound::Int(9))));
let filter = found.filter.expect("a filter over three keys");
assert_eq!(through(&filter, &[Some(5), Some(9), Some(2)]), [true, true, true]);
}
#[test]
fn no_key_that_went_in_is_ever_turned_away() {
let mut plan = Plan::new();
let (expr, schema) = key(&mut plan);
let keyed = Keyed::new(&plan, expr, schema, SessionTimeZone::default());
let keys: Vec<Option<i32>> = (0..4_000).map(|value| Some(value * 7 + 11)).collect();
let found = found(&keyed, &chunks(&keys)).expect("a column of integers");
let filter = found.filter.expect("a filter over four thousand keys");
assert!(through(&filter, &keys).into_iter().all(|held| held), "a key it was given");
}
#[test]
fn a_key_the_build_side_never_held_is_nearly_always_turned_away() {
let mut plan = Plan::new();
let (expr, schema) = key(&mut plan);
let keyed = Keyed::new(&plan, expr, schema, SessionTimeZone::default());
let keys: Vec<Option<i32>> = (0..4_000).map(|value| Some(value * 7 + 11)).collect();
let absent: Vec<Option<i32>> =
(0..4_000).map(|value| Some(value * 7 + 1_000_000)).collect();
let found = found(&keyed, &chunks(&keys)).expect("a column of integers");
let filter = found.filter.expect("a filter over four thousand keys");
let through = through(&filter, &absent).into_iter().filter(|&held| held).count();
assert!(through < absent.len() / 10, "{through} of {} got through", absent.len());
}
#[test]
fn a_null_is_not_a_key_the_filter_holds() {
let mut plan = Plan::new();
let (expr, schema) = key(&mut plan);
let keyed = Keyed::new(&plan, expr, schema, SessionTimeZone::default());
let found = found(&keyed, &[chunk(&[Some(3), None, Some(4)])]).expect("integers");
assert_eq!(found.range, Some((Bound::Int(3), Bound::Int(4))));
assert_eq!(through(&found.filter.expect("a filter"), &[None]), [false]);
}
#[test]
fn an_empty_build_side_turns_every_driving_row_away() {
let mut plan = Plan::new();
let (expr, schema) = key(&mut plan);
let keyed = Keyed::new(&plan, expr, schema, SessionTimeZone::default());
let found = found(&keyed, &[]).expect("nothing to read");
assert_eq!(found.range, None);
assert_eq!(through(&found.filter.expect("a filter of no keys"), &[Some(1)]), [false]);
}
fn driving(text: &str) -> Plan {
Plan::parse(text).expect("the plan text round trips")
}
#[test]
fn a_projection_between_the_join_and_the_scan_renames_the_column_the_filter_is_about() {
let plan = driving(
"Project #1 [#0.1::INTEGER AS k]\n \
TableFunction read_parquet args=['f'::VARCHAR] #0 [a::INTEGER, k::INTEGER]",
);
assert_eq!(
beneath(&plan, plan.root(), ColumnBinding::new(1, 0)),
Some(ColumnBinding::new(0, 1)),
"the scan's own name for the projection's column"
);
}
#[test]
fn a_filter_between_the_two_leaves_the_binding_alone() {
let plan = driving(
"Project #1 [#0.0::INTEGER AS k]\n \
Filter (#0.0::INTEGER > 3::INTEGER)::BOOLEAN\n \
TableFunction read_parquet args=['f'::VARCHAR] #0 [k::INTEGER]",
);
assert_eq!(
beneath(&plan, plan.root(), ColumnBinding::new(1, 0)),
Some(ColumnBinding::new(0, 0))
);
}
#[test]
fn a_computed_column_is_not_a_column_the_filter_can_be_about() {
let plan = driving(
"Project #1 [(#0.0::INTEGER > 3::INTEGER)::BOOLEAN AS k]\n \
TableFunction read_parquet args=['f'::VARCHAR] #0 [k::INTEGER]",
);
assert_eq!(beneath(&plan, plan.root(), ColumnBinding::new(1, 0)), None);
}
#[test]
fn a_walk_that_does_not_reach_the_scan_it_is_about_arms_nothing() {
let plan = driving(
"Limit 5 offset 0\n \
TableFunction read_parquet args=['f'::VARCHAR] #0 [k::INTEGER]",
);
assert_eq!(
beneath(&plan, plan.root(), ColumnBinding::new(0, 0)),
None,
"a node in the way"
);
let plan = driving("TableFunction read_parquet args=['f'::VARCHAR] #0 [k::INTEGER]");
assert_eq!(
beneath(&plan, plan.root(), ColumnBinding::new(3, 0)),
None,
"another table's column"
);
}
}