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
3 changes: 3 additions & 0 deletions data_file_codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,9 @@ func unmarshalAvroDataFileEntry(data []byte, spec PartitionSpec, schema *Schema,
if _, err := s.Decode(data, entry); err != nil {
return nil, fmt.Errorf("iceberg: unmarshalAvroDataFileEntry: %w", err)
}
if err := df.normalizeFormat(); err != nil {
return nil, fmt.Errorf("iceberg: unmarshalAvroDataFileEntry: %w", err)
}
df.specID = int32(spec.ID())
df.fieldNameToID = maps.nameToID
df.fieldIDToLogicalType = maps.idToType
Expand Down
41 changes: 41 additions & 0 deletions data_file_codec_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,47 @@ func TestDataFileCodecWithDroppedPartitionSource(t *testing.T) {
require.Equal(t, map[int]any{1000: nil, 1001: int32(3)}, decoded.Partition())
}

func TestUnmarshalAvroDataFileEntryNormalizesFileFormat(t *testing.T) {
spec := NewPartitionSpec()
schema := NewSchema(0)

tests := []struct {
name string
written FileFormat
expected FileFormat
errorContains string
}{
{name: "spec lowercase", written: "parquet", expected: ParquetFile},
{name: "mixed case", written: "Parquet", expected: ParquetFile},
{name: "uppercase", written: ParquetFile, expected: ParquetFile},
{name: "lowercase orc", written: "orc", expected: OrcFile},
{name: "unknown format", written: "csv", errorContains: "unknown file format: csv"},
{name: "empty format", written: "", errorContains: "unknown file format: "},
}
Comment thread
wroever marked this conversation as resolved.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit — empty-format row asserts a prefix common to every unknown-format error

errorContains "unknown file format: " is a substring of "unknown file format: csv" and of every other unknown-format message, so the row cannot distinguish the empty value from any other rejected spelling. It still verifies that an empty file_format is not silently accepted, which is the point of the row, so this is purely about assertion precision. Same applies to manifest_test.go:2456.


for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
builder, err := NewDataFileBuilder(spec, EntryContentData,
"s3://bucket/table/data.parquet", ParquetFile, nil, nil, nil, 1, 1024)
require.NoError(t, err)
source := builder.Build().(*dataFile)
source.Format = tt.written

encoded, err := source.MarshalAvroEntry(spec, schema, 2)
require.NoError(t, err)

decoded, err := unmarshalAvroDataFileEntry(encoded, spec, schema, 2)
if tt.errorContains != "" {
require.ErrorContains(t, err, tt.errorContains)

return
}
require.NoError(t, err)
require.Equal(t, tt.expected, decoded.FileFormat())
})
}
}

