Skip to content
Merged
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
4 changes: 2 additions & 2 deletions control_plane/examples/audit_clickhouse_corpus.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ fn publication_inputs(schema: &Schema, sql: String) -> ClickHouseSqlWorkload {
let mut precompute_plan =
PrecomputePlan::build_backend_local(envelope.clone(), vec![materialization.clone()])
.unwrap();
let mut transmission_plan = control_plane::physical::compiler::compile_transmission_plan(
let mut transmission_plan = control_plane::physical::compiler::build_transmission_plan(
envelope,
&precompute_plan,
&Default::default(),
Expand All @@ -72,7 +72,7 @@ fn publication_inputs(schema: &Schema, sql: String) -> ClickHouseSqlWorkload {
precompute_plan.summary_catalog = Some(reference.clone());
transmission_plan.summary_catalog = Some(reference);
ClickHouseSqlWorkload {
sds,
summary_catalog: sds,
precompute_plan,
transmission_plan,
tables: std::collections::HashMap::from([("raw_samples".into(), schema.clone())]),
Expand Down
24 changes: 12 additions & 12 deletions control_plane/examples/calibration_candidates.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
//! Export every bindable candidate for isolated measurement, without selecting a winner.
use control_plane::physical::{
compiler::{BackendLocalPlanningSnapshot, PhysicalCompiler},
compiler::{BackendLocalPlanningInput, PhysicalPlanCompiler},
workload_cost,
};
use planner_types::post_asap::{SummaryExpr, SummaryNode};
Expand All @@ -9,7 +9,7 @@ use std::{collections::BTreeMap, rc::Rc};

// Planner IR does not implement Serialize. Preserve actual DAG identity and
// typed variant/edges; leaf metadata uses explicitly labelled Debug encoding.
fn planner_forest(queries: &[control_plane::physical::compiler::PlanningQuery]) -> Value {
fn planner_forest(queries: &[control_plane::physical::compiler::QueryCompilationInput]) -> Value {
fn visit(
node: &Rc<SummaryNode>,
seen: &mut BTreeMap<usize, usize>,
Expand Down Expand Up @@ -118,7 +118,7 @@ fn planner_forest(queries: &[control_plane::physical::compiler::PlanningQuery])
}
let mut seen = BTreeMap::new();
let mut nodes = BTreeMap::new();
let roots:Vec<_>=queries.iter().map(|q|json!({"query_id":q.query_id,"original_promql":q.query_string,"root":visit(&q.post_asap,&mut seen,&mut nodes)})).collect();
let roots:Vec<_>=queries.iter().map(|q|json!({"query_id":q.query_id,"original_promql":q.query_string,"root":visit(&q.selected_plan_root,&mut seen,&mut nodes)})).collect();
json!({"encoding":"structured_graph_with_debug_metadata_v1","scope":"actual candidate Planner post-ASAP input before physical binding; not reconstructed from installed nodes","roots":roots,"nodes":nodes})
}

Expand All @@ -127,26 +127,26 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
.nth(1)
.ok_or("usage: calibration_candidates SNAPSHOT.json [--metricsql]")?;
let metricsql = std::env::args().skip(2).any(|arg| arg == "--metricsql");
let snapshot: BackendLocalPlanningSnapshot = serde_json::from_slice(&std::fs::read(path)?)?;
let (request, environment) = snapshot.planning_request()?;
let snapshot: BackendLocalPlanningInput = serde_json::from_slice(&std::fs::read(path)?)?;
let (request, environment) = snapshot.into_physical_compilation_request()?;
let mut results = Vec::new();
for (index, candidate) in workload_cost::with_exact_alternative(request)?
for (index, candidate) in workload_cost::enumerate_exact_and_materialized_candidates(request)?
.into_iter()
.enumerate()
{
let queries = candidate.queries.clone();
let materialization_policy = candidate.materialization_policy.clone();
let enabled_materialization_keys = candidate.enabled_materialization_keys.clone();
let planner_selected_queries = planner_forest(&queries);
let compiled = if metricsql {
PhysicalCompiler.compile_metricsql(candidate, environment.clone())
PhysicalPlanCompiler.compile_metricsql(candidate, environment.clone())
} else {
PhysicalCompiler.compile(candidate, environment.clone())
PhysicalPlanCompiler.compile_promql(candidate, environment.clone())
};
let plan = match compiled {
Ok(plan) => plan,
Err(error) => {
results.push(
json!({"candidate_index": index, "materialization_policy": materialization_policy, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}),
json!({"candidate_index": index, "materialization_policy": enabled_materialization_keys, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}),
);
continue;
}
Expand All @@ -155,14 +155,14 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
Ok(manifest) => manifest,
Err(error) => {
results.push(
json!({"candidate_index": index, "materialization_policy": materialization_policy, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}),
json!({"candidate_index": index, "materialization_policy": enabled_materialization_keys, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}),
);
continue;
}
};
results.push(json!({
"candidate_index": index,
"materialization_policy": materialization_policy,
"materialization_policy": enabled_materialization_keys,
"planner_selected_queries": planner_selected_queries,
"manifest": manifest,
"lifecycle_estimates": plan.lifecycle_estimates,
Expand Down
8 changes: 4 additions & 4 deletions control_plane/examples/compile_workload_artifact.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
//! Control-plane entry point: cost-select a workload and emit its atomic install request.
use control_plane::physical::compiler::BackendLocalPlanningSnapshot;
use control_plane::physical::compiler::BackendLocalPlanningInput;
use serde_json::json;

fn main() -> Result<(), Box<dyn std::error::Error>> {
Expand All @@ -15,12 +15,12 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
if args.next().is_some() {
return Err("unexpected arguments".into());
}
let snapshot: BackendLocalPlanningSnapshot = serde_json::from_slice(&std::fs::read(path)?)?;
let snapshot: BackendLocalPlanningInput = serde_json::from_slice(&std::fs::read(path)?)?;
let start = std::time::Instant::now();
let plan = if metricsql {
snapshot.compile_metricsql()?
} else {
snapshot.compile()?
snapshot.compile_promql()?
};
let elapsed = start.elapsed().as_nanos();
let comparison = plan
Expand All @@ -33,7 +33,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
"planning_elapsed_ns": elapsed,
"envelope": plan.envelope,
"cost_comparison": comparison,
"logical_selection": plan.logical_selection,
"logical_selection": plan.planner_selection_trace,
"backend_revision": control_plane::physical::compiler::BACKEND_REVISION,
"planner_revision": control_plane::physical::compiler::PLANNER_REVISION,
"lifecycle_estimates": plan.lifecycle_estimates,
Expand Down
12 changes: 6 additions & 6 deletions control_plane/examples/workload_cost_manifest.rs
Original file line number Diff line number Diff line change
@@ -1,21 +1,21 @@
//! Emit pricing requirements; never fabricate quotes or publish a plan.
use control_plane::physical::{
compiler::BackendLocalPlanningSnapshot, compiler::PhysicalCompiler, workload_cost,
compiler::BackendLocalPlanningInput, compiler::PhysicalPlanCompiler, workload_cost,
};

fn main() -> Result<(), Box<dyn std::error::Error>> {
let path = std::env::args()
.nth(1)
.ok_or("usage: workload_cost_manifest SNAPSHOT.json")?;
let snapshot: BackendLocalPlanningSnapshot =
let snapshot: BackendLocalPlanningInput =
serde_json::from_str(&std::fs::read_to_string(path)?)?;
let (request, environment) = snapshot.planning_request()?;
let manifests = workload_cost::with_exact_alternative(request)?
let (request, environment) = snapshot.into_physical_compilation_request()?;
let manifests = workload_cost::enumerate_exact_and_materialized_candidates(request)?
.into_iter()
.filter_map(|candidate| {
let queries = candidate.queries.clone();
PhysicalCompiler
.compile(candidate, environment.clone())
PhysicalPlanCompiler
.compile_promql(candidate, environment.clone())
.and_then(|plan| workload_cost::manifest(&plan, &queries))
.ok()
})
Expand Down
22 changes: 12 additions & 10 deletions control_plane/src/backend_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -549,16 +549,18 @@ mod tests {

#[tokio::test]
async fn catalog_publication_posts_canonical_document_without_legacy_bytes() {
let snapshot: crate::physical::compiler::BackendLocalPlanningSnapshot =
serde_json::from_str(include_str!(
"../../docs/examples/asapquery-planning-snapshot.json"
))
.unwrap();
let publication = crate::physical::compiler::tests::quoted_snapshot(snapshot, false)
.compile()
.unwrap()
.publication()
.unwrap();
let snapshot: crate::physical::compiler::BackendLocalPlanningInput = serde_json::from_str(
include_str!("../../docs/examples/asapquery-planning-snapshot.json"),
)
.unwrap();
let publication = crate::physical::compiler::tests::quoted_snapshot(
snapshot,
crate::physical::compiler::QueryFrontend::PromQl,
)
.compile_promql()
.unwrap()
.to_publication_artifact()
.unwrap();
let hits: StdArc<Mutex<Vec<serde_json::Value>>> = StdArc::new(Mutex::new(Vec::new()));
let route_hits = hits.clone();
let app = Router::new().route(
Expand Down
32 changes: 18 additions & 14 deletions control_plane/src/clickhouse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -205,7 +205,8 @@ pub use asap_frontend_sql::SqlCatalog as ClickHouseSqlCatalog;

#[derive(Debug, Deserialize)]
pub struct ClickHouseSqlWorkload {
pub sds: SummaryCatalog,
#[serde(rename = "sds", alias = "summary_catalog")]
pub summary_catalog: SummaryCatalog,
pub precompute_plan: PrecomputePlan,
pub transmission_plan: TransmissionPlan,
pub tables: HashMap<String, planner_types::pre_asap::Schema>,
Expand Down Expand Up @@ -301,7 +302,7 @@ pub async fn compile_automatic_clickhouse_workload(
.map_err(|error| ClickHousePlanningError::Lower(error.to_string()))?,
);
precompute.executable_dags = installed_dags;
let mut transmission = crate::physical::compiler::compile_transmission_plan(
let mut transmission = crate::physical::compiler::build_transmission_plan(
request.envelope.clone(),
&precompute,
&std::collections::BTreeMap::new(),
Expand Down Expand Up @@ -410,7 +411,7 @@ pub async fn compile_clickhouse_workload(
) -> Result<crate::physical::publication::PhysicalPlanPublication, ClickHousePlanningError> {
request
.precompute_plan
.validate_against_catalog(&request.sds)
.validate_against_catalog(&request.summary_catalog)
.map_err(|error| ClickHousePlanningError::Lower(error.to_string()))?;
request
.transmission_plan
Expand Down Expand Up @@ -441,13 +442,13 @@ pub async fn compile_clickhouse_workload(
let mut precompute_plan = request.precompute_plan.clone();
precompute_plan.executable_dags = installed_dags;
let publication = crate::physical::publication::PhysicalPlanPublication {
summary_catalog: request.sds.clone(),
summary_catalog: request.summary_catalog.clone(),
precompute_plan,
collector_plans: Vec::new(),
transmission_plan: request.transmission_plan.clone(),
query_plan: QueryPlan {
plan_id: request.sds.plan_id,
plan_version: request.sds.plan_version,
plan_id: request.summary_catalog.plan_id,
plan_version: request.summary_catalog.plan_version,
clickhouse_context: Some(ClickHousePlanningContext {
window_templates,
tables: request.tables.clone(),
Expand Down Expand Up @@ -1405,7 +1406,7 @@ mod tests {
let mut precompute =
PrecomputePlan::build_backend_local(envelope.clone(), vec![config]).unwrap();
precompute.summary_catalog = Some(sds.reference().unwrap());
let mut transmission = crate::physical::compiler::compile_transmission_plan(
let mut transmission = crate::physical::compiler::build_transmission_plan(
envelope,
&precompute,
&std::collections::BTreeMap::new(),
Expand All @@ -1423,7 +1424,7 @@ mod tests {
)
};
let mut request = ClickHouseSqlWorkload {
sds,
summary_catalog: sds,
precompute_plan: precompute,
transmission_plan: transmission,
tables: HashMap::from([
Expand Down Expand Up @@ -1453,7 +1454,7 @@ mod tests {
let simple_sql =
"SELECT sum(value) FROM telemetry WHERE timestamp_ms >= 0 AND timestamp_ms < 2000";
let simple = compile_clickhouse_workload(&ClickHouseSqlWorkload {
sds: request.sds.clone(),
summary_catalog: request.summary_catalog.clone(),
precompute_plan: request.precompute_plan.clone(),
transmission_plan: request.transmission_plan.clone(),
tables: request.tables.clone(),
Expand Down Expand Up @@ -1498,7 +1499,7 @@ mod tests {
]
};
let multiple = compile_clickhouse_workload(&ClickHouseSqlWorkload {
sds: request.sds.clone(),
summary_catalog: request.summary_catalog.clone(),
precompute_plan: request.precompute_plan.clone(),
transmission_plan: request.transmission_plan.clone(),
tables: request.tables.clone(),
Expand Down Expand Up @@ -1711,18 +1712,21 @@ mod tests {
value: planner_types::pre_asap::ScalarValue::Utf8("requests".into()),
}],
});
request.sds = SummaryCatalog::from_materializations(71, 1, &[config.clone()]).unwrap();
request.summary_catalog =
SummaryCatalog::from_materializations(71, 1, &[config.clone()]).unwrap();
let envelope = request.precompute_plan.envelope.clone();
request.precompute_plan =
PrecomputePlan::build_backend_local(envelope.clone(), vec![config]).unwrap();
request.precompute_plan.summary_catalog = Some(request.sds.reference().unwrap());
request.transmission_plan = crate::physical::compiler::compile_transmission_plan(
request.precompute_plan.summary_catalog =
Some(request.summary_catalog.reference().unwrap());
request.transmission_plan = crate::physical::compiler::build_transmission_plan(
envelope,
&request.precompute_plan,
&std::collections::BTreeMap::new(),
)
.unwrap();
request.transmission_plan.summary_catalog = Some(request.sds.reference().unwrap());
request.transmission_plan.summary_catalog =
Some(request.summary_catalog.reference().unwrap());
assert!(compile_clickhouse_workload(&request).await.is_ok());
}
}
2 changes: 1 addition & 1 deletion control_plane/src/emit/backend_push.rs
Original file line number Diff line number Diff line change
Expand Up @@ -231,7 +231,7 @@ async fn push_documents_coupled(
warn!(%error, "failed to bind compatibility PrecomputePlan to SummaryCatalog");
return (false, false, 0);
}
let transmission_plan = match crate::physical::compiler::compile_transmission_plan(
let transmission_plan = match crate::physical::compiler::build_transmission_plan(
precompute_plan.envelope.clone(),
&precompute_plan,
&Default::default(),
Expand Down
Loading
Loading