From eb4f5e5b3178da42b8b5137acfded62c35f6e3b2 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Wed, 16 Sep 2026 21:07:29 +0800 Subject: [PATCH] Fix exact-version purges and delete-marker metadata healing Require write-quorum absence for missing purge retries, aggregate removed and already-absent votes only for physical purges, and resolve receiver purges by version across pools. Preserve marker metadata through creation and healing while retaining DELETE response semantics. Add real-disk quorum, pool, callback, metadata, heal and outbound replication regressions. Synchronize capacity fixture installation and restoration with concurrent IAM readers. Signed-off-by: Feng Ruohang (cherry picked from commit 22ba4d426677aa07470927fc349a43afd87e19b6) --- cmd/erasure-delete-marker-purge_test.go | 524 ++++++++++++++++++ cmd/erasure-object.go | 36 +- cmd/erasure-server-pool-consistency.go | 20 + cmd/erasure-server-pool.go | 19 +- cmd/erasure-version-purge.go | 87 +++ cmd/replication-delete-mrf_test.go | 45 +- cmd/xl-storage-delete-marker-metadata_test.go | 148 +++++ cmd/xl-storage-delete-marker.go | 67 +++ cmd/xl-storage-format-v2.go | 2 +- cmd/xl-storage.go | 4 +- 10 files changed, 938 insertions(+), 14 deletions(-) create mode 100644 cmd/erasure-delete-marker-purge_test.go create mode 100644 cmd/erasure-version-purge.go create mode 100644 cmd/xl-storage-delete-marker-metadata_test.go create mode 100644 cmd/xl-storage-delete-marker.go diff --git a/cmd/erasure-delete-marker-purge_test.go b/cmd/erasure-delete-marker-purge_test.go new file mode 100644 index 000000000..de3221aa1 --- /dev/null +++ b/cmd/erasure-delete-marker-purge_test.go @@ -0,0 +1,524 @@ +// Copyright (c) 2026 Feng Ruohang +// +// This file is part of Silo Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "bytes" + "context" + "errors" + "fmt" + "maps" + "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/bucket/replication" + xhttp "github.com/minio/minio/internal/http" + "github.com/minio/minio/internal/once" +) + +func TestDeleteMarkerPurgeIntent(t *testing.T) { + version := mustGetUUID() + for _, tc := range []struct { + name string + edit func(*ObjectOptions) + want bool + }{ + {"ordinary", func(*ObjectOptions) {}, true}, + {"receiver", func(o *ObjectOptions) { o.ReplicationRequest = true; o.SetReplicaStatus(replication.Replica) }, true}, + {"untrusted-replica", func(o *ObjectOptions) { o.SetReplicaStatus(replication.Replica) }, false}, + {"complete-purge", func(o *ObjectOptions) { + o.DeleteReplication.VersionPurgeStatusInternal = string(replication.VersionPurgeComplete) + }, true}, + {"pending-purge", func(o *ObjectOptions) { o.DeleteReplication.VersionPurgeStatusInternal = "PENDING" }, false}, + {"failed-purge", func(o *ObjectOptions) { o.DeleteReplication.VersionPurgeStatusInternal = "FAILED" }, false}, + {"unknown-purge", func(o *ObjectOptions) { o.DeleteReplication.VersionPurgeStatusInternal = "future-state" }, false}, + {"pending-creation", func(o *ObjectOptions) { o.DeleteReplication.ReplicationStatusInternal = "PENDING" }, false}, + {"failed-creation", func(o *ObjectOptions) { o.DeleteReplication.ReplicationStatusInternal = "FAILED" }, false}, + {"completed-creation", func(o *ObjectOptions) { o.DeleteReplication.ReplicationStatusInternal = "COMPLETED" }, false}, + {"create-marker", func(o *ObjectOptions) { o.DeleteMarker = true }, false}, + {"empty", func(o *ObjectOptions) { o.VersionID = "" }, false}, + {"null", func(o *ObjectOptions) { o.VersionID = nullVersionID }, false}, + {"invalid", func(o *ObjectOptions) { o.VersionID = "invalid" }, false}, + {"zero-uuid", func(o *ObjectOptions) { o.VersionID = emptyUUID }, false}, + {"movement", func(o *ObjectOptions) { o.DataMovement = true }, false}, + {"free-version", func(o *ObjectOptions) { o.InclFreeVersions = true }, false}, + {"expiration", func(o *ObjectOptions) { o.Expiration.Expire = true }, false}, + {"transition", func(o *ObjectOptions) { o.Transition.Status = "complete" }, false}, + {"restored-expiration", func(o *ObjectOptions) { o.Transition.ExpireRestored = true }, false}, + } { + t.Run(tc.name, func(t *testing.T) { + o := ObjectOptions{VersionID: version, Versioned: true} + tc.edit(&o) + if got := o.isVersionPurge(); got != tc.want { + t.Fatalf("purge=%v want=%v", got, tc.want) + } + }) + } +} + +// Every other StorageAPI method panics through the nil embedding: a proof may +// read metadata, but must never delete, heal, or ask a different storage API. +type absenceProofDisk struct { + StorageAPI + t *testing.T + err error + reads *atomic.Int32 + deletes *atomic.Int32 +} + +func (d absenceProofDisk) ReadVersion(_ context.Context, _, _, _, _ string, opts ReadOptions) (FileInfo, error) { + if opts.ReadData || opts.Healing { + d.t.Error("absence proof requested data/healing") + } + d.reads.Add(1) + return FileInfo{}, d.err +} + +func (d absenceProofDisk) DeleteVersion(_ context.Context, _, _ string, _ FileInfo, _ bool, _ DeleteOptions) error { + if d.deletes == nil { + d.t.Error("read-only absence proof attempted a deletion") + } else { + d.deletes.Add(1) + } + return d.err +} + +func TestVersionPurgeAbsenceProof(t *testing.T) { + for _, tc := range []struct { + name string + errs []error + quorum, all bool + }{ + {"all-absent", []error{errFileNotFound, errFileVersionNotFound, errFileNotFound, errFileVersionNotFound}, true, true}, + {"majority-absent", []error{errFileNotFound, errFileVersionNotFound, errFileNotFound, nil}, true, false}, + {"half-absent", []error{errFileNotFound, errFileVersionNotFound, errDiskNotFound, errDiskNotFound}, false, false}, + {"corrupt", []error{errFileNotFound, errFileVersionNotFound, errFileCorrupt, errFileCorrupt}, false, false}, + {"permission", []error{errFileNotFound, errFileVersionNotFound, errDiskAccessDenied, errVolumeAccessDenied}, false, false}, + {"volume-missing", []error{errFileNotFound, errFileVersionNotFound, errVolumeNotFound, errVolumeNotFound}, false, false}, + } { + t.Run(tc.name, func(t *testing.T) { + var reads atomic.Int32 + disks := make([]StorageAPI, len(tc.errs)) + for i, err := range tc.errs { + disks[i] = absenceProofDisk{t: t, err: err, reads: &reads} + } + er := erasureObjects{getDisks: func() []StorageAPI { return disks }} + oldQueue := globalMRFState.opCh + globalMRFState.opCh = make(chan PartialOperation, 1) + defer func() { globalMRFState.opCh = oldQueue }() + all, err := er.confirmVersionAbsent(t.Context(), "bucket", "object", mustGetUUID()) + if (err == nil) != tc.quorum || all != tc.all || reads.Load() != int32(len(disks)) || len(globalMRFState.opCh) != 0 { + t.Fatalf("proof all=%v error=%v reads=%d queued=%d", all, err, reads.Load(), len(globalMRFState.opCh)) + } + }) + } +} + +func TestVersionPurgeAggregation(t *testing.T) { + for _, pure := range []bool{false, true} { + for _, fault := range []error{nil, errFileCorrupt, errDiskAccessDenied} { + t.Run(fmt.Sprintf("purge_%v_fault_%v", pure, fault), func(t *testing.T) { + var calls atomic.Int32 + disks := make([]StorageAPI, 4) + for i, err := range []error{nil, errFileNotFound, errFileVersionNotFound, fault} { + disks[i] = absenceProofDisk{t: t, err: err, deletes: &calls} + } + er := erasureObjects{getDisks: func() []StorageAPI { return disks }} + err := er.deleteObjectVersion(t.Context(), "bucket", "object", FileInfo{VersionID: mustGetUUID()}, false, pure) + if (err == nil) != pure || calls.Load() != 4 { + t.Fatalf("aggregation error=%v calls=%d", err, calls.Load()) + } + }) + } + } +} + +func TestDeleteMarkerPurgeOmittedPool(t *testing.T) { + z, bucket := consistencyPools(t) + defer replicationTestCapacity(z)() + const name = "partially-absent" + er := z.serverPools[1].getHashedSet(name) + _, opts := seedPurgeMarker(t, er, bucket, name, false) + opts.ReplicationRequest = false + opts.DeleteReplication = ReplicationState{} + // Pool 0 appears absent at read quorum, but half of it cannot be read. + missing := z.serverPools[0].getHashedSet(name) + original := missing.getDisks + disks := append([]StorageAPI(nil), original()...) + clear(disks[len(disks)/2:]) + missing.getDisks = func() []StorageAPI { return disks } + defer func() { missing.getDisks = original }() + _, err := z.DeleteObject(t.Context(), bucket, name, opts) + if !isErrWriteQuorum(err) { + t.Errorf("omitted pool acknowledged deletion: %T %v", err, err) + } + if countPurgeMarkers(t, er.getDisks(), bucket, name, opts.VersionID) != 16 { + t.Fatal("mutated known copies before checking omitted pool") + } +} + +func TestDeleteMarkerPurgeReceivingPool(t *testing.T) { + z, bucket := consistencyPools(t) + defer replicationTestCapacity(z)() + for _, duplicate := range []bool{false, true} { + t.Run(fmt.Sprintf("duplicate_%v", duplicate), func(t *testing.T) { + name := "receiver-" + mustGetUUID() + _, opts := seedPurgeMarker(t, z.serverPools[0].getHashedSet(name), bucket, name, false) + if duplicate { + create := opts + create.DeleteMarker = true + if _, err := z.serverPools[1].DeleteObject(t.Context(), bucket, name, create); err != nil { + t.Fatal(err) + } + } else { + // Latest-key routing selects pool 1, but the addressed marker is + // in pool 0. A replica purge addresses a version, not the latest. + putConsistencyObject(t, z, bucket, name, 1, "newer-data", ObjectOptions{Versioned: true}) + } + if _, err := z.DeleteObject(t.Context(), bucket, name, opts); err != nil { + t.Errorf("replica purge did not find its addressed version: %v", err) + } + for i, pool := range z.serverPools { + if got := countPurgeMarkers(t, pool.getHashedSet(name).getDisks(), bucket, name, opts.VersionID); got != 0 { + t.Errorf("replica purge left %d marker copies in pool %d", got, i) + } + } + }) + } +} + +func TestDeleteMarkerPurgeCallbackIntent(t *testing.T) { + for _, loseCopies := range []bool{false, true} { + t.Run(fmt.Sprintf("missing_after_callback_%v", loseCopies), func(t *testing.T) { + z, bucket := consistencyPools(t) + defer replicationTestCapacity(z)() + const name = "callback-pending" + er := z.serverPools[0].getHashedSet(name) + _, opts := seedPurgeMarker(t, er, bucket, name, false) + original := er.getDisks + all := append([]StorageAPI(nil), original()...) + defer func() { er.getDisks = original }() + opts.ReplicationRequest = false + opts.DeleteReplication = ReplicationState{} + metadata, retention := 0, 0 + opts.EvalRetentionBypassFn = func(oi ObjectInfo, err error) error { + retention++ + if !oi.DeleteMarker || !isErrMethodNotAllowed(err) { + t.Fatalf("retention callback lost marker: %+v %v", oi, err) + } + return nil + } + opts.EvalMetadataFn = func(*ObjectInfo, error) (ReplicateDecision, error) { + metadata++ + if loseCopies { + // Inject a disk-state change after preflight, before the set's + // metadata-update read. It returns NotFound at read quorum. + for _, disk := range all[:8] { + if err := disk.DeleteVersion(t.Context(), bucket, name, FileInfo{VersionID: opts.VersionID}, false, DeleteOptions{}); err != nil { + t.Fatal(err) + } + } + partial := make([]StorageAPI, len(all)) + copy(partial, all[:8]) + er.getDisks = func() []StorageAPI { return partial } + } + decision := ReplicateDecision{} + decision.Set(newReplicateTargetDecision("arn1", true, false)) + return decision, nil + } + _, err := z.DeleteObject(t.Context(), bucket, name, opts) + if loseCopies { + if !isErrVersionNotFound(err) && !isErrObjectNotFound(err) { + t.Errorf("pending update misclassified as physical purge: %T %v", err, err) + } + } else if err != nil { + t.Errorf("pending update failed: %v", err) + } + want := 16 + if loseCopies { + want = 8 + } + if got := countPurgeMarkers(t, all, bucket, name, opts.VersionID); got != want || metadata != 1 || retention != 1 { + t.Errorf("pending update: copies=%d want=%d metadata=%d retention=%d", got, want, metadata, retention) + } + if !loseCopies { + fi, err := all[0].ReadVersion(t.Context(), "", bucket, name, opts.VersionID, ReadOptions{}) + if err != nil || fi.VersionPurgeStatus() != replication.VersionPurgePending { + t.Errorf("pending state not persisted: %+v %v", fi.ReplicationState, err) + } + } + }) + } +} + +func TestDeleteMarkerReceivingMetadataUpdate(t *testing.T) { + z, bucket := consistencyPools(t) + defer replicationTestCapacity(z)() + for _, state := range []VersionPurgeStatusType{replication.VersionPurgePending, replication.VersionPurgeFailed} { + t.Run(string(state), func(t *testing.T) { + name := "update-" + mustGetUUID() + _, opts := seedPurgeMarker(t, z.serverPools[0].getHashedSet(name), bucket, name, false) + create := opts + create.DeleteMarker = true + if _, err := z.serverPools[1].DeleteObject(t.Context(), bucket, name, create); err != nil { + t.Fatal(err) + } + opts.DeleteReplication = ReplicationState{VersionPurgeStatusInternal: string(state)} + if _, err := z.DeleteObject(t.Context(), bucket, name, opts); err != nil { + t.Fatal(err) + } + for i, pool := range z.serverPools { + if got := countPurgeMarkers(t, pool.getHashedSet(name).getDisks(), bucket, name, opts.VersionID); got != 16 { + t.Errorf("metadata update removed marker copies in pool %d: %d", i, got) + } + } + }) + } +} + +func TestVersionPurgeRetentionGate(t *testing.T) { + z, er, all, bucket := markerPurgeFixture(t, 4) + data, _ := seedPurgeMarker(t, er, bucket, "protected-data", true) + denied := errors.New("retention denied") + opts := ObjectOptions{Versioned: true, VersionID: data, EvalRetentionBypassFn: func(oi ObjectInfo, err error) error { + if err != nil || oi.VersionID != data || oi.DeleteMarker { + t.Errorf("wrong retained version: %+v error=%v", oi, err) + } + return denied + }} + _, err := z.DeleteObject(t.Context(), bucket, "protected-data", opts) + if !errors.Is(err, denied) { + t.Fatalf("retention gate: %v", err) + } + for _, disk := range all { + if _, err := disk.ReadVersion(t.Context(), "", bucket, "protected-data", data, ReadOptions{}); err != nil { + t.Errorf("protected version changed: %v", err) + } + } +} + +func markerPurgeFixture(t *testing.T, n int) (*erasureServerPools, *erasureObjects, []StorageAPI, string) { + t.Helper() + ctx, cancel := context.WithCancel(t.Context()) + obj, dirs, err := prepareErasure(ctx, n) + if err != nil { + cancel() + t.Fatal(err) + } + restore := replicationTestCapacity(obj) + t.Cleanup(func() { + cancel() + restore() + obj.Shutdown(context.Background()) + removeRoots(dirs) + }) + bucket := getRandomBucketName() + if err := obj.MakeBucket(ctx, bucket, MakeBucketOptions{VersioningEnabled: true}); err != nil { + t.Fatal(err) + } + z := obj.(*erasureServerPools) + er := z.serverPools[0].sets[0] + original := er.getDisks + t.Cleanup(func() { er.getDisks = original }) + return z, er, append([]StorageAPI(nil), original()...), bucket +} + +func seedPurgeMarker(t *testing.T, er *erasureObjects, bucket, name string, withData bool) (string, ObjectOptions) { + t.Helper() + dataVersion := "" + if withData { + oi, err := er.PutObject(t.Context(), bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true}) + if err != nil { + t.Fatal(err) + } + dataVersion = oi.VersionID + } + opts := ObjectOptions{Versioned: true, VersionID: mustGetUUID(), DeleteMarker: true, ReplicationRequest: true, MTime: UTCNow().Add(-time.Hour), NoAuditLog: true} + opts.SetReplicaStatus(replication.Replica) + if _, err := er.DeleteObject(t.Context(), bucket, name, opts); err != nil { + t.Fatal(err) + } + opts.DeleteMarker = false + return dataVersion, opts +} + +func countPurgeMarkers(t *testing.T, disks []StorageAPI, bucket, name, version string) int { + t.Helper() + n := 0 + for _, disk := range disks { + fi, err := disk.ReadVersion(t.Context(), "", bucket, name, version, ReadOptions{}) + if err == nil && fi.Deleted { + n++ + } else if err != errFileNotFound && err != errFileVersionNotFound { + t.Fatalf("unexpected disk version: deleted=%v error=%v", fi.Deleted, err) + } + } + return n +} + +func TestDeleteMarkerPurgeQuorum(t *testing.T) { + for _, tc := range []struct{ disks, online int }{{4, 2}, {16, 7}, {16, 8}, {16, 9}} { + t.Run(fmt.Sprintf("%d_disks_%d_online", tc.disks, tc.online), func(t *testing.T) { + z, er, all, bucket := markerPurgeFixture(t, tc.disks) + const name = "marker" + dataVersion, opts := seedPurgeMarker(t, er, bucket, name, true) + partial := make([]StorageAPI, len(all)) + copy(partial, all[:tc.online]) + er.getDisks = func() []StorageAPI { return partial } + _, first := z.DeleteObject(t.Context(), bucket, name, opts) + _, retry := z.DeleteObject(t.Context(), bucket, name, opts) + if tc.online <= tc.disks/2 { + if !isErrWriteQuorum(first) || !isErrWriteQuorum(retry) { + t.Errorf("below write quorum: first=%T %v, retry=%T %v", first, first, retry, retry) + } + } else if first != nil || (retry != nil && !isErrVersionNotFound(retry) && !isErrObjectNotFound(retry)) { + t.Errorf("write quorum available: first=%v retry=%v", first, retry) + } + want := tc.disks - tc.online + if tc.online < tc.disks/2 { + want = tc.disks + } + if got := countPurgeMarkers(t, all, bucket, name, opts.VersionID); got != want { + t.Errorf("partial purge markers=%d want=%d", got, want) + } + er.getDisks = func() []StorageAPI { return all } + _, err := z.DeleteObject(t.Context(), bucket, name, opts) + if err != nil && !isErrVersionNotFound(err) && !isErrObjectNotFound(err) { + t.Errorf("full-online retry: %v", err) + } + if tc.online <= tc.disks/2 { + if got := countPurgeMarkers(t, all, bucket, name, opts.VersionID); got != 0 { + t.Errorf("full-online retry recreated/retained %d marker copies", got) + } + } + // A quorum-confirmed missing retry can leave minority residue. Exercise + // the real dangling-version healer after every disk has returned. + _, _ = er.HealObject(t.Context(), bucket, name, opts.VersionID, madmin.HealOpts{Remove: true}) + if got := countPurgeMarkers(t, all, bucket, name, opts.VersionID); got != 0 { + t.Errorf("marker copies after retry and heal=%d", got) + } + oi, err := z.GetObjectInfo(t.Context(), bucket, name, ObjectOptions{Versioned: true}) + if err != nil || oi.DeleteMarker || oi.VersionID != dataVersion { + t.Errorf("underlying version not visible: %+v error=%v", oi, err) + } + }) + } +} + +func TestDeleteMarkerPurgeMissingKey(t *testing.T) { + for _, pooled := range []bool{false, true} { + t.Run(fmt.Sprintf("pooled_%v", pooled), func(t *testing.T) { + z, er, all, bucket := markerPurgeFixture(t, 4) + _, opts := seedPurgeMarker(t, er, bucket, "marker-only", false) + er.getDisks = func() []StorageAPI { return []StorageAPI{all[0], all[1], nil, nil} } + deleteFn := er.DeleteObject + if pooled { + deleteFn = z.DeleteObject + } + _, _ = deleteFn(t.Context(), bucket, "marker-only", opts) + _, err := deleteFn(t.Context(), bucket, "marker-only", opts) + if !isErrWriteQuorum(err) { + t.Errorf("missing retry with only 2/4 absence votes returned %T %v", err, err) + } + }) + } +} + +func TestDeleteMarkerHealReplicationIdentity(t *testing.T) { + obj, er, all, bucket := markerPurgeFixture(t, 4) + const name = "healed-marker" + _, opts := seedPurgeMarker(t, er, bucket, name, true) + original, err := all[3].ReadVersion(t.Context(), "", bucket, name, opts.VersionID, ReadOptions{}) + if err != nil { + t.Fatal(err) + } + for _, disk := range all[:2] { + if err := disk.DeleteVersion(t.Context(), bucket, name, FileInfo{VersionID: opts.VersionID}, false, DeleteOptions{}); err != nil { + t.Fatal(err) + } + } + if _, err := er.HealObject(t.Context(), bucket, name, opts.VersionID, madmin.HealOpts{}); err != nil { + t.Fatal(err) + } + for i, disk := range all { + fi, err := disk.ReadVersion(t.Context(), "", bucket, name, opts.VersionID, ReadOptions{}) + if err != nil || !maps.Equal(fi.Metadata, original.Metadata) { + t.Errorf("disk %d lost marker metadata: before=%v after=%v error=%v", i, original.Metadata, fi.Metadata, err) + } + } + // Only the healed half remains readable. Read actual disk state before + // invoking scanner replication, rather than constructing an ObjectInfo. + er.getDisks = func() []StorageAPI { return []StorageAPI{all[0], all[1], nil, nil} } + oi, err := er.GetObjectInfo(t.Context(), bucket, name, ObjectOptions{VersionID: opts.VersionID, Versioned: true}) + if !oi.DeleteMarker || !isErrMethodNotAllowed(err) { + t.Fatalf("healed marker unreadable: %+v %v", oi, err) + } + er.getDisks = func() []StorageAPI { return all } + var outbound atomic.Int32 + remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodDelete { + outbound.Add(1) + if r.Header.Get(xhttp.MinIOSourceDeleteMarker) != "true" || r.URL.Query().Get("versionId") != opts.VersionID { + t.Errorf("unexpected outbound request: %s %s", r.Method, r.URL) + } + w.WriteHeader(http.StatusNoContent) + return + } + w.WriteHeader(http.StatusNotFound) + })) + 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) + } + arn := "arn:minio:replication::" + mustGetUUID() + ":bucket" + target := &TargetClient{Client: client, ARN: arn, Bucket: "target"} + oldTargets := globalBucketTargetSys + globalBucketTargetSys = &BucketTargetSys{arnRemotesMap: map[string]arnTarget{arn: {Client: target, lastRefresh: UTCNow()}}, targetsMap: map[string][]madmin.BucketTarget{bucket: {{Arn: arn, TargetBucket: "target"}}}, hc: map[string]epHealth{client.EndpointURL().Host: {Online: true}}} + defer func() { globalBucketTargetSys = oldTargets }() + cfg := configs[0] + cfg.RoleArn = arn + meta, err := globalBucketMetadataSys.Get(bucket) + if err != nil { + t.Fatal(err) + } + meta.replicationConfig = &cfg + globalBucketMetadataSys.Set(bucket, meta) + worker := make(chan ReplicationWorkerOperation, 1) + oldPool := globalReplicationPool + globalReplicationPool = once.NewSingleton[ReplicationPool]() + globalReplicationPool.Set(&ReplicationPool{ctx: t.Context(), objLayer: obj, workers: []chan ReplicationWorkerOperation{worker}, stats: globalReplicationStats.Load(), mrfSaveCh: make(chan MRFReplicateEntry, 1)}) + defer func() { globalReplicationPool = oldPool }() + roi := queueReplicationHeal(t.Context(), bucket, oi, replicationConfig{Config: &cfg, remotes: &madmin.BucketTargets{Targets: []madmin.BucketTarget{{Arn: arn, TargetBucket: "target"}}}}, 0) + select { + case op := <-worker: + d := op.(DeletedObjectReplicationInfo) + result := replicateDeleteToTarget(t.Context(), d, target) + t.Errorf("healed replica scheduled for creation: version=%s existing=%v result=%+v", d.DeleteMarkerVersionID, roi.ExistingObjResync.mustResync(), result) + default: + } + if outbound.Load() != 0 { + t.Errorf("healed replica sent %d new marker creations", outbound.Load()) + } +} diff --git a/cmd/erasure-object.go b/cmd/erasure-object.go index 29c6011cd..0db86e3bc 100644 --- a/cmd/erasure-object.go +++ b/cmd/erasure-object.go @@ -1657,7 +1657,7 @@ func (er erasureObjects) putObject(ctx context.Context, bucket string, object st return fi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended), nil } -func (er erasureObjects) deleteObjectVersion(ctx context.Context, bucket, object string, fi FileInfo, forceDelMarker bool) error { +func (er erasureObjects) deleteObjectVersion(ctx context.Context, bucket, object string, fi FileInfo, forceDelMarker, purge bool) error { disks := er.getDisks() // Assume (N/2 + 1) quorum for Delete() // this is a theoretical assumption such that @@ -1673,7 +1673,13 @@ func (er erasureObjects) deleteObjectVersion(ctx context.Context, bucket, object if disks[index] == nil { return errDiskNotFound } - return disks[index].DeleteVersion(ctx, bucket, object, fi, forceDelMarker, DeleteOptions{}) + err := disks[index].DeleteVersion(ctx, bucket, object, fi, forceDelMarker, DeleteOptions{}) + // Physical removal and reliable already-absent replies are the same + // outcome. Creation and metadata updates must not use this quorum. + if purge && (err == errFileNotFound || err == errFileVersionNotFound) { + return nil + } + return err }, index) } // return errors if any during deletion @@ -2023,6 +2029,11 @@ func (er erasureObjects) DeleteObject(ctx context.Context, bucket, object string if opts.DeleteMarker { versionFound = false } else if !tryDel { + if opts.isVersionPurge() && (isErrObjectNotFound(gerr) || isErrVersionNotFound(gerr)) { + if err := er.checkPurgeAbsent(ctx, bucket, object, opts.VersionID); err != nil { + return objInfo, err + } + } return objInfo, gerr } } @@ -2105,6 +2116,13 @@ func (er erasureObjects) DeleteObject(ctx context.Context, bucket, object string } } + purge := opts.isVersionPurge() + if purge { + // A delete marker describes the version being removed; it is not an + // instruction to create that marker on disks which already lack it. + markDelete, deleteMarker = false, false + } + modTime := opts.MTime if opts.MTime.IsZero() { modTime = UTCNow() @@ -2151,7 +2169,7 @@ func (er erasureObjects) DeleteObject(ctx context.Context, bucket, object string // delete marker. Add delete marker, since we don't have // any version specified explicitly. Or if a particular // version id needs to be replicated. - if err = er.deleteObjectVersion(ctx, bucket, object, fi, opts.DeleteMarker); err != nil { + if err = er.deleteObjectVersion(ctx, bucket, object, fi, opts.DeleteMarker, false); err != nil { return objInfo, toObjectErr(err, bucket, object) } oi := fi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended) @@ -2174,11 +2192,17 @@ func (er erasureObjects) DeleteObject(ctx context.Context, bucket, object string if opts.SkipFreeVersion { dfi.SetSkipTierFreeVersion() } - if err = er.deleteObjectVersion(ctx, bucket, object, dfi, opts.DeleteMarker); err != nil { + if err = er.deleteObjectVersion(ctx, bucket, object, dfi, opts.DeleteMarker, purge); err != nil { return objInfo, toObjectErr(err, bucket, object) } - return dfi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended), nil + oi := dfi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended) + if purge { + // Preserve the DELETE response's identity without reusing it as a disk + // instruction to create a marker. + oi.DeleteMarker = goi.DeleteMarker + } + return oi, nil } // Send the successful but partial upload/delete, however ignore @@ -2467,7 +2491,7 @@ func (er erasureObjects) TransitionObject(ctx context.Context, bucket, object st storageDisks := er.getDisks() - if err = er.deleteObjectVersion(ctx, bucket, object, fi, false); err != nil { + if err = er.deleteObjectVersion(ctx, bucket, object, fi, false, false); err != nil { eventName = event.ObjectTransitionFailed } diff --git a/cmd/erasure-server-pool-consistency.go b/cmd/erasure-server-pool-consistency.go index b4e587516..acac59de5 100644 --- a/cmd/erasure-server-pool-consistency.go +++ b/cmd/erasure-server-pool-consistency.go @@ -316,6 +316,11 @@ func (z *erasureServerPools) retireReplicaCopies(ctx context.Context, bucket, ob func (z *erasureServerPools) deleteObjectReconciled(ctx context.Context, bucket, object string, opts ObjectOptions) (ObjectInfo, error) { copies, err := z.objectPoolInfos(ctx, bucket, object, opts) if err != nil { + if opts.isVersionPurge() && (isErrObjectNotFound(err) || isErrVersionNotFound(err)) { + if quorumErr := z.checkPurgeAbsent(ctx, bucket, object, opts.VersionID); quorumErr != nil { + return ObjectInfo{}, quorumErr + } + } return ObjectInfo{}, err } primary := copies[0] @@ -365,6 +370,21 @@ func (z *erasureServerPools) deleteObjectReconciled(ctx context.Context, bucket, opts.EvalMetadataFn = nil } } + if opts.isVersionPurge() { + // A missing pool may have only read-quorum absence. Verify every pool + // omitted by lookup before mutating the known copies. + for i, pool := range z.serverPools { + found := false + for _, copy := range copies { + found = found || copy.Index == i + } + if !found { + if err := pool.getHashedSet(object).checkPurgeAbsent(ctx, bucket, object, opts.VersionID); err != nil { + return ObjectInfo{}, err + } + } + } + } if opts.VersionID == "" && (opts.Versioned || opts.VersionSuspended) { // A single new delete marker hides the current version. Older versions // remain history and must not be removed by an unqualified DELETE. diff --git a/cmd/erasure-server-pool.go b/cmd/erasure-server-pool.go index 4f9b43cfc..df1d2f22d 100644 --- a/cmd/erasure-server-pool.go +++ b/cmd/erasure-server-pool.go @@ -1239,9 +1239,12 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob return ObjectInfo{}, z.deletePrefix(ctx, bucket, object) } - // Reconcile ordinary addressed-version deletes independently of pool movement. + // Resolve a physical purge by its addressed version even on a replica + // receiver: latest-key routing can select a different pool or miss copies. + // Marker creation and specialized movement/scanner operations retain their + // existing routing. reconcileVersion := opts.VersionID != "" && !opts.DataMovement && - !opts.ReplicationRequest && !opts.Expiration.Expire && !opts.InclFreeVersions + (!opts.ReplicationRequest || opts.isVersionPurge()) && !opts.Expiration.Expire && !opts.InclFreeVersions if !z.SinglePool() && (opts.CheckPrecondFn != nil || reconcileVersion) { return z.deleteObjectReconciled(ctx, bucket, object, opts) } @@ -1254,6 +1257,13 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob if _, ok := err.(InsufficientReadQuorum); ok { return objInfo, InsufficientWriteQuorum{} } + // Lookup can return before the set's purge confirmation. Check here, + // before any callback can change the request into a metadata update. + if opts.isVersionPurge() && (isErrObjectNotFound(err) || isErrVersionNotFound(err)) { + if quorumErr := z.checkPurgeAbsent(ctx, bucket, object, opts.VersionID); quorumErr != nil { + return objInfo, quorumErr + } + } // A conditional (If-Match) delete addressing a specific version treats an // absent key as an absent version. getPoolInfoExistingWithOpts strips // VersionID, so a missing key surfaces ObjectNotFound here even for a @@ -1284,6 +1294,11 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob if verr != nil && (!isErrMethodNotAllowed(verr) || !vi.DeleteMarker) { // Genuine read failure for the addressed version: a missing // version -> VersionNotFound (NoSuchVersion), read-quorum loss, etc. + if opts.isVersionPurge() && (isErrObjectNotFound(verr) || isErrVersionNotFound(verr)) { + if quorumErr := z.checkPurgeAbsent(ctx, bucket, object, opts.VersionID); quorumErr != nil { + return objInfo, quorumErr + } + } return objInfo, verr } // verr is nil for a live version, or MethodNotAllowed with a populated diff --git a/cmd/erasure-version-purge.go b/cmd/erasure-version-purge.go new file mode 100644 index 000000000..9f611ab57 --- /dev/null +++ b/cmd/erasure-version-purge.go @@ -0,0 +1,87 @@ +// Copyright (c) 2026 Feng Ruohang +// +// This file is part of Silo Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "context" + + "github.com/google/uuid" + "github.com/minio/minio/internal/bucket/replication" +) + +// isVersionPurge distinguishes physical removal from marker creation and +// replication-state updates. Call it again after metadata callbacks, which +// can change an ordinary deletion into a pending purge update. +func (o ObjectOptions) isVersionPurge() bool { + if o.VersionID == "" || o.VersionID == nullVersionID || o.DeleteMarker || + o.DeletePrefix || o.DataMovement || o.InclFreeVersions || o.Expiration.Expire || + o.Transition != (TransitionOptions{}) { + return false + } + id, err := uuid.Parse(o.VersionID) + if err != nil || id == uuid.Nil { + return false + } + purge := o.VersionPurgeStatus() + if purge == replication.VersionPurgeComplete { + return true + } + if !purge.Empty() || o.DeleteReplication.VersionPurgeStatusInternal != "" { + return false + } + status := o.DeleteMarkerReplicationStatus() + return (status.Empty() && o.DeleteReplication.ReplicationStatusInternal == "") || + (status == replication.Replica && o.ReplicationRequest) +} + +// confirmVersionAbsent is a read-only proof for an already-missing purge. +// Read quorum (or a synthesized NotFound) cannot acknowledge a write. Do not +// delete here: unreadable minority copies may carry retention we cannot check. +// The caller holds the object lock. This function never heals or enqueues work. +func (er erasureObjects) confirmVersionAbsent(ctx context.Context, bucket, object, versionID string) (allAbsent bool, err error) { + disks := er.getDisks() + _, errs := readAllFileInfo(ctx, disks, "", bucket, object, versionID, false, false) + absent := 0 + for _, err := range errs { + if err == errFileNotFound || err == errFileVersionNotFound { + absent++ + } + } + if absent < len(disks)/2+1 { + return false, InsufficientWriteQuorum{} + } + return absent == len(disks), nil +} + +// checkPurgeAbsent schedules recovery explicitly, outside the read-only proof. +func (er erasureObjects) checkPurgeAbsent(ctx context.Context, bucket, object, versionID string) error { + allAbsent, err := er.confirmVersionAbsent(ctx, bucket, object, versionID) + if !allAbsent { + er.addPartial(bucket, object, versionID) + } + return err +} + +func (z *erasureServerPools) checkPurgeAbsent(ctx context.Context, bucket, object, versionID string) error { + for _, pool := range z.serverPools { + if err := pool.getHashedSet(object).checkPurgeAbsent(ctx, bucket, object, versionID); err != nil { + return err + } + } + return nil +} diff --git a/cmd/replication-delete-mrf_test.go b/cmd/replication-delete-mrf_test.go index 951347665..b0e76daf9 100644 --- a/cmd/replication-delete-mrf_test.go +++ b/cmd/replication-delete-mrf_test.go @@ -82,17 +82,26 @@ type markerRecoveryTarget struct { func replicationTestCapacity(obj ObjectLayer) func() { var restore []func() for _, pool := range obj.(*erasureServerPools).serverPools { + pool.erasureDisksMu.Lock() for _, set := range pool.sets { - original := set.getDisks - disks := append([]StorageAPI(nil), original()...) + original := pool.erasureDisks[set.setIndex] + 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 }) + // GetDisks copies this list under the same mutex. Keep its function + // stable while background IAM scans use the fixture, both when + // installing the adapter and when restoring the original disks. + pool.erasureDisks[set.setIndex] = disks + restore = append(restore, func() { + pool.erasureDisksMu.Lock() + pool.erasureDisks[set.setIndex] = original + pool.erasureDisksMu.Unlock() + }) } + pool.erasureDisksMu.Unlock() } return func() { for _, fn := range restore { @@ -101,6 +110,34 @@ func replicationTestCapacity(obj ObjectLayer) func() { } } +func TestReplicationCapacityConcurrentIAM(t *testing.T) { + z, _ := consistencyPools(t) + if _, _, err := initAPIHandlerTest(t.Context(), z, nil, MakeBucketOptions{}); err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) + defer cancel() + iam := globalIAMSys + started, finished := make(chan struct{}), make(chan error, 1) + go func() { + close(started) + for range 20 { + if err := iam.Load(ctx, false); err != nil { + finished <- err + return + } + } + finished <- nil + }() + <-started + for range 5000 { + replicationTestCapacity(z)() + } + if err := <-finished; err != nil { + t.Fatal(err) + } +} + 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() diff --git a/cmd/xl-storage-delete-marker-metadata_test.go b/cmd/xl-storage-delete-marker-metadata_test.go new file mode 100644 index 000000000..f455ae4b7 --- /dev/null +++ b/cmd/xl-storage-delete-marker-metadata_test.go @@ -0,0 +1,148 @@ +// Copyright (c) 2026 Feng Ruohang +// +// This file is part of Silo Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "bytes" + "io" + "maps" + "net/http" + "testing" + "time" + + "github.com/minio/minio/internal/bucket/replication" +) + +func TestDeleteMarkerReadDataControls(t *testing.T) { + obj, er, disks, bucket := markerPurgeFixture(t, 4) + for _, body := range []string{"", "inline-data-control"} { + name := "data-" + mustGetUUID() + oi, err := er.PutObject(t.Context(), bucket, name, mustGetPutObjReader(t, bytes.NewBufferString(body), int64(len(body)), "", ""), ObjectOptions{Versioned: true}) + if err != nil { + t.Fatal(err) + } + for _, disk := range disks { + fi, err := disk.ReadVersion(t.Context(), "", bucket, name, oi.VersionID, ReadOptions{ReadData: true}) + if err != nil || fi.Deleted || fi.Size != int64(len(body)) || !fi.InlineData() || (body != "" && len(fi.Data) == 0) { + t.Fatalf("data read: size=%d inline=%v bytes=%d error=%v", fi.Size, fi.InlineData(), len(fi.Data), err) + } + // Legacy inline data without its annotation still gets the original + // ReadData behavior; the new marker guard must not affect it. + delete(fi.Metadata, ReservedMetadataPrefixLower+"inline-data") + if err := disk.WriteMetadata(t.Context(), "", bucket, name, fi); err != nil { + t.Fatal(err) + } + fi, err = disk.ReadVersion(t.Context(), "", bucket, name, oi.VersionID, ReadOptions{ReadData: true}) + if err != nil || !fi.InlineData() { + t.Fatalf("legacy inline read: %+v error=%v", fi, err) + } + } + reader, err := obj.GetObjectNInfo(t.Context(), bucket, name, nil, http.Header{}, ObjectOptions{VersionID: oi.VersionID}) + if err != nil { + t.Fatal(err) + } + got, err := io.ReadAll(reader) + reader.Close() + if err != nil || string(got) != body { + t.Fatalf("payload changed: %q error=%v", got, err) + } + } + _, opts := seedPurgeMarker(t, er, bucket, "stored-marker-key", false) + for _, disk := range disks { + fi, err := disk.ReadVersion(t.Context(), "", bucket, "stored-marker-key", opts.VersionID, ReadOptions{}) + if err != nil { + t.Fatal(err) + } + // Preserve even an existing unusual marker key; suppress only the + // manufacture of a new inline annotation by the data-read branch. + fi.Metadata[ReservedMetadataPrefixLower+"inline-data"] = "original" + if err := disk.WriteMetadata(t.Context(), "", bucket, "stored-marker-key", fi); err != nil { + t.Fatal(err) + } + fi, err = disk.ReadVersion(t.Context(), "", bucket, "stored-marker-key", opts.VersionID, ReadOptions{ReadData: true}) + if err != nil || fi.Metadata[ReservedMetadataPrefixLower+"inline-data"] != "original" { + t.Errorf("existing marker key changed: metadata=%v error=%v", fi.Metadata, err) + } + } +} + +func TestDeleteMarkerMetadataRoundTrip(t *testing.T) { + stamp := time.Date(2026, 9, 1, 12, 1, 2, 345, time.UTC) + for _, free := range []bool{false, true} { + t.Run(map[bool]string{false: "replication", true: "free-version"}[free], func(t *testing.T) { + metadata := map[string]string{ + ReservedMetadataPrefixLower + ReplicaStatus: "REPLICA", + ReservedMetadataPrefixLower + ReplicaTimestamp: stamp.Format(time.RFC3339Nano), + ReservedMetadataPrefixLower + ReplicationStatus: "arn1=COMPLETED;arn2=FAILED;", + ReservedMetadataPrefixLower + ReplicationTimestamp: stamp.Add(-time.Hour).Format(time.RFC3339Nano), + VersionPurgeStatusKey: "arn1=PENDING;arn2=FAILED;", + targetResetHeader("arn1"): "original;reset1", + targetResetHeader("arn2"): "original;reset2", + ReservedMetadataPrefixLower + "unknown": "", + } + if free { + metadata[ReservedMetadataPrefixLower+freeVersion] = "" + metadata[metaTierName] = "tier" + metadata[metaTierObjName] = "remote-object" + metadata[metaTierVersionID] = "remote-version" + } + fi := FileInfo{VersionID: mustGetUUID(), Deleted: true, ModTime: stamp, Metadata: maps.Clone(metadata)} + // Deliberately conflicting parsed state must never replace stored + // values or add a key missing from the raw metadata. + fi.ReplicationState = ReplicationState{ReplicaStatus: replication.Failed, ReplicaTimeStamp: stamp.Add(time.Hour), ResetStatusesMap: map[string]string{"new-arn": "invented"}} + fi.SetHealing() + fi.SetDataMov() + fi.SetTierFreeVersionID(mustGetUUID()) + fi.SetTierFreeVersion() + fi.SetSkipTierFreeVersion() + xl := xlMetaV2{} + if err := xl.AddVersion(fi); err != nil { + t.Fatal(err) + } + stored, err := xl.getIdx(0) + if err != nil { + t.Fatal(err) + } + got := make(map[string]string) + for k, v := range stored.DeleteMarker.MetaSys { + got[k] = string(v) + } + if !maps.Equal(got, metadata) { + t.Errorf("stored metadata=%v want=%v", got, metadata) + } + decoded, err := xl.ToFileInfo("bucket", "marker", fi.VersionID, true, true) + if err != nil || !decoded.ModTime.Equal(fi.ModTime) || decoded.TierFreeVersion() != free { + t.Errorf("round trip identity: %+v error=%v", decoded, err) + } + if decoded.ReplicationState.ReplicaStatus != replication.Replica || decoded.ReplicationState.PurgeTargets["arn2"] != replication.VersionPurgeFailed || decoded.ReplicationState.Targets["arn1"] != replication.Completed { + t.Errorf("lost parsed replication state: %+v", decoded.ReplicationState) + } + }) + } +} + +func TestDeleteMarkerCreationMetadata(t *testing.T) { + _, er, disks, bucket := markerPurgeFixture(t, 4) + _, opts := seedPurgeMarker(t, er, bucket, "empty-key", false) + for i, disk := range disks { + fi, err := disk.ReadVersion(t.Context(), "", bucket, "empty-key", opts.VersionID, ReadOptions{}) + if err != nil || fi.Metadata[ReservedMetadataPrefixLower+ReplicaStatus] != "REPLICA" || fi.Metadata[ReservedMetadataPrefixLower+ReplicaTimestamp] != opts.DeleteReplication.ReplicaTimeStamp.Format(time.RFC3339Nano) { + t.Errorf("disk %d lost new marker replica identity: metadata=%v error=%v", i, fi.Metadata, err) + } + } +} diff --git a/cmd/xl-storage-delete-marker.go b/cmd/xl-storage-delete-marker.go new file mode 100644 index 000000000..a5f0a2a4a --- /dev/null +++ b/cmd/xl-storage-delete-marker.go @@ -0,0 +1,67 @@ +// Copyright (c) 2026 Feng Ruohang +// +// This file is part of Silo Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "strings" + "time" +) + +// deleteMarkerMetadata preserves the complete stored marker state during +// healing. The parsed ReplicationState is lossy; even absent raw keys must not +// be reconstructed from it. Only creation writers without stored metadata +// use typed state, and never invent timestamps. +func deleteMarkerMetadata(fi FileInfo) map[string][]byte { + meta := make(map[string][]byte, len(fi.Metadata)) + for k, v := range fi.Metadata { + switch k { + case xMinIOHealing, xMinIODataMov, + ReservedMetadataPrefixLower + tierFVID, + ReservedMetadataPrefixLower + tierFVMarker, + ReservedMetadataPrefixLower + tierSkipFVID: + continue + } + meta[k] = []byte(v) + } + if len(meta) != 0 || fi.Healing() || fi.DataMov() { + return meta + } + rs := fi.ReplicationState + if !rs.ReplicaStatus.Empty() { + meta[ReservedMetadataPrefixLower+ReplicaStatus] = []byte(rs.ReplicaStatus) + if !rs.ReplicaTimeStamp.IsZero() { + meta[ReservedMetadataPrefixLower+ReplicaTimestamp] = []byte(rs.ReplicaTimeStamp.UTC().Format(time.RFC3339Nano)) + } + } + if rs.ReplicationStatusInternal != "" { + meta[ReservedMetadataPrefixLower+ReplicationStatus] = []byte(rs.ReplicationStatusInternal) + if !rs.ReplicationTimeStamp.IsZero() { + meta[ReservedMetadataPrefixLower+ReplicationTimestamp] = []byte(rs.ReplicationTimeStamp.UTC().Format(time.RFC3339Nano)) + } + } + if rs.VersionPurgeStatusInternal != "" { + meta[VersionPurgeStatusKey] = []byte(rs.VersionPurgeStatusInternal) + } + for k, v := range rs.ResetStatusesMap { + if !strings.HasPrefix(k, ReservedMetadataPrefixLower+ReplicationReset) { + k = targetResetHeader(k) + } + meta[k] = []byte(v) + } + return meta +} diff --git a/cmd/xl-storage-format-v2.go b/cmd/xl-storage-format-v2.go index fd6ac1642..3ce72da24 100644 --- a/cmd/xl-storage-format-v2.go +++ b/cmd/xl-storage-format-v2.go @@ -1630,7 +1630,7 @@ func (x *xlMetaV2) AddVersion(fi FileInfo) error { ventry.DeleteMarker = &xlMetaV2DeleteMarker{ VersionID: uv, ModTime: fi.ModTime.UnixNano(), - MetaSys: make(map[string][]byte), + MetaSys: deleteMarkerMetadata(fi), } } else { ventry.Type = ObjectType diff --git a/cmd/xl-storage.go b/cmd/xl-storage.go index f3a6231c3..92d8c8dec 100644 --- a/cmd/xl-storage.go +++ b/cmd/xl-storage.go @@ -1719,7 +1719,9 @@ func (s *xlStorage) ReadVersion(ctx context.Context, origvolume, volume, path, v defer metaDataPoolPut(buf) } - if readData { + // Delete markers have no payload. In particular, do not manufacture an + // inline-data metadata key when healing their zero-size FileInfo. + if readData && !fi.Deleted { if len(fi.Data) > 0 || fi.Size == 0 { if fi.InlineData() { // If written with header we are fine.