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
213 changes: 175 additions & 38 deletions control_plane/src/clickhouse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,9 @@ use crate::query_plan::{
MaterializationBinding, PhysicalGrouping, QueryLanguage, QueryPlan, QueryPlanEntry,
};
use asap_types::summary_catalog::SummaryCatalog;
use serde::{Deserialize, Serialize};
use serde::Deserialize;
use std::collections::HashMap;
use std::rc::Rc;

#[derive(Debug, thiserror::Error)]
pub enum ClickHousePlanningError {
Expand All @@ -43,8 +44,20 @@ pub async fn plan_clickhouse_sql(
// SQL keeps relational parents such as Project and Filter above a
// summary-capable Aggregate. Use ASAPPlanner's recursive selector here;
// the PromQL deployment lowering retains its existing conservative rules.
let cost_model = ControlPlaneCostModel::new(accuracy);
let selected = crate::planner_selection::select_summary(&canonical, &cost_model)?;
let cost_model = ControlPlaneCostModel::new(accuracy.clone());
let selected = crate::planner_selection::select_workload(
vec![(0, Rc::new(canonical.clone()))],
accuracy,
&cost_model,
)?
.into_iter()
.next()
.map(|(_, node)| node)
.ok_or_else(|| {
crate::planner_selection::SelectionError::Workload(
"SQL workload search returned no root".into(),
)
})?;
let physical = PhysicalExpr::committed(selected);
Ok(ClickHousePlannedQuery {
canonical_sql: canonical_sql_identity(&canonical),
Expand Down Expand Up @@ -90,20 +103,9 @@ pub struct ClickHouseSqlWorkloadEntry {
pub cumulative: bool,
}

/// Physical-plan components produced for the normal atomic install path.
#[derive(Debug, Serialize)]
pub struct ClickHouseCompiledBundle {
pub sds: SummaryCatalog,
pub tables: HashMap<String, planner_types::pre_asap::Schema>,
pub accuracy: AccuracyTarget,
pub query_plan: QueryPlan,
pub precompute_plan: PrecomputePlan,
pub transmission_plan: TransmissionPlan,
}

pub async fn compile_clickhouse_workload(
request: &ClickHouseSqlWorkload,
) -> Result<ClickHouseCompiledBundle, ClickHousePlanningError> {
) -> Result<crate::physical::publication::PhysicalPlanPublication, ClickHousePlanningError> {
request
.precompute_plan
.validate_against_catalog(&request.sds)
Expand Down Expand Up @@ -185,10 +187,11 @@ pub async fn compile_clickhouse_workload(
)));
}
}
Ok(ClickHouseCompiledBundle {
sds: request.sds.clone(),
tables: request.tables.clone(),
accuracy: request.accuracy.clone(),
let publication = crate::physical::publication::PhysicalPlanPublication {
summary_catalog: request.sds.clone(),
precompute_plan: request.precompute_plan.clone(),
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,
Expand All @@ -198,9 +201,11 @@ pub async fn compile_clickhouse_workload(
}),
entries,
},
precompute_plan: request.precompute_plan.clone(),
transmission_plan: request.transmission_plan.clone(),
})
};
publication
.validate()
.map_err(ClickHousePlanningError::Lower)?;
Ok(publication)
}

