Skip to content

Commit 7f9361a

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 implements `ExtensionExprFromProto`, which pairs the name it is dispatched by with the constructor that rebuilds it. Its `try_to_proto` writes itself with `extension_expr_node`, which stamps `Self::NAME`, so the encoded name and the registry key cannot drift apart. `register_physical_expr` puts its decoder in the `ProtoDecoderRegistry` the decoding session carries — the same registry extension `ExecutionPlan`s use, keyed on the trait as well as the name. Only extension expressions implement the trait. A built-in has an `ExprType` variant of its own, keeps its inherent `try_from_proto`, and has no wire name to be registered under — so `register_physical_expr::<Column>(..)` does not compile. Nothing about built-in serialization changes, and no caller needs a new import. `extension_expr_node` is a free function rather than a method on `PhysicalExprEncodeCtx` because that context lives in `datafusion-physical-expr-common`, below the crate that can name `ExtensionExprFromProto` — the same layering that makes the decode context generic over its session. 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 c30bb2a commit 7f9361a

7 files changed

Lines changed: 1190 additions & 9 deletions

File tree

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

Lines changed: 217 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,12 @@
1515
// specific language governing permissions and limitations
1616
// under the License.
1717

18-
//! The session side of decoding a `PhysicalExpr` from protobuf.
18+
//! The session an expression decode runs under, and the `PhysicalExpr`
19+
//! face of the session-scoped [`ProtoDecoderRegistry`] that routes extension
20+
//! expressions to their decoders by name — [`register_physical_expr`],
21+
//! [`decode_physical_expr`] and [`physical_expr_names`], all keyed on
22+
//! `dyn PhysicalExpr`. The store itself is shared with every other extension
23+
//! kind, so a session carries one registry.
1924
//!
2025
//! The decode context in `datafusion-physical-expr-common` is generic over its
2126
//! session type: a decoder may need the session, and what it needs —
@@ -25,9 +30,17 @@
2530
//! expressions decode under it. `datafusion-proto` builds one from the
2631
//! `TaskContext` it decodes under.
2732
33+
use std::sync::Arc;
34+
35+
use datafusion_common::Result;
2836
use datafusion_common::config::ConfigOptions;
2937
use datafusion_expr::registry::FunctionRegistry;
38+
use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
3039
pub use datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecodeCtx;
40+
pub use datafusion_physical_expr_common::physical_expr::proto_encode::PhysicalExprEncodeCtx;
41+
pub use datafusion_proto_models::ProtoDecoderRegistry;
42+
use datafusion_proto_models::protobuf::PhysicalExprNode;
43+
use datafusion_proto_models::protobuf::physical_expr_node::ExprType;
3144

