Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions crates/asap-aware-mapping/src/pass/major.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -33,7 +33,8 @@ impl OptimizationPass for MajorPass {
fn optimize(&self, input: OptimizationInput<'_>) -> Result<PlanOutput, OptimizeError> {
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
Expand Down
15 changes: 13 additions & 2 deletions crates/asap-aware-mapping/src/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Box<dyn ReplacementStrategy + 'a>> {
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<Box<dyn ReplacementStrategy + 'a>> {
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,
),
Expand Down
87 changes: 26 additions & 61 deletions crates/planner/tests/summary_sharing.rs
Original file line number Diff line number Diff line change
@@ -1,25 +1,17 @@
//! Structurally identical summary producers chosen by different queries are
//! shared after Pass 1: one `Rc<OperatorNode>` 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::{
Expand Down Expand Up @@ -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<Box<dyn ReplacementStrategy>> =
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<usize> = (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<Rc<OperatorNode>> = 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])));
}
Loading