fix(replication): apply bucket metadata source times under the metadata lock

Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
Feng Ruohang
2026-09-12 11:19:07 +08:00
parent 5c57658163
commit bcc62afe3d
8 changed files with 1134 additions and 218 deletions
+33 -9
View File
@@ -76,7 +76,7 @@ func (a adminAPIHandlers) PutBucketQuotaConfigHandler(w http.ResponseWriter, r *
return
}
quotaConfig, err := parseBucketQuota(bucket, data)
_, err = parseBucketQuota(bucket, data)
if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return
@@ -94,9 +94,6 @@ func (a adminAPIHandlers) PutBucketQuotaConfigHandler(w http.ResponseWriter, r *
Quota: data,
UpdatedAt: updatedAt,
}
if quotaConfig.Size == 0 && quotaConfig.Quota == 0 {
bucketMeta.Quota = nil
}
// Call site replication hook.
replLogIf(ctx, globalSiteReplicationSys.BucketMetaHook(ctx, bucketMeta))
@@ -922,7 +919,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
continue
}
configData, err := json.Marshal(bucketPolicy)
configData, err := canonicalBucketPolicy(bucketPolicy)
if err != nil {
rpt.SetStatus(bucket, fileName, err)
continue
@@ -1081,18 +1078,39 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
continue
}
var merged BucketMetadata
var commitAt time.Time
err := func() error {
lockCtx, unlock, err := lockBucketMetadata(ctx, objectAPI, bucket)
if err != nil {
return err
}
defer unlock()
merged, err = loadBucketMetadataParse(lockCtx, objectAPI, bucket, true)
merged, err = loadBucketMetadataParse(lockCtx, objectAPI, bucket, false)
if err != nil {
return err
}
if err := ensureBucketMetadataCreated(lockCtx, objectAPI, &merged); err != nil {
return err
}
commitAt = UTCNow()
for _, file := range replicatedBucketConfigs {
if _, ok := fields[file]; ok {
commitAt = localBucketConfigUpdatedAt(merged, file, commitAt)
}
}
applyImportedBucketMetadata(&merged, *meta, fields)
return globalBucketMetadataSys.saveMetadata(lockCtx, objectAPI, merged)
for _, file := range replicatedBucketConfigs {
if _, ok := fields[file]; !ok {
continue
}
data, at := replicatedBucketConfig(&merged, file)
payload, _, err := bucketConfigPayload(bucket, file, *data, len(merged.ObjectLockConfigXML) != 0)
if err != nil {
return err
}
*data, *at = payload, commitAt
}
return globalBucketMetadataSys.saveMetadata(lockCtx, objectAPI, &merged)
}()
if err != nil {
rpt.SetStatus(bucket, "", err)
@@ -1100,7 +1118,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
}
*meta = merged
globalNotificationSys.LoadBucketMetadata(bgContext(ctx), bucket)
hook := madmin.SRBucketMeta{Bucket: bucket, UpdatedAt: updatedAt}
hook := madmin.SRBucketMeta{Bucket: bucket, UpdatedAt: commitAt}
var hookNeeded bool
if _, ok := fields[bucketQuotaConfigFile]; ok {
hook.Quota = meta.QuotaConfigJSON
@@ -1108,7 +1126,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
}
if _, ok := fields[bucketPolicyConfig]; ok {
hook.Policy = meta.PolicyConfigJSON
hookNeeded = true
hookNeeded = hookNeeded || len(hook.Policy) != 0
}
if _, ok := fields[bucketVersioningConfig]; ok {
hook.Versioning = enc(meta.VersioningConfigXML)
@@ -1129,6 +1147,12 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r *
if hookNeeded {
err = globalSiteReplicationSys.BucketMetaHook(ctx, hook)
}
if _, ok := fields[bucketPolicyConfig]; ok && len(meta.PolicyConfigJSON) == 0 {
// An omitted bulk Policy cannot express deletion.
err = errors.Join(err, globalSiteReplicationSys.BucketMetaHook(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypePolicy, Bucket: bucket, UpdatedAt: commitAt,
}))
}
if _, ok := fields[bucketCorsConfig]; ok {
// CORS carries its own timestamp, so it replicates through the
// dedicated event rather than the shared bucket metadata hook. It
+318
View File
@@ -0,0 +1,318 @@
// Copyright (c) 2015-2026 MinIO, Inc.
//
// 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.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"context"
"encoding/json"
"errors"
"sort"
"time"
"github.com/minio/minio-go/v7/pkg/tags"
bucketsse "github.com/minio/minio/internal/bucket/encryption"
objectlock "github.com/minio/minio/internal/bucket/object/lock"
"github.com/minio/minio/internal/bucket/versioning"
"github.com/pgsty/silo-pkg/v3/policy"
)
// Only these fields share the site-replication source-time ordering contract.
// Object Lock is applied before Versioning, whose effective document depends on it.
var replicatedBucketConfigs = [...]string{
objectLockConfig, bucketVersioningConfig, bucketPolicyConfig,
bucketTaggingConfig, bucketSSEConfig, bucketQuotaConfigFile,
}
func replicatedBucketConfig(meta *BucketMetadata, file string) (*[]byte, *time.Time) {
switch file {
case bucketPolicyConfig:
return &meta.PolicyConfigJSON, &meta.PolicyConfigUpdatedAt
case bucketTaggingConfig:
return &meta.TaggingConfigXML, &meta.TaggingConfigUpdatedAt
case bucketSSEConfig:
return &meta.EncryptionConfigXML, &meta.EncryptionConfigUpdatedAt
case bucketQuotaConfigFile:
return &meta.QuotaConfigJSON, &meta.QuotaConfigUpdatedAt
case bucketVersioningConfig:
return &meta.VersioningConfigXML, &meta.VersioningConfigUpdatedAt
case objectLockConfig:
return &meta.ObjectLockConfigXML, &meta.ObjectLockConfigUpdatedAt
}
return nil, nil
}
func bucketConfigUpdateOnly(file string) bool {
return file == bucketVersioningConfig || file == objectLockConfig
}
// Reuse the persistence rule before comparison, so an accepted Versioning
// event and the document Save actually writes have the same comparison key.
func effectiveBucketVersioning(data []byte, lockEnabled bool) []byte {
if lockEnabled {
config, err := versioning.ParseConfig(bytes.NewReader(data))
if err != nil || !config.Enabled() || config.PrefixesExcluded() {
return enabledBucketVersioningConfig
}
}
return data
}
// A parsed policy still contains map-backed sets with nondeterministic Marshal
// order. Sort every set array recursively, including statements and conditions.
// RawMessage keeps integer values intact; decoding through float64 would not.
func canonicalBucketPolicyJSON(data json.RawMessage) (json.RawMessage, error) {
data = bytes.TrimSpace(data)
if len(data) == 0 {
return nil, nil
}
switch data[0] {
case '{':
var obj map[string]json.RawMessage
if err := json.Unmarshal(data, &obj); err != nil {
return nil, err
}
for key, value := range obj {
var err error
obj[key], err = canonicalBucketPolicyJSON(value)
if err != nil {
return nil, err
}
}
return json.Marshal(obj)
case '[':
var arr []json.RawMessage
if err := json.Unmarshal(data, &arr); err != nil {
return nil, err
}
for i := range arr {
var err error
arr[i], err = canonicalBucketPolicyJSON(arr[i])
if err != nil {
return nil, err
}
}
sort.Slice(arr, func(i, j int) bool { return bytes.Compare(arr[i], arr[j]) < 0 })
return json.Marshal(arr)
default:
var compact bytes.Buffer
if err := json.Compact(&compact, data); err != nil {
return nil, err
}
return compact.Bytes(), nil
}
}
// Encode the validated policy fields explicitly: BPStatement's required
// Action/Resource tags otherwise try to marshal empty sets for the supported
// NotAction/NotResource alternatives. This stays within the existing schema.
func canonicalBucketPolicy(cfg *policy.BucketPolicy) ([]byte, error) {
if cfg.IsEmpty() {
return nil, nil
}
doc := map[string]any{"Version": cfg.Version}
if cfg.ID != "" {
doc["ID"] = cfg.ID
}
statements := make([]map[string]any, 0, len(cfg.Statements))
for _, st := range cfg.Statements {
statement := map[string]any{"Effect": st.Effect, "Principal": st.Principal}
if st.SID != "" {
statement["Sid"] = st.SID
}
if len(st.Actions) != 0 {
statement["Action"] = st.Actions
}
if len(st.NotActions) != 0 {
statement["NotAction"] = st.NotActions
}
if len(st.Resources) != 0 {
statement["Resource"] = st.Resources
}
if len(st.NotResources) != 0 {
statement["NotResource"] = st.NotResources
}
if len(st.Conditions) != 0 {
statement["Condition"] = st.Conditions
}
statements = append(statements, statement)
}
doc["Statement"] = statements
data, err := json.Marshal(doc)
if err != nil {
return nil, err
}
return canonicalBucketPolicyJSON(data)
}
// Validate with the same parsers as Save. Policy's established empty-policy
// semantics are deletion; a parsed zero quota is still a live document.
func bucketConfigPayload(bucket, file string, data []byte, lockEnabled bool) ([]byte, []byte, error) {
if len(data) == 0 {
return nil, nil, nil
}
var err error
key := data
switch file {
case bucketPolicyConfig:
var cfg *policy.BucketPolicy
cfg, err = policy.ParseBucketPolicyConfig(bytes.NewReader(data), bucket)
if err == nil {
if cfg.IsEmpty() {
return nil, nil, nil
}
key, err = canonicalBucketPolicy(cfg)
}
case bucketQuotaConfigFile:
cfg, parseErr := parseBucketQuota(bucket, data)
err = parseErr
if err == nil {
key, err = json.Marshal(cfg)
}
case bucketTaggingConfig:
_, err = tags.ParseBucketXML(bytes.NewReader(data))
case bucketSSEConfig:
_, err = bucketsse.ParseBucketSSEConfig(bytes.NewReader(data))
case objectLockConfig:
_, err = objectlock.ParseObjectLockConfig(bytes.NewReader(data))
case bucketVersioningConfig:
data = effectiveBucketVersioning(data, lockEnabled)
key = data
_, err = versioning.ParseConfig(bytes.NewReader(data))
}
return data, key, err
}
type bucketConfigState struct {
data, key []byte
at time.Time
real, valid bool
}
func newBucketConfigState(bucket, file string, data []byte, at, created time.Time, lockEnabled bool) (bucketConfigState, error) {
data, key, err := bucketConfigPayload(bucket, file, data, lockEnabled)
if err != nil {
return bucketConfigState{}, err
}
if at.IsZero() {
at = created
}
valid := !created.IsZero() && !at.Before(created)
real := valid && at.After(created)
if bucketConfigUpdateOnly(file) && len(data) == 0 {
real = false
}
return bucketConfigState{data: data, key: key, at: at, real: real, valid: valid}, nil
}
func (s bucketConfigState) candidate() bool {
return s.valid && (s.real || len(s.data) != 0)
}
func compareBucketConfigStates(a, b bucketConfigState) int {
if a.valid != b.valid {
if a.valid {
return 1
}
return -1
}
if a.real != b.real {
if a.real {
return 1
}
return -1
}
if a.real {
if n := a.at.Compare(b.at); n != 0 {
return n
}
// At equal source time a real deletion wins, preventing resurrection.
if (len(a.data) == 0) != (len(b.data) == 0) {
if len(a.data) == 0 {
return 1
}
return -1
}
}
return bytes.Compare(a.key, b.key)
}
func localBucketConfigUpdatedAt(meta BucketMetadata, file string, now time.Time) time.Time {
_, at := replicatedBucketConfig(&meta, file)
for _, lower := range []time.Time{meta.Created, *at} {
if !now.After(lower) {
now = lower.Add(time.Nanosecond)
}
}
return now.UTC()
}
func ensureBucketMetadataCreated(ctx context.Context, obj ObjectLayer, meta *BucketMetadata) error {
if !meta.Created.IsZero() {
return nil
}
info, err := obj.GetBucketInfo(ctx, meta.Name, BucketOptions{NoMetadata: true})
if err != nil {
return err
}
if info.Created.IsZero() {
return errors.New("bucket metadata creation time is unknown")
}
meta.Created = info.Created.UTC()
return nil
}
// applyBucketConfig runs under metadata.lock, on freshly loaded metadata. It
// changes only the selected field; the caller persists once after all checks.
func applyBucketConfig(meta *BucketMetadata, file string, data []byte, at time.Time) (bool, error) {
if bucketConfigUpdateOnly(file) && len(data) == 0 {
return false, nil
}
current, currentAt := replicatedBucketConfig(meta, file)
lockEnabled := len(meta.ObjectLockConfigXML) != 0
incoming, err := newBucketConfigState(meta.Name, file, data, at, meta.Created, lockEnabled)
if err != nil {
return false, err
}
if !incoming.candidate() {
return false, nil
}
local, err := newBucketConfigState(meta.Name, file, *current, *currentAt, meta.Created, lockEnabled)
if err != nil {
return false, err
}
if compareBucketConfigStates(incoming, local) <= 0 {
return false, nil
}
*current, *currentAt = bytes.Clone(incoming.data), incoming.at.UTC()
return true, nil
}
func rebaseBucketConfigDefaults(meta *BucketMetadata, oldCreated time.Time) {
if meta.Created.Equal(oldCreated) {
return
}
// These six fields alone use Created to distinguish a baseline from a
// tombstone. Preserve actual source times when adopting an existing bucket.
for _, file := range replicatedBucketConfigs {
_, at := replicatedBucketConfig(meta, file)
if at.IsZero() || at.Equal(oldCreated) {
*at = meta.Created
}
}
}
+77 -52
View File
@@ -126,19 +126,38 @@ func (sys *BucketMetadataSys) Set(bucket string, meta BucketMetadata) {
}
}
func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string, configFile string, configData []byte, parse, lifecycleDelete bool) (updatedAt time.Time, err error) {
type bucketMetadataUpdate struct {
meta BucketMetadata
updatedAt time.Time
changed bool
}
func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket, configFile string, configData []byte, parse, lifecycleDelete bool) (time.Time, error) {
result, err := sys.updateAndParseMetadata(ctx, bucket, configFile, configData, parse, lifecycleDelete, nil)
return result.updatedAt, err
}
func (sys *BucketMetadataSys) updateAndParseMetadata(ctx context.Context, bucket string, configFile string, configData []byte, parse, lifecycleDelete bool, sourceTime *time.Time) (result bucketMetadataUpdate, err error) {
objAPI := newObjectLayerFn()
if objAPI == nil {
return updatedAt, errServerNotInitialized
return result, errServerNotInitialized
}
if isMinioMetaBucketName(bucket) {
return updatedAt, errInvalidArgument
return result, errInvalidArgument
}
// Load deletions without parsed caches (notably quota), and compare the
// six replicated fields against the raw document under the same lock.
if data, _ := replicatedBucketConfig(&result.meta, configFile); data != nil {
parse = false
if bucketConfigUpdateOnly(configFile) && len(configData) == 0 {
return result, nil
}
}
notifyCtx := ctx
ctx, unlock, err := lockBucketMetadata(ctx, objAPI, bucket)
if err != nil {
return updatedAt, err
return result, err
}
err = func() error {
@@ -158,55 +177,61 @@ func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string,
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 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)
updatedAt := UTCNow()
if data, _ := replicatedBucketConfig(&meta, configFile); data != nil {
if err := ensureBucketMetadataCreated(ctx, objAPI, &meta); err != nil {
return err
}
if sourceTime == nil || sourceTime.IsZero() {
updatedAt = localBucketConfigUpdatedAt(meta, configFile, updatedAt)
} else {
updatedAt = sourceTime.UTC()
}
changed, err := applyBucketConfig(&meta, configFile, configData, updatedAt)
if err != nil {
return err
}
if !changed {
return nil
}
} else {
switch configFile {
case bucketNotificationConfig:
meta.NotificationConfigXML = configData
meta.NotificationConfigUpdatedAt = updatedAt
case bucketLifecycleConfig:
meta.LifecycleConfigXML = configData
meta.LifecycleConfigUpdatedAt = 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)
}
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 := sys.saveMetadata(ctx, objAPI, &meta); err != nil {
return err
}
result = bucketMetadataUpdate{meta: meta, updatedAt: updatedAt, changed: true}
return nil
}()
if err != nil {
return updatedAt, err
return result, err
}
globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), bucket) // Do not use caller context here
return updatedAt, nil
if result.changed {
globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), bucket)
}
return result, nil
}
func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) error {
@@ -219,7 +244,7 @@ func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) err
return errInvalidArgument
}
if err := sys.saveMetadata(ctx, objAPI, meta); err != nil {
if err := sys.saveMetadata(ctx, objAPI, &meta); err != nil {
return err
}
@@ -229,7 +254,7 @@ func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) err
// 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 {
func (sys *BucketMetadataSys) saveMetadata(ctx context.Context, objAPI ObjectLayer, meta *BucketMetadata) error {
// A writer may have queued for metadata.lock before DeleteBucket completed.
// Recheck the physical bucket under that lock, before recreating metadata.
if _, err := objAPI.GetBucketInfo(ctx, meta.Name, BucketOptions{NoMetadata: true}); err != nil {
@@ -238,7 +263,7 @@ func (sys *BucketMetadataSys) saveMetadata(ctx context.Context, objAPI ObjectLay
if err := meta.Save(ctx, objAPI); err != nil {
return err
}
sys.Set(meta.Name, meta)
sys.Set(meta.Name, *meta)
return nil
}
@@ -351,7 +376,7 @@ func (sys *BucketMetadataSys) UpdateExpiryLCConfig(ctx context.Context, bucket s
}
meta.LifecycleConfigXML = configData
meta.LifecycleConfigUpdatedAt = UTCNow()
return sys.saveMetadata(ctx, objAPI, meta)
return sys.saveMetadata(ctx, objAPI, &meta)
}()
if err != nil {
return err
+1 -9
View File
@@ -388,15 +388,7 @@ func (b *BucketMetadata) parseAllConfigs(ctx context.Context, objectAPI ObjectLa
} else {
b.objectLockConfig = nil
}
if b.objectLockConfig != nil {
// Object Lock requires every object to be versioned. Whatever the lock
// document contains, a suspended or prefix-excluded versioning document
// is replaced by plain Enabled versioning; Save persists the result.
config, versioningErr := versioning.ParseConfig(bytes.NewReader(b.VersioningConfigXML))
if versioningErr != nil || !config.Enabled() || config.PrefixesExcluded() {
b.VersioningConfigXML = enabledBucketVersioningConfig
}
}
b.VersioningConfigXML = effectiveBucketVersioning(b.VersioningConfigXML, b.objectLockConfig != nil)
if len(b.VersioningConfigXML) != 0 {
b.versioningConfig, err = versioning.ParseConfig(bytes.NewReader(b.VersioningConfigXML))
+4 -4
View File
@@ -100,13 +100,13 @@ func (api objectAPIHandlers) PutBucketPolicyHandler(w http.ResponseWriter, r *ht
return
}
configData, err := json.Marshal(bucketPolicy)
configData, err := canonicalBucketPolicy(bucketPolicy)
if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return
}
updatedAt, err := globalBucketMetadataSys.Update(ctx, bucket, bucketPolicyConfig, configData)
result, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketPolicyConfig, configData, false, false, nil)
if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return
@@ -116,8 +116,8 @@ func (api objectAPIHandlers) PutBucketPolicyHandler(w http.ResponseWriter, r *ht
replLogIf(ctx, globalSiteReplicationSys.BucketMetaHook(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypePolicy,
Bucket: bucket,
Policy: bucketPolicyBytes,
UpdatedAt: updatedAt,
Policy: result.meta.PolicyConfigJSON,
UpdatedAt: result.updatedAt,
}))
// Success.
+3 -3
View File
@@ -97,7 +97,7 @@ func (api objectAPIHandlers) PutBucketVersioningHandler(w http.ResponseWriter, r
return
}
updatedAt, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, configData)
result, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketVersioningConfig, configData, false, false, nil)
if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return
@@ -107,12 +107,12 @@ func (api objectAPIHandlers) PutBucketVersioningHandler(w http.ResponseWriter, r
//
// We encode the xml bytes as base64 to ensure there are no encoding
// errors.
cfgStr := base64.StdEncoding.EncodeToString(configData)
cfgStr := base64.StdEncoding.EncodeToString(result.meta.VersioningConfigXML)
replLogIf(ctx, globalSiteReplicationSys.BucketMetaHook(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeVersionConfig,
Bucket: bucket,
Versioning: &cfgStr,
UpdatedAt: updatedAt,
UpdatedAt: result.updatedAt,
}))
writeSuccessResponseHeadersOnly(w)
+624
View File
@@ -0,0 +1,624 @@
// Copyright (c) 2015-2026 MinIO, Inc.
//
// 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.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"context"
"encoding/base64"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"path"
"sync"
"sync/atomic"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
)
func TestPeerBucketMetadataSourceTimeAndDeletion(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: testPeerBucketMetadataSourceTimeAndDeletion})
}
func testPeerBucketMetadataSourceTimeAndDeletion(obj ObjectLayer, backend, bucket string, _ http.Handler, cred auth.Credentials, t *testing.T) {
ctx := t.Context()
created := UTCNow().Add(-time.Hour)
putAt := created.Add(10 * time.Minute)
delAt := putAt.Add(time.Minute)
enc := func(s string) *string { v := base64.StdEncoding.EncodeToString([]byte(s)); return &v }
base := newBucketMetadata(bucket)
base.SetCreatedAt(created)
base.defaultTimestamps()
quotaJSON, err := json.Marshal(madmin.BucketQuota{Quota: 1024, Type: madmin.HardQuota})
if err != nil {
t.Fatal(err)
}
cases := []struct {
name, file string
put madmin.SRBucketMeta
value func(BucketMetadata) []byte
stamp func(BucketMetadata) time.Time
exported func(madmin.SRBucketInfo) time.Time
deletable bool
}{
{"policy", bucketPolicyConfig, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypePolicy, Policy: []byte(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket))}, func(m BucketMetadata) []byte { return m.PolicyConfigJSON }, func(m BucketMetadata) time.Time { return m.PolicyConfigUpdatedAt }, func(m madmin.SRBucketInfo) time.Time { return m.PolicyUpdatedAt }, true},
{"tags", bucketTaggingConfig, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypeTags, Tags: enc(`<Tagging><TagSet><Tag><Key>key</Key><Value>old</Value></Tag></TagSet></Tagging>`)}, func(m BucketMetadata) []byte { return m.TaggingConfigXML }, func(m BucketMetadata) time.Time { return m.TaggingConfigUpdatedAt }, func(m madmin.SRBucketInfo) time.Time { return m.TagConfigUpdatedAt }, true},
{"sse", bucketSSEConfig, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypeSSEConfig, SSEConfig: enc(`<ServerSideEncryptionConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ApplyServerSideEncryptionByDefault><SSEAlgorithm>AES256</SSEAlgorithm></ApplyServerSideEncryptionByDefault></Rule></ServerSideEncryptionConfiguration>`)}, func(m BucketMetadata) []byte { return m.EncryptionConfigXML }, func(m BucketMetadata) time.Time { return m.EncryptionConfigUpdatedAt }, func(m madmin.SRBucketInfo) time.Time { return m.SSEConfigUpdatedAt }, true},
{"quota", bucketQuotaConfigFile, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypeQuotaConfig, Quota: quotaJSON}, func(m BucketMetadata) []byte { return m.QuotaConfigJSON }, func(m BucketMetadata) time.Time { return m.QuotaConfigUpdatedAt }, func(m madmin.SRBucketInfo) time.Time { return m.QuotaConfigUpdatedAt }, true},
{"versioning", bucketVersioningConfig, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypeVersionConfig, Versioning: enc(`<VersioningConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Status>Enabled</Status></VersioningConfiguration>`)}, func(m BucketMetadata) []byte { return m.VersioningConfigXML }, func(m BucketMetadata) time.Time { return m.VersioningConfigUpdatedAt }, nil, false},
{"objectlock", objectLockConfig, newSRBucketObjectLockMeta(bucket, enc(`<ObjectLockConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><ObjectLockEnabled>Enabled</ObjectLockEnabled><Rule><DefaultRetention><Mode>GOVERNANCE</Mode><Days>30</Days></DefaultRetention></Rule></ObjectLockConfiguration>`), putAt), func(m BucketMetadata) []byte { return m.ObjectLockConfigXML }, func(m BucketMetadata) time.Time { return m.ObjectLockConfigUpdatedAt }, nil, false},
}
for _, tc := range cases {
t.Run(backend+"/"+tc.name, func(t *testing.T) {
if err := globalBucketMetadataSys.save(ctx, base); err != nil {
t.Fatal(err)
}
item := tc.put
item.Bucket, item.UpdatedAt = bucket, putAt
apply := func(item madmin.SRBucketMeta) {
t.Helper()
rec := applySRBucketMetaViaAdmin(t, cred, item)
if rec.Code != http.StatusOK {
t.Fatalf("admin apply returned %d: %s", rec.Code, rec.Body.String())
}
}
read := func() BucketMetadata {
t.Helper()
m, err := loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
return m
}
apply(item)
put := read()
if len(tc.value(put)) == 0 {
t.Fatal("PUT did not establish live config")
}
if !tc.stamp(put).Equal(putAt) {
t.Errorf("SOURCE_TIME: persisted %s, want source %s", tc.stamp(put), putAt)
}
apply(madmin.SRBucketMeta{Type: item.Type, Bucket: bucket, UpdatedAt: delAt})
after := read()
if !tc.deletable {
if !bytes.Equal(tc.value(put), tc.value(after)) || !tc.stamp(put).Equal(tc.stamp(after)) {
t.Error("NIL_NOOP: update-only config changed")
}
t.Log("nil payload is correctly a no-op")
return
}
if len(tc.value(after)) != 0 {
t.Error("NEWER_DELETE: HTTP 200, but live config remained after newer source DELETE")
}
if _, err := globalBucketMetadataSys.Delete(ctx, bucket, tc.file); err != nil {
t.Fatal(err)
}
tombstone := read()
if len(tc.value(tombstone)) != 0 {
t.Fatal("local DELETE failed to establish tombstone")
}
apply(item)
after = read()
if len(tc.value(after)) != 0 {
t.Error("STALE_RESURRECTION: older source PUT resurrected a locally deleted config")
} else {
t.Log("older source PUT did not resurrect config")
}
})
}
}
func TestPeerBucketMetadataBulkOrdering(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, cred auth.Credentials, t *testing.T) {
t.Run(backend, func(t *testing.T) {
ctx := t.Context()
meta, err := loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
newAt := meta.Created.Add(time.Hour)
newXML := `<Tagging><TagSet><Tag><Key>key</Key><Value>new</Value></Tag></TagSet></Tagging>`
meta.TaggingConfigXML, meta.TaggingConfigUpdatedAt = []byte(newXML), newAt
if err := globalBucketMetadataSys.save(ctx, meta); err != nil {
t.Fatal(err)
}
oldXML := base64.StdEncoding.EncodeToString([]byte(`<Tagging><TagSet><Tag><Key>key</Key><Value>old</Value></Tag></TagSet></Tagging>`))
item := madmin.SRBucketMeta{Bucket: bucket, Tags: &oldXML, UpdatedAt: newAt.Add(-time.Minute)}
rec := applySRBucketMetaViaAdmin(t, cred, item)
if rec.Code != http.StatusOK {
t.Fatalf("bulk apply returned %d: %s", rec.Code, rec.Body.String())
}
after, err := loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
if !bytes.Equal(after.TaggingConfigXML, []byte(newXML)) || !after.TaggingConfigUpdatedAt.Equal(newAt) {
t.Errorf("BULK_STALE_OVERWRITE: older bulk event overwrote newer config, value=%q time=%s", after.TaggingConfigXML, after.TaggingConfigUpdatedAt)
}
})
}})
}
func TestPeerBucketMetadataOrderingUnderLock(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) {
t.Run(backend, func(t *testing.T) {
ctx := t.Context()
meta, err := loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
oldAt, newAt := meta.Created.Add(time.Hour), meta.Created.Add(2*time.Hour)
oldXML := base64.StdEncoding.EncodeToString([]byte(`<Tagging><TagSet><Tag><Key>key</Key><Value>older</Value></Tag></TagSet></Tagging>`))
newXML := []byte(`<Tagging><TagSet><Tag><Key>key</Key><Value>newest</Value></Tag></TagSet></Tagging>`)
lockCtx, unlock, err := lockBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
locked := true
defer func() {
if locked {
unlock()
}
}()
ready := make(chan struct{}, 1)
hook := func(name string) {
if name == bucket {
select {
case ready <- struct{}{}:
default:
}
}
}
lockBucketMetadataAcquireHook.Store(&hook)
defer lockBucketMetadataAcquireHook.Store(nil)
done := make(chan error, 1)
go func() { done <- globalSiteReplicationSys.PeerBucketTaggingHandler(ctx, bucket, &oldXML, oldAt) }()
select {
case <-ready:
case <-time.After(5 * time.Second):
t.Fatal("older event never reached metadata lock")
}
// A newer writer commits while holding the existing metadata lock.
// The old peer event has already checked the pre-commit cache.
meta.TaggingConfigXML, meta.TaggingConfigUpdatedAt = newXML, newAt
if err := globalBucketMetadataSys.saveMetadata(lockCtx, obj, &meta); err != nil {
t.Fatal(err)
}
unlock()
locked = false
select {
case err := <-done:
if err != nil {
t.Fatal(err)
}
case <-time.After(5 * time.Second):
t.Fatal("older event did not finish")
}
after, err := loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
if string(after.TaggingConfigXML) != string(newXML) {
t.Errorf("CHECK_OUTSIDE_LOCK: queued older event replaced newer committed tags with %q", after.TaggingConfigXML)
}
})
}})
}
// Count actual metadata persistence, including writes that would be invisible
// in a value-only assertion after a duplicate or rejected event.
type bucketConfigWriteCounter struct {
ObjectLayer
writes atomic.Int64
}
func (o *bucketConfigWriteCounter) PutObject(ctx context.Context, bucket, object string, data *PutObjReader, opts ObjectOptions) (ObjectInfo, error) {
if bucket == minioMetaBucket && path.Base(object) == bucketMetadataFile {
o.writes.Add(1)
}
return o.ObjectLayer.PutObject(ctx, bucket, object, data, opts)
}
func TestBucketConfigStateOrdering(t *testing.T) {
created := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC)
live := []byte(`<Tagging><TagSet><Tag><Key>k</Key><Value>a</Value></Tag></TagSet></Tagging>`)
other := bytes.ReplaceAll(live, []byte("<Value>a"), []byte("<Value>b"))
state := func(data []byte, at, birth time.Time) bucketConfigState {
t.Helper()
s, err := newBucketConfigState("bucket", bucketTaggingConfig, data, at, birth, false)
if err != nil {
t.Fatal(err)
}
return s
}
baseline := state(nil, created, created)
baselineLive := state(live, created, created)
put := state(live, created.Add(time.Second), created)
deleted := state(nil, put.at, created)
lateBaseline := state(other, created.Add(time.Hour), created.Add(time.Hour))
tests := []struct {
name string
lower, higher bucketConfigState
}{
{"baseline initialization", baseline, baselineLive},
{"real before later baseline", lateBaseline, put},
{"delete wins tie", put, deleted},
{"live key tie", put, state(other, put.at, created)},
{"new PUT after deletion", deleted, state(live, put.at.Add(time.Second), created)},
{"baseline key ignores creation time", state(live, created.Add(time.Hour), created.Add(time.Hour)), state(other, created, created)},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
if compareBucketConfigStates(tc.lower, tc.higher) >= 0 || compareBucketConfigStates(tc.higher, tc.lower) <= 0 {
t.Fatal("ordering is not antisymmetric or chose the wrong winner")
}
})
}
for _, s := range []bucketConfigState{baseline, state(live, created.Add(-time.Second), created), state(live, created, time.Time{})} {
if s.candidate() {
t.Fatal("default, pre-creation, or unknown-generation state became a source")
}
}
for _, file := range []string{bucketVersioningConfig, objectLockConfig} {
s, err := newBucketConfigState("bucket", file, nil, put.at, created, false)
if err != nil || s.candidate() {
t.Fatalf("empty update-only candidate: %+v %v", s, err)
}
}
}
func TestBucketPolicyReplicationKey(t *testing.T) {
// Permute independent set arrays and parse separately, as two sites would.
a := []byte(`{"Version":"2012-10-17","Statement":[{"Sid":"one","Effect":"Allow","Principal":{"AWS":["b","a"]},"Action":["s3:GetObject","s3:PutObject"],"Resource":["arn:aws:s3:::bucket/b*","arn:aws:s3:::bucket/a*"],"Condition":{"StringLike":{"s3:prefix":["b*","a*"]}}},{"Sid":"two","Effect":"Deny","Principal":"*","NotAction":["s3:GetObject","s3:PutObject"],"NotResource":["arn:aws:s3:::bucket/d*","arn:aws:s3:::bucket/c*"]}]}`)
// Use a condition valid for object actions.
a = bytes.ReplaceAll(a, []byte("s3:prefix"), []byte("aws:UserAgent"))
b := []byte(`{"Statement":[{"NotResource":["arn:aws:s3:::bucket/c*","arn:aws:s3:::bucket/d*"],"NotAction":["s3:PutObject","s3:GetObject"],"Principal":"*","Effect":"Deny","Sid":"two"},{"Condition":{"StringLike":{"aws:UserAgent":["a*","b*"]}},"Resource":["arn:aws:s3:::bucket/a*","arn:aws:s3:::bucket/b*"],"Action":["s3:PutObject","s3:GetObject"],"Principal":{"AWS":["a","b"]},"Effect":"Allow","Sid":"one"}],"Version":"2012-10-17"}`)
_, key, err := bucketConfigPayload("bucket", bucketPolicyConfig, a, false)
if err != nil {
t.Fatal(err)
}
for i := range 30 {
_, next, err := bucketConfigPayload("bucket", bucketPolicyConfig, b, false)
if err != nil || !bytes.Equal(key, next) {
t.Fatalf("independent parse %d: %s != %s (%v)", i, key, next, err)
}
}
_, sid, err := bucketConfigPayload("bucket", bucketPolicyConfig, bytes.ReplaceAll(a, []byte(`"one"`), []byte(`"other"`)), false)
if err != nil || bytes.Equal(key, sid) {
t.Fatalf("Sid difference lost: %v", err)
}
n, err := canonicalBucketPolicyJSON([]byte(`{"n":[9007199254740993,9007199254740992]}`))
if err != nil || string(n) != `{"n":[9007199254740992,9007199254740993]}` {
t.Fatalf("integer precision lost: %s %v", n, err)
}
}
func TestPeerBucketMetadataWireAtomicity(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, cred auth.Credentials, t *testing.T) {
t.Run(backend, func(t *testing.T) {
ctx := t.Context()
meta, err := loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
stamp := meta.Created.Add(time.Hour)
policyData := []byte(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket))
meta.PolicyConfigJSON, meta.PolicyConfigUpdatedAt = policyData, stamp
meta.QuotaConfigJSON, meta.QuotaConfigUpdatedAt = []byte(`{"quota":1024,"quotatype":"hard"}`), stamp
meta.TaggingConfigXML, meta.TaggingConfigUpdatedAt = []byte(`<Tagging><TagSet><Tag><Key>k</Key><Value>a</Value></Tag></TagSet></Tagging>`), stamp
if err = globalBucketMetadataSys.save(ctx, meta); err != nil {
t.Fatal(err)
}
counter := &bucketConfigWriteCounter{ObjectLayer: obj}
setObjectLayer(counter)
defer setObjectLayer(obj)
read := func() BucketMetadata {
t.Helper()
m, err := loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
return m
}
apply := func(fragment string, at time.Time, ok bool) {
t.Helper()
raw := fmt.Sprintf(`{"bucket":%q,"updatedAt":%q,%s}`, bucket, at.Format(time.RFC3339Nano), fragment)
var item madmin.SRBucketMeta
if err := json.Unmarshal([]byte(raw), &item); err != nil {
t.Fatal(err)
}
rec := applySRBucketMetaViaAdmin(t, cred, item)
if (rec.Code == http.StatusOK) != ok {
t.Fatalf("%s: %d %s", raw, rec.Code, rec.Body.String())
}
}
// Omitted/null *string and empty update-only payloads are all no-ops.
apply(`"policy":null,"quota":null,"tags":null,"versioningConfig":"","objectLockConfig":""`, stamp.Add(time.Minute), true)
current := read()
if len(current.PolicyConfigJSON) != 0 || len(current.QuotaConfigJSON) == 0 {
t.Fatal("explicit null semantics lost")
}
if q, _, err := globalBucketMetadataSys.GetQuotaConfig(ctx, bucket); err != nil || q.Quota != 0 {
t.Fatalf("zero quota cache: %+v %v", q, err)
}
before := counter.writes.Load()
apply(`"tags":null,"versioningConfig":"","objectLockConfig":""`, stamp.Add(2*time.Minute), true)
if counter.writes.Load() != before {
t.Fatal("omitted fields or empty update-only wrote metadata")
}
// An invalid second field cannot persist the valid deletion in this bulk.
apply(`"tags":"","quota":{"quota":1,"quotatype":"invalid"}`, stamp.Add(2*time.Minute), false)
if counter.writes.Load() != before || len(read().TaggingConfigXML) == 0 {
t.Fatal("invalid bulk partially committed")
}
apply(`"tags":""`, stamp.Add(2*time.Minute), true)
before = counter.writes.Load()
apply(`"tags":""`, stamp.Add(2*time.Minute), true)
apply(`"tags":""`, stamp.Add(time.Minute), true)
if counter.writes.Load() != before {
t.Fatal("duplicate or old event wrote metadata")
}
localAt, err := globalBucketMetadataSys.Update(ctx, bucket, bucketQuotaConfigFile, []byte(`{}`))
if err != nil || !localAt.After(current.QuotaConfigUpdatedAt) {
t.Fatalf("local future monotonic time: %v %v", localAt, err)
}
deletedAt, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketQuotaConfigFile)
if err != nil || !deletedAt.After(localAt) {
t.Fatalf("adjacent local time: %v %v", deletedAt, err)
}
if q, _, err := globalBucketMetadataSys.GetQuotaConfig(ctx, bucket); err != nil || q.Quota != 0 {
t.Fatalf("deleted quota cache: %+v %v", q, err)
}
})
}})
}
func TestPeerBucketAdoptionRebasesOnlyDefaults(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) {
for _, shift := range []time.Duration{-time.Hour, 0, time.Hour} {
t.Run(fmt.Sprintf("%s/%s", backend, shift), func(t *testing.T) {
created := UTCNow().Add(-3 * time.Hour)
meta := newBucketMetadata(bucket)
meta.Created = created
meta.defaultTimestamps()
meta.QuotaConfigUpdatedAt = time.Time{}
meta.TaggingConfigUpdatedAt = created.Add(2 * time.Hour) // real deletion
meta.EncryptionConfigXML = []byte(`<ServerSideEncryptionConfiguration><Rule><ApplyServerSideEncryptionByDefault><SSEAlgorithm>AES256</SSEAlgorithm></ApplyServerSideEncryptionByDefault></Rule></ServerSideEncryptionConfiguration>`)
meta.EncryptionConfigUpdatedAt = created.Add(2 * time.Hour)
if err := globalBucketMetadataSys.save(t.Context(), meta); err != nil {
t.Fatal(err)
}
if err := globalSiteReplicationSys.PeerBucketMakeWithVersioningHandler(t.Context(), bucket, MakeBucketOptions{CreatedAt: created.Add(shift)}); err != nil {
t.Fatal(err)
}
got, err := readBucketMetadata(t.Context(), obj, bucket)
if err != nil {
t.Fatal(err)
}
if !got.PolicyConfigUpdatedAt.Equal(got.Created) || !got.QuotaConfigUpdatedAt.Equal(got.Created) {
t.Fatal("default turned into tombstone or invalid source")
}
if !got.TaggingConfigUpdatedAt.Equal(meta.TaggingConfigUpdatedAt) || !got.EncryptionConfigUpdatedAt.Equal(meta.EncryptionConfigUpdatedAt) || !bytes.Equal(got.EncryptionConfigXML, meta.EncryptionConfigXML) {
t.Fatal("actual state changed during adoption")
}
})
}
}})
}
func recordBucketConfigPeer(t *testing.T, cred auth.Credentials) func() []madmin.SRBucketMeta {
t.Helper()
var mu sync.Mutex
var events []madmin.SRBucketMeta
remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var event madmin.SRBucketMeta
if err := json.NewDecoder(r.Body).Decode(&event); err != nil {
t.Error(err)
w.WriteHeader(http.StatusBadRequest)
return
}
mu.Lock()
events = append(events, event)
mu.Unlock()
w.WriteHeader(http.StatusOK)
}))
t.Cleanup(remote.Close)
serviceCred, err := auth.CreateCredentials("metadata-source-time-svc", "metadata-source-time-service-secret")
if err != nil {
t.Fatal(err)
}
serviceCred.ParentUser = cred.AccessKey
globalSiteReplicatorCred.Set(serviceCred.SecretKey)
t.Cleanup(func() { globalSiteReplicatorCred.Set("") })
if _, err = globalIAMSys.store.AddServiceAccount(t.Context(), serviceCred); err != nil {
t.Fatal(err)
}
t.Cleanup(func() { globalIAMSys.DeleteServiceAccount(context.Background(), serviceCred.AccessKey, false) })
globalSiteReplicationSys.Lock()
oldEnabled, oldState := globalSiteReplicationSys.enabled, globalSiteReplicationSys.state
globalSiteReplicationSys.enabled = true
globalSiteReplicationSys.state = srState{Name: "metadata-source-test", ServiceAccountAccessKey: serviceCred.AccessKey, Peers: map[string]madmin.PeerInfo{
globalDeploymentID(): {Name: "local", DeploymentID: globalDeploymentID()},
"metadata-peer": {Name: "remote", DeploymentID: "metadata-peer", Endpoint: remote.URL},
}}
globalSiteReplicationSys.Unlock()
t.Cleanup(func() {
globalSiteReplicationSys.Lock()
globalSiteReplicationSys.enabled, globalSiteReplicationSys.state = oldEnabled, oldState
globalSiteReplicationSys.Unlock()
})
return func() []madmin.SRBucketMeta {
mu.Lock()
defer mu.Unlock()
return append([]madmin.SRBucketMeta(nil), events...)
}
}
// The final metadata lock is taken after reading the ZIP. Commit another
// writer immediately before that acquisition, using the existing lock wrapper.
type importInterleavingObjectLayer struct {
ObjectLayer
once atomic.Bool
write func()
}
func (o *importInterleavingObjectLayer) NewNSLock(bucket string, objects ...string) RWLocker {
lock := o.ObjectLayer.NewNSLock(bucket, objects...)
if bucket != minioMetaBucket || len(objects) != 1 || path.Base(objects[0]) != "metadata.lock" {
return lock
}
return metadataObservedRWLocker{RWLocker: lock, onLock: func(context.Context) {
if o.once.CompareAndSwap(false, true) {
o.write()
}
}}
}
func TestLocalBucketMetadataCommittedEvents(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, endpoints: []string{"PutBucketPolicy", "GetBucketPolicy", "PutBucketVersioning"}, objAPITest: func(obj ObjectLayer, backend, bucket string, router http.Handler, cred auth.Credentials, t *testing.T) {
t.Run(backend, func(t *testing.T) {
ctx := t.Context()
events := recordBucketConfigPeer(t, cred)
putPolicy := func(data []byte) {
t.Helper()
req, err := newTestSignedRequestV4(http.MethodPut, getPutPolicyURL("", bucket), int64(len(data)), bytes.NewReader(data), cred.AccessKey, cred.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
router.ServeHTTP(rec, req)
if rec.Code != http.StatusNoContent {
t.Fatalf("empty policy PUT: %d %s", rec.Code, rec.Body.String())
}
}
putPolicy([]byte(`{"Version":"2012-10-17","Statement":[]}`))
meta, err := loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
got := events()
if len(got) != 1 || got[0].Type != madmin.SRBucketMetaTypePolicy || len(got[0].Policy) != 0 || len(meta.PolicyConfigJSON) != 0 || !meta.PolicyConfigUpdatedAt.Equal(got[0].UpdatedAt) {
t.Fatalf("empty policy event/disk mismatch: %+v", got)
}
req, err := newTestSignedRequestV4(http.MethodGet, getGetPolicyURL("", bucket), 0, nil, cred.AccessKey, cred.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
router.ServeHTTP(rec, req)
if rec.Code != http.StatusNotFound {
t.Fatalf("empty policy GET: %d %s", rec.Code, rec.Body.String())
}
corsAdminRequest(t, cred, http.MethodPut, "/set-bucket-quota?bucket="+bucket, []byte(`{}`))
meta, err = loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
got = events()
last := got[len(got)-1]
if last.Type != madmin.SRBucketMetaTypeQuotaConfig || len(last.Quota) == 0 || !last.UpdatedAt.Equal(meta.QuotaConfigUpdatedAt) {
t.Fatalf("zero quota event: %+v", last)
}
injectedAt := UTCNow().Add(2 * time.Hour)
interleaving := &importInterleavingObjectLayer{ObjectLayer: obj, write: func() {
data := []byte(`{"quota":2048,"quotatype":"hard"}`)
_, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketQuotaConfigFile, data, false, false, &injectedAt)
if err != nil {
t.Error(err)
}
}}
setObjectLayer(interleaving)
defer setObjectLayer(obj)
before := len(events())
rec = corsAdminRequest(t, cred, http.MethodPut, "/import-bucket-metadata", corsZip(t, map[string][]byte{
bucket + "/" + bucketPolicyConfig: []byte(`{"Version":"2012-10-17","Statement":[]}`),
bucket + "/quota.json": []byte(`{}`),
bucket + "/" + bucketTaggingConfig: []byte(`<Tagging><TagSet><Tag><Key>imported</Key><Value>yes</Value></Tag></TagSet></Tagging>`),
}))
report := corsImportReport(t, rec).Buckets[bucket]
if report.Err != "" || report.Policy.Err != "" || report.Quota.Err != "" || report.Tagging.Err != "" {
t.Fatalf("import: %+v", report)
}
meta, err = loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
if !meta.QuotaConfigUpdatedAt.After(injectedAt) || !meta.TaggingConfigUpdatedAt.Equal(meta.QuotaConfigUpdatedAt) || !meta.PolicyConfigUpdatedAt.Equal(meta.QuotaConfigUpdatedAt) {
t.Fatalf("import time %v must exceed intervening %v; policy=%v tags=%v", meta.QuotaConfigUpdatedAt, injectedAt, meta.PolicyConfigUpdatedAt, meta.TaggingConfigUpdatedAt)
}
got = events()[before:]
if len(got) != 2 {
t.Fatalf("import events: %+v", got)
}
for _, event := range got {
if !event.UpdatedAt.Equal(meta.QuotaConfigUpdatedAt) {
t.Fatal("import hook used ZIP start time")
}
if event.Type == madmin.SRBucketMetaTypePolicy {
if len(event.Policy) != 0 {
t.Fatal("empty policy import is not deletion")
}
} else if event.Tags == nil || len(event.Quota) == 0 {
t.Fatalf("bulk omitted imported state: %+v", event)
}
}
})
}})
}
func TestPeerBucketMetadataNormalizesBeforeComparison(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, cred auth.Credentials, t *testing.T) {
t.Run(backend, func(t *testing.T) {
meta, err := loadBucketMetadata(t.Context(), obj, bucket)
if err != nil {
t.Fatal(err)
}
versioning := base64.StdEncoding.EncodeToString([]byte(`<VersioningConfiguration><Status>Enabled</Status><ExcludedPrefixes><Prefix>skip/</Prefix></ExcludedPrefixes><ExcludeFolders>true</ExcludeFolders></VersioningConfiguration>`))
lock := base64.StdEncoding.EncodeToString(enabledBucketObjectLockConfig)
item := madmin.SRBucketMeta{Bucket: bucket, ObjectLockConfig: &lock, Versioning: &versioning, UpdatedAt: meta.Created.Add(time.Hour)}
counter := &bucketConfigWriteCounter{ObjectLayer: obj}
setObjectLayer(counter)
defer setObjectLayer(obj)
for range 2 {
rec := applySRBucketMetaViaAdmin(t, cred, item)
if rec.Code != http.StatusOK {
t.Fatalf("bulk: %d %s", rec.Code, rec.Body.String())
}
}
if counter.writes.Load() != 1 {
t.Fatalf("normalized duplicate persisted %d times", counter.writes.Load())
}
meta, err = loadBucketMetadata(t.Context(), obj, bucket)
if err != nil || !bytes.Equal(meta.VersioningConfigXML, enabledBucketVersioningConfig) || !meta.VersioningConfigUpdatedAt.Equal(item.UpdatedAt) {
t.Fatalf("normalization: %s %v", meta.VersioningConfigXML, err)
}
// Ordinary local updates use the identical effective document in their
// commit result, even if Object Lock changed since the handler's precheck.
data, _ := base64.StdEncoding.DecodeString(versioning)
result, err := globalBucketMetadataSys.updateAndParseMetadata(t.Context(), bucket, bucketVersioningConfig, data, false, false, nil)
if err != nil || !bytes.Equal(result.meta.VersioningConfigXML, enabledBucketVersioningConfig) {
t.Fatalf("local commit snapshot: %s %v", result.meta.VersioningConfigXML, err)
}
})
}})
}
+74 -141
View File
@@ -906,7 +906,7 @@ func enablePeerBucketVersioning(meta *BucketMetadata, lockEnabled bool) error {
config, err := versioning.ParseConfig(bytes.NewReader(meta.VersioningConfigXML))
if err != nil || (lockEnabled && (config.Suspended() || config.PrefixesExcluded())) {
meta.VersioningConfigXML = enabledBucketVersioningConfig
meta.VersioningConfigUpdatedAt = UTCNow()
meta.VersioningConfigUpdatedAt = localBucketConfigUpdatedAt(*meta, bucketVersioningConfig, UTCNow())
return nil
}
if config.Enabled() {
@@ -917,7 +917,7 @@ func enablePeerBucketVersioning(meta *BucketMetadata, lockEnabled bool) error {
if err != nil {
return err
}
meta.VersioningConfigUpdatedAt = UTCNow()
meta.VersioningConfigUpdatedAt = localBucketConfigUpdatedAt(*meta, bucketVersioningConfig, UTCNow())
return nil
}
@@ -947,7 +947,9 @@ func (c *SiteReplicationSys) PeerBucketMakeWithVersioningHandler(ctx context.Con
if err != nil {
return err
}
oldCreated := meta.Created
meta.SetCreatedAt(opts.CreatedAt)
rebaseBucketConfigDefaults(&meta, oldCreated)
if err = enablePeerBucketVersioning(&meta, opts.LockEnabled || len(meta.ObjectLockConfigXML) != 0); err != nil {
return err
@@ -958,7 +960,7 @@ func (c *SiteReplicationSys) PeerBucketMakeWithVersioningHandler(ctx context.Con
meta.ObjectLockConfigUpdatedAt = meta.Created
}
}
return globalBucketMetadataSys.saveMetadata(bgContext(ctx), objAPI, meta)
return globalBucketMetadataSys.saveMetadata(bgContext(ctx), objAPI, &meta)
}()
if err != nil {
return wrapSRErr(c.annotateErr(makeBucketWithVersion, err))
@@ -1583,24 +1585,18 @@ func (c *SiteReplicationSys) BucketMetaHook(ctx context.Context, item madmin.SRB
// PeerBucketVersioningHandler - updates versioning config to local cluster.
func (c *SiteReplicationSys) PeerBucketVersioningHandler(ctx context.Context, bucket string, versioning *string, updatedAt time.Time) error {
var data []byte
if versioning != nil {
// skip overwrite if local update is newer than peer update.
if !updatedAt.IsZero() {
if _, updateTm, err := globalBucketMetadataSys.GetVersioningConfig(bucket); err == nil && updateTm.After(updatedAt) {
return nil
}
}
configData, err := base64.StdEncoding.DecodeString(*versioning)
var err error
data, err = base64.StdEncoding.DecodeString(*versioning)
if err != nil {
return wrapSRErr(err)
}
_, err = globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, configData)
if err != nil {
return wrapSRErr(err)
}
return nil
}
_, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketVersioningConfig, data, false, false, &updatedAt)
if err != nil {
return wrapSRErr(err)
}
return nil
}
@@ -1650,50 +1646,43 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context
return nil
}
if item.Policy != nil {
meta.PolicyConfigJSON = item.Policy
meta.PolicyConfigUpdatedAt = item.UpdatedAt
// Presence is separate from content: omitted RawMessage is not a delete;
// explicit JSON null still goes through the existing policy/quota parser.
updates := make(map[string][]byte, len(replicatedBucketConfigs))
if len(item.Policy) != 0 {
updates[bucketPolicyConfig] = item.Policy
}
if item.Versioning != nil {
configData, err := base64.StdEncoding.DecodeString(*item.Versioning)
if len(item.Quota) != 0 {
updates[bucketQuotaConfigFile] = item.Quota
}
for file, payload := range map[string]*string{
objectLockConfig: item.ObjectLockConfig, bucketVersioningConfig: item.Versioning,
bucketTaggingConfig: item.Tags, bucketSSEConfig: item.SSEConfig,
} {
if payload != nil {
data, err := base64.StdEncoding.DecodeString(*payload)
if err != nil {
return wrapSRErr(err)
}
updates[file] = data
}
}
if len(updates) != 0 {
if err := ensureBucketMetadataCreated(ctx, objectAPI, &meta); err != nil {
return wrapSRErr(err)
}
}
changed := false
for _, file := range replicatedBucketConfigs {
data, supplied := updates[file]
if !supplied {
continue
}
applied, err := applyBucketConfig(&meta, file, data, item.UpdatedAt)
if err != nil {
return wrapSRErr(err)
}
meta.VersioningConfigXML = configData
meta.VersioningConfigUpdatedAt = item.UpdatedAt
}
if item.Tags != nil {
configData, err := base64.StdEncoding.DecodeString(*item.Tags)
if err != nil {
return wrapSRErr(err)
}
meta.TaggingConfigXML = configData
meta.TaggingConfigUpdatedAt = item.UpdatedAt
}
if item.ObjectLockConfig != nil {
configData, err := base64.StdEncoding.DecodeString(*item.ObjectLockConfig)
if err != nil {
return wrapSRErr(err)
}
meta.ObjectLockConfigXML = configData
meta.ObjectLockConfigUpdatedAt = item.UpdatedAt
}
if item.SSEConfig != nil {
configData, err := base64.StdEncoding.DecodeString(*item.SSEConfig)
if err != nil {
return wrapSRErr(err)
}
meta.EncryptionConfigXML = configData
meta.EncryptionConfigUpdatedAt = item.UpdatedAt
}
if item.Quota != nil {
meta.QuotaConfigJSON = item.Quota
meta.QuotaConfigUpdatedAt = item.UpdatedAt
changed = changed || applied
}
if item.Cors != nil {
@@ -1702,10 +1691,14 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context
if compareCORSReplicationStates(localState, incoming) < 0 {
meta.CorsConfigXML = bytes.Clone(corsConfigData)
meta.CorsConfigUpdatedAt = item.UpdatedAt
changed = true
}
}
if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, meta); err != nil {
if !changed {
return nil
}
if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, &meta); err != nil {
return err
}
unlock()
@@ -1716,62 +1709,35 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context
// PeerBucketPolicyHandler - copies/deletes policy to local cluster.
func (c *SiteReplicationSys) PeerBucketPolicyHandler(ctx context.Context, bucket string, policy *policy.BucketPolicy, updatedAt time.Time) error {
// skip overwrite if local update is newer than peer update.
if !updatedAt.IsZero() {
if _, updateTm, err := globalBucketMetadataSys.GetPolicyConfig(bucket); err == nil && updateTm.After(updatedAt) {
return nil
}
}
var data []byte
if policy != nil {
configData, err := json.Marshal(policy)
var err error
data, err = canonicalBucketPolicy(policy)
if err != nil {
return wrapSRErr(err)
}
_, err = globalBucketMetadataSys.Update(ctx, bucket, bucketPolicyConfig, configData)
if err != nil {
return wrapSRErr(err)
}
return nil
}
// Delete the bucket policy
_, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketPolicyConfig)
_, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketPolicyConfig, data, false, false, &updatedAt)
if err != nil {
return wrapSRErr(err)
}
return nil
}
// PeerBucketTaggingHandler - copies/deletes tags to local cluster.
func (c *SiteReplicationSys) PeerBucketTaggingHandler(ctx context.Context, bucket string, tags *string, updatedAt time.Time) error {
// skip overwrite if local update is newer than peer update.
if !updatedAt.IsZero() {
if _, updateTm, err := globalBucketMetadataSys.GetTaggingConfig(bucket); err == nil && updateTm.After(updatedAt) {
return nil
}
}
var data []byte
if tags != nil {
configData, err := base64.StdEncoding.DecodeString(*tags)
var err error
data, err = base64.StdEncoding.DecodeString(*tags)
if err != nil {
return wrapSRErr(err)
}
_, err = globalBucketMetadataSys.Update(ctx, bucket, bucketTaggingConfig, configData)
if err != nil {
return wrapSRErr(err)
}
return nil
}
// Delete the tags
_, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketTaggingConfig)
_, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketTaggingConfig, data, false, false, &updatedAt)
if err != nil {
return wrapSRErr(err)
}
return nil
}
@@ -1797,51 +1763,32 @@ func (c *SiteReplicationSys) peerBucketObjectLockConfigItem(ctx context.Context,
// PeerBucketObjectLockConfigHandler - sets object lock on local bucket.
func (c *SiteReplicationSys) PeerBucketObjectLockConfigHandler(ctx context.Context, bucket string, objectLockData *string, updatedAt time.Time) error {
var data []byte
if objectLockData != nil {
// skip overwrite if local update is newer than peer update.
if !updatedAt.IsZero() {
if _, updateTm, err := globalBucketMetadataSys.GetObjectLockConfig(bucket); err == nil && updateTm.After(updatedAt) {
return nil
}
}
configData, err := base64.StdEncoding.DecodeString(*objectLockData)
var err error
data, err = base64.StdEncoding.DecodeString(*objectLockData)
if err != nil {
return wrapSRErr(err)
}
_, err = globalBucketMetadataSys.Update(ctx, bucket, objectLockConfig, configData)
if err != nil {
return wrapSRErr(err)
}
return nil
}
_, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, objectLockConfig, data, false, false, &updatedAt)
if err != nil {
return wrapSRErr(err)
}
return nil
}
// PeerBucketSSEConfigHandler - copies/deletes SSE config to local cluster.
func (c *SiteReplicationSys) PeerBucketSSEConfigHandler(ctx context.Context, bucket string, sseConfig *string, updatedAt time.Time) error {
// skip overwrite if local update is newer than peer update.
if !updatedAt.IsZero() {
if _, updateTm, err := globalBucketMetadataSys.GetSSEConfig(bucket); err == nil && updateTm.After(updatedAt) {
return nil
}
}
var data []byte
if sseConfig != nil {
configData, err := base64.StdEncoding.DecodeString(*sseConfig)
var err error
data, err = base64.StdEncoding.DecodeString(*sseConfig)
if err != nil {
return wrapSRErr(err)
}
_, err = globalBucketMetadataSys.Update(ctx, bucket, bucketSSEConfig, configData)
if err != nil {
return wrapSRErr(err)
}
return nil
}
// Delete sse config
_, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketSSEConfig)
_, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketSSEConfig, data, false, false, &updatedAt)
if err != nil {
return wrapSRErr(err)
}
@@ -2046,7 +1993,7 @@ func applyBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket
meta.CorsConfigXML = bytes.Clone(configData)
meta.CorsConfigUpdatedAt = updatedAt
if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, meta); err != nil {
if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, &meta); err != nil {
return time.Time{}, err
}
unlock()
@@ -2078,32 +2025,18 @@ func (c *SiteReplicationSys) PeerBucketCorsConfigHandler(ctx context.Context, bu
// PeerBucketQuotaConfigHandler - copies/deletes policy to local cluster.
func (c *SiteReplicationSys) PeerBucketQuotaConfigHandler(ctx context.Context, bucket string, quota *madmin.BucketQuota, updatedAt time.Time) error {
// skip overwrite if local update is newer than peer update.
if !updatedAt.IsZero() {
if _, updateTm, err := globalBucketMetadataSys.GetQuotaConfig(ctx, bucket); err == nil && updateTm.After(updatedAt) {
return nil
}
}
var data []byte
if quota != nil {
quotaData, err := json.Marshal(quota)
var err error
data, err = json.Marshal(quota)
if err != nil {
return wrapSRErr(err)
}
if _, err = globalBucketMetadataSys.Update(ctx, bucket, bucketQuotaConfigFile, quotaData); err != nil {
return wrapSRErr(err)
}
return nil
}
// Delete the bucket policy
_, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketQuotaConfigFile)
_, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketQuotaConfigFile, data, false, false, &updatedAt)
if err != nil {
return wrapSRErr(err)
}
return nil
}