Skip to content

Commit 86b779f

Browse files
adriangbclaude
andcommitted
feat: decode extension PhysicalExprs through a session-scoped registry
`PhysicalExtensionExprNode` carried no type discriminator, so the `PhysicalExtensionCodec` *was* the discriminator, with the same problems the plan side has: composition is by registration order, and a name collision between two independent crates is undetectable. An extension expression now declares its wire name with `ExtensionExprName`, implements `PhysicalExprFromProto` like a built-in, writes itself with `PhysicalExprEncodeCtx::extension_expr_node` (which stamps the name, so the encoded name and the registry key cannot drift apart), and is registered with `register_physical_expr` into the `ProtoDecoderRegistry` that the decoding session carries — the same registry extension `ExecutionPlan`s use, keyed on the trait as well as the name. The name is a separate trait from the decode contract so that only extension expressions can be registered: a built-in has an `ExprType` variant of its own, and `register_physical_expr::<Column>(..)` does not compile. It lives in `physical-expr-common`, beside the encode context that stamps it, because that crate sits below the one that owns the decode contract; `datafusion_physical_expr::proto` re-exports 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_physical_expr` 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. Documented in the 56.0.0 upgrade guide; tests in `datafusion/proto/tests/cases/plans/expr_registry.rs`. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 9d97df1 commit 86b779f

8 files changed

Lines changed: 1156 additions & 10 deletions

File tree

‎datafusion/physical-expr-common/src/physical_expr.rs‎

Lines changed: 61 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -555,10 +555,36 @@ pub mod proto_encode {
555555
use std::sync::Arc;
556556

557557
use datafusion_common::Result;
558-
use datafusion_proto_models::protobuf::PhysicalExprNode;
558+
use datafusion_proto_models::protobuf::physical_expr_node::ExprType;
559+
use datafusion_proto_models::protobuf::{
560+
PhysicalExprNode, PhysicalExtensionExprNode,
561+
};
559562

560563
use super::PhysicalExpr;
561564

565+
/// The wire name of an extension [`PhysicalExpr`], and the key it is
566+
/// registered and dispatched by.
567+
///
568+
/// Only extension expressions implement this. A built-in expression has an
569+
/// `ExprType` variant of its own and is dispatched by that variant, so it
570+
/// needs no name and cannot be registered.
571+
///
572+
/// It lives here, beside [`PhysicalExprEncodeCtx`], because the encode
573+
/// helper that stamps the name needs it and this crate sits below the one
574+
/// that owns the decode contract. `datafusion_physical_expr::proto`
575+
/// re-exports it, next to `PhysicalExprFromProto`, which is where an
576+
/// expression author should import it from.
577+
pub trait ExtensionExprName: PhysicalExpr {
578+
/// The name this expression type is written and registered under.
579+
///
580+
/// Namespace it with the owning crate — `"my_crate.MyExpr"`, not
581+
/// `"MyExpr"` — so that two independent crates registering into the
582+
/// same session collide at registration time instead of silently
583+
/// decoding each other's nodes. Never write it by hand on the wire:
584+
/// [`PhysicalExprEncodeCtx::extension_expr_node`] stamps it for you.
585+
const NAME: &'static str;
586+
}
587+
562588
/// Encoder context handed to [`super::PhysicalExpr::try_to_proto`].
563589
///
564590
/// Wraps an internal [`PhysicalExprEncode`] trait object so callers see a
@@ -584,6 +610,40 @@ pub mod proto_encode {
584610
self.encoder.encode(expr)
585611
}
586612

613+
/// Build the `PhysicalExprNode` for an extension expression `T`: an
614+
/// `Extension` variant carrying `payload`, the encoded `children`, and
615+
/// `T`'s [`NAME`](ExtensionExprName::NAME).
616+
///
617+
/// The one supported way for an extension expression to write itself.
618+
/// The name on the wire is taken from the same const the registry is
619+
/// keyed by, so an encoder cannot write a name no decoder answers to:
620+
///
621+
/// ```ignore
622+
/// fn try_to_proto(
623+
/// &self,
624+
/// ctx: &PhysicalExprEncodeCtx<'_>,
625+
/// ) -> Result<Option<PhysicalExprNode>> {
626+
/// Ok(Some(ctx.extension_expr_node::<Self>(
627+
/// self.state.to_bytes()?,
628+
/// self.children(),
629+
/// )?))
630+
/// }
631+
/// ```
632+
pub fn extension_expr_node<T: ExtensionExprName>(
633+
&self,
634+
payload: Vec<u8>,
635+
children: Vec<&Arc<dyn PhysicalExpr>>,
636+
) -> Result<PhysicalExprNode> {
637+
Ok(PhysicalExprNode {
638+
expr_type: Some(ExprType::Extension(PhysicalExtensionExprNode {
639+
expr: payload,
640+
inputs: self.encode_children_expressions(children)?,
641+
expr_name: Some(T::NAME.to_string()),
642+
})),
643+
..Default::default()
644+
})
645+
}
646+
587647
/// Encode a sequence of child expressions, preserving order.
588648
///
589649
/// Convenience wrapper over [`Self::encode_child`] for expressions

