From 01aaef2b504dfc2bfbd302469030f19f007661d7 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Sat, 12 Sep 2026 11:22:04 +0800 Subject: [PATCH] fix(replication): converge bucket metadata using deterministic source states Signed-off-by: Feng Ruohang --- cmd/site-replication-metadata-heal_test.go | 192 ++++++++++ cmd/site-replication-metadata.go | 145 +++++++ cmd/site-replication.go | 421 +-------------------- 3 files changed, 343 insertions(+), 415 deletions(-) create mode 100644 cmd/site-replication-metadata-heal_test.go create mode 100644 cmd/site-replication-metadata.go diff --git a/cmd/site-replication-metadata-heal_test.go b/cmd/site-replication-metadata-heal_test.go new file mode 100644 index 000000000..148a11c93 --- /dev/null +++ b/cmd/site-replication-metadata-heal_test.go @@ -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 . + +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(`keyvalue`), + bucketSSEConfig: []byte(`AES256`), + 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) + } + } + }) + } + }) + }}) +} diff --git a/cmd/site-replication-metadata.go b/cmd/site-replication-metadata.go new file mode 100644 index 000000000..22438173a --- /dev/null +++ b/cmd/site-replication-metadata.go @@ -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 . + +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 +} diff --git a/cmd/site-replication.go b/cmd/site-replication.go index fba20be6b..2c70e983c 100644 --- a/cmd/site-replication.go +++ b/cmd/site-replication.go @@ -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) {