Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions crates/types/src/ir/operator/asap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -481,6 +481,11 @@ impl ASAPOp {
"summary merge requires exactly one state column".into(),
));
}
if let Some(family) = self.produced_state().filter(|f| !f.family_merges()) {
return Err(SchemaDerivationError::InvalidScalarSignature(format!(
"summary merge over {family:?} is unsupported: the family has no sound merge"
)));
}
for child in children {
needs_state(child, "SummaryMerge")?;
if child.schema != first.schema {
Expand Down
32 changes: 22 additions & 10 deletions crates/types/src/ir/properties/summary_coverage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ use crate::ir::operator::node::{Operator, OperatorNode};
use crate::ir::operator::non_asap::{NonASAPOp, TimeRangeKind};
use crate::ir::scalar::{CompareOpKind, ScalarValue};
use crate::ir::scalar::{Predicate, ScalarExpr};
use crate::ir::schema::{ColumnId, Schema};
use crate::ir::schema::{ColumnId, Schema, SelectionRelation};

#[derive(Debug, Clone, PartialEq)]
pub struct SummaryCoverage {
Expand Down Expand Up @@ -68,14 +68,17 @@ pub enum CoverageError {
EmptyMerge,
#[error("summary merge inputs compute different things")]
DefinitionMismatch,
#[error("summary merge inputs are not proven disjoint")]
#[error(
"summary merge inputs are not proven disjoint and the family counts shared rows twice"
)]
PossibleOverlap,
}

impl SummaryCoverage {
/// Coverage of a `SummaryAgg` or `SummaryMerge`. A merge fails unless
/// every input has the same definition and their selections are
/// pairwise disjoint.
/// every input has the same definition and their selections relate as
/// the family requires: pairwise disjoint, unless the family's merge is
/// idempotent.
pub fn derive(node: &OperatorNode) -> Result<Self, CoverageError> {
match node.asap() {
Some(ASAPOp::SummaryAgg { .. }) => Ok(of_summary_agg(node)),
Expand All @@ -85,17 +88,23 @@ impl SummaryCoverage {
.map(|child| child.coverage().ok_or(CoverageError::NotSummary))
.collect::<Result<Vec<_>, _>>()?;
let first = inputs.first().ok_or(CoverageError::EmptyMerge)?;
let may_overlap = matches!(
first.definition.asap(),
Some(ASAPOp::SummaryAgg { family, .. })
if family.merge_relation() == Some(SelectionRelation::OverlapAllowed)
);
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)))
});
let overlaps = !may_overlap
&& inputs[..i].iter().any(|other| {
input
.selection
.iter()
.any(|a| other.selection.iter().any(|b| !a.disjoint(b)))
});
if overlaps {
return Err(CoverageError::PossibleOverlap);
}
Expand Down Expand Up @@ -661,6 +670,9 @@ fn union(mut boxes: Vec<SelectionBox>) -> Vec<SelectionBox> {
}

