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..156a0f5bb 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>, @@ -2704,6 +2874,104 @@ 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()], + filters: vec![], + 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(); 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); +}