Skip to content

feat: give PhysicalExtensionCodec the full decode context - #25318

Open
namanjain24-sudo wants to merge 6 commits into
apache:mainfrom
namanjain24-sudo:extension-codec-decode-context
Open

namanjain24-sudo wants to merge 6 commits into
apache:mainfrom
namanjain24-sudo:extension-codec-decode-context

Conversation

@namanjain24-sudo

@namanjain24-sudo namanjain24-sudo commented Sep 15, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Rationale for this change

A custom ExecutionPlan whose codec decodes its own expressions fails to deserialize when one of those expressions is a ScalarSubqueryExpr. It fails even when the plan sits inside the ScalarSubqueryExec that owns the results:

Internal("ScalarSubqueryExpr can only be deserialized as part of a surrounding ScalarSubqueryExec")

The decoder does track the results of the enclosing ScalarSubqueryExec in PhysicalPlanDecodeContext. But PhysicalExtensionCodec::try_decode only receives ctx.task_ctx(), so a codec has to build a fresh PhysicalPlanDecodeContext::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_ctx receives the full PhysicalPlanDecodeContext. Its default implementation calls try_decode(ctx.task_ctx()), so existing codecs compile and behave as before.
  • Extension plans are now decoded through try_decode_with_ctx.
  • ComposedPhysicalExtensionCodec forwards try_decode_with_ctx to the codec that encoded the node. Without this, composing codecs would drop the context again.
  • On the FFI side, FFI_PhysicalExtensionCodec gains a try_decode_with_ctx function pointer that also forwards the caller's active scalar-subquery results scope (if any) across the boundary through a new FFI_ScalarSubqueryResults handle, so a ScalarSubqueryExpr decoded by a foreign codec shares the same populated results as the host's ScalarSubqueryExec. A codec that only implements the new method therefore also works through FFI. For codecs that only implement try_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 calls try_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.rs use an extension plan whose codec decodes its own expressions through try_decode_with_ctx. Its try_decode returns an error, so a caller that drops the context fails the test:

  • a ScalarSubqueryExpr inside the extension plan shares the enclosing ScalarSubqueryExec's results, with both DefaultPhysicalProtoConverter and DeduplicatingProtoConverter
  • the same through ComposedPhysicalExtensionCodec
  • nested ScalarSubqueryExecs: each expression binds to its own scope and not to the other
  • without an enclosing ScalarSubqueryExec, decoding still fails with the existing error

A new FFI test checks that a codec that only decodes through try_decode_with_ctx works behind FFI_PhysicalExtensionCodec. A second, cross-library integration test (datafusion/ffi/tests/ffi_physical_extension_codec.rs, gated by the integration-tests feature) builds that codec inside a separately dlopen'd copy of the datafusion-ffi cdylib, so the scalar-subquery results handle crosses a genuine FFI boundary and exercises ForeignScalarSubqueryResultsBackend for real, rather than only the in-process mock_foreign_marker_id path the unit test uses.

I checked that each part of the change is needed:

  • The first test fails on main with the error from the issue.
  • With the call site reverted to try_decode, all four new proto tests fail.
  • With ComposedPhysicalExtensionCodec::try_decode_with_ctx removed, only the composed-codec test fails.
  • With the FFI wrapper reverted, the new FFI test fails.

Existing suites pass: datafusion-proto (17 lib, 265 integration, 4 doc tests) and datafusion-ffi (123 lib tests, plus the new cross-library test under --features integration-tests). ./dev/rust_lint.sh is clean.

Are there any user-facing changes?

This has two separate kinds of API impact:

  • Rust API (non-breaking): PhysicalExtensionCodec::try_decode_with_ctx is a new provided method. Existing implementations keep compiling and behave the same.
  • FFI ABI (breaking): FFI_PhysicalExtensionCodec is #[repr(C)] and gains a try_decode_with_ctx function-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 the api change label; 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_ScalarSubqueryResults handle added for the scalar-subquery results forwarding described above.

@github-actions github-actions Bot added proto Related to proto crate ffi Changes to the ffi crate labels Sep 15, 2026
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.
@namanjain24-sudo
namanjain24-sudo force-pushed the extension-codec-decode-context branch from a384c28 to cf5431e Compare September 18, 2026 02:58
@codecov-commenter

codecov-commenter commented Sep 18, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 72.49467% with 129 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.73%. Comparing base (edc936f) to head (851d8d0).
⚠️ Report is 272 commits behind head on main.

Files with missing lines Patch % Lines
...tafusion/ffi/src/proto/physical_extension_codec.rs 66.00% 74 Missing and 12 partials ⚠️
...atafusion/ffi/src/proto/scalar_subquery_results.rs 80.57% 12 Missing and 15 partials ⚠️
datafusion/expr/src/physical_planning_context.rs 77.96% 12 Missing and 1 partial ⚠️
datafusion/proto/src/physical_plan/mod.rs 83.33% 3 Missing ⚠️
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.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@namanjain24-sudo

Copy link
Copy Markdown
Contributor Author

Bumping this — still open for review whenever someone has bandwidth.

@kosiew kosiew left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@namanjain24-sudo,

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());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.
@namanjain24-sudo

Copy link
Copy Markdown
Contributor Author

@kosiew addressed: ForeignPhysicalExtensionCodec now forwards the active scalar subquery results scope through an ABI-safe FFI_ScalarSubqueryResults handle (values cross as prost-encoded bytes), and added a forced-foreign end-to-end ScalarSubqueryExec roundtrip test. Also renamed the legacy-decode test per your suggestion.

@github-actions github-actions Bot added logical-expr Logical plan and expressions auto detected api change Auto detected API change labels Oct 2, 2026
- 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.
@namanjain24-sudo

Copy link
Copy Markdown
Contributor Author

Fixed both: moved try_decode_with_ctx to the end of the repr(C) struct so no existing field shifts, and added RefUnwindSafe as a supertrait bound on ScalarSubqueryResultsBackend so ScalarSubqueryResults keeps its auto traits (pinned with a test).

@github-actions github-actions Bot removed the auto detected api change Auto detected API change label Oct 2, 2026
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.
@namanjain24-sudo

Copy link
Copy Markdown
Contributor Author

CI clippy caught needless_pass_by_value on encode_scalar_value (it only borrowed the ScalarValue) — fixed, and verified with ./ci/scripts/rust_clippy.sh locally (clean across the whole workspace).

@kosiew kosiew left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@namanjain24-sudo,

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)?;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 kosiew left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@namanjain24-sudo,

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;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

ffi Changes to the ffi crate logical-expr Logical plan and expressions proto Related to proto crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

PhysicalExtensionCodec should take a PhysicalPlanDecodeContext instead of a TaskContext

3 participants