fn join(a: &SelectionBox, b: &SelectionBox) -> Option<SelectionBox> {
if a == b {
return Some(a.clone());
}
if a.columns == b.columns {
let ((al, au), (bl, bu)) = (a.relative_time.as_ref()?, b.relative_time.as_ref()?);
let meets = |upper: &Bound<i64>, lower: &Bound<i64>| {
Expand Down
50 changes: 50 additions & 0 deletions crates/types/src/ir/schema/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,19 @@ impl From<Field<DataType>> for Field<FieldDataType> {
}
}

/// How the selections of a summary operation's inputs must relate for the
/// result to keep the family's guarantee (#573 §4.4).
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SelectionRelation {
/// No row in two inputs: a shared row would be counted twice.
Disjoint,
/// Inputs may share rows: the merge is idempotent.
OverlapAllowed,
/// The right input's selection lies inside the left's, as subtraction
/// (`SummarySubtract`, reserved) needs.
Contained,
}

/// What a schema field carries: an ordinary readable value, or the summary /
/// exact-accumulator state produced by a `SummaryAgg`.
///
Expand Down Expand Up @@ -159,6 +172,43 @@ impl FieldDataType {
matches!(self, FieldDataType::Plain(_))
}

/// Whether two states of this family over disjoint coverage merge into the
/// state of their union with the family's guarantee intact. Rate/Increase
/// accumulators depend on window edges, and merged heap top-k states have
/// no accuracy model yet; families not listed fail closed.
pub fn family_merges(&self) -> bool {
use state_type::{ExactKind as E, SketchAlgorithm as S};
match self {
FieldDataType::ExactAggregate(kind, _) => {
matches!(kind, E::Sum | E::Count | E::Min | E::Max)
}
FieldDataType::Sketch(kind, _) => matches!(
kind.algorithm(),
S::Kll | S::DDSketch | S::Hll | S::Cms | S::CountSketch | S::UnivMon
),
_ => false,
}
}

/// How the selections of merged states of this family must relate
/// (#573 §4.4); `None` when the family does not merge.
pub fn merge_relation(&self) -> Option<SelectionRelation> {
use state_type::{ExactKind as E, SketchAlgorithm as S};
if !self.family_merges() {
return None;
}
let idempotent = match self {
FieldDataType::ExactAggregate(kind, _) => matches!(kind, E::Min | E::Max),
FieldDataType::Sketch(kind, _) => *kind.algorithm() == S::Hll,
_ => false,
};
Some(if idempotent {
SelectionRelation::OverlapAllowed
} else {
SelectionRelation::Disjoint
})
}

pub fn plain(&self) -> Option<&DataType> {
match self {
FieldDataType::Plain(dtype) => Some(dtype),
Expand Down
203 changes: 199 additions & 4 deletions crates/types/tests/summary_merge_structure.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,15 @@ use std::rc::Rc;
/// KLL over `value` for one `region`, so states of different regions are
/// disjoint.
fn state(k: u32, region: &str) -> Rc<OperatorNode> {
family_state(
FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }),
Default::default(),
),
region,
)
}
fn family_state(family: FieldDataType, region: &str) -> Rc<OperatorNode> {
let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan {
source: Source::Table {
table_ref: "latencies".into(),
Expand All @@ -36,10 +45,7 @@ fn state(k: u32, region: &str) -> Rc<OperatorNode> {
.unwrap();
OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg {
child: only_region,
family: FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }),
Default::default(),
),
family,
input: SummaryUpdate::column(ColumnRef::SampleValue),
reduction: Reduction::by(vec![]),
grouping: GroupingStrategy::default(),
Expand Down Expand Up @@ -70,3 +76,192 @@ fn incompatible_merge_inputs_fail() {
);
}
}

fn merge_regions(family: FieldDataType) -> Result<Rc<OperatorNode>, impl std::fmt::Debug> {
OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge {
children: vec![
family_state(family.clone(), "us"),
family_state(family, "eu"),
],
}))
}
/// Heap top-k states have no sound merge model, so disjoint states still cannot
/// merge; KLL states with the same definition can.
#[test]
fn summary_merge_requires_a_mergeable_family() {
let heap = FieldDataType::Sketch(
SketchKind::new(
SketchAlgorithm::CmsWithHeap,
SketchParams::CmsWithHeap {
width: 64,
depth: 4,
heap_size: 10,
},
),
Default::default(),
);
let error = format!("{:?}", merge_regions(heap).unwrap_err());
assert!(error.contains("no sound merge"), "{error}");
let kll = FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }),
Default::default(),
);
merge_regions(kll).unwrap().validate_structure().unwrap();
}
/// Sound-merge capability is a closed list per family (W6).
#[test]
fn family_merge_capability() {
use asap_types::ir::schema::{ExactKind, ExactParams};
let exact = |kind, params| FieldDataType::ExactAggregate(kind, params);
for family in [
exact(ExactKind::Sum, ExactParams::Sum),
exact(ExactKind::Count, ExactParams::Count),
exact(ExactKind::Min, ExactParams::Min),
exact(ExactKind::Max, ExactParams::Max),
] {
assert!(family.family_merges(), "{family:?}");
}
for family in [
exact(ExactKind::Rate, ExactParams::Rate),
exact(ExactKind::Increase, ExactParams::Increase),
exact(ExactKind::IRate, ExactParams::IRate),
FieldDataType::Plain(DataType::Float64),
] {
assert!(!family.family_merges(), "{family:?}");
}
let sketch = |algorithm, params| {
FieldDataType::Sketch(SketchKind::new(algorithm, params), Default::default())
};
use SketchAlgorithm as A;
use SketchParams as P;
let (width, depth, heap_size) = (64, 4, 10);
for (family, merges) in [
(sketch(A::Kll, P::Kll { k: 200 }), true),
(sketch(A::DDSketch, P::DDSketch { alpha: 0.01 }), true),
(sketch(A::Hll, P::Hll { precision: 12 }), true),
(sketch(A::Cms, P::Cms { width, depth }), true),
(
sketch(A::CountSketch, P::CountSketch { width, depth }),
true,
),
(
sketch(
A::UnivMon,
P::UnivMon {
heap_size,
sketch_rows: depth,
sketch_cols: width,
layers: 8,
},
),
true,
),
(
sketch(
A::CmsWithHeap,
P::CmsWithHeap {
width,
depth,
heap_size,
},
),
false,
),
(
sketch(
A::CountSketchWithHeap,
P::CountSketchWithHeap {
width,
depth,
heap_size,
},
),
false,
),
] {
assert_eq!(family.family_merges(), merges, "{family:?}");
}
}
/// Idempotent families (HLL, exact Min/Max) merge states that may share rows;
/// counting families (KLL, exact Sum) still need disjoint selections.
#[test]
fn idempotent_families_merge_overlapping_states() {
use asap_types::ir::schema::{ExactKind, ExactParams};
let hll = FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::Hll, SketchParams::Hll { precision: 12 }),
Default::default(),
);
let max = FieldDataType::ExactAggregate(ExactKind::Max, ExactParams::Max);
let sum = FieldDataType::ExactAggregate(ExactKind::Sum, ExactParams::Sum);
let kll = FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }),
Default::default(),
);
let same_region = |family: FieldDataType| {
OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge {
children: vec![
family_state(family.clone(), "us"),
family_state(family, "us"),
],
}))
};
for family in [hll, max] {
let merged = same_region(family.clone()).unwrap_or_else(|e| panic!("{family:?}: {e}"));
merged.validate_structure().unwrap();
assert_eq!(
merged.coverage().unwrap().selection,
family_state(family, "us").coverage().unwrap().selection
);
}
for family in [sum, kll] {
assert!(same_region(family).is_err());
}
}
/// The selection relation a merge requires is declared per family.
#[test]
fn family_merge_selection_relation() {
use asap_types::ir::schema::{ExactKind, ExactParams, SelectionRelation};
let exact = |kind, params| FieldDataType::ExactAggregate(kind, params);
let sketch = |algorithm, params| {
FieldDataType::Sketch(SketchKind::new(algorithm, params), Default::default())
};
for (family, relation) in [
(
exact(ExactKind::Sum, ExactParams::Sum),
Some(SelectionRelation::Disjoint),
),
(
exact(ExactKind::Count, ExactParams::Count),
Some(SelectionRelation::Disjoint),
),
(
exact(ExactKind::Min, ExactParams::Min),
Some(SelectionRelation::OverlapAllowed),
),
(
exact(ExactKind::Max, ExactParams::Max),
Some(SelectionRelation::OverlapAllowed),
),
(exact(ExactKind::Rate, ExactParams::Rate), None),
(
sketch(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }),
Some(SelectionRelation::Disjoint),
),
(
sketch(
SketchAlgorithm::Cms,
SketchParams::Cms {
width: 64,
depth: 4,
},
),
Some(SelectionRelation::Disjoint),
),
(
sketch(SketchAlgorithm::Hll, SketchParams::Hll { precision: 12 }),
Some(SelectionRelation::OverlapAllowed),
),
] {
assert_eq!(family.merge_relation(), relation, "{family:?}");
}
}