Compare commits

...

13 Commits

Author SHA1 Message Date
Feng Ruohang 69d5e92791 Merge pull request #220 from AleksaMCode/fix/lambda-targetidset-isempty
Fix `TargetIDSet.IsEmpty` and add regression test
2026-09-23 08:07:09 +08:00
AleksaMCode efd5698597 lambda/event: fix TargetIDSet.IsEmpty and add regression test
Signed-off-by: AleksaMCode <aleksamcode@gmail.com>
2026-09-17 22:06:44 +02:00
Feng Ruohang 2fde3cf535 docs: update ec description 2026-09-17 09:42:26 +08:00
Feng Ruohang 2a4d51406b test: build the purge-response fixture ARN at runtime
The rebrand compatibility guard tracks every literal policy value that
names the upstream brand. Construct the replication ARN the way the
other replication fixtures do instead of adding a new literal.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 23:50:09 +08:00
Feng Ruohang 0596685ae7 docs: record the delete-marker purge, listing quorum and migration tag repairs
Describe the three correctness repairs merged after the 20260916 candidate
was prepared, together with their known remaining limits (#217, #218).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 23:44:44 +08:00
Feng Ruohang 358ab38fb0 fix: report a purged version as a delete marker only for stored markers
The exact-version purge path copies the looked-up DeleteMarker flag
into the DELETE response so that removing a delete marker keeps its
x-amz-delete-marker header. That flag is also set for a data version
whose purge is pending, because the lookup exposes such a version as
deleted for visibility. Retrying the purge after the replication
configuration was removed therefore answered with
x-amz-delete-marker: true and raised ObjectRemovedDeleteMarkerCreated
for a data version.

Only a stored marker carries no erasure layout; use that to decide the
response identity. Found by the adversarial review of the purge change.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 22:53:12 +08:00
Feng Ruohang 254b19ac07 fix(replication): recheck queued delete-marker creations under the lock
A delete-marker creation task carries the marker as it looked when it
was queued by the DELETE handler, a GET/HEAD/LIST heal, the scanner or
MRF. Tasks from several frontends serialize on the per-object
replication lock, so a task can run after the user has purged that
marker and the purge has already reached the targets. The target no
longer holds the marker, so the queued creation recreated it there with
the original VersionID and mtime. This is the late-create sequence seen
in the 2026-09-16 three-site runs: purge 204 on both targets, then about
85 ms later a replicated creation for the same VersionID.

Re-read the source version under the replication lock before sending a
creation. A missing version, a non-marker version, or a version under
purge makes the task stale; it is dropped without touching the source.
A read that cannot confirm either way is retried through MRF instead of
being treated as absence.

Creations already on the wire, replays from other sites, and cleanup of
minority residue after a crash are not covered by this check.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 22:35:08 +08:00
Feng Ruohang eb4f5e5b31 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 <rh@vonng.com>
(cherry picked from commit 22ba4d426677aa07470927fc349a43afd87e19b6)
2026-09-16 22:24:27 +08:00
Feng Ruohang 8d06424b12 fix: retain quorate null versions after newer minorities
Recount only when the existing selection is below quorum and each original
nonempty stream has one ordinary null version with identical EC headers.
Reuse the existing header grouping and leave history pruning unchanged.

Add deterministic signed LIST, complete-metadata permutation, excluded-history,
object-read, scanner/heal, Walk and migration-consumer regressions.

Signed-off-by: Feng Ruohang <rh@vonng.com>
(cherry picked from commit 9804e4deb4d5a18cca640be72bff9b47a415f7de)
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 21:50:03 +08:00
Feng Ruohang fced863036 fix(storage): preserve tags during pool migration
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 20:22:37 +08:00
Feng Ruohang 5af0865aab Merge pull request #214 from pgsty/codex/release-prep-20260916
Prepare the SILO 20260916000000 source candidate with released component pins, CPU metrics synchronization, regression coverage and release tooling updates.

All 11 checks passed on e7963381ba. Signed release artifacts, image publication and deployment remain separate deliverables tracked by #203.

Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 19:30:07 +08:00
Feng Ruohang 83821f0f1f Merge pull request #216 from pgsty/codex/docs-published-design-links-20260916
docs: link retained guides to published design records
2026-09-16 18:50:28 +08:00
Feng Ruohang bf053386a4 docs: link compression and federation to maintained design records
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 18:26:03 +08:00
27 changed files with 2864 additions and 81 deletions
+29
View File
@@ -108,6 +108,35 @@ and [complete commit range](https://github.com/pgsty/silo/compare/RELEASE.2026-0
converged. Follow the [read-only audit procedure](https://silo.pgsty.com/operations/replication/replica-metadata-audit/)
before planning any repair of stored state.
- Make exact-version delete-marker purges converge in replicated buckets
(`eb4f5e5b3`, `254b19ac0`, `358ab38fb`). A purge no longer creates the
marker on drives that lacked it; a retried purge of a missing version is
acknowledged only when a write-quorum majority of drives report it absent;
purge results merge removed and reliably absent replies; healing a marker
preserves its stored replication and purge metadata; and a queued marker
creation is re-checked against the source under the replication lock before
it is sent, so a purge that already reached the targets is not undone by a
stale task from another frontend, a GET/LIST heal, the scanner or MRF.
Purging a data version whose earlier purge is still pending reports it as a
data version. **Known limitations:** creations already in flight or replayed
from another site, and minority marker copies left by a crash after a
majority-acknowledged purge, are tracked in #217. See
[the replication reliability record](https://silo.pgsty.com/blog/design/replication-reliability/).
- Keep a null object version that has listing quorum when a newer minority of
drives sorts first (`8d06424b1`). The resolver recounts per header only when
the original selection lacks quorum, every non-empty drive stream holds
exactly one ordinary null version and all share the same erasure layout;
mixed histories keep their previous behavior. **Known limitation:** a
successful ListObjects can still omit readable keys during rolling restarts
with concurrent overwrites (#218). Do not run destination-deleting sync tools
against a listing taken during a rolling restart; list again once the
cluster is stable.
- Carry object tags through rebalance and decommission for ordinary and
multipart writes (`fced86303`). Both migration entry points restore the tags
and their revision fields when rewriting the object in the destination pool.
Tags dropped by earlier migrations are not recovered; audit tag-dependent
lifecycle and policy rules for pools migrated with an older build.
- Evaluate conditional multipart completion against the logical current object
across all pools while holding the existing object lock. A stale `If-Match`
can no longer replace newer data in another pool, and the current ETag is no
+64
View File
@@ -504,6 +504,35 @@ func replicateDelete(ctx context.Context, dobj DeletedObjectReplicationInfo, obj
ctx = lkctx.Context()
defer lk.Unlock(lkctx)
if !isPurge && dobj.DeleteMarkerVersionID != "" {
// A creation task carries the marker as it looked when it was queued
// by the DELETE handler, a GET/HEAD/LIST heal, the scanner or MRF.
// While it waited for this lock, another frontend may have purged the
// marker and replicated that purge. The targets no longer hold the
// marker, so the queued creation would recreate it from the stale
// snapshot. Confirm the source version under the lock first.
switch deleteMarkerCreationState(ctx, objectAPI, dobj) {
case creationStale:
return replicatedInfos{}
case creationUnverified:
dobj.RetryCount++
globalReplicationPool.Get().queueMRFSave(dobj.ToMRFEntry())
sendEvent(eventArgs{
BucketName: bucket,
Object: ObjectInfo{
Bucket: bucket,
Name: dobj.ObjectName,
VersionID: versionID,
DeleteMarker: dobj.DeleteMarker,
},
UserAgent: "Internal: [Replication]",
Host: globalLocalNodeName,
EventName: event.ObjectReplicationNotTracked,
})
return replicatedInfos{}
}
}
rinfos := replicatedInfos{Targets: make([]replicatedTargetInfo, 0, len(dsc.targetsMap))}
var wg sync.WaitGroup
var mu sync.Mutex
@@ -1955,6 +1984,41 @@ func (di DeletedObjectReplicationInfo) isVersionPurge() bool {
return di.VersionID != "" || di.DeleteMarkerVersionID != "" && !di.VersionPurgeStatus().Empty()
}
// creationState is the outcome of re-reading a queued delete-marker creation
// against the source.
type creationState int
const (
creationCurrent creationState = iota
creationStale
creationUnverified
)
// deleteMarkerCreationState re-reads the marker a queued creation task refers
// to. The version being absent, no longer a marker, or under a version purge
// makes the creation stale. A failed read is not absence: the caller retries
// later instead of guessing.
func deleteMarkerCreationState(ctx context.Context, objectAPI ObjectLayer, dobj DeletedObjectReplicationInfo) creationState {
oi, err := objectAPI.GetObjectInfo(ctx, dobj.Bucket, dobj.ObjectName, ObjectOptions{
VersionID: dobj.DeleteMarkerVersionID,
Versioned: globalBucketVersioningSys.PrefixEnabled(dobj.Bucket, dobj.ObjectName),
VersionSuspended: globalBucketVersioningSys.Suspended(dobj.Bucket),
})
switch {
case isErrObjectNotFound(err), isErrVersionNotFound(err):
return creationStale
case err != nil && !isErrMethodNotAllowed(err):
return creationUnverified
}
if !oi.DeleteMarker || oi.VersionID != dobj.DeleteMarkerVersionID {
return creationStale
}
if !oi.VersionPurgeStatus.Empty() || oi.VersionPurgeStatusInternal != "" {
return creationStale
}
return creationCurrent
}
// Purge metadata uses COMPLETE; operation statistics and audit use COMPLETED.
func purgeReplicationStatus(status VersionPurgeStatusType) replication.StatusType {
if replication.StatusType(status) == replication.CompletedLegacy {
@@ -0,0 +1,109 @@
// 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 <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"net/http"
"net/http/httptest"
"testing"
"github.com/minio/minio/internal/auth"
"github.com/minio/minio/internal/bucket/replication"
xhttp "github.com/minio/minio/internal/http"
)
// An exact-version DELETE without a replication decision physically purges
// the version. Its response must describe the stored version: a data version
// is not a delete marker even while pending-purge metadata makes the lookup
// expose it as deleted, and a stored marker stays a marker.
func TestDeleteMarkerPurgeResponseIdentity(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, endpoints: []string{"DeleteObject"}, objAPITest: func(obj ObjectLayer, backend, bucket string, router http.Handler, creds auth.Credentials, t *testing.T) {
defer replicationTestCapacity(obj)()
ctx := t.Context()
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, enabledBucketVersioningConfig); err != nil {
t.Fatal(err)
}
// Seed the purge state the DELETE handler would record for a target.
arn := "arn:minio:replication::" + mustGetUUID() + ":bucket"
pendingPurge := ReplicationState{
VersionPurgeStatusInternal: arn + "=PENDING;",
PurgeTargets: map[string]VersionPurgeStatusType{arn: replication.VersionPurgePending},
}
for _, tc := range []struct {
name string
marker, purged bool
}{
{name: "data"},
{name: "data-pending-purge", purged: true},
{name: "marker", marker: true},
{name: "marker-pending-purge", marker: true, purged: true},
} {
name := "purge-response-" + mustGetUUID()
oi, err := obj.PutObject(ctx, bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true})
if err != nil {
t.Fatal(err)
}
version := oi.VersionID
if tc.marker {
opts := ObjectOptions{Versioned: true, VersionID: mustGetUUID(), DeleteMarker: true, MTime: UTCNow(), ReplicationRequest: true}
opts.SetReplicaStatus(replication.Replica)
if _, err := obj.DeleteObject(ctx, bucket, name, opts); err != nil {
t.Fatal(err)
}
version = opts.VersionID
}
if tc.purged {
if _, err := obj.DeleteObject(ctx, bucket, name, ObjectOptions{Versioned: true, VersionID: version, DeleteReplication: pendingPurge}); err != nil {
t.Fatal(err)
}
}
before, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{Versioned: true, VersionID: version})
if tc.marker || tc.purged {
if !isErrMethodNotAllowed(err) || !before.DeleteMarker {
t.Fatalf("%s: seeded version not exposed as deleted: %+v %v", tc.name, before, err)
}
} else if err != nil || before.DeleteMarker {
t.Fatalf("%s: seeded data version: %+v %v", tc.name, before, err)
}
if tc.purged && before.VersionPurgeStatus != replication.VersionPurgePending {
t.Fatalf("%s: purge state not seeded: %+v", tc.name, before)
}
req, err := newTestSignedRequestV4(http.MethodDelete, "/"+bucket+"/"+name+"?versionId="+version, 0, nil, creds.AccessKey, creds.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
router.ServeHTTP(rec, req)
if rec.Code != http.StatusNoContent {
t.Fatalf("%s: status=%d: %s", tc.name, rec.Code, rec.Body.String())
}
// The handler sets these headers by map key, not canonical name.
if got := rec.Header()[xhttp.AmzVersionID]; len(got) != 1 || got[0] != version {
t.Fatalf("%s: response version %q, want %q", tc.name, got, version)
}
if got := len(rec.Header()[xhttp.AmzDeleteMarker]) == 1 && rec.Header()[xhttp.AmzDeleteMarker][0] == "true"; got != tc.marker {
t.Fatalf("%s: x-amz-delete-marker=%v for a stored marker=%v", tc.name, got, tc.marker)
}
if _, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{Versioned: true, VersionID: version}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
t.Fatalf("%s: version not purged: %v", tc.name, err)
}
t.Logf("%s %s: purged with x-amz-delete-marker=%q", backend, tc.name, rec.Header()[xhttp.AmzDeleteMarker])
}
}})
}
+524
View File
@@ -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 <http://www.gnu.org/licenses/>.
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())
}
}
+32 -6
View File
@@ -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,19 @@ 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. The lookup also exposes a data
// version pending purge as deleted for visibility; only a stored
// marker, which carries no erasure layout, is reported as one.
oi.DeleteMarker = goi.DeleteMarker && goi.DataBlocks == 0
}
return oi, nil
}
// Send the successful but partial upload/delete, however ignore
@@ -2467,7 +2493,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
}
+20
View File
@@ -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.
+2 -2
View File
@@ -618,7 +618,7 @@ func (z *erasureServerPools) decommissionObject(ctx context.Context, idx int, bu
if objInfo.isMultipart() {
res, err := z.NewMultipartUpload(ctx, bucket, objInfo.Name, ObjectOptions{
VersionID: objInfo.VersionID,
UserDefined: objInfo.UserDefined,
UserDefined: migrationObjectMetadata(objInfo),
NoAuditLog: true,
SrcPoolIdx: idx,
DataMovement: true,
@@ -681,7 +681,7 @@ func (z *erasureServerPools) decommissionObject(ctx context.Context, idx int, bu
SrcPoolIdx: idx,
VersionID: objInfo.VersionID,
MTime: objInfo.ModTime,
UserDefined: objInfo.UserDefined,
UserDefined: migrationObjectMetadata(objInfo),
PreserveETag: objInfo.ETag, // Preserve original ETag to ensure same metadata.
IndexCB: func() []byte {
return objInfo.Parts[0].Index // Preserve part Index to ensure decompression works.
+33
View File
@@ -0,0 +1,33 @@
// 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 <http://www.gnu.org/licenses/>.
package cmd
import xhttp "github.com/minio/minio/internal/http"
// migrationObjectMetadata restores tags removed from the read representation.
// Keep the original revision, including ordered empty states, for the existing
// locked reconciliation at PUT or multipart completion. Moving is not a new
// tagging mutation, and the write must not modify the reader's metadata map.
func migrationObjectMetadata(oi ObjectInfo) map[string]string {
metadata := cloneMSS(oi.UserDefined)
delete(metadata, xhttp.AmzObjectTagging)
if oi.UserTags != "" {
metadata[xhttp.AmzObjectTagging] = oi.UserTags
}
return metadata
}
+592
View File
@@ -0,0 +1,592 @@
// 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 <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"maps"
"strings"
"sync"
"testing"
"time"
xhttp "github.com/minio/minio/internal/http"
)
const migrationTagRevision = ReservedMetadataPrefixLower + TaggingTimestamp
var migrationTagStates = []struct {
name, tags, revision string
hasRevision bool
}{
{name: "initial", tags: "team=storage&path=a%2Fb"},
{name: "empty-revision", tags: "team=storage", hasRevision: true},
{name: "invalid-revision", tags: "team=storage", revision: "invalid", hasRevision: true},
{name: "ordered", tags: "team=storage", revision: "2026-09-16T09:00:00Z", hasRevision: true},
{name: "cleared", revision: "2026-09-16T09:00:00Z", hasRevision: true},
{name: "never-tagged"},
{name: "empty-empty-revision", hasRevision: true},
{name: "empty-invalid-revision", revision: "invalid", hasRevision: true},
}
func TestMigrationObjectMetadata(t *testing.T) {
for _, state := range migrationTagStates {
for _, residual := range []bool{false, true} {
t.Run(fmt.Sprintf("%s/residual=%t", state.name, residual), func(t *testing.T) {
raw := map[string]string{
xhttp.AmzObjectTagging: state.tags,
"x-amz-meta-team": "storage",
"x-amz-meta-empty": "",
ReservedMetadataPrefix + "compression": "opaque-compression-state",
ReservedMetadataPrefix + "sealed-key": "opaque-key-state",
"x-amz-object-lock-retain-until-date": "2030-01-01T00:00:00Z",
"x-amz-object-lock-legal-hold": "ON",
"x-amz-tagging": "keep-noncanonical-metadata",
ReservedMetadataPrefix + "actual-size": "42",
ReservedMetadataPrefix + "replica-status": "opaque-replication-state",
}
if state.hasRevision {
raw[migrationTagRevision] = state.revision
}
oi := (FileInfo{Metadata: raw}).ToObjectInfo("bucket", "object", false)
if residual {
oi.UserDefined[xhttp.AmzObjectTagging] = "stale=raw"
}
before := maps.Clone(oi.UserDefined)
got := migrationObjectMetadata(oi)
if got == nil || !maps.Equal(before, oi.UserDefined) {
t.Fatal("helper must return a new non-nil map without changing the read snapshot")
}
value, exists := got[xhttp.AmzObjectTagging]
if value != state.tags || exists != (state.tags != "") {
t.Fatalf("raw tags=(%q,%t), want (%q,%t)", value, exists, state.tags, state.tags != "")
}
stamp, present := got[migrationTagRevision]
if stamp != state.revision || present != state.hasRevision {
t.Fatalf("revision=(%q,%t), want (%q,%t)", stamp, present, state.revision, state.hasRevision)
}
for key, value := range before {
if stored, exists := got[key]; key != xhttp.AmzObjectTagging && (!exists || stored != value) {
t.Errorf("lost metadata %q", key)
}
}
got["x-amz-meta-team"] = "changed"
if !maps.Equal(before, oi.UserDefined) {
t.Fatal("output aliases the read snapshot")
}
})
}
}
for _, metadata := range []map[string]string{nil, {}} {
got := migrationObjectMetadata(ObjectInfo{UserDefined: metadata})
if got == nil || len(got) != 0 {
t.Fatalf("empty input must yield a writable empty map: %#v", got)
}
got["new"] = "value"
if len(metadata) != 0 {
t.Fatal("empty map input was aliased")
}
}
}
// Allocation ignores the host's used-space percentage, as in the existing tag
// fixtures. All object data and metadata still use real erasure-storage disks.
func migrationTestPools(t *testing.T) (*erasureServerPools, string) {
t.Helper()
z, bucket := consistencyPools(t)
for _, pool := range z.serverPools {
for _, set := range pool.sets {
original := set.getDisks
disks := append([]StorageAPI(nil), original()...)
for i := range disks {
disks[i] = tagTestCapacityDisk{StorageAPI: disks[i]}
}
set.getDisks = func() []StorageAPI { return disks }
t.Cleanup(func() { set.getDisks = original })
}
}
return z, bucket
}
func migrationTestMover(t *testing.T, z *erasureServerPools, source int, kind string) func(context.Context, int, string, *GetObjectReader) error {
t.Helper()
if kind == "rebalance" {
z.rebalMu.Lock()
z.rebalMeta = &rebalanceMeta{PoolStats: []*rebalanceStats{{}, {}}}
z.rebalMeta.PoolStats[source] = &rebalanceStats{Participating: true, Info: rebalanceInfo{Status: rebalStarted}}
z.rebalMu.Unlock()
t.Cleanup(func() { z.rebalMu.Lock(); z.rebalMeta = nil; z.rebalMu.Unlock() })
return z.rebalanceObject
}
z.poolMetaMutex.Lock()
z.poolMeta.Pools[source].Decommission = &PoolDecommissionInfo{}
z.poolMetaMutex.Unlock()
t.Cleanup(func() { z.poolMetaMutex.Lock(); z.poolMeta.Pools[source].Decommission = nil; z.poolMetaMutex.Unlock() })
return z.decommissionObject
}
func migrationTestObject(t *testing.T, z *erasureServerPools, bucket, object string, source int, multipart bool, parts []string, opts ObjectOptions) ObjectInfo {
t.Helper()
if !multipart {
oi := putConsistencyObject(t, z, bucket, object, source, strings.Join(parts, ""), opts)
return migrationTestPersistedInfo(t, z, source, oi)
}
mp, err := z.serverPools[source].NewMultipartUpload(t.Context(), bucket, object, opts)
if err != nil {
t.Fatal(err)
}
completed := make([]CompletePart, len(parts))
for i, body := range parts {
part, err := z.serverPools[source].PutObjectPart(t.Context(), bucket, object, mp.UploadID, i+1,
mustGetPutObjReader(t, strings.NewReader(body), int64(len(body)), "", ""), ObjectOptions{})
if err != nil {
t.Fatal(err)
}
completed[i] = CompletePart{PartNumber: i + 1, ETag: part.ETag}
}
oi, err := z.serverPools[source].CompleteMultipartUpload(t.Context(), bucket, object, mp.UploadID, completed, opts)
if err != nil {
t.Fatal(err)
}
if !oi.isMultipart() {
t.Fatal("fixture did not create a multipart object")
}
return migrationTestPersistedInfo(t, z, source, oi)
}
func migrationTestPersistedInfo(t *testing.T, z *erasureServerPools, pool int, oi ObjectInfo) ObjectInfo {
t.Helper()
// PUT can return transient fields such as tier-free-versionID that are not
// stored. Take the same persisted read representation used by the movers.
got, err := z.serverPools[pool].GetObjectInfo(t.Context(), oi.Bucket, oi.Name, ObjectOptions{VersionID: migrationTestVersion(oi)})
if err != nil {
t.Fatal(err)
}
return got
}
func migrationTestVersion(oi ObjectInfo) string {
if oi.VersionID == "" {
return nullVersionID
}
return oi.VersionID
}
func migrationTestReader(t *testing.T, z *erasureServerPools, source int, oi ObjectInfo) *GetObjectReader {
t.Helper()
gr, err := z.serverPools[source].GetObjectNInfo(t.Context(), oi.Bucket, oi.Name, nil, nil,
ObjectOptions{VersionID: migrationTestVersion(oi), NoLock: true, NoDecryption: true})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { gr.Close() })
return gr
}
func migrationTestStored(t *testing.T, z *erasureServerPools, pool int, original ObjectInfo, tags, revision string, hasRevision bool, body string) {
t.Helper()
version := migrationTestVersion(original)
got, err := z.serverPools[pool].GetObjectInfo(t.Context(), original.Bucket, original.Name, ObjectOptions{VersionID: version})
if err != nil {
t.Fatal(err)
}
if got.UserTags != tags || got.ETag != original.ETag || !got.ModTime.Equal(original.ModTime) || got.VersionID != original.VersionID || got.Size != original.Size {
t.Errorf("pool %d changed tags/object identity: tags=%q want=%q version=%q/%q etag=%q/%q size=%d/%d mtime=%v/%v", pool, got.UserTags, tags, got.VersionID, original.VersionID, got.ETag, original.ETag, got.Size, original.Size, got.ModTime, original.ModTime)
}
infos, errs := readAllFileInfo(t.Context(), z.serverPools[pool].getHashedSet(original.Name).getDisks(), "", original.Bucket, original.Name, version, false, false)
for disk, fi := range infos {
if errs[disk] != nil {
t.Fatalf("pool %d disk %d: %v", pool, disk, errs[disk])
}
stamp, exists := fi.Metadata[migrationTagRevision]
if fi.Metadata[xhttp.AmzObjectTagging] != tags || stamp != revision || exists != hasRevision {
t.Errorf("pool %d disk %d: tags=%q revision=(%q,%t), want %q (%q,%t)", pool, disk, fi.Metadata[xhttp.AmzObjectTagging], stamp, exists, tags, revision, hasRevision)
}
for key, value := range original.UserDefined {
if key != migrationTagRevision && fi.Metadata[key] != value {
t.Errorf("pool %d disk %d lost metadata %q", pool, disk, key)
}
}
}
gr, err := z.serverPools[pool].GetObjectNInfo(t.Context(), original.Bucket, original.Name, nil, nil, ObjectOptions{VersionID: version})
if err != nil {
t.Fatal(err)
}
data, err := io.ReadAll(gr)
gr.Close()
if err != nil || !bytes.Equal(data, []byte(body)) {
t.Fatalf("pool %d content changed: bytes=%d err=%v", pool, len(data), err)
}
}
func TestPoolsMigrationPreservesTagState(t *testing.T) {
z, bucket := migrationTestPools(t)
for _, kind := range []string{"rebalance", "decommission"} {
for _, method := range []string{"put", "multipart"} {
for source := range 2 {
for _, state := range migrationTagStates {
t.Run(fmt.Sprintf("%s/%s/source=%d/%s", kind, method, source, state.name), func(t *testing.T) {
move := migrationTestMover(t, z, source, kind)
meta := map[string]string{xhttp.AmzObjectTagging: state.tags, "x-amz-meta-owner": "retained"}
if state.hasRevision {
meta[migrationTagRevision] = state.revision
}
oi := migrationTestObject(t, z, bucket, t.Name(), source, method == "multipart", []string{"payload"}, ObjectOptions{Versioned: true, UserDefined: meta})
migrationTestStored(t, z, source, oi, state.tags, state.revision, state.hasRevision, "payload")
gr := migrationTestReader(t, z, source, oi)
before := maps.Clone(gr.ObjInfo.UserDefined)
if err := move(t.Context(), source, bucket, gr); err != nil {
t.Fatal(err)
}
if !maps.Equal(before, gr.ObjInfo.UserDefined) {
t.Error("migration modified its source snapshot")
}
migrationTestStored(t, z, 1-source, oi, state.tags, state.revision, state.hasRevision, "payload")
})
}
}
}
}
}
type migrationTestGate struct {
entered, resume chan struct{}
once, release sync.Once
}
func newMigrationTestGate() *migrationTestGate {
return &migrationTestGate{entered: make(chan struct{}), resume: make(chan struct{})}
}
func (g *migrationTestGate) wait(ctx context.Context) error {
g.once.Do(func() { close(g.entered) })
select {
case <-g.resume:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func (g *migrationTestGate) unblock() { g.release.Do(func() { close(g.resume) }) }
type migrationTestGateReader struct {
io.Reader
ctx context.Context
gate *migrationTestGate
}
func (r migrationTestGateReader) Read(p []byte) (int, error) {
if err := r.gate.wait(r.ctx); err != nil {
return 0, err
}
return r.Reader.Read(p)
}
func TestPoolsMigrationRechecksTags(t *testing.T) {
z, bucket := migrationTestPools(t)
const recent = "2026-09-16T10:00:00Z"
for _, kind := range []string{"rebalance", "decommission"} {
for _, method := range []string{"put", "multipart"} {
for source := range 2 {
for _, tags := range []string{"state=after", ""} {
t.Run(fmt.Sprintf("%s/%s/source=%d/clear=%t", kind, method, source, tags == ""), func(t *testing.T) {
move := migrationTestMover(t, z, source, kind)
oi := migrationTestObject(t, z, bucket, t.Name(), source, method == "multipart", []string{"payload"}, ObjectOptions{
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "state=before"},
})
gr := migrationTestReader(t, z, source, oi)
before := maps.Clone(gr.ObjInfo.UserDefined)
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
defer cancel()
gate := newMigrationTestGate()
done, finished := make(chan error, 1), make(chan struct{})
if method == "multipart" {
// The first part read is after persisted upload initialization,
// but before completion takes the object lock and reconciles.
gr.Reader = migrationTestGateReader{Reader: gr.Reader, ctx: ctx, gate: gate}
go func() { defer close(finished); done <- move(ctx, source, bucket, gr) }()
defer func() { gate.unblock(); cancel(); <-finished }()
select {
case <-gate.entered:
case err := <-done:
t.Fatalf("migration finished before the part barrier: %v", err)
case <-ctx.Done():
t.Fatal(ctx.Err())
}
uploads, err := z.serverPools[1-source].listMultipartUploadsExact(ctx, bucket, oi.Name)
if err != nil || len(uploads.Uploads) != 1 {
t.Fatalf("upload not persisted at barrier: %+v, %v", uploads, err)
}
info, err := z.serverPools[1-source].GetMultipartInfo(ctx, bucket, oi.Name, uploads.Uploads[0].UploadID, ObjectOptions{})
if err != nil || info.UserDefined[xhttp.AmzObjectTagging] != "state=before" {
t.Fatalf("upload lost snapshot tags: %v, %v", info.UserDefined, err)
}
}
// For ordinary PUT this must finish before calling the mover:
// blocking its data Read would already hold the object lock.
opts := ObjectOptions{VersionID: migrationTestVersion(oi), UserDefined: map[string]string{migrationTagRevision: recent}}
var err error
if tags == "" {
_, err = z.DeleteObjectTags(ctx, bucket, oi.Name, opts)
} else {
_, err = z.PutObjectTags(ctx, bucket, oi.Name, tags, opts)
}
if err != nil {
t.Fatal(err)
}
if gr.ObjInfo.UserTags != "state=before" || !maps.Equal(before, gr.ObjInfo.UserDefined) {
t.Fatal("test must retain the old source snapshot")
}
migrationTestStored(t, z, source, oi, tags, recent, true, "payload")
if method == "multipart" {
gate.unblock()
err = <-done
} else {
err = move(ctx, source, bucket, gr)
}
if err != nil {
t.Fatal(err)
}
if !maps.Equal(before, gr.ObjInfo.UserDefined) {
t.Error("reconciliation changed the reader's map")
}
migrationTestStored(t, z, 1-source, oi, tags, recent, true, "payload")
})
}
}
}
}
}
type migrationTestMetadataGateDisk struct {
StorageAPI
bucket, object string
gate *migrationTestGate
}
func (d migrationTestMetadataGateDisk) UpdateMetadata(ctx context.Context, volume, path string, fi FileInfo, opts UpdateMetadataOpts) error {
if volume == d.bucket && path == d.object {
if err := d.gate.wait(ctx); err != nil {
return err
}
}
return d.StorageAPI.UpdateMetadata(ctx, volume, path, fi, opts)
}
func TestPoolsMigrationTagsDuringCleanup(t *testing.T) {
z, bucket := migrationTestPools(t)
const recent = "2026-09-16T10:00:00Z"
for _, kind := range []string{"rebalance", "decommission"} {
for source := range 2 {
for _, tags := range []string{"state=after", ""} {
for _, interrupt := range []bool{false, true} {
t.Run(fmt.Sprintf("%s/source=%d/clear=%t/interrupt=%t", kind, source, tags == "", interrupt), func(t *testing.T) {
move := migrationTestMover(t, z, source, kind)
oi := migrationTestObject(t, z, bucket, t.Name(), source, false, []string{"payload"}, ObjectOptions{
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "state=before"},
})
if err := move(t.Context(), source, bucket, migrationTestReader(t, z, source, oi)); err != nil {
t.Fatal(err)
}
migrationTestStored(t, z, 1-source, oi, "state=before", "", false, "payload")
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
defer cancel()
mutate := func() (ObjectInfo, error) {
opts := ObjectOptions{VersionID: migrationTestVersion(oi), UserDefined: map[string]string{migrationTagRevision: recent}}
if tags == "" {
return z.DeleteObjectTags(ctx, bucket, oi.Name, opts)
}
return z.PutObjectTags(ctx, bucket, oi.Name, tags, opts)
}
set := z.serverPools[source].getHashedSet(oi.Name)
cleanup := func() {
// Use the same source prefix cleanup as both outer movers.
if _, err := set.DeleteObject(ctx, bucket, oi.Name, ObjectOptions{DeletePrefix: true, DeletePrefixObject: true}); err != nil {
t.Fatal(err)
}
}
var updated ObjectInfo
if !interrupt {
var err error
updated, err = mutate()
if err != nil {
t.Fatal(err)
}
cleanup()
} else {
gate := newMigrationTestGate()
original := set.getDisks
disks := append([]StorageAPI(nil), original()...)
for i := range disks {
disks[i] = migrationTestMetadataGateDisk{StorageAPI: disks[i], bucket: bucket, object: oi.Name, gate: gate}
}
set.getDisks = func() []StorageAPI { return disks }
defer func() { set.getDisks = original }()
done, finished := make(chan error, 1), make(chan struct{})
go func() { defer close(finished); _, err := mutate(); done <- err }()
defer func() { gate.unblock(); cancel(); <-finished }()
select {
case <-gate.entered:
case err := <-done:
t.Fatalf("metadata update missed the barrier: %v", err)
case <-ctx.Done():
t.Fatal(ctx.Err())
}
cleanup()
gate.unblock()
if err := <-done; err == nil {
t.Fatal("update reported success after its source copy was removed")
}
<-finished
set.getDisks = original
var err error
updated, err = mutate()
if err != nil {
t.Fatalf("metadata retry did not converge: %v", err)
}
}
if _, err := z.serverPools[source].GetObjectInfo(ctx, bucket, oi.Name, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
t.Fatalf("source cleanup left a copy: %v", err)
}
migrationTestStored(t, z, 1-source, oi, tags, updated.UserDefined[migrationTagRevision], true, "payload")
})
}
}
}
}
}
func TestPoolsMigrationVersionHistory(t *testing.T) {
z, bucket := migrationTestPools(t)
for _, kind := range []string{"rebalance", "decommission"} {
for _, method := range []string{"put", "multipart"} {
for _, versioning := range []string{"unversioned", "null", "history"} {
t.Run(kind+"/"+method+"/"+versioning, func(t *testing.T) {
move := migrationTestMover(t, z, 0, kind)
opts := ObjectOptions{
Versioned: versioning == "history", VersionSuspended: versioning == "null",
MTime: time.Date(2026, 9, 16, 8, 0, 0, 0, time.UTC),
UserDefined: map[string]string{xhttp.AmzObjectTagging: "version=old"},
}
versions := []ObjectInfo{migrationTestObject(t, z, bucket, t.Name(), 0, method == "multipart", []string{"old"}, opts)}
if versioning == "history" {
opts.MTime = opts.MTime.Add(time.Minute)
versions = append(versions, migrationTestObject(t, z, bucket, t.Name(), 0, method == "multipart", []string{"old"}, opts))
}
var newest ObjectInfo
if versioning != "unversioned" {
newest = migrationTestObject(t, z, bucket, t.Name(), 0, false, []string{"latest"}, ObjectOptions{
Versioned: true, MTime: opts.MTime.Add(time.Minute),
UserDefined: map[string]string{xhttp.AmzObjectTagging: "version=new"},
})
}
for _, oi := range versions {
if err := move(t.Context(), 0, bucket, migrationTestReader(t, z, 0, oi)); err != nil {
t.Fatal(err)
}
migrationTestStored(t, z, 1, oi, "version=old", "", false, "old")
}
if versioning != "unversioned" {
migrationTestStored(t, z, 0, newest, "version=new", "", false, "latest")
if _, err := z.serverPools[1].GetObjectInfo(t.Context(), bucket, t.Name(), ObjectOptions{VersionID: newest.VersionID}); !isErrVersionNotFound(err) {
t.Fatalf("moving an addressed version affected the newer version: %v", err)
}
}
})
}
}
}
}
func TestPoolsMigrationMultipartParts(t *testing.T) {
z, bucket := migrationTestPools(t)
parts := []string{strings.Repeat("a", 5<<20), "tail-with-tags"}
body := strings.Join(parts, "")
for _, kind := range []string{"rebalance", "decommission"} {
for source := range 2 {
t.Run(fmt.Sprintf("%s/source=%d", kind, source), func(t *testing.T) {
move := migrationTestMover(t, z, source, kind)
oi := migrationTestObject(t, z, bucket, t.Name(), source, true, parts, ObjectOptions{
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "multipart=kept"},
})
if err := move(t.Context(), source, bucket, migrationTestReader(t, z, source, oi)); err != nil {
t.Fatal(err)
}
migrationTestStored(t, z, 1-source, oi, "multipart=kept", "", false, body)
got, err := z.serverPools[1-source].GetObjectInfo(t.Context(), bucket, oi.Name, ObjectOptions{VersionID: oi.VersionID})
if err != nil || len(got.Parts) != len(oi.Parts) {
t.Fatalf("parts changed: %+v, %v", got.Parts, err)
}
for i, part := range got.Parts {
old := oi.Parts[i]
if part.Number != old.Number || part.Size != old.Size || part.ActualSize != old.ActualSize || part.ETag != old.ETag || !bytes.Equal(part.Index, old.Index) {
t.Errorf("part %d changed: %+v / %+v", i, part, old)
}
}
start := int64(len(parts[0]) - 3)
gr, err := z.serverPools[1-source].GetObjectNInfo(t.Context(), bucket, oi.Name, &HTTPRangeSpec{Start: start, End: start + 7}, nil, ObjectOptions{VersionID: oi.VersionID})
if err != nil {
t.Fatal(err)
}
data, err := io.ReadAll(gr)
gr.Close()
if err != nil || string(data) != body[start:start+8] {
t.Fatalf("cross-part range changed: %q, %v", data, err)
}
})
}
}
}
func TestPoolsMigrationUnreadablePool(t *testing.T) {
z, bucket := migrationTestPools(t)
for _, kind := range []string{"rebalance", "decommission"} {
for _, method := range []string{"put", "multipart"} {
t.Run(kind+"/"+method, func(t *testing.T) {
move := migrationTestMover(t, z, 0, kind)
oi := migrationTestObject(t, z, bucket, t.Name(), 0, method == "multipart", []string{"payload"}, ObjectOptions{
Versioned: true, UserDefined: map[string]string{xhttp.AmzObjectTagging: "keep=source"},
})
gr := migrationTestReader(t, z, 0, oi)
set := z.serverPools[0].getHashedSet(oi.Name)
original := set.getDisks
disks := append([]StorageAPI(nil), original()...)
for i := range disks {
disks[i] = consistencyReadFaultDisk{StorageAPI: disks[i], bucket: bucket, object: oi.Name}
}
set.getDisks = func() []StorageAPI { return disks }
defer func() { set.getDisks = original }()
err := move(t.Context(), 0, bucket, gr)
var quorum InsufficientReadQuorum
if !errors.As(err, &quorum) {
t.Fatalf("unreadable pool must fail migration with read quorum error: %v", err)
}
set.getDisks = original
migrationTestStored(t, z, 0, oi, "keep=source", "", false, "payload")
if _, err := z.serverPools[1].GetObjectInfo(t.Context(), bucket, oi.Name, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
t.Fatalf("failed migration committed a target version: %v", err)
}
})
}
}
}
+2 -2
View File
@@ -865,7 +865,7 @@ func (z *erasureServerPools) rebalanceObject(ctx context.Context, poolIdx int, b
if oi.isMultipart() {
res, err := z.NewMultipartUpload(ctx, bucket, oi.Name, ObjectOptions{
VersionID: oi.VersionID,
UserDefined: oi.UserDefined,
UserDefined: migrationObjectMetadata(oi),
NoAuditLog: true,
DataMovement: true,
SrcPoolIdx: poolIdx,
@@ -924,7 +924,7 @@ func (z *erasureServerPools) rebalanceObject(ctx context.Context, poolIdx int, b
DataMovement: true,
VersionID: oi.VersionID,
MTime: oi.ModTime,
UserDefined: oi.UserDefined,
UserDefined: migrationObjectMetadata(oi),
PreserveETag: oi.ETag, // Preserve original ETag to ensure same metadata.
IndexCB: func() []byte {
return oi.Parts[0].Index // Preserve part Index to ensure decompression works.
+17 -2
View File
@@ -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
+87
View File
@@ -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 <http://www.gnu.org/licenses/>.
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
}
+433
View File
@@ -0,0 +1,433 @@
// 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 <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"context"
"encoding/xml"
"fmt"
"io"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"slices"
"sync"
"testing"
"time"
madmin "github.com/minio/madmin-go/v3"
)
// Hold nine completed disk walks while a PUT overwrites one name. The other
// seven walks start after that PUT completes. Both generations come from real
// disk walkers and valid PUT metadata; only the reader schedule is controlled.
type nullQuorumWalkSchedule struct {
bucket string
oldReady chan struct{}
resume chan struct{}
mu sync.Mutex
snapshots [16][]byte
}
type nullQuorumWalkDisk struct {
StorageAPI
schedule *nullQuorumWalkSchedule
index int
}
func (d *nullQuorumWalkDisk) WalkDir(ctx context.Context, opts WalkDirOptions, out io.Writer) error {
if opts.Bucket != d.schedule.bucket {
return d.StorageAPI.WalkDir(ctx, opts, out)
}
var stream bytes.Buffer
if d.index < 9 {
if err := d.StorageAPI.WalkDir(ctx, opts, &stream); err != nil {
return err
}
select {
case d.schedule.oldReady <- struct{}{}:
case <-ctx.Done():
return ctx.Err()
}
}
select {
case <-d.schedule.resume:
case <-ctx.Done():
return ctx.Err()
}
if d.index >= 9 {
if err := d.StorageAPI.WalkDir(ctx, opts, &stream); err != nil {
return err
}
}
d.schedule.mu.Lock()
d.schedule.snapshots[d.index] = bytes.Clone(stream.Bytes())
d.schedule.mu.Unlock()
_, err := out.Write(stream.Bytes())
return err
}
func nullQuorumBackend(t *testing.T) (*erasureServerPools, string, http.Handler) {
t.Helper()
ctx, cancel := context.WithCancel(t.Context())
obj, dirs, err := prepareErasure16(ctx)
if err != nil {
cancel()
t.Fatal(err)
}
z := obj.(*erasureServerPools)
previous := newObjectLayerFn()
setObjectLayer(z)
t.Cleanup(func() {
cancel()
z.Shutdown(context.Background())
removeRoots(dirs)
setObjectLayer(previous)
})
bucket, router, err := initAPIHandlerTest(ctx, z, nil, MakeBucketOptions{})
if err != nil {
t.Fatal(err)
}
return z, bucket, router
}
func nullQuorumRequest(t *testing.T, router http.Handler, method, target string) *httptest.ResponseRecorder {
t.Helper()
req, err := newTestSignedRequestV4(method, target, 0, nil, globalActiveCred.AccessKey, globalActiveCred.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
router.ServeHTTP(rec, req.WithContext(t.Context()))
return rec
}
func TestListObjectsSingleNullQuorumHTTP(t *testing.T) {
z, bucket, router := nullQuorumBackend(t)
oldBody := bytes.Repeat([]byte{'a'}, 8192)
newBody := bytes.Repeat([]byte{'b'}, 8192)
oldTime := time.Now().UTC().Add(-time.Hour)
newTime := oldTime.Add(time.Minute)
var want []string
for i := range 16 {
name := fmt.Sprintf("fixed/%02d", i)
want = append(want, name)
_, err := z.PutObject(t.Context(), bucket, name, mustGetPutObjReader(t, bytes.NewReader(oldBody), int64(len(oldBody)), "", ""), ObjectOptions{MTime: oldTime})
if err != nil {
t.Fatal(err)
}
}
const overwritten = "fixed/10"
schedule := &nullQuorumWalkSchedule{bucket: bucket, oldReady: make(chan struct{}, 16), resume: make(chan struct{})}
set := z.serverPools[0].sets[0]
disks := set.getDisks()
wrapped := make([]StorageAPI, len(disks))
for i, disk := range disks {
wrapped[i] = &nullQuorumWalkDisk{StorageAPI: disk, schedule: schedule, index: i}
}
set.getDisks = func() []StorageAPI { return wrapped }
if globalAPIConfig.getListQuorum() != "strict" {
t.Fatal("test requires strict listing across all sixteen disks")
}
writeDone := make(chan error, 1)
putReader := mustGetPutObjReader(t, bytes.NewReader(newBody), int64(len(newBody)), "", "")
go func() {
defer close(schedule.resume)
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
defer cancel()
for range 9 {
select {
case <-schedule.oldReady:
case <-ctx.Done():
writeDone <- ctx.Err()
return
}
}
_, err := z.PutObject(ctx, bucket, overwritten, putReader, ObjectOptions{MTime: newTime})
writeDone <- err
}()
rec := nullQuorumRequest(t, router, http.MethodGet, getListObjectsV2URL("", bucket, "fixed/", "1000", "", "", ""))
if err := <-writeDone; err != nil {
t.Fatal(err)
}
var list ListObjectsV2Response
if rec.Code != http.StatusOK || xml.Unmarshal(rec.Body.Bytes(), &list) != nil {
t.Fatalf("LIST: %d %s", rec.Code, rec.Body.String())
}
var got []string
for _, object := range list.Contents {
got = append(got, object.Key)
}
t.Logf("LIST HTTP=%d KeyCount=%d IsTruncated=%v keys=%v", rec.Code, list.KeyCount, list.IsTruncated, got)
if list.KeyCount != 16 || list.IsTruncated || !slices.Equal(got, want) {
t.Errorf("LIST omitted a name during overwrite: got %v, want %v", got, want)
}
// Verify the actual reader inputs, including shape and EC, instead of
// assuming the scheduling barrier produced the intended resolver case.
schedule.mu.Lock()
snapshots := schedule.snapshots
schedule.mu.Unlock()
for i, snapshot := range snapshots {
reader := newMetacacheReader(bytes.NewReader(snapshot))
var found bool
for {
entry, err := reader.next()
if err == io.EOF {
break
}
if err != nil {
t.Fatal(err)
}
if entry.name != overwritten {
continue
}
xl, err := entry.xlmeta()
if err != nil || len(xl.versions) != 1 {
t.Fatalf("disk %d: invalid input shape: %v", i, err)
}
header := xl.versions[0].header
wantTime := oldTime
if i >= 9 {
wantTime = newTime
}
if header.ModTime != wantTime.UnixNano() || header.VersionID != [16]byte{} || header.Type != ObjectType || header.FreeVersion() {
t.Fatalf("disk %d: unexpected header %v", i, header)
}
t.Logf("disk=%d versions=1 header=%v", i, header)
found = true
}
reader.Close()
if !found {
t.Fatalf("disk %d did not emit the overwritten name", i)
}
}
get := nullQuorumRequest(t, router, http.MethodGet, getGetObjectURL("", bucket, overwritten))
head := nullQuorumRequest(t, router, http.MethodHead, getGetObjectURL("", bucket, overwritten))
t.Logf("completed PUT readback: GET=%d bytes=%d HEAD=%d length=%s", get.Code, get.Body.Len(), head.Code, head.Header().Get("Content-Length"))
if get.Code != http.StatusOK || !bytes.Equal(get.Body.Bytes(), newBody) || head.Code != http.StatusOK || head.Header().Get("Content-Length") != "8192" {
t.Fatalf("completed PUT was not readable: GET=%d HEAD=%d", get.Code, head.Code)
}
}
func TestSingleNullQuorumReadAndConsumers(t *testing.T) {
z, bucket, router := nullQuorumBackend(t)
set := z.serverPools[0].sets[0]
disks := set.getDisks()
const object = "object"
oldBody, newBody := bytes.Repeat([]byte{'a'}, 8192), bytes.Repeat([]byte{'b'}, 8192)
oldTime := time.Now().UTC().Add(-time.Hour)
put := func(body []byte, modTime time.Time) {
t.Helper()
_, err := z.PutObject(t.Context(), bucket, object, mustGetPutObjReader(t, bytes.NewReader(body), int64(len(body)), "", ""), ObjectOptions{MTime: modTime})
if err != nil {
t.Fatal(err)
}
}
put(oldBody, oldTime)
oldMeta := make([][]byte, len(disks))
for i, disk := range disks {
var err error
oldMeta[i], err = disk.ReadAll(t.Context(), bucket, object+"/"+xlStorageFormatFile)
if err != nil {
t.Fatal(err)
}
}
put(newBody, oldTime.Add(time.Minute))
newMinorityMeta := mustReadNullQuorumMeta(t, disks[14], bucket, object)
// Model an incomplete overwrite using real inline shards: fifteen old
// copies retain a read quorum; the final disk contains the newer minority.
for i := range 15 {
if err := disks[i].WriteAll(t.Context(), bucket, object+"/"+xlStorageFormatFile, oldMeta[i]); err != nil {
t.Fatal(err)
}
}
checkRead := func(body []byte, modTime time.Time) {
t.Helper()
get := nullQuorumRequest(t, router, http.MethodGet, getGetObjectURL("", bucket, object))
head := nullQuorumRequest(t, router, http.MethodHead, getGetObjectURL("", bucket, object))
if get.Code != 200 || !bytes.Equal(get.Body.Bytes(), body) || head.Code != 200 || head.Header().Get("Content-Length") != "8192" {
t.Fatalf("GET/HEAD lost readable quorum: GET=%d HEAD=%d body=%q", get.Code, head.Code, get.Body.String())
}
oi, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{})
if err != nil || !oi.ModTime.Equal(modTime) {
t.Fatalf("wrong readable generation: %v, %v", oi.ModTime, err)
}
}
checkRead(oldBody, oldTime)
// Reading must not reinterpret the minority as absent or delete it.
minority, err := disks[15].ReadVersion(t.Context(), "", bucket, object, "", ReadOptions{})
if err != nil || !minority.ModTime.Equal(oldTime.Add(time.Minute)) {
t.Fatalf("read mutated the minority: %+v, %v", minority, err)
}
for _, kind := range []string{"rebalance", "decommission"} {
t.Run(kind, func(t *testing.T) {
var names []string
consume := func(entry metaCacheEntry) {
versions, err := entry.fileInfoVersions(bucket)
if err != nil || len(versions.Versions) != 1 || !versions.Versions[0].ModTime.Equal(oldTime) {
t.Errorf("invalid migration metadata: %+v, %v", versions, err)
return
}
// Migration reopens the selected name/version on all drives;
// it must not treat the selected listing metadata as disk agreement.
reader, err := set.GetObjectNInfo(t.Context(), bucket, entry.name, nil, nil, ObjectOptions{VersionID: nullVersionID, NoLock: true, NoDecryption: true})
if err != nil {
t.Error(err)
return
}
body, err := io.ReadAll(reader)
reader.Close()
if err != nil || !bytes.Equal(body, oldBody) {
t.Errorf("migration could not read selected object data: %v", err)
}
names = append(names, entry.name)
}
var err error
if kind == "rebalance" {
err = set.listObjectsToRebalance(t.Context(), bucket, consume)
} else {
err = set.listObjectsToDecommission(t.Context(), decomBucketInfo{Name: bucket}, consume)
}
if err != nil || !slices.Equal(names, []string{object}) {
t.Fatalf("migration listing dropped the quorum: %v, %v", names, err)
}
})
}
for _, latestOnly := range []bool{false, true} {
results := make(chan itemOrErr[ObjectInfo], 16)
if err := z.Walk(t.Context(), bucket, "", results, WalkOptions{AskDisks: "strict", LatestOnly: latestOnly}); err != nil {
t.Fatal(err)
}
var names []string
for result := range results {
if result.Err != nil || !result.Item.ModTime.Equal(oldTime) {
t.Fatalf("Walk returned invalid metadata: %+v", result)
}
names = append(names, result.Item.Name)
}
if !slices.Equal(names, []string{object}) {
t.Fatalf("Walk dropped the quorum: %v", names)
}
}
// Exercise the scanner's actual abandoned-child path. Make the scan drive
// miss the object, while keeping an older quorum and a newer minority on
// the remaining disks. Its queued heal must read all disks and repair both.
if err := disks[14].WriteAll(t.Context(), bucket, object+"/"+xlStorageFormatFile, newMinorityMeta); err != nil {
t.Fatal(err)
}
if err := os.RemoveAll(filepath.Join(disks[15].Endpoint().Path, bucket, object)); err != nil {
t.Fatal(err)
}
scanNullQuorumAbandoned(t, z, bucket, object, oldTime)
for i, disk := range disks {
fi, err := disk.ReadVersion(t.Context(), "", bucket, object, "", ReadOptions{Healing: true})
if err != nil || !fi.ModTime.Equal(oldTime) {
t.Fatalf("scanner heal did not reconcile disk %d: %v, %v", i, fi.ModTime, err)
}
}
checkRead(oldBody, oldTime)
put(newBody, oldTime.Add(2*time.Minute))
checkRead(newBody, oldTime.Add(2*time.Minute))
}
func mustReadNullQuorumMeta(t *testing.T, disk StorageAPI, bucket, object string) []byte {
t.Helper()
data, err := disk.ReadAll(t.Context(), bucket, object+"/"+xlStorageFormatFile)
if err != nil {
t.Fatal(err)
}
return data
}
func scanNullQuorumAbandoned(t *testing.T, z *erasureServerPools, bucket, object string, oldTime time.Time) {
t.Helper()
ctx, cancel := context.WithCancel(t.Context())
defer cancel()
disks := z.serverPools[0].sets[0].getDisks()
seq := newBgHealSequence()
defer seq.cancelCtx()
previousState, previousRoutine := globalBackgroundHealState, globalBackgroundHealRoutine
globalBackgroundHealState = &allHealState{healSeqMap: map[string]*healSequence{"test": seq}}
routine := &healRoutine{tasks: make(chan healTask)}
globalBackgroundHealRoutine = routine
defer func() { globalBackgroundHealState, globalBackgroundHealRoutine = previousState, previousRoutine }()
workerDone := make(chan struct{})
var healed []string
go func() {
defer close(workerDone)
for {
select {
case <-ctx.Done():
return
case task := <-routine.tasks:
var result madmin.HealResultItem
var err error
if task.object == "" {
result, err = z.HealBucket(ctx, task.bucket, task.opts)
} else {
healed = append(healed, task.object+"/"+task.versionID)
result, err = z.HealObject(ctx, task.bucket, task.object, task.versionID, task.opts)
}
task.respCh <- healResult{result: result, err: err}
}
}
}()
cache := dataUsageCache{Info: dataUsageCacheInfo{Name: bucket}}
cache.replace(bucket, "", dataUsageEntry{})
cache.replace(bucket+"/"+object, bucket, dataUsageEntry{Size: 8192, Objects: 1, Versions: 1})
scanner := folderScanner{
root: disks[15].Endpoint().Path, oldCache: cache,
newCache: dataUsageCache{Info: cache.Info}, updateCache: dataUsageCache{Info: cache.Info},
disks: disks, disksQuorum: 8, healObjectSelect: 1,
weSleep: func() bool { return false }, shouldHeal: func() bool { return true }, updateCurrentPath: func(string) {},
getSize: func(item scannerItem) (sizeSummary, error) {
if filepath.Base(item.Path) != xlStorageFormatFile {
return sizeSummary{}, errSkipFile
}
data, err := os.ReadFile(item.Path)
if err != nil {
return sizeSummary{}, err
}
var xl xlMetaV2
if err := xl.Load(data); err != nil {
return sizeSummary{}, err
}
fi, err := xl.ToFileInfo(bucket, object, "", false, false)
if err != nil || !fi.ModTime.Equal(oldTime) {
return sizeSummary{}, fmt.Errorf("scanner re-read wrong generation: %v, %v", fi.ModTime, err)
}
return sizeSummary{totalSize: fi.Size, versions: 1}, nil
},
}
var usage dataUsageEntry
err := scanner.scanFolder(ctx, cachedFolder{name: bucket, objectHealProbDiv: 1}, &usage)
cancel()
<-workerDone
if err != nil || len(healed) != 1 || healed[0] != object+"/"+nullVersionID && healed[0] != object+"/" {
t.Fatalf("scanner did not queue a name-based heal: %v, %v", healed, err)
}
flat := scanner.newCache.sizeRecursive(bucket)
if flat == nil || flat.Objects != 1 || flat.Size != 8192 {
t.Fatalf("scanner lost the healed object from usage: %+v", flat)
}
t.Logf("scanner queued %v; healed old quorum, repaired minority/missing copies, and retained one 8192-byte object", healed)
}
+41 -4
View File
@@ -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()
+264
View File
@@ -0,0 +1,264 @@
// 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 <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"context"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio-go/v7"
"github.com/minio/minio/internal/auth"
"github.com/minio/minio/internal/bucket/replication"
xhttp "github.com/minio/minio/internal/http"
"github.com/minio/minio/internal/once"
)
type staleCreationCase struct {
name string
purged, pendingPurge, unreadable bool
}
// A queued delete-marker creation must be checked against the source marker
// under the replication lock before it is sent. Between queueing (DELETE
// handler, GET/HEAD/LIST heal, scanner, MRF) and sending, a user purge of the
// same marker can complete and reach the targets; the stale creation would
// then recreate the marker there.
func TestReplicationDeleteMarkerCreationRevalidated(t *testing.T) {
for _, tc := range []staleCreationCase{
{name: "current"},
{name: "purged", purged: true},
{name: "purge-pending", pendingPurge: true},
{name: "unreadable", unreadable: true},
} {
t.Run(tc.name, func(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, endpoints: []string{"DeleteObject"}, objAPITest: func(obj ObjectLayer, backend, bucket string, router http.Handler, creds auth.Credentials, t *testing.T) {
testReplicationDeleteMarkerCreationRevalidated(t, obj, backend, bucket, router, creds, tc)
}})
})
}
}
type staleCreationTarget struct {
arn, bucket string
creations atomic.Int32
purges atomic.Int32
}
func testReplicationDeleteMarkerCreationRevalidated(t *testing.T, obj ObjectLayer, backend, bucket string, router http.Handler, creds auth.Credentials, tc staleCreationCase) {
ctx, cancel := context.WithCancel(t.Context())
defer cancel()
defer replicationTestCapacity(obj)()
stats := NewReplicationStats(ctx, nil)
oldStats := globalReplicationStats.Swap(stats)
defer globalReplicationStats.Store(oldStats)
oldPool := globalReplicationPool
defer func() { globalReplicationPool = oldPool }()
const name = "marker"
if _, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, enabledBucketVersioningConfig); err != nil {
t.Fatal(err)
}
if _, err := obj.PutObject(ctx, bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true}); err != nil {
t.Fatal(err)
}
target := &staleCreationTarget{arn: "arn:minio:replication::" + mustGetUUID() + ":bucket", bucket: getRandomBucketName()}
if err := obj.MakeBucket(ctx, target.bucket, MakeBucketOptions{VersioningEnabled: true}); err != nil {
t.Fatal(err)
}
if _, err := obj.PutObject(ctx, target.bucket, name, mustGetPutObjReader(t, bytes.NewReader([]byte("data")), 4, "", ""), ObjectOptions{Versioned: true}); err != nil {
t.Fatal(err)
}
// The remote behaves like a real receiver: a marker lookup and a
// replicated DELETE applied to the target bucket on the same fixture.
remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
opts := ObjectOptions{VersionID: r.URL.Query().Get("versionId"), Versioned: true}
switch r.Method {
case http.MethodHead:
oi, err := obj.GetObjectInfo(r.Context(), target.bucket, name, opts)
if oi.DeleteMarker {
w.Header().Set(xhttp.AmzDeleteMarker, "true")
w.Header().Set(xhttp.AmzVersionID, oi.VersionID)
}
if err != nil {
writeErrorResponseHeadersOnly(w, toAPIError(r.Context(), err))
return
}
w.WriteHeader(http.StatusOK)
case http.MethodDelete:
opts.DeleteMarker = r.Header.Get(xhttp.MinIOSourceDeleteMarker) == "true"
if opts.DeleteMarker {
target.creations.Add(1)
} else {
target.purges.Add(1)
}
opts.ReplicationRequest = true
opts.SetReplicaStatus(replication.Replica)
if _, err := obj.DeleteObject(r.Context(), target.bucket, name, opts); err != nil && !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
writeErrorResponse(r.Context(), w, toAPIError(r.Context(), err), r.URL)
return
}
w.WriteHeader(http.StatusNoContent)
default:
t.Errorf("unexpected remote method %s", r.Method)
}
}))
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)
}
globalBucketTargetSys.Lock()
globalBucketTargetSys.arnRemotesMap[target.arn] = arnTarget{Client: &TargetClient{Client: client, ARN: target.arn, Bucket: target.bucket}, lastRefresh: UTCNow()}
globalBucketTargetSys.targetsMap[bucket] = append(globalBucketTargetSys.targetsMap[bucket], madmin.BucketTarget{Arn: target.arn, TargetBucket: target.bucket})
globalBucketTargetSys.Unlock()
globalBucketTargetSys.hMutex.Lock()
globalBucketTargetSys.hc[client.EndpointURL().Host] = epHealth{Online: true}
globalBucketTargetSys.hMutex.Unlock()
rule := configs[0].Rules[0]
rule.Destination = replication.Destination{ARN: target.arn, Bucket: target.bucket}
cfg := replication.Config{RoleArn: target.arn, Rules: []replication.Rule{rule}}
meta, err := globalBucketMetadataSys.Get(bucket)
if err != nil {
t.Fatal(err)
}
meta.replicationConfig = &cfg
globalBucketMetadataSys.Set(bucket, meta)
p := &ReplicationPool{ctx: ctx, objLayer: obj, workers: []chan ReplicationWorkerOperation{make(chan ReplicationWorkerOperation, 8)}, stats: stats, mrfSaveCh: make(chan MRFReplicateEntry, 8)}
globalReplicationPool = once.NewSingleton[ReplicationPool]()
globalReplicationPool.Set(p)
// The real DELETE handler creates the pending marker and queues the
// creation task that a worker would later run.
req, err := newTestSignedRequestV4(http.MethodDelete, "/"+bucket+"/"+name, 0, nil, creds.AccessKey, creds.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
w := httptest.NewRecorder()
router.ServeHTTP(w, req)
if w.Code != http.StatusNoContent {
t.Fatalf("source DELETE status %d: %s", w.Code, w.Body)
}
var creation DeletedObjectReplicationInfo
select {
case op := <-p.workers[0]:
creation = op.(DeletedObjectReplicationInfo)
case <-time.After(3 * time.Second):
t.Fatal("handler queued no creation task")
}
version := creation.DeleteMarkerVersionID
if version == "" || creation.VersionID != "" || !creation.DeleteMarker {
t.Fatalf("handler queued a non-creation task: %+v", creation)
}
before, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true})
if !isErrMethodNotAllowed(err) || !before.DeleteMarker || before.ReplicationStatus != replication.Pending {
t.Fatalf("source marker not pending: %+v %v", before, err)
}
// Meanwhile the user purges the marker through another path. The purge is
// either complete (version gone) or still pending on the source.
switch {
case tc.purged:
if _, err := obj.DeleteObject(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true}); err != nil {
t.Fatal(err)
}
if _, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true}); !isErrVersionNotFound(err) && !isErrObjectNotFound(err) {
t.Fatalf("source purge left the marker: %v", err)
}
case tc.pendingPurge:
opts := ObjectOptions{VersionID: version, Versioned: true, DeleteReplication: ReplicationState{
ReplicateDecisionStr: creation.ReplicationState.ReplicateDecisionStr,
VersionPurgeStatusInternal: target.arn + "=PENDING;",
PurgeTargets: map[string]VersionPurgeStatusType{target.arn: replication.VersionPurgePending},
}}
if _, err := obj.DeleteObject(ctx, bucket, name, opts); err != nil {
t.Fatal(err)
}
oi, err := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true})
if !isErrMethodNotAllowed(err) || !oi.DeleteMarker || oi.VersionPurgeStatus != replication.VersionPurgePending {
t.Fatalf("source marker not pending purge: %+v %v", oi, err)
}
}
source := obj
if tc.unreadable {
source = staleLookupLayer{ObjectLayer: obj}
}
pending, _ := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true})
result := replicateDelete(ctx, creation, source)
after, aerr := obj.GetObjectInfo(ctx, bucket, name, ObjectOptions{VersionID: version, Versioned: true})
_, terr := obj.GetObjectInfo(ctx, target.bucket, name, ObjectOptions{VersionID: version, Versioned: true})
targetHasMarker := isErrMethodNotAllowed(terr)
switch {
case tc.purged, tc.pendingPurge, tc.unreadable:
if len(result.Targets) != 0 || target.creations.Load() != 0 || target.purges.Load() != 0 {
t.Fatalf("%s: stale creation was sent: result=%+v creations=%d purges=%d", tc.name, result, target.creations.Load(), target.purges.Load())
}
if targetHasMarker {
t.Fatalf("%s: target marker recreated from the stale creation", tc.name)
}
if tc.purged {
if !isErrVersionNotFound(aerr) && !isErrObjectNotFound(aerr) {
t.Fatalf("stale creation resurrected the source marker: %+v %v", after, aerr)
}
} else if !isErrMethodNotAllowed(aerr) || !after.DeleteMarker ||
after.ReplicationStatusInternal != pending.ReplicationStatusInternal ||
after.VersionPurgeStatusInternal != pending.VersionPurgeStatusInternal ||
after.UserDefined[ReservedMetadataPrefixLower+ReplicationTimestamp] != pending.UserDefined[ReservedMetadataPrefixLower+ReplicationTimestamp] {
t.Fatalf("skipped creation rewrote the source marker: before=%+v after=%+v %v", pending, after, aerr)
}
if tc.unreadable {
select {
case entry := <-p.mrfSaveCh:
if entry.versionID != version || entry.RetryCount != 1 || entry.Bucket != bucket || entry.Object != name {
t.Fatalf("unverified creation queued wrong MRF entry: %+v", entry)
}
default:
t.Fatal("unverified source read did not queue a retry")
}
}
if len(p.mrfSaveCh) != 0 {
t.Fatalf("%s: unexpected MRF entries queued: %d", tc.name, len(p.mrfSaveCh))
}
default:
if result.ReplicationStatus() != replication.Completed || target.creations.Load() != 1 || !targetHasMarker {
t.Fatalf("current creation not replicated: result=%+v creations=%d targetMarker=%v", result, target.creations.Load(), targetHasMarker)
}
if !isErrMethodNotAllowed(aerr) || !after.DeleteMarker || replicationStatusesMap(after.ReplicationStatusInternal)[target.arn] != replication.Completed {
t.Fatalf("source creation state not completed: %+v %v", after, aerr)
}
if len(p.mrfSaveCh) != 0 {
t.Fatal("completed creation queued MRF work")
}
}
t.Logf("%s: %s checked; creations=%d purges=%d", backend, tc.name, target.creations.Load(), target.purges.Load())
}
// staleLookupLayer cannot confirm the source marker: the read fails without
// saying whether the version is present.
type staleLookupLayer struct{ ObjectLayer }
func (staleLookupLayer) GetObjectInfo(context.Context, string, string, ObjectOptions) (ObjectInfo, error) {
return ObjectInfo{}, InsufficientReadQuorum{}
}
@@ -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 <http://www.gnu.org/licenses/>.
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)
}
}
}
+67
View File
@@ -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 <http://www.gnu.org/licenses/>.
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
}
+274
View File
@@ -0,0 +1,274 @@
// 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 <http://www.gnu.org/licenses/>.
package cmd
import (
"fmt"
"math/rand"
"reflect"
"slices"
"testing"
"time"
)
// Use complete metadata so the selected version can also be decoded by LIST
// and the object read path, rather than testing only synthetic shallow headers.
func quorumNullVersion(t testing.TB, generation int64) xlMetaV2ShallowVersion {
t.Helper()
fi := newFileInfo("fixed/object", 12, 4)
fi.Erasure.Index = 1
fi.DataDir = "11111111-1111-1111-1111-111111111111"
fi.ModTime = time.Unix(generation, 0).UTC()
fi.Size = 8192
fi.Parts = []ObjectPartInfo{{Number: 1, Size: fi.Size, ActualSize: fi.Size}}
fi.Metadata = map[string]string{"etag": fmt.Sprintf("%032d", generation)}
var xl xlMetaV2
if err := xl.AddVersion(fi); err != nil {
t.Fatal(err)
}
return xl.versions[0]
}
// Start with forward, reverse and interleaved orders, then fixed-seed shuffles.
func quorumVersionOrders(input [][]xlMetaV2ShallowVersion) [][][]xlMetaV2ShallowVersion {
orders := [][][]xlMetaV2ShallowVersion{slices.Clone(input), slices.Clone(input)}
slices.Reverse(orders[1])
interleaved := make([][]xlMetaV2ShallowVersion, 0, len(input))
for i, j := 0, len(input)-1; i <= j; i, j = i+1, j-1 {
interleaved = append(interleaved, input[i])
if i != j {
interleaved = append(interleaved, input[j])
}
}
orders = append(orders, interleaved)
for seed := range int64(32) {
order := slices.Clone(input)
rand.New(rand.NewSource(seed)).Shuffle(len(order), func(i, j int) {
order[i], order[j] = order[j], order[i]
})
orders = append(orders, order)
}
return orders
}
func TestMergeXLV2SingleNullQuorum(t *testing.T) {
v := []xlMetaV2ShallowVersion{quorumNullVersion(t, 100), quorumNullVersion(t, 200), quorumNullVersion(t, 300)}
for _, tt := range []struct {
name string
counts [3]int
empty int
quorum int
want int
}{
{"consistent", [3]int{16}, 0, 8, 0},
{"old9-new7", [3]int{9, 7}, 0, 8, 0},
{"old15-new1", [3]int{15, 1}, 0, 8, 0},
{"old3-new1", [3]int{3, 1}, 0, 3, 0},
{"two-quorums", [3]int{8, 8}, 0, 8, 1},
{"empty-streams", [3]int{9, 1}, 6, 8, 0},
{"two-subquorums", [3]int{7, 7}, 2, 8, -1},
{"three-subquorums", [3]int{6, 5, 5}, 0, 8, -1},
} {
t.Run(tt.name, func(t *testing.T) {
var input [][]xlMetaV2ShallowVersion
for g, n := range tt.counts {
for range n {
input = append(input, []xlMetaV2ShallowVersion{v[g]})
}
}
input = append(input, make([][]xlMetaV2ShallowVersion, tt.empty)...)
want := []xlMetaV2ShallowVersion{}
if tt.want >= 0 {
want = append(want, v[tt.want])
}
for order, versions := range quorumVersionOrders(input) {
for _, strict := range []bool{false, true} {
for _, requested := range []int{0, 1} {
got := mergeXLV2Versions(tt.quorum, strict, requested, versions...)
if !reflect.DeepEqual(got, want) {
t.Fatalf("order=%d strict=%v requested=%d: got %#v, want %#v", order, strict, requested, got, want)
}
}
}
}
})
}
}
func TestMergeXLV2SingleNullHeaderGroups(t *testing.T) {
old := quorumNullVersion(t, 100)
newer := quorumNullVersion(t, 200)
differentSig := old
differentSig.header.Signature[0]++
differentFlags := old
differentFlags.header.Flags ^= xlFlagInlineData
for _, tt := range []struct {
name string
input [][]xlMetaV2ShallowVersion
strict bool
want []xlMetaV2ShallowVersion
}{
{"signature-nonstrict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentSig}, {newer}}, false, []xlMetaV2ShallowVersion{differentSig}},
{"signature-strict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentSig}, {newer}}, true, []xlMetaV2ShallowVersion{}},
{"flags-nonstrict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentFlags}, {newer}}, false, []xlMetaV2ShallowVersion{}},
{"flags-strict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentFlags}, {newer}}, true, []xlMetaV2ShallowVersion{}},
} {
t.Run(tt.name, func(t *testing.T) {
got := mergeXLV2Versions(3, tt.strict, 0, tt.input...)
if !reflect.DeepEqual(got, tt.want) {
t.Fatalf("got %#v, want %#v", got, tt.want)
}
})
}
}
func TestResolveSingleNullQuorum(t *testing.T) {
old, newer := quorumNullVersion(t, 100), quorumNullVersion(t, 200)
input := make([][]xlMetaV2ShallowVersion, 16)
for i := range input {
input[i] = []xlMetaV2ShallowVersion{old}
if i >= 9 {
input[i] = []xlMetaV2ShallowVersion{newer}
}
}
for order, versions := range quorumVersionOrders(input) {
entries := make(metaCacheEntries, len(versions))
for i := range versions {
xl := &xlMetaV2{versions: versions[i]}
metadata, err := xl.AppendTo(nil)
if err != nil {
t.Fatal(err)
}
entries[i] = metaCacheEntry{name: "fixed/object", metadata: metadata}
}
selected, ok := entries.resolve(&metadataResolutionParams{objQuorum: 8, dirQuorum: 8, requestedVersions: 1})
if !ok {
t.Fatalf("order %d: a complete null object version has quorum but was omitted", order)
}
xl, err := selected.xlmeta()
if err != nil {
t.Fatal(err)
}
if !reflect.DeepEqual(xl.versions, []xlMetaV2ShallowVersion{old}) {
t.Fatalf("order %d: incorrect complete selected version: %#v", order, xl.versions)
}
fi, err := xl.ToFileInfo("bucket", "fixed/object", "", false, false)
if err != nil || !fi.IsValid() || fi.Size != 8192 || !fi.ModTime.Equal(time.Unix(100, 0)) {
t.Fatalf("order %d: selected metadata cannot be read: %+v, %v", order, fi, err)
}
}
}
// These are compatibility assertions for inputs outside the narrow repair.
// In particular, the mixed-version outputs below are not a new correctness
// contract for version enumeration; that behavior requires a separate change.
func TestMergeXLV2NullHistoriesUnchanged(t *testing.T) {
old, newer := quorumNullVersion(t, 100), quorumNullVersion(t, 300)
middle := quorumNullVersion(t, 200)
middle.header.VersionID = [16]byte{1}
highest := quorumNullVersion(t, 400)
highest.header.VersionID = [16]byte{2}
for _, tt := range []struct {
name string
input [][]xlMetaV2ShallowVersion
want []xlMetaV2ShallowVersion
}{
{"mixed-forward", [][]xlMetaV2ShallowVersion{{old}, {old}, {old}, {newer, middle}, {middle}, {middle}}, []xlMetaV2ShallowVersion{middle}},
{"mixed-reverse", [][]xlMetaV2ShallowVersion{{middle}, {middle}, {newer, middle}, {old}, {old}, {old}}, []xlMetaV2ShallowVersion{old}},
{"history-pruned-first", [][]xlMetaV2ShallowVersion{{highest, old}, {highest, old}, {highest, old}, {highest, newer}}, []xlMetaV2ShallowVersion{highest}},
} {
t.Run(tt.name, func(t *testing.T) {
got := mergeXLV2Versions(3, false, 0, tt.input...)
if !reflect.DeepEqual(got, tt.want) {
t.Fatalf("excluded history changed: got %#v, want %#v", got, tt.want)
}
})
}
for _, kind := range []string{"delete", "free", "legacy", "mixed-ec", "legacy-modern-ec", "nonzero-strict"} {
t.Run(kind, func(t *testing.T) {
a, b := old, newer
switch kind {
case "delete":
b.header.Type = DeleteType
case "free":
b.header.Flags |= xlFlagFreeVersion
case "legacy":
b.header.Type = LegacyType
case "mixed-ec":
b.header.EcN, b.header.EcM = 8, 8
case "legacy-modern-ec":
b.header.EcN, b.header.EcM = 0, 0
case "nonzero-strict":
a.header.VersionID, b.header.VersionID = [16]byte{1}, [16]byte{1}
}
for _, requested := range []int{0, 1} {
got := mergeXLV2Versions(3, true, requested, []xlMetaV2ShallowVersion{a}, []xlMetaV2ShallowVersion{a}, []xlMetaV2ShallowVersion{a}, []xlMetaV2ShallowVersion{b})
if len(got) != 0 {
t.Fatalf("excluded input changed: requested=%d got %#v", requested, got)
}
}
})
}
// All-legacy EC headers are eligible, unlike mixing legacy and modern EC.
old.header.EcN, old.header.EcM = 0, 0
newer.header.EcN, newer.header.EcM = 0, 0
got := mergeXLV2Versions(3, false, 0, []xlMetaV2ShallowVersion{old}, []xlMetaV2ShallowVersion{old}, []xlMetaV2ShallowVersion{old}, []xlMetaV2ShallowVersion{newer})
if !reflect.DeepEqual(got, []xlMetaV2ShallowVersion{old}) {
t.Fatalf("all-legacy EC lost the old quorum: %#v", got)
}
}
func BenchmarkResolveSingleNullQuorumPage(b *testing.B) {
old, newer := quorumNullVersion(b, 100), quorumNullVersion(b, 200)
for _, kind := range []string{"consistent", "old-first", "new-first"} {
b.Run(kind, func(b *testing.B) {
var templates [16]metaCacheEntry
for i := range templates {
v := old
if kind == "old-first" && i >= 9 || kind == "new-first" && i < 7 {
v = newer
}
xl := &xlMetaV2{versions: []xlMetaV2ShallowVersion{v}}
data, err := xl.AppendTo(nil)
if err != nil {
b.Fatal(err)
}
templates[i] = metaCacheEntry{metadata: data}
}
var names [1000]string
for i := range names {
names[i] = fmt.Sprintf("fixed/%04d", i)
}
resolver := metadataResolutionParams{objQuorum: 8, dirQuorum: 8, requestedVersions: 1}
b.ReportAllocs()
b.ResetTimer()
for b.Loop() {
for _, name := range names {
entries := templates
for i := range entries {
entries[i].name = name
}
selected, ok := metaCacheEntries(entries[:]).resolve(&resolver)
if ok && selected.reusable {
metaDataPoolPut(selected.metadata)
}
}
}
})
}
}
+83 -36
View File
@@ -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
@@ -1941,6 +1941,10 @@ func mergeXLV2Versions(quorum int, strict bool, requestedVersions int, versions
// No need for non-strict checks if quorum is 1.
strict = true
}
// Keep the original stream shapes: pruning must not make a versioned object
// eligible for the single null-version recount below.
originalVersions := versions
var checkedSingleNull, singleNull bool
// Shallow copy input
versions = append(make([][]xlMetaV2ShallowVersion, 0, len(versions)), versions...)
@@ -2015,44 +2019,23 @@ func mergeXLV2Versions(quorum int, strict bool, requestedVersions int, versions
// Version IDs match, but otherwise unable to resolve.
// We are either strict, or don't have enough information to match.
// Switch to a pure counting algo.
x := make(map[xlMetaV2VersionHeader]int, len(tops))
for _, a := range tops {
if a.header.VersionID != ver.header.VersionID {
continue
}
if !strict {
// we must match EC, when we are not strict.
if !a.header.matchesEC(ver.header) {
continue
}
a.header.Signature = [4]byte{}
}
x[a.header]++
}
latestCount = 0
for k, v := range x {
if v < latestCount {
continue
}
if v == latestCount && latest.header.sortsBefore(k) {
// Tiebreak, use sort.
continue
}
for _, a := range tops {
hdr := a.header
if !strict {
hdr.Signature = [4]byte{}
}
if hdr == k {
latest = a
}
}
latestCount = v
}
latest, latestCount = countXLV2Versions(tops, ver.header, latest, strict)
break
}
}
if latestCount < quorum {
if !checkedSingleNull {
singleNull = singleNullVersionStreams(originalVersions)
checkedSingleNull = true
}
if singleNull {
// A newer minority at the end can hide an older quorum from
// the selection loop. Recount before discarding the null ID.
if candidate, count := countXLV2Versions(tops, latest.header, latest, strict); count >= quorum {
latest, latestCount = candidate, count
}
}
}
if latestCount >= quorum {
merged = append(merged, latest)
@@ -2113,6 +2096,70 @@ func mergeXLV2Versions(quorum int, strict bool, requestedVersions int, versions
return merged
}
// singleNullVersionStreams excludes histories and other version types from the
// additional recount. Check the original inputs, before any stream is pruned.
func singleNullVersionStreams(versions [][]xlMetaV2ShallowVersion) bool {
var ec xlMetaV2VersionHeader
var haveEC bool
for _, stream := range versions {
if len(stream) == 0 {
continue
}
if len(stream) != 1 {
return false
}
h := stream[0].header
if h.VersionID != [16]byte{} || h.Type != ObjectType || h.FreeVersion() {
return false
}
if haveEC && (h.EcN != ec.EcN || h.EcM != ec.EcM) {
return false
}
ec, haveEC = h, true
}
return haveEC
}
// countXLV2Versions selects the most frequent compatible header for reference's
// VersionID, retaining the existing sort tiebreak and last matching entry.
func countXLV2Versions(tops []xlMetaV2ShallowVersion, reference xlMetaV2VersionHeader, latest xlMetaV2ShallowVersion, strict bool) (xlMetaV2ShallowVersion, int) {
x := make(map[xlMetaV2VersionHeader]int, len(tops))
for _, a := range tops {
if a.header.VersionID != reference.VersionID {
continue
}
if !strict {
// we must match EC, when we are not strict.
if !a.header.matchesEC(reference) {
continue
}
a.header.Signature = [4]byte{}
}
x[a.header]++
}
var latestCount int
for k, v := range x {
if v < latestCount {
continue
}
if v == latestCount && latest.header.sortsBefore(k) {
// Tiebreak, use sort.
continue
}
for _, a := range tops {
hdr := a.header
if !strict {
hdr.Signature = [4]byte{}
}
if hdr == k {
latest = a
}
}
latestCount = v
}
return latest, latestCount
}
type xlMetaBuf []byte
// ToFileInfo converts xlMetaV2 into a common FileInfo datastructure
+3 -1
View File
@@ -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.
+2
View File
@@ -1,5 +1,7 @@
# Compression Guide
For SILO's current SSE-C compression and historical-object boundary, see the [maintained encryption guide](https://silo.pgsty.com/administration/server-side-encryption/server-side-encryption-sse-c/) and [SSE-C replica design](https://silo.pgsty.com/blog/design/ssec-replica-integrity/). Check the documented release boundary before applying main-branch behavior to an older binary.
Silo server allows streaming compression to ensure efficient disk space usage.
Compression happens inflight, i.e objects are compressed before being written to disk(s).
Silo uses [`klauspost/compress/s2`](https://github.com/klauspost/compress/tree/master/s2)
+3 -3
View File
@@ -10,13 +10,13 @@ Silo in distributed mode can help you setup a highly-available storage system wi
Distributed Silo provides protection against multiple node/drive failures and [bit rot](https://github.com/pgsty/silo/blob/main/docs/erasure/README.md#what-is-bit-rot-protection) using [erasure code](https://silo.pgsty.com/operations/concepts/erasure-coding/). As the minimum drives required for distributed Silo is 2 (same as minimum drives required for erasure coding), erasure code automatically kicks in as you launch distributed Silo.
If one or more drives are offline at the start of a PutObject or NewMultipartUpload operation the object will have additional data protection bits added automatically to provide additional safety for these objects.
With the default `availability` storage-class optimization, Silo can increase parity for new objects when drives are offline, up to half the drives in the erasure set. Writes must still meet the quorum for the resulting shard layout. This does not guarantee unchanged availability or keep a two-drive `EC:1` set writable after one drive fails. See the [storage-class reference](https://silo.pgsty.com/reference/minio-server/settings/storage-class/) for details.
### High availability
A stand-alone Silo server would go down if the server hosting the drives goes offline. In contrast, a distributed Silo setup with _m_ servers and _n_ drives will have your data safe as long as _m/2_ servers or _m*n_/2 or more drives are online.
A stand-alone Silo server becomes unavailable if its host goes offline. In a distributed deployment, determine node-failure tolerance from how many drives each failed node removes from every erasure set and from the parity recorded for the affected objects. Aggregate online node or drive counts alone do not establish read or write quorum.
For example, an 16-server distributed setup with 200 drives per node would continue serving files, up to 4 servers can be offline in default configuration i.e around 800 drives down Silo would continue to read and write objects.
For example, if a 16-node pool has 16-drive erasure sets with one drive per node in each set, a fixed `EC:4` layout has a read and write quorum of 12. Losing four nodes leaves 12 drives in every set. This result depends on that drive placement and parity; it cannot be generalized to arbitrary failures of the same total number of drives. In a two-node, two-drive `EC:1` set, losing one node leaves only one drive, so existing intact objects can remain readable but writes cannot continue.
Refer to sizing guide for more understanding on default values chosen depending on your erasure stripe size [here](https://github.com/pgsty/silo/blob/main/docs/distributed/SIZING.md). Parity settings can be changed using [storage classes](https://github.com/pgsty/silo/tree/main/docs/erasure/storage-class).
+6 -6
View File
@@ -1,18 +1,18 @@
# Silo Erasure Code Quickstart Guide
Silo protects data against hardware failures and silent data corruption using erasure code and checksums. With the highest level of redundancy, you may lose up to half (N/2) of the total drives and still be able to recover the data.
Silo protects data against hardware failures and silent data corruption using erasure coding and checksums. Recovery depends on the healthy shards and metadata remaining in each object's erasure set. With maximum parity, an object can tolerate the loss of up to `floor(N/2)` shards in its `N`-drive set; the deployment's total online drive count is not sufficient to determine recoverability.
## What is Erasure Code?
Erasure code is a mathematical algorithm to reconstruct missing or corrupted data. Silo uses Reed-Solomon code to shard objects into variable data and parity blocks. For example, in a 12 drive setup, an object can be sharded to a variable number of data and parity blocks across all the drives - ranging from six data and six parity blocks to ten data and two parity blocks.
Erasure coding uses mathematical algorithms to reconstruct missing or corrupted data. Silo uses Reed-Solomon codes to split objects into data and parity shards. For a 12-drive erasure set, parity can range from `EC:0` (12 data shards and no parity) to `EC:6` (six data shards and six parity shards).
By default, Silo shards the objects across N/2 data and N/2 parity drives. Though, you can use [storage classes](https://github.com/pgsty/silo/tree/main/docs/erasure/storage-class) to use a custom configuration. We recommend N/2 data and parity blocks, as it ensures the best protection from drive failures.
The default parity depends on the erasure set size: `EC:0` for 1 drive, `EC:1` for 2–3 drives, `EC:2` for 4–5 drives, `EC:3` for 6–7 drives, and `EC:4` for 8–16 drives. See the [Silo storage-class reference](https://silo.pgsty.com/reference/minio-server/settings/storage-class/) for configuration details.
In 12 drive example above, with Silo server running in the default configuration, you can lose any of the six drives and still reconstruct the data reliably from the remaining drives.
In the 12-drive example, the default `EC:4` layout uses eight data shards and four parity shards. An existing object can tolerate four unavailable shards if the remaining shards and metadata are intact. Explicit `EC:6` instead uses six data and six parity shards, with a read quorum of six and a write quorum of seven.
## Why is Erasure Code useful?
Erasure code protects data from multiple drives failure, unlike RAID or replication. For example, RAID6 can protect against two drive failure whereas in Silo erasure code you can lose as many as half of drives and still the data remains safe. Further, Silo's erasure code is at the object level and can heal one object at a time. For RAID, healing can be done only at the volume level which translates into high downtime. As Silo encodes each object individually, it can heal objects incrementally. Storage servers once deployed should not require drive replacement or healing for the lifetime of the server. Silo's erasure coded backend is designed for operational efficiency and takes full advantage of hardware acceleration whenever available.
Erasure coding protects objects against drive failures within their erasure set, up to the parity recorded for each object. Silo encodes objects individually and can heal them incrementally when enough healthy shards and consistent metadata remain. Failed drives still need repair or replacement to restore redundancy; erasure coding does not eliminate that maintenance.
![Erasure](https://github.com/pgsty/silo/blob/main/docs/screenshots/erasure-code.jpg?raw=true)
@@ -64,4 +64,4 @@ podman run \
### 3. Test your setup
You may unplug drives randomly and continue to perform I/O on the system.
In a test deployment, take drives offline and verify reads and writes against the quorum required by each erasure set and object layout. Restore the drives after each test; continued I/O depends on the remaining healthy shards and metadata.
+10 -18
View File
@@ -38,31 +38,23 @@ You can calculate _approximate_ storage usage ratio using the formula - total dr
### Allowed values for STANDARD storage class
`STANDARD` storage class implies more parity than `REDUCED_REDUNDANCY` class. So, `STANDARD` parity drives should be
`STANDARD` supports `EC:0` without erasure-code redundancy and nonzero parity values up to `floor(N/2)`, where `N` is the number of drives in the erasure set. When both `STANDARD` and `REDUCED_REDUNDANCY` parity are nonzero, `STANDARD` must be greater than or equal to `REDUCED_REDUNDANCY`; equal parity is allowed.
- Greater than or equal to 2, if `REDUCED_REDUNDANCY` parity is not set.
- Greater than `REDUCED_REDUNDANCY` parity, if it is set.
The default `STANDARD` parity is:
Parity blocks can not be higher than data blocks, so `STANDARD` storage class parity can not be higher than N/2. (N being total number of drives)
The default value for the `STANDARD` storage class depends on the number of volumes in the erasure set:
| Erasure Set Size | Default Parity (EC:N) |
|------------------|-----------------------|
| 5 or fewer | EC:2 |
| 6-7 | EC:3 |
| 8 or more | EC:4 |
| Erasure Set Size | Default Parity (EC:M) |
| --- | --- |
| 1 | EC:0 |
| 2–3 | EC:1 |
| 4–5 | EC:2 |
| 6–7 | EC:3 |
| 8–16 | EC:4 |
For more complete documentation on Erasure Set sizing, see the [Silo Documentation on Erasure Sets](https://silo.pgsty.com/operations/concepts/erasure-coding/#minio-ec-erasure-set).
### Allowed values for REDUCED_REDUNDANCY storage class
`REDUCED_REDUNDANCY` implies lesser parity than `STANDARD` class. So,`REDUCED_REDUNDANCY` parity drives should be
- Less than N/2, if `STANDARD` parity is not set.
- Less than `STANDARD` Parity, if it is set.
Default value for `REDUCED_REDUNDANCY` storage class is `1`.
`REDUCED_REDUNDANCY` parity can be zero or up to `floor(N/2)`. When both storage classes have nonzero parity, its parity must be less than or equal to `STANDARD` parity. The default is `EC:1` for multi-drive sets and `EC:0` for single-drive deployments. See the [Silo storage-class reference](https://silo.pgsty.com/reference/minio-server/settings/storage-class/) for the maintained configuration documentation.
## Get started with Storage Class
+2
View File
@@ -1,5 +1,7 @@
# Federation Quickstart Guide *Federation feature is deprecated and should be avoided for future deployments*
The maintained [federated CopyObject design](https://silo.pgsty.com/blog/design/federated-copy-object/) records the destination encryption, checksum, Object Lock and committed-response contract, including its Server release boundary.
This document explains how to configure Silo with `Bucket lookup from DNS` style federation.
## Cross-deployment copy behavior
+1 -1
View File
@@ -24,7 +24,7 @@ type TargetIDSet map[TargetID]struct{}
// IsEmpty returns true if the set is empty.
func (set TargetIDSet) IsEmpty() bool {
return len(set) != 0
return len(set) == 0
}
// Clone - returns copy of this set.
@@ -22,6 +22,22 @@ import (
"testing"
)
func TestTargetIDSetIsEmpty(t *testing.T) {
testCases := []struct {
set TargetIDSet
expected bool
}{
{NewTargetIDSet(), true},
{NewTargetIDSet(TargetID{"1", "webhook"}), false},
}
for i, testCase := range testCases {
if got := testCase.set.IsEmpty(); got != testCase.expected {
t.Fatalf("test %v: expected: %v, got: %v", i+1, testCase.expected, got)
}
}
}
func TestTargetIDSetClone(t *testing.T) {
testCases := []struct {
set TargetIDSet