Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
690282d
feat(planner): add explicit query execution plans
milindsrivastava1997 Sep 26, 2026
8d0de9d
feat(planner): emit nested aggregation plans
milindsrivastava1997 Sep 27, 2026
c42ef98
feat(query-engine): validate query-time aggregation configs
milindsrivastava1997 Sep 27, 2026
827414c
test(planner): cover nested aggregation stages
milindsrivastava1997 Sep 27, 2026
ad9182f
feat(query-engine): execute query-time aggregation pipelines
milindsrivastava1997 Sep 28, 2026
3b2763d
test(query-engine): add nested aggregation e2e tracer
milindsrivastava1997 Sep 28, 2026
77b6d9c
test(query-engine): cover nested topk e2e
milindsrivastava1997 Sep 28, 2026
a95971c
test(query-engine): cover nested aggregation operators e2e
milindsrivastava1997 Sep 28, 2026
467d149
test(query-engine): assert nested aggregation e2e values
milindsrivastava1997 Sep 28, 2026
cf48026
test(promql): add nested aggregation differential suite
milindsrivastava1997 Sep 28, 2026
f7c3b87
test(query-engine): cover query-time aggregation grouping
milindsrivastava1997 Sep 29, 2026
0f74902
fix(query-config): reject zero topk limits
milindsrivastava1997 Sep 29, 2026
1313532
test(query-engine): cover query-time aggregation boundaries
milindsrivastava1997 Sep 29, 2026
ff280a0
test(query-engine): cover query-time topk ties
milindsrivastava1997 Sep 29, 2026
bbc7f0a
feat(query-engine): execute nested aggregation DAG nodes
milindsrivastava1997 Oct 4, 2026
b1e5cc9
fix(query-engine): preserve missing query-time labels
milindsrivastava1997 Oct 4, 2026
cd5f31e
fix(planner): preserve empty without aggregations
milindsrivastava1997 Oct 4, 2026
35a424f
fix(query-engine): apply nested topk anchors first
milindsrivastava1997 Oct 4, 2026
072ef1e
fix(query-engine): validate nested aggregation plans
milindsrivastava1997 Oct 4, 2026
b0a14f9
fix(query-engine): preserve nested topk labels
milindsrivastava1997 Oct 4, 2026
c0eeb1d
Merge remote-tracking branch 'origin/main' into 738-nested-promql-agg…
milindsrivastava1997 Oct 4, 2026
353d828
fix(query-engine): keep grouped topk buckets contiguous
milindsrivastava1997 Oct 4, 2026
80abed1
Merge remote-tracking branch 'origin/main' into 738-nested-promql-agg…
milindsrivastava1997 Oct 5, 2026
d450f5c
fix(query-engine): align raw topk labels
milindsrivastava1997 Oct 5, 2026
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
95 changes: 93 additions & 2 deletions asap-common/dependencies/rs/asap_types/src/inference_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ use std::io::BufReader;
use crate::aggregation_reference::AggregationReference;
use crate::enums::{CleanupPolicy, QueryLanguage};
use crate::promql_schema::PromQLSchema;
use crate::query_config::QueryConfig;
use crate::query_config::{QueryConfig, QueryTimeAggregation};
use elastic_dsl_utilities::{ElasticIndexSchema, ElasticMappingSchema};
use promql_utilities::data_model::KeyByLabelNames;
use sql_utilities::sqlhelper::{SQLSchema, Table};
Expand Down Expand Up @@ -250,6 +250,18 @@ impl InferenceConfig {
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing query field"))?
.to_string();
let planned_subquery = query_data
.get("planned_subquery")
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing planned_subquery field"))?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Medium: making planned_subquery / query_time_aggregations required breaks the checked-in inference configs.

These configs have neither field and now fail with "Missing planned_subquery field":

  • asap-query-engine/examples/promql/inference_config.yaml
  • asap-query-engine/examples/sql/inference_config.yaml
  • asap-tools/execution-utilities/elastic-asap-benchmarking/inference_config.yaml

src/bin/test_e2e_precompute.rs loads the PromQL example directly, so that binary breaks right away. Older planner outputs fail the same way.

Fix: add the fields to the configs, or default them (planned_subquery = query, empty pipeline).

.to_string();
let query_time_aggregations = query_data
.get("query_time_aggregations")
.ok_or_else(|| anyhow::anyhow!("Missing query_time_aggregations field"))
.and_then(|value| {
serde_yaml::from_value::<Vec<QueryTimeAggregation>>(value.clone())
.map_err(anyhow::Error::from)
})?;

let aggregations = if let Some(aggregations_data) =
query_data.get("aggregations").and_then(|v| v.as_sequence())
Expand Down Expand Up @@ -290,7 +302,12 @@ impl InferenceConfig {
Vec::new()
};

let config = QueryConfig::new(query).with_aggregations(aggregations);
let config =
QueryConfig::with_plan(query, planned_subquery, query_time_aggregations)
.with_aggregations(aggregations);
config
.validate_execution_plan()
.map_err(|error| anyhow::anyhow!("Invalid query execution plan: {error}"))?;
configs.push(config);
}
configs
Expand All @@ -300,3 +317,77 @@ impl InferenceConfig {
Ok(query_configs)
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn rejects_a_query_without_a_planned_subquery() {
let data: Value = serde_yaml::from_str(
r#"
cleanup_policy:
name: no_cleanup
metrics: {}
queries:
- query: "sum(metric)"
query_time_aggregations: []
aggregations: []
"#,
)
.unwrap();

let error = InferenceConfig::from_yaml_data(&data, QueryLanguage::promql)
.expect_err("query plans must name their planned subquery");

assert!(error.to_string().contains("planned_subquery"));
}

#[test]
fn rejects_a_query_without_a_query_time_pipeline() {
let data: Value = serde_yaml::from_str(
r#"
cleanup_policy:
name: no_cleanup
metrics: {}
queries:
- query: "sum(metric)"
planned_subquery: "sum(metric)"
aggregations: []
"#,
)
.unwrap();

let error = InferenceConfig::from_yaml_data(&data, QueryLanguage::promql)
.expect_err("query plans must declare their query-time pipeline");

assert!(error.to_string().contains("query_time_aggregations"));
}

#[test]
fn rejects_an_invalid_query_time_aggregation() {
let data: Value = serde_yaml::from_str(
r#"
cleanup_policy:
name: no_cleanup
metrics: {}
queries:
- query: "quantile(1.5, sum(metric))"
planned_subquery: "sum(metric)"
query_time_aggregations:
- operator: quantile
grouping:
mode: all
labels: []
parameter: 1.5
aggregations: []
"#,
)
.unwrap();

let error = InferenceConfig::from_yaml_data(&data, QueryLanguage::promql)
.expect_err("invalid pipeline parameters must fail during config loading");

assert!(error.to_string().contains("quantile"));
}
}
181 changes: 181 additions & 0 deletions asap-common/dependencies/rs/asap_types/src/query_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,17 +5,133 @@ use crate::aggregation_reference::AggregationReference;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QueryConfig {
pub query: String,
pub planned_subquery: String,
pub query_time_aggregations: Vec<QueryTimeAggregation>,
pub aggregations: Vec<AggregationReference>,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum QueryTimeAggregationOperator {
Sum,
Count,
Avg,
Min,
Max,
Quantile,
Topk,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum QueryTimeGroupingMode {
All,
By,
Without,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QueryTimeGrouping {
pub mode: QueryTimeGroupingMode,
pub labels: Vec<String>,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(untagged)]
pub enum QueryTimeAggregationParameter {
Integer(u64),
Float(f64),
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QueryTimeAggregation {
pub operator: QueryTimeAggregationOperator,
pub grouping: QueryTimeGrouping,
pub parameter: Option<QueryTimeAggregationParameter>,
}

impl QueryTimeAggregation {
pub fn validate(&self) -> Result<(), String> {
match self.grouping.mode {
QueryTimeGroupingMode::All if !self.grouping.labels.is_empty() => {
return Err("all grouping cannot name labels".to_string());
}
QueryTimeGroupingMode::By if self.grouping.labels.is_empty() => {
return Err("by grouping must name at least one label".to_string());
}
_ => {}
}

if self.grouping.labels.iter().any(|label| label.is_empty()) {
return Err("grouping labels cannot be empty".to_string());
}
let unique_label_count = self
.grouping
.labels
.iter()
.collect::<std::collections::HashSet<_>>()
.len();
if unique_label_count != self.grouping.labels.len() {
return Err("grouping labels must be unique".to_string());
}

match (&self.operator, &self.parameter) {
(
QueryTimeAggregationOperator::Topk,
Some(QueryTimeAggregationParameter::Integer(k)),
) if *k > 0 => {}
(
QueryTimeAggregationOperator::Topk,
Some(QueryTimeAggregationParameter::Integer(_)),
) => {
return Err("topk requires a positive integer parameter".to_string());
}
(QueryTimeAggregationOperator::Topk, _) => {
return Err("topk requires an integer parameter".to_string());
}
(
QueryTimeAggregationOperator::Quantile,
Some(QueryTimeAggregationParameter::Float(phi)),
) if phi.is_finite() && (0.0..=1.0).contains(phi) => {}
(QueryTimeAggregationOperator::Quantile, _) => {
return Err("quantile requires a finite parameter between 0 and 1".to_string());
}
(_, None) => {}
_ => return Err("this aggregation does not accept a parameter".to_string()),
}

Ok(())
}
}

impl QueryConfig {
pub fn new(query: String) -> Self {
Self::with_plan(query.clone(), query, Vec::new())
}

pub fn with_plan(
query: String,
planned_subquery: String,
query_time_aggregations: Vec<QueryTimeAggregation>,
) -> Self {
Self {
query,
planned_subquery,
query_time_aggregations,
aggregations: Vec::new(),
}
}

pub fn validate_execution_plan(&self) -> Result<(), String> {
if self.planned_subquery.trim().is_empty() {
return Err("planned_subquery cannot be empty".to_string());
}
for aggregation in &self.query_time_aggregations {
aggregation.validate()?;
}
Ok(())
}

pub fn add_aggregation(mut self, aggregation: AggregationReference) -> Self {
self.aggregations.push(aggregation);
self
Expand All @@ -26,3 +142,68 @@ impl QueryConfig {
self
}
}

#[cfg(test)]
mod tests {
use super::*;

fn topk(k: u64) -> QueryTimeAggregation {
QueryTimeAggregation {
operator: QueryTimeAggregationOperator::Topk,
grouping: QueryTimeGrouping {
mode: QueryTimeGroupingMode::All,
labels: Vec::new(),
},
parameter: Some(QueryTimeAggregationParameter::Integer(k)),
}
}

#[test]
fn topk_requires_a_positive_k() {
assert!(topk(1).validate().is_ok());
assert_eq!(
topk(0).validate(),
Err("topk requires a positive integer parameter".into())
);
}

#[test]
fn quantile_accepts_endpoints_and_rejects_out_of_range_values() {
for phi in [0.0, 1.0] {
assert!(QueryTimeAggregation {
operator: QueryTimeAggregationOperator::Quantile,
grouping: QueryTimeGrouping {
mode: QueryTimeGroupingMode::All,
labels: Vec::new(),
},
parameter: Some(QueryTimeAggregationParameter::Float(phi)),
}
.validate()
.is_ok());
}
assert!(QueryTimeAggregation {
operator: QueryTimeAggregationOperator::Quantile,
grouping: QueryTimeGrouping {
mode: QueryTimeGroupingMode::All,
labels: Vec::new(),
},
parameter: Some(QueryTimeAggregationParameter::Float(1.01)),
}
.validate()
.is_err());
}

#[test]
fn without_empty_labels_is_a_valid_promql_grouping() {
assert!(QueryTimeAggregation {
operator: QueryTimeAggregationOperator::Sum,
grouping: QueryTimeGrouping {
mode: QueryTimeGroupingMode::Without,
labels: Vec::new(),
},
parameter: None,
}
.validate()
.is_ok());
}
}
19 changes: 11 additions & 8 deletions asap-planner-rs/src/elastic_dsl/generator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ use std::collections::HashMap;
use crate::config::input::ElasticDSLControllerConfig;
use crate::error::ControllerError;
use crate::generator::{
build_aggregation_entry, build_queries_yaml, GeneratorOutput, KEY_AGGREGATIONS,
build_aggregation_entry, build_queries_yaml, GeneratorOutput, QueryPlanEntry, KEY_AGGREGATIONS,
KEY_CLEANUP_POLICY, KEY_NAME, KEY_QUERIES,
};
use crate::planner::agg_config::IntermediateAggConfig;
Expand Down Expand Up @@ -82,8 +82,8 @@ pub fn generate_elastic_plan(

// Dedup map: identifying_key -> IntermediateAggConfig
let mut dedup_map: IndexMap<String, IntermediateAggConfig> = IndexMap::new();
// query_string -> Vec<(key, cleanup_param)>
let mut query_keys_map: IndexMap<String, Vec<(String, Option<u64>)>> = IndexMap::new();
// query_string -> explicit physical query plan
let mut query_plan_map: IndexMap<String, QueryPlanEntry> = IndexMap::new();
// index -> schema builder derived from the queries targeting that index
let mut index_schema_builders: IndexMap<String, ElasticIndexSchemaBuilder> = IndexMap::new();

Expand Down Expand Up @@ -127,7 +127,10 @@ pub fn generate_elastic_plan(
keys_for_query.push((key.clone(), cleanup_param));
dedup_map.entry(key).or_insert(config_item);
}
query_keys_map.insert(query_string.clone(), keys_for_query);
query_plan_map.insert(
query_string.clone(),
QueryPlanEntry::fully_planned(query_string.clone(), keys_for_query),
);
}
}

Expand All @@ -140,7 +143,7 @@ pub fn generate_elastic_plan(
let streaming_yaml = build_elastic_streaming_yaml(&dedup_map, &id_map)?;
let inference_yaml = build_elastic_inference_yaml(
cleanup_policy,
&query_keys_map,
&query_plan_map,
&id_map,
&index_schema_builders,
)?;
Expand All @@ -150,7 +153,7 @@ pub fn generate_elastic_plan(
streaming_yaml,
inference_yaml,
aggregation_count: dedup_map.len(),
query_count: query_keys_map.len(),
query_count: query_plan_map.len(),
})
}

Expand All @@ -174,7 +177,7 @@ fn build_elastic_streaming_yaml(

fn build_elastic_inference_yaml(
cleanup_policy: CleanupPolicy,
query_keys_map: &IndexMap<String, Vec<(String, Option<u64>)>>,
query_plan_map: &IndexMap<String, QueryPlanEntry>,
id_map: &HashMap<String, u32>,
index_schema_builders: &IndexMap<String, ElasticIndexSchemaBuilder>,
) -> Result<YamlValue, ControllerError> {
Expand All @@ -191,7 +194,7 @@ fn build_elastic_inference_yaml(
);
root.insert(
YamlValue::String(KEY_QUERIES.to_string()),
YamlValue::Sequence(build_queries_yaml(cleanup_policy, query_keys_map, id_map)),
YamlValue::Sequence(build_queries_yaml(cleanup_policy, query_plan_map, id_map)),
);
root.insert(
YamlValue::String("indices".to_string()),
Expand Down
Loading
Loading