From 75ba0ce402e915fdbd236e9f40a3c48d6b592d63 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Tue, 8 Sep 2026 15:42:50 +0800 Subject: [PATCH] fix: drop the broken bucket-metadata reload publication guard (issue #105 T3) The #105 T3 change (PR #156) tried to keep the resident metadata cache monotonic by guarding peer-reload publication on lastUpdate(). But lastUpdate() is the max of per-config timestamps and cannot order whole records: a node caching {policy@20, CORS@10} that receives a newer CORS@15 still has lastUpdate()==20, so the guard rejects the legitimately-newer record and the periodic refresh (same comparator) cannot repair it. A paused reload could also resurrect deleted resident state. Per the maintainer decision, revert the reload publication to its original unconditional (acceptable-until-refresh) behavior: - remove setReloaded and restore the plain Set plus notification/target registry updates in LoadBucketMetadataHandler; - restore refreshBucketsMetadataLoop's own lastUpdate() staleness check and globalEventNotifier.set / globalBucketTargetSys.set publication; - restore the unconditional GetConfig cache-miss publication; - document the known freshness limitation at the reload site (the periodic refresh is best-effort and cannot repair an equal-maximum-timestamp divergence). The T1 lifecycle merge-under-lock (UpdateExpiryLCConfig) and both T2 fixes (DeleteBucket takes metadata.lock before deleting; saveMetadata and loadBucketMetadataParseUnderLock recheck physical bucket existence) are kept fully intact. Tests: - drop the T3 reproductions (overlapping-reload resident-cache test and the peer-reload-preserves-current-targets publication test); - add lockBucketMetadataAcquireHook, a nil-in-production atomic test hook in the shared metadata.lock path, so tests can deterministically observe a caller (notably DeleteBucket, whose lock is taken through its erasureServerPools receiver and is invisible to an injected object layer) reaching the lock; - rewrite the T2 delete-race ghost test to hold metadata.lock MID-SAVE (past saveMetadata's existence recheck) and synchronize on the delete's actual lock attempt via the hook, so it isolates the lock-before-delete fix: removing only DeleteBucket's metadata.lock (recheck kept) now fails it; - rewrite the cancellation test to observe the delete's actual lock attempt, then cancel and await its error while still holding the lock, so a scheduling-delayed delete stopped by the canceled context can no longer pass on a broken tree. Refs #105. Follows #156. Co-Authored-By: Claude Opus 4.8 Signed-off-by: Feng Ruohang --- cmd/bucket-metadata-publication_test.go | 160 ++++------------- cmd/bucket-metadata-race_test.go | 230 ++++++------------------ cmd/bucket-metadata-sys.go | 64 +++---- cmd/peer-rest-server.go | 22 ++- 4 files changed, 132 insertions(+), 344 deletions(-) diff --git a/cmd/bucket-metadata-publication_test.go b/cmd/bucket-metadata-publication_test.go index 541290592..0bfbb60f4 100644 --- a/cmd/bucket-metadata-publication_test.go +++ b/cmd/bucket-metadata-publication_test.go @@ -4,22 +4,15 @@ package cmd import ( - "bytes" "context" - "encoding/json" "errors" "fmt" - "io" "net/http" - "reflect" "sync" "testing" "time" - "github.com/minio/madmin-go/v3" "github.com/minio/minio/internal/auth" - "github.com/minio/minio/internal/event" - "github.com/minio/minio/internal/grid" ) func TestDeleteBucketMetadataLockCancellation(t *testing.T) { @@ -34,21 +27,43 @@ func testDeleteBucketMetadataLockCancellation(obj ObjectLayer, instanceType, buc } release := sync.OnceFunc(unlock) defer release() - ctx, cancel := context.WithTimeout(t.Context(), 250*time.Millisecond) + + // Observe DeleteBucket's ACTUAL metadata.lock attempt. Set the hook after + // our own acquisition above so it only trips on the delete. + delAtLock := make(chan struct{}) + var once sync.Once + hook := func(b string) { + if b == bucket { + once.Do(func() { close(delAtLock) }) + } + } + lockBucketMetadataAcquireHook.Store(&hook) + defer lockBucketMetadataAcquireHook.Store(nil) + + ctx, cancel := context.WithCancel(t.Context()) defer cancel() done := make(chan error, 1) go func() { done <- obj.DeleteBucket(ctx, bucket, DeleteBucketOptions{Force: true, NoLock: true}) }() - <-ctx.Done() - // A metadata-lock failure must happen before the destructive operation. - if _, err := obj.GetBucketInfo(t.Context(), bucket, BucketOptions{}); err != nil { - t.Errorf("%s: bucket disappeared while metadata.lock was unavailable: %v", instanceType, err) - } - release() - if err := <-done; err == nil { - t.Errorf("%s: canceled deletion succeeded", instanceType) - } - if _, err := readBucketMetadata(t.Context(), obj, bucket); err != nil { - t.Errorf("%s: canceled deletion removed metadata: %v", instanceType, err) + + select { + case <-delAtLock: + // Fixed tree: the delete reached metadata.lock and is blocking on the + // lock we hold. Cancel it and confirm it fails WITHOUT deleting, while we + // still hold the lock (release stays deferred until after the checks). + cancel() + if err := <-done; err == nil { + t.Errorf("%s: canceled deletion succeeded", instanceType) + } + if _, err := obj.GetBucketInfo(t.Context(), bucket, BucketOptions{}); err != nil { + t.Errorf("%s: bucket disappeared while metadata.lock was held: %v", instanceType, err) + } + if _, err := readBucketMetadata(t.Context(), obj, bucket); err != nil { + t.Errorf("%s: canceled deletion removed metadata: %v", instanceType, err) + } + case err := <-done: + // Broken tree: the delete finished without ever taking metadata.lock, + // i.e. it did not serialize the destructive operation behind the lock. + t.Errorf("%s: delete bypassed metadata.lock (err=%v)", instanceType, err) } } @@ -94,110 +109,3 @@ func TestQueuedMetadataUpdateAfterDelete(t *testing.T) { }) } } - -// Pause the first actual peer-handler read, without replacing the handler or -// requiring it to accept a test-only context. -type peerMetadataReadBarrier struct { - ObjectLayer - bucket string - once sync.Once - reading, release chan struct{} -} - -func (o *peerMetadataReadBarrier) GetObjectNInfo(ctx context.Context, bucket, object string, rs *HTTPRangeSpec, h http.Header, opts ObjectOptions) (*GetObjectReader, error) { - first := false - if bucket == minioMetaBucket && object == pathJoin(bucketMetaPrefix, o.bucket, bucketMetadataFile) { - o.once.Do(func() { first = true }) - } - gr, err := o.ObjectLayer.GetObjectNInfo(ctx, bucket, object, rs, h, opts) - if err != nil || !first { - return gr, err - } - data, err := io.ReadAll(gr) - oi := gr.ObjInfo - gr.Close() - if err != nil { - return nil, err - } - close(o.reading) - select { - case <-o.release: - case <-ctx.Done(): - return nil, ctx.Err() - } - return NewGetObjectReaderFromReader(bytes.NewReader(data), oi, opts) -} - -func TestPeerMetadataReloadPreservesCurrentTargets(t *testing.T) { - defer DetectTestLeak(t)() - ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: testPeerMetadataReloadPreservesCurrentTargets}) -} - -func testPeerMetadataReloadPreservesCurrentTargets(obj ObjectLayer, instanceType, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) { - seed := func(revision string) { - t.Helper() - arn := event.ARN{TargetID: event.TargetID{ID: revision, Name: "webhook"}} - notification := []byte(`` + revision + `` + arn.String() + `s3:ObjectCreated:*`) - targets, err := json.Marshal(madmin.BucketTargets{Targets: []madmin.BucketTarget{{ - SourceBucket: bucket, TargetBucket: bucket, Endpoint: "127.0.0.1:9000", Arn: revision, - Credentials: &madmin.Credentials{AccessKey: "fixture", SecretKey: "fixture-secret"}, - }}}) - if err != nil { - t.Fatal(err) - } - for _, config := range []struct { - name string - data []byte - }{{bucketNotificationConfig, notification}, {bucketTargetsFile, targets}} { - if _, err := globalBucketMetadataSys.Update(t.Context(), bucket, config.name, config.data); err != nil { - t.Fatal(err) - } - } - } - reload := func() error { - args := grid.MSS{peerRESTBucket: bucket} - _, err := (&peerRESTServer{}).LoadBucketMetadataHandler(&args) - if err != nil { - return err - } - return nil - } - seed("old") - previous := newObjectLayerFn() - barrier := &peerMetadataReadBarrier{ObjectLayer: obj, bucket: bucket, reading: make(chan struct{}), release: make(chan struct{})} - setObjectLayer(barrier) - defer setObjectLayer(previous) - release := sync.OnceFunc(func() { close(barrier.release) }) - defer release() - done := make(chan error, 1) - go func() { done <- reload() }() - select { - case <-barrier.reading: - case <-time.After(10 * time.Second): - t.Fatal("peer handler did not read metadata") - } - seed("new") - if err := reload(); err != nil { - t.Fatal(err) - } - release() - if err := <-done; err != nil { - t.Fatal(err) - } - meta, err := globalBucketMetadataSys.Get(bucket) - if err != nil { - t.Fatal(err) - } - globalEventNotifier.RLock() - rules := globalEventNotifier.bucketRulesMap[bucket].Clone() - globalEventNotifier.RUnlock() - if !reflect.DeepEqual(rules, meta.notificationConfig.ToRulesMap()) { - t.Errorf("%s: stale peer reload replaced current notification rules", instanceType) - } - globalBucketTargetSys.RLock() - targets := append([]madmin.BucketTarget(nil), globalBucketTargetSys.targetsMap[bucket]...) - globalBucketTargetSys.RUnlock() - if len(targets) != 1 || targets[0].Arn != "new" { - t.Errorf("%s: stale peer reload replaced current replication targets: %+v", instanceType, targets) - } -} diff --git a/cmd/bucket-metadata-race_test.go b/cmd/bucket-metadata-race_test.go index 755ea2ca3..fcd447fae 100644 --- a/cmd/bucket-metadata-race_test.go +++ b/cmd/bucket-metadata-race_test.go @@ -11,11 +11,9 @@ package cmd import ( - "bytes" "context" "encoding/base64" "errors" - "io" "net/http" "sync" "testing" @@ -163,14 +161,24 @@ func testLifecycleExpiryMergeRaceLosesConcurrentTransition(obj ObjectLayer, inst } // --------------------------------------------------------------------------- -// Target 2 (issue #105): DeleteBucket racing an already in-flight metadata -// writer and resurrecting a ghost .metadata.bin record. +// Target 2 (issue #105): DeleteBucket racing an in-flight metadata writer must +// not resurrect a ghost .metadata.bin record. // -// erasureServerPools.DeleteBucket takes only .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. +// Before the fix, erasureServerPools.DeleteBucket took only .lck and +// purged the metadata prefix while config writers (updateAndParse) took only +// metadata.lock, so a writer already MID-SAVE (holding metadata.lock, past +// saveMetadata's existence recheck) could persist .metadata.bin after the purge. +// DeleteBucket now takes metadata.lock before deleting, so it waits for that +// writer and then purges whatever the writer wrote. +// +// This test isolates the DeleteBucket-lock fix specifically: the writer holds +// metadata.lock and is paused at the .metadata.bin PutObject, so the earlier +// saveMetadata existence recheck cannot save it — only serializing the delete +// behind the writer can. Removing just DeleteBucket's metadata.lock (keeping the +// recheck) therefore makes this test fail. The delete's ACTUAL metadata.lock +// attempt is observed with lockBucketMetadataAcquireHook (its lock is taken +// through the erasureServerPools receiver, invisible to the object-layer +// barrier), so the handshake is deterministic with no timing assumption. // --------------------------------------------------------------------------- func TestDeleteBucketResurrectsGhostMetadata(t *testing.T) { @@ -185,6 +193,8 @@ func testDeleteBucketResurrectsGhostMetadata(obj ObjectLayer, instanceType, buck _ http.Handler, _ auth.Credentials, t *testing.T, ) { previousObjectAPI := newObjectLayerFn() + // Writer A holds metadata.lock and pauses at the .metadata.bin PutObject, + // i.e. already past saveMetadata's existence recheck and mid-save. barrier := &metadataRMWBarrierObjectLayer{ ObjectLayer: obj, bucket: bucket, @@ -200,8 +210,9 @@ func testDeleteBucketResurrectsGhostMetadata(obj ObjectLayer, instanceType, buck aCtx := context.WithValue(ctx, metadataRMWWriterKey{}, "A") policyJSON := []byte(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::` + bucket + `/*"}]}`) + aReleased := sync.OnceFunc(func() { close(barrier.aRelease) }) + defer aReleased() 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 @@ -215,33 +226,46 @@ func testDeleteBucketResurrectsGhostMetadata(obj ObjectLayer, instanceType, buck t.Fatalf("%s: writer A never reached metadata save: %v", instanceType, ctx.Err()) } - // Delete the bucket while writer A is paused mid-save. + // A now holds metadata.lock mid-save. Observe DeleteBucket's ACTUAL + // metadata.lock attempt via the acquire hook, set only now so A's earlier + // acquisition does not trip it. + delAtLock := make(chan struct{}) + var once sync.Once + hook := func(b string) { + if b == bucket { + once.Do(func() { close(delAtLock) }) + } + } + lockBucketMetadataAcquireHook.Store(&hook) + defer lockBucketMetadataAcquireHook.Store(nil) + 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) + case <-delAtLock: + // Fixed tree: DeleteBucket reached metadata.lock and blocks on A. Release + // A so it finishes its save and unlocks; the delete then acquires the + // lock and purges the record A wrote. + aReleased() 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) + t.Fatalf("%s: writer A save failed while holding metadata.lock: %v", instanceType, err) } if err := <-delDone; err != nil { t.Fatalf("%s: delete bucket: %v", instanceType, err) } + case err := <-delDone: + // Broken tree: DeleteBucket purged without taking metadata.lock. Release + // A so its mid-save PutObject recreates .metadata.bin (the ghost). + if err != nil { + t.Fatalf("%s: delete bucket: %v", instanceType, err) + } + aReleased() + if err := <-aDone; err != nil { + t.Fatalf("%s: writer A save failed: %v", instanceType, err) + } } // After DeleteBucket, no .metadata.bin record may remain on disk. @@ -252,157 +276,3 @@ func testDeleteBucketResurrectsGhostMetadata(obj ObjectLayer, instanceType, buck 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(`rev0`) - 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(`rev1`) - 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) - } -} diff --git a/cmd/bucket-metadata-sys.go b/cmd/bucket-metadata-sys.go index 60c85b5b7..8096e717d 100644 --- a/cmd/bucket-metadata-sys.go +++ b/cmd/bucket-metadata-sys.go @@ -24,6 +24,7 @@ import ( "fmt" "math/rand" "sync" + "sync/atomic" "time" "github.com/minio/madmin-go/v3" @@ -125,34 +126,6 @@ func (sys *BucketMetadataSys) Set(bucket string, meta BucketMetadata) { } } -// setReloaded publishes the newest known revision and its derived registries -// together. An old reload may finish after a local save or another reload; -// applying its notification/replication targets would regress live behavior -// even if the metadata cache itself rejected that old revision. -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()) { - meta = cur - } else { - sys.metadataMap[bucket] = meta - } - sys.clearLoadFailure(bucket) - // These registry updates only change local state; no peer/network I/O runs - // under the metadata mutex. Keep publication ordered against Set/Remove. - if globalEventNotifier != nil { - if meta.notificationConfig != nil { - globalEventNotifier.AddRulesMap(bucket, meta.notificationConfig.ToRulesMap()) - } else { - globalEventNotifier.RemoveNotification(bucket) - } - } - globalBucketTargetSys.UpdateAllTargets(bucket, meta.bucketTargetConfig) -} - 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 { @@ -273,8 +246,19 @@ func lockBucketMetadata(ctx context.Context, objectAPI ObjectLayer, bucket strin return lockBucketMetadataWithTimeout(ctx, objectAPI, bucket, globalOperationTimeout) } +// lockBucketMetadataAcquireHook, when set, is invoked at the start of every +// metadata.lock acquisition, immediately before the blocking Lock() call. It is +// nil in production (a single atomic load, no behavior change) and exists only +// so tests can deterministically observe a caller reaching the metadata lock — +// notably DeleteBucket, whose lock is taken through its erasureServerPools +// receiver and is therefore invisible to an injected object layer. +var lockBucketMetadataAcquireHook atomic.Pointer[func(bucket string)] + func lockBucketMetadataWithTimeout(ctx context.Context, objectAPI ObjectLayer, bucket string, timeout *dynamicTimeout) (context.Context, func(), error) { lock := objectAPI.NewNSLock(minioMetaBucket, pathJoin(bucketMetaPrefix, bucket, "metadata.lock")) + if hook := lockBucketMetadataAcquireHook.Load(); hook != nil { + (*hook)(bucket) + } lkctx, err := lock.GetLock(ctx, timeout) if err != nil { return nil, nil, err @@ -686,13 +670,7 @@ func (sys *BucketMetadataSys) GetConfig(ctx context.Context, bucket string) (met return meta, false, err } sys.Lock() - 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.metadataMap[bucket] = meta sys.clearLoadFailure(bucket) sys.Unlock() @@ -793,6 +771,8 @@ func (sys *BucketMetadataSys) refreshBucketsMetadataLoop(ctx context.Context) { wait := sleeper.Timer(ctx) bucket := buckets[i].Name + updated := false + meta, err := loadBucketMetadata(ctx, sys.objAPI, bucket) if err != nil { internalLogIf(ctx, err, logger.WarningKind) @@ -803,7 +783,19 @@ func (sys *BucketMetadataSys) refreshBucketsMetadataLoop(ctx context.Context) { continue } - sys.setReloaded(bucket, meta) + sys.Lock() + // Update if the bucket metadata in the memory is older than on-disk one + if lu := sys.metadataMap[bucket].lastUpdate(); lu.Before(meta.lastUpdate()) { + updated = true + sys.metadataMap[bucket] = meta + } + sys.clearLoadFailure(bucket) + sys.Unlock() + + if updated { + globalEventNotifier.set(bucket, meta) + globalBucketTargetSys.set(bucket, meta) + } wait() // wait to proceed to next entry. } diff --git a/cmd/peer-rest-server.go b/cmd/peer-rest-server.go index 2e0f5c5b1..e17bd4b2e 100644 --- a/cmd/peer-rest-server.go +++ b/cmd/peer-rest-server.go @@ -544,8 +544,26 @@ func (s *peerRESTServer) LoadBucketMetadataHandler(mss *grid.MSS) (np grid.NoPay return np, grid.NewRemoteErr(err) } - // Publish the metadata and derived registries from the same current revision. - globalBucketMetadataSys.setReloaded(bucketName, meta) + // Publish the reloaded metadata unconditionally. Overlapping peer reloads + // may briefly leave the resident cache a revision behind, or momentarily + // apply an older reload's derived registries. The resident cache may lag + // until another authoritative update advances it; the periodic metadata + // refresh is only best-effort here and cannot repair a divergence whose + // maximum per-config timestamp already equals the resident record's (its + // staleness check compares the same lastUpdate() and would see no change). + // This acceptable-until-refresh behavior is deliberate: lastUpdate() is the + // max of per-config timestamps and cannot order whole-record revisions, so a + // publication guard keyed on it would wrongly reject legitimately newer + // records (issue #105 follow-up dropped that guard). + globalBucketMetadataSys.Set(bucketName, meta) + + if meta.notificationConfig != nil { + globalEventNotifier.AddRulesMap(bucketName, meta.notificationConfig.ToRulesMap()) + } + + if meta.bucketTargetConfig != nil { + globalBucketTargetSys.UpdateAllTargets(bucketName, meta.bucketTargetConfig) + } return np, nerr }