fix(replication): gate metadata tombstone export during rolling upgrades

Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
Feng Ruohang
2026-09-12 11:40:04 +08:00
parent 01aaef2b50
commit 1ee64a8d89
12 changed files with 435 additions and 86 deletions
+1 -1
View File
@@ -435,7 +435,7 @@ func (a adminAPIHandlers) ExportBucketMetadataHandler(w http.ResponseWriter, r *
writeErrorResponse(ctx, w, exportError(ctx, err, cfgFile, bucket), r.URL)
return
}
configData, err := json.Marshal(config)
configData, err := canonicalBucketPolicy(config)
if err != nil {
writeErrorResponse(ctx, w, exportError(ctx, err, cfgFile, bucket), r.URL)
return
+3 -3
View File
@@ -213,11 +213,11 @@ func newBucketConfigState(bucket, file string, data []byte, at, created time.Tim
at = created
}
valid := !created.IsZero() && !at.Before(created)
real := valid && at.After(created)
modified := valid && at.After(created)
if bucketConfigUpdateOnly(file) && len(data) == 0 {
real = false
modified = false
}
return bucketConfigState{data: data, key: key, at: at, real: real, valid: valid}, nil
return bucketConfigState{data: data, key: key, at: at, real: modified, valid: valid}, nil
}
func (s bucketConfigState) candidate() bool {
+8
View File
@@ -180,13 +180,21 @@ func (sys *BucketMetadataSys) updateAndParseMetadata(ctx context.Context, bucket
updatedAt := UTCNow()
if data, _ := replicatedBucketConfig(&meta, configFile); data != nil {
if err := ensureBucketMetadataCreated(ctx, objAPI, &meta); err != nil {
logBucketConfigReplication(ctx, bucket, configFile, "indeterminate", time.Time{}, meta.Created, err.Error())
return err
}
if sourceTime == nil || sourceTime.IsZero() {
updatedAt = localBucketConfigUpdatedAt(meta, configFile, updatedAt)
if sourceTime != nil {
logBucketConfigReplication(ctx, bucket, configFile, "legacy-zero", *sourceTime, meta.Created, "assigned local source time")
}
} else {
updatedAt = sourceTime.UTC()
}
if updatedAt.Before(meta.Created) {
logBucketConfigReplication(ctx, bucket, configFile, "before-created", updatedAt, meta.Created, "peer event")
return nil
}
changed, err := applyBucketConfig(&meta, configFile, configData, updatedAt)
if err != nil {
return err
+1 -2
View File
@@ -19,7 +19,6 @@ package cmd
import (
"bytes"
"encoding/json"
"io"
"net/http"
@@ -200,7 +199,7 @@ func (api objectAPIHandlers) GetBucketPolicyHandler(w http.ResponseWriter, r *ht
return
}
configData, err := json.Marshal(config)
configData, err := canonicalBucketPolicy(config)
if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return
+1
View File
@@ -900,6 +900,7 @@ func serverHandleEnvVars() {
}
globalEnableSyncBoot = env.Get("MINIO_SYNC_BOOT", config.EnableOff) == config.EnableOn
globalSiteReplicationMetadataTombstones = env.Get("MINIO_SITE_REPLICATION_METADATA_TOMBSTONES", config.EnableOff) == config.EnableOn
}
func loadRootCredentials() auth.Credentials {
+199
View File
@@ -0,0 +1,199 @@
// 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"
"net/http"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
)
func TestBucketMetadataTombstoneExportAndInitialSync(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) {
recordBucketConfigPeer(t, cred)
old := globalSiteReplicationMetadataTombstones
defer func() { globalSiteReplicationMetadataTombstones = old }()
created := UTCNow().Add(-time.Hour)
deletedAt := created.Add(time.Minute)
for _, enabled := range []bool{false, true} {
globalSiteReplicationMetadataTombstones = enabled
meta := newBucketMetadata(bucket)
meta.Created = created
meta.defaultTimestamps()
for _, file := range []string{bucketPolicyConfig, bucketTaggingConfig, bucketSSEConfig, bucketQuotaConfigFile} {
_, at := replicatedBucketConfig(&meta, file)
*at = deletedAt
}
if err := globalBucketMetadataSys.save(t.Context(), meta); err != nil {
t.Fatal(err)
}
// Force a disk reload instead of accepting the just-published cache.
globalBucketMetadataSys.Remove(bucket)
info, err := globalSiteReplicationSys.SiteReplicationMetaInfo(t.Context(), obj, madmin.SRStatusOptions{Buckets: true})
if err != nil {
t.Fatal(err)
}
exported := info.Buckets[bucket]
if !exported.PolicyUpdatedAt.Equal(deletedAt) {
t.Fatal("Policy tombstone hidden by gate")
}
for _, at := range []time.Time{exported.TagConfigUpdatedAt, exported.SSEConfigUpdatedAt, exported.QuotaConfigUpdatedAt} {
if enabled && !at.Equal(deletedAt) || !enabled && !at.IsZero() {
t.Fatalf("gate=%v timestamp=%v", enabled, at)
}
}
for file, data := range bucketConfigTestData(bucket) {
event, send, err := initialBucketConfigReplicationEvent(meta, file)
wantSend := enabled && !bucketConfigUpdateOnly(file)
if err != nil || send != wantSend || send && !event.UpdatedAt.Equal(deletedAt) {
t.Fatalf("%s initial gate=%v: %+v %v %v", file, enabled, event, send, err)
}
baseline := newBucketMetadata(bucket)
baseline.Created = created
baseline.defaultTimestamps()
if _, send, err := initialBucketConfigReplicationEvent(baseline, file); err != nil || send {
t.Fatalf("empty baseline sent: %s", file)
}
value, _ := replicatedBucketConfig(&baseline, file)
*value = data
event, send, err = initialBucketConfigReplicationEvent(baseline, file)
if err != nil || !send || !event.UpdatedAt.Equal(created) {
t.Fatalf("historical baseline-live omitted: %s %v", file, err)
}
}
}
})
}})
}
func TestPeerBucketMetadataLegacyAndGeneration(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) {
old := globalSiteReplicationMetadataTombstones
defer func() { globalSiteReplicationMetadataTombstones = old }()
created := UTCNow().Add(-time.Hour)
for _, enabled := range []bool{false, true} {
globalSiteReplicationMetadataTombstones = enabled
for file, data := range bucketConfigTestData(bucket) {
meta := newBucketMetadata(bucket)
meta.Created = created
meta.defaultTimestamps()
if err := globalBucketMetadataSys.save(t.Context(), meta); err != nil {
t.Fatal(err)
}
counter := &bucketConfigWriteCounter{ObjectLayer: obj}
setObjectLayer(counter)
event := newBucketConfigReplicationEvent(bucket, file, bucketConfigState{data: data, at: created.Add(-time.Second)})
rec := applySRBucketMetaViaAdmin(t, cred, event)
if rec.Code != http.StatusOK || counter.writes.Load() != 0 {
t.Fatalf("pre-creation apply: %s %d writes=%d", file, rec.Code, counter.writes.Load())
}
// A live baseline may initialize. The same-time nil cannot delete it.
event.UpdatedAt = created
rec = applySRBucketMetaViaAdmin(t, cred, event)
if rec.Code != http.StatusOK {
t.Fatalf("baseline apply: %s %s", file, rec.Body.String())
}
before := counter.writes.Load()
rec = applySRBucketMetaViaAdmin(t, cred, madmin.SRBucketMeta{Type: event.Type, Bucket: bucket, UpdatedAt: created})
if rec.Code != http.StatusOK || counter.writes.Load() != before {
t.Fatalf("nil baseline cleared configuration: %s", file)
}
event.UpdatedAt = time.Time{}
for range 2 {
rec = applySRBucketMetaViaAdmin(t, cred, event)
if rec.Code != http.StatusOK {
t.Fatalf("legacy-zero rejected: %s %s", file, rec.Body.String())
}
}
meta, err := loadBucketMetadata(t.Context(), obj, bucket)
if err != nil {
t.Fatal(err)
}
value, at := replicatedBucketConfig(&meta, file)
if len(*value) == 0 || !at.After(created) {
t.Fatalf("legacy-zero not reclocked: %s %v", file, *at)
}
setObjectLayer(obj)
}
}
})
}})
}
type bucketMetadataCreatedObjectLayer struct {
ObjectLayer
created time.Time
missing bool
}
func (o bucketMetadataCreatedObjectLayer) GetBucketInfo(ctx context.Context, bucket string, opts BucketOptions) (BucketInfo, error) {
if opts.NoMetadata {
if o.missing {
return BucketInfo{}, BucketNotFound{Bucket: bucket}
}
return BucketInfo{Name: bucket, Created: o.created}, nil
}
return o.ObjectLayer.GetBucketInfo(ctx, bucket, opts)
}
func TestPeerBucketMetadataUnknownCreated(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) {
defer setObjectLayer(obj)
data := bucketConfigTestData(bucket)[bucketTaggingConfig]
for _, mode := range []string{"unknown", "missing", "physical-created"} {
t.Run(mode, func(t *testing.T) {
meta := newBucketMetadata(bucket)
setObjectLayer(obj)
if err := globalBucketMetadataSys.save(t.Context(), meta); err != nil {
t.Fatal(err)
}
created := UTCNow().Add(-time.Hour)
physical := bucketMetadataCreatedObjectLayer{ObjectLayer: obj, missing: mode == "missing"}
if mode == "physical-created" {
physical.created = created
}
counter := &bucketConfigWriteCounter{ObjectLayer: physical}
setObjectLayer(counter)
stamp := created.Add(time.Minute)
_, err := globalBucketMetadataSys.updateAndParseMetadata(t.Context(), bucket, bucketTaggingConfig, data, false, false, &stamp)
if mode != "physical-created" {
if err == nil || counter.writes.Load() != 0 {
t.Fatalf("unknown generation was invented: %v writes=%d", err, counter.writes.Load())
}
} else {
if err != nil {
t.Fatal(err)
}
got, err := readBucketMetadata(t.Context(), obj, bucket)
if err != nil || !got.Created.Equal(created) || !got.TaggingConfigUpdatedAt.Equal(stamp) || !bytes.Equal(got.TaggingConfigXML, data) {
t.Fatalf("physical creation recovery: %v", err)
}
}
})
}
})
}})
}
+3 -3
View File
@@ -83,9 +83,9 @@ func TestLatestBucketConfigCandidates(t *testing.T) {
t.Fatal("unknown source chosen")
}
liveBaseline := bucketConfigTestInfo("bucket", file, data, created, created)
real := bucketConfigTestInfo("bucket", file, data, created.Add(time.Hour), created)
modified := bucketConfigTestInfo("bucket", file, data, created.Add(time.Hour), created)
oldGeneration := bucketConfigTestInfo("bucket", file, data, created, created.Add(time.Second))
states := []srBucketStatsSummary{liveBaseline, real, oldGeneration}
states := []srBucketStatsSummary{liveBaseline, modified, oldGeneration}
for _, order := range [][3]int{{0, 1, 2}, {0, 2, 1}, {1, 0, 2}, {1, 2, 0}, {2, 0, 1}, {2, 1, 0}} {
bs["a"], bs["b"], bs["c"] = states[order[0]], states[order[1]], states[order[2]]
winner, found := latestBucketConfig("bucket", file, info)
@@ -99,7 +99,7 @@ func TestLatestBucketConfigCandidates(t *testing.T) {
}
if !bucketConfigUpdateOnly(file) {
// A late-created empty baseline cannot beat a real earlier write.
bs["a"], bs["b"], bs["c"] = real, bucketConfigTestInfo("bucket", file, nil, created.Add(2*time.Hour), created.Add(2*time.Hour)), baseline
bs["a"], bs["b"], bs["c"] = modified, bucketConfigTestInfo("bucket", file, nil, created.Add(2*time.Hour), created.Add(2*time.Hour)), baseline
if winner, found := latestBucketConfig("bucket", file, info); !found || len(winner.data) == 0 {
t.Fatal("default became a deletion")
}
+47 -3
View File
@@ -20,12 +20,40 @@ package cmd
import (
"context"
"encoding/base64"
"fmt"
"errors"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/logger"
)
// Read once at startup. Enable only after every participating node is fixed.
var globalSiteReplicationMetadataTombstones bool
func logBucketConfigReplication(ctx context.Context, bucket, file, reason string, at, created time.Time, detail string) {
// LogOnceIf compares error text as well as its key. Keep both stable; changing
// times and peer errors belong in ReqInfo, not in the error's message.
req := &logger.ReqInfo{API: "SiteReplicationMetadata", BucketName: bucket}
req.AppendTags("field", file)
req.AppendTags("sourceTime", at.UTC().Format(time.RFC3339Nano))
req.AppendTags("created", created.UTC().Format(time.RFC3339Nano))
req.AppendTags("detail", detail)
replLogOnceIf(logger.SetReqInfo(ctx, req), errors.New("bucket metadata replication: "+reason),
"bucket-metadata/"+bucket+"/"+file+"/"+reason)
}
func initialBucketConfigReplicationEvent(meta BucketMetadata, file string) (madmin.SRBucketMeta, bool, error) {
data, at := replicatedBucketConfig(&meta, file)
state, err := newBucketConfigState(meta.Name, file, *data, *at, meta.Created, len(meta.ObjectLockConfigXML) != 0)
if err != nil {
return madmin.SRBucketMeta{}, false, err
}
if !state.candidate() || (len(state.data) == 0 && !globalSiteReplicationMetadataTombstones) {
return madmin.SRBucketMeta{}, false, nil
}
return newBucketConfigReplicationEvent(meta.Name, file, state), true, nil
}
func bucketConfigStateFromInfo(bucket, file string, meta madmin.SRBucketInfo) (bucketConfigState, error) {
var payload *string
var at time.Time
@@ -103,6 +131,18 @@ func (c *SiteReplicationSys) healBucketConfig(ctx context.Context, bucket, file
if !c.enabled {
return nil
}
for id := range info.Sites {
if _, present := info.BucketStats[bucket][id]; !present {
logBucketConfigReplication(ctx, bucket, file, "indeterminate", time.Time{}, time.Time{}, "missing peer "+id)
}
}
for id, status := range info.BucketStats[bucket] {
state, err := bucketConfigStateFromInfo(bucket, file, status.meta.SRBucketInfo)
_, known := info.Sites[id]
if !known || id == "" || err != nil || !state.valid {
logBucketConfigReplication(ctx, bucket, file, "indeterminate", state.at, status.meta.CreatedAt, "unusable peer "+id)
}
}
latest, found := latestBucketConfig(bucket, file, info)
if !found {
return nil
@@ -112,7 +152,11 @@ func (c *SiteReplicationSys) healBucketConfig(ctx context.Context, bucket, file
continue
}
target := status.meta.SRBucketInfo
if target.CreatedAt.IsZero() || latest.at.Before(target.CreatedAt) {
if target.CreatedAt.IsZero() {
continue
}
if latest.at.Before(target.CreatedAt) {
logBucketConfigReplication(ctx, bucket, file, "before-created", latest.at, target.CreatedAt, "peer "+id)
continue
}
// Versioning can be normalized differently until Object Lock itself has
@@ -138,7 +182,7 @@ func (c *SiteReplicationSys) healBucketConfig(ctx context.Context, bucket, file
if err != nil {
// A missing credential or unreachable peer must not abandon the other
// targets simply because it happened to be visited first in this map.
replLogIf(ctx, fmt.Errorf("unable to heal bucket %s %s for peer %s: %w", bucket, file, info.Sites[id].Name, err))
logBucketConfigReplication(ctx, bucket, file, "indeterminate", latest.at, target.CreatedAt, "peer "+id+": "+err.Error())
}
}
return nil
+94 -10
View File
@@ -57,15 +57,14 @@ func testPeerBucketMetadataSourceTimeAndDeletion(obj ObjectLayer, backend, bucke
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},
{"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 }, 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 }, 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 }, 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 }, 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 }, 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 }, false},
}
for _, tc := range cases {
t.Run(backend+"/"+tc.name, func(t *testing.T) {
@@ -528,6 +527,19 @@ func TestLocalBucketMetadataCommittedEvents(t *testing.T) {
if rec.Code != http.StatusNotFound {
t.Fatalf("empty policy GET: %d %s", rec.Code, rec.Body.String())
}
// Supported negative sets must survive the same serializer on PUT,
// peer apply, GET and export instead of failing on empty Action.
putPolicy([]byte(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","NotAction":["s3:DeleteObject"],"NotResource":["arn:aws:s3:::%s/private/*"]}]}`, bucket)))
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.StatusOK || !bytes.Contains(rec.Body.Bytes(), []byte(`"NotAction"`)) || !bytes.Contains(rec.Body.Bytes(), []byte(`"NotResource"`)) {
t.Fatalf("negative policy set 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 {
@@ -551,9 +563,11 @@ func TestLocalBucketMetadataCommittedEvents(t *testing.T) {
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>`),
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>`),
bucket + "/" + objectLockConfig: enabledBucketObjectLockConfig,
bucket + "/" + bucketVersioningConfig: []byte(`<VersioningConfiguration><Status>Enabled</Status><ExcludeFolders>true</ExcludeFolders></VersioningConfiguration>`),
}))
report := corsImportReport(t, rec).Buckets[bucket]
if report.Err != "" || report.Policy.Err != "" || report.Quota.Err != "" || report.Tagging.Err != "" {
@@ -566,6 +580,9 @@ func TestLocalBucketMetadataCommittedEvents(t *testing.T) {
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)
}
if !bytes.Equal(meta.VersioningConfigXML, enabledBucketVersioningConfig) || !meta.VersioningConfigUpdatedAt.Equal(meta.QuotaConfigUpdatedAt) || !meta.ObjectLockConfigUpdatedAt.Equal(meta.QuotaConfigUpdatedAt) {
t.Fatal("import normalization or common source time lost")
}
got = events()[before:]
if len(got) != 2 {
t.Fatalf("import events: %+v", got)
@@ -581,6 +598,33 @@ func TestLocalBucketMetadataCommittedEvents(t *testing.T) {
} else if event.Tags == nil || len(event.Quota) == 0 {
t.Fatalf("bulk omitted imported state: %+v", event)
}
if event.Type == "" {
if event.Versioning == nil || event.ObjectLockConfig == nil {
t.Fatal("normalized import fields omitted")
}
payload, err := base64.StdEncoding.DecodeString(*event.Versioning)
if err != nil || !bytes.Equal(payload, enabledBucketVersioningConfig) {
t.Fatal("import sent pre-normalization versioning")
}
}
}
beforeMeta, beforeEvents := meta, len(events())
rec = corsAdminRequest(t, cred, http.MethodPut, "/import-bucket-metadata", corsZip(t, map[string][]byte{
bucket + "/" + bucketTaggingConfig: []byte(`<Tagging><TagSet><Tag><Key>only</Key><Value>tags</Value></Tag></TagSet></Tagging>`),
}))
if status := corsImportReport(t, rec).Buckets[bucket]; status.Err != "" || status.Tagging.Err != "" {
t.Fatalf("tags-only import: %+v", status)
}
meta, err = loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
if !bytes.Equal(meta.PolicyConfigJSON, beforeMeta.PolicyConfigJSON) || !meta.PolicyConfigUpdatedAt.Equal(beforeMeta.PolicyConfigUpdatedAt) || !bytes.Equal(meta.QuotaConfigJSON, beforeMeta.QuotaConfigJSON) || !meta.QuotaConfigUpdatedAt.Equal(beforeMeta.QuotaConfigUpdatedAt) {
t.Fatal("tags-only import changed Policy or Quota")
}
got = events()[beforeEvents:]
if len(got) != 1 || got[0].Tags == nil || got[0].Policy != nil || got[0].Quota != nil || got[0].Versioning != nil || got[0].ObjectLockConfig != nil || !got[0].UpdatedAt.Equal(meta.TaggingConfigUpdatedAt) {
t.Fatalf("tags-only import hook leaked unspecified fields: %+v", got)
}
})
}})
@@ -622,3 +666,43 @@ func TestPeerBucketMetadataNormalizesBeforeComparison(t *testing.T) {
})
}})
}
func TestPeerBucketMetadataEqualTimeArrivalOrders(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) {
created := UTCNow().Add(-time.Hour)
a := []byte(`<Tagging><TagSet><Tag><Key>key</Key><Value>a</Value></Tag></TagSet></Tagging>`)
b := bytes.ReplaceAll(a, []byte("<Value>a"), []byte("<Value>b"))
for _, pair := range [][2][]byte{{a, b}, {b, a}, {a, nil}, {nil, a}} {
meta := newBucketMetadata(bucket)
meta.Created = created
meta.defaultTimestamps()
if err := globalBucketMetadataSys.save(t.Context(), meta); err != nil {
t.Fatal(err)
}
for _, data := range pair {
var payload *string
if data != nil {
encoded := base64.StdEncoding.EncodeToString(data)
payload = &encoded
}
rec := applySRBucketMetaViaAdmin(t, cred, madmin.SRBucketMeta{Bucket: bucket, Type: madmin.SRBucketMetaTypeTags, Tags: payload, UpdatedAt: created.Add(time.Minute)})
if rec.Code != http.StatusOK {
t.Fatalf("equal time apply: %d %s", rec.Code, rec.Body.String())
}
}
got, err := readBucketMetadata(t.Context(), obj, bucket)
if err != nil {
t.Fatal(err)
}
want := b
if pair[0] == nil || pair[1] == nil {
want = nil
}
if !bytes.Equal(got.TaggingConfigXML, want) || !got.TaggingConfigUpdatedAt.Equal(created.Add(time.Minute)) {
t.Fatalf("arrival order changed winner: %q at %v", got.TaggingConfigXML, got.TaggingConfigUpdatedAt)
}
}
})
}})
}
+20 -63
View File
@@ -1643,7 +1643,7 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context
if item.Cors != nil {
replLogOnceIf(ctx, fmt.Errorf("ignoring CORS event for bucket %s from %v before bucket creation at %v", item.Bucket, item.UpdatedAt, meta.Created), "cors-event-before-bucket-creation-"+item.Bucket)
}
return nil
item.Cors = nil
}
// Presence is separate from content: omitted RawMessage is not a delete;
@@ -1669,6 +1669,7 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context
}
if len(updates) != 0 {
if err := ensureBucketMetadataCreated(ctx, objectAPI, &meta); err != nil {
logBucketConfigReplication(ctx, item.Bucket, "bulk", "indeterminate", item.UpdatedAt, meta.Created, err.Error())
return wrapSRErr(err)
}
}
@@ -1678,6 +1679,10 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context
if !supplied {
continue
}
if item.UpdatedAt.Before(meta.Created) {
logBucketConfigReplication(ctx, item.Bucket, file, "before-created", item.UpdatedAt, meta.Created, "bulk event")
continue
}
applied, err := applyBucketConfig(&meta, file, data, item.UpdatedAt)
if err != nil {
return wrapSRErr(err)
@@ -2146,57 +2151,17 @@ func (c *SiteReplicationSys) syncToAllPeers(ctx context.Context, addOpts madmin.
return errSRBucketConfigError(err)
}
// Replicate bucket policy if present.
policyJSON, tm := meta.PolicyConfigJSON, meta.PolicyConfigUpdatedAt
if len(policyJSON) > 0 {
err = c.BucketMetaHook(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypePolicy,
Bucket: bucket,
Policy: policyJSON,
UpdatedAt: tm,
})
// Versioning is bootstrapped by MakeBucketHook and then reconciled by
// heal. Preserve the existing initial-sync fields during rolling upgrades.
for _, file := range []string{bucketPolicyConfig, bucketTaggingConfig, objectLockConfig, bucketSSEConfig, bucketQuotaConfigFile} {
event, send, err := initialBucketConfigReplicationEvent(meta, file)
if err != nil {
return errSRBucketMetaError(err)
}
}
// Replicate bucket tags if present.
tagCfg, tm := meta.TaggingConfigXML, meta.TaggingConfigUpdatedAt
if len(tagCfg) > 0 {
tagCfgStr := base64.StdEncoding.EncodeToString(tagCfg)
err = c.BucketMetaHook(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeTags,
Bucket: bucket,
Tags: &tagCfgStr,
UpdatedAt: tm,
})
if err != nil {
return errSRBucketMetaError(err)
}
}
// Replicate object-lock config if present.
objLockCfgData, tm := meta.ObjectLockConfigXML, meta.ObjectLockConfigUpdatedAt
if len(objLockCfgData) > 0 {
objLockStr := base64.StdEncoding.EncodeToString(objLockCfgData)
err = c.BucketMetaHook(ctx, newSRBucketObjectLockMeta(bucket, &objLockStr, tm))
if err != nil {
return errSRBucketMetaError(err)
}
}
// Replicate existing bucket bucket encryption settings
sseConfigData, tm := meta.EncryptionConfigXML, meta.EncryptionConfigUpdatedAt
if len(sseConfigData) > 0 {
sseConfigStr := base64.StdEncoding.EncodeToString(sseConfigData)
err = c.BucketMetaHook(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeSSEConfig,
Bucket: bucket,
SSEConfig: &sseConfigStr,
UpdatedAt: tm,
})
if err != nil {
return errSRBucketMetaError(err)
if send {
if err := c.BucketMetaHook(ctx, event); err != nil {
return errSRBucketMetaError(err)
}
}
}
@@ -2208,20 +2173,6 @@ func (c *SiteReplicationSys) syncToAllPeers(ctx context.Context, addOpts madmin.
}
}
// Replicate existing bucket quotas settings
quotaConfigJSON, tm := meta.QuotaConfigJSON, meta.QuotaConfigUpdatedAt
if len(quotaConfigJSON) > 0 {
err = c.BucketMetaHook(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeQuotaConfig,
Bucket: bucket,
Quota: quotaConfigJSON,
UpdatedAt: tm,
})
if err != nil {
return errSRBucketMetaError(err)
}
}
// Replicate ILM expiry rules if needed
if addOpts.ReplicateILMExpiry && (meta.lifecycleConfig != nil && meta.lifecycleConfig.HasExpiry()) {
var expLclCfg lifecycle.Lifecycle
@@ -4006,6 +3957,8 @@ func (c *SiteReplicationSys) SiteReplicationMetaInfo(ctx context.Context, objAPI
tagCfgStr := base64.StdEncoding.EncodeToString(meta.TaggingConfigXML)
bms.Tags = &tagCfgStr
bms.TagConfigUpdatedAt = meta.TaggingConfigUpdatedAt
} else if globalSiteReplicationMetadataTombstones && !meta.Created.IsZero() && meta.TaggingConfigUpdatedAt.After(meta.Created) {
bms.TagConfigUpdatedAt = meta.TaggingConfigUpdatedAt
}
if len(meta.VersioningConfigXML) > 0 {
@@ -4024,12 +3977,16 @@ func (c *SiteReplicationSys) SiteReplicationMetaInfo(ctx context.Context, objAPI
quotaConfigStr := base64.StdEncoding.EncodeToString(meta.QuotaConfigJSON)
bms.QuotaConfig = &quotaConfigStr
bms.QuotaConfigUpdatedAt = meta.QuotaConfigUpdatedAt
} else if globalSiteReplicationMetadataTombstones && !meta.Created.IsZero() && meta.QuotaConfigUpdatedAt.After(meta.Created) {
bms.QuotaConfigUpdatedAt = meta.QuotaConfigUpdatedAt
}
if len(meta.EncryptionConfigXML) > 0 {
sseConfigStr := base64.StdEncoding.EncodeToString(meta.EncryptionConfigXML)
bms.SSEConfig = &sseConfigStr
bms.SSEConfigUpdatedAt = meta.EncryptionConfigUpdatedAt
} else if globalSiteReplicationMetadataTombstones && !meta.Created.IsZero() && meta.EncryptionConfigUpdatedAt.After(meta.Created) {
bms.SSEConfigUpdatedAt = meta.EncryptionConfigUpdatedAt
}
bms.CorsConfigUpdatedAt = meta.CorsConfigUpdatedAt