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
408 changes: 404 additions & 4 deletions controller/Cargo.lock

Large diffs are not rendered by default.

2 changes: 2 additions & 0 deletions controller/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@ thiserror = "1"
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt"] }
chrono = { version = "0.4", features = ["serde"] }
sqlparser = "0.61"
promql-parser = "0.8"

[dev-dependencies]
tokio = { version = "1", features = ["full", "test-util"] }
232 changes: 207 additions & 25 deletions controller/src/analyzer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,30 +3,53 @@ use std::time::Duration;
use anyhow::{anyhow, Context};
use serde::{Deserialize, Serialize};

use crate::query_parser;
use crate::types::{AggType, QueryWorkload, SketchType, WorkloadCharacteristics};

// ── Public API ────────────────────────────────────────────────────────────────

/// JSON-friendly representation of a query workload submitted by callers.
///
/// There are two ways to populate a `QuerySpec`:
///
/// 1. **Explicit fields** — supply `metric_name`, `aggregations`,
/// `time_window`, etc. directly. This is the original API.
///
/// 2. **Query string** — supply a raw PromQL or SQL string in
/// `query_string`. The analyzer parses it and fills in `metric_name`,
/// `aggregations`, `group_by_labels`, `label_filters`, and `time_window`
/// automatically. Any explicit fields that are non-empty / non-default
/// **override** the parsed values, so the two approaches compose.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QuerySpec {
/// Raw PromQL or SQL query string to parse (SP-1 automatic extraction).
/// When provided, metric_name / aggregations / time_window may be omitted
/// and will be derived from the query.
#[serde(default)]
pub query_string: Option<String>,

