laskoviymishka commented on code in PR #1913:
URL: https://github.com/apache/iceberg-go/pull/1913#discussion_r3873464310


##########
table/dv_scan_planning_test.go:
##########
@@ -141,6 +141,29 @@ func TestManifestEntries_DVClassification(t *testing.T) {
        assert.Len(t, entries.equalityDeleteEntries, 1)
 }
 
+func TestManifestEntries_MergeRejectsUnknownContent(t *testing.T) {
+       snapshotID := int64(1)
+       validEntry := iceberg.NewManifestEntry(iceberg.EntryStatusADDED, 
&snapshotID, nil, nil, &mockDataFile{
+               path:        "s3://bucket/data/data-001.parquet",
+               contentType: iceberg.EntryContentData,
+       })
+       invalidEntry := iceberg.NewManifestEntry(iceberg.EntryStatusADDED, 
&snapshotID, nil, nil, &mockDataFile{
+               path:        "s3://bucket/data/unknown.bin",
+               contentType: iceberg.ManifestEntryContent(99),
+       })
+       entryAfterInvalid := iceberg.NewManifestEntry(iceberg.EntryStatusADDED, 
&snapshotID, nil, nil, &mockDataFile{
+               path:        "s3://bucket/data/data-002.parquet",
+               contentType: iceberg.EntryContentData,
+       })
+
+       entries := newManifestEntries()
+       err := entries.merge([]iceberg.ManifestEntry{validEntry, invalidEntry, 
entryAfterInvalid})

Review Comment:
   This case only exercises the `kinds == nil` path, since every valid entry 
here is data, so classification stays on the fast 
`classifiedManifestEntriesForKind` branch and the `kinds != nil` backfill loop 
never runs.
   
   I'd add a mixed case like `[data, posDelete, invalid, data]` so the error 
hits with `kinds != nil` and the backfill path gets covered, asserting 
`dataEntries == [e0]` and `positionalDeleteEntries == [e1]`. While we're here, 
worth asserting the other three buckets stay empty too, so a bad entry bleeding 
into the wrong bucket can't pass silently.



##########
table/dv_scan_planning_test.go:
##########
@@ -141,6 +141,29 @@ func TestManifestEntries_DVClassification(t *testing.T) {
        assert.Len(t, entries.equalityDeleteEntries, 1)
 }
 
+func TestManifestEntries_MergeRejectsUnknownContent(t *testing.T) {
+       snapshotID := int64(1)
+       validEntry := iceberg.NewManifestEntry(iceberg.EntryStatusADDED, 
&snapshotID, nil, nil, &mockDataFile{
+               path:        "s3://bucket/data/data-001.parquet",
+               contentType: iceberg.EntryContentData,
+       })
+       invalidEntry := iceberg.NewManifestEntry(iceberg.EntryStatusADDED, 
&snapshotID, nil, nil, &mockDataFile{
+               path:        "s3://bucket/data/unknown.bin",
+               contentType: iceberg.ManifestEntryContent(99),
+       })
+       entryAfterInvalid := iceberg.NewManifestEntry(iceberg.EntryStatusADDED, 
&snapshotID, nil, nil, &mockDataFile{
+               path:        "s3://bucket/data/data-002.parquet",
+               contentType: iceberg.EntryContentData,
+       })
+
+       entries := newManifestEntries()
+       err := entries.merge([]iceberg.ManifestEntry{validEntry, invalidEntry, 
entryAfterInvalid})
+
+       assert.ErrorIs(t, err, ErrInvalidMetadata)
+       assert.ErrorContains(t, err, "unknown DataFileContent type")
+       assert.Equal(t, []iceberg.ManifestEntry{validEntry}, 
entries.dataEntries)

Review Comment:
   One thing worth calling out: this pins the partial-commit-on-error behavior 
as expected, but the real caller (`collectManifestEntriesWithSchema`) returns 
on the first error and discards these entries, so nothing actually depends on 
the partial result landing.
   
   Not asking to change the behavior, just maybe a comment here or in 
`classifyManifestEntries` noting the partial commit is incidental, so this test 
doesn't get read as a hard requirement to preserve. wdyt?



##########
table/scanner.go:
##########
@@ -147,30 +164,128 @@ func newManifestEntries() *manifestEntries {
        }
 }
 
-func (m *manifestEntries) merge(entries []iceberg.ManifestEntry) error {
-       m.mu.Lock()
-       defer m.mu.Unlock()
-
-       for _, entry := range entries {
-               dataFile := entry.DataFile()
-               switch dataFile.ContentType() {
-               case iceberg.EntryContentData:
-                       m.dataEntries = append(m.dataEntries, entry)
-               case iceberg.EntryContentPosDeletes:
-                       if IsDeletionVector(dataFile) {
-                               m.dvEntries = append(m.dvEntries, entry)
-                       } else {
-                               m.positionalDeleteEntries = 
append(m.positionalDeleteEntries, entry)
+func classifyManifestEntry(entry iceberg.ManifestEntry) (manifestEntryKind, 
error) {
+       dataFile := entry.DataFile()
+       switch dataFile.ContentType() {
+       case iceberg.EntryContentData:
+               return manifestEntryData, nil
+       case iceberg.EntryContentPosDeletes:
+               if IsDeletionVector(dataFile) {
+                       return manifestEntryDV, nil
+               }
+
+               return manifestEntryPositionalDelete, nil
+       case iceberg.EntryContentEqDeletes:
+               return manifestEntryEqualityDelete, nil
+       default:
+               return 0, fmt.Errorf("%w: unknown DataFileContent type (%s): 
%s",
+                       ErrInvalidMetadata, dataFile.ContentType(), entry)
+       }
+}
+
+func newClassifiedManifestEntries(counts [manifestEntryKindCount]int) 
classifiedManifestEntries {
+       return classifiedManifestEntries{
+               dataEntries:             make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryData]),
+               positionalDeleteEntries: make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryPositionalDelete]),
+               equalityDeleteEntries:   make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryEqualityDelete]),
+               dvEntries:               make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryDV]),
+       }
+}
+

