Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
33fad08
docs: specify physical candidate handoff and migration acceptance
zzylol Sep 29, 2026
8d089c1
refactor: install compiled SQL relation graphs before query execution
zzylol Sep 29, 2026
ae362bd
Execute native TopK candidates over bound stored outputs
zzylol Sep 29, 2026
93ad0fe
Retain selected native candidates and execute revisioned query ensembles
zzylol Sep 29, 2026
83f7a29
Remove uncompiled relational plan formats and execute native aggregates
zzylol Sep 29, 2026
4f5ce15
fix: validate physical row identity and preserve missing grouping labels
zzylol Sep 29, 2026
1a3efb6
refactor: install Planner-compiled scalar and vector computation
zzylol Sep 29, 2026
c049d6e
refactor: execute connected native value graphs without runtime lowering
zzylol Sep 29, 2026
372bc39
refactor: execute retained native graphs for immutable precomputation
zzylol Sep 29, 2026
cdd23e3
refactor: remove backend precompute interpreter and scheduler
zzylol Sep 29, 2026
ce5d8d2
refactor: execute retained query value graphs and reject stale nodes
zzylol Sep 29, 2026
9ba7c31
fix: preserve shared precompute graphs across stored outputs
zzylol Sep 29, 2026
9280048
fix: publish sliding counter revisions at the selected cadence
zzylol Sep 29, 2026
8afaf03
test: retain summary state through shared physical projections
zzylol Sep 29, 2026
e8960b1
refactor: adopt precompute/query_time vocabulary from #783
zzylol Sep 29, 2026
57da03f
refactor: adopt Planner PostAsapDag names
zzylol Sep 29, 2026
b719fb3
chore(compiler): assert lifecycle plan lists the selected state first
zzylol Sep 29, 2026
6c49478
perf(workload_cost): reuse per-root work across deployment candidates
zzylol Sep 29, 2026
9efc252
perf(asap_types): decode each installed DAG once per validation
zzylol Sep 29, 2026
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
14 changes: 7 additions & 7 deletions Cargo.lock

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

12 changes: 6 additions & 6 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -14,10 +14,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 = "4b0839ce1733aab8231eb90944c876063a4551f9" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "4b0839ce1733aab8231eb90944c876063a4551f9" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "4b0839ce1733aab8231eb90944c876063a4551f9" }
asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "4b0839ce1733aab8231eb90944c876063a4551f9" }
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }

