From a22cd31146cd57237ac753bb748bfb0a239120e0 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Thu, 6 Aug 2026 19:17:36 +0200 Subject: [PATCH 1/3] fix(parquet/pqarrow): validate reader indexes --- parquet/pqarrow/file_reader.go | 11 ++++++++++ parquet/pqarrow/file_reader_test.go | 34 +++++++++++++++++++++++++++++ 2 files changed, 45 insertions(+) diff --git a/parquet/pqarrow/file_reader.go b/parquet/pqarrow/file_reader.go index 96af9f774..bc1ea25f7 100644 --- a/parquet/pqarrow/file_reader.go +++ b/parquet/pqarrow/file_reader.go @@ -212,6 +212,13 @@ func (fr *FileReader) allRowGroupFactory() itrFactory { // // IncludedLeaves and RowGroups are used to specify precisely which leaf indexes and row groups to read a subset of. func (fr *FileReader) GetFieldReader(ctx context.Context, i int, includedLeaves map[int]bool, rowGroups []int) (*ColumnReader, error) { + if i < 0 || i >= len(fr.Manifest.Fields) { + return nil, fmt.Errorf("invalid field index chosen %d, there are only %d fields", i, len(fr.Manifest.Fields)) + } + if err := fr.checkRowGroups(rowGroups); err != nil { + return nil, err + } + ctx = context.WithValue(ctx, rdrCtxKey{}, readerCtx{ rdr: fr.rdr, mem: fr.mem, @@ -280,6 +287,10 @@ func (fr *FileReader) RowGroup(idx int) RowGroupReader { // ReadColumn reads data to create a chunked array only from the requested row groups. func (fr *FileReader) ReadColumn(rowGroups []int, rdr *ColumnReader) (*arrow.Chunked, error) { + if err := fr.checkRowGroups(rowGroups); err != nil { + return nil, err + } + recs := int64(0) for _, rg := range rowGroups { recs += fr.rdr.MetaData().RowGroups[rg].GetNumRows() diff --git a/parquet/pqarrow/file_reader_test.go b/parquet/pqarrow/file_reader_test.go index 45e0a4f37..3f95ccd05 100644 --- a/parquet/pqarrow/file_reader_test.go +++ b/parquet/pqarrow/file_reader_test.go @@ -606,6 +606,40 @@ func TestFileReaderColumnChunkBoundsErrors(t *testing.T) { } } +func TestFileReaderIndexValidation(t *testing.T) { + schema := arrow.NewSchema([]arrow.Field{{Name: "value", Type: arrow.PrimitiveTypes.Int32}}, nil) + record, _, err := array.RecordFromJSON(memory.DefaultAllocator, schema, + strings.NewReader(`[{"value": 1}]`)) + require.NoError(t, err) + defer record.Release() + + var buf bytes.Buffer + writer, err := pqarrow.NewFileWriter(schema, &buf, nil, pqarrow.DefaultWriterProps()) + require.NoError(t, err) + require.NoError(t, writer.Write(record)) + require.NoError(t, writer.Close()) + + fileReader, err := file.NewParquetReader(bytes.NewReader(buf.Bytes())) + require.NoError(t, err) + defer fileReader.Close() + + arrowReader, err := pqarrow.NewFileReader(fileReader, pqarrow.ArrowReadProperties{}, memory.DefaultAllocator) + require.NoError(t, err) + + _, err = arrowReader.GetFieldReader(context.Background(), -1, nil, []int{0}) + require.Error(t, err) + _, err = arrowReader.GetFieldReader(context.Background(), 1, nil, []int{0}) + require.Error(t, err) + _, err = arrowReader.GetFieldReader(context.Background(), 0, nil, []int{1}) + require.Error(t, err) + + columnReader, err := arrowReader.GetColumn(context.Background(), 0) + require.NoError(t, err) + defer columnReader.Release() + _, err = arrowReader.ReadColumn([]int{1}, columnReader) + require.Error(t, err) +} + func TestReadParquetFile(t *testing.T) { dir := os.Getenv("PARQUET_TEST_BAD_DATA") if dir == "" { From ef02ebddb8de507fa8a902401e2333c90095dddc Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Thu, 6 Aug 2026 21:34:57 +0200 Subject: [PATCH 2/3] fix(parquet/pqarrow): type reader index errors --- parquet/pqarrow/file_reader.go | 6 +++--- parquet/pqarrow/file_reader_test.go | 17 +++++++++++++---- 2 files changed, 16 insertions(+), 7 deletions(-) diff --git a/parquet/pqarrow/file_reader.go b/parquet/pqarrow/file_reader.go index bc1ea25f7..f213ce2cb 100644 --- a/parquet/pqarrow/file_reader.go +++ b/parquet/pqarrow/file_reader.go @@ -213,7 +213,7 @@ func (fr *FileReader) allRowGroupFactory() itrFactory { // IncludedLeaves and RowGroups are used to specify precisely which leaf indexes and row groups to read a subset of. func (fr *FileReader) GetFieldReader(ctx context.Context, i int, includedLeaves map[int]bool, rowGroups []int) (*ColumnReader, error) { if i < 0 || i >= len(fr.Manifest.Fields) { - return nil, fmt.Errorf("invalid field index chosen %d, there are only %d fields", i, len(fr.Manifest.Fields)) + return nil, fmt.Errorf("%w: invalid field index chosen %d, there are only %d fields", arrow.ErrIndex, i, len(fr.Manifest.Fields)) } if err := fr.checkRowGroups(rowGroups); err != nil { return nil, err @@ -316,7 +316,7 @@ func (fr *FileReader) ReadTable(ctx context.Context) (arrow.Table, error) { func (fr *FileReader) checkCols(indices []int) (err error) { for _, col := range indices { if col < 0 || col >= fr.rdr.MetaData().Schema.NumColumns() { - err = fmt.Errorf("invalid column index specified %d out of %d", col, fr.rdr.MetaData().Schema.NumColumns()) + err = fmt.Errorf("%w: invalid column index specified %d out of %d", arrow.ErrIndex, col, fr.rdr.MetaData().Schema.NumColumns()) break } } @@ -326,7 +326,7 @@ func (fr *FileReader) checkCols(indices []int) (err error) { func (fr *FileReader) checkRowGroups(indices []int) (err error) { for _, rg := range indices { if rg < 0 || rg >= fr.rdr.NumRowGroups() { - err = fmt.Errorf("invalid row group specified: %d, file only has %d row groups", rg, fr.rdr.NumRowGroups()) + err = fmt.Errorf("%w: invalid row group specified: %d, file only has %d row groups", arrow.ErrIndex, rg, fr.rdr.NumRowGroups()) break } } diff --git a/parquet/pqarrow/file_reader_test.go b/parquet/pqarrow/file_reader_test.go index 3f95ccd05..08199e8ec 100644 --- a/parquet/pqarrow/file_reader_test.go +++ b/parquet/pqarrow/file_reader_test.go @@ -627,17 +627,26 @@ func TestFileReaderIndexValidation(t *testing.T) { require.NoError(t, err) _, err = arrowReader.GetFieldReader(context.Background(), -1, nil, []int{0}) - require.Error(t, err) + require.ErrorIs(t, err, arrow.ErrIndex) _, err = arrowReader.GetFieldReader(context.Background(), 1, nil, []int{0}) - require.Error(t, err) + require.ErrorIs(t, err, arrow.ErrIndex) _, err = arrowReader.GetFieldReader(context.Background(), 0, nil, []int{1}) - require.Error(t, err) + require.ErrorIs(t, err, arrow.ErrIndex) + _, err = arrowReader.GetFieldReader(context.Background(), 0, nil, []int{-1}) + require.ErrorIs(t, err, arrow.ErrIndex) + + fieldReader, err := arrowReader.GetFieldReader(context.Background(), 0, map[int]bool{0: true}, []int{0}) + require.NoError(t, err) + fieldReader.Release() columnReader, err := arrowReader.GetColumn(context.Background(), 0) require.NoError(t, err) defer columnReader.Release() _, err = arrowReader.ReadColumn([]int{1}, columnReader) - require.Error(t, err) + require.ErrorIs(t, err, arrow.ErrIndex) + chunked, err := arrowReader.ReadColumn([]int{0}, columnReader) + require.NoError(t, err) + chunked.Release() } func TestReadParquetFile(t *testing.T) { From 10ddd13aa2d7ab73ed438d9e72d23618b3fe0787 Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Tue, 11 Aug 2026 21:40:39 +0200 Subject: [PATCH 3/3] fix(parquet/pqarrow): type column reader index errors --- parquet/pqarrow/file_reader.go | 2 +- parquet/pqarrow/file_reader_test.go | 2 ++ 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/parquet/pqarrow/file_reader.go b/parquet/pqarrow/file_reader.go index c3c1a0d85..ae659d220 100644 --- a/parquet/pqarrow/file_reader.go +++ b/parquet/pqarrow/file_reader.go @@ -463,7 +463,7 @@ func (fr *FileReader) ReadRowGroups(ctx context.Context, indices, rowGroups []in func (fr *FileReader) getColumnReader(ctx context.Context, i int, colFactory itrFactory) (*ColumnReader, error) { if i < 0 || i >= len(fr.Manifest.Fields) { - return nil, fmt.Errorf("invalid column index chosen %d, there are only %d columns", i, len(fr.Manifest.Fields)) + return nil, fmt.Errorf("%w: invalid column index chosen %d, there are only %d columns", arrow.ErrIndex, i, len(fr.Manifest.Fields)) } ctx = context.WithValue(ctx, rdrCtxKey{}, readerCtx{ diff --git a/parquet/pqarrow/file_reader_test.go b/parquet/pqarrow/file_reader_test.go index 08199e8ec..16c0c9543 100644 --- a/parquet/pqarrow/file_reader_test.go +++ b/parquet/pqarrow/file_reader_test.go @@ -639,6 +639,8 @@ func TestFileReaderIndexValidation(t *testing.T) { require.NoError(t, err) fieldReader.Release() + _, err = arrowReader.GetColumn(context.Background(), -1) + require.ErrorIs(t, err, arrow.ErrIndex) columnReader, err := arrowReader.GetColumn(context.Background(), 0) require.NoError(t, err) defer columnReader.Release()