Skip to main content

parquet/file/
serialized_reader.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! Contains implementations of the reader traits FileReader, RowGroupReader and PageReader
19//! Also contains implementations of the ChunkReader for files (with buffering) and byte arrays (RAM)
20
21use crate::basic::{PageType, Type};
22use crate::bloom_filter::Sbbf;
23use crate::column::page::{Page, PageMetadata, PageReader};
24use crate::compression::{Codec, create_codec};
25#[cfg(feature = "encryption")]
26use crate::encryption::decrypt::{CryptoContext, read_and_decrypt};
27use crate::errors::{ParquetError, Result};
28use crate::file::metadata::thrift::PageHeader;
29use crate::file::page_index::offset_index::{OffsetIndexMetaData, PageLocation};
30use crate::file::statistics;
31use crate::file::{
32    metadata::*,
33    properties::{ReaderProperties, ReaderPropertiesPtr},
34    reader::*,
35};
36#[cfg(feature = "encryption")]
37use crate::parquet_thrift::ThriftSliceInputProtocol;
38use crate::parquet_thrift::{ReadThrift, ThriftReadInputProtocol};
39use crate::record::Row;
40use crate::record::reader::RowIter;
41use crate::schema::types::{SchemaDescPtr, Type as SchemaType};
42use bytes::Bytes;
43use std::collections::VecDeque;
44use std::{fs::File, io::Read, path::Path, sync::Arc};
45
46impl TryFrom<File> for SerializedFileReader<File> {
47    type Error = ParquetError;
48
49    fn try_from(file: File) -> Result<Self> {
50        Self::new(file)
51    }
52}
53
54impl TryFrom<&Path> for SerializedFileReader<File> {
55    type Error = ParquetError;
56
57    fn try_from(path: &Path) -> Result<Self> {
58        let file = File::open(path)?;
59        Self::try_from(file)
60    }
61}
62
63impl TryFrom<String> for SerializedFileReader<File> {
64    type Error = ParquetError;
65
66    fn try_from(path: String) -> Result<Self> {
67        Self::try_from(Path::new(&path))
68    }
69}
70
71impl TryFrom<&str> for SerializedFileReader<File> {
72    type Error = ParquetError;
73
74    fn try_from(path: &str) -> Result<Self> {
75        Self::try_from(Path::new(&path))
76    }
77}
78
79/// Conversion into a [`RowIter`]
80/// using the full file schema over all row groups.
81impl IntoIterator for SerializedFileReader<File> {
82    type Item = Result<Row>;
83    type IntoIter = RowIter<'static>;
84
85    fn into_iter(self) -> Self::IntoIter {
86        RowIter::from_file_into(Box::new(self))
87    }
88}
89
90// ----------------------------------------------------------------------
91// Implementations of file & row group readers
92
93/// A serialized implementation for Parquet [`FileReader`].
94pub struct SerializedFileReader<R: ChunkReader> {
95    chunk_reader: Arc<R>,
96    metadata: Arc<ParquetMetaData>,
97    props: ReaderPropertiesPtr,
98}
99
100/// A predicate for filtering row groups, invoked with the metadata and index
101/// of each row group in the file. Only row groups for which the predicate
102/// evaluates to `true` will be scanned
103pub type ReadGroupPredicate = Box<dyn FnMut(&RowGroupMetaData, usize) -> bool>;
104
105/// A builder for [`ReadOptions`].
106/// For the predicates that are added to the builder,
107/// they will be chained using 'AND' to filter the row groups.
108#[derive(Default)]
109pub struct ReadOptionsBuilder {
110    predicates: Vec<ReadGroupPredicate>,
111    enable_page_index: bool,
112    props: Option<ReaderProperties>,
113    metadata_options: ParquetMetaDataOptions,
114}
115
116impl ReadOptionsBuilder {
117    /// New builder
118    pub fn new() -> Self {
119        Self::default()
120    }
121
122    /// Add a predicate on row group metadata to the reading option,
123    /// Filter only row groups that match the predicate criteria
124    pub fn with_predicate(mut self, predicate: ReadGroupPredicate) -> Self {
125        self.predicates.push(predicate);
126        self
127    }
128
129    /// Add a range predicate on filtering row groups if their midpoints are within
130    /// the Closed-Open range `[start..end) {x | start <= x < end}`
131    ///
132    /// # Panics
133    ///
134    /// Panics if `end <= start`
135    pub fn with_range(mut self, start: i64, end: i64) -> Self {
136        assert!(start < end);
137        let predicate = move |rg: &RowGroupMetaData, _: usize| {
138            let mid = get_midpoint_offset(rg);
139            mid >= start && mid < end
140        };
141        self.predicates.push(Box::new(predicate));
142        self
143    }
144
145    /// Enable reading the page index structures described in
146    /// "[Column Index] Layout to Support Page Skipping"
147    ///
148    /// [Column Index]: https://github.com/apache/parquet-format/blob/master/PageIndex.md
149    pub fn with_page_index(mut self) -> Self {
150        self.enable_page_index = true;
151        self
152    }
153
154    /// Set the [`ReaderProperties`] configuration.
155    pub fn with_reader_properties(mut self, properties: ReaderProperties) -> Self {
156        self.props = Some(properties);
157        self
158    }
159
160    /// Provide a Parquet schema to use when decoding the metadata. The schema in the Parquet
161    /// footer will be skipped.
162    pub fn with_parquet_schema(mut self, schema: SchemaDescPtr) -> Self {
163        self.metadata_options.set_schema(schema);
164        self
165    }
166
167    /// Set whether to convert the [`encoding_stats`] in the Parquet `ColumnMetaData` to a bitmask
168    /// (defaults to `false`).
169    ///
170    /// See [`ColumnChunkMetaData::page_encoding_stats_mask`] for an explanation of why this
171    /// might be desirable.
172    ///
173    /// [`encoding_stats`]:
174    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L917
175    pub fn with_encoding_stats_as_mask(mut self, val: bool) -> Self {
176        self.metadata_options.set_encoding_stats_as_mask(val);
177        self
178    }
179
180    /// Sets the decoding policy for [`encoding_stats`] in the Parquet `ColumnMetaData`.
181    ///
182    /// [`encoding_stats`]:
183    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L917
184    pub fn with_encoding_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
185        self.metadata_options.set_encoding_stats_policy(policy);
186        self
187    }
188
189    /// Sets the decoding policy for [`statistics`] in the Parquet `ColumnMetaData`.
190    ///
191    /// [`statistics`]:
192    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L912
193    pub fn with_column_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
194        self.metadata_options.set_column_stats_policy(policy);
195        self
196    }
197
198    /// Sets the decoding policy for [`size_statistics`] in the Parquet `ColumnMetaData`.
199    ///
200    /// [`size_statistics`]:
201    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L936
202    pub fn with_size_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
203        self.metadata_options.set_size_stats_policy(policy);
204        self
205    }
206
207    /// Seal the builder and return the read options
208    pub fn build(self) -> ReadOptions {
209        let props = self
210            .props
211            .unwrap_or_else(|| ReaderProperties::builder().build());
212        ReadOptions {
213            predicates: self.predicates,
214            enable_page_index: self.enable_page_index,
215            props,
216            metadata_options: self.metadata_options,
217        }
218    }
219}
220
221/// A collection of options for reading a Parquet file.
222///
223/// Predicates are currently only supported on row group metadata.
224/// All predicates will be chained using 'AND' to filter the row groups.
225pub struct ReadOptions {
226    predicates: Vec<ReadGroupPredicate>,
227    enable_page_index: bool,
228    props: ReaderProperties,
229    metadata_options: ParquetMetaDataOptions,
230}
231
232impl<R: 'static + ChunkReader> SerializedFileReader<R> {
233    /// Creates file reader from a Parquet file.
234    /// Returns an error if the Parquet file does not exist or is corrupt.
235    pub fn new(chunk_reader: R) -> Result<Self> {
236        let metadata = ParquetMetaDataReader::new().parse_and_finish(&chunk_reader)?;
237        let props = Arc::new(ReaderProperties::builder().build());
238        Ok(Self {
239            chunk_reader: Arc::new(chunk_reader),
240            metadata: Arc::new(metadata),
241            props,
242        })
243    }
244
245    /// Creates file reader from a Parquet file with read options.
246    /// Returns an error if the Parquet file does not exist or is corrupt.
247    pub fn new_with_options(chunk_reader: R, options: ReadOptions) -> Result<Self> {
248        let mut metadata_builder = ParquetMetaDataReader::new()
249            .with_metadata_options(Some(options.metadata_options.clone()))
250            .parse_and_finish(&chunk_reader)?
251            .into_builder();
252        let mut predicates = options.predicates;
253
254        // Filter row groups based on the predicates
255        for (i, rg_meta) in metadata_builder.take_row_groups().into_iter().enumerate() {
256            let mut keep = true;
257            for predicate in &mut predicates {
258                if !predicate(&rg_meta, i) {
259                    keep = false;
260                    break;
261                }
262            }
263            if keep {
264                metadata_builder = metadata_builder.add_row_group(rg_meta);
265            }
266        }
267
268        let mut metadata = metadata_builder.build();
269
270        // If page indexes are desired, build them with the filtered set of row groups
271        if options.enable_page_index {
272            let mut reader = ParquetMetaDataReader::new_with_metadata(metadata)
273                .with_page_index_policy(PageIndexPolicy::Required);
274            reader.read_page_indexes(&chunk_reader)?;
275            metadata = reader.finish()?;
276        }
277
278        Ok(Self {
279            chunk_reader: Arc::new(chunk_reader),
280            metadata: Arc::new(metadata),
281            props: Arc::new(options.props),
282        })
283    }
284}
285
286/// Get midpoint offset for a row group
287fn get_midpoint_offset(meta: &RowGroupMetaData) -> i64 {
288    let col = meta.column(0);
289    let mut offset = col.data_page_offset();
290    if let Some(dic_offset) = col.dictionary_page_offset()
291        && offset > dic_offset
292    {
293        offset = dic_offset
294    }
295    offset + meta.compressed_size() / 2
296}
297
298impl<R: 'static + ChunkReader> FileReader for SerializedFileReader<R> {
299    fn metadata(&self) -> &ParquetMetaData {
300        &self.metadata
301    }
302
303    fn num_row_groups(&self) -> usize {
304        self.metadata.num_row_groups()
305    }
306
307    fn get_row_group(&self, i: usize) -> Result<Box<dyn RowGroupReader + '_>> {
308        let row_group_metadata = self.metadata.row_group(i);
309        // Row groups should be processed sequentially.
310        let props = Arc::clone(&self.props);
311        let f = Arc::clone(&self.chunk_reader);
312        Ok(Box::new(SerializedRowGroupReader::new(
313            f,
314            row_group_metadata,
315            self.metadata
316                .page_index()
317                .map(|pi| pi.offset_indexes_for_rowgroup(i))
318                .unwrap_or(None),
319            props,
320        )?))
321    }
322
323    fn get_row_iter(&self, projection: Option<SchemaType>) -> Result<RowIter<'_>> {
324        RowIter::from_file(projection, self)
325    }
326}
327
328/// A serialized implementation for Parquet [`RowGroupReader`].
329pub struct SerializedRowGroupReader<'a, R: ChunkReader> {
330    chunk_reader: Arc<R>,
331    metadata: &'a RowGroupMetaData,
332    offset_index: Option<&'a [Option<OffsetIndexMetaData>]>,
333    props: ReaderPropertiesPtr,
334    bloom_filters: Vec<Option<Sbbf>>,
335}
336
337impl<'a, R: ChunkReader> SerializedRowGroupReader<'a, R> {
338    /// Creates new row group reader from a file, row group metadata and custom config.
339    pub fn new(
340        chunk_reader: Arc<R>,
341        metadata: &'a RowGroupMetaData,
342        offset_index: Option<&'a [Option<OffsetIndexMetaData>]>,
343        props: ReaderPropertiesPtr,
344    ) -> Result<Self> {
345        let bloom_filters = if props.read_bloom_filter() {
346            metadata
347                .columns()
348                .iter()
349                .map(|col| Sbbf::read_from_column_chunk(col, &*chunk_reader))
350                .collect::<Result<Vec<_>>>()?
351        } else {
352            std::iter::repeat_n(None, metadata.columns().len()).collect()
353        };
354        Ok(Self {
355            chunk_reader,
356            metadata,
357            offset_index,
358            props,
359            bloom_filters,
360        })
361    }
362}
363
364impl<R: 'static + ChunkReader> RowGroupReader for SerializedRowGroupReader<'_, R> {
365    fn metadata(&self) -> &RowGroupMetaData {
366        self.metadata
367    }
368
369    fn num_columns(&self) -> usize {
370        self.metadata.num_columns()
371    }
372
373    // TODO: fix PARQUET-816
374    fn get_column_page_reader(&self, i: usize) -> Result<Box<dyn PageReader>> {
375        let col = self.metadata.column(i);
376
377        let page_locations = if let Some(offset_index) = self.offset_index {
378            offset_index[i].as_ref().map(|oi| oi.page_locations.clone())
379        } else {
380            None
381        };
382
383        let props = Arc::clone(&self.props);
384        Ok(Box::new(SerializedPageReader::new_with_properties(
385            Arc::clone(&self.chunk_reader),
386            col,
387            usize::try_from(self.metadata.num_rows())?,
388            page_locations,
389            props,
390        )?))
391    }
392
393    /// get bloom filter for the `i`th column
394    fn get_column_bloom_filter(&self, i: usize) -> Option<&Sbbf> {
395        self.bloom_filters[i].as_ref()
396    }
397
398    fn get_row_iter(&self, projection: Option<SchemaType>) -> Result<RowIter<'_>> {
399        RowIter::from_row_group(projection, self)
400    }
401}
402
403/// Decodes a [`Page`] from the provided `buffer`
404pub(crate) fn decode_page(
405    page_header: PageHeader,
406    buffer: Bytes,
407    physical_type: Type,
408    decompressor: Option<&mut Box<dyn Codec>>,
409) -> Result<Page> {
410    // Verify the 32-bit CRC checksum of the page
411    #[cfg(feature = "crc")]
412    if let Some(expected_crc) = page_header.crc {
413        let crc = crc32fast::hash(&buffer);
414        if crc != expected_crc as u32 {
415            return Err(general_err!("Page CRC checksum mismatch"));
416        }
417    }
418
419    // When processing data page v2, depending on enabled compression for the
420    // page, we should account for uncompressed data ('offset') of
421    // repetition and definition levels.
422    //
423    // We always use 0 offset for other pages other than v2, `true` flag means
424    // that compression will be applied if decompressor is defined
425    let (offset, can_decompress): (usize, bool) = match page_header.data_page_header_v2 {
426        Some(ref header_v2) => {
427            if header_v2.definition_levels_byte_length < 0
428                || header_v2.repetition_levels_byte_length < 0
429                || header_v2.definition_levels_byte_length + header_v2.repetition_levels_byte_length
430                    > page_header.uncompressed_page_size
431            {
432                return Err(general_err!(
433                    "DataPage v2 header contains implausible values \
434                        for definition_levels_byte_length ({}) \
435                        and repetition_levels_byte_length ({}) \
436                        given DataPage header provides uncompressed_page_size ({})",
437                    header_v2.definition_levels_byte_length,
438                    header_v2.repetition_levels_byte_length,
439                    page_header.uncompressed_page_size
440                ));
441            }
442            (
443                usize::try_from(
444                    header_v2.definition_levels_byte_length
445                        + header_v2.repetition_levels_byte_length,
446                )?,
447                // When is_compressed flag is missing the page is considered compressed
448                header_v2.is_compressed.unwrap_or(true),
449            )
450        }
451        None => (0, true),
452    };
453
454    let buffer = match decompressor {
455        Some(decompressor) if can_decompress => {
456            let uncompressed_page_size = usize::try_from(page_header.uncompressed_page_size)?;
457            if offset > buffer.len() || offset > uncompressed_page_size {
458                return Err(general_err!("Invalid page header"));
459            }
460            let decompressed_size = uncompressed_page_size - offset;
461            let mut decompressed = Vec::with_capacity(uncompressed_page_size);
462            decompressed.extend_from_slice(&buffer[..offset]);
463            // decompressed size of zero corresponds to a page with no non-null values
464            // see https://github.com/apache/parquet-format/blob/master/README.md#data-pages
465            if decompressed_size > 0 {
466                let compressed = &buffer[offset..];
467                decompressor.decompress(compressed, &mut decompressed, Some(decompressed_size))?;
468            }
469
470            if decompressed.len() != uncompressed_page_size {
471                return Err(general_err!(
472                    "Actual decompressed size doesn't match the expected one ({} vs {})",
473                    decompressed.len(),
474                    uncompressed_page_size
475                ));
476            }
477
478            Bytes::from(decompressed)
479        }
480        _ => buffer,
481    };
482
483    let result = match page_header.r#type {
484        PageType::DICTIONARY_PAGE => {
485            let dict_header = page_header.dictionary_page_header.as_ref().ok_or_else(|| {
486                ParquetError::General("Missing dictionary page header".to_string())
487            })?;
488            let is_sorted = dict_header.is_sorted.unwrap_or(false);
489            Page::DictionaryPage {
490                buf: buffer,
491                num_values: dict_header.num_values.try_into()?,
492                encoding: dict_header.encoding,
493                is_sorted,
494            }
495        }
496        PageType::DATA_PAGE => {
497            let header = page_header
498                .data_page_header
499                .ok_or_else(|| ParquetError::General("Missing V1 data page header".to_string()))?;
500            Page::DataPage {
501                buf: buffer,
502                num_values: header.num_values.try_into()?,
503                encoding: header.encoding,
504                def_level_encoding: header.definition_level_encoding,
505                rep_level_encoding: header.repetition_level_encoding,
506                statistics: statistics::from_thrift_page_stats(physical_type, header.statistics)?,
507            }
508        }
509        PageType::DATA_PAGE_V2 => {
510            let header = page_header
511                .data_page_header_v2
512                .ok_or_else(|| ParquetError::General("Missing V2 data page header".to_string()))?;
513            let is_compressed = header.is_compressed.unwrap_or(true);
514            Page::DataPageV2 {
515                buf: buffer,
516                num_values: header.num_values.try_into()?,
517                encoding: header.encoding,
518                num_nulls: header.num_nulls.try_into()?,
519                num_rows: header.num_rows.try_into()?,
520                def_levels_byte_len: header.definition_levels_byte_length.try_into()?,
521                rep_levels_byte_len: header.repetition_levels_byte_length.try_into()?,
522                is_compressed,
523                statistics: statistics::from_thrift_page_stats(physical_type, header.statistics)?,
524            }
525        }
526        PageType::INDEX_PAGE => {
527            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
528            return Err(general_err!(
529                "Page type {:?} is not supported",
530                page_header.r#type
531            ));
532        }
533    };
534
535    Ok(result)
536}
537
538enum SerializedPageReaderState {
539    Values {
540        /// The current byte offset in the reader
541        /// Note that offset is u64 (i.e., not usize) to support 32-bit architectures such as WASM
542        offset: u64,
543
544        /// The length of the chunk in bytes
545        /// Note that remaining_bytes is u64 (i.e., not usize) to support 32-bit architectures such as WASM
546        remaining_bytes: u64,
547
548        // If the next page header has already been "peeked", we will cache it and it`s length here
549        next_page_header: Option<Box<PageHeader>>,
550
551        /// The index of the data page within this column chunk
552        page_index: usize,
553
554        /// Whether the next page is expected to be a dictionary page
555        require_dictionary: bool,
556    },
557    Pages {
558        /// Remaining page locations
559        page_locations: VecDeque<PageLocation>,
560        /// Remaining dictionary location if any
561        dictionary_page: Option<PageLocation>,
562        /// The total number of rows in this column chunk
563        total_rows: usize,
564        /// The index of the data page within this column chunk
565        page_index: usize,
566    },
567}
568
569#[derive(Default)]
570struct SerializedPageReaderContext {
571    /// Controls decoding of page-level statistics
572    read_stats: bool,
573    /// Crypto context carrying objects required for decryption
574    #[cfg(feature = "encryption")]
575    crypto_context: Option<Arc<CryptoContext>>,
576}
577
578/// A serialized implementation for Parquet [`PageReader`].
579pub struct SerializedPageReader<R: ChunkReader> {
580    /// The chunk reader
581    reader: Arc<R>,
582
583    /// The compression codec for this column chunk. Only set for non-PLAIN codec.
584    decompressor: Option<Box<dyn Codec>>,
585
586    /// Column chunk type.
587    physical_type: Type,
588
589    state: SerializedPageReaderState,
590
591    context: SerializedPageReaderContext,
592}
593
594impl<R: ChunkReader> SerializedPageReader<R> {
595    /// Creates a new serialized page reader from a chunk reader and metadata
596    pub fn new(
597        reader: Arc<R>,
598        column_chunk_metadata: &ColumnChunkMetaData,
599        total_rows: usize,
600        page_locations: Option<Vec<PageLocation>>,
601    ) -> Result<Self> {
602        let props = Arc::new(ReaderProperties::builder().build());
603        SerializedPageReader::new_with_properties(
604            reader,
605            column_chunk_metadata,
606            total_rows,
607            page_locations,
608            props,
609        )
610    }
611
612    /// Stub No-op implementation when encryption is disabled.
613    #[cfg(all(feature = "arrow", not(feature = "encryption")))]
614    pub(crate) fn add_crypto_context(
615        self,
616        _rg_idx: usize,
617        _column_idx: usize,
618        _parquet_meta_data: &ParquetMetaData,
619        _column_chunk_metadata: &ColumnChunkMetaData,
620    ) -> Result<SerializedPageReader<R>> {
621        Ok(self)
622    }
623
624    /// Adds any necessary crypto context to this page reader, if encryption is enabled.
625    #[cfg(feature = "encryption")]
626    pub(crate) fn add_crypto_context(
627        mut self,
628        rg_idx: usize,
629        column_idx: usize,
630        parquet_meta_data: &ParquetMetaData,
631        column_chunk_metadata: &ColumnChunkMetaData,
632    ) -> Result<SerializedPageReader<R>> {
633        let Some(file_decryptor) = parquet_meta_data.file_decryptor() else {
634            return Ok(self);
635        };
636        let Some(crypto_metadata) = column_chunk_metadata.crypto_metadata() else {
637            return Ok(self);
638        };
639        let crypto_context =
640            CryptoContext::for_column(file_decryptor, crypto_metadata, rg_idx, column_idx)?;
641        self.context.crypto_context = Some(Arc::new(crypto_context));
642        Ok(self)
643    }
644
645    /// Creates a new serialized page with custom options.
646    pub fn new_with_properties(
647        reader: Arc<R>,
648        meta: &ColumnChunkMetaData,
649        total_rows: usize,
650        page_locations: Option<Vec<PageLocation>>,
651        props: ReaderPropertiesPtr,
652    ) -> Result<Self> {
653        let decompressor = create_codec(meta.compression(), props.codec_options())?;
654        let (start, len) = meta.byte_range();
655
656        let state = match page_locations {
657            Some(locations) => {
658                // If the offset of the first page doesn't match the start of the column chunk
659                // then the preceding space must contain a dictionary page.
660                let dictionary_page = match locations.first() {
661                    Some(dict_offset) if dict_offset.offset as u64 != start => Some(PageLocation {
662                        offset: start as i64,
663                        compressed_page_size: (dict_offset.offset as u64 - start) as i32,
664                        first_row_index: 0,
665                    }),
666                    _ => None,
667                };
668
669                SerializedPageReaderState::Pages {
670                    page_locations: locations.into(),
671                    dictionary_page,
672                    total_rows,
673                    page_index: 0,
674                }
675            }
676            None => SerializedPageReaderState::Values {
677                offset: start,
678                remaining_bytes: len,
679                next_page_header: None,
680                page_index: 0,
681                require_dictionary: meta.dictionary_page_offset().is_some(),
682            },
683        };
684        let mut context = SerializedPageReaderContext::default();
685        if props.read_page_stats() {
686            context.read_stats = true;
687        }
688        Ok(Self {
689            reader,
690            decompressor,
691            state,
692            physical_type: meta.column_type(),
693            context,
694        })
695    }
696
697    /// Similar to `peek_next_page`, but returns the offset of the next page instead of the page metadata.
698    /// Unlike page metadata, an offset can uniquely identify a page.
699    ///
700    /// This is used when we need to read parquet with row-filter, and we don't want to decompress the page twice.
701    /// This function allows us to check if the next page is being cached or read previously.
702    #[cfg(test)]
703    fn peek_next_page_offset(&mut self) -> Result<Option<u64>> {
704        match &mut self.state {
705            SerializedPageReaderState::Values {
706                offset,
707                remaining_bytes,
708                next_page_header,
709                page_index,
710                require_dictionary,
711            } => {
712                loop {
713                    if *remaining_bytes == 0 {
714                        return Ok(None);
715                    }
716                    return if let Some(header) = next_page_header.as_ref() {
717                        if let Ok(_page_meta) = PageMetadata::try_from(&**header) {
718                            Ok(Some(*offset))
719                        } else {
720                            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
721                            *next_page_header = None;
722                            continue;
723                        }
724                    } else {
725                        let mut read = self.reader.get_read(*offset)?;
726                        let (header_len, header) = Self::read_page_header_len(
727                            &self.context,
728                            &mut read,
729                            *page_index,
730                            *require_dictionary,
731                        )?;
732                        *offset += header_len as u64;
733                        *remaining_bytes -= header_len as u64;
734                        let page_meta = if let Ok(_page_meta) = PageMetadata::try_from(&header) {
735                            Ok(Some(*offset))
736                        } else {
737                            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
738                            continue;
739                        };
740                        *next_page_header = Some(Box::new(header));
741                        page_meta
742                    };
743                }
744            }
745            SerializedPageReaderState::Pages {
746                page_locations,
747                dictionary_page,
748                ..
749            } => {
750                if let Some(page) = dictionary_page {
751                    Ok(Some(page.offset as u64))
752                } else if let Some(page) = page_locations.front() {
753                    Ok(Some(page.offset as u64))
754                } else {
755                    Ok(None)
756                }
757            }
758        }
759    }
760
761    fn read_page_header_len<T: Read>(
762        context: &SerializedPageReaderContext,
763        input: &mut T,
764        page_index: usize,
765        dictionary_page: bool,
766    ) -> Result<(usize, PageHeader)> {
767        /// A wrapper around a [`std::io::Read`] that keeps track of the bytes read
768        struct TrackedRead<R> {
769            inner: R,
770            bytes_read: usize,
771        }
772
773        impl<R: Read> Read for TrackedRead<R> {
774            fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
775                let v = self.inner.read(buf)?;
776                self.bytes_read += v;
777                Ok(v)
778            }
779        }
780
781        let mut tracked = TrackedRead {
782            inner: input,
783            bytes_read: 0,
784        };
785        let header = context.read_page_header(&mut tracked, page_index, dictionary_page)?;
786        Ok((tracked.bytes_read, header))
787    }
788
789    fn read_page_header_len_from_bytes(
790        context: &SerializedPageReaderContext,
791        buffer: &[u8],
792        page_index: usize,
793        dictionary_page: bool,
794    ) -> Result<(usize, PageHeader)> {
795        let mut input = std::io::Cursor::new(buffer);
796        let header = context.read_page_header(&mut input, page_index, dictionary_page)?;
797        let header_len = input.position() as usize;
798        Ok((header_len, header))
799    }
800}
801
802#[cfg(not(feature = "encryption"))]
803impl SerializedPageReaderContext {
804    fn read_page_header<T: Read>(
805        &self,
806        input: &mut T,
807        _page_index: usize,
808        _dictionary_page: bool,
809    ) -> Result<PageHeader> {
810        let mut prot = ThriftReadInputProtocol::new(input);
811        if self.read_stats {
812            Ok(PageHeader::read_thrift(&mut prot)?)
813        } else {
814            Ok(PageHeader::read_thrift_without_stats(&mut prot)?)
815        }
816    }
817
818    fn decrypt_page_data<T>(
819        &self,
820        buffer: T,
821        _page_index: usize,
822        _dictionary_page: bool,
823    ) -> Result<T> {
824        Ok(buffer)
825    }
826}
827
828#[cfg(feature = "encryption")]
829impl SerializedPageReaderContext {
830    fn read_page_header<T: Read>(
831        &self,
832        input: &mut T,
833        page_index: usize,
834        dictionary_page: bool,
835    ) -> Result<PageHeader> {
836        match self.page_crypto_context(page_index, dictionary_page) {
837            None => {
838                let mut prot = ThriftReadInputProtocol::new(input);
839                if self.read_stats {
840                    Ok(PageHeader::read_thrift(&mut prot)?)
841                } else {
842                    use crate::file::metadata::thrift::PageHeader;
843
844                    Ok(PageHeader::read_thrift_without_stats(&mut prot)?)
845                }
846            }
847            Some(page_crypto_context) => {
848                let data_decryptor = page_crypto_context.data_decryptor();
849                let aad = page_crypto_context.create_page_header_aad()?;
850
851                let buf = read_and_decrypt(data_decryptor, input, aad.as_ref()).map_err(|_| {
852                    ParquetError::General(format!(
853                        "Error decrypting page header for column {}, decryption key may be wrong",
854                        page_crypto_context.column_ordinal
855                    ))
856                })?;
857
858                let mut prot = ThriftSliceInputProtocol::new(buf.as_slice());
859                if self.read_stats {
860                    Ok(PageHeader::read_thrift(&mut prot)?)
861                } else {
862                    Ok(PageHeader::read_thrift_without_stats(&mut prot)?)
863                }
864            }
865        }
866    }
867
868    fn decrypt_page_data<T>(&self, buffer: T, page_index: usize, dictionary_page: bool) -> Result<T>
869    where
870        T: AsRef<[u8]>,
871        T: From<Vec<u8>>,
872    {
873        let page_crypto_context = self.page_crypto_context(page_index, dictionary_page);
874        if let Some(page_crypto_context) = page_crypto_context {
875            let decryptor = page_crypto_context.data_decryptor();
876            let aad = page_crypto_context.create_page_aad()?;
877            let decrypted = decryptor.decrypt(buffer.as_ref(), &aad)?;
878            Ok(T::from(decrypted))
879        } else {
880            Ok(buffer)
881        }
882    }
883
884    fn page_crypto_context(
885        &self,
886        page_index: usize,
887        dictionary_page: bool,
888    ) -> Option<Arc<CryptoContext>> {
889        self.crypto_context.as_ref().map(|c| {
890            Arc::new(if dictionary_page {
891                c.for_dictionary_page()
892            } else {
893                c.with_page_ordinal(page_index)
894            })
895        })
896    }
897}
898
899impl<R: ChunkReader> Iterator for SerializedPageReader<R> {
900    type Item = Result<Page>;
901
902    fn next(&mut self) -> Option<Self::Item> {
903        self.get_next_page().transpose()
904    }
905}
906
907fn verify_page_header_len(header_len: usize, remaining_bytes: u64) -> Result<()> {
908    if header_len as u64 > remaining_bytes {
909        return Err(eof_err!("Invalid page header"));
910    }
911    Ok(())
912}
913
914fn verify_page_size(
915    compressed_size: i32,
916    uncompressed_size: i32,
917    remaining_bytes: u64,
918) -> Result<()> {
919    // The page's compressed size should not exceed the remaining bytes that are
920    // available to read. The page's uncompressed size is the expected size
921    // after decompression, which can never be negative.
922    if compressed_size < 0 || compressed_size as u64 > remaining_bytes || uncompressed_size < 0 {
923        return Err(eof_err!("Invalid page header"));
924    }
925    Ok(())
926}
927
928impl<R: ChunkReader> PageReader for SerializedPageReader<R> {
929    fn get_next_page(&mut self) -> Result<Option<Page>> {
930        loop {
931            let page = match &mut self.state {
932                SerializedPageReaderState::Values {
933                    offset,
934                    remaining_bytes: remaining,
935                    next_page_header,
936                    page_index,
937                    require_dictionary,
938                } => {
939                    if *remaining == 0 {
940                        return Ok(None);
941                    }
942
943                    let mut read = self.reader.get_read(*offset)?;
944                    let header = if let Some(header) = next_page_header.take() {
945                        *header
946                    } else {
947                        let (header_len, header) = Self::read_page_header_len(
948                            &self.context,
949                            &mut read,
950                            *page_index,
951                            *require_dictionary,
952                        )?;
953                        verify_page_header_len(header_len, *remaining)?;
954                        *offset += header_len as u64;
955                        *remaining -= header_len as u64;
956                        header
957                    };
958                    verify_page_size(
959                        header.compressed_page_size,
960                        header.uncompressed_page_size,
961                        *remaining,
962                    )?;
963                    let data_len = header.compressed_page_size as usize;
964                    let data_start = *offset;
965                    *offset += data_len as u64;
966                    *remaining -= data_len as u64;
967
968                    if header.r#type == PageType::INDEX_PAGE {
969                        continue;
970                    }
971
972                    let buffer = self.reader.get_bytes(data_start, data_len)?;
973
974                    let buffer =
975                        self.context
976                            .decrypt_page_data(buffer, *page_index, *require_dictionary)?;
977
978                    let page = decode_page(
979                        header,
980                        buffer,
981                        self.physical_type,
982                        self.decompressor.as_mut(),
983                    )?;
984                    if page.is_data_page() {
985                        *page_index += 1;
986                    } else if page.is_dictionary_page() {
987                        *require_dictionary = false;
988                    }
989                    page
990                }
991                SerializedPageReaderState::Pages {
992                    page_locations,
993                    dictionary_page,
994                    page_index,
995                    ..
996                } => {
997                    let (front, is_dictionary_page) = match dictionary_page.take() {
998                        Some(front) => (front, true),
999                        None => match page_locations.pop_front() {
1000                            Some(front) => (front, false),
1001                            None => return Ok(None),
1002                        },
1003                    };
1004
1005                    let page_len = usize::try_from(front.compressed_page_size)?;
1006                    let buffer = self.reader.get_bytes(front.offset as u64, page_len)?;
1007
1008                    let (offset, header) = Self::read_page_header_len_from_bytes(
1009                        &self.context,
1010                        buffer.as_ref(),
1011                        *page_index,
1012                        is_dictionary_page,
1013                    )?;
1014                    let bytes = buffer.slice(offset..);
1015                    let bytes =
1016                        self.context
1017                            .decrypt_page_data(bytes, *page_index, is_dictionary_page)?;
1018
1019                    if !is_dictionary_page {
1020                        *page_index += 1;
1021                    }
1022                    decode_page(
1023                        header,
1024                        bytes,
1025                        self.physical_type,
1026                        self.decompressor.as_mut(),
1027                    )?
1028                }
1029            };
1030
1031            return Ok(Some(page));
1032        }
1033    }
1034
1035    fn peek_next_page(&mut self) -> Result<Option<PageMetadata>> {
1036        match &mut self.state {
1037            SerializedPageReaderState::Values {
1038                offset,
1039                remaining_bytes,
1040                next_page_header,
1041                page_index,
1042                require_dictionary,
1043            } => {
1044                loop {
1045                    if *remaining_bytes == 0 {
1046                        return Ok(None);
1047                    }
1048                    return if let Some(header) = next_page_header.as_ref() {
1049                        if let Ok(page_meta) = (&**header).try_into() {
1050                            Ok(Some(page_meta))
1051                        } else {
1052                            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
1053                            *next_page_header = None;
1054                            continue;
1055                        }
1056                    } else {
1057                        let mut read = self.reader.get_read(*offset)?;
1058                        let (header_len, header) = Self::read_page_header_len(
1059                            &self.context,
1060                            &mut read,
1061                            *page_index,
1062                            *require_dictionary,
1063                        )?;
1064                        verify_page_header_len(header_len, *remaining_bytes)?;
1065                        *offset += header_len as u64;
1066                        *remaining_bytes -= header_len as u64;
1067                        let page_meta = if let Ok(page_meta) = (&header).try_into() {
1068                            Ok(Some(page_meta))
1069                        } else {
1070                            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
1071                            continue;
1072                        };
1073                        *next_page_header = Some(Box::new(header));
1074                        page_meta
1075                    };
1076                }
1077            }
1078            SerializedPageReaderState::Pages {
1079                page_locations,
1080                dictionary_page,
1081                total_rows,
1082                page_index: _,
1083            } => {
1084                if dictionary_page.is_some() {
1085                    Ok(Some(PageMetadata {
1086                        num_rows: None,
1087                        num_levels: None,
1088                        is_dict: true,
1089                    }))
1090                } else if let Some(page) = page_locations.front() {
1091                    let next_rows = page_locations
1092                        .get(1)
1093                        .map(|x| x.first_row_index as usize)
1094                        .unwrap_or(*total_rows);
1095
1096                    Ok(Some(PageMetadata {
1097                        num_rows: Some(next_rows - page.first_row_index as usize),
1098                        num_levels: None,
1099                        is_dict: false,
1100                    }))
1101                } else {
1102                    Ok(None)
1103                }
1104            }
1105        }
1106    }
1107
1108    fn skip_next_page(&mut self) -> Result<()> {
1109        match &mut self.state {
1110            SerializedPageReaderState::Values {
1111                offset,
1112                remaining_bytes,
1113                next_page_header,
1114                page_index,
1115                require_dictionary,
1116            } => {
1117                if let Some(buffered_header) = next_page_header.take() {
1118                    verify_page_size(
1119                        buffered_header.compressed_page_size,
1120                        buffered_header.uncompressed_page_size,
1121                        *remaining_bytes,
1122                    )?;
1123                    // The next page header has already been peeked, so just advance the offset
1124                    *offset += buffered_header.compressed_page_size as u64;
1125                    *remaining_bytes -= buffered_header.compressed_page_size as u64;
1126                } else {
1127                    let mut read = self.reader.get_read(*offset)?;
1128                    let (header_len, header) = Self::read_page_header_len(
1129                        &self.context,
1130                        &mut read,
1131                        *page_index,
1132                        *require_dictionary,
1133                    )?;
1134                    verify_page_header_len(header_len, *remaining_bytes)?;
1135                    verify_page_size(
1136                        header.compressed_page_size,
1137                        header.uncompressed_page_size,
1138                        *remaining_bytes,
1139                    )?;
1140                    let data_page_size = header.compressed_page_size as u64;
1141                    *offset += header_len as u64 + data_page_size;
1142                    *remaining_bytes -= header_len as u64 + data_page_size;
1143                }
1144                if *require_dictionary {
1145                    *require_dictionary = false;
1146                } else {
1147                    *page_index += 1;
1148                }
1149                Ok(())
1150            }
1151            SerializedPageReaderState::Pages {
1152                page_locations,
1153                dictionary_page,
1154                page_index,
1155                ..
1156            } => {
1157                if dictionary_page.is_some() {
1158                    // If a dictionary page exists, consume it by taking it (sets to None)
1159                    dictionary_page.take();
1160                } else {
1161                    // If no dictionary page exists, simply pop the data page from page_locations
1162                    if page_locations.pop_front().is_some() {
1163                        *page_index += 1;
1164                    }
1165                }
1166
1167                Ok(())
1168            }
1169        }
1170    }
1171
1172    fn at_record_boundary(&mut self) -> Result<bool> {
1173        match &mut self.state {
1174            SerializedPageReaderState::Values { .. } => match self.peek_next_page()? {
1175                None => Ok(true),
1176                // V2 data pages must start at record boundaries per the parquet
1177                // spec, so the current page ends at one.
1178                Some(metadata) => Ok(metadata.num_rows.is_some()),
1179            },
1180            SerializedPageReaderState::Pages { .. } => Ok(true),
1181        }
1182    }
1183}
1184
1185#[cfg(test)]
1186mod tests {
1187    use std::collections::HashSet;
1188
1189    use bytes::Buf;
1190
1191    use crate::file::page_index::column_index::{
1192        ByteArrayColumnIndex, ColumnIndexMetaData, PrimitiveColumnIndex,
1193    };
1194    use crate::file::properties::{EnabledStatistics, WriterProperties};
1195
1196    use crate::basic::{self, BoundaryOrder, ColumnOrder, Encoding, SortOrder};
1197    use crate::column::reader::ColumnReader;
1198    use crate::data_type::private::ParquetValueType;
1199    use crate::data_type::{AsBytes, FixedLenByteArrayType, Int32Type};
1200    use crate::file::metadata::thrift::DataPageHeaderV2;
1201    use crate::file::writer::SerializedFileWriter;
1202    use crate::record::RowAccessor;
1203    use crate::schema::parser::parse_message_type;
1204    use crate::util::test_common::file_util::{get_test_file, get_test_path};
1205
1206    use super::*;
1207
1208    #[test]
1209    fn test_decode_page_invalid_offset() {
1210        let page_header = PageHeader {
1211            r#type: PageType::DATA_PAGE_V2,
1212            uncompressed_page_size: 10,
1213            compressed_page_size: 10,
1214            data_page_header: None,
1215            index_page_header: None,
1216            dictionary_page_header: None,
1217            crc: None,
1218            data_page_header_v2: Some(DataPageHeaderV2 {
1219                num_nulls: 0,
1220                num_rows: 0,
1221                num_values: 0,
1222                encoding: Encoding::PLAIN,
1223                definition_levels_byte_length: 11,
1224                repetition_levels_byte_length: 0,
1225                is_compressed: None,
1226                statistics: None,
1227            }),
1228        };
1229
1230        let buffer = Bytes::new();
1231        let err = decode_page(page_header, buffer, Type::INT32, None).unwrap_err();
1232        assert!(
1233            err.to_string()
1234                .contains("DataPage v2 header contains implausible values")
1235        );
1236    }
1237
1238    #[test]
1239    fn test_decode_unsupported_page() {
1240        let mut page_header = PageHeader {
1241            r#type: PageType::INDEX_PAGE,
1242            uncompressed_page_size: 10,
1243            compressed_page_size: 10,
1244            data_page_header: None,
1245            index_page_header: None,
1246            dictionary_page_header: None,
1247            crc: None,
1248            data_page_header_v2: None,
1249        };
1250        let buffer = Bytes::new();
1251        let err = decode_page(page_header.clone(), buffer.clone(), Type::INT32, None).unwrap_err();
1252        assert_eq!(
1253            err.to_string(),
1254            "Parquet error: Page type INDEX_PAGE is not supported"
1255        );
1256
1257        page_header.data_page_header_v2 = Some(DataPageHeaderV2 {
1258            num_nulls: 0,
1259            num_rows: 0,
1260            num_values: 0,
1261            encoding: Encoding::PLAIN,
1262            definition_levels_byte_length: 11,
1263            repetition_levels_byte_length: 0,
1264            is_compressed: None,
1265            statistics: None,
1266        });
1267        let err = decode_page(page_header, buffer, Type::INT32, None).unwrap_err();
1268        assert!(
1269            err.to_string()
1270                .contains("DataPage v2 header contains implausible values")
1271        );
1272    }
1273
1274    #[test]
1275    fn test_cursor_and_file_has_the_same_behaviour() {
1276        let mut buf: Vec<u8> = Vec::new();
1277        get_test_file("alltypes_plain.parquet")
1278            .read_to_end(&mut buf)
1279            .unwrap();
1280        let cursor = Bytes::from(buf);
1281        let read_from_cursor = SerializedFileReader::new(cursor).unwrap();
1282
1283        let test_file = get_test_file("alltypes_plain.parquet");
1284        let read_from_file = SerializedFileReader::new(test_file).unwrap();
1285
1286        let file_iter = read_from_file.get_row_iter(None).unwrap();
1287        let cursor_iter = read_from_cursor.get_row_iter(None).unwrap();
1288
1289        for (a, b) in file_iter.zip(cursor_iter) {
1290            assert_eq!(a.unwrap(), b.unwrap())
1291        }
1292    }
1293
1294    #[test]
1295    fn test_file_reader_try_from() {
1296        // Valid file path
1297        let test_file = get_test_file("alltypes_plain.parquet");
1298        let test_path_buf = get_test_path("alltypes_plain.parquet");
1299        let test_path = test_path_buf.as_path();
1300        let test_path_str = test_path.to_str().unwrap();
1301
1302        let reader = SerializedFileReader::try_from(test_file);
1303        assert!(reader.is_ok());
1304
1305        let reader = SerializedFileReader::try_from(test_path);
1306        assert!(reader.is_ok());
1307
1308        let reader = SerializedFileReader::try_from(test_path_str);
1309        assert!(reader.is_ok());
1310
1311        let reader = SerializedFileReader::try_from(test_path_str.to_string());
1312        assert!(reader.is_ok());
1313
1314        // Invalid file path
1315        let test_path = Path::new("invalid.parquet");
1316        let test_path_str = test_path.to_str().unwrap();
1317
1318        let reader = SerializedFileReader::try_from(test_path);
1319        assert!(reader.is_err());
1320
1321        let reader = SerializedFileReader::try_from(test_path_str);
1322        assert!(reader.is_err());
1323
1324        let reader = SerializedFileReader::try_from(test_path_str.to_string());
1325        assert!(reader.is_err());
1326    }
1327
1328    #[test]
1329    fn test_file_reader_into_iter() {
1330        let path = get_test_path("alltypes_plain.parquet");
1331        let reader = SerializedFileReader::try_from(path.as_path()).unwrap();
1332        let iter = reader.into_iter();
1333        let values: Vec<_> = iter.flat_map(|x| x.unwrap().get_int(0)).collect();
1334
1335        assert_eq!(values, &[4, 5, 6, 7, 2, 3, 0, 1]);
1336    }
1337
1338    #[test]
1339    fn test_file_reader_into_iter_project() {
1340        let path = get_test_path("alltypes_plain.parquet");
1341        let reader = SerializedFileReader::try_from(path.as_path()).unwrap();
1342        let schema = "message schema { OPTIONAL INT32 id; }";
1343        let proj = parse_message_type(schema).ok();
1344        let iter = reader.into_iter().project(proj).unwrap();
1345        let values: Vec<_> = iter.flat_map(|x| x.unwrap().get_int(0)).collect();
1346
1347        assert_eq!(values, &[4, 5, 6, 7, 2, 3, 0, 1]);
1348    }
1349
1350    #[test]
1351    fn test_reuse_file_chunk() {
1352        // This test covers the case of maintaining the correct start position in a file
1353        // stream for each column reader after initializing and moving to the next one
1354        // (without necessarily reading the entire column).
1355        let test_file = get_test_file("alltypes_plain.parquet");
1356        let reader = SerializedFileReader::new(test_file).unwrap();
1357        let row_group = reader.get_row_group(0).unwrap();
1358
1359        let mut page_readers = Vec::new();
1360        for i in 0..row_group.num_columns() {
1361            page_readers.push(row_group.get_column_page_reader(i).unwrap());
1362        }
1363
1364        // Now buffer each col reader, we do not expect any failures like:
1365        // General("underlying Thrift error: end of file")
1366        for mut page_reader in page_readers {
1367            assert!(page_reader.get_next_page().is_ok());
1368        }
1369    }
1370
1371    #[test]
1372    fn test_file_reader() {
1373        let test_file = get_test_file("alltypes_plain.parquet");
1374        let reader_result = SerializedFileReader::new(test_file);
1375        assert!(reader_result.is_ok());
1376        let reader = reader_result.unwrap();
1377
1378        // Test contents in Parquet metadata
1379        let metadata = reader.metadata();
1380        assert_eq!(metadata.num_row_groups(), 1);
1381
1382        // Test contents in file metadata
1383        let file_metadata = metadata.file_metadata();
1384        assert!(file_metadata.created_by().is_some());
1385        assert_eq!(
1386            file_metadata.created_by().unwrap(),
1387            "impala version 1.3.0-INTERNAL (build 8a48ddb1eff84592b3fc06bc6f51ec120e1fffc9)"
1388        );
1389        assert!(file_metadata.key_value_metadata().is_none());
1390        assert_eq!(file_metadata.num_rows(), 8);
1391        assert_eq!(file_metadata.version(), 1);
1392        assert_eq!(file_metadata.column_orders(), None);
1393
1394        // Test contents in row group metadata
1395        let row_group_metadata = metadata.row_group(0);
1396        assert_eq!(row_group_metadata.num_columns(), 11);
1397        assert_eq!(row_group_metadata.num_rows(), 8);
1398        assert_eq!(row_group_metadata.total_byte_size(), 671);
1399        // Check each column order
1400        for i in 0..row_group_metadata.num_columns() {
1401            assert_eq!(file_metadata.column_order(i), ColumnOrder::UNDEFINED);
1402        }
1403
1404        // Test row group reader
1405        let row_group_reader_result = reader.get_row_group(0);
1406        assert!(row_group_reader_result.is_ok());
1407        let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1408        assert_eq!(
1409            row_group_reader.num_columns(),
1410            row_group_metadata.num_columns()
1411        );
1412        assert_eq!(
1413            row_group_reader.metadata().total_byte_size(),
1414            row_group_metadata.total_byte_size()
1415        );
1416
1417        // Test page readers
1418        // TODO: test for every column
1419        let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1420        assert!(page_reader_0_result.is_ok());
1421        let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1422        let mut page_count = 0;
1423        while let Some(page) = page_reader_0.get_next_page().unwrap() {
1424            let is_expected_page = match page {
1425                Page::DictionaryPage {
1426                    buf,
1427                    num_values,
1428                    encoding,
1429                    is_sorted,
1430                } => {
1431                    assert_eq!(buf.len(), 32);
1432                    assert_eq!(num_values, 8);
1433                    assert_eq!(encoding, Encoding::PLAIN_DICTIONARY);
1434                    assert!(!is_sorted);
1435                    true
1436                }
1437                Page::DataPage {
1438                    buf,
1439                    num_values,
1440                    encoding,
1441                    def_level_encoding,
1442                    rep_level_encoding,
1443                    statistics,
1444                } => {
1445                    assert_eq!(buf.len(), 11);
1446                    assert_eq!(num_values, 8);
1447                    assert_eq!(encoding, Encoding::PLAIN_DICTIONARY);
1448                    assert_eq!(def_level_encoding, Encoding::RLE);
1449                    #[expect(deprecated)]
1450                    let expected_rep_level_encoding = Encoding::BIT_PACKED;
1451                    assert_eq!(rep_level_encoding, expected_rep_level_encoding);
1452                    assert!(statistics.is_none());
1453                    true
1454                }
1455                Page::DataPageV2 { .. } => false,
1456            };
1457            assert!(is_expected_page);
1458            page_count += 1;
1459        }
1460        assert_eq!(page_count, 2);
1461    }
1462
1463    #[test]
1464    fn test_file_reader_datapage_v2() {
1465        let test_file = get_test_file("datapage_v2.snappy.parquet");
1466        let reader_result = SerializedFileReader::new(test_file);
1467        assert!(reader_result.is_ok());
1468        let reader = reader_result.unwrap();
1469
1470        // Test contents in Parquet metadata
1471        let metadata = reader.metadata();
1472        assert_eq!(metadata.num_row_groups(), 1);
1473
1474        // Test contents in file metadata
1475        let file_metadata = metadata.file_metadata();
1476        assert!(file_metadata.created_by().is_some());
1477        assert_eq!(
1478            file_metadata.created_by().unwrap(),
1479            "parquet-mr version 1.8.1 (build 4aba4dae7bb0d4edbcf7923ae1339f28fd3f7fcf)"
1480        );
1481        assert!(file_metadata.key_value_metadata().is_some());
1482        assert_eq!(file_metadata.key_value_metadata().unwrap().len(), 1);
1483
1484        assert_eq!(file_metadata.num_rows(), 5);
1485        assert_eq!(file_metadata.version(), 1);
1486        assert_eq!(file_metadata.column_orders(), None);
1487
1488        let row_group_metadata = metadata.row_group(0);
1489
1490        // Check each column order
1491        for i in 0..row_group_metadata.num_columns() {
1492            assert_eq!(file_metadata.column_order(i), ColumnOrder::UNDEFINED);
1493        }
1494
1495        // Test row group reader
1496        let row_group_reader_result = reader.get_row_group(0);
1497        assert!(row_group_reader_result.is_ok());
1498        let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1499        assert_eq!(
1500            row_group_reader.num_columns(),
1501            row_group_metadata.num_columns()
1502        );
1503        assert_eq!(
1504            row_group_reader.metadata().total_byte_size(),
1505            row_group_metadata.total_byte_size()
1506        );
1507
1508        // Test page readers
1509        // TODO: test for every column
1510        let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1511        assert!(page_reader_0_result.is_ok());
1512        let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1513        let mut page_count = 0;
1514        while let Some(page) = page_reader_0.get_next_page().unwrap() {
1515            let is_expected_page = match page {
1516                Page::DictionaryPage {
1517                    buf,
1518                    num_values,
1519                    encoding,
1520                    is_sorted,
1521                } => {
1522                    assert_eq!(buf.len(), 7);
1523                    assert_eq!(num_values, 1);
1524                    assert_eq!(encoding, Encoding::PLAIN);
1525                    assert!(!is_sorted);
1526                    true
1527                }
1528                Page::DataPageV2 {
1529                    buf,
1530                    num_values,
1531                    encoding,
1532                    num_nulls,
1533                    num_rows,
1534                    def_levels_byte_len,
1535                    rep_levels_byte_len,
1536                    is_compressed,
1537                    statistics,
1538                } => {
1539                    assert_eq!(buf.len(), 4);
1540                    assert_eq!(num_values, 5);
1541                    assert_eq!(encoding, Encoding::RLE_DICTIONARY);
1542                    assert_eq!(num_nulls, 1);
1543                    assert_eq!(num_rows, 5);
1544                    assert_eq!(def_levels_byte_len, 2);
1545                    assert_eq!(rep_levels_byte_len, 0);
1546                    assert!(is_compressed);
1547                    assert!(statistics.is_none()); // page stats are no longer read
1548                    true
1549                }
1550                Page::DataPage { .. } => false,
1551            };
1552            assert!(is_expected_page);
1553            page_count += 1;
1554        }
1555        assert_eq!(page_count, 2);
1556    }
1557
1558    #[cfg_attr(miri, ignore)] // calls native Zstd code unsupported by Miri
1559    #[test]
1560    fn test_file_reader_empty_compressed_datapage_v2() {
1561        // this file has a compressed datapage that un-compresses to 0 bytes
1562        let test_file = get_test_file("page_v2_empty_compressed.parquet");
1563        let reader_result = SerializedFileReader::new(test_file);
1564        assert!(reader_result.is_ok());
1565        let reader = reader_result.unwrap();
1566
1567        // Test contents in Parquet metadata
1568        let metadata = reader.metadata();
1569        assert_eq!(metadata.num_row_groups(), 1);
1570
1571        // Test contents in file metadata
1572        let file_metadata = metadata.file_metadata();
1573        assert!(file_metadata.created_by().is_some());
1574        assert_eq!(
1575            file_metadata.created_by().unwrap(),
1576            "parquet-cpp-arrow version 14.0.2"
1577        );
1578        assert!(file_metadata.key_value_metadata().is_some());
1579        assert_eq!(file_metadata.key_value_metadata().unwrap().len(), 1);
1580
1581        assert_eq!(file_metadata.num_rows(), 10);
1582        assert_eq!(file_metadata.version(), 2);
1583        let expected_order = ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED);
1584        assert_eq!(
1585            file_metadata.column_orders(),
1586            Some(vec![expected_order].as_ref())
1587        );
1588
1589        let row_group_metadata = metadata.row_group(0);
1590
1591        // Check each column order
1592        for i in 0..row_group_metadata.num_columns() {
1593            assert_eq!(file_metadata.column_order(i), expected_order);
1594        }
1595
1596        // Test row group reader
1597        let row_group_reader_result = reader.get_row_group(0);
1598        assert!(row_group_reader_result.is_ok());
1599        let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1600        assert_eq!(
1601            row_group_reader.num_columns(),
1602            row_group_metadata.num_columns()
1603        );
1604        assert_eq!(
1605            row_group_reader.metadata().total_byte_size(),
1606            row_group_metadata.total_byte_size()
1607        );
1608
1609        // Test page readers
1610        let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1611        assert!(page_reader_0_result.is_ok());
1612        let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1613        let mut page_count = 0;
1614        while let Some(page) = page_reader_0.get_next_page().unwrap() {
1615            let is_expected_page = match page {
1616                Page::DictionaryPage {
1617                    buf,
1618                    num_values,
1619                    encoding,
1620                    is_sorted,
1621                } => {
1622                    assert_eq!(buf.len(), 0);
1623                    assert_eq!(num_values, 0);
1624                    assert_eq!(encoding, Encoding::PLAIN);
1625                    assert!(!is_sorted);
1626                    true
1627                }
1628                Page::DataPageV2 {
1629                    buf,
1630                    num_values,
1631                    encoding,
1632                    num_nulls,
1633                    num_rows,
1634                    def_levels_byte_len,
1635                    rep_levels_byte_len,
1636                    is_compressed,
1637                    statistics,
1638                } => {
1639                    assert_eq!(buf.len(), 3);
1640                    assert_eq!(num_values, 10);
1641                    assert_eq!(encoding, Encoding::RLE_DICTIONARY);
1642                    assert_eq!(num_nulls, 10);
1643                    assert_eq!(num_rows, 10);
1644                    assert_eq!(def_levels_byte_len, 2);
1645                    assert_eq!(rep_levels_byte_len, 0);
1646                    assert!(is_compressed);
1647                    assert!(statistics.is_none()); // page stats are no longer read
1648                    true
1649                }
1650                Page::DataPage { .. } => false,
1651            };
1652            assert!(is_expected_page);
1653            page_count += 1;
1654        }
1655        assert_eq!(page_count, 2);
1656    }
1657
1658    #[test]
1659    fn test_file_reader_empty_datapage_v2() {
1660        // this file has 0 bytes compressed datapage that un-compresses to 0 bytes
1661        let test_file = get_test_file("datapage_v2_empty_datapage.snappy.parquet");
1662        let reader_result = SerializedFileReader::new(test_file);
1663        assert!(reader_result.is_ok());
1664        let reader = reader_result.unwrap();
1665
1666        // Test contents in Parquet metadata
1667        let metadata = reader.metadata();
1668        assert_eq!(metadata.num_row_groups(), 1);
1669
1670        // Test contents in file metadata
1671        let file_metadata = metadata.file_metadata();
1672        assert!(file_metadata.created_by().is_some());
1673        assert_eq!(
1674            file_metadata.created_by().unwrap(),
1675            "parquet-mr version 1.13.1 (build db4183109d5b734ec5930d870cdae161e408ddba)"
1676        );
1677        assert!(file_metadata.key_value_metadata().is_some());
1678        assert_eq!(file_metadata.key_value_metadata().unwrap().len(), 2);
1679
1680        assert_eq!(file_metadata.num_rows(), 1);
1681        assert_eq!(file_metadata.version(), 1);
1682        let expected_order = ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED);
1683        assert_eq!(
1684            file_metadata.column_orders(),
1685            Some(vec![expected_order].as_ref())
1686        );
1687
1688        let row_group_metadata = metadata.row_group(0);
1689
1690        // Check each column order
1691        for i in 0..row_group_metadata.num_columns() {
1692            assert_eq!(file_metadata.column_order(i), expected_order);
1693        }
1694
1695        // Test row group reader
1696        let row_group_reader_result = reader.get_row_group(0);
1697        assert!(row_group_reader_result.is_ok());
1698        let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1699        assert_eq!(
1700            row_group_reader.num_columns(),
1701            row_group_metadata.num_columns()
1702        );
1703        assert_eq!(
1704            row_group_reader.metadata().total_byte_size(),
1705            row_group_metadata.total_byte_size()
1706        );
1707
1708        // Test page readers
1709        let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1710        assert!(page_reader_0_result.is_ok());
1711        let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1712        let mut page_count = 0;
1713        while let Some(page) = page_reader_0.get_next_page().unwrap() {
1714            let is_expected_page = match page {
1715                Page::DataPageV2 {
1716                    buf,
1717                    num_values,
1718                    encoding,
1719                    num_nulls,
1720                    num_rows,
1721                    def_levels_byte_len,
1722                    rep_levels_byte_len,
1723                    is_compressed,
1724                    statistics,
1725                } => {
1726                    assert_eq!(buf.len(), 2);
1727                    assert_eq!(num_values, 1);
1728                    assert_eq!(encoding, Encoding::PLAIN);
1729                    assert_eq!(num_nulls, 1);
1730                    assert_eq!(num_rows, 1);
1731                    assert_eq!(def_levels_byte_len, 2);
1732                    assert_eq!(rep_levels_byte_len, 0);
1733                    assert!(is_compressed);
1734                    assert!(statistics.is_none());
1735                    true
1736                }
1737                _ => false,
1738            };
1739            assert!(is_expected_page);
1740            page_count += 1;
1741        }
1742        assert_eq!(page_count, 1);
1743    }
1744
1745    fn get_serialized_page_reader<R: ChunkReader>(
1746        file_reader: &SerializedFileReader<R>,
1747        row_group: usize,
1748        column: usize,
1749    ) -> Result<SerializedPageReader<R>> {
1750        let row_group = {
1751            let row_group_metadata = file_reader.metadata.row_group(row_group);
1752            let props = Arc::clone(&file_reader.props);
1753            let f = Arc::clone(&file_reader.chunk_reader);
1754            SerializedRowGroupReader::new(
1755                f,
1756                row_group_metadata,
1757                file_reader
1758                    .metadata
1759                    .page_index()
1760                    .map(|pi| pi.offset_indexes_for_rowgroup(row_group))
1761                    .unwrap_or(None),
1762                props,
1763            )?
1764        };
1765
1766        let col = row_group.metadata.column(column);
1767
1768        let page_locations = if let Some(offset_index) = row_group.offset_index {
1769            offset_index[column]
1770                .as_ref()
1771                .map(|oi| oi.page_locations.clone())
1772        } else {
1773            None
1774        };
1775
1776        let props = Arc::clone(&row_group.props);
1777        SerializedPageReader::new_with_properties(
1778            Arc::clone(&row_group.chunk_reader),
1779            col,
1780            usize::try_from(row_group.metadata.num_rows())?,
1781            page_locations,
1782            props,
1783        )
1784    }
1785
1786    #[test]
1787    fn test_peek_next_page_offset_matches_actual() -> Result<()> {
1788        let test_file = get_test_file("alltypes_plain.parquet");
1789        let reader = SerializedFileReader::new(test_file)?;
1790
1791        let mut offset_set = HashSet::new();
1792        let num_row_groups = reader.metadata.num_row_groups();
1793        for row_group in 0..num_row_groups {
1794            let num_columns = reader.metadata.row_group(row_group).num_columns();
1795            for column in 0..num_columns {
1796                let mut page_reader = get_serialized_page_reader(&reader, row_group, column)?;
1797
1798                while let Ok(Some(page_offset)) = page_reader.peek_next_page_offset() {
1799                    match &page_reader.state {
1800                        SerializedPageReaderState::Pages {
1801                            page_locations,
1802                            dictionary_page,
1803                            ..
1804                        } => {
1805                            if let Some(page) = dictionary_page {
1806                                assert_eq!(page.offset as u64, page_offset);
1807                            } else if let Some(page) = page_locations.front() {
1808                                assert_eq!(page.offset as u64, page_offset);
1809                            } else {
1810                                unreachable!()
1811                            }
1812                        }
1813                        SerializedPageReaderState::Values {
1814                            offset,
1815                            next_page_header,
1816                            ..
1817                        } => {
1818                            assert!(next_page_header.is_some());
1819                            assert_eq!(*offset, page_offset);
1820                        }
1821                    }
1822                    let page = page_reader.get_next_page()?;
1823                    assert!(page.is_some());
1824                    let newly_inserted = offset_set.insert(page_offset);
1825                    assert!(newly_inserted);
1826                }
1827            }
1828        }
1829
1830        Ok(())
1831    }
1832
1833    #[test]
1834    fn test_page_iterator() {
1835        let file = get_test_file("alltypes_plain.parquet");
1836        let file_reader = Arc::new(SerializedFileReader::new(file).unwrap());
1837
1838        let mut page_iterator = FilePageIterator::new(0, file_reader.clone()).unwrap();
1839
1840        // read first page
1841        let page = page_iterator.next();
1842        assert!(page.is_some());
1843        assert!(page.unwrap().is_ok());
1844
1845        // reach end of file
1846        let page = page_iterator.next();
1847        assert!(page.is_none());
1848
1849        let row_group_indices = Box::new(0..1);
1850        let mut page_iterator =
1851            FilePageIterator::with_row_groups(0, row_group_indices, file_reader).unwrap();
1852
1853        // read first page
1854        let page = page_iterator.next();
1855        assert!(page.is_some());
1856        assert!(page.unwrap().is_ok());
1857
1858        // reach end of file
1859        let page = page_iterator.next();
1860        assert!(page.is_none());
1861    }
1862
1863    #[test]
1864    fn test_file_reader_key_value_metadata() {
1865        let file = get_test_file("binary.parquet");
1866        let file_reader = Arc::new(SerializedFileReader::new(file).unwrap());
1867
1868        let metadata = file_reader
1869            .metadata
1870            .file_metadata()
1871            .key_value_metadata()
1872            .unwrap();
1873
1874        assert_eq!(metadata.len(), 3);
1875
1876        assert_eq!(metadata[0].key, "parquet.proto.descriptor");
1877
1878        assert_eq!(metadata[1].key, "writer.model.name");
1879        assert_eq!(metadata[1].value, Some("protobuf".to_owned()));
1880
1881        assert_eq!(metadata[2].key, "parquet.proto.class");
1882        assert_eq!(metadata[2].value, Some("foo.baz.Foobaz$Event".to_owned()));
1883    }
1884
1885    #[test]
1886    fn test_file_reader_optional_metadata() {
1887        // file with optional metadata: bloom filters, encoding stats, column index and offset index.
1888        let file = get_test_file("data_index_bloom_encoding_stats.parquet");
1889        let options = ReadOptionsBuilder::new()
1890            .with_encoding_stats_as_mask(false)
1891            .build();
1892        let file_reader = Arc::new(SerializedFileReader::new_with_options(file, options).unwrap());
1893
1894        let row_group_metadata = file_reader.metadata.row_group(0);
1895        let col0_metadata = row_group_metadata.column(0);
1896
1897        // test optional bloom filter offset
1898        assert_eq!(col0_metadata.bloom_filter_offset().unwrap(), 192);
1899
1900        // test page encoding stats
1901        let page_encoding_stats = &col0_metadata.page_encoding_stats().unwrap()[0];
1902
1903        assert_eq!(page_encoding_stats.page_type, basic::PageType::DATA_PAGE);
1904        assert_eq!(page_encoding_stats.encoding, Encoding::PLAIN);
1905        assert_eq!(page_encoding_stats.count, 1);
1906
1907        // test optional column index offset
1908        assert_eq!(col0_metadata.column_index_offset().unwrap(), 156);
1909        assert_eq!(col0_metadata.column_index_length().unwrap(), 25);
1910
1911        // test optional offset index offset
1912        assert_eq!(col0_metadata.offset_index_offset().unwrap(), 181);
1913        assert_eq!(col0_metadata.offset_index_length().unwrap(), 11);
1914    }
1915
1916    #[test]
1917    fn test_file_reader_page_stats_mask() {
1918        let file = get_test_file("alltypes_tiny_pages.parquet");
1919        let options = ReadOptionsBuilder::new()
1920            .with_encoding_stats_as_mask(true)
1921            .build();
1922        let file_reader = Arc::new(SerializedFileReader::new_with_options(file, options).unwrap());
1923
1924        let row_group_metadata = file_reader.metadata.row_group(0);
1925
1926        // test page encoding stats
1927        let page_encoding_stats = row_group_metadata
1928            .column(0)
1929            .page_encoding_stats_mask()
1930            .unwrap();
1931        assert!(page_encoding_stats.is_only(Encoding::PLAIN));
1932        let page_encoding_stats = row_group_metadata
1933            .column(2)
1934            .page_encoding_stats_mask()
1935            .unwrap();
1936        assert!(page_encoding_stats.is_only(Encoding::PLAIN_DICTIONARY));
1937    }
1938
1939    #[test]
1940    fn test_file_reader_page_stats_skipped() {
1941        let file = get_test_file("alltypes_tiny_pages.parquet");
1942
1943        // test skipping all
1944        let options = ReadOptionsBuilder::new()
1945            .with_encoding_stats_policy(ParquetStatisticsPolicy::SkipAll)
1946            .with_column_stats_policy(ParquetStatisticsPolicy::SkipAll)
1947            .build();
1948        let file_reader = Arc::new(
1949            SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1950        );
1951
1952        let row_group_metadata = file_reader.metadata.row_group(0);
1953        for column in row_group_metadata.columns() {
1954            assert!(column.page_encoding_stats().is_none());
1955            assert!(column.page_encoding_stats_mask().is_none());
1956            assert!(column.statistics().is_none());
1957        }
1958
1959        // test skipping all but one column
1960        let options = ReadOptionsBuilder::new()
1961            .with_encoding_stats_as_mask(true)
1962            .with_encoding_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]))
1963            .with_column_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]))
1964            .build();
1965        let file_reader = Arc::new(
1966            SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1967        );
1968
1969        let row_group_metadata = file_reader.metadata.row_group(0);
1970        for (idx, column) in row_group_metadata.columns().iter().enumerate() {
1971            assert!(column.page_encoding_stats().is_none());
1972            assert_eq!(column.page_encoding_stats_mask().is_some(), idx == 0);
1973            assert_eq!(column.statistics().is_some(), idx == 0);
1974        }
1975    }
1976
1977    #[test]
1978    fn test_file_reader_size_stats_skipped() {
1979        let file = get_test_file("repeated_primitive_no_list.parquet");
1980
1981        // test skipping all
1982        let options = ReadOptionsBuilder::new()
1983            .with_size_stats_policy(ParquetStatisticsPolicy::SkipAll)
1984            .build();
1985        let file_reader = Arc::new(
1986            SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1987        );
1988
1989        let row_group_metadata = file_reader.metadata.row_group(0);
1990        for column in row_group_metadata.columns() {
1991            assert!(column.repetition_level_histogram().is_none());
1992            assert!(column.definition_level_histogram().is_none());
1993            assert!(column.unencoded_byte_array_data_bytes().is_none());
1994        }
1995
1996        // test skipping all but one column
1997        let options = ReadOptionsBuilder::new()
1998            .with_encoding_stats_as_mask(true)
1999            .with_size_stats_policy(ParquetStatisticsPolicy::skip_except(&[1]))
2000            .build();
2001        let file_reader = Arc::new(
2002            SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
2003        );
2004
2005        let row_group_metadata = file_reader.metadata.row_group(0);
2006        for (idx, column) in row_group_metadata.columns().iter().enumerate() {
2007            assert_eq!(column.repetition_level_histogram().is_some(), idx == 1);
2008            assert_eq!(column.definition_level_histogram().is_some(), idx == 1);
2009            assert_eq!(column.unencoded_byte_array_data_bytes().is_some(), idx == 1);
2010        }
2011    }
2012
2013    #[test]
2014    fn test_file_reader_with_no_filter() -> Result<()> {
2015        let test_file = get_test_file("alltypes_plain.parquet");
2016        let origin_reader = SerializedFileReader::new(test_file)?;
2017        // test initial number of row groups
2018        let metadata = origin_reader.metadata();
2019        assert_eq!(metadata.num_row_groups(), 1);
2020        Ok(())
2021    }
2022
2023    #[test]
2024    fn test_file_reader_filter_row_groups_with_predicate() -> Result<()> {
2025        let test_file = get_test_file("alltypes_plain.parquet");
2026        let read_options = ReadOptionsBuilder::new()
2027            .with_predicate(Box::new(|_, _| false))
2028            .build();
2029        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2030        let metadata = reader.metadata();
2031        assert_eq!(metadata.num_row_groups(), 0);
2032        Ok(())
2033    }
2034
2035    #[test]
2036    fn test_file_reader_filter_row_groups_with_range() -> Result<()> {
2037        let test_file = get_test_file("alltypes_plain.parquet");
2038        let origin_reader = SerializedFileReader::new(test_file)?;
2039        // test initial number of row groups
2040        let metadata = origin_reader.metadata();
2041        assert_eq!(metadata.num_row_groups(), 1);
2042        let mid = get_midpoint_offset(metadata.row_group(0));
2043
2044        let test_file = get_test_file("alltypes_plain.parquet");
2045        let read_options = ReadOptionsBuilder::new().with_range(0, mid + 1).build();
2046        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2047        let metadata = reader.metadata();
2048        assert_eq!(metadata.num_row_groups(), 1);
2049
2050        let test_file = get_test_file("alltypes_plain.parquet");
2051        let read_options = ReadOptionsBuilder::new().with_range(0, mid).build();
2052        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2053        let metadata = reader.metadata();
2054        assert_eq!(metadata.num_row_groups(), 0);
2055        Ok(())
2056    }
2057
2058    #[test]
2059    #[cfg_attr(miri, ignore)] // Takes too long
2060    fn test_file_reader_filter_row_groups_and_range() -> Result<()> {
2061        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2062        let origin_reader = SerializedFileReader::new(test_file)?;
2063        let metadata = origin_reader.metadata();
2064        let mid = get_midpoint_offset(metadata.row_group(0));
2065
2066        // true, true predicate
2067        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2068        let read_options = ReadOptionsBuilder::new()
2069            .with_page_index()
2070            .with_predicate(Box::new(|_, _| true))
2071            .with_range(mid, mid + 1)
2072            .build();
2073        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2074        let metadata = reader.metadata();
2075        assert_eq!(metadata.num_row_groups(), 1);
2076        assert!(metadata.page_index().is_some());
2077
2078        // true, false predicate
2079        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2080        let read_options = ReadOptionsBuilder::new()
2081            .with_page_index()
2082            .with_predicate(Box::new(|_, _| true))
2083            .with_range(0, mid)
2084            .build();
2085        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2086        let metadata = reader.metadata();
2087        assert_eq!(metadata.num_row_groups(), 0);
2088        assert!(metadata.page_index().is_none());
2089
2090        // false, true predicate
2091        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2092        let read_options = ReadOptionsBuilder::new()
2093            .with_page_index()
2094            .with_predicate(Box::new(|_, _| false))
2095            .with_range(mid, mid + 1)
2096            .build();
2097        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2098        let metadata = reader.metadata();
2099        assert_eq!(metadata.num_row_groups(), 0);
2100        assert!(metadata.page_index().is_none());
2101
2102        // false, false predicate
2103        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2104        let read_options = ReadOptionsBuilder::new()
2105            .with_page_index()
2106            .with_predicate(Box::new(|_, _| false))
2107            .with_range(0, mid)
2108            .build();
2109        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2110        let metadata = reader.metadata();
2111        assert_eq!(metadata.num_row_groups(), 0);
2112        assert!(metadata.page_index().is_none());
2113        Ok(())
2114    }
2115
2116    #[test]
2117    fn test_file_reader_invalid_metadata() {
2118        let data = [
2119            255, 172, 1, 0, 50, 82, 65, 73, 1, 0, 0, 0, 169, 168, 168, 162, 87, 255, 16, 0, 0, 0,
2120            80, 65, 82, 49,
2121        ];
2122        let ret = SerializedFileReader::new(Bytes::copy_from_slice(&data));
2123        assert_eq!(
2124            ret.err().unwrap().to_string(),
2125            "Parquet error: Expected list element type of Struct but got List"
2126        );
2127    }
2128
2129    #[test]
2130    // Use java parquet-tools get below pageIndex info
2131    // !```
2132    // parquet-tools column-index ./data_index_bloom_encoding_stats.parquet
2133    // row group 0:
2134    // column index for column String:
2135    // Boundary order: ASCENDING
2136    // page-0  :
2137    // null count                 min                                  max
2138    // 0                          Hello                                today
2139    //
2140    // offset index for column String:
2141    // page-0   :
2142    // offset   compressed size       first row index
2143    // 4               152                     0
2144    ///```
2145    //
2146    fn test_page_index_reader() {
2147        let test_file = get_test_file("data_index_bloom_encoding_stats.parquet");
2148        let builder = ReadOptionsBuilder::new();
2149        //enable read page index
2150        let options = builder.with_page_index().build();
2151        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2152        let reader = reader_result.unwrap();
2153
2154        // Test contents in Parquet metadata
2155        let metadata = reader.metadata();
2156        assert_eq!(metadata.num_row_groups(), 1);
2157
2158        let page_index = metadata.page_index().expect("page index should be present");
2159
2160        // only one row group
2161        let Some(ColumnIndexMetaData::BYTE_ARRAY(index)) = page_index.column_index(0, 0) else {
2162            unreachable!()
2163        };
2164
2165        assert_eq!(index.boundary_order, BoundaryOrder::ASCENDING);
2166
2167        //only one page group
2168        assert_eq!(index.num_pages(), 1);
2169
2170        let min = index.min_value(0).unwrap();
2171        let max = index.max_value(0).unwrap();
2172        assert_eq!(b"Hello", min.as_bytes());
2173        assert_eq!(b"today", max.as_bytes());
2174
2175        // only one row group
2176        let offset_index = page_index
2177            .offset_index(0, 0)
2178            .expect("offset index should be present");
2179        let page_offset = offset_index
2180            .page_locations()
2181            .first()
2182            .expect("offset index too small");
2183
2184        assert_eq!(4, page_offset.offset);
2185        assert_eq!(152, page_offset.compressed_page_size);
2186        assert_eq!(0, page_offset.first_row_index);
2187    }
2188
2189    #[test]
2190    #[cfg_attr(miri, ignore)] // Takes too long
2191    fn test_page_index_reader_all_type() {
2192        let test_file = get_test_file("alltypes_tiny_pages_plain.parquet");
2193        let builder = ReadOptionsBuilder::new();
2194        //enable read page index
2195        let options = builder.with_page_index().build();
2196        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2197        let reader = reader_result.unwrap();
2198
2199        // Test contents in Parquet metadata
2200        let metadata = reader.metadata();
2201        assert_eq!(metadata.num_row_groups(), 1);
2202
2203        let page_index = metadata.page_index().unwrap();
2204        let row_group_offset_indexes = page_index.offset_indexes_for_rowgroup(0).unwrap();
2205
2206        // only one row group
2207        let row_group_metadata = metadata.row_group(0);
2208
2209        //col0->id: INT32 UNCOMPRESSED DO:0 FPO:4 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 7299, num_nulls: 0]
2210        let ci = page_index.column_index(0, 0).unwrap();
2211        assert!(!ci.is_sorted());
2212        assert!(matches!(
2213            ci.get_boundary_order(),
2214            Some(BoundaryOrder::UNORDERED)
2215        ));
2216        if let ColumnIndexMetaData::INT32(index) = ci {
2217            check_native_page_index(
2218                index,
2219                325,
2220                get_row_group_min_max_bytes(row_group_metadata, 0),
2221                BoundaryOrder::UNORDERED,
2222            );
2223            assert_eq!(
2224                row_group_offset_indexes[0]
2225                    .as_ref()
2226                    .unwrap()
2227                    .page_locations
2228                    .len(),
2229                325
2230            );
2231        } else {
2232            unreachable!()
2233        }
2234        //col1->bool_col:BOOLEAN UNCOMPRESSED DO:0 FPO:37329 SZ:3022/3022/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: false, max: true, num_nulls: 0]
2235        let ci = page_index.column_index(0, 1).unwrap();
2236        assert!(ci.is_sorted());
2237        if let ColumnIndexMetaData::BOOLEAN(index) = ci {
2238            assert_eq!(index.num_pages(), 82);
2239            assert_eq!(
2240                row_group_offset_indexes[1]
2241                    .as_ref()
2242                    .unwrap()
2243                    .page_locations
2244                    .len(),
2245                82
2246            );
2247        } else {
2248            unreachable!()
2249        }
2250        //col2->tinyint_col: INT32 UNCOMPRESSED DO:0 FPO:40351 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0]
2251        let ci = page_index.column_index(0, 2).unwrap();
2252        assert!(ci.is_sorted());
2253        if let ColumnIndexMetaData::INT32(index) = ci {
2254            check_native_page_index(
2255                index,
2256                325,
2257                get_row_group_min_max_bytes(row_group_metadata, 2),
2258                BoundaryOrder::ASCENDING,
2259            );
2260            assert_eq!(
2261                row_group_offset_indexes[2]
2262                    .as_ref()
2263                    .unwrap()
2264                    .page_locations
2265                    .len(),
2266                325
2267            );
2268        } else {
2269            unreachable!()
2270        }
2271        //col4->smallint_col: INT32 UNCOMPRESSED DO:0 FPO:77676 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0]
2272        let ci = page_index.column_index(0, 3).unwrap();
2273        assert!(ci.is_sorted());
2274        if let ColumnIndexMetaData::INT32(index) = ci {
2275            check_native_page_index(
2276                index,
2277                325,
2278                get_row_group_min_max_bytes(row_group_metadata, 3),
2279                BoundaryOrder::ASCENDING,
2280            );
2281            assert_eq!(
2282                row_group_offset_indexes[3]
2283                    .as_ref()
2284                    .unwrap()
2285                    .page_locations
2286                    .len(),
2287                325
2288            );
2289        } else {
2290            unreachable!()
2291        }
2292        //col5->smallint_col: INT32 UNCOMPRESSED DO:0 FPO:77676 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0]
2293        let ci = page_index.column_index(0, 4).unwrap();
2294        assert!(ci.is_sorted());
2295        if let ColumnIndexMetaData::INT32(index) = ci {
2296            check_native_page_index(
2297                index,
2298                325,
2299                get_row_group_min_max_bytes(row_group_metadata, 4),
2300                BoundaryOrder::ASCENDING,
2301            );
2302            assert_eq!(
2303                row_group_offset_indexes[4]
2304                    .as_ref()
2305                    .unwrap()
2306                    .page_locations
2307                    .len(),
2308                325
2309            );
2310        } else {
2311            unreachable!()
2312        }
2313        //col6->bigint_col: INT64 UNCOMPRESSED DO:0 FPO:152326 SZ:71598/71598/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 90, num_nulls: 0]
2314        let ci = page_index.column_index(0, 5).unwrap();
2315        assert!(!ci.is_sorted());
2316        if let ColumnIndexMetaData::INT64(index) = ci {
2317            check_native_page_index(
2318                index,
2319                528,
2320                get_row_group_min_max_bytes(row_group_metadata, 5),
2321                BoundaryOrder::UNORDERED,
2322            );
2323            assert_eq!(
2324                row_group_offset_indexes[5]
2325                    .as_ref()
2326                    .unwrap()
2327                    .page_locations
2328                    .len(),
2329                528
2330            );
2331        } else {
2332            unreachable!()
2333        }
2334        //col7->float_col: FLOAT UNCOMPRESSED DO:0 FPO:223924 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: -0.0, max: 9.9, num_nulls: 0]
2335        let ci = page_index.column_index(0, 6).unwrap();
2336        assert!(ci.is_sorted());
2337        if let ColumnIndexMetaData::FLOAT(index) = ci {
2338            check_native_page_index(
2339                index,
2340                325,
2341                get_row_group_min_max_bytes(row_group_metadata, 6),
2342                BoundaryOrder::ASCENDING,
2343            );
2344            assert_eq!(
2345                row_group_offset_indexes[6]
2346                    .as_ref()
2347                    .unwrap()
2348                    .page_locations
2349                    .len(),
2350                325
2351            );
2352        } else {
2353            unreachable!()
2354        }
2355        //col8->double_col: DOUBLE UNCOMPRESSED DO:0 FPO:261249 SZ:71598/71598/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: -0.0, max: 90.89999999999999, num_nulls: 0]
2356        let ci = page_index.column_index(0, 7).unwrap();
2357        assert!(!ci.is_sorted());
2358        if let ColumnIndexMetaData::DOUBLE(index) = ci {
2359            check_native_page_index(
2360                index,
2361                528,
2362                get_row_group_min_max_bytes(row_group_metadata, 7),
2363                BoundaryOrder::UNORDERED,
2364            );
2365            assert_eq!(
2366                row_group_offset_indexes[7]
2367                    .as_ref()
2368                    .unwrap()
2369                    .page_locations
2370                    .len(),
2371                528
2372            );
2373        } else {
2374            unreachable!()
2375        }
2376        //col9->date_string_col: BINARY UNCOMPRESSED DO:0 FPO:332847 SZ:111948/111948/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 01/01/09, max: 12/31/10, num_nulls: 0]
2377        let ci = page_index.column_index(0, 8).unwrap();
2378        assert!(!ci.is_sorted());
2379        if let ColumnIndexMetaData::BYTE_ARRAY(index) = ci {
2380            check_byte_array_page_index(
2381                index,
2382                974,
2383                get_row_group_min_max_bytes(row_group_metadata, 8),
2384                BoundaryOrder::UNORDERED,
2385            );
2386            assert_eq!(
2387                row_group_offset_indexes[8]
2388                    .as_ref()
2389                    .unwrap()
2390                    .page_locations
2391                    .len(),
2392                974
2393            );
2394        } else {
2395            unreachable!()
2396        }
2397        //col10->string_col: BINARY UNCOMPRESSED DO:0 FPO:444795 SZ:45298/45298/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0]
2398        let ci = page_index.column_index(0, 9).unwrap();
2399        assert!(ci.is_sorted());
2400        if let ColumnIndexMetaData::BYTE_ARRAY(index) = ci {
2401            check_byte_array_page_index(
2402                index,
2403                352,
2404                get_row_group_min_max_bytes(row_group_metadata, 9),
2405                BoundaryOrder::ASCENDING,
2406            );
2407            assert_eq!(
2408                row_group_offset_indexes[9]
2409                    .as_ref()
2410                    .unwrap()
2411                    .page_locations
2412                    .len(),
2413                352
2414            );
2415        } else {
2416            unreachable!()
2417        }
2418        //col11->timestamp_col: INT96 UNCOMPRESSED DO:0 FPO:490093 SZ:111948/111948/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[num_nulls: 0, min/max not defined]
2419        // this columns lacks an index
2420        assert!(page_index.column_index(0, 10).is_none());
2421        //col12->year: INT32 UNCOMPRESSED DO:0 FPO:602041 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 2009, max: 2010, num_nulls: 0]
2422        let ci = page_index.column_index(0, 11).unwrap();
2423        assert!(ci.is_sorted());
2424        if let ColumnIndexMetaData::INT32(index) = ci {
2425            check_native_page_index(
2426                index,
2427                325,
2428                get_row_group_min_max_bytes(row_group_metadata, 11),
2429                BoundaryOrder::ASCENDING,
2430            );
2431            assert_eq!(
2432                row_group_offset_indexes[11]
2433                    .as_ref()
2434                    .unwrap()
2435                    .page_locations
2436                    .len(),
2437                325
2438            );
2439        } else {
2440            unreachable!()
2441        }
2442        //col13->month: INT32 UNCOMPRESSED DO:0 FPO:639366 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 1, max: 12, num_nulls: 0]
2443        let ci = page_index.column_index(0, 12).unwrap();
2444        assert!(!ci.is_sorted());
2445        if let ColumnIndexMetaData::INT32(index) = ci {
2446            check_native_page_index(
2447                index,
2448                325,
2449                get_row_group_min_max_bytes(row_group_metadata, 12),
2450                BoundaryOrder::UNORDERED,
2451            );
2452            assert_eq!(
2453                row_group_offset_indexes[12]
2454                    .as_ref()
2455                    .unwrap()
2456                    .page_locations
2457                    .len(),
2458                325
2459            );
2460        } else {
2461            unreachable!()
2462        }
2463    }
2464
2465    fn check_native_page_index<T: ParquetValueType>(
2466        row_group_index: &PrimitiveColumnIndex<T>,
2467        page_size: usize,
2468        min_max: (&[u8], &[u8]),
2469        boundary_order: BoundaryOrder,
2470    ) {
2471        assert_eq!(row_group_index.num_pages() as usize, page_size);
2472        assert_eq!(row_group_index.boundary_order, boundary_order);
2473        assert!(row_group_index.min_values().iter().all(|x| {
2474            x >= &T::try_from_le_slice(min_max.0).unwrap()
2475                && x <= &T::try_from_le_slice(min_max.1).unwrap()
2476        }));
2477    }
2478
2479    fn check_byte_array_page_index(
2480        row_group_index: &ByteArrayColumnIndex,
2481        page_size: usize,
2482        min_max: (&[u8], &[u8]),
2483        boundary_order: BoundaryOrder,
2484    ) {
2485        assert_eq!(row_group_index.num_pages() as usize, page_size);
2486        assert_eq!(row_group_index.boundary_order, boundary_order);
2487        for i in 0..row_group_index.num_pages() as usize {
2488            let x = row_group_index.min_value(i).unwrap();
2489            assert!(x >= min_max.0 && x <= min_max.1);
2490        }
2491    }
2492
2493    fn get_row_group_min_max_bytes(r: &RowGroupMetaData, col_num: usize) -> (&[u8], &[u8]) {
2494        let statistics = r.column(col_num).statistics().unwrap();
2495        (
2496            statistics.min_bytes_opt().unwrap_or_default(),
2497            statistics.max_bytes_opt().unwrap_or_default(),
2498        )
2499    }
2500
2501    #[test]
2502    #[cfg_attr(miri, ignore)] // Takes too long
2503    fn test_skip_next_page_with_dictionary_page() {
2504        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2505        let builder = ReadOptionsBuilder::new();
2506        // enable read page index
2507        let options = builder.with_page_index().build();
2508        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2509        let reader = reader_result.unwrap();
2510
2511        let row_group_reader = reader.get_row_group(0).unwrap();
2512
2513        // use 'string_col', Boundary order: UNORDERED, total 352 data pages and 1 dictionary page.
2514        let mut column_page_reader = row_group_reader.get_column_page_reader(9).unwrap();
2515
2516        let mut vec = vec![];
2517
2518        // Step 1: Peek and ensure dictionary page is correctly identified
2519        let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2520        assert!(meta.is_dict);
2521
2522        // Step 2: Call skip_next_page to skip the dictionary page
2523        column_page_reader.skip_next_page().unwrap();
2524
2525        // Step 3: Read the next data page after skipping the dictionary page
2526        let page = column_page_reader.get_next_page().unwrap().unwrap();
2527        assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2528
2529        // Step 4: Continue reading remaining data pages and verify correctness
2530        for _i in 0..351 {
2531            // 352 total pages, 1 dictionary page is skipped
2532            let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2533            assert!(!meta.is_dict); // Verify no dictionary page here
2534            vec.push(meta);
2535
2536            let page = column_page_reader.get_next_page().unwrap().unwrap();
2537            assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2538        }
2539
2540        // Step 5: Check if all pages are read
2541        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2542        assert!(column_page_reader.get_next_page().unwrap().is_none());
2543
2544        // Step 6: Verify the number of data pages read (should be 351 data pages)
2545        assert_eq!(vec.len(), 351);
2546    }
2547
2548    #[test]
2549    #[cfg_attr(miri, ignore)] // Takes too long
2550    fn test_skip_page_with_offset_index() {
2551        let test_file = get_test_file("alltypes_tiny_pages_plain.parquet");
2552        let builder = ReadOptionsBuilder::new();
2553        //enable read page index
2554        let options = builder.with_page_index().build();
2555        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2556        let reader = reader_result.unwrap();
2557
2558        let row_group_reader = reader.get_row_group(0).unwrap();
2559
2560        //use 'int_col', Boundary order: ASCENDING, total 325 pages.
2561        let mut column_page_reader = row_group_reader.get_column_page_reader(4).unwrap();
2562
2563        let mut vec = vec![];
2564
2565        for i in 0..325 {
2566            if i % 2 == 0 {
2567                vec.push(column_page_reader.get_next_page().unwrap().unwrap());
2568            } else {
2569                column_page_reader.skip_next_page().unwrap();
2570            }
2571        }
2572        //check read all pages.
2573        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2574        assert!(column_page_reader.get_next_page().unwrap().is_none());
2575
2576        assert_eq!(vec.len(), 163);
2577    }
2578
2579    #[test]
2580    fn test_skip_page_without_offset_index() {
2581        let test_file = get_test_file("alltypes_tiny_pages_plain.parquet");
2582
2583        // use default SerializedFileReader without read offsetIndex
2584        let reader_result = SerializedFileReader::new(test_file);
2585        let reader = reader_result.unwrap();
2586
2587        let row_group_reader = reader.get_row_group(0).unwrap();
2588
2589        //use 'int_col', Boundary order: ASCENDING, total 325 pages.
2590        let mut column_page_reader = row_group_reader.get_column_page_reader(4).unwrap();
2591
2592        let mut vec = vec![];
2593
2594        for i in 0..325 {
2595            if i % 2 == 0 {
2596                vec.push(column_page_reader.get_next_page().unwrap().unwrap());
2597            } else {
2598                column_page_reader.peek_next_page().unwrap().unwrap();
2599                column_page_reader.skip_next_page().unwrap();
2600            }
2601        }
2602        //check read all pages.
2603        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2604        assert!(column_page_reader.get_next_page().unwrap().is_none());
2605
2606        assert_eq!(vec.len(), 163);
2607    }
2608
2609    #[test]
2610    #[cfg_attr(miri, ignore)] // Takes too long
2611    fn test_peek_page_with_dictionary_page() {
2612        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2613        let builder = ReadOptionsBuilder::new();
2614        //enable read page index
2615        let options = builder.with_page_index().build();
2616        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2617        let reader = reader_result.unwrap();
2618        let row_group_reader = reader.get_row_group(0).unwrap();
2619
2620        //use 'string_col', Boundary order: UNORDERED, total 352 data ages and 1 dictionary page.
2621        let mut column_page_reader = row_group_reader.get_column_page_reader(9).unwrap();
2622
2623        let mut vec = vec![];
2624
2625        let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2626        assert!(meta.is_dict);
2627        let page = column_page_reader.get_next_page().unwrap().unwrap();
2628        assert!(matches!(page.page_type(), basic::PageType::DICTIONARY_PAGE));
2629
2630        for i in 0..352 {
2631            let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2632            // have checked with `parquet-tools column-index   -c string_col  ./alltypes_tiny_pages.parquet`
2633            // page meta has two scenarios(21, 20) of num_rows expect last page has 11 rows.
2634            if i != 351 {
2635                assert!((meta.num_rows == Some(21)) || (meta.num_rows == Some(20)));
2636            } else {
2637                // last page first row index is 7290, total row count is 7300
2638                // because first row start with zero, last page row count should be 10.
2639                assert_eq!(meta.num_rows, Some(10));
2640            }
2641            assert!(!meta.is_dict);
2642            vec.push(meta);
2643            let page = column_page_reader.get_next_page().unwrap().unwrap();
2644            assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2645        }
2646
2647        //check read all pages.
2648        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2649        assert!(column_page_reader.get_next_page().unwrap().is_none());
2650
2651        assert_eq!(vec.len(), 352);
2652    }
2653
2654    #[test]
2655    fn test_peek_page_with_dictionary_page_without_offset_index() {
2656        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2657
2658        let reader_result = SerializedFileReader::new(test_file);
2659        let reader = reader_result.unwrap();
2660        let row_group_reader = reader.get_row_group(0).unwrap();
2661
2662        //use 'string_col', Boundary order: UNORDERED, total 352 data ages and 1 dictionary page.
2663        let mut column_page_reader = row_group_reader.get_column_page_reader(9).unwrap();
2664
2665        let mut vec = vec![];
2666
2667        let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2668        assert!(meta.is_dict);
2669        let page = column_page_reader.get_next_page().unwrap().unwrap();
2670        assert!(matches!(page.page_type(), basic::PageType::DICTIONARY_PAGE));
2671
2672        for i in 0..352 {
2673            let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2674            // have checked with `parquet-tools column-index   -c string_col  ./alltypes_tiny_pages.parquet`
2675            // page meta has two scenarios(21, 20) of num_rows expect last page has 11 rows.
2676            if i != 351 {
2677                assert!((meta.num_levels == Some(21)) || (meta.num_levels == Some(20)));
2678            } else {
2679                // last page first row index is 7290, total row count is 7300
2680                // because first row start with zero, last page row count should be 10.
2681                assert_eq!(meta.num_levels, Some(10));
2682            }
2683            assert!(!meta.is_dict);
2684            vec.push(meta);
2685            let page = column_page_reader.get_next_page().unwrap().unwrap();
2686            assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2687        }
2688
2689        //check read all pages.
2690        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2691        assert!(column_page_reader.get_next_page().unwrap().is_none());
2692
2693        assert_eq!(vec.len(), 352);
2694    }
2695
2696    #[test]
2697    fn test_fixed_length_index() {
2698        let message_type = "
2699        message test_schema {
2700          OPTIONAL FIXED_LEN_BYTE_ARRAY (11) value (DECIMAL(25,2));
2701        }
2702        ";
2703
2704        let schema = parse_message_type(message_type).unwrap();
2705        let mut out = Vec::with_capacity(1024);
2706        let mut writer =
2707            SerializedFileWriter::new(&mut out, Arc::new(schema), Default::default()).unwrap();
2708
2709        let mut r = writer.next_row_group().unwrap();
2710        let mut c = r.next_column().unwrap().unwrap();
2711        c.typed::<FixedLenByteArrayType>()
2712            .write_batch(
2713                &[vec![0; 11].into(), vec![5; 11].into(), vec![3; 11].into()],
2714                Some(&[1, 1, 0, 1]),
2715                None,
2716            )
2717            .unwrap();
2718        c.close().unwrap();
2719        r.close().unwrap();
2720        writer.close().unwrap();
2721
2722        let b = Bytes::from(out);
2723        let options = ReadOptionsBuilder::new().with_page_index().build();
2724        let reader = SerializedFileReader::new_with_options(b, options).unwrap();
2725        let page_index = reader.metadata().page_index().unwrap();
2726
2727        match page_index.column_index(0, 0) {
2728            Some(ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(v)) => {
2729                assert_eq!(v.num_pages(), 1);
2730                assert_eq!(v.null_count(0).unwrap(), 1);
2731                assert_eq!(v.min_value(0).unwrap(), &[0; 11]);
2732                assert_eq!(v.max_value(0).unwrap(), &[5; 11]);
2733            }
2734            _ => unreachable!(),
2735        }
2736    }
2737
2738    #[test]
2739    fn test_multi_gz() {
2740        let file = get_test_file("concatenated_gzip_members.parquet");
2741        let reader = SerializedFileReader::new(file).unwrap();
2742        let row_group_reader = reader.get_row_group(0).unwrap();
2743        match row_group_reader.get_column_reader(0).unwrap() {
2744            ColumnReader::Int64ColumnReader(mut reader) => {
2745                let mut buffer = Vec::with_capacity(1024);
2746                let mut def_levels = Vec::with_capacity(1024);
2747                let (num_records, num_values, num_levels) = reader
2748                    .read_records(1024, Some(&mut def_levels), None, &mut buffer)
2749                    .unwrap();
2750
2751                assert_eq!(num_records, 513);
2752                assert_eq!(num_values, 513);
2753                assert_eq!(num_levels, 513);
2754
2755                let expected: Vec<i64> = (1..514).collect();
2756                assert_eq!(&buffer, &expected);
2757            }
2758            _ => unreachable!(),
2759        }
2760    }
2761
2762    #[test]
2763    #[cfg_attr(miri, ignore)] // Takes too long
2764    fn test_byte_stream_split_extended() {
2765        let path = format!(
2766            "{}/byte_stream_split_extended.gzip.parquet",
2767            arrow::util::test_util::parquet_test_data(),
2768        );
2769        let file = File::open(path).unwrap();
2770        let reader = Box::new(SerializedFileReader::new(file).expect("Failed to create reader"));
2771
2772        // Use full schema as projected schema
2773        let mut iter = reader
2774            .get_row_iter(None)
2775            .expect("Failed to create row iterator");
2776
2777        let mut start = 0;
2778        let end = reader.metadata().file_metadata().num_rows();
2779
2780        let check_row = |row: Result<Row, ParquetError>| {
2781            assert!(row.is_ok());
2782            let r = row.unwrap();
2783            assert_eq!(r.get_float16(0).unwrap(), r.get_float16(1).unwrap());
2784            assert_eq!(r.get_float(2).unwrap(), r.get_float(3).unwrap());
2785            assert_eq!(r.get_double(4).unwrap(), r.get_double(5).unwrap());
2786            assert_eq!(r.get_int(6).unwrap(), r.get_int(7).unwrap());
2787            assert_eq!(r.get_long(8).unwrap(), r.get_long(9).unwrap());
2788            assert_eq!(r.get_bytes(10).unwrap(), r.get_bytes(11).unwrap());
2789            assert_eq!(r.get_decimal(12).unwrap(), r.get_decimal(13).unwrap());
2790        };
2791
2792        while start < end {
2793            match iter.next() {
2794                Some(row) => check_row(row),
2795                None => break,
2796            }
2797            start += 1;
2798        }
2799    }
2800
2801    #[test]
2802    fn test_filtered_rowgroup_metadata() {
2803        let message_type = "
2804            message test_schema {
2805                REQUIRED INT32 a;
2806            }
2807        ";
2808        let schema = Arc::new(parse_message_type(message_type).unwrap());
2809        let props = Arc::new(
2810            WriterProperties::builder()
2811                .set_statistics_enabled(EnabledStatistics::Page)
2812                .build(),
2813        );
2814        let mut file: File = tempfile::tempfile().unwrap();
2815        let mut file_writer = SerializedFileWriter::new(&mut file, schema, props).unwrap();
2816        let data = [1, 2, 3, 4, 5];
2817
2818        // write 5 row groups
2819        for idx in 0..5 {
2820            let data_i: Vec<i32> = data.iter().map(|x| x * (idx + 1)).collect();
2821            let mut row_group_writer = file_writer.next_row_group().unwrap();
2822            if let Some(mut writer) = row_group_writer.next_column().unwrap() {
2823                writer
2824                    .typed::<Int32Type>()
2825                    .write_batch(data_i.as_slice(), None, None)
2826                    .unwrap();
2827                writer.close().unwrap();
2828            }
2829            row_group_writer.close().unwrap();
2830            file_writer.flushed_row_groups();
2831        }
2832        let file_metadata = file_writer.close().unwrap();
2833
2834        assert_eq!(file_metadata.file_metadata().num_rows(), 25);
2835        assert_eq!(file_metadata.num_row_groups(), 5);
2836
2837        // read only the 3rd row group
2838        let read_options = ReadOptionsBuilder::new()
2839            .with_page_index()
2840            .with_predicate(Box::new(|rgmeta, _| rgmeta.ordinal().unwrap_or(0) == 2))
2841            .build();
2842        let reader =
2843            SerializedFileReader::new_with_options(file.try_clone().unwrap(), read_options)
2844                .unwrap();
2845        let metadata = reader.metadata();
2846
2847        // check we got the expected row group
2848        assert_eq!(metadata.num_row_groups(), 1);
2849        assert_eq!(metadata.row_group(0).ordinal(), Some(2));
2850
2851        // check we only got the relevant page indexes
2852        assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
2853        let page_index = metadata.page_index().unwrap();
2854
2855        let col_stats = metadata.row_group(0).column(0).statistics().unwrap();
2856        let pg_idx = page_index.column_index(0, 0);
2857        let off_idx_i = page_index.offset_index(0, 0);
2858
2859        // test that we got the index matching the row group
2860        match pg_idx {
2861            Some(ColumnIndexMetaData::INT32(int_idx)) => {
2862                let min = col_stats.min_bytes_opt().unwrap().get_i32_le();
2863                let max = col_stats.max_bytes_opt().unwrap().get_i32_le();
2864                assert_eq!(int_idx.min_value(0), Some(min).as_ref());
2865                assert_eq!(int_idx.max_value(0), Some(max).as_ref());
2866            }
2867            _ => panic!("wrong stats type"),
2868        }
2869
2870        // check offset index matches too
2871        assert_eq!(
2872            off_idx_i.as_ref().unwrap().page_locations[0].offset,
2873            metadata.row_group(0).column(0).data_page_offset()
2874        );
2875
2876        // read non-contiguous row groups
2877        let read_options = ReadOptionsBuilder::new()
2878            .with_page_index()
2879            .with_predicate(Box::new(|rgmeta, _| rgmeta.ordinal().unwrap_or(0) % 2 == 1))
2880            .build();
2881        let reader =
2882            SerializedFileReader::new_with_options(file.try_clone().unwrap(), read_options)
2883                .unwrap();
2884        let metadata = reader.metadata();
2885
2886        // check we got the expected row groups
2887        assert_eq!(metadata.num_row_groups(), 2);
2888        assert_eq!(metadata.row_group(0).ordinal(), Some(1));
2889        assert_eq!(metadata.row_group(1).ordinal(), Some(3));
2890
2891        // check we only got the relevant page indexes
2892        assert!(metadata.page_index().is_some_and(PageIndex::is_complete));
2893
2894        let page_index = metadata.page_index().unwrap();
2895
2896        for rg_idx in 0..metadata.num_row_groups() {
2897            let col_stats = metadata.row_group(rg_idx).column(0).statistics().unwrap();
2898            let pg_idx = page_index.column_index(rg_idx, 0);
2899            let off_idx_i = page_index.offset_index(rg_idx, 0);
2900
2901            // test that we got the index matching the row group
2902            match pg_idx {
2903                Some(ColumnIndexMetaData::INT32(int_idx)) => {
2904                    let min = col_stats.min_bytes_opt().unwrap().get_i32_le();
2905                    let max = col_stats.max_bytes_opt().unwrap().get_i32_le();
2906                    assert_eq!(int_idx.min_value(0), Some(min).as_ref());
2907                    assert_eq!(int_idx.max_value(0), Some(max).as_ref());
2908                }
2909                _ => panic!("wrong stats type"),
2910            }
2911
2912            // check offset index matches too
2913            assert_eq!(
2914                off_idx_i.as_ref().unwrap().page_locations[0].offset,
2915                metadata.row_group(rg_idx).column(0).data_page_offset()
2916            );
2917        }
2918    }
2919
2920    #[test]
2921    fn test_reuse_schema() {
2922        let file = get_test_file("alltypes_plain.parquet");
2923        let file_reader = SerializedFileReader::new(file.try_clone().unwrap()).unwrap();
2924        let schema = file_reader.metadata().file_metadata().schema_descr_ptr();
2925        let expected = file_reader.metadata;
2926
2927        let options = ReadOptionsBuilder::new()
2928            .with_parquet_schema(schema)
2929            .build();
2930        let file_reader = SerializedFileReader::new_with_options(file, options).unwrap();
2931
2932        assert_eq!(expected.as_ref(), file_reader.metadata.as_ref());
2933        // Should have used the same schema instance
2934        assert!(Arc::ptr_eq(
2935            &expected.file_metadata().schema_descr_ptr(),
2936            &file_reader.metadata.file_metadata().schema_descr_ptr()
2937        ));
2938    }
2939
2940    #[test]
2941    fn test_read_unknown_logical_type() {
2942        let file = get_test_file("unknown-logical-type.parquet");
2943        let reader = SerializedFileReader::new(file).expect("Error opening file");
2944
2945        let schema = reader.metadata().file_metadata().schema_descr();
2946        assert_eq!(
2947            schema.column(0).logical_type_ref(),
2948            Some(&basic::LogicalType::String)
2949        );
2950        assert_eq!(
2951            schema.column(1).logical_type_ref(),
2952            Some(&basic::LogicalType::_Unknown { field_id: 2555 })
2953        );
2954        assert_eq!(schema.column(1).physical_type(), Type::BYTE_ARRAY);
2955
2956        let mut iter = reader
2957            .get_row_iter(None)
2958            .expect("Failed to create row iterator");
2959
2960        let mut num_rows = 0;
2961        while iter.next().is_some() {
2962            num_rows += 1;
2963        }
2964        assert_eq!(num_rows, reader.metadata().file_metadata().num_rows());
2965    }
2966}