diff --git a/cmd/bucket-metadata-publication_test.go b/cmd/bucket-metadata-publication_test.go new file mode 100644 index 000000000..ee053dfde --- /dev/null +++ b/cmd/bucket-metadata-publication_test.go @@ -0,0 +1,201 @@ +// Copyright 2026 PGSTY contributors. +// SPDX-License-Identifier: AGPL-3.0-or-later + +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/grid" +) + +func TestDeleteBucketMetadataLockCancellation(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: testDeleteBucketMetadataLockCancellation}) +} + +func testDeleteBucketMetadataLockCancellation(obj ObjectLayer, instanceType, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) { + _, unlock, err := lockBucketMetadata(t.Context(), obj, bucket) + if err != nil { + t.Fatal(err) + } + release := sync.OnceFunc(unlock) + defer release() + ctx, cancel := context.WithTimeout(t.Context(), 250*time.Millisecond) + 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) + } +} + +func TestQueuedMetadataUpdateAfterDelete(t *testing.T) { + defer DetectTestLeak(t)() + for _, expiry := range []bool{false, true} { + t.Run(fmt.Sprintf("expiry=%v", expiry), func(t *testing.T) { + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, instanceType, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) { + previous := newObjectLayerFn() + barrier := &lcMergeBarrier{ObjectLayer: obj, bucket: bucket, mAtLock: make(chan struct{}), mProceed: make(chan struct{})} + setObjectLayer(barrier) + defer setObjectLayer(previous) + release := sync.OnceFunc(func() { close(barrier.mProceed) }) + defer release() + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + ctx = context.WithValue(ctx, lcMergeWriterKey{}, "M") + done := make(chan error, 1) + go func() { + if expiry { + done <- globalBucketMetadataSys.UpdateExpiryLCConfig(ctx, bucket, nil, UTCNow()) + return + } + _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketTaggingConfig, []byte(``)) + done <- err + }() + select { + case <-barrier.mAtLock: + case <-ctx.Done(): + t.Fatal("writer did not reach metadata.lock") + } + if err := obj.DeleteBucket(t.Context(), bucket, DeleteBucketOptions{Force: true}); err != nil { + t.Fatal(err) + } + release() + if err := <-done; !isErrBucketNotFound(err) { + t.Errorf("%s: queued update should reject a deleted bucket, got %v", instanceType, err) + } + if _, err := readBucketMetadata(t.Context(), obj, bucket); !errors.Is(err, errConfigNotFound) && !isErrBucketNotFound(err) { + t.Errorf("%s: queued update recreated metadata: %v", instanceType, err) + } + }}) + }) + } +} + +// 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() + notification := []byte(`` + revision + `arn:minio:sqs::` + revision + `:webhooks3: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-sys.go b/cmd/bucket-metadata-sys.go index 257e77777..60c85b5b7 100644 --- a/cmd/bucket-metadata-sys.go +++ b/cmd/bucket-metadata-sys.go @@ -125,13 +125,10 @@ 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. +// 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 @@ -139,12 +136,21 @@ func (sys *BucketMetadataSys) setReloaded(bucket string, meta BucketMetadata) { 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 + meta = cur + } else { + sys.metadataMap[bucket] = meta } - 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) { @@ -251,6 +257,11 @@ func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) err // saveMetadata persists and publishes metadata locally. Callers performing a // read-modify-write must hold metadata.lock and release it before peer fan-out. func (sys *BucketMetadataSys) saveMetadata(ctx context.Context, objAPI ObjectLayer, meta BucketMetadata) error { + // A writer may have queued for metadata.lock before DeleteBucket completed. + // Recheck the physical bucket under that lock, before recreating metadata. + if _, err := objAPI.GetBucketInfo(ctx, meta.Name, BucketOptions{NoMetadata: true}); err != nil { + return err + } if err := meta.Save(ctx, objAPI); err != nil { return err } @@ -782,8 +793,6 @@ 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) @@ -794,19 +803,7 @@ func (sys *BucketMetadataSys) refreshBucketsMetadataLoop(ctx context.Context) { continue } - 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) - } + sys.setReloaded(bucket, meta) wait() // wait to proceed to next entry. } diff --git a/cmd/bucket-metadata.go b/cmd/bucket-metadata.go index c93ce9907..0bb59bb9e 100644 --- a/cmd/bucket-metadata.go +++ b/cmd/bucket-metadata.go @@ -303,6 +303,9 @@ func loadBucketMetadataParseUnderLock(ctx context.Context, objectAPI ObjectLayer return newBucketMetadata(bucket), fmt.Errorf("%w: %v", errBucketMetadataMigrationLockUnavailable, err) } defer unlock() + if _, err := objectAPI.GetBucketInfo(ctx, bucket, BucketOptions{NoMetadata: true}); err != nil { + return newBucketMetadata(bucket), err + } return loadBucketMetadataParse(ctx, objectAPI, bucket, parse) } diff --git a/cmd/erasure-server-pool.go b/cmd/erasure-server-pool.go index 584e657d3..722a9dcf8 100644 --- a/cmd/erasure-server-pool.go +++ b/cmd/erasure-server-pool.go @@ -2195,7 +2195,15 @@ func (z *erasureServerPools) DeleteBucket(ctx context.Context, bucket string, op opts.Force = true } - err := z.s3Peer.DeleteBucket(ctx, bucket, opts) + // Take the metadata writer lock before deleting anything. Failure or + // cancellation must leave both the bucket and its metadata intact. + ctx, unlock, err := lockBucketMetadata(ctx, z, bucket) + if err != nil { + return toObjectErr(err, bucket) + } + defer unlock() + + err = z.s3Peer.DeleteBucket(ctx, bucket, opts) if err == nil || isErrBucketNotFound(err) { // If site replication is configured, hold on to deleted bucket state until sites sync if opts.SRDeleteOp == MarkDelete { @@ -2204,18 +2212,9 @@ func (z *erasureServerPools) DeleteBucket(ctx context.Context, bucket string, op } if err == nil { - // 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 .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)) - } + // Finish cleanup after a committed delete even if the client disconnects. + // Both the bucket-name and metadata locks remain held until return. + z.deleteAll(context.WithoutCancel(ctx), minioMetaBucket, pathJoin(bucketMetaPrefix, bucket)) } return toObjectErr(err, bucket) diff --git a/cmd/peer-rest-server.go b/cmd/peer-rest-server.go index f5687cf37..2e0f5c5b1 100644 --- a/cmd/peer-rest-server.go +++ b/cmd/peer-rest-server.go @@ -544,18 +544,9 @@ func (s *peerRESTServer) LoadBucketMetadataHandler(mss *grid.MSS) (np grid.NoPay return np, grid.NewRemoteErr(err) } - // Publish monotonically: an overlapping reload that read an older revision - // must not overwrite a newer resident record (issue #105). + // Publish the metadata and derived registries from the same current revision. globalBucketMetadataSys.setReloaded(bucketName, meta) - if meta.notificationConfig != nil { - globalEventNotifier.AddRulesMap(bucketName, meta.notificationConfig.ToRulesMap()) - } - - if meta.bucketTargetConfig != nil { - globalBucketTargetSys.UpdateAllTargets(bucketName, meta.bucketTargetConfig) - } - return np, nerr }