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..6a7a227a5 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,28 @@ 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. + /// 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 { .. })) + || matches!(self.coverage_cache.0.get(), Some(Some(_))) + { + return Ok(()); + } + let coverage = SummaryCoverage::derive(self)?; + let _ = self.coverage_cache.0.set(Some(coverage)); + Ok(()) } pub fn non_asap(&self) -> Option<&NonASAPOp> { @@ -315,14 +348,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..043d846c5 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -1,127 +1,713 @@ -//! 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.2, 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, Schema}; + +#[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 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 sets their table to its +/// qualifier (none by default). +#[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); + /// 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 region.population.keys().any(String::is_empty) { - return Err(CoverageError::InvalidPopulation); + _ => 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; + + 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, top, &[], &mut lifted)); + let kept: Vec> = chain + .iter() + .enumerate() + .map(|(depth, link)| match link.non_asap() { + Some(NonASAPOp::Filter { pred, .. }) => { + residual(pred, top, &chain[..depth], &mut lifted) } - if self.regions[..index] + _ => 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, top, &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 != *predicates { + 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 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 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 { + 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 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, + }; + Some((column, combined)) + }); + match lift { + Some((column, constraint)) => { + lifted.insert(column, constraint); } + None => rest.push(conjunct.clone()), + } + } + 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)) } - 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); + 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 +} + +/// 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) { + (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(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(); + 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().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)), + ( + 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..aca390d3b 100644 --- a/crates/types/tests/summary_coverage.rs +++ b/crates/types/tests/summary_coverage.rs @@ -1,150 +1,584 @@ -//! 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.2). 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, ProjectItem, 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); +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, + }] + ); } -/// Population partitions can overlap in time without sharing observations. + +/// 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 ); - let mut other = coverage(1, 2, &[]); - other.source = table("other-flows"); +} + +/// 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!( - SummaryCoverage::merge_disjoint(&[coverage(0, 1, &[]), other]), - Err(CoverageError::SourceMismatch) + 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!( - coverage(2, 1, &[]).validate(), - Err(CoverageError::InvalidInterval) + joined.coverage().unwrap().selection, + vec![SelectionBox { + columns: Default::default(), + relative_time: relative(-120_000, 0), + }] + ); + let gapped = merge(vec![first.clone(), pane(120_000, vec![])]).unwrap(); + assert_eq!(gapped.coverage().unwrap().selection.len(), 2); + assert_eq!( + rejected(vec![first.clone(), first]), + CoverageError::PossibleOverlap ); } -/// Coverage is logical state metadata, and input rewrites invalidate its proof. +/// Label matchers on the scan lift like any filter. #[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 scan_predicates_lift() { + let api = pane(0, vec![Predicate(eq(2, "api"))]); + let coverage = api.coverage().unwrap(); + assert_eq!( + 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(); +} + +/// Value ranges are selections: `latency < 100` and `latency >= 100` are +/// disjoint, `latency <= 100` overlaps `latency >= 100`. +#[test] +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; +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!( - SummaryCoverage::merge_disjoint(&[tabular.clone(), other]) - .unwrap() - .regions - .len(), - 2 + 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(); +} + +/// 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!( - SummaryCoverage::merge_disjoint(&[tabular, coverage(5, 6, &[("region", "us")])]), - Err(CoverageError::PossibleOverlap) + 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)); } 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(); -}