@@ -36,7 +36,7 @@ use datafusion_proto::logical_plan::{
3636use datafusion_proto:: protobuf:: LogicalExprList ;
3737use prost:: Message ;
3838
39- use stabby:: vec:: Vec as StabbyVec ;
39+ use stabby:: vec:: Vec as SVec ;
4040use tokio:: runtime:: Handle ;
4141
4242use super :: execution_plan:: FFI_ExecutionPlan ;
@@ -108,8 +108,8 @@ pub struct FFI_TableProvider {
108108 scan : unsafe extern "C" fn (
109109 provider : & Self ,
110110 session : FFI_SessionRef ,
111- projections : FfiOption < StabbyVec < usize > > ,
112- filters_serialized : StabbyVec < u8 > ,
111+ projections : FfiOption < SVec < usize > > ,
112+ filters_serialized : SVec < u8 > ,
113113 limit : FfiOption < usize > ,
114114 ) -> FfiFuture < FFIResult < FFI_ExecutionPlan > > ,
115115
@@ -122,9 +122,8 @@ pub struct FFI_TableProvider {
122122 supports_filters_pushdown : Option <
123123 unsafe extern "C" fn (
124124 provider : & FFI_TableProvider ,
125- filters_serialized : StabbyVec < u8 > ,
126- )
127- -> FFIResult < StabbyVec < FfiTableProviderFilterPushDown > > ,
125+ filters_serialized : SVec < u8 > ,
126+ ) -> FFIResult < SVec < FfiTableProviderFilterPushDown > > ,
128127 > ,
129128
130129 insert_into : unsafe extern "C" fn (
@@ -191,7 +190,7 @@ fn supports_filters_pushdown_internal(
191190 filters_serialized : & [ u8 ] ,
192191 task_ctx : & Arc < TaskContext > ,
193192 codec : & dyn LogicalExtensionCodec ,
194- ) -> Result < StabbyVec < FfiTableProviderFilterPushDown > > {
193+ ) -> Result < SVec < FfiTableProviderFilterPushDown > > {
195194 let filters = match filters_serialized. is_empty ( ) {
196195 true => vec ! [ ] ,
197196 false => {
@@ -203,7 +202,7 @@ fn supports_filters_pushdown_internal(
203202 } ;
204203 let filters_borrowed: Vec < & Expr > = filters. iter ( ) . collect ( ) ;
205204
206- let results: StabbyVec < _ > = provider
205+ let results: SVec < _ > = provider
207206 . supports_filters_pushdown ( & filters_borrowed) ?
208207 . iter ( )
209208 . map ( |v| v. into ( ) )
@@ -214,8 +213,8 @@ fn supports_filters_pushdown_internal(
214213
215214unsafe extern "C" fn supports_filters_pushdown_fn_wrapper (
216215 provider : & FFI_TableProvider ,
217- filters_serialized : StabbyVec < u8 > ,
218- ) -> FFIResult < StabbyVec < FfiTableProviderFilterPushDown > > {
216+ filters_serialized : SVec < u8 > ,
217+ ) -> FFIResult < SVec < FfiTableProviderFilterPushDown > > {
219218 let logical_codec: Arc < dyn LogicalExtensionCodec > = ( & provider. logical_codec ) . into ( ) ;
220219 let task_ctx = rresult_return ! ( <Arc <TaskContext >>:: try_from(
221220 & provider. logical_codec. task_ctx_provider
@@ -233,8 +232,8 @@ unsafe extern "C" fn supports_filters_pushdown_fn_wrapper(
233232unsafe extern "C" fn scan_fn_wrapper (
234233 provider : & FFI_TableProvider ,
235234 session : FFI_SessionRef ,
236- projections : FfiOption < StabbyVec < usize > > ,
237- filters_serialized : StabbyVec < u8 > ,
235+ projections : FfiOption < SVec < usize > > ,
236+ filters_serialized : SVec < u8 > ,
238237 limit : FfiOption < usize > ,
239238) -> FfiFuture < FFIResult < FFI_ExecutionPlan > > {
240239 let task_ctx: Result < Arc < TaskContext > , DataFusionError > =
@@ -270,11 +269,12 @@ unsafe extern "C" fn scan_fn_wrapper(
270269 }
271270 } ;
272271
273- let projections: Vec < _ > = projections. into_iter ( ) . collect ( ) ;
272+ let projections: Option < Vec < usize > > =
273+ projections. into_option ( ) . map ( |p| p. into_iter ( ) . collect ( ) ) ;
274274
275275 let plan = rresult_return ! (
276276 internal_provider
277- . scan( session, Some ( & projections) , & filters, limit. into( ) )
277+ . scan( session, projections. as_ref ( ) , & filters, limit. into( ) )
278278 . await
279279 ) ;
280280
@@ -391,6 +391,9 @@ impl FFI_TableProvider {
391391 runtime : Option < Handle > ,
392392 logical_codec : FFI_LogicalExtensionCodec ,
393393 ) -> Self {
394+ if let Some ( provider) = provider. as_any ( ) . downcast_ref :: < ForeignTableProvider > ( ) {
395+ return provider. 0 . clone ( ) ;
396+ }
394397 let private_data = Box :: new ( ProviderPrivateData { provider, runtime } ) ;
395398
396399 Self {
@@ -462,7 +465,7 @@ impl TableProvider for ForeignTableProvider {
462465 ) -> Result < Arc < dyn ExecutionPlan > > {
463466 let session = FFI_SessionRef :: new ( session, None , self . 0 . logical_codec . clone ( ) ) ;
464467
465- let projections: FfiOption < StabbyVec < usize > > = projection
468+ let projections: FfiOption < SVec < usize > > = projection
466469 . map ( |p| p. iter ( ) . map ( |v| v. to_owned ( ) ) . collect ( ) )
467470 . into ( ) ;
468471
@@ -476,7 +479,7 @@ impl TableProvider for ForeignTableProvider {
476479 let maybe_plan = ( self . 0 . scan ) (
477480 & self . 0 ,
478481 session,
479- projections. unwrap_or_default ( ) ,
482+ projections,
480483 filters_serialized,
481484 limit. into ( ) ,
482485 )
@@ -663,8 +666,9 @@ mod tests {
663666
664667 let provider = Arc :: new ( MemTable :: try_new ( schema, vec ! [ vec![ batch1] ] ) ?) ;
665668
666- let ffi_provider =
669+ let mut ffi_provider =
667670 FFI_TableProvider :: new ( provider, true , None , task_ctx_provider, None ) ;
671+ ffi_provider. library_marker_id = crate :: mock_foreign_marker_id;
668672
669673 let foreign_table_provider: Arc < dyn TableProvider > = ( & ffi_provider) . into ( ) ;
670674
@@ -717,4 +721,62 @@ mod tests {
717721
718722 Ok ( ( ) )
719723 }
724+
725+ #[ tokio:: test]
726+ async fn test_scan_with_none_projection_returns_all_columns ( ) -> Result < ( ) > {
727+ use arrow:: datatypes:: Field ;
728+ use datafusion:: arrow:: array:: Float32Array ;
729+ use datafusion:: arrow:: datatypes:: DataType ;
730+ use datafusion:: arrow:: record_batch:: RecordBatch ;
731+ use datafusion:: datasource:: MemTable ;
732+ use datafusion:: physical_plan:: collect;
733+
734+ let schema = Arc :: new ( Schema :: new ( vec ! [
735+ Field :: new( "a" , DataType :: Float32 , false ) ,
736+ Field :: new( "b" , DataType :: Float32 , false ) ,
737+ Field :: new( "c" , DataType :: Float32 , false ) ,
738+ ] ) ) ;
739+
740+ let batch = RecordBatch :: try_new (
741+ Arc :: clone ( & schema) ,
742+ vec ! [
743+ Arc :: new( Float32Array :: from( vec![ 1.0 , 2.0 ] ) ) ,
744+ Arc :: new( Float32Array :: from( vec![ 3.0 , 4.0 ] ) ) ,
745+ Arc :: new( Float32Array :: from( vec![ 5.0 , 6.0 ] ) ) ,
746+ ] ,
747+ ) ?;
748+
749+ let provider =
750+ Arc :: new ( MemTable :: try_new ( Arc :: clone ( & schema) , vec ! [ vec![ batch] ] ) ?) ;
751+
752+ let ctx = Arc :: new ( SessionContext :: new ( ) ) ;
753+ let task_ctx_provider = Arc :: clone ( & ctx) as Arc < dyn TaskContextProvider > ;
754+ let task_ctx_provider = FFI_TaskContextProvider :: from ( & task_ctx_provider) ;
755+
756+ // Wrap in FFI and force the foreign path (not local bypass)
757+ let mut ffi_provider =
758+ FFI_TableProvider :: new ( provider, true , None , task_ctx_provider, None ) ;
759+ ffi_provider. library_marker_id = crate :: mock_foreign_marker_id;
760+
761+ let foreign_table_provider: Arc < dyn TableProvider > = ( & ffi_provider) . into ( ) ;
762+
763+ // Call scan with projection=None, meaning "return all columns"
764+ let plan = foreign_table_provider
765+ . scan ( & ctx. state ( ) , None , & [ ] , None )
766+ . await ?;
767+ assert_eq ! (
768+ plan. schema( ) . fields( ) . len( ) ,
769+ 3 ,
770+ "scan(projection=None) should return all columns; got {}" ,
771+ plan. schema( ) . fields( ) . len( )
772+ ) ;
773+
774+ // Also verify we can execute and get correct data
775+ let batches = collect ( plan, ctx. task_ctx ( ) ) . await ?;
776+ assert_eq ! ( batches. len( ) , 1 ) ;
777+ assert_eq ! ( batches[ 0 ] . num_columns( ) , 3 ) ;
778+ assert_eq ! ( batches[ 0 ] . num_rows( ) , 2 ) ;
779+
780+ Ok ( ( ) )
781+ }
720782}
0 commit comments