diff --git a/asap-query-engine/src/engines/query_plan.rs b/asap-query-engine/src/engines/query_plan.rs index 2384b75e..4326b6f6 100644 --- a/asap-query-engine/src/engines/query_plan.rs +++ b/asap-query-engine/src/engines/query_plan.rs @@ -5,30 +5,97 @@ 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, + WindowGrid(Arc), + SlidingExactCover(Arc), +} + +impl StoreReadStrategy { + pub(crate) fn for_window(window: Arc, window_type: WindowType) -> Self { + match window_type { + WindowType::Tumbling => Self::WindowGrid(window.clone()), + WindowType::Sliding => Self::SlidingExactCover(window.clone()), + } + } + + pub(crate) fn window(&self) -> &Arc { + match self { + Self::WindowGrid(window) | Self::SlidingExactCover(window) => window, + } + } + + pub(crate) fn window_type(&self) -> WindowType { + match self { + Self::WindowGrid(_) => WindowType::Tumbling, + Self::SlidingExactCover(_) => WindowType::Sliding, + } + } +} + +#[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_size_ms: u64, + pub bucket_step_ms: u64, +} + +#[derive(Debug, Clone)] +pub(crate) struct ValueEstimateInput { + pub aggregation_type: AggregationType, + pub strategy: StoreReadStrategy, +} + +#[derive(Debug, Clone)] +pub(crate) enum KeyInputSpec { + FromValues, + Separate { + aggregation_type: AggregationType, + strategy: StoreReadStrategy, + }, +} + +#[derive(Debug, Clone)] +pub(crate) struct RangeLabelLayout { + pub output_labels: KeyByLabelNames, + pub grouping_labels: KeyByLabelNames, + pub aggregated_labels: KeyByLabelNames, + pub row_label_order: KeyByLabelNames, +} + +#[derive(Debug, Clone)] +pub(crate) struct RangeEstimateSpec { + pub query_range_ms: u64, + pub statistic: Statistic, + pub query_kwargs: std::collections::HashMap, + pub values: ValueEstimateInput, + pub keys: KeyInputSpec, + pub labels: Arc, } #[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, @@ -36,9 +103,7 @@ pub(crate) enum QueryPlanNode { }, Estimate { input: NodeId, - statistic: Statistic, - query_kwargs: std::collections::HashMap, - output_labels: KeyByLabelNames, + spec: RangeEstimateSpec, }, AggregateVector { input: NodeId, @@ -48,7 +113,7 @@ pub(crate) enum QueryPlanNode { LimitTopK { input: NodeId, k: String, - grouping_labels: KeyByLabelNames, + labels: Arc, }, Format { input: NodeId, @@ -113,52 +178,35 @@ 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, - ); - 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, + estimate_spec.values.strategy.clone(), ); - 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 &estimate_spec.keys { + KeyInputSpec::FromValues => None, + KeyInputSpec::Separate { strategy, .. } => { + let query = context.base.store_plan.keys_query.as_ref().ok_or_else(|| { + "Separate-key estimate is missing its store query".to_string() + })?; + let read = + Self::push_read(&mut nodes, query, StoreReadRole::Keys, strategy.clone()); + Some(Self::push_prepare_buckets(&mut nodes, read)) + } + }; let resolved = Self::push(&mut nodes, QueryPlanNode::ResolveKeys { values, keys }); let mut root = Self::push( &mut nodes, QueryPlanNode::Estimate { input: resolved, - 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 { - let k = context - .base - .metadata + if options.limit_topk && estimate_spec.statistic == Statistic::Topk { + let k = estimate_spec .query_kwargs .get("k") .cloned() @@ -170,12 +218,12 @@ impl QueryPlan { QueryPlanNode::LimitTopK { input: root, k, - grouping_labels: context.base.grouping_labels.clone(), + labels: Arc::clone(&estimate_spec.labels), }, ); } let materialize_metric_name = !query_time_aggregations.is_empty() - && context.base.metadata.statistic_to_compute == Statistic::Topk + && estimate_spec.statistic == Statistic::Topk && context.base.metadata.keep_metric_name; if materialize_metric_name { root = Self::push( @@ -187,7 +235,7 @@ impl QueryPlan { }, ); } - let mut labels = context.base.metadata.query_output_labels.clone(); + let mut labels = estimate_spec.labels.output_labels.clone(); for aggregation in query_time_aggregations { root = Self::push( &mut nodes, @@ -204,8 +252,7 @@ impl QueryPlan { &mut nodes, QueryPlanNode::Format { input: root, - include_metric_name: context.base.metadata.statistic_to_compute - == Statistic::Topk + include_metric_name: estimate_spec.statistic == Statistic::Topk && context.base.metadata.keep_metric_name, metric: context.base.metric.clone(), }, @@ -236,6 +283,64 @@ impl QueryPlan { )); } } + if let QueryPlanNode::Estimate { input, spec } = node { + self.validate_estimate_reads(*input, spec)?; + } + } + Ok(()) + } + + fn validate_estimate_reads( + &self, + resolved: NodeId, + spec: &RangeEstimateSpec, + ) -> Result<(), String> { + let QueryPlanNode::ResolveKeys { values, keys } = &self.nodes[resolved.0] else { + return Err(format!( + "Estimate input n{} must resolve value and key reads", + resolved.0 + )); + }; + self.validate_read_window(*values, StoreReadRole::Values, &spec.values.strategy)?; + match (&spec.keys, keys) { + (KeyInputSpec::FromValues, None) => Ok(()), + (KeyInputSpec::Separate { strategy, .. }, Some(keys)) => { + self.validate_read_window(*keys, StoreReadRole::Keys, strategy) + } + (KeyInputSpec::FromValues, Some(_)) => { + Err("Self-keyed estimate must not have a separate key read".to_string()) + } + (KeyInputSpec::Separate { .. }, None) => { + Err("Separate-key estimate must have a key read".to_string()) + } + } + } + + fn validate_read_window( + &self, + prepared: NodeId, + expected_role: StoreReadRole, + expected_strategy: &StoreReadStrategy, + ) -> Result<(), String> { + let QueryPlanNode::PrepareBuckets { input } = &self.nodes[prepared.0] else { + return Err(format!("Read n{} must prepare buckets", prepared.0)); + }; + let QueryPlanNode::StoreRead { role, strategy, .. } = &self.nodes[input.0] else { + return Err(format!( + "Prepared node n{} must read from the store", + prepared.0 + )); + }; + if *role != expected_role { + return Err(format!("Store read n{} has the wrong role", input.0)); + } + if strategy.window_type() != expected_strategy.window_type() + || !Arc::ptr_eq(strategy.window(), expected_strategy.window()) + { + return Err(format!( + "Store read n{} window does not match its estimate", + input.0 + )); } Ok(()) } @@ -285,62 +390,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({}, 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, .. } => { - let mut kwargs: Vec<_> = query_kwargs.iter().collect(); + QueryPlanNode::Estimate { + input, + spec, + } => { + let mut kwargs: Vec<_> = spec.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{}, {}, {kwargs:?}, {})", + input.0, + spec.statistic, + Self::describe_timestamps(&spec.values.strategy.window().output_timestamps) + ) }, QueryPlanNode::LimitTopK { input, k, .. } => { format!("n{index} LimitTopK(n{}, k={k})", input.0) @@ -357,13 +452,93 @@ impl QueryPlan { lines.push(format!("root: n{}", self.root.0)); lines.join("\n") } + + fn describe_read_strategy(strategy: &StoreReadStrategy) -> String { + let window = strategy.window(); + format!( + "{}(lookback={}ms, window={}ms, step={}ms, {})", + match strategy { + StoreReadStrategy::WindowGrid(_) => "WindowGrid", + StoreReadStrategy::SlidingExactCover(_) => "SlidingExactCover", + }, + 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.as_slice()); + let values_window = Arc::new(WindowCompositionSpec { + output_timestamps: Arc::clone(&output_timestamps), + lookback_ms: (context.lookback_bucket_count as u64) * context.tumbling_window_ms, + window_size_ms: context.window_size_ms, + bucket_step_ms: context.tumbling_window_ms, + }); + let values = ValueEstimateInput { + aggregation_type: context.base.agg_info.aggregation_type_for_value, + strategy: StoreReadStrategy::for_window(values_window, context.window_type), + }; + 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, + strategy: StoreReadStrategy::for_window( + Arc::new(WindowCompositionSpec { + output_timestamps, + lookback_ms: context.keys_lookback_ms.ok_or_else(|| { + "Separate keys query is missing its lookback".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() + })?, + }), + context.keys_window_type.ok_or_else(|| { + "Separate keys query is missing its window type".to_string() + })?, + ), + }, + }; + Ok(Self { + query_range_ms: context.query_range_ms, + statistic: context.base.metadata.statistic_to_compute, + query_kwargs: context.base.metadata.query_kwargs.clone(), + values, + keys, + labels: Arc::new(RangeLabelLayout { + output_labels: context.base.metadata.query_output_labels.clone(), + 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 +550,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, .. } @@ -438,7 +613,6 @@ mod tests { }, output_timestamps: vec![1_000], query_range_ms: 1_000, - buckets_per_step: 1, lookback_bucket_count: 1, tumbling_window_ms: 1_000, window_type: WindowType::Tumbling, @@ -460,6 +634,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 +650,12 @@ 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 + .contains("n0 StoreRead(WindowGrid(lookback=1000ms, window=1000ms, step=1000ms")); assert!(explanation.ends_with("root: n5")); } @@ -528,7 +710,182 @@ 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.strategy.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"); + } + + #[test] + fn compiles_off_grid_sliding_counter_plan_for_engine_preflight() { + let mut context = context(); + context.window_type = WindowType::Sliding; + context.output_timestamps = vec![1_500]; + context.base.metadata.statistic_to_compute = Statistic::Rate; + + QueryPlan::compile_range( + &context, + PlanOptions { + limit_topk: false, + format_output: false, + }, + &[], + ) + .expect("the engine, not the compiled plan, classifies off-grid fallback"); + } + + #[test] + fn rejects_a_read_window_that_does_not_match_its_estimate() { + let mut context = context(); + context.window_type = WindowType::Sliding; + let mut plan = QueryPlan::compile_range( + &context, + PlanOptions { + limit_topk: false, + format_output: false, + }, + &[], + ) + .unwrap(); + + let QueryPlanNode::StoreRead { strategy, .. } = &mut plan.nodes[0] else { + panic!("first node must read values"); + }; + *strategy = StoreReadStrategy::WindowGrid(Arc::new(WindowCompositionSpec { + output_timestamps: Arc::from([1_000]), + lookback_ms: 1_000, + window_size_ms: 1_000, + bucket_step_ms: 1_000, + })); + + assert_eq!( + plan.validate() + .expect_err("mismatched plan must fail validation"), + "Store read n0 window does not match its estimate" + ); + } + + 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(Arc::new(WindowCompositionSpec { + output_timestamps: Arc::from([1_000]), + lookback_ms: 2_000, + window_size_ms: 1_000, + bucket_step_ms: 1_000, + })), + ), + ( + StoreReadRole::Keys, + StoreReadStrategy::WindowGrid(Arc::new(WindowCompositionSpec { + output_timestamps: Arc::from([1_000]), + lookback_ms: 3_000, + window_size_ms: 1_000, + bucket_step_ms: 1_000, + })), + ), + ] + ); } #[test] @@ -581,9 +938,7 @@ mod tests { let plan = QueryPlan { nodes: vec![QueryPlanNode::Estimate { input: NodeId(1), - statistic: Statistic::Sum, - query_kwargs: HashMap::new(), - output_labels: KeyByLabelNames::empty(), + spec: RangeEstimateSpec::compile(&context()).unwrap(), }], root: NodeId(0), }; @@ -622,15 +977,15 @@ mod tests { start_timestamp: 0, end_timestamp: 1, }, - strategy: StoreReadStrategy::WindowGrid, - }, - QueryPlanNode::ComposeWindows { - input: NodeId(0), - output_timestamps: vec![1], - lookback_ms: 1, - window_size_ms: 1, - bucket_step_ms: 1, + role: StoreReadRole::Values, + strategy: StoreReadStrategy::WindowGrid(Arc::new(WindowCompositionSpec { + output_timestamps: Arc::from([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..a034828b 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -9,8 +9,11 @@ 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, }; +#[cfg(feature = "native_query_legacy_test_support")] +use crate::engines::query_plan::{RangeLabelLayout, ValueEstimateInput, WindowCompositionSpec}; use crate::engines::query_result::{InstantVectorElement, QueryResult}; use crate::engines::sliding_window_composition::{ plan_exact_cover, CompositionError, SlidingWindowSpec, @@ -142,8 +145,6 @@ pub struct RangeQueryExecutionContext { pub output_timestamps: Vec, /// Exact range-vector duration used for extrapolation, in milliseconds. pub query_range_ms: u64, - /// Number of buckets per step (step / tumbling_window) - pub buckets_per_step: usize, /// Number of buckets in lookback window pub lookback_bucket_count: usize, /// Tumbling window size in ms @@ -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,12 +317,12 @@ impl QueryPlanRuntime for NativePlanRuntime<'_> { "ResolveKeys received incompatible inputs".into(), )), }, - QueryPlanNode::Estimate { output_labels, .. } => match inputs { + QueryPlanNode::Estimate { spec, .. } => match inputs { [NativePlanOutput::Resolved(reads)] => self .engine - .estimate_range_query(self.context, reads.clone()) + .estimate_range_query(spec, reads.clone()) .map(|values| NativePlanOutput::Results { - labels: output_labels.clone(), + labels: spec.labels.output_labels.clone(), values, }), _ => Err(QueryExecutionError::Native( @@ -329,20 +348,11 @@ impl QueryPlanRuntime for NativePlanRuntime<'_> { )), }, QueryPlanNode::LimitTopK { - k, grouping_labels, .. + k, labels: layout, .. } => 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, &layout.row_label_order, &layout.grouping_labels) .map_err(QueryExecutionError::Native) .map(|values| NativePlanOutput::Results { labels: labels.clone(), @@ -892,10 +902,6 @@ impl SimpleEngine { }, output_timestamps: vec![query_time], query_range_ms: lookback_ms, - // Placeholder: no real "step" for a single instant point. Only - // feeds a debug-log string today -- not type-enforced, recheck - // before using it for anything functional. - buckets_per_step: 1, lookback_bucket_count, tumbling_window_ms, window_type, @@ -1214,7 +1220,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 +2495,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 +2550,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 +2572,9 @@ impl SimpleEngine { enable_topk_formatting: bool, ) -> Result, QueryExecutionError> { let reads = self.read_range_query_inputs(context)?; + let estimate_spec = Self::legacy_range_estimate_spec(context)?; let mut results = self.estimate_range_query( - context, + &estimate_spec, ResolvedRangeReads { values: self.compose_range_read(&reads.values), keys: reads @@ -2601,6 +2609,72 @@ impl SimpleEngine { )) } + #[cfg(feature = "native_query_legacy_test_support")] + fn legacy_range_estimate_spec( + context: &RangeQueryExecutionContext, + ) -> Result { + let output_timestamps: Arc<[u64]> = Arc::from(context.output_timestamps.as_slice()); + let values_window = Arc::new(WindowCompositionSpec { + output_timestamps: Arc::clone(&output_timestamps), + lookback_ms: (context.lookback_bucket_count as u64) * context.tumbling_window_ms, + 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(_) => { + let window_type = context.keys_window_type.ok_or_else(|| { + QueryExecutionError::Native( + "Separate keys query is missing its window type".into(), + ) + })?; + let window = Arc::new(WindowCompositionSpec { + output_timestamps, + lookback_ms: context.keys_lookback_ms.ok_or_else(|| { + QueryExecutionError::Native( + "Separate keys query is missing its lookback".into(), + ) + })?, + window_size_ms: context.keys_window_size_ms.ok_or_else(|| { + QueryExecutionError::Native( + "Separate keys query is missing its window size".into(), + ) + })?, + bucket_step_ms: context.keys_tumbling_window_ms.ok_or_else(|| { + QueryExecutionError::Native( + "Separate keys query is missing its bucket step".into(), + ) + })?, + }); + KeyInputSpec::Separate { + aggregation_type: context.base.agg_info.aggregation_type_for_key, + strategy: StoreReadStrategy::for_window(window, window_type), + } + } + }; + Ok(RangeEstimateSpec { + query_range_ms: context.query_range_ms, + statistic: context.base.metadata.statistic_to_compute, + query_kwargs: context.base.metadata.query_kwargs.clone(), + values: ValueEstimateInput { + aggregation_type: context.base.agg_info.aggregation_type_for_value, + strategy: StoreReadStrategy::for_window(values_window, context.window_type), + }, + keys, + labels: Arc::new(RangeLabelLayout { + output_labels: context.base.metadata.query_output_labels.clone(), + grouping_labels: context.base.grouping_labels.clone(), + aggregated_labels: context.base.aggregated_labels.clone(), + row_label_order: Self::topk_row_label_order( + &context.base.metadata, + &context.base.grouping_labels, + &context.base.aggregated_labels, + ), + }), + }) + } + + #[cfg(feature = "native_query_legacy_test_support")] fn read_range_query_inputs( &self, context: &RangeQueryExecutionContext, @@ -2662,26 +2736,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 +2762,7 @@ impl SimpleEngine { fn estimate_range_query( &self, - context: &RangeQueryExecutionContext, + spec: &RangeEstimateSpec, reads: ResolvedRangeReads, ) -> Result, QueryExecutionError> { use crate::engines::query_result::RangeVectorElement; @@ -2699,173 +2772,94 @@ 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_strategy = &spec.values.strategy; + let value_window = value_strategy.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 - // depends on the query SHAPE (dual- vs single-population), not on a - // per-group fallback — mirrors collect_all_results exactly: - // - dual-population (KeysSource::PerStep below, separate - // keys_query present): always expand via the keys aggregation's - // per-step merge (#583). The value accumulator's own get_keys() - // is never consulted, even if the value accumulator itself - // happens to be self-keyed (e.g. a CountMinSketchWithHeap value - // paired with a DeltaSetAggregator keys aggregation is a real - // capability-matched config, see sql.rs). Otherwise a - // self-keyed value accumulator's own (possibly different, - // window-to-window-shifting) keys would silently override the - // keys aggregation's expansion. See #587 review. - // - single-population (KeysSource::Fixed below, no separate - // keys_query): the value accumulator's own get_keys() takes - // priority whenever present (#584, self-keyed accumulators like - // 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 - // further down (#581 stage E.4 review: previously duplicated as the - // raw type at the PerStep site instead of using this alias). + // Separate keys always determine expansion; self-keyed values use + // their own keys after the value window is merged. 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, + strategy: &'a StoreReadStrategy, }, } - // Resolve, for every value group, which groups exist at all (a - // one-time operation — see design doc) and where their expansion - // keys come from. A group with keys data but no value data - // anywhere in the queried range is skipped with a warning instead - // 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 \ + // Groups with keys but no values are skipped rather than failing the + // complete range query. + let groups: Vec<(GroupBucketMap, KeysSource)> = match (&spec.keys, &keys_raw_data) { + ( + KeyInputSpec::Separate { + aggregation_type, + strategy, + .. + }, + 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, + strategy, + }, + )), + 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() - } - // #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 + of failing the whole query", + group_key + ); + None + } + }, + ) + .collect(), + // Keep every group because self-keyed accumulators may use a + // None outer group key. + (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 - // the step-major loop below revisits every group at every output - // timestamp -- doing this per-step instead would repeat identical - // work once per timestamp instead of once per group. - // - // Memory tradeoff vs. the old group-major shape (#581 stage E.2 - // review): every group's value bucket_map is now held simultaneously - // for the whole step-major loop's duration, instead of one group's - // bucket_map at a time (built, used, dropped, next group). Keys-side - // PerStep bucket maps were already built eagerly for every group - // beforehand (see `groups` above), so this brings the value side in - // line with that, not a new pattern -- but for a range query over a - // very high-cardinality label set this is a real (if likely modest) - // increase in peak memory. Inherent to step-major: ranking a - // timestamp's candidates needs every group's bucket_map available at - // that timestamp, so they can't be built lazily one group at a time - // anymore. + // Step-major ranking needs every group's buckets at each timestamp. for (bucket_map, keys_source) in &groups { debug!( "Group with {} start-timestamps ({} keys start-timestamps)", @@ -2881,25 +2875,16 @@ 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.labels.row_label_order; - // Step-major: for each output timestamp, visit every group, not the - // other way around. Required for topk correctness -- ranking a - // timestamp's candidates means seeing every group's value at that - // 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 { + // Visit every group at each timestamp so ranking sees all candidates. + 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(), ) @@ -2920,21 +2905,15 @@ impl SimpleEngine { let mut step_results: Vec<(KeyByLabelValues, f64)> = Vec::new(); for (bucket_map, keys_source) in &groups { - // #583: dual-population groups resolve their expansion keys - // from the keys aggregation, per step — not a single - // snapshot reused for every step. If nothing resolves at - // this step, skip it before ever touching the value merge - // below (avoids wasted merge work on steps outside the - // key's lifetime). Fixed (single-population) groups have no - // separate keys accumulator to merge here at all. + // Separate keys resolve per step; self-keyed groups do not + // merge a separate key accumulator. 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, + strategy: keys_strategy, } => { + let keys_window = keys_strategy.window(); // DeltaSetAggregator's keys window is always // [0, current_time) ("replay from the beginning"), // which saturating_sub's keys_window_start to 0 -- @@ -2943,63 +2922,55 @@ impl SimpleEngine { // (up to ~1e8 positions for a real timestamp) purely // to see what's in keys_bucket_map, an in-memory map // already bounded by real data. Bypass that walk - // entirely for this aggregation type (#581 stage - // E.4 review). + // entirely for this aggregation 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). + if *key_accumulator_type == AggregationType::DeltaSetAggregator { + // Replaying from the beginning has no coherent + // Sliding semantics. assert_eq!( - *keys_window_type, - WindowType::Tumbling, - "DeltaSetAggregator keys config must be Tumbling (#588/#606) -- \ + keys_strategy.window_type(), + WindowType::Tumbling, + "DeltaSetAggregator keys config must be Tumbling -- \ 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 - }; + ); + let replay_end = + if value_strategy.window_type() == WindowType::Sliding { + Self::align_down_with_warning( + current_time, + value_window.bucket_step_ms, + "DeltaSetAggregator replay end", + ) + } else { + current_time + }; Self::collect_bucket_map_entries_before(keys_bucket_map, 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_end = + if keys_strategy.window_type() == WindowType::Sliding { + Self::align_down_with_warning( + current_time, + keys_window.bucket_step_ms, + "Sliding key window end", + ) + } else { + current_time + }; let keys_window_start = - keys_window_end.saturating_sub(*keys_lookback_ms); + keys_window_end.saturating_sub(keys_window.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, + keys_window.bucket_step_ms, + keys_strategy.window_type(), + keys_window.window_size_ms, ) }; - if *keys_window_type == WindowType::Sliding { + if keys_strategy.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 +2988,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 +3003,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_strategy.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 +3018,24 @@ impl SimpleEngine { bucket_map, window_start, window_end, - tumbling_window_ms, - window_type, - context.window_size_ms, + value_window.bucket_step_ms, + value_strategy.window_type(), + value_window.window_size_ms, ); trace!( current_time, - window_type = ?window_type, + window_type = ?value_strategy.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_strategy.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 +3079,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.labels.grouping_labels, + &spec.labels.aggregated_labels, + row_label_order, + &spec.statistic, + &spec.query_kwargs, Some(&query_bounds), ) .into_iter() diff --git a/asap-query-engine/src/engines/simple_engine/promql.rs b/asap-query-engine/src/engines/simple_engine/promql.rs index 0ecc764d..3e79dbca 100644 --- a/asap-query-engine/src/engines/simple_engine/promql.rs +++ b/asap-query-engine/src/engines/simple_engine/promql.rs @@ -640,7 +640,6 @@ impl SimpleEngine { let lookback_ms = Self::widen_query_window(&mut extended_store_plan.values_query, start_ms, end_ms); - let buckets_per_step = (step_ms / tumbling_window_ms) as usize; let lookback_bucket_count = (lookback_ms / tumbling_window_ms) as usize; // #583: widen keys_query the same way, using the instant window @@ -692,7 +691,6 @@ impl SimpleEngine { // start_ms+step_ms, ..., the last value <= end_ms. output_timestamps: (start_ms..=end_ms).step_by(step_ms as usize).collect(), query_range_ms: lookback_ms, - buckets_per_step, lookback_bucket_count, tumbling_window_ms, window_type, 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..c8db2dc9 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,56 @@ 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"); + let (_, empty_keys_result) = + empty_keys_result.expect("values data must pass the read stage when keys are empty"); + assert!(matrix_values(empty_keys_result).is_empty()); + + 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);