fn bind_selected_node(
Expand All @@ -209,13 +214,14 @@ fn bind_selected_node(
query: &ClickHouseSqlWorkloadEntry,
request: &ClickHouseSqlWorkload,
) -> Result<MaterializationBinding, crate::query_plan::QueryPlanError> {
let (metric, source_window, spatial_filter) =
crate::physical::compiler::materialization_leaf_contract(node)
let (table_ref, value_column, source_window, spatial_filter) =
clickhouse_materialization_leaf_contract(node)
.map_err(crate::query_plan::QueryPlanError::Invalid)?;
let expected = crate::physical::compiler::physical_materialization_family(family);
let selected = select_materialization(
&request.precompute_plan.materializations,
&metric,
&table_ref,
&value_column,
&spatial_filter,
&expected,
source_window.unwrap_or((query.end_ms.saturating_sub(query.start_ms)) / 1000),
Expand All @@ -230,15 +236,144 @@ fn bind_selected_node(
})
}

fn clickhouse_materialization_leaf_contract(
node: &planner_types::post_asap::SummaryNode,
) -> Result<(String, String, Option<u64>, String), String> {
use planner_types::{
post_asap::SummaryExpr,
pre_asap::{CompareOpKind, QueryExpr, ScalarValue, Source},
};
let SummaryExpr::SummaryAgg { child, input, .. } = &node.expr else {
return Err("SQL materialization leaf is not a summary aggregate".into());
};
let SummaryExpr::KeepPreAsap(expr) = &child.expr else {
return Err("SQL materialization leaf has no tabular source".into());
};
fn table_scan(
expr: &QueryExpr,
) -> Option<(
&str,
&[planner_types::pre_asap::Predicate],
&planner_types::pre_asap::Schema,
)> {
match expr {
QueryExpr::Project { child, .. } => table_scan(child),
QueryExpr::Scan {
source: Source::Table { table_ref },
predicates,
schema,
} => Some((table_ref, predicates, schema)),
_ => None,
}
}
let (source, explicit_window) = match expr.as_ref() {
QueryExpr::TimeRange { child, range }
if range.as_millis() > 0 && range.as_millis() % 1_000 == 0 =>
{
(child.as_ref(), Some(range.as_secs()))
}
source => (source, None),
};
let Some((table_ref, predicates, schema)) = table_scan(source) else {
return Err("SQL materialization leaf has no table scan".into());
};
let planner_types::post_asap::SummaryInputExpr::Column(value_input) = &input.weight else {
return Err("SQL summary requires a column-valued update".into());
};
let value_column = match value_input {
planner_types::pre_asap::ColumnRef::Named(name)
| planner_types::pre_asap::ColumnRef::Qualified { name, .. } => name.clone(),
_ => return Err("SQL summary requires a named value column".into()),
};

fn comparisons<'a>(expr: &'a QueryExpr, out: &mut Vec<&'a QueryExpr>) {
if let QueryExpr::BoolAnd(children) = expr {
for child in children {
comparisons(child, out);
}
} else {
out.push(expr);
}
}
let mut lower_ms = None;
let mut upper_ms = None;
let mut leaves = Vec::new();
for predicate in predicates {
comparisons(&predicate.0, &mut leaves);
}
for leaf in leaves {
let QueryExpr::Compare { left, op, right } = leaf else {
return Err("SQL materialization predicate is not a comparison".into());
};
let QueryExpr::Column(column) = left.as_ref() else {
return Err("SQL materialization predicate must reference a column".into());
};
let name = schema
.columns
.get(*column)
.map(|column| column.name.as_str())
.ok_or_else(|| {
"SQL materialization predicate references an unknown column".to_string()
})?;
match (name, op, right.as_ref()) {
(
name,
CompareOpKind::Gt | CompareOpKind::Ge,
QueryExpr::Literal(ScalarValue::Int64(value)),
) if schema
.time_index
.is_some_and(|index| schema.columns[index].name == name) =>
{
lower_ms = Some(*value);
}
(
name,
CompareOpKind::Lt | CompareOpKind::Le,
QueryExpr::Literal(ScalarValue::Int64(value)),
) if schema
.time_index
.is_some_and(|index| schema.columns[index].name == name) =>
{
upper_ms = Some(*value);
}
_ => {
return Err(format!(
"SQL population predicate on {name} needs a canonical catalog filter"
))
}
}
}
let inferred_window = match (lower_ms, upper_ms) {
(Some(lower), Some(upper)) if upper > lower && (upper - lower) % 1_000 == 0 => {
Some((upper - lower) as u64 / 1_000)
}
(None, None) => None,
_ => {
return Err("SQL timestamp range must provide compatible lower and upper bounds".into())
}
};
let window_secs = explicit_window.or(inferred_window).ok_or_else(|| {
"SQL table summary requires a positive whole-second timestamp range".to_string()
})?;
Ok((
table_ref.to_owned(),
value_column,
Some(window_secs),
String::new(),
))
}

