diff --git a/asap-dropin/README.md b/asap-dropin/README.md index ce463466..aad43fc6 100644 --- a/asap-dropin/README.md +++ b/asap-dropin/README.md @@ -204,6 +204,10 @@ You should see a line like: INFO query_tracker: planner succeeded — streaming aggregations: N, inference queries: M, punted: P ``` +Set `query_tracker.accuracy_sla` and `query_tracker.latency_sla` in +`config/engine_config.yaml` to control the SLAs applied to every inferred query. +Accuracy must be in `(0, 1]`; latency is in estimated CPU seconds and `0.0` is unconstrained. + From this point on, check the routing in the logs: ```bash diff --git a/asap-dropin/config/engine_config.yaml b/asap-dropin/config/engine_config.yaml index 195d3621..8b7de304 100644 --- a/asap-dropin/config/engine_config.yaml +++ b/asap-dropin/config/engine_config.yaml @@ -21,3 +21,5 @@ ingest: query_tracker: enabled: true observation_window_secs: 60 + accuracy_sla: 0.99 + latency_sla: 0.0 diff --git a/asap-planner-rs/src/config/input.rs b/asap-planner-rs/src/config/input.rs index 8e260323..27608a93 100644 --- a/asap-planner-rs/src/config/input.rs +++ b/asap-planner-rs/src/config/input.rs @@ -4,7 +4,6 @@ use asap_types::streaming_config::StreamingConfig; use asap_types::PromQLSchema; use promql_utilities::data_model::KeyByLabelNames; use serde::{Deserialize, Deserializer}; -use tracing::warn; #[derive(Debug, Clone, Deserialize)] #[serde(deny_unknown_fields)] @@ -32,22 +31,6 @@ pub struct ControllerConfig { } impl ControllerConfig { - /// Warn if any query group has both SLAs at 0.0 (the serde Default), - /// which indicates `controller_options` was omitted from the config. - pub fn warn_default_slas(&self) { - for qg in &self.query_groups { - let opts = &qg.controller_options; - if opts.accuracy_sla == 0.0 && opts.latency_sla == 0.0 { - warn!( - query_group_id = ?qg.id, - "controller_options not set in query group; \ - accuracy_sla=0.0 and latency_sla=0.0 will be used — \ - add controller_options to your config" - ); - } - } - } - /// Build a `PromQLSchema` from the `metrics` hints in this config. /// Returns an empty schema if no hints are present. pub fn schema_from_hints(&self) -> PromQLSchema { @@ -68,7 +51,6 @@ pub struct QueryGroup { pub queries: Vec, #[serde(deserialize_with = "deserialize_positive_u64")] pub repetition_delay_ms: u64, - #[serde(default)] pub controller_options: ControllerOptions, /// Per-group step override (ms). Falls back to `RuntimeOptions::step_ms` when None. #[serde(default)] @@ -80,21 +62,37 @@ pub struct QueryGroup { #[derive(Debug, Clone, Deserialize, Default)] pub struct ControllerOptions { - #[serde(deserialize_with = "deserialize_finite_f64")] + #[serde(deserialize_with = "deserialize_accuracy_sla")] pub accuracy_sla: f64, - #[serde(deserialize_with = "deserialize_finite_f64")] + #[serde(deserialize_with = "deserialize_non_negative_f64")] pub latency_sla: f64, } -fn deserialize_finite_f64<'de, D>(deserializer: D) -> Result +fn deserialize_accuracy_sla<'de, D>(deserializer: D) -> Result +where + D: Deserializer<'de>, +{ + let value = f64::deserialize(deserializer)?; + if value.is_finite() && value > 0.0 && value <= 1.0 { + Ok(value) + } else { + Err(serde::de::Error::custom( + "accuracy_sla must be finite and in (0, 1]", + )) + } +} + +fn deserialize_non_negative_f64<'de, D>(deserializer: D) -> Result where D: Deserializer<'de>, { let value = f64::deserialize(deserializer)?; - if value.is_finite() { + if value.is_finite() && value >= 0.0 { Ok(value) } else { - Err(serde::de::Error::custom("must be a finite number")) + Err(serde::de::Error::custom( + "latency_sla must be finite and non-negative", + )) } } @@ -258,7 +256,7 @@ mod tests { use super::*; #[test] - fn rejects_non_finite_controller_slas() { + fn rejects_invalid_controller_slas() { let yaml = r#" query_groups: - queries: [sum(metric)] @@ -272,7 +270,7 @@ query_groups: .expect_err("non-finite SLA values must be rejected") .to_string(); - assert!(error.contains("must be a finite number")); + assert!(error.contains("accuracy_sla must be finite and in (0, 1]")); } #[test] @@ -292,6 +290,43 @@ query_groups: assert_eq!(options.latency_sla, 1.0); } + #[test] + fn rejects_zero_accuracy_and_negative_latency_slas() { + let zero_accuracy = r#" +query_groups: + - queries: [sum(metric)] + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.0 + latency_sla: 0.0 +"#; + let error = serde_yaml::from_str::(zero_accuracy) + .expect_err("zero accuracy SLA must be rejected") + .to_string(); + assert!(error.contains("accuracy_sla must be finite and in (0, 1]")); + + let negative_latency = zero_accuracy + .replace("accuracy_sla: 0.0", "accuracy_sla: 0.99") + .replace("latency_sla: 0.0", "latency_sla: -1.0"); + let error = serde_yaml::from_str::(&negative_latency) + .expect_err("negative latency SLA must be rejected") + .to_string(); + assert!(error.contains("latency_sla must be finite and non-negative")); + } + + #[test] + fn rejects_query_groups_without_controller_options() { + let yaml = r#" +query_groups: + - queries: [sum(metric)] + repetition_delay_ms: 60000 +"#; + let error = serde_yaml::from_str::(yaml) + .expect_err("controller options are required") + .to_string(); + assert!(error.contains("controller_options")); + } + #[test] fn rejects_zero_repetition_delay() { let yaml = r#" diff --git a/asap-planner-rs/src/optimizer/atomic_costs.rs b/asap-planner-rs/src/optimizer/atomic_costs.rs index ceb3c887..82f32200 100644 --- a/asap-planner-rs/src/optimizer/atomic_costs.rs +++ b/asap-planner-rs/src/optimizer/atomic_costs.rs @@ -87,6 +87,13 @@ pub struct AtomicCostEntry { pub type AtomicCostTable = Vec; +/// Costs and benchmarked accuracy measurements resolved for one candidate. +#[derive(Debug, Clone, PartialEq)] +pub struct ResolvedAtomicCosts { + pub costs: AtomicCosts, + pub query_accuracy: BTreeMap, +} + /// 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. @@ -254,8 +261,23 @@ pub fn resolve_atomic_costs( params: &HashMap, n_grouping_labels: usize, ) -> Option { + resolve_atomic_costs_with_accuracy(table, agg_type, params, n_grouping_labels) + .map(|resolved| resolved.costs) +} + +/// Resolve empirical costs together with the measurements used for SLA eligibility. +pub fn resolve_atomic_costs_with_accuracy( + 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()); + return resolve_cms_heap_costs_with_accuracy( + table, + params, + &CmsHeapCostAssumptions::default(), + ); } if let Some(value_bytes) = trivial_value_bytes(agg_type) { @@ -263,11 +285,14 @@ pub fn resolve_atomic_costs( ?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() + return Some(ResolvedAtomicCosts { + costs: AtomicCosts { + mem_bytes_per_instance: (n_grouping_labels as f64 * LABEL_VALUE_CODE_BYTES + + value_bytes) + * HASH_TABLE_SLACK, + ..AtomicCosts::default() + }, + query_accuracy: BTreeMap::new(), }); } @@ -276,20 +301,26 @@ pub fn resolve_atomic_costs( ?agg_type, "no sketch-bench atomic-cost data for this family; using the flat AtomicCosts stub" ); - return Some(AtomicCosts::default()); + return Some(ResolvedAtomicCosts { + costs: AtomicCosts::default(), + query_accuracy: BTreeMap::new(), + }); }; 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, + .map(|entry| ResolvedAtomicCosts { + costs: 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, + }, + query_accuracy: entry.query_accuracy.clone(), }) } @@ -346,11 +377,20 @@ impl CmsHeapCostAssumptions { /// 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. +#[cfg(test)] fn resolve_cms_heap_costs( table: &AtomicCostTable, params: &HashMap, assumptions: &CmsHeapCostAssumptions, ) -> Option { + resolve_cms_heap_costs_with_accuracy(table, params, assumptions).map(|resolved| resolved.costs) +} + +fn resolve_cms_heap_costs_with_accuracy( + table: &AtomicCostTable, + params: &HashMap, + assumptions: &CmsHeapCostAssumptions, +) -> Option { assumptions.validate(); let agg_type = AggregationType::CountMinSketchWithHeap; @@ -431,7 +471,10 @@ fn resolve_cms_heap_costs( "cms-with-heap atomic cost measurement" ); - Some(costs) + Some(ResolvedAtomicCosts { + costs, + query_accuracy: entry.query_accuracy.clone(), + }) } fn require_u64(params: &HashMap, key: &str, agg_type: AggregationType) -> u64 { diff --git a/asap-planner-rs/src/optimizer/cost_model.rs b/asap-planner-rs/src/optimizer/cost_model.rs index f663bb70..a035df6f 100644 --- a/asap-planner-rs/src/optimizer/cost_model.rs +++ b/asap-planner-rs/src/optimizer/cost_model.rs @@ -125,33 +125,53 @@ pub fn query_cost( 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. + return estimated_query_cpu_secs(item, candidate, costs) * weights.query_cpu; + }; + + let units = stored_units(candidate, agg_config.aggregation_type); + let (cpu, mem) = match &candidate.query_method { + QueryMethod::Direct => ( + estimated_query_cpu_secs(item, candidate, costs), + units * costs.mem_bytes_per_instance, + ), + QueryMethod::Merge { num_windows } => ( + estimated_query_cpu_secs(item, candidate, costs), + *num_windows as f64 * units * costs.mem_bytes_per_instance, + ), + QueryMethod::Subtract => ( + estimated_query_cpu_secs(item, candidate, costs), + 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 +} + +/// Estimated CPU seconds for a single query, excluding objective weights and memory. +pub fn estimated_query_cpu_secs( + item: &OptimizerItem, + candidate: &CandidateConfig, + costs: &AtomicCosts, +) -> f64 { + let Some(agg_config) = &candidate.config else { + return costs.exact_query_cpu_secs; }; 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), + match &candidate.query_method { + QueryMethod::Direct => read_cpu, 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, - ) + units * (*num_windows).saturating_sub(1) as f64 * costs.merge_cpu_secs + read_cpu } 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 * (costs.merge_cpu_secs + costs.subtract_cpu_secs) + read_cpu + } + } } /// Units of `mem_bytes_per_instance` held per window; also scales diff --git a/asap-planner-rs/src/optimizer/eligibility.rs b/asap-planner-rs/src/optimizer/eligibility.rs new file mode 100644 index 00000000..4802f34b --- /dev/null +++ b/asap-planner-rs/src/optimizer/eligibility.rs @@ -0,0 +1,91 @@ +use std::collections::BTreeMap; + +use promql_utilities::query_logics::enums::AggregationType; + +use super::atomic_costs::ResolvedAtomicCosts; +use super::candidate_gen::CandidateConfig; +use super::cost_model::estimated_query_cpu_secs; +use super::solution::OptimizerItem; + +const CMS_ERROR_METRIC: &str = "relative_error_mean"; +const KLL_ERROR_METRIC: &str = "mean_rank_err"; +const HLL_ERROR_METRIC: &str = "relative_error"; +const CMS_HEAP_RECALL_METRIC: &str = "recall_at_k"; + +/// Whether a candidate's empirical measurements satisfy one item's SLAs. +pub fn satisfies_slas( + item: &OptimizerItem, + candidate: &CandidateConfig, + resolved: &ResolvedAtomicCosts, +) -> bool { + satisfies_accuracy_sla(item.accuracy_sla, candidate, &resolved.query_accuracy) + && satisfies_latency_sla(item.latency_sla, item, candidate, &resolved.costs) +} + +fn satisfies_accuracy_sla( + accuracy_sla: f64, + candidate: &CandidateConfig, + measurements: &BTreeMap, +) -> bool { + if accuracy_sla == 0.0 { + return true; + } + + let Some(config) = &candidate.config else { + return true; + }; + + if exact_accumulator(config.aggregation_type) { + return true; + } + + let measurement = match config.aggregation_type { + AggregationType::CountMinSketch => measurements.get(CMS_ERROR_METRIC), + AggregationType::DatasketchesKLL => measurements.get(KLL_ERROR_METRIC), + AggregationType::HLL => measurements.get(HLL_ERROR_METRIC), + AggregationType::CountMinSketchWithHeap => measurements.get(CMS_HEAP_RECALL_METRIC), + _ => return false, + }; + let Some(&measurement) = measurement else { + return false; + }; + if !measurement.is_finite() || !(0.0..=1.0).contains(&measurement) { + return false; + } + + match config.aggregation_type { + AggregationType::CountMinSketch + | AggregationType::DatasketchesKLL + | AggregationType::HLL => measurement <= 1.0 - accuracy_sla, + AggregationType::CountMinSketchWithHeap => measurement >= accuracy_sla, + _ => unreachable!("only measured approximate families reach this branch"), + } +} + +fn satisfies_latency_sla( + latency_sla: f64, + item: &OptimizerItem, + candidate: &CandidateConfig, + costs: &super::cost_model::AtomicCosts, +) -> bool { + if latency_sla == 0.0 { + return true; + } + + let estimate = estimated_query_cpu_secs(item, candidate, costs); + estimate.is_finite() && estimate >= 0.0 && estimate <= latency_sla +} + +fn exact_accumulator(aggregation_type: AggregationType) -> bool { + matches!( + aggregation_type, + AggregationType::Sum + | AggregationType::MultipleSum + | AggregationType::Increase + | AggregationType::MultipleIncrease + | AggregationType::MinMax + | AggregationType::MultipleMinMax + | AggregationType::SetAggregator + | AggregationType::DeltaSetAggregator + ) +} diff --git a/asap-planner-rs/src/optimizer/greedy.rs b/asap-planner-rs/src/optimizer/greedy.rs index dff2b44c..5484baf6 100644 --- a/asap-planner-rs/src/optimizer/greedy.rs +++ b/asap-planner-rs/src/optimizer/greedy.rs @@ -1,8 +1,9 @@ use tracing::debug; -use super::atomic_costs::{resolve_atomic_costs, AtomicCostTable}; +use super::atomic_costs::{resolve_atomic_costs_with_accuracy, AtomicCostTable}; use super::candidate_gen::enumerate_candidates_with_facts; use super::cost_model::{ingest_cost, query_cost, total_cost_rate, AtomicCosts, CostWeights}; +use super::eligibility::satisfies_slas; use super::error::{OptimizerError, UnservableItem}; use super::label_set_facts::ItemFacts; use super::solution::{AQEAssignment, OptimizerItem, OptimizerSolution}; @@ -39,22 +40,35 @@ pub fn greedy_assign( 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 + let candidates_with_costs: Vec<_> = 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( + let resolved = match &c.config { + None => super::atomic_costs::ResolvedAtomicCosts { + costs: AtomicCosts::default(), + query_accuracy: Default::default(), + }, + Some(cfg) => resolve_atomic_costs_with_accuracy( 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)) + Some((c, resolved)) + }) + .collect(); + let had_candidate_before_slas = !candidates_with_costs.is_empty(); + + let Some((best, costs)) = candidates_with_costs + .into_iter() + .filter(|(candidate, resolved)| satisfies_slas(&aqe, candidate, resolved)) + .map(|(candidate, resolved)| { + let cost = + total_cost_rate(&aqe, &candidate, arrival_rate_hz, &resolved.costs, weights); + (candidate, resolved.costs, cost) }) // total_cmp (not partial_cmp().unwrap()) so a stray NaN cost can't panic. .min_by(|(_, _, a), (_, _, b)| a.total_cmp(b)) @@ -67,7 +81,11 @@ pub fn greedy_assign( t_repeat_ms: aqe.t_repeat_ms, accuracy_sla: aqe.accuracy_sla, latency_sla: aqe.latency_sla, - reason: "no candidate remained after structural and atomic-cost filters".into(), + reason: if had_candidate_before_slas { + "no candidate satisfies accuracy_sla or latency_sla".into() + } else { + "no candidate remained after structural and atomic-cost filters".into() + }, }); continue; }; @@ -260,6 +278,63 @@ mod tests { ); } + #[test] + fn cms_heap_candidate_must_meet_accuracy_and_latency_slas() { + 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: 8.0, + query_accuracy: std::collections::BTreeMap::from([("recall_at_k".into(), 0.99)]), + }]; + let mut aqe = make_aqe(Statistic::Topk, 60_000, 60_000, 1.0 / 60.0); + aqe.accuracy_sla = 0.99; + aqe.latency_sla = 1_000.0; + let solution = greedy_assign( + vec![aqe.clone()], + 60_000, + &table, + &CostWeights::default(), + &unit_facts(&[aqe.clone()]), + ) + .expect("candidate meeting both SLAs should be selected"); + assert_eq!(solution.assignments.len(), 1); + + aqe.accuracy_sla = 0.995; + let error = greedy_assign( + vec![aqe.clone()], + 60_000, + &table, + &CostWeights::default(), + &unit_facts(&[aqe.clone()]), + ) + .expect_err("candidate below the accuracy SLA must be rejected"); + let OptimizerError::UnservableItems { items } = error else { + panic!("expected SLA eligibility error"); + }; + assert_eq!( + items[0].reason, + "no candidate satisfies accuracy_sla or latency_sla" + ); + + aqe.accuracy_sla = 0.99; + aqe.latency_sla = 0.000_001; + let error = greedy_assign( + vec![aqe.clone()], + 60_000, + &table, + &CostWeights::default(), + &unit_facts(&[aqe]), + ) + .expect_err("candidate above the latency SLA must be rejected"); + assert!(matches!(error, OptimizerError::UnservableItems { .. })); + } + /// 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] diff --git a/asap-planner-rs/src/optimizer/mod.rs b/asap-planner-rs/src/optimizer/mod.rs index 91326c71..37c71b1b 100644 --- a/asap-planner-rs/src/optimizer/mod.rs +++ b/asap-planner-rs/src/optimizer/mod.rs @@ -3,6 +3,7 @@ pub mod atomic_costs; pub mod candidate_gen; pub mod constants; pub mod cost_model; +mod eligibility; pub mod error; pub mod greedy; pub mod label_set_facts; diff --git a/asap-planner-rs/src/promql/controller.rs b/asap-planner-rs/src/promql/controller.rs index a24d6edc..a0d69887 100644 --- a/asap-planner-rs/src/promql/controller.rs +++ b/asap-planner-rs/src/promql/controller.rs @@ -39,7 +39,6 @@ impl Controller { .validate() .map_err(ControllerError::PlannerError)?; } - config.warn_default_slas(); let all_queries: Vec = config .query_groups .iter() @@ -81,7 +80,6 @@ impl Controller { .validate() .map_err(ControllerError::PlannerError)?; } - config.warn_default_slas(); Ok(Self { config, schema, @@ -101,7 +99,6 @@ impl Controller { .validate() .map_err(ControllerError::PlannerError)?; } - config.warn_default_slas(); Ok(Self { config, schema, diff --git a/asap-planner-rs/src/query_log/converter.rs b/asap-planner-rs/src/query_log/converter.rs index 756dcd6b..91330d45 100644 --- a/asap-planner-rs/src/query_log/converter.rs +++ b/asap-planner-rs/src/query_log/converter.rs @@ -1,6 +1,8 @@ use asap_types::enums::CleanupPolicy; -use crate::config::input::{AggregateCleanupConfig, ControllerConfig, QueryGroup}; +use crate::config::input::{ + AggregateCleanupConfig, ControllerConfig, ControllerOptions, QueryGroup, +}; use super::frequency::{InstantQueryInfo, RangeQueryInfo}; @@ -10,6 +12,15 @@ use super::frequency::{InstantQueryInfo, RangeQueryInfo}; pub fn to_controller_config( instants: Vec, ranges: Vec, +) -> ControllerConfig { + to_controller_config_with_options(instants, ranges, ControllerOptions::default()) +} + +/// Build a `ControllerConfig` from extracted queries with shared SLA requirements. +pub fn to_controller_config_with_options( + instants: Vec, + ranges: Vec, + controller_options: ControllerOptions, ) -> ControllerConfig { let mut query_groups: Vec = Vec::new(); @@ -18,7 +29,7 @@ pub fn to_controller_config( id: None, queries: vec![info.query], repetition_delay_ms: info.repetition_delay_ms, - controller_options: Default::default(), + controller_options: controller_options.clone(), step_ms: None, range_duration_ms: None, }); @@ -29,7 +40,7 @@ pub fn to_controller_config( id: None, queries: vec![info.query], repetition_delay_ms: info.repetition_delay_ms, - controller_options: Default::default(), + controller_options: controller_options.clone(), step_ms: Some(info.step_ms), range_duration_ms: Some(info.range_duration_ms), }); diff --git a/asap-planner-rs/src/query_log/mod.rs b/asap-planner-rs/src/query_log/mod.rs index 135e3fe0..e1d448fa 100644 --- a/asap-planner-rs/src/query_log/mod.rs +++ b/asap-planner-rs/src/query_log/mod.rs @@ -2,6 +2,6 @@ pub mod converter; pub mod frequency; pub mod parser; -pub use converter::to_controller_config; +pub use converter::{to_controller_config, to_controller_config_with_options}; pub use frequency::{infer_queries, InstantQueryInfo, RangeQueryInfo}; pub use parser::{parse_log_file, LogEntry}; diff --git a/asap-planner-rs/tests/integration.rs b/asap-planner-rs/tests/integration.rs index 7c864429..5f231ad3 100644 --- a/asap-planner-rs/tests/integration.rs +++ b/asap-planner-rs/tests/integration.rs @@ -823,6 +823,9 @@ query_groups: queries: - "rate(http_requests_total[5m])" repetition_delay_ms: 300000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 0.0 metrics: - metric: "http_requests_total" labels: ["instance"] @@ -842,6 +845,9 @@ query_groups: queries: - "rate(http_requests_total[5m])" repetition_delay_ms: 300000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 0.0 metrics: - metric: "http_requests_total" labels: ["instance"] @@ -876,6 +882,9 @@ query_groups: queries: - "rate(http_requests_total[5m])" repetition_delay_ms: 300000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 0.0 "#, ) .unwrap(); @@ -906,6 +915,9 @@ query_groups: - "sum by (hint_label) (hinted_metric)" - "rate(unhinted_metric[5m])" repetition_delay_ms: 300000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 0.0 metrics: - metric: "hinted_metric" labels: ["hint_label"] @@ -937,6 +949,9 @@ query_groups: queries: - "rate(http_requests_total[5m])" repetition_delay_ms: 300000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 0.0 "#, ) .unwrap(); diff --git a/asap-query-engine/src/engine_config.rs b/asap-query-engine/src/engine_config.rs index 9388dc25..e6359c26 100644 --- a/asap-query-engine/src/engine_config.rs +++ b/asap-query-engine/src/engine_config.rs @@ -19,6 +19,22 @@ pub fn check_config(config: &EngineConfig) -> Result<(), String> { return Err("query_tracker.enabled=true requires backend.type=prometheus".into()); } + if config.query_tracker.enabled + && (!config.query_tracker.accuracy_sla.is_finite() + || config.query_tracker.accuracy_sla <= 0.0 + || config.query_tracker.accuracy_sla > 1.0) + { + return Err("query_tracker.accuracy_sla must be finite and in (0, 1] when enabled".into()); + } + + if config.query_tracker.enabled + && (!config.query_tracker.latency_sla.is_finite() || config.query_tracker.latency_sla < 0.0) + { + return Err( + "query_tracker.latency_sla must be finite and non-negative when enabled".into(), + ); + } + Ok(()) } @@ -344,6 +360,8 @@ impl Default for PrecomputeSettings { pub struct QueryTrackerSettings { pub enabled: bool, pub observation_window_secs: u64, + pub accuracy_sla: f64, + pub latency_sla: f64, } impl Default for QueryTrackerSettings { @@ -351,6 +369,8 @@ impl Default for QueryTrackerSettings { Self { enabled: false, observation_window_secs: 100, + accuracy_sla: 0.0, + latency_sla: 0.0, } } } @@ -643,6 +663,8 @@ backend: type: "clickhouse" query_tracker: enabled: true + accuracy_sla: 0.99 + latency_sla: 0.0 "#; let config: EngineConfig = Figment::new().merge(Yaml::string(yaml)).extract().unwrap(); assert!(check_config(&config).is_err()); @@ -660,6 +682,8 @@ backend: type: "prometheus" query_tracker: enabled: true + accuracy_sla: 0.99 + latency_sla: 0.0 "#; let config: EngineConfig = Figment::new().merge(Yaml::string(yaml)).extract().unwrap(); assert!(check_config(&config).is_ok()); diff --git a/asap-query-engine/src/main.rs b/asap-query-engine/src/main.rs index cda96f63..7b4e2aea 100644 --- a/asap-query-engine/src/main.rs +++ b/asap-query-engine/src/main.rs @@ -316,6 +316,8 @@ async fn main() -> Result<()> { let tracker_config = QueryTrackerConfig { observation_window_secs: config.query_tracker.observation_window_secs, data_ingestion_interval_ms: config.data_ingestion_interval_ms, + accuracy_sla: config.query_tracker.accuracy_sla, + latency_sla: config.query_tracker.latency_sla, }; let runtime_options = asap_planner::RuntimeOptions { data_ingestion_interval_ms: config.data_ingestion_interval_ms, diff --git a/asap-query-engine/src/query_tracker/tracker.rs b/asap-query-engine/src/query_tracker/tracker.rs index 5f45e08d..e5301ca4 100644 --- a/asap-query-engine/src/query_tracker/tracker.rs +++ b/asap-query-engine/src/query_tracker/tracker.rs @@ -2,7 +2,8 @@ use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, RwLock}; use std::time::Duration; -use asap_planner::query_log::{infer_queries, to_controller_config, LogEntry}; +use asap_planner::config::input::ControllerOptions; +use asap_planner::query_log::{infer_queries, to_controller_config_with_options, LogEntry}; use asap_types::inference_config::InferenceConfig; use asap_types::streaming_config::StreamingConfig; use chrono::{DateTime, Utc}; @@ -18,6 +19,10 @@ pub struct QueryTrackerConfig { pub observation_window_secs: u64, /// Data ingestion interval (ms), passed through to `infer_queries`. pub data_ingestion_interval_ms: u64, + /// Accuracy required for every query inferred from observed traffic. + pub accuracy_sla: f64, + /// Query CPU latency limit for every inferred query; zero is unconstrained. + pub latency_sla: f64, } pub struct QueryTracker { @@ -137,7 +142,14 @@ impl QueryTracker { // Build ControllerConfig, including current configs as context for the planner. // NOTE: existing_* fields are wired through but the planner does not yet act on them. - let mut controller_config = to_controller_config(instants, ranges); + let mut controller_config = to_controller_config_with_options( + instants, + ranges, + ControllerOptions { + accuracy_sla: tracker.config.accuracy_sla, + latency_sla: tracker.config.latency_sla, + }, + ); controller_config.existing_streaming_config = Some((*tracker.streaming_config.read().unwrap().clone()).clone()); controller_config.existing_inference_config = @@ -229,6 +241,8 @@ mod tests { let tracker = make_tracker(QueryTrackerConfig { observation_window_secs: 600, data_ingestion_interval_ms: 15_000, + accuracy_sla: 0.99, + latency_sla: 0.0, }); tracker.record_instant("rate(http_requests_total[5m])", 1700000000.0); tracker.record_instant("rate(http_requests_total[5m])", 1700000060.0); @@ -243,6 +257,8 @@ mod tests { let tracker = make_tracker(QueryTrackerConfig { observation_window_secs: 600, data_ingestion_interval_ms: 15_000, + accuracy_sla: 0.99, + latency_sla: 0.0, }); tracker.record_range( "rate(http_requests_total[5m])", @@ -261,6 +277,8 @@ mod tests { let tracker = Arc::new(make_tracker(QueryTrackerConfig { observation_window_secs: 600, data_ingestion_interval_ms: 15_000, + accuracy_sla: 0.99, + latency_sla: 0.0, })); // Record enough entries for infer_queries to produce results (need >=2 per query). @@ -290,6 +308,8 @@ mod tests { let tracker = Arc::new(make_tracker(QueryTrackerConfig { observation_window_secs: 600, data_ingestion_interval_ms: 15_000, + accuracy_sla: 0.99, + latency_sla: 0.0, })); let mock_client = Arc::new(MockPlannerClient::new()); @@ -309,6 +329,8 @@ mod tests { let tracker = Arc::new(make_tracker(QueryTrackerConfig { observation_window_secs: 600, data_ingestion_interval_ms: 15_000, + accuracy_sla: 0.99, + latency_sla: 0.0, })); // Record enough for infer_queries to produce results. @@ -401,6 +423,8 @@ mod tests { QueryTrackerConfig { observation_window_secs: 600, data_ingestion_interval_ms: 15_000, + accuracy_sla: 0.99, + latency_sla: 0.0, }, sc, ic, @@ -432,5 +456,9 @@ mod tests { config.existing_inference_config.is_some(), "existing_inference_config must be populated from the shared ref" ); + assert!(config.query_groups.iter().all(|group| { + group.controller_options.accuracy_sla == 0.99 + && group.controller_options.latency_sla == 0.0 + })); } } diff --git a/docs/03-how-to-guides/operations/try-asap-planner-promql.md b/docs/03-how-to-guides/operations/try-asap-planner-promql.md index fec12012..b1197cda 100644 --- a/docs/03-how-to-guides/operations/try-asap-planner-promql.md +++ b/docs/03-how-to-guides/operations/try-asap-planner-promql.md @@ -43,8 +43,10 @@ auto-infer label sets per metric. ### `controller_options` -`accuracy_sla` and `latency_sla` are required by the config schema but not currently used by -the planner's decision logic — any numeric values are fine (e.g. the placeholders above). +`accuracy_sla` must be in `(0, 1]`; candidates must meet it using their benchmarked +per-family accuracy measurement. `latency_sla` is a maximum estimated query CPU time in +seconds; `0.0` leaves latency unconstrained. A query is rejected if no candidate meets both +constraints. ### Choosing `repetition_delay_ms` and `--data-ingestion-interval-ms` diff --git a/docs/03-how-to-guides/operations/try-asap-planner-sql.md b/docs/03-how-to-guides/operations/try-asap-planner-sql.md index cdb39838..676d787d 100644 --- a/docs/03-how-to-guides/operations/try-asap-planner-sql.md +++ b/docs/03-how-to-guides/operations/try-asap-planner-sql.md @@ -46,9 +46,10 @@ For each table: ### `controller_options` -`accuracy_sla` and `latency_sla` are required by the config schema but not currently used by -the planner's decision logic — any numeric values are fine (e.g. the placeholders in the -example above). +`accuracy_sla` must be in `(0, 1]`; candidates must meet it using their benchmarked +per-family accuracy measurement. `latency_sla` is a maximum estimated query CPU time in +seconds; `0.0` leaves latency unconstrained. A query is rejected if no candidate meets both +constraints. ### Choosing `repetition_delay` and `--data-ingestion-interval`