From ac2a71c1a5d7f49f67d0dc2dd52fdc81e7930d36 Mon Sep 17 00:00:00 2001 From: Will Roever Date: Tue, 1 Sep 2026 18:58:27 -1000 Subject: [PATCH 1/3] fix(manifest): normalize file_format casing when reading entries Route the decoded file_format value through FileFormatFromString in ReadEntry so spec-conformant lowercase spellings match the FileFormat constants, and reject unknown formats there with the offending file path. Signed-off-by: Will Roever --- manifest.go | 7 ++++++ manifest_test.go | 57 ++++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 64 insertions(+) diff --git a/manifest.go b/manifest.go index bb0a510d3..5243ca8a6 100644 --- a/manifest.go +++ b/manifest.go @@ -917,6 +917,13 @@ func (c *ManifestReader) ReadEntry() (ManifestEntry, error) { if c.isFallback { tmp = tmp.(*fallbackManifestEntry).toEntry() } + if df, ok := tmp.DataFile().(*dataFile); ok { + format, err := FileFormatFromString(string(df.Format)) + if err != nil { + return nil, fmt.Errorf("manifest entry for %q has invalid file format: %w", df.Path, err) + } + df.Format = format + } switch tmp.Status() { case EntryStatusEXISTING, EntryStatusADDED, EntryStatusDELETED: default: diff --git a/manifest_test.go b/manifest_test.go index 1db955fa4..c37e8dfe9 100644 --- a/manifest_test.go +++ b/manifest_test.go @@ -2426,6 +2426,63 @@ 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) + entrySchema, err := internal.NewManifestEntrySchema(partitionSchema, 2) + m.Require().NoError(err) + + 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 avro", written: "avro", expected: AvroFile}, + {name: "unknown format", written: "csv", errorContains: "unknown file format: csv"}, + } + + for _, tt := range tests { + m.Run(tt.name, func() { + entry := *manifestEntryV2Records[0] + file := cloneDataFileAvroFields(entry.Data.(*dataFile)) + file.Format = tt.written + entry.Data = file + + mw := ManifestWriter{version: 2, spec: partitionSpec, schema: testSchema, content: ManifestContentData} + 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) + m.Require().NoError(writer.Encode(&entry)) + m.Require().NoError(writer.Close()) + + manifest := &manifestFile{version: 2, SpecID: 1, Content: ManifestContentData} + 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(), From 9a12014ca8884bdadcc3f02b68d6cbbac8e0defa Mon Sep 17 00:00:00 2001 From: Will Roever Date: Wed, 2 Sep 2026 08:38:26 -1000 Subject: [PATCH 2/3] fix(manifest): normalize file_format in the data file codec Signed-off-by: Will Roever --- data_file_codec.go | 3 +++ data_file_codec_test.go | 40 ++++++++++++++++++++++++++++++++++++++++ manifest.go | 20 ++++++++++++++++---- manifest_test.go | 40 ++++++++++++++++++++++++++++++++++------ 4 files changed, 93 insertions(+), 10 deletions(-) diff --git a/data_file_codec.go b/data_file_codec.go index c7c90d47d..8ad4829ec 100644 --- a/data_file_codec.go +++ b/data_file_codec.go @@ -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 diff --git a/data_file_codec_test.go b/data_file_codec_test.go index 1c7606ca0..3b4b0a064 100644 --- a/data_file_codec_test.go +++ b/data_file_codec_test.go @@ -93,6 +93,46 @@ 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"}, + } + + 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 diff --git a/manifest.go b/manifest.go index 5243ca8a6..ec957c801 100644 --- a/manifest.go +++ b/manifest.go @@ -918,11 +918,9 @@ func (c *ManifestReader) ReadEntry() (ManifestEntry, error) { tmp = tmp.(*fallbackManifestEntry).toEntry() } if df, ok := tmp.DataFile().(*dataFile); ok { - format, err := FileFormatFromString(string(df.Format)) - if err != nil { - return nil, fmt.Errorf("manifest entry for %q has invalid file format: %w", df.Path, err) + if err := df.normalizeFormat(); err != nil { + return nil, err } - df.Format = format } switch tmp.Status() { case EntryStatusEXISTING, EntryStatusADDED, EntryStatusDELETED: @@ -2778,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 has invalid file format: %w", d.FilePath(), err) + } + 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 } diff --git a/manifest_test.go b/manifest_test.go index c37e8dfe9..8365fd336 100644 --- a/manifest_test.go +++ b/manifest_test.go @@ -2432,11 +2432,12 @@ func (m *ManifestTestSuite) TestManifestReaderNormalizesFileFormat() { PartitionField{FieldID: 1001, SourceIDs: []int{2}, Name: "tpep_pickup_datetime", Transform: IdentityTransform{}}) partitionSchema, err := partitionTypeToAvroSchema(partitionSpec.PartitionType(testSchema)) m.Require().NoError(err) - entrySchema, err := internal.NewManifestEntrySchema(partitionSchema, 2) - m.Require().NoError(err) tests := []struct { name string + version int + content ManifestContent + entryContent ManifestEntryContent written FileFormat expected FileFormat errorContains string @@ -2445,17 +2446,40 @@ func (m *ManifestTestSuite) TestManifestReaderNormalizesFileFormat() { {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() { - entry := *manifestEntryV2Records[0] + 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: 2, spec: partitionSpec, schema: testSchema, content: ManifestContentData} + mw := ManifestWriter{version: version, spec: partitionSpec, schema: testSchema, content: content} metadata, err := mw.meta() m.Require().NoError(err) @@ -2463,10 +2487,14 @@ func (m *ManifestTestSuite) TestManifestReaderNormalizesFileFormat() { writer, err := ocf.NewWriter(&buf, entrySchema, ocf.WithSchema(entrySchema.String()), ocf.WithMetadata(metadata)) m.Require().NoError(err) - m.Require().NoError(writer.Encode(&entry)) + 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: 2, SpecID: 1, Content: ManifestContentData} + manifest := &manifestFile{version: version, SpecID: 1, Content: content} reader, err := NewManifestReader(manifest, bytes.NewReader(buf.Bytes())) m.Require().NoError(err) defer reader.Close() From a2867e35660ea6728b2cbcd2da62da225594f077 Mon Sep 17 00:00:00 2001 From: Will Roever Date: Mon, 7 Sep 2026 21:56:49 -0700 Subject: [PATCH 3/3] fix(manifest): drop redundant wrapping in the file format error Signed-off-by: Will Roever --- data_file_codec_test.go | 1 + manifest.go | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/data_file_codec_test.go b/data_file_codec_test.go index 3b4b0a064..8d6398930 100644 --- a/data_file_codec_test.go +++ b/data_file_codec_test.go @@ -108,6 +108,7 @@ func TestUnmarshalAvroDataFileEntryNormalizesFileFormat(t *testing.T) { {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: "}, } for _, tt := range tests { diff --git a/manifest.go b/manifest.go index ec957c801..84bb93d9b 100644 --- a/manifest.go +++ b/manifest.go @@ -2783,7 +2783,7 @@ func (d *dataFile) setFieldIDToDecimalScaleMap(m map[int]int) { func (d *dataFile) normalizeFormat() error { format, err := FileFormatFromString(string(d.Format)) if err != nil { - return fmt.Errorf("data file %q has invalid file format: %w", d.FilePath(), err) + return fmt.Errorf("data file %q: %w", d.FilePath(), err) } d.Format = format