1use 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
62const VALUES_PER_CHUNK: usize = pco::DEFAULT_MAX_PAGE_N;
81
82pub 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 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 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)]
313pub struct Pco;
315
316impl Pco {
317 pub fn try_new(dtype: DType, data: PcoData, validity: Validity) -> VortexResult<PcoArray> {
319 let len = data.len();
320 data.validate(&dtype, len, &validity)?;
321 Ok(unsafe { Self::new_unchecked(dtype, data, validity) })
323 }
324
325 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 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 #[slot(0)]
360 pub validity: Option<ArrayRef>,
361}
362
363pub trait PcoArrayExt: PcoArraySlotsExt {
365 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)]
376pub 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 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 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 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 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![]; let mut chunk_infos = vec![]; 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 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 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 pub(crate) fn decompress(
627 &self,
628 unsliced_validity: &Validity,
629 ctx: &mut ExecutionCtx,
630 ) -> VortexResult<PrimitiveArray> {
631 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 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 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 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 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 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 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 pub fn len(&self) -> usize {
751 self.slice_stop - self.slice_start
752 }
753
754 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 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 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 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}