Skip to main content

parquet/file/metadata/
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#[cfg(feature = "encryption")]
19use crate::encryption::decrypt::FileDecryptionProperties;
20use crate::errors::{ParquetError, Result};
21use crate::file::FOOTER_SIZE;
22use crate::file::metadata::parser::decode_metadata;
23use crate::file::metadata::thrift::parquet_schema_from_bytes;
24use crate::file::metadata::{
25    FooterTail, ParquetMetaData, ParquetMetaDataOptions, ParquetMetaDataPushDecoder,
26};
27use crate::file::reader::ChunkReader;
28use crate::schema::types::SchemaDescriptor;
29use bytes::Bytes;
30use std::sync::Arc;
31use std::{io::Read, ops::Range};
32
33use crate::DecodeResult;
34#[cfg(all(feature = "async", feature = "arrow"))]
35use crate::arrow::async_reader::{MetadataFetch, MetadataSuffixFetch};
36
37/// Reads [`ParquetMetaData`] from a byte stream, with either synchronous or
38/// asynchronous I/O.
39///
40/// There are two flavors of APIs:
41/// * Synchronous: [`Self::try_parse()`], [`Self::try_parse_sized()`], [`Self::parse_and_finish()`], etc.
42/// * Asynchronous (requires `async` and `arrow` features): [`Self::try_load()`], etc
43///
44///  See the [`ParquetMetaDataPushDecoder`] for an API that does not require I/O.
45///
46/// # Format Notes
47///
48/// Parquet metadata is not necessarily contiguous in a Parquet file: a portion is stored
49/// in the footer (the last bytes of the file), but other portions (such as the
50/// PageIndex) can be stored elsewhere.
51/// See [`crate::file::metadata::ParquetMetaDataWriter#output-format`] for more details of
52/// Parquet metadata.
53///
54/// This reader handles reading the footer as well as the non contiguous parts
55/// of the metadata (`PageIndex` and `ColumnIndex`). It does not handle reading Bloom Filters.
56///
57/// # Example
58/// ```no_run
59/// # use parquet::file::metadata::{PageIndexPolicy, ParquetMetaDataReader};
60/// # fn open_parquet_file(path: &str) -> std::fs::File { unimplemented!(); }
61/// // read parquet metadata including page indexes from a file
62/// let file = open_parquet_file("some_path.parquet");
63/// let mut reader = ParquetMetaDataReader::new()
64///     .with_page_index_policy(PageIndexPolicy::Required);
65/// reader.try_parse(&file).unwrap();
66/// let metadata = reader.finish().unwrap();
67/// assert!(metadata.page_index().is_some());
68/// ```
69#[derive(Default, Debug)]
70pub struct ParquetMetaDataReader {
71    metadata: Option<ParquetMetaData>,
72    column_index: PageIndexPolicy,
73    offset_index: PageIndexPolicy,
74    prefetch_hint: Option<usize>,
75    metadata_options: Option<Arc<ParquetMetaDataOptions>>,
76    // Size of the serialized thrift metadata plus the 8 byte footer. Only set if
77    // `self.parse_metadata` is called.
78    metadata_size: Option<usize>,
79    #[cfg(feature = "encryption")]
80    file_decryption_properties: Option<Arc<FileDecryptionProperties>>,
81}
82
83/// Describes the policy for reading page indexes
84#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
85pub enum PageIndexPolicy {
86    /// Do not read the page index.
87    #[default]
88    Skip,
89    /// Read the page index if it exists, otherwise do not error.
90    Optional,
91    /// Require the page index to exist, and error if it does not.
92    ///
93    /// Only enforced for the offset index; this is the same as [`Self::Optional`]
94    /// for the column index.
95    Required,
96}
97
98impl From<bool> for PageIndexPolicy {
99    fn from(value: bool) -> Self {
100        match value {
101            true => Self::Required,
102            false => Self::Skip,
103        }
104    }
105}
106
107impl ParquetMetaDataReader {
108    /// Create a new [`ParquetMetaDataReader`]
109    pub fn new() -> Self {
110        Default::default()
111    }
112
113    /// Create a new [`ParquetMetaDataReader`] populated with a [`ParquetMetaData`] struct
114    /// obtained via other means.
115    pub fn new_with_metadata(metadata: ParquetMetaData) -> Self {
116        Self {
117            metadata: Some(metadata),
118            ..Default::default()
119        }
120    }
121
122    /// Sets the [`PageIndexPolicy`] for the column and offset indexes
123    pub fn with_page_index_policy(self, policy: PageIndexPolicy) -> Self {
124        self.with_column_index_policy(policy)
125            .with_offset_index_policy(policy)
126    }
127
128    /// Sets the [`PageIndexPolicy`] for the column index
129    pub fn with_column_index_policy(mut self, policy: PageIndexPolicy) -> Self {
130        self.column_index = policy;
131        self
132    }
133
134    /// Sets the [`PageIndexPolicy`] for the offset index
135    pub fn with_offset_index_policy(mut self, policy: PageIndexPolicy) -> Self {
136        self.offset_index = policy;
137        self
138    }
139
140    /// Sets the [`ParquetMetaDataOptions`] to use when decoding
141    pub fn with_metadata_options(mut self, options: Option<ParquetMetaDataOptions>) -> Self {
142        self.metadata_options = options.map(Arc::new);
143        self
144    }
145
146    /// Provide a hint as to the number of bytes needed to fully parse the [`ParquetMetaData`].
147    /// Only used for the asynchronous [`Self::try_load()`] method.
148    ///
149    /// By default, the reader will first fetch the last 8 bytes of the input file to obtain the
150    /// size of the footer metadata. A second fetch will be performed to obtain the needed bytes.
151    /// After parsing the footer metadata, a third fetch will be performed to obtain the bytes
152    /// needed to decode the page index structures, if they have been requested. To avoid
153    /// unnecessary fetches, `prefetch` can be set to an estimate of the number of bytes needed
154    /// to fully decode the [`ParquetMetaData`], which can reduce the number of fetch requests and
155    /// reduce latency. Setting `prefetch` too small will not trigger an error, but will result
156    /// in extra fetches being performed.
157    pub fn with_prefetch_hint(mut self, prefetch: Option<usize>) -> Self {
158        self.prefetch_hint = prefetch;
159        self
160    }
161
162    /// Provide the FileDecryptionProperties to use when decrypting the file.
163    ///
164    /// This is only necessary when the file is encrypted.
165    #[cfg(feature = "encryption")]
166    pub fn with_decryption_properties(
167        mut self,
168        properties: Option<std::sync::Arc<FileDecryptionProperties>>,
169    ) -> Self {
170        self.file_decryption_properties = properties;
171        self
172    }
173
174    /// Indicates whether this reader has a [`ParquetMetaData`] internally.
175    pub fn has_metadata(&self) -> bool {
176        self.metadata.is_some()
177    }
178
179    /// Return the parsed [`ParquetMetaData`] struct, leaving `None` in its place.
180    pub fn finish(&mut self) -> Result<ParquetMetaData> {
181        self.metadata
182            .take()
183            .ok_or_else(|| general_err!("could not parse parquet metadata"))
184    }
185
186    /// Given a [`ChunkReader`], parse and return the [`ParquetMetaData`] in a single pass.
187    ///
188    /// If `reader` is [`Bytes`] based, then the buffer must contain sufficient bytes to complete
189    /// the request, and must include the Parquet footer. If page indexes are desired, the buffer
190    /// must contain the entire file, or [`Self::try_parse_sized()`] should be used.
191    ///
192    /// This call will consume `self`.
193    ///
194    /// # Example
195    /// ```no_run
196    /// # use parquet::file::metadata::{PageIndexPolicy, ParquetMetaDataReader};
197    /// # fn open_parquet_file(path: &str) -> std::fs::File { unimplemented!(); }
198    /// // read parquet metadata including page indexes
199    /// let file = open_parquet_file("some_path.parquet");
200    /// let metadata = ParquetMetaDataReader::new()
201    ///     .with_page_index_policy(PageIndexPolicy::Required)
202    ///     .parse_and_finish(&file).unwrap();
203    /// ```
204    pub fn parse_and_finish<R: ChunkReader>(mut self, reader: &R) -> Result<ParquetMetaData> {
205        self.try_parse(reader)?;
206        self.finish()
207    }
208
209    /// Attempts to parse the footer metadata (and optionally page indexes) given a [`ChunkReader`].
210    ///
211    /// If `reader` is [`Bytes`] based, then the buffer must contain sufficient bytes to complete
212    /// the request, and must include the Parquet footer. If page indexes are desired, the buffer
213    /// must contain the entire file, or [`Self::try_parse_sized()`] should be used.
214    pub fn try_parse<R: ChunkReader>(&mut self, reader: &R) -> Result<()> {
215        self.try_parse_sized(reader, reader.len())
216    }
217
218    /// Same as [`Self::try_parse()`], but provide the original file size in the case that `reader`
219    /// is a [`Bytes`] struct that does not contain the entire file. This information is necessary
220    /// when the page indexes are desired. `reader` must have access to the Parquet footer.
221    ///
222    /// Using this function also allows for retrying with a larger buffer.
223    ///
224    /// # Errors
225    ///
226    /// This function will return [`ParquetError::NeedMoreData`] in the event `reader` does not
227    /// provide enough data to fully parse the metadata (see example below). The returned error
228    /// will be populated with a `usize` field indicating the number of bytes required from the
229    /// tail of the file to completely parse the requested metadata.
230    ///
231    /// Other errors returned include [`ParquetError::General`] and [`ParquetError::EOF`].
232    ///
233    /// # Example
234    /// ```no_run
235    /// # use parquet::file::metadata::{PageIndexPolicy, ParquetMetaDataReader};
236    /// # use parquet::errors::ParquetError;
237    /// # use crate::parquet::file::reader::Length;
238    /// # fn get_bytes(file: &std::fs::File, range: std::ops::Range<u64>) -> bytes::Bytes { unimplemented!(); }
239    /// # fn open_parquet_file(path: &str) -> std::fs::File { unimplemented!(); }
240    /// let file = open_parquet_file("some_path.parquet");
241    /// let len = file.len();
242    /// // Speculatively read 1 kilobyte from the end of the file
243    /// let bytes = get_bytes(&file, len - 1024..len);
244    /// let mut reader = ParquetMetaDataReader::new().with_page_index_policy(PageIndexPolicy::Required);
245    /// match reader.try_parse_sized(&bytes, len) {
246    ///     Ok(_) => (),
247    ///     Err(ParquetError::NeedMoreData(needed)) => {
248    ///         // Read the needed number of bytes from the end of the file
249    ///         let bytes = get_bytes(&file, len - needed as u64..len);
250    ///         reader.try_parse_sized(&bytes, len).unwrap();
251    ///     }
252    ///     _ => panic!("unexpected error")
253    /// }
254    /// let metadata = reader.finish().unwrap();
255    /// ```
256    ///
257    /// Note that it is possible for the file metadata to be completely read, but there are
258    /// insufficient bytes available to read the page indexes. [`Self::has_metadata()`] can be used
259    /// to test for this. In the event the file metadata is present, re-parsing of the file
260    /// metadata can be skipped by using [`Self::read_page_indexes_sized()`], as shown below.
261    /// ```no_run
262    /// # use parquet::file::metadata::{PageIndexPolicy, ParquetMetaDataReader};
263    /// # use parquet::errors::ParquetError;
264    /// # use crate::parquet::file::reader::Length;
265    /// # fn get_bytes(file: &std::fs::File, range: std::ops::Range<u64>) -> bytes::Bytes { unimplemented!(); }
266    /// # fn open_parquet_file(path: &str) -> std::fs::File { unimplemented!(); }
267    /// let file = open_parquet_file("some_path.parquet");
268    /// let len = file.len();
269    /// // Speculatively read 1 kilobyte from the end of the file
270    /// let mut bytes = get_bytes(&file, len - 1024..len);
271    /// let mut reader = ParquetMetaDataReader::new().with_page_index_policy(PageIndexPolicy::Required);
272    /// // Loop until `bytes` is large enough
273    /// loop {
274    ///     match reader.try_parse_sized(&bytes, len) {
275    ///         Ok(_) => break,
276    ///         Err(ParquetError::NeedMoreData(needed)) => {
277    ///             // Read the needed number of bytes from the end of the file
278    ///             bytes = get_bytes(&file, len - needed as u64..len);
279    ///             // If file metadata was read only read page indexes, otherwise continue loop
280    ///             if reader.has_metadata() {
281    ///                 reader.read_page_indexes_sized(&bytes, len).unwrap();
282    ///                 break;
283    ///             }
284    ///         }
285    ///         _ => panic!("unexpected error")
286    ///     }
287    /// }
288    /// let metadata = reader.finish().unwrap();
289    /// ```
290    pub fn try_parse_sized<R: ChunkReader>(&mut self, reader: &R, file_size: u64) -> Result<()> {
291        self.metadata = match self.parse_metadata(reader) {
292            Ok(metadata) => Some(metadata),
293            Err(ParquetError::NeedMoreData(needed)) => {
294                // If reader is the same length as `file_size` then presumably there is no more to
295                // read, so return an EOF error.
296                return if file_size == reader.len() || needed as u64 > file_size {
297                    Err(eof_err!(
298                        "Parquet file too small. Size is {} but need {}",
299                        file_size,
300                        needed
301                    ))
302                } else {
303                    // Ask for a larger buffer
304                    Err(ParquetError::NeedMoreData(needed))
305                };
306            }
307            Err(e) => return Err(e),
308        };
309
310        // we can return if page indexes aren't requested
311        if self.column_index == PageIndexPolicy::Skip && self.offset_index == PageIndexPolicy::Skip
312        {
313            return Ok(());
314        }
315
316        self.read_page_indexes_sized(reader, file_size)
317    }
318
319    /// Read the page index structures when a [`ParquetMetaData`] has already been obtained.
320    /// See [`Self::new_with_metadata()`] and [`Self::has_metadata()`].
321    pub fn read_page_indexes<R: ChunkReader>(&mut self, reader: &R) -> Result<()> {
322        self.read_page_indexes_sized(reader, reader.len())
323    }
324
325    /// Read the page index structures when a [`ParquetMetaData`] has already been obtained.
326    /// This variant is used when `reader` cannot access the entire Parquet file (e.g. it is
327    /// a [`Bytes`] struct containing the tail of the file).
328    /// See [`Self::new_with_metadata()`] and [`Self::has_metadata()`]. Like
329    /// [`Self::try_parse_sized()`] this function may return [`ParquetError::NeedMoreData`].
330    pub fn read_page_indexes_sized<R: ChunkReader>(
331        &mut self,
332        reader: &R,
333        file_size: u64,
334    ) -> Result<()> {
335        let Some(metadata) = self.metadata.take() else {
336            return Err(general_err!(
337                "Tried to read page indexes without ParquetMetaData metadata"
338            ));
339        };
340
341        let push_decoder = ParquetMetaDataPushDecoder::try_new_with_metadata(file_size, metadata)?
342            .with_offset_index_policy(self.offset_index)
343            .with_column_index_policy(self.column_index)
344            .with_metadata_options(self.metadata_options.clone());
345        let mut push_decoder = self.prepare_push_decoder(push_decoder);
346
347        // Get bounds needed for page indexes (if any are present in the file).
348        let range = match needs_index_data(&mut push_decoder)? {
349            NeedsIndexData::No(metadata) => {
350                self.metadata = Some(metadata);
351                return Ok(());
352            }
353            NeedsIndexData::Yes(range) => range,
354        };
355
356        // Check to see if needed range is within `file_range`. Checking `range.end` seems
357        // redundant, but it guards against `range_for_page_index()` returning garbage.
358        let file_range = file_size.saturating_sub(reader.len())..file_size;
359        if !(file_range.contains(&range.start) && file_range.contains(&range.end)) {
360            // Requested range starts beyond EOF
361            return if range.end > file_size {
362                Err(eof_err!(
363                    "Parquet file too small. Range {range:?} is beyond file bounds {file_size}",
364                ))
365            } else {
366                // Ask for a larger buffer
367                Err(ParquetError::NeedMoreData(
368                    (file_size - range.start).try_into()?,
369                ))
370            };
371        }
372
373        // Perform extra sanity check to make sure `range` and the footer metadata don't
374        // overlap.
375        if let Some(metadata_size) = self.metadata_size {
376            let metadata_range = file_size.saturating_sub(metadata_size as u64)..file_size;
377            if range.end > metadata_range.start {
378                return Err(eof_err!(
379                    "Parquet file too small. Page index range {range:?} overlaps with file metadata {metadata_range:?}",
380                ));
381            }
382        }
383
384        // add the needed ranges to the decoder
385        let bytes_needed = usize::try_from(range.end - range.start)?;
386        let bytes = reader.get_bytes(range.start - file_range.start, bytes_needed)?;
387
388        push_decoder.push_range(range, bytes)?;
389        let metadata = parse_index_data(&mut push_decoder)?;
390        self.metadata = Some(metadata);
391
392        Ok(())
393    }
394
395    /// Given a [`MetadataFetch`], parse and return the [`ParquetMetaData`] in a single pass.
396    ///
397    /// This call will consume `self`.
398    ///
399    /// See [`Self::with_prefetch_hint`] for a discussion of how to reduce the number of fetches
400    /// performed by this function.
401    #[cfg(all(feature = "async", feature = "arrow"))]
402    pub async fn load_and_finish<F: MetadataFetch>(
403        mut self,
404        fetch: F,
405        file_size: u64,
406    ) -> Result<ParquetMetaData> {
407        self.try_load(fetch, file_size).await?;
408        self.finish()
409    }
410
411    /// Given a [`MetadataSuffixFetch`], parse and return the [`ParquetMetaData`] in a single pass.
412    ///
413    /// This call will consume `self`.
414    ///
415    /// See [`Self::with_prefetch_hint`] for a discussion of how to reduce the number of fetches
416    /// performed by this function.
417    #[cfg(all(feature = "async", feature = "arrow"))]
418    pub async fn load_via_suffix_and_finish<F: MetadataSuffixFetch>(
419        mut self,
420        fetch: F,
421    ) -> Result<ParquetMetaData> {
422        self.try_load_via_suffix(fetch).await?;
423        self.finish()
424    }
425    /// Attempts to (asynchronously) parse the footer metadata (and optionally page indexes)
426    /// given a [`MetadataFetch`].
427    ///
428    /// See [`Self::with_prefetch_hint`] for a discussion of how to reduce the number of fetches
429    /// performed by this function.
430    #[cfg(all(feature = "async", feature = "arrow"))]
431    pub async fn try_load<F: MetadataFetch>(&mut self, mut fetch: F, file_size: u64) -> Result<()> {
432        let (metadata, remainder) = self.load_metadata(&mut fetch, file_size).await?;
433
434        self.metadata = Some(metadata);
435
436        // we can return if page indexes aren't requested
437        if self.column_index == PageIndexPolicy::Skip && self.offset_index == PageIndexPolicy::Skip
438        {
439            return Ok(());
440        }
441
442        self.load_page_index_with_remainder(fetch, remainder).await
443    }
444
445    /// Attempts to (asynchronously) parse the footer metadata (and optionally page indexes)
446    /// given a [`MetadataSuffixFetch`].
447    ///
448    /// See [`Self::with_prefetch_hint`] for a discussion of how to reduce the number of fetches
449    /// performed by this function.
450    #[cfg(all(feature = "async", feature = "arrow"))]
451    pub async fn try_load_via_suffix<F: MetadataSuffixFetch>(
452        &mut self,
453        mut fetch: F,
454    ) -> Result<()> {
455        let (metadata, remainder) = self.load_metadata_via_suffix(&mut fetch).await?;
456
457        self.metadata = Some(metadata);
458
459        // we can return if page indexes aren't requested
460        if self.column_index == PageIndexPolicy::Skip && self.offset_index == PageIndexPolicy::Skip
461        {
462            return Ok(());
463        }
464
465        self.load_page_index_with_remainder(fetch, remainder).await
466    }
467
468    /// Asynchronously fetch the page index structures when a [`ParquetMetaData`] has already
469    /// been obtained. See [`Self::new_with_metadata()`].
470    #[cfg(all(feature = "async", feature = "arrow"))]
471    pub async fn load_page_index<F: MetadataFetch>(&mut self, fetch: F) -> Result<()> {
472        self.load_page_index_with_remainder(fetch, None).await
473    }
474
475    #[cfg(all(feature = "async", feature = "arrow"))]
476    async fn load_page_index_with_remainder<F: MetadataFetch>(
477        &mut self,
478        mut fetch: F,
479        remainder: Option<(usize, Bytes)>,
480    ) -> Result<()> {
481        let Some(metadata) = self.metadata.take() else {
482            return Err(general_err!("Footer metadata is not present"));
483        };
484
485        // in this case we don't actually know what the file size is, so just use u64::MAX
486        // this is ok since the offsets in the metadata are always valid
487        let file_size = u64::MAX;
488        let push_decoder = ParquetMetaDataPushDecoder::try_new_with_metadata(file_size, metadata)?
489            .with_offset_index_policy(self.offset_index)
490            .with_column_index_policy(self.column_index)
491            .with_metadata_options(self.metadata_options.clone());
492        let mut push_decoder = self.prepare_push_decoder(push_decoder);
493
494        // Get bounds needed for page indexes (if any are present in the file).
495        let range = match needs_index_data(&mut push_decoder)? {
496            NeedsIndexData::No(metadata) => {
497                self.metadata = Some(metadata);
498                return Ok(());
499            }
500            NeedsIndexData::Yes(range) => range,
501        };
502
503        let bytes = match &remainder {
504            Some((remainder_start, remainder)) if *remainder_start as u64 <= range.start => {
505                let remainder_start = *remainder_start as u64;
506                let offset = usize::try_from(range.start - remainder_start)?;
507                let end = usize::try_from(range.end - remainder_start)?;
508                if end > remainder.len() {
509                    return Err(general_err!(
510                        "Corrupted parquet file: index data range ({:?}) exceeds remainder length ({})",
511                        range,
512                        remainder.len()
513                    ));
514                }
515                remainder.slice(offset..end)
516            }
517            // Note: this will potentially fetch data already in remainder, this keeps things simple
518            _ => fetch.fetch(range.clone()).await?,
519        };
520
521        // Sanity check
522        if bytes.len() as u64 != range.end - range.start {
523            return Err(general_err!(
524                "Corrupted parquet file: index data length mismatch, expected {}, got {}",
525                range.end - range.start,
526                bytes.len()
527            ));
528        }
529        push_decoder.push_range(range.clone(), bytes)?;
530        let metadata = parse_index_data(&mut push_decoder)?;
531        self.metadata = Some(metadata);
532        Ok(())
533    }
534
535    // One-shot parse of footer.
536    // Side effect: this will set `self.metadata_size`
537    fn parse_metadata<R: ChunkReader>(&mut self, chunk_reader: &R) -> Result<ParquetMetaData> {
538        // check file is large enough to hold footer
539        let file_size = chunk_reader.len();
540        if file_size < (FOOTER_SIZE as u64) {
541            return Err(ParquetError::NeedMoreData(FOOTER_SIZE));
542        }
543
544        let mut footer = [0_u8; FOOTER_SIZE];
545        chunk_reader
546            .get_read(file_size - FOOTER_SIZE as u64)?
547            .read_exact(&mut footer)?;
548
549        let footer = FooterTail::try_new(&footer)?;
550        let metadata_len = footer.metadata_length();
551        let footer_metadata_len = FOOTER_SIZE + metadata_len;
552        self.metadata_size = Some(footer_metadata_len);
553
554        if footer_metadata_len as u64 > file_size {
555            return Err(ParquetError::NeedMoreData(footer_metadata_len));
556        }
557
558        let start = file_size - footer_metadata_len as u64;
559        let bytes = chunk_reader.get_bytes(start, metadata_len)?;
560        self.decode_footer_metadata(bytes, file_size, footer)
561    }
562
563    /// Size of the serialized thrift metadata plus the 8 byte footer. Only set if
564    /// `self.parse_metadata` is called.
565    pub fn metadata_size(&self) -> Option<usize> {
566        self.metadata_size
567    }
568
569    /// Return the number of bytes to read in the initial pass. If `prefetch_size` has
570    /// been provided, then return that value if it is larger than the size of the Parquet
571    /// file footer (8 bytes). Otherwise returns `8`.
572    #[cfg(all(feature = "async", feature = "arrow"))]
573    fn get_prefetch_size(&self) -> usize {
574        if let Some(prefetch) = self.prefetch_hint
575            && prefetch > FOOTER_SIZE
576        {
577            return prefetch;
578        }
579        FOOTER_SIZE
580    }
581
582    #[cfg(all(feature = "async", feature = "arrow"))]
583    async fn load_metadata<F: MetadataFetch>(
584        &self,
585        fetch: &mut F,
586        file_size: u64,
587    ) -> Result<(ParquetMetaData, Option<(usize, Bytes)>)> {
588        let prefetch = self.get_prefetch_size() as u64;
589
590        if file_size < FOOTER_SIZE as u64 {
591            return Err(eof_err!("file size of {} is less than footer", file_size));
592        }
593
594        // If a size hint is provided, read more than the minimum size
595        // to try and avoid a second fetch.
596        // Note: prefetch > file_size is ok since we're using saturating_sub.
597        let footer_start = file_size.saturating_sub(prefetch);
598
599        let suffix = fetch.fetch(footer_start..file_size).await?;
600        let suffix_len = suffix.len();
601        let fetch_len = (file_size - footer_start)
602            .try_into()
603            .expect("footer size should never be larger than u32");
604        if suffix_len < fetch_len {
605            return Err(eof_err!(
606                "metadata requires {} bytes, but could only read {}",
607                fetch_len,
608                suffix_len
609            ));
610        }
611
612        let mut footer = [0; FOOTER_SIZE];
613        footer.copy_from_slice(&suffix[suffix_len - FOOTER_SIZE..suffix_len]);
614
615        let footer = FooterTail::try_new(&footer)?;
616        let length = footer.metadata_length();
617
618        if file_size < (length + FOOTER_SIZE) as u64 {
619            return Err(eof_err!(
620                "file size of {} is less than footer + metadata {}",
621                file_size,
622                length + FOOTER_SIZE
623            ));
624        }
625
626        // Did not fetch the entire file metadata in the initial read, need to make a second request
627        if length > suffix_len - FOOTER_SIZE {
628            let metadata_start = file_size - (length + FOOTER_SIZE) as u64;
629            let meta = fetch
630                .fetch(metadata_start..(file_size - FOOTER_SIZE as u64))
631                .await?;
632            Ok((self.decode_footer_metadata(meta, file_size, footer)?, None))
633        } else {
634            let metadata_start = (file_size - (length + FOOTER_SIZE) as u64 - footer_start)
635                .try_into()
636                .expect("metadata length should never be larger than u32");
637            let slice = suffix.slice(metadata_start..suffix_len - FOOTER_SIZE);
638            Ok((
639                self.decode_footer_metadata(slice, file_size, footer)?,
640                Some((footer_start as usize, suffix.slice(..metadata_start))),
641            ))
642        }
643    }
644
645    #[cfg(all(feature = "async", feature = "arrow"))]
646    async fn load_metadata_via_suffix<F: MetadataSuffixFetch>(
647        &self,
648        fetch: &mut F,
649    ) -> Result<(ParquetMetaData, Option<(usize, Bytes)>)> {
650        let prefetch = self.get_prefetch_size();
651
652        let suffix = fetch.fetch_suffix(prefetch).await?;
653        let suffix_len = suffix.len();
654
655        if suffix_len < FOOTER_SIZE {
656            return Err(eof_err!(
657                "footer metadata requires {} bytes, but could only read {}",
658                FOOTER_SIZE,
659                suffix_len
660            ));
661        }
662
663        let mut footer = [0; FOOTER_SIZE];
664        footer.copy_from_slice(&suffix[suffix_len - FOOTER_SIZE..suffix_len]);
665
666        let footer = FooterTail::try_new(&footer)?;
667        let length = footer.metadata_length();
668        // fake file size as we are only parsing the footer metadata here
669        // (cant be parsing page indexes without the full file size)
670        let file_size = (length + FOOTER_SIZE) as u64;
671
672        // Did not fetch the entire file metadata in the initial read, need to make a second request
673        let metadata_offset = length + FOOTER_SIZE;
674        if length > suffix_len - FOOTER_SIZE {
675            let meta = fetch.fetch_suffix(metadata_offset).await?;
676
677            if meta.len() < metadata_offset {
678                return Err(eof_err!(
679                    "metadata requires {} bytes, but could only read {}",
680                    metadata_offset,
681                    meta.len()
682                ));
683            }
684
685            // need to slice off the footer or decryption fails
686            let meta = meta.slice(0..length);
687            Ok((self.decode_footer_metadata(meta, file_size, footer)?, None))
688        } else {
689            let metadata_start = suffix_len - metadata_offset;
690            let slice = suffix.slice(metadata_start..suffix_len - FOOTER_SIZE);
691            Ok((
692                self.decode_footer_metadata(slice, file_size, footer)?,
693                Some((0, suffix.slice(..metadata_start))),
694            ))
695        }
696    }
697
698    /// Decodes [`ParquetMetaData`] from the provided bytes.
699    ///
700    /// Typically, this is used to decode the metadata from the end of a parquet
701    /// file. The format of `buf` is the Thrift compact binary protocol, as specified
702    /// by the [Parquet Spec].
703    ///
704    /// It does **NOT** include the 8-byte footer.
705    ///
706    /// This method handles using either `decode_metadata` or
707    /// `decode_metadata_with_encryption` depending on whether the encryption
708    /// feature is enabled.
709    ///
710    /// [Parquet Spec]: https://github.com/apache/parquet-format#metadata
711    pub(crate) fn decode_footer_metadata(
712        &self,
713        buf: Bytes,
714        file_size: u64,
715        footer_tail: FooterTail,
716    ) -> Result<ParquetMetaData> {
717        // The push decoder expects the metadata to be at the end of the file
718        // (... data ...) + (metadata) + (footer)
719        // so we need to provide the starting offset of the metadata
720        // within the file.
721        let ending_offset = file_size.checked_sub(FOOTER_SIZE as u64).ok_or_else(|| {
722            general_err!(
723                "file size {file_size} is smaller than footer size {}",
724                FOOTER_SIZE
725            )
726        })?;
727
728        let starting_offset = ending_offset.checked_sub(buf.len() as u64).ok_or_else(|| {
729            general_err!(
730                "file size {file_size} is smaller than buffer size {} + footer size {}",
731                buf.len(),
732                FOOTER_SIZE
733            )
734        })?;
735
736        let range = starting_offset..ending_offset;
737
738        let push_decoder =
739            ParquetMetaDataPushDecoder::try_new_with_footer_tail(file_size, footer_tail)?
740                // NOTE: DO NOT enable page indexes here, they are handled separately
741                .with_page_index_policy(PageIndexPolicy::Skip)
742                .with_metadata_options(self.metadata_options.clone());
743
744        let mut push_decoder = self.prepare_push_decoder(push_decoder);
745        push_decoder.push_range(range, buf)?;
746        match push_decoder.try_decode()? {
747            DecodeResult::Data(metadata) => Ok(metadata),
748            DecodeResult::Finished => Err(general_err!(
749                "could not parse parquet metadata -- previously finished"
750            )),
751            DecodeResult::NeedsData(ranges) => Err(general_err!(
752                "could not parse parquet metadata, needs ranges {:?}",
753                ranges
754            )),
755        }
756    }
757
758    /// Prepares a push decoder and runs it to decode the metadata.
759    #[cfg(feature = "encryption")]
760    fn prepare_push_decoder(
761        &self,
762        push_decoder: ParquetMetaDataPushDecoder,
763    ) -> ParquetMetaDataPushDecoder {
764        push_decoder.with_file_decryption_properties(
765            self.file_decryption_properties
766                .as_ref()
767                .map(std::sync::Arc::clone),
768        )
769    }
770    #[cfg(not(feature = "encryption"))]
771    fn prepare_push_decoder(
772        &self,
773        push_decoder: ParquetMetaDataPushDecoder,
774    ) -> ParquetMetaDataPushDecoder {
775        push_decoder
776    }
777
778    /// Decodes [`ParquetMetaData`] from the provided bytes.
779    ///
780    /// Typically this is used to decode the metadata from the end of a parquet
781    /// file. The format of `buf` is the Thrift compact binary protocol, as specified
782    /// by the [Parquet Spec].
783    ///
784    /// [Parquet Spec]: https://github.com/apache/parquet-format#metadata
785    pub fn decode_metadata(buf: &[u8]) -> Result<ParquetMetaData> {
786        decode_metadata(buf, None)
787    }
788
789    /// Decodes [`ParquetMetaData`] from the provided bytes.
790    ///
791    /// Like [`Self::decode_metadata`] but this also accepts
792    /// metadata parsing options.
793    pub fn decode_metadata_with_options(
794        buf: &[u8],
795        options: Option<&ParquetMetaDataOptions>,
796    ) -> Result<ParquetMetaData> {
797        decode_metadata(buf, options)
798    }
799
800    /// Decodes the schema from the Parquet footer in `buf`. Returned as
801    /// a [`SchemaDescriptor`].
802    pub fn decode_schema(buf: &[u8]) -> Result<Arc<SchemaDescriptor>> {
803        Ok(Arc::new(parquet_schema_from_bytes(buf)?))
804    }
805}
806
807/// The bounds needed to read page indexes
808// this is an internal enum, so it is ok to allow differences in enum size
809enum NeedsIndexData {
810    /// no additional data is needed (e.g. the indexes weren't requested)
811    No(ParquetMetaData),
812    /// Additional data is needed, with the range that are required
813    Yes(Range<u64>),
814}
815
816/// Determines a single combined range of bytes needed to read the page indexes,
817/// or returns the metadata if no additional data is needed (e.g. if no page indexes are requested)
818fn needs_index_data(push_decoder: &mut ParquetMetaDataPushDecoder) -> Result<NeedsIndexData> {
819    match push_decoder.try_decode()? {
820        DecodeResult::NeedsData(ranges) => {
821            let range = ranges
822                .into_iter()
823                .reduce(|a, b| a.start.min(b.start)..a.end.max(b.end))
824                .ok_or_else(|| general_err!("Internal error: no ranges provided"))?;
825            Ok(NeedsIndexData::Yes(range))
826        }
827        DecodeResult::Data(metadata) => Ok(NeedsIndexData::No(metadata)),
828        DecodeResult::Finished => Err(general_err!("Internal error: decoder was finished")),
829    }
830}
831
832/// Given a push decoder that has had the needed ranges pushed to it,
833/// attempt to decode indexes and return the updated metadata.
834fn parse_index_data(push_decoder: &mut ParquetMetaDataPushDecoder) -> Result<ParquetMetaData> {
835    match push_decoder.try_decode()? {
836        DecodeResult::NeedsData(_) => Err(general_err!(
837            "Internal error: decoder still needs data after reading required range"
838        )),
839        DecodeResult::Data(metadata) => Ok(metadata),
840        DecodeResult::Finished => Err(general_err!("Internal error: decoder was finished")),
841    }
842}
843
844#[cfg(test)]
845mod tests {
846    use super::*;
847    use crate::file::reader::Length;
848    use crate::util::test_common::file_util::get_test_file;
849    use std::ops::Range;
850
851    #[test]
852    fn test_parse_metadata_size_smaller_than_footer() {
853        let test_file = tempfile::tempfile().unwrap();
854        let err = ParquetMetaDataReader::new()
855            .parse_metadata(&test_file)
856            .unwrap_err();
857        assert!(matches!(err, ParquetError::NeedMoreData(FOOTER_SIZE)));
858    }
859
860    #[test]
861    fn test_parse_metadata_corrupt_footer() {
862        let data = Bytes::from(vec![1, 2, 3, 4, 5, 6, 7, 8]);
863        let reader_result = ParquetMetaDataReader::new().parse_metadata(&data);
864        assert_eq!(
865            reader_result.unwrap_err().to_string(),
866            "Parquet error: Invalid Parquet file. Corrupt footer"
867        );
868    }
869
870    #[test]
871    fn test_parse_metadata_invalid_start() {
872        let test_file = Bytes::from(vec![255, 0, 0, 0, b'P', b'A', b'R', b'1']);
873        let err = ParquetMetaDataReader::new()
874            .parse_metadata(&test_file)
875            .unwrap_err();
876        assert!(matches!(err, ParquetError::NeedMoreData(263)));
877    }
878
879    #[test]
880    #[cfg_attr(miri, ignore)] // Takes too long
881    fn test_try_parse() {
882        let file = get_test_file("alltypes_tiny_pages.parquet");
883        let len = file.len();
884
885        let mut reader =
886            ParquetMetaDataReader::new().with_page_index_policy(PageIndexPolicy::Required);
887
888        let bytes_for_range = |range: Range<u64>| {
889            file.get_bytes(range.start, (range.end - range.start).try_into().unwrap())
890                .unwrap()
891        };
892
893        // read entire file
894        let bytes = bytes_for_range(0..len);
895        reader.try_parse(&bytes).unwrap();
896        let metadata = reader.finish().unwrap();
897        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
898
899        // read more than enough of file
900        let bytes = bytes_for_range(320000..len);
901        reader.try_parse_sized(&bytes, len).unwrap();
902        let metadata = reader.finish().unwrap();
903        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
904
905        // exactly enough
906        let bytes = bytes_for_range(323583..len);
907        reader.try_parse_sized(&bytes, len).unwrap();
908        let metadata = reader.finish().unwrap();
909        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
910
911        // not enough for page index
912        let bytes = bytes_for_range(323584..len);
913        // should fail
914        match reader.try_parse_sized(&bytes, len).unwrap_err() {
915            // expected error, try again with provided bounds
916            ParquetError::NeedMoreData(needed) => {
917                let bytes = bytes_for_range(len - needed as u64..len);
918                reader.try_parse_sized(&bytes, len).unwrap();
919                let metadata = reader.finish().unwrap();
920                assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
921            }
922            _ => panic!("unexpected error"),
923        }
924
925        // not enough for file metadata, but keep trying until page indexes are read
926        let mut reader =
927            ParquetMetaDataReader::new().with_page_index_policy(PageIndexPolicy::Required);
928        let mut bytes = bytes_for_range(452505..len);
929        loop {
930            match reader.try_parse_sized(&bytes, len) {
931                Ok(()) => break,
932                Err(ParquetError::NeedMoreData(needed)) => {
933                    bytes = bytes_for_range(len - needed as u64..len);
934                    if reader.has_metadata() {
935                        reader.read_page_indexes_sized(&bytes, len).unwrap();
936                        break;
937                    }
938                }
939                _ => panic!("unexpected error"),
940            }
941        }
942        let metadata = reader.finish().unwrap();
943        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
944
945        // not enough for page index but lie about file size
946        let bytes = bytes_for_range(323584..len);
947        let reader_result = reader.try_parse_sized(&bytes, len - 323584).unwrap_err();
948        assert_eq!(
949            reader_result.to_string(),
950            "EOF: Parquet file too small. Range 323583..452504 is beyond file bounds 130649"
951        );
952
953        // not enough for file metadata
954        let mut reader = ParquetMetaDataReader::new();
955        let bytes = bytes_for_range(452505..len);
956        // should fail
957        match reader.try_parse_sized(&bytes, len).unwrap_err() {
958            // expected error, try again with provided bounds
959            ParquetError::NeedMoreData(needed) => {
960                let bytes = bytes_for_range(len - needed as u64..len);
961                reader.try_parse_sized(&bytes, len).unwrap();
962                reader.finish().unwrap();
963            }
964            _ => panic!("unexpected error"),
965        }
966
967        // not enough for file metadata but use try_parse()
968        let reader_result = reader.try_parse(&bytes).unwrap_err();
969        assert_eq!(
970            reader_result.to_string(),
971            "EOF: Parquet file too small. Size is 1728 but need 1729"
972        );
973
974        // read head of file rather than tail
975        let bytes = bytes_for_range(0..1000);
976        let reader_result = reader.try_parse_sized(&bytes, len).unwrap_err();
977        assert_eq!(
978            reader_result.to_string(),
979            "Parquet error: Invalid Parquet file. Corrupt footer"
980        );
981
982        // lie about file size
983        let bytes = bytes_for_range(452510..len);
984        let reader_result = reader.try_parse_sized(&bytes, len - 452505).unwrap_err();
985        assert_eq!(
986            reader_result.to_string(),
987            "EOF: Parquet file too small. Size is 1728 but need 1729"
988        );
989    }
990}
991
992#[cfg(all(feature = "async", feature = "arrow", test))]
993mod async_tests {
994    use super::*;
995
996    use arrow::{array::Int32Array, datatypes::DataType};
997    use arrow_array::RecordBatch;
998    use arrow_schema::{Field, Schema};
999    use bytes::Bytes;
1000    use futures::FutureExt;
1001    use futures::future::BoxFuture;
1002    use std::fs::File;
1003    use std::future::Future;
1004    use std::io::{Read, Seek, SeekFrom};
1005    use std::ops::Range;
1006    use std::sync::Arc;
1007    use std::sync::atomic::{AtomicUsize, Ordering};
1008    use tempfile::NamedTempFile;
1009
1010    use crate::arrow::ArrowWriter;
1011    use crate::file::properties::WriterProperties;
1012    use crate::file::reader::Length;
1013    use crate::util::test_common::file_util::get_test_file;
1014
1015    struct MetadataFetchFn<F>(F);
1016
1017    impl<F, Fut> MetadataFetch for MetadataFetchFn<F>
1018    where
1019        F: FnMut(Range<u64>) -> Fut + Send,
1020        Fut: Future<Output = Result<Bytes>> + Send,
1021    {
1022        fn fetch(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
1023            async move { self.0(range).await }.boxed()
1024        }
1025    }
1026
1027    struct MetadataSuffixFetchFn<F1, F2>(F1, F2);
1028
1029    impl<F1, Fut, F2> MetadataFetch for MetadataSuffixFetchFn<F1, F2>
1030    where
1031        F1: FnMut(Range<u64>) -> Fut + Send,
1032        Fut: Future<Output = Result<Bytes>> + Send,
1033        F2: Send,
1034    {
1035        fn fetch(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
1036            async move { self.0(range).await }.boxed()
1037        }
1038    }
1039
1040    impl<F1, Fut, F2> MetadataSuffixFetch for MetadataSuffixFetchFn<F1, F2>
1041    where
1042        F1: FnMut(Range<u64>) -> Fut + Send,
1043        F2: FnMut(usize) -> Fut + Send,
1044        Fut: Future<Output = Result<Bytes>> + Send,
1045    {
1046        fn fetch_suffix(&mut self, suffix: usize) -> BoxFuture<'_, Result<Bytes>> {
1047            async move { self.1(suffix).await }.boxed()
1048        }
1049    }
1050
1051    fn read_range(file: &mut File, range: Range<u64>) -> Result<Bytes> {
1052        file.seek(SeekFrom::Start(range.start))?;
1053        let len = range.end - range.start;
1054        let mut buf = Vec::with_capacity(len.try_into().unwrap());
1055        file.take(len).read_to_end(&mut buf)?;
1056        Ok(buf.into())
1057    }
1058
1059    fn read_suffix(file: &mut File, suffix: usize) -> Result<Bytes> {
1060        let file_len = file.len();
1061        // Don't seek before beginning of file
1062        file.seek(SeekFrom::End(0 - suffix.min(file_len as _) as i64))?;
1063        let mut buf = Vec::with_capacity(suffix);
1064        file.take(suffix as _).read_to_end(&mut buf)?;
1065        Ok(buf.into())
1066    }
1067
1068    #[tokio::test]
1069    async fn test_simple() {
1070        let mut file = get_test_file("nulls.snappy.parquet");
1071        let len = file.len();
1072
1073        let expected = ParquetMetaDataReader::new()
1074            .parse_and_finish(&file)
1075            .unwrap();
1076        let expected = expected.file_metadata().schema();
1077        let fetch_count = AtomicUsize::new(0);
1078
1079        let mut fetch = |range| {
1080            fetch_count.fetch_add(1, Ordering::SeqCst);
1081            futures::future::ready(read_range(&mut file, range))
1082        };
1083
1084        let input = MetadataFetchFn(&mut fetch);
1085        let actual = ParquetMetaDataReader::new()
1086            .load_and_finish(input, len)
1087            .await
1088            .unwrap();
1089        assert_eq!(actual.file_metadata().schema(), expected);
1090        assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1091
1092        // Metadata hint too small - below footer size
1093        fetch_count.store(0, Ordering::SeqCst);
1094        let input = MetadataFetchFn(&mut fetch);
1095        let actual = ParquetMetaDataReader::new()
1096            .with_prefetch_hint(Some(7))
1097            .load_and_finish(input, len)
1098            .await
1099            .unwrap();
1100        assert_eq!(actual.file_metadata().schema(), expected);
1101        assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1102
1103        // Metadata hint too small
1104        fetch_count.store(0, Ordering::SeqCst);
1105        let input = MetadataFetchFn(&mut fetch);
1106        let actual = ParquetMetaDataReader::new()
1107            .with_prefetch_hint(Some(10))
1108            .load_and_finish(input, len)
1109            .await
1110            .unwrap();
1111        assert_eq!(actual.file_metadata().schema(), expected);
1112        assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1113
1114        // Metadata hint too large
1115        fetch_count.store(0, Ordering::SeqCst);
1116        let input = MetadataFetchFn(&mut fetch);
1117        let actual = ParquetMetaDataReader::new()
1118            .with_prefetch_hint(Some(500))
1119            .load_and_finish(input, len)
1120            .await
1121            .unwrap();
1122        assert_eq!(actual.file_metadata().schema(), expected);
1123        assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1124
1125        // Metadata hint exactly correct
1126        fetch_count.store(0, Ordering::SeqCst);
1127        let input = MetadataFetchFn(&mut fetch);
1128        let actual = ParquetMetaDataReader::new()
1129            .with_prefetch_hint(Some(428))
1130            .load_and_finish(input, len)
1131            .await
1132            .unwrap();
1133        assert_eq!(actual.file_metadata().schema(), expected);
1134        assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1135
1136        let input = MetadataFetchFn(&mut fetch);
1137        let err = ParquetMetaDataReader::new()
1138            .load_and_finish(input, 4)
1139            .await
1140            .unwrap_err()
1141            .to_string();
1142        assert_eq!(err, "EOF: file size of 4 is less than footer");
1143
1144        let input = MetadataFetchFn(&mut fetch);
1145        let err = ParquetMetaDataReader::new()
1146            .load_and_finish(input, 20)
1147            .await
1148            .unwrap_err()
1149            .to_string();
1150        assert_eq!(err, "Parquet error: Invalid Parquet file. Corrupt footer");
1151    }
1152
1153    #[tokio::test]
1154    async fn test_suffix() {
1155        let mut file = get_test_file("nulls.snappy.parquet");
1156        let mut file2 = file.try_clone().unwrap();
1157
1158        let expected = ParquetMetaDataReader::new()
1159            .parse_and_finish(&file)
1160            .unwrap();
1161        let expected = expected.file_metadata().schema();
1162        let fetch_count = AtomicUsize::new(0);
1163        let suffix_fetch_count = AtomicUsize::new(0);
1164
1165        let mut fetch = |range| {
1166            fetch_count.fetch_add(1, Ordering::SeqCst);
1167            futures::future::ready(read_range(&mut file, range))
1168        };
1169        let mut suffix_fetch = |suffix| {
1170            suffix_fetch_count.fetch_add(1, Ordering::SeqCst);
1171            futures::future::ready(read_suffix(&mut file2, suffix))
1172        };
1173
1174        let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1175        let actual = ParquetMetaDataReader::new()
1176            .load_via_suffix_and_finish(input)
1177            .await
1178            .unwrap();
1179        assert_eq!(actual.file_metadata().schema(), expected);
1180        assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1181        assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 2);
1182
1183        // Metadata hint too small - below footer size
1184        fetch_count.store(0, Ordering::SeqCst);
1185        suffix_fetch_count.store(0, Ordering::SeqCst);
1186        let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1187        let actual = ParquetMetaDataReader::new()
1188            .with_prefetch_hint(Some(7))
1189            .load_via_suffix_and_finish(input)
1190            .await
1191            .unwrap();
1192        assert_eq!(actual.file_metadata().schema(), expected);
1193        assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1194        assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 2);
1195
1196        // Metadata hint too small
1197        fetch_count.store(0, Ordering::SeqCst);
1198        suffix_fetch_count.store(0, Ordering::SeqCst);
1199        let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1200        let actual = ParquetMetaDataReader::new()
1201            .with_prefetch_hint(Some(10))
1202            .load_via_suffix_and_finish(input)
1203            .await
1204            .unwrap();
1205        assert_eq!(actual.file_metadata().schema(), expected);
1206        assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1207        assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 2);
1208
1209        // Metadata hint too large
1210        fetch_count.store(0, Ordering::SeqCst);
1211        suffix_fetch_count.store(0, Ordering::SeqCst);
1212        let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1213        let actual = ParquetMetaDataReader::new()
1214            .with_prefetch_hint(Some(500))
1215            .load_via_suffix_and_finish(input)
1216            .await
1217            .unwrap();
1218        assert_eq!(actual.file_metadata().schema(), expected);
1219        assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1220        assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 1);
1221
1222        // Metadata hint exactly correct
1223        fetch_count.store(0, Ordering::SeqCst);
1224        suffix_fetch_count.store(0, Ordering::SeqCst);
1225        let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1226        let actual = ParquetMetaDataReader::new()
1227            .with_prefetch_hint(Some(428))
1228            .load_via_suffix_and_finish(input)
1229            .await
1230            .unwrap();
1231        assert_eq!(actual.file_metadata().schema(), expected);
1232        assert_eq!(fetch_count.load(Ordering::SeqCst), 0);
1233        assert_eq!(suffix_fetch_count.load(Ordering::SeqCst), 1);
1234    }
1235
1236    #[cfg(feature = "encryption")]
1237    #[tokio::test]
1238    async fn test_suffix_with_encryption() {
1239        let mut file = get_test_file("uniform_encryption.parquet.encrypted");
1240        let mut file2 = file.try_clone().unwrap();
1241
1242        let mut fetch = |range| futures::future::ready(read_range(&mut file, range));
1243        let mut suffix_fetch = |suffix| futures::future::ready(read_suffix(&mut file2, suffix));
1244
1245        let input = MetadataSuffixFetchFn(&mut fetch, &mut suffix_fetch);
1246
1247        let key_code: &[u8] = b"0123456789012345";
1248        let decryption_properties = FileDecryptionProperties::builder(key_code.to_vec())
1249            .build()
1250            .unwrap();
1251
1252        // just make sure the metadata is properly decrypted and read
1253        let expected = ParquetMetaDataReader::new()
1254            .with_decryption_properties(Some(decryption_properties))
1255            .load_via_suffix_and_finish(input)
1256            .await
1257            .unwrap();
1258        assert_eq!(expected.num_row_groups(), 1);
1259    }
1260
1261    #[tokio::test]
1262    async fn test_page_index() {
1263        let mut file = get_test_file("alltypes_tiny_pages.parquet");
1264        let len = file.len();
1265        let fetch_count = AtomicUsize::new(0);
1266        let mut fetch = |range| {
1267            fetch_count.fetch_add(1, Ordering::SeqCst);
1268            futures::future::ready(read_range(&mut file, range))
1269        };
1270
1271        let f = MetadataFetchFn(&mut fetch);
1272        let mut loader =
1273            ParquetMetaDataReader::new().with_page_index_policy(PageIndexPolicy::Required);
1274        loader.try_load(f, len).await.unwrap();
1275        assert_eq!(fetch_count.load(Ordering::SeqCst), 3);
1276        let metadata = loader.finish().unwrap();
1277        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
1278
1279        // Prefetch just footer exactly
1280        fetch_count.store(0, Ordering::SeqCst);
1281        let f = MetadataFetchFn(&mut fetch);
1282        let mut loader = ParquetMetaDataReader::new()
1283            .with_page_index_policy(PageIndexPolicy::Required)
1284            .with_prefetch_hint(Some(1729));
1285        loader.try_load(f, len).await.unwrap();
1286        assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1287        let metadata = loader.finish().unwrap();
1288        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
1289
1290        // Prefetch more than footer but not enough
1291        fetch_count.store(0, Ordering::SeqCst);
1292        let f = MetadataFetchFn(&mut fetch);
1293        let mut loader = ParquetMetaDataReader::new()
1294            .with_page_index_policy(PageIndexPolicy::Required)
1295            .with_prefetch_hint(Some(130649));
1296        loader.try_load(f, len).await.unwrap();
1297        assert_eq!(fetch_count.load(Ordering::SeqCst), 2);
1298        let metadata = loader.finish().unwrap();
1299        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
1300
1301        // Prefetch exactly enough
1302        fetch_count.store(0, Ordering::SeqCst);
1303        let f = MetadataFetchFn(&mut fetch);
1304        let metadata = ParquetMetaDataReader::new()
1305            .with_page_index_policy(PageIndexPolicy::Required)
1306            .with_prefetch_hint(Some(130650))
1307            .load_and_finish(f, len)
1308            .await
1309            .unwrap();
1310        assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1311        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
1312
1313        // Prefetch more than enough but less than the entire file
1314        fetch_count.store(0, Ordering::SeqCst);
1315        let f = MetadataFetchFn(&mut fetch);
1316        let metadata = ParquetMetaDataReader::new()
1317            .with_page_index_policy(PageIndexPolicy::Required)
1318            .with_prefetch_hint(Some((len - 1000) as usize)) // prefetch entire file
1319            .load_and_finish(f, len)
1320            .await
1321            .unwrap();
1322        assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1323        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
1324
1325        // Prefetch the entire file
1326        fetch_count.store(0, Ordering::SeqCst);
1327        let f = MetadataFetchFn(&mut fetch);
1328        let metadata = ParquetMetaDataReader::new()
1329            .with_page_index_policy(PageIndexPolicy::Required)
1330            .with_prefetch_hint(Some(len as usize)) // prefetch entire file
1331            .load_and_finish(f, len)
1332            .await
1333            .unwrap();
1334        assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1335        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
1336
1337        // Prefetch more than the entire file
1338        fetch_count.store(0, Ordering::SeqCst);
1339        let f = MetadataFetchFn(&mut fetch);
1340        let metadata = ParquetMetaDataReader::new()
1341            .with_page_index_policy(PageIndexPolicy::Required)
1342            .with_prefetch_hint(Some((len + 1000) as usize)) // prefetch entire file
1343            .load_and_finish(f, len)
1344            .await
1345            .unwrap();
1346        assert_eq!(fetch_count.load(Ordering::SeqCst), 1);
1347        assert!(metadata.page_index().is_some_and(|idx| idx.is_complete()));
1348    }
1349
1350    fn write_parquet_file(offset_index_disabled: bool) -> Result<NamedTempFile> {
1351        let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)]));
1352        let batch = RecordBatch::try_new(
1353            schema.clone(),
1354            vec![Arc::new(Int32Array::from(vec![1, 2, 3]))],
1355        )?;
1356
1357        let file = NamedTempFile::new().unwrap();
1358
1359        // Write properties with page index disabled
1360        let props = WriterProperties::builder()
1361            .set_offset_index_disabled(offset_index_disabled)
1362            .build();
1363
1364        let mut writer = ArrowWriter::try_new(file.reopen()?, schema, Some(props))?;
1365        writer.write(&batch)?;
1366        writer.close()?;
1367
1368        Ok(file)
1369    }
1370
1371    fn read_and_check(file: &File, policy: PageIndexPolicy) -> Result<ParquetMetaData> {
1372        let mut reader = ParquetMetaDataReader::new().with_page_index_policy(policy);
1373        reader.try_parse(file)?;
1374        reader.finish()
1375    }
1376
1377    #[test]
1378    fn test_page_index_policy() {
1379        // With page index
1380        let f = write_parquet_file(false).unwrap();
1381        read_and_check(f.as_file(), PageIndexPolicy::Required).unwrap();
1382        read_and_check(f.as_file(), PageIndexPolicy::Optional).unwrap();
1383        read_and_check(f.as_file(), PageIndexPolicy::Skip).unwrap();
1384
1385        // Without page index
1386        let f = write_parquet_file(true).unwrap();
1387        let res = read_and_check(f.as_file(), PageIndexPolicy::Required);
1388        assert!(matches!(
1389            res,
1390            Err(ParquetError::General(e)) if e == "missing offset index"
1391        ));
1392        read_and_check(f.as_file(), PageIndexPolicy::Optional).unwrap();
1393        read_and_check(f.as_file(), PageIndexPolicy::Skip).unwrap();
1394    }
1395}