From 35ad9cfac1cf9d7c8726f91e6edce832088ddce0 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Sun, 4 Oct 2026 16:17:36 -0400 Subject: [PATCH 1/4] feat(planner): load label-set facts for the optimizer Replace --dataset/--rho with --label-set-facts: series_count per (metric, spatial filter) and cardinality per (metric, spatial filter, grouping labels). Arrival rate is derived as series_count / scrape interval. The label schema now comes from the workload's metrics: hints. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 1 - asap-planner-rs/Cargo.toml | 1 - asap-planner-rs/src/bin/optimizer_cli.rs | 32 +- asap-planner-rs/src/optimizer/dataset.rs | 599 ------------------ asap-planner-rs/src/optimizer/greedy.rs | 65 +- .../src/optimizer/label_set_facts.rs | 481 ++++++++++++++ asap-planner-rs/src/optimizer/mod.rs | 4 +- asap-planner-rs/src/optimizer/pipeline.rs | 116 +++- 8 files changed, 613 insertions(+), 686 deletions(-) delete mode 100644 asap-planner-rs/src/optimizer/dataset.rs create mode 100644 asap-planner-rs/src/optimizer/label_set_facts.rs diff --git a/Cargo.lock b/Cargo.lock index 2e3bd8b1..dc602750 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -135,7 +135,6 @@ dependencies = [ "asap_types", "chrono", "clap 4.6.1", - "csv", "elastic_dsl_utilities", "indexmap 2.14.0", "pretty_assertions", diff --git a/asap-planner-rs/Cargo.toml b/asap-planner-rs/Cargo.toml index a8c7b99e..a94e78d2 100644 --- a/asap-planner-rs/Cargo.toml +++ b/asap-planner-rs/Cargo.toml @@ -37,7 +37,6 @@ anyhow.workspace = true tracing.workspace = true tracing-subscriber.workspace = true clap.workspace = true -csv = "1.4.0" indexmap.workspace = true chrono.workspace = true promql-parser = "0.5.0" diff --git a/asap-planner-rs/src/bin/optimizer_cli.rs b/asap-planner-rs/src/bin/optimizer_cli.rs index b6521b70..591080fa 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -7,7 +7,7 @@ use std::path::PathBuf; use asap_planner::optimizer::{ - load_optional_selected_atomic_cost_table, run_greedy_pipeline, AtomicCostTable, SeriesDataset, + load_optional_selected_atomic_cost_table, run_greedy_pipeline, AtomicCostTable, LabelSetFacts, }; use asap_planner::ControllerConfig; use clap::Parser; @@ -19,21 +19,18 @@ use clap::Parser; )] struct Args { /// Path to a YAML workload config (same format as `asap-planner --input_config`). + /// Its `metrics:` hints are required: they supply each metric's label schema. #[arg(long = "input_config")] input_config: PathBuf, - #[arg(long = "data-ingestion-interval-ms")] + /// Scrape interval; also sets each series' sample rate for arrival rates. + #[arg(long = "data-ingestion-interval-ms", value_parser = clap::value_parser!(u64).range(1..))] data_ingestion_interval_ms: u64, - /// CSV series inventory used to derive metric schemas and label-group counts. - #[arg(long = "dataset")] - dataset: PathBuf, - - /// Placeholder arrival rate (items/sec) applied uniformly to every candidate's - /// IngestCost. Real per-config rates aren't wired up yet — see the open TODOs - /// in .design_docs/optimizer-v1-implementation-plan.md. - #[arg(long = "rho", default_value = "1.0", value_parser = parse_positive_finite)] - rho: f64, + /// 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 @@ -54,14 +51,6 @@ struct Args { verbose: u8, } -fn parse_positive_finite(s: &str) -> Result { - let v: f64 = s.parse().map_err(|_| format!("not a valid number: {s}"))?; - if !v.is_finite() || v <= 0.0 { - return Err(format!("--rho must be a positive finite number, got {v}")); - } - Ok(v) -} - fn main() -> anyhow::Result<()> { let args = Args::parse(); @@ -75,7 +64,7 @@ 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 dataset = SeriesDataset::from_path(&args.dataset)?; + let facts = LabelSetFacts::from_path(&args.label_set_facts)?; let atomic_cost_table = match load_optional_selected_atomic_cost_table( args.atomic_costs.as_deref(), @@ -92,9 +81,8 @@ fn main() -> anyhow::Result<()> { let (streaming, inference) = run_greedy_pipeline( &config, - &dataset, + &facts, args.data_ingestion_interval_ms, - args.rho, &atomic_cost_table, )?; diff --git a/asap-planner-rs/src/optimizer/dataset.rs b/asap-planner-rs/src/optimizer/dataset.rs deleted file mode 100644 index e6adcf92..00000000 --- a/asap-planner-rs/src/optimizer/dataset.rs +++ /dev/null @@ -1,599 +0,0 @@ -use std::collections::{hash_map::Entry, HashMap, HashSet}; -use std::fs::File; -use std::io::Read; -use std::path::Path; - -use asap_types::query_requirements::QueryRequirements; -use asap_types::PromQLSchema; -use csv::StringRecord; -use promql_parser::label::MatchOp; -use promql_parser::parser::{self, Expr}; -use promql_utilities::data_model::KeyByLabelNames; -use thiserror::Error; - -use crate::config::input::MetricDefinition; - -use super::solution::AQE; - -const METRIC_COLUMN: &str = "metric"; - -#[derive(Debug, Error)] -pub enum DatasetError { - #[error("failed to open dataset '{path}': {source}")] - Open { - path: std::path::PathBuf, - source: std::io::Error, - }, - #[error("failed to parse dataset CSV: {0}")] - Csv(#[from] csv::Error), - #[error("dataset must contain a column named 'metric'")] - MissingMetricColumn, - #[error("dataset header contains duplicate column '{0}'")] - DuplicateColumn(String), - #[error("dataset header contains an empty column name")] - EmptyColumn, - #[error("dataset contains no series rows")] - Empty, - #[error("dataset row {row} has an empty metric name")] - EmptyMetric { row: usize }, - #[error( - "metric '{metric}' has inconsistent label columns at row {row}: expected {expected:?}, found {found:?}" - )] - InconsistentMetricSchema { - metric: String, - row: usize, - expected: Vec, - found: Vec, - }, - #[error( - "metric '{metric}' has duplicate series at row {row}; first occurrence was row {first_row}" - )] - DuplicateSeries { - metric: String, - row: usize, - first_row: usize, - }, - #[error("dataset label column '{0}' is not used by any metric")] - UnusedColumn(String), - #[error("dataset is missing workload metric '{0}'")] - MissingMetric(String), - #[error("dataset contains metric '{0}', which is not referenced by the workload")] - ExtraMetric(String), - #[error("metric '{metric}' has no label column '{label}'")] - MissingLabel { metric: String, label: String }, - #[error("metric hint '{metric}' is not present in the dataset")] - ExtraMetricHint { metric: String }, - #[error( - "metric '{metric}' label hint does not match the dataset: hint={hint:?}, dataset={dataset:?}" - )] - MetricHintMismatch { - metric: String, - hint: Vec, - dataset: Vec, - }, - #[error("invalid spatial filter '{filter}': {reason}")] - InvalidFilter { filter: String, reason: String }, - #[error("spatial filter '{filter}' uses unsupported matcher '{matcher}'")] - UnsupportedFilter { filter: String, matcher: String }, - #[error("spatial filter '{filter}' repeats label '{label}'")] - DuplicateFilterLabel { filter: String, label: String }, - #[error("metric '{metric}' has no series matching filter '{filter}'")] - NoMatchingSeries { metric: String, filter: String }, -} - -#[derive(Debug, Clone, Hash, PartialEq, Eq)] -pub struct ProfileKey { - pub metric: String, - pub spatial_filter_normalized: String, - pub grouping_labels: KeyByLabelNames, -} - -impl ProfileKey { - pub fn from_requirements(requirements: &QueryRequirements) -> Self { - Self { - metric: requirements.metric.clone(), - spatial_filter_normalized: requirements.spatial_filter_normalized.clone(), - grouping_labels: requirements.grouping_labels.clone(), - } - } -} - -#[derive(Debug, Clone)] -struct MetricData { - label_names: KeyByLabelNames, - series: Vec>, -} - -#[derive(Debug, Clone)] -pub struct SeriesDataset { - metrics: HashMap, -} - -impl SeriesDataset { - pub fn from_path(path: &Path) -> Result { - let file = File::open(path).map_err(|source| DatasetError::Open { - path: path.to_path_buf(), - source, - })?; - Self::from_reader(file) - } - - pub fn from_reader(reader: R) -> Result { - let mut csv_reader = csv::Reader::from_reader(reader); - let headers = csv_reader.headers()?.clone(); - let header_names = validate_headers(&headers)?; - let metric_index = header_names - .iter() - .position(|name| name == METRIC_COLUMN) - .ok_or(DatasetError::MissingMetricColumn)?; - - let mut raw_rows = Vec::new(); - for (index, record) in csv_reader.records().enumerate() { - let row_number = index + 2; - let record = record?; - let metric = record.get(metric_index).unwrap_or_default().to_string(); - if metric.is_empty() { - return Err(DatasetError::EmptyMetric { row: row_number }); - } - - let values = header_names - .iter() - .enumerate() - .filter_map(|(column_index, column_name)| { - if column_index == metric_index { - return None; - } - let value = record.get(column_index).unwrap_or_default(); - (!value.is_empty()).then(|| (column_name.clone(), value.to_string())) - }) - .collect(); - - raw_rows.push(RawRow { - row_number, - metric, - values, - }); - } - - if raw_rows.is_empty() { - return Err(DatasetError::Empty); - } - - let mut rows_by_metric: HashMap> = HashMap::new(); - for row in raw_rows { - rows_by_metric - .entry(row.metric.clone()) - .or_default() - .push(row); - } - - let mut metrics = HashMap::new(); - let mut used_columns = HashSet::new(); - for (metric, rows) in rows_by_metric { - let expected = sorted_keys(&rows[0].values); - let mut seen_series: HashMap, usize> = HashMap::new(); - let mut series = Vec::with_capacity(rows.len()); - - for row in rows { - let found = sorted_keys(&row.values); - if found != expected { - return Err(DatasetError::InconsistentMetricSchema { - metric: metric.clone(), - row: row.row_number, - expected: expected.clone(), - found, - }); - } - - let key: Vec = expected - .iter() - .map(|label| row.values[label].clone()) - .collect(); - if let Some(first_row) = seen_series.insert(key, row.row_number) { - return Err(DatasetError::DuplicateSeries { - metric: metric.clone(), - row: row.row_number, - first_row, - }); - } - - used_columns.extend(expected.iter().cloned()); - series.push(row.values); - } - - metrics.insert( - metric, - MetricData { - label_names: KeyByLabelNames::new(expected), - series, - }, - ); - } - - for column in header_names { - if column != METRIC_COLUMN && !used_columns.contains(&column) { - return Err(DatasetError::UnusedColumn(column)); - } - } - - Ok(Self { metrics }) - } - - pub fn schema(&self) -> PromQLSchema { - let mut schema = PromQLSchema::new(); - for (metric, data) in &self.metrics { - schema = schema.add_metric(metric.clone(), data.label_names.clone()); - } - schema - } - - pub fn validate_metric_hints( - &self, - hints: Option<&[MetricDefinition]>, - ) -> Result<(), DatasetError> { - let Some(hints) = hints else { - return Ok(()); - }; - - for hint in hints { - let Some(data) = self.metrics.get(&hint.metric) else { - return Err(DatasetError::ExtraMetricHint { - metric: hint.metric.clone(), - }); - }; - let hint_labels = KeyByLabelNames::new(hint.labels.clone()).labels; - if hint_labels != data.label_names.labels { - return Err(DatasetError::MetricHintMismatch { - metric: hint.metric.clone(), - hint: hint_labels, - dataset: data.label_names.labels.clone(), - }); - } - } - Ok(()) - } - - pub fn profile_aqes(&self, aqes: &[AQE]) -> Result, DatasetError> { - let workload_metrics: HashSet<&str> = aqes - .iter() - .map(|aqe| aqe.requirements.metric.as_str()) - .collect(); - - for metric in &workload_metrics { - if !self.metrics.contains_key(*metric) { - return Err(DatasetError::MissingMetric((*metric).to_string())); - } - } - for metric in self.metrics.keys() { - if !workload_metrics.contains(metric.as_str()) { - return Err(DatasetError::ExtraMetric(metric.clone())); - } - } - - let mut profiles = HashMap::new(); - for aqe in aqes { - let key = ProfileKey::from_requirements(&aqe.requirements); - if let Entry::Vacant(entry) = profiles.entry(key) { - entry.insert(self.profile(&aqe.requirements)?); - } - } - Ok(profiles) - } - - pub fn profile(&self, requirements: &QueryRequirements) -> Result { - let data = self - .metrics - .get(&requirements.metric) - .ok_or_else(|| DatasetError::MissingMetric(requirements.metric.clone()))?; - - for label in &requirements.grouping_labels.labels { - if !data.label_names.labels.contains(label) { - return Err(DatasetError::MissingLabel { - metric: requirements.metric.clone(), - label: label.clone(), - }); - } - } - - let matchers = parse_exact_filter( - &requirements.metric, - &requirements.spatial_filter_normalized, - )?; - for (label, _) in &matchers { - if !data.label_names.labels.contains(label) { - return Err(DatasetError::MissingLabel { - metric: requirements.metric.clone(), - label: label.clone(), - }); - } - } - - let matching_series: Vec<&HashMap> = data - .series - .iter() - .filter(|series| { - matchers - .iter() - .all(|(label, value)| series.get(label).is_some_and(|actual| actual == value)) - }) - .collect(); - - if matching_series.is_empty() { - return Err(DatasetError::NoMatchingSeries { - metric: requirements.metric.clone(), - filter: requirements.spatial_filter_normalized.clone(), - }); - } - - if requirements.grouping_labels.is_empty() { - return Ok(1); - } - - let groups: HashSet> = matching_series - .iter() - .map(|series| { - requirements - .grouping_labels - .labels - .iter() - .map(|label| series[label].as_str()) - .collect() - }) - .collect(); - Ok(groups.len() as u64) - } -} - -#[derive(Debug)] -struct RawRow { - row_number: usize, - metric: String, - values: HashMap, -} - -fn validate_headers(headers: &StringRecord) -> Result, DatasetError> { - let mut seen = HashSet::new(); - let mut names = Vec::with_capacity(headers.len()); - for name in headers { - if name.is_empty() { - return Err(DatasetError::EmptyColumn); - } - if !seen.insert(name.to_string()) { - return Err(DatasetError::DuplicateColumn(name.to_string())); - } - names.push(name.to_string()); - } - Ok(names) -} - -fn sorted_keys(values: &HashMap) -> Vec { - let mut keys: Vec = values.keys().cloned().collect(); - keys.sort(); - keys -} - -fn parse_exact_filter(metric: &str, filter: &str) -> Result, DatasetError> { - if filter.is_empty() { - return Ok(Vec::new()); - } - - let selector = format!("{metric}{filter}"); - let expression = parser::parse(&selector).map_err(|error| DatasetError::InvalidFilter { - filter: filter.to_string(), - reason: error.to_string(), - })?; - let Expr::VectorSelector(vector) = expression else { - return Err(DatasetError::InvalidFilter { - filter: filter.to_string(), - reason: "expected a vector selector".to_string(), - }); - }; - - if !vector.matchers.or_matchers.is_empty() { - return Err(DatasetError::UnsupportedFilter { - filter: filter.to_string(), - matcher: "or".to_string(), - }); - } - - let mut seen_labels = HashSet::new(); - let mut matchers = Vec::with_capacity(vector.matchers.matchers.len()); - for matcher in vector.matchers.matchers { - if !matches!(&matcher.op, MatchOp::Equal) { - return Err(DatasetError::UnsupportedFilter { - filter: filter.to_string(), - matcher: matcher.op.to_string(), - }); - } - if !seen_labels.insert(matcher.name.clone()) { - return Err(DatasetError::DuplicateFilterLabel { - filter: filter.to_string(), - label: matcher.name, - }); - } - matchers.push((matcher.name, matcher.value)); - } - Ok(matchers) -} - -#[cfg(test)] -mod tests { - use super::*; - use asap_types::query_requirements::QueryRequirements; - use promql_utilities::data_model::KeyByLabelNames; - - fn requirements(metric: &str, labels: &[&str], filter: &str) -> QueryRequirements { - QueryRequirements { - metric: metric.to_string(), - statistics: vec![], - data_range_ms: 1, - grouping_labels: KeyByLabelNames::new( - labels.iter().map(|label| (*label).to_string()).collect(), - ), - spatial_filter_normalized: filter.to_string(), - topk_count_events: None, - topk_by_labels: None, - } - } - - #[test] - fn derives_metric_schema_and_counts_distinct_group_values() { - let dataset = SeriesDataset::from_reader( - "metric,job,instance,region\nrequests,api,a,us\nrequests,api,b,us\nrequests,worker,c,us\n" - .as_bytes(), - ) - .unwrap(); - - assert_eq!( - dataset.schema().get_labels("requests").unwrap().labels, - vec!["instance", "job", "region"] - ); - assert_eq!( - dataset - .profile(&requirements("requests", &["job"], "")) - .unwrap(), - 2 - ); - assert_eq!( - dataset - .profile(&requirements("requests", &["instance"], "")) - .unwrap(), - 3 - ); - } - - #[test] - fn applies_multiple_exact_filters() { - let dataset = SeriesDataset::from_reader( - "metric,job,instance\nrequests,api,a\nrequests,api,b\nrequests,worker,c\n".as_bytes(), - ) - .unwrap(); - - assert_eq!( - dataset - .profile(&requirements( - "requests", - &["instance"], - "{job=\"api\",instance=\"b\"}" - )) - .unwrap(), - 1 - ); - } - - #[test] - fn empty_grouping_is_one_but_empty_filter_result_fails() { - let dataset = - SeriesDataset::from_reader("metric,job\nrequests,api\nrequests,worker\n".as_bytes()) - .unwrap(); - - assert_eq!( - dataset.profile(&requirements("requests", &[], "")).unwrap(), - 1 - ); - assert!(matches!( - dataset.profile(&requirements("requests", &[], "{job=\"missing\"}")), - Err(DatasetError::NoMatchingSeries { .. }) - )); - } - - #[test] - fn rejects_duplicate_series() { - let error = SeriesDataset::from_reader( - "metric,job,instance\nrequests,api,a\nrequests,api,a\n".as_bytes(), - ) - .unwrap_err(); - - assert!(matches!(error, DatasetError::DuplicateSeries { .. })); - } - - #[test] - fn rejects_schema_changes_within_a_metric() { - let error = SeriesDataset::from_reader( - "metric,job,instance\nrequests,api,a\nrequests,api,\n".as_bytes(), - ) - .unwrap_err(); - - assert!(matches!( - error, - DatasetError::InconsistentMetricSchema { .. } - )); - } - - #[test] - fn rejects_non_exact_filters() { - let dataset = SeriesDataset::from_reader("metric,job\nrequests,api\n".as_bytes()).unwrap(); - - let error = dataset - .profile(&requirements("requests", &[], "{job=~\"api\"}")) - .unwrap_err(); - assert!(matches!(error, DatasetError::UnsupportedFilter { .. })); - } - - #[test] - fn allows_different_stable_schemas_for_different_metrics() { - let dataset = SeriesDataset::from_reader( - "metric,job,instance,region\nrequests,api,a,\nrequests,api,b,\nlatency,,,us-east\nlatency,,,us-west\n" - .as_bytes(), - ) - .unwrap(); - - assert_eq!( - dataset.schema().get_labels("requests").unwrap().labels, - vec!["instance", "job"] - ); - assert_eq!( - dataset.schema().get_labels("latency").unwrap().labels, - vec!["region"] - ); - } - - #[test] - fn rejects_missing_and_extra_workload_metrics() { - let dataset = - SeriesDataset::from_reader("metric,job\nrequests,api\nother,worker\n".as_bytes()) - .unwrap(); - let requests = AQE { - requirements: requirements("requests", &["job"], ""), - query_strings: vec![], - query_frequency_hz: 1.0, - min_t_repeat_ms: 1, - t_repeat_gcd_ms: 1, - }; - - assert!(matches!( - dataset.profile_aqes(&[requests]), - Err(DatasetError::ExtraMetric(metric)) if metric == "other" - )); - } - - #[test] - fn validates_existing_metric_hints_against_dataset_schema() { - let dataset = - SeriesDataset::from_reader("metric,job,instance\nrequests,api,a\n".as_bytes()).unwrap(); - let matching = MetricDefinition { - metric: "requests".into(), - labels: vec!["instance".into(), "job".into()], - }; - let mismatching = MetricDefinition { - metric: "requests".into(), - labels: vec!["job".into()], - }; - - assert!(dataset.validate_metric_hints(Some(&[matching])).is_ok()); - assert!(matches!( - dataset.validate_metric_hints(Some(&[mismatching])), - Err(DatasetError::MetricHintMismatch { .. }) - )); - } - - #[test] - fn rejects_unknown_grouping_labels() { - let dataset = SeriesDataset::from_reader("metric,job\nrequests,api\n".as_bytes()).unwrap(); - - assert!(matches!( - dataset.profile(&requirements("requests", &["instance"], "")), - Err(DatasetError::MissingLabel { .. }) - )); - } -} diff --git a/asap-planner-rs/src/optimizer/greedy.rs b/asap-planner-rs/src/optimizer/greedy.rs index a533a559..df3ba960 100644 --- a/asap-planner-rs/src/optimizer/greedy.rs +++ b/asap-planner-rs/src/optimizer/greedy.rs @@ -5,7 +5,7 @@ use tracing::debug; use super::atomic_costs::{resolve_atomic_costs, AtomicCostTable}; use super::candidate_gen::enumerate_candidates_with_label_group_count; use super::cost_model::{ingest_cost, query_cost, total_cost_rate, AtomicCosts, CostWeights}; -use super::dataset::ProfileKey; +use super::label_set_facts::{ItemFacts, LabelSetKey}; use super::solution::{AQEAssignment, OptimizerSolution, AQE}; /// Greedily assign each AQE to its independently-cheapest candidate config. @@ -14,9 +14,8 @@ use super::solution::{AQEAssignment, OptimizerSolution, AQE}; /// two AQEs could share one. The Phase 3 MIP finds sharing opportunities; this /// is the v1 baseline. /// -/// `arrival_rate_hz` is the per-item arrival rate used for every candidate's IngestCost. -/// Real per-config rates need Prometheus scrape-rate × series-count data, -/// which isn't wired up yet — a single placeholder value is applied uniformly. +/// `facts` supplies each AQE's group cardinality and arrival rate, keyed by +/// its label set; every AQE must have an entry. /// /// Each candidate is costed at its own `(sketch_type, params)` via /// `atomic_cost_table` (see ASAPQuery#524) rather than one cost applied to @@ -25,25 +24,22 @@ use super::solution::{AQEAssignment, OptimizerSolution, AQE}; pub fn greedy_assign( aqes: Vec, scrape_interval_ms: u64, - arrival_rate_hz: f64, atomic_cost_table: &AtomicCostTable, weights: &CostWeights, - label_group_counts: &HashMap, + facts: &HashMap, ) -> OptimizerSolution { let mut solution = OptimizerSolution::empty(); for aqe in aqes { - let profile_key = ProfileKey::from_requirements(&aqe.requirements); - let label_group_count = *label_group_counts.get(&profile_key).unwrap_or_else(|| { - panic!( - "missing dataset profile for metric '{}' and grouping labels {:?}", - aqe.requirements.metric, aqe.requirements.grouping_labels.labels - ) - }); + let key = LabelSetKey::from_requirements(&aqe.requirements); + let item_facts = *facts + .get(&key) + .unwrap_or_else(|| panic!("missing label-set facts for {key}")); + let arrival_rate_hz = item_facts.arrival_rate_per_sec; let candidates = enumerate_candidates_with_label_group_count( &aqe, scrape_interval_ms, - label_group_count, + item_facts.cardinality, ); let (best, costs) = candidates @@ -128,32 +124,34 @@ mod tests { } } + /// One group, one item/sec for every AQE's label set. + fn unit_facts(aqes: &[AQE]) -> HashMap { + aqes.iter() + .map(|aqe| { + ( + LabelSetKey::from_requirements(&aqe.requirements), + ItemFacts { + cardinality: 1, + arrival_rate_per_sec: 1.0, + }, + ) + }) + .collect() + } + #[test] fn assigns_unique_ids_to_each_deployed_config() { let aqes = vec![ make_aqe(Statistic::Min, 300_000, 300_000, 1.0 / 60.0), make_aqe(Statistic::Max, 300_000, 300_000, 1.0 / 60.0), ]; + let facts = unit_facts(&aqes); let solution = greedy_assign( aqes, 60_000, - 1.0, &AtomicCostTable::default(), &CostWeights::default(), - &HashMap::from([ - ( - ProfileKey::from_requirements( - &make_aqe(Statistic::Min, 300_000, 300_000, 1.0 / 60.0).requirements, - ), - 1, - ), - ( - ProfileKey::from_requirements( - &make_aqe(Statistic::Max, 300_000, 300_000, 1.0 / 60.0).requirements, - ), - 1, - ), - ]), + &facts, ); let mut seen_ids: StdHashMap = StdHashMap::new(); @@ -186,10 +184,9 @@ mod tests { let solution = greedy_assign( vec![aqe.clone()], 60_000, - 1.0, &AtomicCostTable::default(), &CostWeights::default(), - &HashMap::from([(ProfileKey::from_requirements(&aqe.requirements), 1)]), + &unit_facts(&[aqe]), ); assert_eq!(solution.num_exact_fallback(), 1); assert!(solution.deployed_configs().is_empty()); @@ -203,10 +200,9 @@ mod tests { let solution = greedy_assign( vec![aqe.clone()], 60_000, - 1.0, &AtomicCostTable::default(), &CostWeights::default(), - &HashMap::from([(ProfileKey::from_requirements(&aqe.requirements), 1)]), + &unit_facts(&[aqe]), ); assert_eq!(solution.num_exact_fallback(), 1); @@ -231,10 +227,9 @@ mod tests { let solution = greedy_assign( vec![aqe.clone()], 60_000, - 1.0, &table, &CostWeights::default(), - &HashMap::from([(ProfileKey::from_requirements(&aqe.requirements), 1)]), + &unit_facts(&[aqe]), ); assert_eq!(solution.num_exact_fallback(), 0); diff --git a/asap-planner-rs/src/optimizer/label_set_facts.rs b/asap-planner-rs/src/optimizer/label_set_facts.rs new file mode 100644 index 00000000..183cb26a --- /dev/null +++ b/asap-planner-rs/src/optimizer/label_set_facts.rs @@ -0,0 +1,481 @@ +//! Externally provided label-set facts: how many raw series feed each +//! (metric, spatial filter), and how many groups each grouping produces. +//! The optimizer never estimates these. + +use std::collections::{HashMap, HashSet}; +use std::path::{Path, PathBuf}; + +use asap_types::query_requirements::QueryRequirements; +use asap_types::utils::normalize_spatial_filter; +use promql_utilities::data_model::KeyByLabelNames; +use serde::Deserialize; +use thiserror::Error; + +use super::solution::AQE; + +#[derive(Debug, Error)] +pub enum LabelSetFactsError { + #[error("failed to read label-set facts '{path}': {source}")] + Read { + path: PathBuf, + source: std::io::Error, + }, + #[error("failed to parse label-set facts: {0}")] + Parse(#[from] serde_yaml::Error), + #[error("duplicate series facts for {0}")] + DuplicateSeries(SeriesKey), + #[error("duplicate group facts for {0}")] + DuplicateGroup(LabelSetKey), + #[error("series_count must be at least 1 for {0}")] + ZeroSeriesCount(SeriesKey), + #[error("group facts for {0} have no matching series facts")] + GroupWithoutSeries(LabelSetKey), + #[error( + "cardinality {cardinality} for {key} must be between 1 and series_count {series_count}" + )] + CardinalityOutOfRange { + key: LabelSetKey, + cardinality: u64, + series_count: u64, + }, + #[error("cardinality for {key} must be 1 when grouping_labels is empty, got {cardinality}")] + FullAggregationCardinality { key: LabelSetKey, cardinality: u64 }, + #[error( + "workload config has no `metrics:` hints; they are required to resolve grouping labels" + )] + MissingMetricHints, + #[error("workload metrics missing from `metrics:` hints: {0:?}")] + MetricsWithoutHints(Vec), + #[error("missing label-set facts (keys shown as the optimizer expects them):\n{}", .0.join("\n"))] + MissingFacts(Vec), +} + +/// (metric, normalized spatial filter): the raw series stream an item reads. +#[derive(Debug, Clone, Hash, PartialEq, Eq)] +pub struct SeriesKey { + pub metric: String, + pub spatial_filter_normalized: String, +} + +impl std::fmt::Display for SeriesKey { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!( + f, + "metric={:?} spatial_filter={:?}", + self.metric, self.spatial_filter_normalized + ) + } +} + +/// (metric, normalized spatial filter, grouping labels): one grouped stream. +#[derive(Debug, Clone, Hash, PartialEq, Eq)] +pub struct LabelSetKey { + pub metric: String, + pub spatial_filter_normalized: String, + pub grouping_labels: KeyByLabelNames, +} + +impl LabelSetKey { + pub fn from_requirements(requirements: &QueryRequirements) -> Self { + Self { + metric: requirements.metric.clone(), + spatial_filter_normalized: requirements.spatial_filter_normalized.clone(), + grouping_labels: requirements.grouping_labels.clone(), + } + } + + fn series_key(&self) -> SeriesKey { + SeriesKey { + metric: self.metric.clone(), + spatial_filter_normalized: self.spatial_filter_normalized.clone(), + } + } +} + +impl std::fmt::Display for LabelSetKey { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!( + f, + "{} grouping_labels={:?}", + self.series_key(), + self.grouping_labels.labels + ) + } +} + +/// Facts for one item's label set, ready for costing. +#[derive(Debug, Clone, Copy, PartialEq)] +pub struct ItemFacts { + /// Distinct grouping-label value combinations. + pub cardinality: u64, + /// Aggregate items/sec into the grouped stream, across all groups. + pub arrival_rate_per_sec: f64, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct FactsFile { + series: Vec, + groups: Vec, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct SeriesFact { + metric: String, + spatial_filter: String, + series_count: u64, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct GroupFact { + metric: String, + spatial_filter: String, + grouping_labels: Vec, + cardinality: u64, +} + +#[derive(Debug, Clone)] +pub struct LabelSetFacts { + series_counts: HashMap, + cardinalities: HashMap, +} + +impl LabelSetFacts { + pub fn from_path(path: &Path) -> Result { + let yaml = std::fs::read_to_string(path).map_err(|source| LabelSetFactsError::Read { + path: path.to_path_buf(), + source, + })?; + Self::from_yaml(&yaml) + } + + pub fn from_yaml(yaml: &str) -> Result { + let file: FactsFile = serde_yaml::from_str(yaml)?; + + let mut series_counts = HashMap::new(); + for fact in file.series { + let key = SeriesKey { + metric: fact.metric, + spatial_filter_normalized: normalize_spatial_filter(&fact.spatial_filter), + }; + if fact.series_count == 0 { + return Err(LabelSetFactsError::ZeroSeriesCount(key)); + } + if series_counts + .insert(key.clone(), fact.series_count) + .is_some() + { + return Err(LabelSetFactsError::DuplicateSeries(key)); + } + } + + let mut cardinalities = HashMap::new(); + for fact in file.groups { + let key = LabelSetKey { + metric: fact.metric, + spatial_filter_normalized: normalize_spatial_filter(&fact.spatial_filter), + grouping_labels: KeyByLabelNames::new(fact.grouping_labels), + }; + let Some(&series_count) = series_counts.get(&key.series_key()) else { + return Err(LabelSetFactsError::GroupWithoutSeries(key)); + }; + if key.grouping_labels.labels.is_empty() && fact.cardinality != 1 { + return Err(LabelSetFactsError::FullAggregationCardinality { + key, + cardinality: fact.cardinality, + }); + } + if fact.cardinality == 0 || fact.cardinality > series_count { + return Err(LabelSetFactsError::CardinalityOutOfRange { + key, + cardinality: fact.cardinality, + series_count, + }); + } + if cardinalities + .insert(key.clone(), fact.cardinality) + .is_some() + { + return Err(LabelSetFactsError::DuplicateGroup(key)); + } + } + + Ok(Self { + series_counts, + cardinalities, + }) + } + + /// Look up facts for every AQE's label set. Errors list every missing key + /// at once; facts no AQE uses are only warned about, so one file can serve + /// several workloads. + /// + /// Arrival rate assumes each series yields one sample per scrape: + /// `series_count / scrape interval`. + // ponytail: overestimates sparse or irregular series; take a measured rate if that matters. + pub fn resolve( + &self, + aqes: &[AQE], + scrape_interval_ms: u64, + ) -> Result, LabelSetFactsError> { + let mut resolved = HashMap::new(); + let mut missing = Vec::new(); + for aqe in aqes { + let key = LabelSetKey::from_requirements(&aqe.requirements); + if resolved.contains_key(&key) { + continue; + } + let series_count = self.series_counts.get(&key.series_key()); + let cardinality = self.cardinalities.get(&key); + if series_count.is_none() { + missing.push(format!("series: {}", key.series_key())); + } + if cardinality.is_none() { + missing.push(format!("groups: {key}")); + } + if let (Some(&series_count), Some(&cardinality)) = (series_count, cardinality) { + let arrival_rate_per_sec = series_count as f64 * 1000.0 / scrape_interval_ms as f64; + resolved.insert( + key, + ItemFacts { + cardinality, + arrival_rate_per_sec, + }, + ); + } + } + if !missing.is_empty() { + missing.sort(); + missing.dedup(); + return Err(LabelSetFactsError::MissingFacts(missing)); + } + + let used_series: HashSet = + resolved.keys().map(LabelSetKey::series_key).collect(); + for key in self.series_counts.keys() { + if !used_series.contains(key) { + tracing::warn!(%key, "series facts match no workload item"); + } + } + for key in self.cardinalities.keys() { + if !resolved.contains_key(key) { + tracing::warn!(%key, "group facts match no workload item"); + } + } + Ok(resolved) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use promql_utilities::query_logics::enums::Statistic; + + fn aqe(metric: &str, filter: &str, labels: &[&str]) -> AQE { + AQE { + requirements: QueryRequirements { + metric: metric.into(), + statistics: vec![Statistic::Sum], + data_range_ms: 60_000, + grouping_labels: KeyByLabelNames::new( + labels.iter().map(|l| l.to_string()).collect(), + ), + spatial_filter_normalized: normalize_spatial_filter(filter), + topk_count_events: None, + topk_by_labels: None, + }, + query_strings: vec!["q".into()], + query_frequency_hz: 1.0 / 60.0, + min_t_repeat_ms: 60_000, + t_repeat_gcd_ms: 60_000, + } + } + + const FACTS: &str = r#" +series: + - metric: http_requests_total + spatial_filter: 'job="api",env="prod"' + series_count: 1000 +groups: + - metric: http_requests_total + spatial_filter: '{env="prod",job="api"}' + grouping_labels: [service, endpoint] + cardinality: 50 + - metric: http_requests_total + spatial_filter: 'job="api",env="prod"' + grouping_labels: [] + cardinality: 1 +"#; + + #[test] + fn resolves_cardinality_and_derives_arrival_rate() { + let facts = LabelSetFacts::from_yaml(FACTS).unwrap(); + let item = aqe( + "http_requests_total", + r#"job="api",env="prod""#, + &["endpoint", "service"], + ); + let resolved = facts.resolve(std::slice::from_ref(&item), 15_000).unwrap(); + let got = resolved[&LabelSetKey::from_requirements(&item.requirements)]; + assert_eq!(got.cardinality, 50); + assert!((got.arrival_rate_per_sec - 1000.0 / 15.0).abs() < 1e-9); + } + + #[test] + fn spatial_filter_matches_after_normalization() { + // File writes the matchers in a different order and with braces. + let facts = LabelSetFacts::from_yaml(FACTS).unwrap(); + let item = aqe("http_requests_total", r#"env="prod",job="api""#, &[]); + assert!(facts.resolve(&[item], 15_000).is_ok()); + } + + #[test] + fn missing_facts_lists_every_missing_key() { + let facts = LabelSetFacts::from_yaml(FACTS).unwrap(); + let err = facts + .resolve( + &[ + aqe( + "http_requests_total", + r#"job="api",env="prod""#, + &["service"], + ), + aqe("other_metric", "", &[]), + ], + 15_000, + ) + .unwrap_err(); + let LabelSetFactsError::MissingFacts(missing) = err else { + panic!("expected MissingFacts, got {err:?}"); + }; + assert_eq!(missing.len(), 3, "{missing:?}"); + assert!(missing.iter().any(|m| m.contains("\"service\""))); + assert!(missing + .iter() + .any(|m| m.starts_with("series:") && m.contains("other_metric"))); + assert!(missing + .iter() + .any(|m| m.starts_with("groups:") && m.contains("other_metric"))); + } + + #[test] + fn group_without_series_is_rejected() { + let yaml = r#" +series: [] +groups: + - {metric: m, spatial_filter: "", grouping_labels: [a], cardinality: 1} +"#; + assert!(matches!( + LabelSetFacts::from_yaml(yaml), + Err(LabelSetFactsError::GroupWithoutSeries(_)) + )); + } + + #[test] + fn cardinality_above_series_count_is_rejected() { + let yaml = r#" +series: + - {metric: m, spatial_filter: "", series_count: 3} +groups: + - {metric: m, spatial_filter: "", grouping_labels: [a], cardinality: 4} +"#; + assert!(matches!( + LabelSetFacts::from_yaml(yaml), + Err(LabelSetFactsError::CardinalityOutOfRange { + cardinality: 4, + series_count: 3, + .. + }) + )); + } + + #[test] + fn zero_cardinality_and_zero_series_count_are_rejected() { + let zero_card = r#" +series: + - {metric: m, spatial_filter: "", series_count: 3} +groups: + - {metric: m, spatial_filter: "", grouping_labels: [a], cardinality: 0} +"#; + assert!(matches!( + LabelSetFacts::from_yaml(zero_card), + Err(LabelSetFactsError::CardinalityOutOfRange { cardinality: 0, .. }) + )); + let zero_series = r#" +series: + - {metric: m, spatial_filter: "", series_count: 0} +groups: [] +"#; + assert!(matches!( + LabelSetFacts::from_yaml(zero_series), + Err(LabelSetFactsError::ZeroSeriesCount(_)) + )); + } + + #[test] + fn full_aggregation_requires_cardinality_one() { + let yaml = r#" +series: + - {metric: m, spatial_filter: "", series_count: 3} +groups: + - {metric: m, spatial_filter: "", grouping_labels: [], cardinality: 2} +"#; + assert!(matches!( + LabelSetFacts::from_yaml(yaml), + Err(LabelSetFactsError::FullAggregationCardinality { cardinality: 2, .. }) + )); + } + + #[test] + fn duplicate_keys_are_rejected() { + // Same key after normalization and label sorting. + let dup_group = r#" +series: + - {metric: m, spatial_filter: 'a="1",b="2"', series_count: 3} +groups: + - {metric: m, spatial_filter: 'a="1",b="2"', grouping_labels: [x, y], cardinality: 1} + - {metric: m, spatial_filter: 'b="2",a="1"', grouping_labels: [y, x], cardinality: 2} +"#; + assert!(matches!( + LabelSetFacts::from_yaml(dup_group), + Err(LabelSetFactsError::DuplicateGroup(_)) + )); + let dup_series = r#" +series: + - {metric: m, spatial_filter: "", series_count: 3} + - {metric: m, spatial_filter: "", series_count: 4} +groups: [] +"#; + assert!(matches!( + LabelSetFacts::from_yaml(dup_series), + Err(LabelSetFactsError::DuplicateSeries(_)) + )); + } + + #[test] + fn omitted_grouping_labels_or_filter_is_a_parse_error() { + // No implicit defaults: an absent field must not silently mean + // "full aggregation" or "unfiltered". + let no_labels = r#" +series: + - {metric: m, spatial_filter: "", series_count: 3} +groups: + - {metric: m, spatial_filter: "", cardinality: 1} +"#; + assert!(matches!( + LabelSetFacts::from_yaml(no_labels), + Err(LabelSetFactsError::Parse(_)) + )); + let no_filter = r#" +series: + - {metric: m, series_count: 3} +groups: [] +"#; + assert!(matches!( + LabelSetFacts::from_yaml(no_filter), + Err(LabelSetFactsError::Parse(_)) + )); + } +} diff --git a/asap-planner-rs/src/optimizer/mod.rs b/asap-planner-rs/src/optimizer/mod.rs index 521db371..e6d8692f 100644 --- a/asap-planner-rs/src/optimizer/mod.rs +++ b/asap-planner-rs/src/optimizer/mod.rs @@ -3,8 +3,8 @@ pub mod atomic_costs; pub mod candidate_gen; pub mod constants; pub mod cost_model; -pub mod dataset; pub mod greedy; +pub mod label_set_facts; pub mod pipeline; pub mod sketch_properties; pub mod solution; @@ -20,8 +20,8 @@ pub use candidate_gen::{ enumerate_candidates, enumerate_candidates_with_label_group_count, CandidateConfig, }; pub use cost_model::{ingest_cost, query_cost, total_cost_rate, AtomicCosts, CostWeights}; -pub use dataset::{DatasetError, ProfileKey, SeriesDataset}; pub use greedy::greedy_assign; +pub use label_set_facts::{ItemFacts, LabelSetFacts, LabelSetFactsError, LabelSetKey, SeriesKey}; pub use pipeline::{run_all_exact_pipeline, run_greedy_pipeline}; pub use sketch_properties::{sketch_properties, SketchProperties}; pub use solution::{AQEAssignment, OptimizerSolution, QueryMethod, AQE}; diff --git a/asap-planner-rs/src/optimizer/pipeline.rs b/asap-planner-rs/src/optimizer/pipeline.rs index 54387c70..8992ceae 100644 --- a/asap-planner-rs/src/optimizer/pipeline.rs +++ b/asap-planner-rs/src/optimizer/pipeline.rs @@ -7,8 +7,8 @@ use crate::config::input::ControllerConfig; use super::aqe_extractor::{extract_aqes, RQE}; use super::atomic_costs::AtomicCostTable; use super::cost_model::CostWeights; -use super::dataset::SeriesDataset; use super::greedy::greedy_assign; +use super::label_set_facts::{LabelSetFacts, LabelSetFactsError}; use super::solution::{OptimizerSolution, AQE}; use super::translator::{translate, TranslationSummary}; @@ -74,39 +74,50 @@ pub fn run_all_exact_pipeline( /// No cross-AQE sharing — every deployed sketch serves exactly one AQE, even /// if two AQEs could share one. The Phase 3 MIP finds sharing opportunities. /// -/// The dataset supplies each AQE's metric schema and label-group count before -/// candidate selection. `arrival_rate_hz` remains a uniform placeholder until -/// per-config scrape-rate data is available. +/// The workload's `metrics:` hints supply the label schema; `facts` supply each +/// AQE's group cardinality and series count. pub fn run_greedy_pipeline( config: &ControllerConfig, - dataset: &SeriesDataset, + facts: &LabelSetFacts, scrape_interval_ms: u64, - arrival_rate_hz: f64, atomic_cost_table: &AtomicCostTable, -) -> Result<(StreamingConfig, InferenceConfig), super::dataset::DatasetError> { - dataset.validate_metric_hints(config.metrics.as_deref())?; - let schema = dataset.schema(); +) -> Result<(StreamingConfig, InferenceConfig), LabelSetFactsError> { + if config.metrics.is_none() { + return Err(LabelSetFactsError::MissingMetricHints); + } + let schema = config.schema_from_hints(); let rqes = config_to_rqes(config); let aqes = extract_aqes(&rqes, &schema, scrape_interval_ms); - let label_group_counts = dataset.profile_aqes(&aqes)?; - for (key, count) in &label_group_counts { + // 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)); + } + + let item_facts = facts.resolve(&aqes, scrape_interval_ms)?; + for (key, item) in &item_facts { tracing::info!( - metric = %key.metric, - spatial_filter = %key.spatial_filter_normalized, - grouping_labels = ?key.grouping_labels.labels, - label_group_count = *count, - "optimizer dataset profile" + %key, + cardinality = item.cardinality, + arrival_rate_per_sec = item.arrival_rate_per_sec, + "optimizer label-set facts" ); } let solution = greedy_assign( aqes, scrape_interval_ms, - arrival_rate_hz, atomic_cost_table, &CostWeights::default(), - &label_group_counts, + &item_facts, ); Ok(finish_pipeline(solution, "greedy")) @@ -169,25 +180,78 @@ mod tests { assert!(streaming.get_all_aggregation_configs().is_empty()); } + fn with_hints(mut config: ControllerConfig, hints: &[(&str, &[&str])]) -> ControllerConfig { + use crate::config::input::MetricDefinition; + + config.metrics = Some( + hints + .iter() + .map(|(metric, labels)| MetricDefinition { + metric: metric.to_string(), + labels: labels.iter().map(|l| l.to_string()).collect(), + }) + .collect(), + ); + config + } + + // `min_over_time` keeps every label, so the grouping is the hinted [job]. + const METRIC_FACTS: &str = r#" +series: + - {metric: metric, spatial_filter: "", series_count: 4} +groups: + - {metric: metric, spatial_filter: "", grouping_labels: [job], cardinality: 4} +"#; + #[test] fn greedy_pipeline_deploys_a_config_for_a_mergeable_aqe() { - let config = make_config(&[("min_over_time(metric[5m])", 60_000)]); - let dataset = SeriesDataset::from_reader("metric,job\nmetric,api\n".as_bytes()).unwrap(); + let config = with_hints( + make_config(&[("min_over_time(metric[5m])", 60_000)]), + &[("metric", &["job"])], + ); + let facts = LabelSetFacts::from_yaml(METRIC_FACTS).unwrap(); let (streaming, inference) = - run_greedy_pipeline(&config, &dataset, 60_000, 1.0, &AtomicCostTable::default()) - .unwrap(); + run_greedy_pipeline(&config, &facts, 60_000, &AtomicCostTable::default()).unwrap(); assert!(!streaming.get_all_aggregation_configs().is_empty()); assert!(!inference.query_configs.is_empty()); } #[test] - fn greedy_pipeline_fails_when_dataset_does_not_match_workload() { + fn greedy_pipeline_requires_metric_hints() { let config = make_config(&[("min_over_time(metric[5m])", 60_000)]); - let dataset = SeriesDataset::from_reader("metric,job\nother,api\n".as_bytes()).unwrap(); + let facts = LabelSetFacts::from_yaml(METRIC_FACTS).unwrap(); + assert!(matches!( + run_greedy_pipeline(&config, &facts, 60_000, &AtomicCostTable::default()), + Err(LabelSetFactsError::MissingMetricHints) + )); + } + #[test] + fn greedy_pipeline_fails_when_workload_metric_has_no_hint() { + let config = with_hints( + make_config(&[ + ("min_over_time(metric[5m])", 60_000), + ("max_over_time(unhinted[5m])", 60_000), + ]), + &[("metric", &["job"])], + ); + let facts = LabelSetFacts::from_yaml(METRIC_FACTS).unwrap(); + assert!(matches!( + run_greedy_pipeline(&config, &facts, 60_000, &AtomicCostTable::default()), + Err(LabelSetFactsError::MetricsWithoutHints(metrics)) if metrics == ["unhinted"] + )); + } + + #[test] + fn greedy_pipeline_fails_when_facts_do_not_cover_workload() { + let config = with_hints( + make_config(&[("sum by (instance) (metric)", 60_000)]), + &[("metric", &["job", "instance"])], + ); + let facts = LabelSetFacts::from_yaml(METRIC_FACTS).unwrap(); assert!(matches!( - run_greedy_pipeline(&config, &dataset, 60_000, 1.0, &AtomicCostTable::default()), - Err(super::super::dataset::DatasetError::MissingMetric(metric)) if metric == "metric" + run_greedy_pipeline(&config, &facts, 60_000, &AtomicCostTable::default()), + Err(LabelSetFactsError::MissingFacts(_)) )); } From 99fdbfebfd0248c4e432ef39d4c486de76f56db8 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Sun, 4 Oct 2026 17:25:02 -0400 Subject: [PATCH 2/4] feat(planner): cost optimizer candidates by label-set cardinality Cardinality now scales cost by how each sketch holds its groups: keyed maps grow per key, fixed-size sketches don't, and grouped queries read one value per output group (one heap read per top-k bucket). Trivial accumulators get an analytical per-key memory estimate, and only their Multiple* forms are proposed. Optimizer configs now use the legacy planner's grouping/aggregated label split, so keyed sketches deploy as one instance rather than one per group, and CMS/HydraKLL deploy their paired DeltaSet key aggregation. Co-Authored-By: Claude Opus 5.5 --- .../src/query_logics/logics.rs | 9 +- asap-planner-rs/src/bin/candidate_gen_dump.rs | 24 +- asap-planner-rs/src/optimizer/atomic_costs.rs | 85 +++++- .../src/optimizer/candidate_gen.rs | 243 ++++++++++++++-- asap-planner-rs/src/optimizer/constants.rs | 20 +- asap-planner-rs/src/optimizer/cost_model.rs | 259 ++++++++++++------ asap-planner-rs/src/optimizer/greedy.rs | 109 ++++++-- .../src/optimizer/label_set_facts.rs | 111 +++++--- asap-planner-rs/src/optimizer/mod.rs | 4 +- asap-planner-rs/src/optimizer/pipeline.rs | 9 +- .../src/optimizer/sketch_properties.rs | 28 +- asap-planner-rs/src/optimizer/solution.rs | 5 + asap-planner-rs/src/optimizer/translator.rs | 13 +- asap-planner-rs/src/planner/agg_config.rs | 14 +- 14 files changed, 706 insertions(+), 227 deletions(-) diff --git a/asap-common/dependencies/rs/promql_utilities/src/query_logics/logics.rs b/asap-common/dependencies/rs/promql_utilities/src/query_logics/logics.rs index a2db46d5..0f94fd9d 100644 --- a/asap-common/dependencies/rs/promql_utilities/src/query_logics/logics.rs +++ b/asap-common/dependencies/rs/promql_utilities/src/query_logics/logics.rs @@ -86,12 +86,11 @@ pub fn does_precompute_operator_support_subpopulations( // own top-k keys via CmsHeapItem. So it supports subpopulations the // same way MultipleSum does: labels go in `aggregated`, not // `grouping` (see `set_subpopulation_labels`). - AggregationType::CountMinSketchWithHeap if matches!(statistic, Statistic::Topk) => true, + AggregationType::CountMinSketchWithHeap => true, - AggregationType::HLL => false, - - // Default: not supported - _ => panic!("Unexpected precompute operator: {}", precompute_operator), + AggregationType::HLL + | AggregationType::SetAggregator + | AggregationType::DeltaSetAggregator => false, } } diff --git a/asap-planner-rs/src/bin/candidate_gen_dump.rs b/asap-planner-rs/src/bin/candidate_gen_dump.rs index b0af6864..8c39e9b6 100644 --- a/asap-planner-rs/src/bin/candidate_gen_dump.rs +++ b/asap-planner-rs/src/bin/candidate_gen_dump.rs @@ -80,7 +80,11 @@ fn main() -> anyhow::Result<()> { println!(" queries: {:?}", aqe.query_strings); let candidates = enumerate_candidates(aqe, args.scrape_interval_ms); - print_candidates_grouped(&candidates, &atomic_cost_table); + print_candidates_grouped( + &candidates, + &atomic_cost_table, + aqe.requirements.grouping_labels.len(), + ); } Ok(()) @@ -93,7 +97,11 @@ fn main() -> anyhow::Result<()> { /// params (M) [× N windows = NM total]: /// ... /// EXACT is printed last as a single line. -fn print_candidates_grouped(candidates: &[CandidateConfig], atomic_cost_table: &AtomicCostTable) { +fn print_candidates_grouped( + candidates: &[CandidateConfig], + atomic_cost_table: &AtomicCostTable, + n_grouping_labels: usize, +) { // Collect unique (agg_type_str, sub_type) keys in first-seen order. let mut group_order: Vec<(String, String)> = Vec::new(); // (agg_type_str, sub_type) -> (agg_type, unique windows, unique params) @@ -197,10 +205,18 @@ fn print_candidates_grouped(candidates: &[CandidateConfig], atomic_cost_table: & } else { let param_map: HashMap = p.iter().map(|(k, v)| (k.clone(), v.clone())).collect(); - match resolve_atomic_costs(atomic_cost_table, *agg_type, ¶m_map) { + match resolve_atomic_costs( + atomic_cost_table, + *agg_type, + ¶m_map, + n_grouping_labels, + ) { Some(costs) => { - let label = if costs == AtomicCosts::default() { + let stub = AtomicCosts::default(); + let label = if costs == stub { "stub" + } else if costs.insert_cpu_secs == stub.insert_cpu_secs { + "analytical mem, stub cpu" } else { "real" }; diff --git a/asap-planner-rs/src/optimizer/atomic_costs.rs b/asap-planner-rs/src/optimizer/atomic_costs.rs index 7dc0bdd3..ceb3c887 100644 --- a/asap-planner-rs/src/optimizer/atomic_costs.rs +++ b/asap-planner-rs/src/optimizer/atomic_costs.rs @@ -17,7 +17,8 @@ use serde_json::Value; use super::constants::{ CMS_HEAP_AVERAGE_KEY_BYTES, CMS_HEAP_COUNTER_BYTES, CMS_HEAP_ENTRY_OVERHEAD_BYTES, - CMS_HEAP_REFERENCE_HEAP_SIZE, EXACT_QUERY_CPU_SECS, SUBTRACT_CPU_SECS, + CMS_HEAP_REFERENCE_HEAP_SIZE, EXACT_QUERY_CPU_SECS, HASH_TABLE_SLACK, INCREASE_VALUE_BYTES, + LABEL_VALUE_CODE_BYTES, MIN_MAX_VALUE_BYTES, SUBTRACT_CPU_SECS, SUM_VALUE_BYTES, }; use super::cost_model::AtomicCosts; @@ -215,6 +216,17 @@ fn sketch_bench_key( } } +/// Bytes of one group's value for the trivial accumulators, whose per-group +/// and `Multiple*` forms store the same entry per group. +fn trivial_value_bytes(agg_type: AggregationType) -> Option { + match agg_type { + AggregationType::Sum | AggregationType::MultipleSum => Some(SUM_VALUE_BYTES), + AggregationType::MinMax | AggregationType::MultipleMinMax => Some(MIN_MAX_VALUE_BYTES), + AggregationType::Increase | AggregationType::MultipleIncrease => Some(INCREASE_VALUE_BYTES), + _ => None, + } +} + /// Resolve the [`AtomicCosts`] a candidate should be costed at. /// /// - `agg_type` outside the benchmarked families (see [`sketch_bench_key`]): @@ -233,15 +245,32 @@ fn sketch_bench_key( /// reference data returns `None` and drops the candidate. /// - Other benchmarked families, no matching row: `None` — drop the candidate, /// per #524. +/// - Trivial accumulators (see [`trivial_value_bytes`]): stub CPU costs, but +/// memory is the analytical size of one group's entry for +/// `n_grouping_labels` labels, so per-group and `Multiple*` twins tie. pub fn resolve_atomic_costs( table: &AtomicCostTable, agg_type: AggregationType, params: &HashMap, + n_grouping_labels: usize, ) -> Option { if agg_type == AggregationType::CountMinSketchWithHeap { return resolve_cms_heap_costs(table, params, &CmsHeapCostAssumptions::default()); } + if let Some(value_bytes) = trivial_value_bytes(agg_type) { + tracing::warn!( + ?agg_type, + "no sketch-bench CPU data for this family; using stub CPU costs and analytical memory" + ); + return Some(AtomicCosts { + mem_bytes_per_instance: (n_grouping_labels as f64 * LABEL_VALUE_CODE_BYTES + + value_bytes) + * HASH_TABLE_SLACK, + ..AtomicCosts::default() + }); + } + let Some((sketch, sketch_params)) = sketch_bench_key(agg_type, params) else { tracing::warn!( ?agg_type, @@ -617,6 +646,7 @@ mod tests { &table, AggregationType::CountMinSketch, &cms_params(3, 1024), + 0, ) .expect("exact grid point must resolve"); assert_eq!(costs.mem_bytes_per_instance, 3.0 * 1024.0 * 4.0); @@ -629,21 +659,47 @@ mod tests { #[test] fn cms_param_point_outside_the_grid_drops_the_candidate() { let table = vec![cms_entry(3, 1024)]; - assert!( - resolve_atomic_costs(&table, AggregationType::CountMinSketch, &cms_params(7, 999)) - .is_none() - ); + assert!(resolve_atomic_costs( + &table, + AggregationType::CountMinSketch, + &cms_params(7, 999), + 0 + ) + .is_none()); } #[test] fn unbenchmarked_family_falls_back_to_the_stub() { let table: AtomicCostTable = vec![]; - let costs = resolve_atomic_costs(&table, AggregationType::Sum, &HashMap::new()) + let costs = resolve_atomic_costs(&table, AggregationType::HydraKLL, &HashMap::new(), 2) .expect("unbenchmarked families still get a usable (stub) cost"); + assert_eq!(costs, AtomicCosts::default()); + } + + #[test] + fn trivial_accumulators_get_analytical_per_group_memory() { + let table: AtomicCostTable = vec![]; + let mem = |agg_type, n_labels| { + resolve_atomic_costs(&table, agg_type, &HashMap::new(), n_labels) + .expect("trivial accumulators always resolve") + .mem_bytes_per_instance + }; + // Two 4-byte label codes + one f64, with 8/7 hash-table slack. + assert!((mem(AggregationType::Sum, 2) - 16.0 * 8.0 / 7.0).abs() < 1e-9); + // Per-group and keyed-map forms store the same entry per group. assert_eq!( - costs.mem_bytes_per_instance, - AtomicCosts::default().mem_bytes_per_instance + mem(AggregationType::Sum, 2), + mem(AggregationType::MultipleSum, 2) + ); + assert_eq!( + mem(AggregationType::MinMax, 3), + mem(AggregationType::MultipleMinMax, 3) + ); + assert_eq!( + mem(AggregationType::Increase, 1), + mem(AggregationType::MultipleIncrease, 1) ); + assert!(mem(AggregationType::Increase, 1) > mem(AggregationType::Sum, 1)); } #[test] @@ -654,7 +710,7 @@ mod tests { let table: AtomicCostTable = vec![]; let params = cms_heap_params(3, 1024, 40); assert!( - resolve_atomic_costs(&table, AggregationType::CountMinSketchWithHeap, ¶ms) + resolve_atomic_costs(&table, AggregationType::CountMinSketchWithHeap, ¶ms, 0) .is_none() ); } @@ -687,8 +743,9 @@ mod tests { fn public_resolver_dispatches_cms_with_heap_to_the_reference_model() { let table = vec![cms_heap_entry(3, 1024)]; let params = cms_heap_params(3, 1024, 32); - let costs = resolve_atomic_costs(&table, AggregationType::CountMinSketchWithHeap, ¶ms) - .expect("public resolver must dispatch CMS-with-heap candidates"); + let costs = + resolve_atomic_costs(&table, AggregationType::CountMinSketchWithHeap, ¶ms, 0) + .expect("public resolver must dispatch CMS-with-heap candidates"); assert_eq!( costs.mem_bytes_per_instance, @@ -722,7 +779,7 @@ mod tests { // different from this case). let table: AtomicCostTable = vec![]; let params = HashMap::from([("width".to_string(), Value::from(1024u64))]); - resolve_atomic_costs(&table, AggregationType::CountMinSketch, ¶ms); + resolve_atomic_costs(&table, AggregationType::CountMinSketch, ¶ms, 0); } #[test] @@ -737,7 +794,7 @@ mod tests { query_accuracy: BTreeMap::new(), }]; let hll_params = HashMap::from([("precision".to_string(), Value::from(14u64))]); - assert!(resolve_atomic_costs(&hll_table, AggregationType::HLL, &hll_params).is_some()); + assert!(resolve_atomic_costs(&hll_table, AggregationType::HLL, &hll_params, 0).is_some()); let kll_table = vec![AtomicCostEntry { sketch: "kll-percall".into(), @@ -750,7 +807,7 @@ mod tests { }]; let kll_params = HashMap::from([("K".to_string(), Value::from(200u64))]); assert!( - resolve_atomic_costs(&kll_table, AggregationType::DatasketchesKLL, &kll_params) + resolve_atomic_costs(&kll_table, AggregationType::DatasketchesKLL, &kll_params, 0) .is_some() ); } diff --git a/asap-planner-rs/src/optimizer/candidate_gen.rs b/asap-planner-rs/src/optimizer/candidate_gen.rs index 807213ac..4ce1fac2 100644 --- a/asap-planner-rs/src/optimizer/candidate_gen.rs +++ b/asap-planner-rs/src/optimizer/candidate_gen.rs @@ -10,8 +10,22 @@ use serde_json::Value; use super::constants::{ CMS_DEPTHS, CMS_HEAP_SIZES, CMS_WIDTHS, HLL_PRECISIONS, HYDRA_COLS, HYDRA_K, HYDRA_ROWS, KLL_KS, }; +use super::label_set_facts::ItemFacts; use super::sketch_properties::sketch_properties; use super::solution::{QueryMethod, AQE}; +use crate::planner::agg_config::needs_key_aggregation; +use crate::planner::labels::set_subpopulation_labels; + +/// Compatible types the optimizer never proposes. Single-group +/// Sum/MinMax/Increase cost the same as their Multiple* twins, so only the +/// Multiple* forms are offered. Filtered here rather than in +/// `compatible_agg_types`, which the engine also uses to match queries to +/// deployed configs. +const OPTIMIZER_SKIPPED_AGG_TYPES: &[AggregationType] = &[ + AggregationType::Sum, + AggregationType::MinMax, + AggregationType::Increase, +]; /// A candidate streaming config for a single AQE, ready for cost evaluation. #[derive(Debug, Clone)] @@ -22,9 +36,14 @@ pub struct CandidateConfig { pub query_method: QueryMethod, /// Number of retained windows used at query time (n for Merge, 1 for Direct/Subtract, 0 for Exact). pub n_windows: u64, - /// Number of distinct label groups represented by this candidate. - /// Subpopulation-aware sketches ignore this value during cost evaluation. - pub label_group_count: u64, + /// Sketch instances the engine creates: one per distinct value of the + /// config's grouping labels (1 when they are empty). + pub instance_count: u64, + /// Distinct value combinations of the AQE's output labels. + pub output_group_count: u64, + /// Paired key aggregation (DeltaSetAggregator) deployed alongside + /// `config` when the value sketch can't list its own keys. + pub key_config: Option, } /// Enumerate all candidate configs for an AQE. @@ -35,24 +54,29 @@ pub struct CandidateConfig { /// Multi-statistic AQEs (e.g. avg = [Sum, Count]) return only EXACT — a single /// sketch family cannot serve two incompatible statistics simultaneously. pub fn enumerate_candidates(aqe: &AQE, scrape_interval_ms: u64) -> Vec { - enumerate_candidates_with_label_group_count(aqe, scrape_interval_ms, 1) + let one_group = ItemFacts { + output_group_count: 1, + topk_by_group_count: aqe.requirements.topk_by_labels.as_ref().map(|_| 1), + arrival_rate_per_sec: 1.0, + }; + enumerate_candidates_with_facts(aqe, scrape_interval_ms, &one_group) } -/// Enumerate candidates with a dataset-derived label-group count. -pub fn enumerate_candidates_with_label_group_count( +/// Enumerate candidates, stamping group counts from the AQE's label-set facts. +pub fn enumerate_candidates_with_facts( aqe: &AQE, scrape_interval_ms: u64, - label_group_count: u64, + facts: &ItemFacts, ) -> Vec { assert!( - label_group_count > 0, - "label_group_count must be greater than zero" + facts.output_group_count > 0, + "output_group_count must be greater than zero" ); let mut candidates = Vec::new(); if aqe.requirements.statistics.len() != 1 { // ponytail: multi-stat AQEs (avg) need two sketches; not supported in v1. - candidates.push(exact_candidate(label_group_count)); + candidates.push(exact_candidate(facts)); return candidates; } @@ -60,6 +84,9 @@ pub fn enumerate_candidates_with_label_group_count( let range_a_ms = aqe.requirements.data_range_ms; for &agg_type in compatible_agg_types(stat) { + if OPTIMIZER_SKIPPED_AGG_TYPES.contains(&agg_type) { + continue; + } let props = sketch_properties(agg_type); // CountMinSketchWithHeap's SUM/COUNT weighting lives in aggregation_sub_type @@ -95,6 +122,7 @@ pub fn enumerate_candidates_with_label_group_count( let config = build_config( aqe, + stat, agg_type, sub_type, ¶ms, @@ -104,26 +132,72 @@ pub fn enumerate_candidates_with_label_group_count( n, ); candidates.push(CandidateConfig { + instance_count: instance_count(&config, aqe, facts), + key_config: needs_key_aggregation(agg_type) + .then(|| build_key_config(&config, range_a_ms)), config: Some(config), query_method: qm, n_windows: n, - label_group_count, + output_group_count: facts.output_group_count, }); } } } } - candidates.push(exact_candidate(label_group_count)); + candidates.push(exact_candidate(facts)); candidates } -fn exact_candidate(label_group_count: u64) -> CandidateConfig { +fn exact_candidate(facts: &ItemFacts) -> CandidateConfig { CandidateConfig { config: None, query_method: QueryMethod::Exact, n_windows: 0, - label_group_count, + instance_count: 0, + output_group_count: facts.output_group_count, + key_config: None, + } +} + +/// The DeltaSetAggregator paired with `value`, as the legacy planner builds +/// it: same labels, Tumbling at the value's slide (DeltaSet is only correct +/// for non-overlapping windows), retaining enough panes to cover the query +/// range. +fn build_key_config(value: &AggregationConfig, range_a_ms: u64) -> AggregationConfig { + let pane_ms = value.slide_interval_ms; + AggregationConfig::new( + 0, // placeholder; overwritten by OptimizerSolution::register_config when deployed + AggregationType::DeltaSetAggregator, + String::new(), + HashMap::new(), + value.grouping_labels.clone(), + value.aggregated_labels.clone(), + KeyByLabelNames::empty(), // rollup_labels + String::new(), // original_yaml + pane_ms, + pane_ms, + WindowType::Tumbling, + value.spatial_filter.clone(), + value.metric.clone(), + Some(range_a_ms.div_ceil(pane_ms)), + None, // read_count_threshold + None, // table_name (SQL only) + None, // value_column (SQL only) + ) +} + +/// Instances the engine creates for `config`: its grouping labels are empty, +/// the AQE's output labels, or (for `topk by`) the bucketing labels. +fn instance_count(config: &AggregationConfig, aqe: &AQE, facts: &ItemFacts) -> u64 { + if config.grouping_labels.is_empty() { + 1 + } else if config.grouping_labels == aqe.requirements.grouping_labels { + facts.output_group_count + } else { + facts + .topk_by_group_count + .expect("grouping labels other than the output labels come from `topk by`") } } @@ -219,6 +293,7 @@ fn determine_query_method( #[allow(clippy::too_many_arguments)] fn build_config( aqe: &AQE, + stat: Statistic, agg_type: AggregationType, sub_type: &str, params: &HashMap, @@ -227,13 +302,28 @@ fn build_config( slide_interval: u64, n_windows: u64, ) -> AggregationConfig { + // Same grouping/aggregated split as the legacy planner: keyed types hold + // the output labels as keys inside one instance; others get one instance + // per group; `topk by` gets one heap per bucket. + let mut grouping = KeyByLabelNames::empty(); + let mut aggregated = KeyByLabelNames::empty(); + set_subpopulation_labels( + stat, + agg_type, + &aqe.requirements.grouping_labels, + aqe.requirements.topk_by_labels.as_ref(), + &mut KeyByLabelNames::empty(), + &mut grouping, + &mut aggregated, + ); + AggregationConfig::new( 0, // placeholder; overwritten by OptimizerSolution::register_config when deployed agg_type, sub_type.to_string(), params.clone(), - aqe.requirements.grouping_labels.clone(), - KeyByLabelNames::empty(), // aggregated_labels (not needed for optimizer feasibility) + grouping, + aggregated, KeyByLabelNames::empty(), // rollup_labels String::new(), // original_yaml w, @@ -364,6 +454,38 @@ mod tests { .any(|c| c.config.is_none() && c.query_method == QueryMethod::Exact)); } + #[test] + fn single_group_trivial_accumulators_are_not_proposed() { + for (stat, skipped, kept) in [ + ( + Statistic::Sum, + AggregationType::Sum, + AggregationType::MultipleSum, + ), + ( + Statistic::Min, + AggregationType::MinMax, + AggregationType::MultipleMinMax, + ), + ( + Statistic::Increase, + AggregationType::Increase, + AggregationType::MultipleIncrease, + ), + ] { + let types: Vec = + enumerate_candidates(&make_aqe(stat, 300_000, 60_000), 15_000) + .into_iter() + .filter_map(|c| c.config.map(|cfg| cfg.aggregation_type)) + .collect(); + assert!( + !types.contains(&skipped), + "{skipped:?} must not be proposed" + ); + assert!(types.contains(&kept), "{kept:?} must still be proposed"); + } + } + #[test] fn multiple_sum_candidates_get_a_non_empty_sub_type() { // MultipleSum's factory now rejects an empty aggregation_sub_type (#503) -- @@ -385,15 +507,92 @@ mod tests { } } + fn labels(names: &[&str]) -> KeyByLabelNames { + KeyByLabelNames::new(names.iter().map(|n| n.to_string()).collect()) + } + + fn facts(output: u64, topk_by: Option) -> ItemFacts { + ItemFacts { + output_group_count: output, + topk_by_group_count: topk_by, + arrival_rate_per_sec: 1.0, + } + } + + /// Before the optimizer reused the legacy planner's label split, every + /// config put the output labels in `grouping_labels`, so the engine built + /// one keyed map/sketch per group (and one top-k heap per series). #[test] - fn stamps_dataset_label_group_count_on_every_candidate() { - let aqe = make_aqe(Statistic::Sum, 300_000, 60_000); - let candidates = enumerate_candidates_with_label_group_count(&aqe, 15_000, 7); + fn keyed_types_hold_groups_inside_one_instance_and_per_group_types_do_not() { + for stat in [Statistic::Sum, Statistic::Quantile] { + let mut aqe = make_aqe(stat, 300_000, 60_000); + aqe.requirements.grouping_labels = labels(&["svc"]); + let candidates = enumerate_candidates_with_facts(&aqe, 15_000, &facts(7, None)); + + for c in &candidates { + assert_eq!(c.output_group_count, 7); + let Some(cfg) = &c.config else { continue }; + match cfg.aggregation_type { + AggregationType::MultipleSum + | AggregationType::CountMinSketch + | AggregationType::HydraKLL => { + assert!(cfg.grouping_labels.is_empty(), "{cfg:?}"); + assert_eq!(cfg.aggregated_labels, labels(&["svc"])); + assert_eq!(c.instance_count, 1); + } + AggregationType::DatasketchesKLL => { + assert_eq!(cfg.grouping_labels, labels(&["svc"])); + assert!(cfg.aggregated_labels.is_empty()); + assert_eq!(c.instance_count, 7); + } + other => panic!("unexpected candidate type {other:?}"), + } + } + } + } - assert!(!candidates.is_empty()); - assert!(candidates + #[test] + fn only_cms_and_hydra_get_a_paired_tumbling_delta_set() { + for stat in [Statistic::Sum, Statistic::Quantile, Statistic::Topk] { + let mut aqe = make_aqe(stat, 600_000, 30_000); + aqe.requirements.grouping_labels = labels(&["svc"]); + for c in enumerate_candidates(&aqe, 30_000) { + let Some(cfg) = &c.config else { continue }; + match (&c.key_config, needs_key_aggregation(cfg.aggregation_type)) { + (Some(key), true) => { + assert_eq!(key.aggregation_type, AggregationType::DeltaSetAggregator); + assert_eq!(key.window_type, WindowType::Tumbling); + assert_eq!(key.window_size_ms, cfg.slide_interval_ms); + assert_eq!(key.grouping_labels, cfg.grouping_labels); + assert_eq!(key.aggregated_labels, cfg.aggregated_labels); + assert_eq!( + key.num_aggregates_to_retain, + Some(600_000 / cfg.slide_interval_ms) + ); + } + (None, false) => {} + (key, _) => panic!("{:?} has key config {key:?}", cfg.aggregation_type), + } + } + } + } + + #[test] + fn topk_by_gets_one_heap_per_bucket() { + let mut aqe = make_aqe(Statistic::Topk, 60_000, 60_000); + aqe.requirements.grouping_labels = labels(&["endpoint", "svc"]); + aqe.requirements.topk_by_labels = Some(labels(&["svc"])); + let candidates = enumerate_candidates_with_facts(&aqe, 15_000, &facts(100, Some(3))); + + let heap = candidates .iter() - .all(|candidate| candidate.label_group_count == 7)); + .find(|c| c.config.is_some()) + .expect("a CMS-with-heap candidate"); + let cfg = heap.config.as_ref().unwrap(); + assert_eq!(cfg.grouping_labels, labels(&["svc"])); + assert_eq!(cfg.aggregated_labels, labels(&["endpoint"])); + assert_eq!(heap.instance_count, 3); + assert_eq!(heap.output_group_count, 100); } #[test] diff --git a/asap-planner-rs/src/optimizer/constants.rs b/asap-planner-rs/src/optimizer/constants.rs index 5ca899d0..02c3ab71 100644 --- a/asap-planner-rs/src/optimizer/constants.rs +++ b/asap-planner-rs/src/optimizer/constants.rs @@ -41,8 +41,18 @@ pub const INGEST_CPU_WEIGHT: f64 = 1.0; pub const QUERY_MEM_WEIGHT: f64 = 1e-9; pub const QUERY_CPU_WEIGHT: f64 = 1.0; -/// Subpopulation count: 1 if subpopulation_aware else the distinct label-group -/// count for this config. The label-group count isn't profiled yet (needs -/// Prometheus series-count data) — use 1 as a placeholder; both branches -/// collapse to the same value until that lands. -pub const SUBPOPULATION_COUNT: f64 = 1.0; +// Analytical per-group memory for trivial accumulators (Sum/MinMax/Increase and +// their Multiple* maps): one dictionary code per grouping-label value plus the +// value, inflated by hash-table slack. The dictionary itself is amortized over +// time and not charged. +pub const LABEL_VALUE_CODE_BYTES: f64 = 4.0; +/// hashbrown's maximum load factor is 7/8. +pub const HASH_TABLE_SLACK: f64 = 8.0 / 7.0; +/// One f64. +pub const SUM_VALUE_BYTES: f64 = 8.0; +/// One f64. +pub const MIN_MAX_VALUE_BYTES: f64 = 8.0; +/// Start/last measurement and timestamp, sample count, reset adjustment, and +/// an empty reset-event Vec. +// ponytail: ignores counter-reset events; add per-event bytes if resets are frequent. +pub const INCREASE_VALUE_BYTES: f64 = 72.0; diff --git a/asap-planner-rs/src/optimizer/cost_model.rs b/asap-planner-rs/src/optimizer/cost_model.rs index b9fb5e75..aca3387e 100644 --- a/asap-planner-rs/src/optimizer/cost_model.rs +++ b/asap-planner-rs/src/optimizer/cost_model.rs @@ -1,11 +1,11 @@ use asap_types::enums::WindowType; -use promql_utilities::query_logics::enums::AggregationType; +use promql_utilities::query_logics::enums::{AggregationType, Statistic}; use super::candidate_gen::CandidateConfig; use super::constants::{ - EXACT_QUERY_CPU_SECS, INGEST_CPU_WEIGHT, INGEST_MEM_WEIGHT, INSERT_CPU_SECS, - MEM_BYTES_PER_INSTANCE, MERGE_CPU_SECS, QUERY_CPU_SECS, QUERY_CPU_WEIGHT, QUERY_MEM_WEIGHT, - SUBPOPULATION_COUNT, SUBTRACT_CPU_SECS, + EXACT_QUERY_CPU_SECS, HASH_TABLE_SLACK, INGEST_CPU_WEIGHT, INGEST_MEM_WEIGHT, INSERT_CPU_SECS, + LABEL_VALUE_CODE_BYTES, MEM_BYTES_PER_INSTANCE, MERGE_CPU_SECS, QUERY_CPU_SECS, + QUERY_CPU_WEIGHT, QUERY_MEM_WEIGHT, SUBTRACT_CPU_SECS, }; use super::sketch_properties::sketch_properties; use super::solution::{QueryMethod, AQE}; @@ -14,6 +14,7 @@ use super::solution::{QueryMethod, AQE}; /// values come from sketch-bench in Phase 3 (see implementation plan, 3c). #[derive(Debug, Clone, Copy, PartialEq)] pub struct AtomicCosts { + /// One instance; for a keyed-unbounded map, one group's entry. pub mem_bytes_per_instance: f64, pub insert_cpu_secs: f64, pub merge_cpu_secs: f64, @@ -76,7 +77,7 @@ pub fn ingest_cost( return 0.0; // EXACT: no streaming config deployed. }; - let subpopulation_count = effective_subpopulation_count(candidate, agg_config.aggregation_type); + let units = stored_units(candidate, agg_config.aggregation_type); // Defensive floor: slide_interval_ms is a plain u64 on a widely-shared struct; // guard against div-by-zero producing `inf` and poisoning cost comparisons. @@ -87,19 +88,38 @@ pub fn ingest_cost( } }; - let mem_active = n_concurrent * subpopulation_count * costs.mem_bytes_per_instance; + let mem_active = n_concurrent * units * costs.mem_bytes_per_instance; let cpu_ingest = match agg_config.window_type { WindowType::Tumbling => arrival_rate_hz * costs.insert_cpu_secs, WindowType::Sliding => arrival_rate_hz * n_concurrent * costs.insert_cpu_secs, }; - weights.ingest_mem * mem_active + weights.ingest_cpu * cpu_ingest + weights.ingest_mem * mem_active + + weights.ingest_cpu * cpu_ingest + + key_tracker_ingest_cost(candidate, arrival_rate_hz, weights) +} + +/// Ingest cost of the paired key aggregation, if any: one key entry per +/// output group in its single tumbling pane, one insert per item. +// ponytail: stub insert CPU and analytical key bytes; use sketch-bench numbers once DeltaSet is measured. +fn key_tracker_ingest_cost( + candidate: &CandidateConfig, + arrival_rate_hz: f64, + weights: &CostWeights, +) -> f64 { + let Some(key_config) = &candidate.key_config else { + return 0.0; + }; + let n_labels = key_config.grouping_labels.len() + key_config.aggregated_labels.len(); + let entry_bytes = n_labels as f64 * LABEL_VALUE_CODE_BYTES * HASH_TABLE_SLACK; + weights.ingest_mem * candidate.output_group_count as f64 * entry_bytes + + weights.ingest_cpu * arrival_rate_hz * INSERT_CPU_SECS } /// QueryCost(a,g): cost of answering one query for `aqe` from `candidate`. pub fn query_cost( - _aqe: &AQE, + aqe: &AQE, candidate: &CandidateConfig, costs: &AtomicCosts, weights: &CostWeights, @@ -108,29 +128,25 @@ pub fn query_cost( return costs.exact_query_cpu_secs * weights.query_cpu; // EXACT: raw query at query time. }; - // Subpopulation count; see ingest_cost comment. - let subpopulation_count = effective_subpopulation_count(candidate, agg_config.aggregation_type); + let units = stored_units(candidate, agg_config.aggregation_type); let props = sketch_properties(agg_config.aggregation_type); + let read_cpu = reads_per_query(aqe, candidate) * costs.query_cpu_secs; let (cpu, mem) = match &candidate.query_method { - QueryMethod::Direct => ( - subpopulation_count * costs.query_cpu_secs, - subpopulation_count * costs.mem_bytes_per_instance, - ), + QueryMethod::Direct => (read_cpu, units * costs.mem_bytes_per_instance), QueryMethod::Merge { num_windows } => { debug_assert!(props.mergeable); let merges = (*num_windows).saturating_sub(1) as f64; ( - subpopulation_count * (merges * costs.merge_cpu_secs + costs.query_cpu_secs), - *num_windows as f64 * subpopulation_count * costs.mem_bytes_per_instance, + units * merges * costs.merge_cpu_secs + read_cpu, + *num_windows as f64 * units * costs.mem_bytes_per_instance, ) } QueryMethod::Subtract => { debug_assert!(props.subtractable); ( - subpopulation_count - * (costs.merge_cpu_secs + costs.subtract_cpu_secs + costs.query_cpu_secs), - 2.0 * subpopulation_count * costs.mem_bytes_per_instance, + units * (costs.merge_cpu_secs + costs.subtract_cpu_secs) + read_cpu, + 2.0 * units * costs.mem_bytes_per_instance, ) } // candidate_gen only ever pairs Exact with config=None, already handled above. @@ -142,18 +158,28 @@ pub fn query_cost( weights.query_cpu * cpu + weights.query_mem * mem } -fn effective_subpopulation_count( - candidate: &CandidateConfig, - aggregation_type: AggregationType, -) -> f64 { - if sketch_properties(aggregation_type).subpopulation_aware { - SUBPOPULATION_COUNT +/// Units of `mem_bytes_per_instance` held per window; also scales +/// whole-structure merge/subtract work. A keyed map stores one entry per +/// output group; anything else stores fixed-size instances. +fn stored_units(candidate: &CandidateConfig, aggregation_type: AggregationType) -> f64 { + assert!( + candidate.instance_count > 0 && candidate.output_group_count > 0, + "candidates require positive group counts" + ); + if sketch_properties(aggregation_type).memory_grows_with_keys { + candidate.output_group_count as f64 } else { - assert!( - candidate.label_group_count > 0, - "non-subpopulation-aware candidates require a positive label_group_count" - ); - candidate.label_group_count as f64 + candidate.instance_count as f64 + } +} + +/// `query_cpu_secs` operations per query: one per output group, except +/// top-k, which reads each heap once (its query cost already covers the heap). +fn reads_per_query(aqe: &AQE, candidate: &CandidateConfig) -> f64 { + if aqe.requirements.statistics == [Statistic::Topk] { + candidate.instance_count as f64 + } else { + candidate.output_group_count as f64 } } @@ -174,10 +200,10 @@ pub fn total_cost_rate( #[cfg(test)] mod tests { use super::*; - use crate::optimizer::candidate_gen::enumerate_candidates; + use crate::optimizer::candidate_gen::{enumerate_candidates, enumerate_candidates_with_facts}; + use crate::optimizer::label_set_facts::ItemFacts; use asap_types::query_requirements::QueryRequirements; use promql_utilities::data_model::KeyByLabelNames; - use promql_utilities::query_logics::enums::Statistic; fn make_aqe(stat: Statistic, range_ms: u64, min_t: u64) -> AQE { AQE { @@ -203,7 +229,9 @@ mod tests { config: None, query_method: QueryMethod::Exact, n_windows: 0, - label_group_count: 1, + instance_count: 1, + output_group_count: 1, + key_config: None, }; let costs = AtomicCosts::default(); let weights = CostWeights::default(); @@ -229,13 +257,17 @@ mod tests { config: Some(template.clone()), query_method: QueryMethod::Merge { num_windows: 2 }, n_windows: 2, - label_group_count: 1, + instance_count: 1, + output_group_count: 1, + key_config: None, }; let c5 = CandidateConfig { config: Some(template), query_method: QueryMethod::Merge { num_windows: 5 }, n_windows: 5, - label_group_count: 1, + instance_count: 1, + output_group_count: 1, + key_config: None, }; assert_eq!( @@ -262,13 +294,17 @@ mod tests { config: Some(template.clone()), query_method: QueryMethod::Merge { num_windows: 5 }, n_windows: 5, - label_group_count: 1, + instance_count: 1, + output_group_count: 1, + key_config: None, }; let subtract = CandidateConfig { config: Some(template), query_method: QueryMethod::Subtract, n_windows: 5, - label_group_count: 1, + instance_count: 1, + output_group_count: 1, + key_config: None, }; assert!( @@ -276,72 +312,115 @@ mod tests { ); } - #[test] - fn non_subpopulation_aware_cost_scales_with_label_group_count() { - let a = make_aqe(Statistic::Sum, 300_000, 300_000); - let candidate = enumerate_candidates(&a, 60_000) + /// The Direct-method `agg_type` candidate for a `stat` AQE grouped by + /// `svc` (and bucketed by it for top-k), with `groups` output groups. + fn direct_candidate( + stat: Statistic, + agg_type: AggregationType, + groups: u64, + ) -> (AQE, CandidateConfig) { + let mut a = make_aqe(stat, 300_000, 300_000); + a.requirements.grouping_labels = KeyByLabelNames::new(vec!["svc".into()]); + let facts = ItemFacts { + output_group_count: groups, + topk_by_group_count: None, + arrival_rate_per_sec: 1.0, + }; + let candidate = enumerate_candidates_with_facts(&a, 60_000, &facts) .into_iter() - .find(|candidate| { - candidate.config.as_ref().is_some_and(|config| { - !sketch_properties(config.aggregation_type).subpopulation_aware - }) + .find(|c| { + c.query_method == QueryMethod::Direct + && c.config + .as_ref() + .is_some_and(|cfg| cfg.aggregation_type == agg_type) }) - .expect("expected a non-subpopulation-aware candidate"); - let one_group = CandidateConfig { - label_group_count: 1, - ..candidate.clone() - }; - let five_groups = CandidateConfig { - label_group_count: 5, - ..candidate - }; + .unwrap_or_else(|| panic!("expected a Direct {agg_type:?} candidate")); + (a, candidate) + } + + const MEM_ONLY: CostWeights = CostWeights { + ingest_mem: 1.0, + ingest_cpu: 0.0, + query_mem: 1.0, + query_cpu: 0.0, + }; + const QUERY_CPU_ONLY: CostWeights = CostWeights { + ingest_mem: 0.0, + ingest_cpu: 0.0, + query_mem: 0.0, + query_cpu: 1.0, + }; + + #[test] + fn per_group_and_keyed_map_costs_scale_with_output_groups() { + // KLL: one instance per group. MultipleSum: one map, one entry per group. let costs = AtomicCosts::default(); - let weights = CostWeights { - ingest_mem: 1.0, - ingest_cpu: 0.0, - query_mem: 1.0, - query_cpu: 0.0, - }; + for (stat, agg_type) in [ + (Statistic::Quantile, AggregationType::DatasketchesKLL), + (Statistic::Sum, AggregationType::MultipleSum), + ] { + let (a, one) = direct_candidate(stat, agg_type, 1); + let (_, five) = direct_candidate(stat, agg_type, 5); + assert_eq!( + ingest_cost(&five, 1.0, &costs, &MEM_ONLY), + 5.0 * ingest_cost(&one, 1.0, &costs, &MEM_ONLY), + "{agg_type:?} ingest memory" + ); + for weights in [MEM_ONLY, QUERY_CPU_ONLY] { + assert_eq!( + query_cost(&a, &five, &costs, &weights), + 5.0 * query_cost(&a, &one, &costs, &weights), + "{agg_type:?} query cost" + ); + } + } + } + #[test] + fn fixed_size_keyed_sketch_memory_ignores_groups_but_query_reads_each() { + let costs = AtomicCosts::default(); + let (a, one) = direct_candidate(Statistic::Sum, AggregationType::CountMinSketch, 1); + let (_, five) = direct_candidate(Statistic::Sum, AggregationType::CountMinSketch, 5); + // Only the paired key tracker grows: one 1-label entry per extra group. + let key_entry_bytes = LABEL_VALUE_CODE_BYTES * HASH_TABLE_SLACK; + assert!( + (ingest_cost(&five, 1.0, &costs, &MEM_ONLY) + - ingest_cost(&one, 1.0, &costs, &MEM_ONLY) + - 4.0 * key_entry_bytes) + .abs() + < 1e-9 + ); assert_eq!( - ingest_cost(&five_groups, 1.0, &costs, &weights), - 5.0 * ingest_cost(&one_group, 1.0, &costs, &weights) + query_cost(&a, &one, &costs, &MEM_ONLY), + query_cost(&a, &five, &costs, &MEM_ONLY) ); assert_eq!( - query_cost(&a, &five_groups, &costs, &weights), - 5.0 * query_cost(&a, &one_group, &costs, &weights) + query_cost(&a, &five, &costs, &QUERY_CPU_ONLY), + 5.0 * query_cost(&a, &one, &costs, &QUERY_CPU_ONLY) ); } #[test] - fn subpopulation_aware_cost_ignores_label_group_count() { - let a = make_aqe(Statistic::Sum, 300_000, 300_000); - let candidate = enumerate_candidates(&a, 60_000) - .into_iter() - .find(|candidate| { - candidate.config.as_ref().is_some_and(|config| { - sketch_properties(config.aggregation_type).subpopulation_aware - }) - }) - .expect("expected a subpopulation-aware candidate"); - let one_group = CandidateConfig { - label_group_count: 1, - ..candidate.clone() - }; - let five_groups = CandidateConfig { - label_group_count: 5, - ..candidate - }; + fn topk_reads_each_heap_once_regardless_of_output_groups() { let costs = AtomicCosts::default(); - let weights = CostWeights::default(); - - assert_eq!( - ingest_cost(&one_group, 1.0, &costs, &weights), - ingest_cost(&five_groups, 1.0, &costs, &weights) - ); + let mut a = make_aqe(Statistic::Topk, 60_000, 60_000); + a.requirements.grouping_labels = KeyByLabelNames::new(vec!["svc".into()]); + let heap = |groups| { + let facts = ItemFacts { + output_group_count: groups, + topk_by_group_count: None, + arrival_rate_per_sec: 1.0, + }; + enumerate_candidates_with_facts(&a, 15_000, &facts) + .into_iter() + .find(|c| c.config.is_some()) + .expect("a CMS-with-heap candidate") + }; + let (few, many) = (heap(1), heap(1000)); + assert_eq!(many.instance_count, 1); assert_eq!( - query_cost(&a, &one_group, &costs, &weights), - query_cost(&a, &five_groups, &costs, &weights) + query_cost(&a, &few, &costs, &QUERY_CPU_ONLY), + query_cost(&a, &many, &costs, &QUERY_CPU_ONLY) ); } } diff --git a/asap-planner-rs/src/optimizer/greedy.rs b/asap-planner-rs/src/optimizer/greedy.rs index df3ba960..2c001f58 100644 --- a/asap-planner-rs/src/optimizer/greedy.rs +++ b/asap-planner-rs/src/optimizer/greedy.rs @@ -1,11 +1,9 @@ -use std::collections::HashMap; - use tracing::debug; use super::atomic_costs::{resolve_atomic_costs, AtomicCostTable}; -use super::candidate_gen::enumerate_candidates_with_label_group_count; +use super::candidate_gen::enumerate_candidates_with_facts; use super::cost_model::{ingest_cost, query_cost, total_cost_rate, AtomicCosts, CostWeights}; -use super::label_set_facts::{ItemFacts, LabelSetKey}; +use super::label_set_facts::ItemFacts; use super::solution::{AQEAssignment, OptimizerSolution, AQE}; /// Greedily assign each AQE to its independently-cheapest candidate config. @@ -14,8 +12,8 @@ use super::solution::{AQEAssignment, OptimizerSolution, AQE}; /// two AQEs could share one. The Phase 3 MIP finds sharing opportunities; this /// is the v1 baseline. /// -/// `facts` supplies each AQE's group cardinality and arrival rate, keyed by -/// its label set; every AQE must have an entry. +/// `facts` supplies each AQE's group counts and arrival rate, index-aligned +/// with `aqes`. /// /// Each candidate is costed at its own `(sketch_type, params)` via /// `atomic_cost_table` (see ASAPQuery#524) rather than one cost applied to @@ -26,21 +24,18 @@ pub fn greedy_assign( scrape_interval_ms: u64, atomic_cost_table: &AtomicCostTable, weights: &CostWeights, - facts: &HashMap, + facts: &[ItemFacts], ) -> OptimizerSolution { + assert_eq!( + aqes.len(), + facts.len(), + "facts must be index-aligned with aqes" + ); let mut solution = OptimizerSolution::empty(); - for aqe in aqes { - let key = LabelSetKey::from_requirements(&aqe.requirements); - let item_facts = *facts - .get(&key) - .unwrap_or_else(|| panic!("missing label-set facts for {key}")); + for (aqe, item_facts) in aqes.into_iter().zip(facts) { let arrival_rate_hz = item_facts.arrival_rate_per_sec; - let candidates = enumerate_candidates_with_label_group_count( - &aqe, - scrape_interval_ms, - item_facts.cardinality, - ); + let candidates = enumerate_candidates_with_facts(&aqe, scrape_interval_ms, item_facts); let (best, costs) = candidates .into_iter() @@ -53,6 +48,7 @@ pub fn greedy_assign( atomic_cost_table, cfg.aggregation_type, &cfg.parameters, + aqe.requirements.grouping_labels.len(), )?, }; let cost = total_cost_rate(&aqe, &c, arrival_rate_hz, &costs, weights); @@ -71,6 +67,9 @@ pub fn greedy_assign( let query_method = best.query_method.clone(); let aggregation_id = best.config.map(|config| solution.register_config(config)); + let key_aggregation_id = best + .key_config + .map(|config| solution.register_config(config)); debug!( metric = %aqe.requirements.metric, @@ -87,6 +86,7 @@ pub fn greedy_assign( solution.assignments.push(AQEAssignment { aqe, aggregation_id, + key_aggregation_id, query_method, estimated_query_cost_per_sec: query_rate, }); @@ -97,10 +97,9 @@ pub fn greedy_assign( #[cfg(test)] mod tests { - use std::collections::HashMap; - use super::*; use crate::optimizer::atomic_costs::AtomicCostEntry; + use crate::optimizer::label_set_facts::ItemFacts; use asap_types::query_requirements::QueryRequirements; use promql_utilities::data_model::KeyByLabelNames; use promql_utilities::query_logics::enums::{AggregationType, Statistic}; @@ -124,17 +123,13 @@ mod tests { } } - /// One group, one item/sec for every AQE's label set. - fn unit_facts(aqes: &[AQE]) -> HashMap { + /// One group, one item/sec for every AQE. + fn unit_facts(aqes: &[AQE]) -> Vec { aqes.iter() - .map(|aqe| { - ( - LabelSetKey::from_requirements(&aqe.requirements), - ItemFacts { - cardinality: 1, - arrival_rate_per_sec: 1.0, - }, - ) + .map(|aqe| ItemFacts { + output_group_count: 1, + topk_by_group_count: aqe.requirements.topk_by_labels.as_ref().map(|_| 1), + arrival_rate_per_sec: 1.0, }) .collect() } @@ -244,4 +239,60 @@ mod tests { AggregationType::CountMinSketchWithHeap ); } + + /// A grouped query served by CMS needs a key aggregation the engine can + /// list groups from, referenced alongside the value in its query config. + #[test] + fn cms_assignment_deploys_a_key_aggregation_the_engine_matches() { + let table = vec![AtomicCostEntry { + sketch: "cms-fastpath-vector2d".into(), + sketch_config: serde_json::json!({ + "algorithm": "cms-fastpath-vector2d", + "params": { "rows": 3, "cols": 512 } + }), + mem_bytes_per_instance: 1.0, + insert_cpu_secs: 0.0, + merge_cpu_secs: 0.0, + query_cpu_secs: 0.0, + query_accuracy: std::collections::BTreeMap::new(), + }]; + let mut aqe = make_aqe(Statistic::Sum, 60_000, 60_000, 1.0 / 60.0); + aqe.requirements.grouping_labels = KeyByLabelNames::new(vec!["svc".into()]); + let solution = greedy_assign( + vec![aqe.clone()], + 60_000, + &table, + &CostWeights::default(), + &unit_facts(std::slice::from_ref(&aqe)), + ); + + let assignment = &solution.assignments[0]; + let value_id = assignment.aggregation_id.expect("CMS deployed"); + let key_id = assignment.key_aggregation_id.expect("key tracker deployed"); + let deployed = solution.deployed_configs(); + assert_eq!( + deployed[&value_id].aggregation_type, + AggregationType::CountMinSketch + ); + assert_eq!( + deployed[&key_id].aggregation_type, + AggregationType::DeltaSetAggregator + ); + + let (_, inference) = crate::optimizer::translate(&solution); + let refs: Vec = inference.query_configs[0] + .aggregations + .iter() + .map(|r| r.aggregation_id) + .collect(); + assert_eq!(refs, vec![value_id, key_id]); + + // The engine's query-config path validates the pair with this check. + assert!( + asap_types::capability_matching::key_agg_compatible_with_value( + &deployed[&value_id], + &deployed[&key_id], + ) + ); + } } diff --git a/asap-planner-rs/src/optimizer/label_set_facts.rs b/asap-planner-rs/src/optimizer/label_set_facts.rs index 183cb26a..2892cc46 100644 --- a/asap-planner-rs/src/optimizer/label_set_facts.rs +++ b/asap-planner-rs/src/optimizer/label_set_facts.rs @@ -103,11 +103,13 @@ impl std::fmt::Display for LabelSetKey { } } -/// Facts for one item's label set, ready for costing. +/// Facts for one AQE, ready for costing. #[derive(Debug, Clone, Copy, PartialEq)] pub struct ItemFacts { - /// Distinct grouping-label value combinations. - pub cardinality: u64, + /// Distinct value combinations of the AQE's output (grouping) labels. + pub output_group_count: u64, + /// Distinct values of the `topk by` labels; `Some` iff the AQE has them. + pub topk_by_group_count: Option, /// Aggregate items/sec into the grouped stream, across all groups. pub arrival_rate_per_sec: f64, } @@ -208,9 +210,10 @@ impl LabelSetFacts { }) } - /// Look up facts for every AQE's label set. Errors list every missing key - /// at once; facts no AQE uses are only warned about, so one file can serve - /// several workloads. + /// Look up facts for every AQE, index-aligned with `aqes`. An AQE needs its + /// series row, its output label set's group row, and for `topk by (L)` a + /// group row for `L`. Errors list every missing key at once; facts no AQE + /// uses are only warned about, so one file can serve several workloads. /// /// Arrival rate assumes each series yields one sample per scrape: /// `series_count / scrape interval`. @@ -219,31 +222,48 @@ impl LabelSetFacts { &self, aqes: &[AQE], scrape_interval_ms: u64, - ) -> Result, LabelSetFactsError> { - let mut resolved = HashMap::new(); + ) -> Result, LabelSetFactsError> { + let mut used_series = HashSet::new(); + let mut used_groups = HashSet::new(); let mut missing = Vec::new(); + let mut lookup_group = |key: LabelSetKey, missing: &mut Vec| { + let found = self.cardinalities.get(&key).copied(); + if found.is_none() { + missing.push(format!("groups: {key}")); + } + used_groups.insert(key); + found + }; + + let mut resolved = Vec::with_capacity(aqes.len()); for aqe in aqes { let key = LabelSetKey::from_requirements(&aqe.requirements); - if resolved.contains_key(&key) { - continue; - } - let series_count = self.series_counts.get(&key.series_key()); - let cardinality = self.cardinalities.get(&key); + let series_key = key.series_key(); + let series_count = self.series_counts.get(&series_key).copied(); if series_count.is_none() { - missing.push(format!("series: {}", key.series_key())); - } - if cardinality.is_none() { - missing.push(format!("groups: {key}")); + missing.push(format!("series: {series_key}")); } - if let (Some(&series_count), Some(&cardinality)) = (series_count, cardinality) { - let arrival_rate_per_sec = series_count as f64 * 1000.0 / scrape_interval_ms as f64; - resolved.insert( - key, - ItemFacts { - cardinality, - arrival_rate_per_sec, + used_series.insert(series_key); + + let topk_by_group_count = aqe.requirements.topk_by_labels.as_ref().map(|labels| { + lookup_group( + LabelSetKey { + grouping_labels: labels.clone(), + ..key.clone() }, - ); + &mut missing, + ) + }); + let output_group_count = lookup_group(key, &mut missing); + + if let (Some(series_count), Some(output_group_count)) = + (series_count, output_group_count) + { + resolved.push(ItemFacts { + output_group_count, + topk_by_group_count: topk_by_group_count.flatten(), + arrival_rate_per_sec: series_count as f64 * 1000.0 / scrape_interval_ms as f64, + }); } } if !missing.is_empty() { @@ -252,15 +272,13 @@ impl LabelSetFacts { return Err(LabelSetFactsError::MissingFacts(missing)); } - let used_series: HashSet = - resolved.keys().map(LabelSetKey::series_key).collect(); for key in self.series_counts.keys() { if !used_series.contains(key) { tracing::warn!(%key, "series facts match no workload item"); } } for key in self.cardinalities.keys() { - if !resolved.contains_key(key) { + if !used_groups.contains(key) { tracing::warn!(%key, "group facts match no workload item"); } } @@ -317,12 +335,43 @@ groups: r#"job="api",env="prod""#, &["endpoint", "service"], ); - let resolved = facts.resolve(std::slice::from_ref(&item), 15_000).unwrap(); - let got = resolved[&LabelSetKey::from_requirements(&item.requirements)]; - assert_eq!(got.cardinality, 50); + let got = facts.resolve(&[item], 15_000).unwrap()[0]; + assert_eq!(got.output_group_count, 50); + assert_eq!(got.topk_by_group_count, None); assert!((got.arrival_rate_per_sec - 1000.0 / 15.0).abs() < 1e-9); } + #[test] + fn topk_by_items_also_resolve_the_by_label_set() { + let facts = LabelSetFacts::from_yaml(FACTS).unwrap(); + let mut item = aqe( + "http_requests_total", + r#"job="api",env="prod""#, + &["endpoint", "service"], + ); + item.requirements.topk_by_labels = Some(KeyByLabelNames::new(vec!["service".into()])); + + // No [service] group row yet: must be reported, not defaulted. + let err = facts + .resolve(std::slice::from_ref(&item), 15_000) + .unwrap_err(); + let LabelSetFactsError::MissingFacts(missing) = err else { + panic!("expected MissingFacts, got {err:?}"); + }; + assert_eq!(missing.len(), 1, "{missing:?}"); + assert!(missing[0].contains(r#"grouping_labels=["service"]"#)); + + let with_by = format!( + "{FACTS} - {{metric: http_requests_total, spatial_filter: 'job=\"api\",env=\"prod\"', grouping_labels: [service], cardinality: 5}}\n" + ); + let got = LabelSetFacts::from_yaml(&with_by) + .unwrap() + .resolve(&[item], 15_000) + .unwrap()[0]; + assert_eq!(got.output_group_count, 50); + assert_eq!(got.topk_by_group_count, Some(5)); + } + #[test] fn spatial_filter_matches_after_normalization() { // File writes the matchers in a different order and with braces. diff --git a/asap-planner-rs/src/optimizer/mod.rs b/asap-planner-rs/src/optimizer/mod.rs index e6d8692f..29927a09 100644 --- a/asap-planner-rs/src/optimizer/mod.rs +++ b/asap-planner-rs/src/optimizer/mod.rs @@ -16,9 +16,7 @@ pub use atomic_costs::{ load_selected_atomic_cost_table, resolve_atomic_costs, AtomicCostEntry, AtomicCostTable, ExternalWorkload, WorkloadDescription, }; -pub use candidate_gen::{ - enumerate_candidates, enumerate_candidates_with_label_group_count, CandidateConfig, -}; +pub use candidate_gen::{enumerate_candidates, enumerate_candidates_with_facts, CandidateConfig}; pub use cost_model::{ingest_cost, query_cost, total_cost_rate, AtomicCosts, CostWeights}; pub use greedy::greedy_assign; pub use label_set_facts::{ItemFacts, LabelSetFacts, LabelSetFactsError, LabelSetKey, SeriesKey}; diff --git a/asap-planner-rs/src/optimizer/pipeline.rs b/asap-planner-rs/src/optimizer/pipeline.rs index 8992ceae..6aa8b405 100644 --- a/asap-planner-rs/src/optimizer/pipeline.rs +++ b/asap-planner-rs/src/optimizer/pipeline.rs @@ -8,7 +8,7 @@ use super::aqe_extractor::{extract_aqes, RQE}; use super::atomic_costs::AtomicCostTable; use super::cost_model::CostWeights; use super::greedy::greedy_assign; -use super::label_set_facts::{LabelSetFacts, LabelSetFactsError}; +use super::label_set_facts::{LabelSetFacts, LabelSetFactsError, LabelSetKey}; use super::solution::{OptimizerSolution, AQE}; use super::translator::{translate, TranslationSummary}; @@ -103,10 +103,11 @@ pub fn run_greedy_pipeline( } let item_facts = facts.resolve(&aqes, scrape_interval_ms)?; - for (key, item) in &item_facts { + for (aqe, item) in aqes.iter().zip(&item_facts) { tracing::info!( - %key, - cardinality = item.cardinality, + key = %LabelSetKey::from_requirements(&aqe.requirements), + output_group_count = item.output_group_count, + topk_by_group_count = ?item.topk_by_group_count, arrival_rate_per_sec = item.arrival_rate_per_sec, "optimizer label-set facts" ); diff --git a/asap-planner-rs/src/optimizer/sketch_properties.rs b/asap-planner-rs/src/optimizer/sketch_properties.rs index 6a56b2e7..435d93a2 100644 --- a/asap-planner-rs/src/optimizer/sketch_properties.rs +++ b/asap-planner-rs/src/optimizer/sketch_properties.rs @@ -6,15 +6,17 @@ pub struct SketchProperties { pub mergeable: bool, /// Element-wise difference is defined; enables the Subtract query method (tumbling only). pub subtractable: bool, - /// One deployed instance handles multiple label-group keys simultaneously. - pub subpopulation_aware: bool, + /// One instance's memory grows with the keys it holds (a keyed map), as + /// opposed to a fixed-size instance. How many instances exist comes from + /// the config's grouping labels, not from here. + pub memory_grows_with_keys: bool, } pub fn sketch_properties(t: AggregationType) -> SketchProperties { - let p = |me, su, sp| SketchProperties { + let p = |me, su, grows| SketchProperties { mergeable: me, subtractable: su, - subpopulation_aware: sp, + memory_grows_with_keys: grows, }; match t { AggregationType::Sum => p(true, true, false), @@ -24,11 +26,11 @@ pub fn sketch_properties(t: AggregationType) -> SketchProperties { AggregationType::MultipleSum => p(true, true, true), AggregationType::MultipleIncrease => p(true, false, true), AggregationType::MultipleMinMax => p(true, false, true), - AggregationType::HydraKLL => p(true, false, true), - AggregationType::CountMinSketch => p(true, true, true), + AggregationType::HydraKLL => p(true, false, false), + AggregationType::CountMinSketch => p(true, true, false), // ponytail: heap top-k lists don't compose across windows; CMS cells do but the // combined type requires the heap, so neither merging nor subtracting is safe here. - AggregationType::CountMinSketchWithHeap => p(false, false, true), + AggregationType::CountMinSketchWithHeap => p(false, false, false), AggregationType::SetAggregator | AggregationType::DeltaSetAggregator => { p(true, false, false) } @@ -41,21 +43,21 @@ mod tests { use super::*; #[test] - fn cms_is_mergeable_subtractable_subpop_aware() { + fn cms_is_mergeable_subtractable_fixed_size() { let p = sketch_properties(AggregationType::CountMinSketch); - assert!(p.mergeable && p.subtractable && p.subpopulation_aware); + assert!(p.mergeable && p.subtractable && !p.memory_grows_with_keys); } #[test] fn cms_with_heap_not_mergeable_not_subtractable() { let p = sketch_properties(AggregationType::CountMinSketchWithHeap); - assert!(!p.mergeable && !p.subtractable && p.subpopulation_aware); + assert!(!p.mergeable && !p.subtractable); } #[test] - fn sum_mergeable_subtractable_not_subpop() { - let p = sketch_properties(AggregationType::Sum); - assert!(p.mergeable && p.subtractable && !p.subpopulation_aware); + fn multiple_sum_memory_grows_with_keys() { + let p = sketch_properties(AggregationType::MultipleSum); + assert!(p.mergeable && p.subtractable && p.memory_grows_with_keys); } #[test] diff --git a/asap-planner-rs/src/optimizer/solution.rs b/asap-planner-rs/src/optimizer/solution.rs index 34516226..ffb2d90f 100644 --- a/asap-planner-rs/src/optimizer/solution.rs +++ b/asap-planner-rs/src/optimizer/solution.rs @@ -69,6 +69,10 @@ pub struct AQEAssignment { /// `None` means the EXACT_a fallback (no streaming config, raw query). pub aggregation_id: Option, + /// ID of the paired key aggregation, for value sketches that can't list + /// their own keys. + pub key_aggregation_id: Option, + /// How this AQE's answer is derived from the assigned config. pub query_method: QueryMethod, @@ -143,6 +147,7 @@ impl OptimizerSolution { .map(|aqe| AQEAssignment { aqe, aggregation_id: None, + key_aggregation_id: None, query_method: QueryMethod::Exact, estimated_query_cost_per_sec: 0.0, }) diff --git a/asap-planner-rs/src/optimizer/translator.rs b/asap-planner-rs/src/optimizer/translator.rs index 8f5a07a5..cf7a8c98 100644 --- a/asap-planner-rs/src/optimizer/translator.rs +++ b/asap-planner-rs/src/optimizer/translator.rs @@ -36,11 +36,18 @@ fn build_inference_config(solution: &OptimizerSolution) -> InferenceConfig { }; let retain = retention_count_for_assignment(&assignment.query_method); let agg_ref = AggregationReference::new(aggregation_id, Some(retain)); + let key_ref = assignment.key_aggregation_id.map(|key_id| { + let key_retain = solution.deployed_configs()[&key_id].num_aggregates_to_retain; + AggregationReference::new(key_id, key_retain) + }); for query_string in &assignment.aqe.query_strings { - inference - .query_configs - .push(QueryConfig::new(query_string.clone()).add_aggregation(agg_ref.clone())); + let mut query_config = + QueryConfig::new(query_string.clone()).add_aggregation(agg_ref.clone()); + if let Some(key_ref) = &key_ref { + query_config = query_config.add_aggregation(key_ref.clone()); + } + inference.query_configs.push(query_config); } } diff --git a/asap-planner-rs/src/planner/agg_config.rs b/asap-planner-rs/src/planner/agg_config.rs index 92dcd60d..04b6b11b 100644 --- a/asap-planner-rs/src/planner/agg_config.rs +++ b/asap-planner-rs/src/planner/agg_config.rs @@ -8,6 +8,15 @@ use std::collections::HashMap; use crate::planner::labels::set_subpopulation_labels; use crate::planner::window::IntermediateWindowConfig; +/// Value aggregations that can't list their own keys, so the engine needs a +/// paired DeltaSetAggregator to enumerate groups at query time. +pub fn needs_key_aggregation(agg_type: AggregationType) -> bool { + matches!( + agg_type, + AggregationType::CountMinSketch | AggregationType::HydraKLL + ) +} + /// Internal representation of an aggregation config before IDs are assigned #[derive(Debug, Clone)] pub struct IntermediateAggConfig { @@ -105,10 +114,7 @@ pub fn build_agg_configs_for_statistics( &mut aggregated, ); - if matches!( - agg_type, - AggregationType::CountMinSketch | AggregationType::HydraKLL - ) { + if needs_key_aggregation(agg_type) { let delta_params = get_params(AggregationType::DeltaSetAggregator, "")?; configs.push(IntermediateAggConfig { aggregation_type: AggregationType::DeltaSetAggregator, From fcd02916785985b4e58257b3fcba30969b9bd4c9 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Sun, 4 Oct 2026 17:34:17 -0400 Subject: [PATCH 3/4] docs(planner): document label-set facts and cardinality costing Co-Authored-By: Claude Opus 5.5 --- .../optimizer-label-set-facts.example.yaml | 11 ++ .../optimizer-v1-implementation-plan.md | 117 ++++++++++++------ .design_docs/optimizer-workload.example.yaml | 10 ++ 3 files changed, 100 insertions(+), 38 deletions(-) create mode 100644 .design_docs/optimizer-label-set-facts.example.yaml create mode 100644 .design_docs/optimizer-workload.example.yaml diff --git a/.design_docs/optimizer-label-set-facts.example.yaml b/.design_docs/optimizer-label-set-facts.example.yaml new file mode 100644 index 00000000..256419e4 --- /dev/null +++ b/.design_docs/optimizer-label-set-facts.example.yaml @@ -0,0 +1,11 @@ +# Label-set facts for optimizer-workload.example.yaml (asap-optimizer-cli --label-set-facts). +series: + - {metric: http_requests_total, spatial_filter: 'env="prod"', series_count: 200} + - {metric: latency_ms, spatial_filter: "", series_count: 50} +groups: + # sum by (job): output labels [job]; also the `topk by (job)` buckets. + - {metric: http_requests_total, spatial_filter: 'env="prod"', grouping_labels: [job], cardinality: 10} + # topk keeps every label: output labels are all of them. + - {metric: http_requests_total, spatial_filter: 'env="prod"', grouping_labels: [env, instance, job], cardinality: 200} + # quantile_over_time keeps every label. + - {metric: latency_ms, spatial_filter: "", grouping_labels: [instance], cardinality: 50} diff --git a/.design_docs/optimizer-v1-implementation-plan.md b/.design_docs/optimizer-v1-implementation-plan.md index 658b172e..a9997716 100644 --- a/.design_docs/optimizer-v1-implementation-plan.md +++ b/.design_docs/optimizer-v1-implementation-plan.md @@ -35,7 +35,12 @@ The optimizer output type is `OptimizerSolution`, which is then translated into |---|---| | `mergeable` | Two instances can be combined → supports Merge query method | | `subtractable` | Element-wise difference defined → supports Subtract (tumbling only) | -| `subpopulation_aware` | One instance handles multiple label-group keys → N(s,g)=1 | +| `memory_grows_with_keys` | One instance is a keyed map whose memory grows per key (`Multiple*`), vs a fixed-size instance | + +How many instances a config has comes from its `grouping_labels`, split the same way as the +legacy planner (`planner/labels.rs::set_subpopulation_labels`): keyed types keep the output +labels in `aggregated_labels` (one instance), per-group types in `grouping_labels` (one instance +per group), and `topk by (L)` gets one heap per value of `L`. ### Query method (not a free variable — derived from ingest_type × W vs range_a × algebra) @@ -76,7 +81,7 @@ asap-planner-rs/src/optimizer/ ├── translator.rs translate(&OptimizerSolution) → (StreamingConfig, InferenceConfig) ├── aqe_extractor.rs Rqe, extract_aqes() ├── pipeline.rs run_all_exact_pipeline(), run_greedy_pipeline() -├── dataset.rs CSV series inventory, schema validation, and N_g profiling [#693] +├── label_set_facts.rs external series counts / group cardinalities → per-AQE facts [#756] ├── sketch_properties.rs algebraic properties per AggregationType [Phase 2a, done] ├── candidate_gen.rs enumerate candidate configs per AQE [Phase 2b, done] ├── cost_model.rs ingest/query cost formulas [Phase 2c, done] @@ -126,20 +131,23 @@ Each AQE gets its own best config independently (no cross-AQE sharing). MIP shar #### 2a — `sketch_properties.rs` (done) ```rust -pub struct SketchProperties { pub mergeable: bool, pub subtractable: bool, pub subpopulation_aware: bool } +pub struct SketchProperties { pub mergeable: bool, pub subtractable: bool, pub memory_grows_with_keys: bool } pub fn sketch_properties(t: AggregationType) -> SketchProperties ``` Implemented values — `CountMinSketchWithHeap` ended up `mergeable=false` (not `true`/unclear as originally guessed): the heap top-k list doesn't compose across merged/subtracted windows even -though the underlying CMS cells would. Everything else matches the original guess: -- `CountMinSketch`: mergeable=T, subtractable=T, subpopulation_aware=T -- `CountMinSketchWithHeap`: mergeable=F, subtractable=F, subpopulation_aware=T -- `DatasketchesKLL`, `HydraKLL`: mergeable=T, subtractable=F, subpopulation_aware=F (true for HydraKLL too) -- `Sum`/`MultipleSum`: mergeable=T, subtractable=T, subpopulation_aware=F/T respectively -- `HLL`, `SetAggregator`, `DeltaSetAggregator`: mergeable=T, subtractable=F, subpopulation_aware=F -- `MinMax`/`MultipleMinMax`, `Increase`/`MultipleIncrease`: mergeable=T, subtractable=F, subpopulation_aware=F/T -- `SingleSubpopulation`/`MultipleSubpopulation` (legacy wrapper types): all false (unknown, treated conservatively) +though the underlying CMS cells would. `memory_grows_with_keys` is true only for the `Multiple*` +maps: +- `CountMinSketch`: mergeable=T, subtractable=T +- `CountMinSketchWithHeap`: mergeable=F, subtractable=F +- `DatasketchesKLL`, `HydraKLL`: mergeable=T, subtractable=F +- `Sum`/`MultipleSum`: mergeable=T, subtractable=T +- `HLL`, `SetAggregator`, `DeltaSetAggregator`: mergeable=T, subtractable=F +- `MinMax`/`MultipleMinMax`, `Increase`/`MultipleIncrease`: mergeable=T, subtractable=F + +The optimizer never proposes single-group `Sum`/`MinMax`/`Increase`: with analytical per-key +memory they cost the same as their `Multiple*` forms (`OPTIMIZER_SKIPPED_AGG_TYPES`). When a type is both `mergeable` and `subtractable` (e.g. `Sum`, `CountMinSketch`), `candidate_gen.rs` prefers `Subtract` — it's O(1) regardless of `n`, strictly cheaper than `Merge`'s O(n). @@ -151,11 +159,18 @@ pub struct CandidateConfig { pub config: Option, // None = EXACT fallback pub query_method: QueryMethod, pub n_windows: u64, - pub label_group_count: u64, // dataset-derived N_g + pub instance_count: u64, // from the config's grouping labels + pub output_group_count: u64, // cardinality of the AQE's output labels + pub key_config: Option, // paired DeltaSet for CMS/HydraKLL } -pub fn enumerate_candidates(aqe: &Aqe, scrape_interval_secs: u64) -> Vec +pub fn enumerate_candidates_with_facts(aqe: &Aqe, scrape_interval_ms: u64, facts: &ItemFacts) -> Vec ``` +CMS and HydraKLL can't list their own keys, so each such candidate carries the same paired +DeltaSetAggregator the legacy planner deploys (`needs_key_aggregation`): same labels, Tumbling +at the value's slide, retaining `⌈range_a / slide⌉` panes. Greedy registers both and the +translator references both in the query config. + Enumeration: `compatible_agg_types(stat)` × param grid (hardcoded small grids per sketch type, see `CMS_DEPTHS`/`CMS_WIDTHS`/etc. constants — replace with sketch-bench sweep results in Phase 3) × window candidates × ingest type: @@ -179,15 +194,23 @@ Label granularity: only proposes configs at the label granularity the AQE itself pub struct AtomicCosts { mem_bytes_per_instance, insert_cpu_secs, merge_cpu_secs, subtract_cpu_secs, query_cpu_secs, exact_query_cpu_secs: f64 } pub struct CostWeights { ingest_mem, ingest_cpu, query_mem, query_cpu: f64 } -pub fn ingest_cost(candidate: &CandidateConfig, rho_g: f64, costs: &AtomicCosts, weights: &CostWeights) -> f64 +pub fn ingest_cost(candidate: &CandidateConfig, arrival_rate_hz: f64, costs: &AtomicCosts, weights: &CostWeights) -> f64 pub fn query_cost(a: &Aqe, candidate: &CandidateConfig, costs: &AtomicCosts, weights: &CostWeights) -> f64 -pub fn total_cost_rate(a: &Aqe, candidate: &CandidateConfig, rho_g: f64, costs: &AtomicCosts, weights: &CostWeights) -> f64 +pub fn total_cost_rate(a: &Aqe, candidate: &CandidateConfig, arrival_rate_hz: f64, costs: &AtomicCosts, weights: &CostWeights) -> f64 ``` -Implements the design doc's formulas with `N(s,g) = 1` for subpopulation-aware sketches and -the dataset-derived `N_g` for per-label-group sketches. The CSV series inventory and exact -filter profiling are implemented in `dataset.rs` (#693); Prometheus-backed profiling remains -the follow-up in #525. +Group cardinality enters per candidate: + +| | Multiplier | +|---|---| +| Memory, merge/subtract CPU | `output_group_count` for keyed maps (`Multiple*`), else `instance_count` | +| Query read CPU | `output_group_count` (one read per output group); top-k: `instance_count` (one heap read each) | +| Insert CPU | arrival rate `λ`, independent of cardinality | +| Paired key tracker | `output_group_count` key entries + one insert per item | + +Trivial accumulators (`Sum`/`MinMax`/`Increase` and `Multiple*`) use an analytical per-group +entry, `(n_labels × 4 + value_bytes) × 8/7`: label values dictionary-encoded as `u32` codes +(dictionary amortized, not charged), hashbrown load-factor slack. CPU stays on the stub. `AtomicCosts` added `exact_query_cpu_secs` (not in the original plan) — without a non-zero cost for the EXACT fallback's query, `IngestCost=0, QueryCost=0` would make EXACT always win trivially. @@ -212,18 +235,19 @@ it then, reusing `window_compatible`, `spatial_filter_compatible`, `topk_weighti #### 2e — `greedy.rs` + `pipeline.rs` + `translator.rs` (done) ```rust -pub fn greedy_assign(aqes: Vec, scrape_interval_secs: u64, rho_g: f64, - costs: &AtomicCosts, weights: &CostWeights) -> OptimizerSolution +pub fn greedy_assign(aqes: Vec, scrape_interval_ms: u64, atomic_cost_table: &AtomicCostTable, + weights: &CostWeights, facts: &[ItemFacts]) -> OptimizerSolution ``` For each AQE independently: `argmin` over `enumerate_candidates(aqe, ...)` by `total_cost_rate`. No feasibility filter needed (see 2d note) — every candidate from `enumerate_candidates` is valid for that AQE by construction. Assigns sequential `aggregation_id`s to deployed configs. -`run_greedy_pipeline(config, dataset, scrape_interval_secs, rho_g) -> Result<(StreamingConfig, InferenceConfig), DatasetError>` -added to `pipeline.rs` alongside `run_all_exact_pipeline()`. The dataset is required for the -greedy path and supplies metric schemas plus the `N_g` profile for each AQE. `rho_g` is -currently a single placeholder value applied uniformly to every candidate — see open TODO below. +`run_greedy_pipeline(config, facts, scrape_interval_ms, atomic_cost_table) -> Result<(StreamingConfig, InferenceConfig), LabelSetFactsError>` +added to `pipeline.rs` alongside `run_all_exact_pipeline()`. The workload's `metrics:` hints +are required and supply the label schema; every workload metric must have one. `facts` +(`LabelSetFacts`) supplies each AQE's group cardinalities and series count; arrival rate is +derived as `series_count × 1000 / scrape_interval_ms`, assuming one sample per series per scrape. `translator.rs::build_inference_config()` now populates `query_configs`: for each assignment with a real `aggregation_id`, emits one `QueryConfig` per original query string, with @@ -319,20 +343,36 @@ standalone `asap-optimizer-cli` binary (`asap-planner-rs/src/bin/optimizer_cli.r ``` cargo run -p asap_planner --bin asap-optimizer-cli -- \ --input_config \ - --dataset \ + --label-set-facts \ --data-ingestion-interval-ms 60000 \ - [--rho 1.0] \ [--atomic-costs \ --atomic-cost-workload ] ``` -Takes the same `ControllerConfig` YAML format as `asap-planner --input_config`. -The required CSV dataset is one row per unique metric series, with a fixed `metric` -column and label columns. It derives the metric schema and each AQE's distinct -group count; no live Prometheus connection is needed. A `metrics:` hints block, when -present, is checked against the dataset and mismatches fail loudly. -Prints deployed streaming configs and query configs to stdout. `--rho` is the -placeholder arrival rate (see TODOs below — not real yet). `--atomic-costs` is +Takes the same `ControllerConfig` YAML format as `asap-planner --input_config`; its +`metrics:` hints are required and must cover every workload metric, since they supply the +label schema. `--label-set-facts` is a YAML file of externally provided facts — the optimizer +never estimates them (example: `optimizer-label-set-facts.example.yaml`, for the workload +`optimizer-workload.example.yaml`): + +```yaml +series: # per (metric, spatial_filter) + - {metric: http_requests_total, spatial_filter: 'env="prod"', series_count: 200} +groups: # per (metric, spatial_filter, grouping_labels) + - {metric: http_requests_total, spatial_filter: 'env="prod"', grouping_labels: [job], cardinality: 10} +``` + +- `spatial_filter` and `grouping_labels` are required; `""` means unfiltered, `[]` means + full aggregation (cardinality must then be 1). Filters are normalized like the workload's. +- `grouping_labels` are the AQE's resolved output labels (all labels for `*_over_time` and + top-k; the `by` labels after resolving `without`). A `topk by (L)` query also needs a + `groups` row for `[L]`. +- `1 ≤ cardinality ≤ series_count`. Missing facts are an error listing every expected key; + unused facts only warn. +- Arrival rate is `series_count × 1000 / data-ingestion-interval-ms`. + +No live Prometheus connection is needed. +Prints deployed streaming configs and query configs to stdout. `--atomic-costs` is optional; when supplied it requires `--atomic-cost-workload`, a JSON file containing the exact `profiles[].workload` value from that benchmark artifact. The loader validates the document schema and rejects a selector that matches @@ -363,12 +403,13 @@ CMS-with-heap candidates warn and are dropped until a matching reference row is # Run the actual optimizer: cargo run -p asap_planner --bin asap-optimizer-cli -- \ --input_config workload.yaml --data-ingestion-interval-ms 60000 \ - --dataset path/to/series-inventory.csv \ + --label-set-facts path/to/label_set_facts.yaml \ --atomic-costs path/to/atomic_costs.json \ --atomic-cost-workload path/to/workload-selector.json ``` -`candidate-gen-dump`'s output labels each resolved params row `[real]` or `[stub]`; candidates +`candidate-gen-dump`'s output labels each resolved params row `[real]`, `[stub]`, or +`[analytical mem, stub cpu]` (trivial accumulators); candidates whose required cost row is missing are shown as `DROPPED`. CMS-with-heap uses its fixed-top-k reference model when the `{rows,cols}` row is present and is dropped otherwise. @@ -384,9 +425,9 @@ instead of `generator::generate_plan()`, likely behind an opt-in flag first. |---|---|---| | `aqe_extractor.rs` | `extract_requirements()` | Duplicates `build_query_requirements_promql` in `asap-query-engine/src/engines/simple_engine/promql.rs:614` — extract to shared free fn in `asap_types::query_requirements` | | `capability_matching.rs` | `labels_compatible()`, line 86 | Relax to superset matching — do in Phase 3b | -| `cost_model.rs` | `ingest_cost()`, `query_cost()` | Prometheus-backed `N_g` profiling remains for #525; dataset-backed `N_g` is implemented in #693 | +| `cost_model.rs` | `key_tracker_ingest_cost()` | Paired DeltaSet uses analytical key bytes and stub insert CPU until sketch-bench measures it | | `cost_model.rs` | `CostWeights::default()` | Self-consistent stub ratio (mem:cpu ≈ 1e-9:1), not real $/byte-sec vs $/cpu-sec calibration | -| `greedy.rs`, `pipeline.rs` | `rho_g` parameter | Single placeholder value applied uniformly; real per-config rates need Prometheus scrape-rate × active-series-count, not wired up | +| `label_set_facts.rs` | `resolve()` | Arrival rate assumes one sample per series per scrape; overestimates sparse or irregular series | | `translator.rs` | `retention_count_for_assignment(Subtract)` | Returns hardcoded `1`; should be the actual checkpoint count needed to cover the full lookback | | `promql/generator.rs` | `generate_plan()` doc comment | Flags that `Controller::generate()` still uses the hardcoded path, not the optimizer — see "Offline Testing" section above | | — | Accuracy constraint | No `Error(a,g) ≤ ε_a` check exists anywhere; nothing stops picking an under-provisioned sketch (Phase 3d) | diff --git a/.design_docs/optimizer-workload.example.yaml b/.design_docs/optimizer-workload.example.yaml new file mode 100644 index 00000000..c6c46670 --- /dev/null +++ b/.design_docs/optimizer-workload.example.yaml @@ -0,0 +1,10 @@ +# Example workload for asap-optimizer-cli; facts in optimizer-label-set-facts.example.yaml. +query_groups: + - queries: + - 'sum by (job) (http_requests_total{env="prod"})' + - 'quantile_over_time(0.99, latency_ms[5m])' + - 'topk by (job) (5, sum_over_time(http_requests_total{env="prod"}[5m]))' + repetition_delay_ms: 60000 +metrics: + - {metric: http_requests_total, labels: [job, instance, env]} + - {metric: latency_ms, labels: [instance]} From 0032c81819d542acbeacbde8d63d9ce2bcd2dcc2 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Sun, 4 Oct 2026 17:55:22 -0400 Subject: [PATCH 4/4] docs(planner): keep top-k in the optimizer example workload Co-Authored-By: Claude Opus 5.5 --- .design_docs/optimizer-label-set-facts.example.yaml | 4 +++- .design_docs/optimizer-v1-implementation-plan.md | 6 ++++-- .design_docs/optimizer-workload.example.yaml | 3 +++ 3 files changed, 10 insertions(+), 3 deletions(-) diff --git a/.design_docs/optimizer-label-set-facts.example.yaml b/.design_docs/optimizer-label-set-facts.example.yaml index a2c1c7d2..256419e4 100644 --- a/.design_docs/optimizer-label-set-facts.example.yaml +++ b/.design_docs/optimizer-label-set-facts.example.yaml @@ -3,7 +3,9 @@ series: - {metric: http_requests_total, spatial_filter: 'env="prod"', series_count: 200} - {metric: latency_ms, spatial_filter: "", series_count: 50} groups: - # sum by (job): output labels [job]. + # sum by (job): output labels [job]; also the `topk by (job)` buckets. - {metric: http_requests_total, spatial_filter: 'env="prod"', grouping_labels: [job], cardinality: 10} + # topk keeps every label: output labels are all of them. + - {metric: http_requests_total, spatial_filter: 'env="prod"', grouping_labels: [env, instance, job], cardinality: 200} # quantile_over_time keeps every label. - {metric: latency_ms, spatial_filter: "", grouping_labels: [instance], cardinality: 50} diff --git a/.design_docs/optimizer-v1-implementation-plan.md b/.design_docs/optimizer-v1-implementation-plan.md index a9997716..0dc16224 100644 --- a/.design_docs/optimizer-v1-implementation-plan.md +++ b/.design_docs/optimizer-v1-implementation-plan.md @@ -353,7 +353,8 @@ Takes the same `ControllerConfig` YAML format as `asap-planner --input_config`; `metrics:` hints are required and must cover every workload metric, since they supply the label schema. `--label-set-facts` is a YAML file of externally provided facts — the optimizer never estimates them (example: `optimizer-label-set-facts.example.yaml`, for the workload -`optimizer-workload.example.yaml`): +`optimizer-workload.example.yaml`; its top-k query needs `--atomic-costs` with a CMS-with-heap +reference row, otherwise the run fails with that item listed as unservable): ```yaml series: # per (metric, spatial_filter) @@ -378,7 +379,8 @@ containing the exact `profiles[].workload` value from that benchmark artifact. The loader validates the document schema and rejects a selector that matches zero or multiple profiles; it never mixes entries across workloads. Omit both flags and ordinary unbenchmarked candidates use the flat stub, while -CMS-with-heap candidates warn and are dropped until a matching reference row is available. +CMS-with-heap candidates warn and are dropped until a matching reference row is available; +an item left with no candidate makes the run fail, listing every unservable item. ### Running with real sketch-bench costs diff --git a/.design_docs/optimizer-workload.example.yaml b/.design_docs/optimizer-workload.example.yaml index b5200999..2372f47a 100644 --- a/.design_docs/optimizer-workload.example.yaml +++ b/.design_docs/optimizer-workload.example.yaml @@ -1,8 +1,11 @@ # Example workload for asap-optimizer-cli; facts in optimizer-label-set-facts.example.yaml. +# The top-k query is only servable with --atomic-costs containing a CMS-with-heap +# reference row; without it the optimizer reports it as unservable. query_groups: - queries: - 'sum by (job) (http_requests_total{env="prod"})' - 'quantile_over_time(0.99, latency_ms[5m])' + - 'topk by (job) (5, sum_over_time(http_requests_total{env="prod"}[5m]))' repetition_delay_ms: 60000 metrics: - {metric: http_requests_total, labels: [job, instance, env]}