From b0ce7a71f01fc9ac0bf852681a012c251cf9a733 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Tue, 6 Oct 2026 16:52:31 -0400 Subject: [PATCH 1/3] feat(planner): build rqe-optimizer Raqes and facts from the workload config Converts the workload's optimizer items into sketch-bench `Raqe`s (one per query occurrence), loads per-metric workload facts, and solves with `rqe_optimizer::milp::minimize`. `asap-optimizer-cli --milp` prints the chosen deployments and plan cost; it writes no configs yet. `AtomicCostEntry` is now rqe-optimizer's type instead of a local copy. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 6 +- asap-planner-rs/src/bin/optimizer_cli.rs | 145 +++++- asap-planner-rs/src/optimizer/atomic_costs.rs | 53 +-- asap-planner-rs/src/optimizer/greedy.rs | 2 + asap-planner-rs/src/optimizer/milp.rs | 431 ++++++++++++++++++ asap-planner-rs/src/optimizer/mod.rs | 4 + asap-planner-rs/src/optimizer/pipeline.rs | 50 +- .../src/optimizer/workload_facts.rs | 188 ++++++++ asap-planner-rs/tests/rqe_optimizer_smoke.rs | 66 --- 9 files changed, 796 insertions(+), 149 deletions(-) create mode 100644 asap-planner-rs/src/optimizer/milp.rs create mode 100644 asap-planner-rs/src/optimizer/workload_facts.rs delete mode 100644 asap-planner-rs/tests/rqe_optimizer_smoke.rs diff --git a/Cargo.lock b/Cargo.lock index 231c874e..5d1bdf8c 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#964ccb5811440245a7d75f93e88649e8463eba86" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#373ea6bc51304f30c9c2afaeca5cacd156dfd318" 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#964ccb5811440245a7d75f93e88649e8463eba86" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#373ea6bc51304f30c9c2afaeca5cacd156dfd318" 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#964ccb5811440245a7d75f93e88649e8463eba86" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#373ea6bc51304f30c9c2afaeca5cacd156dfd318" 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 591080fa..ea4db3f1 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -7,10 +7,12 @@ use std::path::PathBuf; use asap_planner::optimizer::{ - load_optional_selected_atomic_cost_table, run_greedy_pipeline, AtomicCostTable, LabelSetFacts, + build_milp_workload, load_optional_selected_atomic_cost_table, load_workload_facts, + run_greedy_pipeline, solve_milp, AtomicCostTable, LabelSetFacts, }; use asap_planner::ControllerConfig; use clap::Parser; +use rqe_optimizer::milp::Objective; #[derive(Parser, Debug)] #[command( @@ -27,26 +29,53 @@ struct Args { #[arg(long = "data-ingestion-interval-ms", value_parser = clap::value_parser!(u64).range(1..))] data_ingestion_interval_ms: u64, - /// YAML label-set facts: `series_count` per (metric, spatial filter) and - /// `cardinality` per (metric, spatial filter, grouping labels). - #[arg(long = "label-set-facts")] - label_set_facts: PathBuf, - - /// Path to the versioned atomic-cost document sketch-bench's `atomic-costs` - /// subcommand exports. Requires --atomic-cost-workload to select exactly - /// one measured workload profile. Omitted: every - /// benchmarked-family candidate (CMS/HLL/KLL) is dropped, since there is - /// no data to cost it at — only trivial accumulators and EXACT remain - /// selectable. - #[arg(long = "atomic-costs")] + /// Greedy only. YAML label-set facts: `series_count` per (metric, spatial + /// filter) and `cardinality` per (metric, spatial filter, grouping labels). + #[arg(long = "label-set-facts", required_unless_present = "milp")] + label_set_facts: Option, + + /// Greedy: the versioned atomic-cost document sketch-bench's + /// `atomic-costs` subcommand exports; requires --atomic-cost-workload. + /// Omitted: every benchmarked-family candidate (CMS/HLL/KLL) is dropped, + /// leaving only trivial accumulators and EXACT. + /// MILP: the flat cost table `export_rqe_optimizer_costs.sh` writes + /// (`rqe_atomic_costs.json`); required. + #[arg(long = "atomic-costs", required_if_eq("milp", "true"))] atomic_costs: Option, - /// JSON `profiles[].workload` value copied from the sketch-bench atomic-cost - /// document. This makes the empirical workload profile explicit and avoids - /// mixing costs from different traces or time windows. - #[arg(long = "atomic-cost-workload", requires = "atomic_costs")] + /// Greedy only. JSON `profiles[].workload` value copied from the + /// sketch-bench atomic-cost document. This makes the empirical workload + /// profile explicit and avoids mixing costs from different traces or time + /// windows. + #[arg( + long = "atomic-cost-workload", + requires = "atomic_costs", + conflicts_with = "milp" + )] atomic_cost_workload: Option, + /// Plan with sketch-bench's rqe-optimizer MILP and print the plan; writes + /// no configs yet. + #[arg(long)] + milp: bool, + + /// MILP only. YAML workload facts: per metric, `cardinality` per label + /// set, including the set of all its labels (the series count). + #[arg( + long = "workload-facts", + required_if_eq("milp", "true"), + requires = "milp" + )] + workload_facts: Option, + + /// MILP only. Objective weight on CPU-sec/sec. Default: rqe-optimizer's. + #[arg(long = "w-cpu", requires = "milp")] + w_cpu: Option, + + /// MILP only. Objective weight on memory GiB. Default: rqe-optimizer's. + #[arg(long = "w-mem", requires = "milp")] + w_mem: Option, + #[arg(short, long, action = clap::ArgAction::Count)] verbose: u8, } @@ -64,7 +93,14 @@ fn main() -> anyhow::Result<()> { let yaml_str = std::fs::read_to_string(&args.input_config)?; let config: ControllerConfig = serde_yaml::from_str(&yaml_str)?; - let facts = LabelSetFacts::from_path(&args.label_set_facts)?; + if args.milp { + return run_milp(&args, &config); + } + let facts = LabelSetFacts::from_path( + args.label_set_facts + .as_deref() + .expect("clap requires --label-set-facts without --milp"), + )?; let atomic_cost_table = match load_optional_selected_atomic_cost_table( args.atomic_costs.as_deref(), @@ -108,3 +144,76 @@ fn main() -> anyhow::Result<()> { Ok(()) } + +fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { + let hints = config.metrics.as_deref().unwrap_or_default(); + let facts = load_workload_facts( + args.workload_facts + .as_deref() + .expect("clap requires --workload-facts with --milp"), + hints, + args.data_ingestion_interval_ms, + )?; + let costs_path = args + .atomic_costs + .as_deref() + .expect("clap requires --atomic-costs with --milp"); + let costs: AtomicCostTable = serde_json::from_str(&std::fs::read_to_string(costs_path)?) + .map_err(|e| anyhow::anyhow!("parsing cost table {}: {e}", costs_path.display()))?; + let Objective::AUCCost { w_cpu, w_mem } = Objective::default(); + let objective = Objective::AUCCost { + w_cpu: args.w_cpu.unwrap_or(w_cpu), + w_mem: args.w_mem.unwrap_or(w_mem), + }; + 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]; + println!( + " [{d}] {:?} {} config={} metric={} grouping={:?} window={}ms slide={}ms", + dep.capability, + dep.config.sketch, + dep.config.sketch_config, + dep.metric, + dep.grouping_labels, + dep.window_ms, + dep.slide_ms, + ); + } + println!("\n=== Raqes: {} ===", workload.raqes.len()); + for ((raqe, &d), latency_ms) in workload + .raqes + .iter() + .zip(&solution.mapping) + .zip(&solution.plan_cost.query_latency_ms) + { + println!(" {} -> [{d}] latency={latency_ms:.3e}ms", raqe.id); + } + let cost = &solution.plan_cost; + println!( + "\nobjective={:.6e} cpu={:.6e} cpu-sec/sec memory={:.3} MB", + objective.value(cost), + cost.cpu_secs_per_sec(), + cost.memory_bytes() / 1e6, + ); + for (phase, c) in [ + ("ingest", &cost.ingest), + ("merge", &cost.merge), + ("query", &cost.query), + ("storage", &cost.storage), + ] { + println!( + " {phase}: cpu={:.6e} cpu-sec/sec memory={:.3} MB", + c.cpu_secs_per_sec, + c.memory_bytes / 1e6 + ); + } + Ok(()) +} diff --git a/asap-planner-rs/src/optimizer/atomic_costs.rs b/asap-planner-rs/src/optimizer/atomic_costs.rs index b9d03cfa..db769b26 100644 --- a/asap-planner-rs/src/optimizer/atomic_costs.rs +++ b/asap-planner-rs/src/optimizer/atomic_costs.rs @@ -1,14 +1,8 @@ //! The atomic-cost table exported by sketch-bench (sketch-bench#30, //! `scripts/export_atomic_costs.sh`), and the (sketch_type, params) lookup //! that resolves a candidate's [`AtomicCosts`] from it. -//! -//! `AtomicCostEntry`/`AtomicCostTable` are a deliberate duplicate of -//! sketch-bench's `aqpbm_core::atomic_costs` types, not a shared dependency — -//! see ASAPQuery#524 and sketch-bench#30 for why. Keep the two in sync by -//! hand; `atomic_cost_entry_deserializes_sketch_benchs_documented_shape` -//! below is a canary for drift. - -use std::collections::{BTreeMap, HashMap}; + +use std::collections::HashMap; use std::path::Path; use promql_utilities::query_logics::enums::AggregationType; @@ -73,25 +67,7 @@ pub struct ExternalWorkload { pub timestamp_unit: String, } -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct AtomicCostEntry { - pub sketch: String, - pub sketch_config: Value, - pub mem_bytes_per_instance: f64, - pub insert_cpu_secs: f64, - pub merge_cpu_secs: f64, - pub query_cpu_secs: f64, - pub query_accuracy: BTreeMap, - /// The conditions sketch-bench measured the row under (items, keys and - /// value range per instance, merge operand size, distribution; - /// sketch-bench#147). Carried opaquely, like a synthetic workload - /// description; absent in tables written before it existed. - #[serde(default, skip_serializing_if = "Option::is_none")] - pub measured_at: Option, -} - -pub type AtomicCostTable = Vec; +pub use rqe_optimizer::{AtomicCostEntry, AtomicCostTable}; /// Parse a standalone JSON workload selector. The selector is the exact /// `profiles[].workload` value copied from the benchmark artifact, making the @@ -472,6 +448,8 @@ fn valid_cost_entry(entry: &AtomicCostEntry) -> bool { #[cfg(test)] mod tests { + use std::collections::BTreeMap; + use super::*; #[test] @@ -578,8 +556,7 @@ mod tests { assert!(serde_json::from_str::(json).is_err()); } - /// sketch-bench#147 adds `measured_at`; it loads, and older rows without - /// it still do. Other unknown fields are still refused. + /// Rows with `measured_at` load, and older rows without it still do. #[test] fn atomic_cost_entry_accepts_optional_measured_at() { let base = r#""sketch":"kll-percall","sketch_config":null,"mem_bytes_per_instance":1.0,"insert_cpu_secs":1.0,"merge_cpu_secs":1.0,"query_cpu_secs":1.0,"query_accuracy":{}"#; @@ -587,12 +564,9 @@ mod tests { r#"{{{base},"measured_at":{{"items_per_instance":1000000,"keys_per_instance":100000,"value_range":[1.0,100000.0],"merge_operand_items":62500,"distribution":{{"kind":"zipf","skewness":1.1,"population_size":100000,"seed":42}}}}}}"# ); let entry: AtomicCostEntry = serde_json::from_str(&with).unwrap(); - assert_eq!(entry.measured_at.unwrap()["items_per_instance"], 1_000_000); + assert_eq!(entry.measured_at.unwrap().items_per_instance, 1_000_000); let without: AtomicCostEntry = serde_json::from_str(&format!("{{{base}}}")).unwrap(); assert!(without.measured_at.is_none()); - assert!( - serde_json::from_str::(&format!(r#"{{{base},"other":1}}"#)).is_err() - ); } #[test] @@ -617,6 +591,7 @@ mod tests { merge_cpu_secs: 4.5e-4, query_cpu_secs: 7.8e-8, query_accuracy: BTreeMap::new(), + merge_accuracy: BTreeMap::new(), measured_at: None, } } @@ -640,6 +615,7 @@ mod tests { merge_cpu_secs: 4.0, query_cpu_secs: 8.0, query_accuracy: BTreeMap::new(), + merge_accuracy: BTreeMap::new(), measured_at: None, } } @@ -652,15 +628,6 @@ mod tests { ]) } - #[test] - fn atomic_cost_entry_deserializes_sketch_benchs_documented_shape() { - // Pinned against the current sketch-bench atomic-cost entry shape. - let json = r#"{"sketch":"cms-fastpath-vector2d","sketch_config":{"algorithm":"cms-fastpath-vector2d","params":{"cols":1024,"rows":3}},"mem_bytes_per_instance":12288.0,"insert_cpu_secs":8.484689139741214e-9,"merge_cpu_secs":0.00045364040539336466,"query_cpu_secs":7.799774697708031e-8,"query_accuracy":{"relative_error":0.01}}"#; - let entry: AtomicCostEntry = serde_json::from_str(json).expect("documented shape parses"); - assert_eq!(entry.sketch, "cms-fastpath-vector2d"); - assert_eq!(entry.mem_bytes_per_instance, 12288.0); - } - #[test] fn cms_candidate_resolves_by_exact_key_regardless_of_value_type() { // ASAPQuery's grid stores depth/width as u64; sketch-bench's exported @@ -817,6 +784,7 @@ mod tests { merge_cpu_secs: 2.76e-4, query_cpu_secs: 1.23e-4, query_accuracy: BTreeMap::new(), + merge_accuracy: BTreeMap::new(), measured_at: None, }]; let hll_params = HashMap::from([("precision".to_string(), Value::from(14u64))]); @@ -830,6 +798,7 @@ mod tests { merge_cpu_secs: 1.0e-3, query_cpu_secs: 1.6e-4, query_accuracy: BTreeMap::new(), + merge_accuracy: BTreeMap::new(), measured_at: None, }]; let kll_params = HashMap::from([("K".to_string(), Value::from(200u64))]); diff --git a/asap-planner-rs/src/optimizer/greedy.rs b/asap-planner-rs/src/optimizer/greedy.rs index e23d6b8e..45858b0f 100644 --- a/asap-planner-rs/src/optimizer/greedy.rs +++ b/asap-planner-rs/src/optimizer/greedy.rs @@ -237,6 +237,7 @@ mod tests { merge_cpu_secs: 0.0, query_cpu_secs: 0.0, query_accuracy: std::collections::BTreeMap::new(), + merge_accuracy: std::collections::BTreeMap::new(), measured_at: None, }]; let aqe = make_aqe(Statistic::Topk, 60_000, 60_000, 1.0 / 60.0); @@ -276,6 +277,7 @@ mod tests { merge_cpu_secs: 0.0, query_cpu_secs: 0.0, query_accuracy: std::collections::BTreeMap::new(), + merge_accuracy: std::collections::BTreeMap::new(), measured_at: None, }]; let mut aqe = make_aqe(Statistic::Sum, 60_000, 60_000, 1.0 / 60.0); diff --git a/asap-planner-rs/src/optimizer/milp.rs b/asap-planner-rs/src/optimizer/milp.rs new file mode 100644 index 00000000..fc60e17c --- /dev/null +++ b/asap-planner-rs/src/optimizer/milp.rs @@ -0,0 +1,431 @@ +//! Plans a workload with sketch-bench's rqe-optimizer MILP: converts the +//! workload config into `Raqe`s and solves for the cheapest deployments. + +use std::cmp::Ordering; + +use promql_utilities::query_logics::enums::Statistic; +use rqe_optimizer::candidates::{build_all_candidates, eligible_deployments_for}; +use rqe_optimizer::enumerate::unservable; +use rqe_optimizer::milp::{minimize, MilpSolution, Objective}; +use rqe_optimizer::{ + validate_facts, AccuracyDirection, AtomicCostEntry, Capability, Deployment, LabelSet, Raqe, + WorkloadFacts, +}; +use thiserror::Error; + +use crate::config::input::ControllerConfig; + +use super::pipeline::{extract_hinted_items, OptimizerPipelineError}; +use super::solution::OptimizerItem; + +/// `query_accuracy` keys in sketch-bench's cost export. Each is the worst case +/// its comparator reports. +const RELATIVE_ERROR: &str = "relative_error"; +const MAX_RANK_ERROR: &str = "max_rank_err"; +const PRECISION_AT_K: &str = "precision_at_k"; + +#[derive(Debug, Error)] +pub enum MilpError { + #[error(transparent)] + Pipeline(#[from] OptimizerPipelineError), + #[error("query {query:?}: accuracy_sla {accuracy_sla} must be between 0 and 1")] + AccuracySlaOutOfRange { query: String, accuracy_sla: f64 }, + #[error("invalid MILP inputs:\n{}", .0.join("\n"))] + InvalidInputs(Vec), + #[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. +pub fn solve_milp( + workload: &MilpWorkload, + facts: &WorkloadFacts, + costs: &[AtomicCostEntry], + objective: Objective, +) -> Result<(Vec, MilpSolution), MilpError> { + let deployments = build_all_candidates(&workload.raqes, costs, facts, false); + tracing::debug!( + candidates = deployments.len(), + cost_rows = costs.len(), + "milp: built candidates" + ); + for (raqe, eligible) in workload + .raqes + .iter() + .map(|r| (r, eligible_deployments_for(r, &deployments))) + { + tracing::debug!(raqe = %raqe.id, eligible = eligible.len(), "milp: eligible deployments"); + } + let missing = unservable(&workload.raqes, &deployments); + if !missing.is_empty() { + return Err(MilpError::Unservable(missing)); + } + let solution = minimize(&workload.raqes, &deployments, facts, objective) + .map_err(|e| MilpError::Solver(e.to_string()))?; + tracing::debug!( + 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, + "milp: solved" + ); + Ok((deployments, solution)) +} + +/// The MILP's view of a workload. +#[derive(Debug)] +pub struct MilpWorkload { + pub items: Vec, + pub raqes: Vec, + /// `raqes[i]` is an occurrence of `items[raqe_items[i]]`. + pub raqe_items: Vec, +} + +/// One Raqe per query occurrence: a `Raqe` has no frequency weight, so two +/// dashboards running the same query at the same cadence become two +/// identical Raqes, which the MILP serves from one shared deployment. +pub fn build_milp_workload( + config: &ControllerConfig, + facts: &WorkloadFacts, + scrape_interval_ms: u64, +) -> Result { + let mut items = extract_hinted_items(config, scrape_interval_ms)?; + // Stable Raqe order and ids across runs. + items.sort_by(|a, b| { + (&a.query_strings, a.t_repeat_ms) + .cmp(&(&b.query_strings, b.t_repeat_ms)) + .then(a.accuracy_sla.total_cmp(&b.accuracy_sla)) + .then( + a.latency_sla_ms + .partial_cmp(&b.latency_sla_ms) + .unwrap_or(Ordering::Equal), + ) + }); + + let mut raqes = Vec::new(); + let mut raqe_items = Vec::new(); + for (index, item) in items.iter().enumerate() { + let raqe = item_to_raqe(item)?; + let count = occurrence_count(item); + tracing::debug!( + item = index, + queries = ?item.query_strings, + statistics = ?item.requirements.statistics, + occurrences = count, + capability = ?raqe.capability, + metric = %raqe.metric, + spatial_filter = %raqe.spatial_filter, + grouping_labels = ?raqe.grouping_labels, + lookback_ms = raqe.lookback_ms, + interval_ms = raqe.interval_ms, + accuracy_metric = %raqe.accuracy_metric, + accuracy_sla = raqe.accuracy_sla, + accuracy_direction = ?raqe.accuracy_direction, + latency_sla_ms = ?raqe.latency_sla_ms, + "milp inputs: item -> raqe" + ); + for k in 0..count { + raqes.push(Raqe { + id: format!("{}#{k}", raqe.id), + ..raqe.clone() + }); + raqe_items.push(index); + } + } + + tracing::debug!( + items = items.len(), + raqes = raqes.len(), + metrics_with_facts = facts.len(), + "milp inputs: built raqes" + ); + validate_facts(&raqes, facts).map_err(MilpError::InvalidInputs)?; + Ok(MilpWorkload { + items, + raqes, + raqe_items, + }) +} + +/// `query_frequency_hz` is `count * 1000 / t_repeat_ms`. +fn occurrence_count(item: &OptimizerItem) -> usize { + (item.query_frequency_hz * item.t_repeat_ms as f64 / 1000.0).round() as usize +} + +fn item_to_raqe(item: &OptimizerItem) -> Result { + let req = &item.requirements; + let query = item.query_strings.join(" | "); + let [statistic] = req.statistics.as_slice() else { + panic!("optimizer item {query:?} must have exactly one statistic after the avg rewrite"); + }; + let capability = capability(*statistic); + if !(0.0..=1.0).contains(&item.accuracy_sla) { + return Err(MilpError::AccuracySlaOutOfRange { + query, + accuracy_sla: item.accuracy_sla, + }); + } + let (accuracy_metric, accuracy_sla, accuracy_direction) = + accuracy_target(capability, item.accuracy_sla); + + // 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 + .topk_by_labels + .as_ref() + .map(|labels| labels.labels.iter().cloned().collect()) + .unwrap_or_default(), + _ => req.grouping_labels.labels.iter().cloned().collect(), + }; + + Ok(Raqe { + id: query, + capability, + lookback_ms: req.data_range_ms, + interval_ms: item.t_repeat_ms, + metric: req.metric.clone(), + spatial_filter: req.spatial_filter_normalized.clone(), + grouping_labels, + accuracy_metric: accuracy_metric.to_string(), + accuracy_sla, + accuracy_direction, + latency_sla_ms: item.latency_sla_ms, + }) +} + +fn capability(statistic: Statistic) -> Capability { + match statistic { + Statistic::Sum | Statistic::Count => Capability::SumOrCount, + 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, + } +} + +/// `accuracy_sla` is required accuracy: error metrics must stay within +/// `1 - accuracy_sla`, and top-k precision must reach `accuracy_sla`. +fn accuracy_target( + capability: Capability, + accuracy_sla: f64, +) -> (&'static str, f64, AccuracyDirection) { + let max_error = 1.0 - accuracy_sla; + match capability { + Capability::SumOrCount + | 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 => ( + PRECISION_AT_K, + accuracy_sla, + AccuracyDirection::HigherIsBetter, + ), + } +} + +#[cfg(test)] +mod tests { + use std::collections::BTreeMap; + + use super::*; + use crate::optimizer::workload_facts::parse_workload_facts; + + const SCRAPE_MS: u64 = 15_000; + + const FACTS: &str = r#" +metrics: + - metric: http_requests_total + groups: + - labels: [instance, job] + cardinality: 100 + - labels: [job] + cardinality: 10 + - labels: [] + cardinality: 1 +"#; + + 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 { + format!( + " - queries: [\"{query}\"]\n repetition_delay_ms: 60000\n controller_options: {{accuracy_sla: {accuracy_sla}}}\n" + ) + } + + fn facts(config: &ControllerConfig) -> WorkloadFacts { + parse_workload_facts(FACTS, config.metrics.as_deref().unwrap(), SCRAPE_MS).unwrap() + } + + fn workload(groups: &str) -> MilpWorkload { + let config = config(groups); + build_milp_workload(&config, &facts(&config), SCRAPE_MS).unwrap() + } + + fn labels(names: &[&str]) -> LabelSet { + names.iter().map(|s| s.to_string()).collect() + } + + fn cost(sketch: &str, accuracy: &[(&str, f64)]) -> AtomicCostEntry { + AtomicCostEntry { + sketch: sketch.into(), + sketch_config: serde_json::json!({"algorithm": sketch, "params": {}}), + mem_bytes_per_instance: 100.0, + insert_cpu_secs: 1e-7, + merge_cpu_secs: 1e-7, + query_cpu_secs: 1e-7, + query_accuracy: accuracy.iter().map(|(k, v)| (k.to_string(), *v)).collect(), + merge_accuracy: BTreeMap::new(), + measured_at: None, + } + } + + fn costs() -> Vec { + vec![ + cost("exact-sum", &[(RELATIVE_ERROR, 0.0)]), + cost("kll-percall", &[(MAX_RANK_ERROR, 0.005)]), + cost( + "cms-heap-topk-fastpath-vector2d", + &[(PRECISION_AT_K, 0.995)], + ), + ] + } + + #[test] + fn repeated_query_becomes_one_raqe_per_occurrence() { + let query = "sum by (job) (http_requests_total)"; + let w = workload(&(group(query, 0.99) + &group(query, 0.99))); + assert_eq!(w.items.len(), 1); + assert_eq!(w.raqe_items, vec![0, 0]); + let ids: Vec<_> = w.raqes.iter().map(|r| r.id.as_str()).collect(); + assert_eq!(ids, vec![format!("{query}#0"), format!("{query}#1")]); + } + + #[test] + fn sum_maps_to_sum_or_count_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.grouping_labels, labels(&["job"])); + assert_eq!(r.accuracy_metric, RELATIVE_ERROR); + assert_eq!(r.accuracy_direction, AccuracyDirection::LowerIsBetter); + assert!((r.accuracy_sla - 0.01).abs() < 1e-12); + assert_eq!(r.interval_ms, 60_000); + assert_eq!(r.latency_sla_ms, None); + } + + #[test] + fn quantile_uses_max_rank_error() { + let w = workload(&group( + "quantile_over_time(0.99, http_requests_total[5m])", + 0.99, + )); + let r = &w.raqes[0]; + assert_eq!(r.capability, Capability::Quantile); + assert_eq!(r.accuracy_metric, MAX_RANK_ERROR); + assert_eq!(r.lookback_ms, 300_000); + } + + #[test] + 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.accuracy_metric, PRECISION_AT_K); + assert_eq!(r.accuracy_direction, AccuracyDirection::HigherIsBetter); + assert_eq!(r.accuracy_sla, 0.9); + // Not the all-labels output set, which would cost one heap per series. + assert_eq!(r.grouping_labels, labels(&["job"])); + } + + #[test] + fn bare_topk_is_one_global_heap() { + let w = workload(&group("topk(5, http_requests_total)", 0.9)); + assert_eq!(w.raqes[0].grouping_labels, LabelSet::new()); + } + + #[test] + fn accuracy_sla_above_one_is_rejected() { + let config = config(&group("sum by (job) (http_requests_total)", 1.5)); + let err = build_milp_workload(&config, &facts(&config), SCRAPE_MS).unwrap_err(); + assert!(matches!(err, MilpError::AccuracySlaOutOfRange { .. })); + } + + #[test] + fn spatial_filter_is_rejected() { + let config = config(&group( + "sum by (job) (http_requests_total{job='api'})", + 0.99, + )); + let err = build_milp_workload(&config, &facts(&config), SCRAPE_MS).unwrap_err(); + let MilpError::InvalidInputs(problems) = err else { + panic!("expected InvalidInputs, got {err:?}"); + }; + assert!(problems.iter().any(|p| p.contains("spatial filter"))); + } + + #[test] + fn missing_cardinality_is_rejected() { + let config = config(&group("sum by (instance) (http_requests_total)", 0.99)); + let err = build_milp_workload(&config, &facts(&config), SCRAPE_MS).unwrap_err(); + assert!(matches!(err, MilpError::InvalidInputs(_))); + } + + #[test] + fn solve_shares_one_deployment_across_repeated_queries() { + let query = "sum by (job) (http_requests_total)"; + 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"); + } + + #[test] + fn solve_rejects_accuracy_no_sketch_meets() { + // kll's max rank error 0.005 misses a 0.999 requirement (0.001). + let config = config(&group( + "quantile_over_time(0.99, http_requests_total[5m])", + 0.999, + )); + 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(); + assert!(matches!(err, MilpError::Unservable(ids) if ids.len() == 1)); + } + + #[test] + fn solve_picks_a_deployable_family_per_capability() { + let groups = group("sum by (job) (http_requests_total)", 0.99) + + &group("quantile_over_time(0.99, http_requests_total[5m])", 0.99) + + &group("topk(5, http_requests_total)", 0.99); + 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 chosen: BTreeMap = w + .raqes + .iter() + .zip(&solution.mapping) + .map(|(r, &d)| (r.capability, deployments[d].config.sketch.as_str())) + .collect(); + assert_eq!(chosen[&Capability::SumOrCount], "exact-sum"); + assert_eq!(chosen[&Capability::Quantile], "kll-percall"); + assert_eq!(chosen[&Capability::TopK], "cms-heap-topk-fastpath-vector2d"); + } +} diff --git a/asap-planner-rs/src/optimizer/mod.rs b/asap-planner-rs/src/optimizer/mod.rs index 91326c71..9175bf68 100644 --- a/asap-planner-rs/src/optimizer/mod.rs +++ b/asap-planner-rs/src/optimizer/mod.rs @@ -6,10 +6,12 @@ pub mod cost_model; pub mod error; pub mod greedy; pub mod label_set_facts; +pub mod milp; pub mod pipeline; pub mod sketch_properties; pub mod solution; pub mod translator; +pub mod workload_facts; pub use aqe_extractor::{extract_aqes, RQE}; pub use atomic_costs::{ @@ -22,7 +24,9 @@ pub use cost_model::{ingest_cost, query_cost, total_cost_rate, AtomicCosts, Cost 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 pipeline::{run_greedy_pipeline, OptimizerPipelineError}; pub use sketch_properties::{sketch_properties, SketchProperties}; pub use solution::{AQEAssignment, OptimizerItem, OptimizerSolution, QueryMethod}; pub use translator::translate; +pub use workload_facts::{load_workload_facts, parse_workload_facts, WorkloadFactsError}; diff --git a/asap-planner-rs/src/optimizer/pipeline.rs b/asap-planner-rs/src/optimizer/pipeline.rs index 8a817d68..c54aaa93 100644 --- a/asap-planner-rs/src/optimizer/pipeline.rs +++ b/asap-planner-rs/src/optimizer/pipeline.rs @@ -9,7 +9,7 @@ use super::atomic_costs::AtomicCostTable; use super::cost_model::CostWeights; use super::greedy::greedy_assign; use super::label_set_facts::{LabelSetFacts, LabelSetFactsError, LabelSetKey}; -use super::solution::OptimizerSolution; +use super::solution::{OptimizerItem, OptimizerSolution}; use super::translator::{translate, TranslationSummary}; #[derive(Debug, Error)] @@ -51,25 +51,7 @@ pub fn run_greedy_pipeline( scrape_interval_ms: u64, atomic_cost_table: &AtomicCostTable, ) -> Result<(StreamingConfig, InferenceConfig), OptimizerPipelineError> { - if config.metrics.is_none() { - return Err(LabelSetFactsError::MissingMetricHints.into()); - } - let schema = config.schema_from_hints(); - let rqes = config_to_rqes(config); - let aqes = extract_aqes(&rqes, &schema, scrape_interval_ms)?; - - // Requirement extraction treats an unknown metric as having no labels, - // which would silently mis-resolve `without (...)` and plain selectors. - let mut unhinted: Vec = aqes - .iter() - .map(|aqe| aqe.requirements.metric.clone()) - .filter(|metric| schema.get_labels(metric).is_none()) - .collect(); - if !unhinted.is_empty() { - unhinted.sort(); - unhinted.dedup(); - return Err(LabelSetFactsError::MetricsWithoutHints(unhinted).into()); - } + let aqes = extract_hinted_items(config, scrape_interval_ms)?; let item_facts = facts.resolve(&aqes, scrape_interval_ms)?; for (aqe, item) in aqes.iter().zip(&item_facts) { @@ -93,6 +75,34 @@ pub fn run_greedy_pipeline( Ok(finish_pipeline(solution, "greedy")) } +/// The workload's optimizer items, with labels resolved from its `metrics:` +/// hints. Every workload metric must have a hint. +pub(super) fn extract_hinted_items( + config: &ControllerConfig, + scrape_interval_ms: u64, +) -> Result, OptimizerPipelineError> { + if config.metrics.is_none() { + return Err(LabelSetFactsError::MissingMetricHints.into()); + } + let schema = config.schema_from_hints(); + let rqes = config_to_rqes(config); + let aqes = extract_aqes(&rqes, &schema, scrape_interval_ms)?; + + // Requirement extraction treats an unknown metric as having no labels, + // which would silently mis-resolve `without (...)` and plain selectors. + let mut unhinted: Vec = aqes + .iter() + .map(|aqe| aqe.requirements.metric.clone()) + .filter(|metric| schema.get_labels(metric).is_none()) + .collect(); + if !unhinted.is_empty() { + unhinted.sort(); + unhinted.dedup(); + return Err(LabelSetFactsError::MetricsWithoutHints(unhinted).into()); + } + Ok(aqes) +} + /// Convert a `ControllerConfig`'s query groups into a flat list of RQEs. /// Each (query, repetition_delay_ms) pair becomes one RQE. fn config_to_rqes(config: &ControllerConfig) -> Vec { diff --git a/asap-planner-rs/src/optimizer/workload_facts.rs b/asap-planner-rs/src/optimizer/workload_facts.rs new file mode 100644 index 00000000..a5569926 --- /dev/null +++ b/asap-planner-rs/src/optimizer/workload_facts.rs @@ -0,0 +1,188 @@ +//! Externally provided workload facts for the MILP planner: per metric, the +//! cardinality of each label set in use. The entry for all of a metric's +//! labels is its series count. Labels come from the workload's `metrics:` +//! hints and the scrape interval from the caller. + +use std::collections::BTreeMap; +use std::path::{Path, PathBuf}; + +use rqe_optimizer::{LabelSet, MetricFacts, Millis, WorkloadFacts}; +use serde::Deserialize; +use thiserror::Error; + +use crate::config::input::MetricDefinition; + +#[derive(Debug, Error)] +pub enum WorkloadFactsError { + #[error("failed to read workload facts '{path}': {source}")] + Read { + path: PathBuf, + source: std::io::Error, + }, + #[error("failed to parse workload facts: {0}")] + Parse(#[from] serde_yaml::Error), + #[error("duplicate facts for metric {0:?}")] + DuplicateMetric(String), + #[error("metric {metric:?}: duplicate cardinality for labels {labels:?}")] + DuplicateLabels { metric: String, labels: LabelSet }, + #[error("metric {0:?} has facts but no `metrics:` hint giving its labels")] + MetricWithoutHint(String), +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct FactsFile { + metrics: Vec, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct MetricEntry { + metric: String, + groups: Vec, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct GroupEntry { + labels: Vec, + cardinality: u64, +} + +pub fn load_workload_facts( + path: &Path, + hints: &[MetricDefinition], + scrape_interval_ms: Millis, +) -> Result { + let yaml = std::fs::read_to_string(path).map_err(|source| WorkloadFactsError::Read { + path: path.to_path_buf(), + source, + })?; + parse_workload_facts(&yaml, hints, scrape_interval_ms) +} + +/// Missing or zero cardinalities are left to `rqe_optimizer::validate_facts`, +/// which checks exactly the label sets the workload uses. +pub fn parse_workload_facts( + yaml: &str, + hints: &[MetricDefinition], + scrape_interval_ms: Millis, +) -> Result { + let file: FactsFile = serde_yaml::from_str(yaml)?; + let mut facts = WorkloadFacts::new(); + for entry in file.metrics { + let Some(hint) = hints.iter().find(|h| h.metric == entry.metric) else { + return Err(WorkloadFactsError::MetricWithoutHint(entry.metric)); + }; + let mut cardinality = BTreeMap::new(); + for group in entry.groups { + let labels: LabelSet = group.labels.into_iter().collect(); + if cardinality + .insert(labels.clone(), group.cardinality) + .is_some() + { + return Err(WorkloadFactsError::DuplicateLabels { + metric: entry.metric, + labels, + }); + } + } + tracing::debug!( + metric = %entry.metric, + labels = ?hint.labels, + scrape_interval_ms, + cardinality = ?cardinality, + "workload facts: metric" + ); + let metric_facts = MetricFacts { + labels: hint.labels.iter().cloned().collect(), + scrape_interval_ms, + cardinality, + // ponytail: only sizes DDSketch, which ASAPQuery can't deploy yet. + value_range: None, + }; + if facts.insert(entry.metric.clone(), metric_facts).is_some() { + return Err(WorkloadFactsError::DuplicateMetric(entry.metric)); + } + } + Ok(facts) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn hints() -> Vec { + vec![MetricDefinition { + metric: "http_requests_total".into(), + labels: vec!["job".into(), "instance".into()], + }] + } + + fn labels(names: &[&str]) -> LabelSet { + names.iter().map(|s| s.to_string()).collect() + } + + #[test] + fn parses_cardinalities_and_takes_labels_from_hints() { + let yaml = r#" +metrics: + - metric: http_requests_total + groups: + - labels: [instance, job] + cardinality: 1200 + - labels: [job] + cardinality: 10 +"#; + let facts = parse_workload_facts(yaml, &hints(), 15_000).unwrap(); + let m = &facts["http_requests_total"]; + assert_eq!(m.labels, labels(&["job", "instance"])); + assert_eq!(m.scrape_interval_ms, 15_000); + assert_eq!(m.cardinality[&labels(&["job", "instance"])], 1200); + assert_eq!(m.cardinality[&labels(&["job"])], 10); + } + + #[test] + fn rejects_metric_without_hint() { + let yaml = "metrics:\n - metric: other\n groups: []\n"; + let err = parse_workload_facts(yaml, &hints(), 15_000).unwrap_err(); + assert!(matches!(err, WorkloadFactsError::MetricWithoutHint(m) if m == "other")); + } + + #[test] + fn rejects_duplicate_metric() { + let yaml = r#" +metrics: + - metric: http_requests_total + groups: [] + - metric: http_requests_total + groups: [] +"#; + let err = parse_workload_facts(yaml, &hints(), 15_000).unwrap_err(); + assert!(matches!(err, WorkloadFactsError::DuplicateMetric(_))); + } + + #[test] + fn rejects_duplicate_label_set_in_any_order() { + let yaml = r#" +metrics: + - metric: http_requests_total + groups: + - labels: [job, instance] + cardinality: 1200 + - labels: [instance, job] + cardinality: 1100 +"#; + let err = parse_workload_facts(yaml, &hints(), 15_000).unwrap_err(); + assert!(matches!(err, WorkloadFactsError::DuplicateLabels { .. })); + } + + #[test] + fn rejects_old_series_and_groups_format() { + let yaml = "series: []\ngroups: []\n"; + assert!(matches!( + parse_workload_facts(yaml, &hints(), 15_000).unwrap_err(), + WorkloadFactsError::Parse(_) + )); + } +} diff --git a/asap-planner-rs/tests/rqe_optimizer_smoke.rs b/asap-planner-rs/tests/rqe_optimizer_smoke.rs deleted file mode 100644 index 03858aa5..00000000 --- a/asap-planner-rs/tests/rqe_optimizer_smoke.rs +++ /dev/null @@ -1,66 +0,0 @@ -//! Links sketch-bench's rqe-optimizer and solves a one-query MILP with HiGHS. - -use std::collections::BTreeMap; - -use rqe_optimizer::milp::{minimize_cost, MilpBounds}; -use rqe_optimizer::objectives::MachineFamily; -use rqe_optimizer::{ - AccuracyDirection, AtomicCostEntry, Capability, Deployment, LabelSet, LabelSetInfo, Rqe, -}; - -fn deployment(insert_cpu_secs: f64) -> Deployment { - Deployment { - capability: Capability::Quantile, - labels: LabelSet::new(), - config: AtomicCostEntry { - sketch: "kll-percall".into(), - sketch_config: serde_json::json!(null), - mem_bytes_per_instance: 1024.0, - insert_cpu_secs, - merge_cpu_secs: 1.0, - query_cpu_secs: 1.0, - query_accuracy: BTreeMap::from([("err".into(), 0.0)]), - }, - window_secs: 60, - slide_secs: 60, - } -} - -#[test] -fn minimize_cost_picks_the_cheaper_deployment() { - let rqes = vec![Rqe { - id: "q".into(), - capability: Capability::Quantile, - lookback_secs: 60, - interval_secs: 60, - labels: LabelSet::new(), - accuracy_metric: "err".into(), - accuracy_tolerance: 1.0, - accuracy_direction: AccuracyDirection::LowerIsBetter, - }]; - let deployments = vec![deployment(2.0), deployment(1.0)]; - let label_sets = BTreeMap::from([( - LabelSet::new(), - LabelSetInfo { - cardinality: 1, - arrival_rate_per_sec: 1.0, - }, - )]); - let family = MachineFamily { - family: "test".into(), - vcpu: 4.0, - memory_gib: 16.0, - usd_per_hour: 1.0, - }; - - let solution = minimize_cost( - &rqes, - &deployments, - &label_sets, - &MilpBounds::default(), - &family, - ) - .expect("feasible MILP"); - - assert_eq!(solution.mapping, vec![1]); -} From 55617814f8f324ff9badb347173debf8fc8e452f Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Tue, 6 Oct 2026 18:46:05 -0400 Subject: [PATCH 2/3] fix(planner): harden MILP inputs and CLI per review - Raqe ids carry the item index, so items differing only in T or SLA get distinct ids. - `OptimizerItem.occurrences` counts merged leaves exactly; `query_frequency_hz` is derived from it. - An item with more than one statistic is a `MilpError`, not a panic. - Eligibility is computed once per Raqe for the debug log and the unservable check; the item sort is total. - The MILP cost table loader rejects non-finite or negative rows; `--w-cpu`/`--w-mem` must be finite and >= 0. - `--label-set-facts` conflicts with `--milp`; missing `metrics:` hints report `MissingMetricHints`; the MILP path warns on default SLAs. Co-Authored-By: Claude Opus 5.5 --- asap-planner-rs/src/bin/optimizer_cli.rs | 55 ++++++++++++---- .../src/optimizer/aqe_extractor.rs | 15 +++-- asap-planner-rs/src/optimizer/atomic_costs.rs | 36 ++++++++++ .../src/optimizer/candidate_gen.rs | 1 + asap-planner-rs/src/optimizer/cost_model.rs | 1 + asap-planner-rs/src/optimizer/greedy.rs | 2 + .../src/optimizer/label_set_facts.rs | 1 + asap-planner-rs/src/optimizer/milp.rs | 66 ++++++++++++------- asap-planner-rs/src/optimizer/mod.rs | 2 +- asap-planner-rs/src/optimizer/solution.rs | 3 + 10 files changed, 140 insertions(+), 42 deletions(-) diff --git a/asap-planner-rs/src/bin/optimizer_cli.rs b/asap-planner-rs/src/bin/optimizer_cli.rs index ea4db3f1..01a60e61 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -7,8 +7,9 @@ use std::path::PathBuf; use asap_planner::optimizer::{ - build_milp_workload, load_optional_selected_atomic_cost_table, load_workload_facts, - run_greedy_pipeline, solve_milp, AtomicCostTable, LabelSetFacts, + 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, }; use asap_planner::ControllerConfig; use clap::Parser; @@ -31,7 +32,11 @@ struct Args { /// Greedy only. YAML label-set facts: `series_count` per (metric, spatial /// filter) and `cardinality` per (metric, spatial filter, grouping labels). - #[arg(long = "label-set-facts", required_unless_present = "milp")] + #[arg( + long = "label-set-facts", + required_unless_present = "milp", + conflicts_with = "milp" + )] label_set_facts: Option, /// Greedy: the versioned atomic-cost document sketch-bench's @@ -69,11 +74,11 @@ struct Args { workload_facts: Option, /// MILP only. Objective weight on CPU-sec/sec. Default: rqe-optimizer's. - #[arg(long = "w-cpu", requires = "milp")] + #[arg(long = "w-cpu", requires = "milp", value_parser = parse_weight)] w_cpu: Option, /// MILP only. Objective weight on memory GiB. Default: rqe-optimizer's. - #[arg(long = "w-mem", requires = "milp")] + #[arg(long = "w-mem", requires = "milp", value_parser = parse_weight)] w_mem: Option, #[arg(short, long, action = clap::ArgAction::Count)] @@ -145,8 +150,21 @@ fn main() -> anyhow::Result<()> { Ok(()) } +/// Objective weights must be finite and non-negative: a negative weight +/// rewards cost, and NaN poisons every coefficient. +fn parse_weight(s: &str) -> Result { + match s.parse::() { + Ok(w) if w.is_finite() && w >= 0.0 => Ok(w), + Ok(w) => Err(format!("must be finite and >= 0, got {w}")), + Err(e) => Err(e.to_string()), + } +} + fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { - let hints = config.metrics.as_deref().unwrap_or_default(); + config.warn_default_slas(); + let Some(hints) = config.metrics.as_deref() else { + return Err(LabelSetFactsError::MissingMetricHints.into()); + }; let facts = load_workload_facts( args.workload_facts .as_deref() @@ -154,12 +172,11 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { hints, args.data_ingestion_interval_ms, )?; - let costs_path = args - .atomic_costs - .as_deref() - .expect("clap requires --atomic-costs with --milp"); - let costs: AtomicCostTable = serde_json::from_str(&std::fs::read_to_string(costs_path)?) - .map_err(|e| anyhow::anyhow!("parsing cost table {}: {e}", costs_path.display()))?; + let costs = load_flat_atomic_cost_table( + args.atomic_costs + .as_deref() + .expect("clap requires --atomic-costs with --milp"), + )?; let Objective::AUCCost { w_cpu, w_mem } = Objective::default(); let objective = Objective::AUCCost { w_cpu: args.w_cpu.unwrap_or(w_cpu), @@ -217,3 +234,17 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { } Ok(()) } + +#[cfg(test)] +mod tests { + use super::parse_weight; + + #[test] + fn weights_must_be_finite_and_non_negative() { + assert_eq!(parse_weight("0.5"), Ok(0.5)); + assert_eq!(parse_weight("0"), Ok(0.0)); + for bad in ["-1", "NaN", "inf", "x"] { + assert!(parse_weight(bad).is_err(), "{bad}"); + } + } +} diff --git a/asap-planner-rs/src/optimizer/aqe_extractor.rs b/asap-planner-rs/src/optimizer/aqe_extractor.rs index ac825a38..7f1a51fd 100644 --- a/asap-planner-rs/src/optimizer/aqe_extractor.rs +++ b/asap-planner-rs/src/optimizer/aqe_extractor.rs @@ -67,7 +67,8 @@ pub fn extract_aqes( metric_schema: &PromQLSchema, scrape_interval_ms: u64, ) -> Result, OptimizerError> { - let mut acc: HashMap, f64)> = HashMap::new(); + let mut acc: HashMap, usize)> = + HashMap::new(); for rqe in rqes { if rqe.t_repeat_ms == 0 { @@ -82,13 +83,11 @@ pub fn extract_aqes( match extract_requirements(&leaf, metric_schema, scrape_interval_ms) { Ok(req) => { let key = OptimizerItemKey::from_rqe(&req, rqe); - let entry = acc.entry(key).or_insert_with(|| (req, Vec::new(), 0.0)); + let entry = acc.entry(key).or_insert_with(|| (req, Vec::new(), 0)); if !entry.1.contains(&leaf) { entry.1.push(leaf); } - // query_frequency_hz must stay in Hz (queries per real second) - // regardless of t_repeat_ms's internal unit — 1000.0 / ms, not 1.0 / ms. - entry.2 += 1000.0 / rqe.t_repeat_ms as f64; + entry.2 += 1; } Err(reason) => { return Err(OptimizerError::UnsupportedLeaf { @@ -104,10 +103,12 @@ pub fn extract_aqes( Ok(acc .into_iter() .map( - |(key, (requirements, query_strings, query_frequency_hz))| OptimizerItem { + |(key, (requirements, query_strings, occurrences))| OptimizerItem { requirements, query_strings, - query_frequency_hz, + // Hz (queries per real second): 1000.0 / ms, not 1.0 / ms. + query_frequency_hz: occurrences as f64 * 1000.0 / key.t_repeat_ms as f64, + occurrences, t_repeat_ms: key.t_repeat_ms, accuracy_sla: f64::from_bits(key.accuracy_sla_bits), latency_sla_ms: key.latency_sla_ms_bits.map(f64::from_bits), diff --git a/asap-planner-rs/src/optimizer/atomic_costs.rs b/asap-planner-rs/src/optimizer/atomic_costs.rs index db769b26..3167953f 100644 --- a/asap-planner-rs/src/optimizer/atomic_costs.rs +++ b/asap-planner-rs/src/optimizer/atomic_costs.rs @@ -435,6 +435,27 @@ fn require<'a>( }) } +/// Load the flat cost table `export_rqe_optimizer_costs.sh` writes, for the +/// MILP. Unlike the greedy loader, an invalid row is an error, not dropped. +pub fn load_flat_atomic_cost_table(path: &Path) -> anyhow::Result { + let raw = std::fs::read_to_string(path) + .map_err(|e| anyhow::anyhow!("reading cost table {}: {e}", path.display()))?; + let table: AtomicCostTable = serde_json::from_str(&raw) + .map_err(|e| anyhow::anyhow!("parsing cost table {}: {e}", path.display()))?; + let invalid: Vec = table + .iter() + .filter(|entry| !valid_cost_entry(entry)) + .map(|entry| format!("{} {}", entry.sketch, entry.sketch_config)) + .collect(); + if !invalid.is_empty() { + anyhow::bail!( + "cost table {} has non-finite or negative costs in rows: {invalid:?}", + path.display() + ); + } + Ok(table) +} + fn valid_cost_entry(entry: &AtomicCostEntry) -> bool { [ entry.mem_bytes_per_instance, @@ -556,6 +577,21 @@ mod tests { assert!(serde_json::from_str::(json).is_err()); } + #[test] + fn flat_loader_rejects_non_finite_or_negative_costs() { + let row = |insert: &str| { + format!( + r#"{{"sketch":"hll","sketch_config":null,"mem_bytes_per_instance":1.0,"insert_cpu_secs":{insert},"merge_cpu_secs":1.0,"query_cpu_secs":1.0,"query_accuracy":{{}}}}"# + ) + }; + let file = tempfile::NamedTempFile::new().unwrap(); + std::fs::write(file.path(), format!("[{}]", row("1e-7"))).unwrap(); + assert_eq!(load_flat_atomic_cost_table(file.path()).unwrap().len(), 1); + std::fs::write(file.path(), format!("[{},{}]", row("1e-7"), row("-1.0"))).unwrap(); + let err = load_flat_atomic_cost_table(file.path()).unwrap_err(); + assert!(err.to_string().contains("negative"), "{err}"); + } + /// Rows with `measured_at` load, and older rows without it still do. #[test] fn atomic_cost_entry_accepts_optional_measured_at() { diff --git a/asap-planner-rs/src/optimizer/candidate_gen.rs b/asap-planner-rs/src/optimizer/candidate_gen.rs index 5e1d45d6..4f2eb07c 100644 --- a/asap-planner-rs/src/optimizer/candidate_gen.rs +++ b/asap-planner-rs/src/optimizer/candidate_gen.rs @@ -408,6 +408,7 @@ mod tests { }, query_strings: vec!["test_query".into()], query_frequency_hz: 1.0 / 60.0, + occurrences: 1, t_repeat_ms: min_t, accuracy_sla: 0.0, latency_sla_ms: None, diff --git a/asap-planner-rs/src/optimizer/cost_model.rs b/asap-planner-rs/src/optimizer/cost_model.rs index fe9b0f15..006792b5 100644 --- a/asap-planner-rs/src/optimizer/cost_model.rs +++ b/asap-planner-rs/src/optimizer/cost_model.rs @@ -214,6 +214,7 @@ mod tests { }, query_strings: vec!["test_query".into()], query_frequency_hz: 1.0 / 60.0, + occurrences: 1, t_repeat_ms: min_t, accuracy_sla: 0.0, latency_sla_ms: None, diff --git a/asap-planner-rs/src/optimizer/greedy.rs b/asap-planner-rs/src/optimizer/greedy.rs index 45858b0f..eaceb9a4 100644 --- a/asap-planner-rs/src/optimizer/greedy.rs +++ b/asap-planner-rs/src/optimizer/greedy.rs @@ -137,6 +137,7 @@ mod tests { }, query_strings: vec!["test_query".into()], query_frequency_hz: freq_hz, + occurrences: 1, t_repeat_ms: min_t, accuracy_sla: 0.0, latency_sla_ms: None, @@ -194,6 +195,7 @@ mod tests { }, query_strings: vec!["avg_query".into()], query_frequency_hz: 1.0 / 60.0, + occurrences: 1, t_repeat_ms: 60_000, accuracy_sla: 0.0, latency_sla_ms: None, diff --git a/asap-planner-rs/src/optimizer/label_set_facts.rs b/asap-planner-rs/src/optimizer/label_set_facts.rs index 418899fd..0028c291 100644 --- a/asap-planner-rs/src/optimizer/label_set_facts.rs +++ b/asap-planner-rs/src/optimizer/label_set_facts.rs @@ -306,6 +306,7 @@ mod tests { }, query_strings: vec!["q".into()], query_frequency_hz: 1.0 / 60.0, + occurrences: 1, t_repeat_ms: 60_000, accuracy_sla: 0.0, latency_sla_ms: None, diff --git a/asap-planner-rs/src/optimizer/milp.rs b/asap-planner-rs/src/optimizer/milp.rs index fc60e17c..9a049685 100644 --- a/asap-planner-rs/src/optimizer/milp.rs +++ b/asap-planner-rs/src/optimizer/milp.rs @@ -1,11 +1,8 @@ //! Plans a workload with sketch-bench's rqe-optimizer MILP: converts the //! workload config into `Raqe`s and solves for the cheapest deployments. -use std::cmp::Ordering; - use promql_utilities::query_logics::enums::Statistic; use rqe_optimizer::candidates::{build_all_candidates, eligible_deployments_for}; -use rqe_optimizer::enumerate::unservable; use rqe_optimizer::milp::{minimize, MilpSolution, Objective}; use rqe_optimizer::{ validate_facts, AccuracyDirection, AtomicCostEntry, Capability, Deployment, LabelSet, Raqe, @@ -32,6 +29,11 @@ pub enum MilpError { AccuracySlaOutOfRange { query: String, accuracy_sla: f64 }, #[error("invalid MILP inputs:\n{}", .0.join("\n"))] InvalidInputs(Vec), + #[error("query {query:?} needs one statistic, got {statistics:?}")] + MultipleStatistics { + query: String, + statistics: Vec, + }, #[error("no eligible deployment for raqes: {0:?}")] Unservable(Vec), #[error("MILP solve failed: {0}")] @@ -52,14 +54,14 @@ pub fn solve_milp( cost_rows = costs.len(), "milp: built candidates" ); - for (raqe, eligible) in workload - .raqes - .iter() - .map(|r| (r, eligible_deployments_for(r, &deployments))) - { - tracing::debug!(raqe = %raqe.id, eligible = eligible.len(), "milp: eligible deployments"); + let mut missing = Vec::new(); + for raqe in &workload.raqes { + let eligible = eligible_deployments_for(raqe, &deployments).len(); + tracing::debug!(raqe = %raqe.id, eligible, "milp: eligible deployments"); + if eligible == 0 { + missing.push(raqe.id.clone()); + } } - let missing = unservable(&workload.raqes, &deployments); if !missing.is_empty() { return Err(MilpError::Unservable(missing)); } @@ -98,18 +100,14 @@ pub fn build_milp_workload( (&a.query_strings, a.t_repeat_ms) .cmp(&(&b.query_strings, b.t_repeat_ms)) .then(a.accuracy_sla.total_cmp(&b.accuracy_sla)) - .then( - a.latency_sla_ms - .partial_cmp(&b.latency_sla_ms) - .unwrap_or(Ordering::Equal), - ) + .then(latency_key(a).total_cmp(&latency_key(b))) }); let mut raqes = Vec::new(); let mut raqe_items = Vec::new(); for (index, item) in items.iter().enumerate() { let raqe = item_to_raqe(item)?; - let count = occurrence_count(item); + let count = item.occurrences; tracing::debug!( item = index, queries = ?item.query_strings, @@ -129,7 +127,8 @@ pub fn build_milp_workload( ); for k in 0..count { raqes.push(Raqe { - id: format!("{}#{k}", raqe.id), + // The item index keeps ids unique across T and SLAs. + id: format!("{index}:{}#{k}", raqe.id), ..raqe.clone() }); raqe_items.push(index); @@ -150,16 +149,19 @@ pub fn build_milp_workload( }) } -/// `query_frequency_hz` is `count * 1000 / t_repeat_ms`. -fn occurrence_count(item: &OptimizerItem) -> usize { - (item.query_frequency_hz * item.t_repeat_ms as f64 / 1000.0).round() as usize +/// No limit sorts last. +fn latency_key(item: &OptimizerItem) -> f64 { + item.latency_sla_ms.unwrap_or(f64::INFINITY) } fn item_to_raqe(item: &OptimizerItem) -> Result { let req = &item.requirements; let query = item.query_strings.join(" | "); let [statistic] = req.statistics.as_slice() else { - panic!("optimizer item {query:?} must have exactly one statistic after the avg rewrite"); + return Err(MilpError::MultipleStatistics { + query, + statistics: req.statistics.clone(), + }); }; let capability = capability(*statistic); if !(0.0..=1.0).contains(&item.accuracy_sla) { @@ -310,7 +312,27 @@ metrics: assert_eq!(w.items.len(), 1); assert_eq!(w.raqe_items, vec![0, 0]); let ids: Vec<_> = w.raqes.iter().map(|r| r.id.as_str()).collect(); - assert_eq!(ids, vec![format!("{query}#0"), format!("{query}#1")]); + assert_eq!(ids, vec![format!("0:{query}#0"), format!("0:{query}#1")]); + } + + #[test] + fn same_query_at_different_cadences_gets_distinct_ids() { + let query = "sum by (job) (http_requests_total)"; + let slow = group(query, 0.99).replace("60000", "120000"); + let w = workload(&(group(query, 0.99) + &slow)); + assert_eq!(w.raqes.len(), 2); + assert_ne!(w.raqes[0].id, w.raqes[1].id); + } + + #[test] + fn item_with_two_statistics_is_an_error_not_a_panic() { + let w = workload(&group("sum by (job) (http_requests_total)", 0.99)); + let mut item = w.items[0].clone(); + item.requirements.statistics.push(Statistic::Count); + assert!(matches!( + item_to_raqe(&item), + Err(MilpError::MultipleStatistics { .. }) + )); } #[test] diff --git a/asap-planner-rs/src/optimizer/mod.rs b/asap-planner-rs/src/optimizer/mod.rs index 9175bf68..786aab16 100644 --- a/asap-planner-rs/src/optimizer/mod.rs +++ b/asap-planner-rs/src/optimizer/mod.rs @@ -15,7 +15,7 @@ pub mod workload_facts; pub use aqe_extractor::{extract_aqes, RQE}; pub use atomic_costs::{ - load_atomic_cost_table, load_optional_selected_atomic_cost_table, + load_atomic_cost_table, load_flat_atomic_cost_table, load_optional_selected_atomic_cost_table, load_selected_atomic_cost_table, resolve_atomic_costs, AtomicCostEntry, AtomicCostTable, ExternalWorkload, WorkloadDescription, }; diff --git a/asap-planner-rs/src/optimizer/solution.rs b/asap-planner-rs/src/optimizer/solution.rs index 2dfa9930..52b7713c 100644 --- a/asap-planner-rs/src/optimizer/solution.rs +++ b/asap-planner-rs/src/optimizer/solution.rs @@ -20,6 +20,9 @@ pub struct OptimizerItem { /// hitting the sketch. pub query_frequency_hz: f64, + /// Leaf occurrences merged into this item, across all RQEs. + pub occurrences: usize, + /// Repetition interval for every RQE contributing to this item, in ms. pub t_repeat_ms: u64, From a7515f6c93a493df49f6718b4daa90eb52627734 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Tue, 6 Oct 2026 22:14:49 -0400 Subject: [PATCH 3/3] fix(planner): accept accuracy at the SLA boundary; reject an all-zero objective - Accuracy tolerances get a 1e-9 slack so `1 - sla` rounding doesn't reject a cost row measured exactly at the boundary. - `asap-optimizer-cli --milp` rejects `--w-cpu`/`--w-mem` both 0, which would make the solver's pick arbitrary. Co-Authored-By: Claude Opus 5.5 --- asap-planner-rs/src/bin/optimizer_cli.rs | 11 +++++++---- asap-planner-rs/src/optimizer/milp.rs | 24 ++++++++++++++++++++---- 2 files changed, 27 insertions(+), 8 deletions(-) diff --git a/asap-planner-rs/src/bin/optimizer_cli.rs b/asap-planner-rs/src/bin/optimizer_cli.rs index 01a60e61..7a813f8e 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -178,10 +178,13 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { .expect("clap requires --atomic-costs with --milp"), )?; let Objective::AUCCost { w_cpu, w_mem } = Objective::default(); - let objective = Objective::AUCCost { - w_cpu: args.w_cpu.unwrap_or(w_cpu), - w_mem: args.w_mem.unwrap_or(w_mem), - }; + let (w_cpu, w_mem) = (args.w_cpu.unwrap_or(w_cpu), args.w_mem.unwrap_or(w_mem)); + // All-zero weights make every plan cost 0, so the solver's pick is arbitrary. + anyhow::ensure!( + w_cpu > 0.0 || w_mem > 0.0, + "--w-cpu and --w-mem are both 0; at least one must be positive" + ); + let objective = Objective::AUCCost { w_cpu, w_mem }; tracing::debug!(?objective, cost_rows = costs.len(), "milp: inputs loaded"); let workload = build_milp_workload(config, &facts, args.data_ingestion_interval_ms)?; diff --git a/asap-planner-rs/src/optimizer/milp.rs b/asap-planner-rs/src/optimizer/milp.rs index 9a049685..65562336 100644 --- a/asap-planner-rs/src/optimizer/milp.rs +++ b/asap-planner-rs/src/optimizer/milp.rs @@ -20,6 +20,9 @@ use super::solution::OptimizerItem; const RELATIVE_ERROR: &str = "relative_error"; const MAX_RANK_ERROR: &str = "max_rank_err"; const PRECISION_AT_K: &str = "precision_at_k"; +/// Slack on accuracy tolerances so `1 - sla` rounding (`1 - 0.9 = +/// 0.0999...98`) doesn't reject a row measured exactly at the boundary. +const SLA_EPSILON: f64 = 1e-9; #[derive(Debug, Error)] pub enum MilpError { @@ -217,7 +220,7 @@ fn accuracy_target( capability: Capability, accuracy_sla: f64, ) -> (&'static str, f64, AccuracyDirection) { - let max_error = 1.0 - accuracy_sla; + let max_error = 1.0 - accuracy_sla + SLA_EPSILON; match capability { Capability::SumOrCount | Capability::Min @@ -227,7 +230,7 @@ fn accuracy_target( Capability::Quantile => (MAX_RANK_ERROR, max_error, AccuracyDirection::LowerIsBetter), Capability::TopK => ( PRECISION_AT_K, - accuracy_sla, + accuracy_sla - SLA_EPSILON, AccuracyDirection::HigherIsBetter, ), } @@ -315,6 +318,19 @@ metrics: assert_eq!(ids, vec![format!("0:{query}#0"), format!("0:{query}#1")]); } + #[test] + fn accuracy_exactly_at_the_sla_boundary_passes() { + let w = workload(&group( + "quantile_over_time(0.99, http_requests_total[5m])", + 0.9, + )); + assert!(w.raqes[0].accuracy_ok(0.1)); + assert!(!w.raqes[0].accuracy_ok(0.1001)); + let w = workload(&group("topk(5, http_requests_total)", 0.9)); + assert!(w.raqes[0].accuracy_ok(0.9)); + assert!(!w.raqes[0].accuracy_ok(0.8999)); + } + #[test] fn same_query_at_different_cadences_gets_distinct_ids() { let query = "sum by (job) (http_requests_total)"; @@ -343,7 +359,7 @@ metrics: assert_eq!(r.grouping_labels, labels(&["job"])); assert_eq!(r.accuracy_metric, RELATIVE_ERROR); assert_eq!(r.accuracy_direction, AccuracyDirection::LowerIsBetter); - assert!((r.accuracy_sla - 0.01).abs() < 1e-12); + assert!((r.accuracy_sla - 0.01).abs() < 1e-6); assert_eq!(r.interval_ms, 60_000); assert_eq!(r.latency_sla_ms, None); } @@ -367,7 +383,7 @@ metrics: assert_eq!(r.capability, Capability::TopK); assert_eq!(r.accuracy_metric, PRECISION_AT_K); assert_eq!(r.accuracy_direction, AccuracyDirection::HigherIsBetter); - assert_eq!(r.accuracy_sla, 0.9); + assert!((r.accuracy_sla - 0.9).abs() < 1e-6); // Not the all-labels output set, which would cost one heap per series. assert_eq!(r.grouping_labels, labels(&["job"])); }