From eb03a7bcf1be2f94467d91ffcf2d4dd137bac80c Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Wed, 7 Oct 2026 09:58:23 -0400 Subject: [PATCH 1/2] fix(planner): bump rqe-optimizer to S8 and split Sum/Count/TopK capabilities sketch-bench split SumOrCount into Sum and Count and TopK into TopKByValue and TopKByCount, so the two halves of an avg rewrite, and value- vs count-ranked top-k, no longer share one deployment. A topk with unknown weighting is now an error. solve_milp returns S8's MilpSolution (active deployments with retained count, per-Raqe deployment and merged count); the CLI prints both. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 6 +- asap-planner-rs/src/bin/optimizer_cli.rs | 26 ++--- asap-planner-rs/src/optimizer/milp.rs | 121 +++++++++++++++++------ 3 files changed, 108 insertions(+), 45 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 5d1bdf8c..c8b175c3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -112,7 +112,7 @@ checksum = "7f202df86484c868dbad7eaa557ef785d5c66295e41b460ef922eca0723b842c" [[package]] name = "aqpbm-core" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#373ea6bc51304f30c9c2afaeca5cacd156dfd318" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#3b31dfbdc8d0f3981312cfe3f8c3649787bfc381" dependencies = [ "anyhow", "aqpbm-datagen", @@ -128,7 +128,7 @@ dependencies = [ [[package]] name = "aqpbm-datagen" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#373ea6bc51304f30c9c2afaeca5cacd156dfd318" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#3b31dfbdc8d0f3981312cfe3f8c3649787bfc381" dependencies = [ "rand 0.9.4", "rand_distr", @@ -2388,7 +2388,7 @@ dependencies = [ [[package]] name = "rqe-optimizer" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#373ea6bc51304f30c9c2afaeca5cacd156dfd318" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#3b31dfbdc8d0f3981312cfe3f8c3649787bfc381" dependencies = [ "aqpbm-core", "good_lp", diff --git a/asap-planner-rs/src/bin/optimizer_cli.rs b/asap-planner-rs/src/bin/optimizer_cli.rs index 7a813f8e..a0bd9fba 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -188,16 +188,13 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { tracing::debug!(?objective, cost_rows = costs.len(), "milp: inputs loaded"); let workload = build_milp_workload(config, &facts, args.data_ingestion_interval_ms)?; - let (deployments, solution) = solve_milp(&workload, &facts, &costs, objective)?; - - let mut active: Vec = solution.mapping.clone(); - active.sort(); - active.dedup(); - println!("=== Deployments: {} ===", active.len()); - for &d in &active { - let dep = &deployments[d]; + let solution = solve_milp(&workload, &facts, &costs, objective)?; + + println!("=== Deployments: {} ===", solution.deployments.len()); + for (d, planned) in solution.deployments.iter().enumerate() { + let dep = &planned.deployment; println!( - " [{d}] {:?} {} config={} metric={} grouping={:?} window={}ms slide={}ms", + " [{d}] {:?} {} config={} metric={} grouping={:?} window={}ms slide={}ms retained={} key_tracker={}", dep.capability, dep.config.sketch, dep.config.sketch_config, @@ -205,16 +202,21 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { dep.grouping_labels, dep.window_ms, dep.slide_ms, + planned.retained_instance_count, + dep.key_tracker.is_some(), ); } println!("\n=== Raqes: {} ===", workload.raqes.len()); - for ((raqe, &d), latency_ms) in workload + for ((raqe, planned), latency_ms) in workload .raqes .iter() - .zip(&solution.mapping) + .zip(&solution.raqes) .zip(&solution.plan_cost.query_latency_ms) { - println!(" {} -> [{d}] latency={latency_ms:.3e}ms", raqe.id); + println!( + " {} -> [{}] n={} latency={latency_ms:.3e}ms", + raqe.id, planned.deployment, planned.merged_instance_count + ); } let cost = &solution.plan_cost; println!( diff --git a/asap-planner-rs/src/optimizer/milp.rs b/asap-planner-rs/src/optimizer/milp.rs index 65562336..4ff1696c 100644 --- a/asap-planner-rs/src/optimizer/milp.rs +++ b/asap-planner-rs/src/optimizer/milp.rs @@ -5,8 +5,7 @@ use promql_utilities::query_logics::enums::Statistic; use rqe_optimizer::candidates::{build_all_candidates, eligible_deployments_for}; use rqe_optimizer::milp::{minimize, MilpSolution, Objective}; use rqe_optimizer::{ - validate_facts, AccuracyDirection, AtomicCostEntry, Capability, Deployment, LabelSet, Raqe, - WorkloadFacts, + validate_facts, AccuracyDirection, AtomicCostEntry, Capability, LabelSet, Raqe, WorkloadFacts, }; use thiserror::Error; @@ -37,20 +36,22 @@ pub enum MilpError { query: String, statistics: Vec, }, + #[error("query {query:?}: topk must rank by value (sum_over_time) or count (count_over_time)")] + TopkWeightingUnknown { query: String }, #[error("no eligible deployment for raqes: {0:?}")] Unservable(Vec), #[error("MILP solve failed: {0}")] Solver(String), } -/// Candidate deployments and the MILP's choice among them. Only families in -/// sketch-bench's `DEPLOYABLE_FAMILIES` are candidates. +/// The cheapest plan. Only families in sketch-bench's `DEPLOYABLE_FAMILIES` +/// are candidates. pub fn solve_milp( workload: &MilpWorkload, facts: &WorkloadFacts, costs: &[AtomicCostEntry], objective: Objective, -) -> Result<(Vec, MilpSolution), MilpError> { +) -> Result { let deployments = build_all_candidates(&workload.raqes, costs, facts, false); tracing::debug!( candidates = deployments.len(), @@ -59,7 +60,7 @@ pub fn solve_milp( ); let mut missing = Vec::new(); for raqe in &workload.raqes { - let eligible = eligible_deployments_for(raqe, &deployments).len(); + let eligible = eligible_deployments_for(raqe, &deployments, facts).len(); tracing::debug!(raqe = %raqe.id, eligible, "milp: eligible deployments"); if eligible == 0 { missing.push(raqe.id.clone()); @@ -74,10 +75,10 @@ pub fn solve_milp( objective = objective.value(&solution.plan_cost), cpu_secs_per_sec = solution.plan_cost.cpu_secs_per_sec(), memory_bytes = solution.plan_cost.memory_bytes(), - mapping = ?solution.mapping, + deployments = solution.deployments.len(), "milp: solved" ); - Ok((deployments, solution)) + Ok(solution) } /// The MILP's view of a workload. @@ -166,7 +167,9 @@ fn item_to_raqe(item: &OptimizerItem) -> Result { statistics: req.statistics.clone(), }); }; - let capability = capability(*statistic); + let Some(capability) = capability(*statistic, req.topk_count_events) else { + return Err(MilpError::TopkWeightingUnknown { query }); + }; if !(0.0..=1.0).contains(&item.accuracy_sla) { return Err(MilpError::AccuracySlaOutOfRange { query, @@ -179,7 +182,7 @@ fn item_to_raqe(item: &OptimizerItem) -> Result { // TopK keeps one heap per `topk by` bucket; its `grouping_labels` is the // output label set (every label), which would cost one heap per series. let grouping_labels: LabelSet = match capability { - Capability::TopK => req + Capability::TopKByValue | Capability::TopKByCount => req .topk_by_labels .as_ref() .map(|labels| labels.labels.iter().cloned().collect()) @@ -202,16 +205,23 @@ fn item_to_raqe(item: &OptimizerItem) -> Result { }) } -fn capability(statistic: Statistic) -> Capability { - match statistic { - Statistic::Sum | Statistic::Count => Capability::SumOrCount, +/// `None` for a topk whose weighting is unknown: value-ranked +/// (`sum_over_time`) and count-ranked (`count_over_time`) top-k need separate +/// heaps. +fn capability(statistic: Statistic, topk_count_events: Option) -> Option { + Some(match statistic { + Statistic::Sum => Capability::Sum, + Statistic::Count => Capability::Count, Statistic::Rate | Statistic::Increase => Capability::RateOrIncrease, Statistic::Min => Capability::Min, Statistic::Max => Capability::Max, Statistic::Quantile => Capability::Quantile, Statistic::Cardinality => Capability::Cardinality, - Statistic::Topk => Capability::TopK, - } + Statistic::Topk => match topk_count_events? { + true => Capability::TopKByCount, + false => Capability::TopKByValue, + }, + }) } /// `accuracy_sla` is required accuracy: error metrics must stay within @@ -222,13 +232,14 @@ fn accuracy_target( ) -> (&'static str, f64, AccuracyDirection) { let max_error = 1.0 - accuracy_sla + SLA_EPSILON; match capability { - Capability::SumOrCount + Capability::Sum + | Capability::Count | Capability::Min | Capability::Max | Capability::RateOrIncrease | Capability::Cardinality => (RELATIVE_ERROR, max_error, AccuracyDirection::LowerIsBetter), Capability::Quantile => (MAX_RANK_ERROR, max_error, AccuracyDirection::LowerIsBetter), - Capability::TopK => ( + Capability::TopKByValue | Capability::TopKByCount => ( PRECISION_AT_K, accuracy_sla - SLA_EPSILON, AccuracyDirection::HigherIsBetter, @@ -352,10 +363,10 @@ metrics: } #[test] - fn sum_maps_to_sum_or_count_with_relative_error_ceiling() { + fn sum_maps_to_sum_with_relative_error_ceiling() { let w = workload(&group("sum by (job) (http_requests_total)", 0.99)); let r = &w.raqes[0]; - assert_eq!(r.capability, Capability::SumOrCount); + assert_eq!(r.capability, Capability::Sum); assert_eq!(r.grouping_labels, labels(&["job"])); assert_eq!(r.accuracy_metric, RELATIVE_ERROR); assert_eq!(r.accuracy_direction, AccuracyDirection::LowerIsBetter); @@ -364,6 +375,34 @@ metrics: assert_eq!(r.latency_sla_ms, None); } + #[test] + fn count_and_topk_weighting_map_to_their_own_capabilities() { + let cap = |query| workload(&group(query, 0.99)).raqes[0].capability; + assert_eq!( + cap("sum by (job) (count_over_time(http_requests_total[1m]))"), + Capability::Count + ); + assert_eq!( + cap("topk(5, sum_over_time(http_requests_total[1m]))"), + Capability::TopKByValue + ); + assert_eq!( + cap("topk(5, count_over_time(http_requests_total[1m]))"), + Capability::TopKByCount + ); + } + + #[test] + fn topk_with_unknown_weighting_is_an_error() { + let w = workload(&group("topk(5, http_requests_total)", 0.99)); + let mut item = w.items[0].clone(); + item.requirements.topk_count_events = None; + assert!(matches!( + item_to_raqe(&item), + Err(MilpError::TopkWeightingUnknown { .. }) + )); + } + #[test] fn quantile_uses_max_rank_error() { let w = workload(&group( @@ -380,7 +419,7 @@ metrics: fn topk_requires_precision_and_groups_by_its_buckets() { let w = workload(&group("topk by (job) (5, http_requests_total)", 0.9)); let r = &w.raqes[0]; - assert_eq!(r.capability, Capability::TopK); + assert_eq!(r.capability, Capability::TopKByValue); assert_eq!(r.accuracy_metric, PRECISION_AT_K); assert_eq!(r.accuracy_direction, AccuracyDirection::HigherIsBetter); assert!((r.accuracy_sla - 0.9).abs() < 1e-6); @@ -427,10 +466,13 @@ metrics: let config = config(&(group(query, 0.99) + &group(query, 0.99))); let facts = facts(&config); let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); - let (deployments, solution) = - solve_milp(&w, &facts, &costs(), Objective::default()).unwrap(); - assert_eq!(solution.mapping[0], solution.mapping[1]); - assert_eq!(deployments[solution.mapping[0]].config.sketch, "exact-sum"); + let solution = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap(); + assert_eq!(solution.deployments.len(), 1); + assert_eq!(solution.raqes[0].deployment, solution.raqes[1].deployment); + assert_eq!( + solution.deployments[0].deployment.config.sketch, + "exact-sum" + ); } #[test] @@ -454,16 +496,35 @@ metrics: let config = config(&groups); let facts = facts(&config); let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); - let (deployments, solution) = - solve_milp(&w, &facts, &costs(), Objective::default()).unwrap(); + let solution = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap(); let chosen: BTreeMap = w .raqes .iter() - .zip(&solution.mapping) - .map(|(r, &d)| (r.capability, deployments[d].config.sketch.as_str())) + .zip(&solution.raqes) + .map(|(r, p)| { + let sketch = &solution.deployments[p.deployment].deployment.config.sketch; + (r.capability, sketch.as_str()) + }) .collect(); - assert_eq!(chosen[&Capability::SumOrCount], "exact-sum"); + assert_eq!(chosen[&Capability::Sum], "exact-sum"); assert_eq!(chosen[&Capability::Quantile], "kll-percall"); - assert_eq!(chosen[&Capability::TopK], "cms-heap-topk-fastpath-vector2d"); + assert_eq!( + chosen[&Capability::TopKByValue], + "cms-heap-topk-fastpath-vector2d" + ); + } + + #[test] + fn avg_sum_and_count_get_separate_deployments() { + // One exact-sum accumulator can't answer both halves of the avg + // rewrite; sharing would undercount ingest and memory. + let config = config(&group("avg by (job) (http_requests_total)", 0.99)); + let facts = facts(&config); + let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); + let caps: Vec<_> = w.raqes.iter().map(|r| r.capability).collect(); + assert_eq!(caps.len(), 2); + assert!(caps.contains(&Capability::Sum) && caps.contains(&Capability::Count)); + let solution = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap(); + assert_eq!(solution.deployments.len(), 2); } } From 993b0d74c99e0863cfd55ce279359a6f1e50fbf5 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Wed, 7 Oct 2026 10:52:40 -0400 Subject: [PATCH 2/2] fix(planner): reuse rqe-optimizer's unservable check once per item - The unservable pre-check calls rqe_optimizer::enumerate::unservable instead of re-implementing it, and checks one Raqe per item, since an item's occurrences are identical. - TopkWeightingUnknown no longer implies a bare topk is invalid. Co-Authored-By: Claude Opus 5.5 --- asap-planner-rs/src/optimizer/milp.rs | 34 +++++++++++++++++++-------- 1 file changed, 24 insertions(+), 10 deletions(-) diff --git a/asap-planner-rs/src/optimizer/milp.rs b/asap-planner-rs/src/optimizer/milp.rs index 4ff1696c..42183d66 100644 --- a/asap-planner-rs/src/optimizer/milp.rs +++ b/asap-planner-rs/src/optimizer/milp.rs @@ -2,7 +2,8 @@ //! workload config into `Raqe`s and solves for the cheapest deployments. use promql_utilities::query_logics::enums::Statistic; -use rqe_optimizer::candidates::{build_all_candidates, eligible_deployments_for}; +use rqe_optimizer::candidates::build_all_candidates; +use rqe_optimizer::enumerate::unservable; use rqe_optimizer::milp::{minimize, MilpSolution, Objective}; use rqe_optimizer::{ validate_facts, AccuracyDirection, AtomicCostEntry, Capability, LabelSet, Raqe, WorkloadFacts, @@ -36,7 +37,7 @@ pub enum MilpError { query: String, statistics: Vec, }, - #[error("query {query:?}: topk must rank by value (sum_over_time) or count (count_over_time)")] + #[error("query {query:?}: topk ranking (by value or by sample count) is unknown")] TopkWeightingUnknown { query: String }, #[error("no eligible deployment for raqes: {0:?}")] Unservable(Vec), @@ -58,14 +59,16 @@ pub fn solve_milp( cost_rows = costs.len(), "milp: built candidates" ); - let mut missing = Vec::new(); - for raqe in &workload.raqes { - let eligible = eligible_deployments_for(raqe, &deployments, facts).len(); - tracing::debug!(raqe = %raqe.id, eligible, "milp: eligible deployments"); - if eligible == 0 { - missing.push(raqe.id.clone()); - } - } + // Occurrences of one item are identical Raqes (and contiguous), so + // checking the first of each is enough. + let one_per_item: Vec = workload + .raqes + .iter() + .enumerate() + .filter(|&(i, _)| i == 0 || workload.raqe_items[i] != workload.raqe_items[i - 1]) + .map(|(_, raqe)| raqe.clone()) + .collect(); + let missing = unservable(&one_per_item, &deployments, facts); if !missing.is_empty() { return Err(MilpError::Unservable(missing)); } @@ -488,6 +491,17 @@ metrics: assert!(matches!(err, MilpError::Unservable(ids) if ids.len() == 1)); } + #[test] + fn unservable_query_is_reported_once_per_item() { + let quantile = group("quantile_over_time(0.99, http_requests_total[5m])", 0.999); + let config = config(&(quantile.clone() + &quantile)); + let facts = facts(&config); + let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); + assert_eq!(w.raqes.len(), 2); + let err = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap_err(); + assert!(matches!(err, MilpError::Unservable(ids) if ids == [w.raqes[0].id.clone()])); + } + #[test] fn solve_picks_a_deployable_family_per_capability() { let groups = group("sum by (job) (http_requests_total)", 0.99)