mirror of
https://github.com/pgsty/minio.git
synced 2026-09-05 18:16:16 +03:00
Merge pull request #103 from pgsty/codex/issue-102-metadata-lock
fix: serialize whole-record bucket metadata updates
This commit is contained in:
+115
-12
@@ -589,6 +589,43 @@ type importMetaReport struct {
|
|||||||
madmin.BucketMetaImportErrs
|
madmin.BucketMetaImportErrs
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type importMetadataFields map[string]struct{}
|
||||||
|
|
||||||
|
func (f importMetadataFields) add(configFile string) {
|
||||||
|
f[configFile] = struct{}{}
|
||||||
|
}
|
||||||
|
|
||||||
|
func applyImportedBucketMetadata(dst *BucketMetadata, src BucketMetadata, fields importMetadataFields) {
|
||||||
|
for configFile := range fields {
|
||||||
|
switch configFile {
|
||||||
|
case bucketPolicyConfig:
|
||||||
|
dst.PolicyConfigJSON = bytes.Clone(src.PolicyConfigJSON)
|
||||||
|
dst.PolicyConfigUpdatedAt = src.PolicyConfigUpdatedAt
|
||||||
|
case bucketNotificationConfig:
|
||||||
|
dst.NotificationConfigXML = bytes.Clone(src.NotificationConfigXML)
|
||||||
|
dst.NotificationConfigUpdatedAt = src.NotificationConfigUpdatedAt
|
||||||
|
case bucketLifecycleConfig:
|
||||||
|
dst.LifecycleConfigXML = bytes.Clone(src.LifecycleConfigXML)
|
||||||
|
dst.LifecycleConfigUpdatedAt = src.LifecycleConfigUpdatedAt
|
||||||
|
case bucketSSEConfig:
|
||||||
|
dst.EncryptionConfigXML = bytes.Clone(src.EncryptionConfigXML)
|
||||||
|
dst.EncryptionConfigUpdatedAt = src.EncryptionConfigUpdatedAt
|
||||||
|
case bucketTaggingConfig:
|
||||||
|
dst.TaggingConfigXML = bytes.Clone(src.TaggingConfigXML)
|
||||||
|
dst.TaggingConfigUpdatedAt = src.TaggingConfigUpdatedAt
|
||||||
|
case bucketQuotaConfigFile:
|
||||||
|
dst.QuotaConfigJSON = bytes.Clone(src.QuotaConfigJSON)
|
||||||
|
dst.QuotaConfigUpdatedAt = src.QuotaConfigUpdatedAt
|
||||||
|
case objectLockConfig:
|
||||||
|
dst.ObjectLockConfigXML = bytes.Clone(src.ObjectLockConfigXML)
|
||||||
|
dst.ObjectLockConfigUpdatedAt = src.ObjectLockConfigUpdatedAt
|
||||||
|
case bucketVersioningConfig:
|
||||||
|
dst.VersioningConfigXML = bytes.Clone(src.VersioningConfigXML)
|
||||||
|
dst.VersioningConfigUpdatedAt = src.VersioningConfigUpdatedAt
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (i *importMetaReport) SetStatus(bucket, fname string, err error) {
|
func (i *importMetaReport) SetStatus(bucket, fname string, err error) {
|
||||||
st := i.Buckets[bucket]
|
st := i.Buckets[bucket]
|
||||||
var errMsg string
|
var errMsg string
|
||||||
@@ -649,6 +686,16 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
}
|
}
|
||||||
|
|
||||||
bucketMap := make(map[string]*BucketMetadata, len(zr.File))
|
bucketMap := make(map[string]*BucketMetadata, len(zr.File))
|
||||||
|
importedFields := make(map[string]importMetadataFields, len(zr.File))
|
||||||
|
blockedBuckets := make(map[string]struct{})
|
||||||
|
markImported := func(bucket, configFile string) {
|
||||||
|
fields := importedFields[bucket]
|
||||||
|
if fields == nil {
|
||||||
|
fields = make(importMetadataFields)
|
||||||
|
importedFields[bucket] = fields
|
||||||
|
}
|
||||||
|
fields.add(configFile)
|
||||||
|
}
|
||||||
|
|
||||||
updatedAt := UTCNow()
|
updatedAt := UTCNow()
|
||||||
|
|
||||||
@@ -664,6 +711,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
bucketMap[bucket] = &meta
|
bucketMap[bucket] = &meta
|
||||||
} else if err != errConfigNotFound {
|
} else if err != errConfigNotFound {
|
||||||
rpt.SetStatus(bucket, "", err)
|
rpt.SetStatus(bucket, "", err)
|
||||||
|
blockedBuckets[bucket] = struct{}{}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -675,6 +723,9 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
bucket, fileName := slc[0], slc[1]
|
bucket, fileName := slc[0], slc[1]
|
||||||
|
if _, blocked := blockedBuckets[bucket]; blocked {
|
||||||
|
continue
|
||||||
|
}
|
||||||
if fileName == objectLockConfig {
|
if fileName == objectLockConfig {
|
||||||
reader, err := file.Open()
|
reader, err := file.Open()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -708,6 +759,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
|
|
||||||
bucketMap[bucket].ObjectLockConfigXML = configData
|
bucketMap[bucket].ObjectLockConfigXML = configData
|
||||||
bucketMap[bucket].ObjectLockConfigUpdatedAt = updatedAt
|
bucketMap[bucket].ObjectLockConfigUpdatedAt = updatedAt
|
||||||
|
markImported(bucket, fileName)
|
||||||
rpt.SetStatus(bucket, fileName, nil)
|
rpt.SetStatus(bucket, fileName, nil)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -720,6 +772,9 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
bucket, fileName := slc[0], slc[1]
|
bucket, fileName := slc[0], slc[1]
|
||||||
|
if _, blocked := blockedBuckets[bucket]; blocked {
|
||||||
|
continue
|
||||||
|
}
|
||||||
if fileName == bucketVersioningConfig {
|
if fileName == bucketVersioningConfig {
|
||||||
reader, err := file.Open()
|
reader, err := file.Open()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -764,6 +819,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
|
|
||||||
bucketMap[bucket].VersioningConfigXML = configData
|
bucketMap[bucket].VersioningConfigXML = configData
|
||||||
bucketMap[bucket].VersioningConfigUpdatedAt = updatedAt
|
bucketMap[bucket].VersioningConfigUpdatedAt = updatedAt
|
||||||
|
markImported(bucket, fileName)
|
||||||
rpt.SetStatus(bucket, fileName, nil)
|
rpt.SetStatus(bucket, fileName, nil)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -781,6 +837,9 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
bucket, fileName := slc[0], slc[1]
|
bucket, fileName := slc[0], slc[1]
|
||||||
|
if _, blocked := blockedBuckets[bucket]; blocked {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
// create bucket if it does not exist yet.
|
// create bucket if it does not exist yet.
|
||||||
if _, ok := bucketMap[bucket]; !ok {
|
if _, ok := bucketMap[bucket]; !ok {
|
||||||
@@ -813,6 +872,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
|
|
||||||
bucketMap[bucket].NotificationConfigXML = configData
|
bucketMap[bucket].NotificationConfigXML = configData
|
||||||
bucketMap[bucket].NotificationConfigUpdatedAt = updatedAt
|
bucketMap[bucket].NotificationConfigUpdatedAt = updatedAt
|
||||||
|
markImported(bucket, fileName)
|
||||||
rpt.SetStatus(bucket, fileName, nil)
|
rpt.SetStatus(bucket, fileName, nil)
|
||||||
case bucketPolicyConfig:
|
case bucketPolicyConfig:
|
||||||
// Error out if Content-Length is beyond allowed size.
|
// Error out if Content-Length is beyond allowed size.
|
||||||
@@ -847,6 +907,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
|
|
||||||
bucketMap[bucket].PolicyConfigJSON = configData
|
bucketMap[bucket].PolicyConfigJSON = configData
|
||||||
bucketMap[bucket].PolicyConfigUpdatedAt = updatedAt
|
bucketMap[bucket].PolicyConfigUpdatedAt = updatedAt
|
||||||
|
markImported(bucket, fileName)
|
||||||
rpt.SetStatus(bucket, fileName, nil)
|
rpt.SetStatus(bucket, fileName, nil)
|
||||||
case bucketLifecycleConfig:
|
case bucketLifecycleConfig:
|
||||||
bucketLifecycle, err := lifecycle.ParseLifecycleConfig(io.LimitReader(reader, sz))
|
bucketLifecycle, err := lifecycle.ParseLifecycleConfig(io.LimitReader(reader, sz))
|
||||||
@@ -879,6 +940,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
|
|
||||||
bucketMap[bucket].LifecycleConfigXML = configData
|
bucketMap[bucket].LifecycleConfigXML = configData
|
||||||
bucketMap[bucket].LifecycleConfigUpdatedAt = updatedAt
|
bucketMap[bucket].LifecycleConfigUpdatedAt = updatedAt
|
||||||
|
markImported(bucket, fileName)
|
||||||
rpt.SetStatus(bucket, fileName, nil)
|
rpt.SetStatus(bucket, fileName, nil)
|
||||||
case bucketSSEConfig:
|
case bucketSSEConfig:
|
||||||
// Parse bucket encryption xml
|
// Parse bucket encryption xml
|
||||||
@@ -917,6 +979,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
|
|
||||||
bucketMap[bucket].EncryptionConfigXML = configData
|
bucketMap[bucket].EncryptionConfigXML = configData
|
||||||
bucketMap[bucket].EncryptionConfigUpdatedAt = updatedAt
|
bucketMap[bucket].EncryptionConfigUpdatedAt = updatedAt
|
||||||
|
markImported(bucket, fileName)
|
||||||
rpt.SetStatus(bucket, fileName, nil)
|
rpt.SetStatus(bucket, fileName, nil)
|
||||||
case bucketTaggingConfig:
|
case bucketTaggingConfig:
|
||||||
tags, err := tags.ParseBucketXML(io.LimitReader(reader, sz))
|
tags, err := tags.ParseBucketXML(io.LimitReader(reader, sz))
|
||||||
@@ -933,6 +996,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
|
|
||||||
bucketMap[bucket].TaggingConfigXML = configData
|
bucketMap[bucket].TaggingConfigXML = configData
|
||||||
bucketMap[bucket].TaggingConfigUpdatedAt = updatedAt
|
bucketMap[bucket].TaggingConfigUpdatedAt = updatedAt
|
||||||
|
markImported(bucket, fileName)
|
||||||
rpt.SetStatus(bucket, fileName, nil)
|
rpt.SetStatus(bucket, fileName, nil)
|
||||||
case bucketQuotaConfigFile:
|
case bucketQuotaConfigFile:
|
||||||
data, err := io.ReadAll(reader)
|
data, err := io.ReadAll(reader)
|
||||||
@@ -949,6 +1013,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
|
|
||||||
bucketMap[bucket].QuotaConfigJSON = data
|
bucketMap[bucket].QuotaConfigJSON = data
|
||||||
bucketMap[bucket].QuotaConfigUpdatedAt = updatedAt
|
bucketMap[bucket].QuotaConfigUpdatedAt = updatedAt
|
||||||
|
markImported(bucket, fileName)
|
||||||
rpt.SetStatus(bucket, fileName, nil)
|
rpt.SetStatus(bucket, fileName, nil)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -962,22 +1027,60 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
|
|||||||
}
|
}
|
||||||
|
|
||||||
for bucket, meta := range bucketMap {
|
for bucket, meta := range bucketMap {
|
||||||
err := globalBucketMetadataSys.save(ctx, *meta)
|
fields := importedFields[bucket]
|
||||||
|
if len(fields) == 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
var merged BucketMetadata
|
||||||
|
err := func() error {
|
||||||
|
lockCtx, unlock, err := lockBucketMetadata(ctx, objectAPI, bucket)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
defer unlock()
|
||||||
|
merged, err = loadBucketMetadataParse(lockCtx, objectAPI, bucket, true)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
applyImportedBucketMetadata(&merged, *meta, fields)
|
||||||
|
return globalBucketMetadataSys.saveMetadata(lockCtx, objectAPI, merged)
|
||||||
|
}()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
rpt.SetStatus(bucket, "", err)
|
rpt.SetStatus(bucket, "", err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
// Call site replication hook.
|
*meta = merged
|
||||||
if err = globalSiteReplicationSys.BucketMetaHook(ctx, madmin.SRBucketMeta{
|
globalNotificationSys.LoadBucketMetadata(bgContext(ctx), bucket)
|
||||||
Bucket: bucket,
|
hook := madmin.SRBucketMeta{Bucket: bucket, UpdatedAt: updatedAt}
|
||||||
Quota: meta.QuotaConfigJSON,
|
var hookNeeded bool
|
||||||
Policy: meta.PolicyConfigJSON,
|
if _, ok := fields[bucketQuotaConfigFile]; ok {
|
||||||
Versioning: enc(meta.VersioningConfigXML),
|
hook.Quota = meta.QuotaConfigJSON
|
||||||
Tags: enc(meta.TaggingConfigXML),
|
hookNeeded = true
|
||||||
ObjectLockConfig: enc(meta.ObjectLockConfigXML),
|
}
|
||||||
SSEConfig: enc(meta.EncryptionConfigXML),
|
if _, ok := fields[bucketPolicyConfig]; ok {
|
||||||
UpdatedAt: updatedAt,
|
hook.Policy = meta.PolicyConfigJSON
|
||||||
}); err != nil {
|
hookNeeded = true
|
||||||
|
}
|
||||||
|
if _, ok := fields[bucketVersioningConfig]; ok {
|
||||||
|
hook.Versioning = enc(meta.VersioningConfigXML)
|
||||||
|
hookNeeded = true
|
||||||
|
}
|
||||||
|
if _, ok := fields[bucketTaggingConfig]; ok {
|
||||||
|
hook.Tags = enc(meta.TaggingConfigXML)
|
||||||
|
hookNeeded = true
|
||||||
|
}
|
||||||
|
if _, ok := fields[objectLockConfig]; ok {
|
||||||
|
hook.ObjectLockConfig = enc(meta.ObjectLockConfigXML)
|
||||||
|
hookNeeded = true
|
||||||
|
}
|
||||||
|
if _, ok := fields[bucketSSEConfig]; ok {
|
||||||
|
hook.SSEConfig = enc(meta.EncryptionConfigXML)
|
||||||
|
hookNeeded = true
|
||||||
|
}
|
||||||
|
if hookNeeded {
|
||||||
|
err = globalSiteReplicationSys.BucketMetaHook(ctx, hook)
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
rpt.SetStatus(bucket, "", err)
|
rpt.SetStatus(bucket, "", err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,455 @@
|
|||||||
|
// Copyright (c) 2015-2026 MinIO, Inc.
|
||||||
|
// Copyright (c) 2026 PGSTY
|
||||||
|
//
|
||||||
|
// This file is part of MinIO Object Storage stack
|
||||||
|
//
|
||||||
|
// This program is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Affero General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
|
||||||
|
package cmd
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"net/http"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
|
"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{}
|
||||||
|
|
||||||
|
type metadataRMWBarrierObjectLayer struct {
|
||||||
|
ObjectLayer
|
||||||
|
bucket string
|
||||||
|
aReady chan struct{}
|
||||||
|
aRelease chan struct{}
|
||||||
|
bLockAttempt chan struct{}
|
||||||
|
aReadyOnce sync.Once
|
||||||
|
bLockOnce sync.Once
|
||||||
|
cancelOnce sync.Once
|
||||||
|
cancelOnPut context.CancelFunc
|
||||||
|
reads atomic.Int64
|
||||||
|
}
|
||||||
|
|
||||||
|
func (o *metadataRMWBarrierObjectLayer) metadataObject() string {
|
||||||
|
return pathJoin(bucketMetaPrefix, o.bucket, bucketMetadataFile)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (o *metadataRMWBarrierObjectLayer) metadataLock() string {
|
||||||
|
return pathJoin(bucketMetaPrefix, o.bucket, "metadata.lock")
|
||||||
|
}
|
||||||
|
|
||||||
|
func (o *metadataRMWBarrierObjectLayer) GetObjectNInfo(ctx context.Context, bucket, object string, rs *HTTPRangeSpec, h http.Header, opts ObjectOptions) (*GetObjectReader, error) {
|
||||||
|
if bucket == minioMetaBucket && object == o.metadataObject() {
|
||||||
|
o.reads.Add(1)
|
||||||
|
}
|
||||||
|
return o.ObjectLayer.GetObjectNInfo(ctx, bucket, object, rs, h, opts)
|
||||||
|
}
|
||||||
|
|
||||||
|
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 {
|
||||||
|
case <-o.aRelease:
|
||||||
|
case <-ctx.Done():
|
||||||
|
return ObjectInfo{}, ctx.Err()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return o.ObjectLayer.PutObject(ctx, bucket, object, data, opts)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (o *metadataRMWBarrierObjectLayer) NewNSLock(bucket string, objects ...string) RWLocker {
|
||||||
|
lock := o.ObjectLayer.NewNSLock(bucket, objects...)
|
||||||
|
if bucket != minioMetaBucket || len(objects) != 1 || objects[0] != o.metadataLock() {
|
||||||
|
return lock
|
||||||
|
}
|
||||||
|
return metadataObservedRWLocker{RWLocker: lock, onLock: func(ctx context.Context) {
|
||||||
|
if ctx.Value(metadataRMWWriterKey{}) == "B" {
|
||||||
|
o.bLockOnce.Do(func() { close(o.bLockAttempt) })
|
||||||
|
}
|
||||||
|
}}
|
||||||
|
}
|
||||||
|
|
||||||
|
type metadataObservedRWLocker struct {
|
||||||
|
RWLocker
|
||||||
|
onLock func(context.Context)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l metadataObservedRWLocker) GetLock(ctx context.Context, timeout *dynamicTimeout) (LockContext, error) {
|
||||||
|
l.onLock(ctx)
|
||||||
|
return l.RWLocker.GetLock(ctx, timeout)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestBucketMetadataLockPreservesPolicyAndCORS(t *testing.T) {
|
||||||
|
defer DetectTestLeak(t)()
|
||||||
|
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
|
||||||
|
t: t,
|
||||||
|
objAPITest: testBucketMetadataLockPreservesPolicyAndCORS,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func testBucketMetadataLockPreservesPolicyAndCORS(obj ObjectLayer, instanceType, bucket string,
|
||||||
|
_ http.Handler, _ auth.Credentials, t *testing.T,
|
||||||
|
) {
|
||||||
|
policyJSON := fmt.Appendf(nil, `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket)
|
||||||
|
corsXML := []byte(testSiteReplicationCORSDoc)
|
||||||
|
runBucketMetadataRMWConflict(t, obj, bucket,
|
||||||
|
func(ctx context.Context, objectAPI ObjectLayer) error {
|
||||||
|
_, err := globalBucketMetadataSys.Update(ctx, bucket, bucketPolicyConfig, policyJSON)
|
||||||
|
return err
|
||||||
|
},
|
||||||
|
func(ctx context.Context, objectAPI ObjectLayer) error {
|
||||||
|
_, err := updateLocalBucketCORSMetadata(ctx, objectAPI, bucket, corsXML)
|
||||||
|
return err
|
||||||
|
},
|
||||||
|
func(meta BucketMetadata) bool {
|
||||||
|
return bytes.Equal(meta.PolicyConfigJSON, policyJSON) && bytes.Equal(meta.CorsConfigXML, corsXML)
|
||||||
|
}, instanceType+": policy+CORS")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestBucketMetadataLockPreservesTaggingAndSSE(t *testing.T) {
|
||||||
|
defer DetectTestLeak(t)()
|
||||||
|
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
|
||||||
|
t: t,
|
||||||
|
objAPITest: testBucketMetadataLockPreservesTaggingAndSSE,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
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(`<Tagging><TagSet><Tag><Key>local</Key><Value>tag</Value></Tag></TagSet></Tagging>`)
|
||||||
|
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(`<LifecycleConfiguration><Rule><ID>expire</ID><Filter><Prefix>logs/</Prefix></Filter><Status>Enabled</Status><Expiration><Days>30</Days></Expiration></Rule></LifecycleConfiguration>`)
|
||||||
|
if _, err := globalBucketMetadataSys.Update(t.Context(), bucket, bucketLifecycleConfig, lifecycleXML); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
sseXML := []byte(`<ServerSideEncryptionConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ApplyServerSideEncryptionByDefault><SSEAlgorithm>AES256</SSEAlgorithm></ApplyServerSideEncryptionByDefault></Rule></ServerSideEncryptionConfiguration>`)
|
||||||
|
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{
|
||||||
|
t: t,
|
||||||
|
objAPITest: testMakeBucketForceCreatePreservesMetadata,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func testMakeBucketForceCreatePreservesMetadata(obj ObjectLayer, instanceType, bucket string,
|
||||||
|
_ http.Handler, _ auth.Credentials, t *testing.T,
|
||||||
|
) {
|
||||||
|
ctx := t.Context()
|
||||||
|
policyJSON := fmt.Appendf(nil, `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket)
|
||||||
|
corsXML := []byte(testSiteReplicationCORSDoc)
|
||||||
|
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketPolicyConfig, policyJSON); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := updateLocalBucketCORSMetadata(ctx, obj, bucket, corsXML); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
before, err := readBucketMetadata(ctx, obj, bucket)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err = obj.MakeBucket(ctx, bucket, MakeBucketOptions{ForceCreate: true}); err != nil {
|
||||||
|
t.Fatalf("%s: ForceCreate existing bucket: %v", instanceType, err)
|
||||||
|
}
|
||||||
|
after, err := readBucketMetadata(ctx, obj, bucket)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if !after.Created.Equal(before.Created) || !bytes.Equal(after.PolicyConfigJSON, policyJSON) || !bytes.Equal(after.CorsConfigXML, corsXML) {
|
||||||
|
t.Fatalf("%s: ForceCreate replaced metadata: before=%+v after=%+v", instanceType, before, after)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestApplyImportedBucketMetadataPreservesUnspecifiedFields(t *testing.T) {
|
||||||
|
policyJSON := []byte(`{"Version":"2012-10-17","Statement":[]}`)
|
||||||
|
tagXML := []byte(`<Tagging><TagSet><Tag><Key>existing</Key><Value>tag</Value></Tag></TagSet></Tagging>`)
|
||||||
|
src := newBucketMetadata("bucket")
|
||||||
|
src.PolicyConfigJSON = policyJSON
|
||||||
|
src.PolicyConfigUpdatedAt = UTCNow()
|
||||||
|
dst := newBucketMetadata("bucket")
|
||||||
|
dst.TaggingConfigXML = bytes.Clone(tagXML)
|
||||||
|
|
||||||
|
applyImportedBucketMetadata(&dst, src, importMetadataFields{bucketPolicyConfig: {}})
|
||||||
|
if !bytes.Equal(dst.PolicyConfigJSON, policyJSON) || !bytes.Equal(dst.TaggingConfigXML, tagXML) {
|
||||||
|
t.Fatalf("import patch overwrote unspecified metadata: %+v", dst)
|
||||||
|
}
|
||||||
|
src.PolicyConfigJSON[0] = '!'
|
||||||
|
if dst.PolicyConfigJSON[0] == '!' {
|
||||||
|
t.Fatal("import patch retained the source byte slice")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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(`<VersioningConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Status>Suspended</Status></VersioningConfiguration>`)
|
||||||
|
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,
|
||||||
|
) {
|
||||||
|
tagXML := []byte(`<Tagging><TagSet><Tag><Key>key</Key><Value>value</Value></Tag></TagSet></Tagging>`)
|
||||||
|
sseXML := []byte(`<ServerSideEncryptionConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ApplyServerSideEncryptionByDefault><SSEAlgorithm>AES256</SSEAlgorithm></ApplyServerSideEncryptionByDefault></Rule></ServerSideEncryptionConfiguration>`)
|
||||||
|
runBucketMetadataRMWConflict(t, obj, bucket,
|
||||||
|
func(ctx context.Context, objectAPI ObjectLayer) error {
|
||||||
|
_, err := globalBucketMetadataSys.Update(ctx, bucket, bucketTaggingConfig, tagXML)
|
||||||
|
return err
|
||||||
|
},
|
||||||
|
func(ctx context.Context, objectAPI ObjectLayer) error {
|
||||||
|
_, err := globalBucketMetadataSys.Update(ctx, bucket, bucketSSEConfig, sseXML)
|
||||||
|
return err
|
||||||
|
},
|
||||||
|
func(meta BucketMetadata) bool {
|
||||||
|
return bytes.Equal(meta.TaggingConfigXML, tagXML) && bytes.Equal(meta.EncryptionConfigXML, sseXML)
|
||||||
|
}, instanceType+": tagging+SSE")
|
||||||
|
}
|
||||||
|
|
||||||
|
func runBucketMetadataRMWConflict(t *testing.T, obj ObjectLayer, bucket string,
|
||||||
|
writerA, writerB func(context.Context, ObjectLayer) error,
|
||||||
|
complete func(BucketMetadata) bool, name string,
|
||||||
|
) {
|
||||||
|
t.Helper()
|
||||||
|
previousObjectAPI := newObjectLayerFn()
|
||||||
|
barrier := &metadataRMWBarrierObjectLayer{
|
||||||
|
ObjectLayer: obj,
|
||||||
|
bucket: bucket,
|
||||||
|
aReady: make(chan struct{}),
|
||||||
|
aRelease: make(chan struct{}),
|
||||||
|
bLockAttempt: make(chan struct{}),
|
||||||
|
}
|
||||||
|
setObjectLayer(barrier)
|
||||||
|
defer setObjectLayer(previousObjectAPI)
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
aCtx := context.WithValue(ctx, metadataRMWWriterKey{}, "A")
|
||||||
|
bCtx := context.WithValue(ctx, metadataRMWWriterKey{}, "B")
|
||||||
|
aDone := make(chan error, 1)
|
||||||
|
bDone := make(chan error, 1)
|
||||||
|
go func() { aDone <- writerA(aCtx, barrier) }()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-barrier.aReady:
|
||||||
|
case <-ctx.Done():
|
||||||
|
t.Fatalf("%s: writer A did not reach metadata save: %v", name, ctx.Err())
|
||||||
|
}
|
||||||
|
go func() { bDone <- writerB(bCtx, barrier) }()
|
||||||
|
|
||||||
|
var (
|
||||||
|
bErr error
|
||||||
|
bFinished bool
|
||||||
|
)
|
||||||
|
select {
|
||||||
|
case <-barrier.bLockAttempt:
|
||||||
|
if got := barrier.reads.Load(); got != 1 {
|
||||||
|
t.Fatalf("%s: writer B read metadata before acquiring metadata.lock: reads=%d", name, got)
|
||||||
|
}
|
||||||
|
case bErr = <-bDone:
|
||||||
|
bFinished = true
|
||||||
|
case <-ctx.Done():
|
||||||
|
t.Fatalf("%s: writer B neither completed nor attempted metadata.lock: %v", name, ctx.Err())
|
||||||
|
}
|
||||||
|
close(barrier.aRelease)
|
||||||
|
if err := <-aDone; err != nil {
|
||||||
|
t.Fatalf("%s: writer A failed: %v", name, err)
|
||||||
|
}
|
||||||
|
if !bFinished {
|
||||||
|
select {
|
||||||
|
case bErr = <-bDone:
|
||||||
|
case <-ctx.Done():
|
||||||
|
t.Fatalf("%s: writer B did not finish: %v", name, ctx.Err())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if bErr != nil {
|
||||||
|
t.Fatalf("%s: writer B failed: %v", name, bErr)
|
||||||
|
}
|
||||||
|
|
||||||
|
disk, err := readBucketMetadata(ctx, barrier, bucket)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("%s: read disk metadata: %v", name, err)
|
||||||
|
}
|
||||||
|
resident, err := globalBucketMetadataSys.Get(bucket)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("%s: read resident metadata: %v", name, err)
|
||||||
|
}
|
||||||
|
if !complete(disk) || !complete(resident) {
|
||||||
|
t.Fatalf("%s: concurrent updates lost a field: disk=%+v resident=%+v", name, disk, resident)
|
||||||
|
}
|
||||||
|
}
|
||||||
+126
-89
@@ -114,7 +114,7 @@ func (sys *BucketMetadataSys) Set(bucket string, meta BucketMetadata) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string, configFile string, configData []byte, parse bool) (updatedAt time.Time, err error) {
|
func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string, configFile string, configData []byte, parse, lifecycleDelete bool) (updatedAt time.Time, err error) {
|
||||||
objAPI := newObjectLayerFn()
|
objAPI := newObjectLayerFn()
|
||||||
if objAPI == nil {
|
if objAPI == nil {
|
||||||
return updatedAt, errServerNotInitialized
|
return updatedAt, errServerNotInitialized
|
||||||
@@ -123,63 +123,81 @@ func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string,
|
|||||||
if isMinioMetaBucketName(bucket) {
|
if isMinioMetaBucketName(bucket) {
|
||||||
return updatedAt, errInvalidArgument
|
return updatedAt, errInvalidArgument
|
||||||
}
|
}
|
||||||
|
notifyCtx := ctx
|
||||||
meta, err := loadBucketMetadataParse(ctx, objAPI, bucket, parse)
|
ctx, unlock, err := lockBucketMetadata(ctx, objAPI, bucket)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if !globalIsErasure && !globalIsDistErasure && errors.Is(err, errVolumeNotFound) {
|
return updatedAt, err
|
||||||
// Only single drive mode needs this fallback.
|
|
||||||
meta = newBucketMetadata(bucket)
|
|
||||||
} else {
|
|
||||||
return updatedAt, err
|
|
||||||
}
|
|
||||||
}
|
|
||||||
updatedAt = UTCNow()
|
|
||||||
switch configFile {
|
|
||||||
case bucketPolicyConfig:
|
|
||||||
meta.PolicyConfigJSON = configData
|
|
||||||
meta.PolicyConfigUpdatedAt = updatedAt
|
|
||||||
case bucketNotificationConfig:
|
|
||||||
meta.NotificationConfigXML = configData
|
|
||||||
meta.NotificationConfigUpdatedAt = updatedAt
|
|
||||||
case bucketLifecycleConfig:
|
|
||||||
meta.LifecycleConfigXML = configData
|
|
||||||
meta.LifecycleConfigUpdatedAt = updatedAt
|
|
||||||
case bucketSSEConfig:
|
|
||||||
meta.EncryptionConfigXML = configData
|
|
||||||
meta.EncryptionConfigUpdatedAt = updatedAt
|
|
||||||
case bucketTaggingConfig:
|
|
||||||
meta.TaggingConfigXML = configData
|
|
||||||
meta.TaggingConfigUpdatedAt = updatedAt
|
|
||||||
case bucketCorsConfig:
|
|
||||||
meta.CorsConfigXML = configData
|
|
||||||
meta.CorsConfigUpdatedAt = updatedAt
|
|
||||||
case bucketQuotaConfigFile:
|
|
||||||
meta.QuotaConfigJSON = configData
|
|
||||||
meta.QuotaConfigUpdatedAt = updatedAt
|
|
||||||
case objectLockConfig:
|
|
||||||
meta.ObjectLockConfigXML = configData
|
|
||||||
meta.ObjectLockConfigUpdatedAt = updatedAt
|
|
||||||
case bucketVersioningConfig:
|
|
||||||
meta.VersioningConfigXML = configData
|
|
||||||
meta.VersioningConfigUpdatedAt = updatedAt
|
|
||||||
case bucketReplicationConfig:
|
|
||||||
meta.ReplicationConfigXML = configData
|
|
||||||
meta.ReplicationConfigUpdatedAt = updatedAt
|
|
||||||
case bucketTargetsFile:
|
|
||||||
meta.BucketTargetsConfigJSON, meta.BucketTargetsConfigMetaJSON, err = encryptBucketMetadata(ctx, meta.Name, configData, kms.Context{
|
|
||||||
bucket: meta.Name,
|
|
||||||
bucketTargetsFile: bucketTargetsFile,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return updatedAt, fmt.Errorf("Error encrypting bucket target metadata %w", err)
|
|
||||||
}
|
|
||||||
meta.BucketTargetsConfigUpdatedAt = updatedAt
|
|
||||||
meta.BucketTargetsConfigMetaUpdatedAt = updatedAt
|
|
||||||
default:
|
|
||||||
return updatedAt, fmt.Errorf("Unknown bucket %s metadata update requested %s", bucket, configFile)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return updatedAt, sys.save(ctx, meta)
|
err = func() error {
|
||||||
|
defer unlock()
|
||||||
|
meta, err := loadBucketMetadataParse(ctx, objAPI, bucket, parse)
|
||||||
|
if err != nil {
|
||||||
|
if !globalIsErasure && !globalIsDistErasure && errors.Is(err, errVolumeNotFound) {
|
||||||
|
// Only single drive mode needs this fallback.
|
||||||
|
meta = newBucketMetadata(bucket)
|
||||||
|
} else {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if lifecycleDelete {
|
||||||
|
configData, err = lifecycleDeleteConfig(meta.LifecycleConfigXML)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
updatedAt = UTCNow()
|
||||||
|
switch configFile {
|
||||||
|
case bucketPolicyConfig:
|
||||||
|
meta.PolicyConfigJSON = configData
|
||||||
|
meta.PolicyConfigUpdatedAt = updatedAt
|
||||||
|
case bucketNotificationConfig:
|
||||||
|
meta.NotificationConfigXML = configData
|
||||||
|
meta.NotificationConfigUpdatedAt = updatedAt
|
||||||
|
case bucketLifecycleConfig:
|
||||||
|
meta.LifecycleConfigXML = configData
|
||||||
|
meta.LifecycleConfigUpdatedAt = updatedAt
|
||||||
|
case bucketSSEConfig:
|
||||||
|
meta.EncryptionConfigXML = configData
|
||||||
|
meta.EncryptionConfigUpdatedAt = updatedAt
|
||||||
|
case bucketTaggingConfig:
|
||||||
|
meta.TaggingConfigXML = configData
|
||||||
|
meta.TaggingConfigUpdatedAt = updatedAt
|
||||||
|
case bucketCorsConfig:
|
||||||
|
meta.CorsConfigXML = configData
|
||||||
|
meta.CorsConfigUpdatedAt = updatedAt
|
||||||
|
case bucketQuotaConfigFile:
|
||||||
|
meta.QuotaConfigJSON = configData
|
||||||
|
meta.QuotaConfigUpdatedAt = updatedAt
|
||||||
|
case objectLockConfig:
|
||||||
|
meta.ObjectLockConfigXML = configData
|
||||||
|
meta.ObjectLockConfigUpdatedAt = updatedAt
|
||||||
|
case bucketVersioningConfig:
|
||||||
|
meta.VersioningConfigXML = configData
|
||||||
|
meta.VersioningConfigUpdatedAt = updatedAt
|
||||||
|
case bucketReplicationConfig:
|
||||||
|
meta.ReplicationConfigXML = configData
|
||||||
|
meta.ReplicationConfigUpdatedAt = updatedAt
|
||||||
|
case bucketTargetsFile:
|
||||||
|
meta.BucketTargetsConfigJSON, meta.BucketTargetsConfigMetaJSON, err = encryptBucketMetadata(ctx, meta.Name, configData, kms.Context{
|
||||||
|
bucket: meta.Name,
|
||||||
|
bucketTargetsFile: bucketTargetsFile,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("Error encrypting bucket target metadata %w", err)
|
||||||
|
}
|
||||||
|
meta.BucketTargetsConfigUpdatedAt = updatedAt
|
||||||
|
meta.BucketTargetsConfigMetaUpdatedAt = updatedAt
|
||||||
|
default:
|
||||||
|
return fmt.Errorf("Unknown bucket %s metadata update requested %s", bucket, configFile)
|
||||||
|
}
|
||||||
|
return sys.saveMetadata(ctx, objAPI, meta)
|
||||||
|
}()
|
||||||
|
if err != nil {
|
||||||
|
return updatedAt, err
|
||||||
|
}
|
||||||
|
globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), bucket) // Do not use caller context here
|
||||||
|
return updatedAt, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) error {
|
func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) error {
|
||||||
@@ -192,59 +210,78 @@ func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) err
|
|||||||
return errInvalidArgument
|
return errInvalidArgument
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := meta.Save(ctx, objAPI); err != nil {
|
if err := sys.saveMetadata(ctx, objAPI, meta); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
sys.Set(meta.Name, meta)
|
|
||||||
globalNotificationSys.LoadBucketMetadata(bgContext(ctx), meta.Name) // Do not use caller context here
|
globalNotificationSys.LoadBucketMetadata(bgContext(ctx), meta.Name) // Do not use caller context here
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// saveMetadata persists and publishes metadata locally. Callers performing a
|
||||||
|
// read-modify-write must hold metadata.lock and release it before peer fan-out.
|
||||||
|
func (sys *BucketMetadataSys) saveMetadata(ctx context.Context, objAPI ObjectLayer, meta BucketMetadata) error {
|
||||||
|
if err := meta.Save(ctx, objAPI); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
sys.Set(meta.Name, meta)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func lockBucketMetadata(ctx context.Context, objectAPI ObjectLayer, bucket string) (context.Context, func(), error) {
|
||||||
|
return lockBucketMetadataWithTimeout(ctx, objectAPI, bucket, globalOperationTimeout)
|
||||||
|
}
|
||||||
|
|
||||||
|
func lockBucketMetadataWithTimeout(ctx context.Context, objectAPI ObjectLayer, bucket string, timeout *dynamicTimeout) (context.Context, func(), error) {
|
||||||
|
lock := objectAPI.NewNSLock(minioMetaBucket, pathJoin(bucketMetaPrefix, bucket, "metadata.lock"))
|
||||||
|
lkctx, err := lock.GetLock(ctx, timeout)
|
||||||
|
if err != nil {
|
||||||
|
return nil, nil, err
|
||||||
|
}
|
||||||
|
ctx = context.WithValue(lkctx.Context(), bucketMetadataLockContextKey{}, bucket)
|
||||||
|
return ctx, func() { lock.Unlock(lkctx) }, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
type bucketMetadataLockContextKey struct{}
|
||||||
|
|
||||||
|
func bucketMetadataLockHeld(ctx context.Context, bucket string) bool {
|
||||||
|
lockedBucket, _ := ctx.Value(bucketMetadataLockContextKey{}).(string)
|
||||||
|
return lockedBucket == bucket
|
||||||
|
}
|
||||||
|
|
||||||
// Delete delete the bucket metadata for the specified bucket.
|
// Delete delete the bucket metadata for the specified bucket.
|
||||||
// must be used by all callers instead of using Update() with nil configData.
|
// must be used by all callers instead of using Update() with nil configData.
|
||||||
func (sys *BucketMetadataSys) Delete(ctx context.Context, bucket string, configFile string) (updatedAt time.Time, err error) {
|
func (sys *BucketMetadataSys) Delete(ctx context.Context, bucket string, configFile string) (updatedAt time.Time, err error) {
|
||||||
if configFile == bucketLifecycleConfig {
|
return sys.updateAndParse(ctx, bucket, configFile, nil, false, configFile == bucketLifecycleConfig)
|
||||||
// Get bucket config from current site
|
}
|
||||||
meta, e := globalBucketMetadataSys.GetConfigFromDisk(ctx, bucket)
|
|
||||||
if e != nil && !errors.Is(e, errConfigNotFound) {
|
|
||||||
return updatedAt, e
|
|
||||||
}
|
|
||||||
var expiryRuleRemoved bool
|
|
||||||
if len(meta.LifecycleConfigXML) > 0 {
|
|
||||||
var lcCfg lifecycle.Lifecycle
|
|
||||||
if err := xml.Unmarshal(meta.LifecycleConfigXML, &lcCfg); err != nil {
|
|
||||||
return updatedAt, err
|
|
||||||
}
|
|
||||||
// find a single expiry rule set the flag
|
|
||||||
for _, rl := range lcCfg.Rules {
|
|
||||||
if !rl.Expiration.IsNull() || !rl.NoncurrentVersionExpiration.IsNull() {
|
|
||||||
expiryRuleRemoved = true
|
|
||||||
break
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Form empty ILM details with `ExpiryUpdatedAt` field and save
|
func lifecycleDeleteConfig(current []byte) ([]byte, error) {
|
||||||
var cfgData []byte
|
var expiryRuleRemoved bool
|
||||||
if expiryRuleRemoved {
|
if len(current) > 0 {
|
||||||
var lcCfg lifecycle.Lifecycle
|
var lcCfg lifecycle.Lifecycle
|
||||||
currtime := time.Now()
|
if err := xml.Unmarshal(current, &lcCfg); err != nil {
|
||||||
lcCfg.ExpiryUpdatedAt = &currtime
|
return nil, err
|
||||||
cfgData, err = xml.Marshal(lcCfg)
|
}
|
||||||
if err != nil {
|
for _, rl := range lcCfg.Rules {
|
||||||
return updatedAt, err
|
if !rl.Expiration.IsNull() || !rl.NoncurrentVersionExpiration.IsNull() {
|
||||||
|
expiryRuleRemoved = true
|
||||||
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return sys.updateAndParse(ctx, bucket, configFile, cfgData, false)
|
|
||||||
}
|
}
|
||||||
return sys.updateAndParse(ctx, bucket, configFile, nil, false)
|
if !expiryRuleRemoved {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
var lcCfg lifecycle.Lifecycle
|
||||||
|
currtime := time.Now()
|
||||||
|
lcCfg.ExpiryUpdatedAt = &currtime
|
||||||
|
return xml.Marshal(lcCfg)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Update update bucket metadata for the specified bucket.
|
// Update update bucket metadata for the specified bucket.
|
||||||
// The configData data should not be modified after being sent here.
|
// The configData data should not be modified after being sent here.
|
||||||
func (sys *BucketMetadataSys) Update(ctx context.Context, bucket string, configFile string, configData []byte) (updatedAt time.Time, err error) {
|
func (sys *BucketMetadataSys) Update(ctx context.Context, bucket string, configFile string, configData []byte) (updatedAt time.Time, err error) {
|
||||||
return sys.updateAndParse(ctx, bucket, configFile, configData, true)
|
return sys.updateAndParse(ctx, bucket, configFile, configData, true, false)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get metadata for a bucket.
|
// Get metadata for a bucket.
|
||||||
|
|||||||
+41
-3
@@ -249,8 +249,17 @@ func loadBucketMetadataParse(ctx context.Context, objectAPI ObjectLayer, bucket
|
|||||||
}
|
}
|
||||||
|
|
||||||
if len(configs) > 0 {
|
if len(configs) > 0 {
|
||||||
// Old bucket without bucket metadata. Hence we migrate existing settings.
|
if !bucketMetadataLockHeld(ctx, bucket) {
|
||||||
if err = b.convertLegacyConfigs(ctx, objectAPI, configs); err != nil {
|
migrated, lockErr := loadBucketMetadataParseUnderLock(ctx, objectAPI, bucket, parse)
|
||||||
|
if lockErr == nil {
|
||||||
|
return migrated, nil
|
||||||
|
}
|
||||||
|
if !errors.Is(lockErr, errBucketMetadataMigrationLockUnavailable) {
|
||||||
|
return b, lockErr
|
||||||
|
}
|
||||||
|
internalLogOnceIf(ctx, fmt.Errorf("unable to persist bucket metadata migration for %s, using the legacy configuration in memory: %w", bucket, lockErr), "bucket-metadata-migration-lock-"+bucket)
|
||||||
|
b.applyLegacyConfigs(configs)
|
||||||
|
} else if err = b.convertLegacyConfigs(ctx, objectAPI, configs); err != nil {
|
||||||
return b, err
|
return b, err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -270,6 +279,17 @@ func loadBucketMetadataParse(ctx context.Context, objectAPI ObjectLayer, bucket
|
|||||||
}
|
}
|
||||||
|
|
||||||
// migrate unencrypted remote targets
|
// migrate unencrypted remote targets
|
||||||
|
if len(b.BucketTargetsConfigJSON) != 0 && GlobalKMS != nil && len(b.BucketTargetsConfigMetaJSON) == 0 && !bucketMetadataLockHeld(ctx, bucket) {
|
||||||
|
migrated, lockErr := loadBucketMetadataParseUnderLock(ctx, objectAPI, bucket, parse)
|
||||||
|
if lockErr == nil {
|
||||||
|
return migrated, nil
|
||||||
|
}
|
||||||
|
if !errors.Is(lockErr, errBucketMetadataMigrationLockUnavailable) {
|
||||||
|
return b, lockErr
|
||||||
|
}
|
||||||
|
internalLogOnceIf(ctx, fmt.Errorf("unable to persist encrypted bucket target metadata for %s, using the existing configuration in memory: %w", bucket, lockErr), "bucket-metadata-migration-lock-"+bucket)
|
||||||
|
return b, nil
|
||||||
|
}
|
||||||
if err = b.migrateTargetConfig(ctx, objectAPI); err != nil {
|
if err = b.migrateTargetConfig(ctx, objectAPI); err != nil {
|
||||||
return b, err
|
return b, err
|
||||||
}
|
}
|
||||||
@@ -277,6 +297,20 @@ func loadBucketMetadataParse(ctx context.Context, objectAPI ObjectLayer, bucket
|
|||||||
return b, nil
|
return b, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func loadBucketMetadataParseUnderLock(ctx context.Context, objectAPI ObjectLayer, bucket string, parse bool) (BucketMetadata, error) {
|
||||||
|
ctx, unlock, err := lockBucketMetadataWithTimeout(ctx, objectAPI, bucket, bucketMetadataMigrationTimeout)
|
||||||
|
if err != nil {
|
||||||
|
return newBucketMetadata(bucket), fmt.Errorf("%w: %v", errBucketMetadataMigrationLockUnavailable, err)
|
||||||
|
}
|
||||||
|
defer unlock()
|
||||||
|
return loadBucketMetadataParse(ctx, objectAPI, bucket, parse)
|
||||||
|
}
|
||||||
|
|
||||||
|
var (
|
||||||
|
bucketMetadataMigrationTimeout = newDynamicTimeout(5*time.Second, time.Second)
|
||||||
|
errBucketMetadataMigrationLockUnavailable = errors.New("bucket metadata migration lock unavailable")
|
||||||
|
)
|
||||||
|
|
||||||
// loadBucketMetadata loads and migrates to bucket metadata.
|
// loadBucketMetadata loads and migrates to bucket metadata.
|
||||||
func loadBucketMetadata(ctx context.Context, objectAPI ObjectLayer, bucket string) (BucketMetadata, error) {
|
func loadBucketMetadata(ctx context.Context, objectAPI ObjectLayer, bucket string) (BucketMetadata, error) {
|
||||||
return loadBucketMetadataParse(ctx, objectAPI, bucket, true)
|
return loadBucketMetadataParse(ctx, objectAPI, bucket, true)
|
||||||
@@ -440,7 +474,7 @@ func (b *BucketMetadata) getAllLegacyConfigs(ctx context.Context, objectAPI Obje
|
|||||||
return configs, nil
|
return configs, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *BucketMetadata) convertLegacyConfigs(ctx context.Context, objectAPI ObjectLayer, configs map[string][]byte) error {
|
func (b *BucketMetadata) applyLegacyConfigs(configs map[string][]byte) {
|
||||||
for legacyFile, configData := range configs {
|
for legacyFile, configData := range configs {
|
||||||
switch legacyFile {
|
switch legacyFile {
|
||||||
case legacyBucketObjectLockEnabledConfigFile:
|
case legacyBucketObjectLockEnabledConfigFile:
|
||||||
@@ -472,6 +506,10 @@ func (b *BucketMetadata) convertLegacyConfigs(ctx context.Context, objectAPI Obj
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
b.defaultTimestamps()
|
b.defaultTimestamps()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (b *BucketMetadata) convertLegacyConfigs(ctx context.Context, objectAPI ObjectLayer, configs map[string][]byte) error {
|
||||||
|
b.applyLegacyConfigs(configs)
|
||||||
|
|
||||||
if err := b.Save(ctx, objectAPI); err != nil {
|
if err := b.Save(ctx, objectAPI); err != nil {
|
||||||
return err
|
return err
|
||||||
|
|||||||
+47
-13
@@ -882,23 +882,57 @@ func (z *erasureServerPools) MakeBucket(ctx context.Context, bucket string, opts
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
// If it doesn't exist we get a new, so ignore errors
|
if isMinioMetaBucketName(bucket) {
|
||||||
meta := newBucketMetadata(bucket)
|
meta := newBucketMetadata(bucket)
|
||||||
meta.SetCreatedAt(opts.CreatedAt)
|
meta.SetCreatedAt(opts.CreatedAt)
|
||||||
if opts.LockEnabled {
|
if err := meta.Save(context.Background(), z); err != nil {
|
||||||
meta.VersioningConfigXML = enabledBucketVersioningConfig
|
return toObjectErr(err, bucket)
|
||||||
meta.ObjectLockConfigXML = enabledBucketObjectLockConfig
|
}
|
||||||
|
globalBucketMetadataSys.Set(bucket, meta)
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
if opts.VersioningEnabled {
|
ctx, unlock, err := lockBucketMetadata(ctx, z, bucket)
|
||||||
meta.VersioningConfigXML = enabledBucketVersioningConfig
|
if err != nil {
|
||||||
}
|
return toObjectErr(err, bucket)
|
||||||
|
}
|
||||||
if err := meta.Save(context.Background(), z); err != nil {
|
err = func() error {
|
||||||
|
defer unlock()
|
||||||
|
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 err := enablePeerBucketVersioning(&meta); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if len(meta.ObjectLockConfigXML) == 0 {
|
||||||
|
meta.ObjectLockConfigXML = enabledBucketObjectLockConfig
|
||||||
|
meta.ObjectLockConfigUpdatedAt = meta.Created
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if opts.VersioningEnabled {
|
||||||
|
if err := enablePeerBucketVersioning(&meta); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if err = meta.Save(bgContext(ctx), z); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
globalBucketMetadataSys.Set(bucket, meta)
|
||||||
|
return nil
|
||||||
|
}()
|
||||||
|
if err != nil {
|
||||||
return toObjectErr(err, bucket)
|
return toObjectErr(err, bucket)
|
||||||
}
|
}
|
||||||
|
|
||||||
globalBucketMetadataSys.Set(bucket, meta)
|
|
||||||
|
|
||||||
// Success.
|
// Success.
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
+50
-43
@@ -930,35 +930,34 @@ func (c *SiteReplicationSys) PeerBucketMakeWithVersioningHandler(ctx context.Con
|
|||||||
if !ok1 && !ok2 {
|
if !ok1 && !ok2 {
|
||||||
return wrapSRErr(c.annotateErr(makeBucketWithVersion, err))
|
return wrapSRErr(c.annotateErr(makeBucketWithVersion, err))
|
||||||
}
|
}
|
||||||
} else {
|
|
||||||
// Load updated bucket metadata into memory as new
|
|
||||||
// bucket was created.
|
|
||||||
globalNotificationSys.LoadBucketMetadata(GlobalContext, bucket)
|
|
||||||
}
|
}
|
||||||
|
ctx, unlock, err := lockBucketMetadata(ctx, objAPI, bucket)
|
||||||
meta, err := globalBucketMetadataSys.Get(bucket)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return wrapSRErr(c.annotateErr(makeBucketWithVersion, err))
|
return wrapSRErr(c.annotateErr(makeBucketWithVersion, err))
|
||||||
}
|
}
|
||||||
|
err = func() error {
|
||||||
meta.SetCreatedAt(opts.CreatedAt)
|
defer unlock()
|
||||||
|
meta, err := loadBucketMetadataParse(ctx, objAPI, bucket, true)
|
||||||
if err := enablePeerBucketVersioning(&meta); err != nil {
|
if err != nil {
|
||||||
return wrapSRErr(err)
|
return err
|
||||||
}
|
|
||||||
if opts.LockEnabled && len(meta.ObjectLockConfigXML) == 0 {
|
|
||||||
meta.ObjectLockConfigXML = enabledBucketObjectLockConfig
|
|
||||||
if meta.ObjectLockConfigUpdatedAt.IsZero() {
|
|
||||||
meta.ObjectLockConfigUpdatedAt = meta.Created
|
|
||||||
}
|
}
|
||||||
}
|
meta.SetCreatedAt(opts.CreatedAt)
|
||||||
|
|
||||||
if err := meta.Save(context.Background(), objAPI); err != nil {
|
if err = enablePeerBucketVersioning(&meta); err != nil {
|
||||||
return wrapSRErr(err)
|
return err
|
||||||
|
}
|
||||||
|
if opts.LockEnabled && len(meta.ObjectLockConfigXML) == 0 {
|
||||||
|
meta.ObjectLockConfigXML = enabledBucketObjectLockConfig
|
||||||
|
if meta.ObjectLockConfigUpdatedAt.IsZero() {
|
||||||
|
meta.ObjectLockConfigUpdatedAt = meta.Created
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return globalBucketMetadataSys.saveMetadata(bgContext(ctx), objAPI, meta)
|
||||||
|
}()
|
||||||
|
if err != nil {
|
||||||
|
return wrapSRErr(c.annotateErr(makeBucketWithVersion, err))
|
||||||
}
|
}
|
||||||
|
|
||||||
globalBucketMetadataSys.Set(bucket, meta)
|
|
||||||
|
|
||||||
// Load updated bucket metadata into memory as new metadata updated.
|
// Load updated bucket metadata into memory as new metadata updated.
|
||||||
globalNotificationSys.LoadBucketMetadata(GlobalContext, bucket)
|
globalNotificationSys.LoadBucketMetadata(GlobalContext, bucket)
|
||||||
return nil
|
return nil
|
||||||
@@ -1620,13 +1619,18 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context
|
|||||||
if err = validateCORSReplicationPayload(corsConfigData); err != nil {
|
if err = validateCORSReplicationPayload(corsConfigData); err != nil {
|
||||||
return wrapSRErr(err)
|
return wrapSRErr(err)
|
||||||
}
|
}
|
||||||
var unlock func()
|
|
||||||
ctx, unlock, err = lockBucketCORSMetadata(ctx, objectAPI, item.Bucket)
|
|
||||||
if err != nil {
|
|
||||||
return wrapSRErr(err)
|
|
||||||
}
|
|
||||||
defer unlock()
|
|
||||||
}
|
}
|
||||||
|
notifyCtx := ctx
|
||||||
|
ctx, unlock, err := lockBucketMetadata(ctx, objectAPI, item.Bucket)
|
||||||
|
if err != nil {
|
||||||
|
return wrapSRErr(err)
|
||||||
|
}
|
||||||
|
locked := true
|
||||||
|
defer func() {
|
||||||
|
if locked {
|
||||||
|
unlock()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
meta, err := readBucketMetadata(ctx, objectAPI, item.Bucket)
|
meta, err := readBucketMetadata(ctx, objectAPI, item.Bucket)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -1695,7 +1699,13 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return globalBucketMetadataSys.save(ctx, meta)
|
if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, meta); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
unlock()
|
||||||
|
locked = false
|
||||||
|
globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), item.Bucket)
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// PeerBucketPolicyHandler - copies/deletes policy to local cluster.
|
// PeerBucketPolicyHandler - copies/deletes policy to local cluster.
|
||||||
@@ -1960,18 +1970,6 @@ func newBucketCORSReplicationEvent(bucket string, meta BucketMetadata) (madmin.S
|
|||||||
}, true
|
}, true
|
||||||
}
|
}
|
||||||
|
|
||||||
func lockBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket string) (context.Context, func(), error) {
|
|
||||||
// The lock name is deliberately different from .metadata.bin. Saving the
|
|
||||||
// metadata locks that object internally, and namespace locks are not
|
|
||||||
// re-entrant.
|
|
||||||
lock := objectAPI.NewNSLock(minioMetaBucket, pathJoin(bucketMetaPrefix, bucket, "cors-config.lock"))
|
|
||||||
lkctx, err := lock.GetLock(ctx, globalOperationTimeout)
|
|
||||||
if err != nil {
|
|
||||||
return nil, nil, err
|
|
||||||
}
|
|
||||||
return lkctx.Context(), func() { lock.Unlock(lkctx) }, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func updateLocalBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket string, configData []byte) (time.Time, error) {
|
func updateLocalBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket string, configData []byte) (time.Time, error) {
|
||||||
return applyBucketCORSMetadata(ctx, objectAPI, bucket, configData, time.Time{}, true)
|
return applyBucketCORSMetadata(ctx, objectAPI, bucket, configData, time.Time{}, true)
|
||||||
}
|
}
|
||||||
@@ -1984,11 +1982,17 @@ func applyBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket
|
|||||||
return time.Time{}, err
|
return time.Time{}, err
|
||||||
}
|
}
|
||||||
|
|
||||||
ctx, unlock, err := lockBucketCORSMetadata(ctx, objectAPI, bucket)
|
notifyCtx := ctx
|
||||||
|
ctx, unlock, err := lockBucketMetadata(ctx, objectAPI, bucket)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return time.Time{}, err
|
return time.Time{}, err
|
||||||
}
|
}
|
||||||
defer unlock()
|
locked := true
|
||||||
|
defer func() {
|
||||||
|
if locked {
|
||||||
|
unlock()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
var meta BucketMetadata
|
var meta BucketMetadata
|
||||||
if local {
|
if local {
|
||||||
@@ -2026,9 +2030,12 @@ func applyBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket
|
|||||||
|
|
||||||
meta.CorsConfigXML = bytes.Clone(configData)
|
meta.CorsConfigXML = bytes.Clone(configData)
|
||||||
meta.CorsConfigUpdatedAt = updatedAt
|
meta.CorsConfigUpdatedAt = updatedAt
|
||||||
if err = globalBucketMetadataSys.save(ctx, meta); err != nil {
|
if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, meta); err != nil {
|
||||||
return time.Time{}, err
|
return time.Time{}, err
|
||||||
}
|
}
|
||||||
|
unlock()
|
||||||
|
locked = false
|
||||||
|
globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), bucket)
|
||||||
return updatedAt, nil
|
return updatedAt, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -41,9 +41,9 @@ The following are deliberately out of scope:
|
|||||||
|
|
||||||
- redesigning the replication semantics of policy, tags, SSE, quota,
|
- redesigning the replication semantics of policy, tags, SSE, quota,
|
||||||
versioning, or Object Lock;
|
versioning, or Object Lock;
|
||||||
- eliminating lost updates between different bucket-metadata types that all
|
- changing higher-level lifecycle merge semantics or serializing bucket
|
||||||
rewrite `.metadata.bin`; this inherited problem is tracked by
|
deletion against in-flight metadata updates; those follow-ups remain under
|
||||||
[pgsty/silo#77](https://github.com/pgsty/silo/issues/77);
|
[pgsty/silo#102](https://github.com/pgsty/silo/issues/102);
|
||||||
- mixed-version support that permits CORS writes before every site runs a
|
- mixed-version support that permits CORS writes before every site runs a
|
||||||
CORS-aware binary;
|
CORS-aware binary;
|
||||||
- public downgrade, rollback, and global-fallback documentation;
|
- public downgrade, rollback, and global-fallback documentation;
|
||||||
@@ -266,37 +266,44 @@ otherwise:
|
|||||||
```
|
```
|
||||||
|
|
||||||
The admin handler's legacy/default bulk metadata path can also carry a non-nil
|
The admin handler's legacy/default bulk metadata path can also carry a non-nil
|
||||||
CORS field. It therefore takes the same CORS lock, applies strict decoding and
|
CORS field. It therefore takes the shared metadata lock, applies strict
|
||||||
validation, and uses the same state comparison before saving. A nil CORS field
|
decoding and validation, and uses the same state comparison before saving. A
|
||||||
in that untyped legacy shape means "not included" and cannot represent a
|
nil CORS field in that untyped legacy shape means "not included" and cannot
|
||||||
tombstone; current producers use the typed CORS event for deletion.
|
represent a tombstone; current producers use the typed CORS event for deletion.
|
||||||
|
|
||||||
## Concurrency and Locking
|
## Concurrency and Locking
|
||||||
|
|
||||||
The transition lock is:
|
The transition lock is:
|
||||||
|
|
||||||
```text
|
```text
|
||||||
.minio.sys / buckets/<bucket>/cors-config.lock
|
.minio.sys / buckets/<bucket>/metadata.lock
|
||||||
```
|
```
|
||||||
|
|
||||||
It is a virtual distributed namespace lock. The name deliberately differs from
|
It is a virtual distributed namespace lock. The name deliberately differs from
|
||||||
the real `buckets/<bucket>/.metadata.bin` object because the metadata save path
|
the real `buckets/<bucket>/.metadata.bin` object because the metadata save path
|
||||||
locks that object internally and namespace locks are not re-entrant.
|
locks that object internally and namespace locks are not re-entrant.
|
||||||
|
|
||||||
The lock serializes every intentional current-version local, typed-peer,
|
The lock serializes CORS transitions with ordinary `Update`/`Delete`, legacy
|
||||||
legacy-bulk, and local-heal CORS transition across nodes of one cluster. It
|
bulk metadata, imports, bucket creation/adoption, and metadata migrations.
|
||||||
cannot prevent an unrelated whole-record writer from restoring stale CORS
|
Every whole-record writer reads the latest disk state while holding the same
|
||||||
columns. Residual paths include another metadata type's `Update`/`Delete`, a
|
lock, so a writer for another configuration type cannot restore stale CORS
|
||||||
legacy bulk item whose nil CORS field means "not included",
|
columns. Per-type validation, timestamp, and deletion semantics remain
|
||||||
`ImportBucketMetadata`, and bucket-make metadata rewriting. Their inherited
|
independent.
|
||||||
whole-record behavior is the separate architectural problem under issue #77.
|
|
||||||
|
|
||||||
No cross-site admin call or `BucketMetaHook` dispatch is made while holding the
|
No cross-site admin call or `BucketMetaHook` dispatch is made while holding the
|
||||||
CORS lock. The metadata save can perform blocking intra-cluster notification
|
metadata lock. The local disk save and resident-cache update complete under the
|
||||||
fan-out before the lock is released. Local handlers release the lock before
|
lock; intra-cluster metadata reload fan-out happens only after release. This
|
||||||
cross-site dispatch; reordered network delivery is handled by the total-order
|
avoids a peer reload that needs migration from waiting on a lock held by the
|
||||||
|
notifying node. Reordered cross-site delivery is handled by the total-order
|
||||||
join.
|
join.
|
||||||
|
|
||||||
|
The lock name changes from `cors-config.lock` to `metadata.lock`. During a
|
||||||
|
rolling upgrade, old and new nodes therefore do not serialize metadata writers
|
||||||
|
with each other; the shared-lock guarantee begins only after every node in the
|
||||||
|
cluster runs the new binary. Operators should avoid bucket-metadata changes
|
||||||
|
during that window. The on-disk record is unchanged, so rollback remains
|
||||||
|
format-compatible.
|
||||||
|
|
||||||
## Dispatch and Retry
|
## Dispatch and Retry
|
||||||
|
|
||||||
PUT sends a typed `SRBucketMetaTypeCorsConfig` event with canonical base64 XML
|
PUT sends a typed `SRBucketMetaTypeCorsConfig` event with canonical base64 XML
|
||||||
@@ -441,11 +448,12 @@ false tombstone source.
|
|||||||
Rejected. The save path takes the same namespace lock internally; reusing it
|
Rejected. The save path takes the same namespace lock internally; reusing it
|
||||||
would self-deadlock.
|
would self-deadlock.
|
||||||
|
|
||||||
### Redesign every bucket metadata type together
|
### Give every metadata type a new state machine
|
||||||
|
|
||||||
Rejected for issue #75. Neighboring metadata types have related inherited
|
Rejected. The shared lock prevents whole-record lost updates without changing
|
||||||
patterns but different delete, validation, and compatibility semantics. They
|
the independent replication, validation, or deletion semantics of policy,
|
||||||
require focused reproductions under issue #77.
|
tags, SSE, quota, versioning, and Object Lock. Those semantic audits remain
|
||||||
|
separate from the persistence fix in issue #102.
|
||||||
|
|
||||||
## Invariants
|
## Invariants
|
||||||
|
|
||||||
@@ -455,9 +463,9 @@ The implementation is acceptable only while all of these invariants hold:
|
|||||||
2. Nil payload plus non-zero timestamp is a durable tombstone.
|
2. Nil payload plus non-zero timestamp is a durable tombstone.
|
||||||
3. A live payload has canonical base64 on the wire, valid CORS XML, and a
|
3. A live payload has canonical base64 on the wire, valid CORS XML, and a
|
||||||
non-zero source timestamp.
|
non-zero source timestamp.
|
||||||
4. Every intentional current-version CORS state transition is serialized by
|
4. Every intentional current-version CORS state transition and every other
|
||||||
the CORS namespace lock from disk read through state comparison and save;
|
whole-record metadata writer is serialized by `metadata.lock` from the
|
||||||
unrelated whole-record overwrite risk remains explicitly under issue #77.
|
authoritative disk read through save and local cache publication.
|
||||||
5. Peer apply and heal never replace local state with a lower or equal state.
|
5. Peer apply and heal never replace local state with a lower or equal state.
|
||||||
6. Local PUT and DELETE create a state strictly greater than the state observed
|
6. Local PUT and DELETE create a state strictly greater than the state observed
|
||||||
under the lock.
|
under the lock.
|
||||||
@@ -542,10 +550,10 @@ gap. The selected C-prime model incorporated the valid findings while rejecting
|
|||||||
the suggestion to rewrite normal source timestamps.
|
the suggestion to rewrite normal source timestamps.
|
||||||
|
|
||||||
The second review found no P0. Its `GO WITH FIXES` findings were peer semantic
|
The second review found no P0. Its `GO WITH FIXES` findings were peer semantic
|
||||||
validation, the legacy/default admin mutation path bypassing the CORS lock and
|
validation, the legacy/default admin mutation path bypassing the then-current
|
||||||
join, CreatedAt-floor observability, and missing tests for invalid XML, lineage,
|
CORS lock and join, CreatedAt-floor observability, and missing tests for
|
||||||
and concurrent local transitions. Those required changes and tests are now in
|
invalid XML, lineage, and concurrent local transitions. Those required changes
|
||||||
the working tree.
|
and tests are now in the working tree.
|
||||||
|
|
||||||
The final review examined this design and the exact dirty diff, independently
|
The final review examined this design and the exact dirty diff, independently
|
||||||
reran build, vet, lint, normal tests, and race tests, and found no P0 or P1.
|
reran build, vet, lint, normal tests, and race tests, and found no P0 or P1.
|
||||||
|
|||||||
Reference in New Issue
Block a user