Skip to content

refactor(query-engine): make native DAG execution self-contained - #788

Open
milindsrivastava1997 wants to merge 2 commits into
mainfrom
754-refactorquery-engine-make-native-dag-self-contained-at-execution
Open

milindsrivastava1997 wants to merge 2 commits into
mainfrom
754-refactorquery-engine-make-native-dag-self-contained-at-execution

Conversation

@milindsrivastava1997

@milindsrivastava1997 milindsrivastava1997 commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

Why

The native range-query DAG described reads and operators, but execution still depended on the original range context. That made the DAG incomplete as an execution artifact and obscured which node owned each behavior.

Changes

  • Make each StoreRead node own its role and concrete read strategy, including exact-cover parameters for sliding windows.
  • Move range estimation and top-k label-order inputs into typed DAG configuration, so NativePlanRuntime only needs the engine.
  • Preserve value-versus-keys empty-read behavior and existing off-grid counter fallback semantics.
  • Keep the feature-gated legacy pipeline as the differential oracle; it is no longer used by production DAG execution.
  • Add plan coverage for explicit read roles and canonical bucket-derived sliding lookbacks.

Verification

  • Repository hooks passed: formatting, cargo check, Clippy, and tests.
  • Full default and legacy-feature Rust test suites passed before the final review fixes; focused native range and plan tests passed afterward.
  • Docker differential matrix was run. The DAG-focused aggregations, topk-ordering, and quantiles suites passed; temporal, off-grid-rate, aggregations, and olly-bench retain their documented pre-existing failures.

@milindsrivastava1997

Copy link
Copy Markdown
Contributor Author

Code review

The new regression test passes without running either code path.

asap-query-engine/src/tests/native_range_query_tests.rs:712 (also line 716): range_query_tumbling_values_sliding_keys_dag_matches_legacy calls handle_range_query_promql(query, 2.0, 2.0, 1.0). Range setup rejects these arguments for two reasons:

  • Start equals end: validate_range_query_params (simple_engine/mod.rs:2203) requires start < end.
  • Step is off-grid: a 1000ms step isn't a multiple of the 2000ms tumbling window, so finish_range_context returns None.

Both calls therefore return Ok(None) and the test compares None == None. The case this PR adds (Tumbling values with Sliding keys) has no working test.

A second gap: map_local_execution_outcome also turns NoLocalData into Ok(None). So even with valid arguments, an empty read on both sides would still pass.

Fix: use a valid range such as (2.0, 4.0, 2.0), and assert both results are Some and non-empty before comparing them.

The rest of the diff matches the old code path:

  • The value lookback uses lookback_bucket_count * tumbling_window_ms, the same as read_range_query_inputs.
  • The keys_window_type fallback never fires, because the key window type is always set whenever keys_query is Some.
  • Nodes run in index order, so the empty-values check still happens before the key read.
  • reject_off_grid_sliding_counter_query now runs twice per query, which does no harm.

🤖 Generated with Claude Code

@milindsrivastava1997
milindsrivastava1997 force-pushed the 754-refactorquery-engine-make-native-dag-self-contained-at-execution branch from 78bb33b to 337ae1f Compare October 5, 2026 17:02

@milindsrivastava1997 milindsrivastava1997 left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed against #754's intent. All acceptance criteria are met structurally: NativePlanRuntime holds only engine, StoreRead carries its role and strategy, ComposeWindows was renamed to the honest PrepareBuckets, and LimitTopK owns row_label_order.

Two gaps against the issue's stated intent:

  1. Spec shape. #754 says: "Do not copy the existing context wholesale into a node: extract cohesive plan-owned specs with explicit invariants." RangeEstimateSpec is close to a field-for-field copy of the context, and the key-window invariant is enforced by unwrap_or in the compiler and .expect in the estimator rather than by the type. See the inline comments.
  2. Suggested tests. #754 asks for plan-execution tests (execute a compiled plan with no context; value plus separate keys with distinct read specs). The tests added here only check explain() strings.

