rudb_exec/prepared.rs
1//! An expression prepared once for a pipeline and then evaluated over every chunk.
2//!
3//! `spec/engine/04-expressions.md`. [`evaluate`](crate::evaluate) walks the plan's expression tree
4//! on every chunk, which means it does four things per chunk that depend on nothing about the
5//! chunk: it recurses, it resolves every column reference by a linear search through the schema, it
6//! clones a [`LogicalType`] for every node, and it copies the whole column a [`Expr::Column`] names.
7//! Over `hits` at a hundred thousand chunks that is a hundred thousand schema searches per column
8//! reference and a hundred thousand copies of every column any expression mentions.
9//!
10//! This type does all four once. The tree is flattened into a post order array, so evaluating it is
11//! a loop over that array and the recursion is gone with it. Column references are resolved to
12//! positions when the pipeline is built. Types are held here rather than cloned out of the plan.
13//! And a column reference is not a step that produces anything: it is read straight out of the chunk
14//! at the point an operand is wanted, so the column is never copied at all.
15//!
16//! # What is shared and what is not
17//!
18//! [`Prepared`] is immutable after it is built and is `Send` and `Sync`, so one of them serves every
19//! thread running a copy of the pipeline. [`Scratch`] is the per chunk working space and there is
20//! one per pipeline instance. That split is not for this layer's benefit. It is the same split every
21//! operator needs at layer eight, where the scheduler runs one pipeline on as many threads as it has
22//! morsels for, and building it here means the operators above are written against it from the start
23//! rather than retrofitted onto it.
24//!
25//! # What is still allocated per chunk
26//!
27//! Two things, and both are named rather than hidden. A node with four or more operands gathers
28//! references to them into a `Vec<&Vector>` so a kernel can take a slice, which is one allocation of
29//! pointers rather than a copy of any data, and which a node of one, two or three operands does on
30//! the stack instead. And every kernel allocates the vector it returns, because no kernel in
31//! `rudb-kernels` takes an output parameter. The second is much the larger of the two and it is the
32//! one tier 1 fusion removes, which is scheduled after layer six for the reason
33//! `spec/engine/04-expressions.md` gives: once the tree walk is gone what is left to save is pass
34//! count, and at 1024 rows the intermediate vectors are eight kilobytes and stay in L1.
35
36use rudb_common::{
37 Error, ErrorCode, LogicalType, PhysicalType, Result, Session, SessionTimeZone, Span, Value,
38};
39use rudb_kernels::{
40 Comparison, Connective, Found, Held, Lookup, Members, Recipe, cast_in_time_zone, combine,
41 compare_prepared, in_set, is_true, refine_flags, refine_prepared, select_prepared, selection,
42};
43use rudb_plan::{CompareOp, ConjunctionOp, Expr, ExprRef, Plan};
44use rudb_vector::{Assembly, Chunk, Selection, Vector};
45use std::collections::HashMap;
46use std::sync::Arc;
47
48use crate::fused::Fused;
49use crate::lambda::{Lambda, lambda_call};
50use crate::ordering::Ordering;
51use crate::schema::Schema;
52use crate::written::written;
53
54/// The scheduler's half of the expression contract, imposed now rather than at layer eight.
55///
56/// A prepared expression is the immutable half of a pipeline and layer eight hands one of them to
57/// every thread running that pipeline. That is only sound if it holds nothing thread local, and the
58/// way to find out on the commit that breaks it rather than eight layers later is to ask the
59/// compiler here, exactly as [`Chunk`] does for the data plane.
60const _: () = {
61 const fn assert_shareable<T: Send + Sync>() {}
62 assert_shareable::<Prepared>();
63};
64
65/// One or more bound expressions, flattened and resolved against a schema.
66///
67/// Built once per pipeline with [`Prepared::new`] and evaluated per chunk with
68/// [`Prepared::evaluate`] or [`Prepared::evaluate_one`], each of which wants the [`Scratch`] that
69/// [`Prepared::scratch`] hands out.
70#[derive(Debug)]
71pub struct Prepared {
72 /// The nodes in post order, so every node's operands have already been computed when it runs.
73 steps: Vec<Step>,
74 /// The type each step produces, indexed the same way as `steps`.
75 ///
76 /// A parallel array rather than a field in the variant, for the reason [`Expr`] gives: a
77 /// [`LogicalType`] owns a `Vec` for its nested cases and putting one in every variant would make
78 /// the common variants several times larger for the benefit of the rare ones.
79 types: Vec<LogicalType>,
80 /// The source range each step came from, indexed the same way as `steps`.
81 spans: Vec<Span>,
82 /// The operand lists of the steps that have one, as runs of step indices.
83 operands: Vec<usize>,
84 /// The last step that reads each step's slot, or `usize::MAX` for one nothing reads.
85 ///
86 /// A slot is emptied as soon as the step that was the last to read it has run. Keeping every
87 /// intermediate alive to the end of the array instead is what the first measured version of this
88 /// did, and a chain of eight additions was slower prepared than walked because of it: nine live
89 /// intermediates at eight kilobytes each is seventy two kilobytes of working set where the tree
90 /// walk has two, and two is the pair the allocator hands back and forth and that stays in L1.
91 /// Everything else about the prepared form was faster and this one thing paid all of it back.
92 last_use: Vec<usize>,
93 /// The step index each expression this was built from ends at.
94 roots: Vec<usize>,
95 /// Whether each entry of [`roots`](Self::roots) is the last one naming its step.
96 ///
97 /// Two expressions of one projection can end at the same step, because a shared subexpression is
98 /// compiled once, and then the first of them has to copy the answer and the last of them can
99 /// take it. Which is which is a property of `roots` alone, so it is settled here rather than
100 /// counted again on every chunk. Counting it per chunk is what [`evaluate`](Self::evaluate) used
101 /// to do, through a `HashMap` it allocated and hashed every call, and on TPC-H Q1 that map was a
102 /// measurable part of the query for an answer that never changed.
103 last_root: Vec<bool>,
104 /// The step already compiled for each shared plan expression.
105 shared: HashMap<ExprRef, usize>,
106 share: bool,
107 /// Whether a tree of decimal arithmetic is run as one [`Fused`] step. Off only for the steps a
108 /// fused one falls back to, which would otherwise fuse themselves again.
109 fuse: bool,
110 /// The parsed zone used only by casts whose answer depends on the session.
111 time_zone: SessionTimeZone,
112}
113
114/// One node of a flattened expression.
115///
116/// A step refers to its operands by their index in [`Prepared::steps`], which is always smaller than
117/// its own because the array is in post order.
118#[derive(Debug)]
119enum Step {
120 /// A column of the chunk, by resolved position.
121 ///
122 /// This step computes nothing. Its slot stays empty and an operand that names it is read out of
123 /// the chunk, which is the whole of what makes a column reference free rather than a copy.
124 Column(usize),
125 /// A literal, materialized into a constant vector as long as the chunk.
126 Constant(Value),
127 /// A cast to this step's own type.
128 Cast {
129 /// The step being cast.
130 input: usize,
131 /// Whether a failed cast yields null instead of raising.
132 try_cast: bool,
133 },
134 /// A binary comparison.
135 Compare {
136 /// Which comparison.
137 op: Comparison,
138 /// The left operand's step.
139 left: usize,
140 /// The right operand's step.
141 right: usize,
142 /// The side that is a literal, in the one row column the comparison loops read it through,
143 /// and `None` when neither side is one.
144 ///
145 /// Built here because the loops read both sides through a slice, so the constant side has
146 /// to become a column somewhere, and the plan says which side that is. For a string it is
147 /// also where the four byte prefix comes from, which is what almost every row of a string
148 /// comparison is decided by.
149 held: Option<Held>,
150 },
151 /// An `AND` or `OR` over a run of [`Prepared::operands`].
152 Conjunction {
153 /// Which connective.
154 op: Connective,
155 /// Where the operand list starts.
156 start: usize,
157 /// How many operands it has.
158 len: usize,
159 },
160 /// A scalar function over a run of [`Prepared::operands`].
161 Function {
162 /// The call, with the name resolved and whatever the kernel could work out from the
163 /// arguments that were literals already worked out.
164 ///
165 /// Held here so the plan is not consulted per chunk, and built here so that a regular
166 /// expression is compiled once for the query rather than once for each of the hundred
167 /// thousand chunks a pipeline over `hits` runs.
168 recipe: Recipe,
169 /// How the call is written, for the one error message that quotes it.
170 ///
171 /// Rendered when the pipeline is built rather than when a chunk arrives, because the plan
172 /// is here and is not there. It is a short string per function node in the query and it is
173 /// built once, which is a different cost from the tree walk's, where the plan is still to
174 /// hand and the rendering can wait until the row that fails.
175 written: String,
176 /// Where the argument list starts.
177 start: usize,
178 /// How many arguments it has.
179 len: usize,
180 },
181 /// A membership test over a list the query wrote out.
182 ///
183 /// The binder has no `IN` node: `x IN (1, 2, 3)` arrives as an `OR` of three equalities and
184 /// `x NOT IN (1, 2, 3)` as an `AND` of three inequalities. That is the right shape for a binder
185 /// to produce, because nothing after it then needs a second set of rules for null, and it is the
186 /// wrong shape to run, because it is a pass over the column and an output vector per entry.
187 /// This is that shape folded back up, and folding it here rather than after the operands are
188 /// pushed is what keeps the equalities from being run anyway.
189 InSet {
190 /// The step being tested.
191 input: usize,
192 /// The list, as a set, with the null rule and the direction it is read in.
193 members: Members,
194 },
195 /// A searched `CASE`, whose branches are prepared expressions of their own.
196 ///
197 /// Nested rather than flattened into the same array because a branch is not evaluated over the
198 /// chunk, it is evaluated over the rows no earlier arm claimed, and a step in the outer array
199 /// would have no way to say that. The selection threaded form in #57 replaces this whole
200 /// variant, and when it does the branches stop being separate arrays.
201 Case {
202 /// The `WHEN`/`THEN` pairs, in order.
203 arms: Vec<PreparedArm>,
204 /// The `ELSE`, if there is one. Absent means null.
205 otherwise: Option<Prepared>,
206 /// How to answer it as codes, for the shape that can be. Absent means read the values.
207 blend: Option<Blend>,
208 },
209 /// `TRY(x)`, whose operand is a prepared expression of its own because it may have to be run
210 /// again one row at a time.
211 Try {
212 /// The operand.
213 inner: Box<Prepared>,
214 },
215 /// A tree of decimal arithmetic over columns and literals, run as one loop when the columns'
216 /// ranges prove it cannot overflow.
217 ///
218 /// The fallback is the same tree prepared the ordinary way, nested for the reason a case's
219 /// branches are, and it is what runs over a chunk the ranges do not settle.
220 Fused {
221 /// The program.
222 fused: Box<Fused>,
223 /// The steps it replaced.
224 fallback: Box<Prepared>,
225 },
226 /// A call to a function that takes a lambda, whose body is a prepared expression of its own.
227 ///
228 /// Nested for the reason a case's branches are: the body does not run over the chunk, it runs
229 /// over a chunk with a row per element that [`Lambda`] builds, and a step in the outer array has
230 /// no way to say that.
231 Lambda {
232 /// The steps of the call's other arguments: the list and `list_reduce`'s initial value, or
233 /// `invoke`'s parameters.
234 inputs: Vec<usize>,
235 /// The layout of what the body runs over and what to do with its answers.
236 runner: Box<Lambda>,
237 /// The body, prepared against the runner's schema.
238 body: Box<Prepared>,
239 },
240}
241
242/// One `WHEN`/`THEN` pair of a prepared [`Step::Case`].
243#[derive(Debug)]
244struct PreparedArm {
245 /// The condition.
246 when: Prepared,
247 /// The result if the condition is true.
248 then: Prepared,
249}
250
251/// A `CASE` over text whose every branch is a column or a literal, answered as codes.
252///
253/// What the general path does with the branches is read their values and write them into a vector of
254/// their own, which for a text column out of a native file decodes a compressed dictionary block per
255/// row and then throws the dictionary away. An operator above that has to work with strings even
256/// though every string it sees came out of one dictionary it could have kept.
257///
258/// It does not have to. The branches here name values rather than compute them, so if they all name
259/// values of one dictionary then so does the answer, and the answer is the codes: one code per row
260/// copied from the branch that claimed the row, and a literal is one code for all of its rows once
261/// the dictionary has been searched for it. Nothing is read and the dictionary comes out the other
262/// side, so a group by over the `CASE` groups on codes the way a group by over the bare column does.
263///
264/// ClickBench 39 is the query this is for. It groups by `CASE WHEN (SearchEngineID = 0 AND
265/// AdvEngineID = 0) THEN Referer ELSE '' END` beside `URL`, and writing that one column out as
266/// strings was a quarter of the query.
267///
268/// The shape is narrow on purpose. A branch that computes anything is not here, because then the
269/// answer is a value that no dictionary holds. A literal the dictionary does not hold is not here
270/// either, for the same reason, and that is decided per dictionary at run time rather than when the
271/// expression is prepared. And a `CASE` with no `ELSE` is not here, because the rows nothing claims
272/// are null and a null is not a code.
273#[derive(Debug)]
274struct Blend {
275 /// Where each branch takes its value from: one per arm in order, and the `ELSE` last.
276 branches: Vec<Branch>,
277 /// The literals the branches name, each with the search that finds it in a dictionary.
278 literals: Vec<(String, Lookup)>,
279}
280
281/// Where one branch of a [`Blend`] takes its value from.
282#[derive(Debug, Clone, Copy)]
283enum Branch {
284 /// A column of the chunk, by resolved position. Its rows keep the codes they arrived with.
285 Column(usize),
286 /// The literal at this index of [`Blend::literals`]. Its rows all get one code.
287 Literal(usize),
288}
289
290/// The per chunk working space of one [`Prepared`].
291///
292/// One per pipeline instance and never shared, which is the mutable half of the split the module
293/// documentation describes. It is handed back in rather than made inside [`Prepared::evaluate`] so
294/// that the array of slots survives from one chunk to the next instead of being allocated a hundred
295/// thousand times over a scan.
296#[derive(Debug, Default)]
297pub struct Scratch {
298 /// What each step produced, or `None` for a step that produces nothing and for one that has not
299 /// run yet.
300 slots: Vec<Option<Vector>>,
301 /// What each connective step has learned about its operands, indexed by step.
302 ///
303 /// Empty for every step that is not a connective and for a connective a filter has not reached
304 /// yet, since it is built the first time one runs and the shape it needs is not known before
305 /// then. This is the mutable half of the adaptive ordering and it is here rather than in
306 /// [`Prepared`] because a prepared expression is shared by every thread running the pipeline.
307 orders: Vec<Option<Ordering>>,
308}
309
310impl Scratch {
311 /// The order a connective's operands are run in.
312 ///
313 /// For the tests that say the learning reached the walk. Nothing in the engine asks a scratch
314 /// this, because the walk is the only thing that reads an ordering and it reads its own.
315 #[cfg(test)]
316 fn order(&self, step: usize) -> Option<&[usize]> {
317 self.orders[step].as_ref().map(Ordering::order)
318 }
319}
320
321impl Prepared {
322 /// Prepares `exprs` against `schema`.
323 ///
324 /// # Errors
325 ///
326 /// If a column reference names a binding the schema does not have, or if an aggregate appears
327 /// where an ordinary expression was expected. Both are failures of the plan rather than of the
328 /// data, which is why they are found here, once, rather than on some chunk in the middle of a
329 /// scan.
330 pub fn new(plan: &Plan, exprs: &[ExprRef], schema: &Schema) -> Result<Self> {
331 Self::build(plan, exprs, schema, false)
332 }
333
334 /// Prepares expressions whose caller can evaluate a shared expression graph as one unit.
335 pub(crate) fn shared(plan: &Plan, exprs: &[ExprRef], schema: &Schema) -> Result<Self> {
336 Self::build(plan, exprs, schema, true)
337 }
338
339 fn build(plan: &Plan, exprs: &[ExprRef], schema: &Schema, share: bool) -> Result<Self> {
340 Self::built(plan, exprs, schema, share, true)
341 }
342
343 fn built(
344 plan: &Plan,
345 exprs: &[ExprRef],
346 schema: &Schema,
347 share: bool,
348 fuse: bool,
349 ) -> Result<Self> {
350 let mut prepared = Self {
351 steps: Vec::new(),
352 types: Vec::new(),
353 spans: Vec::new(),
354 operands: Vec::new(),
355 last_use: Vec::new(),
356 roots: Vec::new(),
357 last_root: Vec::new(),
358 shared: HashMap::new(),
359 share,
360 fuse,
361 time_zone: SessionTimeZone::default(),
362 };
363 for &expr in exprs {
364 let root = prepared.push(plan, expr, schema)?;
365 prepared.roots.push(root);
366 }
367 prepared.last_use = prepared.last_uses();
368 prepared.last_root = prepared.last_roots();
369 Ok(prepared)
370 }
371
372 /// Uses the zone of the session that owns this prepared expression.
373 #[must_use]
374 pub fn in_session(mut self, session: &Session) -> Self {
375 self.set_time_zone(session.session_time_zone());
376 self
377 }
378
379 /// Sets the zone here and in every lambda body, which is prepared before the session is known.
380 fn set_time_zone(&mut self, time_zone: SessionTimeZone) {
381 self.time_zone = time_zone;
382 for step in &mut self.steps {
383 match step {
384 Step::Lambda { body, .. } => body.set_time_zone(time_zone),
385 Step::Fused { fallback, .. } => fallback.set_time_zone(time_zone),
386 _ => {}
387 }
388 }
389 }
390
391 /// Which step is the last to read each step, computed once when the expression is prepared.
392 ///
393 /// A root is never freed, because the whole point of running the array was to produce it. A
394 /// step nothing reads and that is not a root cannot happen, since every step is pushed by the
395 /// node that wanted it, but saying `usize::MAX` rather than asserting that keeps this a fact
396 /// about the array rather than a claim about the builder.
397 fn last_uses(&self) -> Vec<usize> {
398 let mut last = vec![usize::MAX; self.steps.len()];
399 for index in 0..self.steps.len() {
400 self.for_each_operand(index, |operand| last[operand] = index);
401 }
402 for &root in &self.roots {
403 last[root] = usize::MAX;
404 }
405 last
406 }
407
408 /// Which entries of [`roots`](Self::roots) are the last to name their step. See
409 /// [`last_root`](Self::last_root).
410 ///
411 /// A projection has a handful of roots, so this compares each against the ones after it rather
412 /// than building a map. It runs once per prepared expression.
413 fn last_roots(&self) -> Vec<bool> {
414 (0..self.roots.len()).map(|at| !self.roots[at + 1..].contains(&self.roots[at])).collect()
415 }
416
417 /// Visits the steps one step reads, whatever shape its operands are held in.
418 fn for_each_operand(&self, index: usize, mut visit: impl FnMut(usize)) {
419 match &self.steps[index] {
420 // A case's branches are arrays of their own and read nothing out of this one, and a
421 // fused tree reads its columns straight out of the chunk.
422 Step::Column(_)
423 | Step::Constant(_)
424 | Step::Case { .. }
425 | Step::Fused { .. }
426 | Step::Try { .. } => {}
427 Step::Cast { input, .. } | Step::InSet { input, .. } => visit(*input),
428 Step::Lambda { inputs, .. } => inputs.iter().for_each(|&input| visit(input)),
429 Step::Compare { left, right, .. } => {
430 visit(*left);
431 visit(*right);
432 }
433 Step::Conjunction { start, len, .. } | Step::Function { start, len, .. } => {
434 for &operand in &self.operands[*start..*start + *len] {
435 visit(operand);
436 }
437 }
438 }
439 }
440
441 /// Prepares one expression, which is the common case and saves the caller a slice.
442 ///
443 /// # Errors
444 ///
445 /// Whatever [`Prepared::new`] reports.
446 pub fn one(plan: &Plan, expr: ExprRef, schema: &Schema) -> Result<Self> {
447 Self::new(plan, &[expr], schema)
448 }
449
450 /// Working space sized for this expression.
451 #[must_use]
452 pub fn scratch(&self) -> Scratch {
453 Scratch {
454 slots: (0..self.steps.len()).map(|_| None).collect(),
455 orders: (0..self.steps.len()).map(|_| None).collect(),
456 }
457 }
458
459 /// How many expressions this was built from.
460 #[must_use]
461 pub fn len(&self) -> usize {
462 self.roots.len()
463 }
464
465 /// How many comparisons have their literal side already built.
466 ///
467 /// For the tests, for the same reason as [`Self::sets`]: an answer that moved would be a bug,
468 /// so the only thing a test can look at is whether the building happened.
469 #[cfg(test)]
470 fn literals_built(&self) -> usize {
471 self.steps.iter().filter(|step| matches!(step, Step::Compare { held: Some(_), .. })).count()
472 }
473
474 /// How many of the steps are an `IN` list folded back up.
475 ///
476 /// For the tests, which cannot see the fold in an answer because an answer that changed would
477 /// be a bug.
478 #[cfg(test)]
479 fn sets(&self) -> usize {
480 self.steps.iter().filter(|step| matches!(step, Step::InSet { .. })).count()
481 }
482
483 /// How many of the steps are a tree of decimal arithmetic run as one loop.
484 #[cfg(test)]
485 fn fused(&self) -> usize {
486 self.steps.iter().filter(|step| matches!(step, Step::Fused { .. })).count()
487 }
488
489 /// How many of the function steps worked something out when this was built.
490 ///
491 /// For the tests, which cannot see the hoisting in an answer because an answer that changed
492 /// would be a bug.
493 #[cfg(test)]
494 fn hoisted(&self) -> usize {
495 self.steps
496 .iter()
497 .filter(|step| matches!(step, Step::Function { recipe, .. } if recipe.hoists()))
498 .count()
499 }
500
501 /// Whether it was built from no expressions at all.
502 #[must_use]
503 pub fn is_empty(&self) -> bool {
504 self.roots.is_empty()
505 }
506
507 /// How many of the steps do something to a row.
508 ///
509 /// A column reference and a literal are not among them. A column reference computes nothing at
510 /// all, which is what makes a step that names one free rather than a copy, and a literal is
511 /// materialized once for the whole chunk rather than once a row. What is left is a pass over
512 /// the rows each, so this is roughly what one row costs, counted in the same unit the scan's
513 /// own reading of that row is counted in.
514 ///
515 /// What reads it is the scan, through the weight an operator reports to the pipeline. See
516 /// [`Stream::weight`](rudb_pipeline::Stream::weight).
517 #[must_use]
518 pub fn passes(&self) -> usize {
519 self.steps
520 .iter()
521 .filter(|step| !matches!(step, Step::Column(_) | Step::Constant(_)))
522 .count()
523 }
524
525 /// Evaluates every expression over `chunk`, appending one vector each to `out`.
526 ///
527 /// Appends rather than returns a `Vec`, so a caller in a loop reuses one buffer.
528 ///
529 /// # Errors
530 ///
531 /// Anything a kernel reports, on the first expression that reports it.
532 pub fn evaluate(
533 &self,
534 chunk: &Chunk,
535 scratch: &mut Scratch,
536 out: &mut Vec<Vector>,
537 ) -> Result<()> {
538 self.run(chunk, scratch)?;
539 for (at, &root) in self.roots.iter().enumerate() {
540 // The one place a column is copied, and it is copied because the caller is taking
541 // ownership of a vector that has to outlive the chunk it came from. `SELECT a` is that
542 // shape and a projection of a bare column is the only expression where it happens.
543 match self.steps[root] {
544 Step::Column(position) => out.push(chunk.column(position)?.clone()),
545 _ if self.last_root[at] => {
546 out.push(scratch.slots[root].take().ok_or_else(|| missing(root))?);
547 }
548 _ => {
549 out.push(scratch.slots[root].as_ref().ok_or_else(|| missing(root))?.clone());
550 }
551 }
552 }
553 Ok(())
554 }
555
556 /// [`evaluate`](Self::evaluate) for a caller that is done with `chunk`, which a projection is.
557 ///
558 /// Every step has run before a root is handed over, so nothing reads the chunk after that and a
559 /// root that is a bare column can take the column rather than copy it. A column named by more
560 /// than one root is copied for all but the last of them. `SELECT *` into a table is all bare
561 /// columns, and copying them was most of what its projection did.
562 ///
563 /// # Errors
564 ///
565 /// Whatever [`evaluate`](Self::evaluate) reports.
566 pub fn evaluate_taking(
567 &self,
568 chunk: Chunk,
569 scratch: &mut Scratch,
570 out: &mut Vec<Vector>,
571 ) -> Result<()> {
572 self.run(&chunk, scratch)?;
573 let width = chunk.width();
574 let mut columns: Vec<Option<Vector>> = chunk.into_columns().into_iter().map(Some).collect();
575 let mut uses = vec![0usize; width];
576 for &root in &self.roots {
577 if let Step::Column(position) = self.steps[root]
578 && position < width
579 {
580 uses[position] += 1;
581 }
582 }
583 for (at, &root) in self.roots.iter().enumerate() {
584 if let Step::Column(position) = self.steps[root] {
585 let missing = || {
586 Error::internal(format!(
587 "column {position} of a chunk that has {width} columns"
588 ))
589 };
590 let slot = columns.get_mut(position).ok_or_else(missing)?;
591 let left = &mut uses[position];
592 *left -= 1;
593 let column = if *left == 0 { slot.take() } else { slot.clone() };
594 out.push(column.ok_or_else(missing)?);
595 continue;
596 }
597 if self.last_root[at] {
598 out.push(scratch.slots[root].take().ok_or_else(|| missing(root))?);
599 } else {
600 out.push(scratch.slots[root].as_ref().ok_or_else(|| missing(root))?.clone());
601 }
602 }
603 Ok(())
604 }
605
606 /// Evaluates a single expression over `chunk`, handing back a reference to the answer.
607 ///
608 /// A reference rather than a vector, because the caller of this is a filter, which reads the
609 /// flags to build a selection and then drops them. Nothing about that wants ownership, and a
610 /// predicate that is a bare column reference, which `WHERE flag` is, would otherwise copy the
611 /// column to hand it over.
612 ///
613 /// # Errors
614 ///
615 /// Anything a kernel reports, and an internal error if this was not built from exactly one
616 /// expression.
617 pub fn evaluate_one<'s>(
618 &'s self,
619 chunk: &'s Chunk,
620 scratch: &'s mut Scratch,
621 ) -> Result<&'s Vector> {
622 let [root] = self.roots[..] else {
623 return Err(Error::internal(format!(
624 "evaluate_one over a prepared expression of {} roots",
625 self.roots.len()
626 )));
627 };
628 self.run(chunk, scratch)?;
629 self.operand(root, chunk, &scratch.slots)
630 }
631
632 /// Evaluates a single expression as a filter, handing back the rows it keeps.
633 ///
634 /// The difference between this and [`evaluate_one`](Self::evaluate_one) followed by
635 /// [`selection`] is the whole of what a threaded filter is. An `AND` evaluated as an expression
636 /// runs every conjunct over every row and then combines the flag vectors, so a predicate of four
637 /// conjuncts that each pass a fifth of the rows does five times the work of one that stops
638 /// looking at a row as soon as a conjunct rejects it. TPC-H Q6 is exactly that predicate.
639 ///
640 /// So the conjuncts of a top level `AND` are run one at a time, each over the rows the ones
641 /// before it left, and the moment nothing is left the rest of the predicate is not run at all.
642 /// The order they run in starts as the order the plan gives and then moves, because which
643 /// conjunct is worth running first is a question about the data and the scan is the thing
644 /// holding the answer. The `ordering` module has what is measured and how.
645 ///
646 /// A top level `OR` is threaded the same way against the complement. A row the first branch
647 /// accepts is a row the filter keeps whatever the rest of the predicate says about it, so each
648 /// branch is run over the rows no branch before it accepted, and the moment every row has been
649 /// accepted the rest of the predicate is not run either. That is the mirror of the `AND` case
650 /// and not an approximation of it: the answer is the same set of rows, because `OR` over three
651 /// valued logic is true wherever any branch is true and nothing a later branch says can take a
652 /// row back. It is worth less than the `AND` case in practice, since an `OR` of selective
653 /// branches leaves almost every row in play for the branch after, and it is worth having anyway
654 /// because the cost of finding that out is one merge per branch.
655 ///
656 /// What is threaded is the operand's own comparison rather than the whole of its subtree. A
657 /// conjunct of `a + b > 5` still adds over the whole chunk, because the scalar kernels take a
658 /// vector rather than a selection, and it is the comparison and everything downstream of it that
659 /// reads only the rows still in play. An operand that is a bare column or a function produces
660 /// flags over the chunk and is narrowed with [`refine_flags`], which is what keeps one awkward
661 /// operand from putting the others back on the unthreaded path. An operand that is itself a
662 /// connective recurses, so the two conjuncts of each half of `(a AND b) OR (c AND d)` are
663 /// threaded the same way the halves are.
664 ///
665 /// None of this is available to a projection. `SELECT a > 5 AND b LIKE 'x%'` wants a value per
666 /// row and the rows a selection dropped have no value in it, so [`evaluate`](Self::evaluate) and
667 /// [`evaluate_one`](Self::evaluate_one) evaluate the whole tree over the whole chunk and combine
668 /// flags. The two are separate entry points picked when the pipeline is built rather than one
669 /// path with a flag in it, because conflating them is a wrong answer rather than a slow one.
670 ///
671 /// # Errors
672 ///
673 /// Anything a kernel reports, and an internal error if this was not built from exactly one
674 /// expression.
675 pub fn evaluate_filter(&self, chunk: &Chunk, scratch: &mut Scratch) -> Result<Selection> {
676 let [root] = self.roots[..] else {
677 return Err(Error::internal(format!(
678 "evaluate_filter over a prepared expression of {} roots",
679 self.roots.len()
680 )));
681 };
682 scratch.slots.clear();
683 scratch.slots.resize_with(self.steps.len(), || None);
684 // A predicate that is not a connective at all is the same walk over one operand, which is
685 // where [`thread`](Self::thread) starts: it runs the tree and turns the flags into a
686 // selection, with no narrowing to do because nothing has narrowed anything yet.
687 self.thread(root, 0, chunk, scratch, None)
688 }
689
690 /// How many operands the top level `AND` of a filter has, or `None` when it has no such `AND`.
691 ///
692 /// Operand `i` is the `i`th child of the conjunction in the plan, which is the numbering
693 /// [`evaluate_settled`](Self::evaluate_settled) takes. `None` as well for an expression built
694 /// to share its steps, since a step an operand shares with a later one is a step that has to run
695 /// whether or not the first operand does.
696 #[must_use]
697 pub fn conjuncts(&self) -> Option<usize> {
698 let [root] = self.roots[..] else { return None };
699 match self.steps[root] {
700 Step::Conjunction { op: Connective::And, len, .. } if !self.share => Some(len),
701 _ => None,
702 }
703 }
704
705 /// [`evaluate_filter`](Self::evaluate_filter) with some operands of the top level `AND` known
706 /// to hold on every row of the chunk, which are not run at all.
707 ///
708 /// `settled[i]` is operand `i` in the numbering of [`conjuncts`](Self::conjuncts). What settles
709 /// one is the caller's business and it has to be a proof: an operand left out here is an
710 /// operand that keeps every row, nulls included, so a caller that is wrong about it gets rows
711 /// the query threw away. A scan knows it from the bounds of the part it read.
712 ///
713 /// # Errors
714 ///
715 /// As [`evaluate_filter`](Self::evaluate_filter).
716 pub fn evaluate_settled(
717 &self,
718 chunk: &Chunk,
719 scratch: &mut Scratch,
720 settled: &[bool],
721 ) -> Result<Selection> {
722 if self.conjuncts() != Some(settled.len()) || !settled.contains(&true) {
723 return self.evaluate_filter(chunk, scratch);
724 }
725 let [root] = self.roots[..] else {
726 return Err(Error::internal("a settled filter over several roots"));
727 };
728 scratch.slots.clear();
729 scratch.slots.resize_with(self.steps.len(), || None);
730 self.branches(root, 0, chunk, scratch, None, settled)
731 }
732
733 /// The operands of one connective, run in order, each over the rows the ones before it left.
734 ///
735 /// `live` is the rows this connective has to decide about and `None` means every row of the
736 /// chunk, which is not the same as a selection of all of them: it lets the first operand take
737 /// the unthreaded kernel rather than a pass over an identity selection. The answer is the rows
738 /// out of `live` the connective is true for.
739 ///
740 /// The walk is the same for both connectives and only the bookkeeping differs. `AND` carries the
741 /// rows every operand so far has kept, so each answer replaces it. `OR` carries the rows no
742 /// operand so far has accepted, so each answer comes out of it and the rows the connective keeps
743 /// are the ones that went missing along the way.
744 ///
745 /// The operand is not `steps[begin..=operand]` evaluated and then narrowed. Its subtree is run
746 /// over the whole chunk and it is the operand itself that reads only the rows in play, except
747 /// where the operand is another connective, which recurses and threads its own operands from
748 /// here rather than falling back to a flag vector. That is what makes `(a AND b) OR (c AND d)`
749 /// four threaded comparisons rather than two threaded ones and two flag passes.
750 fn branches(
751 &self,
752 index: usize,
753 begin: usize,
754 chunk: &Chunk,
755 scratch: &mut Scratch,
756 live: Option<&Selection>,
757 settled: &[bool],
758 ) -> Result<Selection> {
759 let Step::Conjunction { op, start, len } = self.steps[index] else {
760 return Err(Error::internal("a connective walk over a step that is not a connective"));
761 };
762 let operands = &self.operands[start..start + len];
763 let rows = chunk.len();
764 // Out of the scratch for the length of the walk, because the walk runs steps and running a
765 // step wants the scratch. It goes back at the end, which is also where it learns. A walk
766 // that fails leaves the slot empty and the next chunk starts the connective over, which is
767 // a history lost on a query that is about to stop running anyway.
768 let mut order = scratch.orders[index]
769 .take()
770 .unwrap_or_else(|| Ordering::new(op, self.weights(operands, begin)));
771 let mut carried: Option<Selection> = live.cloned();
772 // The operands already answered as the other end of a range, see [`Self::range`].
773 let mut ranged: u128 = 0;
774 for slot in 0..len {
775 if carried.as_ref().is_some_and(Selection::is_empty) {
776 break;
777 }
778 let which = order.at(slot);
779 // Known to keep every row, so running it would hand back the rows it was given.
780 if settled.get(which) == Some(&true) || (which < 128 && (ranged >> which) & 1 == 1) {
781 continue;
782 }
783 let operand = operands[which];
784 // The array is in post order and an operand's whole subtree sits between the operand
785 // before it and the operand itself, which is a range the run order cannot move. That is
786 // what lets the operands run in any order at all without a second structure to say
787 // where each one starts.
788 let from = if which == 0 { begin } else { operands[which - 1] + 1 };
789 let given = carried.as_ref().map_or(rows, Selection::len);
790 let fused = match op {
791 Connective::And => {
792 self.range(operands, which, settled, ranged, chunk, scratch, carried.as_ref())?
793 }
794 Connective::Or => None,
795 };
796 let answered = match fused {
797 Some((other, answered)) => {
798 ranged |= 1 << other;
799 order.observed(other, given, answered.len());
800 answered
801 }
802 None => self.thread(operand, from, chunk, scratch, carried.as_ref())?,
803 };
804 order.observed(which, given, answered.len());
805 carried = Some(match (op, carried) {
806 (Connective::And, _) => answered,
807 (Connective::Or, None) => answered.complement(rows),
808 (Connective::Or, Some(carried)) => carried.without(&answered),
809 });
810 // Keep a shared step alive when a later operand still reads it.
811 for step in from..=operand {
812 if self.last_use[step] <= operand {
813 scratch.slots[step] = None;
814 }
815 }
816 }
817 order.relearn();
818 scratch.orders[index] = Some(order);
819 Ok(match (op, carried) {
820 // A connective with no operands, which the binder does not build and which is answered
821 // here rather than left to index arithmetic: an empty `AND` is every row and an empty
822 // `OR` is none.
823 (Connective::And, None) => live.cloned().unwrap_or_else(|| Selection::identity(rows)),
824 (Connective::And, Some(kept)) => kept,
825 (Connective::Or, None) => Selection::empty(),
826 (Connective::Or, Some(missed)) => match live {
827 None => missed.complement(rows),
828 Some(live) => live.without(&missed),
829 },
830 })
831 }
832
833 /// Operand `which` of an `AND` and another operand of it answered together, when the two are a
834 /// low and a high bound on the same column, as the other operand and the rows the two keep.
835 ///
836 /// `l_shipdate >= date '1994-01-01' and l_shipdate < date '1995-01-01'` is two comparisons that
837 /// each keep most of a chunk and together keep a seventh of it, so running one and then the
838 /// other walks most of the chunk twice and builds a selection of most of it in between. Here it
839 /// is one pass with one test a row, see [`rudb_kernels::select_range`]. The plan does not change
840 /// and neither does the order the operands learn, since both of them are told what the pair
841 /// kept. `None` when there is no such pair, when the column has a form the range has no loop
842 /// for, and for steps built to be shared, whose slots a later operand may read.
843 #[expect(clippy::too_many_arguments, reason = "the walk's state, handed over as it stands")]
844 fn range(
845 &self,
846 operands: &[usize],
847 which: usize,
848 settled: &[bool],
849 ranged: u128,
850 chunk: &Chunk,
851 scratch: &Scratch,
852 live: Option<&Selection>,
853 ) -> Result<Option<(usize, Selection)>> {
854 if self.share {
855 return Ok(None);
856 }
857 let Some((column, low)) = self.bound(operands[which]) else { return Ok(None) };
858 let Some(other) = (0..operands.len().min(128)).find(|&other| {
859 other != which
860 && settled.get(other) != Some(&true)
861 && (ranged >> other) & 1 == 0
862 && self.bound(operands[other]) == Some((column, !low))
863 }) else {
864 return Ok(None);
865 };
866 let (lower, upper) = if low {
867 (operands[which], operands[other])
868 } else {
869 (operands[other], operands[which])
870 };
871 let (Some((left, lower)), Some((_, upper))) = (self.end(lower), self.end(upper)) else {
872 return Ok(None);
873 };
874 let values = self.operand(left, chunk, &scratch.slots)?;
875 let live = live.map(Selection::indices);
876 Ok(rudb_kernels::select_range(values, lower, upper, live).map(|kept| (other, kept)))
877 }
878
879 /// The column a comparison of a column with a literal reads, and whether the literal is where
880 /// the column's values start (`>`, `>=`) or where they end (`<`, `<=`).
881 fn bound(&self, operand: usize) -> Option<(usize, bool)> {
882 let Step::Compare { op, left, right, .. } = &self.steps[operand] else { return None };
883 let (Step::Column(column), Step::Constant(value)) =
884 (&self.steps[*left], &self.steps[*right])
885 else {
886 return None;
887 };
888 if value.is_null() || self.types[*left] != self.types[*right] {
889 return None;
890 }
891 match op {
892 Comparison::Greater | Comparison::GreaterOrEqual => Some((*column, true)),
893 Comparison::Less | Comparison::LessOrEqual => Some((*column, false)),
894 _ => None,
895 }
896 }
897
898 /// The column step of a comparison [`Self::bound`] accepted, and its end of the range.
899 fn end(&self, operand: usize) -> Option<(usize, rudb_kernels::Bound<'_>)> {
900 let Step::Compare { op, left, right, held } = &self.steps[operand] else { return None };
901 let Step::Constant(value) = &self.steps[*right] else { return None };
902 Some((*left, rudb_kernels::Bound { op: *op, value, held: held.as_ref() }))
903 }
904
905 /// What each operand of a connective costs to run over a chunk, for the ordering to divide by.
906 ///
907 /// An operand costs what its whole subtree costs, which is the steps from where the operand
908 /// before it ended up to the operand itself.
909 fn weights(&self, operands: &[usize], begin: usize) -> Vec<f64> {
910 let mut costs = Vec::with_capacity(operands.len());
911 let mut from = begin;
912 for &operand in operands {
913 costs.push((from..=operand).map(|step| self.weight(step)).sum());
914 from = operand + 1;
915 }
916 costs
917 }
918
919 /// Roughly what one step costs to run over a chunk, against a comparison of two fixed width
920 /// columns as the unit.
921 ///
922 /// A ranking rather than a prediction. Nothing downstream reads the number itself, only which
923 /// of two of them is larger, and the differences that decide an order are the big ones: a
924 /// column reference costs nothing because it is read in place, a string function costs many
925 /// times what an integer comparison costs, and a comparison over a variable length type costs
926 /// several times what the same comparison over a fixed width one costs. Everything finer than
927 /// that is below the noise of what the window is measuring anyway.
928 fn weight(&self, index: usize) -> f64 {
929 match &self.steps[index] {
930 // Read straight out of the chunk at the point an operand is wanted, so there is no step
931 // to run and nothing to charge for.
932 Step::Column(_) => 0.0,
933 // One vector built per chunk, however many rows the chunk has.
934 Step::Constant(_) => 0.25,
935 // The operands carry the cost of a connective, and they are steps of their own.
936 Step::Conjunction { .. } => 0.0,
937 Step::Cast { input, .. } => 2.0 * touching(&self.types[*input]),
938 Step::Compare { left, .. } => touching(&self.types[*left]),
939 // One hash and one probe a row, whatever the list holds, which is the point of it. It
940 // is dearer than a comparison and much cheaper than the chain of them it replaced.
941 Step::InSet { input, .. } => 2.0 * touching(&self.types[*input]),
942 Step::Function { start, len, .. } => {
943 let widest = self.operands[*start..*start + *len]
944 .iter()
945 .map(|&argument| touching(&self.types[argument]))
946 .fold(1.0, f64::max);
947 4.0 * widest
948 }
949 // A branch per arm, each of which is a prepared expression of its own that this does
950 // not look inside. Charging for the arms alone understates it and says the right thing
951 // about the order, which is that a `CASE` is not what you want in front.
952 Step::Case { arms, .. } => 4.0 * arms.len() as f64,
953 // The operand once, which is what it costs on every chunk that raises nothing.
954 Step::Try { inner } => (0..inner.steps.len()).map(|step| inner.weight(step)).sum(),
955 // A run of the body per element, which is several a row, and a list to take apart and
956 // put back together around it.
957 Step::Lambda { .. } => 16.0,
958 // An integer operation a row per node and no check, which is a quarter of what the
959 // function steps it replaced cost each.
960 Step::Fused { fused, .. } => fused.len() as f64,
961 }
962 }
963
964 /// One operand of a connective, over the rows it is still worth asking about.
965 ///
966 /// `begin` is the first step of the operand's subtree, which the caller knows because the steps
967 /// are in post order.
968 fn thread(
969 &self,
970 index: usize,
971 begin: usize,
972 chunk: &Chunk,
973 scratch: &mut Scratch,
974 live: Option<&Selection>,
975 ) -> Result<Selection> {
976 if matches!(self.steps[index], Step::Conjunction { .. }) {
977 return self.branches(index, begin, chunk, scratch, live, &[]);
978 }
979 for step in begin..index {
980 self.run_step(step, chunk, scratch)?;
981 }
982 // Straight to the rows it keeps, and only among the ones still in play, where the list and
983 // the column allow it. See [`rudb_kernels::select_in`].
984 if let Step::InSet { input, members } = &self.steps[index] {
985 let column = self.operand(*input, chunk, &scratch.slots)?;
986 if let Some(kept) = rudb_kernels::select_in(column, members, live) {
987 return Ok(kept);
988 }
989 }
990 if let Step::Compare { op, left, right, held } = &self.steps[index] {
991 let one = self.operand(*left, chunk, &scratch.slots)?;
992 let other = self.operand(*right, chunk, &scratch.slots)?;
993 let held = held.as_ref();
994 return match live {
995 // The first operand has every row in play, and asking the threaded kernel for that
996 // would be a pass over an identity selection the unthreaded one does not need.
997 None => select_prepared(*op, one, other, held),
998 Some(live) => refine_prepared(*op, one, other, live, held),
999 };
1000 }
1001 // A later LIKE in a threaded filter often sees only a handful of survivors.
1002 // Gather its arguments, not the whole chunk, while preserving the stable
1003 // dictionary behind a gathered string column. The ordinary full-vector
1004 // path remains cheaper when most rows are still live.
1005 if let (Some(live), Step::Function { recipe, written, start, len }) =
1006 (live, &self.steps[index])
1007 && matches!(recipe.name(), "~~" | "!~~" | "~~*" | "!~~*")
1008 && live.len().saturating_mul(4) <= chunk.len()
1009 {
1010 let flags = self
1011 .with_operands(*start, *len, chunk, &scratch.slots, |args| {
1012 let gathered = args
1013 .iter()
1014 .map(|arg| arg.gather(live.indices()))
1015 .collect::<Result<Vec<_>>>()?;
1016 let narrowed = gathered.iter().collect::<Vec<_>>();
1017 rudb_kernels::call_prepared(
1018 recipe,
1019 &narrowed,
1020 &self.types[index],
1021 Some(&|| written.clone()),
1022 )
1023 })
1024 .map_err(|error| error.with_fallback_span(self.spans[index]))?;
1025 return Ok(selection(&flags, live.len()).compose(live));
1026 }
1027 self.run_step(index, chunk, scratch)?;
1028 let flags = self.operand(index, chunk, &scratch.slots)?;
1029 match live {
1030 None => Ok(selection(flags, chunk.len())),
1031 Some(live) => refine_flags(flags, live),
1032 }
1033 }
1034
1035 /// Runs every step in order, filling the slots.
1036 fn run(&self, chunk: &Chunk, scratch: &mut Scratch) -> Result<()> {
1037 scratch.slots.clear();
1038 scratch.slots.resize_with(self.steps.len(), || None);
1039 for index in 0..self.steps.len() {
1040 self.run_step(index, chunk, scratch)?;
1041 }
1042 Ok(())
1043 }
1044
1045 /// Runs one step and empties the slot of every operand this was the last step to read.
1046 fn run_step(&self, index: usize, chunk: &Chunk, scratch: &mut Scratch) -> Result<()> {
1047 let produced = self
1048 .step(index, chunk, &scratch.slots)
1049 .map_err(|error| error.with_fallback_span(self.spans[index]))?;
1050 scratch.slots[index] = produced;
1051 let slots = &mut scratch.slots;
1052 self.for_each_operand(index, |operand| {
1053 if self.last_use[operand] == index {
1054 slots[operand] = None;
1055 }
1056 });
1057 Ok(())
1058 }
1059
1060 /// Runs one step, given what the steps before it produced.
1061 fn step(
1062 &self,
1063 index: usize,
1064 chunk: &Chunk,
1065 slots: &[Option<Vector>],
1066 ) -> Result<Option<Vector>> {
1067 let ty = &self.types[index];
1068 let produced = match &self.steps[index] {
1069 Step::Column(_) => None,
1070 Step::Constant(value) => Some(Vector::constant(ty.clone(), value.clone(), chunk.len())),
1071 Step::Cast { input, try_cast } => Some(cast_in_time_zone(
1072 self.operand(*input, chunk, slots)?,
1073 ty,
1074 *try_cast,
1075 Some(self.time_zone),
1076 )?),
1077 Step::Compare { op, left, right, held } => Some(compare_prepared(
1078 *op,
1079 self.operand(*left, chunk, slots)?,
1080 self.operand(*right, chunk, slots)?,
1081 held.as_ref(),
1082 )?),
1083 Step::Conjunction { op, start, len } => {
1084 Some(
1085 self.with_operands(*start, *len, chunk, slots, |children| {
1086 combine(*op, children)
1087 })?,
1088 )
1089 }
1090 // The one call with no argument to take a row count from, so it is given the chunk's.
1091 Step::Function { recipe, len: 0, .. } if recipe.name() == "random" => {
1092 Some(rudb_kernels::random(chunk.len())?)
1093 }
1094 Step::Function { recipe, written, start, len } => {
1095 Some(self.with_operands(*start, *len, chunk, slots, |args| {
1096 rudb_kernels::call_prepared(recipe, args, ty, Some(&|| written.clone()))
1097 })?)
1098 }
1099 Step::InSet { input, members } => {
1100 Some(in_set(self.operand(*input, chunk, slots)?, members, ty)?)
1101 }
1102 Step::Case { arms, otherwise, blend } => {
1103 Some(self.case(chunk, arms, otherwise.as_ref(), blend.as_ref(), ty)?)
1104 }
1105 Step::Try { inner } => {
1106 let mut scratch = inner.scratch();
1107 Some(attempt(chunk, ty, |rows| inner.evaluate_one(rows, &mut scratch).cloned())?)
1108 }
1109 Step::Fused { fused, fallback } => Some(match fused.run(chunk) {
1110 Some(answer) => answer,
1111 None => fallback.evaluate_one(chunk, &mut fallback.scratch())?.clone(),
1112 }),
1113 Step::Lambda { inputs, runner, body } => {
1114 let mut operands = Vec::with_capacity(inputs.len());
1115 for &input in inputs {
1116 operands.push(self.operand(input, chunk, slots)?);
1117 }
1118 let mut scratch = body.scratch();
1119 Some(runner.run(&operands, chunk, &mut |inner| {
1120 body.evaluate_one(inner, &mut scratch).cloned()
1121 })?)
1122 }
1123 };
1124 Ok(produced)
1125 }
1126
1127 /// The vector a step produced, or the chunk's column if the step is a column reference.
1128 fn operand<'v>(
1129 &self,
1130 index: usize,
1131 chunk: &'v Chunk,
1132 slots: &'v [Option<Vector>],
1133 ) -> Result<&'v Vector> {
1134 if let Step::Column(position) = self.steps[index] {
1135 return chunk.column(position);
1136 }
1137 slots[index].as_ref().ok_or_else(|| missing(index))
1138 }
1139
1140 /// Hands a kernel the references to an operand list, without allocating for the usual widths.
1141 ///
1142 /// One, two and three because those are what a bound tree is made of: every scalar function in
1143 /// the catalog is unary or binary, a comparison is binary, and a conjunction is two or three
1144 /// often enough to be worth a line. A stack array for those means a chain of eight additions
1145 /// makes zero allocations for its operand lists over a chunk instead of eight, and eight
1146 /// allocations a chunk at the rate a pipeline produces chunks is a real number rather than a
1147 /// tidiness argument. Anything wider falls back to [`gather`](Self::gather), which is a `Vec`
1148 /// of pointers and still moves no data.
1149 fn with_operands<'v, T>(
1150 &self,
1151 start: usize,
1152 len: usize,
1153 chunk: &'v Chunk,
1154 slots: &'v [Option<Vector>],
1155 run: impl FnOnce(&[&'v Vector]) -> Result<T>,
1156 ) -> Result<T> {
1157 match self.operands[start..start + len] {
1158 [a] => run(&[self.operand(a, chunk, slots)?]),
1159 [a, b] => run(&[self.operand(a, chunk, slots)?, self.operand(b, chunk, slots)?]),
1160 [a, b, c] => run(&[
1161 self.operand(a, chunk, slots)?,
1162 self.operand(b, chunk, slots)?,
1163 self.operand(c, chunk, slots)?,
1164 ]),
1165 _ => {
1166 let gathered = self.gather(start, len, chunk, slots)?;
1167 run(&gathered)
1168 }
1169 }
1170 }
1171
1172 /// References to an operand list, for a kernel that takes a slice of them.
1173 ///
1174 /// The `Vec` here is the allocation the module documentation names: it holds pointers rather
1175 /// than vectors, so it is a dozen bytes an operand and no data moves.
1176 fn gather<'v>(
1177 &self,
1178 start: usize,
1179 len: usize,
1180 chunk: &'v Chunk,
1181 slots: &'v [Option<Vector>],
1182 ) -> Result<Vec<&'v Vector>> {
1183 let mut gathered = Vec::with_capacity(len);
1184 for &operand in &self.operands[start..start + len] {
1185 gathered.push(self.operand(operand, chunk, slots)?);
1186 }
1187 Ok(gathered)
1188 }
1189
1190 /// A searched `CASE` over the rows no earlier arm claimed.
1191 ///
1192 /// The same shape [`evaluate`](crate::evaluate) has, because the thing that makes it that shape
1193 /// is a correctness rule rather than a performance one: `CASE WHEN x <> 0 THEN 1 / x ELSE 0 END`
1194 /// divides by zero on the rows the arm excludes if the arm is evaluated for them.
1195 ///
1196 /// Each arm answers the rows no earlier arm claimed, so the answers come back short and out of
1197 /// order and have to be put back in the order the rows arrived in. That is what [`Assembly`] is:
1198 /// the arms are laid end to end into one run of data and the interleave is a single typed copy
1199 /// over it. It used to be a `Vec<Value>` filled a row at a time and handed to
1200 /// `Vector::from_values`, which is a heap allocation and a drop for every string in the answer.
1201 /// On the ClickBench query that groups by a `CASE` over `Referer` that was about a quarter of
1202 /// the whole query.
1203 ///
1204 /// What is left of #57 here is the narrowing. An arm still narrows the whole chunk rather than
1205 /// the columns it reads, and the selection threading that replaces the narrowing entirely is
1206 /// the item this one was carved out of.
1207 fn case(
1208 &self,
1209 chunk: &Chunk,
1210 arms: &[PreparedArm],
1211 otherwise: Option<&Prepared>,
1212 blend: Option<&Blend>,
1213 ty: &LogicalType,
1214 ) -> Result<Vector> {
1215 let claimed = self.claims(chunk, arms)?;
1216 if let Some(blend) = blend
1217 && let Some(blended) = blended(chunk, &claimed, blend)?
1218 {
1219 return Ok(blended);
1220 }
1221 let mut built = Assembly::new(ty.clone(), chunk.len())?;
1222 let branches = arms.iter().map(|arm| &arm.then).map(Some).chain([otherwise]);
1223 for (branch, rows) in branches.zip(&claimed) {
1224 let (Some(branch), false) = (branch, rows.is_empty()) else { continue };
1225 // The same cut the conditions skip above, skipped here for the same reason: a branch
1226 // that claimed every row claimed them in order, so narrowing to them is a copy of every
1227 // column in the chunk to arrive back at the chunk.
1228 let cut;
1229 let matched = if rows.len() == chunk.len() {
1230 chunk
1231 } else {
1232 cut = narrow(chunk, rows)?;
1233 &cut
1234 };
1235 let mut scratch = branch.scratch();
1236 let results = branch.evaluate_one(matched, &mut scratch)?;
1237 built.place(&placed(rows)?, results)?;
1238 }
1239 built.finish()
1240 }
1241
1242 /// The rows each branch of a `CASE` answers, one list per arm in order and the `ELSE` last.
1243 ///
1244 /// Only the conditions are run here, which is what keeps the rule the doc above states: an arm's
1245 /// condition is evaluated over the rows no earlier arm claimed, so a condition that would raise
1246 /// on a row an earlier arm took is never asked about it. The results are worked out afterwards,
1247 /// once, from these lists, and both ways of working them out want the same thing, which is the
1248 /// rows of one branch in the order they arrived in.
1249 fn claims(&self, chunk: &Chunk, arms: &[PreparedArm]) -> Result<Vec<Vec<usize>>> {
1250 let mut claimed = Vec::with_capacity(arms.len() + 1);
1251 let mut pending: Vec<usize> = (0..chunk.len()).collect();
1252 for arm in arms {
1253 if pending.is_empty() {
1254 claimed.push(Vec::new());
1255 continue;
1256 }
1257 // `pending` starts as every row in order and only ever shrinks, so the same length is
1258 // the same rows in the same order and there is nothing to cut. That is the whole of the
1259 // first arm of a one armed `CASE`, which is the shape of the ClickBench query this was
1260 // measured on, and cutting it was a copy of every column in the chunk for nothing.
1261 let cut;
1262 let narrowed = if pending.len() == chunk.len() {
1263 chunk
1264 } else {
1265 cut = narrow(chunk, &pending)?;
1266 &cut
1267 };
1268 let mut scratch = arm.when.scratch();
1269 let flags = arm.when.evaluate_one(narrowed, &mut scratch)?;
1270 let mut taken = Vec::new();
1271 let mut still = Vec::new();
1272 // row at a time: splitting the rows an arm claims from the ones it leaves is a test per
1273 // row, and what replaces it is the selection threading the rest of #57 asks for rather
1274 // than anything that can be done here.
1275 for (at, &row) in pending.iter().enumerate() {
1276 if is_true(&flags.value_at(at)) {
1277 taken.push(row);
1278 } else {
1279 still.push(row);
1280 }
1281 }
1282 claimed.push(taken);
1283 pending = still;
1284 }
1285 claimed.push(pending);
1286 Ok(claimed)
1287 }
1288
1289 /// Flattens one expression, appending its steps and returning the index of its last one.
1290 fn push(&mut self, plan: &Plan, expr: ExprRef, schema: &Schema) -> Result<usize> {
1291 if self.share
1292 && let Some(&step) = self.shared.get(&expr)
1293 {
1294 return Ok(step);
1295 }
1296 let ty = plan.expr_type(expr).clone();
1297 if self.fuse
1298 && let Some(fused) = Fused::compile(plan, expr, schema)
1299 {
1300 let fallback = Self::built(plan, &[expr], schema, false, false)?;
1301 let step = Step::Fused { fused: Box::new(fused), fallback: Box::new(fallback) };
1302 return Ok(self.place(plan, expr, step, ty));
1303 }
1304 if let Some((stamp, count)) = stamped_seconds(plan, expr) {
1305 let (start, len) = self.push_list(plan, &[stamp, count], schema)?;
1306 let step = Step::Function {
1307 recipe: Recipe::new("__rudb_stamp_seconds", &self.literals(start, len)),
1308 written: written(plan, expr, schema),
1309 start,
1310 len,
1311 };
1312 return Ok(self.place(plan, expr, step, ty));
1313 }
1314 let step = match *plan.expr(expr) {
1315 Expr::Column(binding) => {
1316 let position = schema.position_of(binding).ok_or_else(|| {
1317 Error::internal(format!(
1318 "column #{}.{} is not in the schema this operator was given",
1319 binding.table, binding.column
1320 ))
1321 })?;
1322 Step::Column(position)
1323 }
1324 Expr::Constant(reference) => Step::Constant(plan.value(reference).clone()),
1325 Expr::Cast { input, try_cast } => {
1326 Step::Cast { input: self.push(plan, input, schema)?, try_cast }
1327 }
1328 Expr::Compare { op, left, right } => {
1329 let left = self.push(plan, left, schema)?;
1330 let right = self.push(plan, right, schema)?;
1331 Step::Compare { op: comparison(op), left, right, held: self.held(left, right) }
1332 }
1333 Expr::Conjunction { op, children } => {
1334 let list = plan.expr_list(children).to_vec();
1335 match self.membership(plan, connective(op), &list, schema)? {
1336 Some(step) => step,
1337 None => {
1338 let (start, len) = self.push_list(plan, &list, schema)?;
1339 Step::Conjunction { op: connective(op), start, len }
1340 }
1341 }
1342 }
1343 Expr::Function { name, args } if lambda_call(plan, args).is_some() => {
1344 let Some((lambda, inputs)) = lambda_call(plan, args) else {
1345 return Err(Error::internal("a lambda call without a lambda"));
1346 };
1347 let Expr::Lambda { body, .. } = *plan.expr(lambda) else {
1348 return Err(Error::internal("a lambda call without a lambda"));
1349 };
1350 let runner = Lambda::new(plan, plan.string(name), lambda, &inputs, schema)?;
1351 let body = Self::one(plan, body, runner.schema())?;
1352 let mut steps = Vec::with_capacity(inputs.len());
1353 for &input in &inputs {
1354 steps.push(self.push(plan, input, schema)?);
1355 }
1356 Step::Lambda { inputs: steps, runner: Box::new(runner), body: Box::new(body) }
1357 }
1358 Expr::LambdaParam(binding) => {
1359 let position = schema.position_of(binding).ok_or_else(|| {
1360 Error::internal(format!(
1361 "lambda parameter @{}.{} is not in the schema its body was given",
1362 binding.table, binding.column
1363 ))
1364 })?;
1365 Step::Column(position)
1366 }
1367 Expr::Lambda { .. } => {
1368 return Err(Error::internal(
1369 "a lambda was evaluated outside the function that takes it",
1370 ));
1371 }
1372 Expr::Function { name, args } if plan.string(name) == "try" => {
1373 let [only] = plan.expr_list(args)[..] else {
1374 return Err(Error::internal("a TRY without exactly one operand"));
1375 };
1376 Step::Try { inner: Box::new(Self::one(plan, only, schema)?) }
1377 }
1378 Expr::Function { name, args } => {
1379 let (start, len) = self.push_list(plan, plan.expr_list(args), schema)?;
1380 Step::Function {
1381 recipe: Recipe::new(plan.string(name), &self.literals(start, len)),
1382 written: written(plan, expr, schema),
1383 start,
1384 len,
1385 }
1386 }
1387 Expr::Aggregate { name, .. } => {
1388 return Err(Error::internal(format!(
1389 "the {} aggregate was evaluated as an ordinary expression",
1390 plan.string(name)
1391 )));
1392 }
1393 Expr::Window { name, .. } => {
1394 return Err(Error::internal(format!(
1395 "the {} window function was evaluated as an ordinary expression",
1396 plan.string(name)
1397 )));
1398 }
1399 Expr::Case { arms, otherwise } => {
1400 let mut prepared = Vec::new();
1401 for &arm in plan.arm_list(arms) {
1402 prepared.push(PreparedArm {
1403 when: Self::one(plan, arm.when, schema)?,
1404 then: Self::one(plan, arm.then, schema)?,
1405 });
1406 }
1407 let otherwise = match otherwise {
1408 Some(otherwise) => Some(Self::one(plan, otherwise, schema)?),
1409 None => None,
1410 };
1411 let blend = blending(&ty, &prepared, otherwise.as_ref());
1412 Step::Case { arms: prepared, otherwise, blend }
1413 }
1414 };
1415 Ok(self.place(plan, expr, step, ty))
1416 }
1417
1418 /// Appends a built step and answers its index.
1419 fn place(&mut self, plan: &Plan, expr: ExprRef, step: Step, ty: LogicalType) -> usize {
1420 self.steps.push(step);
1421 self.types.push(ty);
1422 self.spans.push(plan.expr_span(expr));
1423 let step = self.steps.len() - 1;
1424 if self.share {
1425 self.shared.insert(expr, step);
1426 }
1427 step
1428 }
1429
1430 /// Flattens a list of expressions and records where its operand run starts and how long it is.
1431 ///
1432 /// The operand run is written after every child has been flattened rather than as they go,
1433 /// because a child that is itself a list would otherwise interleave its run with this one.
1434 fn push_list(
1435 &mut self,
1436 plan: &Plan,
1437 exprs: &[ExprRef],
1438 schema: &Schema,
1439 ) -> Result<(usize, usize)> {
1440 let mut indices = Vec::with_capacity(exprs.len());
1441 for &expr in exprs {
1442 indices.push(self.push(plan, expr, schema)?);
1443 }
1444 let start = self.operands.len();
1445 let len = indices.len();
1446 self.operands.extend(indices);
1447 Ok((start, len))
1448 }
1449
1450 /// This connective folded back into the `IN` the user wrote, or `None` when it is not one.
1451 ///
1452 /// What the binder writes for `x IN (1, 2, 3)` is `x = 1 OR x = 2 OR x = 3`, and for
1453 /// `x NOT IN (1, 2, 3)` it is `x <> 1 AND x <> 2 AND x <> 3`. So the shape looked for is every
1454 /// child a comparison of the one direction, every left the same expression, and every right a
1455 /// literal. Anything else is left alone, which covers the `OR` that was written as an `OR` and
1456 /// the one where an `IN` has been flattened together with another branch. The second is a fold
1457 /// this could make and does not, and it is worth having later out of a query that wants it
1458 /// rather than now out of a guess.
1459 ///
1460 /// This runs before the children are pushed, and that is the whole reason it is here rather than
1461 /// as a pass over the finished array. A step that nothing reads is still a step the walk runs,
1462 /// because the walk over a subtree is a range and not a graph, so folding after the fact would
1463 /// leave every equality in place and running.
1464 fn membership(
1465 &mut self,
1466 plan: &Plan,
1467 op: Connective,
1468 children: &[ExprRef],
1469 schema: &Schema,
1470 ) -> Result<Option<Step>> {
1471 let wanted = match op {
1472 Connective::Or => CompareOp::Equal,
1473 Connective::And => CompareOp::NotEqual,
1474 };
1475 let mut subject: Option<ExprRef> = None;
1476 let mut values = Vec::with_capacity(children.len());
1477 for &child in children {
1478 let Expr::Compare { op: found, left, right } = *plan.expr(child) else {
1479 return Ok(None);
1480 };
1481 if found != wanted || !same(plan, *subject.get_or_insert(left), left) {
1482 return Ok(None);
1483 }
1484 let Expr::Constant(reference) = *plan.expr(right) else {
1485 return Ok(None);
1486 };
1487 values.push(plan.value(reference).clone());
1488 }
1489 let (Some(subject), Some(members)) = (subject, Members::of(&values, op == Connective::And))
1490 else {
1491 return Ok(None);
1492 };
1493 Ok(Some(Step::InSet { input: self.push(plan, subject, schema)?, members }))
1494 }
1495
1496 /// The literal side of a comparison, in the one row column the comparison reads it through.
1497 ///
1498 /// The right side first, because that is the side the binder puts a literal on and the side the
1499 /// loops are written for. Two literals is a comparison the optimizer folded, and if it did not
1500 /// then the kernel answers it once for the whole vector and never reads either column, so
1501 /// neither side is built here.
1502 fn held(&self, left: usize, right: usize) -> Option<Held> {
1503 let (at, other) = match (&self.steps[left], &self.steps[right]) {
1504 (Step::Constant(_), Step::Constant(_)) => return None,
1505 (_, Step::Constant(value)) => (right, value),
1506 (Step::Constant(value), _) => (left, value),
1507 _ => return None,
1508 };
1509 Held::of(&self.types[at], other)
1510 }
1511
1512 /// The literal behind each argument in a run of the operand list, and `None` for an argument
1513 /// that is anything else.
1514 ///
1515 /// This is what a [`Recipe`] hoists from. An argument that is a literal in the plan arrives as a
1516 /// constant vector holding exactly this value on every chunk, so what a kernel reads here is
1517 /// what it would have read per chunk. An argument that is a cast of a literal reads as `None`,
1518 /// which is a call the kernel decides per chunk as it always did, and the optimizer folds most
1519 /// of those before the plan gets here anyway.
1520 fn literals(&self, start: usize, len: usize) -> Vec<Option<Value>> {
1521 self.operands[start..start + len]
1522 .iter()
1523 .map(|&operand| match &self.steps[operand] {
1524 Step::Constant(value) => Some(value.clone()),
1525 _ => None,
1526 })
1527 .collect()
1528 }
1529}
1530
1531/// Whether two expressions of one plan are the same expression, written once or written twice.
1532///
1533/// The binder binds the subject of an `IN` once and points every comparison it writes at that one
1534/// reference, so the answer is almost always the first line. A plan that has been through a rewrite,
1535/// and a plan read back from its own text, hold two copies of the same tree instead, and for the
1536/// fold in [`Prepared::membership`] those are the same expression.
1537///
1538/// The four shapes handled are what an `IN` is written over: a column, a literal, a cast of either,
1539/// and a call, which is TPC-H query 22 asking whether the first two digits of a phone number are in
1540/// a list. Anything else answers no, which costs a fold that could have happened rather than a wrong
1541/// one. The walk is bounded by the size of the subject and a subject is small.
1542/// The timestamp and the whole count of `stamp + to_seconds(CAST(count AS DOUBLE))`, the shape the
1543/// benchmark view writes `INTERVAL (EventTime) SECOND` in, and `None` for anything else.
1544///
1545/// It runs as one call, [`rudb_kernels`]'s `__rudb_stamp_seconds`, rather than as a cast to a
1546/// double, an interval per row and a shift by it.
1547fn stamped_seconds(plan: &Plan, expr: ExprRef) -> Option<(ExprRef, ExprRef)> {
1548 let Expr::Function { name, args } = *plan.expr(expr) else { return None };
1549 if plan.string(name) != "+" || plan.expr_type(expr) != &LogicalType::Timestamp {
1550 return None;
1551 }
1552 let &[one, other] = plan.expr_list(args) else { return None };
1553 let (stamp, interval) =
1554 if plan.expr_type(one) == &LogicalType::Timestamp { (one, other) } else { (other, one) };
1555 if plan.expr_type(stamp) != &LogicalType::Timestamp {
1556 return None;
1557 }
1558 let Expr::Function { name, args } = *plan.expr(interval) else { return None };
1559 let &[cast] = plan.expr_list(args) else { return None };
1560 let Expr::Cast { input, try_cast: false } = *plan.expr(cast) else { return None };
1561 let whole = matches!(
1562 plan.expr_type(input),
1563 LogicalType::TinyInt
1564 | LogicalType::SmallInt
1565 | LogicalType::Integer
1566 | LogicalType::BigInt
1567 | LogicalType::UTinyInt
1568 | LogicalType::USmallInt
1569 | LogicalType::UInteger
1570 );
1571 (plan.string(name) == "to_seconds" && plan.expr_type(cast) == &LogicalType::Double && whole)
1572 .then_some((stamp, input))
1573}
1574
1575fn same(plan: &Plan, left: ExprRef, right: ExprRef) -> bool {
1576 if left == right {
1577 return true;
1578 }
1579 if plan.expr_type(left) != plan.expr_type(right) {
1580 return false;
1581 }
1582 match (plan.expr(left), plan.expr(right)) {
1583 (Expr::Column(one), Expr::Column(other)) => one == other,
1584 (Expr::Constant(one), Expr::Constant(other)) => plan.value(*one) == plan.value(*other),
1585 (
1586 Expr::Cast { input: one, try_cast: first },
1587 Expr::Cast { input: other, try_cast: second },
1588 ) => first == second && same(plan, *one, *other),
1589 (
1590 Expr::Function { name: one, args: first },
1591 Expr::Function { name: other, args: second },
1592 ) => {
1593 let (first, second) = (plan.expr_list(*first), plan.expr_list(*second));
1594 plan.string(*one) == plan.string(*other)
1595 && first.len() == second.len()
1596 && first.iter().zip(second).all(|(&one, &other)| same(plan, one, other))
1597 }
1598 _ => false,
1599 }
1600}
1601
1602/// What touching a value of this type costs, against a fixed width one as the unit.
1603///
1604/// A variable length value is a pointer to follow and a length that is not the same twice, and a
1605/// nested one is that per element. Four is not measured, and what it has to be is large enough that
1606/// the ordering puts a fixed width comparison in front of a string one and small enough that it does
1607/// not put one in front of a string comparison that rejects every row.
1608fn touching(ty: &LogicalType) -> f64 {
1609 match ty.physical() {
1610 PhysicalType::Varlen => 4.0,
1611 PhysicalType::List | PhysicalType::Array | PhysicalType::Struct => 8.0,
1612 _ => 1.0,
1613 }
1614}
1615
1616/// The error for a slot that should have held something and did not.
1617///
1618/// This cannot happen while the array is in post order, since every operand's index is smaller than
1619/// the index of the step using it and every step runs in order. It is an error rather than a panic
1620/// because the property it depends on is a property of [`Prepared::push`], and the day somebody
1621/// writes a pass that reorders the array is the day it stops holding.
1622fn missing(index: usize) -> Error {
1623 Error::internal(format!("step {index} was used as an operand before it produced anything"))
1624}
1625
1626/// Chunk rows as the positions an [`Assembly`] places a piece at.
1627///
1628/// A chunk is at most [`VECTOR_SIZE`](rudb_vector::VECTOR_SIZE) rows, so the conversion cannot fail
1629/// in practice. It is checked rather than cast because a silent truncation here would put a value in
1630/// the wrong row, and a wrong row is the one kind of bug nothing downstream can notice.
1631fn placed(rows: &[usize]) -> Result<Vec<u32>> {
1632 rows.iter()
1633 .map(|&row| {
1634 u32::try_from(row).map_err(|_| Error::internal("a chunk of more than u32 rows"))
1635 })
1636 .collect()
1637}
1638
1639/// A `CASE` answered as codes over the dictionary its branches share, or `None` for a chunk that
1640/// cannot be.
1641///
1642/// Declined per chunk rather than once, because whether a column arrives coded is a fact about the
1643/// chunk and not about the expression. The same query reads codes out of a native file and plain
1644/// strings out of rows held in memory, and one file can hand a column over as a dictionary in one
1645/// part and as plain data in the next. Everything that declines does so before a code is written, so
1646/// the caller starts the general path from nothing rather than from a half filled answer.
1647fn blended(chunk: &Chunk, claimed: &[Vec<usize>], blend: &Blend) -> Result<Option<Vector>> {
1648 let Some((dictionary, literals)) = agreed(chunk, blend)? else { return Ok(None) };
1649 let mut codes = vec![0; chunk.len()];
1650 for (branch, rows) in blend.branches.iter().zip(claimed) {
1651 match *branch {
1652 Branch::Column(position) => {
1653 let Some((from, _)) = chunk.column(position)?.stable_dictionary_parts() else {
1654 return Ok(None);
1655 };
1656 for &row in rows {
1657 codes[row] = from[row];
1658 }
1659 }
1660 Branch::Literal(at) => {
1661 for &row in rows {
1662 codes[row] = literals[at];
1663 }
1664 }
1665 }
1666 }
1667 Vector::stable_dictionary(codes, dictionary).map(Some)
1668}
1669
1670/// The one dictionary every branch of a blend names values in, and the code each literal sits at.
1671///
1672/// Three things say no. A column that did not arrive as a stable dictionary has no codes to copy. A
1673/// second column over a different dictionary would have codes that mean something else, and a code
1674/// is a position in one dictionary and nothing anywhere else. And a literal the dictionary does not
1675/// hold has no code at all, which for `ELSE ''` over a column where no row is empty is the honest
1676/// answer rather than a missing one.
1677///
1678/// The null check is the fourth. A dictionary keeps its nulls in the values it points at rather than
1679/// beside its codes, so a column carrying its own validity is one whose codes do not say everything
1680/// the column says, and copying them would turn its nulls into whatever their codes happen to name.
1681fn agreed(chunk: &Chunk, blend: &Blend) -> Result<Option<(Arc<Vector>, Vec<u32>)>> {
1682 let mut held: Option<(&Vector, &Arc<Vector>)> = None;
1683 for branch in &blend.branches {
1684 let Branch::Column(position) = *branch else { continue };
1685 let column = chunk.column(position)?;
1686 let Some((_, dictionary)) = column.stable_dictionary_parts() else { return Ok(None) };
1687 if column.validity().has_nulls(chunk.len()) {
1688 return Ok(None);
1689 }
1690 match held {
1691 Some((_, first)) if !Arc::ptr_eq(first, dictionary) => return Ok(None),
1692 Some(_) => {}
1693 None => held = Some((column, dictionary)),
1694 }
1695 }
1696 let Some((column, dictionary)) = held else { return Ok(None) };
1697 let mut codes = Vec::with_capacity(blend.literals.len());
1698 for (text, lookup) in &blend.literals {
1699 match lookup.find(column, text.as_bytes()) {
1700 Some(Ok(Found::At(code))) => codes.push(code),
1701 Some(Err(error)) => return Err(error),
1702 Some(Ok(Found::Absent)) | None => return Ok(None),
1703 }
1704 }
1705 Ok(Some((Arc::clone(dictionary), codes)))
1706}
1707
1708/// The blend a `CASE` can be answered by, or `None` for one that has to read its branches' values.
1709fn blending(ty: &LogicalType, arms: &[PreparedArm], otherwise: Option<&Prepared>) -> Option<Blend> {
1710 if !matches!(ty, LogicalType::Varchar) {
1711 return None;
1712 }
1713 let otherwise = otherwise?;
1714 let mut branches = Vec::with_capacity(arms.len() + 1);
1715 let mut literals = Vec::new();
1716 for branch in arms.iter().map(|arm| &arm.then).chain([otherwise]) {
1717 branches.push(named(branch, &mut literals)?);
1718 }
1719 // All of them literals means there is no dictionary to name any of them in, and a `CASE` whose
1720 // every branch is a constant is not a thing anybody writes.
1721 let any = branches.iter().any(|branch| matches!(branch, Branch::Column(_)));
1722 any.then_some(Blend { branches, literals })
1723}
1724
1725/// The branch a prepared expression stands for, when it names a value rather than computing one.
1726fn named(prepared: &Prepared, literals: &mut Vec<(String, Lookup)>) -> Option<Branch> {
1727 match prepared.steps.as_slice() {
1728 [Step::Column(position)] => Some(Branch::Column(*position)),
1729 [Step::Constant(Value::Varchar(text))] => {
1730 literals.push((text.clone(), Lookup::default()));
1731 Some(Branch::Literal(literals.len() - 1))
1732 }
1733 _ => None,
1734 }
1735}
1736
1737/// The chunk cut down to the given rows.
1738///
1739/// The reason `CASE` is written with this rather than by evaluating every arm over the whole chunk
1740/// and picking afterwards. `CASE WHEN x <> 0 THEN 1 // x ELSE 0 END` divides by zero on the rows the
1741/// arm does not apply to if the arm is evaluated for them, and a `CASE` that raises on a row it was
1742/// written to exclude is the classic wrong answer this shape prevents.
1743/// `TRY(x)` over a chunk, where `run` evaluates `x` over whatever chunk it is given.
1744///
1745/// The whole chunk first, and only if that raises one of the errors `TRY` catches is it run again a
1746/// row at a time with a null for each row that raises, which is the pin's order. A chunk that
1747/// raises nothing pays for nothing, and any other error goes out as it came in.
1748pub(crate) fn attempt(
1749 chunk: &Chunk,
1750 ty: &LogicalType,
1751 mut run: impl FnMut(&Chunk) -> Result<Vector>,
1752) -> Result<Vector> {
1753 match run(chunk) {
1754 Err(error) if caught(&error) => {}
1755 answer => return answer,
1756 }
1757 let mut values = Vec::with_capacity(chunk.len());
1758 // row at a time: this is the path for a chunk in which some row raised, and finding which one
1759 // is the whole of what it does.
1760 for row in 0..chunk.len() {
1761 values.push(match run(&narrow(chunk, &[row])?) {
1762 Ok(one) => one.try_value_at(0)?,
1763 Err(error) if caught(&error) => Value::Null,
1764 Err(error) => return Err(error),
1765 });
1766 }
1767 Vector::from_values(ty.clone(), &values)
1768}
1769
1770/// Whether `TRY` answers null for this error rather than passing it on, which is the pin's three
1771/// kinds of error a value can cause. The binder's folding has the same list.
1772fn caught(error: &Error) -> bool {
1773 matches!(error.code(), ErrorCode::Conversion | ErrorCode::OutOfRange | ErrorCode::InvalidInput)
1774}
1775
1776pub(crate) fn narrow(chunk: &Chunk, rows: &[usize]) -> Result<Chunk> {
1777 let mut selection = Selection::with_capacity(rows.len());
1778 for &row in rows {
1779 selection.push(row);
1780 }
1781 chunk.clone().select(&selection)
1782}
1783
1784/// The kernels' comparison for the plan's.
1785///
1786/// A translation rather than one shared enum, because the kernels are rank 3 and the plan is rank
1787/// 9. This function is the whole of what that separation costs.
1788pub(crate) fn comparison(op: CompareOp) -> Comparison {
1789 match op {
1790 CompareOp::Equal => Comparison::Equal,
1791 CompareOp::NotEqual => Comparison::NotEqual,
1792 CompareOp::Less => Comparison::Less,
1793 CompareOp::LessOrEqual => Comparison::LessOrEqual,
1794 CompareOp::Greater => Comparison::Greater,
1795 CompareOp::GreaterOrEqual => Comparison::GreaterOrEqual,
1796 CompareOp::DistinctFrom => Comparison::DistinctFrom,
1797 CompareOp::NotDistinctFrom => Comparison::NotDistinctFrom,
1798 }
1799}
1800
1801/// The kernels' connective for the plan's.
1802pub(crate) fn connective(op: ConjunctionOp) -> Connective {
1803 match op {
1804 ConjunctionOp::And => Connective::And,
1805 ConjunctionOp::Or => Connective::Or,
1806 }
1807}
1808
1809#[cfg(test)]
1810mod tests {
1811 use rudb_common::{Field, LogicalType, Value};
1812 use rudb_kernels::is_true;
1813 use rudb_plan::{ExprRef, Node, Plan};
1814 use rudb_vector::{Chunk, Selection, Vector};
1815
1816 use super::{Prepared, narrow};
1817 use crate::expr::evaluate;
1818 use crate::schema::Schema;
1819
1820 /// Two columns with a null in each, because every disagreement between these two evaluators
1821 /// that is worth finding is a disagreement about which rows are null.
1822 fn input() -> (Schema, Chunk) {
1823 let schema = Schema::numbered(
1824 vec![Field::new("x", LogicalType::Integer), Field::new("s", LogicalType::Varchar)],
1825 0,
1826 );
1827 let x = Vector::from_values(
1828 LogicalType::Integer,
1829 &[Value::Integer(3), Value::Integer(1), Value::Null, Value::Integer(2)],
1830 )
1831 .expect("four integers");
1832 let s = Vector::from_values(
1833 LogicalType::Varchar,
1834 &[
1835 Value::Varchar("a".to_string()),
1836 Value::Null,
1837 Value::Varchar("c".to_string()),
1838 Value::Varchar("a".to_string()),
1839 ],
1840 )
1841 .expect("four strings");
1842 (schema, Chunk::new(vec![x, s]).expect("two columns of four rows"))
1843 }
1844
1845 /// The expressions of a projection written in the plan's textual form, over the two columns
1846 /// [`input`] produces.
1847 ///
1848 /// Going through the text rather than the arena builders for the reason the other test module
1849 /// gives: a test that says what it evaluates in the notation a plan dump uses is a test whose
1850 /// failure can be pasted into a plan and vice versa.
1851 fn projection(exprs: &str) -> (Plan, Vec<ExprRef>) {
1852 let text =
1853 format!("Project #1 [{exprs}]\n Get memory.main.t AS t #0 [x::INTEGER, s::VARCHAR]");
1854 let plan = Plan::parse(&text).expect("a well formed plan");
1855 let Node::Project { exprs, .. } = *plan.node(plan.root()) else {
1856 panic!("the root of that text is a projection");
1857 };
1858 let list = plan.expr_list(exprs).to_vec();
1859 (plan, list)
1860 }
1861
1862 /// Every expression shape, evaluated both ways over the same chunk.
1863 ///
1864 /// This is the agreement the module documentation claims and it is the only thing that makes
1865 /// the prepared form safe to put in front of the tree walk. The generated well typed trees the
1866 /// test gate of #57 asks for are a wider version of this and are worth building once the
1867 /// selection threaded shapes exist to disagree about.
1868 fn agrees(exprs: &str) {
1869 let (schema, chunk) = input();
1870 let (plan, list) = projection(exprs);
1871 let prepared = Prepared::new(&plan, &list, &schema).expect("the expressions resolve");
1872 let mut scratch = prepared.scratch();
1873 let mut fast = Vec::new();
1874 prepared.evaluate(&chunk, &mut scratch, &mut fast).expect("the prepared form runs");
1875 for (at, &expr) in list.iter().enumerate() {
1876 let slow = evaluate(&plan, expr, &schema, &chunk).expect("the tree walk runs");
1877 for row in 0..chunk.len() {
1878 assert_eq!(
1879 fast[at].value_at(row),
1880 slow.value_at(row),
1881 "expression {at} of `{exprs}` at row {row}"
1882 );
1883 }
1884 }
1885 }
1886
1887 /// Three decimal columns of TPC-H's shape, in the form `form` puts them in.
1888 fn decimals(prices: &[i128], form: fn(Vector) -> Vector) -> (Schema, Chunk) {
1889 let ty = LogicalType::Decimal { width: 15, scale: 2 };
1890 let schema = Schema::numbered(
1891 vec![
1892 Field::new("p", ty.clone()),
1893 Field::new("d", ty.clone()),
1894 Field::new("t", ty.clone()),
1895 ],
1896 0,
1897 );
1898 let column =
1899 |values: Vec<Value>| form(Vector::from_values(ty.clone(), &values).expect("decimals"));
1900 let decimal = |unscaled| Value::Decimal { unscaled, width: 15, scale: 2 };
1901 let p = column(prices.iter().map(|&v| decimal(v)).collect());
1902 let d = column((0..prices.len() as i128).map(|v| decimal(v % 11)).collect());
1903 let t = column((0..prices.len() as i128).map(|v| decimal(v % 9)).collect());
1904 (schema, Chunk::new(vec![p, d, t]).expect("three columns"))
1905 }
1906
1907 /// q01's charge, as the binder writes it.
1908 const CHARGE: &str = "\"*\"(\"*\"(CAST(#0.0::DECIMAL(15,2))::DECIMAL(18,2), \
1909 CAST(\"-\"(1.00::DECIMAL(16,2), CAST(#0.1::DECIMAL(15,2))::DECIMAL(16,2))::DECIMAL(16,2))\
1910 ::DECIMAL(18,2))::DECIMAL(18,4), CAST(\"+\"(1.00::DECIMAL(16,2), \
1911 CAST(#0.2::DECIMAL(15,2))::DECIMAL(16,2))::DECIMAL(16,2))::DECIMAL(18,2))::DECIMAL(18,6) AS a";
1912
1913 /// The fused answer, the unfused one and the tree walk's, over one chunk.
1914 fn three_ways(chunk: &Chunk, schema: &Schema) -> [rudb_common::Result<Vec<Value>>; 3] {
1915 let text = format!(
1916 "Project #1 [{CHARGE}]\n Get memory.main.t AS t #0 \
1917 [p::DECIMAL(15,2), d::DECIMAL(15,2), t::DECIMAL(15,2)]"
1918 );
1919 let plan = Plan::parse(&text).expect("a well formed plan");
1920 let Node::Project { exprs, .. } = *plan.node(plan.root()) else {
1921 panic!("the root of that text is a projection");
1922 };
1923 let expr = plan.expr_list(exprs)[0];
1924 let values = |vector: &Vector| (0..chunk.len()).map(|row| vector.value_at(row)).collect();
1925 let fused = Prepared::one(&plan, expr, schema).expect("resolves");
1926 assert_eq!(fused.fused(), 1, "the whole tree is one step");
1927 let unfused = Prepared::built(&plan, &[expr], schema, false, false).expect("resolves");
1928 assert_eq!(unfused.fused(), 0);
1929 let run = |prepared: &Prepared| {
1930 prepared.evaluate_one(chunk, &mut prepared.scratch()).map(&values)
1931 };
1932 [run(&fused), run(&unfused), evaluate(&plan, expr, schema, chunk).map(|v| values(&v))]
1933 }
1934
1935 fn all_agree(chunk: &Chunk, schema: &Schema) {
1936 let [fused, unfused, walked] = three_ways(chunk, schema);
1937 let fused = fused.expect("fits");
1938 assert_eq!(fused, unfused.expect("fits"));
1939 assert_eq!(fused, walked.expect("fits"));
1940 }
1941
1942 /// The epoch plus a whole count of seconds runs as one call, and agrees with the cast, the
1943 /// interval and the shift it stands for, on both sides of the count where the double stops
1944 /// being exact and on a count that takes the answer out of range.
1945 #[test]
1946 fn a_timestamp_plus_whole_seconds_agrees_with_the_interval_it_stands_for() {
1947 let schema = Schema::numbered(vec![Field::new("x", LogicalType::BigInt)], 0);
1948 let counts = [
1949 Value::BigInt(1_373_000_000),
1950 Value::BigInt(-5),
1951 Value::Null,
1952 Value::BigInt(9_007_199_254),
1953 Value::BigInt(9_007_199_255),
1954 Value::BigInt(9_000_000_000_123),
1955 ];
1956 let x = Vector::from_values(LogicalType::BigInt, &counts).expect("six counts");
1957 let chunk = Chunk::new(vec![x]).expect("one column");
1958 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]";
1959 let plan = Plan::parse(text).expect("a well formed plan");
1960 let Node::Project { exprs, .. } = *plan.node(plan.root()) else {
1961 panic!("the root of that text is a projection");
1962 };
1963 let list = plan.expr_list(exprs).to_vec();
1964 let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
1965 assert!(
1966 prepared.steps.iter().any(
1967 |step| matches!(step, super::Step::Function { recipe, .. } if recipe.name() == "__rudb_stamp_seconds")
1968 ),
1969 "the shift is one call"
1970 );
1971 let mut scratch = prepared.scratch();
1972 let mut fast = Vec::new();
1973 prepared.evaluate(&chunk, &mut scratch, &mut fast).expect("the prepared form runs");
1974 let slow = evaluate(&plan, list[0], &schema, &chunk).expect("the tree walk runs");
1975 for row in 0..chunk.len() {
1976 assert_eq!(fast[0].value_at(row), slow.value_at(row), "row {row}");
1977 }
1978 assert_eq!(fast[0].value_at(0), Value::Timestamp(1_373_000_000_000_000));
1979
1980 let far = Vector::from_values(LogicalType::BigInt, &[Value::BigInt(9_300_000_000_000)])
1981 .expect("one count");
1982 let chunk = Chunk::new(vec![far]).expect("one column");
1983 let mut fast = Vec::new();
1984 let fused = prepared.evaluate(&chunk, &mut scratch, &mut fast);
1985 let slow = evaluate(&plan, list[0], &schema, &chunk).map(|_| ());
1986 assert!(fused.is_err() && slow.is_err(), "past the last timestamp both raise");
1987 }
1988
1989 #[test]
1990 fn decimal_arithmetic_run_as_one_loop_agrees_in_every_form() {
1991 let prices: Vec<i128> = (0..2500).map(|v| 90_000 + v * 37).collect();
1992 let packed = |vector: Vector| vector.bit_packed().expect("packs");
1993 let coded = |vector: Vector| {
1994 let rows = vector.len();
1995 let codes = (0..rows as u32).rev().collect();
1996 Vector::dictionary(codes, vector.bit_packed().expect("packs")).expect("in range")
1997 };
1998 // Codes too far apart for a block to unpack the run they cover.
1999 let scattered = |vector: Vector| {
2000 let rows = vector.len() as u32;
2001 let codes = (0..rows).map(|row| row * 997 % rows).collect();
2002 Vector::dictionary(codes, vector.bit_packed().expect("packs")).expect("in range")
2003 };
2004 for form in [std::convert::identity, packed, coded, scattered] {
2005 let (schema, chunk) = decimals(&prices, form);
2006 all_agree(&chunk, &schema);
2007 }
2008 }
2009
2010 #[test]
2011 fn a_chunk_the_ranges_cannot_prove_raises_what_the_steps_raise() {
2012 // The large price in the second block, so a flat column gets as far as running the first.
2013 let mut prices = vec![5; 300];
2014 prices.push(999_999_999_999_999);
2015 let packed = |vector: Vector| vector.bit_packed().expect("packs");
2016 for form in [std::convert::identity, packed] {
2017 let (schema, chunk) = decimals(&prices, form);
2018 let [fused, unfused, _] = three_ways(&chunk, &schema);
2019 let (fused, unfused) = (fused.expect_err("overflows"), unfused.expect_err("overflows"));
2020 assert_eq!(fused.message(), unfused.message());
2021 }
2022 }
2023
2024 #[test]
2025 fn a_chunk_with_a_null_goes_through_the_steps() {
2026 let ty = LogicalType::Decimal { width: 15, scale: 2 };
2027 let (schema, mut chunk) = decimals(&[100, 200, 300], std::convert::identity);
2028 let with_null = Vector::from_values(
2029 ty,
2030 &[Value::Decimal { unscaled: 5, width: 15, scale: 2 }, Value::Null, Value::Null],
2031 )
2032 .expect("decimals");
2033 chunk = Chunk::new(vec![
2034 chunk.column(0).expect("p").clone(),
2035 with_null,
2036 chunk.column(2).expect("t").clone(),
2037 ])
2038 .expect("three columns");
2039 all_agree(&chunk, &schema);
2040 }
2041
2042 #[test]
2043 fn a_column_reference_agrees() {
2044 agrees("#0.0::INTEGER AS a, #0.1::VARCHAR AS b");
2045 }
2046
2047 #[test]
2048 fn a_constant_agrees() {
2049 agrees("7::INTEGER AS a, NULL::INTEGER AS b");
2050 }
2051
2052 #[test]
2053 fn a_cast_agrees() {
2054 agrees("CAST(#0.0::INTEGER)::BIGINT AS a, CAST(#0.0::INTEGER)::VARCHAR AS b");
2055 }
2056
2057 #[test]
2058 fn a_comparison_agrees() {
2059 agrees("(#0.0::INTEGER > 1::INTEGER)::BOOLEAN AS a");
2060 }
2061
2062 #[test]
2063 fn a_conjunction_agrees() {
2064 agrees(
2065 "((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 3::INTEGER)::BOOLEAN)\
2066 ::BOOLEAN AS a",
2067 );
2068 }
2069
2070 #[test]
2071 fn a_function_agrees() {
2072 agrees("\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER AS a");
2073 }
2074
2075 /// The two evaluators quote the same expression when a divisor is zero. Per #262.
2076 ///
2077 /// This is the one message in the engine that depends on how an expression is written rather
2078 /// than on what it computes, and the two evaluators render it at different times: the prepared
2079 /// form when the pipeline is built, the tree walk on the row that fails. Same renderer, so the
2080 /// same sentence, and this is what says so.
2081 #[test]
2082 fn both_evaluators_quote_the_same_expression_when_a_divisor_is_zero() {
2083 let (schema, chunk) = input();
2084 let (plan, list) = projection("\"//\"(#0.0::INTEGER, 0::INTEGER)::INTEGER AS a");
2085 let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
2086 let mut scratch = prepared.scratch();
2087 let mut out = Vec::new();
2088 let fast = prepared.evaluate(&chunk, &mut scratch, &mut out).expect_err("divides by zero");
2089 let slow = evaluate(&plan, list[0], &schema, &chunk).expect_err("divides by zero");
2090 assert_eq!(fast.message(), slow.message());
2091 assert!(fast.message().starts_with("Division by zero in expression (x // 0)."), "{fast}");
2092 }
2093
2094 #[test]
2095 fn a_case_agrees() {
2096 agrees(
2097 "CASE WHEN (#0.0::INTEGER > 1::INTEGER)::BOOLEAN THEN 10::INTEGER \
2098 ELSE 20::INTEGER END::INTEGER AS a",
2099 );
2100 }
2101
2102 /// A second arm, which is the first one that sees a cut chunk rather than the whole one.
2103 ///
2104 /// The first arm of any `CASE` runs over every row, so it takes the path that does not cut at
2105 /// all, and a `CASE` of one arm never exercises the other one. Two arms and an `ELSE` puts a
2106 /// different set of rows in front of each of the three.
2107 ///
2108 /// That this is the only test here reaching the cut was checked rather than assumed, by gating a
2109 /// panic on it and rerunning the seven. This one failed and the other six did not.
2110 #[test]
2111 fn a_case_of_two_arms_agrees() {
2112 agrees(
2113 "CASE WHEN (#0.0::INTEGER > 2::INTEGER)::BOOLEAN THEN 10::INTEGER \
2114 WHEN (#0.0::INTEGER > 1::INTEGER)::BOOLEAN THEN 20::INTEGER \
2115 ELSE 30::INTEGER END::INTEGER AS a",
2116 );
2117 }
2118
2119 /// No `ELSE`, so the rows no arm claims are null rather than anything.
2120 ///
2121 /// The case a run of data with a hole in it gets wrong: a null still occupies a position, and an
2122 /// assembly that skipped it would put every value after it one row early.
2123 #[test]
2124 fn a_case_with_no_else_agrees() {
2125 agrees(
2126 "CASE WHEN (#0.0::INTEGER > 2::INTEGER)::BOOLEAN THEN 10::INTEGER \
2127 END::INTEGER AS a",
2128 );
2129 }
2130
2131 /// An arm no row takes, so it contributes nothing to the answer and must not shift it.
2132 #[test]
2133 fn a_case_whose_arm_claims_nothing_agrees() {
2134 agrees(
2135 "CASE WHEN (#0.0::INTEGER > 99::INTEGER)::BOOLEAN THEN 10::INTEGER \
2136 ELSE 20::INTEGER END::INTEGER AS a",
2137 );
2138 }
2139
2140 /// Strings, which is the case that used to allocate one of them per row and drop it afterwards.
2141 ///
2142 /// The arm reads a column and the `ELSE` is a constant, which is the shape of the ClickBench
2143 /// query this path was rewritten for: the arm arrives as views over an arena and the `ELSE` as
2144 /// one value repeated, and the two have to be laid end to end into a single arena.
2145 #[test]
2146 fn a_case_over_strings_agrees() {
2147 agrees(
2148 "CASE WHEN (#0.0::INTEGER > 1::INTEGER)::BOOLEAN THEN #0.1::VARCHAR \
2149 ELSE ''::VARCHAR END::VARCHAR AS a",
2150 );
2151 }
2152
2153 /// A null inside an arm, which is a different thing from a row no arm claimed.
2154 ///
2155 /// Both come out null and they reach the validity mask by different routes, so a mask built for
2156 /// one of them and not the other reads correct on whichever test only has the other in it.
2157 #[test]
2158 fn a_case_whose_arm_answers_null_agrees() {
2159 agrees(
2160 "CASE WHEN (#0.0::INTEGER > 1::INTEGER)::BOOLEAN THEN #0.1::VARCHAR \
2161 ELSE NULL::VARCHAR END::VARCHAR AS a",
2162 );
2163 }
2164
2165 /// A `WHEN` over a column that is null on some rows, which is neither true nor false there.
2166 ///
2167 /// A three valued `WHEN` is what decides whether a row goes to the arm or falls through, and
2168 /// treating unknown as true would claim a row the `ELSE` should have had.
2169 #[test]
2170 fn a_case_whose_test_is_null_on_some_rows_agrees() {
2171 agrees(
2172 "CASE WHEN (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN THEN 10::INTEGER \
2173 ELSE 20::INTEGER END::INTEGER AS a",
2174 );
2175 }
2176
2177 /// The same expression twice, which is where the tree walk copies the column twice and this
2178 /// does not, and the answers still have to be identical.
2179 #[test]
2180 fn a_column_mentioned_three_times_agrees() {
2181 agrees("\"+\"(\"+\"(#0.0::INTEGER, #0.0::INTEGER)::INTEGER, #0.0::INTEGER)::INTEGER AS a");
2182 }
2183
2184 /// The intermediates of a chain are not all held to the end of it.
2185 ///
2186 /// This is the whole difference between the prepared form being faster than the tree walk on a
2187 /// deep chain and being slower than it, and it is a property of the slot array rather than of
2188 /// any answer, so it is asserted here rather than left to the benchmark to catch.
2189 #[test]
2190 fn a_chain_holds_one_intermediate_at_a_time() {
2191 let (schema, chunk) = input();
2192 let mut expr = "#0.0::INTEGER".to_string();
2193 for _ in 0..8 {
2194 expr = format!("\"+\"({expr}, 1::INTEGER)::INTEGER");
2195 }
2196 let (plan, list) = projection(&format!("{expr} AS a"));
2197 let prepared = Prepared::new(&plan, &list, &schema).expect("the chain resolves");
2198 let mut scratch = prepared.scratch();
2199 prepared.run(&chunk, &mut scratch).expect("the chain runs");
2200 let live = scratch.slots.iter().filter(|slot| slot.is_some()).count();
2201 assert_eq!(live, 1, "a chain that has run should be holding its answer and nothing else");
2202 }
2203
2204 /// The rows a threaded filter keeps are the rows the tree walk says the predicate is true for.
2205 ///
2206 /// Every threaded conjunct is a chance to disagree with the unthreaded answer about a null,
2207 /// about a row an earlier conjunct had already dropped, or about a chunk nothing survives, and
2208 /// the answer is a set of row numbers rather than a vector, so this is checked against the tree
2209 /// walk read a row at a time rather than against the prepared form it is part of.
2210 fn filters(predicate: &str) {
2211 let (schema, chunk) = input();
2212 let (plan, list) = projection(&format!("{predicate} AS p"));
2213 let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
2214 let mut scratch = prepared.scratch();
2215 let threaded = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs");
2216 let flags = evaluate(&plan, list[0], &schema, &chunk).expect("the tree walk runs");
2217 let expected = Selection::from_predicate(chunk.len(), |row| is_true(&flags.value_at(row)));
2218 assert_eq!(threaded, expected, "`{predicate}`");
2219 // And running it again over the same scratch is the same answer, because a pipeline calls
2220 // this once a chunk and a slot left behind by the conjunct before would show up here.
2221 let again = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs again");
2222 assert_eq!(again, expected, "`{predicate}` a second time");
2223 }
2224
2225 /// Two bounds on one column are answered as one range, and the rows have to be the ones the
2226 /// tree walk keeps, on a column with no nulls both flat and bit packed, with the bounds either
2227 /// way round, strict or not, beside another conjunct, and meeting nowhere.
2228 #[test]
2229 fn two_bounds_on_one_column_keep_the_rows_the_tree_walk_keeps() {
2230 let schema = Schema::numbered(
2231 vec![Field::new("x", LogicalType::Integer), Field::new("s", LogicalType::Varchar)],
2232 0,
2233 );
2234 let values: Vec<i32> = (0..1000).map(|row| 700 + (row * 37) % 600).collect();
2235 let flat = Vector::flat(LogicalType::Integer, rudb_vector::Data::Int32(values.into()))
2236 .expect("integers are an i32 layout");
2237 let packed = flat.bit_packed().expect("a six hundred wide range packs");
2238 let words = Vector::from_values(
2239 LogicalType::Varchar,
2240 &(0..1000).map(|row| Value::Varchar(["a", "b"][row % 2].into())).collect::<Vec<_>>(),
2241 )
2242 .expect("strings");
2243 let predicates = [
2244 "((#0.0::INTEGER >= 800::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 900::INTEGER)::BOOLEAN)",
2245 "((#0.0::INTEGER <= 900::INTEGER)::BOOLEAN AND (#0.0::INTEGER > 800::INTEGER)::BOOLEAN)",
2246 "((#0.0::INTEGER >= 0::INTEGER)::BOOLEAN AND (#0.0::INTEGER <= 5000::INTEGER)::BOOLEAN)",
2247 "((#0.0::INTEGER > 900::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 800::INTEGER)::BOOLEAN)",
2248 "((#0.0::INTEGER >= 1299::INTEGER)::BOOLEAN AND (#0.0::INTEGER <= 1299::INTEGER)::BOOLEAN)",
2249 "((#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN AND (#0.0::INTEGER >= 750::INTEGER)::BOOLEAN AND \
2250 (#0.0::INTEGER < 1000::INTEGER)::BOOLEAN)",
2251 "((#0.0::INTEGER >= 750::INTEGER)::BOOLEAN AND (#0.0::INTEGER >= 760::INTEGER)::BOOLEAN AND \
2252 (#0.0::INTEGER < 1000::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 990::INTEGER)::BOOLEAN)",
2253 ];
2254 for column in [flat, packed] {
2255 let chunk = Chunk::new(vec![column, words.clone()]).expect("two columns");
2256 // One end settled by the scan leaves the other to run alone, with nothing to pair.
2257 let alone = |predicate: &str, settled: Option<[bool; 2]>| {
2258 let (plan, list) = projection(&format!("{predicate}::BOOLEAN AS p"));
2259 let prepared =
2260 Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
2261 let mut scratch = prepared.scratch();
2262 match settled {
2263 Some(settled) => prepared.evaluate_settled(&chunk, &mut scratch, &settled),
2264 None => prepared.evaluate_filter(&chunk, &mut scratch),
2265 }
2266 .expect("the filter runs")
2267 };
2268 assert_eq!(
2269 alone(predicates[0], Some([true, false])),
2270 alone("(#0.0::INTEGER < 900::INTEGER)", None)
2271 );
2272 assert_eq!(
2273 alone(predicates[0], Some([false, true])),
2274 alone("(#0.0::INTEGER >= 800::INTEGER)", None)
2275 );
2276 for predicate in predicates {
2277 let (plan, list) = projection(&format!("{predicate}::BOOLEAN AS p"));
2278 let prepared =
2279 Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
2280 let mut scratch = prepared.scratch();
2281 let flags = evaluate(&plan, list[0], &schema, &chunk).expect("the tree walk runs");
2282 let expected =
2283 Selection::from_predicate(chunk.len(), |row| is_true(&flags.value_at(row)));
2284 // Enough chunks for the order to learn and move, which puts a different operand of
2285 // the pair in front.
2286 for _ in 0..40 {
2287 let threaded =
2288 prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs");
2289 assert_eq!(threaded, expected, "`{predicate}`");
2290 }
2291 }
2292 }
2293 }
2294
2295 /// A predicate with no `AND` in it is not threaded and has to keep saying the same thing.
2296 #[test]
2297 fn a_single_comparison_filters_the_same_rows() {
2298 filters("(#0.0::INTEGER > 1::INTEGER)::BOOLEAN");
2299 filters("(#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN");
2300 filters("(#0.0::INTEGER IS NOT DISTINCT FROM NULL::INTEGER)::BOOLEAN");
2301 }
2302
2303 #[test]
2304 fn a_chain_of_conjuncts_keeps_what_all_of_them_keep() {
2305 filters(
2306 "((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 3::INTEGER)::BOOLEAN)\
2307 ::BOOLEAN",
2308 );
2309 filters(
2310 "((#0.0::INTEGER >= 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER <= 3::INTEGER)::BOOLEAN \
2311 AND (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN AND (#0.0::INTEGER <> 2::INTEGER)\
2312 ::BOOLEAN)::BOOLEAN",
2313 );
2314 }
2315
2316 /// An operand the caller says is settled is not run, which shows as the rows it would have
2317 /// thrown away coming through: the answer is the other operand's alone. Settling nothing, or
2318 /// handing over the wrong number of operands, is the plain filter.
2319 #[test]
2320 fn a_settled_conjunct_is_left_out_of_the_filter() {
2321 let (schema, chunk) = input();
2322 let both = "((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 3::INTEGER)::BOOLEAN)\
2323 ::BOOLEAN AS p";
2324 let (plan, list) = projection(both);
2325 let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
2326 assert_eq!(prepared.conjuncts(), Some(2));
2327 let mut scratch = prepared.scratch();
2328 let wanted = |predicate: &str| {
2329 let (plan, list) = projection(&format!("{predicate} AS p"));
2330 let flags = evaluate(&plan, list[0], &schema, &chunk).expect("the tree walk runs");
2331 Selection::from_predicate(chunk.len(), |row| is_true(&flags.value_at(row)))
2332 };
2333 let second = prepared.evaluate_settled(&chunk, &mut scratch, &[true, false]);
2334 assert_eq!(
2335 second.expect("the filter runs"),
2336 wanted("(#0.0::INTEGER < 3::INTEGER)::BOOLEAN")
2337 );
2338 let first = prepared.evaluate_settled(&chunk, &mut scratch, &[false, true]);
2339 assert_eq!(
2340 first.expect("the filter runs"),
2341 wanted("(#0.0::INTEGER > 1::INTEGER)::BOOLEAN")
2342 );
2343 let neither = prepared.evaluate_settled(&chunk, &mut scratch, &[true, true]);
2344 assert_eq!(neither.expect("the filter runs"), Selection::identity(chunk.len()));
2345 let whole = wanted(&both[..both.len() - " AS p".len()]);
2346 let none = prepared.evaluate_settled(&chunk, &mut scratch, &[false, false]);
2347 assert_eq!(none.expect("the filter runs"), whole);
2348 let short = prepared.evaluate_settled(&chunk, &mut scratch, &[true]);
2349 assert_eq!(short.expect("the filter runs"), whole, "a list that does not fit is ignored");
2350 }
2351
2352 /// A conjunct that rejects every row, in front of one that would have kept some. The rows are
2353 /// the same either way and the point of the shape is that the second conjunct never runs.
2354 #[test]
2355 fn a_conjunct_that_keeps_nothing_ends_the_predicate() {
2356 filters(
2357 "((#0.0::INTEGER > 9::INTEGER)::BOOLEAN AND (#0.0::INTEGER < 9::INTEGER)::BOOLEAN)\
2358 ::BOOLEAN",
2359 );
2360 }
2361
2362 /// A conjunct whose operands are computed rather than read, which is the shape where the
2363 /// comparison is threaded and the arithmetic under it is not.
2364 #[test]
2365 fn a_conjunct_over_a_computed_operand_keeps_the_same_rows() {
2366 filters(
2367 "((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND \
2368 (\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER < 4::INTEGER)::BOOLEAN)::BOOLEAN",
2369 );
2370 }
2371
2372 /// A conjunct that is not a comparison at all, which is the one that goes through the flag
2373 /// kernel rather than the comparison kernel.
2374 #[test]
2375 fn a_conjunct_that_is_not_a_comparison_is_threaded_too() {
2376 filters(
2377 "((#0.0::INTEGER > 1::INTEGER)::BOOLEAN AND ((#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN \
2378 OR (#0.0::INTEGER = 1::INTEGER)::BOOLEAN)::BOOLEAN)::BOOLEAN",
2379 );
2380 filters(
2381 "(((#0.1::VARCHAR = 'c'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
2382 ::BOOLEAN AND (#0.0::INTEGER <> 1::INTEGER)::BOOLEAN)::BOOLEAN",
2383 );
2384 }
2385
2386 #[test]
2387 fn a_selective_conjunct_evaluates_later_like_on_its_survivors() {
2388 filters(
2389 "((#0.0::INTEGER > 2::INTEGER)::BOOLEAN AND \
2390 \"~~\"(#0.1::VARCHAR, '%a%'::VARCHAR)::BOOLEAN)::BOOLEAN",
2391 );
2392 filters(
2393 "((#0.0::INTEGER > 2::INTEGER)::BOOLEAN AND \
2394 \"!~~\"(#0.1::VARCHAR, '%a%'::VARCHAR)::BOOLEAN)::BOOLEAN",
2395 );
2396 }
2397
2398 /// An `OR` at the top threads the complement: the second branch only sees the rows the first
2399 /// one did not accept, and the rows it accepts are added to them rather than replacing them.
2400 ///
2401 /// The input has a row where the first branch is true, one where the second is, one where both
2402 /// are false and one where the first is null and the second is true, which is the row that says
2403 /// whether the complement was taken over "not true" or over "false".
2404 #[test]
2405 fn an_or_at_the_top_threads_the_complement() {
2406 filters(
2407 "((#0.0::INTEGER > 2::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'c'::VARCHAR)::BOOLEAN)\
2408 ::BOOLEAN",
2409 );
2410 filters(
2411 "((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN \
2412 OR (#0.0::INTEGER > 2::INTEGER)::BOOLEAN)::BOOLEAN",
2413 );
2414 }
2415
2416 /// A branch that accepts every row, in front of one that would have accepted none. The rows are
2417 /// the same either way and the point of the shape is that the second branch never runs.
2418 #[test]
2419 fn a_branch_that_keeps_everything_ends_the_predicate() {
2420 filters(
2421 "((#0.0::INTEGER IS NOT DISTINCT FROM #0.0::INTEGER)::BOOLEAN OR \
2422 (#0.0::INTEGER > 9::INTEGER)::BOOLEAN)::BOOLEAN",
2423 );
2424 }
2425
2426 /// The branches after one that has accepted every row really are skipped.
2427 ///
2428 /// Every other test here says the threaded answer matches the unthreaded one, which it would
2429 /// even if nothing were threaded at all. This one puts a division by zero behind a branch that
2430 /// accepts everything, so the predicate raises if the second branch runs and does not if the
2431 /// walk stopped where it was supposed to.
2432 #[test]
2433 fn a_branch_behind_one_that_accepted_every_row_does_not_run() {
2434 let (schema, chunk) = input();
2435 let predicate = "((#0.0::INTEGER IS NOT DISTINCT FROM #0.0::INTEGER)::BOOLEAN OR \
2436 (\"//\"(#0.0::INTEGER, 0::INTEGER)::INTEGER > 0::INTEGER)::BOOLEAN)\
2437 ::BOOLEAN";
2438 let (plan, list) = projection(&format!("{predicate} AS p"));
2439 let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
2440 let mut scratch = prepared.scratch();
2441 let kept =
2442 prepared.evaluate_filter(&chunk, &mut scratch).expect("the second branch never runs");
2443 assert_eq!(kept, Selection::identity(chunk.len()));
2444 // And the same predicate evaluated as an expression does divide by zero, which is what says
2445 // the test is testing the threading rather than a predicate that happens not to raise.
2446 evaluate(&plan, list[0], &schema, &chunk).expect_err("the tree walk divides by zero");
2447 }
2448
2449 /// The conjunct that rejects the most rows ends up in front of the one that rejects none.
2450 ///
2451 /// The predicate is written the wrong way round on purpose. The plan order costs two passes a
2452 /// chunk where one would do, and after a chunk of watching it the filter runs the selective one
2453 /// first and the other one stops running at all.
2454 #[test]
2455 fn a_filter_learns_which_conjunct_to_run_first() {
2456 let (schema, chunk) = input();
2457 let predicate = "((#0.0::INTEGER > 0::INTEGER)::BOOLEAN AND (#0.0::INTEGER > 9::INTEGER)\
2458 ::BOOLEAN)::BOOLEAN";
2459 let (plan, list) = projection(&format!("{predicate} AS p"));
2460 let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
2461 let mut scratch = prepared.scratch();
2462 let root = prepared.roots[0];
2463 assert_eq!(scratch.order(root), None, "nothing has run yet");
2464 let kept = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs");
2465 assert!(kept.is_empty());
2466 assert_eq!(scratch.order(root), Some(&[1, 0][..]), "the second conjunct rejects the most");
2467 // And it stays there, because the conjunct that now runs first empties the selection and
2468 // the one behind it keeps the history it already had rather than losing it.
2469 let kept = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs again");
2470 assert!(kept.is_empty());
2471 assert_eq!(scratch.order(root), Some(&[1, 0][..]));
2472 }
2473
2474 /// Whatever order it settles on, the rows are the rows.
2475 ///
2476 /// Run for longer than the window is wide, because an order that changes halfway through a scan
2477 /// is the shape where a walk that got the subtree bookkeeping wrong would start reading the
2478 /// wrong steps, and the first chunk would not show it.
2479 #[test]
2480 fn reordering_never_changes_which_rows_survive() {
2481 let (schema, chunk) = input();
2482 let predicate = "((#0.0::INTEGER >= 1::INTEGER)::BOOLEAN AND \
2483 (\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER < 4::INTEGER)::BOOLEAN AND \
2484 (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)::BOOLEAN";
2485 let (plan, list) = projection(&format!("{predicate} AS p"));
2486 let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
2487 let mut scratch = prepared.scratch();
2488 let flags = evaluate(&plan, list[0], &schema, &chunk).expect("the tree walk runs");
2489 let expected = Selection::from_predicate(chunk.len(), |row| is_true(&flags.value_at(row)));
2490 for round in 0..40 {
2491 let kept = prepared.evaluate_filter(&chunk, &mut scratch).expect("the filter runs");
2492 assert_eq!(kept, expected, "round {round}");
2493 }
2494 }
2495
2496 /// A nested connective is threaded rather than evaluated into flags.
2497 ///
2498 /// The inner `AND` keeps nothing, so its second conjunct is never reached and the division by
2499 /// zero in it never happens. Evaluating the branch as an expression and narrowing the flags
2500 /// afterwards, which is what an operand that is not a connective still does, would have run it.
2501 #[test]
2502 fn a_nested_connective_stops_where_the_outer_one_would() {
2503 let (schema, chunk) = input();
2504 let predicate = "((#0.0::INTEGER > 9::INTEGER)::BOOLEAN OR ((#0.0::INTEGER > 9::INTEGER)\
2505 ::BOOLEAN AND (\"//\"(#0.0::INTEGER, 0::INTEGER)::INTEGER > 0::INTEGER)\
2506 ::BOOLEAN)::BOOLEAN)::BOOLEAN";
2507 let (plan, list) = projection(&format!("{predicate} AS p"));
2508 let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
2509 let mut scratch = prepared.scratch();
2510 let kept =
2511 prepared.evaluate_filter(&chunk, &mut scratch).expect("the division never happens");
2512 assert!(kept.is_empty());
2513 evaluate(&plan, list[0], &schema, &chunk).expect_err("the tree walk divides by zero");
2514 }
2515
2516 /// A branch that is not a comparison, which is the one that goes through the flag kernel.
2517 #[test]
2518 fn an_or_branch_that_is_not_a_comparison_is_threaded_too() {
2519 filters(
2520 "((#0.0::INTEGER > 2::INTEGER)::BOOLEAN OR \
2521 \"~~\"(#0.1::VARCHAR, 'a%'::VARCHAR)::BOOLEAN)::BOOLEAN",
2522 );
2523 filters(
2524 "(\"~~\"(#0.1::VARCHAR, 'c%'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 1::INTEGER)\
2525 ::BOOLEAN)::BOOLEAN",
2526 );
2527 }
2528
2529 /// A connective inside a connective, which recurses rather than falling back to flags.
2530 ///
2531 /// Both nestings, because the two carry opposite things: an `AND` under an `OR` starts from the
2532 /// rows no branch has accepted, and an `OR` under an `AND` starts from the rows every conjunct
2533 /// has kept, and getting either one backwards is a wrong set of rows.
2534 #[test]
2535 fn a_connective_inside_a_connective_threads_both_ways() {
2536 filters(
2537 "(((#0.0::INTEGER >= 2::INTEGER)::BOOLEAN AND (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)\
2538 ::BOOLEAN OR ((#0.0::INTEGER < 2::INTEGER)::BOOLEAN AND (#0.1::VARCHAR <> 'c'\
2539 ::VARCHAR)::BOOLEAN)::BOOLEAN)::BOOLEAN",
2540 );
2541 filters(
2542 "(((#0.1::VARCHAR = 'c'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
2543 ::BOOLEAN AND ((#0.0::INTEGER <> 1::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'a'\
2544 ::VARCHAR)::BOOLEAN)::BOOLEAN)::BOOLEAN",
2545 );
2546 // Three deep, since two levels is where an off by one in the subtree bookkeeping can still
2547 // be hidden by the ranges lining up.
2548 filters(
2549 "((#0.0::INTEGER > 9::INTEGER)::BOOLEAN OR ((#0.0::INTEGER >= 1::INTEGER)::BOOLEAN \
2550 AND ((#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 1::INTEGER)\
2551 ::BOOLEAN)::BOOLEAN)::BOOLEAN)::BOOLEAN",
2552 );
2553 }
2554
2555 /// A predicate where one side is null and the other is true, in both orders. `OR` is true there
2556 /// and a complement taken over the rows a branch rejected rather than the rows it accepted
2557 /// would drop the row, which is the one way this can be wrong and is not a wrong vector but a
2558 /// missing row.
2559 #[test]
2560 fn a_null_branch_beside_a_true_one_keeps_the_row() {
2561 filters(
2562 "((#0.0::INTEGER > 2::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'c'::VARCHAR)::BOOLEAN \
2563 OR (#0.0::INTEGER IS NOT DISTINCT FROM NULL::INTEGER)::BOOLEAN)::BOOLEAN",
2564 );
2565 filters(
2566 "((#0.1::VARCHAR > 'b'::VARCHAR)::BOOLEAN OR (#0.0::INTEGER = 1::INTEGER)::BOOLEAN)\
2567 ::BOOLEAN",
2568 );
2569 }
2570
2571 /// A filter over a chunk that has already been narrowed, which is what a second filter in a
2572 /// pipeline sees and is the form pair the threaded kernels have to handle rather than fall
2573 /// through on.
2574 #[test]
2575 fn a_filter_over_a_selected_chunk_keeps_the_same_rows() {
2576 let (schema, chunk) = input();
2577 let predicate = "((#0.0::INTEGER >= 1::INTEGER)::BOOLEAN AND \
2578 (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)::BOOLEAN";
2579 let (plan, list) = projection(&format!("{predicate} AS p"));
2580 let prepared = Prepared::new(&plan, &list, &schema).expect("the predicate resolves");
2581 let mut scratch = prepared.scratch();
2582 let narrowed = narrow(&chunk, &[0, 3]).expect("two of the four rows");
2583 let threaded = prepared.evaluate_filter(&narrowed, &mut scratch).expect("the filter runs");
2584 let flags = evaluate(&plan, list[0], &schema, &narrowed).expect("the tree walk runs");
2585 let expected =
2586 Selection::from_predicate(narrowed.len(), |row| is_true(&flags.value_at(row)));
2587 assert_eq!(threaded, expected);
2588 }
2589
2590 /// Preparing is per pipeline and evaluating is per chunk, so the scratch has to survive being
2591 /// used again and give the same answer the second time.
2592 #[test]
2593 fn a_scratch_used_twice_gives_the_same_answer_twice() {
2594 let (schema, chunk) = input();
2595 let (plan, list) = projection("\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER AS a");
2596 let prepared = Prepared::new(&plan, &list, &schema).expect("the expressions resolve");
2597 let mut scratch = prepared.scratch();
2598 let mut once = Vec::new();
2599 prepared.evaluate(&chunk, &mut scratch, &mut once).expect("the first chunk runs");
2600 let mut twice = Vec::new();
2601 prepared.evaluate(&chunk, &mut scratch, &mut twice).expect("the second chunk runs");
2602 assert_eq!(once, twice);
2603 }
2604
2605 #[test]
2606 fn taking_the_chunk_answers_what_borrowing_it_does() {
2607 let (schema, chunk) = input();
2608 let (plan, list) = projection(
2609 "#0.0::INTEGER AS a, \"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER AS b, #0.0::INTEGER AS c",
2610 );
2611 let prepared = Prepared::new(&plan, &list, &schema).expect("the expressions resolve");
2612 let mut scratch = prepared.scratch();
2613 let mut borrowed = Vec::new();
2614 prepared.evaluate(&chunk, &mut scratch, &mut borrowed).expect("the borrowed chunk runs");
2615 let mut taken = Vec::new();
2616 prepared.evaluate_taking(chunk, &mut scratch, &mut taken).expect("the taken chunk runs");
2617 assert_eq!(borrowed, taken);
2618 }
2619
2620 #[test]
2621 fn a_shared_computed_root_is_compiled_once() {
2622 let (schema, chunk) = input();
2623 let (plan, list) = projection("\"+\"(#0.0::INTEGER, 1::INTEGER)::INTEGER AS a");
2624 let prepared = Prepared::shared(&plan, &[list[0], list[0]], &schema)
2625 .expect("the shared expression resolves");
2626 assert_eq!(prepared.steps.len(), 3);
2627 let mut scratch = prepared.scratch();
2628 let mut answers = Vec::new();
2629 prepared.evaluate(&chunk, &mut scratch, &mut answers).expect("both roots are returned");
2630 assert_eq!(answers[0], answers[1]);
2631 }
2632
2633 /// A chunk shorter than the last one, because a scan's final chunk is that and a constant
2634 /// materialized to the wrong length would be an out of range read rather than a wrong answer.
2635 #[test]
2636 fn a_shorter_chunk_after_a_longer_one_is_evaluated_at_its_own_length() {
2637 let (schema, chunk) = input();
2638 let (plan, list) = projection("7::INTEGER AS a");
2639 let prepared = Prepared::new(&plan, &list, &schema).expect("the expressions resolve");
2640 let mut scratch = prepared.scratch();
2641 let mut full = Vec::new();
2642 prepared.evaluate(&chunk, &mut scratch, &mut full).expect("the full chunk runs");
2643 assert_eq!(full[0].len(), 4);
2644 let short = chunk
2645 .clone()
2646 .select(&{
2647 let mut selection = Selection::with_capacity(2);
2648 selection.push(0);
2649 selection.push(2);
2650 selection
2651 })
2652 .expect("two of the four rows");
2653 let mut cut = Vec::new();
2654 prepared.evaluate(&short, &mut scratch, &mut cut).expect("the short chunk runs");
2655 assert_eq!(cut[0].len(), 2);
2656 }
2657
2658 /// An aggregate is not an expression and saying so when the pipeline is built is better than
2659 /// saying it on the first chunk.
2660 #[test]
2661 fn an_aggregate_is_refused_when_it_is_prepared() {
2662 let (schema, _) = input();
2663 let text = "Aggregate #1 groups=[] aggregates=[sum(#0.0::INTEGER)::HUGEINT]\n \
2664 Get memory.main.t AS t #0 [x::INTEGER, s::VARCHAR]";
2665 let plan = Plan::parse(text).expect("a well formed plan");
2666 let Node::Aggregate { aggregates, .. } = *plan.node(plan.root()) else {
2667 panic!("the root of that text is an aggregate");
2668 };
2669 let list = plan.expr_list(aggregates).to_vec();
2670 let error = Prepared::new(&plan, &list, &schema).expect_err("sum is not a scalar");
2671 assert!(error.message().contains("sum"), "{error}");
2672 }
2673
2674 /// How many of an expression's function steps worked something out when it was prepared, and
2675 /// whether the answer it gives is still the tree walk's answer.
2676 ///
2677 /// The count is the point of the assertion, because an answer that moved would be a bug. The
2678 /// agreement is what says the answer did not move.
2679 fn prepares(expr: &str, lifted: usize) {
2680 let (schema, _) = input();
2681 let projected = format!("{expr} AS a");
2682 let (plan, list) = projection(&projected);
2683 let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
2684 assert_eq!(prepared.hoisted(), lifted, "`{expr}`");
2685 agrees(&projected);
2686 }
2687
2688 /// A pattern the user wrote is compiled where the plan is, which is once.
2689 #[test]
2690 fn a_literal_pattern_is_compiled_when_the_pipeline_is_built() {
2691 prepares("\"~~\"(#0.1::VARCHAR, 'a%'::VARCHAR)::BOOLEAN", 1);
2692 prepares("\"~~*\"(#0.1::VARCHAR, '%A%'::VARCHAR)::BOOLEAN", 1);
2693 }
2694
2695 /// A regular expression, which is the one where the compiling is worth real time.
2696 ///
2697 /// ClickBench query 29 runs one pattern over a hundred million rows, which is a hundred thousand
2698 /// chunks, and before this each of those hundred thousand compiled the pattern again.
2699 #[test]
2700 fn a_regular_expression_is_compiled_when_the_pipeline_is_built() {
2701 prepares("\"regexp_matches\"(#0.1::VARCHAR, '^a'::VARCHAR)::BOOLEAN", 1);
2702 prepares("\"regexp_replace\"(#0.1::VARCHAR, 'a'::VARCHAR, 'b'::VARCHAR)::VARCHAR", 1);
2703 }
2704
2705 /// A pattern that is not a literal, which is legal SQL and is decided per chunk as it was.
2706 #[test]
2707 fn a_pattern_that_is_not_a_literal_is_left_to_the_chunk() {
2708 prepares("\"~~\"(#0.1::VARCHAR, #0.1::VARCHAR)::BOOLEAN", 0);
2709 }
2710
2711 /// A function with nothing to work out, which is almost all of them.
2712 #[test]
2713 fn a_function_with_no_prepare_step_prepares_nothing() {
2714 prepares("\"upper\"(#0.1::VARCHAR)::VARCHAR", 0);
2715 }
2716
2717 /// How many of an expression's steps are a folded `IN`, and whether the answer still agrees.
2718 fn folds(expr: &str, sets: usize) {
2719 let (schema, _) = input();
2720 let projected = format!("{expr} AS a");
2721 let (plan, list) = projection(&projected);
2722 let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
2723 assert_eq!(prepared.sets(), sets, "`{expr}`");
2724 agrees(&projected);
2725 }
2726
2727 /// What the binder writes for `x IN (1, 3)`, folded back into one lookup.
2728 ///
2729 /// The test goes through the plan's text, where the three mentions of the column are three
2730 /// expressions rather than one, which is the case `same` exists for. A plan the binder built has
2731 /// one mention and takes the first line of it.
2732 #[test]
2733 fn an_in_list_becomes_one_lookup() {
2734 folds(
2735 "((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
2736 ::BOOLEAN",
2737 1,
2738 );
2739 folds(
2740 "((#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN OR (#0.1::VARCHAR = 'z'::VARCHAR)::BOOLEAN)\
2741 ::BOOLEAN",
2742 1,
2743 );
2744 }
2745
2746 /// `NOT IN`, which the binder writes as an `AND` of inequalities and which reads the same
2747 /// lookup the other way round.
2748 #[test]
2749 fn a_not_in_list_becomes_the_same_lookup() {
2750 folds(
2751 "((#0.0::INTEGER <> 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER <> 3::INTEGER)::BOOLEAN)\
2752 ::BOOLEAN",
2753 1,
2754 );
2755 }
2756
2757 /// A list with a null in it, which is the rule that makes an `IN` not a set lookup.
2758 ///
2759 /// A row that is not in the list is null rather than false, because it might have equalled the
2760 /// value the null stands for. `agrees` is what says the fold kept that, since the `OR` of
2761 /// comparisons it is checked against gets it from three valued logic for free.
2762 #[test]
2763 fn a_list_with_a_null_in_it_folds_and_keeps_the_null_rule() {
2764 folds(
2765 "((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = NULL::INTEGER)::BOOLEAN \
2766 OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)::BOOLEAN",
2767 1,
2768 );
2769 folds(
2770 "((#0.0::INTEGER <> 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER <> NULL::INTEGER)\
2771 ::BOOLEAN AND (#0.0::INTEGER <> 3::INTEGER)::BOOLEAN)::BOOLEAN",
2772 1,
2773 );
2774 }
2775
2776 /// The connectives that are not an `IN`, each for its own reason.
2777 #[test]
2778 fn a_connective_that_is_not_an_in_list_is_left_alone() {
2779 // Two different columns.
2780 folds(
2781 "((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)\
2782 ::BOOLEAN",
2783 0,
2784 );
2785 // One equality and one of something else.
2786 folds(
2787 "((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER > 3::INTEGER)::BOOLEAN)\
2788 ::BOOLEAN",
2789 0,
2790 );
2791 // The right hand side is a column rather than a literal.
2792 folds(
2793 "((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = #0.0::INTEGER)::BOOLEAN)\
2794 ::BOOLEAN",
2795 0,
2796 );
2797 // An `AND` of equalities is not a `NOT IN`, it is a predicate that is false unless the two
2798 // literals are the same. Folding it as one would answer true where it answers false.
2799 folds(
2800 "((#0.0::INTEGER = 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
2801 ::BOOLEAN",
2802 0,
2803 );
2804 }
2805
2806 /// The same thing in a filter, which is the shape it is written in.
2807 #[test]
2808 fn an_in_list_filters_the_same_rows() {
2809 filters(
2810 "((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
2811 ::BOOLEAN",
2812 );
2813 filters(
2814 "((#0.0::INTEGER <> 1::INTEGER)::BOOLEAN AND (#0.0::INTEGER <> 3::INTEGER)::BOOLEAN)\
2815 ::BOOLEAN",
2816 );
2817 // Inside a larger predicate, where the fold is one operand of the connective above it.
2818 filters(
2819 "(((#0.0::INTEGER = 1::INTEGER)::BOOLEAN OR (#0.0::INTEGER = 3::INTEGER)::BOOLEAN)\
2820 ::BOOLEAN AND (#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN)::BOOLEAN",
2821 );
2822 }
2823
2824 /// The literal side of a comparison is turned into a column when the pipeline is built.
2825 #[test]
2826 fn a_comparison_against_a_literal_builds_it_once() {
2827 let (schema, _) = input();
2828 for (expr, built) in [
2829 ("(#0.1::VARCHAR = 'a'::VARCHAR)::BOOLEAN AS p", 1),
2830 ("(#0.0::INTEGER > 1::INTEGER)::BOOLEAN AS p", 1),
2831 // The literal on the left, which is the same comparison written the other way round.
2832 ("(1::INTEGER < #0.0::INTEGER)::BOOLEAN AS p", 1),
2833 // Two columns, which has no literal side to build.
2834 ("(#0.0::INTEGER = #0.0::INTEGER)::BOOLEAN AS p", 0),
2835 // Two literals, which the kernel answers once for the whole vector without reading a
2836 // column, so building one would be work that nothing reads.
2837 ("(1::INTEGER = 2::INTEGER)::BOOLEAN AS p", 0),
2838 ] {
2839 let (plan, list) = projection(expr);
2840 let prepared = Prepared::new(&plan, &list, &schema).expect("the expression resolves");
2841 assert_eq!(prepared.literals_built(), built, "`{expr}`");
2842 agrees(expr);
2843 }
2844 }
2845
2846 /// A pattern that does not compile still fails where the query said it does.
2847 ///
2848 /// Preparing is not allowed to move an error earlier. Compiling at build time and reporting
2849 /// there would raise before a row had been read, and under a `CASE` arm it would raise on a
2850 /// query whose rows never reach the call at all.
2851 #[test]
2852 fn a_pattern_that_does_not_compile_fails_on_the_chunk_and_not_before() {
2853 let (schema, chunk) = input();
2854 let (plan, list) =
2855 projection("\"regexp_matches\"(#0.1::VARCHAR, 'a('::VARCHAR)::BOOLEAN AS a");
2856 let prepared = Prepared::new(&plan, &list, &schema).expect("preparing does not compile it");
2857 assert_eq!(prepared.hoisted(), 0);
2858 let mut scratch = prepared.scratch();
2859 let mut out = Vec::new();
2860 prepared.evaluate(&chunk, &mut scratch, &mut out).expect_err("the chunk raises");
2861 }
2862}