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
12 changes: 6 additions & 6 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

10 changes: 5 additions & 5 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,10 @@ version = "0.1.0"
[workspace.dependencies]
# Keep Planner frontends, selection, and IR on the same immutable revision.
# Alias upstream asap-types because this workspace also defines asap_types.
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }
asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }
asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }

# Shared external deps (used by 2+ crates)
serde = { version = "1.0", features = ["derive"] }
Expand All @@ -39,7 +39,7 @@ arc-swap = "1.7"
reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] }

# Internal crates
asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" }
asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }
asap_sketch_codec = { path = "crates/asap_sketch_codec" }
asap_summary_state = { path = "crates/asap_summary_state" }
asap_types = { path = "crates/asap_types" }
Expand Down
6 changes: 2 additions & 4 deletions control_plane/examples/calibration_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,6 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
.enumerate()
{
let queries = candidate.queries.clone();
let enabled_materialization_keys = candidate.enabled_materialization_keys.clone();
let planner_selected_queries = planner_forest(&queries);
let compiled = if metricsql {
DeploymentPlanCompiler.compile_metricsql(candidate, environment.clone())
Expand All @@ -137,7 +136,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
Ok(plan) => plan,
Err(error) => {
results.push(
json!({"candidate_index": index, "materialization_policy": enabled_materialization_keys, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}),
json!({"candidate_index": index, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}),
);
continue;
}
Expand All @@ -146,14 +145,13 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
Ok(manifest) => manifest,
Err(error) => {
results.push(
json!({"candidate_index": index, "materialization_policy": enabled_materialization_keys, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}),
json!({"candidate_index": index, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}),
);
continue;
}
};
results.push(json!({
"candidate_index": index,
"materialization_policy": enabled_materialization_keys,
"planner_selected_queries": planner_selected_queries,
"manifest": manifest,
"lifecycle_estimates": plan.lifecycle_estimates,
Expand Down
1 change: 0 additions & 1 deletion control_plane/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -666,7 +666,6 @@ fn compile_physical_plan_request(
allow_mixed_summary_and_exact_execution: request.target
== physical::compiler::PhysicalDeploymentTarget::BackendLocalRemoteWrite,
require_backend_local_execution: false,
enabled_materialization_keys: None,
topk_membership_evidence_by_query_id: request.evidence,
exact_composition_costs: request.exact_composition_costs,
erp: request.erp,
Expand Down
129 changes: 99 additions & 30 deletions control_plane/src/physical/compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ use crate::query_plan::{
use crate::types::AccuracyTarget;
use planner_types::pre_asap::Source;

mod placement;
mod rate_placement;
mod windows;
pub(super) use windows::gcd;
Expand Down Expand Up @@ -184,6 +185,15 @@ pub struct LifecycleUnitCosts {
pub read: f64,
pub retention_per_second: f64,
pub retirement: f64,
/// Summary-store price per retained byte-second. Retaining a state costs
/// its estimated bytes times its retained panes times this price, on top
/// of `retention_per_second`.
#[serde(default, skip_serializing_if = "f64_is_zero")]
pub store_per_byte_second: f64,
}

fn f64_is_zero(value: &f64) -> bool {
*value == 0.0
}

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
Expand All @@ -208,9 +218,6 @@ pub struct PhysicalCompilationRequest {
pub allow_mixed_summary_and_exact_execution: bool,
/// Deployment feasibility: external exact dependencies cannot be bound.
pub require_backend_local_execution: bool,
/// Enabled optional candidate keys: None enables all eligible keys; an
/// explicitly empty set enables none. These are not catalog definition IDs.
pub enabled_materialization_keys: Option<BTreeSet<String>>,
/// Original dashboard demand, in the same order as queries. None is legacy input.
pub query_workload: Option<QueryWorkload>,
/// Source evidence supplied independently from query demand.
Expand Down Expand Up @@ -1149,7 +1156,6 @@ impl BackendLocalPlanningInput {
require_backend_local_execution: self
.physical_inputs
.require_backend_local_execution,
enabled_materialization_keys: None,
query_workload: Some(workload),
data_workload: Some(data_workload),
canonical_roots,
Expand Down Expand Up @@ -1471,6 +1477,7 @@ impl DeploymentPlanCompiler {
for (id, root) in planner_types::post_asap::share_common_summary_subtrees(roots) {
request.queries[id].selected_plan_root = root;
}
let placement = placement::place(&request, &environment, frontend);
let population_operators = super::maintained_population::operators(&request)?;
let mut compiled_materializations = Vec::with_capacity(request.queries.len());
let mut collector_materializations = Vec::with_capacity(request.queries.len());
Expand Down Expand Up @@ -1555,25 +1562,8 @@ impl DeploymentPlanCompiler {
.ok()
.flatten()
.is_some())
&& request.enabled_materialization_keys.as_ref().is_none_or(|policy| {
let key = crate::query_plan::query_time::selected_counter_materialization(
&query.query_string,
&state.node,
)
.ok()
.flatten()
.or_else(|| {
crate::query_plan::query_time::selected_range_max_materialization(
&query.query_string,
&state.node,
)
.ok()
.flatten()
});
// Masks enumerate counter/max choices only. Other selected
// summaries remain required by this physical candidate.
key.is_none_or(|key| policy.contains(&key))
})
// An Ephemeral state is rebuilt at query time, not maintained.
&& !placement.is_ephemeral(query_index, &state.node)
})
.collect::<Vec<_>>();
// An exact native fallback has no maintained state and must not
Expand Down Expand Up @@ -2064,7 +2054,6 @@ impl DeploymentPlanCompiler {
stable_workload_plan_id(&plan_materializations, &request.queries).hash(&mut hash);
"typed-local-residual-v3-counter-index".hash(&mut hash);
super::maintained_population::supported(&request).hash(&mut hash);
request.enabled_materialization_keys.hash(&mut hash);
for query in &request.queries {
format!("{:?}", query.selected_plan_root).hash(&mut hash);
}
Expand Down Expand Up @@ -2224,7 +2213,9 @@ impl DeploymentPlanCompiler {
} else {
false
};
let native_rate = if population_operators[query_index].is_none() {
let native_rate = if population_operators[query_index].is_none()
&& placement.raw_program(query_index).is_none()
{
query
.retained_physical()?
.map(|physical| {
Expand Down Expand Up @@ -2257,7 +2248,9 @@ impl DeploymentPlanCompiler {
} else {
None
};
let mut entry = if let Some((source, program)) = native_rate {
let mut entry = if let Some(raw) = placement.raw_program(query_index) {
Ok(raw_query_time_entry(query, canonical.clone(), raw)?)
} else if let Some((source, program)) = native_rate {
let native_state_binding = if let SummaryExpr::SummaryAgg {
family:
SummaryFamilyType::Sketch(..)
Expand Down Expand Up @@ -2438,7 +2431,9 @@ impl DeploymentPlanCompiler {
},
)
}?;
if request.allow_mixed_summary_and_exact_execution {
if request.allow_mixed_summary_and_exact_execution
&& placement.raw_program(query_index).is_none()
{
// Any Planner-selected leaf without a physical summary binding
// is an exact subtree boundary. Deployed plans never retain a
// backend-local range index leaf.
Expand Down Expand Up @@ -2701,7 +2696,8 @@ impl DeploymentPlanCompiler {
crate::query_plan::QueryPlanNode::ExactFallback { .. }
| crate::query_plan::QueryPlanNode::Logical {
operator: crate::query_plan::query_time::QueryTimeOperator::ExactSubquery { .. }
| crate::query_plan::query_time::QueryTimeOperator::CandidateExactSubquery { .. }, .. })) {
| crate::query_plan::query_time::QueryTimeOperator::CandidateExactSubquery { .. }
| crate::query_plan::query_time::QueryTimeOperator::Scan { .. }, .. })) {
return Err(CompileError::Query {
query_id: entry.query_id.clone(),
reason: "external execution is unavailable in this deployment".into(),
Expand All @@ -2724,11 +2720,74 @@ impl DeploymentPlanCompiler {
storage_routing,
lifecycle_estimates: lifecycle_estimates.into_values().collect(),
cost_comparison: None,
planner_selection_trace: request.planner_selection_trace,
planner_selection_trace: if placement.trace.is_empty() {
request.planner_selection_trace
} else {
let mut trace = (*request.planner_selection_trace).clone();
trace.extend(placement.trace);
trace.into()
},
})
}
}

/// Native query-time execution of a query that keeps no state: each raw input
/// of the retained program is read by its range-selector `Scan`.
fn raw_query_time_entry(
query: &QueryCompilationInput,
canonical: String,
raw: &placement::RawQueryTimeProgram,
) -> Result<QueryPlanEntry, CompileError> {
let mut nodes = BTreeMap::new();
let mut inputs = Vec::new();
let mut source_nodes = Vec::new();
for (ordinal, (slot, scan)) in raw.scans.iter().enumerate() {
let id = crate::query_plan::QueryNodeId(ordinal as u64);
nodes.insert(
id,
crate::query_plan::QueryPlanNode::Logical {
operator: scan.clone(),
inputs: vec![],
},
);
inputs.push(id);
source_nodes.push(*slot);
}
let root = crate::query_plan::QueryNodeId(raw.scans.len() as u64);
nodes.insert(
root,
crate::query_plan::QueryPlanNode::Physical {
inputs,
source_nodes,
max_bytes: 64 * 1024 * 1024,
},
);
let entry = QueryPlanEntry {
physical_dag: Some(
serde_json::from_slice(
&raw.program
.encode()
.map_err(|e| CompileError::Snapshot(e.to_string()))?,
)
.map_err(|e| CompileError::Snapshot(e.to_string()))?,
),
language: crate::query_plan::QueryLanguage::PromQl,
query_id: query.query_id.clone(),
canonical_query: canonical,
fixed_evaluation: None,
root,
nodes,
instant: InstantExecution {
lookback_ms: query.query_lookback_ms,
full_history: false,
cumulative_readout: false,
},
fallback: FallbackPolicy::ExactBackend,
};
entry.recover_vector_physical_dag()?;
Ok(entry)
}

fn summary_agg_metric(node: &SummaryNode) -> Option<String> {
fn walk(node: &SummaryNode, metrics: &mut BTreeSet<String>) {
match &node.expr {
Expand Down Expand Up @@ -5644,6 +5703,7 @@ pub(crate) mod tests {
read: 0.1,
retention_per_second: 0.001,
retirement: 1.0,
store_per_byte_second: 0.0,
},
};
let post_asap = select_post_asap(&parsed, accuracy.clone(), &lifecycle, evidence.as_ref())?;
Expand All @@ -5657,7 +5717,6 @@ pub(crate) mod tests {
canonical_roots: Vec::new(),
allow_mixed_summary_and_exact_execution: false,
require_backend_local_execution: false,
enabled_materialization_keys: None,
query_workload: None,
data_workload: None,
queries: vec![QueryCompilationInput {
Expand Down Expand Up @@ -7089,6 +7148,7 @@ pub(crate) mod tests {
read: 0.0,
retention_per_second: 1.0,
retirement: 0.0,
store_per_byte_second: 0.0,
};
let full = derived_window_cost(
&template,
Expand Down Expand Up @@ -7686,6 +7746,9 @@ pub(crate) mod tests {
continue;
}
let mut snapshot = planning_snapshot();
// Window layouts are properties of retained state; with no query-time
// raw source, sparse or costly maintenance cannot move it to query time.
snapshot.physical_inputs.require_backend_local_execution = true;
let entry = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0];
entry.query = Query("sum_over_time(a[1m])".into());
entry.requirements.accuracy =
Expand Down Expand Up @@ -7731,6 +7794,9 @@ pub(crate) mod tests {
fn different_cadences_share_common_panes_only_when_cheaper() {
for read_cost in [0.0, 1_000_000.0] {
let mut snapshot = planning_snapshot();
// Window layouts are properties of retained state; with no query-time
// raw source, sparse or costly maintenance cannot move it to query time.
snapshot.physical_inputs.require_backend_local_execution = true;
snapshot.physical_inputs.lifecycle_costs.read = read_cost;
snapshot
.physical_inputs
Expand Down Expand Up @@ -7780,6 +7846,9 @@ pub(crate) mod tests {
fn sharing_keeps_profitable_subsets_with_other_phases_or_finer_cadences() {
for (third_phase, third_interval) in [(5_000, 20_000), (0, 1_000)] {
let mut snapshot = planning_snapshot();
// Window layouts are properties of retained state; with no query-time
// raw source, sparse or costly maintenance cannot move it to query time.
snapshot.physical_inputs.require_backend_local_execution = true;
snapshot
.physical_inputs
.lifecycle_costs
Expand Down
Loading
Loading