Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 7 additions & 8 deletions .design_docs/optimizer-mip-formulation.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |

Expand Down Expand Up @@ -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}$$
Expand Down Expand Up @@ -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$.

Expand Down
13 changes: 8 additions & 5 deletions asap-planner-rs/src/bin/candidate_gen_dump.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
83 changes: 82 additions & 1 deletion asap-planner-rs/src/config/input.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down Expand Up @@ -66,6 +66,7 @@ impl ControllerConfig {
pub struct QueryGroup {
pub id: Option<u32>,
pub queries: Vec<String>,
#[serde(deserialize_with = "deserialize_positive_u64")]
pub repetition_delay_ms: u64,
#[serde(default)]
pub controller_options: ControllerOptions,
Expand All @@ -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<f64, D::Error>
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<u64, D::Error>
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,
Expand Down Expand Up @@ -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::<ControllerConfig>(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::<ControllerConfig>(yaml)
.expect_err("zero repeat interval must be rejected")
.to_string();
assert!(error.contains("must be greater than zero"));
}
}
Loading
Loading