mirror of
https://github.com/pgsty/minio.git
synced 2026-09-26 04:45:59 +03:00
Merge accepted bucket metadata reload correction
Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
@@ -4,22 +4,15 @@
|
|||||||
package cmd
|
package cmd
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
|
||||||
"net/http"
|
"net/http"
|
||||||
"reflect"
|
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/minio/madmin-go/v3"
|
|
||||||
"github.com/minio/minio/internal/auth"
|
"github.com/minio/minio/internal/auth"
|
||||||
"github.com/minio/minio/internal/event"
|
|
||||||
"github.com/minio/minio/internal/grid"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestDeleteBucketMetadataLockCancellation(t *testing.T) {
|
func TestDeleteBucketMetadataLockCancellation(t *testing.T) {
|
||||||
@@ -34,21 +27,43 @@ func testDeleteBucketMetadataLockCancellation(obj ObjectLayer, instanceType, buc
|
|||||||
}
|
}
|
||||||
release := sync.OnceFunc(unlock)
|
release := sync.OnceFunc(unlock)
|
||||||
defer release()
|
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()
|
defer cancel()
|
||||||
done := make(chan error, 1)
|
done := make(chan error, 1)
|
||||||
go func() { done <- obj.DeleteBucket(ctx, bucket, DeleteBucketOptions{Force: true, NoLock: true}) }()
|
go func() { done <- obj.DeleteBucket(ctx, bucket, DeleteBucketOptions{Force: true, NoLock: true}) }()
|
||||||
<-ctx.Done()
|
|
||||||
// A metadata-lock failure must happen before the destructive operation.
|
select {
|
||||||
if _, err := obj.GetBucketInfo(t.Context(), bucket, BucketOptions{}); err != nil {
|
case <-delAtLock:
|
||||||
t.Errorf("%s: bucket disappeared while metadata.lock was unavailable: %v", instanceType, err)
|
// 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
|
||||||
release()
|
// still hold the lock (release stays deferred until after the checks).
|
||||||
if err := <-done; err == nil {
|
cancel()
|
||||||
t.Errorf("%s: canceled deletion succeeded", instanceType)
|
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)
|
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(`<NotificationConfiguration><QueueConfiguration><Id>` + revision + `</Id><Queue>` + arn.String() + `</Queue><Event>s3:ObjectCreated:*</Event></QueueConfiguration></NotificationConfiguration>`)
|
|
||||||
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -11,11 +11,9 @@
|
|||||||
package cmd
|
package cmd
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
|
||||||
"context"
|
"context"
|
||||||
"encoding/base64"
|
"encoding/base64"
|
||||||
"errors"
|
"errors"
|
||||||
"io"
|
|
||||||
"net/http"
|
"net/http"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
@@ -163,14 +161,24 @@ func testLifecycleExpiryMergeRaceLosesConcurrentTransition(obj ObjectLayer, inst
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
// Target 2 (issue #105): DeleteBucket racing an already in-flight metadata
|
// Target 2 (issue #105): DeleteBucket racing an in-flight metadata writer must
|
||||||
// writer and resurrecting a ghost .metadata.bin record.
|
// not resurrect a ghost .metadata.bin record.
|
||||||
//
|
//
|
||||||
// erasureServerPools.DeleteBucket takes only <bucket>.lck and purges the whole
|
// Before the fix, erasureServerPools.DeleteBucket took only <bucket>.lck and
|
||||||
// metadata prefix with deleteAll. Config writers (updateAndParse) take only
|
// purged the metadata prefix while config writers (updateAndParse) took only
|
||||||
// metadata.lock. The two locks do not exclude each other, so a writer that is
|
// metadata.lock, so a writer already MID-SAVE (holding metadata.lock, past
|
||||||
// past its read and holding metadata.lock can persist .metadata.bin AFTER the
|
// saveMetadata's existence recheck) could persist .metadata.bin after the purge.
|
||||||
// delete has purged the prefix, resurrecting a record for a deleted bucket.
|
// 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) {
|
func TestDeleteBucketResurrectsGhostMetadata(t *testing.T) {
|
||||||
@@ -185,6 +193,8 @@ func testDeleteBucketResurrectsGhostMetadata(obj ObjectLayer, instanceType, buck
|
|||||||
_ http.Handler, _ auth.Credentials, t *testing.T,
|
_ http.Handler, _ auth.Credentials, t *testing.T,
|
||||||
) {
|
) {
|
||||||
previousObjectAPI := newObjectLayerFn()
|
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{
|
barrier := &metadataRMWBarrierObjectLayer{
|
||||||
ObjectLayer: obj,
|
ObjectLayer: obj,
|
||||||
bucket: bucket,
|
bucket: bucket,
|
||||||
@@ -200,8 +210,9 @@ func testDeleteBucketResurrectsGhostMetadata(obj ObjectLayer, instanceType, buck
|
|||||||
aCtx := context.WithValue(ctx, metadataRMWWriterKey{}, "A")
|
aCtx := context.WithValue(ctx, metadataRMWWriterKey{}, "A")
|
||||||
|
|
||||||
policyJSON := []byte(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::` + bucket + `/*"}]}`)
|
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)
|
aDone := make(chan error, 1)
|
||||||
// Writer A holds metadata.lock and pauses at the .metadata.bin PutObject.
|
|
||||||
go func() {
|
go func() {
|
||||||
_, err := globalBucketMetadataSys.Update(aCtx, bucket, bucketPolicyConfig, policyJSON)
|
_, err := globalBucketMetadataSys.Update(aCtx, bucket, bucketPolicyConfig, policyJSON)
|
||||||
aDone <- err
|
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())
|
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)
|
delDone := make(chan error, 1)
|
||||||
go func() {
|
go func() {
|
||||||
delDone <- obj.DeleteBucket(ctx, bucket, DeleteBucketOptions{Force: true})
|
delDone <- obj.DeleteBucket(ctx, bucket, DeleteBucketOptions{Force: true})
|
||||||
}()
|
}()
|
||||||
|
|
||||||
select {
|
select {
|
||||||
case err := <-delDone:
|
case <-delAtLock:
|
||||||
// Unfixed path: DeleteBucket does not serialize on metadata.lock, so it
|
// Fixed tree: DeleteBucket reached metadata.lock and blocks on A. Release
|
||||||
// purges the prefix immediately. Release A so its save recreates the file.
|
// A so it finishes its save and unlocks; the delete then acquires the
|
||||||
if err != nil {
|
// lock and purges the record A wrote.
|
||||||
t.Fatalf("%s: delete bucket: %v", instanceType, err)
|
aReleased()
|
||||||
}
|
|
||||||
close(barrier.aRelease)
|
|
||||||
if err := <-aDone; err != nil {
|
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)
|
||||||
}
|
|
||||||
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 {
|
if err := <-delDone; err != nil {
|
||||||
t.Fatalf("%s: delete bucket: %v", instanceType, err)
|
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.
|
// 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)
|
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)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
+28
-36
@@ -24,6 +24,7 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
"sync"
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/minio/madmin-go/v3"
|
"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) {
|
func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string, configFile string, configData []byte, parse, lifecycleDelete bool) (updatedAt time.Time, err error) {
|
||||||
objAPI := newObjectLayerFn()
|
objAPI := newObjectLayerFn()
|
||||||
if objAPI == nil {
|
if objAPI == nil {
|
||||||
@@ -273,8 +246,19 @@ func lockBucketMetadata(ctx context.Context, objectAPI ObjectLayer, bucket strin
|
|||||||
return lockBucketMetadataWithTimeout(ctx, objectAPI, bucket, globalOperationTimeout)
|
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) {
|
func lockBucketMetadataWithTimeout(ctx context.Context, objectAPI ObjectLayer, bucket string, timeout *dynamicTimeout) (context.Context, func(), error) {
|
||||||
lock := objectAPI.NewNSLock(minioMetaBucket, pathJoin(bucketMetaPrefix, bucket, "metadata.lock"))
|
lock := objectAPI.NewNSLock(minioMetaBucket, pathJoin(bucketMetaPrefix, bucket, "metadata.lock"))
|
||||||
|
if hook := lockBucketMetadataAcquireHook.Load(); hook != nil {
|
||||||
|
(*hook)(bucket)
|
||||||
|
}
|
||||||
lkctx, err := lock.GetLock(ctx, timeout)
|
lkctx, err := lock.GetLock(ctx, timeout)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, nil, err
|
return nil, nil, err
|
||||||
@@ -686,13 +670,7 @@ func (sys *BucketMetadataSys) GetConfig(ctx context.Context, bucket string) (met
|
|||||||
return meta, false, err
|
return meta, false, err
|
||||||
}
|
}
|
||||||
sys.Lock()
|
sys.Lock()
|
||||||
if cur, ok := sys.metadataMap[bucket]; ok && !cur.lastUpdate().Before(meta.lastUpdate()) {
|
sys.metadataMap[bucket] = meta
|
||||||
// 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.clearLoadFailure(bucket)
|
||||||
sys.Unlock()
|
sys.Unlock()
|
||||||
|
|
||||||
@@ -793,6 +771,8 @@ func (sys *BucketMetadataSys) refreshBucketsMetadataLoop(ctx context.Context) {
|
|||||||
wait := sleeper.Timer(ctx)
|
wait := sleeper.Timer(ctx)
|
||||||
|
|
||||||
bucket := buckets[i].Name
|
bucket := buckets[i].Name
|
||||||
|
updated := false
|
||||||
|
|
||||||
meta, err := loadBucketMetadata(ctx, sys.objAPI, bucket)
|
meta, err := loadBucketMetadata(ctx, sys.objAPI, bucket)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
internalLogIf(ctx, err, logger.WarningKind)
|
internalLogIf(ctx, err, logger.WarningKind)
|
||||||
@@ -803,7 +783,19 @@ func (sys *BucketMetadataSys) refreshBucketsMetadataLoop(ctx context.Context) {
|
|||||||
continue
|
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.
|
wait() // wait to proceed to next entry.
|
||||||
}
|
}
|
||||||
|
|||||||
+20
-2
@@ -544,8 +544,26 @@ func (s *peerRESTServer) LoadBucketMetadataHandler(mss *grid.MSS) (np grid.NoPay
|
|||||||
return np, grid.NewRemoteErr(err)
|
return np, grid.NewRemoteErr(err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Publish the metadata and derived registries from the same current revision.
|
// Publish the reloaded metadata unconditionally. Overlapping peer reloads
|
||||||
globalBucketMetadataSys.setReloaded(bucketName, meta)
|
// 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
|
return np, nerr
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user