fix(replication): converge bucket metadata using deterministic source states

Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
Feng Ruohang
2026-09-12 11:22:04 +08:00
parent bcc62afe3d
commit 01aaef2b50
3 changed files with 343 additions and 415 deletions
+192
View File
@@ -0,0 +1,192 @@
// 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"
"encoding/base64"
"fmt"
"net/http"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
)
func bucketConfigTestData(bucket string) map[string][]byte {
return map[string][]byte{
bucketPolicyConfig: []byte(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket)),
bucketTaggingConfig: []byte(`<Tagging><TagSet><Tag><Key>key</Key><Value>value</Value></Tag></TagSet></Tagging>`),
bucketSSEConfig: []byte(`<ServerSideEncryptionConfiguration><Rule><ApplyServerSideEncryptionByDefault><SSEAlgorithm>AES256</SSEAlgorithm></ApplyServerSideEncryptionByDefault></Rule></ServerSideEncryptionConfiguration>`),
bucketQuotaConfigFile: []byte(`{"quota":1024,"quotatype":"hard"}`),
bucketVersioningConfig: enabledBucketVersioningConfig,
objectLockConfig: enabledBucketObjectLockConfig,
}
}
// Construct independent wire fixtures, including timestamps hidden by legacy
// exporters, without going through the production event or state helpers.
func bucketConfigTestInfo(bucket, file string, data []byte, at, created time.Time) srBucketStatsSummary {
info := madmin.SRBucketInfo{Bucket: bucket, CreatedAt: created}
var payload *string
if len(data) != 0 {
encoded := base64.StdEncoding.EncodeToString(data)
payload = &encoded
}
switch file {
case bucketPolicyConfig:
info.Policy, info.PolicyUpdatedAt = data, at
case bucketTaggingConfig:
info.Tags, info.TagConfigUpdatedAt = payload, at
case bucketSSEConfig:
info.SSEConfig, info.SSEConfigUpdatedAt = payload, at
case bucketQuotaConfigFile:
info.QuotaConfig, info.QuotaConfigUpdatedAt = payload, at
case bucketVersioningConfig:
info.Versioning, info.VersioningConfigUpdatedAt = payload, at
case objectLockConfig:
info.ObjectLockConfig, info.ObjectLockConfigUpdatedAt = payload, at
}
return srBucketStatsSummary{meta: srBucketMetaInfo{SRBucketInfo: info}}
}
func TestLatestBucketConfigCandidates(t *testing.T) {
created := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC)
for file, data := range bucketConfigTestData("bucket") {
t.Run(file, func(t *testing.T) {
info := srStatusInfo{Sites: map[string]madmin.PeerInfo{"a": {}, "b": {}, "c": {}}, BucketStats: map[string]map[string]srBucketStatsSummary{"bucket": {}}}
bs := info.BucketStats["bucket"]
baseline := bucketConfigTestInfo("bucket", file, nil, created, created)
bs["a"], bs["b"], bs["c"] = baseline, baseline, baseline
if _, found := latestBucketConfig("bucket", file, info); found {
t.Fatal("empty baseline chosen as source")
}
// Unknown and unavailable IDs cannot contribute, even with future clocks.
bs[""], bs["unknown"] = bucketConfigTestInfo("bucket", file, data, created.Add(10*time.Hour), created), bucketConfigTestInfo("bucket", file, data, created.Add(20*time.Hour), created)
if _, found := latestBucketConfig("bucket", file, info); found {
t.Fatal("unknown source chosen")
}
liveBaseline := bucketConfigTestInfo("bucket", file, data, created, created)
real := bucketConfigTestInfo("bucket", file, data, created.Add(time.Hour), created)
oldGeneration := bucketConfigTestInfo("bucket", file, data, created, created.Add(time.Second))
states := []srBucketStatsSummary{liveBaseline, real, 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)
if !found || !winner.at.Equal(created.Add(time.Hour)) || !bytes.Equal(winner.data, data) {
t.Fatalf("permutation %v chose %+v found=%v", order, winner, found)
}
}
bs["a"], bs["b"], bs["c"] = baseline, liveBaseline, baseline
if winner, found := latestBucketConfig("bucket", file, info); !found || !bytes.Equal(winner.data, data) || winner.real {
t.Fatal("historical baseline-live cannot initialize")
}
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
if winner, found := latestBucketConfig("bucket", file, info); !found || len(winner.data) == 0 {
t.Fatal("default became a deletion")
}
bs["c"] = bucketConfigTestInfo("bucket", file, nil, created.Add(time.Hour), created)
if winner, found := latestBucketConfig("bucket", file, info); !found || len(winner.data) != 0 || !winner.real {
t.Fatal("real tombstone lost equal-time tie")
}
}
})
}
}
func TestHealBucketConfigSourceAndQuiescence(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()
events := recordBucketConfigPeer(t, cred)
local := globalDeploymentID()
created := UTCNow().Add(-time.Hour)
oldAt, newAt := created.Add(time.Minute), created.Add(2*time.Minute)
counter := &bucketConfigWriteCounter{ObjectLayer: obj}
setObjectLayer(counter)
defer setObjectLayer(obj)
for file, data := range bucketConfigTestData(bucket) {
t.Run(file, func(t *testing.T) {
meta := newBucketMetadata(bucket)
meta.Created = created
meta.defaultTimestamps()
value, at := replicatedBucketConfig(&meta, file)
*value, *at = data, newAt
if err := globalBucketMetadataSys.save(ctx, meta); err != nil {
t.Fatal(err)
}
bs := map[string]srBucketStatsSummary{
local: bucketConfigTestInfo(bucket, file, data, newAt, created),
"metadata-peer": bucketConfigTestInfo(bucket, file, data, oldAt, created),
"": {}, "unknown": bucketConfigTestInfo(bucket, file, nil, newAt.Add(time.Hour), created),
// Known but absent from c.state.Peers, forcing getAdminClient failure.
"broken": bucketConfigTestInfo(bucket, file, nil, created, created),
}
info := srStatusInfo{Sites: map[string]madmin.PeerInfo{local: {}, "metadata-peer": {}, "broken": {}}, BucketStats: map[string]map[string]srBucketStatsSummary{bucket: bs}}
before := len(events())
writes := counter.writes.Load()
if err := globalSiteReplicationSys.healBucketConfig(ctx, bucket, file, info); err != nil {
t.Fatal(err)
}
got := events()[before:]
if len(got) != 1 || !got[0].UpdatedAt.Equal(newAt) {
t.Fatalf("same-payload timestamp heal did not reach healthy peer: %+v", got)
}
if file == bucketTaggingConfig && got[0].Tags == nil {
t.Fatal("tag event missing payload")
}
if counter.writes.Load() != writes {
t.Fatal("matching local state was rewritten")
}
bs["metadata-peer"], bs["broken"] = bs[local], bs[local]
if err := globalSiteReplicationSys.healBucketConfig(ctx, bucket, file, info); err != nil {
t.Fatal(err)
}
if len(events()) != before+1 || counter.writes.Load() != writes {
t.Fatal("stable second heal wrote or broadcast")
}
if bucketConfigUpdateOnly(file) {
return
}
bs["metadata-peer"] = bucketConfigTestInfo(bucket, file, nil, newAt.Add(time.Minute), created)
delete(bs, "broken")
if err := globalSiteReplicationSys.healBucketConfig(ctx, bucket, file, info); err != nil {
t.Fatal(err)
}
after, err := loadBucketMetadata(ctx, obj, bucket)
if err != nil {
t.Fatal(err)
}
value, at = replicatedBucketConfig(&after, file)
if len(*value) != 0 || !at.Equal(newAt.Add(time.Minute)) {
t.Fatalf("local heal lost deletion/source time: %s %v", *value, *at)
}
if file == bucketQuotaConfigFile {
q, _, err := globalBucketMetadataSys.GetQuotaConfig(ctx, bucket)
if err != nil || q.Quota != 0 {
t.Fatalf("quota cache not cleared: %+v %v", q, err)
}
}
})
}
})
}})
}
+145
View File
@@ -0,0 +1,145 @@
// 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 (
"context"
"encoding/base64"
"fmt"
"time"
"github.com/minio/madmin-go/v3"
)
func bucketConfigStateFromInfo(bucket, file string, meta madmin.SRBucketInfo) (bucketConfigState, error) {
var payload *string
var at time.Time
switch file {
case bucketPolicyConfig:
return newBucketConfigState(bucket, file, meta.Policy, meta.PolicyUpdatedAt, meta.CreatedAt, false)
case bucketTaggingConfig:
payload, at = meta.Tags, meta.TagConfigUpdatedAt
case bucketSSEConfig:
payload, at = meta.SSEConfig, meta.SSEConfigUpdatedAt
case bucketQuotaConfigFile:
payload, at = meta.QuotaConfig, meta.QuotaConfigUpdatedAt
case bucketVersioningConfig:
payload, at = meta.Versioning, meta.VersioningConfigUpdatedAt
case objectLockConfig:
payload, at = meta.ObjectLockConfig, meta.ObjectLockConfigUpdatedAt
}
var data []byte
if payload != nil {
var err error
data, err = base64.StdEncoding.DecodeString(*payload)
if err != nil {
return bucketConfigState{}, err
}
}
lockEnabled := meta.ObjectLockConfig != nil && len(*meta.ObjectLockConfig) != 0
return newBucketConfigState(bucket, file, data, at, meta.CreatedAt, lockEnabled)
}
func newBucketConfigReplicationEvent(bucket, file string, state bucketConfigState) madmin.SRBucketMeta {
event := madmin.SRBucketMeta{Bucket: bucket, UpdatedAt: state.at}
var payload *string
if len(state.data) != 0 {
encoded := base64.StdEncoding.EncodeToString(state.data)
payload = &encoded
}
switch file {
case bucketPolicyConfig:
event.Type, event.Policy = madmin.SRBucketMetaTypePolicy, state.data
case bucketTaggingConfig:
event.Type, event.Tags = madmin.SRBucketMetaTypeTags, payload
case bucketSSEConfig:
event.Type, event.SSEConfig = madmin.SRBucketMetaTypeSSEConfig, payload
case bucketQuotaConfigFile:
event.Type, event.Quota = madmin.SRBucketMetaTypeQuotaConfig, state.data
case bucketVersioningConfig:
event.Type, event.Versioning = madmin.SRBucketMetaTypeVersionConfig, payload
case objectLockConfig:
event.Type, event.ObjectLockConfig = madmin.SRBucketMetaTypeObjectLockConfig, payload
}
return event
}
func latestBucketConfig(bucket, file string, info srStatusInfo) (bucketConfigState, bool) {
var latest bucketConfigState
found := false
for id, status := range info.BucketStats[bucket] {
if _, known := info.Sites[id]; !known || id == "" {
continue
}
state, err := bucketConfigStateFromInfo(bucket, file, status.meta.SRBucketInfo)
if err != nil || !state.candidate() {
continue
}
if !found || compareBucketConfigStates(state, latest) > 0 {
latest, found = state, true
}
}
return latest, found
}
func (c *SiteReplicationSys) healBucketConfig(ctx context.Context, bucket, file string, info srStatusInfo) error {
c.RLock()
defer c.RUnlock()
if !c.enabled {
return nil
}
latest, found := latestBucketConfig(bucket, file, info)
if !found {
return nil
}
for id, status := range info.BucketStats[bucket] {
if _, known := info.Sites[id]; !known || id == "" {
continue
}
target := status.meta.SRBucketInfo
if target.CreatedAt.IsZero() || latest.at.Before(target.CreatedAt) {
continue
}
// Versioning can be normalized differently until Object Lock itself has
// converged. Compare what this target would actually persist.
incoming, err := newBucketConfigState(bucket, file, latest.data, latest.at, target.CreatedAt,
target.ObjectLockConfig != nil && len(*target.ObjectLockConfig) != 0)
if err != nil || !incoming.candidate() {
continue
}
current, err := bucketConfigStateFromInfo(bucket, file, target)
if err == nil && compareBucketConfigStates(incoming, current) <= 0 {
continue
}
if id == globalDeploymentID() {
_, err = globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, file, latest.data, false, false, &latest.at)
} else {
var client *madmin.AdminClient
client, err = c.getAdminClient(ctx, id)
if err == nil {
err = client.SRPeerReplicateBucketMeta(ctx, newBucketConfigReplicationEvent(bucket, file, incoming))
}
}
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))
}
}
return nil
}
+6 -415
View File
@@ -4883,366 +4883,23 @@ func (c *SiteReplicationSys) healBucketILMExpiry(ctx context.Context, objAPI Obj
}
func (c *SiteReplicationSys) healTagMetadata(ctx context.Context, objAPI ObjectLayer, bucket string, info srStatusInfo) error {
bs := info.BucketStats[bucket]
c.RLock()
defer c.RUnlock()
if !c.enabled {
return nil
}
var (
latestID, latestPeerName string
lastUpdate time.Time
latestTaggingConfig *string
)
for dID, ss := range bs {
if lastUpdate.IsZero() {
lastUpdate = ss.meta.TagConfigUpdatedAt
latestID = dID
latestTaggingConfig = ss.meta.Tags
}
// avoid considering just created buckets as latest. Perhaps this site
// just joined cluster replication and yet to be sync'd
if ss.meta.CreatedAt.Equal(ss.meta.TagConfigUpdatedAt) {
continue
}
if ss.meta.TagConfigUpdatedAt.After(lastUpdate) {
lastUpdate = ss.meta.TagConfigUpdatedAt
latestID = dID
latestTaggingConfig = ss.meta.Tags
}
}
latestPeerName = info.Sites[latestID].Name
var latestTaggingConfigBytes []byte
var err error
if latestTaggingConfig != nil {
latestTaggingConfigBytes, err = base64.StdEncoding.DecodeString(*latestTaggingConfig)
if err != nil {
return err
}
}
for dID, bStatus := range bs {
if !bStatus.TagMismatch {
continue
}
if isBucketMetadataEqual(latestTaggingConfig, bStatus.meta.Tags) {
continue
}
if dID == globalDeploymentID() {
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketTaggingConfig, latestTaggingConfigBytes); err != nil {
replLogIf(ctx, fmt.Errorf("Unable to heal tagging metadata from peer site %s : %w", latestPeerName, err))
}
continue
}
admClient, err := c.getAdminClient(ctx, dID)
if err != nil {
return wrapSRErr(err)
}
peerName := info.Sites[dID].Name
err = admClient.SRPeerReplicateBucketMeta(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeTags,
Bucket: bucket,
Tags: latestTaggingConfig,
})
if err != nil {
replLogIf(ctx, c.annotatePeerErr(peerName, replicateBucketMetadata,
fmt.Errorf("Unable to heal tagging metadata for peer %s from peer %s : %w", peerName, latestPeerName, err)))
}
}
return nil
return c.healBucketConfig(ctx, bucket, bucketTaggingConfig, info)
}
func (c *SiteReplicationSys) healBucketPolicies(ctx context.Context, objAPI ObjectLayer, bucket string, info srStatusInfo) error {
bs := info.BucketStats[bucket]
c.RLock()
defer c.RUnlock()
if !c.enabled {
return nil
}
var (
latestID, latestPeerName string
lastUpdate time.Time
latestIAMPolicy json.RawMessage
)
for dID, ss := range bs {
if lastUpdate.IsZero() {
lastUpdate = ss.meta.PolicyUpdatedAt
latestID = dID
latestIAMPolicy = ss.meta.Policy
}
// avoid considering just created buckets as latest. Perhaps this site
// just joined cluster replication and yet to be sync'd
if ss.meta.CreatedAt.Equal(ss.meta.PolicyUpdatedAt) {
continue
}
if ss.meta.PolicyUpdatedAt.After(lastUpdate) {
lastUpdate = ss.meta.PolicyUpdatedAt
latestID = dID
latestIAMPolicy = ss.meta.Policy
}
}
latestPeerName = info.Sites[latestID].Name
for dID, bStatus := range bs {
if !bStatus.PolicyMismatch {
continue
}
if strings.EqualFold(string(latestIAMPolicy), string(bStatus.meta.Policy)) {
continue
}
if dID == globalDeploymentID() {
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketPolicyConfig, latestIAMPolicy); err != nil {
replLogIf(ctx, fmt.Errorf("Unable to heal bucket policy metadata from peer site %s : %w", latestPeerName, err))
}
continue
}
admClient, err := c.getAdminClient(ctx, dID)
if err != nil {
return wrapSRErr(err)
}
peerName := info.Sites[dID].Name
if err = admClient.SRPeerReplicateBucketMeta(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypePolicy,
Bucket: bucket,
Policy: latestIAMPolicy,
UpdatedAt: lastUpdate,
}); err != nil {
replLogIf(ctx, c.annotatePeerErr(peerName, replicateBucketMetadata,
fmt.Errorf("Unable to heal bucket policy metadata for peer %s from peer %s : %w",
peerName, latestPeerName, err)))
}
}
return nil
return c.healBucketConfig(ctx, bucket, bucketPolicyConfig, info)
}
func (c *SiteReplicationSys) healBucketQuotaConfig(ctx context.Context, objAPI ObjectLayer, bucket string, info srStatusInfo) error {
bs := info.BucketStats[bucket]
c.RLock()
defer c.RUnlock()
if !c.enabled {
return nil
}
var (
latestID, latestPeerName string
lastUpdate time.Time
latestQuotaConfig *string
latestQuotaConfigBytes []byte
)
for dID, ss := range bs {
if lastUpdate.IsZero() {
lastUpdate = ss.meta.QuotaConfigUpdatedAt
latestID = dID
latestQuotaConfig = ss.meta.QuotaConfig
}
// avoid considering just created buckets as latest. Perhaps this site
// just joined cluster replication and yet to be sync'd
if ss.meta.CreatedAt.Equal(ss.meta.QuotaConfigUpdatedAt) {
continue
}
if ss.meta.QuotaConfigUpdatedAt.After(lastUpdate) {
lastUpdate = ss.meta.QuotaConfigUpdatedAt
latestID = dID
latestQuotaConfig = ss.meta.QuotaConfig
}
}
var err error
if latestQuotaConfig != nil {
latestQuotaConfigBytes, err = base64.StdEncoding.DecodeString(*latestQuotaConfig)
if err != nil {
return err
}
}
latestPeerName = info.Sites[latestID].Name
for dID, bStatus := range bs {
if !bStatus.QuotaCfgMismatch {
continue
}
if isBucketMetadataEqual(latestQuotaConfig, bStatus.meta.QuotaConfig) {
continue
}
if dID == globalDeploymentID() {
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketQuotaConfigFile, latestQuotaConfigBytes); err != nil {
replLogIf(ctx, fmt.Errorf("Unable to heal quota metadata from peer site %s : %w", latestPeerName, err))
}
continue
}
admClient, err := c.getAdminClient(ctx, dID)
if err != nil {
return wrapSRErr(err)
}
peerName := info.Sites[dID].Name
if err = admClient.SRPeerReplicateBucketMeta(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeQuotaConfig,
Bucket: bucket,
Quota: latestQuotaConfigBytes,
UpdatedAt: lastUpdate,
}); err != nil {
replLogIf(ctx, c.annotatePeerErr(peerName, replicateBucketMetadata,
fmt.Errorf("Unable to heal quota config metadata for peer %s from peer %s : %w",
peerName, latestPeerName, err)))
}
}
return nil
return c.healBucketConfig(ctx, bucket, bucketQuotaConfigFile, info)
}
func (c *SiteReplicationSys) healVersioningMetadata(ctx context.Context, objAPI ObjectLayer, bucket string, info srStatusInfo) error {
c.RLock()
defer c.RUnlock()
if !c.enabled {
return nil
}
var (
latestID, latestPeerName string
lastUpdate time.Time
latestVersioningConfig *string
)
bs := info.BucketStats[bucket]
for dID, ss := range bs {
if lastUpdate.IsZero() {
lastUpdate = ss.meta.VersioningConfigUpdatedAt
latestID = dID
latestVersioningConfig = ss.meta.Versioning
}
// avoid considering just created buckets as latest. Perhaps this site
// just joined cluster replication and yet to be sync'd
if ss.meta.CreatedAt.Equal(ss.meta.VersioningConfigUpdatedAt) {
continue
}
if ss.meta.VersioningConfigUpdatedAt.After(lastUpdate) {
lastUpdate = ss.meta.VersioningConfigUpdatedAt
latestID = dID
latestVersioningConfig = ss.meta.Versioning
}
}
latestPeerName = info.Sites[latestID].Name
var latestVersioningConfigBytes []byte
var err error
if latestVersioningConfig != nil {
latestVersioningConfigBytes, err = base64.StdEncoding.DecodeString(*latestVersioningConfig)
if err != nil {
return err
}
}
for dID, bStatus := range bs {
if !bStatus.VersioningConfigMismatch {
continue
}
if isBucketMetadataEqual(latestVersioningConfig, bStatus.meta.Versioning) {
continue
}
if dID == globalDeploymentID() {
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, latestVersioningConfigBytes); err != nil {
replLogIf(ctx, fmt.Errorf("Unable to heal versioning metadata from peer site %s : %w", latestPeerName, err))
}
continue
}
admClient, err := c.getAdminClient(ctx, dID)
if err != nil {
return wrapSRErr(err)
}
peerName := info.Sites[dID].Name
err = admClient.SRPeerReplicateBucketMeta(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeVersionConfig,
Bucket: bucket,
Versioning: latestVersioningConfig,
UpdatedAt: lastUpdate,
})
if err != nil {
replLogIf(ctx, c.annotatePeerErr(peerName, replicateBucketMetadata,
fmt.Errorf("Unable to heal versioning config metadata for peer %s from peer %s : %w",
peerName, latestPeerName, err)))
}
}
return nil
return c.healBucketConfig(ctx, bucket, bucketVersioningConfig, info)
}
func (c *SiteReplicationSys) healSSEMetadata(ctx context.Context, objAPI ObjectLayer, bucket string, info srStatusInfo) error {
c.RLock()
defer c.RUnlock()
if !c.enabled {
return nil
}
var (
latestID, latestPeerName string
lastUpdate time.Time
latestSSEConfig *string
)
bs := info.BucketStats[bucket]
for dID, ss := range bs {
if lastUpdate.IsZero() {
lastUpdate = ss.meta.SSEConfigUpdatedAt
latestID = dID
latestSSEConfig = ss.meta.SSEConfig
}
// avoid considering just created buckets as latest. Perhaps this site
// just joined cluster replication and yet to be sync'd
if ss.meta.CreatedAt.Equal(ss.meta.SSEConfigUpdatedAt) {
continue
}
if ss.meta.SSEConfigUpdatedAt.After(lastUpdate) {
lastUpdate = ss.meta.SSEConfigUpdatedAt
latestID = dID
latestSSEConfig = ss.meta.SSEConfig
}
}
latestPeerName = info.Sites[latestID].Name
var latestSSEConfigBytes []byte
var err error
if latestSSEConfig != nil {
latestSSEConfigBytes, err = base64.StdEncoding.DecodeString(*latestSSEConfig)
if err != nil {
return err
}
}
for dID, bStatus := range bs {
if !bStatus.SSEConfigMismatch {
continue
}
if isBucketMetadataEqual(latestSSEConfig, bStatus.meta.SSEConfig) {
continue
}
if dID == globalDeploymentID() {
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketSSEConfig, latestSSEConfigBytes); err != nil {
replLogIf(ctx, fmt.Errorf("Unable to heal sse metadata from peer site %s : %w", latestPeerName, err))
}
continue
}
admClient, err := c.getAdminClient(ctx, dID)
if err != nil {
return wrapSRErr(err)
}
peerName := info.Sites[dID].Name
err = admClient.SRPeerReplicateBucketMeta(ctx, madmin.SRBucketMeta{
Type: madmin.SRBucketMetaTypeSSEConfig,
Bucket: bucket,
SSEConfig: latestSSEConfig,
UpdatedAt: lastUpdate,
})
if err != nil {
replLogIf(ctx, c.annotatePeerErr(peerName, replicateBucketMetadata,
fmt.Errorf("Unable to heal SSE config metadata for peer %s from peer %s : %w",
peerName, latestPeerName, err)))
}
}
return nil
return c.healBucketConfig(ctx, bucket, bucketSSEConfig, info)
}
func latestCORSConfig(bs map[string]srBucketStatsSummary) (latestID string, latest corsReplicationState, ok bool) {
@@ -5310,73 +4967,7 @@ func (c *SiteReplicationSys) healCORSMetadata(ctx context.Context, objAPI Object
}
func (c *SiteReplicationSys) healOLockConfigMetadata(ctx context.Context, objAPI ObjectLayer, bucket string, info srStatusInfo) error {
bs := info.BucketStats[bucket]
c.RLock()
defer c.RUnlock()
if !c.enabled {
return nil
}
var (
latestID, latestPeerName string
lastUpdate time.Time
latestObjLockConfig *string
)
for dID, ss := range bs {
if lastUpdate.IsZero() {
lastUpdate = ss.meta.ObjectLockConfigUpdatedAt
latestID = dID
latestObjLockConfig = ss.meta.ObjectLockConfig
}
// avoid considering just created buckets as latest. Perhaps this site
// just joined cluster replication and yet to be sync'd
if ss.meta.CreatedAt.Equal(ss.meta.ObjectLockConfigUpdatedAt) {
continue
}
if ss.meta.ObjectLockConfig != nil && ss.meta.ObjectLockConfigUpdatedAt.After(lastUpdate) {
lastUpdate = ss.meta.ObjectLockConfigUpdatedAt
latestID = dID
latestObjLockConfig = ss.meta.ObjectLockConfig
}
}
latestPeerName = info.Sites[latestID].Name
var latestObjLockConfigBytes []byte
var err error
if latestObjLockConfig != nil {
latestObjLockConfigBytes, err = base64.StdEncoding.DecodeString(*latestObjLockConfig)
if err != nil {
return err
}
}
for dID, bStatus := range bs {
if !bStatus.OLockConfigMismatch {
continue
}
if isBucketMetadataEqual(latestObjLockConfig, bStatus.meta.ObjectLockConfig) {
continue
}
if dID == globalDeploymentID() {
if _, err := globalBucketMetadataSys.Update(ctx, bucket, objectLockConfig, latestObjLockConfigBytes); err != nil {
replLogIf(ctx, fmt.Errorf("Unable to heal objectlock config metadata from peer site %s : %w", latestPeerName, err))
}
continue
}
admClient, err := c.getAdminClient(ctx, dID)
if err != nil {
return wrapSRErr(err)
}
peerName := info.Sites[dID].Name
err = admClient.SRPeerReplicateBucketMeta(ctx, newSRBucketObjectLockMeta(bucket, latestObjLockConfig, lastUpdate))
if err != nil {
replLogIf(ctx, c.annotatePeerErr(peerName, replicateBucketMetadata,
fmt.Errorf("Unable to heal object lock config metadata for peer %s from peer %s : %w",
peerName, latestPeerName, err)))
}
}
return nil
return c.healBucketConfig(ctx, bucket, objectLockConfig, info)
}
func (c *SiteReplicationSys) purgeDeletedBucket(ctx context.Context, objAPI ObjectLayer, bucket string) {