From 28e415a10344bb539e7b778cf6dbf9cc95bfa995 Mon Sep 17 00:00:00 2001 From: discord9 Date: Mon, 28 Sep 2026 09:38:12 +0000 Subject: [PATCH] feat(query): implement PromQL @ modifier on vector and matrix selectors (#9224) * feat(query): implement PromQL @ modifier on vector and matrix selectors The planner previously dropped the `at` field of selectors (`at: _`), so `some_metric @ 300` silently returned step-following values instead of the samples anchored at the fixed timestamp. Anchoring follows Prometheus's `setOffsetForAtModifier` + `refetch` semantics: resolve the anchor (`@ `, `@ start()`, `@ end()`), apply `offset` to the anchor, rewrite the selector offset to `eval_start - anchor`, scan only the anchored window, then report the same window at every evaluation step via a grid-wide-replay `InstantManipulate`. `@ start()` / `@ end()` resolve against the statement's evaluation range (`stmt_start` / `stmt_end`), which subquery planning does not rewrite. Timestamps before the Unix epoch are accepted; unrepresentable ones are rejected with `AtModifierTimestampOutOfRange` instead of wrapping. Selectors without `@` are planned exactly as before. Report: .e-agent/greptimedb_promql_compatibility_report_2026-09-16.md P0-2 Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * fix(promql): center predict_linear on the evaluation timestamp predict_linear_impl used the window's last sample time as the evaluation timestamp, and the UDF took only (ts_range, value_range, t) with no channel for the current instant. After a range selector is folded for @ (or with an offset) the same window is replayed at every step, so the prediction stayed constant at the anchor's answer; even a plain window ended before the step produced a stale value. Give prom_predict_linear a 4th argument carrying the step's evaluation instant (ms timestamp), derived in the planner from the row's time index plus the offset the window was folded with (at_offset for @, offset_ms otherwise, recorded on PromPlannerContext). The regression is centered on that instant, matching Prometheus' use of enh.Ts. The mixed float/native-histogram path forwards the same argument. Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * refactor(query): tighten @ modifier planner and predict_linear - range_fold_offset is always a concrete millisecond offset, not an Option; the single-step replay helper no longer wraps an infallible plan in Result. - predict_linear's eval timestamp is always cast to Timestamp(ms) at the call site, so the UDF drops its dead Int64 branch and the extra func_name argument. - Drop two redundant comparison cases from the at_modifier sqlness test and fix the lookback window comment to the half-open (0s, 300s]. Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * test(query): accept parse-time rejection of @ on Windows `@ 1e16` is 10^19 milliseconds, beyond i64::MAX. A Unix SystemTime holds it and the planner rejects the anchor it cannot represent, but a Windows SystemTime tops out near 1.8e12 seconds, so the parser's checked_add fails first and the same literal is rejected while parsing. Accept either rejection path so the test passes on both platforms. Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * fix(query): avoid replaying label_join over rewritten series keys Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * fix(query): narrow anchored range call promotion Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * fix(query): address @ modifier review feedback Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * fix(query): satisfy super import format check Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * docs(query): address at modifier review follow-ups Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> --------- Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> --- src/promql/benches/bench_range_fn.rs | 4 +- src/promql/src/functions/predict_linear.rs | 156 +++- src/query/src/promql/error.rs | 8 + src/query/src/promql/planner.rs | 408 ++++++++-- src/query/src/promql/planner/at_modifier.rs | 326 ++++++++ src/query/src/promql/planner/test.rs | 675 +++++++++++++++- src/query/src/promql/planner/test/delta.rs | 2 +- .../common/promql/at_modifier.result | 743 ++++++++++++++++++ .../standalone/common/promql/at_modifier.sql | 301 +++++++ 9 files changed, 2529 insertions(+), 94 deletions(-) create mode 100644 src/query/src/promql/planner/at_modifier.rs create mode 100644 tests/cases/standalone/common/promql/at_modifier.result create mode 100644 tests/cases/standalone/common/promql/at_modifier.sql diff --git a/src/promql/benches/bench_range_fn.rs b/src/promql/benches/bench_range_fn.rs index 0f0a4867665..b4e026bae3d 100644 --- a/src/promql/benches/bench_range_fn.rs +++ b/src/promql/benches/bench_range_fn.rs @@ -270,7 +270,7 @@ fn make_quantile_input(num_points: usize, window_size: u32) -> Vec Vec { - let (ts_range, val_range, _) = build_sliding_ranges( + let (ts_range, val_range, eval_ts) = build_sliding_ranges( num_points, window_size, 1, @@ -282,6 +282,8 @@ fn make_predict_linear_input(num_points: usize, window_size: u32) -> Vec Result { error::ensure( - input.len() == 3, - DataFusionError::Plan("prom_predict_linear function should have 3 inputs".to_string()), + input.len() == 4, + DataFusionError::Plan("prom_predict_linear function should have 4 inputs".to_string()), )?; let t_col = &input[2]; + let eval_ts_col = &input[3]; let ts_range = extract_range_array(&input[0])?; let value_range = extract_range_array(&input[1])?; @@ -128,6 +131,51 @@ impl PredictLinear { Box::new(t_array.iter()) } }; + // The evaluation instant the window above was folded for. It is the instant the regression + // is centered on, which is *not* the last sample's timestamp: a window may end before the + // step being evaluated, and an `@`-anchored window is replayed at every step while keeping + // its own end. See `PromPlanner::create_range_eval_ts_expr`. + let eval_ts_iter: Box>> = match eval_ts_col { + ColumnarValue::Scalar(eval_ts_scalar) => { + let eval_ts = match eval_ts_scalar { + ScalarValue::TimestampMillisecond(Some(eval_ts), _) => *eval_ts, + // For a NULL or otherwise unusable evaluation timestamp, returns NULL array, + // which conforms to PromQL's behavior. + _ => { + let null_array = Float64Array::new_null(ts_range.len()); + return Ok(ColumnarValue::Array(Arc::new(null_array))); + } + }; + Box::new((0..ts_range.len()).map(move |_| Some(eval_ts))) + } + ColumnarValue::Array(eval_ts_array) => { + error::ensure( + eval_ts_array.len() == ts_range.len(), + DataFusionError::Execution(format!( + "{}: evaluation timestamp array should have the same length as other columns, found {} and {}", + Self::name(), + eval_ts_array.len(), + ts_range.len() + )), + )?; + match eval_ts_array.data_type() { + DataType::Timestamp(TimeUnit::Millisecond, _) => Box::new( + eval_ts_array + .as_any() + .downcast_ref::() + .expect("checked by data type") + .iter(), + ), + other => { + return Err(DataFusionError::Execution(format!( + "{}: expect TimestampMillisecond as evaluation timestamp array's type, found {}", + Self::name(), + other + ))); + } + } + } + }; let all_timestamps = ts_range .values() .as_any() @@ -140,7 +188,12 @@ impl PredictLinear { .downcast_ref::() .unwrap(); let mut result_builder = Float64Builder::with_capacity(ts_range.len()); - for (index, t) in t_iter.enumerate() { + for (index, (t, eval_ts)) in t_iter.zip(eval_ts_iter).enumerate() { + // A step without an evaluation instant has nothing to extrapolate to. + let Some(eval_ts) = eval_ts else { + result_builder.append_null(); + continue; + }; match predict_linear_impl( &ts_range, &value_range, @@ -148,7 +201,7 @@ impl PredictLinear { all_values, index, t.unwrap(), - Self::name(), + eval_ts, )? { Some(value) => result_builder.append_value(value), None => result_builder.append_null(), @@ -167,7 +220,7 @@ fn predict_linear_impl( all_values: &Float64Array, index: usize, t: i64, - func_name: &str, + eval_ts: i64, ) -> Result, DataFusionError> { let (ts_offset, ts_len) = ts_range.get_offset_length(index).unwrap(); let (value_offset, value_len) = value_range.get_offset_length(index).unwrap(); @@ -175,22 +228,28 @@ fn predict_linear_impl( ts_len == value_len, DataFusionError::Execution(format!( "{}: time and value arrays in a group should have the same length, found {} and {}", - func_name, ts_len, value_len + PredictLinear::name(), + ts_len, + value_len )), )?; if ts_len < 2 { return Ok(None); } - // last timestamp is evaluation timestamp - let evaluate_ts = all_timestamps[ts_offset + ts_len - 1]; + // Like Prometheus, the regression is centered on the evaluation timestamp: the returned + // intercept is the value the window's trend predicts at the step being evaluated, and `t` is + // the horizon (in seconds) that is added on top of it. Centering on the last sample instead + // would drop the distance between the window's end and the step, which is non-zero both + // without `@` (the window can end before the step) and with it (the same anchored window is + // replayed at every step). let (slope, intercept) = linear_regression_slices( all_timestamps, ts_offset, all_values, value_offset, value_len, - evaluate_ts, + eval_ts, ); if slope.is_none() || intercept.is_none() { @@ -231,6 +290,11 @@ mod test { (ts_range_array, value_range_array) } + /// The evaluation timestamp at which the window built by [`build_test_range_arrays`] ends. + fn window_end_eval_ts() -> ScalarValue { + ScalarValue::TimestampMillisecond(Some(3000), None) + } + #[test] fn calculate_predict_linear_none() { let ts_array = Arc::new(TimestampMillisecondArray::from_iter( @@ -244,7 +308,7 @@ mod test { PredictLinear::scalar_udf(), ts_array, value_array, - vec![ScalarValue::Int64(Some(0))], + vec![ScalarValue::Int64(Some(0)), window_end_eval_ts()], vec![None, None], ); } @@ -256,7 +320,7 @@ mod test { PredictLinear::scalar_udf(), ts_array, value_array, - vec![ScalarValue::Int64(Some(0))], + vec![ScalarValue::Int64(Some(0)), window_end_eval_ts()], // value at t = 0 vec![Some(38.63636363636364)], ); @@ -269,7 +333,7 @@ mod test { PredictLinear::scalar_udf(), ts_array, value_array, - vec![ScalarValue::Int64(Some(3000))], + vec![ScalarValue::Int64(Some(3000)), window_end_eval_ts()], // value at t = 3000 vec![Some(31856.818181818187)], ); @@ -282,7 +346,7 @@ mod test { PredictLinear::scalar_udf(), ts_array, value_array, - vec![ScalarValue::Int64(Some(4200))], + vec![ScalarValue::Int64(Some(4200)), window_end_eval_ts()], // value at t = 4200 vec![Some(44584.09090909091)], ); @@ -295,7 +359,7 @@ mod test { PredictLinear::scalar_udf(), ts_array, value_array, - vec![ScalarValue::Int64(Some(6600))], + vec![ScalarValue::Int64(Some(6600)), window_end_eval_ts()], // value at t = 6600 vec![Some(70038.63636363638)], ); @@ -308,7 +372,7 @@ mod test { PredictLinear::scalar_udf(), ts_array, value_array, - vec![ScalarValue::Int64(Some(7800))], + vec![ScalarValue::Int64(Some(7800)), window_end_eval_ts()], // value at t = 7800 vec![Some(82765.9090909091)], ); @@ -327,7 +391,7 @@ mod test { PredictLinear::scalar_udf(), ts_array, value_array, - vec![ScalarValue::Int64(Some(0))], + vec![ScalarValue::Int64(Some(0)), window_end_eval_ts()], vec![Some(30.0)], ); } @@ -348,9 +412,69 @@ mod test { ColumnarValue::Array(Arc::new(ts_dict)), ColumnarValue::Array(Arc::new(value_dict)), ColumnarValue::Scalar(ScalarValue::Int64(Some(0))), + ColumnarValue::Scalar(window_end_eval_ts()), ]) .unwrap_err(); assert!(err.to_string().contains("Empty range is not expected")); } + + /// The regression is centered on the evaluation timestamp, so the returned intercept follows + /// the step being evaluated even though the window stays the same. + #[test] + fn calculate_predict_linear_centers_on_eval_ts() { + // One sample per second with values 0, 1, 2: the trend is exactly one unit per second. + for (eval_ts, t, expected) in [ + // The window ends at 2s: evaluating there returns its last value. + (2000, 0, 2.0), + // One second past the window's end the trend is at 3. + (3000, 0, 3.0), + // One second past the window's end, predicting another two seconds ahead. + (3000, 2, 5.0), + ] { + let ts_values = Arc::new(TimestampMillisecondArray::from_iter( + [0i64, 1000, 2000].into_iter().map(Some), + )); + let value_values = Arc::new(Float64Array::from_iter([0.0, 1.0, 2.0])); + let ts_array = RangeArray::from_ranges(ts_values, [(0, 3)]).unwrap(); + let value_array = RangeArray::from_ranges(value_values, [(0, 3)]).unwrap(); + + simple_range_udf_runner( + PredictLinear::scalar_udf(), + ts_array, + value_array, + vec![ + ScalarValue::Int64(Some(t)), + ScalarValue::TimestampMillisecond(Some(eval_ts), None), + ], + vec![Some(expected)], + ); + } + } + + /// The evaluation timestamp arrives as a column, so the UDF accepts it as an array too. + #[test] + fn calculate_predict_linear_accepts_eval_ts_array() { + let (ts_array, value_array) = build_test_range_arrays(); + let eval_ts_array = Arc::new(TimestampMillisecondArray::from_iter_values([3000])); + + let result = PredictLinear::predict_linear(&[ + ColumnarValue::Array(Arc::new(ts_array.into_dict())), + ColumnarValue::Array(Arc::new(value_array.into_dict())), + ColumnarValue::Scalar(ScalarValue::Int64(Some(0))), + ColumnarValue::Array(eval_ts_array), + ]) + .unwrap(); + + let result = result + .into_array(1) + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .clone(); + assert_eq!(result.len(), 1); + // The same value `calculate_predict_linear_test1` gets from a scalar evaluation timestamp. + assert!((result.value(0) - 38.63636363636364).abs() < 0.0001); + } } diff --git a/src/query/src/promql/error.rs b/src/query/src/promql/error.rs index 45e08349623..6cd45d9e57c 100644 --- a/src/query/src/promql/error.rs +++ b/src/query/src/promql/error.rs @@ -195,6 +195,13 @@ pub enum Error { location: Location, }, + #[snafu(display("Timestamp out of range for the `@` modifier: {}", timestamp))] + AtModifierTimestampOutOfRange { + timestamp: String, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display("Timestamp out of range: {} of {:?}", timestamp, unit))] TimestampOutOfRange { timestamp: i64, @@ -252,6 +259,7 @@ impl ErrorExt for Error { | SameLabelSet { .. } | TimestampOutOfRange { .. } | SystemTimeOutOfRange { .. } + | AtModifierTimestampOutOfRange { .. } | InvalidRegularExpression { .. } | InvalidDestinationLabelName { .. } => StatusCode::InvalidArguments, diff --git a/src/query/src/promql/planner.rs b/src/query/src/promql/planner.rs index 08e7005f9fb..be0c1364489 100644 --- a/src/query/src/promql/planner.rs +++ b/src/query/src/promql/planner.rs @@ -12,10 +12,9 @@ // See the License for the specific language governing permissions and // limitations under the License. +mod at_modifier; mod function_plans; - mod island; - mod matching_filters; mod set_operator; @@ -165,6 +164,12 @@ struct PromPlannerContext { end: Millisecond, interval: Millisecond, lookback_delta: Millisecond, + /// Evaluation range of the whole statement, which `@ start()` and `@ end()` refer to. + /// + /// Unlike [`Self::start`] and [`Self::end`], these are never rewritten while planning, so a + /// selector inside a subquery still resolves `@ start()` / `@ end()` against the statement. + stmt_start: Millisecond, + stmt_end: Millisecond, // planner states table_name: Option, @@ -193,6 +198,16 @@ struct PromPlannerContext { schema_name: Option, /// The range in millisecond of range selector. None if there is no range selector. range: Option, + /// The offset in milliseconds the window of the last planned range selector is folded with, + /// or `None` when no range selector has been planned since the last read. + /// + /// [`Self::start`] and the sample timestamps are compared on the shifted timeline: a range + /// payload carries `sample_timestamp + offset`, while the time index column of a folded row + /// stays the evaluation timestamp of its step. A function that reads it right after its input + /// plan is built, like `predict_linear`, consumes it instead of reading the state, so the + /// offset cannot leak from one input to another; see + /// [`PromPlanner::create_function_expr`]. + range_fold_offset: Option, } /// Result labels a vector-vector binary operation derives from its matching modifier, projected @@ -226,11 +241,15 @@ impl BinaryResultLabels { impl PromPlannerContext { fn from_eval_stmt(stmt: &EvalStmt) -> Self { + let start = stmt.start.duration_since(UNIX_EPOCH).unwrap().as_millis() as Millisecond; + let end = stmt.end.duration_since(UNIX_EPOCH).unwrap().as_millis() as Millisecond; Self { - start: stmt.start.duration_since(UNIX_EPOCH).unwrap().as_millis() as _, - end: stmt.end.duration_since(UNIX_EPOCH).unwrap().as_millis() as _, + start, + end, interval: stmt.interval.as_millis() as _, lookback_delta: stmt.lookback_delta.as_millis() as _, + stmt_start: start, + stmt_end: end, ..Default::default() } } @@ -246,6 +265,7 @@ impl PromPlannerContext { self.selector_matcher.clear(); self.schema_name = None; self.range = None; + self.range_fold_offset = None; } /// Reset table name and schema to empty @@ -325,6 +345,16 @@ impl PromPlanner { timestamp_fn: bool, query_engine_state: &QueryEngineState, ) -> Result { + // An anchored range call is step-invariant: evaluate it once, at the start of the + // evaluation, and report its result at every step; see + // [`Self::promote_anchored_range_call`]. + if let Some(plan) = self + .promote_anchored_range_call(prom_expr, timestamp_fn, query_engine_state) + .await? + { + return Ok(plan); + } + let res = match prom_expr { PromExpr::Aggregate(expr) => { self.prom_aggr_expr_to_plan(query_engine_state, expr) @@ -458,6 +488,9 @@ impl PromPlanner { divide_plan, ) .context(DataFusionPlanningSnafu)?; + // A subquery always folds with offset 0, so its payload timestamps are already on the + // evaluation timeline a function above it reads; see [`Self::create_range_eval_ts_expr`]. + self.ctx.range_fold_offset = Some(0); Ok(LogicalPlan::Extension(Extension { node: Arc::new(manipulate), @@ -1579,6 +1612,44 @@ impl PromPlanner { Ok(plan) } + /// The offset of a selector in milliseconds. A positive offset selects samples from an earlier + /// time and moves them forward into the evaluation timeline. + fn offset_millis(offset: &Option) -> Millisecond { + match offset { + Some(Offset::Pos(duration)) => duration.as_millis() as Millisecond, + Some(Offset::Neg(duration)) => -(duration.as_millis() as Millisecond), + None => 0, + } + } + + /// The columns that identify one series, which is the series key expected by the PromQL plan + /// nodes that hold exactly one series per input batch. + fn series_key_columns(&self) -> Vec { + if self.ctx.use_tsid { + vec![DATA_SCHEMA_TSID_COLUMN_NAME.to_string()] + } else { + self.ctx.tag_columns.clone() + } + } + + /// Keep replay and series division on the same effective keys when a call rewrites labels. + fn series_key_columns_for_schema(&self, schema: &DFSchemaRef) -> Vec { + let has_tsid = schema.fields().iter().any(|field| { + field.name() == DATA_SCHEMA_TSID_COLUMN_NAME + && field.data_type() == &ArrowDataType::UInt64 + }); + if has_tsid { + vec![DATA_SCHEMA_TSID_COLUMN_NAME.to_string()] + } else { + self.ctx + .tag_columns + .iter() + .filter(|name| schema.has_column_with_unqualified_name(name)) + .cloned() + .collect() + } + } + async fn prom_vector_selector_to_plan( &mut self, vector_selector: &VectorSelector, @@ -1588,20 +1659,37 @@ impl PromPlanner { name, offset, matchers, - at: _, + at, } = vector_selector; let matchers = self.preprocess_label_matchers(matchers, name)?; if let Some(empty_plan) = self.setup_context().await? { return Ok(empty_plan); } - 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 offset_ms = Self::offset_millis(offset); + // `@` anchors the sample window at a fixed timestamp: the selector selects its samples + // around the anchor once, instead of following the outer evaluation grid. See + // [`Self::at_modifier_offset`]. + let at_offset = self.at_modifier_offset(at, offset)?; + let grid_start = self.ctx.start; + let grid_end = self.ctx.end; + let normalize = match at_offset { + Some(at_offset) => { + // Select the anchored samples at the start of the evaluation, with the offset that + // re-anchors the selector. + // The planner is single-use (one `EvalStmt` produces one plan), so an error + // below aborts the whole planning and `ctx.end` needs no restore-on-error. + self.ctx.end = grid_start; + let plan = self + .selector_to_series_normalize_plan(at_offset, matchers, false) + .await?; + self.ctx.end = grid_end; + plan + } + None => { + self.selector_to_series_normalize_plan(offset_ms, matchers, false) + .await? + } }; - let normalize = self - .selector_to_series_normalize_plan(offset, matchers, false) - .await?; let time_index_column = self.ctx .time_index_column @@ -1683,24 +1771,46 @@ impl PromPlanner { }; let field_column = self.ctx.field_columns.first().cloned(); - let manipulate = InstantManipulate::new( - self.ctx.start, - self.ctx.end, - self.ctx.lookback_delta, - self.ctx.interval, - offset_ms, - time_index_column, - if self.ctx.use_tsid { - vec![DATA_SCHEMA_TSID_COLUMN_NAME.to_string()] - } else { - self.ctx.tag_columns.clone() - }, - field_column, - normalize, - ); - let manipulate = LogicalPlan::Extension(Extension { - node: Arc::new(manipulate), - }); + let series_key_columns = self.series_key_columns(); + let manipulate = match at_offset { + Some(at_offset) => { + // Select the anchored sample once, then report it at every step of the outer + // grid. The samples keep their native anchor-time timestamps, so the manipulate + // must shift them onto the evaluation timeline with the rewritten offset. + let anchored = InstantManipulate::new( + grid_start, + grid_start, + self.ctx.lookback_delta, + self.ctx.interval, + at_offset, + time_index_column.clone(), + series_key_columns, + field_column, + normalize, + ); + self.replay_over_grid( + LogicalPlan::Extension(Extension { + node: Arc::new(anchored), + }), + grid_start, + grid_end, + time_index_column, + ) + } + None => LogicalPlan::Extension(Extension { + node: Arc::new(InstantManipulate::new( + grid_start, + grid_end, + self.ctx.lookback_delta, + self.ctx.interval, + offset_ms, + time_index_column, + series_key_columns, + field_column, + normalize, + )), + }), + }; if let Some(timestamp_value_column) = timestamp_value_column { self.create_timestamp_func_plan(manipulate, ×tamp_value_column) } else { @@ -1758,46 +1868,108 @@ impl PromPlanner { name, offset, matchers, - .. + at, } = vs; let matchers = self.preprocess_label_matchers(matchers, name)?; ensure!(!range.is_zero(), ZeroRangeSelectorSnafu); let range_ms = range.as_millis() as _; self.ctx.range = Some(range_ms); - 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 offset_ms = Self::offset_millis(offset); + + // `@` anchors the range selector's window at a fixed timestamp, so the same window is fed + // to the enclosing function at every evaluation step. See [`Self::at_modifier_offset`]. + let at_offset = self.at_modifier_offset(at, offset)?; + let grid_start = self.ctx.start; + let grid_end = self.ctx.end; // Some functions like rate may require special fields in the RangeManipulate plan // so we can't skip RangeManipulate. - let normalize = match self.setup_context().await? { - Some(empty_plan) => empty_plan, + let (normalize, at_offset) = match self.setup_context().await? { + // An empty metric does not contain any sample, so anchoring cannot change the result. + // The manipulate below folds the empty input with `offset_ms`, and the recorded fold + // offset agrees with it instead of with the anchor the window cannot use. + Some(empty_plan) => { + self.ctx.range_fold_offset = Some(offset_ms); + (empty_plan, None) + } None => { - self.selector_to_series_normalize_plan(offset, matchers, true) - .await? + let normalize = match at_offset { + Some(at_offset) => { + // Fold the anchored window once, at the start of the evaluation. + // Single-use planner: an error below aborts planning, so `ctx.end` + // needs no restore-on-error. + self.ctx.end = grid_start; + let plan = self + .selector_to_series_normalize_plan(at_offset, matchers, true) + .await?; + self.ctx.end = grid_end; + plan + } + None => { + self.selector_to_series_normalize_plan(offset_ms, matchers, true) + .await? + } + }; + // Samples are shifted onto the evaluation timeline with the very same offset while + // the window is folded. Record it so that the function above the selector can + // recover the evaluation instant of the folded window; see + // [`Self::create_range_eval_ts_expr`]. + self.ctx.range_fold_offset = Some(at_offset.unwrap_or(offset_ms)); + (normalize, at_offset) } }; - let manipulate = RangeManipulate::new( - self.ctx.start, - self.ctx.end, - self.ctx.interval, - offset_ms, - // TODO(ruihang): convert via Timestamp datatypes to support different time units - range_ms, - self.ctx - .time_index_column - .clone() - .expect("time index should be set in `setup_context`"), - self.ctx.field_columns.clone(), - normalize, - ) - .context(DataFusionPlanningSnafu)?; + let time_index_column = self + .ctx + .time_index_column + .clone() + .expect("time index should be set in `setup_context`"); + let manipulate = match at_offset { + Some(at_offset) => { + // Fold the anchored window once, then report it at every step of the outer + // grid. The samples keep their native anchor-time timestamps, so the manipulate + // must shift them onto the evaluation timeline with the rewritten offset. + let anchored = RangeManipulate::new( + grid_start, + grid_start, + self.ctx.interval, + at_offset, + // TODO(ruihang): convert via Timestamp datatypes to support different time units + range_ms, + time_index_column.clone(), + self.ctx.field_columns.clone(), + normalize, + ) + .context(DataFusionPlanningSnafu)?; + self.replay_over_grid( + LogicalPlan::Extension(Extension { + node: Arc::new(anchored), + }), + grid_start, + grid_end, + time_index_column, + ) + } + None => { + let manipulate = RangeManipulate::new( + grid_start, + grid_end, + self.ctx.interval, + offset_ms, + // TODO(ruihang): convert via Timestamp datatypes to support different time units + range_ms, + time_index_column, + self.ctx.field_columns.clone(), + normalize, + ) + .context(DataFusionPlanningSnafu)?; - Ok(LogicalPlan::Extension(Extension { - node: Arc::new(manipulate), - })) + LogicalPlan::Extension(Extension { + node: Arc::new(manipulate), + }) + } + }; + + Ok(manipulate) } async fn prom_call_expr_to_plan( @@ -1846,11 +2018,17 @@ impl PromPlanner { ), }) }; + // The input plan records the fold offset of the range selector it is built from. Take it + // here, so that the offset of one input cannot leak into another call, and pass it to + // `create_function_expr`: the function that reads it (`predict_linear`) then depends on an + // argument instead of on planner state written by the selector below it. + let range_fold_offset = self.ctx.range_fold_offset.take(); let (mut func_exprs, new_tags) = self.create_function_expr( func, args.literals.clone(), input.schema(), query_engine_state, + range_fold_offset, )?; func_exprs.insert(0, self.create_time_index_column_expr()?); func_exprs.extend_from_slice(&self.create_tag_column_exprs()?); @@ -2034,7 +2212,7 @@ impl PromPlanner { async fn selector_to_series_normalize_plan( &mut self, - offset: &Option, + offset_duration: Millisecond, label_matchers: Matchers, is_range_selector: bool, ) -> Result { @@ -2044,11 +2222,6 @@ impl PromPlanner { let table_schema = table_scan.schema(); // make filter exprs - let offset_duration = match offset { - Some(Offset::Pos(duration)) => duration.as_millis() as Millisecond, - Some(Offset::Neg(duration)) => -(duration.as_millis() as Millisecond), - None => 0, - }; let mut scan_filters = Self::matchers_to_expr(label_matchers.clone(), table_schema)?; if let Some(time_index_filter) = self.build_time_index_filter(offset_duration, table_schema)? @@ -2147,11 +2320,7 @@ impl PromPlanner { } // make sort plan - let series_key_columns = if self.ctx.use_tsid { - vec![DATA_SCHEMA_TSID_COLUMN_NAME.to_string()] - } else { - self.ctx.tag_columns.clone() - }; + let series_key_columns = self.series_key_columns(); let sort_exprs = if self.ctx.use_tsid { vec![ @@ -2948,6 +3117,7 @@ impl PromPlanner { float_field: &str, histogram_field: &str, input_schema: &DFSchemaRef, + range_fold_offset: Option, ) -> Result>> { let returns_histogram = matches!( func.name, @@ -2987,6 +3157,14 @@ impl PromPlanner { Box::new(other_input_exprs[0].clone()), ArrowDataType::Int64, )); + // Same evaluation instant as in the non-mixed path; it follows the other inputs so the + // float UDF below is called as `predict_linear(ts_range, value_range, t, eval_ts)`. + // The regression is only defined for a folded window, so a missing fold offset is an + // error instead of a default. + other_input_exprs.push_back(self.create_range_eval_ts_expr( + range_fold_offset.context(ExpectRangeSelectorSnafu)?, + input_schema, + )?); } let timestamp_range = DfExpr::Column(Column::from_name( @@ -3045,7 +3223,13 @@ impl PromPlanner { .alias(histogram_field), ] } else { - let display_name = float_expr.schema_name().to_string(); + let display_name = if func.name == "predict_linear" { + // The evaluation instant is the private last argument of the mixed float UDF + // call; keep it out of the output column name like on the non-mixed path. + Self::name_without_last_arg(&float_expr) + } else { + float_expr.schema_name().to_string() + }; self.ctx.field_columns = vec![display_name.clone()]; vec![float_expr.alias(display_name)] }; @@ -3063,6 +3247,7 @@ impl PromPlanner { other_input_exprs: Vec, input_schema: &DFSchemaRef, query_engine_state: &QueryEngineState, + range_fold_offset: Option, ) -> Result<(Vec, Vec)> { // TODO(ruihang): check function args list let mut other_input_exprs: VecDeque = other_input_exprs.into(); @@ -3075,6 +3260,7 @@ impl PromPlanner { &float_field, &histogram_field, input_schema, + range_fold_offset, )? { return Ok((exprs, vec![])); @@ -3276,6 +3462,16 @@ impl PromPlanner { Box::new(other_input_exprs[0].clone()), ArrowDataType::Int64, )); + // The prediction starts at the evaluation instant of the step, which the + // window's fold offset recovers from the row's time index; see + // [`Self::create_range_eval_ts_expr`]. It is appended last, so the UDF is + // called as `predict_linear(ts_range, value_range, t, eval_ts)`. The + // regression is only defined for a folded window, so a missing fold offset + // is an error instead of a default. + other_input_exprs.push_back(self.create_range_eval_ts_expr( + range_fold_offset.context(ExpectRangeSelectorSnafu)?, + input_schema, + )?); ScalarFunc::Udf(Arc::new(PredictLinear::scalar_udf())) } } @@ -3648,7 +3844,17 @@ impl PromPlanner { exprs = exprs .into_iter() .map(|expr| { - let display_name = expr.schema_name().to_string(); + // `predict_linear` appends its private evaluation instant as the last argument + // of the UDF call; the output column is named after the call without it, so + // the injected expression stays out of the user-visible schema. The + // native-histogram drop UDF takes no such argument. + let display_name = if func.name == "predict_linear" + && !all_field_columns_are_native_histogram_ranges + { + Self::name_without_last_arg(&expr) + } else { + expr.schema_name().to_string() + }; new_field_columns.push(display_name.clone()); Ok(expr.alias(display_name)) }) @@ -3917,6 +4123,66 @@ impl PromPlanner { ))) } + /// Builds the evaluation instant the window of the last planned range selector is folded for, + /// as a `Timestamp(Millisecond)` expression. + /// + /// The timestamp payload of a folded window is shifted onto the evaluation timeline by the + /// offset the window is folded with (`fold_offset`), while the time index column of a folded + /// row keeps the evaluation timestamp of its step. Adding the offset back yields the + /// evaluation instant on the payload timeline, which is where the regression of + /// `predict_linear` is centered: neither a plain window (which may end before the step, and + /// is additionally shifted by `offset` on the payload timeline) nor an `@`-anchored one + /// (whose end is the anchor, while the payload is shifted by `at_offset`) ends at the step it + /// is evaluated at. + /// + /// The sum is computed on the millisecond representation and cast back, so that the result + /// keeps the `Timestamp(Millisecond)` type the range functions declare for it. + fn create_range_eval_ts_expr( + &self, + fold_offset: Millisecond, + input_schema: &DFSchemaRef, + ) -> Result { + let eval_ts = self + .create_time_index_column_expr()? + .cast_to( + &ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None), + input_schema, + ) + .context(DataFusionPlanningSnafu)? + .cast_to(&ArrowDataType::Int64, input_schema) + .context(DataFusionPlanningSnafu)?; + DfExpr::BinaryExpr(BinaryExpr { + left: Box::new(eval_ts), + op: Operator::Plus, + right: Box::new(lit(fold_offset)), + }) + .cast_to( + &ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None), + input_schema, + ) + .context(DataFusionPlanningSnafu) + } + + /// The name of `expr` without its last argument. + /// + /// A `predict_linear` call appends the evaluation instant of its window as a private last + /// argument ([`Self::create_range_eval_ts_expr`]). Naming the output column after the call + /// the user wrote, without that argument, keeps the injected expression out of the + /// user-visible schema; the expression itself keeps every argument it needs. + fn name_without_last_arg(expr: &DfExpr) -> String { + if let DfExpr::ScalarFunction(ScalarFunction { func, args }) = expr + && let Some((_, visible_args)) = args.split_last() + { + let visible = ScalarFunction { + func: func.clone(), + args: visible_args.to_vec(), + }; + return DfExpr::ScalarFunction(visible).schema_name().to_string(); + } + + expr.schema_name().to_string() + } + fn create_tag_column_exprs(&self) -> Result> { let mut result = Vec::with_capacity(self.ctx.tag_columns.len()); for tag in &self.ctx.tag_columns { diff --git a/src/query/src/promql/planner/at_modifier.rs b/src/query/src/promql/planner/at_modifier.rs new file mode 100644 index 00000000000..bde7c51bcb1 --- /dev/null +++ b/src/query/src/promql/planner/at_modifier.rs @@ -0,0 +1,326 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Anchored selector planning and replay for the PromQL `@` modifier. + +use std::sync::Arc; +use std::time::{SystemTime, UNIX_EPOCH}; + +use datafusion::logical_expr::{Extension, LogicalPlan, LogicalPlanBuilder}; +use datafusion::prelude::{Column, Expr as DfExpr}; +use promql::extension_plan::{InstantManipulate, Millisecond, SeriesDivide}; +use promql_parser::parser::{ + AtModifier, Call, Expr as PromExpr, MatrixSelector, Offset, ParenExpr, +}; +use snafu::{OptionExt, ResultExt, ensure}; + +use crate::promql::error::{ + AtModifierTimestampOutOfRangeSnafu, DataFusionPlanningSnafu, Result, TimeIndexNotFoundSnafu, +}; +use crate::promql::planner::PromPlanner; +use crate::query_engine::QueryEngineState; + +impl PromPlanner { + /// Resolve the `@` modifier of a vector or matrix selector into the timestamp its sample + /// window is anchored at, in milliseconds since the Unix epoch. + /// + /// Prometheus semantics: + /// - `@ ` anchors at the given timestamp, + /// - `@ start()` / `@ end()` anchor at the evaluation range of the whole statement, + /// - `offset` shifts the anchor backwards: the window ends at `anchor - offset`. + /// + /// Returns `None` when the selector has no `@` modifier. + fn at_ref_time( + &self, + at: &Option, + offset: &Option, + ) -> Result> { + let anchor = match at { + None => return Ok(None), + Some(AtModifier::Start) => self.ctx.stmt_start, + Some(AtModifier::End) => self.ctx.stmt_end, + Some(AtModifier::At(time)) => Self::system_time_to_millis(time)?, + }; + Ok(Some(Self::anchor_sub(anchor, Self::offset_millis(offset))?)) + } + + /// Subtracts `rhs` from `lhs` on the millisecond timeline of an `@` anchor. + /// + /// A negative result is valid: `@` and `offset` accept timestamps before the Unix epoch. A + /// result outside the representable millisecond range is rejected like an unrepresentable + /// anchor ([`Self::system_time_to_millis`]) instead of clamping it, so the same class of + /// input always gets the same answer. + pub(crate) fn anchor_sub(lhs: Millisecond, rhs: Millisecond) -> Result { + lhs.checked_sub(rhs) + .with_context(|| AtModifierTimestampOutOfRangeSnafu { + timestamp: format!("{}ms - {}ms", lhs, rhs), + }) + } + + /// The offset a selector with an `@` modifier is evaluated with. + /// + /// Prometheus anchors such a selector by rewriting its offset to `eval_time - anchor` + /// (`setOffsetForAtModifier`), so that the selector always selects its samples around `anchor` + /// regardless of the step being evaluated. `eval_time` is the start of the evaluation the + /// selector belongs to, which is `ctx.start`. + /// + /// Returns `None` when the selector has no `@` modifier. + pub(crate) fn at_modifier_offset( + &self, + at: &Option, + offset: &Option, + ) -> Result> { + let Some(anchor) = self.at_ref_time(at, offset)? else { + return Ok(None); + }; + Ok(Some(Self::anchor_sub(self.ctx.start, anchor)?)) + } + + /// Whether `expr` is a call that has to be evaluated once for the whole grid, because it folds + /// a range selector anchored by `@` — the range argument of the call's parser signature, which + /// is a [`MatrixSelector`] here; only a call with such an argument can take this path, so no + /// function-name registry is involved. + /// + /// This is the shape that needs Prometheus' `StepInvariantExpr` the most: a range function such + /// as `rate` derives its result from the evaluation instant it is called at, so folding the + /// window once per step would let the outer evaluation grid change the result of a window that + /// `@` fixed. Evaluated once, at the start of the grid, the rewritten offset of + /// [`Self::at_modifier_offset`] places the anchor at that instant, and [`Self::replay_over_grid`] + /// reports the result at every step. + /// + /// Unlike Prometheus, which wraps the whole step-invariant subtree (`preprocessExprHelper`), + /// only the call itself is promoted here. The operators above it are not: they are still planned + /// at every step over the replayed result, which is safe for the row-wise ones and keeps the + /// promotion root narrow. The promotion root has to stay a call over one range selector, because + /// the replay needs one series per batch ([`Self::series_divide_plan`]) and only such a call + /// guarantees that the rows it emits still describe the series it was divided by. The operators + /// left out — an aggregation, a join, or a label rewriting call such as `label_join` — mix or + /// re-label the rows of different series, so they are unsafe as promotion roots even though + /// evaluating them after the replay is fine. A call whose input is an anchored *instant* + /// selector (`abs(some_metric @ 300)`) needs no promotion either: the selector anchors and + /// replays its sample per series on its own, and the call above it is row-wise. + /// + /// `predict_linear` is the exception among the range functions: it predicts from the evaluation + /// instant of each step ([`Self::create_range_eval_ts_expr`]), so it has to stay outside the + /// promoted subtree and follow the grid. The remaining arguments of the call have to be + /// literals, since the replay of the promoted result has no second vector input to divide. + /// Parentheses around the range argument are transparent (`rate((m[5m] @ 300))`), so they are + /// looked through and the call is promoted as if they were absent. Nothing else of the subtree + /// is unwrapped, so an outer parenthesis promotes no operator above the call. + fn promotes_anchored_range_call(expr: &PromExpr) -> bool { + let PromExpr::Call(Call { func, args }) = expr else { + return false; + }; + // See the doc comment: the regression of `predict_linear` follows the evaluation step. + if func.name == "predict_linear" { + return false; + } + let mut anchored_range = false; + for arg in &args.args { + // Parentheses around the range argument are transparent, so the call is promoted the + // same way for `rate((m[5m] @ 300))` as for `rate(m[5m] @ 300)`. Only the parentheses + // directly around this one argument are looked through here: the promotion stays + // confined to a call over one anchored range selector instead of descending into an + // arbitrary parenthesized subtree. + let mut arg = arg.as_ref(); + while let PromExpr::Paren(ParenExpr { expr }) = arg { + arg = expr; + } + match arg { + // The window is pinned by `@`, so every step folds the same samples. + PromExpr::MatrixSelector(MatrixSelector { vs, .. }) if vs.at.is_some() => { + if anchored_range { + return false; + } + anchored_range = true; + } + // A literal argument is the same value at every step. + arg if Self::try_build_literal_expr(arg).is_some() => {} + _ => return false, + } + } + anchored_range + } + + /// Plans the anchored range call `prom_expr` ([`Self::promotes_anchored_range_call`]) as a + /// step-invariant subtree: the call is evaluated on a single evaluation instant (`grid_start`, + /// the start of the outer evaluation) and its result is then reported at every step of + /// `[grid_start, ctx.end]` by [`Self::replay_over_grid`]. + /// + /// This is the planner's counterpart of Prometheus' `StepInvariantExpr` for the one shape it + /// promotes. Evaluating the call once matters for the functions that derive their result from + /// the step being evaluated: `rate(m[5m] @ 300)` folds its window around the anchor once, and + /// the extrapolation boundaries of `rate` must be derived from that same window at every step + /// instead of following the outer evaluation timestamp. + /// + /// Only the call itself is promoted; the operators above it are planned as usual over the + /// replayed result. The result of the promoted call is split into one series per batch before it + /// is replayed ([`Self::series_divide_plan`]), because the row-wise projection of the call does + /// not preserve the batch layout of the selector. + /// + /// Returns `None` when `prom_expr` is not such a call, so that the caller plans it as usual. + /// The selector inside the promoted call keeps its own `@` anchoring (see + /// [`Self::at_modifier_offset`]), and planning it with `ctx.end == ctx.start` folds its window + /// once for that single instant instead of expanding it over the grid, which the replay of the + /// call result above already does. + pub(crate) async fn promote_anchored_range_call( + &mut self, + prom_expr: &PromExpr, + timestamp_fn: bool, + query_engine_state: &QueryEngineState, + ) -> Result> { + let grid_start = self.ctx.start; + let grid_end = self.ctx.end; + // An instant query evaluates a single step, so there is nothing to promote. + if grid_start == grid_end || !Self::promotes_anchored_range_call(prom_expr) { + return Ok(None); + } + + // Plan the subtree on a single evaluation instant: every selector below still anchors its + // window through `@`, and the functions above them derive their result from that one + // instant. The planner is single-use, so `ctx.end` needs no restore-on-error. + self.ctx.end = grid_start; + let anchored = self + .prom_expr_to_plan_inner(prom_expr, timestamp_fn, query_engine_state) + .await?; + self.ctx.end = grid_end; + + let time_index_column = + self.ctx + .time_index_column + .clone() + .with_context(|| TimeIndexNotFoundSnafu { + table: self.ctx.table_name.clone().unwrap_or_default(), + })?; + // The replay reads one series per batch; see [`Self::series_divide_plan`] for why the + // layout of the selector below the promoted call does not survive it. + let anchored = self.series_divide_plan(anchored, &time_index_column)?; + Ok(Some(self.replay_over_grid( + anchored, + grid_start, + grid_end, + time_index_column, + ))) + } + + /// Convert the timestamp of an `@` modifier into milliseconds since the Unix epoch. + fn system_time_to_millis(time: &SystemTime) -> Result { + let (millis, negative) = match time.duration_since(UNIX_EPOCH) { + Ok(duration) => (duration.as_millis(), false), + // The `@` modifier accepts timestamps before the Unix epoch, e.g. `@ -1`. + Err(err) => (err.duration().as_millis(), true), + }; + ensure!( + millis <= i64::MAX as u128, + AtModifierTimestampOutOfRangeSnafu { + timestamp: time + .duration_since(UNIX_EPOCH) + .map(|duration| format!("+{}ms", duration.as_millis())) + .unwrap_or_else(|err| format!("-{}ms", err.duration().as_millis())), + } + ); + let millis = millis as Millisecond; + Ok(if negative { -millis } else { millis }) + } + + /// Report the samples of `anchored` at every step of the evaluation grid + /// `[grid_start, grid_end]`. + /// + /// A selector with an `@` modifier is anchored: the sample window is selected once, around the + /// anchor timestamp, and every evaluation step reports that same window. Prometheus does this + /// by rewriting the selector's offset to `eval_time - anchor` and only fetching the samples on + /// the first step (`setOffsetForAtModifier` plus the `refetch` shortcut in `rangeEval`). + /// + /// The expansion reuses [`InstantManipulate`] with a lookback that spans the whole grid: every + /// step then picks the same sample (or the same already computed value, when `anchored` ends + /// with a function call such as `rate`) and stamps it with the step's timestamp. + /// + /// Every input batch of `anchored` must hold exactly one series, because [`InstantManipulate`] + /// takes a batch as one timeline. A leaf-level replay (`m @ 300`) consumes the [`SeriesDivide`] + /// of its selector directly. A promoted call is guaranteed that layout by + /// [`Self::series_divide_plan`], which is why [`Self::promote_anchored_range_call`] splits + /// its result before calling this method. + pub(crate) fn replay_over_grid( + &self, + anchored: LogicalPlan, + grid_start: Millisecond, + grid_end: Millisecond, + time_index_column: String, + ) -> LogicalPlan { + if grid_start == grid_end { + // A single evaluation step: `anchored` is already stamped with that timestamp. + return anchored; + } + + let series_key_columns = self.series_key_columns_for_schema(anchored.schema()); + let replayed = InstantManipulate::new( + grid_start, + grid_end, + // The lookback must keep the single anchored sample eligible for every step. + grid_end - grid_start + 1, + self.ctx.interval, + 0, + time_index_column, + series_key_columns, + self.ctx.field_columns.first().cloned(), + anchored, + ); + LogicalPlan::Extension(Extension { + node: Arc::new(replayed), + }) + } + + /// Sorts `input` by its series key and time index and splits it into one batch per series. + /// + /// [`InstantManipulate`] reads every input batch as one series (it takes the timeline of the + /// batch and reports the row selected at every step), so a batch holding several series would + /// lose all but one of them. A selector establishes that layout with its own [`SeriesDivide`], + /// but the per-series distribution requirement does not reach a promoted subtree above it: + /// the row-wise projection of a call sits in between, so the batch boundaries of the selector + /// are not preserved — in a distributed plan the promoted result can be delivered as one batch + /// holding every series. Sorting and dividing here restores the layout, exactly like + /// [`Self::prom_matrix_selector_to_plan`] does for the input of a range function. + /// + /// Series keys that are not present in `input` are dropped, since `ctx.tag_columns` may have + /// drifted from the actual output schema. A plan without any series key column is returned + /// unchanged: there is nothing to divide by. + fn series_divide_plan( + &self, + input: LogicalPlan, + time_index_column: &str, + ) -> Result { + let series_key_columns = self.series_key_columns_for_schema(input.schema()); + if series_key_columns.is_empty() { + return Ok(input); + } + + let mut sort_exprs = series_key_columns + .iter() + .map(|name| DfExpr::Column(Column::from_name(name)).sort(true, true)) + .collect::>(); + sort_exprs.push(DfExpr::Column(Column::from_name(time_index_column)).sort(true, true)); + let sort_plan = LogicalPlanBuilder::from(input) + .sort(sort_exprs) + .context(DataFusionPlanningSnafu)? + .build() + .context(DataFusionPlanningSnafu)?; + Ok(LogicalPlan::Extension(Extension { + node: Arc::new(SeriesDivide::new( + series_key_columns, + time_index_column.to_string(), + sort_plan, + )), + })) + } +} diff --git a/src/query/src/promql/planner/test.rs b/src/query/src/promql/planner/test.rs index 1634e6d9e1d..a83a42bae53 100644 --- a/src/query/src/promql/planner/test.rs +++ b/src/query/src/promql/planner/test.rs @@ -1929,6 +1929,600 @@ async fn tsid_is_used_for_series_divide_when_available() { assert!(format!("{exec:?}").contains("reuse_tsid_column: true")); } +async fn build_at_modifier_plan(query: &str, start_secs: u64, end_secs: u64) -> LogicalPlan { + let eval_stmt = build_at_modifier_eval_stmt(query, start_secs, end_secs); + let table_provider = build_test_table_provider( + &[(DEFAULT_SCHEMA_NAME.to_string(), "some_metric".to_string())], + 1, + 1, + ) + .await; + PromPlanner::stmt_to_plan(table_provider, &eval_stmt, &build_query_engine_state()) + .await + .unwrap() +} + +fn build_at_modifier_eval_stmt(query: &str, start_secs: u64, end_secs: u64) -> EvalStmt { + EvalStmt { + expr: parser::parse(query).unwrap(), + start: UNIX_EPOCH + .checked_add(Duration::from_secs(start_secs)) + .unwrap(), + end: UNIX_EPOCH + .checked_add(Duration::from_secs(end_secs)) + .unwrap(), + interval: Duration::from_secs(5), + lookback_delta: Duration::from_secs(1), + } +} + +/// Every selector of an `@` anchored query must scan around the anchor only, instead of the +/// whole evaluation range. +#[tokio::test] +async fn at_modifier_anchors_selector_scan_window() { + // `@ 100` anchors at t=100s; the lookback delta is 1s, so the scan covers (99s, 100s]. + let plan = build_at_modifier_plan("some_metric @ 100", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert!( + plan_str.contains( + "some_metric.timestamp >= TimestampMillisecond(99001, None) AND some_metric.timestamp <= TimestampMillisecond(100000, None)" + ), + "{plan_str}" + ); + // The result is reported at the evaluation timestamps, not at the anchor. + assert!(plan_str.contains("range=[0..1000000]"), "{plan_str}"); + + // `@ start()` / `@ end()` resolve to the evaluation range of the statement. + let plan = build_at_modifier_plan("some_metric @ start()", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert!( + plan_str.contains( + "some_metric.timestamp >= TimestampMillisecond(-999, None) AND some_metric.timestamp <= TimestampMillisecond(0, None)" + ), + "{plan_str}" + ); + + let plan = build_at_modifier_plan("some_metric @ end()", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert!( + plan_str.contains( + "some_metric.timestamp >= TimestampMillisecond(999001, None) AND some_metric.timestamp <= TimestampMillisecond(1000000, None)" + ), + "{plan_str}" + ); + + // `offset` moves the anchor backwards and is not applied twice. + let plan = build_at_modifier_plan("some_metric @ 200 offset 50s", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert!( + plan_str.contains( + "some_metric.timestamp >= TimestampMillisecond(149001, None) AND some_metric.timestamp <= TimestampMillisecond(150000, None)" + ), + "{plan_str}" + ); + + // A timestamp before the Unix epoch is accepted, as in Prometheus. + let plan = build_at_modifier_plan("some_metric @ -1", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert!( + plan_str.contains( + "some_metric.timestamp >= TimestampMillisecond(-1999, None) AND some_metric.timestamp <= TimestampMillisecond(-1000, None)" + ), + "{plan_str}" + ); +} + +/// A call over one range selector anchored by `@` is evaluated once, at the start of the +/// evaluation, and its result is reported at every step: the window is folded around the anchor +/// instead of following the outer evaluation grid. This is the planner's counterpart of +/// Prometheus' `StepInvariantExpr` wrapper; see [`PromPlanner::promotes_anchored_range_call`]. +#[tokio::test] +async fn at_modifier_promotes_anchored_range_call() { + let plan = build_at_modifier_plan("rate(some_metric[5m] @ 300)", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + // The scan is limited to the anchored window (offset by `eval_start - anchor`). + assert!( + plan_str.contains( + "some_metric.timestamp >= TimestampMillisecond(1, None) AND some_metric.timestamp <= TimestampMillisecond(300000, None)" + ), + "{plan_str}" + ); + // A single fold, at the anchor... + assert_eq!( + plan_str + .matches("PromRangeManipulate: req range=[0..0]") + .count(), + 1, + "{plan_str}" + ); + // ... never a fold per step of the outer grid. + assert!( + !plan_str.contains("PromRangeManipulate: req range=[0..1000000]"), + "{plan_str}" + ); + // A single replay of the function result over the whole grid... + assert_eq!( + plan_str.matches("PromInstantManipulate").count(), + 1, + "{plan_str}" + ); + assert!( + plan_str.contains("PromInstantManipulate: range=[0..1000000], lookback=[1000001]"), + "{plan_str}" + ); + // ... so `rate` itself is evaluated below that replay node, on the single evaluation + // instant of the anchored subtree, instead of once per step. + let replay = plan_str + .find("PromInstantManipulate") + .expect("instant manipulate node"); + let rate = plan_str.find("prom_rate(").expect("rate projection"); + assert!( + replay < rate, + "`rate` must be evaluated below the replay node:\n{plan_str}" + ); +} + +/// Parentheses around the range argument are transparent: `rate((some_metric[5m] @ 300))` gets +/// the same fixed-window promotion as `rate(some_metric[5m] @ 300)`, with the anchored window +/// folded once and its result replayed at every step of the grid. Only that one argument is +/// looked through, so a parenthesis above the call promotes no operator of its own and a +/// parenthesized subtree is planned exactly like the bare one. +#[tokio::test] +async fn at_modifier_promotes_parenthesized_range_argument() { + // Each form is planned exactly like its unparenthesized counterpart: parentheses below the + // call are transparent, a parenthesis around the call adds nothing, and a parenthesis above + // it does not widen the promotion. + for (query, plain) in [ + ( + "rate((some_metric[5m] @ 300))", + "rate(some_metric[5m] @ 300)", + ), + ( + "rate(((some_metric[5m] @ 300)))", + "rate(some_metric[5m] @ 300)", + ), + ( + "(rate(some_metric[5m] @ 300))", + "rate(some_metric[5m] @ 300)", + ), + ( + "abs((rate(some_metric[5m] @ 300)))", + "abs(rate(some_metric[5m] @ 300))", + ), + ] { + assert_eq!( + build_at_modifier_plan(query, 0, 1000) + .await + .display_indent_schema() + .to_string(), + build_at_modifier_plan(plain, 0, 1000) + .await + .display_indent_schema() + .to_string(), + "`{query}` must be planned like `{plain}`" + ); + } + + // The parentheses do not push the enclosing operator into the promotion either: `abs` stays + // above the replay of the promoted call, exactly as it does without them. The promotion + // itself (one anchored fold, one replay, `rate` below it) is asserted for the + // unparenthesized form by `at_modifier_promotes_anchored_range_call`, and the form above is + // planned identically to it. + let plan_str = build_at_modifier_plan("abs((rate(some_metric[5m] @ 300)))", 0, 1000) + .await + .display_indent_schema() + .to_string(); + let replay = plan_str + .find("PromInstantManipulate") + .expect("instant manipulate node"); + assert!( + plan_str.find("abs(").expect("`abs` projection") < replay, + "`abs` must be evaluated above the replay of the promoted call:\n{plan_str}" + ); +} + +/// `@ start()` and `@ end()` are fixed anchors for the whole statement, so a call using them is +/// promoted as well. +#[tokio::test] +async fn at_modifier_promotes_start_and_end_anchored_call() { + for query in [ + "rate(some_metric[5m] @ start())", + "rate(some_metric[5m] @ end())", + "max_over_time(some_metric[5m] @ end())", + ] { + let plan = build_at_modifier_plan(query, 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert_eq!( + plan_str + .matches("PromRangeManipulate: req range=[0..0]") + .count(), + 1, + "{query}:\n{plan_str}" + ); + assert!( + plan_str.contains("PromInstantManipulate: range=[0..1000000], lookback=[1000001]"), + "{query}:\n{plan_str}" + ); + } + + // Only the call itself is promoted: a call above it (`abs`) is planned as usual and + // evaluated at every step over the replayed result of the promoted call. The window is + // still folded once, around the anchor. + let plan = build_at_modifier_plan("abs(max_over_time(some_metric[5m] @ end()))", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert_eq!( + plan_str + .matches("PromRangeManipulate: req range=[0..0]") + .count(), + 1, + "{plan_str}" + ); + let abs = plan_str.find("abs(").expect("`abs` projection"); + let replay = plan_str + .find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .expect("replay node"); + assert!( + abs < replay, + "`abs` must be evaluated above the replay of the promoted call:\n{plan_str}" + ); + // The window is the only one folded once, and it sits below the replay: the promoted call + // feeds the grid from there. + let fold = plan_str.find("PromRangeManipulate").expect("range fold"); + assert!( + replay < fold, + "the anchored window must be folded below the replay:\n{plan_str}" + ); + + // An aggregation above the promoted call stays above it as well: `sum` aggregates the + // replayed per-series rows at every step, instead of aggregating the single anchored + // instant and replaying the aggregation — which would also have to replay rows of several + // groups through the one-series-per-batch `InstantManipulate`. + let plan = build_at_modifier_plan("sum(rate(some_metric[5m] @ start()))", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert_eq!( + plan_str + .matches("PromRangeManipulate: req range=[0..0]") + .count(), + 1, + "{plan_str}" + ); + // A single replay, and it sits below the aggregation node: `sum` aggregates the replayed + // per-series rows at every step. + assert_eq!( + plan_str.matches("PromInstantManipulate").count(), + 1, + "{plan_str}" + ); + let aggregate = plan_str.find("Aggregate:").expect("aggregate node"); + let replay = plan_str + .find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .expect("replay node"); + assert!( + aggregate < replay, + "the aggregation must stay above the replay of the promoted call:\n{plan_str}" + ); + + // Without `@` nothing is promoted: the function keeps folding one window per step. + let plan = build_at_modifier_plan("rate(some_metric[5m])", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert!( + plan_str.contains("PromRangeManipulate: req range=[0..1000000]"), + "{plan_str}" + ); + assert!(!plan_str.contains("lookback=[1000001]"), "{plan_str}"); +} + +/// A call or a unary operator above the anchored range call is planned as usual: the inner range +/// call is promoted on its own (it is the direct call over the anchored range selector), and the +/// operator above it is evaluated at every step over the replayed result. +#[tokio::test] +async fn at_modifier_promotes_inner_range_call_below_wrappers() { + for (query, wrapper) in [ + ("abs(rate(some_metric[5m] @ 300))", "abs(prom_rate("), + ("-rate(some_metric[5m] @ 300)", "(- prom_rate("), + ( + "abs(max_over_time(some_metric[5m] @ 300))", + "abs(prom_max_over_time(", + ), + ] { + let plan = build_at_modifier_plan(query, 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + // The anchored window is folded once, and the wrapper sits above its replay. + assert_eq!( + plan_str + .matches("PromRangeManipulate: req range=[0..0]") + .count(), + 1, + "{query}:\n{plan_str}" + ); + assert_eq!( + plan_str.matches("PromInstantManipulate").count(), + 1, + "{query}:\n{plan_str}" + ); + let wrapper = plan_str + .find(wrapper) + .unwrap_or_else(|| panic!("no `{wrapper}` projection in:\n{plan_str}")); + let replay = plan_str + .find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .expect("replay node"); + assert!( + wrapper < replay, + "the wrapper must be evaluated above the replay of the range call:\n{plan_str}" + ); + } +} + +/// `anchored + plain`: only the binary operand that is a call over the anchored range selector +/// is promoted, and the plain side keeps following the evaluation step. +#[tokio::test] +async fn at_modifier_promotes_only_anchored_binary_operand() { + let plan = build_at_modifier_plan("rate(some_metric[5m] @ 300) + some_metric", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + // The anchored operand is folded once and its result replayed over the whole grid. + assert_eq!( + plan_str + .matches("PromRangeManipulate: req range=[0..0]") + .count(), + 1, + "{plan_str}" + ); + assert!( + plan_str.contains("PromInstantManipulate: range=[0..1000000], lookback=[1000001]"), + "{plan_str}" + ); + // The plain operand still selects one sample per step with the default lookback. + assert!( + plan_str.contains("PromInstantManipulate: range=[0..1000000], lookback=[1000]"), + "{plan_str}" + ); + assert_eq!( + plan_str.matches("PromInstantManipulate").count(), + 2, + "{plan_str}" + ); + + // Both operands anchored: the binary expression is not promoted as a whole, because a join + // emits the rows of several series in shared batches and the replay needs one series per + // batch. Each operand anchors and replays on its own instead, and the join runs at every + // step over those per-series results. + let plan = build_at_modifier_plan("some_metric @ 300 + some_metric @ 0", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + // One anchoring node and one replay per operand... + assert_eq!( + plan_str.matches("PromInstantManipulate").count(), + 4, + "{plan_str}" + ); + assert_eq!( + plan_str + .matches("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .count(), + 2, + "{plan_str}" + ); + // ... and the two operands are anchored at different timestamps, so both select their own + // sample. + assert_eq!( + plan_str + .matches("PromInstantManipulate: range=[0..0], lookback=[1000]") + .count(), + 2, + "{plan_str}" + ); +} + +/// Only a direct call over an anchored range selector is promoted. A value function or a unary +/// operator over an anchored *instant* selector needs no promotion: the selector anchors and +/// replays its sample per series on its own, and the operator above it is row-wise, so it can be +/// evaluated at every step over that replay. +#[tokio::test] +async fn at_modifier_does_not_promote_value_calls_over_anchored_selectors() { + for (query, value_expr) in [ + ("abs(some_metric @ 300)", "abs(some_metric.field_0)"), + ("-some_metric @ 300", "(- some_metric.field_0)"), + ] { + let plan = build_at_modifier_plan(query, 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + // The scan is limited to the anchored sample (the lookback delta of this test is 1s). + assert!( + plan_str.contains( + "some_metric.timestamp >= TimestampMillisecond(299001, None) AND some_metric.timestamp <= TimestampMillisecond(300000, None)" + ), + "{query}:\n{plan_str}" + ); + // ... and the result is replayed at every step by the selector itself: the anchored + // selection, then the grid replay, with no promoted subtree on top of the operator. + assert!( + !plan_str.starts_with("PromInstantManipulate"), + "{query}:\n{plan_str}" + ); + assert_eq!( + plan_str + .matches("PromInstantManipulate: range=[0..0], lookback=[1000]") + .count(), + 1, + "{query}:\n{plan_str}" + ); + let replay = plan_str + .find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .unwrap_or_else(|| panic!("replay node:\n{plan_str}")); + let value = plan_str + .find(value_expr) + .unwrap_or_else(|| panic!("no `{value_expr}` projection in:\n{plan_str}")); + assert!( + value < replay, + "`{value_expr}` must be evaluated above the per-series replay:\n{plan_str}" + ); + } +} + +/// A call whose argument merges series (an aggregation or a join) or whose own output reorders +/// the whole vector (`sort*`, the histogram folds) is never the promoted root either: it is +/// planned as usual over the leaf-level anchoring of its selectors, which replays every selector +/// per series and keeps the one-series-per-batch layout the replay needs. +#[tokio::test] +async fn at_modifier_keeps_multi_series_roots_out_of_promoted_subtree() { + // An aggregation below the call emits one row per group in shared batches, so the call is + // not promoted: the replay of the anchored selector stays below the aggregate. + let plan = build_at_modifier_plan("abs(sum(some_metric @ 300))", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + let aggregate = plan_str.find("Aggregate:").expect("aggregate node"); + let replay = plan_str + .find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .expect("replay node"); + assert!( + aggregate < replay, + "the aggregation must stay above the replay:\n{plan_str}" + ); + + // A join of two anchored selectors below the call: neither the join nor the call is + // promoted, so each operand is replayed on its own. + let plan = build_at_modifier_plan("abs(some_metric @ 300 + some_metric @ 0)", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert!(!plan_str.starts_with("PromInstantManipulate"), "{plan_str}"); + assert_eq!( + plan_str + .matches("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .count(), + 2, + "{plan_str}" + ); + + // A call that reorders the whole vector keeps its sort above the per-series replay. + let plan = build_at_modifier_plan("sort_by_label(some_metric @ 300, \"tag_0\")", 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + assert!(plan_str.starts_with("Sort:"), "{plan_str}"); + assert_eq!( + plan_str + .matches("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .count(), + 1, + "{plan_str}" + ); +} + +/// `label_join` rewrites the labels of its input series, so it is never the promoted root: the +/// anchored selector keeps replaying one series per batch, and the join runs at every step above +/// that replay (see [`Self::promotes_anchored_range_call`]). Promoting it would replay the +/// joined rows through the labels the join just rewrote, which merges the distinct input series +/// into one timeline. +#[tokio::test] +async fn at_modifier_does_not_promote_label_join() { + for query in [ + // Directly above the anchored instant selector... + "label_join(some_metric @ 300, \"tag_0\", \"-\", \"\")", + // ... and below another call, which is planned as usual over the join. + "abs(label_join(some_metric @ 300, \"tag_0\", \"-\", \"\"))", + ] { + let plan = build_at_modifier_plan(query, 0, 1000).await; + let plan_str = plan.display_indent_schema().to_string(); + // The join is not wrapped in a replay of its own; the only replay over the grid is the + // one of the anchored selector... + assert!( + !plan_str.starts_with("PromInstantManipulate"), + "{query}:\n{plan_str}" + ); + assert_eq!( + plan_str + .matches("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .count(), + 1, + "{query}:\n{plan_str}" + ); + // ... and the projected join stays above it, evaluated at every step. + let join = plan_str + .find("concat_ws(") + .unwrap_or_else(|| panic!("no `label_join` projection in:\n{plan_str}")); + let replay = plan_str + .find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .expect("replay node"); + assert!( + join < replay, + "`label_join` must be evaluated above the per-series replay:\n{plan_str}" + ); + } + + // A range call below the join is still promoted on its own: the anchored window is folded + // once per series, and the join above it is evaluated at every step over that replay. + let plan = build_at_modifier_plan( + "label_join(rate(some_metric[5m] @ 300), \"tag_0\", \"-\", \"\")", + 0, + 1000, + ) + .await; + let plan_str = plan.display_indent_schema().to_string(); + assert_eq!( + plan_str + .matches("PromRangeManipulate: req range=[0..0]") + .count(), + 1, + "{plan_str}" + ); + assert!(!plan_str.starts_with("PromInstantManipulate"), "{plan_str}"); + let join = plan_str.find("concat_ws(").expect("join projection"); + let replay = plan_str + .find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]") + .expect("replay node"); + assert!( + join < replay, + "the join must be evaluated above the replay of the range call:\n{plan_str}" + ); +} + +#[test] +fn at_modifier_rejects_subtraction_overflow() { + for (anchor, offset) in [(i64::MAX, -1), (i64::MIN, 1)] { + let err = PromPlanner::anchor_sub(anchor, offset).unwrap_err(); + assert_eq!(err.status_code(), StatusCode::InvalidArguments); + assert!( + err.to_string() + .contains("Timestamp out of range for the `@` modifier"), + "{err}" + ); + } + assert_eq!(PromPlanner::anchor_sub(-1, 1).unwrap(), -2); +} + +/// `@` beyond the representable millisecond range is rejected instead of silently wrapping. +/// +/// `@ 1e16` is 10^19 milliseconds, beyond `i64::MAX`. A Unix `SystemTime` can hold it, so +/// the planner rejects the anchor it cannot represent. A Windows `SystemTime` tops out +/// below `i64::MAX` milliseconds, so the same literal is already rejected while parsing. +#[tokio::test] +async fn at_modifier_rejects_unrepresentable_timestamp() { + #[cfg(windows)] + { + let err = parser::parse("some_metric @ 1e16").unwrap_err(); + assert!( + err.to_string() + .contains("timestamp out of bounds for @ modifier"), + "{err}" + ); + } + + #[cfg(not(windows))] + { + let eval_stmt = build_eval_stmt("some_metric @ 1e16"); + let table_provider = build_test_table_provider( + &[(DEFAULT_SCHEMA_NAME.to_string(), "some_metric".to_string())], + 1, + 1, + ) + .await; + let err = + PromPlanner::stmt_to_plan(table_provider, &eval_stmt, &build_query_engine_state()) + .await + .unwrap_err(); + assert!( + err.to_string() + .contains("Timestamp out of range for the `@` modifier"), + "{err}" + ); + assert_eq!(err.status_code(), StatusCode::InvalidArguments); + } +} + #[tokio::test] async fn default_binary_join_uses_tsid_when_available() { let eval_stmt = build_eval_stmt("some_metric / some_alt_metric"); @@ -3421,7 +4015,7 @@ async fn binary_op_column_column() { assert_eq!(plan.display_indent_schema().to_string(), expected); } -async fn indie_query_plan_compare>(query: &str, expected: T) { +async fn indie_query_plan(query: &str) -> String { let prom_expr = parser::parse(query).unwrap(); let eval_stmt = EvalStmt { expr: prom_expr, @@ -3449,7 +4043,12 @@ async fn indie_query_plan_compare>(query: &str, expected: T) { .await .unwrap(); - assert_eq!(plan.display_indent_schema().to_string(), expected.as_ref()); + plan.display_indent_schema().to_string() +} + +async fn indie_query_plan_compare>(query: &str, expected: T) { + let plan = indie_query_plan(query).await; + assert_eq!(plan, expected.as_ref()); } #[tokio::test] @@ -3537,6 +4136,51 @@ async fn increase_aggr() { indie_query_plan_compare(query, expected).await; } +#[tokio::test] +async fn predict_linear_injects_the_eval_timestamp() { + // The fourth argument is the evaluation instant of each row, which the fold offset + // recovers from the row's time index (here: no `@` and no `offset`, so the step itself). + let query = "predict_linear(some_metric[5m], 60)"; + let expected = String::from( + "Filter: prom_predict_linear(timestamp_range,field_0,Float64(60)) IS NOT NULL [timestamp:Timestamp(ms), prom_predict_linear(timestamp_range,field_0,Float64(60)):Float64;N, tag_0:Utf8]\ + \n Projection: some_metric.timestamp, prom_predict_linear(timestamp_range, field_0, CAST(Float64(60) AS Int64), CAST(CAST(some_metric.timestamp AS Int64) + Int64(0) AS Timestamp(ms))) AS prom_predict_linear(timestamp_range,field_0,Float64(60)), some_metric.tag_0 [timestamp:Timestamp(ms), prom_predict_linear(timestamp_range,field_0,Float64(60)):Float64;N, tag_0:Utf8]\ + \n PromRangeManipulate: req range=[0..100000000], interval=[5000], eval range=[300000], 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=[0], time index=[timestamp], filter NaN: [true] [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(-299999, None) AND some_metric.timestamp <= TimestampMillisecond(100000000, 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; +} + +/// The evaluation instant follows the selector's fold offset, not the projection's: an +/// `@`-anchored selector is folded with `at_offset`, which is what the window's timestamps were +/// shifted by. The evaluation starts at 0s here, so `@ 100` anchors 100s in the future and +/// yields `at_offset` = -100000ms. +#[tokio::test] +async fn predict_linear_eval_ts_follows_the_fold_offset() { + for (query, expected_offset) in [ + ( + "predict_linear(some_metric[5m] @ 100, 60)", + "Int64(-100000)", + ), + ( + "predict_linear(some_metric[5m] offset 2m, 60)", + "Int64(120000)", + ), + ] { + let plan = indie_query_plan(query).await; + assert!( + plan.contains(&format!( + "some_metric.timestamp AS Int64) + {expected_offset}" + )), + "{query}\n{plan}" + ); + } +} + async fn native_histogram_plan(query: &str) -> String { let table_provider = build_test_native_histogram_table_provider("some_metric").await; let plan = PromPlanner::stmt_to_plan( @@ -4214,6 +4858,27 @@ async fn mixed_native_histogram_ranges_use_coordinated_udfs() { assert_eq!(plan, expected); } +#[tokio::test] +async fn mixed_native_histogram_predict_linear_forwards_the_eval_timestamp() { + let query = "predict_linear(some_metric[5m], 60)"; + let plan = PromPlanner::stmt_to_plan( + build_test_mixed_native_histogram_table_provider("some_metric").await, + &build_eval_stmt(query), + &build_query_engine_state(), + ) + .await + .unwrap() + .display_indent_schema() + .to_string(); + + assert!( + plan.contains( + "prom_mixed_range_float(Utf8(\"predict_linear\"), timestamp_range, greptime_value, greptime_native_histogram, CAST(Float64(60) AS Int64), CAST(CAST(some_metric.timestamp AS Int64) + Int64(0) AS Timestamp(ms)))" + ), + "{query}\n{plan}" + ); +} + #[tokio::test] async fn mixed_native_histogram_rate_executes_real_ranges() { let schema = Arc::new(ArrowSchema::new(vec![ @@ -4288,7 +4953,7 @@ async fn mixed_native_histogram_rate_executes_real_ranges() { ); let state = build_query_engine_state(); let (mut exprs, _) = planner - .create_function_expr(&call.func, vec![], input.schema(), &state) + .create_function_expr(&call.func, vec![], input.schema(), &state, None) .unwrap(); exprs.insert(0, planner.create_time_index_column_expr().unwrap()); let plan = LogicalPlanBuilder::from(input) @@ -7124,7 +7789,7 @@ async fn test_mixed_or_routes_float_histogram_and_label_functions() { }; let state = build_query_engine_state(); let (mut exprs, _) = planner - .create_function_expr(&call.func, vec![], input.schema(), &state) + .create_function_expr(&call.func, vec![], input.schema(), &state, None) .unwrap(); exprs.insert(0, planner.create_time_index_column_expr().unwrap()); exprs.extend(planner.create_tag_column_exprs().unwrap()); @@ -7177,7 +7842,7 @@ async fn test_mixed_or_routes_float_histogram_and_label_functions() { let args = planner.create_function_args(&call.args.args).unwrap(); let state = build_query_engine_state(); let (mut exprs, _) = planner - .create_function_expr(&call.func, args.literals, input.schema(), &state) + .create_function_expr(&call.func, args.literals, input.schema(), &state, None) .unwrap(); exprs.insert(0, planner.create_time_index_column_expr().unwrap()); exprs.extend(planner.create_tag_column_exprs().unwrap()); diff --git a/src/query/src/promql/planner/test/delta.rs b/src/query/src/promql/planner/test/delta.rs index 3727ea9a30c..a647821e0b3 100644 --- a/src/query/src/promql/planner/test/delta.rs +++ b/src/query/src/promql/planner/test/delta.rs @@ -475,7 +475,7 @@ async fn delta_mixed_ranges_drop_and_float_ranges_sum() { ); let state = build_query_engine_state(); let (mut exprs, _) = planner - .create_function_expr(&call.func, vec![], input.schema(), &state) + .create_function_expr(&call.func, vec![], input.schema(), &state, None) .unwrap(); exprs.insert(0, planner.create_time_index_column_expr().unwrap()); let plan = LogicalPlanBuilder::from(input) diff --git a/tests/cases/standalone/common/promql/at_modifier.result b/tests/cases/standalone/common/promql/at_modifier.result new file mode 100644 index 00000000000..0a19eac692c --- /dev/null +++ b/tests/cases/standalone/common/promql/at_modifier.result @@ -0,0 +1,743 @@ +-- Tests for the PromQL `@` modifier on vector and matrix selectors. +-- +-- `@` anchors the sample selection window at a fixed timestamp instead of the timestamp of each +-- evaluation step: `@ ` uses the given timestamp, `@ start()` / `@ end()` use the +-- start/end of the statement's evaluation range, and `offset` shifts the anchor backwards before +-- the window is selected. The output timestamps still follow the evaluation grid. +-- +-- Every metric below carries two series, host 'a' and host 'b', with different values: the +-- anchored results must therefore report both of them, per step. A series silently dropped, or a +-- batch holding several series handled as one timeline, does not look like a correct single-series +-- answer here but shows up as a missing or shifted row. +-- +-- Sample timestamps below are in milliseconds, while `@` and `TQL EVAL` timestamps are in seconds, +-- as in Prometheus. +CREATE TABLE at_modifier_gauge ( + ts TIMESTAMP TIME INDEX, + val DOUBLE, + host STRING, + PRIMARY KEY(host) +); + +Affected Rows: 0 + +CREATE TABLE at_modifier_counter_total ( + ts TIMESTAMP TIME INDEX, + val DOUBLE, + host STRING, + PRIMARY KEY(host) +); + +Affected Rows: 0 + +-- One sample per minute, so the anchored sample is easy to identify. The default lookback is 5m. +-- Host 'a' reports 0.0 to 6.0, host 'b' reports 7.0 to 13.0: the two series are told apart by their +-- values alone, which stays true after any anchoring. +INSERT INTO at_modifier_gauge VALUES + (0, 0.0, 'a'), + (60000, 1.0, 'a'), + (120000, 2.0, 'a'), + (180000, 3.0, 'a'), + (240000, 4.0, 'a'), + (300000, 5.0, 'a'), + (360000, 6.0, 'a'), + (0, 7.0, 'b'), + (60000, 8.0, 'b'), + (120000, 9.0, 'b'), + (180000, 10.0, 'b'), + (240000, 11.0, 'b'), + (300000, 12.0, 'b'), + (360000, 13.0, 'b'); + +Affected Rows: 14 + +-- Host 'a': a counter increasing by 60 per minute, i.e. by 1 per second. +-- Host 'b': a counter starting at 100 and increasing by 30 per minute, i.e. by 0.5 per second. Its +-- rate differs from host 'a', so a window folded for the wrong series cannot pass as the right one. +INSERT INTO at_modifier_counter_total VALUES + (0, 0.0, 'a'), + (60000, 60.0, 'a'), + (120000, 120.0, 'a'), + (180000, 180.0, 'a'), + (240000, 240.0, 'a'), + (300000, 300.0, 'a'), + (360000, 360.0, 'a'), + (0, 100.0, 'b'), + (60000, 130.0, 'b'), + (120000, 160.0, 'b'), + (180000, 190.0, 'b'), + (240000, 220.0, 'b'), + (300000, 250.0, 'b'), + (360000, 280.0, 'b'); + +Affected Rows: 14 + +-- 1. Instant query anchored at an absolute timestamp: the sample window is `(anchor - lookback, +-- anchor]` = `(0s, 300s]`, so its newest samples are 5.0 for host 'a' and 12.0 for host 'b', both +-- at 300s. The output timestamps stay the evaluation timestamps. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 300; + ++---------------------+------+------+ +| ts | val | host | ++---------------------+------+------+ +| 1970-01-01T00:01:40 | 12.0 | b | +| 1970-01-01T00:01:40 | 5.0 | a | ++---------------------+------+------+ + +-- Control: the same query without `@` looks back from the evaluation timestamp instead: +-- `(-200s, 100s]` selects the samples 1.0 (a) and 8.0 (b) at 60s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge; + ++---------------------+-----+------+ +| ts | val | host | ++---------------------+-----+------+ +| 1970-01-01T00:01:40 | 1.0 | a | +| 1970-01-01T00:01:40 | 8.0 | b | ++---------------------+-----+------+ + +-- The anchor, not the evaluation timestamp, decides which sample is reported. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 360; + ++---------------------+------+------+ +| ts | val | host | ++---------------------+------+------+ +| 1970-01-01T00:01:40 | 13.0 | b | +| 1970-01-01T00:01:40 | 6.0 | a | ++---------------------+------+------+ + +-- Sub-second precision is preserved in the anchor: `240.999` selects the samples at 240s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 240.999; + ++---------------------+------+------+ +| ts | val | host | ++---------------------+------+------+ +| 1970-01-01T00:01:40 | 11.0 | b | +| 1970-01-01T00:01:40 | 4.0 | a | ++---------------------+------+------+ + +-- The window `(-300s, 0s]` includes the samples exactly at the anchor. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 0; + ++---------------------+-----+------+ +| ts | val | host | ++---------------------+-----+------+ +| 1970-01-01T00:01:40 | 0.0 | a | +| 1970-01-01T00:01:40 | 7.0 | b | ++---------------------+-----+------+ + +-- 2. `@ start()`: every step selects its samples around the start of the evaluation range +-- (`(-200s, 100s]`), so every step reports the same value per series (1.0 for 'a', 8.0 for 'b'). +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge @ start(); + ++---------------------+-----+------+ +| ts | val | host | ++---------------------+-----+------+ +| 1970-01-01T00:01:40 | 1.0 | a | +| 1970-01-01T00:01:40 | 8.0 | b | +| 1970-01-01T00:03:20 | 1.0 | a | +| 1970-01-01T00:03:20 | 8.0 | b | +| 1970-01-01T00:05:00 | 1.0 | a | +| 1970-01-01T00:05:00 | 8.0 | b | +| 1970-01-01T00:06:40 | 1.0 | a | +| 1970-01-01T00:06:40 | 8.0 | b | ++---------------------+-----+------+ + +-- Control: without `@` every step looks back from its own timestamp, so the steps differ. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge; + ++---------------------+------+------+ +| ts | val | host | ++---------------------+------+------+ +| 1970-01-01T00:01:40 | 1.0 | a | +| 1970-01-01T00:01:40 | 8.0 | b | +| 1970-01-01T00:03:20 | 10.0 | b | +| 1970-01-01T00:03:20 | 3.0 | a | +| 1970-01-01T00:05:00 | 12.0 | b | +| 1970-01-01T00:05:00 | 5.0 | a | +| 1970-01-01T00:06:40 | 13.0 | b | +| 1970-01-01T00:06:40 | 6.0 | a | ++---------------------+------+------+ + +-- 3. `@ end()`: every step selects its samples around the end of the evaluation range +-- (`(100s, 400s]`), whose newest samples are 6.0 (a) and 13.0 (b) at 360s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge @ end(); + ++---------------------+------+------+ +| ts | val | host | ++---------------------+------+------+ +| 1970-01-01T00:01:40 | 13.0 | b | +| 1970-01-01T00:01:40 | 6.0 | a | +| 1970-01-01T00:03:20 | 13.0 | b | +| 1970-01-01T00:03:20 | 6.0 | a | +| 1970-01-01T00:05:00 | 13.0 | b | +| 1970-01-01T00:05:00 | 6.0 | a | +| 1970-01-01T00:06:40 | 13.0 | b | +| 1970-01-01T00:06:40 | 6.0 | a | ++---------------------+------+------+ + +-- 4. Range selector anchored by `@`: the call directly above the anchored range selector is +-- evaluated once, at the start of the evaluation, and its result is reported at every step — the +-- planner's counterpart of Prometheus' `StepInvariantExpr`. The window `(0s, 300s]` is left-open, +-- so it holds 60 -> 300 for host 'a' and 130 -> 250 for host 'b', i.e. 240 and 120 units. The 60s +-- gap to the window start is below the extrapolation threshold (1.1 x the 60s sampling interval), +-- so `rate` adds it and divides by the 300s window: 240 * 1.25 / 300 = 1.0 and +-- 120 * 1.25 / 300 = 0.5 per second. Every step must report those same per-series values, and both +-- series must be there at every step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate(at_modifier_counter_total[5m] @ 300); + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:00:00 | 0.5 | b | +| 1970-01-01T00:00:00 | 1.0 | a | +| 1970-01-01T00:01:00 | 0.5 | b | +| 1970-01-01T00:01:00 | 1.0 | a | +| 1970-01-01T00:02:00 | 0.5 | b | +| 1970-01-01T00:02:00 | 1.0 | a | +| 1970-01-01T00:03:00 | 0.5 | b | +| 1970-01-01T00:03:00 | 1.0 | a | +| 1970-01-01T00:04:00 | 0.5 | b | +| 1970-01-01T00:04:00 | 1.0 | a | ++---------------------+------------------------------------------+------+ + +-- Consistency: the instant query below shows the anchored values over the same left-open +-- `(0s, 300s]` window (1.0 per second for 'a', 0.5 for 'b'). The `!= 1.0` comparison filters +-- host 'a', but host 'b' remains at every step with its rate of 0.5 per second. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 300, '1s') rate(at_modifier_counter_total[5m]); + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:05:00 | 0.5 | b | +| 1970-01-01T00:05:00 | 1.0 | a | ++---------------------+------------------------------------------+------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate(at_modifier_counter_total[5m] @ 300) != 1.0; + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:00:00 | 0.5 | b | +| 1970-01-01T00:01:00 | 0.5 | b | +| 1970-01-01T00:02:00 | 0.5 | b | +| 1970-01-01T00:03:00 | 0.5 | b | +| 1970-01-01T00:04:00 | 0.5 | b | ++---------------------+------------------------------------------+------+ + +-- Parentheses around the range argument are transparent: the call must be promoted exactly like +-- the unparenthesized one above, so every step reports the same anchored values for both series +-- (1.0 for host 'a', 0.5 for host 'b') instead of folding a window per step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate((at_modifier_counter_total[5m] @ 300)); + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:00:00 | 0.5 | b | +| 1970-01-01T00:00:00 | 1.0 | a | +| 1970-01-01T00:01:00 | 0.5 | b | +| 1970-01-01T00:01:00 | 1.0 | a | +| 1970-01-01T00:02:00 | 0.5 | b | +| 1970-01-01T00:02:00 | 1.0 | a | +| 1970-01-01T00:03:00 | 0.5 | b | +| 1970-01-01T00:03:00 | 1.0 | a | +| 1970-01-01T00:04:00 | 0.5 | b | +| 1970-01-01T00:04:00 | 1.0 | a | ++---------------------+------------------------------------------+------+ + +-- The same holds for `@ start()` and `@ end()`, which are also fixed anchors: they resolve to the +-- statement's evaluation range, so a range call using them is evaluated once as well. +-- The statement starts at 60s: the window `(-240s, 60s]` holds only the samples at 0s and 60s, so +-- host 'a' goes 0 -> 60 and host 'b' 100 -> 130. The far 240s gap to the window start is beyond the +-- extrapolation threshold (1.1 x the 60s interval), so only half an interval (30s) is added there; +-- host 'a' starts its counter at 0, which clamps that back: 60 / 300 = 0.2 and +-- 30 * 1.5 / 300 = 0.15 per second. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (60, 300, '60s') rate(at_modifier_counter_total[5m] @ start()); + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:01:00 | 0.15 | b | +| 1970-01-01T00:01:00 | 0.2 | a | +| 1970-01-01T00:02:00 | 0.15 | b | +| 1970-01-01T00:02:00 | 0.2 | a | +| 1970-01-01T00:03:00 | 0.15 | b | +| 1970-01-01T00:03:00 | 0.2 | a | +| 1970-01-01T00:04:00 | 0.15 | b | +| 1970-01-01T00:04:00 | 0.2 | a | +| 1970-01-01T00:05:00 | 0.15 | b | +| 1970-01-01T00:05:00 | 0.2 | a | ++---------------------+------------------------------------------+------+ + +-- The instant query at 60s folds the same window as every step of the query above. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (60, 60, '1s') rate(at_modifier_counter_total[5m]); + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:01:00 | 0.15 | b | +| 1970-01-01T00:01:00 | 0.2 | a | ++---------------------+------------------------------------------+------+ + +-- The statement ends at 240s: the window `(-60s, 240s]` holds the samples 0s to 240s, so host 'a' +-- goes 0 -> 240 and host 'b' 100 -> 220. The 60s gap to the window start is added for both, but +-- host 'a' starts at zero, which clamps it back (240 / 300 = 0.8 per second) while host 'b' keeps +-- it (120 * 1.25 / 300 = 0.5 per second). +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate(at_modifier_counter_total[5m] @ end()); + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:00:00 | 0.5 | b | +| 1970-01-01T00:00:00 | 0.8 | a | +| 1970-01-01T00:01:00 | 0.5 | b | +| 1970-01-01T00:01:00 | 0.8 | a | +| 1970-01-01T00:02:00 | 0.5 | b | +| 1970-01-01T00:02:00 | 0.8 | a | +| 1970-01-01T00:03:00 | 0.5 | b | +| 1970-01-01T00:03:00 | 0.8 | a | +| 1970-01-01T00:04:00 | 0.5 | b | +| 1970-01-01T00:04:00 | 0.8 | a | ++---------------------+------------------------------------------+------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (240, 240, '1s') rate(at_modifier_counter_total[5m]); + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:04:00 | 0.5 | b | +| 1970-01-01T00:04:00 | 0.8 | a | ++---------------------+------------------------------------------+------+ + +-- Control: without `@` each step folds its own window, so the values keep changing. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate(at_modifier_counter_total[5m]); + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:01:00 | 0.15 | b | +| 1970-01-01T00:01:00 | 0.2 | a | +| 1970-01-01T00:02:00 | 0.25 | b | +| 1970-01-01T00:02:00 | 0.4 | a | +| 1970-01-01T00:03:00 | 0.35000000000000003 | b | +| 1970-01-01T00:03:00 | 0.6000000000000001 | a | +| 1970-01-01T00:04:00 | 0.5 | b | +| 1970-01-01T00:04:00 | 0.8 | a | ++---------------------+------------------------------------------+------+ + +-- `@ start()` on a range selector: the window `(-300s, 0s]` is left-open, so it holds a single +-- sample per series — the one exactly at the anchor — and `count_over_time` returns 1 at every step +-- for both of them. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 300, '100s') count_over_time(at_modifier_counter_total[5m] @ start()); + ++---------------------+------------------------------------+------+ +| ts | prom_count_over_time(ts_range,val) | host | ++---------------------+------------------------------------+------+ +| 1970-01-01T00:00:00 | 1.0 | a | +| 1970-01-01T00:00:00 | 1.0 | b | +| 1970-01-01T00:01:40 | 1.0 | a | +| 1970-01-01T00:01:40 | 1.0 | b | +| 1970-01-01T00:03:20 | 1.0 | a | +| 1970-01-01T00:03:20 | 1.0 | b | +| 1970-01-01T00:05:00 | 1.0 | a | +| 1970-01-01T00:05:00 | 1.0 | b | ++---------------------+------------------------------------+------+ + +-- Control: without `@` the window keeps moving, so the count grows with the step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 300, '100s') count_over_time(at_modifier_counter_total[5m]); + ++---------------------+------------------------------------+------+ +| ts | prom_count_over_time(ts_range,val) | host | ++---------------------+------------------------------------+------+ +| 1970-01-01T00:00:00 | 1.0 | a | +| 1970-01-01T00:00:00 | 1.0 | b | +| 1970-01-01T00:01:40 | 2.0 | a | +| 1970-01-01T00:01:40 | 2.0 | b | +| 1970-01-01T00:03:20 | 4.0 | a | +| 1970-01-01T00:03:20 | 4.0 | b | +| 1970-01-01T00:05:00 | 5.0 | a | +| 1970-01-01T00:05:00 | 5.0 | b | ++---------------------+------------------------------------+------+ + +-- 5. `@` combined with `offset`: the offset moves the anchor backwards before the window is +-- selected, so `@ 300 offset 2m` reports the samples at 180s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 300 offset 2m; + ++---------------------+------+------+ +| ts | val | host | ++---------------------+------+------+ +| 1970-01-01T00:01:40 | 10.0 | b | +| 1970-01-01T00:01:40 | 3.0 | a | ++---------------------+------+------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 300 offset 1m; + ++---------------------+------+------+ +| ts | val | host | ++---------------------+------+------+ +| 1970-01-01T00:01:40 | 11.0 | b | +| 1970-01-01T00:01:40 | 4.0 | a | ++---------------------+------+------+ + +-- Control: the offset applies to the evaluation timestamp when there is no `@`. +TQL EVAL (100, 100, '1s') at_modifier_gauge offset 2m; + ++----+-----+------+ +| ts | val | host | ++----+-----+------+ ++----+-----+------+ + +-- 6. `@` inside a larger expression. The anchored operand contributes the same value at every +-- step, while the plain operand still follows the evaluation step. A binary expression is never +-- promoted: when one side is anchored and the other is not, the anchored side replays on its own, +-- and when both sides are anchored, each side reports its own anchored samples. The join then runs +-- at every step over the replayed per-series inputs. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge @ 300 + at_modifier_gauge @ 0; + ++------+---------------------+-------------------+ +| host | ts | lhs.val + rhs.val | ++------+---------------------+-------------------+ +| a | 1970-01-01T00:01:40 | 5.0 | +| a | 1970-01-01T00:03:20 | 5.0 | +| a | 1970-01-01T00:05:00 | 5.0 | +| a | 1970-01-01T00:06:40 | 5.0 | +| b | 1970-01-01T00:01:40 | 19.0 | +| b | 1970-01-01T00:03:20 | 19.0 | +| b | 1970-01-01T00:05:00 | 19.0 | +| b | 1970-01-01T00:06:40 | 19.0 | ++------+---------------------+-------------------+ + +-- The aggregation above the anchored selector runs at every step over the replayed (still +-- per-series) samples, so it sums both series at every step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') sum(at_modifier_gauge @ 0); + ++---------------------+----------------------------+ +| ts | sum(at_modifier_gauge.val) | ++---------------------+----------------------------+ +| 1970-01-01T00:01:40 | 7.0 | +| 1970-01-01T00:03:20 | 7.0 | +| 1970-01-01T00:05:00 | 7.0 | +| 1970-01-01T00:06:40 | 7.0 | ++---------------------+----------------------------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') abs(at_modifier_gauge @ 0); + ++---------------------+----------+------+ +| ts | abs(val) | host | ++---------------------+----------+------+ +| 1970-01-01T00:01:40 | 0.0 | a | +| 1970-01-01T00:01:40 | 7.0 | b | +| 1970-01-01T00:03:20 | 0.0 | a | +| 1970-01-01T00:03:20 | 7.0 | b | +| 1970-01-01T00:05:00 | 0.0 | a | +| 1970-01-01T00:05:00 | 7.0 | b | +| 1970-01-01T00:06:40 | 0.0 | a | +| 1970-01-01T00:06:40 | 7.0 | b | ++---------------------+----------+------+ + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge @ 300 + at_modifier_gauge; + ++------+---------------------+-------------------+ +| host | ts | lhs.val + rhs.val | ++------+---------------------+-------------------+ +| a | 1970-01-01T00:01:40 | 6.0 | +| a | 1970-01-01T00:03:20 | 8.0 | +| a | 1970-01-01T00:05:00 | 10.0 | +| a | 1970-01-01T00:06:40 | 11.0 | +| b | 1970-01-01T00:01:40 | 20.0 | +| b | 1970-01-01T00:03:20 | 22.0 | +| b | 1970-01-01T00:05:00 | 24.0 | +| b | 1970-01-01T00:06:40 | 25.0 | ++------+---------------------+-------------------+ + +-- 7. `@ -1`, an anchor before the Unix epoch, is accepted. Its window `(-301s, -1s]` holds no +-- sample, so the result is empty. +TQL EVAL (0, 0, '1s') at_modifier_gauge @ -1; + ++----+-----+------+ +| ts | val | host | ++----+-----+------+ ++----+-----+------+ + +-- 8. `@` on a metric that does not exist yields an empty result. +TQL EVAL (100, 100, '1s') at_modifier_missing @ 300; + ++------+-------+ +| time | value | ++------+-------+ ++------+-------+ + +-- The same for a range selector over a missing metric: the anchored window is folded once by the +-- promoted call, but the scan has no row to fold, so the whole grid stays empty. +TQL EVAL (0, 240, '60s') rate(at_modifier_missing[5m] @ 300); + ++------+------------------------------------------------+ +| time | prom_rate(time_range,value,time,Int64(300000)) | ++------+------------------------------------------------+ ++------+------------------------------------------------+ + +-- 9. Multi-series coverage of the promoted call. The call above the anchored range selector is +-- evaluated once per series and its result is replayed per series, so grouping by `host` must +-- report both groups at every step: 1.0 for host 'a' and 0.5 for host 'b'. A promoted result that +-- treated the two series as one timeline would report the first step from the first row of the +-- batch and the later steps from its last row, dropping or mixing a host here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') sum by (host) (rate(at_modifier_counter_total[5m] @ 300)); + ++------+---------------------+-----------------------------------------------+ +| host | ts | sum(prom_rate(ts_range,val,ts,Int64(300000))) | ++------+---------------------+-----------------------------------------------+ +| a | 1970-01-01T00:00:00 | 1.0 | +| a | 1970-01-01T00:01:00 | 1.0 | +| a | 1970-01-01T00:02:00 | 1.0 | +| a | 1970-01-01T00:03:00 | 1.0 | +| a | 1970-01-01T00:04:00 | 1.0 | +| b | 1970-01-01T00:00:00 | 0.5 | +| b | 1970-01-01T00:01:00 | 0.5 | +| b | 1970-01-01T00:02:00 | 0.5 | +| b | 1970-01-01T00:03:00 | 0.5 | +| b | 1970-01-01T00:04:00 | 0.5 | ++------+---------------------+-----------------------------------------------+ + +-- Consistency: comparing the per-host sequence above against 1.0 keeps exactly host 'b' at +-- every step, confirming it stays at 0.5 for the whole grid. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') sum by (host) (rate(at_modifier_counter_total[5m] @ 300)) != 1.0; + ++------+---------------------+-----------------------------------------------+ +| host | ts | sum(prom_rate(ts_range,val,ts,Int64(300000))) | ++------+---------------------+-----------------------------------------------+ +| b | 1970-01-01T00:00:00 | 0.5 | +| b | 1970-01-01T00:01:00 | 0.5 | +| b | 1970-01-01T00:02:00 | 0.5 | +| b | 1970-01-01T00:03:00 | 0.5 | +| b | 1970-01-01T00:04:00 | 0.5 | ++------+---------------------+-----------------------------------------------+ + +-- A value function directly over an anchored instant selector is not promoted: the selector +-- anchors and replays its sample per series, and `abs` is evaluated at every step over that replay. +-- It must keep one row per series per step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') abs(at_modifier_gauge @ 300); + ++---------------------+----------+------+ +| ts | abs(val) | host | ++---------------------+----------+------+ +| 1970-01-01T00:01:40 | 12.0 | b | +| 1970-01-01T00:01:40 | 5.0 | a | +| 1970-01-01T00:03:20 | 12.0 | b | +| 1970-01-01T00:03:20 | 5.0 | a | +| 1970-01-01T00:05:00 | 12.0 | b | +| 1970-01-01T00:05:00 | 5.0 | a | +| 1970-01-01T00:06:40 | 12.0 | b | +| 1970-01-01T00:06:40 | 5.0 | a | ++---------------------+----------+------+ + +-- 10. `predict_linear` over an anchored range selector: the regression is centered on the +-- evaluation instant of the step, not on the last sample of the window. The window `(0s, 300s]` of +-- host 'a' is left-open (the sample at 0s is outside it) and rises by 1.0 per minute, so each step +-- predicts 60s ahead of the trend read at its own instant: 6.0 at 300s and 7.0 at 360s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 360, '60s') predict_linear(at_modifier_gauge{host="a"}[5m] @ 300, 60); + ++---------------------+-----------------------------------------------+------+ +| ts | prom_predict_linear(ts_range,val,Float64(60)) | host | ++---------------------+-----------------------------------------------+------+ +| 1970-01-01T00:05:00 | 6.0 | a | +| 1970-01-01T00:06:00 | 7.0 | a | ++---------------------+-----------------------------------------------+------+ + +-- The anchor and the start of the evaluation differ here, so the window `(-60s, 240s]` — samples 0s +-- to 240s — is shifted by `at_offset` = 300s - 240s = 60s. Every step must follow the evaluation +-- grid (6.0 to 11.0); a regression read at the wrong instant predicts 60s ahead of the window's +-- last sample instead, i.e. 5.0 at every step — the sample at 240s is 4.0 and the trend rises by +-- 1.0 per minute — and shows up as a constant here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 600, '60s') predict_linear(at_modifier_gauge{host="a"}[5m] @ 240, 60); + ++---------------------+-----------------------------------------------+------+ +| ts | prom_predict_linear(ts_range,val,Float64(60)) | host | ++---------------------+-----------------------------------------------+------+ +| 1970-01-01T00:05:00 | 6.0 | a | +| 1970-01-01T00:06:00 | 7.0 | a | +| 1970-01-01T00:07:00 | 8.0 | a | +| 1970-01-01T00:08:00 | 9.0 | a | +| 1970-01-01T00:09:00 | 10.0 | a | +| 1970-01-01T00:10:00 | 11.0 | a | ++---------------------+-----------------------------------------------+------+ + +-- Control: without `@` every step folds its own window, whose last sample is up to a step older +-- than the evaluation instant. Reading the regression off that sample reports 6.0 / 7.0 instead +-- of 6.5 / 7.5 here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (330, 390, '60s') predict_linear(at_modifier_gauge{host="a"}[5m], 60); + ++---------------------+-----------------------------------------------+------+ +| ts | prom_predict_linear(ts_range,val,Float64(60)) | host | ++---------------------+-----------------------------------------------+------+ +| 1970-01-01T00:05:30 | 6.5 | a | +| 1970-01-01T00:06:30 | 7.5 | a | ++---------------------+-----------------------------------------------+------+ + +-- The same holds with an `offset`, which shifts the whole window backwards without moving the +-- evaluation instants: the window of step `T` is the left-open `(T - 1m - 5m, T - 1m]`, and the +-- trend is read at `T`, i.e. one minute past its newest sample. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (330, 390, '60s') predict_linear(at_modifier_gauge{host="a"}[5m] offset 1m, 60); + ++---------------------+-----------------------------------------------+------+ +| ts | prom_predict_linear(ts_range,val,Float64(60)) | host | ++---------------------+-----------------------------------------------+------+ +| 1970-01-01T00:05:30 | 6.5 | a | +| 1970-01-01T00:06:30 | 7.5 | a | ++---------------------+-----------------------------------------------+------+ + +-- `timestamp()` of an anchored instant selector is not promoted either: its argument is an +-- instant selector, not a range selector, so the selector anchors and replays its sample per series +-- on its own and `timestamp()` reports the timestamp of the sample the anchor selected at every +-- step (300s here). Its plan keeps the anchored selection below the replay of the grid and projects +-- the timestamp above it. +-- +-- Only an `@` without an `offset` is asserted here: the value of `timestamp()` for an anchored +-- selector that also carries an `offset` is a pre-existing question of its own, out of scope here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') timestamp(at_modifier_gauge @ 300); + ++---------------------+-------+------+ +| ts | value | host | ++---------------------+-------+------+ +| 1970-01-01T00:00:00 | 300.0 | a | +| 1970-01-01T00:00:00 | 300.0 | b | +| 1970-01-01T00:01:00 | 300.0 | a | +| 1970-01-01T00:01:00 | 300.0 | b | +| 1970-01-01T00:02:00 | 300.0 | a | +| 1970-01-01T00:02:00 | 300.0 | b | +| 1970-01-01T00:03:00 | 300.0 | a | +| 1970-01-01T00:03:00 | 300.0 | b | +| 1970-01-01T00:04:00 | 300.0 | a | +| 1970-01-01T00:04:00 | 300.0 | b | ++---------------------+-------+------+ + +-- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED +-- SQLNESS REPLACE (peers.*) REDACTED +-- SQLNESS REPLACE (Hash.*) REDACTED +-- SQLNESS REPLACE (RepartitionExec:.*) RepartitionExec: REDACTED +TQL EXPLAIN (0, 240, '60s') timestamp(at_modifier_gauge @ 300); + ++---------------+--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ +| plan_type | plan | ++---------------+--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ +| logical_plan | MergeScan [is_placeholder=false, remote_input=[ | +| | Projection: at_modifier_gauge.ts, value AS value, at_modifier_gauge.host | +| | Projection: at_modifier_gauge.ts, __promql_timestamp_value_ AS value, at_modifier_gauge.host | +| | PromInstantManipulate: range=[0..240000], lookback=[240001], interval=[60000], time index=[ts] | +| | PromInstantManipulate: range=[0..0], lookback=[300000], interval=[60000], time index=[ts] | +| | Projection: at_modifier_gauge.ts, at_modifier_gauge.val, at_modifier_gauge.host, CAST(CAST(CAST(CAST(CAST(at_modifier_gauge.ts AS Int64) AS Decimal128(19, 0)) AS Decimal128(21, 0)) + Decimal128(0,19,0) AS Int64) AS Float64) / Float64(1000) AS __promql_timestamp_value_ | +| | PromSeriesNormalize: offset=[-300000], time index=[ts], filter NaN: [false] | +| | PromSeriesDivide: tags=["host"] | +| | Sort: at_modifier_gauge.host ASC NULLS FIRST, at_modifier_gauge.ts ASC NULLS FIRST | +| | Filter: at_modifier_gauge.ts >= TimestampMillisecond(1, None) AND at_modifier_gauge.ts <= TimestampMillisecond(300000, None) | +| | TableScan: at_modifier_gauge, partial_filters=[at_modifier_gauge.ts >= TimestampMillisecond(1, None), at_modifier_gauge.ts <= TimestampMillisecond(300000, None)] | +| | ]] | +| physical_plan | CooperativeExec | +| | MergeScanExec: REDACTED +| | | ++---------------+--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ + +-- 11. `label_join` above an anchored selector. The call rewrites the label the input series are +-- told apart by (`host` becomes the empty string), so it is never promoted: it must stay above the +-- per-series replay of the selector and be evaluated at every step. Both hosts are still reported, +-- with their own value, at every step. A promoted `label_join` would instead be replayed through +-- the labels it just rewrote and report the two series as a single timeline, i.e. one mixed row +-- per step. +-- +-- The invariant asserted here is scoped to the replay: it must not drop or mix the input rows. +-- It is not a claim about the final PromQL semantics of this query: joining `host` to one value +-- leaves two samples with the same label set at the same timestamp, which Prometheus rejects, +-- while `label_join` does not validate that yet (`label_replace` errors on such a rewrite +-- instead). That duplicate-labelset validation gap is pre-existing and out of scope here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 480, '60s') label_join(at_modifier_gauge @ 300, "host", "", ""); + ++---------------------+------+------+ +| ts | val | host | ++---------------------+------+------+ +| 1970-01-01T00:05:00 | 12.0 | | +| 1970-01-01T00:05:00 | 5.0 | | +| 1970-01-01T00:06:00 | 12.0 | | +| 1970-01-01T00:06:00 | 5.0 | | +| 1970-01-01T00:07:00 | 12.0 | | +| 1970-01-01T00:07:00 | 5.0 | | +| 1970-01-01T00:08:00 | 12.0 | | +| 1970-01-01T00:08:00 | 5.0 | | ++---------------------+------+------+ + +-- The same join below another call: neither is promoted, and both are evaluated at every step over +-- the per-series replay of the anchored selector. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 480, '60s') abs(label_join(at_modifier_gauge @ 300, "host", "", "")); + ++---------------------+----------+------+ +| ts | abs(val) | host | ++---------------------+----------+------+ +| 1970-01-01T00:05:00 | 12.0 | | +| 1970-01-01T00:05:00 | 5.0 | | +| 1970-01-01T00:06:00 | 12.0 | | +| 1970-01-01T00:06:00 | 5.0 | | +| 1970-01-01T00:07:00 | 12.0 | | +| 1970-01-01T00:07:00 | 5.0 | | +| 1970-01-01T00:08:00 | 12.0 | | +| 1970-01-01T00:08:00 | 5.0 | | ++---------------------+----------+------+ + +-- A range call below the join is still promoted on its own (it is the direct call over the +-- anchored range selector): the anchored window is folded once per series, and the join above that +-- replay reports the rate of both hosts at every step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 480, '60s') label_join(rate(at_modifier_counter_total[5m] @ 300), "host", "", ""); + ++---------------------+------------------------------------------+------+ +| ts | prom_rate(ts_range,val,ts,Int64(300000)) | host | ++---------------------+------------------------------------------+------+ +| 1970-01-01T00:05:00 | 0.5 | | +| 1970-01-01T00:05:00 | 1.0 | | +| 1970-01-01T00:06:00 | 0.5 | | +| 1970-01-01T00:06:00 | 1.0 | | +| 1970-01-01T00:07:00 | 0.5 | | +| 1970-01-01T00:07:00 | 1.0 | | +| 1970-01-01T00:08:00 | 0.5 | | +| 1970-01-01T00:08:00 | 1.0 | | ++---------------------+------------------------------------------+------+ + +DROP TABLE at_modifier_gauge; + +Affected Rows: 0 + +DROP TABLE at_modifier_counter_total; + +Affected Rows: 0 + diff --git a/tests/cases/standalone/common/promql/at_modifier.sql b/tests/cases/standalone/common/promql/at_modifier.sql new file mode 100644 index 00000000000..9e9bc21a9c1 --- /dev/null +++ b/tests/cases/standalone/common/promql/at_modifier.sql @@ -0,0 +1,301 @@ +-- Tests for the PromQL `@` modifier on vector and matrix selectors. +-- +-- `@` anchors the sample selection window at a fixed timestamp instead of the timestamp of each +-- evaluation step: `@ ` uses the given timestamp, `@ start()` / `@ end()` use the +-- start/end of the statement's evaluation range, and `offset` shifts the anchor backwards before +-- the window is selected. The output timestamps still follow the evaluation grid. +-- +-- Every metric below carries two series, host 'a' and host 'b', with different values: the +-- anchored results must therefore report both of them, per step. A series silently dropped, or a +-- batch holding several series handled as one timeline, does not look like a correct single-series +-- answer here but shows up as a missing or shifted row. +-- +-- Sample timestamps below are in milliseconds, while `@` and `TQL EVAL` timestamps are in seconds, +-- as in Prometheus. + +CREATE TABLE at_modifier_gauge ( + ts TIMESTAMP TIME INDEX, + val DOUBLE, + host STRING, + PRIMARY KEY(host) +); + +CREATE TABLE at_modifier_counter_total ( + ts TIMESTAMP TIME INDEX, + val DOUBLE, + host STRING, + PRIMARY KEY(host) +); + +-- One sample per minute, so the anchored sample is easy to identify. The default lookback is 5m. +-- Host 'a' reports 0.0 to 6.0, host 'b' reports 7.0 to 13.0: the two series are told apart by their +-- values alone, which stays true after any anchoring. +INSERT INTO at_modifier_gauge VALUES + (0, 0.0, 'a'), + (60000, 1.0, 'a'), + (120000, 2.0, 'a'), + (180000, 3.0, 'a'), + (240000, 4.0, 'a'), + (300000, 5.0, 'a'), + (360000, 6.0, 'a'), + (0, 7.0, 'b'), + (60000, 8.0, 'b'), + (120000, 9.0, 'b'), + (180000, 10.0, 'b'), + (240000, 11.0, 'b'), + (300000, 12.0, 'b'), + (360000, 13.0, 'b'); + +-- Host 'a': a counter increasing by 60 per minute, i.e. by 1 per second. +-- Host 'b': a counter starting at 100 and increasing by 30 per minute, i.e. by 0.5 per second. Its +-- rate differs from host 'a', so a window folded for the wrong series cannot pass as the right one. +INSERT INTO at_modifier_counter_total VALUES + (0, 0.0, 'a'), + (60000, 60.0, 'a'), + (120000, 120.0, 'a'), + (180000, 180.0, 'a'), + (240000, 240.0, 'a'), + (300000, 300.0, 'a'), + (360000, 360.0, 'a'), + (0, 100.0, 'b'), + (60000, 130.0, 'b'), + (120000, 160.0, 'b'), + (180000, 190.0, 'b'), + (240000, 220.0, 'b'), + (300000, 250.0, 'b'), + (360000, 280.0, 'b'); + +-- 1. Instant query anchored at an absolute timestamp: the sample window is `(anchor - lookback, +-- anchor]` = `(0s, 300s]`, so its newest samples are 5.0 for host 'a' and 12.0 for host 'b', both +-- at 300s. The output timestamps stay the evaluation timestamps. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 300; + +-- Control: the same query without `@` looks back from the evaluation timestamp instead: +-- `(-200s, 100s]` selects the samples 1.0 (a) and 8.0 (b) at 60s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge; + +-- The anchor, not the evaluation timestamp, decides which sample is reported. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 360; + +-- Sub-second precision is preserved in the anchor: `240.999` selects the samples at 240s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 240.999; + +-- The window `(-300s, 0s]` includes the samples exactly at the anchor. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 0; + +-- 2. `@ start()`: every step selects its samples around the start of the evaluation range +-- (`(-200s, 100s]`), so every step reports the same value per series (1.0 for 'a', 8.0 for 'b'). +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge @ start(); + +-- Control: without `@` every step looks back from its own timestamp, so the steps differ. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge; + +-- 3. `@ end()`: every step selects its samples around the end of the evaluation range +-- (`(100s, 400s]`), whose newest samples are 6.0 (a) and 13.0 (b) at 360s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge @ end(); + +-- 4. Range selector anchored by `@`: the call directly above the anchored range selector is +-- evaluated once, at the start of the evaluation, and its result is reported at every step — the +-- planner's counterpart of Prometheus' `StepInvariantExpr`. The window `(0s, 300s]` is left-open, +-- so it holds 60 -> 300 for host 'a' and 130 -> 250 for host 'b', i.e. 240 and 120 units. The 60s +-- gap to the window start is below the extrapolation threshold (1.1 x the 60s sampling interval), +-- so `rate` adds it and divides by the 300s window: 240 * 1.25 / 300 = 1.0 and +-- 120 * 1.25 / 300 = 0.5 per second. Every step must report those same per-series values, and both +-- series must be there at every step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate(at_modifier_counter_total[5m] @ 300); + +-- Consistency: the instant query below shows the anchored values over the same left-open +-- `(0s, 300s]` window (1.0 per second for 'a', 0.5 for 'b'). The `!= 1.0` comparison filters +-- host 'a', but host 'b' remains at every step with its rate of 0.5 per second. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 300, '1s') rate(at_modifier_counter_total[5m]); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate(at_modifier_counter_total[5m] @ 300) != 1.0; + +-- Parentheses around the range argument are transparent: the call must be promoted exactly like +-- the unparenthesized one above, so every step reports the same anchored values for both series +-- (1.0 for host 'a', 0.5 for host 'b') instead of folding a window per step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate((at_modifier_counter_total[5m] @ 300)); + +-- The same holds for `@ start()` and `@ end()`, which are also fixed anchors: they resolve to the +-- statement's evaluation range, so a range call using them is evaluated once as well. +-- The statement starts at 60s: the window `(-240s, 60s]` holds only the samples at 0s and 60s, so +-- host 'a' goes 0 -> 60 and host 'b' 100 -> 130. The far 240s gap to the window start is beyond the +-- extrapolation threshold (1.1 x the 60s interval), so only half an interval (30s) is added there; +-- host 'a' starts its counter at 0, which clamps that back: 60 / 300 = 0.2 and +-- 30 * 1.5 / 300 = 0.15 per second. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (60, 300, '60s') rate(at_modifier_counter_total[5m] @ start()); + +-- The instant query at 60s folds the same window as every step of the query above. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (60, 60, '1s') rate(at_modifier_counter_total[5m]); + +-- The statement ends at 240s: the window `(-60s, 240s]` holds the samples 0s to 240s, so host 'a' +-- goes 0 -> 240 and host 'b' 100 -> 220. The 60s gap to the window start is added for both, but +-- host 'a' starts at zero, which clamps it back (240 / 300 = 0.8 per second) while host 'b' keeps +-- it (120 * 1.25 / 300 = 0.5 per second). +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate(at_modifier_counter_total[5m] @ end()); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (240, 240, '1s') rate(at_modifier_counter_total[5m]); + +-- Control: without `@` each step folds its own window, so the values keep changing. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') rate(at_modifier_counter_total[5m]); + +-- `@ start()` on a range selector: the window `(-300s, 0s]` is left-open, so it holds a single +-- sample per series — the one exactly at the anchor — and `count_over_time` returns 1 at every step +-- for both of them. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 300, '100s') count_over_time(at_modifier_counter_total[5m] @ start()); + +-- Control: without `@` the window keeps moving, so the count grows with the step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 300, '100s') count_over_time(at_modifier_counter_total[5m]); + +-- 5. `@` combined with `offset`: the offset moves the anchor backwards before the window is +-- selected, so `@ 300 offset 2m` reports the samples at 180s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 300 offset 2m; + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 100, '1s') at_modifier_gauge @ 300 offset 1m; + +-- Control: the offset applies to the evaluation timestamp when there is no `@`. +TQL EVAL (100, 100, '1s') at_modifier_gauge offset 2m; + +-- 6. `@` inside a larger expression. The anchored operand contributes the same value at every +-- step, while the plain operand still follows the evaluation step. A binary expression is never +-- promoted: when one side is anchored and the other is not, the anchored side replays on its own, +-- and when both sides are anchored, each side reports its own anchored samples. The join then runs +-- at every step over the replayed per-series inputs. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge @ 300 + at_modifier_gauge @ 0; + +-- The aggregation above the anchored selector runs at every step over the replayed (still +-- per-series) samples, so it sums both series at every step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') sum(at_modifier_gauge @ 0); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') abs(at_modifier_gauge @ 0); + +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') at_modifier_gauge @ 300 + at_modifier_gauge; + +-- 7. `@ -1`, an anchor before the Unix epoch, is accepted. Its window `(-301s, -1s]` holds no +-- sample, so the result is empty. +TQL EVAL (0, 0, '1s') at_modifier_gauge @ -1; + +-- 8. `@` on a metric that does not exist yields an empty result. +TQL EVAL (100, 100, '1s') at_modifier_missing @ 300; + +-- The same for a range selector over a missing metric: the anchored window is folded once by the +-- promoted call, but the scan has no row to fold, so the whole grid stays empty. +TQL EVAL (0, 240, '60s') rate(at_modifier_missing[5m] @ 300); + +-- 9. Multi-series coverage of the promoted call. The call above the anchored range selector is +-- evaluated once per series and its result is replayed per series, so grouping by `host` must +-- report both groups at every step: 1.0 for host 'a' and 0.5 for host 'b'. A promoted result that +-- treated the two series as one timeline would report the first step from the first row of the +-- batch and the later steps from its last row, dropping or mixing a host here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') sum by (host) (rate(at_modifier_counter_total[5m] @ 300)); + +-- Consistency: comparing the per-host sequence above against 1.0 keeps exactly host 'b' at +-- every step, confirming it stays at 0.5 for the whole grid. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') sum by (host) (rate(at_modifier_counter_total[5m] @ 300)) != 1.0; + +-- A value function directly over an anchored instant selector is not promoted: the selector +-- anchors and replays its sample per series, and `abs` is evaluated at every step over that replay. +-- It must keep one row per series per step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (100, 400, '100s') abs(at_modifier_gauge @ 300); + +-- 10. `predict_linear` over an anchored range selector: the regression is centered on the +-- evaluation instant of the step, not on the last sample of the window. The window `(0s, 300s]` of +-- host 'a' is left-open (the sample at 0s is outside it) and rises by 1.0 per minute, so each step +-- predicts 60s ahead of the trend read at its own instant: 6.0 at 300s and 7.0 at 360s. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 360, '60s') predict_linear(at_modifier_gauge{host="a"}[5m] @ 300, 60); + +-- The anchor and the start of the evaluation differ here, so the window `(-60s, 240s]` — samples 0s +-- to 240s — is shifted by `at_offset` = 300s - 240s = 60s. Every step must follow the evaluation +-- grid (6.0 to 11.0); a regression read at the wrong instant predicts 60s ahead of the window's +-- last sample instead, i.e. 5.0 at every step — the sample at 240s is 4.0 and the trend rises by +-- 1.0 per minute — and shows up as a constant here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 600, '60s') predict_linear(at_modifier_gauge{host="a"}[5m] @ 240, 60); + +-- Control: without `@` every step folds its own window, whose last sample is up to a step older +-- than the evaluation instant. Reading the regression off that sample reports 6.0 / 7.0 instead +-- of 6.5 / 7.5 here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (330, 390, '60s') predict_linear(at_modifier_gauge{host="a"}[5m], 60); + +-- The same holds with an `offset`, which shifts the whole window backwards without moving the +-- evaluation instants: the window of step `T` is the left-open `(T - 1m - 5m, T - 1m]`, and the +-- trend is read at `T`, i.e. one minute past its newest sample. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (330, 390, '60s') predict_linear(at_modifier_gauge{host="a"}[5m] offset 1m, 60); + +-- `timestamp()` of an anchored instant selector is not promoted either: its argument is an +-- instant selector, not a range selector, so the selector anchors and replays its sample per series +-- on its own and `timestamp()` reports the timestamp of the sample the anchor selected at every +-- step (300s here). Its plan keeps the anchored selection below the replay of the grid and projects +-- the timestamp above it. +-- +-- Only an `@` without an `offset` is asserted here: the value of `timestamp()` for an anchored +-- selector that also carries an `offset` is a pre-existing question of its own, out of scope here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (0, 240, '60s') timestamp(at_modifier_gauge @ 300); + +-- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED +-- SQLNESS REPLACE (peers.*) REDACTED +-- SQLNESS REPLACE (Hash.*) REDACTED +-- SQLNESS REPLACE (RepartitionExec:.*) RepartitionExec: REDACTED +TQL EXPLAIN (0, 240, '60s') timestamp(at_modifier_gauge @ 300); + +-- 11. `label_join` above an anchored selector. The call rewrites the label the input series are +-- told apart by (`host` becomes the empty string), so it is never promoted: it must stay above the +-- per-series replay of the selector and be evaluated at every step. Both hosts are still reported, +-- with their own value, at every step. A promoted `label_join` would instead be replayed through +-- the labels it just rewrote and report the two series as a single timeline, i.e. one mixed row +-- per step. +-- +-- The invariant asserted here is scoped to the replay: it must not drop or mix the input rows. +-- It is not a claim about the final PromQL semantics of this query: joining `host` to one value +-- leaves two samples with the same label set at the same timestamp, which Prometheus rejects, +-- while `label_join` does not validate that yet (`label_replace` errors on such a rewrite +-- instead). That duplicate-labelset validation gap is pre-existing and out of scope here. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 480, '60s') label_join(at_modifier_gauge @ 300, "host", "", ""); + +-- The same join below another call: neither is promoted, and both are evaluated at every step over +-- the per-series replay of the anchored selector. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 480, '60s') abs(label_join(at_modifier_gauge @ 300, "host", "", "")); + +-- A range call below the join is still promoted on its own (it is the direct call over the +-- anchored range selector): the anchored window is folded once per series, and the join above that +-- replay reports the rate of both hosts at every step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (300, 480, '60s') label_join(rate(at_modifier_counter_total[5m] @ 300), "host", "", ""); + +DROP TABLE at_modifier_gauge; + +DROP TABLE at_modifier_counter_total;