3245
/// What an expression decoder can see of the session it runs under: the
3346
/// function registry, for resolving UDFs by name, and the configuration
@@ -68,3 +81,206 @@ impl std::fmt::Debug for ExprDecodeSession<'_> {
6881
f.debug_struct("ExprDecodeSession").finish_non_exhaustive()
6982
}
7083
}
84+
/// The wire name of an extension [`PhysicalExpr`], and the constructor that
85+
/// rebuilds it.
86+
///
87+
/// Only extension expressions implement this. A built-in expression has an
88+
/// `ExprType` variant of its own, keeps the inherent `try_from_proto`
89+
/// dispatched from that variant, and has no wire name to be registered under —
90+
/// so [`register_physical_expr`] will not accept one.
91+
///
92+
/// One impl block says what the expression is called and how it decodes; its
93+
/// `try_to_proto` on [`PhysicalExpr`] writes the matching node with
94+
/// [`extension_expr_node`]:
95+
///
96+
/// ```ignore
97+
/// impl PhysicalExpr for MyExpr {
98+
/// // ...
99+
/// fn try_to_proto(
100+
/// &self,
101+
/// ctx: &PhysicalExprEncodeCtx<'_>,
102+
/// ) -> Result<Option<PhysicalExprNode>> {
103+
/// // Stamps `Self::NAME`, so the encoded name and the registry key
104+
/// // cannot drift apart.
105+
/// Ok(Some(extension_expr_node::<Self>(
106+
/// ctx,
107+
/// self.state.to_bytes()?,
108+
/// self.children(),
109+
/// )?))
110+
/// }
111+
/// }
112+
///
113+
/// impl ExtensionExprFromProto for MyExpr {
114+
/// // Namespace the name with the owning crate, so a collision with
115+
/// // another crate is an error at registration and not a wrong decode.
116+
/// const NAME: &'static str = "my_crate.MyExpr";
117+
///
118+
/// fn try_from_proto(
119+
/// node: &PhysicalExprNode,
120+
/// ctx: &PhysicalExprDecodeCtx<'_, ExprDecodeSession<'_>>,
121+
/// ) -> Result<Arc<dyn PhysicalExpr>> {
122+
/// let extension = expect_expr_variant!(
123+
/// node,
124+
/// physical_expr_node::ExprType::Extension,
125+
/// "Extension",
126+
/// );
127+
/// let state = MyState::from_bytes(&extension.expr)?;
128+
/// let children = ctx.decode_children_expressions(&extension.inputs)?;
129+
/// Ok(Arc::new(MyExpr::new(state, children)))
130+
/// }
131+
/// }
132+
///
133+
/// let mut registry = ProtoDecoderRegistry::new();
134+
/// register_physical_expr::<MyExpr>(&mut registry)?;
135+
/// let config = SessionConfig::new().with_extension(Arc::new(registry));
136+
/// ```
137+
pub trait ExtensionExprFromProto: PhysicalExpr + Sized {
138+
/// The name this expression type is written and registered under.
139+
///
140+
/// Namespace it with the owning crate — `"my_crate.MyExpr"`, not
141+
/// `"MyExpr"` — so that two independent crates registering into the same
142+
/// session collide at registration time instead of silently decoding each
143+
/// other's nodes. Never write it by hand on the wire:
144+
/// [`extension_expr_node`] stamps it for you.
145+
const NAME: &'static str;
146+
147+
/// Rebuild the expression from the `PhysicalExprNode` its
148+
/// [`PhysicalExpr::try_to_proto`] wrote.
149+
///
150+
/// Takes the whole node — the exact inverse of what `try_to_proto`
151+
/// returns — so the constructor can also see outer-node fields such as
152+
/// `expr_id`. `ctx.session()` gives the function registry and the
153+
/// configuration options.
154+
fn try_from_proto(
155+
node: &PhysicalExprNode,
156+
ctx: &PhysicalExprDecodeCtx<'_, ExprDecodeSession<'_>>,
157+
) -> Result<Arc<dyn PhysicalExpr>>;
158+
}
159+
160+
/// Build the `PhysicalExprNode` for an extension expression `T`: an
161+
/// `Extension` variant carrying `payload`, the encoded `children`, and `T`'s
162+
/// [`NAME`](ExtensionExprFromProto::NAME).
163+
///
164+
/// The one supported way for an extension expression to write itself. It is a
165+
/// free function rather than a method on [`PhysicalExprEncodeCtx`] because the
166+
/// context lives in `datafusion-physical-expr-common`, below the crate that
167+
/// can name [`ExtensionExprFromProto`] — the same layering that makes the
168+
/// decode context generic over its session.
169+
pub fn extension_expr_node<T: ExtensionExprFromProto>(
170+
ctx: &PhysicalExprEncodeCtx<'_>,
171+
payload: Vec<u8>,
172+
children: Vec<&Arc<dyn PhysicalExpr>>,
173+
) -> Result<PhysicalExprNode> {
174+
Ok(PhysicalExprNode {
175+
expr_type: Some(ExprType::Extension(
176+
datafusion_proto_models::protobuf::PhysicalExtensionExprNode {
177+
expr: payload,
178+
inputs: ctx.encode_children_expressions(children)?,
179+
expr_name: Some(T::NAME.to_string()),
180+
},
181+
)),
182+
..Default::default()
183+
})
184+
}
185+
186+
/// How this facade stores a decoder in the shared registry: a function pointer
187+
/// to the monomorphized [`ExtensionExprFromProto::try_from_proto`].
188+
///
189+
/// Deliberately private, and the same type on both the
190+
/// [`register_physical_expr`] and the [`decode_physical_expr`] side.
191+
/// [`ExtensionExprFromProto`] is the public contract and
192+
/// [`decode_physical_expr`] is the public way to invoke one, so this can become
193+
/// something else — a `dyn` decoder object, to admit stateful or closure
194+
/// decoders, which is what an FFI decoder needs — without a breaking change.
195+
type PhysicalExprDecoder = fn(
196+
&PhysicalExprNode,
197+
&PhysicalExprDecodeCtx<'_, ExprDecodeSession<'_>>,
198+
) -> Result<Arc<dyn PhysicalExpr>>;
199+
200+
/// Register `T` in `registry` under its [`ExtensionExprFromProto::NAME`].
201+
///
202+
/// Registering the same type twice is a no-op. Registering a *different* type
203+
/// under a name already taken is an error, so collisions surface here rather
204+
/// than as a wrong decode later. The name is scoped to `dyn PhysicalExpr`, so a
205+
/// plan or a data source may use the same name in the same registry.
206+
///
207+
/// # Who builds the registry
208+
///
209+
/// The application that owns the session builds it. A session carries at most
210+
/// one registry, and `SessionConfig::with_extension` replaces what is there, so
211+
/// a library must never attach a registry of its own: it would silently discard
212+
/// another library's. A library exposes a function that fills a registry it is
213+
/// handed, and the application composes them. One registry holds every kind, so
214+
/// a library registers its plans and its expressions into the same object.
215+
///
216+
/// # A built-in expression cannot be registered
217+
///
218+
/// The bound is [`ExtensionExprFromProto`], which only extension expressions
219+
/// implement. A built-in keeps an inherent `try_from_proto`, dispatched from
220+
/// its own `ExprType` variant:
221+
///
222+
/// ```
223+
/// use datafusion_physical_expr::expressions::Column;
224+
///
225+
/// let _ = Column::try_from_proto;
226+
/// ```
227+
///
228+
/// but it is not an extension, and registering it does not compile:
229+
///
230+
/// ```compile_fail
231+
/// use datafusion_physical_expr::expressions::Column;
232+
/// use datafusion_physical_expr::proto::register_physical_expr;
233+
/// use datafusion_proto_models::ProtoDecoderRegistry;
234+
///
235+
/// let mut registry = ProtoDecoderRegistry::new();
236+
/// register_physical_expr::<Column>(&mut registry).unwrap();
237+
/// ```
238+
pub fn register_physical_expr<T: ExtensionExprFromProto>(
239+
registry: &mut ProtoDecoderRegistry,
240+
) -> Result<()> {
241+
registry.register_decoder::<dyn PhysicalExpr, T, PhysicalExprDecoder>(
242+
T::NAME,
243+
T::try_from_proto,
244+
)
245+
}
246+
247+
/// Decode `node` with the extension expression decoder registered under the
248+
/// name `node` carries.
249+
///
250+
/// The name is read from the node's `Extension` variant rather than passed in,
251+
/// so a caller cannot pair a node with a name it does not carry.
252+
///
253+
/// `None` means "this node names no extension expression decoder of ours": it
254+
/// is not an extension node, it carries no name, or no registered name matches.
255+
/// The caller then falls back to the `PhysicalExtensionCodec` chain.
256+
/// `Some(Err(..))` means the decoder that *does* own the name failed, which is
257+
/// fatal: falling back there would let a codec decode the payload wrongly, the
258+
/// very thing the name exists to prevent.
259+
pub fn decode_physical_expr(
260+
registry: &ProtoDecoderRegistry,
261+
node: &PhysicalExprNode,
262+
ctx: &PhysicalExprDecodeCtx<'_, ExprDecodeSession<'_>>,
263+
) -> Option<Result<Arc<dyn PhysicalExpr>>> {
264+
let name = expr_name(node)?;
265+
let decoder = registry.decoder::<dyn PhysicalExpr, PhysicalExprDecoder>(name)?;
266+
Some(decoder(node, ctx))
267+
}
268+
269+
/// Every extension expression name registered in `registry`, in arbitrary
270+
/// order.
271+
///
272+
/// Names registered for another kind are not included.
273+
pub fn physical_expr_names(
274+
registry: &ProtoDecoderRegistry,
275+
) -> impl Iterator<Item = &str> {
276+
registry.names::<dyn PhysicalExpr>()
277+
}
278+
279+
/// The registry name `node` was written with, if it is an extension node that
280+
/// carries one.
281+
fn expr_name(node: &PhysicalExprNode) -> Option<&str> {
282+
match node.expr_type.as_ref()? {
283+
ExprType::Extension(extension) => extension.expr_name.as_deref(),
284+
_ => None,
285+
}
286+
}

