timsaucer opened a new issue, #24688:
URL: https://github.com/apache/datafusion/issues/24688
## Describe the bug
`FFI_SessionRef` carries a `logical_codec`, but the `create_physical_plan`
round trip is the one place in `datafusion/ffi/src/session/mod.rs` that never
consults it. Both halves serialize with the codec-less helpers:
- `ForeignSession::create_physical_plan` (session/mod.rs:749) calls
`logical_plan_to_bytes(logical_plan)?`
- `create_physical_plan_fn_wrapper` (session/mod.rs:252) calls
`logical_plan_from_bytes(logical_plan_serialized.as_slice(), task_ctx.as_ref())`
Those are the only two codec-less `logical_plan_to_bytes` /
`logical_plan_from_bytes` calls in the entire `datafusion-ffi` crate. The
neighbouring `optimize_fn_wrapper` (session/mod.rs:221) does the same round
trip correctly, which is what makes this look like an oversight rather than a
design decision:
```rust
let logical_codec: Arc<dyn LogicalExtensionCodec> =
(&session.logical_codec).into();
let logical_plan =
sresult_return!(logical_plan_from_bytes_with_extension_codec(
logical_plan_serialized.as_slice(),
inner.task_ctx().as_ref(),
logical_codec.as_ref(),
));
```
The consequence is that a foreign library holding an `FFI_SessionRef` cannot
call `Session::create_physical_plan` for any plan referencing a table provider
the default codec cannot encode. Serialization fails on the caller's side,
before the FFI call is even made:
```
This feature is not implemented: LogicalExtensionCodec is not provided
at datafusion/proto/src/logical_plan/mod.rs:1534 (Error serializing custom
table)
```
This is not limited to query planners. `FFI_SessionRef` is handed to foreign
libraries from `table_provider.rs` (`scan` and `insert_into`), `udtf.rs`, and
`table_provider_factory.rs` as well, so a foreign `TableProvider` that wants to
plan a sub-query during `scan` hits the same wall.
## To Reproduce
Add to the existing `mod tests` in `datafusion/ffi/src/session/mod.rs`. It
uses the same `crate::util::tests::test_session_and_ctx()` helper as
`test_foreign_session_lazily_loads_planning_state`. The codec below can
round-trip the provider, so the only reason serialization can fail is the
session not consulting it. `optimize` succeeds; `create_physical_plan` fails on
the same plan with the same session.
```rust
// `use super::*` already covers several of these; trim to taste.
use datafusion::catalog::{MemTable, TableProvider};
use datafusion_common::{TableReference, not_impl_err};
use datafusion_expr::Extension;
use datafusion_proto::logical_plan::LogicalExtensionCodec;
/// Round-trips the `MemTable` used below, so the only way serialization can
/// fail is the session declining to use it.
///
/// The table provider encoding is deliberately empty:
`try_decode_table_provider`
/// rebuilds a fresh `MemTable` from the schema this codec already holds, so
there
/// is nothing to carry in the buffer. This is a complete codec for the
purposes of
/// this test, not a stub — it round-trips the only provider in the plan.
#[derive(Debug)]
struct MemTableCodec {
schema: SchemaRef,
}
impl LogicalExtensionCodec for MemTableCodec {
fn try_decode(
&self,
_buf: &[u8],
_inputs: &[LogicalPlan],
_ctx: &TaskContext,
) -> Result<Extension> {
not_impl_err!("no extension nodes in this test")
}
fn try_encode(&self, _node: &Extension, _buf: &mut Vec<u8>) ->
Result<()> {
not_impl_err!("no extension nodes in this test")
}
fn try_decode_table_provider(
&self,
_buf: &[u8],
_table_ref: &TableReference,
_schema: SchemaRef,
_ctx: &TaskContext,
) -> Result<Arc<dyn TableProvider>> {
Ok(Arc::new(MemTable::try_new(
Arc::clone(&self.schema),
vec![vec![]],
)?))
}
fn try_encode_table_provider(
&self,
_table_ref: &TableReference,
_node: Arc<dyn TableProvider>,
_buf: &mut Vec<u8>,
) -> Result<()> {
Ok(())
}
}
#[tokio::test]
async fn test_foreign_session_create_physical_plan_uses_logical_codec()
-> Result<(), DataFusionError> {
let (ctx, task_ctx_provider) =
crate::util::tests::test_session_and_ctx();
let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32,
false)]));
ctx.register_table(
"t",
Arc::new(MemTable::try_new(Arc::clone(&schema), vec![vec![]])?),
)?;
let plan = ctx.table("t").await?.into_unoptimized_plan();
// The session is given a codec that handles this provider.
let logical_codec = FFI_LogicalExtensionCodec::new(
Arc::new(MemTableCodec { schema }),
None,
task_ctx_provider,
);
let state = ctx.state();
let local_session = FFI_SessionRef::new(&state, None, logical_codec);
let foreign_session = ForeignSession::try_from(&local_session)?;
// Passes today: `optimize` consults the codec and round-trips the
provider.
foreign_session.optimize(&plan)?;
// Fails today: `create_physical_plan` does not consult it, and dies
// serializing the provider with
// "This feature is not implemented: LogicalExtensionCodec is not
provided"
foreign_session.create_physical_plan(&plan).await?;
Ok(())
}
```
To assert the current behaviour instead of the desired one, replace the last
call with:
```rust
let err = foreign_session
.create_physical_plan(&plan)
.await
.expect_err("create_physical_plan currently ignores the session codec");
assert!(err.to_string().contains("LogicalExtensionCodec is not provided"));
```
I do not have a build of the DataFusion workspace to hand, so the test above
is written against the conventions of the existing tests module rather than
compiled. Imports may need trimming, since `use super::*` already covers some
of them. The behaviour itself is verified, through the real FFI path, as
described below.
## Expected behavior
`Session::create_physical_plan` should serialize and deserialize the logical
plan with the session's `logical_codec`, exactly as `optimize` already does, so
that plans referencing foreign table providers survive the boundary.
## Suggested fix
Consumer side, `ForeignSession::create_physical_plan`:
```rust
let codec: Arc<dyn LogicalExtensionCodec> =
(&self.session.logical_codec).into();
let logical_plan = logical_plan_to_bytes_with_extension_codec(logical_plan,
codec.as_ref())?;
```
Producer side, `create_physical_plan_fn_wrapper`. The codec has to be read
before `let session = session.inner()` shadows the `FFI_SessionRef`:
```rust
let logical_codec: Arc<dyn LogicalExtensionCodec> =
(&session.logical_codec).into();
let session = session.inner();
let logical_plan =
sresult_return!(logical_plan_from_bytes_with_extension_codec(
logical_plan_serialized.as_slice(),
task_ctx.as_ref(),
logical_codec.as_ref(),
));
```
## Additional context
Found while reviewing FFI query planner support in `datafusion-python`
(apache/datafusion-python#1677). Verified against `datafusion-ffi` 55.0.0-rc3
using that project's three-library FFI examples, where the host, a
table-provider library, and a query-planner library are separate `cdylib`s.
With provider and planner codecs installed on the host session and a
provider-owned table registered, the only variable was whether the foreign
planner built the physical plan itself or asked the session to build it.
Planning locally returned `A: [0,1,2]`; switching that single call to
`session.create_physical_plan(logical_plan)` produced the
`LogicalExtensionCodec is not provided` error above. The codecs demonstrably
reach the planner over FFI, since the first variant round-trips the
provider-owned plan through them successfully.
One caveat if you write the regression test around a foreign **query
planner** rather than a provider: `Session::create_physical_plan` is separately
unusable from inside one, and that is already a documented restriction rather
than a second bug. The module docs (session/mod.rs:20-31) and the
`ForeignSession` docs (session/mod.rs:524-531) both say it:
> C must not call `create_physical_plan`, or invoke the planner returned by
`query_planner`, to delegate planning back to A. Repeating either self-call
recurses until the stack is exhausted.
Concretely, `create_physical_plan_fn_wrapper` calls
`session.inner().create_physical_plan(...)`, which is
`SessionState::create_physical_plan`, which dispatches back through
`state.query_planner()` — the same foreign planner that made the call. Observed
as `SIGSEGV`. When a foreign provider is in the plan the codec error above
fires first and masks it, because serialization happens on the caller's side;
with a plan the default codec can serialize, the stack overflow is what you get.
This matters here only because it shapes the test: fixing this issue alone
would turn a planner-based test from a clean error into a crash, which is why
the test above drives it through a plain `ForeignSession` instead. Separately,
and outside the scope of this issue, it may be worth turning that documented
restriction into a guardrail — a recursion depth counter or a re-entrancy flag
that returns a `DataFusionError` naming the delegation rule, instead of
exhausting the stack.
--
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]