From 5941911d1a5a1785746f8dfc84d164cadf630b91 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 08:05:08 -0600 Subject: [PATCH 1/2] Publish authoritative catalog plans directly --- control_plane/src/backend_client.rs | 57 -------------------------- control_plane/src/emit/backend_push.rs | 40 +++++++++++++----- control_plane/src/main.rs | 15 +++---- 3 files changed, 38 insertions(+), 74 deletions(-) diff --git a/control_plane/src/backend_client.rs b/control_plane/src/backend_client.rs index 2f6a798d..95c5fa43 100644 --- a/control_plane/src/backend_client.rs +++ b/control_plane/src/backend_client.rs @@ -365,63 +365,6 @@ impl BackendClient { } } - /// Publish the backend-facing portions of one PhysicalPlan in a single - /// request, preventing independently retried documents from mixing - /// generations at the backend. - pub async fn post_physical_plan_typed( - &self, - precompute_plan: &crate::physical::compiler::PrecomputePlan, - transmission_plan: &crate::physical::compiler::TransmissionPlan, - query_plan: &crate::query_plan::QueryPlan, - storage_routing: Option, - adaptation_evidence: &[crate::physical::compiler::RuntimeAdaptationEvidence], - ) -> std::result::Result<(), BackendPostError> { - // Compatibility replanner has no Planner-selected query/collector DAG. - // Still publish the actual catalog and bind every provided projection. - let catalog = crate::physical::summary_catalog::SummaryCatalog::from_materializations( - precompute_plan.envelope.plan_id, - precompute_plan.envelope.plan_version, - &precompute_plan.materializations, - ) - .map_err(|error| BackendPostError::Permanent(error.into()))?; - let mut precompute_plan = precompute_plan.clone(); - precompute_plan - .bind_catalog(&catalog) - .map_err(|error| BackendPostError::Permanent(error.into()))?; - let mut transmission_plan = transmission_plan.clone(); - transmission_plan.summary_catalog = Some( - catalog - .reference() - .map_err(|error| BackendPostError::Permanent(error.into()))?, - ); - query_plan - .validate_against_catalog(&catalog) - .map_err(|error| BackendPostError::Permanent(error.into()))?; - let url = derive_physical_plan_url(&self.endpoint); - let response = self - .http - .post(&url) - .json(&serde_json::json!({ - "summary_catalog": catalog, - "collector_plans": [], - "precompute_plan": precompute_plan, - "transmission_plan": transmission_plan, - "query_plan": query_plan, - "storage_routing": storage_routing, - "adaptation_evidence": adaptation_evidence, - })) - .send() - .await - .map_err(classify_reqwest_error)?; - let status = response.status(); - if status.is_success() { - Ok(()) - } else { - let body = response.text().await.unwrap_or_default(); - Err(classify_http_status(status, body, "PhysicalPlan POST")) - } - } - pub async fn discard_staged_physical_plan( &self, plan_id: u64, diff --git a/control_plane/src/emit/backend_push.rs b/control_plane/src/emit/backend_push.rs index feafed81..940ff578 100644 --- a/control_plane/src/emit/backend_push.rs +++ b/control_plane/src/emit/backend_push.rs @@ -244,9 +244,25 @@ async fn push_documents_coupled( clickhouse_context: None, entries: Default::default(), }; + let catalog = match crate::physical::summary_catalog::SummaryCatalog::from_materializations( + precompute_plan.envelope.plan_id, + precompute_plan.envelope.plan_version, + &precompute_plan.materializations, + ) { + Ok(catalog) => catalog, + Err(error) => { + warn!(%error, "failed to build compatibility SummaryCatalog"); + return (false, false, 0); + } + }; + let mut precompute_plan = precompute_plan.clone(); + if let Err(error) = precompute_plan.bind_catalog(&catalog) { + warn!(%error, "failed to bind compatibility PrecomputePlan to SummaryCatalog"); + return (false, false, 0); + } let transmission_plan = match crate::physical::compiler::TransmissionPlan::build( precompute_plan.envelope.clone(), - precompute_plan, + &precompute_plan, &Default::default(), ) { Ok(plan) => plan, @@ -255,16 +271,17 @@ async fn push_documents_coupled( return (false, false, 0); } }; + let publication = crate::physical::publication::PhysicalPlanPublication { + summary_catalog: catalog, + precompute_plan, + collector_plans: Vec::new(), + transmission_plan, + query_plan, + }; for attempt in 1..=RETRY_MAX_ATTEMPTS { match client - .post_physical_plan_typed( - precompute_plan, - &transmission_plan, - &query_plan, - Some(routing.clone()), - &[], - ) + .post_catalog_plan_typed(&publication, Some(routing.clone()), &[]) .await { Ok(()) => return (true, true, attempt), @@ -456,10 +473,13 @@ async fn push_cumulative_entries( return PushOutcome::EmitFailed; } }; - let precompute_plan = match PrecomputePlan::build( + // This compatibility path installs state in the backend-local precompute + // engine. It must not manufacture a distributed collector producer: an + // authoritative publication requires every producer to have a matching + // CollectorPlan, and no collector exists on this path. + let precompute_plan = match PrecomputePlan::build_backend_local( precompute_envelope, materializations, - &["legacy-replanner".into()], ) { Ok(plan) => plan, Err(error) => { diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index 0d79d28a..8c282bb1 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -822,14 +822,15 @@ async fn handle_compile_and_publish_clickhouse_plan( ) .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_physical_plan_typed( - &bundle.precompute_plan, - &bundle.transmission_plan, - &bundle.query_plan, - None, - &[], - ) + .post_catalog_plan_typed(&publication, None, &[]) .await { return (StatusCode::BAD_GATEWAY, error.to_string()).into_response(); From 2d05cfb65a91da4fb3c54d2c3fd05d968e3f79f9 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 08:29:14 -0600 Subject: [PATCH 2/2] Format authoritative publication adapter --- control_plane/src/emit/backend_push.rs | 18 ++++++++---------- 1 file changed, 8 insertions(+), 10 deletions(-) diff --git a/control_plane/src/emit/backend_push.rs b/control_plane/src/emit/backend_push.rs index 940ff578..df8033ea 100644 --- a/control_plane/src/emit/backend_push.rs +++ b/control_plane/src/emit/backend_push.rs @@ -477,16 +477,14 @@ async fn push_cumulative_entries( // engine. It must not manufacture a distributed collector producer: an // authoritative publication requires every producer to have a matching // CollectorPlan, and no collector exists on this path. - let precompute_plan = match PrecomputePlan::build_backend_local( - precompute_envelope, - materializations, - ) { - Ok(plan) => plan, - Err(error) => { - warn!(%error, "failed to validate typed PrecomputePlan"); - return PushOutcome::EmitFailed; - } - }; + let precompute_plan = + match PrecomputePlan::build_backend_local(precompute_envelope, materializations) { + Ok(plan) => plan, + Err(error) => { + warn!(%error, "failed to validate typed PrecomputePlan"); + return PushOutcome::EmitFailed; + } + }; // Storage-routing: the routing classifier (`build_routing_entry` in // `emit/stage_config.rs`) reads `cfg.aggregations` to derive shape