‎datafusion/physical-expr/src/proto.rs‎

Lines changed: 118 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,13 @@
1616
// under the License.
1717

1818
//! The decode half of [`PhysicalExpr::try_to_proto`]: [`PhysicalExprFromProto`],
19-
//! the contract every self-serializing expression implements.
19+
//! the contract every self-serializing expression implements, plus the
20+
//! `PhysicalExpr` face of the session-scoped [`ProtoDecoderRegistry`] that
21+
//! routes extension expressions to theirs by name —
22+
//! [`register_physical_expr`], [`decode_physical_expr`] and
23+
//! [`physical_expr_names`], all keyed on `dyn PhysicalExpr`. The store itself
24+
//! is shared with every other extension kind, so a session carries one
25+
//! registry.
2026
//!
2127
//! These live here rather than next to the encode context in
2228
//! `datafusion-physical-expr-common` because a decoder may need the session,
@@ -33,7 +39,12 @@ use datafusion_common::config::ConfigOptions;
3339
use datafusion_expr::registry::FunctionRegistry;
3440
use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
3541
pub use datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecodeCtx;
42+
pub use datafusion_physical_expr_common::physical_expr::proto_encode::{
43+
ExtensionExprName, PhysicalExprEncodeCtx,
44+
};
45+
pub use datafusion_proto_models::ProtoDecoderRegistry;
3646
use datafusion_proto_models::protobuf::PhysicalExprNode;
47+
use datafusion_proto_models::protobuf::physical_expr_node::ExprType;
3748

