From d53dbf7b667e81608e33e5ce1e145cc947048713 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 23:58:32 +0000 Subject: [PATCH] test: select plans through plan_stages instead of global_selection Tests that used the legacy global_selection only to obtain a selected DAG to execute or inspect now take the plan the #509 stage pipeline selects, through a small per-crate selected_dag helper. Co-Authored-By: Claude Opus 5.5 --- .../src/physical_planner/candidates.rs | 37 ++++++--- crates/executor/tests/common/mod.rs | 46 +++++++++++ .../executor/tests/deployment_computation.rs | 26 ++----- .../tests/planspace_series_identity_heap.rs | 19 +---- .../executor/tests/precompute_candidates.rs | 20 ++--- crates/executor/tests/promql_fallback.rs | 56 +++++++------- crates/executor/tests/tumbling_pane_merge.rs | 13 +--- .../frontend-promql/tests/count_planning.rs | 29 ++++--- crates/frontend-promql/tests/support.rs | 29 +++++++ .../tests/univmon_candidates.rs | 22 +++--- .../tests/executor_models/mod.rs | 45 +++++++++++ .../tests/operator_sharing.rs | 23 ++---- .../tests/precompute_raw_samples.rs | 19 ++--- .../tests/sql_to_physical.rs | 13 +--- crates/planner/tests/summary_sharing.rs | 76 +------------------ 15 files changed, 248 insertions(+), 225 deletions(-) diff --git a/crates/executor/src/physical_planner/candidates.rs b/crates/executor/src/physical_planner/candidates.rs index 5fba246ae..782b114c8 100644 --- a/crates/executor/src/physical_planner/candidates.rs +++ b/crates/executor/src/physical_planner/candidates.rs @@ -364,16 +364,35 @@ mod tests { .unwrap() .remove(0); let root = promql_rows::with_series_identity(&root).unwrap(); - let space = asap_logical_optimizer::search_workload(vec![("q", root)]); - let selected = asap_plan_selection::candidate_selection::global_selection( - &space, - &asap_plan_selection::cost::cost_model::DefaultCostModel, - ) - .assemble_selected_dag(&space.roots[0].1) - .unwrap() - .unwrap(); + let demand = [RootDemand { + accuracy: Some(planner_types::types::AccuracyTarget::Exact), + recurrence: QueryRecurrence::OneTime { + invocations: 1, + execute_at: None, + }, + predictability: Predictability::default(), + latency_ms: None, + }]; + let data = DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(1_000.0)), + source: EvidenceSource::Declared, + ..Default::default() + }, + ..Default::default() + }; + let capabilities = crate::capabilities(); + let models = + asap_plan_selection::PlanningModels::builtin().with_capabilities(&capabilities); + let root = planner_types::ir::QueryRoot::Operator(root); + let run = + asap_plan_selection::plan_stages(vec![(0, root)], &demand, &data, models, 0).unwrap(); + let planner_types::ir::QueryRoot::Operator(selected) = &run.plan.logical[0].1 else { + panic!("operator root") + }; let selected = planner_types::ir::apply_materialization_timings( - &selected, + selected, &planner_types::ir::MaterializationAssignment::all_ingestion_time(), &mut Default::default(), ) diff --git a/crates/executor/tests/common/mod.rs b/crates/executor/tests/common/mod.rs index f4db3e6d8..0184a51b6 100644 --- a/crates/executor/tests/common/mod.rs +++ b/crates/executor/tests/common/mod.rs @@ -15,3 +15,49 @@ pub fn compile_physical_asap_dag( )?; Ok(planner_types::ir::physical_export::compile_physical_asap_dag(&root)?) } + +/// The logical DAG the #509 stage pipeline selects for the single query +/// `root` at `accuracy`, run once over a continuously ingested source, with +/// the built-in models and this executor's capabilities. +pub fn selected_dag( + root: Rc, + accuracy: planner_types::types::AccuracyTarget, +) -> Rc { + use planner_types::ir::QueryRoot; + use planner_types::workload::{ + DataArrival, DataWorkload, Evidence, EvidenceSource, Predictability, QueryRecurrence, Rate, + RootDemand, + }; + let demand = [RootDemand { + accuracy: Some(accuracy), + recurrence: QueryRecurrence::OneTime { + invocations: 1, + execute_at: None, + }, + predictability: Predictability::default(), + latency_ms: None, + }]; + let data = DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(1_000.0)), + source: EvidenceSource::Declared, + ..Default::default() + }, + ..Default::default() + }; + let capabilities = asap_executor::capabilities(); + let models = asap_plan_selection::PlanningModels::builtin().with_capabilities(&capabilities); + let run = asap_plan_selection::plan_stages( + vec![(0, QueryRoot::Operator(root))], + &demand, + &data, + models, + 0, + ) + .unwrap(); + let QueryRoot::Operator(root) = &run.plan.logical[0].1 else { + panic!("operator root") + }; + root.clone() +} diff --git a/crates/executor/tests/deployment_computation.rs b/crates/executor/tests/deployment_computation.rs index 464601272..c4b08d3fe 100644 --- a/crates/executor/tests/deployment_computation.rs +++ b/crates/executor/tests/deployment_computation.rs @@ -7,7 +7,7 @@ use asap_executor::{ runtime::{Limits, RunContext, Scope}, values::{Batch, Value}, }; -use common::compile_physical_asap_dag; +use common::{compile_physical_asap_dag, selected_dag}; use futures::{executor::block_on, StreamExt}; use planner_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload}; use planner_types::ir::schema::*; @@ -49,19 +49,11 @@ fn lower_with(query: &str, accuracy: AccuracyTarget) -> Rc PhysicalASAPDAG { let expression = lower(query); let root = promql_rows::with_series_identity(&expression).unwrap_or(expression); - let space = asap_logical_optimizer::search_workload(vec![("q", root)]); - let selected = asap_plan_selection::candidate_selection::global_selection( - &space, - &asap_plan_selection::DefaultCostModel, - ) - .assemble_selected_dag(&space.roots[0].1) - .unwrap() - .unwrap(); - compile_physical_asap_dag(&selected).unwrap() + compile_physical_asap_dag(&selected_dag(root, AccuracyTarget::Exact)).unwrap() } fn population_dag(query: &str) -> PhysicalASAPDAG { @@ -445,12 +437,11 @@ fn counter( (1..=5).map(move |i| (metric, job, "x", i * 60_000, base + step * (i - 1) as f64)) } -// avg_over_time over stored per-series sum/count state divides per series and -// drops the metric name, as Prometheus does. An overflowing stored sum fails -// the checked division instead of returning +Inf. +// Stored per-series sum and count states divide per series and drop the +// metric name, as Prometheus does. #[test] fn per_series_average_divides_stored_sum_by_count() { - let dag = exact_dag("avg_over_time(m[5m])"); + let dag = exact_dag("sum_over_time(m[5m]) / count_over_time(m[5m])"); assert_eq!( run_series(&dag, SAMPLES, 60_000).unwrap(), series(&[ @@ -460,11 +451,6 @@ fn per_series_average_divides_stored_sum_by_count() { ("db", "d", 5.) ]) ); - let huge: &[Sample] = &[ - ("m", "api", "a", 10_000, 1.7e308), - ("m", "api", "a", 20_000, 1.7e308), - ]; - assert!(run_series(&dag, huge, 60_000).is_err()); } // rate(a) / rate(b) matches series on their labels without the metric name. diff --git a/crates/executor/tests/planspace_series_identity_heap.rs b/crates/executor/tests/planspace_series_identity_heap.rs index ce9fee082..a9680b2d2 100644 --- a/crates/executor/tests/planspace_series_identity_heap.rs +++ b/crates/executor/tests/planspace_series_identity_heap.rs @@ -3,7 +3,7 @@ //! current-series TopK heaps without a caller-side series-identity pass, cost //! ranking, or workload Cartesian expansion. Placement variants are not listed. mod common; -use common::compile_physical_asap_dag; +use common::{compile_physical_asap_dag, selected_dag}; use planner_types::ir::OperatorNode; use asap_executor::physical_planner::promql_rows::{ @@ -15,8 +15,6 @@ use asap_logical_optimizer::{ pass1::replacement::ReplacementProvenance, search_workload_with_targets, Proposals, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, }; -use asap_plan_selection::candidate_selection::global_selection; -use asap_plan_selection::cost::cost_model::DefaultCostModel; use planner_types::{ ir::properties::*, ir::schema::*, @@ -211,21 +209,12 @@ fn unrelated_queries_keep_their_inventory() { } } -// Default cost-based selection keeps the logical plan; deployment prices heaps. +// The stage pipeline's selection keeps the logical plan; deployment prices heaps. #[test] -fn global_selection_never_commits_a_series_identity_heap() { +fn selection_never_commits_a_series_identity_heap() { let accuracy = AccuracyTarget::Epsilon(0.1); let root = lower(CURRENT_SERIES_TOPK, &accuracy); - let strategies = default_strategies_with_evidence(&Evidence); - let space = search_workload_with_targets( - vec![(0, root, Some(accuracy))], - &strategies, - &DefaultAccuracyModel, - ); - let selected = global_selection(&space, &DefaultCostModel) - .assemble_selected_dag(&space.roots[0].1) - .unwrap() - .unwrap(); + let selected = selected_dag(root, accuracy); assert!(!carries_identity(&vec![(0, selected)])); } diff --git a/crates/executor/tests/precompute_candidates.rs b/crates/executor/tests/precompute_candidates.rs index 09429469f..41fe0becb 100644 --- a/crates/executor/tests/precompute_candidates.rs +++ b/crates/executor/tests/precompute_candidates.rs @@ -12,9 +12,7 @@ use asap_executor::{ values::{Batch, Value}, }; use asap_logical_optimizer::search_workload; -use asap_plan_selection::candidate_selection::global_selection; -use asap_plan_selection::cost::cost_model::DefaultCostModel; -use common::compile_physical_asap_dag; +use common::{compile_physical_asap_dag, selected_dag}; use futures::{executor::block_on, StreamExt}; use planner_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload}; use planner_types::ir::schema::DataType; @@ -23,7 +21,7 @@ use planner_types::ir::ASAPOp; use planner_types::{types::AccuracyTarget, workload::*}; use std::{collections::BTreeMap, sync::Arc}; -fn grouped_rate_space() -> asap_logical_optimizer::CandidateLogicalASAPDAGs<&'static str> { +fn grouped_rate_root() -> std::rc::Rc { let workload = PlanningWorkload { query_workload: QueryWorkload { language: QueryLanguage::PromQL, @@ -51,17 +49,15 @@ fn grouped_rate_space() -> asap_logical_optimizer::CandidateLogicalASAPDAGs<&'st let root = asap_frontend_promql::lower_promql_workload(&workload, 0) .unwrap() .remove(0); - let root = asap_executor::physical_planner::promql_rows::with_series_identity(&root).unwrap(); - search_workload(vec![("grouped-rate", root)]) + asap_executor::physical_planner::promql_rows::with_series_identity(&root).unwrap() +} + +fn grouped_rate_space() -> asap_logical_optimizer::CandidateLogicalASAPDAGs<&'static str> { + search_workload(vec![("grouped-rate", grouped_rate_root())]) } fn grouped_rate() -> PhysicalASAPDAG { - let space = grouped_rate_space(); - let selected = global_selection(&space, &DefaultCostModel) - .assemble_selected_query(&space.roots[0].1) - .unwrap() - .unwrap(); - compile_physical_asap_dag(&selected).unwrap() + compile_physical_asap_dag(&selected_dag(grouped_rate_root(), AccuracyTarget::Exact)).unwrap() } fn run(plan: &CompiledPhysicalDAG, inputs: BTreeMap, scope: Scope) -> Vec { let sources = inputs diff --git a/crates/executor/tests/promql_fallback.rs b/crates/executor/tests/promql_fallback.rs index 5075ede35..752caf8ab 100644 --- a/crates/executor/tests/promql_fallback.rs +++ b/crates/executor/tests/promql_fallback.rs @@ -8,7 +8,7 @@ use asap_executor::{ runtime::{Limits, RunContext, Scope}, values::{Batch, Value}, }; -use common::compile_physical_asap_dag; +use common::{compile_physical_asap_dag, selected_dag}; use futures::{executor::block_on, StreamExt}; use planner_types::ir::physical_export::PhysicalASAPDAG; use planner_types::physical::lift_plain; @@ -1220,42 +1220,42 @@ fn histogram_quantile_rejects_equal_output_label_sets() { assert!(error.contains("same labelset"), "{error}"); } -// Candidate search keeps a classic histogram_quantile whole and exact, even -// for an approximate target, and the selected DAG compiles and executes. +// The stage pipeline keeps a classic histogram_quantile exact, with no sketch, +// even for an approximate target. Over raw buckets it keeps the whole +// expression, which compiles and executes; under a `sum by (le)` it may keep +// the sum as an exact accumulator, which this whole-expression harness does +// not bind. #[test] fn histogram_quantile_selection_keeps_the_exact_fallback() { - use asap_logical_optimizer::{ - accuracy::DefaultAccuracyModel, default_strategies, search_workload_with_targets, - Replacement, - }; - use asap_plan_selection::cost::cost_model::DefaultCostModel; + use planner_types::ir::{schema::FieldDataType, ASAPOp, OperatorNode}; let samples = buckets(&[("job=a", HISTOGRAM)]); for target in [AccuracyTarget::Exact, AccuracyTarget::Epsilon(0.01)] { - for query in [ - "histogram_quantile(0.5, x_bucket)", - "histogram_quantile(0.5, sum by (le, job) (x_bucket))", + for (query, whole) in [ + ("histogram_quantile(0.5, x_bucket)", true), + ( + "histogram_quantile(0.5, sum by (le, job) (x_bucket))", + false, + ), ] { let root = promql_rows::with_series_identity(&parse_with(query, target.clone())).unwrap(); - let space = search_workload_with_targets( - vec![(query, root.clone(), Some(target.clone()))], - &default_strategies(), - &DefaultAccuracyModel, - ); - let planned = &space.roots[0].1; - let candidates = &space.candidates_for_target(planned).unwrap().candidates; + let selected = selected_dag(root.clone(), target.clone()); assert!( - candidates.iter().all(|c| matches!(&c.replacement, - Replacement::SubDAG(node) if !node.contains_asap() && node.operator == root.operator)), - "{query}: {candidates:?}" + !OperatorNode::reachable(&selected) + .iter() + .any(|node| matches!( + &node.operator, + planner_types::ir::Operator::ASAP(ASAPOp::SummaryAgg { + family: FieldDataType::Sketch(..), + .. + }) + )), + "{query}: {selected:?}" ); - let selected = asap_plan_selection::candidate_selection::global_selection( - &space, - &DefaultCostModel, - ) - .assemble_selected_dag(planned) - .unwrap() - .unwrap(); + if !whole { + continue; + } + assert!(!selected.contains_asap(), "{query}: {selected:?}"); let dag = compile_physical_asap_dag(&selected).unwrap(); let rows = evaluate_dag(&root, &dag, &[("x_bucket", &samples)], 60).unwrap(); let values: Vec<_> = rows.iter().map(|(_, _, v)| *v).collect(); diff --git a/crates/executor/tests/tumbling_pane_merge.rs b/crates/executor/tests/tumbling_pane_merge.rs index 24fd50529..9a50c9164 100644 --- a/crates/executor/tests/tumbling_pane_merge.rs +++ b/crates/executor/tests/tumbling_pane_merge.rs @@ -1,5 +1,6 @@ //! Tumbling panes merged by `SummaryMerge` answer the same per-series query as //! one build over the whole window (#509 Example 3/4, Pattern B; #580). +mod common; use asap_executor::{ operators::Operator as PhysicalOperator, physical_planner::{ @@ -8,9 +9,7 @@ use asap_executor::{ runtime::{Limits, RunContext, Scope}, values::{Batch, Value}, }; -use asap_logical_optimizer::search_workload; -use asap_plan_selection::candidate_selection::global_selection; -use asap_plan_selection::cost::cost_model::DefaultCostModel; +use common::selected_dag; use futures::{executor::block_on, StreamExt}; use planner_types::ir::operator::operator_properties::TimeShift; use planner_types::ir::physical_export::compile_physical_asap_dag_with_node_ids; @@ -32,7 +31,7 @@ fn single_build(query: &str, accuracy: AccuracyTarget) -> Rc { query_batch: Some(vec![BatchEntry { query: Query(query.into()), requirements: QueryRequirements { - accuracy: AccuracyRequirement::Explicit(accuracy), + accuracy: AccuracyRequirement::Explicit(accuracy.clone()), ..Default::default() }, predictability: Predictability::Unknown, @@ -54,11 +53,7 @@ fn single_build(query: &str, accuracy: AccuracyTarget) -> Rc { .unwrap() .remove(0); let root = asap_executor::physical_planner::promql_rows::with_series_identity(&root).unwrap(); - let space = search_workload(vec![("q", root)]); - global_selection(&space, &DefaultCostModel) - .assemble_selected_query(&space.roots[0].1) - .unwrap() - .unwrap() + selected_dag(root, accuracy) } /// Rewrite the root's per-entity `SummaryAgg` over `TimeRange(5m)` into the diff --git a/crates/frontend-promql/tests/count_planning.rs b/crates/frontend-promql/tests/count_planning.rs index 76b527fdd..1972b43c2 100644 --- a/crates/frontend-promql/tests/count_planning.rs +++ b/crates/frontend-promql/tests/count_planning.rs @@ -4,17 +4,16 @@ use asap_logical_optimizer::{ default_strategies, search_workload_with_targets, ASAPStrategies, Replacement, ReplacementStrategy, TargetSubDAG, }; -use asap_plan_selection::candidate_selection::global_selection; -use asap_plan_selection::cost::cost_model::DefaultCostModel; mod support; use asap_types::ir::physical_export::PhysicalASAPOperatorPayload; use asap_types::ir::schema::{ - ExactKind, FieldDataType, NonNegativeWeightProof, SketchAlgorithm, SummaryInputExpr, - WeightDomain, + ExactKind, FieldDataType, GroupingStrategy, NonNegativeWeightProof, SketchAlgorithm, + SummaryInputExpr, WeightDomain, }; -use asap_types::ir::{ASAPOp, Operator}; +use asap_types::ir::{ASAPOp, Operator, OperatorNode}; use asap_types::types::AccuracyTarget; -use support::{lower_promql, post_asap_dag}; +use std::rc::Rc; +use support::{lower_promql, post_asap_dag, selected_dag}; #[test] fn grouped_count_keeps_uncertified_hydra_candidates_for_backend_review() { @@ -24,7 +23,7 @@ fn grouped_count_keeps_uncertified_hydra_candidates_for_backend_review() { }; let root = lower_promql("count by(job)(up)", target.clone()).unwrap(); let space = search_workload_with_targets( - vec![("count", root, Some(target))], + vec![("count", Rc::clone(&root), Some(target.clone()))], &default_strategies(), &DefaultAccuracyModel, ); @@ -40,11 +39,17 @@ fn grouped_count_keeps_uncertified_hydra_candidates_for_backend_review() { assert!(hydra .iter() .all(|candidate| candidate.has_missing_accuracy_evidence())); - assert!(!global_selection(&space, &DefaultCostModel) - .for_target(planned) - .unwrap() - .chosen - .is_some_and(|candidate| candidate.has_missing_accuracy_evidence())); + // Stage 3 never selects a summary without accuracy evidence. + let selected = selected_dag(root, target); + assert!(!OperatorNode::reachable(&selected) + .iter() + .any(|node| matches!( + &node.operator, + Operator::ASAP(ASAPOp::SummaryAgg { + family: FieldDataType::Sketch(_, GroupingStrategy::SharedMultiSubpopulation { .. }), + .. + }) + ))); } // Exact series and temporal counts must select a count accumulator, not distinct or sum. diff --git a/crates/frontend-promql/tests/support.rs b/crates/frontend-promql/tests/support.rs index 5549c13f2..618d86086 100644 --- a/crates/frontend-promql/tests/support.rs +++ b/crates/frontend-promql/tests/support.rs @@ -109,3 +109,32 @@ pub fn sample_expression(node: &OperatorNode) -> &ScalarExpr { other => panic!("expected sample expression, got {other:?}"), } } + +/// The logical DAG the #509 stage pipeline selects for the single query +/// `root` at `accuracy`, run once, with the built-in models. +#[allow(dead_code)] +pub fn selected_dag(root: Rc, accuracy: AccuracyTarget) -> Rc { + use asap_types::ir::QueryRoot; + use asap_types::workload::{QueryRecurrence, RootDemand}; + let demand = [RootDemand { + accuracy: Some(accuracy), + recurrence: QueryRecurrence::OneTime { + invocations: 1, + execute_at: None, + }, + predictability: Predictability::default(), + latency_ms: None, + }]; + let run = asap_plan_selection::plan_stages( + vec![(0, QueryRoot::Operator(root))], + &demand, + &DataWorkload::default(), + asap_plan_selection::PlanningModels::builtin(), + 0, + ) + .unwrap(); + let QueryRoot::Operator(root) = &run.plan.logical[0].1 else { + panic!("operator root") + }; + root.clone() +} diff --git a/crates/frontend-promql/tests/univmon_candidates.rs b/crates/frontend-promql/tests/univmon_candidates.rs index 3d34603fa..3bc4791db 100644 --- a/crates/frontend-promql/tests/univmon_candidates.rs +++ b/crates/frontend-promql/tests/univmon_candidates.rs @@ -7,8 +7,6 @@ use asap_logical_optimizer::pass1::replacement::{ default_strategies, search_workload_with_targets, }; use asap_logical_optimizer::{ASAPStrategies, Replacement, ReplacementStrategy, TargetSubDAG}; -use asap_plan_selection::candidate_selection::global_selection; -use asap_plan_selection::cost::cost_model::DefaultCostModel; mod support; use asap_types::ir::cse::share_common_sub_dags; use asap_types::ir::properties::{ @@ -17,7 +15,7 @@ use asap_types::ir::properties::{ use asap_types::ir::schema::{FieldDataType, SketchAlgorithm, SketchStatistic, SummaryInputExpr}; use asap_types::ir::{ASAPOp, Operator, OperatorNode}; use asap_types::types::AccuracyTarget; -use support::{lower_promql, post_asap_dag}; +use support::{lower_promql, post_asap_dag, selected_dag}; // Synthetic evidence exercises structural sharing, never runtime accuracy. struct TestEvidence; @@ -161,7 +159,7 @@ fn uncalibrated_frequency_evaluations_do_not_bypass_accuracy_targets() { } else { assert!(unknown > 0); let space = search_workload_with_targets( - vec![("q", Rc::clone(&root), Some(target))], + vec![("q", Rc::clone(&root), Some(target.clone()))], &default_strategies(), &DefaultAccuracyModel, ); @@ -171,11 +169,17 @@ fn uncalibrated_frequency_evaluations_do_not_bypass_accuracy_targets() { .candidates .iter() .any(|candidate| candidate.has_missing_accuracy_evidence())); - assert!(!global_selection(&space, &DefaultCostModel) - .for_target(&space.roots[0].1) - .unwrap() - .chosen - .is_some_and(|candidate| candidate.has_missing_accuracy_evidence())); + // Stage 3 never selects the uncalibrated UnivMon estimate. + let selected = selected_dag(root, target); + assert!(!OperatorNode::reachable(&selected) + .iter() + .any(|node| matches!( + &node.operator, + Operator::ASAP(ASAPOp::SummaryAgg { + family: FieldDataType::Sketch(kind, _), + .. + }) if kind.algorithm() == &SketchAlgorithm::UnivMon + ))); } } } diff --git a/crates/integration-tests/tests/executor_models/mod.rs b/crates/integration-tests/tests/executor_models/mod.rs index 25b92e9f8..702027fb6 100644 --- a/crates/integration-tests/tests/executor_models/mod.rs +++ b/crates/integration-tests/tests/executor_models/mod.rs @@ -10,3 +10,48 @@ static EXECUTOR: LazyLock = LazyLock::new(asap_executor: pub fn executor_models() -> PlanningModels<'static> { PlanningModels::builtin().with_capabilities(&EXECUTOR) } + +/// The logical DAG the #509 stage pipeline selects for the single query +/// `root` at `accuracy`, run once over a continuously ingested source, with +/// [`executor_models`]. +#[allow(dead_code)] +pub fn selected_dag( + root: std::rc::Rc, + accuracy: asap_types::types::AccuracyTarget, +) -> std::rc::Rc { + use asap_types::ir::QueryRoot; + use asap_types::workload::{ + DataArrival, DataWorkload, Evidence, EvidenceSource, Predictability, QueryRecurrence, Rate, + RootDemand, + }; + let demand = [RootDemand { + accuracy: Some(accuracy), + recurrence: QueryRecurrence::OneTime { + invocations: 1, + execute_at: None, + }, + predictability: Predictability::default(), + latency_ms: None, + }]; + let data = DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(1_000.0)), + source: EvidenceSource::Declared, + ..Default::default() + }, + ..Default::default() + }; + let run = asap_plan_selection::plan_stages( + vec![(0, QueryRoot::Operator(root))], + &demand, + &data, + executor_models(), + 0, + ) + .unwrap(); + let QueryRoot::Operator(root) = &run.plan.logical[0].1 else { + panic!("operator root") + }; + root.clone() +} diff --git a/crates/integration-tests/tests/operator_sharing.rs b/crates/integration-tests/tests/operator_sharing.rs index 790c28b67..72787432e 100644 --- a/crates/integration-tests/tests/operator_sharing.rs +++ b/crates/integration-tests/tests/operator_sharing.rs @@ -3,17 +3,15 @@ //! and below summary operators, can share inputs with them, and can carry //! summaries below set operators. //! -//! Each test drives SQL text through `lower_sql` → `search_workload` → -//! global selection → `assemble_selected_dag`, the pipeline -//! `sql_to_post_asap.rs` uses. +//! Each test drives SQL text through `lower_sql` and the #509 stage +//! pipeline (`plan_stages`) and inspects the selected logical DAG. + +mod executor_models; use std::rc::Rc; use asap_frontend_sql::{lower_sql, SqlCatalog}; use asap_integration_tests::post_asap::post_asap_dag; -use asap_logical_optimizer::search_workload; -use asap_plan_selection::candidate_selection::global_selection; -use asap_plan_selection::DefaultCostModel; use asap_types::ir::schema::{DataType, Field, Schema}; use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}; use asap_types::types::AccuracyTarget; @@ -48,18 +46,12 @@ fn catalog() -> SqlCatalog { ) } -/// Lower `sql`, search, select with the default cost model and assemble the -/// selected post-ASAP DAG. +/// Lower `sql` and return the plan the stage pipeline selects. async fn plan(sql: &str, accuracy: AccuracyTarget) -> Rc { - let pre = lower_sql(sql, &catalog(), accuracy) + let pre = lower_sql(sql, &catalog(), accuracy.clone()) .await .unwrap_or_else(|e| panic!("lower failed for {sql:?}: {e}")); - let space = search_workload(vec![("query", pre)]); - let selection = global_selection(&space, &DefaultCostModel); - selection - .assemble_selected_dag(&space.roots[0].1) - .expect("materialization failed") - .expect("root must be discovered") + executor_models::selected_dag(pre, accuracy) } /// Every unique node reachable from `root` whose operator matches `pred`. @@ -165,6 +157,7 @@ async fn exact_aggregate_and_sketch_share_one_scan() { // #468 problem 3: a summary can sit below a set operator — each side of the // UNION ALL holds its own SummaryEstimate. +#[ignore = "Stage 3 selects the raw plan; query-time summaries never cost less until Stage 2 plans materialization: #580"] #[tokio::test] async fn each_side_of_union_all_holds_a_summary_estimate() { let root = plan( diff --git a/crates/integration-tests/tests/precompute_raw_samples.rs b/crates/integration-tests/tests/precompute_raw_samples.rs index 94ba14eae..bd164abcb 100644 --- a/crates/integration-tests/tests/precompute_raw_samples.rs +++ b/crates/integration-tests/tests/precompute_raw_samples.rs @@ -1,5 +1,6 @@ //! Planner-selected summaries over raw samples compile as precompute DAGs //! and produce the same estimates as feeding their kernel sample by sample. +mod executor_models; mod physical_common; use asap_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload}; use asap_types::ir::ASAPOp; @@ -18,11 +19,8 @@ use asap_executor::{ }; use asap_integration_tests::fixtures::lower_promql; use asap_logical_optimizer::{ - search_workload, ASAPStrategies, Replacement, ReplacementStrategy, ReplacementSubDAG, - TargetSubDAG, + ASAPStrategies, Replacement, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, }; -use asap_plan_selection::candidate_selection::global_selection; -use asap_plan_selection::cost::cost_model::DefaultCostModel; use asap_types::ir::operator::Reduction; use asap_types::ir::scalar::ColumnRef; use asap_types::ir::schema::{ @@ -55,10 +53,10 @@ fn canonical(labels: &Series) -> Series { .collect() } -/// Every Planner candidate for `query`: the searched selection plus each -/// summary replacement of the root. +/// Every Planner candidate for `query`: the stage pipeline's selection plus +/// each summary replacement of the root. fn candidates(query: &str, accuracy: AccuracyTarget) -> Vec> { - let root = lower_promql(query, accuracy).expect("lowering failed"); + let root = lower_promql(query, accuracy.clone()).expect("lowering failed"); let mut result = ASAPStrategies::default() .replacements(&TargetSubDAG::new(&root)) .into_iter() @@ -70,12 +68,7 @@ fn candidates(query: &str, accuracy: AccuracyTarget) -> Vec> { _ => None, }) .collect::>(); - let space = search_workload(vec![("query", root)]); - if let Ok(Some(selected)) = - global_selection(&space, &DefaultCostModel).assemble_selected_dag(&space.roots[0].1) - { - result.push(selected); - } + result.push(executor_models::selected_dag(root, accuracy)); result } diff --git a/crates/integration-tests/tests/sql_to_physical.rs b/crates/integration-tests/tests/sql_to_physical.rs index ed03b1372..10f1044e3 100644 --- a/crates/integration-tests/tests/sql_to_physical.rs +++ b/crates/integration-tests/tests/sql_to_physical.rs @@ -1,4 +1,6 @@ -//! SQL frontend, candidate selection, physical compilation and fresh-run execution. +//! SQL frontend, the stage pipeline's selection, physical compilation and +//! fresh-run execution. +mod executor_models; mod physical_common; use asap_executor::{ physical_planner::{compile, InputContract, Source}, @@ -7,9 +9,6 @@ use asap_executor::{ values::{Batch, Value}, }; use asap_frontend_sql::{lower_sql, SqlCatalog}; -use asap_logical_optimizer::search_workload; -use asap_plan_selection::candidate_selection::global_selection; -use asap_plan_selection::DefaultCostModel; use asap_types::ir::physical_export::PhysicalASAPOperatorPayload; use asap_types::ir::schema::{DataType, Field, FieldDataType, Schema}; use asap_types::types::AccuracyTarget; @@ -35,11 +34,7 @@ async fn sql_filter_grouped_sum_executes_and_rebinds() { let logical = lower_sql(query, &catalog, AccuracyTarget::Exact) .await .unwrap(); - let space = search_workload(vec![("sql", logical)]); - let selected = global_selection(&space, &DefaultCostModel) - .assemble_selected_dag(&space.roots[0].1) - .unwrap() - .unwrap(); + let selected = executor_models::selected_dag(logical, AccuracyTarget::Exact); let dag = compile_physical_asap_dag(&selected).unwrap(); let scan = dag .nodes diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs index f2404b061..cc96fa84a 100644 --- a/crates/planner/tests/summary_sharing.rs +++ b/crates/planner/tests/summary_sharing.rs @@ -1,21 +1,13 @@ //! Structurally identical summary producers chosen by different queries are //! shared after Pass 1: one `Rc` across their plans. -use asap_types::ir::cse::share_common_sub_dags; use asap_types::ir::{ASAPOp, OperatorNode}; use std::rc::Rc; -use asap_frontend_promql::lower_promql_workload; use asap_frontend_sql::SqlCatalog; -use asap_logical_optimizer::accuracy::{ - AccuracyModel, DefaultAccuracyModel, EqualSplitAllocator, PropagationStats, -}; +use asap_logical_optimizer::accuracy::{AccuracyModel, DefaultAccuracyModel, PropagationStats}; use asap_logical_optimizer::pass1::replacement::{default_size_params, DEFAULT_DELTA}; -use asap_logical_optimizer::{ - search_workload_with_targets, ASAPStrategies, Replacement, ReplacementStrategy, - ReplacementSubDAG, TargetSubDAG, -}; -use asap_plan_selection::candidate_selection::global_selection; +use asap_logical_optimizer::{Replacement, ReplacementSubDAG, TargetSubDAG}; use asap_plan_selection::PlanningModels; use asap_plan_selection::{CostModel, DefaultCostModel}; use asap_planner::pass::{PlanOutput, QueryPlan}; @@ -84,15 +76,6 @@ const PREFER_LARGE_KLL: PreferSketch = PreferSketch(|kind| match kind.params() { _ => 1.0, }); -/// Prefers UnivMon, which can serve every frequency moment from one state. -const PREFER_UNIVMON: PreferSketch = PreferSketch(|kind| { - if kind.algorithm() == &SketchAlgorithm::UnivMon { - 0.0 - } else { - 1.0 - } -}); - fn requirements(epsilon: f64) -> QueryRequirements { QueryRequirements { accuracy: AccuracyRequirement::Explicit(AccuracyTarget::Epsilon(epsilon)), @@ -501,61 +484,6 @@ impl AccuracyModel for UnivMonEvidence { } } -/// Distinct count, entropy and L2 over one input, certified by an accuracy -/// model and selected by a cost model preferring UnivMon, read one UnivMon state: #515 sharing is the summary-capability rule -/// when the states are identical. This runs the legacy search; the stage -/// pipeline's counterpart is -/// `frequency_moments_share_one_univmon_in_the_stage_pipeline`. -#[test] -fn certified_frequency_evaluations_share_one_univmon_state() { - let queries = [ - ("distinct_over_time(m[5m])", 0.02), - ("entropy_over_time(m[5m])", 0.02), - ("l2_over_time(m[5m])", 0.02), - ]; - let workload = promql_workload(&queries); - let roots = lower_promql_workload(&workload, NOW_MS) - .expect("lowers") - .into_iter() - .zip(queries) - .enumerate() - .map(|(index, (expr, (_, epsilon)))| (index, expr, Some(AccuracyTarget::Epsilon(epsilon)))) - .collect(); - let strategies: Vec> = vec![Box::new( - ASAPStrategies::new_with_planning_inputs(&UnivMonEvidence, &EqualSplitAllocator), - )]; - let space = search_workload_with_targets(roots, &strategies, &UnivMonEvidence); - let selection = global_selection(&space, &PREFER_UNIVMON); - let assembled = space - .roots - .iter() - .map(|(index, root)| { - let dag = selection - .assemble_selected_dag(root) - .expect("assembles") - .expect("root has a group"); - (*index, dag) - }) - .collect(); - let mut states: Vec> = Vec::new(); - for (_, root) in share_common_sub_dags(assembled) { - assert!(root.guarantee.is_some(), "{:?}", root.operator); - let asap_types::ir::Operator::ASAP(ASAPOp::SummaryEstimate { summary_input, .. }) = - &root.operator - else { - panic!("summary evaluation: {:?}", root.operator); - }; - assert!(matches!( - &summary_input.operator, - asap_types::ir::Operator::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) - if kind.algorithm() == &SketchAlgorithm::UnivMon - )); - states.push(Rc::clone(summary_input)); - } - assert_eq!(states.len(), 3); - assert!(states.iter().all(|state| Rc::ptr_eq(state, &states[0]))); -} - /// #509 Example 2 through the stage pipeline: distinct count, entropy and L2 /// of one input over one window, with requirements ε = 0.02, 0.05 and 0.01. /// The summary-capability variant gives the three targets one UnivMon.