Skip to content
268 changes: 238 additions & 30 deletions crates/asap-physical-operators/src/physical_planner/candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,24 +25,42 @@ pub fn compile_candidate(
inputs: BTreeMap<NodeId, InputContract>,
roots: &[NodeId],
frontier: &[NodeId],
) -> Result<PhysicalCandidate, Error> {
cut_candidate(&compile(dag, inputs, roots)?, frontier)
}

/// Derive one frontier's candidate from a complete [`compile`] result by
/// partitioning its operators; nothing is lowered again. A deployment compiles
/// each query DAG once and derives every placement choice from that result.
/// The candidate is identical to [`compile_candidate`] for the same frontier.
pub fn cut_candidate(
compiled: &CompiledPhysicalDag,
frontier: &[NodeId],
) -> Result<PhysicalCandidate, Error> {
if frontier.is_empty() {
return Ok(PhysicalCandidate {
precompute: None,
query: compile(dag, inputs, roots)?,
query: compiled.clone(),
materialized_outputs: BTreeMap::new(),
});
}
let frontier_set: BTreeSet<_> = frontier.iter().copied().collect();
if frontier_set.len() != frontier.len() || frontier.iter().any(|id| inputs.contains_key(id)) {
// `compile` retains only reachable nodes and numbers its helper operators
// above the u32 Planner ID range; only Planner outputs are boundaries.
if frontier_set.len() != frontier.len()
|| frontier
.iter()
.any(|&id| !compiled.is_operator(id) || u32::try_from(id).is_err())
{
return Err(invalid("frontier must contain distinct computed outputs"));
}
let full = compile(dag, inputs.clone(), roots)?;
let precompute = compile(dag, inputs.clone(), frontier)?;
let inputs: BTreeMap<_, _> = compiled
.input_contracts()
.map(|(id, contract)| (id, contract.clone()))
.collect();
let precompute = compiled.cut(&inputs, frontier)?;
let mut materialized_outputs = BTreeMap::new();
for &id in frontier {
// Also proves that the frontier is reachable from the requested roots.
full.output_contract(id)?;
let mut output = precompute.output_contract(id)?;
if output.properties.boundedness != Boundedness::Bounded {
return Err(invalid("materialized output requires bounded execution"));
Expand All @@ -54,7 +72,7 @@ pub fn compile_candidate(
}
let mut query_inputs = inputs;
query_inputs.extend(materialized_outputs.clone());
let query = compile(dag, query_inputs, roots)?;
let query = compiled.cut(&query_inputs, compiled.roots())?;
let used: BTreeSet<_> = query.input_contracts().map(|(id, _)| id).collect();
if !frontier.iter().all(|id| used.contains(id)) {
return Err(invalid(
Expand All @@ -68,6 +86,41 @@ pub fn compile_candidate(
})
}

/// Materialization frontier implied by lifecycle-assigned timing: ingestion-time
/// nodes read by a query-time node, plus the root when it is ingestion-timed.
/// `cut_candidate` of one [`compile`] result with this frontier realizes the
/// assignment, so different assignments are different cuts of one lowering.
/// That holds while timing-dependent lowering (an ingestion-time `Binary`
/// aligns by value column) has the same timing at compile time as here.
/// A query-time node feeding an ingestion-time node has no valid placement.
pub fn frontier_from_timing(dag: &PostAsapDag) -> Result<Vec<NodeId>, Error> {
use planner_types::post_asap::ExecutionTiming::IngestionTime;
let timing = dag
.nodes
.iter()
.map(|node| (node.id, node.output_state.timing))
.collect::<BTreeMap<_, _>>();
let mut frontier = BTreeSet::new();
if timing.get(&dag.root) == Some(&IngestionTime) {
frontier.insert(u64::from(dag.root.0));
}
for edge in &dag.edges {
let (Some(&producer), Some(&consumer)) =
(timing.get(&edge.producer), timing.get(&edge.consumer))
else {
return Err(invalid("timed DAG edge names an unknown node"));
};
match (producer == IngestionTime, consumer == IngestionTime) {
(true, false) => {
frontier.insert(u64::from(edge.producer.0));
}
(false, true) => return Err(invalid("query-time node feeds an ingestion-time node")),
_ => {}
}
}
Ok(frontier.into_iter().collect())
}

/// Enumerate bounded, reachable materialization frontiers above explicit inputs.
/// Each frontier is an antichain: storing an output and its ancestor together
/// would leave the ancestor unused by query execution. Lifecycle eligibility
Expand All @@ -78,43 +131,38 @@ pub fn enumerate_frontiers(
inputs: &BTreeMap<NodeId, InputContract>,
roots: &[NodeId],
max_candidates: usize,
) -> Result<Vec<Vec<NodeId>>, Error> {
enumerate_compiled_frontiers(&compile(dag, inputs.clone(), roots)?, max_candidates)
}

fn enumerate_compiled_frontiers(
compiled: &CompiledPhysicalDag,
max_candidates: usize,
) -> Result<Vec<Vec<NodeId>>, Error> {
if max_candidates == 0 {
return Err(invalid(
"frontier search requires a positive candidate budget",
));
}
let compiled = compile(dag, inputs.clone(), roots)?;
let mut ancestors = BTreeMap::<NodeId, BTreeSet<NodeId>>::new();
let mut eligible = Vec::new();
for node in &dag.nodes {
let id = u64::from(node.id.0);
if inputs.contains_key(&id) {
continue;
}
let Ok(contract) = compiled.output_contract(id) else {
continue;
};
if contract.properties.boundedness != Boundedness::Bounded {
for (id, properties) in compiled.output_properties()? {
if !compiled.is_operator(id)
|| u32::try_from(id).is_err()
|| properties.boundedness != Boundedness::Bounded
{
continue;
}
let mut seen = BTreeSet::new();
let mut pending = vec![id];
while let Some(current) = pending.pop() {
if !seen.insert(current) || inputs.contains_key(&current) {
continue;
if seen.insert(current) {
pending.extend(compiled.dependencies(current));
}
pending.extend(
dag.edges
.iter()
.filter(|edge| u64::from(edge.consumer.0) == current)
.map(|edge| u64::from(edge.producer.0)),
);
}
ancestors.insert(id, seen);
eligible.push(id);
}
eligible.sort_unstable();
let mut frontiers = vec![vec![]];
for id in eligible {
let additions = frontiers
Expand Down Expand Up @@ -142,16 +190,20 @@ pub fn enumerate_frontiers(

/// Lower every maintenance candidate before feasibility/cost evaluation. Keep
/// individual failures visible; do not substitute another computation on error.
/// The DAG is lowered once; each frontier is a [`cut_candidate`] of it.
pub fn compile_candidates(
dag: &PostAsapDag,
inputs: BTreeMap<NodeId, InputContract>,
roots: &[NodeId],
frontiers: &[Vec<NodeId>],
) -> Vec<Result<PhysicalCandidate, Error>> {
frontiers
.iter()
.map(|frontier| compile_candidate(dag, inputs.clone(), roots, frontier))
.collect()
match compile(dag, inputs, roots) {
Ok(compiled) => frontiers
.iter()
.map(|frontier| cut_candidate(&compiled, frontier))
.collect(),
Err(error) => frontiers.iter().map(|_| Err(error.clone())).collect(),
}
}

/// Complete workload cost supplied by scoped optimizer/deployment evidence.
Expand Down Expand Up @@ -273,3 +325,159 @@ impl PhysicalCandidate {
Ok(())
}
}

#[cfg(test)]
mod tests {
use super::*;
use planner_types::workload::*;

fn grouped_rate() -> (PostAsapDag, BTreeMap<NodeId, InputContract>, NodeId) {
let workload = PlanningWorkload {
query_workload: QueryWorkload {
language: QueryLanguage::PromQL,
query_batch: Some(vec![BatchEntry {
query: Query("sum by(job)(rate(m[1m]))".into()),
requirements: QueryRequirements {
accuracy: AccuracyRequirement::Explicit(
planner_types::types::AccuracyTarget::Exact,
),
..Default::default()
},
predictability: Predictability::Unknown,
invocations: 1,
execute_at: None,
time_selection: TimeSelection::default(),
}]),
repeating_queries: None,
},
data_workload: Some(DataWorkload {
data_ingestion_interval: Evidence {
value: Some(DurationMs(1000)),
..Default::default()
},
..Default::default()
}),
};
let root = asap_frontend_promql::lower_promql_workload(&workload, 0)
.unwrap()
.remove(0);
let root = std::rc::Rc::new(promql_rows::with_series_identity(&root).unwrap());
let space = asap_aware_mapping::search_workload(vec![("q", root)]);
let selected = space
.global_selection(&asap_aware_mapping::cost_model::DefaultCostModel)
.assemble_selected_dag(&space.roots[0].1)
.unwrap()
.unwrap();
let dag = planner_types::post_asap::compile_post_asap_dag(&selected).unwrap();
let state = dag
.nodes
.iter()
.find(|node| matches!(node.payload, Payload::SummaryAgg { .. }))
.unwrap();
let inputs = BTreeMap::from([(
u64::from(state.id.0),
InputContract::bounded(Arc::new(state.output_schema.clone())),
)]);
(dag.clone(), inputs, u64::from(dag.root.0))
}

/// Enumerating and cutting every frontier lowers each Planner node once.
#[test]
fn candidates_for_all_frontiers_share_one_lowering() {
let (dag, inputs, root) = grouped_rate();
let lowered = || crate::physical_planner::LOWERED_NODES.with(|count| count.get());
let before = lowered();
let compiled = compile(&dag, inputs, &[root]).unwrap();
let once = lowered() - before;
let frontiers = enumerate_compiled_frontiers(&compiled, 4096).unwrap();
assert!(frontiers.len() >= 3, "{frontiers:?}");
for frontier in &frontiers {
cut_candidate(&compiled, frontier).unwrap();
}
assert!(once > 0);
assert_eq!(lowered() - before, once);
}

fn with_timing(
dag: &PostAsapDag,
timing: impl Fn(&PostAsapDagNode) -> planner_types::post_asap::ExecutionTiming,
) -> PostAsapDag {
let mut timed = dag.clone();
for node in &mut timed.nodes {
node.output_state.timing = timing(node);
}
for edge in &mut timed.edges {
let producer = timed.nodes.iter().find(|node| node.id == edge.producer);
edge.data_state = producer.unwrap().output_state;
}
timed
}

fn raw_input(dag: &PostAsapDag) -> BTreeMap<NodeId, InputContract> {
let raw = dag
.nodes
.iter()
.find(|node| matches!(node.payload, Payload::Fallback { .. }))
.unwrap();
BTreeMap::from([(
u64::from(raw.id.0),
InputContract::bounded(Arc::new(raw.output_schema.clone())),
)])
}

/// Cutting one compilation by a retained-state timing and by the all
/// query-time timing (what ContinuouslyMaintained and Ephemeral assign)
/// lowers each Planner node once and matches `compile_candidate`.
#[test]
fn timing_cuts_share_one_lowering() {
use planner_types::post_asap::ExecutionTiming::QueryTime;
let (retained, _, root) = grouped_rate();
let ephemeral = with_timing(&retained, |_| QueryTime);
let inputs = raw_input(&retained);
let lowered = || crate::physical_planner::LOWERED_NODES.with(|count| count.get());
let before = lowered();
let compiled = compile(&ephemeral, inputs.clone(), &[root]).unwrap();
let once = lowered() - before;
let cuts = [&retained, &ephemeral].map(|timed| {
let frontier = frontier_from_timing(timed).unwrap();
let cut = cut_candidate(&compiled, &frontier).unwrap();
(timed, frontier, cut)
});
assert!(once > 0);
assert_eq!(lowered() - before, once);
assert_eq!(cuts[0].1.len(), 1);
assert!(cuts[1].1.is_empty());
for (timed, frontier, cut) in cuts {
let expected = compile_candidate(timed, inputs.clone(), &[root], &frontier).unwrap();
assert_eq!(
serde_json::to_vec(&cut).unwrap(),
serde_json::to_vec(&expected).unwrap()
);
}
}

/// The frontier is the ingestion-time nodes read at query time; an
/// ingestion-time root is itself the frontier.
#[test]
fn frontier_from_timing_includes_ingestion_root() {
use planner_types::post_asap::ExecutionTiming::IngestionTime;
let (dag, _, root) = grouped_rate();
let timed = with_timing(&dag, |_| IngestionTime);
assert_eq!(frontier_from_timing(&timed).unwrap(), [root]);
}

/// A query-time node feeding an ingestion-time node is rejected.
#[test]
fn frontier_from_timing_rejects_query_time_input_to_ingestion() {
use planner_types::post_asap::ExecutionTiming::{IngestionTime, QueryTime};
let (dag, _, _) = grouped_rate();
let timed = with_timing(&dag, |node| {
if node.id == dag.root {
IngestionTime
} else {
QueryTime
}
});
assert!(frontier_from_timing(&timed).is_err());
}
}
Loading
Loading