Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
150 changes: 149 additions & 1 deletion control_plane/src/query_plan/residual.rs
Original file line number Diff line number Diff line change
Expand Up @@ -328,6 +328,98 @@ fn horizons(expr: &planner_types::pre_asap::QueryExpr, out: &mut Vec<u64>) {

/// Match residuals by semantic IR equality, not display text or source names.
/// This ensures a subtree parsed for physical lowering is the subtree Planner kept.
/// Does this re-parsed subtree denote the same computation as the Planner
/// fragment?
///
/// Not `==`, and deliberately so. Planner documents a PromQL leaf's `schema` as
/// usage-derived: "the `(ts, value)` floor + the labels the query references",
/// and marks it `closed: false` precisely because it does not enumerate the
/// row. A fragment resolved inside the whole query therefore carries every
/// label the *query* mentions, while the same fragment re-parsed on its own
/// carries only the labels *it* mentions.
///
/// So `sum by (label_0) (rate(data[1m]))` yields a residual whose leaf scan has
/// columns `[ts, value, label_0]`, while re-parsing the subtree `rate(data[1m])`
/// yields `[ts, value]`. Identical source, predicates, range, measures and
/// reduction; one extra column that the isolated parse had no way to know
/// about. Requiring equality there asks an isolated parse to reproduce
/// whole-query context, and every other part of the comparison is what actually
/// discriminates: a different matcher, range or metric still fails.
///
/// Open leaf schemas are therefore compared by containment. Containment is a
/// *prefix*, not an arbitrary subset, because `ColumnId`s are positional: the
/// floor comes first and context only appends, so a prefix keeps every column
/// id in `predicates` and grouping keys meaning the same column on both sides.
/// Closed (catalog-backed SQL) schemas do enumerate the row, so they keep exact
/// equality.
fn fragment_matches(
candidate: &planner_types::pre_asap::QueryExpr,
residual: &planner_types::pre_asap::QueryExpr,
) -> bool {
match (
serde_json::to_value(candidate),
serde_json::to_value(residual),
) {
(Ok(candidate), Ok(residual)) => same_modulo_open_leaf_schema(&candidate, &residual),
// Fall back to the strict comparison rather than accepting anything we
// could not inspect.
_ => candidate == residual,
}
}

fn same_modulo_open_leaf_schema(
candidate: &serde_json::Value,
residual: &serde_json::Value,
) -> bool {
use serde_json::Value;
match (candidate, residual) {
(Value::Object(candidate), Value::Object(residual)) => {
if is_open_schema(candidate) && is_open_schema(residual) {
return open_schema_is_widened(candidate, residual);
}
candidate.len() == residual.len()
&& candidate.iter().all(|(key, value)| {
residual
.get(key)
.is_some_and(|other| same_modulo_open_leaf_schema(value, other))
})
}
(Value::Array(candidate), Value::Array(residual)) => {
candidate.len() == residual.len()
&& candidate
.iter()
.zip(residual)
.all(|(a, b)| same_modulo_open_leaf_schema(a, b))
}
_ => candidate == residual,
}
}

fn is_open_schema(value: &serde_json::Map<String, serde_json::Value>) -> bool {
value.get("closed") == Some(&serde_json::Value::Bool(false)) && value.contains_key("columns")
}

/// The isolated parse's columns must be a prefix of the whole-query ones:
/// context appends the labels it references, it never removes or reorders the
/// floor. Every other schema field still has to agree exactly.
fn open_schema_is_widened(
candidate: &serde_json::Map<String, serde_json::Value>,
residual: &serde_json::Map<String, serde_json::Value>,
) -> bool {
let (Some(narrow), Some(wide)) = (
candidate.get("columns").and_then(|v| v.as_array()),
residual.get("columns").and_then(|v| v.as_array()),
) else {
return false;
};
candidate
.iter()
.filter(|(key, _)| key.as_str() != "columns")
.all(|(key, value)| residual.get(key) == Some(value))
&& narrow.len() <= wide.len()
&& narrow.iter().zip(wide).all(|(a, b)| a == b)
}

pub(super) fn residual_nodes(
original: &str,
residual: &planner_types::pre_asap::QueryExpr,
Expand Down Expand Up @@ -389,7 +481,7 @@ pub(super) fn residual_nodes(
accuracy.clone(),
*interval,
) {
if &candidate == residual {
if fragment_matches(&candidate, residual) {
let mut lower = Lower {
nodes: BTreeMap::new(),
seen: BTreeMap::new(),
Expand Down Expand Up @@ -1406,6 +1498,62 @@ mod remote_boundary_regressions {
#[cfg(test)]
mod tests {
use super::*;
// A grouped query's residual carries the grouping label in its leaf scan
// schema, because Planner resolves that schema against the whole query.
// Re-parsing the subtree alone cannot know the label, so requiring equal
// column sets rejected a fragment that is the subtree.
#[test]
fn grouped_query_residual_matches_its_own_subtree() {
for (query, subtree) in [
("sum by (label_0) (rate(data[1m]))", "rate(data[1m])"),
(
"sum by (label_0) (sum_over_time(data[1m]))",
"sum_over_time(data[1m])",
),
] {
let residual = crate::query_parser::parse_query_expr_with_interval(
subtree,
planner_types::types::AccuracyTarget::Exact,
60_000,
)
.unwrap();
assert!(
residual_nodes(query, &residual).is_ok(),
"{query}: residual {subtree} must resolve against its own query"
);
}
}

// Widening is only accepted for the context-derived leaf columns. Anything
// that actually identifies the computation still has to match exactly.
#[test]
fn widened_leaf_schema_does_not_excuse_a_different_computation() {
let residual = crate::query_parser::parse_query_expr_with_interval(
"rate(data[1m])",
planner_types::types::AccuracyTarget::Exact,
60_000,
)
.unwrap();
// Different metric.
assert!(residual_nodes("sum by (label_0) (rate(other[1m]))", &residual).is_err());
// Different range.
assert!(residual_nodes("sum by (label_0) (rate(data[2m]))", &residual).is_err());
// Different function.
assert!(residual_nodes("sum by (label_0) (increase(data[1m]))", &residual).is_err());
// Different matcher.
let filtered = crate::query_parser::parse_query_expr_with_interval(
"rate(data{job=\"api\"}[1m])",
planner_types::types::AccuracyTarget::Exact,
60_000,
)
.unwrap();
assert!(residual_nodes(
"sum by (label_0) (rate(data{job=\"worker\"}[1m]))",
&filtered
)
.is_err());
}

// A workload horizon changes the equality witness, never its filter or explicit range.
#[test]
fn workload_horizon_residual_keeps_semantic_equality() {
Expand Down
Loading