From 359f62bc321784b8c74e847c58bb5ab05155d8cb Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Thu, 1 Oct 2026 23:23:06 +0000 Subject: [PATCH 1/3] feat: share identical summary producers across queries after Pass 1 Two queries can select structurally identical summary producers that pre-ASAP CSE cannot merge: identical ungrouped expressions (no unique key) and readouts over one producer with equal parameters (p50 and p99 of one window and accuracy reading one KLL). Each was costed and deployed as its own state. Selection now groups finalized summary candidates of different targets whose outermost SummaryAgg interns to the same node, costs that state once against the union of the targets' entries, and offers each member an equal split through CandidateCostOverrides. A class whose shared state has no selectable lifecycle stays independent; if selection leaves any member elsewhere, that class reverts to independent costs and selection runs once more. MajorPass assembles every root, runs share_common_summary_subtrees once across them, and plans each root's lifecycle with the entries of every root that reaches its shared states, so a shared SummaryAgg is one Rc with one lifecycle across QueryLifecyclePlans. PlanOutput is not deduplicated; consumers dedupe deployments by Rc pointer. Co-Authored-By: Claude Opus 5.5 --- crates/asap-aware-mapping/src/pass/major.rs | 95 +++++- crates/asap-aware-mapping/src/pass/mod.rs | 5 + crates/asap-aware-mapping/src/replacement.rs | 2 +- .../src/summary_maintenance_lifecycle.rs | 226 +++++++++++-- crates/planner/tests/summary_sharing.rs | 310 ++++++++++++++++++ 5 files changed, 602 insertions(+), 36 deletions(-) create mode 100644 crates/planner/tests/summary_sharing.rs diff --git a/crates/asap-aware-mapping/src/pass/major.rs b/crates/asap-aware-mapping/src/pass/major.rs index d7aa0c7d7..7fc0fb492 100644 --- a/crates/asap-aware-mapping/src/pass/major.rs +++ b/crates/asap-aware-mapping/src/pass/major.rs @@ -8,14 +8,15 @@ use std::rc::Rc; +use asap_types::post_asap::{share_common_summary_subtrees, SummaryNode}; use asap_types::pre_asap::query_expr::QueryExpr; use asap_types::types::AccuracyTarget; use super::{OptimizationInput, OptimizationPass, OptimizeError, PlanOutput, QueryLifecyclePlan}; use crate::replacement::{default_strategies_with_evidence, search_workload_with_targets}; use crate::summary_maintenance_lifecycle::{ - assemble_selected_dag_with_summary_maintenance_lifecycles, - global_selection_with_summary_maintenance_lifecycles, WorkloadDemand, + global_selection_with_summary_maintenance_lifecycles, plan_assembled_dag, shared_state_cost, + summary_states, WorkloadDemand, }; /// The shipped algorithm. Unit struct: its strategy set is the crate default, @@ -79,13 +80,93 @@ impl OptimizationPass for MajorPass { .workload_entries_by_target(demand.workload, &entry_indices) .map_err(|error| OptimizeError::LifecycleSelection(error.into()))?; - let mut plans = Vec::with_capacity(space.roots.len()); + // Assemble every root, then intern structurally identical summary + // producers across them once, so two queries that selected the same + // `SummaryAgg` reach one `Rc` (consumers dedupe states by pointer). + let mut assembled = Vec::with_capacity(space.roots.len()); for (entry_index, root) in &space.roots { - let plan = assemble_selected_dag_with_summary_maintenance_lifecycles( - &selection, + let dag = selection + .assemble_selected_dag(root) + .map_err(|source| OptimizeError::LifecycleAssembly { + entry_index: *entry_index, + source: source.into(), + })? + .ok_or_else(|| self.missing_group(*entry_index))?; + assembled.push(dag); + } + let interned = + share_common_summary_subtrees(assembled.iter().cloned().enumerate().collect()); + let states: Vec<_> = interned + .iter() + .map(|(_, dag)| summary_states(dag)) + .collect(); + + // A state reached from several roots is planned once against all of + // their entries, in every plan that reaches it, so each plan picks + // the same lifecycle for it. When that union cannot be costed the + // roots keep their own, unshared DAG and entries. + let mut shared_entries: Vec<(Rc, Option>)> = Vec::new(); + for (position, (entry_index, root)) in space.roots.iter().enumerate() { + for state in &states[position] { + if shared_entries.iter().any(|(s, _)| Rc::ptr_eq(s, state)) { + continue; + } + let readers: Vec<_> = (0..space.roots.len()) + .filter(|&other| states[other].iter().any(|s| Rc::ptr_eq(s, state))) + .map(|other| &space.roots[other].1) + .collect(); + if readers.iter().all(|reader| Rc::ptr_eq(reader, root)) { + continue; + } + let mut entries: Vec = readers + .iter() + .flat_map(|reader| bindings[&Rc::as_ptr(reader)].iter().copied()) + .collect(); + entries.sort_unstable(); + entries.dedup(); + let cost = shared_state_cost( + state, + WorkloadDemand { + entry_indices: &entries, + ..demand + }, + lifecycle.now_ms, + lifecycle.horizon, + lifecycle.capabilities, + models.cost, + ) + .map_err(|source| OptimizeError::LifecycleAssembly { + entry_index: *entry_index, + source: source.into(), + })?; + shared_entries.push((Rc::clone(state), cost.map(|_| entries))); + } + } + + let mut plans = Vec::with_capacity(space.roots.len()); + for (position, (entry_index, root)) in space.roots.iter().enumerate() { + let mut entries = bindings[&Rc::as_ptr(root)].clone(); + let mut dag = Rc::clone(&interned[position].1); + for (state, shared) in &shared_entries { + if !states[position].iter().any(|s| Rc::ptr_eq(s, state)) { + continue; + } + match shared { + Some(shared) => entries.extend(shared), + None => { + entries = bindings[&Rc::as_ptr(root)].clone(); + dag = Rc::clone(&assembled[position]); + break; + } + } + } + entries.sort_unstable(); + entries.dedup(); + let plan = plan_assembled_dag( + dag, root, WorkloadDemand { - entry_indices: &bindings[&Rc::as_ptr(root)], + entry_indices: &entries, ..demand }, lifecycle.now_ms, @@ -99,7 +180,7 @@ impl OptimizationPass for MajorPass { })?; plans.push(QueryLifecyclePlan { entry_index: *entry_index, - plan: plan.ok_or_else(|| self.missing_group(*entry_index))?, + plan, }); } Ok(PlanOutput::new(plans)) diff --git a/crates/asap-aware-mapping/src/pass/mod.rs b/crates/asap-aware-mapping/src/pass/mod.rs index 4cba3ee09..3dcd85e2c 100644 --- a/crates/asap-aware-mapping/src/pass/mod.rs +++ b/crates/asap-aware-mapping/src/pass/mod.rs @@ -175,6 +175,11 @@ pub struct QueryLifecyclePlan { /// One plan per workload entry, in `QueryWorkload::entries()` order; /// [`check_contract`] enforces that. +/// +/// Plans are not deduplicated across entries: a summary state that several +/// queries share appears in each of their plans as the same `Rc` (with the +/// same lifecycle), so a consumer that deploys or costs the workload must +/// dedupe deployments by `Rc::ptr_eq` on the summary node. #[derive(Debug, Clone)] #[non_exhaustive] pub struct PlanOutput { diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index b40b7549f..7243d8631 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -4256,7 +4256,7 @@ impl PlanSpace { } /// Lifecycle-aware whole-subplan costs keyed by target and candidate identity. -#[derive(Default)] +#[derive(Default, Clone)] pub(crate) struct CandidateCostOverrides { costs: HashMap<(*const QueryExpr, *const ReplacementSubDAG), Cost>, raw_costs: HashMap<*const QueryExpr, Cost>, diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index c29d79608..818d08af8 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -21,11 +21,11 @@ use std::collections::{HashMap, HashSet}; use std::rc::Rc; use asap_types::post_asap::{ - compile_post_asap_dag_with_node_ids, EvaluationSchedule, ExecutionDataStateError, - ExecutionTiming, OutputRepresentation, PostAsapDag, PostAsapDagValidationError, PostAsapNodeId, - ResultGuarantee, SummaryExpr, SummaryMaintenanceLifecycle, - SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, SummaryNode, - SummaryWindowFramework, ValueOperation, + compile_post_asap_dag_with_node_ids, share_common_summary_subtrees, EvaluationSchedule, + ExecutionDataStateError, ExecutionTiming, OutputRepresentation, PostAsapDag, + PostAsapDagValidationError, PostAsapNodeId, ResultGuarantee, SummaryExpr, + SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, + SummaryNode, SummaryWindowFramework, ValueOperation, }; use asap_types::pre_asap::QueryExpr; use asap_types::types::AccuracyTarget; @@ -723,6 +723,14 @@ fn enumerate_with_profile<'a>( /// selection. The candidate space stays compact; only cost overrides are /// attached, so shared `Rc` identity and exact-composition commitments remain /// the responsibility of `GlobalSelection`. +/// +/// Summary candidates of different targets whose outermost `SummaryAgg` is +/// structurally identical (for example p50 and p99 over one KLL) form a class. +/// When [`shared_state_cost`] can cost that state once against the union of +/// the targets' entries, each member is offered an equal split of it instead +/// of its independent cost. If selection then leaves any member of a class on +/// another choice, that class reverts to independent costs and selection runs +/// once more. pub fn global_selection_with_summary_maintenance_lifecycles<'a, Id>( space: &'a PlanSpace, demand: WorkloadDemand<'_>, @@ -745,6 +753,8 @@ pub fn global_selection_with_summary_maintenance_lifecycles<'a, Id>( )?; let bindings = space.workload_entries_by_target(workload, root_workload_entries)?; let mut costs = CandidateCostOverrides::default(); + // Finalized summary candidates, as sharing-class members. + let mut members = Vec::new(); for group in space.target_subdag_candidates() { let Some(entry_indices) = bindings.get(&Rc::as_ptr(&group.target)) else { continue; @@ -780,11 +790,148 @@ pub fn global_selection_with_summary_maintenance_lifecycles<'a, Id>( if let Some(total) = plan.summary_total_cost { costs.insert(&group.target, candidate, total); } + members.push((group, candidate, Rc::clone(summary))); } } } } - Ok(space.global_selection_with_candidate_costs(cost_model, &profiles, horizon, &costs)?) + + // Intern every member once; members whose outermost state (the + // `SummaryAgg` every other state of the candidate feeds) interns to the + // same node share it. Classes are kept in first-member order. + let interned = share_common_summary_subtrees( + members + .iter() + .enumerate() + .map(|(index, (_, _, summary))| (index, Rc::clone(summary))) + .collect(), + ); + let mut classes: Vec<(Rc, Vec)> = Vec::new(); + for (index, root) in interned { + let states = summary_states(&root); + let Some(state) = states + .iter() + .find(|state| summary_states(state).len() == states.len()) + else { + continue; + }; + if !standalone_populations(&root).is_empty() { + continue; + } + match classes.iter_mut().find(|(s, _)| Rc::ptr_eq(s, state)) { + Some((_, class)) => class.push(index), + None => classes.push((Rc::clone(state), vec![index])), + } + } + let mut shared = Vec::new(); + for (state, class) in classes { + let mut targets: Vec<&Rc> = Vec::new(); + for &index in &class { + let target = &members[index].0.target; + if !targets.iter().any(|t| Rc::ptr_eq(t, target)) { + targets.push(target); + } + } + if targets.len() < 2 { + continue; + } + let mut entries: Vec = targets + .iter() + .flat_map(|target| bindings[&Rc::as_ptr(target)].iter().copied()) + .collect(); + entries.sort_unstable(); + entries.dedup(); + let Some(cost) = shared_state_cost( + &state, + WorkloadDemand { + workload, + data_workload, + entry_indices: &entries, + }, + now_ms, + horizon, + capabilities, + cost_model, + )? + else { + continue; + }; + shared.push((class, Cost(cost.0 / targets.len() as f64))); + } + + let with_shared = |kept: &[(Vec, Cost)]| { + let mut costs = costs.clone(); + for (class, split) in kept { + for &index in class { + let (group, candidate, _) = &members[index]; + costs.insert(&group.target, candidate, *split); + } + } + costs + }; + let selection = space.global_selection_with_candidate_costs( + cost_model, + &profiles, + horizon, + &with_shared(&shared), + )?; + let before = shared.len(); + shared.retain(|(class, _)| { + class.iter().all(|&index| { + let target = &members[index].0.target; + let chosen = selection.for_target(target).and_then(|s| s.chosen); + class.iter().any(|&other| { + Rc::ptr_eq(&members[other].0.target, target) + && chosen.is_some_and(|chosen| std::ptr::eq(chosen, members[other].1)) + }) + }) + }); + if shared.len() == before { + return Ok(selection); + } + Ok(space.global_selection_with_candidate_costs( + cost_model, + &profiles, + horizon, + &with_shared(&shared), + )?) +} + +/// Cost of one `SummaryAgg` state maintained once for every entry in +/// `demand`, or `None` when no lifecycle alternative is selectable for it. +/// No comparison target is supplied: the state serves several queries. +pub(crate) fn shared_state_cost( + state: &Rc, + demand: WorkloadDemand<'_>, + now_ms: u64, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + cost_model: &dyn CostModel, +) -> Result, SummaryMaintenanceLifecyclePlanError> { + Ok(enumerate_with_profile( + Rc::clone(state), + demand, + now_ms, + horizon, + capabilities, + cost_model, + None, + None, + )? + .select_cheapest() + .summary_total_cost) +} + +/// Every unique `SummaryAgg` reachable from `root`. +pub(crate) fn summary_states(root: &Rc) -> Vec> { + let mut states = Vec::new(); + collect_states( + root, + &mut HashSet::new(), + &mut states, + StateKind::SummaryAgg, + ); + states } /// Assemble a globally selected phase-valid DAG and attach workload-aware @@ -801,38 +948,61 @@ pub fn assemble_selected_dag_with_summary_maintenance_lifecycles( selection .assemble_selected_dag(target)? .map(|root| { - let mut plan = enumerate_with_profile( + plan_assembled_dag( root, + target, demand, now_ms, horizon, capabilities, cost_model, - None, - Some(target), - )? - .select_cheapest(); - plan.raw_recompute_total_cost = plan - .expected_reads - .and_then(|reads| cost_model.raw_query_recompute_total_cost(target, reads)); - if !plan.selected_raw_recompute - && plan.raw_recompute_total_cost.is_none_or(|raw| { - plan.summary_total_cost - .is_none_or(|summary| raw.0 <= summary.0) - }) - { - plan.root = crate::replacement::keep_pre_asap(target)?; - plan.deployments.clear(); - plan.selected_raw_recompute = true; - plan.selected_window_implementation_id = None; - plan.summary_total_cost = None; - plan.window_accuracy_guarantee = None; - } - Ok(plan) + ) }) .transpose() } +/// The lifecycle half of +/// [`assemble_selected_dag_with_summary_maintenance_lifecycles`], for a root +/// the caller already assembled (and possibly interned across queries). +pub(crate) fn plan_assembled_dag( + root: Rc, + target: &Rc, + demand: WorkloadDemand<'_>, + now_ms: u64, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + cost_model: &dyn CostModel, +) -> Result { + let mut plan = enumerate_with_profile( + root, + demand, + now_ms, + horizon, + capabilities, + cost_model, + None, + Some(target), + )? + .select_cheapest(); + plan.raw_recompute_total_cost = plan + .expected_reads + .and_then(|reads| cost_model.raw_query_recompute_total_cost(target, reads)); + if !plan.selected_raw_recompute + && plan.raw_recompute_total_cost.is_none_or(|raw| { + plan.summary_total_cost + .is_none_or(|summary| raw.0 <= summary.0) + }) + { + plan.root = crate::replacement::keep_pre_asap(target)?; + plan.deployments.clear(); + plan.selected_raw_recompute = true; + plan.selected_window_implementation_id = None; + plan.summary_total_cost = None; + plan.window_accuracy_guarantee = None; + } + Ok(plan) +} + fn workload_facts( workload: &QueryWorkload, data_workload: Option<&DataWorkload>, diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs new file mode 100644 index 000000000..2346537e7 --- /dev/null +++ b/crates/planner/tests/summary_sharing.rs @@ -0,0 +1,310 @@ +//! Structurally identical summary producers chosen by different queries are +//! shared after Pass 1: one `Rc` across their plans, costed once. + +use std::rc::Rc; + +use asap_aware_mapping::cost_model::Cost; +use asap_aware_mapping::pass::{PlanOutput, PlanningModels}; +use asap_aware_mapping::{ + CostModel, CostRate, DefaultCostModel, Horizon, LifecycleInput, SummaryMaintenanceCapabilities, + SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleCostInputs, +}; +use asap_frontend_sql::SqlCatalog; +use asap_planner::{e2e_plan, FrontendInput, UserInput}; +use asap_types::post_asap::{SketchAlgorithm, SummaryNode}; +use asap_types::pre_asap::schema::{Column, DataType, Schema}; +use asap_types::pre_asap::{AggIntent, QueryExpr}; +use asap_types::types::AccuracyTarget; +use asap_types::workload::{ + AccuracyRequirement, DataArrival, DataWorkload, DurationMs, Evidence, LatencyRequirement, + PlanningWorkload, Predictability, Query, QueryLanguage, QueryRequirements, QueryWorkload, Rate, + RepeatedDemand, RepeatingEntry, RepetitionInterval, SqlDialect, TimeSelection, +}; + +const NOW_MS: u64 = 1_700_000_000_000; +const HORIZON_S: f64 = 3_600.0; + +/// A state costs `build` once however often it is read; raw recomputation +/// costs `raw_per_read` per read. +struct FixedCosts { + build: f64, + raw_per_read: f64, +} + +impl CostModel for FixedCosts { + fn rank_candidates( + &self, + intent: &AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + DefaultCostModel.rank_candidates(intent, candidates) + } + + fn summary_maintenance_lifecycle_cost_inputs( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + SummaryMaintenanceLifecycleCostInputs { + build_cost: Some(Cost(self.build)), + maintenance_cost_per_update: Some(Cost::ZERO), + summary_read_cost: Some(Cost::ZERO), + retention_cost_rate: Some(CostRate(0.0)), + retirement_cost: Some(Cost::ZERO), + } + } + + fn summary_maintenance_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + SummaryMaintenanceCapabilities { + incremental_update: true, + merge: true, + delete: true, + } + } + + fn raw_query_recompute_cost(&self, _target: &QueryExpr) -> Option { + Some(Cost(self.raw_per_read)) + } +} + +/// Summaries are far cheaper than raw recomputation, so every query selects +/// one independently and only sharing is under test. +const CHEAP_SUMMARY: FixedCosts = FixedCosts { + build: 1.0, + raw_per_read: 1_000.0, +}; + +fn requirements(epsilon: f64) -> QueryRequirements { + QueryRequirements { + accuracy: AccuracyRequirement::Explicit(AccuracyTarget::Epsilon(epsilon)), + response_latency: LatencyRequirement::Unspecified, + } +} + +/// Every query repeats every ten minutes: six reads each over the horizon. +fn repeating(query: &str, epsilon: f64) -> RepeatingEntry { + RepeatingEntry { + query: Query(query.into()), + demand: RepeatedDemand::FixedInterval(RepetitionInterval(600_000)), + requirements: requirements(epsilon), + predictability: Predictability::Unknown, + time_selection: TimeSelection::default(), + } +} + +fn lifecycle() -> LifecycleInput { + LifecycleInput::new(NOW_MS, SummaryMaintenanceLifecycleCapabilities::default()) + .with_horizon(Horizon(HORIZON_S)) +} + +async fn plan_promql(queries: &[(&str, f64)], costs: &FixedCosts) -> PlanOutput { + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: None, + repeating_queries: Some( + queries + .iter() + .map(|(query, epsilon)| repeating(query, *epsilon)) + .collect(), + ), + }, + data_workload: Some(DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + data_ingestion_interval: Evidence { + value: Some(DurationMs(15_000)), + ..Default::default() + }, + ingestion_rate: Evidence { + value: Some(Rate(1.0)), + ..Default::default() + }, + ..Default::default() + }), + }; + let input = UserInput::new( + &workload, + FrontendInput::Promql { + now_ms: NOW_MS, + histograms: None, + }, + PlanningModels::builtin().with_cost(costs), + lifecycle(), + ); + e2e_plan(input).await.expect("workload plans") +} + +async fn plan_sql(queries: &[&str], costs: &FixedCosts) -> PlanOutput { + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::SQL(SqlDialect::DataFusionSQL), + query_batch: None, + repeating_queries: Some(queries.iter().map(|query| repeating(query, 0.01)).collect()), + }, + data_workload: Some(DataWorkload { + arrival: DataArrival::AtRest, + ..Default::default() + }), + }; + let catalog = SqlCatalog::new().with_table( + "lineitem", + Schema::new(vec![ + Column::new("l_orderkey", DataType::Int64, false), + Column::new("l_extendedprice", DataType::Float64, false), + ]), + ); + let input = UserInput::new( + &workload, + FrontendInput::Sql { catalog: &catalog }, + PlanningModels::builtin().with_cost(costs), + lifecycle(), + ); + e2e_plan(input).await.expect("workload plans") +} + +/// Every summary state each plan deploys. +fn states(output: &PlanOutput) -> Vec>> { + output + .plans + .iter() + .map(|plan| { + assert!(!plan.plan.selected_raw_recompute, "{:?}", plan.plan.root); + assert!(!plan.plan.deployments.is_empty()); + plan.plan + .deployments + .iter() + .map(|deployment| Rc::clone(&deployment.summary)) + .collect() + }) + .collect() +} + +/// Whether the two plans deploy exactly the same states, by pointer. +fn same_states(states: &[Vec>]) -> bool { + states[0].len() == states[1].len() + && states[0] + .iter() + .zip(&states[1]) + .all(|(left, right)| Rc::ptr_eq(left, right)) +} + +/// The deployments a consumer would run, deduplicated by pointer. +fn unique_deployments(output: &PlanOutput) -> usize { + let mut seen: Vec<*const SummaryNode> = Vec::new(); + for plan in &output.plans { + for deployment in &plan.plan.deployments { + let ptr = Rc::as_ptr(&deployment.summary); + if !seen.contains(&ptr) { + seen.push(ptr); + } + } + } + seen.len() +} + +/// p50 and p99 over the same window and accuracy read one KLL: the +/// equal-params subset of summary capability. Both plans hold the same `Rc` +/// with the same lifecycle, so a consumer maintains it once. +#[tokio::test] +async fn quantiles_with_equal_params_share_one_producer() { + let output = plan_promql( + &[ + ("quantile_over_time(0.5, lat[5m])", 0.01), + ("quantile_over_time(0.99, lat[5m])", 0.01), + ], + &CHEAP_SUMMARY, + ) + .await; + assert!(same_states(&states(&output))); + assert!(!Rc::ptr_eq( + &output.plans[0].plan.root, + &output.plans[1].plan.root + )); + assert_eq!(unique_deployments(&output), 1); + let lifecycles: Vec<_> = output + .plans + .iter() + .map(|plan| { + plan.plan.deployments[0] + .summary_maintenance_lifecycle_guarantee + .clone() + }) + .collect(); + assert_eq!(lifecycles[0], lifecycles[1]); + assert!(lifecycles[0].is_some()); + // Each plan is planned against both queries' reads. + for plan in &output.plans { + assert_eq!(plan.plan.expected_reads, Some(12.0)); + } +} + +/// A different window, a stricter accuracy that changes the sketch's +/// parameters, or a different label selector is a different producer. +#[tokio::test] +async fn different_producers_are_not_shared() { + for queries in [ + [ + ("quantile_over_time(0.5, lat[5m])", 0.01), + ("quantile_over_time(0.99, lat[10m])", 0.01), + ], + [ + ("quantile_over_time(0.5, lat[5m])", 0.01), + ("quantile_over_time(0.99, lat[5m])", 0.001), + ], + [ + ("quantile_over_time(0.5, lat{job=\"a\"}[5m])", 0.01), + ("quantile_over_time(0.99, lat{job=\"b\"}[5m])", 0.01), + ], + ] { + 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 { + assert_eq!(plan.plan.expected_reads, Some(6.0), "{queries:?}"); + } + } +} + +/// An ungrouped aggregate has no unique key, so pre-ASAP CSE keeps the two +/// copies apart; their identical producers (rate, then sum) are shared here. +#[tokio::test] +async fn identical_ungrouped_queries_share_their_producers() { + let query = ("sum(rate(x[5m]))", 0.01); + let output = plan_promql(&[query, query], &CHEAP_SUMMARY).await; + assert!(same_states(&states(&output))); + assert_eq!(unique_deployments(&output), 2); +} + +/// The SQL frontend reaches the same sharing for two copies of one filtered +/// percentile. (Its state schema names the query's output column, so p50 and +/// p99 over one SQL sketch are not yet structurally identical.) +#[tokio::test] +async fn identical_sql_percentiles_share_one_producer() { + let query = + "SELECT approx_percentile_cont(l_extendedprice, 0.5) FROM lineitem WHERE l_orderkey > 10"; + let output = plan_sql(&[query, query], &CHEAP_SUMMARY).await; + assert!(same_states(&states(&output))); + assert_eq!(unique_deployments(&output), 1); +} + +/// A state costs 100 and recomputing a query costs 60 over its six reads: +/// alone, the query recomputes raw. Shared by p50 and p99, the state costs 50 +/// per query, so both keep it. +#[tokio::test] +async fn shared_amortization_alone_can_beat_raw_recompute() { + let costs = FixedCosts { + build: 100.0, + raw_per_read: 10.0, + }; + let p50 = ("quantile_over_time(0.5, lat[5m])", 0.01); + let p99 = ("quantile_over_time(0.99, lat[5m])", 0.01); + + let alone = plan_promql(&[p50], &costs).await; + assert!(alone.plans[0].plan.selected_raw_recompute); + + let output = plan_promql(&[p50, p99], &costs).await; + assert!(same_states(&states(&output))); + assert_eq!(unique_deployments(&output), 1); +} From ee30657d508195db1837b1eef80fcebf6e006ff6 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Fri, 2 Oct 2026 00:01:45 +0000 Subject: [PATCH 2/3] test: cover reverting a sharing class that loses selection Co-Authored-By: Claude Opus 5.5 --- .../src/summary_maintenance_lifecycle.rs | 97 +++++++++++++++++++ 1 file changed, 97 insertions(+) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index 818d08af8..ad0594b93 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -2874,6 +2874,103 @@ mod tests { )); } + /// A state costs 10 however often it is read. Recomputing p50 raw costs + /// 1 and p99 costs 8. + struct P50PrefersRaw; + + impl CostModel for P50PrefersRaw { + fn rank_candidates( + &self, + _intent: &AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + + fn summary_maintenance_lifecycle_cost_inputs( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + SummaryMaintenanceLifecycleCostInputs { + build_cost: Some(Cost(10.0)), + maintenance_cost_per_update: Some(Cost::ZERO), + summary_read_cost: Some(Cost::ZERO), + retention_cost_rate: Some(CostRate(0.0)), + retirement_cost: Some(Cost::ZERO), + } + } + + fn summary_maintenance_capabilities( + &self, + summary: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + UnitCosts.summary_maintenance_capabilities(summary) + } + + fn raw_query_recompute_total_cost( + &self, + target: &QueryExpr, + _expected_reads: f64, + ) -> Option { + match target { + QueryExpr::Aggregate { measures, .. } => match measures[..] { + [AggIntent::Quantile { q: 0.5, .. }] => Some(Cost(1.0)), + _ => Some(Cost(8.0)), + }, + _ => None, + } + } + } + + /// p50 and p99 form a sharing class over one state (5 each), but p50's + /// raw recompute (1) still wins. The class reverts, so p99 is reselected + /// at its independent cost (10) and recomputes raw (8), as it does alone. + /// Checked at selection: the assembled plan's own raw comparison would + /// recompute p99 raw either way. + #[test] + fn sharing_class_reverts_when_a_member_selects_elsewhere() { + let quantile = |q| { + Rc::new(QueryExpr::Aggregate { + reduction: Reduction::by(vec![]), + measures: vec![AggIntent::Quantile { + col: None, + q, + accuracy: AccuracyTarget::Epsilon(0.1), + }], + // A shared output name keeps p50 and p99 on one state. + output_names: vec!["value".into()], + having: None, + child: query_root(), + }) + }; + let workload = workload(vec![], vec![repeating(), repeating()], at_rest()); + // Whether each root selected a summary rather than raw recompute. + let summaries = |space: &PlanSpace<&str>, entries: &[usize]| { + let selection = global_selection_with_summary_maintenance_lifecycles( + space, + WorkloadDemand::new_with_data(&workload, &at_rest(), entries), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &P50PrefersRaw, + ) + .unwrap(); + space + .roots + .iter() + .map(|(_, target)| selection.for_target(target).unwrap().chosen.is_some()) + .collect::>() + }; + + let space = crate::replacement::search_workload(vec![ + ("p50", quantile(0.5)), + ("p99", quantile(0.99)), + ]); + let alone = crate::replacement::search_workload(vec![("p99", quantile(0.99))]); + assert_eq!(summaries(&space, &[0, 1]), vec![false, false]); + assert_eq!(summaries(&alone, &[1]), vec![false]); + } + #[test] fn normalized_workload_drives_plan_space_recurrence_profiles() { let root = query_root(); From 80b2c23a3277f874439a87b344341fd76ac4d66b Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Fri, 2 Oct 2026 00:04:57 +0000 Subject: [PATCH 3/3] test: set the aggregate filters field added by #467 Co-Authored-By: Claude Opus 5.5 --- crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index ad0594b93..156a0f5bb 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -2939,6 +2939,7 @@ mod tests { }], // A shared output name keeps p50 and p99 on one state. output_names: vec!["value".into()], + filters: vec![], having: None, child: query_root(), })