laskoviymishka commented on code in PR #2057:
URL: https://github.com/apache/iceberg-go/pull/2057#discussion_r4149598173
##########
table/snapshot_producers.go:
##########
@@ -507,18 +507,32 @@ func (m *manifestMergeManager) createManifest(specID int,
bin []iceberg.Manifest
return nil, err
}
- wr, path, counter, fileCloser, err := m.snap.newManifestWriter(spec)
- if err != nil {
- return nil, err
- }
- defer internal.CheckedClose(fileCloser, &err)
+ var wr *iceberg.ManifestWriter
+ var path string
+ var counter *internal.CountingWriter
+ var fileCloser io.Closer
writerClosed := false
+ // Close the ManifestWriter before the underlying file so a mid-loop
+ // write failure flushes into an open file, not a closed one.
defer func() {
- if !writerClosed {
+ if wr != nil && !writerClosed {
internal.CheckedClose(wr, &err)
}
+ if fileCloser != nil {
+ internal.CheckedClose(fileCloser, &err)
+ }
}()
+ ensureWriter := func() error {
+ if wr != nil {
+ return nil
+ }
+
+ wr, path, counter, fileCloser, err =
m.snap.newManifestWriter(spec)
Review Comment:
Still here from last round: `ensureWriter` writes into the outer
named-return `err` and returns the same value, so every call site
double-assigns it. It's harmless today, but it's a footgun: a call site
switching to `if err := ensureWriter()` would shadow locally and silently leak
the error into the named return and the deferred closes. I'd have the closure
use a local and return it:
```go
ensureWriter := func() error {
if wr != nil {
return nil
}
var err error
wr, path, counter, fileCloser, err = m.snap.newManifestWriter(spec)
return err
}
```
##########
table/snapshot_producers_test.go:
##########
@@ -578,6 +578,130 @@ func TestManifestMergeManagerClosesWriterOnError(t
*testing.T) {
require.ErrorIs(t, err, errLimitedWrite)
}
+func TestManifestMergeManagerClosesWriterBeforeFileOnWriteFailure(t
*testing.T) {
+ spec := iceberg.NewPartitionSpec()
+ schema := simpleSchema()
+
+ // Use a byte-limited IO that fails after the writer is opened and
+ // has written some data, but before all entries are processed.
+ mem := newMemIO(manifestHeaderSize(t, 2, spec, schema), errLimitedWrite)
Review Comment:
This test can't fail on the ordering it's named for. `limitedWriteCloser`
has no write-after-close detection; its `Write` only checks the byte budget, so
if `fileCloser` closed before `wr`, the `wr.Close()` flush would still just
return `errLimitedWrite`, and the final "assertion" (a prose comment, not a
`require`) passes either way. That leaves the round-1 blocker guarded by a test
that can't catch a regression.
I'd switch to `trackingIO`/`trackingWriteCloser`, fail the write after the
header lands, and assert what `TestCommitManifestsCloseFailureReturnsNoUpdates`
already does:
```go
require.NotContains(t, err.Error(), "write after close",
"manifest writer must be closed before its underlying file closer")
```
That's the piece that actually locks the ordering down.
##########
table/snapshot_producers.go:
##########
@@ -543,6 +566,13 @@ func (m *manifestMergeManager) createManifest(specID int,
bin []iceberg.Manifest
}
}
+ // A bin with no live entries produces no manifest. This diverges from
+ // Java/PyIceberg, which write a zero-count manifest; the omission is
Review Comment:
Small thing on the new comment: Java does write the zero-count manifest, but
PyIceberg has merge disabled by default (`MANIFEST_MERGE_ENABLED_DEFAULT =
False`), so I couldn't confirm its merge path does the same. I'd narrow this to
"the Java reference implementation" rather than "Java/PyIceberg" unless
someone's actually traced PyIceberg's merge.
##########
table/snapshot_producers.go:
##########
@@ -1762,16 +1795,13 @@ func (sp *snapshotProducer)
commitManifests(newManifests, addedContent []iceberg
// creates it).
baseHeadID := sp.txn.baseRefSnapshotID(branch)
- return []Update{
- addSnap,
- // Carry over the branch's existing retention settings
so advancing
- // the ref on commit does not silently discard them.
The update
- // encodes exactly the current ref's retention
(settings the branch
- // lacks stay 0 and are dropped by the `omitempty`
tags); the catalog
- // applies a set-snapshot-ref as a pure replace, so
this fully
- // determines the resulting ref rather than merging
with the old one.
- sp.txn.meta.NewRetainingSnapshotRefUpdate(branch,
sp.snapshotID, BranchRef),
- }, []Requirement{
- AssertRefSnapshotID(branch, baseHeadID),
- }, nil
+ // Carry over the branch's existing retention settings so advancing
+ // the ref on commit does not silently discard them. The update
+ // encodes exactly the current ref's retention (settings the branch
+ // lacks stay 0 and are dropped by the `omitempty` tags); the catalog
+ // applies a set-snapshot-ref as a pure replace, so this fully
+ // determines the resulting ref rather than merging with the old one.
+ retainingSnapshotRef :=
sp.txn.meta.NewRetainingSnapshotRefUpdate(branch, sp.snapshotID, BranchRef)
Review Comment:
Same ask as last round: the `commitManifests` reflow (extracting
`retainingSnapshotRef`) and the trailing-comma reformats aren't part of the
fix. I'd split them into a separate tidy commit so the bug-fix diff stays
focused; the trailing commas may be `gofumpt`, but they still belong on their
own.
--
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]