use std::collections::HashMap;
use std::hash::BuildHasherDefault;
use std::sync::Arc;
use rudb_catalog::{Catalog, Parent, QualifiedName, Table};
use rudb_common::{Cancel, Error, Field, LogicalType, Memory, Result, Rule, Session, Value};
use rudb_functions::TableFunction;
use rudb_graph::Link;
use rudb_kernels::Accumulator;
use rudb_metrics::{Counters, Driver, Report};
use rudb_parquet::{Bound, Op};
use rudb_pipeline::{
BufferId, DynSink, DynStream, Pipeline, PipelineId, Source, Watched, root, root_in_order,
};
use rudb_plan::{
BuildSide, ColumnBinding, CompareOp, Expr, ExprRef, JoinKind, Node, NodeRef, PipelineRef, Plan,
ROOT, Shape, Slice, seams_of,
};
use rudb_seam::Settings;
use crate::buffer::Buffered;
use crate::cutoff::{self, Cutoff};
use crate::devicecard::device_card;
use crate::enginenames::{
database_size, dialects, extensions, grammar_extensions, optimizers, platform, user_agent,
version,
};
use crate::entrynames::{
columnnames, databasenames, schemanames, sequencenames, showdatabases, showtables,
showtablesexpanded, tablenames, viewnames,
};
use crate::fetch::{Fetch, TableFetch};
use crate::functionnames::functionnames;
use crate::gather::{Gather, Keep};
use crate::group::{Aggregate, Distinct};
use crate::join::{Broadcast, CrossProduct, Gathered, Join, Marking, Padding, Probe};
use crate::key::Digest;
use crate::keywords::keywords;
use crate::lateral::LateralSeries;
use crate::linkjoin::LinkJoin;
use crate::links::links;
use crate::percent::{LimitPercent, Portion};
use crate::prepared::Prepared;
use crate::query::Query;
use crate::register::registries;
use crate::schema::Schema;
use crate::setop::SetOp;
use crate::settingnames::settingnames;
use crate::sideways::{self, Exact, Keyed, Sideways};
use crate::sort::Sort;
use crate::source::{
Dummy, FileScan, Filters, Frequencies, Pushdown, Scan, Series, Summary, Values,
};
use crate::storagenames::storage_info;
use crate::strategies::strategies;
use crate::stream::{Edge, Filter, Limit, Project};
use crate::topn::TopN;
use crate::typenames::typenames;
use crate::window::{Window, Written};
use crate::writemetrics::{codec_metrics, write_metrics};
pub fn build<'a>(plan: &'a Plan, catalog: &'a Catalog) -> Result<Query<'a>> {
build_with(plan, catalog, &Cancel::new(), &Memory::unlimited(), &Settings::new())
}
pub fn build_with<'a>(
plan: &'a Plan,
catalog: &'a Catalog,
cancel: &Cancel,
memory: &Memory,
seams: &Settings,
) -> Result<Query<'a>> {
build_measured(plan, catalog, cancel, memory, seams, &Session::new(), &Report::new())
}
pub fn build_measured<'a>(
plan: &'a Plan,
catalog: &'a Catalog,
cancel: &Cancel,
memory: &Memory,
seams: &Settings,
session: &Session,
report: &Report,
) -> Result<Query<'a>> {
build_measured_with_sink(
plan,
catalog,
BuildUnder { cancel, memory, seams, session, report },
None,
)
}
pub fn build_measured_into<'a>(
plan: &'a Plan,
catalog: &'a Catalog,
cancel: &Cancel,
memory: &Memory,
seams: &Settings,
session: &Session,
sink: Arc<dyn DynSink + 'a>,
) -> Result<Query<'a>> {
let report = Report::new();
build_measured_with_sink(
plan,
catalog,
BuildUnder { cancel, memory, seams, session, report: &report },
Some(sink),
)
}
struct BuildUnder<'a> {
cancel: &'a Cancel,
memory: &'a Memory,
seams: &'a Settings,
session: &'a Session,
report: &'a Report,
}
#[derive(Clone, Copy, Default)]
struct AggregateBound {
max_groups: Option<usize>,
top_counts: Option<(usize, usize)>,
having_count: Option<(usize, i64)>,
}
fn build_measured_with_sink<'a>(
plan: &'a Plan,
catalog: &'a Catalog,
under: BuildUnder<'_>,
sink: Option<Arc<dyn DynSink + 'a>>,
) -> Result<Query<'a>> {
let BuildUnder { cancel, memory, seams, session, report } = under;
let shape = Shape::of(plan);
for pipeline in shape.all() {
report.pipeline(pipeline);
for waits_for in shape.waits_for(pipeline) {
report.depends(pipeline, *waits_for);
}
}
let mut building = Building {
plan,
catalog,
cancel,
memory,
seams,
session,
report,
shape,
done: Vec::new(),
drivers: Vec::new(),
pruning: Vec::new(),
pushing: None,
sideways: None,
above: Vec::new(),
armed: Vec::new(),
cutoff: None,
top_counts: Vec::new(),
marking: None,
held: Vec::new(),
};
let segment = building.node(plan.root())?;
let schema = segment.schema.clone();
let reader = if let Some(sink) = sink {
building.close(segment, ROOT, sink);
None
} else {
let (sink, reader) = if ordered(plan, plan.root()) {
root(BufferId(0), None)
} else {
root_in_order(BufferId(0), None)
};
building.close(segment, ROOT, Arc::new(sink));
Some(reader)
};
let Building { done, drivers, .. } = building;
Query::new(done, drivers, reader, schema)
}
fn ordered(plan: &Plan, node: NodeRef) -> bool {
match *plan.node(node) {
Node::Sort { .. } | Node::TopN { .. } => true,
Node::Project { input, .. }
| Node::Filter { input, .. }
| Node::Limit { input, .. }
| Node::LimitPercent { input, .. }
| Node::Fetch { input, .. } => ordered(plan, input),
_ => false,
}
}
fn count_top_aggregate(plan: &Plan, input: NodeRef, keys: Slice) -> Option<(NodeRef, usize)> {
let [key] = plan.sort_key_list(keys) else { return None };
if !key.descending {
return None;
}
let Expr::Column(ordered) = *plan.expr(key.expr) else { return None };
let mut aggregate = input;
let mut output = ordered;
loop {
match *plan.node(aggregate) {
Node::Project { input, index, exprs, .. } if output.table == index => {
let projected = *plan.expr_list(exprs).get(output.column as usize)?;
let Expr::Column(next) = *plan.expr(projected) else { return None };
output = next;
aggregate = input;
}
Node::Aggregate { index, .. } if output.table == index => break,
_ => return None,
}
}
let Node::Aggregate { index, groups, aggregates, .. } = *plan.node(aggregate) else {
return None;
};
if output.table != index {
return None;
}
let call = (output.column as usize).checked_sub(plan.expr_list(groups).len())?;
let aggregate_call = *plan.expr_list(aggregates).get(call)?;
let Expr::Aggregate { name, args, distinct, filter } = *plan.expr(aggregate_call) else {
return None;
};
let count_star = plan.string(name) == "count_star"
&& plan.expr_list(args).is_empty()
&& !distinct
&& filter.is_none();
let distinct_count = plan.string(name) == "count"
&& plan.expr_list(args).len() == 1
&& distinct
&& filter.is_none();
(count_star || distinct_count).then_some((aggregate, call))
}
fn count_having_aggregate(
plan: &Plan,
input: NodeRef,
predicate: ExprRef,
) -> Option<(NodeRef, usize, i64)> {
let Node::Aggregate { index, groups, aggregates, .. } = *plan.node(input) else { return None };
let Expr::Compare { op, left, right } = *plan.expr(predicate) else { return None };
let Expr::Column(column) = *plan.expr(left) else { return None };
let Expr::Constant(value) = *plan.expr(right) else { return None };
let Value::BigInt(value) = *plan.value(value) else { return None };
if column.table != index {
return None;
}
let call = (column.column as usize).checked_sub(plan.expr_list(groups).len())?;
let aggregate = *plan.expr_list(aggregates).get(call)?;
let Expr::Aggregate { name, args, distinct, filter } = *plan.expr(aggregate) else {
return None;
};
if plan.string(name) != "count_star"
|| !plan.expr_list(args).is_empty()
|| distinct
|| filter.is_some()
{
return None;
}
let minimum = match op {
CompareOp::Greater => value.checked_add(1)?,
CompareOp::GreaterOrEqual => value,
_ => return None,
};
Some((input, call, minimum))
}
fn mark_binding(plan: &Plan, right: NodeRef, kind: JoinKind) -> Option<usize> {
if kind != JoinKind::Mark {
return None;
}
let Node::Project { exprs, names, .. } = *plan.node(right) else {
return None;
};
let positions: Vec<usize> = plan
.expr_list(exprs)
.iter()
.enumerate()
.filter_map(|(position, &expr)| {
let Expr::Constant(value) = *plan.expr(expr) else {
return None;
};
(*plan.value(value) == Value::Boolean(true)).then_some(position)
})
.collect();
let position = match positions.as_slice() {
[position] => *position,
_ => plan
.name_list(names)
.iter()
.enumerate()
.rev()
.find_map(|(position, &name)| (plan.string(name) == "mark").then_some(position))?,
};
Some(position)
}
struct NativePairFrequencies {
entries: Vec<(Vec<Value>, u64)>,
}
fn native_host_groups(
plan: &Plan,
catalog: &Catalog,
input: NodeRef,
groups: Slice,
aggregates: Slice,
having: Option<(usize, i64)>,
) -> Result<Option<Vec<Vec<Value>>>> {
let Some((1, minimum)) = having else { return Ok(None) };
let Ok(minimum) = u64::try_from(minimum) else { return Ok(None) };
let Node::Filter { input: source, predicate } = *plan.node(input) else {
return Ok(None);
};
let Node::Get { catalog: database, schema, table, index, columns, .. } = *plan.node(source)
else {
return Ok(None);
};
let [group] = plan.expr_list(groups) else { return Ok(None) };
let Expr::Function { name, args } = *plan.expr(*group) else { return Ok(None) };
if plan.string(name) != "regexp_replace" {
return Ok(None);
}
let [subject, pattern, replacement] = plan.expr_list(args) else { return Ok(None) };
let Expr::Column(binding) = *plan.expr(*subject) else { return Ok(None) };
if binding.table != index {
return Ok(None);
}
let (Expr::Constant(pattern), Expr::Constant(replacement)) =
(plan.expr(*pattern), plan.expr(*replacement))
else {
return Ok(None);
};
if plan.value(*pattern) != &Value::Varchar("^https?://(?:www\\.)?([^/]+)/.*$".into())
|| plan.value(*replacement) != &Value::Varchar("\\1".into())
{
return Ok(None);
}
let Expr::Compare { op: CompareOp::NotEqual, left, right } = *plan.expr(predicate) else {
return Ok(None);
};
let filtered = match (plan.expr(left), plan.expr(right)) {
(Expr::Column(held), Expr::Constant(value))
if plan.value(*value) == &Value::Varchar(String::new()) =>
{
held
}
(Expr::Constant(value), Expr::Column(held))
if plan.value(*value) == &Value::Varchar(String::new()) =>
{
held
}
_ => return Ok(None),
};
if filtered != &binding {
return Ok(None);
}
let [average, count, minimum_value] = plan.expr_list(aggregates) else { return Ok(None) };
let check = |reference: &ExprRef, wanted: &str, argument: Option<ExprRef>| {
let Expr::Aggregate { name, args, distinct: false, filter: None } = *plan.expr(*reference)
else {
return false;
};
plan.string(name) == wanted
&& match argument {
None => plan.expr_list(args).is_empty(),
Some(argument) => plan.expr_list(args) == [argument],
}
};
if !check(count, "count_star", None) || !check(minimum_value, "min", Some(*subject)) {
return Ok(None);
}
let Expr::Aggregate { name, args, distinct: false, filter: None } = *plan.expr(*average) else {
return Ok(None);
};
if plan.string(name) != "avg" || plan.expr_list(args).len() != 1 {
return Ok(None);
}
let mut length = plan.expr_list(args)[0];
if let Expr::Cast { input, try_cast: false } = *plan.expr(length) {
length = input;
}
let Expr::Function { name, args } = *plan.expr(length) else { return Ok(None) };
if plan.string(name) != "strlen" || plan.expr_list(args) != [*subject] {
return Ok(None);
}
let Some(field) = plan.field_list(columns).get(binding.column as usize) else {
return Ok(None);
};
let name = QualifiedName::new(plan.string(database), plan.string(schema), plan.string(table));
let table = catalog.table(&name)?;
let Some(column) = table.column_index(&field.name) else { return Ok(None) };
let Some(entries) = table.rows().host_groups(column, minimum)? else { return Ok(None) };
let mut records = Vec::with_capacity(entries.len());
for entry in entries {
let count = i64::try_from(entry.count)
.map_err(|_| Error::internal("a stored host count exceeds BIGINT"))?;
records.push(vec![
Value::Varchar(entry.host),
Accumulator::exact_avg(entry.bytes_sum, count, &LogicalType::Double).finish()?,
Value::BigInt(count),
Value::Varchar(entry.minimum),
]);
}
Ok(Some(records))
}
fn native_pair_frequencies(
plan: &Plan,
catalog: &Catalog,
input: NodeRef,
groups: Slice,
aggregates: Slice,
top: usize,
) -> Result<Option<NativePairFrequencies>> {
let Node::Get { catalog: database, schema, table, index, columns, .. } = *plan.node(input)
else {
return Ok(None);
};
let [first_expr, second_expr] = plan.expr_list(groups) else { return Ok(None) };
let Expr::Column(first) = *plan.expr(*first_expr) else { return Ok(None) };
let Expr::Column(second) = *plan.expr(*second_expr) else { return Ok(None) };
if first.table != index
|| second.table != index
|| plan.expr_type(*first_expr) != &LogicalType::BigInt
|| plan.expr_type(*second_expr) != &LogicalType::Varchar
{
return Ok(None);
}
let [aggregate] = plan.expr_list(aggregates) else { return Ok(None) };
let Expr::Aggregate { name, args, distinct, filter } = *plan.expr(*aggregate) else {
return Ok(None);
};
if top == 0
|| plan.string(name) != "count_star"
|| !plan.expr_list(args).is_empty()
|| distinct
|| filter.is_some()
{
return Ok(None);
}
let fields = plan.field_list(columns);
let Some(first_field) = fields.get(first.column as usize) else { return Ok(None) };
let Some(second_field) = fields.get(second.column as usize) else { return Ok(None) };
let name = QualifiedName::new(plan.string(database), plan.string(schema), plan.string(table));
let table = catalog.table(&name)?;
let Some(first_column) = table.column_index(&first_field.name) else { return Ok(None) };
let Some(second_column) = table.column_index(&second_field.name) else { return Ok(None) };
if let Some(entries) = table.rows().top_pair_frequencies(first_column, second_column, top)? {
return Ok(Some(NativePairFrequencies { entries }));
}
let Some(occurrences) = table.rows().frequency_occurrences(first_column)? else {
return Ok(None);
};
let (anchors, anchor_indices, second_values, dictionary) = if !occurrences
.anchor_indices
.is_empty()
&& occurrences.anchor_indices.len() == occurrences.ordinals.len()
{
let Some((second_values, dictionary)) =
table.rows().stable_codes_at(second_column, &occurrences.ordinals)?
else {
return Ok(None);
};
if occurrences.anchors.iter().any(|value| !matches!(value, Value::BigInt(_) | Value::Null))
|| occurrences
.anchor_indices
.iter()
.any(|&entry| entry as usize >= occurrences.anchors.len())
{
return Err(Error::internal("a BIGINT frequency anchor has another type"));
}
(occurrences.anchors, occurrences.anchor_indices, second_values, dictionary)
} else {
let Some(rows) = table.rows().stable_pair_codes_at(
first_column,
second_column,
&occurrences.ordinals,
)?
else {
return Ok(None);
};
let mut anchors = Vec::new();
let mut by_anchor = HashMap::<Option<i64>, u16, BuildHasherDefault<Digest>>::default();
let mut anchor_indices = Vec::with_capacity(rows.first.len());
for value in rows.first {
let value = value
.map(|value| {
i64::try_from(value).map_err(|_| {
Error::internal("a BIGINT frequency occurrence is out of range")
})
})
.transpose()?;
let entry = match by_anchor.get(&value) {
Some(&entry) => entry,
None => {
let entry = u16::try_from(anchors.len())
.map_err(|_| Error::internal("too many frequency anchors"))?;
anchors.push(value.map_or(Value::Null, Value::BigInt));
by_anchor.insert(value, entry);
entry
}
};
anchor_indices.push(entry);
}
(anchors, anchor_indices, rows.second, rows.dictionary)
};
if anchor_indices.len() != second_values.len() {
return Err(Error::internal("a stable pair fetch returned columns of different lengths"));
}
let mut counts = HashMap::<(u16, Option<u32>), u64, BuildHasherDefault<Digest>>::default();
for (anchor, second) in anchor_indices.into_iter().zip(second_values) {
*counts.entry((anchor, second)).or_default() += 1;
}
let mut boundaries = counts.values().copied().collect::<Vec<_>>();
if boundaries.len() < top {
return Ok(None);
}
boundaries.select_nth_unstable_by(top - 1, |left, right| right.cmp(left));
let boundary = boundaries[top - 1];
if boundary <= occurrences.omitted_max {
return Ok(None);
}
let mut entries = Vec::new();
for ((anchor, second), count) in counts {
if count < boundary {
continue;
}
let first = anchors
.get(anchor as usize)
.cloned()
.ok_or_else(|| Error::internal("a frequency anchor index is outside its values"))?;
let second = match second {
Some(code) => Value::Varchar(
dictionary
.try_text_at(code as usize)?
.ok_or_else(|| Error::internal("a string frequency code is null"))?
.to_owned(),
),
None => Value::Null,
};
entries.push((vec![first, second], count));
}
Ok(Some(NativePairFrequencies { entries }))
}
fn edge(plan: &Plan, bound: rudb_plan::Bound, input: &Schema) -> Result<Edge> {
Ok(match bound {
rudb_plan::Bound::All => Edge::All,
rudb_plan::Bound::Rows(rows) => Edge::Rows(rows),
rudb_plan::Bound::Read(expr) => Edge::Read(Prepared::one(plan, expr, input)?),
})
}
fn portion(plan: &Plan, share: rudb_plan::Share, input: &Schema) -> Result<Portion> {
Ok(match share {
rudb_plan::Share::Percent(percent) => Portion::Percent(percent),
rudb_plan::Share::Read(expr) => Portion::Read(Prepared::one(plan, expr, input)?),
})
}
fn whole_table<'a>(
plan: &Plan,
catalog: &'a Catalog,
node: NodeRef,
) -> Result<Option<(&'a Table, u32, Slice)>> {
let Node::Get { catalog: database, schema, table, index, columns, .. } = *plan.node(node)
else {
return Ok(None);
};
let name = QualifiedName::new(plan.string(database), plan.string(schema), plan.string(table));
let table = catalog.table(&name)?;
Ok(Some((table, index, columns)))
}
fn scanned<'a>(
plan: &Plan,
catalog: &'a Catalog,
node: NodeRef,
index: u32,
) -> Result<Option<(&'a Table, Slice)>> {
if let Some((table, found, columns)) = whole_table(plan, catalog, node)? {
if found == index {
return Ok(Some((table, columns)));
}
}
for child in plan.node(node).children().into_iter().flatten() {
if let Some(found) = scanned(plan, catalog, child, index)? {
return Ok(Some(found));
}
}
Ok(None)
}
fn exact(
plan: &Plan,
catalog: &Catalog,
parent: NodeRef,
key: ExprRef,
driving: NodeRef,
binding: ColumnBinding,
) -> Option<Exact> {
let Expr::Column(key) = *plan.expr(key) else { return None };
let key = traced(plan, parent, key)?;
let (parent_table, parent_columns) = scanned(plan, catalog, parent, key.table).ok()??;
let (child_table, child_columns) = scanned(plan, catalog, driving, binding.table).ok()??;
let parent_rows = parent_table.rows().stored()?;
let parent_column = stored_column(plan, parent_table, key.table, parent_columns, key)?;
let child_column = stored_column(plan, child_table, binding.table, child_columns, binding)?;
let keys = rudb_native::graph::key_map(parent_rows, parent_column)?;
let link = child_table.rows().stored().and_then(|child_rows| {
let edge = rudb_native::graph::Edge {
child: child_table.name().table.clone(),
child_column,
parent: parent_table.name().table.clone(),
parent_column,
};
rudb_native::graph::stored_link(child_rows, parent_rows, &edge)
});
Some(Exact::new(keys, link))
}
fn traced(plan: &Plan, node: NodeRef, binding: ColumnBinding) -> Option<ColumnBinding> {
match *plan.node(node) {
Node::Get { index, .. } => (binding.table == index).then_some(binding),
Node::Project { input, index, exprs, .. } => {
let binding = if binding.table == index {
let at = *plan.expr_list(exprs).get(binding.column as usize)?;
let Expr::Column(inner) = *plan.expr(at) else { return None };
inner
} else {
binding
};
traced(plan, input, binding)
}
_ => plan
.node(node)
.children()
.into_iter()
.flatten()
.find_map(|child| traced(plan, child, binding)),
}
}
fn equated(plan: &Plan, condition: ExprRef) -> Option<[ColumnBinding; 2]> {
let Expr::Compare { op: CompareOp::Equal, left, right } = *plan.expr(condition) else {
return None;
};
match (plan.expr(left), plan.expr(right)) {
(&Expr::Column(left), &Expr::Column(right)) => Some([left, right]),
_ => None,
}
}
fn stored_column(
plan: &Plan,
table: &Table,
index: u32,
columns: Slice,
binding: ColumnBinding,
) -> Option<usize> {
if binding.table != index {
return None;
}
let field = plan.field_list(columns).get(binding.column as usize)?;
table.column_index(&field.name)
}
fn grouped_column<'a>(
plan: &Plan,
catalog: &'a Catalog,
node: NodeRef,
) -> Result<Option<(&'a Table, u32, usize)>> {
let Node::Aggregate { input, index: produced, groups, aggregates } = *plan.node(node) else {
return Ok(None);
};
if !plan.expr_list(aggregates).is_empty() {
return Ok(None);
}
let [group] = plan.expr_list(groups) else { return Ok(None) };
let Expr::Column(binding) = *plan.expr(*group) else { return Ok(None) };
let Some((table, index, columns)) = whole_table(plan, catalog, input)? else {
return Ok(None);
};
Ok(stored_column(plan, table, index, columns, binding).map(|column| (table, produced, column)))
}
struct CertainFilter {
entries: Vec<(Value, u64)>,
omitted_max: u64,
against: Value,
differs: bool,
}
impl CertainFilter {
fn rows(&self) -> Option<u64> {
if self.omitted_max != 0 {
return None;
}
let mut kept = 0_u64;
for (value, count) in &self.entries {
if self.keeps(value)? {
kept = kept.checked_add(*count)?;
}
}
Some(kept)
}
fn keeps(&self, value: &Value) -> Option<bool> {
if value.is_null() {
return Some(false);
}
if value.logical_type() != self.against.logical_type() {
return None;
}
Some((value == &self.against) != self.differs)
}
}
fn certain_filter(plan: &Plan, catalog: &Catalog, node: NodeRef) -> Result<Option<CertainFilter>> {
let Node::Filter { input, predicate } = *plan.node(node) else { return Ok(None) };
let Some((table, index, columns)) = whole_table(plan, catalog, input)? else {
return Ok(None);
};
let Expr::Compare { op, left, right } = *plan.expr(predicate) else { return Ok(None) };
let differs = match op {
CompareOp::Equal => false,
CompareOp::NotEqual => true,
_ => return Ok(None),
};
let (binding, constant) = match (plan.expr(left), plan.expr(right)) {
(&Expr::Column(binding), &Expr::Constant(value))
| (&Expr::Constant(value), &Expr::Column(binding)) => (binding, value),
_ => return Ok(None),
};
let against = plan.value(constant).clone();
if against.is_null() {
return Ok(None);
}
let Some(column) = stored_column(plan, table, index, columns, binding) else {
return Ok(None);
};
let Some(prefix) = table.rows().frequency_prefix(column)? else { return Ok(None) };
let (entries, omitted_max) = (prefix.entries, prefix.omitted_max);
Ok(Some(CertainFilter { entries, omitted_max, against, differs }))
}
fn known_rows(plan: &Plan, catalog: &Catalog, node: NodeRef) -> Result<Option<u64>> {
if let Some((table, _, _)) = whole_table(plan, catalog, node)? {
return Ok(Some(table.rows().len() as u64));
}
if let Some(filter) = certain_filter(plan, catalog, node)? {
return Ok(filter.rows());
}
let Some((table, _, column)) = grouped_column(plan, catalog, node)? else { return Ok(None) };
let Some(distinct) = table.rows().distinct_values(column)? else { return Ok(None) };
let Some(nulls) = table.rows().null_count(column)? else { return Ok(None) };
Ok(Some(distinct.saturating_add(u64::from(nulls > 0))))
}
fn stored_summary(
plan: &Plan,
catalog: &Catalog,
input: NodeRef,
groups: Slice,
aggregates: Slice,
) -> Result<Option<Vec<Value>>> {
if !plan.expr_list(groups).is_empty() || plan.expr_list(aggregates).is_empty() {
return Ok(None);
}
let below = whole_table(plan, catalog, input)?;
let mut values = Vec::with_capacity(plan.expr_list(aggregates).len());
for &aggregate in plan.expr_list(aggregates) {
let Expr::Aggregate { name, args, distinct, filter: None } = *plan.expr(aggregate) else {
return Ok(None);
};
let call = plan.string(name);
let args = plan.expr_list(args);
if distinct {
if call != "count" {
return Ok(None);
}
let [only] = args else { return Ok(None) };
let Expr::Column(binding) = *plan.expr(*only) else { return Ok(None) };
let Some((table, index, columns)) = below else { return Ok(None) };
let Some(column) = stored_column(plan, table, index, columns, binding) else {
return Ok(None);
};
let Some(counted) = table.rows().distinct_values(column)? else { return Ok(None) };
values.push(count(counted)?);
continue;
}
if call == "count_star" && args.is_empty() {
let Some(rows) = known_rows(plan, catalog, input)? else { return Ok(None) };
values.push(count(rows)?);
continue;
}
let [only] = args else { return Ok(None) };
let Expr::Column(binding) = *plan.expr(*only) else { return Ok(None) };
if call == "count" {
if let Some((table, produced, column)) = grouped_column(plan, catalog, input)? {
if binding.table == produced && binding.column == 0 {
let Some(distinct) = table.rows().distinct_values(column)? else {
return Ok(None);
};
values.push(count(distinct)?);
continue;
}
}
}
let Some((table, index, columns)) = below else { return Ok(None) };
let Some(column) = stored_column(plan, table, index, columns, binding) else {
return Ok(None);
};
match call {
"count" => {
let Some(nulls) = table.rows().null_count(column)? else { return Ok(None) };
values.push(count(table.rows().len() as u64 - nulls)?);
}
"min" | "max" => {
let Some(value) = extreme(table, column, call == "min")? else { return Ok(None) };
values.push(value);
}
"sum" | "avg" => {
let Some(field) = table.columns().get(column) else { return Ok(None) };
if !field.ty.is_integer() {
return Ok(None);
}
let Some((total, rows)) = table.rows().exact_sum(column)? else { return Ok(None) };
let returns = plan.expr_type(aggregate);
let state = if call == "sum" {
Accumulator::exact_sum(total, rows > 0, returns)
} else {
let Ok(seen) = i64::try_from(rows) else { return Ok(None) };
Accumulator::exact_avg(total, seen, returns)
};
values.push(state.finish()?);
}
_ => return Ok(None),
}
}
Ok(Some(values))
}
fn summary_schema(plan: &Plan, index: u32, aggregates: Slice) -> Result<Schema> {
let mut fields = Vec::with_capacity(plan.expr_list(aggregates).len());
for &reference in plan.expr_list(aggregates) {
let Expr::Aggregate { name, .. } = *plan.expr(reference) else {
return Err(Error::internal("an aggregate list holds something that is not a call"));
};
fields.push(Field::new(plan.string(name).to_string(), plan.expr_type(reference).clone()));
}
Ok(Schema::numbered(fields, index))
}
fn extreme(table: &Table, column: usize, smallest: bool) -> Result<Option<Value>> {
if let Some((low, high)) = table.rows().text_extremes(column)? {
return Ok(Some(if smallest { low } else { high }));
}
let Some((low, high)) = table.rows().exact_extremes(column)? else { return Ok(None) };
let Some(field) = table.columns().get(column) else { return Ok(None) };
Ok(if smallest { low } else { high }.into_value(&field.ty))
}
fn count(rows: u64) -> Result<Value> {
Ok(Value::BigInt(
i64::try_from(rows).map_err(|_| Error::internal("a stored row count exceeds BIGINT"))?,
))
}
struct Segment<'a> {
source: Arc<dyn Source + 'a>,
streams: Vec<Arc<dyn DynStream + 'a>>,
schema: Schema,
after: Vec<PipelineRef>,
}
impl<'a> Segment<'a> {
fn new(source: Arc<dyn Source + 'a>, schema: Schema) -> Self {
Self { source, streams: Vec::new(), schema, after: Vec::new() }
}
fn reading(source: Arc<dyn Source + 'a>, schema: Schema, after: PipelineRef) -> Self {
Self { source, streams: Vec::new(), schema, after: vec![after] }
}
fn then(mut self, stream: Arc<dyn DynStream + 'a>, schema: Schema) -> Self {
self.streams.push(stream);
self.schema = schema;
self
}
}
struct Building<'a, 'b> {
plan: &'a Plan,
catalog: &'a Catalog,
cancel: &'b Cancel,
memory: &'b Memory,
seams: &'b Settings,
session: &'b Session,
report: &'b Report,
shape: Shape,
done: Vec<Pipeline<'a>>,
drivers: Vec<Arc<Driver>>,
pruning: Vec<(usize, Op, Bound)>,
pushing: Option<Pushdown>,
sideways: Option<Arc<Sideways<'a>>>,
above: Vec<Arc<Sideways<'a>>>,
armed: Vec<Arc<Sideways<'a>>>,
cutoff: Option<Arc<Cutoff>>,
top_counts: Vec<(NodeRef, usize, usize)>,
marking: Option<NodeRef>,
held: Vec<Held>,
}
struct Linked {
link: Arc<Link>,
parent: Arc<Parent>,
projected: Vec<(usize, LogicalType)>,
parent_schema: Schema,
}
struct Held {
cte: u32,
chunks: Buffered,
filling: PipelineRef,
}
fn profiling(session: &Session) -> bool {
session.get("enable_profiling").is_some_and(|format| format != rudb_functions::UNSET)
}
fn pragma_name(plan: &Plan, args: Slice) -> Result<String> {
let [argument] = plan.expr_list(args) else {
return Err(Error::internal("a pragma that resolved to more than one name"));
};
let Expr::Constant(reference) = *plan.expr(*argument) else {
return Err(Error::not_implemented(
"pragma_storage_info() given a name that is not a constant",
));
};
match plan.value(reference) {
Value::Varchar(name) => Ok(name.clone()),
Value::Null => Ok("NULL".to_string()),
other => Err(Error::internal(format!("a pragma name bound as VARCHAR arrived as {other}"))),
}
}
impl<'a> Building<'a, '_> {
fn gathered(&self, node: NodeRef) -> u32 {
self.shape.gathered(node).expect("a node with two inputs has a second operator")
}
fn close(&mut self, segment: Segment<'a>, id: PipelineRef, sink: Arc<dyn DynSink + 'a>) {
let mut pipeline = Pipeline::new(PipelineId(id), segment.source, sink);
for stream in segment.streams {
pipeline = pipeline.then(stream);
}
for after in segment.after {
pipeline = pipeline.after(PipelineId(after));
}
self.done.push(pipeline);
self.drivers.push(self.report.driving(id));
}
fn watch(
&self,
node: NodeRef,
id: u32,
pipeline: u32,
kind: &str,
detail: Option<&str>,
) -> Arc<Counters> {
self.watch_doing(node, None, id, pipeline, kind, detail)
}
fn watch_doing(
&self,
node: NodeRef,
also: Option<NodeRef>,
id: u32,
pipeline: u32,
kind: &str,
detail: Option<&str>,
) -> Arc<Counters> {
let mut counters = Counters::new(id, pipeline, kind)
.charging_cpu(profiling(self.session))
.under(self.consumer(id, also));
if let Some(detail) = detail {
counters = counters.detailed(detail);
}
for node in std::iter::once(node).chain(also) {
for seam in seams_of(self.plan.node(node)) {
if let Some(running) = registries().running(*seam, self.seams) {
counters = counters.chose(seam.name(), &running.name, running.is_reference);
}
}
}
self.report.watch(counters)
}
fn consumer(&self, id: u32, also: Option<NodeRef>) -> Option<u32> {
let swallowed = also.map(|node| self.shape.operator(node));
let mut parent = self.shape.consumer(id);
while parent.is_some() && parent == swallowed {
parent = self.shape.consumer(parent?);
}
parent
}
fn table_function(&mut self, reference: NodeRef) -> Result<Segment<'a>> {
let plan = self.plan;
let id = self.shape.operator(reference);
let pipeline = self.shape.pipeline(reference);
let Node::TableFunction { index, function, args, options, settings, columns } =
*plan.node(reference)
else {
return Err(Error::internal("a table function was built from a node that is not one"));
};
let name = plan.string(function);
let runtime = self.sideways.take();
self.above.clear();
Ok(match TableFunction::lookup(name) {
Some(function @ (TableFunction::ReadParquet | TableFunction::ReadCsv)) => {
let counters = self.watch(reference, id, pipeline, "FileScan", Some(name));
let tests = std::mem::take(&mut self.pruning);
let scan = FileScan::new(
plan, index, function, args, options, settings, columns, tests, runtime,
)?
.watched(counters.clone());
let schema = scan.schema().clone();
Segment::new(Arc::new(Watched::new(scan, counters)), schema)
}
Some(TableFunction::PragmaStorageInfo) => {
let written = pragma_name(plan, args)?;
let table = storage_info(self.catalog, &written, plan, index, columns)?;
let schema = table.schema().clone();
let counters = self.watch(
reference,
id,
pipeline,
"Metadata",
Some(TableFunction::PragmaStorageInfo.name()),
);
Segment::new(Arc::new(Watched::new(table, counters)), schema)
}
Some(TableFunction::RudbDeviceCard) => {
let table = device_card(plan, args, index, columns)?;
let schema = table.schema().clone();
let counters = self.watch(
reference,
id,
pipeline,
"Metadata",
Some(TableFunction::RudbDeviceCard.name()),
);
Segment::new(Arc::new(Watched::new(table, counters)), schema)
}
Some(
function @ (TableFunction::RudbStrategies
| TableFunction::RudbLinks
| TableFunction::RudbWriteMetrics
| TableFunction::RudbCodecMetrics
| TableFunction::DuckdbKeywords
| TableFunction::DuckdbTypes
| TableFunction::DuckdbFunctions
| TableFunction::DuckdbSettings
| TableFunction::DuckdbDatabases
| TableFunction::DuckdbSchemas
| TableFunction::DuckdbTables
| TableFunction::DuckdbViews
| TableFunction::DuckdbSequences
| TableFunction::DuckdbColumns
| TableFunction::DuckdbExtensions
| TableFunction::DuckdbOptimizers
| TableFunction::DuckdbDialects
| TableFunction::DuckdbGrammarExtensions
| TableFunction::PragmaVersion
| TableFunction::PragmaPlatform
| TableFunction::PragmaUserAgent
| TableFunction::PragmaDatabaseSize
| TableFunction::PragmaShowTables
| TableFunction::PragmaShowDatabases
| TableFunction::PragmaShowTablesExpanded),
) => {
let table = match function {
TableFunction::RudbLinks => {
links(self.session, self.catalog, plan, index, columns)?
}
TableFunction::RudbWriteMetrics => write_metrics(plan, index, columns)?,
TableFunction::RudbCodecMetrics => codec_metrics(plan, index, columns)?,
TableFunction::DuckdbKeywords => keywords(plan, index, columns)?,
TableFunction::DuckdbTypes => typenames(plan, index, columns)?,
TableFunction::DuckdbFunctions => functionnames(plan, index, columns)?,
TableFunction::DuckdbSettings => {
settingnames(self.session, plan, index, columns)?
}
TableFunction::DuckdbDatabases => {
databasenames(self.catalog, plan, index, columns)?
}
TableFunction::DuckdbSchemas => {
schemanames(self.catalog, plan, index, columns)?
}
TableFunction::DuckdbTables => tablenames(self.catalog, plan, index, columns)?,
TableFunction::DuckdbViews => viewnames(self.catalog, plan, index, columns)?,
TableFunction::DuckdbSequences => {
sequencenames(self.catalog, plan, index, columns)?
}
TableFunction::DuckdbColumns => {
columnnames(self.catalog, plan, index, columns)?
}
TableFunction::DuckdbExtensions => extensions(plan, index, columns)?,
TableFunction::DuckdbOptimizers => optimizers(plan, index, columns)?,
TableFunction::DuckdbDialects => dialects(plan, index, columns)?,
TableFunction::DuckdbGrammarExtensions => {
grammar_extensions(plan, index, columns)?
}
TableFunction::PragmaVersion => version(plan, index, columns)?,
TableFunction::PragmaPlatform => platform(plan, index, columns)?,
TableFunction::PragmaUserAgent => user_agent(plan, index, columns)?,
TableFunction::PragmaDatabaseSize => {
database_size(self.catalog, self.memory, plan, index, columns)?
}
TableFunction::PragmaShowTables => {
showtables(self.catalog, plan, index, columns)?
}
TableFunction::PragmaShowDatabases => {
showdatabases(self.catalog, plan, index, columns)?
}
TableFunction::PragmaShowTablesExpanded => {
showtablesexpanded(self.catalog, plan, index, columns)?
}
_ => strategies(plan, index, columns)?,
};
let schema = table.schema().clone();
let counters =
self.watch(reference, id, pipeline, "Metadata", Some(function.name()));
Segment::new(Arc::new(Watched::new(table, counters)), schema)
}
_ => {
let series = Series::new(plan, index, name, args)?;
let schema = series.schema().clone();
let counters = self.watch(reference, id, pipeline, "Series", Some(name));
Segment::new(Arc::new(Watched::new(series, counters)), schema)
}
})
}
fn join(&mut self, reference: NodeRef) -> Result<Segment<'a>> {
let plan = self.plan;
let memory = self.memory;
let id = self.shape.operator(reference);
let pipeline = self.shape.pipeline(reference);
let Node::Join { left, right, kind, conditions, build } = *plan.node(reference) else {
return Err(Error::internal("a join was built from a node that is not one"));
};
let marker = mark_binding(plan, right, kind);
let swapped = build == BuildSide::Left;
let (held, driving) = if swapped { (left, right) } else { (right, left) };
let marking = swapped && matches!(kind, JoinKind::Semi | JoinKind::Anti);
let kind = if swapped && !marking {
kind.mirrored().ok_or_else(|| {
Error::internal(format!(
"a {} join was given a build side it has no mirror for",
kind.keyword()
))
})?
} else {
kind
};
let gather_id = self.gathered(reference);
let gathering = self.shape.pipeline(held);
let parent = held;
let above = std::mem::take(&mut self.above);
let held = self.node(held)?;
self.above = above;
let held_schema = held.schema.clone();
let sideways = Sideways::new();
let ordered = kind == JoinKind::Positional;
let (gather, gathered) = Keep::watching(memory, Some(Arc::clone(&sideways)), ordered);
let watched = self.watch(reference, gather_id, gathering, "Gather", None);
self.close(held, gathering, Arc::new(Watched::new(gather, watched)));
self.sideways = Some(Arc::clone(&sideways));
let from = self.armed.len();
let mut left = self.node(driving)?;
self.sideways = None;
self.above.clear();
let below: Vec<Arc<Sideways<'a>>> = self.armed[from..].to_vec();
self.armed.push(Arc::clone(&sideways));
let side = Gathered { schema: &held_schema, chunks: gathered, marker, swapped };
let zone = self.session.session_time_zone();
let (catalog, reducing) =
(self.catalog, self.session.rules().enabled(Rule::GraphReduction));
let arm = |keyed: Vec<(ExprRef, ColumnBinding)>| {
let Some((key, binding)) = keyed.into_iter().find_map(|(key, binding)| {
sideways::beneath(plan, driving, binding).map(|binding| (key, binding))
}) else {
return;
};
if reducing {
if let Some(exact) = exact(plan, catalog, parent, key, driving, binding) {
sideways.exactly(exact);
}
}
sideways.keying(Keyed::new(plan, key, held_schema.clone(), zone));
sideways.about(binding);
};
if marking {
let made =
Marking::new(plan, &left.schema, &side, kind, conditions, self.cancel, memory);
let Some((mark, out)) = made else {
return Err(Error::internal(format!(
"a {} join was given a build side no lookup can answer it from",
kind.keyword()
)));
};
let mark = mark.in_session(self.session);
arm(mark.sideways());
let schema = mark.schema().clone();
let counters = self.watch(reference, id, pipeline, "Mark", None);
let reading = Arc::clone(&counters);
let mark = mark.watched(Arc::clone(&counters));
left.after.push(gathering);
self.close(left, pipeline, Arc::new(Watched::new(mark, counters)));
let reader = Arc::new(Watched::new(out, reading));
return Ok(Segment::reading(reader, schema, pipeline));
}
if let Some(pad) =
Padding::new(plan, &left.schema, &side, kind, conditions, self.cancel, memory)
{
let pad = pad.in_session(self.session);
arm(pad.sideways());
let schema = pad.schema().clone();
let counters = self.watch(reference, id, pipeline, "Pad", None);
let pad = pad.watched(Arc::clone(&counters));
left.after.push(gathering);
return Ok(left.then(Arc::new(Watched::new(pad, counters)), schema));
}
if kind == JoinKind::Single && conditions.is_empty() && !swapped {
let broadcast = Broadcast::new(&left.schema, side.schema, side.chunks.clone());
let schema = broadcast.schema().clone();
let counters = self.watch(reference, id, pipeline, "Broadcast", None);
left.after.push(gathering);
return Ok(left.then(Arc::new(Watched::new(broadcast, counters)), schema));
}
if let Some(probe) =
Probe::new(plan, &left.schema, &side, kind, conditions, self.cancel, memory)
{
let probe = probe.in_session(self.session);
arm(probe.sideways());
let narrowing = probe
.driving_columns()
.into_iter()
.filter_map(|(at, binding)| {
let binding = sideways::beneath(plan, driving, binding)?;
let found = below.iter().find(|below| below.binding() == Some(binding))?;
found.wanted();
Some((at, Arc::clone(found)))
})
.collect();
let probe = probe.narrowed_by(narrowing);
let schema = probe.schema().clone();
let counters = self.watch(reference, id, pipeline, "Probe", None);
let probe = probe.watched(Arc::clone(&counters));
left.after.push(gathering);
return Ok(left.then(Arc::new(Watched::new(probe, counters)), schema));
}
let (join, out) =
Join::new(plan, &left.schema, side, kind, conditions, self.cancel, memory);
let join = join.in_session(self.session);
let schema = join.schema().clone();
let counters = self.watch(reference, id, pipeline, "Join", None);
let reading = Arc::clone(&counters);
let join = join.watched(Arc::clone(&counters));
left.after.push(gathering);
self.close(left, pipeline, Arc::new(Watched::new(join, counters)));
Ok(Segment::reading(Arc::new(Watched::new(out, reading)), schema, pipeline))
}
fn aggregate(
&mut self,
reference: NodeRef,
input: NodeRef,
index: u32,
groups: Slice,
aggregates: Slice,
bound: AggregateBound,
) -> Result<Segment<'a>> {
if bound.max_groups.is_none() && bound.having_count.is_none() {
if let Some(values) =
stored_summary(self.plan, self.catalog, input, groups, aggregates)?
{
let schema = summary_schema(self.plan, index, aggregates)?;
let source = Summary::new(&schema, &values)?;
let id = self.shape.operator(reference);
let pipeline = self.shape.pipeline(reference);
let counters =
self.watch(reference, id, pipeline, "Aggregate", Some("stored summary"));
return Ok(Segment::new(Arc::new(Watched::new(source, counters)), schema));
}
}
self.marking = marks_through(self.plan, input, groups, aggregates).then_some(input);
let below = self.node(input);
self.marking = None;
let below = below?;
let (aggregate, out) =
Aggregate::new(self.plan, &below.schema, index, groups, aggregates, self.memory)?;
let aggregate = aggregate.in_session(self.session);
let aggregate = match bound.max_groups {
Some(limit) => aggregate.limit_groups(limit),
None => aggregate,
};
let aggregate = match self.plan.presized(index).filter(|_| !self.plan.clustered(index)) {
Some(groups) => aggregate.presize(match bound.max_groups {
Some(limit) => groups.min(u64::try_from(limit).unwrap_or(u64::MAX)),
None => groups,
}),
None => aggregate,
};
let aggregate = if self.session.rules().enabled(Rule::MemoryReservation) {
aggregate.reserved()
} else {
aggregate
};
let aggregate = match self.plan.dense(index) {
Some((low, values)) => aggregate.over_range(low, values),
None => aggregate,
};
let aggregate = if self.plan.clustered(index) { aggregate.clustered() } else { aggregate };
let aggregate = match bound.top_counts {
Some((bound, call)) => aggregate.top_counts(bound, call),
None => aggregate,
};
let aggregate = match bound.having_count {
Some((call, minimum)) => aggregate.having_count(call, minimum),
None => aggregate,
};
let schema = aggregate.schema().clone();
let id = self.shape.operator(reference);
let pipeline = self.shape.pipeline(reference);
if bound.max_groups.is_none() {
if let Some(records) = native_host_groups(
self.plan,
self.catalog,
input,
groups,
aggregates,
bound.having_count,
)? {
let source = Frequencies::records(schema.clone(), records)?;
let counters =
self.watch(reference, id, pipeline, "Aggregate", Some("native host groups"));
return Ok(Segment::new(Arc::new(Watched::new(source, counters)), schema));
}
}
if bound.max_groups.is_none() && bound.having_count.is_none() {
let top = bound.top_counts.map(|(bound, _)| bound);
if let Some(top) = top {
if let Some(frequencies) = native_pair_frequencies(
self.plan,
self.catalog,
input,
groups,
aggregates,
top,
)? {
let source = Frequencies::grouped(schema.clone(), frequencies.entries)?;
let counters = self.watch(
reference,
id,
pipeline,
"Aggregate",
Some("native pair frequencies"),
);
return Ok(Segment::new(Arc::new(Watched::new(source, counters)), schema));
}
}
}
let counters = self.watch(reference, id, pipeline, "Aggregate", None);
let reading = Arc::clone(&counters);
self.close(below, pipeline, Arc::new(Watched::new(aggregate, counters)));
Ok(Segment::reading(Arc::new(Watched::new(out, reading)), schema, pipeline))
}
fn linked(
&self,
reference: NodeRef,
child: NodeRef,
parent: NodeRef,
conditions: Slice,
) -> Result<Linked> {
let (plan, catalog) = (self.plan, self.catalog);
let refuse =
|why: &str| Error::internal(format!("a link join over {why}, which cannot be read"));
let Some((parent_table, parent_index, parent_columns)) =
whole_table(plan, catalog, parent)?
else {
return Err(refuse("a parent that is not a stored table"));
};
let keys = plan
.expr_list(conditions)
.iter()
.map(|&condition| equated(plan, condition))
.collect::<Option<Vec<_>>>()
.ok_or_else(|| refuse("something other than equalities between two columns"))?;
let mut oriented = Vec::with_capacity(keys.len());
for [first, second] in keys {
oriented.push(match (first.table == parent_index, second.table == parent_index) {
(false, true) => (first, second),
(true, false) => (second, first),
_ => return Err(refuse("an equality that does not read both of its inputs")),
});
}
let Some(&(child_key, _)) = oriented.first() else {
return Err(refuse("no equality at all"));
};
if oriented.iter().any(|(key, _)| key.table != child_key.table) {
return Err(refuse("a key whose child columns come from two tables"));
}
let Some((child_table, child_columns)) = scanned(plan, catalog, child, child_key.table)?
else {
return Err(refuse("a child that is not a stored table"));
};
let (Some(child_rows), Some(parent_rows)) =
(child_table.rows().stored(), parent_table.rows().stored())
else {
return Err(refuse("a table that is not one committed file"));
};
let mut child_at = Vec::with_capacity(oriented.len());
let mut parent_at = Vec::with_capacity(oriented.len());
for &(child_key, parent_key) in &oriented {
child_at.push(
stored_column(plan, child_table, child_key.table, child_columns, child_key)
.ok_or_else(|| refuse("a child key that is not a stored column"))?,
);
parent_at.push(
stored_column(plan, parent_table, parent_index, parent_columns, parent_key)
.ok_or_else(|| refuse("a parent key that is not a stored column"))?,
);
}
let edge = |child_at: &[usize], parent_at: &[usize]| {
Some(rudb_native::graph::Edge {
child: child_table.name().table.clone(),
child_column: rudb_native::graph::key_of(child_at)?,
parent: parent_table.name().table.clone(),
parent_column: rudb_native::graph::key_of(parent_at)?,
})
};
let mut edges = vec![edge(&child_at, &parent_at)];
child_at.reverse();
parent_at.reverse();
if child_at.len() == 2 {
edges.push(edge(&child_at, &parent_at));
}
let link = edges
.into_iter()
.flatten()
.find_map(|edge| rudb_native::graph::stored_link(child_rows, parent_rows, &edge))
.ok_or_else(|| refuse("a relationship the child's file has no link for"))?;
let reads = !matches!(
*plan.node(reference),
Node::LinkJoin { kind: JoinKind::Semi | JoinKind::Anti, .. }
);
let fields = if reads { plan.field_list(parent_columns).to_vec() } else { Vec::new() };
let mut projected = Vec::with_capacity(fields.len());
for field in &fields {
let at = parent_table
.column_index(&field.name)
.ok_or_else(|| refuse("a parent column the stored table does not have"))?;
projected.push((at, field.ty.clone()));
}
let budget = self.memory.limit().map_or(usize::MAX, |limit| {
usize::try_from(limit.saturating_sub(self.memory.used())).unwrap_or(usize::MAX)
});
Ok(Linked {
link: Arc::new(link),
parent: Arc::new(Parent::new(parent_table.rows().clone(), budget)),
projected,
parent_schema: Schema::numbered(fields, parent_index),
})
}
fn node(&mut self, reference: NodeRef) -> Result<Segment<'a>> {
let plan = self.plan;
let memory = self.memory;
let id = self.shape.operator(reference);
let pipeline = self.shape.pipeline(reference);
if !matches!(
*plan.node(reference),
Node::Get { .. }
| Node::Filter { .. }
| Node::Project { .. }
| Node::TableFunction { .. }
) {
let inherited = self.sideways.take();
if sideways::through(plan.node(reference)).is_some() {
self.above.extend(inherited);
} else {
self.above.clear();
}
}
if !matches!(
*plan.node(reference),
Node::Get { .. } | Node::Filter { .. } | Node::Project { .. }
) {
self.cutoff = None;
}
let segment = match *plan.node(reference) {
Node::Get { catalog: database, schema, table, index, columns, .. } => {
let name = QualifiedName::new(
plan.string(database),
plan.string(schema),
plan.string(table),
);
let filters = Filters {
pruning: std::mem::take(&mut self.pruning),
pushed: self.pushing.take(),
sideways: self.sideways.take(),
also: std::mem::take(&mut self.above),
cutoff: self.cutoff.take(),
};
let moved = filters.pushed.as_ref().map(|pushed| pushed.node);
let counters = self.watch_doing(
reference,
moved,
id,
pipeline,
"Scan",
Some(plan.string(table)),
);
let scan = Scan::new(
plan,
self.catalog.table(&name)?,
index,
columns,
filters,
self.seams,
self.session,
)?
.watched(counters.clone());
let schema = scan.schema().clone();
Segment::new(Arc::new(Watched::new(scan, counters)), schema)
}
Node::Dummy => {
let dummy = Dummy::new();
let schema = dummy.schema().clone();
let counters = self.watch(reference, id, pipeline, "Dummy", None);
Segment::new(Arc::new(Watched::new(dummy, counters)), schema)
}
Node::Values { index, columns, rows } => {
let values = Values::new(plan, index, columns, rows, self.session)?;
let schema = values.schema().clone();
let counters = self.watch(reference, id, pipeline, "Values", None);
Segment::new(Arc::new(Watched::new(values, counters)), schema)
}
Node::TableFunction { .. } => self.table_function(reference)?,
Node::LateralFunction { input, index, function, args, columns, .. } => {
let below = self.node(input)?;
let name = plan.string(function);
let lateral = LateralSeries::new(
plan,
&below.schema,
index,
name,
args,
columns,
self.cancel,
)?
.in_session(self.session);
let schema = lateral.schema().clone();
let counters = self.watch(reference, id, pipeline, "Series", Some(name));
below.then(Arc::new(Watched::new(lateral, counters)), schema)
}
Node::Fetch { input, index, args, columns, row } => {
let below = self.node(input)?;
let counters = self.watch(reference, id, pipeline, "Fetch", None);
let fetch = Fetch::new(plan, &below.schema, index, args, columns, row)?
.in_session(self.session)
.watched(counters.clone());
let schema = fetch.schema().clone();
below.then(Arc::new(Watched::new(fetch, counters)), schema)
}
Node::TableFetch { input, index, catalog, schema, table, columns, row } => {
let below = self.node(input)?;
let name = QualifiedName::new(
plan.string(catalog),
plan.string(schema),
plan.string(table),
);
let counters = self.watch(reference, id, pipeline, "TableFetch", None);
let fetch = TableFetch::new(
plan,
&below.schema,
index,
self.catalog.table(&name)?,
columns,
row,
)?
.in_session(self.session);
let schema = fetch.schema().clone();
below.then(Arc::new(Watched::new(fetch, counters)), schema)
}
Node::Filter { input, predicate } => {
let marks = self.marking == Some(reference);
self.pruning = rudb_opt::bounds::of(plan, input, predicate);
self.pushing = rudb_opt::bounds::into_scan(plan, reference).map(|moved| Pushdown {
node: reference,
predicate,
tests: moved.tests,
whole: moved.whole,
conjuncts: moved.conjuncts,
marks,
});
let offered = self.pushing.is_some();
let below = match count_having_aggregate(plan, input, predicate) {
Some((aggregate, call, minimum)) => {
let Node::Aggregate { input: under, index, groups, aggregates } =
*plan.node(aggregate)
else {
unreachable!("count_having_aggregate returned another node")
};
self.aggregate(
aggregate,
under,
index,
groups,
aggregates,
AggregateBound {
max_groups: None,
top_counts: None,
having_count: Some((call, minimum)),
},
)?
}
None => self.node(input)?,
};
self.pruning = Vec::new();
let taken = offered && self.pushing.take().is_none();
self.pushing = None;
if taken {
return Ok(below);
}
let schema = below.schema.clone();
let filter = Filter::new(plan, reference, predicate, &schema, self.seams)?
.in_session(self.session)
.marking(marks);
let counters = self.watch(reference, id, pipeline, "Filter", None);
below.then(Arc::new(Watched::new(filter, counters)), schema)
}
Node::Project { input, index, exprs, names } => {
let below = self.node(input)?;
let project = Project::new(plan, &below.schema, index, exprs, names)?
.in_session(self.session);
let schema = project.schema().clone();
let counters = self.watch(reference, id, pipeline, "Project", None);
below.then(Arc::new(Watched::new(project, counters)), schema)
}
Node::Aggregate { input, index, groups, aggregates } => {
let top_counts = self.top_counts.iter().find_map(|&(aggregate, bound, call)| {
(aggregate == reference).then_some((bound, call))
});
self.aggregate(
reference,
input,
index,
groups,
aggregates,
AggregateBound { max_groups: None, top_counts, having_count: None },
)?
}
Node::Sort { input, keys } => {
let below = self.node(input)?;
let schema = below.schema.clone();
let (sort, out) = Sort::new(plan, &schema, keys, memory)?;
let sort = sort.in_session(self.session);
let counters = self.watch(reference, id, pipeline, "Sort", None);
let reading = Arc::clone(&counters);
self.close(below, pipeline, Arc::new(Watched::new(sort, counters)));
Segment::reading(Arc::new(Watched::new(out, reading)), schema, pipeline)
}
Node::Limit { input, count, offset } => {
let max_groups = count
.rows()
.zip(offset.rows())
.and_then(|(count, offset)| count.checked_add(offset))
.and_then(|count| usize::try_from(count).ok());
let below = match (plan.node(input).clone(), max_groups) {
(
Node::Aggregate { input: under, index, groups, aggregates },
Some(max_groups),
) => self.aggregate(
input,
under,
index,
groups,
aggregates,
AggregateBound {
max_groups: Some(max_groups),
top_counts: None,
having_count: None,
},
)?,
_ => self.node(input)?,
};
let schema = below.schema.clone();
let limit = Limit::new(edge(plan, count, &schema)?, edge(plan, offset, &schema)?)
.in_session(self.session);
let counters = self.watch(reference, id, pipeline, "Limit", None);
below.then(Arc::new(Watched::new(limit, counters)), schema)
}
Node::LimitPercent { input, percent, offset } => {
let below = self.node(input)?;
let schema = below.schema.clone();
let (limit, out) = LimitPercent::new(
portion(plan, percent, &schema)?,
edge(plan, offset, &schema)?,
memory,
self.session,
);
let counters = self.watch(reference, id, pipeline, "LimitPercent", None);
let reading = Arc::clone(&counters);
self.close(below, pipeline, Arc::new(Watched::new(limit, counters)));
Segment::reading(Arc::new(Watched::new(out, reading)), schema, pipeline)
}
Node::TopN { input, keys, count, offset } => {
if let Some((aggregate, call)) = count_top_aggregate(plan, input, keys) {
let bound = count.saturating_add(offset);
if let Ok(bound) = usize::try_from(bound) {
self.top_counts.push((aggregate, bound, call));
}
}
let cutoff = Cutoff::new();
self.cutoff = Some(Arc::clone(&cutoff));
let below = self.node(input)?;
self.cutoff = None;
if let Some((binding, op)) = cutoff::ordering(plan, keys) {
if let Some(binding) = sideways::beneath(plan, input, binding) {
cutoff.about(binding, op);
}
}
let schema = below.schema.clone();
let (top, out) = TopN::new(plan, &schema, keys, count, offset, memory)?;
let top = top.telling(cutoff).in_session(self.session);
let counters = self.watch(reference, id, pipeline, "TopN", None);
let reading = Arc::clone(&counters);
self.close(below, pipeline, Arc::new(Watched::new(top, counters)));
Segment::reading(Arc::new(Watched::new(out, reading)), schema, pipeline)
}
Node::Distinct { input, on } => {
let below = self.node(input)?;
let schema = below.schema.clone();
let (distinct, out) = Distinct::new(plan, &schema, on, memory)?;
let distinct = distinct.in_session(self.session);
let counters = self.watch(reference, id, pipeline, "Distinct", None);
let reading = Arc::clone(&counters);
self.close(below, pipeline, Arc::new(Watched::new(distinct, counters)));
Segment::reading(Arc::new(Watched::new(out, reading)), schema, pipeline)
}
Node::LinkJoin { child, parent, kind, conditions, rid } => {
let found = self.linked(reference, child, parent, conditions)?;
let below = self.node(child)?;
let operator = LinkJoin::new(
plan,
kind,
found.link,
found.parent,
found.projected,
rid,
&below.schema,
&found.parent_schema,
self.seams,
memory,
self.cancel.clone(),
)?
.in_session(self.session);
let schema = operator.schema().clone();
let counters = self.watch(reference, id, pipeline, "LinkJoin", None);
below.then(Arc::new(Watched::new(operator, counters)), schema)
}
Node::Join { .. } => self.join(reference)?,
Node::CrossProduct { left, right } => {
let keep_id = self.gathered(reference);
let aside = self.shape.pipeline(right);
let right = self.node(right)?;
let right_schema = right.schema.clone();
let (keep, kept) = Keep::new(memory);
let held = self.watch(reference, keep_id, aside, "Keep", None);
self.close(right, aside, Arc::new(Watched::new(keep, held)));
let mut left = self.node(left)?;
let cross = CrossProduct::new(&left.schema, &right_schema, kept);
let schema = cross.schema().clone();
let counters = self.watch(reference, id, pipeline, "CrossProduct", None);
left.after.push(aside);
left.then(Arc::new(Watched::new(cross, counters)), schema)
}
Node::SetOp { left, right, kind, all, index } => {
let gather_id = self.gathered(reference);
let counting = self.shape.pipeline(right);
let right = self.node(right)?;
let (gather, gathered) = Gather::new(memory);
let kept = self.watch(reference, gather_id, counting, "Gather", None);
self.close(right, counting, Arc::new(Watched::new(gather, kept)));
let mut left = self.node(left)?;
let (setop, out) = SetOp::new(&left.schema, gathered, kind, all, index, memory);
let schema = setop.schema().clone();
let counters = self.watch(reference, id, pipeline, "SetOp", None);
let reading = Arc::clone(&counters);
left.after.push(counting);
self.close(left, pipeline, Arc::new(Watched::new(setop, counters)));
Segment::reading(Arc::new(Watched::new(out, reading)), schema, pipeline)
}
Node::DependentJoin { .. } => {
return Err(Error::internal(
"a dependent join reached execution before subquery unnesting",
));
}
Node::Window { input, index, partition, order, frame, expressions } => {
let below = self.node(input)?;
let written = Written { index, partition, order, frame, expressions };
let (window, out) = Window::new(plan, &below.schema, &written, memory)?;
let window = window.in_session(self.session);
let schema = window.schema().clone();
let counters = self.watch(reference, id, pipeline, "Window", None);
let reading = Arc::clone(&counters);
self.close(below, pipeline, Arc::new(Watched::new(window, counters)));
Segment::reading(Arc::new(Watched::new(out, reading)), schema, pipeline)
}
Node::MaterializedCte { definition, body, cte, .. } => {
let held = self.node(definition)?;
let (keep, chunks) = Keep::new(memory);
let counters = self.watch(reference, id, pipeline, "MaterializedCTE", None);
self.close(held, pipeline, Arc::new(Watched::new(keep, counters)));
self.held.push(Held { cte, chunks, filling: pipeline });
let segment = self.node(body);
self.held.pop();
segment?
}
Node::CteScan { index, cte, columns, .. } => {
let Some(source) = self.held.iter().rev().find(|held| held.cte == cte) else {
return Err(Error::internal(
"a read of a materialisation that is not being filled",
));
};
let filling = source.filling;
let source = source.chunks.reader();
let schema = Schema::numbered(plan.field_list(columns).to_vec(), index);
let counters = self.watch(reference, id, pipeline, "CteScan", None);
Segment::reading(Arc::new(Watched::new(source, counters)), schema, filling)
}
};
Ok(segment)
}
}
fn marks_through(plan: &Plan, input: NodeRef, groups: Slice, aggregates: Slice) -> bool {
if !matches!(plan.node(input), Node::Filter { .. }) {
return false;
}
let keys = plan.expr_list(groups);
let calls = plan.expr_list(aggregates);
let plain = calls
.iter()
.all(|&call| matches!(plan.expr(call), Expr::Aggregate { distinct: false, .. }));
!keys.is_empty()
&& plain
&& !keys.iter().chain(calls).any(|&expr| rudb_opt::volatile(plan, expr))
}
#[cfg(test)]
mod tests {
use std::fs;
use std::path::PathBuf;
use std::time::{SystemTime, UNIX_EPOCH};
use rudb_catalog::Catalog;
use rudb_common::{Field, LogicalType, Value};
use rudb_plan::{CompareOp, Expr, Node, Plan};
use rudb_vector::{Chunk, VECTOR_SIZE, Vector};
use super::{count_having_aggregate, count_top_aggregate, native_pair_frequencies};
fn native_path(label: &str) -> PathBuf {
let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
std::env::temp_dir().join(format!("rudb-exec-{label}-{}-{stamp}.rdb", std::process::id()))
}
fn native_catalog(label: &str, rows: &[(i64, String)]) -> (PathBuf, Catalog) {
let path = native_path(label);
let fields = vec![
Field::required("id", LogicalType::BigInt),
Field::required("phrase", LogicalType::Varchar),
];
let mut writer = rudb_native::Writer::create(&path, "items", fields).expect("new file");
for rows in rows.chunks(VECTOR_SIZE) {
let ids = rows.iter().map(|(id, _)| Value::BigInt(*id)).collect::<Vec<_>>();
let phrases =
rows.iter().map(|(_, phrase)| Value::Varchar(phrase.clone())).collect::<Vec<_>>();
let chunk = Chunk::new(vec![
Vector::from_values(LogicalType::BigInt, &ids).expect("big integers"),
Vector::from_values(LogicalType::Varchar, &phrases).expect("strings"),
])
.expect("matching columns");
writer.append(&chunk).expect("rows");
}
writer.finish().expect("commit");
let native = rudb_native::Catalog::open(&path).expect("reopen");
let mut catalog = Catalog::new();
catalog.create_native_table(native.table("items").expect("stored table")).expect("attach");
(path, catalog)
}
fn pair_plan() -> Plan {
Plan::parse(
"Aggregate #1 groups=[#0.0::BIGINT, #0.1::VARCHAR] \
aggregates=[count_star()::BIGINT]\n \
Get memory.main.items AS items #0 [id::BIGINT, phrase::VARCHAR]",
)
.expect("a two-key grouped count")
}
fn pair_frequencies(
plan: &Plan,
catalog: &Catalog,
top: usize,
) -> Option<super::NativePairFrequencies> {
let Node::Aggregate { input, groups, aggregates, .. } = *plan.node(plan.root()) else {
panic!("the root is an aggregate")
};
native_pair_frequencies(plan, catalog, input, groups, aggregates, top)
.expect("metadata reads")
}
#[test]
fn sparse_occurrences_compute_two_key_top_counts_at_query_time() {
let mut rows = Vec::new();
rows.extend(std::iter::repeat_n((1, "a".to_string()), VECTOR_SIZE + 5));
rows.extend(std::iter::repeat_n((1, "b".to_string()), 4));
rows.extend(std::iter::repeat_n((2, "x".to_string()), 3));
rows.push((3, "y".to_string()));
let (path, catalog) = native_catalog("pair-frequencies", &rows);
let answer = pair_frequencies(&pair_plan(), &catalog, 2).expect("query-time result");
assert_eq!(answer.entries.len(), 2);
assert!(answer.entries.contains(&(
vec![Value::BigInt(1), Value::Varchar("a".to_string())],
u64::try_from(VECTOR_SIZE + 5).expect("a small vector width"),
)));
assert!(
answer.entries.contains(&(vec![Value::BigInt(1), Value::Varchar("b".to_string())], 4))
);
fs::remove_file(path).expect("clean up");
}
#[test]
fn sparse_occurrences_refuse_a_pair_tied_with_the_omitted_tail() {
let rows = (0..513_i64).map(|id| (id, format!("phrase {id}"))).collect::<Vec<_>>();
let (path, catalog) = native_catalog("pair-fallback", &rows);
assert!(
pair_frequencies(&pair_plan(), &catalog, 10).is_none(),
"a count of one cannot beat an omitted first-key count of one"
);
fs::remove_file(path).expect("clean up");
}
fn plan(direction: &str) -> Plan {
Plan::parse(&format!(
"TopN 10 offset 0 [#2.2::BIGINT {direction} NULLS LAST]\n \
Project #2 [#1.0::BIGINT AS WatchID, #1.1::INTEGER AS ClientIP, #1.2::BIGINT AS c]\n \
Aggregate #1 groups=[#0.0::BIGINT, #0.1::INTEGER] \
aggregates=[count_star()::BIGINT]\n \
Values #0 [WatchID::BIGINT, ClientIP::INTEGER] rows=[]"
))
.expect("a grouped count plan")
}
#[test]
fn count_descending_topn_marks_its_aggregate() {
let plan = plan("DESC");
let Node::TopN { input, keys, .. } = *plan.node(plan.root()) else {
panic!("the root is a TopN")
};
let (aggregate, call) = count_top_aggregate(&plan, input, keys).expect("the grouped count");
assert!(matches!(plan.node(aggregate), Node::Aggregate { .. }));
assert_eq!(call, 0, "the count is the only call");
}
#[test]
fn count_descending_topn_crosses_several_passthrough_projects() {
let plan = Plan::parse(
"TopN 10 offset 0 [#3.1::BIGINT DESC NULLS LAST]\n \
Project #3 [#2.0::INTEGER AS ClientIP, #2.1::BIGINT AS c]\n \
Project #2 [#1.0::INTEGER AS column0, #1.1::BIGINT AS column1]\n \
Aggregate #1 groups=[#0.0::INTEGER] aggregates=[count_star()::BIGINT]\n \
Values #0 [ClientIP::INTEGER] rows=[]",
)
.expect("a grouped count under two projects");
let Node::TopN { input, keys, .. } = *plan.node(plan.root()) else {
panic!("the root is a TopN")
};
let (aggregate, call) = count_top_aggregate(&plan, input, keys).expect("the grouped count");
assert!(matches!(plan.node(aggregate), Node::Aggregate { .. }));
assert_eq!(call, 0, "the count is the only call");
}
#[test]
fn count_ascending_cannot_discard_large_counts() {
let plan = plan("ASC");
let Node::TopN { input, keys, .. } = *plan.node(plan.root()) else {
panic!("the root is a TopN")
};
assert!(count_top_aggregate(&plan, input, keys).is_none());
}
#[test]
fn distinct_count_descending_topn_marks_its_aggregate() {
let plan = Plan::parse(
"TopN 10 offset 0 [#1.1::BIGINT DESC NULLS LAST]\n \
Aggregate #1 groups=[#0.0::VARCHAR] \
aggregates=[count(DISTINCT #0.1::BIGINT)::BIGINT]\n \
Values #0 [SearchPhrase::VARCHAR, UserID::BIGINT] rows=[]",
)
.expect("a grouped distinct count plan");
let Node::TopN { input, keys, .. } = *plan.node(plan.root()) else {
panic!("the root is a TopN")
};
let (aggregate, call) =
count_top_aggregate(&plan, input, keys).expect("the distinct count");
assert!(matches!(plan.node(aggregate), Node::Aggregate { .. }));
assert_eq!(call, 0, "the distinct count is the only call");
}
#[test]
fn count_descending_topn_finds_a_later_aggregate_call() {
let plan = Plan::parse(
"TopN 10 offset 0 [#1.2::BIGINT DESC NULLS LAST]\n \
Aggregate #1 groups=[#0.0::INTEGER] \
aggregates=[sum(#0.1::SMALLINT)::HUGEINT, count_star()::BIGINT, avg(#0.2::SMALLINT)::DOUBLE, count(DISTINCT #0.3::BIGINT)::BIGINT]\n \
Values #0 [RegionID::INTEGER, AdvEngineID::SMALLINT, ResolutionWidth::SMALLINT, UserID::BIGINT] rows=[]",
)
.expect("a mixed aggregate plan");
let Node::TopN { input, keys, .. } = *plan.node(plan.root()) else {
panic!("the root is a TopN")
};
let (aggregate, call) = count_top_aggregate(&plan, input, keys).expect("the grouped count");
assert!(matches!(plan.node(aggregate), Node::Aggregate { .. }));
assert_eq!(call, 1, "the count is the second of the four calls");
}
#[test]
fn a_count_having_lower_bound_marks_the_count_call() {
let plan = Plan::parse(
"Filter (#1.2::BIGINT > 100::BIGINT)::BOOLEAN\n \
Aggregate #1 groups=[#0.0::BIGINT] \
aggregates=[avg(#0.1::BIGINT)::DOUBLE, count_star()::BIGINT]\n \
Values #0 [key::BIGINT, value::BIGINT] rows=[]",
)
.expect("an aggregate with a HAVING filter");
let Node::Filter { input, predicate } = *plan.node(plan.root()) else {
panic!("the root is a Filter")
};
let (aggregate, call, minimum) =
count_having_aggregate(&plan, input, predicate).expect("the count bound");
assert_eq!(aggregate, input);
assert_eq!((call, minimum), (1, 101));
}
#[test]
fn an_upper_count_having_bound_cannot_drop_aggregate_output() {
let mut plan = Plan::parse(
"Filter (#1.1::BIGINT > 100::BIGINT)::BOOLEAN\n \
Aggregate #1 groups=[#0.0::BIGINT] aggregates=[count_star()::BIGINT]\n \
Values #0 [key::BIGINT] rows=[]",
)
.expect("an aggregate with a HAVING filter");
let Node::Filter { input, predicate } = *plan.node(plan.root()) else {
panic!("the root is a Filter")
};
let Expr::Compare { left, right, .. } = *plan.expr(predicate) else {
panic!("the predicate is a comparison")
};
let less =
plan.add_expr(Expr::Compare { op: CompareOp::Less, left, right }, LogicalType::Boolean);
assert!(count_having_aggregate(&plan, input, less).is_none());
}
}