laskoviymishka commented on code in PR #2140:
URL: https://github.com/apache/iceberg-go/pull/2140#discussion_r4208252312
##########
table/cow_delete_conflict_test.go:
##########
@@ -198,3 +201,34 @@ func
TestCopyOnWriteConflict_ConcurrentDeleteOnUntouchedFileCommits(t *testing.T
}
}
}
+
+// TestCopyOnWriteConflict_ConcurrentRemovalOfSameFileDiverges pins the other
+// half of copy-on-write conflict detection: when a concurrent commit already
+// removed a data file this commit also removes, the retry rebuild's
+// checkRemovedFiles aborts terminally with ErrCommitDiverged, before any
+// validator runs, rather than rebuilding the file from the stale snapshot.
+func TestCopyOnWriteConflict_ConcurrentRemovalOfSameFileDiverges(t *testing.T)
{
+ for _, version := range cowFormatVersions {
+ for _, isolation := range cowIsolations {
+ for _, op := range cowOps {
+ for _, removal := range cowRemovals {
+ name := fmt.Sprintf("v%s/%s/%s/%s",
version, isolation, op, removal.name)
+ t.Run(name, func(t *testing.T) {
+ ctx := context.Background()
+ tbl := appendTenRows(t,
newCoWConflictTestTable(t, version, isolation))
+
+ txn := stageCopyOnWrite(t, tbl,
op, removal.filter)
+
+ // A concurrent copy-on-write
delete of id==4 rewrites the
+ // same data file first.
+ _, err := stageCopyOnWrite(t,
tbl, "delete", iceberg.EqualTo(iceberg.Reference("id"), int64(4))).Commit(ctx)
+ require.NoError(t, err)
+
+ _, err = txn.Commit(ctx)
+ require.ErrorIs(t, err,
table.ErrCommitDiverged)
Review Comment:
`ErrorIs(ErrCommitDiverged)` alone doesn't pin what the doc comment above
claims. That sentinel also falls out of `newConflictContext` (missing branch,
truncated ancestry walk), and since the concurrent commit adds no delete files
the rewrite validators would return nil here anyway, so this stays green
whether `checkRemovedFiles` runs first or a later refactor reorders the
validators ahead of it.
I'd tie it to the actual path: `require.ErrorContains(t, err, "no longer on
the branch head")`, plus `require.NotErrorIs(t, err, table.ErrCommitFailed)`
for the terminal/non-retry contract and `require.NotErrorIs(t, err,
table.ErrConflictingDeleteFiles)` to actually earn the "before any validator
runs" line. Otherwise I'd drop that phrase from the comment.
##########
table/cow_delete_conflict_test.go:
##########
@@ -198,3 +201,34 @@ func
TestCopyOnWriteConflict_ConcurrentDeleteOnUntouchedFileCommits(t *testing.T
}
}
}
+
+// TestCopyOnWriteConflict_ConcurrentRemovalOfSameFileDiverges pins the other
+// half of copy-on-write conflict detection: when a concurrent commit already
+// removed a data file this commit also removes, the retry rebuild's
+// checkRemovedFiles aborts terminally with ErrCommitDiverged, before any
+// validator runs, rather than rebuilding the file from the stale snapshot.
+func TestCopyOnWriteConflict_ConcurrentRemovalOfSameFileDiverges(t *testing.T)
{
+ for _, version := range cowFormatVersions {
+ for _, isolation := range cowIsolations {
+ for _, op := range cowOps {
+ for _, removal := range cowRemovals {
+ name := fmt.Sprintf("v%s/%s/%s/%s",
version, isolation, op, removal.name)
+ t.Run(name, func(t *testing.T) {
+ ctx := context.Background()
+ tbl := appendTenRows(t,
newCoWConflictTestTable(t, version, isolation))
+
+ txn := stageCopyOnWrite(t, tbl,
op, removal.filter)
+
+ // A concurrent copy-on-write
delete of id==4 rewrites the
+ // same data file first.
+ _, err := stageCopyOnWrite(t,
tbl, "delete", iceberg.EqualTo(iceberg.Reference("id"), int64(4))).Commit(ctx)
+ require.NoError(t, err)
+
+ _, err = txn.Commit(ctx)
+ require.ErrorIs(t, err,
table.ErrCommitDiverged)
+ })
Review Comment:
The subtest stops at the error and never checks the table survived the
abort. I'd reload and assert the rows are the concurrent result (ids 1-10 minus
4) with no duplicates. That's what proves the stale rewrite didn't land beside
the concurrent commit, which is the whole reason the abort exists. The sibling
control tests already use `idsInTable`.
##########
table/transaction.go:
##########
@@ -2496,12 +2496,12 @@ func (t *Transaction) performCopyOnWriteDeletion(ctx
context.Context, operation
// without this check a refresh-and-replay would swap the original file
// for a rewrite built from the stale snapshot, dropping the concurrent
// deletes and resurrecting their rows. Mirrors Java's copy-on-write
- // validateNoConflictingDeletes: no isolation gating.
- removed := append(slices.Clip(filesToDelete), filesToRewrite...)
+ // validateNoConflictingDeletes: no isolation gating. A concurrent
+ // removal of one of these files is caught earlier, by the retry
+ // rebuild's checkRemovedFiles (ErrCommitDiverged).
+ removed := slices.Concat(filesToDelete, filesToRewrite)
Review Comment:
`slices.Concat` is the right call and incidentally safer: the old
`append(slices.Clip(filesToDelete), ...)` handed back `filesToDelete`'s own
backing array when `filesToRewrite` was empty, so the validator closure could
alias it. Unreachable given the `len(removed) > 0` guard, but the new form
drops the hazard outright, worth a line in the commit message.
##########
table/conflict_validation.go:
##########
@@ -942,7 +942,7 @@ func validateNoNewDeletesForRewrittenFiles(ctx
*conflictContext, rewrittenFiles
case iceberg.EntryContentPosDeletes:
if path := referencedDataFilePath(df); path != "" {
if _, overlap := rewrittenPaths[path]; overlap {
- return fmt.Errorf("%w: snapshot %d
added pos-delete %s referencing rewritten file %s",
+ return fmt.Errorf("%w: snapshot %d
added pos-delete %s referencing removed data file %s",
Review Comment:
These messages now say "removed data file," but the function is still
`validateNoNewDeletesForRewrittenFiles`, the locals are
`rewrittenPaths`/`rewrittenPartitions`, and the doc comment below still says
"rewritten," so they now disagree. It's all unexported, so I'd either finish
the rename here, or drop the string edits and keep them with the PR that adds
coverage for these paths (none of the three are exercised by this test).
##########
table/conflict_validation.go:
##########
@@ -955,7 +955,7 @@ func validateNoNewDeletesForRewrittenFiles(ctx
*conflictContext, rewrittenFiles
df.FilePath(), df.SpecID(), err)
}
if _, overlap := rewrittenPartitions[key]; overlap {
- return fmt.Errorf("%w: snapshot %d added
partition-scoped pos-delete %s overlapping rewritten partition (spec %d)",
+ return fmt.Errorf("%w: snapshot %d added
partition-scoped pos-delete %s overlapping the partition of a removed data file
(spec %d)",
Review Comment:
small thing while here, the `(spec %d)` dangles after the long clause now.
Reads cleaner as `...overlapping the partition (spec %d) of a removed data
file`.
--
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]