From 8a6e92ec8718d17910b64ec4113491568b0c4381 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Wed, 7 Oct 2026 10:58:46 -0400 Subject: [PATCH 1/2] feat(planner): emit streaming and inference configs from the MILP plan plan_to_planner_output maps each planned deployment's sketch-bench variant to an ASAPQuery aggregation and writes the planner's streaming and inference YAML: exact-sum (sum/count) -> MultipleSum, exact-min/max -> MultipleMinMax, exact-increase -> MultipleIncrease, kll-percall -> DatasketchesKLL, hll -> HLL, cms-heap -> CountMinSketchWithHeap with heapsize = the largest k it serves. Any other variant (e.g. hydra-kll), avg queries, and a query split across deployments are errors. Cleanup is NoCleanup. asap-optimizer-cli --milp gains --output-dir and --allow-undeployable-families (prints the plan, writes no configs). Co-Authored-By: Claude Opus 5.5 --- asap-planner-rs/src/bin/optimizer_cli.rs | 43 +- .../src/optimizer/aqe_extractor.rs | 10 + asap-planner-rs/src/optimizer/milp.rs | 29 +- asap-planner-rs/src/optimizer/milp_output.rs | 422 ++++++++++++++++++ asap-planner-rs/src/optimizer/mod.rs | 2 + asap-planner-rs/src/promql/generator.rs | 4 +- 6 files changed, 490 insertions(+), 20 deletions(-) create mode 100644 asap-planner-rs/src/optimizer/milp_output.rs diff --git a/asap-planner-rs/src/bin/optimizer_cli.rs b/asap-planner-rs/src/bin/optimizer_cli.rs index a0bd9fba..a849266f 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -8,8 +8,8 @@ use std::path::PathBuf; use asap_planner::optimizer::{ build_milp_workload, load_flat_atomic_cost_table, load_optional_selected_atomic_cost_table, - load_workload_facts, run_greedy_pipeline, solve_milp, AtomicCostTable, LabelSetFacts, - LabelSetFactsError, + load_workload_facts, plan_to_planner_output, run_greedy_pipeline, solve_milp, AtomicCostTable, + LabelSetFacts, LabelSetFactsError, }; use asap_planner::ControllerConfig; use clap::Parser; @@ -59,11 +59,24 @@ struct Args { )] atomic_cost_workload: Option, - /// Plan with sketch-bench's rqe-optimizer MILP and print the plan; writes - /// no configs yet. + /// Plan with sketch-bench's rqe-optimizer MILP and print the plan. #[arg(long)] milp: bool, + /// MILP only. Write `streaming_config.yaml` and `inference_config.yaml` + /// for the plan here. + #[arg(long = "output-dir", requires = "milp")] + output_dir: Option, + + /// MILP only. Also plan with families the engine can't deploy; prints the + /// plan and writes no configs. + #[arg( + long = "allow-undeployable-families", + requires = "milp", + conflicts_with = "output_dir" + )] + allow_undeployable_families: bool, + /// MILP only. YAML workload facts: per metric, `cardinality` per label /// set, including the set of all its labels (the series count). #[arg( @@ -188,7 +201,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 solution = solve_milp(&workload, &facts, &costs, objective)?; + let solution = solve_milp( + &workload, + &facts, + &costs, + objective, + args.allow_undeployable_families, + )?; println!("=== Deployments: {} ===", solution.deployments.len()); for (d, planned) in solution.deployments.iter().enumerate() { @@ -237,6 +256,20 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { c.memory_bytes / 1e6 ); } + + if let Some(dir) = &args.output_dir { + let output = plan_to_planner_output(config, &workload, &solution)?; + std::fs::create_dir_all(dir)?; + std::fs::write( + dir.join("streaming_config.yaml"), + output.to_streaming_yaml_string()?, + )?; + std::fs::write( + dir.join("inference_config.yaml"), + output.to_inference_yaml_string()?, + )?; + println!("\nwrote configs to {}", dir.display()); + } Ok(()) } diff --git a/asap-planner-rs/src/optimizer/aqe_extractor.rs b/asap-planner-rs/src/optimizer/aqe_extractor.rs index 7f1a51fd..922dc02b 100644 --- a/asap-planner-rs/src/optimizer/aqe_extractor.rs +++ b/asap-planner-rs/src/optimizer/aqe_extractor.rs @@ -157,6 +157,16 @@ fn decompose_to_leaves(query: &str) -> Vec { leaves } +/// Whether any leaf of `query` is an `avg` that `decompose_to_leaves` rewrites. +pub(crate) fn contains_avg(query: &str) -> bool { + match parse_binary_arms(query) { + Some((lhs, rhs)) => [lhs, rhs] + .into_iter() + .any(|arm| matches!(arm, BinaryArm::Query(q) if contains_avg(&q))), + None => rewrite_avg(query).len() > 1, + } +} + /// Rewrite an `avg` leaf into the sum and count leaves it is computed from: /// `avg by (l) (x)` → `sum by (l) (x)`, `count by (l) (x)`, and /// `avg_over_time(x[5m])` → `sum_over_time(x[5m])`, `count_over_time(x[5m])`. diff --git a/asap-planner-rs/src/optimizer/milp.rs b/asap-planner-rs/src/optimizer/milp.rs index 42183d66..9aae3f64 100644 --- a/asap-planner-rs/src/optimizer/milp.rs +++ b/asap-planner-rs/src/optimizer/milp.rs @@ -46,14 +46,17 @@ pub enum MilpError { } /// The cheapest plan. Only families in sketch-bench's `DEPLOYABLE_FAMILIES` -/// are candidates. +/// are candidates unless `allow_undeployable_families`, whose plan can be +/// studied but not deployed. pub fn solve_milp( workload: &MilpWorkload, facts: &WorkloadFacts, costs: &[AtomicCostEntry], objective: Objective, + allow_undeployable_families: bool, ) -> Result { - let deployments = build_all_candidates(&workload.raqes, costs, facts, false); + let deployments = + build_all_candidates(&workload.raqes, costs, facts, allow_undeployable_families); tracing::debug!( candidates = deployments.len(), cost_rows = costs.len(), @@ -251,13 +254,13 @@ fn accuracy_target( } #[cfg(test)] -mod tests { +pub(super) mod tests { use std::collections::BTreeMap; use super::*; use crate::optimizer::workload_facts::parse_workload_facts; - const SCRAPE_MS: u64 = 15_000; + pub(crate) const SCRAPE_MS: u64 = 15_000; const FACTS: &str = r#" metrics: @@ -271,20 +274,20 @@ metrics: cardinality: 1 "#; - fn config(groups: &str) -> ControllerConfig { + pub(crate) fn config(groups: &str) -> ControllerConfig { let yaml = format!( "query_groups:\n{groups}\nmetrics:\n - metric: http_requests_total\n labels: [instance, job]\n" ); serde_yaml::from_str(&yaml).unwrap() } - fn group(query: &str, accuracy_sla: f64) -> String { + pub(crate) fn group(query: &str, accuracy_sla: f64) -> String { format!( " - queries: [\"{query}\"]\n repetition_delay_ms: 60000\n controller_options: {{accuracy_sla: {accuracy_sla}}}\n" ) } - fn facts(config: &ControllerConfig) -> WorkloadFacts { + pub(crate) fn facts(config: &ControllerConfig) -> WorkloadFacts { parse_workload_facts(FACTS, config.metrics.as_deref().unwrap(), SCRAPE_MS).unwrap() } @@ -311,7 +314,7 @@ metrics: } } - fn costs() -> Vec { + pub(crate) fn costs() -> Vec { vec![ cost("exact-sum", &[(RELATIVE_ERROR, 0.0)]), cost("kll-percall", &[(MAX_RANK_ERROR, 0.005)]), @@ -469,7 +472,7 @@ 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 solution = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap(); + let solution = solve_milp(&w, &facts, &costs(), Objective::default(), false).unwrap(); assert_eq!(solution.deployments.len(), 1); assert_eq!(solution.raqes[0].deployment, solution.raqes[1].deployment); assert_eq!( @@ -487,7 +490,7 @@ metrics: )); let facts = facts(&config); let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); - let err = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap_err(); + let err = solve_milp(&w, &facts, &costs(), Objective::default(), false).unwrap_err(); assert!(matches!(err, MilpError::Unservable(ids) if ids.len() == 1)); } @@ -498,7 +501,7 @@ metrics: 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(); + let err = solve_milp(&w, &facts, &costs(), Objective::default(), false).unwrap_err(); assert!(matches!(err, MilpError::Unservable(ids) if ids == [w.raqes[0].id.clone()])); } @@ -510,7 +513,7 @@ metrics: let config = config(&groups); let facts = facts(&config); let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); - let solution = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap(); + let solution = solve_milp(&w, &facts, &costs(), Objective::default(), false).unwrap(); let chosen: BTreeMap = w .raqes .iter() @@ -538,7 +541,7 @@ metrics: 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(); + let solution = solve_milp(&w, &facts, &costs(), Objective::default(), false).unwrap(); assert_eq!(solution.deployments.len(), 2); } } diff --git a/asap-planner-rs/src/optimizer/milp_output.rs b/asap-planner-rs/src/optimizer/milp_output.rs new file mode 100644 index 00000000..33b6bee3 --- /dev/null +++ b/asap-planner-rs/src/optimizer/milp_output.rs @@ -0,0 +1,422 @@ +//! Turns a solved MILP plan into the planner's streaming and inference YAML. + +use std::collections::HashMap; + +use asap_types::enums::{CleanupPolicy, WindowType}; +use indexmap::map::Entry; +use indexmap::IndexMap; +use promql_parser::parser::{token, Expr}; +use promql_utilities::data_model::KeyByLabelNames; +use promql_utilities::query_logics::enums::AggregationType; +use rqe_optimizer::milp::MilpSolution; +use rqe_optimizer::{Capability, Deployment}; +use serde_json::Value; +use thiserror::Error; + +use crate::config::input::ControllerConfig; +use crate::error::ControllerError; +use crate::generator::{GeneratorOutput, QueryPlanEntry}; +use crate::planner::agg_config::IntermediateAggConfig; +use crate::planner::labels::set_subpopulation_labels; +use crate::planner_output::PlannerOutput; +use crate::promql::generator::{build_inference_yaml, build_streaming_yaml}; + +use super::aqe_extractor::contains_avg; +use super::milp::MilpWorkload; +use super::solution::OptimizerItem; + +#[derive(Debug, Error)] +pub enum MilpOutputError { + #[error("query {0:?} uses avg, which the engine can't answer from sum and count yet")] + AvgQuery(String), + #[error("sketch-bench variant {variant} ({capability:?}) has no ASAPQuery aggregation type")] + UndeployableVariant { + variant: String, + capability: Capability, + }, + #[error("sketch-bench variant {variant}: sketch_config has no numeric {param:?}")] + MissingParam { + variant: String, + param: &'static str, + }, + #[error("topk query {0:?} has no literal k")] + TopkWithoutK(String), + #[error("query {0:?} is served by two deployments; a query string can name only one")] + QueryOnTwoDeployments(String), + #[error(transparent)] + Generator(#[from] ControllerError), +} + +/// Streaming and inference YAML for `solution`. Aggregation ids follow the +/// plan's deployment order, starting at 1. No cleanup policy yet, so the +/// retained instance counts are not emitted. +pub fn plan_to_planner_output( + config: &ControllerConfig, + workload: &MilpWorkload, + solution: &MilpSolution, +) -> Result { + if let Some(query) = config + .query_groups + .iter() + .flat_map(|group| &group.queries) + .find(|query| contains_avg(query)) + { + return Err(MilpOutputError::AvgQuery(query.clone())); + } + + let item_of = |raqe: usize| &workload.items[workload.raqe_items[raqe]]; + let mut served: Vec> = vec![Vec::new(); solution.deployments.len()]; + for (raqe, planned) in solution.raqes.iter().enumerate() { + served[planned.deployment].push(item_of(raqe)); + } + + let mut aggregations: IndexMap = IndexMap::new(); + let mut deployment_keys = Vec::with_capacity(served.len()); + for (planned, items) in solution.deployments.iter().zip(&served) { + let aggregation = aggregation_config(&planned.deployment, items)?; + let key = aggregation.identifying_key(); + aggregations.entry(key.clone()).or_insert(aggregation); + deployment_keys.push(key); + } + + let mut queries: IndexMap = IndexMap::new(); + for (raqe, planned) in solution.raqes.iter().enumerate() { + let key = &deployment_keys[planned.deployment]; + for query in &item_of(raqe).query_strings { + match queries.entry(query.clone()) { + Entry::Vacant(entry) => { + entry.insert(QueryPlanEntry::fully_planned( + query.clone(), + vec![(key.clone(), None)], + )); + } + Entry::Occupied(entry) if entry.get().aggregation_keys[0].0 == *key => {} + Entry::Occupied(_) => { + return Err(MilpOutputError::QueryOnTwoDeployments(query.clone())) + } + } + } + } + + let id_map: HashMap = aggregations + .keys() + .enumerate() + .map(|(index, key)| (key.clone(), index as u32 + 1)) + .collect(); + let schema = config.schema_from_hints(); + Ok(PlannerOutput::from_output(GeneratorOutput { + punted_queries: Vec::new(), + streaming_yaml: build_streaming_yaml(&aggregations, &id_map, &schema)?, + inference_yaml: build_inference_yaml(CleanupPolicy::NoCleanup, &queries, &id_map, &schema)?, + aggregation_count: aggregations.len(), + query_count: queries.len(), + })) +} + +/// `items` are the optimizer items of the Raqes `deployment` serves; they +/// share its capability and grouping, so the first one decides the labels. +fn aggregation_config( + deployment: &Deployment, + items: &[&OptimizerItem], +) -> Result { + let variant = deployment.config.sketch.as_str(); + let param = |name: &'static str| { + deployment.config.sketch_config["params"][name] + .as_u64() + .ok_or_else(|| MilpOutputError::MissingParam { + variant: variant.to_string(), + param: name, + }) + }; + let (aggregation_type, sub_type, parameters) = match (variant, deployment.capability) { + ("exact-sum", Capability::Sum) => (AggregationType::MultipleSum, "sum", vec![]), + ("exact-sum", Capability::Count) => (AggregationType::MultipleSum, "count", vec![]), + ("exact-min", Capability::Min) => (AggregationType::MultipleMinMax, "min", vec![]), + ("exact-max", Capability::Max) => (AggregationType::MultipleMinMax, "max", vec![]), + ("exact-increase", Capability::RateOrIncrease) => { + (AggregationType::MultipleIncrease, "", vec![]) + } + ("kll-percall", Capability::Quantile) => ( + AggregationType::DatasketchesKLL, + "", + vec![("K", param("k")?)], + ), + ("hll", Capability::Cardinality) => ( + AggregationType::HLL, + "", + vec![("precision", param("lg_k")?)], + ), + ( + "cms-heap-topk-fastpath-vector2d", + capability @ (Capability::TopKByValue | Capability::TopKByCount), + ) => ( + AggregationType::CountMinSketchWithHeap, + if capability == Capability::TopKByCount { + "count" + } else { + "sum" + }, + vec![ + ("depth", param("rows")?), + ("width", param("cols")?), + ("heapsize", max_topk_k(items)?), + ], + ), + (_, capability) => { + return Err(MilpOutputError::UndeployableVariant { + variant: variant.to_string(), + capability, + }) + } + }; + + let requirements = &items + .first() + .expect("every planned deployment serves a Raqe") + .requirements; + let mut grouping_labels = KeyByLabelNames::empty(); + let mut aggregated_labels = KeyByLabelNames::empty(); + set_subpopulation_labels( + requirements.statistics[0], + aggregation_type, + &requirements.grouping_labels, + requirements.topk_by_labels.as_ref(), + &mut KeyByLabelNames::empty(), + &mut grouping_labels, + &mut aggregated_labels, + ); + + Ok(IntermediateAggConfig { + aggregation_type, + aggregation_sub_type: sub_type.to_string(), + window_type: if deployment.window_ms == deployment.slide_ms { + WindowType::Tumbling + } else { + WindowType::Sliding + }, + window_size_ms: deployment.window_ms, + slide_interval_ms: deployment.slide_ms, + spatial_filter: deployment.spatial_filter.clone(), + metric: deployment.metric.clone(), + table_name: None, + value_column: None, + parameters: parameters + .into_iter() + .map(|(name, value)| (name.to_string(), Value::from(value))) + .collect(), + rollup_labels: KeyByLabelNames::empty(), + grouping_labels, + aggregated_labels, + }) +} + +/// One heap serves every topk query on the deployment, so it is sized for +/// the largest k. +fn max_topk_k(items: &[&OptimizerItem]) -> Result { + items + .iter() + .flat_map(|item| &item.query_strings) + .map(|query| topk_k(query).ok_or_else(|| MilpOutputError::TopkWithoutK(query.clone()))) + .try_fold(0, |max, k| Ok(max.max(k?))) +} + +fn topk_k(query: &str) -> Option { + let Ok(Expr::Aggregate(aggregate)) = promql_parser::parser::parse(query) else { + return None; + }; + if aggregate.op.id() != token::T_TOPK { + return None; + } + match aggregate.param.as_deref()? { + Expr::NumberLiteral(number) if number.val >= 1.0 => Some(number.val as u64), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use asap_types::aggregation_config::AggregationConfig; + use asap_types::enums::QueryLanguage; + use rqe_optimizer::milp::Objective; + use serde_json::json; + + use super::*; + use crate::optimizer::milp::tests::{config, costs, facts, group, SCRAPE_MS}; + use crate::optimizer::milp::{build_milp_workload, solve_milp}; + + fn plan(groups: &str) -> (ControllerConfig, MilpWorkload, MilpSolution) { + let config = config(groups); + let facts = facts(&config); + let workload = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); + let mut costs = costs(); + for cost in &mut costs { + cost.sketch_config["params"] = match cost.sketch.as_str() { + "kll-percall" => json!({"k": 200}), + "cms-heap-topk-fastpath-vector2d" => json!({"rows": 3, "cols": 2048}), + _ => json!({}), + }; + } + let solution = solve_milp(&workload, &facts, &costs, Objective::default(), false).unwrap(); + (config, workload, solution) + } + + /// The aggregation each query names, parsed back the way the engine reads it. + fn aggregation_for(output: &PlannerOutput, query: &str) -> AggregationConfig { + let inference = output.to_inference_config(QueryLanguage::promql).unwrap(); + let streaming = output.to_streaming_config(QueryLanguage::promql).unwrap(); + let query_config = inference + .query_configs + .iter() + .find(|q| q.query == query) + .unwrap_or_else(|| panic!("no query_config for {query}")); + let [reference] = query_config.aggregations.as_slice() else { + panic!("{query} names {:?}", query_config.aggregations); + }; + streaming + .get_aggregation_config(reference.aggregation_id) + .unwrap() + .clone() + } + + #[test] + fn plan_round_trips_into_engine_configs() { + let sum = "sum by (job) (http_requests_total)"; + let count = "sum by (job) (count_over_time(http_requests_total[1m]))"; + let quantile = "quantile_over_time(0.99, http_requests_total[5m])"; + let topk_value = "topk(5, sum_over_time(http_requests_total[1m]))"; + let topk_count = "topk(5, count_over_time(http_requests_total[1m]))"; + let groups: String = [sum, count, quantile, topk_value, topk_count] + .iter() + .map(|q| group(q, 0.99)) + .collect(); + let (config, workload, solution) = plan(&groups); + let output = plan_to_planner_output(&config, &workload, &solution).unwrap(); + assert_eq!(output.streaming_aggregation_count(), 5); + + let a = aggregation_for(&output, sum); + assert_eq!(a.aggregation_type, AggregationType::MultipleSum); + assert_eq!(a.aggregation_sub_type, "sum"); + assert_eq!( + a.aggregated_labels, + KeyByLabelNames::new(vec!["job".into()]) + ); + assert!(a.grouping_labels.is_empty()); + + let a = aggregation_for(&output, count); + assert_eq!(a.aggregation_type, AggregationType::MultipleSum); + assert_eq!(a.aggregation_sub_type, "count"); + + let a = aggregation_for(&output, quantile); + assert_eq!(a.aggregation_type, AggregationType::DatasketchesKLL); + assert_eq!(a.parameters["K"], json!(200)); + assert_eq!( + a.grouping_labels, + KeyByLabelNames::new(vec!["instance".into(), "job".into()]) + ); + + let a = aggregation_for(&output, topk_value); + assert_eq!(a.aggregation_type, AggregationType::CountMinSketchWithHeap); + assert_eq!(a.aggregation_sub_type, "sum"); + assert_eq!(a.parameters["depth"], json!(3)); + assert_eq!(a.parameters["width"], json!(2048)); + assert_eq!(a.parameters["heapsize"], json!(5)); + assert_eq!( + aggregation_for(&output, topk_count).aggregation_sub_type, + "count" + ); + + for query in [sum, count, quantile, topk_value, topk_count] { + let a = aggregation_for(&output, query); + assert_eq!(a.metric, "http_requests_total"); + // No cleanup policy yet: retention is unbounded. + assert_eq!(a.num_aggregates_to_retain, None); + } + } + + #[test] + fn shared_heap_is_sized_for_the_largest_k() { + let small = "topk(5, sum_over_time(http_requests_total[1m]))"; + let large = "topk(10, sum_over_time(http_requests_total[1m]))"; + let (config, workload, solution) = plan(&(group(small, 0.99) + &group(large, 0.99))); + assert_eq!(solution.deployments.len(), 1); + let output = plan_to_planner_output(&config, &workload, &solution).unwrap(); + for query in [small, large] { + assert_eq!( + aggregation_for(&output, query).parameters["heapsize"], + json!(10) + ); + } + } + + #[test] + fn window_wider_than_slide_is_sliding() { + let query = "quantile_over_time(0.99, http_requests_total[5m])"; + let (config, workload, mut solution) = plan(&group(query, 0.99)); + let deployment = &mut solution.deployments[0].deployment; + deployment.window_ms = 300_000; + deployment.slide_ms = 60_000; + let output = plan_to_planner_output(&config, &workload, &solution).unwrap(); + let a = aggregation_for(&output, query); + assert_eq!(a.window_type, WindowType::Sliding); + assert_eq!((a.window_size_ms, a.slide_interval_ms), (300_000, 60_000)); + } + + #[test] + fn avg_query_is_an_error() { + let query = "avg by (job) (http_requests_total)"; + let (config, workload, solution) = plan(&group(query, 0.99)); + let err = plan_to_planner_output(&config, &workload, &solution) + .err() + .unwrap(); + assert!(matches!(err, MilpOutputError::AvgQuery(q) if q == query)); + } + + #[test] + fn hydra_kll_is_not_deployed() { + // Its DeltaSet key tracker has no ASAPQuery pairing yet. + let (config, workload, mut solution) = plan(&group( + "quantile_over_time(0.99, http_requests_total[5m])", + 0.99, + )); + solution.deployments[0].deployment.config.sketch = "hydra-kll".into(); + let err = plan_to_planner_output(&config, &workload, &solution) + .err() + .unwrap(); + assert!( + matches!(err, MilpOutputError::UndeployableVariant { ref variant, .. } if variant == "hydra-kll") + ); + } + + #[test] + fn query_split_across_deployments_is_an_error() { + // The same query string at two cadences is two items; if they land on + // different deployments, one query_config can't name both. + let query = "sum by (job) (http_requests_total)"; + let slow = group(query, 0.99).replace("60000", "120000"); + let (config, workload, mut solution) = plan(&(group(query, 0.99) + &slow)); + let mut other = solution.deployments[0].clone(); + other.deployment.window_ms *= 2; + other.deployment.slide_ms *= 2; + solution.deployments.push(other); + solution.raqes[1].deployment = 1; + let err = plan_to_planner_output(&config, &workload, &solution) + .err() + .unwrap(); + assert!(matches!(err, MilpOutputError::QueryOnTwoDeployments(q) if q == query)); + } + + #[test] + fn contains_avg_looks_inside_binary_arms() { + assert!(contains_avg( + "rate(http_requests_total[5m]) / avg_over_time(http_requests_total[5m])" + )); + assert!(!contains_avg("sum by (job) (http_requests_total) * 2")); + } + + #[test] + fn topk_k_reads_the_literal() { + assert_eq!(topk_k("topk(7, http_requests_total)"), Some(7)); + assert_eq!(topk_k("topk by (job) (3, http_requests_total)"), Some(3)); + assert_eq!(topk_k("sum(http_requests_total)"), None); + } +} diff --git a/asap-planner-rs/src/optimizer/mod.rs b/asap-planner-rs/src/optimizer/mod.rs index 786aab16..79dc04b2 100644 --- a/asap-planner-rs/src/optimizer/mod.rs +++ b/asap-planner-rs/src/optimizer/mod.rs @@ -7,6 +7,7 @@ pub mod error; pub mod greedy; pub mod label_set_facts; pub mod milp; +pub mod milp_output; pub mod pipeline; pub mod sketch_properties; pub mod solution; @@ -25,6 +26,7 @@ pub use error::{OptimizerError, UnservableItem}; pub use greedy::greedy_assign; pub use label_set_facts::{ItemFacts, LabelSetFacts, LabelSetFactsError, LabelSetKey, SeriesKey}; pub use milp::{build_milp_workload, solve_milp, MilpError, MilpWorkload}; +pub use milp_output::{plan_to_planner_output, MilpOutputError}; pub use pipeline::{run_greedy_pipeline, OptimizerPipelineError}; pub use sketch_properties::{sketch_properties, SketchProperties}; pub use solution::{AQEAssignment, OptimizerItem, OptimizerSolution, QueryMethod}; diff --git a/asap-planner-rs/src/promql/generator.rs b/asap-planner-rs/src/promql/generator.rs index 0481b435..e375f13c 100644 --- a/asap-planner-rs/src/promql/generator.rs +++ b/asap-planner-rs/src/promql/generator.rs @@ -287,7 +287,7 @@ fn collect_binary_leaf_entries( } } -fn build_streaming_yaml( +pub(crate) fn build_streaming_yaml( dedup_map: &IndexMap, id_map: &HashMap, metric_schema: &asap_types::PromQLSchema, @@ -319,7 +319,7 @@ fn build_streaming_yaml( Ok(YamlValue::Mapping(root)) } -fn build_inference_yaml( +pub(crate) fn build_inference_yaml( cleanup_policy: CleanupPolicy, query_plan_map: &IndexMap, id_map: &HashMap, From 0bfda0dcfd3d65358247d81b8a26a50dcf8a86c9 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Wed, 7 Oct 2026 11:36:39 -0400 Subject: [PATCH 2/2] fix(planner): reject avg before solving and write configs atomically - asap-optimizer-cli --output-dir rejects avg queries before the MILP solve; print-only runs still plan them. - Both YAML strings are serialized before either file is written. - Each item is visited once per deployment, not once per occurrence. Co-Authored-By: Claude Opus 5.5 --- asap-planner-rs/src/bin/optimizer_cli.rs | 22 +++++++------- asap-planner-rs/src/optimizer/milp_output.rs | 30 +++++++++++++------- asap-planner-rs/src/optimizer/mod.rs | 2 +- 3 files changed, 33 insertions(+), 21 deletions(-) diff --git a/asap-planner-rs/src/bin/optimizer_cli.rs b/asap-planner-rs/src/bin/optimizer_cli.rs index a849266f..2ca22513 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -8,8 +8,8 @@ use std::path::PathBuf; use asap_planner::optimizer::{ build_milp_workload, load_flat_atomic_cost_table, load_optional_selected_atomic_cost_table, - load_workload_facts, plan_to_planner_output, run_greedy_pipeline, solve_milp, AtomicCostTable, - LabelSetFacts, LabelSetFactsError, + load_workload_facts, plan_to_planner_output, reject_avg_queries, run_greedy_pipeline, + solve_milp, AtomicCostTable, LabelSetFacts, LabelSetFactsError, }; use asap_planner::ControllerConfig; use clap::Parser; @@ -200,6 +200,10 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { let objective = Objective::AUCCost { w_cpu, w_mem }; tracing::debug!(?objective, cost_rows = costs.len(), "milp: inputs loaded"); + // Fail before solving when the plan would be written but can't be. + if args.output_dir.is_some() { + reject_avg_queries(config)?; + } let workload = build_milp_workload(config, &facts, args.data_ingestion_interval_ms)?; let solution = solve_milp( &workload, @@ -259,15 +263,13 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { if let Some(dir) = &args.output_dir { let output = plan_to_planner_output(config, &workload, &solution)?; + // Serialize both before writing either, so a failure can't leave a + // new streaming config next to a stale inference config. + let streaming = output.to_streaming_yaml_string()?; + let inference = output.to_inference_yaml_string()?; std::fs::create_dir_all(dir)?; - std::fs::write( - dir.join("streaming_config.yaml"), - output.to_streaming_yaml_string()?, - )?; - std::fs::write( - dir.join("inference_config.yaml"), - output.to_inference_yaml_string()?, - )?; + std::fs::write(dir.join("streaming_config.yaml"), streaming)?; + std::fs::write(dir.join("inference_config.yaml"), inference)?; println!("\nwrote configs to {}", dir.display()); } Ok(()) diff --git a/asap-planner-rs/src/optimizer/milp_output.rs b/asap-planner-rs/src/optimizer/milp_output.rs index 33b6bee3..0fa8e707 100644 --- a/asap-planner-rs/src/optimizer/milp_output.rs +++ b/asap-planner-rs/src/optimizer/milp_output.rs @@ -1,6 +1,6 @@ //! Turns a solved MILP plan into the planner's streaming and inference YAML. -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use asap_types::enums::{CleanupPolicy, WindowType}; use indexmap::map::Entry; @@ -47,6 +47,19 @@ pub enum MilpOutputError { Generator(#[from] ControllerError), } +/// Errors on the first avg query: its plan can be costed but not deployed. +pub fn reject_avg_queries(config: &ControllerConfig) -> Result<(), MilpOutputError> { + match config + .query_groups + .iter() + .flat_map(|group| &group.queries) + .find(|query| contains_avg(query)) + { + Some(query) => Err(MilpOutputError::AvgQuery(query.clone())), + None => Ok(()), + } +} + /// Streaming and inference YAML for `solution`. Aggregation ids follow the /// plan's deployment order, starting at 1. No cleanup policy yet, so the /// retained instance counts are not emitted. @@ -55,19 +68,16 @@ pub fn plan_to_planner_output( workload: &MilpWorkload, solution: &MilpSolution, ) -> Result { - if let Some(query) = config - .query_groups - .iter() - .flat_map(|group| &group.queries) - .find(|query| contains_avg(query)) - { - return Err(MilpOutputError::AvgQuery(query.clone())); - } + reject_avg_queries(config)?; let item_of = |raqe: usize| &workload.items[workload.raqe_items[raqe]]; + // Each item once per deployment, however many occurrences it has. let mut served: Vec> = vec![Vec::new(); solution.deployments.len()]; + let mut seen = HashSet::new(); for (raqe, planned) in solution.raqes.iter().enumerate() { - served[planned.deployment].push(item_of(raqe)); + if seen.insert((planned.deployment, workload.raqe_items[raqe])) { + served[planned.deployment].push(item_of(raqe)); + } } let mut aggregations: IndexMap = IndexMap::new(); diff --git a/asap-planner-rs/src/optimizer/mod.rs b/asap-planner-rs/src/optimizer/mod.rs index 79dc04b2..f172485a 100644 --- a/asap-planner-rs/src/optimizer/mod.rs +++ b/asap-planner-rs/src/optimizer/mod.rs @@ -26,7 +26,7 @@ pub use error::{OptimizerError, UnservableItem}; pub use greedy::greedy_assign; pub use label_set_facts::{ItemFacts, LabelSetFacts, LabelSetFactsError, LabelSetKey, SeriesKey}; pub use milp::{build_milp_workload, solve_milp, MilpError, MilpWorkload}; -pub use milp_output::{plan_to_planner_output, MilpOutputError}; +pub use milp_output::{plan_to_planner_output, reject_avg_queries, MilpOutputError}; pub use pipeline::{run_greedy_pipeline, OptimizerPipelineError}; pub use sketch_properties::{sketch_properties, SketchProperties}; pub use solution::{AQEAssignment, OptimizerItem, OptimizerSolution, QueryMethod};