mirror of
https://github.com/pgsty/minio.git
synced 2026-09-15 23:14:04 +03:00
fix(replication): retry marker purges through persisted MRF
(cherry picked from commit cf381a7151ef25fc95ace5fedcd767fa19410de2) Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
@@ -418,6 +418,9 @@ func getReplicationState(rinfos replicatedInfos, prevState ReplicationState, vID
|
||||
|
||||
for _, rinfo := range rinfos.Targets {
|
||||
if rinfo.ResyncTimestamp != "" {
|
||||
if rs.ResetStatusesMap == nil {
|
||||
rs.ResetStatusesMap = make(map[string]string)
|
||||
}
|
||||
rs.ResetStatusesMap[targetResetHeader(rinfo.Arn)] = rinfo.ResyncTimestamp
|
||||
}
|
||||
}
|
||||
|
||||
+84
-43
@@ -425,6 +425,7 @@ func checkReplicateDelete(ctx context.Context, bucket string, dobj ObjectToDelet
|
||||
// the mere presence or absence of the target version.
|
||||
func replicateDelete(ctx context.Context, dobj DeletedObjectReplicationInfo, objectAPI ObjectLayer) replicatedInfos {
|
||||
var replicationStatus replication.StatusType
|
||||
isPurge := dobj.isVersionPurge()
|
||||
bucket := dobj.Bucket
|
||||
versionID := dobj.DeleteMarkerVersionID
|
||||
if versionID == "" {
|
||||
@@ -484,6 +485,7 @@ func replicateDelete(ctx context.Context, dobj DeletedObjectReplicationInfo, obj
|
||||
lk := objectAPI.NewNSLock(bucket, "/[replicate]/"+dobj.ObjectName)
|
||||
lkctx, err := lk.GetLock(ctx, globalOperationTimeout)
|
||||
if err != nil {
|
||||
dobj.RetryCount++
|
||||
globalReplicationPool.Get().queueMRFSave(dobj.ToMRFEntry())
|
||||
sendEvent(eventArgs{
|
||||
BucketName: bucket,
|
||||
@@ -546,25 +548,38 @@ func replicateDelete(ctx context.Context, dobj DeletedObjectReplicationInfo, obj
|
||||
replicationStatus = rinfos.ReplicationStatus()
|
||||
prevStatus := dobj.DeleteMarkerReplicationStatus()
|
||||
|
||||
if dobj.VersionID != "" {
|
||||
prevStatus = replication.StatusType(dobj.VersionPurgeStatus())
|
||||
replicationStatus = replication.StatusType(rinfos.VersionPurgeStatus())
|
||||
if isPurge {
|
||||
prevStatus = purgeReplicationStatus(dobj.VersionPurgeStatus())
|
||||
replicationStatus = purgeReplicationStatus(rinfos.VersionPurgeStatus())
|
||||
}
|
||||
|
||||
// to decrement pending count later.
|
||||
for _, rinfo := range rinfos.Targets {
|
||||
if rinfo.ReplicationStatus != rinfo.PrevReplicationStatus {
|
||||
globalReplicationStats.Load().Update(dobj.Bucket, rinfo, replicationStatus,
|
||||
prevStatus)
|
||||
status, previous := rinfo.ReplicationStatus, rinfo.PrevReplicationStatus
|
||||
if isPurge {
|
||||
status = purgeReplicationStatus(rinfo.VersionPurgeStatus)
|
||||
previous = purgeReplicationStatus(dobj.ReplicationState.PurgeTargets[rinfo.Arn])
|
||||
}
|
||||
if status != previous {
|
||||
globalReplicationStats.Load().Update(dobj.Bucket, rinfo, status, previous)
|
||||
}
|
||||
}
|
||||
|
||||
eventName := event.ObjectReplicationComplete
|
||||
if replicationStatus == replication.Failed {
|
||||
eventName = event.ObjectReplicationFailed
|
||||
dobj.RetryCount++
|
||||
globalReplicationPool.Get().queueMRFSave(dobj.ToMRFEntry())
|
||||
}
|
||||
drs := getReplicationState(rinfos, dobj.ReplicationState, dobj.VersionID)
|
||||
if isPurge {
|
||||
// A purge must not rewrite the marker's creation/replica metadata.
|
||||
// Multiple serialized empty target statuses can parse as nonempty,
|
||||
// so explicitly send an empty creation update to the metadata writer.
|
||||
drs.ReplicationStatusInternal = ""
|
||||
drs.Targets = nil
|
||||
drs.ReplicaStatus = ""
|
||||
}
|
||||
if replicationStatus != prevStatus {
|
||||
drs.ReplicationTimeStamp = UTCNow()
|
||||
}
|
||||
@@ -606,6 +621,7 @@ func replicateDelete(ctx context.Context, dobj DeletedObjectReplicationInfo, obj
|
||||
}
|
||||
|
||||
func replicateDeleteToTarget(ctx context.Context, dobj DeletedObjectReplicationInfo, tgt *TargetClient) (rinfo replicatedTargetInfo) {
|
||||
isPurge := dobj.isVersionPurge()
|
||||
versionID := dobj.DeleteMarkerVersionID
|
||||
if versionID == "" {
|
||||
versionID = dobj.VersionID
|
||||
@@ -615,42 +631,50 @@ func replicateDeleteToTarget(ctx context.Context, dobj DeletedObjectReplicationI
|
||||
rinfo.OpType = dobj.OpType
|
||||
rinfo.endpoint = tgt.EndpointURL().Host
|
||||
rinfo.secure = tgt.EndpointURL().Scheme == "https"
|
||||
// A purge leaves ReplicationStatus empty: the metadata writer interprets
|
||||
// that as preserving the entire creation-status block, including targets
|
||||
// outside this fan-out. Only VersionPurgeStatus records the purge outcome.
|
||||
defer func() {
|
||||
if rinfo.ReplicationStatus == replication.Completed && tgt.ResetID != "" && dobj.OpType == replication.ExistingObjectReplicationType {
|
||||
completed := rinfo.ReplicationStatus == replication.Completed
|
||||
if isPurge {
|
||||
completed = rinfo.VersionPurgeStatus == replication.VersionPurgeComplete
|
||||
}
|
||||
if completed && tgt.ResetID != "" && dobj.OpType == replication.ExistingObjectReplicationType {
|
||||
rinfo.ResyncTimestamp = fmt.Sprintf("%s;%s", UTCNow().Format(http.TimeFormat), tgt.ResetID)
|
||||
}
|
||||
}()
|
||||
|
||||
if dobj.VersionID == "" && rinfo.PrevReplicationStatus == replication.Completed && dobj.OpType != replication.ExistingObjectReplicationType {
|
||||
if !isPurge && rinfo.PrevReplicationStatus == replication.Completed && dobj.OpType != replication.ExistingObjectReplicationType {
|
||||
rinfo.ReplicationStatus = rinfo.PrevReplicationStatus
|
||||
return rinfo
|
||||
}
|
||||
if dobj.VersionID != "" && rinfo.VersionPurgeStatus == replication.VersionPurgeComplete {
|
||||
if isPurge && rinfo.VersionPurgeStatus == replication.VersionPurgeComplete {
|
||||
return rinfo
|
||||
}
|
||||
if globalBucketTargetSys.isOffline(tgt.EndpointURL()) {
|
||||
replLogOnceIf(ctx, fmt.Errorf("remote target is offline for bucket:%s arn:%s", dobj.Bucket, tgt.ARN), "replication-target-offline-delete-"+tgt.ARN)
|
||||
rinfo.Err = fmt.Errorf("remote target is offline for bucket:%s arn:%s", dobj.Bucket, tgt.ARN)
|
||||
replLogOnceIf(ctx, rinfo.Err, "replication-target-offline-delete-"+tgt.ARN)
|
||||
sendEvent(eventArgs{
|
||||
BucketName: dobj.Bucket,
|
||||
Object: ObjectInfo{
|
||||
Bucket: dobj.Bucket,
|
||||
Name: dobj.ObjectName,
|
||||
VersionID: dobj.VersionID,
|
||||
VersionID: versionID,
|
||||
DeleteMarker: dobj.DeleteMarker,
|
||||
},
|
||||
UserAgent: "Internal: [Replication]",
|
||||
Host: globalLocalNodeName,
|
||||
EventName: event.ObjectReplicationNotTracked,
|
||||
})
|
||||
if dobj.VersionID == "" {
|
||||
rinfo.ReplicationStatus = replication.Failed
|
||||
} else {
|
||||
if isPurge {
|
||||
rinfo.VersionPurgeStatus = replication.VersionPurgeFailed
|
||||
} else {
|
||||
rinfo.ReplicationStatus = replication.Failed
|
||||
}
|
||||
return rinfo
|
||||
}
|
||||
// early return if already replicated delete marker for existing object replication/ healing delete markers
|
||||
if dobj.DeleteMarkerVersionID != "" {
|
||||
if !isPurge && dobj.DeleteMarkerVersionID != "" {
|
||||
toi, err := tgt.StatObject(ctx, tgt.Bucket, dobj.ObjectName, minio.StatObjectOptions{
|
||||
VersionID: versionID,
|
||||
Internal: minio.AdvancedGetOptions{
|
||||
@@ -662,16 +686,10 @@ func replicateDeleteToTarget(ctx context.Context, dobj DeletedObjectReplicationI
|
||||
switch {
|
||||
case isErrMethodNotAllowed(serr):
|
||||
// delete marker already replicated
|
||||
if dobj.VersionID == "" && rinfo.VersionPurgeStatus.Empty() {
|
||||
rinfo.ReplicationStatus = replication.Completed
|
||||
return rinfo
|
||||
}
|
||||
rinfo.ReplicationStatus = replication.Completed
|
||||
return rinfo
|
||||
case isErrObjectNotFound(serr), isErrVersionNotFound(serr):
|
||||
// version being purged is already not found on target.
|
||||
if !rinfo.VersionPurgeStatus.Empty() {
|
||||
rinfo.VersionPurgeStatus = replication.VersionPurgeComplete
|
||||
return rinfo
|
||||
}
|
||||
// The marker still needs to be created on the target.
|
||||
case isErrReadQuorum(serr), isErrWriteQuorum(serr):
|
||||
// destination has some quorum issues, perform removeObject() anyways
|
||||
// to complete the operation.
|
||||
@@ -691,7 +709,7 @@ func replicateDeleteToTarget(ctx context.Context, dobj DeletedObjectReplicationI
|
||||
rmErr := tgt.RemoveObject(ctx, tgt.Bucket, dobj.ObjectName, minio.RemoveObjectOptions{
|
||||
VersionID: versionID,
|
||||
Internal: minio.AdvancedRemoveOptions{
|
||||
ReplicationDeleteMarker: dobj.DeleteMarkerVersionID != "",
|
||||
ReplicationDeleteMarker: !isPurge && dobj.DeleteMarkerVersionID != "",
|
||||
ReplicationMTime: dobj.DeleteMarkerMTime.Time,
|
||||
ReplicationStatus: minio.ReplicationStatusReplica,
|
||||
ReplicationRequest: true, // always set this to distinguish between `mc mirror` replication and serverside
|
||||
@@ -699,20 +717,20 @@ func replicateDeleteToTarget(ctx context.Context, dobj DeletedObjectReplicationI
|
||||
})
|
||||
if rmErr != nil {
|
||||
rinfo.Err = rmErr
|
||||
if dobj.VersionID == "" {
|
||||
rinfo.ReplicationStatus = replication.Failed
|
||||
} else {
|
||||
if isPurge {
|
||||
rinfo.VersionPurgeStatus = replication.VersionPurgeFailed
|
||||
} else {
|
||||
rinfo.ReplicationStatus = replication.Failed
|
||||
}
|
||||
replLogIf(ctx, fmt.Errorf("unable to replicate delete marker to %s: %s/%s(%s): %w", tgt.EndpointURL(), tgt.Bucket, dobj.ObjectName, versionID, rmErr))
|
||||
if rmErr != nil && minio.IsNetworkOrHostDown(rmErr, true) && !globalBucketTargetSys.isOffline(tgt.EndpointURL()) {
|
||||
globalBucketTargetSys.markOffline(tgt.EndpointURL())
|
||||
}
|
||||
} else {
|
||||
if dobj.VersionID == "" {
|
||||
rinfo.ReplicationStatus = replication.Completed
|
||||
} else {
|
||||
if isPurge {
|
||||
rinfo.VersionPurgeStatus = replication.VersionPurgeComplete
|
||||
} else {
|
||||
rinfo.ReplicationStatus = replication.Completed
|
||||
}
|
||||
}
|
||||
return rinfo
|
||||
@@ -1923,11 +1941,26 @@ func filterReplicationStatusMetadata(metadata map[string]string) map[string]stri
|
||||
// DeletedObjectReplicationInfo has info on deleted object
|
||||
type DeletedObjectReplicationInfo struct {
|
||||
DeletedObject
|
||||
Bucket string
|
||||
EventType string
|
||||
OpType replication.Type
|
||||
ResetID string
|
||||
TargetArn string
|
||||
Bucket string
|
||||
EventType string
|
||||
OpType replication.Type
|
||||
ResetID string
|
||||
TargetArn string
|
||||
RetryCount int
|
||||
}
|
||||
|
||||
// isVersionPurge also recognizes the old marker-shaped purge task. Use the
|
||||
// operation's state, rather than one target's possibly missing purge entry.
|
||||
func (di DeletedObjectReplicationInfo) isVersionPurge() bool {
|
||||
return di.VersionID != "" || di.DeleteMarkerVersionID != "" && !di.VersionPurgeStatus().Empty()
|
||||
}
|
||||
|
||||
// Purge metadata uses COMPLETE; operation statistics and audit use COMPLETED.
|
||||
func purgeReplicationStatus(status VersionPurgeStatusType) replication.StatusType {
|
||||
if replication.StatusType(status) == replication.CompletedLegacy {
|
||||
return replication.Completed
|
||||
}
|
||||
return replication.StatusType(status)
|
||||
}
|
||||
|
||||
// ToMRFEntry returns the relevant info needed by MRF
|
||||
@@ -1937,9 +1970,10 @@ func (di DeletedObjectReplicationInfo) ToMRFEntry() MRFReplicateEntry {
|
||||
versionID = di.VersionID
|
||||
}
|
||||
return MRFReplicateEntry{
|
||||
Bucket: di.Bucket,
|
||||
Object: di.ObjectName,
|
||||
versionID: versionID,
|
||||
Bucket: di.Bucket,
|
||||
Object: di.ObjectName,
|
||||
versionID: versionID,
|
||||
RetryCount: di.RetryCount,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2431,6 +2465,7 @@ func (p *ReplicationPool) queueReplicaDeleteTask(doi DeletedObjectReplicationInf
|
||||
case <-p.ctx.Done():
|
||||
case ch <- doi:
|
||||
default:
|
||||
doi.RetryCount++
|
||||
p.queueMRFSave(doi.ToMRFEntry())
|
||||
p.mu.RLock()
|
||||
prio := p.priority
|
||||
@@ -3794,9 +3829,10 @@ func queueReplicationHeal(ctx context.Context, bucket string, oi ObjectInfo, rcf
|
||||
DeleteMarkerMTime: DeleteMarkerMTime{roi.ModTime},
|
||||
DeleteMarker: roi.DeleteMarker,
|
||||
},
|
||||
Bucket: roi.Bucket,
|
||||
OpType: replication.HealReplicationType,
|
||||
EventType: ReplicateHealDelete,
|
||||
Bucket: roi.Bucket,
|
||||
OpType: replication.HealReplicationType,
|
||||
EventType: ReplicateHealDelete,
|
||||
RetryCount: retryCount,
|
||||
}
|
||||
// heal delete marker replication failure or versioned delete replication failure
|
||||
if roi.ReplicationStatus == replication.Pending ||
|
||||
@@ -4079,7 +4115,12 @@ func (p *ReplicationPool) queueMRFHeal() error {
|
||||
VersionID: vID,
|
||||
})
|
||||
cancel()
|
||||
if err != nil {
|
||||
// A versioned marker lookup returns its metadata with a 405. Only
|
||||
// accept that error with a real, matching marker identity.
|
||||
validMarker := isErrMethodNotAllowed(err) && oi.DeleteMarker &&
|
||||
vID != "" && oi.VersionID == vID && !oi.ModTime.IsZero() &&
|
||||
oi.Bucket == e.Bucket && oi.Name != "" && oi.Name == decodeDirObject(e.Object)
|
||||
if err != nil && !validMarker || oi.Name == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
|
||||
@@ -106,6 +106,7 @@ func TestReplicateDeleteMarkerTargetSemantics(t *testing.T) {
|
||||
|
||||
func testReplicateDeleteMarkerPurge(obj ObjectLayer, instanceType, bucket string, router http.Handler, creds auth.Credentials, t *testing.T, legacy bool) {
|
||||
ctx := t.Context()
|
||||
defer replicationTestCapacity(obj)()
|
||||
const arn = "arn:minio:replication::af470089-d354-4473-934c-9e1f52f6da89:bucket"
|
||||
const name = "marker"
|
||||
version := mustGetUUID()
|
||||
@@ -208,27 +209,12 @@ func testReplicateDeleteMarkerPurge(obj ObjectLayer, instanceType, bucket string
|
||||
t.Errorf("purge scheduled as marker creation: version=%q marker=%q", deletion.VersionID, deletion.DeleteMarkerVersionID)
|
||||
}
|
||||
if legacy {
|
||||
// Reproduce the old producer's state and let the existing scanner/heal
|
||||
// path recover it. Upgrades must also finish purges already left pending.
|
||||
// Old task shapes must complete directly; scanner/MRF recovery is
|
||||
// covered separately by TestReplicationMRFMarkerRecovery.
|
||||
deletion.VersionID, deletion.DeleteMarkerVersionID = "", version
|
||||
replicateDelete(ctx, deletion, obj)
|
||||
oi, _ := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true})
|
||||
if oi.VersionPurgeStatus != replication.VersionPurgePending {
|
||||
t.Fatalf("legacy source purge = %s, want PENDING", oi.VersionPurgeStatus)
|
||||
}
|
||||
targets, err := globalBucketTargetSys.ListBucketTargets(ctx, bucket)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
queueReplicationHeal(ctx, bucket, oi, replicationConfig{Config: &cfg, remotes: targets}, 0)
|
||||
select {
|
||||
case op := <-worker:
|
||||
deletion = op.(DeletedObjectReplicationInfo)
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("legacy pending purge was not scheduled for healing")
|
||||
}
|
||||
}
|
||||
result := replicateDelete(context.Background(), deletion, obj)
|
||||
|
||||
result := replicateDelete(context.Background(), deletion, markerPurgeUpdateLayer{ObjectLayer: obj, t: t})
|
||||
if result.VersionPurgeStatus() != replication.VersionPurgeComplete {
|
||||
t.Errorf("remote purge result = %s, want COMPLETE", result.VersionPurgeStatus())
|
||||
}
|
||||
|
||||
@@ -0,0 +1,614 @@
|
||||
// Copyright (c) 2026 PGSTY
|
||||
// SPDX-License-Identifier: AGPL-3.0-only
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"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/logger"
|
||||
loghttp "github.com/minio/minio/internal/logger/target/http"
|
||||
"github.com/minio/minio/internal/once"
|
||||
xnet "github.com/pgsty/silo-pkg/v3/net"
|
||||
)
|
||||
|
||||
type markerRecoveryCase struct {
|
||||
name string
|
||||
legacy, creation, lostReply, partial, exhaust, lockFirst, invalidMRF, replicaSource, offline, unrecordedPurge bool
|
||||
targets int
|
||||
}
|
||||
|
||||
func TestReplicationMRFMarkerRecovery(t *testing.T) {
|
||||
for _, tc := range []markerRecoveryCase{
|
||||
{name: "canonical", targets: 1},
|
||||
{name: "lock-failure", lockFirst: true, targets: 1},
|
||||
{name: "invalid-MRF-metadata", invalidMRF: true, targets: 1},
|
||||
{name: "legacy", legacy: true, targets: 1},
|
||||
{name: "unrecorded-purge", unrecordedPurge: true, targets: 2},
|
||||
{name: "two-targets", legacy: true, targets: 2},
|
||||
{name: "two-targets-offline", offline: true, targets: 2},
|
||||
{name: "replica-source", replicaSource: true, targets: 1},
|
||||
{name: "lost-reply-canonical", lostReply: true, targets: 1},
|
||||
{name: "lost-reply-legacy", legacy: true, lostReply: true, targets: 1},
|
||||
{name: "creation", creation: true, targets: 1},
|
||||
{name: "partial-creation-block", partial: true, targets: 2},
|
||||
{name: "retry-budget-and-scanner", exhaust: true, targets: 1},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, endpoints: []string{"DeleteObject"}, objAPITest: func(obj ObjectLayer, backend, bucket string, router http.Handler, creds auth.Credentials, t *testing.T) {
|
||||
testReplicationMRFMarkerRecovery(t, obj, backend, bucket, router, creds, tc)
|
||||
}})
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
type markerRecoveryTarget struct {
|
||||
arn, bucket string
|
||||
client *minio.Client
|
||||
reject, loseReply atomic.Bool
|
||||
deletes atomic.Int32
|
||||
}
|
||||
|
||||
// Only adapt host capacity accounting, using the existing real-disk adapter.
|
||||
// Reads, object metadata, MRF files and writes all still use the fixture disks.
|
||||
func replicationTestCapacity(obj ObjectLayer) func() {
|
||||
var restore []func()
|
||||
for _, pool := range obj.(*erasureServerPools).serverPools {
|
||||
for _, set := range pool.sets {
|
||||
original := set.getDisks
|
||||
disks := append([]StorageAPI(nil), original()...)
|
||||
for i, disk := range disks {
|
||||
if disk != nil {
|
||||
disks[i] = tagTestCapacityDisk{StorageAPI: disk}
|
||||
}
|
||||
}
|
||||
set.getDisks = func() []StorageAPI { return disks }
|
||||
restore = append(restore, func() { set.getDisks = original })
|
||||
}
|
||||
}
|
||||
return func() {
|
||||
for _, fn := range restore {
|
||||
fn()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func testReplicationMRFMarkerRecovery(t *testing.T, obj ObjectLayer, backend, bucket string, router http.Handler, creds auth.Credentials, tc markerRecoveryCase) {
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
defer cancel()
|
||||
defer replicationTestCapacity(obj)()
|
||||
stats := NewReplicationStats(ctx, nil)
|
||||
oldStats := globalReplicationStats.Swap(stats)
|
||||
defer globalReplicationStats.Store(oldStats)
|
||||
oldPool := globalReplicationPool
|
||||
defer func() { globalReplicationPool = oldPool }()
|
||||
const name = "marker"
|
||||
version := mustGetUUID()
|
||||
creationTime := UTCNow().Add(-time.Hour).Truncate(time.Second)
|
||||
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, enabledBucketVersioningConfig); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := obj.PutObject(ctx, bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var targets []*markerRecoveryTarget
|
||||
cfg := replication.Config{}
|
||||
creationStates := make(map[string]replication.StatusType)
|
||||
var creationInternal string
|
||||
for i := 0; i < tc.targets; i++ {
|
||||
target := &markerRecoveryTarget{arn: "arn:minio:replication::" + mustGetUUID() + ":bucket", bucket: getRandomBucketName()}
|
||||
if err := obj.MakeBucket(ctx, target.bucket, MakeBucketOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := obj.PutObject(ctx, target.bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !tc.creation {
|
||||
opts := ObjectOptions{VersionID: version, Versioned: true, DeleteMarker: true, ReplicationRequest: true, MTime: creationTime}
|
||||
opts.SetReplicaStatus(replication.Replica)
|
||||
if _, err := obj.DeleteObject(ctx, target.bucket, name, opts); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
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(), target.bucket, 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:
|
||||
target.deletes.Add(1)
|
||||
if got := r.Header.Get(xhttp.MinIOSourceDeleteMarker) == "true"; got != tc.creation {
|
||||
t.Errorf("purge/creation wire flag=%v creation=%v", got, tc.creation)
|
||||
}
|
||||
if opts.VersionID == "" {
|
||||
t.Error("missing remote versionId")
|
||||
}
|
||||
if target.reject.Load() {
|
||||
w.WriteHeader(http.StatusForbidden)
|
||||
fmt.Fprint(w, `<Error><Code>AccessDenied</Code></Error>`)
|
||||
return
|
||||
}
|
||||
opts.DeleteMarker = r.Header.Get(xhttp.MinIOSourceDeleteMarker) == "true"
|
||||
opts.ReplicationRequest = true
|
||||
opts.SetReplicaStatus(replication.Replica)
|
||||
_, err := obj.DeleteObject(r.Context(), target.bucket, name, opts)
|
||||
if err != nil && !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
writeErrorResponse(r.Context(), w, toAPIError(r.Context(), err), r.URL)
|
||||
return
|
||||
}
|
||||
if target.loseReply.Swap(false) {
|
||||
conn, _, err := w.(http.Hijacker).Hijack()
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
return
|
||||
}
|
||||
conn.Close()
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
default:
|
||||
t.Errorf("unexpected remote method %s", r.Method)
|
||||
}
|
||||
}))
|
||||
defer remote.Close()
|
||||
client, err := minio.New(strings.TrimPrefix(remote.URL, "http://"), &minio.Options{Region: "us-east-1", MaxRetries: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
target.client = client
|
||||
globalBucketTargetSys.Lock()
|
||||
globalBucketTargetSys.arnRemotesMap[target.arn] = arnTarget{Client: &TargetClient{Client: client, ARN: target.arn, Bucket: target.bucket}, lastRefresh: UTCNow()}
|
||||
globalBucketTargetSys.targetsMap[bucket] = append(globalBucketTargetSys.targetsMap[bucket], madmin.BucketTarget{Arn: target.arn, TargetBucket: target.bucket})
|
||||
globalBucketTargetSys.Unlock()
|
||||
globalBucketTargetSys.hMutex.Lock()
|
||||
globalBucketTargetSys.hc[client.EndpointURL().Host] = epHealth{Online: true}
|
||||
globalBucketTargetSys.hMutex.Unlock()
|
||||
rule := configs[0].Rules[0]
|
||||
rule.Priority = i + 1
|
||||
rule.Destination = replication.Destination{ARN: target.arn, Bucket: target.bucket}
|
||||
cfg.Rules = append(cfg.Rules, rule)
|
||||
creationStates[target.arn] = replication.Completed
|
||||
creationInternal += target.arn + "=COMPLETED;"
|
||||
targets = append(targets, target)
|
||||
}
|
||||
if tc.targets == 1 {
|
||||
cfg.RoleArn = targets[0].arn
|
||||
}
|
||||
if !tc.creation {
|
||||
opts := ObjectOptions{VersionID: version, Versioned: true, DeleteMarker: true, MTime: creationTime, DeleteReplication: ReplicationState{Targets: creationStates, ReplicationStatusInternal: creationInternal, ReplicationTimeStamp: creationTime}}
|
||||
if tc.replicaSource {
|
||||
opts.DeleteReplication = ReplicationState{}
|
||||
opts.SetReplicaStatus(replication.Replica)
|
||||
}
|
||||
if _, err := obj.DeleteObject(ctx, bucket, name, opts); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
meta, err := globalBucketMetadataSys.Get(bucket)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
meta.replicationConfig = &cfg
|
||||
globalBucketMetadataSys.Set(bucket, meta)
|
||||
newPool := func() *ReplicationPool {
|
||||
p := &ReplicationPool{ctx: ctx, objLayer: obj, workers: []chan ReplicationWorkerOperation{make(chan ReplicationWorkerOperation, 8)}, stats: stats, mrfSaveCh: make(chan MRFReplicateEntry, 8)}
|
||||
globalReplicationPool = once.NewSingleton[ReplicationPool]()
|
||||
globalReplicationPool.Set(p)
|
||||
return p
|
||||
}
|
||||
p := newPool()
|
||||
receive := func() DeletedObjectReplicationInfo {
|
||||
select {
|
||||
case op := <-p.workers[0]:
|
||||
d, ok := op.(DeletedObjectReplicationInfo)
|
||||
if !ok {
|
||||
t.Fatalf("wrong queued operation %T", op)
|
||||
}
|
||||
return d
|
||||
case <-time.After(3 * time.Second):
|
||||
t.Fatal("no marker task from actual MRF/handler/scanner queue")
|
||||
return DeletedObjectReplicationInfo{}
|
||||
}
|
||||
}
|
||||
uri := "/" + bucket + "/" + name
|
||||
if !tc.creation {
|
||||
uri += "?versionId=" + version
|
||||
}
|
||||
req, err := newTestSignedRequestV4(http.MethodDelete, uri, 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("source DELETE status %d: %s", w.Code, w.Body)
|
||||
}
|
||||
deletion := receive()
|
||||
if tc.creation {
|
||||
version = deletion.DeleteMarkerVersionID
|
||||
} else if deletion.VersionID != version || deletion.DeleteMarkerVersionID != "" {
|
||||
t.Fatalf("current producer emitted noncanonical purge: %+v", deletion)
|
||||
}
|
||||
if tc.unrecordedPurge {
|
||||
// Exercise a task carrying purge state while the disk marker still has
|
||||
// only creation metadata. Recreate only this isolated fixture marker.
|
||||
if _, err := obj.DeleteObject(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
opts := ObjectOptions{VersionID: version, Versioned: true, DeleteMarker: true, MTime: creationTime, DeleteReplication: ReplicationState{Targets: creationStates, ReplicationStatusInternal: creationInternal, ReplicationTimeStamp: creationTime}}
|
||||
if _, err := obj.DeleteObject(ctx, bucket, name, opts); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if tc.legacy {
|
||||
deletion.VersionID, deletion.DeleteMarkerVersionID = "", version
|
||||
}
|
||||
if tc.partial {
|
||||
deletion.TargetArn = targets[len(targets)-1].arn
|
||||
}
|
||||
if tc.invalidMRF {
|
||||
testReplicationMRFInvalidLookups(ctx, t, obj, p, newPool, deletion)
|
||||
return
|
||||
}
|
||||
var auditStatuses <-chan string
|
||||
if tc.name == "canonical" {
|
||||
var stopAudit func()
|
||||
auditStatuses, stopAudit = replicationTestAudit(ctx, t, bucket)
|
||||
defer stopAudit()
|
||||
}
|
||||
assertAudit := func(want string) {
|
||||
if auditStatuses == nil {
|
||||
return
|
||||
}
|
||||
select {
|
||||
case got := <-auditStatuses:
|
||||
if got != want {
|
||||
t.Fatalf("replication audit status=%q want %q", got, want)
|
||||
}
|
||||
case <-time.After(3 * time.Second):
|
||||
t.Fatal("missing replication audit event")
|
||||
}
|
||||
}
|
||||
deletion.OpType = replication.HealReplicationType // Exercise operation-status statistics too.
|
||||
failing := targets[len(targets)-1]
|
||||
failing.reject.Store(!tc.lostReply)
|
||||
failing.loseReply.Store(tc.lostReply)
|
||||
if tc.offline {
|
||||
globalBucketTargetSys.hMutex.Lock()
|
||||
globalBucketTargetSys.hc[failing.client.EndpointURL().Host] = epHealth{Online: false}
|
||||
globalBucketTargetSys.hMutex.Unlock()
|
||||
}
|
||||
before, _ := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version})
|
||||
beforeCreation := before.ReplicationStatusInternal
|
||||
beforeStamp := before.UserDefined[ReservedMetadataPrefixLower+ReplicationTimestamp]
|
||||
beforeReplica := before.UserDefined[ReservedMetadataPrefixLower+ReplicaStatus]
|
||||
beforeReplicaStamp := before.UserDefined[ReservedMetadataPrefixLower+ReplicaTimestamp]
|
||||
if !tc.creation && !tc.replicaSource && (len(replicationStatusesMap(beforeCreation)) != tc.targets || beforeStamp == "") {
|
||||
t.Fatalf("missing seeded creation block: %+v", before)
|
||||
}
|
||||
if tc.replicaSource && (beforeReplica != "REPLICA" || beforeReplicaStamp == "") {
|
||||
t.Fatalf("missing replica block: %+v", before)
|
||||
}
|
||||
updateObj := obj
|
||||
if !tc.creation {
|
||||
updateObj = markerPurgeUpdateLayer{ObjectLayer: obj, t: t}
|
||||
}
|
||||
rounds := 2
|
||||
if tc.lostReply {
|
||||
rounds = 1
|
||||
}
|
||||
if tc.exhaust {
|
||||
rounds = mrfRetryLimit + 1
|
||||
}
|
||||
for round := 1; round <= rounds; round++ {
|
||||
callObj := updateObj
|
||||
if tc.lockFirst && round == 1 {
|
||||
callObj = markerLockFailureLayer{ObjectLayer: updateObj}
|
||||
}
|
||||
result := replicateDelete(ctx, deletion, callObj)
|
||||
switch {
|
||||
case tc.lockFirst && round == 1:
|
||||
if len(result.Targets) != 0 {
|
||||
t.Fatal("lock failure attempted a target")
|
||||
}
|
||||
case tc.creation:
|
||||
if result.ReplicationStatus() != replication.Failed {
|
||||
t.Fatalf("creation result: %+v", result)
|
||||
}
|
||||
case result.VersionPurgeStatus() != replication.VersionPurgeFailed:
|
||||
t.Fatalf("purge failure result: %+v", result)
|
||||
}
|
||||
assertAudit("FAILED")
|
||||
if round == 1 && !tc.lockFirst && !tc.creation {
|
||||
stats.RLock()
|
||||
failed := stats.Cache[bucket].Stats[failing.arn].FailStats.SinceUptime
|
||||
stats.RUnlock()
|
||||
if failed.Count != 1 || failed.Bytes != 0 {
|
||||
t.Fatalf("purge failure stats=%+v, want count 1 and zero bytes", failed)
|
||||
}
|
||||
}
|
||||
oi, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version})
|
||||
if !isErrMethodNotAllowed(err) || !oi.DeleteMarker {
|
||||
t.Fatalf("source marker metadata/405 missing: %+v %v", oi, err)
|
||||
}
|
||||
if !tc.creation && (oi.ReplicationStatusInternal != beforeCreation || oi.UserDefined[ReservedMetadataPrefixLower+ReplicationTimestamp] != beforeStamp) {
|
||||
t.Fatalf("purge rewrote creation block: before=%q/%q after=%q/%q", beforeCreation, beforeStamp, oi.ReplicationStatusInternal, oi.UserDefined[ReservedMetadataPrefixLower+ReplicationTimestamp])
|
||||
}
|
||||
if !tc.creation && (oi.UserDefined[ReservedMetadataPrefixLower+ReplicaStatus] != beforeReplica || oi.UserDefined[ReservedMetadataPrefixLower+ReplicaTimestamp] != beforeReplicaStamp) {
|
||||
t.Fatal("purge rewrote replica block")
|
||||
}
|
||||
if tc.targets == 2 && !tc.partial {
|
||||
purges := versionPurgeStatusesMap(oi.VersionPurgeStatusInternal)
|
||||
if purges[targets[0].arn] != replication.VersionPurgeComplete || purges[failing.arn] != replication.VersionPurgeFailed {
|
||||
t.Fatalf("incorrect persisted target states: %v", purges)
|
||||
}
|
||||
if targets[0].deletes.Load() != 1 {
|
||||
t.Fatalf("successful target resent %d times", targets[0].deletes.Load())
|
||||
}
|
||||
}
|
||||
if tc.partial {
|
||||
select {
|
||||
case <-p.mrfSaveCh:
|
||||
default:
|
||||
t.Fatal("partial failure missing MRF")
|
||||
}
|
||||
dsc := deletion.ReplicationState.ReplicateDecisionStr
|
||||
deletion.ReplicationState = oi.ReplicationState()
|
||||
deletion.ReplicationState.ReplicateDecisionStr = dsc
|
||||
continue
|
||||
}
|
||||
if tc.exhaust && round > mrfRetryLimit {
|
||||
if len(p.mrfSaveCh) != 0 || atomic.LoadUint64(&stats.mrfStats.TotalDroppedCount) != 1 {
|
||||
t.Fatalf("retry budget not applied: queued=%d drops=%d", len(p.mrfSaveCh), stats.mrfStats.TotalDroppedCount)
|
||||
}
|
||||
break
|
||||
}
|
||||
var entry MRFReplicateEntry
|
||||
select {
|
||||
case entry = <-p.mrfSaveCh:
|
||||
default:
|
||||
t.Fatal("failure did not enter MRF")
|
||||
}
|
||||
if entry.RetryCount != round || entry.versionID != version {
|
||||
t.Fatalf("MRF entry lost identity/budget: %+v, round=%d", entry, round)
|
||||
}
|
||||
p.saveMRFEntries(ctx, map[string]MRFReplicateEntry{entry.versionID: entry})
|
||||
record, err := p.loadMRF()
|
||||
if err != nil || len(record.Entries) != 1 {
|
||||
t.Fatalf("disk MRF missing: %+v %v", record, err)
|
||||
}
|
||||
if got := record.Entries[version]; got.RetryCount != round || got.Object != name || got.Bucket != bucket {
|
||||
t.Fatalf("disk MRF mismatch: %+v", got)
|
||||
}
|
||||
p.saveMRFEntries(ctx, record.Entries) // loadMRF consumes the file; each replay uses a fresh disk record.
|
||||
p = newPool() // no in-memory entries carried to the replacement pool.
|
||||
if err := p.queueMRFHeal(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
deletion = receive()
|
||||
if deletion.RetryCount != round {
|
||||
t.Fatalf("MRF retry count=%d want %d", deletion.RetryCount, round)
|
||||
}
|
||||
if !tc.creation && (deletion.VersionID != version || deletion.DeleteMarkerVersionID != "") {
|
||||
t.Fatalf("MRF emitted wrong purge: %+v", deletion)
|
||||
}
|
||||
}
|
||||
if tc.partial {
|
||||
return
|
||||
}
|
||||
failing.reject.Store(false)
|
||||
globalBucketTargetSys.hMutex.Lock()
|
||||
globalBucketTargetSys.hc[failing.client.EndpointURL().Host] = epHealth{Online: true}
|
||||
globalBucketTargetSys.hMutex.Unlock()
|
||||
if tc.exhaust {
|
||||
oi, _ := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version})
|
||||
QueueReplicationHeal(ctx, bucket, oi, 0) // Separately prove the existing scanner fallback after MRF exhaustion.
|
||||
deletion = receive()
|
||||
if deletion.RetryCount != 0 {
|
||||
t.Fatal("scanner did not start a fresh retry budget")
|
||||
}
|
||||
}
|
||||
result := replicateDelete(ctx, deletion, updateObj)
|
||||
assertAudit("COMPLETED")
|
||||
if tc.creation {
|
||||
if result.ReplicationStatus() != replication.Completed {
|
||||
t.Fatalf("creation recovery failed: %+v", result)
|
||||
}
|
||||
} else if result.VersionPurgeStatus() != replication.VersionPurgeComplete {
|
||||
t.Fatalf("purge recovery failed: %+v", result)
|
||||
}
|
||||
for _, b := range append([]string{bucket}, func() []string {
|
||||
var b []string
|
||||
for _, target := range targets {
|
||||
b = append(b, target.bucket)
|
||||
}
|
||||
return b
|
||||
}()...) {
|
||||
oi, err := obj.GetObjectInfo(ctx, b, name, ObjectOptions{VersionID: version})
|
||||
if tc.creation {
|
||||
if !oi.DeleteMarker || !isErrMethodNotAllowed(err) {
|
||||
t.Fatalf("creation missing in %s: %+v %v", b, oi, err)
|
||||
}
|
||||
} else if !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
t.Fatalf("purge left marker in %s: %+v %v", b, oi, err)
|
||||
}
|
||||
}
|
||||
if !tc.creation {
|
||||
// Re-deliver an ambiguous purge directly. An absent marker must stay absent.
|
||||
duplicate := deletion
|
||||
duplicate.VersionID, duplicate.DeleteMarkerVersionID = "", version
|
||||
duplicate.ReplicationState.PurgeTargets = map[string]VersionPurgeStatusType{failing.arn: replication.VersionPurgePending}
|
||||
duplicate.ReplicationState.VersionPurgeStatusInternal = ""
|
||||
got := replicateDeleteToTarget(ctx, duplicate, &TargetClient{Client: failing.client, ARN: failing.arn, Bucket: failing.bucket})
|
||||
if got.VersionPurgeStatus != replication.VersionPurgeComplete {
|
||||
t.Fatalf("duplicate purge: %+v", got)
|
||||
}
|
||||
if oi, err := obj.GetObjectInfo(ctx, failing.bucket, name, ObjectOptions{VersionID: version}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
|
||||
t.Fatalf("duplicate recreated marker: %+v %v", oi, err)
|
||||
}
|
||||
}
|
||||
if len(p.mrfSaveCh) != 0 || len(p.workers[0]) != 0 {
|
||||
t.Fatalf("completion left queued work: mrf=%d worker=%d", len(p.mrfSaveCh), len(p.workers[0]))
|
||||
}
|
||||
if !tc.creation && !tc.exhaust && !tc.lostReply {
|
||||
stats.RLock()
|
||||
got := stats.Cache[bucket].Stats[failing.arn].ReplicatedCount
|
||||
size := stats.Cache[bucket].Stats[failing.arn].ReplicatedSize
|
||||
stats.RUnlock()
|
||||
if got != 1 || size != 0 {
|
||||
t.Fatalf("purge completion statistic=%d want 1", got)
|
||||
}
|
||||
}
|
||||
t.Logf("%s: %s recovered; %d target(s), persisted MRF, source/target metadata checked", backend, tc.name, tc.targets)
|
||||
}
|
||||
|
||||
func TestReplicationDeleteQueueFullRetryBudget(t *testing.T) {
|
||||
stats := NewReplicationStats(t.Context(), nil)
|
||||
p := &ReplicationPool{ctx: t.Context(), objLayer: &replicationMRFTestObjectLayer{}, workers: []chan ReplicationWorkerOperation{make(chan ReplicationWorkerOperation)}, stats: stats, mrfSaveCh: make(chan MRFReplicateEntry, 1), priority: "slow"}
|
||||
d := DeletedObjectReplicationInfo{Bucket: "bucket", DeletedObject: DeletedObject{ObjectName: "marker", VersionID: mustGetUUID()}, RetryCount: 2}
|
||||
p.queueReplicaDeleteTask(d)
|
||||
if got := <-p.mrfSaveCh; got.RetryCount != 3 {
|
||||
t.Fatalf("queue-full retry count=%d", got.RetryCount)
|
||||
}
|
||||
d.RetryCount = mrfRetryLimit
|
||||
p.queueReplicaDeleteTask(d)
|
||||
if len(p.mrfSaveCh) != 0 || atomic.LoadUint64(&stats.mrfStats.TotalDroppedCount) != 1 {
|
||||
t.Fatal("queue-full retry exceeded budget without a visible drop")
|
||||
}
|
||||
}
|
||||
|
||||
type markerLockFailureLayer struct{ ObjectLayer }
|
||||
|
||||
func (markerLockFailureLayer) NewNSLock(string, ...string) RWLocker { return markerFailedLock{} }
|
||||
|
||||
type markerFailedLock struct{ RWLocker }
|
||||
|
||||
func (markerFailedLock) GetLock(context.Context, *dynamicTimeout) (LockContext, error) {
|
||||
return LockContext{}, context.DeadlineExceeded
|
||||
}
|
||||
|
||||
type markerLookupLayer struct {
|
||||
ObjectLayer
|
||||
mutate func(ObjectInfo, error) (ObjectInfo, error)
|
||||
lookedUp chan struct{}
|
||||
}
|
||||
|
||||
func (l markerLookupLayer) GetObjectInfo(ctx context.Context, bucket, object string, opts ObjectOptions) (ObjectInfo, error) {
|
||||
oi, err := l.ObjectLayer.GetObjectInfo(ctx, bucket, object, opts)
|
||||
defer close(l.lookedUp)
|
||||
return l.mutate(oi, err)
|
||||
}
|
||||
|
||||
func testReplicationMRFInvalidLookups(ctx context.Context, t *testing.T, obj ObjectLayer, p *ReplicationPool, newPool func() *ReplicationPool, d DeletedObjectReplicationInfo) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
mutate func(ObjectInfo, error) (ObjectInfo, error)
|
||||
}{
|
||||
{"empty-info", func(_ ObjectInfo, e error) (ObjectInfo, error) { return ObjectInfo{}, e }},
|
||||
{"not-a-marker", func(o ObjectInfo, e error) (ObjectInfo, error) { o.DeleteMarker = false; return o, e }},
|
||||
{"wrong-version", func(o ObjectInfo, e error) (ObjectInfo, error) { o.VersionID = mustGetUUID(); return o, e }},
|
||||
{"wrong-bucket", func(o ObjectInfo, e error) (ObjectInfo, error) { o.Bucket = "different-bucket"; return o, e }},
|
||||
{"wrong-object", func(o ObjectInfo, e error) (ObjectInfo, error) { o.Name = "different-object"; return o, e }},
|
||||
{"zero-modtime", func(o ObjectInfo, e error) (ObjectInfo, error) { o.ModTime = time.Time{}; return o, e }},
|
||||
{"missing", func(o ObjectInfo, _ error) (ObjectInfo, error) {
|
||||
return o, ObjectNotFound{Bucket: o.Bucket, Object: o.Name}
|
||||
}},
|
||||
{"read-error", func(o ObjectInfo, _ error) (ObjectInfo, error) { return o, InsufficientReadQuorum{} }},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
entry := d.ToMRFEntry()
|
||||
p.saveMRFEntries(ctx, map[string]MRFReplicateEntry{entry.versionID: entry})
|
||||
fresh := newPool()
|
||||
looked := make(chan struct{})
|
||||
fresh.objLayer = markerLookupLayer{ObjectLayer: obj, mutate: tc.mutate, lookedUp: looked}
|
||||
if err := fresh.queueMRFHeal(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
select {
|
||||
case <-looked:
|
||||
case <-time.After(3 * time.Second):
|
||||
t.Fatal("disk MRF entry was not read")
|
||||
}
|
||||
select {
|
||||
case op := <-fresh.workers[0]:
|
||||
t.Fatalf("invalid lookup scheduled %T", op)
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Capture the actual internal audit event through the supported webhook sink.
|
||||
func replicationTestAudit(ctx context.Context, t *testing.T, bucket string) (<-chan string, func()) {
|
||||
t.Helper()
|
||||
if len(logger.AuditTargets()) != 0 {
|
||||
t.Fatal("unexpected pre-existing test audit targets")
|
||||
}
|
||||
statuses := make(chan string, 16)
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
var entry struct {
|
||||
API struct{ Name, Bucket, Status string }
|
||||
}
|
||||
if err := json.NewDecoder(r.Body).Decode(&entry); err != nil {
|
||||
t.Error(err)
|
||||
} else if entry.API.Name == ReplicateDeleteAPI && entry.API.Bucket == bucket {
|
||||
statuses <- entry.API.Status
|
||||
}
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
endpoint, err := xnet.ParseHTTPURL(server.URL)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if errs := logger.UpdateAuditWebhooks(ctx, map[string]loghttp.Config{"r6": {Enabled: true, Name: "r6", Endpoint: endpoint, BatchSize: 1, QueueSize: 128, MaxRetry: 1, RetryIntvl: time.Millisecond, HTTPTimeout: time.Second}}); len(errs) > 0 {
|
||||
t.Fatal(errs)
|
||||
}
|
||||
targets := logger.AuditTargets()
|
||||
return statuses, func() {
|
||||
logger.UpdateAuditWebhooks(ctx, nil)
|
||||
for _, target := range targets {
|
||||
target.Cancel()
|
||||
}
|
||||
server.Close()
|
||||
}
|
||||
}
|
||||
|
||||
type markerPurgeUpdateLayer struct {
|
||||
ObjectLayer
|
||||
t *testing.T
|
||||
}
|
||||
|
||||
func (l markerPurgeUpdateLayer) DeleteObject(ctx context.Context, bucket, object string, opts ObjectOptions) (ObjectInfo, error) {
|
||||
l.t.Helper()
|
||||
rs := opts.DeleteReplication
|
||||
if rs.ReplicationStatusInternal != "" || rs.Targets != nil || rs.ReplicaStatus != "" || !rs.CompositeReplicationStatus().Empty() {
|
||||
l.t.Fatalf("purge has a nonempty creation update: %+v", rs)
|
||||
}
|
||||
if rs.CompositeVersionPurgeStatus().Empty() {
|
||||
l.t.Fatal("purge update lost its purge status")
|
||||
}
|
||||
return l.ObjectLayer.DeleteObject(ctx, bucket, object, opts)
|
||||
}
|
||||
@@ -0,0 +1,201 @@
|
||||
// Copyright (c) 2026 PGSTY
|
||||
// SPDX-License-Identifier: AGPL-3.0-only
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/minio/minio-go/v7"
|
||||
"github.com/minio/minio/internal/bucket/replication"
|
||||
xhttp "github.com/minio/minio/internal/http"
|
||||
)
|
||||
|
||||
// Purge uses the version-delete wire form regardless of the in-memory shape.
|
||||
// Creation status, purge status and successful resync markers are independent.
|
||||
func TestReplicateDeleteOperationExits(t *testing.T) {
|
||||
for _, shape := range []string{"marker-creation", "legacy-marker-purge", "marker-purge", "object-purge"} {
|
||||
purge := shape != "marker-creation"
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
creation replication.StatusType
|
||||
priorPurge VersionPurgeStatusType
|
||||
headCode, deleteCode int
|
||||
headError string
|
||||
offline, resync bool
|
||||
}{
|
||||
{name: "pending-success", creation: replication.Pending, priorPurge: replication.VersionPurgePending, headCode: 404, deleteCode: 204},
|
||||
{name: "completed-creation", creation: replication.Completed, priorPurge: replication.VersionPurgePending, headCode: 405, deleteCode: 204},
|
||||
{name: "existing-marker", creation: replication.Pending, priorPurge: replication.VersionPurgePending, headCode: 405, deleteCode: 204},
|
||||
// Use the quorum S3 code without 503, which is classified as backend-down first.
|
||||
{name: "head-read-quorum", creation: replication.Pending, priorPurge: replication.VersionPurgePending, headCode: 400, headError: "SlowDownRead", deleteCode: 204},
|
||||
{name: "replica-creation-status", creation: replication.Replica, priorPurge: replication.VersionPurgeFailed, headCode: 405, deleteCode: 204},
|
||||
{name: "head-unavailable", creation: replication.Pending, priorPurge: replication.VersionPurgePending, headCode: 503, deleteCode: 204},
|
||||
{name: "head-forbidden", creation: replication.Pending, priorPurge: replication.VersionPurgePending, headCode: 403, deleteCode: 204},
|
||||
{name: "delete-forbidden", creation: replication.Pending, priorPurge: replication.VersionPurgePending, headCode: 404, deleteCode: 403},
|
||||
{name: "delete-method-rejected", creation: replication.Pending, priorPurge: replication.VersionPurgePending, headCode: 404, deleteCode: 405},
|
||||
{name: "delete-unavailable", creation: replication.Pending, priorPurge: replication.VersionPurgePending, headCode: 404, deleteCode: 503},
|
||||
{name: "offline", creation: replication.Pending, priorPurge: replication.VersionPurgePending, headCode: 404, deleteCode: 204, offline: true},
|
||||
{name: "retry-failed", creation: replication.Completed, priorPurge: replication.VersionPurgeFailed, headCode: 404, deleteCode: 204},
|
||||
{name: "purge-complete", creation: replication.Pending, priorPurge: replication.VersionPurgeComplete, headCode: 404, deleteCode: 403},
|
||||
{name: "resync-success", creation: replication.Completed, priorPurge: replication.VersionPurgePending, headCode: 405, deleteCode: 204, resync: true},
|
||||
{name: "resync-already-purged", creation: replication.Pending, priorPurge: replication.VersionPurgeComplete, headCode: 404, deleteCode: 403, resync: true},
|
||||
{name: "resync-failure", creation: replication.Completed, priorPurge: replication.VersionPurgePending, headCode: 404, deleteCode: 403, resync: true},
|
||||
} {
|
||||
t.Run(shape+"/"+tc.name, func(t *testing.T) {
|
||||
version := mustGetUUID()
|
||||
var heads, deletes atomic.Int32
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Query().Get("versionId") != version {
|
||||
t.Errorf("wrong request version: %s", r.URL)
|
||||
}
|
||||
switch r.Method {
|
||||
case http.MethodHead:
|
||||
heads.Add(1)
|
||||
w.Header().Set(xhttp.AmzVersionID, version)
|
||||
w.Header().Set(xhttp.LastModified, time.Now().UTC().Format(http.TimeFormat))
|
||||
if tc.headCode == 405 {
|
||||
w.Header().Set(xhttp.AmzDeleteMarker, "true")
|
||||
}
|
||||
if tc.headError != "" {
|
||||
w.Header().Set("x-minio-error-code", tc.headError)
|
||||
}
|
||||
w.WriteHeader(tc.headCode)
|
||||
case http.MethodDelete:
|
||||
deletes.Add(1)
|
||||
if got := r.Header.Get(xhttp.MinIOSourceDeleteMarker) == "true"; got == purge {
|
||||
t.Errorf("source delete-marker header=%v, purge=%v", got, purge)
|
||||
}
|
||||
w.WriteHeader(tc.deleteCode)
|
||||
if tc.deleteCode != 204 {
|
||||
fmt.Fprintf(w, `<Error><Code>%s</Code><Message>injected rejection</Message></Error>`, map[int]string{403: "AccessDenied", 405: "MethodNotAllowed", 503: "ServiceUnavailable"}[tc.deleteCode])
|
||||
}
|
||||
default:
|
||||
t.Errorf("unexpected method %s", r.Method)
|
||||
}
|
||||
}))
|
||||
defer server.Close()
|
||||
client, err := minio.New(strings.TrimPrefix(server.URL, "http://"), &minio.Options{Region: "us-east-1", MaxRetries: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
old := globalBucketTargetSys
|
||||
globalBucketTargetSys = &BucketTargetSys{hc: map[string]epHealth{client.EndpointURL().Host: {Online: !tc.offline}}}
|
||||
defer func() { globalBucketTargetSys = old }()
|
||||
d := DeletedObjectReplicationInfo{Bucket: "source", DeletedObject: DeletedObject{ObjectName: "marker", DeleteMarker: shape != "object-purge"}}
|
||||
if shape == "marker-creation" || shape == "legacy-marker-purge" {
|
||||
d.DeleteMarkerVersionID = version
|
||||
} else {
|
||||
d.VersionID = version
|
||||
}
|
||||
d.ReplicationState.Targets = map[string]replication.StatusType{"arn1": tc.creation}
|
||||
d.ReplicationState.ResetStatusesMap = map[string]string{"arn1": "previous-reset"}
|
||||
if purge {
|
||||
d.ReplicationState.PurgeTargets = map[string]VersionPurgeStatusType{"arn1": tc.priorPurge}
|
||||
}
|
||||
if tc.resync {
|
||||
d.OpType = replication.ExistingObjectReplicationType
|
||||
}
|
||||
got := replicateDeleteToTarget(t.Context(), d, &TargetClient{Client: client, ARN: "arn1", Bucket: "target", ResetID: "current-reset"})
|
||||
var wantCreation replication.StatusType
|
||||
var wantPurge VersionPurgeStatusType
|
||||
var wantHeads, wantDeletes int32
|
||||
success := false
|
||||
if purge {
|
||||
wantCreation = ""
|
||||
wantPurge = replication.VersionPurgeComplete
|
||||
switch {
|
||||
case tc.priorPurge == replication.VersionPurgeComplete:
|
||||
success = true
|
||||
case tc.offline:
|
||||
wantPurge = replication.VersionPurgeFailed
|
||||
default:
|
||||
wantDeletes = 1
|
||||
success = tc.deleteCode == 204
|
||||
if !success {
|
||||
wantPurge = replication.VersionPurgeFailed
|
||||
}
|
||||
}
|
||||
} else {
|
||||
switch {
|
||||
case tc.creation == replication.Completed && !tc.resync:
|
||||
success = true
|
||||
case tc.offline:
|
||||
default:
|
||||
wantHeads = 1
|
||||
switch tc.headCode {
|
||||
case 405:
|
||||
success = true
|
||||
case 403, 503:
|
||||
default:
|
||||
wantDeletes = 1
|
||||
success = tc.deleteCode == 204
|
||||
}
|
||||
}
|
||||
if success {
|
||||
wantCreation = replication.Completed
|
||||
} else {
|
||||
wantCreation = replication.Failed
|
||||
}
|
||||
}
|
||||
if got.PrevReplicationStatus != tc.creation {
|
||||
t.Error("previous creation state changed")
|
||||
}
|
||||
if got.ReplicationStatus != wantCreation || got.VersionPurgeStatus != wantPurge {
|
||||
t.Errorf("status=%+v, want creation=%s purge=%s", got, wantCreation, wantPurge)
|
||||
}
|
||||
if heads.Load() != wantHeads || deletes.Load() != wantDeletes {
|
||||
t.Errorf("HEAD/DELETE=%d/%d, want %d/%d", heads.Load(), deletes.Load(), wantHeads, wantDeletes)
|
||||
}
|
||||
if (got.Err != nil) == success {
|
||||
t.Errorf("success=%v, error=%v", success, got.Err)
|
||||
}
|
||||
if tc.resync && success {
|
||||
if !strings.HasSuffix(got.ResyncTimestamp, ";current-reset") {
|
||||
t.Errorf("successful resync missing reset: %q", got.ResyncTimestamp)
|
||||
}
|
||||
} else if got.ResyncTimestamp != "previous-reset" {
|
||||
t.Errorf("unexpected reset: %q", got.ResyncTimestamp)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestReplicateDeletePurgeMissingTargetState(t *testing.T) {
|
||||
// Classification belongs to the operation, even if this target has no
|
||||
// previous purge entry (another target supplies the tracked purge state).
|
||||
var deletes atomic.Int32
|
||||
version := mustGetUUID()
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodDelete || r.URL.Query().Get("versionId") != version || r.Header.Get(xhttp.MinIOSourceDeleteMarker) == "true" {
|
||||
t.Errorf("wrong purge request: %s %s %v", r.Method, r.URL, r.Header)
|
||||
}
|
||||
deletes.Add(1)
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}))
|
||||
defer server.Close()
|
||||
client, err := minio.New(strings.TrimPrefix(server.URL, "http://"), &minio.Options{Region: "us-east-1", MaxRetries: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
old := globalBucketTargetSys
|
||||
globalBucketTargetSys = &BucketTargetSys{hc: map[string]epHealth{client.EndpointURL().Host: {Online: true}}}
|
||||
defer func() { globalBucketTargetSys = old }()
|
||||
d := DeletedObjectReplicationInfo{Bucket: "source", DeletedObject: DeletedObject{ObjectName: "marker", DeleteMarker: true, DeleteMarkerVersionID: version, ReplicationState: ReplicationState{Targets: map[string]replication.StatusType{"arn1": replication.Completed}, PurgeTargets: map[string]VersionPurgeStatusType{"arn2": replication.VersionPurgePending}}}}
|
||||
got := replicateDeleteToTarget(t.Context(), d, &TargetClient{Client: client, ARN: "arn1", Bucket: "target", ResetID: "reset"})
|
||||
if deletes.Load() != 1 || got.VersionPurgeStatus != replication.VersionPurgeComplete || got.ReplicationStatus != "" {
|
||||
t.Fatalf("operation misclassified: %+v, deletes=%d", got, deletes.Load())
|
||||
}
|
||||
got.ResyncTimestamp = "resync;reset"
|
||||
state := getReplicationState(replicatedInfos{Targets: []replicatedTargetInfo{got}}, ReplicationState{}, "")
|
||||
if state.ResetStatusesMap[targetResetHeader("arn1")] != got.ResyncTimestamp {
|
||||
t.Fatal("resync timestamp not preserved with nil reset map")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user