From e699985f5098e165bf682013915a9dc535cd3535 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Wed, 7 Oct 2026 23:00:43 +0000 Subject: [PATCH 1/6] feat(ir): derive summary coverage as definition + selection MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Coverage is derived from the node, never declared. For a SummaryAgg, predicate conjuncts that lift through Filter, a range TimeRange, a TimeShift without @ and direct-column Project items, and are a value set or an interval on one column, form the selection; everything else stays in the definition (the SummaryAgg with the selection removed). A TimeRange(w) over TimeShift(s) gives the relative window (-(s+w), -s]. A SummaryMerge is valid only when its inputs have structurally equal definitions (ignoring timing and guarantee) and pairwise disjoint selections; OperatorNode::new and validate_structure enforce it. Its coverage is the shared definition and the union of the selections, joining adjacent windows and value sets. OperatorNode.coverage becomes a private cache behind coverage(): ignored by equality, skipped by serde, emptied on clone. with_coverage, requires_coverage and CoverageError::Missing are removed, as are the coverage fields of CSE keys and FlatNode. Design: docs/design_docs/proposals/asap-primitive-schema.md §4 (#573). Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- crates/types/src/ir/cse.rs | 13 +- crates/types/src/ir/flat.rs | 4 - crates/types/src/ir/node.rs | 75 +- crates/types/src/ir/summary_coverage.rs | 750 +++++++++++++++--- crates/types/tests/schema_rebuilding.rs | 40 +- crates/types/tests/structure_contract.rs | 17 - crates/types/tests/summary_coverage.rs | 488 +++++++++--- crates/types/tests/summary_merge_structure.rs | 80 +- 8 files changed, 1108 insertions(+), 359 deletions(-) diff --git a/crates/types/src/ir/cse.rs b/crates/types/src/ir/cse.rs index b9046a1ab..c229772f7 100644 --- a/crates/types/src/ir/cse.rs +++ b/crates/types/src/ir/cse.rs @@ -112,7 +112,6 @@ pub fn structural_hash(node: &OperatorNode, cache: &mut HashCache) -> u64 { &node.schema, &node.guarantee, node.timing, - &node.coverage, ); serde_json::to_string(&own) .unwrap_or_default() @@ -168,7 +167,6 @@ fn same_node(left: &OperatorNode, right: &OperatorNode, memo: &mut EqMemo) -> bo && left.result_kind == right.result_kind && left.schema == right.schema && left.timing == right.timing - && left.coverage == right.coverage && same_value(&left.guarantee, &right.guarantee) && same_value(&own_fields(left), &own_fields(right)) } @@ -250,14 +248,9 @@ fn intern_bottom_up( let operator = node .operator .map_children(|child| intern_bottom_up(table, visited, child)); - let rebuilt = OperatorNode { - operator, - result_kind: node.result_kind, - schema: node.schema.clone(), - guarantee: node.guarantee.clone(), - timing: node.timing, - coverage: node.coverage.clone(), - }; + let rebuilt = OperatorNode::with_schema(operator, node.schema.clone()) + .with_guarantee(node.guarantee.clone()) + .with_timing(node.timing); let interned = table.intern(rebuilt); visited.insert(Rc::as_ptr(node), (Rc::clone(node), Rc::clone(&interned))); interned diff --git a/crates/types/src/ir/flat.rs b/crates/types/src/ir/flat.rs index 830a294db..434a3406f 100644 --- a/crates/types/src/ir/flat.rs +++ b/crates/types/src/ir/flat.rs @@ -15,7 +15,6 @@ use serde::{Deserialize, Serialize}; use super::node::{Operator, OperatorNode, OperatorResultKind}; use super::query::QueryRoot; -use super::summary_coverage::SummaryCoverage; use crate::post_asap::execution_data_state::ExecutionTiming; use crate::post_asap::guarantee::ResultGuarantee; use crate::pre_asap::schema::Schema; @@ -31,8 +30,6 @@ pub struct FlatNode { pub schema: Schema, pub guarantee: Option, pub timing: Option, - #[serde(default)] - pub coverage: Option, } #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -80,7 +77,6 @@ fn visit( schema: node.schema.clone(), guarantee: node.guarantee.clone(), timing: node.timing, - coverage: node.coverage.clone(), }); originals.push(Rc::clone(node)); ids.insert(Rc::as_ptr(node), id); diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index 47ffa0ea3..d58c362cd 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -3,6 +3,7 @@ //! every traversal needs (output category and schema, accuracy guarantee, //! execution timing). +use std::cell::OnceCell; use std::collections::HashSet; use std::rc::Rc; @@ -10,7 +11,7 @@ use serde::{Deserialize, Serialize}; use super::asap::ASAPOp; use super::non_asap::NonASAPOp; -use super::summary_coverage::{CoverageError, SummaryCoverage}; +use super::summary_coverage::SummaryCoverage; use crate::ir::SchemaDerivationError; use crate::post_asap::execution_data_state::ExecutionTiming; use crate::post_asap::guarantee::ResultGuarantee; @@ -100,8 +101,33 @@ pub struct OperatorNode { pub schema: Schema, pub guarantee: Option, pub timing: Option, - #[serde(default)] - pub coverage: Option, + /// Cache for [`Self::coverage`], derived from `operator`. + #[serde(skip)] + coverage_cache: CoverageCache, +} + +/// A lazily filled [`SummaryCoverage`]. It is not part of a node's value: +/// equality ignores it, serialization skips it, and a clone starts empty so +/// a clone whose operator is then edited cannot read a stale entry. +#[derive(Default)] +struct CoverageCache(OnceCell>); + +impl Clone for CoverageCache { + fn clone(&self) -> Self { + Self::default() + } +} + +impl PartialEq for CoverageCache { + fn eq(&self, _: &Self) -> bool { + true + } +} + +impl std::fmt::Debug for CoverageCache { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str("CoverageCache") + } } impl OperatorNode { @@ -110,7 +136,9 @@ impl OperatorNode { /// ASAP operator, ...). pub fn new(operator: Operator) -> Result { let schema = operator.output_schema()?; - Ok(Self::with_schema(operator, schema)) + let node = Self::with_schema(operator, schema); + node.check_merge()?; + Ok(node) } /// Build a node with caller-supplied output names and qualifiers. For @@ -124,7 +152,7 @@ impl OperatorNode { schema, guarantee: None, timing: None, - coverage: None, + coverage_cache: CoverageCache::default(), } } @@ -145,23 +173,22 @@ impl OperatorNode { self } - /// Attach caller-established coverage. Required on summary nodes; see - /// [`Self::requires_coverage`]. - pub fn with_coverage( - mut self, - coverage: SummaryCoverage, - ) -> Result { - coverage.validate()?; - if self.result_kind != OperatorResultKind::State { - return Err(CoverageError::NotState.into()); - } - self.coverage = Some(coverage); - Ok(self) + /// What this summary state covers; `None` for a node that is not a + /// `SummaryAgg` or a valid `SummaryMerge`. Derived on first use. + pub fn coverage(&self) -> Option<&SummaryCoverage> { + self.coverage_cache + .0 + .get_or_init(|| SummaryCoverage::derive(self).ok()) + .as_ref() } - /// Summary nodes whose state can be composed must declare coverage. - pub fn requires_coverage(&self) -> bool { - matches!(self.asap(), Some(ASAPOp::SummaryAgg { .. })) + /// A `SummaryMerge` is valid only over inputs with the same definition + /// and disjoint selections. + fn check_merge(&self) -> Result<(), SchemaDerivationError> { + if matches!(self.asap(), Some(ASAPOp::SummaryMerge { .. })) { + SummaryCoverage::derive(self)?; + } + Ok(()) } pub fn non_asap(&self) -> Option<&NonASAPOp> { @@ -315,14 +342,8 @@ impl OperatorNode { "invalid time or identity column in schema".into(), )); } - match &node.coverage { - Some(coverage) => { - (*node.as_ref()).clone().with_coverage(coverage.clone())?; - } - None if node.requires_coverage() => return Err(CoverageError::Missing.into()), - None => {} - } node.operator.validate_inputs()?; + node.check_merge()?; if node.result_kind != node.operator.output_kind() { return Err(SchemaDerivationError::InvalidScalarSignature( "retained result kind disagrees with operation".into(), diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index 3f3798a88..da65128ee 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -1,127 +1,675 @@ -//! Joint time/population coverage for summary composition, independent of schema. -//! Equality predicates are a deliberately narrow proof vocabulary. Unsupported -//! predicates cannot be declared disjoint merely by giving them different names. -use crate::pre_asap::Source; -use serde::{Deserialize, Serialize}; -use std::collections::BTreeMap; -use std::ops::Range; +//! What a summary state covers: `definition` (what it computes) and +//! `selection` (which output rows of that computation it took). Design: +//! `docs/design_docs/proposals/asap-primitive-schema.md` §4, after +//! Goldstein & Larson's view matching. +//! +//! Coverage is derived from the node, never declared. Walking down from a +//! `SummaryAgg`, a predicate conjunct moves into `selection` when it can be +//! lifted to the `SummaryAgg` (through `Filter`, a range `TimeRange`, a +//! `TimeShift` without `@` and direct-column `Project` items) and it is a +//! value set or an interval on one column. Everything else stays in +//! `definition` as a residual, so two states are merged only when they +//! compute the same thing over disjoint rows. +use std::cmp::Ordering; +use std::collections::{BTreeMap, HashSet}; +use std::ops::Bound; +use std::rc::Rc; + use thiserror::Error; -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] +use super::asap::ASAPOp; +use super::node::{Operator, OperatorNode}; +use super::non_asap::{NonASAPOp, TimeRangeKind}; +use super::scalar::{Predicate, ScalarExpr}; +use crate::pre_asap::expr_ir::{CompareOpKind, ScalarValue}; +use crate::pre_asap::schema::ColumnId; + +#[derive(Debug, Clone, PartialEq)] pub struct SummaryCoverage { - /// Observation data source, as named by `Scan`: a table or a time series. - /// Region time bounds refer to its time column. - pub source: Source, - /// Union of joint regions; never the Cartesian product of independent bounds. - pub regions: Vec, -} - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct CoverageRegion { - /// Half-open absolute bounds on the source's time column, in milliseconds. - /// Only a materialized state (for example an ingested pane) has absolute - /// bounds; a state in a logical plan that is not yet evaluated uses `None`, - /// as does a source without a time column. `None` means no time restriction. - pub time_ms: Option>, - /// Conjunction of `field = 'text'` equality predicates (PromQL label - /// matchers, SQL text columns), keyed by field name. Only Utf8 equality is - /// represented, so values are strings; any other restriction is not - /// recorded here. Empty means unrestricted. - pub population: BTreeMap, + /// The `SummaryAgg` with the selection removed from its sub-DAG. + pub definition: Rc, + /// Union of boxes over the output rows of the definition's child. + pub selection: Vec, +} + +/// A conjunction of per-column constraints and an optional time window. +/// A column or time it does not mention is unrestricted. +#[derive(Debug, Clone, PartialEq, Default)] +pub struct SelectionBox { + pub columns: BTreeMap, + /// Offsets from the evaluation time in milliseconds: a PromQL range + /// `TimeRange(w)` over `TimeShift(s)` is `(-(s + w), -s]`. Absolute + /// time is an ordinary interval on the timestamp column. + pub relative_time: Option<(Bound, Bound)>, +} + +/// A column of the definition child's output, by `(table, name)`. +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] +pub struct ColumnIdentity { + pub table: Option, + pub name: String, +} + +#[derive(Debug, Clone, PartialEq)] +pub enum Constraint { + In(Vec), + NotIn(Vec), + Interval { + lower: Bound, + upper: Bound, + }, } #[derive(Debug, Clone, PartialEq, Eq, Error)] pub enum CoverageError { - #[error("coverage interval must have start < end")] - InvalidInterval, - #[error("population dimension names cannot be empty")] - InvalidPopulation, - #[error("summary coverage sources differ")] - SourceMismatch, - #[error("coverage overlap is not proven absent")] - PossibleOverlap, - #[error("coverage merge requires at least one input")] + #[error("coverage is derived only for SummaryAgg and SummaryMerge")] + NotSummary, + #[error("summary merge requires at least one input")] EmptyMerge, - #[error("summary coverage requires state output")] - NotState, - #[error("summary node requires coverage")] - Missing, + #[error("summary merge inputs compute different things")] + DefinitionMismatch, + #[error("summary merge inputs are not proven disjoint")] + PossibleOverlap, } impl SummaryCoverage { - pub fn validate(&self) -> Result<(), CoverageError> { - for (index, region) in self.regions.iter().enumerate() { - if region.time_ms.as_ref().is_some_and(Range::is_empty) { - return Err(CoverageError::InvalidInterval); - } - if region.population.keys().any(String::is_empty) { - return Err(CoverageError::InvalidPopulation); + /// Coverage of a `SummaryAgg` or `SummaryMerge`. A merge fails unless + /// every input has the same definition and their selections are + /// pairwise disjoint. + pub fn derive(node: &OperatorNode) -> Result { + match node.asap() { + Some(ASAPOp::SummaryAgg { .. }) => Ok(of_summary_agg(node)), + Some(ASAPOp::SummaryMerge { children }) => { + let inputs = children + .iter() + .map(|child| child.coverage().ok_or(CoverageError::NotSummary)) + .collect::, _>>()?; + let first = inputs.first().ok_or(CoverageError::EmptyMerge)?; + let mut proven = HashSet::new(); + for (i, input) in inputs.iter().enumerate() { + if !same_definition(&first.definition, &input.definition, &mut proven) { + return Err(CoverageError::DefinitionMismatch); + } + let overlaps = inputs[..i].iter().any(|other| { + input + .selection + .iter() + .any(|a| other.selection.iter().any(|b| !a.disjoint(b))) + }); + if overlaps { + return Err(CoverageError::PossibleOverlap); + } + } + Ok(Self { + definition: Rc::clone(&first.definition), + selection: union(inputs.iter().flat_map(|c| c.selection.clone()).collect()), + }) } - if self.regions[..index] + _ => Err(CoverageError::NotSummary), + } + } +} + +/// Structural equality ignoring planning metadata (`timing`, `guarantee`). +/// `proven` memoizes node pairs already found equal. +fn same_definition( + a: &Rc, + b: &Rc, + proven: &mut HashSet<(*const OperatorNode, *const OperatorNode)>, +) -> bool { + let key = (Rc::as_ptr(a), Rc::as_ptr(b)); + if Rc::ptr_eq(a, b) || proven.contains(&key) { + return true; + } + let (ac, bc) = (a.children(), b.children()); + let equal = a.operator.map_children(|_| ()) == b.operator.map_children(|_| ()) + && a.result_kind == b.result_kind + && a.schema == b.schema + && ac.len() == bc.len() + && ac + .iter() + .zip(&bc) + .all(|(x, y)| same_definition(x, y, proven)); + if equal { + proven.insert(key); + } + equal +} + +fn of_summary_agg(node: &OperatorNode) -> SummaryCoverage { + let Some(ASAPOp::SummaryAgg { + child, + family, + input, + reduction, + grouping, + filter, + }) = node.asap() + else { + unreachable!("called on a SummaryAgg"); + }; + + // The chain a predicate can be lifted through, top to bottom, and the + // first node below it. + let mut chain: Vec<&Rc> = Vec::new(); + let mut base = child; + while let Some(op) = base.non_asap() { + let passes = match op { + NonASAPOp::Filter { .. } | NonASAPOp::Project { .. } => true, + NonASAPOp::TimeRange { kind, .. } => *kind == TimeRangeKind::Range, + NonASAPOp::TimeShift { shift, .. } => shift.at.is_none(), + _ => false, + }; + if !passes { + break; + } + chain.push(base); + base = base.children()[0]; + } + // Time is lifted only from exactly one range window. + let lift_time = chain + .iter() + .filter(|link| matches!(link.non_asap(), Some(NonASAPOp::TimeRange { .. }))) + .count() + == 1; + + // Top-down: decide which conjuncts are lifted. `kept[d]` is the residual + // of `chain[d]` when it is a `Filter`. + let mut lifted = BTreeMap::new(); + let agg_filter = filter + .as_ref() + .and_then(|pred| residual(pred, &[], &mut lifted)); + let kept: Vec> = chain + .iter() + .enumerate() + .map(|(depth, link)| match link.non_asap() { + Some(NonASAPOp::Filter { pred, .. }) => residual(pred, &chain[..depth], &mut lifted), + _ => None, + }) + .collect(); + let scan_kept = match base.non_asap() { + Some(NonASAPOp::Scan { predicates, .. }) => Some( + predicates .iter() - .any(|other| region.may_overlap(other)) - { - return Err(CoverageError::PossibleOverlap); + .filter_map(|pred| residual(pred, &chain, &mut lifted)) + .collect::>(), + ), + _ => None, + }; + + // Bottom-up: rebuild without what was lifted. A node whose input and + // operator are unchanged keeps its `Rc`. + let mut rebuilt = Rc::clone(base); + if let ( + Some(kept), + Some(NonASAPOp::Scan { + source, + predicates, + schema, + }), + ) = (scan_kept, base.non_asap()) + { + if kept.len() != predicates.len() { + rebuilt = rebuild( + base, + NonASAPOp::Scan { + source: source.clone(), + predicates: kept, + schema: schema.clone(), + }, + ); + } + } + let (mut window_ms, mut shift_ms) = (0i64, 0i64); + for (depth, link) in chain.iter().enumerate().rev() { + let child = Rc::clone(&rebuilt); + let op = match link.non_asap().expect("chain nodes are non-ASAP") { + NonASAPOp::Filter { pred, .. } => match &kept[depth] { + None => continue, + Some(rest) if rest == pred && Rc::ptr_eq(&child, link.children()[0]) => { + rebuilt = Rc::clone(link); + continue; + } + Some(rest) => NonASAPOp::Filter { + pred: rest.clone(), + child, + }, + }, + NonASAPOp::TimeRange { range, .. } if lift_time => { + window_ms = range.as_millis() as i64; + continue; + } + NonASAPOp::TimeShift { shift, .. } if lift_time => { + shift_ms += shift.offset_ms; + continue; + } + _ if Rc::ptr_eq(&child, link.children()[0]) => { + rebuilt = Rc::clone(link); + continue; + } + NonASAPOp::TimeRange { range, kind, .. } => NonASAPOp::TimeRange { + range: *range, + kind: *kind, + child, + }, + NonASAPOp::TimeShift { shift, .. } => NonASAPOp::TimeShift { + shift: *shift, + child, + }, + NonASAPOp::Project { + cols, qualifier, .. + } => NonASAPOp::Project { + cols: cols.clone(), + qualifier: qualifier.clone(), + child, + }, + _ => unreachable!("only liftable operators are on the chain"), + }; + rebuilt = rebuild(link, op); + } + let definition = rebuild_asap( + node, + ASAPOp::SummaryAgg { + child: rebuilt, + family: family.clone(), + input: input.clone(), + reduction: reduction.clone(), + grouping: grouping.clone(), + filter: agg_filter, + }, + ); + + let top = &child.schema; + let selection = SelectionBox { + columns: lifted + .into_iter() + .map(|(column, constraint)| { + let field = &top.fields[column]; + let identity = ColumnIdentity { + table: field.table.clone(), + name: field.name.clone(), + }; + (identity, constraint) + }) + .collect(), + relative_time: lift_time.then_some(( + Bound::Excluded(-(shift_ms + window_ms)), + Bound::Included(-shift_ms), + )), + }; + SummaryCoverage { + definition, + selection: vec![selection], + } +} + +/// A node with the same output schema over a new operator. The definition +/// describes what was computed; it is not re-validated as a plan. +fn rebuild(original: &OperatorNode, op: NonASAPOp) -> Rc { + Rc::new(OperatorNode::with_schema( + Operator::NonASAP(op), + original.schema.clone(), + )) +} + +fn rebuild_asap(original: &OperatorNode, op: ASAPOp) -> Rc { + Rc::new(OperatorNode::with_schema( + Operator::ASAP(op), + original.schema.clone(), + )) +} + +/// Lift what `pred`'s conjuncts can into `lifted` (keyed by the agg child's +/// column) and return the rest. `above` are the chain nodes between the +/// predicate and the `SummaryAgg`, top to bottom. +fn residual( + pred: &Predicate, + above: &[&Rc], + lifted: &mut BTreeMap, +) -> Option { + let mut rest: Vec = Vec::new(); + for conjunct in pred.0.conjuncts() { + let lift = constraint_of(conjunct).and_then(|(column, constraint)| { + let column = column_at_top(column, above)?; + let combined = match lifted.get(&column) { + Some(existing) => existing.intersect(&constraint)?, + None => constraint, + }; + Some((column, combined)) + }); + match lift { + Some((column, constraint)) => { + lifted.insert(column, constraint); } + None => rest.push(conjunct.clone()), } - Ok(()) - } - - /// Every observation in a region is assumed to contribute once to the state. - /// Compose once-per-observation summaries only when their joint regions are - /// provably disjoint. Update/reduction compatibility, family merge capability - /// and accuracy are checked by `SummaryMerge`, not here. - pub fn merge_disjoint(inputs: &[Self]) -> Result { - let first = inputs.first().ok_or(CoverageError::EmptyMerge)?; - let mut merged = first.clone(); - merged.regions.clear(); - for input in inputs { - input.validate()?; - if input.source != first.source { - return Err(CoverageError::SourceMismatch); + } + match rest.len() { + 0 => None, + 1 => rest.pop().map(Predicate), + _ => Some(Predicate(ScalarExpr::BoolAnd(rest))), + } +} + +/// Map a column of the input of `above.last()` up to the agg child's output: +/// a `Project` passes it only as a direct column item. +fn column_at_top(mut column: ColumnId, above: &[&Rc]) -> Option { + for link in above.iter().rev() { + if let Some(NonASAPOp::Project { cols, .. }) = link.non_asap() { + column = cols + .iter() + .position(|item| item.expr == ScalarExpr::Column(column))?; + } + } + Some(column) +} + +/// `column = v`, `!=`, `<`, `<=`, `>`, `>=`, `[NOT] IN (...)` and `OR` of +/// equalities on one column, with non-null literals. +fn constraint_of(expr: &ScalarExpr) -> Option<(ColumnId, Constraint)> { + let literal = |e: &ScalarExpr| match e { + ScalarExpr::Literal(v) if *v != ScalarValue::Null => Some(v.clone()), + _ => None, + }; + match expr { + ScalarExpr::Compare { + left, op, right, .. + } => { + let (column, value, op) = match (left.as_ref(), right.as_ref()) { + (ScalarExpr::Column(c), r) => (*c, literal(r)?, op.clone()), + (l, ScalarExpr::Column(c)) => (*c, literal(l)?, flip(op)?), + _ => return None, + }; + let constraint = match op { + CompareOpKind::Eq => Constraint::In(vec![value]), + CompareOpKind::Ne => Constraint::NotIn(vec![value]), + CompareOpKind::Lt => interval(Bound::Unbounded, Bound::Excluded(value)), + CompareOpKind::Le => interval(Bound::Unbounded, Bound::Included(value)), + CompareOpKind::Gt => interval(Bound::Excluded(value), Bound::Unbounded), + CompareOpKind::Ge => interval(Bound::Included(value), Bound::Unbounded), + _ => return None, + }; + Some((column, constraint)) + } + ScalarExpr::InList { + expr, + list, + negated, + } => { + let ScalarExpr::Column(column) = expr.as_ref() else { + return None; + }; + let values = list.iter().map(literal).collect::>>()?; + let values = dedup(values); + Some(( + *column, + if *negated { + Constraint::NotIn(values) + } else { + Constraint::In(values) + }, + )) + } + ScalarExpr::BoolOr(disjuncts) => { + let mut column = None; + let mut values = Vec::new(); + for disjunct in disjuncts { + let (c, Constraint::In(v)) = constraint_of(disjunct)? else { + return None; + }; + if column.replace(c).is_some_and(|prev| prev != c) { + return None; + } + values.extend(v); } - merged.regions.extend(input.regions.iter().cloned()); + Some((column?, Constraint::In(dedup(values)))) } - merged.validate()?; - // Coalesce adjacent intervals only for identical population predicates. - merged.regions.sort_by(|a, b| { - a.population.cmp(&b.population).then( - a.time_ms - .as_ref() - .map(|t| t.start) - .cmp(&b.time_ms.as_ref().map(|t| t.start)), - ) - }); - let mut normalized: Vec = Vec::new(); - for region in merged.regions { - if let Some(last) = normalized.last_mut() { - if let (Some(last_time), Some(time)) = (&mut last.time_ms, ®ion.time_ms) { - if last.population == region.population && last_time.end == time.start { - last_time.end = time.end; - continue; + _ => None, + } +} + +fn flip(op: &CompareOpKind) -> Option { + Some(match op { + CompareOpKind::Eq => CompareOpKind::Eq, + CompareOpKind::Ne => CompareOpKind::Ne, + CompareOpKind::Lt => CompareOpKind::Gt, + CompareOpKind::Le => CompareOpKind::Ge, + CompareOpKind::Gt => CompareOpKind::Lt, + CompareOpKind::Ge => CompareOpKind::Le, + _ => return None, + }) +} + +fn interval(lower: Bound, upper: Bound) -> Constraint { + Constraint::Interval { lower, upper } +} + +fn dedup(values: Vec) -> Vec { + let mut out: Vec = Vec::new(); + for v in values { + if !out.contains(&v) { + out.push(v); + } + } + out +} + +/// Order of two literals of the same type; `None` across types or for NaN. +fn compare(a: &ScalarValue, b: &ScalarValue) -> Option { + match (a, b) { + (ScalarValue::Int64(a), ScalarValue::Int64(b)) => Some(a.cmp(b)), + (ScalarValue::Float64(a), ScalarValue::Float64(b)) => a.partial_cmp(b), + (ScalarValue::Utf8(a), ScalarValue::Utf8(b)) => Some(a.cmp(b)), + (ScalarValue::Boolean(a), ScalarValue::Boolean(b)) => Some(a.cmp(b)), + _ => None, + } +} + +/// Whether `v` lies inside the interval; `None` when it cannot be compared. +fn inside( + v: &T, + (lower, upper): (&Bound, &Bound), + cmp: impl Fn(&T, &T) -> Option, +) -> Option { + let above = match lower { + Bound::Unbounded => true, + Bound::Included(l) => cmp(v, l)? != Ordering::Less, + Bound::Excluded(l) => cmp(v, l)? == Ordering::Greater, + }; + let below = match upper { + Bound::Unbounded => true, + Bound::Included(u) => cmp(v, u)? != Ordering::Greater, + Bound::Excluded(u) => cmp(v, u)? == Ordering::Less, + }; + Some(above && below) +} + +/// Whether no value is both at most `upper` and at least `lower`. +fn ends_before( + upper: &Bound, + lower: &Bound, + cmp: impl Fn(&T, &T) -> Option, +) -> bool { + match (upper, lower) { + (Bound::Included(u), Bound::Included(l)) => cmp(u, l) == Some(Ordering::Less), + (Bound::Included(u) | Bound::Excluded(u), Bound::Included(l) | Bound::Excluded(l)) => { + matches!(cmp(u, l), Some(Ordering::Less | Ordering::Equal)) + } + _ => false, + } +} + +fn intervals_disjoint( + (al, au): (&Bound, &Bound), + (bl, bu): (&Bound, &Bound), + cmp: impl Fn(&T, &T) -> Option + Copy, +) -> bool { + ends_before(au, bl, cmp) || ends_before(bu, al, cmp) +} + +/// The tighter of two lower (`want = Greater`) or upper (`Less`) bounds. +fn tighter( + a: &Bound, + b: &Bound, + want: Ordering, +) -> Option> { + let value = |bound: &Bound| match bound { + Bound::Included(v) | Bound::Excluded(v) => Some(v.clone()), + Bound::Unbounded => None, + }; + Some(match (value(a), value(b)) { + (None, _) => b.clone(), + (_, None) => a.clone(), + (Some(x), Some(y)) => match compare(&x, &y)? { + Ordering::Equal if matches!(a, Bound::Excluded(_)) => a.clone(), + Ordering::Equal => b.clone(), + order if order == want => a.clone(), + _ => b.clone(), + }, + }) +} + +impl Constraint { + /// The conjunction of two constraints on one column, when it is again a + /// single constraint. + fn intersect(&self, other: &Self) -> Option { + use Constraint::*; + Some(match (self, other) { + (In(a), In(b)) => In(a.iter().filter(|v| b.contains(v)).cloned().collect()), + (In(a), NotIn(b)) | (NotIn(b), In(a)) => { + In(a.iter().filter(|v| !b.contains(v)).cloned().collect()) + } + (NotIn(a), NotIn(b)) => NotIn(dedup(a.iter().chain(b).cloned().collect())), + (In(a), Interval { lower, upper }) | (Interval { lower, upper }, In(a)) => { + let mut kept = Vec::new(); + for v in a { + if inside(v, (lower, upper), compare)? { + kept.push(v.clone()); } } + In(kept) } - normalized.push(region); + ( + Interval { + lower: al, + upper: au, + }, + Interval { + lower: bl, + upper: bu, + }, + ) => Interval { + lower: tighter(al, bl, Ordering::Greater)?, + upper: tighter(au, bu, Ordering::Less)?, + }, + (NotIn(_), Interval { .. }) | (Interval { .. }, NotIn(_)) => return None, + }) + } + + /// Whether no value satisfies both; `false` when unsure. + fn disjoint(&self, other: &Self) -> bool { + use Constraint::*; + match (self, other) { + (In(a), In(b)) => !a.iter().any(|v| b.contains(v)), + (In(a), NotIn(b)) | (NotIn(b), In(a)) => a.iter().all(|v| b.contains(v)), + (In(a), Interval { lower, upper }) | (Interval { lower, upper }, In(a)) => a + .iter() + .all(|v| inside(v, (lower, upper), compare) == Some(false)), + ( + Interval { + lower: al, + upper: au, + }, + Interval { + lower: bl, + upper: bu, + }, + ) => intervals_disjoint((al, au), (bl, bu), compare), + (NotIn(_), _) | (_, NotIn(_)) => false, } - merged.regions = normalized; - Ok(merged) } } -impl CoverageRegion { - fn may_overlap(&self, other: &Self) -> bool { - let time_overlaps = match (&self.time_ms, &other.time_ms) { - (Some(a), Some(b)) => a.start < b.end && b.start < a.end, - _ => true, + +impl SelectionBox { + /// Whether no row lies in both boxes: some column, or the time window, + /// is restricted by both and the restrictions are disjoint. + fn disjoint(&self, other: &Self) -> bool { + let time = match (&self.relative_time, &other.relative_time) { + (Some((al, au)), Some((bl, bu))) => { + intervals_disjoint((al, au), (bl, bu), |a: &i64, b: &i64| Some(a.cmp(b))) + } + _ => false, }; - time_overlaps - && !self.population.iter().any(|(dimension, value)| { - other - .population - .get(dimension) - .is_some_and(|other| other != value) - }) + time || self + .columns + .iter() + .any(|(column, a)| other.columns.get(column).is_some_and(|b| a.disjoint(b))) + } +} + +/// Join boxes that differ only in adjacent time windows or only in the +/// values of one `In` column; gaps stay separate boxes. +fn union(mut boxes: Vec) -> Vec { + let mut i = 0; + while i < boxes.len() { + let joined = (i + 1..boxes.len()).find_map(|j| join(&boxes[i], &boxes[j]).map(|b| (j, b))); + match joined { + Some((j, joined)) => { + boxes.remove(j); + boxes[i] = joined; + i = 0; + } + None => i += 1, + } + } + boxes +} + +fn join(a: &SelectionBox, b: &SelectionBox) -> Option { + if a.columns == b.columns { + let ((al, au), (bl, bu)) = (a.relative_time.as_ref()?, b.relative_time.as_ref()?); + let meets = |upper: &Bound, lower: &Bound| { + matches!( + (upper, lower), + (Bound::Included(u), Bound::Excluded(l)) | (Bound::Excluded(u), Bound::Included(l)) + if u == l + ) + }; + let time = if meets(au, bl) { + (*al, *bu) + } else if meets(bu, al) { + (*bl, *au) + } else { + return None; + }; + return Some(SelectionBox { + columns: a.columns.clone(), + relative_time: Some(time), + }); + } + if a.relative_time != b.relative_time || a.columns.len() != b.columns.len() { + return None; + } + let mut differing = a + .columns + .iter() + .filter(|(column, constraint)| b.columns.get(*column) != Some(*constraint)); + let (column, Constraint::In(values)) = differing.next()? else { + return None; + }; + if differing.next().is_some() { + return None; } + let Some(Constraint::In(more)) = b.columns.get(column) else { + return None; + }; + let mut columns = a.columns.clone(); + columns.insert( + column.clone(), + Constraint::In(dedup(values.iter().chain(more).cloned().collect())), + ); + Some(SelectionBox { + columns, + relative_time: a.relative_time, + }) } diff --git a/crates/types/tests/schema_rebuilding.rs b/crates/types/tests/schema_rebuilding.rs index b0fc4c666..bb9f07526 100644 --- a/crates/types/tests/schema_rebuilding.rs +++ b/crates/types/tests/schema_rebuilding.rs @@ -1,4 +1,3 @@ -use asap_types::ir::summary_coverage::{CoverageRegion, SummaryCoverage}; use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}; use asap_types::post_asap::{ ExactKind, ExactParams, ExecutionTiming, GroupingStrategy, ResultGuarantee, SummaryUpdate, @@ -8,27 +7,6 @@ use asap_types::pre_asap::{ }; use std::rc::Rc; -fn coverage() -> SummaryCoverage { - SummaryCoverage { - source: Source::Table { - table_ref: "t".into(), - }, - regions: vec![CoverageRegion { - time_ms: None, - population: Default::default(), - }], - } -} -/// Rewrites clear coverage; a rewriter must declare it again for summary nodes. -fn redeclare(node: OperatorNode) -> Rc { - let node = Rc::new(node); - if !node.requires_coverage() { - return node; - } - assert!(node.validate_structure().is_err()); - Rc::new((*node).clone().with_coverage(coverage()).unwrap()) -} - fn scan(key_type: DataType, name: &str) -> Rc { OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { source: Source::Table { @@ -62,12 +40,7 @@ fn aggregate(child: Rc, asap: bool) -> Rc { having: None, }) }; - let node = OperatorNode::new(operator).unwrap(); - Rc::new(if asap { - node.with_coverage(coverage()).unwrap() - } else { - node - }) + OperatorNode::new_shared(operator).unwrap() } /// Rewrites follow changed input types and inherited names for either category. @@ -77,7 +50,7 @@ fn rebuilding_rederives_schema_for_both_categories() { let original = aggregate(scan(DataType::Int64, "key"), asap); original.validate_structure().unwrap(); let replacement = scan(DataType::Utf8, "new_key"); - let rebuilt = redeclare(original.with_new_children(|_| replacement.clone()).unwrap()); + let rebuilt = Rc::new(original.with_new_children(|_| replacement.clone()).unwrap()); assert_eq!(rebuilt.schema, rebuilt.operator.output_schema().unwrap()); rebuilt.validate_structure().unwrap(); } @@ -91,14 +64,13 @@ fn rebuilding_preserves_only_explicit_naming_overrides() { let mut schema = original.schema.clone(); schema.fields[0].name = "alias".into(); schema.fields[0].table = Some("result".into()); - let mut renamed = OperatorNode::with_schema(original.operator.clone(), schema) + let renamed = OperatorNode::with_schema(original.operator.clone(), schema) .with_guarantee(Some(ResultGuarantee::exact("fixture"))) .with_timing(Some(ExecutionTiming::QueryTime)); - renamed.coverage = original.coverage.clone(); let original = Rc::new(renamed); original.validate_structure().unwrap(); let replacement = scan(DataType::Utf8, "new_key"); - let rebuilt = redeclare(original.with_new_children(|_| replacement.clone()).unwrap()); + let rebuilt = Rc::new(original.with_new_children(|_| replacement.clone()).unwrap()); assert_eq!(rebuilt.schema.fields[0].name, "alias"); assert_eq!(rebuilt.schema.fields[0].table.as_deref(), Some("result")); assert_eq!( @@ -137,9 +109,7 @@ fn validation_rejects_structural_overrides_for_both_categories() { schema.fields.pop(); invalid.push(schema); for schema in invalid { - let mut forged = OperatorNode::with_schema(original.operator.clone(), schema); - forged.coverage = original.coverage.clone(); - let forged = Rc::new(forged); + let forged = Rc::new(OperatorNode::with_schema(original.operator.clone(), schema)); assert!( forged.validate_structure().is_err(), "accepted structural override: {:?}", diff --git a/crates/types/tests/structure_contract.rs b/crates/types/tests/structure_contract.rs index d7c144f4e..f5f8f5cc2 100644 --- a/crates/types/tests/structure_contract.rs +++ b/crates/types/tests/structure_contract.rs @@ -13,19 +13,6 @@ fn scan() -> Rc { })) .unwrap() } -/// Tabular coverage for a whole-table summary. -fn whole_table() -> asap_types::ir::summary_coverage::SummaryCoverage { - use asap_types::ir::summary_coverage::{CoverageRegion, SummaryCoverage}; - SummaryCoverage { - source: Source::Table { - table_ref: "t".into(), - }, - regions: vec![CoverageRegion { - time_ms: None, - population: Default::default(), - }], - } -} /// Resolved filters cannot hide invalid scalar types or out-of-scope columns. #[test] fn invalid_predicates_are_rejected() { @@ -116,8 +103,6 @@ fn state_evaluations_and_passthrough_keep_their_contracts() { grouping: GroupingStrategy::default(), filter: None, })) - .unwrap() - .with_coverage(whole_table()) .unwrap(), ); state.validate_structure().unwrap(); @@ -197,8 +182,6 @@ fn shared_construction_derives_both_operator_categories() { grouping: GroupingStrategy::default(), filter: None, })) - .unwrap() - .with_coverage(whole_table()) .unwrap(), ); assert_eq!(state.result_kind, OperatorResultKind::State); diff --git a/crates/types/tests/summary_coverage.rs b/crates/types/tests/summary_coverage.rs index e36c3d7a4..3ef5890f3 100644 --- a/crates/types/tests/summary_coverage.rs +++ b/crates/types/tests/summary_coverage.rs @@ -1,150 +1,412 @@ -//! Coverage composition preserves gaps and rejects duplicate observations. -use asap_types::{ - ir::operator_properties::Reduction, - ir::summary_coverage::*, - post_asap::SummaryUpdate, - pre_asap::{ColumnRef, Source}, +//! Summary coverage is `definition + selection`, derived from the sub-DAG +//! (`docs/design_docs/proposals/asap-primitive-schema.md` §4). Merges need +//! the same definition and disjoint selections. +use std::ops::Bound; +use std::rc::Rc; +use std::time::Duration; + +use asap_types::ir::operator_properties::{Reduction, Source}; +use asap_types::ir::summary_coverage::{ColumnIdentity, Constraint, CoverageError, SelectionBox}; +use asap_types::ir::{ + ASAPOp, ExprSemantics, NonASAPOp, Operator, OperatorNode, Predicate, ScalarExpr, + SchemaDerivationError, TimeRangeKind, +}; +use asap_types::post_asap::{ + ExecutionTiming, GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate, }; -fn table(name: &str) -> Source { - Source::Table { - table_ref: name.into(), +use asap_types::pre_asap::expr_ir::{ArithmeticOpKind, CompareOpKind, ScalarValue}; +use asap_types::pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema, TimeShift}; + +const JOB: usize = 0; +const REGION: usize = 1; +const LATENCY: usize = 2; + +fn table() -> Rc { + table_with(vec![ + Field::plain("job", DataType::Utf8, false), + Field::plain("region", DataType::Utf8, false), + Field::plain("latency", DataType::Float64, false), + Field::plain("size", DataType::Float64, false), + ]) +} + +fn table_with(fields: Vec) -> Rc { + OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: Source::Table { + table_ref: "t".into(), + }, + predicates: vec![], + schema: Schema::new(fields), + })) + .unwrap() +} + +/// PromQL leaf `m{predicates}` with columns `[ts, value, job]`. +fn series(predicates: Vec) -> Rc { + OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: Source::TimeSeries { metric: "m".into() }, + predicates, + schema: Schema::with_time_index( + vec![ + Field::plain("ts", DataType::Timestamp, false), + Field::plain("value", DataType::Float64, false), + Field::plain("job", DataType::Utf8, true), + ], + 0, + vec![], + ), + })) + .unwrap() +} + +fn compare(column: usize, op: CompareOpKind, value: ScalarValue) -> ScalarExpr { + ScalarExpr::Compare { + left: Box::new(ScalarExpr::Column(column)), + op, + right: Box::new(ScalarExpr::Literal(value)), + semantics: ExprSemantics::Sql, + } +} + +fn eq(column: usize, value: &str) -> ScalarExpr { + compare(column, CompareOpKind::Eq, ScalarValue::Utf8(value.into())) +} + +fn filter(child: Rc, pred: ScalarExpr) -> Rc { + OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Filter { + pred: Predicate(pred), + child, + })) + .unwrap() +} + +fn kll(child: Rc, column: &str) -> Rc { + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { + child, + family: FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), + Default::default(), + ), + input: SummaryUpdate::column(ColumnRef::Named(column.into())), + reduction: Reduction::by(vec![JOB]), + grouping: GroupingStrategy::default(), + filter: None, + })) + .unwrap() +} + +/// KLL of latency by job over `t` restricted by `pred`. +fn latency_where(pred: ScalarExpr) -> Rc { + kll(filter(table(), pred), "latency") +} + +/// Tumbling pane: `TimeRange(1m)` over `TimeShift(shift)` over `m`. +fn pane(shift_ms: i64, predicates: Vec) -> Rc { + let shifted = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeShift { + shift: TimeShift { + offset_ms: shift_ms, + at: None, + }, + child: series(predicates), + })) + .unwrap(); + let range = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange { + range: Duration::from_millis(60_000), + kind: TimeRangeKind::Range, + child: shifted, + })) + .unwrap(); + kll(range, "value") +} + +fn merge(children: Vec>) -> Result, SchemaDerivationError> { + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })) +} + +fn rejected(children: Vec>) -> CoverageError { + match merge(children) { + Err(SchemaDerivationError::Coverage(error)) => error, + other => panic!("expected a coverage error, got {other:?}"), } } -fn coverage(start: i64, end: i64, population: &[(&str, &str)]) -> SummaryCoverage { - SummaryCoverage { - source: table("flows"), - regions: vec![CoverageRegion { - time_ms: Some(start..end), - population: population - .iter() - .map(|(k, v)| (k.to_string(), v.to_string())) - .collect(), - }], + +fn column(name: &str) -> ColumnIdentity { + ColumnIdentity { + table: None, + name: name.into(), } } -/// Adjacent panes coalesce; gaps remain disconnected rather than becoming a hull. + +fn utf8(values: &[&str]) -> Constraint { + Constraint::In( + values + .iter() + .map(|v| ScalarValue::Utf8((*v).into())) + .collect(), + ) +} + +fn relative(from_ms: i64, to_ms: i64) -> Option<(Bound, Bound)> { + Some((Bound::Excluded(from_ms), Bound::Included(to_ms))) +} + +/// A lifted filter leaves the definition over the bare scan. #[test] -fn time_union_preserves_gaps() { - let merged = - SummaryCoverage::merge_disjoint(&[coverage(0, 1, &[]), coverage(1, 2, &[])]).unwrap(); - assert_eq!(merged.regions[0].time_ms, Some(0..2)); - assert_eq!(merged.regions.len(), 1); - let gapped = - SummaryCoverage::merge_disjoint(&[coverage(0, 1, &[]), coverage(2, 3, &[])]).unwrap(); - assert_eq!(gapped.regions.len(), 2); -} -/// Population partitions can overlap in time without sharing observations. +fn filters_lift_into_the_selection() { + let scan = table(); + let us = kll(filter(Rc::clone(&scan), eq(REGION, "us")), "latency"); + let coverage = us.coverage().unwrap(); + assert!(Rc::ptr_eq(coverage.definition.children()[0], &scan)); + assert_eq!( + coverage.selection, + vec![SelectionBox { + columns: [(column("region"), utf8(&["us"]))].into(), + relative_time: None, + }] + ); +} + +/// Different populations merge into one box; the same population twice +/// would count every row twice. #[test] -fn population_and_joint_union() { - let merged = SummaryCoverage::merge_disjoint(&[ - coverage(0, 2, &[("region", "us")]), - coverage(0, 2, &[("region", "eu")]), - ]) - .unwrap(); - assert_eq!(merged.regions.len(), 2); - let joint = SummaryCoverage::merge_disjoint(&[ - coverage(0, 1, &[("region", "us")]), - coverage(1, 2, &[("region", "eu")]), +fn populations_merge_when_disjoint() { + let merged = merge(vec![ + latency_where(eq(REGION, "us")), + latency_where(eq(REGION, "eu")), ]) .unwrap(); - assert_eq!(joint.regions.len(), 2); - let decoded: SummaryCoverage = - serde_json::from_str(&serde_json::to_string(&joint).unwrap()).unwrap(); - assert_eq!(decoded, joint); + assert_eq!( + merged.coverage().unwrap().selection[0].columns[&column("region")], + utf8(&["us", "eu"]) + ); + assert_eq!( + rejected(vec![ + latency_where(eq(REGION, "us")), + latency_where(eq(REGION, "us")), + ]), + CoverageError::PossibleOverlap + ); } -/// Intersecting predicates and windows cannot authorize once-per-observation merge. + +/// Equal schemas do not make a KLL over latency and one over size mergeable. #[test] -fn overlap_and_identity_fail_closed() { +fn different_inputs_do_not_merge() { assert_eq!( - SummaryCoverage::merge_disjoint(&[coverage(0, 2, &[]), coverage(1, 3, &[])]), - Err(CoverageError::PossibleOverlap) + rejected(vec![ + kll(filter(table(), eq(REGION, "us")), "latency"), + kll(filter(table(), eq(REGION, "eu")), "size"), + ]), + CoverageError::DefinitionMismatch ); +} + +/// The same update over scans of different columns is a different computation. +#[test] +fn the_same_update_over_different_scan_schemas_does_not_merge() { + let scan = |measure: &str| { + table_with(vec![ + Field::plain("job", DataType::Utf8, false), + Field::plain("region", DataType::Utf8, false), + Field::plain(measure, DataType::Float64, false), + Field::plain("value", DataType::Float64, false), + ]) + }; assert_eq!( - SummaryCoverage::merge_disjoint(&[ - coverage(0, 2, &[("region", "us")]), - coverage(0, 2, &[("tier", "premium")]) + rejected(vec![ + kll(filter(scan("latency"), eq(REGION, "us")), "value"), + kll(filter(scan("size"), eq(REGION, "eu")), "value"), ]), - Err(CoverageError::PossibleOverlap) + CoverageError::DefinitionMismatch ); +} + +/// `shipping.region` and `billing.region` are different columns, so +/// restricting each says nothing about overlap. +#[test] +fn qualified_columns_stay_distinct() { + let scan = || { + let qualified = |table: &str| Field { + table: Some(table.into()), + ..Field::plain("region", DataType::Utf8, false) + }; + table_with(vec![ + Field::plain("job", DataType::Utf8, false), + qualified("shipping"), + qualified("billing"), + Field::plain("latency", DataType::Float64, false), + ]) + }; assert_eq!( - SummaryCoverage::merge_disjoint(&[ - coverage(0, 2, &[("region", "us")]), - coverage(0, 2, &[("region", "us")]) + rejected(vec![ + kll(filter(scan(), eq(1, "us")), "latency"), + kll(filter(scan(), eq(2, "eu")), "latency"), ]), - Err(CoverageError::PossibleOverlap) + CoverageError::PossibleOverlap + ); +} + +/// Panes take their time from `TimeRange` over `TimeShift`: adjacent panes +/// join, a gap stays two boxes, and the same pane twice overlaps. +#[test] +fn time_panes_merge_and_keep_gaps() { + let first = pane(0, vec![]); + assert_eq!( + first.coverage().unwrap().selection[0].relative_time, + relative(-60_000, 0) + ); + assert!(Rc::ptr_eq( + first.coverage().unwrap().definition.children()[0], + first.children()[0].children()[0].children()[0] + )); + let joined = merge(vec![first.clone(), pane(60_000, vec![])]).unwrap(); + assert_eq!( + joined.coverage().unwrap().selection, + vec![SelectionBox { + columns: Default::default(), + relative_time: relative(-120_000, 0), + }] ); - let mut other = coverage(1, 2, &[]); - other.source = table("other-flows"); + let gapped = merge(vec![first.clone(), pane(120_000, vec![])]).unwrap(); + assert_eq!(gapped.coverage().unwrap().selection.len(), 2); assert_eq!( - SummaryCoverage::merge_disjoint(&[coverage(0, 1, &[]), other]), - Err(CoverageError::SourceMismatch) + rejected(vec![first.clone(), first]), + CoverageError::PossibleOverlap ); +} + +/// Label matchers on the scan lift like any filter. +#[test] +fn scan_predicates_lift() { + let api = pane(0, vec![Predicate(eq(2, "api"))]); + let coverage = api.coverage().unwrap(); assert_eq!( - coverage(2, 1, &[]).validate(), - Err(CoverageError::InvalidInterval) + coverage.selection[0].columns[&column("job")], + utf8(&["api"]) ); + let scan = coverage.definition.children()[0]; + assert!(matches!( + scan.non_asap(), + Some(NonASAPOp::Scan { predicates, .. }) if predicates.is_empty() + )); + merge(vec![api, pane(0, vec![Predicate(eq(2, "web"))])]).unwrap(); } -/// Coverage is logical state metadata, and input rewrites invalidate its proof. +/// Value ranges are selections: `latency < 100` and `latency >= 100` are +/// disjoint, `latency <= 100` overlaps `latency >= 100`. #[test] -fn node_coverage_is_required_checked_and_cleared_by_rewrites() { - use asap_types::{ - ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}, - post_asap::{SketchAlgorithm, SketchKind, SketchParams}, - pre_asap::{DataType, Field, FieldDataType, Schema}, +fn value_ranges_merge_when_disjoint() { + let latency = |op| latency_where(compare(LATENCY, op, ScalarValue::Float64(100.0))); + merge(vec![latency(CompareOpKind::Lt), latency(CompareOpKind::Ge)]).unwrap(); + assert_eq!( + rejected(vec![latency(CompareOpKind::Le), latency(CompareOpKind::Ge)]), + CoverageError::PossibleOverlap + ); +} + +/// A predicate that is not a value set or interval on one column stays in +/// the definition, so states with different residuals do not merge. +#[test] +fn residual_predicates_stay_in_the_definition() { + let doubled_above = |threshold: f64, region: &str| { + let doubled = ScalarExpr::Arithmetic { + op: ArithmeticOpKind::Mul, + left: Box::new(ScalarExpr::Column(LATENCY)), + right: Box::new(ScalarExpr::Literal(ScalarValue::Float64(2.0))), + semantics: ExprSemantics::Sql, + }; + latency_where(ScalarExpr::BoolAnd(vec![ + ScalarExpr::Compare { + left: Box::new(doubled), + op: CompareOpKind::Gt, + right: Box::new(ScalarExpr::Literal(ScalarValue::Float64(threshold))), + semantics: ExprSemantics::Sql, + }, + eq(REGION, region), + ])) }; - let raw = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { - source: table("flows"), - predicates: vec![], - schema: Schema::new(vec![Field::plain("latency", DataType::Float64, false)]), + let state = doubled_above(10.0, "us"); + let definition = &state.coverage().unwrap().definition; + assert!(matches!( + definition.children()[0].non_asap(), + Some(NonASAPOp::Filter { .. }) + )); + merge(vec![state.clone(), doubled_above(10.0, "eu")]).unwrap(); + assert_eq!( + rejected(vec![state, doubled_above(20.0, "eu")]), + CoverageError::DefinitionMismatch + ); +} + +/// An instant selector picks the latest sample per series, which is not a +/// selection of rows, so it stays in the definition. +#[test] +fn instant_selectors_stay_in_the_definition() { + let instant = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange { + range: Duration::from_millis(300_000), + kind: TimeRangeKind::Instant, + child: series(vec![]), })) .unwrap(); - let declared = coverage(0, 1, &[]); - assert!((*raw).clone().with_coverage(declared.clone()).is_err()); - let state = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { - child: raw, - family: FieldDataType::Sketch( - SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), - Default::default(), - ), - input: SummaryUpdate::column(ColumnRef::Named("latency".into())), - reduction: Reduction::by(vec![]), - grouping: Default::default(), - filter: None, - })) + let state = kll(Rc::clone(&instant), "value"); + let coverage = state.coverage().unwrap(); + assert_eq!(coverage.selection[0].relative_time, None); + assert!(Rc::ptr_eq(coverage.definition.children()[0], &instant)); +} + +/// A merge is a state with coverage, so merges nest. +#[test] +fn merges_nest() { + let us_eu = merge(vec![ + latency_where(eq(REGION, "us")), + latency_where(eq(REGION, "eu")), + ]) .unwrap(); - // Summary nodes cannot validate without coverage. + merge(vec![us_eu.clone(), latency_where(eq(REGION, "apac"))]).unwrap(); + assert_eq!( + rejected(vec![us_eu, latency_where(eq(REGION, "eu"))]), + CoverageError::PossibleOverlap + ); +} + +/// When a state is computed is not what it computes. +#[test] +fn timing_does_not_block_merges() { + let ingested = Rc::new( + (*pane(60_000, vec![])) + .clone() + .with_timing(Some(ExecutionTiming::IngestionTime)), + ); + let queried = Rc::new( + (*pane(0, vec![])) + .clone() + .with_timing(Some(ExecutionTiming::QueryTime)), + ); + merge(vec![ingested, queried]).unwrap(); +} + +/// A merge built around `new` is still checked by `validate_structure`. +#[test] +fn forged_merges_fail_validation() { + let us = latency_where(eq(REGION, "us")); + let schema = us.schema.clone(); + let forged = Rc::new(OperatorNode::with_schema( + Operator::ASAP(ASAPOp::SummaryMerge { + children: vec![us.clone(), us], + }), + schema, + )); + assert!(forged.coverage().is_none()); assert!(matches!( - std::rc::Rc::new(state.clone()).validate_structure(), - Err(asap_types::ir::SchemaDerivationError::Coverage( - CoverageError::Missing + forged.validate_structure(), + Err(SchemaDerivationError::Coverage( + CoverageError::PossibleOverlap )) )); - let state = state.with_coverage(declared.clone()).unwrap(); - std::rc::Rc::new(state.clone()) - .validate_structure() - .unwrap(); - let rebuilt = state.with_new_children(Clone::clone).unwrap(); - assert!(rebuilt.coverage.is_none()); } -/// Sources without a time column declare no time bounds; such a region overlaps -/// any region it is not population-disjoint from. +/// Only summary states have coverage. #[test] -fn regions_without_time_bounds() { - let mut tabular = coverage(0, 1, &[("region", "us")]); - tabular.regions[0].time_ms = None; - let mut other = coverage(0, 1, &[("region", "eu")]); - other.regions[0].time_ms = None; - assert_eq!( - SummaryCoverage::merge_disjoint(&[tabular.clone(), other]) - .unwrap() - .regions - .len(), - 2 - ); - assert_eq!( - SummaryCoverage::merge_disjoint(&[tabular, coverage(5, 6, &[("region", "us")])]), - Err(CoverageError::PossibleOverlap) - ); +fn non_summary_nodes_have_no_coverage() { + assert!(table().coverage().is_none()); + assert!(filter(table(), eq(REGION, "us")).coverage().is_none()); } diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index 7b7993ff4..b621d8403 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -1,55 +1,54 @@ //! Window composition merges compatible summary states without consuming raw rows. use asap_types::{ ir::operator_properties::{Reduction, Source}, - ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, SchemaDerivationError}, + ir::{ASAPOp, ExprSemantics, NonASAPOp, Operator, OperatorNode, Predicate, ScalarExpr}, post_asap::{GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate}, + pre_asap::expr_ir::{CompareOpKind, ScalarValue}, pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema}, }; use std::rc::Rc; -fn state(k: u32) -> Rc { - state_with(FieldDataType::Sketch( - SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }), - Default::default(), - )) -} - -fn state_with(family: FieldDataType) -> Rc { +/// KLL over `value` for one `region`, so states of different regions are +/// disjoint. +fn state(k: u32, region: &str) -> Rc { let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { source: Source::Table { table_ref: "latencies".into(), }, predicates: vec![], - schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]), + schema: Schema::new(vec![ + Field::plain("region", DataType::Utf8, false), + Field::plain("value", DataType::Float64, false), + ]), })) .unwrap(); - let summary = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { + let only_region = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Filter { + pred: Predicate(ScalarExpr::Compare { + left: Box::new(ScalarExpr::Column(0)), + op: CompareOpKind::Eq, + right: Box::new(ScalarExpr::Literal(ScalarValue::Utf8(region.into()))), + semantics: ExprSemantics::Sql, + }), child: scan, - family, + })) + .unwrap(); + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { + child: only_region, + family: FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }), + Default::default(), + ), input: SummaryUpdate::column(ColumnRef::SampleValue), reduction: Reduction::by(vec![]), grouping: GroupingStrategy::default(), filter: None, })) - .unwrap(); - std::rc::Rc::new( - summary - .with_coverage(asap_types::ir::summary_coverage::SummaryCoverage { - source: Source::Table { - table_ref: "latencies".into(), - }, - regions: vec![asap_types::ir::summary_coverage::CoverageRegion { - time_ms: Some(0..1), - population: Default::default(), - }], - }) - .unwrap(), - ) + .unwrap() } /// Two KLL panes compose into one typed logical state without timing assignment. #[test] fn compatible_panes_merge_structurally() { let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { - children: vec![state(200), shifted_state(200, 1, 2)], + children: vec![state(200, "us"), state(200, "eu")], })) .unwrap(); root.validate_structure().unwrap(); @@ -60,34 +59,11 @@ fn compatible_panes_merge_structurally() { fn incompatible_merge_inputs_fail() { for children in [ vec![], - vec![state(200), state(300)], - vec![state(200).children()[0].clone()], + vec![state(200, "us"), state(300, "eu")], + vec![state(200, "us").children()[0].clone()], ] { assert!( OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })).is_err() ); } } - -fn shifted_state(k: u32, start: i64, end: i64) -> Rc { - shifted(&state(k), start, end) -} - -fn shifted(state: &Rc, start: i64, end: i64) -> Rc { - let mut node = (**state).clone(); - let region = &mut node.coverage.as_mut().unwrap().regions[0]; - region.time_ms = Some(start..end); - Rc::new(node) -} - -fn merge(children: Vec>) -> Result, SchemaDerivationError> { - OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })) -} - -/// A merge is a state, so merges nest structurally. -#[test] -fn merges_nest() { - let inner = merge(vec![state(200)]).unwrap(); - let outer = merge(vec![inner, shifted_state(200, 1, 2)]).unwrap(); - outer.validate_structure().unwrap(); -} From 1bbd2a026408f1d67e819ed894a388ea6d0f8a87 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Wed, 7 Oct 2026 23:07:37 +0000 Subject: [PATCH 2/6] fix(ir): keep ambiguous columns and untyped-equal literals out of selections From an independent review of the derivation: - a column whose (table, name) is not unique in the agg child's schema cannot be named in a selection, so its conjuncts stay in the definition; - value sets compare literals by typed order, so 1 and 1.0 (or NaN) are never proven different; - a partly lifted Scan predicate is rebuilt with only its residual. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- crates/types/src/ir/summary_coverage.rs | 65 ++++++++++++++++++------ crates/types/tests/summary_coverage.rs | 67 ++++++++++++++++++++++++- 2 files changed, 116 insertions(+), 16 deletions(-) diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index da65128ee..31b000711 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -22,7 +22,7 @@ use super::node::{Operator, OperatorNode}; use super::non_asap::{NonASAPOp, TimeRangeKind}; use super::scalar::{Predicate, ScalarExpr}; use crate::pre_asap::expr_ir::{CompareOpKind, ScalarValue}; -use crate::pre_asap::schema::ColumnId; +use crate::pre_asap::schema::{ColumnId, Schema}; #[derive(Debug, Clone, PartialEq)] pub struct SummaryCoverage { @@ -173,17 +173,20 @@ fn of_summary_agg(node: &OperatorNode) -> SummaryCoverage { .count() == 1; + let top = &child.schema; // Top-down: decide which conjuncts are lifted. `kept[d]` is the residual // of `chain[d]` when it is a `Filter`. let mut lifted = BTreeMap::new(); let agg_filter = filter .as_ref() - .and_then(|pred| residual(pred, &[], &mut lifted)); + .and_then(|pred| residual(pred, top, &[], &mut lifted)); let kept: Vec> = chain .iter() .enumerate() .map(|(depth, link)| match link.non_asap() { - Some(NonASAPOp::Filter { pred, .. }) => residual(pred, &chain[..depth], &mut lifted), + Some(NonASAPOp::Filter { pred, .. }) => { + residual(pred, top, &chain[..depth], &mut lifted) + } _ => None, }) .collect(); @@ -191,7 +194,7 @@ fn of_summary_agg(node: &OperatorNode) -> SummaryCoverage { Some(NonASAPOp::Scan { predicates, .. }) => Some( predicates .iter() - .filter_map(|pred| residual(pred, &chain, &mut lifted)) + .filter_map(|pred| residual(pred, top, &chain, &mut lifted)) .collect::>(), ), _ => None, @@ -209,7 +212,7 @@ fn of_summary_agg(node: &OperatorNode) -> SummaryCoverage { }), ) = (scan_kept, base.non_asap()) { - if kept.len() != predicates.len() { + if kept != *predicates { rebuilt = rebuild( base, NonASAPOp::Scan { @@ -279,7 +282,6 @@ fn of_summary_agg(node: &OperatorNode) -> SummaryCoverage { }, ); - let top = &child.schema; let selection = SelectionBox { columns: lifted .into_iter() @@ -319,11 +321,14 @@ fn rebuild_asap(original: &OperatorNode, op: ASAPOp) -> Rc { )) } -/// Lift what `pred`'s conjuncts can into `lifted` (keyed by the agg child's -/// column) and return the rest. `above` are the chain nodes between the -/// predicate and the `SummaryAgg`, top to bottom. +/// Lift what `pred`'s conjuncts can into `lifted` (keyed by the column of +/// `top`, the agg child's schema) and return the rest. `above` are the chain +/// nodes between the predicate and the `SummaryAgg`, top to bottom. A column +/// whose `(table, name)` is not unique in `top` cannot be named in a +/// selection, so its conjuncts stay. fn residual( pred: &Predicate, + top: &Schema, above: &[&Rc], lifted: &mut BTreeMap, ) -> Option { @@ -331,6 +336,15 @@ fn residual( for conjunct in pred.0.conjuncts() { let lift = constraint_of(conjunct).and_then(|(column, constraint)| { let column = column_at_top(column, above)?; + let field = &top.fields[column]; + let namesakes = top + .fields + .iter() + .filter(|f| f.name == field.name && f.table == field.table) + .count(); + if namesakes != 1 { + return None; + } let combined = match lifted.get(&column) { Some(existing) => existing.intersect(&constraint)?, None => constraint, @@ -454,6 +468,29 @@ fn dedup(values: Vec) -> Vec { out } +/// Whether `v` equals one of `values`; `None` when some pair cannot be +/// compared (different types, NaN), since `1` and `1.0` may select the same +/// rows. +fn among(v: &ScalarValue, values: &[ScalarValue]) -> Option { + for w in values { + if compare(v, w)? == Ordering::Equal { + return Some(true); + } + } + Some(false) +} + +/// The values of `a` that are (`keep = true`) or are not in `b`. +fn filter_values(a: &[ScalarValue], b: &[ScalarValue], keep: bool) -> Option> { + let mut out = Vec::new(); + for v in a { + if among(v, b)? == keep { + out.push(v.clone()); + } + } + Some(out) +} + /// Order of two literals of the same type; `None` across types or for NaN. fn compare(a: &ScalarValue, b: &ScalarValue) -> Option { match (a, b) { @@ -535,10 +572,8 @@ impl Constraint { fn intersect(&self, other: &Self) -> Option { use Constraint::*; Some(match (self, other) { - (In(a), In(b)) => In(a.iter().filter(|v| b.contains(v)).cloned().collect()), - (In(a), NotIn(b)) | (NotIn(b), In(a)) => { - In(a.iter().filter(|v| !b.contains(v)).cloned().collect()) - } + (In(a), In(b)) => In(filter_values(a, b, true)?), + (In(a), NotIn(b)) | (NotIn(b), In(a)) => In(filter_values(a, b, false)?), (NotIn(a), NotIn(b)) => NotIn(dedup(a.iter().chain(b).cloned().collect())), (In(a), Interval { lower, upper }) | (Interval { lower, upper }, In(a)) => { let mut kept = Vec::new(); @@ -570,8 +605,8 @@ impl Constraint { fn disjoint(&self, other: &Self) -> bool { use Constraint::*; match (self, other) { - (In(a), In(b)) => !a.iter().any(|v| b.contains(v)), - (In(a), NotIn(b)) | (NotIn(b), In(a)) => a.iter().all(|v| b.contains(v)), + (In(a), In(b)) => a.iter().all(|v| among(v, b) == Some(false)), + (In(a), NotIn(b)) | (NotIn(b), In(a)) => a.iter().all(|v| among(v, b) == Some(true)), (In(a), Interval { lower, upper }) | (Interval { lower, upper }, In(a)) => a .iter() .all(|v| inside(v, (lower, upper), compare) == Some(false)), diff --git a/crates/types/tests/summary_coverage.rs b/crates/types/tests/summary_coverage.rs index 3ef5890f3..7ce8a280b 100644 --- a/crates/types/tests/summary_coverage.rs +++ b/crates/types/tests/summary_coverage.rs @@ -8,7 +8,7 @@ use std::time::Duration; use asap_types::ir::operator_properties::{Reduction, Source}; use asap_types::ir::summary_coverage::{ColumnIdentity, Constraint, CoverageError, SelectionBox}; use asap_types::ir::{ - ASAPOp, ExprSemantics, NonASAPOp, Operator, OperatorNode, Predicate, ScalarExpr, + ASAPOp, ExprSemantics, NonASAPOp, Operator, OperatorNode, Predicate, ProjectItem, ScalarExpr, SchemaDerivationError, TimeRangeKind, }; use asap_types::post_asap::{ @@ -410,3 +410,68 @@ fn non_summary_nodes_have_no_coverage() { assert!(table().coverage().is_none()); assert!(filter(table(), eq(REGION, "us")).coverage().is_none()); } + +/// Two output columns with the same `(table, name)` cannot be told apart in +/// a selection, so restrictions on them stay in the definition. +#[test] +fn ambiguous_column_names_do_not_lift() { + let renamed = |restricted: usize, value: &str| { + let item = |column: usize, alias: Option<&str>| ProjectItem { + alias: alias.map(Into::into), + expr: ScalarExpr::Column(column), + }; + let project = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Project { + cols: vec![ + item(JOB, None), + item(REGION, Some("k")), + item(LATENCY, Some("k")), + item(LATENCY, None), + ], + qualifier: None, + child: table(), + })) + .unwrap(); + kll(filter(project, eq(restricted, value)), "latency") + }; + assert_eq!( + rejected(vec![renamed(1, "us"), renamed(2, "eu")]), + CoverageError::DefinitionMismatch + ); +} + +/// `latency = 1` and `latency = 1.0` select the same rows of a Float64 column. +#[test] +fn literals_of_different_types_are_not_proven_different() { + let one = |value| latency_where(compare(LATENCY, CompareOpKind::Eq, value)); + assert!(merge(vec![ + one(ScalarValue::Int64(1)), + one(ScalarValue::Float64(1.0)) + ]) + .is_err()); +} + +/// A scan predicate whose conjuncts only partly lift keeps just the rest. +#[test] +fn partly_lifted_scan_predicates_keep_only_the_rest() { + let doubled_positive = ScalarExpr::Compare { + left: Box::new(ScalarExpr::Arithmetic { + op: ArithmeticOpKind::Mul, + left: Box::new(ScalarExpr::Column(1)), + right: Box::new(ScalarExpr::Literal(ScalarValue::Float64(2.0))), + semantics: ExprSemantics::Promql, + }), + op: CompareOpKind::Gt, + right: Box::new(ScalarExpr::Literal(ScalarValue::Float64(0.0))), + semantics: ExprSemantics::Promql, + }; + let job = |name: &str| { + pane( + 0, + vec![Predicate(ScalarExpr::BoolAnd(vec![ + eq(2, name), + doubled_positive.clone(), + ]))], + ) + }; + merge(vec![job("api"), job("web")]).unwrap(); +} From 1af815b0a3830a29aec1bcaf2b12543b402e6fd7 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Fri, 9 Oct 2026 21:23:27 +0000 Subject: [PATCH 3/6] test(ir): cover the design doc's worked coverage example MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Point the module docs at #573 §4.2 and note that a Project renames columns in their identity. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- crates/types/src/ir/summary_coverage.rs | 5 +- crates/types/tests/summary_coverage.rs | 109 +++++++++++++++++++++++- 2 files changed, 111 insertions(+), 3 deletions(-) diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index 31b000711..e3a90563c 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -1,6 +1,6 @@ //! What a summary state covers: `definition` (what it computes) and //! `selection` (which output rows of that computation it took). Design: -//! `docs/design_docs/proposals/asap-primitive-schema.md` §4, after +//! `docs/design_docs/proposals/asap-primitive-schema.md` §4.2, after //! Goldstein & Larson's view matching. //! //! Coverage is derived from the node, never declared. Walking down from a @@ -43,7 +43,8 @@ pub struct SelectionBox { pub relative_time: Option<(Bound, Bound)>, } -/// A column of the definition child's output, by `(table, name)`. +/// A column of the definition child's output, by `(table, name)`. A +/// `Project` below renames its columns and drops their table. #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] pub struct ColumnIdentity { pub table: Option, diff --git a/crates/types/tests/summary_coverage.rs b/crates/types/tests/summary_coverage.rs index 7ce8a280b..aca390d3b 100644 --- a/crates/types/tests/summary_coverage.rs +++ b/crates/types/tests/summary_coverage.rs @@ -1,5 +1,5 @@ //! Summary coverage is `definition + selection`, derived from the sub-DAG -//! (`docs/design_docs/proposals/asap-primitive-schema.md` §4). Merges need +//! (`docs/design_docs/proposals/asap-primitive-schema.md` §4.2). Merges need //! the same definition and disjoint selections. use std::ops::Bound; use std::rc::Rc; @@ -475,3 +475,110 @@ fn partly_lifted_scan_predicates_keep_only_the_rest() { }; merge(vec![job("api"), job("web")]).unwrap(); } + +/// The worked example of the design doc (§4.2.2): conditions on the +/// `SummaryAgg` filter, through a renaming `Project`, and the time window +/// move into the selection; the condition on an expression stays. +#[test] +fn design_doc_worked_example() { + let (value, job, region) = (1, 2, 3); + let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: Source::TimeSeries { metric: "m".into() }, + predicates: vec![], + schema: Schema::with_time_index( + vec![ + Field::plain("ts", DataType::Timestamp, false), + Field::plain("value", DataType::Float64, false), + Field::plain("job", DataType::Utf8, true), + Field::plain("region", DataType::Utf8, true), + ], + 0, + vec![], + ), + })) + .unwrap(); + let shifted = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeShift { + shift: TimeShift { + offset_ms: 120_000, + at: None, + }, + child: Rc::clone(&scan), + })) + .unwrap(); + let range = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange { + range: Duration::from_millis(60_000), + kind: TimeRangeKind::Range, + child: shifted, + })) + .unwrap(); + let doubled_above_10 = ScalarExpr::Compare { + left: Box::new(ScalarExpr::Arithmetic { + op: ArithmeticOpKind::Mul, + left: Box::new(ScalarExpr::Column(value)), + right: Box::new(ScalarExpr::Literal(ScalarValue::Float64(2.0))), + semantics: ExprSemantics::Sql, + }), + op: CompareOpKind::Gt, + right: Box::new(ScalarExpr::Literal(ScalarValue::Float64(10.0))), + semantics: ExprSemantics::Sql, + }; + let filtered = filter( + range, + ScalarExpr::BoolAnd(vec![eq(region, "us"), doubled_above_10.clone()]), + ); + let item = |column: usize, alias: Option<&str>| ProjectItem { + alias: alias.map(Into::into), + expr: ScalarExpr::Column(column), + }; + let project = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Project { + cols: vec![item(job, None), item(region, Some("r")), item(value, None)], + qualifier: None, + child: filtered, + })) + .unwrap(); + let below_100 = compare(2, CompareOpKind::Lt, ScalarValue::Float64(100.0)); + let state = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { + child: project, + family: FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), + Default::default(), + ), + input: SummaryUpdate::column(ColumnRef::Named("value".into())), + reduction: Reduction::by(vec![0]), + grouping: GroupingStrategy::default(), + filter: Some(Predicate(below_100)), + })) + .unwrap(); + + let coverage = state.coverage().unwrap(); + assert_eq!( + coverage.selection, + vec![SelectionBox { + columns: [ + (column("r"), utf8(&["us"])), + ( + column("value"), + Constraint::Interval { + lower: Bound::Unbounded, + upper: Bound::Excluded(ScalarValue::Float64(100.0)), + }, + ), + ] + .into(), + relative_time: relative(-180_000, -120_000), + }] + ); + // definition: SummaryAgg (no filter) over Project over + // Filter(value * 2 > 10) over the bare Scan. + let Some(ASAPOp::SummaryAgg { filter, child, .. }) = coverage.definition.asap() else { + panic!("definition is a SummaryAgg"); + }; + assert!(filter.is_none()); + assert!(matches!(child.non_asap(), Some(NonASAPOp::Project { .. }))); + let kept = child.children()[0]; + assert!(matches!( + kept.non_asap(), + Some(NonASAPOp::Filter { pred, .. }) if pred.0 == doubled_above_10 + )); + assert!(Rc::ptr_eq(kept.children()[0], &scan)); +} From 1926328315a5e7904c7de357526e236fb53a210d Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Fri, 9 Oct 2026 21:45:45 +0000 Subject: [PATCH 4/6] perf(ir): cache a merge's coverage when it is checked OperatorNode::new and validate_structure derived a SummaryMerge's coverage to check it and dropped the result, so coverage() derived it again. Store it in the cache instead. Also correct two comments (Project qualifier, absolute time needs timestamp literals). Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- crates/types/src/ir/node.rs | 10 ++++++++-- crates/types/src/ir/summary_coverage.rs | 6 ++++-- 2 files changed, 12 insertions(+), 4 deletions(-) diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index d58c362cd..6a7a227a5 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -184,10 +184,16 @@ impl OperatorNode { /// A `SummaryMerge` is valid only over inputs with the same definition /// and disjoint selections. + /// The derived coverage is cached, so a merge built by `new` is not + /// derived again by `coverage()` or `validate_structure`. fn check_merge(&self) -> Result<(), SchemaDerivationError> { - if matches!(self.asap(), Some(ASAPOp::SummaryMerge { .. })) { - SummaryCoverage::derive(self)?; + if !matches!(self.asap(), Some(ASAPOp::SummaryMerge { .. })) + || matches!(self.coverage_cache.0.get(), Some(Some(_))) + { + return Ok(()); } + let coverage = SummaryCoverage::derive(self)?; + let _ = self.coverage_cache.0.set(Some(coverage)); Ok(()) } diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index e3a90563c..043d846c5 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -39,12 +39,14 @@ pub struct SelectionBox { pub columns: BTreeMap, /// Offsets from the evaluation time in milliseconds: a PromQL range /// `TimeRange(w)` over `TimeShift(s)` is `(-(s + w), -s]`. Absolute - /// time is an ordinary interval on the timestamp column. + /// time will be an ordinary interval on the timestamp column once the IR + /// has timestamp literals. pub relative_time: Option<(Bound, Bound)>, } /// A column of the definition child's output, by `(table, name)`. A -/// `Project` below renames its columns and drops their table. +/// `Project` below renames its columns and sets their table to its +/// qualifier (none by default). #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] pub struct ColumnIdentity { pub table: Option, From 6f6f869b8d38a568ee0b0d7cfd5dce6f9d92ac91 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 10 Oct 2026 01:05:33 +0000 Subject: [PATCH 5/6] test(ir): worked example follows the SQL query of the design doc MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The #573 worked example (§4.2.2) is now a SQL query; build the DAG the SQL frontend lowers it to (WHERE folded into Scan.predicates, no time window). Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- crates/types/tests/summary_coverage.rs | 88 +++++++++++--------------- 1 file changed, 36 insertions(+), 52 deletions(-) diff --git a/crates/types/tests/summary_coverage.rs b/crates/types/tests/summary_coverage.rs index aca390d3b..28799f123 100644 --- a/crates/types/tests/summary_coverage.rs +++ b/crates/types/tests/summary_coverage.rs @@ -476,74 +476,60 @@ fn partly_lifted_scan_predicates_keep_only_the_rest() { merge(vec![job("api"), job("web")]).unwrap(); } -/// The worked example of the design doc (§4.2.2): conditions on the -/// `SummaryAgg` filter, through a renaming `Project`, and the time window -/// move into the selection; the condition on an expression stays. +/// The worked example of the design doc (§4.2.2), as the SQL frontend +/// lowers it: the `FILTER` condition and `region = 'us'` (in +/// `Scan.predicates`, through a renaming `Project`) move into the +/// selection; the condition on an expression stays. #[test] fn design_doc_worked_example() { - let (value, job, region) = (1, 2, 3); - let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { - source: Source::TimeSeries { metric: "m".into() }, - predicates: vec![], - schema: Schema::with_time_index( - vec![ - Field::plain("ts", DataType::Timestamp, false), - Field::plain("value", DataType::Float64, false), - Field::plain("job", DataType::Utf8, true), - Field::plain("region", DataType::Utf8, true), - ], - 0, - vec![], - ), - })) - .unwrap(); - let shifted = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeShift { - shift: TimeShift { - offset_ms: 120_000, - at: None, - }, - child: Rc::clone(&scan), - })) - .unwrap(); - let range = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange { - range: Duration::from_millis(60_000), - kind: TimeRangeKind::Range, - child: shifted, - })) - .unwrap(); let doubled_above_10 = ScalarExpr::Compare { left: Box::new(ScalarExpr::Arithmetic { op: ArithmeticOpKind::Mul, - left: Box::new(ScalarExpr::Column(value)), - right: Box::new(ScalarExpr::Literal(ScalarValue::Float64(2.0))), + left: Box::new(ScalarExpr::Column(LATENCY)), + right: Box::new(ScalarExpr::Literal(ScalarValue::Int64(2))), semantics: ExprSemantics::Sql, }), op: CompareOpKind::Gt, - right: Box::new(ScalarExpr::Literal(ScalarValue::Float64(10.0))), + right: Box::new(ScalarExpr::Literal(ScalarValue::Int64(10))), semantics: ExprSemantics::Sql, }; - let filtered = filter( - range, - ScalarExpr::BoolAnd(vec![eq(region, "us"), doubled_above_10.clone()]), - ); + let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: Source::Table { + table_ref: "t".into(), + }, + predicates: vec![Predicate(ScalarExpr::BoolAnd(vec![ + eq(REGION, "us"), + doubled_above_10.clone(), + ]))], + schema: Schema::new(vec![ + Field::plain("job", DataType::Utf8, false), + Field::plain("region", DataType::Utf8, false), + Field::plain("latency", DataType::Float64, false), + ]), + })) + .unwrap(); let item = |column: usize, alias: Option<&str>| ProjectItem { alias: alias.map(Into::into), expr: ScalarExpr::Column(column), }; let project = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Project { - cols: vec![item(job, None), item(region, Some("r")), item(value, None)], + cols: vec![ + item(JOB, None), + item(REGION, Some("r")), + item(LATENCY, None), + ], qualifier: None, - child: filtered, + child: scan, })) .unwrap(); - let below_100 = compare(2, CompareOpKind::Lt, ScalarValue::Float64(100.0)); + let below_100 = compare(2, CompareOpKind::Lt, ScalarValue::Int64(100)); let state = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { child: project, family: FieldDataType::Sketch( SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), Default::default(), ), - input: SummaryUpdate::column(ColumnRef::Named("value".into())), + input: SummaryUpdate::column(ColumnRef::Named("latency".into())), reduction: Reduction::by(vec![0]), grouping: GroupingStrategy::default(), filter: Some(Predicate(below_100)), @@ -557,28 +543,26 @@ fn design_doc_worked_example() { columns: [ (column("r"), utf8(&["us"])), ( - column("value"), + column("latency"), Constraint::Interval { lower: Bound::Unbounded, - upper: Bound::Excluded(ScalarValue::Float64(100.0)), + upper: Bound::Excluded(ScalarValue::Int64(100)), }, ), ] .into(), - relative_time: relative(-180_000, -120_000), + relative_time: None, }] ); // definition: SummaryAgg (no filter) over Project over - // Filter(value * 2 > 10) over the bare Scan. + // Scan t with only `latency * 2 > 10` left in its predicates. let Some(ASAPOp::SummaryAgg { filter, child, .. }) = coverage.definition.asap() else { panic!("definition is a SummaryAgg"); }; assert!(filter.is_none()); assert!(matches!(child.non_asap(), Some(NonASAPOp::Project { .. }))); - let kept = child.children()[0]; assert!(matches!( - kept.non_asap(), - Some(NonASAPOp::Filter { pred, .. }) if pred.0 == doubled_above_10 + child.children()[0].non_asap(), + Some(NonASAPOp::Scan { predicates, .. }) if *predicates == vec![Predicate(doubled_above_10)] )); - assert!(Rc::ptr_eq(kept.children()[0], &scan)); } From e548a0b675fcd4d8ec8cd58e2a43cc521f460d4d Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 10 Oct 2026 02:00:30 +0000 Subject: [PATCH 6/6] feat(ir): join touching value ranges in a merged selection Value ranges on one column now join like touching time windows, so every selection dimension follows Goldstein & Larson's one-range-per-column form. The joined range stays explicit: neither input took NULL rows. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- crates/types/src/ir/summary_coverage.rs | 71 ++++++++++++++++--------- crates/types/tests/summary_coverage.rs | 32 +++++++++++ 2 files changed, 78 insertions(+), 25 deletions(-) diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index 043d846c5..1bbb2978a 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -645,8 +645,9 @@ impl SelectionBox { } } -/// Join boxes that differ only in adjacent time windows or only in the -/// values of one `In` column; gaps stay separate boxes. +/// Join boxes that differ only in one dimension whose union is again one +/// constraint: touching time windows or value ranges, or the values of one +/// `In` column. Gaps stay separate boxes. fn union(mut boxes: Vec) -> Vec { let mut i = 0; while i < boxes.len() { @@ -666,20 +667,7 @@ fn union(mut boxes: Vec) -> Vec { fn join(a: &SelectionBox, b: &SelectionBox) -> Option { if a.columns == b.columns { let ((al, au), (bl, bu)) = (a.relative_time.as_ref()?, b.relative_time.as_ref()?); - let meets = |upper: &Bound, lower: &Bound| { - matches!( - (upper, lower), - (Bound::Included(u), Bound::Excluded(l)) | (Bound::Excluded(u), Bound::Included(l)) - if u == l - ) - }; - let time = if meets(au, bl) { - (*al, *bu) - } else if meets(bu, al) { - (*bl, *au) - } else { - return None; - }; + let time = touching((al, au), (bl, bu), |x: &i64, y: &i64| Some(x.cmp(y)))?; return Some(SelectionBox { columns: a.columns.clone(), relative_time: Some(time), @@ -692,22 +680,55 @@ fn join(a: &SelectionBox, b: &SelectionBox) -> Option { .columns .iter() .filter(|(column, constraint)| b.columns.get(*column) != Some(*constraint)); - let (column, Constraint::In(values)) = differing.next()? else { - return None; - }; + let (column, constraint) = differing.next()?; if differing.next().is_some() { return None; } - let Some(Constraint::In(more)) = b.columns.get(column) else { - return None; + let joined = match (constraint, b.columns.get(column)?) { + (Constraint::In(values), Constraint::In(more)) => { + Constraint::In(dedup(values.iter().chain(more).cloned().collect())) + } + ( + Constraint::Interval { + lower: al, + upper: au, + }, + Constraint::Interval { + lower: bl, + upper: bu, + }, + ) => { + let (lower, upper) = touching((al, au), (bl, bu), compare)?; + Constraint::Interval { lower, upper } + } + _ => return None, }; let mut columns = a.columns.clone(); - columns.insert( - column.clone(), - Constraint::In(dedup(values.iter().chain(more).cloned().collect())), - ); + columns.insert(column.clone(), joined); Some(SelectionBox { columns, relative_time: a.relative_time, }) } + +/// The union of two intervals when one ends exactly where the other starts, +/// with the shared end point in exactly one of them. +fn touching( + (al, au): (&Bound, &Bound), + (bl, bu): (&Bound, &Bound), + cmp: impl Fn(&T, &T) -> Option, +) -> Option<(Bound, Bound)> { + let meets = |upper: &Bound, lower: &Bound| match (upper, lower) { + (Bound::Included(u), Bound::Excluded(l)) | (Bound::Excluded(u), Bound::Included(l)) => { + cmp(u, l) == Some(Ordering::Equal) + } + _ => false, + }; + if meets(au, bl) { + Some((al.clone(), bu.clone())) + } else if meets(bu, al) { + Some((bl.clone(), au.clone())) + } else { + None + } +} diff --git a/crates/types/tests/summary_coverage.rs b/crates/types/tests/summary_coverage.rs index 28799f123..97be84b56 100644 --- a/crates/types/tests/summary_coverage.rs +++ b/crates/types/tests/summary_coverage.rs @@ -303,6 +303,38 @@ fn value_ranges_merge_when_disjoint() { ); } +/// Touching value ranges join like touching time windows; the joined range +/// stays explicit, since neither input took NULL rows. A gap stays apart. +#[test] +fn touching_value_ranges_join() { + let latency = |op, v| latency_where(compare(LATENCY, op, ScalarValue::Float64(v))); + let joined = merge(vec![ + latency(CompareOpKind::Lt, 100.0), + latency(CompareOpKind::Ge, 100.0), + ]) + .unwrap(); + assert_eq!( + joined.coverage().unwrap().selection, + vec![SelectionBox { + columns: [( + column("latency"), + Constraint::Interval { + lower: Bound::Unbounded, + upper: Bound::Unbounded, + }, + )] + .into(), + relative_time: None, + }] + ); + let gap = merge(vec![ + latency(CompareOpKind::Lt, 100.0), + latency(CompareOpKind::Gt, 100.0), + ]) + .unwrap(); + assert_eq!(gap.coverage().unwrap().selection.len(), 2); +} + /// A predicate that is not a value set or interval on one column stays in /// the definition, so states with different residuals do not merge. #[test]