diff --git a/weed/plugin/worker/iceberg/exec_test.go b/weed/plugin/worker/iceberg/exec_test.go index 3f6b9534f..cbb4c088a 100644 --- a/weed/plugin/worker/iceberg/exec_test.go +++ b/weed/plugin/worker/iceberg/exec_test.go @@ -38,7 +38,7 @@ type fakeFilerServer struct { mu sync.Mutex entries map[string]map[string]*filer_pb.Entry // dir → name → entry - beforeUpdate func(*fakeFilerServer, *filer_pb.UpdateEntryRequest) + beforeUpdate func(*fakeFilerServer, *filer_pb.UpdateEntryRequest) error // Counters for assertions createCalls int @@ -145,7 +145,9 @@ func (f *fakeFilerServer) UpdateEntry(_ context.Context, req *filer_pb.UpdateEnt f.mu.Unlock() if beforeUpdate != nil { - beforeUpdate(f, req) + if err := beforeUpdate(f, req); err != nil { + return nil, err + } } f.mu.Lock() @@ -161,7 +163,13 @@ func (f *fakeFilerServer) UpdateEntry(_ context.Context, req *filer_pb.UpdateEnt } for key, expectedValue := range req.ExpectedExtended { actualValue, ok := current.Extended[key] - if !ok || !bytes.Equal(actualValue, expectedValue) { + if ok { + if !bytes.Equal(actualValue, expectedValue) { + return nil, status.Errorf(codes.FailedPrecondition, "extended attribute %q changed", key) + } + continue + } + if len(expectedValue) > 0 { return nil, status.Errorf(codes.FailedPrecondition, "extended attribute %q changed", key) } } @@ -307,7 +315,8 @@ func populateTable(t *testing.T, fs *fakeFilerServer, setup tableSetup) table.Me Name: setup.TableName, IsDirectory: true, Extended: map[string][]byte{ - s3tables.ExtendedKeyMetadata: xattr, + s3tables.ExtendedKeyMetadata: xattr, + s3tables.ExtendedKeyMetadataVersion: metadataVersionXattr(metadataVersion), }, }) @@ -1527,30 +1536,32 @@ func TestMetadataVersionCASDetectsConcurrentUpdate(t *testing.T) { populateTable(t, fs, setup) tableDir := path.Join(s3tables.TablesPath, setup.BucketName, setup.tablePath()) - fs.beforeUpdate = func(f *fakeFilerServer, req *filer_pb.UpdateEntryRequest) { + fs.beforeUpdate = func(f *fakeFilerServer, req *filer_pb.UpdateEntryRequest) error { entry := f.getEntry(path.Dir(tableDir), path.Base(tableDir)) if entry == nil { - t.Fatal("table entry not found before concurrent update") + return fmt.Errorf("table entry not found before concurrent update") } updatedEntry := cloneEntryForTest(t, entry) var internalMeta map[string]json.RawMessage if err := json.Unmarshal(updatedEntry.Extended[s3tables.ExtendedKeyMetadata], &internalMeta); err != nil { - t.Fatalf("unmarshal xattr: %v", err) + return fmt.Errorf("unmarshal xattr: %w", err) } versionJSON, err := json.Marshal(2) if err != nil { - t.Fatalf("marshal version: %v", err) + return fmt.Errorf("marshal version: %w", err) } internalMeta["metadataVersion"] = versionJSON updatedXattr, err := json.Marshal(internalMeta) if err != nil { - t.Fatalf("marshal xattr: %v", err) + return fmt.Errorf("marshal xattr: %w", err) } updatedEntry.Extended[s3tables.ExtendedKeyMetadata] = updatedXattr + updatedEntry.Extended[s3tables.ExtendedKeyMetadataVersion] = metadataVersionXattr(2) f.putEntry(path.Dir(tableDir), path.Base(tableDir), updatedEntry) + return nil } err := updateTableMetadataXattr(context.Background(), client, tableDir, 1, []byte(`{}`), "metadata/v2.metadata.json") diff --git a/weed/plugin/worker/iceberg/filer_io.go b/weed/plugin/worker/iceberg/filer_io.go index 40eae8b93..406ae4560 100644 --- a/weed/plugin/worker/iceberg/filer_io.go +++ b/weed/plugin/worker/iceberg/filer_io.go @@ -8,6 +8,7 @@ import ( "fmt" "io" "path" + "strconv" "strings" "sync" "time" @@ -364,12 +365,14 @@ func updateTableMetadataXattr(ctx context.Context, client filer_pb.SeaweedFilerC return fmt.Errorf("marshal updated xattr: %w", err) } + expectedVersionXattr := resp.Entry.Extended[s3tables.ExtendedKeyMetadataVersion] resp.Entry.Extended[s3tables.ExtendedKeyMetadata] = updatedXattr + resp.Entry.Extended[s3tables.ExtendedKeyMetadataVersion] = metadataVersionXattr(newVersion) _, err = client.UpdateEntry(ctx, &filer_pb.UpdateEntryRequest{ Directory: parentDir, Entry: resp.Entry, ExpectedExtended: map[string][]byte{ - s3tables.ExtendedKeyMetadata: existingXattr, + s3tables.ExtendedKeyMetadataVersion: expectedVersionXattr, }, }) if err != nil { @@ -381,6 +384,10 @@ func updateTableMetadataXattr(ctx context.Context, client filer_pb.SeaweedFilerC return nil } +func metadataVersionXattr(version int) []byte { + return []byte(strconv.Itoa(version)) +} + // generateIcebergVersionToken produces a random hex token, mirroring the // logic in s3tables.generateVersionToken (which is unexported). func generateIcebergVersionToken() string { diff --git a/weed/s3api/s3tables/handler.go b/weed/s3api/s3tables/handler.go index 1ead82309..2572a49c3 100644 --- a/weed/s3api/s3tables/handler.go +++ b/weed/s3api/s3tables/handler.go @@ -21,10 +21,11 @@ const ( DefaultRegion = "us-east-1" // Extended entry attributes for metadata storage - ExtendedKeyTableBucket = "s3tables.tableBucket" - ExtendedKeyMetadata = "s3tables.metadata" - ExtendedKeyPolicy = "s3tables.policy" - ExtendedKeyTags = "s3tables.tags" + ExtendedKeyTableBucket = "s3tables.tableBucket" + ExtendedKeyMetadata = "s3tables.metadata" + ExtendedKeyMetadataVersion = "s3tables.metadataVersion" + ExtendedKeyPolicy = "s3tables.policy" + ExtendedKeyTags = "s3tables.tags" // Maximum request body size (10MB) maxRequestBodySize = 10 * 1024 * 1024 diff --git a/weed/server/filer_grpc_server.go b/weed/server/filer_grpc_server.go index 076c3e9e9..1b7b050c0 100644 --- a/weed/server/filer_grpc_server.go +++ b/weed/server/filer_grpc_server.go @@ -252,7 +252,13 @@ func validateUpdateEntryPreconditions(entry *filer.Entry, expectedExtended map[s if entry != nil { actualValue, ok = entry.Extended[key] } - if !ok || !bytes.Equal(actualValue, expectedValue) { + if ok { + if !bytes.Equal(actualValue, expectedValue) { + return status.Errorf(codes.FailedPrecondition, "extended attribute %q changed", key) + } + continue + } + if len(expectedValue) > 0 { return status.Errorf(codes.FailedPrecondition, "extended attribute %q changed", key) } }