From fec0c51b90bfd903c242cc376ec6f56cf1054034 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 08:35:59 -0600 Subject: [PATCH 1/2] feat(clickhouse): execute SQL joins in the shared query DAG --- Cargo.lock | 36 ++--- control_plane/Cargo.toml | 10 +- .../examples/calibration_candidates.rs | 10 ++ .../examples/offline_planner_replay.rs | 5 + control_plane/src/emit/mod.rs | 1 + .../src/physical/colored_dag/allocator.rs | 7 + .../src/physical/colored_dag/emitter.rs | 1 + control_plane/src/physical/compiler.rs | 6 + control_plane/src/physical/post_asap/tests.rs | 1 + control_plane/src/query_plan.rs | 33 +++- control_plane/src/query_plan/logical.rs | 4 + crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +- .../asap_clickhouse_query_engine/execution.rs | 149 ++++++++++++++++++ .../relational_adapter.rs | 99 ++++++++++++ .../asap_query_engine/post_asap_readout.rs | 3 +- .../asap_query_engine/summary_exec.rs | 1 + 17 files changed, 344 insertions(+), 30 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 9acd1f3b..4c18ff77 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -374,7 +374,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" dependencies = [ "asap-types", "serde", @@ -385,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-metricsql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" dependencies = [ "asap-types", "metricsql_parser", @@ -395,7 +395,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" dependencies = [ "asap-types", "promql-parser 0.10.0", @@ -404,7 +404,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -426,12 +426,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" dependencies = [ "serde", "serde_json", @@ -1711,7 +1711,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -2260,7 +2260,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.6.3", + "socket2 0.5.10", "tokio", "tower-service", "tracing", @@ -2453,7 +2453,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" dependencies = [ "hermit-abi 0.5.2", "libc", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -2801,7 +2801,7 @@ dependencies = [ [[package]] name = "metricsql_common" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" dependencies = [ "chrono", ] @@ -2809,7 +2809,7 @@ dependencies = [ [[package]] name = "metricsql_parser" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=26e580710cf9b44b282598719d6c5ea9a9cc62be#26e580710cf9b44b282598719d6c5ea9a9cc62be" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" dependencies = [ "ahash", "chrono", @@ -3433,7 +3433,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "be769465445e8c1474e9c5dac2018218498557af32d9ed057325ec9a41ae81bf" dependencies = [ "heck 0.5.0", - "itertools 0.13.0", + "itertools 0.10.5", "log", "multimap", "once_cell", @@ -3453,7 +3453,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d" dependencies = [ "anyhow", - "itertools 0.13.0", + "itertools 0.10.5", "proc-macro2", "quote", "syn 2.0.117", @@ -3551,7 +3551,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls 0.23.40", - "socket2 0.6.3", + "socket2 0.5.10", "thiserror 2.0.18", "tokio", "tracing", @@ -3588,7 +3588,7 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.6.3", + "socket2 0.5.10", "tracing", "windows-sys 0.52.0", ] @@ -3904,7 +3904,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -4438,7 +4438,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix 1.1.4", - "windows-sys 0.61.2", + "windows-sys 0.52.0", ] [[package]] @@ -5260,7 +5260,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.48.0", ] [[package]] diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 47d3bbf1..8d1b09f5 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -91,8 +91,8 @@ asap_types.workspace = true # scaffolding, unaware that `data_plane`'s `summary_executor.rs` in *this* # repo is a real one. Vendored locally instead of chased upstream -- see # `data_plane/src/query_engines/asap_query_engine/summary_exec.rs`. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } # L1 adoption (design-target-architecture.md Part B): the PromQL front # end itself, replacing control_plane's own query_parser/promql.rs. @@ -100,9 +100,9 @@ asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = # `planner-types`/`asap-aware-mapping` above -- these three MUST move # together (two revs of the same upstream repo's types in one workspace # resolve to distinct Rust types that won't unify). -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } -asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } +asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/control_plane/examples/calibration_candidates.rs b/control_plane/examples/calibration_candidates.rs index 64ce8d3c..607bc150 100644 --- a/control_plane/examples/calibration_candidates.rs +++ b/control_plane/examples/calibration_candidates.rs @@ -74,6 +74,16 @@ fn planner_forest(queries: &[control_plane::physical::compiler::PlanningQuery]) vec![outer, inner], json!({"key_debug":format!("{key:?}"),"family_debug":format!("{family:?}")}), ), + SummaryExpr::RelationalJoin { + left, + right, + kind, + pred, + } => ( + "RelationalJoin", + vec![left, right], + json!({"kind_debug":format!("{kind:?}"),"predicate_debug":format!("{pred:?}")}), + ), SummaryExpr::SummarySubtract { left, right } => { ("SummarySubtract", vec![left, right], json!({})) } diff --git a/control_plane/examples/offline_planner_replay.rs b/control_plane/examples/offline_planner_replay.rs index 0ccf9d05..f0e31758 100644 --- a/control_plane/examples/offline_planner_replay.rs +++ b/control_plane/examples/offline_planner_replay.rs @@ -73,6 +73,11 @@ fn inspect( left: lhs, right: rhs, } + | SummaryExpr::RelationalJoin { + left: lhs, + right: rhs, + .. + } | SummaryExpr::BinaryOp { lhs, rhs, .. } => { inspect(lhs, model, seen, states, raw); inspect(rhs, model, seen, states, raw); diff --git a/control_plane/src/emit/mod.rs b/control_plane/src/emit/mod.rs index af4fe16f..9732e9de 100644 --- a/control_plane/src/emit/mod.rs +++ b/control_plane/src/emit/mod.rs @@ -326,6 +326,7 @@ fn extract_from_node(node: &Rc) -> Option { // Not surfaced by any `Bind*` path yet (gated on rules that // haven't landed — see `deployment_expr.rs`'s module docs). SummaryExpr::BinaryOp { .. } + | SummaryExpr::RelationalJoin { .. } | SummaryExpr::SummaryJoin { .. } | SummaryExpr::SummarySubtract { .. } | SummaryExpr::SummaryDelete { .. } diff --git a/control_plane/src/physical/colored_dag/allocator.rs b/control_plane/src/physical/colored_dag/allocator.rs index b66dcff3..5cbf74a2 100644 --- a/control_plane/src/physical/colored_dag/allocator.rs +++ b/control_plane/src/physical/colored_dag/allocator.rs @@ -184,6 +184,13 @@ impl ThreeStageWalker { } StageId::Backend } + SummaryExpr::RelationalJoin { left, right, .. } => { + for child in [left, right] { + let (cid, _) = self.visit_l4node(child)?; + self.dag.edges.push((id, cid)); + } + StageId::Backend + } SummaryExpr::ValueOperation { child, timing, .. } => { let (cid, child_stage) = self.visit_l4node(child)?; self.dag.edges.push((id, cid)); diff --git a/control_plane/src/physical/colored_dag/emitter.rs b/control_plane/src/physical/colored_dag/emitter.rs index 2bd4bcf2..9d230e89 100644 --- a/control_plane/src/physical/colored_dag/emitter.rs +++ b/control_plane/src/physical/colored_dag/emitter.rs @@ -117,6 +117,7 @@ fn classify(expr: &PhysicalExpr) -> NodeKind<'_> { SummaryExpr::SummaryEstimate { query, .. } => NodeKind::SketchEstimate { query }, SummaryExpr::SummaryMerge { .. } => NodeKind::SketchMerge, SummaryExpr::BinaryOp { .. } + | SummaryExpr::RelationalJoin { .. } | SummaryExpr::CandidateTopK { .. } | SummaryExpr::ValueOperation { .. } | SummaryExpr::SummaryJoin { .. } diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 421f8201..a75029ef 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -2885,6 +2885,7 @@ fn summary_agg_metric(node: &SummaryNode) -> Option { inner: right, .. } + | SummaryExpr::RelationalJoin { left, right, .. } | SummaryExpr::SummarySubtract { left, right } | SummaryExpr::BinaryOp { lhs: left, @@ -3032,6 +3033,7 @@ fn requires_exact_erp_fallback( walk(inner, out); } SummaryExpr::SummarySubtract { left, right } + | SummaryExpr::RelationalJoin { left, right, .. } | SummaryExpr::BinaryOp { lhs: left, rhs: right, @@ -3845,6 +3847,10 @@ fn collect_selected_materializations( SummaryExpr::ValueOperation { child, .. } => { walk(child, readout, composable, grouping.clone(), selected)?; } + SummaryExpr::RelationalJoin { left, right, .. } => { + walk(left, readout, composable, grouping.clone(), selected)?; + walk(right, readout, composable, grouping.clone(), selected)?; + } SummaryExpr::BinaryOp { lhs, rhs, .. } if composable || crate::query_plan::exact_value_executable(node) => { diff --git a/control_plane/src/physical/post_asap/tests.rs b/control_plane/src/physical/post_asap/tests.rs index 745b2ab2..3735861c 100644 --- a/control_plane/src/physical/post_asap/tests.rs +++ b/control_plane/src/physical/post_asap/tests.rs @@ -120,6 +120,7 @@ fn node_is_archive(node: &Rc) -> bool { candidates, values, .. } => node_is_archive(candidates) || node_is_archive(values), SummaryExpr::SummarySubtract { left, right } + | SummaryExpr::RelationalJoin { left, right, .. } | SummaryExpr::BinaryOp { lhs: left, rhs: right, diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index 8d3a8192..4a9741ab 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -636,6 +636,14 @@ pub struct ExternalExactRequest { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(tag = "op", rename_all = "snake_case", deny_unknown_fields)] pub enum QueryPlanNode { + RelationalJoin { + inputs: [QueryNodeId; 2], + join_kind: planner_types::pre_asap::JoinKind, + pred: serde_json::Value, + left_schema: planner_types::post_asap::SummarySchema, + right_schema: planner_types::post_asap::SummarySchema, + output_schema: planner_types::post_asap::SummarySchema, + }, Relational { input: QueryNodeId, /// Serialized planner-owned operation. Keeping the wire form here makes @@ -699,7 +707,7 @@ impl QueryPlanNode { Self::Scalar { .. } | Self::ReadMaterialization { .. } | Self::ExactFallback { .. } => { &[] } - Self::Binary { inputs, .. } => inputs, + Self::Binary { inputs, .. } | Self::RelationalJoin { inputs, .. } => inputs, Self::ReduceSum { input, .. } | Self::Relational { input, .. } | Self::SummaryEstimate { input, .. } @@ -806,7 +814,8 @@ where } } QueryPlanNode::CandidateTopK { inputs, .. } - | QueryPlanNode::Binary { inputs, .. } => { + | QueryPlanNode::Binary { inputs, .. } + | QueryPlanNode::RelationalJoin { inputs, .. } => { for input in inputs { *input = remap[input]; } @@ -872,6 +881,26 @@ where } let physical = match &node.expr { + SummaryExpr::RelationalJoin { + left, + right, + kind, + pred, + } if self.preserve_relational => QueryPlanNode::RelationalJoin { + inputs: [self.lower(left)?, self.lower(right)?], + join_kind: kind.clone(), + pred: serde_json::to_value(pred).map_err(|error| { + QueryPlanError::Invalid(format!( + "cannot serialize relational join predicate: {error}" + )) + })?, + left_schema: left.schema.clone(), + right_schema: right.schema.clone(), + output_schema: node.schema.clone(), + }, + SummaryExpr::RelationalJoin { .. } => QueryPlanNode::ExactFallback { + reason: "read-time relational join requires the relational compiler".into(), + }, SummaryExpr::ValueOperation { child, operation, .. } if self.preserve_relational diff --git a/control_plane/src/query_plan/logical.rs b/control_plane/src/query_plan/logical.rs index 903f2e98..05d47471 100644 --- a/control_plane/src/query_plan/logical.rs +++ b/control_plane/src/query_plan/logical.rs @@ -1289,6 +1289,10 @@ pub fn materialization_candidate_keys( visit(original, lhs, keys)?; visit(original, rhs, keys)?; } + SummaryExpr::RelationalJoin { left, right, .. } => { + visit(original, left, keys)?; + visit(original, right, keys)?; + } SummaryExpr::CandidateTopK { candidates, values, .. } => { diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index 80ac2800..b968c7ef 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -32,4 +32,4 @@ sha2 = "0.10" # exactly (`control_plane/Cargo.toml`) -- two different revs of the same # git dependency in one workspace resolve to two distinct Rust types that # won't unify. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 38619a2b..45b77aab 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -38,8 +38,8 @@ control_plane = { path = "../control_plane" } # reduction: Reduction, .. }`) are `pre_asap` types, in the same crate now # (not a separate `asap-ir` import). Query serving consumes the compiled # QueryPlan; these types are used at physical-plan compilation boundaries. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } -asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } +asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } # Shared external (workspace) serde.workspace = true @@ -130,7 +130,7 @@ crc32fast = "1.4" # none of them. [dev-dependencies] -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "26e580710cf9b44b282598719d6c5ea9a9cc62be" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } tempfile = "3.20.0" criterion = { version = "0.5", features = ["html_reports"] } tokio-tungstenite = "0.21" diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs index 61985220..31369947 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/execution.rs @@ -24,6 +24,102 @@ pub enum ClickHouseDagFallback { }, ResultEncoding(String), } + +fn apply_relational_operation( + operation: serde_json::Value, + output_schema: &planner_types::post_asap::SummarySchema, + relation: ClickHouseRelation, +) -> Result { + let adapter = ClickHouseRelationalAdapter; + if let Some(filter) = operation.get("Filter") { + let predicate = filter + .get("pred") + .cloned() + .ok_or_else(|| "published Filter lacks pred".to_owned()) + .and_then(|value| serde_json::from_value(value).map_err(|error| error.to_string()))?; + return adapter + .apply_filter(&predicate, relation) + .map_err(|error| error.to_string()); + } + let operation: ValueOperation = + serde_json::from_value(operation).map_err(|error| error.to_string())?; + adapter + .apply_operation(&operation, output_schema, relation) + .map_err(|error| error.to_string()) +} + +fn execute_relation_subtree( + index: &SketchStore, + entry: &QueryPlanEntry, + root: QueryNodeId, + expected_schema: &planner_types::post_asap::SummarySchema, + t0_ms: u64, + t1_ms: u64, + is_cumulative: bool, +) -> Result { + match entry.nodes.get(&root) { + Some(QueryPlanNode::Relational { + input, + operation, + input_schema, + output_schema, + }) => { + let input = execute_relation_subtree( + index, + entry, + *input, + input_schema, + t0_ms, + t1_ms, + is_cumulative, + )?; + apply_relational_operation(operation.clone(), output_schema, input) + } + Some(QueryPlanNode::RelationalJoin { + inputs, + join_kind, + pred, + left_schema, + right_schema, + output_schema, + }) => { + if !matches!(join_kind, planner_types::pre_asap::JoinKind::Inner) { + return Err("only inner relational joins are executable".into()); + } + let left = execute_relation_subtree( + index, + entry, + inputs[0], + left_schema, + t0_ms, + t1_ms, + is_cumulative, + )?; + let right = execute_relation_subtree( + index, + entry, + inputs[1], + right_schema, + t0_ms, + t1_ms, + is_cumulative, + )?; + let pred = serde_json::from_value(pred.clone()).map_err(|error| error.to_string())?; + ClickHouseRelationalAdapter + .apply_inner_equi_join(&pred, output_schema, left, right) + .map_err(|error| error.to_string()) + } + Some(_) => { + validate_reachable(entry, root)?; + let outcome = + execute_query_plan_from_readout(index, entry, root, t0_ms, t1_ms, is_cumulative) + .map_err(|error| format!("{error:?}"))?; + ClickHouseRelation::from_series_rows(expected_schema, outcome.series, outcome.coverage) + .map_err(|error| error.to_string()) + } + None => Err(format!("published DAG references missing node {}", root.0)), + } +} pub enum ClickHouseDagOutcome { Accelerated(ClickHouseQueryResult), Fallback(ClickHouseDagFallback), @@ -55,6 +151,59 @@ pub fn execute_sql_dag( error.to_string(), )); } + let has_relational_join = entry.topological_order().is_ok_and(|ids| { + ids.iter().any(|id| { + matches!( + entry.nodes.get(id), + Some(QueryPlanNode::RelationalJoin { .. }) + ) + }) + }); + if has_relational_join { + let root_schema = match entry.nodes.get(&entry.root) { + Some(QueryPlanNode::Relational { output_schema, .. }) + | Some(QueryPlanNode::RelationalJoin { output_schema, .. }) => output_schema, + _ => { + return ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::UnsupportedPlan( + "relational join plan root has no relation schema".into(), + )) + } + }; + let relation = match execute_relation_subtree( + index, + entry, + entry.root, + root_schema, + t0_ms, + t1_ms, + is_cumulative, + ) { + Ok(relation) => relation, + Err(error) => { + return ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::UnsupportedPlan( + error, + )) + } + }; + let pane_ms = entry + .materialization_bindings() + .iter() + .map(|binding| binding.window_ms) + .max() + .unwrap_or(0); + if !complete_pane_coverage(relation.coverage, (t0_ms, t1_ms), pane_ms) { + return ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::IncompleteCoverage { + requested: (t0_ms, t1_ms), + observed: relation.coverage, + }); + } + return match relation.into_result() { + Ok(result) => ClickHouseDagOutcome::Accelerated(result), + Err(error) => ClickHouseDagOutcome::Fallback(ClickHouseDagFallback::ResultEncoding( + error.to_string(), + )), + }; + } let mut base_root = entry.root; let mut relational = Vec::new(); loop { diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs index 69402b65..7ddef54c 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs @@ -93,6 +93,39 @@ impl ClickHouseRelation { pub struct ClickHouseRelationalAdapter; impl ClickHouseRelationalAdapter { + pub fn apply_inner_equi_join( + &self, + pred: &planner_types::pre_asap::Predicate, + output_schema: &SummarySchema, + left: ClickHouseRelation, + right: ClickHouseRelation, + ) -> Result { + let coverage = match (left.coverage, right.coverage) { + (Some((left_start, left_end)), Some((right_start, right_end))) => { + let start = left_start.max(right_start); + let end = left_end.min(right_end); + (start <= end).then_some((start, end)) + } + _ => None, + }; + let mut rows = Vec::new(); + for left_row in &left.rows { + for right_row in &right.rows { + let mut joined = Vec::with_capacity(left_row.len() + right_row.len()); + joined.extend(left_row.iter().cloned()); + joined.extend(right_row.iter().cloned()); + if matches!(eval(&pred.0, &joined)?, Cell::Bool(true)) { + rows.push(joined); + } + } + } + Ok(ClickHouseRelation { + rows, + fields: fields_from_schema(output_schema), + coverage, + }) + } + pub fn apply_filter( &self, pred: &planner_types::pre_asap::Predicate, @@ -646,4 +679,70 @@ mod tests { let error = eval(&QueryExpr::BoolAnd(vec![]), &row).unwrap_err(); assert!(matches!(error, ClickHouseRelationalError::Unsupported(_))); } + + #[test] + fn inner_equi_join_feeds_typed_ratio_projection() { + let side_schema = schema(&[("service", DataType::Utf8), ("value", DataType::Float64)]); + let left = ClickHouseRelation { + rows: vec![vec![Cell::Utf8("api".into()), Cell::Float64(2.0)]], + fields: fields_from_schema(&side_schema), + coverage: Some((300_000, 600_000)), + }; + let right = ClickHouseRelation { + rows: vec![vec![Cell::Utf8("api".into()), Cell::Float64(10.0)]], + fields: fields_from_schema(&side_schema), + coverage: Some((300_000, 600_000)), + }; + let joined_schema = schema(&[ + ("service", DataType::Utf8), + ("left_value", DataType::Float64), + ("service", DataType::Utf8), + ("right_value", DataType::Float64), + ]); + let pred = planner_types::pre_asap::Predicate(Rc::new(QueryExpr::Compare { + left: Rc::new(QueryExpr::Column(0)), + op: CompareOpKind::Eq, + right: Rc::new(QueryExpr::Column(2)), + })); + let joined = ClickHouseRelationalAdapter + .apply_inner_equi_join(&pred, &joined_schema, left, right) + .unwrap(); + let output_schema = schema(&[("service", DataType::Utf8), ("ratio", DataType::Float64)]); + let projected = ClickHouseRelationalAdapter + .apply_operation( + &ValueOperation::Project { + cols: vec![ + ProjectItem { + alias: Some("service".into()), + expr: QueryExpr::Column(0), + }, + ProjectItem { + alias: Some("ratio".into()), + expr: QueryExpr::Arithmetic { + op: ArithmeticOpKind::Div, + left: Rc::new(QueryExpr::Column(1)), + right: Rc::new(QueryExpr::Column(3)), + }, + }, + ], + qualifier: None, + }, + &output_schema, + joined, + ) + .unwrap(); + assert_eq!(projected.coverage, Some((300_000, 600_000))); + let result = projected.into_result().unwrap(); + let batch = &result.batches[0]; + assert_eq!(batch.schema().field(1).name(), "ratio"); + assert_eq!( + batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap() + .value(0), + 0.2 + ); + } } diff --git a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs index 9914b702..a908595d 100644 --- a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs +++ b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs @@ -288,7 +288,8 @@ impl QueryNodeRuntime for PhysicalQueryRuntime<'_> { QueryPlanNode::Logical { .. } | QueryPlanNode::CandidateTopK { .. } | QueryPlanNode::Relational { .. } - | QueryPlanNode::ExternalExact { .. } => Err(PhysicalNodeError::Fallback( + | QueryPlanNode::ExternalExact { .. } + | QueryPlanNode::RelationalJoin { .. } => Err(PhysicalNodeError::Fallback( "logical node requires installed logical runtime".into(), )), QueryPlanNode::ExactFallback { reason } => { diff --git a/data_plane/src/query_engines/asap_query_engine/summary_exec.rs b/data_plane/src/query_engines/asap_query_engine/summary_exec.rs index 5f8c3029..4ab734ea 100644 --- a/data_plane/src/query_engines/asap_query_engine/summary_exec.rs +++ b/data_plane/src/query_engines/asap_query_engine/summary_exec.rs @@ -244,6 +244,7 @@ pub fn execute( } SummaryExpr::SummaryJoin { .. } => Err(ExecError::NotYetSupported("SummaryJoin")), + SummaryExpr::RelationalJoin { .. } => Err(ExecError::NotYetSupported("RelationalJoin")), // CandidateTopK is lowered to the deployed QueryPlan DAG, where both // row inputs retain labels for intersection and exact reranking. This // legacy generic adapter exposes opaque GroupKey values and cannot From 20734eaf9e4f8b2994b7023fdece0b0908a362fa Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 11:36:34 -0600 Subject: [PATCH 2/2] fix: consume canonical planner lifecycle contract --- Cargo.lock | 16 ++++++++-------- control_plane/Cargo.toml | 10 +++++----- control_plane/src/physical/compiler.rs | 8 ++++---- crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +++--- 5 files changed, 21 insertions(+), 21 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 4c18ff77..96eaa3f9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -374,7 +374,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "asap-types", "serde", @@ -385,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-metricsql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "asap-types", "metricsql_parser", @@ -395,7 +395,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "asap-types", "promql-parser 0.10.0", @@ -404,7 +404,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -426,12 +426,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "serde", "serde_json", @@ -2801,7 +2801,7 @@ dependencies = [ [[package]] name = "metricsql_common" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "chrono", ] @@ -2809,7 +2809,7 @@ dependencies = [ [[package]] name = "metricsql_parser" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=984897fba9bdd82336606a1d3092c54f493aaf76#984897fba9bdd82336606a1d3092c54f493aaf76" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=f99237d#f99237d1643e64db75cd9c9632c2383f272f582e" dependencies = [ "ahash", "chrono", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 8d1b09f5..6a30d061 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -91,8 +91,8 @@ asap_types.workspace = true # scaffolding, unaware that `data_plane`'s `summary_executor.rs` in *this* # repo is a real one. Vendored locally instead of chased upstream -- see # `data_plane/src/query_engines/asap_query_engine/summary_exec.rs`. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } # L1 adoption (design-target-architecture.md Part B): the PromQL front # end itself, replacing control_plane's own query_parser/promql.rs. @@ -100,9 +100,9 @@ asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = # `planner-types`/`asap-aware-mapping` above -- these three MUST move # together (two revs of the same upstream repo's types in one workspace # resolve to distinct Rust types that won't unify). -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } -asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index a75029ef..32e524ef 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -3476,12 +3476,12 @@ fn select_lifecycle( reason: "latest ASAPPlanner selected no window framework from the supplied physical evidence".into(), })?; Ok(PlannerPhysicalSelection { - window_implementation_id: plan.selected_physical_plan_id.clone().ok_or_else(|| { - CompileError::Lifecycle { + window_implementation_id: plan.selected_window_implementation_id.clone().ok_or_else( + || CompileError::Lifecycle { query_id: query.query_id.clone(), reason: "Planner returned no concrete window implementation identity".into(), - } - })?, + }, + )?, expected_reads: plan.expected_reads.ok_or_else(|| CompileError::Lifecycle { query_id: query.query_id.clone(), reason: "missing joint read demand".into(), diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index b968c7ef..24380517 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -32,4 +32,4 @@ sha2 = "0.10" # exactly (`control_plane/Cargo.toml`) -- two different revs of the same # git dependency in one workspace resolve to two distinct Rust types that # won't unify. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 45b77aab..d5621796 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -38,8 +38,8 @@ control_plane = { path = "../control_plane" } # reduction: Reduction, .. }`) are `pre_asap` types, in the same crate now # (not a separate `asap-ir` import). Query serving consumes the compiled # QueryPlan; these types are used at physical-plan compilation boundaries. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } -asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } +asap-frontend-metricsql = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } # Shared external (workspace) serde.workspace = true @@ -130,7 +130,7 @@ crc32fast = "1.4" # none of them. [dev-dependencies] -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "984897fba9bdd82336606a1d3092c54f493aaf76" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "f99237d" } tempfile = "3.20.0" criterion = { version = "0.5", features = ["html_reports"] } tokio-tungstenite = "0.21"