The remaining comments are smaller efficiency and test-precision points.

);
let values = Self::push_prepare_buckets(&mut nodes, values_read);
let keys = context.base.store_plan.keys_query.as_ref().map(|query| {
let keys_lookback_ms = context.keys_lookback_ms.unwrap_or(context.query_range_ms);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Silent defaults for required key-window config. keys_lookback_ms.unwrap_or(query_range_ms), keys_window_size_ms.unwrap_or(window_size_ms), keys_tumbling_window_ms.unwrap_or(tumbling_window_ms) and (below) keys_window_type.unwrap_or(window_type) replace the explicit errors the legacy reader returned ("Sliding keys query is missing its lookback", etc.).

Both production builders (simple_engine/mod.rs:870-920, promql.rs:655-702) set all four keys_* fields together, so this can't be reached today. But if a future builder leaves one unset, the DAG reads the key aggregation with the value aggregation's W/S/lookback instead of failing. It may also pick SlidingExactCover where legacy would have used WindowGrid, and estimate_range_query would then panic on keys_lookback_ms.expect(...). This breaks the "fail loud" rule (code-design-review §1: "config resolution that silently picks a default when a required value is missing").

Suggested fix: see the comment on RangeEstimateSpec. A single Option<KeyWindowSpec> built once at compile time removes these fallbacks and the downstream .expects.

}

#[derive(Debug, Clone)]
pub(crate) struct RangeEstimateSpec {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

#754: "Do not copy the existing context wholesale into a node: extract cohesive plan-owned specs with explicit invariants."

This struct is close to a field-for-field copy of RangeQueryExecutionContext, including four parallel keys_* Options that must be all-Some or all-None. Nothing enforces that, so the estimator .expect()s each one separately and the compiler falls back with unwrap_or.

Suggest a single keys: Option<KeyWindowSpec { window_type, window_size_ms, lookback_ms, bucket_step_ms }> (and the same for the value window). Build it once in compile_range, failing if any field is missing, and use it for both the key StoreRead strategy and the estimator. That makes the half-set state unrepresentable, and the read and estimate can't disagree.

query_time_aggregations: &[QueryTimeAggregation],
) -> Result<Self, String> {
let mut nodes = Vec::new();
let values_lookback_ms =

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lookback_bucket_count * tumbling_window_ms is now computed here, in estimate_range_query (mod.rs:2721), and in legacy read_range_query_inputs. The read strategy and the estimator each hold their own copy of lookback, window size and step.

If one copy changes (e.g. the bucket-lookback rounding this PR just fixed), the reads fetch a different cover than the estimator expects, and steps get silently skipped as "incomplete Sliding value cover". This works against #754's "each node consumes the fields that define its semantics": let the estimator take the window spec from the plan instead of recomputing it.

statistic: context.base.metadata.statistic_to_compute,
query_kwargs: context.base.metadata.query_kwargs.clone(),
output_labels: context.base.metadata.query_output_labels.clone(),
spec: estimate_spec.clone(),

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: estimate_spec can be moved here once row_label_order is taken out first. Also, output_timestamps is now cloned into each SlidingExactCover strategy and into the spec, and explain() prints the full vector up to three times. On long range queries with debug logging on, debug!(plan = %plan.explain()) gets very long. Consider summarizing the timestamps in explain() (count, first, last) or sharing them behind an Arc.

keys: keys_raw_data,
} = reads;
let lookback_ms = (context.lookback_bucket_count as u64) * context.tumbling_window_ms;
Self::reject_off_grid_sliding_counter_query(spec, statistic)?;

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

reject_off_grid_sliding_counter_query now runs twice per query: at the top of execute_observed_range_query_pipeline (L2507) and again here, after all store reads are done. The second call can't fire in production. Separately, L2507 builds a full RangeEstimateSpec::from(context) (cloning timestamps and label sets, computing topk_row_label_order) just for this check and then throws it away.

Suggest keeping one check, run early, that takes the few fields it needs (window_type, tumbling_window_ms, output_timestamps) as arguments.

*bucket_step_ms,
)?,
};
if data.is_empty() && role == StoreReadRole::Values {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This role-dependent empty-read rule (empty values → NoLocalData, empty keys allowed) and the switch from one combined read to per-node reads have no runtime-level test. The new tests only check explain() strings.

#754's suggested tests ask for exactly this: "Compile a range plan, execute it with a runtime that has no query context" and "Exercise two different read/window specifications in one plan shape (value plus separate keys) to prove each node uses plan-owned semantics." A regression that swaps roles, or reads keys with the value strategy, would pass every test in this PR.

assert!(explanation.contains("n2 StoreRead(SlidingExactCover, requests#8"));
assert!(explanation.contains("role=Values"));
assert!(explanation.contains("role=Keys"));
assert!(explanation.contains("n2 StoreRead(SlidingExactCover {"));

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This assertion doesn't check the key read's parameters (lookback, window, step). The fixture (L537) sets keys_window_type = Some(Sliding) but leaves keys_lookback_ms and the other key fields as None, which production builders never produce and legacy execution rejected. So the test passes on the silent unwrap_or fallbacks and locks in that invalid configuration. Suggest filling in all key fields in the fixture with values distinct from the value window, and asserting them in the key StoreRead.

@milindsrivastava1997 milindsrivastava1997 left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Automated code review: 8 findings, posted inline. They were not checked in a separate verification pass.

);
let values = Self::push_prepare_buckets(&mut nodes, values_read);
let keys = context.base.store_plan.keys_query.as_ref().map(|query| {
let keys_lookback_ms = context.keys_lookback_ms.unwrap_or(context.query_range_ms);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Missing keys settings now fall back silently. The keys StoreRead resolves keys_lookback_ms, keys_window_size_ms and keys_tumbling_window_ms with unwrap_or fallbacks to the values query's settings. The removed read_range_query_inputs returned explicit errors in these cases.

Example: with keys_window_type=Some(Sliding) and keys_lookback_ms=None, the old code failed with "Sliding keys query is missing its lookback". Now the read plans a cover from the values' query_range_ms, window size and step, so it fetches the wrong windows. estimate_range_query then panics on keys_lookback_ms.expect(...). These fallbacks used to sit only in the explain-only ComposeWindows node; now they decide real store reads. This breaks the fail-loud rule in the code-design-review checklist.

context.keys_window_type.unwrap_or(context.window_type),
StoreReadRole::Keys,
StoreReadStrategy::for_window(
context.keys_window_type.unwrap_or(context.window_type),

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The DAG and the legacy oracle choose different keys read strategies. Here the strategy comes from keys_window_type.unwrap_or(context.window_type); the legacy oracle checks keys_window_type == Some(Sliding).

When values are Sliding, keys_query is Some and keys_window_type is None, the DAG does a SlidingExactCover read and the oracle does a WindowGrid read. They read different buckets, so the differential oracle no longer checks the same logic. Both should derive the strategy from one shared function.

pub tumbling_window_ms: u64,
pub window_type: WindowType,
pub window_size_ms: u64,
pub keys_window_type: Option<WindowType>,

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The keys-window settings are stored twice. RangeEstimateSpec copies keys_window_type, keys_window_size_ms, keys_lookback_ms, keys_tumbling_window_ms and output_timestamps, which the keys StoreRead strategy already holds.

The read resolves these with unwrap_or defaults, while the estimator calls .expect() on its own copies. If one side changes and the other does not, the read fetches one set of windows and the estimator looks for buckets from another, giving silently incomplete covers.

keys: keys_raw_data,
} = reads;
let lookback_ms = (context.lookback_bucket_count as u64) * context.tumbling_window_ms;
Self::reject_off_grid_sliding_counter_query(spec, statistic)?;

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The off-grid check runs twice. reject_off_grid_sliding_counter_query runs at the top of execute_observed_range_query_pipeline and again inside estimate_range_query. The first call also builds a whole RangeEstimateSpec (around line 2507) only to run the check, then discards it.

Suggest keeping only the early check, which runs before any store reads, and passing it (window_type, tumbling_window_ms, &output_timestamps, statistic) directly.

statistic: context.base.metadata.statistic_to_compute,
query_kwargs: context.base.metadata.query_kwargs.clone(),
output_labels: context.base.metadata.query_output_labels.clone(),
spec: estimate_spec.clone(),

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The timestamp and row-order lists are copied several times. compile_range copies output_timestamps into every Sliding StoreRead strategy and into the Estimate spec, and copies row_label_order into both the spec and LimitTopK. estimate_range_query then clones spec.row_label_order again (~line 2903), where a borrow would do.

A query with 10k steps ends up holding three copies of the timestamp vector per plan. Sharing them through an Arc<[u64]> or borrowing avoids this.

kwargs.sort_unstable_by_key(|(key, _)| *key);
format!("n{index} Estimate(n{}, {statistic}, {kwargs:?})", input.0)
format!(
"n{index} Estimate(n{}, {statistic}, {kwargs:?}, outputs={:?})",

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

explain() prints the full timestamp list. It prints all of output_timestamps for each Sliding StoreRead (through the Debug output of SlidingExactCover) and again for Estimate, and the native path logs the plan at debug level.

For a query with thousands of steps, one debug log line holds the list up to three times. Suggest printing a summary instead: the count plus the first and last timestamp.

window_size_ms,
bucket_step_ms,
} => self.engine.execute_sliding_cover_query(
query,

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A debug log was dropped from production builds. The "Range query: fetched N keys, M total buckets" log from read_range_query_inputs now exists only under native_query_legacy_test_support.

Production debug traces therefore no longer show how many groups and buckets a range read returned, which makes empty-cover and partial-cover problems harder to diagnose. Worth adding it back on the new read path.

assert!(explanation.contains("outputs=[1000, 2000, 3000]"));
}

#[test]

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The lookback test is brittle, and the empty-read behavior is untested. sliding_value_read_uses_bucket_lookback checks the lookback only by substring-matching the Debug output of explain(), so a formatting change can break or pass it for unrelated reasons.

No native-path test checks that an empty Keys read succeeds while an empty Values read returns NoLocalData. The PR description says this behavior is preserved, but a regression in the role == StoreReadRole::Values check would go uncaught.

@milindsrivastava1997 milindsrivastava1997 left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Automated code review (Claude Code).

On an unchanged line, so not inline: RangeQueryExecutionContext.buckets_per_step (simple_engine/mod.rs:147) no longer has any readers. This PR removed the step_overlap_mode debug log, which was its only consumer. It's still computed in finish_range_context (step_ms / tumbling_window_ms) and set to a placeholder 1 in build_instant_range_context. Please either delete it or restore the debug output.

query_time_aggregations: &[asap_types::query_config::QueryTimeAggregation],
) -> Result<RangePipelineOutput, QueryExecutionError> {
Self::reject_off_grid_sliding_counter_query(context)?;
Self::reject_off_grid_sliding_counter_query(

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The off-grid Sliding counter guard is still outside the plan. reject_off_grid_sliding_counter_query reads RangeQueryExecutionContext, so the compiled DAG doesn't enforce it. A caller that runs plan.execute() directly (the context-free path this PR aims for, as in compiled_plan_executes_..._without_context) skips it. A Sliding increase/rate query with off-grid output timestamps then returns estimates from misaligned buckets instead of NoLocalData. Estimate.spec already has values.window.{window_type, bucket_step_ms, output_timestamps} and the statistic, so the check could live in the Estimate node or the plan compiler.

[NativePlanOutput::Resolved(reads)] => self
.engine
.estimate_range_query(self.context, reads.clone())
.estimate_range_query(spec, *statistic, query_kwargs, reads.clone())

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Bucket maps are cloned several times. plan.execute clones every input (outputs[input.0].clone()), and Estimate clones again here with reads.clone(). The Read, Composed and Resolved maps each get copied, and every copy stays in outputs until the root finishes. For a range query over many groups and buckets, that's about 3–4 full copies of HashMap<Option<Key>, BucketMap> held at peak. Wrapping outputs in Arc, or moving the final input instead of cloning it, avoids this.

pub(crate) enum StoreReadStrategy {
WindowGrid,
SlidingExactCover,
SlidingExactCover(WindowCompositionSpec),

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Window state is stored twice. SlidingExactCover(WindowCompositionSpec) copies spec.values.window / spec.keys.window, and LimitTopK.row_label_order copies Estimate.spec.row_label_order. Nothing keeps each pair in sync. A later edit or a hand-built plan (as in the tests) could give the read a different lookback/window/step than the one Estimate composes with. The cover fetch would then miss windows, and those steps would be silently dropped as an "incomplete Sliding cover". A single source of truth (the strategy refers to the spec's window by role) would remove this kind of mismatch.

);
Some(Self::push_prepare_buckets(&mut nodes, read))
}
(Some(_), KeyInputSpec::FromValues) => {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These mismatch arms can't be reached. (Some, FromValues) and (None, Separate) can't occur because RangeEstimateSpec::compile derives KeyInputSpec from the same keys_query. The matching error arms in estimate_range_query are unreachable too. If KeyInputSpec::Separate carried the keys StoreQueryParams (or the key read were compiled from the spec), the mismatch couldn't be represented and these branches could go.

let keys_window_buckets = if *key_accumulator_type
== AggregationType::DeltaSetAggregator
{
// #588/#606 force DeltaSetAggregator's own

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: comments that this diff re-indents still include issue numbers and review history (#588/#606, (#581 stage E.4 review), (#583) in the warn! around line 2812). AGENTS.md says: "Code comments should not include git history, or issue numbers or historical facts." Since these lines are being touched anyway, now's a good time to drop the references.

.handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0)
.expect("an empty keys read must not fail native execution");
assert!(
empty_keys_result.is_some(),

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test checks too little. It only asserts empty_keys_result.is_some() and never looks at the samples. If a regression made the empty-keys path fall back to the value groups' keys (emitting rows for host-a/evt-1), this would still pass. Asserting the exact (probably empty) result set would pin down the "empty keys produce no samples" behavior the test describes.


fn describe_read_strategy(strategy: &StoreReadStrategy) -> String {
match strategy {
StoreReadStrategy::WindowGrid => "WindowGrid".to_string(),

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

explain() doesn't show WindowGrid window parameters. WindowGrid reads print no lookback/window/step, and Estimate prints only the value outputs. A Tumbling keys read's lookback and step, which the estimator now gets only from spec.keys, never appear in the plan dump. The old ComposeWindows line printed them. Without them, a misconfigured keys window (e.g. DeltaSet keys with lookback = query end) can't be diagnosed from the "Compiled native query plan" log.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

refactor(query-engine): make native DAG self-contained at execution

1 participant