From 6f5cbd11b5172493acdb695a83fbc77170b2b82b Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Wed, 30 Sep 2026 21:20:30 -0400 Subject: [PATCH 1/2] refactor(query-engine): expose PromQL execution errors --- .../src/drivers/query/servers/http.rs | 28 ++- asap-query-engine/src/engines/mod.rs | 2 +- .../src/engines/simple_engine/mod.rs | 83 ++++-- .../src/engines/simple_engine/promql.rs | 238 +++++++++++------- asap-query-engine/src/lib.rs | 2 +- .../src/tests/dispatch_arithmetic_tests.rs | 30 ++- .../exact_window_grid_adversarial_tests.rs | 21 +- .../native_binary_arithmetic_plan_tests.rs | 32 ++- .../src/tests/native_binary_instant_tests.rs | 33 ++- .../src/tests/native_range_query_tests.rs | 119 ++++++--- .../src/tests/query_equivalence_tests.rs | 38 ++- .../src/tests/range_query_arithmetic_tests.rs | 55 ++-- ...stage_e_instant_range_equivalence_tests.rs | 7 + .../window_semantics_consistency_tests.rs | 33 ++- .../tests/e2e_precompute_equivalence.rs | 2 + 15 files changed, 521 insertions(+), 202 deletions(-) diff --git a/asap-query-engine/src/drivers/query/servers/http.rs b/asap-query-engine/src/drivers/query/servers/http.rs index 7fe38a39..dd185b43 100644 --- a/asap-query-engine/src/drivers/query/servers/http.rs +++ b/asap-query-engine/src/drivers/query/servers/http.rs @@ -15,7 +15,7 @@ use tokio::net::TcpListener; use tracing::{debug, info}; use crate::drivers::query::adapters::{create_http_adapter, AdapterConfig, HttpProtocolAdapter}; -use crate::engines::SimpleEngine; +use crate::engines::{QueryExecutionError, SimpleEngine}; use crate::query_tracker::QueryTracker; use crate::stores::Store; @@ -44,6 +44,22 @@ struct AppState { fallback: Option>, } +async fn format_native_execution_error(state: &AppState, error: QueryExecutionError) -> Response { + match state + .adapter + .format_error_response( + &crate::drivers::query::adapters::AdapterError::ProtocolError(error.to_string()), + ) + .await + { + Ok(mut response) => { + *response.status_mut() = StatusCode::INTERNAL_SERVER_ERROR; + response + } + Err(status) => status.into_response(), + } +} + impl HttpServer { pub fn new( config: HttpServerConfig, @@ -193,7 +209,7 @@ async fn process_query_request( .query_engine .handle_query(parsed_request.query.clone(), parsed_request.time) { - Some((query_output_labels, query_result)) => { + Ok(Some((query_output_labels, query_result))) => { let query_duration = query_start_time.elapsed(); debug!("=== QUERY ENGINE SUCCESS ==="); debug!( @@ -233,7 +249,7 @@ async fn process_query_request( Err(status) => status.into_response(), } } - None => { + Ok(None) => { let total_duration = start_time.elapsed(); debug!("=== QUERY ENGINE RETURNED NONE ==="); debug!( @@ -271,6 +287,7 @@ async fn process_query_request( } } } + Err(error) => format_native_execution_error(state, error).await, } } @@ -524,7 +541,7 @@ async fn process_range_query_request( parsed_request.end, parsed_request.step, ) { - Some((query_output_labels, query_result)) => { + Ok(Some((query_output_labels, query_result))) => { let query_duration = query_start_time.elapsed(); debug!( "Range query execution took: {:.2}ms", @@ -547,7 +564,7 @@ async fn process_range_query_request( Err(status) => status.into_response(), } } - None => { + Ok(None) => { let total_duration = start_time.elapsed(); debug!("Range query returned None - query not supported"); @@ -578,6 +595,7 @@ async fn process_range_query_request( } } } + Err(error) => format_native_execution_error(state, error).await, } } diff --git a/asap-query-engine/src/engines/mod.rs b/asap-query-engine/src/engines/mod.rs index 9d333640..a33d545f 100644 --- a/asap-query-engine/src/engines/mod.rs +++ b/asap-query-engine/src/engines/mod.rs @@ -6,5 +6,5 @@ pub(crate) mod sliding_window_composition; pub mod window_merger; pub use query_result::{InstantVector, QueryResult, RangeVector, RangeVectorElement, Sample}; -pub use simple_engine::SimpleEngine; +pub use simple_engine::{QueryExecutionError, SimpleEngine}; pub use window_merger::{create_window_merger, NaiveMerger, WindowMerger}; diff --git a/asap-query-engine/src/engines/simple_engine/mod.rs b/asap-query-engine/src/engines/simple_engine/mod.rs index b5b4180f..c5c2c032 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -37,6 +37,13 @@ use serde_json::Value; #[allow(dead_code)] type MergedOutputsMap = HashMap, Box>; +const NO_LOCAL_DATA_ERROR_PREFIXES: [&str; 4] = [ + "No precomputed outputs found", + "No data found", + "Incomplete Sliding-window cover", + "Exact Prometheus counter bounds are unavailable for off-grid", +]; + /// Metadata extracted from a query, independent of query language #[derive(Debug, Clone)] pub struct QueryMetadata { @@ -54,6 +61,13 @@ pub struct QueryMetadata { pub keep_metric_name: bool, } +/// A native query was accepted but could not be executed locally. +#[derive(Debug, thiserror::Error)] +pub enum QueryExecutionError { + #[error("native query execution failed: {0}")] + Native(String), +} + /// Parameters for a single store query #[derive(Debug, Clone)] pub struct StoreQueryParams { @@ -1536,27 +1550,60 @@ impl SimpleEngine { enable_topk_limiting: bool, enable_topk_formatting: bool, ) -> Option<(KeyByLabelNames, QueryResult)> { - let results = self - .execute_query_pipeline(&context, enable_topk_limiting, enable_topk_formatting) - .map_err(|e| { - warn!("Query execution failed: {}", e); - e - }) - .ok()?; - Some(( + self.execute_context_result(context, enable_topk_limiting, enable_topk_formatting) + .map_err(|error| warn!("Query execution failed: {error}")) + .ok() + .flatten() + } + + fn execute_context_result( + &self, + context: QueryExecutionContext, + enable_topk_limiting: bool, + enable_topk_formatting: bool, + ) -> Result, QueryExecutionError> { + let Some(results) = Self::classify_native_execution(self.execute_query_pipeline( + &context, + enable_topk_limiting, + enable_topk_formatting, + ))? + else { + return Ok(None); + }; + Ok(Some(( context.metadata.query_output_labels, QueryResult::vector(results, context.query_time), - )) + ))) + } + + fn classify_native_execution( + result: Result, + ) -> Result, QueryExecutionError> { + match result { + Ok(value) => Ok(Some(value)), + Err(error) + if NO_LOCAL_DATA_ERROR_PREFIXES + .iter() + .any(|prefix| error.starts_with(prefix)) => + { + Ok(None) + } + Err(error) => Err(QueryExecutionError::Native(error)), + } } /// Handle a query following Python's unified architecture // pub async fn handle_query( - pub fn handle_query(&self, query: String, time: f64) -> Option<(KeyByLabelNames, QueryResult)> { + pub fn handle_query( + &self, + query: String, + time: f64, + ) -> Result, QueryExecutionError> { match self.query_language { QueryLanguage::promql => self.handle_query_promql(query, time), - QueryLanguage::sql => self.handle_query_sql(query, time), - QueryLanguage::elastic_querydsl => self.handle_query_elastic(query, time), - QueryLanguage::elastic_sql => self.handle_query_sql(query, time), + QueryLanguage::sql => Ok(self.handle_query_sql(query, time)), + QueryLanguage::elastic_querydsl => Ok(self.handle_query_elastic(query, time)), + QueryLanguage::elastic_sql => Ok(self.handle_query_sql(query, time)), } } @@ -4088,7 +4135,7 @@ mod sketch_query_tests { // fn test_sketch_instant_entropy_over_time() { // let engine = engine_with_sketch_data("mymetric"); // // Query at time 0.1s (= 100ms) with a 100ms range - // let result = engine.handle_query_promql("entropy_over_time(mymetric[100s])".into(), 0.1); + // let result = engine.handle_query_promql("entropy_over_time(mymetric[100s])".into(), 0.1).expect("native query execution should not fail"); // assert!(result.is_some(), "entropy_over_time should return a result"); // let (labels, qr) = result.unwrap(); // assert!(!labels.labels.is_empty()); @@ -4105,7 +4152,7 @@ mod sketch_query_tests { // fn test_sketch_instant_quantile_over_time() { // let engine = engine_with_sketch_data("mymetric"); // let result = - // engine.handle_query_promql("quantile_over_time(0.5, mymetric[100s])".into(), 0.1); + // engine.handle_query_promql("quantile_over_time(0.5, mymetric[100s])".into(), 0.1).expect("native query execution should not fail"); // assert!( // result.is_some(), // "quantile_over_time should return a result" @@ -4128,7 +4175,7 @@ mod sketch_query_tests { // #[test] // fn test_sketch_instant_avg_over_time() { // let engine = engine_with_sketch_data("cpu"); - // let result = engine.handle_query_promql("avg_over_time(cpu[100s])".into(), 0.1); + // let result = engine.handle_query_promql("avg_over_time(cpu[100s])".into(), 0.1).expect("native query execution should not fail"); // assert!(result.is_some(), "avg_over_time should return a result"); // let (_labels, qr) = result.unwrap(); // if let crate::engines::query_result::QueryResult::Vector(iv) = qr { @@ -4188,7 +4235,7 @@ mod sketch_query_tests { // 0.01, // 0.1, // 0.01, - // ); + // ).expect("native query execution should not fail"); // assert!( // result.is_some(), // "sketch range query should return a result" @@ -4958,6 +5005,7 @@ mod stage_e4_instant_wrapper_equivalence_tests { ); let (_, qr) = engine .handle_query_promql("sum(cpu_load) by (host) * 5".to_string(), 3.0) + .expect("native query execution should not fail") .expect("scalar binary-expr query should resolve"); assert_eq!(matrix_metric(qr), 500.0); } @@ -5045,6 +5093,7 @@ mod stage_e4_instant_wrapper_equivalence_tests { "sum(metric_a) by (host) + sum(metric_b) by (host)".to_string(), 3.0, ) + .expect("native query execution should not fail") .expect("vector-vector binary-expr query should resolve"); assert_eq!(matrix_metric(qr), 30.0); } diff --git a/asap-query-engine/src/engines/simple_engine/promql.rs b/asap-query-engine/src/engines/simple_engine/promql.rs index 324977d5..a00e1c25 100644 --- a/asap-query-engine/src/engines/simple_engine/promql.rs +++ b/asap-query-engine/src/engines/simple_engine/promql.rs @@ -4,7 +4,10 @@ //! dispatch, range-query handling, and query dispatch. use super::SimpleEngine; -use super::{QueryExecutionContext, QueryMetadata, QueryTimestamps, RangeQueryExecutionContext}; +use super::{ + QueryExecutionContext, QueryExecutionError, QueryMetadata, QueryTimestamps, + RangeQueryExecutionContext, +}; use crate::data_model::{AggregationIdInfo, KeyByLabelValues, QueryConfig, SchemaConfig}; use crate::engines::query_result::{InstantVectorElement, QueryResult, RangeVectorElement}; use asap_types::query_requirements::build_query_requirements_promql; @@ -19,6 +22,8 @@ use tracing::{debug, warn}; pub(super) const METRIC_NAME_LABEL: &str = "__name__"; +type BinaryInstantArm = (Vec, Vec); + /// Detects whether either side of a PromQL binary expression is a scalar /// (numeric literal), returning the scalar value, the other (vector) arm, /// and whether the scalar was on the left. Shared by instant and range @@ -452,57 +457,52 @@ impl SimpleEngine { &self, arm_ast: &promql_parser::parser::Expr, time: f64, - ) -> Option<(Vec, Vec)> { + ) -> Result, QueryExecutionError> { use promql_parser::parser::Expr; match arm_ast { - Expr::NumberLiteral(_) => None, // caller handles scalars + Expr::NumberLiteral(_) => Ok(None), // caller handles scalars Expr::Paren(paren) => self.evaluate_binary_arm(&paren.expr, time), Expr::Binary(binary) => { if binary.modifier.is_some() { - return None; + return Ok(None); } if !is_supported_binary_arithmetic_op(&binary.op) { - return None; + return Ok(None); } // Nested binary expression — recurse on both sides - let (lhs_results, lhs_labels) = self.evaluate_binary_arm(&binary.lhs, time)?; - let (rhs_results, rhs_labels) = self.evaluate_binary_arm(&binary.rhs, time)?; - let combined = combine_vector_vector( + let Some((lhs_results, lhs_labels)) = + self.evaluate_binary_arm(&binary.lhs, time)? + else { + return Ok(None); + }; + let Some((rhs_results, rhs_labels)) = + self.evaluate_binary_arm(&binary.rhs, time)? + else { + return Ok(None); + }; + let Some(combined) = combine_vector_vector( lhs_results, &lhs_labels, rhs_results, &rhs_labels, &binary.op, - )?; - Some((combined, lhs_labels)) + ) else { + return Ok(None); + }; + Ok(Some((combined, lhs_labels))) } _ => { - let (ctx, label_names) = self.resolve_arm_leaf_context(arm_ast, time)?; - // Unlike DataFusion's PrecomputedSummaryReadExec (which streamed - // whatever rows existed, including zero, so a currently-empty arm - // used to return Some(empty vector)), execute_query_pipeline errors - // when the store has no precomputed outputs at all for this arm — - // that propagates to None here, triggering a full Prometheus - // fallback for the whole expression instead of an empty result - // for just this arm. Accepted behavior change (#567); warn loudly - // so it's visible rather than silent. - let results = self - // Binary arms need Topk limiting, but must remain in the - // unformatted intermediate label representation until - // after the binary join. - .execute_query_pipeline(&ctx, true, false) - .map_err(|e| { - warn!( - "Binary-expr arm for metric '{}' failed ({}) — \ - falls back to Prometheus for the whole expression rather than \ - returning an empty result for just this arm", - ctx.metric, e - ); - e - }) - .ok()?; - Some((results, label_names)) + let Some((ctx, label_names)) = self.resolve_arm_leaf_context(arm_ast, time) else { + return Ok(None); + }; + let Some(results) = Self::classify_native_execution( + self.execute_query_pipeline(&ctx, true, false), + )? + else { + return Ok(None); + }; + Ok(Some((results, label_names))) } } } @@ -515,21 +515,21 @@ impl SimpleEngine { &self, ast: &promql_parser::parser::Expr, time: f64, - ) -> Option<(KeyByLabelNames, QueryResult)> { + ) -> Result, QueryExecutionError> { use promql_parser::parser::Expr; let query_time = Self::convert_query_time_to_data_time(time); let binary = match ast { Expr::Binary(b) => b, - _ => return None, + _ => return Ok(None), }; if !is_supported_binary_arithmetic_op(&binary.op) { - return None; + return Ok(None); } if binary.modifier.is_some() { - return None; + return Ok(None); } let lhs = binary.lhs.as_ref(); @@ -537,21 +537,34 @@ impl SimpleEngine { let op = &binary.op; if let Some((scalar, vector_arm, scalar_on_left)) = detect_scalar_arm(lhs, rhs) { - let (vector_results, label_names) = self.evaluate_binary_arm(vector_arm, time)?; + let Some((vector_results, label_names)) = self.evaluate_binary_arm(vector_arm, time)? + else { + return Ok(None); + }; let combined = combine_scalar(vector_results, scalar, op, scalar_on_left); - return Some(( + return Ok(Some(( KeyByLabelNames::new(label_names), QueryResult::vector(combined, query_time), - )); + ))); } // Vector–vector - let (lhs_results, lhs_labels) = self.evaluate_binary_arm(lhs, time)?; - let (rhs_results, rhs_labels) = self.evaluate_binary_arm(rhs, time)?; - let combined = - combine_vector_vector(lhs_results, &lhs_labels, rhs_results, &rhs_labels, op)?; + let Some((lhs_results, lhs_labels)) = self.evaluate_binary_arm(lhs, time)? else { + return Ok(None); + }; + let Some((rhs_results, rhs_labels)) = self.evaluate_binary_arm(rhs, time)? else { + return Ok(None); + }; + let Some(combined) = + combine_vector_vector(lhs_results, &lhs_labels, rhs_results, &rhs_labels, op) + else { + return Ok(None); + }; let output_labels = KeyByLabelNames::new(lhs_labels); - Some((output_labels, QueryResult::vector(combined, query_time))) + Ok(Some(( + output_labels, + QueryResult::vector(combined, query_time), + ))) } /// Applies a PromQL binary arithmetic operator to two f64 values. @@ -718,19 +731,19 @@ impl SimpleEngine { start: f64, end: f64, step: f64, - ) -> Option<(KeyByLabelNames, QueryResult)> { + ) -> Result, QueryExecutionError> { use promql_parser::parser::Expr; let binary = match ast { Expr::Binary(b) => b, - _ => return None, + _ => return Ok(None), }; if !is_supported_binary_arithmetic_op(&binary.op) { - return None; + return Ok(None); } if binary.modifier.is_some() { - return None; + return Ok(None); } let lhs = binary.lhs.as_ref(); @@ -738,13 +751,19 @@ impl SimpleEngine { let op = &binary.op; if let Some((scalar, vector_arm, scalar_on_left)) = detect_scalar_arm(lhs, rhs) { - let (ctx, labels) = self.build_arm_range_context(vector_arm, start, end, step)?; + let Some((ctx, labels)) = self.build_arm_range_context(vector_arm, start, end, step) + else { + return Ok(None); + }; // Binary arms need Topk limiting, but must remain in the // unformatted intermediate label representation until after the // arithmetic operation. - let results = self - .execute_observed_range_query_pipeline(&ctx, true, false) - .ok()?; + let Some(results) = Self::classify_native_execution( + self.execute_observed_range_query_pipeline(&ctx, true, false), + )? + else { + return Ok(None); + }; let combined: Vec = results .into_iter() .map(|mut elem| { @@ -758,7 +777,10 @@ impl SimpleEngine { elem }) .collect(); - return Some((KeyByLabelNames::new(labels), QueryResult::matrix(combined))); + return Ok(Some(( + KeyByLabelNames::new(labels), + QueryResult::matrix(combined), + ))); } // Vector-vector: evaluate both arms, join by label key, apply op per matching timestamp. @@ -767,18 +789,30 @@ impl SimpleEngine { // KeyByLabelValues equality below is only safe once the label *names* // match (they're canonically sorted by KeyByLabelNames::new(), so two // arms with the same label set always order their values the same way). - let (lhs_ctx, lhs_labels) = self.build_arm_range_context(lhs, start, end, step)?; - let (rhs_ctx, rhs_labels) = self.build_arm_range_context(rhs, start, end, step)?; + let Some((lhs_ctx, lhs_labels)) = self.build_arm_range_context(lhs, start, end, step) + else { + return Ok(None); + }; + let Some((rhs_ctx, rhs_labels)) = self.build_arm_range_context(rhs, start, end, step) + else { + return Ok(None); + }; if lhs_labels != rhs_labels { - return None; + return Ok(None); } // Binary arms need Topk limiting, but not final presentation formatting. - let lhs_results = self - .execute_observed_range_query_pipeline(&lhs_ctx, true, false) - .ok()?; - let rhs_results = self - .execute_observed_range_query_pipeline(&rhs_ctx, true, false) - .ok()?; + let Some(lhs_results) = Self::classify_native_execution( + self.execute_observed_range_query_pipeline(&lhs_ctx, true, false), + )? + else { + return Ok(None); + }; + let Some(rhs_results) = Self::classify_native_execution( + self.execute_observed_range_query_pipeline(&rhs_ctx, true, false), + )? + else { + return Ok(None); + }; // Build lookup: label_key -> {timestamp -> value} for rhs let mut rhs_map: HashMap> = HashMap::new(); @@ -810,7 +844,7 @@ impl SimpleEngine { } let output_labels = KeyByLabelNames::new(lhs_labels); - Some((output_labels, QueryResult::matrix(combined))) + Ok(Some((output_labels, QueryResult::matrix(combined)))) } // /// Try to extract sketch query components from a PromQL query string. @@ -1053,7 +1087,7 @@ impl SimpleEngine { &self, query: String, time: f64, - ) -> Option<(KeyByLabelNames, QueryResult)> { + ) -> Result, QueryExecutionError> { let query_start_time = Instant::now(); debug!("Handling query: {} at time {}", query, time); @@ -1061,7 +1095,7 @@ impl SimpleEngine { Ok(ast) => ast, Err(e) => { warn!("Failed to parse PromQL query '{}': {}", query, e); - return None; + return Ok(None); } }; @@ -1077,7 +1111,10 @@ impl SimpleEngine { return result; } - let context = self.build_query_execution_context_from_parsed(&ast, &query, time)?; + let Some(context) = self.build_query_execution_context_from_parsed(&ast, &query, time) + else { + return Ok(None); + }; debug!( "Querying store for metric: {}, aggregation_id: {}, range: [{}, {}]", @@ -1087,7 +1124,7 @@ impl SimpleEngine { context.store_plan.values_query.end_timestamp ); - let result = self.execute_context(context, true, true); + let result = self.execute_context_result(context, true, true)?; // Determine query routing order based on function type. // USampling functions prefer the precomputed path first (sketch fallback), @@ -1160,7 +1197,7 @@ impl SimpleEngine { "Total query handling took: {:.2}ms (no results)", total_query_duration.as_secs_f64() * 1000.0 ); - result + Ok(result) } pub fn build_query_execution_context_promql( @@ -1311,7 +1348,7 @@ impl SimpleEngine { start: f64, end: f64, step: f64, - ) -> Option<(KeyByLabelNames, QueryResult)> { + ) -> Result, QueryExecutionError> { let query_start_time = Instant::now(); debug!( "Handling range query: {} from {} to {} step {}", @@ -1322,7 +1359,7 @@ impl SimpleEngine { Ok(ast) => ast, Err(e) => { warn!("Failed to parse PromQL query '{}': {}", query, e); - return None; + return Ok(None); } }; @@ -1337,19 +1374,21 @@ impl SimpleEngine { return result; } - let context = - self.build_range_query_execution_context_from_parsed(&ast, &query, start, end, step)?; + let Some(context) = + self.build_range_query_execution_context_from_parsed(&ast, &query, start, end, step) + else { + return Ok(None); + }; // Execute range query pipeline. (true, true): self-gated, same as // instant's handle_query_promql -- both flags are no-ops unless this // query's statistic is Topk. - let results: Vec = self - .execute_observed_range_query_pipeline(&context, true, true) - .map_err(|e| { - warn!("Range query execution failed: {}", e); - e - }) - .ok()?; + let Some(results): Option> = Self::classify_native_execution( + self.execute_observed_range_query_pipeline(&context, true, true), + )? + else { + return Ok(None); + }; // // Determine query routing order based on function type. // // USampling functions prefer the precomputed path first (sketch fallback), @@ -1407,10 +1446,10 @@ impl SimpleEngine { total_duration.as_secs_f64() * 1000.0 ); - Some(( + Ok(Some(( context.base.metadata.query_output_labels, QueryResult::matrix(results), - )) + ))) } } @@ -1560,6 +1599,30 @@ mod topk_pipeline_tests { (engine, store) } + #[test] + fn public_promql_methods_report_capability_misses_as_ok_none() { + let (engine, _store) = build_topk_engine(); + + assert!(matches!( + engine.handle_query_promql("topk(".to_string(), QUERY_TIME), + Ok(None) + )); + assert!(matches!( + engine.handle_range_query_promql( + "topk(".to_string(), + QUERY_TIME - 1.0, + QUERY_TIME, + 1.0 + ), + Ok(None) + )); + let no_local_data = engine.handle_query_promql(TOPK_QUERY.to_string(), QUERY_TIME); + assert!( + matches!(no_local_data, Ok(None)), + "missing local data should be a capability miss, got {no_local_data:?}" + ); + } + #[test] fn detects_topk_and_resolves_self_keyed_heap() { let (engine, _store) = build_topk_engine(); @@ -1706,8 +1769,10 @@ mod topk_pipeline_tests { .insert_precomputed_output(output, Box::new(sketch)) .expect("insert should succeed"); - let (_, query_result) = engine - .handle_query_promql(format!("{TOPK_QUERY} + 0"), QUERY_TIME) + let result = engine.handle_query_promql(format!("{TOPK_QUERY} + 0"), QUERY_TIME); + assert!(matches!(&result, Ok(Some(_)))); + let (_, query_result) = result + .expect("native query execution should not fail") .expect("binary-expr-wrapped topk should still resolve"); let results = match query_result { @@ -1940,6 +2005,7 @@ mod topk_pipeline_tests { format!("{TOPK_OVER_SUM_OVER_TIME_QUERY} + 0"), TOPK_OVER_SUM_OVER_TIME_QUERY_TIME, ) + .expect("native query execution should not fail") .expect("binary-expr-wrapped topk-over-sum_over_time should still resolve"); assert_eq!( diff --git a/asap-query-engine/src/lib.rs b/asap-query-engine/src/lib.rs index 7419fafe..b86d8148 100644 --- a/asap-query-engine/src/lib.rs +++ b/asap-query-engine/src/lib.rs @@ -33,7 +33,7 @@ pub use precompute_operators::{ pub use stores::{SimpleMapStore, Store, StoreResult}; -pub use engines::{InstantVector, QueryResult, SimpleEngine}; +pub use engines::{InstantVector, QueryExecutionError, QueryResult, SimpleEngine}; pub use drivers::{HttpServer, HttpServerConfig, OtlpReceiver, OtlpReceiverConfig}; diff --git a/asap-query-engine/src/tests/dispatch_arithmetic_tests.rs b/asap-query-engine/src/tests/dispatch_arithmetic_tests.rs index 69c888e7..7bfa27b5 100644 --- a/asap-query-engine/src/tests/dispatch_arithmetic_tests.rs +++ b/asap-query-engine/src/tests/dispatch_arithmetic_tests.rs @@ -35,10 +35,12 @@ mod tests { "sum(requests_total) by (host)", ); - let result = engine.handle_query_promql( - "sum(errors_total) by (host) / sum(requests_total) by (host)".to_string(), - QUERY_TIME, - ); + let result = engine + .handle_query_promql( + "sum(errors_total) by (host) / sum(requests_total) by (host)".to_string(), + QUERY_TIME, + ) + .expect("native query execution should not fail"); assert!(result.is_some(), "Binary query should return Some"); let (labels, qr) = result.unwrap(); assert!(!labels.labels.is_empty(), "Should have output label names"); @@ -65,10 +67,12 @@ mod tests { ); // foo() is not a supported PromQL function → arm lookup fails → returns None - let result = engine.handle_query_promql( - "foo(errors_total[5m]) / sum(requests_total) by (host)".to_string(), - QUERY_TIME, - ); + let result = engine + .handle_query_promql( + "foo(errors_total[5m]) / sum(requests_total) by (host)".to_string(), + QUERY_TIME, + ) + .expect("native query execution should not fail"); assert!( result.is_none(), "Should return None for non-acceleratable arm (graceful fallback)" @@ -88,8 +92,9 @@ mod tests { "sum(errors_total) by (host)", ); - let result = - engine.handle_query_promql("sum(errors_total) by (host) * 100".to_string(), QUERY_TIME); + let result = engine + .handle_query_promql("sum(errors_total) by (host) * 100".to_string(), QUERY_TIME) + .expect("native query execution should not fail"); assert!(result.is_some(), "Scalar binary should return Some"); let (_, qr) = result.unwrap(); let elements = match qr { @@ -120,8 +125,9 @@ mod tests { "sum(http_requests) by (host)", ); - let result = - engine.handle_query_promql("sum(http_requests) by (host)".to_string(), QUERY_TIME); + let result = engine + .handle_query_promql("sum(http_requests) by (host)".to_string(), QUERY_TIME) + .expect("native query execution should not fail"); assert!( result.is_some(), "Single-metric query should still work after binary dispatch" diff --git a/asap-query-engine/src/tests/exact_window_grid_adversarial_tests.rs b/asap-query-engine/src/tests/exact_window_grid_adversarial_tests.rs index 3ad9ed45..b42e3c26 100644 --- a/asap-query-engine/src/tests/exact_window_grid_adversarial_tests.rs +++ b/asap-query-engine/src/tests/exact_window_grid_adversarial_tests.rs @@ -568,7 +568,9 @@ mod tests { WindowType::Sliding, ); - let result = engine.handle_range_query_promql(query.to_string(), 3.0, 5.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 3.0, 5.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); let samples = host_a_samples(&elements); @@ -623,6 +625,7 @@ mod tests { let (_, qr) = engine .handle_query_promql(query.to_string(), 4.0) + .expect("native query execution should not fail") .expect("instant query failed"); let values = vector_values(qr); assert_eq!(values.len(), 1, "expected exactly one series for host-a"); @@ -683,7 +686,9 @@ mod tests { WindowType::Sliding, ); - let result = engine.handle_range_query_promql(query.to_string(), 2.0, 4.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 2.0, 4.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); let samples = host_a_samples(&elements); @@ -738,7 +743,9 @@ mod tests { WindowType::Tumbling, ); - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 5.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 5.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); let samples = host_a_samples(&elements); @@ -901,7 +908,9 @@ mod tests { let query_time_sec = base_ts as f64 / 1000.0; let call_start = Instant::now(); - let result = engine.handle_query_promql(query.to_string(), query_time_sec); + let result = engine + .handle_query_promql(query.to_string(), query_time_sec) + .expect("native query execution should not fail"); let elapsed = call_start.elapsed(); let (_, qr) = result.expect("query failed to resolve real data near a huge timestamp"); @@ -961,7 +970,9 @@ mod tests { // exercises execute_range_query_pipeline's per-step keys merge // exactly once, same shape as the instant case, through the range // entry point instead. - let result = engine.handle_range_query_promql(query.to_string(), t, t + 1.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), t, t + 1.0, 1.0) + .expect("native query execution should not fail"); let elapsed = call_start.elapsed(); let (_, qr) = diff --git a/asap-query-engine/src/tests/native_binary_arithmetic_plan_tests.rs b/asap-query-engine/src/tests/native_binary_arithmetic_plan_tests.rs index 69a119ea..98de85a2 100644 --- a/asap-query-engine/src/tests/native_binary_arithmetic_plan_tests.rs +++ b/asap-query-engine/src/tests/native_binary_arithmetic_plan_tests.rs @@ -63,7 +63,9 @@ mod tests { ); let query = "sum(errors_total) by (host) / sum(requests_total) by (host)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); assert!(result.is_some(), "Expected Some result for binary query"); let (_, qr) = result.unwrap(); let elements = match qr { @@ -99,7 +101,9 @@ mod tests { ); let query = "sum(metric_a) by (host) * sum(metric_b) by (host)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result"); let elements = match qr { crate::engines::query_result::QueryResult::Vector(iv) => iv.values, @@ -130,7 +134,9 @@ mod tests { ); let query = "sum(metric_a) by (host) + sum(metric_b) by (host)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result"); let elements = match qr { crate::engines::query_result::QueryResult::Vector(iv) => iv.values, @@ -161,7 +167,9 @@ mod tests { ); let query = "sum(metric_a) by (host) - sum(metric_b) by (host)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result"); let elements = match qr { crate::engines::query_result::QueryResult::Vector(iv) => iv.values, @@ -207,7 +215,9 @@ mod tests { ); let query = "sum(errors_total) by (host) / sum(requests_total) by (host)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result"); let elements = match qr { crate::engines::query_result::QueryResult::Vector(iv) => iv.values, @@ -244,7 +254,9 @@ mod tests { ); let query = "sum(errors_total) by (host) * 100"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result for scalar-right multiply"); let elements = match qr { crate::engines::query_result::QueryResult::Vector(iv) => iv.values, @@ -278,7 +290,9 @@ mod tests { ); let query = "1 - sum(success_total) by (host)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result for scalar-left subtract"); let elements = match qr { crate::engines::query_result::QueryResult::Vector(iv) => iv.values, @@ -347,7 +361,9 @@ mod tests { ); let query = "(sum(metric_a) by (host) + sum(metric_b) by (host)) / sum(metric_c) by (host)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result for nested binary expression"); let elements = match qr { crate::engines::query_result::QueryResult::Vector(iv) => iv.values, diff --git a/asap-query-engine/src/tests/native_binary_instant_tests.rs b/asap-query-engine/src/tests/native_binary_instant_tests.rs index 9aef669f..57c6a308 100644 --- a/asap-query-engine/src/tests/native_binary_instant_tests.rs +++ b/asap-query-engine/src/tests/native_binary_instant_tests.rs @@ -72,6 +72,7 @@ mod tests { let query = format!("sum(metric_a) by (host) {op} sum(metric_b) by (host)"); let (_, qr) = engine .handle_query_promql(query, QUERY_TIME) + .expect("native query execution should not fail") .unwrap_or_else(|| panic!("query failed for op {op}")); let values = vector_values(qr); @@ -108,6 +109,7 @@ mod tests { let query = "sum(metric_a) by (host) ^ sum(metric_b) by (host)"; let (_, qr) = engine .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail") .expect("query failed"); let values = vector_values(qr); assert!((values[0].1 - 1024.0).abs() < 1e-6, "2^10 = 1024"); @@ -132,6 +134,7 @@ mod tests { ] { let (_, qr) = engine .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail") .unwrap_or_else(|| panic!("query failed for {query}")); let values = vector_values(qr); assert!((values[0].1 - 700.0).abs() < 1e-10, "{query}"); @@ -171,6 +174,7 @@ mod tests { let query = "(sum(metric_a) by (host) + sum(metric_b) by (host)) * sum(metric_c) by (host)"; let (_, qr) = engine .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail") .expect("query failed"); let values = vector_values(qr); assert!((values[0].1 - 90.0).abs() < 1e-10, "(10+20)*3 = 90"); @@ -201,7 +205,9 @@ mod tests { ); let query = "sum(metric_a) by (host) + sum(metric_b) by (host)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); assert!( result.is_none(), "arm with no current precomputed data falls back to Prometheus" @@ -226,6 +232,7 @@ mod tests { let query = "foo(errors_total[5m]) / sum(requests_total) by (host)"; assert!(engine .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail") .is_none()); } @@ -255,6 +262,7 @@ mod tests { let query = "count(event_frequency) by (host, event) + 0"; let (_, qr) = engine .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail") .expect("query failed"); assert!(!sorted(vector_values(qr)).is_empty()); } @@ -318,6 +326,7 @@ mod tests { let query = "count(event_frequency) by (region, host, event)"; let (_, qr) = engine .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail") .expect( "instant query should succeed by skipping the value-less region=orphan \ group, not fail the entire query because of it", @@ -403,6 +412,7 @@ mod tests { let query = "count(event_frequency) by (region, host, event) + 0"; let (_, qr) = engine .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail") .expect( "instant query should succeed by skipping the unresolvable region=broken group", ); @@ -461,7 +471,7 @@ mod tests { ); let query = "sum(event_frequency) by (host, event) + 0"; - let (_, qr) = engine.handle_query_promql(query.to_string(), QUERY_TIME).expect( + let (_, qr) = engine.handle_query_promql(query.to_string(), QUERY_TIME).expect("native query execution should not fail").expect( "instant query should succeed by skipping the one key missing from the value accumulator", ); let values = sorted(vector_values(qr)); @@ -507,6 +517,7 @@ mod tests { let query = format!("{leaf_query} + 0"); let (_, qr) = engine .handle_query_promql(query, QUERY_TIME) + .expect("native query execution should not fail") .expect("query failed"); let values = vector_values(qr); assert_eq!(values.len(), 1); @@ -536,8 +547,9 @@ mod tests { "sum(errors_total) by (host)", ); - let result = - engine.handle_query_promql("sum(errors_total) by (host) * 2".to_string(), QUERY_TIME); + let result = engine + .handle_query_promql("sum(errors_total) by (host) * 2".to_string(), QUERY_TIME) + .expect("native query execution should not fail"); assert!(result.is_some()); } @@ -578,7 +590,9 @@ mod tests { ); let query = "sum(metric_a) by (host) + sum(metric_b) by (region)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); assert!( result.is_none(), @@ -611,7 +625,9 @@ mod tests { ); let query = "sum(metric_a) by (host, dc) + sum(metric_b) by (region, zone)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); assert!( result.is_none(), @@ -645,7 +661,9 @@ mod tests { ); let query = "sum(metric_a) by (host) + sum(metric_b) by (region)"; - let result = engine.handle_query_promql(query.to_string(), QUERY_TIME); + let result = engine + .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail"); assert!( result.is_none(), @@ -680,6 +698,7 @@ mod tests { let query = "sum(metric_a) by (host) + sum(metric_b) by (host)"; let (_, qr) = engine .handle_query_promql(query.to_string(), QUERY_TIME) + .expect("native query execution should not fail") .expect("query failed"); assert_eq!(vector_values(qr), Vec::new()); 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 dd71c954..28d6734a 100644 --- a/asap-query-engine/src/tests/native_range_query_tests.rs +++ b/asap-query-engine/src/tests/native_range_query_tests.rs @@ -135,6 +135,7 @@ mod tests { 1_002.0, 1.0, ) + .expect("native query execution should not fail") .expect("range rate query failed"); let elements = matrix_values(result); let samples = &elements @@ -209,11 +210,13 @@ mod tests { 12.5, 1.0, ) + .expect("native query execution should not fail") .is_none() ); let (_, on_grid_result) = engine .handle_range_query_promql("rate(http_requests_total[2s])".to_string(), 11.0, 12.0, 1.0) + .expect("native query execution should not fail") .expect("on-grid Sliding counter query should use native execution"); let on_grid_samples = matrix_values(on_grid_result) .into_iter() @@ -584,7 +587,9 @@ mod tests { ); let query = "sum by (host, event) (count_over_time(event_frequency[2s]))"; - let result = engine.handle_range_query_promql(query.to_string(), 5.0, 5.5, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 5.0, 5.5, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -619,12 +624,14 @@ mod tests { 1_000, ); - let result = engine.handle_range_query_promql( - "count(event_frequency) by (host, event)".to_string(), - 1.0, - 1.5, - 1.0, - ); + let result = engine + .handle_range_query_promql( + "count(event_frequency) by (host, event)".to_string(), + 1.0, + 1.5, + 1.0, + ) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); assert!(labels_have_sample_at( @@ -668,12 +675,14 @@ mod tests { 1_000, ); - let result = engine.handle_range_query_promql( - "count(event_frequency) by (host, event)".to_string(), - 1.0, - 2.0, - 1.0, - ); + let result = engine + .handle_range_query_promql( + "count(event_frequency) by (host, event)".to_string(), + 1.0, + 2.0, + 1.0, + ) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); assert!(labels_have_sample_at( @@ -760,7 +769,9 @@ mod tests { ); let query = "sum by (host, event) (count_over_time(event_frequency[2s]))"; - let result = engine.handle_range_query_promql(query.to_string(), 3.0, 4.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 3.0, 4.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -816,7 +827,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); assert!( !matrix_values(qr).is_empty(), @@ -883,7 +896,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -970,7 +985,9 @@ mod tests { WindowType::Sliding, ); - let result = engine.handle_range_query_promql(query.to_string(), 1000.0, 1000.5, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1000.0, 1000.5, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); assert_eq!(elements.len(), 1, "expected one merged series for host-a"); @@ -1016,7 +1033,9 @@ mod tests { WindowType::Sliding, ); - let result = engine.handle_range_query_promql(query.to_string(), 1000.0, 1000.5, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1000.0, 1000.5, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); assert_eq!(elements.len(), 1, "expected one merged series for host-a"); @@ -1052,7 +1071,9 @@ mod tests { WindowType::Sliding, ); - let result = engine.handle_range_query_promql(query.to_string(), 1000.0, 1000.5, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1000.0, 1000.5, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); assert_eq!(elements.len(), 1); @@ -1098,7 +1119,9 @@ mod tests { WindowType::Sliding, ); - let result = engine.handle_range_query_promql(query.to_string(), 2.0, 2.5, 2.0); + let result = engine + .handle_range_query_promql(query.to_string(), 2.0, 2.5, 2.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); assert_eq!(elements.len(), 1, "expected one series for host-a"); @@ -1146,7 +1169,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); assert!( @@ -1203,7 +1228,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -1280,7 +1307,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -1362,7 +1391,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event) * 1"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -1460,7 +1491,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 5.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 5.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -1522,7 +1555,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event) * 1"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -1619,7 +1654,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 3.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 3.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -1696,7 +1733,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -1821,7 +1860,9 @@ mod tests { ); let query = "count(event_frequency) by (region, host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -1968,7 +2009,9 @@ mod tests { ); let query = "count(event_frequency) by (region, host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -2078,7 +2121,9 @@ mod tests { ); let query = "count(event_frequency) by (region, host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect( "range query should succeed by skipping the value-less region=orphan group, \ not fail the entire query because of it", @@ -2206,7 +2251,9 @@ mod tests { let engine = create_range_engine_self_keyed("transfer_events", "srcip", vec![(1000, sketch)], query); - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -2305,7 +2352,9 @@ mod tests { ); let query = "count(event_frequency) by (host, event)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); @@ -2435,7 +2484,9 @@ mod tests { QueryLanguage::promql, ); - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("range query failed"); let elements = matrix_values(qr); assert_eq!( diff --git a/asap-query-engine/src/tests/query_equivalence_tests.rs b/asap-query-engine/src/tests/query_equivalence_tests.rs index 01d9966e..60f115c3 100644 --- a/asap-query-engine/src/tests/query_equivalence_tests.rs +++ b/asap-query-engine/src/tests/query_equivalence_tests.rs @@ -14,7 +14,7 @@ use crate::tests::test_utilities::{assert_execution_context_equivalent, TestConf use std::collections::HashMap; use std::sync::Arc; -/// Minimal no-op store that panics if queried +/// Minimal no-op store that fails if queried. /// /// This ensures that tests don't accidentally query the store. /// Context building should not require store access. @@ -28,7 +28,7 @@ impl Store for NoOpStore { _start_timestamp: u64, _end_timestamp: u64, ) -> Result> { - panic!("NoOpStore: query_precomputed_output should not be called in equivalence tests"); + Err(std::io::Error::other("test store read failure").into()) } fn query_precomputed_output_exact( @@ -38,9 +38,7 @@ impl Store for NoOpStore { _exact_start: u64, _exact_end: u64, ) -> Result> { - panic!( - "NoOpStore: query_precomputed_output_exact should not be called in equivalence tests" - ); + Err(std::io::Error::other("test store read failure").into()) } fn query_precomputed_output_exact_batch( @@ -49,9 +47,7 @@ impl Store for NoOpStore { _aggregation_id: u64, _windows: &[crate::stores::TimestampRange], ) -> Result> { - panic!( - "NoOpStore: query_precomputed_output_exact_batch should not be called in equivalence tests" - ); + Err(std::io::Error::other("test store read failure").into()) } fn insert_precomputed_output( @@ -89,6 +85,32 @@ impl Store for NoOpStore { mod tests { use super::*; + #[test] + fn public_promql_method_returns_store_failures() { + let query = "sum(requests)"; + let (inference_config, streaming_config) = TestConfigBuilder::new("requests") + .add_spatial_query(query, "SELECT SUM(value) FROM requests", 1) + .build(); + let engine = SimpleEngine::new( + Arc::new(NoOpStore), + inference_config, + streaming_config, + 1_000, + QueryLanguage::promql, + ); + + assert!(engine.handle_query_promql(query.to_string(), 1.0).is_err()); + assert!(engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .is_err()); + assert!(engine + .handle_query_promql(format!("{query} + 1"), 1.0) + .is_err()); + assert!(engine + .handle_range_query_promql(format!("{query} + 1"), 1.0, 2.0, 1.0) + .is_err()); + } + #[test] fn sql_executes_wider_sliding_window_exact_cover() { use crate::data_model::{CleanupPolicy, KeyByLabelValues, PrecomputedOutput}; diff --git a/asap-query-engine/src/tests/range_query_arithmetic_tests.rs b/asap-query-engine/src/tests/range_query_arithmetic_tests.rs index 1a912940..a18c6497 100644 --- a/asap-query-engine/src/tests/range_query_arithmetic_tests.rs +++ b/asap-query-engine/src/tests/range_query_arithmetic_tests.rs @@ -158,7 +158,9 @@ mod tests { ); let query = "sum(errors_total) by (host) / sum(requests_total) by (host)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result for range vector-vector query"); let elements = matrix_values(qr); assert_eq!(elements.len(), 1, "Expected 1 series (host-a)"); @@ -191,7 +193,9 @@ mod tests { ); let query = "sum(metric_a) by (host) + sum(metric_b) by (host)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result"); let elements = matrix_values(qr); assert_eq!(elements.len(), 1); @@ -219,7 +223,9 @@ mod tests { ); let query = "sum(metric_a) by (host) * 100"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result for scalar range query"); let elements = matrix_values(qr); assert_eq!(elements.len(), 1); @@ -246,7 +252,9 @@ mod tests { ); let query = "1 - sum(metric_a) by (host)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result for scalar-left range query"); let elements = matrix_values(qr); assert_eq!(elements.len(), 1); @@ -280,7 +288,9 @@ mod tests { ); let query = "sum(metric_a) by (host) + sum(metric_b) by (region)"; - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); assert!( result.is_none(), "BUG: arms grouped by different label sets must not join, even when their \ @@ -423,6 +433,7 @@ mod tests { let (_, qr) = engine .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail") .expect("Expected result for range topk/plain query"); let elements = matrix_values(qr); @@ -543,6 +554,7 @@ mod tests { let (_, qr) = engine .handle_query_promql(query.to_string(), 1.0) + .expect("native query execution should not fail") .expect("Expected result for instant topk/topk query"); let elements = match qr { QueryResult::Vector(vector) => vector.values, @@ -566,6 +578,7 @@ mod tests { let (_, qr) = engine .handle_query_promql(query.to_string(), 1.0) + .expect("native query execution should not fail") .expect("Expected result for instant topk/plain query"); let elements = match qr { QueryResult::Vector(vector) => vector.values, @@ -600,7 +613,9 @@ mod tests { &[("host-a", 5.0), ("host-b", 200.0), ("host-c", 300.0)], ); - let result = engine.handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0); + let result = engine + .handle_range_query_promql(query.to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("Expected result for topk/topk range query"); let elements = matrix_values(qr); @@ -663,6 +678,7 @@ mod tests { let query = format!("topk(2, metric_a) {op} topk(2, metric_b)"); let (_, qr) = engine .handle_range_query_promql(query, 1.0, 2.0, 1.0) + .expect("native query execution should not fail") .unwrap_or_else(|| panic!("Expected result for operator {op}")); let elements = matrix_values(qr); @@ -693,6 +709,7 @@ mod tests { 2.0, 1.0, ) + .expect("native query execution should not fail") .expect("Expected result for power operator"); let elements = matrix_values(qr); let values: HashMap = elements @@ -713,8 +730,9 @@ mod tests { &[("host-c", 300.0), ("host-d", 200.0)], ); - let result = - engine.handle_query_promql("topk(2, metric_a) + topk(2, metric_b)".to_string(), 1.0); + let result = engine + .handle_query_promql("topk(2, metric_a) + topk(2, metric_b)".to_string(), 1.0) + .expect("native query execution should not fail"); let (_, qr) = result.expect("valid no-match join should return an empty vector"); let elements = match qr { QueryResult::Vector(vector) => vector.values, @@ -732,12 +750,14 @@ mod tests { &[("host-c", 300.0), ("host-d", 200.0)], ); - let result = engine.handle_range_query_promql( - "topk(2, metric_a) + topk(2, metric_b)".to_string(), - 1.0, - 2.0, - 1.0, - ); + let result = engine + .handle_range_query_promql( + "topk(2, metric_a) + topk(2, metric_b)".to_string(), + 1.0, + 2.0, + 1.0, + ) + .expect("native query execution should not fail"); let (_, qr) = result.expect("valid no-match join should return an empty matrix"); assert!(matrix_values(qr).is_empty()); } @@ -753,6 +773,7 @@ mod tests { let (_, qr) = engine .handle_range_query_promql("topk(2, metric_a) + 2".to_string(), 1.0, 2.0, 1.0) + .expect("native query execution should not fail") .expect("Expected result for range topk/scalar query"); let elements = matrix_values(qr); @@ -780,6 +801,7 @@ mod tests { assert!(engine .handle_query_promql("topk(2, metric_a) == topk(2, metric_b)".to_string(), 1.0,) + .expect("native query execution should not fail") .is_none()); } @@ -799,6 +821,7 @@ mod tests { 2.0, 1.0, ) + .expect("native query execution should not fail") .is_none()); } @@ -816,6 +839,7 @@ mod tests { "(topk(2, metric_a) == topk(2, metric_b)) + topk(2, metric_a)".to_string(), 1.0, ) + .expect("native query execution should not fail") .is_none()); } @@ -833,6 +857,7 @@ mod tests { "topk(2, metric_a) + on(__name__) topk(2, metric_b)".to_string(), 1.0, ) + .expect("native query execution should not fail") .is_none()); } @@ -852,6 +877,7 @@ mod tests { 2.0, 1.0, ) + .expect("native query execution should not fail") .is_none()); } @@ -870,6 +896,7 @@ mod tests { .to_string(), 1.0, ) + .expect("native query execution should not fail") .is_none()); } } diff --git a/asap-query-engine/src/tests/stage_e_instant_range_equivalence_tests.rs b/asap-query-engine/src/tests/stage_e_instant_range_equivalence_tests.rs index fa84ef0f..510ba17c 100644 --- a/asap-query-engine/src/tests/stage_e_instant_range_equivalence_tests.rs +++ b/asap-query-engine/src/tests/stage_e_instant_range_equivalence_tests.rs @@ -275,6 +275,7 @@ mod tests { let (_, instant_result) = engine .handle_query_promql(shape.query.to_string(), shape.query_time_s) + .expect("native query execution should not fail") .unwrap_or_else(|| panic!("{case_name}: instant query returned None")); let mut instant_pairs: Vec<(Vec, f64)> = match instant_result { QueryResult::Vector(v) => v @@ -297,6 +298,7 @@ mod tests { shape.query_time_s + 0.5, 1.0, ) + .expect("native query execution should not fail") .unwrap_or_else(|| panic!("{case_name}: range query returned None")); let mut range_pairs: Vec<(Vec, f64)> = match range_result { QueryResult::Matrix(m) => m @@ -516,6 +518,7 @@ mod tests { let (_, result) = engine .handle_range_query_promql(query.to_string(), 1.0, 2.5, 1.0) + .expect("native query execution should not fail") .expect("range topk query failed"); let elements = match result { QueryResult::Matrix(m) => m.values, @@ -556,6 +559,7 @@ mod tests { let (_, instant_result) = engine .handle_query_promql(query.to_string(), 1.0) + .expect("native query execution should not fail") .expect("instant topk query failed"); let mut instant_pairs: Vec<(Vec, f64)> = match instant_result { QueryResult::Vector(v) => v @@ -568,6 +572,7 @@ mod tests { let (_, range_result) = engine .handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0) + .expect("native query execution should not fail") .expect("range topk query failed"); let mut range_pairs: Vec<(Vec, f64)> = match range_result { QueryResult::Matrix(m) => m @@ -604,6 +609,7 @@ mod tests { let (_, result) = engine .handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0) + .expect("native query execution should not fail") .expect("range topk query failed"); let elements = match result { QueryResult::Matrix(m) => m.values, @@ -642,6 +648,7 @@ mod tests { let (_, result) = engine .handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0) + .expect("native query execution should not fail") .expect("range topk query failed"); let elements = match result { QueryResult::Matrix(m) => m.values, diff --git a/asap-query-engine/src/tests/window_semantics_consistency_tests.rs b/asap-query-engine/src/tests/window_semantics_consistency_tests.rs index 6b2144bc..93e06013 100644 --- a/asap-query-engine/src/tests/window_semantics_consistency_tests.rs +++ b/asap-query-engine/src/tests/window_semantics_consistency_tests.rs @@ -166,6 +166,7 @@ mod tests { let range_result = engine .handle_range_query_promql(query.to_string(), 3.0, 5.0, 1.0) + .expect("native query execution should not fail") .expect("range query failed"); let range_samples = host_a_samples(&matrix_values(range_result.1)); @@ -180,6 +181,7 @@ mod tests { { let instant_result = engine .handle_query_promql(query.to_string(), query_time_sec) + .expect("native query execution should not fail") .unwrap_or_else(|| panic!("instant query at t={query_time_sec} failed")); let instant_value = single_host_a_value(instant_result.1); assert_close( @@ -232,6 +234,7 @@ mod tests { let result = engine .handle_query_promql(query.to_string(), 6.0) + .expect("native query execution should not fail") .expect("the wider Sliding query should be accelerated"); assert_close( @@ -242,6 +245,7 @@ mod tests { let misaligned_result = engine .handle_query_promql(query.to_string(), 6.5) + .expect("native query execution should not fail") .expect("the endpoint should align down to the latest complete slide boundary"); assert_close( single_host_a_value(misaligned_result.1), @@ -283,7 +287,10 @@ mod tests { ); assert!( - engine.handle_query_promql(query.to_string(), 6.0).is_none(), + engine + .handle_query_promql(query.to_string(), 6.0) + .expect("native query execution should not fail") + .is_none(), "a partial exact cover must fall back instead of returning partial data" ); } @@ -315,6 +322,7 @@ mod tests { let result = engine .handle_range_query_promql(query.to_string(), 6.0, 7.0, 1.0) + .expect("native query execution should not fail") .expect("the wider Sliding range query should be accelerated"); assert_eq!( @@ -375,11 +383,13 @@ mod tests { let instant_result = engine .handle_query_promql(query.to_string(), 4.0) + .expect("native query execution should not fail") .expect("instant query failed"); let instant_value = single_host_a_value(instant_result.1); let range_result = engine .handle_range_query_promql(query.to_string(), 4.0, 4.5, 1.0) + .expect("native query execution should not fail") .expect("range query failed"); let range_samples = host_a_samples(&matrix_values(range_result.1)); @@ -534,6 +544,7 @@ mod tests { let (_, tumbling_qr) = tumbling_engine .handle_query_promql(query.to_string(), 5.0) + .expect("native query execution should not fail") .expect("tumbling-keys instant query failed"); let tumbling_values = vector_values(tumbling_qr); @@ -546,6 +557,7 @@ mod tests { assert!( sliding_engine .handle_query_promql(query.to_string(), 5.0) + .expect("native query execution should not fail") .is_none(), "a pre-bound SetAggregator must share the value aggregation's window grid" ); @@ -635,6 +647,7 @@ mod tests { let (_, result) = engine .handle_query_promql(query.to_string(), 6.0) + .expect("native query execution should not fail") .expect("the wider dual-population Sliding query should be accelerated"); let mut values = vector_values(result); values.sort_by(|left, right| left.0.cmp(&right.0)); @@ -742,6 +755,7 @@ mod tests { let (_, result) = engine .handle_query_promql(query.to_string(), 4.0) + .expect("native query execution should not fail") .expect("compatible Tumbling DeltaSet keys should resolve Sliding values"); let mut values = vector_values(result); values.sort_by(|left, right| left.0.cmp(&right.0)); @@ -756,6 +770,7 @@ mod tests { let (_, misaligned_result) = engine .handle_query_promql(query.to_string(), 4.5) + .expect("native query execution should not fail") .expect("misaligned evaluation should use the latest complete value grid point"); let mut misaligned_values = vector_values(misaligned_result); misaligned_values.sort_by(|left, right| left.0.cmp(&right.0)); @@ -826,6 +841,7 @@ mod tests { let (_, first_qr) = engine .handle_query_promql(query.to_string(), 2.0) + .expect("native query execution should not fail") .expect("instant query for the fully-paned window [0,2000) failed"); assert_close( single_host_a_value(first_qr), @@ -833,7 +849,9 @@ mod tests { "window [0,2000) has both required panes (5+7) and must resolve exactly", ); - let gap_result = engine.handle_query_promql(query.to_string(), 3.0); + let gap_result = engine + .handle_query_promql(query.to_string(), 3.0) + .expect("native query execution should not fail"); let gap_values = match gap_result { Some((_, qr)) => vector_values(qr), None => Vec::new(), @@ -848,6 +866,7 @@ mod tests { let (_, third_qr) = engine .handle_query_promql(query.to_string(), 5.0) + .expect("native query execution should not fail") .expect("instant query for the fully-paned window [3000,5000) failed"); assert_close( single_host_a_value(third_qr), @@ -885,6 +904,7 @@ mod tests { let (_, instant_qr) = engine .handle_query_promql(query.to_string(), 1.0) + .expect("native query execution should not fail") .expect("instant query failed"); let instant_value = single_host_a_value(instant_qr); assert_close( @@ -895,6 +915,7 @@ mod tests { let (_, range_qr) = engine .handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0) + .expect("native query execution should not fail") .expect("range query failed"); let range_samples = host_a_samples(&matrix_values(range_qr)); assert_eq!( @@ -930,7 +951,9 @@ mod tests { WindowType::Tumbling, ); - let instant_result = engine.handle_query_promql(query.to_string(), 1.0); + let instant_result = engine + .handle_query_promql(query.to_string(), 1.0) + .expect("native query execution should not fail"); let instant_values = match instant_result { Some((_, qr)) => vector_values(qr), None => Vec::new(), @@ -943,7 +966,9 @@ mod tests { for host-a, not a fabricated/partial value: got {instant_values:?}" ); - let range_result = engine.handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0); + let range_result = engine + .handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0) + .expect("native query execution should not fail"); let range_has_sample_at_1000 = match range_result { Some((_, qr)) => host_a_samples(&matrix_values(qr)) .iter() diff --git a/asap-query-engine/tests/e2e_precompute_equivalence.rs b/asap-query-engine/tests/e2e_precompute_equivalence.rs index b7ce962d..fd671b15 100644 --- a/asap-query-engine/tests/e2e_precompute_equivalence.rs +++ b/asap-query-engine/tests/e2e_precompute_equivalence.rs @@ -235,6 +235,7 @@ impl PromqlPrecomputeFixture<'_> { query_engine .handle_query_promql(self.query.to_string(), self.evaluation_time_seconds) + .expect("native query execution should not fail") .unwrap_or_else(|| panic!("precomputed query should succeed: {}", self.query)) .1 } @@ -320,6 +321,7 @@ async fn e2e_sliding_precompute_outputs_compose_a_wider_query() { let (_, result) = query_engine .handle_query_promql(query.to_string(), 10.0) + .expect("native query execution should not fail") .expect("worker-emitted Sliding windows should answer the wider query"); let QueryResult::Vector(vector) = result else { panic!("expected instant vector result"); From 7a63690e6c7fbba858ef19c3998b9d70177493fb Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Wed, 30 Sep 2026 22:38:26 -0400 Subject: [PATCH 2/2] fix(query-engine): classify local data outcomes --- .../src/drivers/query/servers/http.rs | 24 ++- .../src/engines/simple_engine/mod.rs | 157 ++++++++++++------ .../src/engines/simple_engine/promql.rs | 12 +- 3 files changed, 129 insertions(+), 64 deletions(-) diff --git a/asap-query-engine/src/drivers/query/servers/http.rs b/asap-query-engine/src/drivers/query/servers/http.rs index dd185b43..94ce7fd7 100644 --- a/asap-query-engine/src/drivers/query/servers/http.rs +++ b/asap-query-engine/src/drivers/query/servers/http.rs @@ -12,7 +12,7 @@ use std::collections::HashMap; use std::sync::Arc; use std::time::Instant; use tokio::net::TcpListener; -use tracing::{debug, info}; +use tracing::{debug, info, warn}; use crate::drivers::query::adapters::{create_http_adapter, AdapterConfig, HttpProtocolAdapter}; use crate::engines::{QueryExecutionError, SimpleEngine}; @@ -287,7 +287,16 @@ async fn process_query_request( } } } - Err(error) => format_native_execution_error(state, error).await, + Err(error) => { + let total_duration = start_time.elapsed(); + warn!(query = %parsed_request.query, error = %error, "Native query execution failed"); + info!( + "query='{}' destination=none_native_error total_latency_ms={:.2}", + parsed_request.query, + total_duration.as_secs_f64() * 1000.0 + ); + format_native_execution_error(state, error).await + } } } @@ -595,7 +604,16 @@ async fn process_range_query_request( } } } - Err(error) => format_native_execution_error(state, error).await, + Err(error) => { + let total_duration = start_time.elapsed(); + warn!(query = %parsed_request.query, error = %error, "Native range query execution failed"); + info!( + "query='{}' destination=none_native_error total_latency_ms={:.2}", + parsed_request.query, + total_duration.as_secs_f64() * 1000.0 + ); + format_native_execution_error(state, error).await + } } } diff --git a/asap-query-engine/src/engines/simple_engine/mod.rs b/asap-query-engine/src/engines/simple_engine/mod.rs index c5c2c032..cca337a7 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -10,7 +10,9 @@ use crate::data_model::{ }; use crate::engines::query_plan::{PlanOptions, QueryPlan}; use crate::engines::query_result::{InstantVectorElement, QueryResult}; -use crate::engines::sliding_window_composition::{plan_exact_cover, SlidingWindowSpec}; +use crate::engines::sliding_window_composition::{ + plan_exact_cover, CompositionError, SlidingWindowSpec, +}; // use crate::stores::promsketch_store::{ // self, is_usampling_function, metrics as ps_metrics, PromSketchStore, // }; @@ -37,13 +39,6 @@ use serde_json::Value; #[allow(dead_code)] type MergedOutputsMap = HashMap, Box>; -const NO_LOCAL_DATA_ERROR_PREFIXES: [&str; 4] = [ - "No precomputed outputs found", - "No data found", - "Incomplete Sliding-window cover", - "Exact Prometheus counter bounds are unavailable for off-grid", -]; - /// Metadata extracted from a query, independent of query language #[derive(Debug, Clone)] pub struct QueryMetadata { @@ -64,6 +59,8 @@ pub struct QueryMetadata { /// A native query was accepted but could not be executed locally. #[derive(Debug, thiserror::Error)] pub enum QueryExecutionError { + #[error("no local data: {0}")] + NoLocalData(String), #[error("native query execution failed: {0}")] Native(String), } @@ -813,7 +810,7 @@ impl SimpleEngine { lookback_ms: u64, window_size_ms: u64, slide_interval_ms: u64, - ) -> Result { + ) -> Result { let spec = SlidingWindowSpec { window_size_ms, slide_interval_ms, @@ -824,10 +821,16 @@ impl SimpleEngine { let mut windows = BTreeSet::new(); for &output_timestamp in output_timestamps { let cover = plan_exact_cover(output_timestamp, lookback_ms, spec).map_err(|error| { - format!( + let message = format!( "Cannot compose Sliding aggregation {} for lookback {}ms (W={}ms, S={}ms): {:?}", params.aggregation_id, lookback_ms, window_size_ms, slide_interval_ms, error - ) + ); + match error { + CompositionError::LookbackBeforeEpoch { .. } => { + QueryExecutionError::NoLocalData(message) + } + _ => QueryExecutionError::Native(message), + } })?; windows.extend(cover.windows); } @@ -850,13 +853,13 @@ impl SimpleEngine { .store .query_precomputed_output_exact_batch(¶ms.metric, params.aggregation_id, &windows) .map_err(|error| { - format!( + QueryExecutionError::Native(format!( "Error querying store for metric {}, agg {}, {} exact Sliding windows: {}", params.metric, params.aggregation_id, windows.len(), error - ) + )) })?; // An instant query has one all-or-nothing cover. Range queries defer @@ -868,7 +871,7 @@ impl SimpleEngine { let found: BTreeSet<_> = buckets.iter().map(|(range, _)| *range).collect(); if found != required { let missing: Vec<_> = required.difference(&found).copied().collect(); - return Err(format!( + return Err(QueryExecutionError::NoLocalData(format!( "Incomplete Sliding-window cover for metric {}, agg {}, group {:?}: \ requested {} exact windows, missing {:?}", params.metric, @@ -876,7 +879,7 @@ impl SimpleEngine { group_key, windows.len(), missing - )); + ))); } } } @@ -1343,14 +1346,24 @@ impl SimpleEngine { enable_topk_limiting: bool, enable_topk_formatting: bool, ) -> Result, String> { + self.execute_query_pipeline_result(context, enable_topk_limiting, enable_topk_formatting) + .map_err(|error| error.to_string()) + } + + fn execute_query_pipeline_result( + &self, + context: &QueryExecutionContext, + enable_topk_limiting: bool, + enable_topk_formatting: bool, + ) -> Result, QueryExecutionError> { let query_time = context.query_time; let range_context = self .build_instant_range_context(context.clone(), query_time) .ok_or_else(|| { - format!( + QueryExecutionError::NoLocalData(format!( "Failed to build instant-as-range context for metric: {}", context.metric - ) + )) })?; let range_results = self.execute_observed_range_query_pipeline( @@ -1562,7 +1575,7 @@ impl SimpleEngine { enable_topk_limiting: bool, enable_topk_formatting: bool, ) -> Result, QueryExecutionError> { - let Some(results) = Self::classify_native_execution(self.execute_query_pipeline( + let Some(results) = Self::map_local_execution_outcome(self.execute_query_pipeline_result( &context, enable_topk_limiting, enable_topk_formatting, @@ -1576,19 +1589,13 @@ impl SimpleEngine { ))) } - fn classify_native_execution( - result: Result, + fn map_local_execution_outcome( + result: Result, ) -> Result, QueryExecutionError> { match result { Ok(value) => Ok(Some(value)), - Err(error) - if NO_LOCAL_DATA_ERROR_PREFIXES - .iter() - .any(|prefix| error.starts_with(prefix)) => - { - Ok(None) - } - Err(error) => Err(QueryExecutionError::Native(error)), + Err(QueryExecutionError::NoLocalData(_)) => Ok(None), + Err(error) => Err(error), } } @@ -2119,14 +2126,15 @@ impl SimpleEngine { context: &RangeQueryExecutionContext, enable_topk_limiting: bool, enable_topk_formatting: bool, - ) -> Result, String> { + ) -> Result, QueryExecutionError> { let plan = QueryPlan::compile_range( context, PlanOptions { limit_topk: enable_topk_limiting, format_output: enable_topk_formatting, }, - )?; + ) + .map_err(QueryExecutionError::Native)?; debug!(plan = %plan.explain(), "Compiled native query plan"); self.execute_range_query_pipeline(context, enable_topk_limiting, enable_topk_formatting) } @@ -2136,7 +2144,7 @@ impl SimpleEngine { context: &RangeQueryExecutionContext, enable_topk_limiting: bool, enable_topk_formatting: bool, - ) -> Result, String> { + ) -> Result, QueryExecutionError> { use crate::engines::query_result::RangeVectorElement; use crate::engines::window_merger::create_window_merger; @@ -2152,11 +2160,11 @@ impl SimpleEngine { .iter() .find(|&×tamp| !timestamp.is_multiple_of(context.tumbling_window_ms)) { - return Err(format!( + 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 - )); + ))); } } @@ -2174,11 +2182,15 @@ impl SimpleEngine { context.tumbling_window_ms, )? } else { - self.execute_store_query(&context.base.store_plan.values_query)? + self.execute_store_query(&context.base.store_plan.values_query) + .map_err(QueryExecutionError::Native)? }; if all_data.is_empty() { - return Err(format!("No data found for metric: {}", context.base.metric)); + return Err(QueryExecutionError::NoLocalData(format!( + "No data found for metric: {}", + context.base.metric + ))); } debug!( @@ -2197,22 +2209,31 @@ impl SimpleEngine { // 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) => Some( - self.execute_sliding_cover_query( + Some(keys_query) if context.keys_window_type == Some(WindowType::Sliding) => { + Some(self.execute_sliding_cover_query( keys_query, &context.output_timestamps, - context - .keys_lookback_ms - .ok_or("Sliding keys query is missing its lookback")?, - context - .keys_window_size_ms - .ok_or("Sliding keys query is missing its window size")?, - context - .keys_tumbling_window_ms - .ok_or("Sliding keys query is missing its slide interval")?, - )?, + context.keys_lookback_ms.ok_or_else(|| { + QueryExecutionError::Native( + "Sliding keys query is missing its lookback".to_string(), + ) + })?, + context.keys_window_size_ms.ok_or_else(|| { + QueryExecutionError::Native( + "Sliding keys query is missing its window size".to_string(), + ) + })?, + context.keys_tumbling_window_ms.ok_or_else(|| { + QueryExecutionError::Native( + "Sliding keys query is missing its slide interval".to_string(), + ) + })?, + )?) + } + Some(keys_query) => Some( + self.execute_store_query(keys_query) + .map_err(QueryExecutionError::Native)?, ), - Some(keys_query) => Some(self.execute_store_query(keys_query)?), None => None, }; @@ -2402,7 +2423,10 @@ impl SimpleEngine { 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)?) + Some( + Self::parse_topk_limit(&context.base.metadata.query_kwargs) + .map_err(QueryExecutionError::Native)?, + ) } else { None }; @@ -2420,14 +2444,24 @@ impl SimpleEngine { // (#581). One loop shape for topk and non-topk alike, rather than // maintaining two. for ¤t_time in &context.output_timestamps { - let current_time_i64 = i64::try_from(current_time) - .map_err(|_| "Output timestamp exceeds signed timestamp range".to_string())?; - let query_range_ms = i64::try_from(context.query_range_ms) - .map_err(|_| "Query range exceeds signed timestamp range".to_string())?; + 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(|_| { + QueryExecutionError::Native( + "Query range exceeds signed timestamp range".to_string(), + ) + })?; let query_bounds = QueryBounds::new( current_time_i64 .checked_sub(query_range_ms) - .ok_or("Query range underflows timestamp range".to_string())?, + .ok_or_else(|| { + QueryExecutionError::Native( + "Query range underflows timestamp range".to_string(), + ) + })?, current_time_i64, ); // This timestamp's (key, value) pairs from every group, each @@ -2682,7 +2716,7 @@ impl SimpleEngine { #[cfg(test)] mod topk_metadata_tests { - use super::SimpleEngine; + use super::{QueryExecutionError, SimpleEngine}; use std::collections::HashMap; #[test] @@ -2700,6 +2734,19 @@ mod topk_metadata_tests { Ok(3) ); } + + #[test] + fn local_execution_outcomes_are_classified_by_variant() { + let no_local_data = SimpleEngine::map_local_execution_outcome::<()>(Err( + QueryExecutionError::NoLocalData("an arbitrary diagnostic".to_string()), + )); + assert!(matches!(no_local_data, Ok(None))); + + let native_error = SimpleEngine::map_local_execution_outcome::<()>(Err( + QueryExecutionError::Native("No data found, but the operation failed".to_string()), + )); + assert!(matches!(native_error, Err(QueryExecutionError::Native(_)))); + } } #[cfg(test)] diff --git a/asap-query-engine/src/engines/simple_engine/promql.rs b/asap-query-engine/src/engines/simple_engine/promql.rs index a00e1c25..65632b90 100644 --- a/asap-query-engine/src/engines/simple_engine/promql.rs +++ b/asap-query-engine/src/engines/simple_engine/promql.rs @@ -496,8 +496,8 @@ impl SimpleEngine { let Some((ctx, label_names)) = self.resolve_arm_leaf_context(arm_ast, time) else { return Ok(None); }; - let Some(results) = Self::classify_native_execution( - self.execute_query_pipeline(&ctx, true, false), + let Some(results) = Self::map_local_execution_outcome( + self.execute_query_pipeline_result(&ctx, true, false), )? else { return Ok(None); @@ -758,7 +758,7 @@ impl SimpleEngine { // Binary arms need Topk limiting, but must remain in the // unformatted intermediate label representation until after the // arithmetic operation. - let Some(results) = Self::classify_native_execution( + let Some(results) = Self::map_local_execution_outcome( self.execute_observed_range_query_pipeline(&ctx, true, false), )? else { @@ -801,13 +801,13 @@ impl SimpleEngine { return Ok(None); } // Binary arms need Topk limiting, but not final presentation formatting. - let Some(lhs_results) = Self::classify_native_execution( + let Some(lhs_results) = Self::map_local_execution_outcome( self.execute_observed_range_query_pipeline(&lhs_ctx, true, false), )? else { return Ok(None); }; - let Some(rhs_results) = Self::classify_native_execution( + let Some(rhs_results) = Self::map_local_execution_outcome( self.execute_observed_range_query_pipeline(&rhs_ctx, true, false), )? else { @@ -1383,7 +1383,7 @@ impl SimpleEngine { // Execute range query pipeline. (true, true): self-gated, same as // instant's handle_query_promql -- both flags are no-ops unless this // query's statistic is Topk. - let Some(results): Option> = Self::classify_native_execution( + let Some(results): Option> = Self::map_local_execution_outcome( self.execute_observed_range_query_pipeline(&context, true, true), )? else {