Review Comment:
   Minor: the homogeneous path assigns the caller's slice straight into the 
result (`classified.dataEntries = entries`) with no copy, so the returned 
struct aliases the `openManifest` slice that `merge` was handed. It's safe 
today since `merge` only reads it through `append`, but the aliasing is 
intentional and undocumented.
   
   A one-line comment would save the next person who tries to mutate one of 
these buckets in place from a surprise.



##########
table/scanner.go:
##########
@@ -147,30 +164,128 @@ func newManifestEntries() *manifestEntries {
        }
 }
 
-func (m *manifestEntries) merge(entries []iceberg.ManifestEntry) error {
-       m.mu.Lock()
-       defer m.mu.Unlock()
-
-       for _, entry := range entries {
-               dataFile := entry.DataFile()
-               switch dataFile.ContentType() {
-               case iceberg.EntryContentData:
-                       m.dataEntries = append(m.dataEntries, entry)
-               case iceberg.EntryContentPosDeletes:
-                       if IsDeletionVector(dataFile) {
-                               m.dvEntries = append(m.dvEntries, entry)
-                       } else {
-                               m.positionalDeleteEntries = 
append(m.positionalDeleteEntries, entry)
+func classifyManifestEntry(entry iceberg.ManifestEntry) (manifestEntryKind, 
error) {
+       dataFile := entry.DataFile()
+       switch dataFile.ContentType() {
+       case iceberg.EntryContentData:
+               return manifestEntryData, nil
+       case iceberg.EntryContentPosDeletes:
+               if IsDeletionVector(dataFile) {
+                       return manifestEntryDV, nil
+               }
+
+               return manifestEntryPositionalDelete, nil
+       case iceberg.EntryContentEqDeletes:
+               return manifestEntryEqualityDelete, nil
+       default:
+               return 0, fmt.Errorf("%w: unknown DataFileContent type (%s): 
%s",
+                       ErrInvalidMetadata, dataFile.ContentType(), entry)
+       }
+}
+
+func newClassifiedManifestEntries(counts [manifestEntryKindCount]int) 
classifiedManifestEntries {
+       return classifiedManifestEntries{
+               dataEntries:             make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryData]),
+               positionalDeleteEntries: make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryPositionalDelete]),
+               equalityDeleteEntries:   make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryEqualityDelete]),
+               dvEntries:               make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryDV]),
+       }
+}
+
+func classifiedManifestEntriesForKind(kind manifestEntryKind, entries 
[]iceberg.ManifestEntry) classifiedManifestEntries {
+       classified := classifiedManifestEntries{}
+       switch kind {
+       case manifestEntryData:
+               classified.dataEntries = entries
+       case manifestEntryPositionalDelete:
+               classified.positionalDeleteEntries = entries
+       case manifestEntryEqualityDelete:
+               classified.equalityDeleteEntries = entries
+       case manifestEntryDV:
+               classified.dvEntries = entries
+       }
+
+       return classified
+}
+
+func classifyManifestEntries(entries []iceberg.ManifestEntry) 
(classifiedManifestEntries, error) {
+       if len(entries) == 0 {
+               return classifiedManifestEntries{}, nil
+       }
+
+       firstKind, err := classifyManifestEntry(entries[0])
+       if err != nil {
+               return classifiedManifestEntries{}, err
+       }
+
+       var (
+               kinds  []manifestEntryKind
+               counts [manifestEntryKindCount]int
+       )
+       for i := 1; i < len(entries); i++ {
+               kind, err := classifyManifestEntry(entries[i])
+               if err != nil {
+                       if kinds == nil {
+                               return 
classifiedManifestEntriesForKind(firstKind, entries[:i]), err
                        }
-               case iceberg.EntryContentEqDeletes:
-                       m.equalityDeleteEntries = 
append(m.equalityDeleteEntries, entry)
-               default:
-                       return fmt.Errorf("%w: unknown DataFileContent type 
(%s): %s",
-                               ErrInvalidMetadata, dataFile.ContentType(), 
entry)
+
+                       classified := newClassifiedManifestEntries(counts)
+                       for j, validEntry := range entries[:i] {
+                               appendClassifiedManifestEntry(&classified, 
kinds[j], validEntry)
+                       }
+
+                       return classified, err
+               }
+
+               if kinds == nil && kind != firstKind {
+                       kinds = make([]manifestEntryKind, len(entries))
+                       for j := range i {
+                               kinds[j] = firstKind
+                       }
+                       counts[firstKind] = i
+               }
+               if kinds != nil {
+                       kinds[i] = kind
+                       counts[kind]++
                }
        }
 
-       return nil
+       if kinds == nil {
+               return classifiedManifestEntriesForKind(firstKind, entries), nil
+       }
+
+       classified := newClassifiedManifestEntries(counts)
+       for i, entry := range entries {
+               appendClassifiedManifestEntry(&classified, kinds[i], entry)
+       }
+
+       return classified, nil
+}
+

