use std::sync::Arc;
use rudb_catalog::{Catalog, QualifiedName};
use rudb_common::{Cancel, Memory, Result};
use rudb_functions::TableFunction;
use rudb_metrics::{Counters, Report};
use rudb_pipeline::{Source, Watched};
use rudb_plan::{Node, NodeRef, Plan, Shape};
use crate::adapt::{Broken, Fed, Paired, Pulled, Streamed};
use crate::cancel::Guarded;
use crate::gather::{Gather, Keep};
use crate::group::{Aggregate, Distinct};
use crate::join::{CrossProduct, Gathered, Join};
use crate::operator::Operator;
use crate::schema::Schema;
use crate::setop::SetOp;
use crate::sort::Sort;
use crate::source::{Dummy, FileScan, Scan, Series, Values};
use crate::strategies::Strategies;
use crate::stream::{Filter, Limit, Project};
use crate::topn::TopN;
pub fn build<'a>(plan: &'a Plan, catalog: &'a Catalog) -> Result<Box<dyn Operator + 'a>> {
build_with(plan, catalog, &Cancel::new(), &Memory::unlimited())
}
pub fn build_with<'a>(
plan: &'a Plan,
catalog: &'a Catalog,
cancel: &Cancel,
memory: &Memory,
) -> Result<Box<dyn Operator + 'a>> {
build_measured(plan, catalog, cancel, memory, &Report::new())
}
pub fn build_measured<'a>(
plan: &'a Plan,
catalog: &'a Catalog,
cancel: &Cancel,
memory: &Memory,
report: &Report,
) -> Result<Box<dyn Operator + 'a>> {
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 building = Building { plan, catalog, cancel, memory, report, shape };
building.node(plan.root())
}
struct Building<'a, 'b> {
plan: &'a Plan,
catalog: &'a Catalog,
cancel: &'b Cancel,
memory: &'b Memory,
report: &'b Report,
shape: Shape,
}
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 watch(&self, id: u32, pipeline: u32, kind: &str, detail: Option<&str>) -> Arc<Counters> {
let counters = Counters::new(id, pipeline, kind).reference();
let counters = match detail {
Some(detail) => counters.detailed(detail),
None => counters,
};
self.report.watch(counters)
}
fn node(&self, reference: NodeRef) -> Result<Box<dyn Operator + 'a>> {
let plan = self.plan;
let memory = self.memory;
let id = self.shape.operator(reference);
let pipeline = self.shape.pipeline(reference);
let inner: Box<dyn Operator + 'a> = 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 scan = Scan::new(plan, self.catalog.table(&name)?, index, columns)?;
let schema = scan.schema().clone();
let counters = self.watch(id, pipeline, "Scan", Some(plan.string(table)));
pulled(Watched::new(scan, counters), schema)
}
Node::Dummy => {
let dummy = Dummy::new();
let schema = dummy.schema().clone();
pulled(Watched::new(dummy, self.watch(id, pipeline, "Dummy", None)), schema)
}
Node::Values { index, columns, rows } => {
let values = Values::new(plan, index, columns, rows)?;
let schema = values.schema().clone();
pulled(Watched::new(values, self.watch(id, pipeline, "Values", None)), schema)
}
Node::TableFunction { index, function, args, options, settings, columns } => {
let name = plan.string(function);
match TableFunction::lookup(name) {
Some(function @ (TableFunction::ReadParquet | TableFunction::ReadCsv)) => {
let scan =
FileScan::new(plan, index, function, args, options, settings, columns)?;
let schema = scan.schema().clone();
let counters = self.watch(id, pipeline, "FileScan", Some(name));
pulled(Watched::new(scan, counters), schema)
}
Some(TableFunction::RudbStrategies) => {
let table = Strategies::new(plan, index, columns)?;
let schema = table.schema().clone();
let counters = self.watch(id, pipeline, "Strategies", None);
pulled(Watched::new(table, counters), schema)
}
_ => {
let series = Series::new(plan, index, name, args)?;
let schema = series.schema().clone();
let counters = self.watch(id, pipeline, "Series", Some(name));
pulled(Watched::new(series, counters), schema)
}
}
}
Node::Filter { input, predicate } => {
let input = self.node(input)?;
let schema = input.schema().clone();
let filter = Filter::new(plan, predicate, &schema)?;
let counters = self.watch(id, pipeline, "Filter", None);
Box::new(Streamed::new(input, Watched::new(filter, counters), schema))
}
Node::Project { input, index, exprs, names } => {
let input = self.node(input)?;
let project = Project::new(plan, input.schema(), index, exprs, names)?;
let schema = project.schema().clone();
let counters = self.watch(id, pipeline, "Project", None);
Box::new(Streamed::new(input, Watched::new(project, counters), schema))
}
Node::Aggregate { input, index, groups, aggregates } => {
let input = self.node(input)?;
let (aggregate, out) =
Aggregate::new(plan, input.schema(), index, groups, aggregates, memory)?;
let schema = aggregate.schema().clone();
let counters = self.watch(id, pipeline, "Aggregate", None);
let made = Arc::clone(&counters);
Box::new(Broken::new(input, Watched::new(aggregate, counters), out, made, schema))
}
Node::Sort { input, keys } => {
let input = self.node(input)?;
let schema = input.schema().clone();
let (sort, out) = Sort::new(plan, &schema, keys, memory)?;
let counters = self.watch(id, pipeline, "Sort", None);
let made = Arc::clone(&counters);
Box::new(Broken::new(input, Watched::new(sort, counters), out, made, schema))
}
Node::Limit { input, count, offset } => {
let input = self.node(input)?;
let schema = input.schema().clone();
let limit = Limit::new(count, offset);
let counters = self.watch(id, pipeline, "Limit", None);
Box::new(Streamed::new(input, Watched::new(limit, counters), schema))
}
Node::TopN { input, keys, count, offset } => {
let input = self.node(input)?;
let schema = input.schema().clone();
let (top, out) = TopN::new(plan, &schema, keys, count, offset, memory)?;
let counters = self.watch(id, pipeline, "TopN", None);
let made = Arc::clone(&counters);
Box::new(Broken::new(input, Watched::new(top, counters), out, made, schema))
}
Node::Distinct { input, on } => {
let input = self.node(input)?;
let schema = input.schema().clone();
let (distinct, out) = Distinct::new(plan, &schema, on, memory)?;
let counters = self.watch(id, pipeline, "Distinct", None);
let made = Arc::clone(&counters);
Box::new(Broken::new(input, Watched::new(distinct, counters), out, made, schema))
}
Node::Join { left, right, kind, conditions } => {
let gather_id = self.gathered(reference);
let gathering = self.shape.pipeline(right);
let right = self.node(right)?;
let left = self.node(left)?;
let (gather, gathered) = Gather::new(memory);
let side = Gathered { schema: right.schema(), rows: gathered };
let (join, out) =
Join::new(plan, left.schema(), side, kind, conditions, self.cancel, memory);
let schema = join.schema().clone();
let kept = self.watch(gather_id, gathering, "Gather", None);
let counters = self.watch(id, pipeline, "Join", None);
let made = Arc::clone(&counters);
Box::new(Paired::new(
right,
Watched::new(gather, kept),
left,
Watched::new(join, counters),
out,
made,
schema,
))
}
Node::CrossProduct { left, right } => {
let keep_id = self.gathered(reference);
let aside = self.shape.pipeline(right);
let right = self.node(right)?;
let left = self.node(left)?;
let (keep, kept) = Keep::new(memory);
let cross = CrossProduct::new(left.schema(), right.schema(), kept);
let schema = cross.schema().clone();
let held = self.watch(keep_id, aside, "Keep", None);
let counters = self.watch(id, pipeline, "CrossProduct", None);
Box::new(Fed::new(
right,
Watched::new(keep, held),
Streamed::new(left, 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 left = self.node(left)?;
let (gather, gathered) = Gather::new(memory);
let (setop, out) = SetOp::new(left.schema(), gathered, kind, all, index, memory);
let schema = setop.schema().clone();
let kept = self.watch(gather_id, counting, "Gather", None);
let counters = self.watch(id, pipeline, "SetOp", None);
let made = Arc::clone(&counters);
Box::new(Paired::new(
right,
Watched::new(gather, kept),
left,
Watched::new(setop, counters),
out,
made,
schema,
))
}
};
Ok(Box::new(Guarded::new(inner, self.cancel.clone())))
}
}
fn pulled<'a, S: Source + 'a>(source: S, schema: Schema) -> Box<dyn Operator + 'a> {
Box::new(Pulled::new(source, schema))
}