Skip to main content

lance_encoding/
repdef.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4//! Utilities for rep-def levels
5//!
6//! Repetition and definition levels are a way to encode multipile validity / offsets arrays
7//! into a single buffer.  They are a form of "zipping" buffers together that takes advantage
8//! of the fact that, if the outermost array is invalid, then the validity of the inner items
9//! is irrelevant.
10//!
11//! Note: the concept of repetition & definition levels comes from the Dremel paper and has
12//! been implemented in Apache Parquet.  However, the implementation here is not necessarily
13//! compatible with Parquet.  For example, we use 0 to represent the "inner-most" item and
14//! Parquet uses 0 to represent the "outer-most" item.
15//!
16//! # Repetition Levels
17//!
18//! With repetition levels we convert a sparse array of offsets into a dense array of levels.
19//! These levels are marked non-zero whenever a new list begins.  In other words, given the
20//! list array with 3 rows [{<0,1>, <>, <2>}, {<3>}, {}], [], [{<4>}] we would have three
21//! offsets arrays:
22//!
23//! Outer-most ([]): [0, 3, 3, 4]
24//! Middle     ({}): [0, 3, 4, 4, 5]
25//! Inner      (<>): [0, 2, 2, 3, 4, 5]
26//! Values         : [0, 1, 2, 3, 4]
27//!
28//! We can convert these into repetition levels as follows:
29//!
30//! | Values | Repetition |
31//! | ------ | ---------- |
32//! |      0 |          3 | // Start of outer-most list
33//! |      1 |          0 | // Continues inner-most list (no new lists)
34//! |      - |          1 | // Start of new inner-most list (empty list)
35//! |      2 |          1 | // Start of new inner-most list
36//! |      3 |          2 | // Start of new middle list
37//! |      - |          2 | // Start of new inner-most list (empty list)
38//! |      - |          3 | // Start of new outer-most list (empty list)
39//! |      4 |          0 | // Start of new outer-most list
40//!
41//! Note: We actually have MORE repetition levels than values.  This is because the repetition
42//! levels need to be able to represent empty lists.
43//!
44//! # Definition Levels
45//!
46//! Definition levels are simpler.  We can think of them as zipping together various validity bitmaps
47//! (from different levels of nesting) into a single buffer.  For example, we could zip the arrays
48//! [1, 1, 0, 0] and [1, 0, 1, 0] into [11, 10, 01, 00].  However, 00 and 01 are redundant.  If the
49//! outer level is null then the validity of the inner levels is irrelevant.  To save space we instead
50//! encode a "level" which is the "depth" of the null.  Let's look at a more complete example:
51//!
52//! Array: [{"middle": {"inner": 1]}}, NULL, {"middle": NULL}, {"middle": {"inner": NULL}}]
53//!
54//! In Arrow we would have the following validity arrays:
55//! Outer validity : 1, 0, 1, 1
56//! Middle validity: 1, ?, 0, 1
57//! Inner validity : 1, ?, ?, 0
58//! Values         : 1, ?, ?, ?
59//!
60//! The ? values are undefined in the Arrow format.  We can convert these into definition levels as follows:
61//!
62//! | Values | Definition |
63//! | ------ | ---------- |
64//! |      1 |          0 | // Valid at all levels
65//! |      - |          3 | // Null at outer level
66//! |      - |          2 | // Null at middle level
67//! |      - |          1 | // Null at inner level
68//!
69//! # Compression
70//!
71//! Note that we only need 2 bits of definition levels to represent 3 levels of nesting.  Definition
72//! levels are always more compact than the input validity arrays.  However, compressed levels are not
73//! necessarily more compact than the compressed validity arrays.
74//!
75//! Repetition levels are more complex.  If there are very large lists then a sparse array of offsets
76//! (which has one element per list) might be more compact than a dense array of repetition levels
77//! (which has one element per list value, possibly even more if there are empty lists).
78//!
79//! However, both repetition levels and definition levels are typically very compressible with RLE.
80//!
81//! However, in Lance we don't always take advantage of that compression because we want to be able
82//! to zip rep-def levels together with our values.  This gives us fewer IOPS when accessing row values.
83//!
84//! # Utilities in this Module
85//!
86//! - `RepDefBuilder` - Extracts validity and offset information from Arrow arrays.  We use this as we
87//!   shred the incoming data into primitive leaf arrays.  We don't immediately convert into rep-def because
88//!   we need to share and cheaply clone the builder when we have structs (each struct child shares some parent
89//!   validity / offset information)  The `serialize` method is called once all data has been received to create
90//!   the final rep-def levels.
91//!
92//! - `SerializerContext` - This is an internal utility that helps with serializing rep-def levels.
93//!
94//! - `CompositeRepDefUnraveler` - This structure is used to reverse the process.  It starts with a set of
95//!   rep-def levels and then uses the `unravel_validity` and `unravel_offsets` methods to produce validity
96//!   buffers and offset buffers.  It is "composite" because we may be combining sets of rep-def buffers from
97//!   multiple locations (e.g. multiple blocks in a mini-block encoded file).
98//!
99//! - `RepDefSlicer` - This is a utility that helps with slicing rep-def buffers.  These buffers are "kind of"
100//!   transparent (maps 1:1 with the values in the array) but not exactly because of the special (empty/null) lists.
101//!   The slicer helps with this issue and is used when slicing a rep-def buffer into mini-blocks.
102//!
103//! - `build_control_word_iterator` - This takes in rep-def levels and returns an iterator that returns byte-padded
104//!   "control words" which are used when creating full-zip encoded data.
105//!
106//! - `ControlWordParser` - This parser can parse the control words returned by `build_control_word_iterator` and is
107//!   used when decoding full-zip encoded data.
108
109use 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/// A top-level-row range whose dense rep/def stream fits one mini-block page.
126#[derive(Debug, Clone, PartialEq, Eq)]
127pub(crate) struct MiniBlockRepDefSplit {
128    /// Top-level row offset, relative to the original unsplit page.
129    pub(crate) row_start: u64,
130    /// Number of top-level rows in this split.
131    pub(crate) num_rows: u64,
132    /// Rep/def level range, relative to the original unsplit page.
133    pub(crate) level_range: Range<usize>,
134    /// Visible value offset, relative to the original unsplit page.
135    pub(crate) value_start: u64,
136    /// Number of visible values in this split.
137    pub(crate) num_values: u64,
138}
139
140/// Dense mini-block rep/def budget result for one accumulated page.
141#[derive(Debug, Clone, PartialEq, Eq)]
142pub(crate) enum MiniBlockRepDefBudget {
143    /// The dense rep/def stream fits one mini-block structural page.
144    WithinBudget,
145    /// The dense rep/def stream fits after splitting on top-level row boundaries.
146    RequiresPageSplit(Vec<MiniBlockRepDefSplit>),
147    /// A single top-level row has this many rep/def levels and exceeds the budget.
148    SingleRowOverBudget(u64),
149}
150
151// As we build def levels we add this to special values to indicate that they
152// are special so that we can skip over them when processing lower levels.
153//
154// We assume 16 bits is good enough for rep-def levels.  This _would_ give
155// us 65536 levels of struct nesting and list nesting.  However, we cut that
156// in half for SPECIAL_THRESHOLD because we use the top bit to indicate if an
157// item is a special value (null list / empty list) during construction.
158//
159// We subtract this off at the end of construction to get the actual definition
160// levels.
161const SPECIAL_THRESHOLD: u16 = u16::MAX / 2;
162
163/// Represents information that we extract from a list array as we are
164/// encoding
165#[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/// Represents validity information that we extract from non-list arrays (that
175/// have nulls) as we are encoding
176#[derive(Clone, Debug)]
177struct ValidityDesc {
178    validity: Option<BooleanBuffer>,
179    num_values: usize,
180}
181
182/// Represents validity information that we extract from FSL arrays.  This is
183/// just validity (no offsets) but we also record the dimension of the FSL array
184/// as that will impact the next layer
185#[derive(Clone, Debug)]
186struct FslDesc {
187    validity: Option<BooleanBuffer>,
188    dimension: usize,
189    num_values: usize,
190}
191
192// As we build up rep/def from arrow arrays we record a
193// series of RawRepDef objects.  Each one corresponds to layer
194// in the array structure
195#[derive(Clone, Debug)]
196enum RawRepDef {
197    Offsets(OffsetDesc),
198    Validity(ValidityDesc),
199    Fsl(FslDesc),
200}
201
202impl RawRepDef {
203    // Are there any nulls in this layer
204    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    // How many values are in this layer
213    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    /// How many empty/null lists are in this layer
222    fn num_specials(&self) -> usize {
223        match self {
224            Self::Offsets(OffsetDesc { num_specials, .. }) => *num_specials,
225            _ => 0,
226        }
227    }
228
229    /// How many definition levels do we need for this layer
230    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    /// How many repetition levels do we need for this layer
254    fn max_rep(&self) -> u16 {
255        match self {
256            Self::Offsets(_) => 1,
257            _ => 0,
258        }
259    }
260}
261
262/// Represents repetition and definition levels that have been
263/// serialized into a pair of (optional) level buffers
264#[derive(Debug)]
265pub struct SerializedRepDefs {
266    /// The repetition levels, one per item
267    ///
268    /// If None, there are no lists
269    pub repetition_levels: Option<Arc<[u16]>>,
270    /// The definition levels, one per item
271    ///
272    /// If None, there are no nulls
273    pub definition_levels: Option<Arc<[u16]>>,
274    /// The meaning of each definition level
275    pub def_meaning: Vec<DefinitionInterpretation>,
276    /// The maximum level that is "visible" from the lowest level
277    ///
278    /// This is the last level before we encounter a list level of some kind.  Once we've
279    /// hit a list level then nulls in any level beyond do not map to actual items.
280    ///
281    /// This is None if there are no lists
282    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    /// Creates an empty SerializedRepDefs (no repetition, all valid)
328    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/// Slices a level buffer into pieces
356///
357/// This is needed to handle the fact that a level buffer may have more
358/// levels than values due to special (empty/null) lists.
359///
360/// As a result, a call to `slice_next(10)` may return 10 levels or it may
361/// return more than 10 levels if any special values are encountered.
362#[derive(Debug)]
363pub struct RepDefSlicer<'a> {
364    repdef: &'a SerializedRepDefs,
365    to_slice: LanceBuffer,
366    current: usize,
367}
368
369// TODO: All of this logic will need some changing when we compress rep/def levels.
370impl<'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    /// Returns the rest of the levels not yet sliced
392    ///
393    /// This must be called instead of `slice_next` on the final iteration.
394    /// This is because anytime we slice there may be empty/null lists on the
395    /// boundary that are "free" and the current behavior in `slice_next` is to
396    /// leave them for the next call.
397    ///
398    /// `slice_rest` will slice all remaining levels and return them.
399    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    /// Returns enough levels to satisfy the next `num_values` values
407    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            // No lists, should be 1:1 mapping from levels to values
411            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            // There are lists and there are def levels.  That means there may be
416            // more rep/def levels than values.  We need to scan the def levels to figure
417            // out which items are "invisible" and skip over them
418            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            // No def levels, should be 1:1 mapping from levels to values
432            self.current = start + num_values;
433            self.to_slice.slice_with_length(start * 2, num_values * 2)
434        }
435    }
436}
437
438/// This tells us how an array handles definition.  Given a stack of
439/// these and a nested array and a set of definition levels we can calculate
440/// how we should interpret the definition levels.
441///
442/// For example, if the interpretation is [AllValidItem, NullableItem] then
443/// a 0 means "valid item" and a 1 means "null struct".  If the interpretation
444/// is [NullableItem, NullableItem] then a 0 means "valid item" and a 1 means
445/// "null item" and a 2 means "null struct".
446///
447/// Lists are tricky because we might use up to two definition levels for a
448/// single layer of list nesting because we need one value to indicate "empty list"
449/// and another value to indicate "null list".
450#[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    /// How many definition levels do we need for this layer
462    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    /// Does this layer have nulls?
474    pub fn is_all_valid(&self) -> bool {
475        matches!(
476            self,
477            Self::AllValidItem | Self::AllValidList | Self::EmptyableList
478        )
479    }
480
481    /// Does this layer represent a list?
482    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/// The RepDefBuilder is used to collect offsets & validity buffers
494/// from arrow structures.  Once we have those we use the SerializerContext
495/// to build the actual repetition and definition levels.
496///
497/// We know ahead of time how many rep/def levels we will need (number of items
498/// in inner-most array + the number of empty/null lists in any parent arrays).
499///
500/// As a result we try and avoid any re-allocations by pre-allocating the buffers
501/// up front.  We allocate two copies of each buffer which allows us to avoid unsafe
502/// code caused by reading and writing to the same buffer (also, it's unavoidable
503/// because there are times we need to write 'faster' than we read)
504#[derive(Debug)]
505struct SerializerContext {
506    // This is built from outer-to-inner and then reversed at the end
507    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        // We write into the spare buffers and read from the active buffers
585        // and then swap at the end.  This way we don't write over what we
586        // are reading.
587
588        let mut new_len = 0;
589        let expected_len = offset_desc.num_values + self.current_num_specials;
590        if expected_len == 0 {
591            // Offsets [0] mean no list values, so no levels.
592            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                // len can't be 0 because then we'd have def levels
602                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                // Copy over any higher-level special values in place
625                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                // If the rep_level is 0 then we are the first list level
637                // otherwise we are starting a higher level list so keep
638                // existing rep level
639                let list_level = if rep == 0 { rep_level } else { rep };
640
641                if def == 0 && len > 0 {
642                    // New valid list, write a rep level and then add new 0/0 items
643                    *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                    // Empty list, insert new special
654                    *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                    // Either the list is null or one of its struct parents
659                    // is null.  Promote it to a special value.
660                    *def_write_itr.next().unwrap() = def + SPECIAL_THRESHOLD;
661                    *rep_write_itr.next().unwrap() = list_level;
662                    new_len += 1;
663                }
664            }
665
666            // If we have any special values at the end, we need to copy them over
667            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        // All non-special values will be broadcasted by the multiplier.  Special values are copied as-is.
719        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            // All valid with no rep/def levels, nothing to do
724            return;
725        } else if self.rep_levels.is_empty() {
726            assert!(self.def_levels.len() >= self.current_len);
727            // No rep levels, just multiply the def levels
728            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            // No def levels, just multiply the rep levels
743            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        // Extremely sparse lists can have many rep/def levels for very few
811        // visible leaf values.  If this ratio becomes too skewed then a
812        // mini-block rep/def chunk can exceed its packed metadata budget even
813        // though the value buffers are small.  We detect that case while
814        // normalizing special def levels and split on top-level row boundaries
815        // so each emitted dense mini-block page stays within the budget.
816        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        // Need to reverse the def meaning since we build rep / def levels in reverse
1015        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        // Need to reverse the def meaning since we build rep / def levels in reverse
1044        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/// A structure used to collect validity buffers and offsets from arrow
1079/// arrays and eventually create repetition and definition levels
1080///
1081/// As we are encoding the structural encoders are given this struct and
1082/// will record the arrow information into it.  Once we hit a leaf node we
1083/// serialize the data into rep/def levels and write these into the page.
1084#[derive(Clone, Default, Debug)]
1085pub struct RepDefBuilder {
1086    // The rep/def info we have collected so far
1087    repdefs: Vec<RawRepDef>,
1088    // The current length, can get larger as we traverse lists (e.g. an
1089    // array might have 5 lists which results in 50 items)
1090    //
1091    // Starts uninitialized until we see the first rep/def item
1092    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            // First validity buffer we've seen
1101            self.len = Some(incoming_len);
1102        }
1103    }
1104
1105    fn num_layers(&self) -> usize {
1106        self.repdefs.len()
1107    }
1108
1109    /// The builder is "empty" if there is no repetition and no nulls.  In this case we don't need
1110    /// to store anything to disk (except the description)
1111    pub fn is_empty(&self) -> bool {
1112        self.repdefs
1113            .iter()
1114            .all(|r| matches!(r, RawRepDef::Validity(ValidityDesc { validity: None, .. })))
1115    }
1116
1117    /// Returns true if there is only a single layer of definition
1118    pub fn is_simple_validity(&self) -> bool {
1119        self.repdefs.len() == 1 && matches!(self.repdefs[0], RawRepDef::Validity(_))
1120    }
1121
1122    /// Registers a nullable validity bitmap
1123    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    /// Registers an all-valid validity layer
1136    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    /// Adds a layer of offsets
1219    ///
1220    /// Offsets are casted to a common type (i64) and also normalized.  Null lists are
1221    /// always represented by a zero-length (identical) pair of offsets and so the caller
1222    /// should filter out any garbage items before encoding them.  To assist with this the
1223    /// method will return true if any non-empty null lists were found.
1224    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    // When we are encoding data it arrives in batches.  For each batch we create a RepDefBuilder and collect the
1244    // various validity buffers and offset buffers from that batch.  Once we have enough batches to write a page we
1245    // need to take this collection of RepDefBuilders and concatenate them and then serialize them into rep/def levels.
1246    //
1247    // TODO: In the future, we may concatenate and serialize at the same time?
1248    //
1249    // This method takes care of the concatenation part.  First we collect all of layer 0 from each builder, then we
1250    // call this method.  Then we collect all of layer 1 from each builder and call this method.  And so on.
1251    //
1252    // That means this method should get a collection of `RawRepDef` where each item is the same kind (all validity or
1253    // all offsets) though the nullability / lengths may be different in each layer.
1254    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        // We make two passes through the layers.  The first determines if we need to pay the cost of allocating
1265        // buffers.  The second pass actually adds the values.
1266        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        // Shortcut if there are no nulls
1298        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        // Only allocate if needed
1318        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    /// Converts the validity / offsets buffers that have been gathered so far
1409    /// into repetition and definition levels
1410    pub fn serialize(builders: Vec<Self>) -> SerializedRepDefs {
1411        Self::serialize_builders(builders).0.build()
1412    }
1413
1414    /// Converts gathered structural buffers into rep/def levels and a mini-block budget result.
1415    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            // No repetition, all-valid
1433            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/// Starts with serialized repetition and definition levels and unravels
1511/// them into validity buffers and offsets buffers
1512///
1513/// This is used during decoding to create the necessary arrow structures
1514#[derive(Debug)]
1515pub struct RepDefUnraveler {
1516    rep_levels: Option<LevelBuffer>,
1517    def_levels: Option<LevelBuffer>,
1518    // Maps from definition level to the rep level at which that definition level is visible
1519    levels_to_rep: Vec<u16>,
1520    def_meaning: Arc<[DefinitionInterpretation]>,
1521    // Current definition level to compare to.
1522    current_def_cmp: u16,
1523    // Current rep level, determines which specials we can see
1524    current_rep_cmp: u16,
1525    // Current layer index, 0 means inner-most layer and it counts up from there.  Used to index
1526    // into special_defs
1527    current_layer: usize,
1528    // Number of items in the inner-most layer (needed if the definition levels are not present)
1529    num_items: u64,
1530}
1531
1532impl RepDefUnraveler {
1533    /// Creates a new unraveler from serialized repetition and definition information
1534    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        // Level=0 is always visible and means valid item
1543        levels_to_rep.push(0);
1544        for meaning in def_meaning.as_ref() {
1545            match meaning {
1546                DefinitionInterpretation::AllValidItem | DefinitionInterpretation::AllValidList => {
1547                    // There is no corresponding level, so nothing to put in levels_to_rep
1548                }
1549                DefinitionInterpretation::NullableItem => {
1550                    // Some null structs are not visible at inner rep levels in cases like LIST<STRUCT<LIST<...>>>
1551                    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    /// If the current level is a repetition layer then this returns the number of lists
1585    /// at this level.
1586    ///
1587    /// This is not valid to call when the current level is a struct/primitive layer because
1588    /// in some cases there may be no rep or def information to know this.
1589    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            // Worst case every rep item is max_rep and a new list
1596            .map(|levels| levels.len())
1597            .unwrap_or(0)
1598    }
1599
1600    /// Unravels a layer of offsets from the unraveler into the given offset width
1601    ///
1602    /// When decoding a list the caller should first unravel the offsets and then
1603    /// unravel the validity (this is the opposite order used during encoding)
1604    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        // This is the highest def level that is still visible.  Once we hit a list then
1633        // we stop looking because any null / empty list (or list masked by a higher level
1634        // null) will not be visible
1635        let mut max_level = null_level.max(empty_level).max(valid_level);
1636        // Anything higher than this (but less than max_level) is a null struct masking our
1637        // list.  We will materialize this is a null list.
1638        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        // If offsets is empty this is a no-op.  If offsets is not empty that means we already
1654        // added a set of offsets.  For example, we might have added [0, 3, 5] (2 lists).  Now
1655        // say we want to add [0, 1, 4] (2 lists).  We should get [0, 3, 5, 6, 9] (4 lists).  If
1656        // we don't pop here we get [0, 3, 5, 5, 6, 9] which is wrong.
1657        //
1658        // Or, to think about it another way, if every unraveler adds the starting 0 and the trailing
1659        // length then we have N + unravelers.len() values instead of N + 1.
1660        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            // It's possible validity is None even if we have def levels.  For example, we might have
1670            // empty lists (which require def levels) but no nulls.
1671            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            // This is a strange access pattern.  We are iterating over the rep/def levels and
1677            // at the same time writing the rep/def levels.  This means we need both a mutable
1678            // and immutable reference to the rep/def levels.
1679            let mut read_idx = 0;
1680            let mut write_idx = 0;
1681            while read_idx < rep_levels.len() {
1682                // SAFETY: We assert that rep_levels and def_levels have the same
1683                // len and read_idx and write_idx can never go past the end.
1684                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                        // Copy over
1689                        *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                            // This is a valid list
1695                            offsets.push(to_offset(curlen)?);
1696                            curlen += 1;
1697                            push_validity(true);
1698                        } else if def_val > max_level {
1699                            // This is not visible at this rep level, do not add to offsets, but keep in repdef
1700                        } else if def_val == null_level || def_val > upper_null {
1701                            // This is a null list (or a list masked by a null struct)
1702                            offsets.push(to_offset(curlen)?);
1703                            push_validity(false);
1704                        } else if def_val == empty_level {
1705                            // This is an empty list
1706                            offsets.push(to_offset(curlen)?);
1707                            push_validity(true);
1708                        } else {
1709                            // New valid list starting with null item
1710                            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            // SAFETY: See above loop
1726            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                // SAFETY: read_idx / write_idx cannot go past rep_levels.len()
1731                unsafe {
1732                    let rep_val = *rep_levels.get_unchecked(read_idx);
1733                    if rep_val != 0 {
1734                        // Finish the current list
1735                        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            // Truncate to the number of lists THIS unraveler produced (write_idx),
1746            // not `offsets.len() - 1` — the latter includes offsets contributed by
1747            // earlier unravelers in a multi-page read, which would leave too many
1748            // rep levels for the next (outer) layer and over-count its lists.
1749            rep_levels.truncate(write_idx);
1750            if let Some(validity) = validity {
1751                // Even though we don't have validity it is possible another unraveler did and so we need
1752                // to push all valids
1753                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    /// Unravels a layer of validity from the definition levels
1765    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            // If we need to support this then I think we need to walk through the rep def levels to find
1793            // the spots at which we keep.  E.g. if we have:
1794            //  rep: 1 0 0 1 0 1 0 0 0 1 0 0
1795            //  def: 1 1 1 0 1 0 1 1 0 1 1 0
1796            //  dimension: 2
1797            //
1798            // The output should be:
1799            //  rep: 1 0 0 1 0 0 0
1800            //  def: 1 1 1 0 1 1 0
1801            //
1802            // Maybe there's some special logic for empty/null lists?  I'll save the headache for future me.
1803            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/// As we decode we may extract rep/def information from multiple pages (or multiple
1822/// chunks within a page).
1823///
1824/// For each chunk we create an unraveler.  Each unraveler can have a completely different
1825/// interpretation (e.g. one page might contain null items but no null structs and the next
1826/// page might have null structs but no null items).
1827///
1828/// Concatenating these unravelers would be tricky and expensive so instead we have a
1829/// composite unraveler which unravels across multiple unravelers.
1830///
1831/// Note: this class should be used even if there is only one page / unraveler.  This is
1832/// because the `RepDefUnraveler`'s API is more complex (it's meant to be called by this
1833/// class)
1834#[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    /// Unravels a layer of validity
1845    ///
1846    /// Returns None if there are no null items in this layer
1847    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    /// Unravels a layer of offsets (and the validity for that layer)
1879    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            // Note: This is probably an over-estimate and potentially even an under-estimate.  We only know
1893            // right now how many items we have and not how many rows.  (TODO: Shouldn't we know the # of rows?)
1894            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/// A [`ControlWordIterator`] when there are both repetition and definition levels
1911///
1912/// The iterator will put the repetition level in the upper bits and the definition
1913/// level in the lower bits.  The number of bits used for each level is determined
1914/// by the width of the repetition and definition levels.
1915#[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/// A [`ControlWordIterator`] when there are only definition levels or only repetition levels
1986#[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            // Either there is no rep, in which case there are no invisible items
2005            // or there is no def, in which case there are no invisible items
2006            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/// A [`ControlWordIterator`] when there are no repetition or definition levels
2048#[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
2069/// Helper function to get a bit mask of the given width
2070fn get_mask(width: u16) -> u16 {
2071    (1 << width) - 1
2072}
2073
2074// We're really going out of our way to avoid boxing here but this will be called on a per-value basis
2075// so it is in the critical path.
2076type SpecificBinaryControlWordIterator<'a, T> = BinaryControlWordIterator<
2077    Zip<Copied<std::slice::Iter<'a, u16>>, Copied<std::slice::Iter<'a, u16>>>,
2078    T,
2079>;
2080
2081/// An iterator that generates control words from repetition and definition levels
2082///
2083/// "Control word" is just a fancy term for a single u8/u16/u32 that contains both
2084/// the repetition and definition in it.
2085///
2086/// In the large majority of case we only need a single byte to represent both the
2087/// repetition and definition levels.  However, if there is deep nesting then we may
2088/// need two bytes.  In the worst case we need 4 bytes though this suggests hundreds of
2089/// levels of nesting which seems unlikely to encounter in practice.
2090#[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/// Describes the properties of a control word
2102#[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    /// Appends the next control word to the buffer
2111    ///
2112    /// Returns true if this is the start of a new item (i.e. the repetition level is maxed out)
2113    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    /// Return true if the control word iterator has repetition levels
2126    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    /// Returns the number of bytes per control word
2137    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    /// Returns the number of bits used for the repetition level
2150    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    /// Returns the number of bits used for the definition level
2163    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
2176/// Builds a [`ControlWordIterator`] from repetition and definition levels
2177/// by first calculating the width needed and then creating the iterator
2178/// with the appropriate width
2179pub 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/// A parser to unwrap control words into repetition and definition levels
2309///
2310/// This is the inverse of the [`ControlWordIterator`].
2311#[derive(Copy, Clone, Debug)]
2312pub enum ControlWordParser {
2313    // First item is the bits to shift, second is the mask to apply (the mask can be
2314    // calculated from the bits to shift but we don't want to calculate it each time)
2315    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    /// Returns the number of bytes per control word
2473    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    /// Appends the next control word to the rep & def buffers
2489    ///
2490    /// `src` should be pointing at the first byte (little endian) of the control word
2491    ///
2492    /// `dst_rep` and `dst_def` are the buffers to append the rep and def levels to.
2493    /// They will not be appended to if not needed.
2494    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    /// Return true if the control words contain repetition information
2516    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    /// Temporarily parses the control word to inspect its properties but does not append to any buffers
2529    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    /// Creates a new parser from the number of bits used for the repetition and definition levels
2567    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        // Empty offsets should serialize without panicking.
2626        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        // Basic case, rep & def
2636        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        // [[I], [I, I]], NULL, [[NULL, NULL], NULL, [NULL, I, I, NULL]]
2657
2658        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        // Note: validity doesn't exactly round-trip because repdef normalizes some of the
2666        // redundant validity values
2667        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        // Null list and empty list should be serialized mostly the same
2696
2697        // Null case
2698        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        // Empty case
2710        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        // Regresses a failure we encountered when the last item was an empty list
2722        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        // List nulls are allowed to have non-empty offsets and garbage values
2745        // and the add_offsets call should normalize this
2746        let mut builder = RepDefBuilder::default();
2747        builder.add_offsets(
2748            offsets_32(&[0, 2, 5, 8]),
2749            Some(validity(&[true, false, true])),
2750        );
2751        // Note: we pass 5 here and not 8.  If add_offsets tells us there is garbage nulls they
2752        // should be removed before continuing
2753        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        // Sliced lists may have offsets that don't start with zero.  The
2860        // add_offsets call needs to normalize these to operate correctly.
2861        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        // Tests when we have some empty lists but no null lists.  This case
2910        // caused some bugs because we have definition but no nulls
2911        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        // Single builder: List<List<i32>>, 3 rows.
2968        //   outer [0,2,3,5] -> rows have 2,1,2 inner lists
2969        //   inner [0,1,3,5,7,9] -> 5 inner lists, lengths 1,2,2,2,2 (9 leaf)
2970        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        // Same logical data split into two batches:
2979        //   batch0 = rows 0,1 : outer [0,2,3], inner [0,1,3,5] (3 inner, 5 leaf)
2980        //   batch1 = row 2    : outer [0,2],   inner [0,2,4]   (2 inner, 4 leaf)
2981        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        // This regresses a bug
3089
3090        let rep = vec![1, 0, 0, 0];
3091        let def = vec![2, 0, 2, 2];
3092        // AllValidList<NullableStruct<NullableItem>>
3093        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        // Basic case, rep & def
3205        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        // First 5 items include a null list so we get 6 levels (12 bytes)
3261        assert_eq!(rep_slicer.slice_next(5).len(), 12);
3262        // Next 20 are all plain
3263        assert_eq!(rep_slicer.slice_next(20).len(), 40);
3264        // Last 5 include an empty list so we get 6 levels (12 bytes)
3265        assert_eq!(rep_slicer.slice_rest().len(), 12);
3266
3267        let mut def_slicer = repdefs.rep_slicer().unwrap();
3268
3269        // First 5 items include a null list so we get 6 levels (12 bytes)
3270        assert_eq!(def_slicer.slice_next(5).len(), 12);
3271        // Next 20 are all plain
3272        assert_eq!(def_slicer.slice_next(20).len(), 40);
3273        // Last 5 include an empty list so we get 6 levels (12 bytes)
3274        assert_eq!(def_slicer.slice_rest().len(), 12);
3275    }
3276
3277    #[test]
3278    fn test_control_words() {
3279        // Convert to control words, verify expected, convert back, verify same as original
3280        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        // Each will need 4 bits and so we should get 1-byte control words
3331        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, // 0, 5
3335            0b01110011, // 7, 3
3336            0b00110001, // 3, 1
3337            0b00100010, // 2, 2
3338            0b10011100, // 9, 12
3339            0b10001111, // 8, 15
3340            0b11000000, // 12, 0
3341            0b01010010, // 5, 2
3342        ];
3343        check(rep, def, expected, 1, 4, 4);
3344
3345        // Now we need 5 bits for def so we get 2-byte control words
3346        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, // 0, 5
3350            0b11100011, 0b00000000, // 7, 3
3351            0b01100001, 0b00000000, // 3, 1
3352            0b01000010, 0b00000000, // 2, 2
3353            0b00101100, 0b00000001, // 9, 12
3354            0b00010110, 0b00000001, // 8, 22
3355            0b10000000, 0b00000001, // 12, 0
3356            0b10100010, 0b00000000, // 5, 2
3357        ];
3358        check(rep, def, expected, 2, 4, 5);
3359
3360        // Just rep, 4 bits so 1 byte each
3361        let levels = &[0_u16, 7, 3, 2, 9, 8, 12, 5];
3362        let expected = vec![
3363            0b00000000, // 0
3364            0b00000111, // 7
3365            0b00000011, // 3
3366            0b00000010, // 2
3367            0b00001001, // 9
3368            0b00001000, // 8
3369            0b00001100, // 12
3370            0b00000101, // 5
3371        ];
3372        check(levels, &[], expected.clone(), 1, 4, 0);
3373
3374        // Just def
3375        check(&[], levels, expected, 1, 0, 4);
3376
3377        // No rep, no def, no bytes
3378        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                /*max_visible_def=*/ 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        // 2 means new list
3417        let rep = &[2_u16, 1, 0, 2, 2, 0, 1, 1, 0, 2, 0];
3418        // These values don't matter for this test
3419        let def = &[0_u16, 0, 0, 3, 1, 1, 2, 1, 0, 0, 1];
3420
3421        // Rep & def
3422        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        // Rep only
3433        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        // No repetition
3442        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        // No repetition, no definition
3451        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        // This regresses a case where we had 3 null lists inside a struct
3464        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        // This tests cases where the validity is different between two different pages
3530        // because one page has nulls and the other doesn't.
3531
3532        // Simple case with one layer of validity and no repetition
3533        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        // More complex case with two layers of validity and repetition
3556        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        // A page can keep nullable layer metadata even when all definition levels are 0
3601        // and no definition buffer needs to be materialized. This should decode as all-valid.
3602        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}