Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
220 changes: 166 additions & 54 deletions crates/integration-tests/tests/planner_layering_example1.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,9 @@
//!
//! MVP scope: Stage 1 = Pass 1 + the identical-expression rule only (no
//! window-composition variants); Stage 2 = physical operator implementation
//! only (no materialization). Expected counts are 1 → 6 → 6 → 1, a subset of
//! the doc's 1 → 54 → 156 → 1.
//! only (no materialization). Counts follow the planner's output (user
//! decision): 1 → 24 → 24 → 1, because Pass 1 also offers exact accumulators.
//! The doc's 1 → 54 → 156 → 1 needs window composition and materialization.

use std::collections::{BTreeMap, BTreeSet, HashSet};

Expand Down Expand Up @@ -342,6 +343,7 @@ fn example1_workload() -> PlanningWorkload {
enum Q2Option {
Exact,
CountMinHeapPerJob,
CountSketchHeapPerJob,
Hydra,
}

Expand Down Expand Up @@ -430,13 +432,34 @@ fn sketch_options(dag: &impl ExportedDag, nodes: &HashSet<NodeId>) -> BTreeSet<Q
{
Q2Option::CountMinHeapPerJob
}
GroupingStrategy::PerSubpopulationInstance
if *kind.algorithm() == SketchAlgorithm::CountSketchWithHeap =>
{
Q2Option::CountSketchHeapPerJob
}
other => panic!("summary family outside Example 1: {kind:?} {other:?}"),
}),
_ => None,
})
.collect()
}

/// Exact accumulator kinds built in `nodes` (e.g. `["Rate", "Sum"]`).
fn exact_accumulators(dag: &impl ExportedDag, nodes: &HashSet<NodeId>) -> Vec<String> {
let mut kinds: Vec<_> = nodes
.iter()
.filter_map(|&id| match payload(dag, id) {
Operator::ASAP(ASAPOp::SummaryAgg {
family: FieldDataType::ExactAggregate(kind, _),
..
}) => Some(format!("{kind:?}")),
_ => None,
})
.collect();
kinds.sort();
kinds
}

