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
}