diff --git a/crates/devtools/src/bin/sketch_coverage.rs b/crates/devtools/src/bin/sketch_coverage.rs index 3a3e604c4..6c5865c95 100644 --- a/crates/devtools/src/bin/sketch_coverage.rs +++ b/crates/devtools/src/bin/sketch_coverage.rs @@ -2,36 +2,24 @@ // // Lowers every query in every corpus we have (mirrors `variant_coverage`'s // corpus list exactly, so the two reports are directly comparable) with an -// *approximate* `AccuracyTarget`, runs `asap_logical_optimizer::explain_replacements` -// over each corpus as one workload, and reports the MVP demo's query-coverage -// metric: of the queries that lowered successfully, what fraction got -// -// - a `SketchApproximation` candidate (a genuine sketch alternative was -// found for at least one aggregate in the query — the KLL-vs-DDSketch -// kind of degree of freedom), and/or -// - a `CommonSubexpressionReuse` candidate (the query shares a sub-DAG, -// inside itself or with another query in the same corpus, that a -// build-once-and-share candidate was found for). +// *approximate* `AccuracyTarget`, runs Stage 1's Pass 1 +// (`enumerate_local_logical_candidates`) over each query, and reports which +// sketch alternatives Pass 1 offers for it, plus the fraction of lowered +// queries with at least one. // // `--epsilon ` (default 0.01) sets the `AccuracyTarget` every query in -// every corpus lowers with. Without an approximate target, -// `ASAPStrategies` never has a genuine sketch alternative to -// report. -// -// This is a workload-level count (`explain_replacements` runs -// `search_workload` once per corpus, over every query in it together), not -// just a per-query re-run of the single-target path — so cross-query CSE -// reuse inside one corpus shows up here the same way it would in the -// dag-viewer's Union mode. +// every corpus lowers with. Without an approximate target, Pass 1 offers no +// sketch alternative. use asap_devtools::lower_promql_with_data_ingestion_interval; use asap_frontend_sql::{lower_sql_dialect, SqlCatalog}; -use asap_logical_optimizer::{explain_replacements, ExplanationKind}; -use asap_types::ir::schema::{DataType, Field, Schema}; -use asap_types::ir::OperatorNode; +use asap_logical_optimizer::pass1::logical_candidates::enumerate_local_logical_candidates; +use asap_logical_optimizer::Realization; +use asap_types::ir::schema::{DataType, Field, GroupingStrategy, Schema}; +use asap_types::ir::{OperatorNode, QueryRoot}; use asap_types::types::AccuracyTarget; use asap_types::workload::SqlDialect; -use std::collections::BTreeSet; +use std::collections::{BTreeMap, BTreeSet}; use std::rc::Rc; /// Line-based `#`/`--` comment stripping, then split on `;` — the shape every @@ -123,9 +111,10 @@ struct CorpusCoverage { name: &'static str, lowered: usize, failed: usize, - sketch_covered: usize, - cse_covered: usize, - either_covered: usize, + /// Lowered queries Pass 1 rejected. + rejected: usize, + /// Per lowered query Pass 1 accepted, the sketch alternatives it offers. + sketches: Vec<(String, BTreeSet)>, } fn pct(n: usize, total: usize) -> String { @@ -136,86 +125,69 @@ fn pct(n: usize, total: usize) -> String { } } -/// `explain_replacements`' `location` is a comma-joined list of breadcrumbs -/// (`collect_locations` in `explanation.rs`), one per path from a workload -/// root to the target — e.g. `root "q3"` or `root "q3" > lhs`. A location -/// "covers" `label` (a bare `root "qN"` breadcrumb) if `label` is exactly one -/// of those comma-separated entries or the prefix of one that goes deeper — -/// i.e. the explanation's target is reachable from that query's root at all. -fn covers(location: &str, label: &str) -> bool { - location - .split(", ") - .any(|loc| loc == label || loc.starts_with(&format!("{label} > "))) -} - -fn root_label(id: &str) -> String { - format!("root {id:?}") -} - -/// Run `explain_replacements` over one corpus's already-lowered roots as one -/// workload, then attribute each finding back to the query root(s) it's -/// reachable from. +/// Run Pass 1 over each of one corpus's already-lowered queries and collect +/// the sketch alternatives it offers for any of the query's targets. A Hydra +/// alternative is listed as `Hydra()`. fn analyze_corpus( name: &'static str, roots: Vec<(String, Rc)>, failed: usize, ) -> CorpusCoverage { let lowered = roots.len(); - let labels: Vec = roots.iter().map(|(id, _)| root_label(id)).collect(); - let explanations = explain_replacements(roots); - - let mut sketch_covered: BTreeSet = BTreeSet::new(); - let mut cse_covered: BTreeSet = BTreeSet::new(); - for explanation in &explanations { - for (i, label) in labels.iter().enumerate() { - if !covers(&explanation.location, label) { - continue; - } - match explanation.kind { - ExplanationKind::SketchApproximation => { - sketch_covered.insert(i); - } - ExplanationKind::CommonSubexpressionReuse => { - cse_covered.insert(i); + let mut rejected = 0; + let mut sketches = Vec::new(); + for (id, root) in roots { + let Ok(inventory) = enumerate_local_logical_candidates( + vec![(id.clone(), QueryRoot::Operator(root))], + &BTreeMap::new(), + ) else { + rejected += 1; + continue; + }; + let offered = inventory + .targets + .iter() + .flat_map(|target| target.alternatives.iter().zip(&target.groupings)) + .filter_map(|(alternative, grouping)| match alternative { + Realization::Sketch(kind) if *grouping == GroupingStrategy::default() => { + Some(format!("{:?}", kind.algorithm())) } - // `#[non_exhaustive]`: a future kind just doesn't count - // toward either bucket here until this tool is taught about it. - _ => {} - } - } + Realization::Sketch(kind) => Some(format!("Hydra({:?})", kind.algorithm())), + _ => None, + }) + .collect(); + sketches.push((id, offered)); } - let either_covered = sketch_covered.union(&cse_covered).count(); - CorpusCoverage { name, lowered, failed, - sketch_covered: sketch_covered.len(), - cse_covered: cse_covered.len(), - either_covered, + rejected, + sketches, } } +fn covered(r: &CorpusCoverage) -> usize { + r.sketches.iter().filter(|(_, s)| !s.is_empty()).count() +} + fn report(r: &CorpusCoverage) { println!("--- {} ---", r.name); - println!("lowered: {}, failed: {}", r.lowered, r.failed); println!( - "sketch-approximable: {}/{} ({})", - r.sketch_covered, - r.lowered, - pct(r.sketch_covered, r.lowered) - ); - println!( - "CSE-shareable: {}/{} ({})", - r.cse_covered, - r.lowered, - pct(r.cse_covered, r.lowered) + "lowered: {}, failed: {}, rejected by Pass 1: {}", + r.lowered, r.failed, r.rejected ); + for (id, offered) in &r.sketches { + if !offered.is_empty() { + let offered: Vec<_> = offered.iter().map(String::as_str).collect(); + println!(" {id}: {}", offered.join(", ")); + } + } println!( - "either (coverage): {}/{} ({})", - r.either_covered, + "sketch-approximable: {}/{} ({})", + covered(r), r.lowered, - pct(r.either_covered, r.lowered) + pct(covered(r), r.lowered) ); println!(); } @@ -332,22 +304,16 @@ async fn main() { let total_lowered: usize = results.iter().map(|r| r.lowered).sum(); let total_failed: usize = results.iter().map(|r| r.failed).sum(); - let total_sketch: usize = results.iter().map(|r| r.sketch_covered).sum(); - let total_cse: usize = results.iter().map(|r| r.cse_covered).sum(); - let total_either: usize = results.iter().map(|r| r.either_covered).sum(); + let total_rejected: usize = results.iter().map(|r| r.rejected).sum(); + let total_sketch: usize = results.iter().map(covered).sum(); println!("=== global (epsilon = {epsilon}) ==="); - println!("total lowered: {total_lowered}, total failed: {total_failed}"); println!( - "sketch-approximable: {total_sketch}/{total_lowered} ({})", - pct(total_sketch, total_lowered) + "total lowered: {total_lowered}, total failed: {total_failed}, \ + rejected by Pass 1: {total_rejected}" ); println!( - "CSE-shareable: {total_cse}/{total_lowered} ({})", - pct(total_cse, total_lowered) - ); - println!( - "either (coverage): {total_either}/{total_lowered} ({})", - pct(total_either, total_lowered) + "sketch-approximable: {total_sketch}/{total_lowered} ({})", + pct(total_sketch, total_lowered) ); } diff --git a/crates/executor/tests/common/mod.rs b/crates/executor/tests/common/mod.rs index 0184a51b6..c07ec62fc 100644 --- a/crates/executor/tests/common/mod.rs +++ b/crates/executor/tests/common/mod.rs @@ -61,3 +61,30 @@ pub fn selected_dag( }; root.clone() } + +/// Every whole-query candidate Stage 1's Pass 1 lists for `root` (up to +/// 4096), composed; choices that do not compose are skipped. +pub fn stage1_candidates(root: &Rc) -> Vec> { + use asap_logical_optimizer::pass1::logical_candidates::{ + compose_logical_candidate, enumerate_choices, enumerate_local_logical_candidates, + }; + use planner_types::ir::QueryRoot; + let inventory = enumerate_local_logical_candidates( + vec![(0, QueryRoot::Operator(Rc::clone(root)))], + &Default::default(), + ) + .unwrap(); + enumerate_choices(&inventory, 4096) + .iter() + .filter_map(|choice| { + match compose_logical_candidate(&inventory, choice) + .ok()? + .remove(0) + .1 + { + QueryRoot::Operator(node) => Some(node), + QueryRoot::Scalar(_) => None, + } + }) + .collect() +} diff --git a/crates/executor/tests/deployment_computation.rs b/crates/executor/tests/deployment_computation.rs index c4b08d3fe..c3ab78e79 100644 --- a/crates/executor/tests/deployment_computation.rs +++ b/crates/executor/tests/deployment_computation.rs @@ -7,12 +7,11 @@ use asap_executor::{ runtime::{Limits, RunContext, Scope}, values::{Batch, Value}, }; -use common::{compile_physical_asap_dag, selected_dag}; +use common::{compile_physical_asap_dag, selected_dag, stage1_candidates}; use futures::{executor::block_on, StreamExt}; use planner_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload}; -use planner_types::ir::schema::*; use planner_types::ir::ASAPOp; -use planner_types::{types::AccuracyTarget, workload::*}; +use planner_types::{ir::schema::*, types::AccuracyTarget, workload::*}; use std::{collections::BTreeMap, rc::Rc, sync::Arc}; fn lower(query: &str) -> Rc { @@ -611,31 +610,26 @@ fn population_sums_and_averages_are_compensated() { #[test] fn stored_count_min_bare_count_compiles_to_a_evaluation() { use asap_executor::summary_kernels::CountMinSketchAccumulator; - use asap_logical_optimizer::{Replacement, ReplacementStrategy, TargetSubDAG}; let root = lower_with("count(up)", AccuracyTarget::Epsilon(0.02)); - let dag = asap_logical_optimizer::ASAPStrategies::default() - .replacements(&TargetSubDAG::new(&root)) + let dag = stage1_candidates(&root) .into_iter() - .find_map(|candidate| match candidate.replacement { - Replacement::SubDAG(node) => { - let dag = compile_physical_asap_dag(&node).ok()?; - let bare_count = dag.nodes.iter().any(|n| { - matches!( - &n.payload, - PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryEstimate { - query: SketchStatistic::PointCount { value: None, .. }, - .. - }) - ) - }); - let count_min = dag.nodes.iter().any(|n| { - matches!(&n.payload, PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg { + .find_map(|node| { + let dag = compile_physical_asap_dag(&node).ok()?; + let bare_count = dag.nodes.iter().any(|n| { + matches!( + &n.payload, + PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryEstimate { + query: SketchStatistic::PointCount { value: None, .. }, + .. + }) + ) + }); + let count_min = dag.nodes.iter().any(|n| { + matches!(&n.payload, PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) if kind.algorithm() == &SketchAlgorithm::Cms) - }); - (bare_count && count_min).then_some(dag) - } - _ => None, + }); + (bare_count && count_min).then_some(dag) }) .expect("Planner lists a Count-Min candidate for count(up)"); let state = dag diff --git a/crates/executor/tests/planspace_series_identity_heap.rs b/crates/executor/tests/planspace_series_identity_heap.rs index a9680b2d2..f2aa24c03 100644 --- a/crates/executor/tests/planspace_series_identity_heap.rs +++ b/crates/executor/tests/planspace_series_identity_heap.rs @@ -1,23 +1,11 @@ -//! Logical heap alternatives that need the PromQL series identity are part of -//! Planner's search space: `enumerate_candidate_dags_for_root` lists -//! current-series TopK heaps without a caller-side series-identity pass, cost -//! ranking, or workload Cartesian expansion. Placement variants are not listed. +//! The stage pipeline's selection never commits a heap that needs the PromQL +//! series identity; a deployment prices such heaps. mod common; use common::{compile_physical_asap_dag, selected_dag}; use planner_types::ir::OperatorNode; -use asap_executor::physical_planner::promql_rows::{ - compile_current_series_evaluation, SERIES_IDENTITY_COLUMN, -}; -use asap_logical_optimizer::{ - accuracy::AccuracyEvidenceProvider, accuracy::DefaultAccuracyModel, accuracy::PropagationStats, - pass1::replacement::default_strategies_with_evidence, - pass1::replacement::ReplacementProvenance, search_workload_with_targets, Proposals, - ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, -}; +use asap_executor::physical_planner::promql_rows::SERIES_IDENTITY_COLUMN; use planner_types::{ - ir::properties::*, - ir::schema::*, types::AccuracyTarget, workload::{ AccuracyRequirement, BatchEntry, DataWorkload, DurationMs, Evidence as WorkloadEvidence, @@ -27,47 +15,6 @@ use planner_types::{ }; use std::rc::Rc; -struct Evidence; -impl AccuracyEvidenceProvider for Evidence { - fn topk_max_distinct_items(&self, _: &OperatorNode) -> Option { - Some(1000) - } - fn propagation_stats( - &self, - op: &CompositionOperator, - _: &FieldDataType, - _: Option<&SketchStatistic>, - ) -> PropagationStats { - if matches!(op, CompositionOperator::TopKSelection) { - PropagationStats { - topk_selected_lower_bound: Some(101.), - topk_excluded_upper_bound: Some(100.), - topk_interval_failure_probability: Some(0.001), - ..Default::default() - } - } else { - Default::default() - } - } -} - -/// Forwards everything except whole-root proposals: the pre-change search. -struct LogicalOnly(Box); -impl ReplacementStrategy for LogicalOnly { - fn name(&self) -> &'static str { - self.0.name() - } - fn matches(&self, target: &TargetSubDAG<'_>) -> bool { - self.0.matches(target) - } - fn replacements(&self, target: &TargetSubDAG<'_>) -> Vec { - self.0.replacements(target) - } - fn propose(&self, target: &TargetSubDAG<'_>) -> Proposals { - self.0.propose(target) - } -} - fn lower(query: &str, accuracy: &AccuracyTarget) -> Rc { let workload = PlanningWorkload { query_workload: QueryWorkload { @@ -98,152 +45,28 @@ fn lower(query: &str, accuracy: &AccuracyTarget) -> Rc { .remove(0) } -type InventoryDAG = Vec<(usize, Rc)>; - -/// Candidate DAGs for query 1 of a two-query workload, with and without -/// whole-root proposals. Query 0 is a bystander that must not multiply them. -fn inventories(query: &str, accuracy: AccuracyTarget) -> (Vec, Vec) { - let roots = vec![ - ( - 0, - lower("sum by(job)(m)", &AccuracyTarget::Exact), - Some(AccuracyTarget::Exact), - ), - (1, lower(query, &accuracy), Some(accuracy)), - ]; - let full = default_strategies_with_evidence(&Evidence); - let logical: Vec> = default_strategies_with_evidence(&Evidence) - .into_iter() - .map(|strategy| Box::new(LogicalOnly(strategy)) as Box) - .collect(); - let enumerate = |strategies: &[Box]| { - search_workload_with_targets(roots.clone(), strategies, &DefaultAccuracyModel) - .enumerate_candidate_dags_for_root(&1, 65_536) - .unwrap() - .candidates - }; - (enumerate(&full), enumerate(&logical)) -} - /// Whether an operator below the root carries series identity. A top-k root /// returns its selected series' identity whatever realizes it. -fn carries_identity(dag: &InventoryDAG) -> bool { - dag.iter().any(|(_, root)| { - let dag = compile_physical_asap_dag(root).unwrap(); - dag.nodes - .iter() - .filter(|node| !dag.roots.contains(&node.id)) - .any(|node| { - node.output_schema - .fields - .iter() - .any(|field| field.name == SERIES_IDENTITY_COLUMN) - }) - }) -} - -/// Shared acceptance checks; returns the added identity-carrying alternatives. -fn added_alternatives( - query: &str, - accuracy: AccuracyTarget, -) -> Vec> { - let (full, logical) = inventories(query, accuracy); - for (index, dag) in full.iter().enumerate() { - assert_eq!(dag.len(), 1, "one root per candidate, no workload product"); - assert!( - !full[..index].contains(dag), - "{query}: identical DAG listed twice" - ); - } - let (added, kept): (Vec<_>, Vec<_>) = full.into_iter().partition(carries_identity); - assert_eq!( - kept, logical, - "{query}: existing candidates must be unchanged" - ); - added.into_iter().map(|mut dag| dag.remove(0).1).collect() +fn carries_identity(root: &Rc) -> bool { + let dag = compile_physical_asap_dag(root).unwrap(); + dag.nodes + .iter() + .filter(|node| !dag.roots.contains(&node.id)) + .any(|node| { + node.output_schema + .fields + .iter() + .any(|field| field.name == SERIES_IDENTITY_COLUMN) + }) } const CURRENT_SERIES_TOPK: &str = "topk by(job)(1, m)"; -// Instant-vector TopK lists finalized current-series heap evaluations. -#[test] -fn current_series_topk_lists_heap_evaluations() { - let added = added_alternatives(CURRENT_SERIES_TOPK, AccuracyTarget::Epsilon(0.1)); - assert!(!added.is_empty()); - for root in added { - assert!(!matches!( - root.operator, - planner_types::ir::Operator::ASAP(planner_types::ir::ASAPOp::SummaryAgg { .. }) - )); - assert!( - compile_current_series_evaluation(&root).is_ok(), - "unbindable alternative {root:?}" - ); - } -} - -// Rate queries gain no fixed-window or query-time placement variants. -#[test] -fn rate_placement_variants_are_not_listed() { - for (query, accuracy) in [ - ("sum by(job)(rate(m[1m]))", AccuracyTarget::Exact), - ("topk by(job)(2, rate(m[1m]))", AccuracyTarget::Epsilon(0.1)), - ("rate(m[1m])", AccuracyTarget::Exact), - ] { - assert!(added_alternatives(query, accuracy).is_empty(), "{query}"); - } -} - -// Queries without a current-series heap realization are unchanged. -#[test] -fn unrelated_queries_keep_their_inventory() { - for (query, accuracy) in [ - ("sum by(job)(m)", AccuracyTarget::Exact), - ( - "quantile_over_time(0.9, m[1m])", - AccuracyTarget::Epsilon(0.05), - ), - ("max_over_time(m[1m])", AccuracyTarget::Exact), - ] { - assert!(added_alternatives(query, accuracy).is_empty(), "{query}"); - } -} - // The stage pipeline's selection keeps the logical plan; deployment prices heaps. #[test] fn selection_never_commits_a_series_identity_heap() { let accuracy = AccuracyTarget::Epsilon(0.1); let root = lower(CURRENT_SERIES_TOPK, &accuracy); let selected = selected_dag(root, accuracy); - assert!(!carries_identity(&vec![(0, selected)])); -} - -// A query repeated in the workload is proposed once, not once per copy. -#[test] -fn repeated_roots_do_not_duplicate_alternatives() { - let accuracy = AccuracyTarget::Epsilon(0.1); - let strategies = default_strategies_with_evidence(&Evidence); - let count = |copies: usize| { - let roots = (0..copies) - .map(|id| { - ( - id, - lower(CURRENT_SERIES_TOPK, &accuracy), - Some(accuracy.clone()), - ) - }) - .collect(); - let space = search_workload_with_targets(roots, &strategies, &DefaultAccuracyModel); - space - .candidates_for_target(&space.roots[0].1) - .unwrap() - .candidates - .iter() - .filter(|candidate| { - candidate.provenance == ReplacementProvenance::RootPhysicalRealization - }) - .count() - }; - assert!(count(1) > 0); - assert_eq!(count(2), count(1)); + assert!(!carries_identity(&selected)); } diff --git a/crates/executor/tests/precompute_candidates.rs b/crates/executor/tests/precompute_candidates.rs index 41fe0becb..a4823d5bb 100644 --- a/crates/executor/tests/precompute_candidates.rs +++ b/crates/executor/tests/precompute_candidates.rs @@ -11,8 +11,7 @@ use asap_executor::{ runtime::{Limits, RunContext, Scope}, values::{Batch, Value}, }; -use asap_logical_optimizer::search_workload; -use common::{compile_physical_asap_dag, selected_dag}; +use common::{compile_physical_asap_dag, selected_dag, stage1_candidates}; use futures::{executor::block_on, StreamExt}; use planner_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload}; use planner_types::ir::schema::DataType; @@ -52,10 +51,6 @@ fn grouped_rate_root() -> std::rc::Rc { 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 { compile_physical_asap_dag(&selected_dag(grouped_rate_root(), AccuracyTarget::Exact)).unwrap() } @@ -410,11 +405,9 @@ fn bounded_inventory_exposes_grouped_rate_physical_frontiers() { #[test] fn enumerated_grouped_rate_candidates_execute_numeric_query_outputs() { - let inventory = grouped_rate_space().enumerate_candidate_dags(4096).unwrap(); let mut executed = 0; - for forest in inventory.candidates { - let root = &forest[0].1; - let dag = compile_physical_asap_dag(root).unwrap(); + for root in stage1_candidates(&grouped_rate_root()) { + let dag = compile_physical_asap_dag(&root).unwrap(); let Some(state) = dag.nodes.iter().find(|node| { matches!( node.payload, diff --git a/crates/executor/tests/weighted_topk_binding.rs b/crates/executor/tests/weighted_topk_binding.rs index 756f66359..835404155 100644 --- a/crates/executor/tests/weighted_topk_binding.rs +++ b/crates/executor/tests/weighted_topk_binding.rs @@ -6,87 +6,40 @@ use asap_executor::dag::{ values::{Batch, Value}, Limits, RunContext, Scope, }; -use asap_logical_optimizer::{ - accuracy::AccuracyEvidenceProvider, accuracy::DefaultAccuracyModel, - accuracy::EqualSplitAllocator, accuracy::PropagationStats, ASAPStrategies, Replacement, - ReplacementStrategy, TargetSubDAG, -}; -use common::compile_physical_asap_dag; +use common::{compile_physical_asap_dag, stage1_candidates}; use futures::{executor::block_on, StreamExt}; use planner_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload}; use planner_types::ir::properties::*; -use planner_types::ir::schema::DataType; -use planner_types::ir::schema::*; +use planner_types::ir::schema::{DataType, *}; use planner_types::ir::ASAPOp; use planner_types::types::AccuracyTarget; use std::{collections::BTreeMap, rc::Rc, sync::Arc}; -struct Evidence; -impl AccuracyEvidenceProvider for Evidence { - fn topk_max_distinct_items(&self, _: &planner_types::ir::OperatorNode) -> Option { - Some(1000) - } - fn propagation_stats( - &self, - op: &CompositionOperator, - _: &FieldDataType, - _: Option<&SketchStatistic>, - ) -> PropagationStats { - if matches!(op, CompositionOperator::TopKSelection) { - PropagationStats { - topk_selected_lower_bound: Some(101.), - topk_excluded_upper_bound: Some(100.), - topk_interval_failure_probability: Some(0.001), - ..Default::default() - } - } else { - Default::default() - } - } -} -// The evidence here exercises binding; it is not inferred from the sample data. #[test] +#[ignore = "Stage 1's whole-expression heap keeps the outer top-k's partition keys, which \ + index the inner aggregate's output, over the raw rates it reads"] fn planner_weighted_topk_binds_at_either_deployment_phase() { - assert_weighted_binding(&Evidence, SketchAlgorithm::CmsWithHeap); - assert_weighted_binding(&Evidence, SketchAlgorithm::CountSketchWithHeap); + assert_weighted_binding(SketchAlgorithm::CmsWithHeap); + assert_weighted_binding(SketchAlgorithm::CountSketchWithHeap); } -// Binding validates representation, while deployment owns evidence acceptance. -#[test] -fn physical_binding_does_not_impose_an_accuracy_acceptance_policy() { - assert_weighted_binding( - &asap_logical_optimizer::accuracy::NoAccuracyEvidence, - SketchAlgorithm::CmsWithHeap, - ); - assert_weighted_binding( - &asap_logical_optimizer::accuracy::NoAccuracyEvidence, - SketchAlgorithm::CountSketchWithHeap, - ); +/// Whether `dag` builds a heap sketch of `algorithm`. +fn builds(dag: &PhysicalASAPDAG, algorithm: &SketchAlgorithm) -> bool { + dag.nodes.iter().any(|node| matches!(&node.payload, + PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) if kind.algorithm() == algorithm)) } -fn assert_weighted_binding(evidence: &dyn AccuracyEvidenceProvider, algorithm: SketchAlgorithm) { +fn assert_weighted_binding(algorithm: SketchAlgorithm) { let root = lower_promql( "topk by(job)(2, sum by(service, job)(rate(m[1m])))", AccuracyTarget::Epsilon(0.1), ) .unwrap(); - let strategy = ASAPStrategies::new_with_planning_inputs_and_evidence( - &DefaultAccuracyModel, - &EqualSplitAllocator, - evidence, - ); - let plan = strategy - .replacements(&TargetSubDAG::new(&root)) - .into_iter() - .find_map(|candidate| match candidate.replacement { - Replacement::SubDAG(node) - if candidate.rationale.contains(&format!("{algorithm:?}")) => - { - Some(node) - } - _ => None, - }) + // The whole-expression heap: it ranks the rates, absorbing the grouped Sum. + let dag = stage1_candidates(&root) + .iter() + .map(|candidate| compile_physical_asap_dag(candidate).unwrap()) + .find(|dag| builds(dag, &algorithm)) .unwrap(); - let dag = compile_physical_asap_dag(&plan).unwrap(); let build=dag.nodes.iter().find(|node|matches!(&node.payload,PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg {family:FieldDataType::Sketch(kind,_),..})if kind.algorithm()==&algorithm)).unwrap(); let rate_id = dag .edges @@ -261,388 +214,73 @@ fn rate_updates_cannot_enter_integer_heap_factory() { .is_err()); } -/// A catalog-resolved per-series rate can feed a heap sketch directly, without -/// requiring an otherwise unnecessary grouped Sum between Rate and TopK. -#[test] -fn direct_rate_topk_exposes_heap_candidates_with_complete_series_identity() { - check_direct_rate_topk(false); -} - -// Unreferenced labels still distinguish series throughout Rate and heap evaluation. +// A per-series rate feeds a heap sketch directly, without a grouped Sum +// between Rate and TopK. Unreferenced labels still distinguish series +// throughout Rate and heap evaluation. #[test] fn direct_rate_topk_preserves_dynamic_unreferenced_labels() { - check_direct_rate_topk(true); -} - -fn check_direct_rate_topk(dynamic: bool) { use asap_executor::physical_planner::promql_rows::{ decode_series_identity, series_row, with_series_identity, SERIES_IDENTITY_COLUMN, }; - let mut logical = + let logical = lower_promql("topk by(job)(2, rate(m[1m]))", AccuracyTarget::Epsilon(0.1)).unwrap(); - fn resolve_catalog(node: &mut planner_types::ir::OperatorNode) { - match &mut node.operator { - planner_types::ir::Operator::NonASAP( - planner_types::ir::NonASAPOp::Aggregate { child, .. } - | planner_types::ir::NonASAPOp::TimeRange { child, .. }, - ) => resolve_catalog(Rc::make_mut(child)), - planner_types::ir::Operator::NonASAP(planner_types::ir::NonASAPOp::Scan { - schema, - .. - }) => { - schema.closed = true; - schema.fields.push(planner_types::ir::schema::Field::plain( - "service", - DataType::Utf8, - false, - )); - } - _ => panic!("unexpected input shape: {node:?}"), - } - node.schema = node.operator.output_schema().unwrap(); - } - if dynamic { - logical = with_series_identity(&logical).unwrap(); - } else { - resolve_catalog(Rc::make_mut(&mut logical)); - } - let root = logical; - let strategy = ASAPStrategies::new_with_planning_inputs_and_evidence( - &DefaultAccuracyModel, - &EqualSplitAllocator, - &Evidence, - ); - let candidates = strategy.replacements(&TargetSubDAG::new(&root)); - for algorithm in [ - SketchAlgorithm::CmsWithHeap, - SketchAlgorithm::CountSketchWithHeap, - ] { - let candidate = candidates - .iter() - .find_map(|candidate| match &candidate.replacement { - Replacement::SubDAG(node) - if candidate.rationale.contains(&format!("{algorithm:?}")) => - { - Some(node) - } - _ => None, - }) - .unwrap_or_else(|| panic!("missing {algorithm:?} over direct Rate")); - if dynamic { - let (source, ranked) = - asap_executor::physical_planner::promql_rows::compile_rate_ranking(candidate) - .unwrap(); - assert!(matches!( - source.operator, - planner_types::ir::Operator::ASAP( - planner_types::ir::ASAPOp::FinalizeExactAccumulator { .. } - ) - )); - assert_eq!(ranked.input_contracts().count(), 1); - let encoded = String::from_utf8(serde_json::to_vec(&ranked).unwrap()).unwrap(); - assert!(encoded.contains("KeyedSummaryBuild")); - assert!(encoded.contains("KeyedEvaluation")); - assert!( - !encoded.contains("\"Rate\""), - "Rate must be supplied by its exact stored-state evaluation" - ); - } - let dag = compile_physical_asap_dag(candidate).unwrap(); - assert!(dag.nodes.iter().any(|node| matches!(&node.payload, - PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) if kind.algorithm() == &algorithm))); - let build = dag.nodes.iter().find(|node| matches!(&node.payload, - PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) if kind.algorithm() == &algorithm)).unwrap(); - let input_id = dag - .edges - .iter() - .find(|edge| edge.consumer == build.id) - .unwrap() - .producer; - let schema = Arc::new( - dag.nodes - .iter() - .find(|node| node.id == input_id) - .unwrap() - .output_schema - .clone(), - ); - let raw = dag - .nodes - .iter() - .find(|node| { - matches!( - &node.payload, - PhysicalASAPOperatorPayload::NonASAP( - planner_types::ir::NonASAPOp::TimeRange { .. } - ) - ) - }) - .unwrap_or_else(|| panic!("no raw counter source: {dag:?}")); - let raw_schema = Arc::new(raw.output_schema.clone()); - let raw_compiled = compile( - &dag, - BTreeMap::from([(raw.id as u64, InputContract::bounded(raw_schema.clone()))]), - &[dag.roots[0] as u64], - ) - .unwrap(); - let bytes = serde_json::to_vec(&raw_compiled).unwrap(); - let raw_compiled = - serde_json::from_slice::(&bytes) - .unwrap(); - // Each evaluation receives a complete raw window. A reset, a stopped - // series and an expired leader must not retain last run's heap weights. - for (end, series, expected) in [ - ( - 60_000, - vec![ - ("auth", vec![10., 30., 50.]), - ("checkout", vec![10., 50., 90.]), - ("search", vec![10., 70., 130.]), - ], - vec![11. / 6., 8. / 3.], - ), - ( - 120_000, - vec![ - ("auth", vec![100., 10., 50.]), - ("checkout", vec![100., 100., 100.]), - ], - vec![0., 1.25], - ), - ] { - let mut raw_rows = Vec::new(); - for (service, samples) in series { - for (offset, value) in [10_000, 30_000, 50_000].into_iter().zip(samples) { - if dynamic { - raw_rows.push( - series_row( - &raw_schema, - &BTreeMap::from([ - ("job".into(), "api".into()), - ("service".into(), service.into()), - ("unreferenced".into(), format!("{service}-extra")), - ]), - end - 60_000 + offset, - value, - ) - .unwrap(), - ); - continue; - } - raw_rows.push( - raw_schema - .fields - .iter() - .map(|field| match field.name.as_str() { - "service" => Value::Utf8(service.into()), - "job" => Value::Utf8("api".into()), - "value" => Value::Float64(value), - "ts" => Value::Timestamp(end - 60_000 + offset), - _ => panic!("unexpected raw field"), - }) - .collect(), - ); - } - } - let raw_batch = Batch::try_new(raw_schema.clone(), raw_rows).unwrap(); - for scope in [ - Scope::Ingestion { - window_start_ms: end - 60_000, - window_end_ms: end, - revision: 1, - }, - Scope::Query { - evaluation_time_ms: end, - revision: 1, - }, - ] { - let source = Box::new( - Operator::source(raw_schema.clone(), vec![raw_batch.clone()]).unwrap(), - ) as Source<'static>; - let physical_dag = raw_compiled - .instantiate(BTreeMap::from([(raw.id as u64, source)])) - .unwrap(); - let context = RunContext::new(scope, Limits::default()).unwrap(); - let mut raw_scores = block_on(async { - let mut scores = Vec::new(); - let mut stream = physical_dag - .execute(&[dag.roots[0] as u64], context) - .unwrap() - .remove(0); - while let Some(batch) = stream.next().await { - let batch = batch.unwrap(); - for row in batch.rows() { - if dynamic { - let column = batch - .schema() - .fields - .iter() - .position(|field| field.name == SERIES_IDENTITY_COLUMN) - .unwrap(); - let Value::Utf8(encoded) = &row[column] else { - panic!("identity lost"); - }; - let labels = decode_series_identity(encoded).unwrap(); - assert_eq!(labels["job"], "api"); - assert_eq!( - labels["unreferenced"], - format!("{}-extra", labels["service"]) - ); - } - assert!(row.iter().any( - |value| matches!(value, Value::Timestamp(time) if *time == end) - )); - scores.extend(row.iter().filter_map(|value| match value { - Value::Float64(value) => Some(*value), - _ => None, - })); - } - } - scores - }); - raw_scores.sort_by(f64::total_cmp); - assert_eq!(raw_scores.len(), expected.len()); - for (actual, expected) in raw_scores.iter().zip(&expected) { - assert!( - (actual - expected).abs() < 1e-12, - "raw counter semantics must precede heap ranking: {raw_scores:?}" - ); - } - } - } - let compiled = compile( - &dag, - BTreeMap::from([(input_id as u64, InputContract::bounded(schema.clone()))]), - &[dag.roots[0] as u64], - ) - .unwrap(); - for (time, values, expected) in [ - ( - 60_000, - vec![("auth", 3.), ("checkout", 2.), ("search", 1.)], - vec![2., 3.], - ), - ( - 61_000, - vec![("auth", 0.), ("checkout", 2.), ("search", 4.)], - vec![2., 4.], - ), - (62_000, vec![("auth", 0.), ("checkout", 2.)], vec![0., 2.]), - ] { - let rows = values - .into_iter() - .map(|(service, value)| { - if dynamic { - return series_row( - &schema, - &BTreeMap::from([ - ("job".into(), "api".into()), - ("service".into(), service.into()), - ]), - time, - value, - ) - .unwrap(); - } - schema - .fields - .iter() - .map(|field| match field.name.as_str() { - "service" => Value::Utf8(service.into()), - "job" => Value::Utf8("api".into()), - "value" => Value::Float64(value), - "ts" => Value::Timestamp(time), - _ => panic!("unexpected rate field {field:?}"), + let root = with_series_identity(&logical).unwrap(); + let candidates = stage1_candidates(&root); + // Stage 1 proves no sign for rate values, so a Count-Min heap over them + // is not executable; only Count Sketch is. + let algorithm = SketchAlgorithm::CountSketchWithHeap; + // The heap over the Rate realized as an exact accumulator. + let candidate = candidates + .iter() + .find(|candidate| { + let dag = compile_physical_asap_dag(candidate).unwrap(); + builds(&dag, &algorithm) + && dag.nodes.iter().any(|node| { + matches!( + &node.payload, + PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg { + family: FieldDataType::ExactAggregate(ExactKind::Rate, _), + .. }) - .collect() + ) }) - .collect(); - let batch = Batch::try_new(schema.clone(), rows).unwrap(); - for scope in [ - Scope::Query { - evaluation_time_ms: time, - revision: 1, - }, - Scope::Ingestion { - window_start_ms: time - 60_000, - window_end_ms: time, - revision: 1, - }, - ] { - let source = - Box::new(Operator::source(schema.clone(), vec![batch.clone()]).unwrap()) - as Source<'static>; - let physical_dag = compiled - .instantiate(BTreeMap::from([(input_id as u64, source)])) - .unwrap(); - let context = RunContext::new(scope, Limits::default()).unwrap(); - let mut scores = block_on(async { - let mut scores = vec![]; - let mut stream = physical_dag - .execute(&[dag.roots[0] as u64], context) - .unwrap() - .remove(0); - while let Some(batch) = stream.next().await { - let batch = batch.unwrap(); - for row in batch.rows() { - assert!(row.iter().any( - |value| matches!(value, Value::Timestamp(actual) if *actual == time) - )); - scores.push( - row.iter() - .find_map(|value| { - if let Value::Float64(value) = value { - Some(*value) - } else { - None - } - }) - .unwrap(), - ); - } - } - scores - }); - scores.sort_by(f64::total_cmp); - assert_eq!( - scores, expected, - "heap snapshots must not accumulate across evaluations" - ); - } - } - } -} - -// Spatial ranking consumes one eligible instant vector. Signed values require -// CountSketch; a raw metric does not establish the non-negative CMS contract. -#[test] -fn spatial_topk_exposes_signed_heap_candidate_over_complete_snapshot() { - use asap_executor::physical_planner::promql_rows::{ - decode_series_identity, series_row, with_series_identity, SERIES_IDENTITY_COLUMN, - }; - let logical = lower_promql("topk by(job)(1, m)", AccuracyTarget::Epsilon(0.1)).unwrap(); - let root = Rc::new(with_series_identity(&logical).unwrap()); - let strategy = ASAPStrategies::new_with_planning_inputs_and_evidence( - &DefaultAccuracyModel, - &EqualSplitAllocator, - &Evidence, + }) + .unwrap_or_else(|| panic!("missing {algorithm:?} over direct Rate")); + let (source, ranked) = + asap_executor::physical_planner::promql_rows::compile_rate_ranking(candidate).unwrap(); + assert!(matches!( + source.operator, + planner_types::ir::Operator::ASAP( + planner_types::ir::ASAPOp::FinalizeExactAccumulator { .. } + ) + )); + assert_eq!(ranked.input_contracts().count(), 1); + let encoded = String::from_utf8(serde_json::to_vec(&ranked).unwrap()).unwrap(); + assert!(encoded.contains("KeyedSummaryBuild")); + assert!(encoded.contains("KeyedEvaluation")); + assert!( + !encoded.contains("\"Rate\""), + "Rate must be supplied by its exact stored-state evaluation" ); - let candidates = strategy - .current_series_topk_candidates(&root, &AccuracyTarget::Epsilon(0.1)) - .candidates; - assert!(!candidates - .iter() - .any(|c| c.rationale.contains("CmsWithHeap"))); - let selected = candidates + let dag = compile_physical_asap_dag(candidate).unwrap(); + assert!(dag.nodes.iter().any(|node| matches!(&node.payload, + PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) if kind.algorithm() == &algorithm))); + let build = dag.nodes.iter().find(|node| matches!(&node.payload, + PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) if kind.algorithm() == &algorithm)).unwrap(); + let input_id = dag + .edges .iter() - .find_map(|candidate| match &candidate.replacement { - Replacement::SubDAG(node) if candidate.rationale.contains("CountSketchWithHeap") => { - Some(node) - } - _ => None, - }) - .expect("signed spatial TopK must expose CountSketch with heap"); - let dag = compile_physical_asap_dag(selected).unwrap(); + .find(|edge| edge.consumer == build.id) + .unwrap() + .producer; + let schema = Arc::new( + dag.nodes + .iter() + .find(|node| node.id == input_id) + .unwrap() + .output_schema + .clone(), + ); let raw = dag .nodes .iter() @@ -654,93 +292,200 @@ fn spatial_topk_exposes_signed_heap_candidate_over_complete_snapshot() { ) ) }) - .unwrap(); - let schema = Arc::new(raw.output_schema.clone()); - let program = compile( + .unwrap_or_else(|| panic!("no raw counter source: {dag:?}")); + let raw_schema = Arc::new(raw.output_schema.clone()); + let raw_compiled = compile( &dag, - BTreeMap::from([(raw.id as u64, InputContract::bounded(schema.clone()))]), + BTreeMap::from([(raw.id as u64, InputContract::bounded(raw_schema.clone()))]), &[dag.roots[0] as u64], ) .unwrap(); - let snapshot_program = - asap_executor::physical_planner::promql_rows::compile_current_series_evaluation(selected) + let bytes = serde_json::to_vec(&raw_compiled).unwrap(); + let raw_compiled = + serde_json::from_slice::(&bytes) .unwrap(); - let encoded: serde_json::Value = - serde_json::from_slice(&serde_json::to_vec(&snapshot_program).unwrap()).unwrap(); - assert!(!encoded.to_string().contains("CurrentSeries")); - assert!(encoded.to_string().contains("KeyedSummaryBuild")); - assert!(encoded.to_string().contains("KeyedEvaluation")); - for (values, expected, score) in [ - ([100., 20.], "a", 100.), - ([1., 20.], "b", 20.), - ([-10., -2.], "b", -2.), + // Each evaluation receives a complete raw window. A reset, a stopped + // series and an expired leader must not retain last run's heap weights. + for (end, series, expected) in [ + ( + 60_000, + vec![ + ("auth", vec![10., 30., 50.]), + ("checkout", vec![10., 50., 90.]), + ("search", vec![10., 70., 130.]), + ], + vec![11. / 6., 8. / 3.], + ), + ( + 120_000, + vec![ + ("auth", vec![100., 10., 50.]), + ("checkout", vec![100., 100., 100.]), + ], + vec![0., 1.25], + ), + ] { + let mut raw_rows = Vec::new(); + for (service, samples) in series { + for (offset, value) in [10_000, 30_000, 50_000].into_iter().zip(samples) { + raw_rows.push( + series_row( + &raw_schema, + &BTreeMap::from([ + ("job".into(), "api".into()), + ("service".into(), service.into()), + ("unreferenced".into(), format!("{service}-extra")), + ]), + end - 60_000 + offset, + value, + ) + .unwrap(), + ); + } + } + let raw_batch = Batch::try_new(raw_schema.clone(), raw_rows).unwrap(); + for scope in [ + Scope::Ingestion { + window_start_ms: end - 60_000, + window_end_ms: end, + revision: 1, + }, + Scope::Query { + evaluation_time_ms: end, + revision: 1, + }, + ] { + let source = + Box::new(Operator::source(raw_schema.clone(), vec![raw_batch.clone()]).unwrap()) + as Source<'static>; + let physical_dag = raw_compiled + .instantiate(BTreeMap::from([(raw.id as u64, source)])) + .unwrap(); + let context = RunContext::new(scope, Limits::default()).unwrap(); + let mut raw_scores = block_on(async { + let mut scores = Vec::new(); + let mut stream = physical_dag + .execute(&[dag.roots[0] as u64], context) + .unwrap() + .remove(0); + while let Some(batch) = stream.next().await { + let batch = batch.unwrap(); + for row in batch.rows() { + let column = batch + .schema() + .fields + .iter() + .position(|field| field.name == SERIES_IDENTITY_COLUMN) + .unwrap(); + let Value::Utf8(encoded) = &row[column] else { + panic!("identity lost"); + }; + let labels = decode_series_identity(encoded).unwrap(); + assert_eq!(labels["job"], "api"); + assert_eq!( + labels["unreferenced"], + format!("{}-extra", labels["service"]) + ); + scores.extend(row.iter().filter_map(|value| match value { + Value::Float64(value) => Some(*value), + _ => None, + })); + } + } + scores + }); + raw_scores.sort_by(f64::total_cmp); + assert_eq!(raw_scores.len(), expected.len()); + for (actual, expected) in raw_scores.iter().zip(&expected) { + assert!( + (actual - expected).abs() < 1e-12, + "raw counter semantics must precede heap ranking: {raw_scores:?}" + ); + } + } + } + let compiled = compile( + &dag, + BTreeMap::from([(input_id as u64, InputContract::bounded(schema.clone()))]), + &[dag.roots[0] as u64], + ) + .unwrap(); + for (time, values, expected) in [ + ( + 60_000, + vec![("auth", 3.), ("checkout", 2.), ("search", 1.)], + vec![2., 3.], + ), + ( + 61_000, + vec![("auth", 0.), ("checkout", 2.), ("search", 4.)], + vec![2., 4.], + ), + (62_000, vec![("auth", 0.), ("checkout", 2.)], vec![0., 2.]), ] { - let rows = ["a", "b"] + let rows = values .into_iter() - .zip(values) - .map(|(instance, value)| { + .map(|(service, value)| { series_row( &schema, &BTreeMap::from([ ("job".into(), "api".into()), - ("unreferenced".into(), instance.into()), + ("service".into(), service.into()), ]), - 60_000, + time, value, ) .unwrap() }) .collect(); let batch = Batch::try_new(schema.clone(), rows).unwrap(); - let physical_dag = program - .instantiate(BTreeMap::from([( - raw.id as u64, - Box::new(Operator::source(schema.clone(), vec![batch]).unwrap()) as Source<'_>, - )])) - .unwrap(); - block_on(async { - let context = RunContext::new( - Scope::Query { - evaluation_time_ms: 60_000, - revision: 0, - }, - Limits::default(), - ) - .unwrap(); - let mut stream = physical_dag - .execute(program.roots(), context) - .unwrap() - .remove(0); - let mut result = Vec::new(); - while let Some(batch) = stream.next().await { - let batch = batch.unwrap(); - let identity = batch - .schema() - .fields - .iter() - .position(|f| f.name == SERIES_IDENTITY_COLUMN) - .unwrap(); - let value = batch - .schema() - .fields - .iter() - .position(|f| f.name == "value") - .unwrap(); - for row in batch.rows() { - let Value::Utf8(labels) = &row[identity] else { - panic!() - }; - let Value::Float64(v) = row[value] else { - panic!() - }; - result.push(( - decode_series_identity(labels).unwrap()["unreferenced"].clone(), - v, - )); + for scope in [ + Scope::Query { + evaluation_time_ms: time, + revision: 1, + }, + Scope::Ingestion { + window_start_ms: time - 60_000, + window_end_ms: time, + revision: 1, + }, + ] { + let source = Box::new(Operator::source(schema.clone(), vec![batch.clone()]).unwrap()) + as Source<'static>; + let physical_dag = compiled + .instantiate(BTreeMap::from([(input_id as u64, source)])) + .unwrap(); + let context = RunContext::new(scope, Limits::default()).unwrap(); + let mut scores = block_on(async { + let mut scores = vec![]; + let mut stream = physical_dag + .execute(&[dag.roots[0] as u64], context) + .unwrap() + .remove(0); + while let Some(batch) = stream.next().await { + let batch = batch.unwrap(); + for row in batch.rows() { + scores.push( + row.iter() + .find_map(|value| { + if let Value::Float64(value) = value { + Some(*value) + } else { + None + } + }) + .unwrap(), + ); + } } - } - assert_eq!(result, vec![(expected.into(), score)]); - }); + scores + }); + scores.sort_by(f64::total_cmp); + assert_eq!( + scores, expected, + "heap snapshots must not accumulate across evaluations" + ); + } } } @@ -806,20 +551,25 @@ fn maintained_rate_heap_compiles_fixed_window_precompute() { ) .unwrap(), ); - let strategy = ASAPStrategies::new_with_planning_inputs_and_evidence( - &DefaultAccuracyModel, - &EqualSplitAllocator, - &Evidence, - ); - let candidates = strategy - .replacements(&TargetSubDAG::new(&root)) + // The heap over the Rate realized as an exact accumulator. Stage 1 proves + // no sign for rate values, so only the Count Sketch heap is executable. + let candidates = stage1_candidates(&root) .into_iter() - .filter_map(|candidate| match candidate.replacement { - Replacement::SubDAG(root) if candidate.rationale.contains("WithHeap") => Some(root), - _ => None, + .filter(|candidate| { + let dag = compile_physical_asap_dag(candidate).unwrap(); + builds(&dag, &SketchAlgorithm::CountSketchWithHeap) + && dag.nodes.iter().any(|node| { + matches!( + &node.payload, + PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg { + family: FieldDataType::ExactAggregate(ExactKind::Rate, _), + .. + }) + ) + }) }) .collect::>(); - assert_eq!(candidates.len(), 2); + assert_eq!(candidates.len(), 1); for root in candidates { let dag = continuously_maintained_dag(&root); let state = dag diff --git a/crates/frontend-promql/tests/count_planning.rs b/crates/frontend-promql/tests/count_planning.rs index 1972b43c2..3832d5d03 100644 --- a/crates/frontend-promql/tests/count_planning.rs +++ b/crates/frontend-promql/tests/count_planning.rs @@ -1,9 +1,4 @@ //! Query text through summary selection: counts use observations, never value weights. -use asap_logical_optimizer::accuracy::DefaultAccuracyModel; -use asap_logical_optimizer::{ - default_strategies, search_workload_with_targets, ASAPStrategies, Replacement, - ReplacementStrategy, TargetSubDAG, -}; mod support; use asap_types::ir::physical_export::PhysicalASAPOperatorPayload; use asap_types::ir::schema::{ @@ -13,43 +8,36 @@ use asap_types::ir::schema::{ use asap_types::ir::{ASAPOp, Operator, OperatorNode}; use asap_types::types::AccuracyTarget; use std::rc::Rc; -use support::{lower_promql, post_asap_dag, selected_dag}; +use support::{lower_promql, post_asap_dag, selected_dag, stage1_candidates}; #[test] -fn grouped_count_keeps_uncertified_hydra_candidates_for_backend_review() { +fn grouped_count_offers_hydra_but_selection_keeps_per_group_state() { let target = AccuracyTarget::EpsilonDelta { epsilon: 0.01, delta: 0.01, }; - let root = lower_promql("count by(job)(up)", target.clone()).unwrap(); - let space = search_workload_with_targets( - vec![("count", Rc::clone(&root), Some(target.clone()))], - &default_strategies(), - &DefaultAccuracyModel, - ); - let planned = &space.roots[0].1; - let hydra: Vec<_> = space - .candidates_for_target(planned) - .unwrap() - .candidates - .iter() - .filter(|candidate| candidate.strategy == "HydraGroupingStrategy") - .collect(); - assert_eq!(hydra.len(), 2); - assert!(hydra - .iter() - .all(|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 { .. }), - .. - }) - ))); + // Hydra hashes a non-null item per row: the series identity (PromQL + // labels are nullable). + let root = asap_types::ir::schema_support::with_promql_series_identity( + &lower_promql("count by(job)(up)", target.clone()).unwrap(), + ) + .unwrap(); + let hydra = |node: &Rc| { + OperatorNode::reachable(node).iter().any(|node| { + matches!( + &node.operator, + Operator::ASAP(ASAPOp::SummaryAgg { + family: FieldDataType::Sketch( + _, + GroupingStrategy::SharedMultiSubpopulation { .. } + ), + .. + }) + ) + }) + }; + assert!(stage1_candidates(&root).iter().any(hydra)); + assert!(!hydra(&selected_dag(root, target))); } // Exact series and temporal counts must select a count accumulator, not distinct or sum. @@ -57,12 +45,18 @@ fn grouped_count_keeps_uncertified_hydra_candidates_for_backend_review() { fn exact_counts_select_count_accumulators() { for query in ["count(up)", "count by(job)(up)", "count_over_time(up[5m])"] { let root = lower_promql(query, AccuracyTarget::Exact).unwrap(); - let candidates = ASAPStrategies::default().replacements(&TargetSubDAG::new(&root)); + let candidates = stage1_candidates(&root); assert!( candidates.iter().any(|candidate| { - matches!(&candidate.replacement, Replacement::SubDAG(node) - if matches!(&node.operator, Operator::ASAP(ASAPOp::SummaryAgg { - family: FieldDataType::ExactAggregate(ExactKind::Count, _), .. }))) + OperatorNode::reachable(candidate).iter().any(|node| { + matches!( + &node.operator, + Operator::ASAP(ASAPOp::SummaryAgg { + family: FieldDataType::ExactAggregate(ExactKind::Count, _), + .. + }) + ) + }) }), "{query}: {candidates:?}" ); @@ -74,13 +68,18 @@ fn exact_counts_select_count_accumulators() { fn frequency_count_candidates_use_unit_weights() { for query in ["count_over_time(up[5m])", "count(up)"] { let root = lower_promql(query, AccuracyTarget::Epsilon(0.02)).unwrap(); - let candidates = ASAPStrategies::default().replacements(&TargetSubDAG::new(&root)); + let candidates = stage1_candidates(&root); let mut algorithms = Vec::new(); - for candidate in &candidates { - let Replacement::SubDAG(node) = &candidate.replacement else { - continue; - }; - let Operator::ASAP(ASAPOp::SummaryEstimate { summary_input, .. }) = &node.operator + for node in &candidates { + let Some(summary_input) = + OperatorNode::reachable(node) + .into_iter() + .find_map(|node| match &node.operator { + Operator::ASAP(ASAPOp::SummaryEstimate { summary_input, .. }) => { + Some(summary_input.clone()) + } + _ => None, + }) else { continue; }; @@ -231,13 +230,9 @@ fn count_over_time_counts_scrapes_not_sample_values() { fn cms_count_updates_total_ten_for_zero_positive_and_negative_samples() { use asap_types::ir::scalar::ColumnRef; let root = lower_promql("count_over_time(up[5m])", AccuracyTarget::Epsilon(0.02)).unwrap(); - let candidates = ASAPStrategies::default().replacements(&TargetSubDAG::new(&root)); - let dag = candidates + let dag = stage1_candidates(&root) .iter() - .find_map(|candidate| { - let Replacement::SubDAG(node) = &candidate.replacement else { - return None; - }; + .find_map(|node| { let dag = post_asap_dag(node); dag.nodes .iter() diff --git a/crates/frontend-promql/tests/observability/metrics_observability.rs b/crates/frontend-promql/tests/observability/metrics_observability.rs index 24178c20f..8e2fb9fe2 100644 --- a/crates/frontend-promql/tests/observability/metrics_observability.rs +++ b/crates/frontend-promql/tests/observability/metrics_observability.rs @@ -4,18 +4,11 @@ //! Prometheus rule YAML. The preprocessing script produces the line-oriented //! fixtures consumed here. -use std::rc::Rc; - use asap_frontend_promql::PromqlError; -use asap_logical_optimizer::pass1::replacement::{retain_exact, RealizationError}; -use asap_logical_optimizer::{ - ASAPStrategies, Replacement, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, -}; #[path = "../support.rs"] mod support; -use asap_types::ir::OperatorNode; use asap_types::types::AccuracyTarget; -use support::lower_promql; +use support::{lower_promql, stage1_plan}; const CORPORA: &[(&str, &str)] = &[ ( @@ -63,21 +56,6 @@ fn queries(corpus: &str) -> impl Iterator { .filter(|line| !line.is_empty() && !line.starts_with('#')) } -fn post_asap_candidate(root: &Rc) -> Result, RealizationError> { - let target = TargetSubDAG::new(root); - match ASAPStrategies::default() - .replacements(&target) - .into_iter() - .next() - { - Some(ReplacementSubDAG { - replacement: Replacement::SubDAG(node), - .. - }) => Ok(node), - _ => retain_exact(root), - } -} - #[test] fn benchmark_corpora_are_total_and_report_coverage() { for (name, corpus) in CORPORA { @@ -91,12 +69,12 @@ fn benchmark_corpora_are_total_and_report_coverage() { for query in queries(corpus) { total += 1; - // Approximate accuracy exercises the sketch-replacement boundary - // used by the existing post-ASAP corpus binding test. + // Approximate accuracy makes Stage 1 offer sketch alternatives, + // as in the PromQL corpus test. match lower_promql(query, AccuracyTarget::Epsilon(0.01)) { Ok(expr) => { lowered += 1; - match post_asap_candidate(&expr) { + match stage1_plan(&expr) { // An ASAP operator bound somewhere below the root. Ok(node) if node.contains_asap() => post_asap_candidates += 1, Ok(_) => { diff --git a/crates/frontend-promql/tests/observability/promql_corpus.rs b/crates/frontend-promql/tests/observability/promql_corpus.rs index 606b2dbab..44a29726f 100644 --- a/crates/frontend-promql/tests/observability/promql_corpus.rs +++ b/crates/frontend-promql/tests/observability/promql_corpus.rs @@ -13,39 +13,11 @@ //! that is the guarantee. A coverage floor guards against a change silently //! tanking how much of the corpus we can lower. -use std::rc::Rc; - use asap_frontend_promql::PromqlError as LoweringError; -use asap_logical_optimizer::pass1::replacement::{retain_exact, RealizationError}; -use asap_logical_optimizer::{ - ASAPStrategies, Replacement, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, -}; #[path = "../support.rs"] mod support; -use asap_types::ir::OperatorNode; use asap_types::types::AccuracyTarget; -use support::lower_promql; - -/// This crate has no "bind me one dag" public API any more — -/// `ASAPStrategies::replacements` always returns every candidate, and -/// a caller decides what to keep. This test-only helper reproduces the -/// take-the-first-(`cost_model`-preferred)-candidate pattern so [`bind_tally`] -/// gets one representative `Result` per query, matching what a totality -/// check over the whole corpus wants. -fn bind(root: &Rc) -> Result, RealizationError> { - let target = TargetSubDAG::new(root); - match ASAPStrategies::default() - .replacements(&target) - .into_iter() - .next() - { - Some(ReplacementSubDAG { - replacement: Replacement::SubDAG(node), - .. - }) => Ok(node), - _ => retain_exact(root), - } -} +use support::{lower_promql, stage1_plan}; const DOCS: &str = include_str!("data/promql_corpus_docs.txt"); const TESTDATA: &str = include_str!("data/promql_corpus_testdata.txt"); @@ -87,11 +59,9 @@ fn tally(corpus: &str) -> Tally { t } -/// Every query that lowers, additionally run through the pre-ASAP → -/// post-ASAP `asap-logical-optimizer` binding pass (issue #98), at an -/// approximate accuracy target so the sketch-selection boundary actually -/// fires (an `Exact` target would only ever exercise the exact-accumulator -/// arm). +/// Every query that lowers, additionally run through Stage 1 ([`stage1_plan`]) +/// at an approximate accuracy target so sketch alternatives are offered (an +/// `Exact` target would only ever exercise the exact-accumulator arm). #[derive(Default, Debug)] struct BindTally { /// An ASAP operator was bound somewhere below the root — the pass did @@ -99,7 +69,7 @@ struct BindTally { transformed: usize, /// The kept pre-ASAP dag — the pass left the query untouched. unchanged: usize, - /// [`bind`] returned `Err` (schema derivation failed). + /// [`stage1_plan`] returned `Err`. errored: usize, } @@ -109,7 +79,7 @@ fn bind_tally(corpus: &str, accuracy: AccuracyTarget) -> BindTally { let Ok(dag) = lower_promql(q, accuracy.clone()) else { continue; }; - match bind(&dag) { + match stage1_plan(&dag) { Ok(bound) if !bound.contains_asap() => t.unchanged += 1, Ok(_) => t.transformed += 1, Err(_) => t.errored += 1, @@ -126,7 +96,7 @@ fn binding_is_total_over_the_entire_corpus() { eprintln!("docs corpus post-ASAP binding: {docs:?}"); eprintln!("testdata corpus post-ASAP binding: {td:?}"); - // Same totality guarantee as lowering: reaching here means `bind` + // Same totality guarantee as lowering: reaching here means Stage 1 // never panicked over any lowerable query in the corpus. assert_eq!(docs.errored, 0, "post-ASAP binding errored: {docs:?}"); assert_eq!(td.errored, 0, "post-ASAP binding errored: {td:?}"); diff --git a/crates/frontend-promql/tests/support.rs b/crates/frontend-promql/tests/support.rs index 618d86086..244f47204 100644 --- a/crates/frontend-promql/tests/support.rs +++ b/crates/frontend-promql/tests/support.rs @@ -138,3 +138,67 @@ pub fn selected_dag(root: Rc, accuracy: AccuracyTarget) -> Rc) -> Vec> { + use asap_logical_optimizer::pass1::logical_candidates::{ + compose_logical_candidate, enumerate_choices, enumerate_local_logical_candidates, + }; + use asap_types::ir::QueryRoot; + let inventory = enumerate_local_logical_candidates( + vec![(0, QueryRoot::Operator(Rc::clone(root)))], + &Default::default(), + ) + .unwrap(); + enumerate_choices(&inventory, 4096) + .iter() + .filter_map(|choice| { + match compose_logical_candidate(&inventory, choice) + .ok()? + .remove(0) + .1 + { + QueryRoot::Operator(node) => Some(node), + QueryRoot::Scalar(_) => None, + } + }) + .collect() +} + +/// A candidate Stage 1 composes for `root`: the first of its (up to 64) +/// whole-query choices that composes and builds a summary, else the +/// pass-through choice. `Err` when Pass 1 rejects the query or the +/// pass-through choice does not compose. A choice may legitimately not +/// compose (an alternative its input cannot feed); selection skips it. +#[allow(dead_code)] +pub fn stage1_plan( + root: &Rc, +) -> Result< + Rc, + asap_logical_optimizer::pass1::logical_candidates::LogicalCandidateError, +> { + use asap_logical_optimizer::pass1::logical_candidates::{ + compose_logical_candidate, enumerate_choices, enumerate_local_logical_candidates, + }; + use asap_types::ir::QueryRoot; + let inventory = enumerate_local_logical_candidates( + vec![(0, QueryRoot::Operator(Rc::clone(root)))], + &Default::default(), + )?; + let compose = |choice: &[usize]| { + compose_logical_candidate(&inventory, choice).map(|mut roots| match roots.remove(0).1 { + QueryRoot::Operator(node) => node, + QueryRoot::Scalar(_) => unreachable!("an operator root composes to an operator root"), + }) + }; + let summarized = enumerate_choices(&inventory, 64) + .iter() + .filter_map(|choice| compose(choice).ok()) + .find(|node| node.contains_asap()); + match summarized { + Some(node) => Ok(node), + None => compose(&vec![0; inventory.targets.len()]), + } +} diff --git a/crates/frontend-promql/tests/univmon_candidates.rs b/crates/frontend-promql/tests/univmon_candidates.rs index 3bc4791db..6b16667c4 100644 --- a/crates/frontend-promql/tests/univmon_candidates.rs +++ b/crates/frontend-promql/tests/univmon_candidates.rs @@ -1,67 +1,29 @@ use std::rc::Rc; -use asap_logical_optimizer::accuracy::{ - AccuracyModel, DefaultAccuracyModel, EqualSplitAllocator, PropagationStats, -}; -use asap_logical_optimizer::pass1::replacement::{ - default_strategies, search_workload_with_targets, -}; -use asap_logical_optimizer::{ASAPStrategies, Replacement, ReplacementStrategy, TargetSubDAG}; +use asap_logical_optimizer::accuracy::{AccuracyModel, DefaultAccuracyModel}; mod support; use asap_types::ir::cse::share_common_sub_dags; -use asap_types::ir::properties::{ - AccuracyError, BoundExpr, CompositionOperator, ErrorMetric, ProbabilityExpr, ResultGuarantee, -}; +use asap_types::ir::properties::ErrorMetric; 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, selected_dag}; +use support::{lower_promql, post_asap_dag, selected_dag, stage1_candidates}; -// Synthetic evidence exercises structural sharing, never runtime accuracy. -struct TestEvidence; -impl AccuracyModel for TestEvidence { - fn local_guarantee( - &self, - family: &FieldDataType, - query: &SketchStatistic, - ) -> Option { - if matches!(family, FieldDataType::Sketch(kind, _) if kind.algorithm() == &SketchAlgorithm::UnivMon) - && !matches!(query, SketchStatistic::PointCount { .. }) - { - let mut guarantee = ResultGuarantee::exact("SYNTHETIC test evidence; not measured"); - guarantee.metric = ErrorMetric::RelativeValue; - guarantee.bound = BoundExpr::Constant { value: 0.01 }; - guarantee.failure_probability = ProbabilityExpr::Constant { value: 0.01 }; - Some(guarantee) - } else { - DefaultAccuracyModel.local_guarantee(family, query) - } - } - fn propagate( - &self, - op: &CompositionOperator, - inputs: &[ResultGuarantee], - local: Option<&ResultGuarantee>, - stats: &PropagationStats, - ) -> Result { - DefaultAccuracyModel.propagate(op, inputs, local, stats) - } - fn satisfies(&self, guarantee: &ResultGuarantee, target: &AccuracyTarget) -> bool { - DefaultAccuracyModel.satisfies(guarantee, target) - } +fn reads_univmon(node: &OperatorNode) -> bool { + let Operator::ASAP(ASAPOp::SummaryEstimate { summary_input, .. }) = &node.operator else { + return false; + }; + matches!(&summary_input.operator, Operator::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) + if kind.algorithm() == &SketchAlgorithm::UnivMon) } +/// Stage 1's candidate for `query` whose root reads a UnivMon summary. fn candidate(query: &str, accuracy: AccuracyTarget) -> Rc { let root = lower_promql(query, accuracy).unwrap(); - ASAPStrategies::new_with_planning_inputs(&TestEvidence, &EqualSplitAllocator) - .replacements(&TargetSubDAG::new(&root)) + stage1_candidates(&root) .into_iter() - .find_map(|candidate| { - let Replacement::SubDAG(node) = candidate.replacement else { return None }; - let Operator::ASAP(ASAPOp::SummaryEstimate { summary_input, .. }) = &node.operator else { return None }; - matches!(&summary_input.operator, Operator::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) - if kind.algorithm() == &SketchAlgorithm::UnivMon).then_some(node) - }).expect("UnivMon candidate") + .find(|node| reads_univmon(node)) + .expect("UnivMon candidate") } #[test] @@ -70,14 +32,14 @@ fn four_evaluations_share_one_value_frequency_state_and_keep_honest_guarantees() // independently of evaluation: UnivMon is sized for L2 whatever it reads. let accuracy = AccuracyTarget::Epsilon(0.02); let roots: Vec<_> = [ - ("distinct_over_time(m[5m])", accuracy.clone()), - ("count_over_time(m[5m])", accuracy.clone()), - ("l2_over_time(m[5m])", accuracy.clone()), - ("entropy_over_time(m[5m])", accuracy), + "distinct_over_time(m[5m])", + "count_over_time(m[5m])", + "l2_over_time(m[5m])", + "entropy_over_time(m[5m])", ] .into_iter() .enumerate() - .map(|(id, (query, accuracy))| (id, candidate(query, accuracy))) + .map(|(id, query)| (id, candidate(query, accuracy.clone()))) .collect(); let roots = share_common_sub_dags(roots); let mut first_state = None; @@ -97,7 +59,8 @@ fn four_evaluations_share_one_value_frequency_state_and_keep_honest_guarantees() } else { first_state = Some(Rc::clone(summary_input)); } - let Operator::ASAP(ASAPOp::SummaryAgg { input, .. }) = &summary_input.operator else { + let Operator::ASAP(ASAPOp::SummaryAgg { family, input, .. }) = &summary_input.operator + else { panic!() }; assert!(matches!(input.item, Some(SummaryInputExpr::Column(_)))); @@ -107,12 +70,7 @@ fn four_evaluations_share_one_value_frequency_state_and_keep_honest_guarantees() query, SketchStatistic::PointCount { value: None, .. } )); - assert!(root.guarantee.as_ref().is_some_and(|g| g.is_exact())); } else { - assert!(!root.guarantee.as_ref().unwrap().is_exact()); - let Operator::ASAP(ASAPOp::SummaryAgg { family, .. }) = &summary_input.operator else { - panic!() - }; // Production certifies L2 from layer 0's F₂, but has no // calibrated bound for distinct count or entropy. assert_eq!( @@ -131,7 +89,7 @@ fn four_evaluations_share_one_value_frequency_state_and_keep_honest_guarantees() fn uncalibrated_frequency_evaluations_do_not_bypass_accuracy_targets() { // An unmeasured heuristic remains inspectable but is never certified or // automatically selected for a caller-visible bounded-error result. - for query in ["entropy_over_time(m[5m])"] { + for query in ["entropy_over_time(m[5m])", "l2_over_time(m[5m])"] { for target in [ AccuracyTarget::Exact, AccuracyTarget::Epsilon(0.02), @@ -141,34 +99,13 @@ fn uncalibrated_frequency_evaluations_do_not_bypass_accuracy_targets() { }, ] { let root = lower_promql(query, target.clone()).unwrap(); - let candidates = ASAPStrategies::default().replacements(&TargetSubDAG::new(&root)); - let unknown = candidates + let offered = stage1_candidates(&root) .iter() - .filter(|candidate| { - matches!( - &candidate.replacement, - Replacement::SubDAG(node) - if matches!(&node.operator, Operator::ASAP(ASAPOp::SummaryEstimate { .. })) - && node.guarantee.is_none() - && candidate.has_missing_accuracy_evidence() - ) - }) - .count(); + .any(|node| reads_univmon(node)); if target == AccuracyTarget::Exact { - assert_eq!(unknown, 0); + assert!(!offered); } else { - assert!(unknown > 0); - let space = search_workload_with_targets( - vec![("q", Rc::clone(&root), Some(target.clone()))], - &default_strategies(), - &DefaultAccuracyModel, - ); - assert!(space - .candidates_for_target(&space.roots[0].1) - .unwrap() - .candidates - .iter() - .any(|candidate| candidate.has_missing_accuracy_evidence())); + assert!(offered); // Stage 3 never selects the uncalibrated UnivMon estimate. let selected = selected_dag(root, target); assert!(!OperatorNode::reachable(&selected) diff --git a/crates/frontend-sql/tests/pearson_corr.rs b/crates/frontend-sql/tests/pearson_corr.rs index effc7cb0a..0765931cc 100644 --- a/crates/frontend-sql/tests/pearson_corr.rs +++ b/crates/frontend-sql/tests/pearson_corr.rs @@ -148,8 +148,29 @@ async fn corr_filter_is_a_measure_filter() { #[tokio::test] async fn corr_survives_exact_plan_compilation() { let query = lower("SELECT corr(x, y) AS r FROM a").await; - let plan = asap_logical_optimizer::pass1::replacement::retain_exact(&query).unwrap(); - assert!(plan.guarantee.as_ref().unwrap().is_exact()); + use asap_logical_optimizer::pass1::logical_candidates::{ + compose_logical_candidate, enumerate_local_logical_candidates, + }; + use asap_logical_optimizer::Realization; + use asap_types::ir::QueryRoot; + let inventory = enumerate_local_logical_candidates( + vec![(0, QueryRoot::Operator(Rc::clone(&query)))], + &Default::default(), + ) + .unwrap(); + // corr has no summary realization: Stage 1 offers only pass-through. + assert!(inventory + .targets + .iter() + .all(|target| target.alternatives == [Realization::PassThrough])); + let choice = vec![0; inventory.targets.len()]; + let QueryRoot::Operator(plan) = compose_logical_candidate(&inventory, &choice) + .unwrap() + .remove(0) + .1 + else { + panic!("operator root") + }; // The exact fallback is the query's own operator DAG, no ASAP node added. assert!(!plan.contains_asap(), "expected exact fallback"); assert_eq!(aggregate(&plan).0, aggregate(&query).0); diff --git a/crates/integration-tests/tests/cse.rs b/crates/integration-tests/tests/cse.rs deleted file mode 100644 index e1b35029c..000000000 --- a/crates/integration-tests/tests/cse.rs +++ /dev/null @@ -1,178 +0,0 @@ -//! End-to-end pre-ASAP CSE → workload-wide search pin (issue #212, #222, -//! #223). -//! -//! Drives the full staged pipeline this issue lands: two independently -//! lowered `OperatorNode` DAGs → `share_common_sub_dags` (stage 1, -//! `asap-types::ir::cse`, run internally by `search_workload`) → -//! `search_workload` (stage 2, `asap-logical-optimizer`) — and asserts the -//! sharing that stage 1 decides survives into stage 2's discovered -//! `CandidateLogicalASAPDAGs` as one genuinely shared `TargetSubDAGCandidates`, not just one shared -//! `Rc`. This is the "real caller" the issue's landing plan -//! requires before `share_common_sub_dags` is allowed to exist at all (its -//! predecessor, `asap-plan::cse::dedupe_sub-DAGs`, was deleted in #192 for -//! being unwired dead code). -//! -//! Committing to one final, physically-materialized answer for a whole -//! workload (the former `implement_workload`/`implement_workload_with`, -//! which this test file used to drive instead of `search_workload`) is out -//! of `asap-logical-optimizer`'s scope — see that crate's `lib.rs` `## Status` -//! section — so these tests assert on the discovered `CandidateLogicalASAPDAGs` shape -//! directly, the same way `asap-logical-optimizer::pass1::replacement`'s own -//! `shared_aggregate_across_two_roots_gets_both_strategies_candidates` test -//! does, just exercised through the crate's public API from this external -//! integration-test crate. - -use std::rc::Rc; - -use asap_integration_tests::fixtures::lower_promql; -use asap_logical_optimizer::{is_logical_rewrite, search_workload, Replacement}; -use asap_types::ir::NonASAPOp; -use asap_types::types::AccuracyTarget; - -/// Two workload entries that happen to submit the exact same query (a -/// realistic case — two dashboards, or a query fired both standalone and as -/// part of a larger batch) collapse onto one shared `Rc` after -/// `search_workload`'s internal `share_common_sub_dags` pass, and onto one -/// genuinely-shared [`TargetSubDAGCandidates`](asap_logical_optimizer::TargetSubDAGCandidates) — carrying -/// every candidate discovered for it exactly once, not once per root — no -/// second structural-equality pass at the post-ASAP layer needed for this -/// kind of sharing. -#[test] -fn duplicate_workload_queries_collapse_onto_one_memo_group() { - // Grouped (`by (job)`), so the shared `Aggregate`'s output schema carries - // a provable unique key — the legality gate `share_common_sub_dags` - // enforces (see `asap-types::ir::cse`'s module doc) — and its - // `ExactAggregate(Sum)` realization is deterministic regardless of the - // accuracy target, so this pins the sharing mechanism itself rather than - // any one particular summary-family choice. - let query = "sum by (job) (http_requests_total)"; - let a = lower_promql(query, AccuracyTarget::Exact).expect("query a failed to lower"); - let b = lower_promql(query, AccuracyTarget::Exact).expect("query b failed to lower"); - - // Independently lowered: not yet sharing any `Rc`, even though they are - // structurally identical (`resolve_root` gives each call its own fresh - // DAG). - assert_eq!( - a, b, - "fixture sanity: identical query text lowers identically" - ); - - let space = search_workload(vec![("a", a), ("b", b)]); - - // roots[0] and roots[1] must have merged onto the same Rc — the - // `share_common_sub_dags` pass `search_workload` runs internally. - assert!( - Rc::ptr_eq(&space.roots[0].1, &space.roots[1].1), - "search_workload must collapse the two identical roots onto one Rc" - ); - - // The single shared root is one discovered TargetSubDAG, holding one - // TargetSubDAGCandidates with consumer_count 2 — ASAPStrategies's one - // ExactAggregate candidate *and* SharedSubDAGStrategy's share-vs- - // recompute pair, exactly as `shared_aggregate_across_two_roots_gets_both_strategies_candidates` - // (asap-logical-optimizer::pass1::replacement's own equivalent, internal test) - // pins for the same fixture shape. - let group = space - .candidates_for_target(&space.roots[0].1) - .expect("shared root must be a discovered target"); - assert_eq!(group.consumer_count, 2); - assert_eq!( - group.candidates.len(), - 3, - "1 ExactAggregate summary + 2 logical rewrites (share/recompute): {:?}", - group.candidates - ); - - // A bound summary is a `Subtree` with an ASAP operator in it; a logical - // rewrite is a `Subtree` with none (`is_logical_rewrite`). - let summary_count = group - .candidates - .iter() - .filter(|c| matches!(&c.replacement, Replacement::SubDAG(n) if n.contains_asap())) - .count(); - let rewrite_count = group - .candidates - .iter() - .filter(|c| matches!(&c.replacement, Replacement::SubDAG(n) if is_logical_rewrite(n))) - .count(); - assert_eq!(summary_count, 1); - assert_eq!(rewrite_count, 2); - - // The two Rewrite candidates must NOT have collapsed into one (the - // "false-positive dedup" failure mode `is_duplicate_rewrite` exists to - // prevent): one shares the group's own target `Rc`, the other is a - // structurally-identical but independently-built `Rc`. - let one_is_the_target = group.candidates.iter().any(|c| { - matches!(&c.replacement, Replacement::SubDAG(rc) - if is_logical_rewrite(rc) && Rc::ptr_eq(rc, &group.target)) - }); - let one_is_not = group.candidates.iter().any(|c| { - matches!(&c.replacement, Replacement::SubDAG(rc) - if is_logical_rewrite(rc) && !Rc::ptr_eq(rc, &group.target)) - }); - assert!(one_is_the_target && one_is_not); -} - -/// The negative control: two workload entries whose queries are NOT -/// structurally identical must not be conflated — `search_workload` -/// discovers two independent roots, each its own `TargetSubDAG` with its -/// own `TargetSubDAGCandidates` and `consumer_count == 1`. -#[test] -fn distinct_workload_queries_get_independent_memo_groups() { - let a = lower_promql("sum by (job) (http_requests_total)", AccuracyTarget::Exact) - .expect("query a failed to lower"); - let b = lower_promql("sum by (job) (http_response_total)", AccuracyTarget::Exact) - .expect("query b failed to lower"); - assert_ne!(a, b, "fixture sanity: the two queries differ"); - - let space = search_workload(vec![("a", a), ("b", b)]); - assert!(!Rc::ptr_eq(&space.roots[0].1, &space.roots[1].1)); - - let group_a = space - .candidates_for_target(&space.roots[0].1) - .expect("root a must be a discovered target"); - let group_b = space - .candidates_for_target(&space.roots[1].1) - .expect("root b must be a discovered target"); - assert!( - !Rc::ptr_eq(&group_a.target, &group_b.target), - "distinct queries must have distinct TargetSubDAGCandidates entries" - ); - assert_eq!(group_a.consumer_count, 1); - assert_eq!(group_b.consumer_count, 1); -} - -/// Single-query CSE (a repeated sub-expression within one query) also -/// survives through `search_workload`: the two grouped-`Aggregate` branches -/// of a `BinaryOp` collapse to one shared `Rc` in the internal -/// `share_common_sub_dags` pass, and to one shared `TargetSubDAGCandidates` (with -/// `consumer_count == 2`, one per branch) here. -#[test] -fn single_query_repeated_subexpression_shares_one_memo_group() { - let query = "sum by (job) (http_requests_total) / sum by (job) (http_requests_total)"; - let expr = lower_promql(query, AccuracyTarget::Exact).expect("query failed to lower"); - - let space = search_workload(vec![("q", expr)]); - let [(_, root)] = space.roots.as_slice() else { - panic!("expected 1 root"); - }; - let Some(NonASAPOp::BinaryOp { lhs, rhs, .. }) = root.non_asap() else { - panic!("expected a BinaryOp root, got {root:?}"); - }; - assert!( - Rc::ptr_eq(lhs, rhs), - "the two identical sum-by-job branches must collapse onto one Rc" - ); - - let group = space - .candidates_for_target(lhs) - .expect("the shared branch must be a discovered target"); - assert_eq!( - group.consumer_count, 2, - "the shared branch is referenced from both BinaryOp operand positions" - ); - assert!( - !group.candidates.is_empty(), - "the shared branch must carry at least one candidate" - ); -} diff --git a/crates/integration-tests/tests/precompute_raw_samples.rs b/crates/integration-tests/tests/precompute_raw_samples.rs index bd164abcb..5aa0c051d 100644 --- a/crates/integration-tests/tests/precompute_raw_samples.rs +++ b/crates/integration-tests/tests/precompute_raw_samples.rs @@ -3,8 +3,7 @@ mod executor_models; mod physical_common; use asap_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload}; -use asap_types::ir::ASAPOp; -use asap_types::ir::OperatorNode; +use asap_types::ir::{ASAPOp, OperatorNode, QueryRoot}; use physical_common::compile_maintained_physical_asap_dag; use std::{collections::BTreeMap, collections::BTreeSet, rc::Rc, sync::Arc}; @@ -18,14 +17,14 @@ use asap_executor::{ AggregateCore, KeyByLabelValues, Statistic, }; use asap_integration_tests::fixtures::lower_promql; -use asap_logical_optimizer::{ - ASAPStrategies, Replacement, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, +use asap_logical_optimizer::pass1::logical_candidates::{ + compose_logical_candidate, enumerate_choices, enumerate_local_logical_candidates, }; use asap_types::ir::operator::Reduction; use asap_types::ir::scalar::ColumnRef; use asap_types::ir::schema::{ EntityIdentity, ExactKind, FieldDataType, SketchAlgorithm, SketchStatistic, SummaryInputExpr, - SummaryUpdate, + SummaryUpdate, WeightDomain, }; use asap_types::types::AccuracyTarget; use futures::{executor::block_on, StreamExt}; @@ -54,18 +53,36 @@ fn canonical(labels: &Series) -> Series { } /// Every Planner candidate for `query`: the stage pipeline's selection plus -/// each summary replacement of the root. +/// each candidate Stage 1 composes that Stage 3 could admit (Count-Min only +/// over weights proven non-negative). fn candidates(query: &str, accuracy: AccuracyTarget) -> Vec> { let root = lower_promql(query, accuracy.clone()).expect("lowering failed"); - let mut result = ASAPStrategies::default() - .replacements(&TargetSubDAG::new(&root)) - .into_iter() - .filter_map(|candidate| match candidate { - ReplacementSubDAG { - replacement: Replacement::SubDAG(node), - .. - } => Some(node), - _ => None, + let inventory = enumerate_local_logical_candidates( + vec![(0, QueryRoot::Operator(Rc::clone(&root)))], + &Default::default(), + ) + .unwrap(); + let mut result = enumerate_choices(&inventory, 4096) + .iter() + .filter_map(|choice| { + match compose_logical_candidate(&inventory, choice) + .ok()? + .remove(0) + .1 + { + QueryRoot::Operator(node) => Some(node), + QueryRoot::Scalar(_) => None, + } + }) + .filter(|node| { + OperatorNode::reachable(node).iter().all(|node| { + !matches!(&node.operator, asap_types::ir::Operator::ASAP(ASAPOp::SummaryAgg { + family: FieldDataType::Sketch(kind, _), + input, + .. + }) if matches!(kind.algorithm(), SketchAlgorithm::Cms | SketchAlgorithm::CmsWithHeap) + && !matches!(input.weight_domain, WeightDomain::NonNegative { .. })) + }) }) .collect::>(); result.push(executor_models::selected_dag(root, accuracy)); diff --git a/crates/integration-tests/tests/promql_numeric_regressions.rs b/crates/integration-tests/tests/promql_numeric_regressions.rs index 1d1aac8ab..e0c130376 100644 --- a/crates/integration-tests/tests/promql_numeric_regressions.rs +++ b/crates/integration-tests/tests/promql_numeric_regressions.rs @@ -2,26 +2,43 @@ //! The count/sum interpreter below verifies planner update semantics, not a deployed backend. use asap_integration_tests::fixtures::lower_promql; use asap_integration_tests::post_asap::post_asap_dag; -use asap_logical_optimizer::pass1::replacement::is_logical_rewrite; -use asap_logical_optimizer::{ASAPStrategies, Replacement, ReplacementStrategy, TargetSubDAG}; +use asap_logical_optimizer::pass1::logical_candidates::{ + compose_logical_candidate, enumerate_choices, enumerate_local_logical_candidates, + LocalLogicalCandidates, +}; use asap_types::ir::operator::Reduction; use asap_types::ir::scalar::ColumnRef; use asap_types::ir::schema::{ExactKind, FieldDataType, SummaryInputExpr, SummaryUpdate}; -use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}; +use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, QueryRoot}; use asap_types::types::AccuracyTarget; use std::rc::Rc; -fn plan(query: &str, accuracy: AccuracyTarget) -> Rc { +/// Stage 1's Pass 1 alternatives for `query`. +fn inventory(query: &str, accuracy: AccuracyTarget) -> LocalLogicalCandidates { let pre = lower_promql(query, accuracy).unwrap(); - ASAPStrategies::default() - .replacements(&TargetSubDAG::new(&pre)) - .into_iter() - .find_map(|r| match r.replacement { - // A bound decision: a summary DAG or a kept (exact) sub-DAG. - Replacement::SubDAG(n) if !is_logical_rewrite(&n) => Some(n), - _ => None, - }) - .unwrap_or_else(|| asap_logical_optimizer::pass1::replacement::retain_exact(&pre).unwrap()) + enumerate_local_logical_candidates(vec![(0, QueryRoot::Operator(pre))], &Default::default()) + .unwrap() +} +fn compose(inventory: &LocalLogicalCandidates, choice: &[usize]) -> Rc { + match compose_logical_candidate(inventory, choice) + .unwrap() + .remove(0) + .1 + { + QueryRoot::Operator(node) => node, + QueryRoot::Scalar(_) => unreachable!("operator root"), + } +} +/// The Stage 1 candidate where every target takes its first alternative other +/// than pass-through (an exact accumulator, else the first sketch). +fn plan(query: &str, accuracy: AccuracyTarget) -> Rc { + let inventory = inventory(query, accuracy); + let choice: Vec<_> = inventory + .targets + .iter() + .map(|target| usize::from(target.alternatives.len() > 1)) + .collect(); + compose(&inventory, &choice) } fn aggregate(node: &OperatorNode) -> (&FieldDataType, &SummaryUpdate, &Reduction) { match &node.operator { @@ -40,6 +57,17 @@ fn aggregate(node: &OperatorNode) -> (&FieldDataType, &SummaryUpdate, &Reduction other => panic!("not a maintained accumulator: {other:?}"), } } +/// The family and update of the summary `node` builds, if any. +fn summary(node: &Rc) -> Option<(FieldDataType, SummaryUpdate)> { + OperatorNode::reachable(node) + .iter() + .find_map(|node| match &node.operator { + Operator::ASAP(ASAPOp::SummaryAgg { family, input, .. }) => { + Some((family.clone(), input.clone())) + } + _ => None, + }) +} fn contribution(family: &FieldDataType, update: &SummaryUpdate, value: f64) -> f64 { if matches!(family, FieldDataType::ExactAggregate(ExactKind::Count, _)) { return 1.; @@ -102,7 +130,6 @@ fn sum_rate_and_increase_have_real_exact_accumulator_nodes() { let node = plan(query, AccuracyTarget::Exact); let (family, _, _) = aggregate(&node); assert!(matches!(family, FieldDataType::ExactAggregate(k, _) if *k == kind)); - assert!(node.guarantee.as_ref().unwrap().is_exact()); post_asap_dag(&node); } } @@ -119,7 +146,6 @@ fn checked_ratio_must_not_certify_cross_zero_interpolation() { && node.contains_asap(), "direct quantile ratio should remain an available candidate" ); - assert!(node.guarantee.is_none()); post_asap_dag(&node); // Keep the actual signed-sketch counterexample: division guards alone pass // even though the quantile interpolation does not preserve relative error. @@ -151,7 +177,7 @@ fn checked_ratio_must_not_certify_cross_zero_interpolation() { fn quantile_over_temporal_average_keeps_a_legal_candidate() { for query in [ "quantile(0.9, avg_over_time(a[5m]))", - "quantile(0.9, avg_over_time(a[5m]) + avg_over_time(b[5m]))", + "quantile(0.9, avg_over_time(a[5m]) + avg_over_time(a[5m] offset 5m))", ] { let node = plan(query, AccuracyTarget::Epsilon(0.01)); assert!( @@ -172,62 +198,40 @@ fn quantile_over_temporal_average_keeps_a_legal_candidate() { } } -struct OneKeyTopKEvidence; -impl asap_logical_optimizer::accuracy::AccuracyEvidenceProvider for OneKeyTopKEvidence { - fn propagation_stats( - &self, - op: &asap_types::ir::properties::CompositionOperator, - _family: &FieldDataType, - _query: Option<&asap_types::ir::schema::SketchStatistic>, - ) -> asap_logical_optimizer::accuracy::PropagationStats { - // Single-key fixture: no excluded keys; bounds cover every value below. - if matches!( - op, - asap_types::ir::properties::CompositionOperator::TopKSelection - ) { - asap_logical_optimizer::accuracy::PropagationStats { - topk_selected_lower_bound: Some(-1000.), - topk_excluded_upper_bound: Some(-1001.), - topk_interval_failure_probability: Some(0.001), - ..Default::default() - } - } else { - Default::default() - } - } -} - #[test] fn sketch_counts_use_unit_weights_and_signed_sums_keep_value_weights() { - use asap_logical_optimizer::accuracy::{DefaultAccuracyModel, EqualSplitAllocator}; use asap_types::ir::schema::{NonNegativeWeightProof, SketchAlgorithm, WeightDomain}; - let strategy = ASAPStrategies::new_with_planning_inputs_and_evidence( - &DefaultAccuracyModel, - &EqualSplitAllocator, - &OneKeyTopKEvidence, - ); for is_count in [true, false] { let query = if is_count { "topk(1, count_over_time(up[5m]))" } else { "topk(1, sum_over_time(up[5m]))" }; - let pre = lower_promql(query, AccuracyTarget::Epsilon(0.01)).unwrap(); - let candidates = strategy.replacements(&TargetSubDAG::new(&pre)); + let inventory = inventory(query, AccuracyTarget::Epsilon(0.01)); + let candidates: Vec<_> = enumerate_choices(&inventory, 4096) + .iter() + .map(|choice| compose(&inventory, choice)) + .collect(); let wanted = if is_count { SketchAlgorithm::CmsWithHeap } else { SketchAlgorithm::CountSketchWithHeap }; + // The heap that absorbs the inner aggregate reads its input rows: the + // candidate's only summary, with no aggregate left beneath it. let node = candidates .iter() - .find_map(|c| { - let Replacement::SubDAG(node) = &c.replacement else { - return None; - }; - let (family, _, _) = aggregate(node); - matches!(family, FieldDataType::Sketch(kind, _) if kind.algorithm() == &wanted) - .then_some(node) + .find(|node| { + let reachable = OperatorNode::reachable(node); + summary(node).is_some_and(|(family, _)| matches!(family, FieldDataType::Sketch(kind, _) if kind.algorithm() == &wanted)) + && reachable + .iter() + .filter(|node| matches!(node.operator, Operator::ASAP(ASAPOp::SummaryAgg { .. }))) + .count() + == 1 + && reachable + .iter() + .all(|node| !matches!(node.non_asap(), Some(NonASAPOp::Aggregate { .. }))) }) .expect("weighted sketch candidate"); let (family, update, _) = aggregate(node); @@ -244,11 +248,11 @@ fn sketch_counts_use_unit_weights_and_signed_sums_keep_value_weights() { update.weight, SummaryInputExpr::Column(ColumnRef::SampleValue) ); - for c in &candidates { - if let Replacement::SubDAG(n) = &c.replacement { - assert!( - !matches!(aggregate(n).0, FieldDataType::Sketch(kind, _) if kind.algorithm() == &SketchAlgorithm::CmsWithHeap) - ); + // Signed sums never prove the non-negative weights Count-Min needs. + for (family, update) in candidates.iter().filter_map(summary) { + if matches!(family, FieldDataType::Sketch(kind, _) if kind.algorithm() == &SketchAlgorithm::CmsWithHeap) + { + assert_eq!(update.weight_domain, WeightDomain::UnknownOrSigned); } } }