diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 14f1efd58..573bf6764 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -4758,22 +4758,21 @@ fn collect_selected_materializations_with( } if let Some(readout) = readout { let mut parameters = sketch_params_json(kind.params()); - let mut item_label = None; + use planner_types::post_asap::SummaryInputExpr; + let item_label = match &input.item { + Some(SummaryInputExpr::Column( + planner_types::pre_asap::ColumnRef::Named(label), + )) => Some(label.clone()), + Some(SummaryInputExpr::Column( + planner_types::pre_asap::ColumnRef::Qualified { name, .. }, + )) => Some(name.clone()), + // Legacy/direct TopK plans did not carry an item + // projection. Keep their generic item-key behavior; + // typed Planner plans name the inner aggregate + // identity explicitly through SummaryUpdate.item. + _ => None, + }; if matches!(readout, SketchQuery::TopK { .. }) { - use planner_types::post_asap::SummaryInputExpr; - item_label = match &input.item { - Some(SummaryInputExpr::Column( - planner_types::pre_asap::ColumnRef::Named(label), - )) => Some(label.clone()), - Some(SummaryInputExpr::Column( - planner_types::pre_asap::ColumnRef::Qualified { name, .. }, - )) => Some(name.clone()), - // Legacy/direct TopK plans did not carry an item - // projection. Keep their generic item-key behavior; - // typed Planner plans name the inner aggregate - // identity explicitly through SummaryUpdate.item. - _ => None, - }; let mode = match &input.weight { SummaryInputExpr::Constant(value) if *value == 1.0 => "count", SummaryInputExpr::Column( @@ -9174,6 +9173,44 @@ pub(crate) mod tests { } } + // Both plan transports retain nonempty typed DAGs and numeric binding keys. + #[test] + fn installed_dag_plan_round_trips_in_json_and_yaml() { + let plan = DeploymentPlanCompiler + .compile_promql( + request("roundtrip", "sum_over_time(m[1m])"), + environment(10_000), + ) + .unwrap() + .precompute_plan; + assert!(!plan.executable_dags.is_empty()); + let json = serde_json::to_value(&plan).unwrap(); + let from_json: PrecomputePlan = serde_json::from_value(json.clone()).unwrap(); + let from_yaml: PrecomputePlan = + serde_yaml::from_str(&serde_yaml::to_string(&plan).unwrap()).unwrap(); + assert_eq!(serde_json::to_value(from_json).unwrap(), json); + assert_eq!(serde_json::to_value(from_yaml).unwrap(), json); + } + + // Non-TopK keyed sketches retain the same item dimension as their Planner update. + #[test] + fn keyed_non_topk_materialization_preserves_item_label() { + let request = request("keyed", "quantile_over_time(0.90, m[1m])"); + let mut root = request.queries[0].selected_plan_root.as_ref().clone(); + let SummaryExpr::SummaryEstimate { summary_input, .. } = &mut root.expr else { + panic!("expected estimate"); + }; + let SummaryExpr::SummaryAgg { input, .. } = &mut Rc::make_mut(summary_input).expr else { + panic!("expected summary state"); + }; + input.item = Some(planner_types::post_asap::SummaryInputExpr::Column( + planner_types::pre_asap::ColumnRef::Named("customer".into()), + )); + let root = Rc::new(root); + let selected = collect_selected_materializations(&root, false).unwrap(); + assert_eq!(selected[0].item_label.as_deref(), Some("customer")); + } + #[test] fn multiple_readouts_share_one_precompute_materialization() { let mut compilation_request = request("q-p90", "quantile_over_time(0.90, m[1m])"); diff --git a/crates/asap_types/src/precompute_plan.rs b/crates/asap_types/src/precompute_plan.rs index f9bc7fa6f..793cdbfff 100644 --- a/crates/asap_types/src/precompute_plan.rs +++ b/crates/asap_types/src/precompute_plan.rs @@ -121,6 +121,130 @@ pub struct PrecomputePlan { pub executable_dags: BTreeMap, } +/// Temporary decoded indexes for a complete plan validation/installation pass. +/// Keeping this separate from the wire plan preserves its Send/Sync contract. +pub struct PrecomputePlanLookup<'a> { + families: HashMap, + dags: Vec<( + &'a crate::executable_plan::InstalledPostAsapDag, + planner_types::post_asap::PostAsapDag, + )>, +} + +impl PrecomputePlanLookup<'_> { + pub fn state_family(&self, output: crate::sds::StoredOutputId) -> Option<&SummaryFamilyType> { + self.families.get(&output).copied() + } + + #[allow(clippy::type_complexity)] + pub fn summary_producer( + &self, + output: crate::sds::StoredOutputId, + ) -> Result< + Option<( + planner_types::post_asap::PostAsapDagNode, + Option, + )>, + String, + > { + use planner_types::post_asap::PostAsapOperatorPayload; + let mut selected: Option<( + planner_types::post_asap::PostAsapDagNode, + Option, + )> = None; + for (installed, dag) in &self.dags { + let Some(node) = installed.binding.nodes.iter().find_map(|(node, binding)| { + matches!(binding, crate::executable_plan::BackendNodeBinding::Materialization { stored_output } if *stored_output == output).then_some(*node) + }) else { + continue; + }; + let producer = dag + .nodes + .iter() + .find(|candidate| candidate.id == node) + .ok_or("bound materialization node is absent from its DAG")?; + if !matches!(producer.payload, PostAsapOperatorPayload::SummaryAgg { .. }) { + continue; + } + let inputs: Vec<_> = dag.edges.iter().filter(|e| e.consumer == node).collect(); + let source = match inputs.as_slice() { + [edge] => dag + .nodes + .iter() + .find(|candidate| candidate.id == edge.producer) + .and_then(|source| match &source.payload { + PostAsapOperatorPayload::Fallback { expression } => { + Some(expression.clone()) + } + _ => None, + }), + _ => None, + }; + if let Some((first, first_source)) = &selected { + // Shared outputs may be read by scans that project different + // columns; what must agree is the update and its population. + let population = + |node: &planner_types::post_asap::PostAsapDagNode, + source: &Option| { + let PostAsapOperatorPayload::SummaryAgg { family, .. } = &node.payload + else { + return None; + }; + source.as_ref().map(|expression| { + raw_time_series_input_contract( + expression, + matches!(family, SummaryFamilyType::ExactAggregate(..)), + ) + .map(|(metric, _, filter)| { + (metric, crate::utils::normalize_spatial_filter(&filter)) + }) + }) + }; + if first.payload != producer.payload + || population(first, first_source) != population(producer, &source) + { + return Err(format!( + "stored output {} is produced by DAGs that disagree on its computation", + output.as_u64() + )); + } + } else { + selected = Some((producer.clone(), source)); + } + } + Ok(selected) + } + + /// Canonical population predicate of `config`'s input: the typed table + /// population, or the label filter of its Planner DAG time-series scan. + pub fn population_filter( + &self, + config: &crate::PrecomputeMaterialization, + ) -> Result { + let table = config.table_population_canonical()?; + if config.table_name.is_some() || config.derived_input.is_some() { + return Ok(table); + } + let expression = match self.summary_producer(config.stored_output_id)? { + Some((_, Some(expression))) => expression, + // A bound raw output whose input cannot be read must not be + // treated as unfiltered. + Some((_, None)) => { + return Err("raw output's Planner producer does not read a source scan".into()) + } + // Only plans without a Planner DAG for this output (imported or + // fixture outputs) have no predicate to read. + None => return Ok(String::new()), + }; + let family = self.state_family(config.stored_output_id); + let (_, _, filter) = raw_time_series_input_contract( + &expression, + matches!(family, Some(SummaryFamilyType::ExactAggregate(..))), + )?; + Ok(crate::utils::normalize_spatial_filter(&filter)) + } +} + /// Serde mirror of [`PrecomputePlan`]; the compiler checks it stays in sync. #[derive(Deserialize)] #[serde(remote = "PrecomputePlan")] @@ -141,11 +265,11 @@ impl<'de> Deserialize<'de> for PrecomputePlan { /// version rather than with its first unrecognized field. fn deserialize>(deserializer: D) -> Result { use serde::de::Error; - let value = serde_json::Value::deserialize(deserializer)?; + let value = serde_yaml::Value::deserialize(deserializer)?; let compat = value .get("envelope") .and_then(|envelope| envelope.get("backend_compat")) - .and_then(serde_json::Value::as_str); + .and_then(serde_yaml::Value::as_str); if compat != Some(BACKEND_COMPAT) { return Err(D::Error::custom(format!( "unsupported installed precompute plan schema {}: this backend requires \ @@ -153,7 +277,15 @@ impl<'de> Deserialize<'de> for PrecomputePlan { compat.unwrap_or("") ))); } - PrecomputePlanDef::deserialize(value).map_err(D::Error::custom) + // Preserve YAML enum tags and numeric mapping keys. JSON uses string + // keys for numeric node IDs, so its buffered form needs JSON's key decoder. + PrecomputePlanDef::deserialize(value.clone()) + .map_err(|error| error.to_string()) + .or_else(|_| { + let json = serde_json::to_value(value).map_err(|error| error.to_string())?; + PrecomputePlanDef::deserialize(json).map_err(|error| error.to_string()) + }) + .map_err(D::Error::custom) } } @@ -471,11 +603,24 @@ impl PrecomputePlan { .map(|schema| &schema.family) } - /// The Planner SummaryAgg node that produces `output`, with the source - /// expression feeding it when that input is a raw source. `None` when no - /// installed DAG produces the output with a SummaryAgg. Every DAG that - /// produces the output must agree on both, or the plan is rejected: the - /// runtime routes, filters and catalogs the output once. + /// Decode each Planner DAG once for a validation or installation pass. + /// The lookup stays local because Planner payloads contain process-local `Rc`s. + pub fn lookup(&self) -> Result, String> { + Ok(PrecomputePlanLookup { + families: self + .schemas + .iter() + .map(|schema| (schema.materialization, &schema.family)) + .collect(), + dags: self + .executable_dags + .values() + .map(|installed| installed.document.decode().map(|dag| (installed, dag))) + .collect::>()?, + }) + } + + /// The authoritative producer and raw input of one stored output. #[allow(clippy::type_complexity)] pub fn summary_producer( &self, @@ -487,102 +632,15 @@ impl PrecomputePlan { )>, String, > { - use planner_types::post_asap::PostAsapOperatorPayload; - let mut selected: Option<( - planner_types::post_asap::PostAsapDagNode, - Option, - )> = None; - for installed in self.executable_dags.values() { - let Some(node) = installed.binding.nodes.iter().find_map(|(node, binding)| { - matches!(binding, crate::executable_plan::BackendNodeBinding::Materialization { stored_output } if *stored_output == output).then_some(*node) - }) else { - continue; - }; - let dag = installed.document.decode()?; - let producer = dag - .nodes - .iter() - .find(|candidate| candidate.id == node) - .ok_or("bound materialization node is absent from its DAG")?; - if !matches!(producer.payload, PostAsapOperatorPayload::SummaryAgg { .. }) { - continue; - } - let inputs: Vec<_> = dag.edges.iter().filter(|e| e.consumer == node).collect(); - let source = match inputs.as_slice() { - [edge] => dag - .nodes - .iter() - .find(|candidate| candidate.id == edge.producer) - .and_then(|source| match &source.payload { - PostAsapOperatorPayload::Fallback { expression } => { - Some(expression.clone()) - } - _ => None, - }), - _ => None, - }; - if let Some((first, first_source)) = &selected { - // Shared outputs may be read by scans that project different - // columns; what must agree is the update and its population. - let population = - |node: &planner_types::post_asap::PostAsapDagNode, - source: &Option| { - let PostAsapOperatorPayload::SummaryAgg { family, .. } = &node.payload - else { - return None; - }; - source.as_ref().map(|expression| { - raw_time_series_input_contract( - expression, - matches!(family, SummaryFamilyType::ExactAggregate(..)), - ) - .map(|(metric, _, filter)| { - (metric, crate::utils::normalize_spatial_filter(&filter)) - }) - }) - }; - if first.payload != producer.payload - || population(first, first_source) != population(producer, &source) - { - return Err(format!( - "stored output {} is produced by DAGs that disagree on its computation", - output.as_u64() - )); - } - } else { - selected = Some((producer.clone(), source)); - } - } - Ok(selected) + self.lookup()?.summary_producer(output) } - /// Canonical population predicate of `config`'s input: the typed table - /// population, or the label filter of its Planner DAG time-series scan. + /// Canonical input predicate of one output. pub fn population_filter( &self, config: &crate::PrecomputeMaterialization, ) -> Result { - let table = config.table_population_canonical()?; - if config.table_name.is_some() || config.derived_input.is_some() { - return Ok(table); - } - let expression = match self.summary_producer(config.stored_output_id)? { - Some((_, Some(expression))) => expression, - // A bound raw output whose input cannot be read must not be - // treated as unfiltered. - Some((_, None)) => { - return Err("raw output's Planner producer does not read a source scan".into()) - } - // Only plans without a Planner DAG for this output (imported or - // fixture outputs) have no predicate to read. - None => return Ok(String::new()), - }; - let family = self.state_family(config.stored_output_id); - let (_, _, filter) = raw_time_series_input_contract( - &expression, - matches!(family, Some(SummaryFamilyType::ExactAggregate(..))), - )?; - Ok(crate::utils::normalize_spatial_filter(&filter)) + self.lookup()?.population_filter(config) } pub fn runtime_materializations( @@ -675,6 +733,9 @@ impl PrecomputePlan { if !valid_ingest { return Err(PrecomputePlanError::UnsupportedIngestEndpoint); } + let lookup = self + .lookup() + .map_err(PrecomputePlanError::CatalogContract)?; let canonical_cohort_members: BTreeSet<_> = self .materializations .iter() @@ -685,12 +746,6 @@ impl PrecomputePlan { .chain(input.inputs.iter().copied()) }) .collect(); - // Decode each DAG at most once; errors still surface only where used. - let dags = self - .executable_dags - .values() - .map(|installed| (installed, std::cell::OnceCell::new())) - .collect::>(); for config in &self.materializations { if !config.population_key_encoding.is_legacy() && (self.ingest.protocol != IngestProtocol::PrometheusRemoteWriteV1 @@ -741,11 +796,7 @@ impl PrecomputePlan { } validated_source_window_cohort(config, &sources)?; let mut matched = false; - for (installed, dag) in &dags { - let dag = dag - .get_or_init(|| installed.document.decode()) - .as_ref() - .map_err(|error| PrecomputePlanError::CatalogContract(error.clone()))?; + for (installed, dag) in &lookup.dags { for sink in &installed.binding.precompute_sinks { if !matches!(installed.binding.node(*sink), Some(crate::executable_plan::BackendNodeBinding::Materialization { stored_output }) @@ -794,7 +845,7 @@ impl PrecomputePlan { .map(|contract| !asap_physical_operators::physical_planner::precompute::is_population_schema(&contract.schema))? .then_some(native); let exact = |config: &crate::PrecomputeMaterialization, kinds: &[ExactKind]| { - matches!(self.state_family(config.stored_output_id), + matches!(lookup.state_family(config.stored_output_id), Some(SummaryFamilyType::ExactAggregate(kind, _)) if kinds.contains(kind)) }; let source_kinds: &[ExactKind] = if native.is_some() { @@ -812,7 +863,7 @@ impl PrecomputePlan { validate_maintenance_reduction(config, target_node) .map_err(PrecomputePlanError::CatalogContract)?; } else if !exact(config, &[ExactKind::Sum]) - && !matches!(self.state_family(config.stored_output_id), + && !matches!(lookup.state_family(config.stored_output_id), Some(SummaryFamilyType::Sketch(kind, _)) if matches!(kind.algorithm(), SketchAlgorithm::CmsWithHeap | SketchAlgorithm::CountSketchWithHeap)) { @@ -911,7 +962,7 @@ impl PrecomputePlan { return Err(invalid()); } } - for (query_id, installed) in &self.executable_dags { + for ((query_id, _), (installed, dag)) in self.executable_dags.iter().zip(&lookup.dags) { if query_id != &installed.document.query_id { return Err(PrecomputePlanError::CatalogContract( "post-ASAP DAG map key differs from document query ID".into(), @@ -920,10 +971,6 @@ impl PrecomputePlan { installed .validate() .map_err(PrecomputePlanError::CatalogContract)?; - let dag = installed - .document - .decode() - .map_err(PrecomputePlanError::CatalogContract)?; for node in &dag.nodes { let Some(crate::executable_plan::BackendNodeBinding::Materialization { stored_output, @@ -946,7 +993,7 @@ impl PrecomputePlan { .. } = &node.payload { - if self.state_family(config.stored_output_id) != Some(family) { + if lookup.state_family(config.stored_output_id) != Some(family) { return Err(PrecomputePlanError::CatalogContract( "stored output schema family differs from its Planner producer".into(), )); @@ -1095,7 +1142,7 @@ impl PrecomputePlan { // read every series of its metric. if raw && !self.executable_dags.is_empty() - && self + && lookup .summary_producer(schema.materialization) .map_err(invalid_computation)? .is_none() @@ -1104,7 +1151,8 @@ impl PrecomputePlan { "raw output has no Planner producer".into(), )); } - self.population_filter(materialization) + lookup + .population_filter(materialization) .map_err(invalid_computation)?; let source = materialization.table_name.as_ref().map_or_else( || Source::TimeSeries { @@ -1544,6 +1592,21 @@ mod source_window_cohort_tests { } } + // YAML transport preserves tagged Planner families and numeric DAG node keys. + #[test] + fn installed_plan_yaml_round_trip_preserves_typed_fields() { + let mut config = full_window(); + config.slide_interval = config.window_size; + config.window_type = crate::WindowKind::Tumbling; + let plan = PrecomputePlan::build_backend_local(envelope(1), vec![(config, sum())]).unwrap(); + let yaml = serde_yaml::to_string(&plan).unwrap(); + let decoded: PrecomputePlan = serde_yaml::from_str(&yaml).unwrap(); + assert_eq!( + serde_json::to_value(decoded).unwrap(), + serde_json::to_value(plan).unwrap() + ); + } + // A plan compiled for the v1 schema, which duplicated computation fields // on materializations, is rejected by version before its body is read. #[test] diff --git a/crates/asap_types/src/precompute_plan/catalog.rs b/crates/asap_types/src/precompute_plan/catalog.rs index eaf49df13..a463b74e9 100644 --- a/crates/asap_types/src/precompute_plan/catalog.rs +++ b/crates/asap_types/src/precompute_plan/catalog.rs @@ -50,12 +50,7 @@ impl PrecomputePlan { "stored output semantics differ from installed writer computation", )); } - // Decode each DAG at most once; errors still surface only where used. - let dags = self - .executable_dags - .values() - .map(|installed| (installed, std::cell::OnceCell::new())) - .collect::>(); + let lookup = self.lookup().map_err(invalid)?; for config in &self.materializations { if config.derived_input.is_some() && config.semantic_fragment.is_none() { return Err(invalid( @@ -64,11 +59,7 @@ impl PrecomputePlan { } if let Some(expected) = &config.semantic_fragment { let mut found = false; - for (installed, dag) in &dags { - let dag = dag - .get_or_init(|| installed.document.decode()) - .as_ref() - .map_err(|error| invalid(error.clone()))?; + for (installed, dag) in &lookup.dags { for (id, binding) in &installed.binding.nodes { if matches!(binding, crate::executable_plan::BackendNodeBinding::Materialization { stored_output } if stored_output.fingerprint() == config.policy_fingerprint()) @@ -113,7 +104,7 @@ impl PrecomputePlan { for config in &self.materializations { let id = StoredOutputId::from(config.policy_fingerprint()); let binding = &catalog.outputs[&id]; - let family = self + let family = lookup .state_family(id) .ok_or(PrecomputePlanError::SchemaSetMismatch)?; let expected = @@ -147,7 +138,7 @@ impl PrecomputePlan { || data.source != expected_source || &data.value_projection != expected_projection || data.population_filter_canonical - != self.population_filter(config).map_err(invalid)? + != lookup.population_filter(config).map_err(invalid)? || data.group_by_keys != config.grouping_labels { return Err(invalid("source/population/grouping differs from catalog")); diff --git a/crates/asap_types/src/summary_catalog.rs b/crates/asap_types/src/summary_catalog.rs index c83ef3ded..ed5f50ddd 100644 --- a/crates/asap_types/src/summary_catalog.rs +++ b/crates/asap_types/src/summary_catalog.rs @@ -103,14 +103,15 @@ impl SummaryCatalog { pub fn from_plan( plan: &crate::precompute_plan::PrecomputePlan, ) -> Result { + let lookup = plan.lookup().map_err(SummaryCatalogError::Descriptor)?; let outputs = plan .materializations .iter() .map(|config| { - let family = plan.state_family(config.stored_output_id).ok_or( + let family = lookup.state_family(config.stored_output_id).ok_or( SummaryCatalogError::MissingDescriptor(config.stored_output_id.as_u64()), )?; - let filter = plan + let filter = lookup .population_filter(config) .map_err(SummaryCatalogError::Descriptor)?; Ok((config, family, filter)) diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 14fe449b4..6228d5769 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -1893,7 +1893,7 @@ impl PrometheusRemoteWriteReceiver { if let Some(previous) = committed.get(&output.0) { let config = plan.precompute_plan.materializations.iter().find(|c| c.policy_fp_u64() == output.0).ok_or("committed raw output missing")?; let mut groups: BTreeMap, BTreeMap<(u64,u64),Arc>> = BTreeMap::new(); - for record in previous { groups.entry(record.group.clone()).or_default().insert((record.start_ms,record.end_ms),crate::precompute_engine::revisions::decode_state(record, plan.precompute_plan.state_family(config.stored_output_id).ok_or("committed raw output has no state schema")?)?); } + for record in previous { groups.entry(record.group.clone()).or_default().insert((record.start_ms,record.end_ms),crate::precompute_engine::revisions::decode_state(record, plan.installed_precompute_plan.state_family(config.stored_output_id).ok_or("committed raw output has no state schema")?)?); } for (group,windows) in groups { frozen.push(crate::storage_engines::sketch_db::index::FrozenExactWindows { stored_output_reference: plan.installed_precompute_plan.stored_output_reference(output).ok_or("committed raw binding missing")?, storage_handle: output.0, definition: output, generation:Arc::new(generation.clone()), group, windows, singleton_population_complete:true, diff --git a/data_plane/src/precompute_engine/revisions.rs b/data_plane/src/precompute_engine/revisions.rs index 14ef43c4a..6403ec199 100644 --- a/data_plane/src/precompute_engine/revisions.rs +++ b/data_plane/src/precompute_engine/revisions.rs @@ -637,7 +637,7 @@ impl RevisionRuntime { { decode_state( record, - plan.precompute_plan + plan.installed_precompute_plan .state_family(config.stored_output_id) .ok_or("recovered revision has no state schema")?, )?; diff --git a/data_plane/src/query_engines/asap_query_engine/engine.rs b/data_plane/src/query_engines/asap_query_engine/engine.rs index 8b57629bc..8bb3bfec0 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -431,9 +431,7 @@ impl ASAPQueryEngine { asap_types::query_plan::QueryPlanNode::ExternalExact { .. } => false, asap_types::query_plan::QueryPlanNode::Logical { operator, .. } => !matches!( operator, - asap_types::query_plan::query_time::QueryTimeOperator::ExactSubquery { .. } - | asap_types::query_plan::query_time::QueryTimeOperator::CandidateExactSubquery { .. } - | asap_types::query_plan::query_time::QueryTimeOperator::Scan { .. } + asap_types::query_plan::query_time::QueryTimeOperator::Scan { .. } ), _ => true, }) { diff --git a/data_plane/src/query_engines/asap_query_engine/exact_subqueries.rs b/data_plane/src/query_engines/asap_query_engine/exact_subqueries.rs index 7ed9b1ae0..455a7ad51 100644 --- a/data_plane/src/query_engines/asap_query_engine/exact_subqueries.rs +++ b/data_plane/src/query_engines/asap_query_engine/exact_subqueries.rs @@ -16,17 +16,10 @@ fn miss(message: impl Into) -> EngineError { EngineError::capability_miss("exact_subquery", message.into()) } -/// Traverse only the installed graph. -#[derive(Debug, Clone)] -enum ExactLeaf { - Legacy(QueryTimeOperator), - External(ExternalExactRequest), -} - fn leaves( entry: &QueryPlanEntry, times: &[u64], -) -> Result, EngineError> { +) -> Result, EngineError> { let mut pending = Vec::new(); for at in times { pending.push(( @@ -54,7 +47,9 @@ fn leaves( } QueryTimeOperator::ExactSubquery { .. } | QueryTimeOperator::CandidateExactSubquery { .. } => { - result.insert((id, at), ExactLeaf::Legacy(operator.clone())); + return Err(miss( + "legacy exact operators are not executable; recompile the plan", + )); } _ => pending.extend(inputs.iter().map(|input| (*input, at))), }, @@ -82,7 +77,7 @@ fn leaves( QueryPlanNode::ExternalExact { request, inputs } => { pending.extend(inputs.iter().map(|input| (*input, at))); - result.insert((id, at), ExactLeaf::External(request.clone())); + result.insert((id, at), request.clone()); } // Existing bound summary subtrees are read by the synchronous callback. _ => {} @@ -96,24 +91,13 @@ pub(super) fn external_dependencies( times: &[u64], ) -> Result, EngineError> { let mut result = Vec::new(); - for ((id, at), leaf) in leaves(entry, times)? { - if matches!(leaf, ExactLeaf::External(_)) { - result.extend( - entry.nodes[&id] - .inputs() - .iter() - .map(|input| (id, *input, at)), - ); - } else if matches!( - leaf, - ExactLeaf::Legacy(QueryTimeOperator::CandidateExactSubquery { .. }) - ) { - let input = *entry.nodes[&id] + for ((id, at), _) in leaves(entry, times)? { + result.extend( + entry.nodes[&id] .inputs() - .first() - .ok_or_else(|| miss("candidate exact subtree has no membership input"))?; - result.push((id, input, at)); - } + .iter() + .map(|input| (id, *input, at)), + ); } Ok(result) } @@ -292,36 +276,23 @@ pub(super) async fn prepare_external( let mut remote_cache = HashMap::<(QueryLanguage, String, i64), Value>::new(); for ((id, at), leaf) in leaves(entry, times)? { u64::try_from(at).map_err(|_| miss("subquery predates epoch"))?; - let (language, query, candidate_input) = match &leaf { - ExactLeaf::Legacy(QueryTimeOperator::ExactSubquery { query }) => { - (QueryLanguage::PromQl, query.clone(), None) - } - ExactLeaf::Legacy(QueryTimeOperator::CandidateExactSubquery { query, item_label }) => ( - QueryLanguage::PromQl, - query.clone(), - Some((entry.nodes[&id].inputs()[0], item_label.as_str())), - ), - ExactLeaf::External(request) => { - if !matches!( - request.language, - QueryLanguage::PromQl | QueryLanguage::MetricsQl - ) { - return Err(miss(format!( - "external exact language {:?} has no installed adapter", - request.language - ))); - } - let candidate = match request.input_contracts.as_slice() { - [] => None, - [ExternalExactInput::CandidateMembership { item_label }] => { - Some((entry.nodes[&id].inputs()[0], item_label.as_str())) - } - _ => return Err(miss("unsupported external exact input contract")), - }; - (request.language, request.expression.clone(), candidate) + if !matches!( + leaf.language, + QueryLanguage::PromQl | QueryLanguage::MetricsQl + ) { + return Err(miss(format!( + "external exact language {:?} has no installed adapter", + leaf.language + ))); + } + let candidate_input = match leaf.input_contracts.as_slice() { + [] => None, + [ExternalExactInput::CandidateMembership { item_label }] => { + Some((entry.nodes[&id].inputs()[0], item_label.as_str())) } - _ => return Err(miss("prepared leaf is not an exact subtree")), + _ => return Err(miss("unsupported external exact input contract")), }; + let (language, query) = (leaf.language, leaf.expression.clone()); let (query, candidate_filtered) = if let Some((candidate_input, item_label)) = candidate_input { @@ -809,6 +780,21 @@ mod tests { )); server.abort(); } + fn external_exact_leaf(expression: &str) -> QueryPlanNode { + QueryPlanNode::ExternalExact { + request: ExternalExactRequest { + language: QueryLanguage::PromQl, + expression: expression.into(), + output: asap_types::query_plan::ExternalExactOutput::InstantVector, + parameters: BTreeMap::new(), + start_parameter: None, + end_parameter: None, + input_contracts: vec![], + }, + inputs: vec![], + } + } + /// Planner's label-map division, the computation over a prepared exact leaf. fn planner_division() -> Vec { use planner_types::{ @@ -863,13 +849,7 @@ mod tests { QueryNodeId(1), QueryPlanNode::SummaryMerge { inputs: vec![] }, ), - ( - QueryNodeId(2), - QueryPlanNode::Logical { - operator: QueryTimeOperator::ExactSubquery { query: "b".into() }, - inputs: vec![], - }, - ), + (QueryNodeId(2), external_exact_leaf("b")), ])); let leaves = prepare( &entry, @@ -905,13 +885,9 @@ mod tests { assert_eq!(stats.summary_readout_evaluations, 1); assert_eq!(calls.load(Ordering::SeqCst), 1); let mut repeated = entry.clone(); - repeated.nodes.insert( - QueryNodeId(1), - QueryPlanNode::Logical { - operator: QueryTimeOperator::ExactSubquery { query: "b".into() }, - inputs: vec![], - }, - ); + repeated + .nodes + .insert(QueryNodeId(1), external_exact_leaf("b")); let prepared = prepare( &repeated, &[1000], @@ -992,15 +968,7 @@ mod tests { pruning: None, }, ), - ( - QueryNodeId(1), - QueryPlanNode::Logical { - operator: QueryTimeOperator::ExactSubquery { - query: exact_query.into(), - }, - inputs: vec![], - }, - ), + (QueryNodeId(1), external_exact_leaf(exact_query)), ( QueryNodeId(2), QueryPlanNode::ExactReadout { @@ -1069,6 +1037,27 @@ mod tests { assert_eq!(calls.load(Ordering::SeqCst), 0); server.abort(); } + // Legacy exact operators never reach the runtime or trigger a PromQL reparse. + #[test] + fn legacy_exact_leaves_are_rejected_before_execution() { + for operator in [ + QueryTimeOperator::ExactSubquery { query: "m".into() }, + QueryTimeOperator::CandidateExactSubquery { + query: "m".into(), + item_label: "job".into(), + }, + ] { + let entry = entry(BTreeMap::from([( + QueryNodeId(0), + QueryPlanNode::Logical { + operator, + inputs: vec![], + }, + )])); + assert!(leaves(&entry, &[1000]).is_err()); + } + } + #[test] fn deployed_raw_scan_is_rejected_before_execution() { let entry = entry(BTreeMap::from([( diff --git a/data_plane/src/storage_engines/sketch_db/index/mod.rs b/data_plane/src/storage_engines/sketch_db/index/mod.rs index 19037238a..e06209f7e 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -6980,7 +6980,7 @@ impl SketchStore { output.catalog_generation = Some(Arc::clone(&generation)); output.stored_output_reference = Some(record.reference.clone()); let family = plan - .precompute_plan + .installed_precompute_plan .state_family(definition) .ok_or("revision view lacks output state schema")?; let kind = plan diff --git a/data_plane/src/storage_engines/types/installed_precompute_plan.rs b/data_plane/src/storage_engines/types/installed_precompute_plan.rs index d616a0307..d9235d0aa 100644 --- a/data_plane/src/storage_engines/types/installed_precompute_plan.rs +++ b/data_plane/src/storage_engines/types/installed_precompute_plan.rs @@ -49,10 +49,11 @@ impl InstalledPrecomputePlan { /// Production construction always validates the executable DAG and bindings. pub fn from_precompute_plan(plan: asap_types::precompute_plan::PrecomputePlan) -> Result { let materializations = plan.runtime_materializations()?; + let lookup = plan.lookup().map_err(anyhow::Error::msg)?; let outputs = materializations .into_values() .map(|config| { - let family = plan + let family = lookup .state_family(config.stored_output_id) .cloned() .ok_or_else(|| anyhow::anyhow!("stored output has no state schema"))?; @@ -62,7 +63,8 @@ impl InstalledPrecomputePlan { let population_filters = outputs .iter() .map(|(config, _)| { - plan.population_filter(config) + lookup + .population_filter(config) .map(|filter| (config.stored_output_id, filter)) .map_err(anyhow::Error::msg) }) @@ -71,7 +73,7 @@ impl InstalledPrecomputePlan { for (config, _) in &outputs { use planner_types::post_asap::{PostAsapOperatorPayload, SummaryInputExpr}; use planner_types::pre_asap::ColumnRef; - let Some((node, _)) = plan + let Some((node, _)) = lookup .summary_producer(config.stored_output_id) .map_err(anyhow::Error::msg)? else { @@ -107,6 +109,7 @@ impl InstalledPrecomputePlan { .map_err(anyhow::Error::msg)?; programs.insert(config.policy_fp_u64(), std::sync::Arc::new(program)); } + drop(lookup); let mut view = Self::derived_view(outputs); view.population_filters = population_filters; view.item_labels = item_labels; diff --git a/docs/design_docs/physical-operators.md b/docs/design_docs/physical-operators.md index 5d00727e2..8c9cd3155 100644 --- a/docs/design_docs/physical-operators.md +++ b/docs/design_docs/physical-operators.md @@ -25,7 +25,14 @@ Planner does not own storage formats. The backend owns them: The query and precompute integration PRs both use the independent ASAP DAG runtime. Deployment code binds input sources, storage, ingestion windows, publication and protocol outputs. Execution phase belongs to the physical node's data state, not the operator payload. -Computation has the same semantics during precomputation and query execution. +The Planner maintenance lifecycle is the only source of node timing. Computation +has the same semantics during precomputation and query execution. + +`PrecomputeMaterialization` contains deployment fields: output identity, source +binding, partitioning, windows, retention and runtime policy. Families come from +state-schema contracts; updates, item labels and input predicates come from the +bound Planner DAG. Removed computation fields are rejected by the v2 backend +compatibility contract. JSON and YAML transports preserve the same typed plan. Candidate pruning uses a general semi-join with explicit matching keys, followed by grouped Sort and grouped Limit. The completeness certificate belongs to the @@ -42,9 +49,10 @@ Planner tests cover shared producers, phase assignment, typed batches and composed candidate pruning. Backend tests cover source binding, installed plan validation, storage compatibility and query responses. -The migration does not supply a local raw Scan. That deployment capability -remains deferred. Library tests supplied with raw batches are not evidence that -the backend can execute arbitrary raw-only installed plans. +Ephemeral plans bind raw-series inputs to the configured Prometheus source at +query evaluation time. Mixed plans bind those raw inputs beside complete stored +revisions within the configured lag bound; unavailable inputs follow the installed +exact fallback. ## Stack integration @@ -56,9 +64,9 @@ structure, #742 checks selection with synthetic costs, and #775 checks installed plans against data-plane results. Production evidence and performance validation remain separate from these correctness tests. -The library also owns stored-summary decoding, delta reconstruction and -family-specific readout kernels. Deployment adapters select compatible panes -and translate inputs and outputs; they do not copy those computations. +Planner owns pane-state merge and family-specific readout kernels. Backend owns +stored-summary decoding, delta reconstruction and pane selection, then binds the +decoded state to those kernels. The shared `physical_planner::compile` API validates concrete operators and typed input contracts before opening sources. Deployment resolves those inputs and diff --git a/docs/design_docs/precompute-dag-execution.md b/docs/design_docs/precompute-dag-execution.md index 9bd7f6050..67782516f 100644 --- a/docs/design_docs/precompute-dag-execution.md +++ b/docs/design_docs/precompute-dag-execution.md @@ -97,7 +97,10 @@ version; population and window selection identify the records to read. The backend derives an immutable `InstalledPrecomputePlan` from the supplied PrecomputePlan. This runtime representation avoids validating the graph and -binding operators again for every input batch. +binding operators again for every input batch. A temporary lookup decodes each +Planner DAG once per validation or installation pass and indexes state families +by stored output. The installed router retains only the derived deployment +indexes; it does not deserialize a DAG for each incoming sample. Installation performs four steps: