This is an automated email from the ASF dual-hosted git repository.

alamb pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git


The following commit(s) were added to refs/heads/main by this push:
     new 53620e0880 parquet: return an error instead of panicking when a pushed 
buffer's length does not match its range (#10564)
53620e0880 is described below

commit 53620e08800090cababb9e4529f2567198a102c2
Author: ranflarion <[email protected]>
AuthorDate: Thu Aug 6 09:25:35 2026 -0400

    parquet: return an error instead of panicking when a pushed buffer's length 
does not match its range (#10564)
    
    # Which issue does this PR close?
    
    Closes #10563.
    
    # Rationale for this change
    
    A short read pushed into the parquet push decoder panics at
    `PushBuffers::push_range` ("Range length must match buffer length")
    although every public entry point that feeds it returns `Result`.
    Through the async reader this turns a transient store fault into a
    reader-thread panic the caller cannot classify or retry; the json and
    avro readers surface the equivalent fault as a decode error.
    
    # What changes are included in this PR?
    
    `PushBuffers::push_range` and `push_ranges` return `Result<(),
    ParquetError>` instead of asserting (buffer/range length mismatch, and
    ranges/buffers count mismatch in `push_ranges`). The error propagates
    through the crate-internal chain:
    `ParquetMetaDataPushDecoder::push_range`/`push_ranges` (already
    `Result`, now use `?`), and `ParquetDecoderState::push_data` ->
    `RemainingRowGroups::push_data` -> `RowGroupReaderBuilder::push_data`,
    the last two becoming fallible; all their callers were already in
    `Result` contexts.
    
    # Are these changes tested?
    
    New unit tests in `push_buffers.rs` cover the accepted case, the
    short-buffer error, and the count-mismatch error. The existing parquet
    test suite passes.
    
    # Are there any user-facing changes?
    
    No stable API change: the public decoder entry points keep their
    `Result` signatures. `PushBuffers` (exported only under the
    `experimental` feature, which carries no stability guarantees) changes
    `push_range`/`push_ranges` from `fn(...)` to `fn(...) -> Result<(),
    ParquetError>`. Code that previously panicked on mismatched pushes now
    receives an `Err`.
---
 parquet/src/arrow/push_decoder/mod.rs              |  4 +-
 .../src/arrow/push_decoder/reader_builder/mod.rs   |  8 +-
 parquet/src/arrow/push_decoder/remaining.rs        |  8 +-
 parquet/src/file/metadata/push_decoder.rs          |  4 +-
 parquet/src/util/push_buffers.rs                   | 87 ++++++++++++++++++----
 5 files changed, 90 insertions(+), 21 deletions(-)

diff --git a/parquet/src/arrow/push_decoder/mod.rs 
b/parquet/src/arrow/push_decoder/mod.rs
index 93994071fc..54f4e2d331 100644
--- a/parquet/src/arrow/push_decoder/mod.rs
+++ b/parquet/src/arrow/push_decoder/mod.rs
@@ -788,7 +788,7 @@ impl ParquetDecoderState {
                 mut remaining_row_groups,
             } => {
                 // Push data to the RowGroupReaderBuilder
-                remaining_row_groups.push_data(ranges, data);
+                remaining_row_groups.push_data(ranges, data)?;
                 Ok(ParquetDecoderState::ReadingRowGroup {
                     remaining_row_groups,
                 })
@@ -798,7 +798,7 @@ impl ParquetDecoderState {
                 record_batch_reader,
                 mut remaining_row_groups,
             } => {
-                remaining_row_groups.push_data(ranges, data);
+                remaining_row_groups.push_data(ranges, data)?;
                 Ok(ParquetDecoderState::DecodingRowGroup {
                     record_batch_reader,
                     remaining_row_groups,
diff --git a/parquet/src/arrow/push_decoder/reader_builder/mod.rs 
b/parquet/src/arrow/push_decoder/reader_builder/mod.rs
index 60b987828b..ffc1038244 100644
--- a/parquet/src/arrow/push_decoder/reader_builder/mod.rs
+++ b/parquet/src/arrow/push_decoder/reader_builder/mod.rs
@@ -351,8 +351,12 @@ impl RowGroupReaderBuilder {
     }
 
     /// Push new data buffers that can be used to satisfy pending requests
-    pub fn push_data(&mut self, ranges: Vec<Range<u64>>, buffers: Vec<Bytes>) {
-        self.buffers.push_ranges(ranges, buffers);
+    pub fn push_data(
+        &mut self,
+        ranges: Vec<Range<u64>>,
+        buffers: Vec<Bytes>,
+    ) -> Result<(), ParquetError> {
+        self.buffers.push_ranges(ranges, buffers)
     }
 
     /// True iff the inner state is `Finished`. This is the only state in
diff --git a/parquet/src/arrow/push_decoder/remaining.rs 
b/parquet/src/arrow/push_decoder/remaining.rs
index d2658189b7..09a1cdc95b 100644
--- a/parquet/src/arrow/push_decoder/remaining.rs
+++ b/parquet/src/arrow/push_decoder/remaining.rs
@@ -290,8 +290,12 @@ impl RemainingRowGroups {
     }
 
     /// Push new data buffers that can be used to satisfy pending requests
-    pub fn push_data(&mut self, ranges: Vec<Range<u64>>, buffers: Vec<Bytes>) {
-        self.row_group_reader_builder.push_data(ranges, buffers);
+    pub fn push_data(
+        &mut self,
+        ranges: Vec<Range<u64>>,
+        buffers: Vec<Bytes>,
+    ) -> Result<(), ParquetError> {
+        self.row_group_reader_builder.push_data(ranges, buffers)
     }
 
     /// Return the total number of bytes buffered so far
diff --git a/parquet/src/file/metadata/push_decoder.rs 
b/parquet/src/file/metadata/push_decoder.rs
index de812b3e87..4212b58ce4 100644
--- a/parquet/src/file/metadata/push_decoder.rs
+++ b/parquet/src/file/metadata/push_decoder.rs
@@ -343,7 +343,7 @@ impl ParquetMetaDataPushDecoder {
                 "ParquetMetaDataPushDecoder: cannot push data after decoding 
is finished"
             ));
         }
-        self.buffers.push_ranges(ranges, buffers);
+        self.buffers.push_ranges(ranges, buffers)?;
         Ok(())
     }
 
@@ -354,7 +354,7 @@ impl ParquetMetaDataPushDecoder {
                 "ParquetMetaDataPushDecoder: cannot push data after decoding 
is finished"
             ));
         }
-        self.buffers.push_range(range, buffer);
+        self.buffers.push_range(range, buffer)?;
         Ok(())
     }
 
diff --git a/parquet/src/util/push_buffers.rs b/parquet/src/util/push_buffers.rs
index 32c9ba8304..d2404d50ab 100644
--- a/parquet/src/util/push_buffers.rs
+++ b/parquet/src/util/push_buffers.rs
@@ -93,26 +93,48 @@ impl PushBuffers {
     }
 
     /// Push all the ranges and buffers
-    pub fn push_ranges(&mut self, ranges: Vec<Range<u64>>, buffers: 
Vec<Bytes>) {
-        assert_eq!(
-            ranges.len(),
-            buffers.len(),
-            "Number of ranges must match number of buffers"
-        );
+    ///
+    /// # Errors
+    /// Returns an error if the number of ranges does not match the number of
+    /// buffers, or if any buffer's length does not match its range (see
+    /// [`Self::push_range`]).
+    pub fn push_ranges(
+        &mut self,
+        ranges: Vec<Range<u64>>,
+        buffers: Vec<Bytes>,
+    ) -> Result<(), ParquetError> {
+        if ranges.len() != buffers.len() {
+            return Err(general_err!(
+                "Number of ranges ({}) must match number of buffers ({})",
+                ranges.len(),
+                buffers.len()
+            ));
+        }
         for (range, buffer) in ranges.into_iter().zip(buffers) {
-            self.push_range(range, buffer);
+            self.push_range(range, buffer)?;
         }
+        Ok(())
     }
 
     /// Push a new range and its associated buffer
-    pub fn push_range(&mut self, range: Range<u64>, buffer: Bytes) {
-        assert_eq!(
-            (range.end - range.start) as usize,
-            buffer.len(),
-            "Range length must match buffer length"
-        );
+    ///
+    /// # Errors
+    /// Returns an error if the buffer's length does not match the range's
+    /// length, e.g. when a truncated (short) read is pushed.
+    pub fn push_range(&mut self, range: Range<u64>, buffer: Bytes) -> 
Result<(), ParquetError> {
+        let expected = range.end.saturating_sub(range.start);
+        if expected != buffer.len() as u64 {
+            return Err(general_err!(
+                "Buffer length ({}) does not match length ({}) of range 
{}..{}",
+                buffer.len(),
+                expected,
+                range.start,
+                range.end
+            ));
+        }
         self.ranges.push(range);
         self.buffers.push(buffer);
+        Ok(())
     }
 
     /// Returns true if the Buffers contains data for the given range
@@ -226,3 +248,42 @@ impl ChunkReader for PushBuffers {
         Err(ParquetError::NeedMoreDataRange(start..requested_end))
     }
 }
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+
+    #[test]
+    fn push_range_accepts_matching_length() {
+        let mut buffers = PushBuffers::new(100);
+        buffers
+            .push_range(10..14, Bytes::from_static(b"abcd"))
+            .unwrap();
+        assert!(buffers.has_range(&(10..14)));
+    }
+
+    #[test]
+    fn push_range_rejects_short_buffer() {
+        let mut buffers = PushBuffers::new(100);
+        let err = buffers
+            .push_range(10..20, Bytes::from_static(b"abcd"))
+            .unwrap_err();
+        assert_eq!(
+            err.to_string(),
+            "Parquet error: Buffer length (4) does not match length (10) of 
range 10..20"
+        );
+        assert!(!buffers.has_range(&(10..20)));
+    }
+
+    #[test]
+    fn push_ranges_rejects_mismatched_counts() {
+        let mut buffers = PushBuffers::new(100);
+        let err = buffers
+            .push_ranges(vec![0..4, 4..8], vec![Bytes::from_static(b"abcd")])
+            .unwrap_err();
+        assert_eq!(
+            err.to_string(),
+            "Parquet error: Number of ranges (2) must match number of buffers 
(1)"
+        );
+    }
+}

Reply via email to