From e132c2673394e5e375d0257a9a0abb9ac9673b2d Mon Sep 17 00:00:00 2001 From: Aarav Date: Fri, 25 Sep 2026 12:59:31 +0530 Subject: [PATCH] fix(promql): apply offset to subquery evaluation window `prom_subquery_expr_to_plan` destructured `SubqueryExpr` without reading `offset`, so `[range:step] offset ` planned exactly the same window as the un-offset form and silently returned data for the wrong time range. The plain vector/matrix-selector paths already threaded the offset through `selector_to_series_normalize_plan` and `RangeManipulate`. Shift the inner evaluation window back by the offset and pass the offset to the subquery's `RangeManipulate`, which maps the inner samples forward onto the evaluation timeline before bucketing them into ranges. This matches Prometheus, whose `evaluator.subqueryTimeRange` evaluates the inner expression over `(start - offset - range, end - offset]` and whose `evalSubquery` then hands the samples to the outer range-vector function as a `MatrixSelector` that still carries the subquery offset. An offset on the inner selector composes additively, as `subqueryTimes` documents. `RangeManipulate`'s protobuf message has no offset field and recovers it on decode from an immediately underlying `SeriesNormalize`. Since `RangeManipulate` is commutative in `dist_plan` and can be pushed below a `MergeScan`, insert that carrier node so the offset survives distributed planning instead of decoding as zero. Known divergence, unchanged by this commit: Prometheus anchors subquery step points on absolute epoch multiples of the step, while GreptimeDB anchors them on the evaluation start. The two agree whenever the offset is a multiple of the subquery step; the added sqlness cases stay within that range. Closes #9330 Co-Authored-By: Claude Opus 5 (1M context) Signed-off-by: Aarav --- src/query/src/promql/planner.rs | 88 ++++++++++++++++- .../standalone/common/promql/subquery.result | 99 +++++++++++++++++++ .../standalone/common/promql/subquery.sql | 55 +++++++++++ 3 files changed, 238 insertions(+), 4 deletions(-) diff --git a/src/query/src/promql/planner.rs b/src/query/src/promql/planner.rs index 73332c237c8..f5bdc51eb7d 100644 --- a/src/query/src/promql/planner.rs +++ b/src/query/src/promql/planner.rs @@ -548,18 +548,55 @@ impl PromPlanner { subquery_expr: &SubqueryExpr, ) -> Result { let SubqueryExpr { - expr, range, step, .. + expr, + range, + step, + offset, + .. } = subquery_expr; + // Prometheus shifts a subquery's own evaluation window back by its offset: the inner + // expression is evaluated over `(start - offset - range, end - offset]` + // (`evaluator.subqueryTimeRange` in `promql/engine.go`), producing samples that keep + // their real, un-shifted timestamps. `evalSubquery` then hands those samples to the + // outer range-vector function as a `MatrixSelector` whose `VectorSelector` still + // carries `subq.Offset`, so the function slices `(ts - offset - range, ts - offset]`. + // The offset is therefore applied once semantically, never twice. + // + // Here the same shift is expressed as: evaluate the inner plan over the shifted + // window, then let `RangeManipulate` map the resulting samples forward by `offset_ms` + // onto the evaluation timeline before bucketing them into ranges -- exactly what the + // plain matrix-selector path above does. + // + // An offset on the inner selector composes additively with this one, because the inner + // selector subtracts its own offset from the already shifted step timestamps + // (`subqueryTimes` documents this as "the sum of offsets ... of all subqueries in the + // path"). + // + // Divergence worth knowing: Prometheus anchors subquery step points on absolute epoch + // multiples of the step, so an offset that is not a multiple of the step does not + // rotate the grid. GreptimeDB anchors the grid on the evaluation start instead (see + // the `ctx.start` arithmetic below, which predates this offset handling), so a + // sub-step offset does rotate it. The two agree whenever the offset is a multiple of + // the subquery step. + let offset_ms = match offset { + Some(Offset::Pos(duration)) => duration.as_millis() as Millisecond, + Some(Offset::Neg(duration)) => -(duration.as_millis() as Millisecond), + None => 0, + }; + let current_interval = self.ctx.interval; if let Some(step) = step { self.ctx.interval = step.as_millis() as _; } let current_start = self.ctx.start; - self.ctx.start -= range.as_millis() as i64 - self.ctx.interval; + let current_end = self.ctx.end; + self.ctx.start -= offset_ms + range.as_millis() as i64 - self.ctx.interval; + self.ctx.end -= offset_ms; let input = self.prom_expr_to_plan(expr, query_engine_state).await?; self.ctx.interval = current_interval; self.ctx.start = current_start; + self.ctx.end = current_end; ensure!(!range.is_zero(), ZeroRangeSelectorSnafu); let range_ms = range.as_millis() as _; @@ -616,17 +653,38 @@ impl PromPlanner { .context(DataFusionPlanningSnafu)?; let divide_plan = LogicalPlan::Extension(Extension { node: Arc::new(SeriesDivide::new( - series_key_columns, + series_key_columns.clone(), time_index_column.clone(), sort_plan, )), }); + // `RangeManipulate`'s protobuf message has no offset field: on the decode path it + // recovers the offset from an immediately underlying `SeriesNormalize` (`local_offset`). + // `RangeManipulate` is `Commutative` in `dist_plan`, so a subquery's node can be pushed + // below a `MergeScan` and round-tripped through that encoding; without the carrier its + // offset would silently decode as zero and reintroduce this bug in distributed mode. + // Stale-marker filtering stays off: the input here is a computed inner result, not raw + // storage samples. + let divide_plan = if offset_ms == 0 { + divide_plan + } else { + LogicalPlan::Extension(Extension { + node: Arc::new(SeriesNormalize::new( + offset_ms, + time_index_column.clone(), + false, + series_key_columns, + divide_plan, + )), + }) + }; + let manipulate = RangeManipulate::new( self.ctx.start, self.ctx.end, self.ctx.interval, - 0, + offset_ms, range_ms, time_index_column, self.ctx.field_columns.clone(), @@ -12194,6 +12252,28 @@ mod test { indie_query_plan_compare(query, expected).await; } + /// `offset` on a subquery must shift the inner evaluation window back and be + /// carried into the outer range manipulation. See + /// . + #[tokio::test] + async fn count_over_time_subquery_with_offset() { + let query = "count_over_time(some_metric[10m:1m] offset 5m)"; + let expected = String::from( + "Filter: prom_count_over_time(timestamp_range,field_0) IS NOT NULL [timestamp:Timestamp(ms), prom_count_over_time(timestamp_range,field_0):Float64;N, tag_0:Utf8]\ + \n Projection: some_metric.timestamp, prom_count_over_time(timestamp_range, field_0) AS prom_count_over_time(timestamp_range,field_0), some_metric.tag_0 [timestamp:Timestamp(ms), prom_count_over_time(timestamp_range,field_0):Float64;N, tag_0:Utf8]\ + \n PromRangeManipulate: req range=[0..100000000], interval=[5000], eval range=[600000], time index=[timestamp], values=[\"field_0\"] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Dictionary(Int64, Float64);N, timestamp_range:Dictionary(Int64, Timestamp(ms))]\ + \n PromSeriesNormalize: offset=[300000], time index=[timestamp], filter NaN: [false] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ + \n PromSeriesDivide: tags=[\"tag_0\"] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ + \n Sort: some_metric.tag_0 ASC NULLS FIRST, some_metric.timestamp ASC NULLS FIRST [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ + \n PromInstantManipulate: range=[-840000..99700000], lookback=[1000], interval=[60000], time index=[timestamp] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ + \n PromSeriesDivide: tags=[\"tag_0\"] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ + \n Sort: some_metric.tag_0 ASC NULLS FIRST, some_metric.timestamp ASC NULLS FIRST [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ + \n Filter: some_metric.timestamp >= TimestampMillisecond(-840999, None) AND some_metric.timestamp <= TimestampMillisecond(99700000, None) [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ + \n TableScan: some_metric [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]", + ); + indie_query_plan_compare(query, expected).await; + } + #[tokio::test] async fn test_hash_join() { let mut eval_stmt = EvalStmt { diff --git a/tests/cases/standalone/common/promql/subquery.result b/tests/cases/standalone/common/promql/subquery.result index 12e65c4310c..cf734d2ed1b 100644 --- a/tests/cases/standalone/common/promql/subquery.result +++ b/tests/cases/standalone/common/promql/subquery.result @@ -63,3 +63,102 @@ drop table metric_total; Affected Rows: 0 +-- Offset on a subquery shifts the subquery's own evaluation window back by the offset. +-- Reference: Prometheus `evaluator.subqueryTimeRange` (promql/engine.go) evaluates the inner +-- expression over `(start - offset - range, end - offset]`; `evalSubquery` then passes the +-- resulting samples, which keep their real timestamps, to the outer range-vector function as a +-- `MatrixSelector` that still carries the subquery offset. See also +-- promql/promqltest/testdata/subquery.test. +-- +-- Every offset below is a whole multiple of the subquery step. Prometheus anchors subquery step +-- points on absolute epoch multiples of the step, while GreptimeDB anchors them on the +-- evaluation start, so the two only agree for step-multiple offsets; sub-step offsets are +-- deliberately not pinned here. +create table subquery_offset_total ( + ts timestamp time index, + host string primary key, + val double, +); + +Affected Rows: 0 + +insert into subquery_offset_total values + (0, 'a', 1), + (10000, 'a', 2), + (20000, 'a', 3), + (30000, 'a', 4), + (40000, 'a', 5), + (50000, 'a', 6), + (60000, 'a', 7); + +Affected Rows: 7 + +-- baseline: no offset at t=60 covers the 10s subquery points in (40s, 60s] -> 6 + 7 +tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s]); + ++---------------------+----------------------------------+------+ +| ts | prom_sum_over_time(ts_range,val) | host | ++---------------------+----------------------------------+------+ +| 1970-01-01T00:01:00 | 13.0 | a | ++---------------------+----------------------------------+------+ + +-- the same subquery evaluated at t=30 -> 3 + 4 +tql eval (30, 30, '1s') sum_over_time(subquery_offset_total[20s:10s]); + ++---------------------+----------------------------------+------+ +| ts | prom_sum_over_time(ts_range,val) | host | ++---------------------+----------------------------------+------+ +| 1970-01-01T00:00:30 | 7.0 | a | ++---------------------+----------------------------------+------+ + +-- `offset 30s` at t=60 must equal the un-offset subquery at t=30 +tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s] offset 30s); + ++---------------------+----------------------------------+------+ +| ts | prom_sum_over_time(ts_range,val) | host | ++---------------------+----------------------------------+------+ +| 1970-01-01T00:01:00 | 7.0 | a | ++---------------------+----------------------------------+------+ + +-- a zero subquery offset is not expressible: the shared duration check in `promql-parser` +-- rejects any zero duration literal, matching Prometheus's own `parseDuration` +-- ("duration must be greater than 0", cf. its `foo[0m]` parser test). The no-op case is +-- therefore covered by the offset-free baseline above. +tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s] offset 0s); + +Error: 2000(InvalidSyntax), duration must be greater than 0 + +-- a negative offset looks ahead of the evaluation time +tql eval (30, 30, '1s') sum_over_time(subquery_offset_total[20s:10s] offset -30s); + ++---------------------+----------------------------------+------+ +| ts | prom_sum_over_time(ts_range,val) | host | ++---------------------+----------------------------------+------+ +| 1970-01-01T00:00:30 | 13.0 | a | ++---------------------+----------------------------------+------+ + +-- an offset on the inner selector composes additively with the subquery offset: Prometheus +-- `subqueryTimes` accumulates "the sum of offsets and ranges of all subqueries in the path", +-- and the inner selector subtracts its own offset from the already shifted step timestamps. +-- 20s + 10s therefore behaves like the un-offset subquery at t=30. +tql eval (60, 60, '1s') sum_over_time((subquery_offset_total offset 10s)[20s:10s] offset 20s); + ++---------------------+----------------------------------+------+ +| ts | prom_sum_over_time(ts_range,val) | host | ++---------------------+----------------------------------+------+ +| 1970-01-01T00:01:00 | 7.0 | a | ++---------------------+----------------------------------+------+ + +-- ... and the inner offset alone accounts for the same total shift +tql eval (60, 60, '1s') sum_over_time((subquery_offset_total offset 30s)[20s:10s]); + ++---------------------+----------------------------------+------+ +| ts | prom_sum_over_time(ts_range,val) | host | ++---------------------+----------------------------------+------+ +| 1970-01-01T00:01:00 | 7.0 | a | ++---------------------+----------------------------------+------+ + +drop table subquery_offset_total; + +Affected Rows: 0 + diff --git a/tests/cases/standalone/common/promql/subquery.sql b/tests/cases/standalone/common/promql/subquery.sql index 95215b8488b..e95ae6a113e 100644 --- a/tests/cases/standalone/common/promql/subquery.sql +++ b/tests/cases/standalone/common/promql/subquery.sql @@ -20,3 +20,58 @@ tql eval (10, 10, '1s') rate(metric_total[20s:10s]); tql eval (20, 20, '1s') rate(metric_total[20s:5s]); drop table metric_total; + +-- Offset on a subquery shifts the subquery's own evaluation window back by the offset. +-- Reference: Prometheus `evaluator.subqueryTimeRange` (promql/engine.go) evaluates the inner +-- expression over `(start - offset - range, end - offset]`; `evalSubquery` then passes the +-- resulting samples, which keep their real timestamps, to the outer range-vector function as a +-- `MatrixSelector` that still carries the subquery offset. See also +-- promql/promqltest/testdata/subquery.test. +-- +-- Every offset below is a whole multiple of the subquery step. Prometheus anchors subquery step +-- points on absolute epoch multiples of the step, while GreptimeDB anchors them on the +-- evaluation start, so the two only agree for step-multiple offsets; sub-step offsets are +-- deliberately not pinned here. +create table subquery_offset_total ( + ts timestamp time index, + host string primary key, + val double, +); + +insert into subquery_offset_total values + (0, 'a', 1), + (10000, 'a', 2), + (20000, 'a', 3), + (30000, 'a', 4), + (40000, 'a', 5), + (50000, 'a', 6), + (60000, 'a', 7); + +-- baseline: no offset at t=60 covers the 10s subquery points in (40s, 60s] -> 6 + 7 +tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s]); + +-- the same subquery evaluated at t=30 -> 3 + 4 +tql eval (30, 30, '1s') sum_over_time(subquery_offset_total[20s:10s]); + +-- `offset 30s` at t=60 must equal the un-offset subquery at t=30 +tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s] offset 30s); + +-- a zero subquery offset is not expressible: the shared duration check in `promql-parser` +-- rejects any zero duration literal, matching Prometheus's own `parseDuration` +-- ("duration must be greater than 0", cf. its `foo[0m]` parser test). The no-op case is +-- therefore covered by the offset-free baseline above. +tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s] offset 0s); + +-- a negative offset looks ahead of the evaluation time +tql eval (30, 30, '1s') sum_over_time(subquery_offset_total[20s:10s] offset -30s); + +-- an offset on the inner selector composes additively with the subquery offset: Prometheus +-- `subqueryTimes` accumulates "the sum of offsets and ranges of all subqueries in the path", +-- and the inner selector subtracts its own offset from the already shifted step timestamps. +-- 20s + 10s therefore behaves like the un-offset subquery at t=30. +tql eval (60, 60, '1s') sum_over_time((subquery_offset_total offset 10s)[20s:10s] offset 20s); + +-- ... and the inner offset alone accounts for the same total shift +tql eval (60, 60, '1s') sum_over_time((subquery_offset_total offset 30s)[20s:10s]); + +drop table subquery_offset_total;