‎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;
82+
pub use datafusion_physical_expr::proto::{
83+
ExprDecodeSession, ExtensionExprFromProto, decode_physical_expr, extension_expr_node,
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: 65 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,9 @@ 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;
44+
use datafusion_physical_plan::proto::{
45+
ExprDecodeSession, ProtoDecoderRegistry, decode_physical_expr, physical_expr_names,
46+
};
4547
use datafusion_physical_plan::repartition::RangeExpr;
4648
use datafusion_physical_plan::windows::{create_window_expr, schema_add_window_field};
4749
use datafusion_physical_plan::{Partitioning, PhysicalExpr, WindowExpr};
@@ -364,16 +366,42 @@ pub fn parse_physical_expr_with_converter(
364366
SqlSimilarToPattern::try_from_proto(proto, &decode_ctx)?
365367
}
366368
ExprType::Extension(extension) => {
369+
// Lookup-order policy, the same one extension plans use: an
370+
// expression that named itself on the wire and is registered on
371+
// this session decodes itself, no codec involved. Anything else —
372+
// no name, or a name this session has no decoder for — takes the
373+
// codec path, which is also every node written before `expr_name`
374+
// existed. `None` from the registry means no decoder claims the
375+
// name; a decode *failure* is returned as-is rather than falling
376+
// through to the codec.
377+
let registry = ctx
378+
.task_ctx()
379+
.session_config()
380+
.get_extension::<ProtoDecoderRegistry>();
381+
if let Some(registry) = registry.as_ref()
382+
&& let Some(decoded) = decode_physical_expr(registry, proto, &decode_ctx)
383+
{
384+
return decoded;
385+
}
386+
367387
let inputs: Vec<Arc<dyn PhysicalExpr>> = extension
368388
.inputs
369389
.iter()
370390
.map(|e| proto_converter.proto_to_physical_expr(e, input_schema, ctx))
371391
.collect::<Result<_>>()?;
372-
ctx.codec().try_decode_expr(
373-
extension.expr.as_slice(),
374-
&inputs,
375-
&decode_ctx,
376-
)? as _
392+
ctx.codec()
393+
.try_decode_expr(extension.expr.as_slice(), &inputs, &decode_ctx)
394+
.map_err(|e| match extension.expr_name.as_deref() {
395+
// The writer named the expression but this session has no
396+
// decoder for it: say so, rather than leaving only the
397+
// codec's error to explain a missing registration.
398+
Some(expr_name) => unregistered_extension_expr_context(
399+
expr_name,
400+
registry.as_deref(),
401+
e,
402+
),
403+
None => e,
404+
})?
377405
}
378406
ExprType::Lambda(_) => LambdaExpr::try_from_proto(proto, &decode_ctx)?,
379407
ExprType::LambdaVariable(_) => {
@@ -483,6 +511,36 @@ struct ConverterDecoder<'a, 'b> {
483511
proto_converter: &'a dyn PhysicalProtoConverterExtension,
484512
}
485513

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

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

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

17501750
/// Decode a custom extension expression from `buf`.
17511751
///
1752+
/// Expressions can instead name themselves on the wire and be decoded by a
1753+
/// per-type decoder registered on the session — see
1754+
/// [`PhysicalExprFromProto`] and [`register_physical_expr`]. That path
1755+
/// resolves by name rather than by codec registration order, so two crates
1756+
/// claiming the same name collide at registration. This method stays the
1757+
/// fallback for everything unnamed or unregistered.
1758+
///
17521759
/// `inputs` holds the already-decoded children carried in the
17531760
/// `PhysicalExtensionExprNode.inputs` field. If the codec instead embeds
17541761
/// nested `PhysicalExprNode`s *inside* `buf`, decode them through
@@ -1761,6 +1768,8 @@ pub trait PhysicalExtensionCodec: Debug + Send + Sync + Any {
17611768
/// cache-hits on its `expr_id` and re-shares one `Arc<dyn PhysicalExpr>`.
17621769
///
17631770
/// [`parse_physical_expr`]: crate::physical_plan::from_proto::parse_physical_expr
1771+
/// [`PhysicalExprFromProto`]: datafusion_physical_plan::proto::PhysicalExprFromProto
1772+
/// [`register_physical_expr`]: datafusion_physical_plan::proto::register_physical_expr
17641773
fn try_decode_expr(
17651774
&self,
17661775
_buf: &[u8],

0 commit comments

Comments
 (0)