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


##########
table/snapshot_producers_test.go:
##########
@@ -578,6 +578,208 @@ func TestManifestMergeManagerClosesWriterOnError(t 
*testing.T) {
        require.ErrorIs(t, err, errLimitedWrite)
 }
 
+func TestManifestMergeManagerClosesWriterBeforeFileOnReadFailure(t *testing.T) 
{
+       spec := iceberg.NewPartitionSpec()
+       schema := simpleSchema()
+       trackIO := newTrackingIO()
+       txn := createTestTransaction(t, trackIO, spec)
+       sp := newFastAppendFilesProducer(OpAppend, txn, trackIO, nil, nil)
+       df := newTestDataFile(t, spec, "file://data-1.parquet", nil)
+       manifest := writeTestManifestWithContent(t, trackIO, spec, schema, 
sp.snapshotID,
+               "table-location/metadata/manifest-1.avro", 
iceberg.ManifestContentData,
+               []iceberg.ManifestEntry{
+                       iceberg.NewManifestEntry(iceberg.EntryStatusADDED, 
&sp.snapshotID, nil, nil, df),
+               })
+       missingManifest := iceberg.NewManifestFile(2, 
"table-location/metadata/missing.avro", 1,
+               int32(spec.ID()), sp.snapshotID).Build()
+
+       mgr := manifestMergeManager{snap: sp}
+       _, err := mgr.createManifest(spec.ID(), 
[]iceberg.ManifestFile{manifest, missingManifest})
+       require.ErrorIs(t, err, fs.ErrNotExist)
+       require.NotContains(t, err.Error(), "write after close",
+               "buffered entries must flush before the underlying file closes")
+       require.Equal(t, 1, trackIO.GetWriterCount())
+       require.Empty(t, trackIO.GetUnclosedWriters())
+}
+
+func TestManifestMergeReturnsReadErrorBeforeCreatingWriter(t *testing.T) {
+       spec := iceberg.NewPartitionSpec()
+       schema := simpleSchema()
+       trackIO := newTrackingIO()
+       txn := createTestTransaction(t, trackIO, spec)
+       sp := newFastAppendFilesProducer(OpAppend, txn, trackIO, nil, nil)
+
+       oldSnapshotID := sp.snapshotID - 1
+       sequenceNumber := int64(1)
+       df := newTestDataFile(t, spec, "file://deleted.parquet", nil)
+       deletedManifest := writeTestManifestWithContent(t, trackIO, spec, 
schema, oldSnapshotID,
+               "table-location/metadata/historical-delete.avro", 
iceberg.ManifestContentData,
+               []iceberg.ManifestEntry{
+                       iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, 
&oldSnapshotID, &sequenceNumber, nil, df),
+               })
+       missingManifest := iceberg.NewManifestFile(2, 
"table-location/metadata/missing.avro", 1,
+               int32(spec.ID()), sp.snapshotID).Build()
+
+       mgr := manifestMergeManager{snap: sp}
+       _, err := mgr.createManifest(spec.ID(), 
[]iceberg.ManifestFile{deletedManifest, missingManifest})
+       require.ErrorIs(t, err, fs.ErrNotExist)
+       require.Zero(t, trackIO.GetWriterCount())
+}
+
+func TestManifestMergeSkipsHistoricalDeletedOnlyManifest(t *testing.T) {
+       spec := iceberg.NewPartitionSpec()
+       schema := simpleSchema()
+       trackIO := newTrackingIO()
+       txn := createTestTransaction(t, trackIO, spec)
+       sp := newFastAppendFilesProducer(OpAppend, txn, trackIO, nil, nil)
+
+       oldSnapshotID := sp.snapshotID - 1
+       sequenceNumber := int64(1)
+       df := newTestDataFile(t, spec, "file://deleted.parquet", nil)
+       manifestFile := writeTestManifestWithContent(t, trackIO, spec, schema, 
oldSnapshotID,
+               "table-location/metadata/historical-delete.avro", 
iceberg.ManifestContentData,
+               []iceberg.ManifestEntry{
+                       iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, 
&oldSnapshotID, &sequenceNumber, nil, df),
+               })
+
+       trackIO.writers = make(map[string]*trackingWriteCloser)

Review Comment:
   `writeTestManifestWithContent` persists through `fs.WriteFile`, never 