3849
/// What an expression decoder can see of the session it runs under: the
3950
/// function registry, for resolving UDFs by name, and the configuration
@@ -81,6 +92,10 @@ impl std::fmt::Debug for ExprDecodeSession<'_> {
8192
/// Every self-serializing expression implements this. Built-in expressions
8293
/// are dispatched to their `try_from_proto` by their `ExprType` variant.
8394
///
95+
/// An extension expression has no variant of its own, so it is dispatched by
96+
/// name instead and implements [`ExtensionExprName`] as well.
97+
/// [`register_physical_expr`] carries the worked example.
98+
///
8499
/// ```ignore
85100
/// impl PhysicalExprFromProto for MyExpr {
86101
/// fn try_from_proto(
@@ -114,3 +129,105 @@ pub trait PhysicalExprFromProto: PhysicalExpr + Sized {
114129
ctx: &PhysicalExprDecodeCtx<'_, ExprDecodeSession<'_>>,
115130
) -> Result<Arc<dyn PhysicalExpr>>;
116131
}
132+
133+
/// How this facade stores a decoder in the shared registry: a function pointer
134+
/// to the monomorphized [`PhysicalExprFromProto::try_from_proto`].
135+
///
136+
/// Deliberately private, and the same type on both the
137+
/// [`register_physical_expr`] and the [`decode_physical_expr`] side.
138+
/// [`PhysicalExprFromProto`] is the public contract and
139+
/// [`decode_physical_expr`] is the public way to invoke one, so this can become
140+
/// something else — a `dyn` decoder object, to admit stateful or closure
141+
/// decoders, which is what an FFI decoder needs — without a breaking change.
142+
type PhysicalExprDecoder = fn(
143+
&PhysicalExprNode,
144+
&PhysicalExprDecodeCtx<'_, ExprDecodeSession<'_>>,
145+
) -> Result<Arc<dyn PhysicalExpr>>;
146+
147+
/// Register `T` in `registry` under its [`ExtensionExprName::NAME`].
148+
///
149+
/// Registering the same type twice is a no-op. Registering a *different* type
150+
/// under a name already taken is an error, so collisions surface here rather
151+
/// than as a wrong decode later. The name is scoped to `dyn PhysicalExpr`, so a
152+
/// plan or a data source may use the same name in the same registry.
153+
///
154+
/// # Who builds the registry
155+
///
156+
/// The application that owns the session builds it. A session carries at most
157+
/// one registry, and `SessionConfig::with_extension` replaces what is there, so
158+
/// a library must never attach a registry of its own: it would silently discard
159+
/// another library's. A library exposes a function that fills a registry it is
160+
/// handed, and the application composes them. One registry holds every kind, so
161+
/// a library registers its plans and its expressions into the same object.
162+
///
163+
/// # A built-in expression cannot be registered
164+
///
165+
/// The bound is [`ExtensionExprName`], which only extension expressions
166+
/// implement. A built-in implements the decode contract:
167+
///
168+
/// ```
169+
/// use datafusion_physical_expr::expressions::Column;
170+
/// use datafusion_physical_expr::proto::PhysicalExprFromProto;
171+
///
172+
/// let _ = <Column as PhysicalExprFromProto>::try_from_proto;
173+
/// ```
174+
///
175+
/// but it has no wire name, and registering it does not compile:
176+
///
177+
/// ```compile_fail
178+
/// use datafusion_physical_expr::expressions::Column;
179+
/// use datafusion_physical_expr::proto::register_physical_expr;
180+
/// use datafusion_proto_models::ProtoDecoderRegistry;
181+
///
182+
/// let mut registry = ProtoDecoderRegistry::new();
183+
/// register_physical_expr::<Column>(&mut registry).unwrap();
184+
/// ```
185+
pub fn register_physical_expr<T: ExtensionExprName + PhysicalExprFromProto>(
186+
registry: &mut ProtoDecoderRegistry,
187+
) -> Result<()> {
188+
registry.register_decoder::<dyn PhysicalExpr, T, PhysicalExprDecoder>(
189+
T::NAME,
190+
T::try_from_proto,
191+
)
192+
}
193+
194+
/// Decode `node` with the extension expression decoder registered under the
195+
/// name `node` carries.
196+
///
197+
/// The name is read from the node's `Extension` variant rather than passed in,
198+
/// so a caller cannot pair a node with a name it does not carry.
199+
///
200+
/// `None` means "this node names no extension expression decoder of ours": it
201+
/// is not an extension node, it carries no name, or no registered name matches.
202+
/// The caller then falls back to the `PhysicalExtensionCodec` chain.
203+
/// `Some(Err(..))` means the decoder that *does* own the name failed, which is
204+
/// fatal: falling back there would let a codec decode the payload wrongly, the
205+
/// very thing the name exists to prevent.
206+
pub fn decode_physical_expr(
207+
registry: &ProtoDecoderRegistry,
208+
node: &PhysicalExprNode,
209+
ctx: &PhysicalExprDecodeCtx<'_, ExprDecodeSession<'_>>,
210+
) -> Option<Result<Arc<dyn PhysicalExpr>>> {
211+
let name = expr_name(node)?;
212+
let decoder = registry.decoder::<dyn PhysicalExpr, PhysicalExprDecoder>(name)?;
213+
Some(decoder(node, ctx))
214+
}
215+
216+
/// Every extension expression name registered in `registry`, in arbitrary
217+
/// order.
218+
///
219+
/// Names registered for another kind are not included.
220+
pub fn physical_expr_names(
221+
registry: &ProtoDecoderRegistry,
222+
) -> impl Iterator<Item = &str> {
223+
registry.names::<dyn PhysicalExpr>()
224+
}
225+
226+
/// The registry name `node` was written with, if it is an extension node that
227+
/// carries one.
228+
fn expr_name(node: &PhysicalExprNode) -> Option<&str> {
229+
match node.expr_type.as_ref()? {
230+
ExprType::Extension(extension) => extension.expr_name.as_deref(),
231+
_ => None,
232+
}
233+
}

