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 @@ -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 {
Expand Down
114 changes: 105 additions & 9 deletions crates/types/src/ir/properties/summary_coverage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Range<i64>>,
pub time_ms: Option<CoverageTime>,
/// Conjunction of non-null equality predicates; empty means unrestricted.
pub population: BTreeMap<String, String>,
}

/// 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<i64>),
/// Offsets from the evaluation time. Tumbling pane `i` of width `w`
/// covers `-(i + 1) * w..-i * w`.
RelativeToEvaluation(Range<i64>),
}

impl CoverageTime {
pub fn range(&self) -> &Range<i64> {
match self {
Self::Absolute(range) | Self::RelativeToEvaluation(range) => range,
}
}

fn is_relative(&self) -> bool {
matches!(self, Self::RelativeToEvaluation(_))
}
}

impl From<Range<i64>> for CoverageTime {
fn from(range: Range<i64>) -> 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<i64>),
Relative(RelativeTimeWire),
}

#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct RelativeTimeWire {
relative_to_evaluation: Range<i64>,
}

impl From<CoverageTimeWire> 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<CoverageTime> 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")]
Expand All @@ -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")]
Expand All @@ -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) {
Expand Down Expand Up @@ -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<CoverageRegion> = 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, &region.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, &region.time_ms)
{
if last.population == region.population && last_time.end == time.range().start {
last_time.end = time.range().end;
continue;
}
}
Expand All @@ -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
Expand Down
18 changes: 18 additions & 0 deletions crates/types/src/ir/schema/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
2 changes: 1 addition & 1 deletion crates/types/tests/logical_export.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
},
],
Expand Down
2 changes: 1 addition & 1 deletion crates/types/tests/physical_export.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ fn plan() -> Rc<OperatorNode> {
.with_coverage(SummaryCoverage {
source,
regions: vec![CoverageRegion {
time_ms: Some(0..60_000),
time_ms: Some((0..60_000).into()),
population: Default::default(),
}],
})
Expand Down
82 changes: 80 additions & 2 deletions crates/types/tests/summary_coverage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()))
Expand All @@ -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();
Expand Down Expand Up @@ -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::<SummaryCoverage>(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::<SummaryCoverage>(legacy).unwrap(),
absolute
);
}
2 changes: 1 addition & 1 deletion crates/types/tests/summary_coverage_examples.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ fn region_is(value: &str) -> Predicate {

fn region(time_ms: Option<std::ops::Range<i64>>, 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()))
Expand Down
Loading