From 3fdde1e36b14ba7bdf3e1d6e1fbbf0ce47f0024b Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 17:31:02 -0600 Subject: [PATCH 1/3] refactor: share installed executable DAG contracts --- control_plane/src/physical/compiler.rs | 13 +- .../src/physical/executable_binding.rs | 292 +---------------- control_plane/src/query_plan.rs | 6 +- crates/asap_types/src/executable_plan.rs | 296 ++++++++++++++++++ crates/asap_types/src/lib.rs | 1 + .../precompute_engine/maintenance_runtime.rs | 2 +- .../src/precompute_engine/subdag_scheduler.rs | 2 +- .../summary-catalog-sds-architecture.md | 10 + 8 files changed, 332 insertions(+), 290 deletions(-) create mode 100644 crates/asap_types/src/executable_plan.rs diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 3ad55df0..fdca6ad8 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -2761,7 +2761,7 @@ impl PhysicalCompiler { } precompute_sinks.sort(); let installed = crate::physical::executable_binding::InstalledPostAsapDag { - document: super::executable_binding::PostAsapDagDocument::from_executable( + document: super::executable_binding::OwnedPostAsapDag::from_executable( query_id.clone(), &compiled.dag, ) @@ -2830,12 +2830,12 @@ impl PhysicalCompiler { query_id: query_id.clone(), reason: "installed post-ASAP DAG has no query-plan entry".into(), })?; - installed - .validate_query_plan(entry) - .map_err(|reason| CompileError::Query { + super::executable_binding::validate_query_plan(installed, entry).map_err(|reason| { + CompileError::Query { query_id: query_id.clone(), reason, - })?; + } + })?; } let summary_catalog = super::summary_catalog::SummaryCatalog::from_materializations( envelope.plan_id, @@ -4402,8 +4402,7 @@ mod tests { .get(&entry.query_id) .expect("compiled query retains its Planner DAG and backend placement"); installed.validate().expect("typed DAG document"); - installed - .validate_query_plan(entry) + crate::physical::executable_binding::validate_query_plan(installed, entry) .expect("query node bindings"); assert_eq!(installed.binding.query_plan_sink, entry.root); assert!(installed.binding.nodes.values().any(|placement| matches!( diff --git a/control_plane/src/physical/executable_binding.rs b/control_plane/src/physical/executable_binding.rs index 8959790c..6aff751f 100644 --- a/control_plane/src/physical/executable_binding.rs +++ b/control_plane/src/physical/executable_binding.rs @@ -1,286 +1,26 @@ -use std::collections::{BTreeMap, BTreeSet}; +//! Compiler checks relating the shared executable contract to QueryPlan. -use asap_types::sds::SummaryDefinitionId; -use planner_types::post_asap::{ - EdgeRole, ExecutableDag, ExecutableDagEdge, ExecutableDagNode, ExecutableOperator, - ExecutionDataState, ExecutionTiming, GroupingEdgeCompatibility, PostAsapNodeId, - WindowEdgeCompatibility, -}; -use serde::{Deserialize, Serialize}; +pub use asap_types::executable_plan::*; -use crate::query_plan::{QueryNodeId, QueryPlanEntry}; - -pub const POST_ASAP_DAG_DOCUMENT_SCHEMA_VERSION: u32 = 1; - -/// Versioned, language-neutral Planner DAG persisted with an installed plan. -/// Plan lifecycle belongs to the enclosing `PrecomputePlan`; this document -/// carries semantic identity only and does not duplicate its envelope. -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(deny_unknown_fields)] -pub struct PostAsapDagDocument { - pub schema_version: u32, - pub query_id: String, - pub nodes: Vec, - pub edges: Vec, - pub root: PostAsapNodeId, -} - -/// Send/Sync wire form of one Planner node. Planner's in-memory payload may -/// contain `Rc`; the canonical JSON payload preserves its tagged type without -/// leaking that process-local ownership choice into installed runtime state. -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(deny_unknown_fields)] -pub struct PostAsapDagNodeDocument { - pub id: PostAsapNodeId, - pub operator: ExecutableOperator, - pub payload: serde_json::Value, - pub output_state: ExecutionDataState, - pub output_schema: serde_json::Value, - pub guarantee: Option, -} - -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(deny_unknown_fields)] -pub struct PostAsapDagEdgeDocument { - pub producer: PostAsapNodeId, - pub consumer: PostAsapNodeId, - pub role: EdgeRole, - pub intermediate_schema: serde_json::Value, - pub data_state: ExecutionDataState, - pub grouping: GroupingEdgeCompatibility, - pub window: WindowEdgeCompatibility, -} - -impl PostAsapDagDocument { - pub fn from_executable(query_id: String, dag: &ExecutableDag) -> Result { - let nodes = dag - .nodes - .iter() - .map(|node| { - Ok(PostAsapDagNodeDocument { - id: node.id, - operator: node.operator, - payload: serde_json::to_value(&node.payload).map_err(|e| e.to_string())?, - output_state: node.output_state, - output_schema: serde_json::to_value(&node.output_schema) - .map_err(|e| e.to_string())?, - guarantee: node - .guarantee - .as_ref() - .map(serde_json::to_value) - .transpose() - .map_err(|e| e.to_string())?, - }) - }) - .collect::, String>>()?; - let edges = dag - .edges - .iter() - .map(|edge| { - Ok(PostAsapDagEdgeDocument { - producer: edge.producer, - consumer: edge.consumer, - role: edge.role, - intermediate_schema: serde_json::to_value(&edge.intermediate_schema) - .map_err(|e| e.to_string())?, - data_state: edge.data_state, - grouping: edge.grouping, - window: edge.window, - }) - }) - .collect::, String>>()?; - Ok(Self { - schema_version: POST_ASAP_DAG_DOCUMENT_SCHEMA_VERSION, - query_id, - nodes, - edges, - root: dag.root, - }) - } - - pub fn decode(&self) -> Result { - let node_ids = self - .nodes - .iter() - .map(|node| node.id) - .collect::>(); - if node_ids.len() != self.nodes.len() { - return Err("duplicate post-ASAP DAG node ID".into()); - } - let edge_ids = self - .edges - .iter() - .map(|edge| { - ( - edge.producer, - edge.consumer, - serde_json::to_string(&edge.role).expect("EdgeRole is serializable"), - ) - }) - .collect::>(); - if edge_ids.len() != self.edges.len() { - return Err("duplicate post-ASAP DAG edge".into()); - } - let nodes = self - .nodes - .iter() - .map(|node| { - let payload: planner_types::post_asap::ExecutableOperatorPayload = - serde_json::from_value(node.payload.clone()).map_err(|e| e.to_string())?; - if payload.operator() != node.operator { - return Err(format!( - "post-ASAP node {} operator disagrees with payload", - node.id.0 - )); - } - Ok(ExecutableDagNode { - id: node.id, - operator: node.operator, - payload, - output_state: node.output_state, - output_schema: serde_json::from_value(node.output_schema.clone()) - .map_err(|e| e.to_string())?, - guarantee: node - .guarantee - .clone() - .map(serde_json::from_value) - .transpose() - .map_err(|e| e.to_string())?, - }) - }) - .collect::, String>>()?; - let edges = self - .edges - .iter() - .map(|edge| { - Ok(ExecutableDagEdge { - producer: edge.producer, - consumer: edge.consumer, - role: edge.role, - intermediate_schema: serde_json::from_value(edge.intermediate_schema.clone()) - .map_err(|e| e.to_string())?, - data_state: edge.data_state, - grouping: edge.grouping, - window: edge.window, - }) - }) - .collect::, String>>()?; - Ok(ExecutableDag { - nodes, - edges, - root: self.root, - }) +pub fn validate_query_plan( + installed: &InstalledPostAsapDag, + query: &crate::query_plan::QueryPlanEntry, +) -> Result<(), String> { + if query.query_id != installed.document.query_id { + return Err("post-ASAP document and query plan identity disagree".into()); } -} - -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -#[serde(deny_unknown_fields)] -pub struct InstalledPostAsapDag { - pub document: PostAsapDagDocument, - pub binding: BackendExecutableBinding, -} - -impl InstalledPostAsapDag { - pub fn validate(&self) -> Result<(), String> { - if self.document.schema_version != POST_ASAP_DAG_DOCUMENT_SCHEMA_VERSION - || self.document.query_id.trim().is_empty() - { - return Err("invalid post-ASAP DAG document identity/version".into()); - } - self.binding.validate(&self.document.decode()?) + if query.root != installed.binding.query_plan_sink { + return Err("backend query-plan sink disagrees with installed query root".into()); } - - pub fn validate_query_plan(&self, query: &QueryPlanEntry) -> Result<(), String> { - if query.query_id != self.document.query_id { - return Err("post-ASAP document and query plan identity disagree".into()); - } - if query.root != self.binding.query_plan_sink { - return Err("backend query-plan sink disagrees with installed query root".into()); - } - for placement in self.binding.nodes.values() { - if let BackendNodeBinding::Query { query_node } = placement { - if !query.nodes.contains_key(query_node) { - return Err(format!( - "backend binding references missing query node {}", - query_node.0 - )); - } - } - } - Ok(()) - } -} - -/// Backend-owned physical identities and sink placement for one Planner -/// semantic DAG. Planner node IDs remain stable semantic references; the -/// backend never assumes they equal its independently allocated query IDs. -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] -#[serde(deny_unknown_fields)] -pub struct BackendExecutableBinding { - pub nodes: BTreeMap, - pub query_sink: PostAsapNodeId, - pub query_plan_sink: QueryNodeId, - pub precompute_sinks: Vec, -} - -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] -#[serde(tag = "placement", rename_all = "snake_case", deny_unknown_fields)] -pub enum BackendNodeBinding { - Query { - query_node: QueryNodeId, - }, - /// Semantic read node absorbed into a larger installed query node, such - /// as an external exact subtree boundary. - QueryInput, - MaintenanceInput, - Materialization { - summary_definition: SummaryDefinitionId, - }, -} - -impl BackendExecutableBinding { - pub fn node(&self, id: PostAsapNodeId) -> Option<&BackendNodeBinding> { - self.nodes.get(&id) - } - - pub fn validate(&self, dag: &ExecutableDag) -> Result<(), String> { - let semantic = dag - .nodes - .iter() - .map(|node| node.id.0) - .collect::>(); - if !semantic.contains(&self.query_sink.0) { - return Err("backend query sink is absent from semantic DAG".into()); - } - if self.nodes.keys().map(|id| id.0).collect::>() != semantic { - return Err("backend executable binding does not cover semantic DAG exactly".into()); - } - for node in &dag.nodes { - match (node.output_state.timing, self.node(node.id).unwrap()) { - (ExecutionTiming::ReadTime, BackendNodeBinding::Query { .. }) - | (ExecutionTiming::ReadTime, BackendNodeBinding::QueryInput) - | (ExecutionTiming::MaintenanceTime, BackendNodeBinding::MaintenanceInput) - | (ExecutionTiming::MaintenanceTime, BackendNodeBinding::Materialization { .. }) => { - } - _ => { - return Err(format!( - "backend placement disagrees with node {} mode", - node.id.0 - )) - } - } - } - for sink in &self.precompute_sinks { - if !matches!( - self.node(*sink), - Some(BackendNodeBinding::Materialization { .. }) - ) { + for placement in installed.binding.nodes.values() { + if let BackendNodeBinding::Query { query_node } = placement { + if !query.nodes.contains_key(query_node) { return Err(format!( - "backend precompute sink {} is not materialized", - sink.0 + "backend binding references missing query node {}", + query_node.0 )); } } - Ok(()) } + Ok(()) } diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index 1540f81e..ff9d0e78 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -205,11 +205,7 @@ impl QueryPlan { } } -/// Stable identity inside one query entry. Edges are IDs so common -/// subexpressions remain shared after serialization. -#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord, Hash)] -#[serde(transparent)] -pub struct QueryNodeId(pub u64); +pub use asap_types::executable_plan::QueryNodeId; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(deny_unknown_fields)] diff --git a/crates/asap_types/src/executable_plan.rs b/crates/asap_types/src/executable_plan.rs new file mode 100644 index 00000000..fe59854a --- /dev/null +++ b/crates/asap_types/src/executable_plan.rs @@ -0,0 +1,296 @@ +//! Shared backend installation contract for semantic DAGs and physical bindings. +//! +//! `OwnedPostAsapDag` is the Send/Sync installed representation of Planner IR; +//! it is distinct from Planner's `PostAsapDagDocument` transport envelope. +//! Planner payloads contain process-local `Rc` values, so they are decoded only +//! when executing or validating a DAG. Compilation and QueryPlan cross-checks +//! remain control-plane responsibilities. + +use std::collections::{BTreeMap, BTreeSet}; + +use crate::sds::SummaryDefinitionId; +use planner_types::post_asap::{ + EdgeRole, ExecutableDag, ExecutableDagEdge, ExecutableDagNode, ExecutableOperator, + ExecutionDataState, ExecutionTiming, GroupingEdgeCompatibility, PostAsapNodeId, + WindowEdgeCompatibility, +}; +use serde::{Deserialize, Serialize}; + +/// Identity within an installed query entry, separate from semantic Planner IDs. +#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord, Hash)] +#[serde(transparent)] +pub struct QueryNodeId(pub u64); + +pub const OWNED_POST_ASAP_DAG_SCHEMA_VERSION: u32 = 1; + +/// Versioned, language-neutral Planner DAG persisted with an installed plan. +/// Plan lifecycle belongs to the enclosing `PrecomputePlan`; this document +/// carries semantic identity only and does not duplicate its envelope. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct OwnedPostAsapDag { + pub schema_version: u32, + pub query_id: String, + pub nodes: Vec, + pub edges: Vec, + pub root: PostAsapNodeId, +} + +/// Send/Sync wire form of one Planner node. Planner's in-memory payload may +/// contain `Rc`; the canonical JSON payload preserves its tagged type without +/// leaking that process-local ownership choice into installed runtime state. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct OwnedPostAsapNode { + pub id: PostAsapNodeId, + pub operator: ExecutableOperator, + pub payload: serde_json::Value, + pub output_state: ExecutionDataState, + pub output_schema: serde_json::Value, + pub guarantee: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct OwnedPostAsapEdge { + pub producer: PostAsapNodeId, + pub consumer: PostAsapNodeId, + pub role: EdgeRole, + pub intermediate_schema: serde_json::Value, + pub data_state: ExecutionDataState, + pub grouping: GroupingEdgeCompatibility, + pub window: WindowEdgeCompatibility, +} + +impl OwnedPostAsapDag { + pub fn from_executable(query_id: String, dag: &ExecutableDag) -> Result { + let nodes = dag + .nodes + .iter() + .map(|node| { + Ok(OwnedPostAsapNode { + id: node.id, + operator: node.operator, + payload: serde_json::to_value(&node.payload).map_err(|e| e.to_string())?, + output_state: node.output_state, + output_schema: serde_json::to_value(&node.output_schema) + .map_err(|e| e.to_string())?, + guarantee: node + .guarantee + .as_ref() + .map(serde_json::to_value) + .transpose() + .map_err(|e| e.to_string())?, + }) + }) + .collect::, String>>()?; + let edges = dag + .edges + .iter() + .map(|edge| { + Ok(OwnedPostAsapEdge { + producer: edge.producer, + consumer: edge.consumer, + role: edge.role, + intermediate_schema: serde_json::to_value(&edge.intermediate_schema) + .map_err(|e| e.to_string())?, + data_state: edge.data_state, + grouping: edge.grouping, + window: edge.window, + }) + }) + .collect::, String>>()?; + Ok(Self { + schema_version: OWNED_POST_ASAP_DAG_SCHEMA_VERSION, + query_id, + nodes, + edges, + root: dag.root, + }) + } + + pub fn decode(&self) -> Result { + let node_ids = self + .nodes + .iter() + .map(|node| node.id) + .collect::>(); + if node_ids.len() != self.nodes.len() { + return Err("duplicate post-ASAP DAG node ID".into()); + } + let edge_ids = self + .edges + .iter() + .map(|edge| { + ( + edge.producer, + edge.consumer, + serde_json::to_string(&edge.role).expect("EdgeRole is serializable"), + ) + }) + .collect::>(); + if edge_ids.len() != self.edges.len() { + return Err("duplicate post-ASAP DAG edge".into()); + } + let nodes = self + .nodes + .iter() + .map(|node| { + let payload: planner_types::post_asap::ExecutableOperatorPayload = + serde_json::from_value(node.payload.clone()).map_err(|e| e.to_string())?; + if payload.operator() != node.operator { + return Err(format!( + "post-ASAP node {} operator disagrees with payload", + node.id.0 + )); + } + Ok(ExecutableDagNode { + id: node.id, + operator: node.operator, + payload, + output_state: node.output_state, + output_schema: serde_json::from_value(node.output_schema.clone()) + .map_err(|e| e.to_string())?, + guarantee: node + .guarantee + .clone() + .map(serde_json::from_value) + .transpose() + .map_err(|e| e.to_string())?, + }) + }) + .collect::, String>>()?; + let edges = self + .edges + .iter() + .map(|edge| { + Ok(ExecutableDagEdge { + producer: edge.producer, + consumer: edge.consumer, + role: edge.role, + intermediate_schema: serde_json::from_value(edge.intermediate_schema.clone()) + .map_err(|e| e.to_string())?, + data_state: edge.data_state, + grouping: edge.grouping, + window: edge.window, + }) + }) + .collect::, String>>()?; + Ok(ExecutableDag { + nodes, + edges, + root: self.root, + }) + } +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct InstalledPostAsapDag { + pub document: OwnedPostAsapDag, + pub binding: BackendExecutableBinding, +} + +impl InstalledPostAsapDag { + pub fn validate(&self) -> Result<(), String> { + if self.document.schema_version != OWNED_POST_ASAP_DAG_SCHEMA_VERSION + || self.document.query_id.trim().is_empty() + { + return Err("invalid post-ASAP DAG document identity/version".into()); + } + self.binding.validate(&self.document.decode()?) + } +} + +/// Backend-owned physical identities and sink placement for one Planner +/// semantic DAG. Planner node IDs remain stable semantic references; the +/// backend never assumes they equal its independently allocated query IDs. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct BackendExecutableBinding { + pub nodes: BTreeMap, + pub query_sink: PostAsapNodeId, + pub query_plan_sink: QueryNodeId, + pub precompute_sinks: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(tag = "placement", rename_all = "snake_case", deny_unknown_fields)] +pub enum BackendNodeBinding { + Query { + query_node: QueryNodeId, + }, + /// Semantic read node absorbed into a larger installed query node, such + /// as an external exact subtree boundary. + QueryInput, + MaintenanceInput, + Materialization { + summary_definition: SummaryDefinitionId, + }, +} + +impl BackendExecutableBinding { + pub fn node(&self, id: PostAsapNodeId) -> Option<&BackendNodeBinding> { + self.nodes.get(&id) + } + + pub fn validate(&self, dag: &ExecutableDag) -> Result<(), String> { + let semantic = dag + .nodes + .iter() + .map(|node| node.id.0) + .collect::>(); + if !semantic.contains(&self.query_sink.0) { + return Err("backend query sink is absent from semantic DAG".into()); + } + if self.nodes.keys().map(|id| id.0).collect::>() != semantic { + return Err("backend executable binding does not cover semantic DAG exactly".into()); + } + for node in &dag.nodes { + match (node.output_state.timing, self.node(node.id).unwrap()) { + (ExecutionTiming::ReadTime, BackendNodeBinding::Query { .. }) + | (ExecutionTiming::ReadTime, BackendNodeBinding::QueryInput) + | (ExecutionTiming::MaintenanceTime, BackendNodeBinding::MaintenanceInput) + | (ExecutionTiming::MaintenanceTime, BackendNodeBinding::Materialization { .. }) => { + } + _ => { + return Err(format!( + "backend placement disagrees with node {} mode", + node.id.0 + )) + } + } + } + for sink in &self.precompute_sinks { + if !matches!( + self.node(*sink), + Some(BackendNodeBinding::Materialization { .. }) + ) { + return Err(format!( + "backend precompute sink {} is not materialized", + sink.0 + )); + } + } + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + // Shared installation metadata must be safe to retain in cross-thread + // ActivePhysicalPlan snapshots without importing the compiler crate. + #[test] + fn installed_contract_is_send_sync_and_preserves_wire_identity() { + fn send_sync() {} + send_sync::(); + let wire = serde_json::json!({ + "schema_version": 1, "query_id": "q", "nodes": [], "edges": [], "root": 0 + }); + let document: OwnedPostAsapDag = serde_json::from_value(wire.clone()).unwrap(); + assert_eq!(serde_json::to_value(document).unwrap(), wire); + assert_eq!(serde_json::to_value(QueryNodeId(9)).unwrap(), 9); + } +} diff --git a/crates/asap_types/src/lib.rs b/crates/asap_types/src/lib.rs index d70625de..f1d9f3dd 100644 --- a/crates/asap_types/src/lib.rs +++ b/crates/asap_types/src/lib.rs @@ -3,6 +3,7 @@ pub mod accuracy; pub mod aggregation_config; pub mod aggregation_type; pub mod enums; +pub mod executable_plan; pub mod key_by_label_names; pub mod monitor_spec; pub mod policy_fingerprint; diff --git a/data_plane/src/precompute_engine/maintenance_runtime.rs b/data_plane/src/precompute_engine/maintenance_runtime.rs index 18db59e3..78951851 100644 --- a/data_plane/src/precompute_engine/maintenance_runtime.rs +++ b/data_plane/src/precompute_engine/maintenance_runtime.rs @@ -6,7 +6,7 @@ use super::subdag_scheduler::{ PrecomputeOperatorRegistry, ScheduleError, }; use crate::storage_engines::types::{AggregateCore, HotReloadStreamingConfig, PrecomputedOutput}; -use control_plane::physical::executable_binding::{BackendExecutableBinding, BackendNodeBinding}; +use asap_types::executable_plan::{BackendExecutableBinding, BackendNodeBinding}; use planner_types::post_asap::{ExecutableDagNode, ExecutableOperatorPayload, PostAsapNodeId}; use std::collections::{BTreeMap, BTreeSet}; use std::sync::{Arc, Mutex}; diff --git a/data_plane/src/precompute_engine/subdag_scheduler.rs b/data_plane/src/precompute_engine/subdag_scheduler.rs index 2d0036f6..3a803853 100644 --- a/data_plane/src/precompute_engine/subdag_scheduler.rs +++ b/data_plane/src/precompute_engine/subdag_scheduler.rs @@ -1,4 +1,4 @@ -use control_plane::physical::executable_binding::{BackendExecutableBinding, BackendNodeBinding}; +use asap_types::executable_plan::{BackendExecutableBinding, BackendNodeBinding}; use planner_types::post_asap::PostAsapNodeId; use planner_types::post_asap::{ExecutableDag, ExecutableDagNode, ExecutionDataState}; use std::{ diff --git a/docs/design_docs/summary-catalog-sds-architecture.md b/docs/design_docs/summary-catalog-sds-architecture.md index e900c36f..44131f62 100644 --- a/docs/design_docs/summary-catalog-sds-architecture.md +++ b/docs/design_docs/summary-catalog-sds-architecture.md @@ -148,6 +148,16 @@ readout and fallback routing, and the common deployment envelope carries their shared plan identity. Consumers atomically install one catalog snapshot with the plans that reference it. +`asap_types::executable_plan` owns the installed semantic-DAG representation, +physical node bindings, and `QueryNodeId`. Its `OwnedPostAsapDag` is a Send/Sync +representation for shared runtime snapshots; it is not Planner's +`PostAsapDagDocument` envelope. The owned representation preserves semantic +node IDs and typed operator tags while serializing Planner payloads that contain +process-local `Rc` pointers. The control plane constructs it and checks its +bindings against QueryPlan; precompute execution consumes the shared contract. +`QueryPlan` and `PrecomputePlan` definitions still reside in the control-plane +crate while their remaining compilation methods are separated from wire types. + The implemented ownership split is: 1. Move the SDS catalog contract into `asap_types`. From 310b23f3e191856572013201be6c60dc8341a766 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 18:06:09 -0600 Subject: [PATCH 2/3] test: use shared DAG name in maintenance fixtures --- data_plane/src/precompute_engine/maintenance_runtime.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/data_plane/src/precompute_engine/maintenance_runtime.rs b/data_plane/src/precompute_engine/maintenance_runtime.rs index 6cd34585..6acc9f25 100644 --- a/data_plane/src/precompute_engine/maintenance_runtime.rs +++ b/data_plane/src/precompute_engine/maintenance_runtime.rs @@ -697,7 +697,7 @@ mod tests { ActivePhysicalPlan, HotReloadActivePhysicalPlan, StreamingConfig, }; use control_plane::physical::executable_binding::{ - InstalledPostAsapDag, PostAsapDagDocument, + InstalledPostAsapDag, OwnedPostAsapDag, }; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -774,7 +774,7 @@ mod tests { bundle.precompute_plan.executable_dags = BTreeMap::from([( "retry".into(), InstalledPostAsapDag { - document: PostAsapDagDocument::from_executable("retry".into(), &dag).unwrap(), + document: OwnedPostAsapDag::from_executable("retry".into(), &dag).unwrap(), binding, }, )]); From c1cc416d0db372cb6a54d6df9c0b83f05f0938bf Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 18:11:01 -0600 Subject: [PATCH 3/3] style: format shared DAG fixture import --- data_plane/src/precompute_engine/maintenance_runtime.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/data_plane/src/precompute_engine/maintenance_runtime.rs b/data_plane/src/precompute_engine/maintenance_runtime.rs index 6acc9f25..d5e8b6f7 100644 --- a/data_plane/src/precompute_engine/maintenance_runtime.rs +++ b/data_plane/src/precompute_engine/maintenance_runtime.rs @@ -696,9 +696,7 @@ mod tests { use crate::storage_engines::types::{ ActivePhysicalPlan, HotReloadActivePhysicalPlan, StreamingConfig, }; - use control_plane::physical::executable_binding::{ - InstalledPostAsapDag, OwnedPostAsapDag, - }; + use control_plane::physical::executable_binding::{InstalledPostAsapDag, OwnedPostAsapDag}; use std::sync::atomic::{AtomicUsize, Ordering}; #[derive(Default)]