gabotechs commented on code in PR #24631: URL: https://github.com/apache/datafusion/pull/24631#discussion_r4141787781
########## 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<dyn Any + Send + Sync>, + 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::<MyExec>(registry)?; +/// register_physical_expr::<MyExpr>(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>, +} + Review Comment: :+1: nice ########## 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<Option<PhysicalPlanNode>> { +/// // Stamps `Self::NAME`, so the encoded name and the registry key +/// // cannot drift apart. +/// Ok(Some(ctx.extension_node::<Self>(self.payload()?, self.children())?)) +/// } +/// } +/// +/// impl ExtensionPlanFromProto for MyExec { +/// const NAME: &'static str = "my-crate.MyExec"; +/// +/// fn try_from_proto( +/// node: &PhysicalPlanNode, +/// ctx: &ExecutionPlanDecodeCtx<'_>, +/// ) -> Result<Arc<dyn ExecutionPlan>> { +/// 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<Arc<dyn ExecutionPlan>>; +} + +/// 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<Arc<dyn ExecutionPlan>>; + +/// 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::<FilterExec>(&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<T: ExtensionPlanFromProto>( + registry: &mut ProtoDecoderRegistry, +) -> Result<()> { + registry.register_decoder::<dyn ExecutionPlan, T, ExecutionPlanDecoder>( + 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<Result<Arc<dyn ExecutionPlan>>> { + let name = plan_name(node)?; Review Comment: This is a bit weird from an API standpoint. As a consumer of the API the first thing that comes to mind is we would the registry not be part of the context. ########## 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<Option<PhysicalPlanNode>> { +/// // Stamps `Self::NAME`, so the encoded name and the registry key +/// // cannot drift apart. +/// Ok(Some(ctx.extension_node::<Self>(self.payload()?, self.children())?)) +/// } +/// } +/// +/// impl ExtensionPlanFromProto for MyExec { +/// const NAME: &'static str = "my-crate.MyExec"; +/// +/// fn try_from_proto( +/// node: &PhysicalPlanNode, +/// ctx: &ExecutionPlanDecodeCtx<'_>, +/// ) -> Result<Arc<dyn ExecutionPlan>> { +/// 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<Arc<dyn ExecutionPlan>>; +} + +/// 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<Arc<dyn ExecutionPlan>>; + +/// 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::<FilterExec>(&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<T: ExtensionPlanFromProto>( + registry: &mut ProtoDecoderRegistry, +) -> Result<()> { + registry.register_decoder::<dyn ExecutionPlan, T, ExecutionPlanDecoder>( + 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<Result<Arc<dyn ExecutionPlan>>> { + let name = plan_name(node)?; Review Comment: :thinking: this does not seem that users would typically need to interact with anyway, seems like it's just `pub` because it needs to be used in a separate crate in this workspace. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
