From 103b11ac9b988ed0f79c131947c93fccf0633fae Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Mon, 5 Oct 2026 02:11:46 +0000 Subject: [PATCH 1/2] refactor(logical-optimizer): move Stage 1's realization helpers out of the legacy search Stage 1 (pass1/logical_candidates.rs, pass2/*) and the accuracy estimators imported their realization catalogue and summary-input rules from the legacy search modules (pass1/replacement.rs, pass1/grouping.rs, pass2/reconciliation.rs). Move those items into a new pass1/realization.rs module, and hydra_guarantee into accuracy/estimators, so the stage pipeline no longer reaches the legacy search. No behaviour change. Co-Authored-By: Claude Opus 5.5 --- crates/frontend-sql/tests/frequency_l2.rs | 2 +- .../tests/planner_layering_example2.rs | 2 +- .../src/accuracy/estimators/cms.rs | 4 +- .../src/accuracy/estimators/count_sketch.rs | 2 +- .../src/accuracy/estimators/hll.rs | 2 +- .../src/accuracy/estimators/kll.rs | 2 +- .../src/accuracy/estimators/mod.rs | 94 +++- crates/logical-optimizer/src/accuracy/mod.rs | 2 +- crates/logical-optimizer/src/lib.rs | 12 +- .../src/pass1/exact_composition.rs | 5 +- .../logical-optimizer/src/pass1/grouping.rs | 149 +----- .../src/pass1/logical_candidates.rs | 25 +- crates/logical-optimizer/src/pass1/mod.rs | 1 + .../src/pass1/realization.rs | 432 ++++++++++++++++++ .../src/pass1/replacement.rs | 375 +-------------- .../src/pass2/reconciliation.rs | 36 +- .../src/pass2/summary_capability.rs | 8 +- .../src/pass2/window_composition.rs | 2 +- crates/plan-selection/src/cost/cost_model.rs | 6 +- .../plan-selection/src/cost/empirical_cost.rs | 5 +- crates/planner/tests/summary_sharing.rs | 2 +- 21 files changed, 589 insertions(+), 579 deletions(-) create mode 100644 crates/logical-optimizer/src/pass1/realization.rs diff --git a/crates/frontend-sql/tests/frequency_l2.rs b/crates/frontend-sql/tests/frequency_l2.rs index ff4c05638..748cd1d87 100644 --- a/crates/frontend-sql/tests/frequency_l2.rs +++ b/crates/frontend-sql/tests/frequency_l2.rs @@ -91,7 +91,7 @@ async fn frequency_l2_preserves_accuracy_target() { #[tokio::test] async fn stage1_offers_exact_and_summary_l2_alternatives() { use asap_logical_optimizer::pass1::logical_candidates::enumerate_local_logical_candidates; - use asap_logical_optimizer::pass1::replacement::Realization; + use asap_logical_optimizer::pass1::realization::Realization; use asap_types::ir::QueryRoot; let target = AccuracyTarget::EpsilonDelta { epsilon: 0.01, diff --git a/crates/integration-tests/tests/planner_layering_example2.rs b/crates/integration-tests/tests/planner_layering_example2.rs index fbfa89c94..b0bc8969e 100644 --- a/crates/integration-tests/tests/planner_layering_example2.rs +++ b/crates/integration-tests/tests/planner_layering_example2.rs @@ -4,7 +4,7 @@ mod executor_models; mod physical_common; use asap_executor::values::Value; use asap_frontend_sql::{lower_sql, SqlCatalog}; -use asap_logical_optimizer::pass1::replacement::Realization; +use asap_logical_optimizer::pass1::realization::Realization; use asap_plan_selection::plan_stages; use asap_types::ir::operator::AggIntent; use asap_types::ir::schema::{DataType, Field, Schema, SketchAlgorithm}; diff --git a/crates/logical-optimizer/src/accuracy/estimators/cms.rs b/crates/logical-optimizer/src/accuracy/estimators/cms.rs index 7a0c43477..6208b0fa5 100644 --- a/crates/logical-optimizer/src/accuracy/estimators/cms.rs +++ b/crates/logical-optimizer/src/accuracy/estimators/cms.rs @@ -41,7 +41,7 @@ mod tests { #[test] fn local_guarantee_inverts_frequency_sizing() { - use crate::pass1::replacement::default_size_params; + use crate::pass1::realization::default_size_params; use asap_types::ir::schema::{GroupingStrategy, SketchKind}; let c = asap_types::ir::operator::agg_intent::default_cardinality(); let params = default_size_params(SketchAlgorithm::Cms, &c, 0.01, 0.001); @@ -69,7 +69,7 @@ mod tests { /// for ε/2 and δ/2 (Pass 1's split) it meets it. #[test] fn hydra_guarantee_adds_the_shared_grid_term() { - use crate::pass1::replacement::default_size_params; + use crate::pass1::realization::default_size_params; use asap_types::ir::schema::{ default_hydra_params, GroupingStrategy, HydraKind, SketchKind, }; diff --git a/crates/logical-optimizer/src/accuracy/estimators/count_sketch.rs b/crates/logical-optimizer/src/accuracy/estimators/count_sketch.rs index 9e9fc66c0..919f920a2 100644 --- a/crates/logical-optimizer/src/accuracy/estimators/count_sketch.rs +++ b/crates/logical-optimizer/src/accuracy/estimators/count_sketch.rs @@ -55,7 +55,7 @@ mod tests { #[test] fn count_sketch_uses_an_l2_guarantee() { - use crate::pass1::replacement::default_size_params; + use crate::pass1::realization::default_size_params; use asap_types::ir::operator::agg_intent::default_cardinality; use asap_types::ir::schema::{GroupingStrategy, SketchKind}; let intent = default_cardinality(); diff --git a/crates/logical-optimizer/src/accuracy/estimators/hll.rs b/crates/logical-optimizer/src/accuracy/estimators/hll.rs index c7c26b435..836f2c574 100644 --- a/crates/logical-optimizer/src/accuracy/estimators/hll.rs +++ b/crates/logical-optimizer/src/accuracy/estimators/hll.rs @@ -230,7 +230,7 @@ mod tests { } #[test] fn generic_rse_sizing_does_not_certify_confidence() { - use crate::pass1::replacement::default_size_params; + use crate::pass1::realization::default_size_params; use asap_types::ir::operator::agg_intent::default_cardinality; use asap_types::ir::schema::{GroupingStrategy, SketchKind}; let c = default_cardinality(); diff --git a/crates/logical-optimizer/src/accuracy/estimators/kll.rs b/crates/logical-optimizer/src/accuracy/estimators/kll.rs index efa3feaa0..2993af5ed 100644 --- a/crates/logical-optimizer/src/accuracy/estimators/kll.rs +++ b/crates/logical-optimizer/src/accuracy/estimators/kll.rs @@ -42,7 +42,7 @@ mod tests { #[test] fn local_guarantee_inverts_rank_sizing() { - use crate::pass1::replacement::default_size_params; + use crate::pass1::realization::default_size_params; use asap_types::ir::operator::agg_intent::default_quantile; use asap_types::ir::schema::{GroupingStrategy, SketchKind}; let q = default_quantile(0.99); diff --git a/crates/logical-optimizer/src/accuracy/estimators/mod.rs b/crates/logical-optimizer/src/accuracy/estimators/mod.rs index c3406b830..3ffeba42d 100644 --- a/crates/logical-optimizer/src/accuracy/estimators/mod.rs +++ b/crates/logical-optimizer/src/accuracy/estimators/mod.rs @@ -91,7 +91,7 @@ pub(super) fn local_guarantee( hydra_shared_grid_failure_probability: Some((-f64::from(*shared_rows)).exp()), ..Default::default() }; - Some(crate::pass1::grouping::hydra_guarantee(&inner, &stats)) + Some(hydra_guarantee(&inner, &stats)) } // No accuracy model for the other shared groupings. FieldDataType::Sketch(..) => None, @@ -186,7 +186,7 @@ impl<'a> EstimatorAccuracy<'a> { target: Option<&AccuracyTarget>, ) -> Self { let (epsilon, delta) = target - .map(crate::pass1::replacement::accuracy_budget) + .map(crate::pass1::realization::accuracy_budget) .unwrap_or((0.0, 0.0)); Self { base, @@ -257,3 +257,93 @@ impl AccuracyModel for EstimatorAccuracy<'_> { self.base.answers(statistic, guarantee) } } + +/// Compose the inner per-subpopulation guarantee with Hydra's outer shared +/// grid. The paper's collision term depends on deployment/data statistics; +/// keeping those leaves symbolic makes the formula explicit while ensuring +/// target satisfaction fails closed until a caller supplies them. +pub(crate) fn hydra_guarantee( + inner: &ResultGuarantee, + stats: &PropagationStats, +) -> ResultGuarantee { + let mut provenance = inner.provenance.clone(); + provenance.extend(stats.evidence_provenance.clone()); + provenance.push(GuaranteeSource::ChildGuarantee { + input_index: 0, + guarantee: Box::new(inner.clone()), + }); + if stats.hydra_shared_grid_collision_bound.is_none() { + provenance.push(GuaranteeSource::UnavailableStatistic { + statistic: "hydra_shared_grid_collision_bound".into(), + }); + } + if stats.hydra_shared_grid_failure_probability.is_none() { + provenance.push(GuaranteeSource::UnavailableStatistic { + statistic: "hydra_shared_grid_failure_probability".into(), + }); + } + provenance.push(GuaranteeSource::CompositionStep { + operator: CompositionOperator::ApproximateAggregate, + rule: "hydra_shared_grid_union_bound".into(), + }); + ResultGuarantee { + metric: inner.metric, + bound: BoundExpr::Sum { + terms: vec![ + inner.bound.clone(), + stats.hydra_shared_grid_collision_bound.map_or_else( + || BoundExpr::Unknown { + statistic: "hydra_shared_grid_collision_bound".into(), + }, + |value| BoundExpr::Constant { value }, + ), + ], + }, + failure_probability: ProbabilityExpr::UnionBound { + terms: vec![ + inner.failure_probability.clone(), + stats.hydra_shared_grid_failure_probability.map_or_else( + || ProbabilityExpr::Unknown { + statistic: "hydra_shared_grid_failure_probability".into(), + }, + |value| ProbabilityExpr::Constant { value }, + ), + ], + }, + provenance, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn hydra_composes_inner_and_shared_grid_error_symbolically() { + let inner = ResultGuarantee { + metric: ErrorMetric::Frequency, + bound: BoundExpr::Constant { value: 0.01 }, + failure_probability: ProbabilityExpr::Constant { value: 0.02 }, + provenance: vec![], + }; + let composed = hydra_guarantee(&inner, &PropagationStats::default()); + + assert_eq!(composed.metric, ErrorMetric::Frequency); + assert!(matches!( + composed.bound, + BoundExpr::Sum { ref terms } + if matches!(terms.as_slice(), [ + BoundExpr::Constant { value }, + BoundExpr::Unknown { statistic }, + ] if *value == 0.01 && statistic == "hydra_shared_grid_collision_bound") + )); + assert!(matches!( + composed.failure_probability, + ProbabilityExpr::UnionBound { ref terms } + if matches!(terms.as_slice(), [ + ProbabilityExpr::Constant { value }, + ProbabilityExpr::Unknown { statistic }, + ] if *value == 0.02 && statistic == "hydra_shared_grid_failure_probability") + )); + } +} diff --git a/crates/logical-optimizer/src/accuracy/mod.rs b/crates/logical-optimizer/src/accuracy/mod.rs index 762ab96b3..b609bd034 100644 --- a/crates/logical-optimizer/src/accuracy/mod.rs +++ b/crates/logical-optimizer/src/accuracy/mod.rs @@ -45,7 +45,7 @@ pub trait AccuracyModel { /// The guarantee of reading `query` out of a summary of family `family` /// built over an **exact** input — derived from the family's committed /// parameters by inverting the same sizing formulas - /// [`crate::pass1::replacement::default_size_params`] uses. `None` when this + /// [`crate::pass1::realization::default_size_params`] uses. `None` when this /// model has no error model for the family (the default has none for /// `Sample`/`Wavelet`/`StatModel`). fn local_guarantee( diff --git a/crates/logical-optimizer/src/lib.rs b/crates/logical-optimizer/src/lib.rs index d5d6f458a..6ca519c8d 100644 --- a/crates/logical-optimizer/src/lib.rs +++ b/crates/logical-optimizer/src/lib.rs @@ -66,14 +66,14 @@ pub use pass1::exact_composition::{ pub use pass1::explanation::{ explain_replacements, explain_replacements_with, ExplanationKind, ReplacementExplanation, }; -pub use pass1::grouping::{has_subpopulations, HydraGroupingStrategy}; +pub use pass1::grouping::HydraGroupingStrategy; +pub use pass1::realization::{has_subpopulations, summary_candidates, Realization}; pub use pass1::replacement::{ default_strategies, is_logical_rewrite, search_workload, search_workload_with, - search_workload_with_targets, summary_candidates, ASAPStrategies, CandidateLogicalASAPDAGs, - GlobalSelection, Matcher, Proposals, Realization, RealizationError, RejectedCandidate, - Replacement, ReplacementProvenance, ReplacementStrategy, ReplacementSubDAG, - SharedSubDAGStrategy, TargetSubDAG, TargetSubDAGCandidates, TargetSubDAGSelection, - MAX_SEARCH_ITERATIONS, + search_workload_with_targets, ASAPStrategies, CandidateLogicalASAPDAGs, GlobalSelection, + Matcher, Proposals, RealizationError, RejectedCandidate, Replacement, ReplacementProvenance, + ReplacementStrategy, ReplacementSubDAG, SharedSubDAGStrategy, TargetSubDAG, + TargetSubDAGCandidates, TargetSubDAGSelection, MAX_SEARCH_ITERATIONS, }; pub use pass1::rewrite::{AvgToSumOverCountStrategy, SemanticEquivalentRewriteStrategy}; pub use pass2::reconciliation::AccuracyReconciliationStrategy; diff --git a/crates/logical-optimizer/src/pass1/exact_composition.rs b/crates/logical-optimizer/src/pass1/exact_composition.rs index 8b5fab902..db5a0654e 100644 --- a/crates/logical-optimizer/src/pass1/exact_composition.rs +++ b/crates/logical-optimizer/src/pass1/exact_composition.rs @@ -77,9 +77,10 @@ use asap_types::physical::execution_data_state::lift_plain; use asap_types::physical::ExactOperationSchemaError; use asap_types::types::AccuracyTarget; +use crate::pass1::realization::Realization; use crate::pass1::replacement::{ - bindable_intent, describe_intent, realizations_for_intent, Realization, RealizationError, - Replacement, ReplacementProvenance, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, + bindable_intent, describe_intent, realizations_for_intent, RealizationError, Replacement, + ReplacementProvenance, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, }; use crate::{AccuracyModel, DefaultAccuracyModel, PropagationStats}; diff --git a/crates/logical-optimizer/src/pass1/grouping.rs b/crates/logical-optimizer/src/pass1/grouping.rs index 3b5b1ca29..5ec864cd3 100644 --- a/crates/logical-optimizer/src/pass1/grouping.rs +++ b/crates/logical-optimizer/src/pass1/grouping.rs @@ -72,45 +72,26 @@ use std::rc::Rc; use asap_types::ir::operator::agg_intent::AggIntent; -use asap_types::ir::operator::operator_properties::Reduction; -use asap_types::ir::properties::{ - AccuracyError, BoundExpr, CompositionOperator, GuaranteeSource, ProbabilityExpr, - ResultGuarantee, -}; +use asap_types::ir::properties::{AccuracyError, CompositionOperator}; use asap_types::ir::schema::{ default_hydra_params, hydra_kind_for, FieldDataType, GroupingStrategy, HydraKind, SketchAlgorithm, SketchParams, }; use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}; +use crate::accuracy::estimators::hydra_guarantee; use crate::accuracy::{ AccuracyBudgetAllocator, AccuracyEvidenceProvider, AccuracyModel, PropagationStats, }; +use crate::pass1::realization::{ + accuracy_target, has_subpopulations, summary_candidates, Realization, +}; use crate::pass1::replacement::{ - accuracy_target, bindable_intent, construct_summary_with, describe_intent, - realizations_for_intent, summary_candidates, CandidatePlanningInputs, Proposals, Realization, - RejectedCandidate, Replacement, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, + bindable_intent, construct_summary_with, describe_intent, realizations_for_intent, + CandidatePlanningInputs, Proposals, RejectedCandidate, Replacement, ReplacementStrategy, + ReplacementSubDAG, TargetSubDAG, }; -/// Whether `reduction` has a genuine subpopulation concept for -/// `GroupingStrategy::SharedMultiSubpopulation` to multiplex across — the -/// non-empty-`by` legality condition issue #256 requires. -/// -/// - [`Reduction::PerEntity`]: no grouping concept at all (never merges -/// across entities) — `false`. -/// - [`Reduction::Reduce`] with an empty, non-`without` `by`: a genuine full -/// reduction, one output row, no subpopulations — `false`. -/// - [`Reduction::Reduce`] with a non-empty `by`, or any `without(...)` -/// exclusion grouping (which groups by whatever labels remain, even -/// `without([])` — "group by every label"): a real subpopulation concept -/// — `true`. -pub fn has_subpopulations(reduction: &Reduction) -> bool { - match reduction.group_keys() { - None => false, - Some(keys) => keys.is_without() || !keys.is_empty(), - } -} - /// Wraps the `GroupingStrategy` axis (issue #256) as a /// [`ReplacementStrategy`]: for a target `ASAPStrategies` /// already has an opinion on, offers an additional @@ -424,127 +405,15 @@ fn with_grouping( } } -/// Compose the inner per-subpopulation guarantee with Hydra's outer shared -/// grid. The paper's collision term depends on deployment/data statistics; -/// keeping those leaves symbolic makes the formula explicit while ensuring -/// target satisfaction fails closed until a caller supplies them. -pub(crate) fn hydra_guarantee( - inner: &ResultGuarantee, - stats: &PropagationStats, -) -> ResultGuarantee { - let mut provenance = inner.provenance.clone(); - provenance.extend(stats.evidence_provenance.clone()); - provenance.push(GuaranteeSource::ChildGuarantee { - input_index: 0, - guarantee: Box::new(inner.clone()), - }); - if stats.hydra_shared_grid_collision_bound.is_none() { - provenance.push(GuaranteeSource::UnavailableStatistic { - statistic: "hydra_shared_grid_collision_bound".into(), - }); - } - if stats.hydra_shared_grid_failure_probability.is_none() { - provenance.push(GuaranteeSource::UnavailableStatistic { - statistic: "hydra_shared_grid_failure_probability".into(), - }); - } - provenance.push(GuaranteeSource::CompositionStep { - operator: CompositionOperator::ApproximateAggregate, - rule: "hydra_shared_grid_union_bound".into(), - }); - ResultGuarantee { - metric: inner.metric, - bound: BoundExpr::Sum { - terms: vec![ - inner.bound.clone(), - stats.hydra_shared_grid_collision_bound.map_or_else( - || BoundExpr::Unknown { - statistic: "hydra_shared_grid_collision_bound".into(), - }, - |value| BoundExpr::Constant { value }, - ), - ], - }, - failure_probability: ProbabilityExpr::UnionBound { - terms: vec![ - inner.failure_probability.clone(), - stats.hydra_shared_grid_failure_probability.map_or_else( - || ProbabilityExpr::Unknown { - statistic: "hydra_shared_grid_failure_probability".into(), - }, - |value| ProbabilityExpr::Constant { value }, - ), - ], - }, - provenance, - } -} - #[cfg(test)] mod tests { use super::*; use crate::accuracy::{DefaultAccuracyModel, EqualSplitAllocator}; use crate::test_support::{agg, agg_per_entity, metric_scan}; use asap_types::ir::operator::agg_intent::{default_cardinality, default_quantile}; - use asap_types::ir::properties::ErrorMetric; + use asap_types::ir::operator::operator_properties::Reduction; use asap_types::types::AccuracyTarget; - // ── has_subpopulations ──────────────────────────────────────────────── - - #[test] - fn per_entity_has_no_subpopulation_concept() { - assert!(!has_subpopulations(&Reduction::PerEntity)); - } - - #[test] - fn empty_by_reduction_has_no_subpopulation_concept() { - assert!(!has_subpopulations(&Reduction::by(vec![]))); - } - - #[test] - fn non_empty_by_reduction_has_a_subpopulation_concept() { - assert!(has_subpopulations(&Reduction::by(vec![2]))); - } - - #[test] - fn without_grouping_has_a_subpopulation_concept_even_when_empty() { - use asap_types::ir::operator::operator_properties::GroupKeys; - // `without([])` groups by every remaining label — a real - // subpopulation concept, unlike `by([])`'s genuine full reduction. - assert!(has_subpopulations(&Reduction::Reduce(GroupKeys::without( - vec![] - )))); - } - - #[test] - fn hydra_composes_inner_and_shared_grid_error_symbolically() { - let inner = ResultGuarantee { - metric: ErrorMetric::Frequency, - bound: BoundExpr::Constant { value: 0.01 }, - failure_probability: ProbabilityExpr::Constant { value: 0.02 }, - provenance: vec![], - }; - let composed = hydra_guarantee(&inner, &PropagationStats::default()); - - assert_eq!(composed.metric, ErrorMetric::Frequency); - assert!(matches!( - composed.bound, - BoundExpr::Sum { ref terms } - if matches!(terms.as_slice(), [ - BoundExpr::Constant { value }, - BoundExpr::Unknown { statistic }, - ] if *value == 0.01 && statistic == "hydra_shared_grid_collision_bound") - )); - assert!(matches!( - composed.failure_probability, - ProbabilityExpr::UnionBound { ref terms } - if matches!(terms.as_slice(), [ - ProbabilityExpr::Constant { value }, - ProbabilityExpr::Unknown { statistic }, - ] if *value == 0.02 && statistic == "hydra_shared_grid_failure_probability") - )); - } - // ── HydraGroupingStrategy ───────────────────────────────────────────── #[test] diff --git a/crates/logical-optimizer/src/pass1/logical_candidates.rs b/crates/logical-optimizer/src/pass1/logical_candidates.rs index e499fd6ca..d773fca81 100644 --- a/crates/logical-optimizer/src/pass1/logical_candidates.rs +++ b/crates/logical-optimizer/src/pass1/logical_candidates.rs @@ -20,8 +20,10 @@ use asap_types::types::AccuracyTarget; use asap_types::workload::MetricType; use thiserror::Error; -use crate::pass1::replacement::{ - accuracy_budget, accuracy_target, default_size_params, summary_candidates, Realization, +use crate::pass1::realization::{ + accuracy_budget, accuracy_target, column_ref, default_size_params, has_subpopulations, + realize_keyed_additive_summary_input, summarised_input, summary_candidates, + PhysicalSummaryInputRuleResult, Realization, }; use crate::pass2::window_composition::{tumbling_state, WindowForm}; @@ -244,7 +246,7 @@ pub fn add_hydra_alternatives( return Ok(()); }; if *accuracy == AccuracyTarget::Exact - || !crate::pass1::grouping::has_subpopulations(reduction) + || !has_subpopulations(reduction) || reduction.group_keys().is_some_and(|keys| keys.is_without()) || hydra_update(reduction, &child.schema).is_none() { @@ -275,9 +277,9 @@ pub fn add_hydra_alternatives( /// a [`count_item`]. fn hydra_update(reduction: &Reduction, child: &Schema) -> Option { Some(SummaryUpdate { - item: Some(SummaryInputExpr::Column( - crate::pass1::replacement::column_ref(count_item(reduction, child)?), - )), + item: Some(SummaryInputExpr::Column(column_ref(count_item( + reduction, child, + )?))), weight: SummaryInputExpr::Constant(1.0), weight_domain: WeightDomain::NonNegative { proof: NonNegativeWeightProof::UnitCount, @@ -314,9 +316,6 @@ fn count_item<'a>( /// apply. Rows that carry the full series identity rank it as a column, as /// [`summary_update`] does. fn whole_expression_input(target: &OperatorNode) -> Option<(Rc, SummaryUpdate)> { - use crate::pass1::replacement::{ - realize_keyed_additive_summary_input, PhysicalSummaryInputRuleResult, - }; let Some(NonASAPOp::Aggregate { child, reduction, @@ -778,7 +777,7 @@ fn summary_update( return sql_row_count_update(algorithm.is_some(), reduction, child); } } - let weight = crate::pass1::replacement::summarised_input(intent, child) + let weight = summarised_input(intent, child) .map_err(|_| LogicalCandidateError::Unsupported("input column outside child schema"))?; Ok(match (intent, algorithm) { (AggIntent::TopK { .. }, Some(_)) => { @@ -799,7 +798,7 @@ fn summary_update( .into_iter() .flat_map(|keys| keys.iter()) .filter_map(|&index| child.fields.get(index)) - .map(crate::pass1::replacement::column_ref) + .map(column_ref) .collect(); // Rows that carry the full series identity rank it as a column, // the item form the runtime builds keyed summaries from. @@ -846,9 +845,7 @@ fn sql_row_count_update( let column = count_item(reduction, child).ok_or(LogicalCandidateError::Unsupported( "a COUNT(*) sketch needs a non-null item column", ))?; - Some(SummaryInputExpr::Column( - crate::pass1::replacement::column_ref(column), - )) + Some(SummaryInputExpr::Column(column_ref(column))) } else { None }; diff --git a/crates/logical-optimizer/src/pass1/mod.rs b/crates/logical-optimizer/src/pass1/mod.rs index bf32ceee4..95bd58f66 100644 --- a/crates/logical-optimizer/src/pass1/mod.rs +++ b/crates/logical-optimizer/src/pass1/mod.rs @@ -8,6 +8,7 @@ pub(crate) mod function_rules; pub mod grouping; pub mod logical_candidates; pub mod maintained_population; +pub mod realization; pub mod replacement; pub mod rewrite; pub mod rollup; diff --git a/crates/logical-optimizer/src/pass1/realization.rs b/crates/logical-optimizer/src/pass1/realization.rs new file mode 100644 index 000000000..cafc290b0 --- /dev/null +++ b/crates/logical-optimizer/src/pass1/realization.rs @@ -0,0 +1,432 @@ +//! How one aggregate intent may be realized, and the rules Stage 1 uses to +//! size a summary and bind its input: the summary families per intent, the +//! accuracy budget a target resolves to, and the update a summary reads. + +use std::rc::Rc; + +use asap_types::ir::operator::{AggIntent, Reduction}; +use asap_types::ir::scalar::ColumnRef; +use asap_types::ir::schema::{ + EntityIdentity, ExactKind, ExactParams, Field, FieldDataType, NonNegativeWeightProof, + SamplingKind, SamplingParams, Schema, SketchAlgorithm, SketchKind, SketchParams, StatModelKind, + StatModelParams, SummaryInputExpr, SummaryUpdate, WaveletKind, WaveletParams, WeightDomain, +}; +use asap_types::ir::{NonASAPOp, OperatorNode}; +use asap_types::types::AccuracyTarget; + +/// How an [`AggIntent`] may be realised at post-ASAP binding time (issue +/// #98): by an approximate summary (sketch, sample, wavelet, statistical +/// model, …), by an exact mergeable accumulator, or by an ordinary exact +/// operator (pass-through). This is a post-ASAP concern — the pre-ASAP IR +/// carries only the intent + accuracy target, never the realization — and +/// it's a per-node decision, made once per `AggIntent`, not a plan-wide one. +/// +/// Stage 1 lists every realization of a target +/// ([`local_realizations_for_intent`](crate::pass1::logical_candidates::local_realizations_for_intent)). +#[derive(Debug, Clone, PartialEq)] +pub enum Realization { + /// An exact **mergeable** accumulator (partial state ≡ the value + /// itself: `Sum` / `Count` / `Min` / `Max` / `Rate` / `Increase`). The + /// built state *is* the answer already — no `SummaryEstimate` evaluation + /// step. + ExactAggregate { + kind: ExactKind, + params: ExactParams, + }, + /// An approximate sketch sized to the intent's [`AccuracyTarget`]. + /// Needs a `SummaryEstimate` evaluation to recover a value. Already + /// classified into its [`SketchKind`] category (`SketchKind::new` + /// having been called) — construction always goes through that + /// classifier, never this variant directly. + Sketch(SketchKind), + /// A sampling-based summary (a retained row subset). Needs a + /// `SummaryEstimate` evaluation. Not chosen by any core `AggIntent` + /// dispatch today — see the module docs. + Sample { + kind: SamplingKind, + params: SamplingParams, + }, + /// A wavelet-transform summary. Needs a `SummaryEstimate` evaluation. Not + /// chosen by any core `AggIntent` dispatch today — see the module docs. + Wavelet { + kind: WaveletKind, + params: WaveletParams, + }, + /// A fitted statistical/parametric-model summary. Needs a + /// `SummaryEstimate` evaluation. Not chosen by any core `AggIntent` + /// dispatch today — see the module docs. + StatModel { + kind: StatModelKind, + params: StatModelParams, + }, + /// No summary form — the node stays a logical pre-ASAP operator and is + /// executed exactly (per-series transforms, non-mergeable reducers, exact + /// quantile/top-k/cardinality, classic-bucket `HistogramQuantile`, …). + PassThrough, +} + +/// Confidence δ assumed when the target carries only an ε +/// (`AccuracyTarget::Epsilon`): the (ε, δ)-parameterised sketches (CMS) need +/// one. `ln(1/0.01) → depth 5`, matching the conventional CMS sizing. +pub const DEFAULT_DELTA: f64 = 0.01; + +/// The sketch kinds that can serve an intent, most-preferred first. +/// This is the `AggIntent → SketchAlgorithm` map of issue #98; +/// Stage 1 sizes every entry to the target's accuracy. Listed here so the +/// candidate set has one home. +pub fn summary_candidates(intent: &AggIntent) -> &'static [SketchAlgorithm] { + match intent { + AggIntent::Quantile { .. } => &[SketchAlgorithm::Kll, SketchAlgorithm::DDSketch], + // A distinct-tuple count hashes the whole tuple as one item + // (`SummaryInputExpr::Tuple`), which the distinct-count sketches take + // unchanged. UnivMon is dropped there: it estimates frequency moments + // over a single value stream, and `realize_value_frequency_summary_input` + // would feed it one column of the tuple. + AggIntent::Cardinality { cols, .. } if cols.len() > 1 => &[ + SketchAlgorithm::Hll, + SketchAlgorithm::Theta, + SketchAlgorithm::Kmv, + ], + AggIntent::Cardinality { .. } => &[ + SketchAlgorithm::Hll, + SketchAlgorithm::Theta, + SketchAlgorithm::Kmv, + SketchAlgorithm::UnivMon, + ], + AggIntent::FrequencyL2 { .. } | AggIntent::FrequencyEntropy { .. } => { + &[SketchAlgorithm::UnivMon] + } + // Count-Sketch-with-heap is CMS-with-heap's balanced/zero-mean-error + // alternative for the same heavy-hitter shape. + AggIntent::TopK { .. } => &[ + SketchAlgorithm::CmsWithHeap, + SketchAlgorithm::CountSketchWithHeap, + ], + AggIntent::Count { .. } => &[ + SketchAlgorithm::Cms, + SketchAlgorithm::CountSketch, + SketchAlgorithm::UnivMon, + ], + _ => &[], + } +} + +/// The [`AccuracyTarget`] threaded onto an approximate-capable intent +/// (`Quantile`/`Cardinality`/`Count`/`TopK`), or `None` for every other +/// intent (no sketch candidate applies). +pub fn accuracy_target(intent: &AggIntent) -> Option<&AccuracyTarget> { + match intent { + AggIntent::Quantile { accuracy, .. } + | AggIntent::Cardinality { accuracy, .. } + | AggIntent::FrequencyL2 { accuracy, .. } + | AggIntent::FrequencyEntropy { accuracy, .. } + | AggIntent::Count { accuracy } + | AggIntent::TopK { accuracy, .. } => Some(accuracy), + _ => None, + } +} + +/// Resolve an [`AccuracyTarget`] into the `(eps, delta)` budget sketch +/// sizing needs — one place this resolution happens, so nothing can drift +/// apart on it. +/// +/// `Exact` has no sketch realization; it degrades to the tightest parameters +/// for a caller that resolves it directly anyway. +pub fn accuracy_budget(accuracy: &AccuracyTarget) -> (f64, f64) { + match accuracy { + AccuracyTarget::Exact => (f64::MIN_POSITIVE, DEFAULT_DELTA), + AccuracyTarget::Epsilon(e) => (*e, DEFAULT_DELTA), + AccuracyTarget::EpsilonDelta { epsilon, delta } => (*epsilon, *delta), + } +} + +/// `asap-plan`'s built-in `SketchParams` sizing, keyed off the resolved +/// `(eps, delta)` accuracy budget. +/// +/// Each formula inverts the sketch family's standard error bound to the +/// smallest parameter satisfying the target, clamped to the family's sane +/// range. A non-positive ε saturates to the clamp maximum (tightest +/// allowed). +pub fn default_size_params( + kind: SketchAlgorithm, + intent: &AggIntent, + eps: f64, + delta: f64, +) -> SketchParams { + crate::accuracy::estimators::size_params(kind, intent, eps, delta) +} + +/// The physical input consumed by one summary realization. Most summaries +/// consume the logical aggregate's immediate child and summarize its declared +/// input value. Composite realizations can instead consume a larger +/// logical sub-DAG and bind a different key or value. +pub(crate) struct PhysicalSummaryInput { + pub(crate) child: Rc, + pub(crate) input: SummaryUpdate, +} + +pub(crate) enum PhysicalSummaryInputRuleResult { + NotApplicable, + Realized(PhysicalSummaryInput), + Unsupported(&'static str), +} + +/// Realize the composite heavy-hitter realization for +/// `TopK(Count GROUP BY key)`. The heap sketch consumes the raw keyed stream; +/// it does not consume an independently materialized Count result. +pub(crate) fn realize_keyed_additive_summary_input( + intent: &AggIntent, + family: &FieldDataType, + output_reduction: &Reduction, + child: &Rc, +) -> PhysicalSummaryInputRuleResult { + if !matches!(intent, AggIntent::TopK { .. }) { + return PhysicalSummaryInputRuleResult::NotApplicable; + } + let FieldDataType::Sketch(kind, _) = family else { + return PhysicalSummaryInputRuleResult::NotApplicable; + }; + let heap_algorithm = kind.algorithm(); + if !matches!( + heap_algorithm, + SketchAlgorithm::CmsWithHeap | SketchAlgorithm::CountSketchWithHeap + ) { + return PhysicalSummaryInputRuleResult::NotApplicable; + } + let Some(NonASAPOp::Aggregate { + reduction, + measures, + having: None, + child: raw_child, + .. + }) = child.non_asap() + else { + return PhysicalSummaryInputRuleResult::NotApplicable; + }; + let counter_input = matches!(measures.as_slice(), [AggIntent::Sum { .. }]) + && matches!(raw_child.non_asap(), Some(NonASAPOp::Aggregate { measures, .. }) + if matches!(measures.as_slice(), [AggIntent::Rate | AggIntent::Increase])); + let weight = match measures.as_slice() { + [AggIntent::Count { .. }] => SummaryInputExpr::Constant(1.0), + [AggIntent::Sum { .. }] if counter_input => { + SummaryInputExpr::Column(ColumnRef::SampleValue) + } + [AggIntent::Sum { col }] => SummaryInputExpr::Column(match col { + None => ColumnRef::SampleValue, + Some(index) => match schema_column_ref(raw_child, *index) { + Some(column) => column, + None => { + return PhysicalSummaryInputRuleResult::Unsupported( + "sum-ranked Top-K value column is outside the raw input schema", + ) + } + }, + }), + _ => return PhysicalSummaryInputRuleResult::NotApplicable, + }; + let weight_domain = match measures.as_slice() { + [AggIntent::Count { .. }] => WeightDomain::NonNegative { + proof: NonNegativeWeightProof::UnitCount, + }, + [AggIntent::Sum { .. }] if counter_input => WeightDomain::NonNegative { + proof: NonNegativeWeightProof::ResetAwareCounterDerivative, + }, + _ => WeightDomain::UnknownOrSigned, + }; + if matches!(heap_algorithm, SketchAlgorithm::CmsWithHeap) + && !matches!(weight_domain, WeightDomain::NonNegative { .. }) + { + return PhysicalSummaryInputRuleResult::Unsupported( + "value-weighted CMS requires non-negative update evidence; use CountSketch for arbitrary values", + ); + } + let subpopulation_columns = match output_reduction { + Reduction::PerEntity => vec![], + Reduction::Reduce(keys) => keys + .iter() + .filter_map(|index| schema_column_ref(child, *index)) + .collect(), + }; + let item = match reduction { + Reduction::PerEntity => SummaryInputExpr::EntityIdentity(EntityIdentity::PromqlLabelSet { + excluding: subpopulation_columns, + }), + Reduction::Reduce(keys) if !keys.is_without() && !keys.is_empty() => { + let Some(columns) = keys + .iter() + .map(|index| schema_column_ref(raw_child, *index)) + .collect::>>() + else { + return PhysicalSummaryInputRuleResult::Unsupported( + "ranked item column is outside the raw input schema", + ); + }; + let item_columns: Vec<_> = columns + .into_iter() + .filter(|column| !subpopulation_columns.contains(column)) + .collect(); + match item_columns.as_slice() { + [] => { + return PhysicalSummaryInputRuleResult::Unsupported( + "subpopulation columns consume the complete ranked item identity", + ) + } + [column] => SummaryInputExpr::Column(column.clone()), + _ => SummaryInputExpr::Tuple( + item_columns + .into_iter() + .map(SummaryInputExpr::Column) + .collect(), + ), + } + } + Reduction::Reduce(_) => { + return PhysicalSummaryInputRuleResult::Unsupported( + "an empty or without grouping does not identify ranked items", + ) + } + }; + PhysicalSummaryInputRuleResult::Realized(PhysicalSummaryInput { + child: Rc::clone(raw_child), + input: SummaryUpdate { + item: Some(item), + weight, + weight_domain, + }, + }) +} + +pub(crate) fn schema_column_ref(child: &OperatorNode, index: usize) -> Option { + let column = child.schema.fields.get(index)?; + Some(match &column.table { + Some(table) => ColumnRef::Qualified { + table: table.clone(), + name: column.name.clone(), + }, + None => ColumnRef::Named(column.name.clone()), + }) +} + +/// The column fed into a *single-column* summary: the intent's leading +/// positional input resolved to a name against the child schema, or the PromQL +/// sample value when it reads none. Callers are responsible for only reaching +/// here with a one-column intent — [`summarised_input`] is the general form. +pub(crate) fn summarised_column(intent: &AggIntent, child_schema: &Schema) -> ColumnRef { + match intent + .input_cols() + .first() + .and_then(|id| child_schema.fields.get(*id)) + { + Some(c) => column_ref(c), + None => ColumnRef::SampleValue, + } +} + +pub(crate) fn column_ref(column: &Field) -> ColumnRef { + match &column.table { + Some(t) => ColumnRef::Qualified { + table: t.clone(), + name: column.name.clone(), + }, + None => ColumnRef::Named(column.name.clone()), + } +} + +/// What the summary consumes per input row. An intent that reads one column (or +/// none) feeds that column; `COUNT(DISTINCT a, b)` feeds the whole tuple as one +/// item, so the distinct-count sketch hashes `(a, b)` rather than `a` — the +/// difference between tuple cardinality and single-column cardinality. +/// +/// A tuple leg outside the child schema is an error rather than +/// [`summarised_column`]'s sample-value fallback: a leg has no sample-value +/// reading, and silently dropping one would under-count. +pub(crate) fn summarised_input( + intent: &AggIntent, + child_schema: &Schema, +) -> Result { + let cols = intent.input_cols(); + if cols.len() < 2 { + return Ok(SummaryInputExpr::Column(summarised_column( + intent, + child_schema, + ))); + } + let legs = cols + .iter() + .map(|id| child_schema.fields.get(*id).map(column_ref)) + .collect::>>() + .ok_or("a tuple column is outside the input schema")?; + Ok(SummaryInputExpr::Tuple( + legs.into_iter().map(SummaryInputExpr::Column).collect(), + )) +} + +/// Whether `reduction` has a genuine subpopulation concept for +/// `GroupingStrategy::SharedMultiSubpopulation` to multiplex across — the +/// non-empty-`by` legality condition issue #256 requires. +/// +/// - [`Reduction::PerEntity`]: no grouping concept at all (never merges +/// across entities) — `false`. +/// - [`Reduction::Reduce`] with an empty, non-`without` `by`: a genuine full +/// reduction, one output row, no subpopulations — `false`. +/// - [`Reduction::Reduce`] with a non-empty `by`, or any `without(...)` +/// exclusion grouping (which groups by whatever labels remain, even +/// `without([])` — "group by every label"): a real subpopulation concept +/// — `true`. +pub fn has_subpopulations(reduction: &Reduction) -> bool { + match reduction.group_keys() { + None => false, + Some(keys) => keys.is_without() || !keys.is_empty(), + } +} + +/// Would a build sized to `tighter`'s accuracy requirement also satisfy +/// `looser`'s? A tighter `(eps, delta)` bound implies the looser one, so +/// this is a Pareto check: both +/// sides resolve through [`accuracy_budget`] to concrete `(eps, delta)` +/// numbers, and `tighter` dominates `looser` iff neither of its two numbers +/// is larger. +/// +/// `AccuracyTarget::Exact` on either side always returns `false` — never a +/// dominator, never dominated. Numerically, `accuracy_budget(Exact)` +/// resolves to a budget that would Pareto-dominate everything (zero error), +/// but `Exact` is realized exactly, not by a sketch — a different +/// `Realization` family, not a point on the same sizing curve. +pub(crate) fn dominates(tighter: &AccuracyTarget, looser: &AccuracyTarget) -> bool { + if matches!(tighter, AccuracyTarget::Exact) || matches!(looser, AccuracyTarget::Exact) { + return false; + } + let (tighter_eps, tighter_delta) = accuracy_budget(tighter); + let (looser_eps, looser_delta) = accuracy_budget(looser); + tighter_eps <= looser_eps && tighter_delta <= looser_delta +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn per_entity_has_no_subpopulation_concept() { + assert!(!has_subpopulations(&Reduction::PerEntity)); + } + + #[test] + fn empty_by_reduction_has_no_subpopulation_concept() { + assert!(!has_subpopulations(&Reduction::by(vec![]))); + } + + #[test] + fn non_empty_by_reduction_has_a_subpopulation_concept() { + assert!(has_subpopulations(&Reduction::by(vec![2]))); + } + + #[test] + fn without_grouping_has_a_subpopulation_concept_even_when_empty() { + use asap_types::ir::operator::operator_properties::GroupKeys; + // `without([])` groups by every remaining label — a real + // subpopulation concept, unlike `by([])`'s genuine full reduction. + assert!(has_subpopulations(&Reduction::Reduce(GroupKeys::without( + vec![] + )))); + } +} diff --git a/crates/logical-optimizer/src/pass1/replacement.rs b/crates/logical-optimizer/src/pass1/replacement.rs index d048203c1..a672c5510 100644 --- a/crates/logical-optimizer/src/pass1/replacement.rs +++ b/crates/logical-optimizer/src/pass1/replacement.rs @@ -359,9 +359,8 @@ use asap_types::ir::properties::{ExecutionDataStateError, ExecutionTiming}; use asap_types::ir::scalar::{ArithmeticOpKind, ColumnRef}; use asap_types::ir::schema::{ ColumnId, EntityIdentity, ExactKind, ExactParams, Field, FieldDataType, GroupingStrategy, - NonNegativeWeightProof, SamplingKind, SamplingParams, Schema, SketchAlgorithm, SketchKind, - SketchParams, SketchStatistic as PostAsapSketchStatistic, StatModelKind, StatModelParams, - SummaryInputExpr, SummaryUpdate, WaveletKind, WaveletParams, WeightDomain, + NonNegativeWeightProof, Schema, SketchAlgorithm, SketchKind, SketchParams, + SketchStatistic as PostAsapSketchStatistic, SummaryInputExpr, SummaryUpdate, WeightDomain, }; use asap_types::ir::SchemaDerivationError; use asap_types::ir::{ @@ -381,6 +380,12 @@ use crate::pass1::exact_composition::{ ExactComposition, ExactCompositionStrategy, OperationPlacement, }; use crate::pass1::grouping::HydraGroupingStrategy; +use crate::pass1::realization::{ + accuracy_budget, accuracy_target, column_ref, default_size_params, + realize_keyed_additive_summary_input, schema_column_ref, summarised_column, summarised_input, + summary_candidates, PhysicalSummaryInput, PhysicalSummaryInputRuleResult, Realization, + DEFAULT_DELTA, +}; use crate::pass1::rollup::RollupStrategy; use crate::pass2::reconciliation::AccuracyReconciliationStrategy; use crate::pass2::topk_reuse::TopKLimitReuseStrategy; @@ -655,64 +660,6 @@ pub trait ReplacementStrategy { } } -// ── Realization: how one AggIntent may be realised ─────────────────────── - -/// How an [`AggIntent`] may be realised at post-ASAP binding time (issue -/// #98): by an approximate summary (sketch, sample, wavelet, statistical -/// model, …), by an exact mergeable accumulator, or by an ordinary exact -/// operator (pass-through). This is a post-ASAP concern — the pre-ASAP IR -/// carries only the intent + accuracy target, never the realization — and -/// it's a per-node decision, made once per `AggIntent`, not a plan-wide one. -/// -/// [`realizations_for_intent`] is where every valid realization gets -/// enumerated, exhaustive and ranked (most-preferred first) — this crate has -/// no separate function that computes just "the one" `Realization` -/// independently of that list. [`ASAPStrategies`] is the sole -/// consumer: it wraps every entry of this list into its own bound -/// [`OperatorNode`] and returns all of them, ranked — a caller wanting a -/// single answer keeps the first one itself (see the module docs above). -#[derive(Debug, Clone, PartialEq)] -pub enum Realization { - /// An exact **mergeable** accumulator (partial state ≡ the value - /// itself: `Sum` / `Count` / `Min` / `Max` / `Rate` / `Increase`). The - /// built state *is* the answer already — no `SummaryEstimate` evaluation - /// step. - ExactAggregate { - kind: ExactKind, - params: ExactParams, - }, - /// An approximate sketch sized to the intent's [`AccuracyTarget`]. - /// Needs a `SummaryEstimate` evaluation to recover a value. Already - /// classified into its [`SketchKind`] category (`SketchKind::new` - /// having been called) — construction always goes through that - /// classifier, never this variant directly. - Sketch(SketchKind), - /// A sampling-based summary (a retained row subset). Needs a - /// `SummaryEstimate` evaluation. Not chosen by any core `AggIntent` - /// dispatch today — see the module docs. - Sample { - kind: SamplingKind, - params: SamplingParams, - }, - /// A wavelet-transform summary. Needs a `SummaryEstimate` evaluation. Not - /// chosen by any core `AggIntent` dispatch today — see the module docs. - Wavelet { - kind: WaveletKind, - params: WaveletParams, - }, - /// A fitted statistical/parametric-model summary. Needs a - /// `SummaryEstimate` evaluation. Not chosen by any core `AggIntent` - /// dispatch today — see the module docs. - StatModel { - kind: StatModelKind, - params: StatModelParams, - }, - /// No summary form — the node stays a logical pre-ASAP operator and is - /// executed exactly (per-series transforms, non-mergeable reducers, exact - /// quantile/top-k/cardinality, classic-bucket `HistogramQuantile`, …). - PassThrough, -} - /// Does an already-**available** [`Realization`] — e.g. a summary /// instance a downstream deployment already materialized somewhere, found /// via whatever inventory/index that deployment keeps — satisfy a @@ -752,70 +699,6 @@ pub trait Matcher { fn is_satisfied_by(&self, required: &Realization, available: &Realization) -> bool; } -/// Confidence δ assumed when the target carries only an ε -/// (`AccuracyTarget::Epsilon`): the (ε, δ)-parameterised sketches (CMS) need -/// one. `ln(1/0.01) → depth 5`, matching the conventional CMS sizing. -pub const DEFAULT_DELTA: f64 = 0.01; - -/// The sketch kinds that can serve an intent, most-preferred first. -/// This is the `AggIntent → SketchAlgorithm` map of issue #98; -/// [`realizations_for_intent`] sizes and ranks every entry via `cost_model`. -/// Listed here so the candidate set has one home. -pub fn summary_candidates(intent: &AggIntent) -> &'static [SketchAlgorithm] { - match intent { - AggIntent::Quantile { .. } => &[SketchAlgorithm::Kll, SketchAlgorithm::DDSketch], - // A distinct-tuple count hashes the whole tuple as one item - // (`SummaryInputExpr::Tuple`), which the distinct-count sketches take - // unchanged. UnivMon is dropped there: it estimates frequency moments - // over a single value stream, and `realize_value_frequency_summary_input` - // would feed it one column of the tuple. - AggIntent::Cardinality { cols, .. } if cols.len() > 1 => &[ - SketchAlgorithm::Hll, - SketchAlgorithm::Theta, - SketchAlgorithm::Kmv, - ], - AggIntent::Cardinality { .. } => &[ - SketchAlgorithm::Hll, - SketchAlgorithm::Theta, - SketchAlgorithm::Kmv, - SketchAlgorithm::UnivMon, - ], - AggIntent::FrequencyL2 { .. } | AggIntent::FrequencyEntropy { .. } => { - &[SketchAlgorithm::UnivMon] - } - // Count-Sketch-with-heap is CMS-with-heap's balanced/zero-mean-error - // alternative for the same heavy-hitter shape. - AggIntent::TopK { .. } => &[ - SketchAlgorithm::CmsWithHeap, - SketchAlgorithm::CountSketchWithHeap, - ], - AggIntent::Count { .. } => &[ - SketchAlgorithm::Cms, - SketchAlgorithm::CountSketch, - SketchAlgorithm::UnivMon, - ], - _ => &[], - } -} - -/// The [`AccuracyTarget`] threaded onto an approximate-capable intent -/// (`Quantile`/`Cardinality`/`Count`/`TopK`), or `None` for every other -/// intent (no sketch candidate applies — [`realizations_for_intent`]'s own -/// match routes those elsewhere). Exposed so callers resolve the exact same -/// accuracy target [`realizations_for_intent`] does, without re-deriving it -/// from scratch. -pub fn accuracy_target(intent: &AggIntent) -> Option<&AccuracyTarget> { - match intent { - AggIntent::Quantile { accuracy, .. } - | AggIntent::Cardinality { accuracy, .. } - | AggIntent::FrequencyL2 { accuracy, .. } - | AggIntent::FrequencyEntropy { accuracy, .. } - | AggIntent::Count { accuracy } - | AggIntent::TopK { accuracy, .. } => Some(accuracy), - _ => None, - } -} - /// Every valid [`Realization`] for `intent`, exhaustive, in /// [`summary_candidates`]' static order — the *only* place this crate /// decides what an `AggIntent` may become. Nothing in this crate computes @@ -948,22 +831,6 @@ fn exact_accumulator(intent: &AggIntent, kind: ExactKind, params: ExactParams) - Realization::ExactAggregate { kind, params } } -/// Resolve an [`AccuracyTarget`] into the `(eps, delta)` budget sketch -/// sizing needs. Shared by [`sketch_realizations`] and -/// this crate's own sizing — one place this resolution happens, so nothing -/// can drift apart on it. -/// -/// `Exact` is unreachable via [`realizations_for_intent`] (which routes -/// `Exact` to [`exact_realization`] instead); degrades to the tightest -/// parameters for a caller that resolves it directly anyway. -pub fn accuracy_budget(accuracy: &AccuracyTarget) -> (f64, f64) { - match accuracy { - AccuracyTarget::Exact => (f64::MIN_POSITIVE, DEFAULT_DELTA), - AccuracyTarget::Epsilon(e) => (*e, DEFAULT_DELTA), - AccuracyTarget::EpsilonDelta { epsilon, delta } => (*epsilon, *delta), - } -} - /// Every candidate sketch [`Realization`] for an approximate-capable /// intent, sized analytically to `accuracy`, in [`summary_candidates`]' /// order — [`realizations_for_intent`]'s Sketch branch. @@ -1025,22 +892,6 @@ pub fn sketch_state_bytes(params: &SketchParams) -> Option { .checked_add(u64::from(heap_size).checked_mul(64)?) } -/// `asap-plan`'s built-in `SketchParams` sizing, keyed off the resolved -/// `(eps, delta)` accuracy budget. -/// -/// Each formula inverts the sketch family's standard error bound to the -/// smallest parameter satisfying the target, clamped to the family's sane -/// range. A non-positive ε saturates to the clamp maximum (tightest -/// allowed). -pub fn default_size_params( - kind: SketchAlgorithm, - intent: &AggIntent, - eps: f64, - delta: f64, -) -> SketchParams { - crate::accuracy::estimators::size_params(kind, intent, eps, delta) -} - /// A deployment's explicit bet about how "typical" (non-adversarial) its /// workload's collision pattern is expected to be, consumed only by /// [`posterior_aware_size_params`]. @@ -2685,21 +2536,6 @@ fn summary_family(realization: Realization) -> Option<(FieldDataType, bool)> { }) } -/// The physical input consumed by one summary realization. Most summaries -/// consume the logical aggregate's immediate child and summarize its declared -/// input value. Composite realizations can instead consume a larger -/// logical sub-DAG and bind a different key or value. -pub(crate) struct PhysicalSummaryInput { - pub(crate) child: Rc, - pub(crate) input: SummaryUpdate, -} - -pub(crate) enum PhysicalSummaryInputRuleResult { - NotApplicable, - Realized(PhysicalSummaryInput), - Unsupported(&'static str), -} - type PhysicalSummaryInputRule = fn(&AggIntent, &FieldDataType, &Reduction, &Rc) -> PhysicalSummaryInputRuleResult; @@ -2776,7 +2612,8 @@ fn realize_physical_summary_input( child: Rc::clone(child), input: SummaryUpdate { item: None, - weight: summarised_input(intent, child_schema)?, + weight: summarised_input(intent, child_schema) + .map_err(RealizationError::PhysicalRealization)?, weight_domain: WeightDomain::UnknownOrSigned, }, }) @@ -3456,142 +3293,6 @@ fn realize_current_series_summary_input( }) } -/// Realize the composite heavy-hitter realization for -/// `TopK(Count GROUP BY key)`. The heap sketch consumes the raw keyed stream; -/// it does not consume an independently materialized Count result. -pub(crate) fn realize_keyed_additive_summary_input( - intent: &AggIntent, - family: &FieldDataType, - output_reduction: &Reduction, - child: &Rc, -) -> PhysicalSummaryInputRuleResult { - if !matches!(intent, AggIntent::TopK { .. }) { - return PhysicalSummaryInputRuleResult::NotApplicable; - } - let FieldDataType::Sketch(kind, _) = family else { - return PhysicalSummaryInputRuleResult::NotApplicable; - }; - let heap_algorithm = kind.algorithm(); - if !matches!( - heap_algorithm, - SketchAlgorithm::CmsWithHeap | SketchAlgorithm::CountSketchWithHeap - ) { - return PhysicalSummaryInputRuleResult::NotApplicable; - } - let Some(NonASAPOp::Aggregate { - reduction, - measures, - having: None, - child: raw_child, - .. - }) = child.non_asap() - else { - return PhysicalSummaryInputRuleResult::NotApplicable; - }; - let counter_input = matches!(measures.as_slice(), [AggIntent::Sum { .. }]) - && matches!(raw_child.non_asap(), Some(NonASAPOp::Aggregate { measures, .. }) - if matches!(measures.as_slice(), [AggIntent::Rate | AggIntent::Increase])); - let weight = match measures.as_slice() { - [AggIntent::Count { .. }] => SummaryInputExpr::Constant(1.0), - [AggIntent::Sum { .. }] if counter_input => { - SummaryInputExpr::Column(ColumnRef::SampleValue) - } - [AggIntent::Sum { col }] => SummaryInputExpr::Column(match col { - None => ColumnRef::SampleValue, - Some(index) => match schema_column_ref(raw_child, *index) { - Some(column) => column, - None => { - return PhysicalSummaryInputRuleResult::Unsupported( - "sum-ranked Top-K value column is outside the raw input schema", - ) - } - }, - }), - _ => return PhysicalSummaryInputRuleResult::NotApplicable, - }; - let weight_domain = match measures.as_slice() { - [AggIntent::Count { .. }] => WeightDomain::NonNegative { - proof: NonNegativeWeightProof::UnitCount, - }, - [AggIntent::Sum { .. }] if counter_input => WeightDomain::NonNegative { - proof: NonNegativeWeightProof::ResetAwareCounterDerivative, - }, - _ => WeightDomain::UnknownOrSigned, - }; - if matches!(heap_algorithm, SketchAlgorithm::CmsWithHeap) - && !matches!(weight_domain, WeightDomain::NonNegative { .. }) - { - return PhysicalSummaryInputRuleResult::Unsupported( - "value-weighted CMS requires non-negative update evidence; use CountSketch for arbitrary values", - ); - } - let subpopulation_columns = match output_reduction { - Reduction::PerEntity => vec![], - Reduction::Reduce(keys) => keys - .iter() - .filter_map(|index| schema_column_ref(child, *index)) - .collect(), - }; - let item = match reduction { - Reduction::PerEntity => SummaryInputExpr::EntityIdentity(EntityIdentity::PromqlLabelSet { - excluding: subpopulation_columns, - }), - Reduction::Reduce(keys) if !keys.is_without() && !keys.is_empty() => { - let Some(columns) = keys - .iter() - .map(|index| schema_column_ref(raw_child, *index)) - .collect::>>() - else { - return PhysicalSummaryInputRuleResult::Unsupported( - "ranked item column is outside the raw input schema", - ); - }; - let item_columns: Vec<_> = columns - .into_iter() - .filter(|column| !subpopulation_columns.contains(column)) - .collect(); - match item_columns.as_slice() { - [] => { - return PhysicalSummaryInputRuleResult::Unsupported( - "subpopulation columns consume the complete ranked item identity", - ) - } - [column] => SummaryInputExpr::Column(column.clone()), - _ => SummaryInputExpr::Tuple( - item_columns - .into_iter() - .map(SummaryInputExpr::Column) - .collect(), - ), - } - } - Reduction::Reduce(_) => { - return PhysicalSummaryInputRuleResult::Unsupported( - "an empty or without grouping does not identify ranked items", - ) - } - }; - PhysicalSummaryInputRuleResult::Realized(PhysicalSummaryInput { - child: Rc::clone(raw_child), - input: SummaryUpdate { - item: Some(item), - weight, - weight_domain, - }, - }) -} - -fn schema_column_ref(child: &OperatorNode, index: usize) -> Option { - let column = child.schema.fields.get(index)?; - Some(match &column.table { - Some(table) => ColumnRef::Qualified { - table: table.clone(), - name: column.name.clone(), - }, - None => ColumnRef::Named(column.name.clone()), - }) -} - /// The guarantee of the value a `family` node produces over `child` — /// [`AccuracyModel::propagate`] under the [`CompositionOperator`] this family /// applies to its child's values — checked against `intent`'s own @@ -3698,62 +3399,6 @@ fn summary_col_index(out_schema: &Schema, reduction: &Reduction, measures: usize } } -/// The column fed into a *single-column* summary: the intent's leading -/// positional input resolved to a name against the child schema, or the PromQL -/// sample value when it reads none. Callers are responsible for only reaching -/// here with a one-column intent — [`summarised_input`] is the general form. -fn summarised_column(intent: &AggIntent, child_schema: &Schema) -> ColumnRef { - match intent - .input_cols() - .first() - .and_then(|id| child_schema.fields.get(*id)) - { - Some(c) => column_ref(c), - None => ColumnRef::SampleValue, - } -} - -pub(crate) fn column_ref(column: &Field) -> ColumnRef { - match &column.table { - Some(t) => ColumnRef::Qualified { - table: t.clone(), - name: column.name.clone(), - }, - None => ColumnRef::Named(column.name.clone()), - } -} - -/// What the summary consumes per input row. An intent that reads one column (or -/// none) feeds that column; `COUNT(DISTINCT a, b)` feeds the whole tuple as one -/// item, so the distinct-count sketch hashes `(a, b)` rather than `a` — the -/// difference between tuple cardinality and single-column cardinality. -/// -/// A tuple leg outside the child schema is an error rather than -/// [`summarised_column`]'s sample-value fallback: a leg has no sample-value -/// reading, and silently dropping one would under-count. -pub(crate) fn summarised_input( - intent: &AggIntent, - child_schema: &Schema, -) -> Result { - let cols = intent.input_cols(); - if cols.len() < 2 { - return Ok(SummaryInputExpr::Column(summarised_column( - intent, - child_schema, - ))); - } - let legs = cols - .iter() - .map(|id| child_schema.fields.get(*id).map(column_ref)) - .collect::>>() - .ok_or(RealizationError::PhysicalRealization( - "a tuple column is outside the input schema", - ))?; - Ok(SummaryInputExpr::Tuple( - legs.into_iter().map(SummaryInputExpr::Column).collect(), - )) -} - /// The `SummaryEstimate` evaluation for a summary-bound intent. fn evaluation(intent: &AggIntent, input: &SummaryUpdate) -> PostAsapSketchStatistic { match intent { diff --git a/crates/logical-optimizer/src/pass2/reconciliation.rs b/crates/logical-optimizer/src/pass2/reconciliation.rs index 983999b63..36ad637f4 100644 --- a/crates/logical-optimizer/src/pass2/reconciliation.rs +++ b/crates/logical-optimizer/src/pass2/reconciliation.rs @@ -34,7 +34,7 @@ //! 1. Both are the same bindable shape [`crate::pass1::replacement::ASAPStrategies`] //! itself targets — a single measure, no `HAVING` (`bindable_intent`'s own //! scope) — **and** that one measure is one of the four accuracy-bearing -//! [`AggIntent`] variants ([`crate::pass1::replacement::accuracy_target`]'s own +//! [`AggIntent`] variants ([`crate::pass1::realization::accuracy_target`]'s own //! scope: `Count` / `Quantile` / `Cardinality` / `TopK`). Every other //! intent has no `AccuracyTarget` to reconcile in the first place. //! 2. Same `reduction` (grouping), same `output_names`, and the same shared @@ -61,7 +61,7 @@ //! //! ## Safety of tightening: why reading the tighter build is always sound //! -//! [`crate::pass1::replacement::accuracy_budget`] resolves *every* `AccuracyTarget` +//! [`crate::pass1::realization::accuracy_budget`] resolves *every* `AccuracyTarget` //! (`Epsilon`/`EpsilonDelta`) to the literal `(eps, delta)` pair //! `realizations_for_intent`'s `sketch_realizations` feeds into the //! analytical sizing — the same numbers `default_size_params`' @@ -156,9 +156,9 @@ use asap_types::ir::operator::operator_properties::Reduction; use asap_types::ir::{NonASAPOp, OperatorNode}; use asap_types::types::AccuracyTarget; +use crate::pass1::realization::{accuracy_budget, accuracy_target, dominates}; use crate::pass1::replacement::{ - accuracy_budget, accuracy_target, Replacement, ReplacementProvenance, ReplacementStrategy, - ReplacementSubDAG, TargetSubDAG, + Replacement, ReplacementProvenance, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, }; /// `bindable_accuracy_aggregate`'s return shape: `(reduction, intent, @@ -177,7 +177,7 @@ type BindableAccuracyAggregate<'a> = ( /// [`crate::pass1::replacement::ASAPStrategies`] targets (see that /// module's private `bindable_intent`), further narrowed to a measure whose /// intent actually carries an [`AccuracyTarget`] -/// ([`crate::pass1::replacement::accuracy_target`]'s own scope: `Count` / +/// ([`crate::pass1::realization::accuracy_target`]'s own scope: `Count` / /// `Quantile` / `Cardinality` / `TopK`). `None` for anything else, including /// a multi-measure or `HAVING` aggregate, a non-`Aggregate` node, or an /// accuracy-free intent (`Sum`, `Avg`, …). @@ -224,30 +224,6 @@ fn same_intent_except_accuracy(a: &AggIntent, b: &AggIntent) -> bool { } } -/// Would a build sized to `tighter`'s accuracy requirement also satisfy -/// `looser`'s? See the module docs' "Safety of tightening" section for the -/// full argument; this is the Pareto check that argument reduces to: both -/// sides resolve through [`accuracy_budget`] to concrete `(eps, delta)` -/// numbers, and `tighter` dominates `looser` iff neither of its two numbers -/// is larger. -/// -/// `AccuracyTarget::Exact` on either side always returns `false` — never a -/// dominator, never dominated. Numerically, `accuracy_budget(Exact)` -/// resolves to a budget that would Pareto-dominate everything (zero error), -/// but `realizations_for_intent` realizes `Exact` through a wholly -/// different code path (`exact_realization`, never `sketch_realizations`) -/// — a different `Realization` family, not a point on the same sizing -/// curve — so the numeric comparison alone does not mean what it means for -/// two approximate targets. See the module docs for the full reasoning. -pub(crate) fn dominates(tighter: &AccuracyTarget, looser: &AccuracyTarget) -> bool { - if matches!(tighter, AccuracyTarget::Exact) || matches!(looser, AccuracyTarget::Exact) { - return false; - } - let (tighter_eps, tighter_delta) = accuracy_budget(tighter); - let (looser_eps, looser_delta) = accuracy_budget(looser); - tighter_eps <= looser_eps && tighter_delta <= looser_delta -} - /// `dominates(a, b)` and the two resolved budgets aren't equal — the /// **strict** ordering [`AccuracyReconciliationStrategy`] actually needs. /// Without strictness, two `AccuracyTarget`s that resolve to the identical @@ -530,7 +506,7 @@ mod tests { let epsilon_only = AccuracyTarget::Epsilon(0.01); let equivalent_epsilon_delta = AccuracyTarget::EpsilonDelta { epsilon: 0.01, - delta: crate::pass1::replacement::DEFAULT_DELTA, + delta: crate::pass1::realization::DEFAULT_DELTA, }; assert!(dominates(&epsilon_only, &equivalent_epsilon_delta)); assert!(dominates(&equivalent_epsilon_delta, &epsilon_only)); diff --git a/crates/logical-optimizer/src/pass2/summary_capability.rs b/crates/logical-optimizer/src/pass2/summary_capability.rs index 3c0835bad..1001f8365 100644 --- a/crates/logical-optimizer/src/pass2/summary_capability.rs +++ b/crates/logical-optimizer/src/pass2/summary_capability.rs @@ -6,8 +6,7 @@ //! The rule is all-or-nothing per key (#580 decision W5): every approximate //! target of a key is re-sized for the key's strictest requirement, so their //! summary producers are identical and the shared variant's merge after -//! composition reaches one state. Sizing reuses the legacy reconciliation's -//! argument ([`super::reconciliation`]): requirements resolve through +//! composition reaches one state. Requirements resolve through //! [`accuracy_budget`] and every shipped sizing formula is monotonic in //! `(ε, δ)`, so a summary sized for a requirement that dominates every //! consumer's meets each of them. Stage 3 still checks each query against its @@ -23,11 +22,10 @@ use asap_types::ir::schema::ColumnId; use asap_types::ir::{NonASAPOp, OperatorNode}; use asap_types::types::AccuracyTarget; -use super::reconciliation::dominates; use crate::pass1::logical_candidates::{ local_realizations_for_intent, LocalLogicalCandidates, LogicalCandidateError, }; -use crate::pass1::replacement::{accuracy_budget, accuracy_target}; +use crate::pass1::realization::{accuracy_budget, accuracy_target, dominates}; /// The estimates one summary serves over one column. #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -206,7 +204,7 @@ fn single_intent(inventory: &LocalLogicalCandidates, t: usize) -> &AggIn mod tests { use super::*; use crate::pass1::logical_candidates::enumerate_local_logical_candidates; - use crate::pass1::replacement::Realization; + use crate::pass1::realization::Realization; use crate::test_support::lower_promql; use asap_types::ir::schema::{SketchAlgorithm, SketchParams}; use asap_types::ir::QueryRoot; diff --git a/crates/logical-optimizer/src/pass2/window_composition.rs b/crates/logical-optimizer/src/pass2/window_composition.rs index aa6a825ba..d31ce5ce4 100644 --- a/crates/logical-optimizer/src/pass2/window_composition.rs +++ b/crates/logical-optimizer/src/pass2/window_composition.rs @@ -39,7 +39,7 @@ use super::summary_capability::{strictest, with_accuracy}; use crate::pass1::logical_candidates::{ local_realizations_for_intent, LocalLogicalCandidates, LogicalCandidateError, }; -use crate::pass1::replacement::{accuracy_target, Realization}; +use crate::pass1::realization::{accuracy_target, Realization}; /// How a summary alternative covers its query's window. #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] diff --git a/crates/plan-selection/src/cost/cost_model.rs b/crates/plan-selection/src/cost/cost_model.rs index 315f439ea..783c1363c 100644 --- a/crates/plan-selection/src/cost/cost_model.rs +++ b/crates/plan-selection/src/cost/cost_model.rs @@ -403,7 +403,7 @@ pub trait CostModel { } /// Rank `candidates` (as returned by - /// [`summary_candidates`](asap_logical_optimizer::pass1::replacement::summary_candidates)) for + /// [`summary_candidates`](asap_logical_optimizer::pass1::realization::summary_candidates)) for /// `intent`, best choice first. /// /// Implementations MAY reorder freely, but MUST return exactly the input @@ -771,7 +771,7 @@ fn hydra_grid_cells(params: &HydraParams) -> f64 { /// The default cost model: preserves [`summary_candidates`]'s built-in static /// order. /// -/// [`summary_candidates`]: asap_logical_optimizer::pass1::replacement::summary_candidates +/// [`summary_candidates`]: asap_logical_optimizer::pass1::realization::summary_candidates pub struct DefaultCostModel; impl CostModel for DefaultCostModel { @@ -880,7 +880,7 @@ impl CostModel for DefaultCostModel { #[cfg(test)] mod tests { use super::*; - use asap_logical_optimizer::pass1::replacement::summary_candidates; + use asap_logical_optimizer::pass1::realization::summary_candidates; use asap_types::ir::operator::agg_intent::default_cardinality; #[test] diff --git a/crates/plan-selection/src/cost/empirical_cost.rs b/crates/plan-selection/src/cost/empirical_cost.rs index fc77bb85a..c81904e7b 100644 --- a/crates/plan-selection/src/cost/empirical_cost.rs +++ b/crates/plan-selection/src/cost/empirical_cost.rs @@ -7,9 +7,10 @@ use asap_types::ir::schema::{SketchAlgorithm, SketchParams}; use serde::{Deserialize, Serialize}; use crate::cost::cost_model::{CostModel, DefaultCostModel}; -use asap_logical_optimizer::pass1::replacement::{ - accuracy_budget, accuracy_target, default_size_params, ReplacementSubDAG, TargetSubDAG, +use asap_logical_optimizer::pass1::realization::{ + accuracy_budget, accuracy_target, default_size_params, }; +use asap_logical_optimizer::pass1::replacement::{ReplacementSubDAG, TargetSubDAG}; pub const EVIDENCE_SCHEMA_VERSION: u32 = 1; pub const EVIDENCE_MODEL_VERSION: &str = "empirical-update-cpu-v1"; diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs index cc96fa84a..f9281ee46 100644 --- a/crates/planner/tests/summary_sharing.rs +++ b/crates/planner/tests/summary_sharing.rs @@ -6,7 +6,7 @@ use std::rc::Rc; use asap_frontend_sql::SqlCatalog; use asap_logical_optimizer::accuracy::{AccuracyModel, DefaultAccuracyModel, PropagationStats}; -use asap_logical_optimizer::pass1::replacement::{default_size_params, DEFAULT_DELTA}; +use asap_logical_optimizer::pass1::realization::{default_size_params, DEFAULT_DELTA}; use asap_logical_optimizer::{Replacement, ReplacementSubDAG, TargetSubDAG}; use asap_plan_selection::PlanningModels; use asap_plan_selection::{CostModel, DefaultCostModel}; From 816e755908931169f759f2e564e175deaccea889 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Mon, 5 Oct 2026 04:40:47 +0000 Subject: [PATCH 2/2] fix: integrate with #621: import default_size_params from pass1::realization Co-Authored-By: Claude Opus 5.5 --- crates/executor/src/summary_kernels/univmon.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/executor/src/summary_kernels/univmon.rs b/crates/executor/src/summary_kernels/univmon.rs index dab550427..1e46f3b95 100644 --- a/crates/executor/src/summary_kernels/univmon.rs +++ b/crates/executor/src/summary_kernels/univmon.rs @@ -280,7 +280,7 @@ mod tests { /// sizing it is within 1% of a Zipf stream's true L2 norm. #[test] fn l2_is_layer0_f2_within_the_certified_bound() { - use asap_logical_optimizer::pass1::replacement::default_size_params; + use asap_logical_optimizer::pass1::realization::default_size_params; use planner_types::ir::operator::agg_intent::default_cardinality; use planner_types::ir::schema::{SketchAlgorithm, SketchParams}; let SketchParams::UnivMon {