mirror of
https://github.com/pgsty/minio.git
synced 2026-09-30 23:05:59 +03:00
Compare commits
13 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 69d5e92791 | |||
| efd5698597 | |||
| 2fde3cf535 | |||
| 2a4d51406b | |||
| 0596685ae7 | |||
| 358ab38fb0 | |||
| 254b19ac07 | |||
| eb4f5e5b31 | |||
| 8d06424b12 | |||
| fced863036 | |||
| 5af0865aab | |||
| 83821f0f1f | |||
| bf053386a4 |
@@ -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
|
||||
|
||||
@@ -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])
|
||||
}
|
||||
}})
|
||||
}
|
||||
@@ -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
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
@@ -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
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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).
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||

|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user