kosiew commented on code in PR #25318: URL: https://github.com/apache/datafusion/pull/25318#discussion_r4192326627
########## datafusion/ffi/src/proto/scalar_subquery_results.rs: ########## @@ -0,0 +1,237 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use std::ffi::c_void; +use std::sync::Arc; + +use datafusion_common::error::Result; +use datafusion_common::{DataFusionError, ScalarValue}; +use datafusion_expr::physical_planning_context::{ + ScalarSubqueryResults, ScalarSubqueryResultsBackend, SubqueryIndex, +}; +use prost::Message; + +use stabby::vec::Vec as SVec; + +use crate::util::{FFI_Option, FFI_Result}; +use crate::{df_result, sresult, sresult_return}; + +/// A stable struct for sharing the shared results container of a +/// `ScalarSubqueryExec` across FFI boundaries. +/// +/// Unlike most FFI wrappers in this crate, which stand in for a whole +/// `dyn Trait` object, this one stands in for a single +/// [`ScalarSubqueryResults`] value: the active scope that a +/// `PhysicalExtensionCodec` on the near side of the boundary passes to +/// `try_decode_with_ctx` so a `ScalarSubqueryExpr` decoded on the far side +/// reads the same populated values as the `ScalarSubqueryExec` that owns +/// them. +/// +/// `ScalarValue`s cross the boundary as prost-encoded bytes, the same wire +/// format `PhysicalExtensionCodec::try_decode`/`try_encode` already use, so +/// this struct never assumes a stable memory layout for `ScalarValue` itself. +#[repr(C)] +#[derive(Debug)] +pub struct FFI_ScalarSubqueryResults { + get: unsafe extern "C" fn(&Self, index: u64) -> FFI_Result<FFI_Option<SVec<u8>>>, + + set: unsafe extern "C" fn(&Self, index: u64, value: SVec<u8>) -> FFI_Result<()>, + + clear: unsafe extern "C" fn(&Self), + + /// Used to create a clone on the provider of the results container. This + /// should only need to be called by the receiver of the handle. + clone: unsafe extern "C" fn(&Self) -> Self, + + /// Release the memory of the private data when it is no longer being used. + release: unsafe extern "C" fn(&mut Self), + + private_data: *mut c_void, + + /// Utility to identify when FFI objects are accessed locally through + /// the foreign interface. + library_marker_id: extern "C" fn() -> usize, +} + +unsafe impl Send for FFI_ScalarSubqueryResults {} +unsafe impl Sync for FFI_ScalarSubqueryResults {} + +impl FFI_ScalarSubqueryResults { + fn inner(&self) -> &ScalarSubqueryResults { + let private_data = self.private_data as *const ScalarSubqueryResults; + unsafe { &*private_data } + } + + /// Creates a new [`FFI_ScalarSubqueryResults`] sharing `results`. + pub fn new(results: ScalarSubqueryResults) -> Self { + let private_data = Box::new(results); + + Self { + get: get_fn_wrapper, + set: set_fn_wrapper, + clear: clear_fn_wrapper, + clone: clone_fn_wrapper, + release: release_fn_wrapper, + private_data: Box::into_raw(private_data).cast::<c_void>(), + library_marker_id: crate::get_library_marker_id, + } + } +} + +fn encode_scalar_value(value: &ScalarValue) -> Result<SVec<u8>> { + let proto: datafusion_proto::protobuf::ScalarValue = + value.try_into().map_err(DataFusionError::from)?; Review 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. -- 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]
