From ea000d82890cbc4e42184132fe479975e0278157 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 18 Sep 2026 15:25:32 +0000 Subject: [PATCH 1/4] feat(control-plane): lower PromQL with ingestion cadence --- Cargo.lock | 54 +++++--- Cargo.toml | 11 +- control_plane/src/physical/compiler.rs | 116 ++++++++++++++++-- control_plane/src/query_parser/mod.rs | 44 +++++-- ...asapquery-compatibility-demo-snapshot.json | 2 + .../examples/asapquery-planning-snapshot.json | 2 + 6 files changed, 188 insertions(+), 41 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index e173aa68..da10811a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -364,9 +364,10 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=029ff2fe041172c94c2d32c90b185bc83c5e8a57#029ff2fe041172c94c2d32c90b185bc83c5e8a57" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" dependencies = [ "asap-types", + "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?rev=da3635a80f8f854d47b772d49d5a9e5fb6927d8e)", "serde", "serde_json", "thiserror 2.0.20", @@ -375,7 +376,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=029ff2fe041172c94c2d32c90b185bc83c5e8a57#029ff2fe041172c94c2d32c90b185bc83c5e8a57" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" dependencies = [ "asap-types", "promql-parser 0.10.0 (git+https://github.com/ProjectASAP/promql-parser?rev=9fede7eecca923c9882fe256484d00d37f8706cb)", @@ -384,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=029ff2fe041172c94c2d32c90b185bc83c5e8a57#029ff2fe041172c94c2d32c90b185bc83c5e8a57" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -397,7 +398,7 @@ name = "asap-precompute-rs" version = "0.1.0" source = "git+https://github.com/ProjectASAP/ASAPCollector?branch=main#1d8efd07e40fc151cbd4678a5c6aa9774b1aed34" dependencies = [ - "asap_sketchlib", + "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?branch=main)", "prost", "serde", "serde_json", @@ -407,12 +408,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=029ff2fe041172c94c2d32c90b185bc83c5e8a57#029ff2fe041172c94c2d32c90b185bc83c5e8a57" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=029ff2fe041172c94c2d32c90b185bc83c5e8a57#029ff2fe041172c94c2d32c90b185bc83c5e8a57" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" dependencies = [ "serde", "serde_json", @@ -447,6 +448,23 @@ dependencies = [ "xxhash-rust", ] +[[package]] +name = "asap_sketchlib" +version = "0.3.0" +source = "git+https://github.com/ProjectASAP/asap_sketchlib?rev=da3635a80f8f854d47b772d49d5a9e5fb6927d8e#da3635a80f8f854d47b772d49d5a9e5fb6927d8e" +dependencies = [ + "bytes", + "prost", + "rand 0.9.5", + "rmp-serde", + "serde", + "serde-big-array", + "serde_bytes", + "smallvec", + "twox-hash 2.1.4", + "xxhash-rust", +] + [[package]] name = "asap_types" version = "0.1.0" @@ -1144,7 +1162,7 @@ dependencies = [ "asap-precompute-rs", "asap-types", "asap_otel_proto", - "asap_sketchlib", + "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?branch=main)", "asap_types", "async-trait", "axum", @@ -1642,7 +1660,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]] @@ -2073,7 +2091,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.6.5", + "socket2 0.5.10", "tokio", "tower-service", "tracing", @@ -2259,7 +2277,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" dependencies = [ "hermit-abi", "libc", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -2606,7 +2624,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -3184,7 +3202,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls", - "socket2 0.6.5", + "socket2 0.5.10", "thiserror 2.0.20", "tokio", "tracing", @@ -3222,9 +3240,9 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.6.5", + "socket2 0.5.10", "tracing", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -3500,7 +3518,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -3945,7 +3963,7 @@ dependencies = [ "getrandom 0.4.3", "once_cell", "rustix 1.1.4", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -4457,7 +4475,7 @@ version = "2.1.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5283634e518fe9e82c7b20520bb4bc209009fd16c82077c802f8111ecbb0117a" dependencies = [ - "rand 0.10.2", + "rand 0.9.5", ] [[package]] @@ -4720,7 +4738,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.52.0", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index ae7e90f8..4bfbbfb8 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -18,12 +18,12 @@ version = "0.1.0" asap_sketchlib = { git = "https://github.com/ProjectASAP/asap_sketchlib", branch = "main" } [workspace.dependencies] -# Keep Planner frontends, selection, and IR on the same immutable revision (current-series Planner PR). +# Keep Planner frontends, selection, and IR on the same immutable revision. # Alias upstream asap-types because this workspace also defines asap_types. -planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "029ff2fe041172c94c2d32c90b185bc83c5e8a57" } -asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "029ff2fe041172c94c2d32c90b185bc83c5e8a57" } -asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "029ff2fe041172c94c2d32c90b185bc83c5e8a57" } -asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "029ff2fe041172c94c2d32c90b185bc83c5e8a57" } +planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "b56d0e2d21c878169ee26c628646bd2b112b10a1" } +asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "b56d0e2d21c878169ee26c628646bd2b112b10a1" } +asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "b56d0e2d21c878169ee26c628646bd2b112b10a1" } +asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "b56d0e2d21c878169ee26c628646bd2b112b10a1" } # Shared external deps (used by 2+ crates) serde = { version = "1.0", features = ["derive"] } @@ -46,4 +46,3 @@ reqwest = { version = "0.12", default-features = false, features = ["json", "rus asap_types = { path = "crates/asap_types" } asap_otel_proto = { path = "crates/asap_otel_proto" } indexmap = { version = "2.0", features = ["serde"] } - diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index c55d2a8c..d82c3c62 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -587,7 +587,18 @@ impl BackendLocalPlanningInput { )); } } - workload.data_workload = Some(self.data_workload.clone()); + let mut data_workload = self.data_workload.clone(); + if data_workload.data_ingestion_interval.value.is_none() + && self.physical_inputs.scrape_interval_ms > 0 + { + data_workload.data_ingestion_interval = Evidence { + value: Some(DurationMs(self.physical_inputs.scrape_interval_ms)), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }; + } + workload.data_workload = Some(data_workload.clone()); workload .validate() .map_err(|error| CompileError::Snapshot(error.to_string()))?; @@ -596,8 +607,7 @@ impl BackendLocalPlanningInput { "ASAPQuery compatibility profile accepts PromQL workloads only".into(), )); } - let ingestion_rate = self - .data_workload + let ingestion_rate = data_workload .ingestion_rate .value_at(self.environment.observed_at_unix_ms) .copied() @@ -610,10 +620,12 @@ impl BackendLocalPlanningInput { "QueryWorkload must contain at least one query".into(), )); } + let canonical_queries = asap_frontend_promql::lower_promql_workload(&workload) + .map_err(|error| CompileError::Snapshot(error.to_string()))?; let mut queries = Vec::with_capacity(entries.len()); let mut canonical_roots = Vec::with_capacity(entries.len()); let mut topk_evidence_by_id = HashMap::new(); - for (index, entry) in entries.into_iter().enumerate() { + for (index, (entry, parsed)) in entries.into_iter().zip(canonical_queries).enumerate() { let evaluation_interval_ms = match entry.recurrence { QueryRecurrence::Repeated(RepeatedDemand::FixedIntervalAt { interval, .. @@ -636,9 +648,6 @@ impl BackendLocalPlanningInput { } let accuracy = entry.requirements.accuracy.target(); let query_string = entry.query.0; - let parsed = - crate::query_parser::parse_query_expr_canonical(&query_string, accuracy.clone()) - .map_err(|error| CompileError::Snapshot(format!("query {index}: {error}")))?; let lookback_ms = query_history_window_ms(&parsed, self.physical_inputs.scrape_interval_ms) .map_err(|error| CompileError::Snapshot(format!("query {index}: {error}")))?; @@ -2934,6 +2943,15 @@ fn select_lifecycle( ), data_workload: Some(DataWorkload { arrival: DataArrival::ContinuouslyIngesting, + // This workload is an internal lifecycle projection, not a + // plan-ready query input. Use the compatibility profile cadence + // when there is no original workload to carry source evidence. + data_ingestion_interval: Evidence { + value: Some(DurationMs(1_000)), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }, ingestion_rate: Evidence { value: Some(Rate( query.summary_lifecycle_inputs.ingestion_rate_per_second, @@ -2952,6 +2970,21 @@ fn select_lifecycle( // evidence comes from the same lifecycle input as window pricing. if original.data_workload.is_none() { original.data_workload = workload.data_workload; + } else if original + .data_workload + .as_ref() + .is_some_and(|data| data.data_ingestion_interval.value.is_none()) + { + original + .data_workload + .as_mut() + .unwrap() + .data_ingestion_interval = workload + .data_workload + .as_ref() + .unwrap() + .data_ingestion_interval + .clone(); } (original, indices) } @@ -5036,7 +5069,8 @@ pub(crate) mod tests { let query = "mad_over_time(m[1m])"; let mut workload = request("vm-q", "last_over_time(m[1m])"); let accuracy = workload.queries[0].accuracy_target.clone(); - let canonical = asap_frontend_promql::lower_promql(query, accuracy.clone()).unwrap(); + let canonical = + crate::query_parser::parse_query_expr_canonical(query, accuracy.clone()).unwrap(); workload.queries[0].query_string = query.into(); workload.queries[0].selected_plan_root = crate::planner_selection::keep_pre_asap(&canonical).unwrap(); @@ -6179,6 +6213,72 @@ pub(crate) mod tests { .query_lookback_seconds } + fn set_data_ingestion_interval( + snapshot: &mut BackendLocalPlanningInput, + interval_ms: Option, + ) { + let evidence = Evidence { + value: interval_ms.map(DurationMs), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }; + snapshot.data_workload.data_ingestion_interval = evidence.clone(); + snapshot + .query_workload + .data_workload + .as_mut() + .unwrap() + .data_ingestion_interval = evidence; + } + + // Plan-ready lowering uses workload cadence, not the backend-local migration field. + #[test] + fn instant_aggregate_uses_data_ingestion_interval() { + let mut snapshot = planning_snapshot(); + snapshot.query_workload.repeating_queries.as_mut().unwrap()[0].query = + Query("sum(data)".into()); + set_data_ingestion_interval(&mut snapshot, Some(1_000)); + + let (request, _) = snapshot.into_physical_compilation_request().unwrap(); + assert_eq!(request.queries[0].query_lookback_seconds, 1); + } + + // Legacy snapshots migrate their explicit scrape cadence into DataWorkload. + #[test] + fn missing_data_ingestion_interval_migrates_scrape_interval() { + let mut snapshot = planning_snapshot(); + snapshot.query_workload.repeating_queries.as_mut().unwrap()[0].query = + Query("sum(data)".into()); + set_data_ingestion_interval(&mut snapshot, None); + + let (request, _) = snapshot.into_physical_compilation_request().unwrap(); + assert_eq!( + request + .query_workload + .unwrap() + .data_workload + .unwrap() + .data_ingestion_interval + .value, + Some(DurationMs(5_000)) + ); + assert_eq!(request.queries[0].query_lookback_seconds, 5); + } + + // An explicitly invalid interval must not be replaced by the migration value. + #[test] + fn zero_data_ingestion_interval_is_rejected() { + let mut snapshot = planning_snapshot(); + set_data_ingestion_interval(&mut snapshot, Some(0)); + + let error = snapshot + .into_physical_compilation_request() + .expect_err("zero interval must fail") + .to_string(); + assert!(error.contains("data_ingestion_interval must be greater than zero")); + } + // Every concatenated histogram result contributes its source history. #[test] fn derived_history_visits_concat_branches() { diff --git a/control_plane/src/query_parser/mod.rs b/control_plane/src/query_parser/mod.rs index 9853800a..eb2883e1 100644 --- a/control_plane/src/query_parser/mod.rs +++ b/control_plane/src/query_parser/mod.rs @@ -1,4 +1,4 @@ -//! PromQL workload extraction through `asap_frontend_promql::lower_promql`. +//! PromQL workload extraction through `asap_frontend_promql::lower_promql_workload`. //! //! ASAPPlanner owns parsing and intent classification. Both entry points require //! an explicit `AccuracyTarget` because lowering can size summaries. @@ -14,6 +14,12 @@ use crate::types::AggType; use planner_types::pre_asap::AggIntent; use planner_types::pre_asap::{ColumnId, CompareOpKind, ScalarValue}; use planner_types::pre_asap::{Predicate, QueryExpr, Source}; +use planner_types::workload::{ + AccuracyRequirement, BatchEntry, DataWorkload, DurationMs, Evidence, Query, QueryLanguage, + QueryRequirements, QueryWorkload, TimeSelection, +}; + +const COMPATIBILITY_DATA_INGESTION_INTERVAL_MS: u64 = 1_000; // ── Output types (legacy — consumed by analyzer and planner) ────────────────── @@ -68,18 +74,38 @@ pub enum QueryHint { /// Parse a PromQL query string into the **canonical** L3 /// [`query_expr::QueryExpr`](planner_types::pre_asap::QueryExpr) IR. /// -/// This is the single public algebra-IR entry point — a direct call into -/// `asap_frontend_promql::lower_promql`, which does the full L1 parse → -/// L2 relational tree → L3 canonical conversion in one call. No local -/// parser, no local L2 tree; `accuracy` is threaded onto every -/// accuracy-bearing intent the same way ASAPController's own PromQL -/// front end threads it. +/// Compatibility entry point for callers that do not own a complete workload. +/// Plan-ready compilation lowers the complete caller-supplied workload instead. pub fn parse_query_expr_canonical( query: &str, accuracy: AccuracyTarget, ) -> anyhow::Result { - let canonical = asap_frontend_promql::lower_promql(query.trim(), accuracy)?; - Ok(canonical) + let workload = QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: Some(vec![BatchEntry { + query: Query(query.trim().into()), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit(accuracy), + ..Default::default() + }, + predictability: Default::default(), + invocations: 1, + execute_at: None, + time_selection: TimeSelection::default(), + }]), + repeating_queries: None, + data_workload: Some(DataWorkload { + data_ingestion_interval: Evidence { + value: Some(DurationMs(COMPATIBILITY_DATA_INGESTION_INTERVAL_MS)), + ..Default::default() + }, + ..Default::default() + }), + }; + asap_frontend_promql::lower_promql_workload(&workload)? + .into_iter() + .next() + .ok_or_else(|| anyhow::anyhow!("PromQL compatibility workload produced no query")) } /// Parse a PromQL query string into a [`ParsedQuery`]. diff --git a/docs/examples/asapquery-compatibility-demo-snapshot.json b/docs/examples/asapquery-compatibility-demo-snapshot.json index 0c466eff..f617f8ce 100644 --- a/docs/examples/asapquery-compatibility-demo-snapshot.json +++ b/docs/examples/asapquery-compatibility-demo-snapshot.json @@ -58,6 +58,7 @@ ], "data_workload": { "arrival": "continuously_ingesting", + "data_ingestion_interval": { "value": 5000, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, "ingestion_volume": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, "ingestion_rate": { "value": 100.0, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, "input_cardinality": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, @@ -66,6 +67,7 @@ }, "data_workload": { "arrival": "continuously_ingesting", + "data_ingestion_interval": { "value": 5000, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, "ingestion_volume": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, "ingestion_rate": { "value": 100.0, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, "input_cardinality": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, diff --git a/docs/examples/asapquery-planning-snapshot.json b/docs/examples/asapquery-planning-snapshot.json index 4f0ed030..39017909 100644 --- a/docs/examples/asapquery-planning-snapshot.json +++ b/docs/examples/asapquery-planning-snapshot.json @@ -25,6 +25,7 @@ ], "data_workload": { "arrival": "continuously_ingesting", + "data_ingestion_interval": { "value": 5000, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, "ingestion_volume": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, "ingestion_rate": { "value": 100.0, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, "input_cardinality": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, @@ -33,6 +34,7 @@ }, "data_workload": { "arrival": "continuously_ingesting", + "data_ingestion_interval": { "value": 5000, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, "ingestion_volume": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, "ingestion_rate": { "value": 100.0, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, "input_cardinality": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, From 661cc1251cbf44a175ebbd3e79c7e6969ea028c4 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 18 Sep 2026 17:04:36 +0000 Subject: [PATCH 2/4] Integrate peer workloads and preserve planner ingestion horizons --- Cargo.lock | 10 +- Cargo.toml | 8 +- control_plane/src/main.rs | 90 ++++-- control_plane/src/physical/compiler.rs | 285 ++++++++++-------- .../src/physical/maintained_population.rs | 3 +- control_plane/src/physical/post_asap/lower.rs | 6 + control_plane/src/physical/workload_cost.rs | 16 +- control_plane/src/query_parser/mod.rs | 42 ++- control_plane/src/query_plan/residual.rs | 178 ++++++++--- control_plane/tests/discovery_snapshot.rs | 4 + .../src/query_plan/current_series.rs | 33 +- crates/asap_types/src/table_population.rs | 6 +- .../relational_adapter.rs | 21 ++ .../sketch_db/backfill/clickhouse_reader.rs | 8 +- .../sketch_db/current_series.rs | 22 ++ .../tests/support/current_series_process.rs | 34 ++- ...asapquery-compatibility-demo-snapshot.json | 10 +- .../examples/asapquery-planning-snapshot.json | 10 +- docs/user_guide/asapquery-profile.md | 9 + tools/o11y-execution/discover_snapshot.py | 5 +- 20 files changed, 552 insertions(+), 248 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index da10811a..2f8454a8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -364,7 +364,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?rev=da3635a80f8f854d47b772d49d5a9e5fb6927d8e)", @@ -376,7 +376,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" dependencies = [ "asap-types", "promql-parser 0.10.0 (git+https://github.com/ProjectASAP/promql-parser?rev=9fede7eecca923c9882fe256484d00d37f8706cb)", @@ -385,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -408,12 +408,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b56d0e2d21c878169ee26c628646bd2b112b10a1#b56d0e2d21c878169ee26c628646bd2b112b10a1" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" dependencies = [ "serde", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index 4bfbbfb8..15c27b38 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -20,10 +20,10 @@ asap_sketchlib = { git = "https://github.com/ProjectASAP/asap_sketchlib", branch [workspace.dependencies] # Keep Planner frontends, selection, and IR on the same immutable revision. # Alias upstream asap-types because this workspace also defines asap_types. -planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "b56d0e2d21c878169ee26c628646bd2b112b10a1" } -asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "b56d0e2d21c878169ee26c628646bd2b112b10a1" } -asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "b56d0e2d21c878169ee26c628646bd2b112b10a1" } -asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "b56d0e2d21c878169ee26c628646bd2b112b10a1" } +planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" } +asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" } +asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" } +asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" } # Shared external deps (used by 2+ crates) serde = { version = "1.0", features = ["derive"] } diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index bb77234e..ef76bf73 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -198,6 +198,7 @@ struct CompileAndPublishPhysicalPlanRequest { #[serde(default)] workload_cost_evidence: Option, queries: Vec, + data_workload: planner_types::workload::DataWorkload, #[serde(rename = "collector_ids", alias = "target_collector_ids")] target_collector_ids: Vec, capability_snapshot_id: String, @@ -505,11 +506,7 @@ fn compile_physical_plan_request( .duration_since(std::time::UNIX_EPOCH) .unwrap_or_default() .as_millis() as u64; - let mut queries = Vec::with_capacity(request.queries.len()); - let mut canonical_roots = Vec::with_capacity(request.queries.len()); - let mut window_models = Vec::new(); - let mut workload_entries = Vec::new(); - for query in request.queries { + for query in &request.queries { if query.query_id.trim().is_empty() || query.metric.trim().is_empty() || query.window_secs == 0 @@ -521,17 +518,11 @@ fn compile_physical_plan_request( .into(), )); } - let expr = match frontend.parse(&query.query_string, query.accuracy.clone()) { - Ok(expr) => expr, - Err(error) => return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into())), - }; - let post_asap = match control_plane::planner_selection::keep_pre_asap(&expr) { - Ok(plan) => plan, - Err(error) => return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into())), - }; - canonical_roots.push(std::rc::Rc::new(expr)); - window_models.push(query.window_cost_model); - workload_entries.push(planner_types::workload::RepeatingEntry { + } + let workload_entries = request + .queries + .iter() + .map(|query| planner_types::workload::RepeatingEntry { query: planner_types::workload::Query(query.query_string.clone()), demand: planner_types::workload::RepeatedDemand::FixedIntervalAt { interval: planner_types::workload::RepetitionInterval( @@ -552,7 +543,44 @@ fn compile_physical_plan_request( )), ..Default::default() }, - }); + }) + .collect(); + let query_workload = planner_types::workload::QueryWorkload { + language: planner_types::workload::QueryLanguage::PromQL, + query_batch: None, + repeating_queries: Some(workload_entries), + }; + request + .data_workload + .validate() + .map_err(|error| (StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into()))?; + let lowered = match frontend { + QueryFrontend::PromQl => asap_frontend_promql::lower_promql_workload( + &planner_types::workload::PlanningWorkload { + query_workload: query_workload.clone(), + data_workload: Some(request.data_workload.clone()), + }, + now, + ) + .map_err(|error| (StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into()))?, + QueryFrontend::MetricsQl => request + .queries + .iter() + .map(|query| frontend.parse(&query.query_string, query.accuracy.clone())) + .collect::, _>>() + .map_err(|error| (StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into()))?, + }; + + let mut queries = Vec::with_capacity(request.queries.len()); + let mut canonical_roots = Vec::with_capacity(request.queries.len()); + let mut window_models = Vec::new(); + for (query, expr) in request.queries.into_iter().zip(lowered) { + let post_asap = match control_plane::planner_selection::keep_pre_asap(&expr) { + Ok(plan) => plan, + Err(error) => return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into())), + }; + canonical_roots.push(std::rc::Rc::new(expr)); + window_models.push(query.window_cost_model); queries.push(physical::compiler::QueryCompilationInput { query_id: query.query_id, query_string: query.query_string, @@ -571,7 +599,7 @@ fn compile_physical_plan_request( let planner_selection_trace = match physical::compiler::select_logical_roots_with_trace( &mut queries, - canonical_roots, + canonical_roots.clone(), &request.evidence, &request.exact_composition_costs, request.erp.as_ref(), @@ -586,12 +614,9 @@ fn compile_physical_plan_request( } let compilation_request = physical::compiler::PhysicalCompilationRequest { planner_selection_trace, - query_workload: Some(planner_types::workload::QueryWorkload { - language: planner_types::workload::QueryLanguage::PromQL, - query_batch: None, - repeating_queries: Some(workload_entries), - data_workload: None, - }), + query_workload: Some(query_workload), + data_workload: Some(request.data_workload), + canonical_roots, queries, allow_mixed_summary_and_exact_execution: request.target == physical::compiler::PhysicalDeploymentTarget::BackendLocalRemoteWrite, @@ -820,6 +845,7 @@ mod api_tests { .unwrap() .as_millis() as u64; let mut request_body = serde_json::json!({ + "data_workload": snapshot.data_workload, "queries": [{ "query_id": query.query_id, "query_string": query.query_string, "metric": "m", "window_secs": 60, "accuracy": query.accuracy_target, @@ -895,6 +921,7 @@ mod api_tests { panic!("expected time series fixture"); }; let value = serde_json::json!({ + "data_workload": snapshot.data_workload, "target": "backend_local_remote_write", "queries": [{ "query_id": query.query_id, "query_string": query.query_string, @@ -917,6 +944,21 @@ mod api_tests { physical::compiler::IngestProtocol::PrometheusRemoteWriteV1 ); assert!(!plan.precompute_plan.materializations.is_empty()); + // HTTP lowering uses the current planning clock, not a timeless compatibility parse. + let mut stale = value.clone(); + stale["data_workload"]["data_ingestion_interval"]["observed_at_ms"] = serde_json::json!(0); + stale["data_workload"]["data_ingestion_interval"]["valid_for_ms"] = serde_json::json!(0); + let error = compile_physical_plan_request( + serde_json::from_value(stale).unwrap(), + false, + QueryFrontend::PromQl, + ) + .expect_err("expired cadence must fail HTTP planning"); + assert_eq!(error.0, StatusCode::UNPROCESSABLE_ENTITY); + assert!(error + .1 + .to_string() + .contains("evidence is unavailable at planning time")); let mut distributed = value; distributed["target"] = serde_json::json!("distributed_collectors"); assert!(compile_physical_plan_request( diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index d82c3c62..4b9dd1d5 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -163,6 +163,10 @@ pub struct PhysicalCompilationRequest { pub enabled_materialization_keys: Option>, /// Original dashboard demand, in the same order as queries. None is legacy input. pub query_workload: Option, + /// Source evidence supplied independently from query demand. + pub data_workload: Option, + /// Workload-lowered roots retained across physical alternative enumeration. + pub canonical_roots: Vec>, pub queries: Vec, pub topk_membership_evidence_by_query_id: HashMap, /// Fresh measured costs for Planner exact/summary composition sites, @@ -579,14 +583,7 @@ impl BackendLocalPlanningInput { "compatibility workload requires backend_local_remote_write target".into(), )); } - let mut workload = self.query_workload; - if let Some(embedded) = &workload.data_workload { - if embedded != &self.data_workload { - return Err(CompileError::Snapshot( - "embedded and standalone DataWorkload snapshots disagree".into(), - )); - } - } + let workload = self.query_workload; let mut data_workload = self.data_workload.clone(); if data_workload.data_ingestion_interval.value.is_none() && self.physical_inputs.scrape_interval_ms > 0 @@ -598,7 +595,9 @@ impl BackendLocalPlanningInput { valid_for_ms: None, }; } - workload.data_workload = Some(data_workload.clone()); + data_workload + .validate() + .map_err(|error| CompileError::Snapshot(error.to_string()))?; workload .validate() .map_err(|error| CompileError::Snapshot(error.to_string()))?; @@ -620,8 +619,14 @@ impl BackendLocalPlanningInput { "QueryWorkload must contain at least one query".into(), )); } - let canonical_queries = asap_frontend_promql::lower_promql_workload(&workload) - .map_err(|error| CompileError::Snapshot(error.to_string()))?; + let canonical_queries = asap_frontend_promql::lower_promql_workload( + &planner_types::workload::PlanningWorkload { + query_workload: workload.clone(), + data_workload: Some(data_workload.clone()), + }, + self.environment.observed_at_unix_ms, + ) + .map_err(|error| CompileError::Snapshot(error.to_string()))?; let mut queries = Vec::with_capacity(entries.len()); let mut canonical_roots = Vec::with_capacity(entries.len()); let mut topk_evidence_by_id = HashMap::new(); @@ -709,7 +714,7 @@ impl BackendLocalPlanningInput { } let planner_selection_trace = select_logical_roots_with_trace( &mut queries, - canonical_roots, + canonical_roots.clone(), &topk_evidence_by_id, &exact_costs_by_id, self.physical_inputs.erp.as_ref(), @@ -729,6 +734,8 @@ impl BackendLocalPlanningInput { allow_mixed_summary_and_exact_execution: true, enabled_materialization_keys: None, query_workload: Some(workload), + data_workload: Some(data_workload), + canonical_roots, queries, topk_membership_evidence_by_query_id: topk_evidence_by_id, exact_composition_costs: exact_costs_by_id, @@ -810,11 +817,30 @@ fn has_unsafe_raw_entity_leaf( } } +fn original_root( + query: &QueryCompilationInput, + index: usize, + roots: &[Rc], +) -> Result { + if let Some(root) = roots.get(index) { + return Ok(root.as_ref().clone()); + } + crate::query_parser::parse_query_expr_canonical( + &query.query_string, + query.accuracy_target.clone(), + ) + .map_err(|error| CompileError::Query { + query_id: query.query_id.clone(), + reason: error.to_string(), + }) +} + /// Preserve native execution for raw states that cannot preserve source semantics. fn preserve_native_unsafe_raw_roots( queries: &mut [QueryCompilationInput], + canonical_roots: &[Rc], ) -> Result<(), CompileError> { - for query in queries { + for (index, query) in queries.iter_mut().enumerate() { let selected = collect_selected_materializations(&query.selected_plan_root, false) .map_err(|reason| CompileError::Query { query_id: query.query_id.clone(), @@ -830,14 +856,7 @@ fn preserve_native_unsafe_raw_roots( let unsafe_entities = has_unsafe_raw_entity_leaf(&query.selected_plan_root, &selected_nodes, false); if unsafe_entities { - let parsed = crate::query_parser::parse_query_expr_canonical( - &query.query_string, - query.accuracy_target.clone(), - ) - .map_err(|error| CompileError::Query { - query_id: query.query_id.clone(), - reason: error.to_string(), - })?; + let parsed = original_root(query, index, canonical_roots)?; query.selected_plan_root = crate::planner_selection::keep_pre_asap(&parsed).map_err(|error| { CompileError::Query { @@ -856,9 +875,10 @@ fn preserve_native_unsafe_raw_roots( /// maintenance DAG), but the original query remains a valid exact plan. fn preserve_invalid_exact_fallback_roots( queries: &mut [QueryCompilationInput], + canonical_roots: &[Rc], composable: bool, ) -> Result<(), CompileError> { - for query in queries { + for (index, query) in queries.iter_mut().enumerate() { let selected = collect_selected_materializations(&query.selected_plan_root, composable) .map_err(|reason| CompileError::Query { query_id: query.query_id.clone(), @@ -869,14 +889,7 @@ fn preserve_invalid_exact_fallback_roots( if invalid_executable && !matches!(query.selected_plan_root.expr, SummaryExpr::KeepPreAsap(_)) { - let parsed = crate::query_parser::parse_query_expr_canonical( - &query.query_string, - query.accuracy_target.clone(), - ) - .map_err(|error| CompileError::Query { - query_id: query.query_id.clone(), - reason: error.to_string(), - })?; + let parsed = original_root(query, index, canonical_roots)?; query.selected_plan_root = crate::planner_selection::keep_pre_asap(&parsed).map_err(|error| { CompileError::Query { @@ -895,9 +908,10 @@ fn preserve_invalid_exact_fallback_roots( /// summaries and let residual lowering cut only the counter branches. fn preserve_metricsql_counter_only_roots( queries: &mut [QueryCompilationInput], + canonical_roots: &[Rc], composable: bool, ) -> Result<(), CompileError> { - for query in queries { + for (index, query) in queries.iter_mut().enumerate() { let selected = collect_selected_materializations(&query.selected_plan_root, composable) .map_err(|reason| CompileError::Query { query_id: query.query_id.clone(), @@ -917,14 +931,7 @@ fn preserve_metricsql_counter_only_roots( { continue; } - let parsed = crate::query_parser::parse_query_expr_canonical( - &query.query_string, - query.accuracy_target.clone(), - ) - .map_err(|error| CompileError::Query { - query_id: query.query_id.clone(), - reason: error.to_string(), - })?; + let parsed = original_root(query, index, canonical_roots)?; query.selected_plan_root = crate::planner_selection::keep_pre_asap(&parsed).map_err(|error| { CompileError::Query { @@ -959,6 +966,10 @@ impl PhysicalPlanCompiler { environment: PhysicalDeploymentContext, frontend: QueryFrontend, ) -> Result { + if let Some(data) = &request.data_workload { + data.validate() + .map_err(|error| CompileError::Snapshot(error.to_string()))?; + } if super::maintained_population::supported(&request) && (frontend != QueryFrontend::PromQl || environment.target != PhysicalDeploymentTarget::BackendLocalRemoteWrite) @@ -1014,18 +1025,20 @@ impl PhysicalPlanCompiler { if frontend == QueryFrontend::MetricsQl { preserve_metricsql_counter_only_roots( &mut request.queries, + &request.canonical_roots, request.allow_mixed_summary_and_exact_execution, )?; } preserve_invalid_exact_fallback_roots( &mut request.queries, + &request.canonical_roots, request.allow_mixed_summary_and_exact_execution, )?; if environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite && !request.allow_mixed_summary_and_exact_execution { - preserve_native_unsafe_raw_roots(&mut request.queries)?; + preserve_native_unsafe_raw_roots(&mut request.queries, &request.canonical_roots)?; } let roots = request @@ -1315,10 +1328,13 @@ impl PhysicalPlanCompiler { &model, &environment, &state_consumers, - request - .query_workload - .as_ref() - .map(|workload| (workload, consumer_indices.clone())), + request.query_workload.as_ref().map(|workload| { + ( + workload, + request.data_workload.as_ref(), + consumer_indices.clone(), + ) + }), )?; let window_implementation = query.window_realization_candidates.iter() .find(|candidate| candidate.realization_id == planner_selection.window_realization_id @@ -1717,14 +1733,7 @@ impl PhysicalPlanCompiler { // Retain its native boundary without discarding other workload roots. let native_root = request.allow_mixed_summary_and_exact_execution && if let SummaryExpr::KeepPreAsap(expr) = &query.selected_plan_root.expr { - let original = crate::query_parser::parse_query_expr_canonical( - &query.query_string, - query.accuracy_target.clone(), - ) - .map_err(|error| CompileError::Query { - query_id: query.query_id.clone(), - reason: error.to_string(), - })?; + let original = original_root(query, query_index, &request.canonical_roots)?; expr.as_ref() == &original && crate::query_plan::residual::compile_logical( query.query_id.clone(), @@ -1898,9 +1907,8 @@ impl PhysicalPlanCompiler { validate_retained_summary_footprint( &materializations, request - .query_workload + .data_workload .as_ref() - .and_then(|workload| workload.data_workload.as_ref()) .and_then(|data| { data.input_cardinality .value_at(environment.observed_at_unix_ms) @@ -2892,7 +2900,7 @@ fn select_lifecycle( model: &ControlPlaneCostModel, environment: &PhysicalDeploymentContext, consumers: &[&QueryCompilationInput], - original_workload: Option<(&QueryWorkload, Vec)>, + original_workload: Option<(&QueryWorkload, Option<&DataWorkload>, Vec)>, ) -> Result { // Current lifecycle evidence is per producer with one unit read cost. // Conflicting source/rate/horizon/cost snapshots cannot be averaged into @@ -2941,58 +2949,35 @@ fn select_lifecycle( }) .collect(), ), - data_workload: Some(DataWorkload { - arrival: DataArrival::ContinuouslyIngesting, - // This workload is an internal lifecycle projection, not a - // plan-ready query input. Use the compatibility profile cadence - // when there is no original workload to carry source evidence. - data_ingestion_interval: Evidence { - value: Some(DurationMs(1_000)), - source: EvidenceSource::Declared, - observed_at_ms: None, - valid_for_ms: None, - }, - ingestion_rate: Evidence { - value: Some(Rate( - query.summary_lifecycle_inputs.ingestion_rate_per_second, - )), - source: EvidenceSource::Observed, - observed_at_ms: Some(query.summary_lifecycle_inputs.evidence_observed_at_unix_ms), - valid_for_ms: Some(query.summary_lifecycle_inputs.evidence_valid_for_ms), - }, - ..DataWorkload::default() - }), }; - let (workload, indices) = match original_workload { - Some((original, indices)) => { - let mut original = original.clone(); - // HTTP demand carries phases and requirements; its per-state source - // evidence comes from the same lifecycle input as window pricing. - if original.data_workload.is_none() { - original.data_workload = workload.data_workload; - } else if original - .data_workload - .as_ref() - .is_some_and(|data| data.data_ingestion_interval.value.is_none()) - { - original - .data_workload - .as_mut() - .unwrap() - .data_ingestion_interval = workload - .data_workload - .as_ref() - .unwrap() - .data_ingestion_interval - .clone(); - } - (original, indices) - } - None => (workload, (0..consumers.len()).collect()), + let fallback_data = DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + // This workload is an internal lifecycle projection, not a + // plan-ready query input. Use the compatibility profile cadence + // when there is no original workload to carry source evidence. + data_ingestion_interval: Evidence { + value: Some(DurationMs(1_000)), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }, + ingestion_rate: Evidence { + value: Some(Rate( + query.summary_lifecycle_inputs.ingestion_rate_per_second, + )), + source: EvidenceSource::Observed, + observed_at_ms: Some(query.summary_lifecycle_inputs.evidence_observed_at_unix_ms), + valid_for_ms: Some(query.summary_lifecycle_inputs.evidence_valid_for_ms), + }, + ..DataWorkload::default() + }; + let (workload, data, indices) = match original_workload { + Some((original, data, indices)) => (original, data.unwrap_or(&fallback_data), indices), + None => (&workload, &fallback_data, (0..consumers.len()).collect()), }; let plan = plan_summary_maintenance_lifecycles( Rc::new(node.clone()), - WorkloadDemand::new(&workload, &indices), + WorkloadDemand::new_with_data(workload, data, &indices), environment.observed_at_unix_ms, Some(Horizon(query.summary_lifecycle_inputs.horizon_seconds)), SummaryMaintenanceLifecycleCapabilities { @@ -3857,6 +3842,7 @@ pub(crate) mod tests { )) .unwrap(); snapshot.schema_version = 2; + snapshot.data_workload.data_ingestion_interval.value = Some(DurationMs(1_000)); let template = snapshot.query_workload.repeating_queries.as_ref().unwrap()[0].clone(); let queries = [ "quantile by (job) (0.5, a)", @@ -3902,10 +3888,9 @@ pub(crate) mod tests { .. } = node { - // Scrape cadence bounds input lag; selector staleness keeps - // the independently declared Planner lookback. - assert_eq!(population.max_input_lag_ms, 5_000); - assert_eq!(population.lookback_ms, 300_000); + // A longer scrape cadence cannot extend the declared selector horizon. + assert_eq!(population.max_input_lag_ms, 1_000); + assert_eq!(population.lookback_ms, 1_000); assert_eq!(population.max_k, 5); assert!(population.quantiles); populations.insert(population.key()); @@ -3964,6 +3949,10 @@ pub(crate) mod tests { let planner_types::pre_asap::QueryExpr::Aggregate { child, .. } = &mut root else { unreachable!() }; + // SQL rows have no PromQL instant-selector time range. + if let QueryExpr::TimeRange { child: scan, .. } = child.as_ref() { + *child = Rc::clone(scan); + } let planner_types::pre_asap::QueryExpr::Scan { source, schema, .. } = Rc::make_mut(child) else { unreachable!() @@ -4163,8 +4152,7 @@ pub(crate) mod tests { for text in [ "avg_over_time(data[5m])", "min_over_time(data[5m])", - "quantile_over_time(0.9,data[5m])/quantile_over_time(0.5,data[5m])", - "avg_over_time(data[5m])/quantile_over_time(0.5,data[5m])", + "quantile_over_time(0.9,data[5m])", ] { let mut snapshot = planning_snapshot(); let entry = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0]; @@ -4197,6 +4185,34 @@ pub(crate) mod tests { } } + // Quantile rank error does not certify relative error of a quotient. + #[test] + fn uncertified_quantile_ratios_retain_exact_execution() { + for text in [ + "quantile_over_time(0.9,data[5m])/quantile_over_time(0.5,data[5m])", + "avg_over_time(data[5m])/quantile_over_time(0.5,data[5m])", + ] { + let mut snapshot = planning_snapshot(); + snapshot.query_workload.repeating_queries.as_mut().unwrap()[0].query = + Query(text.into()); + let (request, environment) = snapshot.into_physical_compilation_request().unwrap(); + let plan = PhysicalPlanCompiler + .compile_promql(request, environment) + .unwrap(); + assert!(plan.precompute_plan.materializations.is_empty(), "{text}"); + assert!( + plan.query_plan + .entries + .values() + .all(|entry| entry.nodes.values().any(|node| matches!( + node, + crate::query_plan::QueryPlanNode::ExactFallback { .. } + ))), + "{text}" + ); + } + } + /// Optional counter masks must retain the workload's mandatory sketch bindings. #[test] fn costed_mixed_workload_retains_sketches_and_counter_readouts() { @@ -4283,7 +4299,8 @@ pub(crate) mod tests { for query in [ "sum_over_time(m[1m])", "quantile_over_time(0.99, m[1m])", - "sum_over_time(m[1m]) / count_over_time(m[1m])", + // Planner's certified sum/count decomposition, including its division guard. + "avg_over_time(m[1m])", ] { let mut environment = environment(10_000); environment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; @@ -4735,9 +4752,11 @@ pub(crate) mod tests { } Ok(PhysicalCompilationRequest { planner_selection_trace: Vec::new(), + canonical_roots: Vec::new(), allow_mixed_summary_and_exact_execution: false, enabled_materialization_keys: None, query_workload: None, + data_workload: None, queries: vec![QueryCompilationInput { query_id: query_id.into(), query_string: promql.into(), @@ -5455,7 +5474,6 @@ pub(crate) mod tests { }) .collect(), ), - data_workload: None, } } @@ -5774,7 +5792,8 @@ pub(crate) mod tests { SummaryExpr::KeepPreAsap(_) )); - preserve_invalid_exact_fallback_roots(&mut request.queries, true).unwrap(); + preserve_invalid_exact_fallback_roots(&mut request.queries, &request.canonical_roots, true) + .unwrap(); assert!(matches!( request.queries[0].selected_plan_root.expr, @@ -5953,7 +5972,6 @@ pub(crate) mod tests { value["query_workload"]["repeating_queries"][0]["query"] = json!("quantile(0.9, sum_over_time(m[1m]))"); value["data_workload"]["ingestion_rate"]["value"] = json!(rate); - value["query_workload"]["data_workload"]["ingestion_rate"]["value"] = json!(rate); let snapshot: BackendLocalPlanningInput = serde_json::from_value(value).unwrap(); let (request, env) = snapshot.into_physical_compilation_request().unwrap(); let plan = PhysicalPlanCompiler.compile_promql(request, env).unwrap(); @@ -6223,13 +6241,7 @@ pub(crate) mod tests { observed_at_ms: None, valid_for_ms: None, }; - snapshot.data_workload.data_ingestion_interval = evidence.clone(); - snapshot - .query_workload - .data_workload - .as_mut() - .unwrap() - .data_ingestion_interval = evidence; + snapshot.data_workload.data_ingestion_interval = evidence; } // Plan-ready lowering uses workload cadence, not the backend-local migration field. @@ -6244,6 +6256,30 @@ pub(crate) mod tests { assert_eq!(request.queries[0].query_lookback_seconds, 1); } + // Snapshot lowering must use the environment clock for cadence freshness. + #[test] + fn snapshot_rejects_unavailable_cadence_evidence() { + for (observed, valid_for) in [ + (Some(0), Some(0)), + (Some(u64::MAX), None), + (None, Some(100)), + ] { + let mut snapshot = planning_snapshot(); + snapshot + .data_workload + .data_ingestion_interval + .observed_at_ms = observed; + snapshot.data_workload.data_ingestion_interval.valid_for_ms = valid_for; + let error = snapshot.into_physical_compilation_request().unwrap_err(); + assert!( + error + .to_string() + .contains("evidence is unavailable at planning time"), + "{error}" + ); + } + } + // Legacy snapshots migrate their explicit scrape cadence into DataWorkload. #[test] fn missing_data_ingestion_interval_migrates_scrape_interval() { @@ -6254,13 +6290,7 @@ pub(crate) mod tests { let (request, _) = snapshot.into_physical_compilation_request().unwrap(); assert_eq!( - request - .query_workload - .unwrap() - .data_workload - .unwrap() - .data_ingestion_interval - .value, + request.data_workload.unwrap().data_ingestion_interval.value, Some(DurationMs(5_000)) ); assert_eq!(request.queries[0].query_lookback_seconds, 5); @@ -7306,7 +7336,6 @@ pub(crate) mod tests { as_of: None, }, }]), - data_workload: Some(data_workload.clone()), }; let mut environment = environment(10_000); environment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; diff --git a/control_plane/src/physical/maintained_population.rs b/control_plane/src/physical/maintained_population.rs index 2b1270b1..e1e6d61c 100644 --- a/control_plane/src/physical/maintained_population.rs +++ b/control_plane/src/physical/maintained_population.rs @@ -92,7 +92,8 @@ pub(super) fn operator( .scrape_interval_ms .unwrap_or(60_000) .saturating_add(request.query_retention_margin_ms) - .clamp(1, 300_000), + .max(1) + .min(input.lookback_ms), }; population.validate()?; let readout = match readout { diff --git a/control_plane/src/physical/post_asap/lower.rs b/control_plane/src/physical/post_asap/lower.rs index b9c69f72..bd981854 100644 --- a/control_plane/src/physical/post_asap/lower.rs +++ b/control_plane/src/physical/post_asap/lower.rs @@ -122,6 +122,12 @@ fn bind_recursive( bind_recursive(&pushed, cost_model) } + // Workload-aware instant selectors retain their source horizon even + // when there is no aggregate to bind. + QueryExpr::TimeRange { .. } => Ok(PostAsapPlan::Summary( + crate::planner_selection::keep_pre_asap(expr)?, + )), + // Exact Count cannot use this deployment's Sum accumulator: it counts // values rather than samples. Keep it logical for archive execution. QueryExpr::Aggregate { diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index 8b5eff65..5816ce8e 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -807,12 +807,16 @@ fn materialization_alternatives( let mut exact = request.clone(); exact.allow_mixed_summary_and_exact_execution = false; exact.enabled_materialization_keys = None; - for query in &mut exact.queries { - let parsed = crate::query_parser::parse_query_expr_canonical( - &query.query_string, - query.accuracy_target.clone(), - ) - .map_err(|error| invalid(error.to_string()))?; + for (index, query) in exact.queries.iter_mut().enumerate() { + let parsed = if let Some(root) = request.canonical_roots.get(index) { + root.as_ref().clone() + } else { + crate::query_parser::parse_query_expr_canonical( + &query.query_string, + query.accuracy_target.clone(), + ) + .map_err(|error| invalid(error.to_string()))? + }; query.selected_plan_root = crate::planner_selection::keep_pre_asap(&parsed) .map_err(|error| invalid(error.to_string()))?; } diff --git a/control_plane/src/query_parser/mod.rs b/control_plane/src/query_parser/mod.rs index eb2883e1..8913635b 100644 --- a/control_plane/src/query_parser/mod.rs +++ b/control_plane/src/query_parser/mod.rs @@ -80,29 +80,39 @@ pub fn parse_query_expr_canonical( query: &str, accuracy: AccuracyTarget, ) -> anyhow::Result { - let workload = QueryWorkload { - language: QueryLanguage::PromQL, - query_batch: Some(vec![BatchEntry { - query: Query(query.trim().into()), - requirements: QueryRequirements { - accuracy: AccuracyRequirement::Explicit(accuracy), - ..Default::default() - }, - predictability: Default::default(), - invocations: 1, - execute_at: None, - time_selection: TimeSelection::default(), - }]), - repeating_queries: None, + parse_query_expr_with_interval(query, accuracy, COMPATIBILITY_DATA_INGESTION_INTERVAL_MS) +} + +pub(crate) fn parse_query_expr_with_interval( + query: &str, + accuracy: AccuracyTarget, + interval_ms: u64, +) -> anyhow::Result { + let workload = planner_types::workload::PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: Some(vec![BatchEntry { + query: Query(query.trim().into()), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit(accuracy), + ..Default::default() + }, + predictability: Default::default(), + invocations: 1, + execute_at: None, + time_selection: TimeSelection::default(), + }]), + repeating_queries: None, + }, data_workload: Some(DataWorkload { data_ingestion_interval: Evidence { - value: Some(DurationMs(COMPATIBILITY_DATA_INGESTION_INTERVAL_MS)), + value: Some(DurationMs(interval_ms)), ..Default::default() }, ..Default::default() }), }; - asap_frontend_promql::lower_promql_workload(&workload)? + asap_frontend_promql::lower_promql_workload(&workload, 0)? .into_iter() .next() .ok_or_else(|| anyhow::anyhow!("PromQL compatibility workload produced no query")) diff --git a/control_plane/src/query_plan/residual.rs b/control_plane/src/query_plan/residual.rs index cd3f5c98..83fcdda0 100644 --- a/control_plane/src/query_plan/residual.rs +++ b/control_plane/src/query_plan/residual.rs @@ -291,6 +291,41 @@ pub fn compile_logical( Ok(entry) } +fn horizons(expr: &planner_types::pre_asap::QueryExpr, out: &mut Vec) { + use planner_types::pre_asap::QueryExpr; + if let QueryExpr::TimeRange { range, .. } = expr { + if let Ok(ms) = u64::try_from(range.as_millis()) { + out.push(ms); + } + } + match expr { + QueryExpr::PromqlScalarBridge(child) + | QueryExpr::PromqlVectorFromScalar(child) + | QueryExpr::PromqlScalarFromVector(child) + | QueryExpr::PromqlRelabel { child, .. } + | QueryExpr::PromqlSeriesSample { child, .. } + | QueryExpr::PromqlInfoEnrich { child, .. } + | QueryExpr::Filter { child, .. } + | QueryExpr::Project { child, .. } + | QueryExpr::Aggregate { child, .. } + | QueryExpr::Dedup { child, .. } + | QueryExpr::Sort { child, .. } + | QueryExpr::Limit { child, .. } + | QueryExpr::PromqlSubquery { child, .. } + | QueryExpr::TimeRange { child, .. } + | QueryExpr::TimeShift { child, .. } => horizons(child, out), + QueryExpr::BinaryOp { lhs, rhs, .. } => { + horizons(lhs, out); + horizons(rhs, out); + } + QueryExpr::Join { left, right, .. } | QueryExpr::SetOp { left, right, .. } => { + horizons(left, out); + horizons(right, out); + } + _ => {} + } +} + /// Match residuals by semantic IR equality, not display text or source names. /// This ensures a subtree parsed for physical lowering is the subtree Planner kept. pub(super) fn residual_nodes( @@ -319,18 +354,28 @@ pub(super) fn residual_nodes( let original = parser::parse(original).map_err(|e| invalid(e.to_string()))?; let mut expressions = Vec::new(); visit(&original, &mut expressions); + // Reconstruct equality witnesses with the selected IR's source horizon, + // not the compatibility parser's default. Explicit matrix ranges remain + // query-owned and equality still checks the complete tree. + let mut intervals = vec![1_000]; + horizons(residual, &mut intervals); + intervals.sort_unstable(); + intervals.dedup(); for expression in expressions { - if let Ok(candidate) = crate::query_parser::parse_query_expr_canonical( - &expression.to_string(), - planner_types::types::AccuracyTarget::Exact, - ) { - if &candidate == residual { - let mut lower = Lower { - nodes: BTreeMap::new(), - seen: BTreeMap::new(), - }; - let root = lower.lower(expression)?; - return Ok((root, lower.nodes)); + for interval in &intervals { + if let Ok(candidate) = crate::query_parser::parse_query_expr_with_interval( + &expression.to_string(), + planner_types::types::AccuracyTarget::Exact, + *interval, + ) { + if &candidate == residual { + let mut lower = Lower { + nodes: BTreeMap::new(), + seen: BTreeMap::new(), + }; + let root = lower.lower(expression)?; + return Ok((root, lower.nodes)); + } } } } @@ -424,33 +469,78 @@ pub(crate) fn selected_residual_nodes( let parsed = parser::parse(original).map_err(|e| invalid(e.to_string()))?; let mut expressions = Vec::new(); visit(&parsed, &mut expressions); + fn selected_horizons(node: &planner_types::post_asap::SummaryNode, out: &mut Vec) { + use planner_types::post_asap::SummaryExpr; + match &node.expr { + SummaryExpr::KeepPreAsap(expr) => horizons(expr, out), + SummaryExpr::ValueOperation { child, .. } | SummaryExpr::SummaryAgg { child, .. } => { + selected_horizons(child, out) + } + SummaryExpr::SummaryEstimate { summary_input, .. } + | SummaryExpr::SummaryDelete { summary_input, .. } => { + selected_horizons(summary_input, out) + } + SummaryExpr::BinaryOp { + lhs: left, + rhs: right, + .. + } + | SummaryExpr::CandidateTopK { + candidates: left, + values: right, + .. + } + | SummaryExpr::RelationalJoin { left, right, .. } + | SummaryExpr::SummaryJoin { + outer: left, + inner: right, + .. + } + | SummaryExpr::SummarySubtract { left, right } => { + selected_horizons(left, out); + selected_horizons(right, out); + } + SummaryExpr::SummaryMerge { children } => { + for child in children { + selected_horizons(child, out); + } + } + } + } + let mut intervals = vec![1_000]; + selected_horizons(selected, &mut intervals); + intervals.sort_unstable(); + intervals.dedup(); let mut matched = None; for expression in expressions { - let Ok(canonical) = crate::query_parser::parse_query_expr_canonical( - &expression.to_string(), - planner_types::types::AccuracyTarget::Exact, - ) else { - continue; - }; - let Ok(witness) = crate::planner_selection::select_summary_default(&canonical) else { - continue; - }; - if witness.as_ref() == selected { - let mut lower = Lower { - nodes: BTreeMap::new(), - seen: BTreeMap::new(), + for interval in &intervals { + let Ok(canonical) = crate::query_parser::parse_query_expr_with_interval( + &expression.to_string(), + planner_types::types::AccuracyTarget::Exact, + *interval, + ) else { + continue; }; - let root = lower.lower(expression)?; - let candidate = (root, lower.nodes); - if matched - .as_ref() - .is_some_and(|previous| previous != &candidate) - { - return Err(invalid( - "ambiguous original subtrees share a Planner summary representation", - )); + let Ok(witness) = crate::planner_selection::select_summary_default(&canonical) else { + continue; + }; + if witness.as_ref() == selected { + let mut lower = Lower { + nodes: BTreeMap::new(), + seen: BTreeMap::new(), + }; + let root = lower.lower(expression)?; + let candidate = (root, lower.nodes); + if matched + .as_ref() + .is_some_and(|previous| previous != &candidate) + { + return Err(invalid( + "ambiguous original subtrees share a Planner summary representation", + )); + } + matched = Some(candidate); } - matched = Some(candidate); } } matched.ok_or_else(|| { @@ -1249,6 +1339,26 @@ pub use eligible_materialization_keys as materialization_candidate_keys; #[cfg(test)] mod tests { use super::*; + // A workload horizon changes the equality witness, never its filter or explicit range. + #[test] + fn workload_horizon_residual_keeps_semantic_equality() { + let residual = crate::query_parser::parse_query_expr_with_interval( + "sum(m{job=\"api\"})", + planner_types::types::AccuracyTarget::Exact, + 5_000, + ) + .unwrap(); + assert!(residual_nodes("sum(m{job=\"api\"})", &residual).is_ok()); + assert!(residual_nodes("sum(m{job=\"worker\"})", &residual).is_err()); + let range = crate::query_parser::parse_query_expr_with_interval( + "sum_over_time(m[1m])", + planner_types::types::AccuracyTarget::Exact, + 5_000, + ) + .unwrap(); + assert!(residual_nodes("sum_over_time(m[2m])", &range).is_err()); + } + fn instant() -> InstantExecution { InstantExecution { lookback_ms: 300_000, diff --git a/control_plane/tests/discovery_snapshot.rs b/control_plane/tests/discovery_snapshot.rs index 03b88164..c6e4a69b 100644 --- a/control_plane/tests/discovery_snapshot.rs +++ b/control_plane/tests/discovery_snapshot.rs @@ -46,6 +46,10 @@ fn discovered_snapshot_plans_with_observed_cadence_and_promql_history() { let snapshot: BackendLocalPlanningInput = serde_json::from_str(&fs::read_to_string(&output).unwrap()).unwrap(); assert_eq!(snapshot.physical_inputs.scrape_interval_ms, 60_000); + assert_eq!( + snapshot.data_workload.data_ingestion_interval.value, + Some(planner_types::workload::DurationMs(60_000)) + ); let (request, _) = snapshot .clone() .into_physical_compilation_request() diff --git a/crates/asap_types/src/query_plan/current_series.rs b/crates/asap_types/src/query_plan/current_series.rs index 8dc53abf..a4b7480f 100644 --- a/crates/asap_types/src/query_plan/current_series.rs +++ b/crates/asap_types/src/query_plan/current_series.rs @@ -24,7 +24,8 @@ impl SeriesPopulation { } pub fn validate(&self) -> Result<(), QueryPlanError> { if self.metric.is_empty() - || self.lookback_ms != 300_000 + || self.lookback_ms == 0 + || self.lookback_ms > i64::MAX as u64 || self.max_input_lag_ms == 0 || self.max_input_lag_ms > self.lookback_ms || self.max_series == 0 @@ -50,3 +51,33 @@ pub enum SeriesReadout { Count, Average, } + +#[cfg(test)] +mod tests { + use super::*; + // Planner-declared horizons are valid independently of the old five-minute default. + #[test] + fn accepts_declared_horizon_with_bounded_lag() { + let mut population = SeriesPopulation { + metric: "a".into(), + matchers: vec![], + grouping: Grouping { + labels: vec![], + without: false, + }, + lookback_ms: 1_000, + max_input_lag_ms: 1_000, + max_series: 100, + max_bytes: 1_000_000, + max_k: 3, + quantiles: true, + }; + population.validate().unwrap(); + population.lookback_ms = 0; + assert!(population.validate().is_err()); + population.lookback_ms = u64::MAX; + assert!(population.validate().is_err()); + population.lookback_ms = 999; + assert!(population.validate().is_err()); + } +} diff --git a/crates/asap_types/src/table_population.rs b/crates/asap_types/src/table_population.rs index 946813cf..4bb37138 100644 --- a/crates/asap_types/src/table_population.rs +++ b/crates/asap_types/src/table_population.rs @@ -32,8 +32,10 @@ impl TablePopulation { ) { return Err("table population comparison is unsupported".into()); } - if matches!(predicate.value, ScalarValue::Null) - || matches!(predicate.value, ScalarValue::Float64(value) if !value.is_finite()) + if matches!( + predicate.value, + ScalarValue::Null | ScalarValue::Interval { .. } + ) || matches!(predicate.value, ScalarValue::Float64(value) if !value.is_finite()) { return Err("table population requires a finite non-null literal".into()); } 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 e3395668..c47b7b48 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 @@ -61,6 +61,9 @@ fn json_cell( )) }; match dtype { + DataType::Interval | DataType::Date => Err(ClickHouseRelationalError::Unsupported( + "temporal value transport".into(), + )), DataType::Null if value.is_null() => Ok(Cell::Null), DataType::Null => Err(invalid()), DataType::List { element } => { @@ -330,6 +333,7 @@ fn clickhouse_type_matches(actual: Option<&str>, expected: &DataType, nullable: return false; } match expected { + DataType::Interval | DataType::Date => false, DataType::Null => actual == "Nothing", DataType::List { element } => { !nullable @@ -573,6 +577,11 @@ fn eval( )) } QueryExpr::Literal(value) => Ok(match value { + ScalarValue::Interval { .. } => { + return Err(ClickHouseRelationalError::Unsupported( + "interval literal".into(), + )) + } ScalarValue::Int64(value) => Cell::Int64(*value), ScalarValue::Float64(value) => Cell::Float64(*value), ScalarValue::Utf8(value) => Cell::Utf8(value.clone()), @@ -752,6 +761,11 @@ fn default_collection_element( return Ok(Cell::Null); } Ok(match dtype { + DataType::Interval | DataType::Date => { + return Err(ClickHouseRelationalError::Unsupported( + "temporal value transport".into(), + )) + } DataType::Null => Cell::Null, DataType::Int64 => Cell::Int64(0), DataType::Float64 => Cell::Float64(0.0), @@ -955,6 +969,8 @@ fn cell_cmp(left: &Cell, right: &Cell) -> Option { fn arrow_type(dtype: &DataType) -> ArrowDataType { match dtype { + DataType::Date => ArrowDataType::Date32, + DataType::Interval => ArrowDataType::Interval(arrow::datatypes::IntervalUnit::MonthDayNano), DataType::Null => ArrowDataType::Null, DataType::List { element } => ArrowDataType::List(Arc::new(Field::new( &element.name, @@ -1015,6 +1031,11 @@ fn build_array( }}; } Ok(match dtype { + DataType::Interval | DataType::Date => { + return Err(ClickHouseRelationalError::Unsupported( + "temporal value transport".into(), + )) + } DataType::Null => { if rows .iter() diff --git a/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs b/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs index d19aad0e..acea6498 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs @@ -129,7 +129,9 @@ impl ClickHouseReader { ScalarValue::Int64(_) => "Int64", ScalarValue::Float64(_) => "Float64", ScalarValue::Boolean(_) => "Bool", - ScalarValue::Null => unreachable!("validated table literal"), + ScalarValue::Null | ScalarValue::Interval { .. } => { + unreachable!("validated table literal") + } }; format!( "{} {operator} {{population_{index}:{kind}}}", @@ -302,7 +304,9 @@ impl RawSampleReader for ClickHouseReader { ScalarValue::Int64(value) => value.to_string(), ScalarValue::Float64(value) => value.to_string(), ScalarValue::Boolean(value) => value.to_string(), - ScalarValue::Null => unreachable!("validated table literal"), + ScalarValue::Null | ScalarValue::Interval { .. } => { + unreachable!("validated table literal") + } }; request = request.query(&[(format!("param_population_{index}"), value)]); } diff --git a/data_plane/src/storage_engines/sketch_db/current_series.rs b/data_plane/src/storage_engines/sketch_db/current_series.rs index f4a1c5d3..2ba9e584 100644 --- a/data_plane/src/storage_engines/sketch_db/current_series.rs +++ b/data_plane/src/storage_engines/sketch_db/current_series.rs @@ -401,6 +401,28 @@ mod tests { quantiles: true, } } + + // A declared one-second horizon expires the left-boundary sample, not five minutes later. + #[test] + fn declared_one_second_horizon_expires_members() { + let mut population = definition(); + population.lookback_ms = 1_000; + population.max_input_lag_ms = 1_000; + population.validate().unwrap(); + let plan = plan(&population); + let mut store = CurrentSeriesStore::default(); + store.ingest(&plan, &[sample("old", "api", 0, Some(10.))]); + store.ingest(&plan, &[sample("new", "api", 500, Some(3.))]); + let values = store + .read((7, 1), &population, &SeriesReadout::Sum, 1_000) + .unwrap(); + assert_eq!(values.len(), 1); + assert_eq!(values[0].1, 3.); + let values = store + .read((7, 1), &population, &SeriesReadout::Sum, 1_500) + .unwrap(); + assert!(values.is_empty()); + } fn plan(p: &SeriesPopulation) -> QueryPlan { let mut plan = QueryPlan::empty(); plan.plan_id = 7; diff --git a/data_plane/tests/support/current_series_process.rs b/data_plane/tests/support/current_series_process.rs index b4e36951..be23ae92 100644 --- a/data_plane/tests/support/current_series_process.rs +++ b/data_plane/tests/support/current_series_process.rs @@ -26,7 +26,14 @@ async fn current_series_quantiles_topk_share_and_replace_values() { )) .unwrap(); snapshot.schema_version = 2; - snapshot.physical_inputs.scrape_interval_ms = 60_000; + // Exercise a short declared horizon; native differential mode keeps its five-minute contract. + let horizon_ms: i64 = if native.is_some() { 300_000 } else { 5_000 }; + let scrape_ms = horizon_ms / 5; + snapshot.physical_inputs.scrape_interval_ms = scrape_ms as u64; + snapshot.data_workload.data_ingestion_interval = planner_types::workload::Evidence { + value: Some(planner_types::workload::DurationMs(horizon_ms as u64)), + ..Default::default() + }; let template = snapshot.query_workload.repeating_queries.as_ref().unwrap()[0].clone(); let mut queries = vec![]; for q in [0.5, 0.9, 0.95, 0.99] { @@ -93,6 +100,21 @@ async fn current_series_quantiles_topk_share_and_replace_values() { quotes, }); let planned = snapshot.clone().compile_promql().unwrap(); + for entry in planned.query_plan.entries.values() { + for node in entry.nodes.values() { + if let asap_types::query_plan::QueryPlanNode::Logical { + operator: + asap_types::query_plan::residual::ResidualQueryOperator::CurrentSeries { + population, + .. + }, + .. + } = node + { + assert_eq!(population.lookback_ms, horizon_ms as u64); + } + } + } assert!( planned .query_plan @@ -141,7 +163,7 @@ async fn current_series_quantiles_topk_share_and_replace_values() { .duration_since(std::time::UNIX_EPOCH) .unwrap() .as_millis() as i64; - let times: Vec<_> = (0..=5).map(|i| end - 300_000 + i * 60_000).collect(); + let times: Vec<_> = (0..=5).map(|i| end - horizon_ms + i * scrape_ms).collect(); let wire = WriteRequest { timeseries: [ ("x", "api", 1.), @@ -346,8 +368,8 @@ async fn current_series_quantiles_topk_share_and_replace_values() { } // Keep one series fresh while all members last seen at `end` hit the exact // left lookback boundary. Native Prometheus 3.5 must agree for every readout. - let updates: Vec<_> = (60_000..=300_000) - .step_by(60_000) + let updates: Vec<_> = (scrape_ms.max(3_000)..=horizon_ms) + .step_by(scrape_ms as usize) .map(|offset| (end + offset, 9.0)) .collect(); let wire = WriteRequest { @@ -362,14 +384,14 @@ async fn current_series_quantiles_topk_share_and_replace_values() { } assert_eq!(remote_write(&client, &base, &wire).await, 204); for text in &queries { - let body = query(&client, &base, text, end + 300_000).await; + let body = query(&client, &base, text, end + horizon_ms).await; assert!(is_warm(&body), "{text}: {body}"); assert_eq!( body["data"]["result"].as_array().unwrap().len(), 1, "{text}: {body}" ); - compare_native(&client, &native, text, end + 300_000, &body).await; + compare_native(&client, &native, text, end + horizon_ms, &body).await; } task.abort(); } diff --git a/docs/examples/asapquery-compatibility-demo-snapshot.json b/docs/examples/asapquery-compatibility-demo-snapshot.json index f617f8ce..e381d536 100644 --- a/docs/examples/asapquery-compatibility-demo-snapshot.json +++ b/docs/examples/asapquery-compatibility-demo-snapshot.json @@ -55,15 +55,7 @@ "predictability": { "predictable": { "known_at": null } }, "time_selection": { "scope": "real_time", "lookback": null, "as_of": null } } - ], - "data_workload": { - "arrival": "continuously_ingesting", - "data_ingestion_interval": { "value": 5000, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, - "ingestion_volume": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, - "ingestion_rate": { "value": 100.0, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, - "input_cardinality": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, - "distribution": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null } - } + ] }, "data_workload": { "arrival": "continuously_ingesting", diff --git a/docs/examples/asapquery-planning-snapshot.json b/docs/examples/asapquery-planning-snapshot.json index 39017909..471c8ad0 100644 --- a/docs/examples/asapquery-planning-snapshot.json +++ b/docs/examples/asapquery-planning-snapshot.json @@ -22,15 +22,7 @@ "as_of": null } } - ], - "data_workload": { - "arrival": "continuously_ingesting", - "data_ingestion_interval": { "value": 5000, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, - "ingestion_volume": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, - "ingestion_rate": { "value": 100.0, "source": "declared", "observed_at_ms": null, "valid_for_ms": null }, - "input_cardinality": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null }, - "distribution": { "value": null, "source": "unknown", "observed_at_ms": null, "valid_for_ms": null } - } + ] }, "data_workload": { "arrival": "continuously_ingesting", diff --git a/docs/user_guide/asapquery-profile.md b/docs/user_guide/asapquery-profile.md index fba4e022..8a9a2a0d 100644 --- a/docs/user_guide/asapquery-profile.md +++ b/docs/user_guide/asapquery-profile.md @@ -88,6 +88,15 @@ ASAPPlanner, compiles matching SummaryCatalog/PrecomputePlan/QueryPlan views, an installs the resulting immutable snapshot before accepting traffic. It does not create or wait for a CollectorPlan. +Supply data evidence only in the top-level `data_workload`; nesting it inside +`query_workload` is no longer accepted. `data_ingestion_interval` declares the +instant-selector horizon in milliseconds. Snapshot planning checks its freshness +at `environment.observed_at_unix_ms`; HTTP physical planning checks it at the +current planning time and requires the same top-level data evidence. +Current-series plans retain this horizon for membership expiry, and their input +lag allowance never exceeds it. Missing cadence in a legacy snapshot is derived +from its explicit scrape interval; expired or invalid supplied evidence is rejected. + `implementation.max_retained_summary_bytes` limits the estimated total encoded summary footprint across every retained pane and partition. It defaults to 2 GiB for older snapshots and is clamped to `--persistence-memory-limit-mb` at diff --git a/tools/o11y-execution/discover_snapshot.py b/tools/o11y-execution/discover_snapshot.py index 7308b925..9b2ec9a6 100644 --- a/tools/o11y-execution/discover_snapshot.py +++ b/tools/o11y-execution/discover_snapshot.py @@ -76,7 +76,7 @@ def evidence(value): "occurrence_count": frequencies[query], "expected_evaluations": frequencies[query] * args.repetitions, "declared_interval_ms": interval, "evaluation_phase_ms": phase, "lookback_method": "derived by the backend compiler from PromQL and the declared scrape cadence"}) - snapshot["query_workload"].update(repeating_queries=registrations, data_workload=data, query_batch=None) + snapshot["query_workload"].update(repeating_queries=registrations, query_batch=None) snapshot["snapshot_version"] = 2 snapshot.pop("workload_cost_evidence", None) implementation = snapshot["implementation"] @@ -87,6 +87,9 @@ def evidence(value): if not isinstance(scrape_interval_ms, int) or scrape_interval_ms <= 0 or scrape_interval_ms % 1000: raise ValueError("scrape_interval_ms must be a positive whole number of seconds") implementation["scrape_interval_ms"] = scrape_interval_ms + data["data_ingestion_interval"] = evidence(scrape_interval_ms) + if source_sample_interval_ms is None: + data["data_ingestion_interval"]["source"] = "declared" # Finite replay evaluates the oldest repetition first after loading the # complete input. Preserve that admitted historical-query span separately # from each query's PromQL range selector. From a8f75721714c9a74e9df0d20b72cc8805cf75fac Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 18 Sep 2026 18:02:17 +0000 Subject: [PATCH 3/4] Separate accelerated issue workloads from uncertified ratio fallbacks --- .../tests/support/issue_701_702_process.rs | 120 ++++++++++-------- 1 file changed, 68 insertions(+), 52 deletions(-) diff --git a/data_plane/tests/support/issue_701_702_process.rs b/data_plane/tests/support/issue_701_702_process.rs index 4baff6d4..d1fe869e 100644 --- a/data_plane/tests/support/issue_701_702_process.rs +++ b/data_plane/tests/support/issue_701_702_process.rs @@ -71,24 +71,17 @@ fn queries() -> Vec<(String, u64, u64)> { 300, 30, )); - queries.push(( - "quantile_over_time(0.9, issue701_data[5m]) / quantile_over_time(0.5, issue701_data[5m])" - .into(), - 300, - 60, - )); - queries.push(( - "avg_over_time(issue701_data[5m]) / quantile_over_time(0.5, issue701_data[5m])".into(), - 300, - 30, - )); queries } // A single mixed workload covers moving windows, current series, minimum/average, -// and ratio accuracy. Optional native URL adds a real Prometheus differential oracle. +// without uncertified ratios. Optional native URL adds a real Prometheus differential oracle. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn issue_workloads_execute_warm_at_successive_evaluations() { + run_warm_workload(queries()).await; +} + +async fn run_warm_workload(queries: Vec<(String, u64, u64)>) { let native = std::env::var("ASAP_CURRENT_SERIES_PROMETHEUS_URL").ok(); if let Some(url) = &native { let info: Value = reqwest::get(format!("{url}/api/v1/status/buildinfo")) @@ -116,12 +109,12 @@ async fn issue_workloads_execute_warm_at_successive_evaluations() { .await .unwrap(); }); - let queries = queries(); let mut fixture: Value = serde_json::from_str(include_str!( "../../../docs/examples/asapquery-planning-snapshot.json" )) .unwrap(); fixture["implementation"]["scrape_interval_ms"] = 1000.into(); + fixture["data_workload"]["data_ingestion_interval"]["value"] = 1000.into(); fixture["implementation"]["horizon_seconds"] = 3600.into(); let template = fixture["query_workload"]["repeating_queries"][0].clone(); fixture["query_workload"]["repeating_queries"] = queries @@ -147,7 +140,8 @@ async fn issue_workloads_execute_warm_at_successive_evaluations() { control_plane::query_plan::QueryPlanNode::ExactFallback { .. } | control_plane::query_plan::QueryPlanNode::ExternalExact { .. } | control_plane::query_plan::QueryPlanNode::Logical { - operator: control_plane::query_plan::residual::ResidualQueryOperator::ExactSubquery { .. }, .. + operator: control_plane::query_plan::residual::ResidualQueryOperator::ExactSubquery { .. } + | control_plane::query_plan::residual::ResidualQueryOperator::CandidateExactSubquery { .. }, .. } ))) }; @@ -253,6 +247,14 @@ async fn issue_workloads_execute_warm_at_successive_evaluations() { ) .await; assert!(is_warm(&actual), "{query}: {actual}"); + assert!( + actual["infos"] + .as_array() + .unwrap() + .iter() + .any(|info| info == "execution: asap"), + "{query}: {actual}" + ); if let Some(url) = &native { let expected: Value = client .get(format!("{url}/api/v1/query")) @@ -297,46 +299,60 @@ async fn issue_workloads_execute_warm_at_successive_evaluations() { } } } - if let Some(url) = &native { - let wire = WriteRequest { - timeseries: [("a", "api"), ("b", "api"), ("c", "db")] - .into_iter() - .map(|(pod, job)| { - let samples: Vec<_> = (962..=1261).map(|i| (origin + i * 1000, 0.0)).collect(); - series_with_labels("issue701_data", &[("pod", pod), ("job", job)], &samples) - }) - .collect(), - }; - assert_eq!(remote_write(&client, url, &wire).await, 204); - assert_eq!(remote_write(&client, &backend, &wire).await, 204); - let at = (origin + 1260 * 1000) as f64 / 1000.0; - for (query, _, _) in queries.iter().filter(|(q, _, _)| q.contains(" / ")) { - let params = [("query", query.clone()), ("time", at.to_string())]; - let expected: Value = client - .get(format!("{url}/api/v1/query")) - .query(¶ms) - .send() - .await - .unwrap() - .json() - .await - .unwrap(); - let actual: Value = client - .get(format!("{backend}/api/v1/query")) - .query(¶ms) - .send() - .await - .unwrap() - .json() - .await +} + +// Each supported #702 query must independently compile and execute without an exact service. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn issue_702_individual_queries_execute_without_fallback() { + for query in [ + "quantile by (job) (0.9, issue701_data)", + "sum by (job) (issue701_data)", + "sum(issue701_data)", + "avg_over_time(issue701_data[5m])", + "count by (job) (issue701_data)", + "count(issue701_data)", + "avg by (job) (issue701_data)", + "sum by (job) (sum_over_time(issue701_data[5m]))", + ] { + eprintln!("Checking acceleration: {query}"); + let temporal = query.contains('['); + run_warm_workload(vec![( + query.into(), + if temporal { 300 } else { 1 }, + if temporal { 30 } else { 1 }, + )]) + .await; + } +} + +// Rank-error guarantees do not certify division; keep this limitation explicit in CI. +#[test] +fn issue_701_702_uncertified_ratios_require_exact_fallback() { + for query in [ + "quantile_over_time(0.9, issue701_data[5m]) / quantile_over_time(0.5, issue701_data[5m])", + "avg_over_time(issue701_data[5m]) / quantile_over_time(0.5, issue701_data[5m])", + ] { + let mut fixture: Value = serde_json::from_str(include_str!( + "../../../docs/examples/asapquery-planning-snapshot.json" + )) + .unwrap(); + fixture["query_workload"]["repeating_queries"][0]["query"] = query.into(); + fixture["data_workload"]["data_ingestion_interval"]["value"] = 1000.into(); + let snapshot: BackendLocalPlanningInput = serde_json::from_value(fixture).unwrap(); + let (request, environment) = snapshot.into_physical_compilation_request().unwrap(); + let candidates = workload_cost::with_exact_alternative(request).unwrap(); + assert!(!candidates.is_empty()); + for candidate in candidates { + let plan = PhysicalCompiler + .compile_promql(candidate, environment.clone()) .unwrap(); + assert!(plan.precompute_plan.materializations.is_empty(), "{query}"); assert!( - !is_warm(&actual), - "undefined relative error must fall back: {query}" - ); - assert_eq!( - actual["data"], expected["data"], - "zero denominator: {query}" + plan.query_plan.entries.values().all(|entry| matches!( + entry.nodes.get(&entry.root), + Some(control_plane::query_plan::QueryPlanNode::ExactFallback { .. }) + )), + "{query}" ); } } From f176857bb54cd259fac3bbf00d6c509d0f509eb8 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 18 Sep 2026 18:17:37 +0000 Subject: [PATCH 4/4] Pin merged ASAPPlanner ingestion interval support --- Cargo.lock | 10 +++++----- Cargo.toml | 8 ++++---- 2 files changed, 9 insertions(+), 9 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 2f8454a8..aca9e44a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -364,7 +364,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f#c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?rev=da3635a80f8f854d47b772d49d5a9e5fb6927d8e)", @@ -376,7 +376,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f#c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f" dependencies = [ "asap-types", "promql-parser 0.10.0 (git+https://github.com/ProjectASAP/promql-parser?rev=9fede7eecca923c9882fe256484d00d37f8706cb)", @@ -385,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f#c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -408,12 +408,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f#c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a#1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f#c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f" dependencies = [ "serde", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index 15c27b38..dc0e70f0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -20,10 +20,10 @@ asap_sketchlib = { git = "https://github.com/ProjectASAP/asap_sketchlib", branch [workspace.dependencies] # Keep Planner frontends, selection, and IR on the same immutable revision. # Alias upstream asap-types because this workspace also defines asap_types. -planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" } -asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" } -asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" } -asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "1f1f0993ce0b273ed3bbf2ef6d214aa617e8978a" } +planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f" } +asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f" } +asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f" } +asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c9b2aa78aa8baa60ddc2f4d880bee4aaab944c5f" } # Shared external deps (used by 2+ crates) serde = { version = "1.0", features = ["derive"] }