diff --git a/.design_docs/optimizer-mip-formulation.md b/.design_docs/optimizer-mip-formulation.md index 4a648b3b..a457fc16 100644 --- a/.design_docs/optimizer-mip-formulation.md +++ b/.design_docs/optimizer-mip-formulation.md @@ -32,13 +32,13 @@ Symbols computed during optimizer setup (before the MIP is solved). Formulas are | Symbol | Derived from | Definition | |--------|-------------|-----------| | $T_r$ | Query workload | Repeat interval for RQE $r$ (ms) | -| $A = \{a_1, \ldots, a_m\}$ | Query workload | All AQEs (Atomic Query Expressions), deduplicated across all RQEs | -| $R_a \subseteq R$ | Query workload | RQEs that reference AQE $a$ | +| $A = \{a_1, \ldots, a_m\}$ | Query workload | Optimizer items, each keyed by AQE requirements, repeat interval, and SLAs | +| $R_a \subseteq R$ | Query workload | RQEs that contribute to item $a$ | | $\text{range}_a$ | Query workload | Lookback duration baked into $a$'s range vector | | $n_g$ | Configs | Number of $g$-windows used at query time. For tumbling, $n_g = d_g$. For sliding, $n_g$ is a generation parameter used to size $d_g$; the optimizer does not reference it directly. | | $d_g$ | Configs | Physical storage depth — number of completed windows $g$ retains. For tumbling $d_g = n_g$; for sliding $d_g \geq (n_g - 1)(W_g/S_g) + 1$. | | $n(a,g)$ | Query workload, configs | $\lceil \text{range}_a / W_g \rceil$ — number of $g$-windows needed to cover $a$'s lookback range. | -| $f_a$ | Query workload | $\sum_{r \in R_a} 1/T_r$ — aggregate query rate for AQE $a$ (queries/sec). | +| $f_a$ | Query workload | $|R_a| / T_a$ — aggregate query rate for item $a$ (queries/sec). | | $N(s,g)$ | | Label-group multiplier: 1 if $\text{subpop-aware}(s)$, else $N_g$. | | Query method for $(a,g)$ | | One of Direct / Merge / Subtract / Exact — see Derived Quantities. | @@ -123,11 +123,10 @@ A sketch's window size must not be greater than the query range. *Implication for Direct queries:* $n(a,g) = 1$ requires $\lceil \text{range}_a / W_g \rceil = 1$, i.e. $W_g \geq \text{range}_a$. Combined with (WIN) this means the Direct method is only possible when $W_g = \text{range}_a$. Thus. Direct configs are uniquely sized per AQE. -**Freshness** (ingest-type specific): -$$x_{a,g} = 1 \Rightarrow W_g \leq \min_{r \in R_a} T_r \quad \tau_g = \text{tumbling} \tag{FRESHt}$$ -$$x_{a,g} = 1 \Rightarrow S_g \leq \min_{r \in R_a} T_r \quad \tau_g = \text{sliding} \tag{FRESHs}$$ +**Window compatibility:** +$$x_{a,g} = 1 \Rightarrow \text{range}_a \bmod W_g = 0 \land T_a \bmod S_g = 0 \land W_g \bmod S_g = 0 \tag{WINDOW}$$ -For tumbling, a completed window must exist for every query cycle, so $W \leq T_r$. For sliding, a new completed window appears every $S$ seconds, so the binding constraint is $S \leq T_r$ — $W$ can exceed $T_r$ for sliding. +For tumbling, $S_g = W_g$. Candidate dimensions must also be aligned to the scrape interval. **Accuracy:** $$x_{a,g} = 1 \Rightarrow \text{Error}(a, g, \theta_a) \leq \varepsilon_a \tag{ACC}$$ @@ -191,7 +190,7 @@ $$x_{a,g} \leq y_g \qquad \forall a \in A,\ g \in G \tag{2}$$ $$x_{a,g},\ y_g \in \{0,1\} \tag{3}$$ -**(1)** Every AQE is assigned to exactly one config. Always satisfiable since $\text{EXACT}_a$ satisfies all feasibility conditions for any $a$. +**(1)** Every item is assigned to exactly one eligible config. A workload with an item having no eligible config is an optimizer error. **(2)** An AQE cannot be served by a config that is not deployed. **(3)** Integrality. Feasibility is enforced by restricting $x_{a,g}$ to the domain where $\text{Feasible}(a,g) = 1$. diff --git a/asap-planner-rs/src/bin/candidate_gen_dump.rs b/asap-planner-rs/src/bin/candidate_gen_dump.rs index b0af6864..7a2d536e 100644 --- a/asap-planner-rs/src/bin/candidate_gen_dump.rs +++ b/asap-planner-rs/src/bin/candidate_gen_dump.rs @@ -60,21 +60,24 @@ fn main() -> anyhow::Result<()> { 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: qg.controller_options.latency_sla, }) }) .collect(); - let aqes = extract_aqes(&rqes, &schema, args.scrape_interval_ms); - println!("=== {} AQE(s) ===", aqes.len()); + 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--- AQE #{i}: metric={} stat={:?} range={}ms min_t={}ms gcd_t={}ms freq={:.4}Hz ---", + "\n--- Item #{i}: metric={} stat={:?} range={}ms T={}ms accuracy_sla={} latency_sla={} freq={:.4}Hz ---", aqe.requirements.metric, aqe.requirements.statistics, aqe.requirements.data_range_ms, - aqe.min_t_repeat_ms, - aqe.t_repeat_gcd_ms, + aqe.t_repeat_ms, + aqe.accuracy_sla, + aqe.latency_sla, aqe.query_frequency_hz, ); println!(" queries: {:?}", aqe.query_strings); diff --git a/asap-planner-rs/src/config/input.rs b/asap-planner-rs/src/config/input.rs index 18fa1a7f..8e260323 100644 --- a/asap-planner-rs/src/config/input.rs +++ b/asap-planner-rs/src/config/input.rs @@ -3,7 +3,7 @@ use asap_types::inference_config::InferenceConfig; use asap_types::streaming_config::StreamingConfig; use asap_types::PromQLSchema; use promql_utilities::data_model::KeyByLabelNames; -use serde::Deserialize; +use serde::{Deserialize, Deserializer}; use tracing::warn; #[derive(Debug, Clone, Deserialize)] @@ -66,6 +66,7 @@ impl ControllerConfig { pub struct QueryGroup { pub id: Option, pub queries: Vec, + #[serde(deserialize_with = "deserialize_positive_u64")] pub repetition_delay_ms: u64, #[serde(default)] pub controller_options: ControllerOptions, @@ -79,10 +80,36 @@ pub struct QueryGroup { #[derive(Debug, Clone, Deserialize, Default)] pub struct ControllerOptions { + #[serde(deserialize_with = "deserialize_finite_f64")] pub accuracy_sla: f64, + #[serde(deserialize_with = "deserialize_finite_f64")] pub latency_sla: f64, } +fn deserialize_finite_f64<'de, D>(deserializer: D) -> Result +where + D: Deserializer<'de>, +{ + let value = f64::deserialize(deserializer)?; + if value.is_finite() { + Ok(value) + } else { + Err(serde::de::Error::custom("must be a finite number")) + } +} + +fn deserialize_positive_u64<'de, D>(deserializer: D) -> Result +where + D: Deserializer<'de>, +{ + let value = u64::deserialize(deserializer)?; + if value > 0 { + Ok(value) + } else { + Err(serde::de::Error::custom("must be greater than zero")) + } +} + #[derive(Debug, Clone, Deserialize)] pub struct MetricDefinition { pub metric: String, @@ -225,3 +252,57 @@ pub struct ElasticDSLQueryGroup { pub time_field: String, pub controller_options: ControllerOptions, } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn rejects_non_finite_controller_slas() { + let yaml = r#" +query_groups: + - queries: [sum(metric)] + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: .nan + latency_sla: .inf +"#; + + let error = serde_yaml::from_str::(yaml) + .expect_err("non-finite SLA values must be rejected") + .to_string(); + + assert!(error.contains("must be a finite number")); + } + + #[test] + fn accepts_finite_controller_slas() { + let yaml = r#" +query_groups: + - queries: [sum(metric)] + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +"#; + + let config: ControllerConfig = serde_yaml::from_str(yaml).unwrap(); + let options = &config.query_groups[0].controller_options; + assert_eq!(options.accuracy_sla, 0.99); + assert_eq!(options.latency_sla, 1.0); + } + + #[test] + fn rejects_zero_repetition_delay() { + let yaml = r#" +query_groups: + - queries: [sum(metric)] + repetition_delay_ms: 0 +"#; + + let error = serde_yaml::from_str::(yaml) + .expect_err("zero repeat interval must be rejected") + .to_string(); + assert!(error.contains("must be greater than zero")); + } +} diff --git a/asap-planner-rs/src/optimizer/aqe_extractor.rs b/asap-planner-rs/src/optimizer/aqe_extractor.rs index 52869d1b..3e4a2d48 100644 --- a/asap-planner-rs/src/optimizer/aqe_extractor.rs +++ b/asap-planner-rs/src/optimizer/aqe_extractor.rs @@ -4,12 +4,12 @@ use asap_types::query_requirements::{build_query_requirements_promql, QueryRequi use asap_types::PromQLSchema; use promql_utilities::data_model::KeyByLabelNames; use promql_utilities::query_logics::enums::Statistic; -use tracing::warn; use crate::planner::patterns::build_patterns; use crate::planner::promql::{parse_binary_arms, BinaryArm}; -use super::solution::AQE; +use super::error::OptimizerError; +use super::solution::OptimizerItem; /// One repeating query expression: a PromQL query string and its repetition /// interval (e.g. the refresh interval of the dashboard panel it belongs to). @@ -17,13 +17,13 @@ use super::solution::AQE; pub struct RQE { pub query_string: String, pub t_repeat_ms: u64, + pub accuracy_sla: f64, + pub latency_sla: f64, } -/// Stable deduplication key for an AQE. -/// Two leaf queries that produce identical requirements are treated as the same -/// AQE regardless of which RQE they came from. +/// Stable key for merging identical optimizer demand. #[derive(Debug, Clone, Hash, PartialEq, Eq)] -struct AQEKey { +struct OptimizerItemKey { metric: String, /// Statistics are produced in a stable order by get_statistics_to_compute. statistics: Vec, @@ -31,10 +31,14 @@ struct AQEKey { grouping_labels: KeyByLabelNames, spatial_filter_normalized: String, topk_count_events: Option, + topk_by_labels: Option, + t_repeat_ms: u64, + accuracy_sla_bits: u64, + latency_sla_bits: u64, } -impl AQEKey { - fn from_requirements(req: &QueryRequirements) -> Self { +impl OptimizerItemKey { + fn from_rqe(req: &QueryRequirements, rqe: &RQE) -> Self { Self { metric: req.metric.clone(), statistics: req.statistics.clone(), @@ -42,91 +46,80 @@ impl AQEKey { grouping_labels: req.grouping_labels.clone(), spatial_filter_normalized: req.spatial_filter_normalized.clone(), topk_count_events: req.topk_count_events, + topk_by_labels: req.topk_by_labels.clone(), + t_repeat_ms: rqe.t_repeat_ms, + accuracy_sla_bits: normalized_f64_bits(rqe.accuracy_sla), + latency_sla_bits: normalized_f64_bits(rqe.latency_sla), } } } -/// Extract and deduplicate AQEs from a set of RQEs. +/// Extract and deduplicate optimizer items from a set of RQEs. /// /// Each RQE is decomposed into leaf query expressions (recursively splitting /// binary arithmetic operators), then each leaf is pattern-matched to produce -/// a `QueryRequirements`. AQEs with identical requirements are merged. -/// -/// Three frequency-related values are computed per AQE: -/// - `query_frequency_hz`: Σ 1/T_r — total query load for the MIP objective. -/// - `min_t_repeat_ms`: min(T_r) — freshness bound on window size W ≤ min_t. -/// - `t_repeat_gcd_ms`: GCD(T_r) — natural slide interval S for candidate -/// generation (windows completing every GCD ms align with all dashboards). -/// -/// Leaf queries that do not match any supported pattern (e.g. unsupported -/// functions, parse errors) are skipped with a warning. +/// a `QueryRequirements`. Occurrences merge only when requirements, cadence, +/// and both SLAs agree. Their frequency is `count * 1000 / T_ms`. pub fn extract_aqes( rqes: &[RQE], metric_schema: &PromQLSchema, scrape_interval_ms: u64, -) -> Vec { - // (key) -> (requirements, query_strings, sum_freq, min_t, gcd_t) - let mut acc: HashMap, f64, u64, u64)> = HashMap::new(); +) -> Result, OptimizerError> { + let mut acc: HashMap, f64)> = HashMap::new(); for rqe in rqes { if rqe.t_repeat_ms == 0 { - warn!( - query = %rqe.query_string, - "aqe_extractor: skipping RQE with repetition_delay_ms=0 \ - (would produce infinite query frequency and corrupt GCD)" - ); - continue; + return Err(OptimizerError::InvalidRepeatInterval { + query: rqe.query_string.clone(), + }); } let leaves = decompose_to_leaves(&rqe.query_string); for leaf in leaves { match extract_requirements(&leaf, metric_schema, scrape_interval_ms) { - Some(req) => { - let key = AQEKey::from_requirements(&req); - let entry = acc - .entry(key) - .or_insert_with(|| (req, Vec::new(), 0.0, u64::MAX, 0)); + Ok(req) => { + let key = OptimizerItemKey::from_rqe(&req, rqe); + let entry = acc.entry(key).or_insert_with(|| (req, Vec::new(), 0.0)); if !entry.1.contains(&leaf) { entry.1.push(leaf); } // query_frequency_hz must stay in Hz (queries per real second) // regardless of t_repeat_ms's internal unit — 1000.0 / ms, not 1.0 / ms. entry.2 += 1000.0 / rqe.t_repeat_ms as f64; - entry.3 = entry.3.min(rqe.t_repeat_ms); - entry.4 = if entry.4 == 0 { - rqe.t_repeat_ms - } else { - gcd(entry.4, rqe.t_repeat_ms) - }; } - None => { - warn!( - query = %leaf, - "aqe_extractor: skipping unsupported or unparseable leaf query" - ); + Err(reason) => { + return Err(OptimizerError::UnsupportedLeaf { + query: rqe.query_string.clone(), + leaf, + reason, + }) } } } } - acc.into_values() + Ok(acc + .into_iter() .map( - |( - requirements, - query_strings, - query_frequency_hz, - min_t_repeat_ms, - t_repeat_gcd_ms, - )| AQE { + |(key, (requirements, query_strings, query_frequency_hz))| OptimizerItem { requirements, query_strings, query_frequency_hz, - min_t_repeat_ms, - t_repeat_gcd_ms, + t_repeat_ms: key.t_repeat_ms, + accuracy_sla: f64::from_bits(key.accuracy_sla_bits), + latency_sla: f64::from_bits(key.latency_sla_bits), }, ) - .collect() + .collect()) +} + +fn normalized_f64_bits(value: f64) -> u64 { + if value == 0.0 { + 0.0f64.to_bits() + } else { + value.to_bits() + } } /// Euclidean GCD. `num-integer` is not in the workspace; this two-liner is @@ -166,18 +159,21 @@ fn extract_requirements( query: &str, metric_schema: &PromQLSchema, data_ingestion_interval_ms: u64, -) -> Option { - let ast = promql_parser::parser::parse(query).ok()?; +) -> Result { + let ast = promql_parser::parser::parse(query).map_err(|error| error.to_string())?; let patterns = build_patterns(); - let match_result = patterns.iter().find_map(|pat| { - let r = pat.matches(&ast); - if r.matches { - Some(r) - } else { - None - } - })?; + let match_result = patterns + .iter() + .find_map(|pat| { + let r = pat.matches(&ast); + if r.matches { + Some(r) + } else { + None + } + }) + .ok_or_else(|| "no supported query pattern matched".to_string())?; build_query_requirements_promql( query, @@ -185,6 +181,7 @@ fn extract_requirements( metric_schema, data_ingestion_interval_ms, ) + .ok_or_else(|| "query requirements could not be derived".to_string()) } #[cfg(test)] @@ -199,13 +196,15 @@ mod tests { RQE { query_string: query.to_string(), t_repeat_ms: t_ms, + accuracy_sla: 0.0, + latency_sla: 0.0, } } #[test] fn single_temporal_query() { let rqes = vec![rqe("sum_over_time(metric[5m])", 60_000)]; - let aqes = extract_aqes(&rqes, &empty_schema(), 15_000); + let aqes = extract_aqes(&rqes, &empty_schema(), 15_000).unwrap(); assert_eq!(aqes.len(), 1); assert!((aqes[0].query_frequency_hz - 1.0 / 60.0).abs() < 1e-9); assert_eq!(aqes[0].requirements.metric, "metric"); @@ -218,45 +217,94 @@ mod tests { "sum_over_time(metric_a[5m]) / sum_over_time(metric_b[5m])", 60_000, )]; - let aqes = extract_aqes(&rqes, &empty_schema(), 15_000); + let aqes = extract_aqes(&rqes, &empty_schema(), 15_000).unwrap(); assert_eq!(aqes.len(), 2); } #[test] fn binary_with_scalar_produces_one_aqe() { let rqes = vec![rqe("sum_over_time(metric[5m]) * 100", 60_000)]; - let aqes = extract_aqes(&rqes, &empty_schema(), 15_000); + let aqes = extract_aqes(&rqes, &empty_schema(), 15_000).unwrap(); assert_eq!(aqes.len(), 1); } #[test] - fn same_aqe_in_two_rqes_deduplicates_and_sums_frequency() { + fn identical_items_deduplicate_and_sum_frequency() { let rqes = vec![ rqe("sum_over_time(metric[5m])", 60_000), - rqe("sum_over_time(metric[5m])", 30_000), + rqe("sum_over_time(metric[5m])", 60_000), ]; - let aqes = extract_aqes(&rqes, &empty_schema(), 15_000); + let aqes = extract_aqes(&rqes, &empty_schema(), 15_000).unwrap(); assert_eq!(aqes.len(), 1); // query_frequency_hz = sum of rates (total query load for the MIP objective) - let expected_freq = 1.0 / 60.0 + 1.0 / 30.0; + let expected_freq = 2.0 / 60.0; assert!((aqes[0].query_frequency_hz - expected_freq).abs() < 1e-9); - // min_t and gcd_t used for windowing constraints - assert_eq!(aqes[0].min_t_repeat_ms, 30_000); - assert_eq!(aqes[0].t_repeat_gcd_ms, 30_000); // gcd(60_000, 30_000) = 30_000 - assert_eq!(aqes[0].query_strings.len(), 1); // same string, deduplicated + assert_eq!(aqes[0].t_repeat_ms, 60_000); + assert_eq!(aqes[0].query_strings.len(), 1); + } + + #[test] + fn different_repeat_intervals_become_distinct_items() { + let rqes = vec![ + rqe("sum_over_time(metric[5m])", 60_000), + rqe("sum_over_time(metric[5m])", 30_000), + ]; + let mut items = extract_aqes(&rqes, &empty_schema(), 15_000).unwrap(); + items.sort_by_key(|item| item.t_repeat_ms); + + assert_eq!(items.len(), 2); + assert_eq!(items[0].t_repeat_ms, 30_000); + assert_eq!(items[1].t_repeat_ms, 60_000); + assert!((items[0].query_frequency_hz - 1.0 / 30.0).abs() < 1e-9); + assert!((items[1].query_frequency_hz - 1.0 / 60.0).abs() < 1e-9); + } + + #[test] + fn different_slas_become_distinct_items() { + let mut strict = rqe("sum_over_time(metric[5m])", 60_000); + strict.accuracy_sla = 0.99; + let mut relaxed = rqe("sum_over_time(metric[5m])", 60_000); + relaxed.accuracy_sla = 0.9; + + let items = extract_aqes(&[strict, relaxed], &empty_schema(), 15_000).unwrap(); + assert_eq!(items.len(), 2); + } + + #[test] + fn signed_zero_slas_merge_into_one_item() { + let mut negative_zero = rqe("sum_over_time(metric[5m])", 60_000); + negative_zero.latency_sla = -0.0; + let items = extract_aqes( + &[negative_zero, rqe("sum_over_time(metric[5m])", 60_000)], + &empty_schema(), + 15_000, + ) + .unwrap(); + assert_eq!(items.len(), 1); + } + + #[test] + fn zero_repeat_interval_is_rejected() { + let rqes = vec![rqe("sum_over_time(metric[5m])", 0)]; + assert!(matches!( + extract_aqes(&rqes, &empty_schema(), 15_000), + Err(OptimizerError::InvalidRepeatInterval { .. }) + )); } #[test] - fn unsupported_query_is_skipped() { + fn unsupported_query_is_rejected() { let rqes = vec![rqe("not_a_real_function(metric[5m])", 60_000)]; - let aqes = extract_aqes(&rqes, &empty_schema(), 15_000); - assert_eq!(aqes.len(), 0); + assert!(matches!( + extract_aqes(&rqes, &empty_schema(), 15_000), + Err(OptimizerError::UnsupportedLeaf { .. }) + )); } #[test] fn spatial_only_query_sets_range_to_scrape_interval() { let rqes = vec![rqe("sum(metric)", 60_000)]; - let aqes = extract_aqes(&rqes, &empty_schema(), 15_000); + let aqes = extract_aqes(&rqes, &empty_schema(), 15_000).unwrap(); assert_eq!(aqes.len(), 1); assert_eq!(aqes[0].requirements.data_range_ms, 15_000); } diff --git a/asap-planner-rs/src/optimizer/candidate_gen.rs b/asap-planner-rs/src/optimizer/candidate_gen.rs index 807213ac..fc69bd61 100644 --- a/asap-planner-rs/src/optimizer/candidate_gen.rs +++ b/asap-planner-rs/src/optimizer/candidate_gen.rs @@ -11,12 +11,12 @@ use super::constants::{ CMS_DEPTHS, CMS_HEAP_SIZES, CMS_WIDTHS, HLL_PRECISIONS, HYDRA_COLS, HYDRA_K, HYDRA_ROWS, KLL_KS, }; use super::sketch_properties::sketch_properties; -use super::solution::{QueryMethod, AQE}; +use super::solution::{OptimizerItem, QueryMethod}; -/// A candidate streaming config for a single AQE, ready for cost evaluation. +/// A candidate streaming config for one optimizer item, ready for cost evaluation. #[derive(Debug, Clone)] pub struct CandidateConfig { - /// None = EXACT fallback (no streaming config; raw Prometheus query at query time). + /// 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, @@ -27,20 +27,18 @@ pub struct CandidateConfig { pub label_group_count: u64, } -/// Enumerate all candidate configs for an AQE. +/// Enumerate all structurally valid candidate configs for an optimizer item. /// /// Iterates over compatible agg types × parameter grid × valid window sizes × -/// {Tumbling, Sliding}. Always appends an EXACT candidate last (always feasible). -/// -/// Multi-statistic AQEs (e.g. avg = [Sum, Count]) return only EXACT — a single -/// sketch family cannot serve two incompatible statistics simultaneously. -pub fn enumerate_candidates(aqe: &AQE, scrape_interval_ms: u64) -> Vec { - enumerate_candidates_with_label_group_count(aqe, scrape_interval_ms, 1) +/// {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 { + enumerate_candidates_with_label_group_count(item, scrape_interval_ms, 1) } /// Enumerate candidates with a dataset-derived label-group count. pub fn enumerate_candidates_with_label_group_count( - aqe: &AQE, + item: &OptimizerItem, scrape_interval_ms: u64, label_group_count: u64, ) -> Vec { @@ -50,14 +48,12 @@ pub fn enumerate_candidates_with_label_group_count( ); let mut candidates = Vec::new(); - if aqe.requirements.statistics.len() != 1 { - // ponytail: multi-stat AQEs (avg) need two sketches; not supported in v1. - candidates.push(exact_candidate(label_group_count)); + if item.requirements.statistics.len() != 1 { return candidates; } - let stat = aqe.requirements.statistics[0]; - let range_a_ms = aqe.requirements.data_range_ms; + let stat = item.requirements.statistics[0]; + let range_a_ms = item.requirements.data_range_ms; for &agg_type in compatible_agg_types(stat) { let props = sketch_properties(agg_type); @@ -67,7 +63,7 @@ pub fn enumerate_candidates_with_label_group_count( // dimension params -- enumerate both weightings when the query doesn't pin one. let sub_type_variants: Vec = if agg_type == AggregationType::CountMinSketchWithHeap { - match aqe.requirements.topk_count_events { + 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()], @@ -79,7 +75,7 @@ pub fn enumerate_candidates_with_label_group_count( 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, aqe.t_repeat_gcd_ms, scrape_interval_ms) + 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 @@ -94,7 +90,7 @@ pub fn enumerate_candidates_with_label_group_count( }; let config = build_config( - aqe, + item, agg_type, sub_type, ¶ms, @@ -114,37 +110,17 @@ pub fn enumerate_candidates_with_label_group_count( } } - candidates.push(exact_candidate(label_group_count)); candidates } -fn exact_candidate(label_group_count: u64) -> CandidateConfig { - CandidateConfig { - config: None, - query_method: QueryMethod::Exact, - n_windows: 0, - label_group_count, - } -} - /// Window candidates: (WindowType, W_ms, slide_interval_ms, n_windows). /// -/// Tumbling: W must divide GCD(range_a, t_repeat_gcd_ms) and be a multiple of scrape_interval. -/// Dividing the GCD ensures (a) n complete windows cover range_a exactly, and -/// (b) window completions are harmonically aligned with every dashboard refresh cycle. -/// Slide interval = W (a tumbling window "slides" by its own width). -/// Sliding: W = range_a / k for each W that is a multiple of scrape_interval and divides -/// range_a (k = range_a / W). At query time k staggered readings spaced W apart -/// are merged or subtracted to cover range_a. -/// S must satisfy three constraints: -/// (a) S | W — so W-spaced snapshots land on slide boundaries (multi-window correctness) -/// (b) S | t_repeat_gcd — so slide boundaries align with every dashboard refresh cycle -/// (c) S < W — S=W is excluded because slide==window is tumbling (duplicate candidate) -/// S is enumerated as multiples of scrape_interval < W; the divisibility check on -/// gcd(W, t_repeat_gcd) rejects values that fail (a) or (b) without a separate bound. +/// 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_gcd_ms: u64, + t_repeat_ms: u64, scrape_interval_ms: u64, ) -> Vec<(WindowType, u64, u64, u64)> { let range_a = range_a_ms; @@ -154,11 +130,7 @@ fn window_candidates( let mut out = Vec::new(); - // Tumbling: W divides GCD(range_a, t_repeat_gcd) and is a multiple of scrape_interval. - // W | t_repeat_gcd ensures window completions align harmonically with all dashboards. - // W | range_a (implied since t_repeat_gcd | range_a is checked at generation time, but - // we verify explicitly via the gcd) ensures n windows cover range_a exactly. - let tumbling_divisor = super::aqe_extractor::gcd(range_a, t_repeat_gcd_ms); + 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) { @@ -168,15 +140,11 @@ fn window_candidates( w += scrape_interval_ms; } - // Sliding: W = range_a / k for each valid W (multiple of scrape_interval, divides range_a). - // S doubles from scrape_interval up to min(W, min_t_repeat_ms). n_windows = k. let mut w = scrape_interval_ms; while w <= range_a { if range_a.is_multiple_of(w) { let k = range_a / w; - // Valid S: S | gcd(W, t_repeat_gcd). Iterate up to W (exclusive); the - // divisibility check rejects anything above gcd automatically. - let slide_divisor = super::aqe_extractor::gcd(w, t_repeat_gcd_ms); + 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) { @@ -218,7 +186,7 @@ fn determine_query_method( /// overwrites it with a real id when (if) a solver deploys this candidate. #[allow(clippy::too_many_arguments)] fn build_config( - aqe: &AQE, + item: &OptimizerItem, agg_type: AggregationType, sub_type: &str, params: &HashMap, @@ -232,15 +200,15 @@ fn build_config( agg_type, sub_type.to_string(), params.clone(), - aqe.requirements.grouping_labels.clone(), + item.requirements.grouping_labels.clone(), KeyByLabelNames::empty(), // aggregated_labels (not needed for optimizer feasibility) KeyByLabelNames::empty(), // rollup_labels String::new(), // original_yaml w, slide_interval, window_type, - aqe.requirements.spatial_filter_normalized.clone(), - aqe.requirements.metric.clone(), + item.requirements.spatial_filter_normalized.clone(), + item.requirements.metric.clone(), Some(n_windows), None, // read_count_threshold None, // table_name (SQL only) @@ -336,9 +304,9 @@ mod tests { use asap_types::enums::WindowType; use promql_utilities::data_model::KeyByLabelNames; - fn make_aqe(stat: Statistic, range_ms: u64, min_t: u64) -> AQE { + fn make_aqe(stat: Statistic, range_ms: u64, min_t: u64) -> OptimizerItem { use asap_types::query_requirements::QueryRequirements; - AQE { + OptimizerItem { requirements: QueryRequirements { metric: "test_metric".into(), statistics: vec![stat], @@ -350,18 +318,17 @@ mod tests { }, query_strings: vec!["test_query".into()], query_frequency_hz: 1.0 / 60.0, - min_t_repeat_ms: min_t, - t_repeat_gcd_ms: min_t, + t_repeat_ms: min_t, + accuracy_sla: 0.0, + latency_sla: 0.0, } } #[test] - fn always_includes_exact_fallback() { + 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() - .any(|c| c.config.is_none() && c.query_method == QueryMethod::Exact)); + assert!(candidates.iter().all(|c| c.config.is_some())); } #[test] @@ -474,7 +441,7 @@ mod tests { #[test] fn sliding_full_width_direct_generated_when_range_exceeds_t_repeat() { - // range_a=600_000ms > t_repeat_gcd=30_000ms. W=range_a is valid for sliding since + // 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); @@ -485,13 +452,13 @@ mod tests { }) && c.query_method == QueryMethod::Direct && c.n_windows == 1 }), - "full-width Sliding Direct should be generated even when range_a > min_t_repeat" + "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_repeat_gcd=5_000, scrape=1_000. + // 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). @@ -532,7 +499,7 @@ mod tests { #[test] fn sliding_slide_must_divide_window_size() { - // W=6_000, t_repeat_gcd=6_000: slide_divisor = gcd(6_000, 6_000) = 6_000. + // 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); diff --git a/asap-planner-rs/src/optimizer/cost_model.rs b/asap-planner-rs/src/optimizer/cost_model.rs index b9fb5e75..f2044a1e 100644 --- a/asap-planner-rs/src/optimizer/cost_model.rs +++ b/asap-planner-rs/src/optimizer/cost_model.rs @@ -8,7 +8,7 @@ use super::constants::{ SUBPOPULATION_COUNT, SUBTRACT_CPU_SECS, }; use super::sketch_properties::sketch_properties; -use super::solution::{QueryMethod, AQE}; +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). @@ -99,7 +99,7 @@ pub fn ingest_cost( /// QueryCost(a,g): cost of answering one query for `aqe` from `candidate`. pub fn query_cost( - _aqe: &AQE, + _item: &OptimizerItem, candidate: &CandidateConfig, costs: &AtomicCosts, weights: &CostWeights, @@ -132,11 +132,7 @@ pub fn query_cost( * (costs.merge_cpu_secs + costs.subtract_cpu_secs + costs.query_cpu_secs), 2.0 * subpopulation_count * costs.mem_bytes_per_instance, ) - } - // candidate_gen only ever pairs Exact with config=None, already handled above. - QueryMethod::Exact => { - unreachable!("Exact query_method must not be paired with Some(config)") - } + } // candidate_gen only ever pairs Exact with config=None, already handled above. }; weights.query_cpu * cpu + weights.query_mem * mem @@ -161,14 +157,14 @@ fn effective_subpopulation_count( /// `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( - aqe: &AQE, + item: &OptimizerItem, candidate: &CandidateConfig, arrival_rate_hz: f64, costs: &AtomicCosts, weights: &CostWeights, ) -> f64 { ingest_cost(candidate, arrival_rate_hz, costs, weights) - + aqe.query_frequency_hz * query_cost(aqe, candidate, costs, weights) + + item.query_frequency_hz * query_cost(item, candidate, costs, weights) } #[cfg(test)] @@ -179,8 +175,8 @@ mod tests { use promql_utilities::data_model::KeyByLabelNames; use promql_utilities::query_logics::enums::Statistic; - fn make_aqe(stat: Statistic, range_ms: u64, min_t: u64) -> AQE { - AQE { + fn make_aqe(stat: Statistic, range_ms: u64, min_t: u64) -> OptimizerItem { + OptimizerItem { requirements: QueryRequirements { metric: "test_metric".into(), statistics: vec![stat], @@ -192,26 +188,12 @@ mod tests { }, query_strings: vec!["test_query".into()], query_frequency_hz: 1.0 / 60.0, - min_t_repeat_ms: min_t, - t_repeat_gcd_ms: min_t, + t_repeat_ms: min_t, + accuracy_sla: 0.0, + latency_sla: 0.0, } } - #[test] - fn exact_has_zero_ingest_cost_and_nonzero_query_cost() { - let candidate = CandidateConfig { - config: None, - query_method: QueryMethod::Exact, - n_windows: 0, - label_group_count: 1, - }; - let costs = AtomicCosts::default(); - let weights = CostWeights::default(); - let a = make_aqe(Statistic::Sum, 300_000, 300_000); - assert_eq!(ingest_cost(&candidate, 1.0, &costs, &weights), 0.0); - assert!(query_cost(&a, &candidate, &costs, &weights) > 0.0); - } - #[test] fn ingest_cost_independent_of_n_windows() { // Retained windows are a transient per-query cost (Mem_query in query_cost), diff --git a/asap-planner-rs/src/optimizer/dataset.rs b/asap-planner-rs/src/optimizer/dataset.rs index e6adcf92..3e31e475 100644 --- a/asap-planner-rs/src/optimizer/dataset.rs +++ b/asap-planner-rs/src/optimizer/dataset.rs @@ -13,7 +13,7 @@ use thiserror::Error; use crate::config::input::MetricDefinition; -use super::solution::AQE; +use super::solution::OptimizerItem; const METRIC_COLUMN: &str = "metric"; @@ -253,7 +253,10 @@ impl SeriesDataset { Ok(()) } - pub fn profile_aqes(&self, aqes: &[AQE]) -> Result, DatasetError> { + pub fn profile_aqes( + &self, + aqes: &[OptimizerItem], + ) -> Result, DatasetError> { let workload_metrics: HashSet<&str> = aqes .iter() .map(|aqe| aqe.requirements.metric.as_str()) @@ -553,12 +556,13 @@ mod tests { let dataset = SeriesDataset::from_reader("metric,job\nrequests,api\nother,worker\n".as_bytes()) .unwrap(); - let requests = AQE { + let requests = OptimizerItem { requirements: requirements("requests", &["job"], ""), query_strings: vec![], query_frequency_hz: 1.0, - min_t_repeat_ms: 1, - t_repeat_gcd_ms: 1, + t_repeat_ms: 1, + accuracy_sla: 0.0, + latency_sla: 0.0, }; assert!(matches!( diff --git a/asap-planner-rs/src/optimizer/error.rs b/asap-planner-rs/src/optimizer/error.rs new file mode 100644 index 00000000..9269c5c7 --- /dev/null +++ b/asap-planner-rs/src/optimizer/error.rs @@ -0,0 +1,29 @@ +use promql_utilities::query_logics::enums::Statistic; +use thiserror::Error; + +#[derive(Debug, Error)] +pub enum OptimizerError { + #[error("query '{query}' has repetition_delay_ms=0")] + InvalidRepeatInterval { query: String }, + + #[error("cannot optimize leaf '{leaf}' in query '{query}': {reason}")] + UnsupportedLeaf { + query: String, + 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: f64, + pub reason: String, +} diff --git a/asap-planner-rs/src/optimizer/greedy.rs b/asap-planner-rs/src/optimizer/greedy.rs index a533a559..ecfcb9ad 100644 --- a/asap-planner-rs/src/optimizer/greedy.rs +++ b/asap-planner-rs/src/optimizer/greedy.rs @@ -6,7 +6,8 @@ use super::atomic_costs::{resolve_atomic_costs, AtomicCostTable}; use super::candidate_gen::enumerate_candidates_with_label_group_count; use super::cost_model::{ingest_cost, query_cost, total_cost_rate, AtomicCosts, CostWeights}; use super::dataset::ProfileKey; -use super::solution::{AQEAssignment, OptimizerSolution, AQE}; +use super::error::{OptimizerError, UnservableItem}; +use super::solution::{AQEAssignment, OptimizerItem, OptimizerSolution}; /// Greedily assign each AQE to its independently-cheapest candidate config. /// @@ -23,14 +24,15 @@ use super::solution::{AQEAssignment, OptimizerSolution, AQE}; /// 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, + aqes: Vec, scrape_interval_ms: u64, arrival_rate_hz: f64, atomic_cost_table: &AtomicCostTable, weights: &CostWeights, label_group_counts: &HashMap, -) -> OptimizerSolution { +) -> Result { let mut solution = OptimizerSolution::empty(); + let mut unservable_items = Vec::new(); for aqe in aqes { let profile_key = ProfileKey::from_requirements(&aqe.requirements); @@ -46,7 +48,7 @@ pub fn greedy_assign( label_group_count, ); - let (best, costs) = candidates + let Some((best, costs)) = candidates .into_iter() .filter_map(|c| { // EXACT (config: None) always costs at the flat stub — it has @@ -65,16 +67,27 @@ pub fn greedy_assign( // 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)) - .expect( - "enumerate_candidates always returns at least the EXACT fallback, \ - which always resolves (flat stub, no table lookup)", - ); + 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: aqe.latency_sla, + 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 = best.config.map(|config| solution.register_config(config)); + let aggregation_id = solution.register_config( + best.config + .expect("candidate configs are streaming configs"), + ); debug!( metric = %aqe.requirements.metric, @@ -89,14 +102,20 @@ pub fn greedy_assign( solution.estimated_total_cost_per_sec += ingest + query_rate; solution.assignments.push(AQEAssignment { - aqe, + item: aqe, aggregation_id, query_method, estimated_query_cost_per_sec: query_rate, }); } - solution + if unservable_items.is_empty() { + Ok(solution) + } else { + Err(OptimizerError::UnservableItems { + items: unservable_items, + }) + } } #[cfg(test)] @@ -110,8 +129,8 @@ mod tests { 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) -> AQE { - AQE { + 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], @@ -123,8 +142,9 @@ mod tests { }, query_strings: vec!["test_query".into()], query_frequency_hz: freq_hz, - min_t_repeat_ms: min_t, - t_repeat_gcd_ms: min_t, + t_repeat_ms: min_t, + accuracy_sla: 0.0, + latency_sla: 0.0, } } @@ -154,7 +174,8 @@ mod tests { 1, ), ]), - ); + ) + .unwrap(); let mut seen_ids: StdHashMap = StdHashMap::new(); for id in solution.deployed_configs().keys() { @@ -167,8 +188,8 @@ mod tests { } #[test] - fn unsupported_multi_statistic_aqe_falls_back_to_exact() { - let aqe = AQE { + 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 @@ -180,10 +201,11 @@ mod tests { }, query_strings: vec!["avg_query".into()], query_frequency_hz: 1.0 / 60.0, - min_t_repeat_ms: 60_000, - t_repeat_gcd_ms: 60_000, + t_repeat_ms: 60_000, + accuracy_sla: 0.0, + latency_sla: 0.0, }; - let solution = greedy_assign( + let error = greedy_assign( vec![aqe.clone()], 60_000, 1.0, @@ -191,16 +213,15 @@ mod tests { &CostWeights::default(), &HashMap::from([(ProfileKey::from_requirements(&aqe.requirements), 1)]), ); - assert_eq!(solution.num_exact_fallback(), 1); - assert!(solution.deployed_configs().is_empty()); + assert!(matches!(error, Err(OptimizerError::UnservableItems { .. }))); } #[test] - fn missing_cms_with_heap_reference_cost_falls_back_to_exact() { + 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 solution = greedy_assign( + let error = greedy_assign( vec![aqe.clone()], 60_000, 1.0, @@ -209,8 +230,7 @@ mod tests { &HashMap::from([(ProfileKey::from_requirements(&aqe.requirements), 1)]), ); - assert_eq!(solution.num_exact_fallback(), 1); - assert!(solution.deployed_configs().is_empty()); + assert!(matches!(error, Err(OptimizerError::UnservableItems { .. }))); } #[test] @@ -235,9 +255,9 @@ mod tests { &table, &CostWeights::default(), &HashMap::from([(ProfileKey::from_requirements(&aqe.requirements), 1)]), - ); + ) + .unwrap(); - assert_eq!(solution.num_exact_fallback(), 0); assert_eq!(solution.deployed_configs().len(), 1); assert_eq!( solution diff --git a/asap-planner-rs/src/optimizer/mod.rs b/asap-planner-rs/src/optimizer/mod.rs index 521db371..d48e842f 100644 --- a/asap-planner-rs/src/optimizer/mod.rs +++ b/asap-planner-rs/src/optimizer/mod.rs @@ -4,6 +4,7 @@ pub mod candidate_gen; pub mod constants; pub mod cost_model; pub mod dataset; +pub mod error; pub mod greedy; pub mod pipeline; pub mod sketch_properties; @@ -21,8 +22,9 @@ pub use candidate_gen::{ }; pub use cost_model::{ingest_cost, query_cost, total_cost_rate, AtomicCosts, CostWeights}; pub use dataset::{DatasetError, ProfileKey, SeriesDataset}; +pub use error::{OptimizerError, UnservableItem}; pub use greedy::greedy_assign; -pub use pipeline::{run_all_exact_pipeline, run_greedy_pipeline}; +pub use pipeline::run_greedy_pipeline; pub use sketch_properties::{sketch_properties, SketchProperties}; -pub use solution::{AQEAssignment, OptimizerSolution, QueryMethod, AQE}; +pub use solution::{AQEAssignment, OptimizerItem, OptimizerSolution, QueryMethod}; pub use translator::translate; diff --git a/asap-planner-rs/src/optimizer/pipeline.rs b/asap-planner-rs/src/optimizer/pipeline.rs index 54387c70..9fc96909 100644 --- a/asap-planner-rs/src/optimizer/pipeline.rs +++ b/asap-planner-rs/src/optimizer/pipeline.rs @@ -1,6 +1,6 @@ use asap_types::inference_config::InferenceConfig; use asap_types::streaming_config::StreamingConfig; -use asap_types::PromQLSchema; +use thiserror::Error; use crate::config::input::ControllerConfig; @@ -9,23 +9,15 @@ use super::atomic_costs::AtomicCostTable; use super::cost_model::CostWeights; use super::dataset::SeriesDataset; use super::greedy::greedy_assign; -use super::solution::{OptimizerSolution, AQE}; +use super::solution::OptimizerSolution; use super::translator::{translate, TranslationSummary}; -/// Shared shell for optimizer pipelines: RQEs → AQEs → `solve` → deployment -/// artifacts, with a uniform log line. `solver_name` only affects the log. -fn run_pipeline( - config: &ControllerConfig, - schema: &PromQLSchema, - scrape_interval_ms: u64, - solver_name: &str, - solve: impl FnOnce(Vec) -> OptimizerSolution, -) -> (StreamingConfig, InferenceConfig) { - let rqes = config_to_rqes(config); - let aqes = extract_aqes(&rqes, schema, scrape_interval_ms); - let solution = solve(aqes); - - finish_pipeline(solution, solver_name) +#[derive(Debug, Error)] +pub enum OptimizerPipelineError { + #[error(transparent)] + Dataset(#[from] super::dataset::DatasetError), + #[error(transparent)] + Optimizer(#[from] super::error::OptimizerError), } fn finish_pipeline( @@ -37,7 +29,6 @@ fn finish_pipeline( solver = solver_name, num_deployed_configs = summary.num_deployed_configs, num_sketch_assignments = summary.num_sketch_assignments, - num_exact_fallbacks = summary.num_exact_fallbacks, 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" @@ -46,35 +37,13 @@ fn finish_pipeline( translate(&solution) } -/// Run the all-EXACT optimizer pipeline (Phase 1 scaffolding). -/// -/// Converts a `ControllerConfig` into `(StreamingConfig, InferenceConfig)` via -/// the optimizer path: RQEs → AQEs → all-EXACT solution → deployment artifacts. -/// -/// No streaming configs are deployed — every AQE falls back to raw data at -/// query time. This validates the end-to-end pipeline plumbing before real -/// sketch selection logic is added in Phase 2. -pub fn run_all_exact_pipeline( - config: &ControllerConfig, - schema: &PromQLSchema, - scrape_interval_ms: u64, -) -> (StreamingConfig, InferenceConfig) { - run_pipeline( - config, - schema, - scrape_interval_ms, - "all-EXACT", - OptimizerSolution::all_exact, - ) -} - -/// Run the greedy optimizer pipeline (Phase 2): each AQE is assigned, independently, -/// to its cheapest feasible candidate config (or to the EXACT fallback). +/// Run the greedy optimizer pipeline: each item is assigned independently to +/// its cheapest eligible streaming config. /// -/// No cross-AQE sharing — every deployed sketch serves exactly one AQE, even +/// 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 dataset supplies each AQE's metric schema and label-group count before +/// The dataset supplies each item's metric schema and label-group count before /// candidate selection. `arrival_rate_hz` remains a uniform placeholder until /// per-config scrape-rate data is available. pub fn run_greedy_pipeline( @@ -83,11 +52,11 @@ pub fn run_greedy_pipeline( scrape_interval_ms: u64, arrival_rate_hz: f64, atomic_cost_table: &AtomicCostTable, -) -> Result<(StreamingConfig, InferenceConfig), super::dataset::DatasetError> { +) -> Result<(StreamingConfig, InferenceConfig), OptimizerPipelineError> { dataset.validate_metric_hints(config.metrics.as_deref())?; let schema = dataset.schema(); let rqes = config_to_rqes(config); - let aqes = extract_aqes(&rqes, &schema, scrape_interval_ms); + let aqes = extract_aqes(&rqes, &schema, scrape_interval_ms)?; let label_group_counts = dataset.profile_aqes(&aqes)?; for (key, count) in &label_group_counts { @@ -107,7 +76,7 @@ pub fn run_greedy_pipeline( atomic_cost_table, &CostWeights::default(), &label_group_counts, - ); + )?; Ok(finish_pipeline(solution, "greedy")) } @@ -122,6 +91,8 @@ fn config_to_rqes(config: &ControllerConfig) -> Vec { 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: qg.controller_options.latency_sla, }) }) .collect() @@ -130,6 +101,7 @@ fn config_to_rqes(config: &ControllerConfig) -> Vec { #[cfg(test)] mod tests { use super::*; + use asap_types::PromQLSchema; fn make_config(queries: &[(&str, u64)]) -> ControllerConfig { use crate::config::input::QueryGroup; @@ -157,18 +129,6 @@ mod tests { } } - #[test] - fn all_exact_pipeline_produces_empty_streaming_config() { - let config = make_config(&[ - ("sum_over_time(metric[5m])", 60_000), - ("sum(other_metric)", 30_000), - ]); - let schema = PromQLSchema::new(); - let (streaming, _inference) = run_all_exact_pipeline(&config, &schema, 15_000); - // All-EXACT: no streaming configs deployed. - assert!(streaming.get_all_aggregation_configs().is_empty()); - } - #[test] fn greedy_pipeline_deploys_a_config_for_a_mergeable_aqe() { let config = make_config(&[("min_over_time(metric[5m])", 60_000)]); @@ -187,7 +147,7 @@ mod tests { assert!(matches!( run_greedy_pipeline(&config, &dataset, 60_000, 1.0, &AtomicCostTable::default()), - Err(super::super::dataset::DatasetError::MissingMetric(metric)) if metric == "metric" + Err(OptimizerPipelineError::Dataset(super::super::dataset::DatasetError::MissingMetric(metric))) if metric == "metric" )); } @@ -195,7 +155,7 @@ mod tests { 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); + 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); } diff --git a/asap-planner-rs/src/optimizer/solution.rs b/asap-planner-rs/src/optimizer/solution.rs index 34516226..e680c8e2 100644 --- a/asap-planner-rs/src/optimizer/solution.rs +++ b/asap-planner-rs/src/optimizer/solution.rs @@ -3,40 +3,34 @@ use std::collections::HashMap; use asap_types::aggregation_config::AggregationConfig; use asap_types::query_requirements::QueryRequirements; -/// An atomic query expression: one leaf aggregation extracted from a QE tree, -/// together with the optimizer-level metadata needed to assign it a config. +/// One optimizer demand item: an AQE's requirements at one cadence and SLA pair. #[derive(Debug, Clone)] -pub struct AQE { +pub struct OptimizerItem { /// What the query needs (metric, statistics, range, labels, spatial filter). pub requirements: QueryRequirements, - /// Original PromQL/SQL query strings from all RQEs that contain this AQE. + /// Original query strings from RQEs that contribute to this item. /// Preserved for use by the translator when building InferenceConfig. pub query_strings: Vec, - /// Query frequency in Hz: this field = Σ_{r ∈ R_a} 1/T_r. + /// 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, - /// Minimum repeat interval across all RQEs that reference this AQE (ms). - /// Determines the freshness constraint on the streaming config: for Tumbling, - /// W ≤ min_t_repeat ensures a completed window is available every cycle; for - /// Sliding, S ≤ min_t_repeat is the binding constraint (fresh answers arrive - /// every slide interval, not every W). When multiple RQEs share this AQE, - /// the fastest dashboard is the binding constraint. - pub min_t_repeat_ms: u64, - - /// GCD of all repeat intervals across RQEs that reference this AQE (ms). - /// The natural candidate for the slide interval S: windows that complete - /// every GCD ms align harmonically with all dashboard refresh cycles, - /// ensuring every dashboard can always be served a fresh result on-cycle. - pub t_repeat_gcd_ms: u64, + /// Repetition interval for every RQE contributing to this item, in ms. + pub t_repeat_ms: u64, + + /// Required accuracy for every RQE contributing to this item. + pub accuracy_sla: f64, + + /// Required query latency for every RQE contributing to this item. + pub latency_sla: f64, } -/// How an AQE is answered from its assigned streaming config. +/// 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. @@ -54,38 +48,31 @@ pub enum QueryMethod { /// W < range_a, sketch is subtractable: subtract two prefix-sum checkpoints. /// O(1) cost regardless of range_a/W. Subtract, - - /// No streaming config deployed for this AQE — query raw/exact data at - /// query time. Corresponds to the EXACT_a fallback (IngestCost = 0, Error = 0). - Exact, } -/// The assignment of a single AQE to a streaming config (or EXACT fallback). +/// The assignment of a single optimizer item to a streaming config. #[derive(Debug, Clone)] pub struct AQEAssignment { - pub aqe: AQE, + pub item: OptimizerItem, - /// ID of the deployed config that serves this AQE. - /// `None` means the EXACT_a fallback (no streaming config, raw query). - pub aggregation_id: Option, + /// ID of the deployed config that serves this item. + pub aggregation_id: u64, - /// How this AQE's answer is derived from the assigned config. + /// 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`. - /// Zero for Exact assignments (IngestCost is also zero). 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 -/// AQE to one of those configs (or to the EXACT fallback). A thin translator +/// 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. - /// Empty for all-EXACT solutions (Phase 1 scaffolding). /// /// 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 @@ -95,7 +82,7 @@ pub struct OptimizerSolution { /// Next id `register_config` will hand out. next_id: u64, - /// One entry per deduplicated AQE across the full RQE workload. + /// One entry per optimizer item across the full RQE workload. pub assignments: Vec, /// Estimated steady-state ingestion cost rate across all deployed configs @@ -134,36 +121,9 @@ impl OptimizerSolution { &self.deployed_configs } - /// Construct an all-EXACT solution: every AQE falls back to raw data, - /// no streaming configs are deployed. Used as the Phase 1 scaffolding baseline. - pub fn all_exact(aqes: Vec) -> Self { - let mut solution = Self::empty(); - solution.assignments = aqes - .into_iter() - .map(|aqe| AQEAssignment { - aqe, - aggregation_id: None, - query_method: QueryMethod::Exact, - estimated_query_cost_per_sec: 0.0, - }) - .collect(); - solution - } - - /// Number of AQEs served by an approximate sketch (not EXACT fallback). + /// Number of optimizer items served by a streaming sketch. pub fn num_sketch_served(&self) -> usize { - self.assignments - .iter() - .filter(|a| a.query_method != QueryMethod::Exact) - .count() - } - - /// Number of AQEs falling back to exact/raw computation. - pub fn num_exact_fallback(&self) -> usize { - self.assignments - .iter() - .filter(|a| a.query_method == QueryMethod::Exact) - .count() + self.assignments.len() } } diff --git a/asap-planner-rs/src/optimizer/translator.rs b/asap-planner-rs/src/optimizer/translator.rs index 8f5a07a5..05ea60ee 100644 --- a/asap-planner-rs/src/optimizer/translator.rs +++ b/asap-planner-rs/src/optimizer/translator.rs @@ -8,9 +8,6 @@ use super::solution::{OptimizerSolution, QueryMethod}; /// Translate an `OptimizerSolution` into the deployment artifacts consumed by /// Arroyo and the query engine. /// -/// Phase 1 (all-EXACT): deployed_configs is empty, all assignments are Exact, -/// so both output structs are empty/stub. Real translation logic fills in as -/// Phase 2/3 add sketch configs to the solution. pub fn translate(solution: &OptimizerSolution) -> (StreamingConfig, InferenceConfig) { let streaming_config = build_streaming_config(solution); let inference_config = build_inference_config(solution); @@ -27,17 +24,12 @@ fn build_inference_config(solution: &OptimizerSolution) -> InferenceConfig { let mut inference = InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::NoCleanup); - // For Phase 1 (all-EXACT), every assignment has aggregation_id = None, so - // this loop emits nothing — the inference engine falls back to raw - // querying for all AQEs, matching the all-EXACT solution. for assignment in &solution.assignments { - let Some(aggregation_id) = assignment.aggregation_id else { - continue; - }; + let aggregation_id = assignment.aggregation_id; let retain = retention_count_for_assignment(&assignment.query_method); let agg_ref = AggregationReference::new(aggregation_id, Some(retain)); - for query_string in &assignment.aqe.query_strings { + for query_string in &assignment.item.query_strings { inference .query_configs .push(QueryConfig::new(query_string.clone()).add_aggregation(agg_ref.clone())); @@ -59,7 +51,6 @@ pub fn retention_count_for_assignment(query_method: &QueryMethod) -> u64 { // (see candidate_gen.rs's n_windows, a separate concept: the deployed // AggregationConfig's retention depth). QueryMethod::Subtract => 2, - QueryMethod::Exact => 0, } } @@ -68,7 +59,6 @@ pub fn retention_count_for_assignment(query_method: &QueryMethod) -> u64 { pub struct TranslationSummary { pub num_deployed_configs: usize, pub num_sketch_assignments: usize, - pub num_exact_fallbacks: usize, } impl TranslationSummary { @@ -76,7 +66,6 @@ impl TranslationSummary { Self { num_deployed_configs: solution.deployed_configs().len(), num_sketch_assignments: solution.num_sketch_served(), - num_exact_fallbacks: solution.num_exact_fallback(), } } }