diff --git a/vortex-array/src/arrays/validation_tests.rs b/vortex-array/src/arrays/validation_tests.rs index 30758a5c29e..fc7a6ee0a30 100644 --- a/vortex-array/src/arrays/validation_tests.rs +++ b/vortex-array/src/arrays/validation_tests.rs @@ -16,6 +16,8 @@ mod tests { use vortex_error::VortexError; use crate::IntoArray; + use crate::VortexSessionExecute; + use crate::array_session; use crate::arrays::ChunkedArray; use crate::arrays::DecimalArray; use crate::arrays::FixedSizeListArray; @@ -171,6 +173,7 @@ mod tests { Arc::new([]), DType::Utf8(Nullability::NonNullable), Validity::NonNullable, + &mut array_session().create_execution_ctx(), ); assert!(result.is_ok()); } @@ -190,6 +193,7 @@ mod tests { buffers, DType::Binary(Nullability::NonNullable), Validity::NonNullable, + &mut array_session().create_execution_ctx(), ); assert!(matches!(result, Err(VortexError::InvalidArgument(_, _)))); diff --git a/vortex-array/src/arrays/varbinview/array.rs b/vortex-array/src/arrays/varbinview/array.rs index 9e97257b3be..8a2daee623d 100644 --- a/vortex-array/src/arrays/varbinview/array.rs +++ b/vortex-array/src/arrays/varbinview/array.rs @@ -19,6 +19,7 @@ use vortex_error::vortex_panic; use crate::ArrayRef; use crate::ArraySlots; +use crate::ExecutionCtx; use crate::VortexSessionExecute; use crate::array::Array; use crate::array::ArrayParts; @@ -146,8 +147,9 @@ impl VarBinViewData { buffers: Arc<[ByteBuffer]>, dtype: DType, validity: Validity, + ctx: &mut ExecutionCtx, ) -> Self { - Self::try_new(views, buffers, dtype, validity) + Self::try_new(views, buffers, dtype, validity, ctx) .vortex_expect("VarBinViewArray construction failed") } @@ -180,13 +182,42 @@ impl VarBinViewData { buffers: Arc<[ByteBuffer]>, dtype: DType, validity: Validity, + ctx: &mut ExecutionCtx, ) -> VortexResult { + let views = Self::replace_invalid_views(views, &validity, ctx)?; Self::validate(&views, &buffers, &dtype, &validity)?; // SAFETY: validate ensures all invariants are met. Ok(unsafe { Self::new_unchecked(views, buffers, dtype, validity) }) } + // Replace views at invalid (null) slots with empty views + pub(crate) fn replace_invalid_views( + views: Buffer, + validity: &Validity, + ctx: &mut ExecutionCtx, + ) -> VortexResult> { + let mask = validity.execute_mask(views.len(), ctx)?; + if mask.all_true() { + return Ok(views); + } + let empty = BinaryView::empty_view(); + let Some(first) = views + .iter() + .zip(mask.iter()) + .position(|(view, valid)| !valid && *view != empty) + else { + return Ok(views); + }; + let mut views = views.into_mut(); + for (view, valid) in views.iter_mut().zip(mask.iter()).skip(first) { + if !valid { + *view = empty; + } + } + Ok(views.freeze()) + } + /// Constructs a new `VarBinViewArray`. /// /// See `VarBinViewArray::new_unchecked` for more information. @@ -672,8 +703,9 @@ impl Array { buffers: Arc<[ByteBuffer]>, dtype: DType, validity: Validity, + ctx: &mut ExecutionCtx, ) -> VortexResult { - let data = VarBinViewData::try_new(views, buffers, dtype.clone(), validity.clone())?; + let data = VarBinViewData::try_new(views, buffers, dtype.clone(), validity.clone(), ctx)?; let slots = VarBinViewData::make_slots(&validity, data.len()); Ok(Self::from_prevalidated_data(dtype, data, slots)) } diff --git a/vortex-array/src/arrays/varbinview/tests.rs b/vortex-array/src/arrays/varbinview/tests.rs index 981d4f602a1..33289df51ff 100644 --- a/vortex-array/src/arrays/varbinview/tests.rs +++ b/vortex-array/src/arrays/varbinview/tests.rs @@ -1,11 +1,30 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors +use std::sync::Arc; + +use vortex_buffer::BitBuffer; +use vortex_buffer::Buffer; +use vortex_buffer::ByteBuffer; +use vortex_buffer::ByteBufferMut; +use vortex_error::VortexResult; +use vortex_error::vortex_err; +use vortex_session::registry::ReadContext; + +use crate::ArrayContext; +use crate::IntoArray; use crate::VortexSessionExecute; use crate::array_session; +use crate::arrays::VarBinView; use crate::arrays::VarBinViewArray; use crate::arrays::varbinview::BinaryView; +use crate::arrays::varbinview::VarBinViewData; use crate::assert_arrays_eq; +use crate::dtype::DType; +use crate::dtype::Nullability; +use crate::serde::SerializeOptions; +use crate::serde::SerializedArray; +use crate::validity::Validity; #[test] pub fn varbin_view() { @@ -49,3 +68,72 @@ pub fn binary_view_size_and_alignment() { assert_eq!(size_of::(), 16); assert_eq!(align_of::(), 16); } + +#[test] +pub fn replace_invalid_views() -> VortexResult<()> { + let mut ctx = array_session().create_execution_ctx(); + let views = Buffer::::copy_from(vec![ + BinaryView::new_inlined(b"ololo"), + BinaryView::new_ref(13, *b"AAAA", 0xDEAD_BEEF, 0xF000_0000), + ]); + let buffer = BitBuffer::from_iter([true, false]); + let validity = Validity::from_bit_buffer(buffer, Nullability::Nullable); + + let replaced = VarBinViewData::replace_invalid_views(views.clone(), &validity, &mut ctx)?; + assert_eq!(replaced[0], views[0]); + assert_eq!(replaced[1], BinaryView::empty_view()); + + let replaced = VarBinViewData::replace_invalid_views(views, &Validity::AllInvalid, &mut ctx)?; + assert!( + replaced + .iter() + .all(|view| *view == BinaryView::empty_view()) + ); + Ok(()) +} + +#[test] +pub fn deserialize_null_views() -> VortexResult<()> { + let views = Buffer::::copy_from(vec![ + BinaryView::new_ref(14, *b"hell", 0, 0), + BinaryView::new_ref(13, *b"AAAA", 0xDEAD_BEEF, 0xF000_0000), + ]); + let buffers = Arc::new([ByteBuffer::from(b"hello world ololo".to_vec())]); + let dtype = DType::Utf8(Nullability::Nullable); + let buffer = BitBuffer::from_iter([true, false]); + let validity = Validity::from_bit_buffer(buffer, Nullability::Nullable); + let session = array_session(); + let array = VarBinViewArray::try_new( + views.clone(), + buffers, + dtype.clone(), + validity, + &mut session.create_execution_ctx(), + )?; + + let array_ctx = ArrayContext::empty(); + let serialized = + array + .clone() + .into_array() + .serialize(&array_ctx, &session, &SerializeOptions::default())?; + + let mut concat = ByteBufferMut::empty(); + for buf in serialized { + concat.extend_from_slice(buf.as_ref()); + } + let parts = SerializedArray::try_from(concat.freeze())?; + let decoded = parts.decode( + &dtype, + array.len(), + &ReadContext::new(array_ctx.to_ids()), + &session, + )?; + + let decoded = decoded + .as_opt::() + .ok_or_else(|| vortex_err!("expected VarBinView"))?; + assert_eq!(decoded.views()[0], views[0]); + assert_eq!(decoded.views()[1], BinaryView::empty_view()); + Ok(()) +} diff --git a/vortex-array/src/arrays/varbinview/vtable/mod.rs b/vortex-array/src/arrays/varbinview/vtable/mod.rs index 33e9ce289a5..0d2963cd6b2 100644 --- a/vortex-array/src/arrays/varbinview/vtable/mod.rs +++ b/vortex-array/src/arrays/varbinview/vtable/mod.rs @@ -19,6 +19,7 @@ use crate::ArrayRef; use crate::EqMode; use crate::ExecutionCtx; use crate::ExecutionResult; +use crate::VortexSessionExecute; use crate::array::Array; use crate::array::ArrayId; use crate::array::ArrayView; @@ -168,7 +169,7 @@ impl VTable for VarBinView { buffers: &[BufferHandle], children: &dyn ArrayChildren, - _session: &VortexSession, + session: &VortexSession, ) -> VortexResult> { if !metadata.is_empty() { vortex_bail!( @@ -224,6 +225,7 @@ impl VTable for VarBinView { Arc::from(data_buffers), dtype.clone(), validity.clone(), + &mut session.create_execution_ctx(), )?; let slots = VarBinViewData::make_slots(&validity, len); Ok(ArrayParts::new(self.clone(), dtype.clone(), len, data).with_slots(slots)) diff --git a/vortex-cuda/src/arrow/canonical.rs b/vortex-cuda/src/arrow/canonical.rs index 362fef04b55..71e1a88f468 100644 --- a/vortex-cuda/src/arrow/canonical.rs +++ b/vortex-cuda/src/arrow/canonical.rs @@ -1439,6 +1439,8 @@ mod tests { use rstest::rstest; use vortex::array::ArrayRef; use vortex::array::IntoArray; + use vortex::array::VortexSessionExecute; + use vortex::array::array_session; use vortex::array::arrays::BoolArray; use vortex::array::arrays::ChunkedArray; use vortex::array::arrays::DecimalArray; @@ -1501,7 +1503,7 @@ mod tests { } fn cuda_ctx_with_varbin_layout(layout: VarBinExportLayout) -> VortexResult { - let session = vortex::array::array_session() + let session = array_session() .with_some(CudaSession::try_default()?.with_varbin_export_layout(layout)); CudaSession::create_execution_ctx(&session) } @@ -1669,6 +1671,7 @@ mod tests { Arc::from([first, second]), dtype, Validity::NonNullable, + &mut array_session().create_execution_ctx(), ) .vortex_expect("valid multi-buffer VarBinViewArray") .into_array(); diff --git a/vortex-test/e2e-cuda/src/lib.rs b/vortex-test/e2e-cuda/src/lib.rs index 1c6530ce612..a18bdb89802 100644 --- a/vortex-test/e2e-cuda/src/lib.rs +++ b/vortex-test/e2e-cuda/src/lib.rs @@ -35,6 +35,7 @@ use futures::executor::block_on; use vortex::array::ArrayRef as VortexArrayRef; use vortex::array::IntoArray; use vortex::array::VortexSessionExecute; +use vortex::array::array_session; use vortex::array::arrays::BoolArray; use vortex::array::arrays::DecimalArray; use vortex::array::arrays::DictArray as VortexDictArray; @@ -69,7 +70,7 @@ use vortex_cuda::arrow::DeviceArrayStreamExt; const PRIMITIVE_DTYPE_ENV: &str = "VORTEX_CUDF_PRIMITIVE_DTYPE"; static SESSION: LazyLock = LazyLock::new(|| { - vortex::array::array_session() + array_session() .with::() .with::() .with::() @@ -201,6 +202,7 @@ fn multi_buffer_varbinview(dtype: DType) -> VortexArrayRef { Arc::from([first, second]), dtype, Validity::NonNullable, + &mut array_session().create_execution_ctx(), ) .expect("multi-buffer VarBinViewArray") .into_array()