namanjain24-sudo commented on code in PR #25318:
URL: https://github.com/apache/datafusion/pull/25318#discussion_r4205822597
##########
datafusion/ffi/src/proto/physical_extension_codec.rs:
##########
@@ -756,4 +922,219 @@ pub(crate) mod tests {
.expect("rebound codec resolves");
assert_eq!(task_ctx.session_id(), ctx_b.task_ctx().session_id());
}
+
+ /// An extension plan that carries a single physical expression, so its
+ /// codec has to decode that expression itself.
+ #[derive(Debug)]
+ struct ScalarSubqueryExprExec {
+ expr: Arc<dyn PhysicalExpr>,
+ child: Arc<dyn ExecutionPlan>,
+ }
+
+ impl DisplayAs for ScalarSubqueryExprExec {
+ fn fmt_as(
+ &self,
+ _t: DisplayFormatType,
+ f: &mut std::fmt::Formatter,
+ ) -> std::fmt::Result {
+ write!(f, "ScalarSubqueryExprExec")
+ }
+ }
+
+ impl ExecutionPlan for ScalarSubqueryExprExec {
+ fn name(&self) -> &str {
+ "ScalarSubqueryExprExec"
+ }
+
+ fn properties(&self) -> &Arc<PlanProperties> {
+ self.child.properties()
+ }
+
+ fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
+ vec![&self.child]
+ }
+
+ fn apply_expressions(
+ &self,
+ f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) ->
Result<TreeNodeRecursion>,
+ ) -> Result<TreeNodeRecursion> {
+ apply_expression_roots(std::slice::from_ref(&self.expr), f)
+ }
+
+ fn replace_children(
+ self: Arc<Self>,
+ _: Vec<Arc<dyn ExecutionPlan>>,
+ _: ReplaceChildrenOptions,
+ ) -> Result<Arc<dyn ExecutionPlan>> {
+ unreachable!()
+ }
+
+ fn with_new_children(
+ self: Arc<Self>,
+ children: Vec<Arc<dyn ExecutionPlan>>,
+ ) -> Result<Arc<dyn ExecutionPlan>> {
+ self.replace_children(
+ children,
+ ReplaceChildrenOptions::new(ChildrenPropertiesMode::Recompute),
+ )
+ }
+
+ fn execute(
+ &self,
+ _partition: usize,
+ _context: Arc<TaskContext>,
+ ) -> Result<SendableRecordBatchStream> {
+ unreachable!()
+ }
+ }
+
+ #[derive(Clone, PartialEq, Message)]
+ struct ScalarSubqueryExprExecProto {
+ #[prost(message, optional, tag = "1")]
+ expr: Option<PhysicalExprNode>,
+ }
+
+ /// Decodes [`ScalarSubqueryExprExec`] through `try_decode_with_ctx`,
+ /// passing the decode context on to its expression. `try_decode` fails so
+ /// a caller that drops the context (by routing through the plain
+ /// `try_decode` FFI entry point instead of `try_decode_with_ctx`) is
+ /// caught by this test.
+ #[derive(Debug)]
+ struct ScalarSubqueryExprExecCodec;
+
+ impl PhysicalExtensionCodec for ScalarSubqueryExprExecCodec {
+ fn try_decode(
+ &self,
+ _buf: &[u8],
+ _inputs: &[Arc<dyn ExecutionPlan>],
+ _ctx: &TaskContext,
+ _proto_converter: &dyn PhysicalProtoConverterExtension,
+ ) -> Result<Arc<dyn ExecutionPlan>> {
+ exec_err!("ScalarSubqueryExprExecCodec decodes through
try_decode_with_ctx")
+ }
+
+ fn try_decode_with_ctx(
+ &self,
+ buf: &[u8],
+ inputs: &[Arc<dyn ExecutionPlan>],
+ ctx: &PhysicalPlanDecodeContext<'_>,
+ proto_converter: &dyn PhysicalProtoConverterExtension,
+ ) -> Result<Arc<dyn ExecutionPlan>> {
+ let proto = ScalarSubqueryExprExecProto::decode(buf).map_err(|e| {
+ internal_datafusion_err!("failed to decode
ScalarSubqueryExprExec: {e}")
+ })?;
+ let expr_proto = proto.expr.ok_or_else(|| {
+ internal_datafusion_err!("ScalarSubqueryExprExec is missing
its expr")
+ })?;
+ let schema = inputs[0].schema();
+ let expr =
+ proto_converter.proto_to_physical_expr(&expr_proto, &schema,
ctx)?;
+ Ok(Arc::new(ScalarSubqueryExprExec {
+ expr,
+ child: Arc::clone(&inputs[0]),
+ }))
+ }
+
+ fn try_encode(
+ &self,
+ node: Arc<dyn ExecutionPlan>,
+ buf: &mut Vec<u8>,
+ proto_converter: &dyn PhysicalProtoConverterExtension,
+ ) -> Result<()> {
+ let exec =
+ node.downcast_ref::<ScalarSubqueryExprExec>()
+ .ok_or_else(|| {
+ internal_datafusion_err!("expected
ScalarSubqueryExprExec")
+ })?;
+ let proto = ScalarSubqueryExprExecProto {
+ expr: Some(proto_converter.physical_expr_to_proto(&exec.expr,
self)?),
+ };
+ proto.encode(buf).map_err(|e| {
+ internal_datafusion_err!("failed to encode
ScalarSubqueryExprExec: {e}")
+ })
+ }
+ }
+
+ /// The decode context's active scalar subquery results scope must reach a
+ /// `ScalarSubqueryExpr` decoded by a codec that is forced foreign through
+ /// the FFI boundary, not just a local one. Without the fix, the far side
+ /// always decodes with a root context (no scope), so the embedded
+ /// `ScalarSubqueryExpr` fails to deserialize.
+ #[test]
+ fn ffi_physical_extension_codec_forced_foreign_scalar_subquery_roundtrip()
+ -> Result<()> {
+ let schema = Arc::new(Schema::new(vec![Field::new("a",
DataType::Int64, false)]));
+ let subquery_schema =
+ Arc::new(Schema::new(vec![Field::new("x", DataType::Int64,
true)]));
+
+ let results = ScalarSubqueryResults::new(1);
+ let sq_expr: Arc<dyn PhysicalExpr> = Arc::new(ScalarSubqueryExpr::new(
+ DataType::Int64,
+ true,
+ SubqueryIndex::new(0),
+ results.clone(),
+ ));
+ let extension_plan: Arc<dyn ExecutionPlan> =
Arc::new(ScalarSubqueryExprExec {
+ expr: sq_expr,
+ child: Arc::new(RealEmptyExec::new(Arc::clone(&schema))),
+ });
+ let plan: Arc<dyn ExecutionPlan> = Arc::new(ScalarSubqueryExec::new(
+ extension_plan,
+ vec![ScalarSubqueryLink {
+ plan: Arc::new(RealEmptyExec::new(subquery_schema)),
+ index: SubqueryIndex::new(0),
+ }],
+ results,
+ ));
+
+ let bytes = physical_plan_to_bytes_with_proto_converter(
+ Arc::clone(&plan),
+ &ScalarSubqueryExprExecCodec,
+ &DefaultPhysicalProtoConverter {},
+ )?;
+
+ let (ctx, task_ctx_provider) =
crate::util::tests::test_session_and_ctx();
+ let mut ffi_codec = FFI_PhysicalExtensionCodec::new(
+ Arc::new(ScalarSubqueryExprExecCodec),
+ None,
+ task_ctx_provider,
+ );
+ ffi_codec.library_marker_id = crate::mock_foreign_marker_id;
Review Comment:
Added a feature-gated cross-library test
(`datafusion/ffi/tests/ffi_physical_extension_codec.rs`, `integration-tests`
feature): it decodes inside a separately `dlopen`'d copy of the cdylib, so the
results handle genuinely crosses the FFI boundary through
`ForeignScalarSubqueryResultsBackend` instead of the in-process marker mock.
##########
datafusion/ffi/src/proto/physical_extension_codec.rs:
##########
@@ -111,6 +112,21 @@ pub struct FFI_PhysicalExtensionCodec {
/// Utility to identify when FFI objects are accessed locally through
/// the foreign interface.
pub library_marker_id: extern "C" fn() -> usize,
+
+ /// Decode bytes into an execution plan, forwarding the active scalar
+ /// subquery results scope (if any) from the caller's decode context so a
+ /// `ScalarSubqueryExpr` decoded on this side of the boundary shares the
+ /// same populated results as the `ScalarSubqueryExec` that owns them.
+ ///
+ /// Added after the original fields, at the end of the struct, so that
+ /// older code built against this `repr(C)` struct without this field
+ /// still sees every field it knows about at its original offset.
+ try_decode_with_ctx: unsafe extern "C" fn(
Review Comment:
You're right, the struct's size does change. I don't have label permissions
on this repo, so I've updated the PR description to call out the ABI break
separately from the additive Rust API change — could you apply the `api change`
label?
--
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]