fn roots(query_roots: &[NodeId]) -> (NodeId, NodeId) {
assert_eq!(query_roots.len(), 2, "every candidate covers Q1 and Q2");
(query_roots[0], query_roots[1])
Expand Down Expand Up @@ -482,16 +505,41 @@ fn pipeline() -> (
(workload, logical, physical)
}

/// The six Example 1 MVP combinations: three Q2 options × separate/shared input.
fn expected_combinations() -> BTreeSet<(Q2Option, bool)> {
[
/// One candidate's local choices: Q1's exact accumulators, Q2's top-k option
/// and Q2's exact accumulators.
type Choices = (Vec<String>, Q2Option, Vec<String>);

fn choices(dag: &impl ExportedDag, query_roots: &[NodeId]) -> Choices {
let (q1, q2) = roots(query_roots);
let (c1, c2) = (closure(dag, q1), closure(dag, q2));
let (option, _) = classify(dag, query_roots);
(
exact_accumulators(dag, &c1),
option,
exact_accumulators(dag, &c2),
)
}

/// The 24 Pass 1 combinations: Q1's rate and sum each raw or an exact
/// accumulator (4) × Q2's top-k exact, Count-Min + heap or CountSketch + heap
/// (3) × Q2's sum_over_time raw or an exact accumulator (2).
fn expected_choices() -> BTreeSet<Choices> {
let q1 = [vec![], vec!["Rate"], vec!["Sum"], vec!["Rate", "Sum"]];
let q2 = [
Q2Option::Exact,
Q2Option::CountMinHeapPerJob,
Q2Option::Hydra,
]
.into_iter()
.flat_map(|o| [(o, false), (o, true)])
.collect()
Q2Option::CountSketchHeapPerJob,
];
let owned = |kinds: &[&str]| kinds.iter().map(|k| k.to_string()).collect::<Vec<_>>();
let mut all = BTreeSet::new();
for a in &q1 {
for &option in &q2 {
for b in [vec![], vec!["Sum"]] {
all.insert((owned(a), option, owned(&b)));
}
}
}
all
}

// ── Workload ─────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -560,21 +608,24 @@ fn stage0_queries_do_not_share_nodes() {

// ── Stage 1 ──────────────────────────────────────────────────────────────

/// Stage 1 outputs exactly the 3 Q2 options × {separate, shared input} = 6 candidates.
/// Stage 1 outputs exactly the 24 Pass 1 combinations, each once, with
/// separate inputs (no Pass 2 sharing yet).
#[test]
#[ignore = "Pass 1 yields 24 candidates (exact-accumulator alternatives for Q1's rate and sum and Q2's sum, and CountSketch+heap for Q2), there is no Hydra and no Pass 2 shared-input variant"]
fn stage1_has_six_candidates_covering_every_combination() {
fn stage1_has_24_candidates_covering_every_combination() {
let (_, logical, _) = pipeline();
assert_eq!(logical.len(), 6);
assert_eq!(logical.len(), 24);
let found: BTreeSet<_> = logical
.iter()
.map(|c| classify(&c.dag, &c.query_roots))
.map(|c| choices(&c.dag, &c.query_roots))
.collect();
assert_eq!(
found,
expected_combinations(),
"each combination exactly once"
);
assert_eq!(found, expected_choices(), "each combination exactly once");
for c in &logical {
assert!(
!classify(&c.dag, &c.query_roots).1,
"{}: shared input",
c.id
);
}
}

/// Q1 is exact in every Stage 1 candidate: no summary is reachable from its root.
Expand All @@ -592,10 +643,11 @@ fn stage1_q1_is_always_exact() {
}
}

/// Q2's summary families are exactly Count-Min + heap per job and Hydra over all jobs.
/// Q2's summary families are exactly Count-Min + heap and CountSketch + heap
/// per job, and Hydra over all jobs.
#[test]
#[ignore = "Pass 1 also offers CountSketchWithHeap for Q2, which sketch_options rejects as outside Example 1 (spec ambiguity 8); Pass 1 has no Hydra alternative"]
fn stage1_q2_summary_families_are_count_min_heap_and_hydra() {
#[ignore = "missing feature: Pass 1 has no Hydra alternative"]
fn stage1_q2_summary_families_are_heap_sketches_and_hydra() {
let (_, logical, _) = pipeline();
let families: BTreeSet<_> = logical
.iter()
Expand All @@ -606,13 +658,17 @@ fn stage1_q2_summary_families_are_count_min_heap_and_hydra() {
.collect();
assert_eq!(
families,
BTreeSet::from([Q2Option::CountMinHeapPerJob, Q2Option::Hydra])
BTreeSet::from([
Q2Option::CountMinHeapPerJob,
Q2Option::CountSketchHeapPerJob,
Q2Option::Hydra
])
);
}

/// Sharing adds a variant and keeps the independent one, for every Q2 option.
#[test]
#[ignore = "Pass 1 also offers CountSketchWithHeap for Q2, which sketch_options rejects as outside Example 1 (spec ambiguity 8); Pass 2 (identical-expression sharing) is not implemented, so there is no shared variant"]
#[ignore = "missing features: Pass 2 (identical-expression sharing) is not implemented, so there is no shared-input variant; Pass 1 has no Hydra alternative"]
fn stage1_keeps_independent_and_shared_variants() {
let (_, logical, _) = pipeline();
let found: Vec<_> = logical
Expand All @@ -622,6 +678,7 @@ fn stage1_keeps_independent_and_shared_variants() {
for option in [
Q2Option::Exact,
Q2Option::CountMinHeapPerJob,
Q2Option::CountSketchHeapPerJob,
Q2Option::Hydra,
] {
assert!(
Expand Down Expand Up @@ -666,12 +723,11 @@ fn stage1_candidates_are_valid_and_uniquely_named() {

// ── Stage 2 ──────────────────────────────────────────────────────────────

/// No candidate is discarded before Stage 3: Stage 2 maps the 6 logical candidates one-to-one.
/// No candidate is discarded before Stage 3: Stage 2 maps the 24 logical candidates one-to-one.
#[test]
#[ignore = "Stage 2 maps the 24 Pass 1 candidates one-to-one (from_logical is a bijection), but the spec expects 6"]
fn stage2_keeps_every_logical_candidate() {
let (_, logical, physical) = pipeline();
assert_eq!(physical.len(), 6);
assert_eq!(physical.len(), 24);
let sources: BTreeSet<_> = physical.iter().map(|p| p.from_logical.as_str()).collect();
let logical_ids: BTreeSet<_> = logical.iter().map(|c| c.id.as_str()).collect();
assert_eq!(sources, logical_ids);
Expand All @@ -681,7 +737,6 @@ fn stage2_keeps_every_logical_candidate() {

/// Stage 2 preserves each logical candidate's Q2 option and input sharing.
#[test]
#[ignore = "Pass 1 also offers CountSketchWithHeap for Q2, which sketch_options rejects as outside Example 1 (spec ambiguity 8)"]
fn stage2_preserves_logical_choices() {
let (_, logical, physical) = pipeline();
for p in &physical {
Expand All @@ -698,14 +753,13 @@ fn stage2_preserves_logical_choices() {

/// Exact TopK is implemented as a sort followed by a limit.
#[test]
#[ignore = "Pass 1 also offers CountSketchWithHeap for Q2, which sketch_options rejects as outside Example 1 (spec ambiguity 8); also 8, not 2, candidates have an exact Q2"]
fn stage2_exact_topk_is_sort_then_limit() {
let (_, _, physical) = pipeline();
let exact: Vec<_> = physical
.iter()
.filter(|p| classify(&p.dag, &p.query_roots).0 == Q2Option::Exact)
.collect();
assert_eq!(exact.len(), 2);
assert_eq!(exact.len(), 8);
for p in exact {
let sort_then_limit = p.dag.edges.iter().any(|e| {
relational(payload(&p.dag, e.producer)).as_deref() == Some("sort")
Expand All @@ -717,7 +771,6 @@ fn stage2_exact_topk_is_sort_then_limit() {

/// A summary Q2 is a build node feeding a top-10 estimation node, with no merge.
#[test]
#[ignore = "Pass 1 also offers CountSketchWithHeap for Q2, which sketch_options rejects as outside Example 1 (spec ambiguity 8)"]
fn stage2_summary_topk_is_build_then_estimate() {
let (_, _, physical) = pipeline();
for p in &physical {
Expand Down Expand Up @@ -816,20 +869,59 @@ fn compile_in_runtime(p: &PhysicalCandidate) -> Result<(), String> {
}

/// Runtime capability check (added by the implementer, not part of the
/// spec): the physical planner compiles every Stage 2 candidate.
/// spec): the physical planner compiles every candidate Stage 3 finds valid,
/// and rejects the invalid ones (Count-Min over weights not proven
/// non-negative) for the same reason Stage 3 gives.
#[test]
#[ignore = "runtime gaps: the physical planner rejects CMS+heap without a non-negative \
weight contract (Stage 3 rejects these as invalid too), and rejects every \
CountSketch+heap top-k because the IR's SummaryEstimate{TopK} output schema \
(partition keys + one encoded top-k column) is not the keyed-evaluation shape \
it builds (input columns + item + Float64 score); 8 of 24 compile, all exact-Q2"]
fn stage2_every_candidate_compiles_in_the_physical_planner() {
fn stage2_runtime_compiles_exactly_the_candidates_stage3_finds_valid() {
let (workload, _, physical) = pipeline();
let selection = stage3_select(&workload, &physical, PlanningModels::builtin());
let invalid: BTreeMap<_, _> = selection
.rejected
.iter()
.filter(|r| !r.valid)
.map(|r| (r.id.as_str(), r.reason.as_str()))
.collect();
assert_eq!(invalid.len(), 8, "the Count-Min + heap candidates");
for p in &physical {
let compiled = compile_in_runtime(p);
match invalid.get(p.id.as_str()) {
None => compiled.unwrap(),
Some(reason) => {
assert!(
reason.contains("CmsWithHeap needs non-negative update weights"),
"{}: {reason}",
p.id
);
let error = compiled.expect_err(&p.id);
assert!(
error.contains("CMS requires a nonnegative weight contract"),
"{error}"
);
}
}
}
}

/// A CountSketch+heap top-k readout compiles: the IR's derived readout schema
/// is the ranked-rows shape the runtime's keyed evaluation produces.
#[test]
fn stage2_count_sketch_heap_topk_compiles_in_the_physical_planner() {
let (_, _, physical) = pipeline();
let failures: Vec<_> = physical
let count_sketch: Vec<_> = physical
.iter()
.filter_map(|p| compile_in_runtime(p).err())
.filter(|p| {
p.dag.nodes.iter().any(|n| {
matches!(&n.payload, Operator::ASAP(ASAPOp::SummaryAgg {
family: FieldDataType::Sketch(kind, _), ..
}) if *kind.algorithm() == SketchAlgorithm::CountSketchWithHeap)
})
})
.collect();
assert!(failures.is_empty(), "{failures:#?}");
assert!(!count_sketch.is_empty());
for p in count_sketch {
compile_in_runtime(p).unwrap();
}
}

/// The plan Stage 3 selects compiles in the physical planner (added by the
Expand Down Expand Up @@ -884,13 +976,23 @@ fn stage3_selects_cheapest_valid() {
}
}

/// Every node is charged exactly once, so a shared input is costed once for both queries.
/// Every node is charged exactly once, so a shared input is costed once for
/// both queries. Stage 3 prices valid candidates only (user decision).
#[test]
#[ignore = "Stage 3 prices valid candidates only (user decision); the 8 CMS+heap candidates are rejected as invalid (weights not proven non-negative) and have no cost"]
fn stage3_charges_each_node_once() {
let (workload, _, physical) = pipeline();
let selection = stage3_select(&workload, &physical, PlanningModels::builtin());
let invalid: BTreeSet<_> = selection
.rejected
.iter()
.filter(|r| !r.valid)
.map(|r| r.id.as_str())
.collect();
for p in &physical {
if invalid.contains(p.id.as_str()) {
assert!(!selection.costs.contains_key(&p.id), "{} is priced", p.id);
continue;
}
let cost = &selection.costs[&p.id];
let nodes: BTreeSet<_> = p.dag.nodes.iter().map(|n| n.id).collect();
let charged: BTreeSet<_> = cost.per_node.keys().copied().collect();
Expand All @@ -908,29 +1010,39 @@ fn stage3_charges_each_node_once() {
}
}

/// Sharing the input never costs more than reading it separately.
/// Sharing the input never costs more than reading it separately, for the
/// same local choices. Count-Min + heap is invalid here and has no cost.
#[test]
#[ignore = "Pass 1 also offers CountSketchWithHeap for Q2, which sketch_options rejects as outside Example 1 (spec ambiguity 8); there are no shared-input variants to compare"]
#[ignore = "missing features: Pass 2 (identical-expression sharing) is not implemented, so there are no shared-input variants to compare; Pass 1 has no Hydra alternative"]
fn stage3_shared_input_is_not_costlier() {
let (workload, _, physical) = pipeline();
let selection = stage3_select(&workload, &physical, PlanningModels::builtin());
let by_combo: BTreeMap<_, _> = physical
.iter()
.map(|p| {
(
classify(&p.dag, &p.query_roots),
selection.costs[&p.id].total,
)
.filter_map(|p| {
let cost = selection.costs.get(&p.id)?.total;
let shared = classify(&p.dag, &p.query_roots).1;
Some(((choices(&p.dag, &p.query_roots), shared), cost))
})
.collect();
for option in [
Q2Option::Exact,
Q2Option::CountMinHeapPerJob,
Q2Option::CountSketchHeapPerJob,
Q2Option::Hydra,
] {
assert!(
by_combo[&(option, true)] <= by_combo[&(option, false)],
"{option:?}"
by_combo.keys().any(|((_, o, _), _)| *o == option),
"{option:?} has no priced candidate"
);
}
for ((choice, shared), cost) in &by_combo {
if *shared {
let separate = by_combo[&(choice.clone(), false)];
assert!(*cost <= separate, "{choice:?}");
}
}
assert!(
by_combo.keys().any(|(_, shared)| *shared),
"no shared variant"
);
}
Loading