From 3c4663acfcfa602e584b0e5cbaa402372fe6a9b7 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:52:09 +0000 Subject: [PATCH 1/2] test: acceptance spec and pending tests for planner-layering Example 1 Define MVP acceptance for #509 Example 1 before the Phase C stage APIs exist: 1 -> 6 -> 6 -> 1 candidates (Pass 1 x identical-expression sharing; physical operator implementation only, no materialization). - Spec with per-stage candidate tables, invariants and doc ambiguities. - Integration tests against todo!() stage stubs, ignored until the stages land; workload and Stage 0 frontend-shape tests run now. - Expected-only asap-stage-pipeline/v1 fixture for the DAG viewer. Co-Authored-By: Claude Opus 5.5 --- .../tests/planner_layering_example1.rs | 671 +++++++ .../planner-layering-example1-acceptance.md | 150 ++ .../planner-layering-example1.expected.json | 1603 +++++++++++++++++ 3 files changed, 2424 insertions(+) create mode 100644 crates/integration-tests/tests/planner_layering_example1.rs create mode 100644 docs/design_docs/proposals/planner-layering-example1-acceptance.md create mode 100644 tools/dag-viewer/examples/planner-layering-example1.expected.json diff --git a/crates/integration-tests/tests/planner_layering_example1.rs b/crates/integration-tests/tests/planner_layering_example1.rs new file mode 100644 index 00000000..ad2c1a0b --- /dev/null +++ b/crates/integration-tests/tests/planner_layering_example1.rs @@ -0,0 +1,671 @@ +//! Acceptance tests for #509 "Example 1: Aggregation over dimensions", MVP scope. +//! +//! Spec: `docs/design_docs/proposals/planner-layering-example1-acceptance.md`. +//! Written by the test designer before the Phase C stage APIs exist: every +//! stage test is `#[ignore]`d and calls the `todo!()` stubs in [`stages`], +//! which the implementer replaces with the real entry points. +//! +//! 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. + +use std::collections::{BTreeMap, BTreeSet, HashSet}; + +use asap_aware_mapping::PlanningModels; +use asap_types::ir::export::{LogicalASAPDAG, LogicalASAPNodeId, LogicalASAPOperatorPayload}; +use asap_types::post_asap::sketch::{GroupingStrategy, HydraKind, SketchAlgorithm}; +use asap_types::pre_asap::schema::FieldDataType; +use asap_types::types::AccuracyTarget; +use asap_types::workload::{ + AccuracyRequirement, DataArrival, DataDistribution, DataWorkload, DurationMs, Evidence, + EvidenceSource, LatencyRequirement, PlanningWorkload, Predictability, Query, QueryLanguage, + QueryRequirements, QueryTimeScope, QueryWorkload, Rate, RepeatedDemand, RepeatingEntry, + RepetitionInterval, TimeSelection, +}; + +/// Hypothetical Phase C stage API. Replace each `todo!()` with the real call. +#[allow(dead_code, unused_variables)] +mod stages { + use super::*; + + /// One whole-workload candidate. `LogicalASAPDAG::root` names one query, + /// so `query_roots` carries one root per workload entry, in + /// `QueryWorkload::entries()` order (`[q1, q2]`), until the export + /// supports several roots. + #[derive(Debug, Clone)] + pub struct LogicalCandidate { + pub id: String, + pub label: String, + pub dag: LogicalASAPDAG, + pub query_roots: Vec, + } + + /// Placeholder for `asap_types::ir::export::PhysicalASAPDAG` (phase A). + /// Its payload is an alias of the logical payload, so payload-kind checks + /// carry over unchanged. + pub type PhysicalASAPDAG = LogicalASAPDAG; + + /// One Stage 2 candidate, derived from exactly one Stage 1 candidate. + #[derive(Debug, Clone)] + pub struct PhysicalCandidate { + pub id: String, + pub from_logical: String, + pub label: String, + pub dag: PhysicalASAPDAG, + pub query_roots: Vec, + } + + /// Whole-workload cost of one physical candidate; `per_node` has one + /// entry per DAG node, so a shared node is charged once. + #[derive(Debug, Clone)] + pub struct CandidateCost { + pub total: f64, + pub per_node: BTreeMap, + } + + /// A candidate Stage 3 did not select. `valid == false` means it failed + /// an accuracy, latency or capability check; `true` means it lost on cost. + #[derive(Debug, Clone)] + pub struct Rejection { + pub id: String, + pub valid: bool, + pub reason: String, + } + + /// Stage 3 output. Costs are keyed by physical candidate id; the viewer + /// document attaches them to the Stage 2 entries. + #[derive(Debug, Clone)] + pub struct Selection { + pub selected: String, + pub costs: BTreeMap, + pub rejected: Vec, + } + + /// Stage 0: frontends lower every query into one summary-free workload DAG. + pub fn stage0_logical(workload: &PlanningWorkload) -> LogicalCandidate { + todo!("Phase C: Stage 0 workload LogicalDAG") + } + + /// Stage 1: Pass 1 local alternatives × Pass 2 identical-expression + /// sharing. Takes the workload for accuracy requirements and time selection. + pub fn stage1_logical_asap( + workload: &PlanningWorkload, + logical: &LogicalCandidate, + ) -> Vec { + todo!("Phase C: Stage 1 CandidateLogicalASAPDAGs") + } + + /// Stage 2: physical operator implementation of every logical candidate + /// (no materialization in the MVP). + pub fn stage2_physical( + workload: &PlanningWorkload, + logical: &[LogicalCandidate], + ) -> Vec { + todo!("Phase C: Stage 2 CandidatePhysicalASAPDAGs") + } + + /// Stage 3: reject invalid candidates, cost the rest for the whole + /// workload, select the cheapest. + pub fn stage3_select( + workload: &PlanningWorkload, + physical: &[PhysicalCandidate], + models: PlanningModels<'_>, + ) -> Selection { + todo!("Phase C: Stage 3 plan selection") + } + + /// Stands in for `node.output_state.timing == IngestionTime` until + /// `PhysicalASAPDAG` lands. + pub fn runs_at_ingestion(candidate: &PhysicalCandidate, node: LogicalASAPNodeId) -> bool { + todo!("Phase C: read PhysicalASAPDAGNode::output_state.timing") + } +} + +use stages::*; + +// ── Example 1 workload ─────────────────────────────────────────────────── + +const Q1: &str = "sum by (job) (rate(http_requests_total[1m]))"; +const Q2: &str = "topk by (job) (10, sum_over_time(http_requests_total[1m]))"; + +fn declared(value: T) -> Evidence { + Evidence { + value: Some(value), + source: EvidenceSource::Declared, + ..Default::default() + } +} + +fn dashboard_panel(query: &str, requirements: QueryRequirements) -> RepeatingEntry { + RepeatingEntry { + query: Query(query.into()), + demand: RepeatedDemand::FixedInterval(RepetitionInterval(10_000)), + requirements, + predictability: Predictability::Predictable { known_at: None }, + time_selection: TimeSelection { + scope: QueryTimeScope::RealTime, + lookback: Some(DurationMs(60_000)), + as_of: None, + }, + } +} + +/// Example 1 queries over the shared data workload of #509. +fn example1_workload() -> PlanningWorkload { + PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: None, + repeating_queries: Some(vec![ + dashboard_panel( + Q1, + QueryRequirements { + accuracy: AccuracyRequirement::Explicit(AccuracyTarget::Exact), + response_latency: LatencyRequirement::Unspecified, + }, + ), + dashboard_panel( + Q2, + QueryRequirements { + accuracy: AccuracyRequirement::Explicit(AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.001, + }), + response_latency: LatencyRequirement::ExplicitMaxMs(100.0), + }, + ), + ]), + }, + data_workload: Some(DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + data_ingestion_interval: declared(DurationMs(15_000)), + ingestion_volume: Evidence::default(), + ingestion_rate: declared(Rate(1_000_000.0 / 15.0)), + input_cardinality: declared(1_000_000), + distribution: declared(DataDistribution::Zipf), + }), + } +} + +// ── DAG helpers ────────────────────────────────────────────────────────── + +/// Q2's local option, read off its summary build node. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +enum Q2Option { + Exact, + CountMinHeapPerJob, + Hydra, +} + +/// Every node `root` depends on, including itself. +fn closure(dag: &LogicalASAPDAG, root: LogicalASAPNodeId) -> HashSet { + let mut seen = HashSet::from([root]); + let mut stack = vec![root]; + while let Some(node) = stack.pop() { + for edge in dag.edges.iter().filter(|e| e.consumer == node) { + if seen.insert(edge.producer) { + stack.push(edge.producer); + } + } + } + seen +} + +fn payload(dag: &LogicalASAPDAG, id: LogicalASAPNodeId) -> &LogicalASAPOperatorPayload { + &dag.nodes.iter().find(|n| n.id == id).expect("node").payload +} + +fn is_summary(payload: &LogicalASAPOperatorPayload) -> bool { + matches!( + payload, + LogicalASAPOperatorPayload::SummaryAgg { + family: FieldDataType::Sketch(..), + .. + } | LogicalASAPOperatorPayload::SummaryEstimate { .. } + | LogicalASAPOperatorPayload::SummaryMerge + ) +} + +/// Sketch families built in `nodes`, as Example 1's Q2 options. +fn sketch_options(dag: &LogicalASAPDAG, nodes: &HashSet) -> BTreeSet { + nodes + .iter() + .filter_map(|&id| match payload(dag, id) { + LogicalASAPOperatorPayload::SummaryAgg { + family: FieldDataType::Sketch(kind, grouping), + .. + } => Some(match grouping { + GroupingStrategy::SharedMultiSubpopulation { + kind: HydraKind::HydraCms, + .. + } => Q2Option::Hydra, + GroupingStrategy::PerSubpopulationInstance + if *kind.algorithm() == SketchAlgorithm::CmsWithHeap => + { + Q2Option::CountMinHeapPerJob + } + other => panic!("summary family outside Example 1: {kind:?} {other:?}"), + }), + _ => None, + }) + .collect() +} + +fn roots(query_roots: &[LogicalASAPNodeId]) -> (LogicalASAPNodeId, LogicalASAPNodeId) { + assert_eq!(query_roots.len(), 2, "every candidate covers Q1 and Q2"); + (query_roots[0], query_roots[1]) +} + +/// Q2's option and whether Q1 and Q2 share any node, for one candidate. +fn classify(dag: &LogicalASAPDAG, query_roots: &[LogicalASAPNodeId]) -> (Q2Option, bool) { + let (q1, q2) = roots(query_roots); + let (c1, c2) = (closure(dag, q1), closure(dag, q2)); + let options = sketch_options(dag, &c2); + assert!(options.len() <= 1, "Q2 uses one local option: {options:?}"); + let option = options.into_iter().next().unwrap_or(Q2Option::Exact); + (option, !c1.is_disjoint(&c2)) +} + +/// The relational operator's wire `kind` (`"scan"`, `"sort"`, …); the +/// operator enum itself is not public outside `asap-types`. +fn relational(payload: &LogicalASAPOperatorPayload) -> Option { + match payload { + LogicalASAPOperatorPayload::Relational { .. } => { + let json = serde_json::to_value(payload).expect("payload serializes"); + json["operator"]["kind"].as_str().map(str::to_owned) + } + _ => None, + } +} + +fn pipeline() -> ( + PlanningWorkload, + Vec, + Vec, +) { + let workload = example1_workload(); + let logical = stage1_logical_asap(&workload, &stage0_logical(&workload)); + let physical = stage2_physical(&workload, &logical); + (workload, logical, physical) +} + +/// The six Example 1 MVP combinations: three Q2 options × separate/shared input. +fn expected_combinations() -> BTreeSet<(Q2Option, bool)> { + [ + Q2Option::Exact, + Q2Option::CountMinHeapPerJob, + Q2Option::Hydra, + ] + .into_iter() + .flat_map(|o| [(o, false), (o, true)]) + .collect() +} + +// ── Workload ───────────────────────────────────────────────────────────── + +/// The encoded workload is valid and normalizes to Q1 then Q2. +#[test] +fn workload_encodes_example1() { + let workload = example1_workload(); + workload.validate().expect("valid workload"); + let entries: Vec<_> = workload.query_workload.entries().collect(); + assert_eq!(entries.len(), 2); + assert_eq!(entries[0].query.0, Q1); + assert_eq!(entries[1].query.0, Q2); +} + +// ── Stage 0 ────────────────────────────────────────────────────────────── + +/// Today's PromQL frontend already lowers each query to the doc's Stage 0 chain. +#[test] +fn stage0_frontend_lowers_each_query_to_doc_chain() { + let roots = asap_frontend_promql::unified::lower_promql_query_workload(&example1_workload(), 0) + .expect("Example 1 lowers"); + let chains: Vec> = roots + .iter() + .map(|root| { + let dag = asap_types::ir::export::compile_logical_asap_query(root).expect("compiles"); + let json = serde_json::to_value(&dag).expect("serializes"); + let mut ops: Vec = json["nodes"] + .as_array() + .unwrap() + .iter() + .map(|n| { + let op = &n["payload"]["operator"]; + match op["measures"][0]["kind"].as_str() { + Some(measure) => format!("aggregate:{measure}"), + None => op["kind"].as_str().unwrap_or("?").to_owned(), + } + }) + .collect(); + ops.sort(); // node ids are assigned in post-order; compare as a set of operations + ops + }) + .collect(); + assert_eq!( + chains, + [ + ["aggregate:rate", "aggregate:sum", "scan", "time_range"], + ["aggregate:sum", "aggregate:top_k", "scan", "time_range"], + ] + ); +} + +/// Stage 0 yields one workload DAG with one root per query and no summaries. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage0_one_summary_free_workload_dag() { + let stage0 = stage0_logical(&example1_workload()); + stage0.dag.validate().expect("valid DAG"); + roots(&stage0.query_roots); + assert!(stage0.dag.nodes.iter().all(|n| !is_summary(&n.payload))); +} + +/// Stage 0 keeps Q1 and Q2 separate; sharing is a Stage 1 decision. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage0_queries_do_not_share_nodes() { + let stage0 = stage0_logical(&example1_workload()); + let (q1, q2) = roots(&stage0.query_roots); + assert!(closure(&stage0.dag, q1).is_disjoint(&closure(&stage0.dag, q2))); +} + +// ── Stage 1 ────────────────────────────────────────────────────────────── + +/// Stage 1 outputs exactly the 3 Q2 options × {separate, shared input} = 6 candidates. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage1_has_six_candidates_covering_every_combination() { + let (_, logical, _) = pipeline(); + assert_eq!(logical.len(), 6); + let found: BTreeSet<_> = logical + .iter() + .map(|c| classify(&c.dag, &c.query_roots)) + .collect(); + assert_eq!( + found, + expected_combinations(), + "each combination exactly once" + ); +} + +/// Q1 is exact in every Stage 1 candidate: no summary is reachable from its root. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage1_q1_is_always_exact() { + let (_, logical, _) = pipeline(); + for c in &logical { + let (q1, _) = roots(&c.query_roots); + let reach = closure(&c.dag, q1); + assert!( + reach.iter().all(|&id| !is_summary(payload(&c.dag, id))), + "{}: Q1 reaches a summary", + c.id + ); + } +} + +/// Q2's summary families are exactly Count-Min + heap per job and Hydra over all jobs. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage1_q2_summary_families_are_count_min_heap_and_hydra() { + let (_, logical, _) = pipeline(); + let families: BTreeSet<_> = logical + .iter() + .flat_map(|c| { + let (_, q2) = roots(&c.query_roots); + sketch_options(&c.dag, &closure(&c.dag, q2)) + }) + .collect(); + assert_eq!( + families, + BTreeSet::from([Q2Option::CountMinHeapPerJob, Q2Option::Hydra]) + ); +} + +/// Sharing adds a variant and keeps the independent one, for every Q2 option. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage1_keeps_independent_and_shared_variants() { + let (_, logical, _) = pipeline(); + let found: Vec<_> = logical + .iter() + .map(|c| classify(&c.dag, &c.query_roots)) + .collect(); + for option in [ + Q2Option::Exact, + Q2Option::CountMinHeapPerJob, + Q2Option::Hydra, + ] { + assert!( + found.contains(&(option, false)), + "{option:?} independent missing" + ); + assert!(found.contains(&(option, true)), "{option:?} shared missing"); + } +} + +/// Only the raw input is shared between Q1 and Q2; no summary is shared. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage1_shares_input_but_never_a_summary() { + let (_, logical, _) = pipeline(); + for c in &logical { + let (q1, q2) = roots(&c.query_roots); + let shared = &closure(&c.dag, q1) & &closure(&c.dag, q2); + for id in shared { + let p = payload(&c.dag, id); + assert!( + matches!(relational(p).as_deref(), Some("scan" | "time_range")), + "{}: shared node {id:?} is not the range selector input: {p:?}", + c.id + ); + } + } +} + +/// Stage 1 candidates are valid DAGs with unique ids. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage1_candidates_are_valid_and_uniquely_named() { + let (_, logical, _) = pipeline(); + let ids: BTreeSet<_> = logical.iter().map(|c| c.id.as_str()).collect(); + assert_eq!(ids.len(), logical.len()); + for c in &logical { + c.dag.validate().unwrap_or_else(|e| panic!("{}: {e}", c.id)); + } +} + +// ── Stage 2 ────────────────────────────────────────────────────────────── + +/// No candidate is discarded before Stage 3: Stage 2 maps the 6 logical candidates one-to-one. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage2_keeps_every_logical_candidate() { + let (_, logical, physical) = pipeline(); + assert_eq!(physical.len(), 6); + 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); + let ids: BTreeSet<_> = physical.iter().map(|p| p.id.as_str()).collect(); + assert_eq!(ids.len(), physical.len()); +} + +/// Stage 2 preserves each logical candidate's Q2 option and input sharing. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage2_preserves_logical_choices() { + let (_, logical, physical) = pipeline(); + for p in &physical { + let source = logical.iter().find(|c| c.id == p.from_logical).unwrap(); + assert_eq!( + classify(&p.dag, &p.query_roots), + classify(&source.dag, &source.query_roots), + "{} vs {}", + p.id, + source.id + ); + } +} + +/// Exact TopK is implemented as a sort followed by a limit. +#[test] +#[ignore = "pending Phase C stage APIs"] +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); + for p in exact { + let sort_then_limit = p.dag.edges.iter().any(|e| { + relational(payload(&p.dag, e.producer)).as_deref() == Some("sort") + && relational(payload(&p.dag, e.consumer)).as_deref() == Some("limit") + }); + assert!(sort_then_limit, "{}: no sort → limit", p.id); + } +} + +/// A summary Q2 is a build node feeding a top-10 estimation node, with no merge. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage2_summary_topk_is_build_then_estimate() { + let (_, _, physical) = pipeline(); + for p in &physical { + if classify(&p.dag, &p.query_roots).0 == Q2Option::Exact { + continue; + } + let kinds: Vec<_> = p.dag.nodes.iter().map(|n| &n.payload).collect(); + assert!( + !kinds + .iter() + .any(|k| matches!(k, LogicalASAPOperatorPayload::SummaryMerge)), + "{}: no window summaries in the MVP, so no merge", + p.id + ); + let build_to_estimate = p.dag.edges.iter().any(|e| { + matches!( + payload(&p.dag, e.producer), + LogicalASAPOperatorPayload::SummaryAgg { + family: FieldDataType::Sketch(..), + .. + } + ) && matches!( + payload(&p.dag, e.consumer), + LogicalASAPOperatorPayload::SummaryEstimate { .. } + ) + }); + assert!(build_to_estimate, "{}: no build → estimate", p.id); + } +} + +/// With no materialization in the MVP, every node runs at query time. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage2_everything_runs_at_query_time() { + let (_, _, physical) = pipeline(); + for p in &physical { + for n in &p.dag.nodes { + assert!( + !runs_at_ingestion(p, n.id), + "{}: {:?} at ingestion", + p.id, + n.id + ); + } + } +} + +// ── Stage 3 ────────────────────────────────────────────────────────────── + +/// Stage 3 selects one candidate and gives every other one a reason. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage3_selects_one_and_explains_the_rest() { + let (workload, _, physical) = pipeline(); + let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let all: BTreeSet<_> = physical.iter().map(|p| p.id.clone()).collect(); + let mut accounted: BTreeSet<_> = selection.rejected.iter().map(|r| r.id.clone()).collect(); + assert_eq!( + accounted.len(), + selection.rejected.len(), + "rejected once each" + ); + assert!(selection.rejected.iter().all(|r| !r.reason.is_empty())); + assert!( + accounted.insert(selection.selected.clone()), + "selected is not rejected" + ); + assert_eq!(accounted, all); +} + +/// The selected plan is the cheapest valid candidate for the whole workload. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage3_selects_cheapest_valid() { + 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(); + let best = selection.costs[&selection.selected].total; + for p in physical.iter().filter(|p| !invalid.contains(p.id.as_str())) { + assert!(best <= selection.costs[&p.id].total, "{} is cheaper", p.id); + } +} + +/// Every node is charged exactly once, so a shared input is costed once for both queries. +#[test] +#[ignore = "pending Phase C stage APIs"] +fn stage3_charges_each_node_once() { + let (workload, _, physical) = pipeline(); + let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + for p in &physical { + 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(); + assert_eq!( + charged, nodes, + "{}: per-node costs cover each node once", + p.id + ); + let sum: f64 = cost.per_node.values().sum(); + assert!( + (cost.total - sum).abs() <= 1e-9 * sum.abs().max(1.0), + "{}", + p.id + ); + } +} + +/// Sharing the input never costs more than reading it separately. +#[test] +#[ignore = "pending Phase C stage APIs"] +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, + ) + }) + .collect(); + for option in [ + Q2Option::Exact, + Q2Option::CountMinHeapPerJob, + Q2Option::Hydra, + ] { + assert!( + by_combo[&(option, true)] <= by_combo[&(option, false)], + "{option:?}" + ); + } +} diff --git a/docs/design_docs/proposals/planner-layering-example1-acceptance.md b/docs/design_docs/proposals/planner-layering-example1-acceptance.md new file mode 100644 index 00000000..2ce1057a --- /dev/null +++ b/docs/design_docs/proposals/planner-layering-example1-acceptance.md @@ -0,0 +1,150 @@ +# Planner layering, Example 1: acceptance spec (MVP) + +Audience: planner designers and the Phase C implementer. +Source: [planner-layering.md](planner-layering.md), "Example 1: Aggregation over +dimensions". Tests: `crates/integration-tests/tests/planner_layering_example1.rs`. +Viewer fixture shape: `tools/dag-viewer/examples/planner-layering-example1.expected.json`. + +This spec was written by a test designer who did not implement the stages. + +## Scope + +| Stage | Doc (#509) | MVP (this spec) | +|---|---|---| +| 0. Frontends | 1 workload `LogicalDAG` | same | +| 1. Logical ASAP | Pass 1 (3) × identical-expression rule (2) × window-composition rule (3 × 3) = **54** | Pass 1 × identical-expression rule = **6**. Window composition is blocked on #511 (#518, #522). | +| 2. Physical ASAP | Materialization options per window form = **156** | Physical operator implementation only, no materialization = **6** (the doc's Raw/Raw plans) | +| 3. Selection | 1 plan | 1 plan | + +Every MVP candidate is one of the doc's candidates: in Stage 1, the one with no +window summary for either query; in Stage 2, the one where both queries are +**Raw** (state rebuilt from the last 1 min of raw samples at every refresh). + +## Workload + +Shared data workload: `continuously_ingesting`, 15 s ingestion interval, +1,000,000 series, about 66,667 samples/s, `zipf`, volume unknown. + +| Query | Repeats | `lookback` | `as_of` | Accuracy | Latency | +|---|---|---|---|---|---| +| Q1 `sum by (job) (rate(http_requests_total[1m]))` | 10 s | 1 m | evaluation time | exact | none | +| Q2 `topk by (job) (10, sum_over_time(http_requests_total[1m]))` | 10 s | 1 m | evaluation time | ε = 0.01, δ = 0.001 | ≤ 100 ms | + +## Stage 0: 1 candidate + +| Query | Chain | +|---|---| +| Q1 | scan `http_requests_total` → range 1m → rate (per series) → sum by (job) | +| Q2 | scan `http_requests_total` → range 1m → sum_over_time (per series) → topk by (job) (10) | + +Invariants: + +* Exactly one candidate, with one root per query, in workload entry order. +* No summary node (sketch `summary_agg`, `summary_estimate`, `summary_merge`). +* Q1 and Q2 share no node. Sharing is a Stage 1 decision. + +## Stage 1: 6 candidates + +| Id | Q1 | Q2 | Input | +|---|---|---|---| +| L1 | exact | exact (per-series sum, top 10 per job) | separate | +| L2 | exact | exact | shared | +| L3 | exact | Count-Min + top-*k* heap, one per job | separate | +| L4 | exact | Count-Min + top-*k* heap, one per job | shared | +| L5 | exact | Hydra (CMS) over all (job, series) keys, heap per job | separate | +| L6 | exact | Hydra (CMS) | shared | + +"Shared" means Q1 and Q2 read one `http_requests_total[1m]` input node +(scan and range). Ids and order are illustrative; tests identify candidates by +structure. + +Invariants: + +* Exactly 6 candidates, one per (Q2 option, separate or shared). No duplicates. +* Every candidate covers both queries. +* Q1 never reaches a summary node. +* The sketch families in Q2 are exactly {`CmsWithHeap` per subpopulation, + Hydra `HydraCms`}. Exact Q2 has no sketch. +* For each Q2 option, both the independent and the shared variant are present. +* Only the scan and range nodes may be shared. No summary is shared. +* Nothing is pruned: Count-Min and Hydra are admitted by ε = 0.01, δ = 0.001, + and Q1's exact target admits only the exact option. + +## Stage 2: 6 candidates + +| Id | From | Q2 physical operators | Timing | +|---|---|---|---| +| P1 | L1 | sort (partition by job) → limit 10 | all query time | +| P2 | L2 | sort → limit 10 | all query time | +| P3 | L3 | CMS + heap build → top-10 estimate | all query time | +| P4 | L4 | CMS + heap build → top-10 estimate | all query time | +| P5 | L5 | Hydra build → top-10 estimate per job | all query time | +| P6 | L6 | Hydra build → top-10 estimate per job | all query time | + +Q1 is the same in all six: scan → range → per-series rate → sum by (job). + +Invariants: + +* Exactly one physical candidate per logical candidate. `from_logical` is a + bijection onto the Stage 1 ids, so no valid candidate is dropped before + Stage 3. +* Each physical candidate keeps its logical Q2 option and input sharing. +* Exact TopK is implemented as a sort followed by a limit. +* A summary Q2 is a sketch build node feeding an estimation node. There is no + merge node, because the MVP has no window summaries. +* No node runs at ingestion time, because the MVP has no materialization. + +## Stage 3: 1 plan + +Invariants: + +* One selected id. Every other candidate is listed once as rejected, with a + reason, and marked invalid (accuracy, latency or capability) or valid but + costlier. +* Each candidate's cost has one entry per DAG node, and `total` is their sum. + A shared node is therefore charged once, for all its consumers. +* The selected plan is no more expensive than any valid candidate. +* For each Q2 option, the shared-input candidate costs no more than its + separate counterpart. + +Expected outcome, not asserted because it depends on the cost and accuracy +models: a shared-input candidate wins, matching the doc's "Raw with a shared +input" winner. An exact Q2 (P1, P2) may be rejected for missing the 100 ms +bound. The doc's other typical winners need window forms or materialization and +are out of MVP scope. + +## Ambiguities and MVP deviations + +1. **Stage 2 is materialization-only in the doc.** Example 1 lists only + window-form or materialization options. The MVP reads "physical operator + implementation" from the Stage 2 section: "TopK as a sort followed by a + limit". This gives one physical candidate per logical candidate. The doc + names no alternative implementation (for example hash versus sort + aggregation) for this example. +2. **Where exact TopK becomes sort + limit.** The Pass 1 diagram already + draws exact Q2 as "sort + limit 10 per job", but the Stage 2 section calls + this a physical choice. Today the frontend emits `aggregate[top_k]`. The + spec requires sort → limit only by Stage 2. +3. **Minimal Pass 2.** Only the identical-expression rule applies. The + window-composition rule adds 48 of the doc's 54 candidates, and its + tumbling and sliding variants wait on #511. +4. **"Shared summary charged once"** does not apply literally: the doc says no + summary is shared in Example 1. The tests check the general form (each node + charged once) on the shared input node. +5. **Workload DAG with two roots.** `LogicalASAPDAG` has one `root`, but + Example 1 needs one DAG for both queries. The stubs carry `query_roots` + until the export supports several roots. +6. **Where cost is computed.** The viewer contract puts `cost` on Stage 2 + candidates, but the doc says only Stage 3 uses the cost model. The spec has + Stage 3 produce costs, and the viewer document attaches them to Stage 2 + entries. +7. **Hydra's inner sketch.** The doc says "Hydra over the whole `job` + column" without naming the inner sketch. The spec assumes `HydraCms`, the + construction with the proven guarantee for frequencies. +8. **Other top-*k* families.** The code also knows `CountSketchWithHeap`. The + doc lists only Count-Min + heap and Hydra, so a third family fails the + tests until the doc adds it. +9. **Latency.** The doc says an exact top 10 rebuilt at every refresh "may + miss" 100 ms. Whether P1 and P2 are rejected is left to the cost model. +10. **Accuracy of Count-Min heap merges** matters only with tumbling windows, + which are out of MVP scope. diff --git a/tools/dag-viewer/examples/planner-layering-example1.expected.json b/tools/dag-viewer/examples/planner-layering-example1.expected.json new file mode 100644 index 00000000..c655fc33 --- /dev/null +++ b/tools/dag-viewer/examples/planner-layering-example1.expected.json @@ -0,0 +1,1603 @@ +{ + "format": "asap-stage-pipeline/v1", + "expected_only": true, + "_notes": [ + "Shape reference for planner-layering Example 1 (MVP): see docs/design_docs/proposals/planner-layering-example1-acceptance.md.", + "Nodes carry only payload kinds and a label, not full schemas. Real documents have the full LogicalASAPDAG/PhysicalASAPDAG nodes.", + "query_roots stands in for root: one workload DAG serves both queries, and LogicalASAPDAG has a single root.", + "Costs are null: the cost model decides them. The selection is illustrative (the doc's 'Raw with a shared input' winner); tests only require the cheapest valid candidate." + ], + "workload": { + "queries": [ + { + "id": "q1", + "language": "promql", + "text": "sum by (job) (rate(http_requests_total[1m]))", + "repeat_every_s": 10, + "lookback_s": 60, + "accuracy": "exact" + }, + { + "id": "q2", + "language": "promql", + "text": "topk by (job) (10, sum_over_time(http_requests_total[1m]))", + "repeat_every_s": 10, + "lookback_s": 60, + "accuracy": { + "epsilon": 0.01, + "delta": 0.001 + }, + "max_latency_ms": 100 + } + ] + }, + "stage0_logical": { + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)" + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)" + }, + { + "id": 4, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 5, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 6, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum_over_time (per series)" + }, + { + "id": 7, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "topk by (job) (10)" + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 4, + "consumer": 5 + }, + { + "producer": 5, + "consumer": 6 + }, + { + "producer": 6, + "consumer": 7 + } + ], + "query_roots": { + "q1": 3, + "q2": 7 + } + } + }, + "stage1_logical_asap": { + "candidates": [ + { + "id": "L1", + "label": "Q1 exact · Q2 exact · separate input", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)" + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)" + }, + { + "id": 4, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 5, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 6, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum_over_time (per series)" + }, + { + "id": 7, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "topk by (job) (10)" + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 4, + "consumer": 5 + }, + { + "producer": 5, + "consumer": 6 + }, + { + "producer": 6, + "consumer": 7 + } + ], + "query_roots": { + "q1": 3, + "q2": 7 + } + } + }, + { + "id": "L2", + "label": "Q1 exact · Q2 exact · shared input", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)" + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)" + }, + { + "id": 4, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum_over_time (per series)" + }, + { + "id": 5, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "topk by (job) (10)" + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 1, + "consumer": 4 + }, + { + "producer": 4, + "consumer": 5 + } + ], + "query_roots": { + "q1": 3, + "q2": 5 + } + } + }, + { + "id": "L3", + "label": "Q1 exact · Q2 Count-Min + heap · separate input", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)" + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)" + }, + { + "id": 4, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 5, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 6, + "payload": { + "kind": "summary_agg" + }, + "label": "Count-Min + top-k heap, one per job" + }, + { + "id": 7, + "payload": { + "kind": "summary_estimate" + }, + "label": "top 10 per job" + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 4, + "consumer": 5 + }, + { + "producer": 5, + "consumer": 6 + }, + { + "producer": 6, + "consumer": 7 + } + ], + "query_roots": { + "q1": 3, + "q2": 7 + } + } + }, + { + "id": "L4", + "label": "Q1 exact · Q2 Count-Min + heap · shared input", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)" + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)" + }, + { + "id": 4, + "payload": { + "kind": "summary_agg" + }, + "label": "Count-Min + top-k heap, one per job" + }, + { + "id": 5, + "payload": { + "kind": "summary_estimate" + }, + "label": "top 10 per job" + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 1, + "consumer": 4 + }, + { + "producer": 4, + "consumer": 5 + } + ], + "query_roots": { + "q1": 3, + "q2": 5 + } + } + }, + { + "id": "L5", + "label": "Q1 exact · Q2 Hydra · separate input", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)" + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)" + }, + { + "id": 4, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 5, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 6, + "payload": { + "kind": "summary_agg" + }, + "label": "Hydra (CMS) over all (job, series) keys + heap per job" + }, + { + "id": 7, + "payload": { + "kind": "summary_estimate" + }, + "label": "top 10 per job" + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 4, + "consumer": 5 + }, + { + "producer": 5, + "consumer": 6 + }, + { + "producer": 6, + "consumer": 7 + } + ], + "query_roots": { + "q1": 3, + "q2": 7 + } + } + }, + { + "id": "L6", + "label": "Q1 exact · Q2 Hydra · shared input", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total" + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m" + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)" + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)" + }, + { + "id": 4, + "payload": { + "kind": "summary_agg" + }, + "label": "Hydra (CMS) over all (job, series) keys + heap per job" + }, + { + "id": 5, + "payload": { + "kind": "summary_estimate" + }, + "label": "top 10 per job" + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 1, + "consumer": 4 + }, + { + "producer": 4, + "consumer": 5 + } + ], + "query_roots": { + "q1": 3, + "q2": 5 + } + } + } + ] + }, + "stage2_physical_asap": { + "candidates": [ + { + "id": "P1", + "from_logical": "L1", + "label": "Q1 exact · Q2 exact · separate input · sort + limit · query time", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 4, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 5, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 6, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum_over_time (per series)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 7, + "payload": { + "kind": "relational", + "operator": { + "kind": "sort" + } + }, + "label": "sort by value, partition by job", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 8, + "payload": { + "kind": "relational", + "operator": { + "kind": "limit" + } + }, + "label": "limit 10", + "output_state": { + "timing": "QueryTime" + } + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 4, + "consumer": 5 + }, + { + "producer": 5, + "consumer": 6 + }, + { + "producer": 6, + "consumer": 7 + }, + { + "producer": 7, + "consumer": 8 + } + ], + "query_roots": { + "q1": 3, + "q2": 8 + } + }, + "cost": { + "total": null, + "unit": "cpu_ms_per_s", + "per_node": {} + } + }, + { + "id": "P2", + "from_logical": "L2", + "label": "Q1 exact · Q2 exact · shared input · sort + limit · query time", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 4, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum_over_time (per series)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 5, + "payload": { + "kind": "relational", + "operator": { + "kind": "sort" + } + }, + "label": "sort by value, partition by job", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 6, + "payload": { + "kind": "relational", + "operator": { + "kind": "limit" + } + }, + "label": "limit 10", + "output_state": { + "timing": "QueryTime" + } + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 1, + "consumer": 4 + }, + { + "producer": 4, + "consumer": 5 + }, + { + "producer": 5, + "consumer": 6 + } + ], + "query_roots": { + "q1": 3, + "q2": 6 + } + }, + "cost": { + "total": null, + "unit": "cpu_ms_per_s", + "per_node": {} + } + }, + { + "id": "P3", + "from_logical": "L3", + "label": "Q1 exact · Q2 Count-Min + heap · separate input · build + estimate · query time", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 4, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 5, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 6, + "payload": { + "kind": "summary_agg" + }, + "label": "Count-Min + top-k heap, one per job", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 7, + "payload": { + "kind": "summary_estimate" + }, + "label": "top 10 per job", + "output_state": { + "timing": "QueryTime" + } + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 4, + "consumer": 5 + }, + { + "producer": 5, + "consumer": 6 + }, + { + "producer": 6, + "consumer": 7 + } + ], + "query_roots": { + "q1": 3, + "q2": 7 + } + }, + "cost": { + "total": null, + "unit": "cpu_ms_per_s", + "per_node": {} + } + }, + { + "id": "P4", + "from_logical": "L4", + "label": "Q1 exact · Q2 Count-Min + heap · shared input · build + estimate · query time", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 4, + "payload": { + "kind": "summary_agg" + }, + "label": "Count-Min + top-k heap, one per job", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 5, + "payload": { + "kind": "summary_estimate" + }, + "label": "top 10 per job", + "output_state": { + "timing": "QueryTime" + } + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 1, + "consumer": 4 + }, + { + "producer": 4, + "consumer": 5 + } + ], + "query_roots": { + "q1": 3, + "q2": 5 + } + }, + "cost": { + "total": null, + "unit": "cpu_ms_per_s", + "per_node": {} + } + }, + { + "id": "P5", + "from_logical": "L5", + "label": "Q1 exact · Q2 Hydra · separate input · build + estimate · query time", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 4, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 5, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 6, + "payload": { + "kind": "summary_agg" + }, + "label": "Hydra (CMS) over all (job, series) keys + heap per job", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 7, + "payload": { + "kind": "summary_estimate" + }, + "label": "top 10 per job", + "output_state": { + "timing": "QueryTime" + } + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 4, + "consumer": 5 + }, + { + "producer": 5, + "consumer": 6 + }, + { + "producer": 6, + "consumer": 7 + } + ], + "query_roots": { + "q1": 3, + "q2": 7 + } + }, + "cost": { + "total": null, + "unit": "cpu_ms_per_s", + "per_node": {} + } + }, + { + "id": "P6", + "from_logical": "L6", + "label": "Q1 exact · Q2 Hydra · shared input · build + estimate · query time", + "dag": { + "nodes": [ + { + "id": 0, + "payload": { + "kind": "relational", + "operator": { + "kind": "scan" + } + }, + "label": "http_requests_total", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 1, + "payload": { + "kind": "relational", + "operator": { + "kind": "time_range" + } + }, + "label": "range 1m", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 2, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "rate (per series)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 3, + "payload": { + "kind": "relational", + "operator": { + "kind": "aggregate" + } + }, + "label": "sum by (job)", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 4, + "payload": { + "kind": "summary_agg" + }, + "label": "Hydra (CMS) over all (job, series) keys + heap per job", + "output_state": { + "timing": "QueryTime" + } + }, + { + "id": 5, + "payload": { + "kind": "summary_estimate" + }, + "label": "top 10 per job", + "output_state": { + "timing": "QueryTime" + } + } + ], + "edges": [ + { + "producer": 0, + "consumer": 1 + }, + { + "producer": 1, + "consumer": 2 + }, + { + "producer": 2, + "consumer": 3 + }, + { + "producer": 1, + "consumer": 4 + }, + { + "producer": 4, + "consumer": 5 + } + ], + "query_roots": { + "q1": 3, + "q2": 5 + } + }, + "cost": { + "total": null, + "unit": "cpu_ms_per_s", + "per_node": {} + } + } + ] + }, + "stage3_selection": { + "selected": "P6", + "rejected": [ + { + "id": "P1", + "reason": "valid; costlier than P6 (illustrative)" + }, + { + "id": "P2", + "reason": "valid; costlier than P6 (illustrative)" + }, + { + "id": "P3", + "reason": "valid; costlier than P6 (illustrative)" + }, + { + "id": "P4", + "reason": "valid; costlier than P6 (illustrative)" + }, + { + "id": "P5", + "reason": "valid; costlier than P6 (illustrative)" + } + ] + } +} From f620555b5ee47fcb32293d6487980b8c3e7d58ed Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 21:00:36 +0000 Subject: [PATCH 2/2] test: run the Example 1 acceptance tests against the real stages Replace the todo!() stubs with adapters over Stage 0-3, un-ignore the 11 tests that pass, and give each still-ignored test the precise difference between the implementation and the spec. Add runtime capability checks: the selected plan compiles in the physical planner; compiling every candidate stays ignored with the two runtime gaps it finds. Co-Authored-By: Claude Opus 5.5 --- .../tests/planner_layering_example1.rs | 361 +++++++++++++++--- 1 file changed, 308 insertions(+), 53 deletions(-) diff --git a/crates/integration-tests/tests/planner_layering_example1.rs b/crates/integration-tests/tests/planner_layering_example1.rs index ad2c1a0b..e497caa7 100644 --- a/crates/integration-tests/tests/planner_layering_example1.rs +++ b/crates/integration-tests/tests/planner_layering_example1.rs @@ -1,9 +1,10 @@ //! Acceptance tests for #509 "Example 1: Aggregation over dimensions", MVP scope. //! //! Spec: `docs/design_docs/proposals/planner-layering-example1-acceptance.md`. -//! Written by the test designer before the Phase C stage APIs exist: every -//! stage test is `#[ignore]`d and calls the `todo!()` stubs in [`stages`], -//! which the implementer replaces with the real entry points. +//! Written by the test designer before the Phase C stage APIs existed; the +//! implementer replaced the stubs in [`stages`] with adapters over the real +//! stages. Tests that still fail because the implementation differs from the +//! spec stay `#[ignore]`d with the difference as the reason. //! //! MVP scope: Stage 1 = Pass 1 + the identical-expression rule only (no //! window-composition variants); Stage 2 = physical operator implementation @@ -24,27 +25,32 @@ use asap_types::workload::{ RepetitionInterval, TimeSelection, }; -/// Hypothetical Phase C stage API. Replace each `todo!()` with the real call. -#[allow(dead_code, unused_variables)] +/// Adapters from the Phase C stage APIs to the shapes these tests were written +/// against. They only convert; every decision is the real stage's. +#[allow(dead_code)] mod stages { - use super::*; + use std::rc::Rc; - /// One whole-workload candidate. `LogicalASAPDAG::root` names one query, - /// so `query_roots` carries one root per workload entry, in - /// `QueryWorkload::entries()` order (`[q1, q2]`), until the export - /// supports several roots. + use super::*; + use asap_aware_mapping::logical_candidates::{ + compose_logical_candidate, enumerate_local_logical_candidates, + }; + use asap_types::ir::export::{compile_logical_asap_workload, LogicalASAPQueryRoot}; + use asap_types::ir::{OperatorNode, QueryRoot}; + + /// One whole-workload candidate. `query_roots` holds one root per + /// workload entry, in `QueryWorkload::entries()` order (`[q1, q2]`). + /// `roots` are the in-memory roots the next stage consumes. #[derive(Debug, Clone)] pub struct LogicalCandidate { pub id: String, pub label: String, pub dag: LogicalASAPDAG, pub query_roots: Vec, + pub roots: Vec>, } - /// Placeholder for `asap_types::ir::export::PhysicalASAPDAG` (phase A). - /// Its payload is an alias of the logical payload, so payload-kind checks - /// carry over unchanged. - pub type PhysicalASAPDAG = LogicalASAPDAG; + pub type PhysicalASAPDAG = asap_types::ir::export::PhysicalASAPDAG; /// One Stage 2 candidate, derived from exactly one Stage 1 candidate. #[derive(Debug, Clone)] @@ -54,6 +60,7 @@ mod stages { pub label: String, pub dag: PhysicalASAPDAG, pub query_roots: Vec, + pub stage2: asap_aware_mapping::physical_candidates::PhysicalCandidate, } /// Whole-workload cost of one physical candidate; `per_node` has one @@ -73,8 +80,7 @@ mod stages { pub reason: String, } - /// Stage 3 output. Costs are keyed by physical candidate id; the viewer - /// document attaches them to the Stage 2 entries. + /// Stage 3 output. Costs are keyed by physical candidate id. #[derive(Debug, Clone)] pub struct Selection { pub selected: String, @@ -82,27 +88,115 @@ mod stages { pub rejected: Vec, } + /// The frontend DAG with each PromQL series' full identity as a column, + /// the row representation per-series state needs at runtime. + fn lower(workload: &PlanningWorkload) -> Vec { + asap_frontend_promql::lower_promql_query_workload(workload, 0) + .expect("Example 1 lowers") + .into_iter() + .map(|root| match root { + QueryRoot::Operator(node) => QueryRoot::Operator( + asap_types::ir::schema_support::with_promql_series_identity(&node) + .expect("series identity"), + ), + QueryRoot::Scalar(_) => panic!("Example 1 has operator roots"), + }) + .collect() + } + + fn candidate(id: String, label: String, roots: Vec) -> LogicalCandidate { + let dag = compile_logical_asap_workload(&roots).expect("logical export"); + let query_roots = dag + .roots + .iter() + .map(|root| match root { + LogicalASAPQueryRoot::Operator(id) => *id, + LogicalASAPQueryRoot::Scalar(_) => panic!("Example 1 has operator roots"), + }) + .collect(); + let roots = roots + .into_iter() + .map(|root| match root { + QueryRoot::Operator(node) => node, + QueryRoot::Scalar(_) => panic!("Example 1 has operator roots"), + }) + .collect(); + LogicalCandidate { + id, + label, + dag, + query_roots, + roots, + } + } + /// Stage 0: frontends lower every query into one summary-free workload DAG. pub fn stage0_logical(workload: &PlanningWorkload) -> LogicalCandidate { - todo!("Phase C: Stage 0 workload LogicalDAG") + candidate("S0".into(), "frontend".into(), lower(workload)) } - /// Stage 1: Pass 1 local alternatives × Pass 2 identical-expression - /// sharing. Takes the workload for accuracy requirements and time selection. + /// Stage 1: every combination of Pass 1 local alternatives (Pass 2 is + /// not implemented). Lowers `workload` again: Pass 1 reads the in-memory + /// DAG, not the Stage 0 export. pub fn stage1_logical_asap( workload: &PlanningWorkload, - logical: &LogicalCandidate, + _logical: &LogicalCandidate, ) -> Vec { - todo!("Phase C: Stage 1 CandidateLogicalASAPDAGs") + let inventory = + enumerate_local_logical_candidates(lower(workload).into_iter().enumerate().collect()) + .expect("Pass 1"); + let mut choices = vec![vec![]]; + for target in &inventory.targets { + choices = choices + .into_iter() + .flat_map(|prefix: Vec| { + (0..target.alternatives.len()).map(move |i| { + let mut choice = prefix.clone(); + choice.push(i); + choice + }) + }) + .collect(); + } + choices + .into_iter() + .enumerate() + .map(|(index, choice)| { + let roots = compose_logical_candidate(&inventory, &choice) + .expect("composes") + .into_iter() + .map(|(_, root)| root) + .collect(); + candidate(format!("L{}", index + 1), format!("{choice:?}"), roots) + }) + .collect() } /// Stage 2: physical operator implementation of every logical candidate /// (no materialization in the MVP). pub fn stage2_physical( - workload: &PlanningWorkload, + _workload: &PlanningWorkload, logical: &[LogicalCandidate], ) -> Vec { - todo!("Phase C: Stage 2 CandidatePhysicalASAPDAGs") + logical + .iter() + .enumerate() + .map(|(index, l)| { + let mut stage2 = + asap_aware_mapping::physical_candidates::stage2_physical(&l.id, &l.roots) + .unwrap_or_else(|e| panic!("{}: {e}", l.id)); + stage2.id = format!("P{}", index + 1); + stage2.label = l.label.clone(); + PhysicalCandidate { + id: stage2.id.clone(), + from_logical: stage2.from_logical.clone(), + label: stage2.label.clone(), + dag: stage2.dag.clone(), + query_roots: stage2.dag.roots.clone(), + stage2, + } + }) + .collect() } /// Stage 3: reject invalid candidates, cost the rest for the whole @@ -112,13 +206,61 @@ mod stages { physical: &[PhysicalCandidate], models: PlanningModels<'_>, ) -> Selection { - todo!("Phase C: Stage 3 plan selection") + let targets: Vec<_> = workload + .query_workload + .entries() + .map(|entry| Some(entry.requirements.accuracy.target())) + .collect(); + let candidates: Vec<_> = physical.iter().map(|p| p.stage2.clone()).collect(); + let selection = asap_aware_mapping::plan_selection::stage3_select( + &candidates, + &targets, + workload.data_workload.as_ref().expect("data workload"), + models, + ) + .expect("Stage 3 selects"); + Selection { + selected: selection.selected, + costs: selection + .costs + .into_iter() + .map(|(id, cost)| { + let per_node = cost + .per_node + .into_iter() + .map(|(node, c)| (node, c.cost)) + .collect(); + ( + id, + CandidateCost { + total: cost.total, + per_node, + }, + ) + }) + .collect(), + rejected: selection + .rejected + .into_iter() + .map(|r| Rejection { + id: r.id, + valid: r.valid, + reason: r.reason, + }) + .collect(), + } } - /// Stands in for `node.output_state.timing == IngestionTime` until - /// `PhysicalASAPDAG` lands. pub fn runs_at_ingestion(candidate: &PhysicalCandidate, node: LogicalASAPNodeId) -> bool { - todo!("Phase C: read PhysicalASAPDAGNode::output_state.timing") + candidate + .dag + .nodes + .iter() + .find(|n| n.id == node) + .expect("node") + .output_state + .timing + == asap_types::post_asap::ExecutionTiming::IngestionTime } } @@ -198,22 +340,59 @@ enum Q2Option { Hydra, } +/// The logical and physical exports share node ids, payloads and edge +/// endpoints; the helpers below read only those. +trait ExportedDag { + fn producers(&self, consumer: LogicalASAPNodeId) -> Vec; + fn node_payload(&self, id: LogicalASAPNodeId) -> &LogicalASAPOperatorPayload; +} + +impl ExportedDag for LogicalASAPDAG { + fn producers(&self, consumer: LogicalASAPNodeId) -> Vec { + let edges = self.edges.iter().filter(|e| e.consumer == consumer); + edges.map(|e| e.producer).collect() + } + fn node_payload(&self, id: LogicalASAPNodeId) -> &LogicalASAPOperatorPayload { + &self + .nodes + .iter() + .find(|n| n.id == id) + .expect("node") + .payload + } +} + +impl ExportedDag for PhysicalASAPDAG { + fn producers(&self, consumer: LogicalASAPNodeId) -> Vec { + let edges = self.edges.iter().filter(|e| e.consumer == consumer); + edges.map(|e| e.producer).collect() + } + fn node_payload(&self, id: LogicalASAPNodeId) -> &LogicalASAPOperatorPayload { + &self + .nodes + .iter() + .find(|n| n.id == id) + .expect("node") + .payload + } +} + /// Every node `root` depends on, including itself. -fn closure(dag: &LogicalASAPDAG, root: LogicalASAPNodeId) -> HashSet { +fn closure(dag: &impl ExportedDag, root: LogicalASAPNodeId) -> HashSet { let mut seen = HashSet::from([root]); let mut stack = vec![root]; while let Some(node) = stack.pop() { - for edge in dag.edges.iter().filter(|e| e.consumer == node) { - if seen.insert(edge.producer) { - stack.push(edge.producer); + for producer in dag.producers(node) { + if seen.insert(producer) { + stack.push(producer); } } } seen } -fn payload(dag: &LogicalASAPDAG, id: LogicalASAPNodeId) -> &LogicalASAPOperatorPayload { - &dag.nodes.iter().find(|n| n.id == id).expect("node").payload +fn payload(dag: &impl ExportedDag, id: LogicalASAPNodeId) -> &LogicalASAPOperatorPayload { + dag.node_payload(id) } fn is_summary(payload: &LogicalASAPOperatorPayload) -> bool { @@ -228,7 +407,10 @@ fn is_summary(payload: &LogicalASAPOperatorPayload) -> bool { } /// Sketch families built in `nodes`, as Example 1's Q2 options. -fn sketch_options(dag: &LogicalASAPDAG, nodes: &HashSet) -> BTreeSet { +fn sketch_options( + dag: &impl ExportedDag, + nodes: &HashSet, +) -> BTreeSet { nodes .iter() .filter_map(|&id| match payload(dag, id) { @@ -258,7 +440,7 @@ fn roots(query_roots: &[LogicalASAPNodeId]) -> (LogicalASAPNodeId, LogicalASAPNo } /// Q2's option and whether Q1 and Q2 share any node, for one candidate. -fn classify(dag: &LogicalASAPDAG, query_roots: &[LogicalASAPNodeId]) -> (Q2Option, bool) { +fn classify(dag: &impl ExportedDag, query_roots: &[LogicalASAPNodeId]) -> (Q2Option, bool) { let (q1, q2) = roots(query_roots); let (c1, c2) = (closure(dag, q1), closure(dag, q2)); let options = sketch_options(dag, &c2); @@ -320,7 +502,7 @@ fn workload_encodes_example1() { /// Today's PromQL frontend already lowers each query to the doc's Stage 0 chain. #[test] fn stage0_frontend_lowers_each_query_to_doc_chain() { - let roots = asap_frontend_promql::unified::lower_promql_query_workload(&example1_workload(), 0) + let roots = asap_frontend_promql::lower_promql_query_workload(&example1_workload(), 0) .expect("Example 1 lowers"); let chains: Vec> = roots .iter() @@ -354,7 +536,6 @@ fn stage0_frontend_lowers_each_query_to_doc_chain() { /// Stage 0 yields one workload DAG with one root per query and no summaries. #[test] -#[ignore = "pending Phase C stage APIs"] fn stage0_one_summary_free_workload_dag() { let stage0 = stage0_logical(&example1_workload()); stage0.dag.validate().expect("valid DAG"); @@ -364,7 +545,6 @@ fn stage0_one_summary_free_workload_dag() { /// Stage 0 keeps Q1 and Q2 separate; sharing is a Stage 1 decision. #[test] -#[ignore = "pending Phase C stage APIs"] fn stage0_queries_do_not_share_nodes() { let stage0 = stage0_logical(&example1_workload()); let (q1, q2) = roots(&stage0.query_roots); @@ -375,7 +555,7 @@ fn stage0_queries_do_not_share_nodes() { /// Stage 1 outputs exactly the 3 Q2 options × {separate, shared input} = 6 candidates. #[test] -#[ignore = "pending Phase C stage APIs"] +#[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() { let (_, logical, _) = pipeline(); assert_eq!(logical.len(), 6); @@ -392,7 +572,6 @@ fn stage1_has_six_candidates_covering_every_combination() { /// Q1 is exact in every Stage 1 candidate: no summary is reachable from its root. #[test] -#[ignore = "pending Phase C stage APIs"] fn stage1_q1_is_always_exact() { let (_, logical, _) = pipeline(); for c in &logical { @@ -408,7 +587,7 @@ fn stage1_q1_is_always_exact() { /// Q2's summary families are exactly Count-Min + heap per job and Hydra over all jobs. #[test] -#[ignore = "pending Phase C stage APIs"] +#[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() { let (_, logical, _) = pipeline(); let families: BTreeSet<_> = logical @@ -426,7 +605,7 @@ fn stage1_q2_summary_families_are_count_min_heap_and_hydra() { /// Sharing adds a variant and keeps the independent one, for every Q2 option. #[test] -#[ignore = "pending Phase C stage APIs"] +#[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"] fn stage1_keeps_independent_and_shared_variants() { let (_, logical, _) = pipeline(); let found: Vec<_> = logical @@ -448,7 +627,6 @@ fn stage1_keeps_independent_and_shared_variants() { /// Only the raw input is shared between Q1 and Q2; no summary is shared. #[test] -#[ignore = "pending Phase C stage APIs"] fn stage1_shares_input_but_never_a_summary() { let (_, logical, _) = pipeline(); for c in &logical { @@ -467,7 +645,6 @@ fn stage1_shares_input_but_never_a_summary() { /// Stage 1 candidates are valid DAGs with unique ids. #[test] -#[ignore = "pending Phase C stage APIs"] fn stage1_candidates_are_valid_and_uniquely_named() { let (_, logical, _) = pipeline(); let ids: BTreeSet<_> = logical.iter().map(|c| c.id.as_str()).collect(); @@ -481,7 +658,7 @@ fn stage1_candidates_are_valid_and_uniquely_named() { /// No candidate is discarded before Stage 3: Stage 2 maps the 6 logical candidates one-to-one. #[test] -#[ignore = "pending Phase C stage APIs"] +#[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); @@ -494,7 +671,7 @@ fn stage2_keeps_every_logical_candidate() { /// Stage 2 preserves each logical candidate's Q2 option and input sharing. #[test] -#[ignore = "pending Phase C stage APIs"] +#[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 { @@ -511,7 +688,7 @@ fn stage2_preserves_logical_choices() { /// Exact TopK is implemented as a sort followed by a limit. #[test] -#[ignore = "pending Phase C stage APIs"] +#[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 @@ -530,7 +707,7 @@ 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 = "pending Phase C stage APIs"] +#[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 { @@ -563,7 +740,6 @@ fn stage2_summary_topk_is_build_then_estimate() { /// With no materialization in the MVP, every node runs at query time. #[test] -#[ignore = "pending Phase C stage APIs"] fn stage2_everything_runs_at_query_time() { let (_, _, physical) = pipeline(); for p in &physical { @@ -578,11 +754,91 @@ fn stage2_everything_runs_at_query_time() { } } +/// Compile `p` in the physical planner (the runtime capability check). Inputs +/// are what a deployment supplies: raw series for each sub-DAG the planner +/// runs as a retained PromQL expression, and the samples of each time range a +/// native operator reads. +fn compile_in_runtime(p: &PhysicalCandidate) -> Result<(), String> { + use asap_physical_operators::physical_planner::{compile, promql_fallback, InputContract}; + use asap_types::ir::export::compile_physical_asap_workload_with_node_ids; + use std::sync::Arc; + let ids = compile_physical_asap_workload_with_node_ids(&p.stage2.roots) + .expect("re-export") + .node_ids; + let mut inputs = BTreeMap::new(); + let mut pending = p.dag.roots.clone(); + let mut seen = HashSet::new(); + while let Some(id) = pending.pop() { + if !seen.insert(id) { + continue; + } + let node = ids.operator_node(id).expect("node"); + let time_range = relational(payload(&p.dag, id)).as_deref() == Some("time_range"); + // Raw samples a summary reads are an input, not a retained expression. + let summary_input = time_range + && p.dag.edges.iter().any(|e| { + e.producer == id + && matches!( + payload(&p.dag, e.consumer), + LogicalASAPOperatorPayload::SummaryAgg { .. } + ) + }); + let fallback = (!node.contains_asap() && !summary_input) + .then(|| promql_fallback::raw_series(node).ok()) + .flatten(); + if let Some(selectors) = fallback { + for (i, (_, schema)) in selectors.into_iter().enumerate() { + let slot = promql_fallback::raw_series_input(u64::from(id.0), i); + inputs.insert(slot, InputContract::bounded(schema)); + } + } else if time_range { + let schema = p.dag.nodes.iter().find(|n| n.id == id).unwrap(); + let schema = Arc::new(schema.output_schema.clone()); + inputs.insert(u64::from(id.0), InputContract::bounded(schema)); + } else { + pending.extend(p.dag.producers(id)); + } + } + let roots: Vec = p.dag.roots.iter().map(|id| u64::from(id.0)).collect(); + compile(&p.dag, inputs, &roots) + .map(|_| ()) + .map_err(|e| format!("{} ({}): {e}", p.id, p.label)) +} + +/// Runtime capability check (added by the implementer, not part of the +/// spec): the physical planner compiles every Stage 2 candidate. +#[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() { + let (_, _, physical) = pipeline(); + let failures: Vec<_> = physical + .iter() + .filter_map(|p| compile_in_runtime(p).err()) + .collect(); + assert!(failures.is_empty(), "{failures:#?}"); +} + +/// The plan Stage 3 selects compiles in the physical planner (added by the +/// implementer, not part of the spec). +#[test] +fn stage3_selected_plan_compiles_in_the_physical_planner() { + let (workload, _, physical) = pipeline(); + let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); + let selected = physical + .iter() + .find(|p| p.id == selection.selected) + .unwrap(); + compile_in_runtime(selected).unwrap(); +} + // ── Stage 3 ────────────────────────────────────────────────────────────── /// Stage 3 selects one candidate and gives every other one a reason. #[test] -#[ignore = "pending Phase C stage APIs"] fn stage3_selects_one_and_explains_the_rest() { let (workload, _, physical) = pipeline(); let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); @@ -603,7 +859,6 @@ fn stage3_selects_one_and_explains_the_rest() { /// The selected plan is the cheapest valid candidate for the whole workload. #[test] -#[ignore = "pending Phase C stage APIs"] fn stage3_selects_cheapest_valid() { let (workload, _, physical) = pipeline(); let selection = stage3_select(&workload, &physical, PlanningModels::builtin()); @@ -621,7 +876,7 @@ fn stage3_selects_cheapest_valid() { /// Every node is charged exactly once, so a shared input is costed once for both queries. #[test] -#[ignore = "pending Phase C stage APIs"] +#[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()); @@ -645,7 +900,7 @@ fn stage3_charges_each_node_once() { /// Sharing the input never costs more than reading it separately. #[test] -#[ignore = "pending Phase C stage APIs"] +#[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"] fn stage3_shared_input_is_not_costlier() { let (workload, _, physical) = pipeline(); let selection = stage3_select(&workload, &physical, PlanningModels::builtin());