Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions codec/file_scan_task.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,9 @@ func EncodeFileScanTask(task table.FileScanTask, spec iceberg.PartitionSpec, sch
if err := checkDataFileSpecID(task.File, spec); err != nil {
return nil, fmt.Errorf("codec: EncodeFileScanTask: %w", err)
}
if err := validateScanRange(task.Start, task.Length, task.File.FileSizeBytes()); err != nil {
return nil, fmt.Errorf("codec: EncodeFileScanTask: %w", err)
}
fileBytes, err := EncodeDataFile(task.File, spec, schema, version)
if err != nil {
return nil, fmt.Errorf("codec: EncodeFileScanTask: file: %w", err)
Expand Down Expand Up @@ -132,6 +135,9 @@ func DecodeFileScanTask(data []byte, spec iceberg.PartitionSpec, schema *iceberg
if err != nil {
return table.FileScanTask{}, fmt.Errorf("codec: DecodeFileScanTask: file: %w", err)
}
if err := validateScanRange(envelope.Start, envelope.Length, file.FileSizeBytes()); err != nil {
return table.FileScanTask{}, fmt.Errorf("codec: DecodeFileScanTask: %w", err)
}
del, err := decodeDataFileSlice(envelope.DeleteFiles, spec, schema, version)
if err != nil {
return table.FileScanTask{}, fmt.Errorf("codec: DecodeFileScanTask: delete files: %w", err)
Expand Down Expand Up @@ -162,6 +168,14 @@ func DecodeFileScanTask(data []byte, spec iceberg.PartitionSpec, schema *iceberg
}, nil
}

func validateScanRange(start, length, fileSize int64) error {
if start > fileSize || length > fileSize-start {
return fmt.Errorf("scan range start=%d length=%d exceeds file size %d", start, length, fileSize)
}

return nil
}

// fileScanTaskShape is a compile-time drift guard for FileScanTask.
// Go only permits struct conversion between two types that have
// identical underlying field sequences (names, types, and order; tags
Expand Down
41 changes: 41 additions & 0 deletions codec/file_scan_task_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,37 @@ package codec

import (
"encoding/binary"
"math"
"testing"

"github.com/apache/iceberg-go"
"github.com/stretchr/testify/require"
)

func TestValidateScanRange(t *testing.T) {
for _, tt := range []struct {
name string
start, length int64
fileSize int64
shouldReturnError bool
}{
{name: "full file", length: 100, fileSize: 100},
{name: "empty range at EOF", start: 100, fileSize: 100},
{name: "start after EOF", start: 101, fileSize: 100, shouldReturnError: true},
{name: "end after EOF", start: 99, length: 2, fileSize: 100, shouldReturnError: true},
{name: "overflowing end", start: math.MaxInt64, length: 1, fileSize: 100, shouldReturnError: true},
} {
t.Run(tt.name, func(t *testing.T) {
err := validateScanRange(tt.start, tt.length, tt.fileSize)
if tt.shouldReturnError {
require.Error(t, err)
} else {
require.NoError(t, err)
}
})
}
}

// TestDecodeFileScanTaskInnerErrorCarriesMarker covers an inner decode path:
// the outer Avro envelope decodes fine, but the framed primary-file blob is not
// a valid DataFile encoding. The wrapped error must still carry the
Expand Down Expand Up @@ -110,3 +135,19 @@ func TestDecodeFileScanTaskRejectsNegativeScanRanges(t *testing.T) {
require.Contains(t, err.Error(), "length must be non-negative")
})
}

func TestDecodeFileScanTaskRejectsRangeBeyondFileSize(t *testing.T) {
spec := *iceberg.UnpartitionedSpec
builder, err := iceberg.NewDataFileBuilder(spec, iceberg.EntryContentData,
"data.parquet", iceberg.ParquetFile, nil, nil, nil, 1, 100)
require.NoError(t, err)
file, err := EncodeDataFile(builder.Build(), spec, nil, 2)
require.NoError(t, err)
envelope, err := fileScanTaskSchema.Encode(&fileScanTaskEnvelope{
File: file, Start: 90, Length: 11,
})
require.NoError(t, err)

_, err = DecodeFileScanTask(envelope, spec, nil, 2)
require.ErrorContains(t, err, "scan range start=90 length=11 exceeds file size 100")
}
32 changes: 32 additions & 0 deletions codec/file_scan_task_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
package codec_test

import (
"math"
"strconv"
"testing"

Expand Down Expand Up @@ -221,6 +222,37 @@ func TestEncodeFileScanTaskRejectsNegativeScanRanges(t *testing.T) {
})
}

func TestEncodeFileScanTaskValidatesRangeAgainstFileSize(t *testing.T) {
spec, schema, base := fullyPopulatedFileScanTask(t, 2)
fileSize := base.File.FileSizeBytes()
tests := []struct {
name string
start int64
length int64
shouldError bool
}{
{"full file", 0, fileSize, false},
{"suffix ending at EOF", fileSize - 1, 1, false},
{"zero length at EOF", fileSize, 0, false},
{"start after EOF", fileSize + 1, 0, true},
{"range ends after EOF", fileSize - 1, 2, true},
{"overflowing range", math.MaxInt64, 1, true},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
task := base
task.Start, task.Length = tt.start, tt.length
_, err := codec.EncodeFileScanTask(task, spec, schema, 2)
if tt.shouldError {
require.ErrorContains(t, err, "exceeds file size")
} else {
require.NoError(t, err)
}
})
}
}

func fullyPopulatedFileScanTask(t *testing.T, version int) (iceberg.PartitionSpec, *iceberg.Schema, table.FileScanTask) {
t.Helper()
schema := iceberg.NewSchema(123,
Expand Down
Loading