diff --git a/encodings/fastlanes/src/rle/mod.rs b/encodings/fastlanes/src/rle/mod.rs index 25610742c19..57b37c7201c 100644 --- a/encodings/fastlanes/src/rle/mod.rs +++ b/encodings/fastlanes/src/rle/mod.rs @@ -9,6 +9,7 @@ pub use array::RLESlots; mod compute; mod kernel; +mod probe; mod vtable; pub use vtable::RLE; diff --git a/encodings/fastlanes/src/rle/probe.rs b/encodings/fastlanes/src/rle/probe.rs new file mode 100644 index 00000000000..8f80aa09228 --- /dev/null +++ b/encodings/fastlanes/src/rle/probe.rs @@ -0,0 +1,86 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! RLE scalar access through the indices, offsets, and values slots. + +use vortex_array::ExecutionCtx; +use vortex_array::ProbeState; +use vortex_array::scalar::Scalar; +use vortex_error::VortexResult; +use vortex_error::vortex_err; + +use crate::FL_CHUNK_SIZE; +use crate::RLE; +use crate::rle::RLEArrayExt; +use crate::rle::RLESlots; + +/// The slice's base value offset, retained across repeated lookups. +#[derive(Default)] +pub struct RleProbeState { + base: Option, +} + +pub(crate) fn scalar_at( + state: &mut ProbeState<'_, RLE>, + index: usize, + ctx: &mut ExecutionCtx, +) -> VortexResult { + let array = state.array(); + let logical_index = array.offset() + index; + // Hold the retained state while reading children: the base offset is cached across rows. + let (mut retained, mut children) = state.split(); + let code = children + .slot(RLESlots::INDICES)? + .ok_or_else(|| vortex_err!("RLE indices slot is missing"))? + .execute_scalar(logical_index, ctx)?; + let Some(code) = code.as_primitive().as_::() else { + return Ok(Scalar::null(array.dtype().clone())); + }; + + let chunk = logical_index / FL_CHUNK_SIZE; + let offset = if chunk == 0 { + 0 + } else { + let base = match retained.as_deref().and_then(|state| state.base) { + Some(base) => base, + None => { + let value = read_offset( + children + .slot(RLESlots::VALUES_IDX_OFFSETS)? + .ok_or_else(|| vortex_err!("RLE offsets slot is missing"))? + .execute_scalar(0, ctx)?, + )?; + if let Some(state) = retained.as_mut() { + state.base = Some(value); + } + value + } + }; + read_offset( + children + .slot(RLESlots::VALUES_IDX_OFFSETS)? + .ok_or_else(|| vortex_err!("RLE offsets slot is missing"))? + .execute_scalar(chunk, ctx)?, + )? + .checked_sub(base) + .ok_or_else(|| vortex_err!("RLE offsets precede the slice base"))? + }; + let value_index = offset + .checked_add(code) + .ok_or_else(|| vortex_err!("RLE value index overflow"))?; + let scalar = children + .slot(RLESlots::VALUES)? + .ok_or_else(|| vortex_err!("RLE values slot is missing"))? + .execute_scalar(value_index, ctx)?; + Scalar::try_new(array.dtype().clone(), scalar.into_value()) +} + +fn read_offset(scalar: Scalar) -> VortexResult { + scalar + .as_primitive() + .as_::() + .ok_or_else(|| vortex_err!("RLE offset must be a non-null usize")) +} + +#[cfg(test)] +mod tests; diff --git a/encodings/fastlanes/src/rle/probe/tests.rs b/encodings/fastlanes/src/rle/probe/tests.rs new file mode 100644 index 00000000000..fec027edf58 --- /dev/null +++ b/encodings/fastlanes/src/rle/probe/tests.rs @@ -0,0 +1,94 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use rstest::rstest; +use vortex_array::ArrayProbe; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::RepeatedArrayProbe; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::assert_arrays_eq; +use vortex_array::builders::builder_with_capacity_in; +use vortex_array::builtins::ArrayBuiltins; +use vortex_array::scalar_fn::fns::operators::Operator; +use vortex_array::validity::Validity; +use vortex_error::VortexResult; + +use crate::RLE; +use crate::RLEData; + +/// A one-off or retained probe over `array`, the retained one living in `retained`. +fn probe_for<'a>( + array: &'a ArrayRef, + retained: &'a mut Option, + repeated: bool, +) -> ArrayProbe<'a> { + if repeated { + retained.insert(array.repeated_probe()).as_probe() + } else { + array.probe() + } +} + +#[rstest] +fn random_access_across_chunks_and_nulls( + #[values(false, true)] repeated: bool, + #[values(false, true)] sliced: bool, +) -> VortexResult<()> { + let mut ctx = crate::test::SESSION.create_execution_ctx(); + let input = + PrimitiveArray::from_option_iter((0..8192u32).map(|i| (i % 11 != 0).then_some(i / 16))); + let encoded = RLEData::encode(input.as_view(), &mut ctx)?.into_array(); + let range = if sliced { 1777..7333 } else { 0..8192 }; + let source = if sliced { + encoded + .slice(range.clone())? + .execute::(&mut ctx)? + } else { + encoded + }; + assert!(source.is::()); + let input = input.slice(range)?; + let indices = [0u32, 1, 1023, 1024, 2048, 2047, 4097, 11, 33, 17, 0]; + let mut actual = builder_with_capacity_in(source.dtype(), indices.len(), ctx.allocator()); + let mut retained = None; + let mut probe = probe_for(&source, &mut retained, repeated); + for index in indices { + actual.append_scalar(&probe.execute_scalar(index as usize, &mut ctx)?)?; + } + assert_arrays_eq!( + actual.finish(), + input.take(PrimitiveArray::from_iter(indices).into_array())?, + &mut ctx + ); + assert!(probe.execute_scalar(source.len(), &mut ctx).is_err()); + Ok(()) +} + +#[rstest] +fn lazy_validity_does_not_evaluate_unrequested_rows( + #[values(false, true)] repeated: bool, +) -> VortexResult<()> { + let mut ctx = crate::test::SESSION.create_execution_ctx(); + let numerators = PrimitiveArray::from_iter(vec![1u32; 1024]).into_array(); + let denominators = + PrimitiveArray::from_iter((0..1024).map(|i| u32::from(i != 1023))).into_array(); + let validity = numerators + .binary(denominators, Operator::Div)? + .binary(numerators, Operator::Eq)?; + let array = RLE::try_new( + PrimitiveArray::from_iter([42u32]).into_array(), + PrimitiveArray::new(vec![0u16; 1024], Validity::Array(validity)).into_array(), + PrimitiveArray::from_iter([0u64]).into_array(), + 0, + 1024, + )? + .into_array(); + let expected = array.execute_scalar(0, &mut ctx)?; + let mut retained = None; + let mut probe = probe_for(&array, &mut retained, repeated); + assert_eq!(probe.execute_scalar(0, &mut ctx)?, expected); + assert!(probe.execute_scalar(1023, &mut ctx).is_err()); + Ok(()) +} diff --git a/encodings/fastlanes/src/rle/vtable/operations.rs b/encodings/fastlanes/src/rle/vtable/operations.rs index ba6a624f1f5..92060060d23 100644 --- a/encodings/fastlanes/src/rle/vtable/operations.rs +++ b/encodings/fastlanes/src/rle/vtable/operations.rs @@ -3,42 +3,32 @@ use vortex_array::ArrayView; use vortex_array::ExecutionCtx; +use vortex_array::ProbeState; use vortex_array::scalar::Scalar; use vortex_array::vtable::OperationsVTable; -use vortex_error::VortexExpect; use vortex_error::VortexResult; use super::RLE; -use crate::FL_CHUNK_SIZE; -use crate::rle::RLEArrayExt; -use crate::rle::RLEArraySlotsExt; +use crate::rle::probe; +use crate::rle::probe::RleProbeState; impl OperationsVTable for RLE { - type ProbeState = (); + type ProbeState = RleProbeState; + + fn probe_scalar( + state: &mut ProbeState<'_, RLE>, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + probe::scalar_at(state, index, ctx) + } fn scalar_at( array: ArrayView<'_, RLE>, index: usize, ctx: &mut ExecutionCtx, ) -> VortexResult { - let offset_in_chunk = array.offset(); - let chunk_relative_idx = array - .indices() - .execute_scalar(offset_in_chunk + index, ctx)?; - - let chunk_relative_idx = chunk_relative_idx - .as_primitive() - .as_::() - .vortex_expect("Index must not be null"); - - let chunk_id = (offset_in_chunk + index) / FL_CHUNK_SIZE; - let value_idx_offset = array.values_idx_offset(chunk_id, ctx); - - let scalar = array - .values() - .execute_scalar(value_idx_offset + chunk_relative_idx, ctx)?; - - Scalar::try_new(array.dtype().clone(), scalar.into_value()) + probe::scalar_at(&mut ProbeState::once(array), index, ctx) } } @@ -54,9 +44,9 @@ mod tests { use vortex_array::validity::Validity; use vortex_buffer::Buffer; use vortex_buffer::buffer; + use vortex_error::VortexExpect; use vortex_session::VortexSession; - use super::*; use crate::RLE; use crate::RLEArray; use crate::RLEData; diff --git a/encodings/pco/benches/scalar.rs b/encodings/pco/benches/scalar.rs index eff55fceb8c..6a928c4b4a0 100644 --- a/encodings/pco/benches/scalar.rs +++ b/encodings/pco/benches/scalar.rs @@ -1,8 +1,9 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright the Vortex contributors -//! Scalar reads out of a PCO array. Cases are `(access_count, nullable, scattered)`; clustered -//! indices stay inside one page so a retained decode can be reused, scattered ones cross pages. +//! Scalar reads out of a PCO array through one repeated probe. Cases are +//! `(access_count, nullable, scattered)`; clustered indices stay inside one page so the retained +//! decode is reused, scattered ones cross pages. use std::sync::LazyLock; @@ -65,11 +66,17 @@ fn scalar_access(bencher: Bencher, (count, nullable, scattered): (usize, bool, b let array = pco(nullable); let indices = indices(count, scattered); bencher - .with_inputs(|| (SESSION.create_execution_ctx(), Vec::with_capacity(count))) - .bench_refs(|(ctx, scalars)| { + .with_inputs(|| { + ( + SESSION.create_execution_ctx(), + array.repeated_probe(), + Vec::with_capacity(count), + ) + }) + .bench_refs(|(ctx, probe, scalars)| { for &index in &indices { scalars.push( - array + probe .execute_scalar(index, ctx) .vortex_expect("scalar access"), ); diff --git a/encodings/pco/src/array.rs b/encodings/pco/src/array.rs index 97a17fe6a3c..2d9d563f473 100644 --- a/encodings/pco/src/array.rs +++ b/encodings/pco/src/array.rs @@ -29,6 +29,7 @@ use vortex_array::EqMode; use vortex_array::ExecutionCtx; use vortex_array::ExecutionResult; use vortex_array::IntoArray; +use vortex_array::ProbeState; use vortex_array::TypedArrayRef; use vortex_array::array_slots; use vortex_array::arrays::Primitive; @@ -778,19 +779,23 @@ impl ValidityVTable for Pco { } impl OperationsVTable for Pco { - type ProbeState = (); + type ProbeState = crate::probe::PcoProbeState; + + fn probe_scalar( + state: &mut ProbeState<'_, Pco>, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + let array = state.array(); + crate::probe::scalar_at(array, index, state.retained(), ctx) + } fn scalar_at( array: ArrayView<'_, Pco>, index: usize, ctx: &mut ExecutionCtx, ) -> VortexResult { - let unsliced_validity = array.unsliced_validity(); - array - ._slice(index, index + 1) - .decompress(&unsliced_validity, ctx)? - .into_array() - .execute_scalar(0, ctx) + crate::probe::scalar_at(array, index, None, ctx) } } diff --git a/encodings/pco/src/lib.rs b/encodings/pco/src/lib.rs index 601bd0826bd..70d71abf7d2 100644 --- a/encodings/pco/src/lib.rs +++ b/encodings/pco/src/lib.rs @@ -20,6 +20,7 @@ mod array; mod compute; +mod probe; mod rules; mod slice; diff --git a/encodings/pco/src/probe.rs b/encodings/pco/src/probe.rs new file mode 100644 index 00000000000..2012d8abf51 --- /dev/null +++ b/encodings/pco/src/probe.rs @@ -0,0 +1,164 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! PCO probes retain every decoded page and index the compacted, non-null value positions. + +use std::ops::Range; + +use pco::data_types::Number; +use pco::data_types::NumberType; +use pco::match_number_enum; +use pco::wrapped::FileDecompressor; +use vortex_array::ArrayView; +use vortex_array::ExecutionCtx; +use vortex_array::IntoArray; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::dtype::NativePType; +use vortex_array::dtype::half; +use vortex_array::match_each_native_ptype; +use vortex_array::scalar::Scalar; +use vortex_array::validity::Validity; +use vortex_buffer::BufferMut; +use vortex_error::VortexExpect; +use vortex_error::VortexResult; +use vortex_mask::Mask; + +use crate::Pco; +use crate::PcoArrayExt; +use crate::array::number_type_from_ptype; +use crate::array::vortex_err_from_pco; + +const RANK_STRIDE: usize = 512; + +/// State retained by repeated PCO probes. Construction performs no allocation. +/// +/// The mask and rank index use unsliced logical rows; page boundaries use compacted value +/// positions. Every page decoded so far is retained, so the state grows with the pages touched +/// and never beyond the decompressed array. +#[derive(Default)] +pub struct PcoProbeState { + validity: Option, + rank: Vec, + pages: Vec, + /// The page the previous probe landed in, `None` until the first probe resolves one. + last_page: Option, +} + +struct Page { + values: Range, + chunk: usize, + decoded: Option, +} + +pub(crate) fn scalar_at( + array: ArrayView<'_, Pco>, + index: usize, + state: Option<&mut PcoProbeState>, + ctx: &mut ExecutionCtx, +) -> VortexResult { + let Some(state) = state else { + let validity = array.unsliced_validity(); + return array + ._slice(index, index + 1) + .decompress(&validity, ctx)? + .into_array() + .execute_scalar(0, ctx); + }; + let mask = match &mut state.validity { + Some(mask) => mask, + slot @ None => slot.insert( + array + .unsliced_validity() + .execute_mask(array.unsliced_n_rows(), ctx)?, + ), + }; + let logical_index = array.slice_start() + index; + + let value_index = match mask { + Mask::AllTrue(_) => logical_index, + Mask::AllFalse(_) => unreachable!("probe dispatch checks validity"), + Mask::Values(values) => { + let bits = values.bit_buffer(); + if state.rank.is_empty() { + state.rank.reserve(bits.len().div_ceil(RANK_STRIDE)); + let mut count = 0; + for start in (0..bits.len()).step_by(RANK_STRIDE) { + state.rank.push(count); + count += bits.count_range(start, (start + RANK_STRIDE).min(bits.len())); + } + } + let block = logical_index / RANK_STRIDE; + state.rank[block] + bits.count_range(block * RANK_STRIDE, logical_index) + } + }; + + if state.pages.is_empty() { + let mut start = 0; + for (chunk, metadata) in array.metadata.chunks.iter().enumerate() { + for page in &metadata.pages { + let end = start + page.n_values as usize; + state.pages.push(Page { + values: start..end, + chunk, + decoded: None, + }); + start = end; + } + } + } + let page_index = match state.last_page { + Some(last) if state.pages[last].values.contains(&value_index) => last, + _ => state + .pages + .partition_point(|page| page.values.end <= value_index), + }; + state.last_page = Some(page_index); + let page = state + .pages + .get_mut(page_index) + .vortex_expect("PCO pages cover every valid value index"); + let values = match &mut page.decoded { + Some(decoded) => decoded, + slot @ None => { + let decoded = match_number_enum!( + number_type_from_ptype(array.dtype().as_ptype()), + NumberType => { + decode_page::(array, page.chunk, page.values.len(), array.pages[page_index].as_slice(), ctx)? + } + ); + slot.insert(decoded) + } + }; + Ok(match_each_native_ptype!(values.ptype(), |T| { + Scalar::primitive( + values.as_slice::()[value_index - page.values.start], + array.dtype().nullability(), + ) + })) +} + +fn decode_page( + array: ArrayView<'_, Pco>, + chunk: usize, + n_values: usize, + buffer: &[u8], + ctx: &mut ExecutionCtx, +) -> VortexResult { + let (file, _) = + FileDecompressor::new(array.metadata.header.as_slice()).map_err(vortex_err_from_pco)?; + let (mut chunk, _) = file + .chunk_decompressor::(array.chunk_metas[chunk].as_ref()) + .map_err(vortex_err_from_pco)?; + let mut decoder = chunk + .page_decompressor(buffer, n_values) + .map_err(vortex_err_from_pco)?; + let mut values = BufferMut::::with_capacity_in(n_values, ctx.allocator().clone()); + // SAFETY: the buffer reserves `n_values` elements, and the page decompressor was built for + // exactly that count, so `read` writes every element before anything observes it. + unsafe { values.set_len(n_values) }; + decoder.read(&mut values).map_err(vortex_err_from_pco)?; + Ok(PrimitiveArray::new(values.freeze(), Validity::NonNullable)) +} + +#[cfg(test)] +mod tests; diff --git a/encodings/pco/src/probe/tests.rs b/encodings/pco/src/probe/tests.rs new file mode 100644 index 00000000000..f4c9740cf2c --- /dev/null +++ b/encodings/pco/src/probe/tests.rs @@ -0,0 +1,107 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use rstest::rstest; +use vortex_array::ArrayProbe; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::RepeatedArrayProbe; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::assert_arrays_eq; +use vortex_array::builders::builder_with_capacity_in; +use vortex_array::validity::Validity; +use vortex_error::VortexResult; + +use crate::Pco; + +/// A one-off or retained probe over `array`, the retained one living in `retained`. +fn probe_for<'a>( + array: &'a ArrayRef, + retained: &'a mut Option, + repeated: bool, +) -> ArrayProbe<'a> { + if repeated { + retained.insert(array.repeated_probe()).as_probe() + } else { + array.probe() + } +} + +#[rstest] +fn random_access( + #[values(false, true)] repeated: bool, + #[values(false, true)] nullable: bool, + #[values(false, true)] sliced: bool, +) -> VortexResult<()> { + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let input = if nullable { + PrimitiveArray::from_option_iter((0..4096i32).map(|i| (i % 7 != 0).then_some(i * 19))) + } else { + PrimitiveArray::from_iter((0..4096i32).map(|i| i * 19)) + }; + let encoded = Pco::from_primitive(input.as_view(), 3, 128, &mut ctx)?.into_array(); + let range = if sliced { 777..3333 } else { 0..4096 }; + let source = encoded.slice(range.clone())?; + assert!(source.is::()); + let input = input.slice(range)?; + let indices = [0u32, 1, 127, 128, 512, 2048, 129, 2, 1024, 0]; + let mut actual = builder_with_capacity_in(source.dtype(), indices.len(), ctx.allocator()); + let mut retained = None; + let mut probe = probe_for(&source, &mut retained, repeated); + for index in indices { + actual.append_scalar(&probe.execute_scalar(index as usize, &mut ctx)?)?; + } + assert_arrays_eq!( + actual.finish(), + input.take(PrimitiveArray::from_iter(indices).into_array())?, + &mut ctx + ); + assert!(probe.execute_scalar(source.len(), &mut ctx).is_err()); + Ok(()) +} + +#[rstest] +#[case(PrimitiveArray::from_iter([1u16, 9, 32768, 65535]))] +#[case(PrimitiveArray::from_iter([i64::MIN, -1, 0, i64::MAX]))] +#[case(PrimitiveArray::from_iter([1.25f64, -2.5, 0.0, f64::INFINITY]))] +fn preserves_physical_type(#[case] input: PrimitiveArray) -> VortexResult<()> { + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let encoded = Pco::from_primitive(input.as_view(), 3, 2, &mut ctx)?.into_array(); + let mut probe = encoded.repeated_probe(); + let mut actual = builder_with_capacity_in(input.dtype(), input.len(), ctx.allocator()); + for i in 0..input.len() { + actual.append_scalar(&probe.execute_scalar(i, &mut ctx)?)?; + } + assert_arrays_eq!(actual.finish(), input, &mut ctx); + Ok(()) +} + +#[test] +fn all_null_access_returns_null() -> VortexResult<()> { + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let input = PrimitiveArray::new(vec![0i32; 128], Validity::AllInvalid); + let encoded = Pco::from_primitive(input.as_view(), 3, 128, &mut ctx)?; + let encoded = encoded.into_array(); + let mut probe = encoded.repeated_probe(); + assert!(probe.execute_scalar(42, &mut ctx)?.is_null()); + Ok(()) +} + +#[test] +fn probe_outlives_source_and_moves() -> VortexResult<()> { + let mut ctx = vortex_array::array_session().create_execution_ctx(); + let input = PrimitiveArray::from_iter(0..4096u32); + let array = Pco::from_primitive(input.as_view(), 3, 512, &mut ctx)?.into_array(); + let probe = RepeatedArrayProbe::new(array.clone()); + drop(array); + let mut moved = (probe, ()); + for index in [1, 5, 127, 255, 511, 1023, 0, 1, 4095] { + assert_eq!( + moved.0.execute_scalar(index, &mut ctx)?, + u32::try_from(index)?.into() + ); + } + assert!(moved.0.execute_scalar(4096, &mut ctx).is_err()); + Ok(()) +} diff --git a/encodings/runend/src/lib.rs b/encodings/runend/src/lib.rs index b9afd20151c..98c661e6e36 100644 --- a/encodings/runend/src/lib.rs +++ b/encodings/runend/src/lib.rs @@ -15,6 +15,7 @@ pub mod decompress_bool; mod iter; mod kernel; pub mod ops; +mod probe; mod rules; #[cfg(test)] #[cfg(not(codspeed))] diff --git a/encodings/runend/src/ops.rs b/encodings/runend/src/ops.rs index 46949d84e1b..637fa8ecd25 100644 --- a/encodings/runend/src/ops.rs +++ b/encodings/runend/src/ops.rs @@ -4,6 +4,7 @@ use vortex_array::ArrayRef; use vortex_array::ArrayView; use vortex_array::ExecutionCtx; +use vortex_array::ProbeState; use vortex_array::match_each_unsigned_integer_ptype; use vortex_array::scalar::Scalar; use vortex_array::search_sorted::SearchResult; @@ -14,19 +15,24 @@ use vortex_array::vtable::OperationsVTable; use vortex_error::VortexResult; use crate::RunEnd; -use crate::array::RunEndArrayExt; -use crate::array::RunEndArraySlotsExt; impl OperationsVTable for RunEnd { type ProbeState = (); + fn probe_scalar( + state: &mut ProbeState<'_, RunEnd>, + index: usize, + ctx: &mut ExecutionCtx, + ) -> VortexResult { + crate::probe::scalar_at(state, index, ctx) + } + fn scalar_at( array: ArrayView<'_, RunEnd>, index: usize, ctx: &mut ExecutionCtx, ) -> VortexResult { - let physical_index = array.find_physical_index(index, ctx)?; - array.values().execute_scalar(physical_index, ctx) + crate::probe::scalar_at(&mut ProbeState::once(array), index, ctx) } } diff --git a/encodings/runend/src/probe.rs b/encodings/runend/src/probe.rs new file mode 100644 index 00000000000..abcb3d22e0d --- /dev/null +++ b/encodings/runend/src/probe.rs @@ -0,0 +1,50 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! RunEnd probes retain their ends and values child probes across lookups. + +use vortex_array::ExecutionCtx; +use vortex_array::ProbeState; +use vortex_array::scalar::Scalar; +use vortex_error::VortexResult; +use vortex_error::vortex_err; + +use crate::RunEnd; +use crate::RunEndArrayExt; +use crate::RunEndSlots; + +pub(crate) fn scalar_at( + state: &mut ProbeState<'_, RunEnd>, + index: usize, + ctx: &mut ExecutionCtx, +) -> VortexResult { + let array = state.array(); + let logical_index = array + .offset() + .checked_add(index) + .ok_or_else(|| vortex_err!("RunEnd logical index overflow"))?; + // Search for the first end strictly greater than the logical index. Every comparison + // uses the same ends probe, so a retained read keeps the child's preparation within and + // between searches. + let mut ends = state + .slot(RunEndSlots::ENDS)? + .ok_or_else(|| vortex_err!("RunEnd ends slot is missing"))?; + let mut left = 0; + let mut right = ends.array().len(); + while left < right { + let mid = left + (right - left) / 2; + let end = usize::try_from(&ends.execute_scalar(mid, ctx)?)?; + if end <= logical_index { + left = mid + 1; + } else { + right = mid; + } + } + state + .slot(RunEndSlots::VALUES)? + .ok_or_else(|| vortex_err!("RunEnd values slot is missing"))? + .execute_scalar(left, ctx) +} + +#[cfg(test)] +mod tests; diff --git a/encodings/runend/src/probe/tests.rs b/encodings/runend/src/probe/tests.rs new file mode 100644 index 00000000000..c57b930b61e --- /dev/null +++ b/encodings/runend/src/probe/tests.rs @@ -0,0 +1,102 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +use rstest::rstest; +use vortex_array::ArrayProbe; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::RepeatedArrayProbe; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::PrimitiveArray; +use vortex_array::arrays::VarBinViewArray; +use vortex_array::assert_arrays_eq; +use vortex_array::builders::builder_with_capacity_in; +use vortex_array::builtins::ArrayBuiltins; +use vortex_array::dtype::DType; +use vortex_array::dtype::Nullability; +use vortex_array::scalar_fn::fns::operators::Operator; +use vortex_array::validity::Validity; +use vortex_buffer::Buffer; +use vortex_buffer::buffer; +use vortex_error::VortexResult; + +use crate::RunEnd; +use crate::tests::SESSION; + +/// A one-off or retained probe over `array`, the retained one living in `retained`. +fn probe_for<'a>( + array: &'a ArrayRef, + retained: &'a mut Option, + repeated: bool, +) -> ArrayProbe<'a> { + if repeated { + retained.insert(array.repeated_probe()).as_probe() + } else { + array.probe() + } +} + +#[rstest] +fn sliced_strings_and_null_runs(#[values(false, true)] repeated: bool) -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let array = RunEnd::try_new_offset_length( + buffer![2u8, 5, 9].into_array(), + VarBinViewArray::from_iter( + [Some("first"), None, Some("last")], + DType::Utf8(Nullability::Nullable), + ) + .into_array(), + 1, + 7, + &mut ctx, + )? + .into_array(); + let mut retained = None; + let mut probe = probe_for(&array, &mut retained, repeated); + let mut actual = builder_with_capacity_in(array.dtype(), 5, ctx.allocator()); + for index in [6, 0, 1, 3, 4] { + actual.append_scalar(&probe.execute_scalar(index, &mut ctx)?)?; + } + assert_arrays_eq!( + actual.finish(), + VarBinViewArray::from_iter( + [Some("last"), Some("first"), None, None, Some("last")], + DType::Utf8(Nullability::Nullable) + ), + &mut ctx + ); + assert!(probe.execute_scalar(array.len(), &mut ctx).is_err()); + Ok(()) +} + +#[test] +fn empty_probe_checks_bounds() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let array = RunEnd::try_new( + Buffer::::empty().into_array(), + Buffer::::empty().into_array(), + &mut ctx, + )? + .into_array(); + assert!(array.probe().execute_scalar(0, &mut ctx).is_err()); + assert!(array.repeated_probe().execute_scalar(0, &mut ctx).is_err()); + Ok(()) +} + +#[test] +fn lazy_value_validity_only_evaluates_requested_runs() -> VortexResult<()> { + let mut ctx = SESSION.create_execution_ctx(); + let numerators = buffer![1u32, 1].into_array(); + let validity = numerators + .binary(buffer![1u32, 0].into_array(), Operator::Div)? + .binary(numerators, Operator::Eq)?; + let values = PrimitiveArray::new(buffer![42u32, 99], Validity::Array(validity)); + let array = + RunEnd::try_new(buffer![4u32, 8].into_array(), values.into_array(), &mut ctx)?.into_array(); + let expected = array.execute_scalar(0, &mut ctx)?; + let mut probe = array.repeated_probe(); + assert_eq!(probe.execute_scalar(0, &mut ctx)?, expected); + assert_eq!(probe.execute_scalar(3, &mut ctx)?, expected); + assert!(probe.execute_scalar(4, &mut ctx).is_err()); + Ok(()) +}