From e2da91a720be7883b489605ba1d41c8b410eb22b Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 17:05:17 +0000 Subject: [PATCH 1/2] feat(types): evaluation-relative coverage time CoverageRegion.time_ms is now Option: Absolute or RelativeToEvaluation (tumbling pane i of width w covers -(i+1)w..-iw). The two anchors are incomparable, so validate/merge_disjoint reject a mix with MixedTimeAnchors; within one anchor, adjacent regions coalesce and gaps are kept as before. Absolute time keeps the plain {start,end} JSON encoding, so legacy coverage still deserializes. Co-Authored-By: Claude Opus 5.5 --- .../src/ir/properties/summary_coverage.rs | 114 ++++++++++++++++-- crates/types/tests/logical_export.rs | 2 +- crates/types/tests/physical_export.rs | 2 +- crates/types/tests/summary_coverage.rs | 82 ++++++++++++- .../types/tests/summary_coverage_examples.rs | 2 +- crates/types/tests/summary_merge_structure.rs | 8 +- 6 files changed, 192 insertions(+), 18 deletions(-) diff --git a/crates/types/src/ir/properties/summary_coverage.rs b/crates/types/src/ir/properties/summary_coverage.rs index 19ae10eb..a2494653 100644 --- a/crates/types/src/ir/properties/summary_coverage.rs +++ b/crates/types/src/ir/properties/summary_coverage.rs @@ -22,11 +22,79 @@ pub struct SummaryCoverage { pub struct CoverageRegion { /// Half-open bounds on the source's time column, in milliseconds. `None` /// means no time restriction, e.g. a source without a time column. - pub time_ms: Option>, + pub time_ms: Option, /// Conjunction of non-null equality predicates; empty means unrestricted. pub population: BTreeMap, } +/// Half-open time bounds in milliseconds and what they are measured from. +/// Absolute and evaluation-relative bounds are incomparable: without an +/// evaluation time, no overlap between them can be ruled out. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(from = "CoverageTimeWire", into = "CoverageTimeWire")] +pub enum CoverageTime { + /// Timestamps on the source's time column. + Absolute(Range), + /// Offsets from the evaluation time. Tumbling pane `i` of width `w` + /// covers `-(i + 1) * w..-i * w`. + RelativeToEvaluation(Range), +} + +impl CoverageTime { + pub fn range(&self) -> &Range { + match self { + Self::Absolute(range) | Self::RelativeToEvaluation(range) => range, + } + } + + fn is_relative(&self) -> bool { + matches!(self, Self::RelativeToEvaluation(_)) + } +} + +impl From> for CoverageTime { + fn from(range: Range) -> Self { + Self::Absolute(range) + } +} + +/// Absolute time keeps the plain `{start, end}` encoding, so coverage written +/// before relative time existed still deserializes. +#[derive(Serialize, Deserialize)] +#[serde(untagged)] +enum CoverageTimeWire { + Absolute(Range), + Relative(RelativeTimeWire), +} + +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct RelativeTimeWire { + relative_to_evaluation: Range, +} + +impl From for CoverageTime { + fn from(wire: CoverageTimeWire) -> Self { + match wire { + CoverageTimeWire::Absolute(range) => Self::Absolute(range), + CoverageTimeWire::Relative(relative) => { + Self::RelativeToEvaluation(relative.relative_to_evaluation) + } + } + } +} + +impl From for CoverageTimeWire { + fn from(time: CoverageTime) -> Self { + match time { + CoverageTime::Absolute(range) => Self::Absolute(range), + CoverageTime::RelativeToEvaluation(range) => Self::Relative(RelativeTimeWire { + relative_to_evaluation: range, + }), + } + } +} + #[derive(Debug, Clone, PartialEq, Eq, Error)] pub enum CoverageError { #[error("coverage interval must have start < end")] @@ -35,6 +103,10 @@ pub enum CoverageError { InvalidPopulation, #[error("summary coverage sources differ")] SourceMismatch, + #[error( + "coverage mixes absolute and evaluation-relative time, so disjointness cannot be proven" + )] + MixedTimeAnchors, #[error("coverage overlap is not proven absent")] PossibleOverlap, #[error("coverage merge requires at least one input")] @@ -51,8 +123,21 @@ pub enum CoverageError { impl SummaryCoverage { pub fn validate(&self) -> Result<(), CoverageError> { + let mut anchors = self + .regions + .iter() + .filter_map(|region| region.time_ms.as_ref().map(CoverageTime::is_relative)); + if let Some(first) = anchors.next() { + if anchors.any(|relative| relative != first) { + return Err(CoverageError::MixedTimeAnchors); + } + } for (index, region) in self.regions.iter().enumerate() { - if region.time_ms.as_ref().is_some_and(Range::is_empty) { + if region + .time_ms + .as_ref() + .is_some_and(|time| time.range().is_empty()) + { return Err(CoverageError::InvalidInterval); } if region.population.keys().any(String::is_empty) { @@ -84,21 +169,29 @@ impl SummaryCoverage { merged.regions.extend(input.regions.iter().cloned()); } merged.validate()?; - // Coalesce adjacent intervals only for identical population predicates. + // Coalesce adjacent intervals only for identical population predicates; + // `validate` has already ensured every bounded region shares one anchor. 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)), + .map(|t| t.range().start) + .cmp(&b.time_ms.as_ref().map(|t| t.range().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; + if let ( + Some( + CoverageTime::Absolute(last_time) + | CoverageTime::RelativeToEvaluation(last_time), + ), + Some(time), + ) = (&mut last.time_ms, ®ion.time_ms) + { + if last.population == region.population && last_time.end == time.range().start { + last_time.end = time.range().end; continue; } } @@ -112,7 +205,10 @@ impl SummaryCoverage { 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, + (Some(a), Some(b)) if a.is_relative() == b.is_relative() => { + let (a, b) = (a.range(), b.range()); + a.start < b.end && b.start < a.end + } _ => true, }; time_overlaps diff --git a/crates/types/tests/logical_export.rs b/crates/types/tests/logical_export.rs index f490733a..9c80d4f1 100644 --- a/crates/types/tests/logical_export.rs +++ b/crates/types/tests/logical_export.rs @@ -105,7 +105,7 @@ fn merged_summary_preserves_typed_state() { }, regions: vec![ asap_types::ir::properties::summary_coverage::CoverageRegion { - time_ms: Some(start..start + 1), + time_ms: Some((start..start + 1).into()), population: Default::default(), }, ], diff --git a/crates/types/tests/physical_export.rs b/crates/types/tests/physical_export.rs index a3397e2f..42325ce0 100644 --- a/crates/types/tests/physical_export.rs +++ b/crates/types/tests/physical_export.rs @@ -40,7 +40,7 @@ fn plan() -> Rc { .with_coverage(SummaryCoverage { source, regions: vec![CoverageRegion { - time_ms: Some(0..60_000), + time_ms: Some((0..60_000).into()), population: Default::default(), }], }) diff --git a/crates/types/tests/summary_coverage.rs b/crates/types/tests/summary_coverage.rs index b1a8149d..bb1e2e37 100644 --- a/crates/types/tests/summary_coverage.rs +++ b/crates/types/tests/summary_coverage.rs @@ -13,7 +13,7 @@ fn coverage(start: i64, end: i64, population: &[(&str, &str)]) -> SummaryCoverag SummaryCoverage { source: table("flows"), regions: vec![CoverageRegion { - time_ms: Some(start..end), + time_ms: Some((start..end).into()), population: population .iter() .map(|(k, v)| (k.to_string(), v.to_string())) @@ -26,7 +26,7 @@ fn coverage(start: i64, end: i64, population: &[(&str, &str)]) -> SummaryCoverag 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[0].time_ms, Some((0..2).into())); assert_eq!(merged.regions.len(), 1); let gapped = SummaryCoverage::merge_disjoint(&[coverage(0, 1, &[]), coverage(2, 3, &[])]).unwrap(); @@ -146,3 +146,81 @@ fn regions_without_time_bounds() { Err(CoverageError::PossibleOverlap) ); } + +fn relative(start: i64, end: i64) -> SummaryCoverage { + let mut coverage = coverage(0, 1, &[]); + coverage.regions[0].time_ms = Some(CoverageTime::RelativeToEvaluation(start..end)); + coverage +} +const MINUTE: i64 = 60_000; +/// Pane i of width w covers `[-(i+1)w, -iw)`; five adjacent 1 m panes +/// coalesce into the 5 m window ending at evaluation time. +#[test] +fn adjacent_relative_panes_coalesce() { + let panes: Vec<_> = (0..5) + .map(|i| relative(-(i + 1) * MINUTE, -i * MINUTE)) + .collect(); + let merged = SummaryCoverage::merge_disjoint(&panes).unwrap(); + assert_eq!( + merged.regions[0].time_ms, + Some(CoverageTime::RelativeToEvaluation(-5 * MINUTE..0)) + ); + assert_eq!(merged.regions.len(), 1); +} +/// A missing relative pane leaves a gap instead of a hull. +#[test] +fn relative_gap_is_kept() { + let merged = SummaryCoverage::merge_disjoint(&[ + relative(-3 * MINUTE, -2 * MINUTE), + relative(-MINUTE, 0), + ]) + .unwrap(); + assert_eq!(merged.regions.len(), 2); +} +/// Overlapping relative panes would count observations twice. +#[test] +fn overlapping_relative_regions_are_rejected() { + assert_eq!( + SummaryCoverage::merge_disjoint(&[relative(-2 * MINUTE, 0), relative(-MINUTE, 0)]), + Err(CoverageError::PossibleOverlap) + ); +} +/// Relative and absolute time are incomparable without an evaluation time, so +/// even numerically disjoint ranges cannot be proven disjoint. +#[test] +fn mixing_relative_and_absolute_time_is_rejected() { + assert_eq!( + SummaryCoverage::merge_disjoint(&[relative(-MINUTE, 0), coverage(0, MINUTE, &[])]), + Err(CoverageError::MixedTimeAnchors) + ); + let mut mixed = relative(-MINUTE, 0); + mixed.regions.extend(coverage(0, MINUTE, &[]).regions); + assert_eq!(mixed.validate(), Err(CoverageError::MixedTimeAnchors)); +} +/// Relative time round-trips through serde, and a legacy plain `time_ms` +/// range still deserializes as absolute. +#[test] +fn time_anchor_serde_round_trip_and_legacy_json() { + let pane = relative(-MINUTE, 0); + let json = serde_json::to_value(&pane).unwrap(); + assert_eq!( + json["regions"][0]["time_ms"], + serde_json::json!({"relative_to_evaluation": {"start": -MINUTE, "end": 0}}) + ); + assert_eq!( + serde_json::from_value::(json).unwrap(), + pane + ); + let absolute = coverage(0, MINUTE, &[]); + let json = serde_json::to_value(&absolute).unwrap(); + assert_eq!( + json["regions"][0]["time_ms"], + serde_json::json!({"start": 0, "end": MINUTE}) + ); + let legacy = r#"{"source":{"Table":{"table_ref":"flows"}}, + "regions":[{"time_ms":{"start":0,"end":60000},"population":{}}]}"#; + assert_eq!( + serde_json::from_str::(legacy).unwrap(), + absolute + ); +} diff --git a/crates/types/tests/summary_coverage_examples.rs b/crates/types/tests/summary_coverage_examples.rs index 5e2a7c05..b5a81841 100644 --- a/crates/types/tests/summary_coverage_examples.rs +++ b/crates/types/tests/summary_coverage_examples.rs @@ -51,7 +51,7 @@ fn region_is(value: &str) -> Predicate { fn region(time_ms: Option>, population: &[(&str, &str)]) -> CoverageRegion { CoverageRegion { - time_ms, + time_ms: time_ms.map(Into::into), population: population .iter() .map(|(k, v)| (k.to_string(), v.to_string())) diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index 27bfcf90..2647d182 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -37,7 +37,7 @@ fn state(k: u32) -> Rc { }, regions: vec![ asap_types::ir::properties::summary_coverage::CoverageRegion { - time_ms: Some(0..1), + time_ms: Some((0..1).into()), population: Default::default(), }, ], @@ -57,7 +57,7 @@ fn compatible_panes_merge_structurally() { assert_eq!(root.schema.fields.len(), 1); assert_eq!( root.coverage.as_ref().unwrap().regions[0].time_ms, - Some(0..2) + Some((0..2).into()) ); } /// An empty merge, raw rows and differently sized state cannot masquerade as compatible panes. @@ -77,7 +77,7 @@ fn incompatible_merge_inputs_fail() { fn shifted_state(k: u32, start: i64, end: i64) -> Rc { let mut node = (*state(k)).clone(); let region = &mut node.coverage.as_mut().unwrap().regions[0]; - region.time_ms = Some(start..end); + region.time_ms = Some((start..end).into()); Rc::new(node) } /// Schema equality cannot authorize overlapping or unknown observation coverage. @@ -104,6 +104,6 @@ fn merge_derives_coverage_and_validates_retained_metadata() { .unwrap(); assert_eq!(root.coverage.as_ref().unwrap().regions.len(), 2); let mut forged = (*root).clone(); - forged.coverage.as_mut().unwrap().regions[0].time_ms = Some(0..2); + forged.coverage.as_mut().unwrap().regions[0].time_ms = Some((0..2).into()); assert!(Rc::new(forged).validate_structure().is_err()); } From 6b31d2dc45ffd90a03228dd9a0a75b8876022c0f Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 17:05:18 +0000 Subject: [PATCH 2/2] feat(types): reject SummaryMerge over families without a sound merge FieldDataType::family_merges() is true for exact Sum/Count/Min/Max, KLL, DDSketch, HLL, CMS, CountSketch and UnivMon; every other family fails closed (Rate/Increase accumulators and heap top-k sketches included, per W6). SummaryMerge::validate_inputs rejects the rest. Co-Authored-By: Claude Opus 5.5 --- crates/types/src/ir/operator/asap.rs | 5 + crates/types/src/ir/schema/mod.rs | 18 +++ crates/types/tests/summary_merge_structure.rs | 122 +++++++++++++++++- 3 files changed, 140 insertions(+), 5 deletions(-) diff --git a/crates/types/src/ir/operator/asap.rs b/crates/types/src/ir/operator/asap.rs index d1baee5e..9cec3992 100644 --- a/crates/types/src/ir/operator/asap.rs +++ b/crates/types/src/ir/operator/asap.rs @@ -486,6 +486,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 { diff --git a/crates/types/src/ir/schema/mod.rs b/crates/types/src/ir/schema/mod.rs index 3ea89792..67dbb8ec 100644 --- a/crates/types/src/ir/schema/mod.rs +++ b/crates/types/src/ir/schema/mod.rs @@ -159,6 +159,24 @@ 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, + } + } + pub fn plain(&self) -> Option<&DataType> { match self { FieldDataType::Plain(dtype) => Some(dtype), diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index 2647d182..827fda0e 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -8,6 +8,15 @@ use asap_types::ir::schema::{ use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}; use std::rc::Rc; fn state(k: u32) -> Rc { + family_state( + FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }), + Default::default(), + ), + 0..1, + ) +} +fn family_state(family: FieldDataType, time: std::ops::Range) -> Rc { let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { source: Source::Table { table_ref: "latencies".into(), @@ -18,10 +27,7 @@ fn state(k: u32) -> Rc { .unwrap(); let summary = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { child: scan, - 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(), @@ -37,7 +43,7 @@ fn state(k: u32) -> Rc { }, regions: vec![ asap_types::ir::properties::summary_coverage::CoverageRegion { - time_ms: Some((0..1).into()), + time_ms: Some(time.into()), population: Default::default(), }, ], @@ -107,3 +113,109 @@ fn merge_derives_coverage_and_validates_retained_metadata() { forged.coverage.as_mut().unwrap().regions[0].time_ms = Some((0..2).into()); assert!(Rc::new(forged).validate_structure().is_err()); } + +fn merge_panes(family: FieldDataType) -> Result, impl std::fmt::Debug> { + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { + children: vec![ + family_state(family.clone(), 0..1), + family_state(family, 1..2), + ], + })) +} +/// Heap top-k states have no sound merge model, so disjoint panes still cannot +/// merge; KLL panes over the same coverage 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_panes(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_panes(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:?}"); + } +}