From 10f3a8590be166637118c6cbc21dc97fbdde21ab Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Wed, 2 Sep 2026 01:15:53 +0800 Subject: [PATCH] fix: harden bucket metadata creation edge cases Preserve existing records only for ForceCreate, reject ghost metadata on genuine creation, keep object-lock versioning invariants, and complete metadata saves after caller cancellation. Expand deterministic coverage for peer bulk, lifecycle delete, ghost creation, and cancellation.\n\nRefs: #102 Signed-off-by: Feng Ruohang --- cmd/bucket-metadata-lock_test.go | 187 +++++++++++++++++++++++++++++++ cmd/erasure-server-pool.go | 27 +++-- cmd/site-replication.go | 8 +- 3 files changed, 207 insertions(+), 15 deletions(-) diff --git a/cmd/bucket-metadata-lock_test.go b/cmd/bucket-metadata-lock_test.go index cdcacefaa..4e2a4af83 100644 --- a/cmd/bucket-metadata-lock_test.go +++ b/cmd/bucket-metadata-lock_test.go @@ -20,7 +20,10 @@ import ( "testing" "time" + "github.com/minio/madmin-go/v3" "github.com/minio/minio/internal/auth" + "github.com/minio/minio/internal/bucket/lifecycle" + "github.com/minio/minio/internal/bucket/versioning" ) type metadataRMWWriterKey struct{} @@ -33,6 +36,8 @@ type metadataRMWBarrierObjectLayer struct { bLockAttempt chan struct{} aReadyOnce sync.Once bLockOnce sync.Once + cancelOnce sync.Once + cancelOnPut context.CancelFunc reads atomic.Int64 } @@ -52,6 +57,9 @@ func (o *metadataRMWBarrierObjectLayer) GetObjectNInfo(ctx context.Context, buck } func (o *metadataRMWBarrierObjectLayer) PutObject(ctx context.Context, bucket, object string, data *PutObjReader, opts ObjectOptions) (ObjectInfo, error) { + if bucket == minioMetaBucket && object == o.metadataObject() && o.cancelOnPut != nil { + o.cancelOnce.Do(o.cancelOnPut) + } if bucket == minioMetaBucket && object == o.metadataObject() && ctx.Value(metadataRMWWriterKey{}) == "A" { o.aReadyOnce.Do(func() { close(o.aReady) }) select { @@ -120,6 +128,69 @@ func TestBucketMetadataLockPreservesTaggingAndSSE(t *testing.T) { }) } +func TestBucketMetadataLockPreservesPeerBulkAndLocalUpdate(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + objAPITest: testBucketMetadataLockPreservesPeerBulkAndLocalUpdate, + }) +} + +func testBucketMetadataLockPreservesPeerBulkAndLocalUpdate(obj ObjectLayer, instanceType, bucket string, + _ http.Handler, _ auth.Credentials, t *testing.T, +) { + meta, err := readBucketMetadata(t.Context(), obj, bucket) + if err != nil { + t.Fatal(err) + } + policyJSON := fmt.Appendf(nil, `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket) + tagXML := []byte(`localtag`) + runBucketMetadataRMWConflict(t, obj, bucket, + func(ctx context.Context, objectAPI ObjectLayer) error { + return globalSiteReplicationSys.PeerBucketMetadataUpdateHandler(ctx, madmin.SRBucketMeta{ + Bucket: bucket, Policy: policyJSON, UpdatedAt: meta.Created.Add(time.Second), + }) + }, + func(ctx context.Context, objectAPI ObjectLayer) error { + _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketTaggingConfig, tagXML) + return err + }, + func(meta BucketMetadata) bool { + return bytes.Equal(meta.PolicyConfigJSON, policyJSON) && bytes.Equal(meta.TaggingConfigXML, tagXML) + }, instanceType+": peer bulk+local tagging") +} + +func TestBucketMetadataLockPreservesLifecycleDeleteAndSSE(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + objAPITest: testBucketMetadataLockPreservesLifecycleDeleteAndSSE, + }) +} + +func testBucketMetadataLockPreservesLifecycleDeleteAndSSE(obj ObjectLayer, instanceType, bucket string, + _ http.Handler, _ auth.Credentials, t *testing.T, +) { + lifecycleXML := []byte(`expirelogs/Enabled30`) + if _, err := globalBucketMetadataSys.Update(t.Context(), bucket, bucketLifecycleConfig, lifecycleXML); err != nil { + t.Fatal(err) + } + sseXML := []byte(`AES256`) + runBucketMetadataRMWConflict(t, obj, bucket, + func(ctx context.Context, objectAPI ObjectLayer) error { + _, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketLifecycleConfig) + return err + }, + func(ctx context.Context, objectAPI ObjectLayer) error { + _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketSSEConfig, sseXML) + return err + }, + func(meta BucketMetadata) bool { + cfg, err := lifecycle.ParseLifecycleConfig(bytes.NewReader(meta.LifecycleConfigXML)) + return err == nil && cfg.ExpiryUpdatedAt != nil && len(cfg.Rules) == 0 && bytes.Equal(meta.EncryptionConfigXML, sseXML) + }, instanceType+": lifecycle delete+SSE") +} + func TestMakeBucketForceCreatePreservesMetadata(t *testing.T) { defer DetectTestLeak(t)() ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ @@ -175,6 +246,122 @@ func TestApplyImportedBucketMetadataPreservesUnspecifiedFields(t *testing.T) { } } +func TestMakeBucketDoesNotAdoptGhostMetadata(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + objAPITest: testMakeBucketDoesNotAdoptGhostMetadata, + }) +} + +func testMakeBucketDoesNotAdoptGhostMetadata(obj ObjectLayer, instanceType, _ string, + _ http.Handler, _ auth.Credentials, t *testing.T, +) { + ctx := t.Context() + bucket := getRandomBucketName() + if err := obj.MakeBucket(ctx, bucket, MakeBucketOptions{}); err != nil { + t.Fatal(err) + } + policyJSON := fmt.Appendf(nil, `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket) + if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketPolicyConfig, policyJSON); err != nil { + t.Fatal(err) + } + oldMeta, err := readBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + z, ok := obj.(*erasureServerPools) + if !ok { + t.Fatalf("%s: object layer is %T, want *erasureServerPools", instanceType, obj) + } + if err = z.s3Peer.DeleteBucket(ctx, bucket, DeleteBucketOptions{Force: true}); err != nil { + t.Fatalf("%s: delete bucket volume only: %v", instanceType, err) + } + if err = obj.MakeBucket(ctx, bucket, MakeBucketOptions{}); err != nil { + t.Fatalf("%s: recreate bucket: %v", instanceType, err) + } + newMeta, err := readBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + if bytes.Equal(newMeta.PolicyConfigJSON, policyJSON) || newMeta.Created.Equal(oldMeta.Created) { + t.Fatalf("%s: new bucket adopted ghost metadata: old=%+v new=%+v", instanceType, oldMeta, newMeta) + } +} + +func TestMakeBucketForceCreateLockEnablesVersioning(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + objAPITest: testMakeBucketForceCreateLockEnablesVersioning, + }) +} + +func TestPeerBucketMetadataSaveSurvivesCallerCancellation(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + objAPITest: testPeerBucketMetadataSaveSurvivesCallerCancellation, + }) +} + +func testPeerBucketMetadataSaveSurvivesCallerCancellation(obj ObjectLayer, instanceType, bucket string, + _ http.Handler, _ auth.Credentials, t *testing.T, +) { + previousObjectAPI := newObjectLayerFn() + ctx, cancel := context.WithCancel(t.Context()) + barrier := &metadataRMWBarrierObjectLayer{ + ObjectLayer: obj, + bucket: bucket, + cancelOnPut: cancel, + } + setObjectLayer(barrier) + defer setObjectLayer(previousObjectAPI) + + err := globalSiteReplicationSys.PeerBucketMakeWithVersioningHandler(ctx, bucket, MakeBucketOptions{VersioningEnabled: true}) + if err != nil { + t.Fatalf("%s: peer metadata save failed after caller cancellation: %v", instanceType, err) + } + if ctx.Err() != context.Canceled { + t.Fatalf("%s: metadata write did not trigger caller cancellation", instanceType) + } + meta, err := readBucketMetadata(t.Context(), obj, bucket) + if err != nil { + t.Fatal(err) + } + cfg, err := versioning.ParseConfig(bytes.NewReader(meta.VersioningConfigXML)) + if err != nil { + t.Fatal(err) + } + if !cfg.Enabled() { + t.Fatalf("%s: peer metadata save lost versioning after cancellation", instanceType) + } +} + +func testMakeBucketForceCreateLockEnablesVersioning(obj ObjectLayer, instanceType, bucket string, + _ http.Handler, _ auth.Credentials, t *testing.T, +) { + ctx := t.Context() + suspended := []byte(`Suspended`) + if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, suspended); err != nil { + t.Fatal(err) + } + if err := obj.MakeBucket(ctx, bucket, MakeBucketOptions{ForceCreate: true, LockEnabled: true}); err != nil { + t.Fatalf("%s: ForceCreate with object lock: %v", instanceType, err) + } + meta, err := readBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + cfg, err := versioning.ParseConfig(bytes.NewReader(meta.VersioningConfigXML)) + if err != nil { + t.Fatal(err) + } + if !cfg.Enabled() || len(meta.ObjectLockConfigXML) == 0 { + t.Fatalf("%s: object lock state lacks enabled versioning: metadata=%+v", instanceType, meta) + } +} + func testBucketMetadataLockPreservesTaggingAndSSE(obj ObjectLayer, instanceType, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T, ) { diff --git a/cmd/erasure-server-pool.go b/cmd/erasure-server-pool.go index 262a5aae3..1e85cf039 100644 --- a/cmd/erasure-server-pool.go +++ b/cmd/erasure-server-pool.go @@ -898,30 +898,33 @@ func (z *erasureServerPools) MakeBucket(ctx context.Context, bucket string, opts } err = func() error { defer unlock() - meta, err := readBucketMetadata(ctx, z, bucket) - if errors.Is(err, errConfigNotFound) { - meta = newBucketMetadata(bucket) - } else if err != nil { - return err + meta := newBucketMetadata(bucket) + if opts.ForceCreate { + existing, err := loadBucketMetadataParse(ctx, z, bucket, true) + if err == nil { + meta = existing + } else if !errors.Is(err, errConfigNotFound) { + return err + } } if meta.Created.IsZero() { meta.SetCreatedAt(opts.CreatedAt) } if opts.LockEnabled { - if len(meta.VersioningConfigXML) == 0 { - meta.VersioningConfigXML = enabledBucketVersioningConfig - meta.VersioningConfigUpdatedAt = meta.Created + if err := enablePeerBucketVersioning(&meta); err != nil { + return err } if len(meta.ObjectLockConfigXML) == 0 { meta.ObjectLockConfigXML = enabledBucketObjectLockConfig meta.ObjectLockConfigUpdatedAt = meta.Created } } - if opts.VersioningEnabled && len(meta.VersioningConfigXML) == 0 { - meta.VersioningConfigXML = enabledBucketVersioningConfig - meta.VersioningConfigUpdatedAt = meta.Created + if opts.VersioningEnabled { + if err := enablePeerBucketVersioning(&meta); err != nil { + return err + } } - if err = meta.Save(ctx, z); err != nil { + if err = meta.Save(bgContext(ctx), z); err != nil { return err } globalBucketMetadataSys.Set(bucket, meta) diff --git a/cmd/site-replication.go b/cmd/site-replication.go index eadaef93a..386d321b7 100644 --- a/cmd/site-replication.go +++ b/cmd/site-replication.go @@ -952,7 +952,7 @@ func (c *SiteReplicationSys) PeerBucketMakeWithVersioningHandler(ctx context.Con meta.ObjectLockConfigUpdatedAt = meta.Created } } - return globalBucketMetadataSys.saveMetadata(ctx, objAPI, meta) + return globalBucketMetadataSys.saveMetadata(bgContext(ctx), objAPI, meta) }() if err != nil { return wrapSRErr(c.annotateErr(makeBucketWithVersion, err)) @@ -1620,6 +1620,7 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context return wrapSRErr(err) } } + notifyCtx := ctx ctx, unlock, err := lockBucketMetadata(ctx, objectAPI, item.Bucket) if err != nil { return wrapSRErr(err) @@ -1703,7 +1704,7 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context } unlock() locked = false - globalNotificationSys.LoadBucketMetadata(bgContext(ctx), item.Bucket) + globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), item.Bucket) return nil } @@ -1981,6 +1982,7 @@ func applyBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket return time.Time{}, err } + notifyCtx := ctx ctx, unlock, err := lockBucketMetadata(ctx, objectAPI, bucket) if err != nil { return time.Time{}, err @@ -2033,7 +2035,7 @@ func applyBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket } unlock() locked = false - globalNotificationSys.LoadBucketMetadata(bgContext(ctx), bucket) + globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), bucket) return updatedAt, nil }