From 9d7094b770fcb7a38399a8a6a1eb956df852d4d7 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Wed, 9 Sep 2026 14:32:09 +0800 Subject: [PATCH 1/5] fix(replication): complete single-object delete marker purges Signed-off-by: Feng Ruohang --- cmd/object-handlers.go | 2 +- cmd/replication-delete-marker_test.go | 236 ++++++++++++++++++++++++++ 2 files changed, 237 insertions(+), 1 deletion(-) create mode 100644 cmd/replication-delete-marker_test.go diff --git a/cmd/object-handlers.go b/cmd/object-handlers.go index 83d312ba3..f4cb876b3 100644 --- a/cmd/object-handlers.go +++ b/cmd/object-handlers.go @@ -3170,7 +3170,7 @@ func (api objectAPIHandlers) DeleteObjectHandler(w http.ResponseWriter, r *http. if objInfo.ReplicationStatus == replication.Pending || objInfo.VersionPurgeStatus == replication.VersionPurgePending { dmVersionID := "" versionID := "" - if objInfo.DeleteMarker { + if objInfo.DeleteMarker && objInfo.VersionPurgeStatus.Empty() { dmVersionID = objInfo.VersionID } else { versionID = objInfo.VersionID diff --git a/cmd/replication-delete-marker_test.go b/cmd/replication-delete-marker_test.go new file mode 100644 index 000000000..778d1823b --- /dev/null +++ b/cmd/replication-delete-marker_test.go @@ -0,0 +1,236 @@ +// Copyright (c) 2026 PGSTY +// SPDX-License-Identifier: AGPL-3.0-only + +package cmd + +import ( + "bytes" + "context" + "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 TestReplicateDeleteMarkerPurge(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) { + testReplicateDeleteMarkerPurge(obj, instanceType, bucket, router, creds, t, legacy) + }, + }) + }) + } +} + +func TestReplicateDeleteMarkerTargetSemantics(t *testing.T) { + for _, tt := range []struct { + name string + purge bool + deleteCode int + wantDelete bool + wantFailed bool + }{ + {name: "existing marker is idempotent"}, + {name: "purge removes an existing marker", purge: true, deleteCode: 204, wantDelete: true}, + {name: "purge forbidden", purge: true, deleteCode: 403, wantDelete: true, wantFailed: true}, + {name: "purge method rejected", purge: true, deleteCode: 405, wantDelete: true, wantFailed: true}, + {name: "purge unavailable", purge: true, deleteCode: 503, wantDelete: true, wantFailed: true}, + } { + t.Run(tt.name, func(t *testing.T) { + var deletes atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodHead { + w.Header().Set(xhttp.AmzDeleteMarker, "true") + w.Header().Set(xhttp.AmzVersionID, r.URL.Query().Get("versionId")) + w.WriteHeader(http.StatusMethodNotAllowed) + return + } + deletes.Add(1) + w.WriteHeader(tt.deleteCode) + if tt.wantFailed { + code := map[int]string{403: "AccessDenied", 405: "MethodNotAllowed", 503: "ServiceUnavailable"}[tt.deleteCode] + fmt.Fprintf(w, `%sinjected delete rejection`, code) + } + })) + 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 }() + deletion := DeletedObjectReplicationInfo{Bucket: "source", DeletedObject: DeletedObject{ + ObjectName: "marker", DeleteMarker: true, DeleteMarkerVersionID: mustGetUUID(), + }} + if tt.purge { + deletion.VersionID = deletion.DeleteMarkerVersionID + deletion.DeleteMarkerVersionID = "" + deletion.ReplicationState.PurgeTargets = map[string]VersionPurgeStatusType{"arn1": replication.VersionPurgePending} + } + result := replicateDeleteToTarget(t.Context(), deletion, &TargetClient{Client: client, ARN: "arn1", Bucket: "target"}) + if (deletes.Load() != 0) != tt.wantDelete { + t.Fatalf("remote DELETE count = %d, want delete %v", deletes.Load(), tt.wantDelete) + } + if tt.purge { + want := replication.VersionPurgeComplete + if tt.wantFailed { + want = replication.VersionPurgeFailed + } + if result.VersionPurgeStatus != want || (result.Err != nil) != tt.wantFailed { + t.Errorf("purge result = %+v, want %s", result, want) + } + } else if result.ReplicationStatus != replication.Completed { + t.Errorf("existing marker result = %+v, want Completed", result) + } + }) + } +} + +func testReplicateDeleteMarkerPurge(obj ObjectLayer, instanceType, bucket string, router http.Handler, creds auth.Credentials, t *testing.T, legacy bool) { + ctx := t.Context() + const arn = "arn:minio:replication::00000000-0000-4000-8000-000000000001:replica" + 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) + } + } + 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: + 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) + } + 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. + 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) + if result.VersionPurgeStatus() != replication.VersionPurgeComplete { + t.Errorf("remote purge result = %s, want COMPLETE", result.VersionPurgeStatus()) + } + for _, b := range []string{bucket, remoteBucket} { + oi, err := obj.GetObjectInfo(ctx, b, name, ObjectOptions{VersionID: version, Versioned: true}) + if !isErrVersionNotFound(err) && !isErrObjectNotFound(err) { + t.Errorf("%s: marker remains in %s: %+v, err=%v", instanceType, b, oi, err) + } + } +} From 63aace409954326bbc354cca70bbab5ca60083cf Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Wed, 9 Sep 2026 14:32:57 +0800 Subject: [PATCH 2/5] fix(replication): expose bounded MRF queue drops Signed-off-by: Feng Ruohang --- cmd/bucket-replication.go | 6 +- cmd/bucket-stats.go | 8 +- cmd/metrics-v2.go | 22 ++++ cmd/metrics-v3-replication.go | 8 ++ cmd/metrics-v3.go | 2 + cmd/replication-mrf-observability_test.go | 126 ++++++++++++++++++++++ docs/metrics/prometheus/list.md | 2 + docs/metrics/v3.md | 2 + 8 files changed, 172 insertions(+), 4 deletions(-) create mode 100644 cmd/replication-mrf-observability_test.go diff --git a/cmd/bucket-replication.go b/cmd/bucket-replication.go index def8a74f8..ce7fc69b0 100644 --- a/cmd/bucket-replication.go +++ b/cmd/bucket-replication.go @@ -2354,9 +2354,9 @@ func (p *ReplicationPool) queueReplicaTask(ri ReplicateObjectInfo) { switch ri.OpType { case replication.HealReplicationType, replication.ExistingObjectReplicationType: ch = p.mrfReplicaCh - healCh = p.getWorkerCh(ri.Name, ri.Bucket, ri.Size) + healCh = p.getWorkerCh(ri.Bucket, ri.Name, ri.Size) default: - ch = p.getWorkerCh(ri.Name, ri.Bucket, ri.Size) + ch = p.getWorkerCh(ri.Bucket, ri.Name, ri.Size) } if ch == nil && healCh == nil { return @@ -3820,6 +3820,7 @@ func (p *ReplicationPool) queueMRFSave(entry MRFReplicateEntry) { if entry.RetryCount > mrfRetryLimit { // let scanner catch up if retry count exceeded atomic.AddUint64(&p.stats.mrfStats.TotalDroppedCount, 1) atomic.AddUint64(&p.stats.mrfStats.TotalDroppedBytes, uint64(entry.sz)) + replLogOnceIf(GlobalContext, errors.New("Replication MRF retry limit reached; further repair is deferred to the scanner"), "replication-mrf-retry-limit", logger.WarningKind) return } @@ -3834,6 +3835,7 @@ func (p *ReplicationPool) queueMRFSave(entry MRFReplicateEntry) { default: atomic.AddUint64(&p.stats.mrfStats.TotalDroppedCount, 1) atomic.AddUint64(&p.stats.mrfStats.TotalDroppedBytes, uint64(entry.sz)) + replLogOnceIf(GlobalContext, errors.New("Replication MRF queue is full; dropped entries will need scanner repair"), "replication-mrf-queue-full", logger.WarningKind) } } } diff --git a/cmd/bucket-stats.go b/cmd/bucket-stats.go index e2eb00aa4..573c1caf1 100644 --- a/cmd/bucket-stats.go +++ b/cmd/bucket-stats.go @@ -319,7 +319,9 @@ func (r *ReplicationStats) getNodeQueueStats(bucket string) (qs ReplQNodeStats) qs.QStats = r.qCache.getBucketStats(bucket) qs.TgtXferStats = make(map[string]map[RMetricName]XferStats) qs.MRFStats = ReplicationMRFStats{ - LastFailedCount: atomic.LoadUint64(&r.mrfStats.LastFailedCount), + LastFailedCount: atomic.LoadUint64(&r.mrfStats.LastFailedCount), + TotalDroppedCount: atomic.LoadUint64(&r.mrfStats.TotalDroppedCount), + TotalDroppedBytes: atomic.LoadUint64(&r.mrfStats.TotalDroppedBytes), } r.RLock() @@ -410,7 +412,9 @@ func (r *ReplicationStats) getNodeQueueStatsSummary() (qs ReplQNodeStats) { qs.XferStats = make(map[RMetricName]XferStats) qs.QStats = r.qCache.getSiteStats() qs.MRFStats = ReplicationMRFStats{ - LastFailedCount: atomic.LoadUint64(&r.mrfStats.LastFailedCount), + LastFailedCount: atomic.LoadUint64(&r.mrfStats.LastFailedCount), + TotalDroppedCount: atomic.LoadUint64(&r.mrfStats.TotalDroppedCount), + TotalDroppedBytes: atomic.LoadUint64(&r.mrfStats.TotalDroppedBytes), } r.RLock() defer r.RUnlock() diff --git a/cmd/metrics-v2.go b/cmd/metrics-v2.go index b8d5c1290..120db3494 100644 --- a/cmd/metrics-v2.go +++ b/cmd/metrics-v2.go @@ -992,6 +992,26 @@ func getClusterReplMRFFailedOperationsMD() MetricDescription { } } +func getClusterReplMRFDroppedOperationsMD() MetricDescription { + return MetricDescription{ + Namespace: nodeMetricNamespace, + Subsystem: replicationSubsystem, + Name: "mrf_dropped_operations_total", + Help: "Total number of replication MRF entries dropped since server start; entries may refer to the same object", + Type: counterMetric, + } +} + +func getClusterReplMRFDroppedBytesMD() MetricDescription { + return MetricDescription{ + Namespace: nodeMetricNamespace, + Subsystem: replicationSubsystem, + Name: "mrf_dropped_bytes_total", + Help: "Total known bytes of replication MRF entries dropped since server start; delete entries count as zero bytes", + Type: counterMetric, + } +} + func getClusterRepCredentialErrorsMD(namespace MetricNamespace) MetricDescription { return MetricDescription{ Namespace: namespace, @@ -2423,6 +2443,8 @@ func getReplicationNodeMetrics(opts MetricsGroupOpts) *MetricsGroupV2 { avgTransferRate, maxTransferRate, mrfCount, + {Description: getClusterReplMRFDroppedOperationsMD(), Value: float64(qs.MRFStats.TotalDroppedCount)}, + {Description: getClusterReplMRFDroppedBytesMD(), Value: float64(qs.MRFStats.TotalDroppedBytes)}, } } for ep, health := range globalBucketTargetSys.healthStats() { diff --git a/cmd/metrics-v3-replication.go b/cmd/metrics-v3-replication.go index 44a8e87ae..adaf59ccd 100644 --- a/cmd/metrics-v3-replication.go +++ b/cmd/metrics-v3-replication.go @@ -35,6 +35,8 @@ const ( replicationMaxQueuedCount = "max_queued_count" replicationMaxDataTransferRate = "max_data_transfer_rate" replicationRecentBacklogCount = "recent_backlog_count" + replicationMRFDroppedOperations = "mrf_dropped_operations_total" + replicationMRFDroppedBytes = "mrf_dropped_bytes_total" ) var ( @@ -64,6 +66,10 @@ var ( "Maximum replication data transfer rate in bytes/sec seen since server start") replicationRecentBacklogCountMD = NewGaugeMD(replicationRecentBacklogCount, "Total number of objects seen in replication backlog in the last 5 minutes") + replicationMRFDroppedOperationsMD = NewCounterMD(replicationMRFDroppedOperations, + "Total number of replication MRF entries dropped since server start; entries may refer to the same object") + replicationMRFDroppedBytesMD = NewCounterMD(replicationMRFDroppedBytes, + "Total known bytes of replication MRF entries dropped since server start; delete entries count as zero bytes") ) // loadClusterReplicationMetrics - `MetricsLoaderFn` for cluster replication metrics @@ -96,6 +102,8 @@ func loadClusterReplicationMetrics(ctx context.Context, m MetricValues, c *metri m.Set(replicationMaxDataTransferRate, tots.Peak) } m.Set(replicationRecentBacklogCount, float64(qs.MRFStats.LastFailedCount)) + m.Set(replicationMRFDroppedOperations, float64(qs.MRFStats.TotalDroppedCount)) + m.Set(replicationMRFDroppedBytes, float64(qs.MRFStats.TotalDroppedBytes)) return nil } diff --git a/cmd/metrics-v3.go b/cmd/metrics-v3.go index 5b578448d..501589a26 100644 --- a/cmd/metrics-v3.go +++ b/cmd/metrics-v3.go @@ -342,6 +342,8 @@ func newMetricGroups(r *prometheus.Registry) *metricsV3Collection { replicationMaxQueuedCountMD, replicationMaxDataTransferRateMD, replicationRecentBacklogCountMD, + replicationMRFDroppedOperationsMD, + replicationMRFDroppedBytesMD, }, loadClusterReplicationMetrics, ) diff --git a/cmd/replication-mrf-observability_test.go b/cmd/replication-mrf-observability_test.go new file mode 100644 index 000000000..5d0c5af7b --- /dev/null +++ b/cmd/replication-mrf-observability_test.go @@ -0,0 +1,126 @@ +// Copyright (c) 2026 PGSTY +// SPDX-License-Identifier: AGPL-3.0-only + +package cmd + +import ( + "encoding/json" + "fmt" + "sync/atomic" + "testing" + + "github.com/minio/minio/internal/bucket/replication" + "github.com/prometheus/client_golang/prometheus" +) + +// A full MRF queue and an exhausted retry budget drop queue entries, not the +// source objects. Both drops must remain visible in the admin API snapshots. +func TestReplicationMRFDropsVisible(t *testing.T) { + stats := NewReplicationStats(t.Context(), nil) + old := globalReplicationStats.Swap(stats) + t.Cleanup(func() { globalReplicationStats.Store(old) }) + p := &ReplicationPool{ + objLayer: &replicationMRFTestObjectLayer{}, + stats: stats, + mrfSaveCh: make(chan MRFReplicateEntry, 1), + mrfStopCh: make(chan struct{}), + } + p.queueMRFSave(MRFReplicateEntry{sz: 10}) + p.queueMRFSave(MRFReplicateEntry{sz: 20}) + p.queueMRFSave(MRFReplicateEntry{sz: 30, RetryCount: mrfRetryLimit + 1}) + if len(p.mrfSaveCh) != 1 { + t.Fatalf("queue length = %d, want 1", len(p.mrfSaveCh)) + } + if got := atomic.LoadUint64(&stats.mrfStats.TotalDroppedCount); got != 2 { + t.Fatalf("dropped entries = %d, want 2", got) + } + for name, qs := range map[string]ReplQNodeStats{ + "bucket": stats.getNodeQueueStats("test-bucket"), + "node": stats.getNodeQueueStatsSummary(), + } { + t.Run(name, func(t *testing.T) { + data, err := json.Marshal(qs) + if err != nil { + t.Fatal(err) + } + var decoded ReplQNodeStats + if err := json.Unmarshal(data, &decoded); err != nil { + t.Fatal(err) + } + if got := decoded.MRFStats; got.TotalDroppedCount != 2 || got.TotalDroppedBytes != 50 { + t.Fatalf("API MRF stats = %+v, want dropped count 2 and bytes 50", got) + } + }) + } + oldTargets := globalBucketTargetSys + globalBucketTargetSys = &BucketTargetSys{} + t.Cleanup(func() { globalBucketTargetSys = oldTargets }) + want := map[string]float64{"mrf_dropped_operations_total": 2, "mrf_dropped_bytes_total": 50} + seen := make(map[string]bool) + for _, metric := range getReplicationNodeMetrics(MetricsGroupOpts{}).Get() { + name := string(metric.Description.Name) + if value, ok := want[name]; ok { + seen[name] = true + if metric.Value != value || metric.Description.Type != counterMetric { + t.Errorf("v2 %s = %+v, want counter %v", name, metric, value) + } + } + } + for name := range want { + if !seen[name] { + t.Errorf("v2 does not expose %s", name) + } + } + groups := newMetricGroups(prometheus.NewRegistry()) + families, err := groups.mgGatherers[replicationCollectorPath].Gather() + if err != nil { + t.Fatal(err) + } + seen = make(map[string]bool) + for _, family := range families { + for name, value := range want { + if family.GetName() != "minio_replication_"+name { + continue + } + seen[name] = true + if len(family.Metric) != 1 || family.Metric[0].Counter == nil || family.Metric[0].GetCounter().GetValue() != value { + t.Errorf("v3 %s = %v, want counter %v", name, family, value) + } + } + } + for name := range want { + if !seen[name] { + t.Errorf("v3 registry does not expose %s", name) + } + } +} + +type replicationMRFTestObjectLayer struct{ ObjectLayer } + +func TestReplicationObjectDeleteWorkerAffinity(t *testing.T) { + p := &ReplicationPool{ctx: t.Context(), workers: make([]chan ReplicationWorkerOperation, 8)} + for i := range p.workers { + p.workers[i] = make(chan ReplicationWorkerOperation, 2) + } + for _, op := range []replication.Type{replication.ObjectReplicationType, replication.HealReplicationType, replication.ExistingObjectReplicationType} { + for i := range 10 { + name := fmt.Sprintf("object-%d", i) + p.queueReplicaTask(ReplicateObjectInfo{Bucket: "bucket", Name: name, OpType: op}) + p.queueReplicaDeleteTask(DeletedObjectReplicationInfo{Bucket: "bucket", DeletedObject: DeletedObject{ObjectName: name}}) + objectWorker, deleteWorker := -1, -1 + for idx, ch := range p.workers { + for len(ch) > 0 { + switch (<-ch).(type) { + case ReplicateObjectInfo: + objectWorker = idx + case DeletedObjectReplicationInfo: + deleteWorker = idx + } + } + } + if objectWorker < 0 || objectWorker != deleteWorker { + t.Fatalf("operation %d, %s: object worker %d != delete worker %d", op, name, objectWorker, deleteWorker) + } + } + } +} diff --git a/docs/metrics/prometheus/list.md b/docs/metrics/prometheus/list.md index 050e2577b..acd977825 100644 --- a/docs/metrics/prometheus/list.md +++ b/docs/metrics/prometheus/list.md @@ -132,6 +132,8 @@ For deployments with [bucket](https://silo.pgsty.com/administration/bucket-repli | `minio_node_replication_max_queued_bytes` | Maximum number of bytes queued for replication seen since server start | | `minio_node_replication_max_queued_count` | Maximum number of objects queued for replication seen since server start | | `minio_node_replication_recent_backlog_count` | Total number of objects seen in replication backlog in the last 5 minutes | +| `minio_node_replication_mrf_dropped_operations_total` | Cumulative MRF entries dropped due to queue capacity or retry exhaustion; may count the same object more than once. Scanner repair remains available. | +| `minio_node_replication_mrf_dropped_bytes_total` | Cumulative known bytes of dropped MRF entries; delete entries count as zero bytes. | ## Healing Metrics diff --git a/docs/metrics/v3.md b/docs/metrics/v3.md index 192c2fda0..7583680f3 100644 --- a/docs/metrics/v3.md +++ b/docs/metrics/v3.md @@ -278,6 +278,8 @@ Metrics about Silo site and bucket replication. | `minio_replication_max_queued_count` | Maximum number of objects queued for replication since server start.

Type: gauge | `server` | | `minio_replication_max_data_transfer_rate` | Maximum replication data transfer rate in bytes/sec since server start.

Type: gauge | `server` | | `minio_replication_recent_backlog_count` | Total number of objects seen in replication backlog in the last 5 minutes

Type: gauge | `server` | +| `minio_replication_mrf_dropped_operations_total` | MRF entries dropped since server start due to queue capacity or retry exhaustion; entries may refer to the same object. Source objects remain eligible for scanner repair.

Type: counter | `server` | +| `minio_replication_mrf_dropped_bytes_total` | Known bytes of dropped MRF entries since server start; delete entries count as zero bytes.

Type: counter | `server` | #### `/bucket/replication` | Name | Description | Labels | From 702f113f5119114291893bb52b35aaff59c23cba Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Wed, 9 Sep 2026 14:34:42 +0800 Subject: [PATCH 3/5] fix(replication): scope resync cancellation and drain worker lifecycle Signed-off-by: Feng Ruohang --- cmd/bucket-replication-utils.go | 8 +- cmd/bucket-replication.go | 215 ++++++++----- cmd/bucket-replication_test.go | 45 ++- cmd/replication-resync-cancel_test.go | 442 ++++++++++++++++++++++++++ cmd/site-replication-utils.go | 22 +- cmd/site-replication.go | 33 +- 6 files changed, 636 insertions(+), 129 deletions(-) create mode 100644 cmd/replication-resync-cancel_test.go diff --git a/cmd/bucket-replication-utils.go b/cmd/bucket-replication-utils.go index 41b5dbc63..8c0ca7041 100644 --- a/cmd/bucket-replication-utils.go +++ b/cmd/bucket-replication-utils.go @@ -640,10 +640,10 @@ type VersionPurgeStatusType = replication.VersionPurgeStatusType type replicationResyncer struct { // map of bucket to their resync status - statusMap map[string]BucketReplicationResyncStatus - workerSize int - resyncCancelCh chan struct{} - workerCh chan struct{} + statusMap map[string]BucketReplicationResyncStatus + workerSize int + cancelResyncs map[resyncOpts]context.CancelCauseFunc + workerCh chan struct{} sync.RWMutex } diff --git a/cmd/bucket-replication.go b/cmd/bucket-replication.go index ce7fc69b0..d9c97cd2b 100644 --- a/cmd/bucket-replication.go +++ b/cmd/bucket-replication.go @@ -2964,10 +2964,10 @@ const ( func newresyncer() *replicationResyncer { rs := replicationResyncer{ - statusMap: make(map[string]BucketReplicationResyncStatus), - workerSize: resyncWorkerCnt, - resyncCancelCh: make(chan struct{}, resyncWorkerCnt), - workerCh: make(chan struct{}, resyncWorkerCnt), + statusMap: make(map[string]BucketReplicationResyncStatus), + workerSize: resyncWorkerCnt, + cancelResyncs: make(map[resyncOpts]context.CancelCauseFunc), + workerCh: make(chan struct{}, resyncWorkerCnt), } for i := 0; i < rs.workerSize; i++ { rs.workerCh <- struct{}{} @@ -2975,13 +2975,67 @@ func newresyncer() *replicationResyncer { return &rs } +var errResyncCanceled = errors.New("replication resync canceled") + +// Registration and cancellation share the status lock. A queued run therefore +// cannot miss cancellation between publishing its status and taking a slot. +func (s *replicationResyncer) registerResync(parent context.Context, opts resyncOpts) (context.Context, context.CancelCauseFunc, bool) { + s.Lock() + defer s.Unlock() + st, ok := s.statusMap[opts.bucket].TargetsMap[opts.arn] + if !ok || st.ResyncID != opts.resyncID { + return nil, nil, false + } + if _, running := s.cancelResyncs[opts]; running { + return nil, nil, false + } + ctx, cancel := context.WithCancelCause(parent) + if s.cancelResyncs == nil { + s.cancelResyncs = make(map[resyncOpts]context.CancelCauseFunc) + } + s.cancelResyncs[opts] = cancel + if st.ResyncStatus == ResyncCanceled { + cancel(errResyncCanceled) + } + return ctx, cancel, true +} + +func (s *replicationResyncer) cancelResyncID(resyncID string) { + s.Lock() + defer s.Unlock() + for bucket, m := range s.statusMap { + for arn, st := range m.TargetsMap { + if st.ResyncID == resyncID && (st.ResyncStatus == ResyncPending || st.ResyncStatus == ResyncStarted) { + st.ResyncStatus = ResyncCanceled + st.LastUpdate = UTCNow() + m.TargetsMap[arn] = st + m.LastUpdate = st.LastUpdate + } + } + s.statusMap[bucket] = m + } + for opts, cancel := range s.cancelResyncs { + if opts.resyncID == resyncID { + cancel(errResyncCanceled) + } + } +} + // mark status of replication resync on remote target for the bucket -func (s *replicationResyncer) markStatus(status ResyncStatusType, opts resyncOpts, objAPI ObjectLayer) { +func (s *replicationResyncer) markStatus(status ResyncStatusType, opts resyncOpts, objAPI ObjectLayer) ResyncStatusType { s.Lock() defer s.Unlock() m := s.statusMap[opts.bucket] - st := m.TargetsMap[opts.arn] + st, ok := m.TargetsMap[opts.arn] + if !ok || st.ResyncID != opts.resyncID { + return NoResync + } + // A cancel may win the lock after the finalizer checked its context. + // Persist that cancellation, never a stale Started/Completed result. + if st.ResyncStatus == ResyncCanceled { + status = ResyncCanceled + } st.LastUpdate = UTCNow() st.ResyncStatus = status m.TargetsMap[opts.arn] = st @@ -2991,14 +3045,18 @@ func (s *replicationResyncer) markStatus(status ResyncStatusType, opts resyncOpt ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() saveResyncStatus(ctx, opts.bucket, m, objAPI) + return status } // update replication resync stats for bucket's remote target -func (s *replicationResyncer) incStats(ts TargetReplicationResyncStatus, opts resyncOpts) { +func (s *replicationResyncer) incStats(ts TargetReplicationResyncStatus, opts resyncOpts) bool { s.Lock() defer s.Unlock() m := s.statusMap[opts.bucket] - st := m.TargetsMap[opts.arn] + st, ok := m.TargetsMap[opts.arn] + if !ok || st.ResyncID != opts.resyncID || st.ResyncStatus == ResyncCanceled { + return false + } st.Object = ts.Object st.ReplicatedCount += ts.ReplicatedCount st.FailedCount += ts.FailedCount @@ -3007,6 +3065,7 @@ func (s *replicationResyncer) incStats(ts TargetReplicationResyncStatus, opts re m.TargetsMap[opts.arn] = st m.LastUpdate = UTCNow() s.statusMap[opts.bucket] = m + return true } // resyncResults consumes the per-object outcomes produced by the resync worker @@ -3023,8 +3082,9 @@ type resyncResults struct { // result into the bucket's resync status. func (s *replicationResyncer) newResyncResults(opts resyncOpts) *resyncResults { return startResyncResults(func(r TargetReplicationResyncStatus) { - s.incStats(r, opts) - globalSiteResyncMetrics.updateMetric(r, opts.resyncID) + if s.incStats(r, opts) { + globalSiteResyncMetrics.updateMetric(r, opts.resyncID) + } }) } @@ -3063,29 +3123,22 @@ func (rr *resyncResults) finish(workers []chan ReplicateObjectInfo, workerWg *sy } // sendResyncResult delivers a worker's computed per-object result to ch, -// returning false if the worker must stop first. On the resync-cancel signal it -// records the abort - the already-computed result is dropped - so -// finalResyncStatus can downgrade a Completed run; on ctx cancellation it stops -// without recording, since finalResyncStatus's parent-context check covers that. -func (s *replicationResyncer) sendResyncResult(ctx context.Context, ch chan<- TargetReplicationResyncStatus, st TargetReplicationResyncStatus, workerAborted *atomic.Bool) bool { +// returning false if the run's context was canceled first. +func (s *replicationResyncer) sendResyncResult(ctx context.Context, ch chan<- TargetReplicationResyncStatus, st TargetReplicationResyncStatus) bool { select { case <-ctx.Done(): return false - case <-s.resyncCancelCh: - workerAborted.Store(true) - return false case ch <- st: return true } } -// finalResyncStatus downgrades a Completed status to Failed when the run could -// not have observed every object: the parent context was canceled (workers then -// return without sending their computed result) or a worker dropped a result on -// the resync-cancel signal. Without this a persisted Completed would misrepresent -// an incomplete resync. -func finalResyncStatus(status ResyncStatusType, ctxErr error, workerAborted bool) ResyncStatusType { - if status == ResyncCompleted && (ctxErr != nil || workerAborted) { +// User cancellation is terminal; interrupted completion remains retryable. +func finalResyncStatus(status ResyncStatusType, cause error) ResyncStatusType { + if errors.Is(cause, errResyncCanceled) { + return ResyncCanceled + } + if status == ResyncCompleted && cause != nil { return ResyncFailed } return status @@ -3156,26 +3209,31 @@ func objectNeedsResyncForARN(roi ReplicateObjectInfo, arn string) bool { // resyncBucket resyncs all qualifying objects as per replication rules for the target // ARN func (s *replicationResyncer) resyncBucket(ctx context.Context, objectAPI ObjectLayer, heal bool, opts resyncOpts) { + ctx, cancel, registered := s.registerResync(ctx, opts) + if !registered { + return + } + defer func() { + cancel(nil) + s.Lock() + delete(s.cancelResyncs, opts) + s.Unlock() + }() select { case <-s.workerCh: // block till a worker is available case <-ctx.Done(): + if errors.Is(context.Cause(ctx), errResyncCanceled) { + status := s.markStatus(ResyncCanceled, opts, objectAPI) + globalSiteResyncMetrics.incBucket(opts, status) + } return } resyncStatus := ResyncFailed - // workerAborted records that a worker dropped an already-computed result on - // the resync-cancel signal. With a canceled parent context (which makes - // workers return without sending their result), it means a Completed run did - // not actually observe every object - see finalResyncStatus below. - var workerAborted atomic.Bool defer func() { - // Downgrade a Completed status whose counts are incomplete, so the - // persisted status is not a misleading Completed. Runs after results.finish - // drains (LIFO) and before markStatus persists - markStatus uses its own - // background context, so a parent cancellation during the drain would - // otherwise still record Completed. - resyncStatus = finalResyncStatus(resyncStatus, ctx.Err(), workerAborted.Load()) - s.markStatus(resyncStatus, opts, objectAPI) + // Runs after workers/results drain and before our own deferred cancel. + resyncStatus = finalResyncStatus(resyncStatus, context.Cause(ctx)) + resyncStatus = s.markStatus(resyncStatus, opts, objectAPI) globalSiteResyncMetrics.incBucket(opts, resyncStatus) s.workerCh <- struct{}{} }() @@ -3210,8 +3268,8 @@ func (s *replicationResyncer) resyncBucket(ctx context.Context, objectAPI Object return } // mark resync status as resync started - if !heal { - s.markStatus(ResyncStarted, opts, objectAPI) + if !heal && s.markStatus(ResyncStarted, opts, objectAPI) != ResyncStarted { + return } // Walk through all object versions - Walk() is always in ascending order needed to ensure @@ -3237,7 +3295,14 @@ func (s *replicationResyncer) resyncBucket(ctx context.Context, objectAPI Object // cannot race the last incStats. Registered after the markStatus finalizer, so // LIFO runs finish first. results := s.newResyncResults(opts) - defer results.finish(workers, &wg) + defer func() { + if resyncStatus != ResyncCompleted { + cancel(nil) + } + // Exactly one finish: success drains with a live context; errors stop + // blocked workers first. Both drain counts before persisting status. + results.finish(workers, &wg) + }() for i := range resyncParallelRoutines { wg.Add(1) workers[i] = make(chan ReplicateObjectInfo, 100) @@ -3248,7 +3313,6 @@ func (s *replicationResyncer) resyncBucket(ctx context.Context, objectAPI Object select { case <-ctx.Done(): return - case <-s.resyncCancelCh: default: } traceFn := s.trace(tgt.ResetID, fmt.Sprintf("%s/%s (%s)", opts.bucket, roi.Name, roi.VersionID)) @@ -3295,26 +3359,29 @@ func (s *replicationResyncer) resyncBucket(ctx context.Context, objectAPI Object } } traceFn(traceSize, traceErr) - if !s.sendResyncResult(ctx, results.ch, st, &workerAborted) { + if !s.sendResyncResult(ctx, results.ch, st) { return } } }(ctx, i) } - for res := range objInfoCh { +walkLoop: + for { + var res itemOrErr[ObjectInfo] + select { + case <-ctx.Done(): + return + case item, ok := <-objInfoCh: + if !ok { + break walkLoop + } + res = item + } if res.Err != nil { resyncStatus = ResyncFailed replLogIf(ctx, res.Err) return } - select { - case <-s.resyncCancelCh: - resyncStatus = ResyncCanceled - return - case <-ctx.Done(): - return - default: - } if heal && lastCheckpoint != "" && lastCheckpoint != res.Item.Name { continue } @@ -3328,14 +3395,11 @@ func (s *replicationResyncer) resyncBucket(ctx context.Context, objectAPI Object if !objectNeedsResyncForARN(roi, opts.arn) { continue } + h := xxh3.HashString(roi.Bucket + roi.Name) select { - case <-s.resyncCancelCh: - return case <-ctx.Done(): return - default: - h := xxh3.HashString(roi.Bucket + roi.Name) - workers[h%uint64(resyncParallelRoutines)] <- roi + case workers[h%uint64(resyncParallelRoutines)] <- roi: } } resyncStatus = ResyncCompleted @@ -3363,9 +3427,9 @@ func (s *replicationResyncer) start(ctx context.Context, objAPI ObjectLayer, opt if len(tgtArns) == 0 { return fmt.Errorf("arn %s specified for resync not found in replication config", opts.arn) } - globalReplicationPool.Get().resyncer.RLock() - data, ok := globalReplicationPool.Get().resyncer.statusMap[opts.bucket] - globalReplicationPool.Get().resyncer.RUnlock() + s.Lock() + defer s.Unlock() + data, ok := s.statusMap[opts.bucket] if !ok { data, err = loadBucketResyncMetadata(ctx, opts.bucket, objAPI) if err != nil { @@ -3386,23 +3450,14 @@ func (s *replicationResyncer) start(ctx context.Context, objAPI ObjectLayer, opt ResyncStatus: ResyncPending, Bucket: opts.bucket, } + data.TargetsMap = data.cloneTgtStats() data.TargetsMap[opts.arn] = status if err = saveResyncStatus(ctx, opts.bucket, data, objAPI); err != nil { return err } - globalReplicationPool.Get().resyncer.Lock() - defer globalReplicationPool.Get().resyncer.Unlock() - brs, ok := globalReplicationPool.Get().resyncer.statusMap[opts.bucket] - if !ok { - brs = BucketReplicationResyncStatus{ - Version: resyncMetaVersion, - TargetsMap: make(map[string]TargetReplicationResyncStatus), - } - } - brs.TargetsMap[opts.arn] = status - globalReplicationPool.Get().resyncer.statusMap[opts.bucket] = brs - go globalReplicationPool.Get().resyncer.resyncBucket(GlobalContext, objAPI, false, opts) + s.statusMap[opts.bucket] = data + go s.resyncBucket(GlobalContext, objAPI, false, opts) return nil } @@ -3476,6 +3531,9 @@ func (p *ReplicationPool) loadResync(ctx context.Context, buckets []string, objA // Make sure only one node running resync on the cluster. ctx, cancel := globalLeaderLock.GetLock(ctx) defer cancel() + var workers sync.WaitGroup + // Keep the merged leader context alive until every resumed run exits. + defer workers.Wait() for index := range buckets { bucket := buckets[index] @@ -3489,19 +3547,24 @@ func (p *ReplicationPool) loadResync(ctx context.Context, buckets []string, objA } p.resyncer.Lock() - p.resyncer.statusMap[bucket] = meta - p.resyncer.Unlock() - + if current, ok := p.resyncer.statusMap[bucket]; ok { + // A concurrent start/cancel is newer than the disk snapshot. + meta = current + } else { + p.resyncer.statusMap[bucket] = meta + } tgts := meta.cloneTgtStats() + p.resyncer.Unlock() for arn, st := range tgts { switch st.ResyncStatus { case ResyncFailed, ResyncStarted, ResyncPending: - go p.resyncer.resyncBucket(ctx, objAPI, true, resyncOpts{ + opts := resyncOpts{ bucket: bucket, arn: arn, resyncID: st.ResyncID, resyncBefore: st.ResyncBeforeDate, - }) + } + workers.Go(func() { p.resyncer.resyncBucket(ctx, objAPI, true, opts) }) } } } diff --git a/cmd/bucket-replication_test.go b/cmd/bucket-replication_test.go index 23a6fd2b7..214cdadbd 100644 --- a/cmd/bucket-replication_test.go +++ b/cmd/bucket-replication_test.go @@ -26,7 +26,6 @@ import ( "path" "strings" "sync" - "sync/atomic" "testing" "testing/synctest" "time" @@ -328,20 +327,19 @@ func TestReplicationValidationObjectUsesRulePrefix(t *testing.T) { func newTestResyncer(bucket, arn string) (*replicationResyncer, resyncOpts) { s := &replicationResyncer{ - statusMap: map[string]BucketReplicationResyncStatus{}, - resyncCancelCh: make(chan struct{}, resyncWorkerCnt), + statusMap: map[string]BucketReplicationResyncStatus{}, } brs := newBucketResyncStatus(bucket) - brs.TargetsMap[arn] = TargetReplicationResyncStatus{ResyncStatus: ResyncStarted} + brs.TargetsMap[arn] = TargetReplicationResyncStatus{ResyncStatus: ResyncStarted, ResyncID: "reset-" + bucket} s.statusMap[bucket] = brs return s, resyncOpts{bucket: bucket, arn: arn, resyncID: "reset-" + bucket} } // TestResyncBucketFinalize round-trips the terminal status through a real // ObjectLayer: a clean run persists Completed with every result, while a run -// whose parent context was canceled during the drain, or in which a worker -// dropped a result on the cancel signal, is downgraded to Failed so a persisted -// Completed never misrepresents an incomplete resync. +// whose parent context was canceled during the drain is downgraded to Failed; +// a user-canceled run persists Canceled. Completed never misrepresents an +// incomplete resync. func TestResyncBucketFinalize(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() @@ -354,9 +352,9 @@ func TestResyncBucketFinalize(t *testing.T) { // persistTerminal applies resyncBucket's finalizer logic (finalResyncStatus // then markStatus, which persists) and reads the status back the way the // resync status API does. - persistTerminal := func(t *testing.T, s *replicationResyncer, opts resyncOpts, status ResyncStatusType, ctxErr error, aborted bool) TargetReplicationResyncStatus { + persistTerminal := func(t *testing.T, s *replicationResyncer, opts resyncOpts, status ResyncStatusType, cause error) TargetReplicationResyncStatus { t.Helper() - s.markStatus(finalResyncStatus(status, ctxErr, aborted), opts, objAPI) + s.markStatus(finalResyncStatus(status, cause), opts, objAPI) brs, err := loadBucketResyncMetadata(ctx, opts.bucket, objAPI) if err != nil { t.Fatalf("load persisted resync metadata: %v", err) @@ -376,7 +374,7 @@ func TestResyncBucketFinalize(t *testing.T) { var wg sync.WaitGroup // no producer workers for this case results.finish(nil, &wg) - st := persistTerminal(t, s, opts, ResyncCompleted, nil, false) + st := persistTerminal(t, s, opts, ResyncCompleted, nil) if st.ResyncStatus != ResyncCompleted { t.Fatalf("persisted status = %s, want Completed", st.ResyncStatus) } @@ -398,29 +396,24 @@ func TestResyncBucketFinalize(t *testing.T) { cctx, ccancel := context.WithCancel(context.Background()) ccancel() - st := persistTerminal(t, s, opts, ResyncCompleted, cctx.Err(), false) + st := persistTerminal(t, s, opts, ResyncCompleted, context.Cause(cctx)) if st.ResyncStatus != ResyncFailed { t.Fatalf("persisted status = %s, want Failed (parent canceled during drain)", st.ResyncStatus) } }) - // 3. A worker dropped a computed result on the resync-cancel token (parent - // still alive) -> sendResyncResult records the abort and Completed is - // downgraded to Failed. - t.Run("worker abort downgrades to failed", func(t *testing.T) { + // 3. A user-canceled worker cannot report a completed resync. + t.Run("user cancel persists canceled", func(t *testing.T) { s, opts := newTestResyncer("finalize-worker-abort", "arn1") - s.resyncCancelCh <- struct{}{} // cancel token waiting - ch := make(chan TargetReplicationResyncStatus) // no reader: the send would block - var aborted atomic.Bool - if s.sendResyncResult(context.Background(), ch, TargetReplicationResyncStatus{Object: "dropped", ReplicatedCount: 1}, &aborted) { - t.Fatal("sendResyncResult reported success despite the cancel token") + ctx, cancel := context.WithCancelCause(context.Background()) + cancel(errResyncCanceled) + ch := make(chan TargetReplicationResyncStatus) + if s.sendResyncResult(ctx, ch, TargetReplicationResyncStatus{Object: "dropped", ReplicatedCount: 1}) { + t.Fatal("sendResyncResult reported success after cancellation") } - if !aborted.Load() { - t.Fatal("worker abort was not recorded") - } - st := persistTerminal(t, s, opts, ResyncCompleted, nil, aborted.Load()) - if st.ResyncStatus != ResyncFailed { - t.Fatalf("persisted status = %s, want Failed (worker dropped a result)", st.ResyncStatus) + st := persistTerminal(t, s, opts, ResyncCompleted, context.Cause(ctx)) + if st.ResyncStatus != ResyncCanceled { + t.Fatalf("persisted status = %s, want Canceled", st.ResyncStatus) } }) } diff --git a/cmd/replication-resync-cancel_test.go b/cmd/replication-resync-cancel_test.go new file mode 100644 index 000000000..a524601c6 --- /dev/null +++ b/cmd/replication-resync-cancel_test.go @@ -0,0 +1,442 @@ +// Copyright (c) 2026 PGSTY +// SPDX-License-Identifier: AGPL-3.0-only + +package cmd + +import ( + "bytes" + "context" + "encoding/binary" + "errors" + "fmt" + "io" + "net/http" + "sync" + "testing" + "testing/synctest" + "time" + + "github.com/minio/madmin-go/v3" + "github.com/minio/minio-go/v7" + "github.com/minio/minio/internal/once" +) + +func TestSiteResyncCancelState(t *testing.T) { + globalSiteReplicationSys.Lock() + old := globalSiteReplicationSys.enabled + globalSiteReplicationSys.enabled = true + globalSiteReplicationSys.Unlock() + t.Cleanup(func() { + globalSiteReplicationSys.Lock() + globalSiteReplicationSys.enabled = old + globalSiteReplicationSys.Unlock() + }) + rs := newSiteResyncStatus("peer", []BucketInfo{{Name: "one"}, {Name: "two"}}) + sm := &siteResyncMetrics{ + resyncStatus: map[string]SiteResyncStatus{rs.ResyncID: rs.clone()}, + peerResyncMap: map[string]resyncState{"peer": {resyncID: rs.ResyncID}}, + } + rs.Status = ResyncCanceled + if err := sm.updateState(rs); err != nil { + t.Fatal(err) + } + got, err := sm.status("peer") + if err != nil { + t.Fatal(err) + } + if got.Status != ResyncCanceled { + t.Fatalf("cancel returned without updating site state: got %s", got.Status) + } + for _, bucket := range []string{"one", "two"} { + sm.incBucket(resyncOpts{bucket: bucket, resyncID: rs.ResyncID}, ResyncCanceled) + } + got, err = sm.status("peer") + if err != nil { + t.Fatal(err) + } + for bucket, status := range got.BucketStatuses { + if status != ResyncCanceled { + t.Errorf("bucket %s = %s, want Canceled", bucket, status) + } + } + // An already-running worker must not resurrect a canceled site. + sm.incBucket(resyncOpts{bucket: "one", resyncID: rs.ResyncID}, ResyncCompleted) + got, _ = sm.status("peer") + if got.Status != ResyncCanceled || got.BucketStatuses["one"] != ResyncCanceled { + t.Fatalf("late completion overwrote canceled state: %+v", got) + } +} + +// The fixture runs the real resyncBucket control flow with in-memory metadata +// persistence. Walk is deliberately controllable so cancellation does not +// depend on disk speed, network timing or the walker closing its output. +type resyncCancelObjectLayer struct { + ObjectLayer + walk func(context.Context, chan<- itemOrErr[ObjectInfo]) error + lock RWLocker + mu sync.Mutex + saved map[string]BucketReplicationResyncStatus +} + +func (o *resyncCancelObjectLayer) NewNSLock(bucket string, objects ...string) RWLocker { + return o.lock +} + +func (o *resyncCancelObjectLayer) Walk(ctx context.Context, bucket, prefix string, results chan<- itemOrErr[ObjectInfo], opts WalkOptions) error { + return o.walk(ctx, results) +} + +func (o *resyncCancelObjectLayer) PutObject(_ context.Context, bucket, object string, r *PutObjReader, opts ObjectOptions) (ObjectInfo, error) { + data, err := io.ReadAll(r) + if err == nil && o.saved != nil { + var status BucketReplicationResyncStatus + _, err = status.UnmarshalMsg(data[4:]) + o.mu.Lock() + o.saved[object] = status + o.mu.Unlock() + } + return ObjectInfo{}, err +} + +func (o *resyncCancelObjectLayer) GetObjectNInfo(_ context.Context, bucket, object string, _ *HTTPRangeSpec, _ http.Header, opts ObjectOptions) (*GetObjectReader, error) { + o.mu.Lock() + defer o.mu.Unlock() + status, ok := o.saved[object] + if !ok { + return nil, ObjectNotFound{Bucket: bucket, Object: object} + } + data := make([]byte, 4) + binary.LittleEndian.PutUint16(data[:2], resyncMetaFormat) + binary.LittleEndian.PutUint16(data[2:], resyncMetaVersion) + data, err := status.MarshalMsg(data) + if err != nil { + return nil, err + } + return NewGetObjectReaderFromReader(bytes.NewReader(data), ObjectInfo{Size: int64(len(data))}, opts) +} + +func TestResyncRecoveryOwnsLeaderContext(t *testing.T) { + for _, loseLeader := range []bool{false, true} { + t.Run(fmt.Sprintf("lose_leader_%v", loseLeader), func(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + var walkCtx context.Context + var output chan<- itemOrErr[ObjectInfo] + obj := &resyncCancelObjectLayer{saved: make(map[string]BucketReplicationResyncStatus), walk: func(ctx context.Context, ch chan<- itemOrErr[ObjectInfo]) error { + walkCtx, output = ctx, ch + return nil + }} + s, opts := setupResyncCancelTest(t, obj) + if err := saveResyncStatus(t.Context(), opts.bucket, s.statusMap[opts.bucket], obj); err != nil { + t.Fatal(err) + } + leaderCtx, lose := context.WithCancel(t.Context()) + defer lose() + leader := &sharedLock{lockContext: make(chan LockContext, 1)} + leader.lockContext <- LockContext{ctx: leaderCtx} + oldLeader := globalLeaderLock + globalLeaderLock = leader + defer func() { globalLeaderLock = oldLeader }() + done := make(chan error, 1) + go func() { done <- (&ReplicationPool{resyncer: s}).loadResync(t.Context(), []string{opts.bucket}, obj) }() + synctest.Wait() + if walkCtx == nil || walkCtx.Err() != nil { + t.Fatal("recovery canceled its worker at startup") + } + select { + case <-done: + t.Fatal("recovery released its leader context before the run finished") + default: + } + want := ResyncCompleted + if loseLeader { + want = ResyncFailed + lose() + } else { + close(output) + } + synctest.Wait() + if err := <-done; err != nil { + t.Fatal(err) + } + if walkCtx.Err() == nil { + t.Fatal("finished recovery leaked the Walk context") + } + if got := s.statusMap[opts.bucket].TargetsMap[opts.arn].ResyncStatus; got != want { + t.Fatalf("recovery status = %s, want %s", got, want) + } + }) + }) + } +} + +func TestResyncCancelRouting(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + parent, stop := context.WithCancel(t.Context()) + defer stop() + walks := make(chan context.Context, 4) + obj := &resyncCancelObjectLayer{walk: func(ctx context.Context, output chan<- itemOrErr[ObjectInfo]) error { + walks <- ctx + return nil + }} + s, opts := setupResyncCancelTest(t, obj) + add := func(bucket, id string) resyncOpts { + o := opts + o.bucket, o.resyncID = bucket, id + s.Lock() + m := newBucketResyncStatus(bucket) + m.TargetsMap[o.arn] = TargetReplicationResyncStatus{ResyncID: id, ResyncStatus: ResyncPending} + s.statusMap[bucket] = m + s.Unlock() + meta, _ := globalBucketMetadataSys.Get(opts.bucket) + globalBucketMetadataSys.Set(bucket, meta) + globalBucketTargetSys.Lock() + globalBucketTargetSys.targetsMap[bucket] = []madmin.BucketTarget{{Arn: o.arn}} + globalBucketTargetSys.Unlock() + return o + } + run := func(o resyncOpts) chan struct{} { + done := make(chan struct{}) + go func() { s.resyncBucket(parent, obj, false, o); close(done) }() + return done + } + first := run(opts) + synctest.Wait() + firstCtx := <-walks + queuedOpts := add("queued", opts.resyncID) + queued := run(queuedOpts) + otherOpts := add("other", "other-id") + other := run(otherOpts) + synctest.Wait() + s.cancelResyncID(opts.resyncID) + synctest.Wait() + for name, done := range map[string]chan struct{}{"active": first, "queued": queued} { + select { + case <-done: + default: + t.Fatalf("%s matching run survived cancellation", name) + } + } + if !errors.Is(context.Cause(firstCtx), errResyncCanceled) { + t.Fatal("active Walk did not receive user cancellation") + } + select { + case <-other: + t.Fatal("unrelated run was canceled") + default: + } + otherCtx := <-walks + if otherCtx.Err() != nil { + t.Fatal("unrelated Walk was canceled") + } + // There is no token/tombstone left for a subsequently registered run. + freshOpts := add("fresh", opts.resyncID) + fresh := run(freshOpts) + synctest.Wait() + select { + case <-fresh: + t.Fatal("fresh run inherited an old cancellation") + default: + } + stop() + synctest.Wait() + <-other + <-fresh + s.RLock() + defer s.RUnlock() + if len(s.cancelResyncs) != 0 || len(s.workerCh) != 1 { + t.Fatal("run registrations or worker slot leaked") + } + if s.statusMap[queuedOpts.bucket].TargetsMap[opts.arn].ResyncStatus != ResyncCanceled { + t.Fatal("queued user cancellation was not recorded") + } + if s.statusMap[freshOpts.bucket].TargetsMap[opts.arn].ResyncStatus != ResyncPending { + t.Fatal("shutdown changed a queued run's resumable Pending status") + } + }) +} + +func TestResyncCancellationWinsFinalization(t *testing.T) { + s, opts := newTestResyncer("finalize-cancel", "arn1") + obj := &resyncCancelObjectLayer{saved: make(map[string]BucketReplicationResyncStatus)} + ctx, cancel, registered := s.registerResync(t.Context(), opts) + if !registered { + t.Fatal("registration failed") + } + defer cancel(nil) + // Simulate cancellation after the finalizer computed Completed but before + // it acquired the persistence lock. + status := finalResyncStatus(ResyncCompleted, context.Cause(ctx)) + s.cancelResyncID(opts.resyncID) + if got := s.markStatus(status, opts, obj); got != ResyncCanceled { + t.Fatalf("final status = %s, want Canceled", got) + } + for _, saved := range obj.saved { + if saved.TargetsMap[opts.arn].ResyncStatus != ResyncCanceled { + t.Fatal("persisted Completed over cancellation") + } + } + if len(obj.saved) != 1 { + t.Fatal("canceled status was not persisted") + } + m := s.statusMap[opts.bucket] + m.TargetsMap[opts.arn] = TargetReplicationResyncStatus{ResyncID: "new-run", ResyncStatus: ResyncStarted} + s.statusMap[opts.bucket] = m + if s.markStatus(ResyncCompleted, opts, obj) != NoResync || s.incStats(TargetReplicationResyncStatus{ReplicatedCount: 1}, opts) { + t.Fatal("old run changed the replacement run") + } + if st := s.statusMap[opts.bucket].TargetsMap[opts.arn]; st.ResyncStatus != ResyncStarted || st.ReplicatedCount != 0 { + t.Fatalf("replacement state changed: %+v", st) + } + delete(s.statusMap, opts.bucket) + if s.markStatus(ResyncFailed, opts, obj) != NoResync { + t.Fatal("deleted bucket was recreated") + } +} + +func setupResyncCancelTest(t *testing.T, obj *resyncCancelObjectLayer) (*replicationResyncer, resyncOpts) { + t.Helper() + s, opts := newTestResyncer("cancel-bucket", "arn1") + s.workerCh = make(chan struct{}, 1) + s.workerCh <- struct{}{} + cfg := configs[0] + cfg.RoleArn = opts.arn + meta := newBucketMetadata(opts.bucket) + meta.replicationConfig = &cfg + oldMeta, oldTargets, oldObj := globalBucketMetadataSys, globalBucketTargetSys, newObjectLayerFn() + oldPool := globalReplicationPool + oldNotifier := globalEventNotifier + globalEventNotifier = &EventNotifier{} + globalReplicationPool = once.NewSingleton[ReplicationPool]() + globalReplicationPool.Set(&ReplicationPool{}) + globalBucketMetadataSys = NewBucketMetadataSys() + globalBucketMetadataSys.Set(opts.bucket, meta) + client, err := minio.New("127.0.0.1:1", &minio.Options{Region: "us-east-1"}) + if err != nil { + t.Fatal(err) + } + globalBucketTargetSys = &BucketTargetSys{ + arnRemotesMap: map[string]arnTarget{opts.arn: {Client: &TargetClient{Client: client, ARN: opts.arn}}}, + targetsMap: map[string][]madmin.BucketTarget{opts.bucket: {{Arn: opts.arn}}}, + } + setObjectLayer(obj) + t.Cleanup(func() { + globalBucketMetadataSys, globalBucketTargetSys = oldMeta, oldTargets + globalReplicationPool = oldPool + globalEventNotifier = oldNotifier + setObjectLayer(oldObj) + }) + return s, opts +} + +type blockedResyncLock struct { + RWLocker + started chan struct{} +} + +func (l *blockedResyncLock) GetLock(ctx context.Context, _ *dynamicTimeout) (LockContext, error) { + close(l.started) + <-ctx.Done() + return LockContext{}, ctx.Err() +} + +func TestResyncCancelFullWorkerQueue(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + lock := &blockedResyncLock{started: make(chan struct{})} + produced := 0 + walkerDone := make(chan struct{}) + obj := &resyncCancelObjectLayer{lock: lock} + s, opts := setupResyncCancelTest(t, obj) + obj.walk = func(ctx context.Context, output chan<- itemOrErr[ObjectInfo]) error { + go func() { + defer close(walkerDone) + defer close(output) + for range 120 { + select { + case <-ctx.Done(): + return + case output <- itemOrErr[ObjectInfo]{Item: ObjectInfo{Bucket: opts.bucket, Name: "same-key", VersionID: mustGetUUID(), DeleteMarker: true, ModTime: time.Now()}}: + produced++ + } + } + }() + return nil + } + done := make(chan struct{}) + go func() { s.resyncBucket(ctx, obj, false, opts); close(done) }() + synctest.Wait() + select { + case <-lock.started: + default: + t.Fatal("worker did not enter replication") + } + if produced < 101 || produced >= 120 { + t.Fatalf("expected a full 100-entry worker queue to block dispatch; Walk produced %d", produced) + } + cancel() + synctest.Wait() + select { + case <-done: + default: + t.Fatal("canceled dispatcher deadlocked sending to a full worker queue") + } + select { + case <-walkerDone: + default: + t.Fatal("canceled Walk producer leaked") + } + if len(s.workerCh) != 1 { + t.Fatal("resync worker slot was not returned") + } + }) +} + +func TestResyncCancelBlockedWalkReceive(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + parent, cancel := context.WithCancel(t.Context()) + defer cancel() + var output chan<- itemOrErr[ObjectInfo] + obj := &resyncCancelObjectLayer{walk: func(ctx context.Context, ch chan<- itemOrErr[ObjectInfo]) error { + output = ch + return nil + }} + s, opts := setupResyncCancelTest(t, obj) + done := make(chan struct{}) + go func() { s.resyncBucket(parent, obj, false, opts); close(done) }() + synctest.Wait() + if output == nil { + t.Fatal("resync did not reach Walk") + } + cancel() + synctest.Wait() + select { + case <-done: + default: + t.Error("resync remained blocked on Walk output after cancellation") + } + close(output) // release the old implementation on failure + synctest.Wait() + if len(s.workerCh) != 1 { + t.Error("resync worker slot was not released") + } + }) +} + +func TestResyncCancelsOwnedWalkOnError(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + var walkCtx context.Context + obj := &resyncCancelObjectLayer{walk: func(ctx context.Context, ch chan<- itemOrErr[ObjectInfo]) error { + walkCtx = ctx + return errors.New("injected walk failure") + }} + s, opts := setupResyncCancelTest(t, obj) + s.resyncBucket(t.Context(), obj, false, opts) + if walkCtx == nil || walkCtx.Err() == nil { + t.Fatal("resync exit did not cancel the context it passed to Walk") + } + if t.Context().Err() != nil { + t.Fatal("resync canceled its caller's context") + } + }) +} diff --git a/cmd/site-replication-utils.go b/cmd/site-replication-utils.go index 192275845..9c14d4dd3 100644 --- a/cmd/site-replication-utils.go +++ b/cmd/site-replication-utils.go @@ -204,16 +204,24 @@ func (sm *siteResyncMetrics) updateState(s SiteResyncStatus) error { switch s.Status { case ResyncStarted: sm.peerResyncMap[s.DeplID] = resyncState{resyncID: s.ResyncID, LastSaved: time.Time{}} - sm.resyncStatus[s.ResyncID] = s + sm.resyncStatus[s.ResyncID] = s.clone() case ResyncCompleted, ResyncCanceled, ResyncFailed: st, ok := sm.resyncStatus[s.ResyncID] if ok { st.LastUpdate = s.LastUpdate st.Status = s.Status + if s.Status == ResyncCanceled { + for bucket, status := range st.BucketStatuses { + if status == ResyncPending || status == ResyncStarted { + st.BucketStatuses[bucket] = ResyncCanceled + } + } + } + sm.resyncStatus[s.ResyncID] = st return nil } - sm.resyncStatus[s.ResyncID] = st - return saveSiteResyncMetadata(GlobalContext, st, newObjectLayerFn()) + sm.resyncStatus[s.ResyncID] = s.clone() + return saveSiteResyncMetadata(GlobalContext, s, newObjectLayerFn()) } return nil } @@ -230,7 +238,15 @@ func (sm *siteResyncMetrics) incBucket(o resyncOpts, bktStatus ResyncStatusType) if st.BucketStatuses == nil { st.BucketStatuses = map[string]ResyncStatusType{} } + if st.BucketStatuses[o.bucket] == ResyncCanceled { + return + } switch bktStatus { + case ResyncCanceled: + st.BucketStatuses[o.bucket] = ResyncCanceled + st.Status = ResyncCanceled + st.LastUpdate = UTCNow() + sm.resyncStatus[o.resyncID] = st case ResyncCompleted: st.BucketStatuses[o.bucket] = ResyncCompleted st.Status = siteResyncStatus(st.Status, st.BucketStatuses) diff --git a/cmd/site-replication.go b/cmd/site-replication.go index e04f481ab..31e843e6f 100644 --- a/cmd/site-replication.go +++ b/cmd/site-replication.go @@ -202,6 +202,7 @@ func wrapSRErr(err error) SRError { // SiteReplicationSys - manages cluster-level replication. type SiteReplicationSys struct { sync.RWMutex + resyncMu sync.Mutex // serialize site resync configuration start/cancel enabled bool @@ -6202,6 +6203,8 @@ func (c *SiteReplicationSys) getPeerForUpload(deplID string) (pi srPeerInfo, loc // is maintained in .minio.sys/buckets/site-replication/resync/, while collecting // individual bucket resync status in .minio.sys/buckets//replication/resync.bin func (c *SiteReplicationSys) startResync(ctx context.Context, objAPI ObjectLayer, peer madmin.PeerInfo) (res madmin.SRResyncOpStatus, err error) { + c.resyncMu.Lock() + defer c.resyncMu.Unlock() if !c.isEnabled() { return res, errSRNotEnabled } @@ -6321,6 +6324,8 @@ func (c *SiteReplicationSys) startResync(ctx context.Context, objAPI ObjectLayer // cancelResync stops an ongoing site level resync for the peer specified. func (c *SiteReplicationSys) cancelResync(ctx context.Context, objAPI ObjectLayer, peer madmin.PeerInfo) (res madmin.SRResyncOpStatus, err error) { + c.resyncMu.Lock() + defer c.resyncMu.Unlock() if !c.isEnabled() { return res, errSRNotEnabled } @@ -6386,34 +6391,22 @@ func (c *SiteReplicationSys) cancelResync(ctx context.Context, objAPI ObjectLaye }) continue } - // update resync state for the bucket - globalReplicationPool.Get().resyncer.Lock() - m, ok := globalReplicationPool.Get().resyncer.statusMap[bucket] - if !ok { - m = newBucketResyncStatus(bucket) - } - if st, ok := m.TargetsMap[t.Arn]; ok { - st.LastUpdate = UTCNow() - st.ResyncStatus = ResyncCanceled - m.TargetsMap[t.Arn] = st - m.LastUpdate = UTCNow() - } - globalReplicationPool.Get().resyncer.statusMap[bucket] = m - globalReplicationPool.Get().resyncer.Unlock() } } + // Configuration errors must not leave active or queued buckets running. + globalReplicationPool.Get().resyncer.cancelResyncID(rs.ResyncID) rs.Status = ResyncCanceled rs.LastUpdate = UTCNow() + for bucket, status := range rs.BucketStatuses { + if status == ResyncPending || status == ResyncStarted { + rs.BucketStatuses[bucket] = ResyncCanceled + } + } + globalSiteResyncMetrics.updateState(rs) if err := saveSiteResyncMetadata(ctx, rs, objAPI); err != nil { return res, err } - select { - case globalReplicationPool.Get().resyncer.resyncCancelCh <- struct{}{}: - case <-ctx.Done(): - } - - globalSiteResyncMetrics.updateState(rs) res.Status = rs.Status.String() return res, nil From a1141a43f250ace07aa804ec2cc9185e89bfdf38 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Wed, 9 Sep 2026 14:35:31 +0800 Subject: [PATCH 4/5] style: format replication regression fixture Signed-off-by: Feng Ruohang --- cmd/replication-delete-marker_test.go | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/cmd/replication-delete-marker_test.go b/cmd/replication-delete-marker_test.go index 778d1823b..dc31bcc8d 100644 --- a/cmd/replication-delete-marker_test.go +++ b/cmd/replication-delete-marker_test.go @@ -176,8 +176,13 @@ func testReplicateDeleteMarkerPurge(obj ObjectLayer, instanceType, bucket string 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)} + 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) From 66fe61ff65c83d68b74baa637a11623015c7aa21 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Wed, 9 Sep 2026 14:44:02 +0800 Subject: [PATCH 5/5] test: reuse replication compatibility fixtures Signed-off-by: Feng Ruohang --- cmd/replication-delete-marker_test.go | 2 +- cmd/replication-mrf-observability_test.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/cmd/replication-delete-marker_test.go b/cmd/replication-delete-marker_test.go index dc31bcc8d..e1ac57e6f 100644 --- a/cmd/replication-delete-marker_test.go +++ b/cmd/replication-delete-marker_test.go @@ -106,7 +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() - const arn = "arn:minio:replication::00000000-0000-4000-8000-000000000001:replica" + const arn = "arn:minio:replication::af470089-d354-4473-934c-9e1f52f6da89:bucket" const name = "marker" version := mustGetUUID() remoteBucket := getRandomBucketName() diff --git a/cmd/replication-mrf-observability_test.go b/cmd/replication-mrf-observability_test.go index 5d0c5af7b..bde6b1d65 100644 --- a/cmd/replication-mrf-observability_test.go +++ b/cmd/replication-mrf-observability_test.go @@ -79,7 +79,7 @@ func TestReplicationMRFDropsVisible(t *testing.T) { seen = make(map[string]bool) for _, family := range families { for name, value := range want { - if family.GetName() != "minio_replication_"+name { + if family.GetName() != replicationCollectorPath.metricPrefix()+"_"+name { continue } seen[name] = true