diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index f5d91043..b3788499 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -1239,18 +1239,18 @@ impl PhysicalCompiler { query_id: query.query_id.clone(), reason: "derived source config missing".into(), })?; - if runtime_materialization.window_size != source_config.window_size - || runtime_materialization.slide_interval != source_config.slide_interval - || runtime_materialization.pane_origin_ms != source_config.pane_origin_ms - || runtime_materialization.window_size - != runtime_materialization.slide_interval - || runtime_materialization.stored_window_ms() - != source_config.stored_window_ms() - || source_config.stored_window_ms() - != source_config.window_size.saturating_mul(1_000) + asap_types::precompute_plan::validated_source_window_cohort( + &runtime_materialization, + &[source_config], + ) + .map_err(|error| CompileError::Query { + query_id: query.query_id.clone(), + reason: error.to_string(), + })?; + if runtime_materialization.window_size != runtime_materialization.slide_interval { return Err(CompileError::Query { query_id: query.query_id.clone(), - reason: "immutable scalar composition requires matching full nonoverlapping windows".into() }); + reason: "immutable scalar composition runtime requires nonoverlapping windows".into() }); } let SummaryExpr::SummaryAgg { child, .. } = &selected.node.expr else { unreachable!() diff --git a/crates/asap_types/src/precompute_plan.rs b/crates/asap_types/src/precompute_plan.rs index 543fae8d..818195a3 100644 --- a/crates/asap_types/src/precompute_plan.rs +++ b/crates/asap_types/src/precompute_plan.rs @@ -9,6 +9,55 @@ use serde::{Deserialize, Serialize}; use std::collections::{BTreeMap, BTreeSet, HashMap}; use thiserror::Error; +/// Validate the physical windows consumed by one derived input program. +/// This returns the existing source definitions, not a second serialized +/// contract. It does not prove completion or authorize runtime execution. +/// Full-window sliding is lossless here; pane merging requires a separate +/// explicit operator and is deliberately not inferred from window sizes. +pub fn validated_source_window_cohort<'a>( + target: &crate::PrecomputeMaterialization, + sources: &[&'a crate::PrecomputeMaterialization], +) -> Result, PrecomputePlanError> { + let invalid = || { + PrecomputePlanError::CatalogContract( + "derived inputs require matching explicit full stored windows".into(), + ) + }; + let full_ms = target.window_size.checked_mul(1000).ok_or_else(invalid)?; + if sources.is_empty() + || full_ms == 0 + || target.slide_interval == 0 + || target.slide_interval > target.window_size + || (target.slide_interval < target.window_size && target.pane_origin_ms.is_none()) + || target.stored_window_ms() != full_ms + || (target.slide_interval < target.window_size + && !matches!( + target.window_layout, + crate::WindowMaterializationLayout::FullWindow + )) + { + return Err(invalid()); + } + let mut identities = BTreeSet::new(); + for source in sources { + if source.derived_input.is_some() + || !identities.insert(source.policy_fingerprint()) + || source.window_size != target.window_size + || source.slide_interval != target.slide_interval + || source.pane_origin_ms != target.pane_origin_ms + || source.stored_window_ms() != full_ms + || (source.slide_interval < source.window_size + && !matches!( + source.window_layout, + crate::WindowMaterializationLayout::FullWindow + )) + { + return Err(invalid()); + } + } + Ok(sources.to_vec()) +} + pub const BACKEND_COMPAT: &str = "asap-query-backend.v1"; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] @@ -381,14 +430,11 @@ impl PrecomputePlan { .iter() .find(|candidate| candidate.policy_fingerprint() == source_id.fingerprint()) .ok_or_else(invalid)?; - if source.derived_input.is_some() - || source.window_size != config.window_size - || source.slide_interval != config.slide_interval - || source.window_size != source.slide_interval - || source.pane_origin_ms != config.pane_origin_ms - || source.stored_window_ms() != config.stored_window_ms() - || source.window_size.checked_mul(1000) != Some(source.stored_window_ms()) - { + validated_source_window_cohort(config, &[source])?; + // Current installed runtime capability remains nonoverlapping. + // The shared cohort contract also describes explicit full-window + // sliding for consumers which separately prove its completion. + if source.window_size != source.slide_interval { return Err(invalid()); } let mut matched = false; @@ -700,3 +746,54 @@ pub(crate) fn state_encodings(family: &SummaryFamilyType) -> Vec _ => Vec::new(), } } + +#[cfg(test)] +mod source_window_cohort_tests { + use super::*; + fn full_window() -> crate::PrecomputeMaterialization { + serde_json::from_value(serde_json::json!({ + "aggregation_type":"Sum", "aggregation_sub_type":"", "parameters":{}, + "grouping_labels":{"labels":[]}, "aggregated_labels":{"labels":[]}, + "rollup_labels":{"labels":[]}, "original_yaml":"", + "window_size":60, "slide_interval":10, "window_type":"sliding", + "window_layout":{"kind":"full_window"}, "pane_origin_ms":0, + "spatial_filter":"", "spatial_filter_normalized":"", "metric":"m", + "num_aggregates_to_retain":null, "table_name":null, "value_projection":null + })) + .unwrap() + } + #[test] + fn full_sliding_cohort_preserves_explicit_windows_and_identity() { + let target = full_window(); + let source = full_window(); + let mut missing_origin = source.clone(); + missing_origin.pane_origin_ms = None; + assert!(validated_source_window_cohort(&missing_origin, &[&missing_origin]).is_err()); + missing_origin.slide_interval = missing_origin.window_size; + assert!(validated_source_window_cohort(&missing_origin, &[&missing_origin]).is_ok()); + let result = validated_source_window_cohort(&target, &[&source]).unwrap(); + assert!(std::ptr::eq(result[0], &source)); + let mut other = source.clone(); + other.metric = "other".into(); + assert_eq!( + validated_source_window_cohort(&target, &[&source, &other]) + .unwrap() + .len(), + 2 + ); + assert!(validated_source_window_cohort(&target, &[&source, &source]).is_err()); + for mutation in 0..3 { + let mut changed = source.clone(); + match mutation { + 0 => changed.pane_origin_ms = Some(1), + 1 => changed.slide_interval = 20, + _ => { + changed.window_layout = + crate::WindowMaterializationLayout::Pane { pane_secs: 10 } + } + } + assert_ne!(source.policy_fingerprint(), changed.policy_fingerprint()); + assert!(validated_source_window_cohort(&target, &[&changed]).is_err()); + } + } +}