diff --git a/asap-query-engine/src/engines/query_plan.rs b/asap-query-engine/src/engines/query_plan.rs index 2384b75e..c826984b 100644 --- a/asap-query-engine/src/engines/query_plan.rs +++ b/asap-query-engine/src/engines/query_plan.rs @@ -5,30 +5,77 @@ 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 std::sync::Arc; 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(WindowCompositionSpec), +} + +impl StoreReadStrategy { + fn for_window(window: &WindowCompositionSpec) -> Self { + match window.window_type { + WindowType::Tumbling => Self::WindowGrid, + WindowType::Sliding => Self::SlidingExactCover(window.clone()), + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum StoreReadRole { + Values, + 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 query_range_ms: u64, + pub values: ValueEstimateInput, + pub keys: KeyInputSpec, + 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 +86,7 @@ pub(crate) enum QueryPlanNode { statistic: Statistic, query_kwargs: std::collections::HashMap, output_labels: KeyByLabelNames, + spec: RangeEstimateSpec, }, AggregateVector { input: NodeId, @@ -49,6 +97,7 @@ pub(crate) enum QueryPlanNode { input: NodeId, k: String, grouping_labels: KeyByLabelNames, + row_label_order: KeyByLabelNames, }, Format { input: NodeId, @@ -113,38 +162,32 @@ impl QueryPlan { query_time_aggregations: &[QueryTimeAggregation], ) -> Result { let mut nodes = Vec::new(); + let estimate_spec = RangeEstimateSpec::compile(context)?; let values_read = Self::push_read( &mut nodes, &context.base.store_plan.values_query, - context.window_type, + StoreReadRole::Values, + StoreReadStrategy::for_window(&estimate_spec.values.window), ); - let values = Self::push_compose( - &mut nodes, - values_read, - &context.output_timestamps, - context.query_range_ms, - context.window_size_ms, - context.tumbling_window_ms, - ); - let keys = context.base.store_plan.keys_query.as_ref().map(|query| { - let read = Self::push_read( - &mut nodes, - query, - context.keys_window_type.unwrap_or(context.window_type), - ); - 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), - ) - }); + let values = Self::push_prepare_buckets(&mut nodes, values_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 mut root = Self::push( &mut nodes, @@ -153,6 +196,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 +215,7 @@ impl QueryPlan { input: root, k, grouping_labels: context.base.grouping_labels.clone(), + row_label_order: estimate_spec.row_label_order.clone(), }, ); } @@ -285,62 +330,54 @@ 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({}, role={role:?}, {}#{}, [{}, {}])", + Self::describe_read_strategy(strategy), 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:?}, {})", + input.0, + Self::describe_timestamps(&spec.values.window.output_timestamps) + ) }, QueryPlanNode::LimitTopK { input, k, .. } => { format!("n{index} LimitTopK(n{}, k={k})", input.0) @@ -357,13 +394,84 @@ 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 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, + 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( + &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 +483,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, .. } @@ -460,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, @@ -473,7 +584,10 @@ 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(lookback=2000ms, window=1000ms, step=500ms")); assert!(explanation.ends_with("root: n5")); } @@ -528,7 +642,126 @@ mod tests { .unwrap() .explain(); - assert!(explanation.contains("outputs=[1000, 2000, 3000]")); + assert!(explanation.contains("outputs=count:3, first:1000, last: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 plan = QueryPlan::compile_range( + &context, + PlanOptions { + limit_topk: false, + format_output: false, + }, + &[], + ) + .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"), + } + } + + #[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] @@ -584,6 +817,7 @@ mod tests { statistic: Statistic::Sum, query_kwargs: HashMap::new(), output_labels: KeyByLabelNames::empty(), + spec: RangeEstimateSpec::compile(&context()).unwrap(), }], root: NodeId(0), }; @@ -622,15 +856,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..e131695d 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -9,7 +9,8 @@ use crate::data_model::{ StreamingConfig, }; use crate::engines::query_plan::{ - NodeId, PlanOptions, QueryPlan, QueryPlanExecutionError, QueryPlanNode, QueryPlanRuntime, + KeyInputSpec, 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,35 @@ 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(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!( + "No data found for metric: {}", + query.metric + ))); } - Ok(self - .reads - .borrow() - .as_ref() - .expect("reads initialized") - .clone()) + 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) } } @@ -259,25 +284,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 +317,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 +354,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 +1233,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 +2508,12 @@ 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( + context.window_type, + context.tumbling_window_ms, + &context.output_timestamps, + context.base.metadata.statistic_to_compute, + )?; #[cfg(feature = "native_query_legacy_test_support")] if matches!( self.native_range_execution_mode, @@ -2539,11 +2563,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 +2585,12 @@ impl SimpleEngine { enable_topk_formatting: bool, ) -> Result, QueryExecutionError> { let reads = self.read_range_query_inputs(context)?; + let estimate_spec = + RangeEstimateSpec::compile(context).map_err(QueryExecutionError::Native)?; 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 +2625,7 @@ impl SimpleEngine { )) } + #[cfg(feature = "native_query_legacy_test_support")] fn read_range_query_inputs( &self, context: &RangeQueryExecutionContext, @@ -2662,26 +2687,25 @@ impl SimpleEngine { } fn reject_off_grid_sliding_counter_query( - context: &RangeQueryExecutionContext, + window_type: WindowType, + bucket_step_ms: u64, + output_timestamps: &[u64], + 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 window_type != WindowType::Sliding + || bucket_step_ms == 0 + || !matches!(statistic, Statistic::Increase | Statistic::Rate) { return Ok(()); } - if let Some(×tamp) = context - .output_timestamps + if let Some(×tamp) = output_timestamps .iter() - .find(|&×tamp| !timestamp.is_multiple_of(context.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, context.tumbling_window_ms + timestamp, bucket_step_ms ))); } Ok(()) @@ -2689,7 +2713,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,43 +2725,22 @@ 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; + 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 = &context.base.agg_info.aggregation_type_for_value; - let key_accumulator_type = context.base.agg_info.aggregation_type_for_key; + let accumulator_type = &spec.values.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; - - // 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={}", - context.output_timestamps.len(), - context.output_timestamps.first().copied().unwrap_or(0), - context.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 @@ -2757,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 @@ -2774,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, }, } @@ -2792,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 context.keys_lookback_ms / - // context.keys_tumbling_window_ms 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 @@ -2881,11 +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 = Self::topk_row_label_order( - &context.base.metadata, - &context.base.grouping_labels, - &context.base.aggregated_labels, - ); + 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 @@ -2893,13 +2881,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 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(), ) })?; - 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(), ) @@ -2930,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"), @@ -2945,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={}", @@ -3017,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), @@ -3032,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 { @@ -3047,24 +3035,24 @@ impl SimpleEngine { bucket_map, window_start, window_end, - tumbling_window_ms, - window_type, - context.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 = context.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 / context.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={}", @@ -3108,11 +3096,11 @@ impl SimpleEngine { Some(merged.as_ref()), keys_precompute.as_deref(), &fallback_key, - &context.base.grouping_labels, - &context.base.aggregated_labels, - &row_label_order, - &context.base.metadata.statistic_to_compute, - &context.base.metadata.query_kwargs, + &spec.grouping_labels, + &spec.aggregated_labels, + row_label_order, + &statistic, + query_kwargs, Some(&query_bounds), ) .into_iter() 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);