diff --git a/control_plane/src/clickhouse.rs b/control_plane/src/clickhouse.rs index 5a3c92eb1..f5d822dbf 100644 --- a/control_plane/src/clickhouse.rs +++ b/control_plane/src/clickhouse.rs @@ -295,7 +295,7 @@ pub async fn compile_automatic_clickhouse_workload( installed_dags.insert( query.sql.clone(), installed - .maintenance_projection() + .precompute_projection() .map_err(ClickHousePlanningError::Lower)?, ); } @@ -450,7 +450,7 @@ pub async fn compile_clickhouse_workload( installed_dags.insert( query.sql.clone(), installed - .maintenance_projection() + .precompute_projection() .map_err(ClickHousePlanningError::Lower)?, ); let identity = diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 5e12548d8..f32a4f6d2 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -737,7 +737,7 @@ impl BackendLocalPlanningInput { self.physical_inputs.query_retention_margin_ms, )?; } - // Composable lowering residualizes unsafe leaves individually; retain Planner siblings. + // Composable lowering defers unsafe leaves to query time individually; retain Planner siblings. Ok(( PhysicalCompilationRequest { planner_selection_trace, @@ -915,7 +915,7 @@ fn preserve_invalid_exact_fallback_roots( /// A MetricsQL query whose only selected states are Prometheus-specific /// counter readouts has no backend materialization to bind. Keep the original /// query as one native exact root. Mixed queries retain their other selected -/// summaries and let residual lowering cut only the counter branches. +/// summaries and let query-time lowering cut only the counter branches. fn preserve_metricsql_counter_only_roots( queries: &mut [QueryCompilationInput], canonical_roots: &[Rc], @@ -1133,7 +1133,7 @@ impl DeploymentPlanCompiler { planner_types::post_asap::ExactKind::Max, _ ) - ) || crate::query_plan::residual::selected_range_max_materialization( + ) || crate::query_plan::query_time::selected_range_max_materialization( &query.query_string, &state.node, ) @@ -1141,14 +1141,14 @@ impl DeploymentPlanCompiler { .flatten() .is_some()) && request.enabled_materialization_keys.as_ref().is_none_or(|policy| { - let key = crate::query_plan::residual::selected_counter_materialization( + let key = crate::query_plan::query_time::selected_counter_materialization( &query.query_string, &state.node, ) .ok() .flatten() .or_else(|| { - crate::query_plan::residual::selected_range_max_materialization( + crate::query_plan::query_time::selected_range_max_materialization( &query.query_string, &state.node, ) @@ -1773,13 +1773,13 @@ impl DeploymentPlanCompiler { cumulative_readout: true, }; // A whole-query native fallback need not be expressible in the local - // residual algebra (for example an ERP-rejected entropy readout). + // query-time algebra (for example an ERP-rejected entropy readout). // Retain its native boundary without discarding other workload roots. let native_root = request.allow_mixed_summary_and_exact_execution && if let SummaryExpr::KeepPreAsap(expr) = &query.selected_plan_root.expr { let original = original_root(query, query_index, &request.canonical_roots)?; expr.as_ref() == &original - && crate::query_plan::residual::compile_logical( + && crate::query_plan::query_time::compile_logical( query.query_id.clone(), canonical.clone(), instant, @@ -1883,7 +1883,7 @@ impl DeploymentPlanCompiler { // Any Planner-selected leaf without a physical summary binding // is an exact subtree boundary. Deployed plans never retain a // backend-local range index leaf. - crate::query_plan::residual::finalize_residuals(&mut entry)?; + crate::query_plan::query_time::finalize_query_time_nodes(&mut entry)?; } if frontend == QueryFrontend::MetricsQl { entry.language = crate::query_plan::QueryLanguage::MetricsQl; @@ -1990,7 +1990,7 @@ impl DeploymentPlanCompiler { .filter(|(_, installed)| !installed.binding.precompute_sinks.is_empty()) .map(|(query_id, installed)| { installed - .maintenance_projection() + .precompute_projection() .map(|projected| (query_id.clone(), projected)) .map_err(|reason| CompileError::Query { query_id, reason }) }) @@ -3166,7 +3166,7 @@ fn select_lifecycle( } /// A warm producer may consume only a source whose semantics its precompute accumulator -/// implements. Predicates and shifted ranges remain executable residual nodes. +/// implements. Predicates and shifted ranges remain executable query-time nodes. pub(crate) fn raw_materialization_input_contract( node: &SummaryNode, ) -> Result<(String, Option, String), String> { @@ -3901,16 +3901,24 @@ pub(crate) mod tests { .ok() }) .collect(); - let plan = plans.iter().find(|plan| plan.query_plan.entries.values().all(|entry| - entry.nodes.values().any(|node| matches!(node, crate::query_plan::QueryPlanNode::Logical { - operator: crate::query_plan::residual::ResidualQueryOperator::CurrentSeries { .. }, .. - })))).expect("no shared current-series candidate"); + let plan = plans + .iter() + .find(|plan| { + plan.query_plan.entries.values().all(|entry| { + entry.nodes.values().any(|node| { + matches!(node, crate::query_plan::QueryPlanNode::Logical { + operator: crate::query_plan::query_time::QueryTimeOperator::CurrentSeries { .. }, .. + }) + }) + }) + }) + .expect("no shared current-series candidate"); let mut populations = BTreeSet::new(); for entry in plan.query_plan.entries.values() { for node in entry.nodes.values() { if let crate::query_plan::QueryPlanNode::Logical { operator: - crate::query_plan::residual::ResidualQueryOperator::CurrentSeries { + crate::query_plan::query_time::QueryTimeOperator::CurrentSeries { population, .. }, @@ -4023,8 +4031,9 @@ pub(crate) mod tests { .any(|node| matches!( node, crate::query_plan::QueryPlanNode::Logical { - operator: asap_types::query_plan::residual::ResidualQueryOperator::Binary { - operation: asap_types::query_plan::residual::BinaryOperation::FiniteDiv, + operator: asap_types::query_plan::query_time::QueryTimeOperator::Binary { + operation: + asap_types::query_plan::query_time::BinaryOperation::FiniteDiv, .. }, .. @@ -4318,7 +4327,7 @@ pub(crate) mod tests { ) .unwrap(); installed.document.schema_version = - asap_types::executable_plan::MAINTENANCE_DAG_SCHEMA_VERSION; + asap_types::executable_plan::PRECOMPUTE_DAG_SCHEMA_VERSION; assert!(plan .precompute_plan .validate() @@ -4540,11 +4549,11 @@ pub(crate) mod tests { .compile_promql(request, environment) .unwrap(); let entry = plan.query_plan.lookup(query).unwrap(); - use crate::query_plan::{residual::ResidualQueryOperator, QueryPlanNode}; + use crate::query_plan::{query_time::QueryTimeOperator, QueryPlanNode}; assert!(matches!( &entry.nodes[&entry.root], QueryPlanNode::Logical { - operator: ResidualQueryOperator::Limit { n: 2, .. }, + operator: QueryTimeOperator::Limit { n: 2, .. }, .. } )); @@ -4553,7 +4562,7 @@ pub(crate) mod tests { .nodes .values() .filter(|node| matches!(node, QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { query }, .. + operator: QueryTimeOperator::ExactSubquery { query }, .. } if query == "sum by (job) (rate(m[1m]))")) .count(), 1 @@ -5090,7 +5099,7 @@ pub(crate) mod tests { crate::query_plan::QueryPlanNode::ExternalExact { .. } | crate::query_plan::QueryPlanNode::Logical { operator: - crate::query_plan::residual::ResidualQueryOperator::ExactSubquery { .. }, + crate::query_plan::query_time::QueryTimeOperator::ExactSubquery { .. }, .. } )) @@ -5709,7 +5718,7 @@ pub(crate) mod tests { .any(|entry| entry.nodes.values().any(|node| matches!( node, crate::query_plan::QueryPlanNode::Logical { - operator: crate::query_plan::residual::ResidualQueryOperator::Binary { .. }, + operator: crate::query_plan::query_time::QueryTimeOperator::Binary { .. }, .. } )))); @@ -6778,7 +6787,7 @@ pub(crate) mod tests { // with each operand keeping its own range. #[test] fn composable_binary_summarizes_each_prometheus_filtered_operand() { - use crate::query_plan::{residual::ResidualQueryOperator, QueryPlanNode}; + use crate::query_plan::{query_time::QueryTimeOperator, QueryPlanNode}; let mut snapshot: BackendLocalPlanningInput = serde_json::from_str(include_str!( "../../../docs/examples/asapquery-planning-snapshot.json" )) @@ -6792,9 +6801,9 @@ pub(crate) mod tests { let query = plan.query_plan.entries.values().next().unwrap(); let bindings = query.materialization_bindings(); // Both operands now hold a summary. The filtered denominator is no - // longer a typed residual: its 5m range has a candidate, so it gets + // longer query-time only: its 5m range has a candidate, so it gets // its own summary over the filtered population rather than exact - // execution. Nothing about the filter forced the residual -- the + // execution. Nothing about the filter forced query-time execution -- the // missing 5m window candidate did, and this test previously pinned // that artifact as intended behavior. let bound = bindings @@ -6823,7 +6832,7 @@ pub(crate) mod tests { assert!(!query.nodes.values().any(|node| matches!( node, QueryPlanNode::Logical { - operator: ResidualQueryOperator::Scan { .. }, + operator: QueryTimeOperator::Scan { .. }, .. } ))); diff --git a/control_plane/src/physical/maintained_population.rs b/control_plane/src/physical/maintained_population.rs index dd1f856ba..1b0e1d97c 100644 --- a/control_plane/src/physical/maintained_population.rs +++ b/control_plane/src/physical/maintained_population.rs @@ -2,7 +2,7 @@ use super::compiler::{CompileError, PhysicalCompilationRequest, QueryCompilationInput}; use asap_types::query_plan::{ current_series::{SeriesPopulation, SeriesReadout}, - residual::{Grouping, LabelMatch, LabelMatcher, ResidualQueryOperator}, + query_time::{Grouping, LabelMatch, LabelMatcher, QueryTimeOperator}, }; use planner_types::post_asap::{ maintained_population::*, SummaryExpr, SummaryNode, ValueOperation, @@ -43,7 +43,7 @@ pub(super) fn supported(request: &PhysicalCompilationRequest) -> bool { pub(super) fn operator( request: &PhysicalCompilationRequest, query: &QueryCompilationInput, -) -> Result, CompileError> { +) -> Result, CompileError> { let Some((spec, readout)) = selected(&query.selected_plan_root) else { return Ok(None); }; @@ -103,7 +103,7 @@ pub(super) fn operator( PopulationReadout::Count => SeriesReadout::Count, PopulationReadout::Average => SeriesReadout::Average, }; - Ok(Some(ResidualQueryOperator::CurrentSeries { + Ok(Some(QueryTimeOperator::CurrentSeries { population, readout, })) diff --git a/control_plane/src/physical/plan_dot.rs b/control_plane/src/physical/plan_dot.rs index e0753c716..22dc9ed99 100644 --- a/control_plane/src/physical/plan_dot.rs +++ b/control_plane/src/physical/plan_dot.rs @@ -5,7 +5,7 @@ use super::compiler::CompiledPhysicalPlan; use crate::query_plan::QueryPlanNode; -use asap_types::query_plan::residual::ResidualQueryOperator; +use asap_types::query_plan::query_time::QueryTimeOperator; /// Render the selected precompute and query execution DAGs as deterministic DOT. pub fn render(plan: &CompiledPhysicalPlan) -> String { @@ -151,7 +151,9 @@ fn query_node_label(node: &QueryPlanNode) -> String { } QueryPlanNode::RelationalJoin { .. } => "RelationalJoin".into(), QueryPlanNode::Relational { .. } => "Relational".into(), - QueryPlanNode::Logical { operator, .. } => format!("Logical\n{}", residual_label(operator)), + QueryPlanNode::Logical { operator, .. } => { + format!("Logical\n{}", query_time_label(operator)) + } QueryPlanNode::Scalar { value } => format!("Scalar\n{value}"), QueryPlanNode::Binary { operator, .. } => format!("Binary\n{operator:?}"), QueryPlanNode::ReduceSum { .. } => "ReduceSum".into(), @@ -169,21 +171,21 @@ fn query_node_label(node: &QueryPlanNode) -> String { } } -fn residual_label(operator: &ResidualQueryOperator) -> &'static str { +fn query_time_label(operator: &QueryTimeOperator) -> &'static str { match operator { - ResidualQueryOperator::CurrentSeries { .. } => "CurrentSeries", - ResidualQueryOperator::ExactSubquery { .. } => "ExactSubquery", - ResidualQueryOperator::CandidateExactSubquery { .. } => "CandidateExactSubquery", - ResidualQueryOperator::Scan { .. } => "Scan", - ResidualQueryOperator::UnaryNegate => "UnaryNegate", - ResidualQueryOperator::VectorToScalar => "VectorToScalar", - ResidualQueryOperator::Aggregate { .. } => "Aggregate", - ResidualQueryOperator::Limit { .. } => "Limit", - ResidualQueryOperator::Binary { .. } => "Binary", - ResidualQueryOperator::Temporal { .. } => "Temporal", - ResidualQueryOperator::Sort { .. } => "Sort", - ResidualQueryOperator::HistogramQuantile => "HistogramQuantile", - ResidualQueryOperator::Subquery { .. } => "Subquery", + QueryTimeOperator::CurrentSeries { .. } => "CurrentSeries", + QueryTimeOperator::ExactSubquery { .. } => "ExactSubquery", + QueryTimeOperator::CandidateExactSubquery { .. } => "CandidateExactSubquery", + QueryTimeOperator::Scan { .. } => "Scan", + QueryTimeOperator::UnaryNegate => "UnaryNegate", + QueryTimeOperator::VectorToScalar => "VectorToScalar", + QueryTimeOperator::Aggregate { .. } => "Aggregate", + QueryTimeOperator::Limit { .. } => "Limit", + QueryTimeOperator::Binary { .. } => "Binary", + QueryTimeOperator::Temporal { .. } => "Temporal", + QueryTimeOperator::Sort { .. } => "Sort", + QueryTimeOperator::HistogramQuantile => "HistogramQuantile", + QueryTimeOperator::Subquery { .. } => "Subquery", } } diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index 963ebccac..8176d91f7 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -250,7 +250,7 @@ pub fn manifest( for node in entry.nodes.values() { if let crate::query_plan::QueryPlanNode::Logical { operator: - crate::query_plan::residual::ResidualQueryOperator::CurrentSeries { + crate::query_plan::query_time::QueryTimeOperator::CurrentSeries { population, .. }, .. @@ -271,7 +271,7 @@ pub fn manifest( if matches!( node, crate::query_plan::QueryPlanNode::Logical { - operator: crate::query_plan::residual::ResidualQueryOperator::Scan { .. }, + operator: crate::query_plan::query_time::QueryTimeOperator::Scan { .. }, .. } ) { @@ -281,8 +281,8 @@ pub fn manifest( } if let crate::query_plan::QueryPlanNode::Logical { operator: - crate::query_plan::residual::ResidualQueryOperator::ExactSubquery { query } - | crate::query_plan::residual::ResidualQueryOperator::CandidateExactSubquery { + crate::query_plan::query_time::QueryTimeOperator::ExactSubquery { query } + | crate::query_plan::query_time::QueryTimeOperator::CandidateExactSubquery { query, .. }, @@ -831,7 +831,7 @@ fn materialization_candidates( } let mut keys = BTreeSet::new(); for query in &request.queries { - match crate::query_plan::residual::eligible_materialization_keys( + match crate::query_plan::query_time::eligible_materialization_keys( &query.query_string, &query.selected_plan_root, ) { diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index a1e70c590..04a7afd85 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -2,7 +2,7 @@ //! Serving consumes asap_types::query_plan; compilation stays in this component. mod clickhouse_exact; -pub mod residual; +pub mod query_time; pub use asap_types::query_plan::*; #[cfg(test)] @@ -132,7 +132,7 @@ where instant, fallback, }; - residual::finalize_residuals(&mut entry)?; + query_time::finalize_query_time_nodes(&mut entry)?; Ok(entry) } @@ -304,9 +304,9 @@ where let id = QueryNodeId(self.next_id); self.next_id += 1; self.seen.insert(identity, id); - let residual = match (&self.logical_source, &node.expr) { + let query_time = match (&self.logical_source, &node.expr) { (Some(original), SummaryExpr::KeepPreAsap(expr)) => { - Some(residual::residual_nodes(original, expr)?) + Some(query_time::query_time_nodes(original, expr)?) } (Some(original), SummaryExpr::SummaryAgg { child, .. }) if matches!(child.expr, SummaryExpr::KeepPreAsap(_)) @@ -315,11 +315,11 @@ where Ok((_, Some(_), _)) ) => { - Some(residual::selected_residual_nodes(original, node)?) + Some(query_time::selected_query_time_nodes(original, node)?) } _ => None, }; - if let Some((root, nodes)) = residual { + if let Some((root, nodes)) = query_time { let id = self.graft(id, root, nodes)?; if let Some(lowered) = &mut self.lowered { lowered(node, id); @@ -387,11 +387,11 @@ where } if measures.len() == 1 => { use planner_types::pre_asap::AggIntent; let operation = match &measures[0] { - AggIntent::Sum { .. } => Some(residual::Aggregation::Sum), - AggIntent::Count { .. } => Some(residual::Aggregation::Count), - AggIntent::Min { .. } => Some(residual::Aggregation::Min), - AggIntent::Max { .. } => Some(residual::Aggregation::Max), - AggIntent::Avg { .. } => Some(residual::Aggregation::Avg), + AggIntent::Sum { .. } => Some(query_time::Aggregation::Sum), + AggIntent::Count { .. } => Some(query_time::Aggregation::Count), + AggIntent::Min { .. } => Some(query_time::Aggregation::Min), + AggIntent::Max { .. } => Some(query_time::Aggregation::Max), + AggIntent::Avg { .. } => Some(query_time::Aggregation::Avg), _ => { return Err(QueryPlanError::Invalid( "unsupported exact value aggregation".into(), @@ -419,11 +419,11 @@ where }) }) .collect::, _>>()?; - let grouping = residual::Grouping { + let grouping = query_time::Grouping { labels, without: keys.is_without(), }; - let operator = residual::ResidualQueryOperator::Aggregate { + let operator = query_time::QueryTimeOperator::Aggregate { operation: operation.expect("aggregate operation"), grouping, }; @@ -457,7 +457,8 @@ where let value_input = if let Some(original) = self.logical_source.as_ref().filter(|_| pruning.is_some()) { - let exact_expression = residual::selected_native_expression(original, values)?; + let exact_expression = + query_time::selected_native_expression(original, values)?; let value_id = QueryNodeId(self.next_id); self.next_id += 1; self.nodes.insert( @@ -524,7 +525,7 @@ where || operator.checked_relative_division || operator.checked_finite_division => { - let operator = residual::binary_operator(operator)?; + let operator = query_time::binary_operator(operator)?; QueryPlanNode::Logical { operator, inputs: vec![self.lower(lhs)?, self.lower(rhs)?], @@ -544,7 +545,7 @@ where planner_types::post_asap::ExactKind::Sum | planner_types::post_asap::ExactKind::Count ) { - let operator = residual::selected_aggregate_operator( + let operator = query_time::selected_aggregate_operator( self.logical_source.as_deref().unwrap(), node, )?; @@ -559,8 +560,8 @@ where return Ok(id); } let operation = match kind { - planner_types::post_asap::ExactKind::Sum => residual::Aggregation::Sum, - planner_types::post_asap::ExactKind::Count => residual::Aggregation::Count, + planner_types::post_asap::ExactKind::Sum => query_time::Aggregation::Sum, + planner_types::post_asap::ExactKind::Count => query_time::Aggregation::Count, _ => { return Err(QueryPlanError::Invalid( "unsupported aggregation over selected summary values".into(), @@ -587,9 +588,9 @@ where }) .collect::, _>>()?; QueryPlanNode::Logical { - operator: residual::ResidualQueryOperator::Aggregate { + operator: query_time::QueryTimeOperator::Aggregate { operation, - grouping: residual::Grouping { + grouping: query_time::Grouping { labels, without: keys.is_without(), }, @@ -725,7 +726,7 @@ where Err(error) => { if let Some(original) = &self.logical_source { let (root, nodes) = - residual::selected_residual_nodes(original, node)?; + query_time::selected_query_time_nodes(original, node)?; return self.graft(id, root, nodes); } return Err(error); @@ -1216,7 +1217,7 @@ mod tests { } .unwrap(); let QueryPlanNode::Logical { - operator: residual::ResidualQueryOperator::Binary { operation, .. }, + operator: query_time::QueryTimeOperator::Binary { operation, .. }, .. } = &entry.nodes[&entry.root] else { @@ -1228,9 +1229,9 @@ mod tests { assert_eq!( *operation, if relative { - residual::BinaryOperation::CheckedDiv + query_time::BinaryOperation::CheckedDiv } else { - residual::BinaryOperation::FiniteDiv + query_time::BinaryOperation::FiniteDiv } ); assert!(!entry.materialization_bindings().is_empty()); diff --git a/control_plane/src/query_plan/residual.rs b/control_plane/src/query_plan/query_time.rs similarity index 89% rename from control_plane/src/query_plan/residual.rs rename to control_plane/src/query_plan/query_time.rs index 39926848c..7f1afc89e 100644 --- a/control_plane/src/query_plan/residual.rs +++ b/control_plane/src/query_plan/query_time.rs @@ -1,4 +1,8 @@ -//! Typed residual operations compiled once by the control plane, never parsed at serving time. +//! The query-time half of a plan: everything a summary did not replace. +//! +//! Typed operations compiled once by the control plane, never parsed at +//! serving time. This is the counterpart to the precompute half, which runs +//! ahead of the query and ends at stored outputs. use super::{ FallbackPolicy, InstantExecution, QueryNodeId, QueryPlanEntry, QueryPlanError, QueryPlanNode, }; @@ -8,7 +12,7 @@ use promql_parser::{ }; use std::collections::BTreeMap; -pub use asap_types::query_plan::residual::*; +pub use asap_types::query_plan::query_time::*; /// Stable identity of a Planner-authorized materializable DAG leaf. This is a /// workload-selection key, not another physical materialization definition. @@ -54,7 +58,7 @@ impl Lower { } fn operation( &mut self, - operator: ResidualQueryOperator, + operator: QueryTimeOperator, inputs: Vec, ) -> Result { operator.validate(inputs.len())?; @@ -84,7 +88,7 @@ impl Lower { }) .collect(); self.operation( - ResidualQueryOperator::Scan { + QueryTimeOperator::Scan { metric: s.name.clone(), matchers, range_ms, @@ -101,7 +105,7 @@ impl Lower { Expr::Paren(p) => self.lower(&p.expr), Expr::Unary(u) => { let input = self.lower(&u.expr)?; - self.operation(ResidualQueryOperator::UnaryNegate, vec![input]) + self.operation(QueryTimeOperator::UnaryNegate, vec![input]) } Expr::VectorSelector(s) => self.scan(s, None), Expr::MatrixSelector(s) => self.scan(&s.vs, Some(millis(s.range)?)), @@ -111,7 +115,7 @@ impl Lower { } let input = self.lower(&s.expr)?; self.operation( - ResidualQueryOperator::Subquery { + QueryTimeOperator::Subquery { range_ms: millis(s.range)?, // Prometheus uses its configured default evaluation // interval when `[range:]` omits the resolution. The @@ -163,7 +167,7 @@ impl Lower { self.nodes = nodes_before; self.seen = seen_before; self.operation( - ResidualQueryOperator::ExactSubquery { + QueryTimeOperator::ExactSubquery { query: a.expr.to_string(), }, vec![], @@ -171,14 +175,14 @@ impl Lower { } }; let sorted = self.operation( - ResidualQueryOperator::Sort { + QueryTimeOperator::Sort { descending: true, grouping: grouping.clone(), }, vec![input], )?; return self.operation( - ResidualQueryOperator::Limit { + QueryTimeOperator::Limit { n: u64::try_from(k).unwrap_or(0), offset: 0, grouping, @@ -199,7 +203,7 @@ impl Lower { }; let input = self.lower(&a.expr)?; self.operation( - ResidualQueryOperator::Aggregate { + QueryTimeOperator::Aggregate { operation, grouping, }, @@ -208,23 +212,23 @@ impl Lower { } Expr::Call(c) => { let operator = match c.func.name { - "scalar" => ResidualQueryOperator::VectorToScalar, - "histogram_quantile" => ResidualQueryOperator::HistogramQuantile, - "sort" => ResidualQueryOperator::Sort { + "scalar" => QueryTimeOperator::VectorToScalar, + "histogram_quantile" => QueryTimeOperator::HistogramQuantile, + "sort" => QueryTimeOperator::Sort { descending: false, grouping: Grouping { labels: vec![], without: false, }, }, - "sort_desc" => ResidualQueryOperator::Sort { + "sort_desc" => QueryTimeOperator::Sort { descending: true, grouping: Grouping { labels: vec![], without: false, }, }, - name => ResidualQueryOperator::Temporal { + name => QueryTimeOperator::Temporal { operation: match name { "rate" => TemporalOperation::Rate, "increase" => TemporalOperation::Increase, @@ -271,7 +275,7 @@ impl Lower { }; let inputs = vec![self.lower(&b.lhs)?, self.lower(&b.rhs)?]; self.operation( - ResidualQueryOperator::Binary { + QueryTimeOperator::Binary { operation, return_bool: b.return_bool(), }, @@ -283,7 +287,7 @@ impl Lower { } } -/// Lower a Planner-authorized native residual into typed backend operations. +/// Lower a Planner-authorized native fragment into typed backend operations. /// Callers retain a separate external-native alternative for cost comparison. pub fn compile_logical( query_id: String, @@ -346,7 +350,7 @@ fn horizons(expr: &planner_types::pre_asap::QueryExpr, out: &mut Vec) { } } -/// Match residuals by semantic IR equality, not display text or source names. +/// Match fragments by semantic IR equality, not display text or source names. /// This ensures a subtree parsed for physical lowering is the subtree Planner kept. /// Does this re-parsed subtree denote the same computation as the Planner /// fragment? @@ -358,7 +362,7 @@ fn horizons(expr: &planner_types::pre_asap::QueryExpr, out: &mut Vec) { /// label the *query* mentions, while the same fragment re-parsed on its own /// carries only the labels *it* mentions. /// -/// So `sum by (label_0) (rate(data[1m]))` yields a residual whose leaf scan has +/// So `sum by (label_0) (rate(data[1m]))` yields a fragment whose leaf scan has /// columns `[ts, value, label_0]`, while re-parsing the subtree `rate(data[1m])` /// yields `[ts, value]`. Identical source, predicates, range, measures and /// reduction; one extra column that the isolated parse had no way to know @@ -374,44 +378,44 @@ fn horizons(expr: &planner_types::pre_asap::QueryExpr, out: &mut Vec) { /// equality. fn fragment_matches( candidate: &planner_types::pre_asap::QueryExpr, - residual: &planner_types::pre_asap::QueryExpr, + fragment: &planner_types::pre_asap::QueryExpr, ) -> bool { match ( serde_json::to_value(candidate), - serde_json::to_value(residual), + serde_json::to_value(fragment), ) { - (Ok(candidate), Ok(residual)) => same_modulo_open_leaf_schema(&candidate, &residual), + (Ok(candidate), Ok(fragment)) => same_modulo_open_leaf_schema(&candidate, &fragment), // Fall back to the strict comparison rather than accepting anything we // could not inspect. - _ => candidate == residual, + _ => candidate == fragment, } } fn same_modulo_open_leaf_schema( candidate: &serde_json::Value, - residual: &serde_json::Value, + fragment: &serde_json::Value, ) -> bool { use serde_json::Value; - match (candidate, residual) { - (Value::Object(candidate), Value::Object(residual)) => { - if is_open_schema(candidate) && is_open_schema(residual) { - return open_schema_is_widened(candidate, residual); + match (candidate, fragment) { + (Value::Object(candidate), Value::Object(fragment)) => { + if is_open_schema(candidate) && is_open_schema(fragment) { + return open_schema_is_widened(candidate, fragment); } - candidate.len() == residual.len() + candidate.len() == fragment.len() && candidate.iter().all(|(key, value)| { - residual + fragment .get(key) .is_some_and(|other| same_modulo_open_leaf_schema(value, other)) }) } - (Value::Array(candidate), Value::Array(residual)) => { - candidate.len() == residual.len() + (Value::Array(candidate), Value::Array(fragment)) => { + candidate.len() == fragment.len() && candidate .iter() - .zip(residual) + .zip(fragment) .all(|(a, b)| same_modulo_open_leaf_schema(a, b)) } - _ => candidate == residual, + _ => candidate == fragment, } } @@ -424,25 +428,25 @@ fn is_open_schema(value: &serde_json::Map) -> bool { /// floor. Every other schema field still has to agree exactly. fn open_schema_is_widened( candidate: &serde_json::Map, - residual: &serde_json::Map, + fragment: &serde_json::Map, ) -> bool { let (Some(narrow), Some(wide)) = ( candidate.get("columns").and_then(|v| v.as_array()), - residual.get("columns").and_then(|v| v.as_array()), + fragment.get("columns").and_then(|v| v.as_array()), ) else { return false; }; candidate .iter() .filter(|(key, _)| key.as_str() != "columns") - .all(|(key, value)| residual.get(key) == Some(value)) + .all(|(key, value)| fragment.get(key) == Some(value)) && narrow.len() <= wide.len() && narrow.iter().zip(wide).all(|(a, b)| a == b) } -pub(super) fn residual_nodes( +pub(super) fn query_time_nodes( original: &str, - residual: &planner_types::pre_asap::QueryExpr, + fragment: &planner_types::pre_asap::QueryExpr, ) -> Result<(QueryNodeId, BTreeMap), QueryPlanError> { fn visit<'a>(expr: &'a Expr, out: &mut Vec<&'a Expr>) { out.push(expr); @@ -470,13 +474,13 @@ pub(super) fn residual_nodes( // not the compatibility parser's default. Explicit matrix ranges remain // query-owned and equality still checks the complete tree. let mut intervals = vec![1_000]; - horizons(residual, &mut intervals); + horizons(fragment, &mut intervals); intervals.sort_unstable(); intervals.dedup(); // Accuracy annotations select a candidate, but exact execution still // implements that candidate's computation. Reconstruct the same typed IR // before comparing it; do not erase operators or source predicates. - let accuracy = match residual { + let accuracy = match fragment { planner_types::pre_asap::QueryExpr::Aggregate { measures, .. } => measures .iter() .find_map(|intent| { @@ -501,7 +505,7 @@ pub(super) fn residual_nodes( accuracy.clone(), *interval, ) { - if fragment_matches(&candidate, residual) { + if fragment_matches(&candidate, fragment) { let mut lower = Lower { nodes: BTreeMap::new(), seen: BTreeMap::new(), @@ -513,13 +517,13 @@ pub(super) fn residual_nodes( } } Err(invalid( - "Planner residual does not match any original query subtree", + "Planner fragment does not match any original query subtree", )) } pub(super) fn binary_operator( operator: &planner_types::post_asap::BinaryOperator, -) -> Result { +) -> Result { if operator.checked_relative_division || operator.checked_finite_division { if (operator.checked_relative_division && operator.checked_finite_division) || operator.vector_match.is_some() @@ -532,7 +536,7 @@ pub(super) fn binary_operator( { return Err(invalid("invalid Planner checked division contract")); } - return Ok(ResidualQueryOperator::Binary { + return Ok(QueryTimeOperator::Binary { operation: if operator.checked_finite_division { BinaryOperation::FiniteDiv } else { @@ -542,7 +546,7 @@ pub(super) fn binary_operator( }); } if operator.vector_match.is_some() { - return Err(invalid("explicit residual vector matching unsupported")); + return Err(invalid("explicit fragment vector matching unsupported")); } let operation = match operator.kind.to_string().as_str() { "+" => BinaryOperation::Add, @@ -563,7 +567,7 @@ pub(super) fn binary_operator( ))) } }; - Ok(ResidualQueryOperator::Binary { + Ok(QueryTimeOperator::Binary { operation, return_bool: false, }) @@ -571,7 +575,7 @@ pub(super) fn binary_operator( /// Prove a physical-native substitute represents exactly the selected summary leaf. /// A second Planner invocation is an equality witness, not a replacement selection. -pub(crate) fn selected_residual_nodes( +pub(crate) fn selected_query_time_nodes( original: &str, selected: &planner_types::post_asap::SummaryNode, ) -> Result<(QueryNodeId, BTreeMap), QueryPlanError> { @@ -592,7 +596,7 @@ pub(super) fn selected_native_expression( ) -> Result { if !selected.guarantee.as_ref().is_some_and(|g| g.is_exact()) { return Err(invalid( - "native residual substitution requires an exact selected value", + "native fragment substitution requires an exact selected value", )); } let selected = match &selected.expr { @@ -704,11 +708,11 @@ pub(super) fn selected_native_expression( pub(super) fn selected_aggregate_operator( original: &str, selected: &planner_types::post_asap::SummaryNode, -) -> Result { - let (root, nodes) = selected_residual_nodes(original, selected)?; +) -> Result { + let (root, nodes) = selected_query_time_nodes(original, selected)?; match nodes.get(&root) { Some(QueryPlanNode::Logical { - operator: operator @ ResidualQueryOperator::Aggregate { .. }, + operator: operator @ QueryTimeOperator::Aggregate { .. }, .. }) => Ok(operator.clone()), _ => Err(invalid( @@ -792,21 +796,21 @@ mod hybrid_tests { assert!(!entry.nodes.values().any(|node| matches!( node, QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { .. }, + operator: QueryTimeOperator::ExactSubquery { .. }, .. } ))); assert!(!entry.nodes.values().any(|node| matches!( node, QueryPlanNode::Logical { - operator: ResidualQueryOperator::Scan { .. }, + operator: QueryTimeOperator::Scan { .. }, .. } ))); assert!(matches!( entry.nodes[&entry.root], QueryPlanNode::Logical { - operator: ResidualQueryOperator::Binary { .. }, + operator: QueryTimeOperator::Binary { .. }, .. } )); @@ -832,7 +836,7 @@ mod hybrid_tests { .unwrap(); let selected = crate::planner_selection::plan_test_query(&canonical).unwrap(); assert!( - selected_residual_nodes("sum_over_time(m{job=\"worker\"}[5m])", &selected).is_err() + selected_query_time_nodes("sum_over_time(m{job=\"worker\"}[5m])", &selected).is_err() ); } } @@ -869,7 +873,7 @@ mod planner_workload_tests { assert_eq!(plan.operator_name(plan.roots()[0]), Some("Limit")); } QueryPlanNode::Logical { - operator: ResidualQueryOperator::Limit { .. }, + operator: QueryTimeOperator::Limit { .. }, .. } => {} _ => panic!("expected local Limit, got {node:?}"), @@ -975,7 +979,7 @@ mod planner_workload_tests { let operator = selected_aggregate_operator(query, selected).unwrap(); assert!(matches!( operator, - ResidualQueryOperator::Aggregate { + QueryTimeOperator::Aggregate { operation: Aggregation::Max, .. } @@ -997,7 +1001,7 @@ mod planner_workload_tests { ) .unwrap(); let maximum = crate::planner_selection::plan_test_query(&maximum).unwrap(); - let result = selected_residual_nodes("min(m) + max(m)", &selected); + let result = selected_query_time_nodes("min(m) + max(m)", &selected); if selected == maximum { assert!(result.is_err()); } else { @@ -1005,7 +1009,7 @@ mod planner_workload_tests { assert!(matches!( nodes[&root], QueryPlanNode::Logical { - operator: ResidualQueryOperator::Aggregate { + operator: QueryTimeOperator::Aggregate { operation: Aggregation::Min, .. }, @@ -1041,10 +1045,10 @@ pub(crate) fn selected_range_max_materialization( ) { return Ok(None); } - let (root, nodes) = selected_residual_nodes(original, node)?; + let (root, nodes) = selected_query_time_nodes(original, node)?; let Some(QueryPlanNode::Logical { operator: - ResidualQueryOperator::Temporal { + QueryTimeOperator::Temporal { operation: TemporalOperation::Max, }, inputs, @@ -1057,7 +1061,7 @@ pub(crate) fn selected_range_max_materialization( } let Some(QueryPlanNode::Logical { operator: - ResidualQueryOperator::Scan { + QueryTimeOperator::Scan { metric: Some(metric), matchers, range_ms: Some(range_ms), @@ -1158,7 +1162,7 @@ fn counter_contract( nodes: &BTreeMap, ) -> Option { let QueryPlanNode::Logical { - operator: ResidualQueryOperator::Temporal { operation }, + operator: QueryTimeOperator::Temporal { operation }, inputs, } = nodes.get(&root)? else { @@ -1173,7 +1177,7 @@ fn counter_contract( } let QueryPlanNode::Logical { operator: - ResidualQueryOperator::Scan { + QueryTimeOperator::Scan { metric: Some(metric), matchers, range_ms: Some(range_ms), @@ -1208,7 +1212,7 @@ pub(crate) fn selected_counter_materialization( ) { return Ok(None); } - let (root, nodes) = selected_residual_nodes(original, node)?; + let (root, nodes) = selected_query_time_nodes(original, node)?; counter_contract(root, &nodes) .map(materialization_candidate_key) .transpose() @@ -1227,9 +1231,9 @@ fn prune(entry: &mut QueryPlanEntry) { entry.nodes.retain(|id, _| seen.contains(id)); } -/// Finish the installed DAG by externalizing every residual raw subtree. -pub fn finalize_residuals(entry: &mut QueryPlanEntry) -> Result<(), QueryPlanError> { - externalize_residuals(entry)?; +/// Finish the installed DAG by externalizing every fragment raw subtree. +pub fn finalize_query_time_nodes(entry: &mut QueryPlanEntry) -> Result<(), QueryPlanError> { + externalize_query_time_nodes(entry)?; assign_retention(entry) } @@ -1294,7 +1298,7 @@ fn expression_shape( .get(&id) .ok_or_else(|| invalid("missing expression node"))?; if let QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { query }, + operator: QueryTimeOperator::ExactSubquery { query }, .. } = node { @@ -1321,9 +1325,9 @@ fn expression_shape( Ok(format!("{value}({})", children.join(";"))) } -/// Collapse only maximal exact residual subtrees whose full typed expression is +/// Collapse only maximal exact fragment subtrees whose full typed expression is /// witnessed in the original query. Matrix boundaries remain inside Prometheus. -pub fn externalize_residuals(entry: &mut QueryPlanEntry) -> Result<(), QueryPlanError> { +pub fn externalize_query_time_nodes(entry: &mut QueryPlanEntry) -> Result<(), QueryPlanError> { fn gather(expr: &Expr, out: &mut Vec) { if !matches!(expr, Expr::MatrixSelector(_) | Expr::Subquery(_)) { out.push(expr.clone()); @@ -1371,7 +1375,7 @@ pub fn externalize_residuals(entry: &mut QueryPlanEntry) -> Result<(), QueryPlan ) || matches!( node, QueryPlanNode::Logical { - operator: ResidualQueryOperator::Limit { .. }, + operator: QueryTimeOperator::Limit { .. }, .. } ); @@ -1379,9 +1383,9 @@ pub fn externalize_residuals(entry: &mut QueryPlanEntry) -> Result<(), QueryPlan if let QueryPlanNode::Logical { operator, .. } = node { exact = matches!( operator, - ResidualQueryOperator::Scan { .. } - | ResidualQueryOperator::ExactSubquery { .. } - | ResidualQueryOperator::CandidateExactSubquery { .. } + QueryTimeOperator::Scan { .. } + | QueryTimeOperator::ExactSubquery { .. } + | QueryTimeOperator::CandidateExactSubquery { .. } ); } for child in node.inputs() { @@ -1396,8 +1400,8 @@ pub fn externalize_residuals(entry: &mut QueryPlanEntry) -> Result<(), QueryPlan if matches!( entry.nodes.get(&id), Some(QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { .. } - | ResidualQueryOperator::CandidateExactSubquery { .. }, + operator: QueryTimeOperator::ExactSubquery { .. } + | QueryTimeOperator::CandidateExactSubquery { .. }, .. }) ) { @@ -1410,7 +1414,7 @@ pub fn externalize_residuals(entry: &mut QueryPlanEntry) -> Result<(), QueryPlan entry.nodes.insert( id, QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { + operator: QueryTimeOperator::ExactSubquery { query: query.clone(), }, inputs: vec![], @@ -1427,7 +1431,7 @@ pub fn externalize_residuals(entry: &mut QueryPlanEntry) -> Result<(), QueryPlan matches!( node, QueryPlanNode::Logical { - operator: ResidualQueryOperator::Scan { .. }, + operator: QueryTimeOperator::Scan { .. }, .. } ) @@ -1453,7 +1457,7 @@ fn assign_retention(entry: &mut QueryPlanEntry) -> Result<(), QueryPlanError> { .ok_or_else(|| invalid("missing index ancestor"))?; let mut child_depth = depth; if let QueryPlanNode::Logical { operator, .. } = node { - if let ResidualQueryOperator::Subquery { + if let QueryTimeOperator::Subquery { range_ms, offset_ms, .. @@ -1527,7 +1531,7 @@ mod remote_boundary_regressions { #[cfg(test)] mod tests { use super::*; - // A grouped query's residual carries the grouping label in its leaf scan + // A grouped query's fragment carries the grouping label in its leaf scan // schema, because Planner resolves that schema against the whole query. // Re-parsing the subtree alone cannot know the label, so requiring equal // column sets rejected a fragment that is the subtree. @@ -1540,15 +1544,15 @@ mod tests { "sum_over_time(data[1m])", ), ] { - let residual = crate::query_parser::parse_query_expr_with_interval( + let fragment = crate::query_parser::parse_query_expr_with_interval( subtree, planner_types::types::AccuracyTarget::Exact, 60_000, ) .unwrap(); assert!( - residual_nodes(query, &residual).is_ok(), - "{query}: residual {subtree} must resolve against its own query" + query_time_nodes(query, &fragment).is_ok(), + "{query}: fragment {subtree} must resolve against its own query" ); } } @@ -1557,18 +1561,18 @@ mod tests { // that actually identifies the computation still has to match exactly. #[test] fn widened_leaf_schema_does_not_excuse_a_different_computation() { - let residual = crate::query_parser::parse_query_expr_with_interval( + let fragment = crate::query_parser::parse_query_expr_with_interval( "rate(data[1m])", planner_types::types::AccuracyTarget::Exact, 60_000, ) .unwrap(); // Different metric. - assert!(residual_nodes("sum by (label_0) (rate(other[1m]))", &residual).is_err()); + assert!(query_time_nodes("sum by (label_0) (rate(other[1m]))", &fragment).is_err()); // Different range. - assert!(residual_nodes("sum by (label_0) (rate(data[2m]))", &residual).is_err()); + assert!(query_time_nodes("sum by (label_0) (rate(data[2m]))", &fragment).is_err()); // Different function. - assert!(residual_nodes("sum by (label_0) (increase(data[1m]))", &residual).is_err()); + assert!(query_time_nodes("sum by (label_0) (increase(data[1m]))", &fragment).is_err()); // Different matcher. let filtered = crate::query_parser::parse_query_expr_with_interval( "rate(data{job=\"api\"}[1m])", @@ -1576,7 +1580,7 @@ mod tests { 60_000, ) .unwrap(); - assert!(residual_nodes( + assert!(query_time_nodes( "sum by (label_0) (rate(data{job=\"worker\"}[1m]))", &filtered ) @@ -1586,21 +1590,21 @@ mod tests { // A workload horizon changes the equality witness, never its filter or explicit range. #[test] fn workload_horizon_residual_keeps_semantic_equality() { - let residual = crate::query_parser::parse_query_expr_with_interval( + let fragment = crate::query_parser::parse_query_expr_with_interval( "sum(m{job=\"api\"})", planner_types::types::AccuracyTarget::Exact, 5_000, ) .unwrap(); - assert!(residual_nodes("sum(m{job=\"api\"})", &residual).is_ok()); - assert!(residual_nodes("sum(m{job=\"worker\"})", &residual).is_err()); + assert!(query_time_nodes("sum(m{job=\"api\"})", &fragment).is_ok()); + assert!(query_time_nodes("sum(m{job=\"worker\"})", &fragment).is_err()); let range = crate::query_parser::parse_query_expr_with_interval( "sum_over_time(m[1m])", planner_types::types::AccuracyTarget::Exact, 5_000, ) .unwrap(); - assert!(residual_nodes("sum_over_time(m[2m])", &range).is_err()); + assert!(query_time_nodes("sum_over_time(m[2m])", &range).is_err()); } fn instant() -> InstantExecution { @@ -1618,7 +1622,7 @@ mod tests { serde_json::from_str(include_str!("../../tests/fixtures/o11y_queries.json")).unwrap(); for row in corpus["queries"].as_array().unwrap() { let query = row["query"].as_str().unwrap(); - let entry = crate::query_plan::residual::compile_logical( + let entry = crate::query_plan::query_time::compile_logical( row["id"].as_str().unwrap().into(), query.into(), instant(), @@ -1635,22 +1639,22 @@ mod tests { } } #[test] - fn residual_mapping_preserves_filters_and_rejects_different_sources() { + fn query_time_mapping_preserves_filters_and_rejects_different_sources() { // Physical lowering must prove correspondence with the Planner-kept semantic subtree. let query = "sum(rate(requests_total{job=\"api\"}[5m]))"; - let residual = crate::query_parser::parse_query_expr_canonical( + let fragment = crate::query_parser::parse_query_expr_canonical( query, planner_types::types::AccuracyTarget::Exact, ) .unwrap(); - let (_, nodes) = residual_nodes(query, &residual).unwrap(); - assert!(nodes.values().any(|node| matches!(node, QueryPlanNode::Logical { operator: ResidualQueryOperator::Scan { matchers, .. }, .. } if matchers.iter().any(|m| m.name == "job" && m.value == "api")))); - assert!(residual_nodes("sum(rate(other_total[5m]))", &residual).is_err()); + let (_, nodes) = query_time_nodes(query, &fragment).unwrap(); + assert!(nodes.values().any(|node| matches!(node, QueryPlanNode::Logical { operator: QueryTimeOperator::Scan { matchers, .. }, .. } if matchers.iter().any(|m| m.name == "job" && m.value == "api")))); + assert!(query_time_nodes("sum(rate(other_total[5m]))", &fragment).is_err()); } #[test] fn repeated_subexpressions_share_node_identity() { // Serialized edges must retain CSE rather than duplicating raw work. - let entry = crate::query_plan::residual::compile_logical( + let entry = crate::query_plan::query_time::compile_logical( "q".into(), "sum(up) / sum(up)".into(), instant(), @@ -1665,10 +1669,8 @@ mod tests { #[test] fn malformed_operator_arity_is_rejected_at_installation() { // A serialized graph cannot bypass the operation's input contract. - assert!(ResidualQueryOperator::HistogramQuantile - .validate(1) - .is_err()); - assert!(ResidualQueryOperator::Subquery { + assert!(QueryTimeOperator::HistogramQuantile.validate(1).is_err()); + assert!(QueryTimeOperator::Subquery { range_ms: 60_000, step_ms: 0, offset_ms: 0 @@ -1698,7 +1700,7 @@ mod tests { 3, ), ] { - let entry = crate::query_plan::residual::compile_logical( + let entry = crate::query_plan::query_time::compile_logical( "topk".into(), query.into(), instant(), @@ -1708,7 +1710,7 @@ mod tests { assert!(matches!( entry.nodes[&entry.root], QueryPlanNode::Logical { - operator: ResidualQueryOperator::Limit { n: actual, .. }, + operator: QueryTimeOperator::Limit { n: actual, .. }, .. } if actual == k )); @@ -1717,7 +1719,7 @@ mod tests { #[test] fn topk_keeps_unsupported_child_as_exact_leaf() { - let entry = crate::query_plan::residual::compile_logical( + let entry = crate::query_plan::query_time::compile_logical( "topk-subquery".into(), "topk(3, label_replace(memory_bytes, \"dst\", \"$1\", \"src\", \"(.*)\"))".into(), instant(), @@ -1727,14 +1729,14 @@ mod tests { assert!(matches!( entry.nodes[&entry.root], QueryPlanNode::Logical { - operator: ResidualQueryOperator::Limit { n: 3, .. }, + operator: QueryTimeOperator::Limit { n: 3, .. }, .. } )); assert!(entry.nodes.values().any(|node| matches!( node, QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { .. }, + operator: QueryTimeOperator::ExactSubquery { .. }, .. } ))); @@ -1746,7 +1748,7 @@ mod tests { ("topk by (cluster) (2, m)", vec!["cluster"], false), ("topk without (pod) (2, m)", vec!["pod"], true), ] { - let entry = crate::query_plan::residual::compile_logical( + let entry = crate::query_plan::query_time::compile_logical( "topk-group".into(), query.into(), instant(), @@ -1756,7 +1758,7 @@ mod tests { assert!(matches!( &entry.nodes[&entry.root], QueryPlanNode::Logical { - operator: ResidualQueryOperator::Limit { grouping, .. }, + operator: QueryTimeOperator::Limit { grouping, .. }, .. } if grouping.labels == labels && grouping.without == without )); diff --git a/crates/asap_types/src/derived_input.rs b/crates/asap_types/src/derived_input.rs index b3556e69d..e7aba6906 100644 --- a/crates/asap_types/src/derived_input.rs +++ b/crates/asap_types/src/derived_input.rs @@ -39,7 +39,7 @@ impl DerivedInputIdentity { ) -> Result { if ![ crate::executable_plan::OWNED_POST_ASAP_DAG_SCHEMA_VERSION, - crate::executable_plan::MAINTENANCE_DAG_SCHEMA_VERSION, + crate::executable_plan::PRECOMPUTE_DAG_SCHEMA_VERSION, ] .contains(&document.schema_version) { diff --git a/crates/asap_types/src/executable_plan.rs b/crates/asap_types/src/executable_plan.rs index c475584a6..a562dc1e2 100644 --- a/crates/asap_types/src/executable_plan.rs +++ b/crates/asap_types/src/executable_plan.rs @@ -21,7 +21,7 @@ use serde::{Deserialize, Serialize}; pub struct QueryNodeId(pub u64); pub const OWNED_POST_ASAP_DAG_SCHEMA_VERSION: u32 = 5; -pub const MAINTENANCE_DAG_SCHEMA_VERSION: u32 = 6; +pub const PRECOMPUTE_DAG_SCHEMA_VERSION: u32 = 6; /// Versioned, language-neutral Planner DAG persisted with an installed plan. /// Plan lifecycle belongs to the enclosing `PrecomputePlan`; this document @@ -187,15 +187,15 @@ impl InstalledPostAsapDag { if self.document.query_id.trim().is_empty() { return Err("invalid post-ASAP DAG document identity/version".into()); } - if self.document.schema_version != MAINTENANCE_DAG_SCHEMA_VERSION { + if self.document.schema_version != PRECOMPUTE_DAG_SCHEMA_VERSION { return Err("unsupported maintenance DAG document version".into()); } - self.binding.validate_maintenance(&self.document.decode()?) + self.binding.validate_precompute(&self.document.decode()?) } /// Project the selected semantic DAG onto the maintenance ancestors of its /// stored outputs. - pub fn maintenance_projection(mut self) -> Result { + pub fn precompute_projection(mut self) -> Result { if self.document.schema_version != OWNED_POST_ASAP_DAG_SCHEMA_VERSION { return Err("selected DAG has an unsupported document version".into()); } @@ -230,7 +230,7 @@ impl InstalledPostAsapDag { .retain(|edge| included.contains(&edge.producer) && included.contains(&edge.consumer)); self.binding.nodes.retain(|id, _| included.contains(id)); self.document.root = *self.binding.precompute_sinks.first().unwrap(); - self.document.schema_version = MAINTENANCE_DAG_SCHEMA_VERSION; + self.document.schema_version = PRECOMPUTE_DAG_SCHEMA_VERSION; self.validate()?; Ok(self) } @@ -268,7 +268,7 @@ impl BackendExecutableBinding { self.nodes.get(&id) } - pub fn validate_maintenance(&self, dag: &ExecutableDag) -> Result<(), String> { + pub fn validate_precompute(&self, dag: &ExecutableDag) -> Result<(), String> { let ids = dag .nodes .iter() diff --git a/crates/asap_types/src/query_plan.rs b/crates/asap_types/src/query_plan.rs index 452943cca..c4c2516c8 100644 --- a/crates/asap_types/src/query_plan.rs +++ b/crates/asap_types/src/query_plan.rs @@ -6,7 +6,7 @@ //! searching for compatible materializations. pub mod current_series; -pub mod residual; +pub mod query_time; use std::collections::{BTreeMap, BTreeSet}; @@ -694,7 +694,7 @@ pub enum QueryPlanNode { output_schema: planner_types::post_asap::SummarySchema, }, Logical { - operator: residual::ResidualQueryOperator, + operator: query_time::QueryTimeOperator, inputs: Vec, }, Scalar { diff --git a/crates/asap_types/src/query_plan/current_series.rs b/crates/asap_types/src/query_plan/current_series.rs index 4eb69cf18..dc175183c 100644 --- a/crates/asap_types/src/query_plan/current_series.rs +++ b/crates/asap_types/src/query_plan/current_series.rs @@ -1,6 +1,6 @@ //! Maintained current-value populations, shared independently of q and k. use super::{ - residual::{Grouping, LabelMatcher}, + query_time::{Grouping, LabelMatcher}, QueryPlanError, }; use serde::{Deserialize, Serialize}; diff --git a/crates/asap_types/src/query_plan/residual.rs b/crates/asap_types/src/query_plan/query_time.rs similarity index 95% rename from crates/asap_types/src/query_plan/residual.rs rename to crates/asap_types/src/query_plan/query_time.rs index 002f98dab..a5b3fa9ef 100644 --- a/crates/asap_types/src/query_plan/residual.rs +++ b/crates/asap_types/src/query_plan/query_time.rs @@ -1,4 +1,4 @@ -//! Typed installed residual operators; no Planner selection or AST lowering. +//! Typed installed query-time operators; no Planner selection or AST lowering. use super::QueryPlanError; use promql_parser::parser::{self, Expr}; use serde::{Deserialize, Serialize}; @@ -8,7 +8,7 @@ fn invalid(message: impl Into) -> QueryPlanError { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] -pub enum ResidualQueryOperator { +pub enum QueryTimeOperator { /// Readout over a bounded current-value population maintained at ingest. CurrentSeries { population: super::current_series::SeriesPopulation, @@ -123,7 +123,7 @@ pub enum TemporalOperation { Count, } -impl ResidualQueryOperator { +impl QueryTimeOperator { pub fn validate(&self, inputs: usize) -> Result<(), QueryPlanError> { if let Self::CurrentSeries { population, @@ -184,5 +184,5 @@ impl ResidualQueryOperator { } // Compatibility imports; new callers use the domain names above. -#[deprecated(note = "Use ResidualQueryOperator")] -pub use ResidualQueryOperator as LogicalOperator; +#[deprecated(note = "Use QueryTimeOperator")] +pub use QueryTimeOperator as LogicalOperator; diff --git a/data_plane/src/precompute_engine/engine.rs b/data_plane/src/precompute_engine/engine.rs index fa451ca36..b61a87602 100644 --- a/data_plane/src/precompute_engine/engine.rs +++ b/data_plane/src/precompute_engine/engine.rs @@ -108,7 +108,7 @@ impl PrecomputeEngine { pub async fn run(mut self) -> Result<(), Box> { let num_workers = self.config.num_workers; let output_sink: Arc = Arc::new( - crate::precompute_engine::maintenance_runtime::MaintenanceDagSink::new( + crate::precompute_engine::maintenance_runtime::PrecomputeDagSink::new( Arc::clone(&self.output_sink), self.hot_reload_config.clone(), ), diff --git a/data_plane/src/precompute_engine/maintenance_runtime.rs b/data_plane/src/precompute_engine/maintenance_runtime.rs index bbed2b039..e2410c2be 100644 --- a/data_plane/src/precompute_engine/maintenance_runtime.rs +++ b/data_plane/src/precompute_engine/maintenance_runtime.rs @@ -1781,14 +1781,14 @@ impl IdempotentCommitSink for CommitRegistry { /// Decorates the ordinary store sink with installed maintenance DAG execution. /// With no matching DAG, the source output is forwarded unchanged. -pub struct MaintenanceDagSink { +pub struct PrecomputeDagSink { inner: Arc, plans: StreamingConfigHandle, commits: CommitRegistry, batch_guard: Mutex<()>, } -impl MaintenanceDagSink { +impl PrecomputeDagSink { pub fn new(inner: Arc, plans: StreamingConfigHandle) -> Self { Self { inner, @@ -2000,7 +2000,7 @@ fn dependencies_until( seen } -impl OutputSink for MaintenanceDagSink { +impl OutputSink for PrecomputeDagSink { fn emit_batch( &self, outputs: Vec<(PrecomputedOutput, Box)>, @@ -2568,7 +2568,7 @@ mod tests { ) .unwrap(), ); - document.schema_version = asap_types::executable_plan::MAINTENANCE_DAG_SCHEMA_VERSION; + document.schema_version = asap_types::executable_plan::PRECOMPUTE_DAG_SCHEMA_VERSION; let mut durable_binding = scheduled_binding.clone(); durable_binding.nodes.insert( PostAsapNodeId(3), @@ -2946,7 +2946,7 @@ mod tests { ) .unwrap(), ); - document.schema_version = asap_types::executable_plan::MAINTENANCE_DAG_SCHEMA_VERSION; + document.schema_version = asap_types::executable_plan::PRECOMPUTE_DAG_SCHEMA_VERSION; let mut binding = binding.clone(); for (node, stored_output) in [ (1, first_id), @@ -3903,7 +3903,7 @@ mod tests { let mut document = OwnedPostAsapDag::from_executable("retry".into(), &dag).unwrap(); document.schema_version = - asap_types::executable_plan::MAINTENANCE_DAG_SCHEMA_VERSION; + asap_types::executable_plan::PRECOMPUTE_DAG_SCHEMA_VERSION; document }, binding, @@ -3927,7 +3927,7 @@ mod tests { fail_at, ..Default::default() }); - let sink = MaintenanceDagSink::new( + let sink = PrecomputeDagSink::new( downstream.clone(), StreamingConfigHandle::from_active_physical_plan(ActivePhysicalPlanHandle::new( active.clone(), diff --git a/data_plane/src/precompute_engine/subdag_scheduler.rs b/data_plane/src/precompute_engine/subdag_scheduler.rs index 660ef8906..e64f98973 100644 --- a/data_plane/src/precompute_engine/subdag_scheduler.rs +++ b/data_plane/src/precompute_engine/subdag_scheduler.rs @@ -75,7 +75,7 @@ where ))); } binding - .validate_maintenance(dag) + .validate_precompute(dag) .map_err(ScheduleError::Invalid)?; if !binding.precompute_sinks.contains(&sink_node) || !matches!( diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index ae9db38fe..baceebff1 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -4447,7 +4447,7 @@ mod dag_execution_tests { ) .unwrap(); installed.document.schema_version = - asap_types::executable_plan::MAINTENANCE_DAG_SCHEMA_VERSION; + asap_types::executable_plan::PRECOMPUTE_DAG_SCHEMA_VERSION; let error = StreamingConfig::from_precompute_plan(plan) .unwrap_err() .to_string(); 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 05d1e8c26..073cbf5ee 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -390,7 +390,7 @@ impl ASAPQueryEngine { |root, evaluation_ms| { if let Some(asap_types::query_plan::QueryPlanNode::Logical { operator: - asap_types::query_plan::residual::ResidualQueryOperator::CurrentSeries { + asap_types::query_plan::query_time::QueryTimeOperator::CurrentSeries { population, readout, }, 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 aed8dd8b9..eb5d0221b 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 @@ -2,7 +2,7 @@ use super::logical_dag::{PreparedLeaf, PreparedLeaves, Value}; use crate::query_engines::EngineError; use asap_types::query_plan::{ - residual::ResidualQueryOperator, ExternalExactInput, ExternalExactRequest, QueryLanguage, + query_time::QueryTimeOperator, ExternalExactInput, ExternalExactRequest, QueryLanguage, QueryNodeId, QueryPlanEntry, QueryPlanNode, }; use std::collections::{BTreeMap, BTreeSet, HashMap}; @@ -17,7 +17,7 @@ fn miss(message: impl Into) -> EngineError { /// Traverse only the installed graph, including epoch-aligned nested subquery grids. #[derive(Debug, Clone)] enum ExactLeaf { - Legacy(ResidualQueryOperator), + Legacy(QueryTimeOperator), External(ExternalExactRequest), } @@ -47,14 +47,14 @@ fn leaves( .ok_or_else(|| miss("missing installed node"))?; match node { QueryPlanNode::Logical { operator, inputs } => match operator { - ResidualQueryOperator::Scan { .. } => { + QueryTimeOperator::Scan { .. } => { return Err(miss("local raw Scan is forbidden in deployed plans")) } - ResidualQueryOperator::ExactSubquery { .. } - | ResidualQueryOperator::CandidateExactSubquery { .. } => { + QueryTimeOperator::ExactSubquery { .. } + | QueryTimeOperator::CandidateExactSubquery { .. } => { result.insert((id, at), ExactLeaf::Legacy(operator.clone())); } - ResidualQueryOperator::Subquery { + QueryTimeOperator::Subquery { range_ms, step_ms, offset_ms, @@ -119,7 +119,7 @@ pub(super) fn external_dependencies( ); } else if matches!( leaf, - ExactLeaf::Legacy(ResidualQueryOperator::CandidateExactSubquery { .. }) + ExactLeaf::Legacy(QueryTimeOperator::CandidateExactSubquery { .. }) ) { let input = *entry.nodes[&id] .inputs() @@ -301,13 +301,10 @@ pub(super) async fn prepare_external( 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(ResidualQueryOperator::ExactSubquery { query }) => { + ExactLeaf::Legacy(QueryTimeOperator::ExactSubquery { query }) => { (QueryLanguage::PromQl, query.clone(), None) } - ExactLeaf::Legacy(ResidualQueryOperator::CandidateExactSubquery { - query, - item_label, - }) => ( + ExactLeaf::Legacy(QueryTimeOperator::CandidateExactSubquery { query, item_label }) => ( QueryLanguage::PromQl, query.clone(), Some((entry.nodes[&id].inputs()[0], item_label.as_str())), @@ -687,10 +684,10 @@ mod tests { entry.nodes.insert( QueryNodeId(3), QueryPlanNode::Logical { - operator: ResidualQueryOperator::Limit { + operator: QueryTimeOperator::Limit { offset: 0, n: 2, - grouping: asap_types::query_plan::residual::Grouping { + grouping: asap_types::query_plan::query_time::Grouping { labels: vec![], without: false, }, @@ -805,7 +802,7 @@ mod tests { #[tokio::test] async fn exact_leaf_calls_prometheus_and_combines_with_prepared_summary() { // A successful exact branch remains an intermediate, not a whole-root fallback. - use asap_types::query_plan::residual::BinaryOperation; + use asap_types::query_plan::query_time::BinaryOperation; use std::sync::{ atomic::{AtomicUsize, Ordering}, Arc, @@ -825,7 +822,7 @@ mod tests { ( QueryNodeId(0), QueryPlanNode::Logical { - operator: ResidualQueryOperator::Binary { + operator: QueryTimeOperator::Binary { operation: BinaryOperation::Div, return_bool: false, }, @@ -839,7 +836,7 @@ mod tests { ( QueryNodeId(2), QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { query: "b".into() }, + operator: QueryTimeOperator::ExactSubquery { query: "b".into() }, inputs: vec![], }, ), @@ -881,7 +878,7 @@ mod tests { repeated.nodes.insert( QueryNodeId(1), QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { query: "b".into() }, + operator: QueryTimeOperator::ExactSubquery { query: "b".into() }, inputs: vec![], }, ); @@ -921,7 +918,7 @@ mod tests { use crate::storage_engines::types::{KeyByLabelValues, Measurement}; use asap_physical_operators::summary_kernels::IncreaseAccumulator; use asap_types::query_plan::{ - residual::BinaryOperation, ExactReadout, MaterializationBinding, PhysicalGrouping, + query_time::BinaryOperation, ExactReadout, MaterializationBinding, PhysicalGrouping, }; use std::sync::{ atomic::{AtomicUsize, Ordering}, @@ -963,7 +960,7 @@ mod tests { ( QueryNodeId(0), QueryPlanNode::Logical { - operator: ResidualQueryOperator::Binary { + operator: QueryTimeOperator::Binary { operation: BinaryOperation::Div, return_bool: false, }, @@ -973,7 +970,7 @@ mod tests { ( QueryNodeId(1), QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { + operator: QueryTimeOperator::ExactSubquery { query: exact_query.into(), }, inputs: vec![], @@ -1093,7 +1090,7 @@ mod tests { let entry = entry(BTreeMap::from([( QueryNodeId(0), QueryPlanNode::Logical { - operator: ResidualQueryOperator::Scan { + operator: QueryTimeOperator::Scan { metric: Some("m".into()), matchers: vec![], range_ms: None, 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 b14fe9d17..1dea68693 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 @@ -5,8 +5,8 @@ use crate::query_engines::{ EngineError, }; use crate::storage_engines::types::KeyByLabelValues; -use asap_types::query_plan::residual::{ - Aggregation, BinaryOperation, Grouping, ResidualQueryOperator, TemporalOperation, +use asap_types::query_plan::query_time::{ + Aggregation, BinaryOperation, Grouping, QueryTimeOperator, TemporalOperation, }; use asap_types::query_plan::{CandidateCompleteness, QueryNodeId, QueryPlanEntry, QueryPlanNode}; use std::collections::{BTreeMap, BTreeSet}; @@ -233,7 +233,7 @@ impl Result> Evaluator<' QueryPlanNode::Scalar { value } => Value::Scalar(value), QueryPlanNode::Logical { - operator: ResidualQueryOperator::CurrentSeries { .. }, + operator: QueryTimeOperator::CurrentSeries { .. }, .. } => { self.stats.summary_readout_evaluations += 1; @@ -245,9 +245,9 @@ impl Result> Evaluator<' QueryPlanNode::Logical { operator, inputs } => { if matches!( operator, - ResidualQueryOperator::Scan { .. } - | ResidualQueryOperator::ExactSubquery { .. } - | ResidualQueryOperator::CandidateExactSubquery { .. } + QueryTimeOperator::Scan { .. } + | QueryTimeOperator::ExactSubquery { .. } + | QueryTimeOperator::CandidateExactSubquery { .. } ) { return Err(miss( "installed Prometheus leaf was not prepared; backend raw execution is forbidden", @@ -297,7 +297,7 @@ impl Result> Evaluator<' } fn logical( &mut self, - operator: ResidualQueryOperator, + operator: QueryTimeOperator, inputs: &[QueryNodeId], at: i64, ) -> Result { @@ -308,17 +308,17 @@ impl Result> Evaluator<' .ok_or_else(|| miss("missing logical input")) }; match operator { - ResidualQueryOperator::ExactSubquery { .. } - | ResidualQueryOperator::CandidateExactSubquery { .. } => { + QueryTimeOperator::ExactSubquery { .. } + | QueryTimeOperator::CandidateExactSubquery { .. } => { Err(miss("Prometheus exact leaf was not prepared")) } - ResidualQueryOperator::CurrentSeries { .. } => Err(miss( + QueryTimeOperator::CurrentSeries { .. } => Err(miss( "current-series leaf must use its installed node identity", )), - ResidualQueryOperator::Scan { .. } => { + QueryTimeOperator::Scan { .. } => { Err(miss("local raw Scan is forbidden in deployed plans")) } - ResidualQueryOperator::UnaryNegate => match self.eval(input(0)?, at)? { + QueryTimeOperator::UnaryNegate => match self.eval(input(0)?, at)? { Value::Scalar(value) => Ok(Value::Scalar(-value)), Value::Vector(values) => Ok(Value::Vector( values @@ -328,7 +328,7 @@ impl Result> Evaluator<' )), _ => Err(miss("cannot negate range vector")), }, - ResidualQueryOperator::VectorToScalar => { + QueryTimeOperator::VectorToScalar => { let values = vector(self.eval(input(0)?, at)?)?; Ok(Value::Scalar(if values.len() == 1 { values[0].1 @@ -336,14 +336,14 @@ impl Result> Evaluator<' f64::NAN })) } - ResidualQueryOperator::Aggregate { + QueryTimeOperator::Aggregate { operation, grouping, } => { let values = vector(self.eval(input(0)?, at)?)?; Ok(Value::Vector(aggregate(operation, &grouping, values))) } - ResidualQueryOperator::Limit { + QueryTimeOperator::Limit { n, offset, grouping, @@ -357,7 +357,7 @@ impl Result> Evaluator<' self.context.clone(), )?)) } - ResidualQueryOperator::Binary { + QueryTimeOperator::Binary { operation, return_bool, } => { @@ -365,7 +365,7 @@ impl Result> Evaluator<' let right = self.eval(input(1)?, at)?; binary(operation, return_bool, left, right) } - ResidualQueryOperator::Temporal { operation } => { + QueryTimeOperator::Temporal { operation } => { let Value::Matrix(values, start, end) = self.eval(input(0)?, at)? else { return Err(miss("temporal operator requires range vector")); }; @@ -426,7 +426,7 @@ impl Result> Evaluator<' .collect(), )) } - ResidualQueryOperator::Sort { + QueryTimeOperator::Sort { descending, grouping, } => { @@ -438,7 +438,7 @@ impl Result> Evaluator<' self.context.clone(), )?)) } - ResidualQueryOperator::HistogramQuantile => { + QueryTimeOperator::HistogramQuantile => { let Value::Scalar(quantile) = self.eval(input(0)?, at)? else { return Err(miss("quantile requires scalar")); }; @@ -455,7 +455,7 @@ impl Result> Evaluator<' .collect(), )) } - ResidualQueryOperator::Subquery { + QueryTimeOperator::Subquery { range_ms, step_ms, offset_ms, @@ -1028,7 +1028,7 @@ mod topk_tests { #[test] fn installed_topk_combines_with_prometheus_exact_child() { - let mut entry = control_plane::query_plan::residual::compile_logical( + let mut entry = control_plane::query_plan::query_time::compile_logical( "hybrid-topk".into(), "topk(2, m)".into(), InstantExecution { @@ -1039,7 +1039,7 @@ mod topk_tests { FallbackPolicy::ExactBackend, ) .unwrap(); - control_plane::query_plan::residual::finalize_residuals(&mut entry).unwrap(); + control_plane::query_plan::query_time::finalize_query_time_nodes(&mut entry).unwrap(); let leaf = entry .nodes .iter() @@ -1047,7 +1047,7 @@ mod topk_tests { matches!( node, QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { .. }, + operator: QueryTimeOperator::ExactSubquery { .. }, .. } ) @@ -1114,7 +1114,7 @@ mod topk_tests { ( QueryNodeId(0), QueryPlanNode::Logical { - operator: ResidualQueryOperator::ExactSubquery { + operator: QueryTimeOperator::ExactSubquery { query: "m[1s]".into(), }, inputs: vec![], @@ -1123,7 +1123,7 @@ mod topk_tests { ( QueryNodeId(1), QueryPlanNode::Logical { - operator: ResidualQueryOperator::Temporal { operation }, + operator: QueryTimeOperator::Temporal { operation }, inputs: vec![QueryNodeId(0)], }, ), @@ -1198,7 +1198,7 @@ mod topk_tests { ( root, QueryPlanNode::Logical { - operator: ResidualQueryOperator::Limit { + operator: QueryTimeOperator::Limit { offset: 0, n: 2, grouping: Grouping { @@ -1212,7 +1212,7 @@ mod topk_tests { ( QueryNodeId(98), QueryPlanNode::Logical { - operator: ResidualQueryOperator::Sort { + operator: QueryTimeOperator::Sort { descending: true, grouping: Grouping { labels: vec![], @@ -1431,7 +1431,7 @@ mod topk_tests { ( root, QueryPlanNode::Logical { - operator: ResidualQueryOperator::Limit { + operator: QueryTimeOperator::Limit { offset: 0, n: 1, grouping: Grouping { @@ -1445,7 +1445,7 @@ mod topk_tests { ( QueryNodeId(98), QueryPlanNode::Logical { - operator: ResidualQueryOperator::Sort { + operator: QueryTimeOperator::Sort { descending: true, grouping: Grouping { labels: vec![], diff --git a/data_plane/src/storage_engines/sketch_db/current_series.rs b/data_plane/src/storage_engines/sketch_db/current_series.rs index efd3c0127..2c58bece6 100644 --- a/data_plane/src/storage_engines/sketch_db/current_series.rs +++ b/data_plane/src/storage_engines/sketch_db/current_series.rs @@ -2,7 +2,7 @@ use crate::drivers::ingest::prometheus_remote_write::CanonicalSample; use asap_types::query_plan::{ current_series::{SeriesPopulation, SeriesReadout}, - residual::{LabelMatch, ResidualQueryOperator}, + query_time::{LabelMatch, QueryTimeOperator}, QueryPlan, QueryPlanNode, }; use std::collections::{BTreeMap, BTreeSet}; @@ -221,7 +221,7 @@ impl Population { match readout { SeriesReadout::Quantile { q } => { // `values` is only populated for a quantile-carrying population. - // `ResidualQueryOperator::validate` rejects the mismatched pairing at + // `QueryTimeOperator::validate` rejects the mismatched pairing at // install, so this is defensive: answer like Prometheus does for // an empty group rather than underflow `values.len() - 1` while // holding the lock every remote-write batch waits on. @@ -286,7 +286,7 @@ impl CurrentSeriesStore { for entry in plan.entries.values() { for node in entry.nodes.values() { if let QueryPlanNode::Logical { - operator: ResidualQueryOperator::CurrentSeries { population, .. }, + operator: QueryTimeOperator::CurrentSeries { population, .. }, .. } = node { @@ -389,7 +389,7 @@ mod tests { SeriesPopulation { metric: "a".into(), matchers: vec![], - grouping: asap_types::query_plan::residual::Grouping { + grouping: asap_types::query_plan::query_time::Grouping { labels: vec!["job".into()], without: false, }, @@ -438,7 +438,7 @@ mod tests { nodes: BTreeMap::from([( QueryNodeId(0), QueryPlanNode::Logical { - operator: ResidualQueryOperator::CurrentSeries { + operator: QueryTimeOperator::CurrentSeries { population: p.clone(), readout: SeriesReadout::Quantile { q: 0.5 }, }, diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index 99fae4b0a..e3dfa40da 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -1007,7 +1007,7 @@ async fn run_shared_dashboard(multi_pane: bool) { .executable_dags .values() .all(|installed| installed.document.schema_version - == asap_types::executable_plan::MAINTENANCE_DAG_SCHEMA_VERSION)); + == asap_types::executable_plan::PRECOMPUTE_DAG_SCHEMA_VERSION)); if multi_pane { assert!(plan.lifecycle_estimates[0] .window_realization_id diff --git a/data_plane/tests/support/current_series_process.rs b/data_plane/tests/support/current_series_process.rs index 4f0dc6992..86bb390b5 100644 --- a/data_plane/tests/support/current_series_process.rs +++ b/data_plane/tests/support/current_series_process.rs @@ -105,7 +105,7 @@ async fn current_series_quantiles_topk_share_and_replace_values() { for node in entry.nodes.values() { if let asap_types::query_plan::QueryPlanNode::Logical { operator: - asap_types::query_plan::residual::ResidualQueryOperator::CurrentSeries { + asap_types::query_plan::query_time::QueryTimeOperator::CurrentSeries { population, .. }, diff --git a/data_plane/tests/support/issue_701_702_process.rs b/data_plane/tests/support/issue_701_702_process.rs index f2dbe6b6c..533b0d4c4 100644 --- a/data_plane/tests/support/issue_701_702_process.rs +++ b/data_plane/tests/support/issue_701_702_process.rs @@ -142,8 +142,8 @@ async fn run_warm_workload(queries: Vec<(String, u64, u64)>) { control_plane::query_plan::QueryPlanNode::ExactFallback { .. } | control_plane::query_plan::QueryPlanNode::ExternalExact { .. } | control_plane::query_plan::QueryPlanNode::Logical { - operator: control_plane::query_plan::residual::ResidualQueryOperator::ExactSubquery { .. } - | control_plane::query_plan::residual::ResidualQueryOperator::CandidateExactSubquery { .. }, .. + operator: control_plane::query_plan::query_time::QueryTimeOperator::ExactSubquery { .. } + | control_plane::query_plan::query_time::QueryTimeOperator::CandidateExactSubquery { .. }, .. } ))) }; @@ -400,8 +400,8 @@ async fn temporal_average_overflow_falls_back_after_state_is_warm() { .any(|node| matches!( node, control_plane::query_plan::QueryPlanNode::Logical { - operator: control_plane::query_plan::residual::ResidualQueryOperator::Binary { - operation: control_plane::query_plan::residual::BinaryOperation::FiniteDiv, + operator: control_plane::query_plan::query_time::QueryTimeOperator::Binary { + operation: control_plane::query_plan::query_time::BinaryOperation::FiniteDiv, .. }, ..