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
7 changes: 6 additions & 1 deletion control_plane/src/physical/compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2436,7 +2436,7 @@ impl DeploymentPlanCompiler {
if frontend == QueryFrontend::MetricsQl {
entry.language = crate::query_plan::QueryLanguage::MetricsQl;
}
super::maintained_population::install_native_topk(
super::maintained_population::install_population_readout(
&mut entry,
query
.retained_physical()?
Expand Down Expand Up @@ -4867,6 +4867,11 @@ pub(crate) mod tests {
}
}
assert_eq!(populations.len(), 1);
// Storage returns the members; every readout is the entry's Planner program.
for entry in plan.query_plan.entries.values() {
assert!(entry.population_snapshot().is_some(), "{entry:?}");
entry.recover_population_physical_dag().unwrap();
}
let installed = serde_json::to_string(&plan.precompute_plan.executable_dags).unwrap();
assert!(
!installed.contains("MaintainPopulation"),
Expand Down
71 changes: 24 additions & 47 deletions control_plane/src/physical/maintained_population.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,17 +4,17 @@ use super::compiler::QueryCompilationInput;
use super::compiler::{CompileError, PhysicalCompilationRequest};
use asap_types::physical_plan_codec::PhysicalPlanCodec;
use asap_types::query_plan::{
current_series::{SeriesPopulation, SeriesReadout},
current_series::SeriesPopulation,
query_time::{Grouping, LabelMatch, LabelMatcher, QueryTimeOperator},
};
use planner_types::post_asap::{
maintained_population::*, SummaryExpr, SummaryNode, ValueOperation,
};

fn selected(node: &SummaryNode) -> Option<(MaintainedPopulation, PopulationReadout)> {
fn selected(node: &SummaryNode) -> Option<MaintainedPopulation> {
if let SummaryExpr::ValueOperation {
child,
operation: ValueOperation::ReadPopulation { readout },
operation: ValueOperation::ReadPopulation { .. },
..
} = &node.expr
{
Expand All @@ -23,13 +23,13 @@ fn selected(node: &SummaryNode) -> Option<(MaintainedPopulation, PopulationReado
..
} = &child.expr
{
return Some((population.clone(), readout.clone()));
return Some(population.clone());
}
}
// The source remains a maintained population when Planner places a heap,
// projection and ranking above it. Backend binds that source only.
let SummaryExpr::ValueOperation {
operation: ValueOperation::Limit { n, offset: 0, .. },
operation: ValueOperation::Limit { offset: 0, .. },
..
} = &node.expr
else {
Expand All @@ -54,12 +54,14 @@ fn selected(node: &SummaryNode) -> Option<(MaintainedPopulation, PopulationReado
&std::rc::Rc::new(node.clone()),
)
.ok()?;
Some(((*population).clone(), PopulationReadout::TopK { k: *n }))
Some((*population).clone())
}

/// A current-series population Planner can read out. Planner cannot yet
/// project `without` groups from the series identity.
pub(super) fn supported_node(node: &SummaryNode) -> bool {
selected(node).is_some_and(|(population, _)| {
matches!(population.input, PopulationInput::CurrentSeries(_))
selected(node).is_some_and(|population| {
matches!(&population.input, PopulationInput::CurrentSeries(spec) if !spec.without)
})
}

Expand All @@ -83,17 +85,15 @@ pub(super) fn operators(
let populations: std::collections::BTreeSet<_> = selected
.iter()
.flatten()
.map(|(population, _)| {
serde_json::to_string(population).expect("typed population serializes")
})
.map(|population| serde_json::to_string(population).expect("typed population serializes"))
.collect();
let max_bytes = request
.retained_summary_memory_budget_bytes
.unwrap_or(64 * 1024 * 1024)
.min(1_073_741_824)
/ populations.len().max(1) as u64;
selected.into_iter().zip(&request.queries).map(|(selected, query)| {
let Some((spec, readout)) = selected else { return Ok(None); };
let Some(spec) = selected else { return Ok(None); };
let PopulationInput::CurrentSeries(input) = &spec.input else {
return Err(CompileError::Query { query_id: query.query_id.clone(), reason: "maintained table-row populations require a row-update executor; remote-write current-series state is incompatible".into() });
};
Expand Down Expand Up @@ -131,17 +131,7 @@ pub(super) fn operators(
.min(input.lookback_ms),
};
population.validate()?;
let readout = match &readout {
PopulationReadout::Quantile { q } => SeriesReadout::Quantile { q: *q },
PopulationReadout::TopK { k } => SeriesReadout::TopK { k: *k as u64 },
PopulationReadout::Sum => SeriesReadout::Sum,
PopulationReadout::Count => SeriesReadout::Count,
PopulationReadout::Average => SeriesReadout::Average,
};
Ok(Some(QueryTimeOperator::CurrentSeries {
population,
readout,
}))
Ok(Some(QueryTimeOperator::CurrentSeries { population }))
}).collect()
}

Expand All @@ -158,30 +148,25 @@ pub(super) fn operator(
Ok(operators(request)?.remove(index))
}

/// The maintained population is a deployment source; ranking is compiled by
/// Planner before this candidate is priced or installed.
pub(super) fn install_native_topk(
/// The maintained population is a deployment source; its readout is the
/// Planner program compiled before this candidate is priced or installed.
pub(super) fn install_population_readout(
entry: &mut asap_types::query_plan::QueryPlanEntry,
compiled: Option<&asap_physical_operators::physical_planner::CompiledPhysicalDag>,
) -> Result<(), CompileError> {
use asap_types::query_plan::QueryPlanNode;
let Some(QueryPlanNode::Logical {
operator:
QueryTimeOperator::CurrentSeries {
population,
readout: SeriesReadout::TopK { .. },
},
..
}) = entry.nodes.get(&entry.root)
else {
return Ok(());
};
if population.grouping.without {
if !matches!(
entry.nodes.get(&entry.root),
Some(QueryPlanNode::Logical {
operator: QueryTimeOperator::CurrentSeries { .. },
..
})
) {
return Ok(());
}
let compiled = compiled.ok_or_else(|| CompileError::Query {
query_id: entry.query_id.clone(),
reason: "selected TopK candidate has no retained physical DAG".into(),
reason: "selected population readout has no retained physical DAG".into(),
})?;
let encoded = compiled.encode().map_err(|error| CompileError::Query {
query_id: entry.query_id.clone(),
Expand All @@ -191,14 +176,6 @@ pub(super) fn install_native_topk(
serde_json::from_slice(&encoded)
.map_err(|error| CompileError::Snapshot(error.to_string()))?,
);
let Some(QueryPlanNode::Logical {
operator: QueryTimeOperator::CurrentSeries { readout, .. },
..
}) = entry.nodes.get_mut(&entry.root)
else {
unreachable!()
};
*readout = SeriesReadout::Snapshot;
entry.recover_population_physical_dag()?;
Ok(())
}
88 changes: 68 additions & 20 deletions control_plane/src/physical/workload_cost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -834,14 +834,30 @@ fn enumerate_frontier_candidates(
_ => unreachable!("native candidate retains canonical roots"),
})
.collect();
let strategy =
asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new(&roots);
let maintained_roots: Vec<_> = roots
// Populations are selected over roots typed with the complete series
// identity, so Planner can compile every readout over the snapshot rows.
let typed = roots
.iter()
.map(|root| {
strategy
.candidate(root)
.filter(|node| super::maintained_population::supported_node(node))
asap_physical_operators::physical_planner::promql_rows::with_series_identity(root)
.map(std::rc::Rc::new)
.ok()
})
.collect::<Vec<_>>();
let strategy = asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new(
&typed.iter().flatten().cloned().collect::<Vec<_>>(),
);
let maintained_roots: Vec<_> = typed
.iter()
.map(|root| {
root.as_ref()
.and_then(|root| strategy.candidate(root))
.filter(|node| {
// Every readout must compile to the Planner program installed with it.
super::maintained_population::supported_node(node)
&& asap_physical_operators::physical_planner::promql_rows::compile_current_series_readout(node)
.is_ok()
})
})
.collect();
if maintained_roots.iter().any(Option::is_some) {
Expand Down Expand Up @@ -1037,6 +1053,37 @@ mod tests {
})));
}

/// Planner cannot yet read out `without` groups, so such a population is
/// never selected, while the `by` form is.
#[test]
fn without_grouping_selects_no_population() {
for (query, supported) in [
("quantile without (pod) (0.5, m)", false),
("quantile by (job) (0.5, m)", true),
] {
let root = std::rc::Rc::new(
asap_physical_operators::physical_planner::promql_rows::with_series_identity(
&crate::query_parser::parse_query_expr_canonical(
query,
crate::types::AccuracyTarget::Exact,
)
.unwrap(),
)
.unwrap(),
);
let strategy =
asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new(
std::slice::from_ref(&root),
);
let candidate = strategy.candidate(&root).unwrap();
assert_eq!(
super::super::maintained_population::supported_node(&candidate),
supported,
"{query}"
);
}
}

/// Instant counts select current membership, never accumulated observations.
#[test]
fn local_grouped_count_has_a_bindable_candidate() {
Expand All @@ -1050,20 +1097,21 @@ mod tests {
queries.truncate(1);
queries[0].query = planner_types::workload::Query("count by(job)(m)".into());
let plan = with_unit_quotes(input).compile_promql().unwrap();
assert!(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 {
readout: asap_types::query_plan::current_series::SeriesReadout::Count,
..
},
..
}
))));
// The population supplies members; Planner counts them per job.
for entry in plan.query_plan.entries.values() {
assert!(entry.population_snapshot().is_some(), "{entry:?}");
let program = entry.recover_population_physical_dag().unwrap();
let output = program.output_contract(program.roots()[0]).unwrap();
assert_eq!(
output
.schema
.fields
.iter()
.map(|field| field.name.as_str())
.collect::<Vec<_>>(),
["job", "count"]
);
}
}

fn fixture() -> BackendLocalPlanningInput {
Expand Down
16 changes: 0 additions & 16 deletions crates/asap_types/src/query_plan/current_series.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,22 +46,6 @@ impl SeriesPopulation {
}
}

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
pub enum SeriesReadout {
/// All eligible members for a Planner-compiled physical readout.
Snapshot,
Quantile {
q: f64,
},
TopK {
k: u64,
},
Sum,
Count,
Average,
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
34 changes: 26 additions & 8 deletions crates/asap_types/src/query_plan/native.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,11 +11,7 @@ impl QueryPlanEntry {
}
match self.nodes.get(&self.root) {
Some(QueryPlanNode::Logical {
operator:
query_time::QueryTimeOperator::CurrentSeries {
population,
readout: current_series::SeriesReadout::Snapshot,
},
operator: query_time::QueryTimeOperator::CurrentSeries { population },
inputs,
}) if inputs.is_empty() => Some(population),
_ => None,
Expand Down Expand Up @@ -49,12 +45,34 @@ impl QueryPlanEntry {
let output = dag
.output_contract(*root)
.map_err(|error| QueryPlanError::Invalid(error.to_string()))?;
if input.schema != output.schema {
use planner_types::{post_asap::SummaryFamilyType, pre_asap::DataType};
// A ranking returns complete source rows; an aggregate returns one
// value per group of the population's grouping labels.
let numeric = |field: &planner_types::post_asap::SummaryField| {
matches!(
field.dtype,
SummaryFamilyType::Plain(DataType::Float64 | DataType::Int64)
)
};
let labels = output
.schema
.fields
.iter()
.filter(|field| !numeric(field))
.map(|field| {
(field.dtype == SummaryFamilyType::Plain(DataType::Utf8)).then_some(&field.name)
})
.collect::<Option<BTreeSet<_>>>();
let grouped_value = output.schema.time_index.is_none()
&& output.schema.fields.iter().filter(|f| numeric(f)).count() == 1
&& labels.is_some_and(|labels| {
labels == population.grouping.labels.iter().collect::<BTreeSet<_>>()
});
if input.schema != output.schema && !grouped_value {
return Err(invalid(
"population ranking must preserve complete source rows",
"population readout must return source rows or one value per group",
));
}
use planner_types::{post_asap::SummaryFamilyType, pre_asap::DataType};
let fields = &input.schema.fields;
let column = |name: &str, dtype: DataType| {
fields.iter().any(|field| {
Expand Down
23 changes: 3 additions & 20 deletions crates/asap_types/src/query_plan/query_time.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,10 @@ fn invalid(message: impl Into<String>) -> QueryPlanError {
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
pub enum QueryTimeOperator {
/// Readout over a bounded current-value population maintained at ingest.
/// Every member of a bounded current-value population maintained at
/// ingest; the entry's Planner program computes the readout.
CurrentSeries {
population: super::current_series::SeriesPopulation,
readout: super::current_series::SeriesReadout,
},
/// A maximal exact scalar/vector subtree evaluated by Prometheus.
ExactSubquery { query: String },
Expand Down Expand Up @@ -51,25 +51,8 @@ pub enum LabelMatch {
}
impl QueryTimeOperator {
pub fn validate(&self, inputs: usize) -> Result<(), QueryPlanError> {
if let Self::CurrentSeries {
population,
readout,
} = self
{
if let Self::CurrentSeries { population } = self {
population.validate()?;
match readout {
super::current_series::SeriesReadout::Quantile { q }
if !q.is_finite() || !population.quantiles =>
{
return Err(invalid(
"quantile readout requires finite q and a quantile population",
))
}
super::current_series::SeriesReadout::TopK { k } if *k > population.max_k => {
return Err(invalid("TopK readout exceeds shared population capacity"))
}
_ => {}
}
}
let expected = match self {
Self::Scan { .. } | Self::ExactSubquery { .. } | Self::CurrentSeries { .. } => 0,
Expand Down
Loading
Loading