Skip to main content

query/promql/planner/
at_modifier.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Anchored selector planning and replay for the PromQL `@` modifier.
16
17use std::sync::Arc;
18use std::time::{SystemTime, UNIX_EPOCH};
19
20use datafusion::logical_expr::{Extension, LogicalPlan, LogicalPlanBuilder};
21use datafusion::prelude::{Column, Expr as DfExpr};
22use promql::extension_plan::{InstantManipulate, Millisecond, SeriesDivide};
23use promql_parser::parser::{
24    AtModifier, Call, Expr as PromExpr, MatrixSelector, Offset, ParenExpr,
25};
26use snafu::{OptionExt, ResultExt, ensure};
27
28use crate::promql::error::{
29    AtModifierTimestampOutOfRangeSnafu, DataFusionPlanningSnafu, Result, TimeIndexNotFoundSnafu,
30};
31use crate::promql::planner::PromPlanner;
32use crate::query_engine::QueryEngineState;
33
34impl PromPlanner {
35    /// Resolve the `@` modifier of a vector or matrix selector into the timestamp its sample
36    /// window is anchored at, in milliseconds since the Unix epoch.
37    ///
38    /// Prometheus semantics:
39    /// - `@ <unix_ts>` anchors at the given timestamp,
40    /// - `@ start()` / `@ end()` anchor at the evaluation range of the whole statement,
41    /// - `offset` shifts the anchor backwards: the window ends at `anchor - offset`.
42    ///
43    /// Returns `None` when the selector has no `@` modifier.
44    fn at_ref_time(
45        &self,
46        at: &Option<AtModifier>,
47        offset: &Option<Offset>,
48    ) -> Result<Option<Millisecond>> {
49        let anchor = match at {
50            None => return Ok(None),
51            Some(AtModifier::Start) => self.ctx.stmt_start,
52            Some(AtModifier::End) => self.ctx.stmt_end,
53            Some(AtModifier::At(time)) => Self::system_time_to_millis(time)?,
54        };
55        Ok(Some(Self::anchor_sub(anchor, Self::offset_millis(offset))?))
56    }
57
58    /// Subtracts `rhs` from `lhs` on the millisecond timeline of an `@` anchor.
59    ///
60    /// A negative result is valid: `@` and `offset` accept timestamps before the Unix epoch. A
61    /// result outside the representable millisecond range is rejected like an unrepresentable
62    /// anchor ([`Self::system_time_to_millis`]) instead of clamping it, so the same class of
63    /// input always gets the same answer.
64    pub(crate) fn anchor_sub(lhs: Millisecond, rhs: Millisecond) -> Result<Millisecond> {
65        lhs.checked_sub(rhs)
66            .with_context(|| AtModifierTimestampOutOfRangeSnafu {
67                timestamp: format!("{}ms - {}ms", lhs, rhs),
68            })
69    }
70
71    /// The offset a selector with an `@` modifier is evaluated with.
72    ///
73    /// Prometheus anchors such a selector by rewriting its offset to `eval_time - anchor`
74    /// (`setOffsetForAtModifier`), so that the selector always selects its samples around `anchor`
75    /// regardless of the step being evaluated. `eval_time` is the start of the evaluation the
76    /// selector belongs to, which is `ctx.start`.
77    ///
78    /// Returns `None` when the selector has no `@` modifier.
79    pub(crate) fn at_modifier_offset(
80        &self,
81        at: &Option<AtModifier>,
82        offset: &Option<Offset>,
83    ) -> Result<Option<Millisecond>> {
84        let Some(anchor) = self.at_ref_time(at, offset)? else {
85            return Ok(None);
86        };
87        Ok(Some(Self::anchor_sub(self.ctx.start, anchor)?))
88    }
89
90    /// Whether `expr` is a call that has to be evaluated once for the whole grid, because it folds
91    /// a range selector anchored by `@` — the range argument of the call's parser signature, which
92    /// is a [`MatrixSelector`] here; only a call with such an argument can take this path, so no
93    /// function-name registry is involved.
94    ///
95    /// This is the shape that needs Prometheus' `StepInvariantExpr` the most: a range function such
96    /// as `rate` derives its result from the evaluation instant it is called at, so folding the
97    /// window once per step would let the outer evaluation grid change the result of a window that
98    /// `@` fixed. Evaluated once, at the start of the grid, the rewritten offset of
99    /// [`Self::at_modifier_offset`] places the anchor at that instant, and [`Self::replay_over_grid`]
100    /// reports the result at every step.
101    ///
102    /// Unlike Prometheus, which wraps the whole step-invariant subtree (`preprocessExprHelper`),
103    /// only the call itself is promoted here. The operators above it are not: they are still planned
104    /// at every step over the replayed result, which is safe for the row-wise ones and keeps the
105    /// promotion root narrow. The promotion root has to stay a call over one range selector, because
106    /// the replay needs one series per batch ([`Self::series_divide_plan`]) and only such a call
107    /// guarantees that the rows it emits still describe the series it was divided by. The operators
108    /// left out — an aggregation, a join, or a label rewriting call such as `label_join` — mix or
109    /// re-label the rows of different series, so they are unsafe as promotion roots even though
110    /// evaluating them after the replay is fine. A call whose input is an anchored *instant*
111    /// selector (`abs(some_metric @ 300)`) needs no promotion either: the selector anchors and
112    /// replays its sample per series on its own, and the call above it is row-wise.
113    ///
114    /// `predict_linear` is the exception among the range functions: it predicts from the evaluation
115    /// instant of each step ([`Self::create_range_eval_ts_expr`]), so it has to stay outside the
116    /// promoted subtree and follow the grid. The remaining arguments of the call have to be
117    /// literals, since the replay of the promoted result has no second vector input to divide.
118    /// Parentheses around the range argument are transparent (`rate((m[5m] @ 300))`), so they are
119    /// looked through and the call is promoted as if they were absent. Nothing else of the subtree
120    /// is unwrapped, so an outer parenthesis promotes no operator above the call.
121    fn promotes_anchored_range_call(expr: &PromExpr) -> bool {
122        let PromExpr::Call(Call { func, args }) = expr else {
123            return false;
124        };
125        // See the doc comment: the regression of `predict_linear` follows the evaluation step.
126        if func.name == "predict_linear" {
127            return false;
128        }
129        let mut anchored_range = false;
130        for arg in &args.args {
131            // Parentheses around the range argument are transparent, so the call is promoted the
132            // same way for `rate((m[5m] @ 300))` as for `rate(m[5m] @ 300)`. Only the parentheses
133            // directly around this one argument are looked through here: the promotion stays
134            // confined to a call over one anchored range selector instead of descending into an
135            // arbitrary parenthesized subtree.
136            let mut arg = arg.as_ref();
137            while let PromExpr::Paren(ParenExpr { expr }) = arg {
138                arg = expr;
139            }
140            match arg {
141                // The window is pinned by `@`, so every step folds the same samples.
142                PromExpr::MatrixSelector(MatrixSelector { vs, .. }) if vs.at.is_some() => {
143                    if anchored_range {
144                        return false;
145                    }
146                    anchored_range = true;
147                }
148                // A literal argument is the same value at every step.
149                arg if Self::try_build_literal_expr(arg).is_some() => {}
150                _ => return false,
151            }
152        }
153        anchored_range
154    }
155
156    /// Plans the anchored range call `prom_expr` ([`Self::promotes_anchored_range_call`]) as a
157    /// step-invariant subtree: the call is evaluated on a single evaluation instant (`grid_start`,
158    /// the start of the outer evaluation) and its result is then reported at every step of
159    /// `[grid_start, ctx.end]` by [`Self::replay_over_grid`].
160    ///
161    /// This is the planner's counterpart of Prometheus' `StepInvariantExpr` for the one shape it
162    /// promotes. Evaluating the call once matters for the functions that derive their result from
163    /// the step being evaluated: `rate(m[5m] @ 300)` folds its window around the anchor once, and
164    /// the extrapolation boundaries of `rate` must be derived from that same window at every step
165    /// instead of following the outer evaluation timestamp.
166    ///
167    /// Only the call itself is promoted; the operators above it are planned as usual over the
168    /// replayed result. The result of the promoted call is split into one series per batch before it
169    /// is replayed ([`Self::series_divide_plan`]), because the row-wise projection of the call does
170    /// not preserve the batch layout of the selector.
171    ///
172    /// Returns `None` when `prom_expr` is not such a call, so that the caller plans it as usual.
173    /// The selector inside the promoted call keeps its own `@` anchoring (see
174    /// [`Self::at_modifier_offset`]), and planning it with `ctx.end == ctx.start` folds its window
175    /// once for that single instant instead of expanding it over the grid, which the replay of the
176    /// call result above already does.
177    pub(crate) async fn promote_anchored_range_call(
178        &mut self,
179        prom_expr: &PromExpr,
180        timestamp_fn: bool,
181        query_engine_state: &QueryEngineState,
182    ) -> Result<Option<LogicalPlan>> {
183        let grid_start = self.ctx.start;
184        let grid_end = self.ctx.end;
185        // An instant query evaluates a single step, so there is nothing to promote.
186        if grid_start == grid_end || !Self::promotes_anchored_range_call(prom_expr) {
187            return Ok(None);
188        }
189
190        // Plan the subtree on a single evaluation instant: every selector below still anchors its
191        // window through `@`, and the functions above them derive their result from that one
192        // instant. The planner is single-use, so `ctx.end` needs no restore-on-error.
193        self.ctx.end = grid_start;
194        let anchored = self
195            .prom_expr_to_plan_inner(prom_expr, timestamp_fn, query_engine_state)
196            .await?;
197        self.ctx.end = grid_end;
198
199        let time_index_column =
200            self.ctx
201                .time_index_column
202                .clone()
203                .with_context(|| TimeIndexNotFoundSnafu {
204                    table: self.ctx.table_name.clone().unwrap_or_default(),
205                })?;
206        // The replay reads one series per batch; see [`Self::series_divide_plan`] for why the
207        // layout of the selector below the promoted call does not survive it.
208        let anchored = self.series_divide_plan(anchored, &time_index_column)?;
209        Ok(Some(self.replay_over_grid(
210            anchored,
211            grid_start,
212            grid_end,
213            time_index_column,
214        )))
215    }
216
217    /// Convert the timestamp of an `@` modifier into milliseconds since the Unix epoch.
218    fn system_time_to_millis(time: &SystemTime) -> Result<Millisecond> {
219        let (millis, negative) = match time.duration_since(UNIX_EPOCH) {
220            Ok(duration) => (duration.as_millis(), false),
221            // The `@` modifier accepts timestamps before the Unix epoch, e.g. `@ -1`.
222            Err(err) => (err.duration().as_millis(), true),
223        };
224        ensure!(
225            millis <= i64::MAX as u128,
226            AtModifierTimestampOutOfRangeSnafu {
227                timestamp: time
228                    .duration_since(UNIX_EPOCH)
229                    .map(|duration| format!("+{}ms", duration.as_millis()))
230                    .unwrap_or_else(|err| format!("-{}ms", err.duration().as_millis())),
231            }
232        );
233        let millis = millis as Millisecond;
234        Ok(if negative { -millis } else { millis })
235    }
236
237    /// Report the samples of `anchored` at every step of the evaluation grid
238    /// `[grid_start, grid_end]`.
239    ///
240    /// A selector with an `@` modifier is anchored: the sample window is selected once, around the
241    /// anchor timestamp, and every evaluation step reports that same window. Prometheus does this
242    /// by rewriting the selector's offset to `eval_time - anchor` and only fetching the samples on
243    /// the first step (`setOffsetForAtModifier` plus the `refetch` shortcut in `rangeEval`).
244    ///
245    /// The expansion reuses [`InstantManipulate`] with a lookback that spans the whole grid: every
246    /// step then picks the same sample (or the same already computed value, when `anchored` ends
247    /// with a function call such as `rate`) and stamps it with the step's timestamp.
248    ///
249    /// Every input batch of `anchored` must hold exactly one series, because [`InstantManipulate`]
250    /// takes a batch as one timeline. A leaf-level replay (`m @ 300`) consumes the [`SeriesDivide`]
251    /// of its selector directly. A promoted call is guaranteed that layout by
252    /// [`Self::series_divide_plan`], which is why [`Self::promote_anchored_range_call`] splits
253    /// its result before calling this method.
254    pub(crate) fn replay_over_grid(
255        &self,
256        anchored: LogicalPlan,
257        grid_start: Millisecond,
258        grid_end: Millisecond,
259        time_index_column: String,
260    ) -> LogicalPlan {
261        if grid_start == grid_end {
262            // A single evaluation step: `anchored` is already stamped with that timestamp.
263            return anchored;
264        }
265
266        let series_key_columns = self.series_key_columns_for_schema(anchored.schema());
267        let replayed = InstantManipulate::new(
268            grid_start,
269            grid_end,
270            // The lookback must keep the single anchored sample eligible for every step.
271            grid_end - grid_start + 1,
272            self.ctx.interval,
273            0,
274            time_index_column,
275            series_key_columns,
276            self.ctx.field_columns.first().cloned(),
277            anchored,
278        );
279        LogicalPlan::Extension(Extension {
280            node: Arc::new(replayed),
281        })
282    }
283
284    /// Sorts `input` by its series key and time index and splits it into one batch per series.
285    ///
286    /// [`InstantManipulate`] reads every input batch as one series (it takes the timeline of the
287    /// batch and reports the row selected at every step), so a batch holding several series would
288    /// lose all but one of them. A selector establishes that layout with its own [`SeriesDivide`],
289    /// but the per-series distribution requirement does not reach a promoted subtree above it:
290    /// the row-wise projection of a call sits in between, so the batch boundaries of the selector
291    /// are not preserved — in a distributed plan the promoted result can be delivered as one batch
292    /// holding every series. Sorting and dividing here restores the layout, exactly like
293    /// [`Self::prom_matrix_selector_to_plan`] does for the input of a range function.
294    ///
295    /// Series keys that are not present in `input` are dropped, since `ctx.tag_columns` may have
296    /// drifted from the actual output schema. A plan without any series key column is returned
297    /// unchanged: there is nothing to divide by.
298    fn series_divide_plan(
299        &self,
300        input: LogicalPlan,
301        time_index_column: &str,
302    ) -> Result<LogicalPlan> {
303        let series_key_columns = self.series_key_columns_for_schema(input.schema());
304        if series_key_columns.is_empty() {
305            return Ok(input);
306        }
307
308        let mut sort_exprs = series_key_columns
309            .iter()
310            .map(|name| DfExpr::Column(Column::from_name(name)).sort(true, true))
311            .collect::<Vec<_>>();
312        sort_exprs.push(DfExpr::Column(Column::from_name(time_index_column)).sort(true, true));
313        let sort_plan = LogicalPlanBuilder::from(input)
314            .sort(sort_exprs)
315            .context(DataFusionPlanningSnafu)?
316            .build()
317            .context(DataFusionPlanningSnafu)?;
318        Ok(LogicalPlan::Extension(Extension {
319            node: Arc::new(SeriesDivide::new(
320                series_key_columns,
321                time_index_column.to_string(),
322                sort_plan,
323            )),
324        }))
325    }
326}