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
259 changes: 106 additions & 153 deletions control_plane/src/physical/compiler.rs

Large diffs are not rendered by default.

206 changes: 205 additions & 1 deletion control_plane/src/physical/compiler/placement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
use super::*;
use asap_aware_mapping::{enumerate_summary_maintenance_lifecycles, CostModel};
use asap_physical_operators::physical_planner::{
compile, promql_fallback, CompiledPhysicalDag, InputContract,
compile, promql_fallback, promql_rows, CompiledPhysicalDag, InputContract, PhysicalCandidate,
};
use planner_types::post_asap::PostAsapOperatorPayload;
use planner_types::pre_asap::{AggIntent, CompareOpKind, ScalarValue};
Expand Down Expand Up @@ -418,6 +418,210 @@ fn range_selector_scan(selector: &QueryExpr) -> Result<QueryTimeOperator, String
})
}

/// Planner's fixed-window realization: a heap or grouped Sum over complete
/// per-series Rate states maintained at ingestion, per the DAG's timing.
pub(super) fn fixed_window_candidate(
dag: &planner_types::post_asap::PostAsapDag,
) -> Result<PhysicalCandidate, asap_physical_operators::Error> {
promql_rows::compile_fixed_window_rate_aggregation(dag)
}

/// [`fixed_window_candidate`] under the timing written into `root`.
pub(super) fn root_fixed_window_candidate(
root: &Rc<SummaryNode>,
) -> Result<PhysicalCandidate, asap_physical_operators::Error> {
let dag = planner_types::post_asap::compile_post_asap_dag(root)
.map_err(|e| asap_physical_operators::Error::Invalid(e.to_string()))?;
fixed_window_candidate(&dag)
}

/// Planner's query-time realization above a maintained population or exact
/// per-series Rate readouts; the input is bound by the backend.
pub(super) fn query_time_candidate(root: &Rc<SummaryNode>) -> Option<PhysicalCandidate> {
promql_rows::compile_current_series_readout(root)
.or_else(|_| promql_rows::compile_rate_ranking(root).map(|(_, dag)| dag))
.ok()
.map(|query| PhysicalCandidate {
precompute: None,
query,
materialized_outputs: BTreeMap::new(),
})
}

/// A summary built from another retained state's finalized readouts, such as
/// a heap or grouped Sum over per-series Rate.
fn over_readouts(summary: &SummaryNode) -> bool {
matches!(&summary.expr, SummaryExpr::SummaryAgg { child, .. }
if matches!(&child.expr, SummaryExpr::ValueOperation {
operation: planner_types::post_asap::ValueOperation::FinalizeExactAccumulator,
child,
..
} if matches!(child.expr, SummaryExpr::SummaryAgg { .. })))
}

/// Write `timing` onto the readouts feeding `target`, rebuilding its path.
fn retime(
node: &Rc<SummaryNode>,
target: &Rc<SummaryNode>,
timing: planner_types::post_asap::ExecutionTiming,
) -> Option<Rc<SummaryNode>> {
let mut next = node.as_ref().clone();
if Rc::ptr_eq(node, target) {
let SummaryExpr::SummaryAgg { child, .. } = &mut next.expr else {
return None;
};
let mut readout = child.as_ref().clone();
let SummaryExpr::ValueOperation { timing: old, .. } = &mut readout.expr else {
return None;
};
*old = timing;
*child = Rc::new(readout);
return Some(Rc::new(next));
}
match &mut next.expr {
SummaryExpr::ValueOperation { child, .. } | SummaryExpr::SummaryAgg { child, .. } => {
*child = retime(child, target, timing)?
}
SummaryExpr::SummaryEstimate { summary_input, .. } => {
*summary_input = retime(summary_input, target, timing)?
}
_ => return None,
}
Some(Rc::new(next))
}

/// A Planner candidate with a native physical realization after its
/// readout-built states are placed by lifecycle.
pub(super) struct TimedCandidate {
pub(super) root: Rc<SummaryNode>,
pub(super) physical: PhysicalCandidate,
pub(super) trace: Vec<Value>,
}

