Skip to main content

vortex_pco/
array.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright the Vortex contributors
3
4use std::cmp;
5use std::fmt::Debug;
6use std::fmt::Display;
7use std::fmt::Formatter;
8use std::hash::Hash;
9use std::hash::Hasher;
10
11use pco::ChunkConfig;
12use pco::PagingSpec;
13use pco::data_types::Number;
14use pco::data_types::NumberType;
15use pco::errors::PcoError;
16use pco::match_number_enum;
17use pco::wrapped::ChunkDecompressor;
18use pco::wrapped::FileCompressor;
19use pco::wrapped::FileDecompressor;
20use prost::Message;
21use vortex_array::Array;
22use vortex_array::ArrayEq;
23use vortex_array::ArrayHash;
24use vortex_array::ArrayId;
25use vortex_array::ArrayParts;
26use vortex_array::ArrayRef;
27use vortex_array::ArrayView;
28use vortex_array::EqMode;
29use vortex_array::ExecutionCtx;
30use vortex_array::ExecutionResult;
31use vortex_array::IntoArray;
32use vortex_array::TypedArrayRef;
33use vortex_array::array_slots;
34use vortex_array::arrays::Primitive;
35use vortex_array::arrays::PrimitiveArray;
36use vortex_array::buffer::BufferHandle;
37use vortex_array::dtype::DType;
38use vortex_array::dtype::PType;
39use vortex_array::dtype::half;
40use vortex_array::scalar::Scalar;
41use vortex_array::serde::ArrayChildren;
42use vortex_array::validity::Validity;
43use vortex_array::vtable::OperationsVTable;
44use vortex_array::vtable::VTable;
45use vortex_array::vtable::ValidityVTable;
46use vortex_array::vtable::child_to_validity;
47use vortex_array::vtable::validity_to_child;
48use vortex_buffer::BufferMut;
49use vortex_buffer::ByteBuffer;
50use vortex_error::VortexError;
51use vortex_error::VortexResult;
52use vortex_error::vortex_bail;
53use vortex_error::vortex_ensure;
54use vortex_error::vortex_err;
55use vortex_session::VortexSession;
56use vortex_session::registry::CachedId;
57
58use crate::PcoChunkInfo;
59use crate::PcoMetadata;
60use crate::PcoPageInfo;
61
62// Overall approach here:
63// Chunk the array into Pco chunks (currently using the default recommended size
64// for good compression), and into finer-grained Pco pages. As we go, write each
65// ChunkMeta as a buffer, followed by each of that chunk's pages as a buffer. We
66// store metadata for each of these "components" (chunk or page). At
67// decompression time, we figure out which components we need to read and only
68// process those. We only compress and decompress valid values.
69
70// Visually, during decompression, we have an interval of pages we're
71// decompressing and a tighter interval of the slice we actually care about.
72// |=============values (all valid elements)==============|
73// |<-n_skipped_values->|----decompressed_values------|
74//                          |----slice_values----|
75//                          ^                    ^
76// |<---slice_value_start-->|<--slice_n_values-->|
77// We then insert these values to the correct position using a primitive array
78// constructor.
79
80const VALUES_PER_CHUNK: usize = pco::DEFAULT_MAX_PAGE_N;
81
82/// A [`Pco`]-encoded Vortex array.
83pub type PcoArray = Array<Pco>;
84
85impl ArrayHash for PcoData {
86    fn array_hash<H: Hasher>(&self, state: &mut H, accuracy: EqMode) {
87        self.unsliced_n_rows.hash(state);
88        self.slice_start.hash(state);
89        self.slice_stop.hash(state);
90        // Hash chunk_metas and pages using pointer-based hashing
91        for chunk_meta in &self.chunk_metas {
92            chunk_meta.array_hash(state, accuracy);
93        }
94        for page in &self.pages {
95            page.array_hash(state, accuracy);
96        }
97    }
98}
99
100impl ArrayEq for PcoData {
101    fn array_eq(&self, other: &Self, accuracy: EqMode) -> bool {
102        if self.unsliced_n_rows != other.unsliced_n_rows
103            || self.slice_start != other.slice_start
104            || self.slice_stop != other.slice_stop
105            || self.chunk_metas.len() != other.chunk_metas.len()
106            || self.pages.len() != other.pages.len()
107        {
108            return false;
109        }
110        for (a, b) in self.chunk_metas.iter().zip(&other.chunk_metas) {
111            if !a.array_eq(b, accuracy) {
112                return false;
113            }
114        }
115        for (a, b) in self.pages.iter().zip(&other.pages) {
116            if !a.array_eq(b, accuracy) {
117                return false;
118            }
119        }
120        true
121    }
122}
123
124impl VTable for Pco {
125    type TypedArrayData = PcoData;
126
127    type OperationsVTable = Self;
128    type ValidityVTable = Self;
129
130    fn id(&self) -> ArrayId {
131        static ID: CachedId = CachedId::new("vortex.pco");
132        *ID
133    }
134
135    fn validate(
136        &self,
137        data: &PcoData,
138        dtype: &DType,
139        len: usize,
140        slots: &[Option<ArrayRef>],
141    ) -> VortexResult<()> {
142        let validity = child_to_validity(
143            PcoSlotsView::from_slots(slots).validity,
144            dtype.nullability(),
145        );
146        data.validate(dtype, len, &validity)
147    }
148
149    fn nbuffers(array: ArrayView<'_, Self>) -> usize {
150        array.chunk_metas.len() + array.pages.len()
151    }
152
153    fn buffer(array: ArrayView<'_, Self>, idx: usize) -> BufferHandle {
154        if idx < array.chunk_metas.len() {
155            BufferHandle::new_host(array.chunk_metas[idx].clone())
156        } else {
157            let page_idx = idx - array.chunk_metas.len();
158            BufferHandle::new_host(array.pages[page_idx].clone())
159        }
160    }
161
162    fn buffer_name(array: ArrayView<'_, Self>, idx: usize) -> Option<String> {
163        if idx < array.chunk_metas.len() {
164            Some(format!("chunk_meta_{idx}"))
165        } else {
166            Some(format!("page_{}", idx - array.chunk_metas.len()))
167        }
168    }
169
170    fn with_buffers(
171        &self,
172        array: ArrayView<'_, Self>,
173        buffers: &[BufferHandle],
174    ) -> VortexResult<ArrayParts<Self>> {
175        let mut data = array.data().clone();
176        let chunk_metas_len = data.metadata.chunks.len();
177        vortex_ensure!(buffers.len() >= chunk_metas_len);
178        data.chunk_metas = buffers[..chunk_metas_len]
179            .iter()
180            .map(|buffer| buffer.clone().try_to_host_sync())
181            .collect::<VortexResult<Vec<_>>>()?;
182        data.pages = buffers[chunk_metas_len..]
183            .iter()
184            .map(|buffer| buffer.clone().try_to_host_sync())
185            .collect::<VortexResult<Vec<_>>>()?;
186        Ok(
187            ArrayParts::new(self.clone(), array.dtype().clone(), array.len(), data)
188                .with_slots(array.slots().iter().cloned().collect()),
189        )
190    }
191
192    fn serialize(
193        array: ArrayView<'_, Self>,
194        _session: &VortexSession,
195    ) -> VortexResult<Option<Vec<u8>>> {
196        Ok(Some(array.metadata.clone().encode_to_vec()))
197    }
198
199    fn deserialize(
200        &self,
201        dtype: &DType,
202        len: usize,
203        metadata: &[u8],
204        buffers: &[BufferHandle],
205        children: &dyn ArrayChildren,
206        _session: &VortexSession,
207    ) -> VortexResult<ArrayParts<Self>> {
208        let metadata = PcoMetadata::decode(metadata)?;
209        let validity = if children.is_empty() {
210            Validity::from(dtype.nullability())
211        } else if children.len() == 1 {
212            let validity = children.get(0, &Validity::DTYPE, len)?;
213            Validity::Array(validity)
214        } else {
215            vortex_bail!("PcoArray expected 0 or 1 child, got {}", children.len());
216        };
217
218        vortex_ensure!(buffers.len() >= metadata.chunks.len());
219        let chunk_metas = buffers[..metadata.chunks.len()]
220            .iter()
221            .map(|b| b.clone().try_to_host_sync())
222            .collect::<VortexResult<Vec<_>>>()?;
223        let pages = buffers[metadata.chunks.len()..]
224            .iter()
225            .map(|b| b.clone().try_to_host_sync())
226            .collect::<VortexResult<Vec<_>>>()?;
227
228        let expected_n_pages = metadata
229            .chunks
230            .iter()
231            .map(|info| info.pages.len())
232            .sum::<usize>();
233        vortex_ensure!(pages.len() == expected_n_pages);
234
235        let slots = PcoSlots {
236            validity: validity_to_child(&validity, len),
237        }
238        .into_slots();
239        // SAFETY: `Array::try_from_parts`, which consumes these parts, validates the data before
240        // publishing the array.
241        let data =
242            unsafe { PcoData::new_unchecked(chunk_metas, pages, dtype.as_ptype(), metadata, len) };
243        Ok(ArrayParts::new(self.clone(), dtype.clone(), len, data).with_slots(slots))
244    }
245
246    fn slot_name(_array: ArrayView<'_, Self>, idx: usize) -> String {
247        PcoSlots::NAMES[idx].to_string()
248    }
249
250    fn execute(array: Array<Self>, ctx: &mut ExecutionCtx) -> VortexResult<ExecutionResult> {
251        let unsliced_validity = array.unsliced_validity();
252        Ok(ExecutionResult::done(
253            array
254                .data()
255                .decompress(&unsliced_validity, ctx)?
256                .into_array(),
257        ))
258    }
259
260    fn reduce_parent(
261        array: ArrayView<'_, Self>,
262        parent: &ArrayRef,
263        child_idx: usize,
264    ) -> VortexResult<Option<ArrayRef>> {
265        crate::rules::RULES.evaluate(array, parent, child_idx)
266    }
267}
268
269pub(crate) fn number_type_from_dtype(dtype: &DType) -> NumberType {
270    number_type_from_ptype(dtype.as_ptype())
271}
272
273pub(crate) fn number_type_from_ptype(ptype: PType) -> NumberType {
274    match ptype {
275        PType::F16 => NumberType::F16,
276        PType::F32 => NumberType::F32,
277        PType::F64 => NumberType::F64,
278        PType::I16 => NumberType::I16,
279        PType::I32 => NumberType::I32,
280        PType::I64 => NumberType::I64,
281        PType::U16 => NumberType::U16,
282        PType::U32 => NumberType::U32,
283        PType::U64 => NumberType::U64,
284        _ => unreachable!("PType not supported by Pco: {:?}", ptype),
285    }
286}
287
288fn collect_valid(
289    parray: ArrayView<'_, Primitive>,
290    ctx: &mut ExecutionCtx,
291) -> VortexResult<PrimitiveArray> {
292    let mask = parray
293        .array()
294        .validity()?
295        .execute_mask(parray.array().len(), ctx)?;
296    let result = parray
297        .array()
298        .filter(mask)?
299        .execute::<PrimitiveArray>(ctx)?;
300    Ok(result)
301}
302
303pub(crate) fn vortex_err_from_pco(err: PcoError) -> VortexError {
304    use pco::errors::ErrorKind::*;
305    match err.kind {
306        Io(io_kind) => VortexError::from(std::io::Error::new(io_kind, err.message)),
307        InvalidArgument => vortex_err!(InvalidArgument: "{}", err.message),
308        other => vortex_err!("Pco {:?} error: {}", other, err.message),
309    }
310}
311
312#[derive(Clone, Debug)]
313/// Pco array encoding marker.
314pub struct Pco;
315
316impl Pco {
317    /// Constructs a Pco array after validating its data and validity invariants.
318    pub fn try_new(dtype: DType, data: PcoData, validity: Validity) -> VortexResult<PcoArray> {
319        let len = data.len();
320        data.validate(&dtype, len, &validity)?;
321        // SAFETY: `validate` checked the dtype, length, validity, metadata, and buffer invariants.
322        Ok(unsafe { Self::new_unchecked(dtype, data, validity) })
323    }
324
325    /// Constructs a Pco array without validating its data and validity invariants.
326    ///
327    /// # Safety
328    ///
329    /// The caller must ensure that [`PcoData::validate`] would succeed for `dtype`, `data.len()`,
330    /// and `validity`.
331    pub unsafe fn new_unchecked(dtype: DType, data: PcoData, validity: Validity) -> PcoArray {
332        let len = data.len();
333        let slots = PcoSlots {
334            validity: validity_to_child(&validity, data.unsliced_n_rows()),
335        }
336        .into_slots();
337        unsafe {
338            Array::from_parts_unchecked(ArrayParts::new(Pco, dtype, len, data).with_slots(slots))
339        }
340    }
341
342    /// Compress a primitive array using pcodec.
343    pub fn from_primitive(
344        parray: ArrayView<'_, Primitive>,
345        level: usize,
346        values_per_page: usize,
347        ctx: &mut ExecutionCtx,
348    ) -> VortexResult<PcoArray> {
349        let dtype = parray.dtype().clone();
350        let validity = parray.validity()?;
351        let data = PcoData::from_primitive(parray, level, values_per_page, ctx)?;
352        Self::try_new(dtype, data, validity)
353    }
354}
355
356#[array_slots(Pco)]
357pub struct PcoSlots {
358    /// The validity bitmap indicating which elements are non-null.
359    #[slot(0)]
360    pub validity: Option<ArrayRef>,
361}
362
363/// Additional typed accessors for Pco arrays.
364pub trait PcoArrayExt: PcoArraySlotsExt {
365    /// Reconstruct the unsliced [`Validity`] from the validity slot.
366    fn unsliced_validity(&self) -> Validity {
367        child_to_validity(
368            self.as_ref().slots()[PcoSlots::VALIDITY].as_ref(),
369            self.as_ref().dtype().nullability(),
370        )
371    }
372}
373impl<T: TypedArrayRef<Pco>> PcoArrayExt for T {}
374
375#[derive(Clone, Debug)]
376/// Encoding-specific data for a [`PcoArray`].
377pub struct PcoData {
378    pub(crate) chunk_metas: Vec<ByteBuffer>,
379    pub(crate) pages: Vec<ByteBuffer>,
380    pub(crate) metadata: PcoMetadata,
381    ptype: PType,
382    unsliced_n_rows: usize,
383    slice_start: usize,
384    slice_stop: usize,
385}
386
387impl Display for PcoData {
388    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
389        write!(
390            f,
391            "ptype: {}, nrows: {}, slice: {}..{}",
392            self.ptype, self.unsliced_n_rows, self.slice_start, self.slice_stop
393        )
394    }
395}
396
397impl PcoData {
398    /// Validate dtype, validity, slice, and Pco component invariants.
399    pub fn validate(&self, dtype: &DType, len: usize, validity: &Validity) -> VortexResult<()> {
400        let _ = number_type_from_ptype(self.ptype);
401        vortex_ensure!(
402            dtype.as_ptype() == self.ptype,
403            "expected ptype {}, got {}",
404            self.ptype,
405            dtype.as_ptype()
406        );
407        vortex_ensure!(
408            dtype.nullability() == validity.nullability(),
409            "expected nullability {}, got {}",
410            validity.nullability(),
411            dtype.nullability()
412        );
413        vortex_ensure!(
414            self.slice_start <= self.slice_stop && self.slice_stop <= self.unsliced_n_rows,
415            "invalid slice range {}..{} for {} rows",
416            self.slice_start,
417            self.slice_stop,
418            self.unsliced_n_rows
419        );
420        vortex_ensure!(
421            self.slice_stop - self.slice_start == len,
422            "expected len {len}, got {}",
423            self.slice_stop - self.slice_start
424        );
425        if let Some(validity_len) = validity.maybe_len() {
426            vortex_ensure!(
427                validity_len == self.unsliced_n_rows,
428                "expected validity len {}, got {}",
429                self.unsliced_n_rows,
430                validity_len
431            );
432        }
433        vortex_ensure!(
434            self.chunk_metas.len() == self.metadata.chunks.len(),
435            "expected {} chunk metas, got {}",
436            self.metadata.chunks.len(),
437            self.chunk_metas.len()
438        );
439        vortex_ensure!(
440            self.pages.len()
441                == self
442                    .metadata
443                    .chunks
444                    .iter()
445                    .map(|chunk| chunk.pages.len())
446                    .sum::<usize>(),
447            "page count does not match metadata"
448        );
449
450        let mut n_values = 0usize;
451        for (chunk_idx, chunk) in self.metadata.chunks.iter().enumerate() {
452            let mut chunk_n_values = 0usize;
453            for page in &chunk.pages {
454                let page_n_values = page.n_values as usize;
455                vortex_ensure!(
456                    page_n_values != 0,
457                    "Pco chunk {chunk_idx} contains an empty page"
458                );
459                chunk_n_values = chunk_n_values.checked_add(page_n_values).ok_or_else(|| {
460                    vortex_err!("Pco chunk {chunk_idx} value count overflows usize")
461                })?;
462            }
463            vortex_ensure!(
464                chunk_n_values <= VALUES_PER_CHUNK,
465                "Pco chunk {chunk_idx} contains {chunk_n_values} values, exceeding the maximum of {VALUES_PER_CHUNK}"
466            );
467            n_values = n_values
468                .checked_add(chunk_n_values)
469                .ok_or_else(|| vortex_err!("Pco value count overflows usize"))?;
470        }
471        vortex_ensure!(
472            n_values <= self.unsliced_n_rows,
473            "Pco contains {n_values} values for only {} rows",
474            self.unsliced_n_rows
475        );
476        if validity.definitely_no_nulls() {
477            vortex_ensure!(
478                n_values == self.unsliced_n_rows,
479                "Pco contains {n_values} values for {} non-null rows",
480                self.unsliced_n_rows
481            );
482        } else if validity.definitely_all_null() {
483            vortex_ensure!(n_values == 0, "Pco contains values for an all-null array");
484        }
485        Ok(())
486    }
487
488    /// Constructs unsliced Pco data without validating its metadata and buffers.
489    ///
490    /// # Safety
491    ///
492    /// The returned data must pass [`Self::validate`] before it is decompressed or published as
493    /// part of an array.
494    pub unsafe fn new_unchecked(
495        chunk_metas: Vec<ByteBuffer>,
496        pages: Vec<ByteBuffer>,
497        ptype: PType,
498        metadata: PcoMetadata,
499        len: usize,
500    ) -> Self {
501        Self {
502            chunk_metas,
503            pages,
504            metadata,
505            ptype,
506            unsliced_n_rows: len,
507            slice_start: 0,
508            slice_stop: len,
509        }
510    }
511
512    /// Compress a primitive array into Pco data.
513    pub fn from_primitive(
514        parray: ArrayView<'_, Primitive>,
515        level: usize,
516        values_per_page: usize,
517        ctx: &mut ExecutionCtx,
518    ) -> VortexResult<Self> {
519        Self::from_primitive_with_values_per_chunk(
520            parray,
521            level,
522            VALUES_PER_CHUNK,
523            values_per_page,
524            ctx,
525        )
526    }
527
528    pub(crate) fn from_primitive_with_values_per_chunk(
529        parray: ArrayView<'_, Primitive>,
530        level: usize,
531        values_per_chunk: usize,
532        values_per_page: usize,
533        ctx: &mut ExecutionCtx,
534    ) -> VortexResult<Self> {
535        let number_type = number_type_from_dtype(parray.dtype());
536        let values_per_page = if values_per_page == 0 {
537            values_per_chunk
538        } else {
539            values_per_page
540        };
541
542        // perhaps one day we can make this more configurable
543        let chunk_config = ChunkConfig::default()
544            .with_compression_level(level)
545            .with_paging_spec(PagingSpec::EqualPagesUpTo(values_per_page));
546
547        let values = collect_valid(parray, ctx)?;
548        let n_values = values.len();
549
550        let fc = FileCompressor::default();
551        let mut header = vec![];
552        fc.write_header(&mut header).map_err(vortex_err_from_pco)?;
553
554        let mut chunk_meta_buffers = vec![]; // the Pco component
555        let mut chunk_infos = vec![]; // the Vortex metadata
556        let mut page_buffers = vec![];
557        for chunk_start in (0..n_values).step_by(values_per_chunk) {
558            let chunk_end = cmp::min(n_values, chunk_start + values_per_chunk);
559            let mut cc = match_number_enum!(
560                number_type,
561                NumberType<T> => {
562                    let values = values.to_buffer::<T>();
563                    let chunk = &values.as_slice()[chunk_start..chunk_end];
564                    fc
565                        .chunk_compressor(chunk, &chunk_config)
566                        .map_err(vortex_err_from_pco)?
567                }
568            );
569
570            let mut chunk_meta_buffer = Vec::with_capacity(cc.meta_size_hint());
571            cc.write_meta(&mut chunk_meta_buffer)
572                .map_err(vortex_err_from_pco)?;
573            chunk_meta_buffers.push(ByteBuffer::from(chunk_meta_buffer));
574
575            let mut page_infos = vec![];
576            for (page_idx, page_n_values) in cc.n_per_page().into_iter().enumerate() {
577                let mut page = Vec::with_capacity(cc.page_size_hint(page_idx));
578                cc.write_page(page_idx, &mut page)
579                    .map_err(vortex_err_from_pco)?;
580                page_buffers.push(ByteBuffer::from(page));
581                page_infos.push(PcoPageInfo {
582                    n_values: u32::try_from(page_n_values)?,
583                });
584            }
585            chunk_infos.push(PcoChunkInfo { pages: page_infos })
586        }
587
588        let metadata = PcoMetadata {
589            header,
590            chunks: chunk_infos,
591        };
592        // SAFETY: The compressor produced matching chunk metadata and pages from `parray`, with
593        // the same ptype, logical length, and valid-value count.
594        Ok(unsafe {
595            PcoData::new_unchecked(
596                chunk_meta_buffers,
597                page_buffers,
598                parray.dtype().as_ptype(),
599                metadata,
600                parray.len(),
601            )
602        })
603    }
604
605    /// Downcast and compress an array into Pco data.
606    ///
607    /// # Errors
608    ///
609    /// Returns an error if the input is not a primitive array or compression fails.
610    pub fn from_array(
611        array: ArrayRef,
612        level: usize,
613        nums_per_page: usize,
614        ctx: &mut ExecutionCtx,
615    ) -> VortexResult<Self> {
616        let parray = array.try_downcast::<Primitive>().map_err(|a| {
617            vortex_err!(
618                "Pco can only encode primitive arrays, got {}",
619                a.encoding_id()
620            )
621        })?;
622        Self::from_primitive(parray.as_view(), level, nums_per_page, ctx)
623    }
624
625    /// Decompress this Pco data into a primitive array.
626    pub(crate) fn decompress(
627        &self,
628        unsliced_validity: &Validity,
629        ctx: &mut ExecutionCtx,
630    ) -> VortexResult<PrimitiveArray> {
631        // To start, we figure out which chunks and pages we need to decompress, and with
632        // what value offset into the first such page.
633        let number_type = number_type_from_ptype(self.ptype);
634        let values_byte_buffer = match_number_enum!(
635            number_type,
636            NumberType<T> => {
637              self.decompress_values_typed::<T>(unsliced_validity, ctx)?
638            }
639        );
640
641        Ok(PrimitiveArray::from_values_byte_buffer(
642            values_byte_buffer,
643            self.ptype,
644            unsliced_validity.slice(self.slice_start..self.slice_stop)?,
645            self.slice_stop - self.slice_start,
646            ctx,
647        ))
648    }
649
650    fn decompress_values_typed<T: Number>(
651        &self,
652        unsliced_validity: &Validity,
653        ctx: &mut ExecutionCtx,
654    ) -> VortexResult<ByteBuffer> {
655        // To start, we figure out what range of values we need to decompress.
656        let slice_value_indices = unsliced_validity
657            .execute_mask(self.unsliced_n_rows, ctx)?
658            .valid_counts_for_indices(&[self.slice_start, self.slice_stop]);
659        let slice_value_start = slice_value_indices[0];
660        let slice_value_stop = slice_value_indices[1];
661        let slice_n_values = slice_value_stop - slice_value_start;
662
663        // Then we decompress those pages into a buffer. Note that these values
664        // may exceed the bounds of the slice, so we need to slice later.
665        let (fd, _) =
666            FileDecompressor::new(self.metadata.header.as_slice()).map_err(vortex_err_from_pco)?;
667        let mut decompressed_values =
668            BufferMut::<T>::with_capacity(slice_n_values.min(VALUES_PER_CHUNK));
669        // Validation bounds every page-count prefix and selected-page subtotal by
670        // `unsliced_n_rows`, so the arithmetic below cannot overflow.
671        let mut page_idx = 0;
672        let mut page_value_start = 0usize;
673        let mut n_skipped_values = 0;
674        for (chunk_info, chunk_meta) in self.metadata.chunks.iter().zip(&self.chunk_metas) {
675            // lazily initialize chunk decompressor
676            let mut chunk_decompressor: Option<ChunkDecompressor<T>> = None;
677            for page_info in &chunk_info.pages {
678                let page_n_values = page_info.n_values as usize;
679                let page_value_stop = page_value_start + page_n_values;
680
681                if page_value_start >= slice_value_stop {
682                    break;
683                }
684
685                if page_value_stop > slice_value_start {
686                    // we need this page
687                    let old_len = decompressed_values.len();
688                    let new_len = old_len + page_n_values;
689                    decompressed_values.reserve(page_n_values);
690                    unsafe {
691                        decompressed_values.set_len(new_len);
692                    }
693                    let page: &[u8] = self
694                        .pages
695                        .get(page_idx)
696                        .ok_or_else(|| vortex_err!("Missing Pco page {page_idx}"))?
697                        .as_ref();
698
699                    let mut cd = match chunk_decompressor.take() {
700                        Some(d) => d,
701                        None => {
702                            let (new_cd, _) = fd
703                                .chunk_decompressor(chunk_meta.as_ref())
704                                .map_err(vortex_err_from_pco)?;
705                            new_cd
706                        }
707                    };
708
709                    let mut pd = cd
710                        .page_decompressor(page, page_n_values)
711                        .map_err(vortex_err_from_pco)?;
712                    pd.read(&mut decompressed_values[old_len..new_len])
713                        .map_err(vortex_err_from_pco)?;
714
715                    chunk_decompressor = Some(cd);
716                } else {
717                    n_skipped_values += page_n_values;
718                }
719
720                page_value_start = page_value_stop;
721                page_idx += 1;
722            }
723        }
724
725        // Slice only the values requested.
726        // Skipped pages end before `slice_value_start`, and the resulting stop is at most
727        // `slice_value_stop`, which is bounded by `unsliced_n_rows`.
728        let value_offset = slice_value_start - n_skipped_values;
729        let value_stop = value_offset + slice_n_values;
730        vortex_ensure!(
731            value_stop <= decompressed_values.len(),
732            "Pco contains {} decompressed values, but the requested range ends at {value_stop}",
733            decompressed_values.len()
734        );
735        Ok(decompressed_values
736            .freeze()
737            .slice(value_offset..value_stop)
738            .into_byte_buffer())
739    }
740
741    pub(crate) fn _slice(&self, start: usize, stop: usize) -> Self {
742        PcoData {
743            slice_start: self.slice_start + start,
744            slice_stop: self.slice_start + stop,
745            ..self.clone()
746        }
747    }
748
749    /// Returns the number of elements in the array.
750    pub fn len(&self) -> usize {
751        self.slice_stop - self.slice_start
752    }
753
754    /// Returns `true` if the array contains no elements.
755    pub fn is_empty(&self) -> bool {
756        self.slice_stop == self.slice_start
757    }
758
759    pub(crate) fn slice_start(&self) -> usize {
760        self.slice_start
761    }
762
763    pub(crate) fn slice_stop(&self) -> usize {
764        self.slice_stop
765    }
766
767    pub(crate) fn unsliced_n_rows(&self) -> usize {
768        self.unsliced_n_rows
769    }
770}
771
772impl ValidityVTable<Pco> for Pco {
773    fn validity(array: ArrayView<'_, Pco>) -> VortexResult<Validity> {
774        array
775            .unsliced_validity()
776            .slice(array.slice_start()..array.slice_stop())
777    }
778}
779
780impl OperationsVTable<Pco> for Pco {
781    fn scalar_at(
782        array: ArrayView<'_, Pco>,
783        index: usize,
784        ctx: &mut ExecutionCtx,
785    ) -> VortexResult<Scalar> {
786        let unsliced_validity = array.unsliced_validity();
787        array
788            ._slice(index, index + 1)
789            .decompress(&unsliced_validity, ctx)?
790            .into_array()
791            .execute_scalar(0, ctx)
792    }
793}
794
795#[cfg(test)]
796mod tests {
797    use vortex_array::IntoArray;
798    use vortex_array::VortexSessionExecute;
799    use vortex_array::array_session;
800    use vortex_array::arrays::PrimitiveArray;
801    use vortex_array::assert_arrays_eq;
802    use vortex_array::validity::Validity;
803    use vortex_buffer::buffer;
804    use vortex_error::VortexResult;
805
806    use super::VALUES_PER_CHUNK;
807    use crate::Pco;
808
809    #[test]
810    fn test_slice_nullable() {
811        let mut ctx = array_session().create_execution_ctx();
812        // Create a nullable array with some nulls
813        let values = PrimitiveArray::new(
814            buffer![10u32, 20, 30, 40, 50, 60],
815            Validity::from_iter([false, true, true, true, true, false]),
816        );
817        let pco = Pco::from_primitive(values.as_view(), 0, 128, &mut ctx).unwrap();
818        assert_arrays_eq!(
819            pco,
820            PrimitiveArray::from_option_iter([
821                None,
822                Some(20u32),
823                Some(30),
824                Some(40),
825                Some(50),
826                None
827            ]),
828            &mut ctx
829        );
830
831        // Slice to get only the non-null values in the middle
832        let sliced = pco.slice(1..5).unwrap();
833        let expected =
834            PrimitiveArray::from_option_iter([Some(20u32), Some(30), Some(40), Some(50)])
835                .into_array();
836        assert_arrays_eq!(sliced, expected, &mut ctx);
837    }
838
839    #[test]
840    fn test_decompress_bounds_initial_allocation() -> VortexResult<()> {
841        let mut ctx = array_session().create_execution_ctx();
842        let values = PrimitiveArray::from_iter([42u32]);
843        let pco = Pco::from_primitive(values.as_view(), 0, 128, &mut ctx)?;
844        let mut data = pco.data().clone();
845
846        // Simulate a corrupt logical length whose validity claims far more values than the
847        // encoded pages contain. This must not be used directly as an allocation size.
848        data.unsliced_n_rows = usize::MAX;
849        data.slice_stop = usize::MAX;
850
851        let result = data.decompress_values_typed::<u32>(&Validity::NonNullable, &mut ctx);
852        assert!(result.is_err());
853        Ok(())
854    }
855
856    #[test]
857    fn test_validate_rejects_oversized_chunk() -> VortexResult<()> {
858        let mut ctx = array_session().create_execution_ctx();
859        let values = PrimitiveArray::from_iter([42u32]);
860        let pco = Pco::from_primitive(values.as_view(), 0, 128, &mut ctx)?;
861        let mut data = pco.data().clone();
862        data.metadata.chunks[0].pages[0].n_values = u32::try_from(VALUES_PER_CHUNK + 1)?;
863
864        assert!(
865            data.validate(pco.dtype(), pco.len(), &Validity::NonNullable)
866                .is_err()
867        );
868        Ok(())
869    }
870}