From 589d1412065dd447b141b655ef7f28c90247f421 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Fri, 2 Oct 2026 01:30:03 +0000 Subject: [PATCH] feat: size shared summaries for the strictest consumer Pass 1 also sizes each sketch candidate for the strictest sibling that reads the same summary input (same child, grouping, filters and intent apart from accuracy and quantile rank), so post-ASAP CSE can share one state across consumers with different accuracy targets (#509 summary- capability rule). Co-Authored-By: Claude Opus 5.5 --- crates/asap-aware-mapping/src/replacement.rs | 124 ++++++++++- crates/planner/tests/summary_sharing.rs | 204 ++++++++++++++++++- crates/types/src/post_asap/cse.rs | 7 + 3 files changed, 321 insertions(+), 14 deletions(-) diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index 386bd028..80995599 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -440,10 +440,16 @@ pub enum RealizationError { /// strategy against one node in isolation. A strategy that only cares about /// `root`'s shape (for example, [`SketchAlgorithmStrategy`]) can ignore the /// count; [`SharedSubtreeStrategy`] consults it directly. +/// +/// `strictest_sibling_accuracy` is the strictest accuracy among workload +/// siblings that read the same summary input as `root`, when stricter than +/// `root`'s own. [`search_workload_with`] sets it; [`SketchAlgorithmStrategy`] +/// also sizes a candidate to it. #[derive(Debug, Clone, Copy)] pub struct TargetSubDAG<'a> { pub root: &'a Rc, pub consumer_count: usize, + pub strictest_sibling_accuracy: Option<&'a AccuracyTarget>, } impl<'a> TargetSubDAG<'a> { @@ -453,6 +459,7 @@ impl<'a> TargetSubDAG<'a> { Self { root, consumer_count: 1, + strictest_sibling_accuracy: None, } } @@ -462,6 +469,7 @@ impl<'a> TargetSubDAG<'a> { Self { root, consumer_count, + strictest_sibling_accuracy: None, } } } @@ -1356,7 +1364,7 @@ impl<'a> SketchAlgorithmStrategy<'a> { having: None, child: Rc::clone(child), }); - self.propose_with(&ranked, None) + self.propose_with(&ranked, None, None) } pub(crate) fn from_planning_inputs(planning_inputs: CandidatePlanningInputs<'a>) -> Self { @@ -1365,8 +1373,16 @@ impl<'a> SketchAlgorithmStrategy<'a> { /// The whole enumeration for one target, with `intent_override` /// substituting the target's own intent (only ever its `AccuracyTarget` - /// differs — see [`realize_child_with`]). - fn propose_with(&self, root: &Rc, intent_override: Option<&AggIntent>) -> Proposals { + /// differs — see [`realize_child_with`]). `strictest_sibling` adds each + /// sketch resized to that stricter sibling accuracy (#509 summary + /// capability): alone it only costs more, but post-ASAP CSE shares it + /// with the sibling that needs it. + fn propose_with( + &self, + root: &Rc, + intent_override: Option<&AggIntent>, + strictest_sibling: Option<&AccuracyTarget>, + ) -> Proposals { let mut proposals = Proposals::default(); // A selected logical rewrite otherwise remains KeepPreAsap during DAG // assembly. Also expose its concrete summary realization for selection. @@ -1461,6 +1477,32 @@ impl<'a> SketchAlgorithmStrategy<'a> { ), ); + // Sized for the strictest sibling reading the same summary input, + // when that changes the parameters. + if let (Some(stricter), Realization::Sketch(kind)) = (strictest_sibling, &realization) { + let (eps, delta) = accuracy_budget(stricter); + let algorithm = kind.algorithm().clone(); + let params = + planning_inputs + .cost + .size_params(algorithm.clone(), intent, eps, delta); + if params != *kind.params() { + proposals.record( + format!( + "{rationale}; sized for the strictest sibling consumer {stricter:?}" + ), + construct_summary_with( + root, + &override_accuracy(intent, stricter), + Realization::Sketch(SketchKind::new(algorithm, params)), + planning_inputs, + None, + None, + ), + ); + } + } + // Budget-split alternatives (issue #172, PR 2): re-size this // layer and the approximate child under each allocation of this // node's target across every approximate layer. @@ -1613,7 +1655,7 @@ impl ReplacementStrategy for SketchAlgorithmStrategy<'_> { } fn propose(&self, target: &TargetSubDAG<'_>) -> Proposals { - self.propose_with(target.root, None) + self.propose_with(target.root, None, target.strictest_sibling_accuracy) } /// Heap realizations of an instant-vector ranking (current-series TopK). @@ -1871,7 +1913,7 @@ pub(crate) fn realize_child_with( } }); match SketchAlgorithmStrategy::from_planning_inputs(planning_inputs) - .propose_with(root, overridden.as_ref()) + .propose_with(root, overridden.as_ref(), None) .candidates .into_iter() .next() @@ -6540,6 +6582,74 @@ pub fn search_workload_with_targets<'s, Id>( space } +/// The strictest accuracy among `siblings` that read the same summary input +/// as `root` — same child, grouping and filters, and the same intent apart +/// from its accuracy (and a quantile's rank, a readout parameter) — when +/// stricter than `root`'s own. One summary sized for the strictest consumer +/// serves every sibling: #509's summary-capability rule. +fn strictest_sibling_accuracy( + root: &QueryExpr, + siblings: &[Rc], +) -> Option { + fn approximate(intent: &AggIntent) -> Option<&AccuracyTarget> { + accuracy_target(intent).filter(|accuracy| !matches!(accuracy, AccuracyTarget::Exact)) + } + let QueryExpr::Aggregate { + reduction, + filters, + child, + .. + } = root + else { + return None; + }; + let intent = bindable_intent(root)?; + let own = accuracy_budget(approximate(intent)?); + let (mut eps, mut delta) = own; + for sibling in siblings { + let QueryExpr::Aggregate { + reduction: sibling_reduction, + filters: sibling_filters, + child: sibling_child, + .. + } = sibling.as_ref() + else { + continue; + }; + let Some(other) = bindable_intent(sibling) else { + continue; + }; + let Some(accuracy) = approximate(other) else { + continue; + }; + let same_intent = match (intent, other) { + (AggIntent::Quantile { col, .. }, AggIntent::Quantile { col: other_col, .. }) => { + col == other_col + } + _ => override_accuracy(intent, accuracy) == *other, + }; + if same_intent + && sibling_reduction == reduction + && sibling_filters == filters + && (Rc::ptr_eq(sibling_child, child) || sibling_child == child) + { + let (sibling_eps, sibling_delta) = accuracy_budget(accuracy); + eps = eps.min(sibling_eps); + delta = delta.min(sibling_delta); + } + } + if (eps, delta) == own { + None + } else if delta == DEFAULT_DELTA { + Some(AccuracyTarget::Epsilon(eps)) + } else { + Some(AccuracyTarget::EpsilonDelta { + epsilon: eps, + delta, + }) + } +} + fn cse_workload(roots: Vec<(Id, Rc)>) -> Vec<(Id, Rc)> { // `share_common_subtrees` wants owned `QueryExpr`s, not already-`Rc` // roots — the same `Rc::try_unwrap`-with-clone-fallback pattern @@ -6615,7 +6725,9 @@ fn search_cse_workload_with<'s, Id>( let group = &groups[ptr]; (Rc::clone(&group.target), group.consumer_count) }; - let target = TargetSubDAG::with_consumer_count(&root, consumer_count); + let strictest = strictest_sibling_accuracy(&root, &siblings); + let mut target = TargetSubDAG::with_consumer_count(&root, consumer_count); + target.strictest_sibling_accuracy = strictest.as_ref(); let mut proposed = Vec::new(); let mut rejected = Vec::new(); diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs index c344296b..829c38a9 100644 --- a/crates/planner/tests/summary_sharing.rs +++ b/crates/planner/tests/summary_sharing.rs @@ -3,15 +3,31 @@ use std::rc::Rc; +use asap_aware_mapping::accuracy::{ + AccuracyModel, DefaultAccuracyModel, EqualSplitAllocator, PropagationStats, +}; use asap_aware_mapping::cost_model::Cost; use asap_aware_mapping::pass::{PlanOutput, PlanningModels}; +use asap_aware_mapping::replacement::{default_size_params, DEFAULT_DELTA}; +use asap_aware_mapping::{ + global_selection_with_summary_maintenance_lifecycles, search_workload_with_targets, + ReplacementStrategy, SketchAlgorithmStrategy, WorkloadDemand, +}; use asap_aware_mapping::{ CostModel, CostRate, DefaultCostModel, Horizon, LifecycleInput, SummaryMaintenanceCapabilities, SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleCostInputs, }; +use asap_frontend_promql::lower_promql_workload; use asap_frontend_sql::SqlCatalog; use asap_planner::{e2e_plan, FrontendInput, UserInput}; -use asap_types::post_asap::{SketchAlgorithm, SummaryNode}; +use asap_types::post_asap::{ + share_common_summary_subtrees, AccuracyError, BoundExpr, CompositionOperator, ErrorMetric, + ProbabilityExpr, ResultGuarantee, SketchQuery, +}; +use asap_types::post_asap::{ + SketchAlgorithm, SketchParams, SummaryExpr, SummaryFamilyType, SummaryNode, +}; +use asap_types::pre_asap::agg_intent::default_quantile; use asap_types::pre_asap::schema::{Column, DataType, Schema}; use asap_types::pre_asap::{AggIntent, QueryExpr}; use asap_types::types::AccuracyTarget; @@ -99,8 +115,8 @@ fn lifecycle() -> LifecycleInput { .with_horizon(Horizon(HORIZON_S)) } -async fn plan_promql(queries: &[(&str, f64)], costs: &FixedCosts) -> PlanOutput { - let workload = PlanningWorkload { +fn promql_workload(queries: &[(&str, f64)]) -> PlanningWorkload { + PlanningWorkload { query_workload: QueryWorkload { language: QueryLanguage::PromQL, query_batch: None, @@ -123,7 +139,11 @@ async fn plan_promql(queries: &[(&str, f64)], costs: &FixedCosts) -> PlanOutput }, ..Default::default() }), - }; + } +} + +async fn plan_promql(queries: &[(&str, f64)], costs: &FixedCosts) -> PlanOutput { + let workload = promql_workload(queries); let input = UserInput::new( &workload, FrontendInput::Promql { @@ -240,8 +260,8 @@ async fn quantiles_with_equal_params_share_one_producer() { } } -/// A different window, a stricter accuracy that changes the sketch's -/// parameters, or a different label selector is a different producer. +/// A different window or label selector is a different producer, even when +/// one query asks for a stricter accuracy than the other. #[tokio::test] async fn different_producers_are_not_shared() { for queries in [ @@ -251,7 +271,11 @@ async fn different_producers_are_not_shared() { ], [ ("quantile_over_time(0.5, lat[5m])", 0.01), - ("quantile_over_time(0.99, lat[5m])", 0.001), + ("quantile_over_time(0.99, lat[10m])", 0.001), + ], + [ + ("quantile_over_time(0.5, lat{job=\"a\"}[5m])", 0.01), + ("quantile_over_time(0.99, lat{job=\"b\"}[5m])", 0.001), ], [ ("quantile_over_time(0.5, lat{job=\"a\"}[5m])", 0.01), @@ -261,12 +285,69 @@ async fn different_producers_are_not_shared() { let output = plan_promql(&queries, &CHEAP_SUMMARY).await; assert!(!same_states(&states(&output)), "{queries:?}"); assert_eq!(unique_deployments(&output), 2, "{queries:?}"); - for plan in &output.plans { + for (plan, (_, epsilon)) in output.plans.iter().zip(queries) { assert_eq!(plan.plan.expected_reads, Some(6.0), "{queries:?}"); + assert_eq!(kll_k(plan), kll_k_for(epsilon), "{queries:?}"); } } } +/// The KLL `k` of the one state a plan deploys. +fn kll_k(plan: &asap_aware_mapping::pass::QueryLifecyclePlan) -> u32 { + let [deployment] = plan.plan.deployments.as_slice() else { + panic!("one state: {:?}", plan.plan.deployments.len()); + }; + let SummaryExpr::SummaryAgg { + family: SummaryFamilyType::Sketch(kind, _), + .. + } = &deployment.summary.expr + else { + panic!("sketch state: {:?}", deployment.summary.expr); + }; + let SketchParams::Kll { k } = kind.params() else { + panic!("KLL state: {kind:?}"); + }; + *k +} + +/// The KLL `k` sized for `epsilon`. +fn kll_k_for(epsilon: f64) -> u32 { + let SketchParams::Kll { k } = default_size_params( + SketchAlgorithm::Kll, + &default_quantile(0.5), + epsilon, + DEFAULT_DELTA, + ) else { + unreachable!() + }; + k +} + +/// p50 at ε=0.01 and p99 at ε=0.001 over the same input share one KLL sized +/// for the strictest consumer; each reader's guarantee meets its own target. +#[tokio::test] +async fn quantiles_share_one_producer_sized_for_the_strictest_consumer() { + let p50 = ("quantile_over_time(0.5, lat[5m])", 0.01); + let p99 = ("quantile_over_time(0.99, lat[5m])", 0.001); + assert!(kll_k_for(0.001) > kll_k_for(0.01)); + + let output = plan_promql(&[p50, p99], &CHEAP_SUMMARY).await; + assert!(same_states(&states(&output))); + assert_eq!(unique_deployments(&output), 1); + for (plan, (_, epsilon)) in output.plans.iter().zip([p50, p99]) { + assert_eq!(kll_k(plan), kll_k_for(0.001)); + let guarantee = plan.plan.root.guarantee.as_ref().expect("certified"); + assert!( + guarantee.bound.evaluate().unwrap() <= epsilon, + "{guarantee:?}" + ); + } + + // Alone, the looser query keeps its own, smaller KLL. + let alone = plan_promql(&[p50], &CHEAP_SUMMARY).await; + assert_eq!(kll_k(&alone.plans[0]), kll_k_for(0.01)); +} + /// Cross-series quantiles name their KLL state after the input column, not the /// quantile, so p50 and p99 over one selector share it. #[tokio::test] @@ -369,3 +450,110 @@ async fn shared_amortization_alone_can_beat_raw_recompute() { assert!(same_states(&states(&output))); assert_eq!(unique_deployments(&output), 1); } + +/// Synthetic evidence certifying UnivMon readouts; it exercises sharing, never +/// runtime accuracy. +struct UnivMonEvidence; + +impl AccuracyModel for UnivMonEvidence { + fn local_guarantee( + &self, + family: &SummaryFamilyType, + query: &SketchQuery, + ) -> Option { + if matches!(family, SummaryFamilyType::Sketch(kind, _) if kind.algorithm() == &SketchAlgorithm::UnivMon) + { + 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) + } +} + +/// Distinct count, entropy and L2 over one input, certified by an accuracy +/// model, read one UnivMon state: #515 sharing is the summary-capability rule +/// when the states are identical. `MajorPass` builds candidates with the +/// built-in accuracy model, so this runs its pipeline with the test model. +#[test] +fn certified_frequency_readouts_share_one_univmon_state() { + let queries = [ + ("distinct_over_time(m[5m])", 0.02), + ("entropy_over_time(m[5m])", 0.02), + ("l2_over_time(m[5m])", 0.02), + ]; + let workload = promql_workload(&queries); + let roots = lower_promql_workload(&workload, NOW_MS) + .expect("lowers") + .into_iter() + .zip(queries) + .enumerate() + .map(|(index, (expr, (_, epsilon)))| { + (index, Rc::new(expr), Some(AccuracyTarget::Epsilon(epsilon))) + }) + .collect(); + let strategies: Vec> = + vec![Box::new(SketchAlgorithmStrategy::new_with_planning_inputs( + &CHEAP_SUMMARY, + &UnivMonEvidence, + &EqualSplitAllocator, + ))]; + let space = search_workload_with_targets(roots, &strategies, &UnivMonEvidence); + let entry_indices: Vec = (0..queries.len()).collect(); + let selection = global_selection_with_summary_maintenance_lifecycles( + &space, + WorkloadDemand { + workload: &workload.query_workload, + data_workload: workload.data_workload.as_ref(), + entry_indices: &entry_indices, + }, + NOW_MS, + Some(Horizon(HORIZON_S)), + SummaryMaintenanceLifecycleCapabilities::default(), + &CHEAP_SUMMARY, + ) + .expect("selects"); + let assembled = space + .roots + .iter() + .map(|(index, root)| { + let dag = selection + .assemble_selected_dag(root) + .expect("assembles") + .expect("root has a group"); + (*index, dag) + }) + .collect(); + let mut states: Vec> = Vec::new(); + for (_, root) in share_common_summary_subtrees(assembled) { + assert!(root.guarantee.is_some(), "{:?}", root.expr); + let SummaryExpr::SummaryEstimate { summary_input, .. } = &root.expr else { + panic!("summary readout: {:?}", root.expr); + }; + assert!(matches!( + &summary_input.expr, + SummaryExpr::SummaryAgg { family: SummaryFamilyType::Sketch(kind, _), .. } + if kind.algorithm() == &SketchAlgorithm::UnivMon + )); + states.push(Rc::clone(summary_input)); + } + assert_eq!(states.len(), 3); + assert!(states.iter().all(|state| Rc::ptr_eq(state, &states[0]))); +} diff --git a/crates/types/src/post_asap/cse.rs b/crates/types/src/post_asap/cse.rs index 47b10701..1872243f 100644 --- a/crates/types/src/post_asap/cse.rs +++ b/crates/types/src/post_asap/cse.rs @@ -177,6 +177,13 @@ fn same_node(left: &SummaryNode, right: &SummaryNode) -> bool { /// coercions are performed. All roots must belong to the same data snapshot or /// maintenance scope. Downstream realization must still check physical /// implementation compatibility. Use separate calls for independent executions. +/// +/// When the selected states are identical, this is the planner's +/// summary-capability rule (#509 Pass 2): one summary build node feeds every +/// readout it supports, e.g. one KLL for p50 and p99, or one UnivMon for +/// distinct count, entropy and L2. Candidate generation sizes a variant for +/// the strictest sibling consumer so differing accuracy targets can reach +/// identical states here. pub fn share_common_summary_subtrees( roots: Vec<(Id, Rc)>, ) -> Vec<(Id, Rc)> {