diff --git a/control_plane/src/opamp/mod.rs b/control_plane/src/opamp/mod.rs index 47817887..8143d6cf 100644 --- a/control_plane/src/opamp/mod.rs +++ b/control_plane/src/opamp/mod.rs @@ -1027,9 +1027,9 @@ mod tests { abstract_window_framework: planner_types::post_asap::SummaryWindowFramework::Tumbling, window_implementation_id: "collector-tumbling-v1".into(), - pane_secs: 60, + slide_secs: 60, pane_origin_ms: Some(0), - state_layout: "anchored-pane-v1".into(), + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 60 }, evidence_source: None, lifecycle: crate::physical::compiler::CollectorLifecycle { kind: "continuously_maintained".into(), diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 421f8201..d3402b79 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -90,12 +90,17 @@ pub struct ImplementationCostEvidence { pub observed_at_unix_ms: u64, pub valid_for_ms: u64, pub horizon_seconds: f64, + /// Measured build plus incremental-update CPU over `horizon_seconds`. pub cpu_cost: f64, pub peak_memory_bytes: u64, pub network_bytes: u64, pub storage_bytes: u64, pub source_scan_bytes: u64, - /// Dimensionally calibrated scalar passed to Planner for comparison. + /// Dimensionally calibrated scalar passed to Planner for comparison. The + /// evidence producer must price update CPU, query-time merges at the + /// workload's read frequency, retained memory, storage, source scans, and + /// network traffic. Keeping the components alongside this quote makes the + /// selected tradeoff auditable without teaching Planner backend units. pub weighted_cost: f64, } @@ -106,8 +111,8 @@ pub struct WindowImplementationCandidate { pub implementation_id: String, pub framework: SummaryWindowFramework, pub window_secs: u64, - pub pane_secs: u64, - pub state_layout: String, + pub slide_secs: u64, + pub layout: asap_types::WindowMaterializationLayout, pub cost: ImplementationCostEvidence, } @@ -218,7 +223,6 @@ pub struct BackendLocalImplementation { pub evidence_valid_for_ms: u64, pub horizon_seconds: f64, pub window_implementation_id: String, - pub state_layout: String, pub implementation_cost: ImplementationCostEvidence, #[serde(default, skip_serializing_if = "Option::is_none")] pub source_sample_interval_ms: Option, @@ -271,14 +275,14 @@ pub struct CollectorMaterialization { pub window_secs: u64, pub abstract_window_framework: SummaryWindowFramework, pub window_implementation_id: String, - pub pane_secs: u64, + pub slide_secs: u64, #[serde( default, alias = "paneOriginMs", skip_serializing_if = "Option::is_none" )] pub pane_origin_ms: Option, - pub state_layout: String, + pub window_layout: asap_types::WindowMaterializationLayout, pub evidence_source: Option, pub lifecycle: CollectorLifecycle, } @@ -467,6 +471,11 @@ pub enum PrecomputePlanError { InvalidSchema { schema_id: String }, #[error("materialization {0} uses a summary family unsupported by the runtime schema")] UnsupportedFamily(u64), + #[error("materialization {materialization} has an invalid window layout: {reason}")] + InvalidWindowLayout { + materialization: u64, + reason: String, + }, #[error("producer {producer_id} references an unknown materialization or schema")] InvalidProducer { producer_id: String }, #[error("materialization {0} has no registered producer")] @@ -622,6 +631,24 @@ impl PrecomputePlan { } let mut materializations = BTreeSet::new(); for materialization in &self.materializations { + materialization + .window_layout + .validate(materialization.window_size, materialization.slide_interval) + .map_err(|reason| PrecomputePlanError::InvalidWindowLayout { + materialization: materialization.policy_fp_u64(), + reason, + })?; + let expected_kind = if materialization.slide_interval == materialization.window_size { + asap_types::WindowKind::Tumbling + } else { + asap_types::WindowKind::Sliding + }; + if materialization.window_type != expected_kind { + return Err(PrecomputePlanError::InvalidWindowLayout { + materialization: materialization.policy_fp_u64(), + reason: "window kind disagrees with size and slide".into(), + }); + } // HLL is supported as an ingested sketch envelope, not as a raw // accumulator. Validate here so external installs cannot bypass it. if self.ingest.protocol == IngestProtocol::PrometheusRemoteWriteV1 @@ -1874,8 +1901,10 @@ impl BackendLocalPlanningSnapshot { implementation_id: self.implementation.window_implementation_id.clone(), framework: SummaryWindowFramework::Tumbling, window_secs: lookback_ms / 1_000, - pane_secs: lookback_ms / 1_000, - state_layout: self.implementation.state_layout.clone(), + slide_secs: lookback_ms / 1_000, + layout: asap_types::WindowMaterializationLayout::Pane { + pane_secs: lookback_ms / 1_000, + }, cost, }] }), @@ -2343,33 +2372,27 @@ impl PhysicalCompiler { })?; // The logical consumer group is identified above. The installed state // identity includes the actual pane width selected by Planner. - aggregation.window_secs = if environment.target - == PhysicalDeploymentTarget::BackendLocalRemoteWrite - && matches!( - &selected.family, - SummaryFamilyType::ExactAggregate( - planner_types::post_asap::ExactKind::Increase - | planner_types::post_asap::ExactKind::Rate - | planner_types::post_asap::ExactKind::MinMax, - _ - ) - ) { - // Exact temporal readouts may only merge panes wholly - // contained by the requested interval. Use the runtime's - // smallest supported pane and let the data plane anchor it - // to the first event-time sample; common scrape/repetition - // cadences then preserve exact boundaries without raw data. - exact_temporal_pane_secs( - window_implementation.pane_secs, - query.lifecycle.evaluation_interval_ms, - request.source_sample_interval_ms, - ) - } else { - window_implementation.pane_secs - }; + // Preserve semantic window and evaluation cadence independently + // from the selected storage representation. + aggregation.window_secs = window_implementation.window_secs; let mut runtime_materialization = aggregation_config_for_materialization(&aggregation)?; - let pane_width_ms = runtime_materialization.slide_interval.saturating_mul(1_000); + runtime_materialization.window_size = window_implementation.window_secs; + runtime_materialization.slide_interval = window_implementation.slide_secs; + runtime_materialization.window_type = + if window_implementation.slide_secs == window_implementation.window_secs { + asap_types::WindowKind::Tumbling + } else { + asap_types::WindowKind::Sliding + }; + runtime_materialization.window_layout = window_implementation.layout.clone(); + let pane_width_ms = match &window_implementation.layout { + asap_types::WindowMaterializationLayout::FullWindow => { + window_implementation.slide_secs + } + layout => layout.base_pane_secs(), + } + .saturating_mul(1_000); runtime_materialization.pane_origin_ms = shared_pane_origin_ms( request.query_workload.as_ref(), consumers[&materialization].iter().copied(), @@ -2450,9 +2473,9 @@ impl PhysicalCompiler { window_secs: runtime_materialization.window_size, abstract_window_framework: planner_selection.window_framework.clone(), window_implementation_id: window_implementation.implementation_id.clone(), - pane_secs: runtime_materialization.slide_interval, + slide_secs: runtime_materialization.slide_interval, pane_origin_ms: runtime_materialization.pane_origin_ms, - state_layout: window_implementation.state_layout.clone(), + window_layout: window_implementation.layout.clone(), evidence_source: evidence.map(|e| e.source.clone()), lifecycle: planner_selection.lifecycle.clone(), }; @@ -2605,11 +2628,17 @@ impl PhysicalCompiler { fingerprint.0 )) })?; - let pane_ms = precompute_plan + let stored_interval_ms = precompute_plan .materializations .iter() .find(|candidate| candidate.policy_fingerprint() == fingerprint) - .map(|candidate| candidate.slide_interval.saturating_mul(1_000)) + .map(|candidate| match &candidate.window_layout { + asap_types::WindowMaterializationLayout::FullWindow => { + candidate.window_size + } + layout => layout.base_pane_secs(), + } + .saturating_mul(1_000)) .ok_or_else(|| { crate::query_plan::QueryPlanError::Invalid(format!( "compiled binding {} has no precompute materialization", @@ -2639,7 +2668,7 @@ impl PhysicalCompiler { materialization.grouping_labels.labels.clone(), ), item_labels: materialization.aggregated_labels.labels.clone(), - window_ms: pane_ms, + window_ms: stored_interval_ms, pane_origin_ms: materialization.pane_origin_ms, }) }; @@ -2770,11 +2799,11 @@ impl PhysicalCompiler { .filter_map(|binding| binding.readout_lookback_ms) .max(); if let Some(lookback_ms) = max_lookback_ms { - let pane_ms = materialization.slide_interval.saturating_mul(1_000).max(1); - materialization.num_aggregates_to_retain = Some(retained_window_count( + materialization.num_aggregates_to_retain = Some(retained_state_count( lookback_ms, request.query_staleness_margin_ms, - pane_ms, + materialization.slide_interval.saturating_mul(1_000), + &materialization.window_layout, )); } } @@ -3191,7 +3220,6 @@ fn validate_window_implementations( let valid = evidence.observed_at_unix_ms <= environment.observed_at_unix_ms && !candidate.implementation_id.trim().is_empty() && ids.insert(candidate.implementation_id.clone()) - && !candidate.state_layout.trim().is_empty() && !evidence.model_version.trim().is_empty() && !evidence.workload_fingerprint.trim().is_empty() && evidence.valid_for_ms != 0 @@ -3203,16 +3231,32 @@ fn validate_window_implementations( && evidence.weighted_cost.is_finite() && evidence.weighted_cost >= 0.0 && candidate.window_secs == query.window_secs - && candidate.pane_secs != 0 - && candidate.pane_secs <= candidate.window_secs - && candidate.window_secs % candidate.pane_secs == 0 - && candidate.framework == SummaryWindowFramework::Tumbling + && candidate + .layout + .validate(candidate.window_secs, candidate.slide_secs) + .is_ok() + && match (&candidate.framework, &candidate.layout) { + ( + SummaryWindowFramework::Tumbling | SummaryWindowFramework::Sliding, + asap_types::WindowMaterializationLayout::Pane { .. }, + ) + | ( + SummaryWindowFramework::Sliding, + asap_types::WindowMaterializationLayout::FullWindow, + ) => true, + ( + SummaryWindowFramework::Extension(name), + asap_types::WindowMaterializationLayout::HierarchicalRollup { .. }, + ) => name == "backend.exact-hierarchical-rollup.v1", + _ => false, + } && match environment.target { PhysicalDeploymentTarget::DistributedCollectors => { - candidate.pane_secs == candidate.window_secs + matches!( + candidate.layout, + asap_types::WindowMaterializationLayout::FullWindow + ) || matches!(candidate.layout, asap_types::WindowMaterializationLayout::Pane { pane_secs } if pane_secs == candidate.window_secs) } - // Cadence does not establish phase alignment. Serving checks each - // actual interval and falls back when whole panes cannot cover it. PhysicalDeploymentTarget::BackendLocalRemoteWrite => true, }; if !valid { @@ -3239,37 +3283,27 @@ fn validate_window_implementations( Ok(candidates) } -fn exact_temporal_pane_secs( - candidate_pane_secs: u64, - evaluation_interval_ms: u32, - source_sample_interval_ms: Option, +fn retained_state_count( + lookback_ms: u64, + staleness_margin_ms: u64, + slide_ms: u64, + layout: &asap_types::WindowMaterializationLayout, ) -> u64 { - fn gcd(mut left: u64, mut right: u64) -> u64 { - while right != 0 { - (left, right) = (right, left % right); + match layout { + asap_types::WindowMaterializationLayout::FullWindow => staleness_margin_ms + .div_ceil(slide_ms.max(1)) + .saturating_add(1), + asap_types::WindowMaterializationLayout::Pane { pane_secs } => lookback_ms + .saturating_add(staleness_margin_ms) + .div_ceil(pane_secs.saturating_mul(1_000).max(1)) + .saturating_add(1), + asap_types::WindowMaterializationLayout::HierarchicalRollup { base_pane_secs, .. } => { + lookback_ms + .saturating_add(staleness_margin_ms) + .div_ceil(base_pane_secs.saturating_mul(1_000).max(1)) + .saturating_add(1) } - left } - - let minimum = crate::emit::stage_config::MIN_WINDOW_SECS; - let Some(source_ms) = source_sample_interval_ms else { - return minimum; - }; - if source_ms % 1_000 != 0 || u64::from(evaluation_interval_ms) % 1_000 != 0 { - return minimum; - } - let aligned = gcd( - candidate_pane_secs, - gcd(source_ms / 1_000, u64::from(evaluation_interval_ms) / 1_000), - ); - aligned.max(minimum) -} - -fn retained_window_count(lookback_ms: u64, staleness_margin_ms: u64, pane_ms: u64) -> u64 { - lookback_ms - .saturating_add(staleness_margin_ms) - .div_ceil(pane_ms.max(1)) - .saturating_add(1) } /// Estimate the encoded bytes retained by one physical state. Sketch matrix @@ -4111,8 +4145,8 @@ mod tests { assert_eq!(plan.precompute_plan.materializations.len(), 1); assert_eq!( plan.precompute_plan.materializations[0].num_aggregates_to_retain, - Some(13), - "one-minute readout retains twelve five-second panes plus a boundary pane" + Some(2), + "the explicitly selected one-minute pane plus a boundary pane is retained" ); assert!(plan .query_plan @@ -4522,8 +4556,8 @@ mod tests { implementation_id: "collector-tumbling-v1".into(), framework: SummaryWindowFramework::Tumbling, window_secs: 60, - pane_secs: 60, - state_layout: "anchored-pane-v1".into(), + slide_secs: 60, + layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 60 }, cost: ImplementationCostEvidence { model_version: "test-cost-v1".into(), workload_fingerprint: "test-workload".into(), @@ -4921,13 +4955,6 @@ mod tests { .any(|node| matches!(node, crate::query_plan::QueryPlanNode::ExactReadout { .. }))); } - #[test] - fn exact_temporal_pane_uses_observed_source_and_query_cadence() { - assert_eq!(exact_temporal_pane_secs(86_400, 60_000, Some(30_000)), 30); - assert_eq!(exact_temporal_pane_secs(43_200, 60_000, Some(30_000)), 30); - assert_eq!(exact_temporal_pane_secs(21_600, 60_000, None), 5); - } - fn phase_workload(phases: &[u64]) -> QueryWorkload { QueryWorkload { language: QueryLanguage::PromQL, @@ -4975,11 +5002,24 @@ mod tests { #[test] fn query_staleness_extends_retention_without_changing_readout_lookback() { + let panes = asap_types::WindowMaterializationLayout::Pane { pane_secs: 30 }; assert_eq!( - retained_window_count(6 * 60 * 60_000, 19 * 60_000, 30_000), + retained_state_count(6 * 60 * 60_000, 19 * 60_000, 30_000, &panes), 759 ); - assert_eq!(retained_window_count(6 * 60 * 60_000, 0, 30_000), 721); + assert_eq!( + retained_state_count(6 * 60 * 60_000, 0, 30_000, &panes), + 721 + ); + assert_eq!( + retained_state_count( + 6 * 60 * 60_000, + 19 * 60_000, + 30_000, + &asap_types::WindowMaterializationLayout::FullWindow, + ), + 39 + ); } // Both production adapters preserve canonical root identity and select the @@ -5182,26 +5222,17 @@ mod tests { #[test] fn shared_materialization_rejects_conflicting_deployment_contracts() { - for field in ["implementation", "layout"] { - let mut workload = request("q90", "quantile_over_time(0.90, m[1m])"); - let mut other = request("q99", "quantile_over_time(0.99, m[1m])"); - let implementation = &mut other.queries[0].window_implementations[0]; - if field == "implementation" { - implementation.implementation_id = "another-implementation".into(); - } else { - implementation.state_layout = "another-layout".into(); - } - workload.queries.extend(other.queries); - let error = PhysicalCompiler - .compile(workload, environment(10_000)) - .expect_err("conflicting shared state must fail before publication"); - assert!( - error - .to_string() - .contains("conflicting deployment contracts"), - "{error}" - ); - } + let mut workload = request("q90", "quantile_over_time(0.90, m[1m])"); + let mut other = request("q99", "quantile_over_time(0.99, m[1m])"); + other.queries[0].window_implementations[0].implementation_id = + "another-implementation".into(); + workload.queries.extend(other.queries); + let error = PhysicalCompiler + .compile(workload, environment(10_000)) + .expect_err("conflicting shared state must fail before publication"); + assert!(error + .to_string() + .contains("conflicting deployment contracts")); } #[test] @@ -5401,7 +5432,9 @@ mod tests { let mut five_minutes = query.window_implementations[0].clone(); five_minutes.implementation_id = "five-minute-evidence".into(); five_minutes.window_secs = 300; - five_minutes.pane_secs = 300; + five_minutes.framework = SummaryWindowFramework::Sliding; + five_minutes.slide_secs = 60; + five_minutes.layout = asap_types::WindowMaterializationLayout::Pane { pane_secs: 60 }; query.window_implementations.push(five_minutes); let plan = PhysicalCompiler.compile(request, env).unwrap(); let bindings = plan @@ -5796,7 +5829,7 @@ mod tests { plan.materializations[0].window_implementation_id, "collector-tumbling-v1" ); - assert_eq!(plan.materializations[0].pane_secs, 60); + assert_eq!(plan.materializations[0].slide_secs, 60); assert_eq!( plan.materializations[0].lifecycle, CollectorLifecycle { @@ -5932,7 +5965,6 @@ mod tests { evidence_valid_for_ms: 60_000, horizon_seconds: 300.0, window_implementation_id: "backend-tumbling-v1".into(), - state_layout: "anchored-pane-v1".into(), implementation_cost: template.window_implementations[0].cost.clone(), source_sample_interval_ms: None, query_staleness_margin_ms: 0, @@ -6250,13 +6282,15 @@ mod tests { // Concrete candidates with the same framework survive selection; changing // their quoted costs changes installed state, not the query's lookback. #[test] - fn tumbling_sizes_are_selected_by_cost_and_installed() { - for (small_cost, expected_secs, expected_id) in [(0.1, 10, "small"), (10.0, 60, "large")] { + fn pane_sizes_are_selected_by_cost_without_changing_semantic_window() { + for (small_cost, expected_pane_secs, expected_id) in + [(0.1, 10, "small"), (10.0, 60, "large")] + { let mut request = request("q", "sum(sum_over_time(m[1m]))"); let query = &mut request.queries[0]; let mut small = query.window_implementations[0].clone(); small.implementation_id = "small".into(); - small.pane_secs = 10; + small.layout = asap_types::WindowMaterializationLayout::Pane { pane_secs: 10 }; small.cost.weighted_cost = small_cost; query.window_implementations[0].implementation_id = "large".into(); query.window_implementations.push(small); @@ -6265,7 +6299,11 @@ mod tests { env.collector_ids.clear(); let bundle = PhysicalCompiler.compile(request, env).unwrap(); let materialization = bundle.precompute_plan.materializations.first().unwrap(); - assert_eq!(materialization.window_size, expected_secs); + assert_eq!(materialization.window_size, 60); + assert_eq!( + materialization.window_layout.base_pane_secs(), + expected_pane_secs + ); let entry = bundle .query_plan .lookup("sum(sum_over_time(m[1m]))") @@ -6273,7 +6311,7 @@ mod tests { assert_eq!(entry.instant.lookback_ms, 60_000); assert_eq!( entry.materialization_bindings()[0].window_ms, - expected_secs * 1000 + expected_pane_secs * 1000 ); assert_eq!( bundle.lifecycle_estimates[0].window_implementation_id, @@ -6282,7 +6320,59 @@ mod tests { } } - // Distinct logical cohorts cannot silently coalesce using only the first quote. + #[test] + fn sliding_layout_cost_selects_query_or_update_optimized_state() { + for (full_cost, expected_id, full_selected) in [ + (0.1, "full-window", true), + (100.0, "mergeable-panes", false), + ] { + let mut request = request("q", "sum(sum_over_time(m[1m]))"); + let query = &mut request.queries[0]; + let mut panes = query.window_implementations[0].clone(); + panes.implementation_id = "mergeable-panes".into(); + panes.framework = SummaryWindowFramework::Sliding; + panes.slide_secs = 10; + panes.layout = asap_types::WindowMaterializationLayout::Pane { pane_secs: 10 }; + panes.cost.weighted_cost = 1.0; + panes.cost.cpu_cost = 0.2; + panes.cost.storage_bytes = 1_024; + + let mut full = panes.clone(); + full.implementation_id = "full-window".into(); + full.layout = asap_types::WindowMaterializationLayout::FullWindow; + full.cost.weighted_cost = full_cost; + // Full windows spend more update CPU and retained bytes, while + // avoiding query-time pane merges. The provider's weighted quote + // includes the workload's measured read frequency and cardinality. + full.cost.cpu_cost = 20.0; + full.cost.storage_bytes = 64 * 1_024; + query.window_implementations = vec![panes, full]; + + let mut env = environment(10_000); + env.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; + env.collector_ids.clear(); + let bundle = PhysicalCompiler.compile(request, env).unwrap(); + assert_eq!( + bundle.lifecycle_estimates[0].window_implementation_id, + expected_id + ); + assert_eq!( + matches!( + bundle.precompute_plan.materializations[0].window_layout, + asap_types::WindowMaterializationLayout::FullWindow + ), + full_selected + ); + assert_eq!(bundle.precompute_plan.materializations[0].window_size, 60); + assert_eq!( + bundle.precompute_plan.materializations[0].slide_interval, + 10 + ); + } + } + + // Distinct semantic windows remain distinct definitions even when they use + // an equal base-pane width and implementation provider. #[test] fn selected_panes_reject_unpriced_cross_cohort_coalescing() { let mut workload = request("q20", "sum(sum_over_time(m[20s]))"); @@ -6296,22 +6386,25 @@ mod tests { second.window_implementations[0].window_secs = 40; workload.queries.push(second); for query in &mut workload.queries { - query.window_implementations[0].pane_secs = 10; + query.window_implementations[0].framework = SummaryWindowFramework::Sliding; + query.window_implementations[0].slide_secs = 10; + query.window_implementations[0].layout = + asap_types::WindowMaterializationLayout::Pane { pane_secs: 10 }; query.window_implementations[0].implementation_id = "shared-ten-second-pane".into(); } let mut env = environment(10_000); env.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; env.collector_ids.clear(); - assert!( - matches!(PhysicalCompiler.compile(workload, env), Err(CompileError::Lifecycle { reason, .. }) if reason.contains("distinct logical consumer cohorts")) - ); + let bundle = PhysicalCompiler.compile(workload, env).unwrap(); + assert_eq!(bundle.precompute_plan.materializations.len(), 2); } // Non-divisor panes cannot reconstruct a lookback from whole states. #[test] fn tumbling_sizes_reject_non_divisors() { let mut request = request("q", "sum(sum_over_time(m[1m]))"); - request.queries[0].window_implementations[0].pane_secs = 7; + request.queries[0].window_implementations[0].layout = + asap_types::WindowMaterializationLayout::Pane { pane_secs: 7 }; let mut env = environment(10_000); env.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; assert!(PhysicalCompiler.compile(request, env).is_err()); diff --git a/crates/asap_types/src/aggregation_config.rs b/crates/asap_types/src/aggregation_config.rs index ce66d9b9..9f16833d 100644 --- a/crates/asap_types/src/aggregation_config.rs +++ b/crates/asap_types/src/aggregation_config.rs @@ -10,6 +10,79 @@ use crate::utils::normalize_spatial_filter; use crate::AggregationType; use crate::KeyByLabelNames; +/// Physical maintenance layout for one semantic windowed summary. +/// +/// `window_size` and `slide_interval` on [`PrecomputeMaterialization`] retain +/// the query's window and evaluation cadence. This enum describes how that +/// semantic window is represented in storage; it must never be inferred by +/// overloading either semantic duration. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)] +#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] +pub enum WindowMaterializationLayout { + /// Store disjoint mergeable states and compose a query window at read time. + Pane { pane_secs: u64 }, + /// Maintain one complete state for every evaluation point. + FullWindow, + /// Store base panes plus coarser mergeable rollups. Every level is a + /// duration in seconds and is an integer multiple of its predecessor. + HierarchicalRollup { + base_pane_secs: u64, + levels_secs: Vec, + }, +} + +impl WindowMaterializationLayout { + pub fn base_pane_secs(&self) -> u64 { + match self { + Self::Pane { pane_secs } => *pane_secs, + Self::FullWindow => 0, + Self::HierarchicalRollup { base_pane_secs, .. } => *base_pane_secs, + } + } + + pub fn validate(&self, window_secs: u64, slide_secs: u64) -> Result<(), String> { + if window_secs == 0 || slide_secs == 0 || slide_secs > window_secs { + return Err( + "window and slide must be positive and slide must not exceed window".into(), + ); + } + match self { + Self::Pane { pane_secs } => { + if *pane_secs == 0 || window_secs % pane_secs != 0 || slide_secs % pane_secs != 0 { + return Err("pane size must divide both window size and slide".into()); + } + } + Self::FullWindow => {} + Self::HierarchicalRollup { + base_pane_secs, + levels_secs, + } => { + if *base_pane_secs == 0 + || window_secs % base_pane_secs != 0 + || slide_secs % base_pane_secs != 0 + || levels_secs.is_empty() + { + return Err( + "rollup base pane must divide window and slide, with at least one level" + .into(), + ); + } + let mut previous = *base_pane_secs; + for level in levels_secs { + if *level <= previous || *level % previous != 0 || window_secs % level != 0 { + return Err( + "rollup levels must increase by integral factors and divide the window" + .into(), + ); + } + previous = *level; + } + } + } + Ok(()) + } +} + /// Per-aggregation policy carried in the streaming config. /// /// **PR 5 (merged-sid-identity refactor)** retired the @@ -34,6 +107,7 @@ pub struct PrecomputeMaterialization { pub window_size: u64, // Window size in seconds (e.g., 900s for 15m) pub slide_interval: u64, // Slide/hop interval in seconds (e.g., 30s) pub window_type: WindowKind, // Tumbling or Sliding + pub window_layout: WindowMaterializationLayout, /// Unix millisecond timestamp on the pane-boundary grid selected from /// the consuming query workload. Missing on legacy definitions, which /// must not be used for certified pane-only reads. @@ -121,6 +195,13 @@ impl PrecomputeMaterialization { window_size, slide_interval, window_type, + window_layout: WindowMaterializationLayout::Pane { + pane_secs: if slide_interval == 0 { + window_size + } else { + slide_interval + }, + }, pane_origin_ms: None, spatial_filter, spatial_filter_normalized, @@ -420,6 +501,38 @@ impl SerializableToSink for PrecomputeMaterialization { } } +#[cfg(test)] +mod window_layout_tests { + use super::WindowMaterializationLayout; + + #[test] + fn validates_multiple_slides_and_rejects_uncomposable_panes() { + for slide in [5, 10, 15, 30] { + WindowMaterializationLayout::Pane { pane_secs: 5 } + .validate(60, slide) + .unwrap(); + WindowMaterializationLayout::FullWindow + .validate(60, slide) + .unwrap(); + } + assert!(WindowMaterializationLayout::Pane { pane_secs: 7 } + .validate(60, 10) + .is_err()); + assert!(WindowMaterializationLayout::HierarchicalRollup { + base_pane_secs: 5, + levels_secs: vec![10, 30], + } + .validate(60, 10) + .is_ok()); + assert!(WindowMaterializationLayout::HierarchicalRollup { + base_pane_secs: 5, + levels_secs: vec![12], + } + .validate(60, 10) + .is_err()); + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/asap_types/src/policy_fingerprint.rs b/crates/asap_types/src/policy_fingerprint.rs index 606796c9..eea167c8 100644 --- a/crates/asap_types/src/policy_fingerprint.rs +++ b/crates/asap_types/src/policy_fingerprint.rs @@ -148,6 +148,16 @@ impl PolicyFingerprint { ); buf.push(0); + // Physical representation is part of state identity. A full overlapping + // window and a mergeable pane layout may share semantic descriptors but + // never share payload instances or lifecycle accounting. + buf.extend_from_slice( + serde_json::to_string(&cfg.window_layout) + .unwrap_or_default() + .as_bytes(), + ); + buf.push(0); + // 9. pane origin. Presence is explicit so a legacy definition with // unknown phase cannot alias an epoch-aligned definition. match cfg.pane_origin_ms { @@ -239,6 +249,24 @@ mod tests { ); } + #[test] + fn physical_window_layout_is_part_of_state_identity() { + let panes = cfg( + "http_lat", + AggregationType::Sum, + HashMap::new(), + vec!["zone"], + 60, + "", + ); + let mut full = panes.clone(); + full.window_layout = crate::WindowMaterializationLayout::FullWindow; + assert_ne!( + PolicyFingerprint::from_config(&panes), + PolicyFingerprint::from_config(&full) + ); + } + #[test] fn different_metric_yields_different_fingerprint() { let a = cfg( diff --git a/crates/asap_types/src/summary_catalog.rs b/crates/asap_types/src/summary_catalog.rs index 7274dc1b..90095a33 100644 --- a/crates/asap_types/src/summary_catalog.rs +++ b/crates/asap_types/src/summary_catalog.rs @@ -10,6 +10,7 @@ use crate::sds::{ SummaryDescriptor, SummaryDescriptorId, ValueProjectionIdentity, }; use crate::PolicyFingerprint; +use crate::WindowMaterializationLayout; use serde::{Deserialize, Serialize}; pub const SUMMARY_CATALOG_SCHEMA_VERSION: u32 = 2; @@ -21,6 +22,10 @@ pub const SUMMARY_CATALOG_SCHEMA_VERSION: u32 = 2; pub struct SummaryDefinitionIdentity { pub summary_descriptor_id: SummaryDescriptorId, pub data_descriptor_id: DataDescriptorId, + /// Backend-selected physical representation. It is catalog-visible so + /// producers, readers, lifecycle management, and recovery agree on the + /// concrete state being referenced. + pub window_layout: WindowMaterializationLayout, /// Pane boundary selected from the shared consumer workload. Legacy /// snapshots deserialize as unknown and fail closed at pane-only reads. #[serde( @@ -123,6 +128,7 @@ impl SummaryCatalog { config.policy_fingerprint(), summary, data, + config.window_layout.clone(), config.pane_origin_ms, )) }) @@ -133,14 +139,23 @@ impl SummaryCatalog { pub fn build( plan_id: u64, plan_version: u64, - entries: impl IntoIterator, + entries: impl IntoIterator< + Item = ( + PolicyFingerprint, + SummaryDescriptor, + DataDescriptor, + WindowMaterializationLayout, + ), + >, ) -> Result { Self::build_with_origins( plan_id, plan_version, entries .into_iter() - .map(|(fingerprint, summary, data)| (fingerprint, summary, data, None)), + .map(|(fingerprint, summary, data, layout)| { + (fingerprint, summary, data, layout, None) + }), ) } @@ -152,6 +167,7 @@ impl SummaryCatalog { PolicyFingerprint, SummaryDescriptor, DataDescriptor, + WindowMaterializationLayout, Option, ), >, @@ -164,7 +180,7 @@ impl SummaryCatalog { data_descriptors: BTreeMap::new(), materializations: BTreeMap::new(), }; - for (fingerprint, summary, data, pane_origin_ms) in entries { + for (fingerprint, summary, data, window_layout, pane_origin_ms) in entries { let materialization = SummaryDefinitionId::from(fingerprint); summary .validate() @@ -174,6 +190,7 @@ impl SummaryCatalog { let binding = SummaryDefinitionIdentity { summary_descriptor_id: summary.id().clone(), data_descriptor_id: data.id().clone(), + window_layout, pane_origin_ms, }; if catalog @@ -374,8 +391,18 @@ mod tests { 1, 1, [ - (first.policy_fingerprint(), summary.clone(), data[0].clone()), - (first.policy_fingerprint(), summary, data[1].clone()), + ( + first.policy_fingerprint(), + summary.clone(), + data[0].clone(), + first.window_layout.clone(), + ), + ( + first.policy_fingerprint(), + summary, + data[1].clone(), + first.window_layout.clone(), + ), ], ) .unwrap_err(); diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index bf8473e6..573018d8 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -365,7 +365,6 @@ impl PrometheusRemoteWriteReceiver { .ingest .samples_ingested .fetch_add(new_samples.len() as u64, Ordering::Relaxed); - crate::precompute_engine::metrics::record_accepted_samples(new_samples.len() as u64); Ok(()) } } @@ -840,6 +839,7 @@ mod tests { window_size: 60, slide_interval: 60, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 60 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), @@ -892,6 +892,7 @@ mod tests { window_size: 60, slide_interval: 60, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 60 }, pane_origin_ms: Some(0), spatial_filter: String::new(), spatial_filter_normalized: String::new(), diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index dd6625cc..d0dbd6a6 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -14,8 +14,6 @@ use std::time::{Duration, Instant}; use tokio::net::TcpListener; use tracing::{debug, info, warn}; -const INGEST_DIAG_INTERVAL: Duration = Duration::from_secs(30); - use crate::drivers::query::adapters::{create_http_adapter, AdapterConfig, HttpProtocolAdapter}; use crate::drivers::query::servers::metrics as srv_metrics; use crate::query_engines::routing::{ @@ -201,48 +199,6 @@ pub struct HttpServer { remote_write: Option, } -/// Aborts the diagnostics task when the HTTP server exits or is cancelled. -struct AbortOnDrop(tokio::task::JoinHandle<()>); - -impl Drop for AbortOnDrop { - fn drop(&mut self) { - self.0.abort(); - } -} - -async fn log_ingest_throughput() { - let mut interval = tokio::time::interval(INGEST_DIAG_INTERVAL); - interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); - interval.tick().await; - let mut last = Instant::now(); - let mut previous = crate::precompute_engine::metrics::throughput_totals(); - loop { - interval.tick().await; - let now = Instant::now(); - let elapsed = now.duration_since(last); - last = now; - // Keep all counters monotonic. Observers derive interval deltas from - // loads, so logging and /metrics scraping cannot alter one another. - let current = crate::precompute_engine::metrics::throughput_totals(); - let delta = current.since(previous); - previous = current; - let secs = elapsed.as_secs_f64(); - debug!( - accepted_samples = delta.accepted_samples, - processed_updates = delta.processed_updates, - materialized_outputs = delta.materialized_outputs, - elapsed_seconds = secs, - accepted_per_second = delta.accepted_samples as f64 / secs, - processed_per_second = delta.processed_updates as f64 / secs, - materialized_per_second = delta.materialized_outputs as f64 / secs, - cumulative_accepted = current.accepted_samples, - cumulative_processed = current.processed_updates, - cumulative_materialized = current.materialized_outputs, - "precompute throughput interval" - ); - } -} - #[derive(Clone)] struct AppState { config: HttpServerConfig, @@ -577,13 +533,6 @@ impl HttpServer { let listener = TcpListener::bind(format!("0.0.0.0:{}", self.config.port)).await?; info!("HTTP server listening on port {}", self.config.port); - // Bind first so a listener failure cannot leave an orphaned task. - // Keep the guard alive through serve so cancellation also aborts it. - let _ingest_ticker = self - .remote_write - .as_ref() - .map(|_| AbortOnDrop(tokio::spawn(log_ingest_throughput()))); - axum::serve(listener, app).await?; Ok(()) } @@ -3713,6 +3662,7 @@ aggregations: window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), @@ -7166,29 +7116,6 @@ mod logical_provenance_tests { let mut value = serde_json::json!({"warnings":["asap_execution:asap", "asap_logical_stats:raw=1,summary=1,memo_hits=0"]}); assert_eq!(extract_logical_provenance(&mut value), Some(Err(()))); } - - #[tokio::test] - async fn ingest_ticker_is_cancelled_when_guard_is_dropped() { - let counter = Arc::new(std::sync::atomic::AtomicU64::new(0)); - let task_counter = counter.clone(); - let guard = AbortOnDrop(tokio::spawn(async move { - loop { - task_counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst); - tokio::time::sleep(Duration::from_millis(5)).await; - } - })); - - while counter.load(std::sync::atomic::Ordering::SeqCst) == 0 { - tokio::task::yield_now().await; - } - drop(guard); - let count_at_drop = counter.load(std::sync::atomic::Ordering::SeqCst); - tokio::time::sleep(Duration::from_millis(100)).await; - assert_eq!( - counter.load(std::sync::atomic::Ordering::SeqCst), - count_at_drop - ); - } } #[cfg(test)] diff --git a/data_plane/src/precompute_engine/output_sink.rs b/data_plane/src/precompute_engine/output_sink.rs index 5acb9eab..47d5cb31 100644 --- a/data_plane/src/precompute_engine/output_sink.rs +++ b/data_plane/src/precompute_engine/output_sink.rs @@ -180,7 +180,6 @@ impl OutputSink for SketchStoreSink { ) .into()); } - crate::precompute_engine::metrics::record_materialized_outputs(output_count as u64); Ok(()) } } @@ -286,6 +285,7 @@ mod tests { window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), diff --git a/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs b/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs index da96295c..a5feb2f2 100644 --- a/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs +++ b/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs @@ -267,6 +267,7 @@ mod tests { window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), diff --git a/data_plane/src/storage_engines/sketch_db/sds.rs b/data_plane/src/storage_engines/sketch_db/sds.rs index 6cbdd0ff..e54cac85 100644 --- a/data_plane/src/storage_engines/sketch_db/sds.rs +++ b/data_plane/src/storage_engines/sketch_db/sds.rs @@ -370,6 +370,7 @@ mod tests { asap_types::PolicyFingerprint(7), summary.clone(), data.clone(), + asap_types::WindowMaterializationLayout::Pane { pane_secs: 60 }, )], ) .unwrap(); diff --git a/data_plane/src/tests/test_utilities/engine_factories.rs b/data_plane/src/tests/test_utilities/engine_factories.rs index 63fc9fd2..b6cae5b0 100644 --- a/data_plane/src/tests/test_utilities/engine_factories.rs +++ b/data_plane/src/tests/test_utilities/engine_factories.rs @@ -100,6 +100,7 @@ pub fn create_engine_single_pop_with_aggregated( window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), @@ -193,6 +194,7 @@ pub fn create_engine_dual_input( window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), @@ -216,6 +218,7 @@ pub fn create_engine_dual_input( window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), @@ -304,6 +307,7 @@ pub fn create_engine_two_metrics( window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), @@ -326,6 +330,7 @@ pub fn create_engine_two_metrics( window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), @@ -424,6 +429,7 @@ pub fn create_engine_three_metrics( window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), @@ -499,6 +505,7 @@ pub fn create_engine_multi_timestamp( window_size: 1, slide_interval: 1, window_type: WindowKind::Tumbling, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), @@ -566,6 +573,7 @@ pub fn create_engine_multi_timestamp_with_window( window_size, slide_interval: 1, window_type, + window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 1 }, pane_origin_ms: None, spatial_filter: String::new(), spatial_filter_normalized: String::new(), diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index c0e6a61d..e1c994a0 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -710,7 +710,7 @@ async fn run_shared_dashboard(multi_pane: bool) { let mut candidates = query.window_implementations; let mut small = candidates[0].clone(); small.implementation_id = "five-second-pane".into(); - small.pane_secs = 5; + small.layout = asap_types::WindowMaterializationLayout::Pane { pane_secs: 5 }; small.cost.weighted_cost = 0.0; candidates[0].cost.weighted_cost = 10.0; candidates.push(small); diff --git a/docs/examples/asapquery-compatibility-demo-snapshot.json b/docs/examples/asapquery-compatibility-demo-snapshot.json index e427b248..e345114c 100644 --- a/docs/examples/asapquery-compatibility-demo-snapshot.json +++ b/docs/examples/asapquery-compatibility-demo-snapshot.json @@ -83,7 +83,6 @@ "evidence_valid_for_ms": 60000, "horizon_seconds": 300.0, "window_implementation_id": "backend-tumbling-v1", - "state_layout": "anchored-pane-v1", "implementation_cost": { "model_version": "compat-cost-v1", "workload_fingerprint": "asapquery-compatibility-demo", diff --git a/docs/examples/asapquery-planning-snapshot.json b/docs/examples/asapquery-planning-snapshot.json index df3071e5..99c2a284 100644 --- a/docs/examples/asapquery-planning-snapshot.json +++ b/docs/examples/asapquery-planning-snapshot.json @@ -50,7 +50,6 @@ "evidence_valid_for_ms": 60000, "horizon_seconds": 300.0, "window_implementation_id": "backend-tumbling-v1", - "state_layout": "anchored-pane-v1", "implementation_cost": { "model_version": "compat-cost-v1", "workload_fingerprint": "replaced-with-canonical-promql",