From 337ae1f44e503661694ffb68d22442e6bab2ef6e Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Mon, 5 Oct 2026 11:41:11 -0400 Subject: [PATCH 1/2] refactor(query-engine): make native DAG execution self-contained --- asap-query-engine/src/engines/query_plan.rs | 238 ++++++++++++------ .../src/engines/simple_engine/mod.rs | 191 +++++++------- 2 files changed, 270 insertions(+), 159 deletions(-) diff --git a/asap-query-engine/src/engines/query_plan.rs b/asap-query-engine/src/engines/query_plan.rs index 2384b75e..1257f1c1 100644 --- a/asap-query-engine/src/engines/query_plan.rs +++ b/asap-query-engine/src/engines/query_plan.rs @@ -5,30 +5,78 @@ use crate::engines::simple_engine::{RangeQueryExecutionContext, StoreQueryParams use asap_types::enums::WindowType; use asap_types::query_config::QueryTimeAggregation; use promql_utilities::data_model::KeyByLabelNames; -use promql_utilities::query_logics::enums::Statistic; +use promql_utilities::query_logics::enums::{AggregationType, Statistic}; use tracing::debug; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) struct NodeId(usize); -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, PartialEq, Eq)] pub(crate) enum StoreReadStrategy { WindowGrid, - SlidingExactCover, + SlidingExactCover { + output_timestamps: Vec, + lookback_ms: u64, + window_size_ms: u64, + bucket_step_ms: u64, + }, +} + +impl StoreReadStrategy { + fn for_window( + window_type: WindowType, + output_timestamps: &[u64], + lookback_ms: u64, + window_size_ms: u64, + bucket_step_ms: u64, + ) -> Self { + match window_type { + WindowType::Tumbling => Self::WindowGrid, + WindowType::Sliding => Self::SlidingExactCover { + output_timestamps: output_timestamps.to_vec(), + lookback_ms, + window_size_ms, + bucket_step_ms, + }, + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum StoreReadRole { + Values, + Keys, +} + +#[derive(Debug, Clone)] +pub(crate) struct RangeEstimateSpec { + pub output_timestamps: Vec, + pub query_range_ms: u64, + pub buckets_per_step: usize, + pub lookback_bucket_count: usize, + pub tumbling_window_ms: u64, + pub window_type: WindowType, + pub window_size_ms: u64, + pub keys_window_type: Option, + pub keys_window_size_ms: Option, + pub keys_lookback_ms: Option, + pub keys_tumbling_window_ms: Option, + pub value_aggregation_type: AggregationType, + pub key_aggregation_type: AggregationType, + pub grouping_labels: KeyByLabelNames, + pub aggregated_labels: KeyByLabelNames, + pub row_label_order: KeyByLabelNames, } #[derive(Debug, Clone)] pub(crate) enum QueryPlanNode { StoreRead { query: StoreQueryParams, + role: StoreReadRole, strategy: StoreReadStrategy, }, - ComposeWindows { + PrepareBuckets { input: NodeId, - output_timestamps: Vec, - lookback_ms: u64, - window_size_ms: u64, - bucket_step_ms: u64, }, ResolveKeys { values: NodeId, @@ -39,6 +87,7 @@ pub(crate) enum QueryPlanNode { statistic: Statistic, query_kwargs: std::collections::HashMap, output_labels: KeyByLabelNames, + spec: RangeEstimateSpec, }, AggregateVector { input: NodeId, @@ -49,6 +98,7 @@ pub(crate) enum QueryPlanNode { input: NodeId, k: String, grouping_labels: KeyByLabelNames, + row_label_order: KeyByLabelNames, }, Format { input: NodeId, @@ -113,39 +163,45 @@ impl QueryPlan { query_time_aggregations: &[QueryTimeAggregation], ) -> Result { let mut nodes = Vec::new(); + let values_lookback_ms = + (context.lookback_bucket_count as u64) * context.tumbling_window_ms; let values_read = Self::push_read( &mut nodes, &context.base.store_plan.values_query, - context.window_type, - ); - let values = Self::push_compose( - &mut nodes, - values_read, - &context.output_timestamps, - context.query_range_ms, - context.window_size_ms, - context.tumbling_window_ms, + StoreReadRole::Values, + StoreReadStrategy::for_window( + context.window_type, + &context.output_timestamps, + values_lookback_ms, + context.window_size_ms, + context.tumbling_window_ms, + ), ); + let values = Self::push_prepare_buckets(&mut nodes, values_read); let keys = context.base.store_plan.keys_query.as_ref().map(|query| { + let keys_lookback_ms = context.keys_lookback_ms.unwrap_or(context.query_range_ms); + let keys_window_size_ms = context + .keys_window_size_ms + .unwrap_or(context.window_size_ms); + let keys_bucket_step_ms = context + .keys_tumbling_window_ms + .unwrap_or(context.tumbling_window_ms); let read = Self::push_read( &mut nodes, query, - context.keys_window_type.unwrap_or(context.window_type), + StoreReadRole::Keys, + StoreReadStrategy::for_window( + context.keys_window_type.unwrap_or(context.window_type), + &context.output_timestamps, + keys_lookback_ms, + keys_window_size_ms, + keys_bucket_step_ms, + ), ); - Self::push_compose( - &mut nodes, - read, - &context.output_timestamps, - context.keys_lookback_ms.unwrap_or(context.query_range_ms), - context - .keys_window_size_ms - .unwrap_or(context.window_size_ms), - context - .keys_tumbling_window_ms - .unwrap_or(context.tumbling_window_ms), - ) + Self::push_prepare_buckets(&mut nodes, read) }); let resolved = Self::push(&mut nodes, QueryPlanNode::ResolveKeys { values, keys }); + let estimate_spec = RangeEstimateSpec::from(context); let mut root = Self::push( &mut nodes, QueryPlanNode::Estimate { @@ -153,6 +209,7 @@ impl QueryPlan { statistic: context.base.metadata.statistic_to_compute, query_kwargs: context.base.metadata.query_kwargs.clone(), output_labels: context.base.metadata.query_output_labels.clone(), + spec: estimate_spec.clone(), }, ); if options.limit_topk && context.base.metadata.statistic_to_compute == Statistic::Topk { @@ -171,6 +228,7 @@ impl QueryPlan { input: root, k, grouping_labels: context.base.grouping_labels.clone(), + row_label_order: estimate_spec.row_label_order.clone(), }, ); } @@ -285,62 +343,52 @@ impl QueryPlan { fn push_read( nodes: &mut Vec, query: &StoreQueryParams, - window_type: WindowType, + role: StoreReadRole, + strategy: StoreReadStrategy, ) -> NodeId { - let strategy = match window_type { - WindowType::Tumbling => StoreReadStrategy::WindowGrid, - WindowType::Sliding => StoreReadStrategy::SlidingExactCover, - }; Self::push( nodes, QueryPlanNode::StoreRead { query: query.clone(), + role, strategy, }, ) } - fn push_compose( - nodes: &mut Vec, - input: NodeId, - output_timestamps: &[u64], - lookback_ms: u64, - window_size_ms: u64, - bucket_step_ms: u64, - ) -> NodeId { - Self::push( - nodes, - QueryPlanNode::ComposeWindows { - input, - output_timestamps: output_timestamps.to_vec(), - lookback_ms, - window_size_ms, - bucket_step_ms, - }, - ) + fn push_prepare_buckets(nodes: &mut Vec, input: NodeId) -> NodeId { + Self::push(nodes, QueryPlanNode::PrepareBuckets { input }) } pub(crate) fn explain(&self) -> String { let mut lines = Vec::with_capacity(self.nodes.len() + 1); for (index, node) in self.nodes.iter().enumerate() { let line = match node { - QueryPlanNode::StoreRead { query, strategy } => format!( - "n{index} StoreRead({strategy:?}, {}#{}, [{}, {}])", + QueryPlanNode::StoreRead { query, role, strategy } => format!( + "n{index} StoreRead({strategy:?}, role={role:?}, {}#{}, [{}, {}])", query.metric, query.aggregation_id, query.start_timestamp, query.end_timestamp ), - QueryPlanNode::ComposeWindows { input, output_timestamps, lookback_ms, window_size_ms, bucket_step_ms } => format!( - "n{index} ComposeWindows(n{}, outputs={:?}, lookback={lookback_ms}ms, window={window_size_ms}ms, step={bucket_step_ms}ms)", - input.0, output_timestamps - ), + QueryPlanNode::PrepareBuckets { input } => { + format!("n{index} PrepareBuckets(n{})", input.0) + } QueryPlanNode::ResolveKeys { values, keys } => format!( "n{index} ResolveKeys(values=n{}, keys={})", values.0, keys.map(|id| format!("n{}", id.0)).unwrap_or_else(|| "self".to_string()) ), - QueryPlanNode::Estimate { input, statistic, query_kwargs, .. } => { + QueryPlanNode::Estimate { + input, + statistic, + query_kwargs, + spec, + .. + } => { let mut kwargs: Vec<_> = query_kwargs.iter().collect(); kwargs.sort_unstable_by_key(|(key, _)| *key); - format!("n{index} Estimate(n{}, {statistic}, {kwargs:?})", input.0) + format!( + "n{index} Estimate(n{}, {statistic}, {kwargs:?}, outputs={:?})", + input.0, spec.output_timestamps + ) }, QueryPlanNode::LimitTopK { input, k, .. } => { format!("n{index} LimitTopK(n{}, k={k})", input.0) @@ -359,11 +407,38 @@ impl QueryPlan { } } +impl From<&RangeQueryExecutionContext> for RangeEstimateSpec { + fn from(context: &RangeQueryExecutionContext) -> Self { + Self { + output_timestamps: context.output_timestamps.clone(), + query_range_ms: context.query_range_ms, + buckets_per_step: context.buckets_per_step, + lookback_bucket_count: context.lookback_bucket_count, + tumbling_window_ms: context.tumbling_window_ms, + window_type: context.window_type, + window_size_ms: context.window_size_ms, + keys_window_type: context.keys_window_type, + keys_window_size_ms: context.keys_window_size_ms, + keys_lookback_ms: context.keys_lookback_ms, + keys_tumbling_window_ms: context.keys_tumbling_window_ms, + value_aggregation_type: context.base.agg_info.aggregation_type_for_value, + key_aggregation_type: context.base.agg_info.aggregation_type_for_key, + grouping_labels: context.base.grouping_labels.clone(), + aggregated_labels: context.base.aggregated_labels.clone(), + row_label_order: crate::engines::simple_engine::SimpleEngine::topk_row_label_order( + &context.base.metadata, + &context.base.grouping_labels, + &context.base.aggregated_labels, + ), + } + } +} + impl QueryPlanNode { fn kind(&self) -> &'static str { match self { Self::StoreRead { .. } => "StoreRead", - Self::ComposeWindows { .. } => "ComposeWindows", + Self::PrepareBuckets { .. } => "PrepareBuckets", Self::ResolveKeys { .. } => "ResolveKeys", Self::Estimate { .. } => "Estimate", Self::AggregateVector { .. } => "AggregateVector", @@ -375,7 +450,7 @@ impl QueryPlanNode { fn inputs(&self) -> Vec { match self { Self::StoreRead { .. } => Vec::new(), - Self::ComposeWindows { input, .. } + Self::PrepareBuckets { input, .. } | Self::Estimate { input, .. } | Self::AggregateVector { input, .. } | Self::LimitTopK { input, .. } @@ -473,7 +548,9 @@ mod tests { .explain(); assert!(explanation.contains("n4 ResolveKeys(values=n1, keys=n3)")); - assert!(explanation.contains("n2 StoreRead(SlidingExactCover, requests#8")); + assert!(explanation.contains("role=Values")); + assert!(explanation.contains("role=Keys")); + assert!(explanation.contains("n2 StoreRead(SlidingExactCover {")); assert!(explanation.ends_with("root: n5")); } @@ -531,6 +608,29 @@ mod tests { assert!(explanation.contains("outputs=[1000, 2000, 3000]")); } + #[test] + fn sliding_value_read_uses_bucket_lookback() { + let mut context = context(); + context.window_type = WindowType::Sliding; + context.query_range_ms = 1_500; + context.lookback_bucket_count = 1; + + let explanation = QueryPlan::compile_range( + &context, + PlanOptions { + limit_topk: false, + format_output: false, + }, + &[], + ) + .unwrap() + .explain(); + + // Store reads must match the estimator's whole-bucket lookback. + assert!(explanation.contains("lookback_ms: 1000")); + assert!(!explanation.contains("lookback_ms: 1500")); + } + #[test] fn topk_formatting_is_the_plan_root() { let mut context = context(); @@ -584,6 +684,7 @@ mod tests { statistic: Statistic::Sum, query_kwargs: HashMap::new(), output_labels: KeyByLabelNames::empty(), + spec: RangeEstimateSpec::from(&context()), }], root: NodeId(0), }; @@ -622,15 +723,10 @@ mod tests { start_timestamp: 0, end_timestamp: 1, }, + role: StoreReadRole::Values, strategy: StoreReadStrategy::WindowGrid, }, - QueryPlanNode::ComposeWindows { - input: NodeId(0), - output_timestamps: vec![1], - lookback_ms: 1, - window_size_ms: 1, - bucket_step_ms: 1, - }, + QueryPlanNode::PrepareBuckets { input: NodeId(0) }, ], root: NodeId(1), }; diff --git a/asap-query-engine/src/engines/simple_engine/mod.rs b/asap-query-engine/src/engines/simple_engine/mod.rs index 347d5d1a..1084dd0f 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -10,6 +10,7 @@ use crate::data_model::{ }; use crate::engines::query_plan::{ NodeId, PlanOptions, QueryPlan, QueryPlanExecutionError, QueryPlanNode, QueryPlanRuntime, + RangeEstimateSpec, StoreReadRole, StoreReadStrategy, }; use crate::engines::query_result::{InstantVectorElement, QueryResult}; use crate::engines::sliding_window_composition::{ @@ -182,6 +183,7 @@ pub struct RangeQueryExecutionContext { pub keys_tumbling_window_ms: Option, } +#[cfg(feature = "native_query_legacy_test_support")] #[derive(Clone)] struct RangeQueryReads { values: TimestampedBucketsMap, @@ -221,12 +223,15 @@ struct RangePipelineOutput { struct NativePlanRuntime<'a> { engine: &'a SimpleEngine, - context: &'a RangeQueryExecutionContext, - reads: std::cell::RefCell>, } impl NativePlanRuntime<'_> { - fn reads(&self) -> Result { + fn read( + &self, + query: &StoreQueryParams, + role: StoreReadRole, + strategy: &StoreReadStrategy, + ) -> Result { #[cfg(feature = "native_query_legacy_test_support")] if matches!( self.engine.native_range_execution_mode, @@ -236,15 +241,31 @@ impl NativePlanRuntime<'_> { "test-only native store failure".to_string(), )); } - if self.reads.borrow().is_none() { - *self.reads.borrow_mut() = Some(self.engine.read_range_query_inputs(self.context)?); + let data = match strategy { + StoreReadStrategy::WindowGrid => self + .engine + .execute_store_query(query) + .map_err(QueryExecutionError::Native)?, + StoreReadStrategy::SlidingExactCover { + output_timestamps, + lookback_ms, + window_size_ms, + bucket_step_ms, + } => self.engine.execute_sliding_cover_query( + query, + output_timestamps, + *lookback_ms, + *window_size_ms, + *bucket_step_ms, + )?, + }; + if data.is_empty() && role == StoreReadRole::Values { + return Err(QueryExecutionError::NoLocalData(format!( + "No data found for metric: {}", + query.metric + ))); } - Ok(self - .reads - .borrow() - .as_ref() - .expect("reads initialized") - .clone()) + Ok(data) } } @@ -259,25 +280,19 @@ impl QueryPlanRuntime for NativePlanRuntime<'_> { inputs: &[Self::Output], ) -> Result { match node { - QueryPlanNode::StoreRead { query, strategy: _ } => { - let reads = self.reads()?; - if query.aggregation_id == self.context.base.store_plan.values_query.aggregation_id - { - Ok(NativePlanOutput::Read(reads.values)) - } else { - reads.keys.map(NativePlanOutput::Read).ok_or_else(|| { - QueryExecutionError::Native( - "Query plan requested missing key read".to_string(), - ) - }) - } - } - QueryPlanNode::ComposeWindows { .. } => match inputs { + QueryPlanNode::StoreRead { + query, + role, + strategy, + } => self + .read(query, *role, strategy) + .map(NativePlanOutput::Read), + QueryPlanNode::PrepareBuckets { .. } => match inputs { [NativePlanOutput::Read(data)] => Ok(NativePlanOutput::Composed( self.engine.compose_range_read(data), )), _ => Err(QueryExecutionError::Native( - "ComposeWindows expected store data".into(), + "PrepareBuckets expected store data".into(), )), }, QueryPlanNode::ResolveKeys { keys, .. } => match (inputs, keys) { @@ -298,10 +313,16 @@ impl QueryPlanRuntime for NativePlanRuntime<'_> { "ResolveKeys received incompatible inputs".into(), )), }, - QueryPlanNode::Estimate { output_labels, .. } => match inputs { + QueryPlanNode::Estimate { + statistic, + query_kwargs, + output_labels, + spec, + .. + } => match inputs { [NativePlanOutput::Resolved(reads)] => self .engine - .estimate_range_query(self.context, reads.clone()) + .estimate_range_query(spec, *statistic, query_kwargs, reads.clone()) .map(|values| NativePlanOutput::Results { labels: output_labels.clone(), values, @@ -329,20 +350,14 @@ impl QueryPlanRuntime for NativePlanRuntime<'_> { )), }, QueryPlanNode::LimitTopK { - k, grouping_labels, .. + k, + grouping_labels, + row_label_order, + .. } => match inputs { [NativePlanOutput::Results { labels, values }] => self .engine - .limit_range_topk( - values, - k, - &SimpleEngine::topk_row_label_order( - &self.context.base.metadata, - &self.context.base.grouping_labels, - &self.context.base.aggregated_labels, - ), - grouping_labels, - ) + .limit_range_topk(values, k, row_label_order, grouping_labels) .map_err(QueryExecutionError::Native) .map(|values| NativePlanOutput::Results { labels: labels.clone(), @@ -1214,7 +1229,7 @@ impl SimpleEngine { /// bare-row topk has no named-output-label concept at all), so this /// falls back to the aggregation config's own `grouping_labels ++ /// aggregated_labels` order in that case. - fn topk_row_label_order( + pub(crate) fn topk_row_label_order( metadata: &QueryMetadata, grouping_labels: &KeyByLabelNames, aggregated_labels: &KeyByLabelNames, @@ -2489,7 +2504,10 @@ impl SimpleEngine { enable_topk_formatting: bool, query_time_aggregations: &[asap_types::query_config::QueryTimeAggregation], ) -> Result { - Self::reject_off_grid_sliding_counter_query(context)?; + Self::reject_off_grid_sliding_counter_query( + &RangeEstimateSpec::from(context), + context.base.metadata.statistic_to_compute, + )?; #[cfg(feature = "native_query_legacy_test_support")] if matches!( self.native_range_execution_mode, @@ -2539,11 +2557,7 @@ impl SimpleEngine { ) .map_err(QueryExecutionError::Native)?; debug!(plan = %plan.explain(), "Compiled native query plan"); - let runtime = NativePlanRuntime { - engine: self, - context, - reads: std::cell::RefCell::new(None), - }; + let runtime = NativePlanRuntime { engine: self }; match plan.execute(&runtime).map_err(|error| match error { QueryPlanExecutionError::InvalidPlan(reason) => QueryExecutionError::Native(reason), QueryPlanExecutionError::Node { source, .. } => source, @@ -2565,8 +2579,11 @@ impl SimpleEngine { enable_topk_formatting: bool, ) -> Result, QueryExecutionError> { let reads = self.read_range_query_inputs(context)?; + let estimate_spec = RangeEstimateSpec::from(context); let mut results = self.estimate_range_query( - context, + &estimate_spec, + context.base.metadata.statistic_to_compute, + &context.base.metadata.query_kwargs, ResolvedRangeReads { values: self.compose_range_read(&reads.values), keys: reads @@ -2601,6 +2618,7 @@ impl SimpleEngine { )) } + #[cfg(feature = "native_query_legacy_test_support")] fn read_range_query_inputs( &self, context: &RangeQueryExecutionContext, @@ -2662,26 +2680,24 @@ impl SimpleEngine { } fn reject_off_grid_sliding_counter_query( - context: &RangeQueryExecutionContext, + spec: &RangeEstimateSpec, + statistic: Statistic, ) -> Result<(), QueryExecutionError> { - if context.window_type != WindowType::Sliding - || context.tumbling_window_ms == 0 - || !matches!( - context.base.metadata.statistic_to_compute, - Statistic::Increase | Statistic::Rate - ) + if spec.window_type != WindowType::Sliding + || spec.tumbling_window_ms == 0 + || !matches!(statistic, Statistic::Increase | Statistic::Rate) { return Ok(()); } - if let Some(×tamp) = context + if let Some(×tamp) = spec .output_timestamps .iter() - .find(|&×tamp| !timestamp.is_multiple_of(context.tumbling_window_ms)) + .find(|&×tamp| !timestamp.is_multiple_of(spec.tumbling_window_ms)) { return Err(QueryExecutionError::NoLocalData(format!( "Exact Prometheus counter bounds are unavailable for off-grid Sliding \ timestamp {} (grid interval {}ms)", - timestamp, context.tumbling_window_ms + timestamp, spec.tumbling_window_ms ))); } Ok(()) @@ -2689,7 +2705,9 @@ impl SimpleEngine { fn estimate_range_query( &self, - context: &RangeQueryExecutionContext, + spec: &RangeEstimateSpec, + statistic: Statistic, + query_kwargs: &HashMap, reads: ResolvedRangeReads, ) -> Result, QueryExecutionError> { use crate::engines::query_result::RangeVectorElement; @@ -2699,23 +2717,24 @@ impl SimpleEngine { values: ComposedRangeRead { groups: all_data }, keys: keys_raw_data, } = reads; - let lookback_ms = (context.lookback_bucket_count as u64) * context.tumbling_window_ms; + Self::reject_off_grid_sliding_counter_query(spec, statistic)?; + let lookback_ms = (spec.lookback_bucket_count as u64) * spec.tumbling_window_ms; let mut results: HashMap = HashMap::new(); // Determine accumulator type for merger selection - let accumulator_type = &context.base.agg_info.aggregation_type_for_value; - let key_accumulator_type = context.base.agg_info.aggregation_type_for_key; + let accumulator_type = &spec.value_aggregation_type; + let key_accumulator_type = spec.key_aggregation_type; // Calculate step parameters - let buckets_per_step = context.buckets_per_step; - let lookback_bucket_count = context.lookback_bucket_count; - let tumbling_window_ms = context.tumbling_window_ms; - let window_type = context.window_type; - let keys_lookback_ms = context.keys_lookback_ms; - let keys_tumbling_window_ms = context.keys_tumbling_window_ms; - let keys_window_type = context.keys_window_type; - let keys_window_size_ms = context.keys_window_size_ms; + let buckets_per_step = spec.buckets_per_step; + let lookback_bucket_count = spec.lookback_bucket_count; + let tumbling_window_ms = spec.tumbling_window_ms; + let window_type = spec.window_type; + let keys_lookback_ms = spec.keys_lookback_ms; + let keys_tumbling_window_ms = spec.keys_tumbling_window_ms; + let keys_window_type = spec.keys_window_type; + let keys_window_size_ms = spec.keys_window_size_ms; // Named distinctly from `WindowType` (Sliding/Tumbling, picks how a // step's window is composed from `bucket_map` below) -- this describes step-to-step overlap in @@ -2729,9 +2748,9 @@ impl SimpleEngine { debug!( "Range query params: {} output timestamp(s) [{}..{}], tumbling_window_ms={}, \ buckets_per_step (slide)={}, lookback_bucket_count (size)={}, mode={}", - context.output_timestamps.len(), - context.output_timestamps.first().copied().unwrap_or(0), - context.output_timestamps.last().copied().unwrap_or(0), + spec.output_timestamps.len(), + spec.output_timestamps.first().copied().unwrap_or(0), + spec.output_timestamps.last().copied().unwrap_or(0), tumbling_window_ms, buckets_per_step, lookback_bucket_count, @@ -2794,8 +2813,8 @@ impl SimpleEngine { // group). See #582 review for collect_results_separate_keys parity. let groups: Vec<(GroupBucketMap, KeysSource)> = match &keys_raw_data { Some(keys_map) => { - // keys_raw_data is Some, so context.keys_lookback_ms / - // context.keys_tumbling_window_ms are guaranteed Some too + // keys_raw_data is Some, so its plan spec's key-window + // fields are guaranteed Some too. // (both derived from the same keys_query.is_some() check in // finish_range_context) -- resolved once here instead of // re-unwrapped per group per step. @@ -2881,11 +2900,7 @@ impl SimpleEngine { // this is actually a topk query with limiting requested -- gates // both the per-step sort/truncate below and nothing else, so a // non-topk query pays zero cost for this. - let row_label_order = Self::topk_row_label_order( - &context.base.metadata, - &context.base.grouping_labels, - &context.base.aggregated_labels, - ); + let row_label_order = spec.row_label_order.clone(); // Step-major: for each output timestamp, visit every group, not the // other way around. Required for topk correctness -- ranking a @@ -2893,13 +2908,13 @@ impl SimpleEngine { // timestamp before truncating, which a group-major loop can't do // (#581). One loop shape for topk and non-topk alike, rather than // maintaining two. - for ¤t_time in &context.output_timestamps { + for ¤t_time in &spec.output_timestamps { let current_time_i64 = i64::try_from(current_time).map_err(|_| { QueryExecutionError::Native( "Output timestamp exceeds signed timestamp range".to_string(), ) })?; - let query_range_ms = i64::try_from(context.query_range_ms).map_err(|_| { + let query_range_ms = i64::try_from(spec.query_range_ms).map_err(|_| { QueryExecutionError::Native( "Query range exceeds signed timestamp range".to_string(), ) @@ -3049,7 +3064,7 @@ impl SimpleEngine { window_end, tumbling_window_ms, window_type, - context.window_size_ms, + spec.window_size_ms, ); trace!( @@ -3058,13 +3073,13 @@ impl SimpleEngine { window_start, window_end, grid_step_ms = tumbling_window_ms, - stored_window_size_ms = context.window_size_ms, + stored_window_size_ms = spec.window_size_ms, selected_bucket_count = window_buckets.len(), "Composed query output window from stored buckets" ); if window_type == WindowType::Sliding { - let expected = (lookback_ms / context.window_size_ms) as usize; + let expected = (lookback_ms / spec.window_size_ms) as usize; if window_buckets.len() < expected { debug!( "Skipping incomplete Sliding value cover at t={}", @@ -3108,11 +3123,11 @@ impl SimpleEngine { Some(merged.as_ref()), keys_precompute.as_deref(), &fallback_key, - &context.base.grouping_labels, - &context.base.aggregated_labels, + &spec.grouping_labels, + &spec.aggregated_labels, &row_label_order, - &context.base.metadata.statistic_to_compute, - &context.base.metadata.query_kwargs, + &statistic, + query_kwargs, Some(&query_bounds), ) .into_iter() From b53ad94446f6c9385369f15cbde79e0231220d98 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Wed, 7 Oct 2026 21:15:22 -0400 Subject: [PATCH 2/2] refactor(query-engine): make range plan specs cohesive --- asap-query-engine/src/engines/query_plan.rs | 317 ++++++++++++----- .../src/engines/simple_engine/mod.rs | 323 ++++++++---------- .../src/tests/native_range_query_tests.rs | 51 +++ 3 files changed, 424 insertions(+), 267 deletions(-) diff --git a/asap-query-engine/src/engines/query_plan.rs b/asap-query-engine/src/engines/query_plan.rs index 1257f1c1..c826984b 100644 --- a/asap-query-engine/src/engines/query_plan.rs +++ b/asap-query-engine/src/engines/query_plan.rs @@ -6,6 +6,7 @@ use asap_types::enums::WindowType; use asap_types::query_config::QueryTimeAggregation; use promql_utilities::data_model::KeyByLabelNames; use promql_utilities::query_logics::enums::{AggregationType, Statistic}; +use std::sync::Arc; use tracing::debug; #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -14,30 +15,14 @@ pub(crate) struct NodeId(usize); #[derive(Debug, Clone, PartialEq, Eq)] pub(crate) enum StoreReadStrategy { WindowGrid, - SlidingExactCover { - output_timestamps: Vec, - lookback_ms: u64, - window_size_ms: u64, - bucket_step_ms: u64, - }, + SlidingExactCover(WindowCompositionSpec), } impl StoreReadStrategy { - fn for_window( - window_type: WindowType, - output_timestamps: &[u64], - lookback_ms: u64, - window_size_ms: u64, - bucket_step_ms: u64, - ) -> Self { - match window_type { + fn for_window(window: &WindowCompositionSpec) -> Self { + match window.window_type { WindowType::Tumbling => Self::WindowGrid, - WindowType::Sliding => Self::SlidingExactCover { - output_timestamps: output_timestamps.to_vec(), - lookback_ms, - window_size_ms, - bucket_step_ms, - }, + WindowType::Sliding => Self::SlidingExactCover(window.clone()), } } } @@ -48,21 +33,35 @@ pub(crate) enum StoreReadRole { Keys, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct WindowCompositionSpec { + pub output_timestamps: Arc<[u64]>, + pub lookback_ms: u64, + pub window_type: WindowType, + pub window_size_ms: u64, + pub bucket_step_ms: u64, +} + +#[derive(Debug, Clone)] +pub(crate) struct ValueEstimateInput { + pub aggregation_type: AggregationType, + pub window: WindowCompositionSpec, +} + +#[derive(Debug, Clone)] +pub(crate) enum KeyInputSpec { + FromValues, + Separate { + aggregation_type: AggregationType, + window: WindowCompositionSpec, + }, +} + #[derive(Debug, Clone)] pub(crate) struct RangeEstimateSpec { - pub output_timestamps: Vec, pub query_range_ms: u64, - pub buckets_per_step: usize, - pub lookback_bucket_count: usize, - pub tumbling_window_ms: u64, - pub window_type: WindowType, - pub window_size_ms: u64, - pub keys_window_type: Option, - pub keys_window_size_ms: Option, - pub keys_lookback_ms: Option, - pub keys_tumbling_window_ms: Option, - pub value_aggregation_type: AggregationType, - pub key_aggregation_type: AggregationType, + pub values: ValueEstimateInput, + pub keys: KeyInputSpec, pub grouping_labels: KeyByLabelNames, pub aggregated_labels: KeyByLabelNames, pub row_label_order: KeyByLabelNames, @@ -163,45 +162,33 @@ impl QueryPlan { query_time_aggregations: &[QueryTimeAggregation], ) -> Result { let mut nodes = Vec::new(); - let values_lookback_ms = - (context.lookback_bucket_count as u64) * context.tumbling_window_ms; + let estimate_spec = RangeEstimateSpec::compile(context)?; let values_read = Self::push_read( &mut nodes, &context.base.store_plan.values_query, StoreReadRole::Values, - StoreReadStrategy::for_window( - context.window_type, - &context.output_timestamps, - values_lookback_ms, - context.window_size_ms, - context.tumbling_window_ms, - ), + StoreReadStrategy::for_window(&estimate_spec.values.window), ); let values = Self::push_prepare_buckets(&mut nodes, values_read); - let keys = context.base.store_plan.keys_query.as_ref().map(|query| { - let keys_lookback_ms = context.keys_lookback_ms.unwrap_or(context.query_range_ms); - let keys_window_size_ms = context - .keys_window_size_ms - .unwrap_or(context.window_size_ms); - let keys_bucket_step_ms = context - .keys_tumbling_window_ms - .unwrap_or(context.tumbling_window_ms); - let read = Self::push_read( - &mut nodes, - query, - StoreReadRole::Keys, - StoreReadStrategy::for_window( - context.keys_window_type.unwrap_or(context.window_type), - &context.output_timestamps, - keys_lookback_ms, - keys_window_size_ms, - keys_bucket_step_ms, - ), - ); - Self::push_prepare_buckets(&mut nodes, read) - }); + let keys = match (&context.base.store_plan.keys_query, &estimate_spec.keys) { + (None, KeyInputSpec::FromValues) => None, + (Some(query), KeyInputSpec::Separate { window, .. }) => { + let read = Self::push_read( + &mut nodes, + query, + StoreReadRole::Keys, + StoreReadStrategy::for_window(window), + ); + Some(Self::push_prepare_buckets(&mut nodes, read)) + } + (Some(_), KeyInputSpec::FromValues) => { + return Err("Query plan has a keys read without a keys specification".to_string()); + } + (None, KeyInputSpec::Separate { .. }) => { + return Err("Query plan has a keys specification without a keys read".to_string()); + } + }; let resolved = Self::push(&mut nodes, QueryPlanNode::ResolveKeys { values, keys }); - let estimate_spec = RangeEstimateSpec::from(context); let mut root = Self::push( &mut nodes, QueryPlanNode::Estimate { @@ -365,7 +352,8 @@ impl QueryPlan { for (index, node) in self.nodes.iter().enumerate() { let line = match node { QueryPlanNode::StoreRead { query, role, strategy } => format!( - "n{index} StoreRead({strategy:?}, role={role:?}, {}#{}, [{}, {}])", + "n{index} StoreRead({}, role={role:?}, {}#{}, [{}, {}])", + Self::describe_read_strategy(strategy), query.metric, query.aggregation_id, query.start_timestamp, query.end_timestamp ), QueryPlanNode::PrepareBuckets { input } => { @@ -386,8 +374,9 @@ impl QueryPlan { let mut kwargs: Vec<_> = query_kwargs.iter().collect(); kwargs.sort_unstable_by_key(|(key, _)| *key); format!( - "n{index} Estimate(n{}, {statistic}, {kwargs:?}, outputs={:?})", - input.0, spec.output_timestamps + "n{index} Estimate(n{}, {statistic}, {kwargs:?}, {})", + input.0, + Self::describe_timestamps(&spec.values.window.output_timestamps) ) }, QueryPlanNode::LimitTopK { input, k, .. } => { @@ -405,24 +394,68 @@ impl QueryPlan { lines.push(format!("root: n{}", self.root.0)); lines.join("\n") } + + fn describe_read_strategy(strategy: &StoreReadStrategy) -> String { + match strategy { + StoreReadStrategy::WindowGrid => "WindowGrid".to_string(), + StoreReadStrategy::SlidingExactCover(window) => format!( + "SlidingExactCover(lookback={}ms, window={}ms, step={}ms, {})", + window.lookback_ms, + window.window_size_ms, + window.bucket_step_ms, + Self::describe_timestamps(&window.output_timestamps), + ), + } + } + + fn describe_timestamps(timestamps: &[u64]) -> String { + format!( + "outputs=count:{}, first:{}, last:{}", + timestamps.len(), + timestamps.first().copied().unwrap_or(0), + timestamps.last().copied().unwrap_or(0), + ) + } } -impl From<&RangeQueryExecutionContext> for RangeEstimateSpec { - fn from(context: &RangeQueryExecutionContext) -> Self { - Self { - output_timestamps: context.output_timestamps.clone(), +impl RangeEstimateSpec { + pub(crate) fn compile(context: &RangeQueryExecutionContext) -> Result { + let output_timestamps: Arc<[u64]> = Arc::from(context.output_timestamps.clone()); + let values = ValueEstimateInput { + aggregation_type: context.base.agg_info.aggregation_type_for_value, + window: WindowCompositionSpec { + output_timestamps: Arc::clone(&output_timestamps), + lookback_ms: (context.lookback_bucket_count as u64) * context.tumbling_window_ms, + window_type: context.window_type, + window_size_ms: context.window_size_ms, + bucket_step_ms: context.tumbling_window_ms, + }, + }; + let keys = match &context.base.store_plan.keys_query { + None => KeyInputSpec::FromValues, + Some(_) => KeyInputSpec::Separate { + aggregation_type: context.base.agg_info.aggregation_type_for_key, + window: WindowCompositionSpec { + output_timestamps, + lookback_ms: context + .keys_lookback_ms + .ok_or_else(|| "Separate keys query is missing its lookback".to_string())?, + window_type: context.keys_window_type.ok_or_else(|| { + "Separate keys query is missing its window type".to_string() + })?, + window_size_ms: context.keys_window_size_ms.ok_or_else(|| { + "Separate keys query is missing its window size".to_string() + })?, + bucket_step_ms: context.keys_tumbling_window_ms.ok_or_else(|| { + "Separate keys query is missing its bucket step".to_string() + })?, + }, + }, + }; + Ok(Self { query_range_ms: context.query_range_ms, - buckets_per_step: context.buckets_per_step, - lookback_bucket_count: context.lookback_bucket_count, - tumbling_window_ms: context.tumbling_window_ms, - window_type: context.window_type, - window_size_ms: context.window_size_ms, - keys_window_type: context.keys_window_type, - keys_window_size_ms: context.keys_window_size_ms, - keys_lookback_ms: context.keys_lookback_ms, - keys_tumbling_window_ms: context.keys_tumbling_window_ms, - value_aggregation_type: context.base.agg_info.aggregation_type_for_value, - key_aggregation_type: context.base.agg_info.aggregation_type_for_key, + values, + keys, grouping_labels: context.base.grouping_labels.clone(), aggregated_labels: context.base.aggregated_labels.clone(), row_label_order: crate::engines::simple_engine::SimpleEngine::topk_row_label_order( @@ -430,7 +463,7 @@ impl From<&RangeQueryExecutionContext> for RangeEstimateSpec { &context.base.grouping_labels, &context.base.aggregated_labels, ), - } + }) } } @@ -535,6 +568,9 @@ mod tests { end_timestamp: 1_000, }); context.keys_window_type = Some(WindowType::Sliding); + context.keys_lookback_ms = Some(2_000); + context.keys_window_size_ms = Some(1_000); + context.keys_tumbling_window_ms = Some(500); let explanation = QueryPlan::compile_range( &context, @@ -550,7 +586,8 @@ mod tests { assert!(explanation.contains("n4 ResolveKeys(values=n1, keys=n3)")); assert!(explanation.contains("role=Values")); assert!(explanation.contains("role=Keys")); - assert!(explanation.contains("n2 StoreRead(SlidingExactCover {")); + assert!(explanation + .contains("n2 StoreRead(SlidingExactCover(lookback=2000ms, window=1000ms, step=500ms")); assert!(explanation.ends_with("root: n5")); } @@ -605,7 +642,7 @@ mod tests { .unwrap() .explain(); - assert!(explanation.contains("outputs=[1000, 2000, 3000]")); + assert!(explanation.contains("outputs=count:3, first:1000, last:3000")); } #[test] @@ -615,7 +652,7 @@ mod tests { context.query_range_ms = 1_500; context.lookback_bucket_count = 1; - let explanation = QueryPlan::compile_range( + let plan = QueryPlan::compile_range( &context, PlanOptions { limit_topk: false, @@ -623,12 +660,108 @@ mod tests { }, &[], ) - .unwrap() - .explain(); + .unwrap(); + + match &plan.nodes[0] { + QueryPlanNode::StoreRead { + strategy: StoreReadStrategy::SlidingExactCover(window), + .. + } => assert_eq!(window.lookback_ms, 1_000), + _ => panic!("value read must use a sliding exact cover"), + } + + match &plan.nodes[3] { + QueryPlanNode::Estimate { spec, .. } => { + assert_eq!(spec.values.window.lookback_ms, 1_000); + } + _ => panic!("fourth node must estimate the resolved value read"), + } + } - // Store reads must match the estimator's whole-bucket lookback. - assert!(explanation.contains("lookback_ms: 1000")); - assert!(!explanation.contains("lookback_ms: 1500")); + #[test] + fn rejects_separate_keys_with_incomplete_window_settings() { + let mut context = context(); + context.base.store_plan.keys_query = Some(context.base.store_plan.values_query.clone()); + context.keys_window_type = Some(WindowType::Sliding); + + let error = QueryPlan::compile_range( + &context, + PlanOptions { + limit_topk: false, + format_output: false, + }, + &[], + ) + .expect_err("incomplete separate key settings must fail loudly"); + + assert_eq!(error, "Separate keys query is missing its lookback"); + } + + struct PlanOwnedReadRuntime(RefCell>); + + impl QueryPlanRuntime for PlanOwnedReadRuntime { + type Output = usize; + type Error = std::convert::Infallible; + + fn execute_node( + &self, + _id: NodeId, + node: &QueryPlanNode, + inputs: &[Self::Output], + ) -> Result { + if let QueryPlanNode::StoreRead { role, strategy, .. } = node { + self.0.borrow_mut().push((*role, strategy.clone())); + } + Ok(1 + inputs.iter().sum::()) + } + } + + #[test] + fn compiled_plan_executes_distinct_value_and_key_read_specs_without_context() { + let mut context = context(); + context.window_type = WindowType::Sliding; + context.lookback_bucket_count = 2; + context.base.store_plan.keys_query = Some(StoreQueryParams { + metric: "requests".into(), + aggregation_id: 8, + start_timestamp: 0, + end_timestamp: 1_000, + }); + context.keys_window_type = Some(WindowType::Tumbling); + context.keys_lookback_ms = Some(3_000); + context.keys_window_size_ms = Some(1_000); + context.keys_tumbling_window_ms = Some(1_000); + + let plan = QueryPlan::compile_range( + &context, + PlanOptions { + limit_topk: false, + format_output: false, + }, + &[], + ) + .unwrap(); + let runtime = PlanOwnedReadRuntime(RefCell::new(Vec::new())); + + plan.execute(&runtime) + .expect("compiled plan must execute with a context-free runtime"); + + assert_eq!( + runtime.0.into_inner(), + vec![ + ( + StoreReadRole::Values, + StoreReadStrategy::SlidingExactCover(WindowCompositionSpec { + output_timestamps: Arc::from([1_000]), + lookback_ms: 2_000, + window_type: WindowType::Sliding, + window_size_ms: 1_000, + bucket_step_ms: 1_000, + }), + ), + (StoreReadRole::Keys, StoreReadStrategy::WindowGrid), + ] + ); } #[test] @@ -684,7 +817,7 @@ mod tests { statistic: Statistic::Sum, query_kwargs: HashMap::new(), output_labels: KeyByLabelNames::empty(), - spec: RangeEstimateSpec::from(&context()), + spec: RangeEstimateSpec::compile(&context()).unwrap(), }], root: NodeId(0), }; diff --git a/asap-query-engine/src/engines/simple_engine/mod.rs b/asap-query-engine/src/engines/simple_engine/mod.rs index 1084dd0f..e131695d 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -9,8 +9,8 @@ use crate::data_model::{ StreamingConfig, }; use crate::engines::query_plan::{ - NodeId, PlanOptions, QueryPlan, QueryPlanExecutionError, QueryPlanNode, QueryPlanRuntime, - RangeEstimateSpec, StoreReadRole, StoreReadStrategy, + KeyInputSpec, NodeId, PlanOptions, QueryPlan, QueryPlanExecutionError, QueryPlanNode, + QueryPlanRuntime, RangeEstimateSpec, StoreReadRole, StoreReadStrategy, }; use crate::engines::query_result::{InstantVectorElement, QueryResult}; use crate::engines::sliding_window_composition::{ @@ -246,18 +246,15 @@ impl NativePlanRuntime<'_> { .engine .execute_store_query(query) .map_err(QueryExecutionError::Native)?, - StoreReadStrategy::SlidingExactCover { - output_timestamps, - lookback_ms, - window_size_ms, - bucket_step_ms, - } => self.engine.execute_sliding_cover_query( - query, - output_timestamps, - *lookback_ms, - *window_size_ms, - *bucket_step_ms, - )?, + StoreReadStrategy::SlidingExactCover(window) => { + self.engine.execute_sliding_cover_query( + query, + &window.output_timestamps, + window.lookback_ms, + window.window_size_ms, + window.bucket_step_ms, + )? + } }; if data.is_empty() && role == StoreReadRole::Values { return Err(QueryExecutionError::NoLocalData(format!( @@ -265,6 +262,13 @@ impl NativePlanRuntime<'_> { query.metric ))); } + let bucket_count = data.values().map(Vec::len).sum::(); + debug!( + role = ?role, + group_count = data.len(), + bucket_count, + "Range query store read completed" + ); Ok(data) } } @@ -2505,7 +2509,9 @@ impl SimpleEngine { query_time_aggregations: &[asap_types::query_config::QueryTimeAggregation], ) -> Result { Self::reject_off_grid_sliding_counter_query( - &RangeEstimateSpec::from(context), + context.window_type, + context.tumbling_window_ms, + &context.output_timestamps, context.base.metadata.statistic_to_compute, )?; #[cfg(feature = "native_query_legacy_test_support")] @@ -2579,7 +2585,8 @@ impl SimpleEngine { enable_topk_formatting: bool, ) -> Result, QueryExecutionError> { let reads = self.read_range_query_inputs(context)?; - let estimate_spec = RangeEstimateSpec::from(context); + let estimate_spec = + RangeEstimateSpec::compile(context).map_err(QueryExecutionError::Native)?; let mut results = self.estimate_range_query( &estimate_spec, context.base.metadata.statistic_to_compute, @@ -2680,24 +2687,25 @@ impl SimpleEngine { } fn reject_off_grid_sliding_counter_query( - spec: &RangeEstimateSpec, + window_type: WindowType, + bucket_step_ms: u64, + output_timestamps: &[u64], statistic: Statistic, ) -> Result<(), QueryExecutionError> { - if spec.window_type != WindowType::Sliding - || spec.tumbling_window_ms == 0 + if window_type != WindowType::Sliding + || bucket_step_ms == 0 || !matches!(statistic, Statistic::Increase | Statistic::Rate) { return Ok(()); } - if let Some(×tamp) = spec - .output_timestamps + if let Some(×tamp) = output_timestamps .iter() - .find(|&×tamp| !timestamp.is_multiple_of(spec.tumbling_window_ms)) + .find(|&×tamp| !timestamp.is_multiple_of(bucket_step_ms)) { return Err(QueryExecutionError::NoLocalData(format!( "Exact Prometheus counter bounds are unavailable for off-grid Sliding \ timestamp {} (grid interval {}ms)", - timestamp, spec.tumbling_window_ms + timestamp, bucket_step_ms ))); } Ok(()) @@ -2717,44 +2725,22 @@ impl SimpleEngine { values: ComposedRangeRead { groups: all_data }, keys: keys_raw_data, } = reads; - Self::reject_off_grid_sliding_counter_query(spec, statistic)?; - let lookback_ms = (spec.lookback_bucket_count as u64) * spec.tumbling_window_ms; + let value_window = &spec.values.window; + let lookback_ms = value_window.lookback_ms; let mut results: HashMap = HashMap::new(); // Determine accumulator type for merger selection - let accumulator_type = &spec.value_aggregation_type; - let key_accumulator_type = spec.key_aggregation_type; + let accumulator_type = &spec.values.aggregation_type; // Calculate step parameters - let buckets_per_step = spec.buckets_per_step; - let lookback_bucket_count = spec.lookback_bucket_count; - let tumbling_window_ms = spec.tumbling_window_ms; - let window_type = spec.window_type; - let keys_lookback_ms = spec.keys_lookback_ms; - let keys_tumbling_window_ms = spec.keys_tumbling_window_ms; - let keys_window_type = spec.keys_window_type; - let keys_window_size_ms = spec.keys_window_size_ms; - - // Named distinctly from `WindowType` (Sliding/Tumbling, picks how a - // step's window is composed from `bucket_map` below) -- this describes step-to-step overlap in - // the OUTPUT iteration, an unrelated concept that happens to reuse - // the words "sliding"/"hopping". See #581. - let step_overlap_mode = if buckets_per_step <= lookback_bucket_count { - "sliding (slide <= size)" - } else { - "hopping (slide > size)" - }; debug!( - "Range query params: {} output timestamp(s) [{}..{}], tumbling_window_ms={}, \ - buckets_per_step (slide)={}, lookback_bucket_count (size)={}, mode={}", - spec.output_timestamps.len(), - spec.output_timestamps.first().copied().unwrap_or(0), - spec.output_timestamps.last().copied().unwrap_or(0), - tumbling_window_ms, - buckets_per_step, - lookback_bucket_count, - step_overlap_mode + output_count = value_window.output_timestamps.len(), + first_output = value_window.output_timestamps.first().copied().unwrap_or(0), + last_output = value_window.output_timestamps.last().copied().unwrap_or(0), + lookback_ms, + bucket_step_ms = value_window.bucket_step_ms, + "Range query estimate parameters" ); // Whether the value accumulator's own get_keys() is even consulted @@ -2776,16 +2762,6 @@ impl SimpleEngine { // top-k), evaluated AFTER merging the window's value buckets // (a top-k heap's keys can depend on that window's data), // falling back to the store-level group key otherwise. - // PerStep bundles everything a dual-population group's per-step - // resolution needs (bucket_map, lookback_ms, tumbling_window_ms) in - // one place, built once at group-construction time — rather than - // three separate Option fields at function scope that only - // happened to be Some together by convention, each re-unwrapped via - // .expect() on every iteration of the per-step loop. Making the - // invalid state (PerStep present but one companion value missing) - // unrepresentable is the same reasoning that motivated this enum - // over two raw Option fields in the first place — just applied all - // the way through instead of partway. // Named alias purely to keep declarations under // clippy::type_complexity -- used by both KeysSource::PerStep's own // bucket_map field below and the step-major `groups` binding @@ -2793,14 +2769,12 @@ impl SimpleEngine { // raw type at the PerStep site instead of using this alias). type GroupBucketMap = BucketMap; - enum KeysSource { + enum KeysSource<'a> { Fixed(Option), PerStep { bucket_map: GroupBucketMap, - lookback_ms: u64, - tumbling_window_ms: u64, - stored_window_size_ms: u64, - window_type: WindowType, + aggregation_type: AggregationType, + window: &'a crate::engines::query_plan::WindowCompositionSpec, }, } @@ -2811,61 +2785,60 @@ impl SimpleEngine { // of failing the whole range query (#583; previously // `.ok_or_else(...)?` here hard-failed everything for one missing // group). See #582 review for collect_results_separate_keys parity. - let groups: Vec<(GroupBucketMap, KeysSource)> = match &keys_raw_data { - Some(keys_map) => { - // keys_raw_data is Some, so its plan spec's key-window - // fields are guaranteed Some too. - // (both derived from the same keys_query.is_some() check in - // finish_range_context) -- resolved once here instead of - // re-unwrapped per group per step. - let keys_lookback_ms = - keys_lookback_ms.expect("keys_raw_data implies keys_lookback_ms is Some"); - let keys_tumbling_window_ms = keys_tumbling_window_ms - .expect("keys_raw_data implies keys_tumbling_window_ms is Some"); - let keys_window_type = - keys_window_type.expect("keys_raw_data implies keys_window_type is Some"); - let keys_window_size_ms = - keys_window_size_ms.expect("keys_raw_data implies keys_window_size_ms is Some"); - keys_map - .groups - .iter() - .filter_map( - |(group_key, key_bucket_map)| match all_data.get(group_key) { - Some(value_bucket_map) => Some(( - value_bucket_map.clone(), - KeysSource::PerStep { - bucket_map: key_bucket_map.clone(), - lookback_ms: keys_lookback_ms, - tumbling_window_ms: keys_tumbling_window_ms, - stored_window_size_ms: keys_window_size_ms, - window_type: keys_window_type, - }, - )), - None => { - warn!( - "Range query: group {:?} has keys data but no value data \ + let groups: Vec<(GroupBucketMap, KeysSource)> = match (&spec.keys, &keys_raw_data) { + ( + KeyInputSpec::Separate { + aggregation_type, + window, + }, + Some(keys_map), + ) => keys_map + .groups + .iter() + .filter_map( + |(group_key, key_bucket_map)| match all_data.get(group_key) { + Some(value_bucket_map) => Some(( + value_bucket_map.clone(), + KeysSource::PerStep { + bucket_map: key_bucket_map.clone(), + aggregation_type: *aggregation_type, + window, + }, + )), + None => { + warn!( + "Range query: group {:?} has keys data but no value data \ anywhere in the queried range — skipping this group instead \ of failing the whole query (#583)", - group_key - ); - None - } - }, - ) - .collect() - } + group_key + ); + None + } + }, + ) + .collect(), // #584/#587: keep every group, including group_key=None — that's // exactly where a self-keyed single-population accumulator // (e.g. top-k) is typically stored. An empty fallback list here // is fine; the per-step loop below tries the value // accumulator's own get_keys() first and only falls back to // this list. - None => all_data + (KeyInputSpec::FromValues, None) => all_data .iter() .map(|(group_key, bucket_map)| { (bucket_map.clone(), KeysSource::Fixed(group_key.clone())) }) .collect(), + (KeyInputSpec::Separate { .. }, None) => { + return Err(QueryExecutionError::Native( + "Invalid plan: separate key specification has no key reads".into(), + )); + } + (KeyInputSpec::FromValues, Some(_)) => { + return Err(QueryExecutionError::Native( + "Invalid plan: self-keyed specification has separate key reads".into(), + )); + } }; // Precompute per-group setup (bucket_map, keys_source) once, before @@ -2900,7 +2873,7 @@ impl SimpleEngine { // this is actually a topk query with limiting requested -- gates // both the per-step sort/truncate below and nothing else, so a // non-topk query pays zero cost for this. - let row_label_order = spec.row_label_order.clone(); + let row_label_order = &spec.row_label_order; // Step-major: for each output timestamp, visit every group, not the // other way around. Required for topk correctness -- ranking a @@ -2908,7 +2881,7 @@ impl SimpleEngine { // timestamp before truncating, which a group-major loop can't do // (#581). One loop shape for topk and non-topk alike, rather than // maintaining two. - for ¤t_time in &spec.output_timestamps { + for ¤t_time in value_window.output_timestamps.iter() { let current_time_i64 = i64::try_from(current_time).map_err(|_| { QueryExecutionError::Native( "Output timestamp exceeds signed timestamp range".to_string(), @@ -2945,10 +2918,8 @@ impl SimpleEngine { let keys_precompute: Option> = match keys_source { KeysSource::PerStep { bucket_map: keys_bucket_map, - lookback_ms: keys_lookback_ms, - tumbling_window_ms: keys_tumbling_window_ms, - stored_window_size_ms: keys_stored_window_size_ms, - window_type: keys_window_type, + aggregation_type: key_accumulator_type, + window: keys_window, } => { // DeltaSetAggregator's keys window is always // [0, current_time) ("replay from the beginning"), @@ -2960,61 +2931,63 @@ impl SimpleEngine { // already bounded by real data. Bypass that walk // entirely for this aggregation type (#581 stage // E.4 review). - let keys_window_buckets = - if key_accumulator_type == AggregationType::DeltaSetAggregator { - // #588/#606 force DeltaSetAggregator's own - // config to Tumbling at planning time -- but - // that's a planner convention, not a runtime - // invariant this code can trust blindly. - // AggregationConfig can be (and in this crate's - // own tests routinely is) constructed directly, - // bypassing the planner. A Sliding DeltaSetAgg - // has no coherent "replay from the beginning" - // semantics to begin with, so this asserts - // rather than silently reinterpreting it (#581 - // stage E.4 review). - assert_eq!( - *keys_window_type, + let keys_window_buckets = if *key_accumulator_type + == AggregationType::DeltaSetAggregator + { + // #588/#606 force DeltaSetAggregator's own + // config to Tumbling at planning time -- but + // that's a planner convention, not a runtime + // invariant this code can trust blindly. + // AggregationConfig can be (and in this crate's + // own tests routinely is) constructed directly, + // bypassing the planner. A Sliding DeltaSetAgg + // has no coherent "replay from the beginning" + // semantics to begin with, so this asserts + // rather than silently reinterpreting it (#581 + // stage E.4 review). + assert_eq!( + keys_window.window_type, WindowType::Tumbling, "DeltaSetAggregator keys config must be Tumbling (#588/#606) -- \ the replay-from-the-beginning fast path has no correct meaning \ for Sliding" ); - let replay_end = if window_type == WindowType::Sliding { - Self::align_down_with_warning( - current_time, - tumbling_window_ms, - "DeltaSetAggregator replay end", - ) - } else { - current_time - }; - Self::collect_bucket_map_entries_before(keys_bucket_map, replay_end) + let replay_end = if value_window.window_type == WindowType::Sliding { + Self::align_down_with_warning( + current_time, + value_window.bucket_step_ms, + "DeltaSetAggregator replay end", + ) } else { - let keys_window_end = if *keys_window_type == WindowType::Sliding { - Self::align_down_with_warning( - current_time, - *keys_tumbling_window_ms, - "Sliding key window end", - ) - } else { - current_time - }; - let keys_window_start = - keys_window_end.saturating_sub(*keys_lookback_ms); - Self::window_buckets_for_step( - keys_bucket_map, - keys_window_start, - keys_window_end, - *keys_tumbling_window_ms, - *keys_window_type, - *keys_stored_window_size_ms, + current_time + }; + Self::collect_bucket_map_entries_before(keys_bucket_map, replay_end) + } else { + let keys_window_end = if keys_window.window_type == WindowType::Sliding + { + Self::align_down_with_warning( + current_time, + keys_window.bucket_step_ms, + "Sliding key window end", ) + } else { + current_time }; - - if *keys_window_type == WindowType::Sliding { + let keys_window_start = + keys_window_end.saturating_sub(keys_window.lookback_ms); + Self::window_buckets_for_step( + keys_bucket_map, + keys_window_start, + keys_window_end, + keys_window.bucket_step_ms, + keys_window.window_type, + keys_window.window_size_ms, + ) + }; + + if keys_window.window_type == WindowType::Sliding { let expected = - (*keys_lookback_ms / *keys_stored_window_size_ms) as usize; + (keys_window.lookback_ms / keys_window.window_size_ms) as usize; if keys_window_buckets.len() < expected { debug!( "Skipping incomplete Sliding key cover at t={}", @@ -3032,7 +3005,7 @@ impl SimpleEngine { continue; } - let mut key_merger = create_window_merger(key_accumulator_type); + let mut key_merger = create_window_merger(*key_accumulator_type); key_merger.initialize(keys_window_buckets); match key_merger.get_merged() { Ok(merged_keys) => Some(merged_keys), @@ -3047,10 +3020,10 @@ impl SimpleEngine { // Window covers [current_time - lookback_ms, current_time) // This means we look at buckets that START within this range - let window_end = if window_type == WindowType::Sliding { + let window_end = if value_window.window_type == WindowType::Sliding { Self::align_down_with_warning( current_time, - tumbling_window_ms, + value_window.bucket_step_ms, "Sliding value window end", ) } else { @@ -3062,24 +3035,24 @@ impl SimpleEngine { bucket_map, window_start, window_end, - tumbling_window_ms, - window_type, - spec.window_size_ms, + value_window.bucket_step_ms, + value_window.window_type, + value_window.window_size_ms, ); trace!( current_time, - window_type = ?window_type, + window_type = ?value_window.window_type, window_start, window_end, - grid_step_ms = tumbling_window_ms, - stored_window_size_ms = spec.window_size_ms, + grid_step_ms = value_window.bucket_step_ms, + stored_window_size_ms = value_window.window_size_ms, selected_bucket_count = window_buckets.len(), "Composed query output window from stored buckets" ); - if window_type == WindowType::Sliding { - let expected = (lookback_ms / spec.window_size_ms) as usize; + if value_window.window_type == WindowType::Sliding { + let expected = (lookback_ms / value_window.window_size_ms) as usize; if window_buckets.len() < expected { debug!( "Skipping incomplete Sliding value cover at t={}", @@ -3125,7 +3098,7 @@ impl SimpleEngine { &fallback_key, &spec.grouping_labels, &spec.aggregated_labels, - &row_label_order, + row_label_order, &statistic, query_kwargs, Some(&query_bounds), diff --git a/asap-query-engine/src/tests/native_range_query_tests.rs b/asap-query-engine/src/tests/native_range_query_tests.rs index c900a7e9..1405adb0 100644 --- a/asap-query-engine/src/tests/native_range_query_tests.rs +++ b/asap-query-engine/src/tests/native_range_query_tests.rs @@ -712,6 +712,57 @@ mod tests { )); } + #[tokio::test(flavor = "multi_thread")] + async fn range_query_empty_values_fail_but_empty_keys_are_allowed() { + let mut value = CountMinSketchAccumulator::new(2, 3); + value.inner.update("host-a;evt-1", 1.0); + let mut keys = SetAggregatorAccumulator::new(); + keys.add_key(KeyByLabelValues { + labels: vec!["host-a".to_string(), "evt-1".to_string()], + }); + let query = "count(event_frequency) by (host, event)"; + + let empty_keys_engine = create_range_engine_dual_input_with_windows( + "event_frequency", + AggregationType::CountMinSketch, + AggregationType::SetAggregator, + vec![], + vec!["host", "event"], + vec![(1_000, None, Box::new(value) as Box)], + vec![], + query, + 1_000, + 1_000, + ); + let empty_keys_result = empty_keys_engine + .handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0) + .expect("an empty keys read must not fail native execution"); + assert!( + empty_keys_result.is_some(), + "empty keys may produce no samples, but values data must still pass the read stage" + ); + + let empty_values_engine = create_range_engine_dual_input_with_windows( + "event_frequency", + AggregationType::CountMinSketch, + AggregationType::SetAggregator, + vec![], + vec!["host", "event"], + vec![], + vec![(1_000, None, Box::new(keys) as Box)], + query, + 1_000, + 1_000, + ); + let empty_values_result = empty_values_engine + .handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0) + .expect("empty values are reported as no local data, not an execution error"); + assert!( + empty_values_result.is_none(), + "values reads must retain the NoLocalData behavior" + ); + } + #[tokio::test(flavor = "multi_thread")] async fn range_query_tumbling_dual_population_keeps_key_steps_isolated() { let mut value_1 = CountMinSketchAccumulator::new(2, 3);