Review Comment:
   Neither this switch nor the one in `classifiedManifestEntriesForKind` has a 
default, so if a fifth `manifestEntryKind` ever gets added and someone forgets 
to wire it here, the entry is silently dropped with no error. `exhaustive` 
isn't enabled, so nothing catches that statically, and the non-data fast paths 
aren't covered by a test either, so a drop would be invisible.
   
   Not blocking, but I'd add a `default: panic(fmt.Sprintf("unhandled 
manifestEntryKind %d", kind))` to both switches so a future miss fails loudly. 
wdyt?



##########
table/scanner.go:
##########
@@ -147,30 +164,128 @@ func newManifestEntries() *manifestEntries {
        }
 }
 
-func (m *manifestEntries) merge(entries []iceberg.ManifestEntry) error {
-       m.mu.Lock()
-       defer m.mu.Unlock()
-
-       for _, entry := range entries {
-               dataFile := entry.DataFile()
-               switch dataFile.ContentType() {
-               case iceberg.EntryContentData:
-                       m.dataEntries = append(m.dataEntries, entry)
-               case iceberg.EntryContentPosDeletes:
-                       if IsDeletionVector(dataFile) {
-                               m.dvEntries = append(m.dvEntries, entry)
-                       } else {
-                               m.positionalDeleteEntries = 
append(m.positionalDeleteEntries, entry)
+func classifyManifestEntry(entry iceberg.ManifestEntry) (manifestEntryKind, 
error) {
+       dataFile := entry.DataFile()
+       switch dataFile.ContentType() {
+       case iceberg.EntryContentData:
+               return manifestEntryData, nil
+       case iceberg.EntryContentPosDeletes:
+               if IsDeletionVector(dataFile) {
+                       return manifestEntryDV, nil
+               }
+
+               return manifestEntryPositionalDelete, nil
+       case iceberg.EntryContentEqDeletes:
+               return manifestEntryEqualityDelete, nil
+       default:
+               return 0, fmt.Errorf("%w: unknown DataFileContent type (%s): 
%s",
+                       ErrInvalidMetadata, dataFile.ContentType(), entry)
+       }
+}
+
+func newClassifiedManifestEntries(counts [manifestEntryKindCount]int) 
classifiedManifestEntries {
+       return classifiedManifestEntries{
+               dataEntries:             make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryData]),
+               positionalDeleteEntries: make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryPositionalDelete]),
+               equalityDeleteEntries:   make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryEqualityDelete]),
+               dvEntries:               make([]iceberg.ManifestEntry, 0, 
counts[manifestEntryDV]),
+       }
+}
+
+func classifiedManifestEntriesForKind(kind manifestEntryKind, entries 
[]iceberg.ManifestEntry) classifiedManifestEntries {
+       classified := classifiedManifestEntries{}
+       switch kind {
+       case manifestEntryData:
+               classified.dataEntries = entries
+       case manifestEntryPositionalDelete:
+               classified.positionalDeleteEntries = entries
+       case manifestEntryEqualityDelete:
+               classified.equalityDeleteEntries = entries
+       case manifestEntryDV:
+               classified.dvEntries = entries
+       }
+
+       return classified
+}
+
+func classifyManifestEntries(entries []iceberg.ManifestEntry) 
(classifiedManifestEntries, error) {
+       if len(entries) == 0 {
+               return classifiedManifestEntries{}, nil
+       }
+
+       firstKind, err := classifyManifestEntry(entries[0])
+       if err != nil {
+               return classifiedManifestEntries{}, err
+       }
+
+       var (
+               kinds  []manifestEntryKind
+               counts [manifestEntryKindCount]int
+       )
+       for i := 1; i < len(entries); i++ {
+               kind, err := classifyManifestEntry(entries[i])
+               if err != nil {
+                       if kinds == nil {
+                               return 
classifiedManifestEntriesForKind(firstKind, entries[:i]), err
                        }
-               case iceberg.EntryContentEqDeletes:
-                       m.equalityDeleteEntries = 
append(m.equalityDeleteEntries, entry)
-               default:
-                       return fmt.Errorf("%w: unknown DataFileContent type 
(%s): %s",
-                               ErrInvalidMetadata, dataFile.ContentType(), 
entry)
+
+                       classified := newClassifiedManifestEntries(counts)
+                       for j, validEntry := range entries[:i] {
+                               appendClassifiedManifestEntry(&classified, 
kinds[j], validEntry)
+                       }
+
+                       return classified, err
+               }
+
+               if kinds == nil && kind != firstKind {
+                       kinds = make([]manifestEntryKind, len(entries))
+                       for j := range i {
+                               kinds[j] = firstKind
+                       }
+                       counts[firstKind] = i
+               }
+               if kinds != nil {
+                       kinds[i] = kind
+                       counts[kind]++
                }
        }
 
-       return nil
+       if kinds == nil {
+               return classifiedManifestEntriesForKind(firstKind, entries), nil
+       }
+
+       classified := newClassifiedManifestEntries(counts)
+       for i, entry := range entries {
+               appendClassifiedManifestEntry(&classified, kinds[i], entry)
+       }
+
+       return classified, nil
+}
+
+func appendClassifiedManifestEntry(classified *classifiedManifestEntries, kind 
manifestEntryKind, entry iceberg.ManifestEntry) {
+       switch kind {
+       case manifestEntryData:
+               classified.dataEntries = append(classified.dataEntries, entry)
+       case manifestEntryPositionalDelete:
+               classified.positionalDeleteEntries = 
append(classified.positionalDeleteEntries, entry)
+       case manifestEntryEqualityDelete:
+               classified.equalityDeleteEntries = 
append(classified.equalityDeleteEntries, entry)
+       case manifestEntryDV:
+               classified.dvEntries = append(classified.dvEntries, entry)
+       }
+}
+
+func (m *manifestEntries) merge(entries []iceberg.ManifestEntry) error {
+       classified, err := classifyManifestEntries(entries)

Review Comment:
   I'd restore `defer m.mu.Unlock()` here rather than unlocking explicitly at 
the end.
   
   Realistically these four appends won't panic, so this isn't a live deadlock, 
but it's the one spot that dropped the defer idiom every other method on `m.mu` 
still uses, and if any append ever panics (OOM, or a future change) the mutex 
stays held and every subsequent `merge` deadlocks. It's free to keep the 
guarantee: `m.mu.Lock()` then `defer m.mu.Unlock()` right at the top.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to