# Shared external deps (used by 2+ crates)
serde = { version = "1.0", features = ["derive"] }
Expand All @@ -37,8 +37,8 @@ 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 = "4b0839ce1733aab8231eb90944c876063a4551f9" }
asap_sketch_codec = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "4b0839ce1733aab8231eb90944c876063a4551f9" }
asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap_sketch_codec = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "76fbbf16cc44b19f56a780bfdb47a95327e84711" }
asap_types = { path = "crates/asap_types" }
asap_otel_proto = { path = "crates/asap_otel_proto" }
indexmap = { version = "2.0", features = ["serde"] }
2 changes: 1 addition & 1 deletion control_plane/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ path = "src/main.rs"
tokio = { version = "1", features = ["full"] }
axum = { version = "0.7", features = ["ws"] }
futures-util = "0.3"
serde = { version = "1", features = ["derive"] }
serde = { version = "1", features = ["derive", "rc"] }
serde_json = "1"
sha2 = "0.10"
serde_yaml = "0.9"
Expand Down
2 changes: 1 addition & 1 deletion control_plane/src/clickhouse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -534,7 +534,7 @@ where
"SQL did not produce a summary DAG".into(),
));
};
let semantic = planner_types::post_asap::compile_executable_dag_with_node_ids(&root)
let semantic = planner_types::post_asap::compile_post_asap_dag_with_node_ids(&root)
.map_err(|error| ClickHousePlanningError::Lower(error.to_string()))?;
let mut materialization_nodes = std::collections::BTreeMap::new();
let mut query_nodes = std::collections::BTreeMap::new();
Expand Down
75 changes: 60 additions & 15 deletions control_plane/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -204,8 +204,12 @@ struct CompileAndPublishPhysicalPlanRequest {
target_collector_ids: Vec<String>,
capability_snapshot_id: String,
#[serde(default)]
data_snapshot_id: Option<String>,
#[serde(default)]
evidence: HashMap<String, physical::compiler::TopKMembershipEvidence>,
#[serde(default)]
accuracy_evidence: HashMap<String, physical::compiler::ScopedAccuracyEvidence>,
#[serde(default)]
exact_composition_costs:
HashMap<String, Vec<physical::post_asap::cost_model::ExactCompositionCostEvidence>>,
#[serde(default)]
Expand Down Expand Up @@ -378,7 +382,7 @@ async fn compile_and_publish_physical_plan(

Json(CompileAndPublishPhysicalPlanResponse {
cost_comparison: bundle.cost_comparison,
planner_selection_trace: bundle.planner_selection_trace,
planner_selection_trace: bundle.planner_selection_trace.as_ref().clone(),
plan_id: bundle.envelope.plan_id,
plan_version: bundle.envelope.plan_version,
status: "active",
Expand Down Expand Up @@ -586,10 +590,11 @@ fn compile_physical_plan_request(
query_id: query.query_id,
query_string: query.query_string,
selected_plan_root: post_asap,
physical_candidate: None,
legacy_query_source: planner_types::pre_asap::Source::TimeSeries {
metric: query.metric,
},
query_lookback_seconds: query.window_secs,
query_lookback_ms: query.window_secs.saturating_mul(1_000),
group_by_labels: query.group_by,
accuracy_target: query.accuracy,
summary_lifecycle_inputs: query.lifecycle,
Expand All @@ -598,29 +603,69 @@ fn compile_physical_plan_request(
});
}

let planner_selection_trace = match physical::compiler::select_logical_roots_with_trace(
&mut queries,
canonical_roots.clone(),
&request.evidence,
&request.exact_composition_costs,
request.erp.as_ref(),
) {
Ok(trace) => trace,
Err(error) => return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into())),
};
let scoped_snapshot_id = request.data_snapshot_id.as_deref().or_else(|| {
request
.workload_cost_evidence
.as_ref()
.map(|evidence| evidence.data_snapshot_id.as_str())
});
if request.data_snapshot_id.as_ref().is_some_and(|id| {
request
.workload_cost_evidence
.as_ref()
.is_some_and(|evidence| evidence.data_snapshot_id != *id)
}) {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"accuracy evidence data snapshot differs from workload cost evidence".into(),
));
}
for (query_id, evidence) in &request.accuracy_evidence {
let Some(query) = queries.iter().find(|query| &query.query_id == query_id) else {
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
format!("accuracy evidence names unknown query {query_id}").into(),
));
};
evidence
.validate(
query_id,
&query.query_string,
&request.data_workload,
scoped_snapshot_id,
now,
request.max_evidence_age_ms,
)
.map_err(|error| (StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into()))?;
}
let planner_selection_trace =
match physical::compiler::select_logical_roots_with_scoped_evidence_and_trace(
&mut queries,
canonical_roots.clone(),
&request.evidence,
&request.accuracy_evidence,
&request.exact_composition_costs,
request.erp.as_ref(),
now,
) {
Ok(trace) => trace,
Err(error) => return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into())),
};

for (query, model) in queries.iter_mut().zip(window_models) {
physical::compiler::prepare_window_implementations(query, &model, request.target, 0)
.map_err(|error| (StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into()))?;
}
let compilation_request = physical::compiler::PhysicalCompilationRequest {
planner_selection_trace,
planner_candidate_forests: Vec::new(),
planner_selection_trace: planner_selection_trace.into(),
query_workload: Some(query_workload),
data_workload: Some(request.data_workload),
canonical_roots,
queries,
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,
Expand All @@ -646,7 +691,7 @@ fn compile_physical_plan_request(
compilation_request.clone(),
)
.map_err(|error| (StatusCode::UNPROCESSABLE_ENTITY, error.to_string().into()))?;
let planner_selection_trace = compilation_request.planner_selection_trace.clone();
let planner_selection_trace = compilation_request.planner_selection_trace.as_ref().clone();
let (manifests, candidate_evaluations) =
physical::workload_cost::compile_candidates_for_pricing(
candidates.clone(),
Expand Down Expand Up @@ -929,7 +974,7 @@ mod api_tests {
"target": "backend_local_remote_write",
"queries": [{
"query_id": query.query_id, "query_string": query.query_string,
"metric": metric, "window_secs": query.query_lookback_seconds, "accuracy": query.accuracy_target,
"metric": metric, "window_secs": query.query_lookback_ms / 1_000, "accuracy": query.accuracy_target,
"lifecycle": query.summary_lifecycle_inputs, "evaluation_phase_ms": 0, "window_cost_model": { "implementation_id": "test", "cost": query.window_realization_candidates[0].cost }
}],
"dataset_identity": snapshot.environment.dataset_identity,
Expand Down
Loading