/// Metric name override. Required when `query_string` is absent.
#[serde(default)]
pub metric_name: String,
#[serde(default)]
pub label_filters: HashMap<String, String>,
#[serde(default)]
pub group_by_labels: Vec<String>,
/// Aggregation type overrides ("quantile", "cardinality", "frequency").
/// Required when `query_string` is absent.
#[serde(default)]
pub aggregations: Vec<String>,
/// Time window override (e.g. "5m"). Required when `query_string` is absent.
#[serde(default)]
pub time_window: String,
#[serde(default)]
pub repeat_every: Option<String>,
pub accuracy_sla: f64,
pub latency_sla: Option<String>,
/// Optional: pin a specific sketch type, bypassing the cost-model planner.
/// Useful when the target collector supports only a subset of sketches.
pub sketch_type: Option<SketchType>,
/// Observable data-stream characteristics used for delta transmission
/// decisions and raw-vs-sketch bandwidth comparison.
/// Omit to use conservative defaults (1 000 series, 100 Hz, 100 B/sample,
/// Zipf distribution, no memory budget).
/// Observable data-stream characteristics used for delta / raw-vs-sketch
/// bandwidth comparison. Omit to use conservative defaults.
#[serde(default)]
pub workload: WorkloadCharacteristics,
}
Expand All @@ -37,24 +60,75 @@ impl Analyzer {
pub fn new() -> Self { Self }

pub fn analyze(&self, spec: QuerySpec) -> anyhow::Result<QueryWorkload> {
if spec.metric_name.trim().is_empty() {
return Err(anyhow!("metric_name is required"));
}
if spec.aggregations.is_empty() {
return Err(anyhow!("at least one aggregation is required"));
}
if !(0.0..=1.0).contains(&spec.accuracy_sla) {
return Err(anyhow!("accuracy_sla must be in [0,1], got {}", spec.accuracy_sla));
}

let aggs = parse_agg_types(&spec.aggregations)?;
// ── Step 1: parse query_string if provided ─────────────────────────
let parsed = spec.query_string.as_deref()
.map(|q| query_parser::parse_query(q))
.transpose()
.with_context(|| "failed to parse query_string")?;

let time_window = parse_duration(&spec.time_window)
.with_context(|| format!("invalid time_window {:?}", spec.time_window))?;
if time_window.is_zero() {
return Err(anyhow!("time_window must be positive"));
}
// ── Step 2: resolve metric_name ────────────────────────────────────
let metric_name = if !spec.metric_name.trim().is_empty() {
spec.metric_name.clone()
} else if let Some(ref p) = parsed {
p.metric_name.clone()
} else {
return Err(anyhow!(
"metric_name is required (or provide query_string)"
));
};

// ── Step 3: resolve aggregations ───────────────────────────────────
let aggregations = if !spec.aggregations.is_empty() {
parse_agg_types(&spec.aggregations)?
} else if let Some(ref p) = parsed {
if p.aggregations.is_empty() && !p.exact_required {
return Err(anyhow!(
"could not infer aggregation type from query_string; \
provide explicit aggregations"
));
}
p.aggregations.clone()
} else {
return Err(anyhow!("at least one aggregation is required"));
};

// ── Step 4: resolve time_window ────────────────────────────────────
let time_window = if !spec.time_window.trim().is_empty() {
let d = parse_duration(&spec.time_window)
.with_context(|| format!("invalid time_window {:?}", spec.time_window))?;
if d.is_zero() {
return Err(anyhow!("time_window must be positive"));
}
d
} else if let Some(ref p) = parsed {
p.time_window
} else {
return Err(anyhow!("time_window is required (or provide query_string)"));
};

// ── Step 5: resolve dimensions (group_by + label_filter keys) ──────
// Parsed values are the base; explicit spec fields override / extend.
let parsed_group_by = parsed.as_ref().map(|p| p.group_by_labels.as_slice()).unwrap_or(&[]);
let parsed_filters: HashMap<String, String> =
parsed.as_ref().map(|p| p.label_filters.clone()).unwrap_or_default();

let merged_filters: HashMap<String, String> = {
let mut m = parsed_filters;
m.extend(spec.label_filters.clone()); // explicit overrides parsed
m
};

let filter_keys: Vec<String> = merged_filters.keys().cloned().collect();
let all_group_by: Vec<String> = dedup_dims(
&dedup_dims(parsed_group_by, &spec.group_by_labels),
&filter_keys,
);

// ── Step 6: scalar fields ──────────────────────────────────────────
let repeat_every = spec.repeat_every.as_deref()
.map(parse_duration)
.transpose()
Expand All @@ -65,21 +139,21 @@ impl Analyzer {
.transpose()
.with_context(|| "invalid latency_sla")?;

// Merge group_by_labels and label_filter keys, deduplicating while
// preserving the group_by_labels order first.
let filter_keys: Vec<String> = spec.label_filters.keys().cloned().collect();
let dims = dedup_dims(&spec.group_by_labels, &filter_keys);
let exact_required = parsed.as_ref().map(|p| p.exact_required).unwrap_or(false);
let quantiles = parsed.as_ref().map(|p| p.quantiles.clone()).unwrap_or_default();

Ok(QueryWorkload {
metric_name: spec.metric_name,
label_filters: spec.label_filters,
group_by_labels: dims,
aggregations: aggs,
metric_name,
label_filters: merged_filters,
group_by_labels: all_group_by,
aggregations,
time_window,
repeat_every,
accuracy_sla: spec.accuracy_sla,
latency_sla,
sketch_type_override: spec.sketch_type,
exact_required,
quantiles,
})
}
}
Expand Down Expand Up @@ -158,6 +232,7 @@ mod tests {

fn basic_spec() -> QuerySpec {
QuerySpec {
query_string: None,
metric_name: "request_latency".into(),
label_filters: [("service".into(), "web".into())].into(),
group_by_labels: vec!["host.name".into()],
Expand Down Expand Up @@ -267,4 +342,111 @@ mod tests {
fn trailing_digits_error() {
assert!(parse_duration("5").is_err());
}

// ── query_string path ─────────────────────────────────────────────────────

/// Build a minimal QuerySpec driven entirely by a query_string.
fn qs_only(query: &str) -> QuerySpec {
QuerySpec {
query_string: Some(query.into()),
metric_name: "".into(),
label_filters: Default::default(),
group_by_labels: vec![],
aggregations: vec![],
time_window: "".into(),
repeat_every: None,
accuracy_sla: 0.01,
latency_sla: None,
sketch_type: None,
workload: Default::default(),
}
}

/// PromQL query_string auto-populates metric_name, aggregations,
/// time_window, and quantiles — no explicit fields required.
#[test]
fn query_string_promql_populates_workload() {
let w = Analyzer::new()
.analyze(qs_only("sum by (host) (quantile_over_time(0.99, latency[5m]))"))
.unwrap();
assert_eq!(w.metric_name, "latency");
assert_eq!(w.aggregations, vec![AggType::Quantile]);
assert_eq!(w.time_window, Duration::from_secs(300));
assert_eq!(w.quantiles, vec![0.99]);
assert!(!w.exact_required);
}

/// SQL query_string auto-populates metric_name, aggregations,
/// and group_by_labels.
#[test]
fn query_string_sql_populates_workload() {
let w = Analyzer::new()
.analyze(qs_only(
"SELECT symbol, COUNT(*) FROM financial_last_trade_price GROUP BY symbol",
))
.unwrap();
assert_eq!(w.metric_name, "financial_last_trade_price");
assert_eq!(w.aggregations, vec![AggType::Frequency]);
assert!(w.group_by_labels.contains(&"symbol".to_string()));
}

/// Explicit metric_name overrides the name derived from query_string.
#[test]
fn explicit_metric_name_overrides_parsed() {
let mut spec = qs_only("sum by (host) (avg_over_time(cpu[5m]))");
spec.metric_name = "my_custom_metric".into();
let w = Analyzer::new().analyze(spec).unwrap();
assert_eq!(w.metric_name, "my_custom_metric");
// aggregations still come from parse (avg → DDSketch → Quantile)
assert_eq!(w.aggregations, vec![AggType::Quantile]);
}

/// Explicit time_window overrides the window derived from query_string.
#[test]
fn explicit_time_window_overrides_parsed() {
let mut spec = qs_only("sum by (host) (avg_over_time(cpu[5m]))");
spec.time_window = "1h".into();
let w = Analyzer::new().analyze(spec).unwrap();
assert_eq!(w.time_window, Duration::from_secs(3600));
}

/// Explicit aggregations override those derived from query_string.
#[test]
fn explicit_aggregations_override_parsed() {
let mut spec = qs_only("sum by (host) (avg_over_time(cpu[5m]))"); // → Quantile
spec.aggregations = vec!["cardinality".into()];
let w = Analyzer::new().analyze(spec).unwrap();
assert_eq!(w.aggregations, vec![AggType::Cardinality]);
}

/// sum_over_time is a stateful exact aggregation; exact_required is set.
#[test]
fn query_string_exact_required_propagated() {
let w = Analyzer::new()
.analyze(qs_only("sum by (service) (sum_over_time(request_bytes[1h]))"))
.unwrap();
assert!(w.exact_required, "sum_over_time must set exact_required");
assert_eq!(w.aggregations, vec![]);
}

/// DDSketch quantile φ values are surfaced through the workload.
#[test]
fn query_string_quantiles_populated() {
let w = Analyzer::new()
.analyze(qs_only("sum by (host) (quantile_over_time(0.5, latency[5m]))"))
.unwrap();
assert_eq!(w.quantiles, vec![0.5]);
}

/// Existing callers that supply all fields explicitly and omit
/// query_string continue to work unchanged (backward compatibility).
#[test]
fn backward_compat_no_query_string() {
let w = Analyzer::new().analyze(basic_spec()).unwrap();
assert_eq!(w.metric_name, "request_latency");
assert_eq!(w.aggregations, vec![AggType::Quantile]);
assert_eq!(w.time_window, Duration::from_secs(300));
assert!(!w.exact_required);
assert!(w.quantiles.is_empty());
}
}
2 changes: 2 additions & 0 deletions controller/src/config/precompute.rs
Original file line number Diff line number Diff line change
Expand Up @@ -168,6 +168,8 @@ mod tests {
accuracy_sla: 0.01,
latency_sla,
sketch_type_override: None,
exact_required: false,
quantiles: vec![],
}
}

Expand Down
1 change: 1 addition & 0 deletions controller/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ mod config;
mod monitor;
mod opamp;
mod planner;
mod query_parser;
mod store;
mod types;

Expand Down
2 changes: 2 additions & 0 deletions controller/src/planner/cost_model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -285,6 +285,8 @@ mod tests {
accuracy_sla: 0.01,
latency_sla: None,
sketch_type_override: None,
exact_required: false,
quantiles: vec![],
}
}

Expand Down
2 changes: 2 additions & 0 deletions controller/src/planner/delta_cost_model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -472,6 +472,8 @@ mod tests {
accuracy_sla: 0.01,
latency_sla: None,
sketch_type_override: None,
exact_required: false,
quantiles: vec![],
}
}

Expand Down
Loading