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.