fn select_materialization<'a>(
materializations: &'a [asap_types::PrecomputeMaterialization],
metric: &str,
table_ref: &str,
value_column: &str,
spatial_filter: &str,
expected: &planner_types::post_asap::SummaryFamilyType,
semantic_window_seconds: u64,
) -> Result<&'a asap_types::PrecomputeMaterialization, crate::query_plan::QueryPlanError> {
let mut matches = materializations.iter().filter(|candidate| {
candidate.metric == metric
candidate.table_name.as_deref() == Some(table_ref)
&& candidate.value_column.as_deref() == Some(value_column)
&& candidate.spatial_filter_normalized == spatial_filter
&& candidate
.accumulator_spec()
Expand All @@ -250,12 +385,12 @@ fn select_materialization<'a>(
});
let selected = matches.next().ok_or_else(|| {
crate::query_plan::QueryPlanError::Invalid(format!(
"no precompute materialization matches {metric}/{expected:?}"
"no precompute materialization matches {table_ref}.{value_column}/{expected:?}"
))
})?;
if matches.next().is_some() {
return Err(crate::query_plan::QueryPlanError::Invalid(format!(
"ambiguous precompute materializations match {metric}/{expected:?}"
"ambiguous precompute materializations match {table_ref}.{value_column}/{expected:?}"
)));
}
Ok(selected)
Expand All @@ -268,7 +403,7 @@ mod tests {

fn materialization(
agg: AggregationType,
metric: &str,
value_column: &str,
window: u64,
slide: u64,
parameter: (&str, serde_json::Value),
Expand All @@ -285,10 +420,10 @@ mod tests {
slide,
WindowKind::Tumbling,
String::new(),
metric.into(),
None,
None,
format!("telemetry.{value_column}"),
None,
Some("telemetry".into()),
Some(value_column.into()),
);
value.pane_origin_ms = Some(0);
value
Expand Down Expand Up @@ -349,29 +484,31 @@ mod tests {
let sum_family = sum_60.accumulator_spec().unwrap().family;
let count_family = count_60.accumulator_spec().unwrap().family;
assert_eq!(
select_materialization(&configs, "requests", "", &sum_family, 60)
select_materialization(&configs, "telemetry", "requests", "", &sum_family, 60)
.unwrap()
.policy_fingerprint(),
sum_60.policy_fingerprint()
);
let dd_family = dd_2.accumulator_spec().unwrap().family;
assert_eq!(
select_materialization(&configs, "requests", "", &dd_family, 60)
select_materialization(&configs, "telemetry", "requests", "", &dd_family, 60)
.unwrap()
.policy_fingerprint(),
dd_2.policy_fingerprint()
);
assert_eq!(
select_materialization(&configs, "requests", "", &count_family, 60)
select_materialization(&configs, "telemetry", "requests", "", &count_family, 60)
.unwrap()
.policy_fingerprint(),
count_60.policy_fingerprint()
);
assert!(select_materialization(&configs, "missing", "", &sum_family, 60).is_err());
assert!(
select_materialization(&configs, "telemetry", "missing", "", &sum_family, 60).is_err()
);
let mut ambiguous = configs.clone();
ambiguous.push(sum_60);
assert!(
select_materialization(&ambiguous, "requests", "", &sum_family, 60)
select_materialization(&ambiguous, "telemetry", "requests", "", &sum_family, 60)
.unwrap_err()
.to_string()
.contains("ambiguous")
Expand Down
15 changes: 4 additions & 11 deletions control_plane/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -809,26 +809,19 @@ async fn handle_compile_and_publish_clickhouse_plan(
State(state): State<AppState>,
Json(request): Json<clickhouse::ClickHouseSqlWorkload>,
) -> impl IntoResponse {
let bundle = match clickhouse::compile_clickhouse_workload(&request).await {
Ok(bundle) => bundle,
let publication = match clickhouse::compile_clickhouse_workload(&request).await {
Ok(publication) => publication,
Err(error) => return (StatusCode::BAD_REQUEST, error.to_string()).into_response(),
};
let plan_id = bundle.sds.plan_id;
let plan_version = bundle.sds.plan_version;
let plan_id = publication.summary_catalog.plan_id;
let plan_version = publication.summary_catalog.plan_version;
let Some(client) = state.backend_client.as_ref() else {
return (
StatusCode::SERVICE_UNAVAILABLE,
"backend publication is not configured",
)
.into_response();
};
let publication = physical::publication::PhysicalPlanPublication {
summary_catalog: bundle.sds,
precompute_plan: bundle.precompute_plan,
collector_plans: Vec::new(),
transmission_plan: bundle.transmission_plan,
query_plan: bundle.query_plan,
};
if let Err(error) = client
.post_catalog_plan_typed(&publication, None, &[])
.await
Expand Down
Loading