Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions control_plane/src/clickhouse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)?,
);
}
Expand Down Expand Up @@ -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 =
Expand Down
63 changes: 36 additions & 27 deletions control_plane/src/physical/compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<QueryExpr>],
Expand Down Expand Up @@ -1133,22 +1133,22 @@ 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,
)
.ok()
.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,
)
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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 })
})
Expand Down Expand Up @@ -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<u64>, String), String> {
Expand Down Expand Up @@ -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,
..
},
Expand Down Expand Up @@ -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,
..
},
..
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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, .. },
..
}
));
Expand All @@ -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
Expand Down Expand Up @@ -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 { .. },
..
}
))
Expand Down Expand Up @@ -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 { .. },
..
}
))));
Expand Down Expand Up @@ -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"
))
Expand All @@ -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
Expand Down Expand Up @@ -6823,7 +6832,7 @@ pub(crate) mod tests {
assert!(!query.nodes.values().any(|node| matches!(
node,
QueryPlanNode::Logical {
operator: ResidualQueryOperator::Scan { .. },
operator: QueryTimeOperator::Scan { .. },
..
}
)));
Expand Down
6 changes: 3 additions & 3 deletions control_plane/src/physical/maintained_population.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -43,7 +43,7 @@ pub(super) fn supported(request: &PhysicalCompilationRequest) -> bool {
pub(super) fn operator(
request: &PhysicalCompilationRequest,
query: &QueryCompilationInput,
) -> Result<Option<ResidualQueryOperator>, CompileError> {
) -> Result<Option<QueryTimeOperator>, CompileError> {
let Some((spec, readout)) = selected(&query.selected_plan_root) else {
return Ok(None);
};
Expand Down Expand Up @@ -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,
}))
Expand Down
34 changes: 18 additions & 16 deletions control_plane/src/physical/plan_dot.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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(),
Expand All @@ -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",
}
}

Expand Down
10 changes: 5 additions & 5 deletions control_plane/src/physical/workload_cost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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, ..
},
..
Expand All @@ -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 { .. },
..
}
) {
Expand All @@ -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,
..
},
Expand Down Expand Up @@ -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,
) {
Expand Down
Loading
Loading