diff --git a/asap-query-engine/Cargo.toml b/asap-query-engine/Cargo.toml index 73249126..ed413a18 100644 --- a/asap-query-engine/Cargo.toml +++ b/asap-query-engine/Cargo.toml @@ -97,3 +97,4 @@ jemalloc = ["dep:tikv-jemallocator"] lock_profiling = [] # Enable extra debugging output extra_debugging = [] +native_query_legacy_test_support = [] diff --git a/asap-query-engine/src/drivers/query/adapters/traits.rs b/asap-query-engine/src/drivers/query/adapters/traits.rs index 44031b62..edfec9f2 100644 --- a/asap-query-engine/src/drivers/query/adapters/traits.rs +++ b/asap-query-engine/src/drivers/query/adapters/traits.rs @@ -60,6 +60,7 @@ impl std::error::Error for AdapterError {} /// Trait for parsing incoming HTTP requests into internal query format /// Handles Axum extractors directly for different request types (GET/POST) +#[allow(clippy::double_must_use)] #[async_trait] pub trait QueryRequestAdapter: Send + Sync { /// Parse a GET request with query parameters @@ -106,6 +107,7 @@ pub trait QueryRequestAdapter: Send + Sync { } /// Trait for formatting query results into protocol-specific HTTP responses +#[allow(clippy::double_must_use)] #[async_trait] pub trait QueryResponseAdapter: Send + Sync { /// Format a successful query result into protocol response @@ -135,6 +137,7 @@ pub trait QueryResponseAdapter: Send + Sync { /// define separate adapter traits. /// /// Note: Fallback logic is handled separately via FallbackClient +#[allow(clippy::double_must_use)] #[async_trait] pub trait HttpProtocolAdapter: QueryRequestAdapter + QueryResponseAdapter + Send + Sync { /// Get a descriptive name for this adapter (for logging/debugging) diff --git a/asap-query-engine/src/drivers/query/fallback/mod.rs b/asap-query-engine/src/drivers/query/fallback/mod.rs index 5add3333..c379b415 100644 --- a/asap-query-engine/src/drivers/query/fallback/mod.rs +++ b/asap-query-engine/src/drivers/query/fallback/mod.rs @@ -37,6 +37,7 @@ impl IntoResponse for FallbackResponse { } /// Client for forwarding unsupported queries to a fallback backend +#[allow(clippy::double_must_use)] #[async_trait] pub trait FallbackClient: Send + Sync { /// Execute a query against the fallback backend diff --git a/asap-query-engine/src/engines/mod.rs b/asap-query-engine/src/engines/mod.rs index a33d545f..27898c14 100644 --- a/asap-query-engine/src/engines/mod.rs +++ b/asap-query-engine/src/engines/mod.rs @@ -6,5 +6,7 @@ pub(crate) mod sliding_window_composition; pub mod window_merger; pub use query_result::{InstantVector, QueryResult, RangeVector, RangeVectorElement, Sample}; +#[cfg(feature = "native_query_legacy_test_support")] +pub use simple_engine::NativeRangeExecutionMode; pub use simple_engine::{QueryExecutionError, SimpleEngine}; pub use window_merger::{create_window_merger, NaiveMerger, WindowMerger}; diff --git a/asap-query-engine/src/engines/query_plan.rs b/asap-query-engine/src/engines/query_plan.rs index 4499d6ad..629d1464 100644 --- a/asap-query-engine/src/engines/query_plan.rs +++ b/asap-query-engine/src/engines/query_plan.rs @@ -2,7 +2,9 @@ use crate::engines::simple_engine::{RangeQueryExecutionContext, StoreQueryParams}; use asap_types::enums::WindowType; +use promql_utilities::data_model::KeyByLabelNames; use promql_utilities::query_logics::enums::Statistic; +use tracing::debug; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) struct NodeId(usize); @@ -38,10 +40,12 @@ pub(crate) enum QueryPlanNode { LimitTopK { input: NodeId, k: String, + grouping_labels: KeyByLabelNames, }, Format { input: NodeId, include_metric_name: bool, + metric: String, }, } @@ -57,7 +61,44 @@ pub(crate) struct PlanOptions { pub format_output: bool, } +pub(crate) trait QueryPlanRuntime { + type Output: Clone; + type Error: std::fmt::Display; + + fn execute_node( + &self, + id: NodeId, + node: &QueryPlanNode, + inputs: &[Self::Output], + ) -> Result; +} + +#[derive(Debug)] +pub(crate) enum QueryPlanExecutionError { + InvalidPlan(String), + Node { id: NodeId, source: E }, +} + +impl std::fmt::Display for QueryPlanExecutionError { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::InvalidPlan(error) => write!(formatter, "invalid query plan: {error}"), + Self::Node { id, source } => { + write!(formatter, "Query plan node n{} failed: {source}", id.0) + } + } + } +} + impl QueryPlan { + #[cfg(feature = "native_query_legacy_test_support")] + pub(crate) fn malformed_for_test() -> Self { + Self { + nodes: Vec::new(), + root: NodeId(0), + } + } + pub(crate) fn compile_range( context: &RangeQueryExecutionContext, options: PlanOptions, @@ -114,18 +155,90 @@ impl QueryPlan { .ok_or_else(|| "Topk query is missing required `k` parameter".to_string())?; k.parse::() .map_err(|_| "Topk query has an invalid `k` parameter".to_string())?; - root = Self::push(&mut nodes, QueryPlanNode::LimitTopK { input: root, k }); + root = Self::push( + &mut nodes, + QueryPlanNode::LimitTopK { + input: root, + k, + grouping_labels: context.base.grouping_labels.clone(), + }, + ); } if options.format_output { root = Self::push( &mut nodes, QueryPlanNode::Format { input: root, - include_metric_name: context.base.metadata.keep_metric_name, + include_metric_name: context.base.metadata.statistic_to_compute + == Statistic::Topk + && context.base.metadata.keep_metric_name, + metric: context.base.metric.clone(), }, ); } - Ok(Self { nodes, root }) + let plan = Self { nodes, root }; + plan.validate()?; + Ok(plan) + } + + /// Rejects plans whose node dependencies cannot be executed safely. + pub(crate) fn validate(&self) -> Result<(), String> { + if self.nodes.is_empty() { + return Err("Query plan has no nodes".to_string()); + } + if self.root.0 != self.nodes.len() - 1 { + return Err(format!( + "Query plan root n{} does not include every node", + self.root.0 + )); + } + for (index, node) in self.nodes.iter().enumerate() { + for input in node.inputs() { + if input.0 >= index { + return Err(format!( + "Query plan node n{index} references unavailable input n{}", + input.0 + )); + } + } + } + Ok(()) + } + + pub(crate) fn execute( + &self, + runtime: &R, + ) -> Result> { + self.validate() + .map_err(QueryPlanExecutionError::InvalidPlan)?; + let mut outputs: Vec = Vec::with_capacity(self.nodes.len()); + for (index, node) in self.nodes.iter().enumerate() { + let inputs = node + .inputs() + .into_iter() + .map(|input| outputs[input.0].clone()) + .collect::>(); + debug!( + node_id = index, + node_kind = node.kind(), + input_count = inputs.len(), + "Executing native query plan node" + ); + let output = runtime + .execute_node(NodeId(index), node, &inputs) + .map_err(|source| QueryPlanExecutionError::Node { + id: NodeId(index), + source, + })?; + debug!( + node_id = index, + node_kind = node.kind(), + "Completed native query plan node" + ); + outputs.push(output); + } + debug!(root_node_id = self.root.0, "Completed native query plan"); + Ok(outputs[self.root.0].clone()) } fn push(nodes: &mut Vec, node: QueryPlanNode) -> NodeId { @@ -194,9 +307,11 @@ impl QueryPlan { kwargs.sort_unstable_by_key(|(key, _)| *key); format!("n{index} Estimate(n{}, {statistic}, {kwargs:?})", input.0) }, - QueryPlanNode::LimitTopK { input, k } => format!("n{index} LimitTopK(n{}, k={k})", input.0), - QueryPlanNode::Format { input, include_metric_name } => format!( - "n{index} Format(n{}, include_metric_name={include_metric_name})", input.0 + QueryPlanNode::LimitTopK { input, k, .. } => { + format!("n{index} LimitTopK(n{}, k={k})", input.0) + } + QueryPlanNode::Format { input, include_metric_name, metric } => format!( + "n{index} Format(n{}, include_metric_name={include_metric_name}) metric={metric}", input.0 ), }; lines.push(line); @@ -206,6 +321,35 @@ impl QueryPlan { } } +impl QueryPlanNode { + fn kind(&self) -> &'static str { + match self { + Self::StoreRead { .. } => "StoreRead", + Self::ComposeWindows { .. } => "ComposeWindows", + Self::ResolveKeys { .. } => "ResolveKeys", + Self::Estimate { .. } => "Estimate", + Self::LimitTopK { .. } => "LimitTopK", + Self::Format { .. } => "Format", + } + } + + fn inputs(&self) -> Vec { + match self { + Self::StoreRead { .. } => Vec::new(), + Self::ComposeWindows { input, .. } + | Self::Estimate { input, .. } + | Self::LimitTopK { input, .. } + | Self::Format { input, .. } => vec![*input], + Self::ResolveKeys { values, keys } => { + keys.iter().copied().fold(vec![*values], |mut inputs, key| { + inputs.push(key); + inputs + }) + } + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -213,6 +357,7 @@ mod tests { use crate::engines::simple_engine::{QueryExecutionContext, QueryMetadata, StoreQueryPlan}; use promql_utilities::data_model::KeyByLabelNames; use promql_utilities::query_logics::enums::AggregationType; + use std::cell::RefCell; use std::collections::HashMap; fn context() -> RangeQueryExecutionContext { @@ -347,4 +492,66 @@ mod tests { assert_eq!(error, "Topk query is missing required `k` parameter"); } + + #[test] + fn rejects_a_node_that_references_a_later_node() { + let plan = QueryPlan { + nodes: vec![QueryPlanNode::Estimate { + input: NodeId(1), + statistic: Statistic::Sum, + query_kwargs: HashMap::new(), + }], + root: NodeId(0), + }; + + assert_eq!( + plan.validate().expect_err("invalid plan must fail loudly"), + "Query plan node n0 references unavailable input n1" + ); + } + + struct RecordingRuntime(RefCell>); + + impl QueryPlanRuntime for RecordingRuntime { + type Output = usize; + type Error = std::convert::Infallible; + + fn execute_node( + &self, + id: NodeId, + _node: &QueryPlanNode, + inputs: &[Self::Output], + ) -> Result { + self.0.borrow_mut().push(id.0); + Ok(1 + inputs.iter().sum::()) + } + } + + #[test] + fn executes_nodes_once_in_dependency_order() { + let plan = QueryPlan { + nodes: vec![ + QueryPlanNode::StoreRead { + query: StoreQueryParams { + metric: "requests".into(), + aggregation_id: 7, + 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, + }, + ], + root: NodeId(1), + }; + let runtime = RecordingRuntime(RefCell::new(Vec::new())); + assert_eq!(plan.execute(&runtime).unwrap(), 2); + assert_eq!(*runtime.0.borrow(), vec![0, 1]); + } } diff --git a/asap-query-engine/src/engines/simple_engine/mod.rs b/asap-query-engine/src/engines/simple_engine/mod.rs index cca337a7..2f42cc86 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -8,7 +8,9 @@ use crate::data_model::{ AggregationIdInfo, InferenceConfig, KeyByLabelValues, QueryBounds, QueryConfig, QueryLanguage, StreamingConfig, }; -use crate::engines::query_plan::{PlanOptions, QueryPlan}; +use crate::engines::query_plan::{ + NodeId, PlanOptions, QueryPlan, QueryPlanExecutionError, QueryPlanNode, QueryPlanRuntime, +}; use crate::engines::query_result::{InstantVectorElement, QueryResult}; use crate::engines::sliding_window_composition::{ plan_exact_cover, CompositionError, SlidingWindowSpec, @@ -65,6 +67,15 @@ pub enum QueryExecutionError { Native(String), } +#[cfg(feature = "native_query_legacy_test_support")] +#[derive(Clone, Copy)] +pub enum NativeRangeExecutionMode { + Dag, + Legacy, + MalformedPlan, + FailingStore, +} + /// Parameters for a single store query #[derive(Debug, Clone)] pub struct StoreQueryParams { @@ -171,6 +182,161 @@ pub struct RangeQueryExecutionContext { pub keys_tumbling_window_ms: Option, } +#[derive(Clone)] +struct RangeQueryReads { + values: TimestampedBucketsMap, + keys: Option, +} + +type BucketMap = HashMap>>; +type TopkPartition = (u64, Vec); +type TopkCandidates<'a> = HashMap>; + +#[derive(Clone)] +struct ComposedRangeRead { + groups: HashMap, BucketMap>, +} + +#[derive(Clone)] +struct ResolvedRangeReads { + values: ComposedRangeRead, + keys: Option, +} + +#[derive(Clone)] +enum NativePlanOutput { + Read(TimestampedBucketsMap), + Composed(ComposedRangeRead), + Resolved(ResolvedRangeReads), + Results(Vec), +} + +struct NativePlanRuntime<'a> { + engine: &'a SimpleEngine, + context: &'a RangeQueryExecutionContext, + reads: std::cell::RefCell>, +} + +impl NativePlanRuntime<'_> { + fn reads(&self) -> Result { + #[cfg(feature = "native_query_legacy_test_support")] + if matches!( + self.engine.native_range_execution_mode, + NativeRangeExecutionMode::FailingStore + ) { + return Err(QueryExecutionError::Native( + "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)?); + } + Ok(self + .reads + .borrow() + .as_ref() + .expect("reads initialized") + .clone()) + } +} + +impl QueryPlanRuntime for NativePlanRuntime<'_> { + type Output = NativePlanOutput; + type Error = QueryExecutionError; + + fn execute_node( + &self, + _id: NodeId, + node: &QueryPlanNode, + 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 { + [NativePlanOutput::Read(data)] => Ok(NativePlanOutput::Composed( + self.engine.compose_range_read(data), + )), + _ => Err(QueryExecutionError::Native( + "ComposeWindows expected store data".into(), + )), + }, + QueryPlanNode::ResolveKeys { keys, .. } => match (inputs, keys) { + ([NativePlanOutput::Composed(values)], None) => { + Ok(NativePlanOutput::Resolved(ResolvedRangeReads { + values: values.clone(), + keys: None, + })) + } + ( + [NativePlanOutput::Composed(values), NativePlanOutput::Composed(keys)], + Some(_), + ) => Ok(NativePlanOutput::Resolved(ResolvedRangeReads { + values: values.clone(), + keys: Some(keys.clone()), + })), + _ => Err(QueryExecutionError::Native( + "ResolveKeys received incompatible inputs".into(), + )), + }, + QueryPlanNode::Estimate { .. } => match inputs { + [NativePlanOutput::Resolved(reads)] => self + .engine + .estimate_range_query(self.context, reads.clone()) + .map(NativePlanOutput::Results), + _ => Err(QueryExecutionError::Native( + "Estimate expected resolved reads".into(), + )), + }, + QueryPlanNode::LimitTopK { + k, grouping_labels, .. + } => match inputs { + [NativePlanOutput::Results(results)] => self + .engine + .limit_range_topk( + results, + k, + &SimpleEngine::topk_row_label_order( + &self.context.base.metadata, + &self.context.base.grouping_labels, + &self.context.base.aggregated_labels, + ), + grouping_labels, + ) + .map_err(QueryExecutionError::Native) + .map(NativePlanOutput::Results), + _ => Err(QueryExecutionError::Native( + "LimitTopK expected estimates".into(), + )), + }, + QueryPlanNode::Format { + include_metric_name, + metric, + .. + } => match inputs { + [NativePlanOutput::Results(results)] => Ok(NativePlanOutput::Results( + self.engine + .format_range_results(results, *include_metric_name, metric), + )), + _ => Err(QueryExecutionError::Native( + "result node expected estimates".into(), + )), + }, + } + } +} + // /// Parsed components of a sketch query, extracted either via the PromQL AST // /// parser (for standard functions) or via regex (for custom functions like // /// `entropy_over_time` that the promql-parser crate doesn't recognize). @@ -195,6 +361,8 @@ pub struct SimpleEngine { data_ingestion_interval_ms: u64, controller_patterns: Vec, query_language: QueryLanguage, + #[cfg(feature = "native_query_legacy_test_support")] + native_range_execution_mode: NativeRangeExecutionMode, } impl SimpleEngine { @@ -368,9 +536,20 @@ impl SimpleEngine { data_ingestion_interval_ms, controller_patterns, query_language, + #[cfg(feature = "native_query_legacy_test_support")] + native_range_execution_mode: NativeRangeExecutionMode::Dag, } } + #[cfg(feature = "native_query_legacy_test_support")] + pub fn with_native_range_execution_mode_for_test( + mut self, + mode: NativeRangeExecutionMode, + ) -> Self { + self.native_range_execution_mode = mode; + self + } + /// Replace the inference config at runtime. Called by the applier task after /// the planner fires. /// @@ -1322,6 +1501,86 @@ impl SimpleEngine { key.labels.insert(0, metric.to_string()); } + fn limit_range_topk( + &self, + results: &[crate::engines::query_result::RangeVectorElement], + k: &str, + row_label_order: &KeyByLabelNames, + grouping_labels: &KeyByLabelNames, + ) -> Result, String> { + use crate::engines::query_result::RangeVectorElement; + + let k = Self::parse_topk_limit(&HashMap::from([("k".to_string(), k.to_string())]))?; + let mut retained: HashMap> = HashMap::new(); + let grouping_positions: Vec = grouping_labels + .labels + .iter() + .map(|label| { + row_label_order + .labels + .iter() + .position(|candidate| candidate == label) + .ok_or_else(|| format!("Topk grouping label '{label}' is absent from output")) + }) + .collect::>()?; + let mut candidates: TopkCandidates<'_> = HashMap::new(); + for result in results { + let grouping_key = grouping_positions + .iter() + .map(|&position| { + result.labels.labels.get(position).cloned().ok_or_else(|| { + format!( + "Topk result has {} labels but grouping position {} was requested", + result.labels.labels.len(), + position + ) + }) + }) + .collect::, _>>()?; + for sample in &result.samples { + candidates + .entry((sample.timestamp, grouping_key.clone())) + .or_default() + .push((&result.labels, sample.value)); + } + } + for ((timestamp, _), mut candidates) in candidates { + candidates + .sort_by(|a, b| Self::cmp_topk_value_desc(a.1, &a.0.labels, b.1, &b.0.labels)); + for (labels, _) in candidates.into_iter().take(k) { + retained.entry(labels.clone()).or_default().push(timestamp); + } + } + Ok(results + .iter() + .filter_map(|result| { + let timestamps = retained.get(&result.labels)?; + let mut limited = RangeVectorElement::new(result.labels.clone()); + for sample in &result.samples { + if timestamps.contains(&sample.timestamp) { + limited.add_sample(sample.timestamp, sample.value); + } + } + (!limited.samples.is_empty()).then_some(limited) + }) + .collect()) + } + + fn format_range_results( + &self, + results: &[crate::engines::query_result::RangeVectorElement], + include_metric_name: bool, + metric: &str, + ) -> Vec { + let mut results = results.to_vec(); + if include_metric_name { + for result in &mut results { + Self::prepend_metric_name(metric, &mut result.labels); + } + } + results + } + /// Executes the complete query pipeline: plan, execute, collect, and format. /// /// The two top-k flags are deliberately separate because the two engines @@ -1970,16 +2229,26 @@ impl SimpleEngine { /// be merged, not just the last one collected here. Used identically by /// `execute_range_query_pipeline` for both the value side and (#583) /// the keys side. - fn build_bucket_map( - buckets: &[crate::stores::TimestampedBucket], - ) -> HashMap> { - let mut bucket_map: HashMap> = HashMap::new(); + fn build_bucket_map(buckets: &[crate::stores::TimestampedBucket]) -> BucketMap { + let mut bucket_map: BucketMap = HashMap::new(); for ((start, _), bucket) in buckets { - bucket_map.entry(*start).or_default().push(bucket.as_ref()); + bucket_map + .entry(*start) + .or_default() + .push(Arc::clone(bucket)); } bucket_map } + fn compose_range_read(&self, data: &TimestampedBucketsMap) -> ComposedRangeRead { + ComposedRangeRead { + groups: data + .iter() + .map(|(key, buckets)| (key.clone(), Self::build_bucket_map(buckets))) + .collect(), + } + } + /// Collects every bucket in `bucket_map` whose start falls in /// `[window_start, window_end)`, stepping by `step_increment`. Missing /// buckets at a given start are skipped (partial data is okay). Used @@ -1996,7 +2265,7 @@ impl SimpleEngine { /// path MUST use `collect_bucket_map_entries_before` instead, never this /// (#581 stage E.4 review; see that function's doc for why). fn sum_window( - bucket_map: &HashMap>, + bucket_map: &BucketMap, window_start: u64, window_end: u64, step_increment: u64, @@ -2068,10 +2337,10 @@ impl SimpleEngine { /// inherited from the store's own sort, same as it always was for /// `sum_window` (#581 stage E.4 review). fn collect_bucket_map_entries_before( - bucket_map: &HashMap>, + bucket_map: &BucketMap, before: u64, ) -> Vec> { - let mut entries: Vec<(u64, &&dyn AggregateCore)> = bucket_map + let mut entries: Vec<(u64, &Arc)> = bucket_map .iter() .filter(|(&t, _)| t < before) .flat_map(|(&t, buckets)| buckets.iter().map(move |b| (t, b))) @@ -2093,7 +2362,7 @@ impl SimpleEngine { /// that path must call `collect_bucket_map_entries_before` directly /// instead (see its doc comment). fn window_buckets_for_step( - bucket_map: &HashMap>, + bucket_map: &BucketMap, window_start: u64, window_end: u64, step_increment: u64, @@ -2127,6 +2396,35 @@ impl SimpleEngine { enable_topk_limiting: bool, enable_topk_formatting: bool, ) -> Result, QueryExecutionError> { + Self::reject_off_grid_sliding_counter_query(context)?; + #[cfg(feature = "native_query_legacy_test_support")] + if matches!( + self.native_range_execution_mode, + NativeRangeExecutionMode::Legacy + ) { + return self.execute_legacy_range_query_pipeline( + context, + enable_topk_limiting, + enable_topk_formatting, + ); + } + #[cfg(feature = "native_query_legacy_test_support")] + let plan = if matches!( + self.native_range_execution_mode, + NativeRangeExecutionMode::MalformedPlan + ) { + QueryPlan::malformed_for_test() + } else { + QueryPlan::compile_range( + context, + PlanOptions { + limit_topk: enable_topk_limiting, + format_output: enable_topk_formatting, + }, + ) + .map_err(QueryExecutionError::Native)? + }; + #[cfg(not(feature = "native_query_legacy_test_support"))] let plan = QueryPlan::compile_range( context, PlanOptions { @@ -2136,44 +2434,72 @@ impl SimpleEngine { ) .map_err(QueryExecutionError::Native)?; debug!(plan = %plan.explain(), "Compiled native query plan"); - self.execute_range_query_pipeline(context, enable_topk_limiting, enable_topk_formatting) + let runtime = NativePlanRuntime { + engine: self, + context, + reads: std::cell::RefCell::new(None), + }; + match plan.execute(&runtime).map_err(|error| match error { + QueryPlanExecutionError::InvalidPlan(reason) => QueryExecutionError::Native(reason), + QueryPlanExecutionError::Node { source, .. } => source, + })? { + NativePlanOutput::Results(results) => Ok(results), + _ => Err(QueryExecutionError::Native( + "Query plan root did not produce results".to_string(), + )), + } } - fn execute_range_query_pipeline( + #[cfg(feature = "native_query_legacy_test_support")] + fn execute_legacy_range_query_pipeline( &self, context: &RangeQueryExecutionContext, enable_topk_limiting: bool, enable_topk_formatting: bool, ) -> Result, QueryExecutionError> { - use crate::engines::query_result::RangeVectorElement; - use crate::engines::window_merger::create_window_merger; - - if context.window_type == WindowType::Sliding - && context.tumbling_window_ms > 0 - && matches!( - context.base.metadata.statistic_to_compute, - Statistic::Increase | Statistic::Rate - ) - { - if let Some(&off_grid_timestamp) = context - .output_timestamps - .iter() - .find(|&×tamp| !timestamp.is_multiple_of(context.tumbling_window_ms)) - { - return Err(QueryExecutionError::NoLocalData(format!( - "Exact Prometheus counter bounds are unavailable for off-grid Sliding \ - timestamp {} (grid interval {}ms)", - off_grid_timestamp, context.tumbling_window_ms - ))); - } + let reads = self.read_range_query_inputs(context)?; + let mut results = self.estimate_range_query( + context, + ResolvedRangeReads { + values: self.compose_range_read(&reads.values), + keys: reads + .keys + .as_ref() + .map(|keys| self.compose_range_read(keys)), + }, + )?; + if enable_topk_limiting && context.base.metadata.statistic_to_compute == Statistic::Topk { + let k = context.base.metadata.query_kwargs.get("k").ok_or_else(|| { + QueryExecutionError::Native("Topk query is missing required `k` parameter".into()) + })?; + results = self + .limit_range_topk( + &results, + k, + &Self::topk_row_label_order( + &context.base.metadata, + &context.base.grouping_labels, + &context.base.aggregated_labels, + ), + &context.base.grouping_labels, + ) + .map_err(QueryExecutionError::Native)?; } + Ok(self.format_range_results( + &results, + enable_topk_formatting + && context.base.metadata.statistic_to_compute == Statistic::Topk + && context.base.metadata.keep_metric_name, + &context.base.metric, + )) + } + fn read_range_query_inputs( + &self, + context: &RangeQueryExecutionContext, + ) -> Result { let lookback_ms = (context.lookback_bucket_count as u64) * context.tumbling_window_ms; - - // Step 1: Fetch all data needed for the entire range. Sliding - // aggregates are already full, overlapping windows in the store, so - // request only the W-spaced exact cover for each output timestamp. - let all_data = if context.window_type == WindowType::Sliding { + let values = if context.window_type == WindowType::Sliding { self.execute_sliding_cover_query( &context.base.store_plan.values_query, &context.output_timestamps, @@ -2185,57 +2511,88 @@ impl SimpleEngine { self.execute_store_query(&context.base.store_plan.values_query) .map_err(QueryExecutionError::Native)? }; - - if all_data.is_empty() { + if values.is_empty() { return Err(QueryExecutionError::NoLocalData(format!( "No data found for metric: {}", context.base.metric ))); } - debug!( "Range query: fetched {} keys, {} total buckets", - all_data.len(), - all_data.values().map(|v| v.len()).sum::() + values.len(), + values.values().map(|buckets| buckets.len()).sum::() ); - // #583: fetch keys raw (no merge). Unlike keys, values have always - // been fetched raw here and merged per-step below (see the loop); - // keys used to go through fetch_and_merge_keys, which collapses - // every fetched bucket into ONE snapshot before this function ever - // sees it. That collapse is the bug: once buckets are merged - // together there's no way to ask what the key set looked like at - // any specific earlier timestamp. Fetching raw and merging per-step, - // mirroring the values loop, is the fix. - let keys_raw_data: Option = match &context.base.store_plan.keys_query - { - Some(keys_query) if context.keys_window_type == Some(WindowType::Sliding) => { + let keys = match &context.base.store_plan.keys_query { + Some(query) if context.keys_window_type == Some(WindowType::Sliding) => { Some(self.execute_sliding_cover_query( - keys_query, + query, &context.output_timestamps, context.keys_lookback_ms.ok_or_else(|| { QueryExecutionError::Native( - "Sliding keys query is missing its lookback".to_string(), + "Sliding keys query is missing its lookback".into(), ) })?, context.keys_window_size_ms.ok_or_else(|| { QueryExecutionError::Native( - "Sliding keys query is missing its window size".to_string(), + "Sliding keys query is missing its window size".into(), ) })?, context.keys_tumbling_window_ms.ok_or_else(|| { QueryExecutionError::Native( - "Sliding keys query is missing its slide interval".to_string(), + "Sliding keys query is missing its slide interval".into(), ) })?, )?) } - Some(keys_query) => Some( - self.execute_store_query(keys_query) + Some(query) => Some( + self.execute_store_query(query) .map_err(QueryExecutionError::Native)?, ), None => None, }; + Ok(RangeQueryReads { values, keys }) + } + + fn reject_off_grid_sliding_counter_query( + context: &RangeQueryExecutionContext, + ) -> Result<(), QueryExecutionError> { + if context.window_type != WindowType::Sliding + || context.tumbling_window_ms == 0 + || !matches!( + context.base.metadata.statistic_to_compute, + Statistic::Increase | Statistic::Rate + ) + { + return Ok(()); + } + if let Some(×tamp) = context + .output_timestamps + .iter() + .find(|&×tamp| !timestamp.is_multiple_of(context.tumbling_window_ms)) + { + return Err(QueryExecutionError::NoLocalData(format!( + "Exact Prometheus counter bounds are unavailable for off-grid Sliding \ + timestamp {} (grid interval {}ms)", + timestamp, context.tumbling_window_ms + ))); + } + Ok(()) + } + + fn estimate_range_query( + &self, + context: &RangeQueryExecutionContext, + reads: ResolvedRangeReads, + ) -> Result, QueryExecutionError> { + use crate::engines::query_result::RangeVectorElement; + use crate::engines::window_merger::create_window_merger; + + let ResolvedRangeReads { + values: ComposedRangeRead { groups: all_data }, + keys: keys_raw_data, + } = reads; + let lookback_ms = (context.lookback_bucket_count as u64) * context.tumbling_window_ms; let mut results: HashMap = HashMap::new(); @@ -2308,12 +2665,12 @@ impl SimpleEngine { // 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). - type GroupBucketMap<'a> = HashMap>; + type GroupBucketMap = BucketMap; - enum KeysSource<'a> { + enum KeysSource { Fixed(Option), PerStep { - bucket_map: GroupBucketMap<'a>, + bucket_map: GroupBucketMap, lookback_ms: u64, tumbling_window_ms: u64, stored_window_size_ms: u64, @@ -2328,8 +2685,7 @@ 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<(&Vec, KeysSource)> = match &keys_raw_data - { + 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 @@ -2345,13 +2701,14 @@ impl SimpleEngine { 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, raw_keys_buckets)| match all_data.get(group_key) { - Some(timestamped_buckets) => Some(( - timestamped_buckets, + |(group_key, key_bucket_map)| match all_data.get(group_key) { + Some(value_bucket_map) => Some(( + value_bucket_map.clone(), KeysSource::PerStep { - bucket_map: Self::build_bucket_map(raw_keys_buckets), + 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, @@ -2379,7 +2736,9 @@ impl SimpleEngine { // this list. None => all_data .iter() - .map(|(group_key, buckets)| (buckets, KeysSource::Fixed(group_key.clone()))) + .map(|(group_key, bucket_map)| { + (bucket_map.clone(), KeysSource::Fixed(group_key.clone())) + }) .collect(), }; @@ -2400,37 +2759,21 @@ impl SimpleEngine { // 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. - let groups: Vec<(GroupBucketMap, KeysSource)> = groups - .into_iter() - .map(|(timestamped_buckets, keys_source)| { - let bucket_map = Self::build_bucket_map(timestamped_buckets); - debug!( - "Group with {} start-timestamps ({} keys start-timestamps)", - bucket_map.len(), - match &keys_source { - KeysSource::PerStep { bucket_map, .. } => bucket_map.len(), - KeysSource::Fixed(_) => 0, - } - ); - (bucket_map, keys_source) - }) - .collect(); + for (bucket_map, keys_source) in &groups { + debug!( + "Group with {} start-timestamps ({} keys start-timestamps)", + bucket_map.len(), + match keys_source { + KeysSource::PerStep { bucket_map, .. } => bucket_map.len(), + KeysSource::Fixed(_) => 0, + } + ); + } // Top-k's k, parsed once rather than per timestamp. Some only when // 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 topk_k: Option = if enable_topk_limiting - && context.base.metadata.statistic_to_compute == Statistic::Topk - { - Some( - Self::parse_topk_limit(&context.base.metadata.query_kwargs) - .map_err(QueryExecutionError::Native)?, - ) - } else { - None - }; - let row_label_order = Self::topk_row_label_order( &context.base.metadata, &context.base.grouping_labels, @@ -2653,7 +2996,7 @@ impl SimpleEngine { // single-population groups let the value accumulator's own // get_keys() take priority once merged, falling back to // fallback_key otherwise. Same resolver the instant path uses. - let mut group_results: Vec<(KeyByLabelValues, f64)> = self + let group_results: Vec<(KeyByLabelValues, f64)> = self .resolve_and_query_group( Some(merged.as_ref()), keys_precompute.as_deref(), @@ -2678,13 +3021,6 @@ impl SimpleEngine { // against each other, not against other jobs' candidates. // Tie-broken by label for determinism (HashMap iteration // order isn't stable across runs). - if let Some(k) = topk_k { - group_results.sort_by(|a, b| { - Self::cmp_topk_value_desc(a.1, &a.0.labels, b.1, &b.0.labels) - }); - group_results.truncate(k); - } - step_results.extend(group_results); } @@ -2701,15 +3037,6 @@ impl SimpleEngine { // not once per timestep -- a separate pass over the final results, // after every timestamp's ranking above has already decided which // groups/samples survive. - if enable_topk_formatting - && context.base.metadata.statistic_to_compute == Statistic::Topk - && context.base.metadata.keep_metric_name - { - for elem in results.values_mut() { - Self::prepend_metric_name(&context.base.metric, &mut elem.labels); - } - } - Ok(results.into_values().collect()) } } diff --git a/asap-query-engine/src/lib.rs b/asap-query-engine/src/lib.rs index b86d8148..a458be54 100644 --- a/asap-query-engine/src/lib.rs +++ b/asap-query-engine/src/lib.rs @@ -33,6 +33,8 @@ pub use precompute_operators::{ pub use stores::{SimpleMapStore, Store, StoreResult}; +#[cfg(feature = "native_query_legacy_test_support")] +pub use engines::NativeRangeExecutionMode; pub use engines::{InstantVector, QueryExecutionError, QueryResult, SimpleEngine}; pub use drivers::{HttpServer, HttpServerConfig, OtlpReceiver, OtlpReceiverConfig}; diff --git a/asap-query-engine/src/planner_client.rs b/asap-query-engine/src/planner_client.rs index 5d3f3e7c..f7496948 100644 --- a/asap-query-engine/src/planner_client.rs +++ b/asap-query-engine/src/planner_client.rs @@ -14,6 +14,7 @@ pub struct PlannerResult { pub punted_queries: Vec, } +#[allow(clippy::double_must_use)] #[async_trait::async_trait] pub trait PlannerClient: Send + Sync { async fn plan(&self, config: ControllerConfig) -> Result; diff --git a/asap-query-engine/src/precompute_engine/ingest_source.rs b/asap-query-engine/src/precompute_engine/ingest_source.rs index 3b4670fa..4e2977ff 100644 --- a/asap-query-engine/src/precompute_engine/ingest_source.rs +++ b/asap-query-engine/src/precompute_engine/ingest_source.rs @@ -130,6 +130,7 @@ fn compile_spatial_filter(config: &AggregationConfig) -> Result, St /// /// Implementors decode incoming data (HTTP, file, etc.) and push it /// into the engine via [`route_decoded_samples`]. +#[allow(clippy::double_must_use)] #[async_trait::async_trait] pub trait IngestSource: Send + Sync { async fn run( 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 28d6734a..c900a7e9 100644 --- a/asap-query-engine/src/tests/native_range_query_tests.rs +++ b/asap-query-engine/src/tests/native_range_query_tests.rs @@ -37,6 +37,8 @@ mod tests { use crate::stores::Store; use crate::tests::test_utilities::engine_factories::create_engine_multi_timestamp_with_window; use crate::AggregateCore; + #[cfg(feature = "native_query_legacy_test_support")] + use crate::NativeRangeExecutionMode; use promql_utilities::data_model::KeyByLabelNames; use std::collections::HashMap; use std::sync::Arc; @@ -224,6 +226,20 @@ mod tests { .expect("host-a result missing") .samples; assert_eq!(on_grid_samples.len(), 2); + + #[cfg(feature = "native_query_legacy_test_support")] + assert!( + engine + .with_native_range_execution_mode_for_test(NativeRangeExecutionMode::FailingStore) + .handle_range_query_promql( + "rate(http_requests_total[2s])".to_string(), + 10.5, + 12.5, + 1.0, + ) + .expect("off-grid counter query should not read the store") + .is_none() + ); } /// Checks every `(label_values, ts, expected_present, reason)` case @@ -284,6 +300,61 @@ mod tests { ) } + fn create_oscillating_delta_set_engine() -> SimpleEngine { + let value_data: TimeSeriesData = (1..=5) + .map(|i| { + ( + i * 1000, + None, + Box::new(CountMinSketchAccumulator::new(2, 3)) as Box, + ) + }) + .collect(); + + let mut keys_add = DeltaSetAggregatorAccumulator::new(); + keys_add.add_key(KeyByLabelValues { + labels: vec!["host-a".to_string(), "evt-1".to_string()], + }); + let mut keys_remove = DeltaSetAggregatorAccumulator::new(); + keys_remove.remove_key(KeyByLabelValues { + labels: vec!["host-a".to_string(), "evt-1".to_string()], + }); + let keys_data: TimeSeriesData = vec![ + ( + 1000, + None, + Box::new(keys_add.clone()) as Box, + ), + ( + 2000, + None, + Box::new(keys_remove.clone()) as Box, + ), + ( + 3000, + None, + Box::new(keys_add.clone()) as Box, + ), + ( + 4000, + None, + Box::new(keys_remove.clone()) as Box, + ), + (5000, None, Box::new(keys_add) as Box), + ]; + + create_range_engine_dual_input( + "event_frequency", + AggregationType::CountMinSketch, + AggregationType::DeltaSetAggregator, + vec![], + vec!["host", "event"], + value_data, + keys_data, + "count(event_frequency) by (host, event)", + ) + } + /// Same as `create_range_engine_dual_input`, but lets the value and key /// aggregations use different bucket widths — the value aggregation's /// `tumbling_window_ms` must not be assumed to also be the key @@ -1437,58 +1508,7 @@ mod tests { // delta is an add) and reuses it for every step — so it would wrongly // show host-a present at every step, including the two "removed" // windows (t=2000, t=4000). - let value_data: TimeSeriesData = (1..=5) - .map(|i| { - ( - i * 1000, - None, - Box::new(CountMinSketchAccumulator::new(2, 3)) as Box, - ) - }) - .collect(); - - let mut keys_add = DeltaSetAggregatorAccumulator::new(); - keys_add.add_key(KeyByLabelValues { - labels: vec!["host-a".to_string(), "evt-1".to_string()], - }); - let mut keys_remove = DeltaSetAggregatorAccumulator::new(); - keys_remove.remove_key(KeyByLabelValues { - labels: vec!["host-a".to_string(), "evt-1".to_string()], - }); - let keys_data: TimeSeriesData = vec![ - ( - 1000, - None, - Box::new(keys_add.clone()) as Box, - ), - ( - 2000, - None, - Box::new(keys_remove.clone()) as Box, - ), - ( - 3000, - None, - Box::new(keys_add.clone()) as Box, - ), - ( - 4000, - None, - Box::new(keys_remove.clone()) as Box, - ), - (5000, None, Box::new(keys_add) as Box), - ]; - - let engine = create_range_engine_dual_input( - "event_frequency", - AggregationType::CountMinSketch, - AggregationType::DeltaSetAggregator, - vec![], - vec!["host", "event"], - value_data, - keys_data, - "count(event_frequency) by (host, event)", - ); + let engine = create_oscillating_delta_set_engine(); let query = "count(event_frequency) by (host, event)"; let result = engine @@ -1518,6 +1538,25 @@ mod tests { ); } + #[cfg(feature = "native_query_legacy_test_support")] + #[test] + fn range_query_delta_set_replay_dag_matches_legacy() { + let query = "count(event_frequency) by (host, event)"; + let dag = create_oscillating_delta_set_engine() + .handle_range_query_promql(query.to_string(), 1.0, 5.0, 1.0) + .expect("DAG execution failed"); + let legacy = create_oscillating_delta_set_engine() + .with_native_range_execution_mode_for_test(NativeRangeExecutionMode::Legacy) + .handle_range_query_promql(query.to_string(), 1.0, 5.0, 1.0) + .expect("legacy execution failed"); + + assert_eq!( + serde_json::to_value(dag).expect("DAG result should serialize"), + serde_json::to_value(legacy).expect("legacy result should serialize"), + "DAG execution must preserve DeltaSet replay semantics" + ); + } + #[tokio::test(flavor = "multi_thread")] async fn range_query_binary_expr_arm_set_aggregator_earlier_key_not_silently_dropped() { // SetAggregator counterpart to diff --git a/asap-query-engine/src/tests/prometheus_forwarding_tests.rs b/asap-query-engine/src/tests/prometheus_forwarding_tests.rs index 03a7b4c7..605bf2d3 100644 --- a/asap-query-engine/src/tests/prometheus_forwarding_tests.rs +++ b/asap-query-engine/src/tests/prometheus_forwarding_tests.rs @@ -1,11 +1,21 @@ #[cfg(test)] -use crate::data_model::{CleanupPolicy, InferenceConfig, QueryLanguage, StreamingConfig}; +use crate::data_model::{ + AggregationConfig, AggregationReference, CleanupPolicy, InferenceConfig, PromQLSchema, + QueryConfig, QueryLanguage, SchemaConfig, StreamingConfig, WindowType, +}; use crate::drivers::query::adapters::AdapterConfig; use crate::drivers::query::servers::http::{HttpServer, HttpServerConfig}; use crate::engines::SimpleEngine; use crate::stores::simple_map_store::SimpleMapStore; +#[cfg(feature = "native_query_legacy_test_support")] +use crate::NativeRangeExecutionMode; +use promql_utilities::data_model::KeyByLabelNames; +use promql_utilities::query_logics::enums::AggregationType; use reqwest::Client; use serde_json::Value; +use std::collections::HashMap; +#[cfg(feature = "native_query_legacy_test_support")] +use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; use tokio::net::TcpListener; use tokio::time::{sleep, Duration}; @@ -100,6 +110,35 @@ async fn start_mock_prometheus_server() -> Result (u16, Arc) { + use axum::{response::Json, routing::get, Router}; + use serde_json::json; + + let requests = Arc::new(AtomicUsize::new(0)); + let request_counter = requests.clone(); + let app = Router::new().route( + "/api/v1/query_range", + get(move || { + let request_counter = request_counter.clone(); + async move { + request_counter.fetch_add(1, Ordering::SeqCst); + Json(json!({"status": "success", "data": {"resultType": "matrix", "result": []}})) + } + }), + ); + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("mock fallback should bind"); + let port = listener.local_addr().expect("mock fallback address").port(); + tokio::spawn(async move { + axum::serve(listener, app) + .await + .expect("mock fallback should run"); + }); + (port, requests) +} + async fn setup_test_server(prometheus_port: u16) -> (HttpServer, u16) { let config = HttpServerConfig { port: 0, // Use random port @@ -135,6 +174,80 @@ async fn setup_test_server(prometheus_port: u16) -> (HttpServer, u16) { (server, actual_port) } +#[cfg(feature = "native_query_legacy_test_support")] +fn native_range_error_configs() -> (InferenceConfig, Arc) { + let metric = "native_metric"; + let streaming_config = Arc::new(StreamingConfig::new(HashMap::from([( + 1, + AggregationConfig { + aggregation_id: 1, + aggregation_type: AggregationType::Sum, + aggregation_sub_type: String::new(), + parameters: HashMap::new(), + grouping_labels: KeyByLabelNames::empty(), + aggregated_labels: KeyByLabelNames::empty(), + rollup_labels: KeyByLabelNames::empty(), + original_yaml: String::new(), + window_size_ms: 1_000, + slide_interval_ms: 1_000, + window_type: WindowType::Tumbling, + spatial_filter: String::new(), + spatial_filter_normalized: String::new(), + metric: metric.to_string(), + num_aggregates_to_retain: None, + read_count_threshold: None, + table_name: None, + value_column: None, + }, + )]))); + let inference_config = InferenceConfig { + schema: SchemaConfig::PromQL( + PromQLSchema::new().add_metric(metric.to_string(), KeyByLabelNames::empty()), + ), + query_configs: vec![QueryConfig::new("sum(native_metric)".to_string()) + .add_aggregation(AggregationReference::new(1, None))], + cleanup_policy: CleanupPolicy::NoCleanup, + }; + (inference_config, streaming_config) +} + +#[cfg(feature = "native_query_legacy_test_support")] +async fn setup_test_server_with_native_range_mode( + prometheus_port: u16, + mode: NativeRangeExecutionMode, +) -> (HttpServer, u16) { + let config = HttpServerConfig { + port: 0, + handle_http_requests: true, + adapter_config: AdapterConfig::prometheus_promql( + format!("http://127.0.0.1:{prometheus_port}"), + true, + 30, + ), + }; + let (inference_config, streaming_config) = native_range_error_configs(); + let store = Arc::new(SimpleMapStore::new( + streaming_config.clone(), + CleanupPolicy::NoCleanup, + )); + let query_engine = Arc::new( + SimpleEngine::new( + store.clone(), + inference_config, + streaming_config, + 15000, + QueryLanguage::promql, + ) + .with_native_range_execution_mode_for_test(mode), + ); + let server = HttpServer::new(config, query_engine, store, None); + let port = server + .start_test_server() + .await + .expect("test query server should start"); + (server, port) +} + #[tokio::test] async fn test_prometheus_forwarding_instant_query() { // Start mock Prometheus server @@ -340,6 +453,42 @@ async fn test_prometheus_forwarding_range_query() { assert_eq!(result["values"][1][1], "43.0"); } +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn native_range_execution_error_is_local_and_does_not_fallback() { + for mode in [ + NativeRangeExecutionMode::MalformedPlan, + NativeRangeExecutionMode::FailingStore, + ] { + let (prometheus_port, fallback_requests) = start_counting_mock_prometheus_server().await; + let (_server, server_port) = + setup_test_server_with_native_range_mode(prometheus_port, mode).await; + + let response = Client::new() + .get(format!("http://127.0.0.1:{server_port}/api/v1/query_range")) + .query(&[ + ("query", "sum(native_metric)"), + ("start", "60"), + ("end", "61"), + ("step", "1"), + ]) + .send() + .await + .expect("range request should complete"); + let body: Value = response + .json() + .await + .expect("error response should be JSON"); + + assert_eq!(body["status"], "error"); + assert_eq!( + fallback_requests.load(Ordering::SeqCst), + 0, + "native execution errors must not be forwarded to Prometheus" + ); + } +} + #[tokio::test] async fn test_range_query_forwarding_disabled() { let config = HttpServerConfig { diff --git a/asap-query-engine/tests/e2e_precompute_equivalence.rs b/asap-query-engine/tests/e2e_precompute_equivalence.rs index fd671b15..6a20ed3c 100644 --- a/asap-query-engine/tests/e2e_precompute_equivalence.rs +++ b/asap-query-engine/tests/e2e_precompute_equivalence.rs @@ -7,8 +7,10 @@ //! 3. Advances the watermark past the window boundary to close it //! 4. Drains captured outputs and queries them +use asap_planner::{Controller, RuntimeOptions, StreamingEngine}; use asap_types::aggregation_config::AggregationConfig; use asap_types::enums::{AggregationType, CleanupPolicy, QueryLanguage, WindowType}; +use promql_utilities::data_model::KeyByLabelNames; use prost::Message; use serde_json::json; use std::collections::HashMap; @@ -23,6 +25,8 @@ use query_engine_rust::drivers::ingest::prometheus_remote_write::{ use query_engine_rust::precompute_engine::config::{LateDataPolicy, PrecomputeEngineConfig}; use query_engine_rust::precompute_engine::output_sink::CapturingOutputSink; use query_engine_rust::precompute_engine::{HttpIngestConfig, HttpIngestSource, PrecomputeEngine}; +#[cfg(feature = "native_query_legacy_test_support")] +use query_engine_rust::NativeRangeExecutionMode; use query_engine_rust::{QueryResult, SimpleEngine, SimpleMapStore, Store}; // ─── helpers ──────────────────────────────────────────────────────────────── @@ -155,7 +159,95 @@ fn engine_config() -> PrecomputeEngineConfig { } } -struct PromqlPrecomputeFixture<'a> { +fn plan_promql_query( + metric: &str, + labels: Vec, + query: &str, + interval_ms: u64, +) -> (Arc, InferenceConfig) { + let controller_config = format!( + r#" +query_groups: + - id: 1 + queries: + - "{query}" + repetition_delay_ms: {interval_ms} + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +"# + ); + let planner = Controller::from_yaml_with_schema( + &controller_config, + PromQLSchema::new().add_metric(metric.to_string(), KeyByLabelNames::new(labels)), + RuntimeOptions { + data_ingestion_interval_ms: interval_ms, + streaming_engine: StreamingEngine::Precompute, + enable_punting: false, + range_duration_ms: interval_ms, + step_ms: interval_ms, + }, + ) + .expect("planner configuration should be valid"); + let output = planner.generate().expect("planner should support query"); + let inference_config = output + .to_inference_config(QueryLanguage::promql) + .expect("planner should produce inference config"); + let streaming_config = output + .to_streaming_config(QueryLanguage::promql) + .expect("planner should produce streaming config"); + (Arc::new(streaming_config), inference_config) +} + +async fn build_engine_from_configs( + port: u16, + streaming_config: Arc, + inference_config: InferenceConfig, + samples: Vec, + base_interval_ms: u64, +) -> SimpleEngine { + let sink = Arc::new(CapturingOutputSink::new()); + let engine = PrecomputeEngine::new( + engine_config(), + streaming_config.clone(), + sink.clone(), + vec![Box::new(HttpIngestSource::new(HttpIngestConfig { port }))], + ); + tokio::spawn(async move { + engine + .run() + .await + .expect("precompute engine should keep running"); + }); + tokio::time::sleep(tokio::time::Duration::from_millis(300)).await; + + let client = reqwest::Client::new(); + for sample in samples { + send_remote_write(&client, port, vec![sample]).await; + } + tokio::time::sleep(tokio::time::Duration::from_millis(600)).await; + + let store = Arc::new(SimpleMapStore::new( + streaming_config.clone(), + CleanupPolicy::NoCleanup, + )); + for (output, accumulator) in sink.drain() { + store + .insert_precomputed_output(output, accumulator) + .unwrap(); + } + + SimpleEngine::new( + store, + inference_config, + streaming_config, + base_interval_ms, + QueryLanguage::promql, + ) +} + +#[derive(Clone)] +struct NativeDagScenario<'a> { port: u16, metric: &'a str, query: &'a str, @@ -166,8 +258,8 @@ struct PromqlPrecomputeFixture<'a> { base_interval_ms: u64, } -impl PromqlPrecomputeFixture<'_> { - async fn run(self) -> QueryResult { +impl NativeDagScenario<'_> { + async fn build_engine(self) -> (SimpleEngine, String) { let aggregation_ids: Vec = self .aggregation_configs .iter() @@ -189,7 +281,10 @@ impl PromqlPrecomputeFixture<'_> { }))], ); tokio::spawn(async move { - let _ = engine.run().await; + engine + .run() + .await + .expect("precompute engine should keep running"); }); tokio::time::sleep(tokio::time::Duration::from_millis(300)).await; @@ -233,14 +328,40 @@ impl PromqlPrecomputeFixture<'_> { QueryLanguage::promql, ); + (query_engine, self.query.to_string()) + } + + async fn run(self) -> QueryResult { + let evaluation_time_seconds = self.evaluation_time_seconds; + let (query_engine, query) = self.build_engine().await; query_engine - .handle_query_promql(self.query.to_string(), self.evaluation_time_seconds) + .handle_query_promql(query.clone(), evaluation_time_seconds) .expect("native query execution should not fail") - .unwrap_or_else(|| panic!("precomputed query should succeed: {}", self.query)) + .unwrap_or_else(|| panic!("precomputed query should succeed: {query}")) .1 } } +#[cfg(feature = "native_query_legacy_test_support")] +fn assert_range_results_match( + dag: Option<(promql_utilities::data_model::KeyByLabelNames, QueryResult)>, + legacy: Option<(promql_utilities::data_model::KeyByLabelNames, QueryResult)>, +) { + let mut dag = dag.expect("DAG path should execute natively"); + let mut legacy = legacy.expect("legacy path should execute natively"); + for (_, result) in [&mut dag, &mut legacy] { + if let QueryResult::Matrix(matrix) = result { + matrix + .values + .sort_by(|left, right| left.labels.labels.cmp(&right.labels.labels)); + } + } + assert_eq!( + serde_json::to_value(Some(dag)).unwrap(), + serde_json::to_value(Some(legacy)).unwrap() + ); +} + #[tokio::test] async fn e2e_sliding_precompute_outputs_compose_a_wider_query() { let port = 19402u16; @@ -362,7 +483,7 @@ async fn e2e_promql_sum_uses_open_closed_evaluation_window() { .into_iter() .map(|(timestamp_ms, value)| make_timeseries(metric, vec![], timestamp_ms, value)) .collect(); - let result = PromqlPrecomputeFixture { + let result = NativeDagScenario { port, metric, query, @@ -382,6 +503,443 @@ async fn e2e_promql_sum_uses_open_closed_evaluation_window() { assert_eq!(vector.values[0].value, 5.0); } +/// A native leaf must give the same value at the end of a range query as an +/// instant query at that timestamp. This is the baseline that the DAG +/// executor must preserve during the cutover. +#[tokio::test] +async fn e2e_native_leaf_range_matches_instant_at_range_end() { + let port = 19408u16; + let agg_id = 8u64; + let window_size_ms = 1_000u64; + let metric = "dag_requests"; + let query = "sum(dag_requests)"; + let scenario = NativeDagScenario { + port, + metric, + query, + aggregation_configs: vec![make_agg_config( + agg_id, + metric, + AggregationType::Sum, + "", + window_size_ms, + 0, + vec![], + )], + schema_labels: vec![], + samples: vec![ + make_timeseries(metric, vec![], 1_000, 100.0), + make_timeseries(metric, vec![], 1_500, 2.0), + make_timeseries(metric, vec![], 2_000, 3.0), + make_timeseries(metric, vec![], 3_500, 0.0), + ], + evaluation_time_seconds: 2.0, + base_interval_ms: window_size_ms, + }; + + let (engine, query) = scenario.build_engine().await; + let (_, instant) = engine + .handle_query_promql(query.clone(), 2.0) + .expect("instant native query should not fail") + .expect("instant native query should succeed"); + let (_, range) = engine + .handle_range_query_promql(query, 1.0, 2.0, 1.0) + .expect("range native query should not fail") + .expect("range native query should succeed"); + + let QueryResult::Vector(instant) = instant else { + panic!("expected instant vector result"); + }; + let QueryResult::Matrix(range) = range else { + panic!("expected range vector result"); + }; + assert_eq!(instant.values.len(), 1); + assert_eq!(range.values.len(), 1); + let final_sample = range.values[0] + .samples + .last() + .expect("range result should contain the end timestamp"); + assert_eq!(final_sample.timestamp, 2_000); + assert_eq!(final_sample.value, instant.values[0].value); +} + +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn e2e_native_dag_range_matches_legacy_range() { + let scenario = NativeDagScenario { + port: 19409, + metric: "dag_differential_requests", + query: "sum_over_time(dag_differential_requests[2s])", + aggregation_configs: vec![make_agg_config( + 9, + "dag_differential_requests", + AggregationType::Sum, + "", + 1_000, + 0, + vec![], + )], + schema_labels: vec![], + samples: vec![ + make_timeseries("dag_differential_requests", vec![], 1_000, 1.0), + make_timeseries("dag_differential_requests", vec![], 2_000, 2.0), + make_timeseries("dag_differential_requests", vec![], 3_000, 3.0), + make_timeseries("dag_differential_requests", vec![], 5_000, 0.0), + ], + evaluation_time_seconds: 3.0, + base_interval_ms: 1_000, + }; + let mut legacy_scenario = scenario.clone(); + legacy_scenario.port = 19410; + let (dag, query) = scenario.build_engine().await; + let (legacy, _) = legacy_scenario.build_engine().await; + let dag = dag + .handle_range_query_promql(query.clone(), 2.0, 3.0, 1.0) + .expect("DAG execution should not fail"); + let legacy = legacy + .with_native_range_execution_mode_for_test(NativeRangeExecutionMode::Legacy) + .handle_range_query_promql(query, 2.0, 3.0, 1.0) + .expect("legacy execution should not fail"); + assert_range_results_match(dag, legacy); +} + +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn e2e_sparse_range_dag_matches_legacy_range() { + let scenario = NativeDagScenario { + port: 19411, + metric: "sparse_dag_differential", + query: "sum(sparse_dag_differential)", + aggregation_configs: vec![make_agg_config( + 10, + "sparse_dag_differential", + AggregationType::Sum, + "", + 1_000, + 0, + vec![], + )], + schema_labels: vec![], + samples: vec![ + make_timeseries("sparse_dag_differential", vec![], 1_000, 1.0), + make_timeseries("sparse_dag_differential", vec![], 2_000, 1.0), + make_timeseries("sparse_dag_differential", vec![], 8_000, 1.0), + make_timeseries("sparse_dag_differential", vec![], 9_000, 1.0), + make_timeseries("sparse_dag_differential", vec![], 12_000, 0.0), + ], + evaluation_time_seconds: 9.0, + base_interval_ms: 1_000, + }; + let mut legacy_scenario = scenario.clone(); + legacy_scenario.port = 19412; + let (dag, query) = scenario.build_engine().await; + let (legacy, _) = legacy_scenario.build_engine().await; + let dag = dag + .handle_range_query_promql(query.clone(), 1.0, 9.0, 1.0) + .expect("DAG execution should not fail"); + let legacy = legacy + .with_native_range_execution_mode_for_test(NativeRangeExecutionMode::Legacy) + .handle_range_query_promql(query, 1.0, 9.0, 1.0) + .expect("legacy execution should not fail"); + assert_range_results_match(dag, legacy); +} + +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn e2e_keyed_count_range_dag_matches_legacy_range() { + let metric = "keyed_dag_differential"; + let mut values = make_agg_config_full( + 11, + metric, + AggregationType::CountMinSketch, + "count", + 1_000, + 0, + vec![], + vec!["host"], + ); + values.parameters.insert("depth".to_string(), json!(3_u64)); + values + .parameters + .insert("width".to_string(), json!(128_u64)); + let scenario = NativeDagScenario { + port: 19413, + metric, + query: "count(keyed_dag_differential) by (host)", + aggregation_configs: vec![ + values, + make_agg_config_full( + 12, + metric, + AggregationType::SetAggregator, + "", + 1_000, + 0, + vec![], + vec!["host"], + ), + ], + schema_labels: vec!["host".to_string()], + samples: vec![ + make_timeseries(metric, vec![("host", "a")], 1_000, 1.0), + make_timeseries(metric, vec![("host", "b")], 2_000, 1.0), + make_timeseries(metric, vec![("host", "a")], 3_000, 1.0), + make_timeseries(metric, vec![], 5_000, 0.0), + ], + evaluation_time_seconds: 3.0, + base_interval_ms: 1_000, + }; + let mut legacy_scenario = scenario.clone(); + legacy_scenario.port = 19414; + let (dag, query) = scenario.build_engine().await; + let (legacy, _) = legacy_scenario.build_engine().await; + let dag = dag + .handle_range_query_promql(query.clone(), 1.0, 3.0, 1.0) + .unwrap(); + let legacy = legacy + .with_native_range_execution_mode_for_test(NativeRangeExecutionMode::Legacy) + .handle_range_query_promql(query, 1.0, 3.0, 1.0) + .unwrap(); + assert_range_results_match(dag, legacy); +} + +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn e2e_self_keyed_topk_dag_matches_legacy_range() { + let metric = "topk_dag_differential"; + let mut config = make_agg_config_full( + 13, + metric, + AggregationType::CountMinSketchWithHeap, + "count", + 1_000, + 0, + vec![], + vec!["host"], + ); + config.parameters.insert("depth".to_string(), json!(3_u64)); + config + .parameters + .insert("width".to_string(), json!(128_u64)); + config + .parameters + .insert("heapsize".to_string(), json!(16_u64)); + let scenario = NativeDagScenario { + port: 19415, + metric, + query: "topk(2, topk_dag_differential)", + aggregation_configs: vec![config], + schema_labels: vec!["host".to_string()], + samples: vec![ + make_timeseries(metric, vec![("host", "a")], 1_000, 3.0), + make_timeseries(metric, vec![("host", "b")], 1_000, 3.0), + make_timeseries(metric, vec![("host", "c")], 1_000, 3.0), + make_timeseries(metric, vec![], 3_000, 0.0), + ], + evaluation_time_seconds: 1.0, + base_interval_ms: 1_000, + }; + let mut legacy_scenario = scenario.clone(); + legacy_scenario.port = 19416; + let (dag, query) = scenario.build_engine().await; + let (legacy, _) = legacy_scenario.build_engine().await; + let dag = dag + .handle_range_query_promql(query.clone(), 1.0, 2.0, 1.0) + .unwrap(); + let legacy = legacy + .with_native_range_execution_mode_for_test(NativeRangeExecutionMode::Legacy) + .handle_range_query_promql(query, 1.0, 2.0, 1.0) + .unwrap(); + assert_range_results_match(dag, legacy); +} + +#[cfg(feature = "native_query_legacy_test_support")] +async fn assert_grouped_topk_range_dag_matches_legacy_range( + metric: &str, + query: &str, + dag_port: u16, + legacy_port: u16, +) { + let labels = vec!["job".to_string(), "instance".to_string()]; + let samples: Vec = ["frontend", "backend", "worker"] + .into_iter() + .flat_map(|job| { + (1_i64..=4).map(move |rank| { + make_timeseries( + metric, + vec![("job", job), ("instance", format!("i-{rank}").as_str())], + 1_000, + rank as f64, + ) + }) + }) + .chain(["frontend", "backend", "worker"].into_iter().map(|job| { + make_timeseries( + metric, + vec![("job", job), ("instance", "flush")], + 3_000, + 0.0, + ) + })) + .collect(); + let (dag_streaming_config, dag_inference_config) = + plan_promql_query(metric, labels.clone(), query, 1_000); + let (legacy_streaming_config, legacy_inference_config) = + plan_promql_query(metric, labels, query, 1_000); + let dag_engine = build_engine_from_configs( + dag_port, + dag_streaming_config, + dag_inference_config, + samples.clone(), + 1_000, + ) + .await; + assert!( + dag_engine + .build_range_query_execution_context_promql(query.to_string(), 1.0, 2.0, 1.0) + .is_some(), + "planner output should build a native range context" + ); + let dag = dag_engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .unwrap(); + let legacy_engine = build_engine_from_configs( + legacy_port, + legacy_streaming_config, + legacy_inference_config, + samples, + 1_000, + ) + .await; + assert!( + legacy_engine + .build_range_query_execution_context_promql(query.to_string(), 1.0, 2.0, 1.0) + .is_some(), + "planner output should build a native range context" + ); + let legacy = legacy_engine + .with_native_range_execution_mode_for_test(NativeRangeExecutionMode::Legacy) + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .unwrap(); + assert_range_results_match(dag.clone(), legacy); + let Some((_, result)) = dag else { + panic!("grouped topk should execute natively, got {dag:?}"); + }; + let row_count = match result { + QueryResult::Matrix(matrix) => matrix.values.len(), + QueryResult::Vector(vector) => vector.values.len(), + }; + // Each job has four differently frequent instances. The lowest-ranked + // instance per job is removed, leaving three rows in each of three jobs. + assert_eq!(row_count, 9); +} + +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn e2e_grouped_topk_range_dag_matches_legacy_range() { + assert_grouped_topk_range_dag_matches_legacy_range( + "grouped_topk_dag_differential", + "topk by (job) (3, grouped_topk_dag_differential)", + 19421, + 19422, + ) + .await; +} + +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn e2e_grouped_topk_sum_over_time_dag_matches_legacy_range() { + assert_grouped_topk_range_dag_matches_legacy_range( + "grouped_topk_sum_over_time_dag_differential", + "topk by (job) (3, sum_over_time(grouped_topk_sum_over_time_dag_differential[1s]))", + 19423, + 19424, + ) + .await; +} + +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn e2e_grouped_topk_count_over_time_dag_matches_legacy_range() { + assert_grouped_topk_range_dag_matches_legacy_range( + "grouped_topk_count_over_time_dag_differential", + "topk by (job) (3, count_over_time(grouped_topk_count_over_time_dag_differential[1s]))", + 19425, + 19426, + ) + .await; +} + +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn e2e_malformed_native_plan_returns_local_error() { + let metric = "dag_requests"; + let (engine, query) = NativeDagScenario { + port: 19417, + metric, + query: "sum(dag_requests)", + aggregation_configs: vec![make_agg_config( + 14, + metric, + AggregationType::Sum, + "", + 1_000, + 0, + vec![], + )], + schema_labels: vec![], + samples: vec![ + make_timeseries(metric, vec![], 1_000, 100.0), + make_timeseries(metric, vec![], 1_500, 2.0), + make_timeseries(metric, vec![], 2_000, 3.0), + make_timeseries(metric, vec![], 3_500, 0.0), + ], + evaluation_time_seconds: 2.0, + base_interval_ms: 1_000, + } + .build_engine() + .await; + assert!(engine + .with_native_range_execution_mode_for_test(NativeRangeExecutionMode::MalformedPlan) + .handle_range_query_promql(query, 1.0, 2.0, 1.0) + .is_err()); +} + +#[cfg(feature = "native_query_legacy_test_support")] +#[tokio::test] +async fn e2e_native_store_failure_returns_local_error() { + let metric = "store_failure_differential"; + let (engine, query) = NativeDagScenario { + port: 19418, + metric, + query: "sum(store_failure_differential)", + aggregation_configs: vec![make_agg_config( + 15, + metric, + AggregationType::Sum, + "", + 1_000, + 0, + vec![], + )], + schema_labels: vec![], + samples: vec![ + make_timeseries(metric, vec![], 1_000, 1.0), + make_timeseries(metric, vec![], 1_500, 2.0), + make_timeseries(metric, vec![], 2_000, 3.0), + make_timeseries(metric, vec![], 3_500, 0.0), + ], + evaluation_time_seconds: 2.0, + base_interval_ms: 1_000, + } + .build_engine() + .await; + assert!(engine + .with_native_range_execution_mode_for_test(NativeRangeExecutionMode::FailingStore) + .handle_range_query_promql(query, 1.0, 2.0, 1.0) + .is_err()); +} + /// The #698 boundary contract applies independently to a query's value and /// key precomputes. An endpoint series can only appear when both sides assign /// its sample to the window ending at the evaluation timestamp. @@ -430,7 +988,7 @@ async fn e2e_promql_count_uses_open_closed_value_and_key_windows() { .into_iter() .map(|(timestamp_ms, host)| make_timeseries(metric, vec![("host", host)], timestamp_ms, 1.0)) .collect(); - let result = PromqlPrecomputeFixture { + let result = NativeDagScenario { port, metric, query, @@ -498,7 +1056,7 @@ async fn e2e_quantile_over_time_uses_open_closed_evaluation_window() { .into_iter() .map(|(timestamp_ms, value)| make_timeseries(metric, vec![], timestamp_ms, value)) .collect(); - let result = PromqlPrecomputeFixture { + let result = NativeDagScenario { port, metric, query, @@ -518,6 +1076,65 @@ async fn e2e_quantile_over_time_uses_open_closed_evaluation_window() { assert_eq!(vector.values[0].value, 4.0); } +/// Regression: grouped quantiles retain their grouping labels. The native DAG +/// must not prepend the metric name; only PromQL topk has that output shape. +#[tokio::test] +async fn e2e_grouped_quantile_preserves_output_label_shape() { + let port = 19420u16; + let metric = "grouped_latency"; + let query = "quantile by (job) (0.99, grouped_latency)"; + let mut config = make_agg_config( + 16, + metric, + AggregationType::DatasketchesKLL, + "", + 1_000, + 0, + vec!["job"], + ); + config.parameters.insert("K".to_string(), json!(200_u64)); + let samples = [("frontend", 100.0), ("backend", 200.0)] + .into_iter() + .flat_map(|(job, value)| { + [ + make_timeseries(metric, vec![("job", job)], 1_500, value), + make_timeseries(metric, vec![("job", job)], 3_500, 0.0), + ] + }) + .collect(); + let (engine, query) = NativeDagScenario { + port, + metric, + query, + aggregation_configs: vec![config], + schema_labels: vec!["job".to_string()], + samples, + evaluation_time_seconds: 2.0, + base_interval_ms: 1_000, + } + .build_engine() + .await; + + let (output_labels, result) = engine + .handle_query_promql(query, 2.0) + .expect("grouped quantile should execute") + .expect("grouped quantile should match configured inference"); + assert_eq!(output_labels.labels, vec!["job"]); + let QueryResult::Vector(vector) = result else { + panic!("expected instant vector result"); + }; + let mut returned_labels: Vec<_> = vector + .values + .into_iter() + .map(|element| element.labels.labels) + .collect(); + returned_labels.sort(); + assert_eq!( + returned_labels, + vec![vec!["backend".to_string()], vec!["frontend".to_string()]] + ); +} + /// Sliding precomputes keep their existing exact-cover composition while /// samples on every slide boundary move to the pane ending at that boundary. /// The shared 6s boundary must be counted once, not once per stored window. @@ -550,7 +1167,7 @@ async fn e2e_sliding_query_uses_open_closed_boundaries_without_double_counting() .into_iter() .map(|(timestamp_ms, value)| make_timeseries(metric, vec![], timestamp_ms, value)) .collect(); - let result = PromqlPrecomputeFixture { + let result = NativeDagScenario { port, metric, query, diff --git a/docs/729-native-query-dag-todos.md b/docs/729-native-query-dag-todos.md new file mode 100644 index 00000000..f6949129 --- /dev/null +++ b/docs/729-native-query-dag-todos.md @@ -0,0 +1,45 @@ +# Native query DAG: remaining work + +This is the completion checklist for the native range-query DAG cutover. + +## Required before removing the temporary comparison path + +- [x] Re-run the Docker differential matrix from an isolated Compose lifecycle and record the result for every case. + Each case now has an isolated Compose project, base timestamp, and host-port trio. The 2026-10-01 reports are in `/tmp/asapquery-differential-reports`. +- [x] Re-run the Docker differential matrix against `main` and classify the result per suite. + `quantiles` passes on both revisions. The existing `request-rate` temporal case passes on both; the PR's newly added final-window `sum_over_time` cases fail and characterize a pre-existing native-query parity gap. The new off-grid rate case fails because native execution does not fall back at non-grid timestamps. `olly-bench` retains its documented non-CI planner-coverage failures. + The initial `aggregations` comparison had 14 new failures relative to `main`: + grouped `topk` over a bare selector, `sum_over_time`, and `count_over_time` + returned too few series. The pre-existing rate-based `topk` and `sum` failures + remained. The DAG-focused aggregation suite passed, so the current coverage did + not reproduce the regression; add a legacy-versus-DAG E2E case for grouped topk + before fixing it. +- [x] Record the final isolated Docker matrix after the grouped-topk fix. + The 2026-10-03 run is in `/tmp/asapquery-differential-reports-2026-10-03-final`: + `native-dag-aggregations` passes 4/4 and `quantiles` passes 35/35. + `aggregations` has 5/47 failures, all pre-existing rate-based `topk`/`sum` + cases; the 14 grouped-topk regressions are resolved. `temporal` retains 2/3 + final-window `sum_over_time` mismatches, `off-grid-rate` retains its 1/1 + non-grid `rate` mismatch, and `olly-bench` retains its documented 22/25 + planner-coverage failures. `make run-all` therefore exits nonzero by design. +- [x] Add public E2E characterization for each classified DAG regression, then fix only DAG-caused regressions. + Grouped `topk by (...)` now has legacy-versus-DAG coverage for a bare selector, + `sum_over_time`, and `count_over_time`. The limiter ranks candidates within each + timestamp/group partition, fixing the 14 grouped-topk Docker regressions. + +## Public error API + +- [x] Complete the agreed API migration: `handle_query_promql` and `handle_range_query_promql` return `Result, QueryExecutionError>` (merged separately in #749). +- [x] Keep the error taxonomy explicit: `Ok(None)` means unsupported and fallback is allowed; `Err(QueryExecutionError)` means an accepted native execution failed and fallback is forbidden. + +## Retain temporary cutover machinery while staging DAG execution + +- [x] Keep DAG as the production default. + `NativeRangeExecutionMode::Dag` is the default, and all legacy modes are compiled + only by the `native_query_legacy_test_support` feature. +- [ ] Retain `native_query_legacy_test_support`, `NativeRangeExecutionMode`, the legacy executor, and legacy-versus-DAG E2E tests until a future cutover decision. +- [x] Retain DAG-only E2E behavior tests, malformed-plan graph validation, local HTTP error/no-fallback coverage, and Docker compliance suites. + +## Follow-up scope + +- [ ] Implement complete native-query DAG execution in #743: represent arithmetic, constants, label matching, timestamp alignment, and fallback decisions in the plan rather than only executing each native query arm through a DAG. diff --git a/promql-compliance/docker-compose.yml b/promql-compliance/docker-compose.yml index 0a20620d..df3877e4 100644 --- a/promql-compliance/docker-compose.yml +++ b/promql-compliance/docker-compose.yml @@ -52,7 +52,7 @@ services: - "${ASAP_QUERY_PORT:-18088}:8088" - "${ASAP_INGEST_PORT:-19091}:9091" environment: - RUST_LOG: INFO + RUST_LOG: ${RUST_LOG:-INFO} RUST_BACKTRACE: "1" volumes: - differential-planner-output:/asap-planner-output:ro diff --git a/promql-compliance/runner/Makefile b/promql-compliance/runner/Makefile index 12ca1591..a844a202 100644 --- a/promql-compliance/runner/Makefile +++ b/promql-compliance/runner/Makefile @@ -3,9 +3,14 @@ DATASET ?= ../datasets/single-rate.yaml SUITE ?= ../suites/temporal.yaml REPORT_DIR ?= /tmp/asapquery-differential-reports +RUN_ID ?= $(shell date +%s%N) +COMPOSE_PROJECT ?= asapquery-differential-$(RUN_ID) +BASE_TIME_MS ?= $(shell now_ms=$$(date +%s%3N); echo $$((now_ms / 300000 * 300000 - 1800000))) CI_CASES := \ single-rate-temporal:../datasets/single-rate.yaml:../suites/temporal.yaml \ + single-rate-off-grid-rate:../datasets/single-rate.yaml:../suites/off-grid-rate.yaml \ + aggregations-native-dag:../datasets/aggregations.yaml:../suites/native-dag-aggregations.yaml \ aggregations:../datasets/aggregations.yaml:../suites/aggregations.yaml \ quantiles:../datasets/quantiles.yaml:../suites/quantiles.yaml @@ -20,7 +25,9 @@ run: go run ./cmd/differential-runner \ --dataset $(DATASET) \ --suite $(SUITE) \ - --compose-file ../docker-compose.yml + --compose-file ../docker-compose.yml \ + --compose-project $(COMPOSE_PROJECT) \ + --base-time-ms $(BASE_TIME_MS) run-all: @$(MAKE) --no-print-directory run-cases CASES="$(CI_CASES) $(NON_CI_CASES)" @@ -32,16 +39,34 @@ run-cases: @set -eu; \ mkdir -p "$(REPORT_DIR)"; \ result=0; \ + run_id="$$(date +%s%N)"; \ + now_ms="$$(date +%s%3N)"; \ + base_time_ms="$$((now_ms / 300000 * 300000 - 1800000))"; \ + port_seed="$$((run_id % 10000))"; \ + case_index=0; \ for test_case in $(CASES); do \ name="$${test_case%%:*}"; \ rest="$${test_case#*:}"; \ dataset="$${rest%%:*}"; \ suite="$${rest#*:}"; \ + project="asapquery-differential-$$name-$$run_id"; \ + case_base_time_ms="$$((base_time_ms + case_index * 60000))"; \ + port_offset="$$((port_seed + case_index))"; \ + prometheus_port="$$((20000 + port_offset))"; \ + query_port="$$((30000 + port_offset))"; \ + ingest_port="$$((40000 + port_offset))"; \ + case_index="$$((case_index + 1))"; \ echo "Running $$name ($$dataset, $$suite)"; \ - if ! go run ./cmd/differential-runner \ + if ! PROMETHEUS_PORT="$$prometheus_port" ASAP_QUERY_PORT="$$query_port" ASAP_INGEST_PORT="$$ingest_port" go run ./cmd/differential-runner \ --dataset "$$dataset" \ --suite "$$suite" \ --compose-file ../docker-compose.yml \ + --compose-project "$$project" \ + --reference-url "http://localhost:$$prometheus_port" \ + --reference-write-url "http://localhost:$$prometheus_port" \ + --test-url "http://localhost:$$query_port" \ + --test-write-url "http://localhost:$$ingest_port" \ + --base-time-ms "$$case_base_time_ms" \ --output "$(REPORT_DIR)/$$name.json"; then \ echo "FAILED: $$name" >&2; \ result=1; \ diff --git a/promql-compliance/runner/checked_in_fixture_validation_test.go b/promql-compliance/runner/checked_in_fixture_validation_test.go index 4da62761..e4a08898 100644 --- a/promql-compliance/runner/checked_in_fixture_validation_test.go +++ b/promql-compliance/runner/checked_in_fixture_validation_test.go @@ -29,3 +29,52 @@ func TestCheckedInFixturesAndSuites(t *testing.T) { } } } + +func TestCheckedInTemporalSuite(t *testing.T) { + if _, err := seeder.LoadFixture("../datasets/single-rate.yaml"); err != nil { + t.Fatalf("LoadFixture: %v", err) + } + suite, err := LoadSuiteFile("../suites/temporal.yaml") + if err != nil { + t.Fatalf("LoadSuiteFile: %v", err) + } + parser := promqlparser.NewParser(promqlparser.Options{}) + for _, query := range suite.Queries { + if _, err := parser.ParseExpr(query.Expr); err != nil { + t.Fatalf("parse %q: %v", query.Expr, err) + } + } +} + +func TestCheckedInOffGridRateSuite(t *testing.T) { + if _, err := seeder.LoadFixture("../datasets/single-rate.yaml"); err != nil { + t.Fatalf("LoadFixture: %v", err) + } + suite, err := LoadSuiteFile("../suites/off-grid-rate.yaml") + if err != nil { + t.Fatalf("LoadSuiteFile: %v", err) + } + parser := promqlparser.NewParser(promqlparser.Options{}) + for _, query := range suite.Queries { + if _, err := parser.ParseExpr(query.Expr); err != nil { + t.Fatalf("parse %q: %v", query.Expr, err) + } + } +} + +func TestCheckedInNativeDagSuites(t *testing.T) { + parser := promqlparser.NewParser(promqlparser.Options{}) + for _, path := range []string{ + "../suites/native-dag-aggregations.yaml", + } { + suite, err := LoadSuiteFile(path) + if err != nil { + t.Fatalf("LoadSuiteFile(%q): %v", path, err) + } + for _, query := range suite.Queries { + if _, err := parser.ParseExpr(query.Expr); err != nil { + t.Fatalf("parse %q: %v", query.Expr, err) + } + } + } +} diff --git a/promql-compliance/runner/run.go b/promql-compliance/runner/run.go index 3900b667..d4aebfbc 100644 --- a/promql-compliance/runner/run.go +++ b/promql-compliance/runner/run.go @@ -284,6 +284,8 @@ func (l *composeLifecycle) Start(ctx context.Context) error { command.Env = append(os.Environ(), l.env...) output, err := command.CombinedOutput() if err != nil { + l.started = true + l.Stop() return fmt.Errorf("start compose project %q: %w\n%s", l.project, err, output) } l.started = true diff --git a/promql-compliance/suites/native-dag-aggregations.yaml b/promql-compliance/suites/native-dag-aggregations.yaml new file mode 100644 index 00000000..94d7557d --- /dev/null +++ b/promql-compliance/suites/native-dag-aggregations.yaml @@ -0,0 +1,35 @@ +name: native-dag-aggregations +comparison_defaults: + value_tolerance: + relative: 0.01 + absolute: 0.000001 +queries: + - name: tumbling-sum + expr: sum(data) + instant_offsets_seconds: [300, 600, 900, 1200] + range: + start_offset_seconds: 300 + end_offset_seconds: 1200 + step_seconds: 60 + - name: keyed-count + expr: count(data) by (job) + instant_offsets_seconds: [300, 600, 900, 1200] + range: + start_offset_seconds: 300 + end_offset_seconds: 1200 + step_seconds: 60 + # Regression: native formatting must preserve the grouped label values. + - name: grouped-quantile-labels + expr: quantile by (job) (0.99, data) + instant_offsets_seconds: [300, 600, 900, 1200] + range: + start_offset_seconds: 300 + end_offset_seconds: 1200 + step_seconds: 60 + - name: self-keyed-topk + expr: topk(2, data) + instant_offsets_seconds: [300, 600, 900, 1200] + range: + start_offset_seconds: 300 + end_offset_seconds: 1200 + step_seconds: 60 diff --git a/promql-compliance/suites/off-grid-rate.yaml b/promql-compliance/suites/off-grid-rate.yaml new file mode 100644 index 00000000..48878ad2 --- /dev/null +++ b/promql-compliance/suites/off-grid-rate.yaml @@ -0,0 +1,13 @@ +name: off-grid-rate +comparison_defaults: + value_tolerance: + relative: 0 + absolute: 0.000001 +queries: + - name: off-grid-rate-capability-fallback + expr: rate(http_requests_total[5m]) + instant_offsets_seconds: [315, 630, 900] + range: + start_offset_seconds: 315 + end_offset_seconds: 900 + step_seconds: 45 diff --git a/promql-compliance/suites/temporal.yaml b/promql-compliance/suites/temporal.yaml index 69740cfc..54efa4cf 100644 --- a/promql-compliance/suites/temporal.yaml +++ b/promql-compliance/suites/temporal.yaml @@ -4,6 +4,20 @@ comparison_defaults: relative: 0 absolute: 0.000001 queries: + - name: sum-over-time + expr: sum_over_time(http_requests_total[5m]) + instant_offsets_seconds: [300, 600, 1200] + range: + start_offset_seconds: 300 + end_offset_seconds: 1200 + step_seconds: 60 + - name: sum-over-time-plus-one + expr: sum_over_time(http_requests_total[5m]) + 1 + instant_offsets_seconds: [300, 600, 1200] + range: + start_offset_seconds: 300 + end_offset_seconds: 1200 + step_seconds: 60 - name: request-rate expr: rate(http_requests_total[5m]) instant_offsets_seconds: [300, 600, 1200]