diff --git a/encodings/alp/src/alp/ops.rs b/encodings/alp/src/alp/ops.rs index a8850744056..58b5ee99360 100644 --- a/encodings/alp/src/alp/ops.rs +++ b/encodings/alp/src/alp/ops.rs @@ -15,6 +15,8 @@ use crate::ALPFloat; use crate::match_each_alp_float_ptype; impl OperationsVTable for ALP { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, ALP>, index: usize, diff --git a/encodings/alp/src/alp_rd/ops.rs b/encodings/alp/src/alp_rd/ops.rs index edb2fb21186..1433966a01e 100644 --- a/encodings/alp/src/alp_rd/ops.rs +++ b/encodings/alp/src/alp_rd/ops.rs @@ -14,6 +14,8 @@ use crate::ALPRDArrayExt; use crate::ALPRDArraySlotsExt; impl OperationsVTable for ALPRD { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, ALPRD>, index: usize, diff --git a/encodings/bytebool/src/array.rs b/encodings/bytebool/src/array.rs index faa1fea81fa..02bb4fcb2b8 100644 --- a/encodings/bytebool/src/array.rs +++ b/encodings/bytebool/src/array.rs @@ -312,6 +312,8 @@ impl ValidityVTable for ByteBool { } impl OperationsVTable for ByteBool { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, ByteBool>, index: usize, diff --git a/encodings/datetime-parts/src/ops.rs b/encodings/datetime-parts/src/ops.rs index d99e55c7542..0ee4f9e315c 100644 --- a/encodings/datetime-parts/src/ops.rs +++ b/encodings/datetime-parts/src/ops.rs @@ -17,6 +17,8 @@ use crate::timestamp; use crate::timestamp::TimestampParts; impl OperationsVTable for DateTimeParts { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, DateTimeParts>, index: usize, diff --git a/encodings/decimal-byte-parts/src/decimal_byte_parts/mod.rs b/encodings/decimal-byte-parts/src/decimal_byte_parts/mod.rs index d5b0024f5b7..e1d081ddaca 100644 --- a/encodings/decimal-byte-parts/src/decimal_byte_parts/mod.rs +++ b/encodings/decimal-byte-parts/src/decimal_byte_parts/mod.rs @@ -289,6 +289,8 @@ fn to_canonical_decimal( } impl OperationsVTable for DecimalByteParts { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, DecimalByteParts>, index: usize, diff --git a/encodings/fastlanes/src/bitpacking/vtable/operations.rs b/encodings/fastlanes/src/bitpacking/vtable/operations.rs index e14b27323c1..2816407ac03 100644 --- a/encodings/fastlanes/src/bitpacking/vtable/operations.rs +++ b/encodings/fastlanes/src/bitpacking/vtable/operations.rs @@ -11,6 +11,8 @@ use crate::BitPacked; use crate::bitpack_decompress; use crate::bitpacking::array::BitPackedArrayExt; impl OperationsVTable for BitPacked { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, BitPacked>, index: usize, diff --git a/encodings/fastlanes/src/delta/vtable/operations.rs b/encodings/fastlanes/src/delta/vtable/operations.rs index 7ed57a0886d..b37760621ea 100644 --- a/encodings/fastlanes/src/delta/vtable/operations.rs +++ b/encodings/fastlanes/src/delta/vtable/operations.rs @@ -11,6 +11,8 @@ use vortex_error::VortexResult; use super::Delta; impl OperationsVTable for Delta { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Delta>, index: usize, diff --git a/encodings/fastlanes/src/for/vtable/operations.rs b/encodings/fastlanes/src/for/vtable/operations.rs index 36dac998cbe..361020824d5 100644 --- a/encodings/fastlanes/src/for/vtable/operations.rs +++ b/encodings/fastlanes/src/for/vtable/operations.rs @@ -13,6 +13,8 @@ use super::FoR; use crate::r#for::array::FoRArrayExt; use crate::r#for::array::FoRArraySlotsExt; impl OperationsVTable for FoR { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, FoR>, index: usize, diff --git a/encodings/fastlanes/src/rle/vtable/operations.rs b/encodings/fastlanes/src/rle/vtable/operations.rs index ca4d2d39545..ba6a624f1f5 100644 --- a/encodings/fastlanes/src/rle/vtable/operations.rs +++ b/encodings/fastlanes/src/rle/vtable/operations.rs @@ -14,6 +14,8 @@ use crate::rle::RLEArrayExt; use crate::rle::RLEArraySlotsExt; impl OperationsVTable for RLE { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, RLE>, index: usize, diff --git a/encodings/fastlanes/src/transposed_bool.rs b/encodings/fastlanes/src/transposed_bool.rs index efb3443cdfc..92da3b3b79c 100644 --- a/encodings/fastlanes/src/transposed_bool.rs +++ b/encodings/fastlanes/src/transposed_bool.rs @@ -256,6 +256,8 @@ impl VTable for TransposedBool { } impl OperationsVTable for TransposedBool { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, TransposedBool>, index: usize, diff --git a/encodings/fsst/src/ops.rs b/encodings/fsst/src/ops.rs index b630508ed9e..03560b2c8b8 100644 --- a/encodings/fsst/src/ops.rs +++ b/encodings/fsst/src/ops.rs @@ -14,6 +14,8 @@ use crate::FSST; use crate::FSSTArrayExt; impl OperationsVTable for FSST { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, FSST>, index: usize, diff --git a/encodings/onpair/src/ops.rs b/encodings/onpair/src/ops.rs index 728e5a0e6f2..082bc55e199 100644 --- a/encodings/onpair/src/ops.rs +++ b/encodings/onpair/src/ops.rs @@ -18,6 +18,8 @@ use crate::decode::code_boundary_at; use crate::decode::collect_widened; impl OperationsVTable for OnPair { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, OnPair>, index: usize, diff --git a/encodings/parquet-variant/src/operations.rs b/encodings/parquet-variant/src/operations.rs index 8517df27f22..eab80bf63e1 100644 --- a/encodings/parquet-variant/src/operations.rs +++ b/encodings/parquet-variant/src/operations.rs @@ -31,6 +31,8 @@ use crate::ParquetVariantArraySlotsExt; use crate::vtable::ParquetVariant; impl OperationsVTable for ParquetVariant { + type ProbeState = (); + /// Resolves one row according to the Parquet Variant shredding rules. /// /// For valid data, a row with both `value` and struct `typed_value` is a partially diff --git a/encodings/pco/src/array.rs b/encodings/pco/src/array.rs index 44a6a8e4045..97a17fe6a3c 100644 --- a/encodings/pco/src/array.rs +++ b/encodings/pco/src/array.rs @@ -778,6 +778,8 @@ impl ValidityVTable for Pco { } impl OperationsVTable for Pco { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Pco>, index: usize, diff --git a/encodings/runend/src/ops.rs b/encodings/runend/src/ops.rs index e2c2e3fc99b..46949d84e1b 100644 --- a/encodings/runend/src/ops.rs +++ b/encodings/runend/src/ops.rs @@ -18,6 +18,8 @@ use crate::array::RunEndArrayExt; use crate::array::RunEndArraySlotsExt; impl OperationsVTable for RunEnd { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, RunEnd>, index: usize, diff --git a/encodings/sequence/src/array.rs b/encodings/sequence/src/array.rs index a36a1f6edab..f3f44e64f6e 100644 --- a/encodings/sequence/src/array.rs +++ b/encodings/sequence/src/array.rs @@ -425,6 +425,8 @@ impl VTable for Sequence { } impl OperationsVTable for Sequence { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Sequence>, index: usize, diff --git a/encodings/sparse/src/ops.rs b/encodings/sparse/src/ops.rs index 568d8d377d1..acabbb103d5 100644 --- a/encodings/sparse/src/ops.rs +++ b/encodings/sparse/src/ops.rs @@ -11,6 +11,8 @@ use crate::Sparse; use crate::SparseExt as _; impl OperationsVTable for Sparse { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Sparse>, index: usize, diff --git a/encodings/zigzag/src/array.rs b/encodings/zigzag/src/array.rs index ef3165132f4..09076fcbbdf 100644 --- a/encodings/zigzag/src/array.rs +++ b/encodings/zigzag/src/array.rs @@ -231,6 +231,8 @@ impl Default for ZigZagData { } impl OperationsVTable for ZigZag { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, ZigZag>, index: usize, diff --git a/encodings/zstd/src/array.rs b/encodings/zstd/src/array.rs index fb3551e539b..dafe8e820a9 100644 --- a/encodings/zstd/src/array.rs +++ b/encodings/zstd/src/array.rs @@ -1588,6 +1588,8 @@ impl ValidityVTable for Zstd { } impl OperationsVTable for Zstd { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Zstd>, index: usize, diff --git a/encodings/zstd/src/zstd_buffers.rs b/encodings/zstd/src/zstd_buffers.rs index f21deee6f64..51d338777ca 100644 --- a/encodings/zstd/src/zstd_buffers.rs +++ b/encodings/zstd/src/zstd_buffers.rs @@ -520,6 +520,8 @@ impl VTable for ZstdBuffers { } impl OperationsVTable for ZstdBuffers { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, ZstdBuffers>, index: usize, diff --git a/vortex-array/src/array/erased.rs b/vortex-array/src/array/erased.rs index 76c0b33cd18..f0e6b569cb6 100644 --- a/vortex-array/src/array/erased.rs +++ b/vortex-array/src/array/erased.rs @@ -34,6 +34,8 @@ use crate::array::ArrayId; use crate::array::ArrayInner; use crate::array::ArraySlots; use crate::array::DynArrayData; +use crate::array::probe::ArrayProbe; +use crate::array::probe::RepeatedArrayProbe; use crate::arrays::Constant; use crate::arrays::DictArray; use crate::arrays::FilterArray; @@ -273,31 +275,61 @@ impl ArrayRef { } /// Execute the array to extract a scalar at the given index. + /// + /// A one-off read; the same as `self.probe().execute_scalar(index, ctx)`. + // TODO(joe): deprecate this in favour of `probe()`. + #[inline] pub fn execute_scalar(&self, index: usize, ctx: &mut ExecutionCtx) -> VortexResult { - vortex_ensure!(index < self.len(), OutOfBounds: index, 0, self.len()); - if self.dtype().is_nullable() && self.is_invalid(index, ctx)? { - return Ok(Scalar::null(self.dtype().clone())); - } - let scalar = self.0.data.execute_scalar(self, index, ctx)?; - debug_assert_eq!(self.dtype(), scalar.dtype(), "Scalar dtype mismatch"); - Ok(scalar) + self.probe().execute_scalar(index, ctx) + } + + /// A one-off row accessor over this array. It borrows the handle and retains nothing; for + /// many reads of the same array use [`Self::repeated_probe`]. + /// + /// ``` + /// use vortex_array::{IntoArray, VortexSessionExecute}; + /// use vortex_array::arrays::PrimitiveArray; + /// + /// let array = PrimitiveArray::from_iter([10i32, 20, 30]).into_array(); + /// let mut ctx = vortex_array::array_session().create_execution_ctx(); + /// assert_eq!(array.probe().execute_scalar(2, &mut ctx)?, 30i32.into()); + /// # Ok::<(), vortex_error::VortexError>(()) + /// ``` + #[inline] + pub fn probe(&self) -> ArrayProbe<'_> { + ArrayProbe::Once(self) + } + + /// A row accessor that owns a handle to this array and keeps encoding state, its validity + /// probe and child probes between reads. + /// + /// ``` + /// use vortex_array::{IntoArray, VortexSessionExecute}; + /// use vortex_array::arrays::PrimitiveArray; + /// + /// let array = PrimitiveArray::from_iter([10i32, 20, 30]).into_array(); + /// let mut ctx = vortex_array::array_session().create_execution_ctx(); + /// let mut probe = array.repeated_probe(); + /// assert_eq!(probe.execute_scalar(2, &mut ctx)?, 30i32.into()); + /// assert_eq!(probe.execute_scalar(0, &mut ctx)?, 10i32.into()); + /// # Ok::<(), vortex_error::VortexError>(()) + /// ``` + pub fn repeated_probe(&self) -> RepeatedArrayProbe { + RepeatedArrayProbe::new(self.clone()) } /// Returns whether the item at `index` is valid. + /// + /// A one-off read; the same as `self.probe().execute_is_valid(index, ctx)`. + // TODO(joe): deprecate this in favour of `probe/repeated_probe()`. + #[inline] pub fn is_valid(&self, index: usize, ctx: &mut ExecutionCtx) -> VortexResult { - vortex_ensure!(index < self.len(), OutOfBounds: index, 0, self.len()); - match self.validity()? { - Validity::NonNullable | Validity::AllValid => Ok(true), - Validity::AllInvalid => Ok(false), - Validity::Array(a) => a - .execute_scalar(index, ctx)? - .as_bool() - .value() - .ok_or_else(|| vortex_err!("validity value at index {} is null", index)), - } + self.probe().execute_is_valid(index, ctx) } /// Returns whether the item at `index` is invalid. + // TODO(joe): deprecate this. + #[inline] pub fn is_invalid(&self, index: usize, ctx: &mut ExecutionCtx) -> VortexResult { Ok(!self.is_valid(index, ctx)?) } diff --git a/vortex-array/src/array/mod.rs b/vortex-array/src/array/mod.rs index 8b9b2812bfa..c7f709e6f56 100644 --- a/vortex-array/src/array/mod.rs +++ b/vortex-array/src/array/mod.rs @@ -30,6 +30,8 @@ pub use erased::*; mod plugin; pub use plugin::*; +mod probe; +pub use probe::*; mod foreign; pub(crate) use foreign::*; @@ -220,15 +222,26 @@ pub(crate) trait DynArrayData: 'static + private::Sealed + Send + Sync + Debug { ctx: &mut ExecutionCtx, ) -> VortexResult; - /// Execute the scalar at the given index. + /// Read a non-null scalar at `index` for a one-off access; nothing is retained. /// - /// This method panics if the index is out of bounds for the array. - fn execute_scalar( + /// Kept apart from [`Self::probe_scalar_retained`] so this entry is a bare trampoline into + /// the encoding: sharing one function made the one-off path pay the retained branch's + /// register saves before its tail call. + fn probe_scalar_once( &self, this: &ArrayRef, index: usize, ctx: &mut ExecutionCtx, ) -> VortexResult; + + /// Read a non-null scalar at `index`, keeping preparation in `state` for later reads. + fn probe_scalar_retained( + &self, + this: &ArrayRef, + index: usize, + state: &mut Option>, + ctx: &mut ExecutionCtx, + ) -> VortexResult; } /// Trait for converting a type into a Vortex [`ArrayRef`]. @@ -490,14 +503,32 @@ impl DynArrayData for ArrayData { V::execute(typed, ctx) } - fn execute_scalar( + fn probe_scalar_once( + &self, + this: &ArrayRef, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + // SAFETY: this adapter belongs to the ArrayData stored in `this`. + let view = unsafe { ArrayView::new_unchecked(this, &self.data) }; + >::probe_scalar( + &mut ProbeState::once(view), + index, + ctx, + ) + } + + fn probe_scalar_retained( &self, this: &ArrayRef, index: usize, + state: &mut Option>, ctx: &mut ExecutionCtx, ) -> VortexResult { + // SAFETY: this adapter belongs to the ArrayData stored in `this`. let view = unsafe { ArrayView::new_unchecked(this, &self.data) }; - >::scalar_at(view, index, ctx) + let mut state = ProbeState::repeated(view, repeated_state(state)?); + >::probe_scalar(&mut state, index, ctx) } } diff --git a/vortex-array/src/array/probe/array.rs b/vortex-array/src/array/probe/array.rs new file mode 100644 index 00000000000..d53afdbe3e2 --- /dev/null +++ b/vortex-array/src/array/probe/array.rs @@ -0,0 +1,333 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use std::mem::needs_drop; + +use vortex_error::VortexError; +use vortex_error::VortexResult; +use vortex_error::vortex_err; + +use crate::ArrayRef; +use crate::ExecutionCtx; +use crate::array::ArrayView; +use crate::array::VTable; +use crate::array::probe::RepeatedArrayProbe; +use crate::array::probe::RepeatedState; +use crate::array::probe::repeated::child_probe; +use crate::arrays::Primitive; +use crate::scalar::Scalar; +use crate::vtable::OperationsVTable; + +/// A borrowed row accessor. +/// +/// Either a one-off reader over a borrowed array, which retains nothing, or a borrow of a +/// [`RepeatedArrayProbe`] that keeps state between reads. Both variants hold only references, +/// so an `ArrayProbe` has no destructor and building one per read, including for every child +/// slot of a nested array, costs nothing. +pub enum ArrayProbe<'a> { + /// A single read of a borrowed array; nothing outlives the call. + Once(&'a ArrayRef), + /// A read through a retained probe. + Repeated(&'a mut RepeatedArrayProbe), +} + +// Everything built per read must be free to drop, or the unwind path pins it in memory and +// blocks the tail call into the encoding. Keep these at compile time. +const _: () = assert!(!needs_drop::>()); +const _: () = assert!(!needs_drop::>()); + +impl ArrayProbe<'_> { + /// The array this probe reads from. + #[inline] + pub fn array(&self) -> &ArrayRef { + match self { + Self::Once(array) => array, + Self::Repeated(probe) => probe.array(), + } + } + + /// Read the scalar at `index`, including its nullness. + #[inline] + pub fn execute_scalar(&mut self, index: usize, ctx: &mut ExecutionCtx) -> VortexResult { + match self { + Self::Once(array) => execute_scalar_once(array, index, ctx), + Self::Repeated(probe) => probe.execute_scalar(index, ctx), + } + } + + /// Whether the row at `index` is valid. + #[inline] + pub fn execute_is_valid(&mut self, index: usize, ctx: &mut ExecutionCtx) -> VortexResult { + match self { + Self::Once(array) => execute_is_valid_once(array, index, ctx), + Self::Repeated(probe) => probe.execute_is_valid(index, ctx), + } + } + + /// Whether the row at `index` is null. + #[inline] + pub fn execute_is_invalid( + &mut self, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + Ok(!self.execute_is_valid(index, ctx)?) + } +} + +/// One-off scalar read: nothing outlives the call. +#[inline] +fn execute_scalar_once( + array: &ArrayRef, + index: usize, + ctx: &mut ExecutionCtx, +) -> VortexResult { + if !execute_is_valid_once(array, index, ctx)? { + return Ok(Scalar::null(array.dtype().clone())); + } + check_dtype( + array, + array.dyn_array().probe_scalar_once(array, index, ctx), + ) +} + +/// One-off validity read: nothing outlives the call. +#[inline] +fn execute_is_valid_once( + array: &ArrayRef, + index: usize, + ctx: &mut ExecutionCtx, +) -> VortexResult { + check_bounds(array, index)?; + if !array.dtype().is_nullable() { + return Ok(true); + } + // Matching the validity directly keeps this path free of any retained temporary. + array.validity()?.execute_is_valid(index, ctx) +} + +/// Pass an encoding's result through, checking its dtype in debug builds. +/// +/// Unwrapping and re-wrapping the result here would cost an extra copy of the scalar on every +/// read. +#[inline] +pub(super) fn check_dtype(array: &ArrayRef, result: VortexResult) -> VortexResult { + result.inspect(|scalar| { + debug_assert_eq!(scalar.dtype(), array.dtype(), "Scalar dtype mismatch"); + }) +} + +#[inline] +pub(super) fn check_bounds(array: &ArrayRef, index: usize) -> VortexResult<()> { + if index >= array.len() { + return Err(out_of_bounds(index, array.len())); + } + Ok(()) +} + +/// Kept out of line so the error path, which captures a backtrace, does not count against the +/// hot probe bodies when the inliner sizes them. +#[cold] +#[inline(never)] +fn out_of_bounds(index: usize, len: usize) -> VortexError { + vortex_err!(OutOfBounds: index, 0, len) +} + +/// The encoding state type of `V`'s operations vtable. +pub type EncodingProbeState = + <::OperationsVTable as OperationsVTable>::ProbeState; + +/// Everything an encoding's `probe_scalar` runs with: the typed view of the array being read +/// and, for a repeated read, a borrow of the state its [`RepeatedArrayProbe`] keeps. +/// +/// Passed to [`OperationsVTable::probe_scalar`]. +/// A one-off read gets [`ProbeState::once`], which holds only the view. A repeated read borrows +/// the encoding's own state and the lazily created child probes. Neither owns anything, so +/// building one per read is free. +/// +/// Encodings never inspect the policy: [`ProbeState::slot`] hands out a probe over a child under +/// whichever policy is in force, [`ProbeState::retained`] hands out the encoding state only when +/// it is kept, and [`ProbeState::split`] gives both at once. +pub struct ProbeState<'a, V: VTable> { + array: ArrayView<'a, V>, + retained: Option<&'a mut RepeatedState>>, +} + +impl<'a, V: VTable> ProbeState<'a, V> { + /// State for a single read of `array`. + /// + /// Encodings use it to run their `probe_scalar` path from the deprecated `scalar_at`; it + /// goes away with `scalar_at`. + #[inline] + pub fn once(array: ArrayView<'a, V>) -> Self { + Self { + array, + retained: None, + } + } + + /// State for a read through a [`RepeatedArrayProbe`], borrowing what it keeps. + #[inline] + pub(crate) fn repeated( + array: ArrayView<'a, V>, + retained: &'a mut RepeatedState>, + ) -> Self { + Self { + array, + retained: Some(retained), + } + } + + /// The typed view of the array being read. + #[inline] + pub fn array(&self) -> ArrayView<'a, V> { + self.array + } + + /// The encoding's retained state, or `None` for a one-off read. + /// + /// Encodings whose repeated algorithm differs from their one-off one branch on this. To hold + /// the state while reading children, use [`Self::split`]. + #[inline] + pub fn retained(&mut self) -> Option<&mut EncodingProbeState> { + self.retained.as_deref_mut().map(RepeatedState::state_mut) + } + + /// The encoding's retained state and the array's children, borrowed apart, so a cache can + /// stay borrowed while children are read. + #[inline] + pub fn split(&mut self) -> (Option<&mut EncodingProbeState>, ProbeChildren<'_, 'a, V>) { + let array = self.array; + match &mut self.retained { + None => (None, ProbeChildren { array, slots: None }), + Some(repeated) => { + let (state, slots) = repeated.split_mut(); + let children = ProbeChildren { + array, + slots: Some(slots), + }; + (Some(state), children) + } + } + } + + /// A probe over the array's child in `slot`, under this state's policy. + /// + /// Shorthand for [`Self::split`] followed by [`ProbeChildren::slot`]. `None` if the slot is + /// absent, an error if it is out of bounds. + #[inline] + pub fn slot(&mut self, slot: usize) -> VortexResult>> { + let parent = self.array.array(); + match &mut self.retained { + None => Ok(child_of(parent, slot)?.map(ArrayProbe::Once)), + Some(repeated) => { + Ok(child_probe(repeated.split_mut().1, parent, slot)?.map(ArrayProbe::Repeated)) + } + } + } +} + +/// The child slots of the array a [`ProbeState`] reads, borrowed apart from the encoding state +/// by [`ProbeState::split`]. +pub struct ProbeChildren<'s, 'a, V: VTable> { + array: ArrayView<'a, V>, + /// The retained slot table, or `None` for a one-off read. + slots: Option<&'s mut Vec>>, +} + +impl ProbeChildren<'_, '_, V> { + /// A probe over the child in `slot`, under the read's policy. + /// + /// For a one-off read this borrows the child; for a repeated read it is the + /// [`RepeatedArrayProbe`] kept in the slot table, created on first use. `None` if the slot + /// is absent, an error if it is out of bounds. + #[inline] + pub fn slot(&mut self, slot: usize) -> VortexResult>> { + let parent = self.array.array(); + match &mut self.slots { + None => Ok(child_of(parent, slot)?.map(ArrayProbe::Once)), + Some(slots) => Ok(child_probe(slots, parent, slot)?.map(ArrayProbe::Repeated)), + } + } +} + +/// The child of `parent` in `slot`: `None` if the slot is absent, an error if it is out of bounds. +#[inline] +pub(super) fn child_of(parent: &ArrayRef, slot: usize) -> VortexResult> { + Ok(parent + .slots() + .get(slot) + .ok_or_else(|| vortex_err!("Probe slot {slot} is out of bounds"))? + .as_ref()) +} +#[cfg(test)] +mod tests { + use vortex_error::VortexResult; + use vortex_error::vortex_err; + + use super::*; + use crate::VortexSessionExecute; + use crate::array::IntoArray; + use crate::arrays::PrimitiveArray; + use crate::arrays::Struct; + use crate::arrays::StructArray; + + fn nullable_ints() -> ArrayRef { + PrimitiveArray::from_option_iter([Some(10i32), None, Some(30)]).into_array() + } + + fn check_reads(mut probe: ArrayProbe<'_>, ctx: &mut ExecutionCtx) -> VortexResult<()> { + assert!(probe.execute_scalar(3, ctx).is_err()); + assert!(probe.execute_is_valid(3, ctx).is_err()); + assert!(probe.execute_scalar(1, ctx)?.is_null()); + assert!(!probe.execute_is_valid(1, ctx)?); + assert!(probe.execute_is_invalid(1, ctx)?); + assert_eq!(probe.execute_scalar(2, ctx)?, Scalar::from(Some(30i32))); + assert_eq!(probe.execute_scalar(0, ctx)?, Scalar::from(Some(10i32))); + Ok(()) + } + + #[test] + fn once_probe_checks_bounds_and_nulls() -> VortexResult<()> { + let mut ctx = crate::array_session().create_execution_ctx(); + let array = nullable_ints(); + check_reads(array.probe(), &mut ctx) + } + + #[test] + fn once_state_reads_children_without_retaining() -> VortexResult<()> { + let mut ctx = crate::array_session().create_execution_ctx(); + let array = struct_of_two_fields()?; + let typed = array + .as_opt::() + .ok_or_else(|| vortex_err!("expected a struct"))?; + let mut state = ProbeState::once(typed); + + assert!(state.slot(5).is_err()); + // Slot 0 is the struct's absent validity. + assert!(state.slot(0)?.is_none()); + { + let mut field = state.slot(2)?.ok_or_else(|| vortex_err!("missing field"))?; + assert!(matches!(field, ArrayProbe::Once(_))); + assert_eq!(field.array().len(), 2); + assert_eq!(field.execute_scalar(1, &mut ctx)?, Scalar::from(4i64)); + assert!(field.execute_is_valid(0, &mut ctx)?); + } + assert!(state.retained().is_none()); + let (retained, mut children) = state.split(); + assert!(retained.is_none()); + let mut field = children + .slot(1)? + .ok_or_else(|| vortex_err!("missing field"))?; + assert_eq!(field.execute_scalar(1, &mut ctx)?, Scalar::from(2i32)); + Ok(()) + } + + fn struct_of_two_fields() -> VortexResult { + Ok(StructArray::from_fields(&[ + ("a", PrimitiveArray::from_iter([1i32, 2]).into_array()), + ("b", PrimitiveArray::from_iter([3i64, 4]).into_array()), + ])? + .into_array()) + } +} diff --git a/vortex-array/src/array/probe/mod.rs b/vortex-array/src/array/probe/mod.rs new file mode 100644 index 00000000000..7b1e1e4a481 --- /dev/null +++ b/vortex-array/src/array/probe/mod.rs @@ -0,0 +1,15 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Row access over arrays. +//! +//! [`ArrayProbe`] is a borrowed reader: a one-off read that retains nothing, or a borrow of a +//! [`RepeatedArrayProbe`], which owns its array and keeps encoding state, validity and child +//! probes between reads. Encodings implement a single +//! [`probe_scalar`](crate::vtable::OperationsVTable::probe_scalar) that serves both through +//! [`ProbeState`]. + +mod array; +pub use array::*; +mod repeated; +pub use repeated::*; diff --git a/vortex-array/src/array/probe/repeated.rs b/vortex-array/src/array/probe/repeated.rs new file mode 100644 index 00000000000..3a6533b938e --- /dev/null +++ b/vortex-array/src/array/probe/repeated.rs @@ -0,0 +1,266 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use std::any::Any; + +use vortex_error::VortexResult; +use vortex_error::vortex_err; + +use crate::ArrayRef; +use crate::ExecutionCtx; +use crate::array::probe::ArrayProbe; +use crate::array::probe::array::check_bounds; +use crate::array::probe::array::check_dtype; +use crate::array::probe::array::child_of; +use crate::scalar::Scalar; +use crate::validity::Validity; + +/// A row accessor that owns its array and keeps preparation between reads. +/// +/// The encoding's state, the validity probe and the probes over child slots are created on +/// first use and reused by every following read. Dropping the probe drops all of them. Array +/// handles share their buffers, so the probe can outlive the handle it was built from. Probes +/// are local to a thread. +pub struct RepeatedArrayProbe { + array: ArrayRef, + /// The encoding's [`RepeatedState`], type-erased, created on the first read. + state: Option>, + /// Resolved on the first read of a nullable array: `Some(valid)` for uniform validity, + /// otherwise `validity` holds a probe over the validity array, boxed because it is a probe. + uniform_validity: Option, + validity: Option>, +} + +impl RepeatedArrayProbe { + /// Own an array for repeated reads. + pub fn new(array: ArrayRef) -> Self { + Self { + array, + state: None, + uniform_validity: None, + validity: None, + } + } + + /// The array this probe reads from. + pub fn array(&self) -> &ArrayRef { + &self.array + } + + /// Borrow this probe as an [`ArrayProbe`]. + #[inline] + pub fn as_probe(&mut self) -> ArrayProbe<'_> { + ArrayProbe::Repeated(self) + } + + /// Read the scalar at `index`, including its nullness, reusing retained preparation. + pub fn execute_scalar(&mut self, index: usize, ctx: &mut ExecutionCtx) -> VortexResult { + if !self.execute_is_valid(index, ctx)? { + return Ok(Scalar::null(self.array.dtype().clone())); + } + let result = + self.array + .dyn_array() + .probe_scalar_retained(&self.array, index, &mut self.state, ctx); + check_dtype(&self.array, result) + } + + /// Whether the row at `index` is valid, through the retained validity. + pub fn execute_is_valid(&mut self, index: usize, ctx: &mut ExecutionCtx) -> VortexResult { + check_bounds(&self.array, index)?; + if !self.array.dtype().is_nullable() { + return Ok(true); + } + if let Some(valid) = self.uniform_validity { + return Ok(valid); + } + if self.validity.is_none() { + match self.array.validity()? { + Validity::NonNullable | Validity::AllValid => { + self.uniform_validity = Some(true); + return Ok(true); + } + Validity::AllInvalid => { + self.uniform_validity = Some(false); + return Ok(false); + } + Validity::Array(array) => { + self.validity = Some(Box::new(RepeatedArrayProbe::new(array))); + } + } + } + self.validity + .as_mut() + .ok_or_else(|| vortex_err!("validity probe was just initialized"))? + .execute_scalar(index, ctx)? + .as_bool() + .value() + .ok_or_else(|| vortex_err!("validity value at index {index} is null")) + } + + /// Whether the row at `index` is null. + pub fn execute_is_invalid( + &mut self, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + Ok(!self.execute_is_valid(index, ctx)?) + } +} + +/// Get or create the [`RepeatedState`] for encoding state `S` in a probe's erased slot. +pub(crate) fn repeated_state( + slot: &mut Option>, +) -> VortexResult<&mut RepeatedState> { + slot.get_or_insert_with(|| Box::new(RepeatedState::::default())) + .downcast_mut::>() + .ok_or_else(|| vortex_err!("Probe state type mismatch")) +} + +/// What a [`RepeatedArrayProbe`] keeps for its encoding between reads. +pub(crate) struct RepeatedState { + state: S, + /// Probes over the source's child slots, allocated on first use. Slots that are never + /// requested stay empty. + slots: Vec>, +} + +impl Default for RepeatedState { + fn default() -> Self { + Self { + state: S::default(), + slots: Vec::new(), + } + } +} + +impl RepeatedState { + #[inline] + pub(super) fn state_mut(&mut self) -> &mut S { + &mut self.state + } + + /// The encoding state and the child slot table as disjoint borrows. + #[inline] + pub(super) fn split_mut(&mut self) -> (&mut S, &mut Vec>) { + (&mut self.state, &mut self.slots) + } +} + +/// Get or create the retained probe over `parent`'s child in `slot`. +pub(super) fn child_probe<'s>( + slots: &'s mut Vec>, + parent: &ArrayRef, + slot: usize, +) -> VortexResult> { + let Some(child) = child_of(parent, slot)? else { + return Ok(None); + }; + if slots.is_empty() { + slots.resize_with(parent.slots().len(), || None); + } + Ok(Some(slots[slot].get_or_insert_with(|| { + RepeatedArrayProbe::new(child.clone()) + }))) +} + +#[cfg(test)] +mod tests { + use vortex_error::VortexResult; + use vortex_error::vortex_err; + + use super::*; + use crate::VortexSessionExecute; + use crate::array::IntoArray; + use crate::array::probe::ProbeState; + use crate::arrays::PrimitiveArray; + use crate::arrays::Struct; + use crate::arrays::StructArray; + + fn nullable_ints() -> ArrayRef { + PrimitiveArray::from_option_iter([Some(10i32), None, Some(30)]).into_array() + } + + fn check_reads(mut probe: ArrayProbe<'_>, ctx: &mut ExecutionCtx) -> VortexResult<()> { + assert!(probe.execute_scalar(3, ctx).is_err()); + assert!(probe.execute_is_valid(3, ctx).is_err()); + assert!(probe.execute_scalar(1, ctx)?.is_null()); + assert!(!probe.execute_is_valid(1, ctx)?); + assert!(probe.execute_is_invalid(1, ctx)?); + assert_eq!(probe.execute_scalar(2, ctx)?, Scalar::from(Some(30i32))); + assert_eq!(probe.execute_scalar(0, ctx)?, Scalar::from(Some(10i32))); + Ok(()) + } + + #[test] + fn repeated_probe_initializes_lazily_and_outlives_handle() -> VortexResult<()> { + let mut ctx = crate::array_session().create_execution_ctx(); + let array = nullable_ints(); + let mut probe = RepeatedArrayProbe::new(array.clone()); + drop(array); + assert!(probe.state.is_none()); + assert!(probe.validity.is_none()); + + check_reads(probe.as_probe(), &mut ctx)?; + + assert!(probe.state.is_some()); + assert!(probe.validity.is_some()); + assert!(probe.uniform_validity.is_none()); + Ok(()) + } + + #[test] + fn repeated_probe_resolves_uniform_validity_once() -> VortexResult<()> { + let mut ctx = crate::array_session().create_execution_ctx(); + let array = PrimitiveArray::from_option_iter([Some(1i32), Some(2)]).into_array(); + let mut probe = array.repeated_probe(); + assert!(probe.execute_is_valid(1, &mut ctx)?); + assert_eq!(probe.uniform_validity, Some(true)); + assert!(probe.validity.is_none()); + Ok(()) + } + + #[test] + fn repeated_state_creates_children_on_demand() -> VortexResult<()> { + let mut ctx = crate::array_session().create_execution_ctx(); + let array = struct_of_two_fields()?; + let typed = array + .as_opt::() + .ok_or_else(|| vortex_err!("expected a struct"))?; + let mut repeated = RepeatedState::<()>::default(); + { + let mut state = ProbeState::repeated(typed, &mut repeated); + assert!(state.slot(5).is_err()); + assert!(state.slot(0)?.is_none()); + assert!(state.retained().is_some()); + let mut field = state.slot(2)?.ok_or_else(|| vortex_err!("missing field"))?; + assert!(matches!(field, ArrayProbe::Repeated(_))); + assert_eq!(field.execute_scalar(1, &mut ctx)?, Scalar::from(4i64)); + assert_eq!(field.execute_scalar(0, &mut ctx)?, Scalar::from(3i64)); + } + + assert_eq!(repeated.slots.len(), array.slots().len()); + assert!(repeated.slots[1].is_none()); + let child = repeated.slots[2] + .as_ref() + .ok_or_else(|| vortex_err!("missing child probe"))?; + assert_eq!(child.array().len(), 2); + Ok(()) + } + + fn struct_of_two_fields() -> VortexResult { + Ok(StructArray::from_fields(&[ + ("a", PrimitiveArray::from_iter([1i32, 2]).into_array()), + ("b", PrimitiveArray::from_iter([3i64, 4]).into_array()), + ])? + .into_array()) + } + + #[test] + fn state_slot_rejects_mismatched_type() -> VortexResult<()> { + let mut slot: Option> = None; + repeated_state::<()>(&mut slot)?; + assert!(repeated_state::(&mut slot).is_err()); + Ok(()) + } +} diff --git a/vortex-array/src/array/vtable/operations.rs b/vortex-array/src/array/vtable/operations.rs index 7f49e683640..471ea6e3f88 100644 --- a/vortex-array/src/array/vtable/operations.rs +++ b/vortex-array/src/array/vtable/operations.rs @@ -7,6 +7,7 @@ use vortex_error::vortex_bail; use crate::ExecutionCtx; use crate::array::ArrayView; use crate::array::VTable; +use crate::array::probe::ProbeState; use crate::scalar::Scalar; use crate::vtable::NotSupported; @@ -17,6 +18,33 @@ use crate::vtable::NotSupported; /// [`ArrayRef`](crate::ArrayRef) /// methods perform common checks before dispatching here. pub trait OperationsVTable { + /// Encoding-specific state retained by repeated scalar access. + /// + /// Built once per repeated probe and never for a one-off read. Preparation belongs in + /// [`Self::probe_scalar`]. State owns its preparation and may hold shared buffer or array + /// handles. Use `()` when no state is needed. + type ProbeState: Default + 'static; + + /// Read the non-null scalar at `index` of the array in `state`. + /// + /// Bounds and validity have been checked; the row is non-null. `state` carries the typed + /// view of the array and, for a read through a + /// [`RepeatedArrayProbe`](crate::RepeatedArrayProbe), the state that probe keeps. Read + /// children through [`ProbeState::slot`], which follows the read's policy without the + /// encoding having to know it. Take encoding state from [`ProbeState::retained`]. The scalar must retain the source's + /// logical dtype, including nullability. + /// + /// The default preserves the existing scalar path without adding caching. + fn probe_scalar( + state: &mut ProbeState<'_, V>, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + // FIXME: Remove this default once all encodings have migrated to probe_scalar. + Self::scalar_at(state.array(), index, ctx) + } + + // FIXME: Deprecate scalar_at once encodings have migrated to probe_scalar. /// Fetch the scalar at the given index. /// /// ## Preconditions @@ -35,6 +63,8 @@ pub trait OperationsVTable { } impl OperationsVTable for NotSupported { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, V>, _index: usize, diff --git a/vortex-array/src/arrays/bool/vtable/operations.rs b/vortex-array/src/arrays/bool/vtable/operations.rs index c29ab20331b..e1ec02ecbc8 100644 --- a/vortex-array/src/arrays/bool/vtable/operations.rs +++ b/vortex-array/src/arrays/bool/vtable/operations.rs @@ -11,6 +11,8 @@ use crate::arrays::bool::BoolArrayExt; use crate::scalar::Scalar; impl OperationsVTable for Bool { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Bool>, index: usize, diff --git a/vortex-array/src/arrays/chunked/vtable/operations.rs b/vortex-array/src/arrays/chunked/vtable/operations.rs index 8f9e0867a88..395e8aff70a 100644 --- a/vortex-array/src/arrays/chunked/vtable/operations.rs +++ b/vortex-array/src/arrays/chunked/vtable/operations.rs @@ -11,6 +11,8 @@ use crate::arrays::chunked::ChunkedArrayExt; use crate::scalar::Scalar; impl OperationsVTable for Chunked { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Chunked>, index: usize, diff --git a/vortex-array/src/arrays/constant/vtable/operations.rs b/vortex-array/src/arrays/constant/vtable/operations.rs index e3568a9c39f..9c49094f324 100644 --- a/vortex-array/src/arrays/constant/vtable/operations.rs +++ b/vortex-array/src/arrays/constant/vtable/operations.rs @@ -10,6 +10,8 @@ use crate::arrays::Constant; use crate::scalar::Scalar; impl OperationsVTable for Constant { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Constant>, _index: usize, diff --git a/vortex-array/src/arrays/decimal/vtable/operations.rs b/vortex-array/src/arrays/decimal/vtable/operations.rs index 257f26127ae..aecf6638bd2 100644 --- a/vortex-array/src/arrays/decimal/vtable/operations.rs +++ b/vortex-array/src/arrays/decimal/vtable/operations.rs @@ -12,6 +12,8 @@ use crate::scalar::DecimalValue; use crate::scalar::Scalar; impl OperationsVTable for Decimal { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Decimal>, index: usize, diff --git a/vortex-array/src/arrays/dict/vtable/operations.rs b/vortex-array/src/arrays/dict/vtable/operations.rs index 1982a1e0870..d497db0f2f8 100644 --- a/vortex-array/src/arrays/dict/vtable/operations.rs +++ b/vortex-array/src/arrays/dict/vtable/operations.rs @@ -12,6 +12,8 @@ use crate::arrays::dict::DictArraySlotsExt; use crate::scalar::Scalar; impl OperationsVTable for Dict { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Dict>, index: usize, diff --git a/vortex-array/src/arrays/extension/vtable/operations.rs b/vortex-array/src/arrays/extension/vtable/operations.rs index 66de94b596a..519ef6088f7 100644 --- a/vortex-array/src/arrays/extension/vtable/operations.rs +++ b/vortex-array/src/arrays/extension/vtable/operations.rs @@ -11,6 +11,8 @@ use crate::arrays::extension::ExtensionArrayExt; use crate::scalar::Scalar; impl OperationsVTable for Extension { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Extension>, index: usize, diff --git a/vortex-array/src/arrays/filter/vtable.rs b/vortex-array/src/arrays/filter/vtable.rs index 260abe0046c..f5a61ba79d7 100644 --- a/vortex-array/src/arrays/filter/vtable.rs +++ b/vortex-array/src/arrays/filter/vtable.rs @@ -200,6 +200,8 @@ impl VTable for Filter { } } impl OperationsVTable for Filter { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Filter>, index: usize, diff --git a/vortex-array/src/arrays/fixed_size_list/vtable/operations.rs b/vortex-array/src/arrays/fixed_size_list/vtable/operations.rs index 9f4cf02fbf8..eca0e45ad95 100644 --- a/vortex-array/src/arrays/fixed_size_list/vtable/operations.rs +++ b/vortex-array/src/arrays/fixed_size_list/vtable/operations.rs @@ -11,6 +11,8 @@ use crate::arrays::fixed_size_list::FixedSizeListArrayExt; use crate::scalar::Scalar; impl OperationsVTable for FixedSizeList { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, FixedSizeList>, index: usize, diff --git a/vortex-array/src/arrays/interleave/mod.rs b/vortex-array/src/arrays/interleave/mod.rs index d29981249d4..fb649cf6e19 100644 --- a/vortex-array/src/arrays/interleave/mod.rs +++ b/vortex-array/src/arrays/interleave/mod.rs @@ -389,6 +389,8 @@ impl VTable for Interleave { } impl OperationsVTable for Interleave { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Interleave>, index: usize, diff --git a/vortex-array/src/arrays/list/vtable/operations.rs b/vortex-array/src/arrays/list/vtable/operations.rs index 02c686cd1f1..668cc18a32f 100644 --- a/vortex-array/src/arrays/list/vtable/operations.rs +++ b/vortex-array/src/arrays/list/vtable/operations.rs @@ -13,6 +13,8 @@ use crate::arrays::list::ListArrayExt; use crate::scalar::Scalar; impl OperationsVTable for List { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, List>, index: usize, diff --git a/vortex-array/src/arrays/listview/vtable/operations.rs b/vortex-array/src/arrays/listview/vtable/operations.rs index f0cb9539cc3..8608985463d 100644 --- a/vortex-array/src/arrays/listview/vtable/operations.rs +++ b/vortex-array/src/arrays/listview/vtable/operations.rs @@ -13,6 +13,8 @@ use crate::arrays::listview::ListViewArrayExt; use crate::scalar::Scalar; impl OperationsVTable for ListView { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, ListView>, index: usize, diff --git a/vortex-array/src/arrays/map/vtable/operations.rs b/vortex-array/src/arrays/map/vtable/operations.rs index d6e8fe87f12..66b8e47a75a 100644 --- a/vortex-array/src/arrays/map/vtable/operations.rs +++ b/vortex-array/src/arrays/map/vtable/operations.rs @@ -13,6 +13,8 @@ use crate::arrays::struct_::StructArrayExt; use crate::scalar::Scalar; impl OperationsVTable for Map { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Map>, index: usize, diff --git a/vortex-array/src/arrays/masked/vtable/operations.rs b/vortex-array/src/arrays/masked/vtable/operations.rs index c82d0bf03ed..a418e23f217 100644 --- a/vortex-array/src/arrays/masked/vtable/operations.rs +++ b/vortex-array/src/arrays/masked/vtable/operations.rs @@ -11,6 +11,8 @@ use crate::arrays::masked::MaskedArraySlotsExt; use crate::scalar::Scalar; impl OperationsVTable for Masked { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Masked>, index: usize, diff --git a/vortex-array/src/arrays/null/mod.rs b/vortex-array/src/arrays/null/mod.rs index 8479b7137cc..b4d4536c9cf 100644 --- a/vortex-array/src/arrays/null/mod.rs +++ b/vortex-array/src/arrays/null/mod.rs @@ -175,6 +175,8 @@ impl Array { } impl OperationsVTable for Null { + type ProbeState = (); + fn scalar_at( _array: ArrayView<'_, Null>, _index: usize, diff --git a/vortex-array/src/arrays/patched/vtable/operations.rs b/vortex-array/src/arrays/patched/vtable/operations.rs index 51dd1fc9e3c..bd7440965ab 100644 --- a/vortex-array/src/arrays/patched/vtable/operations.rs +++ b/vortex-array/src/arrays/patched/vtable/operations.rs @@ -14,6 +14,8 @@ use crate::optimizer::ArrayOptimizer; use crate::scalar::Scalar; impl OperationsVTable for Patched { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Patched>, index: usize, diff --git a/vortex-array/src/arrays/piecewise_sequence/vtable.rs b/vortex-array/src/arrays/piecewise_sequence/vtable.rs index d68a00815f3..d39bd6eccac 100644 --- a/vortex-array/src/arrays/piecewise_sequence/vtable.rs +++ b/vortex-array/src/arrays/piecewise_sequence/vtable.rs @@ -145,6 +145,8 @@ impl VTable for PiecewiseSequence { } impl OperationsVTable for PiecewiseSequence { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, PiecewiseSequence>, index: usize, diff --git a/vortex-array/src/arrays/primitive/vtable/operations.rs b/vortex-array/src/arrays/primitive/vtable/operations.rs index ddeaa386485..0f77ed22d9d 100644 --- a/vortex-array/src/arrays/primitive/vtable/operations.rs +++ b/vortex-array/src/arrays/primitive/vtable/operations.rs @@ -6,18 +6,30 @@ use vortex_error::VortexResult; use crate::ExecutionCtx; use crate::array::ArrayView; use crate::array::OperationsVTable; +use crate::array::ProbeState; use crate::arrays::Primitive; use crate::match_each_native_ptype; use crate::scalar::Scalar; impl OperationsVTable for Primitive { - fn scalar_at( - array: ArrayView<'_, Primitive>, + type ProbeState = (); + + fn probe_scalar( + state: &mut ProbeState<'_, Primitive>, index: usize, _ctx: &mut ExecutionCtx, ) -> VortexResult { + let array = state.array(); Ok(match_each_native_ptype!(array.ptype(), |T| { Scalar::primitive(array.as_slice::()[index], array.dtype().nullability()) })) } + + fn scalar_at( + array: ArrayView<'_, Primitive>, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + Self::probe_scalar(&mut ProbeState::once(array), index, ctx) + } } diff --git a/vortex-array/src/arrays/scalar_fn/vtable/operations.rs b/vortex-array/src/arrays/scalar_fn/vtable/operations.rs index 40d75906356..21af5728572 100644 --- a/vortex-array/src/arrays/scalar_fn/vtable/operations.rs +++ b/vortex-array/src/arrays/scalar_fn/vtable/operations.rs @@ -15,6 +15,8 @@ use crate::scalar::Scalar; use crate::scalar_fn::VecExecutionArgs; impl OperationsVTable for ScalarFn { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, ScalarFn>, index: usize, diff --git a/vortex-array/src/arrays/shared/vtable.rs b/vortex-array/src/arrays/shared/vtable.rs index 758d091b557..b1a9a4d7b15 100644 --- a/vortex-array/src/arrays/shared/vtable.rs +++ b/vortex-array/src/arrays/shared/vtable.rs @@ -125,6 +125,8 @@ impl VTable for Shared { } } impl OperationsVTable for Shared { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Shared>, index: usize, diff --git a/vortex-array/src/arrays/slice/vtable.rs b/vortex-array/src/arrays/slice/vtable.rs index c88c28a9a4d..c09d81ee4f1 100644 --- a/vortex-array/src/arrays/slice/vtable.rs +++ b/vortex-array/src/arrays/slice/vtable.rs @@ -169,6 +169,8 @@ impl VTable for Slice { } } impl OperationsVTable for Slice { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Slice>, index: usize, diff --git a/vortex-array/src/arrays/struct_/vtable/operations.rs b/vortex-array/src/arrays/struct_/vtable/operations.rs index b491e231520..a139545c17d 100644 --- a/vortex-array/src/arrays/struct_/vtable/operations.rs +++ b/vortex-array/src/arrays/struct_/vtable/operations.rs @@ -2,27 +2,39 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors use vortex_error::VortexResult; +use vortex_error::vortex_err; use crate::ExecutionCtx; use crate::array::ArrayView; use crate::array::OperationsVTable; +use crate::array::ProbeState; use crate::arrays::Struct; use crate::arrays::struct_::StructArrayExt; +use crate::arrays::struct_::StructSlots; use crate::scalar::Scalar; use crate::scalar::ScalarValue; impl OperationsVTable for Struct { - fn scalar_at( - array: ArrayView<'_, Struct>, + type ProbeState = (); + + fn probe_scalar( + state: &mut ProbeState<'_, Struct>, index: usize, ctx: &mut ExecutionCtx, ) -> VortexResult { - let field_values = array - .iter_unmasked_fields() - .map(|field| field.execute_scalar(index, ctx).map(Scalar::into_value)) - .collect::>>()?; + let array = state.array(); + let nfields = array.iter_unmasked_fields().len(); + let mut field_values = Vec::with_capacity(nfields); + for field in 0..nfields { + let slot = StructSlots::FIELDS_OFFSET + field; + let value = state + .slot(slot)? + .ok_or_else(|| vortex_err!("Struct field slot {slot} is absent"))? + .execute_scalar(index, ctx)?; + field_values.push(value.into_value()); + } // SAFETY: The vtable guarantees index is in-bounds and non-null before this is called. - // Each field's scalar_at returns a value with the field's own dtype. + // Each field read returns a value with the field's own dtype. Ok(unsafe { Scalar::new_unchecked( array.dtype().clone(), @@ -30,4 +42,12 @@ impl OperationsVTable for Struct { ) }) } + + fn scalar_at( + array: ArrayView<'_, Struct>, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + Self::probe_scalar(&mut ProbeState::once(array), index, ctx) + } } diff --git a/vortex-array/src/arrays/union/vtable/operations.rs b/vortex-array/src/arrays/union/vtable/operations.rs index 39ba95ed4ee..78b9759ab1c 100644 --- a/vortex-array/src/arrays/union/vtable/operations.rs +++ b/vortex-array/src/arrays/union/vtable/operations.rs @@ -14,6 +14,8 @@ use crate::arrays::union::UnionArraySlotsExt; use crate::scalar::Scalar; impl OperationsVTable for Union { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Union>, index: usize, diff --git a/vortex-array/src/arrays/varbin/vtable/operations.rs b/vortex-array/src/arrays/varbin/vtable/operations.rs index e11043e605a..cc5090d8912 100644 --- a/vortex-array/src/arrays/varbin/vtable/operations.rs +++ b/vortex-array/src/arrays/varbin/vtable/operations.rs @@ -12,6 +12,8 @@ use crate::arrays::varbin::varbin_scalar; use crate::scalar::Scalar; impl OperationsVTable for VarBin { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, VarBin>, index: usize, diff --git a/vortex-array/src/arrays/varbinview/vtable/operations.rs b/vortex-array/src/arrays/varbinview/vtable/operations.rs index 1a1f20a0dbe..53866d29fb5 100644 --- a/vortex-array/src/arrays/varbinview/vtable/operations.rs +++ b/vortex-array/src/arrays/varbinview/vtable/operations.rs @@ -11,6 +11,8 @@ use crate::arrays::varbin::varbin_scalar; use crate::scalar::Scalar; impl OperationsVTable for VarBinView { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, VarBinView>, index: usize, diff --git a/vortex-array/src/arrays/variant/vtable/operations.rs b/vortex-array/src/arrays/variant/vtable/operations.rs index c7f7d36fd97..ed61223d637 100644 --- a/vortex-array/src/arrays/variant/vtable/operations.rs +++ b/vortex-array/src/arrays/variant/vtable/operations.rs @@ -12,6 +12,8 @@ use crate::arrays::variant::VariantArraySlotsExt; use crate::scalar::Scalar; impl OperationsVTable for Variant { + type ProbeState = (); + fn scalar_at( array: ArrayView<'_, Variant>, index: usize, diff --git a/vortex-array/src/validity.rs b/vortex-array/src/validity.rs index a2c10e30ca7..efe4a9172d9 100644 --- a/vortex-array/src/validity.rs +++ b/vortex-array/src/validity.rs @@ -163,6 +163,7 @@ impl Validity { } /// Returns whether the `index` item is valid, using `ctx` to execute the validity array. + // Todo(joe): deprecate this #[inline] pub fn execute_is_valid(&self, index: usize, ctx: &mut ExecutionCtx) -> VortexResult { Ok(match self { @@ -177,6 +178,7 @@ impl Validity { } /// Returns whether the `index` item is null, using `ctx` to execute the validity array. + // Todo(joe): deprecate this #[inline] pub fn execute_is_null(&self, index: usize, ctx: &mut ExecutionCtx) -> VortexResult { Ok(!self.execute_is_valid(index, ctx)?) diff --git a/vortex-python/src/arrays/py/vtable.rs b/vortex-python/src/arrays/py/vtable.rs index 5ec749d3a5f..5697d5a7574 100644 --- a/vortex-python/src/arrays/py/vtable.rs +++ b/vortex-python/src/arrays/py/vtable.rs @@ -122,6 +122,8 @@ impl VTable for PythonVTable { } impl OperationsVTable for PythonVTable { + type ProbeState = (); + fn scalar_at( _array: ArrayView<'_, PythonVTable>, _index: usize,