/// Place each state of `root` built from retained readouts: maintained, it
/// consumes complete per-series states at ingestion; ephemeral, it is rebuilt
/// for each query from the retained state's readouts. Planner's lifecycle plan
/// times the DAG and its native compiler reads that timing. Other states stay
/// retained here; compilation places them. `None` means no native realization.
pub(super) fn time_native_candidate(
root: &Rc<SummaryNode>,
query: &QueryCompilationInput,
workload: &QueryWorkload,
data: &DataWorkload,
index: usize,
environment: &PhysicalDeploymentContext,
) -> Option<TimedCandidate> {
use planner_types::post_asap::ExecutionTiming;
let untimed = || {
query_time_candidate(root).map(|physical| TimedCandidate {
root: Rc::clone(root),
physical,
trace: Vec::new(),
})
};
let lifecycle = &query.summary_lifecycle_inputs;
let model = LifecycleCosts {
costs: &lifecycle.costs,
evaluation_interval_ms: lifecycle.evaluation_interval_ms,
input_cardinality: data
.input_cardinality
.value_at(environment.observed_at_unix_ms)
.copied(),
delete: environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite,
};
let enumerate = || {
enumerate_summary_maintenance_lifecycles(
Rc::clone(root),
WorkloadDemand::new_with_data(workload, data, std::slice::from_ref(&index)),
environment.observed_at_unix_ms,
Some(Horizon(lifecycle.horizon_seconds)),
SummaryMaintenanceLifecycleCapabilities {
supports_ephemeral: true,
supports_prepared: false,
supports_shared: false,
supports_continuously_maintained: true,
},
&model,
)
.ok()
};
let candidates = enumerate()?;
let mut choices = Vec::new();
let mut retimed = Rc::clone(root);
let mut trace = Vec::new();
let mut maintained_over_readouts = false;
for deployment in candidates.deployments() {
let retained = alternative_cost(
deployment,
&SummaryMaintenanceLifecycle::ContinuouslyMaintained,
);
let lifecycle = if over_readouts(&deployment.summary) {
let rebuilt = alternative_cost(deployment, &SummaryMaintenanceLifecycle::Ephemeral);
let ephemeral = rebuild_is_cheaper(retained, rebuilt);
maintained_over_readouts |= !ephemeral;
let timing = if ephemeral {
ExecutionTiming::QueryTime
} else {
ExecutionTiming::IngestionTime
};
retimed = retime(&retimed, &find(&retimed, &deployment.summary)?, timing)?;
// Recorded per candidate forest; compilation records the final
// placement of the states it installs as `lifecycle_placement`.
trace.push(json!({
"stage": "deployment.native_candidate_placement",
"query_ids": [&query.query_id],
"candidate_root_id": crate::planner_selection::explained_root_id(root, &query.accuracy_target),
"logical_root_id": crate::planner_selection::explained_root_id(&deployment.summary, &query.accuracy_target),
"ephemeral_bindable": true,
"continuously_maintained_cost": retained.map(|cost| cost.0),
"ephemeral_cost": rebuilt.map(|cost| cost.0),
"selected": if ephemeral { "ephemeral" } else { "continuously_maintained" },
}));
if ephemeral {
SummaryMaintenanceLifecycle::Ephemeral
} else {
SummaryMaintenanceLifecycle::ContinuouslyMaintained
}
} else {
SummaryMaintenanceLifecycle::ContinuouslyMaintained
};
choices.push((deployment.post_asap_node_id, lifecycle));
}
if trace.is_empty() {
return untimed();
}
let timed = candidates
.select(&choices)
.ok()?
.execution_timed_dag()
.ok()?;
let physical = if maintained_over_readouts {
fixed_window_candidate(&timed).ok()?
} else {
query_time_candidate(&retimed)?
};
Some(TimedCandidate {
root: retimed,
physical,
trace,
})
}

/// The node of `root` structurally equal to `summary`, found after earlier
/// rewrites replaced the original `Rc`s.
fn find(root: &Rc<SummaryNode>, summary: &SummaryNode) -> Option<Rc<SummaryNode>> {
if root.as_ref() == summary {
return Some(Rc::clone(root));
}
match &root.expr {
SummaryExpr::ValueOperation { child, .. } | SummaryExpr::SummaryAgg { child, .. } => {
find(child, summary)
}
SummaryExpr::SummaryEstimate { summary_input, .. } => find(summary_input, summary),
_ => None,
}
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
Loading
Loading