1use std::{
110 iter::{Copied, Zip},
111 ops::Range,
112 sync::Arc,
113};
114
115use arrow_array::OffsetSizeTrait;
116use arrow_buffer::{
117 ArrowNativeType, BooleanBuffer, BooleanBufferBuilder, NullBuffer, OffsetBuffer, ScalarBuffer,
118};
119use lance_core::{Error, Result, utils::bit::log_2_ceil};
120
121use crate::buffer::LanceBuffer;
122
123pub type LevelBuffer = Vec<u16>;
124
125#[derive(Debug, Clone, PartialEq, Eq)]
127pub(crate) struct MiniBlockRepDefSplit {
128 pub(crate) row_start: u64,
130 pub(crate) num_rows: u64,
132 pub(crate) level_range: Range<usize>,
134 pub(crate) value_start: u64,
136 pub(crate) num_values: u64,
138}
139
140#[derive(Debug, Clone, PartialEq, Eq)]
142pub(crate) enum MiniBlockRepDefBudget {
143 WithinBudget,
145 RequiresPageSplit(Vec<MiniBlockRepDefSplit>),
147 SingleRowOverBudget(u64),
149}
150
151const SPECIAL_THRESHOLD: u16 = u16::MAX / 2;
162
163#[derive(Clone, Debug)]
166struct OffsetDesc {
167 offsets: Arc<[i64]>,
168 validity: Option<BooleanBuffer>,
169 has_empty_lists: bool,
170 num_values: usize,
171 num_specials: usize,
172}
173
174#[derive(Clone, Debug)]
177struct ValidityDesc {
178 validity: Option<BooleanBuffer>,
179 num_values: usize,
180}
181
182#[derive(Clone, Debug)]
186struct FslDesc {
187 validity: Option<BooleanBuffer>,
188 dimension: usize,
189 num_values: usize,
190}
191
192#[derive(Clone, Debug)]
196enum RawRepDef {
197 Offsets(OffsetDesc),
198 Validity(ValidityDesc),
199 Fsl(FslDesc),
200}
201
202impl RawRepDef {
203 fn has_nulls(&self) -> bool {
205 match self {
206 Self::Offsets(OffsetDesc { validity, .. }) => validity.is_some(),
207 Self::Validity(ValidityDesc { validity, .. }) => validity.is_some(),
208 Self::Fsl(FslDesc { validity, .. }) => validity.is_some(),
209 }
210 }
211
212 fn num_values(&self) -> usize {
214 match self {
215 Self::Offsets(OffsetDesc { num_values, .. }) => *num_values,
216 Self::Validity(ValidityDesc { num_values, .. }) => *num_values,
217 Self::Fsl(FslDesc { num_values, .. }) => *num_values,
218 }
219 }
220
221 fn num_specials(&self) -> usize {
223 match self {
224 Self::Offsets(OffsetDesc { num_specials, .. }) => *num_specials,
225 _ => 0,
226 }
227 }
228
229 fn max_def(&self) -> u16 {
231 match self {
232 Self::Offsets(OffsetDesc {
233 has_empty_lists,
234 validity,
235 ..
236 }) => {
237 let mut max_def = 0;
238 if *has_empty_lists {
239 max_def += 1;
240 }
241 if validity.is_some() {
242 max_def += 1;
243 }
244 max_def
245 }
246 Self::Validity(ValidityDesc { validity: None, .. }) => 0,
247 Self::Validity(ValidityDesc { .. }) => 1,
248 Self::Fsl(FslDesc { validity: None, .. }) => 0,
249 Self::Fsl(FslDesc { .. }) => 1,
250 }
251 }
252
253 fn max_rep(&self) -> u16 {
255 match self {
256 Self::Offsets(_) => 1,
257 _ => 0,
258 }
259 }
260}
261
262#[derive(Debug)]
265pub struct SerializedRepDefs {
266 pub repetition_levels: Option<Arc<[u16]>>,
270 pub definition_levels: Option<Arc<[u16]>>,
274 pub def_meaning: Vec<DefinitionInterpretation>,
276 pub max_visible_level: Option<u16>,
283 has_fsl: bool,
284}
285
286impl SerializedRepDefs {
287 fn max_visible_level(def_meaning: &[DefinitionInterpretation]) -> Option<u16> {
288 let first_list = def_meaning.iter().position(|level| level.is_list());
289 first_list.map(|first_list| {
290 def_meaning
291 .iter()
292 .map(|level| level.num_def_levels())
293 .take(first_list)
294 .sum::<u16>()
295 })
296 }
297
298 pub fn new(
299 repetition_levels: Option<LevelBuffer>,
300 definition_levels: Option<LevelBuffer>,
301 def_meaning: Vec<DefinitionInterpretation>,
302 ) -> Self {
303 Self::new_with_fixed_size_list_levels(
304 repetition_levels,
305 definition_levels,
306 def_meaning,
307 false,
308 )
309 }
310
311 pub(crate) fn new_with_fixed_size_list_levels(
312 repetition_levels: Option<LevelBuffer>,
313 definition_levels: Option<LevelBuffer>,
314 def_meaning: Vec<DefinitionInterpretation>,
315 has_fsl: bool,
316 ) -> Self {
317 let max_visible_level = Self::max_visible_level(&def_meaning);
318 Self {
319 repetition_levels: repetition_levels.map(Arc::from),
320 definition_levels: definition_levels.map(Arc::from),
321 def_meaning,
322 max_visible_level,
323 has_fsl,
324 }
325 }
326
327 pub fn empty(def_meaning: Vec<DefinitionInterpretation>) -> Self {
329 Self {
330 repetition_levels: None,
331 definition_levels: None,
332 def_meaning,
333 max_visible_level: None,
334 has_fsl: false,
335 }
336 }
337
338 pub fn rep_slicer(&self) -> Option<RepDefSlicer<'_>> {
339 self.repetition_levels
340 .as_ref()
341 .map(|rep| RepDefSlicer::new(self, rep.clone()))
342 }
343
344 pub fn def_slicer(&self) -> Option<RepDefSlicer<'_>> {
345 self.definition_levels
346 .as_ref()
347 .map(|def| RepDefSlicer::new(self, def.clone()))
348 }
349
350 pub(crate) fn has_fixed_size_list_levels(&self) -> bool {
351 self.has_fsl
352 }
353}
354
355#[derive(Debug)]
363pub struct RepDefSlicer<'a> {
364 repdef: &'a SerializedRepDefs,
365 to_slice: LanceBuffer,
366 current: usize,
367}
368
369impl<'a> RepDefSlicer<'a> {
371 fn new(repdef: &'a SerializedRepDefs, levels: Arc<[u16]>) -> Self {
372 Self {
373 repdef,
374 to_slice: LanceBuffer::reinterpret_slice(levels),
375 current: 0,
376 }
377 }
378
379 pub fn num_levels(&self) -> usize {
380 self.to_slice.len() / 2
381 }
382
383 pub fn num_levels_remaining(&self) -> usize {
384 self.num_levels() - self.current
385 }
386
387 pub fn all_levels(&self) -> &LanceBuffer {
388 &self.to_slice
389 }
390
391 pub fn slice_rest(&mut self) -> LanceBuffer {
400 let start = self.current;
401 let remaining = self.num_levels_remaining();
402 self.current = self.num_levels();
403 self.to_slice.slice_with_length(start * 2, remaining * 2)
404 }
405
406 pub fn slice_next(&mut self, num_values: usize) -> LanceBuffer {
408 let start = self.current;
409 let Some(max_visible_level) = self.repdef.max_visible_level else {
410 self.current = start + num_values;
412 return self.to_slice.slice_with_length(start * 2, num_values * 2);
413 };
414 if let Some(def) = self.repdef.definition_levels.as_ref() {
415 let mut def_itr = def[start..].iter();
419 let mut num_taken = 0;
420 let mut num_passed = 0;
421 while num_taken < num_values {
422 let def_level = *def_itr.next().unwrap();
423 if def_level <= max_visible_level {
424 num_taken += 1;
425 }
426 num_passed += 1;
427 }
428 self.current = start + num_passed;
429 self.to_slice.slice_with_length(start * 2, num_passed * 2)
430 } else {
431 self.current = start + num_values;
433 self.to_slice.slice_with_length(start * 2, num_values * 2)
434 }
435 }
436}
437
438#[derive(Debug, Copy, Clone, PartialEq, Eq)]
451pub enum DefinitionInterpretation {
452 AllValidItem,
453 AllValidList,
454 NullableItem,
455 NullableList,
456 EmptyableList,
457 NullableAndEmptyableList,
458}
459
460impl DefinitionInterpretation {
461 pub fn num_def_levels(&self) -> u16 {
463 match self {
464 Self::AllValidItem => 0,
465 Self::AllValidList => 0,
466 Self::NullableItem => 1,
467 Self::NullableList => 1,
468 Self::EmptyableList => 1,
469 Self::NullableAndEmptyableList => 2,
470 }
471 }
472
473 pub fn is_all_valid(&self) -> bool {
475 matches!(
476 self,
477 Self::AllValidItem | Self::AllValidList | Self::EmptyableList
478 )
479 }
480
481 pub fn is_list(&self) -> bool {
483 matches!(
484 self,
485 Self::AllValidList
486 | Self::NullableList
487 | Self::EmptyableList
488 | Self::NullableAndEmptyableList
489 )
490 }
491}
492
493#[derive(Debug)]
505struct SerializerContext {
506 def_meaning: Vec<DefinitionInterpretation>,
508 rep_levels: LevelBuffer,
509 spare_rep: LevelBuffer,
510 def_levels: LevelBuffer,
511 spare_def: LevelBuffer,
512 current_rep: u16,
513 current_def: u16,
514 current_len: usize,
515 current_num_specials: usize,
516 has_fsl: bool,
517}
518
519impl SerializerContext {
520 fn new(len: usize, num_layers: usize, max_rep: u16, max_def: u16) -> Self {
521 let def_meaning = Vec::with_capacity(num_layers);
522 Self {
523 rep_levels: if max_rep > 0 {
524 vec![0; len]
525 } else {
526 LevelBuffer::default()
527 },
528 spare_rep: if max_rep > 0 {
529 vec![0; len]
530 } else {
531 LevelBuffer::default()
532 },
533 def_levels: if max_def > 0 {
534 vec![0; len]
535 } else {
536 LevelBuffer::default()
537 },
538 spare_def: if max_def > 0 {
539 vec![0; len]
540 } else {
541 LevelBuffer::default()
542 },
543 def_meaning,
544 current_rep: max_rep,
545 current_def: max_def,
546 current_len: 0,
547 current_num_specials: 0,
548 has_fsl: false,
549 }
550 }
551
552 fn checkout_def(&mut self, meaning: DefinitionInterpretation) -> u16 {
553 let def = self.current_def;
554 self.current_def -= meaning.num_def_levels();
555 self.def_meaning.push(meaning);
556 def
557 }
558
559 fn record_offsets(&mut self, offset_desc: &OffsetDesc) {
560 let rep_level = self.current_rep;
561 let (null_list_level, empty_list_level) =
562 match (offset_desc.validity.is_some(), offset_desc.has_empty_lists) {
563 (true, true) => {
564 let level =
565 self.checkout_def(DefinitionInterpretation::NullableAndEmptyableList);
566 (level - 1, level)
567 }
568 (true, false) => (self.checkout_def(DefinitionInterpretation::NullableList), 0),
569 (false, true) => (
570 0,
571 self.checkout_def(DefinitionInterpretation::EmptyableList),
572 ),
573 (false, false) => {
574 self.checkout_def(DefinitionInterpretation::AllValidList);
575 (0, 0)
576 }
577 };
578 self.current_rep -= 1;
579
580 if let Some(validity) = &offset_desc.validity {
581 self.do_record_validity(validity, null_list_level);
582 }
583
584 let mut new_len = 0;
589 let expected_len = offset_desc.num_values + self.current_num_specials;
590 if expected_len == 0 {
591 self.current_len = 0;
593 return;
594 }
595 assert!(self.rep_levels.len() >= expected_len - 1);
596 if self.def_levels.is_empty() {
597 let mut write_itr = self.spare_rep.iter_mut();
598 let mut read_iter = self.rep_levels.iter().copied();
599 for w in offset_desc.offsets.windows(2) {
600 let len = w[1] - w[0];
601 assert!(len > 0);
603 let rep = read_iter.next().unwrap();
604 let list_level = if rep == 0 { rep_level } else { rep };
605 *write_itr.next().unwrap() = list_level;
606
607 for _ in 1..len {
608 *write_itr.next().unwrap() = 0;
609 }
610 new_len += len as usize;
611 }
612 std::mem::swap(&mut self.rep_levels, &mut self.spare_rep);
613 } else {
614 assert!(self.def_levels.len() >= expected_len - 1);
615 let mut def_write_itr = self.spare_def.iter_mut();
616 let mut rep_write_itr = self.spare_rep.iter_mut();
617 let mut rep_read_itr = self.rep_levels.iter().copied();
618 let mut def_read_itr = self.def_levels.iter().copied();
619 let specials_to_pass = self.current_num_specials;
620 let mut specials_passed = 0;
621
622 for w in offset_desc.offsets.windows(2) {
623 let mut def = def_read_itr.next().unwrap();
624 while def > SPECIAL_THRESHOLD {
626 *def_write_itr.next().unwrap() = def;
627 *rep_write_itr.next().unwrap() = rep_read_itr.next().unwrap();
628 def = def_read_itr.next().unwrap();
629 new_len += 1;
630 specials_passed += 1;
631 }
632
633 let len = w[1] - w[0];
634 let rep = rep_read_itr.next().unwrap();
635
636 let list_level = if rep == 0 { rep_level } else { rep };
640
641 if def == 0 && len > 0 {
642 *def_write_itr.next().unwrap() = 0;
644 *rep_write_itr.next().unwrap() = list_level;
645
646 for _ in 1..len {
647 *def_write_itr.next().unwrap() = 0;
648 *rep_write_itr.next().unwrap() = 0;
649 }
650
651 new_len += len as usize;
652 } else if def == 0 {
653 *def_write_itr.next().unwrap() = empty_list_level + SPECIAL_THRESHOLD;
655 *rep_write_itr.next().unwrap() = list_level;
656 new_len += 1;
657 } else {
658 *def_write_itr.next().unwrap() = def + SPECIAL_THRESHOLD;
661 *rep_write_itr.next().unwrap() = list_level;
662 new_len += 1;
663 }
664 }
665
666 while specials_passed < specials_to_pass {
668 *def_write_itr.next().unwrap() = def_read_itr.next().unwrap();
669 *rep_write_itr.next().unwrap() = rep_read_itr.next().unwrap();
670 new_len += 1;
671 specials_passed += 1;
672 }
673 std::mem::swap(&mut self.def_levels, &mut self.spare_def);
674 std::mem::swap(&mut self.rep_levels, &mut self.spare_rep);
675 }
676
677 self.current_len = new_len;
678 self.current_num_specials += offset_desc.num_specials;
679 }
680
681 fn do_record_validity(&mut self, validity: &BooleanBuffer, null_level: u16) {
682 assert!(self.def_levels.len() >= validity.len() + self.current_num_specials);
683 debug_assert!(
684 self.current_len == 0 || self.current_len == validity.len() + self.current_num_specials
685 );
686 self.current_len = validity.len();
687
688 let mut def_read_itr = self.def_levels.iter().copied();
689 let mut def_write_itr = self.spare_def.iter_mut();
690
691 let specials_to_pass = self.current_num_specials;
692 let mut specials_passed = 0;
693
694 for incoming_validity in validity.iter() {
695 let mut def = def_read_itr.next().unwrap();
696 while def > SPECIAL_THRESHOLD {
697 *def_write_itr.next().unwrap() = def;
698 def = def_read_itr.next().unwrap();
699 specials_passed += 1;
700 }
701 if def == 0 && !incoming_validity {
702 *def_write_itr.next().unwrap() = null_level;
703 } else {
704 *def_write_itr.next().unwrap() = def;
705 }
706 }
707
708 while specials_passed < specials_to_pass {
709 *def_write_itr.next().unwrap() = def_read_itr.next().unwrap();
710 specials_passed += 1;
711 }
712
713 std::mem::swap(&mut self.def_levels, &mut self.spare_def);
714 }
715
716 fn multiply_levels(&mut self, multiplier: usize) {
717 let old_len = self.current_len;
718 self.current_len =
720 (self.current_len - self.current_num_specials) * multiplier + self.current_num_specials;
721
722 if self.rep_levels.is_empty() && self.def_levels.is_empty() {
723 return;
725 } else if self.rep_levels.is_empty() {
726 assert!(self.def_levels.len() >= self.current_len);
727 let mut def_read_itr = self.def_levels.iter().copied();
729 let mut def_write_itr = self.spare_def.iter_mut();
730 for _ in 0..old_len {
731 let mut def = def_read_itr.next().unwrap();
732 while def > SPECIAL_THRESHOLD {
733 *def_write_itr.next().unwrap() = def;
734 def = def_read_itr.next().unwrap();
735 }
736 for _ in 0..multiplier {
737 *def_write_itr.next().unwrap() = def;
738 }
739 }
740 } else if self.def_levels.is_empty() {
741 assert!(self.rep_levels.len() >= self.current_len);
742 let mut rep_read_itr = self.rep_levels.iter().copied();
744 let mut rep_write_itr = self.spare_rep.iter_mut();
745 for _ in 0..old_len {
746 let rep = rep_read_itr.next().unwrap();
747 for _ in 0..multiplier {
748 *rep_write_itr.next().unwrap() = rep;
749 }
750 }
751 } else {
752 assert!(self.rep_levels.len() >= self.current_len);
753 assert!(self.def_levels.len() >= self.current_len);
754 let mut rep_read_itr = self.rep_levels.iter().copied();
755 let mut def_read_itr = self.def_levels.iter().copied();
756 let mut rep_write_itr = self.spare_rep.iter_mut();
757 let mut def_write_itr = self.spare_def.iter_mut();
758 for _ in 0..old_len {
759 let mut def = def_read_itr.next().unwrap();
760 while def > SPECIAL_THRESHOLD {
761 *def_write_itr.next().unwrap() = def;
762 *rep_write_itr.next().unwrap() = rep_read_itr.next().unwrap();
763 def = def_read_itr.next().unwrap();
764 }
765 let rep = rep_read_itr.next().unwrap();
766 for _ in 0..multiplier {
767 *def_write_itr.next().unwrap() = def;
768 *rep_write_itr.next().unwrap() = rep;
769 }
770 }
771 }
772 std::mem::swap(&mut self.def_levels, &mut self.spare_def);
773 std::mem::swap(&mut self.rep_levels, &mut self.spare_rep);
774 }
775
776 fn record_validity_buf(&mut self, validity: &Option<BooleanBuffer>) {
777 if let Some(validity) = validity {
778 let def_level = self.checkout_def(DefinitionInterpretation::NullableItem);
779 self.do_record_validity(validity, def_level);
780 } else {
781 self.checkout_def(DefinitionInterpretation::AllValidItem);
782 }
783 }
784
785 fn record_validity(&mut self, validity_desc: &ValidityDesc) {
786 self.record_validity_buf(&validity_desc.validity)
787 }
788
789 fn record_fsl(&mut self, fsl_desc: &FslDesc) {
790 self.has_fsl = true;
791 self.record_validity_buf(&fsl_desc.validity);
792 self.multiply_levels(fsl_desc.dimension);
793 }
794
795 fn normalize_specials(&mut self) {
796 for def in self.def_levels.iter_mut() {
797 if *def > SPECIAL_THRESHOLD {
798 *def -= SPECIAL_THRESHOLD;
799 }
800 }
801 }
802
803 fn normalize_specials_and_plan_splits(
804 &mut self,
805 def_meaning: &[DefinitionInterpretation],
806 max_levels_per_page: Option<u64>,
807 num_rows: u64,
808 num_values: u64,
809 ) -> Result<MiniBlockRepDefBudget> {
810 if self.def_levels.is_empty() {
817 return Ok(MiniBlockRepDefBudget::WithinBudget);
818 }
819
820 if self.rep_levels.is_empty() {
821 self.normalize_specials();
822 return Ok(MiniBlockRepDefBudget::WithinBudget);
823 }
824
825 if self.rep_levels.len() != self.def_levels.len() {
826 return Err(Error::internal(format!(
827 "Cannot plan structural page splits with mismatched rep/def lengths: rep={}, def={}",
828 self.rep_levels.len(),
829 self.def_levels.len()
830 )));
831 }
832
833 let Some(max_levels_per_page) = max_levels_per_page else {
834 self.normalize_specials();
835 return Ok(MiniBlockRepDefBudget::WithinBudget);
836 };
837
838 if num_values == 0 {
839 self.normalize_specials();
840 return Ok(MiniBlockRepDefBudget::WithinBudget);
841 }
842
843 let max_schema_rep = def_meaning.iter().filter(|level| level.is_list()).count() as u16;
844 let max_visible_level = SerializedRepDefs::max_visible_level(def_meaning);
845 let should_plan = !self.has_fsl && max_schema_rep > 0 && max_visible_level.is_some();
846
847 if !should_plan {
848 self.normalize_specials();
849 return Ok(MiniBlockRepDefBudget::WithinBudget);
850 }
851
852 let max_visible_level = max_visible_level.unwrap();
853 let mut splits = Vec::new();
854 let mut counted_rows = 0u64;
855 let mut counted_values = 0u64;
856 let mut saw_structural_overhead = false;
857 let mut single_row_over_budget_levels = None;
858
859 let mut current_row_level_start = None;
860 let mut current_row_num_values = 0u64;
861
862 let mut current_page_row_start = 0u64;
863 let mut current_page_num_rows = 0u64;
864 let mut current_page_level_start = 0usize;
865 let mut current_page_level_end = 0usize;
866 let mut current_page_value_start = 0u64;
867 let mut current_page_num_values = 0u64;
868 let mut current_page_num_levels = 0u64;
869 let mut current_page_has_structural_overhead = false;
870
871 let mut finish_row =
872 |row_level_start: usize, row_level_end: usize, row_num_values: u64| -> Result<()> {
873 let row_num_levels = (row_level_end - row_level_start) as u64;
874 let row_has_structural_overhead = row_num_levels > row_num_values;
875 saw_structural_overhead |= row_has_structural_overhead;
876
877 if row_has_structural_overhead && row_num_levels > max_levels_per_page {
878 single_row_over_budget_levels = Some(row_num_levels);
879 }
880
881 if current_page_num_rows > 0
882 && (current_page_has_structural_overhead || row_has_structural_overhead)
883 && current_page_num_levels + row_num_levels > max_levels_per_page
884 {
885 splits.push(MiniBlockRepDefSplit {
886 row_start: current_page_row_start,
887 num_rows: current_page_num_rows,
888 level_range: current_page_level_start..current_page_level_end,
889 value_start: current_page_value_start,
890 num_values: current_page_num_values,
891 });
892 current_page_row_start = counted_rows;
893 current_page_num_rows = 0;
894 current_page_level_start = row_level_start;
895 current_page_value_start = counted_values;
896 current_page_num_values = 0;
897 current_page_num_levels = 0;
898 current_page_has_structural_overhead = false;
899 }
900
901 if current_page_num_rows == 0 {
902 current_page_level_start = row_level_start;
903 }
904 current_page_num_rows += 1;
905 current_page_level_end = row_level_end;
906 current_page_num_values += row_num_values;
907 current_page_num_levels += row_num_levels;
908 current_page_has_structural_overhead |= row_has_structural_overhead;
909 counted_rows += 1;
910 counted_values += row_num_values;
911 Ok(())
912 };
913
914 for (idx, (rep_level, def_level)) in self
915 .rep_levels
916 .iter()
917 .copied()
918 .zip(self.def_levels.iter_mut())
919 .enumerate()
920 {
921 if *def_level > SPECIAL_THRESHOLD {
922 *def_level -= SPECIAL_THRESHOLD;
923 }
924
925 if rep_level == max_schema_rep {
926 if let Some(level_start) = current_row_level_start {
927 finish_row(level_start, idx, current_row_num_values)?;
928 current_row_num_values = 0;
929 } else if idx != 0 {
930 return Err(Error::internal(format!(
931 "Cannot plan structural page splits: first top-level row starts at level {}, expected 0",
932 idx
933 )));
934 }
935 current_row_level_start = Some(idx);
936 }
937
938 if current_row_level_start.is_none() {
939 return Err(Error::internal(
940 "Cannot plan structural page splits: found levels before the first top-level row start",
941 ));
942 }
943 if *def_level <= max_visible_level {
944 current_row_num_values += 1;
945 }
946 }
947
948 let Some(level_start) = current_row_level_start else {
949 return Err(Error::internal(
950 "Cannot plan structural page splits: found no top-level row starts",
951 ));
952 };
953 finish_row(level_start, self.rep_levels.len(), current_row_num_values)?;
954
955 if counted_rows != num_rows {
956 return Err(Error::internal(format!(
957 "Cannot plan structural page splits: expected {} top-level row starts, found {}",
958 num_rows, counted_rows
959 )));
960 }
961 if counted_values != num_values {
962 return Err(Error::internal(format!(
963 "Cannot plan structural page splits: counted {} visible values, expected {}",
964 counted_values, num_values
965 )));
966 }
967 if !saw_structural_overhead {
968 return Ok(MiniBlockRepDefBudget::WithinBudget);
969 }
970 if let Some(row_num_levels) = single_row_over_budget_levels {
971 return Ok(MiniBlockRepDefBudget::SingleRowOverBudget(row_num_levels));
972 }
973
974 if current_page_num_rows > 0 {
975 splits.push(MiniBlockRepDefSplit {
976 row_start: current_page_row_start,
977 num_rows: current_page_num_rows,
978 level_range: current_page_level_start..current_page_level_end,
979 value_start: current_page_value_start,
980 num_values: current_page_num_values,
981 });
982 }
983
984 if splits.len() > 1 {
985 Ok(MiniBlockRepDefBudget::RequiresPageSplit(splits))
986 } else {
987 Ok(MiniBlockRepDefBudget::WithinBudget)
988 }
989 }
990
991 fn build(mut self) -> SerializedRepDefs {
992 if self.current_len == 0 {
993 return SerializedRepDefs::new_with_fixed_size_list_levels(
994 None,
995 None,
996 self.def_meaning,
997 self.has_fsl,
998 );
999 }
1000
1001 self.normalize_specials();
1002
1003 let definition_levels = if self.def_levels.is_empty() {
1004 None
1005 } else {
1006 Some(self.def_levels)
1007 };
1008 let repetition_levels = if self.rep_levels.is_empty() {
1009 None
1010 } else {
1011 Some(self.rep_levels)
1012 };
1013
1014 let def_meaning = self.def_meaning.into_iter().rev().collect::<Vec<_>>();
1016
1017 SerializedRepDefs::new_with_fixed_size_list_levels(
1018 repetition_levels,
1019 definition_levels,
1020 def_meaning,
1021 self.has_fsl,
1022 )
1023 }
1024
1025 fn build_with_miniblock_repdef_budget(
1026 mut self,
1027 max_levels_per_page: Option<u64>,
1028 num_rows: u64,
1029 num_values: u64,
1030 ) -> Result<(SerializedRepDefs, MiniBlockRepDefBudget)> {
1031 if self.current_len == 0 {
1032 return Ok((
1033 SerializedRepDefs::new_with_fixed_size_list_levels(
1034 None,
1035 None,
1036 self.def_meaning,
1037 self.has_fsl,
1038 ),
1039 MiniBlockRepDefBudget::WithinBudget,
1040 ));
1041 }
1042
1043 let def_meaning = std::mem::take(&mut self.def_meaning)
1045 .into_iter()
1046 .rev()
1047 .collect::<Vec<_>>();
1048 let budget = self.normalize_specials_and_plan_splits(
1049 &def_meaning,
1050 max_levels_per_page,
1051 num_rows,
1052 num_values,
1053 )?;
1054
1055 let definition_levels = if self.def_levels.is_empty() {
1056 None
1057 } else {
1058 Some(self.def_levels)
1059 };
1060 let repetition_levels = if self.rep_levels.is_empty() {
1061 None
1062 } else {
1063 Some(self.rep_levels)
1064 };
1065
1066 Ok((
1067 SerializedRepDefs::new_with_fixed_size_list_levels(
1068 repetition_levels,
1069 definition_levels,
1070 def_meaning,
1071 self.has_fsl,
1072 ),
1073 budget,
1074 ))
1075 }
1076}
1077
1078#[derive(Clone, Default, Debug)]
1085pub struct RepDefBuilder {
1086 repdefs: Vec<RawRepDef>,
1088 len: Option<usize>,
1093}
1094
1095impl RepDefBuilder {
1096 fn check_validity_len(&mut self, incoming_len: usize) {
1097 if let Some(len) = self.len {
1098 assert_eq!(incoming_len, len);
1099 } else {
1100 self.len = Some(incoming_len);
1102 }
1103 }
1104
1105 fn num_layers(&self) -> usize {
1106 self.repdefs.len()
1107 }
1108
1109 pub fn is_empty(&self) -> bool {
1112 self.repdefs
1113 .iter()
1114 .all(|r| matches!(r, RawRepDef::Validity(ValidityDesc { validity: None, .. })))
1115 }
1116
1117 pub fn is_simple_validity(&self) -> bool {
1119 self.repdefs.len() == 1 && matches!(self.repdefs[0], RawRepDef::Validity(_))
1120 }
1121
1122 pub fn add_validity_bitmap(&mut self, validity: NullBuffer) {
1124 self.check_validity_len(validity.len());
1125 if validity.null_count() == 0 {
1126 self.add_no_null(validity.len());
1127 return;
1128 }
1129 self.repdefs.push(RawRepDef::Validity(ValidityDesc {
1130 num_values: validity.len(),
1131 validity: Some(validity.into_inner()),
1132 }));
1133 }
1134
1135 pub fn add_no_null(&mut self, len: usize) {
1137 self.check_validity_len(len);
1138 self.repdefs.push(RawRepDef::Validity(ValidityDesc {
1139 validity: None,
1140 num_values: len,
1141 }));
1142 }
1143
1144 pub fn add_fsl(&mut self, validity: Option<NullBuffer>, dimension: usize, num_values: usize) {
1145 if let Some(len) = self.len {
1146 assert_eq!(num_values, len);
1147 }
1148 self.len = Some(num_values * dimension);
1149 debug_assert!(validity.is_none() || validity.as_ref().unwrap().len() == num_values);
1150 self.repdefs.push(RawRepDef::Fsl(FslDesc {
1151 num_values,
1152 validity: validity.map(|v| v.into_inner()),
1153 dimension,
1154 }))
1155 }
1156
1157 fn check_offset_len(&mut self, offsets: &[i64]) {
1158 if let Some(len) = self.len {
1159 assert!(offsets.len() == len + 1);
1160 }
1161 self.len = Some(offsets[offsets.len() - 1] as usize);
1162 }
1163
1164 fn do_add_offsets(
1165 &mut self,
1166 lengths: impl Iterator<Item = i64>,
1167 validity: Option<NullBuffer>,
1168 capacity: usize,
1169 ) -> bool {
1170 let mut num_specials = 0;
1171 let mut has_empty_lists = false;
1172 let mut has_garbage_values = false;
1173 let mut last_off: i64 = 0;
1174
1175 let mut normalized_offsets = Vec::with_capacity(capacity);
1176 normalized_offsets.push(0);
1177
1178 if let Some(ref validity) = validity {
1179 for (len, is_valid) in lengths.zip(validity.iter()) {
1180 match (is_valid, len == 0) {
1181 (false, is_empty) => {
1182 num_specials += 1;
1183 has_garbage_values |= !is_empty;
1184 }
1185 (true, true) => {
1186 num_specials += 1;
1187 has_empty_lists = true;
1188 }
1189 _ => {
1190 last_off += len;
1191 }
1192 }
1193 normalized_offsets.push(last_off);
1194 }
1195 } else {
1196 for len in lengths {
1197 if len == 0 {
1198 num_specials += 1;
1199 has_empty_lists = true;
1200 }
1201 last_off += len;
1202 normalized_offsets.push(last_off);
1203 }
1204 }
1205
1206 self.check_offset_len(&normalized_offsets);
1207 self.repdefs.push(RawRepDef::Offsets(OffsetDesc {
1208 num_values: normalized_offsets.len() - 1,
1209 offsets: normalized_offsets.into(),
1210 validity: validity.map(|v| v.into_inner()),
1211 has_empty_lists,
1212 num_specials: num_specials as usize,
1213 }));
1214
1215 has_garbage_values
1216 }
1217
1218 pub fn add_offsets<O: OffsetSizeTrait>(
1225 &mut self,
1226 offsets: OffsetBuffer<O>,
1227 validity: Option<NullBuffer>,
1228 ) -> bool {
1229 let inner = offsets.into_inner();
1230 let buffer_len = inner.len();
1231
1232 if O::IS_LARGE {
1233 let i64_buff = ScalarBuffer::<i64>::new(inner.into_inner(), 0, buffer_len);
1234 let lengths = i64_buff.windows(2).map(|off| off[1] - off[0]);
1235 self.do_add_offsets(lengths, validity, buffer_len)
1236 } else {
1237 let i32_buff = ScalarBuffer::<i32>::new(inner.into_inner(), 0, buffer_len);
1238 let lengths = i32_buff.windows(2).map(|off| (off[1] - off[0]) as i64);
1239 self.do_add_offsets(lengths, validity, buffer_len)
1240 }
1241 }
1242
1243 fn concat_layers<'a>(
1255 layers: impl Iterator<Item = &'a RawRepDef>,
1256 num_layers: usize,
1257 ) -> RawRepDef {
1258 enum LayerKind {
1259 Validity,
1260 Fsl,
1261 Offsets,
1262 }
1263
1264 let mut collected = Vec::with_capacity(num_layers);
1267 let mut has_nulls = false;
1268 let mut layer_kind = LayerKind::Validity;
1269 let mut total_num_specials = 0;
1270 let mut all_dimension = 0;
1271 let mut all_has_empty_lists = false;
1272 let mut all_num_values = 0;
1273 for layer in layers {
1274 has_nulls |= layer.has_nulls();
1275 match layer {
1276 RawRepDef::Validity(_) => {
1277 layer_kind = LayerKind::Validity;
1278 }
1279 RawRepDef::Offsets(OffsetDesc {
1280 num_specials,
1281 has_empty_lists,
1282 ..
1283 }) => {
1284 all_has_empty_lists |= *has_empty_lists;
1285 layer_kind = LayerKind::Offsets;
1286 total_num_specials += num_specials;
1287 }
1288 RawRepDef::Fsl(FslDesc { dimension, .. }) => {
1289 layer_kind = LayerKind::Fsl;
1290 all_dimension = *dimension;
1291 }
1292 }
1293 collected.push(layer);
1294 all_num_values += layer.num_values();
1295 }
1296
1297 if !has_nulls {
1299 match layer_kind {
1300 LayerKind::Validity => {
1301 return RawRepDef::Validity(ValidityDesc {
1302 validity: None,
1303 num_values: all_num_values,
1304 });
1305 }
1306 LayerKind::Fsl => {
1307 return RawRepDef::Fsl(FslDesc {
1308 validity: None,
1309 num_values: all_num_values,
1310 dimension: all_dimension,
1311 });
1312 }
1313 LayerKind::Offsets => {}
1314 }
1315 }
1316
1317 let mut validity_builder = if has_nulls {
1319 BooleanBufferBuilder::new(all_num_values)
1320 } else {
1321 BooleanBufferBuilder::new(0)
1322 };
1323 let mut all_offsets = if matches!(layer_kind, LayerKind::Offsets) {
1324 let mut all_offsets = Vec::with_capacity(all_num_values);
1325 all_offsets.push(0);
1326 all_offsets
1327 } else {
1328 Vec::new()
1329 };
1330
1331 for layer in collected {
1332 match layer {
1333 RawRepDef::Validity(ValidityDesc {
1334 validity: Some(validity),
1335 ..
1336 }) => {
1337 validity_builder.append_buffer(validity);
1338 }
1339 RawRepDef::Validity(ValidityDesc {
1340 validity: None,
1341 num_values,
1342 }) => {
1343 validity_builder.append_n(*num_values, true);
1344 }
1345 RawRepDef::Fsl(FslDesc {
1346 validity,
1347 num_values,
1348 ..
1349 }) => {
1350 if let Some(validity) = validity {
1351 validity_builder.append_buffer(validity);
1352 } else {
1353 validity_builder.append_n(*num_values, true);
1354 }
1355 }
1356 RawRepDef::Offsets(OffsetDesc {
1357 offsets,
1358 validity: Some(validity),
1359 has_empty_lists,
1360 ..
1361 }) => {
1362 all_has_empty_lists |= has_empty_lists;
1363 validity_builder.append_buffer(validity);
1364 let last = *all_offsets.last().unwrap();
1365 all_offsets.extend(offsets.iter().skip(1).map(|off| *off + last));
1366 }
1367 RawRepDef::Offsets(OffsetDesc {
1368 offsets,
1369 validity: None,
1370 has_empty_lists,
1371 num_values,
1372 ..
1373 }) => {
1374 all_has_empty_lists |= has_empty_lists;
1375 if has_nulls {
1376 validity_builder.append_n(*num_values, true);
1377 }
1378 let last = *all_offsets.last().unwrap();
1379 all_offsets.extend(offsets.iter().skip(1).map(|off| *off + last));
1380 }
1381 }
1382 }
1383 let validity = if has_nulls {
1384 Some(validity_builder.finish())
1385 } else {
1386 None
1387 };
1388 match layer_kind {
1389 LayerKind::Fsl => RawRepDef::Fsl(FslDesc {
1390 validity,
1391 num_values: all_num_values,
1392 dimension: all_dimension,
1393 }),
1394 LayerKind::Validity => RawRepDef::Validity(ValidityDesc {
1395 validity,
1396 num_values: all_num_values,
1397 }),
1398 LayerKind::Offsets => RawRepDef::Offsets(OffsetDesc {
1399 offsets: all_offsets.into(),
1400 validity,
1401 has_empty_lists: all_has_empty_lists,
1402 num_values: all_num_values,
1403 num_specials: total_num_specials,
1404 }),
1405 }
1406 }
1407
1408 pub fn serialize(builders: Vec<Self>) -> SerializedRepDefs {
1411 Self::serialize_builders(builders).0.build()
1412 }
1413
1414 pub(crate) fn serialize_with_miniblock_repdef_budget(
1416 builders: Vec<Self>,
1417 max_levels_for_bits: impl FnOnce(u64) -> u64,
1418 num_rows: u64,
1419 num_values: u64,
1420 ) -> Result<(SerializedRepDefs, MiniBlockRepDefBudget)> {
1421 let (context, bits_per_level) = Self::serialize_builders(builders);
1422 context.build_with_miniblock_repdef_budget(
1423 bits_per_level.map(max_levels_for_bits),
1424 num_rows,
1425 num_values,
1426 )
1427 }
1428
1429 fn serialize_builders(builders: Vec<Self>) -> (SerializerContext, Option<u64>) {
1430 assert!(!builders.is_empty());
1431 if builders.iter().all(|b| b.is_empty()) {
1432 let def_meaning = builders
1434 .first()
1435 .unwrap()
1436 .repdefs
1437 .iter()
1438 .map(|_| DefinitionInterpretation::AllValidItem)
1439 .collect::<Vec<_>>();
1440 return (
1441 SerializerContext {
1442 def_meaning,
1443 rep_levels: LevelBuffer::default(),
1444 spare_rep: LevelBuffer::default(),
1445 def_levels: LevelBuffer::default(),
1446 spare_def: LevelBuffer::default(),
1447 current_rep: 0,
1448 current_def: 0,
1449 current_len: 0,
1450 current_num_specials: 0,
1451 has_fsl: false,
1452 },
1453 None,
1454 );
1455 }
1456
1457 let num_layers = builders[0].num_layers();
1458 let combined_layers = (0..num_layers)
1459 .map(|layer_index| {
1460 Self::concat_layers(
1461 builders.iter().map(|b| &b.repdefs[layer_index]),
1462 builders.len(),
1463 )
1464 })
1465 .collect::<Vec<_>>();
1466 debug_assert!(
1467 builders
1468 .iter()
1469 .all(|b| b.num_layers() == builders[0].num_layers())
1470 );
1471
1472 let total_len = combined_layers.last().unwrap().num_values()
1473 + combined_layers
1474 .iter()
1475 .map(|l| l.num_specials())
1476 .sum::<usize>();
1477 let max_rep = combined_layers.iter().map(|l| l.max_rep()).sum::<u16>();
1478 let max_def = combined_layers.iter().map(|l| l.max_def()).sum::<u16>();
1479 let bits_per_rep = if max_rep > 0 {
1480 u64::from(u16::BITS - max_rep.leading_zeros())
1481 } else {
1482 0
1483 };
1484 let bits_per_def = if max_def > 0 {
1485 u64::from(u16::BITS - max_def.leading_zeros())
1486 } else {
1487 0
1488 };
1489 let bits_per_level =
1490 (bits_per_rep + bits_per_def > 0).then_some(bits_per_rep + bits_per_def);
1491
1492 let mut context = SerializerContext::new(total_len, num_layers, max_rep, max_def);
1493 for layer in combined_layers.into_iter() {
1494 match layer {
1495 RawRepDef::Validity(def) => {
1496 context.record_validity(&def);
1497 }
1498 RawRepDef::Offsets(rep) => {
1499 context.record_offsets(&rep);
1500 }
1501 RawRepDef::Fsl(fsl) => {
1502 context.record_fsl(&fsl);
1503 }
1504 }
1505 }
1506 (context, bits_per_level)
1507 }
1508}
1509
1510#[derive(Debug)]
1515pub struct RepDefUnraveler {
1516 rep_levels: Option<LevelBuffer>,
1517 def_levels: Option<LevelBuffer>,
1518 levels_to_rep: Vec<u16>,
1520 def_meaning: Arc<[DefinitionInterpretation]>,
1521 current_def_cmp: u16,
1523 current_rep_cmp: u16,
1525 current_layer: usize,
1528 num_items: u64,
1530}
1531
1532impl RepDefUnraveler {
1533 pub fn new(
1535 rep_levels: Option<LevelBuffer>,
1536 def_levels: Option<LevelBuffer>,
1537 def_meaning: Arc<[DefinitionInterpretation]>,
1538 num_items: u64,
1539 ) -> Self {
1540 let mut levels_to_rep = Vec::with_capacity(def_meaning.len());
1541 let mut rep_counter = 0;
1542 levels_to_rep.push(0);
1544 for meaning in def_meaning.as_ref() {
1545 match meaning {
1546 DefinitionInterpretation::AllValidItem | DefinitionInterpretation::AllValidList => {
1547 }
1549 DefinitionInterpretation::NullableItem => {
1550 levels_to_rep.push(rep_counter);
1552 }
1553 DefinitionInterpretation::NullableList => {
1554 rep_counter += 1;
1555 levels_to_rep.push(rep_counter);
1556 }
1557 DefinitionInterpretation::EmptyableList => {
1558 rep_counter += 1;
1559 levels_to_rep.push(rep_counter);
1560 }
1561 DefinitionInterpretation::NullableAndEmptyableList => {
1562 rep_counter += 1;
1563 levels_to_rep.push(rep_counter);
1564 levels_to_rep.push(rep_counter);
1565 }
1566 }
1567 }
1568 Self {
1569 rep_levels,
1570 def_levels,
1571 current_def_cmp: 0,
1572 current_rep_cmp: 0,
1573 levels_to_rep,
1574 current_layer: 0,
1575 def_meaning,
1576 num_items,
1577 }
1578 }
1579
1580 pub fn is_all_valid(&self) -> bool {
1581 self.def_levels.is_none() || self.def_meaning[self.current_layer].is_all_valid()
1582 }
1583
1584 pub fn max_lists(&self) -> usize {
1590 debug_assert!(
1591 self.def_meaning[self.current_layer] != DefinitionInterpretation::NullableItem
1592 );
1593 self.rep_levels
1594 .as_ref()
1595 .map(|levels| levels.len())
1597 .unwrap_or(0)
1598 }
1599
1600 pub fn unravel_offsets<T: ArrowNativeType>(
1605 &mut self,
1606 offsets: &mut Vec<T>,
1607 validity: Option<&mut BooleanBufferBuilder>,
1608 ) -> Result<()> {
1609 let rep_levels = self
1610 .rep_levels
1611 .as_mut()
1612 .expect("Expected repetition level but data didn't contain repetition");
1613 let valid_level = self.current_def_cmp;
1614 let (null_level, empty_level) = match self.def_meaning[self.current_layer] {
1615 DefinitionInterpretation::NullableList => {
1616 self.current_def_cmp += 1;
1617 (valid_level + 1, 0)
1618 }
1619 DefinitionInterpretation::EmptyableList => {
1620 self.current_def_cmp += 1;
1621 (0, valid_level + 1)
1622 }
1623 DefinitionInterpretation::NullableAndEmptyableList => {
1624 self.current_def_cmp += 2;
1625 (valid_level + 1, valid_level + 2)
1626 }
1627 DefinitionInterpretation::AllValidList => (0, 0),
1628 _ => unreachable!(),
1629 };
1630 self.current_layer += 1;
1631
1632 let mut max_level = null_level.max(empty_level).max(valid_level);
1636 let upper_null = max_level;
1639 for level in self.def_meaning[self.current_layer..].iter() {
1640 match level {
1641 DefinitionInterpretation::NullableItem => {
1642 max_level += 1;
1643 }
1644 DefinitionInterpretation::AllValidItem => {}
1645 _ => {
1646 break;
1647 }
1648 }
1649 }
1650
1651 let mut curlen: usize = offsets.last().map(|o| o.as_usize()).unwrap_or(0);
1652
1653 offsets.pop();
1661
1662 let to_offset = |val: usize| {
1663 T::from_usize(val)
1664 .ok_or_else(|| Error::invalid_input("A single batch had more than i32::MAX values and so a large container type is required"))
1665 };
1666 self.current_rep_cmp += 1;
1667 if let Some(def_levels) = &mut self.def_levels {
1668 assert!(rep_levels.len() == def_levels.len());
1669 let mut push_validity: Box<dyn FnMut(bool)> = if let Some(validity) = validity {
1672 Box::new(|is_valid| validity.append(is_valid))
1673 } else {
1674 Box::new(|_| {})
1675 };
1676 let mut read_idx = 0;
1680 let mut write_idx = 0;
1681 while read_idx < rep_levels.len() {
1682 unsafe {
1685 let rep_val = *rep_levels.get_unchecked(read_idx);
1686 if rep_val != 0 {
1687 let def_val = *def_levels.get_unchecked(read_idx);
1688 *rep_levels.get_unchecked_mut(write_idx) = rep_val - 1;
1690 *def_levels.get_unchecked_mut(write_idx) = def_val;
1691 write_idx += 1;
1692
1693 if def_val == 0 {
1694 offsets.push(to_offset(curlen)?);
1696 curlen += 1;
1697 push_validity(true);
1698 } else if def_val > max_level {
1699 } else if def_val == null_level || def_val > upper_null {
1701 offsets.push(to_offset(curlen)?);
1703 push_validity(false);
1704 } else if def_val == empty_level {
1705 offsets.push(to_offset(curlen)?);
1707 push_validity(true);
1708 } else {
1709 offsets.push(to_offset(curlen)?);
1711 curlen += 1;
1712 push_validity(true);
1713 }
1714 } else {
1715 curlen += 1;
1716 }
1717 read_idx += 1;
1718 }
1719 }
1720 offsets.push(to_offset(curlen)?);
1721 rep_levels.truncate(write_idx);
1722 def_levels.truncate(write_idx);
1723 Ok(())
1724 } else {
1725 let mut read_idx = 0;
1727 let mut write_idx = 0;
1728 let old_offsets_len = offsets.len();
1729 while read_idx < rep_levels.len() {
1730 unsafe {
1732 let rep_val = *rep_levels.get_unchecked(read_idx);
1733 if rep_val != 0 {
1734 offsets.push(to_offset(curlen)?);
1736 *rep_levels.get_unchecked_mut(write_idx) = rep_val - 1;
1737 write_idx += 1;
1738 }
1739 curlen += 1;
1740 read_idx += 1;
1741 }
1742 }
1743 let num_new_lists = offsets.len() - old_offsets_len;
1744 offsets.push(to_offset(curlen)?);
1745 rep_levels.truncate(write_idx);
1750 if let Some(validity) = validity {
1751 validity.append_n(num_new_lists, true);
1754 }
1755 Ok(())
1756 }
1757 }
1758
1759 pub fn skip_validity(&mut self) {
1760 debug_assert!(self.is_all_valid());
1761 self.current_layer += 1;
1762 }
1763
1764 pub fn unravel_validity(&mut self, validity: &mut BooleanBufferBuilder) {
1766 let meaning = self.def_meaning[self.current_layer];
1767 if meaning == DefinitionInterpretation::AllValidItem || self.def_levels.is_none() {
1768 self.current_layer += 1;
1769 validity.append_n(self.num_items as usize, true);
1770 return;
1771 }
1772
1773 self.current_layer += 1;
1774 let def_levels = &self.def_levels.as_ref().unwrap();
1775
1776 let current_def_cmp = self.current_def_cmp;
1777 self.current_def_cmp += 1;
1778
1779 for is_valid in def_levels.iter().filter_map(|&level| {
1780 if self.levels_to_rep[level as usize] <= self.current_rep_cmp {
1781 Some(level <= current_def_cmp)
1782 } else {
1783 None
1784 }
1785 }) {
1786 validity.append(is_valid);
1787 }
1788 }
1789
1790 pub fn decimate(&mut self, dimension: usize) {
1791 if self.rep_levels.is_some() {
1792 todo!("Not yet supported FSL<...List<...>>");
1804 }
1805 let Some(def_levels) = self.def_levels.as_mut() else {
1806 return;
1807 };
1808 let mut read_idx = 0;
1809 let mut write_idx = 0;
1810 while read_idx < def_levels.len() {
1811 unsafe {
1812 *def_levels.get_unchecked_mut(write_idx) = *def_levels.get_unchecked(read_idx);
1813 }
1814 write_idx += 1;
1815 read_idx += dimension;
1816 }
1817 def_levels.truncate(write_idx);
1818 }
1819}
1820
1821#[derive(Debug)]
1835pub struct CompositeRepDefUnraveler {
1836 unravelers: Vec<RepDefUnraveler>,
1837}
1838
1839impl CompositeRepDefUnraveler {
1840 pub fn new(unravelers: Vec<RepDefUnraveler>) -> Self {
1841 Self { unravelers }
1842 }
1843
1844 pub fn unravel_validity(&mut self, num_values: usize) -> Option<NullBuffer> {
1848 let is_all_valid = self
1849 .unravelers
1850 .iter()
1851 .all(|unraveler| unraveler.is_all_valid());
1852
1853 if is_all_valid {
1854 for unraveler in self.unravelers.iter_mut() {
1855 unraveler.skip_validity();
1856 }
1857 None
1858 } else {
1859 let mut validity = BooleanBufferBuilder::new(num_values);
1860 for unraveler in self.unravelers.iter_mut() {
1861 unraveler.unravel_validity(&mut validity);
1862 }
1863 Some(NullBuffer::new(validity.finish()))
1864 }
1865 }
1866
1867 pub fn unravel_fsl_validity(
1868 &mut self,
1869 num_values: usize,
1870 dimension: usize,
1871 ) -> Option<NullBuffer> {
1872 for unraveler in self.unravelers.iter_mut() {
1873 unraveler.decimate(dimension);
1874 }
1875 self.unravel_validity(num_values)
1876 }
1877
1878 pub fn unravel_offsets<T: ArrowNativeType>(
1880 &mut self,
1881 ) -> Result<(OffsetBuffer<T>, Option<NullBuffer>)> {
1882 let mut is_all_valid = true;
1883 let mut max_num_lists = 0;
1884 for unraveler in self.unravelers.iter() {
1885 is_all_valid &= unraveler.is_all_valid();
1886 max_num_lists += unraveler.max_lists();
1887 }
1888
1889 let mut validity = if is_all_valid {
1890 None
1891 } else {
1892 Some(BooleanBufferBuilder::new(max_num_lists))
1895 };
1896
1897 let mut offsets = Vec::with_capacity(max_num_lists + 1);
1898
1899 for unraveler in self.unravelers.iter_mut() {
1900 unraveler.unravel_offsets(&mut offsets, validity.as_mut())?;
1901 }
1902
1903 Ok((
1904 OffsetBuffer::new(ScalarBuffer::from(offsets)),
1905 validity.map(|mut v| NullBuffer::new(v.finish())),
1906 ))
1907 }
1908}
1909
1910#[derive(Debug)]
1916pub struct BinaryControlWordIterator<I: Iterator<Item = (u16, u16)>, W> {
1917 repdef: I,
1918 def_width: usize,
1919 max_rep: u16,
1920 max_visible_def: u16,
1921 rep_mask: u16,
1922 def_mask: u16,
1923 bits_rep: u8,
1924 bits_def: u8,
1925 phantom: std::marker::PhantomData<W>,
1926}
1927
1928impl<I: Iterator<Item = (u16, u16)>> BinaryControlWordIterator<I, u8> {
1929 fn append_next(&mut self, buf: &mut Vec<u8>) -> Option<ControlWordDesc> {
1930 let next = self.repdef.next()?;
1931 let control_word: u8 =
1932 (((next.0 & self.rep_mask) as u8) << self.def_width) + ((next.1 & self.def_mask) as u8);
1933 buf.push(control_word);
1934 let is_new_row = next.0 == self.max_rep;
1935 let is_visible = next.1 <= self.max_visible_def;
1936 let is_valid_item = next.1 == 0;
1937 Some(ControlWordDesc {
1938 is_new_row,
1939 is_visible,
1940 is_valid_item,
1941 })
1942 }
1943}
1944
1945impl<I: Iterator<Item = (u16, u16)>> BinaryControlWordIterator<I, u16> {
1946 fn append_next(&mut self, buf: &mut Vec<u8>) -> Option<ControlWordDesc> {
1947 let next = self.repdef.next()?;
1948 let control_word: u16 =
1949 ((next.0 & self.rep_mask) << self.def_width) + (next.1 & self.def_mask);
1950 let control_word = control_word.to_le_bytes();
1951 buf.push(control_word[0]);
1952 buf.push(control_word[1]);
1953 let is_new_row = next.0 == self.max_rep;
1954 let is_visible = next.1 <= self.max_visible_def;
1955 let is_valid_item = next.1 == 0;
1956 Some(ControlWordDesc {
1957 is_new_row,
1958 is_visible,
1959 is_valid_item,
1960 })
1961 }
1962}
1963
1964impl<I: Iterator<Item = (u16, u16)>> BinaryControlWordIterator<I, u32> {
1965 fn append_next(&mut self, buf: &mut Vec<u8>) -> Option<ControlWordDesc> {
1966 let next = self.repdef.next()?;
1967 let control_word: u32 = (((next.0 & self.rep_mask) as u32) << self.def_width)
1968 + ((next.1 & self.def_mask) as u32);
1969 let control_word = control_word.to_le_bytes();
1970 buf.push(control_word[0]);
1971 buf.push(control_word[1]);
1972 buf.push(control_word[2]);
1973 buf.push(control_word[3]);
1974 let is_new_row = next.0 == self.max_rep;
1975 let is_visible = next.1 <= self.max_visible_def;
1976 let is_valid_item = next.1 == 0;
1977 Some(ControlWordDesc {
1978 is_new_row,
1979 is_visible,
1980 is_valid_item,
1981 })
1982 }
1983}
1984
1985#[derive(Debug)]
1987pub struct UnaryControlWordIterator<I: Iterator<Item = u16>, W> {
1988 repdef: I,
1989 level_mask: u16,
1990 bits_rep: u8,
1991 bits_def: u8,
1992 max_rep: u16,
1993 phantom: std::marker::PhantomData<W>,
1994}
1995
1996impl<I: Iterator<Item = u16>> UnaryControlWordIterator<I, u8> {
1997 fn append_next(&mut self, buf: &mut Vec<u8>) -> Option<ControlWordDesc> {
1998 let next = self.repdef.next()?;
1999 buf.push((next & self.level_mask) as u8);
2000 let is_new_row = self.max_rep == 0 || next == self.max_rep;
2001 let is_valid_item = next == 0 || self.bits_def == 0;
2002 Some(ControlWordDesc {
2003 is_new_row,
2004 is_visible: true,
2007 is_valid_item,
2008 })
2009 }
2010}
2011
2012impl<I: Iterator<Item = u16>> UnaryControlWordIterator<I, u16> {
2013 fn append_next(&mut self, buf: &mut Vec<u8>) -> Option<ControlWordDesc> {
2014 let next = self.repdef.next().unwrap() & self.level_mask;
2015 let control_word = next.to_le_bytes();
2016 buf.push(control_word[0]);
2017 buf.push(control_word[1]);
2018 let is_new_row = self.max_rep == 0 || next == self.max_rep;
2019 let is_valid_item = next == 0 || self.bits_def == 0;
2020 Some(ControlWordDesc {
2021 is_new_row,
2022 is_visible: true,
2023 is_valid_item,
2024 })
2025 }
2026}
2027
2028impl<I: Iterator<Item = u16>> UnaryControlWordIterator<I, u32> {
2029 fn append_next(&mut self, buf: &mut Vec<u8>) -> Option<ControlWordDesc> {
2030 let next = self.repdef.next()?;
2031 let next = (next & self.level_mask) as u32;
2032 let control_word = next.to_le_bytes();
2033 buf.push(control_word[0]);
2034 buf.push(control_word[1]);
2035 buf.push(control_word[2]);
2036 buf.push(control_word[3]);
2037 let is_new_row = self.max_rep == 0 || next as u16 == self.max_rep;
2038 let is_valid_item = next == 0 || self.bits_def == 0;
2039 Some(ControlWordDesc {
2040 is_new_row,
2041 is_visible: true,
2042 is_valid_item,
2043 })
2044 }
2045}
2046
2047#[derive(Debug)]
2049pub struct NilaryControlWordIterator {
2050 len: usize,
2051 idx: usize,
2052}
2053
2054impl NilaryControlWordIterator {
2055 fn append_next(&mut self) -> Option<ControlWordDesc> {
2056 if self.idx == self.len {
2057 None
2058 } else {
2059 self.idx += 1;
2060 Some(ControlWordDesc {
2061 is_new_row: true,
2062 is_visible: true,
2063 is_valid_item: true,
2064 })
2065 }
2066 }
2067}
2068
2069fn get_mask(width: u16) -> u16 {
2071 (1 << width) - 1
2072}
2073
2074type SpecificBinaryControlWordIterator<'a, T> = BinaryControlWordIterator<
2077 Zip<Copied<std::slice::Iter<'a, u16>>, Copied<std::slice::Iter<'a, u16>>>,
2078 T,
2079>;
2080
2081#[derive(Debug)]
2091pub enum ControlWordIterator<'a> {
2092 Binary8(SpecificBinaryControlWordIterator<'a, u8>),
2093 Binary16(SpecificBinaryControlWordIterator<'a, u16>),
2094 Binary32(SpecificBinaryControlWordIterator<'a, u32>),
2095 Unary8(UnaryControlWordIterator<Copied<std::slice::Iter<'a, u16>>, u8>),
2096 Unary16(UnaryControlWordIterator<Copied<std::slice::Iter<'a, u16>>, u16>),
2097 Unary32(UnaryControlWordIterator<Copied<std::slice::Iter<'a, u16>>, u32>),
2098 Nilary(NilaryControlWordIterator),
2099}
2100
2101#[derive(Debug)]
2103pub struct ControlWordDesc {
2104 pub is_new_row: bool,
2105 pub is_visible: bool,
2106 pub is_valid_item: bool,
2107}
2108
2109impl ControlWordIterator<'_> {
2110 pub fn append_next(&mut self, buf: &mut Vec<u8>) -> Option<ControlWordDesc> {
2114 match self {
2115 Self::Binary8(iter) => iter.append_next(buf),
2116 Self::Binary16(iter) => iter.append_next(buf),
2117 Self::Binary32(iter) => iter.append_next(buf),
2118 Self::Unary8(iter) => iter.append_next(buf),
2119 Self::Unary16(iter) => iter.append_next(buf),
2120 Self::Unary32(iter) => iter.append_next(buf),
2121 Self::Nilary(iter) => iter.append_next(),
2122 }
2123 }
2124
2125 pub fn has_repetition(&self) -> bool {
2127 match self {
2128 Self::Binary8(_) | Self::Binary16(_) | Self::Binary32(_) => true,
2129 Self::Unary8(iter) => iter.bits_rep > 0,
2130 Self::Unary16(iter) => iter.bits_rep > 0,
2131 Self::Unary32(iter) => iter.bits_rep > 0,
2132 Self::Nilary(_) => false,
2133 }
2134 }
2135
2136 pub fn bytes_per_word(&self) -> usize {
2138 match self {
2139 Self::Binary8(_) => 1,
2140 Self::Binary16(_) => 2,
2141 Self::Binary32(_) => 4,
2142 Self::Unary8(_) => 1,
2143 Self::Unary16(_) => 2,
2144 Self::Unary32(_) => 4,
2145 Self::Nilary(_) => 0,
2146 }
2147 }
2148
2149 pub fn bits_rep(&self) -> u8 {
2151 match self {
2152 Self::Binary8(iter) => iter.bits_rep,
2153 Self::Binary16(iter) => iter.bits_rep,
2154 Self::Binary32(iter) => iter.bits_rep,
2155 Self::Unary8(iter) => iter.bits_rep,
2156 Self::Unary16(iter) => iter.bits_rep,
2157 Self::Unary32(iter) => iter.bits_rep,
2158 Self::Nilary(_) => 0,
2159 }
2160 }
2161
2162 pub fn bits_def(&self) -> u8 {
2164 match self {
2165 Self::Binary8(iter) => iter.bits_def,
2166 Self::Binary16(iter) => iter.bits_def,
2167 Self::Binary32(iter) => iter.bits_def,
2168 Self::Unary8(iter) => iter.bits_def,
2169 Self::Unary16(iter) => iter.bits_def,
2170 Self::Unary32(iter) => iter.bits_def,
2171 Self::Nilary(_) => 0,
2172 }
2173 }
2174}
2175
2176pub fn build_control_word_iterator<'a>(
2180 rep: Option<&'a [u16]>,
2181 max_rep: u16,
2182 def: Option<&'a [u16]>,
2183 max_def: u16,
2184 max_visible_def: u16,
2185 len: usize,
2186) -> ControlWordIterator<'a> {
2187 let rep_width = if max_rep == 0 {
2188 0
2189 } else {
2190 log_2_ceil(max_rep as u32) as u16
2191 };
2192 let rep_mask = if max_rep == 0 { 0 } else { get_mask(rep_width) };
2193 let def_width = if max_def == 0 {
2194 0
2195 } else {
2196 log_2_ceil(max_def as u32) as u16
2197 };
2198 let def_mask = if max_def == 0 { 0 } else { get_mask(def_width) };
2199 let total_width = rep_width + def_width;
2200 match (rep, def) {
2201 (Some(rep), Some(def)) => {
2202 let iter = rep.iter().copied().zip(def.iter().copied());
2203 let def_width = def_width as usize;
2204 if total_width <= 8 {
2205 ControlWordIterator::Binary8(BinaryControlWordIterator {
2206 repdef: iter,
2207 rep_mask,
2208 def_mask,
2209 def_width,
2210 max_rep,
2211 max_visible_def,
2212 bits_rep: rep_width as u8,
2213 bits_def: def_width as u8,
2214 phantom: std::marker::PhantomData,
2215 })
2216 } else if total_width <= 16 {
2217 ControlWordIterator::Binary16(BinaryControlWordIterator {
2218 repdef: iter,
2219 rep_mask,
2220 def_mask,
2221 def_width,
2222 max_rep,
2223 max_visible_def,
2224 bits_rep: rep_width as u8,
2225 bits_def: def_width as u8,
2226 phantom: std::marker::PhantomData,
2227 })
2228 } else {
2229 ControlWordIterator::Binary32(BinaryControlWordIterator {
2230 repdef: iter,
2231 rep_mask,
2232 def_mask,
2233 def_width,
2234 max_rep,
2235 max_visible_def,
2236 bits_rep: rep_width as u8,
2237 bits_def: def_width as u8,
2238 phantom: std::marker::PhantomData,
2239 })
2240 }
2241 }
2242 (Some(lev), None) => {
2243 let iter = lev.iter().copied();
2244 if total_width <= 8 {
2245 ControlWordIterator::Unary8(UnaryControlWordIterator {
2246 repdef: iter,
2247 level_mask: rep_mask,
2248 bits_rep: total_width as u8,
2249 bits_def: 0,
2250 max_rep,
2251 phantom: std::marker::PhantomData,
2252 })
2253 } else if total_width <= 16 {
2254 ControlWordIterator::Unary16(UnaryControlWordIterator {
2255 repdef: iter,
2256 level_mask: rep_mask,
2257 bits_rep: total_width as u8,
2258 bits_def: 0,
2259 max_rep,
2260 phantom: std::marker::PhantomData,
2261 })
2262 } else {
2263 ControlWordIterator::Unary32(UnaryControlWordIterator {
2264 repdef: iter,
2265 level_mask: rep_mask,
2266 bits_rep: total_width as u8,
2267 bits_def: 0,
2268 max_rep,
2269 phantom: std::marker::PhantomData,
2270 })
2271 }
2272 }
2273 (None, Some(lev)) => {
2274 let iter = lev.iter().copied();
2275 if total_width <= 8 {
2276 ControlWordIterator::Unary8(UnaryControlWordIterator {
2277 repdef: iter,
2278 level_mask: def_mask,
2279 bits_rep: 0,
2280 bits_def: total_width as u8,
2281 max_rep: 0,
2282 phantom: std::marker::PhantomData,
2283 })
2284 } else if total_width <= 16 {
2285 ControlWordIterator::Unary16(UnaryControlWordIterator {
2286 repdef: iter,
2287 level_mask: def_mask,
2288 bits_rep: 0,
2289 bits_def: total_width as u8,
2290 max_rep: 0,
2291 phantom: std::marker::PhantomData,
2292 })
2293 } else {
2294 ControlWordIterator::Unary32(UnaryControlWordIterator {
2295 repdef: iter,
2296 level_mask: def_mask,
2297 bits_rep: 0,
2298 bits_def: total_width as u8,
2299 max_rep: 0,
2300 phantom: std::marker::PhantomData,
2301 })
2302 }
2303 }
2304 (None, None) => ControlWordIterator::Nilary(NilaryControlWordIterator { len, idx: 0 }),
2305 }
2306}
2307
2308#[derive(Copy, Clone, Debug)]
2312pub enum ControlWordParser {
2313 BOTH8(u8, u32),
2316 BOTH16(u8, u32),
2317 BOTH32(u8, u32),
2318 REP8,
2319 REP16,
2320 REP32,
2321 DEF8,
2322 DEF16,
2323 DEF32,
2324 NIL,
2325}
2326
2327impl ControlWordParser {
2328 fn parse_both<const WORD_SIZE: u8>(
2329 src: &[u8],
2330 dst_rep: &mut Vec<u16>,
2331 dst_def: &mut Vec<u16>,
2332 bits_to_shift: u8,
2333 mask_to_apply: u32,
2334 ) {
2335 match WORD_SIZE {
2336 1 => {
2337 let word = src[0];
2338 let rep = word >> bits_to_shift;
2339 let def = word & (mask_to_apply as u8);
2340 dst_rep.push(rep as u16);
2341 dst_def.push(def as u16);
2342 }
2343 2 => {
2344 let word = u16::from_le_bytes([src[0], src[1]]);
2345 let rep = word >> bits_to_shift;
2346 let def = word & mask_to_apply as u16;
2347 dst_rep.push(rep);
2348 dst_def.push(def);
2349 }
2350 4 => {
2351 let word = u32::from_le_bytes([src[0], src[1], src[2], src[3]]);
2352 let rep = word >> bits_to_shift;
2353 let def = word & mask_to_apply;
2354 dst_rep.push(rep as u16);
2355 dst_def.push(def as u16);
2356 }
2357 _ => unreachable!(),
2358 }
2359 }
2360
2361 fn parse_desc_both<const WORD_SIZE: u8>(
2362 src: &[u8],
2363 bits_to_shift: u8,
2364 mask_to_apply: u32,
2365 max_rep: u16,
2366 max_visible_def: u16,
2367 ) -> ControlWordDesc {
2368 match WORD_SIZE {
2369 1 => {
2370 let word = src[0];
2371 let rep = word >> bits_to_shift;
2372 let def = word & (mask_to_apply as u8);
2373 let is_visible = def as u16 <= max_visible_def;
2374 let is_new_row = rep as u16 == max_rep;
2375 let is_valid_item = def == 0;
2376 ControlWordDesc {
2377 is_visible,
2378 is_new_row,
2379 is_valid_item,
2380 }
2381 }
2382 2 => {
2383 let word = u16::from_le_bytes([src[0], src[1]]);
2384 let rep = word >> bits_to_shift;
2385 let def = word & mask_to_apply as u16;
2386 let is_visible = def <= max_visible_def;
2387 let is_new_row = rep == max_rep;
2388 let is_valid_item = def == 0;
2389 ControlWordDesc {
2390 is_visible,
2391 is_new_row,
2392 is_valid_item,
2393 }
2394 }
2395 4 => {
2396 let word = u32::from_le_bytes([src[0], src[1], src[2], src[3]]);
2397 let rep = word >> bits_to_shift;
2398 let def = word & mask_to_apply;
2399 let is_visible = def as u16 <= max_visible_def;
2400 let is_new_row = rep as u16 == max_rep;
2401 let is_valid_item = def == 0;
2402 ControlWordDesc {
2403 is_visible,
2404 is_new_row,
2405 is_valid_item,
2406 }
2407 }
2408 _ => unreachable!(),
2409 }
2410 }
2411
2412 fn parse_one<const WORD_SIZE: u8>(src: &[u8], dst: &mut Vec<u16>) {
2413 match WORD_SIZE {
2414 1 => {
2415 let word = src[0];
2416 dst.push(word as u16);
2417 }
2418 2 => {
2419 let word = u16::from_le_bytes([src[0], src[1]]);
2420 dst.push(word);
2421 }
2422 4 => {
2423 let word = u32::from_le_bytes([src[0], src[1], src[2], src[3]]);
2424 dst.push(word as u16);
2425 }
2426 _ => unreachable!(),
2427 }
2428 }
2429
2430 fn parse_rep_desc_one<const WORD_SIZE: u8>(src: &[u8], max_rep: u16) -> ControlWordDesc {
2431 match WORD_SIZE {
2432 1 => ControlWordDesc {
2433 is_new_row: src[0] as u16 == max_rep,
2434 is_visible: true,
2435 is_valid_item: true,
2436 },
2437 2 => ControlWordDesc {
2438 is_new_row: u16::from_le_bytes([src[0], src[1]]) == max_rep,
2439 is_visible: true,
2440 is_valid_item: true,
2441 },
2442 4 => ControlWordDesc {
2443 is_new_row: u32::from_le_bytes([src[0], src[1], src[2], src[3]]) as u16 == max_rep,
2444 is_visible: true,
2445 is_valid_item: true,
2446 },
2447 _ => unreachable!(),
2448 }
2449 }
2450
2451 fn parse_def_desc_one<const WORD_SIZE: u8>(src: &[u8]) -> ControlWordDesc {
2452 match WORD_SIZE {
2453 1 => ControlWordDesc {
2454 is_new_row: true,
2455 is_visible: true,
2456 is_valid_item: src[0] == 0,
2457 },
2458 2 => ControlWordDesc {
2459 is_new_row: true,
2460 is_visible: true,
2461 is_valid_item: u16::from_le_bytes([src[0], src[1]]) == 0,
2462 },
2463 4 => ControlWordDesc {
2464 is_new_row: true,
2465 is_visible: true,
2466 is_valid_item: u32::from_le_bytes([src[0], src[1], src[2], src[3]]) as u16 == 0,
2467 },
2468 _ => unreachable!(),
2469 }
2470 }
2471
2472 pub fn bytes_per_word(&self) -> usize {
2474 match self {
2475 Self::BOTH8(..) => 1,
2476 Self::BOTH16(..) => 2,
2477 Self::BOTH32(..) => 4,
2478 Self::REP8 => 1,
2479 Self::REP16 => 2,
2480 Self::REP32 => 4,
2481 Self::DEF8 => 1,
2482 Self::DEF16 => 2,
2483 Self::DEF32 => 4,
2484 Self::NIL => 0,
2485 }
2486 }
2487
2488 pub fn parse(&self, src: &[u8], dst_rep: &mut Vec<u16>, dst_def: &mut Vec<u16>) {
2495 match self {
2496 Self::BOTH8(bits_to_shift, mask_to_apply) => {
2497 Self::parse_both::<1>(src, dst_rep, dst_def, *bits_to_shift, *mask_to_apply)
2498 }
2499 Self::BOTH16(bits_to_shift, mask_to_apply) => {
2500 Self::parse_both::<2>(src, dst_rep, dst_def, *bits_to_shift, *mask_to_apply)
2501 }
2502 Self::BOTH32(bits_to_shift, mask_to_apply) => {
2503 Self::parse_both::<4>(src, dst_rep, dst_def, *bits_to_shift, *mask_to_apply)
2504 }
2505 Self::REP8 => Self::parse_one::<1>(src, dst_rep),
2506 Self::REP16 => Self::parse_one::<2>(src, dst_rep),
2507 Self::REP32 => Self::parse_one::<4>(src, dst_rep),
2508 Self::DEF8 => Self::parse_one::<1>(src, dst_def),
2509 Self::DEF16 => Self::parse_one::<2>(src, dst_def),
2510 Self::DEF32 => Self::parse_one::<4>(src, dst_def),
2511 Self::NIL => {}
2512 }
2513 }
2514
2515 pub fn has_rep(&self) -> bool {
2517 match self {
2518 Self::BOTH8(..)
2519 | Self::BOTH16(..)
2520 | Self::BOTH32(..)
2521 | Self::REP8
2522 | Self::REP16
2523 | Self::REP32 => true,
2524 Self::DEF8 | Self::DEF16 | Self::DEF32 | Self::NIL => false,
2525 }
2526 }
2527
2528 pub fn parse_desc(&self, src: &[u8], max_rep: u16, max_visible_def: u16) -> ControlWordDesc {
2530 match self {
2531 Self::BOTH8(bits_to_shift, mask_to_apply) => Self::parse_desc_both::<1>(
2532 src,
2533 *bits_to_shift,
2534 *mask_to_apply,
2535 max_rep,
2536 max_visible_def,
2537 ),
2538 Self::BOTH16(bits_to_shift, mask_to_apply) => Self::parse_desc_both::<2>(
2539 src,
2540 *bits_to_shift,
2541 *mask_to_apply,
2542 max_rep,
2543 max_visible_def,
2544 ),
2545 Self::BOTH32(bits_to_shift, mask_to_apply) => Self::parse_desc_both::<4>(
2546 src,
2547 *bits_to_shift,
2548 *mask_to_apply,
2549 max_rep,
2550 max_visible_def,
2551 ),
2552 Self::REP8 => Self::parse_rep_desc_one::<1>(src, max_rep),
2553 Self::REP16 => Self::parse_rep_desc_one::<2>(src, max_rep),
2554 Self::REP32 => Self::parse_rep_desc_one::<4>(src, max_rep),
2555 Self::DEF8 => Self::parse_def_desc_one::<1>(src),
2556 Self::DEF16 => Self::parse_def_desc_one::<2>(src),
2557 Self::DEF32 => Self::parse_def_desc_one::<4>(src),
2558 Self::NIL => ControlWordDesc {
2559 is_new_row: true,
2560 is_valid_item: true,
2561 is_visible: true,
2562 },
2563 }
2564 }
2565
2566 pub fn new(bits_rep: u8, bits_def: u8) -> Self {
2568 let total_bits = bits_rep + bits_def;
2569
2570 enum WordSize {
2571 One,
2572 Two,
2573 Four,
2574 }
2575
2576 let word_size = if total_bits <= 8 {
2577 WordSize::One
2578 } else if total_bits <= 16 {
2579 WordSize::Two
2580 } else {
2581 WordSize::Four
2582 };
2583
2584 match (bits_rep > 0, bits_def > 0, word_size) {
2585 (false, false, _) => Self::NIL,
2586 (false, true, WordSize::One) => Self::DEF8,
2587 (false, true, WordSize::Two) => Self::DEF16,
2588 (false, true, WordSize::Four) => Self::DEF32,
2589 (true, false, WordSize::One) => Self::REP8,
2590 (true, false, WordSize::Two) => Self::REP16,
2591 (true, false, WordSize::Four) => Self::REP32,
2592 (true, true, WordSize::One) => Self::BOTH8(bits_def, get_mask(bits_def as u16) as u32),
2593 (true, true, WordSize::Two) => Self::BOTH16(bits_def, get_mask(bits_def as u16) as u32),
2594 (true, true, WordSize::Four) => {
2595 Self::BOTH32(bits_def, get_mask(bits_def as u16) as u32)
2596 }
2597 }
2598 }
2599}
2600
2601#[cfg(test)]
2602mod tests {
2603 use arrow_buffer::{NullBuffer, OffsetBuffer, ScalarBuffer};
2604
2605 use crate::repdef::{
2606 CompositeRepDefUnraveler, DefinitionInterpretation, RepDefUnraveler, SerializedRepDefs,
2607 };
2608
2609 use super::RepDefBuilder;
2610
2611 fn validity(values: &[bool]) -> NullBuffer {
2612 NullBuffer::from_iter(values.iter().copied())
2613 }
2614
2615 fn offsets_32(values: &[i32]) -> OffsetBuffer<i32> {
2616 OffsetBuffer::<i32>::new(ScalarBuffer::from_iter(values.iter().copied()))
2617 }
2618
2619 fn offsets_64(values: &[i64]) -> OffsetBuffer<i64> {
2620 OffsetBuffer::<i64>::new(ScalarBuffer::from_iter(values.iter().copied()))
2621 }
2622
2623 #[test]
2624 fn test_repdef_empty_offsets() {
2625 let mut builder = RepDefBuilder::default();
2627 builder.add_offsets(offsets_32(&[0]), None);
2628 let repdefs = RepDefBuilder::serialize(vec![builder]);
2629 assert!(repdefs.repetition_levels.is_none());
2630 assert!(repdefs.definition_levels.is_none());
2631 }
2632
2633 #[test]
2634 fn test_repdef_basic() {
2635 let mut builder = RepDefBuilder::default();
2637 builder.add_offsets(
2638 offsets_64(&[0, 2, 2, 5]),
2639 Some(validity(&[true, false, true])),
2640 );
2641 builder.add_offsets(
2642 offsets_64(&[0, 1, 3, 5, 5, 9]),
2643 Some(validity(&[true, true, true, false, true])),
2644 );
2645 builder.add_validity_bitmap(validity(&[
2646 true, true, true, false, false, false, true, true, false,
2647 ]));
2648
2649 let repdefs = RepDefBuilder::serialize(vec![builder]);
2650 let rep = repdefs.repetition_levels.unwrap();
2651 let def = repdefs.definition_levels.unwrap();
2652
2653 assert_eq!(vec![0, 0, 0, 3, 1, 1, 2, 1, 0, 0, 1], *def);
2654 assert_eq!(vec![2, 1, 0, 2, 2, 0, 1, 1, 0, 0, 0], *rep);
2655
2656 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
2659 Some(rep.as_ref().to_vec()),
2660 Some(def.as_ref().to_vec()),
2661 repdefs.def_meaning.into(),
2662 9,
2663 )]);
2664
2665 assert_eq!(
2668 unraveler.unravel_validity(9),
2669 Some(validity(&[
2670 true, true, true, false, false, false, true, true, false
2671 ]))
2672 );
2673 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
2674 assert_eq!(off.inner(), offsets_32(&[0, 1, 3, 5, 5, 9]).inner());
2675 assert_eq!(val, Some(validity(&[true, true, true, false, true])));
2676 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
2677 assert_eq!(off.inner(), offsets_32(&[0, 2, 2, 5]).inner());
2678 assert_eq!(val, Some(validity(&[true, false, true])));
2679 }
2680
2681 #[test]
2682 fn test_repdef_simple_null_empty_list() {
2683 let check = |repdefs: SerializedRepDefs, last_def: DefinitionInterpretation| {
2684 let rep = repdefs.repetition_levels.unwrap();
2685 let def = repdefs.definition_levels.unwrap();
2686
2687 assert_eq!([1, 0, 1, 1, 0, 0], *rep);
2688 assert_eq!([0, 0, 2, 0, 1, 0], *def);
2689 assert_eq!(
2690 vec![DefinitionInterpretation::NullableItem, last_def,],
2691 repdefs.def_meaning
2692 );
2693 };
2694
2695 let mut builder = RepDefBuilder::default();
2699 builder.add_offsets(
2700 offsets_32(&[0, 2, 2, 5]),
2701 Some(validity(&[true, false, true])),
2702 );
2703 builder.add_validity_bitmap(validity(&[true, true, true, false, true]));
2704
2705 let repdefs = RepDefBuilder::serialize(vec![builder]);
2706
2707 check(repdefs, DefinitionInterpretation::NullableList);
2708
2709 let mut builder = RepDefBuilder::default();
2711 builder.add_offsets(offsets_32(&[0, 2, 2, 5]), None);
2712 builder.add_validity_bitmap(validity(&[true, true, true, false, true]));
2713
2714 let repdefs = RepDefBuilder::serialize(vec![builder]);
2715
2716 check(repdefs, DefinitionInterpretation::EmptyableList);
2717 }
2718
2719 #[test]
2720 fn test_repdef_empty_list_at_end() {
2721 let mut builder = RepDefBuilder::default();
2723 builder.add_offsets(offsets_32(&[0, 2, 5, 5]), None);
2724 builder.add_validity_bitmap(validity(&[true, true, true, false, true]));
2725
2726 let repdefs = RepDefBuilder::serialize(vec![builder]);
2727
2728 let rep = repdefs.repetition_levels.unwrap();
2729 let def = repdefs.definition_levels.unwrap();
2730
2731 assert_eq!([1, 0, 1, 0, 0, 1], *rep);
2732 assert_eq!([0, 0, 0, 1, 0, 2], *def);
2733 assert_eq!(
2734 vec![
2735 DefinitionInterpretation::NullableItem,
2736 DefinitionInterpretation::EmptyableList,
2737 ],
2738 repdefs.def_meaning
2739 );
2740 }
2741
2742 #[test]
2743 fn test_repdef_abnormal_nulls() {
2744 let mut builder = RepDefBuilder::default();
2747 builder.add_offsets(
2748 offsets_32(&[0, 2, 5, 8]),
2749 Some(validity(&[true, false, true])),
2750 );
2751 builder.add_no_null(5);
2754
2755 let repdefs = RepDefBuilder::serialize(vec![builder]);
2756
2757 let rep = repdefs.repetition_levels.unwrap();
2758 let def = repdefs.definition_levels.unwrap();
2759
2760 assert_eq!([1, 0, 1, 1, 0, 0], *rep);
2761 assert_eq!([0, 0, 1, 0, 0, 0], *def);
2762
2763 assert_eq!(
2764 vec![
2765 DefinitionInterpretation::AllValidItem,
2766 DefinitionInterpretation::NullableList,
2767 ],
2768 repdefs.def_meaning
2769 );
2770 }
2771
2772 #[test]
2773 fn test_repdef_fsl() {
2774 let mut builder = RepDefBuilder::default();
2775 builder.add_fsl(Some(validity(&[true, false])), 2, 2);
2776 builder.add_fsl(None, 2, 4);
2777 builder.add_validity_bitmap(validity(&[
2778 true, false, true, false, true, false, true, false,
2779 ]));
2780
2781 let repdefs = RepDefBuilder::serialize(vec![builder]);
2782
2783 assert_eq!(
2784 vec![
2785 DefinitionInterpretation::NullableItem,
2786 DefinitionInterpretation::AllValidItem,
2787 DefinitionInterpretation::NullableItem
2788 ],
2789 repdefs.def_meaning
2790 );
2791
2792 assert!(repdefs.repetition_levels.is_none());
2793
2794 let def = repdefs.definition_levels.unwrap();
2795
2796 assert_eq!([0, 1, 0, 1, 2, 2, 2, 2], *def);
2797
2798 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
2799 None,
2800 Some(def.as_ref().to_vec()),
2801 repdefs.def_meaning.into(),
2802 8,
2803 )]);
2804
2805 assert_eq!(
2806 unraveler.unravel_validity(8),
2807 Some(validity(&[
2808 true, false, true, false, false, false, false, false
2809 ]))
2810 );
2811 assert_eq!(unraveler.unravel_fsl_validity(4, 2), None);
2812 assert_eq!(
2813 unraveler.unravel_fsl_validity(2, 2),
2814 Some(validity(&[true, false]))
2815 );
2816 }
2817
2818 #[test]
2819 fn test_repdef_fsl_allvalid_item() {
2820 let mut builder = RepDefBuilder::default();
2821 builder.add_fsl(Some(validity(&[true, false])), 2, 2);
2822 builder.add_fsl(None, 2, 4);
2823 builder.add_no_null(8);
2824
2825 let repdefs = RepDefBuilder::serialize(vec![builder]);
2826
2827 assert_eq!(
2828 vec![
2829 DefinitionInterpretation::AllValidItem,
2830 DefinitionInterpretation::AllValidItem,
2831 DefinitionInterpretation::NullableItem
2832 ],
2833 repdefs.def_meaning
2834 );
2835
2836 assert!(repdefs.repetition_levels.is_none());
2837
2838 let def = repdefs.definition_levels.unwrap();
2839
2840 assert_eq!([0, 0, 0, 0, 1, 1, 1, 1], *def);
2841
2842 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
2843 None,
2844 Some(def.as_ref().to_vec()),
2845 repdefs.def_meaning.into(),
2846 8,
2847 )]);
2848
2849 assert_eq!(unraveler.unravel_validity(8), None);
2850 assert_eq!(unraveler.unravel_fsl_validity(4, 2), None);
2851 assert_eq!(
2852 unraveler.unravel_fsl_validity(2, 2),
2853 Some(validity(&[true, false]))
2854 );
2855 }
2856
2857 #[test]
2858 fn test_repdef_sliced_offsets() {
2859 let mut builder = RepDefBuilder::default();
2862 builder.add_offsets(
2863 offsets_32(&[5, 7, 7, 10]),
2864 Some(validity(&[true, false, true])),
2865 );
2866 builder.add_no_null(5);
2867
2868 let repdefs = RepDefBuilder::serialize(vec![builder]);
2869
2870 let rep = repdefs.repetition_levels.unwrap();
2871 let def = repdefs.definition_levels.unwrap();
2872
2873 assert_eq!([1, 0, 1, 1, 0, 0], *rep);
2874 assert_eq!([0, 0, 1, 0, 0, 0], *def);
2875
2876 assert_eq!(
2877 vec![
2878 DefinitionInterpretation::AllValidItem,
2879 DefinitionInterpretation::NullableList,
2880 ],
2881 repdefs.def_meaning
2882 );
2883 }
2884
2885 #[test]
2886 fn test_repdef_complex_null_empty() {
2887 let mut builder = RepDefBuilder::default();
2888 builder.add_offsets(
2889 offsets_32(&[0, 4, 4, 4, 6]),
2890 Some(validity(&[true, false, true, true])),
2891 );
2892 builder.add_offsets(
2893 offsets_32(&[0, 1, 1, 2, 2, 2, 3]),
2894 Some(validity(&[true, false, true, false, true, true])),
2895 );
2896 builder.add_no_null(3);
2897
2898 let repdefs = RepDefBuilder::serialize(vec![builder]);
2899
2900 let rep = repdefs.repetition_levels.unwrap();
2901 let def = repdefs.definition_levels.unwrap();
2902
2903 assert_eq!([2, 1, 1, 1, 2, 2, 2, 1], *rep);
2904 assert_eq!([0, 1, 0, 1, 3, 4, 2, 0], *def);
2905 }
2906
2907 #[test]
2908 fn test_repdef_empty_list_no_null() {
2909 let mut builder = RepDefBuilder::default();
2912 builder.add_offsets(offsets_32(&[0, 4, 4, 4, 6]), None);
2913 builder.add_no_null(6);
2914
2915 let repdefs = RepDefBuilder::serialize(vec![builder]);
2916
2917 let rep = repdefs.repetition_levels.unwrap();
2918 let def = repdefs.definition_levels.unwrap();
2919
2920 assert_eq!([1, 0, 0, 0, 1, 1, 1, 0], *rep);
2921 assert_eq!([0, 0, 0, 0, 1, 1, 0, 0], *def);
2922
2923 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
2924 Some(rep.as_ref().to_vec()),
2925 Some(def.as_ref().to_vec()),
2926 repdefs.def_meaning.into(),
2927 8,
2928 )]);
2929
2930 assert_eq!(unraveler.unravel_validity(6), None);
2931 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
2932 assert_eq!(off.inner(), offsets_32(&[0, 4, 4, 4, 6]).inner());
2933 assert_eq!(val, None);
2934 }
2935
2936 #[test]
2937 fn test_repdef_all_valid() {
2938 let mut builder = RepDefBuilder::default();
2939 builder.add_offsets(offsets_64(&[0, 2, 3, 5]), None);
2940 builder.add_offsets(offsets_64(&[0, 1, 3, 5, 7, 9]), None);
2941 builder.add_no_null(9);
2942
2943 let repdefs = RepDefBuilder::serialize(vec![builder]);
2944 let rep = repdefs.repetition_levels.unwrap();
2945 assert!(repdefs.definition_levels.is_none());
2946
2947 assert_eq!([2, 1, 0, 2, 0, 2, 0, 1, 0], *rep);
2948
2949 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
2950 Some(rep.as_ref().to_vec()),
2951 None,
2952 repdefs.def_meaning.into(),
2953 9,
2954 )]);
2955
2956 assert_eq!(unraveler.unravel_validity(9), None);
2957 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
2958 assert_eq!(off.inner(), offsets_32(&[0, 1, 3, 5, 7, 9]).inner());
2959 assert_eq!(val, None);
2960 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
2961 assert_eq!(off.inner(), offsets_32(&[0, 2, 3, 5]).inner());
2962 assert_eq!(val, None);
2963 }
2964
2965 #[test]
2966 fn test_repdef_nested_list_multibatch_matches_single() {
2967 let mut single = RepDefBuilder::default();
2971 single.add_offsets(offsets_64(&[0, 2, 3, 5]), None);
2972 single.add_offsets(offsets_64(&[0, 1, 3, 5, 7, 9]), None);
2973 single.add_no_null(9);
2974 let single_rep = RepDefBuilder::serialize(vec![single])
2975 .repetition_levels
2976 .unwrap();
2977
2978 let mut b0 = RepDefBuilder::default();
2982 b0.add_offsets(offsets_64(&[0, 2, 3]), None);
2983 b0.add_offsets(offsets_64(&[0, 1, 3, 5]), None);
2984 b0.add_no_null(5);
2985 let mut b1 = RepDefBuilder::default();
2986 b1.add_offsets(offsets_64(&[0, 2]), None);
2987 b1.add_offsets(offsets_64(&[0, 2, 4]), None);
2988 b1.add_no_null(4);
2989 let multi_rep = RepDefBuilder::serialize(vec![b0, b1])
2990 .repetition_levels
2991 .unwrap();
2992
2993 assert_eq!(
2994 *single_rep, *multi_rep,
2995 "multi-batch nested-list rep levels must equal single-batch"
2996 );
2997 }
2998
2999 #[test]
3000 fn test_only_empty_lists() {
3001 let mut builder = RepDefBuilder::default();
3002 builder.add_offsets(offsets_32(&[0, 4, 4, 4, 6]), None);
3003 builder.add_no_null(6);
3004
3005 let repdefs = RepDefBuilder::serialize(vec![builder]);
3006
3007 let rep = repdefs.repetition_levels.unwrap();
3008 let def = repdefs.definition_levels.unwrap();
3009
3010 assert_eq!([1, 0, 0, 0, 1, 1, 1, 0], *rep);
3011 assert_eq!([0, 0, 0, 0, 1, 1, 0, 0], *def);
3012
3013 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
3014 Some(rep.as_ref().to_vec()),
3015 Some(def.as_ref().to_vec()),
3016 repdefs.def_meaning.into(),
3017 8,
3018 )]);
3019
3020 assert_eq!(unraveler.unravel_validity(6), None);
3021 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
3022 assert_eq!(off.inner(), offsets_32(&[0, 4, 4, 4, 6]).inner());
3023 assert_eq!(val, None);
3024 }
3025
3026 #[test]
3027 fn test_only_null_lists() {
3028 let mut builder = RepDefBuilder::default();
3029 builder.add_offsets(
3030 offsets_32(&[0, 4, 4, 4, 6]),
3031 Some(validity(&[true, false, false, true])),
3032 );
3033 builder.add_no_null(6);
3034
3035 let repdefs = RepDefBuilder::serialize(vec![builder]);
3036
3037 let rep = repdefs.repetition_levels.unwrap();
3038 let def = repdefs.definition_levels.unwrap();
3039
3040 assert_eq!([1, 0, 0, 0, 1, 1, 1, 0], *rep);
3041 assert_eq!([0, 0, 0, 0, 1, 1, 0, 0], *def);
3042
3043 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
3044 Some(rep.as_ref().to_vec()),
3045 Some(def.as_ref().to_vec()),
3046 repdefs.def_meaning.into(),
3047 8,
3048 )]);
3049
3050 assert_eq!(unraveler.unravel_validity(6), None);
3051 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
3052 assert_eq!(off.inner(), offsets_32(&[0, 4, 4, 4, 6]).inner());
3053 assert_eq!(val, Some(validity(&[true, false, false, true])));
3054 }
3055
3056 #[test]
3057 fn test_null_and_empty_lists() {
3058 let mut builder = RepDefBuilder::default();
3059 builder.add_offsets(
3060 offsets_32(&[0, 4, 4, 4, 6]),
3061 Some(validity(&[true, false, true, true])),
3062 );
3063 builder.add_no_null(6);
3064
3065 let repdefs = RepDefBuilder::serialize(vec![builder]);
3066
3067 let rep = repdefs.repetition_levels.unwrap();
3068 let def = repdefs.definition_levels.unwrap();
3069
3070 assert_eq!([1, 0, 0, 0, 1, 1, 1, 0], *rep);
3071 assert_eq!([0, 0, 0, 0, 1, 2, 0, 0], *def);
3072
3073 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
3074 Some(rep.as_ref().to_vec()),
3075 Some(def.as_ref().to_vec()),
3076 repdefs.def_meaning.into(),
3077 8,
3078 )]);
3079
3080 assert_eq!(unraveler.unravel_validity(6), None);
3081 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
3082 assert_eq!(off.inner(), offsets_32(&[0, 4, 4, 4, 6]).inner());
3083 assert_eq!(val, Some(validity(&[true, false, true, true])));
3084 }
3085
3086 #[test]
3087 fn test_repdef_null_struct_valid_list() {
3088 let rep = vec![1, 0, 0, 0];
3091 let def = vec![2, 0, 2, 2];
3092 let def_meaning = vec![
3094 DefinitionInterpretation::NullableItem,
3095 DefinitionInterpretation::NullableItem,
3096 DefinitionInterpretation::AllValidList,
3097 ];
3098 let num_items = 4;
3099
3100 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
3101 Some(rep),
3102 Some(def),
3103 def_meaning.into(),
3104 num_items,
3105 )]);
3106
3107 assert_eq!(
3108 unraveler.unravel_validity(4),
3109 Some(validity(&[false, true, false, false]))
3110 );
3111 assert_eq!(
3112 unraveler.unravel_validity(4),
3113 Some(validity(&[false, true, false, false]))
3114 );
3115 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
3116 assert_eq!(off.inner(), offsets_32(&[0, 4]).inner());
3117 assert_eq!(val, None);
3118 }
3119
3120 #[test]
3121 fn test_repdef_no_rep() {
3122 let mut builder = RepDefBuilder::default();
3123 builder.add_no_null(5);
3124 builder.add_validity_bitmap(validity(&[false, false, true, true, true]));
3125 builder.add_validity_bitmap(validity(&[false, true, true, true, false]));
3126
3127 let repdefs = RepDefBuilder::serialize(vec![builder]);
3128 assert!(repdefs.repetition_levels.is_none());
3129 let def = repdefs.definition_levels.unwrap();
3130
3131 assert_eq!([2, 2, 0, 0, 1], *def);
3132
3133 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
3134 None,
3135 Some(def.as_ref().to_vec()),
3136 repdefs.def_meaning.into(),
3137 5,
3138 )]);
3139
3140 assert_eq!(
3141 unraveler.unravel_validity(5),
3142 Some(validity(&[false, false, true, true, false]))
3143 );
3144 assert_eq!(
3145 unraveler.unravel_validity(5),
3146 Some(validity(&[false, false, true, true, true]))
3147 );
3148 assert_eq!(unraveler.unravel_validity(5), None);
3149 }
3150
3151 #[test]
3152 fn test_composite_unravel() {
3153 let mut builder = RepDefBuilder::default();
3154 builder.add_offsets(
3155 offsets_64(&[0, 2, 2, 5]),
3156 Some(validity(&[true, false, true])),
3157 );
3158 builder.add_no_null(5);
3159 let repdef1 = RepDefBuilder::serialize(vec![builder]);
3160
3161 let mut builder = RepDefBuilder::default();
3162 builder.add_offsets(offsets_64(&[0, 1, 3, 5, 7, 9]), None);
3163 builder.add_no_null(9);
3164 let repdef2 = RepDefBuilder::serialize(vec![builder]);
3165
3166 let rep1 = repdef1.repetition_levels.clone().unwrap();
3167 let def1 = repdef1.definition_levels.clone().unwrap();
3168 let rep2 = repdef2.repetition_levels.clone().unwrap();
3169 assert!(repdef2.definition_levels.is_none());
3170
3171 assert_eq!([1, 0, 1, 1, 0, 0], *rep1);
3172 assert_eq!([0, 0, 1, 0, 0, 0], *def1);
3173 assert_eq!([1, 1, 0, 1, 0, 1, 0, 1, 0], *rep2);
3174
3175 let unravel1 = RepDefUnraveler::new(
3176 repdef1.repetition_levels.map(|l| l.to_vec()),
3177 repdef1.definition_levels.map(|l| l.to_vec()),
3178 repdef1.def_meaning.into(),
3179 5,
3180 );
3181 let unravel2 = RepDefUnraveler::new(
3182 repdef2.repetition_levels.map(|l| l.to_vec()),
3183 repdef2.definition_levels.map(|l| l.to_vec()),
3184 repdef2.def_meaning.into(),
3185 9,
3186 );
3187
3188 let mut unraveler = CompositeRepDefUnraveler::new(vec![unravel1, unravel2]);
3189
3190 assert!(unraveler.unravel_validity(9).is_none());
3191 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
3192 assert_eq!(
3193 off.inner(),
3194 offsets_32(&[0, 2, 2, 5, 6, 8, 10, 12, 14]).inner()
3195 );
3196 assert_eq!(
3197 val,
3198 Some(validity(&[true, false, true, true, true, true, true, true]))
3199 );
3200 }
3201
3202 #[test]
3203 fn test_repdef_multiple_builders() {
3204 let mut builder1 = RepDefBuilder::default();
3206 builder1.add_offsets(offsets_64(&[0, 2]), None);
3207 builder1.add_offsets(offsets_64(&[0, 1, 3]), None);
3208 builder1.add_validity_bitmap(validity(&[true, true, true]));
3209
3210 let mut builder2 = RepDefBuilder::default();
3211 builder2.add_offsets(offsets_64(&[0, 0, 3]), Some(validity(&[false, true])));
3212 builder2.add_offsets(
3213 offsets_64(&[0, 2, 2, 6]),
3214 Some(validity(&[true, false, true])),
3215 );
3216 builder2.add_validity_bitmap(validity(&[false, false, false, true, true, false]));
3217
3218 let repdefs = RepDefBuilder::serialize(vec![builder1, builder2]);
3219
3220 let rep = repdefs.repetition_levels.unwrap();
3221 let def = repdefs.definition_levels.unwrap();
3222
3223 assert_eq!([2, 1, 0, 2, 2, 0, 1, 1, 0, 0, 0], *rep);
3224 assert_eq!([0, 0, 0, 3, 1, 1, 2, 1, 0, 0, 1], *def);
3225 }
3226
3227 #[test]
3228 fn test_all_valid_validity_bitmap_serializes_as_no_null() {
3229 let mut from_bitmap = RepDefBuilder::default();
3230 from_bitmap.add_validity_bitmap(validity(&[true, true, true, true]));
3231
3232 let mut from_no_null = RepDefBuilder::default();
3233 from_no_null.add_no_null(4);
3234
3235 let from_bitmap = RepDefBuilder::serialize(vec![from_bitmap]);
3236 let from_no_null = RepDefBuilder::serialize(vec![from_no_null]);
3237
3238 assert!(from_bitmap.repetition_levels.is_none());
3239 assert!(from_bitmap.definition_levels.is_none());
3240 assert_eq!(from_bitmap.def_meaning, from_no_null.def_meaning);
3241 assert_eq!(
3242 from_bitmap.max_visible_level,
3243 from_no_null.max_visible_level
3244 );
3245 }
3246
3247 #[test]
3248 fn test_slicer() {
3249 let mut builder = RepDefBuilder::default();
3250 builder.add_offsets(
3251 offsets_64(&[0, 2, 2, 30, 30]),
3252 Some(validity(&[true, false, true, true])),
3253 );
3254 builder.add_no_null(30);
3255
3256 let repdefs = RepDefBuilder::serialize(vec![builder]);
3257
3258 let mut rep_slicer = repdefs.rep_slicer().unwrap();
3259
3260 assert_eq!(rep_slicer.slice_next(5).len(), 12);
3262 assert_eq!(rep_slicer.slice_next(20).len(), 40);
3264 assert_eq!(rep_slicer.slice_rest().len(), 12);
3266
3267 let mut def_slicer = repdefs.rep_slicer().unwrap();
3268
3269 assert_eq!(def_slicer.slice_next(5).len(), 12);
3271 assert_eq!(def_slicer.slice_next(20).len(), 40);
3273 assert_eq!(def_slicer.slice_rest().len(), 12);
3275 }
3276
3277 #[test]
3278 fn test_control_words() {
3279 fn check(
3281 rep: &[u16],
3282 def: &[u16],
3283 expected_values: Vec<u8>,
3284 expected_bytes_per_word: usize,
3285 expected_bits_rep: u8,
3286 expected_bits_def: u8,
3287 ) {
3288 let num_vals = rep.len().max(def.len());
3289 let max_rep = rep.iter().max().copied().unwrap_or(0);
3290 let max_def = def.iter().max().copied().unwrap_or(0);
3291
3292 let in_rep = if rep.is_empty() { None } else { Some(rep) };
3293 let in_def = if def.is_empty() { None } else { Some(def) };
3294
3295 let mut iter = super::build_control_word_iterator(
3296 in_rep,
3297 max_rep,
3298 in_def,
3299 max_def,
3300 max_def + 1,
3301 expected_values.len(),
3302 );
3303 assert_eq!(iter.bytes_per_word(), expected_bytes_per_word);
3304 assert_eq!(iter.bits_rep(), expected_bits_rep);
3305 assert_eq!(iter.bits_def(), expected_bits_def);
3306 let mut cw_vec = Vec::with_capacity(num_vals * iter.bytes_per_word());
3307
3308 for _ in 0..num_vals {
3309 iter.append_next(&mut cw_vec);
3310 }
3311 assert!(iter.append_next(&mut cw_vec).is_none());
3312
3313 assert_eq!(expected_values, cw_vec);
3314
3315 let parser = super::ControlWordParser::new(expected_bits_rep, expected_bits_def);
3316
3317 let mut rep_out = Vec::with_capacity(num_vals);
3318 let mut def_out = Vec::with_capacity(num_vals);
3319
3320 if expected_bytes_per_word > 0 {
3321 for slice in cw_vec.chunks_exact(expected_bytes_per_word) {
3322 parser.parse(slice, &mut rep_out, &mut def_out);
3323 }
3324 }
3325
3326 assert_eq!(rep, rep_out.as_slice());
3327 assert_eq!(def, def_out.as_slice());
3328 }
3329
3330 let rep = &[0_u16, 7, 3, 2, 9, 8, 12, 5];
3332 let def = &[5_u16, 3, 1, 2, 12, 15, 0, 2];
3333 let expected = vec![
3334 0b00000101, 0b01110011, 0b00110001, 0b00100010, 0b10011100, 0b10001111, 0b11000000, 0b01010010, ];
3343 check(rep, def, expected, 1, 4, 4);
3344
3345 let rep = &[0_u16, 7, 3, 2, 9, 8, 12, 5];
3347 let def = &[5_u16, 3, 1, 2, 12, 22, 0, 2];
3348 let expected = vec![
3349 0b00000101, 0b00000000, 0b11100011, 0b00000000, 0b01100001, 0b00000000, 0b01000010, 0b00000000, 0b00101100, 0b00000001, 0b00010110, 0b00000001, 0b10000000, 0b00000001, 0b10100010, 0b00000000, ];
3358 check(rep, def, expected, 2, 4, 5);
3359
3360 let levels = &[0_u16, 7, 3, 2, 9, 8, 12, 5];
3362 let expected = vec![
3363 0b00000000, 0b00000111, 0b00000011, 0b00000010, 0b00001001, 0b00001000, 0b00001100, 0b00000101, ];
3372 check(levels, &[], expected.clone(), 1, 4, 0);
3373
3374 check(&[], levels, expected, 1, 0, 4);
3376
3377 check(&[], &[], Vec::default(), 0, 0, 0);
3379 }
3380
3381 #[test]
3382 fn test_control_words_rep_index() {
3383 fn check(
3384 rep: &[u16],
3385 def: &[u16],
3386 expected_new_rows: Vec<bool>,
3387 expected_is_visible: Vec<bool>,
3388 ) {
3389 let num_vals = rep.len().max(def.len());
3390 let max_rep = rep.iter().max().copied().unwrap_or(0);
3391 let max_def = def.iter().max().copied().unwrap_or(0);
3392
3393 let in_rep = if rep.is_empty() { None } else { Some(rep) };
3394 let in_def = if def.is_empty() { None } else { Some(def) };
3395
3396 let mut iter = super::build_control_word_iterator(
3397 in_rep,
3398 max_rep,
3399 in_def,
3400 max_def,
3401 2,
3402 expected_new_rows.len(),
3403 );
3404
3405 let mut cw_vec = Vec::with_capacity(num_vals * iter.bytes_per_word());
3406 let mut expected_new_rows = expected_new_rows.iter().copied();
3407 let mut expected_is_visible = expected_is_visible.iter().copied();
3408 for _ in 0..expected_new_rows.len() {
3409 let word_desc = iter.append_next(&mut cw_vec).unwrap();
3410 assert_eq!(word_desc.is_new_row, expected_new_rows.next().unwrap());
3411 assert_eq!(word_desc.is_visible, expected_is_visible.next().unwrap());
3412 }
3413 assert!(iter.append_next(&mut cw_vec).is_none());
3414 }
3415
3416 let rep = &[2_u16, 1, 0, 2, 2, 0, 1, 1, 0, 2, 0];
3418 let def = &[0_u16, 0, 0, 3, 1, 1, 2, 1, 0, 0, 1];
3420
3421 check(
3423 rep,
3424 def,
3425 vec![
3426 true, false, false, true, true, false, false, false, false, true, false,
3427 ],
3428 vec![
3429 true, true, true, false, true, true, true, true, true, true, true,
3430 ],
3431 );
3432 check(
3434 rep,
3435 &[],
3436 vec![
3437 true, false, false, true, true, false, false, false, false, true, false,
3438 ],
3439 vec![true; 11],
3440 );
3441 check(
3443 &[],
3444 def,
3445 vec![
3446 true, true, true, true, true, true, true, true, true, true, true,
3447 ],
3448 vec![true; 11],
3449 );
3450 check(
3452 &[],
3453 &[],
3454 vec![
3455 true, true, true, true, true, true, true, true, true, true, true,
3456 ],
3457 vec![true; 11],
3458 );
3459 }
3460
3461 #[test]
3462 fn regress_empty_list_case() {
3463 let mut builder = RepDefBuilder::default();
3465 builder.add_validity_bitmap(validity(&[true, false, true]));
3466 builder.add_offsets(
3467 offsets_32(&[0, 0, 0, 0]),
3468 Some(validity(&[false, false, false])),
3469 );
3470 builder.add_no_null(0);
3471
3472 let repdefs = RepDefBuilder::serialize(vec![builder]);
3473 let rep = repdefs.repetition_levels.unwrap();
3474 let def = repdefs.definition_levels.unwrap();
3475
3476 assert_eq!([1, 1, 1], *rep);
3477 assert_eq!([1, 2, 1], *def);
3478
3479 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
3480 Some(rep.as_ref().to_vec()),
3481 Some(def.as_ref().to_vec()),
3482 repdefs.def_meaning.into(),
3483 0,
3484 )]);
3485
3486 assert_eq!(unraveler.unravel_validity(0), None);
3487 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
3488 assert_eq!(off.inner(), offsets_32(&[0, 0, 0, 0]).inner());
3489 assert_eq!(val, Some(validity(&[false, false, false])));
3490 let val = unraveler.unravel_validity(3).unwrap();
3491 assert_eq!(val.inner(), validity(&[true, false, true]).inner());
3492 }
3493
3494 #[test]
3495 fn regress_list_ends_null_case() {
3496 let mut builder = RepDefBuilder::default();
3497 builder.add_offsets(
3498 offsets_64(&[0, 1, 2, 2]),
3499 Some(validity(&[true, true, false])),
3500 );
3501 builder.add_offsets(offsets_64(&[0, 1, 1]), Some(validity(&[true, false])));
3502 builder.add_no_null(1);
3503
3504 let repdefs = RepDefBuilder::serialize(vec![builder]);
3505 let rep = repdefs.repetition_levels.unwrap();
3506 let def = repdefs.definition_levels.unwrap();
3507
3508 assert_eq!([2, 2, 2], *rep);
3509 assert_eq!([0, 1, 2], *def);
3510
3511 let mut unraveler = CompositeRepDefUnraveler::new(vec![RepDefUnraveler::new(
3512 Some(rep.as_ref().to_vec()),
3513 Some(def.as_ref().to_vec()),
3514 repdefs.def_meaning.into(),
3515 1,
3516 )]);
3517
3518 assert_eq!(unraveler.unravel_validity(1), None);
3519 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
3520 assert_eq!(off.inner(), offsets_32(&[0, 1, 1]).inner());
3521 assert_eq!(val, Some(validity(&[true, false])));
3522 let (off, val) = unraveler.unravel_offsets::<i32>().unwrap();
3523 assert_eq!(off.inner(), offsets_32(&[0, 1, 2, 2]).inner());
3524 assert_eq!(val, Some(validity(&[true, true, false])));
3525 }
3526
3527 #[test]
3528 fn test_mixed_unraveler() {
3529 let mut unraveler = CompositeRepDefUnraveler::new(vec![
3534 RepDefUnraveler::new(
3535 None,
3536 Some(vec![0, 1, 0, 1]),
3537 vec![DefinitionInterpretation::NullableItem].into(),
3538 4,
3539 ),
3540 RepDefUnraveler::new(
3541 None,
3542 None,
3543 vec![DefinitionInterpretation::AllValidItem].into(),
3544 4,
3545 ),
3546 ]);
3547
3548 assert_eq!(
3549 unraveler.unravel_validity(8),
3550 Some(validity(&[
3551 true, false, true, false, true, true, true, true
3552 ]))
3553 );
3554
3555 let def1 = Some(vec![0, 1, 2]);
3557 let rep1 = Some(vec![1, 0, 1]);
3558
3559 let def2 = Some(vec![1, 0, 0]);
3560 let rep2 = Some(vec![1, 1, 0]);
3561
3562 let mut unraveler = CompositeRepDefUnraveler::new(vec![
3563 RepDefUnraveler::new(
3564 rep1,
3565 def1,
3566 vec![
3567 DefinitionInterpretation::NullableItem,
3568 DefinitionInterpretation::EmptyableList,
3569 ]
3570 .into(),
3571 2,
3572 ),
3573 RepDefUnraveler::new(
3574 rep2,
3575 def2,
3576 vec![
3577 DefinitionInterpretation::AllValidItem,
3578 DefinitionInterpretation::NullableList,
3579 ]
3580 .into(),
3581 2,
3582 ),
3583 ]);
3584
3585 assert_eq!(
3586 unraveler.unravel_validity(4),
3587 Some(validity(&[true, false, true, true]))
3588 );
3589 assert_eq!(
3590 unraveler.unravel_offsets::<i32>().unwrap(),
3591 (
3592 offsets_32(&[0, 2, 2, 2, 4]),
3593 Some(validity(&[true, true, false, true]))
3594 )
3595 );
3596 }
3597
3598 #[test]
3599 fn test_mixed_unraveler_nullable_without_def_levels() {
3600 let mut unraveler = CompositeRepDefUnraveler::new(vec![
3603 RepDefUnraveler::new(
3604 None,
3605 Some(vec![0, 1, 0, 1]),
3606 vec![DefinitionInterpretation::NullableItem].into(),
3607 4,
3608 ),
3609 RepDefUnraveler::new(
3610 None,
3611 None,
3612 vec![DefinitionInterpretation::NullableItem].into(),
3613 4,
3614 ),
3615 ]);
3616
3617 assert_eq!(
3618 unraveler.unravel_validity(8),
3619 Some(validity(&[
3620 true, false, true, false, true, true, true, true
3621 ]))
3622 );
3623 }
3624}