Repository navigation
feat: give PhysicalExtensionCodec the full decode context - #25318
namanjain24-sudo wants to merge 6 commits into
Conversation
PhysicalExtensionCodec::try_decode only receives the TaskContext, so a codec that decodes its own nested expressions has to build a fresh PhysicalPlanDecodeContext. That context has no scalar subquery results, so a ScalarSubqueryExpr inside an extension plan fails to deserialize even under the ScalarSubqueryExec that owns the results. Add try_decode_with_ctx, which receives the full PhysicalPlanDecodeContext and by default calls try_decode. Decode extension plans through it, forward it in ComposedPhysicalExtensionCodec, and call it with a root context on the FFI provider side.
a384c28 to
cf5431e
Compare
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #25318 +/- ##
==========================================
+ Coverage 82.33% 82.73% +0.39%
==========================================
Files 1137 1148 +11
Lines 432498 449662 +17164
Branches 432498 449662 +17164
==========================================
+ Hits 356116 372009 +15893
- Misses 54843 54973 +130
- Partials 21539 22680 +1141 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
Bumping this — still open for review whenever someone has bandwidth. |
kosiew
left a comment
There was a problem hiding this comment.
Thanks for working on this. The context-aware decode path looks good for local codecs, but the FFI path still loses the active scalar subquery scope. I left one blocking comment and one small test-naming suggestion.
| let plan = sresult_return!(codec.try_decode( | ||
| // The caller's decode context cannot cross the FFI boundary, so decode | ||
| // with a root context for this side's codec. | ||
| let decode_ctx = PhysicalPlanDecodeContext::new(task_ctx.as_ref(), codec.as_ref()); |
There was a problem hiding this comment.
The FFI path still drops the active ScalarSubqueryResults scope, so a foreign codec decoding an embedded ScalarSubqueryExpr can still hit the original error. Please preserve that scope through an ABI-safe, versioned FFI decode path, without passing PhysicalPlanDecodeContext itself, and add a forced-foreign end-to-end ScalarSubqueryExec roundtrip test.
| ffi_codec.library_marker_id = crate::mock_foreign_marker_id; | ||
| let foreign_codec: Arc<dyn PhysicalExtensionCodec> = (&ffi_codec).into(); | ||
|
|
||
| let returned_exec = foreign_codec.try_decode( |
There was a problem hiding this comment.
Could you clarify the test name or comment to show that it intentionally calls the foreign adapter's legacy try_decode, which then reaches the producer codec's try_decode_with_ctx override? For example, ffi_physical_extension_codec_legacy_decode_uses_context_aware_default would make that behavior clearer.
…ndary ForeignPhysicalExtensionCodec now forwards the active ScalarSubqueryResults scope through an ABI-safe FFI_ScalarSubqueryResults handle instead of dropping it, so a ScalarSubqueryExpr decoded by a genuinely foreign codec shares the same results as the surrounding ScalarSubqueryExec. ScalarSubqueryResults gains an alternate backend representation (ScalarSubqueryResultsBackend) to support this proxy, with ScalarValues crossing the boundary as prost-encoded bytes like the rest of the codec's wire format. Adds a forced-foreign end-to-end ScalarSubqueryExec roundtrip test, and clarifies the legacy-decode test's name/intent.
|
@kosiew addressed: |
- Move the new FFI_PhysicalExtensionCodec::try_decode_with_ctx field to the end of the repr(C) struct instead of inserting it in the middle, so every pre-existing field keeps its original offset. - Add RefUnwindSafe as a supertrait bound on ScalarSubqueryResultsBackend so ScalarSubqueryResults keeps implementing UnwindSafe/RefUnwindSafe, with a regression test pinning it.
|
Fixed both: moved |
CI's clippy (ci/scripts/rust_clippy.sh, not --all-features) flagged encode_scalar_value for taking ScalarValue by value while only borrowing it. Takes &ScalarValue instead.
|
CI clippy caught |
kosiew
left a comment
There was a problem hiding this comment.
Thanks for working through the earlier feedback. The decode-context forwarding and related cleanup look improved, but there is still a type-preservation issue in the new FFI scalar-subquery results bridge that needs to be fixed before this is safe to merge.
|
|
||
| fn encode_scalar_value(value: &ScalarValue) -> Result<SVec<u8>> { | ||
| let proto: datafusion_proto::protobuf::ScalarValue = | ||
| value.try_into().map_err(DataFusionError::from)?; |
There was a problem hiding this comment.
This protobuf conversion does not preserve non-null Float16: it is encoded as Float32Value and comes back as ScalarValue::Float32, so a foreign ScalarSubqueryExpr can return a value whose type no longer matches its declared Float16 output. Please use a type-preserving transport here, or extend the encoding to retain Float16, and add forced-foreign non-null Float16 get/set regression coverage.
There was a problem hiding this comment.
Fixed — the bytes now carry the original DataType (as ArrowType) alongside the encoded value, and decode casts back to it, so Float16 round-trips correctly instead of widening to Float32. Added forced-foreign non-null Float16 get/set regression coverage in scalar_subquery_results.rs.
…uery bridge datafusion_proto::protobuf::ScalarValue has no dedicated variant for every DataType: a non-null Float16 is encoded as Float32Value, so a foreign ScalarSubqueryExpr reading it back got a Float32 instead. Wrap the encoded bytes with the original DataType (as an ArrowType) and cast back to it on decode, so the value crossing FFI always matches the type it was declared with. Adds forced-foreign non-null Float16 get/set regression coverage.
kosiew
left a comment
There was a problem hiding this comment.
Thanks for the updates. The Float16 transport issue is fixed, and the context forwarding is in place. I still have one FFI integration coverage gap that I think should be addressed before merging, plus an ABI metadata/documentation issue.
| None, | ||
| task_ctx_provider, | ||
| ); | ||
| ffi_codec.library_marker_id = crate::mock_foreign_marker_id; |
There was a problem hiding this 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.
There was a problem hiding this 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.
| /// 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( |
There was a problem hiding this 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.
There was a problem hiding this 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?
…bridge The existing unit test only forces the codec foreign through mock_foreign_marker_id, so the results handle still takes the local-bypass fast path and never exercises ForeignScalarSubqueryResultsBackend. Adds an integration-tests-gated test that decodes inside a separately dlopen'd copy of the cdylib, so the handle genuinely crosses the boundary. Shares the scalar-subquery fixtures between the unit test and the new integration test via a fixtures module gated on cfg(any(test, feature = "integration-tests")). Also flags the FFI ABI break from the appended try_decode_with_ctx field for a maintainer to label, since the struct's size changes even though existing try_decode-only implementations keep compiling.
Which issue does this PR close?
PhysicalExtensionCodecshould take aPhysicalPlanDecodeContextinstead of aTaskContext#25089.Rationale for this change
A custom
ExecutionPlanwhose codec decodes its own expressions fails to deserialize when one of those expressions is aScalarSubqueryExpr. It fails even when the plan sits inside theScalarSubqueryExecthat owns the results:The decoder does track the results of the enclosing
ScalarSubqueryExecinPhysicalPlanDecodeContext. ButPhysicalExtensionCodec::try_decodeonly receivesctx.task_ctx(), so a codec has to build a freshPhysicalPlanDecodeContext::new(..)to decode its expressions, and that fresh context has no results.What changes are included in this PR?
This takes the non-breaking option from the issue:
PhysicalExtensionCodec::try_decode_with_ctxreceives the fullPhysicalPlanDecodeContext. Its default implementation callstry_decode(ctx.task_ctx()), so existing codecs compile and behave as before.try_decode_with_ctx.ComposedPhysicalExtensionCodecforwardstry_decode_with_ctxto the codec that encoded the node. Without this, composing codecs would drop the context again.FFI_PhysicalExtensionCodecgains atry_decode_with_ctxfunction pointer that also forwards the caller's active scalar-subquery results scope (if any) across the boundary through a newFFI_ScalarSubqueryResultshandle, so aScalarSubqueryExprdecoded by a foreign codec shares the same populated results as the host'sScalarSubqueryExec. A codec that only implements the new method therefore also works through FFI. For codecs that only implementtry_decode, nothing changes.I went with the new method rather than changing
try_decode's signature because of the API health policy. The breaking change is simpler if reviewers prefer it, and I'm happy to switch.This touches the same function as the draft #24631, which adds a per-type registry in front of the codec. That PR's registry path already passes decoders the full context (through
ConverterPlanDecoder), but its codec fallback still callstry_decode(ctx.task_ctx()), which is the line this PR changes. Whichever lands second needs a one-line rebase.What is the testing strategy for this PR?
New tests in
datafusion/proto/tests/cases/plans/scalar_subquery.rsuse an extension plan whose codec decodes its own expressions throughtry_decode_with_ctx. Itstry_decodereturns an error, so a caller that drops the context fails the test:ScalarSubqueryExprinside the extension plan shares the enclosingScalarSubqueryExec's results, with bothDefaultPhysicalProtoConverterandDeduplicatingProtoConverterComposedPhysicalExtensionCodecScalarSubqueryExecs: each expression binds to its own scope and not to the otherScalarSubqueryExec, decoding still fails with the existing errorA new FFI test checks that a codec that only decodes through
try_decode_with_ctxworks behindFFI_PhysicalExtensionCodec. A second, cross-library integration test (datafusion/ffi/tests/ffi_physical_extension_codec.rs, gated by theintegration-testsfeature) builds that codec inside a separatelydlopen'd copy of thedatafusion-fficdylib, so the scalar-subquery results handle crosses a genuine FFI boundary and exercisesForeignScalarSubqueryResultsBackendfor real, rather than only the in-processmock_foreign_marker_idpath the unit test uses.I checked that each part of the change is needed:
mainwith the error from the issue.try_decode, all four new proto tests fail.ComposedPhysicalExtensionCodec::try_decode_with_ctxremoved, only the composed-codec test fails.Existing suites pass:
datafusion-proto(17 lib, 265 integration, 4 doc tests) anddatafusion-ffi(123 lib tests, plus the new cross-library test under--features integration-tests)../dev/rust_lint.shis clean.Are there any user-facing changes?
This has two separate kinds of API impact:
PhysicalExtensionCodec::try_decode_with_ctxis a new provided method. Existing implementations keep compiling and behave the same.FFI_PhysicalExtensionCodecis#[repr(C)]and gains atry_decode_with_ctxfunction-pointer field, appended after the existing fields so a producer compiled against the old layout still sees every field it knows about at its original offset. The struct's size still changes, so this is an ABI break: a producer and a consumer built against different layouts of this struct must not be mixed — both sides need to be rebuilt against the same DataFusion version. This should carry theapi changelabel; I don't have permission to add labels on this repo, so flagging it here for a maintainer to apply.The same two-sided split applies to the
FFI_ScalarSubqueryResultshandle added for the scalar-subquery results forwarding described above.