mirror of
https://github.com/pgsty/minio.git
synced 2026-10-10 03:35:59 +03:00
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 <noreply@anthropic.com> Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
@@ -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(`<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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user