Skip to content

Commit 2628423

Browse files
zzylolzz_yclaude
authored
feat(lower): support *_over_time over a sub-query (#27) (#42)
`max_over_time(rate(m[5m])[1h:])` and the rest of the `*_over_time`/`quantile_over_time`-of-sub-query family were rejected — `extract_matrix` only accepts a (parenthesised) matrix selector, not the `PromQLSubquery` a sub-query argument produces. `walk_call` now intercepts a `*_over_time` call whose argument is a sub-query, lowers the sub-query recursively, and reduces it per series: max_over_time(rate(m[5m])[1h:]) -> Aggregate{Max} -> Subquery{1h} -> Aggregate{Rate} -> TimeRange{5m} -> Scan Non-sub-query calls fall through to the existing flat template unchanged, so plain `max_over_time(m[w])` and heavy-hitter `topk(k, count_over_time(m[w]))` are untouched. Per-series correctness: an empty-keys aggregate over a `Subquery` is now recognised as a label-preserving per-series range reduction, mirroring the existing `TimeRange` marker. This is sound because a genuine cross-series aggregation operator over a range vector (`sum(rate(m[5m])[1h:])`) is a PromQL type error the parser already rejects — so an `Aggregate` directly over a `Subquery` only ever arises from this per-series `*_over_time` shape. Verified by `sum by (job) (max_over_time(rate(m[5m])[1h:]))`: the inner `Max` preserves `job` for the outer `sum by (job)` to group on. Flips the `over_time_of_subquery_is_rejected__GAP` conformance test into three passing tests (per-series reduction, `quantile_over_time` phi, outer aggregation label preservation). Full workspace suite green; clippy clean. Co-authored-by: zz_y <zz_y@node0.zz-y-308294.softmeasure-pg0.clemson.cloudlab.us> Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent b324676 commit 2628423

3 files changed

Lines changed: 124 additions & 11 deletions

File tree

‎crates/core/src/intent_algebra/query_expr.rs‎

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -428,11 +428,18 @@ impl QueryExpr {
428428

429429
// Per-series range reduction: `rate`/`increase` (is_per_series)
430430
// OR any single aggregate whose direct child is a `TimeRange`
431-
// (`*_over_time` functions). Both produce one value per series
432-
// and are label-preserving — the `TimeRange` child is the
433-
// structural marker that confers per-series semantics on
434-
// otherwise cross-series intents like `Avg`/`Sum`/`Count`.
435-
let is_range_child = matches!(child.as_ref(), QueryExpr::TimeRange { .. });
431+
// (`*_over_time` functions) or a `Subquery` (`*_over_time` over a
432+
// sub-query, e.g. `max_over_time(rate(m[5m])[1h:])`). All produce
433+
// one value per series and are label-preserving — the range child
434+
// is the structural marker that confers per-series semantics on
435+
// otherwise cross-series intents like `Avg`/`Sum`/`Count`. A
436+
// cross-series aggregation operator over a range vector is a
437+
// PromQL type error the parser rejects, so an `Aggregate` over a
438+
// `Subquery` is only ever this per-series `*_over_time` shape.
439+
let is_range_child = matches!(
440+
child.as_ref(),
441+
QueryExpr::TimeRange { .. } | QueryExpr::Subquery { .. }
442+
);
436443
if by.is_empty() && aggs.len() == 1 && (aggs[0].is_per_series() || is_range_child) {
437444
return Ok(per_series_reduction_schema(&in_schema, &aggs[0]));
438445
}

‎crates/lower/src/promql.rs‎

Lines changed: 59 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -152,7 +152,7 @@ fn walk(expr: &Expr) -> Result<L2> {
152152
match expr {
153153
Expr::Aggregate(agg) => walk_aggregate(agg),
154154
Expr::Call(call) if call.func.name == "histogram_quantile" => walk_histogram_quantile(call),
155-
Expr::Call(call) => build(lower_inner_call(call)?, vec![], Outer::None),
155+
Expr::Call(call) => walk_call(call),
156156
Expr::Binary(bin) => walk_binary(bin),
157157
Expr::Paren(p) => walk(&p.expr),
158158
// `UnaryExpr` is built only by negation (`Neg`); unary `+` is folded to
@@ -190,6 +190,64 @@ fn walk(expr: &Expr) -> Result<L2> {
190190
}
191191
}
192192

193+
/// Lower a bare function call (`rate(m[5m])`, `max_over_time(m[5m])`, …).
194+
///
195+
/// The common case routes through the flat `lower_inner_call` template. The one
196+
/// exception is a `*_over_time`/`quantile_over_time` function applied to a
197+
/// **sub-query** (`max_over_time(rate(m[5m])[1h:])`): its argument is a
198+
/// `PromQLSubquery`, not a matrix selector, so the flat template's
199+
/// `extract_matrix` can't accept it. Lower the sub-query recursively and reduce
200+
/// it per series (issue #27).
201+
fn walk_call(call: &Call) -> Result<L2> {
202+
if let Some((func, arg_expr)) = over_time_reducer(call)? {
203+
if is_subquery(arg_expr) {
204+
// `<agg>_over_time(<subquery>)` reduces the sub-query's range vector
205+
// *per series* over time — label-preserving, no group keys. The
206+
// `PromQLSubquery` child is the structural range marker that confers
207+
// per-series semantics on the otherwise cross-series reducer (mirrors
208+
// the `TimeRange` marker for `*_over_time(m[w])`). A cross-series
209+
// aggregation operator (`sum`/`avg`) over a range vector is a PromQL
210+
// type error the parser already rejects, so an `Aggregate` directly
211+
// over a `PromQLSubquery` only ever arises here.
212+
return Ok(outer_aggregate(vec![], func, walk(arg_expr)?));
213+
}
214+
}
215+
build(lower_inner_call(call)?, vec![], Outer::None)
216+
}
217+
218+
/// The per-series reducer for a `*_over_time` range-vector function together
219+
/// with the expression in its matrix/sub-query argument slot (`φ` for
220+
/// `quantile_over_time` is read from arg 0). Returns `None` for any other call,
221+
/// so non-`*_over_time` functions fall through to the flat template.
222+
fn over_time_reducer(call: &Call) -> Result<Option<(AggFunc, &Expr)>> {
223+
let simple = |f: InnerFunc| -> Result<Option<(AggFunc, &Expr)>> {
224+
Ok(Some((inner_func(&f), arg(call, 0)?)))
225+
};
226+
match call.func.name {
227+
"avg_over_time" => simple(InnerFunc::Avg),
228+
"min_over_time" => simple(InnerFunc::Min),
229+
"max_over_time" => simple(InnerFunc::Max),
230+
"sum_over_time" => simple(InnerFunc::Sum),
231+
"stddev_over_time" => simple(InnerFunc::StdDev),
232+
"stdvar_over_time" => simple(InnerFunc::Variance),
233+
"count_over_time" => simple(InnerFunc::Count),
234+
"quantile_over_time" => {
235+
let phi = quantile_param(num_arg(call, 0)?)?;
236+
Ok(Some((inner_func(&InnerFunc::Quantile(phi)), arg(call, 1)?)))
237+
}
238+
_ => Ok(None),
239+
}
240+
}
241+
242+
/// A (parenthesised) PromQL sub-query — `<inst>[range:res]`.
243+
fn is_subquery(expr: &Expr) -> bool {
244+
match expr {
245+
Expr::Subquery(_) => true,
246+
Expr::Paren(p) => is_subquery(&p.expr),
247+
_ => false,
248+
}
249+
}
250+
193251
fn walk_aggregate(agg: &AggregateExpr) -> Result<L2> {
194252
let keys = resolve_group(agg)?;
195253
let outer = outer_kind(agg)?;

‎crates/lower/tests/promql_conformance.rs‎

Lines changed: 53 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -648,11 +648,59 @@ fn subquery_wraps_inner_query() {
648648
}
649649

650650
#[test]
651-
fn over_time_of_subquery_is_rejected__GAP() {
652-
// SEMANTICS (PromQL): `max_over_time(rate(...)[1h:])` chains a subquery into
653-
// a range-vector function. `extract_matrix` doesn't accept a subquery arg,
654-
// so this canonical pattern is rejected today.
655-
let _ = rejected("max_over_time(rate(demo_api_request_duration_seconds_count[5m])[1h:])");
651+
fn over_time_of_subquery_reduces_per_series() {
652+
// SEMANTICS (PromQL): `max_over_time(rate(...)[1h:])` chains a sub-query into
653+
// a range-vector function — the sub-query evaluates `rate` across a 1h range,
654+
// then `max_over_time` takes the max of those samples *per series*. It lowers
655+
// to a per-series `Max` reduction over a `Subquery` (issue #27).
656+
let qe = ok("max_over_time(rate(demo_api_request_duration_seconds_count[5m])[1h:])");
657+
let QueryExpr::Aggregate { by, aggs, child, .. } = &qe else {
658+
panic!("expected an Aggregate at the root, got {qe:?}");
659+
};
660+
assert!(by.is_empty(), "`*_over_time` has no grouping — reduces per series");
661+
assert!(matches!(aggs.as_slice(), [AggIntent::Max { .. }]));
662+
// The reduction rides directly on the sub-query (the structural range marker
663+
// that keeps it label-preserving), which wraps the inner `rate`.
664+
assert!(
665+
matches!(child.as_ref(), QueryExpr::Subquery { .. }),
666+
"the `Max` reduces over a Subquery, got {child:?}"
667+
);
668+
assert!(intents(&qe).iter().any(|i| matches!(i, AggIntent::Rate)));
669+
}
670+
671+
#[test]
672+
fn quantile_over_time_of_subquery_carries_phi() {
673+
// The `quantile_over_time` φ parameter is read from arg 0; the sub-query is
674+
// arg 1. It lowers to a per-series `Quantile(φ)` over the `Subquery`.
675+
let qe = ok("quantile_over_time(0.9, rate(demo[5m])[1h:])");
676+
let QueryExpr::Aggregate { aggs, child, .. } = &qe else {
677+
panic!("expected an Aggregate, got {qe:?}");
678+
};
679+
assert!(matches!(aggs.as_slice(), [AggIntent::Quantile { q, .. }] if (*q - 0.9).abs() < 1e-9));
680+
assert!(matches!(child.as_ref(), QueryExpr::Subquery { .. }));
681+
}
682+
683+
#[test]
684+
fn aggregation_over_over_time_of_subquery_keeps_labels() {
685+
// `sum by (job) (max_over_time(rate(m[5m])[1h:]))` — the inner
686+
// `max_over_time` is per-series (label-preserving), so the `job` label
687+
// survives for the OUTER cross-series `sum by (job)` to group on. If the
688+
// inner `Max` collapsed labels, `job` would not resolve here.
689+
let qe = ok("sum by (job) (max_over_time(rate(demo{job=\"api\"}[5m])[1h:]))");
690+
let QueryExpr::Aggregate { by, aggs, child, .. } = &qe else {
691+
panic!("expected outer Aggregate, got {qe:?}");
692+
};
693+
assert!(!by.is_empty(), "outer `sum by (job)` groups on a label");
694+
assert!(matches!(aggs.as_slice(), [AggIntent::Sum { .. }]));
695+
// Inner node is the per-series `max_over_time` reduction over the subquery.
696+
let QueryExpr::Aggregate { by: inner_by, aggs: inner_aggs, child: inner_child, .. } =
697+
child.as_ref()
698+
else {
699+
panic!("expected inner Aggregate, got {child:?}");
700+
};
701+
assert!(inner_by.is_empty());
702+
assert!(matches!(inner_aggs.as_slice(), [AggIntent::Max { .. }]));
703+
assert!(matches!(inner_child.as_ref(), QueryExpr::Subquery { .. }));
656704
}
657705

658706
// ─────────────────────────────────────────────────────────────────────────────

0 commit comments

Comments
 (0)