Skip to content

Commit 47f989b

Browse files
fix(parquet/metadata): return page index decode errors (#1032)
### Rationale for this change `RowGroupPageIndexReader` exposes error-returning methods, but malformed serialized column or offset indexes could panic inside the lower-level constructors. A corrupt Parquet file could therefore terminate the caller instead of producing a read error. ### What changes are included in this PR? * Convert page-index deserialization panics into errors at the row-group reader boundary. * Wrap failures with `arrow.ErrInvalid` and identify the affected index kind. * Preserve the existing lower-level constructor signatures for compatibility. ### Are these changes tested? Yes. Regression coverage exercises malformed serialized column and offset indexes. The full `parquet/metadata` package passes.
1 parent 970a12b commit 47f989b

2 files changed

Lines changed: 119 additions & 2 deletions

File tree

parquet/metadata/page_index.go

Lines changed: 35 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ import (
2020
"fmt"
2121
"io"
2222
"math"
23+
"runtime"
2324
"sync"
2425

2526
"github.com/apache/arrow-go/v18/arrow"
@@ -244,6 +245,32 @@ func NewOffsetIndex(serializedIndex []byte, _ *parquet.ReaderProperties, decrypt
244245
return &offsetIndex
245246
}
246247

248+
func deserializeColumnIndex(descr *schema.Column, serializedIndex []byte, props *parquet.ReaderProperties) (idx ColumnIndex, err error) {
249+
defer func() {
250+
if recovered := recover(); recovered != nil {
251+
if _, ok := recovered.(runtime.Error); ok {
252+
panic(recovered)
253+
}
254+
err = fmt.Errorf("%w: %w", arrow.ErrInvalid,
255+
shared_utils.FormatRecoveredError("invalid column index", recovered))
256+
}
257+
}()
258+
return NewColumnIndex(descr, serializedIndex, props, nil), nil
259+
}
260+
261+
func deserializeOffsetIndex(serializedIndex []byte, props *parquet.ReaderProperties) (idx OffsetIndex, err error) {
262+
defer func() {
263+
if recovered := recover(); recovered != nil {
264+
if _, ok := recovered.(runtime.Error); ok {
265+
panic(recovered)
266+
}
267+
err = fmt.Errorf("%w: %w", arrow.ErrInvalid,
268+
shared_utils.FormatRecoveredError("invalid offset index", recovered))
269+
}
270+
}()
271+
return NewOffsetIndex(serializedIndex, props, nil), nil
272+
}
273+
247274
type readRange struct {
248275
Offset, Length int64
249276
}
@@ -358,7 +385,10 @@ func (r *RowGroupPageIndexReader) GetColumnIndex(i int) (ColumnIndex, error) {
358385
serializedIndex = decrypted
359386
}
360387

361-
idx := NewColumnIndex(descr, serializedIndex, r.props, nil)
388+
idx, err := deserializeColumnIndex(descr, serializedIndex, r.props)
389+
if err != nil {
390+
return nil, err
391+
}
362392
r.colIndexes[i] = idx
363393
return idx, nil
364394
}
@@ -418,7 +448,10 @@ func (r *RowGroupPageIndexReader) GetOffsetIndex(i int) (OffsetIndex, error) {
418448
serializedIndex = decrypted
419449
}
420450

421-
oidx := NewOffsetIndex(serializedIndex, r.props, nil)
451+
oidx, err := deserializeOffsetIndex(serializedIndex, r.props)
452+
if err != nil {
453+
return nil, err
454+
}
422455
r.offsetIndices[i] = oidx
423456
return oidx, nil
424457
}
Lines changed: 84 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,84 @@
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+
package metadata
19+
20+
import (
21+
"bytes"
22+
"errors"
23+
"testing"
24+
25+
"github.com/apache/arrow-go/v18/arrow"
26+
"github.com/apache/arrow-go/v18/parquet"
27+
format "github.com/apache/arrow-go/v18/parquet/internal/gen-go/parquet"
28+
"github.com/apache/arrow-go/v18/parquet/internal/thrift"
29+
"github.com/apache/arrow-go/v18/parquet/schema"
30+
"github.com/stretchr/testify/assert"
31+
"github.com/stretchr/testify/require"
32+
)
33+
34+
func TestDeserializePageIndexReturnsErrors(t *testing.T) {
35+
descr := schema.NewColumn(schema.NewInt32Node("values", parquet.Repetitions.Required, -1), 0, 0)
36+
37+
_, err := deserializeColumnIndex(descr, []byte{0xff}, nil)
38+
assert.Error(t, err)
39+
assert.True(t, errors.Is(err, arrow.ErrInvalid))
40+
41+
_, err = deserializeOffsetIndex([]byte{0xff}, nil)
42+
assert.Error(t, err)
43+
assert.True(t, errors.Is(err, arrow.ErrInvalid))
44+
}
45+
46+
func TestRowGroupPageIndexReaderReturnsDecodeErrors(t *testing.T) {
47+
meta := constructFakeMetadata([]PageIndexRanges{{
48+
ColIndexOffset: 0, ColIndexLen: 1,
49+
OffsetIndexOffset: 1, OffsetIndexLen: 1,
50+
}})
51+
reader := &PageIndexReader{
52+
Input: bytes.NewReader([]byte{0xff, 0xff}),
53+
FileMetadata: meta,
54+
Props: parquet.NewReaderProperties(nil),
55+
}
56+
rgReader, err := reader.RowGroup(0)
57+
require.NoError(t, err)
58+
require.NotNil(t, rgReader)
59+
60+
_, err = rgReader.GetColumnIndex(0)
61+
assert.ErrorIs(t, err, arrow.ErrInvalid)
62+
assertCausePreserved(t, err)
63+
64+
_, err = rgReader.GetOffsetIndex(0)
65+
assert.ErrorIs(t, err, arrow.ErrInvalid)
66+
assertCausePreserved(t, err)
67+
}
68+
69+
func assertCausePreserved(t *testing.T, err error) {
70+
t.Helper()
71+
unwrapper, ok := err.(interface{ Unwrap() []error })
72+
require.True(t, ok)
73+
require.Len(t, unwrapper.Unwrap(), 2)
74+
}
75+
76+
func TestDeserializeColumnIndexRepanicsRuntimeErrors(t *testing.T) {
77+
var serialized bytes.Buffer
78+
_, err := thrift.NewThriftSerializer().Serialize(&format.ColumnIndex{}, &serialized, nil)
79+
require.NoError(t, err)
80+
81+
assert.Panics(t, func() {
82+
deserializeColumnIndex(nil, serialized.Bytes(), nil)
83+
})
84+
}

0 commit comments

Comments
 (0)