// Copyright (c) 2026 PGSTY // SPDX-License-Identifier: AGPL-3.0-only package cmd import ( "bytes" "fmt" "net/http" "net/http/httptest" "strings" "sync/atomic" "testing" "time" "github.com/minio/madmin-go/v3" "github.com/minio/minio-go/v7" "github.com/minio/minio/internal/auth" "github.com/minio/minio/internal/bucket/replication" xhttp "github.com/minio/minio/internal/http" "github.com/minio/minio/internal/once" ) // Exercise the actual single-object DELETE handler and erasure metadata. A // marker purge must delete the target version and finish the source purge; // merely seeing a 405 on HEAD is not evidence of a completed permanent delete. func TestReviewR6FailureMRF(t *testing.T) { defer DetectTestLeak(t)() for _, legacy := range []bool{false, true} { t.Run(fmt.Sprintf("recover_legacy_%v", legacy), func(t *testing.T) { ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ t: t, endpoints: []string{"DeleteObject"}, objAPITest: func(obj ObjectLayer, instanceType, bucket string, router http.Handler, creds auth.Credentials, t *testing.T) { testReviewR6FailureMRF(obj, instanceType, bucket, router, creds, t, legacy) }, }) }) } } func testReviewR6FailureMRF(obj ObjectLayer, instanceType, bucket string, router http.Handler, creds auth.Credentials, t *testing.T, legacy bool) { ctx := t.Context() oldStats := globalReplicationStats.Swap(NewReplicationStats(ctx, nil)) defer globalReplicationStats.Store(oldStats) if pools, ok := obj.(*erasureServerPools); ok { for _, pool := range pools.serverPools { for _, set := range pool.sets { disks := set.getDisks() wrapped := make([]StorageAPI, len(disks)) for i,d := range disks { if d!=nil {wrapped[i]=tagTestCapacityDisk{StorageAPI:d}} } set.getDisks = func() []StorageAPI {return wrapped} }} } const arn = "arn:minio:replication::af470089-d354-4473-934c-9e1f52f6da89:bucket" const name = "marker" version := mustGetUUID() remoteBucket := getRandomBucketName() if err := obj.MakeBucket(ctx, remoteBucket, MakeBucketOptions{}); err != nil { t.Fatal(err) } if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, enabledBucketVersioningConfig); err != nil { t.Fatal(err) } for _, b := range []string{bucket, remoteBucket} { if _, err := obj.PutObject(ctx, b, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true}); err != nil { t.Fatal(err) } opts := ObjectOptions{VersionID: version, Versioned: true, DeleteMarker: true, ReplicationRequest: true, MTime: UTCNow()} opts.SetReplicaStatus(replication.Replica) if _, err := obj.DeleteObject(ctx, b, name, opts); err != nil { t.Fatalf("%s: seed marker in %s: %v", instanceType, b, err) } } var reject atomic.Bool reject.Store(true) remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { opts := ObjectOptions{VersionID: r.URL.Query().Get("versionId"), Versioned: true} switch r.Method { case http.MethodHead: oi, err := obj.GetObjectInfo(r.Context(), remoteBucket, name, opts) if oi.DeleteMarker { w.Header().Set(xhttp.AmzDeleteMarker, "true") w.Header().Set(xhttp.AmzVersionID, oi.VersionID) } if err != nil { writeErrorResponseHeadersOnly(w, toAPIError(r.Context(), err)) return } w.WriteHeader(http.StatusOK) case http.MethodDelete: if reject.Load() {w.WriteHeader(403); fmt.Fprint(w, `AccessDenied`);return} opts.DeleteMarker = r.Header.Get(xhttp.MinIOSourceDeleteMarker) == "true" opts.SetReplicaStatus(replication.Replica) _, err := obj.DeleteObject(r.Context(), remoteBucket, name, opts) if err != nil && !isErrVersionNotFound(err) && !isErrObjectNotFound(err) { writeErrorResponse(r.Context(), w, toAPIError(r.Context(), err), r.URL) return } w.WriteHeader(http.StatusNoContent) default: t.Errorf("unexpected remote method %s", r.Method) w.WriteHeader(http.StatusBadRequest) } })) defer remote.Close() client, err := minio.New(strings.TrimPrefix(remote.URL, "http://"), &minio.Options{Region: "us-east-1"}) if err != nil { t.Fatal(err) } target := &TargetClient{Client: client, ARN: arn, Bucket: remoteBucket} globalBucketTargetSys.Lock() globalBucketTargetSys.arnRemotesMap[arn] = arnTarget{Client: target, lastRefresh: UTCNow()} globalBucketTargetSys.targetsMap[bucket] = []madmin.BucketTarget{{Arn: arn, TargetBucket: remoteBucket}} globalBucketTargetSys.Unlock() globalBucketTargetSys.hMutex.Lock() globalBucketTargetSys.hc[client.EndpointURL().Host] = epHealth{Online: true} globalBucketTargetSys.hMutex.Unlock() meta, err := globalBucketMetadataSys.Get(bucket) if err != nil { t.Fatal(err) } cfg := configs[0] cfg.RoleArn = arn meta.replicationConfig = &cfg globalBucketMetadataSys.Set(bucket, meta) worker := make(chan ReplicationWorkerOperation, 1) p := &ReplicationPool{ ctx: ctx, objLayer: obj, workers: []chan ReplicationWorkerOperation{worker}, stats: globalReplicationStats.Load(), mrfSaveCh: make(chan MRFReplicateEntry, 1), } oldPool := globalReplicationPool globalReplicationPool = once.NewSingleton[ReplicationPool]() globalReplicationPool.Set(p) defer func() { globalReplicationPool = oldPool }() req, err := newTestSignedRequestV4(http.MethodDelete, "/"+bucket+"/"+name+"?versionId="+version, 0, nil, creds.AccessKey, creds.SecretKey, nil) if err != nil { t.Fatal(err) } w := httptest.NewRecorder() router.ServeHTTP(w, req) if w.Code != http.StatusNoContent { t.Fatalf("DELETE status %d: %s", w.Code, w.Body.String()) } var deletion DeletedObjectReplicationInfo select { case op := <-worker: deletion = op.(DeletedObjectReplicationInfo) case <-time.After(time.Second): t.Fatal("DELETE did not schedule replication") } if deletion.VersionID != version || deletion.DeleteMarkerVersionID != "" { t.Errorf("purge scheduled as marker creation: version=%q marker=%q", deletion.VersionID, deletion.DeleteMarkerVersionID) } deletion.ReplicationState.Targets = map[string]replication.StatusType{arn: replication.Completed} deletion.ReplicationState.ReplicationStatusInternal = arn+"=COMPLETED;" if legacy {deletion.VersionID,deletion.DeleteMarkerVersionID="",version} result:=replicateDelete(ctx,deletion,obj) t.Logf("legacy=%v result creation=%s purge=%s MRF=%d",legacy,result.ReplicationStatus(),result.VersionPurgeStatus(),len(p.mrfSaveCh)) if result.VersionPurgeStatus()!=replication.VersionPurgeFailed || result.ReplicationStatus()!=replication.Completed {t.Error("failure stored in incorrect status field")} var entry MRFReplicateEntry select {case entry= <-p.mrfSaveCh:default:t.Fatal("failed purge did not enter MRF")} oi,lookupErr:=obj.GetObjectInfo(ctx,bucket,name,ObjectOptions{VersionID:version}) t.Logf("real source lookup: name=%s version=%s marker=%v purge=%s err=%v",oi.Name,oi.VersionID,oi.DeleteMarker,oi.VersionPurgeStatus,lookupErr) if !isErrMethodNotAllowed(lookupErr)||!oi.DeleteMarker {t.Fatal("expected real marker plus 405")} p.saveMRFEntries(ctx,map[string]MRFReplicateEntry{entry.versionID:entry}) fresh:= &ReplicationPool{ctx:ctx,objLayer:obj,workers:[]chan ReplicationWorkerOperation{worker},stats:globalReplicationStats.Load(),mrfSaveCh:make(chan MRFReplicateEntry,1)} globalReplicationPool=once.NewSingleton[ReplicationPool]();globalReplicationPool.Set(fresh) reject.Store(false) if err:=fresh.queueMRFHeal();err!=nil{t.Fatal(err)} select {case op:= <-worker:deletion=op.(DeletedObjectReplicationInfo);case <-time.After(time.Second):t.Fatal("persisted MRF marker lookup was skipped")} result=replicateDelete(ctx,deletion,obj) if result.VersionPurgeStatus()!=replication.VersionPurgeComplete {t.Fatalf("MRF retry result: %+v",result)} for _,b:=range []string{bucket,remoteBucket}{ oi,err:=obj.GetObjectInfo(ctx,b,name,ObjectOptions{VersionID:version}) if !isErrVersionNotFound(err)&&!isErrObjectNotFound(err){t.Errorf("marker remains in %s: %+v %v",b,oi,err)} } }