diff --git a/.design_docs/optimizer-label-set-facts.example.yaml b/.design_docs/optimizer-label-set-facts.example.yaml deleted file mode 100644 index 256419e4..00000000 --- a/.design_docs/optimizer-label-set-facts.example.yaml +++ /dev/null @@ -1,11 +0,0 @@ -# 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-workload.example.yaml b/.design_docs/optimizer-workload.example.yaml deleted file mode 100644 index 2372f47a..00000000 --- a/.design_docs/optimizer-workload.example.yaml +++ /dev/null @@ -1,12 +0,0 @@ -# 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]} - - {metric: latency_ms, labels: [instance]} diff --git a/asap-planner-rs/Cargo.toml b/asap-planner-rs/Cargo.toml index 3918f0b1..5e9bf33d 100644 --- a/asap-planner-rs/Cargo.toml +++ b/asap-planner-rs/Cargo.toml @@ -15,10 +15,6 @@ path = "src/main.rs" name = "asap-optimizer-cli" path = "src/bin/optimizer_cli.rs" -[[bin]] -name = "candidate-gen-dump" -path = "src/bin/candidate_gen_dump.rs" - [[bin]] name = "benchmark-promql-status" path = "src/bin/benchmark_promql_status.rs" diff --git a/asap-planner-rs/Dockerfile b/asap-planner-rs/Dockerfile index b040253f..c6d68ee6 100644 --- a/asap-planner-rs/Dockerfile +++ b/asap-planner-rs/Dockerfile @@ -22,7 +22,6 @@ COPY asap-planner-rs/Cargo.toml ./asap-planner-rs/ RUN mkdir -p asap-query-engine/src && echo "fn main() {}" > asap-query-engine/src/main.rs && \ mkdir -p asap-planner-rs/src/bin && echo "fn main() {}" > asap-planner-rs/src/main.rs && \ echo "fn main() {}" > asap-planner-rs/src/bin/optimizer_cli.rs && \ - echo "fn main() {}" > asap-planner-rs/src/bin/candidate_gen_dump.rs && \ echo "fn main() {}" > asap-planner-rs/src/bin/benchmark_promql_status.rs && \ echo "pub fn placeholder() {}" >> asap-planner-rs/src/lib.rs diff --git a/asap-planner-rs/src/bin/candidate_gen_dump.rs b/asap-planner-rs/src/bin/candidate_gen_dump.rs deleted file mode 100644 index 80eb59c4..00000000 --- a/asap-planner-rs/src/bin/candidate_gen_dump.rs +++ /dev/null @@ -1,249 +0,0 @@ -//! Dump all candidate configs produced by enumerate_candidates for each AQE, -//! before greedy selection. Useful for auditing what the candidate space looks like. - -use std::collections::HashMap; -use std::path::PathBuf; - -use asap_planner::{ - optimizer::{ - enumerate_candidates, extract_aqes, load_optional_selected_atomic_cost_table, - resolve_atomic_costs, AtomicCostTable, AtomicCosts, CandidateConfig, RQE, - }, - ControllerConfig, -}; -use asap_types::enums::WindowType; -use clap::Parser; -use promql_utilities::query_logics::enums::AggregationType; -use serde_json::Value; - -#[derive(Parser)] -#[command( - name = "candidate-gen-dump", - about = "Print all candidates from enumerate_candidates (pre-selection) for a workload config" -)] -struct Args { - #[arg(long = "input_config")] - input_config: PathBuf, - - #[arg(long = "data-ingestion-interval-ms")] - scrape_interval_ms: u64, - - /// Path to sketch-bench's versioned atomic-cost document. - /// When given, each params row also prints its resolved AtomicCosts -- - /// real (from the table) or the flat stub (unbenchmarked family, or this - /// exact param point missing from the table) -- labeled which. - #[arg(long = "atomic-costs")] - atomic_costs: Option, - - /// JSON `profiles[].workload` value selecting exactly one measured profile. - #[arg(long = "atomic-cost-workload", requires = "atomic_costs")] - atomic_cost_workload: Option, -} - -fn main() -> anyhow::Result<()> { - let args = Args::parse(); - - let atomic_cost_table = load_optional_selected_atomic_cost_table( - args.atomic_costs.as_deref(), - args.atomic_cost_workload.as_deref(), - )? - .unwrap_or_default(); - - let yaml_str = std::fs::read_to_string(&args.input_config)?; - let config: ControllerConfig = serde_yaml::from_str(&yaml_str)?; - let schema = config.schema_from_hints(); - - let rqes: Vec = config - .query_groups - .iter() - .flat_map(|qg| { - qg.queries.iter().map(|q| RQE { - query_string: q.clone(), - t_repeat_ms: qg.repetition_delay_ms, - accuracy_sla: qg.controller_options.accuracy_sla, - latency_sla_ms: qg.controller_options.latency_sla_ms, - }) - }) - .collect(); - - let aqes = extract_aqes(&rqes, &schema, args.scrape_interval_ms)?; - println!("=== {} optimizer item(s) ===", aqes.len()); - - for (i, aqe) in aqes.iter().enumerate() { - println!( - "\n--- Item #{i}: metric={} stat={:?} range={}ms T={}ms accuracy_sla={} latency_sla_ms={:?} freq={:.4}Hz ---", - aqe.requirements.metric, - aqe.requirements.statistics, - aqe.requirements.data_range_ms, - aqe.t_repeat_ms, - aqe.accuracy_sla, - aqe.latency_sla_ms, - aqe.query_frequency_hz, - ); - println!(" queries: {:?}", aqe.query_strings); - - let candidates = enumerate_candidates(aqe, args.scrape_interval_ms); - print_candidates_grouped( - &candidates, - &atomic_cost_table, - aqe.requirements.grouping_labels.len(), - ); - } - - Ok(()) -} - -/// Group the flat candidate list by (agg_type, sub_type) and print hierarchically: -/// SketchType (sub_type) -/// windows (N): -/// ... -/// params (M) [× N windows = NM total]: -/// ... -/// EXACT is printed last as a single line. -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) - type WindowKey = (WindowType, u64, u64, u64); // (type, size, slide, n) - #[allow(clippy::type_complexity)] - let mut groups: HashMap< - (String, String), - (AggregationType, Vec, Vec>), - > = HashMap::new(); - - let mut has_exact = false; - - for c in candidates { - let Some(cfg) = &c.config else { - has_exact = true; - continue; - }; - - let key = ( - format!("{:?}", cfg.aggregation_type), - cfg.aggregation_sub_type.clone(), - ); - - let entry = groups.entry(key.clone()).or_insert_with(|| { - group_order.push(key); - (cfg.aggregation_type, Vec::new(), Vec::new()) - }); - - let wk: WindowKey = ( - cfg.window_type, - cfg.window_size_ms, - cfg.slide_interval_ms, - c.n_windows, - ); - if !entry.1.contains(&wk) { - entry.1.push(wk); - } - - let mut params: Vec<(String, Value)> = cfg - .parameters - .iter() - .map(|(k, v)| (k.clone(), v.clone())) - .collect(); - params.sort_by(|(a, _), (b, _)| a.cmp(b)); - if !entry.2.contains(¶ms) { - entry.2.push(params); - } - } - - let total_sketch: usize = group_order - .iter() - .map(|k| { - let (_, ws, ps) = &groups[k]; - ws.len() * ps.len() - }) - .sum(); - let total = total_sketch + if has_exact { 1 } else { 0 }; - println!( - " {} candidate(s) across {} sketch group(s):", - total, - group_order.len() - ); - - for key in &group_order { - let (agg_type, windows, params) = &groups[key]; - let (agg_type_str, sub_type) = key; - let sub = if sub_type.is_empty() { - String::new() - } else { - format!(" ({})", sub_type) - }; - println!( - "\n {}{} — {} window(s) × {} param(s) = {} candidate(s):", - agg_type_str, - sub, - windows.len(), - params.len(), - windows.len() * params.len(), - ); - - println!(" windows:"); - for (wtype, size, slide, n) in windows { - let window_str = match wtype { - WindowType::Tumbling => format!("Tumbling {}ms", size), - WindowType::Sliding => format!("Sliding {}ms / slide {}ms", size, slide), - }; - println!( - " {} n={} → {:?}", - window_str, - n, - // Re-derive method label from n and window type for display. - if *n == 1 { "Direct" } else { "Merge/Subtract" } - ); - } - - println!(" params:"); - for p in params { - let kv: Vec<_> = p.iter().map(|(k, v)| format!("{k}: {v}")).collect(); - let costs_str = if atomic_cost_table.is_empty() { - String::new() - } 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, - n_grouping_labels, - ) { - Some(costs) => { - 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" - }; - format!(" -> [{label}] {}", format_costs(&costs)) - } - None => " -> DROPPED (no matching table row)".to_string(), - } - }; - println!(" {{{}}}{costs_str}", kv.join(", ")); - } - } - - if has_exact { - println!("\n [EXACT]"); - } -} - -fn format_costs(costs: &AtomicCosts) -> String { - format!( - "mem={:.0}B insert={:.3e}s merge={:.3e}s subtract={:.3e}s query={:.3e}s", - costs.mem_bytes_per_instance, - costs.insert_cpu_secs, - costs.merge_cpu_secs, - costs.subtract_cpu_secs, - costs.query_cpu_secs, - ) -} diff --git a/asap-planner-rs/src/bin/optimizer_cli.rs b/asap-planner-rs/src/bin/optimizer_cli.rs index e2c7cd02..d9db9193 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -1,15 +1,11 @@ -//! Offline runner for the optimization-based sketch/config selector. -//! -//! Standalone: not wired into `asap-planner`/`Controller::generate()` yet. Lets -//! you exercise `run_greedy_pipeline` against real workload configs while the -//! optimizer module is still under development (Phase 2 of issue #405). +//! Offline runner for the rqe-optimizer MILP planner: prints the plan and its +//! cost, and optionally writes the streaming and inference configs. use std::path::PathBuf; use asap_planner::optimizer::{ - build_milp_workload, load_flat_atomic_cost_table, load_optional_selected_atomic_cost_table, - load_workload_facts, plan_to_planner_output, reject_avg_queries, run_greedy_pipeline, - solve_milp, AtomicCostTable, LabelSetFacts, LabelSetFactsError, + build_milp_workload, load_flat_atomic_cost_table, load_workload_facts, plan_to_planner_output, + reject_avg_queries, solve_milp, MilpError, }; use asap_planner::ControllerConfig; use clap::Parser; @@ -18,7 +14,7 @@ use rqe_optimizer::milp::Objective; #[derive(Parser, Debug)] #[command( name = "asap-optimizer-cli", - about = "Offline runner for the optimization-based sketch/config selector (not wired into asap-planner yet)" + about = "Offline runner for the rqe-optimizer MILP planner" )] struct Args { /// Path to a YAML workload config (same format as `asap-planner --input_config`). @@ -30,68 +26,32 @@ struct Args { #[arg(long = "data-ingestion-interval-ms", value_parser = clap::value_parser!(u64).range(1..))] data_ingestion_interval_ms: u64, - /// Greedy only. YAML label-set facts: `series_count` per (metric, spatial - /// filter) and `cardinality` per (metric, spatial filter, grouping labels). - #[arg( - long = "label-set-facts", - required_unless_present = "milp", - conflicts_with = "milp" - )] - label_set_facts: Option, + /// The flat cost table `export_rqe_optimizer_costs.sh` writes + /// (`rqe_atomic_costs.json`). + #[arg(long = "atomic-costs")] + atomic_costs: PathBuf, - /// Greedy: the versioned atomic-cost document sketch-bench's - /// `atomic-costs` subcommand exports; requires --atomic-cost-workload. - /// Omitted: every benchmarked-family candidate (CMS/HLL/KLL) is dropped, - /// leaving only trivial accumulators and EXACT. - /// MILP: the flat cost table `export_rqe_optimizer_costs.sh` writes - /// (`rqe_atomic_costs.json`); required. - #[arg(long = "atomic-costs", required_if_eq("milp", "true"))] - atomic_costs: Option, - - /// Greedy only. JSON `profiles[].workload` value copied from the - /// sketch-bench atomic-cost document. This makes the empirical workload - /// profile explicit and avoids mixing costs from different traces or time - /// windows. - #[arg( - long = "atomic-cost-workload", - requires = "atomic_costs", - conflicts_with = "milp" - )] - atomic_cost_workload: Option, - - /// Plan with sketch-bench's rqe-optimizer MILP and print the plan. - #[arg(long)] - milp: bool, - - /// MILP only. Write `streaming_config.yaml` and `inference_config.yaml` - /// for the plan here. - #[arg(long = "output-dir", requires = "milp")] + /// Write `streaming_config.yaml` and `inference_config.yaml` for the plan + /// here. + #[arg(long = "output-dir")] output_dir: Option, - /// MILP only. Also plan with families the engine can't deploy; prints the - /// plan and writes no configs. - #[arg( - long = "allow-undeployable-families", - requires = "milp", - conflicts_with = "output_dir" - )] + /// Also plan with families the engine can't deploy; prints the plan and + /// writes no configs. + #[arg(long = "allow-undeployable-families", conflicts_with = "output_dir")] allow_undeployable_families: bool, - /// MILP only. YAML workload facts: per metric, positive `value_range` - /// and `cardinality` per label set, including all labels (the series count). - #[arg( - long = "workload-facts", - required_if_eq("milp", "true"), - requires = "milp" - )] - workload_facts: Option, + /// YAML workload facts: per metric, positive `value_range` and + /// `cardinality` per label set, including all labels (the series count). + #[arg(long = "workload-facts")] + workload_facts: PathBuf, - /// MILP only. Objective weight on CPU-sec/sec. Default: rqe-optimizer's. - #[arg(long = "w-cpu", requires = "milp", value_parser = parse_weight)] + /// Objective weight on CPU-sec/sec. Default: rqe-optimizer's. + #[arg(long = "w-cpu", value_parser = parse_weight)] w_cpu: Option, - /// MILP only. Objective weight on memory GiB. Default: rqe-optimizer's. - #[arg(long = "w-mem", requires = "milp", value_parser = parse_weight)] + /// Objective weight on memory GiB. Default: rqe-optimizer's. + #[arg(long = "w-mem", value_parser = parse_weight)] w_mem: Option, #[arg(short, long, action = clap::ArgAction::Count)] @@ -111,56 +71,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)?; - if args.milp { - return run_milp(&args, &config); - } - let facts = LabelSetFacts::from_path( - args.label_set_facts - .as_deref() - .expect("clap requires --label-set-facts without --milp"), - )?; - - let atomic_cost_table = match load_optional_selected_atomic_cost_table( - args.atomic_costs.as_deref(), - args.atomic_cost_workload.as_deref(), - )? { - Some(table) => table, - None => { - tracing::warn!( - "no --atomic-costs supplied; CMS/HLL/KLL candidates will never be selected" - ); - AtomicCostTable::default() - } - }; - - let (streaming, inference) = run_greedy_pipeline( - &config, - &facts, - args.data_ingestion_interval_ms, - &atomic_cost_table, - )?; - - let deployed = streaming.get_all_aggregation_configs(); - println!("=== Deployed streaming configs: {} ===", deployed.len()); - for (id, cfg) in deployed { - println!( - " [{id}] {} sub_type={:?} window={}ms slide={}ms type={:?} metric={} params={:?}", - cfg.aggregation_type, - cfg.aggregation_sub_type, - cfg.window_size_ms, - cfg.slide_interval_ms, - cfg.window_type, - cfg.metric, - cfg.parameters, - ); - } - - println!("\n=== Query configs: {} ===", inference.query_configs.len()); - for qc in &inference.query_configs { - println!(" \"{}\" -> {:?}", qc.query, qc.aggregations); - } - - Ok(()) + run_milp(&args, &config) } /// Objective weights must be finite and non-negative: a negative weight @@ -176,20 +87,10 @@ fn parse_weight(s: &str) -> Result { fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { config.warn_default_slas(); let Some(hints) = config.metrics.as_deref() else { - return Err(LabelSetFactsError::MissingMetricHints.into()); + return Err(MilpError::MissingMetricHints.into()); }; - let facts = load_workload_facts( - args.workload_facts - .as_deref() - .expect("clap requires --workload-facts with --milp"), - hints, - args.data_ingestion_interval_ms, - )?; - let costs = load_flat_atomic_cost_table( - args.atomic_costs - .as_deref() - .expect("clap requires --atomic-costs with --milp"), - )?; + let facts = load_workload_facts(&args.workload_facts, hints, args.data_ingestion_interval_ms)?; + let costs = load_flat_atomic_cost_table(&args.atomic_costs)?; let Objective::AUCCost { w_cpu, w_mem } = Objective::default(); let (w_cpu, w_mem) = (args.w_cpu.unwrap_or(w_cpu), args.w_mem.unwrap_or(w_mem)); // All-zero weights make every plan cost 0, so the solver's pick is arbitrary. diff --git a/asap-planner-rs/src/optimizer/aqe_extractor.rs b/asap-planner-rs/src/optimizer/aqe_extractor.rs index 922dc02b..1a7b3262 100644 --- a/asap-planner-rs/src/optimizer/aqe_extractor.rs +++ b/asap-planner-rs/src/optimizer/aqe_extractor.rs @@ -125,16 +125,6 @@ fn normalized_f64_bits(value: f64) -> u64 { } } -/// Euclidean GCD. `num-integer` is not in the workspace; this two-liner is -/// sufficient and avoids a dependency. -pub(super) fn gcd(a: u64, b: u64) -> u64 { - if b == 0 { - a - } else { - gcd(b, a % b) - } -} - /// Recursively decompose a PromQL expression into non-binary leaf queries. /// /// Binary arithmetic expressions (e.g. `rate(a[5m]) / rate(b[5m])`) are split diff --git a/asap-planner-rs/src/optimizer/atomic_costs.rs b/asap-planner-rs/src/optimizer/atomic_costs.rs index 3167953f..6b7f4046 100644 --- a/asap-planner-rs/src/optimizer/atomic_costs.rs +++ b/asap-planner-rs/src/optimizer/atomic_costs.rs @@ -1,442 +1,12 @@ -//! The atomic-cost table exported by sketch-bench (sketch-bench#30, -//! `scripts/export_atomic_costs.sh`), and the (sketch_type, params) lookup -//! that resolves a candidate's [`AtomicCosts`] from it. +//! The flat atomic-cost table sketch-bench's `export_rqe_optimizer_costs.sh` +//! writes for the MILP. -use std::collections::HashMap; use std::path::Path; -use promql_utilities::query_logics::enums::AggregationType; -use serde::{Deserialize, Serialize}; -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, 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; - -const CMS_HEAP_BENCHMARK: &str = "cms-heap-topk-regularpath-vector2d"; -pub const ATOMIC_COST_SCHEMA_VERSION: u32 = 1; - -/// Versioned atomic-cost document emitted by `approxbench atomic-costs`. -/// -/// A profile is deliberately selected before candidate resolution: costs from -/// different input workloads must never be mixed by a flat lookup. -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct AtomicCostDocument { - pub schema_version: u32, - pub profiles: Vec, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct AtomicCostProfile { - pub workload: WorkloadDescription, - pub entries: Vec, -} - -/// The provenance of the input data on which atomic costs were measured. -/// Synthetic descriptions remain opaque because their generator schema evolves -/// independently; external profiles are represented explicitly for selection. -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -#[serde(rename_all = "lowercase")] -pub enum WorkloadDescription { - Synthetic { description: Value }, - External(ExternalWorkload), -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct ExternalWorkload { - pub source: String, - pub dataset: String, - pub mode: String, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub key_columns: Vec, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub group_columns: Vec, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub variate: Option, - pub value_column: String, - pub window_start_ns: i64, - pub window_end_ns: i64, - pub records_loaded: u64, - pub source_timestamp_unit: String, - pub timestamp_unit: String, -} - pub use rqe_optimizer::{AtomicCostEntry, AtomicCostTable}; -/// Parse a standalone JSON workload selector. The selector is the exact -/// `profiles[].workload` value copied from the benchmark artifact, making the -/// selected empirical input explicit in an offline planning run. -pub fn load_workload_selector(path: &Path) -> anyhow::Result { - let raw = std::fs::read_to_string(path).map_err(|e| { - anyhow::anyhow!( - "reading atomic-cost workload selector {}: {e}", - path.display() - ) - })?; - serde_json::from_str(&raw).map_err(|e| { - anyhow::anyhow!( - "parsing atomic-cost workload selector {}: {e}", - path.display() - ) - }) -} - -/// Read a versioned `sketch-bench atomic-costs` document and return entries -/// from exactly one requested workload profile. -pub fn load_atomic_cost_table( - path: &Path, - workload: &WorkloadDescription, -) -> anyhow::Result { - let raw = std::fs::read_to_string(path) - .map_err(|e| anyhow::anyhow!("reading atomic-cost table {}: {e}", path.display()))?; - let document: AtomicCostDocument = serde_json::from_str(&raw) - .map_err(|e| anyhow::anyhow!("parsing atomic-cost document {}: {e}", path.display()))?; - - if document.schema_version != ATOMIC_COST_SCHEMA_VERSION { - anyhow::bail!( - "unsupported atomic-cost schema_version {} in {} (this planner supports {})", - document.schema_version, - path.display(), - ATOMIC_COST_SCHEMA_VERSION - ); - } - - let matches: Vec<_> = document - .profiles - .iter() - .filter(|profile| profile.workload == *workload) - .collect(); - match matches.as_slice() { - [profile] => Ok(profile.entries.clone()), - [] => anyhow::bail!( - "no atomic-cost profile in {} matches workload selector {}", - path.display(), - serde_json::to_string(workload).unwrap_or_else(|_| "".into()) - ), - _ => anyhow::bail!( - "{} atomic-cost profiles in {} match workload selector {}; expected exactly one", - matches.len(), - path.display(), - serde_json::to_string(workload).unwrap_or_else(|_| "".into()) - ), - } -} - -/// Load the selector artifact and return the corresponding empirical table. -/// Offline callers use this single interface so selector validation cannot -/// drift between planner tools. -pub fn load_selected_atomic_cost_table( - document_path: &Path, - selector_path: &Path, -) -> anyhow::Result { - let workload = load_workload_selector(selector_path)?; - load_atomic_cost_table(document_path, &workload) -} - -/// Resolve the optional atomic-cost CLI inputs as one unit. A document without -/// its workload selector is invalid; callers can keep their no-document -/// fallback without duplicating that validation. -pub fn load_optional_selected_atomic_cost_table( - document_path: Option<&Path>, - selector_path: Option<&Path>, -) -> anyhow::Result> { - match (document_path, selector_path) { - (Some(document_path), Some(selector_path)) => { - load_selected_atomic_cost_table(document_path, selector_path).map(Some) - } - (Some(_), None) => anyhow::bail!( - "--atomic-cost-workload is required with --atomic-costs; \ - it must contain the selected profiles[].workload JSON value" - ), - (None, Some(_)) => anyhow::bail!("--atomic-cost-workload requires --atomic-costs"), - (None, None) => Ok(None), - } -} - -/// sketch-bench's (algorithm, params) key for one of ASAPQuery's benchmarked -/// families, or `None` if `agg_type` isn't one sketch-bench measures at all -/// (trivial O(1) accumulators — Sum/Increase/MinMax/... — and sketch types -/// sketch-bench has no wrapper for yet — HydraKLL). -/// -/// Field names differ from ASAPQuery's own `parameters` map by design: each -/// side picked its own config vocabulary independently, so this is a real -/// translation, not a passthrough. The field names read here (`"depth"`, -/// `"width"`, `"precision"`, `"K"`) are duplicated from `candidate_gen.rs`'s -/// `param_grid()` — not derived from it — so a rename on either side without -/// the other silently breaks this lookup. Guarded by panicking below rather -/// than treating a missing key the same as "not a benchmarked family": a -/// `CountMinSketch`/`HLL`/`DatasketchesKLL` candidate is only ever built by -/// `param_grid()`, which always sets these keys, so their absence means the -/// two have drifted, not that there's no data for this family. -fn sketch_bench_key( - agg_type: AggregationType, - params: &HashMap, -) -> Option<(&'static str, Value)> { - match agg_type { - AggregationType::CountMinSketch => Some(( - "cms-fastpath-vector2d", - serde_json::json!({ - "rows": require(params, "depth", agg_type), - "cols": require(params, "width", agg_type), - }), - )), - AggregationType::HLL => Some(( - "hll", - serde_json::json!({ "lg_k": require(params, "precision", agg_type) }), - )), - AggregationType::DatasketchesKLL => Some(( - "kll-percall", - serde_json::json!({ "k": require(params, "K", agg_type) }), - )), - _ => None, - } -} - -/// 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`]): -/// `Some(AtomicCosts::default())` — the flat stub, unchanged from before -/// this table existed. Logged, since it's silently wrong for anything -/// sketch-bench could plausibly measure later. -/// TODO(#524): remove this fallback once every family the optimizer can -/// select has a real sketch-bench entry; costing should end up 100% -/// empirical, with nothing left reading `AtomicCosts::default()`. -/// - Benchmarked family, matching table row found: `Some(costs)` built from -/// it (`subtract_cpu_secs`/`exact_query_cpu_secs` still come from the -/// stub — the table has neither: subtract isn't implemented upstream yet, -/// and EXACT isn't a sketch sketch-bench could measure). -/// - `CountMinSketchWithHeap`: resolve through the temporary fixed-top-k -/// reference model in [`resolve_cms_heap_costs`]. Missing or malformed -/// 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, - "no sketch-bench atomic-cost data for this family; using the flat AtomicCosts stub" - ); - return Some(AtomicCosts::default()); - }; - - let expected_config = serde_json::json!({ "algorithm": sketch, "params": sketch_params }); - table - .iter() - .find(|e| e.sketch == sketch && e.sketch_config == expected_config) - .map(|entry| AtomicCosts { - mem_bytes_per_instance: entry.mem_bytes_per_instance, - insert_cpu_secs: entry.insert_cpu_secs, - merge_cpu_secs: entry.merge_cpu_secs, - subtract_cpu_secs: SUBTRACT_CPU_SECS, - query_cpu_secs: entry.query_cpu_secs, - exact_query_cpu_secs: EXACT_QUERY_CPU_SECS, - }) -} - -/// Temporary cost model for the runtime CMS-with-heap implementation. -/// -/// sketch-bench currently measures a fixed top-k=32 wrapper, while the -/// runtime's heap size is a candidate parameter. CPU costs therefore scale -/// linearly from the matching regular-path/top-k benchmark row. Memory is -/// computed from the runtime's i64 CMS counters and an explicit estimate for -/// each heap entry; the benchmark's i32-only matrix memory is not reused. -#[derive(Debug, Clone, Copy, PartialEq)] -struct CmsHeapCostAssumptions { - reference_heap_size: u64, - counter_bytes: f64, - average_key_bytes: f64, - heap_entry_overhead_bytes: f64, -} - -impl Default for CmsHeapCostAssumptions { - fn default() -> Self { - Self { - reference_heap_size: CMS_HEAP_REFERENCE_HEAP_SIZE, - counter_bytes: CMS_HEAP_COUNTER_BYTES, - average_key_bytes: CMS_HEAP_AVERAGE_KEY_BYTES, - heap_entry_overhead_bytes: CMS_HEAP_ENTRY_OVERHEAD_BYTES, - } - } -} - -impl CmsHeapCostAssumptions { - fn validate(self) { - assert!( - self.reference_heap_size > 0, - "CMS-with-heap reference_heap_size must be greater than zero" - ); - assert!( - self.counter_bytes.is_finite() && self.counter_bytes >= 0.0, - "CMS-with-heap counter_bytes must be finite and non-negative" - ); - assert!( - self.average_key_bytes.is_finite() && self.average_key_bytes >= 0.0, - "CMS-with-heap average_key_bytes must be finite and non-negative" - ); - assert!( - self.heap_entry_overhead_bytes.is_finite() && self.heap_entry_overhead_bytes >= 0.0, - "CMS-with-heap heap_entry_overhead_bytes must be finite and non-negative" - ); - } -} - -/// Resolve a CMS-with-heap candidate from the fixed-top-k benchmark reference. -/// -/// The helper deliberately returns `None` when the reference row is absent or -/// malformed. The caller then drops this candidate, leaving the always-feasible -/// EXACT candidate available. TODO(#651): turn these temporary warning paths -/// into hard errors once sketch-bench sweeps cover the candidate grid. -fn resolve_cms_heap_costs( - table: &AtomicCostTable, - params: &HashMap, - assumptions: &CmsHeapCostAssumptions, -) -> Option { - assumptions.validate(); - - let agg_type = AggregationType::CountMinSketchWithHeap; - let depth = require_u64(params, "depth", agg_type); - let width = require_u64(params, "width", agg_type); - let heap_size = require_u64(params, "heapsize", agg_type); - let heap_size_f64 = heap_size as f64; - let scale = heap_size_f64 / assumptions.reference_heap_size as f64; - let expected_config = serde_json::json!({ - "algorithm": CMS_HEAP_BENCHMARK, - "params": { "rows": depth, "cols": width }, - }); - - let Some(entry) = table - .iter() - .find(|entry| entry.sketch == CMS_HEAP_BENCHMARK && entry.sketch_config == expected_config) - else { - tracing::info!( - status = "missing_reference", - sketch = CMS_HEAP_BENCHMARK, - depth, - width, - heap_size, - "cms-with-heap atomic cost measurement" - ); - tracing::warn!( - sketch = CMS_HEAP_BENCHMARK, - depth, - width, - heap_size, - "no CMS-with-heap reference cost for candidate; dropping candidate; \ - TODO(#651): fail loudly once sketch-bench sweeps cover this grid" - ); - return None; - }; - - if !valid_cost_entry(entry) { - tracing::info!( - status = "invalid_reference", - sketch = CMS_HEAP_BENCHMARK, - depth, - width, - heap_size, - "cms-with-heap atomic cost measurement" - ); - tracing::warn!( - sketch = CMS_HEAP_BENCHMARK, - depth, - width, - heap_size, - "invalid CMS-with-heap reference cost for candidate; dropping candidate; \ - TODO(#651): fail loudly once sketch-bench sweeps cover this grid" - ); - return None; - } - - let costs = AtomicCosts { - mem_bytes_per_instance: depth as f64 * width as f64 * assumptions.counter_bytes - + heap_size_f64 - * (assumptions.average_key_bytes + assumptions.heap_entry_overhead_bytes), - insert_cpu_secs: entry.insert_cpu_secs * scale, - merge_cpu_secs: entry.merge_cpu_secs * scale, - subtract_cpu_secs: SUBTRACT_CPU_SECS, - query_cpu_secs: entry.query_cpu_secs * scale, - exact_query_cpu_secs: EXACT_QUERY_CPU_SECS, - }; - - tracing::info!( - status = "modeled", - sketch = CMS_HEAP_BENCHMARK, - depth, - width, - heap_size, - mem_bytes_per_instance = costs.mem_bytes_per_instance, - insert_cpu_secs = costs.insert_cpu_secs, - merge_cpu_secs = costs.merge_cpu_secs, - query_cpu_secs = costs.query_cpu_secs, - "cms-with-heap atomic cost measurement" - ); - - Some(costs) -} - -fn require_u64(params: &HashMap, key: &str, agg_type: AggregationType) -> u64 { - require(params, key, agg_type) - .as_u64() - .unwrap_or_else(|| panic!("{agg_type:?} candidate has non-integer \"{key}\" param")) -} - -fn require<'a>( - params: &'a HashMap, - key: &str, - agg_type: AggregationType, -) -> &'a Value { - params.get(key).unwrap_or_else(|| { - panic!( - "{agg_type:?} candidate has no \"{key}\" param; candidate_gen.rs's param_grid() \ - has drifted" - ) - }) -} - /// Load the flat cost table `export_rqe_optimizer_costs.sh` writes, for the -/// MILP. Unlike the greedy loader, an invalid row is an error, not dropped. +/// MILP. An invalid row is an error, not dropped. pub fn load_flat_atomic_cost_table(path: &Path) -> anyhow::Result { let raw = std::fs::read_to_string(path) .map_err(|e| anyhow::anyhow!("reading cost table {}: {e}", path.display()))?; @@ -469,108 +39,8 @@ fn valid_cost_entry(entry: &AtomicCostEntry) -> bool { #[cfg(test)] mod tests { - use std::collections::BTreeMap; - use super::*; - #[test] - fn loader_selects_only_the_requested_external_profile() { - let requested = WorkloadDescription::External(ExternalWorkload { - source: "google".into(), - dataset: "google/task_usage.csv.gz".into(), - mode: "grouped".into(), - key_columns: vec![], - group_columns: vec!["machine_id".into()], - variate: None, - value_column: "cpu_rate".into(), - window_start_ns: 10, - window_end_ns: 20, - records_loaded: 100, - source_timestamp_unit: "microseconds".into(), - timestamp_unit: "nanoseconds".into(), - }); - let other = WorkloadDescription::External(ExternalWorkload { - window_end_ns: 30, - ..match requested.clone() { - WorkloadDescription::External(workload) => workload, - WorkloadDescription::Synthetic { .. } => unreachable!(), - } - }); - let document = serde_json::json!({ - "schema_version": 1, - "profiles": [ - {"workload": other, "entries": []}, - {"workload": requested, "entries": [{ - "sketch": "kll-percall", - "sketch_config": {"algorithm": "kll-percall", "params": {"k": 200}}, - "mem_bytes_per_instance": 6400.0, - "insert_cpu_secs": 1e-8, - "merge_cpu_secs": 1e-3, - "query_cpu_secs": 1e-4, - "query_accuracy": {"mean_rank_err": 0.01} - }]} - ] - }); - let file = tempfile::NamedTempFile::new().unwrap(); - std::fs::write(file.path(), document.to_string()).unwrap(); - - let table = load_atomic_cost_table(file.path(), &requested).unwrap(); - - assert_eq!(table.len(), 1); - assert_eq!(table[0].sketch, "kll-percall"); - } - - #[test] - fn loader_rejects_an_ambiguous_or_incompatible_document() { - let workload = WorkloadDescription::Synthetic { - description: serde_json::json!({"name": "one"}), - }; - let file = tempfile::NamedTempFile::new().unwrap(); - - std::fs::write( - file.path(), - serde_json::json!({ - "schema_version": 2, - "profiles": [] - }) - .to_string(), - ) - .unwrap(); - let err = load_atomic_cost_table(file.path(), &workload).unwrap_err(); - assert!(err.to_string().contains("schema_version 2")); - assert!(err.to_string().contains("supports 1")); - - std::fs::write( - file.path(), - serde_json::json!({ - "schema_version": 1, - "profiles": [ - {"workload": workload, "entries": []}, - {"workload": workload, "entries": []} - ] - }) - .to_string(), - ) - .unwrap(); - let err = load_atomic_cost_table(file.path(), &workload).unwrap_err(); - assert!(err.to_string().contains("2 atomic-cost profiles")); - - std::fs::write( - file.path(), - serde_json::json!({ - "schema_version": 1, - "profiles": [{ - "workload": {"synthetic": {"description": {"name": "other"}}}, - "entries": [] - }] - }) - .to_string(), - ) - .unwrap(); - let err = load_atomic_cost_table(file.path(), &workload).unwrap_err(); - assert!(err.to_string().contains("no atomic-cost profile")); - } - #[test] fn atomic_cost_entry_requires_query_accuracy() { let json = r#"{"sketch":"kll-percall","sketch_config":null,"mem_bytes_per_instance":1.0,"insert_cpu_secs":1.0,"merge_cpu_secs":1.0,"query_cpu_secs":1.0}"#; @@ -604,243 +74,4 @@ mod tests { let without: AtomicCostEntry = serde_json::from_str(&format!("{{{base}}}")).unwrap(); assert!(without.measured_at.is_none()); } - - #[test] - fn optional_loader_rejects_an_unselected_document() { - let document = tempfile::NamedTempFile::new().unwrap(); - let err = - load_optional_selected_atomic_cost_table(Some(document.path()), None).unwrap_err(); - assert!(err - .to_string() - .contains("--atomic-cost-workload is required")); - } - - fn cms_entry(depth: i64, width: i64) -> AtomicCostEntry { - AtomicCostEntry { - sketch: "cms-fastpath-vector2d".into(), - sketch_config: serde_json::json!({ - "algorithm": "cms-fastpath-vector2d", - "params": { "cols": width, "rows": depth } - }), - mem_bytes_per_instance: (depth * width * 4) as f64, - insert_cpu_secs: 8e-9, - merge_cpu_secs: 4.5e-4, - query_cpu_secs: 7.8e-8, - query_accuracy: BTreeMap::new(), - merge_accuracy: BTreeMap::new(), - measured_at: None, - } - } - - fn cms_params(depth: u64, width: u64) -> HashMap { - HashMap::from([ - ("depth".to_string(), Value::from(depth)), - ("width".to_string(), Value::from(width)), - ]) - } - - fn cms_heap_entry(depth: u64, width: u64) -> AtomicCostEntry { - AtomicCostEntry { - sketch: CMS_HEAP_BENCHMARK.into(), - sketch_config: serde_json::json!({ - "algorithm": CMS_HEAP_BENCHMARK, - "params": { "rows": depth, "cols": width } - }), - mem_bytes_per_instance: 1.0, - insert_cpu_secs: 2.0, - merge_cpu_secs: 4.0, - query_cpu_secs: 8.0, - query_accuracy: BTreeMap::new(), - merge_accuracy: BTreeMap::new(), - measured_at: None, - } - } - - fn cms_heap_params(depth: u64, width: u64, heap_size: u64) -> HashMap { - HashMap::from([ - ("depth".to_string(), Value::from(depth)), - ("width".to_string(), Value::from(width)), - ("heapsize".to_string(), Value::from(heap_size)), - ]) - } - - #[test] - fn cms_candidate_resolves_by_exact_key_regardless_of_value_type() { - // ASAPQuery's grid stores depth/width as u64; sketch-bench's exported - // JSON round-trips CLI-parsed integers as i64. The lookup must not - // care which Rust integer type produced the JSON number. - let table = vec![cms_entry(3, 1024), cms_entry(5, 2048)]; - let costs = resolve_atomic_costs( - &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); - assert_eq!(costs.insert_cpu_secs, 8e-9); - // Not from the table -- sketch-bench has neither, so these stay stub. - assert_eq!(costs.subtract_cpu_secs, SUBTRACT_CPU_SECS); - assert_eq!(costs.exact_query_cpu_secs, EXACT_QUERY_CPU_SECS); - } - - #[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), - 0 - ) - .is_none()); - } - - #[test] - fn unbenchmarked_family_falls_back_to_the_stub() { - let table: AtomicCostTable = vec![]; - 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!( - 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] - fn cms_with_heap_without_reference_cost_drops_the_candidate() { - // Until sketch-bench has a matching reference row, CMS-with-heap must - // not inherit the flat stub: the optimizer should retain EXACT as its - // visible fallback instead of silently selecting an uncosted sketch. - let table: AtomicCostTable = vec![]; - let params = cms_heap_params(3, 1024, 40); - assert!( - resolve_atomic_costs(&table, AggregationType::CountMinSketchWithHeap, ¶ms, 0) - .is_none() - ); - } - - #[test] - fn cms_with_heap_scales_cpu_and_models_runtime_memory() { - let table = vec![cms_heap_entry(3, 1024)]; - let assumptions = CmsHeapCostAssumptions { - reference_heap_size: 32, - counter_bytes: 8.0, - average_key_bytes: 10.0, - heap_entry_overhead_bytes: 6.0, - }; - let params = cms_heap_params(3, 1024, 64); - let costs = resolve_cms_heap_costs(&table, ¶ms, &assumptions) - .expect("matching CMS-with-heap reference row must resolve"); - - assert_eq!( - costs.mem_bytes_per_instance, - 3.0 * 1024.0 * 8.0 + 64.0 * 16.0 - ); - assert_eq!(costs.insert_cpu_secs, 4.0); - assert_eq!(costs.merge_cpu_secs, 8.0); - assert_eq!(costs.query_cpu_secs, 16.0); - assert_eq!(costs.subtract_cpu_secs, SUBTRACT_CPU_SECS); - assert_eq!(costs.exact_query_cpu_secs, EXACT_QUERY_CPU_SECS); - } - - #[test] - 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, 0) - .expect("public resolver must dispatch CMS-with-heap candidates"); - - assert_eq!( - costs.mem_bytes_per_instance, - 3.0 * 1024.0 * 8.0 + 32.0 * 64.0 - ); - assert_eq!(costs.insert_cpu_secs, 2.0); - assert_eq!(costs.merge_cpu_secs, 4.0); - assert_eq!(costs.query_cpu_secs, 8.0); - } - - #[test] - #[should_panic(expected = "reference_heap_size must be greater than zero")] - fn cms_with_heap_rejects_invalid_assumptions() { - let table = vec![cms_heap_entry(3, 1024)]; - let params = cms_heap_params(3, 1024, 40); - let assumptions = CmsHeapCostAssumptions { - reference_heap_size: 0, - ..CmsHeapCostAssumptions::default() - }; - resolve_cms_heap_costs(&table, ¶ms, &assumptions); - } - - #[test] - #[should_panic(expected = "has no \"depth\" param")] - fn cms_candidate_missing_its_expected_param_panics_instead_of_silently_stubbing() { - // A CountMinSketch candidate only ever comes from candidate_gen.rs's - // param_grid(), which always sets "depth"/"width". Landing here without - // one means sketch_bench_key's field names have drifted from - // param_grid()'s -- a real bug, not "this family has no data" (which - // the CMS-with-heap resolver handles separately and must stay visibly - // 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, 0); - } - - #[test] - fn hll_and_kll_translate_and_resolve() { - let hll_table = vec![AtomicCostEntry { - sketch: "hll".into(), - sketch_config: serde_json::json!({"algorithm": "hll", "params": {"lg_k": 14}}), - mem_bytes_per_instance: 16384.0, - insert_cpu_secs: 1.68e-9, - merge_cpu_secs: 2.76e-4, - query_cpu_secs: 1.23e-4, - query_accuracy: BTreeMap::new(), - merge_accuracy: BTreeMap::new(), - measured_at: None, - }]; - let hll_params = HashMap::from([("precision".to_string(), Value::from(14u64))]); - assert!(resolve_atomic_costs(&hll_table, AggregationType::HLL, &hll_params, 0).is_some()); - - let kll_table = vec![AtomicCostEntry { - sketch: "kll-percall".into(), - sketch_config: serde_json::json!({"algorithm": "kll-percall", "params": {"k": 200}}), - mem_bytes_per_instance: 6400.0, - insert_cpu_secs: 1.6e-8, - merge_cpu_secs: 1.0e-3, - query_cpu_secs: 1.6e-4, - query_accuracy: BTreeMap::new(), - merge_accuracy: BTreeMap::new(), - measured_at: None, - }]; - let kll_params = HashMap::from([("K".to_string(), Value::from(200u64))]); - assert!( - 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 deleted file mode 100644 index 4f2eb07c..00000000 --- a/asap-planner-rs/src/optimizer/candidate_gen.rs +++ /dev/null @@ -1,807 +0,0 @@ -use std::collections::HashMap; - -use asap_types::aggregation_config::AggregationConfig; -use asap_types::capability_matching::{compatible_agg_types, key_agg_window_valid}; -use asap_types::enums::WindowType; -use promql_utilities::data_model::KeyByLabelNames; -use promql_utilities::query_logics::enums::{AggregationType, Statistic}; -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::{OptimizerItem, QueryMethod}; -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, - // No parameter grid or cost rows yet. - AggregationType::DDSketch, -]; - -/// A candidate streaming config for one optimizer item, ready for cost evaluation. -#[derive(Debug, Clone)] -pub struct CandidateConfig { - /// A streaming config; candidates without one are no longer generated. - pub config: Option, - /// Query method derived from (ingest type × W vs range_a × sketch algebra). - 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, - /// 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 item'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 structurally valid candidate configs for an optimizer item. -/// -/// Iterates over compatible agg types × parameter grid × valid window sizes × -/// {Tumbling, Sliding}. Multi-statistic items yield no candidates because a -/// single sketch cannot serve incompatible statistics simultaneously. -pub fn enumerate_candidates(item: &OptimizerItem, scrape_interval_ms: u64) -> Vec { - let one_group = ItemFacts { - output_group_count: 1, - topk_by_group_count: item.requirements.topk_by_labels.as_ref().map(|_| 1), - arrival_rate_per_sec: 1.0, - }; - enumerate_candidates_with_facts(item, scrape_interval_ms, &one_group) -} - -/// Enumerate candidates, stamping group counts from the item's label-set facts. -pub fn enumerate_candidates_with_facts( - item: &OptimizerItem, - scrape_interval_ms: u64, - facts: &ItemFacts, -) -> Vec { - assert!( - facts.output_group_count > 0, - "output_group_count must be greater than zero" - ); - let mut candidates = Vec::new(); - - if item.requirements.statistics.len() != 1 { - return candidates; - } - - let stat = item.requirements.statistics[0]; - let range_a_ms = item.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 - // (#670), so unlike other types it can vary independently of the sketch's - // dimension params -- enumerate both weightings when the query doesn't pin one. - let sub_type_variants: Vec = if agg_type == AggregationType::CountMinSketchWithHeap - { - match item.requirements.topk_count_events { - Some(true) => vec!["count".to_string()], - Some(false) => vec!["sum".to_string()], - None => vec!["count".to_string(), "sum".to_string()], - } - } else { - vec![derive_sub_type(stat, agg_type)] - }; - - for sub_type in &sub_type_variants { - for params in param_grid(agg_type) { - for (window_type, w, slide_interval, n) in - window_candidates(range_a_ms, item.t_repeat_ms, scrape_interval_ms) - { - // DeltaSetAggregator only tracks added/removed keys since the - // last window, so it's only correct for non-overlapping - // (tumbling) windows (#588) -- same invariant enforced by - // capability_matching's window_compatible() at query time. - if !key_agg_window_valid(agg_type, window_type) { - continue; - } - - let Some(qm) = determine_query_method(n, &props) else { - continue; - }; - - let config = build_config( - item, - stat, - agg_type, - sub_type, - ¶ms, - window_type, - w, - slide_interval, - n, - ); - candidates.push(CandidateConfig { - instance_count: instance_count(&config, item, 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, - output_group_count: facts.output_group_count, - }); - } - } - } - } - - candidates -} - -/// 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 item's output labels, or (for `topk by`) the bucketing labels. -fn instance_count(config: &AggregationConfig, item: &OptimizerItem, facts: &ItemFacts) -> u64 { - if config.grouping_labels.is_empty() { - 1 - } else if config.grouping_labels == item.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`") - } -} - -/// Window candidates: (WindowType, W_ms, slide_interval_ms, n_windows). -/// -/// Every candidate satisfies `L % W == 0`, `T % S == 0`, and `W % S == 0`. -/// Tumbling uses `S = W`; sliding uses `S < W`. Both dimensions remain aligned -/// to the scrape interval. -fn window_candidates( - range_a_ms: u64, - t_repeat_ms: u64, - scrape_interval_ms: u64, -) -> Vec<(WindowType, u64, u64, u64)> { - let range_a = range_a_ms; - if range_a == 0 || scrape_interval_ms == 0 { - return vec![]; - } - - let mut out = Vec::new(); - - let tumbling_divisor = super::aqe_extractor::gcd(range_a, t_repeat_ms); - let mut w = scrape_interval_ms; - while w <= tumbling_divisor { - if tumbling_divisor.is_multiple_of(w) { - let n = range_a / w; - out.push((WindowType::Tumbling, w, w, n)); - } - w += scrape_interval_ms; - } - - let mut w = scrape_interval_ms; - while w <= range_a { - if range_a.is_multiple_of(w) { - let k = range_a / w; - let slide_divisor = super::aqe_extractor::gcd(w, t_repeat_ms); - let mut s = scrape_interval_ms; - while s < w { - if slide_divisor.is_multiple_of(s) { - out.push((WindowType::Sliding, w, s, k)); - } - s += scrape_interval_ms; - } - } - w += scrape_interval_ms; - } - - out -} - -/// Determine query method from (n_windows, sketch algebra). -/// Returns None when the combination is infeasible (W < range_a + neither merge nor subtract). -fn determine_query_method( - n_windows: u64, - props: &super::sketch_properties::SketchProperties, -) -> Option { - if n_windows == 1 { - // W = range_a (or spatial-only): one completed window covers the query range exactly. - return Some(QueryMethod::Direct); - } - // n > 1: partial-width windows (W < range_a); valid for both Tumbling and Sliding. - if props.subtractable { - Some(QueryMethod::Subtract) - } else if props.mergeable { - Some(QueryMethod::Merge { - num_windows: n_windows, - }) - } else { - None - } -} - -/// Build an AggregationConfig from candidate parameters. aggregation_id = 0 is a -/// placeholder never used past cost evaluation — OptimizerSolution::register_config -/// overwrites it with a real id when (if) a solver deploys this candidate. -#[allow(clippy::too_many_arguments)] -fn build_config( - item: &OptimizerItem, - stat: Statistic, - agg_type: AggregationType, - sub_type: &str, - params: &HashMap, - window_type: WindowType, - w: u64, - 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, - &item.requirements.grouping_labels, - item.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(), - grouping, - aggregated, - KeyByLabelNames::empty(), // rollup_labels - String::new(), // original_yaml - w, - slide_interval, - window_type, - item.requirements.spatial_filter_normalized.clone(), - item.requirements.metric.clone(), - Some(n_windows), - None, // read_count_threshold - None, // table_name (SQL only) - None, // value_column (SQL only) - ) -} - -/// aggregation_sub_type string expected by the streaming engine and capability matching. -/// Not called for `CountMinSketchWithHeap` -- its sub_type carries SUM/COUNT -/// weighting, enumerated separately in the caller (#670). -fn derive_sub_type(stat: Statistic, agg_type: AggregationType) -> String { - match (stat, agg_type) { - (Statistic::Min, _) => "min", - (Statistic::Max, _) => "max", - (Statistic::Sum, AggregationType::CountMinSketch | AggregationType::MultipleSum) => "sum", - (Statistic::Count, AggregationType::CountMinSketch) => "count", - _ => "", - } - .to_string() -} - -fn param_grid(agg_type: AggregationType) -> Vec> { - match agg_type { - AggregationType::CountMinSketch => { - let mut grids = Vec::new(); - for &d in CMS_DEPTHS { - for &w in CMS_WIDTHS { - let mut m = HashMap::new(); - m.insert("depth".into(), Value::from(d)); - m.insert("width".into(), Value::from(w)); - grids.push(m); - } - } - grids - } - - AggregationType::CountMinSketchWithHeap => { - let mut grids = Vec::new(); - for &d in CMS_DEPTHS { - for &w in CMS_WIDTHS { - for &h in CMS_HEAP_SIZES { - let mut m = HashMap::new(); - m.insert("depth".into(), Value::from(d)); - m.insert("width".into(), Value::from(w)); - m.insert("heapsize".into(), Value::from(h)); - grids.push(m); - } - } - } - grids - } - - AggregationType::DatasketchesKLL => KLL_KS - .iter() - .map(|&k| { - let mut m = HashMap::new(); - m.insert("K".into(), Value::from(k)); - m - }) - .collect(), - - AggregationType::HydraKLL => { - let mut grids = Vec::new(); - for &r in HYDRA_ROWS { - for &c in HYDRA_COLS { - let mut m = HashMap::new(); - m.insert("row_num".into(), Value::from(r)); - m.insert("col_num".into(), Value::from(c)); - m.insert("k".into(), Value::from(HYDRA_K)); - grids.push(m); - } - } - grids - } - - AggregationType::HLL => HLL_PRECISIONS - .iter() - .map(|&p| { - let mut m = HashMap::new(); - m.insert("precision".into(), Value::from(p)); - m - }) - .collect(), - - // Parameterless types: one empty-params entry per type. - _ => vec![HashMap::new()], - } -} - -#[cfg(test)] -mod tests { - use super::*; - use asap_types::enums::WindowType; - use promql_utilities::data_model::KeyByLabelNames; - - fn make_aqe(stat: Statistic, range_ms: u64, min_t: u64) -> OptimizerItem { - use asap_types::query_requirements::QueryRequirements; - OptimizerItem { - requirements: QueryRequirements { - metric: "test_metric".into(), - statistics: vec![stat], - data_range_ms: range_ms, - grouping_labels: KeyByLabelNames::empty(), - spatial_filter_normalized: String::new(), - topk_count_events: None, - topk_by_labels: None, - }, - query_strings: vec!["test_query".into()], - query_frequency_hz: 1.0 / 60.0, - occurrences: 1, - t_repeat_ms: min_t, - accuracy_sla: 0.0, - latency_sla_ms: None, - } - } - - #[test] - fn does_not_include_exact_fallback() { - let aqe = make_aqe(Statistic::Sum, 300_000, 60_000); - let candidates = enumerate_candidates(&aqe, 15_000); - assert!(candidates.iter().all(|c| c.config.is_some())); - } - - #[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) -- - // the optimizer must derive "sum" for it, same as it already does for - // CountMinSketch, or every MultipleSum candidate fails to ingest. - let aqe = make_aqe(Statistic::Sum, 300_000, 60_000); - let candidates = enumerate_candidates(&aqe, 15_000); - let multiple_sum_configs: Vec<_> = candidates - .iter() - .filter_map(|c| c.config.as_ref()) - .filter(|cfg| cfg.aggregation_type == AggregationType::MultipleSum) - .collect(); - assert!( - !multiple_sum_configs.is_empty(), - "expected at least one MultipleSum candidate" - ); - for cfg in multiple_sum_configs { - assert_eq!(cfg.aggregation_sub_type, "sum"); - } - } - - 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 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:?}"), - } - } - } - } - - #[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() - .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] - fn spatial_only_produces_direct_candidates() { - // Spatial-only: range = scrape_interval (set by extract_aqes). One Direct window. - let aqe = make_aqe(Statistic::Sum, 15_000, 60_000); - let candidates = enumerate_candidates(&aqe, 15_000); - for c in candidates.iter().filter(|c| c.config.is_some()) { - assert_eq!(c.query_method, QueryMethod::Direct); - } - } - - #[test] - fn tumbling_w_equals_range_produces_neither() { - // range_a = 60_000ms, scrape = 60_000ms → only W=60_000, n=1 → Direct - let aqe = make_aqe(Statistic::Sum, 60_000, 60_000); - let candidates = enumerate_candidates(&aqe, 60_000); - for c in candidates.iter().filter(|c| c.config.is_some()) { - assert_eq!(c.query_method, QueryMethod::Direct); - } - } - - #[test] - fn mergeable_sketch_with_multiple_windows_produces_merge() { - // Min → MinMax (mergeable, not subtractable): range_a=300_000ms, scrape=60_000ms, min_t=300_000ms - // → W=60_000 → n=5, Merge{5}. (Sum would prefer Subtract since it's also subtractable.) - let aqe = make_aqe(Statistic::Min, 300_000, 300_000); - let candidates = enumerate_candidates(&aqe, 60_000); - let merge_candidates: Vec<_> = candidates - .iter() - .filter(|c| matches!(c.query_method, QueryMethod::Merge { .. })) - .collect(); - assert!( - !merge_candidates.is_empty(), - "expected at least one Merge candidate" - ); - } - - #[test] - fn cms_with_heap_only_neither_no_merge() { - // CMS+Heap is neither mergeable nor subtractable → only n=1 (Direct) valid. - let aqe = make_aqe(Statistic::Topk, 300_000, 300_000); - let candidates = enumerate_candidates(&aqe, 60_000); - for c in candidates.iter().filter(|c| c.config.is_some()) { - assert_eq!( - c.query_method, - QueryMethod::Direct, - "CMS+Heap should only produce Direct candidates" - ); - } - } - - #[test] - fn partial_width_sliding_candidates_are_generated() { - // range_a=600_000ms, min_t=30_000ms, scrape=30_000ms. - // W=300_000 (k=2) with S=30_000 should be emitted alongside the full-width W=600_000. - let aqe = make_aqe(Statistic::Min, 600_000, 30_000); - let candidates = enumerate_candidates(&aqe, 30_000); - let partial = candidates.iter().find(|c| { - c.config.as_ref().is_some_and(|cfg| { - cfg.window_type == WindowType::Sliding - && cfg.window_size_ms == 300_000 - && c.n_windows == 2 - }) - }); - assert!( - partial.is_some(), - "expected a partial-width Sliding candidate with W=300_000ms, k=2" - ); - assert!( - matches!( - partial.unwrap().query_method, - QueryMethod::Merge { num_windows: 2 } - ), - "partial Sliding with a mergeable-only sketch should produce Merge{{2}}" - ); - } - - #[test] - fn sliding_full_width_direct_generated_when_range_exceeds_t_repeat() { - // range_a=600_000ms > T=30_000ms. W=range_a is valid for sliding since - // freshness is governed by S (not W). S=30_000 | gcd(600_000, 30_000)=30_000 → emitted. - let aqe = make_aqe(Statistic::Sum, 600_000, 30_000); - let candidates = enumerate_candidates(&aqe, 30_000); - assert!( - candidates.iter().any(|c| { - c.config.as_ref().is_some_and(|cfg| { - cfg.window_type == WindowType::Sliding && cfg.window_size_ms == 600_000 - }) && c.query_method == QueryMethod::Direct - && c.n_windows == 1 - }), - "full-width Sliding Direct should be generated even when range_a > T" - ); - } - - #[test] - fn sliding_slide_must_divide_gcd_of_window_and_t_repeat() { - // range_a=20_000, T=5_000, scrape=1_000. - // W=10_000 (k=2): slide_divisor = gcd(10_000, 5_000) = 5_000. - // Valid S: divisors of 5_000 that are multiples of 1_000 and < 10_000 → {1_000, 5_000}. - // Invalid: S=2_000 (5_000 % 2_000 ≠ 0), S=4_000 (5_000 % 4_000 ≠ 0). - let aqe = make_aqe(Statistic::Sum, 20_000, 5_000); - let candidates = enumerate_candidates(&aqe, 1_000); - - let sliding_w10: Vec<_> = candidates - .iter() - .filter(|c| { - c.config.as_ref().is_some_and(|cfg| { - cfg.window_type == WindowType::Sliding && cfg.window_size_ms == 10_000 - }) - }) - .collect(); - - let slides: Vec = sliding_w10 - .iter() - .map(|c| c.config.as_ref().unwrap().slide_interval_ms) - .collect(); - - assert!( - slides.contains(&1_000), - "S=1_000 should be valid (divides 5_000)" - ); - assert!( - slides.contains(&5_000), - "S=5_000 should be valid (divides 5_000)" - ); - assert!( - !slides.contains(&2_000), - "S=2_000 should be rejected (5_000 % 2_000 ≠ 0)" - ); - assert!( - !slides.contains(&4_000), - "S=4_000 should be rejected (5_000 % 4_000 ≠ 0)" - ); - } - - #[test] - fn sliding_slide_must_divide_window_size() { - // W=6_000, T=6_000: slide_divisor = gcd(6_000, 6_000) = 6_000. - // S=4_000: 6_000 % 4_000 = 2_000 ≠ 0 → rejected even though 4_000 < 6_000. - // S=2_000: 6_000 % 2_000 = 0 → valid. - let aqe = make_aqe(Statistic::Sum, 12_000, 6_000); - let candidates = enumerate_candidates(&aqe, 1_000); - - let slides_w6: Vec = candidates - .iter() - .filter(|c| { - c.config.as_ref().is_some_and(|cfg| { - cfg.window_type == WindowType::Sliding && cfg.window_size_ms == 6_000 - }) - }) - .map(|c| c.config.as_ref().unwrap().slide_interval_ms) - .collect(); - - assert!( - slides_w6.contains(&2_000), - "S=2_000 should be valid (6_000 % 2_000 = 0)" - ); - assert!( - !slides_w6.contains(&4_000), - "S=4_000 should be rejected (6_000 % 4_000 ≠ 0)" - ); - } - - #[test] - fn partial_sliding_subtractable_sketch_gets_subtract() { - // Sum → CMS (subtractable): partial Sliding with k=2 should produce Subtract. - let aqe = make_aqe(Statistic::Sum, 600_000, 30_000); - let candidates = enumerate_candidates(&aqe, 30_000); - assert!( - candidates.iter().any(|c| { - c.config - .as_ref() - .is_some_and(|cfg| cfg.window_type == WindowType::Sliding && c.n_windows == 2) - && c.query_method == QueryMethod::Subtract - }), - "partial Sliding with subtractable sketch should produce Subtract" - ); - } - - /// Issue #588: DeltaSetAggregator only tracks added/removed keys since - /// the last window, so it's only correct for non-overlapping (tumbling) - /// windows. `Statistic::Cardinality` is compatible with DeltaSetAggregator - /// (see `compatible_agg_types`), and with the same range/t_repeat/scrape - /// parameters as `partial_width_sliding_candidates_are_generated` above, - /// the window-candidate grid does include Sliding entries -- the - /// optimizer must never turn one of those into a DeltaSetAggregator - /// candidate. - #[test] - fn delta_set_aggregator_never_gets_a_sliding_candidate() { - let aqe = make_aqe(Statistic::Cardinality, 600_000, 30_000); - let candidates = enumerate_candidates(&aqe, 30_000); - - let sliding_delta = candidates.iter().find(|c| { - c.config.as_ref().is_some_and(|cfg| { - cfg.aggregation_type == AggregationType::DeltaSetAggregator - && cfg.window_type == WindowType::Sliding - }) - }); - - assert!( - sliding_delta.is_none(), - "DeltaSetAggregator must never be enumerated as a Sliding candidate, found: {sliding_delta:?}" - ); - } - - /// Regression guard: DeltaSetAggregator must still be enumerated as a - /// Tumbling candidate (the fix filters Sliding, not the agg type itself). - #[test] - fn delta_set_aggregator_still_gets_tumbling_candidates() { - let aqe = make_aqe(Statistic::Cardinality, 600_000, 30_000); - let candidates = enumerate_candidates(&aqe, 30_000); - - assert!( - candidates.iter().any(|c| { - c.config.as_ref().is_some_and(|cfg| { - cfg.aggregation_type == AggregationType::DeltaSetAggregator - && cfg.window_type == WindowType::Tumbling - }) - }), - "DeltaSetAggregator should still get Tumbling candidates" - ); - } - - /// Regression guard: sibling Cardinality-compatible agg types that ARE - /// safe under Sliding windows (SetAggregator, HLL) must be unaffected by - /// the DeltaSetAggregator-specific filter. - #[test] - fn set_aggregator_and_hll_still_get_sliding_candidates() { - let aqe = make_aqe(Statistic::Cardinality, 600_000, 30_000); - let candidates = enumerate_candidates(&aqe, 30_000); - - for agg_type in [AggregationType::SetAggregator, AggregationType::HLL] { - assert!( - candidates.iter().any(|c| { - c.config.as_ref().is_some_and(|cfg| { - cfg.aggregation_type == agg_type && cfg.window_type == WindowType::Sliding - }) - }), - "{agg_type:?} should still be able to get Sliding candidates" - ); - } - } -} diff --git a/asap-planner-rs/src/optimizer/constants.rs b/asap-planner-rs/src/optimizer/constants.rs deleted file mode 100644 index 02c3ab71..00000000 --- a/asap-planner-rs/src/optimizer/constants.rs +++ /dev/null @@ -1,58 +0,0 @@ -//! Tunable constants for candidate generation and cost modeling, centralized -//! so they're easy to find and swap for calibrated/profiled values later. - -// ponytail: small representative grids; replace with sketch-bench sweep results in Phase 3. -pub const CMS_DEPTHS: &[u64] = &[3, 5]; -pub const CMS_WIDTHS: &[u64] = &[512, 1024, 2048]; -pub const CMS_HEAP_SIZES: &[u64] = &[40, 200, 1000]; -// Temporary CMS-with-heap cost-model assumptions. The reference CPU costs -// come from the fixed-top-k sketch-bench wrapper; these values describe the -// runtime representation until sketch-bench sweeps the actual heap sizes. -pub const CMS_HEAP_REFERENCE_HEAP_SIZE: u64 = 32; -pub const CMS_HEAP_COUNTER_BYTES: f64 = 8.0; -pub const CMS_HEAP_AVERAGE_KEY_BYTES: f64 = 32.0; -pub const CMS_HEAP_ENTRY_OVERHEAD_BYTES: f64 = 32.0; -pub const KLL_KS: &[u64] = &[200, 500]; -pub const HYDRA_ROWS: &[u64] = &[3, 5]; -pub const HYDRA_COLS: &[u64] = &[512, 1024]; -pub const HYDRA_K: u64 = 20; -pub const HLL_PRECISIONS: &[u64] = &[12, 14]; - -// Per-operation costs for one sketch instance (`AtomicCosts` defaults). -// Stub values for v1 — real numbers come from sketch-bench in Phase 3. -pub const MEM_BYTES_PER_INSTANCE: f64 = 1024.0; -pub const INSERT_CPU_SECS: f64 = 1e-7; -pub const MERGE_CPU_SECS: f64 = 1e-5; -pub const SUBTRACT_CPU_SECS: f64 = 1e-6; -pub const QUERY_CPU_SECS: f64 = 1e-5; -/// Cost of one raw/exact query execution (the EXACT_a fallback's QueryCost). -/// Without this, EXACT always wins since its IngestCost and QueryCost would -/// otherwise both be zero. -pub const EXACT_QUERY_CPU_SECS: f64 = 1e-3; - -// Global objective weights (w1..w4 in the design doc), `CostWeights` defaults. -// Real calibration (from actual cloud $/byte-sec and $/cpu-sec) is punted -// post-v1; defaults reflect that RAM-held-over-time is several orders of -// magnitude cheaper per unit than CPU-time (e.g. ~$5/GB-month vs -// ~$0.04/vCPU-hour is roughly a 1e6 ratio), so memory weights are scaled -// down accordingly rather than left equal to CPU weights. -pub const INGEST_MEM_WEIGHT: f64 = 1e-9; -pub const INGEST_CPU_WEIGHT: f64 = 1.0; -pub const QUERY_MEM_WEIGHT: f64 = 1e-9; -pub const QUERY_CPU_WEIGHT: 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 deleted file mode 100644 index 006792b5..00000000 --- a/asap-planner-rs/src/optimizer/cost_model.rs +++ /dev/null @@ -1,407 +0,0 @@ -use asap_types::enums::WindowType; -use promql_utilities::query_logics::enums::{AggregationType, Statistic}; - -use super::candidate_gen::CandidateConfig; -use super::constants::{ - 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::{OptimizerItem, QueryMethod}; - -/// Per-operation costs for one sketch instance. Stub defaults for v1 — real -/// 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, - pub subtract_cpu_secs: f64, - pub query_cpu_secs: f64, - /// Cost of one raw/exact query execution (the EXACT_a fallback's QueryCost). - /// Without this, EXACT always wins since its IngestCost and QueryCost would - /// otherwise both be zero. - pub exact_query_cpu_secs: f64, -} - -impl Default for AtomicCosts { - fn default() -> Self { - Self { - mem_bytes_per_instance: MEM_BYTES_PER_INSTANCE, - insert_cpu_secs: INSERT_CPU_SECS, - merge_cpu_secs: MERGE_CPU_SECS, - subtract_cpu_secs: SUBTRACT_CPU_SECS, - query_cpu_secs: QUERY_CPU_SECS, - exact_query_cpu_secs: EXACT_QUERY_CPU_SECS, - } - } -} - -/// Global objective weights (w1..w4 in the design doc). Real calibration (from -/// actual cloud $/byte-sec and $/cpu-sec) is punted post-v1; defaults below -/// just reflect that RAM-held-over-time is several orders of magnitude -/// cheaper per unit than CPU-time (e.g. ~$5/GB-month vs ~$0.04/vCPU-hour is -/// roughly a 1e6 ratio), so memory weights are scaled down accordingly rather -/// than left equal to CPU weights. -#[derive(Debug, Clone, Copy)] -pub struct CostWeights { - pub ingest_mem: f64, - pub ingest_cpu: f64, - pub query_mem: f64, - pub query_cpu: f64, -} - -impl Default for CostWeights { - fn default() -> Self { - Self { - ingest_mem: INGEST_MEM_WEIGHT, - ingest_cpu: INGEST_CPU_WEIGHT, - query_mem: QUERY_MEM_WEIGHT, - query_cpu: QUERY_CPU_WEIGHT, - } - } -} - -/// IngestCost(g): steady-state cost rate of keeping `candidate` deployed, -/// independent of which AQEs query it (facility-location requirement). -/// `arrival_rate_hz` is the arrival rate (items/sec) for this config's metric+filter. -pub fn ingest_cost( - candidate: &CandidateConfig, - arrival_rate_hz: f64, - costs: &AtomicCosts, - weights: &CostWeights, -) -> f64 { - let Some(agg_config) = &candidate.config else { - return 0.0; // EXACT: no streaming config deployed. - }; - - 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. - let n_concurrent = match agg_config.window_type { - WindowType::Tumbling => 1.0, - WindowType::Sliding => { - (agg_config.window_size_ms as f64 / agg_config.slide_interval_ms.max(1) as f64).ceil() - } - }; - - 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 - + 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( - item: &OptimizerItem, - candidate: &CandidateConfig, - costs: &AtomicCosts, - weights: &CostWeights, -) -> f64 { - let Some(agg_config) = &candidate.config else { - return costs.exact_query_cpu_secs * weights.query_cpu; // EXACT: raw query at query time. - }; - - let units = stored_units(candidate, agg_config.aggregation_type); - let props = sketch_properties(agg_config.aggregation_type); - let read_cpu = reads_per_query(item, candidate) * costs.query_cpu_secs; - - let (cpu, mem) = match &candidate.query_method { - 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; - ( - units * merges * costs.merge_cpu_secs + read_cpu, - *num_windows as f64 * units * costs.mem_bytes_per_instance, - ) - } - QueryMethod::Subtract => { - debug_assert!(props.subtractable); - ( - 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. - }; - - weights.query_cpu * cpu + weights.query_mem * mem -} - -/// 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 { - 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(item: &OptimizerItem, candidate: &CandidateConfig) -> f64 { - if item.requirements.statistics == [Statistic::Topk] { - candidate.instance_count as f64 - } else { - candidate.output_group_count as f64 - } -} - -/// Total cost rate contributed by assigning AQE `aqe` (with frequency -/// `aqe.query_frequency_hz`) to `candidate`: IngestCost(g) + frequency * QueryCost(a,g). -/// This is the per-(a,g) term the greedy/MIP solver minimizes. -pub fn total_cost_rate( - item: &OptimizerItem, - candidate: &CandidateConfig, - arrival_rate_hz: f64, - costs: &AtomicCosts, - weights: &CostWeights, -) -> f64 { - ingest_cost(candidate, arrival_rate_hz, costs, weights) - + item.query_frequency_hz * query_cost(item, candidate, costs, weights) -} - -#[cfg(test)] -mod tests { - use super::*; - 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; - - fn make_aqe(stat: Statistic, range_ms: u64, min_t: u64) -> OptimizerItem { - OptimizerItem { - requirements: QueryRequirements { - metric: "test_metric".into(), - statistics: vec![stat], - data_range_ms: range_ms, - grouping_labels: KeyByLabelNames::empty(), - spatial_filter_normalized: String::new(), - topk_count_events: None, - topk_by_labels: None, - }, - query_strings: vec!["test_query".into()], - query_frequency_hz: 1.0 / 60.0, - occurrences: 1, - t_repeat_ms: min_t, - accuracy_sla: 0.0, - latency_sla_ms: None, - } - } - - #[test] - fn ingest_cost_independent_of_n_windows() { - // Retained windows are a transient per-query cost (Mem_query in query_cost), - // not a continuous allocation — so ingest_cost must not vary with n_windows. - let a = make_aqe(Statistic::Sum, 300_000, 300_000); - let candidates = enumerate_candidates(&a, 60_000); - let template = candidates - .iter() - .find_map(|c| c.config.clone()) - .expect("expected at least one deployed candidate"); - - let costs = AtomicCosts::default(); - let weights = CostWeights::default(); - let c2 = CandidateConfig { - config: Some(template.clone()), - query_method: QueryMethod::Merge { num_windows: 2 }, - n_windows: 2, - 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, - instance_count: 1, - output_group_count: 1, - key_config: None, - }; - - assert_eq!( - ingest_cost(&c2, 1.0, &costs, &weights), - ingest_cost(&c5, 1.0, &costs, &weights), - ); - } - - #[test] - fn subtract_is_cheaper_than_merge_for_the_same_window_count() { - // Calibration-independent: Subtract is O(1) (one subtract + one read) while - // Merge is O(n) (n-1 merges + one read), so for the same n and same - // underlying config, Subtract must cost less regardless of weight tuning. - let a = make_aqe(Statistic::Sum, 300_000, 300_000); - let candidates = enumerate_candidates(&a, 60_000); - let template = candidates - .iter() - .find_map(|c| c.config.clone()) - .expect("expected at least one deployed candidate"); - - let costs = AtomicCosts::default(); - let weights = CostWeights::default(); - let merge = CandidateConfig { - config: Some(template.clone()), - query_method: QueryMethod::Merge { num_windows: 5 }, - n_windows: 5, - instance_count: 1, - output_group_count: 1, - key_config: None, - }; - let subtract = CandidateConfig { - config: Some(template), - query_method: QueryMethod::Subtract, - n_windows: 5, - instance_count: 1, - output_group_count: 1, - key_config: None, - }; - - assert!( - query_cost(&a, &subtract, &costs, &weights) < query_cost(&a, &merge, &costs, &weights) - ); - } - - /// 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, - ) -> (OptimizerItem, 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(|c| { - c.query_method == QueryMethod::Direct - && c.config - .as_ref() - .is_some_and(|cfg| cfg.aggregation_type == agg_type) - }) - .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(); - 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!( - query_cost(&a, &one, &costs, &MEM_ONLY), - query_cost(&a, &five, &costs, &MEM_ONLY) - ); - assert_eq!( - query_cost(&a, &five, &costs, &QUERY_CPU_ONLY), - 5.0 * query_cost(&a, &one, &costs, &QUERY_CPU_ONLY) - ); - } - - #[test] - fn topk_reads_each_heap_once_regardless_of_output_groups() { - let costs = AtomicCosts::default(); - 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, &few, &costs, &QUERY_CPU_ONLY), - query_cost(&a, &many, &costs, &QUERY_CPU_ONLY) - ); - } -} diff --git a/asap-planner-rs/src/optimizer/error.rs b/asap-planner-rs/src/optimizer/error.rs index da7933b5..ee61bc8e 100644 --- a/asap-planner-rs/src/optimizer/error.rs +++ b/asap-planner-rs/src/optimizer/error.rs @@ -1,4 +1,3 @@ -use promql_utilities::query_logics::enums::Statistic; use thiserror::Error; #[derive(Debug, Error)] @@ -12,18 +11,4 @@ pub enum OptimizerError { leaf: String, reason: String, }, - - #[error("{items:?}")] - UnservableItems { items: Vec }, -} - -#[derive(Debug, Clone, PartialEq)] -pub struct UnservableItem { - pub metric: String, - pub statistics: Vec, - pub data_range_ms: u64, - pub t_repeat_ms: u64, - pub accuracy_sla: f64, - pub latency_sla_ms: Option, - pub reason: String, } diff --git a/asap-planner-rs/src/optimizer/greedy.rs b/asap-planner-rs/src/optimizer/greedy.rs deleted file mode 100644 index eaceb9a4..00000000 --- a/asap-planner-rs/src/optimizer/greedy.rs +++ /dev/null @@ -1,325 +0,0 @@ -use tracing::debug; - -use super::atomic_costs::{resolve_atomic_costs, AtomicCostTable}; -use super::candidate_gen::enumerate_candidates_with_facts; -use super::cost_model::{ingest_cost, query_cost, total_cost_rate, AtomicCosts, CostWeights}; -use super::error::{OptimizerError, UnservableItem}; -use super::label_set_facts::ItemFacts; -use super::solution::{AQEAssignment, OptimizerItem, OptimizerSolution}; - -/// Greedily assign each AQE to its independently-cheapest candidate config. -/// -/// 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; this -/// is the v1 baseline. -/// -/// `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 -/// every candidate; a candidate whose config has no matching table entry is -/// dropped from consideration (`resolve_atomic_costs` returns `None`). -pub fn greedy_assign( - aqes: Vec, - scrape_interval_ms: u64, - atomic_cost_table: &AtomicCostTable, - weights: &CostWeights, - facts: &[ItemFacts], -) -> Result { - assert_eq!( - aqes.len(), - facts.len(), - "facts must be index-aligned with aqes" - ); - let mut solution = OptimizerSolution::empty(); - let mut unservable_items = Vec::new(); - - 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_facts(&aqe, scrape_interval_ms, item_facts); - - let Some((best, costs)) = candidates - .into_iter() - .filter_map(|c| { - // EXACT (config: None) always costs at the flat stub — it has - // no sketch_type/params for the table to key on. - let costs = match &c.config { - None => AtomicCosts::default(), - Some(cfg) => resolve_atomic_costs( - 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); - Some((c, costs, cost)) - }) - // total_cmp (not partial_cmp().unwrap()) so a stray NaN cost can't panic. - .min_by(|(_, _, a), (_, _, b)| a.total_cmp(b)) - .map(|(c, costs, _)| (c, costs)) - else { - unservable_items.push(UnservableItem { - metric: aqe.requirements.metric.clone(), - statistics: aqe.requirements.statistics.clone(), - data_range_ms: aqe.requirements.data_range_ms, - t_repeat_ms: aqe.t_repeat_ms, - accuracy_sla: aqe.accuracy_sla, - latency_sla_ms: aqe.latency_sla_ms, - reason: "no candidate remained after structural and atomic-cost filters".into(), - }); - continue; - }; - - let ingest = ingest_cost(&best, arrival_rate_hz, &costs, weights); - let query_rate = aqe.query_frequency_hz * query_cost(&aqe, &best, &costs, weights); - let query_method = best.query_method.clone(); - - let aggregation_id = solution.register_config( - best.config - .expect("candidate configs are streaming configs"), - ); - let key_aggregation_id = best - .key_config - .map(|config| solution.register_config(config)); - - debug!( - metric = %aqe.requirements.metric, - aggregation_id = ?aggregation_id, - query_method = ?query_method, - ingest_cost_per_sec = ingest, - query_cost_per_sec = query_rate, - "greedy: assigned AQE" - ); - - solution.estimated_ingest_cost_per_sec += ingest; - solution.estimated_total_cost_per_sec += ingest + query_rate; - - solution.assignments.push(AQEAssignment { - item: aqe, - aggregation_id, - key_aggregation_id, - query_method, - estimated_query_cost_per_sec: query_rate, - }); - } - - if unservable_items.is_empty() { - Ok(solution) - } else { - Err(OptimizerError::UnservableItems { - items: unservable_items, - }) - } -} - -#[cfg(test)] -mod tests { - 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}; - use std::collections::HashMap as StdHashMap; - - fn make_aqe(stat: Statistic, range_ms: u64, min_t: u64, freq_hz: f64) -> OptimizerItem { - OptimizerItem { - requirements: QueryRequirements { - metric: "test_metric".into(), - statistics: vec![stat], - data_range_ms: range_ms, - grouping_labels: KeyByLabelNames::empty(), - spatial_filter_normalized: String::new(), - topk_count_events: None, - topk_by_labels: None, - }, - query_strings: vec!["test_query".into()], - query_frequency_hz: freq_hz, - occurrences: 1, - t_repeat_ms: min_t, - accuracy_sla: 0.0, - latency_sla_ms: None, - } - } - - /// One group, one item/sec for every item. - fn unit_facts(aqes: &[OptimizerItem]) -> Vec { - aqes.iter() - .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() - } - - #[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, - &AtomicCostTable::default(), - &CostWeights::default(), - &facts, - ) - .unwrap(); - - let mut seen_ids: StdHashMap = StdHashMap::new(); - for id in solution.deployed_configs().keys() { - assert!( - seen_ids.insert(*id, ()).is_none(), - "duplicate aggregation_id" - ); - } - assert_eq!(solution.assignments.len(), 2); - } - - #[test] - fn unsupported_multi_statistic_item_is_unservable() { - let aqe = OptimizerItem { - requirements: QueryRequirements { - metric: "test_metric".into(), - statistics: vec![Statistic::Sum, Statistic::Count], // avg-style, unsupported - data_range_ms: 60_000, - grouping_labels: KeyByLabelNames::empty(), - spatial_filter_normalized: String::new(), - topk_count_events: None, - topk_by_labels: None, - }, - query_strings: vec!["avg_query".into()], - query_frequency_hz: 1.0 / 60.0, - occurrences: 1, - t_repeat_ms: 60_000, - accuracy_sla: 0.0, - latency_sla_ms: None, - }; - let error = greedy_assign( - vec![aqe.clone()], - 60_000, - &AtomicCostTable::default(), - &CostWeights::default(), - &unit_facts(&[aqe]), - ); - assert!(matches!(error, Err(OptimizerError::UnservableItems { .. }))); - } - - #[test] - fn missing_cms_with_heap_reference_cost_is_unservable() { - // Regression coverage for #651: an uncosted CMS-with-heap candidate - // must be dropped rather than inheriting the flat stub and winning. - let aqe = make_aqe(Statistic::Topk, 60_000, 60_000, 1.0 / 60.0); - let error = greedy_assign( - vec![aqe.clone()], - 60_000, - &AtomicCostTable::default(), - &CostWeights::default(), - &unit_facts(&[aqe]), - ); - - assert!(matches!(error, Err(OptimizerError::UnservableItems { .. }))); - } - - #[test] - fn matching_cms_with_heap_reference_cost_can_be_deployed() { - let table = vec![AtomicCostEntry { - sketch: "cms-heap-topk-regularpath-vector2d".into(), - sketch_config: serde_json::json!({ - "algorithm": "cms-heap-topk-regularpath-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(), - merge_accuracy: std::collections::BTreeMap::new(), - measured_at: None, - }]; - let aqe = make_aqe(Statistic::Topk, 60_000, 60_000, 1.0 / 60.0); - let solution = greedy_assign( - vec![aqe.clone()], - 60_000, - &table, - &CostWeights::default(), - &unit_facts(&[aqe]), - ) - .unwrap(); - - assert_eq!(solution.deployed_configs().len(), 1); - assert_eq!( - solution - .deployed_configs() - .values() - .next() - .expect("one CMS-with-heap config deployed") - .aggregation_type, - 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(), - merge_accuracy: std::collections::BTreeMap::new(), - measured_at: None, - }]; - 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)), - ) - .unwrap(); - - let assignment = &solution.assignments[0]; - let value_id = assignment.aggregation_id; - 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 deleted file mode 100644 index 0028c291..00000000 --- a/asap-planner-rs/src/optimizer/label_set_facts.rs +++ /dev/null @@ -1,532 +0,0 @@ -//! 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::OptimizerItem; - -#[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 optimizer item, ready for costing. -#[derive(Debug, Clone, Copy, PartialEq)] -pub struct ItemFacts { - /// Distinct value combinations of the item's output (grouping) labels. - pub output_group_count: u64, - /// Distinct values of the `topk by` labels; `Some` iff the item has them. - pub topk_by_group_count: Option, - /// 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 item, index-aligned with `items`. An item 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 - /// item 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, - items: &[OptimizerItem], - scrape_interval_ms: u64, - ) -> 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(items.len()); - for item in items { - let key = LabelSetKey::from_requirements(&item.requirements); - 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: {series_key}")); - } - used_series.insert(series_key); - - let topk_by_group_count = item.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() { - missing.sort(); - missing.dedup(); - return Err(LabelSetFactsError::MissingFacts(missing)); - } - - 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 !used_groups.contains(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]) -> OptimizerItem { - OptimizerItem { - 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, - occurrences: 1, - t_repeat_ms: 60_000, - accuracy_sla: 0.0, - latency_sla_ms: None, - } - } - - 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 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. - 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/milp.rs b/asap-planner-rs/src/optimizer/milp.rs index 7269ad3d..9f2d4a9f 100644 --- a/asap-planner-rs/src/optimizer/milp.rs +++ b/asap-planner-rs/src/optimizer/milp.rs @@ -12,7 +12,8 @@ use thiserror::Error; use crate::config::input::ControllerConfig; -use super::pipeline::{extract_hinted_items, OptimizerPipelineError}; +use super::aqe_extractor::{extract_aqes, RQE}; +use super::error::OptimizerError; use super::solution::OptimizerItem; /// Slack on accuracy tolerances so `1 - sla` rounding (`1 - 0.9 = @@ -21,8 +22,14 @@ const SLA_EPSILON: f64 = 1e-9; #[derive(Debug, Error)] pub enum MilpError { + #[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(transparent)] - Pipeline(#[from] OptimizerPipelineError), + Extraction(#[from] OptimizerError), #[error("query {query:?}: accuracy_sla {accuracy_sla} must be between 0 and 1")] AccuracySlaOutOfRange { query: String, accuracy_sla: f64 }, #[error("invalid MILP inputs:\n{}", .0.join("\n"))] @@ -163,6 +170,51 @@ pub fn build_milp_workload( }) } +/// The workload's optimizer items, with labels resolved from its `metrics:` +/// hints. Every workload metric must have a hint. +fn extract_hinted_items( + config: &ControllerConfig, + scrape_interval_ms: u64, +) -> Result, MilpError> { + if config.metrics.is_none() { + return Err(MilpError::MissingMetricHints); + } + let schema = config.schema_from_hints(); + let rqes = config_to_rqes(config); + let aqes = extract_aqes(&rqes, &schema, scrape_interval_ms)?; + + // Requirement extraction treats an unknown metric as having no labels, + // which would silently mis-resolve `without (...)` and plain selectors. + let mut unhinted: Vec = aqes + .iter() + .map(|aqe| aqe.requirements.metric.clone()) + .filter(|metric| schema.get_labels(metric).is_none()) + .collect(); + if !unhinted.is_empty() { + unhinted.sort(); + unhinted.dedup(); + return Err(MilpError::MetricsWithoutHints(unhinted)); + } + Ok(aqes) +} + +/// Convert a `ControllerConfig`'s query groups into a flat list of RQEs. +/// Each (query, repetition_delay_ms) pair becomes one RQE. +fn config_to_rqes(config: &ControllerConfig) -> Vec { + config + .query_groups + .iter() + .flat_map(|qg| { + qg.queries.iter().map(|q| RQE { + query_string: q.clone(), + t_repeat_ms: qg.repetition_delay_ms, + accuracy_sla: qg.controller_options.accuracy_sla, + latency_sla_ms: qg.controller_options.latency_sla_ms, + }) + }) + .collect() +} + /// No limit sorts last. fn latency_key(item: &OptimizerItem) -> f64 { item.latency_sla_ms.unwrap_or(f64::INFINITY) @@ -537,4 +589,54 @@ metrics: let solution = solve_milp(&w, &facts, &costs(), Objective::default(), false).unwrap(); assert_eq!(solution.deployments.len(), 2); } + + #[test] + fn workload_requires_metric_hints() { + let config: ControllerConfig = serde_yaml::from_str(&format!( + "query_groups:\n{}", + group("sum(http_requests_total)", 0.99) + )) + .unwrap(); + assert!(matches!( + extract_hinted_items(&config, SCRAPE_MS), + Err(MilpError::MissingMetricHints) + )); + } + + #[test] + fn workload_metric_without_a_hint_is_rejected() { + let config = config(&format!( + "{}{}", + group("sum(http_requests_total)", 0.99), + group("max_over_time(unhinted[5m])", 0.99) + )); + assert!(matches!( + extract_hinted_items(&config, SCRAPE_MS), + Err(MilpError::MetricsWithoutHints(metrics)) if metrics == ["unhinted"] + )); + } + + #[test] + fn spatial_only_aqe_gets_the_scrape_interval_as_its_range() { + let config = config(&group("sum(http_requests_total)", 0.99)); + let items = extract_hinted_items(&config, SCRAPE_MS).unwrap(); + assert_eq!(items.len(), 1); + assert_eq!(items[0].requirements.data_range_ms, SCRAPE_MS); + } + + #[test] + fn config_to_rqes_flattens_groups() { + let group_at = |query: &str, t_repeat_ms: u64| { + format!(" - queries: [\"{query}\"]\n repetition_delay_ms: {t_repeat_ms}\n") + }; + let config = config(&format!( + "{}{}", + group_at("sum_over_time(a[5m])", 60_000), + group_at("sum_over_time(b[5m])", 30_000) + )); + let rqes = config_to_rqes(&config); + assert_eq!(rqes.len(), 2); + assert_eq!(rqes[0].t_repeat_ms, 60_000); + assert_eq!(rqes[1].t_repeat_ms, 30_000); + } } diff --git a/asap-planner-rs/src/optimizer/mod.rs b/asap-planner-rs/src/optimizer/mod.rs index f172485a..f7e1cea5 100644 --- a/asap-planner-rs/src/optimizer/mod.rs +++ b/asap-planner-rs/src/optimizer/mod.rs @@ -1,34 +1,15 @@ pub mod aqe_extractor; pub mod atomic_costs; -pub mod candidate_gen; -pub mod constants; -pub mod cost_model; pub mod error; -pub mod greedy; -pub mod label_set_facts; pub mod milp; pub mod milp_output; -pub mod pipeline; -pub mod sketch_properties; pub mod solution; -pub mod translator; pub mod workload_facts; pub use aqe_extractor::{extract_aqes, RQE}; -pub use atomic_costs::{ - load_atomic_cost_table, load_flat_atomic_cost_table, load_optional_selected_atomic_cost_table, - load_selected_atomic_cost_table, resolve_atomic_costs, AtomicCostEntry, AtomicCostTable, - ExternalWorkload, WorkloadDescription, -}; -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 error::{OptimizerError, UnservableItem}; -pub use greedy::greedy_assign; -pub use label_set_facts::{ItemFacts, LabelSetFacts, LabelSetFactsError, LabelSetKey, SeriesKey}; +pub use atomic_costs::{load_flat_atomic_cost_table, AtomicCostEntry, AtomicCostTable}; +pub use error::OptimizerError; pub use milp::{build_milp_workload, solve_milp, MilpError, MilpWorkload}; pub use milp_output::{plan_to_planner_output, reject_avg_queries, MilpOutputError}; -pub use pipeline::{run_greedy_pipeline, OptimizerPipelineError}; -pub use sketch_properties::{sketch_properties, SketchProperties}; -pub use solution::{AQEAssignment, OptimizerItem, OptimizerSolution, QueryMethod}; -pub use translator::translate; +pub use solution::OptimizerItem; pub use workload_facts::{load_workload_facts, parse_workload_facts, WorkloadFactsError}; diff --git a/asap-planner-rs/src/optimizer/pipeline.rs b/asap-planner-rs/src/optimizer/pipeline.rs deleted file mode 100644 index c54aaa93..00000000 --- a/asap-planner-rs/src/optimizer/pipeline.rs +++ /dev/null @@ -1,255 +0,0 @@ -use asap_types::inference_config::InferenceConfig; -use asap_types::streaming_config::StreamingConfig; -use thiserror::Error; - -use crate::config::input::ControllerConfig; - -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, LabelSetKey}; -use super::solution::{OptimizerItem, OptimizerSolution}; -use super::translator::{translate, TranslationSummary}; - -#[derive(Debug, Error)] -pub enum OptimizerPipelineError { - #[error(transparent)] - LabelSetFacts(#[from] LabelSetFactsError), - #[error(transparent)] - Optimizer(#[from] super::error::OptimizerError), -} - -fn finish_pipeline( - solution: OptimizerSolution, - solver_name: &str, -) -> (StreamingConfig, InferenceConfig) { - let summary = TranslationSummary::from_solution(&solution); - tracing::info!( - solver = solver_name, - num_deployed_configs = summary.num_deployed_configs, - num_sketch_assignments = summary.num_sketch_assignments, - estimated_ingest_cost_per_sec = solution.estimated_ingest_cost_per_sec, - estimated_total_cost_per_sec = solution.estimated_total_cost_per_sec, - "optimizer pipeline: solution produced" - ); - - translate(&solution) -} - -/// Run the greedy optimizer pipeline: each item is assigned independently to -/// its cheapest eligible streaming config. -/// -/// No cross-item sharing — every deployed sketch serves exactly one item, even -/// if two AQEs could share one. The Phase 3 MIP finds sharing opportunities. -/// -/// The workload's `metrics:` hints supply the label schema; `facts` supply each -/// item's group cardinality and series count. -pub fn run_greedy_pipeline( - config: &ControllerConfig, - facts: &LabelSetFacts, - scrape_interval_ms: u64, - atomic_cost_table: &AtomicCostTable, -) -> Result<(StreamingConfig, InferenceConfig), OptimizerPipelineError> { - let aqes = extract_hinted_items(config, scrape_interval_ms)?; - - let item_facts = facts.resolve(&aqes, scrape_interval_ms)?; - for (aqe, item) in aqes.iter().zip(&item_facts) { - tracing::info!( - 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" - ); - } - - let solution = greedy_assign( - aqes, - scrape_interval_ms, - atomic_cost_table, - &CostWeights::default(), - &item_facts, - )?; - - Ok(finish_pipeline(solution, "greedy")) -} - -/// The workload's optimizer items, with labels resolved from its `metrics:` -/// hints. Every workload metric must have a hint. -pub(super) fn extract_hinted_items( - config: &ControllerConfig, - scrape_interval_ms: u64, -) -> Result, OptimizerPipelineError> { - if config.metrics.is_none() { - return Err(LabelSetFactsError::MissingMetricHints.into()); - } - let schema = config.schema_from_hints(); - let rqes = config_to_rqes(config); - let aqes = extract_aqes(&rqes, &schema, scrape_interval_ms)?; - - // Requirement extraction treats an unknown metric as having no labels, - // which would silently mis-resolve `without (...)` and plain selectors. - let mut unhinted: Vec = aqes - .iter() - .map(|aqe| aqe.requirements.metric.clone()) - .filter(|metric| schema.get_labels(metric).is_none()) - .collect(); - if !unhinted.is_empty() { - unhinted.sort(); - unhinted.dedup(); - return Err(LabelSetFactsError::MetricsWithoutHints(unhinted).into()); - } - Ok(aqes) -} - -/// Convert a `ControllerConfig`'s query groups into a flat list of RQEs. -/// Each (query, repetition_delay_ms) pair becomes one RQE. -fn config_to_rqes(config: &ControllerConfig) -> Vec { - config - .query_groups - .iter() - .flat_map(|qg| { - qg.queries.iter().map(|q| RQE { - query_string: q.clone(), - t_repeat_ms: qg.repetition_delay_ms, - accuracy_sla: qg.controller_options.accuracy_sla, - latency_sla_ms: qg.controller_options.latency_sla_ms, - }) - }) - .collect() -} - -#[cfg(test)] -mod tests { - use super::*; - use asap_types::PromQLSchema; - - fn make_config(queries: &[(&str, u64)]) -> ControllerConfig { - use crate::config::input::QueryGroup; - - let query_groups = queries - .iter() - .map(|(q, t)| QueryGroup { - id: None, - queries: vec![q.to_string()], - repetition_delay_ms: *t, - controller_options: Default::default(), - step_ms: None, - range_duration_ms: None, - }) - .collect(); - - ControllerConfig { - query_groups, - windowing: None, - sketch_parameters: None, - aggregate_cleanup: None, - metrics: None, - existing_streaming_config: None, - existing_inference_config: None, - } - } - - 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 = 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, &facts, 60_000, &AtomicCostTable::default()).unwrap(); - assert!(!streaming.get_all_aggregation_configs().is_empty()); - assert!(!inference.query_configs.is_empty()); - } - - #[test] - fn greedy_pipeline_requires_metric_hints() { - let config = make_config(&[("min_over_time(metric[5m])", 60_000)]); - let facts = LabelSetFacts::from_yaml(METRIC_FACTS).unwrap(); - assert!(matches!( - run_greedy_pipeline(&config, &facts, 60_000, &AtomicCostTable::default()), - Err(OptimizerPipelineError::LabelSetFacts( - 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(OptimizerPipelineError::LabelSetFacts( - 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, &facts, 60_000, &AtomicCostTable::default()), - Err(OptimizerPipelineError::LabelSetFacts( - LabelSetFactsError::MissingFacts(_) - )) - )); - } - - #[test] - fn spatial_only_aqe_gets_explicit_range_from_pipeline() { - let config = make_config(&[("sum(metric)", 60_000)]); - let rqes = config_to_rqes(&config); - let aqes = extract_aqes(&rqes, &PromQLSchema::new(), 15_000).unwrap(); - assert_eq!(aqes.len(), 1); - assert_eq!(aqes[0].requirements.data_range_ms, 15_000); - } - - #[test] - fn config_to_rqes_flattens_groups() { - let config = make_config(&[ - ("sum_over_time(a[5m])", 60_000), - ("sum_over_time(b[5m])", 30_000), - ]); - let rqes = config_to_rqes(&config); - assert_eq!(rqes.len(), 2); - assert_eq!(rqes[0].t_repeat_ms, 60_000); - assert_eq!(rqes[1].t_repeat_ms, 30_000); - } -} diff --git a/asap-planner-rs/src/optimizer/sketch_properties.rs b/asap-planner-rs/src/optimizer/sketch_properties.rs deleted file mode 100644 index f914e3a0..00000000 --- a/asap-planner-rs/src/optimizer/sketch_properties.rs +++ /dev/null @@ -1,70 +0,0 @@ -use promql_utilities::query_logics::enums::AggregationType; - -#[derive(Debug, Clone, Copy)] -pub struct SketchProperties { - /// Two instances can be combined into one representing the union. - pub mergeable: bool, - /// Element-wise difference is defined; enables the Subtract query method (tumbling only). - pub subtractable: 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, grows| SketchProperties { - mergeable: me, - subtractable: su, - memory_grows_with_keys: grows, - }; - match t { - AggregationType::Sum => p(true, true, false), - AggregationType::Increase => p(true, false, false), - AggregationType::MinMax => p(true, false, false), - AggregationType::DatasketchesKLL => p(true, false, false), - // Bucket counts could be subtracted, but the accumulator only merges. - AggregationType::DDSketch => p(true, false, false), - AggregationType::MultipleSum => p(true, true, true), - AggregationType::MultipleIncrease => p(true, false, true), - AggregationType::MultipleMinMax => p(true, false, 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, false), - AggregationType::SetAggregator | AggregationType::DeltaSetAggregator => { - p(true, false, false) - } - AggregationType::HLL => p(true, false, false), - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn cms_is_mergeable_subtractable_fixed_size() { - let p = sketch_properties(AggregationType::CountMinSketch); - 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); - } - - #[test] - fn multiple_sum_memory_grows_with_keys() { - let p = sketch_properties(AggregationType::MultipleSum); - assert!(p.mergeable && p.subtractable && p.memory_grows_with_keys); - } - - #[test] - fn kll_mergeable_not_subtractable() { - let p = sketch_properties(AggregationType::DatasketchesKLL); - assert!(p.mergeable && !p.subtractable); - } -} diff --git a/asap-planner-rs/src/optimizer/solution.rs b/asap-planner-rs/src/optimizer/solution.rs index 52b7713c..44d79537 100644 --- a/asap-planner-rs/src/optimizer/solution.rs +++ b/asap-planner-rs/src/optimizer/solution.rs @@ -1,6 +1,3 @@ -use std::collections::HashMap; - -use asap_types::aggregation_config::AggregationConfig; use asap_types::query_requirements::QueryRequirements; /// One optimizer demand item: an AQE's requirements at one cadence and SLA pair. @@ -10,12 +7,10 @@ pub struct OptimizerItem { pub requirements: QueryRequirements, /// Original query strings from RQEs that contribute to this item. - /// Preserved for use by the translator when building InferenceConfig. + /// Used to build InferenceConfig query configs. pub query_strings: Vec, /// Query frequency in Hz: `count * 1000 / t_repeat_ms`. - /// Used in the MIP objective to convert per-query QueryCost into a cost - /// rate (cost/sec) commensurate with the continuously-accruing IngestCost. /// Represents the total query load from all dashboards independently /// hitting the sketch. pub query_frequency_hz: f64, @@ -33,156 +28,3 @@ pub struct OptimizerItem { /// `None` means no limit. pub latency_sla_ms: Option, } - -/// How an optimizer item is answered from its assigned streaming config. -/// -/// Determined by (ingest_type, W vs range_a, sketch algebra) — not a free -/// decision variable. See the compatibility table in the design doc. -#[derive(Debug, Clone, PartialEq, Eq)] -pub enum QueryMethod { - /// W = range_a: one completed window covers the query range exactly. - /// Direct read, no merge or subtract needed. - Direct, - - /// W < range_a, sketch is mergeable: combine `num_windows` retained - /// sub-windows at query time (Tumbling or partial-width Sliding). - /// Cost scales linearly with num_windows. - Merge { num_windows: u64 }, - - /// W < range_a, sketch is subtractable: subtract two prefix-sum checkpoints. - /// O(1) cost regardless of range_a/W. - Subtract, -} - -/// The assignment of a single optimizer item to a streaming config. -#[derive(Debug, Clone)] -pub struct AQEAssignment { - pub item: OptimizerItem, - - /// ID of the deployed config that serves this item. - pub aggregation_id: u64, - - /// ID of the paired key aggregation, for value sketches that can't list - /// their own keys. - pub key_aggregation_id: Option, - - /// How this item's answer is derived from the assigned config. - pub query_method: QueryMethod, - - /// Estimated cost rate for this assignment: QueryCost(a, g) * `aqe.query_frequency_hz`. - pub estimated_query_cost_per_sec: f64, -} - -/// The output of the optimizer: a complete plan for a given RQE workload. -/// -/// Contains the set of streaming configs to deploy and the assignment of every -/// item to one of those configs. A thin translator -/// converts this into `StreamingConfig + InferenceConfig` deployment artifacts. -#[derive(Debug, Clone)] -pub struct OptimizerSolution { - /// Deployed streaming configs (y_g = 1 in the MIP). Keyed by aggregation_id. - /// - /// Private: the only way to add an entry is `register_config`, which - /// assigns the id. This keeps candidate_gen.rs's placeholder id (0) from - /// ever reaching a deployed config — see ASAPQuery#564. - deployed_configs: HashMap, - - /// Next id `register_config` will hand out. - next_id: u64, - - /// One entry per optimizer item across the full RQE workload. - pub assignments: Vec, - - /// Estimated steady-state ingestion cost rate across all deployed configs - /// (Σ_{g: y_g=1} IngestCost(g)). - pub estimated_ingest_cost_per_sec: f64, - - /// Estimated total cost rate: ingest + query components combined. - pub estimated_total_cost_per_sec: f64, -} - -impl OptimizerSolution { - /// An empty solution with no assignments or deployed configs yet, ready - /// to be built up incrementally (e.g. by a solver's assignment loop). - pub fn empty() -> Self { - Self { - deployed_configs: HashMap::new(), - next_id: 1, - assignments: Vec::new(), - estimated_ingest_cost_per_sec: 0.0, - estimated_total_cost_per_sec: 0.0, - } - } - - /// Register a candidate config as deployed: assigns it a fresh unique id - /// (overwriting whatever placeholder candidate_gen.rs set), stores it, - /// and returns the id. The only way to populate `deployed_configs`. - pub fn register_config(&mut self, mut config: AggregationConfig) -> u64 { - let id = self.next_id; - self.next_id += 1; - config.aggregation_id = id; - self.deployed_configs.insert(id, config); - id - } - - pub fn deployed_configs(&self) -> &HashMap { - &self.deployed_configs - } - - /// Number of optimizer items served by a streaming sketch. - pub fn num_sketch_served(&self) -> usize { - self.assignments.len() - } -} - -#[cfg(test)] -mod tests { - use super::*; - use asap_types::enums::WindowType; - use promql_utilities::data_model::KeyByLabelNames; - use promql_utilities::query_logics::enums::AggregationType; - - fn candidate_config() -> AggregationConfig { - // aggregation_id: 0, matching candidate_gen.rs's placeholder — the - // thing register_config must always overwrite (ASAPQuery#564). - AggregationConfig::new( - 0, - AggregationType::CountMinSketch, - "sum".into(), - HashMap::new(), - KeyByLabelNames::empty(), - KeyByLabelNames::empty(), - KeyByLabelNames::empty(), - String::new(), - 60_000, - 60_000, - WindowType::Tumbling, - String::new(), - "test_metric".into(), - Some(1), - None, - None, - None, - ) - } - - #[test] - fn register_config_never_leaves_the_placeholder_id() { - let mut solution = OptimizerSolution::empty(); - let id = solution.register_config(candidate_config()); - assert_ne!( - id, 0, - "register_config must not hand out the placeholder id" - ); - assert_eq!(solution.deployed_configs()[&id].aggregation_id, id); - } - - #[test] - fn register_config_assigns_distinct_ids() { - let mut solution = OptimizerSolution::empty(); - let id1 = solution.register_config(candidate_config()); - let id2 = solution.register_config(candidate_config()); - assert_ne!(id1, id2); - assert_eq!(solution.deployed_configs().len(), 2); - } -} diff --git a/asap-planner-rs/src/optimizer/translator.rs b/asap-planner-rs/src/optimizer/translator.rs deleted file mode 100644 index 465e5d0b..00000000 --- a/asap-planner-rs/src/optimizer/translator.rs +++ /dev/null @@ -1,79 +0,0 @@ -use asap_types::aggregation_reference::AggregationReference; -use asap_types::inference_config::InferenceConfig; -use asap_types::query_config::QueryConfig; -use asap_types::streaming_config::StreamingConfig; - -use super::solution::{OptimizerSolution, QueryMethod}; - -/// Translate an `OptimizerSolution` into the deployment artifacts consumed by -/// Arroyo and the query engine. -/// -pub fn translate(solution: &OptimizerSolution) -> (StreamingConfig, InferenceConfig) { - let streaming_config = build_streaming_config(solution); - let inference_config = build_inference_config(solution); - (streaming_config, inference_config) -} - -fn build_streaming_config(solution: &OptimizerSolution) -> StreamingConfig { - // Deployed configs map directly to AggregationConfigs — the types are the same. - StreamingConfig::new(solution.deployed_configs().clone()) -} - -fn build_inference_config(solution: &OptimizerSolution) -> InferenceConfig { - use asap_types::enums::{CleanupPolicy, QueryLanguage}; - - let mut inference = InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::NoCleanup); - - for assignment in &solution.assignments { - let aggregation_id = assignment.aggregation_id; - 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.item.query_strings { - let mut query_config = - QueryConfig::with_plan(query_string.clone(), query_string.clone(), vec![]) - .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); - } - } - - inference -} - -/// For a Merge assignment, the number of retained windows to configure in the -/// inference config (num_aggregates_to_retain on the AggregationReference). -pub fn retention_count_for_assignment(query_method: &QueryMethod) -> u64 { - match query_method { - QueryMethod::Direct => 1, - QueryMethod::Merge { num_windows } => *num_windows, - // Subtract combines exactly 2 prefix-sum checkpoints per query - // (current cumulative, and the one from range_a ago), regardless of - // how many checkpoints the engine retains to make that pair available - // (see candidate_gen.rs's n_windows, a separate concept: the deployed - // AggregationConfig's retention depth). - QueryMethod::Subtract => 2, - } -} - -/// Summary of what the translator produced, for logging/debugging. -#[derive(Debug)] -pub struct TranslationSummary { - pub num_deployed_configs: usize, - pub num_sketch_assignments: usize, -} - -impl TranslationSummary { - pub fn from_solution(solution: &OptimizerSolution) -> Self { - Self { - num_deployed_configs: solution.deployed_configs().len(), - num_sketch_assignments: solution.num_sketch_served(), - } - } -} diff --git a/asap-planner-rs/src/promql/generator.rs b/asap-planner-rs/src/promql/generator.rs index e375f13c..11f59da8 100644 --- a/asap-planner-rs/src/promql/generator.rs +++ b/asap-planner-rs/src/promql/generator.rs @@ -23,9 +23,8 @@ type LeafEntries = Vec<(String, Vec<(String, Option)>)>; /// Run the full planning pipeline and produce YAML outputs. /// /// This is the hardcoded sketch/window selection path that `Controller::generate()` -/// currently calls. `crate::optimizer` implements an optimization-based replacement -/// (issue #405) but is not wired in here yet — see `bin/optimizer_cli.rs` for an -/// offline runner and `.design_docs/optimizer-v1-implementation-plan.md` for status. +/// calls. `crate::optimizer` plans with the rqe-optimizer MILP instead; see +/// `bin/optimizer_cli.rs`. pub fn generate_plan( controller_config: &ControllerConfig, schema: &PromQLSchema,