use rudb_common::{
Error, LogicalType, PhysicalType, Result, Session, SessionTimeZone, Span, Value,
};
use rudb_kernels::{
Comparison, Connective, Found, Held, Lookup, Members, Recipe, cast_in_time_zone, combine,
compare_prepared, in_set, is_true, refine_flags, refine_prepared, select_prepared, selection,
};
use rudb_plan::{CompareOp, ConjunctionOp, Expr, ExprRef, Plan};
use rudb_vector::{Assembly, Chunk, Selection, Vector};
use std::collections::HashMap;
use std::sync::Arc;
use crate::fused::Fused;
use crate::lambda::{Lambda, lambda_call};
use crate::ordering::Ordering;
use crate::schema::Schema;
use crate::written::written;
const _: () = {
const fn assert_shareable<T: Send + Sync>() {}
assert_shareable::<Prepared>();
};
#[derive(Debug)]
pub struct Prepared {
steps: Vec<Step>,
types: Vec<LogicalType>,
spans: Vec<Span>,
operands: Vec<usize>,
last_use: Vec<usize>,
roots: Vec<usize>,
shared: HashMap<ExprRef, usize>,
share: bool,
fuse: bool,
time_zone: SessionTimeZone,
}
#[derive(Debug)]
enum Step {
Column(usize),
Constant(Value),
Cast {
input: usize,
try_cast: bool,
},
Compare {
op: Comparison,
left: usize,
right: usize,
held: Option<Held>,
},
Conjunction {
op: Connective,
start: usize,
len: usize,
},
Function {
recipe: Recipe,
written: String,
start: usize,
len: usize,
},
InSet {
input: usize,
members: Members,
},
Case {
arms: Vec<PreparedArm>,
otherwise: Option<Prepared>,
blend: Option<Blend>,
},
Fused {
fused: Box<Fused>,
fallback: Box<Prepared>,
},
Lambda {
inputs: Vec<usize>,
runner: Box<Lambda>,
body: Box<Prepared>,
},
}
#[derive(Debug)]
struct PreparedArm {
when: Prepared,
then: Prepared,
}
#[derive(Debug)]
struct Blend {
branches: Vec<Branch>,
literals: Vec<(String, Lookup)>,
}
#[derive(Debug, Clone, Copy)]
enum Branch {
Column(usize),
Literal(usize),
}
#[derive(Debug, Default)]
pub struct Scratch {
slots: Vec<Option<Vector>>,
orders: Vec<Option<Ordering>>,
}
impl Scratch {
#[cfg(test)]
fn order(&self, step: usize) -> Option<&[usize]> {
self.orders[step].as_ref().map(Ordering::order)
}
}
impl Prepared {
pub fn new(plan: &Plan, exprs: &[ExprRef], schema: &Schema) -> Result<Self> {
Self::build(plan, exprs, schema, false)
}
pub(crate) fn shared(plan: &Plan, exprs: &[ExprRef], schema: &Schema) -> Result<Self> {
Self::build(plan, exprs, schema, true)
}
fn build(plan: &Plan, exprs: &[ExprRef], schema: &Schema, share: bool) -> Result<Self> {
Self::built(plan, exprs, schema, share, true)
}
fn built(
plan: &Plan,
exprs: &[ExprRef],
schema: &Schema,
share: bool,
fuse: bool,
) -> Result<Self> {
let mut prepared = Self {
steps: Vec::new(),
types: Vec::new(),
spans: Vec::new(),
operands: Vec::new(),
last_use: Vec::new(),
roots: Vec::new(),
shared: HashMap::new(),
share,
fuse,
time_zone: SessionTimeZone::default(),
};
for &expr in exprs {
let root = prepared.push(plan, expr, schema)?;
prepared.roots.push(root);
}
prepared.last_use = prepared.last_uses();
Ok(prepared)
}
#[must_use]
pub fn in_session(mut self, session: &Session) -> Self {
self.set_time_zone(session.session_time_zone());
self
}
fn set_time_zone(&mut self, time_zone: SessionTimeZone) {
self.time_zone = time_zone;
for step in &mut self.steps {
match step {
Step::Lambda { body, .. } => body.set_time_zone(time_zone),
Step::Fused { fallback, .. } => fallback.set_time_zone(time_zone),
_ => {}
}
}
}
fn last_uses(&self) -> Vec<usize> {
let mut last = vec![usize::MAX; self.steps.len()];
for index in 0..self.steps.len() {
self.for_each_operand(index, |operand| last[operand] = index);
}
for &root in &self.roots {
last[root] = usize::MAX;
}
last
}
fn for_each_operand(&self, index: usize, mut visit: impl FnMut(usize)) {
match &self.steps[index] {
Step::Column(_) | Step::Constant(_) | Step::Case { .. } | Step::Fused { .. } => {}
Step::Cast { input, .. } | Step::InSet { input, .. } => visit(*input),
Step::Lambda { inputs, .. } => inputs.iter().for_each(|&input| visit(input)),
Step::Compare { left, right, .. } => {
visit(*left);
visit(*right);
}
Step::Conjunction { start, len, .. } | Step::Function { start, len, .. } => {
for &operand in &self.operands[*start..*start + *len] {
visit(operand);
}
}
}
}
pub fn one(plan: &Plan, expr: ExprRef, schema: &Schema) -> Result<Self> {
Self::new(plan, &[expr], schema)
}
#[must_use]
pub fn scratch(&self) -> Scratch {
Scratch {
slots: (0..self.steps.len()).map(|_| None).collect(),
orders: (0..self.steps.len()).map(|_| None).collect(),
}
}
#[must_use]
pub fn len(&self) -> usize {
self.roots.len()
}
#[cfg(test)]
fn literals_built(&self) -> usize {
self.steps.iter().filter(|step| matches!(step, Step::Compare { held: Some(_), .. })).count()
}
#[cfg(test)]
fn sets(&self) -> usize {
self.steps.iter().filter(|step| matches!(step, Step::InSet { .. })).count()
}
#[cfg(test)]
fn fused(&self) -> usize {
self.steps.iter().filter(|step| matches!(step, Step::Fused { .. })).count()
}
#[cfg(test)]
fn hoisted(&self) -> usize {
self.steps
.iter()
.filter(|step| matches!(step, Step::Function { recipe, .. } if recipe.hoists()))
.count()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.roots.is_empty()
}
#[must_use]
pub fn passes(&self) -> usize {
self.steps
.iter()
.filter(|step| !matches!(step, Step::Column(_) | Step::Constant(_)))
.count()
}
pub fn evaluate(
&self,
chunk: &Chunk,
scratch: &mut Scratch,
out: &mut Vec<Vector>,
) -> Result<()> {
self.run(chunk, scratch)?;
let mut remaining: HashMap<usize, usize> = HashMap::new();
for &root in &self.roots {
*remaining.entry(root).or_default() += 1;
}
for &root in &self.roots {
match self.steps[root] {
Step::Column(position) => out.push(chunk.column(position)?.clone()),
_ => {
let Some(left) = remaining.get_mut(&root) else {
return Err(Error::internal("a prepared root was not counted"));
};
*left -= 1;
if *left == 0 {
out.push(scratch.slots[root].take().ok_or_else(|| missing(root))?);
} else {
out.push(
scratch.slots[root].as_ref().ok_or_else(|| missing(root))?.clone(),
);
}
}
}
}
Ok(())
}
pub fn evaluate_taking(
&self,
chunk: Chunk,
scratch: &mut Scratch,
out: &mut Vec<Vector>,
) -> Result<()> {
self.run(&chunk, scratch)?;
let width = chunk.width();
let mut columns: Vec<Option<Vector>> = chunk.into_columns().into_iter().map(Some).collect();
let mut uses = vec![0usize; width];
let mut remaining: HashMap<usize, usize> = HashMap::new();
for &root in &self.roots {
match self.steps[root] {
Step::Column(position) if position < width => uses[position] += 1,
_ => *remaining.entry(root).or_default() += 1,
}
}
for &root in &self.roots {
if let Step::Column(position) = self.steps[root] {
let missing = || {
Error::internal(format!(
"column {position} of a chunk that has {width} columns"
))
};
let slot = columns.get_mut(position).ok_or_else(missing)?;
let left = &mut uses[position];
*left -= 1;
let column = if *left == 0 { slot.take() } else { slot.clone() };
out.push(column.ok_or_else(missing)?);
continue;
}
let Some(left) = remaining.get_mut(&root) else {
return Err(Error::internal("a prepared root was not counted"));
};
*left -= 1;
if *left == 0 {
out.push(scratch.slots[root].take().ok_or_else(|| missing(root))?);
} else {
out.push(scratch.slots[root].as_ref().ok_or_else(|| missing(root))?.clone());
}
}
Ok(())
}
pub fn evaluate_one<'s>(
&'s self,
chunk: &'s Chunk,
scratch: &'s mut Scratch,
) -> Result<&'s Vector> {
let [root] = self.roots[..] else {
return Err(Error::internal(format!(
"evaluate_one over a prepared expression of {} roots",
self.roots.len()
)));
};
self.run(chunk, scratch)?;
self.operand(root, chunk, &scratch.slots)
}
pub fn evaluate_filter(&self, chunk: &Chunk, scratch: &mut Scratch) -> Result<Selection> {
let [root] = self.roots[..] else {
return Err(Error::internal(format!(
"evaluate_filter over a prepared expression of {} roots",
self.roots.len()
)));
};
scratch.slots.clear();
scratch.slots.resize_with(self.steps.len(), || None);
self.thread(root, 0, chunk, scratch, None)
}
#[must_use]
pub fn conjuncts(&self) -> Option<usize> {
let [root] = self.roots[..] else { return None };
match self.steps[root] {
Step::Conjunction { op: Connective::And, len, .. } if !self.share => Some(len),
_ => None,
}
}
pub fn evaluate_settled(
&self,
chunk: &Chunk,
scratch: &mut Scratch,
settled: &[bool],
) -> Result<Selection> {
if self.conjuncts() != Some(settled.len()) || !settled.contains(&true) {
return self.evaluate_filter(chunk, scratch);
}
let [root] = self.roots[..] else {
return Err(Error::internal("a settled filter over several roots"));
};
scratch.slots.clear();
scratch.slots.resize_with(self.steps.len(), || None);
self.branches(root, 0, chunk, scratch, None, settled)
}
fn branches(
&self,
index: usize,
begin: usize,
chunk: &Chunk,
scratch: &mut Scratch,
live: Option<&Selection>,
settled: &[bool],
) -> Result<Selection> {
let Step::Conjunction { op, start, len } = self.steps[index] else {
return Err(Error::internal("a connective walk over a step that is not a connective"));
};
let operands = &self.operands[start..start + len];
let rows = chunk.len();
let mut order = scratch.orders[index]
.take()
.unwrap_or_else(|| Ordering::new(op, self.weights(operands, begin)));
let mut carried: Option<Selection> = live.cloned();
for slot in 0..len {
if carried.as_ref().is_some_and(Selection::is_empty) {
break;
}
let which = order.at(slot);
if settled.get(which) == Some(&true) {
continue;
}
let operand = operands[which];
let from = if which == 0 { begin } else { operands[which - 1] + 1 };
let given = carried.as_ref().map_or(rows, Selection::len);
let answered = self.thread(operand, from, chunk, scratch, carried.as_ref())?;
order.observed(which, given, answered.len());
carried = Some(match (op, carried) {
(Connective::And, _) => answered,
(Connective::Or, None) => answered.complement(rows),
(Connective::Or, Some(carried)) => carried.without(&answered),
});
for step in from..=operand {
if self.last_use[step] <= operand {
scratch.slots[step] = None;
}
}
}
order.relearn();
scratch.orders[index] = Some(order);
Ok(match (op, carried) {
(Connective::And, None) => live.cloned().unwrap_or_else(|| Selection::identity(rows)),
(Connective::And, Some(kept)) => kept,
(Connective::Or, None) => Selection::empty(),
(Connective::Or, Some(missed)) => match live {
None => missed.complement(rows),
Some(live) => live.without(&missed),
},
})
}
fn weights(&self, operands: &[usize], begin: usize) -> Vec<f64> {
let mut costs = Vec::with_capacity(operands.len());
let mut from = begin;
for &operand in operands {
costs.push((from..=operand).map(|step| self.weight(step)).sum());
from = operand + 1;
}
costs
}
fn weight(&self, index: usize) -> f64 {
match &self.steps[index] {
Step::Column(_) => 0.0,
Step::Constant(_) => 0.25,
Step::Conjunction { .. } => 0.0,
Step::Cast { input, .. } => 2.0 * touching(&self.types[*input]),
Step::Compare { left, .. } => touching(&self.types[*left]),
Step::InSet { input, .. } => 2.0 * touching(&self.types[*input]),
Step::Function { start, len, .. } => {
let widest = self.operands[*start..*start + *len]
.iter()
.map(|&argument| touching(&self.types[argument]))
.fold(1.0, f64::max);
4.0 * widest
}
Step::Case { arms, .. } => 4.0 * arms.len() as f64,
Step::Lambda { .. } => 16.0,
Step::Fused { fused, .. } => fused.len() as f64,
}
}
fn thread(
&self,
index: usize,
begin: usize,
chunk: &Chunk,
scratch: &mut Scratch,
live: Option<&Selection>,
) -> Result<Selection> {
if matches!(self.steps[index], Step::Conjunction { .. }) {
return self.branches(index, begin, chunk, scratch, live, &[]);
}
for step in begin..index {
self.run_step(step, chunk, scratch)?;
}
if let Step::Compare { op, left, right, held } = &self.steps[index] {
let one = self.operand(*left, chunk, &scratch.slots)?;
let other = self.operand(*right, chunk, &scratch.slots)?;
let held = held.as_ref();
return match live {
None => select_prepared(*op, one, other, held),
Some(live) => refine_prepared(*op, one, other, live, held),
};
}
if let (Some(live), Step::Function { recipe, written, start, len }) =
(live, &self.steps[index])
{
if matches!(recipe.name(), "~~" | "!~~" | "~~*" | "!~~*")
&& live.len().saturating_mul(4) <= chunk.len()
{
let flags = self
.with_operands(*start, *len, chunk, &scratch.slots, |args| {
let gathered = args
.iter()
.map(|arg| arg.gather(live.indices()))
.collect::<Result<Vec<_>>>()?;
let narrowed = gathered.iter().collect::<Vec<_>>();
rudb_kernels::call_prepared(
recipe,
&narrowed,
&self.types[index],
Some(&|| written.clone()),
)
})
.map_err(|error| error.with_fallback_span(self.spans[index]))?;
return Ok(selection(&flags, live.len()).compose(live));
}
}
self.run_step(index, chunk, scratch)?;
let flags = self.operand(index, chunk, &scratch.slots)?;
match live {
None => Ok(selection(flags, chunk.len())),
Some(live) => refine_flags(flags, live),
}
}
fn run(&self, chunk: &Chunk, scratch: &mut Scratch) -> Result<()> {
scratch.slots.clear();
scratch.slots.resize_with(self.steps.len(), || None);
for index in 0..self.steps.len() {
self.run_step(index, chunk, scratch)?;
}
Ok(())
}
fn run_step(&self, index: usize, chunk: &Chunk, scratch: &mut Scratch) -> Result<()> {
let produced = self
.step(index, chunk, &scratch.slots)
.map_err(|error| error.with_fallback_span(self.spans[index]))?;
scratch.slots[index] = produced;
let slots = &mut scratch.slots;
self.for_each_operand(index, |operand| {
if self.last_use[operand] == index {
slots[operand] = None;
}
});
Ok(())
}
fn step(
&self,
index: usize,
chunk: &Chunk,
slots: &[Option<Vector>],
) -> Result<Option<Vector>> {
let ty = &self.types[index];
let produced = match &self.steps[index] {
Step::Column(_) => None,
Step::Constant(value) => Some(Vector::constant(ty.clone(), value.clone(), chunk.len())),
Step::Cast { input, try_cast } => Some(cast_in_time_zone(
self.operand(*input, chunk, slots)?,
ty,
*try_cast,
Some(self.time_zone),
)?),
Step::Compare { op, left, right, held } => Some(compare_prepared(
*op,
self.operand(*left, chunk, slots)?,
self.operand(*right, chunk, slots)?,
held.as_ref(),
)?),
Step::Conjunction { op, start, len } => {
Some(
self.with_operands(*start, *len, chunk, slots, |children| {
combine(*op, children)
})?,
)
}
Step::Function { recipe, written, start, len } => {
Some(self.with_operands(*start, *len, chunk, slots, |args| {
rudb_kernels::call_prepared(recipe, args, ty, Some(&|| written.clone()))
})?)
}
Step::InSet { input, members } => {
Some(in_set(self.operand(*input, chunk, slots)?, members, ty)?)
}
Step::Case { arms, otherwise, blend } => {
Some(self.case(chunk, arms, otherwise.as_ref(), blend.as_ref(), ty)?)
}
Step::Fused { fused, fallback } => Some(match fused.run(chunk) {
Some(answer) => answer,
None => fallback.evaluate_one(chunk, &mut fallback.scratch())?.clone(),
}),
Step::Lambda { inputs, runner, body } => {
let mut operands = Vec::with_capacity(inputs.len());
for &input in inputs {
operands.push(self.operand(input, chunk, slots)?);
}
let mut scratch = body.scratch();
Some(runner.run(&operands, chunk, &mut |inner| {
body.evaluate_one(inner, &mut scratch).cloned()
})?)
}
};
Ok(produced)
}
fn operand<'v>(
&self,
index: usize,
chunk: &'v Chunk,
slots: &'v [Option<Vector>],
) -> Result<&'v Vector> {
if let Step::Column(position) = self.steps[index] {
return chunk.column(position);
}
slots[index].as_ref().ok_or_else(|| missing(index))
}
fn with_operands<'v, T>(
&self,
start: usize,
len: usize,
chunk: &'v Chunk,
slots: &'v [Option<Vector>],
run: impl FnOnce(&[&'v Vector]) -> Result<T>,
) -> Result<T> {
match self.operands[start..start + len] {
[a] => run(&[self.operand(a, chunk, slots)?]),
[a, b] => run(&[self.operand(a, chunk, slots)?, self.operand(b, chunk, slots)?]),
[a, b, c] => run(&[
self.operand(a, chunk, slots)?,
self.operand(b, chunk, slots)?,
self.operand(c, chunk, slots)?,
]),
_ => {
let gathered = self.gather(start, len, chunk, slots)?;
run(&gathered)
}
}
}
fn gather<'v>(
&self,
start: usize,
len: usize,
chunk: &'v Chunk,
slots: &'v [Option<Vector>],
) -> Result<Vec<&'v Vector>> {
let mut gathered = Vec::with_capacity(len);
for &operand in &self.operands[start..start + len] {
gathered.push(self.operand(operand, chunk, slots)?);
}
Ok(gathered)
}
fn case(
&self,
chunk: &Chunk,
arms: &[PreparedArm],
otherwise: Option<&Prepared>,
blend: Option<&Blend>,
ty: &LogicalType,
) -> Result<Vector> {
let claimed = self.claims(chunk, arms)?;
if let Some(blend) = blend {
if let Some(blended) = blended(chunk, &claimed, blend)? {
return Ok(blended);
}
}
let mut built = Assembly::new(ty.clone(), chunk.len())?;
let branches = arms.iter().map(|arm| &arm.then).map(Some).chain([otherwise]);
for (branch, rows) in branches.zip(&claimed) {
let (Some(branch), false) = (branch, rows.is_empty()) else { continue };
let cut;
let matched = if rows.len() == chunk.len() {
chunk
} else {
cut = narrow(chunk, rows)?;
&cut
};
let mut scratch = branch.scratch();
let results = branch.evaluate_one(matched, &mut scratch)?;
built.place(&placed(rows)?, results)?;
}
built.finish()
}
fn claims(&self, chunk: &Chunk, arms: &[PreparedArm]) -> Result<Vec<Vec<usize>>> {
let mut claimed = Vec::with_capacity(arms.len() + 1);
let mut pending: Vec<usize> = (0..chunk.len()).collect();
for arm in arms {
if pending.is_empty() {
claimed.push(Vec::new());
continue;
}
let cut;
let narrowed = if pending.len() == chunk.len() {
chunk
} else {
cut = narrow(chunk, &pending)?;
&cut
};
let mut scratch = arm.when.scratch();
let flags = arm.when.evaluate_one(narrowed, &mut scratch)?;
let mut taken = Vec::new();
let mut still = Vec::new();
for (at, &row) in pending.iter().enumerate() {
if is_true(&flags.value_at(at)) {
taken.push(row);
} else {
still.push(row);
}
}
claimed.push(taken);
pending = still;
}
claimed.push(pending);
Ok(claimed)
}
fn push(&mut self, plan: &Plan, expr: ExprRef, schema: &Schema) -> Result<usize> {
if self.share {
if let Some(&step) = self.shared.get(&expr) {
return Ok(step);
}
}
let ty = plan.expr_type(expr).clone();
if self.fuse {
if let Some(fused) = Fused::compile(plan, expr, schema) {
let fallback = Self::built(plan, &[expr], schema, false, false)?;
let step = Step::Fused { fused: Box::new(fused), fallback: Box::new(fallback) };
return Ok(self.place(plan, expr, step, ty));
}
}
if let Some((stamp, count)) = stamped_seconds(plan, expr) {
let (start, len) = self.push_list(plan, &[stamp, count], schema)?;
let step = Step::Function {
recipe: Recipe::new("__rudb_stamp_seconds", &self.literals(start, len)),
written: written(plan, expr, schema),
start,
len,
};
return Ok(self.place(plan, expr, step, ty));
}
let step = match *plan.expr(expr) {
Expr::Column(binding) => {
let position = schema.position_of(binding).ok_or_else(|| {
Error::internal(format!(
"column #{}.{} is not in the schema this operator was given",
binding.table, binding.column
))
})?;
Step::Column(position)
}
Expr::Constant(reference) => Step::Constant(plan.value(reference).clone()),
Expr::Cast { input, try_cast } => {
Step::Cast { input: self.push(plan, input, schema)?, try_cast }
}
Expr::Compare { op, left, right } => {
let left = self.push(plan, left, schema)?;
let right = self.push(plan, right, schema)?;
Step::Compare { op: comparison(op), left, right, held: self.held(left, right) }
}
Expr::Conjunction { op, children } => {
let list = plan.expr_list(children).to_vec();
match self.membership(plan, connective(op), &list, schema)? {
Some(step) => step,
None => {
let (start, len) = self.push_list(plan, &list, schema)?;
Step::Conjunction { op: connective(op), start, len }
}
}
}
Expr::Function { name, args } if lambda_call(plan, args).is_some() => {
let Some((lambda, inputs)) = lambda_call(plan, args) else {
return Err(Error::internal("a lambda call without a lambda"));
};
let Expr::Lambda { body, .. } = *plan.expr(lambda) else {
return Err(Error::internal("a lambda call without a lambda"));
};
let runner = Lambda::new(plan, plan.string(name), lambda, &inputs, schema)?;
let body = Self::one(plan, body, runner.schema())?;
let mut steps = Vec::with_capacity(inputs.len());
for &input in &inputs {
steps.push(self.push(plan, input, schema)?);
}
Step::Lambda { inputs: steps, runner: Box::new(runner), body: Box::new(body) }
}
Expr::LambdaParam(binding) => {
let position = schema.position_of(binding).ok_or_else(|| {
Error::internal(format!(
"lambda parameter @{}.{} is not in the schema its body was given",
binding.table, binding.column
))
})?;
Step::Column(position)
}
Expr::Lambda { .. } => {
return Err(Error::internal(
"a lambda was evaluated outside the function that takes it",
));
}
Expr::Function { name, args } => {
let (start, len) = self.push_list(plan, plan.expr_list(args), schema)?;
Step::Function {
recipe: Recipe::new(plan.string(name), &self.literals(start, len)),
written: written(plan, expr, schema),
start,
len,
}
}
Expr::Aggregate { name, .. } => {
return Err(Error::internal(format!(
"the {} aggregate was evaluated as an ordinary expression",
plan.string(name)
)));
}
Expr::Window { name, .. } => {
return Err(Error::internal(format!(
"the {} window function was evaluated as an ordinary expression",
plan.string(name)
)));
}
Expr::Case { arms, otherwise } => {
let mut prepared = Vec::new();
for &arm in plan.arm_list(arms) {
prepared.push(PreparedArm {
when: Self::one(plan, arm.when, schema)?,
then: Self::one(plan, arm.then, schema)?,
});
}
let otherwise = match otherwise {
Some(otherwise) => Some(Self::one(plan, otherwise, schema)?),
None => None,
};
let blend = blending(&ty, &prepared, otherwise.as_ref());
Step::Case { arms: prepared, otherwise, blend }
}
};
Ok(self.place(plan, expr, step, ty))
}
fn place(&mut self, plan: &Plan, expr: ExprRef, step: Step, ty: LogicalType) -> usize {
self.steps.push(step);
self.types.push(ty);
self.spans.push(plan.expr_span(expr));
let step = self.steps.len() - 1;
if self.share {
self.shared.insert(expr, step);
}
step
}
fn push_list(
&mut self,
plan: &Plan,
exprs: &[ExprRef],
schema: &Schema,
) -> Result<(usize, usize)> {
let mut indices = Vec::with_capacity(exprs.len());
for &expr in exprs {
indices.push(self.push(plan, expr, schema)?);
}
let start = self.operands.len();
let len = indices.len();
self.operands.extend(indices);
Ok((start, len))
}
fn membership(
&mut self,
plan: &Plan,
op: Connective,
children: &[ExprRef],
schema: &Schema,
) -> Result<Option<Step>> {
let wanted = match op {
Connective::Or => CompareOp::Equal,
Connective::And => CompareOp::NotEqual,
};
let mut subject: Option<ExprRef> = None;
let mut values = Vec::with_capacity(children.len());
for &child in children {
let Expr::Compare { op: found, left, right } = *plan.expr(child) else {
return Ok(None);
};
if found != wanted || !same(plan, *subject.get_or_insert(left), left) {
return Ok(None);
}
let Expr::Constant(reference) = *plan.expr(right) else {
return Ok(None);
};
values.push(plan.value(reference).clone());
}
let (Some(subject), Some(members)) = (subject, Members::of(&values, op == Connective::And))
else {
return Ok(None);
};
Ok(Some(Step::InSet { input: self.push(plan, subject, schema)?, members }))
}
fn held(&self, left: usize, right: usize) -> Option<Held> {
let (at, other) = match (&self.steps[left], &self.steps[right]) {
(Step::Constant(_), Step::Constant(_)) => return None,
(_, Step::Constant(value)) => (right, value),
(Step::Constant(value), _) => (left, value),
_ => return None,
};
Held::of(&self.types[at], other)
}
fn literals(&self, start: usize, len: usize) -> Vec<Option<Value>> {
self.operands[start..start + len]
.iter()
.map(|&operand| match &self.steps[operand] {
Step::Constant(value) => Some(value.clone()),
_ => None,
})
.collect()
}
}
fn stamped_seconds(plan: &Plan, expr: ExprRef) -> Option<(ExprRef, ExprRef)> {
let Expr::Function { name, args } = *plan.expr(expr) else { return None };
if plan.string(name) != "+" || plan.expr_type(expr) != &LogicalType::Timestamp {
return None;
}
let &[one, other] = plan.expr_list(args) else { return None };
let (stamp, interval) =
if plan.expr_type(one) == &LogicalType::Timestamp { (one, other) } else { (other, one) };
if plan.expr_type(stamp) != &LogicalType::Timestamp {
return None;
}
let Expr::Function { name, args } = *plan.expr(interval) else { return None };
let &[cast] = plan.expr_list(args) else { return None };
let Expr::Cast { input, try_cast: false } = *plan.expr(cast) else { return None };
let whole = matches!(
plan.expr_type(input),
LogicalType::TinyInt
| LogicalType::SmallInt
| LogicalType::Integer
| LogicalType::BigInt
| LogicalType::UTinyInt
| LogicalType::USmallInt
| LogicalType::UInteger
);
(plan.string(name) == "to_seconds" && plan.expr_type(cast) == &LogicalType::Double && whole)
.then_some((stamp, input))
}
fn same(plan: &Plan, left: ExprRef, right: ExprRef) -> bool {
if left == right {
return true;
}
if plan.expr_type(left) != plan.expr_type(right) {
return false;
}
match (plan.expr(left), plan.expr(right)) {
(Expr::Column(one), Expr::Column(other)) => one == other,
(Expr::Constant(one), Expr::Constant(other)) => plan.value(*one) == plan.value(*other),
(
Expr::Cast { input: one, try_cast: first },
Expr::Cast { input: other, try_cast: second },
) => first == second && same(plan, *one, *other),
(
Expr::Function { name: one, args: first },
Expr::Function { name: other, args: second },
) => {
let (first, second) = (plan.expr_list(*first), plan.expr_list(*second));
plan.string(*one) == plan.string(*other)
&& first.len() == second.len()
&& first.iter().zip(second).all(|(&one, &other)| same(plan, one, other))
}
_ => false,
}
}
fn touching(ty: &LogicalType) -> f64 {
match ty.physical() {
PhysicalType::Varlen => 4.0,
PhysicalType::List | PhysicalType::Array | PhysicalType::Struct => 8.0,
_ => 1.0,
}
}
fn missing(index: usize) -> Error {
Error::internal(format!("step {index} was used as an operand before it produced anything"))
}
fn placed(rows: &[usize]) -> Result<Vec<u32>> {
rows.iter()
.map(|&row| {
u32::try_from(row).map_err(|_| Error::internal("a chunk of more than u32 rows"))
})
.collect()
}
fn blended(chunk: &Chunk, claimed: &[Vec<usize>], blend: &Blend) -> Result<Option<Vector>> {
let Some((dictionary, literals)) = agreed(chunk, blend)? else { return Ok(None) };
let mut codes = vec![0; chunk.len()];
for (branch, rows) in blend.branches.iter().zip(claimed) {
match *branch {
Branch::Column(position) => {
let Some((from, _)) = chunk.column(position)?.stable_dictionary_parts() else {
return Ok(None);
};
for &row in rows {
codes[row] = from[row];
}
}
Branch::Literal(at) => {
for &row in rows {
codes[row] = literals[at];
}
}
}
}
Vector::stable_dictionary(codes, dictionary).map(Some)
}
fn agreed(chunk: &Chunk, blend: &Blend) -> Result<Option<(Arc<Vector>, Vec<u32>)>> {
let mut held: Option<(&Vector, &Arc<Vector>)> = None;
for branch in &blend.branches {
let Branch::Column(position) = *branch else { continue };
let column = chunk.column(position)?;
let Some((_, dictionary)) = column.stable_dictionary_parts() else { return Ok(None) };
if column.validity().has_nulls(chunk.len()) {
return Ok(None);
}
match held {
Some((_, first)) if !Arc::ptr_eq(first, dictionary) => return Ok(None),
Some(_) => {}
None => held = Some((column, dictionary)),
}
}
let Some((column, dictionary)) = held else { return Ok(None) };
let mut codes = Vec::with_capacity(blend.literals.len());
for (text, lookup) in &blend.literals {
match lookup.find(column, text.as_bytes()) {
Some(Ok(Found::At(code))) => codes.push(code),
Some(Err(error)) => return Err(error),
Some(Ok(Found::Absent)) | None => return Ok(None),
}
}
Ok(Some((Arc::clone(dictionary), codes)))
}
fn blending(ty: &LogicalType, arms: &[PreparedArm], otherwise: Option<&Prepared>) -> Option<Blend> {
if !matches!(ty, LogicalType::Varchar) {
return None;
}
let otherwise = otherwise?;
let mut branches = Vec::with_capacity(arms.len() + 1);
let mut literals = Vec::new();
for branch in arms.iter().map(|arm| &arm.then).chain([otherwise]) {
branches.push(named(branch, &mut literals)?);
}
let any = branches.iter().any(|branch| matches!(branch, Branch::Column(_)));
any.then_some(Blend { branches, literals })
}
fn named(prepared: &Prepared, literals: &mut Vec<(String, Lookup)>) -> Option<Branch> {
match prepared.steps.as_slice() {
[Step::Column(position)] => Some(Branch::Column(*position)),
[Step::Constant(Value::Varchar(text))] => {
literals.push((text.clone(), Lookup::default()));
Some(Branch::Literal(literals.len() - 1))
}
_ => None,
}
}
pub(crate) fn narrow(chunk: &Chunk, rows: &[usize]) -> Result<Chunk> {
let mut selection = Selection::with_capacity(rows.len());
for &row in rows {
selection.push(row);
}
chunk.clone().select(&selection)
}
pub(crate) fn comparison(op: CompareOp) -> Comparison {
match op {
CompareOp::Equal => Comparison::Equal,
CompareOp::NotEqual => Comparison::NotEqual,
CompareOp::Less => Comparison::Less,
CompareOp::LessOrEqual => Comparison::LessOrEqual,
CompareOp::Greater => Comparison::Greater,
CompareOp::GreaterOrEqual => Comparison::GreaterOrEqual,
CompareOp::DistinctFrom => Comparison::DistinctFrom,
CompareOp::NotDistinctFrom => Comparison::NotDistinctFrom,
}
}
pub(crate) fn connective(op: ConjunctionOp) -> Connective {
match op {
ConjunctionOp::And => Connective::And,
ConjunctionOp::Or => Connective::Or,
}
}
#[cfg(test)]
mod tests {
use rudb_common::{Field, LogicalType, Value};
use rudb_kernels::is_true;
use rudb_plan::{ExprRef, Node, Plan};
use rudb_vector::{Chunk, Selection, Vector};
use super::{Prepared, narrow};
use crate::expr::evaluate;
use crate::schema::Schema;
fn input() -> (Schema, Chunk) {
let schema = Schema::numbered(
vec![Field::new("x", LogicalType::Integer), Field::new("s", LogicalType::Varchar)],
0,
);
let x = Vector::from_values(
LogicalType::Integer,
&[Value::Integer(3), Value::Integer(1), Value::Null, Value::Integer(2)],
)
.expect("four integers");
let s = Vector::from_values(
LogicalType::Varchar,
&[
Value::Varchar("a".to_string()),
Value::Null,
Value::Varchar("c".to_string()),
Value::Varchar("a".to_string()),
],
)
.expect("four strings");
(schema, Chunk::new(vec![x, s]).expect("two columns of four rows"))
}
fn projection(exprs: &str) -> (Plan, Vec<ExprRef>) {
let text =
format!("Project #1 [{exprs}]\n Get memory.main.t AS t #0 [x::INTEGER, s::VARCHAR]");
let plan = Plan::parse(&text).expect("a well formed plan");
let Node::Project { exprs, .. } = *plan.node(plan.root()) else {
panic!("the root of that text is a projection");
};
let list = plan.expr_list(exprs).to_vec();
(plan, list)
}
fn agrees(exprs: &str) {
let (schema, chunk) = input();
let (plan, list) = projection(exprs);
let prepared = Prepared::new(&plan, &list, &schema).expect("the expressions resolve");
let mut scratch = prepared.scratch();
let mut fast = Vec::new();
prepared.evaluate(&chunk, &mut scratch, &mut fast).expect("the prepared form runs");
for (at, &expr) in list.iter().enumerate() {
let slow = evaluate(&plan, expr, &schema, &chunk).expect("the tree walk runs");
for row in 0..chunk.len() {
assert_eq!(
fast[at].value_at(row),
slow.value_at(row),
"expression {at} of `{exprs}` at row {row}"
);
}
}
}
fn decimals(prices: &[i128], form: fn(Vector) -> Vector) -> (Schema, Chunk) {
let ty = LogicalType::Decimal { width: 15, scale: 2 };
let schema = Schema::numbered(
vec![
Field::new("p", ty.clone()),
Field::new("d", ty.clone()),
Field::new("t", ty.clone()),
],
0,
);
let column =
|values: Vec<Value>| form(Vector::from_values(ty.clone(), &values).expect("decimals"));
let decimal = |unscaled| Value::Decimal { unscaled, width: 15, scale: 2 };
let p = column(prices.iter().map(|&v| decimal(v)).collect());
let d = column((0..prices.len() as i128).map(|v| decimal(v % 11)).collect());
let t = column((0..prices.len() as i128).map(|v| decimal(v % 9)).collect());
(schema, Chunk::new(vec![p, d, t]).expect("three columns"))
}
const CHARGE: &str = "\"*\"(\"*\"(CAST(#0.0::DECIMAL(15,2))::DECIMAL(18,2), \
CAST(\"-\"(1.00::DECIMAL(16,2), CAST(#0.1::DECIMAL(15,2))::DECIMAL(16,2))::DECIMAL(16,2))\
::DECIMAL(18,2))::DECIMAL(18,4), CAST(\"+\"(1.00::DECIMAL(16,2), \
CAST(#0.2::DECIMAL(15,2))::DECIMAL(16,2))::DECIMAL(16,2))::DECIMAL(18,2))::DECIMAL(18,6) AS a";
fn three_ways(chunk: &Chunk, schema: &Schema) -> [rudb_common::Result<Vec<Value>>; 3] {
let text = format!(
"Project #1 [{CHARGE}]\n Get memory.main.t AS t #0 \
[p::DECIMAL(15,2), d::DECIMAL(15,2), t::DECIMAL(15,2)]"
);
let plan = Plan::parse(&text).expect("a well formed plan");
let Node::Project { exprs, .. } = *plan.node(plan.root()) else {
panic!("the root of that text is a projection");
};
let expr = plan.expr_list(exprs)[0];
let values = |vector: &Vector| (0..chunk.len()).map(|row| vector.value_at(row)).collect();
let fused = Prepared::one(&plan, expr, schema).expect("resolves");
assert_eq!(fused.fused(), 1, "the whole tree is one step");
let unfused = Prepared::built(&plan, &[expr], schema, false, false).expect("resolves");
assert_eq!(unfused.fused(), 0);
let run = |prepared: &Prepared| {
prepared.evaluate_one(chunk, &mut prepared.scratch()).map(&values)
};
[run(&fused), run(&unfused), evaluate(&plan, expr, schema, chunk).map(|v| values(&v))]
}
fn all_agree(chunk: &Chunk, schema: &Schema) {
let [fused, unfused, walked] = three_ways(chunk, schema);
let fused = fused.expect("fits");
assert_eq!(fused, unfused.expect("fits"));
assert_eq!(fused, walked.expect("fits"));
}
#[test]
fn a_timestamp_plus_whole_seconds_agrees_with_the_interval_it_stands_for() {
let schema = Schema::numbered(vec![Field::new("x", LogicalType::BigInt)], 0);
let counts = [
Value::BigInt(1_373_000_000),
Value::BigInt(-5),
Value::Null,
Value::BigInt(9_007_199_254),
Value::BigInt(9_007_199_255),
Value::BigInt(9_000_000_000_123),
];
let x = Vector::from_values(LogicalType::BigInt, &counts).expect("six counts");
let chunk = Chunk::new(vec![x]).expect("one column");
let text = "Project #1 [\"+\"(0::TIMESTAMP, to_seconds(CAST(#0.0::BIGINT)::DOUBLE)::INTERVAL)::TIMESTAMP AS e]\n Get memory.main.t AS t #0 [x::BIGINT]";
let plan = Plan::parse(text).expect("a well formed plan");
let Node::Project { exprs, .. } = *plan.node(plan.root()) else {
panic!("the root of that text is a projection");
};
let list = plan.expr_list(exprs).to_vec();
let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
assert!(
prepared.steps.iter().any(
|step| matches!(step, super::Step::Function { recipe, .. } if recipe.name() == "__rudb_stamp_seconds")
),
"the shift is one call"
);
let mut scratch = prepared.scratch();
let mut fast = Vec::new();
prepared.evaluate(&chunk, &mut scratch, &mut fast).expect("the prepared form runs");
let slow = evaluate(&plan, list[0], &schema, &chunk).expect("the tree walk runs");
for row in 0..chunk.len() {
assert_eq!(fast[0].value_at(row), slow.value_at(row), "row {row}");
}
assert_eq!(fast[0].value_at(0), Value::Timestamp(1_373_000_000_000_000));
let far = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(9_300_000_000_000)])
.expect("one count");
let chunk = Chunk::new(vec![far]).expect("one column");
let mut fast = Vec::new();
let fused = prepared.evaluate(&chunk, &mut scratch, &mut fast);
let slow = evaluate(&plan, list[0], &schema, &chunk).map(|_| ());
assert!(fused.is_err() && slow.is_err(), "past the last timestamp both raise");
}
#[test]
fn decimal_arithmetic_run_as_one_loop_agrees_in_every_form() {
let prices: Vec<i128> = (0..2500).map(|v| 90_000 + v * 37).collect();
let packed = |vector: Vector| vector.bit_packed().expect("packs");
let coded = |vector: Vector| {
let rows = vector.len();
let codes = (0..rows as u32).rev().collect();
Vector::dictionary(codes, vector.bit_packed().expect("packs")).expect("in range")
};
let scattered = |vector: Vector| {
let rows = vector.len() as u32;
let codes = (0..rows).map(|row| row * 997 % rows).collect();
Vector::dictionary(codes, vector.bit_packed().expect("packs")).expect("in range")
};
for form in [std::convert::identity, packed, coded, scattered] {
let (schema, chunk) = decimals(&prices, form);
all_agree(&chunk, &schema);
}
}
#[test]
fn a_chunk_the_ranges_cannot_prove_raises_what_the_steps_raise() {
let mut prices = vec![5; 300];
prices.push(999_999_999_999_999);
let packed = |vector: Vector| vector.bit_packed().expect("packs");
for form in [std::convert::identity, packed] {
let (schema, chunk) = decimals(&prices, form);
let [fused, unfused, _] = three_ways(&chunk, &schema);
let (fused, unfused) = (fused.expect_err("overflows"), unfused.expect_err("overflows"));
assert_eq!(fused.message(), unfused.message());
}
}
#[test]
fn a_chunk_with_a_null_goes_through_the_steps() {
let ty = LogicalType::Decimal { width: 15, scale: 2 };
let (schema, mut chunk) = decimals(&[100, 200, 300], std::convert::identity);
let with_null = Vector::from_values(
ty,
&[Value::Decimal { unscaled: 5, width: 15, scale: 2 }, Value::Null, Value::Null],
)
.expect("decimals");
chunk = Chunk::new(vec![
chunk.column(0).expect("p").clone(),
with_null,
chunk.column(2).expect("t").clone(),
])
.expect("three columns");
all_agree(&chunk, &schema);
}
#[test]
fn a_column_reference_agrees() {
agrees("#0.0::INTEGER AS a, #0.1::VARCHAR AS b");
}
#[test]
fn a_constant_agrees() {
agrees("7::INTEGER AS a, NULL::INTEGER AS b");
}
#[test]
fn a_cast_agrees() {
agrees("CAST(#0.0::INTEGER)::BIGINT AS a, CAST(#0.0::INTEGER)::VARCHAR AS b");
}
#[test]
fn a_comparison_agrees() {
agrees("(#0.0::INTEGER > 1::INTEGER)::BOOLEAN AS a");
}
#[test]
fn a_conjunction_agrees() {
agrees(
"((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 3::INTEGER)::BOOLEAN)\
::BOOLEAN AS a",
);
}
#[test]
fn a_function_agrees() {
agrees("\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER AS a");
}
#[test]
fn both_evaluators_quote_the_same_expression_when_a_divisor_is_zero() {
let (schema, chunk) = input();
let (plan, list) = projection("\"//\"(#0.0::INTEGER, 0::INTEGER)::INTEGER AS a");
let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
let mut scratch = prepared.scratch();
let mut out = Vec::new();
let fast = prepared.evaluate(&chunk, &mut scratch, &mut out).expect_err("divides by zero");
let slow = evaluate(&plan, list[0], &schema, &chunk).expect_err("divides by zero");
assert_eq!(fast.message(), slow.message());
assert!(fast.message().starts_with("Division by zero in expression (x // 0)."), "{fast}");
}
#[test]
fn a_case_agrees() {
agrees(
"CASE WHEN (#0.0::INTEGER > 1::INTEGER)::BOOLEAN THEN 10::INTEGER \
ELSE 20::INTEGER END::INTEGER AS a",
);
}
#[test]
fn a_case_of_two_arms_agrees() {
agrees(
"CASE WHEN (#0.0::INTEGER > 2::INTEGER)::BOOLEAN THEN 10::INTEGER \
WHEN (#0.0::INTEGER > 1::INTEGER)::BOOLEAN THEN 20::INTEGER \
ELSE 30::INTEGER END::INTEGER AS a",
);
}
#[test]
fn a_case_with_no_else_agrees() {
agrees(
"CASE WHEN (#0.0::INTEGER > 2::INTEGER)::BOOLEAN THEN 10::INTEGER \
END::INTEGER AS a",
);
}
#[test]
fn a_case_whose_arm_claims_nothing_agrees() {
agrees(
"CASE WHEN (#0.0::INTEGER > 99::INTEGER)::BOOLEAN THEN 10::INTEGER \
ELSE 20::INTEGER END::INTEGER AS a",
);
}
#[test]
fn a_case_over_strings_agrees() {
agrees(
"CASE WHEN (#0.0::INTEGER > 1::INTEGER)::BOOLEAN THEN #0.1::VARCHAR \
ELSE ''::VARCHAR END::VARCHAR AS a",
);
}
#[test]
fn a_case_whose_arm_answers_null_agrees() {
agrees(
"CASE WHEN (#0.0::INTEGER > 1::INTEGER)::BOOLEAN THEN #0.1::VARCHAR \
ELSE NULL::VARCHAR END::VARCHAR AS a",
);
}
#[test]
fn a_case_whose_test_is_null_on_some_rows_agrees() {
agrees(
"CASE WHEN (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN THEN 10::INTEGER \
ELSE 20::INTEGER END::INTEGER AS a",
);
}
#[test]
fn a_column_mentioned_three_times_agrees() {
agrees("\"+\"(\"+\"(#0.0::INTEGER, #0.0::INTEGER)::INTEGER, #0.0::INTEGER)::INTEGER AS a");
}
#[test]
fn a_chain_holds_one_intermediate_at_a_time() {
let (schema, chunk) = input();
let mut expr = "#0.0::INTEGER".to_string();
for _ in 0..8 {
expr = format!("\"+\"({expr}, 1::INTEGER)::INTEGER");
}
let (plan, list) = projection(&format!("{expr} AS a"));
let prepared = Prepared::new(&plan, &list, &schema).expect("the chain resolves");
let mut scratch = prepared.scratch();
prepared.run(&chunk, &mut scratch).expect("the chain runs");
let live = scratch.slots.iter().filter(|slot| slot.is_some()).count();
assert_eq!(live, 1, "a chain that has run should be holding its answer and nothing else");
}
fn filters(predicate: &str) {
let (schema, chunk) = input();
let (plan, list) = projection(&format!("{predicate} AS p"));
let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
let mut scratch = prepared.scratch();
let threaded = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs");
let flags = evaluate(&plan, list[0], &schema, &chunk).expect("the tree walk runs");
let expected = Selection::from_predicate(chunk.len(), |row| is_true(&flags.value_at(row)));
assert_eq!(threaded, expected, "`{predicate}`");
let again = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs again");
assert_eq!(again, expected, "`{predicate}` a second time");
}
#[test]
fn a_single_comparison_filters_the_same_rows() {
filters("(#0.0::INTEGER > 1::INTEGER)::BOOLEAN");
filters("(#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN");
filters("(#0.0::INTEGER IS NOT DISTINCT FROM NULL::INTEGER)::BOOLEAN");
}
#[test]
fn a_chain_of_conjuncts_keeps_what_all_of_them_keep() {
filters(
"((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 3::INTEGER)::BOOLEAN)\
::BOOLEAN",
);
filters(
"((#0.0::INTEGER >= 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER <= 3::INTEGER)::BOOLEAN \
AND (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN AND (#0.0::INTEGER <> 2::INTEGER)\
::BOOLEAN)::BOOLEAN",
);
}
#[test]
fn a_settled_conjunct_is_left_out_of_the_filter() {
let (schema, chunk) = input();
let both = "((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 3::INTEGER)::BOOLEAN)\
::BOOLEAN AS p";
let (plan, list) = projection(both);
let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
assert_eq!(prepared.conjuncts(), Some(2));
let mut scratch = prepared.scratch();
let wanted = |predicate: &str| {
let (plan, list) = projection(&format!("{predicate} AS p"));
let flags = evaluate(&plan, list[0], &schema, &chunk).expect("the tree walk runs");
Selection::from_predicate(chunk.len(), |row| is_true(&flags.value_at(row)))
};
let second = prepared.evaluate_settled(&chunk, &mut scratch, &[true, false]);
assert_eq!(
second.expect("the filter runs"),
wanted("(#0.0::INTEGER < 3::INTEGER)::BOOLEAN")
);
let first = prepared.evaluate_settled(&chunk, &mut scratch, &[false, true]);
assert_eq!(
first.expect("the filter runs"),
wanted("(#0.0::INTEGER > 1::INTEGER)::BOOLEAN")
);
let neither = prepared.evaluate_settled(&chunk, &mut scratch, &[true, true]);
assert_eq!(neither.expect("the filter runs"), Selection::identity(chunk.len()));
let whole = wanted(&both[..both.len() - " AS p".len()]);
let none = prepared.evaluate_settled(&chunk, &mut scratch, &[false, false]);
assert_eq!(none.expect("the filter runs"), whole);
let short = prepared.evaluate_settled(&chunk, &mut scratch, &[true]);
assert_eq!(short.expect("the filter runs"), whole, "a list that does not fit is ignored");
}
#[test]
fn a_conjunct_that_keeps_nothing_ends_the_predicate() {
filters(
"((#0.0::INTEGER > 9::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 9::INTEGER)::BOOLEAN)\
::BOOLEAN",
);
}
#[test]
fn a_conjunct_over_a_computed_operand_keeps_the_same_rows() {
filters(
"((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND \
(\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER < 4::INTEGER)::BOOLEAN)::BOOLEAN",
);
}
#[test]
fn a_conjunct_that_is_not_a_comparison_is_threaded_too() {
filters(
"((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND ((#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN \
OR (#0.0::INTEGER = 1::INTEGER)::BOOLEAN)::BOOLEAN)::BOOLEAN",
);
filters(
"(((#0.1::VARCHAR = 'c'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
::BOOLEAN AND (#0.0::INTEGER <> 1::INTEGER)::BOOLEAN)::BOOLEAN",
);
}
#[test]
fn a_selective_conjunct_evaluates_later_like_on_its_survivors() {
filters(
"((#0.0::INTEGER > 2::INTEGER)::BOOLEAN AND \
\"~~\"(#0.1::VARCHAR, '%a%'::VARCHAR)::BOOLEAN)::BOOLEAN",
);
filters(
"((#0.0::INTEGER > 2::INTEGER)::BOOLEAN AND \
\"!~~\"(#0.1::VARCHAR, '%a%'::VARCHAR)::BOOLEAN)::BOOLEAN",
);
}
#[test]
fn an_or_at_the_top_threads_the_complement() {
filters(
"((#0.0::INTEGER > 2::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'c'::VARCHAR)::BOOLEAN)\
::BOOLEAN",
);
filters(
"((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN \
OR (#0.0::INTEGER > 2::INTEGER)::BOOLEAN)::BOOLEAN",
);
}
#[test]
fn a_branch_that_keeps_everything_ends_the_predicate() {
filters(
"((#0.0::INTEGER IS NOT DISTINCT FROM #0.0::INTEGER)::BOOLEAN OR \
(#0.0::INTEGER > 9::INTEGER)::BOOLEAN)::BOOLEAN",
);
}
#[test]
fn a_branch_behind_one_that_accepted_every_row_does_not_run() {
let (schema, chunk) = input();
let predicate = "((#0.0::INTEGER IS NOT DISTINCT FROM #0.0::INTEGER)::BOOLEAN OR \
(\"//\"(#0.0::INTEGER, 0::INTEGER)::INTEGER > 0::INTEGER)::BOOLEAN)\
::BOOLEAN";
let (plan, list) = projection(&format!("{predicate} AS p"));
let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
let mut scratch = prepared.scratch();
let kept =
prepared.evaluate_filter(&chunk, &mut scratch).expect("the second branch never runs");
assert_eq!(kept, Selection::identity(chunk.len()));
evaluate(&plan, list[0], &schema, &chunk).expect_err("the tree walk divides by zero");
}
#[test]
fn a_filter_learns_which_conjunct_to_run_first() {
let (schema, chunk) = input();
let predicate = "((#0.0::INTEGER > 0::INTEGER)::BOOLEAN AND (#0.0::INTEGER > 9::INTEGER)\
::BOOLEAN)::BOOLEAN";
let (plan, list) = projection(&format!("{predicate} AS p"));
let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
let mut scratch = prepared.scratch();
let root = prepared.roots[0];
assert_eq!(scratch.order(root), None, "nothing has run yet");
let kept = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs");
assert!(kept.is_empty());
assert_eq!(scratch.order(root), Some(&[1, 0][..]), "the second conjunct rejects the most");
let kept = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs again");
assert!(kept.is_empty());
assert_eq!(scratch.order(root), Some(&[1, 0][..]));
}
#[test]
fn reordering_never_changes_which_rows_survive() {
let (schema, chunk) = input();
let predicate = "((#0.0::INTEGER >= 1::INTEGER)::BOOLEAN AND \
(\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER < 4::INTEGER)::BOOLEAN AND \
(#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)::BOOLEAN";
let (plan, list) = projection(&format!("{predicate} AS p"));
let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
let mut scratch = prepared.scratch();
let flags = evaluate(&plan, list[0], &schema, &chunk).expect("the tree walk runs");
let expected = Selection::from_predicate(chunk.len(), |row| is_true(&flags.value_at(row)));
for round in 0..40 {
let kept = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs");
assert_eq!(kept, expected, "round {round}");
}
}
#[test]
fn a_nested_connective_stops_where_the_outer_one_would() {
let (schema, chunk) = input();
let predicate = "((#0.0::INTEGER > 9::INTEGER)::BOOLEAN OR ((#0.0::INTEGER > 9::INTEGER)\
::BOOLEAN AND (\"//\"(#0.0::INTEGER, 0::INTEGER)::INTEGER > 0::INTEGER)\
::BOOLEAN)::BOOLEAN)::BOOLEAN";
let (plan, list) = projection(&format!("{predicate} AS p"));
let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
let mut scratch = prepared.scratch();
let kept =
prepared.evaluate_filter(&chunk, &mut scratch).expect("the division never happens");
assert!(kept.is_empty());
evaluate(&plan, list[0], &schema, &chunk).expect_err("the tree walk divides by zero");
}
#[test]
fn an_or_branch_that_is_not_a_comparison_is_threaded_too() {
filters(
"((#0.0::INTEGER > 2::INTEGER)::BOOLEAN OR \
\"~~\"(#0.1::VARCHAR, 'a%'::VARCHAR)::BOOLEAN)::BOOLEAN",
);
filters(
"(\"~~\"(#0.1::VARCHAR, 'c%'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 1::INTEGER)\
::BOOLEAN)::BOOLEAN",
);
}
#[test]
fn a_connective_inside_a_connective_threads_both_ways() {
filters(
"(((#0.0::INTEGER >= 2::INTEGER)::BOOLEAN AND (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)\
::BOOLEAN OR ((#0.0::INTEGER < 2::INTEGER)::BOOLEAN AND (#0.1::VARCHAR <> 'c'\
::VARCHAR)::BOOLEAN)::BOOLEAN)::BOOLEAN",
);
filters(
"(((#0.1::VARCHAR = 'c'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
::BOOLEAN AND ((#0.0::INTEGER <> 1::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'a'\
::VARCHAR)::BOOLEAN)::BOOLEAN)::BOOLEAN",
);
filters(
"((#0.0::INTEGER > 9::INTEGER)::BOOLEAN OR ((#0.0::INTEGER >= 1::INTEGER)::BOOLEAN \
AND ((#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 1::INTEGER)\
::BOOLEAN)::BOOLEAN)::BOOLEAN)::BOOLEAN",
);
}
#[test]
fn a_null_branch_beside_a_true_one_keeps_the_row() {
filters(
"((#0.0::INTEGER > 2::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'c'::VARCHAR)::BOOLEAN \
OR (#0.0::INTEGER IS NOT DISTINCT FROM NULL::INTEGER)::BOOLEAN)::BOOLEAN",
);
filters(
"((#0.1::VARCHAR > 'b'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 1::INTEGER)::BOOLEAN)\
::BOOLEAN",
);
}
#[test]
fn a_filter_over_a_selected_chunk_keeps_the_same_rows() {
let (schema, chunk) = input();
let predicate = "((#0.0::INTEGER >= 1::INTEGER)::BOOLEAN AND \
(#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)::BOOLEAN";
let (plan, list) = projection(&format!("{predicate} AS p"));
let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
let mut scratch = prepared.scratch();
let narrowed = narrow(&chunk, &[0, 3]).expect("two of the four rows");
let threaded = prepared.evaluate_filter(&narrowed, &mut scratch).expect("the filter runs");
let flags = evaluate(&plan, list[0], &schema, &narrowed).expect("the tree walk runs");
let expected =
Selection::from_predicate(narrowed.len(), |row| is_true(&flags.value_at(row)));
assert_eq!(threaded, expected);
}
#[test]
fn a_scratch_used_twice_gives_the_same_answer_twice() {
let (schema, chunk) = input();
let (plan, list) = projection("\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER AS a");
let prepared = Prepared::new(&plan, &list, &schema).expect("the expressions resolve");
let mut scratch = prepared.scratch();
let mut once = Vec::new();
prepared.evaluate(&chunk, &mut scratch, &mut once).expect("the first chunk runs");
let mut twice = Vec::new();
prepared.evaluate(&chunk, &mut scratch, &mut twice).expect("the second chunk runs");
assert_eq!(once, twice);
}
#[test]
fn taking_the_chunk_answers_what_borrowing_it_does() {
let (schema, chunk) = input();
let (plan, list) = projection(
"#0.0::INTEGER AS a, \"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER AS b, #0.0::INTEGER AS c",
);
let prepared = Prepared::new(&plan, &list, &schema).expect("the expressions resolve");
let mut scratch = prepared.scratch();
let mut borrowed = Vec::new();
prepared.evaluate(&chunk, &mut scratch, &mut borrowed).expect("the borrowed chunk runs");
let mut taken = Vec::new();
prepared.evaluate_taking(chunk, &mut scratch, &mut taken).expect("the taken chunk runs");
assert_eq!(borrowed, taken);
}
#[test]
fn a_shared_computed_root_is_compiled_once() {
let (schema, chunk) = input();
let (plan, list) = projection("\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER AS a");
let prepared = Prepared::shared(&plan, &[list[0], list[0]], &schema)
.expect("the shared expression resolves");
assert_eq!(prepared.steps.len(), 3);
let mut scratch = prepared.scratch();
let mut answers = Vec::new();
prepared.evaluate(&chunk, &mut scratch, &mut answers).expect("both roots are returned");
assert_eq!(answers[0], answers[1]);
}
#[test]
fn a_shorter_chunk_after_a_longer_one_is_evaluated_at_its_own_length() {
let (schema, chunk) = input();
let (plan, list) = projection("7::INTEGER AS a");
let prepared = Prepared::new(&plan, &list, &schema).expect("the expressions resolve");
let mut scratch = prepared.scratch();
let mut full = Vec::new();
prepared.evaluate(&chunk, &mut scratch, &mut full).expect("the full chunk runs");
assert_eq!(full[0].len(), 4);
let short = chunk
.clone()
.select(&{
let mut selection = Selection::with_capacity(2);
selection.push(0);
selection.push(2);
selection
})
.expect("two of the four rows");
let mut cut = Vec::new();
prepared.evaluate(&short, &mut scratch, &mut cut).expect("the short chunk runs");
assert_eq!(cut[0].len(), 2);
}
#[test]
fn an_aggregate_is_refused_when_it_is_prepared() {
let (schema, _) = input();
let text = "Aggregate #1 groups=[] aggregates=[sum(#0.0::INTEGER)::HUGEINT]\n \
Get memory.main.t AS t #0 [x::INTEGER, s::VARCHAR]";
let plan = Plan::parse(text).expect("a well formed plan");
let Node::Aggregate { aggregates, .. } = *plan.node(plan.root()) else {
panic!("the root of that text is an aggregate");
};
let list = plan.expr_list(aggregates).to_vec();
let error = Prepared::new(&plan, &list, &schema).expect_err("sum is not a scalar");
assert!(error.message().contains("sum"), "{error}");
}
fn prepares(expr: &str, lifted: usize) {
let (schema, _) = input();
let projected = format!("{expr} AS a");
let (plan, list) = projection(&projected);
let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
assert_eq!(prepared.hoisted(), lifted, "`{expr}`");
agrees(&projected);
}
#[test]
fn a_literal_pattern_is_compiled_when_the_pipeline_is_built() {
prepares("\"~~\"(#0.1::VARCHAR, 'a%'::VARCHAR)::BOOLEAN", 1);
prepares("\"~~*\"(#0.1::VARCHAR, '%A%'::VARCHAR)::BOOLEAN", 1);
}
#[test]
fn a_regular_expression_is_compiled_when_the_pipeline_is_built() {
prepares("\"regexp_matches\"(#0.1::VARCHAR, '^a'::VARCHAR)::BOOLEAN", 1);
prepares("\"regexp_replace\"(#0.1::VARCHAR, 'a'::VARCHAR, 'b'::VARCHAR)::VARCHAR", 1);
}
#[test]
fn a_pattern_that_is_not_a_literal_is_left_to_the_chunk() {
prepares("\"~~\"(#0.1::VARCHAR, #0.1::VARCHAR)::BOOLEAN", 0);
}
#[test]
fn a_function_with_no_prepare_step_prepares_nothing() {
prepares("\"upper\"(#0.1::VARCHAR)::VARCHAR", 0);
}
fn folds(expr: &str, sets: usize) {
let (schema, _) = input();
let projected = format!("{expr} AS a");
let (plan, list) = projection(&projected);
let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
assert_eq!(prepared.sets(), sets, "`{expr}`");
agrees(&projected);
}
#[test]
fn an_in_list_becomes_one_lookup() {
folds(
"((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
::BOOLEAN",
1,
);
folds(
"((#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN OR (#0.1::VARCHAR = 'z'::VARCHAR)::BOOLEAN)\
::BOOLEAN",
1,
);
}
#[test]
fn a_not_in_list_becomes_the_same_lookup() {
folds(
"((#0.0::INTEGER <> 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER <> 3::INTEGER)::BOOLEAN)\
::BOOLEAN",
1,
);
}
#[test]
fn a_list_with_a_null_in_it_folds_and_keeps_the_null_rule() {
folds(
"((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = NULL::INTEGER)::BOOLEAN \
OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)::BOOLEAN",
1,
);
folds(
"((#0.0::INTEGER <> 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER <> NULL::INTEGER)\
::BOOLEAN AND (#0.0::INTEGER <> 3::INTEGER)::BOOLEAN)::BOOLEAN",
1,
);
}
#[test]
fn a_connective_that_is_not_an_in_list_is_left_alone() {
folds(
"((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)\
::BOOLEAN",
0,
);
folds(
"((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER > 3::INTEGER)::BOOLEAN)\
::BOOLEAN",
0,
);
folds(
"((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = #0.0::INTEGER)::BOOLEAN)\
::BOOLEAN",
0,
);
folds(
"((#0.0::INTEGER = 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
::BOOLEAN",
0,
);
}
#[test]
fn an_in_list_filters_the_same_rows() {
filters(
"((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
::BOOLEAN",
);
filters(
"((#0.0::INTEGER <> 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER <> 3::INTEGER)::BOOLEAN)\
::BOOLEAN",
);
filters(
"(((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
::BOOLEAN AND (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)::BOOLEAN",
);
}
#[test]
fn a_comparison_against_a_literal_builds_it_once() {
let (schema, _) = input();
for (expr, built) in [
("(#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN AS p", 1),
("(#0.0::INTEGER > 1::INTEGER)::BOOLEAN AS p", 1),
("(1::INTEGER < #0.0::INTEGER)::BOOLEAN AS p", 1),
("(#0.0::INTEGER = #0.0::INTEGER)::BOOLEAN AS p", 0),
("(1::INTEGER = 2::INTEGER)::BOOLEAN AS p", 0),
] {
let (plan, list) = projection(expr);
let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
assert_eq!(prepared.literals_built(), built, "`{expr}`");
agrees(expr);
}
}
#[test]
fn a_pattern_that_does_not_compile_fails_on_the_chunk_and_not_before() {
let (schema, chunk) = input();
let (plan, list) =
projection("\"regexp_matches\"(#0.1::VARCHAR, 'a('::VARCHAR)::BOOLEAN AS a");
let prepared = Prepared::new(&plan, &list, &schema).expect("preparing does not compile it");
assert_eq!(prepared.hoisted(), 0);
let mut scratch = prepared.scratch();
let mut out = Vec::new();
prepared.evaluate(&chunk, &mut scratch, &mut out).expect_err("the chunk raises");
}
}