fix: close residual bucket-metadata races (issue #105 audit)

Audit of the three deferred #105 follow-ups. Each reproduces with a
deterministic red test in cmd/bucket-metadata-race_test.go, and each fix is
the minimal change that turns its test green while preserving the
<bucket>.lck -> metadata.lock -> .metadata.bin lock order established by #103.

1. Lifecycle expiry merge lost update (persistent). PeerBucketLCConfigHandler
   and healBucketILMExpiry read the current lifecycle with an unlocked
   GetConfigFromDisk, merged the replicated expiry rules with the local
   transition rules, then wrote the pre-computed blob via Update. Any lifecycle
   transition change committed between the merge read and the merge write was
   silently lost. New BucketMetadataSys.UpdateExpiryLCConfig performs the read,
   merge, and save under one metadata.lock; mergeExpiryWithLCConfig now takes
   the locked snapshot and validates object-lock retention from it instead of
   re-reading (avoids a re-entrant metadata load under the lock).

2. DeleteBucket ghost .metadata.bin (persistent). DeleteBucket took only
   <bucket>.lck while config writers take only metadata.lock, so a writer that
   was mid-save could re-create .metadata.bin after the prefix purge. The purge
   now runs under metadata.lock, with a best-effort unlocked fallback so a
   delete is never blocked from completing.

3. Overlapping peer reloads publishing a stale resident cache (freshness only;
   the persisted record stays correct). LoadBucketMetadataHandler and the
   GetConfig cache-miss path published with an unconditional Set, so a reload
   that read an older revision could overwrite a newer resident record until the
   next refresh. New BucketMetadataSys.setReloaded (and a matching GetConfig
   guard) refuses to regress a newer resident record, mirroring
   refreshBucketsMetadataLoop.

Verification: go build -tags kqueue,dev ./...; go vet ./cmd; gofmt clean;
rebrand-guard baseline unchanged; go test -tags kqueue,dev ./cmd (207s) green;
new tests plus the #103 metadata suite green under -race.

Refs #105. Parent #102. Foundation #103.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01L7qJqWwy8oFA6aCXWRzXQe
Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
Feng Ruohang
2026-09-07 10:19:28 +08:00
parent ad873c7357
commit e7654d470c
5 changed files with 523 additions and 28 deletions
+408
View File
@@ -0,0 +1,408 @@
// Copyright (c) 2015-2026 MinIO, Inc.
// Copyright (c) 2026 PGSTY
//
// This file is part of MinIO Object Storage stack
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
package cmd
import (
"bytes"
"context"
"encoding/base64"
"errors"
"io"
"net/http"
"sync"
"testing"
"time"
"github.com/minio/minio/internal/auth"
)
// ---------------------------------------------------------------------------
// Target 1 (issue #105): a higher-level lifecycle XML merge racing another
// bucket-metadata transition outside BucketMetadataSys.Delete.
//
// Before the fix, PeerBucketLCConfigHandler / healBucketILMExpiry read the
// current lifecycle document without metadata.lock, merged the replicated expiry
// rules with the local transition rules, and then persisted the merged blob with
// BucketMetadataSys.Update. Update re-read the record under metadata.lock but
// overwrote LifecycleConfigXML wholesale with the pre-computed blob, so any
// lifecycle transition change committed between the merge read and the merge
// write was silently lost on disk. BucketMetadataSys.UpdateExpiryLCConfig now
// performs the read, merge, and save under a single metadata.lock. This test
// drives PeerBucketLCConfigHandler and asserts the concurrent change survives.
// ---------------------------------------------------------------------------
type lcMergeWriterKey struct{}
// lcMergeBarrier pauses the merge writer (context value "M") exactly when it
// tries to take metadata.lock for its persisting Update. By that point the
// merge has already read the stale lifecycle document, so the test can commit a
// concurrent transition change before releasing the merge write.
type lcMergeBarrier struct {
ObjectLayer
bucket string
mAtLock chan struct{}
mProceed chan struct{}
mOnce sync.Once
}
func (o *lcMergeBarrier) metadataLock() string {
return pathJoin(bucketMetaPrefix, o.bucket, "metadata.lock")
}
func (o *lcMergeBarrier) NewNSLock(bucket string, objects ...string) RWLocker {
lock := o.ObjectLayer.NewNSLock(bucket, objects...)
if bucket != minioMetaBucket || len(objects) != 1 || objects[0] != o.metadataLock() {
return lock
}
return metadataObservedRWLocker{RWLocker: lock, onLock: func(ctx context.Context) {
if ctx.Value(lcMergeWriterKey{}) != "M" {
return
}
o.mOnce.Do(func() { close(o.mAtLock) })
select {
case <-o.mProceed:
case <-ctx.Done():
}
}}
}
func TestLifecycleExpiryMergeRaceLosesConcurrentTransition(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testLifecycleExpiryMergeRaceLosesConcurrentTransition,
})
}
func testLifecycleExpiryMergeRaceLosesConcurrentTransition(obj ObjectLayer, instanceType, bucket string,
_ http.Handler, _ auth.Credentials, t *testing.T,
) {
// Revision N: a single transition-only rule.
baseXML := []byte(`<LifecycleConfiguration><Rule><ID>keep</ID><Filter><Prefix>data/</Prefix></Filter><Status>Enabled</Status><Transition><Days>30</Days><StorageClass>WARM</StorageClass></Transition></Rule></LifecycleConfiguration>`)
if _, err := globalBucketMetadataSys.Update(t.Context(), bucket, bucketLifecycleConfig, baseXML); err != nil {
t.Fatalf("%s: seed lifecycle: %v", instanceType, err)
}
// A replicated expiry-only rule arriving from a peer site.
expXML := `<LifecycleConfiguration><Rule><ID>expire</ID><Filter><Prefix>tmp/</Prefix></Filter><Status>Enabled</Status><Expiration><Days>7</Days></Expiration></Rule></LifecycleConfiguration>`
expLCConfig := base64.StdEncoding.EncodeToString([]byte(expXML))
previousObjectAPI := newObjectLayerFn()
barrier := &lcMergeBarrier{
ObjectLayer: obj,
bucket: bucket,
mAtLock: make(chan struct{}),
mProceed: make(chan struct{}),
}
setObjectLayer(barrier)
defer setObjectLayer(previousObjectAPI)
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
defer cancel()
mCtx := context.WithValue(ctx, lcMergeWriterKey{}, "M")
mDone := make(chan error, 1)
go func() {
mDone <- globalSiteReplicationSys.PeerBucketLCConfigHandler(mCtx, bucket, &expLCConfig, UTCNow())
}()
// Wait until the merge writer has read revision N and is about to persist.
select {
case <-barrier.mAtLock:
case err := <-mDone:
t.Fatalf("%s: merge writer finished before persisting: %v", instanceType, err)
case <-ctx.Done():
t.Fatalf("%s: merge writer never reached metadata.lock: %v", instanceType, ctx.Err())
}
// Concurrent local lifecycle transition change commits revision N+1.
concurrentXML := []byte(`<LifecycleConfiguration><Rule><ID>keep</ID><Filter><Prefix>data/</Prefix></Filter><Status>Enabled</Status><Transition><Days>10</Days><StorageClass>COLD</StorageClass></Transition></Rule></LifecycleConfiguration>`)
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketLifecycleConfig, concurrentXML); err != nil {
t.Fatalf("%s: concurrent transition update: %v", instanceType, err)
}
// Release the merge write so it lands after the concurrent commit.
close(barrier.mProceed)
if err := <-mDone; err != nil {
t.Fatalf("%s: merge writer failed: %v", instanceType, err)
}
// The persisted lifecycle must contain both the replicated expiry rule and
// the concurrent transition change.
cfg, _, err := globalBucketMetadataSys.GetLifecycleConfig(bucket)
if err != nil {
t.Fatalf("%s: read merged lifecycle: %v", instanceType, err)
}
var keep, expire bool
for i := range cfg.Rules {
switch cfg.Rules[i].ID {
case "keep":
keep = true
if cfg.Rules[i].Transition.Days != 10 || cfg.Rules[i].Transition.StorageClass != "COLD" {
t.Fatalf("%s: lifecycle merge overwrote the concurrent transition change: got Days=%d StorageClass=%q, want Days=10 StorageClass=COLD",
instanceType, cfg.Rules[i].Transition.Days, cfg.Rules[i].Transition.StorageClass)
}
case "expire":
expire = true
}
}
if !expire {
t.Fatalf("%s: merged lifecycle dropped the replicated expiry rule: %+v", instanceType, cfg.Rules)
}
if !keep {
t.Fatalf("%s: merged lifecycle dropped the transition rule entirely: %+v", instanceType, cfg.Rules)
}
}
// ---------------------------------------------------------------------------
// Target 2 (issue #105): DeleteBucket racing an already in-flight metadata
// writer and resurrecting a ghost .metadata.bin record.
//
// erasureServerPools.DeleteBucket takes only <bucket>.lck and purges the whole
// metadata prefix with deleteAll. Config writers (updateAndParse) take only
// metadata.lock. The two locks do not exclude each other, so a writer that is
// past its read and holding metadata.lock can persist .metadata.bin AFTER the
// delete has purged the prefix, resurrecting a record for a deleted bucket.
// ---------------------------------------------------------------------------
func TestDeleteBucketResurrectsGhostMetadata(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testDeleteBucketResurrectsGhostMetadata,
})
}
func testDeleteBucketResurrectsGhostMetadata(obj ObjectLayer, instanceType, bucket string,
_ http.Handler, _ auth.Credentials, t *testing.T,
) {
previousObjectAPI := newObjectLayerFn()
barrier := &metadataRMWBarrierObjectLayer{
ObjectLayer: obj,
bucket: bucket,
aReady: make(chan struct{}),
aRelease: make(chan struct{}),
bLockAttempt: make(chan struct{}),
}
setObjectLayer(barrier)
defer setObjectLayer(previousObjectAPI)
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
defer cancel()
aCtx := context.WithValue(ctx, metadataRMWWriterKey{}, "A")
policyJSON := []byte(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::` + bucket + `/*"}]}`)
aDone := make(chan error, 1)
// Writer A holds metadata.lock and pauses at the .metadata.bin PutObject.
go func() {
_, err := globalBucketMetadataSys.Update(aCtx, bucket, bucketPolicyConfig, policyJSON)
aDone <- err
}()
select {
case <-barrier.aReady:
case err := <-aDone:
t.Fatalf("%s: writer A finished before persisting: %v", instanceType, err)
case <-ctx.Done():
t.Fatalf("%s: writer A never reached metadata save: %v", instanceType, ctx.Err())
}
// Delete the bucket while writer A is paused mid-save.
delDone := make(chan error, 1)
go func() {
delDone <- obj.DeleteBucket(ctx, bucket, DeleteBucketOptions{Force: true})
}()
select {
case err := <-delDone:
// Unfixed path: DeleteBucket does not serialize on metadata.lock, so it
// purges the prefix immediately. Release A so its save recreates the file.
if err != nil {
t.Fatalf("%s: delete bucket: %v", instanceType, err)
}
close(barrier.aRelease)
if err := <-aDone; err != nil {
t.Fatalf("%s: writer A failed: %v", instanceType, err)
}
case <-time.After(3 * time.Second):
// Fixed path: DeleteBucket is blocked waiting for metadata.lock held by
// A. Release A so it finishes, then the purge runs after the save.
close(barrier.aRelease)
if err := <-aDone; err != nil {
t.Fatalf("%s: writer A failed: %v", instanceType, err)
}
if err := <-delDone; err != nil {
t.Fatalf("%s: delete bucket: %v", instanceType, err)
}
}
// After DeleteBucket, no .metadata.bin record may remain on disk.
if meta, err := readBucketMetadata(ctx, obj, bucket); err == nil {
t.Fatalf("%s: ghost .metadata.bin resurrected after DeleteBucket: name=%q created=%s policyLen=%d",
instanceType, meta.Name, meta.Created, len(meta.PolicyConfigJSON))
} else if !errors.Is(err, errConfigNotFound) && !isErrBucketNotFound(err) && !errors.Is(err, errVolumeNotFound) {
t.Fatalf("%s: unexpected error reading deleted bucket metadata: %v", instanceType, err)
}
}
// ---------------------------------------------------------------------------
// Target 3 (issue #105): overlapping peer reloads publishing an older cache
// record after a newer save, leaving the resident cache one revision behind
// until refresh.
//
// LoadBucketMetadataHandler does loadBucketMetadata (read) followed by a publish
// into the resident cache. Before the fix that publish was an unconditional
// BucketMetadataSys.Set, so a reload that read an older on-disk revision could
// overwrite a newer resident record when it published late (unlike
// refreshBucketsMetadataLoop, which guards on lastUpdate()). The publish is now
// BucketMetadataSys.setReloaded, which refuses to regress a newer resident copy.
//
// The peer handler hardcodes context.Background(), so its reload core is mirrored
// here to inject deterministic ordering. This is a resident-cache freshness
// defect only: the persisted record always stays correct.
// ---------------------------------------------------------------------------
type reloadWriterKey struct{}
// reloadStaleBarrier captures the old on-disk metadata for the reload tagged
// "R1" and holds R1's read open until released, so a newer revision can be
// written and published before R1 finishes publishing the stale copy.
type reloadStaleBarrier struct {
ObjectLayer
bucket string
r1Reading chan struct{}
r1Release chan struct{}
once sync.Once
}
func (o *reloadStaleBarrier) metadataObject() string {
return pathJoin(bucketMetaPrefix, o.bucket, bucketMetadataFile)
}
func (o *reloadStaleBarrier) GetObjectNInfo(ctx context.Context, bucket, object string, rs *HTTPRangeSpec, h http.Header, opts ObjectOptions) (*GetObjectReader, error) {
if bucket == minioMetaBucket && object == o.metadataObject() && ctx.Value(reloadWriterKey{}) == "R1" {
gr, err := o.ObjectLayer.GetObjectNInfo(ctx, bucket, object, rs, h, opts)
if err != nil {
return nil, err
}
data, rerr := io.ReadAll(gr)
oi := gr.ObjInfo
gr.Close()
if rerr != nil {
return nil, rerr
}
o.once.Do(func() { close(o.r1Reading) })
select {
case <-o.r1Release:
case <-ctx.Done():
return nil, ctx.Err()
}
return NewGetObjectReaderFromReader(bytes.NewReader(data), oi, opts)
}
return o.ObjectLayer.GetObjectNInfo(ctx, bucket, object, rs, h, opts)
}
func TestOverlappingReloadPublishesStaleResidentCache(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testOverlappingReloadPublishesStaleResidentCache,
})
}
func residentTagValue(t *testing.T, instanceType, bucket string) string {
t.Helper()
cfg, _, err := globalBucketMetadataSys.GetTaggingConfig(bucket)
if err != nil {
t.Fatalf("%s: read resident tagging: %v", instanceType, err)
}
return cfg.ToMap()["rev"]
}
func testOverlappingReloadPublishesStaleResidentCache(obj ObjectLayer, instanceType, bucket string,
_ http.Handler, _ auth.Credentials, t *testing.T,
) {
// Revision 0 on disk and in the resident cache.
tag0 := []byte(`<Tagging><TagSet><Tag><Key>rev</Key><Value>0</Value></Tag></TagSet></Tagging>`)
if _, err := globalBucketMetadataSys.Update(t.Context(), bucket, bucketTaggingConfig, tag0); err != nil {
t.Fatalf("%s: seed rev0: %v", instanceType, err)
}
previousObjectAPI := newObjectLayerFn()
barrier := &reloadStaleBarrier{
ObjectLayer: obj,
bucket: bucket,
r1Reading: make(chan struct{}),
r1Release: make(chan struct{}),
}
setObjectLayer(barrier)
defer setObjectLayer(previousObjectAPI)
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
defer cancel()
r1Ctx := context.WithValue(ctx, reloadWriterKey{}, "R1")
// reload mirrors the core of LoadBucketMetadataHandler (which hardcodes
// context.Background() and cannot take an injected context). The publish
// uses setReloaded, exactly as the fixed handler does.
reload := func(rctx context.Context) error {
meta, err := loadBucketMetadata(rctx, newObjectLayerFn(), bucket)
if err != nil {
return err
}
globalBucketMetadataSys.setReloaded(bucket, meta)
return nil
}
// R1: an overlapping peer reload that reads rev0 and then stalls before it
// publishes.
r1Done := make(chan error, 1)
go func() { r1Done <- reload(r1Ctx) }()
select {
case <-barrier.r1Reading:
case err := <-r1Done:
t.Fatalf("%s: R1 finished before publishing: %v", instanceType, err)
case <-ctx.Done():
t.Fatalf("%s: R1 never read metadata: %v", instanceType, ctx.Err())
}
// A newer save commits revision 1 to disk and the resident cache.
tag1 := []byte(`<Tagging><TagSet><Tag><Key>rev</Key><Value>1</Value></Tag></TagSet></Tagging>`)
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketTaggingConfig, tag1); err != nil {
t.Fatalf("%s: commit rev1: %v", instanceType, err)
}
// R2: a second overlapping peer reload that reads rev1 and publishes it.
if err := reload(ctx); err != nil {
t.Fatalf("%s: R2 reload: %v", instanceType, err)
}
if got := residentTagValue(t, instanceType, bucket); got != "1" {
t.Fatalf("%s: precondition failed, resident cache should be rev1 before R1 publishes, got %q", instanceType, got)
}
// Release R1 so its stale publish lands after the newer save.
close(barrier.r1Release)
if err := <-r1Done; err != nil {
t.Fatalf("%s: R1 reload: %v", instanceType, err)
}
// The persisted record must stay at rev1 (persistent correctness).
if dcfg, perr := globalBucketMetadataSys.GetConfigFromDisk(ctx, bucket); perr != nil {
t.Fatalf("%s: reload disk metadata: %v", instanceType, perr)
} else if got := dcfg.taggingConfig.ToMap()["rev"]; got != "1" {
t.Fatalf("%s: persisted record regressed to rev %q, want rev1", instanceType, got)
}
// The resident cache must not be left behind at rev0.
if got := residentTagValue(t, instanceType, bucket); got != "1" {
t.Fatalf("%s: overlapping reload left resident cache at rev %q, want rev1", instanceType, got)
}
}
+80 -1
View File
@@ -125,6 +125,28 @@ func (sys *BucketMetadataSys) Set(bucket string, meta BucketMetadata) {
}
}
// setReloaded publishes a record freshly loaded from disk into the resident
// cache without letting an older on-disk revision overwrite a newer resident
// one. Overlapping peer reloads (LoadBucketMetadataHandler) and cache-miss
// loads can finish out of order, so an unconditional Set can leave the resident
// cache a revision behind until the next refresh (issue #105). The authoritative
// read-modify-write save path uses Set directly and always wins because it
// stamps a fresh updatedAt. Only a shallow copy is stored, exactly like Set.
func (sys *BucketMetadataSys) setReloaded(bucket string, meta BucketMetadata) {
if isMinioMetaBucketName(bucket) {
return
}
sys.Lock()
defer sys.Unlock()
if cur, ok := sys.metadataMap[bucket]; ok && !cur.lastUpdate().Before(meta.lastUpdate()) {
// A resident revision that is at least as new is already published; do
// not regress it to the older reload.
return
}
sys.metadataMap[bucket] = meta
sys.clearLoadFailure(bucket)
}
func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string, configFile string, configData []byte, parse, lifecycleDelete bool) (updatedAt time.Time, err error) {
objAPI := newObjectLayerFn()
if objAPI == nil {
@@ -292,6 +314,57 @@ func (sys *BucketMetadataSys) Update(ctx context.Context, bucket string, configF
return sys.updateAndParse(ctx, bucket, configFile, configData, true, false)
}
// UpdateExpiryLCConfig merges a replicated ILM expiry configuration with the
// bucket's current lifecycle document and persists the merged result while
// holding metadata.lock across the read, merge, and save. The site-replication
// expiry heal and peer-apply paths must use this instead of computing the merge
// from an unlocked GetConfigFromDisk read and then writing it with Update: that
// two-step sequence drops any lifecycle transition change committed in between
// (issue #105). Lock order stays <bucket>.lck -> metadata.lock -> .metadata.bin;
// the merge and save run under metadata.lock and the peer fan-out runs after it
// is released.
func (sys *BucketMetadataSys) UpdateExpiryLCConfig(ctx context.Context, bucket string, expLCConfig *string, updatedAt time.Time) error {
objAPI := newObjectLayerFn()
if objAPI == nil {
return errServerNotInitialized
}
if isMinioMetaBucketName(bucket) {
return errInvalidArgument
}
notifyCtx := ctx
ctx, unlock, err := lockBucketMetadata(ctx, objAPI, bucket)
if err != nil {
return err
}
err = func() error {
defer unlock()
meta, err := loadBucketMetadataParse(ctx, objAPI, bucket, true)
if err != nil {
if !globalIsErasure && !globalIsDistErasure && errors.Is(err, errVolumeNotFound) {
// Only single drive mode needs this fallback.
meta = newBucketMetadata(bucket)
} else {
return err
}
}
configData, err := mergeExpiryWithLCConfig(bucket, meta, expLCConfig, updatedAt)
if err != nil {
return err
}
meta.LifecycleConfigXML = configData
meta.LifecycleConfigUpdatedAt = UTCNow()
return sys.saveMetadata(ctx, objAPI, meta)
}()
if err != nil {
return err
}
globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), bucket) // Do not use caller context here
return nil
}
// Get metadata for a bucket.
// If no metadata exists errConfigNotFound is returned and a new metadata is returned.
// Only a shallow copy is returned, so referenced data should not be modified,
@@ -602,7 +675,13 @@ func (sys *BucketMetadataSys) GetConfig(ctx context.Context, bucket string) (met
return meta, false, err
}
sys.Lock()
sys.metadataMap[bucket] = meta
if cur, ok := sys.metadataMap[bucket]; ok && !cur.lastUpdate().Before(meta.lastUpdate()) {
// A concurrent publish installed a resident revision at least as new as
// this cache-miss load; return it instead of regressing (issue #105).
meta = cur
} else {
sys.metadataMap[bucket] = meta
}
sys.clearLoadFailure(bucket)
sys.Unlock()
+12 -2
View File
@@ -2204,8 +2204,18 @@ func (z *erasureServerPools) DeleteBucket(ctx context.Context, bucket string, op
}
if err == nil {
// Purge the entire bucket metadata entirely.
z.deleteAll(context.Background(), minioMetaBucket, pathJoin(bucketMetaPrefix, bucket))
// Purge the entire bucket metadata entirely. Hold metadata.lock across
// the purge so a bucket metadata writer that is mid-save cannot
// resurrect a ghost .metadata.bin after the prefix has been removed
// (issue #105). Lock order stays <bucket>.lck -> metadata.lock; a lock
// acquisition failure falls back to a best-effort unlocked purge so a
// delete is never blocked from completing.
if lctx, unlock, lerr := lockBucketMetadata(context.Background(), z, bucket); lerr == nil {
z.deleteAll(lctx, minioMetaBucket, pathJoin(bucketMetaPrefix, bucket))
unlock()
} else {
z.deleteAll(context.Background(), minioMetaBucket, pathJoin(bucketMetaPrefix, bucket))
}
}
return toObjectErr(err, bucket)
+3 -1
View File
@@ -544,7 +544,9 @@ func (s *peerRESTServer) LoadBucketMetadataHandler(mss *grid.MSS) (np grid.NoPay
return np, grid.NewRemoteErr(err)
}
globalBucketMetadataSys.Set(bucketName, meta)
// Publish monotonically: an overlapping reload that read an older revision
// must not overwrite a newer resident record (issue #105).
globalBucketMetadataSys.setReloaded(bucketName, meta)
if meta.notificationConfig != nil {
globalEventNotifier.AddRulesMap(bucketName, meta.notificationConfig.ToRulesMap())
+20 -24
View File
@@ -45,6 +45,7 @@ import (
"github.com/minio/minio/internal/auth"
"github.com/minio/minio/internal/bucket/cors"
"github.com/minio/minio/internal/bucket/lifecycle"
objectlock "github.com/minio/minio/internal/bucket/object/lock"
sreplication "github.com/minio/minio/internal/bucket/replication"
"github.com/minio/minio/internal/bucket/versioning"
"github.com/minio/minio/internal/logger"
@@ -2115,12 +2116,7 @@ func (c *SiteReplicationSys) PeerBucketLCConfigHandler(ctx context.Context, buck
}
if expLCConfig != nil {
configData, err := mergeWithCurrentLCConfig(ctx, bucket, expLCConfig, updatedAt)
if err != nil {
return wrapSRErr(err)
}
_, err = globalBucketMetadataSys.Update(ctx, bucket, bucketLifecycleConfig, configData)
if err != nil {
if err := globalBucketMetadataSys.UpdateExpiryLCConfig(ctx, bucket, expLCConfig, updatedAt); err != nil {
return wrapSRErr(err)
}
return nil
@@ -4922,13 +4918,12 @@ func (c *SiteReplicationSys) healBucketILMExpiry(ctx context.Context, objAPI Obj
continue
}
finalConfigData, err := mergeWithCurrentLCConfig(ctx, bucket, latestExpLCConfig, lastUpdate)
if err != nil {
return wrapSRErr(err)
}
if dID == globalDeploymentID() {
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketLifecycleConfig, finalConfigData); err != nil {
// Merge and persist atomically under metadata.lock so a concurrent
// lifecycle transition change is not lost (issue #105). The merged
// blob is recomputed from the locked current document here rather
// than from an earlier unlocked read.
if err := globalBucketMetadataSys.UpdateExpiryLCConfig(ctx, bucket, latestExpLCConfig, lastUpdate); err != nil {
replLogIf(ctx, fmt.Errorf("Unable to heal bucket ILM expiry data from peer site %s : %w", latestPeerName, err))
}
continue
@@ -6594,13 +6589,12 @@ func (c *SiteReplicationSys) getSiteMetrics(ctx context.Context) (madmin.SRMetri
return sm, nil
}
// mergeWithCurrentLCConfig - merges the given ilm expiry configuration with existing for the current site and returns
func mergeWithCurrentLCConfig(ctx context.Context, bucket string, expLCCfg *string, updatedAt time.Time) ([]byte, error) {
// Get bucket config from current site
meta, e := globalBucketMetadataSys.GetConfigFromDisk(ctx, bucket)
if e != nil && !errors.Is(e, errConfigNotFound) {
return []byte{}, e
}
// mergeExpiryWithLCConfig merges the given ILM expiry configuration with the
// supplied current bucket lifecycle document and returns the merged XML. The
// current document must be read by the caller while holding metadata.lock (see
// BucketMetadataSys.UpdateExpiryLCConfig) so a concurrent lifecycle transition
// change cannot be lost between the read and the merged save (issue #105).
func mergeExpiryWithLCConfig(bucket string, meta BucketMetadata, expLCCfg *string, updatedAt time.Time) ([]byte, error) {
rMap := make(map[string]lifecycle.Rule)
var xmlName xml.Name
if len(meta.LifecycleConfigXML) > 0 {
@@ -6611,7 +6605,7 @@ func mergeWithCurrentLCConfig(ctx context.Context, bucket string, expLCCfg *stri
for _, rl := range lcCfg.Rules {
rMap[rl.ID] = rl
}
xmlName = meta.lifecycleConfig.XMLName
xmlName = lcCfg.XMLName
}
// get latest expiry rules
@@ -6683,11 +6677,13 @@ func mergeWithCurrentLCConfig(ctx context.Context, bucket string, expLCCfg *stri
ExpiryUpdatedAt: &updatedAt,
}
rcfg, err := globalBucketObjectLockSys.Get(bucket)
if err != nil {
return nil, err
// Validate against the object-lock retention from the locked metadata
// snapshot rather than re-reading it, which would risk a re-entrant
// metadata load while metadata.lock is held (issue #105).
var rcfg objectlock.Retention
if meta.objectLockConfig != nil {
rcfg = meta.objectLockConfig.ToRetention()
}
if err := finalLcCfg.Validate(rcfg); err != nil {
return []byte{}, err
}