‎datafusion/physical-plan/src/proto/mod.rs‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -79,7 +79,10 @@ use datafusion_execution::TaskContext;
7979
use datafusion_expr::physical_planning_context::ScalarSubqueryResults;
8080
use datafusion_expr::{AggregateUDF, ScalarUDF, WindowUDF};
8181
use datafusion_physical_expr::PhysicalExpr;
82-
pub use datafusion_physical_expr::proto::{ExprDecodeSession, PhysicalExprFromProto};
82+
pub use datafusion_physical_expr::proto::{
83+
ExprDecodeSession, ExtensionExprName, PhysicalExprFromProto, decode_physical_expr,
84+
physical_expr_names, register_physical_expr,
85+
};
8386
use datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecode;
8487
pub use datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecodeCtx;
8588
use datafusion_physical_expr_common::physical_expr::proto_encode::PhysicalExprEncode;

‎datafusion/proto/src/physical_plan/from_proto.rs‎

Lines changed: 66 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ use arrow::array::RecordBatch;
2323
use arrow::compute::SortOptions;
2424
use arrow::datatypes::{Field, Schema};
2525
use arrow::ipc::reader::StreamReader;
26-
use datafusion_common::{Result, internal_datafusion_err, not_impl_err};
26+
use datafusion_common::{DataFusionError, Result, internal_datafusion_err, not_impl_err};
2727
use datafusion_datasource::TableSchema;
2828
use datafusion_datasource::file::FileSource;
2929
use datafusion_datasource::file_scan_config::FileScanConfig;
@@ -41,7 +41,10 @@ use datafusion_physical_plan::expressions::{
4141
};
4242
use datafusion_physical_plan::joins::HashExpr;
4343
use datafusion_physical_plan::proto::ExecutionPlanDecodeCtx;
44-
use datafusion_physical_plan::proto::{ExprDecodeSession, PhysicalExprFromProto};
44+
use datafusion_physical_plan::proto::{
45+
ExprDecodeSession, PhysicalExprFromProto, ProtoDecoderRegistry, decode_physical_expr,
46+
physical_expr_names,
47+
};
4548
use datafusion_physical_plan::repartition::RangeExpr;
4649
use datafusion_physical_plan::windows::{create_window_expr, schema_add_window_field};
4750
use datafusion_physical_plan::{Partitioning, PhysicalExpr, WindowExpr};
@@ -364,16 +367,42 @@ pub fn parse_physical_expr_with_converter(
364367
SqlSimilarToPattern::try_from_proto(proto, &decode_ctx)?
365368
}
366369
ExprType::Extension(extension) => {
370+
// Lookup-order policy, the same one extension plans use: an
371+
// expression that named itself on the wire and is registered on
372+
// this session decodes itself, no codec involved. Anything else —
373+
// no name, or a name this session has no decoder for — takes the
374+
// codec path, which is also every node written before `expr_name`
375+
// existed. `None` from the registry means no decoder claims the
376+
// name; a decode *failure* is returned as-is rather than falling
377+
// through to the codec.
378+
let registry = ctx
379+
.task_ctx()
380+
.session_config()
381+
.get_extension::<ProtoDecoderRegistry>();
382+
if let Some(registry) = registry.as_ref()
383+
&& let Some(decoded) = decode_physical_expr(registry, proto, &decode_ctx)
384+
{
385+
return decoded;
386+
}
387+
367388
let inputs: Vec<Arc<dyn PhysicalExpr>> = extension
368389
.inputs
369390
.iter()
370391
.map(|e| proto_converter.proto_to_physical_expr(e, input_schema, ctx))
371392
.collect::<Result<_>>()?;
372-
ctx.codec().try_decode_expr(
373-
extension.expr.as_slice(),
374-
&inputs,
375-
&decode_ctx,
376-
)? as _
393+
ctx.codec()
394+
.try_decode_expr(extension.expr.as_slice(), &inputs, &decode_ctx)
395+
.map_err(|e| match extension.expr_name.as_deref() {
396+
// The writer named the expression but this session has no
397+
// decoder for it: say so, rather than leaving only the
398+
// codec's error to explain a missing registration.
399+
Some(expr_name) => unregistered_extension_expr_context(
400+
expr_name,
401+
registry.as_deref(),
402+
e,
403+
),
404+
None => e,
405+
})?
377406
}
378407
ExprType::Lambda(_) => LambdaExpr::try_from_proto(proto, &decode_ctx)?,
379408
ExprType::LambdaVariable(_) => {
@@ -483,6 +512,36 @@ struct ConverterDecoder<'a, 'b> {
483512
proto_converter: &'a dyn PhysicalProtoConverterExtension,
484513
}
485514

515+
/// Error for a `PhysicalExtensionExprNode` that names its expression type but
516+
/// finds no decoder registered for that name.
517+
///
518+
/// Lists what *is* registered, because the usual cause is a session configured
519+
/// on the writing side but not on the reading one. The codec's own error is
520+
/// kept as the cause rather than reworded: the codec may have failed for a
521+
/// reason that has nothing to do with the missing registration, and a caller
522+
/// matching on the error kind must still see the kind the codec chose.
523+
fn unregistered_extension_expr_context(
524+
expr_name: &str,
525+
registry: Option<&ProtoDecoderRegistry>,
526+
codec_error: DataFusionError,
527+
) -> DataFusionError {
528+
let mut registered: Vec<&str> = registry
529+
.map(|registry| physical_expr_names(registry).collect())
530+
.unwrap_or_default();
531+
registered.sort_unstable();
532+
let registered = if registered.is_empty() {
533+
"none".to_string()
534+
} else {
535+
registered.join(", ")
536+
};
537+
codec_error.context(format!(
538+
"No decoder is registered for the extension PhysicalExpr '{expr_name}'. Register \
539+
the expression in the ProtoDecoderRegistry attached to the decoding session's \
540+
SessionConfig, or supply a PhysicalExtensionCodec that handles it. Registered \
541+
extension expressions: {registered}"
542+
))
543+
}
544+
486545
impl datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecode
487546
for ConverterDecoder<'_, '_>
488547
{

‎datafusion/proto/src/physical_plan/mod.rs‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1745,6 +1745,13 @@ pub trait PhysicalExtensionCodec: Debug + Send + Sync + Any {
17451745

17461746
/// Decode a custom extension expression from `buf`.
17471747
///
1748+
/// Expressions can instead name themselves on the wire and be decoded by a
1749+
/// per-type decoder registered on the session — see
1750+
/// [`PhysicalExprFromProto`] and [`register_physical_expr`]. That path
1751+
/// resolves by name rather than by codec registration order, so two crates
1752+
/// claiming the same name collide at registration. This method stays the
1753+
/// fallback for everything unnamed or unregistered.
1754+
///
17481755
/// `inputs` holds the already-decoded children carried in the
17491756
/// `PhysicalExtensionExprNode.inputs` field. If the codec instead embeds
17501757
/// nested `PhysicalExprNode`s *inside* `buf`, decode them through
@@ -1757,6 +1764,8 @@ pub trait PhysicalExtensionCodec: Debug + Send + Sync + Any {
17571764
/// cache-hits on its `expr_id` and re-shares one `Arc<dyn PhysicalExpr>`.
17581765
///
17591766
/// [`parse_physical_expr`]: crate::physical_plan::from_proto::parse_physical_expr
1767+
/// [`PhysicalExprFromProto`]: datafusion_physical_plan::proto::PhysicalExprFromProto
1768+
/// [`register_physical_expr`]: datafusion_physical_plan::proto::register_physical_expr
17601769
fn try_decode_expr(
17611770
&self,
17621771
_buf: &[u8],

0 commit comments

Comments
 (0)