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" } diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 09ff1c1e5..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; @@ -1024,11 +1100,76 @@ 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 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_and(|error| error.to_string().contains(SERIES_IDENTITY_REQUIRED)) + }; + if queries + .iter() + .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| { + 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(); + match 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, + ) { + 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, + ) { + Ok(()) => planner_candidate_forests.push(typed), + Err(reason) => planner_selection_trace.push(reject(reason)), + } + } + } + } // 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 +1273,26 @@ 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 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, + ), + }; + Ok((request, self.environment)) } } @@ -5095,6 +5232,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,31 +5288,175 @@ 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 has a candidate that 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() { + fn per_series_arithmetic_executes_as_a_planner_fragment() { + for query in [ + "avg_over_time(data[5m])", + "rate(a[5m]) / rate(b[5m])", + "rate(data[5m]) * 2", + "sum_over_time(a[1m]) + sum_over_time(a{job=\"x\"}[1m])", + ] { + 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}"); + } + } + + fn exact_workload_snapshot(queries: &[&str]) -> BackendLocalPlanningInput { 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())); + 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. 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/src/query_engines/asap_query_engine/logical_dag.rs b/data_plane/src/query_engines/asap_query_engine/logical_dag.rs index a5e57c3dd..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 @@ -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,103 @@ 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 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"); + 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( + &[ + (&[("__name__", "a"), ("job", "api")], 1.5), + (&[("__name__", "a"), ("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( diff --git a/data_plane/tests/support/issue_701_702_process.rs b/data_plane/tests/support/issue_701_702_process.rs index 4989e6bc8..d1f9907e7 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,52 @@ 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. +/// 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_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()); @@ -396,20 +440,33 @@ async fn temporal_average_forwards_exactly_while_sum_and_count_stay_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(); - 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(); + 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" ); let output = tempfile::tempdir().unwrap(); let path = output.path().join("snapshot.json"); @@ -462,20 +519,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 +540,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(); diff --git a/docs/developer_docs/control-plane/physical-compiler.md b/docs/developer_docs/control-plane/physical-compiler.md index 1769efaf5..e99d41a2a 100644 --- a/docs/developer_docs/control-plane/physical-compiler.md +++ b/docs/developer_docs/control-plane/physical-compiler.md @@ -235,6 +235,19 @@ 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 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. The shared physical DAG runtime creates one producer per reachable node and