kosiew commented on code in PR #25318:
URL: https://github.com/apache/datafusion/pull/25318#discussion_r4203508653
##########
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:
Appending this callback preserves the existing field offsets, but it still
changes the size of this `repr(C)` type and therefore changes the FFI ABI.
Please add the required `api change` label and update the PR description to
distinguish the additive Rust API change from the FFI ABI change, including a
concise rebuild/ABI note.
##########
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:
This test verifies context forwarding, but only forces the codec foreign.
The results handle keeps the local marker, so it bypasses
`ForeignScalarSubqueryResultsBackend`. Please add feature-gated cross-library
coverage that invokes the context-aware callback with an active results handle
and verifies that the foreign-decoded scalar-subquery expression observes a
host-populated value through the remote backend.
--
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]