diff --git a/crates/asap-aware-mapping/src/pass/major.rs b/crates/asap-aware-mapping/src/pass/major.rs index 7be53e0f0..6deffdfcd 100644 --- a/crates/asap-aware-mapping/src/pass/major.rs +++ b/crates/asap-aware-mapping/src/pass/major.rs @@ -13,7 +13,7 @@ use asap_types::ir::OperatorNode; 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::replacement::{default_strategies_with_models, search_workload_with_targets}; use crate::summary_maintenance_lifecycle::{ global_selection_with_summary_maintenance_lifecycles, plan_assembled_dag, shared_state_cost, summary_states, WorkloadDemand, @@ -33,7 +33,8 @@ impl OptimizationPass for MajorPass { fn optimize(&self, input: OptimizationInput<'_>) -> Result { let workload = input.workload; let models = input.models; - let strategies = default_strategies_with_evidence(models.cost, models.evidence); + let strategies = + default_strategies_with_models(models.cost, models.accuracy, models.evidence); // `Id` is the entry's position in `QueryWorkload::entries()`, so the // search result carries the workload binding the lifecycle stage and diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index fecc69025..22693954f 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -6428,18 +6428,29 @@ pub fn default_strategies_with<'a>( pub fn default_strategies_with_evidence<'a>( cost_model: &'a dyn CostModel, evidence: &'a dyn AccuracyEvidenceProvider, +) -> Vec> { + default_strategies_with_models(cost_model, &DEFAULT_ACCURACY_MODEL, evidence) +} + +/// Default planning strategies using the deployment's cost, accuracy, and +/// typed evidence models. Pass 1 must use the same accuracy algebra as final +/// root validation or it can discard candidates the deployment can certify. +pub fn default_strategies_with_models<'a>( + cost_model: &'a dyn CostModel, + accuracy_model: &'a dyn AccuracyModel, + evidence: &'a dyn AccuracyEvidenceProvider, ) -> Vec> { vec![ Box::new(ASAPStrategies::new_with_planning_inputs_and_evidence( cost_model, - &DEFAULT_ACCURACY_MODEL, + accuracy_model, &DEFAULT_ALLOCATOR, evidence, )), Box::new( HydraGroupingStrategy::new_with_planning_inputs_and_evidence( cost_model, - &DEFAULT_ACCURACY_MODEL, + accuracy_model, &DEFAULT_ALLOCATOR, evidence, ), diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs index 7c198493b..9b7a204ef 100644 --- a/crates/planner/tests/summary_sharing.rs +++ b/crates/planner/tests/summary_sharing.rs @@ -1,25 +1,17 @@ //! Structurally identical summary producers chosen by different queries are //! shared after Pass 1: one `Rc` across their plans, costed once. -use asap_types::ir::cse::share_common_sub_dags; use asap_types::ir::{ASAPOp, OperatorNode}; use std::rc::Rc; -use asap_aware_mapping::accuracy::{ - AccuracyModel, DefaultAccuracyModel, EqualSplitAllocator, PropagationStats, -}; +use asap_aware_mapping::accuracy::{AccuracyModel, DefaultAccuracyModel, 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, - ASAPStrategies, ReplacementStrategy, 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::{ @@ -488,72 +480,45 @@ impl AccuracyModel for UnivMonEvidence { } } -/// 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_evaluations_share_one_univmon_state() { +/// The public planner passes its supplied accuracy model into Pass 1 and shares one state. +#[tokio::test] +async fn e2e_certified_frequency_evaluations_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, expr, Some(AccuracyTarget::Epsilon(epsilon)))) - .collect(); - let strategies: Vec> = - vec![Box::new(ASAPStrategies::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, + let output = e2e_plan(UserInput::new( + &workload, + FrontendInput::Promql { + now_ms: NOW_MS, + histograms: None, }, - 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_sub_dags(assembled) { - assert!(root.guarantee.is_some(), "{:?}", root.operator); + PlanningModels::builtin() + .with_cost(&CHEAP_SUMMARY) + .with_accuracy(&UnivMonEvidence), + lifecycle(), + )) + .await + .expect("workload plans"); + assert_eq!(output.plans.len(), 3); + assert_eq!(unique_deployments(&output), 1); + let states = states(&output); + assert!(states.iter().all(|states| states.len() == 1)); + assert!(same_states(&states)); + for (plan, (_, epsilon)) in output.plans.iter().zip(queries) { let asap_types::ir::Operator::ASAP(ASAPOp::SummaryEstimate { summary_input, .. }) = - &root.operator + &plan.plan.root.operator else { - panic!("summary evaluation: {:?}", root.operator); + panic!("summary evaluation: {:?}", plan.plan.root.operator); }; assert!(matches!( &summary_input.operator, asap_types::ir::Operator::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. }) if kind.algorithm() == &SketchAlgorithm::UnivMon )); - states.push(Rc::clone(summary_input)); + let guarantee = plan.plan.root.guarantee.as_ref().expect("certified"); + assert!(DefaultAccuracyModel.satisfies(guarantee, &AccuracyTarget::Epsilon(epsilon))); } - assert_eq!(states.len(), 3); - assert!(states.iter().all(|state| Rc::ptr_eq(state, &states[0]))); }