`Create`, so no `trackingWriteCloser` is ever registered during setup and this 
reset just clears an already-empty map. It reads like setup writes go through 
`Create` and need clearing, which future tests might cargo-cult into a spot 
where it masks a real writer leak. I'd drop the line; the `require.Zero(t, 
trackIO.GetWriterCount())` below already covers it (and it also mutates the 
field without `writersMu`).



##########
table/snapshot_producers.go:
##########
@@ -527,12 +542,21 @@ func (m *manifestMergeManager) createManifest(specID int, 
bin []iceberg.Manifest
 
                        switch {
                        case entry.Status() == iceberg.EntryStatusDELETED && 
entry.SnapshotID() == m.snap.snapshotID:
+                               if err = ensureWriter(); err != nil {

Review Comment:
   The dead write into the named `err` is gone now that `ensureWriter` returns 
its own `writerErr`, nice. The one thing left is that the loop-scoped `err` 
from `for entry, err := range` at line 538 shadows the named return, and these 
new `ensureWriter` assignments make that shadow easier to misread. Renaming it 
to `iterErr` (or a one-line note) would settle it. Not blocking.



##########
table/snapshot_producers_test.go:
##########
@@ -578,6 +578,208 @@ func TestManifestMergeManagerClosesWriterOnError(t 
*testing.T) {
        require.ErrorIs(t, err, errLimitedWrite)
 }
 
+func TestManifestMergeManagerClosesWriterBeforeFileOnReadFailure(t *testing.T) 
{
+       spec := iceberg.NewPartitionSpec()
+       schema := simpleSchema()
+       trackIO := newTrackingIO()
+       txn := createTestTransaction(t, trackIO, spec)
+       sp := newFastAppendFilesProducer(OpAppend, txn, trackIO, nil, nil)
+       df := newTestDataFile(t, spec, "file://data-1.parquet", nil)
+       manifest := writeTestManifestWithContent(t, trackIO, spec, schema, 
sp.snapshotID,
+               "table-location/metadata/manifest-1.avro", 
iceberg.ManifestContentData,
+               []iceberg.ManifestEntry{
+                       iceberg.NewManifestEntry(iceberg.EntryStatusADDED, 
&sp.snapshotID, nil, nil, df),
+               })
+       missingManifest := iceberg.NewManifestFile(2, 
"table-location/metadata/missing.avro", 1,
+               int32(spec.ID()), sp.snapshotID).Build()
+
+       mgr := manifestMergeManager{snap: sp}
+       _, err := mgr.createManifest(spec.ID(), 
[]iceberg.ManifestFile{manifest, missingManifest})
+       require.ErrorIs(t, err, fs.ErrNotExist)
+       require.NotContains(t, err.Error(), "write after close",
+               "buffered entries must flush before the underlying file closes")
+       require.Equal(t, 1, trackIO.GetWriterCount())
+       require.Empty(t, trackIO.GetUnclosedWriters())
+}
+
+func TestManifestMergeReturnsReadErrorBeforeCreatingWriter(t *testing.T) {
+       spec := iceberg.NewPartitionSpec()
+       schema := simpleSchema()
+       trackIO := newTrackingIO()
+       txn := createTestTransaction(t, trackIO, spec)
+       sp := newFastAppendFilesProducer(OpAppend, txn, trackIO, nil, nil)
+
+       oldSnapshotID := sp.snapshotID - 1
+       sequenceNumber := int64(1)
+       df := newTestDataFile(t, spec, "file://deleted.parquet", nil)
+       deletedManifest := writeTestManifestWithContent(t, trackIO, spec, 
schema, oldSnapshotID,
+               "table-location/metadata/historical-delete.avro", 
iceberg.ManifestContentData,
+               []iceberg.ManifestEntry{
+                       iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, 
&oldSnapshotID, &sequenceNumber, nil, df),
+               })
+       missingManifest := iceberg.NewManifestFile(2, 
"table-location/metadata/missing.avro", 1,
+               int32(spec.ID()), sp.snapshotID).Build()
+
+       mgr := manifestMergeManager{snap: sp}
+       _, err := mgr.createManifest(spec.ID(), 
[]iceberg.ManifestFile{deletedManifest, missingManifest})
+       require.ErrorIs(t, err, fs.ErrNotExist)
+       require.Zero(t, trackIO.GetWriterCount())
+}
+
+func TestManifestMergeSkipsHistoricalDeletedOnlyManifest(t *testing.T) {
+       spec := iceberg.NewPartitionSpec()
+       schema := simpleSchema()
+       trackIO := newTrackingIO()
+       txn := createTestTransaction(t, trackIO, spec)
+       sp := newFastAppendFilesProducer(OpAppend, txn, trackIO, nil, nil)
+
+       oldSnapshotID := sp.snapshotID - 1
+       sequenceNumber := int64(1)
+       df := newTestDataFile(t, spec, "file://deleted.parquet", nil)
+       manifestFile := writeTestManifestWithContent(t, trackIO, spec, schema, 
oldSnapshotID,
+               "table-location/metadata/historical-delete.avro", 
iceberg.ManifestContentData,
+               []iceberg.ManifestEntry{
+                       iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, 
&oldSnapshotID, &sequenceNumber, nil, df),
+               })
+
+       trackIO.writers = make(map[string]*trackingWriteCloser)
+
+       mgr := manifestMergeManager{snap: sp}
+       created, err := mgr.createManifest(spec.ID(), 
[]iceberg.ManifestFile{manifestFile})
+       require.NoError(t, err)
+       require.Nil(t, created)
+       require.Zero(t, trackIO.GetWriterCount())
+}
+
+func TestManifestMergeKeepsCurrentSnapshotDeletedEntries(t *testing.T) {
+       spec := iceberg.NewPartitionSpec()
+       schema := simpleSchema()
+       txn, wfs := createTestTransactionWithMemIO(t, spec)
+       sp := newFastAppendFilesProducer(OpAppend, txn, wfs, nil, nil)
+
+       sequenceNumber := int64(1)
+       df := newTestDataFile(t, spec, "file://deleted.parquet", nil)
+       manifestFile := writeTestManifestWithContent(t, wfs, spec, schema, 
sp.snapshotID,
+               "mem://default/table-location/metadata/current-delete.avro", 
iceberg.ManifestContentData,
+               []iceberg.ManifestEntry{
+                       iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, 
&sp.snapshotID, &sequenceNumber, nil, df),
+               })
+
+       mgr := manifestMergeManager{snap: sp}
+       created, err := mgr.createManifest(spec.ID(), 
[]iceberg.ManifestFile{manifestFile})
+       require.NoError(t, err)
+       require.NotNil(t, created)
+
+       var entries []iceberg.ManifestEntry
+       for entry, err := range created.Entries(wfs, false) {
+               require.NoError(t, err)
+               entries = append(entries, entry)
+       }
+       require.Len(t, entries, 1)
+       require.Equal(t, iceberg.EntryStatusDELETED, entries[0].Status())
+       require.Equal(t, sp.snapshotID, entries[0].SnapshotID())
+       require.Equal(t, df.FilePath(), entries[0].DataFile().FilePath())
+}
+
+func TestManifestMergeGroupDropsEmptyMergedBin(t *testing.T) {
+       spec := iceberg.NewPartitionSpec()
+       schema := simpleSchema()
+       txn, wfs := createTestTransactionWithMemIO(t, spec)
+       sp := newFastAppendFilesProducer(OpAppend, txn, wfs, nil, nil)
+
+       oldSnapshotID := sp.snapshotID - 1
+       sequenceNumber := int64(1)
+       manifests := make([]iceberg.ManifestFile, 0, 2)
+       for i := range 2 {
+               df := newTestDataFile(t, spec, 
"file://deleted-"+strconv.Itoa(i)+".parquet", nil)
+               manifests = append(manifests, writeTestManifestWithContent(t, 
wfs, spec, schema, oldSnapshotID,
+                       
"mem://default/table-location/metadata/historical-delete-"+strconv.Itoa(i)+".avro",
+                       iceberg.ManifestContentData,
+                       []iceberg.ManifestEntry{
+                               
iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, &oldSnapshotID, 
&sequenceNumber, nil, df),
+                       }))
+       }
+
+       mgr := manifestMergeManager{
+               targetSizeBytes:  manifests[0].Length() + manifests[1].Length(),
+               minCountToMerge:  1,
+               mergeEnabled:     true,
+               mergeConcurrency: 1,
+               snap:             sp,
+       }
+
+       merged, err := mgr.mergeGroup(manifests[0], spec.ID(), manifests)
+       require.NoError(t, err)
+       require.Empty(t, merged)
+}
+
+func TestMergeAppendCommitSkipsHistoricalDeleteOnlyBin(t *testing.T) {
+       spec := iceberg.NewPartitionSpec()
+       txn, wfs := createTestTransactionWithMemIO(t, spec)
+       schema := simpleSchema()
+       parentSnapshotID := txn.meta.newSnapshotID()
+       sequenceNumber := int64(1)
+       historicalManifests := make([]iceberg.ManifestFile, 0, 2)
+       for i := range 2 {
+               dataFile := newTestDataFile(t, spec, 
"file://deleted-"+strconv.Itoa(i)+".parquet", nil)
+               historicalManifests = append(historicalManifests, 
writeTestManifestWithContent(t, wfs, spec, schema,
+                       parentSnapshotID, 
"mem://default/table-location/metadata/historical-delete-bin-"+strconv.Itoa(i)+".avro",
+                       iceberg.ManifestContentData, []iceberg.ManifestEntry{
+                               
iceberg.NewManifestEntry(iceberg.EntryStatusDELETED, &parentSnapshotID, 
&sequenceNumber, nil, dataFile),
+                       }))
+       }
+       manifestListPath := 
"mem://default/table-location/metadata/snap-historical-delete-bin.avro"
+       manifestList, err := wfs.Create(manifestListPath)
+       require.NoError(t, err)
+       require.NoError(t, iceberg.WriteManifestList(
+               2, manifestList, parentSnapshotID, nil, &sequenceNumber, 0, 
historicalManifests))
+       require.NoError(t, manifestList.Close())
+       parentSnapshot := &Snapshot{
+               SnapshotID:     parentSnapshotID,
+               SequenceNumber: sequenceNumber,
+               TimestampMs:    time.Now().UnixMilli(),
+               ManifestList:   manifestListPath,
+               Summary:        &Summary{Operation: OpAppend, Properties: 
iceberg.Properties{}},
+       }
+       require.NoError(t, txn.meta.AddSnapshot(parentSnapshot))
+       require.NoError(t, txn.meta.SetSnapshotRef(MainBranch, 
parentSnapshotID, BranchRef))
+       meta, err := txn.meta.Build()
+       require.NoError(t, err)
+
+       mergeTxn := New(Identifier{"db", "tbl"}, meta, "metadata.json", 
func(context.Context) (iceio.IO, error) {
+               return wfs, nil
+       }, nil).NewTransaction()
+       if mergeTxn.meta.props == nil {
+               mergeTxn.meta.props = make(iceberg.Properties)
+       }
+       mergeTxn.meta.props[ManifestMergeEnabledKey] = "true"
+       mergeTxn.meta.props[ManifestMinMergeCountKey] = "2"
+       mergeTxn.meta.props[ManifestTargetSizeBytesKey] = strconv.FormatInt(
+               
historicalManifests[0].Length()+historicalManifests[1].Length(), 10)
+       mergeAppend := newMergeAppendFilesProducer(OpAppend, mergeTxn, wfs, 
nil, nil)
+       newFile := newTestDataFile(t, spec, "file://new-data.parquet", nil)
+       mergeAppend.appendDataFile(newFile)
+       updates, requirements, err := mergeAppend.commit(context.Background())
+       require.NoError(t, err)
+       require.NoError(t, mergeTxn.apply(updates, requirements))
+       meta, err = mergeTxn.meta.Build()
+       require.NoError(t, err)
+
+       outputManifests, err := meta.CurrentSnapshot().Manifests(wfs)
+       require.NoError(t, err)
+       require.Len(t, outputManifests, 1)
+       require.Equal(t, int32(1), outputManifests[0].AddedDataFiles())
+       var outputEntries []iceberg.ManifestEntry
+       for entry, entryErr := range outputManifests[0].Entries(wfs, false) {
+               require.NoError(t, entryErr)
+               outputEntries = append(outputEntries, entry)
+       }
+       require.Len(t, outputEntries, 1)
+       require.Equal(t, iceberg.EntryStatusADDED, outputEntries[0].Status())
+       require.Equal(t, newFile.FilePath(), 
outputEntries[0].DataFile().FilePath())
+       t.Logf("snapshot commit succeeded with %d output manifest containing 
the new data file", len(outputManifests))

Review Comment:
   `len(outputManifests)` is already pinned by `require.Len(..., 1)` just 
above, so this log is noise on the success path. I'd drop it.



-- 
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