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

laskoviymishka pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-go.git


The following commit(s) were added to refs/heads/main by this push:
     new 327ec1e39 fix(table): validate position-delete columns before access 
(#1583)
327ec1e39 is described below

commit 327ec1e3989c07a57daf5af65bf02ef13bbce407
Author: Minh Vu <[email protected]>
AuthorDate: Fri Jul 31 16:22:48 2026 +0200

    fix(table): validate position-delete columns before access (#1583)
    
    ## What changed
    
    Validate that a position-delete table contains exactly one `file_path`
    column and exactly one `pos` column before indexing either field.
    
    ## Why
    
    The scanner indexed the first result from `FieldIndices` without
    checking its length. Missing required columns therefore caused an
    index-out-of-range panic, while duplicate columns were ambiguous.
    
    Regression coverage includes valid schemas, the optional `row` column,
    each missing-column case, neither column, and duplicate `pos` columns.
    
    ## Testing
    
    - `go test ./table`
    - `go vet ./table`
    
    ---------
    
    Signed-off-by: Minh Vu <[email protected]>
---
 table/arrow_scanner.go                           | 33 ++++++++++++-
 table/arrow_scanner_posdelete_regression_test.go | 63 ++++++++++++++++++++++++
 2 files changed, 94 insertions(+), 2 deletions(-)

diff --git a/table/arrow_scanner.go b/table/arrow_scanner.go
index ffcf987ed..027018f6a 100644
--- a/table/arrow_scanner.go
+++ b/table/arrow_scanner.go
@@ -449,8 +449,12 @@ func readDeletes(ctx context.Context, fs iceio.IO, 
dataFile iceberg.DataFile) (_
        }
        defer tbl.Release()
 
-       filePathCol := 
tbl.Column(tbl.Schema().FieldIndices("file_path")[0]).Data()
-       posCol := tbl.Column(tbl.Schema().FieldIndices("pos")[0]).Data()
+       filePathIndex, posIndex, err := 
positionDeleteColumnIndices(tbl.Schema())
+       if err != nil {
+               return nil, err
+       }
+       filePathCol := tbl.Column(filePathIndex).Data()
+       posCol := tbl.Column(posIndex).Data()
        if posCol.NullN() > 0 {
                return nil, fmt.Errorf("%w: null pos in position delete file", 
iceberg.ErrInvalidSchema)
        }
@@ -458,6 +462,31 @@ func readDeletes(ctx context.Context, fs iceio.IO, 
dataFile iceberg.DataFile) (_
        return groupPosDeletesByFilePath(ctx, filePathCol, posCol)
 }
 
+func positionDeleteColumnIndices(schema *arrow.Schema) (int, int, error) {
+       // Position-delete columns are resolved by their spec-defined names 
because
+       // Arrow schemas read from external files do not always retain Iceberg 
IDs.
+       requiredColumn := func(name string) (int, error) {
+               indices := schema.FieldIndices(name)
+               if len(indices) != 1 {
+                       return 0, fmt.Errorf("%w: position delete file must 
contain exactly one %q column, found %d",
+                               iceberg.ErrInvalidSchema, name, len(indices))
+               }
+
+               return indices[0], nil
+       }
+
+       filePathIndex, err := requiredColumn("file_path")
+       if err != nil {
+               return 0, 0, err
+       }
+       posIndex, err := requiredColumn("pos")
+       if err != nil {
+               return 0, 0, err
+       }
+
+       return filePathIndex, posIndex, nil
+}
+
 type set[T comparable] map[T]struct{}
 
 // combinePositionalDeletes builds the surviving-row index list for a single 
record
diff --git a/table/arrow_scanner_posdelete_regression_test.go 
b/table/arrow_scanner_posdelete_regression_test.go
index 52d6634eb..be59f0e18 100644
--- a/table/arrow_scanner_posdelete_regression_test.go
+++ b/table/arrow_scanner_posdelete_regression_test.go
@@ -32,6 +32,69 @@ import (
        "github.com/stretchr/testify/require"
 )
 
+func TestPositionDeleteColumnIndices(t *testing.T) {
+       t.Parallel()
+
+       tests := []struct {
+               name              string
+               fields            []arrow.Field
+               wantFilePathIndex int
+               wantPosIndex      int
+               wantErr           string
+       }{
+               {name: "valid", fields: []arrow.Field{{Name: "file_path", Type: 
arrow.BinaryTypes.String}, {Name: "pos", Type: arrow.PrimitiveTypes.Int64}}, 
wantPosIndex: 1},
+               {name: "valid with row", fields: []arrow.Field{{Name: 
"file_path", Type: arrow.BinaryTypes.String}, {Name: "pos", Type: 
arrow.PrimitiveTypes.Int64}, {Name: "row", Type: arrow.BinaryTypes.String}}, 
wantPosIndex: 1},
+               {name: "valid reversed", fields: []arrow.Field{{Name: "pos", 
Type: arrow.PrimitiveTypes.Int64}, {Name: "file_path", Type: 
arrow.BinaryTypes.String}}, wantFilePathIndex: 1},
+               {name: "missing file_path", fields: []arrow.Field{{Name: "pos", 
Type: arrow.PrimitiveTypes.Int64}}, wantErr: `exactly one "file_path" column, 
found 0`},
+               {name: "missing pos", fields: []arrow.Field{{Name: "file_path", 
Type: arrow.BinaryTypes.String}}, wantErr: `exactly one "pos" column, found 0`},
+               {name: "missing both", fields: nil, wantErr: `exactly one 
"file_path" column, found 0`},
+               {name: "duplicate pos", fields: []arrow.Field{{Name: 
"file_path", Type: arrow.BinaryTypes.String}, {Name: "pos", Type: 
arrow.PrimitiveTypes.Int64}, {Name: "pos", Type: arrow.PrimitiveTypes.Int64}}, 
wantErr: `exactly one "pos" column, found 2`},
+               {name: "duplicate file_path", fields: []arrow.Field{{Name: 
"file_path", Type: arrow.BinaryTypes.String}, {Name: "file_path", Type: 
arrow.BinaryTypes.String}, {Name: "pos", Type: arrow.PrimitiveTypes.Int64}}, 
wantErr: `exactly one "file_path" column, found 2`},
+       }
+
+       for _, test := range tests {
+               t.Run(test.name, func(t *testing.T) {
+                       filePathIndex, posIndex, err := 
positionDeleteColumnIndices(arrow.NewSchema(test.fields, nil))
+                       if test.wantErr != "" {
+                               require.ErrorIs(t, err, 
iceberg.ErrInvalidSchema)
+                               require.ErrorContains(t, err, test.wantErr)
+
+                               return
+                       }
+
+                       require.NoError(t, err)
+                       assert.Equal(t, test.wantFilePathIndex, filePathIndex)
+                       assert.Equal(t, test.wantPosIndex, posIndex)
+               })
+       }
+}
+
+func TestReadDeletesRejectsMissingFilePath(t *testing.T) {
+       mem := memory.NewCheckedAllocator(memory.DefaultAllocator)
+       ctx := compute.WithAllocator(t.Context(), mem)
+       defer mem.AssertSize(t, 0)
+
+       deleteSchema := arrow.NewSchema([]arrow.Field{{Name: "pos", Type: 
arrow.PrimitiveTypes.Int64}}, nil)
+       deletePath := "mem://bucket/deletes/missing-file-path.parquet"
+       rec := mustLoadRecordBatchFromJSON(deleteSchema, `[{"pos": 1}]`)
+       defer rec.Release()
+       tbl := array.NewTableFromRecords(deleteSchema, []arrow.RecordBatch{rec})
+       defer tbl.Release()
+
+       memFS := iceio.NewMemFS()
+       fw, err := memFS.Create(deletePath)
+       require.NoError(t, err)
+       require.NoError(t, pqarrow.WriteTable(tbl, fw, rec.NumRows(),
+               parquet.NewWriterProperties(parquet.WithStats(true)),
+               pqarrow.DefaultWriterProps()))
+       require.NoError(t, fw.Close())
+
+       deletes, err := readDeletes(ctx, memFS, newPosDeleteFile(t, deletePath, 
1, 128))
+       require.ErrorIs(t, err, iceberg.ErrInvalidSchema)
+       assert.Nil(t, deletes)
+       assert.Contains(t, err.Error(), `exactly one "file_path" column, found 
0`)
+}
+
 func TestGroupPosDeletesByFilePathSupportsStringLayouts(t *testing.T) {
        for _, tc := range []struct {
                name                  string

Reply via email to