From b6fcf066045145e8819ab56bcffc0c939e5df0e8 Mon Sep 17 00:00:00 2001 From: Zeying Zhu <50204836+zzylol@users.noreply.github.com> Date: Thu, 10 Sep 2026 13:27:46 -0400 Subject: [PATCH 1/5] Add costed sliding-window materialization layouts (#580) --- control_plane/src/opamp/mod.rs | 4 +- control_plane/src/physical/compiler.rs | 341 +++++++++++------- crates/asap_types/src/aggregation_config.rs | 113 ++++++ crates/asap_types/src/policy_fingerprint.rs | 28 ++ crates/asap_types/src/summary_catalog.rs | 37 +- .../drivers/ingest/prometheus_remote_write.rs | 3 +- data_plane/src/drivers/query/servers/http.rs | 75 +--- .../src/precompute_engine/output_sink.rs | 2 +- .../sketch_db/lifecycle/eviction.rs | 1 + .../src/storage_engines/sketch_db/sds.rs | 1 + .../tests/test_utilities/engine_factories.rs | 8 + .../asapquery_compatibility_process_e2e.rs | 2 +- ...asapquery-compatibility-demo-snapshot.json | 1 - .../examples/asapquery-planning-snapshot.json | 1 - 14 files changed, 407 insertions(+), 210 deletions(-) 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", From 3392f69e29d2b9875488273590cc4791118b4977 Mon Sep 17 00:00:00 2001 From: Zeying Zhu <50204836+zzylol@users.noreply.github.com> Date: Thu, 10 Sep 2026 13:39:06 -0400 Subject: [PATCH 2/5] Execute selected sliding window layouts in workers (#583) --- .../src/precompute_engine/window_manager.rs | 51 ++- data_plane/src/precompute_engine/worker.rs | 323 +++++++++++------- 2 files changed, 245 insertions(+), 129 deletions(-) diff --git a/data_plane/src/precompute_engine/window_manager.rs b/data_plane/src/precompute_engine/window_manager.rs index 51b92308..89dd5bbd 100644 --- a/data_plane/src/precompute_engine/window_manager.rs +++ b/data_plane/src/precompute_engine/window_manager.rs @@ -8,6 +8,8 @@ pub struct WindowManager { window_size_ms: i64, /// Slide interval in milliseconds (== window_size_ms for tumbling windows). slide_interval_ms: i64, + /// Width of independently stored non-overlapping panes. + pane_interval_ms: i64, /// Planned event-time phase of this materialization definition. origin_ms: Option, } @@ -35,10 +37,24 @@ impl WindowManager { Self { window_size_ms, slide_interval_ms, + pane_interval_ms: slide_interval_ms, origin_ms, } } + pub fn with_layout( + window_size_secs: u64, + slide_interval_secs: u64, + origin_ms: Option, + layout: &asap_types::WindowMaterializationLayout, + ) -> Self { + let mut manager = Self::with_origin(window_size_secs, slide_interval_secs, origin_ms); + if !matches!(layout, asap_types::WindowMaterializationLayout::FullWindow) { + manager.pane_interval_ms = (layout.base_pane_secs() * 1_000) as i64; + } + manager + } + pub fn window_size_ms(&self) -> i64 { self.window_size_ms } @@ -95,13 +111,12 @@ impl WindowManager { return Vec::new(); } let mut panes = Vec::new(); - let mut start = - self.window_start_for(previous_wm.saturating_sub(self.slide_interval_ms - 1)); - while start.saturating_add(self.slide_interval_ms) <= current_wm { - if start.saturating_add(self.slide_interval_ms) > previous_wm { + let mut start = self.pane_start_for(previous_wm.saturating_sub(self.pane_interval_ms - 1)); + while start.saturating_add(self.pane_interval_ms) <= current_wm { + if start.saturating_add(self.pane_interval_ms) > previous_wm { panes.push(start); } - start = start.saturating_add(self.slide_interval_ms); + start = start.saturating_add(self.pane_interval_ms); } panes } @@ -126,7 +141,7 @@ impl WindowManager { } pub fn pane_bounds(&self, pane_start: i64) -> (i64, i64) { - (pane_start, pane_start + self.slide_interval_ms) + (pane_start, pane_start + self.pane_interval_ms) } /// Slide interval accessor. @@ -137,16 +152,17 @@ impl WindowManager { /// Pane start for a timestamp. Panes are aligned to the slide_interval grid, /// which is the same grid as `window_start_for`. pub fn pane_start_for(&self, timestamp_ms: i64) -> i64 { - self.window_start_for(timestamp_ms) + let origin = self.origin_ms.unwrap_or(0); + origin + (timestamp_ms - origin).div_euclid(self.pane_interval_ms) * self.pane_interval_ms } /// All pane starts composing a window, in ascending order. /// A window `[ws, ws + window_size)` is composed of /// `window_size / slide_interval` consecutive panes. pub fn panes_for_window(&self, window_start: i64) -> Vec { - let num_panes = self.window_size_ms / self.slide_interval_ms; + let num_panes = self.window_size_ms / self.pane_interval_ms; (0..num_panes) - .map(|i| window_start + i * self.slide_interval_ms) + .map(|i| window_start + i * self.pane_interval_ms) .collect() } } @@ -213,6 +229,23 @@ mod tests { assert_eq!(wm.pane_bounds(20_000), (20_000, 30_000)); } + #[test] + fn physical_pane_width_is_independent_of_query_slide() { + let wm = WindowManager::with_layout( + 60, + 30, + Some(0), + &asap_types::WindowMaterializationLayout::Pane { pane_secs: 10 }, + ); + assert_eq!(wm.pane_start_for(29_999), 20_000); + assert_eq!(wm.pane_bounds(20_000), (20_000, 30_000)); + assert_eq!( + wm.panes_for_window(0), + vec![0, 10_000, 20_000, 30_000, 40_000, 50_000] + ); + assert_eq!(wm.closed_panes(15_000, 35_000), vec![10_000, 20_000]); + } + #[test] fn test_no_close_when_watermark_stagnant() { let wm = WindowManager::new(60, 0); diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index 58177388..6cec5063 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -84,6 +84,37 @@ struct PaneWallClock { } impl GroupState { + fn stores_full_windows(&self) -> bool { + matches!( + self.config.window_layout, + asap_types::WindowMaterializationLayout::FullWindow + ) + } + + fn bucket_starts_for(&self, timestamp_ms: i64) -> Vec { + if self.stores_full_windows() { + self.window_manager.window_starts_containing(timestamp_ms) + } else { + vec![self.window_manager.pane_start_for(timestamp_ms)] + } + } + + fn closed_buckets(&self, previous_ms: i64, current_ms: i64) -> Vec { + if self.stores_full_windows() { + self.window_manager.closed_windows(previous_ms, current_ms) + } else { + self.window_manager.closed_panes(previous_ms, current_ms) + } + } + + fn bucket_bounds(&self, start_ms: i64) -> (i64, i64) { + if self.stores_full_windows() { + self.window_manager.window_bounds(start_ms) + } else { + self.window_manager.pane_bounds(start_ms) + } + } + fn touch_pane(&mut self, pane_start_ms: i64, now_ms: i64) { self.pane_wall_clock .entry(pane_start_ms) @@ -365,10 +396,11 @@ impl Worker { let cfg = snap.get_aggregation_config(policy_fp.as_u64())?; let config = Arc::new(cfg.clone()); let gs = GroupState { - window_manager: WindowManager::with_origin( + window_manager: WindowManager::with_layout( config.window_size, config.slide_interval, config.pane_origin_ms, + &config.window_layout, ), config, policy_fp, @@ -459,74 +491,13 @@ impl Worker { let mut emit_batch: Vec<(PrecomputedOutput, Box)> = Vec::new(); - // Route each sample to its pane + // The selected physical layout owns update fanout. Pane and rollup + // layouts update one non-overlapping base pane; FullWindow updates + // every overlapping semantic window that contains the sample. for (series_key, ts, val) in &samples { let too_late = previous_event_time != i64::MIN && pane_timestamp(*ts) < watermark_for_event_time(previous_event_time, allowed_lateness_ms); - let pane_start = state.window_manager.pane_start_for(pane_timestamp(*ts)); - let pane_end = pane_start + state.window_manager.slide_interval_ms(); - let pane_closed = !state.active_panes.contains_key(&pane_start) - && previous_closure_watermark >= pane_end; - - if too_late || pane_closed { - let window_start = pane_start; - let window_end = pane_end; - match late_data_policy { - LateDataPolicy::Drop => { - record_late_input("drop", "raw_sample"); - debug!( - "Worker {} dropping late sample for sid={} (group={}): \ - ts={} observed_event_time={} pane=[{}, {})", - worker_id, - sid, - group_key, - ts, - previous_event_time, - pane_start, - pane_end - ); - continue; - } - LateDataPolicy::ForwardToStore => { - // A late cumulative sample cannot be converted to a - // derivative without its time-adjacent neighbours. - // Never feed the raw counter value into a membership - // heap; the authoritative ExactCounter branch remains - // responsible for the visible result. - if matches!( - state.config.sample_update_rule(), - SampleUpdateRule::CounterDelta { .. } - ) { - record_late_input("drop", "counter_delta_membership"); - continue; - } - record_late_input("append_correction", "raw_sample"); - let mut updater = create_accumulator_updater(&state.config); - apply_sample(&mut *updater, series_key, *val, *ts, &state.config); - let key = build_group_key_label_values(group_key); - let output = precomputed_output_for_group( - window_start as u64, - window_end as u64, - key, - PolicyFingerprint::from_config(&state.config), - group_key, - ); - emit_batch.push((output, updater.take_accumulator())); - debug!( - "Forwarding late sample to store for evicted pane [{}, {})", - pane_start, pane_end - ); - continue; - } - } - } - - // Normal path: route sample to its single pane accumulator. - // Refresh the pane's wall-clock last-touch time so the fallback - // only closes an idle pane, not a long-running bulk ingest whose - // records share one event timestamp. - state.touch_pane(pane_start, now_ms); let value = if let SampleUpdateRule::CounterDelta { .. } = state.config.sample_update_rule() { @@ -539,11 +510,75 @@ impl Worker { } else { *val }; - let updater = state - .active_panes - .entry(pane_start) - .or_insert_with(|| create_accumulator_updater(&state.config)); - apply_sample(&mut **updater, series_key, value, *ts, &state.config); + for bucket_start in state.bucket_starts_for(pane_timestamp(*ts)) { + let (_, bucket_end) = state.bucket_bounds(bucket_start); + let bucket_closed = !state.active_panes.contains_key(&bucket_start) + && previous_closure_watermark >= bucket_end; + + if too_late || bucket_closed { + let window_start = bucket_start; + let window_end = bucket_end; + match late_data_policy { + LateDataPolicy::Drop => { + record_late_input("drop", "raw_sample"); + debug!( + "Worker {} dropping late sample for sid={} (group={}): \ + ts={} observed_event_time={} pane=[{}, {})", + worker_id, + sid, + group_key, + ts, + previous_event_time, + bucket_start, + bucket_end + ); + continue; + } + LateDataPolicy::ForwardToStore => { + // A late cumulative sample cannot be converted to a + // derivative without its time-adjacent neighbours. + // Never feed the raw counter value into a membership + // heap; the authoritative ExactCounter branch remains + // responsible for the visible result. + if matches!( + state.config.sample_update_rule(), + SampleUpdateRule::CounterDelta { .. } + ) { + record_late_input("drop", "counter_delta_membership"); + continue; + } + record_late_input("append_correction", "raw_sample"); + let mut updater = create_accumulator_updater(&state.config); + apply_sample(&mut *updater, series_key, *val, *ts, &state.config); + let key = build_group_key_label_values(group_key); + let output = precomputed_output_for_group( + window_start as u64, + window_end as u64, + key, + PolicyFingerprint::from_config(&state.config), + group_key, + ); + emit_batch.push((output, updater.take_accumulator())); + debug!( + "Forwarding late sample to store for evicted pane [{}, {})", + bucket_start, bucket_end + ); + continue; + } + } + } + + // Normal path: route sample to its single pane accumulator. + // Refresh the pane's wall-clock last-touch time so the fallback + // only closes an idle pane, not a long-running bulk ingest whose + // records share one event timestamp. + state.touch_pane(bucket_start, now_ms); + let updater = state + .active_panes + .entry(bucket_start) + .or_insert_with(|| create_accumulator_updater(&state.config)); + apply_sample(&mut **updater, series_key, value, *ts, &state.config); + } } // Check for closed windows @@ -552,12 +587,10 @@ impl Worker { } else { previous_closure_watermark }; - let closed = state - .window_manager - .closed_panes(closure_scan_start, event_watermark); + let closed = state.closed_buckets(closure_scan_start, event_watermark); for window_start in &closed { - let (_, window_end) = state.window_manager.pane_bounds(*window_start); + let (_, window_end) = state.bucket_bounds(*window_start); let pane_starts = [*window_start]; if let Some(accumulator) = merge_panes_for_window(&mut state.active_panes, &pane_starts) @@ -650,12 +683,20 @@ impl Worker { // Late-arrival check against the existing watermark. let too_late = previous_event_time != i64::MIN && timestamp_ms < watermark_for_event_time(previous_event_time, allowed_lateness_ms); - let pane_start = state.window_manager.pane_start_for(timestamp_ms); - let pane_end = pane_start + state.window_manager.slide_interval_ms(); - let pane_closed = - !state.sketch_panes.contains_key(&pane_start) && previous_closure_watermark >= pane_end; + // A pre-built accumulator is already a materialized interval, not a + // point sample. Never duplicate its state across overlapping windows. + // Full-window producers stamp the corresponding slide boundary. + let bucket_starts = if state.stores_full_windows() { + vec![state.window_manager.window_start_for(timestamp_ms)] + } else { + state.bucket_starts_for(timestamp_ms) + }; + let any_bucket_closed = bucket_starts.iter().any(|start| { + let (_, end) = state.bucket_bounds(*start); + !state.sketch_panes.contains_key(start) && previous_closure_watermark >= end + }); - if too_late || pane_closed { + if too_late || any_bucket_closed { match late_data_policy { LateDataPolicy::Drop => { record_late_input("drop", "prebuilt_sketch"); @@ -666,22 +707,18 @@ impl Worker { } LateDataPolicy::ForwardToStore => { record_late_input("append_correction", "prebuilt_sketch"); - let window_start = pane_start; - let window_end = pane_end; - let key = build_group_key_label_values(group_key); - let output = precomputed_output_for_group( - window_start as u64, - window_end as u64, - key, - PolicyFingerprint::from_config(&state.config), - group_key, - ); - emit_batch.push((output, incoming)); - debug!( - "Forwarding late accumulator input to store for evicted pane [{}, {})", - pane_start, - pane_start + state.window_manager.slide_interval_ms() - ); + for window_start in bucket_starts { + let (_, window_end) = state.bucket_bounds(window_start); + let key = build_group_key_label_values(group_key); + let output = precomputed_output_for_group( + window_start as u64, + window_end as u64, + key, + PolicyFingerprint::from_config(&state.config), + group_key, + ); + emit_batch.push((output, incoming.clone_boxed_core())); + } self.output_sink.emit_batch(emit_batch)?; } } @@ -690,27 +727,27 @@ impl Worker { // Refresh the pane's wall-clock last-touch time so an active sketch // stream with a fixed event timestamp is not force-closed mid-ingest. - state.touch_pane(pane_start, now_ms); - - // Merge into the sketch pane covering this timestamp. - match state.sketch_panes.remove(&pane_start) { - Some(existing) => { - let merged = existing - .merge_with(incoming.as_ref()) - .map_err(|e| format!("merge_with failed for pane {pane_start}: {e}"))?; - state.sketch_panes.insert(pane_start, merged); - } - None => { - state.sketch_panes.insert(pane_start, incoming); + for bucket_start in bucket_starts { + state.touch_pane(bucket_start, now_ms); + match state.sketch_panes.remove(&bucket_start) { + Some(existing) => { + let merged = existing + .merge_with(incoming.as_ref()) + .map_err(|e| format!("merge_with failed for bucket {bucket_start}: {e}"))?; + state.sketch_panes.insert(bucket_start, merged); + } + None => { + state + .sketch_panes + .insert(bucket_start, incoming.clone_boxed_core()); + } } } // Check for closed windows and emit merged outputs. - let closed = state - .window_manager - .closed_panes(previous_closure_watermark, event_watermark); + let closed = state.closed_buckets(previous_closure_watermark, event_watermark); for window_start in &closed { - let (_, window_end) = state.window_manager.pane_bounds(*window_start); + let (_, window_end) = state.bucket_bounds(*window_start); let pane_starts = [*window_start]; // Emit from the raw-sample pane map (in case both sources are @@ -912,12 +949,10 @@ impl Worker { } } - let closed = state - .window_manager - .closed_panes(state.closure_watermark_ms, effective_wm); + let closed = state.closed_buckets(state.closure_watermark_ms, effective_wm); for window_start in &closed { - let (_, window_end) = state.window_manager.pane_bounds(*window_start); + let (_, window_end) = state.bucket_bounds(*window_start); let pane_starts = [*window_start]; if let Some(accumulator) = @@ -1011,15 +1046,13 @@ impl Worker { (None, Some(b)) => b, (None, None) => continue, // no open panes }; - let force_wm = max_pane.saturating_add(state.window_manager.slide_interval_ms()); + let (_, force_wm) = state.bucket_bounds(max_pane); let group_key = state.group_key.clone(); - let closed = state - .window_manager - .closed_panes(state.closure_watermark_ms, force_wm); + let closed = state.closed_buckets(state.closure_watermark_ms, force_wm); for window_start in &closed { - let (_, window_end) = state.window_manager.pane_bounds(*window_start); + let (_, window_end) = state.bucket_bounds(*window_start); let pane_starts = [*window_start]; if let Some(accumulator) = @@ -2047,6 +2080,56 @@ mod tests { } } + #[test] + fn full_window_layout_materializes_each_overlapping_slide() { + let mut config = make_agg_config( + 2, + "cpu", + AggregationType::SingleSubpopulation, + "Sum", + 30, + 10, + vec![], + ); + config.window_layout = asap_types::WindowMaterializationLayout::FullWindow; + let policy = config.policy_fingerprint(); + let sink = Arc::new(CapturingOutputSink::new()); + let mut worker = make_worker( + HashMap::from([(policy.as_u64(), config)]), + sink.clone(), + false, + 0, + LateDataPolicy::Drop, + ); + + worker + .process_group_samples(2, policy, "", group_samples("cpu", vec![(15_000, 42.0)])) + .unwrap(); + worker + .process_group_samples(2, policy, "", group_samples("cpu", vec![(45_000, 0.0)])) + .unwrap(); + + let captured = sink.drain(); + assert_eq!(captured.len(), 2); + assert_eq!( + captured + .iter() + .map(|(output, _)| output.start_timestamp) + .collect::>(), + vec![0, 10_000] + ); + for (_, accumulator) in captured { + assert_eq!( + accumulator + .as_any() + .downcast_ref::() + .unwrap() + .sum, + 42.0 + ); + } + } + // ----------------------------------------------------------------------- // Test: MultipleSubpopulation — keyed accumulator with aggregated labels // Matches planner output: grouping=[], aggregated=[host] From c6b37bd3a71c690e87aab2e5c40706e2998c672a Mon Sep 17 00:00:00 2001 From: Zeying Zhu <50204836+zzylol@users.noreply.github.com> Date: Thu, 10 Sep 2026 13:43:06 -0400 Subject: [PATCH 3/5] Execute ClickHouse joins in the shared query DAG (#569) * feat(clickhouse): execute SQL joins in the shared query DAG * fix: consume canonical planner lifecycle contract --- Cargo.lock | 36 ++--- control_plane/Cargo.toml | 10 +- .../examples/calibration_candidates.rs | 10 ++ .../examples/offline_planner_replay.rs | 5 + control_plane/src/emit/mod.rs | 1 + .../src/physical/colored_dag/allocator.rs | 7 + .../src/physical/colored_dag/emitter.rs | 1 + control_plane/src/physical/compiler.rs | 14 +- control_plane/src/physical/post_asap/tests.rs | 1 + control_plane/src/query_plan.rs | 33 +++- control_plane/src/query_plan/logical.rs | 4 + crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +- .../asap_clickhouse_query_engine/execution.rs | 149 ++++++++++++++++++ .../relational_adapter.rs | 99 ++++++++++++ .../asap_query_engine/post_asap_readout.rs | 3 +- .../asap_query_engine/summary_exec.rs | 1 + 17 files changed, 348 insertions(+), 34 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 9acd1f3b..96eaa3f9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -374,7 +374,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "asap-types", "serde", @@ -385,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-metricsql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "asap-types", "metricsql_parser", @@ -395,7 +395,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "asap-types", "promql-parser 0.10.0", @@ -404,7 +404,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -426,12 +426,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "serde", "serde_json", @@ -1711,7 +1711,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -2260,7 +2260,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.6.3", + "socket2 0.5.10", "tokio", "tower-service", "tracing", @@ -2453,7 +2453,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" dependencies = [ "hermit-abi 0.5.2", "libc", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -2801,7 +2801,7 @@ dependencies = [ [[package]] name = "metricsql_common" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "chrono", ] @@ -2809,7 +2809,7 @@ dependencies = [ [[package]] name = "metricsql_parser" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "ahash", "chrono", @@ -3433,7 +3433,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "be769465445e8c1474e9c5dac2018218498557af32d9ed057325ec9a41ae81bf" dependencies = [ "heck 0.5.0", - "itertools 0.13.0", + "itertools 0.10.5", "log", "multimap", "once_cell", @@ -3453,7 +3453,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d" dependencies = [ "anyhow", - "itertools 0.13.0", + "itertools 0.10.5", "proc-macro2", "quote", "syn 2.0.117", @@ -3551,7 +3551,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls 0.23.40", - "socket2 0.6.3", + "socket2 0.5.10", "thiserror 2.0.18", "tokio", "tracing", @@ -3588,7 +3588,7 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.6.3", + "socket2 0.5.10", "tracing", "windows-sys 0.52.0", ] @@ -3904,7 +3904,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -4438,7 +4438,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix 1.1.4", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -5260,7 +5260,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.48.0", ] [[package]] diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 47d3bbf1..6a30d061 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -91,8 +91,8 @@ asap_types.workspace = true # scaffolding, unaware that `data_plane`'s `summary_executor.rs` in *this* # repo is a real one. Vendored locally instead of chased upstream -- see # `data_plane/src/query_engines/asap_query_engine/summary_exec.rs`. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } # L1 adoption (design-target-architecture.md Part B): the PromQL front # end itself, replacing control_plane's own query_parser/promql.rs. @@ -100,9 +100,9 @@ asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = # `planner-types`/`asap-aware-mapping` above -- these three MUST move # together (two revs of the same upstream repo's types in one workspace # resolve to distinct Rust types that won't unify). -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } -asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/control_plane/examples/calibration_candidates.rs b/control_plane/examples/calibration_candidates.rs index 64ce8d3c..607bc150 100644 --- a/control_plane/examples/calibration_candidates.rs +++ b/control_plane/examples/calibration_candidates.rs @@ -74,6 +74,16 @@ fn planner_forest(queries: &[control_plane::physical::compiler::PlanningQuery]) vec![outer, inner], json!({"key_debug":format!("{key:?}"),"family_debug":format!("{family:?}")}), ), + SummaryExpr::RelationalJoin { + left, + right, + kind, + pred, + } => ( + "RelationalJoin", + vec![left, right], + json!({"kind_debug":format!("{kind:?}"),"predicate_debug":format!("{pred:?}")}), + ), SummaryExpr::SummarySubtract { left, right } => { ("SummarySubtract", vec![left, right], json!({})) } diff --git a/control_plane/examples/offline_planner_replay.rs b/control_plane/examples/offline_planner_replay.rs index 0ccf9d05..f0e31758 100644 --- a/control_plane/examples/offline_planner_replay.rs +++ b/control_plane/examples/offline_planner_replay.rs @@ -73,6 +73,11 @@ fn inspect( left: lhs, right: rhs, } + | SummaryExpr::RelationalJoin { + left: lhs, + right: rhs, + .. + } | SummaryExpr::BinaryOp { lhs, rhs, .. } => { inspect(lhs, model, seen, states, raw); inspect(rhs, model, seen, states, raw); diff --git a/control_plane/src/emit/mod.rs b/control_plane/src/emit/mod.rs index af4fe16f..9732e9de 100644 --- a/control_plane/src/emit/mod.rs +++ b/control_plane/src/emit/mod.rs @@ -326,6 +326,7 @@ fn extract_from_node(node: &Rc) -> Option { // Not surfaced by any `Bind*` path yet (gated on rules that // haven't landed — see `deployment_expr.rs`'s module docs). SummaryExpr::BinaryOp { .. } + | SummaryExpr::RelationalJoin { .. } | SummaryExpr::SummaryJoin { .. } | SummaryExpr::SummarySubtract { .. } | SummaryExpr::SummaryDelete { .. } diff --git a/control_plane/src/physical/colored_dag/allocator.rs b/control_plane/src/physical/colored_dag/allocator.rs index b66dcff3..5cbf74a2 100644 --- a/control_plane/src/physical/colored_dag/allocator.rs +++ b/control_plane/src/physical/colored_dag/allocator.rs @@ -184,6 +184,13 @@ impl ThreeStageWalker { } StageId::Backend } + SummaryExpr::RelationalJoin { left, right, .. } => { + for child in [left, right] { + let (cid, _) = self.visit_l4node(child)?; + self.dag.edges.push((id, cid)); + } + StageId::Backend + } SummaryExpr::ValueOperation { child, timing, .. } => { let (cid, child_stage) = self.visit_l4node(child)?; self.dag.edges.push((id, cid)); diff --git a/control_plane/src/physical/colored_dag/emitter.rs b/control_plane/src/physical/colored_dag/emitter.rs index 2bd4bcf2..9d230e89 100644 --- a/control_plane/src/physical/colored_dag/emitter.rs +++ b/control_plane/src/physical/colored_dag/emitter.rs @@ -117,6 +117,7 @@ fn classify(expr: &PhysicalExpr) -> NodeKind<'_> { SummaryExpr::SummaryEstimate { query, .. } => NodeKind::SketchEstimate { query }, SummaryExpr::SummaryMerge { .. } => NodeKind::SketchMerge, SummaryExpr::BinaryOp { .. } + | SummaryExpr::RelationalJoin { .. } | SummaryExpr::CandidateTopK { .. } | SummaryExpr::ValueOperation { .. } | SummaryExpr::SummaryJoin { .. } diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index d3402b79..5551a036 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -2914,6 +2914,7 @@ fn summary_agg_metric(node: &SummaryNode) -> Option { inner: right, .. } + | SummaryExpr::RelationalJoin { left, right, .. } | SummaryExpr::SummarySubtract { left, right } | SummaryExpr::BinaryOp { lhs: left, @@ -3061,6 +3062,7 @@ fn requires_exact_erp_fallback( walk(inner, out); } SummaryExpr::SummarySubtract { left, right } + | SummaryExpr::RelationalJoin { left, right, .. } | SummaryExpr::BinaryOp { lhs: left, rhs: right, @@ -3508,12 +3510,12 @@ fn select_lifecycle( reason: "latest ASAPPlanner selected no window framework from the supplied physical evidence".into(), })?; Ok(PlannerPhysicalSelection { - window_implementation_id: plan.selected_physical_plan_id.clone().ok_or_else(|| { - CompileError::Lifecycle { + window_implementation_id: plan.selected_window_implementation_id.clone().ok_or_else( + || CompileError::Lifecycle { query_id: query.query_id.clone(), reason: "Planner returned no concrete window implementation identity".into(), - } - })?, + }, + )?, expected_reads: plan.expected_reads.ok_or_else(|| CompileError::Lifecycle { query_id: query.query_id.clone(), reason: "missing joint read demand".into(), @@ -3879,6 +3881,10 @@ fn collect_selected_materializations( SummaryExpr::ValueOperation { child, .. } => { walk(child, readout, composable, grouping.clone(), selected)?; } + SummaryExpr::RelationalJoin { left, right, .. } => { + walk(left, readout, composable, grouping.clone(), selected)?; + walk(right, readout, composable, grouping.clone(), selected)?; + } SummaryExpr::BinaryOp { lhs, rhs, .. } if composable || crate::query_plan::exact_value_executable(node) => { diff --git a/control_plane/src/physical/post_asap/tests.rs b/control_plane/src/physical/post_asap/tests.rs index 745b2ab2..3735861c 100644 --- a/control_plane/src/physical/post_asap/tests.rs +++ b/control_plane/src/physical/post_asap/tests.rs @@ -120,6 +120,7 @@ fn node_is_archive(node: &Rc) -> bool { candidates, values, .. } => node_is_archive(candidates) || node_is_archive(values), SummaryExpr::SummarySubtract { left, right } + | SummaryExpr::RelationalJoin { left, right, .. } | SummaryExpr::BinaryOp { lhs: left, rhs: right, diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index 8d3a8192..4a9741ab 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -636,6 +636,14 @@ pub struct ExternalExactRequest { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(tag = "op", rename_all = "snake_case", deny_unknown_fields)] pub enum QueryPlanNode { + RelationalJoin { + inputs: [QueryNodeId; 2], + join_kind: planner_types::pre_asap::JoinKind, + pred: serde_json::Value, + left_schema: planner_types::post_asap::SummarySchema, + right_schema: planner_types::post_asap::SummarySchema, + output_schema: planner_types::post_asap::SummarySchema, + }, Relational { input: QueryNodeId, /// Serialized planner-owned operation. Keeping the wire form here makes @@ -699,7 +707,7 @@ impl QueryPlanNode { Self::Scalar { .. } | Self::ReadMaterialization { .. } | Self::ExactFallback { .. } => { &[] } - Self::Binary { inputs, .. } => inputs, + Self::Binary { inputs, .. } | Self::RelationalJoin { inputs, .. } => inputs, Self::ReduceSum { input, .. } | Self::Relational { input, .. } | Self::SummaryEstimate { input, .. } @@ -806,7 +814,8 @@ where } } QueryPlanNode::CandidateTopK { inputs, .. } - | QueryPlanNode::Binary { inputs, .. } => { + | QueryPlanNode::Binary { inputs, .. } + | QueryPlanNode::RelationalJoin { inputs, .. } => { for input in inputs { *input = remap[input]; } @@ -872,6 +881,26 @@ where } let physical = match &node.expr { + SummaryExpr::RelationalJoin { + left, + right, + kind, + pred, + } if self.preserve_relational => QueryPlanNode::RelationalJoin { + inputs: [self.lower(left)?, self.lower(right)?], + join_kind: kind.clone(), + pred: serde_json::to_value(pred).map_err(|error| { + QueryPlanError::Invalid(format!( + "cannot serialize relational join predicate: {error}" + )) + })?, + left_schema: left.schema.clone(), + right_schema: right.schema.clone(), + output_schema: node.schema.clone(), + }, + SummaryExpr::RelationalJoin { .. } => QueryPlanNode::ExactFallback { + reason: "read-time relational join requires the relational compiler".into(), + }, SummaryExpr::ValueOperation { child, operation, .. } if self.preserve_relational diff --git a/control_plane/src/query_plan/logical.rs b/control_plane/src/query_plan/logical.rs index 903f2e98..05d47471 100644 --- a/control_plane/src/query_plan/logical.rs +++ b/control_plane/src/query_plan/logical.rs @@ -1289,6 +1289,10 @@ pub fn materialization_candidate_keys( visit(original, lhs, keys)?; visit(original, rhs, keys)?; } + SummaryExpr::RelationalJoin { left, right, .. } => { + visit(original, left, keys)?; + visit(original, right, keys)?; + } SummaryExpr::CandidateTopK { candidates, values, .. } => { diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index 80ac2800..24380517 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -32,4 +32,4 @@ sha2 = "0.10" # exactly (`control_plane/Cargo.toml`) -- two different revs of the same # git dependency in one workspace resolve to two distinct Rust types that # won't unify. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 38619a2b..d5621796 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -38,8 +38,8 @@ control_plane = { path = "../control_plane" } # reduction: Reduction, .. }`) are `pre_asap` types, in the same crate now # (not a separate `asap-ir` import). Query serving consumes the compiled # QueryPlan; these types are used at physical-plan compilation boundaries. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } -asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } # Shared external (workspace) serde.workspace = true @@ -130,7 +130,7 @@ crc32fast = "1.4" # none of them. [dev-dependencies] -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } tempfile = "3.20.0" criterion = { version = "0.5", features = ["html_reports"] } tokio-tungstenite = "0.21" diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs index 61985220..31369947 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs @@ -24,6 +24,102 @@ pub enum ClickHouseDagFallback { }, ResultEncoding(String), } + +fn apply_relational_operation( + operation: serde_json::Value, + output_schema: &planner_types::post_asap::SummarySchema, + relation: ClickHouseRelation, +) -> Result { + let adapter = ClickHouseRelationalAdapter; + if let Some(filter) = operation.get("Filter") { + let predicate = filter + .get("pred") + .cloned() + .ok_or_else(|| "published Filter lacks pred".to_owned()) + .and_then(|value| serde_json::from_value(value).map_err(|error| error.to_string()))?; + return adapter + .apply_filter(&predicate, relation) + .map_err(|error| error.to_string()); + } + let operation: ValueOperation = + serde_json::from_value(operation).map_err(|error| error.to_string())?; + adapter + .apply_operation(&operation, output_schema, relation) + .map_err(|error| error.to_string()) +} + +fn execute_relation_subtree( + index: &SketchStore, + entry: &QueryPlanEntry, + root: QueryNodeId, + expected_schema: &planner_types::post_asap::SummarySchema, + t0_ms: u64, + t1_ms: u64, + is_cumulative: bool, +) -> Result { + match entry.nodes.get(&root) { + Some(QueryPlanNode::Relational { + input, + operation, + input_schema, + output_schema, + }) => { + let input = execute_relation_subtree( + index, + entry, + *input, + input_schema, + t0_ms, + t1_ms, + is_cumulative, + )?; + apply_relational_operation(operation.clone(), output_schema, input) + } + Some(QueryPlanNode::RelationalJoin { + inputs, + join_kind, + pred, + left_schema, + right_schema, + output_schema, + }) => { + if !matches!(join_kind, planner_types::pre_asap::JoinKind::Inner) { + return Err("only inner relational joins are executable".into()); + } + let left = execute_relation_subtree( + index, + entry, + inputs[0], + left_schema, + t0_ms, + t1_ms, + is_cumulative, + )?; + let right = execute_relation_subtree( + index, + entry, + inputs[1], + right_schema, + t0_ms, + t1_ms, + is_cumulative, + )?; + let pred = serde_json::from_value(pred.clone()).map_err(|error| error.to_string())?; + ClickHouseRelationalAdapter + .apply_inner_equi_join(&pred, output_schema, left, right) + .map_err(|error| error.to_string()) + } + Some(_) => { + validate_reachable(entry, root)?; + let outcome = + execute_query_plan_from_readout(index, entry, root, t0_ms, t1_ms, is_cumulative) + .map_err(|error| format!("{error:?}"))?; + ClickHouseRelation::from_series_rows(expected_schema, outcome.series, outcome.coverage) + .map_err(|error| error.to_string()) + } + None => Err(format!("published DAG references missing node {}", root.0)), + } +} pub enum ClickHouseDagOutcome { Accelerated(ClickHouseQueryResult), Fallback(ClickHouseDagFallback), @@ -55,6 +151,59 @@ pub fn execute_sql_dag( error.to_string(), )); } + let has_relational_join = entry.topological_order().is_ok_and(|ids| { + ids.iter().any(|id| { + matches!( + entry.nodes.get(id), + Some(QueryPlanNode::RelationalJoin { .. }) + ) + }) + }); + if has_relational_join { + let root_schema = match entry.nodes.get(&entry.root) { + Some(QueryPlanNode::Relational { output_schema, .. }) + | Some(QueryPlanNode::RelationalJoin { output_schema, .. }) => output_schema, + _ => { + return ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::UnsupportedPlan( + "relational join plan root has no relation schema".into(), + )) + } + }; + let relation = match execute_relation_subtree( + index, + entry, + entry.root, + root_schema, + t0_ms, + t1_ms, + is_cumulative, + ) { + Ok(relation) => relation, + Err(error) => { + return ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::UnsupportedPlan( + error, + )) + } + }; + let pane_ms = entry + .materialization_bindings() + .iter() + .map(|binding| binding.window_ms) + .max() + .unwrap_or(0); + if !complete_pane_coverage(relation.coverage, (t0_ms, t1_ms), pane_ms) { + return ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::IncompleteCoverage { + requested: (t0_ms, t1_ms), + observed: relation.coverage, + }); + } + return match relation.into_result() { + Ok(result) => ClickHouseDagOutcome::Accelerated(result), + Err(error) => ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::ResultEncoding( + error.to_string(), + )), + }; + } let mut base_root = entry.root; let mut relational = Vec::new(); loop { diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs index 69402b65..7ddef54c 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs @@ -93,6 +93,39 @@ impl ClickHouseRelation { pub struct ClickHouseRelationalAdapter; impl ClickHouseRelationalAdapter { + pub fn apply_inner_equi_join( + &self, + pred: &planner_types::pre_asap::Predicate, + output_schema: &SummarySchema, + left: ClickHouseRelation, + right: ClickHouseRelation, + ) -> Result { + let coverage = match (left.coverage, right.coverage) { + (Some((left_start, left_end)), Some((right_start, right_end))) => { + let start = left_start.max(right_start); + let end = left_end.min(right_end); + (start <= end).then_some((start, end)) + } + _ => None, + }; + let mut rows = Vec::new(); + for left_row in &left.rows { + for right_row in &right.rows { + let mut joined = Vec::with_capacity(left_row.len() + right_row.len()); + joined.extend(left_row.iter().cloned()); + joined.extend(right_row.iter().cloned()); + if matches!(eval(&pred.0, &joined)?, Cell::Bool(true)) { + rows.push(joined); + } + } + } + Ok(ClickHouseRelation { + rows, + fields: fields_from_schema(output_schema), + coverage, + }) + } + pub fn apply_filter( &self, pred: &planner_types::pre_asap::Predicate, @@ -646,4 +679,70 @@ mod tests { let error = eval(&QueryExpr::BoolAnd(vec![]), &row).unwrap_err(); assert!(matches!(error, ClickHouseRelationalError::Unsupported(_))); } + + #[test] + fn inner_equi_join_feeds_typed_ratio_projection() { + let side_schema = schema(&[("service", DataType::Utf8), ("value", DataType::Float64)]); + let left = ClickHouseRelation { + rows: vec![vec![Cell::Utf8("api".into()), Cell::Float64(2.0)]], + fields: fields_from_schema(&side_schema), + coverage: Some((300_000, 600_000)), + }; + let right = ClickHouseRelation { + rows: vec![vec![Cell::Utf8("api".into()), Cell::Float64(10.0)]], + fields: fields_from_schema(&side_schema), + coverage: Some((300_000, 600_000)), + }; + let joined_schema = schema(&[ + ("service", DataType::Utf8), + ("left_value", DataType::Float64), + ("service", DataType::Utf8), + ("right_value", DataType::Float64), + ]); + let pred = planner_types::pre_asap::Predicate(Rc::new(QueryExpr::Compare { + left: Rc::new(QueryExpr::Column(0)), + op: CompareOpKind::Eq, + right: Rc::new(QueryExpr::Column(2)), + })); + let joined = ClickHouseRelationalAdapter + .apply_inner_equi_join(&pred, &joined_schema, left, right) + .unwrap(); + let output_schema = schema(&[("service", DataType::Utf8), ("ratio", DataType::Float64)]); + let projected = ClickHouseRelationalAdapter + .apply_operation( + &ValueOperation::Project { + cols: vec![ + ProjectItem { + alias: Some("service".into()), + expr: QueryExpr::Column(0), + }, + ProjectItem { + alias: Some("ratio".into()), + expr: QueryExpr::Arithmetic { + op: ArithmeticOpKind::Div, + left: Rc::new(QueryExpr::Column(1)), + right: Rc::new(QueryExpr::Column(3)), + }, + }, + ], + qualifier: None, + }, + &output_schema, + joined, + ) + .unwrap(); + assert_eq!(projected.coverage, Some((300_000, 600_000))); + let result = projected.into_result().unwrap(); + let batch = &result.batches[0]; + assert_eq!(batch.schema().field(1).name(), "ratio"); + assert_eq!( + batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap() + .value(0), + 0.2 + ); + } } diff --git a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs index 9914b702..a908595d 100644 --- a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs +++ b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs @@ -288,7 +288,8 @@ impl QueryNodeRuntime for PhysicalQueryRuntime<'_> { QueryPlanNode::Logical { .. } | QueryPlanNode::CandidateTopK { .. } | QueryPlanNode::Relational { .. } - | QueryPlanNode::ExternalExact { .. } => Err(PhysicalNodeError::Fallback( + | QueryPlanNode::ExternalExact { .. } + | QueryPlanNode::RelationalJoin { .. } => Err(PhysicalNodeError::Fallback( "logical node requires installed logical runtime".into(), )), QueryPlanNode::ExactFallback { reason } => { diff --git a/data_plane/src/query_engines/asap_query_engine/summary_exec.rs b/data_plane/src/query_engines/asap_query_engine/summary_exec.rs index 5f8c3029..4ab734ea 100644 --- a/data_plane/src/query_engines/asap_query_engine/summary_exec.rs +++ b/data_plane/src/query_engines/asap_query_engine/summary_exec.rs @@ -244,6 +244,7 @@ pub fn execute( } SummaryExpr::SummaryJoin { .. } => Err(ExecError::NotYetSupported("SummaryJoin")), + SummaryExpr::RelationalJoin { .. } => Err(ExecError::NotYetSupported("RelationalJoin")), // CandidateTopK is lowered to the deployed QueryPlan DAG, where both // row inputs retain labels for intersection and exact reranking. This // legacy generic adapter exposes opaque GroupKey values and cannot From d661dc6967a394bb89111f0ed3046ae42142cef5 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 09:00:55 -0600 Subject: [PATCH 4/5] fix(clickhouse): validate each join input coverage --- .../asap_clickhouse_query_engine/execution.rs | 48 ++++++++++++++++++- 1 file changed, 47 insertions(+), 1 deletion(-) diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs index 31369947..c4424e2b 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs @@ -113,7 +113,47 @@ fn execute_relation_subtree( validate_reachable(entry, root)?; let outcome = execute_query_plan_from_readout(index, entry, root, t0_ms, t1_ms, is_cumulative) - .map_err(|error| format!("{error:?}"))?; + .map_err(|error| format!("incomplete leaf coverage: {error:?}"))?; + let reachable = entry + .topological_order_from(root) + .map_err(|error| format!("invalid leaf DAG: {error}"))?; + for (leaf_id, binding) in reachable.iter().filter_map(|id| match entry.nodes.get(id) { + Some(QueryPlanNode::ReadMaterialization { binding }) => Some((*id, binding)), + _ => None, + }) { + let leaf_outcome = execute_query_plan_from_readout( + index, + entry, + leaf_id, + t0_ms, + t1_ms, + is_cumulative, + ) + .map_err(|error| format!("incomplete leaf coverage: {error:?}"))?; + let origin = binding + .pane_origin_ms + .ok_or_else(|| "incomplete leaf coverage: missing pane origin".to_owned())?; + let start = i64::try_from(t0_ms) + .map_err(|_| "incomplete leaf coverage: start exceeds i64".to_owned())?; + let end = i64::try_from(t1_ms) + .map_err(|_| "incomplete leaf coverage: end exceeds i64".to_owned())?; + let pane = i64::try_from(binding.window_ms) + .map_err(|_| "incomplete leaf coverage: pane exceeds i64".to_owned())?; + if pane <= 0 + || (start - origin).rem_euclid(pane) != 0 + || (end - origin).rem_euclid(pane) != 0 + || !complete_pane_coverage( + leaf_outcome.coverage, + (t0_ms, t1_ms), + binding.window_ms, + ) + { + return Err(format!( + "incomplete leaf coverage: requested ({t0_ms}, {t1_ms}), observed {:?}, pane {} origin {}", + leaf_outcome.coverage, binding.window_ms, origin + )); + } + } ClickHouseRelation::from_series_rows(expected_schema, outcome.series, outcome.coverage) .map_err(|error| error.to_string()) } @@ -179,6 +219,12 @@ pub fn execute_sql_dag( is_cumulative, ) { Ok(relation) => relation, + Err(error) if error.contains("incomplete leaf coverage") => { + return ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::IncompleteCoverage { + requested: (t0_ms, t1_ms), + observed: None, + }) + } Err(error) => { return ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::UnsupportedPlan( error, From 9543e52431af72d41e9a5563cbc36e9830a5399c Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 11:47:06 -0600 Subject: [PATCH 5/5] chore: pin merged planner join contract --- Cargo.lock | 26 +++++++++++++------------- control_plane/Cargo.toml | 10 +++++----- crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +++--- 4 files changed, 22 insertions(+), 22 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 96eaa3f9..ecca1c75 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -374,7 +374,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3dda28032c2a2badb64bf5358e0f40fef9918227#3dda28032c2a2badb64bf5358e0f40fef9918227" dependencies = [ "asap-types", "serde", @@ -385,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-metricsql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3dda28032c2a2badb64bf5358e0f40fef9918227#3dda28032c2a2badb64bf5358e0f40fef9918227" dependencies = [ "asap-types", "metricsql_parser", @@ -395,7 +395,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3dda28032c2a2badb64bf5358e0f40fef9918227#3dda28032c2a2badb64bf5358e0f40fef9918227" dependencies = [ "asap-types", "promql-parser 0.10.0", @@ -404,7 +404,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3dda28032c2a2badb64bf5358e0f40fef9918227#3dda28032c2a2badb64bf5358e0f40fef9918227" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -426,12 +426,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3dda28032c2a2badb64bf5358e0f40fef9918227#3dda28032c2a2badb64bf5358e0f40fef9918227" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3dda28032c2a2badb64bf5358e0f40fef9918227#3dda28032c2a2badb64bf5358e0f40fef9918227" dependencies = [ "serde", "serde_json", @@ -1711,7 +1711,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -2453,7 +2453,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" dependencies = [ "hermit-abi 0.5.2", "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -2801,7 +2801,7 @@ dependencies = [ [[package]] name = "metricsql_common" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3dda28032c2a2badb64bf5358e0f40fef9918227#3dda28032c2a2badb64bf5358e0f40fef9918227" dependencies = [ "chrono", ] @@ -2809,7 +2809,7 @@ dependencies = [ [[package]] name = "metricsql_parser" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3dda28032c2a2badb64bf5358e0f40fef9918227#3dda28032c2a2badb64bf5358e0f40fef9918227" dependencies = [ "ahash", "chrono", @@ -3904,7 +3904,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -4438,7 +4438,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix 1.1.4", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -5260,7 +5260,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.48.0", + "windows-sys 0.61.2", ] [[package]] diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 6a30d061..71c9a282 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -91,8 +91,8 @@ asap_types.workspace = true # scaffolding, unaware that `data_plane`'s `summary_executor.rs` in *this* # repo is a real one. Vendored locally instead of chased upstream -- see # `data_plane/src/query_engines/asap_query_engine/summary_exec.rs`. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "3dda28032c2a2badb64bf5358e0f40fef9918227" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "3dda28032c2a2badb64bf5358e0f40fef9918227" } # L1 adoption (design-target-architecture.md Part B): the PromQL front # end itself, replacing control_plane's own query_parser/promql.rs. @@ -100,9 +100,9 @@ asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = # `planner-types`/`asap-aware-mapping` above -- these three MUST move # together (two revs of the same upstream repo's types in one workspace # resolve to distinct Rust types that won't unify). -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } -asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "3dda28032c2a2badb64bf5358e0f40fef9918227" } +asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "3dda28032c2a2badb64bf5358e0f40fef9918227" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "3dda28032c2a2badb64bf5358e0f40fef9918227" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index 24380517..b533726f 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -32,4 +32,4 @@ sha2 = "0.10" # exactly (`control_plane/Cargo.toml`) -- two different revs of the same # git dependency in one workspace resolve to two distinct Rust types that # won't unify. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "3dda28032c2a2badb64bf5358e0f40fef9918227" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index d5621796..808e1d7c 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -38,8 +38,8 @@ control_plane = { path = "../control_plane" } # reduction: Reduction, .. }`) are `pre_asap` types, in the same crate now # (not a separate `asap-ir` import). Query serving consumes the compiled # QueryPlan; these types are used at physical-plan compilation boundaries. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } -asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "3dda28032c2a2badb64bf5358e0f40fef9918227" } +asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "3dda28032c2a2badb64bf5358e0f40fef9918227" } # Shared external (workspace) serde.workspace = true @@ -130,7 +130,7 @@ crc32fast = "1.4" # none of them. [dev-dependencies] -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "3dda28032c2a2badb64bf5358e0f40fef9918227" } tempfile = "3.20.0" criterion = { version = "0.5", features = ["html_reports"] } tokio-tungstenite = "0.21"