diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 9fb0a03e..c4967cc9 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -1204,7 +1204,7 @@ impl PhysicalCompiler { query_id: query.query_id.clone(), reason, })?; - if let Some(source) = immutable_materialization_source(&selected.node) { + if let Some(source) = immutable_materialization_sources(&selected.node) { if environment.target != PhysicalDeploymentTarget::BackendLocalRemoteWrite { return Err(CompileError::Query { query_id: query.query_id.clone(), @@ -1213,24 +1213,24 @@ impl PhysicalCompiler { }); } let compiled = executable_dags[query_index].as_ref().expect("compiled DAG"); - let source_node = - compiled - .node_ids - .node_id(&source) - .ok_or_else(|| CompileError::Query { + let mut frontiers = BTreeMap::new(); + let mut source_configs = Vec::new(); + for source in source { + let source_node = compiled.node_ids.node_id(&source).ok_or_else(|| { + CompileError::Query { query_id: query.query_id.clone(), reason: "derived source absent from selected DAG".into(), - })?; - let source_id = node_bindings - .get(&(query_index, source_node)) - .copied() - .ok_or_else(|| CompileError::Query { - query_id: query.query_id.clone(), - reason: "derived input source was not installed before its consumer" - .into(), + } })?; - let source_config: &asap_types::PrecomputeMaterialization = - compiled_materializations + let source_id = node_bindings + .get(&(query_index, source_node)) + .copied() + .ok_or_else(|| CompileError::Query { + query_id: query.query_id.clone(), + reason: "derived source was not installed before consumer".into(), + })?; + frontiers.insert(source_node, source_id.into()); + let config = compiled_materializations .iter() .find(|config: &&asap_types::PrecomputeMaterialization| { config.policy_fingerprint() == source_id @@ -1239,9 +1239,17 @@ impl PhysicalCompiler { query_id: query.query_id.clone(), reason: "derived source config missing".into(), })?; + if !source_configs.iter().any( + |existing: &&asap_types::PrecomputeMaterialization| { + existing.policy_fingerprint() == source_id + }, + ) { + source_configs.push(config); + } + } asap_types::precompute_plan::validated_source_window_cohort( &runtime_materialization, - &[source_config], + &source_configs, ) .map_err(|error| CompileError::Query { query_id: query.query_id.clone(), @@ -1269,9 +1277,7 @@ impl PhysicalCompiler { })?; runtime_materialization.derived_input = Some( asap_types::derived_input::DerivedInputIdentity::from_dag( - &document, - input_node, - &BTreeMap::from([(source_node, source_id.into())]), + &document, input_node, &frontiers, ) .map_err(|reason| CompileError::Query { query_id: query.query_id.clone(), @@ -1435,11 +1441,13 @@ impl PhysicalCompiler { reason: format!("promql-compatible identity: {error}"), })?; let binding = |node: &Rc, node_family: &SummaryFamilyType| -> Result { - summary_agg_metric(node).ok_or_else(|| { - crate::query_plan::QueryPlanError::Invalid( - "materialized node has no unique time-series source".into(), - ) - })?; + if immutable_materialization_sources(node).is_none() { + summary_agg_metric(node).ok_or_else(|| { + crate::query_plan::QueryPlanError::Invalid( + "materialized node has no unique time-series source".into(), + ) + })?; + } let node_id = executable_dags[query_index] .as_ref() .ok_or_else(|| { @@ -2582,10 +2590,10 @@ fn raw_time_series_input_contract( } } -/// The first immutable-input capability accepts one exact accumulator readout. -/// Population/window closure is checked by the installed runtime, not inferred -/// from the presence of this syntax. -fn immutable_materialization_source(node: &SummaryNode) -> Option> { +/// Immutable inputs may combine explicit exact accumulator readouts using +/// maintenance-time arithmetic. Population and window closure are checked by +/// the installed runtime, not inferred from the presence of this syntax. +fn immutable_materialization_sources(node: &SummaryNode) -> Option>> { use planner_types::post_asap::{ExactKind, ExecutionTiming, SummaryInputExpr}; let SummaryExpr::SummaryAgg { child, @@ -2607,27 +2615,46 @@ fn immutable_materialization_source(node: &SummaryNode) -> Option, sources: &mut Vec>) -> Option<()> { + match &node.expr { + SummaryExpr::BinaryOp { + lhs, + rhs, + operator, + timing: ExecutionTiming::MaintenanceTime, + } if operator.vector_match.is_none() + && matches!( + operator.kind, + planner_types::pre_asap::BinaryOpKind::Arithmetic(_) + ) => + { + collect(lhs, sources)?; + collect(rhs, sources)?; + } + SummaryExpr::ValueOperation { + child: source, + operation: planner_types::post_asap::ValueOperation::FinalizeExactAccumulator, + timing: ExecutionTiming::MaintenanceTime, + } if matches!(&source.expr, + SummaryExpr::SummaryAgg { family: SummaryFamilyType::ExactAggregate(ExactKind::Sum | ExactKind::Count, _), child, .. } + if matches!(child.expr, SummaryExpr::KeepPreAsap(_))) => + { + if !sources.iter().any(|old| Rc::ptr_eq(old, source)) { + sources.push(source.clone()); + } + } + _ => return None, + } + Some(()) } - Some(Rc::clone(source)) + let mut sources = Vec::new(); + collect(child, &mut sources)?; + Some(sources) } fn selected_input_contract(node: &SummaryNode) -> Result<(String, Option, String), String> { - if let Some(source) = immutable_materialization_source(node) { - materialization_leaf_contract(&source) + if let Some(source) = immutable_materialization_sources(node) { + materialization_leaf_contract(&source[0]) } else { materialization_leaf_contract(node) } @@ -2778,7 +2805,7 @@ fn materialization_consumers( &state.node, )?; let program = - immutable_materialization_source(&state.node).map(|_| Rc::as_ptr(&state.node)); + immutable_materialization_sources(&state.node).map(|_| Rc::as_ptr(&state.node)); if let Some(previous) = cohort_programs.insert(config.policy_fingerprint(), program) { if previous != program && (previous.is_some() || program.is_some()) { return Err(CompileError::Query { @@ -2894,8 +2921,10 @@ fn collect_selected_materializations( } else { None }; - if let Some(source) = immutable_materialization_source(node) { - walk(&source, None, composable, None, selected)?; + if let Some(source) = immutable_materialization_sources(node) { + for source in source { + walk(&source, None, composable, None, selected)?; + } } match &node.expr { SummaryExpr::CandidateTopK { @@ -2936,7 +2965,7 @@ fn collect_selected_materializations( } SummaryExpr::SummaryAgg { child, .. } if !matches!(child.expr, SummaryExpr::KeepPreAsap(_)) - && immutable_materialization_source(node).is_none() => {} + && immutable_materialization_sources(node).is_none() => {} SummaryExpr::SummaryAgg { family: SummaryFamilyType::ExactAggregate(planner_types::post_asap::ExactKind::Count, _), @@ -3714,8 +3743,8 @@ mod tests { let states = collect_selected_materializations(&workload.queries[0].post_asap, true).unwrap(); assert_eq!(states.len(), 2, "source and consumer must both be selected"); - assert!(immutable_materialization_source(&states[0].node).is_none()); - assert!(immutable_materialization_source(&states[1].node).is_some()); + assert!(immutable_materialization_sources(&states[0].node).is_none()); + assert!(immutable_materialization_sources(&states[1].node).is_some()); let mut deployment = environment(10_000); deployment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; deployment.collector_ids.clear(); @@ -3747,6 +3776,48 @@ mod tests { .values() .any(|node| matches!(node, crate::query_plan::QueryPlanNode::ExactFallback { .. }))); } + #[test] + fn immutable_two_sources_keep_actual_frontiers_and_bindings() { + let mut workload = request( + "nested", + "quantile(0.9, sum_over_time(m[1m]) + sum_over_time(n[1m]))", + ); + workload.hybrid_execution = true; + let states = + collect_selected_materializations(&workload.queries[0].post_asap, true).unwrap(); + assert_eq!(states.len(), 3, "source and consumer must both be selected"); + assert!(immutable_materialization_sources(&states[0].node).is_none()); + assert!(immutable_materialization_sources(&states[2].node).is_some()); + let mut deployment = environment(10_000); + deployment.target = PhysicalDeploymentTarget::BackendLocalRemoteWrite; + deployment.collector_ids.clear(); + let plan = PhysicalCompiler.compile(workload, deployment).unwrap(); + assert_eq!(plan.precompute_plan.materializations.len(), 3); + let derived = plan + .precompute_plan + .materializations + .iter() + .find(|m| m.derived_input.is_some()) + .unwrap(); + let sources = plan + .precompute_plan + .materializations + .iter() + .filter(|m| m.derived_input.is_none()) + .map(|m| m.policy_fingerprint().into()) + .collect::>(); + assert_eq!(derived.derived_input.as_ref().unwrap().inputs, sources); + assert_eq!(sources.len(), 2); + let entry = plan.query_plan.entries.values().next().unwrap(); + assert!(entry.nodes.values().any(|node| matches!( + node, + crate::query_plan::QueryPlanNode::ReadMaterialization { .. } + ))); + assert!(!entry + .nodes + .values() + .any(|node| matches!(node, crate::query_plan::QueryPlanNode::ExactFallback { .. }))); + } /// Distinct range queries retain a per-series HLL selected by Planner. #[test] diff --git a/crates/asap_types/src/precompute_plan.rs b/crates/asap_types/src/precompute_plan.rs index 818195a3..ef4621ba 100644 --- a/crates/asap_types/src/precompute_plan.rs +++ b/crates/asap_types/src/precompute_plan.rs @@ -420,21 +420,30 @@ impl PrecomputePlan { ) }; if self.ingest.protocol != IngestProtocol::PrometheusRemoteWriteV1 - || derived.inputs.len() != 1 + || derived.inputs.is_empty() { return Err(invalid()); } - let source_id = *derived.inputs.first().unwrap(); - let source = self - .materializations + let sources = derived + .inputs .iter() - .find(|candidate| candidate.policy_fingerprint() == source_id.fingerprint()) - .ok_or_else(invalid)?; - validated_source_window_cohort(config, &[source])?; - // Current installed runtime capability remains nonoverlapping. - // The shared cohort contract also describes explicit full-window - // sliding for consumers which separately prove its completion. - if source.window_size != source.slide_interval { + .map(|id| { + self.materializations + .iter() + .find(|candidate| candidate.policy_fingerprint() == id.fingerprint()) + .ok_or_else(invalid) + }) + .collect::, _>>()?; + validated_source_window_cohort(config, &sources)?; + if sources.iter().any(|source| { + !matches!( + source.aggregation_type, + crate::AggregationType::Sum + ) + }) { + return Err(invalid()); + } + if config.window_size != config.slide_interval { return Err(invalid()); } let mut matched = false; @@ -458,10 +467,19 @@ impl PrecomputePlan { let [edge] = inputs.as_slice() else { return Err(invalid()); }; - let frontiers = installed.binding.nodes.iter().filter_map(|(node,binding)| { - matches!(binding, crate::executable_plan::BackendNodeBinding::Materialization { summary_definition } - if *summary_definition == source_id).then_some((*node,source_id)) - }).collect(); + let frontiers = installed + .binding + .nodes + .iter() + .filter_map(|(node, binding)| match binding { + crate::executable_plan::BackendNodeBinding::Materialization { + summary_definition, + } if derived.inputs.contains(summary_definition) => { + Some((*node, *summary_definition)) + } + _ => None, + }) + .collect(); let actual = crate::derived_input::DerivedInputIdentity::from_dag( &installed.document, edge.producer, @@ -471,20 +489,64 @@ impl PrecomputePlan { if &actual != derived { return Err(invalid()); } - let input = dag - .nodes - .iter() - .find(|node| node.id == edge.producer) - .ok_or_else(invalid)?; - if !matches!( - input.payload, - planner_types::post_asap::ExecutableOperatorPayload::Value { - operation: - planner_types::post_asap::ValueOperation::FinalizeExactAccumulator, - timing: planner_types::post_asap::ExecutionTiming::MaintenanceTime, + let mut pending = vec![edge.producer]; + let mut visited = BTreeSet::new(); + while let Some(id) = pending.pop() { + if !visited.insert(id) { + continue; + } + let node = dag + .nodes + .iter() + .find(|node| node.id == id) + .ok_or_else(invalid)?; + let children: Vec<_> = dag + .edges + .iter() + .filter(|edge| edge.consumer == id) + .collect(); + use planner_types::post_asap::{ + ExecutableOperatorPayload as Payload, ExecutionTiming, ValueOperation, + }; + if node.output_state + != planner_types::post_asap::ExecutionDataState::MAINTENANCE_ROWS + { + return Err(invalid()); + } + match &node.payload { + Payload::Value { + operation: ValueOperation::FinalizeExactAccumulator, + timing: ExecutionTiming::MaintenanceTime, + } if children.len() == 1 + && frontiers.contains_key(&children[0].producer) => {} + Payload::Binary { + operator, + timing: ExecutionTiming::MaintenanceTime, + } if children.len() == 2 + && children + .iter() + .filter(|edge| { + edge.role == planner_types::post_asap::EdgeRole::Left + }) + .count() + == 1 + && children + .iter() + .filter(|edge| { + edge.role == planner_types::post_asap::EdgeRole::Right + }) + .count() + == 1 + && operator.vector_match.is_none() + && matches!( + operator.kind, + planner_types::pre_asap::BinaryOpKind::Arithmetic(_) + ) => + { + pending.extend(children.iter().map(|edge| edge.producer)); + } + _ => return Err(invalid()), } - ) { - return Err(invalid()); } matched = true; } diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 73b54134..914ec4c3 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -4295,8 +4295,7 @@ aggregations: async fn http_routes_asap_tier_metric_to_simple_engine() { // Default (no hot-reload) → `SketchStore`. The handler // takes the direct `ASAPQueryEngine::handle_query` path; the - // response's `infos` array carries `data_source: asap_query` - // so callers can byte-compare which engine answered. + // empty local engine reports no result without a warm annotation. let server_port = setup_test_server_with_router(StorageBackend::SketchStore, Vec::new()).await; let client = Client::new(); @@ -4312,7 +4311,12 @@ aggregations: resp.status() ); let body: serde_json::Value = resp.json().await.unwrap(); - assert_data_source(&body, "asap_query"); + // Empty local fixtures exercise routing, not successful execution. + assert_eq!( + body, + serde_json::json!({"status":"error", "data":null, + "errorType":"bad_data", "error":"No result for query"}) + ); } #[tokio::test] @@ -4375,7 +4379,12 @@ aggregations: resp.status() ); let body: serde_json::Value = resp.json().await.unwrap(); - assert_data_source(&body, "asap_query"); + // Empty local fixtures exercise routing, not successful execution. + assert_eq!( + body, + serde_json::json!({"status":"error", "data":null, + "errorType":"bad_data", "error":"No result for query"}) + ); } #[tokio::test] @@ -4614,7 +4623,12 @@ aggregations: resp.status(), ); let body: serde_json::Value = resp.json().await.unwrap(); - assert_data_source(&body, "asap_query"); + // Empty local fixtures exercise routing, not successful execution. + assert_eq!( + body, + serde_json::json!({"status":"error", "data":null, + "errorType":"bad_data", "error":"No result for query"}) + ); } #[tokio::test] @@ -4762,7 +4776,12 @@ aggregations: resp.status() ); let body: serde_json::Value = resp.json().await.unwrap(); - assert_data_source(&body, "asap_query"); + // Empty local fixtures exercise routing, not successful execution. + assert_eq!( + body, + serde_json::json!({"status":"error", "data":null, + "errorType":"bad_data", "error":"No result for query"}) + ); assert_eq!( gorilla_calls.load(Ordering::SeqCst), 0, @@ -4941,7 +4960,12 @@ aggregations: .expect("Failed to send request"); assert!(resp.status().is_success(), "default routing must still 2xx"); let body: serde_json::Value = resp.json().await.unwrap(); - assert_data_source(&body, "asap_query"); + // Empty local fixtures exercise routing, not successful execution. + assert_eq!( + body, + serde_json::json!({"status":"error", "data":null, + "errorType":"bad_data", "error":"No result for query"}) + ); assert_eq!( gorilla_calls.load(Ordering::SeqCst), 0, diff --git a/data_plane/tests/support/immutable_maintenance_process.rs b/data_plane/tests/support/immutable_maintenance_process.rs index f7357bc9..e329e9a2 100644 --- a/data_plane/tests/support/immutable_maintenance_process.rs +++ b/data_plane/tests/support/immutable_maintenance_process.rs @@ -3,13 +3,27 @@ use super::*; #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn single_source_maintenance_is_automatic_and_durable() { - const QUERY: &str = "quantile(0.9, sum_over_time(immutable_value[1m]))"; + run_maintenance_process(false).await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn two_source_maintenance_is_automatic_and_durable() { + run_maintenance_process(true).await; +} + +async fn run_maintenance_process(multi_source: bool) { + let query = if multi_source { + "quantile(0.9, sum_over_time(immutable_value[1m]) + sum_over_time(immutable_other[1m]))" + } else { + "quantile(0.9, sum_over_time(immutable_value[1m]))" + }; + let expected = if multi_source { 20.0 } else { 10.0 }; let mut fixture: Value = serde_json::from_str(include_str!( "../../../docs/examples/asapquery-compatibility-demo-snapshot.json" )) .unwrap(); let mut entry = fixture["query_workload"]["repeating_queries"][3].clone(); - entry["query"] = QUERY.into(); + entry["query"] = query.into(); entry["demand"]["fixed_interval_at"]["interval"] = 60_000.into(); entry["demand"]["fixed_interval_at"]["evaluation_phase"] = 0.into(); entry["time_selection"]["lookback"] = 60_000.into(); @@ -17,7 +31,10 @@ async fn single_source_maintenance_is_automatic_and_durable() { let snapshot: control_plane::physical::compiler::BackendLocalPlanningSnapshot = serde_json::from_value(fixture.clone()).unwrap(); let plan = snapshot.compile().unwrap(); - assert_eq!(plan.precompute_plan.materializations.len(), 2); + assert_eq!( + plan.precompute_plan.materializations.len(), + if multi_source { 3 } else { 2 } + ); let source = plan .precompute_plan .materializations @@ -36,7 +53,12 @@ async fn single_source_maintenance_is_automatic_and_durable() { assert_eq!(derived.pane_origin_ms, Some(0)); assert_eq!( derived.derived_input.as_ref().unwrap().inputs, - std::collections::BTreeSet::from([source.policy_fingerprint().into()]) + plan.precompute_plan + .materializations + .iter() + .filter(|m| m.derived_input.is_none()) + .map(|m| m.policy_fingerprint().into()) + .collect() ); // The production cost model may choose DDSketch or KLL. Preserve that // choice and use its actual value contract for this singleton oracle. @@ -63,7 +85,10 @@ async fn single_source_maintenance_is_automatic_and_durable() { }; // Two independent deployments: singleton is supported; a second physical // input series must never be mistaken for a complete singleton population. - for count in [1, 2] { + for (count, missing_source) in [(1, false), (2, false), (1, true)] { + if missing_source && !multi_source { + continue; + } let mut directory = tempfile::tempdir().unwrap(); eprintln!("IMMUTABLE_PROCESS_ARTIFACT {}", directory.path().display()); directory.disable_cleanup(true); @@ -113,7 +138,7 @@ async fn single_source_maintenance_is_automatic_and_durable() { let backend = format!("http://127.0.0.1:{port}"); let mut first = spawn(port); wait_until_ready(&client, &format!("{backend}/api/v1/health"), &mut first.0).await; - let series = (0..count) + let mut series: Vec<_> = (0..count) .map(|i| { series_with_labels( "immutable_value", @@ -125,6 +150,18 @@ async fn single_source_maintenance_is_automatic_and_durable() { ) }) .collect(); + if multi_source && !missing_source { + for i in 0..count { + series.push(series_with_labels( + "immutable_other", + &[ + ("instance", if i == 0 { "a" } else { "b" }), + ("job", "worker"), + ], + &[(1_000, 2.0), (2_000, 3.0), (60_000, 5.0)], + )); + } + } assert_eq!( remote_write(&client, &backend, &WriteRequest { timeseries: series }).await, 204 @@ -134,7 +171,7 @@ async fn single_source_maintenance_is_automatic_and_durable() { .send() .await .unwrap(); - if count == 1 { + if count == 1 && !missing_source { assert!( drain.status().is_success(), "{}", @@ -143,14 +180,14 @@ async fn single_source_maintenance_is_automatic_and_durable() { } let response: Value = client .get(format!("{backend}/api/v1/query")) - .query(&[("query", QUERY), ("time", "60")]) + .query(&[("query", query), ("time", "60")]) .send() .await .unwrap() .json() .await .unwrap(); - if count == 2 { + if count == 2 || missing_source { assert!( !is_warm(&response), "multi-series population was incorrectly admitted: {response}" @@ -173,7 +210,7 @@ async fn single_source_maintenance_is_automatic_and_durable() { .parse::() .unwrap(); assert!( - estimate.is_finite() && (estimate - 10.0).abs() / 10.0 <= max_relative_error, + estimate.is_finite() && (estimate - expected).abs() / expected <= max_relative_error, "selected singleton quantile exceeded its value contract: {response}" ); drop(first); @@ -188,7 +225,7 @@ async fn single_source_maintenance_is_automatic_and_durable() { .await; let after: Value = client .get(format!("{backend}/api/v1/query")) - .query(&[("query", QUERY), ("time", "60")]) + .query(&[("query", query), ("time", "60")]) .send() .await .unwrap() @@ -245,7 +282,7 @@ async fn single_source_maintenance_is_automatic_and_durable() { .await; let stale: Value = client .get(format!("{backend}/api/v1/query")) - .query(&[("query", QUERY), ("time", "60")]) + .query(&[("query", query), ("time", "60")]) .send() .await .unwrap()