// TestManifestEntrySchemaForDistinguishesSameSpecIDDifferentPartitionTypes
// guards against the cache-key collision that would occur if the cache
// were keyed only by spec.ID(). Two tables can both use
Expand Down
19 changes: 19 additions & 0 deletions manifest.go
Original file line number Diff line number Diff line change
Expand Up @@ -917,6 +917,11 @@ func (c *ManifestReader) ReadEntry() (ManifestEntry, error) {
if c.isFallback {
tmp = tmp.(*fallbackManifestEntry).toEntry()
}
if df, ok := tmp.DataFile().(*dataFile); ok {
if err := df.normalizeFormat(); err != nil {
return nil, err
}
}
switch tmp.Status() {
case EntryStatusEXISTING, EntryStatusADDED, EntryStatusDELETED:
default:
Expand Down Expand Up @@ -2771,6 +2776,20 @@ func (d *dataFile) setFieldIDToDecimalScaleMap(m map[int]int) {
d.fieldIDToDecimalScale = m
}

// normalizeFormat sets d.Format to the FileFormat constant matching its
// decoded spelling, or returns an error if it names no known format.
// Every Avro decode path must call it before the value is compared or
// exposed.
func (d *dataFile) normalizeFormat() error {
format, err := FileFormatFromString(string(d.Format))
if err != nil {
return fmt.Errorf("data file %q: %w", d.FilePath(), err)
}
Comment thread
wroever marked this conversation as resolved.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit — Corrected error message is pinned by no assertion

The delta's whole payload is the wording of this format string, but no test asserts it. Reintroducing the redundant 'has invalid file format:' prefix keeps the suite green, so the improvement can silently regress. Tightening one row per table to assert the full 'data file "...": unknown file format: csv' would lock it in. Cosmetic-only, so not worth a round-trip on its own.

d.Format = format

return nil
}

func (d *dataFile) ContentType() ManifestEntryContent { return d.Content }
func (d *dataFile) FilePath() string { return d.Path }
func (d *dataFile) FileFormat() FileFormat { return d.Format }
Expand Down
85 changes: 85 additions & 0 deletions manifest_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2426,6 +2426,91 @@ func (m *ManifestTestSuite) TestManifestReaderRejectsInvalidEntries() {
}
}

func (m *ManifestTestSuite) TestManifestReaderNormalizesFileFormat() {
partitionSpec := NewPartitionSpecID(1,
PartitionField{FieldID: 1000, SourceIDs: []int{1}, Name: "VendorID", Transform: IdentityTransform{}},
PartitionField{FieldID: 1001, SourceIDs: []int{2}, Name: "tpep_pickup_datetime", Transform: IdentityTransform{}})
partitionSchema, err := partitionTypeToAvroSchema(partitionSpec.PartitionType(testSchema))
m.Require().NoError(err)

tests := []struct {
name string
version int
content ManifestContent
entryContent ManifestEntryContent
written FileFormat
expected FileFormat
errorContains string
}{
{name: "spec lowercase", written: "parquet", expected: ParquetFile},
{name: "mixed case", written: "Parquet", expected: ParquetFile},
{name: "uppercase", written: ParquetFile, expected: ParquetFile},
{name: "lowercase avro", written: "avro", expected: AvroFile},
{name: "lowercase orc", written: "orc", expected: OrcFile},
{
name: "lowercase puffin in delete manifest", content: ManifestContentDeletes,
entryContent: EntryContentPosDeletes, written: "puffin", expected: PuffinFile,
},
{name: "v1 fallback entry", version: 1, written: "parquet", expected: ParquetFile},
{name: "unknown format", written: "csv", errorContains: "unknown file format: csv"},
{name: "empty format", written: "", errorContains: "unknown file format: "},
}

for _, tt := range tests {
m.Run(tt.name, func() {
version := tt.version
if version == 0 {
version = 2
}
content := tt.content
if content == 0 {
content = ManifestContentData
}
entrySchema, err := internal.NewManifestEntrySchema(partitionSchema, version)
m.Require().NoError(err)

records := manifestEntryV2Records
if version == 1 {
records = manifestEntryV1Records
}
entry := *records[0]
file := cloneDataFileAvroFields(entry.Data.(*dataFile))
file.Format = tt.written
file.Content = tt.entryContent
entry.Data = file

mw := ManifestWriter{version: version, spec: partitionSpec, schema: testSchema, content: content}
metadata, err := mw.meta()
m.Require().NoError(err)

var buf bytes.Buffer
writer, err := ocf.NewWriter(&buf, entrySchema,
ocf.WithSchema(entrySchema.String()), ocf.WithMetadata(metadata))
m.Require().NoError(err)
var toEncode any = &entry
if version == 1 {
toEncode = &fallbackManifestEntry{manifestEntry: entry}
}
m.Require().NoError(writer.Encode(toEncode))
m.Require().NoError(writer.Close())

manifest := &manifestFile{version: version, SpecID: 1, Content: content}
reader, err := NewManifestReader(manifest, bytes.NewReader(buf.Bytes()))
m.Require().NoError(err)
defer reader.Close()

got, err := reader.ReadEntry()
if tt.errorContains != "" {
m.ErrorContains(err, tt.errorContains)

return
}
m.Require().NoError(err)
m.Equal(tt.expected, got.DataFile().FileFormat())
})
}
}

func (m *ManifestTestSuite) TestManifestEntryBuilder() {
dataFileBuilder, err := NewDataFileBuilder(
NewPartitionSpec(),
Expand Down
Loading