Skip to content
Closed
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
4 changes: 4 additions & 0 deletions asap-dropin/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions asap-dropin/config/engine_config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -21,3 +21,5 @@ ingest:
query_tracker:
enabled: true
observation_window_secs: 60
accuracy_sla: 0.99
latency_sla: 0.0
85 changes: 60 additions & 25 deletions asap-planner-rs/src/config/input.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down Expand Up @@ -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 {
Expand All @@ -68,7 +51,6 @@ pub struct QueryGroup {
pub queries: Vec<String>,
#[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)]
Expand All @@ -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<f64, D::Error>
fn deserialize_accuracy_sla<'de, D>(deserializer: D) -> Result<f64, D::Error>
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<f64, D::Error>
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",
))
}
}

Expand Down Expand Up @@ -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)]
Expand All @@ -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]
Expand All @@ -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::<ControllerConfig>(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::<ControllerConfig>(&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::<ControllerConfig>(yaml)
.expect_err("controller options are required")
.to_string();
assert!(error.contains("controller_options"));
}

#[test]
fn rejects_zero_repetition_delay() {
let yaml = r#"
Expand Down
73 changes: 58 additions & 15 deletions asap-planner-rs/src/optimizer/atomic_costs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,13 @@ pub struct AtomicCostEntry {

pub type AtomicCostTable = Vec<AtomicCostEntry>;

/// Costs and benchmarked accuracy measurements resolved for one candidate.
#[derive(Debug, Clone, PartialEq)]
pub struct ResolvedAtomicCosts {
pub costs: AtomicCosts,
pub query_accuracy: BTreeMap<String, f64>,
}

/// 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.
Expand Down Expand Up @@ -254,20 +261,38 @@ pub fn resolve_atomic_costs(
params: &HashMap<String, Value>,
n_grouping_labels: usize,
) -> Option<AtomicCosts> {
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<String, Value>,
n_grouping_labels: usize,
) -> Option<ResolvedAtomicCosts> {
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) {
tracing::warn!(
?agg_type,
"no sketch-bench CPU data for this family; using stub CPU costs and analytical memory"
);
return Some(AtomicCosts {
mem_bytes_per_instance: (n_grouping_labels as f64 * LABEL_VALUE_CODE_BYTES
+ value_bytes)
* HASH_TABLE_SLACK,
..AtomicCosts::default()
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(),
});
}

Expand All @@ -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(),
})
}

Expand Down Expand Up @@ -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<String, Value>,
assumptions: &CmsHeapCostAssumptions,
) -> Option<AtomicCosts> {
resolve_cms_heap_costs_with_accuracy(table, params, assumptions).map(|resolved| resolved.costs)
}

fn resolve_cms_heap_costs_with_accuracy(
table: &AtomicCostTable,
params: &HashMap<String, Value>,
assumptions: &CmsHeapCostAssumptions,
) -> Option<ResolvedAtomicCosts> {
assumptions.validate();

let agg_type = AggregationType::CountMinSketchWithHeap;
Expand Down Expand Up @@ -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<String, Value>, key: &str, agg_type: AggregationType) -> u64 {
Expand Down
52 changes: 36 additions & 16 deletions asap-planner-rs/src/optimizer/cost_model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading