From 1833c77a75286885e993870e49da222de206fa0c Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Tue, 15 Sep 2026 07:58:31 -0500 Subject: [PATCH 1/4] feat: add an optional `plan_name` discriminator to `PhysicalExtensionNode` MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Additive and proto3-compatible: on the binary wire format old writers omit the field and old readers skip it. The generated JSON codec is stricter — `pbjson` rejects a field it does not know — so a named node written by a new writer cannot be read as JSON by an older reader. An unnamed node is unaffected, and nothing writes a name yet. Nothing reads it either; the codec path keeps writing `None`, since there the codec remains the discriminator. Co-Authored-By: Claude Fable 5.1 Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_017UzVkTrKmRRc2fKCdMkqBp --- .../adapter_serialization.rs | 5 +++++ datafusion/proto-models/proto/datafusion.proto | 6 ++++++ .../proto-models/src/generated/pbjson.rs | 18 ++++++++++++++++++ datafusion/proto-models/src/generated/prost.rs | 7 +++++++ datafusion/proto/src/physical_plan/mod.rs | 8 +++++++- 5 files changed, 43 insertions(+), 1 deletion(-) diff --git a/datafusion-examples/examples/custom_data_source/adapter_serialization.rs b/datafusion-examples/examples/custom_data_source/adapter_serialization.rs index a4e5b06ec1137..127deaf3e27e3 100644 --- a/datafusion-examples/examples/custom_data_source/adapter_serialization.rs +++ b/datafusion-examples/examples/custom_data_source/adapter_serialization.rs @@ -364,6 +364,11 @@ impl PhysicalProtoConverterExtension for AdapterPreservingCodec { PhysicalExtensionNode { node: payload_bytes, inputs: vec![inner_proto], + // This converter recognizes its own nodes by the marker + // inside the payload, so it needs no `plan_name` + // discriminator. Plans that decode through the session + // registry set it to their `PLAN_NAME` instead. + plan_name: None, }, )), }); diff --git a/datafusion/proto-models/proto/datafusion.proto b/datafusion/proto-models/proto/datafusion.proto index 04b9464c9be66..dac9b8a1bf81a 100644 --- a/datafusion/proto-models/proto/datafusion.proto +++ b/datafusion/proto-models/proto/datafusion.proto @@ -1045,6 +1045,12 @@ message ListUnnest { message PhysicalExtensionNode { bytes node = 1; repeated PhysicalPlanNode inputs = 2; + // Globally unique, namespaced name of the extension plan type (e.g. + // `my-crate.MyExec`), written by plans that serialize themselves through + // `ExecutionPlan::try_to_proto`. When present and registered on the decoding + // session, it selects the decoder directly; otherwise decoding falls back to + // the `PhysicalExtensionCodec` chain. Absent for codec-encoded plans. + optional string plan_name = 3; } // physical expressions diff --git a/datafusion/proto-models/src/generated/pbjson.rs b/datafusion/proto-models/src/generated/pbjson.rs index c3b5d1d3740cb..78a1d2d340440 100644 --- a/datafusion/proto-models/src/generated/pbjson.rs +++ b/datafusion/proto-models/src/generated/pbjson.rs @@ -20121,6 +20121,9 @@ impl serde::Serialize for PhysicalExtensionNode { if !self.inputs.is_empty() { len += 1; } + if self.plan_name.is_some() { + len += 1; + } let mut struct_ser = serializer.serialize_struct("datafusion.PhysicalExtensionNode", len)?; if !self.node.is_empty() { #[allow(clippy::needless_borrow)] @@ -20130,6 +20133,9 @@ impl serde::Serialize for PhysicalExtensionNode { if !self.inputs.is_empty() { struct_ser.serialize_field("inputs", &self.inputs)?; } + if let Some(v) = self.plan_name.as_ref() { + struct_ser.serialize_field("planName", v)?; + } struct_ser.end() } } @@ -20142,12 +20148,15 @@ impl<'de> serde::Deserialize<'de> for PhysicalExtensionNode { const FIELDS: &[&str] = &[ "node", "inputs", + "plan_name", + "planName", ]; #[allow(clippy::enum_variant_names)] enum GeneratedField { Node, Inputs, + PlanName, } impl<'de> serde::Deserialize<'de> for GeneratedField { fn deserialize(deserializer: D) -> std::result::Result @@ -20171,6 +20180,7 @@ impl<'de> serde::Deserialize<'de> for PhysicalExtensionNode { match value { "node" => Ok(GeneratedField::Node), "inputs" => Ok(GeneratedField::Inputs), + "planName" | "plan_name" => Ok(GeneratedField::PlanName), _ => Err(serde::de::Error::unknown_field(value, FIELDS)), } } @@ -20192,6 +20202,7 @@ impl<'de> serde::Deserialize<'de> for PhysicalExtensionNode { { let mut node__ = None; let mut inputs__ = None; + let mut plan_name__ = None; while let Some(k) = map_.next_key()? { match k { GeneratedField::Node => { @@ -20208,11 +20219,18 @@ impl<'de> serde::Deserialize<'de> for PhysicalExtensionNode { } inputs__ = Some(map_.next_value()?); } + GeneratedField::PlanName => { + if plan_name__.is_some() { + return Err(serde::de::Error::duplicate_field("planName")); + } + plan_name__ = map_.next_value()?; + } } } Ok(PhysicalExtensionNode { node: node__.unwrap_or_default(), inputs: inputs__.unwrap_or_default(), + plan_name: plan_name__, }) } } diff --git a/datafusion/proto-models/src/generated/prost.rs b/datafusion/proto-models/src/generated/prost.rs index ec01eb1e67998..894e2aa08cc79 100644 --- a/datafusion/proto-models/src/generated/prost.rs +++ b/datafusion/proto-models/src/generated/prost.rs @@ -1619,6 +1619,13 @@ pub struct PhysicalExtensionNode { pub node: ::prost::alloc::vec::Vec, #[prost(message, repeated, tag = "2")] pub inputs: ::prost::alloc::vec::Vec, + /// Globally unique, namespaced name of the extension plan type (e.g. + /// `my-crate.MyExec`), written by plans that serialize themselves through + /// `ExecutionPlan::try_to_proto`. When present and registered on the decoding + /// session, it selects the decoder directly; otherwise decoding falls back to + /// the `PhysicalExtensionCodec` chain. Absent for codec-encoded plans. + #[prost(string, optional, tag = "3")] + pub plan_name: ::core::option::Option<::prost::alloc::string::String>, } /// physical expressions #[derive(Clone, PartialEq, ::prost::Message)] diff --git a/datafusion/proto/src/physical_plan/mod.rs b/datafusion/proto/src/physical_plan/mod.rs index 2b349035b3292..464d5a2aa7666 100644 --- a/datafusion/proto/src/physical_plan/mod.rs +++ b/datafusion/proto/src/physical_plan/mod.rs @@ -1381,7 +1381,13 @@ pub trait PhysicalPlanNodeExt: Sized { Ok(protobuf::PhysicalPlanNode { physical_plan_type: Some(PhysicalPlanType::Extension( - protobuf::PhysicalExtensionNode { node: buf, inputs }, + protobuf::PhysicalExtensionNode { + node: buf, + inputs, + // Codec-encoded plans stay anonymous on the wire: + // the codec that wrote them decodes them back. + plan_name: None, + }, )), }) } From 18b810ff105beece82f8cc49afe0dfc66888c0b1 Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Wed, 16 Sep 2026 23:04:36 -0500 Subject: [PATCH 2/4] feat: add `ProtoDecoderRegistry`, one decoder store for every extension kind DataFusion serializes several kinds of polymorphic value: `ExecutionPlan`, `PhysicalExpr`, and in time the `DataSource`, `DataSink` and `LazyBatchGenerator` families. A built-in gets its own wire variant, so the wire format names it. An extension shares one catch-all variant with every other extension of its kind, so the payload carries a name and the decoding session maps that name back to a decoder. `ProtoDecoderRegistry` is that map, once, for every kind. The key is a pair: the trait the decoder produces and the name on the wire. Keying on the trait as well as the name means a session carries one registry rather than one per kind, so an application composes every library's registrations into a single object and attaches it once. Two kinds may also use the same name without a collision. The registry never names a decode context or any trait it dispatches to: it stores each decoder type-erased and hands it back on a downcast. That is what lets it live in `datafusion-proto-models`, below every crate that owns one of those traits. Each kind adds a small typed facade next to its own trait, so a caller never writes a `TypeId` or a downcast by hand. Nothing uses it yet; the extension `ExecutionPlan` facade is the next commit. Co-Authored-By: Claude Opus 5 (1M context) --- datafusion/proto-models/src/lib.rs | 6 + datafusion/proto-models/src/registry.rs | 326 ++++++++++++++++++++++++ 2 files changed, 332 insertions(+) create mode 100644 datafusion/proto-models/src/registry.rs diff --git a/datafusion/proto-models/src/lib.rs b/datafusion/proto-models/src/lib.rs index f5469cbfec16e..4a8e7b7ee2bb2 100644 --- a/datafusion/proto-models/src/lib.rs +++ b/datafusion/proto-models/src/lib.rs @@ -34,6 +34,10 @@ //! itself — see [`from_proto`] and [`to_proto`]. It is the schema source of //! truth for [`datafusion-proto`]. //! +//! It also hosts [`ProtoDecoderRegistry`](registry::ProtoDecoderRegistry), the +//! one store of extension decoders every serializable kind shares, for the same +//! layering reason: it sits below every crate that owns one of those traits. +//! //! Most users should depend on [`datafusion-proto`] instead, which re-exports //! these types under [`datafusion_proto::protobuf`]. //! @@ -43,6 +47,7 @@ pub mod from_proto; pub mod generated; +pub mod registry; pub mod to_proto; /// All DataFusion protobuf model types. @@ -57,6 +62,7 @@ pub mod protobuf { /// Re-export of the `datafusion_proto_common` types as exposed through this /// crate's generated module, for callers that want the common-only namespace. pub use generated::datafusion_common; +pub use registry::ProtoDecoderRegistry; #[cfg(all(test, feature = "json"))] mod tests { diff --git a/datafusion/proto-models/src/registry.rs b/datafusion/proto-models/src/registry.rs new file mode 100644 index 0000000000000..e2852e57c682d --- /dev/null +++ b/datafusion/proto-models/src/registry.rs @@ -0,0 +1,326 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! A name-keyed store of decoders for extension types. +//! +//! DataFusion serializes several kinds of polymorphic value: `ExecutionPlan`, +//! `PhysicalExpr`, and in time the `DataSource`, `DataSink` and +//! `LazyBatchGenerator` families. A built-in gets its own wire variant, so the +//! wire format names it. An extension shares one catch-all variant with every +//! other extension of its kind, so the payload has to carry a name and the +//! decoding session has to map that name back to a decoder. +//! +//! [`ProtoDecoderRegistry`] is that map, once, for every kind. The key is a +//! pair: the trait the decoder produces, and the name on the wire. Keying on +//! the trait as well as the name means one session carries one registry rather +//! than one per kind, and two kinds may use the same name without a collision. +//! +//! The registry never names a decode context, or any trait it dispatches to. +//! It stores each decoder type-erased and hands it back on a downcast, which is +//! what lets it sit in this crate, below every crate that owns one of those +//! traits. Each kind supplies a small typed facade next to its own trait — +//! `datafusion-physical-plan` for `ExecutionPlan`, and so on — so that a caller +//! never writes a `TypeId` or a downcast by hand. + +use std::any::{Any, TypeId, type_name}; +use std::collections::HashMap; +use std::collections::hash_map::Entry; +use std::fmt; +use std::sync::Arc; + +use datafusion_common::{Result, config_err}; + +/// One registered decoder, plus the identity that makes re-registering the +/// same type a no-op while a genuine name collision is an error. +#[derive(Clone)] +struct Registered { + decoder: Arc, + type_id: TypeId, + type_name: &'static str, +} + +/// A store of extension decoders, keyed by the trait they produce and the name +/// the encoder wrote on the wire. +/// +/// Build one, let every library that ships extension types register into it, +/// and attach it to the decoding session with `SessionConfig::with_extension`. +/// A session carries at most one, and attaching another replaces it, so the +/// application that owns the session is the one that builds it. A library +/// exposes a function that fills a registry it is handed: +/// +/// ```rust,ignore +/// // in each library +/// pub fn register(registry: &mut ProtoDecoderRegistry) -> Result<()> { +/// register_execution_plan::(registry)?; +/// register_physical_expr::(registry) +/// } +/// +/// // in the application that owns the session +/// let mut registry = ProtoDecoderRegistry::new(); +/// lib_one::register(&mut registry)?; +/// lib_two::register(&mut registry)?; +/// let config = SessionConfig::new().with_extension(Arc::new(registry)); +/// ``` +/// +/// The libraries need to know nothing about each other, and the order they are +/// called in does not matter: lookup is by name, and a real collision is an +/// error from [`register_decoder`](Self::register_decoder). +/// +/// Callers normally reach this type through a kind's typed facade rather than +/// through the methods here. +#[derive(Clone, Default)] +pub struct ProtoDecoderRegistry { + decoders: HashMap<(TypeId, String), Registered>, +} + +impl ProtoDecoderRegistry { + /// Create an empty registry. + pub fn new() -> Self { + Self::default() + } + + /// Register `decoder` under `name`, for extensions dispatched as `Kind`. + /// + /// `Kind` is the trait object the decoder produces, such as + /// `dyn ExecutionPlan`. `T` is the concrete type being registered; it is + /// used only for identity, so that registering the same type twice is a + /// no-op while a *different* type under a taken name is an error. `D` is + /// the decoder itself, whose type the matching + /// [`decoder`](Self::decoder) call must name exactly. + pub fn register_decoder(&mut self, name: &str, decoder: D) -> Result<()> + where + Kind: ?Sized + 'static, + T: ?Sized + 'static, + D: Any + Send + Sync, + { + let registered = Registered { + decoder: Arc::new(decoder), + type_id: TypeId::of::(), + type_name: type_name::(), + }; + if name.is_empty() { + return config_err!( + "Cannot register the extension {} decoder for {} under an empty name", + kind_label::(), + registered.type_name + ); + } + match self + .decoders + .entry((TypeId::of::(), name.to_string())) + { + Entry::Vacant(entry) => { + entry.insert(registered); + Ok(()) + } + // Re-registering the same type is a no-op: sessions are often + // configured by more than one layer of an application. + Entry::Occupied(entry) if entry.get().type_id == registered.type_id => Ok(()), + Entry::Occupied(entry) => config_err!( + "Extension {} name '{}' is already registered by {}, cannot register {}. \ + Namespace the name with the owning crate to avoid the collision.", + kind_label::(), + entry.key().1, + entry.get().type_name, + registered.type_name + ), + } + } + + /// The decoder registered under `name` for extensions dispatched as `Kind`. + /// + /// `D` must be the same type the matching + /// [`register_decoder`](Self::register_decoder) call supplied. A different + /// `D` reads as "not registered" rather than as an error, because the two + /// calls are always made by the same facade. + pub fn decoder(&self, name: &str) -> Option<&D> + where + Kind: ?Sized + 'static, + D: Any + Send + Sync, + { + let registered = self + .decoders + .get(&(TypeId::of::(), name.to_string()))?; + registered.decoder.downcast_ref::() + } + + /// Every name registered for `Kind`, in arbitrary order. + pub fn names(&self) -> impl Iterator { + let kind = TypeId::of::(); + self.decoders + .keys() + .filter(move |(registered_kind, _)| *registered_kind == kind) + .map(|(_, name)| name.as_str()) + } +} + +impl fmt::Debug for ProtoDecoderRegistry { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + // The decoders are opaque; the registered types are the useful part. + let mut entries: Vec<&str> = + self.decoders.values().map(|r| r.type_name).collect(); + entries.sort_unstable(); + f.debug_struct("ProtoDecoderRegistry") + .field("registered", &entries) + .finish() + } +} + +/// The last path segment of a trait object's type name, for error messages: +/// `dyn datafusion_physical_plan::execution_plan::ExecutionPlan` reads as +/// `ExecutionPlan`. +fn kind_label() -> &'static str { + let name = type_name::(); + name.rsplit("::").next().unwrap_or(name) +} + +#[cfg(test)] +mod tests { + use super::*; + + trait Plan {} + trait Expr {} + + struct OnePlan; + struct OtherPlan; + + type PlanDecoder = fn(u8) -> u8; + type ExprDecoder = fn(u8) -> u16; + + fn plan_decoder(v: u8) -> u8 { + v + } + fn other_plan_decoder(v: u8) -> u8 { + v + 1 + } + fn expr_decoder(v: u8) -> u16 { + u16::from(v) + 1000 + } + + #[test] + fn register_and_look_up_a_decoder() -> Result<()> { + let mut registry = ProtoDecoderRegistry::new(); + registry.register_decoder::( + "a-crate.OnePlan", + plan_decoder, + )?; + + let decoder = registry + .decoder::("a-crate.OnePlan") + .expect("the decoder must be registered"); + assert_eq!(decoder(7), 7); + assert!( + registry + .decoder::("a-crate.Missing") + .is_none() + ); + Ok(()) + } + + #[test] + fn the_same_name_under_two_kinds_does_not_collide() -> Result<()> { + // The point of keying on the trait as well as the name: one session + // carries one registry, and two kinds never see each other's names. + let mut registry = ProtoDecoderRegistry::new(); + registry + .register_decoder::("shared", plan_decoder)?; + registry.register_decoder::( + "shared", + expr_decoder, + )?; + + assert_eq!( + registry.decoder::("shared").unwrap()(1), + 1 + ); + assert_eq!( + registry.decoder::("shared").unwrap()(1), + 1001 + ); + assert_eq!( + registry.names::().collect::>(), + vec!["shared"] + ); + assert_eq!( + registry.names::().collect::>(), + vec!["shared"] + ); + Ok(()) + } + + #[test] + fn register_is_idempotent_for_the_same_type() -> Result<()> { + let mut registry = ProtoDecoderRegistry::new(); + registry + .register_decoder::("one", plan_decoder)?; + registry + .register_decoder::("one", plan_decoder)?; + + assert_eq!(registry.names::().count(), 1); + Ok(()) + } + + #[test] + fn register_rejects_a_name_collision_between_types() -> Result<()> { + let mut registry = ProtoDecoderRegistry::new(); + registry + .register_decoder::("one", plan_decoder)?; + + let err = registry + .register_decoder::( + "one", + other_plan_decoder, + ) + .expect_err("a colliding registration must fail"); + let err = err.to_string(); + assert!( + err.contains("already registered"), + "unexpected error: {err}" + ); + assert!(err.contains("Plan"), "the kind must be named: {err}"); + assert!( + err.contains("OnePlan") && err.contains("OtherPlan"), + "{err}" + ); + Ok(()) + } + + #[test] + fn register_rejects_an_empty_name() { + let mut registry = ProtoDecoderRegistry::new(); + let err = registry + .register_decoder::("", plan_decoder) + .expect_err("an empty name must fail"); + assert!( + err.to_string().contains("empty name"), + "unexpected error: {err}" + ); + assert_eq!(registry.names::().count(), 0); + } + + #[test] + fn a_decoder_of_another_type_reads_as_absent() -> Result<()> { + let mut registry = ProtoDecoderRegistry::new(); + registry + .register_decoder::("one", plan_decoder)?; + + // A facade always pairs the same `D` with the same `Kind`; a mismatch + // means the caller asked the wrong facade. + assert!(registry.decoder::("one").is_none()); + Ok(()) + } +} From 5e7c50e32556613afd5b85552bea602406942f6a Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Tue, 15 Sep 2026 08:02:33 -0500 Subject: [PATCH 3/4] feat: decode extension `ExecutionPlan`s through a session-scoped registry MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `PhysicalExtensionNode` carried no type discriminator, so the codec *was* the discriminator: `ComposedPhysicalExtensionCodec` writes the position of the encoding codec into the payload, the decoding side must register the same codecs in the same order, and a name collision between two crates is undetectable. An extension plan now implements `ExtensionPlanFromProto`, which pairs the name it is dispatched by with the constructor that rebuilds it. Its `try_to_proto` writes itself with `ExecutionPlanEncodeCtx::extension_node`, which stamps `Self::NAME`, so the encoded name and the registry key cannot drift apart. `register_execution_plan` puts its decoder in the `ProtoDecoderRegistry` the decoding session carries. Only extension plans implement the trait. A built-in has a `PhysicalPlanType` variant of its own, keeps the inherent `try_from_proto` it has had since 55.0.0, and has no wire name to be registered under — so `register_execution_plan::(..)` does not compile. Nothing about built-in serialization changes, and no caller needs a new import. This crate supplies only the `ExecutionPlan` face of the registry — `register_execution_plan`, `decode_execution_plan` and `execution_plan_names`, keyed on `dyn ExecutionPlan`. The store is shared with every other extension kind, so a session carries one registry and an application composes every library's registrations into it. Decode rule: a named node the session has a decoder for goes to that decoder, and a failure there is fatal (falling through would let a codec decode a payload it never wrote); anything else takes the `PhysicalExtensionCodec` chain exactly as before. If the codec fails too, the missing registration is added as *context* on the codec's own error, so the kind the codec chose still reaches a caller that matches on it. `decode_execution_plan` reads the name off the node rather than taking it as an argument, so a caller cannot pair a node with a name it does not carry. Decoders are stored as a private `fn` pointer so the storage can become a `dyn` decoder object (what an FFI decoder needs) without a break; registration is keyed by `TypeId` (the same type twice is a no-op, a different type under a taken name is an error). Adds `proto/extension_plan_registry.rs` to the examples: the same two-crate tree `composed_extension_codec.rs` builds, decoded by name with no composed codec. Documented in the 56.0.0 upgrade guide; tests in `datafusion/proto/tests/cases/plans/extensions.rs`. Co-Authored-By: Claude Fable 5.1 Co-Authored-By: Claude Opus 5 (1M context) --- datafusion-examples/README.md | 1 + .../examples/proto/extension_plan_registry.rs | 252 ++++++ datafusion-examples/examples/proto/main.rs | 10 +- .../src/{proto.rs => proto/mod.rs} | 140 +++- .../physical-plan/src/proto/registry.rs | 275 +++++++ .../physical-plan/src/proto_test_util.rs | 147 ++++ datafusion/proto-models/src/lib.rs | 2 +- datafusion/proto/src/physical_plan/mod.rs | 86 +- .../proto/tests/cases/plans/extensions.rs | 775 ++++++++++++++++++ datafusion/proto/tests/cases/plans/mod.rs | 1 + .../library-user-guide/upgrading/56.0.0.md | 86 ++ 11 files changed, 1763 insertions(+), 12 deletions(-) create mode 100644 datafusion-examples/examples/proto/extension_plan_registry.rs rename datafusion/physical-plan/src/{proto.rs => proto/mod.rs} (74%) create mode 100644 datafusion/physical-plan/src/proto/registry.rs create mode 100644 datafusion/proto/tests/cases/plans/extensions.rs diff --git a/datafusion-examples/README.md b/datafusion-examples/README.md index 6c7aad31a9233..2137fb30d508f 100644 --- a/datafusion-examples/README.md +++ b/datafusion-examples/README.md @@ -171,6 +171,7 @@ cargo run --example dataframe -- dataframe | ------------------------ | --------------------------------------------------------------------------------- | ----------------------------------------------------------------------------- | | composed_extension_codec | [`proto/composed_extension_codec.rs`](examples/proto/composed_extension_codec.rs) | Use multiple extension codecs for serialization/deserialization | | expression_deduplication | [`proto/expression_deduplication.rs`](examples/proto/expression_deduplication.rs) | Example of expression caching/deduplication using the codec decorator pattern | +| extension_plan_registry | [`proto/extension_plan_registry.rs`](examples/proto/extension_plan_registry.rs) | Decode two crates' extension plans by name, with no composed codec | ## Query Planning Examples diff --git a/datafusion-examples/examples/proto/extension_plan_registry.rs b/datafusion-examples/examples/proto/extension_plan_registry.rs new file mode 100644 index 0000000000000..c50823be20913 --- /dev/null +++ b/datafusion-examples/examples/proto/extension_plan_registry.rs @@ -0,0 +1,252 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! See `main.rs` for how to run it. +//! +//! The same problem as `composed_extension_codec.rs`, solved without composing +//! codecs: two plans from two independent crates in one tree. +//! +//! ```text +//! ParentExec (from crate A) +//! ChildExec (from crate B) +//! ``` +//! +//! With `ComposedPhysicalExtensionCodec` the decoding side has to register the +//! same codecs in the same *order* the encoder used, because the composed codec +//! writes the position of the encoding codec into the payload. Here each plan +//! names itself on the wire instead, so order does not matter and a name two +//! crates both claim is an error at registration rather than a wrong decode. + +use std::fmt::Formatter; +use std::sync::Arc; + +use datafusion::common::tree_node::TreeNodeRecursion; +use datafusion::common::{Result, internal_err}; +use datafusion::execution::TaskContext; +use datafusion::physical_expr::PhysicalExpr; +use datafusion::physical_plan::proto::{ + ExecutionPlanDecodeCtx, ExecutionPlanEncodeCtx, ExtensionPlanFromProto, + ProtoDecoderRegistry, register_execution_plan, +}; +use datafusion::physical_plan::{ + ChildrenPropertiesMode, DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties, + ReplaceChildrenOptions, SendableRecordBatchStream, +}; +use datafusion::prelude::{SessionConfig, SessionContext}; +use datafusion_proto::physical_plan::{AsExecutionPlan, DefaultPhysicalExtensionCodec}; +use datafusion_proto::protobuf::PhysicalPlanNode; +use datafusion_proto::protobuf::physical_plan_node::PhysicalPlanType; + +pub fn extension_plan_registry() -> Result<()> { + let plan: Arc = Arc::new(ParentExec { + input: Arc::new(ChildExec {}), + }); + + // Each crate exposes a function that fills a registry it is handed; the + // application that owns the session composes them. The two crates need to + // know nothing about each other, and the call order does not matter. + let mut registry = ProtoDecoderRegistry::new(); + crate_a::register(&mut registry)?; + crate_b::register(&mut registry)?; + + let ctx = SessionContext::new_with_config( + SessionConfig::new().with_extension(Arc::new(registry)), + ); + + // `DefaultPhysicalExtensionCodec` decodes nothing at all, so whatever + // survives the round trip came from the registry. + let codec = DefaultPhysicalExtensionCodec {}; + let node = PhysicalPlanNode::try_from_physical_plan(Arc::clone(&plan), &codec)?; + let decoded = node.try_into_physical_plan(ctx.task_ctx().as_ref(), &codec)?; + + assert_eq!(format!("{plan:?}"), format!("{decoded:?}")); + + // A session that does not register the plans cannot decode them, and the + // error names what is missing rather than only what the codec said. Children + // are decoded before their parent falls back to the codec, so the innermost + // unregistered plan is the one reported: + // + // No decoder is registered for the extension ExecutionPlan + // 'example-crate-b.ChildExec'. Register the plan in the + // ProtoDecoderRegistry attached to the decoding session's + // SessionConfig, ... Registered extension plans: none + let bare = SessionContext::new(); + let err = node + .try_into_physical_plan(bare.task_ctx().as_ref(), &codec) + .expect_err("an unregistered plan must not decode"); + assert!(err.to_string().contains("No decoder is registered")); + assert!(err.to_string().contains(ChildExec::NAME)); + + Ok(()) +} + +/// What the crate that owns `ParentExec` would export. +mod crate_a { + use super::*; + + pub fn register(registry: &mut ProtoDecoderRegistry) -> Result<()> { + register_execution_plan::(registry) + } +} + +/// What the crate that owns `ChildExec` would export. +mod crate_b { + use super::*; + + pub fn register(registry: &mut ProtoDecoderRegistry) -> Result<()> { + register_execution_plan::(registry) + } +} + +#[derive(Debug)] +struct ParentExec { + input: Arc, +} + +impl ExtensionPlanFromProto for ParentExec { + // Namespaced with the owning crate, so a collision with another crate is + // an error at registration rather than a wrong decode. + const NAME: &'static str = "example-crate-a.ParentExec"; + + fn try_from_proto( + node: &PhysicalPlanNode, + ctx: &ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + let Some(PhysicalPlanType::Extension(extension)) = &node.physical_plan_type + else { + return internal_err!("PhysicalPlanNode is not an Extension"); + }; + let mut children = ctx.decode_children(&extension.inputs)?; + if children.len() != 1 { + return internal_err!("ParentExec expects exactly one child"); + } + Ok(Arc::new(ParentExec { + input: children.remove(0), + })) + } +} + +#[derive(Debug)] +struct ChildExec {} + +impl ExtensionPlanFromProto for ChildExec { + const NAME: &'static str = "example-crate-b.ChildExec"; + + fn try_from_proto( + _node: &PhysicalPlanNode, + _ctx: &ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + Ok(Arc::new(ChildExec {})) + } +} + +/// Both plans are serde-only: this example never executes them. +macro_rules! serde_only_body { + ($name:ident) => { + fn name(&self) -> &str { + stringify!($name) + } + + fn properties(&self) -> &Arc { + unreachable!("this example only serializes") + } + + fn apply_expressions( + &self, + _f: &mut dyn FnMut(&Arc) -> Result, + ) -> Result { + Ok(TreeNodeRecursion::Continue) + } + + fn with_new_children( + self: Arc, + children: Vec>, + ) -> Result> { + self.replace_children( + children, + ReplaceChildrenOptions::new(ChildrenPropertiesMode::Recompute), + ) + } + + fn execute( + &self, + _partition: usize, + _context: Arc, + ) -> Result { + internal_err!("{} is a serde-only example plan", stringify!($name)) + } + + /// `extension_node` stamps `Self::NAME`, so the name on the wire and + /// the registry key cannot drift apart. Neither plan has state of its + /// own, hence the empty payload. + fn try_to_proto( + &self, + ctx: &ExecutionPlanEncodeCtx<'_>, + ) -> Result> { + Ok(Some(ctx.extension_node::(vec![], self.children())?)) + } + }; +} + +impl DisplayAs for ParentExec { + fn fmt_as(&self, _t: DisplayFormatType, f: &mut Formatter) -> std::fmt::Result { + write!(f, "ParentExec") + } +} + +impl ExecutionPlan for ParentExec { + serde_only_body!(ParentExec); + + fn children(&self) -> Vec<&Arc> { + vec![&self.input] + } + + fn replace_children( + self: Arc, + mut children: Vec>, + _: ReplaceChildrenOptions, + ) -> Result> { + if children.len() != 1 { + return internal_err!("ParentExec expects exactly one child"); + } + Ok(Arc::new(ParentExec { + input: children.remove(0), + })) + } +} + +impl DisplayAs for ChildExec { + fn fmt_as(&self, _t: DisplayFormatType, f: &mut Formatter) -> std::fmt::Result { + write!(f, "ChildExec") + } +} + +impl ExecutionPlan for ChildExec { + serde_only_body!(ChildExec); + + fn children(&self) -> Vec<&Arc> { + vec![] + } + + fn replace_children( + self: Arc, + _children: Vec>, + _: ReplaceChildrenOptions, + ) -> Result> { + Ok(self) + } +} diff --git a/datafusion-examples/examples/proto/main.rs b/datafusion-examples/examples/proto/main.rs index d534eda24ba64..c3b6c22701d2b 100644 --- a/datafusion-examples/examples/proto/main.rs +++ b/datafusion-examples/examples/proto/main.rs @@ -21,7 +21,7 @@ //! //! ## Usage //! ```bash -//! cargo run --example proto -- [all|composed_extension_codec|expression_deduplication] +//! cargo run --example proto -- [all|composed_extension_codec|expression_deduplication|extension_plan_registry] //! ``` //! //! Each subcommand runs a corresponding example: @@ -32,9 +32,13 @@ //! //! - `expression_deduplication` //! (file: expression_deduplication.rs, desc: Example of expression caching/deduplication using the codec decorator pattern) +//! +//! - `extension_plan_registry` +//! (file: extension_plan_registry.rs, desc: Decode two crates' extension plans by name, with no composed codec) mod composed_extension_codec; mod expression_deduplication; +mod extension_plan_registry; use datafusion::error::{DataFusionError, Result}; use strum::{IntoEnumIterator, VariantNames}; @@ -46,6 +50,7 @@ enum ExampleKind { All, ComposedExtensionCodec, ExpressionDeduplication, + ExtensionPlanRegistry, } impl ExampleKind { @@ -69,6 +74,9 @@ impl ExampleKind { ExampleKind::ExpressionDeduplication => { expression_deduplication::expression_deduplication()? } + ExampleKind::ExtensionPlanRegistry => { + extension_plan_registry::extension_plan_registry()? + } } Ok(()) } diff --git a/datafusion/physical-plan/src/proto.rs b/datafusion/physical-plan/src/proto/mod.rs similarity index 74% rename from datafusion/physical-plan/src/proto.rs rename to datafusion/physical-plan/src/proto/mod.rs index 7640d76c3e010..b3fae8df50b59 100644 --- a/datafusion/physical-plan/src/proto.rs +++ b/datafusion/physical-plan/src/proto/mod.rs @@ -56,8 +56,21 @@ //! `physical-expr-common`, *below* `datafusion-expr`) cannot do this, which is //! why `ScalarFunctionExpr` remains special-cased there. //! +//! # Extension plans +//! +//! Third-party plans have no `PhysicalPlanType` variant of their own; they all +//! share `PhysicalExtensionNode`, whose payload is opaque bytes. Such a plan +//! rides the same hooks: its `try_to_proto` builds a `PhysicalExtensionNode` +//! through [`ExecutionPlanEncodeCtx::extension_node`], which stamps the node +//! with the plan's [`ExtensionPlanFromProto::NAME`], and its `try_from_proto` reads +//! that node back. The name lets the decoding session select the decoder +//! directly rather than by trying every registered `PhysicalExtensionCodec` in +//! turn. See [`registry`]. +//! //! [`ExecutionPlan`]: crate::ExecutionPlan +pub mod registry; + use std::sync::Arc; use arrow::datatypes::Schema; @@ -72,10 +85,19 @@ use datafusion_physical_expr_common::physical_expr::proto_decode::{ use datafusion_physical_expr_common::physical_expr::proto_encode::{ PhysicalExprEncode, PhysicalExprEncodeCtx, }; -use datafusion_proto_models::protobuf::{PhysicalExprNode, PhysicalPlanNode}; +use datafusion_proto_models::protobuf::physical_plan_node::PhysicalPlanType; +use datafusion_proto_models::protobuf::{ + PhysicalExprNode, PhysicalExtensionNode, PhysicalPlanNode, +}; use crate::ExecutionPlan; +pub use datafusion_proto_models::ProtoDecoderRegistry; +pub use registry::{ + ExtensionPlanFromProto, decode_execution_plan, execution_plan_names, + register_execution_plan, +}; + /// Internal dispatch trait backing [`ExecutionPlanEncodeCtx`]. /// /// Implemented by `datafusion-proto`. Plan authors never name this trait; they @@ -187,6 +209,41 @@ impl<'a> ExecutionPlanEncodeCtx<'a> { plans.into_iter().map(|p| self.encode_child(p)).collect() } + /// Build the `PhysicalPlanNode` for an extension plan `T`: an + /// `Extension` variant carrying `payload`, the encoded `children`, and + /// `T`'s [`NAME`](ExtensionPlanFromProto::NAME). + /// + /// The one supported way for an extension plan to write itself. The name + /// on the wire is taken from the same const the registry is keyed by, so + /// an encoder cannot write a name no decoder answers to: + /// + /// ```ignore + /// fn try_to_proto( + /// &self, + /// ctx: &ExecutionPlanEncodeCtx<'_>, + /// ) -> Result> { + /// Ok(Some(ctx.extension_node::( + /// my_payload_bytes(self)?, + /// self.children(), + /// )?)) + /// } + /// ``` + pub fn extension_node( + &self, + payload: Vec, + children: Vec<&Arc>, + ) -> Result { + Ok(PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::Extension( + PhysicalExtensionNode { + node: payload, + inputs: self.encode_children(children)?, + plan_name: Some(T::NAME.to_string()), + }, + )), + }) + } + /// Serialize a single physical expression. pub fn encode_expr(&self, expr: &Arc) -> Result { self.encoder.encode_expr(expr) @@ -260,6 +317,14 @@ impl<'a> ExecutionPlanDecodeCtx<'a> { self.decoder.decode_plan(node) } + /// Deserialize a slice of child plans. + pub fn decode_children( + &self, + nodes: &[PhysicalPlanNode], + ) -> Result>> { + nodes.iter().map(|node| self.decode_child(node)).collect() + } + /// Deserialize a child plan with `results` active for scalar subquery /// expressions in that plan's subtree. pub fn decode_child_with_scalar_subquery_results( @@ -384,3 +449,76 @@ macro_rules! expect_plan_variant { } }}; } + +#[cfg(test)] +mod tests { + use super::*; + use crate::proto_test_util::registry_test_plan::RegisteredExec; + use crate::proto_test_util::{ + StubPlanDecoder, StubPlanEncoder, encoded_child_node, stub_child, + }; + use datafusion_proto_models::protobuf::PhysicalExtensionNode; + use datafusion_proto_models::protobuf::physical_plan_node::PhysicalPlanType; + + /// The wire node a `RegisteredExec` with one child encodes to. + fn encode(plan: &RegisteredExec, encoder: &StubPlanEncoder) -> PhysicalExtensionNode { + let node = plan + .try_to_proto(&ExecutionPlanEncodeCtx::new(encoder)) + .expect("encode must succeed") + .expect("an extension plan must serialize itself"); + match node.physical_plan_type { + Some(PhysicalPlanType::Extension(extension)) => extension, + other => panic!("expected an extension node, got {other:?}"), + } + } + + #[test] + fn an_extension_plan_stamps_its_name_and_recurses_into_children() { + let encoder = StubPlanEncoder::ok(); + let extension = encode( + &RegisteredExec::new("payload", vec![stub_child()]), + &encoder, + ); + + assert_eq!(extension.plan_name.as_deref(), Some(RegisteredExec::NAME)); + assert_eq!(extension.node, b"payload"); + // The child rode the central serializer rather than being dropped. + assert_eq!(encoder.plan_calls(), 1); + assert_eq!(extension.inputs, vec![encoded_child_node()]); + } + + #[test] + fn an_extension_plan_propagates_a_child_failure() { + let encoder = StubPlanEncoder::failing_on_plan(1); + let err = RegisteredExec::new("payload", vec![stub_child()]) + .try_to_proto(&ExecutionPlanEncodeCtx::new(&encoder)) + .expect_err("a failing child encode must fail the plan"); + assert!( + err.to_string().contains("stub plan encode failure"), + "unexpected error: {err}" + ); + } + + #[test] + fn extension_plan_round_trips_through_the_hooks() { + let encoder = StubPlanEncoder::ok(); + let extension = encode( + &RegisteredExec::new("payload", vec![stub_child()]), + &encoder, + ); + let node = PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::Extension(extension)), + }; + + let decoder = StubPlanDecoder::ok(); + let decoded = + RegisteredExec::try_from_proto(&node, &ExecutionPlanDecodeCtx::new(&decoder)) + .expect("decode must succeed"); + + let decoded = decoded + .downcast_ref::() + .expect("decoded plan must be a RegisteredExec"); + assert_eq!(decoded.payload, "payload"); + assert_eq!(decoded.children.len(), 1); + } +} diff --git a/datafusion/physical-plan/src/proto/registry.rs b/datafusion/physical-plan/src/proto/registry.rs new file mode 100644 index 0000000000000..b604b9a91c5f0 --- /dev/null +++ b/datafusion/physical-plan/src/proto/registry.rs @@ -0,0 +1,275 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Name-keyed decoding for extension [`ExecutionPlan`]s. +//! +//! A built-in plan has a `PhysicalPlanType` variant of its own, so the wire +//! format names it. Extension plans all share `PhysicalExtensionNode`, which +//! carried no discriminator, so the `PhysicalExtensionCodec` *was* the +//! discriminator — and `ComposedPhysicalExtensionCodec` resolves by +//! registration order, which two independent crates cannot agree on. +//! +//! An extension plan implements [`ExtensionPlanFromProto`] instead, and +//! [`register_execution_plan`] puts its decoder in the +//! [`ProtoDecoderRegistry`] the decoding session carries. Lookup is by name, +//! so registration order does not matter. Built-ins cannot be registered; +//! anything unnamed or unregistered still takes the codec chain, unchanged. +//! +//! This is only the `ExecutionPlan` face of the registry: the store is shared +//! with every other extension kind, keyed by the trait as well as the name. +//! `datafusion-examples/examples/proto/extension_plan_registry.rs` is a +//! worked example. + +use std::sync::Arc; + +use datafusion_common::Result; +use datafusion_proto_models::ProtoDecoderRegistry; +use datafusion_proto_models::protobuf::PhysicalPlanNode; +use datafusion_proto_models::protobuf::physical_plan_node::PhysicalPlanType; + +use crate::ExecutionPlan; +use crate::proto::ExecutionPlanDecodeCtx; + +/// The wire name of an extension [`ExecutionPlan`], and the constructor that +/// rebuilds it. +/// +/// Only extension plans implement this, so a built-in cannot be registered by +/// mistake. One impl block carries both halves, and +/// [`ExecutionPlan::try_to_proto`] writes the matching node: +/// +/// ```ignore +/// impl ExecutionPlan for MyExec { +/// fn try_to_proto( +/// &self, +/// ctx: &ExecutionPlanEncodeCtx<'_>, +/// ) -> Result> { +/// // Stamps `Self::NAME`, so the encoded name and the registry key +/// // cannot drift apart. +/// Ok(Some(ctx.extension_node::(self.payload()?, self.children())?)) +/// } +/// } +/// +/// impl ExtensionPlanFromProto for MyExec { +/// const NAME: &'static str = "my-crate.MyExec"; +/// +/// fn try_from_proto( +/// node: &PhysicalPlanNode, +/// ctx: &ExecutionPlanDecodeCtx<'_>, +/// ) -> Result> { +/// let extension = expect_plan_variant!(node, PhysicalPlanType::Extension, "Extension"); +/// let children = ctx.decode_children(&extension.inputs)?; +/// my_plan_from_bytes(&extension.node, children) +/// } +/// } +/// ``` +pub trait ExtensionPlanFromProto: ExecutionPlan + Sized { + /// The name this plan type is written and registered under. + /// + /// Namespace it with the owning crate (`"my-crate.MyExec"`) so that a + /// collision between two independent crates surfaces as a registration + /// error rather than as a wrong decode. Never write it by hand on the + /// wire: [`ExecutionPlanEncodeCtx::extension_node`] stamps it for you. + /// + /// [`ExecutionPlanEncodeCtx::extension_node`]: crate::proto::ExecutionPlanEncodeCtx::extension_node + const NAME: &'static str; + + /// Rebuild the plan from the `PhysicalPlanNode` its + /// [`ExecutionPlan::try_to_proto`] wrote. + /// + /// Match the `Extension` variant (`expect_plan_variant!`), read the + /// payload from `node`, and decode `inputs` with + /// [`ExecutionPlanDecodeCtx::decode_children`]. `ctx` also carries the + /// decoding session, so a plan that rebuilds session state at decode time + /// can do so. + fn try_from_proto( + node: &PhysicalPlanNode, + ctx: &ExecutionPlanDecodeCtx<'_>, + ) -> Result>; +} + +/// How this facade stores a decoder in the shared registry: a function pointer +/// to the monomorphized [`ExtensionPlanFromProto::try_from_proto`]. +/// +/// Deliberately private, and the same type on both the +/// [`register_execution_plan`] and the [`decode_execution_plan`] side. +/// [`ExtensionPlanFromProto`] is the public contract and +/// [`decode_execution_plan`] is the public way to invoke one, so this can +/// become something else — a `dyn` decoder object, to admit stateful or +/// closure decoders, which is what an FFI decoder needs — without a breaking +/// change. +type ExecutionPlanDecoder = + fn(&PhysicalPlanNode, &ExecutionPlanDecodeCtx<'_>) -> Result>; + +/// Register `T` in `registry` under its [`ExtensionPlanFromProto::NAME`]. +/// +/// Registering the same type twice is a no-op; a *different* type under a +/// taken name is an error, so collisions surface here rather than as a wrong +/// decode. The name is scoped to `dyn ExecutionPlan`, so an expression or a +/// data source may use the same name in the same registry. +/// +/// The application that owns the session builds the registry and attaches it +/// once with `SessionConfig::with_extension`; a library exposes a function +/// that fills a registry it is handed. See [`ProtoDecoderRegistry`]. +/// +/// A built-in plan is not an extension and cannot be registered: +/// +/// ```compile_fail +/// use datafusion_physical_plan::filter::FilterExec; +/// use datafusion_physical_plan::proto::{ProtoDecoderRegistry, register_execution_plan}; +/// +/// let mut registry = ProtoDecoderRegistry::new(); +/// register_execution_plan::(&mut registry).unwrap(); +/// ``` +/// +/// The same imports without that call compile, so the failure above is the +/// missing bound and not a bad path: +/// +/// ``` +/// use datafusion_physical_plan::filter::FilterExec; +/// #[allow(unused_imports)] +/// use datafusion_physical_plan::proto::{ProtoDecoderRegistry, register_execution_plan}; +/// +/// let mut registry = ProtoDecoderRegistry::new(); +/// let _ = (FilterExec::try_from_proto, &mut registry); +/// ``` +pub fn register_execution_plan( + registry: &mut ProtoDecoderRegistry, +) -> Result<()> { + registry.register_decoder::( + T::NAME, + T::try_from_proto, + ) +} + +/// Decode `node` with the extension plan decoder registered under the name +/// `node` carries. +/// +/// The name is read from the node's `Extension` variant rather than passed in, +/// so a caller cannot pair a node with a name it does not carry. +/// +/// `None` means "this node names no extension plan decoder of ours": it is not +/// an extension node, it carries no name, or no registered name matches. The +/// caller then falls back to the `PhysicalExtensionCodec` chain. +/// `Some(Err(..))` means the decoder that *does* own the name failed, which is +/// fatal: falling back there would let another codec decode the payload +/// wrongly, the very thing the name exists to prevent. +pub fn decode_execution_plan( + registry: &ProtoDecoderRegistry, + node: &PhysicalPlanNode, + ctx: &ExecutionPlanDecodeCtx<'_>, +) -> Option>> { + let name = plan_name(node)?; + let decoder = registry.decoder::(name)?; + Some(decoder(node, ctx)) +} + +/// Every extension plan name registered in `registry`, in arbitrary order. +/// +/// Names registered for another kind are not included. +pub fn execution_plan_names( + registry: &ProtoDecoderRegistry, +) -> impl Iterator { + registry.names::() +} + +/// The registry name `node` was written with, if it is an extension node that +/// carries one. +fn plan_name(node: &PhysicalPlanNode) -> Option<&str> { + match node.physical_plan_type.as_ref()? { + PhysicalPlanType::Extension(extension) => extension.plan_name.as_deref(), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::proto_test_util::registry_test_plan::{ + OtherRegisteredExec, RegisteredExec, RegisteredExecClone, UnnamedExec, + }; + + #[test] + fn register_is_idempotent_for_the_same_type() -> Result<()> { + let mut registry = ProtoDecoderRegistry::new(); + register_execution_plan::(&mut registry)?; + register_execution_plan::(&mut registry)?; + + assert_eq!(execution_plan_names(®istry).count(), 1); + assert!(execution_plan_names(®istry).any(|name| name == RegisteredExec::NAME)); + Ok(()) + } + + #[test] + fn register_rejects_a_name_collision_between_types() -> Result<()> { + let mut registry = ProtoDecoderRegistry::new(); + register_execution_plan::(&mut registry)?; + + // A distinct type declaring the same `NAME`: the collision two + // independent crates could hit, caught at registration. + let err = register_execution_plan::(&mut registry) + .expect_err("colliding registration must fail"); + assert!( + err.to_string().contains(RegisteredExec::NAME) + && err.to_string().contains("already registered"), + "unexpected error: {err}" + ); + Ok(()) + } + + #[test] + fn register_keeps_distinct_names_apart() -> Result<()> { + let mut registry = ProtoDecoderRegistry::new(); + register_execution_plan::(&mut registry)?; + register_execution_plan::(&mut registry)?; + + assert_eq!(execution_plan_names(®istry).count(), 2); + let mut names: Vec<&str> = execution_plan_names(®istry).collect(); + names.sort_unstable(); + assert_eq!(names, vec![OtherRegisteredExec::NAME, RegisteredExec::NAME]); + assert!(!execution_plan_names(®istry).any(|name| name == "not.registered")); + Ok(()) + } + + #[test] + fn register_rejects_an_empty_plan_name() { + let mut registry = ProtoDecoderRegistry::new(); + let err = register_execution_plan::(&mut registry) + .expect_err("an empty NAME must fail"); + assert!( + err.to_string().contains("empty name"), + "unexpected error: {err}" + ); + assert!(execution_plan_names(®istry).next().is_none()); + } + + #[test] + fn a_name_registered_for_another_kind_is_not_an_execution_plan() -> Result<()> { + // One session carries one registry, so the `ExecutionPlan` face must + // not see a name another kind registered, and must not collide with it. + trait OtherKind {} + let mut registry = ProtoDecoderRegistry::new(); + registry.register_decoder::(RegisteredExec::NAME, 7)?; + + assert!(execution_plan_names(®istry).next().is_none()); + register_execution_plan::(&mut registry)?; + assert_eq!( + execution_plan_names(®istry).collect::>(), + vec![RegisteredExec::NAME] + ); + Ok(()) + } +} diff --git a/datafusion/physical-plan/src/proto_test_util.rs b/datafusion/physical-plan/src/proto_test_util.rs index c1f27ad7fe6d2..a3d2d3bef9a67 100644 --- a/datafusion/physical-plan/src/proto_test_util.rs +++ b/datafusion/physical-plan/src/proto_test_util.rs @@ -320,6 +320,153 @@ impl ExecutionPlanDecode for StubPlanDecoder { } } +/// Minimal extension plans for the [`crate::proto::registry`] tests. +/// +/// Deliberately tiny: they exist to exercise name registration, collision +/// detection and the extension-node encode / decode pair, not to +/// execute. +pub(crate) mod registry_test_plan { + use super::stub_schema; + use crate::execution_plan::{Boundedness, EmissionType}; + use crate::proto::{ + ExecutionPlanDecodeCtx, ExecutionPlanEncodeCtx, ExtensionPlanFromProto, + }; + use crate::{ + DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties, + SendableRecordBatchStream, + }; + use datafusion_common::tree_node::TreeNodeRecursion; + use datafusion_common::{Result, internal_err}; + use datafusion_execution::TaskContext; + use datafusion_physical_expr::{EquivalenceProperties, Partitioning}; + use datafusion_proto_models::protobuf::PhysicalPlanNode; + use std::fmt::Formatter; + use std::sync::Arc; + + fn stub_properties() -> Arc { + Arc::new(PlanProperties::new( + EquivalenceProperties::new(stub_schema()), + Partitioning::UnknownPartitioning(1), + EmissionType::Incremental, + Boundedness::Bounded, + )) + } + + /// Defines an extension plan carrying a `String` payload and its children, + /// serialized through the extension helpers on the plan contexts. + macro_rules! extension_plan { + ($name:ident, $plan_name:literal) => { + #[derive(Debug)] + pub(crate) struct $name { + pub(crate) payload: String, + pub(crate) children: Vec>, + properties: Arc, + } + + impl $name { + pub(crate) fn new( + payload: impl Into, + children: Vec>, + ) -> Self { + Self { + payload: payload.into(), + children, + properties: stub_properties(), + } + } + } + + impl DisplayAs for $name { + fn fmt_as( + &self, + _t: DisplayFormatType, + f: &mut Formatter, + ) -> std::fmt::Result { + write!(f, "{}({})", $plan_name, self.payload) + } + } + + impl ExecutionPlan for $name { + fn name(&self) -> &str { + $plan_name + } + + fn properties(&self) -> &Arc { + &self.properties + } + + fn children(&self) -> Vec<&Arc> { + self.children.iter().collect() + } + + fn apply_expressions( + &self, + _f: &mut dyn FnMut( + &Arc, + ) -> Result, + ) -> Result { + Ok(TreeNodeRecursion::Continue) + } + + fn with_new_children( + self: Arc, + children: Vec>, + ) -> Result> { + Ok(Arc::new(Self::new(self.payload.clone(), children))) + } + + fn execute( + &self, + _partition: usize, + _context: Arc, + ) -> Result { + internal_err!("{} is a serde-only test stub", $plan_name) + } + + fn try_to_proto( + &self, + ctx: &ExecutionPlanEncodeCtx<'_>, + ) -> Result> { + Ok(Some(ctx.extension_node::( + self.payload.as_bytes().to_vec(), + self.children(), + )?)) + } + } + + impl ExtensionPlanFromProto for $name { + const NAME: &'static str = $plan_name; + + fn try_from_proto( + node: &PhysicalPlanNode, + ctx: &ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + let extension = crate::expect_plan_variant!( + node, + datafusion_proto_models::protobuf::physical_plan_node::PhysicalPlanType::Extension, + "Extension" + ); + let children = ctx.decode_children(&extension.inputs)?; + let Ok(payload) = std::str::from_utf8(&extension.node) else { + return internal_err!("{} payload is not UTF-8", $plan_name); + }; + Ok(Arc::new(Self::new(payload, children))) + } + } + + }; + } + + extension_plan!(RegisteredExec, "datafusion-test.RegisteredExec"); + // A *different* Rust type claiming the same name: the collision two + // independent crates could hit. + extension_plan!(RegisteredExecClone, "datafusion-test.RegisteredExec"); + extension_plan!(OtherRegisteredExec, "datafusion-test.OtherRegisteredExec"); + // A plan whose author forgot to fill in `NAME`: registration must + // reject it rather than let an unnamed node reach the wire. + extension_plan!(UnnamedExec, ""); +} + /// Decoder that must never run: asserts that the reject paths of a /// `try_from_proto` (wrong node variant, missing required child) bail out /// before any decoding happens. diff --git a/datafusion/proto-models/src/lib.rs b/datafusion/proto-models/src/lib.rs index 4a8e7b7ee2bb2..c953a37fc6f9f 100644 --- a/datafusion/proto-models/src/lib.rs +++ b/datafusion/proto-models/src/lib.rs @@ -34,7 +34,7 @@ //! itself — see [`from_proto`] and [`to_proto`]. It is the schema source of //! truth for [`datafusion-proto`]. //! -//! It also hosts [`ProtoDecoderRegistry`](registry::ProtoDecoderRegistry), the +//! It also hosts [`ProtoDecoderRegistry`], the //! one store of extension decoders every serializable kind shares, for the same //! layering reason: it sits below every crate that owns one of those traits. //! diff --git a/datafusion/proto/src/physical_plan/mod.rs b/datafusion/proto/src/physical_plan/mod.rs index 464d5a2aa7666..e9b7fd647b302 100644 --- a/datafusion/proto/src/physical_plan/mod.rs +++ b/datafusion/proto/src/physical_plan/mod.rs @@ -71,7 +71,8 @@ use datafusion_physical_plan::placeholder_row::PlaceholderRowExec; use datafusion_physical_plan::projection::ProjectionExec; use datafusion_physical_plan::proto::{ ExecutionPlanDecode, ExecutionPlanDecodeCtx, ExecutionPlanEncode, - ExecutionPlanEncodeCtx, + ExecutionPlanEncodeCtx, ProtoDecoderRegistry, decode_execution_plan, + execution_plan_names, }; use datafusion_physical_plan::repartition::RepartitionExec; use datafusion_physical_plan::scalar_subquery::ScalarSubqueryExec; @@ -1403,20 +1404,57 @@ pub trait PhysicalPlanNodeExt: Sized { ctx: &PhysicalPlanDecodeContext<'_>, proto_converter: &dyn PhysicalProtoConverterExtension, ) -> Result> { + // Lookup-order policy, mirroring the one function decode already uses a + // layer down (payload -> codec; else registry -> codec fallback): a + // plan that named itself on the wire and is registered on this session + // decodes itself, no codec involved. Anything else — an unnamed node + // from a codec-encoded writer, or a name this session does not know — + // takes the codec chain exactly as before. + let registry = ctx + .task_ctx() + .session_config() + .get_extension::(); + if let Some(registry) = registry.as_ref() { + let plan_decoder = ConverterPlanDecoder { + ctx, + proto_converter, + }; + // The decoder receives the whole node, like every built-in + // `try_from_proto`, reads its own name off it and decodes its own + // children through the ctx. `None` here means no decoder claims + // the node; a decode *failure* is returned as-is rather than + // falling through to the codec. + if let Some(decoded) = decode_execution_plan( + registry, + self.node(), + &ExecutionPlanDecodeCtx::new(&plan_decoder), + ) { + return decoded; + } + } + let inputs: Vec> = extension .inputs .iter() .map(|i| proto_converter.proto_to_execution_plan(i, ctx)) .collect::>()?; - let extension_node = ctx.codec().try_decode( - extension.node.as_slice(), - &inputs, - ctx.task_ctx(), - proto_converter, - )?; - - Ok(extension_node) + ctx.codec() + .try_decode( + extension.node.as_slice(), + &inputs, + ctx.task_ctx(), + proto_converter, + ) + .map_err(|e| match extension.plan_name.as_deref() { + // The writer named the plan but this session has no decoder for + // it. Add that as context, rather than leaving only the codec's + // "unsupported plan" error to explain a missing registration. + Some(plan_name) => { + unregistered_extension_plan_context(plan_name, registry.as_deref(), e) + } + None => e, + }) } fn generate_series_name_to_str(name: protobuf::GenerateSeriesName) -> &'static str { @@ -2146,6 +2184,36 @@ impl PhysicalExtensionCodec for ComposedPhysicalExtensionCodec { } } +/// Add "nothing is registered under this name" context to the error a +/// `PhysicalExtensionCodec` returned for a node that *does* name its plan type. +/// +/// Lists what *is* registered, because the usual cause is a session configured +/// on the writing side but not on the reading one. The codec's own error is +/// kept as the cause rather than reworded: the codec may have failed for a +/// reason that has nothing to do with the missing registration, and a caller +/// matching on the error kind must still see the kind the codec chose. +fn unregistered_extension_plan_context( + plan_name: &str, + registry: Option<&ProtoDecoderRegistry>, + codec_error: DataFusionError, +) -> DataFusionError { + let mut registered: Vec<&str> = registry + .map(|registry| execution_plan_names(registry).collect()) + .unwrap_or_default(); + registered.sort_unstable(); + let registered = if registered.is_empty() { + "none".to_string() + } else { + registered.join(", ") + }; + codec_error.context(format!( + "No decoder is registered for the extension ExecutionPlan '{plan_name}'. Register \ + the plan in the ProtoDecoderRegistry attached to the decoding session's \ + SessionConfig, or supply a PhysicalExtensionCodec that handles it. Registered \ + extension plans: {registered}" + )) +} + /// Adapter backing [`ExecutionPlanEncodeCtx`] for plans migrated to the /// `try_to_proto` hook (#22419). Routes child-plan and child-expr encoding back /// through the central converter so nested plans honor their own hooks. diff --git a/datafusion/proto/tests/cases/plans/extensions.rs b/datafusion/proto/tests/cases/plans/extensions.rs new file mode 100644 index 0000000000000..ded9cf4d0d6dd --- /dev/null +++ b/datafusion/proto/tests/cases/plans/extensions.rs @@ -0,0 +1,775 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Extension `ExecutionPlan`s that decode through the session's name-keyed +//! registry rather than through a `PhysicalExtensionCodec`. +//! +//! The registry closes the loop `try_to_proto` opened: encode already lived on +//! the plan, decode did not. These tests pin the parts the colocated unit tests +//! in `datafusion-physical-plan` structurally cannot reach — that the central +//! dispatch really routes to the registered decoder, that a codec is genuinely +//! not required, and that codec-encoded plans keep decoding exactly as before. + +use std::fmt::Formatter; +use std::sync::Arc; + +use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef}; +use datafusion::execution::TaskContext; +use datafusion::physical_plan::empty::EmptyExec; +use datafusion::physical_plan::proto::{ + ExecutionPlanDecodeCtx, ExecutionPlanEncodeCtx, ExtensionPlanFromProto, + ProtoDecoderRegistry, register_execution_plan, +}; +use datafusion::physical_plan::{ + DisplayAs, DisplayFormatType, ExecutionPlan, PhysicalExpr, PlanProperties, + SendableRecordBatchStream, +}; +use datafusion::prelude::{SessionConfig, SessionContext}; +use datafusion_common::tree_node::TreeNodeRecursion; +use datafusion_common::{ + DataFusionError, Result, internal_datafusion_err, internal_err, resources_err, +}; +use datafusion_proto::physical_plan::{ + AsExecutionPlan, DefaultPhysicalExtensionCodec, PhysicalExtensionCodec, + PhysicalProtoConverterExtension, +}; +use datafusion_proto::protobuf::physical_plan_node::PhysicalPlanType; +use datafusion_proto::protobuf::{PhysicalExtensionNode, PhysicalPlanNode}; +use prost::Message; + +fn schema() -> SchemaRef { + Arc::new(Schema::new(vec![Field::new("a", DataType::Int64, false)])) +} + +fn child() -> Arc { + Arc::new(EmptyExec::new(schema())) +} + +/// Session state an extension plan needs at decode time but never puts on the +/// wire — the connection pool a distributed engine rebuilds per worker. +#[derive(Debug, PartialEq, Eq)] +struct WorkerPool { + endpoints: Vec, +} + +/// Boilerplate shared by the test plans: a single child, properties borrowed +/// from it, and no expressions of their own. +macro_rules! single_child_exec { + // A plan that serializes itself through the hook. + ($name:ident, $display:literal, self_serializing) => { + single_child_exec!(@shape $name, $display, + fn try_to_proto( + &self, + ctx: &ExecutionPlanEncodeCtx<'_>, + ) -> Result> { + Ok(Some(self.encode_self(ctx)?)) + } + ); + }; + // A plan that leaves serialization to a `PhysicalExtensionCodec`. + ($name:ident, $display:literal) => { + single_child_exec!(@shape $name, $display,); + }; + (@shape $name:ident, $display:literal, $($hook:item)*) => { + impl DisplayAs for $name { + fn fmt_as(&self, _t: DisplayFormatType, f: &mut Formatter) -> std::fmt::Result { + write!(f, $display) + } + } + + impl ExecutionPlan for $name { + fn name(&self) -> &str { + $display + } + + fn properties(&self) -> &Arc { + self.child.properties() + } + + fn children(&self) -> Vec<&Arc> { + vec![&self.child] + } + + fn apply_expressions( + &self, + _f: &mut dyn FnMut(&Arc) -> Result, + ) -> Result { + Ok(TreeNodeRecursion::Continue) + } + + fn with_new_children( + self: Arc, + mut children: Vec>, + ) -> Result> { + Ok(Arc::new(self.with_child(children.remove(0)))) + } + + fn execute( + &self, + _partition: usize, + _context: Arc, + ) -> Result { + internal_err!("{} is a serde-only test plan", $display) + } + + $($hook)* + } + }; +} + +/// Exactly one child, or a decode error naming the plan. +fn only_child( + plan_name: &str, + mut children: Vec>, +) -> Result> { + if children.len() != 1 { + return internal_err!( + "{plan_name} expects exactly one input, got {}", + children.len() + ); + } + Ok(children.remove(0)) +} + +// --------------------------------------------------------------------------- +// A session-dependent extension plan, modeled on `NetworkShuffleExec`: the +// stage number is on the wire, the worker pool is rebuilt from the decoding +// session. +// --------------------------------------------------------------------------- + +#[derive(Debug)] +struct ShuffleExec { + stage: u32, + pool: Arc, + child: Arc, +} + +impl ShuffleExec { + fn with_child(&self, child: Arc) -> Self { + Self { + stage: self.stage, + pool: Arc::clone(&self.pool), + child, + } + } +} + +#[derive(Clone, PartialEq, Message)] +struct ShuffleExecProto { + #[prost(uint32, tag = "1")] + stage: u32, +} + +single_child_exec!(ShuffleExec, "ShuffleExec", self_serializing); + +impl ShuffleExec { + fn encode_self(&self, ctx: &ExecutionPlanEncodeCtx<'_>) -> Result { + let mut payload = vec![]; + ShuffleExecProto { stage: self.stage } + .encode(&mut payload) + .map_err(|e| internal_datafusion_err!("failed to encode ShuffleExec: {e}"))?; + ctx.extension_node::(payload, self.children()) + } +} + +impl ExtensionPlanFromProto for ShuffleExec { + const NAME: &'static str = "datafusion-test.ShuffleExec"; + + fn try_from_proto( + node: &PhysicalPlanNode, + ctx: &ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + let Some(PhysicalPlanType::Extension(extension)) = &node.physical_plan_type + else { + return internal_err!("expected an extension node"); + }; + let children = ctx.decode_children(&extension.inputs)?; + let payload = extension.node.as_slice(); + let proto = ShuffleExecProto::decode(payload) + .map_err(|e| internal_datafusion_err!("failed to decode ShuffleExec: {e}"))?; + // Session-dependent decode: the pool never crossed the wire. + let pool = ctx + .task_ctx() + .session_config() + .get_extension::() + .ok_or_else(|| { + internal_datafusion_err!("no WorkerPool configured on this session") + })?; + Ok(Arc::new(ShuffleExec { + stage: proto.stage, + pool, + child: only_child(Self::NAME, children)?, + })) + } +} + +// --------------------------------------------------------------------------- +// A second, entirely independent extension plan. Registering both needs no +// composition step: the two names never interact. +// --------------------------------------------------------------------------- + +#[derive(Debug)] +struct SamplerExec { + fraction: f64, + child: Arc, +} + +impl SamplerExec { + fn with_child(&self, child: Arc) -> Self { + Self { + fraction: self.fraction, + child, + } + } +} + +#[derive(Clone, PartialEq, Message)] +struct SamplerExecProto { + #[prost(double, tag = "1")] + fraction: f64, +} + +single_child_exec!(SamplerExec, "SamplerExec", self_serializing); + +impl SamplerExec { + fn encode_self(&self, ctx: &ExecutionPlanEncodeCtx<'_>) -> Result { + let mut payload = vec![]; + SamplerExecProto { + fraction: self.fraction, + } + .encode(&mut payload) + .map_err(|e| internal_datafusion_err!("failed to encode SamplerExec: {e}"))?; + ctx.extension_node::(payload, self.children()) + } +} + +impl ExtensionPlanFromProto for SamplerExec { + const NAME: &'static str = "other-crate.SamplerExec"; + + fn try_from_proto( + node: &PhysicalPlanNode, + ctx: &ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + let Some(PhysicalPlanType::Extension(extension)) = &node.physical_plan_type + else { + return internal_err!("expected an extension node"); + }; + let children = ctx.decode_children(&extension.inputs)?; + let payload = extension.node.as_slice(); + let proto = SamplerExecProto::decode(payload) + .map_err(|e| internal_datafusion_err!("failed to decode SamplerExec: {e}"))?; + Ok(Arc::new(SamplerExec { + fraction: proto.fraction, + child: only_child(Self::NAME, children)?, + })) + } +} + +// --------------------------------------------------------------------------- +// An unmigrated extension plan: no hook, no name on the wire, decoded by a +// `PhysicalExtensionCodec` exactly as before this feature existed. +// --------------------------------------------------------------------------- + +#[derive(Debug)] +struct LegacyExec { + child: Arc, +} + +impl LegacyExec { + fn with_child(&self, child: Arc) -> Self { + Self { child } + } +} + +single_child_exec!(LegacyExec, "LegacyExec"); + +#[derive(Debug)] +struct LegacyCodec {} + +impl PhysicalExtensionCodec for LegacyCodec { + fn try_decode( + &self, + buf: &[u8], + inputs: &[Arc], + _ctx: &TaskContext, + _proto_converter: &dyn PhysicalProtoConverterExtension, + ) -> Result> { + if buf != b"legacy" { + return internal_err!("LegacyCodec does not recognize this payload"); + } + Ok(Arc::new(LegacyExec { + child: only_child("LegacyExec", inputs.to_vec())?, + })) + } + + fn try_encode( + &self, + node: Arc, + buf: &mut Vec, + _proto_converter: &dyn PhysicalProtoConverterExtension, + ) -> Result<()> { + if node.downcast_ref::().is_none() { + return internal_err!("LegacyCodec only encodes LegacyExec"); + } + buf.extend_from_slice(b"legacy"); + Ok(()) + } +} + +/// A codec that claims *every* extension payload, standing in for a codec whose +/// name-blind `try_decode` would happily decode another crate's plan wrongly. +#[derive(Debug)] +struct GreedyCodec {} + +impl PhysicalExtensionCodec for GreedyCodec { + fn try_decode( + &self, + _buf: &[u8], + inputs: &[Arc], + _ctx: &TaskContext, + _proto_converter: &dyn PhysicalProtoConverterExtension, + ) -> Result> { + Ok(Arc::new(LegacyExec { + child: only_child("LegacyExec", inputs.to_vec())?, + })) + } + + fn try_encode( + &self, + _node: Arc, + buf: &mut Vec, + _proto_converter: &dyn PhysicalProtoConverterExtension, + ) -> Result<()> { + buf.extend_from_slice(b"greedy"); + Ok(()) + } +} + +/// A codec that recognizes a payload but fails while decoding it, with an error +/// kind of its own. Stands in for every codec failure whose cause is *not* a +/// missing registration. +#[derive(Debug)] +struct ExhaustedCodec {} + +impl PhysicalExtensionCodec for ExhaustedCodec { + fn try_decode( + &self, + _buf: &[u8], + _inputs: &[Arc], + _ctx: &TaskContext, + _proto_converter: &dyn PhysicalProtoConverterExtension, + ) -> Result> { + resources_err!("the decode buffer pool is empty") + } + + fn try_encode( + &self, + _node: Arc, + _buf: &mut Vec, + _proto_converter: &dyn PhysicalProtoConverterExtension, + ) -> Result<()> { + internal_err!("ExhaustedCodec does not encode") + } +} + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +/// A session with `WorkerPool` attached and `plans` registered. +/// A registry populated by `register`, ready to attach with +/// `SessionConfig::with_extension`. +fn registry( + register: impl FnOnce(&mut ProtoDecoderRegistry) -> Result<()>, +) -> Result> { + let mut registry = ProtoDecoderRegistry::new(); + register(&mut registry)?; + Ok(Arc::new(registry)) +} + +fn session_with(pool_endpoints: &[&str]) -> SessionConfig { + SessionConfig::new().with_extension(Arc::new(WorkerPool { + endpoints: pool_endpoints.iter().map(|e| e.to_string()).collect(), + })) +} + +fn encode( + plan: Arc, + codec: &dyn PhysicalExtensionCodec, +) -> Result { + PhysicalPlanNode::try_from_physical_plan(plan, codec) +} + +fn decode( + node: &PhysicalPlanNode, + ctx: &SessionContext, + codec: &dyn PhysicalExtensionCodec, +) -> Result> { + node.try_into_physical_plan(ctx.task_ctx().as_ref(), codec) +} + +/// The `PhysicalExtensionNode` at the root of `node`. +fn extension_of(node: &PhysicalPlanNode) -> &PhysicalExtensionNode { + match &node.physical_plan_type { + Some(PhysicalPlanType::Extension(extension)) => extension, + other => panic!("expected an extension node, got {other:?}"), + } +} + +// --------------------------------------------------------------------------- +// Tests +// --------------------------------------------------------------------------- + +#[test] +fn extension_plan_round_trips_with_no_codec_at_all() -> Result<()> { + let config = + session_with(&["worker-1", "worker-2"]).with_extension(registry(|registry| { + register_execution_plan::(registry) + })?); + let ctx = SessionContext::new_with_config(config); + // The default codec decodes nothing: whatever survives came from the + // registry. + let codec = DefaultPhysicalExtensionCodec {}; + + let plan = Arc::new(ShuffleExec { + stage: 7, + pool: Arc::new(WorkerPool { + endpoints: vec!["worker-1".to_string(), "worker-2".to_string()], + }), + child: child(), + }); + + let node = encode(plan.clone(), &codec)?; + assert_eq!( + extension_of(&node).plan_name.as_deref(), + Some(ShuffleExec::NAME), + "the encode helper must stamp the plan name" + ); + + let decoded = decode(&node, &ctx, &codec)?; + let decoded = decoded + .downcast_ref::() + .expect("decoded plan must be a ShuffleExec"); + assert_eq!(decoded.stage, 7); + assert_eq!(decoded.child.name(), "EmptyExec"); + Ok(()) +} + +#[test] +fn extension_plan_decode_reads_session_state() -> Result<()> { + let codec = DefaultPhysicalExtensionCodec {}; + let writer = SessionContext::new_with_config( + session_with(&["writer-only"]) + .with_extension(registry(register_execution_plan::)?), + ); + let node = encode( + Arc::new(ShuffleExec { + stage: 1, + pool: Arc::clone( + &writer + .copied_config() + .get_extension::() + .expect("pool"), + ), + child: child(), + }), + &codec, + )?; + + // Decoding on a *different* session rebuilds the pool from that session, + // proving the pool never rode the wire. + let reader = SessionContext::new_with_config( + session_with(&["reader-a", "reader-b"]).with_extension(registry(|registry| { + register_execution_plan::(registry) + })?), + ); + let decoded = decode(&node, &reader, &codec)?; + let decoded = decoded + .downcast_ref::() + .expect("decoded plan must be a ShuffleExec"); + assert_eq!(decoded.pool.endpoints, vec!["reader-a", "reader-b"]); + Ok(()) +} + +#[test] +fn independent_extension_plans_coexist_without_a_composed_codec() -> Result<()> { + let config = session_with(&["worker-1"]).with_extension(registry(|registry| { + register_execution_plan::(registry)?; + register_execution_plan::(registry) + })?); + let ctx = SessionContext::new_with_config(config); + let codec = DefaultPhysicalExtensionCodec {}; + + // Two extension plans from two notional crates, nested. + let plan: Arc = Arc::new(ShuffleExec { + stage: 3, + pool: Arc::new(WorkerPool { + endpoints: vec!["worker-1".to_string()], + }), + child: Arc::new(SamplerExec { + fraction: 0.25, + child: child(), + }), + }); + + let decoded = decode(&encode(Arc::clone(&plan), &codec)?, &ctx, &codec)?; + let shuffle = decoded + .downcast_ref::() + .expect("outer plan must be a ShuffleExec"); + assert_eq!(shuffle.stage, 3); + let sampler = shuffle + .child + .downcast_ref::() + .expect("inner plan must be a SamplerExec"); + assert_eq!(sampler.fraction, 0.25); + Ok(()) +} + +#[test] +fn codec_encoded_extension_plans_are_unaffected() -> Result<()> { + // A plan with no hook: no name on the wire, decoded by its codec, exactly + // as before the registry existed. The session registers an unrelated plan + // to prove the registry does not interfere. + let ctx = SessionContext::new_with_config( + SessionConfig::new() + .with_extension(registry(register_execution_plan::)?), + ); + let codec = LegacyCodec {}; + + let node = encode(Arc::new(LegacyExec { child: child() }), &codec)?; + assert_eq!( + extension_of(&node).plan_name, + None, + "codec-encoded plans must stay anonymous on the wire" + ); + + let decoded = decode(&node, &ctx, &codec)?; + assert!(decoded.downcast_ref::().is_some()); + Ok(()) +} + +#[test] +fn a_registered_name_wins_over_a_codec_that_claims_everything() -> Result<()> { + let config = session_with(&["worker-1"]).with_extension(registry(|registry| { + register_execution_plan::(registry) + })?); + let ctx = SessionContext::new_with_config(config); + // This codec would decode *any* payload into a `LegacyExec`; the name must + // route around it. + let codec = GreedyCodec {}; + + let plan = Arc::new(ShuffleExec { + stage: 9, + pool: Arc::new(WorkerPool { + endpoints: vec!["worker-1".to_string()], + }), + child: child(), + }); + + let decoded = decode(&encode(plan, &codec)?, &ctx, &codec)?; + assert!( + decoded.downcast_ref::().is_some(), + "the registry must take precedence over the codec, got {decoded:?}" + ); + Ok(()) +} + +#[test] +fn a_named_plan_the_session_does_not_know_reports_the_missing_registration() -> Result<()> +{ + let codec = DefaultPhysicalExtensionCodec {}; + let writer = SessionContext::new_with_config( + session_with(&["worker-1"]) + .with_extension(registry(register_execution_plan::)?), + ); + let node = encode( + Arc::new(ShuffleExec { + stage: 1, + pool: Arc::clone( + &writer + .copied_config() + .get_extension::() + .expect("pool"), + ), + child: child(), + }), + &codec, + )?; + + // A reader that registered a *different* plan: the fallback runs, fails, + // and the error has to point at the missing registration rather than at the + // codec's generic complaint. + let reader = SessionContext::new_with_config( + SessionConfig::new() + .with_extension(registry(register_execution_plan::)?), + ); + let err = decode(&node, &reader, &codec) + .expect_err("an unregistered plan name must fail to decode"); + let err = err.to_string(); + assert!(err.contains(ShuffleExec::NAME), "unexpected error: {err}"); + assert!( + err.contains("ProtoDecoderRegistry"), + "unexpected error: {err}" + ); + assert!( + !err.contains("bug in DataFusion"), + "a missing registration is a session-configuration mistake, not a bug: {err}" + ); + assert!(err.contains(SamplerExec::NAME), "unexpected error: {err}"); + Ok(()) +} + +#[test] +fn a_registered_decoder_failing_does_not_fall_back_to_the_codec() -> Result<()> { + // The point of the name is that it *decides* the decoder. Falling back to + // the codec after a registered decoder failed would resurrect exactly the + // wrong decode this replaces: `GreedyCodec` would happily turn a corrupt + // ShuffleExec into a LegacyExec. + let config = session_with(&["worker-1"]).with_extension(registry(|registry| { + register_execution_plan::(registry) + })?); + let ctx = SessionContext::new_with_config(config); + + let node = PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::Extension(PhysicalExtensionNode { + // Not a valid `ShuffleExecProto`: field 1 is declared `uint32` but + // this is a length-delimited payload, so prost rejects it. + node: vec![0x0a, 0x01, 0x00], + inputs: vec![encode(child(), &DefaultPhysicalExtensionCodec {})?], + plan_name: Some(ShuffleExec::NAME.to_string()), + })), + }; + + let err = decode(&node, &ctx, &GreedyCodec {}) + .expect_err("a registered decoder's failure must be fatal"); + assert!( + err.to_string().contains("failed to decode ShuffleExec"), + "the registered decoder's own error must survive, got: {err}" + ); + Ok(()) +} + +#[test] +fn the_plan_name_survives_the_wire_and_is_absent_for_old_writers() -> Result<()> { + let config = session_with(&["worker-1"]).with_extension(registry(|registry| { + register_execution_plan::(registry) + })?); + let ctx = SessionContext::new_with_config(config); + let codec = DefaultPhysicalExtensionCodec {}; + + let node = encode( + Arc::new(ShuffleExec { + stage: 11, + pool: Arc::new(WorkerPool { + endpoints: vec!["worker-1".to_string()], + }), + child: child(), + }), + &codec, + )?; + + // Through real bytes, not just the in-memory struct. + let decoded = PhysicalPlanNode::decode(node.encode_to_vec().as_slice()) + .map_err(|e| internal_datafusion_err!("failed to decode the plan node: {e}"))?; + assert_eq!( + extension_of(&decoded).plan_name.as_deref(), + Some(ShuffleExec::NAME) + ); + assert_eq!( + decode(&decoded, &ctx, &codec)? + .downcast_ref::() + .expect("decoded plan must be a ShuffleExec") + .stage, + 11 + ); + + // A writer that predates the field emits the same bytes without field 3; + // dropping it must read back as "unnamed", i.e. the codec path. + let mut old_writer = node.clone(); + match &mut old_writer.physical_plan_type { + Some(PhysicalPlanType::Extension(extension)) => extension.plan_name = None, + other => panic!("expected an extension node, got {other:?}"), + } + let old_writer = PhysicalPlanNode::decode(old_writer.encode_to_vec().as_slice()) + .map_err(|e| internal_datafusion_err!("failed to decode the plan node: {e}"))?; + assert_eq!(extension_of(&old_writer).plan_name, None); + assert!( + decode(&old_writer, &ctx, &codec).is_err(), + "an unnamed node must reach the codec, which cannot decode it" + ); + Ok(()) +} + +#[test] +fn a_named_plan_can_still_be_decoded_by_a_codec() -> Result<()> { + // A name this session did not register is not fatal on its own, even + // when the session has a registry with other plans in it: a codec that + // understands the payload still decodes it, which is what keeps a mixed + // fleet (some nodes upgraded, some not) working. + let config = SessionConfig::new().with_extension(registry(|registry| { + register_execution_plan::(registry) + })?); + let ctx = SessionContext::new_with_config(config); + let node = PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::Extension(PhysicalExtensionNode { + node: b"legacy".to_vec(), + inputs: vec![encode(child(), &DefaultPhysicalExtensionCodec {})?], + plan_name: Some("some-crate.UnknownExec".to_string()), + })), + }; + + let decoded = decode(&node, &ctx, &LegacyCodec {})?; + assert!(decoded.downcast_ref::().is_some()); + Ok(()) +} + +#[test] +fn the_codecs_own_error_survives_the_missing_registration_hint() -> Result<()> { + // The hint is added to whatever the codec returned, never substituted for + // it: a codec can fail for a reason that has nothing to do with a missing + // registration, and a caller matching on the error kind must still see the + // kind the codec chose. + let ctx = SessionContext::new_with_config( + SessionConfig::new() + .with_extension(registry(register_execution_plan::)?), + ); + let node = PhysicalPlanNode { + physical_plan_type: Some(PhysicalPlanType::Extension(PhysicalExtensionNode { + node: b"whatever".to_vec(), + inputs: vec![encode(child(), &DefaultPhysicalExtensionCodec {})?], + plan_name: Some("some-crate.UnknownExec".to_string()), + })), + }; + + let err = decode(&node, &ctx, &ExhaustedCodec {}) + .expect_err("the codec must fail for this payload"); + assert!( + matches!(err.find_root(), DataFusionError::ResourcesExhausted(_)), + "the codec's error kind must survive, got: {err:?}" + ); + let text = err.to_string(); + assert!( + text.contains("the decode buffer pool is empty"), + "the codec's own message must survive: {text}" + ); + assert!( + text.contains("No decoder is registered") + && text.contains("some-crate.UnknownExec"), + "the missing-registration hint must be added: {text}" + ); + Ok(()) +} diff --git a/datafusion/proto/tests/cases/plans/mod.rs b/datafusion/proto/tests/cases/plans/mod.rs index 40839d4b3a31a..1faa634380289 100644 --- a/datafusion/proto/tests/cases/plans/mod.rs +++ b/datafusion/proto/tests/cases/plans/mod.rs @@ -37,6 +37,7 @@ mod aggregates; mod dispatch; mod dynamic_filters; mod exprs; +mod extensions; mod filters; mod joins; mod leaves; diff --git a/docs/source/library-user-guide/upgrading/56.0.0.md b/docs/source/library-user-guide/upgrading/56.0.0.md index ef27ef0a1ee37..1c10c7f0dcefa 100644 --- a/docs/source/library-user-guide/upgrading/56.0.0.md +++ b/docs/source/library-user-guide/upgrading/56.0.0.md @@ -839,3 +839,89 @@ let plan = provider .await? .into_inner(); ``` + +### Extension `ExecutionPlan`s can decode without a `PhysicalExtensionCodec` + +`ExecutionPlan::try_to_proto` let a plan serialize itself, but there was no way back: decoding an extension node always routed through `PhysicalExtensionCodec::try_decode`. `PhysicalExtensionNode` carried no type discriminator, so the codec _was_ the discriminator. `ComposedPhysicalExtensionCodec` works around that by writing the position of the encoding codec into the payload, which means the decoding side must register the same codecs in the same order, and a name collision between two independent crates is undetectable. + +`PhysicalExtensionNode` now has an optional `plan_name`, and a session can map that name to a decoder. An extension plan implements `ExtensionPlanFromProto`, which pairs the name it is dispatched by with the constructor that rebuilds it, and is registered on the session that will decode it: + +```rust,ignore +use datafusion::physical_plan::proto::{ + ExecutionPlanDecodeCtx, ExecutionPlanEncodeCtx, ExtensionPlanFromProto, + ProtoDecoderRegistry, register_execution_plan, +}; + +impl ExecutionPlan for MyExec { + // ... + fn try_to_proto( + &self, + ctx: &ExecutionPlanEncodeCtx<'_>, + ) -> Result> { + // `extension_node` stamps `Self::NAME` on the node, so the name on the + // wire and the registry key cannot drift apart. + Ok(Some(ctx.extension_node::( + my_payload_bytes(self)?, + self.children(), + )?)) + } + +} + +impl ExtensionPlanFromProto for MyExec { + // Namespace the name so a collision with another crate is an error at + // registration rather than a wrong decode. + const NAME: &'static str = "my-crate.MyExec"; + + fn try_from_proto( + node: &PhysicalPlanNode, + ctx: &ExecutionPlanDecodeCtx<'_>, + ) -> Result> { + let extension = expect_plan_variant!(node, PhysicalPlanType::Extension, "Extension"); + let children = ctx.decode_children(&extension.inputs)?; + // `ctx.task_ctx()` is available here, so a plan that rebuilds session + // state at decode time (a connection pool, say) can do so. + my_plan_from_bytes(&extension.node, children) + } +} + +// On every session that decodes the plan — for a distributed engine, the +// workers as well as the coordinator. +let mut registry = ProtoDecoderRegistry::new(); +register_execution_plan::(&mut registry)?; +let config = SessionConfig::new().with_extension(Arc::new(registry)); +``` + +Only extension plans implement `ExtensionPlanFromProto`, and `register_execution_plan` takes that trait as its bound. A built-in plan has a wire variant of its own and keeps its inherent `try_from_proto`, so `register_execution_plan::(&mut registry)` does not compile. + +`ProtoDecoderRegistry` is shared by every kind of extension DataFusion can decode, not only plans. Its key is the trait plus the name, so a plan and an expression may use the same name, and one session carries one registry rather than one per kind. + +The application that owns the session builds the registry. A session carries at most one, and `with_extension` replaces what is there, so a library must not attach a registry of its own: it would silently discard another library's. A library exposes a function that fills a registry it is handed, and the application composes them: + +```rust,ignore +// In each library that ships extension types: +pub fn register(registry: &mut ProtoDecoderRegistry) -> Result<()> { + register_execution_plan::(registry) +} + +// In the application that owns the session: +let mut registry = ProtoDecoderRegistry::new(); +shuffle_lib::register(&mut registry)?; +sampling_lib::register(&mut registry)?; +let config = SessionConfig::new().with_extension(Arc::new(registry)); +``` + +The libraries need to know nothing about each other, and the order they are called in does not matter. Dispatch is by name, and a real name collision is an error from `register_execution_plan` rather than a wrong decode later. + +A registered plan needs no `PhysicalExtensionCodec` at all, and two crates' plans coexist on one session without a `ComposedPhysicalExtensionCodec`. + +**Who is affected:** + +- Nobody is required to change anything. This is opt-in per plan type: a `PhysicalExtensionNode` with no name, or with a name the decoding session does not know, takes the `PhysicalExtensionCodec` chain exactly as before. On the binary wire format the new field is additive: old writers omit it and old readers skip it. +- The generated JSON codec is stricter than the binary one. `pbjson` rejects a field it does not know, so a named node written by a new writer cannot be read as JSON by a reader on an older DataFusion. A plan that does not opt in writes no name and is unaffected. If you serialize plans as JSON across versions, upgrade the readers before you opt a plan in. +- Users constructing `PhysicalExtensionNode` with an exhaustive struct literal must set the new `plan_name` field (`None` reproduces the previous behavior). +- Replacing a `PhysicalExtensionCodec` with a registration changes only _how_ the payload is dispatched, not what the payload is. Readers running an older DataFusion ignore `plan_name` and hand the payload to their codec, so a mixed fleet keeps working only while the bytes your `try_to_proto` writes stay decodable by the codec you removed. + +`PhysicalExtensionCodec` is not deprecated: extension `PhysicalExpr`s and UDF payloads still need it, and it remains the fallback for unmigrated plans. + +See [issue #24625](https://github.com/apache/datafusion/issues/24625) for details. From ca602d12d545d2ff70cdddebf525dd1f0473f13d Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Wed, 30 Sep 2026 14:32:06 -0500 Subject: [PATCH 4/4] Hide decode_execution_plan from docs It is only public because datafusion-proto calls it across a crate boundary. Co-Authored-By: Claude Opus 5.5 --- datafusion/physical-plan/src/proto/mod.rs | 6 ++++-- datafusion/physical-plan/src/proto/registry.rs | 8 +++++++- 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/datafusion/physical-plan/src/proto/mod.rs b/datafusion/physical-plan/src/proto/mod.rs index b3fae8df50b59..edbad4c71d66a 100644 --- a/datafusion/physical-plan/src/proto/mod.rs +++ b/datafusion/physical-plan/src/proto/mod.rs @@ -93,9 +93,11 @@ use datafusion_proto_models::protobuf::{ use crate::ExecutionPlan; pub use datafusion_proto_models::ProtoDecoderRegistry; +// Not public API: see `decode_execution_plan`. +#[doc(hidden)] +pub use registry::decode_execution_plan; pub use registry::{ - ExtensionPlanFromProto, decode_execution_plan, execution_plan_names, - register_execution_plan, + ExtensionPlanFromProto, execution_plan_names, register_execution_plan, }; /// Internal dispatch trait backing [`ExecutionPlanEncodeCtx`]. diff --git a/datafusion/physical-plan/src/proto/registry.rs b/datafusion/physical-plan/src/proto/registry.rs index b604b9a91c5f0..4f12f89e6f5a3 100644 --- a/datafusion/physical-plan/src/proto/registry.rs +++ b/datafusion/physical-plan/src/proto/registry.rs @@ -107,7 +107,8 @@ pub trait ExtensionPlanFromProto: ExecutionPlan + Sized { /// Deliberately private, and the same type on both the /// [`register_execution_plan`] and the [`decode_execution_plan`] side. /// [`ExtensionPlanFromProto`] is the public contract and -/// [`decode_execution_plan`] is the public way to invoke one, so this can +/// [`decode_execution_plan`] (hidden, internal to DataFusion) is the only way +/// to invoke one, so this can /// become something else — a `dyn` decoder object, to admit stateful or /// closure decoders, which is what an FFI decoder needs — without a breaking /// change. @@ -167,6 +168,11 @@ pub fn register_execution_plan( /// `Some(Err(..))` means the decoder that *does* own the name failed, which is /// fatal: falling back there would let another codec decode the payload /// wrongly, the very thing the name exists to prevent. +/// +/// Not intended as public API: `datafusion-proto` drives decoding and is the +/// only intended caller. It is `pub` only because that call crosses a crate +/// boundary, so it is hidden from the docs and may change without notice. +#[doc(hidden)] pub fn decode_execution_plan( registry: &ProtoDecoderRegistry, node: &PhysicalPlanNode,