Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
39 commits
Select commit Hold shift + click to select a range
8776a4a
test(query-engine): characterize native range leaf execution
milindsrivastava1997 Sep 26, 2026
b263b0a
refactor(query-engine): validate native query DAG dependencies
milindsrivastava1997 Sep 26, 2026
f2d930c
refactor(query-engine): isolate range query store reads
milindsrivastava1997 Sep 26, 2026
6926597
refactor(query-engine): execute ranges from prepared reads
milindsrivastava1997 Sep 26, 2026
8bf319b
refactor(query-engine): execute native range DAGs
milindsrivastava1997 Sep 26, 2026
f828a4e
test(query-engine): cover native DAG execution order
milindsrivastava1997 Sep 26, 2026
5ae29e3
refactor(query-engine): expose native PromQL execution errors
milindsrivastava1997 Sep 26, 2026
1c7ab85
fix(query-engine): return native execution failures locally
milindsrivastava1997 Sep 26, 2026
2bc840a
fix(query-engine): return range execution failures locally
milindsrivastava1997 Sep 26, 2026
48b77c1
refactor(query-engine): stage native range DAG execution
milindsrivastava1997 Sep 26, 2026
869f2d4
test(query-engine): cover native DAG differential cases
milindsrivastava1997 Sep 27, 2026
1ddffb4
test(query-engine): align off-grid rate fixture
milindsrivastava1997 Sep 27, 2026
236b4f7
test(query-engine): keep differential cases leaf-only
milindsrivastava1997 Sep 27, 2026
882677d
test(query-engine): keep temporal fixture planner-compatible
milindsrivastava1997 Sep 27, 2026
5733afb
test(query-engine): restore temporal differential fixture
milindsrivastava1997 Sep 27, 2026
22ddb15
test(query-engine): add temporal DAG differential cases
milindsrivastava1997 Sep 27, 2026
fc94ded
test(query-engine): cover off-grid rate fallback
milindsrivastava1997 Sep 27, 2026
b27a152
test(query-engine): isolate off-grid rate differential
milindsrivastava1997 Sep 27, 2026
0d58d12
test(query-engine): add native DAG differential suites
milindsrivastava1997 Sep 27, 2026
8385920
test(query-engine): add legacy range execution selector
milindsrivastava1997 Sep 27, 2026
d012716
test(query-engine): compare legacy and DAG range execution
milindsrivastava1997 Sep 27, 2026
7ec25cc
test(query-engine): characterize sparse legacy DAG parity
milindsrivastava1997 Sep 27, 2026
933a2af
test(query-engine): characterize keyed legacy DAG parity
milindsrivastava1997 Sep 27, 2026
17f506d
test(query-engine): characterize topk legacy DAG parity
milindsrivastava1997 Sep 27, 2026
36abea3
test(query-engine): cover malformed native plan errors
milindsrivastava1997 Sep 28, 2026
bce5cc2
test(query-engine): cover native store failures
milindsrivastava1997 Sep 28, 2026
77d6b2d
test(query-engine): characterize topk tie parity
milindsrivastava1997 Sep 28, 2026
a2364e7
test(query-engine): compare DeltaSet replay DAG execution
milindsrivastava1997 Sep 28, 2026
98c9744
test(query-engine): cover local native execution failures
milindsrivastava1997 Sep 28, 2026
fc72694
docs(query-engine): record native DAG follow-up work
milindsrivastava1997 Sep 28, 2026
87eb352
fix(compliance): isolate differential compose runs
milindsrivastava1997 Sep 28, 2026
d10a1cf
fix(compliance): seed differential data in the past
milindsrivastava1997 Sep 28, 2026
691945b
feat(query-engine): trace native DAG execution
milindsrivastava1997 Sep 28, 2026
171d516
fix(query-engine): preserve grouped quantile labels
milindsrivastava1997 Sep 29, 2026
5bed5d4
refactor(query-engine): preserve typed DAG execution errors
milindsrivastava1997 Oct 1, 2026
5886eb4
test(compliance): drop removed sparse fixture
milindsrivastava1997 Oct 1, 2026
4e76dfa
test(compliance): isolate Docker ports per matrix case
milindsrivastava1997 Oct 1, 2026
e03318e
docs(query-engine): record native DAG matrix results
milindsrivastava1997 Oct 1, 2026
1fafe7f
docs(query-engine): correct DAG matrix classification
milindsrivastava1997 Oct 3, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions asap-query-engine/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -97,3 +97,4 @@ jemalloc = ["dep:tikv-jemallocator"]
lock_profiling = []
# Enable extra debugging output
extra_debugging = []
native_query_legacy_test_support = []
2 changes: 2 additions & 0 deletions asap-query-engine/src/engines/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
196 changes: 192 additions & 4 deletions asap-query-engine/src/engines/query_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
use crate::engines::simple_engine::{RangeQueryExecutionContext, StoreQueryParams};
use asap_types::enums::WindowType;
use promql_utilities::query_logics::enums::Statistic;
use tracing::debug;

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct NodeId(usize);
Expand Down Expand Up @@ -42,6 +43,7 @@ pub(crate) enum QueryPlanNode {
Format {
input: NodeId,
include_metric_name: bool,
metric: String,
},
}

Expand All @@ -57,6 +59,35 @@ 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<Self::Output, Self::Error>;
}

#[derive(Debug)]
pub(crate) enum QueryPlanExecutionError<E> {
InvalidPlan(String),
Node { id: NodeId, source: E },
}

impl<E: std::fmt::Display> std::fmt::Display for QueryPlanExecutionError<E> {
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 {
pub(crate) fn compile_range(
context: &RangeQueryExecutionContext,
Expand Down Expand Up @@ -121,11 +152,76 @@ impl QueryPlan {
&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<R: QueryPlanRuntime>(
&self,
runtime: &R,
) -> Result<R::Output, QueryPlanExecutionError<R::Error>> {
self.validate()
.map_err(QueryPlanExecutionError::InvalidPlan)?;
let mut outputs: Vec<R::Output> = 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::<Vec<_>>();
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<QueryPlanNode>, node: QueryPlanNode) -> NodeId {
Expand Down Expand Up @@ -195,8 +291,8 @@ impl QueryPlan {
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::Format { input, include_metric_name, metric } => format!(
"n{index} Format(n{}, include_metric_name={include_metric_name}) metric={metric}", input.0
),
};
lines.push(line);
Expand All @@ -206,13 +302,43 @@ 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<NodeId> {
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::*;
use crate::data_model::AggregationIdInfo;
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 {
Expand Down Expand Up @@ -347,4 +473,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<Vec<usize>>);

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::Output, Self::Error> {
self.0.borrow_mut().push(id.0);
Ok(1 + inputs.iter().sum::<usize>())
}
}

#[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]);
}
}
Loading
Loading