Repository navigation
refactor(query-engine): make native DAG execution self-contained - #788
milindsrivastava1997 wants to merge 2 commits into
Conversation
Code reviewThe new regression test passes without running either code path.
Both calls therefore return A second gap: Fix: use a valid range such as The rest of the diff matches the old code path:
🤖 Generated with Claude Code |
78bb33b to
337ae1f
Compare
milindsrivastava1997
left a comment
There was a problem hiding this comment.
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:
- Spec shape. #754 says: "Do not copy the existing context wholesale into a node: extract cohesive plan-owned specs with explicit invariants."
RangeEstimateSpecis close to a field-for-field copy of the context, and the key-window invariant is enforced byunwrap_orin the compiler and.expectin the estimator rather than by the type. See the inline comments. - 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); |
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
#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 = |
There was a problem hiding this comment.
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(), |
There was a problem hiding this comment.
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)?; |
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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 {")); |
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
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), |
There was a problem hiding this comment.
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>, |
There was a problem hiding this comment.
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)?; |
There was a problem hiding this comment.
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(), |
There was a problem hiding this comment.
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={:?})", |
There was a problem hiding this comment.
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, |
There was a problem hiding this comment.
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] |
There was a problem hiding this comment.
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.
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
Verification