From 06e657c41f48bc164d92d19a8eb63536bd564350 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 08:24:01 +0000 Subject: [PATCH 1/6] test: expect warm per-series arithmetic and panic-free without planning Restores the issue 701/702 warm avg_over_time assertions and the finite-division overflow test that #798 rewrote, and adds control-plane acceptance tests. These fail at the current Planner pin. Co-Authored-By: Claude Opus 5.5 --- control_plane/src/physical/compiler.rs | 48 +++++----- control_plane/src/physical/workload_cost.rs | 13 +++ .../tests/support/issue_701_702_process.rs | 87 ++++++++++++------- 3 files changed, 93 insertions(+), 55 deletions(-) diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 09ff1c1e5..98997cd75 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -5095,6 +5095,7 @@ pub(crate) mod tests { #[test] fn issue_701_702_temporal_workloads_have_warm_candidates() { for text in [ + "avg_over_time(data[5m])", "min_over_time(data[5m])", "quantile_over_time(0.9,data[5m])", ] { @@ -5150,30 +5151,31 @@ pub(crate) mod tests { assert!(plan.precompute_plan.materializations.is_empty()); } - // avg_over_time divides two per-series readouts. Planner does not yet - // match per-series rows in a Binary, so no candidate keeps local state and - // the exact engine evaluates the query. + // Per-series arithmetic over stored readouts runs as one Planner fragment + // over those readouts, not as an exact fallback. #[test] - fn per_series_average_has_no_warm_candidate_until_planner_matches_series() { - let mut snapshot = planning_snapshot(); - let entry = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0]; - entry.query = Query("avg_over_time(data[5m])".into()); - entry.requirements.accuracy = AccuracyRequirement::Explicit(AccuracyTarget::Exact); - let (request, environment) = snapshot.into_physical_compilation_request().unwrap(); - for candidate in - super::super::workload_cost::enumerate_exact_and_materialized_candidates(request) - .unwrap() - { - let Ok(plan) = DeploymentPlanCompiler.compile_promql(candidate, environment.clone()) - else { - continue; - }; - assert!(plan.precompute_plan.materializations.is_empty()); - assert!(plan - .query_plan - .entries - .values() - .all(|entry| entry.materialization_bindings().is_empty())); + fn per_series_arithmetic_executes_as_a_planner_fragment() { + for query in [ + "avg_over_time(a[1m])", + "rate(a[5m]) / rate(b[5m])", + "rate(a[5m]) * 2", + "sum_over_time(a[1m]) + sum_over_time(a{job=\"x\"}[1m])", + ] { + let mut environment = environment(10_000); + environment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; + environment.target_collector_ids.clear(); + let plan = DeploymentPlanCompiler + .compile_promql(request("arithmetic", query), environment) + .unwrap(); + let entry = plan.query_plan.lookup(query).unwrap(); + assert!( + matches!( + &entry.nodes[&entry.root], + crate::query_plan::QueryPlanNode::PhysicalFragment { .. } + ), + "{query}: {entry:?}" + ); + assert!(!entry.materialization_bindings().is_empty(), "{query}"); } } diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index 5a623f1d6..4ff84cdfd 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -1084,6 +1084,19 @@ mod tests { } } + /// `without` aggregations plan on the workload-cost fixture without a + /// Planner panic. + #[test] + fn without_aggregations_plan_on_the_workload_cost_fixture() { + for query in ["sum without (pod) (m)", "quantile without (pod) (0.5, m)"] { + let mut input = fixture(); + input.query_workload.repeating_queries.as_mut().unwrap()[0].query = + planner_types::workload::Query(query.into()); + let plan = with_unit_quotes(input).compile_promql().unwrap(); + plan.query_plan.lookup(query).unwrap(); + } + } + /// Instant counts select current membership, never accumulated observations. #[test] fn local_grouped_count_has_a_bindable_candidate() { diff --git a/data_plane/tests/support/issue_701_702_process.rs b/data_plane/tests/support/issue_701_702_process.rs index 4989e6bc8..e617e575b 100644 --- a/data_plane/tests/support/issue_701_702_process.rs +++ b/data_plane/tests/support/issue_701_702_process.rs @@ -1,5 +1,6 @@ //! Issue workloads execute their selected Planner DAG on the production HTTP path. use super::*; +use asap_types::physical_plan_codec::PhysicalPlanCodec; use control_plane::physical::{ compiler::{ BackendLocalPlanningInput, DeploymentPlanCompiler, BACKEND_REVISION, PLANNER_REVISION, @@ -54,7 +55,7 @@ fn queries() -> Vec<(String, u64, u64)> { )); queries.push((format!("quantile by(job)({q}, issue701_data)"), 1, 1)); } - for operation in ["sum", "count", "min", "max"] { + for operation in ["sum", "count", "avg", "min", "max"] { queries.push((format!("{operation}_over_time(issue701_data[5m])"), 300, 30)); } for operation in ["sum", "count", "avg"] { @@ -76,7 +77,7 @@ fn queries() -> Vec<(String, u64, u64)> { queries } -// A single mixed workload covers moving windows, current series and extrema, +// A single mixed workload covers moving windows, current series, minimum/average, // 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() { @@ -311,6 +312,7 @@ async fn issue_702_individual_queries_execute_without_fallback() { "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)", @@ -367,10 +369,9 @@ fn issue_701_702_uncertified_ratios_require_exact_fallback() { } } -// A per-series average is forwarded exactly while its sum and count stay warm; -// Planner does not yet match per-series rows in a division. +// Finite input can overflow sum; the installed average must fall back while zero stays warm. #[tokio::test] -async fn temporal_average_forwards_exactly_while_sum_and_count_stay_warm() { +async fn temporal_average_overflow_falls_back_after_state_is_warm() { let native = std::env::var("ASAP_CURRENT_SERIES_PROMETHEUS_URL").ok(); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let mock_url = format!("http://{}", listener.local_addr().unwrap()); @@ -398,18 +399,26 @@ async fn temporal_average_forwards_exactly_while_sum_and_count_stay_warm() { .into(); let snapshot = quote_snapshot_for_test(serde_json::from_value(fixture).unwrap()); let plan = snapshot.clone().compile_promql().unwrap(); - let average = plan - .query_plan - .entries - .values() - .find(|entry| entry.canonical_query.starts_with("avg_over_time")) - .unwrap(); assert!( - matches!( - average.nodes.get(&average.root), - Some(control_plane::query_plan::QueryPlanNode::ExactFallback { .. }) - ), - "a per-series average has no Planner-compiled local plan: {average:?}" + plan.query_plan + .entries + .values() + .flat_map(|entry| entry.nodes.values()) + .any(|node| { + let control_plane::query_plan::QueryPlanNode::PhysicalFragment { dag, .. } = node + else { + return false; + }; + asap_physical_operators::physical_planner::CompiledPhysicalDag::decode(dag) + .unwrap(); + let document: serde_json::Value = serde_json::from_slice(dag).unwrap(); + document["nodes"].as_object().unwrap().values().any(|node| { + node["Operator"]["operator"]["kind"]["VectorBinary"]["operator"] + ["checked_finite_division"] + == true + }) + }), + "average must retain its native finite-division contract" ); let output = tempfile::tempdir().unwrap(); let path = output.path().join("snapshot.json"); @@ -462,20 +471,20 @@ async fn temporal_average_forwards_exactly_while_sum_and_count_stay_warm() { .await; } let query = "avg_over_time(average_overflow[5s])"; - let params = [("query", query.to_string()), ("time", at.to_string())]; - let actual: Value = client - .get(format!("{backend}/api/v1/query")) - .query(¶ms) - .send() - .await - .unwrap() - .json() - .await - .unwrap(); - assert!(!is_warm(&actual), "average must be forwarded: {actual}"); - if let Some(url) = &native { - let expected: Value = client - .get(format!("{url}/api/v1/query")) + if value == 0.0 { + let result = wait_for_issue_warm_instant( + &client, + &backend, + query, + at, + &output.path().join("query_engine.log"), + ) + .await; + assert_eq!(first_value(&result, "value"), Some(0.0)); + } else { + let params = [("query", query.to_string()), ("time", at.to_string())]; + let actual: Value = client + .get(format!("{backend}/api/v1/query")) .query(¶ms) .send() .await @@ -483,9 +492,23 @@ async fn temporal_average_forwards_exactly_while_sum_and_count_stay_warm() { .json() .await .unwrap(); - assert_eq!(actual["data"], expected["data"]); - } else { + assert!( + !is_warm(&actual), + "overflowed average must fall back: {actual}" + ); assert_eq!(first_value(&actual, "value"), Some(1e308), "{actual}"); + if let Some(url) = &native { + let expected: Value = client + .get(format!("{url}/api/v1/query")) + .query(¶ms) + .send() + .await + .unwrap() + .json() + .await + .unwrap(); + assert_eq!(actual["data"], expected["data"]); + } } } mock.abort(); From ca7d885557dc146dc79d0deaf890855ba1138cb3 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 08:27:11 +0000 Subject: [PATCH 2/6] chore: repin ASAPPlanner to integration rev e7c64ab Picks up per-series Binary compilation over stored readouts, the without-aggregation state column fix, and Prometheus-compensated sums. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 12 ++++++------ Cargo.toml | 10 +++++----- 2 files changed, 11 insertions(+), 11 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index ea23af070..6d9b9eeb4 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=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib)", @@ -376,7 +376,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" 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=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -396,7 +396,7 @@ dependencies = [ [[package]] name = "asap-physical-operators" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?rev=5f03ccbd798ed5fec62bdd839bcb331123cab369)", @@ -410,12 +410,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=ca422df610009bdcc4f05b9becf4b2f5ee0c56b2#ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e7c64ab20a18592c462357c9fff954892b0d4a08#e7c64ab20a18592c462357c9fff954892b0d4a08" dependencies = [ "serde", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index 63cf197d0..ecc80e668 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -16,10 +16,10 @@ version = "0.1.0" [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 = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } -asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } -asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } -asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } +planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "e7c64ab20a18592c462357c9fff954892b0d4a08" } +asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "e7c64ab20a18592c462357c9fff954892b0d4a08" } +asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "e7c64ab20a18592c462357c9fff954892b0d4a08" } +asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "e7c64ab20a18592c462357c9fff954892b0d4a08" } # Shared external deps (used by 2+ crates) serde = { version = "1.0", features = ["derive"] } @@ -39,7 +39,7 @@ arc-swap = "1.7" reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] } # Internal crates -asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" } +asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "e7c64ab20a18592c462357c9fff954892b0d4a08" } asap_sketch_codec = { path = "crates/asap_sketch_codec" } asap_summary_state = { path = "crates/asap_summary_state" } asap_types = { path = "crates/asap_types" } From 4d47809d6e1c80c4efc2e0ffafee168f5d8e17bb Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 08:54:08 +0000 Subject: [PATCH 3/6] feat(data-plane): bind the series identity in Planner fragments Per-series Planner fragments read $promql_series_identity. Fill it from each readout series' labels and decode result labels from it. Co-Authored-By: Claude Opus 5.5 --- .../asap_query_engine/logical_dag.rs | 97 ++++++++++++++++++- .../logical_dag/native_values.rs | 39 +++++++- .../asap_query_engine/test_plan.rs | 16 ++- 3 files changed, 149 insertions(+), 3 deletions(-) diff --git a/data_plane/src/query_engines/asap_query_engine/logical_dag.rs b/data_plane/src/query_engines/asap_query_engine/logical_dag.rs index a5e57c3dd..d45feaae8 100644 --- a/data_plane/src/query_engines/asap_query_engine/logical_dag.rs +++ b/data_plane/src/query_engines/asap_query_engine/logical_dag.rs @@ -814,7 +814,9 @@ mod join_tests { #[cfg(test)] mod planner_computation_tests { use super::*; - use crate::query_engines::asap_query_engine::test_plan::planner_computed_entry; + use crate::query_engines::asap_query_engine::test_plan::{ + planner_computed_entry, planner_series_entry, + }; fn readout(values: &[(&str, f64)], at: u64) -> QueryResult { QueryResult::vector( @@ -919,6 +921,99 @@ mod planner_computation_tests { assert_eq!(stats.summary_readout_evaluations, 1); } + fn series_readout(series: &[(&[(&str, &str)], f64)], at: u64) -> QueryResult { + QueryResult::vector( + series + .iter() + .map(|(labels, value)| { + InstantVectorElement::new( + KeyByLabelValues::new_with_labels( + labels.iter().map(|(_, v)| (*v).into()).collect(), + ), + *value, + ) + .with_label_keys_override(labels.iter().map(|(k, _)| (*k).into()).collect()) + }) + .collect(), + at, + ) + } + + // Per-series division over two readouts matches series on their labels + // without the metric name, drops unmatched series and the metric name. + #[test] + fn per_series_ratio_matches_readout_series() { + let entry = planner_series_entry("rate(a[5m]) / rate(b[5m])"); + assert!(matches!( + entry.nodes[&entry.root], + QueryPlanNode::PhysicalFragment { .. } + )); + let [a, b] = readouts(&entry).try_into().unwrap(); + let (result, _) = execute_installed(&entry, &BTreeMap::new(), 300_000, |id, at| { + Ok(if id == a { + series_readout( + &[ + (&[("__name__", "a"), ("job", "api")], 6.0), + (&[("__name__", "a"), ("job", "db")], 1.0), + ], + at, + ) + } else { + assert_eq!(id, b); + series_readout( + &[ + (&[("__name__", "b"), ("job", "api")], 3.0), + (&[("__name__", "b"), ("job", "web")], 1.0), + ], + at, + ) + }) + }) + .unwrap(); + let QueryResult::Vector(result) = result else { + panic!("instant vector expected") + }; + assert_eq!( + result + .values + .iter() + .map(|point| (point.labels.labels.clone(), point.value)) + .collect::>(), + vec![(vec!["api".to_string()], 2.0)] + ); + } + + // A literal operand scales every per-series readout value. + #[test] + fn per_series_scalar_arithmetic_scales_each_series() { + let entry = planner_series_entry("rate(a[5m]) * 2"); + let [a] = readouts(&entry).try_into().unwrap(); + let (result, _) = execute_installed(&entry, &BTreeMap::new(), 300_000, |id, at| { + assert_eq!(id, a); + Ok(series_readout( + &[(&[("job", "api")], 1.5), (&[("job", "db")], 4.0)], + at, + )) + }) + .unwrap(); + let QueryResult::Vector(result) = result else { + panic!("instant vector expected") + }; + let mut values = result + .values + .iter() + .map(|point| (point.labels.labels.clone(), point.value)) + .collect::>(); + values.sort_by(|a, b| a.0.cmp(&b.0)); + assert_eq!( + values, + vec![ + (vec!["api".to_string()], 3.0), + (vec!["db".to_string()], 8.0) + ] + ); + } + // Source failures keep their routing classification across the shared runtime. #[test] fn source_error_classification_survives_execution() { diff --git a/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs b/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs index c6c711562..4447f17b2 100644 --- a/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs +++ b/data_plane/src/query_engines/asap_query_engine/logical_dag/native_values.rs @@ -185,7 +185,10 @@ pub(super) fn physical( at: i64, context: dag::RunContext, ) -> Result { - use asap_physical_operators::physical_planner::{CompiledPhysicalDag, Source}; + use asap_physical_operators::physical_planner::{ + promql_rows::{decode_series_identity, encode_series_identity, SERIES_IDENTITY_COLUMN}, + CompiledPhysicalDag, Source, + }; use futures::{FutureExt, StreamExt}; use std::collections::{BTreeMap, VecDeque}; let compiled = CompiledPhysicalDag::decode(encoded)?; @@ -222,6 +225,13 @@ pub(super) fn physical( Err(asap_physical_operators::Error::Invalid("protocol sample cannot represent the required Int64 input exactly".into()).into()) } } + // Planner matches per-series rows by their + // complete label set. + SummaryFamilyType::Plain(DataType::Utf8) + if field.name == SERIES_IDENTITY_COLUMN => + { + Ok(Value::Utf8(encode_series_identity(labels)?.into())) + } SummaryFamilyType::Plain(DataType::Utf8) => { Ok(labels.get(&field.name).map_or_else( || { @@ -259,6 +269,12 @@ pub(super) fn physical( ); } let output_schema = compiled.output_contract(compiled.roots()[0])?.schema; + // A per-series result names its series by identity; its other label + // columns are projections of that identity. + let output_identity = output_schema + .fields + .iter() + .position(|field| field.name == SERIES_IDENTITY_COLUMN); let graph = compiled.instantiate(sources)?; let mut streams = graph.execute(compiled.roots(), context)?; if streams.len() != 1 { @@ -273,6 +289,27 @@ pub(super) fn physical( match stream.next().now_or_never() { Some(Some(batch)) => { for row in batch?.rows() { + if let Some(identity) = output_identity { + let Value::Utf8(encoded) = &row[identity] else { + return Err(miss("invalid physical series identity")); + }; + let sample = output_schema + .fields + .iter() + .zip(row) + .find_map(|(field, cell)| match cell { + Value::Float64(value) + if field.dtype + == SummaryFamilyType::Plain(DataType::Float64) => + { + Some(*value) + } + _ => None, + }) + .ok_or_else(|| miss("native vector output has no numeric value"))?; + result.push((decode_series_identity(encoded)?, sample)); + continue; + } if row_input.is_none() { let mut labels = Labels::new(); let mut sample = None; diff --git a/data_plane/src/query_engines/asap_query_engine/test_plan.rs b/data_plane/src/query_engines/asap_query_engine/test_plan.rs index 617c7794a..fa3ccfcc4 100644 --- a/data_plane/src/query_engines/asap_query_engine/test_plan.rs +++ b/data_plane/src/query_engines/asap_query_engine/test_plan.rs @@ -196,12 +196,26 @@ pub(super) fn bound_reference( /// An exact PromQL query lowered by the control plane: backend readouts over /// fixture bindings, with Planner-compiled computation above them. pub(super) fn planner_computed_entry(query: &str) -> QueryPlanEntry { + computed_entry(query, false) +} + +/// [`planner_computed_entry`] selected over the identity-typed root, as the +/// control plane does for per-series arithmetic. +pub(super) fn planner_series_entry(query: &str) -> QueryPlanEntry { + computed_entry(query, true) +} + +fn computed_entry(query: &str, series_identity: bool) -> QueryPlanEntry { let canonical = canonical_promql(query).unwrap(); - let expr = control_plane::query_parser::parse_query_expr_canonical( + let mut expr = control_plane::query_parser::parse_query_expr_canonical( &canonical, planner_types::types::AccuracyTarget::Exact, ) .unwrap(); + if series_identity { + expr = asap_physical_operators::physical_planner::promql_rows::with_series_identity(&expr) + .unwrap(); + } let selected = control_plane::planner_selection::select_query( &expr, &control_plane::physical::post_asap::cost_model::ControlPlaneCostModel::new( From 56a245de636b1c5bdc49c8798bcccf0bd7c97f00 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 08:54:08 +0000 Subject: [PATCH 4/6] feat: select per-series arithmetic over identity-typed roots When a selected query-time computation does not compile over canonical readouts, select the workload again over roots typed with the series identity. Prefer it when it deploys; otherwise keep it as a candidate forest. The per-series test now checks the snapshot path, and the overflow test finds the checked-division flag anywhere in the fragment. Co-Authored-By: Claude Opus 5.5 --- control_plane/src/physical/compiler.rs | 174 +++++++++++++----- .../tests/support/issue_701_702_process.rs | 17 +- .../control-plane/physical-compiler.md | 9 + 3 files changed, 150 insertions(+), 50 deletions(-) diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 98997cd75..057bfff5d 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -1024,11 +1024,78 @@ impl BackendLocalPlanningInput { )?; query.retain_physical_candidate()?; } + let mut planner_candidate_forests = Vec::new(); + // Planner matches per-series rows in query-time arithmetic only by the + // series identity. When a selected computation does not compile over + // the canonical roots' readouts but does over identity-typed roots, + // the workload selected over typed roots is an alternative. Every + // query is selected over typed roots so the states they share keep + // one semantic definition. + let mut typed_workload = None; + let uncompiled = |root: &Rc| { + crate::query_plan::is_query_computation(root) + && crate::query_plan::compile_query_computation(root).is_err() + }; + if queries + .iter() + .any(|query| uncompiled(&query.selected_plan_root)) + { + let typed_roots: Vec<_> = canonical_roots + .iter() + .map(|root| { + asap_physical_operators::physical_planner::promql_rows::with_series_identity( + root, + ) + .map_or_else(|_| Rc::clone(root), Rc::new) + }) + .collect(); + let retyped: Vec<_> = typed_roots + .iter() + .zip(&canonical_roots) + .map(|(typed, canonical)| !Rc::ptr_eq(typed, canonical)) + .collect(); + let mut typed = queries.clone(); + if let Some(trace) = select_logical_roots_with_scoped_evidence_and_trace( + &mut typed, + typed_roots, + &topk_evidence_by_id, + &scoped_evidence_by_id, + &exact_costs_by_id, + self.physical_inputs.erp.as_ref(), + self.environment.observed_at_unix_ms, + ) + .ok() + .filter(|_| { + queries.iter().zip(&typed).any(|(canonical, typed)| { + uncompiled(&canonical.selected_plan_root) + && !uncompiled(&typed.selected_plan_root) + }) && typed.iter_mut().all(|query| { + prepare_window_implementations( + query, + &self.physical_inputs.window_cost_model, + self.environment.target, + self.physical_inputs.query_retention_margin_ms, + ) + .is_ok() + }) + }) { + // Native realizations of typed roots are bound only through + // their lifecycle placement below. + for (query, retyped) in typed.iter_mut().zip(retyped) { + if retyped { + query.retain(None)?; + } else { + query.retain_physical_candidate()?; + } + } + planner_selection_trace.extend(trace); + typed_workload = Some(typed); + } + } // Native physical realizations need the complete series identity in // their rows. Planner's PlanSpace proposes them for the identity-typed // root; each is a logical alternative whose readout-built states are // placed by lifecycle, then substituted into the preferred workload. - let mut planner_candidate_forests = Vec::new(); for (index, root) in canonical_roots.iter().enumerate() { let Ok(typed) = asap_physical_operators::physical_planner::promql_rows::with_series_identity(root) @@ -1132,30 +1199,42 @@ impl BackendLocalPlanningInput { } } // Composable lowering residualizes unsafe leaves individually; retain Planner siblings. - Ok(( - PhysicalCompilationRequest { - planner_candidate_forests, - planner_selection_trace: planner_selection_trace.into(), - allow_mixed_summary_and_exact_execution: true, - require_backend_local_execution: self - .physical_inputs - .require_backend_local_execution, - 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, - erp: self.physical_inputs.erp, - planner_revision: PLANNER_REVISION.into(), - scrape_interval_ms: Some(self.physical_inputs.scrape_interval_ms), - query_retention_margin_ms: self.physical_inputs.query_retention_margin_ms, - retained_summary_memory_budget_bytes: Some( - self.physical_inputs.retained_summary_memory_budget_bytes, - ), - }, - self.environment, - )) + let mut request = PhysicalCompilationRequest { + planner_candidate_forests, + planner_selection_trace: planner_selection_trace.into(), + allow_mixed_summary_and_exact_execution: true, + require_backend_local_execution: self.physical_inputs.require_backend_local_execution, + 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, + erp: self.physical_inputs.erp, + planner_revision: PLANNER_REVISION.into(), + scrape_interval_ms: Some(self.physical_inputs.scrape_interval_ms), + query_retention_margin_ms: self.physical_inputs.query_retention_margin_ms, + retained_summary_memory_budget_bytes: Some( + self.physical_inputs.retained_summary_memory_budget_bytes, + ), + }; + // The typed workload is preferred when it deploys; typing can make an + // unrelated root infeasible, so it otherwise remains an alternative. + if let Some(typed) = typed_workload { + let mut preferred = request.clone(); + preferred.queries = typed.clone(); + preferred.planner_candidate_forests.clear(); + if DeploymentPlanCompiler + .compile_promql(preferred, self.environment.clone()) + .is_ok() + { + let canonical = std::mem::replace(&mut request.queries, typed); + request.planner_candidate_forests.insert(0, canonical); + } else { + request.planner_candidate_forests.insert(0, typed); + } + } + Ok((request, self.environment)) } } @@ -5151,31 +5230,38 @@ pub(crate) mod tests { assert!(plan.precompute_plan.materializations.is_empty()); } - // Per-series arithmetic over stored readouts runs as one Planner fragment - // over those readouts, not as an exact fallback. + // Per-series arithmetic over stored readouts has a candidate that runs + // as one Planner fragment over those readouts, not as an exact fallback. #[test] fn per_series_arithmetic_executes_as_a_planner_fragment() { for query in [ - "avg_over_time(a[1m])", + "avg_over_time(data[5m])", "rate(a[5m]) / rate(b[5m])", - "rate(a[5m]) * 2", + "rate(data[5m]) * 2", "sum_over_time(a[1m]) + sum_over_time(a{job=\"x\"}[1m])", ] { - let mut environment = environment(10_000); - environment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; - environment.target_collector_ids.clear(); - let plan = DeploymentPlanCompiler - .compile_promql(request("arithmetic", query), environment) - .unwrap(); - let entry = plan.query_plan.lookup(query).unwrap(); - assert!( - matches!( - &entry.nodes[&entry.root], - crate::query_plan::QueryPlanNode::PhysicalFragment { .. } - ), - "{query}: {entry:?}" - ); - assert!(!entry.materialization_bindings().is_empty(), "{query}"); + let mut snapshot = planning_snapshot(); + let entry = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0]; + entry.query = Query(query.into()); + entry.requirements.accuracy = AccuracyRequirement::Explicit(AccuracyTarget::Exact); + let (request, environment) = snapshot.into_physical_compilation_request().unwrap(); + let warm = + super::super::workload_cost::enumerate_exact_and_materialized_candidates(request) + .unwrap() + .into_iter() + .filter_map(|candidate| { + DeploymentPlanCompiler + .compile_promql(candidate, environment.clone()) + .ok() + }) + .any(|plan| { + let entry = plan.query_plan.lookup(query).unwrap(); + matches!( + &entry.nodes[&entry.root], + crate::query_plan::QueryPlanNode::PhysicalFragment { .. } + ) && !entry.materialization_bindings().is_empty() + }); + assert!(warm, "{query}"); } } diff --git a/data_plane/tests/support/issue_701_702_process.rs b/data_plane/tests/support/issue_701_702_process.rs index e617e575b..7d2e32b0e 100644 --- a/data_plane/tests/support/issue_701_702_process.rs +++ b/data_plane/tests/support/issue_701_702_process.rs @@ -411,12 +411,17 @@ async fn temporal_average_overflow_falls_back_after_state_is_warm() { }; asap_physical_operators::physical_planner::CompiledPhysicalDag::decode(dag) .unwrap(); - let document: serde_json::Value = serde_json::from_slice(dag).unwrap(); - document["nodes"].as_object().unwrap().values().any(|node| { - node["Operator"]["operator"]["kind"]["VectorBinary"]["operator"] - ["checked_finite_division"] - == true - }) + fn checked(value: &Value) -> bool { + match value { + Value::Object(map) => map.iter().any(|(key, value)| { + (key == "checked_finite_division" && value == &Value::Bool(true)) + || checked(value) + }), + Value::Array(items) => items.iter().any(checked), + _ => false, + } + } + checked(&serde_json::from_slice(dag).unwrap()) }), "average must retain its native finite-division contract" ); diff --git a/docs/developer_docs/control-plane/physical-compiler.md b/docs/developer_docs/control-plane/physical-compiler.md index 1769efaf5..c79738cd0 100644 --- a/docs/developer_docs/control-plane/physical-compiler.md +++ b/docs/developer_docs/control-plane/physical-compiler.md @@ -235,6 +235,15 @@ PromQL selector, keeps no state and forwards the whole query as whether PromQL drops `__name__` from its result, since Planner keeps it in the series identity. +Planner matches per-series rows in query-time arithmetic (`avg_over_time` as +sum/count, `rate(a) / rate(b)`, `rate(x) * 2`) only by the series identity +column. When a selected computation does not compile over the canonical roots, +the workload is selected again over roots typed with the series identity +(`promql_rows::with_series_identity`). Every query is retyped so shared states +keep one semantic definition. The typed workload is preferred if it deploys; +otherwise it remains a candidate forest. The adapter fills the identity column +from each readout series' labels and decodes result labels from it. + Graph traversal is separate from node definitions and store semantics. Activation validates roots, edges, bindings, reachability, and cycles. The shared physical DAG runtime creates one producer per reachable node and From 25279908e3161242ada39bd57ead329e6c7da350 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 09:46:13 +0000 Subject: [PATCH 5/6] fix: add the typed reselection only as a candidate forest Remove the trial compile that made the identity-typed workload the primary request: every forest reaches the same enumeration and pricing, and the trial compiled a request that differed from the deployed one. Reselect only on Planner's missing-identity error, reject a typed forest whose unretyped root shares a state with a retyped root, and record rejections and tagged typed-selection entries in the selection trace. The overflow process test now prices stateful candidates cheapest instead of relying on forest order. Co-Authored-By: Claude Opus 5.5 --- control_plane/src/physical/compiler.rs | 296 +++++++++++++++--- .../tests/support/issue_701_702_process.rs | 45 ++- .../control-plane/physical-compiler.md | 16 +- 3 files changed, 299 insertions(+), 58 deletions(-) diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 057bfff5d..5b53924fc 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -782,6 +782,82 @@ impl AccuracyEvidenceProvider for QueryEvidence<'_> { } } +/// Planner's row-binary error when per-series rows carry no series identity. +/// Planner reports it only as text, so reselection matches the message. +const SERIES_IDENTITY_REQUIRED: &str = + "row binary requires grouped rows or rows with a series identity"; + +/// Admit the workload selected over identity-typed roots as a candidate +/// forest, or give the reason it is rejected. +fn typed_reselection_forest( + typed: &mut [QueryCompilationInput], + canonical: &[QueryCompilationInput], + retyped: &[bool], + needs_identity: &dyn Fn(&Rc) -> bool, + inputs: &BackendLocalPhysicalInputs, + target: PhysicalDeploymentTarget, +) -> Result<(), String> { + if !canonical.iter().zip(&*typed).any(|(canonical, typed)| { + needs_identity(&canonical.selected_plan_root) && !needs_identity(&typed.selected_plan_root) + }) { + return Err("no computation compiles over identity-typed roots".into()); + } + for query in typed.iter_mut() { + prepare_window_implementations( + query, + &inputs.window_cost_model, + target, + inputs.query_retention_margin_ms, + ) + .map_err(|error| format!("{}: {error}", query.query_id))?; + } + // A root that cannot be retyped keeps its canonical states. Sharing one + // with a retyped root would give one deployed output two definitions. + let fingerprints = |retyped_side: bool| -> Result, String> { + let mut all = BTreeSet::new(); + for (query, _) in typed + .iter() + .zip(retyped) + .filter(|(_, &retyped)| retyped == retyped_side) + { + all.extend(state_fingerprints(query, target)?); + } + Ok(all) + }; + if !fingerprints(true)?.is_disjoint(&fingerprints(false)?) { + return Err("a root without a series identity shares state with a retyped root".into()); + } + // Native realizations of typed roots are bound only through their + // lifecycle placement. + for (query, &retyped) in typed.iter_mut().zip(retyped) { + if retyped { + query.retain(None) + } else { + query.retain_physical_candidate() + } + .map_err(|error| error.to_string())?; + } + Ok(()) +} + +/// Deployed-output fingerprints of the states a query's selected root reads. +fn state_fingerprints( + query: &QueryCompilationInput, + target: PhysicalDeploymentTarget, +) -> Result, String> { + collect_selected_materializations(&query.selected_plan_root, true)? + .iter() + .map(|state| { + scoped_materialization( + &physical_aggregation(query, state, query.query_id.clone(), target), + &state.node, + ) + .map(|config| config.policy_fingerprint()) + .map_err(|error| error.to_string()) + }) + .collect() +} + #[derive(Debug, Default)] pub struct DeploymentPlanCompiler; @@ -1026,20 +1102,24 @@ impl BackendLocalPlanningInput { } let mut planner_candidate_forests = Vec::new(); // Planner matches per-series rows in query-time arithmetic only by the - // series identity. When a selected computation does not compile over - // the canonical roots' readouts but does over identity-typed roots, - // the workload selected over typed roots is an alternative. Every - // query is selected over typed roots so the states they share keep - // one semantic definition. - let mut typed_workload = None; - let uncompiled = |root: &Rc| { + // series identity. When a selected computation fails to compile for + // lack of it, the workload selected over identity-typed roots is one + // more forest; deployment pricing chooses among all forests. + let needs_identity = |root: &Rc| { crate::query_plan::is_query_computation(root) - && crate::query_plan::compile_query_computation(root).is_err() + && crate::query_plan::compile_query_computation(root) + .is_err_and(|error| error.to_string().contains(SERIES_IDENTITY_REQUIRED)) }; if queries .iter() - .any(|query| uncompiled(&query.selected_plan_root)) + .any(|query| needs_identity(&query.selected_plan_root)) { + let reject = |reason: String| { + serde_json::json!({ + "stage": "planner.series_identity_reselection", + "status": "rejected", "reason": reason, + }) + }; let typed_roots: Vec<_> = canonical_roots .iter() .map(|root| { @@ -1055,7 +1135,7 @@ impl BackendLocalPlanningInput { .map(|(typed, canonical)| !Rc::ptr_eq(typed, canonical)) .collect(); let mut typed = queries.clone(); - if let Some(trace) = select_logical_roots_with_scoped_evidence_and_trace( + match select_logical_roots_with_scoped_evidence_and_trace( &mut typed, typed_roots, &topk_evidence_by_id, @@ -1063,33 +1143,27 @@ impl BackendLocalPlanningInput { &exact_costs_by_id, self.physical_inputs.erp.as_ref(), self.environment.observed_at_unix_ms, - ) - .ok() - .filter(|_| { - queries.iter().zip(&typed).any(|(canonical, typed)| { - uncompiled(&canonical.selected_plan_root) - && !uncompiled(&typed.selected_plan_root) - }) && typed.iter_mut().all(|query| { - prepare_window_implementations( - query, - &self.physical_inputs.window_cost_model, + ) { + Err(error) => planner_selection_trace.push(reject(error.to_string())), + Ok(trace) => { + planner_selection_trace.extend(trace.into_iter().map(|mut event| { + if let Some(event) = event.as_object_mut() { + event.insert("reselection".into(), "series_identity".into()); + } + event + })); + match typed_reselection_forest( + &mut typed, + &queries, + &retyped, + &needs_identity, + &self.physical_inputs, self.environment.target, - self.physical_inputs.query_retention_margin_ms, - ) - .is_ok() - }) - }) { - // Native realizations of typed roots are bound only through - // their lifecycle placement below. - for (query, retyped) in typed.iter_mut().zip(retyped) { - if retyped { - query.retain(None)?; - } else { - query.retain_physical_candidate()?; + ) { + Ok(()) => planner_candidate_forests.push(typed), + Err(reason) => planner_selection_trace.push(reject(reason)), } } - planner_selection_trace.extend(trace); - typed_workload = Some(typed); } } // Native physical realizations need the complete series identity in @@ -1199,7 +1273,7 @@ impl BackendLocalPlanningInput { } } // Composable lowering residualizes unsafe leaves individually; retain Planner siblings. - let mut request = PhysicalCompilationRequest { + let request = PhysicalCompilationRequest { planner_candidate_forests, planner_selection_trace: planner_selection_trace.into(), allow_mixed_summary_and_exact_execution: true, @@ -1218,22 +1292,6 @@ impl BackendLocalPlanningInput { self.physical_inputs.retained_summary_memory_budget_bytes, ), }; - // The typed workload is preferred when it deploys; typing can make an - // unrelated root infeasible, so it otherwise remains an alternative. - if let Some(typed) = typed_workload { - let mut preferred = request.clone(); - preferred.queries = typed.clone(); - preferred.planner_candidate_forests.clear(); - if DeploymentPlanCompiler - .compile_promql(preferred, self.environment.clone()) - .is_ok() - { - let canonical = std::mem::replace(&mut request.queries, typed); - request.planner_candidate_forests.insert(0, canonical); - } else { - request.planner_candidate_forests.insert(0, typed); - } - } Ok((request, self.environment)) } } @@ -5265,6 +5323,142 @@ pub(crate) mod tests { } } + fn exact_workload_snapshot(queries: &[&str]) -> BackendLocalPlanningInput { + let mut snapshot = planning_snapshot(); + let entries = snapshot.query_workload.repeating_queries.as_mut().unwrap(); + entries[0].requirements.accuracy = AccuracyRequirement::Explicit(AccuracyTarget::Exact); + let template = entries[0].clone(); + entries.clear(); + for query in queries { + let mut entry = template.clone(); + entry.query = Query((*query).into()); + entries.push(entry); + } + snapshot + } + + fn needs_series_identity(root: &Rc) -> bool { + crate::query_plan::compile_query_computation(root) + .is_err_and(|error| error.to_string().contains(SERIES_IDENTITY_REQUIRED)) + } + + // The identity-typed reselection is only a candidate forest: the primary + // workload stays the canonical selection, and its trace entries are tagged. + #[test] + fn typed_reselection_is_only_a_candidate_forest() { + let (request, _) = exact_workload_snapshot(&["avg_over_time(data[5m])"]) + .into_physical_compilation_request() + .unwrap(); + assert!(needs_series_identity( + &request.queries[0].selected_plan_root + )); + let typed = &request.planner_candidate_forests[0]; + assert!(crate::query_plan::compile_query_computation(&typed[0].selected_plan_root).is_ok()); + assert!(request + .planner_selection_trace + .iter() + .any(|event| event["reselection"] == "series_identity")); + } + + // Workloads mixing per-series arithmetic with aggregates of the same + // source keep a candidate that serves every query from stored state. + #[test] + fn mixed_per_series_workloads_have_all_warm_candidates() { + for queries in [ + ["avg_over_time(data[5m])", "sum by (job) (rate(data[5m]))"], + ["rate(data[5m]) * 2", "sum(rate(data[5m]))"], + ["avg_over_time(data[5m])", "sum_over_time(data[5m])"], + ] { + let (request, environment) = exact_workload_snapshot(&queries) + .into_physical_compilation_request() + .unwrap(); + let warm = + super::super::workload_cost::enumerate_exact_and_materialized_candidates(request) + .unwrap() + .into_iter() + .filter_map(|candidate| { + DeploymentPlanCompiler + .compile_promql(candidate, environment.clone()) + .ok() + }) + .any(|plan| { + queries.iter().all(|query| { + !plan + .query_plan + .lookup(query) + .unwrap() + .materialization_bindings() + .is_empty() + }) + }); + assert!(warm, "{queries:?}"); + } + } + + // A state shared by a retyped and a canonical root has the same deployed + // fingerprint either way, so it must keep one definition. + #[test] + fn shared_state_fingerprint_is_independent_of_series_identity() { + let (request, _) = + exact_workload_snapshot(&["avg_over_time(data[5m])", "sum_over_time(data[5m])"]) + .into_physical_compilation_request() + .unwrap(); + let target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; + let canonical = state_fingerprints(&request.queries[1], target).unwrap(); + let typed = &request.planner_candidate_forests[0]; + assert!(!canonical.is_empty()); + assert_eq!(canonical, state_fingerprints(&typed[1], target).unwrap()); + assert!(canonical.is_subset(&state_fingerprints(&typed[0], target).unwrap())); + } + + // A typed workload whose unretyped root reads a state of a retyped root + // is rejected rather than deploying two definitions of that state. + #[test] + fn typed_reselection_rejects_state_shared_with_unretyped_root() { + let (request, _) = + exact_workload_snapshot(&["avg_over_time(data[5m])", "sum_over_time(data[5m])"]) + .into_physical_compilation_request() + .unwrap(); + let inputs = planning_snapshot().physical_inputs; + let target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; + let mut typed = request.planner_candidate_forests[0].clone(); + assert!(typed_reselection_forest( + &mut typed.clone(), + &request.queries, + &[true, true], + &needs_series_identity, + &inputs, + target, + ) + .is_ok()); + // The canonical `sum_over_time` stands in for a root that cannot be retyped. + typed[1] = request.queries[1].clone(); + let error = typed_reselection_forest( + &mut typed, + &request.queries, + &[true, false], + &needs_series_identity, + &inputs, + target, + ) + .unwrap_err(); + assert!(error.contains("shares state"), "{error}"); + } + + // Reselection is gated on the missing series identity: a computation that + // fails to compile for another reason is not reselected. + #[test] + fn reselection_requires_the_series_identity_error() { + let query = "quantile_over_time(0.9,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(query.into()); + let (request, _) = snapshot.into_physical_compilation_request().unwrap(); + assert!(request.planner_selection_trace.iter().all(|event| { + event["reselection"].is_null() + && event["stage"] != "planner.series_identity_reselection" + })); + } + // Quantile rank error does not certify relative error of a quotient. #[test] fn uncertified_quantile_ratios_retain_exact_execution() { diff --git a/data_plane/tests/support/issue_701_702_process.rs b/data_plane/tests/support/issue_701_702_process.rs index 7d2e32b0e..d1f9907e7 100644 --- a/data_plane/tests/support/issue_701_702_process.rs +++ b/data_plane/tests/support/issue_701_702_process.rs @@ -369,6 +369,49 @@ fn issue_701_702_uncertified_ratios_require_exact_fallback() { } } +/// Price candidates that read stored state for every query below the rest, +/// so selection does not depend on candidate order. +fn quote_stateful_snapshot(mut snapshot: BackendLocalPlanningInput) -> BackendLocalPlanningInput { + let (request, environment) = snapshot + .clone() + .into_physical_compilation_request() + .unwrap(); + let quotes = workload_cost::enumerate_exact_and_materialized_candidates(request) + .unwrap() + .into_iter() + .filter_map(|candidate| { + let plan = DeploymentPlanCompiler + .compile_promql(candidate.clone(), environment.clone()) + .ok()?; + let stateful = plan + .query_plan + .entries + .values() + .all(|entry| !entry.materialization_bindings().is_empty()); + let manifest = workload_cost::manifest(&plan, &candidate.queries).unwrap(); + Some(WorkloadQuote { + unit_costs: manifest + .components + .keys() + .map(|key| (key.clone(), if stateful { 1.0 } else { 1e12 })) + .collect(), + manifest, + executable: true, + }) + }) + .collect(); + snapshot.workload_cost_evidence = Some(WorkloadCostEvidence { + backend_revision: BACKEND_REVISION.into(), + planner_revision: PLANNER_REVISION.into(), + data_snapshot_id: "issue-701-702-process".into(), + model_version: "synthetic-correctness-quotes".into(), + observed_at_unix_ms: environment.observed_at_unix_ms, + valid_for_ms: environment.max_evidence_age_ms, + quotes, + }); + snapshot +} + // Finite input can overflow sum; the installed average must fall back while zero stays warm. #[tokio::test] async fn temporal_average_overflow_falls_back_after_state_is_warm() { @@ -397,7 +440,7 @@ async fn temporal_average_overflow_falls_back_after_state_is_warm() { }) .to_vec() .into(); - let snapshot = quote_snapshot_for_test(serde_json::from_value(fixture).unwrap()); + let snapshot = quote_stateful_snapshot(serde_json::from_value(fixture).unwrap()); let plan = snapshot.clone().compile_promql().unwrap(); assert!( plan.query_plan diff --git a/docs/developer_docs/control-plane/physical-compiler.md b/docs/developer_docs/control-plane/physical-compiler.md index c79738cd0..e99d41a2a 100644 --- a/docs/developer_docs/control-plane/physical-compiler.md +++ b/docs/developer_docs/control-plane/physical-compiler.md @@ -237,12 +237,16 @@ series identity. Planner matches per-series rows in query-time arithmetic (`avg_over_time` as sum/count, `rate(a) / rate(b)`, `rate(x) * 2`) only by the series identity -column. When a selected computation does not compile over the canonical roots, -the workload is selected again over roots typed with the series identity -(`promql_rows::with_series_identity`). Every query is retyped so shared states -keep one semantic definition. The typed workload is preferred if it deploys; -otherwise it remains a candidate forest. The adapter fills the identity column -from each readout series' labels and decodes result labels from it. +column. When a selected computation fails to compile for lack of that identity, +the workload is selected again over roots typed with it +(`promql_rows::with_series_identity`) and added as one more candidate forest; +deployment pricing chooses among all forests. Every query is retyped so shared +states keep one semantic definition. A root that cannot be retyped keeps its +canonical states, so the typed forest is rejected if such a root shares a state +with a retyped one. Rejections and the typed selection's trace entries (tagged +`"reselection": "series_identity"`) are kept in the selection trace. The +adapter fills the identity column from each readout series' labels and decodes +result labels from it. Graph traversal is separate from node definitions and store semantics. Activation validates roots, edges, bindings, reachability, and cycles. From c36b435a50ce5bfef5b305040c6132e8b9556530 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 09:46:13 +0000 Subject: [PATCH 6/6] test: check per-series scalar arithmetic drops readout metric names Co-Authored-By: Claude Opus 5.5 --- .../src/query_engines/asap_query_engine/logical_dag.rs | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/data_plane/src/query_engines/asap_query_engine/logical_dag.rs b/data_plane/src/query_engines/asap_query_engine/logical_dag.rs index d45feaae8..9fa65f67d 100644 --- a/data_plane/src/query_engines/asap_query_engine/logical_dag.rs +++ b/data_plane/src/query_engines/asap_query_engine/logical_dag.rs @@ -983,7 +983,8 @@ mod planner_computation_tests { ); } - // A literal operand scales every per-series readout value. + // A literal operand scales every per-series readout value and drops the + // metric name the readouts carry. #[test] fn per_series_scalar_arithmetic_scales_each_series() { let entry = planner_series_entry("rate(a[5m]) * 2"); @@ -991,7 +992,10 @@ mod planner_computation_tests { let (result, _) = execute_installed(&entry, &BTreeMap::new(), 300_000, |id, at| { assert_eq!(id, a); Ok(series_readout( - &[(&[("job", "api")], 1.5), (&[("job", "db")], 4.0)], + &[ + (&[("__name__", "a"), ("job", "api")], 1.5), + (&[("__name__", "a"), ("job", "db")], 4.0), + ], at, )) })