use crate::node::Node;
use crate::plan::Plan;
use crate::{NodeRef, OperatorRef, PipelineRef};
#[derive(Debug, Clone)]
pub struct Shape {
of: Vec<Option<Placed>>,
waits: Vec<Vec<PipelineRef>>,
operators: OperatorRef,
}
#[derive(Debug, Clone, Copy)]
struct Placed {
operator: OperatorRef,
gathered: Option<OperatorRef>,
pipeline: PipelineRef,
}
impl Shape {
#[must_use]
pub fn of(plan: &Plan) -> Self {
let mut shape =
Self { of: vec![None; plan.node_count()], waits: vec![Vec::new()], operators: 0 };
shape.walk(plan, plan.root(), ROOT);
shape
}
#[must_use]
pub fn pipelines(&self) -> usize {
self.waits.len()
}
#[must_use]
pub fn operators(&self) -> OperatorRef {
self.operators
}
#[must_use]
pub fn operator(&self, node: NodeRef) -> OperatorRef {
self.placed(node).operator
}
#[must_use]
pub fn operator_of(&self, node: NodeRef) -> Option<OperatorRef> {
self.of.get(node as usize).copied().flatten().map(|placed| placed.operator)
}
#[must_use]
pub fn gathered(&self, node: NodeRef) -> Option<OperatorRef> {
self.placed(node).gathered
}
#[must_use]
pub fn pipeline(&self, node: NodeRef) -> PipelineRef {
self.placed(node).pipeline
}
#[must_use]
pub fn waits_for(&self, pipeline: PipelineRef) -> &[PipelineRef] {
&self.waits[pipeline as usize]
}
pub fn all(&self) -> impl Iterator<Item = PipelineRef> {
0..u32::try_from(self.waits.len()).unwrap_or(u32::MAX)
}
fn placed(&self, node: NodeRef) -> Placed {
self.of[node as usize].expect("a node under the root of the plan it was walked from")
}
fn fresh(&mut self) -> PipelineRef {
self.waits.push(Vec::new());
u32::try_from(self.waits.len() - 1).unwrap_or(u32::MAX)
}
fn waits_on(&mut self, pipeline: PipelineRef, on: PipelineRef) {
self.waits[pipeline as usize].push(on);
}
fn number(&mut self) -> OperatorRef {
let id = self.operators;
self.operators += 1;
id
}
fn walk(&mut self, plan: &Plan, node: NodeRef, pipeline: PipelineRef) {
let operator = self.number();
match *plan.node(node) {
Node::Aggregate { input, .. }
| Node::Sort { input, .. }
| Node::TopN { input, .. }
| Node::Distinct { input, .. } => {
let below = self.fresh();
self.waits_on(pipeline, below);
self.of[node as usize] = Some(Placed { operator, gathered: None, pipeline: below });
self.walk(plan, input, below);
}
Node::Join { left, right, .. } | Node::SetOp { left, right, .. } => {
let gathered = self.number();
let first = self.fresh();
let second = self.fresh();
self.waits_on(second, first);
self.waits_on(pipeline, second);
self.of[node as usize] =
Some(Placed { operator, gathered: Some(gathered), pipeline: second });
self.walk(plan, right, first);
self.walk(plan, left, second);
}
Node::CrossProduct { left, right } => {
let gathered = self.number();
let aside = self.fresh();
self.waits_on(pipeline, aside);
self.of[node as usize] =
Some(Placed { operator, gathered: Some(gathered), pipeline });
self.walk(plan, right, aside);
self.walk(plan, left, pipeline);
}
ref other => {
self.of[node as usize] = Some(Placed { operator, gathered: None, pipeline });
for child in other.children().into_iter().flatten() {
self.walk(plan, child, pipeline);
}
}
}
}
}
pub const ROOT: PipelineRef = 0;
#[cfg(test)]
mod tests {
use super::Shape;
use crate::plan::Plan;
fn shaped(text: &str) -> (Plan, Shape) {
let plan =
Plan::parse(text).unwrap_or_else(|error| panic!("{text} did not parse: {error}"));
let shape = Shape::of(&plan);
(plan, shape)
}
#[test]
fn a_plan_with_nothing_that_buffers_is_one_pipeline() {
let (plan, shape) = shaped(concat!(
"Filter (#0.0::INTEGER > 1::INTEGER)::BOOLEAN\n",
" Get memory.main.t AS t #0 [a::INTEGER]\n",
));
assert_eq!(shape.pipelines(), 1);
assert_eq!(shape.pipeline(plan.root()), 0);
assert!(shape.waits_for(0).is_empty());
}
#[test]
fn a_sort_ends_the_pipeline_below_it_and_the_one_above_waits() {
let (plan, shape) = shaped(concat!(
"Sort [#0.0::INTEGER ASC NULLS LAST]\n",
" Get memory.main.t AS t #0 [a::INTEGER]\n",
));
assert_eq!(shape.pipelines(), 2);
assert_eq!(shape.pipeline(plan.root()), 1, "the sort is the sink of the one below");
assert_eq!(shape.waits_for(0), [1]);
assert!(shape.waits_for(1).is_empty());
}
#[test]
fn a_join_is_two_pipelines_in_the_order_they_have_to_run() {
let (plan, shape) = shaped(concat!(
"Join INNER on=[(#0.0::INTEGER = #1.0::INTEGER)::BOOLEAN]\n",
" Get memory.main.l AS l #0 [a::INTEGER]\n",
" Get memory.main.r AS r #1 [a::INTEGER]\n",
));
let [left, right] = plan.node(plan.root()).children();
assert_eq!(shape.pipelines(), 3);
assert_eq!(shape.pipeline(right.unwrap()), 1, "the gathered side runs first");
assert_eq!(shape.pipeline(left.unwrap()), 2, "the probing side is the second");
assert_eq!(shape.pipeline(plan.root()), 2, "and the join is its sink");
assert_eq!(shape.waits_for(2), [1]);
assert_eq!(shape.waits_for(0), [2]);
}
#[test]
fn a_cross_product_keeps_its_left_side_where_it_was() {
let (plan, shape) = shaped(concat!(
"CrossProduct\n",
" Get memory.main.l AS l #0 [a::INTEGER]\n",
" Get memory.main.r AS r #1 [a::INTEGER]\n",
));
let [left, right] = plan.node(plan.root()).children();
assert_eq!(shape.pipelines(), 2);
assert_eq!(shape.pipeline(plan.root()), 0, "the product streams");
assert_eq!(shape.pipeline(left.unwrap()), 0, "and so does the side it streams");
assert_eq!(shape.pipeline(right.unwrap()), 1, "the side that is kept is its own");
assert_eq!(shape.waits_for(0), [1]);
}
#[test]
fn two_sorts_under_one_another_are_three_pipelines_in_a_line() {
let (plan, shape) = shaped(concat!(
"Sort [#0.0::INTEGER ASC NULLS LAST]\n",
" Limit 10 offset 0\n",
" Sort [#0.0::INTEGER DESC NULLS FIRST]\n",
" Get memory.main.t AS t #0 [a::INTEGER]\n",
));
assert_eq!(shape.pipelines(), 3);
assert_eq!(shape.pipeline(plan.root()), 1);
assert_eq!(shape.waits_for(0), [1]);
assert_eq!(shape.waits_for(1), [2]);
assert!(shape.waits_for(2).is_empty());
}
#[test]
fn a_parent_is_numbered_before_everything_under_it() {
let (plan, shape) = shaped(concat!(
"Sort [#0.0::INTEGER ASC NULLS LAST]\n",
" Filter (#0.0::INTEGER > 1::INTEGER)::BOOLEAN\n",
" Get memory.main.t AS t #0 [a::INTEGER]\n",
));
let filter = plan.node(plan.root()).children()[0].unwrap();
let get = plan.node(filter).children()[0].unwrap();
assert_eq!(shape.operator(plan.root()), 0);
assert_eq!(shape.operator(filter), 1);
assert_eq!(shape.operator(get), 2);
assert_eq!(shape.operators(), 3);
assert_eq!(shape.gathered(plan.root()), None, "one input, nothing to hold");
}
#[test]
fn a_node_with_two_inputs_is_two_operators_and_the_first_side_is_numbered_first() {
let (plan, shape) = shaped(concat!(
"Join INNER on=[(#0.0::INTEGER = #1.0::INTEGER)::BOOLEAN]\n",
" Get memory.main.l AS l #0 [a::INTEGER]\n",
" Get memory.main.r AS r #1 [a::INTEGER]\n",
));
let [left, right] = plan.node(plan.root()).children();
assert_eq!(shape.operator(plan.root()), 0);
assert_eq!(shape.gathered(plan.root()), Some(1), "the gather is an operator of its own");
assert_eq!(shape.operator(right.unwrap()), 2, "the side that has to finish first");
assert_eq!(shape.operator(left.unwrap()), 3);
assert_eq!(shape.operators(), 4);
}
}