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


##########
table/cow_delete_conflict_test.go:
##########
@@ -0,0 +1,200 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package table_test
+
+import (
+       "context"
+       "fmt"
+       "path/filepath"
+       "testing"
+
+       "github.com/apache/arrow-go/v18/arrow"
+       "github.com/apache/arrow-go/v18/arrow/array"
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/iceberg-go"
+       iceio "github.com/apache/iceberg-go/io"
+       "github.com/apache/iceberg-go/table"
+       "github.com/stretchr/testify/require"
+)
+
+// newCoWConflictTestTable builds a table that deletes via merge-on-read (so a
+// concurrent delete writes delete files rather than rewriting data) and
+// retries commits, so a lost CAS race is refreshed and replayed. Both the
+// delete and update isolation levels are set to isolation.
+func newCoWConflictTestTable(t *testing.T, formatVersion string, isolation 
table.IsolationLevel) *table.Table {
+       t.Helper()
+
+       location := filepath.ToSlash(t.TempDir())
+       schema := iceberg.NewSchema(0,
+               iceberg.NestedField{ID: 1, Name: "id", Type: 
iceberg.PrimitiveTypes.Int64, Required: true},
+               iceberg.NestedField{ID: 2, Name: "data", Type: 
iceberg.PrimitiveTypes.String, Required: false},
+       )
+       meta, err := table.NewMetadata(schema, iceberg.UnpartitionedSpec, 
table.UnsortedSortOrder, location,
+               iceberg.Properties{
+                       table.PropertyFormatVersion:        formatVersion,
+                       table.WriteDeleteModeKey:           
table.WriteModeMergeOnRead,
+                       table.WriteDeleteIsolationLevelKey: string(isolation),
+                       table.WriteUpdateIsolationLevelKey: string(isolation),
+                       table.CommitNumRetriesKey:          "2",
+                       table.CommitMinRetryWaitMsKey:      "1",
+                       table.CommitMaxRetryWaitMsKey:      "2",
+                       table.CommitTotalRetryTimeoutMsKey: "1000",
+               })
+       require.NoError(t, err)
+
+       metaLoc := location + "/metadata/v1.metadata.json"
+       fsF := func(context.Context) (iceio.IO, error) { return 
iceio.LocalFS{}, nil }
+       cat := &concurrentTestCatalog{metadata: meta, location: metaLoc, fsF: 
fsF}
+
+       return table.New(table.Identifier{"db", "cow_delete_conflict"}, meta, 
metaLoc, fsF, cat)
+}
+
+func cowTestRecords(t *testing.T, rowsJSON string) array.RecordReader {
+       t.Helper()
+
+       arrowSchema := arrow.NewSchema([]arrow.Field{
+               {Name: "id", Type: arrow.PrimitiveTypes.Int64, Nullable: false},
+               {Name: "data", Type: arrow.BinaryTypes.String, Nullable: true},
+       }, nil)
+       data, err := array.TableFromJSON(memory.DefaultAllocator, arrowSchema, 
[]string{rowsJSON})
+       require.NoError(t, err)
+       t.Cleanup(data.Release)
+
+       return array.NewTableReader(data, -1)
+}
+
+// stageCopyOnWrite stages, without committing, a copy-on-write removal of the
+// rows matching filter through either Delete or Overwrite (with no new rows).
+func stageCopyOnWrite(t *testing.T, tbl *table.Table, op string, filter 
iceberg.BooleanExpression) *table.Transaction {
+       t.Helper()
+       ctx := context.Background()
+
+       txn := tbl.NewTransaction()
+       switch op {
+       case "delete":
+               require.NoError(t, txn.SetProperties(iceberg.Properties{
+                       table.WriteDeleteModeKey: table.WriteModeCopyOnWrite,
+               }))
+               require.NoError(t, txn.Delete(ctx, filter, nil))
+       case "overwrite":
+               require.NoError(t, txn.Overwrite(ctx, cowTestRecords(t, `[]`), 
nil, table.WithOverwriteFilter(filter)))
+       default:
+               t.Fatalf("unknown op %q", op)
+       }
+
+       return txn
+}
+
+var (
+       cowFormatVersions = []string{"2", "3"}
+       cowIsolations     = []table.IsolationLevel{table.IsolationSerializable, 
table.IsolationSnapshot}
+       cowOps            = []string{"delete", "overwrite"}
+)
+
+// TestCopyOnWriteConflict_ConcurrentDeleteOnRemovedFile proves a copy-on-write
+// Delete or Overwrite is rejected, at every isolation level, when a concurrent
+// commit added a delete against a data file it removes. Without the validator,
+// refresh-and-replay swaps the original file for a rewrite built from the 
stale
+// snapshot, dropping the concurrent delete and resurrecting its row (#2090).
+func TestCopyOnWriteConflict_ConcurrentDeleteOnRemovedFile(t *testing.T) {
+       removals := []struct {
+               name   string
+               filter iceberg.BooleanExpression
+       }{
+               // id==2 matches part of the file, so it is rewritten.
+               {"rewritten", iceberg.EqualTo(iceberg.Reference("id"), 
int64(2))},
+               // id>=1 matches every row, so the file is dropped outright.
+               {"fully deleted", 
iceberg.GreaterThanEqual(iceberg.Reference("id"), int64(1))},
+       }
+
+       for _, version := range cowFormatVersions {
+               for _, isolation := range cowIsolations {
+                       for _, op := range cowOps {
+                               for _, removal := range removals {
+                                       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 merge-on-read 
delete of id==4 commits first.
+                                               _, err := tbl.Delete(ctx, 
iceberg.EqualTo(iceberg.Reference("id"), int64(4)), nil)

Review Comment:
   The concurrent writer here only ever stages a merge-on-read delete, so on v2 
that's a position delete and on v3 a DV; the validator's equality-delete branch 
never fires in any of these runs. If we can stage a concurrent eq-delete 
writer, a case there would cover the one path the matrix misses; if iceberg-go 
can't write eq-deletes yet, calling that branch out as untested is fine.



##########
table/cow_delete_conflict_test.go:
##########
@@ -0,0 +1,200 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package table_test
+
+import (
+       "context"
+       "fmt"
+       "path/filepath"
+       "testing"
+
+       "github.com/apache/arrow-go/v18/arrow"
+       "github.com/apache/arrow-go/v18/arrow/array"
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/iceberg-go"
+       iceio "github.com/apache/iceberg-go/io"
+       "github.com/apache/iceberg-go/table"
+       "github.com/stretchr/testify/require"
+)
+
+// newCoWConflictTestTable builds a table that deletes via merge-on-read (so a
+// concurrent delete writes delete files rather than rewriting data) and
+// retries commits, so a lost CAS race is refreshed and replayed. Both the
+// delete and update isolation levels are set to isolation.
+func newCoWConflictTestTable(t *testing.T, formatVersion string, isolation 
table.IsolationLevel) *table.Table {
+       t.Helper()
+
+       location := filepath.ToSlash(t.TempDir())
+       schema := iceberg.NewSchema(0,
+               iceberg.NestedField{ID: 1, Name: "id", Type: 
iceberg.PrimitiveTypes.Int64, Required: true},
+               iceberg.NestedField{ID: 2, Name: "data", Type: 
iceberg.PrimitiveTypes.String, Required: false},
+       )
+       meta, err := table.NewMetadata(schema, iceberg.UnpartitionedSpec, 
table.UnsortedSortOrder, location,
+               iceberg.Properties{
+                       table.PropertyFormatVersion:        formatVersion,
+                       table.WriteDeleteModeKey:           
table.WriteModeMergeOnRead,
+                       table.WriteDeleteIsolationLevelKey: string(isolation),
+                       table.WriteUpdateIsolationLevelKey: string(isolation),
+                       table.CommitNumRetriesKey:          "2",
+                       table.CommitMinRetryWaitMsKey:      "1",
+                       table.CommitMaxRetryWaitMsKey:      "2",
+                       table.CommitTotalRetryTimeoutMsKey: "1000",
+               })
+       require.NoError(t, err)
+
+       metaLoc := location + "/metadata/v1.metadata.json"
+       fsF := func(context.Context) (iceio.IO, error) { return 
iceio.LocalFS{}, nil }
+       cat := &concurrentTestCatalog{metadata: meta, location: metaLoc, fsF: 
fsF}
+
+       return table.New(table.Identifier{"db", "cow_delete_conflict"}, meta, 
metaLoc, fsF, cat)
+}
+
+func cowTestRecords(t *testing.T, rowsJSON string) array.RecordReader {
+       t.Helper()
+
+       arrowSchema := arrow.NewSchema([]arrow.Field{
+               {Name: "id", Type: arrow.PrimitiveTypes.Int64, Nullable: false},
+               {Name: "data", Type: arrow.BinaryTypes.String, Nullable: true},
+       }, nil)
+       data, err := array.TableFromJSON(memory.DefaultAllocator, arrowSchema, 
[]string{rowsJSON})
+       require.NoError(t, err)
+       t.Cleanup(data.Release)
+
+       return array.NewTableReader(data, -1)

Review Comment:
   `array.NewTableReader` retains `data`, but only `data.Release` is registered 
for cleanup, so the reader itself is never released. Harmless under the default 
allocator, but it leaks a ref under a `CheckedAllocator`. I'd create the reader 
first and register `t.Cleanup(rdr.Release)` alongside the table release, unless 
the `Append`/`Overwrite` callers already release it.



##########
table/transaction.go:
##########
@@ -2478,6 +2478,20 @@ func (t *Transaction) performCopyOnWriteDeletion(ctx 
context.Context, operation
                }
        }
 
+       // Reject the commit if a concurrent snapshot added deletes against any
+       // data file this operation removes. The overwriteFiles producer only
+       // checks added data files (and only under serializable isolation);
+       // 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...)

Review Comment:
   The `slices.Clip` plus `append` is correct: clipping caps `filesToDelete` so 
the append forces a fresh backing array and we never write into its spare 
capacity. It reads as non-obvious enough that I'd either add a one-line comment 
on why the clip is load-bearing, or switch to `slices.Concat(filesToDelete, 
filesToRewrite)` (1.22+), which says the same thing without the idiom.



##########
table/transaction.go:
##########
@@ -2478,6 +2478,20 @@ func (t *Transaction) performCopyOnWriteDeletion(ctx 
context.Context, operation
                }
        }
 
+       // Reject the commit if a concurrent snapshot added deletes against any
+       // data file this operation removes. The overwriteFiles producer only
+       // checks added data files (and only under serializable isolation);
+       // 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...)
+       if len(removed) > 0 {
+               t.addValidator(func(cc *conflictContext) error {

Review Comment:
   We register the new-deletes half of the check here but not the 
deleted-data-files half. `performMergeOnReadDeletion` also registers 
`validateDataFilesExist`, and Java's `validateNoConflictingDeletes` rejects 
when a concurrent commit removed or rewrote a file this operation is also 
removing. Without that, if a concurrent CoW delete or compaction rewrites one 
of these files first, the replay re-removes a file that's already gone and 
rebuilds it from the stale snapshot. I'd confirm `overwriteFiles` already 
rejects that on its own; if it doesn't, I'd add `validateDataFilesExist(cc, 
removed)` next to this registration or track it as a follow-up.



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