From 1cf529ce8a4cd492c65ffac7fb832165ef9786b2 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Tue, 15 Sep 2026 06:54:29 +0800 Subject: [PATCH] fix(storage): reconcile ordinary version DELETE across pools Delete every copy of an explicitly addressed UUID, null version or delete marker under the pool lock. Preserve retention and replication callbacks, report unreadable pools and cleanup failures, and keep movement, incoming replication, expiration and free-version cleanup on their existing paths. Retain the separately developed general DELETE repair and replace its access-mover-only coverage with a real interrupted rebalance copy followed by HTTP deletion. Cover unqualified directory-marker DELETE and document the existing pool-order-dependent 503 behavior that this makes consistent. Signed-off-by: Feng Ruohang --- CHANGELOG.md | 9 + cmd/erasure-server-pool-consistency.go | 32 +- cmd/erasure-server-pool-consistency_test.go | 607 +++++++++++++++++- cmd/erasure-server-pool.go | 7 +- .../lifecycle/access-tiering-removal.md | 7 + 5 files changed, 648 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2d26b4e2a..8152b0186 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,15 @@ and [complete commit range](https://github.com/pgsty/silo/compare/RELEASE.2026-0 ### Object storage and replication +- Reconcile ordinary single-object version DELETE across all pools, including + null versions, delete markers and unqualified directory-marker DELETE. This + prevents movement leftovers from surviving a successful response. Unreadable + pools now consistently return 503 instead of depending on pool traversal + order; this extends the existing failure surface. Retry after recovery. + Cleanup failures also return an error. Batch deletion already fans out across + pools; replication and scanner cleanup keep their existing contracts. See + [scope and limitations](docs/bucket/lifecycle/access-tiering-removal.md#version-deletion-scope). + - Remove the opt-in GET-frequency pool-tiering feature from PR #60, including its tracker, mover, scanner hooks, configuration, XML actions and metrics. Accept and ignore retired configuration/XML and preserve ordinary statistics diff --git a/cmd/erasure-server-pool-consistency.go b/cmd/erasure-server-pool-consistency.go index 415f1fa19..99c14a223 100644 --- a/cmd/erasure-server-pool-consistency.go +++ b/cmd/erasure-server-pool-consistency.go @@ -287,33 +287,47 @@ func (z *erasureServerPools) retireReplicaCopies(ctx context.Context, bucket, ob return nil } -// deleteObjectConditional evaluates the condition once against the logical +// deleteObjectReconciled evaluates any condition once against the logical // version, then removes all its copies under the same lock as pooled writers. -func (z *erasureServerPools) deleteObjectConditional(ctx context.Context, bucket, object string, opts ObjectOptions) (ObjectInfo, error) { +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 { return ObjectInfo{}, err } primary := copies[0] - if opts.CheckPrecondFn(primary.ObjInfo) { + if opts.CheckPrecondFn != nil && opts.CheckPrecondFn(primary.ObjInfo) { return ObjectInfo{}, PreConditionFailed{} } opts.CheckPrecondFn = nil opts.NoLock = true if opts.EvalRetentionBypassFn != nil || opts.EvalMetadataFn != nil { - versions, err := z.metadataPoolInfos(ctx, bucket, object, opts) - if err != nil { - return ObjectInfo{}, err + logical := primary.ObjInfo + var gerr error + if logical.DeleteMarker { + // Markers can be deleted by version ID. Match the set layer's + // callback inputs instead of rejecting them as metadata updates. + gerr = toObjectErr(errMethodNotAllowed, bucket, object) + if opts.VersionID == "" || opts.DeleteMarker { + gerr = toObjectErr(errFileNotFound, bucket, object) + } + } else { + versions, err := z.metadataPoolInfos(ctx, bucket, object, opts) + if err != nil { + return ObjectInfo{}, err + } + logical = mergedPoolObjectInfo(versions) } - logical := mergedPoolObjectInfo(versions) + // Keep the retention gate first. These callbacks independently evaluate + // the logical version and run once before any deletion; the handler only + // sweeps metadata's transition state after a successful delete. if opts.EvalRetentionBypassFn != nil { - if err := opts.EvalRetentionBypassFn(logical, nil); err != nil { + if err := opts.EvalRetentionBypassFn(logical, gerr); err != nil { return ObjectInfo{}, err } opts.EvalRetentionBypassFn = nil } if opts.EvalMetadataFn != nil { - decision, err := opts.EvalMetadataFn(&logical, nil) + decision, err := opts.EvalMetadataFn(&logical, gerr) if err != nil { return ObjectInfo{}, err } diff --git a/cmd/erasure-server-pool-consistency_test.go b/cmd/erasure-server-pool-consistency_test.go index f85d0a2d1..227b1436f 100644 --- a/cmd/erasure-server-pool-consistency_test.go +++ b/cmd/erasure-server-pool-consistency_test.go @@ -24,12 +24,16 @@ import ( "fmt" "io" "maps" + "net/http" + "net/http/httptest" "strings" "sync" + "sync/atomic" "testing" "time" madmin "github.com/minio/madmin-go/v3" + "github.com/minio/minio/internal/bucket/replication" xhttp "github.com/minio/minio/internal/http" ) @@ -62,6 +66,477 @@ func putConsistencyObject(t *testing.T, z *erasureServerPools, bucket, object st return oi } +func consistencyRequest(t *testing.T, router http.Handler, method, bucket, object, version string) *httptest.ResponseRecorder { + t.Helper() + url := getGetObjectURL("", bucket, object) + if version != "" { + url += "?versionId=" + version + } + req, err := newTestSignedRequestV4(method, url, 0, nil, globalActiveCred.AccessKey, globalActiveCred.SecretKey, nil) + if err != nil { + t.Fatal(err) + } + rec := httptest.NewRecorder() + router.ServeHTTP(rec, req) + return rec +} + +// Exercise the handler's real replication/retention callbacks, including the +// MethodNotAllowed metadata returned for an explicitly addressed delete marker. +func TestPoolsDeleteVersionAPI(t *testing.T) { + z, _ := consistencyPools(t) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + bucket, router, err := initAPIHandlerTest(ctx, z, nil, MakeBucketOptions{VersioningEnabled: true}) + if err != nil { + t.Fatal(err) + } + for _, kind := range []string{"uuid", "null", "marker"} { + for primary := range 2 { + t.Run(fmt.Sprintf("%s/primary=%d", kind, primary), func(t *testing.T) { + object := fmt.Sprintf("%s-%d", kind, primary) + opts := ObjectOptions{Versioned: kind != "null", MTime: UTCNow().Add(-time.Hour)} + var addressed ObjectInfo + if kind == "marker" { + opts.VersionID, opts.DeleteMarker = mustGetUUID(), true + addressed, err = z.serverPools[primary].DeleteObject(ctx, bucket, object, opts) + if err != nil { + t.Fatal(err) + } + } else { + addressed = putConsistencyObject(t, z, bucket, object, primary, "payload", opts) + } + version := addressed.VersionID + if version == "" { + version = nullVersionID + } + opts.VersionID = version + opts.MTime = addressed.ModTime.Add(-time.Minute) + if kind == "marker" { + opts.DeleteMarker = true + if _, err := z.serverPools[1-primary].DeleteObject(ctx, bucket, object, opts); err != nil { + t.Fatal(err) + } + } else { + putConsistencyObject(t, z, bucket, object, 1-primary, "payload", opts) + } + latest := putConsistencyObject(t, z, bucket, object, 1-primary, "keep-latest", ObjectOptions{Versioned: true}) + rec := consistencyRequest(t, router, http.MethodDelete, bucket, object, version) + if rec.Code != http.StatusNoContent { + t.Fatalf("DELETE: %d %s", rec.Code, rec.Body.String()) + } + if kind == "marker" && (strings.Join(rec.Header()[xhttp.AmzDeleteMarker], "") != "true" || strings.Join(rec.Header()[xhttp.AmzVersionID], "") != version) { + t.Errorf("lost deleted marker response headers: %v", rec.Header()) + } + for i, pool := range z.serverPools { + if _, err := pool.GetObjectInfo(ctx, bucket, object, ObjectOptions{VersionID: version}); !isErrVersionNotFound(err) { + t.Errorf("pool %d retained addressed version: %v", i, err) + } + } + for _, method := range []string{http.MethodGet, http.MethodHead} { + if rec := consistencyRequest(t, router, method, bucket, object, version); rec.Code != http.StatusNotFound { + t.Errorf("%s after DELETE: %d %s", method, rec.Code, rec.Body.String()) + } + } + if got, err := z.GetObjectInfo(ctx, bucket, object, ObjectOptions{}); err != nil || got.VersionID != latest.VersionID { + t.Errorf("DELETE changed another version: %+v, %v", got, err) + } + if rec := consistencyRequest(t, router, http.MethodDelete, bucket, object, version); rec.Code != http.StatusNoContent { + t.Errorf("idempotent retry: %d %s", rec.Code, rec.Body.String()) + } + }) + } + } + if rec := consistencyRequest(t, router, http.MethodDelete, bucket, "absent-key", mustGetUUID()); rec.Code != http.StatusNoContent { + t.Errorf("absent key: %d %s", rec.Code, rec.Body.String()) + } +} + +type consistencyReadFaultDisk struct { + StorageAPI + bucket, object string +} + +func (d consistencyReadFaultDisk) ReadVersion(ctx context.Context, origvolume, volume, path, version string, opts ReadOptions) (FileInfo, error) { + if volume == d.bucket && path == d.object { + return FileInfo{}, errDiskNotFound + } + return d.StorageAPI.ReadVersion(ctx, origvolume, volume, path, version, opts) +} + +func (d consistencyReadFaultDisk) ReadXL(ctx context.Context, volume, path string, readData bool) (RawFileInfo, error) { + if volume == d.bucket && path == d.object { + return RawFileInfo{}, errDiskNotFound + } + return d.StorageAPI.ReadXL(ctx, volume, path, readData) +} + +func TestPoolsDeleteVersionUnreadablePool(t *testing.T) { + z, _ := consistencyPools(t) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + bucket, router, err := initAPIHandlerTest(ctx, z, nil, MakeBucketOptions{VersioningEnabled: true}) + if err != nil { + t.Fatal(err) + } + for _, holding := range []bool{false, true} { + for failed := range 2 { + t.Run(fmt.Sprintf("holding=%t/pool=%d", holding, failed), func(t *testing.T) { + object := fmt.Sprintf("unreadable-%t-%d", holding, failed) + oi := putConsistencyObject(t, z, bucket, object, 1-failed, "payload", ObjectOptions{Versioned: true}) + if holding { + putConsistencyObject(t, z, bucket, object, failed, "payload", ObjectOptions{Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime}) + } + set := z.serverPools[failed].getHashedSet(object) + getDisks := set.getDisks + faulty := append([]StorageAPI(nil), getDisks()...) + for i := range faulty { + faulty[i] = consistencyReadFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object} + } + set.getDisks = func() []StorageAPI { return faulty } + defer func() { set.getDisks = getDisks }() + if rec := consistencyRequest(t, router, http.MethodDelete, bucket, object, oi.VersionID); rec.Code != http.StatusServiceUnavailable { + t.Errorf("unreadable pool DELETE: %d %s", rec.Code, rec.Body.String()) + } + if _, err := z.serverPools[1-failed].GetObjectInfo(ctx, bucket, object, ObjectOptions{VersionID: oi.VersionID}); err != nil { + t.Errorf("DELETE lost the readable copy: %v", err) + } + set.getDisks = getDisks + if rec := consistencyRequest(t, router, http.MethodDelete, bucket, object, oi.VersionID); rec.Code != http.StatusNoContent { + t.Errorf("recovered pool retry: %d %s", rec.Code, rec.Body.String()) + } + for i, pool := range z.serverPools { + if _, err := pool.GetObjectInfo(ctx, bucket, object, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) { + t.Errorf("pool %d retained version after retry: %v", i, err) + } + } + }) + } + } +} + +type consistencyDeleteCountDisk struct { + StorageAPI + bucket, object string + deletes *atomic.Int32 +} + +func (d consistencyDeleteCountDisk) DeleteVersion(ctx context.Context, volume, path string, fi FileInfo, force bool, opts DeleteOptions) error { + if volume == d.bucket && path == d.object { + d.deletes.Add(1) + } + return d.StorageAPI.DeleteVersion(ctx, volume, path, fi, force, opts) +} + +func TestPoolsDeleteVersionSingleCopy(t *testing.T) { + z, bucket := consistencyPools(t) + for primary := range 2 { + t.Run(fmt.Sprintf("pool=%d", primary), func(t *testing.T) { + object := fmt.Sprintf("single-copy-%d", primary) + oi := putConsistencyObject(t, z, bucket, object, primary, "payload", ObjectOptions{Versioned: true}) + var deletes [2]atomic.Int32 + for i, pool := range z.serverPools { + set := pool.getHashedSet(object) + getDisks := set.getDisks + disks := append([]StorageAPI(nil), getDisks()...) + for j := range disks { + disks[j] = consistencyDeleteCountDisk{StorageAPI: disks[j], bucket: bucket, object: object, deletes: &deletes[i]} + } + set.getDisks = func() []StorageAPI { return disks } + defer func() { set.getDisks = getDisks }() + } + metadata, retention := 0, 0 + _, err := z.DeleteObject(t.Context(), bucket, object, ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, + EvalMetadataFn: func(current *ObjectInfo, err error) (ReplicateDecision, error) { + metadata++ + if err != nil || current.VersionID != oi.VersionID { + t.Errorf("metadata callback: %+v, %v", current, err) + } + return ReplicateDecision{}, nil + }, + EvalRetentionBypassFn: func(current ObjectInfo, err error) error { + retention++ + return err + }, + }) + if err != nil || metadata != 1 || retention != 1 { + t.Fatalf("DELETE: %v, metadata=%d retention=%d", err, metadata, retention) + } + if deletes[primary].Load() != 16 || deletes[1-primary].Load() != 0 { + t.Errorf("physical deletes per pool = %d, %d; want once per holding disk only", deletes[0].Load(), deletes[1].Load()) + } + }) + } +} + +func TestPoolsDeleteVersionReplicationPurge(t *testing.T) { + z, bucket := consistencyPools(t) + for _, marker := range []bool{false, true} { + t.Run(fmt.Sprintf("marker=%t", marker), func(t *testing.T) { + object := fmt.Sprintf("local-replication-purge-%t", marker) + version := mustGetUUID() + for _, pool := range z.serverPools { + opts := ObjectOptions{Versioned: true, VersionID: version, DeleteMarker: marker} + var err error + if marker { + _, err = pool.DeleteObject(t.Context(), bucket, object, opts) + } else { + _, err = pool.PutObject(t.Context(), bucket, object, mustGetPutObjReader(t, strings.NewReader("payload"), 7, "", ""), opts) + } + if err != nil { + t.Fatal(err) + } + } + dsc := ReplicateDecision{} + dsc.Set(newReplicateTargetDecision("arn1", true, false)) + opts := ObjectOptions{Versioned: true, VersionID: version} + for attempt := range 2 { + calls := 0 + opts.EvalMetadataFn = func(current *ObjectInfo, gerr error) (ReplicateDecision, error) { + calls++ + wantMethodNotAllowed := marker || attempt > 0 + if (wantMethodNotAllowed && !isErrMethodNotAllowed(gerr)) || (!wantMethodNotAllowed && gerr != nil) { + t.Errorf("pending purge callback lost set-layer read error: %v", gerr) + } + return dsc, nil + } + got, err := z.DeleteObject(t.Context(), bucket, object, opts) + if err != nil || calls != 1 || got.VersionPurgeStatus != replication.VersionPurgePending || got.replicationDecision != dsc.String() { + t.Fatalf("replication scheduling result: %+v, %v, calls=%d", got, err, calls) + } + for i, pool := range z.serverPools { + got, err := pool.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: version}) + if (err != nil && (!got.DeleteMarker || !isErrMethodNotAllowed(err))) || got.VersionPurgeStatus != replication.VersionPurgePending { + t.Errorf("pool %d lost pending purge: %+v, %v", i, got, err) + } + } + } + // The local worker completes the purge without ReplicationRequest. + // Both pending copies must be removed, including a pending marker purge. + opts.EvalMetadataFn = nil + opts.DeleteReplication = ReplicationState{ + VersionPurgeStatusInternal: "arn1=COMPLETED;", + PurgeTargets: map[string]VersionPurgeStatusType{"arn1": replication.VersionPurgeComplete}, + } + if _, err := z.DeleteObject(t.Context(), bucket, object, opts); err != nil { + t.Fatal(err) + } + for i, pool := range z.serverPools { + if _, err := pool.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: version}); !isErrVersionNotFound(err) { + t.Errorf("pool %d retained completed replication purge: %v", i, err) + } + } + }) + } +} + +func TestPoolsDeleteVersionCleanupFailure(t *testing.T) { + z, bucket := consistencyPools(t) + for primary := range 2 { + t.Run(fmt.Sprintf("primary=%d", primary), func(t *testing.T) { + object := fmt.Sprintf("cleanup-failure-%d", primary) + oi := putConsistencyObject(t, z, bucket, object, primary, "payload", ObjectOptions{Versioned: true}) + putConsistencyObject(t, z, bucket, object, 1-primary, "payload", ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime.Add(-time.Minute), + }) + set := z.serverPools[1-primary].getHashedSet(object) + getDisks := set.getDisks + faulty := append([]StorageAPI(nil), getDisks()...) + for i := range faulty { + faulty[i] = consistencyDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object, version: oi.VersionID} + } + set.getDisks = func() []StorageAPI { return faulty } + defer func() { set.getDisks = getDisks }() + opts := ObjectOptions{Versioned: true, VersionID: oi.VersionID} + if _, err := z.DeleteObject(t.Context(), bucket, object, opts); err == nil { + t.Error("acknowledged DELETE despite failed secondary cleanup") + } + if got, err := z.serverPools[primary].GetObjectInfo(t.Context(), bucket, object, opts); err != nil || got.ETag != oi.ETag { + t.Errorf("lost authoritative copy after secondary failure: %+v, %v", got, err) + } + set.getDisks = getDisks + if _, err := z.DeleteObject(t.Context(), bucket, object, opts); err != nil { + t.Fatal(err) + } + for i, pool := range z.serverPools { + if _, err := pool.GetObjectInfo(t.Context(), bucket, object, opts); !isErrVersionNotFound(err) { + t.Errorf("pool %d retained version after retry: %v", i, err) + } + } + }) + } +} + +func TestPoolsDeleteVersionCallbacks(t *testing.T) { + z, bucket := consistencyPools(t) + for _, marker := range []bool{false, true} { + t.Run(fmt.Sprintf("marker=%t", marker), func(t *testing.T) { + object := fmt.Sprintf("callbacks-%t", marker) + old, recent := "2026-09-09T09:00:00Z", "2026-09-09T10:00:00Z" + oi := putConsistencyObject(t, z, bucket, object, 0, "payload", ObjectOptions{ + Versioned: true, UserDefined: poolLockMetadata("GOVERNANCE", "OFF", recent, old), + }) + putConsistencyObject(t, z, bucket, object, 1, "payload", ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime, + UserDefined: poolLockMetadata("", "ON", old, recent), + }) + if marker { + markerOpts := ObjectOptions{Versioned: true, VersionID: mustGetUUID(), DeleteMarker: true, MTime: UTCNow()} + for _, pool := range z.serverPools { + var err error + oi, err = pool.DeleteObject(t.Context(), bucket, object, markerOpts) + if err != nil { + t.Fatal(err) + } + } + } + check := func(current ObjectInfo, err error) { + t.Helper() + if current.VersionID != oi.VersionID || current.DeleteMarker != marker || (marker && !isErrMethodNotAllowed(err)) || (!marker && err != nil) { + t.Errorf("wrong callback version or read error: %+v, %v", current, err) + } + if !marker { + state := storedObjectLockState(current.UserDefined) + if state.mode != "GOVERNANCE" || state.legalHold != "ON" { + t.Errorf("callback did not reconcile independently ordered lock state: %+v", state) + } + } + } + for _, reject := range []string{"retention", "metadata", "none"} { + metadata, retention := 0, 0 + denied := errors.New("callback denied deletion") + opts := ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, + EvalRetentionBypassFn: func(current ObjectInfo, err error) error { + retention++ + check(current, err) + if reject == "retention" { + return denied + } + return nil + }, + EvalMetadataFn: func(current *ObjectInfo, err error) (ReplicateDecision, error) { + metadata++ + check(*current, err) + if reject == "metadata" { + return ReplicateDecision{}, denied + } + return ReplicateDecision{}, nil + }, + } + _, err := z.DeleteObject(t.Context(), bucket, object, opts) + if reject == "none" { + if err != nil || metadata != 1 || retention != 1 { + t.Fatalf("DELETE: %v, metadata=%d retention=%d", err, metadata, retention) + } + } else if !errors.Is(err, denied) { + t.Fatalf("%s rejection was lost: %v", reject, err) + } + for i, pool := range z.serverPools { + current, err := pool.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: oi.VersionID}) + if reject == "none" { + if !isErrVersionNotFound(err) { + t.Errorf("pool %d retained deleted version: %v", i, err) + } + } else if current.VersionID != oi.VersionID || (err != nil && (!marker || !isErrMethodNotAllowed(err))) { + t.Errorf("pool %d changed despite %s rejection: %+v, %v", i, reject, current, err) + } + } + } + }) + } +} + +func TestPoolsDeleteVersionSpecialCalls(t *testing.T) { + z, bucket := consistencyPools(t) + t.Run("incoming-marker-replication", func(t *testing.T) { + const object = "incoming-marker" + opts := ObjectOptions{Versioned: true, VersionID: mustGetUUID(), DeleteMarker: true, ReplicationRequest: true} + opts.SetReplicaStatus(replication.Replica) + if _, err := z.DeleteObject(t.Context(), bucket, object, opts); err != nil { + t.Fatalf("replication could not create an absent marker: %v", err) + } + if got, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: opts.VersionID}); !isErrMethodNotAllowed(err) || !got.DeleteMarker { + t.Fatalf("replicated marker missing: %+v, %v", got, err) + } + }) + for _, rebalance := range []bool{false, true} { + t.Run(fmt.Sprintf("movement/rebalance=%t", rebalance), func(t *testing.T) { + object := fmt.Sprintf("movement-%t", rebalance) + putConsistencyObject(t, z, bucket, object, 1, "destination", ObjectOptions{Versioned: true}) + opts := ObjectOptions{Versioned: true, VersionID: mustGetUUID(), DeleteMarker: true, MTime: UTCNow()} + if _, err := z.serverPools[0].DeleteObject(t.Context(), bucket, object, opts); err != nil { + t.Fatal(err) + } + if rebalance { + z.rebalMu.Lock() + z.rebalMeta = &rebalanceMeta{PoolStats: []*rebalanceStats{{Participating: true, Info: rebalanceInfo{Status: rebalStarted}}, {}}} + z.rebalMu.Unlock() + defer func() { z.rebalMu.Lock(); z.rebalMeta = nil; z.rebalMu.Unlock() }() + } else { + z.poolMetaMutex.Lock() + z.poolMeta.Pools[0].Decommission = &PoolDecommissionInfo{} + z.poolMetaMutex.Unlock() + defer func() { z.poolMetaMutex.Lock(); z.poolMeta.Pools[0].Decommission = nil; z.poolMetaMutex.Unlock() }() + } + // These are the marker-copy options used by rebalance/decommission; + // Source cleanup belongs to the mover. + opts.DataMovement, opts.SrcPoolIdx = true, 0 + opts.SkipRebalancing, opts.SkipDecommissioned = rebalance, !rebalance + if _, err := z.DeleteObject(t.Context(), bucket, object, opts); err != nil { + t.Fatal(err) + } + for i, pool := range z.serverPools { + if got, err := pool.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: opts.VersionID}); !isErrMethodNotAllowed(err) || !got.DeleteMarker { + t.Errorf("movement lost marker in pool %d: %+v, %v", i, got, err) + } + } + }) + } + for _, expiration := range []bool{false, true} { + t.Run(fmt.Sprintf("scanner/expiration=%t", expiration), func(t *testing.T) { + object := fmt.Sprintf("scanner-%t", expiration) + oi := putConsistencyObject(t, z, bucket, object, 0, "payload", ObjectOptions{Versioned: true}) + set := z.serverPools[1].getHashedSet(object) + getDisks := set.getDisks + faulty := append([]StorageAPI(nil), getDisks()...) + for i := range faulty { + faulty[i] = consistencyReadFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object} + } + set.getDisks = func() []StorageAPI { return faulty } + defer func() { set.getDisks = getDisks }() + opts := ObjectOptions{Versioned: true, VersionID: oi.VersionID, InclFreeVersions: !expiration, Expiration: ExpirationOptions{Expire: expiration}} + _, err := z.DeleteObject(t.Context(), bucket, object, opts) + if expiration { + // No lifecycle rule authorizes expiration. Preserve its existing + // version-not-found result, even with an unrelated unreadable pool. + if !isErrVersionNotFound(err) { + t.Errorf("expiration changed its scanner contract: %v", err) + } + } else if err != nil { + t.Errorf("free-version cleanup was forced through all-pool resolution: %v", err) + } + }) + } +} + +func TestPoolsDeleteUnversionedFanout(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "unversioned-fanout" + for i := range z.serverPools { + putConsistencyObject(t, z, bucket, object, i, "payload", ObjectOptions{}) + } + if _, err := z.DeleteObject(t.Context(), bucket, object, ObjectOptions{}); err != nil { + t.Fatal(err) + } + for i, pool := range z.serverPools { + if _, err := pool.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{}); !isErrObjectNotFound(err) { + t.Errorf("unversioned fanout left pool %d readable: %v", i, err) + } + } +} + func TestPoolsConditionalDeleteVersionSelection(t *testing.T) { z, bucket := consistencyPools(t) const object = "split-versions" @@ -581,6 +1056,8 @@ func TestPoolsRetiringCopyPreservesSharedTierObject(t *testing.T) { failPrimaryDelete bool differentRemote bool restored bool + unconditional bool + skipFreeVersion bool }{ {name: "metadata-copy"}, {name: "restored-metadata-copy", restored: true}, @@ -588,6 +1065,9 @@ func TestPoolsRetiringCopyPreservesSharedTierObject(t *testing.T) { {name: "failed-primary-delete", deleting: true, failPrimaryDelete: true}, {name: "failed-primary-delete-distinct-reference", deleting: true, failPrimaryDelete: true, differentRemote: true}, {name: "successful-delete", deleting: true}, + {name: "ordinary-failed-primary-delete", deleting: true, unconditional: true, failPrimaryDelete: true}, + {name: "ordinary-successful-delete", deleting: true, unconditional: true}, + {name: "ordinary-skip-free-version", deleting: true, unconditional: true, skipFreeVersion: true}, } { t.Run(test.name, func(t *testing.T) { z, bucket := consistencyPools(t) @@ -627,10 +1107,14 @@ func TestPoolsRetiringCopyPreservesSharedTierObject(t *testing.T) { defer func() { set.getDisks = getDisks }() } if test.deleting { - _, err := z.DeleteObject(t.Context(), bucket, object, ObjectOptions{ - Versioned: true, VersionID: oi.VersionID, + opts := ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, SkipFreeVersion: test.skipFreeVersion, CheckPrecondFn: func(info ObjectInfo) bool { return info.ETag != current.ETag }, - }) + } + if test.unconditional { + opts.CheckPrecondFn = nil + } + _, err := z.DeleteObject(t.Context(), bucket, object, opts) if (err != nil) != test.failPrimaryDelete { t.Fatalf("unexpected authoritative delete result: %v", err) } @@ -655,6 +1139,7 @@ func TestPoolsRetiringCopyPreservesSharedTierObject(t *testing.T) { t.Fatalf("lost retained authoritative copy: %v", err) } for pool, wantFree := range []bool{primaryDeleted, test.differentRemote} { + wantFree = wantFree && !test.skipFreeVersion for _, disk := range z.serverPools[pool].getHashedSet(object).getDisks() { data, err := disk.ReadAll(t.Context(), bucket, pathJoin(object, xlStorageFormatFile)) if errors.Is(err, errFileNotFound) && !wantFree { @@ -692,3 +1177,119 @@ func (d consistencyDeleteFaultDisk) DeleteVersion(ctx context.Context, volume, p } return d.StorageAPI.DeleteVersion(ctx, volume, path, fi, forceDelMarker, opts) } + +// Exercise the surviving production rebalance copy, then interrupt the workflow +// before its later source cleanup. Repeated versions are reachable without the +// retired access-tier mover, and an API delete must remove both copies. +func TestPoolsDeleteVersionAfterInterruptedRebalance(t *testing.T) { + z, _ := consistencyPools(t) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + bucket, router, err := initAPIHandlerTest(ctx, z, nil, MakeBucketOptions{VersioningEnabled: true}) + if err != nil { + t.Fatal(err) + } + for source := range 2 { + t.Run(fmt.Sprintf("source=%d", source), func(t *testing.T) { + object := fmt.Sprintf("interrupted-rebalance-%d", source) + original := putConsistencyObject(t, z, bucket, object, source, "original", ObjectOptions{Versioned: true, MTime: UTCNow().Add(-time.Hour)}) + stats := []*rebalanceStats{{}, {}} + stats[source] = &rebalanceStats{Participating: true, Info: rebalanceInfo{Status: rebalStarted}} + z.rebalMu.Lock() + z.rebalMeta = &rebalanceMeta{PoolStats: stats} + z.rebalMu.Unlock() + defer func() { z.rebalMu.Lock(); z.rebalMeta = nil; z.rebalMu.Unlock() }() + gr, err := z.serverPools[source].GetObjectNInfo(ctx, bucket, object, nil, nil, ObjectOptions{VersionID: original.VersionID, NoLock: true}) + if err != nil { + t.Fatal(err) + } + if err := z.rebalanceObject(ctx, source, bucket, gr); err != nil { + t.Fatalf("rebalance copy: %v", err) + } + // Simulate interruption after rebalanceObject, before the enclosing loop + // removes the source version stack. Stop rebalance before the API request. + z.rebalMu.Lock() + z.rebalMeta = nil + z.rebalMu.Unlock() + for i, pool := range z.serverPools { + got, err := pool.GetObjectInfo(ctx, bucket, object, ObjectOptions{VersionID: original.VersionID}) + if err != nil || got.VersionID != original.VersionID || !got.ModTime.Equal(original.ModTime) { + t.Fatalf("copy missing/changed in pool %d: %+v, %v", i, got, err) + } + } + later, err := z.PutObject(ctx, bucket, object, mustGetPutObjReader(t, bytes.NewBufferString("later"), 5, "", ""), ObjectOptions{Versioned: true}) + if err != nil { + t.Fatal(err) + } + rec := consistencyRequest(t, router, http.MethodDelete, bucket, object, original.VersionID) + if rec.Code != http.StatusNoContent { + t.Fatalf("DELETE: %d %s", rec.Code, rec.Body.String()) + } + for i, pool := range z.serverPools { + if _, err := pool.GetObjectInfo(ctx, bucket, object, ObjectOptions{VersionID: original.VersionID}); !isErrVersionNotFound(err) { + t.Errorf("pool %d retained old version: %v", i, err) + } + } + if got, err := z.GetObjectInfo(ctx, bucket, object, ObjectOptions{}); err != nil || got.VersionID != later.VersionID { + t.Fatalf("other version changed: %+v, %v", got, err) + } + for _, method := range []string{http.MethodGet, http.MethodHead} { + if rec := consistencyRequest(t, router, method, bucket, object, original.VersionID); rec.Code != http.StatusNotFound { + t.Errorf("%s returned %d", method, rec.Code) + } + } + }) + } +} + +func TestPoolsDeleteDirectoryMarker(t *testing.T) { + z, _ := consistencyPools(t) + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + bucket, router, err := initAPIHandlerTest(ctx, z, nil, MakeBucketOptions{VersioningEnabled: true}) + if err != nil { + t.Fatal(err) + } + for holding := range 2 { + for _, unreadable := range []bool{false, true} { + t.Run(fmt.Sprintf("holding=%d/unreadable=%t", holding, unreadable), func(t *testing.T) { + object := fmt.Sprintf("directory-%d-%t/", holding, unreadable) + encoded := encodeDirObject(object) + putConsistencyObject(t, z, bucket, encoded, holding, "", ObjectOptions{}) + set := z.serverPools[1-holding].getHashedSet(encoded) + getDisks := set.getDisks + defer func() { set.getDisks = getDisks }() + if unreadable { + disks := append([]StorageAPI(nil), getDisks()...) + for i := range disks { + disks[i] = consistencyReadFaultDisk{StorageAPI: disks[i], bucket: bucket, object: encoded} + } + set.getDisks = func() []StorageAPI { return disks } + } + // No explicit versionId: delOpts permanently deletes the null version. + rec := consistencyRequest(t, router, http.MethodDelete, bucket, object, "") + if unreadable { + if rec.Code != http.StatusServiceUnavailable { + t.Fatalf("unreadable pool: %d %s", rec.Code, rec.Body.String()) + } + if _, err := z.serverPools[holding].GetObjectInfo(ctx, bucket, encoded, ObjectOptions{VersionID: nullVersionID}); err != nil { + t.Fatalf("readable directory marker lost: %v", err) + } + set.getDisks = getDisks + rec = consistencyRequest(t, router, http.MethodDelete, bucket, object, "") + } + if rec.Code != http.StatusNoContent { + t.Fatalf("DELETE: %d %s", rec.Code, rec.Body.String()) + } + for i, pool := range z.serverPools { + if _, err := pool.GetObjectInfo(ctx, bucket, encoded, ObjectOptions{VersionID: nullVersionID}); !isErrVersionNotFound(err) { + t.Errorf("pool %d retained null version: %v", i, err) + } + } + if rec := consistencyRequest(t, router, http.MethodHead, bucket, object, ""); rec.Code != http.StatusNotFound { + t.Errorf("directory still visible: %d", rec.Code) + } + }) + } + } +} diff --git a/cmd/erasure-server-pool.go b/cmd/erasure-server-pool.go index 5b3c5b324..bb1f66d23 100644 --- a/cmd/erasure-server-pool.go +++ b/cmd/erasure-server-pool.go @@ -1205,8 +1205,11 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob return ObjectInfo{}, z.deletePrefix(ctx, bucket, object) } - if !z.SinglePool() && opts.CheckPrecondFn != nil { - return z.deleteObjectConditional(ctx, bucket, object, opts) + // Reconcile ordinary addressed-version deletes independently of pool movement. + reconcileVersion := opts.VersionID != "" && !opts.DataMovement && + !opts.ReplicationRequest && !opts.Expiration.Expire && !opts.InclFreeVersions + if !z.SinglePool() && (opts.CheckPrecondFn != nil || reconcileVersion) { + return z.deleteObjectReconciled(ctx, bucket, object, opts) } gopts := opts diff --git a/docs/bucket/lifecycle/access-tiering-removal.md b/docs/bucket/lifecycle/access-tiering-removal.md index ded0e8698..5a0e50d0f 100644 --- a/docs/bucket/lifecycle/access-tiering-removal.md +++ b/docs/bucket/lifecycle/access-tiering-removal.md @@ -25,3 +25,10 @@ The published Server 20260903 predates this feature. These instructions concern Interrupted rebalance/decommission can leave the same version in more than one pool independently of access tiering. Removing the scheduler does not remove such existing copies. General Object Lock, conditional-delete, metadata reconciliation and shared remote-tier reference protections remain in place. +## Version deletion scope + +Ordinary single-object `DELETE ?versionId=...` reconciles the addressed UUID, null version or delete marker across pools. Unqualified DELETE of a directory marker (a key ending in `/`) also addresses its null version and uses this path. A successful response means the addressed copies were removed; other versions remain. + +If a pool is unreadable, these requests can return 503 even when another pool has a readable copy. This extends an existing failure surface: previously the result could depend on whether the unreadable pool preceded the successful pool in traversal order; it now fails consistently. Retry after recovery. Cleanup failures also return an error. Ordinary unqualified DELETE retains its existing semantics. Batch `DeleteObjects` already fans out across pools. + +Incoming replicated deletes, lifecycle expiration, free-version cleanup and movement-internal calls retain their existing contracts. In particular, an incoming replicated version delete can leave movement duplicates in other pools; this change does not solve that separate case. Expiration scanners process their own pools and may remove duplicate expired copies in later cycles; free-version cleanup remains local to a pool. Do not treat the ordinary DELETE repair as a guarantee for every source of deletion.