mirror of
https://github.com/pgsty/minio.git
synced 2026-09-08 11:34:04 +03:00
fix: make access tier moves preserve versions and isolate writes
Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
@@ -0,0 +1,524 @@
|
||||
// Copyright 2026 PGSTY contributors.
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/md5"
|
||||
"encoding/base64"
|
||||
"encoding/xml"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"maps"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/minio/minio/internal/dsync"
|
||||
"github.com/minio/minio/internal/hash"
|
||||
xhttp "github.com/minio/minio/internal/http"
|
||||
)
|
||||
|
||||
func accessMovePools(t *testing.T) (*erasureServerPools, string) {
|
||||
t.Helper()
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
dirs, err := getRandomDisks(32)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
endpoints := mustGetPoolEndpoints(0, dirs[:16]...)
|
||||
endpoints = append(endpoints, mustGetPoolEndpoints(1, dirs[16:]...)...)
|
||||
obj, _, err := initObjectLayer(ctx, endpoints)
|
||||
if err != nil {
|
||||
cancel()
|
||||
removeRoots(dirs)
|
||||
t.Fatal(err)
|
||||
}
|
||||
z := obj.(*erasureServerPools)
|
||||
previous := newObjectLayerFn()
|
||||
setObjectLayer(z)
|
||||
t.Cleanup(func() { cancel(); z.Shutdown(context.Background()); removeRoots(dirs); setObjectLayer(previous) })
|
||||
bucket := "access-move-test"
|
||||
if err := z.MakeBucket(ctx, bucket, MakeBucketOptions{VersioningEnabled: true}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return z, bucket
|
||||
}
|
||||
|
||||
func putAccessMoveVersion(t *testing.T, z *erasureServerPools, bucket, object string, pool int, body string, moved bool) ObjectInfo {
|
||||
t.Helper()
|
||||
metadata := map[string]string{"test-value": body}
|
||||
if moved {
|
||||
metadata[accessTierMetadataKey] = accessTierStamp(pool, time.Now().UnixNano())
|
||||
}
|
||||
cs := hash.NewChecksumFromData(hash.ChecksumCRC32C, []byte(body))
|
||||
oi, err := z.serverPools[pool].PutObject(t.Context(), bucket, object,
|
||||
mustGetPutObjReader(t, bytes.NewBufferString(body), int64(len(body)), "", ""),
|
||||
ObjectOptions{Versioned: true, UserDefined: metadata, WantChecksum: cs})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return oi
|
||||
}
|
||||
|
||||
func assertAccessMoveVersion(t *testing.T, z *erasureServerPools, bucket, object string, pool int, want ObjectInfo, body string) {
|
||||
t.Helper()
|
||||
gr, err := z.serverPools[pool].GetObjectNInfo(t.Context(), bucket, object, nil, http.Header{}, ObjectOptions{VersionID: want.VersionID})
|
||||
if err != nil {
|
||||
t.Errorf("pool %d version %s: %v", pool, want.VersionID, err)
|
||||
return
|
||||
}
|
||||
got, err := io.ReadAll(gr)
|
||||
gr.Close()
|
||||
if err != nil || string(got) != body {
|
||||
t.Errorf("pool %d payload = %q, err = %v, want %q", pool, got, err, body)
|
||||
}
|
||||
if gr.ObjInfo.ETag != want.ETag || !gr.ObjInfo.ModTime.Equal(want.ModTime) {
|
||||
t.Error("move changed ETag or modification time")
|
||||
}
|
||||
if len(want.Checksum) == 0 {
|
||||
t.Fatal("fixture has no checksum")
|
||||
}
|
||||
if !bytes.Equal(gr.ObjInfo.Checksum, want.Checksum) {
|
||||
t.Error("move changed or dropped the stored checksum")
|
||||
}
|
||||
if gr.ObjInfo.UserDefined["test-value"] != body {
|
||||
t.Error("move changed user metadata")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAccessMoveVersionStack(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
const object = "versions"
|
||||
old := putAccessMoveVersion(t, z, bucket, object, 1, "old", false)
|
||||
dm, err := z.serverPools[1].DeleteObject(t.Context(), bucket, object, ObjectOptions{Versioned: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
latest := putAccessMoveVersion(t, z, bucket, object, 1, "new", false)
|
||||
for _, pair := range [][2]int{{1, 0}, {0, 1}} {
|
||||
if _, err := moveObjectPool(t.Context(), z, bucket, object, pair[0], pair[1], nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, pair[1], old, "old")
|
||||
assertAccessMoveVersion(t, z, bucket, object, pair[1], latest, "new")
|
||||
versions, err := accessObjectVersions(t.Context(), z, pair[1], bucket, object)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
found := false
|
||||
for _, version := range versions {
|
||||
if version.VersionID == dm.VersionID && version.Deleted {
|
||||
found = true
|
||||
}
|
||||
}
|
||||
if !found || len(versions) != 3 {
|
||||
t.Errorf("move lost a version or delete marker: %+v", versions)
|
||||
}
|
||||
if _, err := z.serverPools[pair[0]].GetObjectInfo(t.Context(), bucket, object, ObjectOptions{}); !isErrObjectNotFound(err) {
|
||||
t.Errorf("source remains after completed move: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAccessMovePreservesNestedObject(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
parent := putAccessMoveVersion(t, z, bucket, "parent", 1, "parent", false)
|
||||
child := putAccessMoveVersion(t, z, bucket, "parent/child", 1, "child", false)
|
||||
if _, err := moveObjectPool(t.Context(), z, bucket, "parent", 1, 0, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, "parent", 0, parent, "parent")
|
||||
assertAccessMoveVersion(t, z, bucket, "parent/child", 1, child, "child")
|
||||
}
|
||||
|
||||
func TestAccessMoveNullVersion(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
const object, body = "null-version", "unversioned"
|
||||
// A null version can remain in a bucket after versioning is enabled.
|
||||
oi, err := z.serverPools[1].PutObject(t.Context(), bucket, object,
|
||||
mustGetPutObjReader(t, bytes.NewBufferString(body), int64(len(body)), "", ""), ObjectOptions{
|
||||
UserDefined: map[string]string{"test-value": body},
|
||||
WantChecksum: hash.NewChecksumFromData(hash.ChecksumCRC32C, []byte(body)),
|
||||
})
|
||||
if err != nil || oi.VersionID != "" {
|
||||
t.Fatalf("null version fixture failed: %v %q", err, oi.VersionID)
|
||||
}
|
||||
if _, err := moveObjectPool(t.Context(), z, bucket, object, 1, 0, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, oi, body)
|
||||
if _, err := z.serverPools[1].GetObjectInfo(t.Context(), bucket, object, ObjectOptions{}); !isErrObjectNotFound(err) {
|
||||
t.Fatalf("null version source was not removed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Source removal can partially commit before a process or disk fails. The
|
||||
// destination can therefore hold the only remaining copy of an older version.
|
||||
func TestAccessMoveRetryPreservesDestinationVersions(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
const object = "retry"
|
||||
onlyAtDestination := putAccessMoveVersion(t, z, bucket, object, 0, "unique-old", true)
|
||||
source := putAccessMoveVersion(t, z, bucket, object, 1, "source-new", false)
|
||||
if _, err := moveObjectPool(t.Context(), z, bucket, object, 1, 0, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, onlyAtDestination, "unique-old")
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, source, "source-new")
|
||||
}
|
||||
|
||||
type accessMoveFaultDisk struct {
|
||||
StorageAPI
|
||||
failVersion string
|
||||
}
|
||||
|
||||
func (d accessMoveFaultDisk) RenameData(ctx context.Context, srcVolume, srcPath string, fi FileInfo, dstVolume, dstPath string, opts RenameOptions) (RenameDataResp, error) {
|
||||
if fi.VersionID == d.failVersion {
|
||||
return RenameDataResp{}, errDiskFull
|
||||
}
|
||||
return d.StorageAPI.RenameData(ctx, srcVolume, srcPath, fi, dstVolume, dstPath, opts)
|
||||
}
|
||||
|
||||
func TestAccessMoveWriteFailureCanResume(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
const object = "write-failure"
|
||||
old := putAccessMoveVersion(t, z, bucket, object, 1, "old", false)
|
||||
latest := putAccessMoveVersion(t, z, bucket, object, 1, "new", false)
|
||||
existing := putAccessMoveVersion(t, z, bucket, object, 0, "unique-destination", true)
|
||||
set := z.serverPools[0].getHashedSet(object)
|
||||
getDisks := set.getDisks
|
||||
disks := getDisks()
|
||||
faulty := make([]StorageAPI, len(disks))
|
||||
for i, disk := range disks {
|
||||
faulty[i] = accessMoveFaultDisk{StorageAPI: disk, failVersion: latest.VersionID}
|
||||
}
|
||||
set.getDisks = func() []StorageAPI { return faulty }
|
||||
_, err := moveObjectPool(t.Context(), z, bucket, object, 1, 0, nil)
|
||||
set.getDisks = getDisks
|
||||
if err == nil {
|
||||
t.Fatal("injected destination write failure was ignored")
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 1, old, "old")
|
||||
assertAccessMoveVersion(t, z, bucket, object, 1, latest, "new")
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, existing, "unique-destination")
|
||||
if _, err := moveObjectPool(t.Context(), z, bucket, object, 1, 0, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, old, "old")
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, latest, "new")
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, existing, "unique-destination")
|
||||
}
|
||||
|
||||
func TestAccessMoveExcludesSourceWriter(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
const object = "concurrent"
|
||||
old := putAccessMoveVersion(t, z, bucket, object, 1, "old", false)
|
||||
_, err := moveObjectPool(t.Context(), z, bucket, object, 1, 0, func(ObjectInfo, uint64) error {
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 200*time.Millisecond)
|
||||
defer cancel()
|
||||
_, err := z.serverPools[1].PutObject(ctx, bucket, object,
|
||||
mustGetPutObjReader(t, bytes.NewBufferString("racing-write"), 12, "", ""), ObjectOptions{Versioned: true})
|
||||
if err == nil {
|
||||
t.Error("source writer committed while move was in progress")
|
||||
}
|
||||
if ctx.Err() == nil {
|
||||
return errors.New("writer did not wait for the move lock")
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, old, "old")
|
||||
}
|
||||
|
||||
func TestAccessMoveResumesPartialCopy(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
const object = "partial-copy"
|
||||
old := putAccessMoveVersion(t, z, bucket, object, 1, "old", false)
|
||||
latest := putAccessMoveVersion(t, z, bucket, object, 1, "new", false)
|
||||
gr, err := z.serverPools[1].GetObjectNInfo(t.Context(), bucket, object, nil, http.Header{}, ObjectOptions{VersionID: old.VersionID, NoDecryption: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Model a process stopping after the first copy commits, before cleanup.
|
||||
if err := moveAccessTierVersion(t.Context(), z, 1, 0, bucket, gr, time.Now().UnixNano()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := moveObjectPool(t.Context(), z, bucket, object, 1, 0, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, old, "old")
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, latest, "new")
|
||||
}
|
||||
|
||||
func TestAccessMoveRefusesConflictingVersion(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
const object = "conflicting-version"
|
||||
source := putAccessMoveVersion(t, z, bucket, object, 1, "source", false)
|
||||
destination, err := z.serverPools[0].PutObject(t.Context(), bucket, object,
|
||||
mustGetPutObjReader(t, bytes.NewBufferString("target"), 6, "", ""), ObjectOptions{
|
||||
VersionID: source.VersionID, MTime: source.ModTime,
|
||||
UserDefined: map[string]string{"test-value": "target"},
|
||||
WantChecksum: hash.NewChecksumFromData(hash.ChecksumCRC32C, []byte("target")),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := moveObjectPool(t.Context(), z, bucket, object, 1, 0, nil); !errors.Is(err, errAccessTierNotEligible) {
|
||||
t.Fatalf("conflicting version should be left intact: %v", err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 1, source, "source")
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, destination, "target")
|
||||
}
|
||||
|
||||
type accessMoveDeleteFaultDisk struct {
|
||||
StorageAPI
|
||||
bucket, object, version string
|
||||
}
|
||||
|
||||
func (d accessMoveDeleteFaultDisk) DeleteVersion(ctx context.Context, volume, path string, fi FileInfo, forceDelMarker bool, opts DeleteOptions) error {
|
||||
if volume == d.bucket && path == d.object && fi.VersionID == d.version {
|
||||
return errDiskFull
|
||||
}
|
||||
return d.StorageAPI.DeleteVersion(ctx, volume, path, fi, forceDelMarker, opts)
|
||||
}
|
||||
|
||||
func TestAccessMoveSourceDeleteFailureCanResume(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
const object = "source-delete-failure"
|
||||
old := putAccessMoveVersion(t, z, bucket, object, 1, "old", false)
|
||||
latest := putAccessMoveVersion(t, z, bucket, object, 1, "new", false)
|
||||
set := z.serverPools[1].getHashedSet(object)
|
||||
getDisks := set.getDisks
|
||||
faulty := append([]StorageAPI(nil), getDisks()...)
|
||||
// The old version is purged. One disk then completes deletion of the
|
||||
// latest version; the remaining disks reject it.
|
||||
for i := 1; i < len(faulty); i++ {
|
||||
faulty[i] = accessMoveDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object, version: latest.VersionID}
|
||||
}
|
||||
set.getDisks = func() []StorageAPI { return faulty }
|
||||
_, err := moveObjectPool(t.Context(), z, bucket, object, 1, 0, nil)
|
||||
set.getDisks = getDisks
|
||||
if err == nil {
|
||||
t.Fatal("injected partial source purge was ignored")
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, old, "old")
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, latest, "new")
|
||||
if _, err := moveObjectPool(t.Context(), z, bucket, object, 1, 0, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, old, "old")
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, latest, "new")
|
||||
}
|
||||
|
||||
// Exercise the public multipart API and compare whole-object and part reads
|
||||
// across both directions of a real pool move, including raw SSE-C ciphertext.
|
||||
func TestAccessMoveMultipartChecksums(t *testing.T) {
|
||||
z, _ := accessMovePools(t)
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
defer cancel()
|
||||
bucket, router, err := initAPIHandlerTest(ctx, z, nil, MakeBucketOptions{VersioningEnabled: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
previousTLS := globalIsTLS
|
||||
globalIsTLS = true
|
||||
defer func() { globalIsTLS = previousTLS }()
|
||||
partData, full := multipartChecksumTestData()
|
||||
for _, tc := range []struct {
|
||||
typ hash.ChecksumType
|
||||
ssec bool
|
||||
}{{hash.ChecksumCRC32, false}, {hash.ChecksumCRC32C, false}, {hash.ChecksumCRC64NVME, false}, {hash.ChecksumCRC32C, true}} {
|
||||
t.Run(fmt.Sprintf("%s/ssec=%t", tc.typ, tc.ssec), func(t *testing.T) {
|
||||
object := getRandomObjectName()
|
||||
headers := map[string]string{}
|
||||
if tc.ssec {
|
||||
key := bytes.Repeat([]byte{0x42}, 32)
|
||||
keyMD5 := md5.Sum(key)
|
||||
headers[xhttp.AmzServerSideEncryptionCustomerAlgorithm] = xhttp.AmzEncryptionAES
|
||||
headers[xhttp.AmzServerSideEncryptionCustomerKey] = base64.StdEncoding.EncodeToString(key)
|
||||
headers[xhttp.AmzServerSideEncryptionCustomerKeyMD5] = base64.StdEncoding.EncodeToString(keyMD5[:])
|
||||
}
|
||||
do := func(method, url string, body []byte, extra map[string]string) *httptest.ResponseRecorder {
|
||||
t.Helper()
|
||||
h := maps.Clone(headers)
|
||||
maps.Copy(h, extra)
|
||||
req, err := newTestSignedRequestV4(method, url, int64(len(body)), bytes.NewReader(body), globalActiveCred.AccessKey, globalActiveCred.SecretKey, h)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rec := httptest.NewRecorder()
|
||||
router.ServeHTTP(rec, req)
|
||||
if rec.Code != http.StatusOK && rec.Code != http.StatusPartialContent {
|
||||
t.Fatalf("%s %s: %d %s", method, url, rec.Code, rec.Body.String())
|
||||
}
|
||||
return rec
|
||||
}
|
||||
init := do(http.MethodPost, getNewMultipartURL("", bucket, object), nil, map[string]string{
|
||||
xhttp.AmzChecksumAlgo: tc.typ.String(), xhttp.AmzChecksumType: xhttp.AmzChecksumTypeFullObject,
|
||||
})
|
||||
var upload InitiateMultipartUploadResponse
|
||||
if err := xml.Unmarshal(init.Body.Bytes(), &upload); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
parts := make([]CompletePart, len(partData))
|
||||
for i, data := range partData {
|
||||
rec := do(http.MethodPut, getPutObjectPartURL("", bucket, object, upload.UploadID, fmt.Sprint(i+1)), data,
|
||||
map[string]string{tc.typ.Key(): mustChecksum(t, tc.typ, data)})
|
||||
parts[i] = CompletePart{PartNumber: i + 1, ETag: canonicalizeETag(rec.Header()[xhttp.ETag][0])}
|
||||
}
|
||||
body, err := xml.Marshal(CompleteMultipartUpload{Parts: parts})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
do(http.MethodPost, getCompleteMultipartUploadURL("", bucket, object, upload.UploadID), body,
|
||||
map[string]string{tc.typ.Key(): mustChecksum(t, tc.typ, full), xhttp.AmzChecksumType: xhttp.AmzChecksumTypeFullObject})
|
||||
source, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{})
|
||||
if err != nil || len(source.Checksum) == 0 {
|
||||
t.Fatalf("source checksum missing: %v", err)
|
||||
}
|
||||
src := 0
|
||||
if _, err := z.serverPools[src].GetObjectInfo(t.Context(), bucket, object, ObjectOptions{}); isErrObjectNotFound(err) {
|
||||
src = 1
|
||||
}
|
||||
for i := 0; i < 2; i++ {
|
||||
dst := 1 - src
|
||||
if _, err := moveObjectPool(t.Context(), z, bucket, object, src, dst, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got, err := z.serverPools[dst].GetObjectInfo(t.Context(), bucket, object, ObjectOptions{})
|
||||
if err != nil || !bytes.Equal(got.Checksum, source.Checksum) || got.ETag != source.ETag || got.VersionID != source.VersionID || !got.ModTime.Equal(source.ModTime) {
|
||||
t.Fatalf("move changed version metadata: %v checksumEqual=%t ETag=%q/%q version=%q/%q modTimeEqual=%t", err, bytes.Equal(got.Checksum, source.Checksum), got.ETag, source.ETag, got.VersionID, source.VersionID, got.ModTime.Equal(source.ModTime))
|
||||
}
|
||||
rec := do(http.MethodGet, getGetObjectURL("", bucket, object), nil, map[string]string{xhttp.AmzChecksumMode: "ENABLED"})
|
||||
if !bytes.Equal(rec.Body.Bytes(), full) || rec.Header().Get(tc.typ.Key()) != mustChecksum(t, tc.typ, full) {
|
||||
t.Fatal("GET payload or checksum changed after move")
|
||||
}
|
||||
for j, data := range partData {
|
||||
rec := do(http.MethodGet, getGetObjectURL("", bucket, object)+fmt.Sprintf("?partNumber=%d", j+1), nil, nil)
|
||||
if !bytes.Equal(rec.Body.Bytes(), data) {
|
||||
t.Fatalf("part %d changed after move", j+1)
|
||||
}
|
||||
}
|
||||
src = dst
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
type accessMoveNamedLocker struct {
|
||||
*localLocker
|
||||
address string
|
||||
}
|
||||
|
||||
func (l accessMoveNamedLocker) String() string { return l.address }
|
||||
|
||||
func TestAccessMoveSharedDistributedLockers(t *testing.T) {
|
||||
// Three overlapping sets of real dsync lock servers. Taking independent
|
||||
// quorum locks for the same object would contend with our own earlier lock.
|
||||
peers := make([]dsync.NetLocker, 5)
|
||||
for i := range peers {
|
||||
peers[i] = accessMoveNamedLocker{localLocker: newLocker(), address: fmt.Sprint(i)}
|
||||
}
|
||||
z := &erasureServerPools{}
|
||||
for i := 0; i < 3; i++ {
|
||||
setPeers := peers[i : i+3]
|
||||
set := &erasureObjects{nsMutex: &nsLockMap{isDistErasure: true}, getLockers: func() ([]dsync.NetLocker, string) {
|
||||
return setPeers, "access-move-fixture"
|
||||
}}
|
||||
z.serverPools = append(z.serverPools, &erasureSets{sets: []*erasureObjects{set}, distributionAlgo: formatErasureVersionV3DistributionAlgoV3})
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second)
|
||||
defer cancel()
|
||||
_, unlock, err := lockAccessTierObject(ctx, z, "bucket", "object", 1, 2)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
release := sync.OnceFunc(unlock)
|
||||
defer release()
|
||||
for i, pool := range z.serverPools {
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond)
|
||||
lock := pool.NewNSLock("bucket", "object")
|
||||
lc, err := lock.GetLock(ctx, globalOperationTimeout)
|
||||
cancel()
|
||||
if err == nil {
|
||||
lock.Unlock(lc)
|
||||
t.Fatalf("pool %d writer acquired its quorum during the move", i)
|
||||
}
|
||||
}
|
||||
release()
|
||||
// Every acquired peer lock is released, so ordinary writes can resume.
|
||||
for _, peer := range peers {
|
||||
// Distributed Unlock sends releases asynchronously.
|
||||
lock := (&nsLockMap{isDistErasure: true}).NewNSLock(func() ([]dsync.NetLocker, string) {
|
||||
return []dsync.NetLocker{peer}, "access-move-fixture"
|
||||
}, "bucket", "object")
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 2*time.Second)
|
||||
lc, err := lock.GetLock(ctx, globalOperationTimeout)
|
||||
cancel()
|
||||
if err != nil {
|
||||
t.Fatalf("peer %s leaked a lock: %v", peer.String(), err)
|
||||
}
|
||||
lock.Unlock(lc)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAccessMoveConfiguredPromotionAndDemotion(t *testing.T) {
|
||||
z, bucket := accessMovePools(t)
|
||||
const object = "configured-move"
|
||||
version := putAccessMoveVersion(t, z, bucket, object, 1, "configured", false)
|
||||
oldCfg, oldTracker := globalILMConfig.accessCfg(), globalAccessTracker
|
||||
defer func() { globalILMConfig.update(oldCfg); globalAccessTracker = oldTracker }()
|
||||
cfg := oldCfg
|
||||
cfg.AccessTiering, cfg.AccessPools = true, []int{0, 1}
|
||||
cfg.AccessMinResidency, cfg.AccessPromoteWatermark = 0, 99
|
||||
globalILMConfig.update(cfg)
|
||||
globalAccessTracker = newAccessTracker()
|
||||
lc := accessLifecycleForTest(t)
|
||||
data, err := xml.Marshal(lc)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := globalBucketMetadataSys.Update(t.Context(), bucket, bucketLifecycleConfig, data); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
state := newAccessTierState(t.Context())
|
||||
state.z, state.objAPI = z, z
|
||||
task := accessTierTask{ctx: t.Context(), bucket: bucket, object: object, src: 1, dst: 0, direction: accessTierPromote, bytes: uint64(version.Size)}
|
||||
// A queued move is rechecked when configuration is disabled dynamically.
|
||||
cfg.AccessTiering = false
|
||||
globalILMConfig.update(cfg)
|
||||
if err := state.processTask(task); !errors.Is(err, errAccessTierNotEligible) {
|
||||
t.Fatalf("disabled mover = %v", err)
|
||||
}
|
||||
cfg.AccessTiering = true
|
||||
globalILMConfig.update(cfg)
|
||||
globalAccessTracker.merged.Store(&mergedAccess{binWidth: 60, entries: map[string]accessEntry{
|
||||
accessKey(bucket, object): {Bins: []uint32{100}, LastAt: time.Now().Unix()},
|
||||
}})
|
||||
if reason, ok := state.reservePromotion(task, cfg, 0); !ok {
|
||||
t.Fatalf("promotion reservation failed: %s", reason)
|
||||
}
|
||||
if err := state.processTask(task); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 0, version, "configured")
|
||||
// Advance the counter view beyond the rule's idle period.
|
||||
globalAccessTracker.merged.Store(&mergedAccess{binWidth: 60, entries: map[string]accessEntry{
|
||||
accessKey(bucket, object): {Bins: []uint32{0}, LastAt: time.Now().Add(-2 * time.Hour).Unix()},
|
||||
}})
|
||||
task.src, task.dst, task.direction, task.bytes = 0, 1, accessTierDemote, 0
|
||||
if err := state.processTask(task); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assertAccessMoveVersion(t, z, bucket, object, 1, version, "configured")
|
||||
if state.promotions.Load() != 1 || state.demotions.Load() != 1 || state.bytesMoved.Load() != 2*uint64(version.Size) || len(state.pending) != 0 || len(state.reserved) != 0 {
|
||||
t.Fatal("promotion/demotion metrics or reservation cleanup are incorrect")
|
||||
}
|
||||
}
|
||||
+128
-50
@@ -18,12 +18,14 @@
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"maps"
|
||||
"net/http"
|
||||
"slices"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
@@ -32,6 +34,7 @@ import (
|
||||
|
||||
"github.com/minio/minio/internal/bucket/lifecycle"
|
||||
"github.com/minio/minio/internal/config/ilm"
|
||||
"github.com/minio/minio/internal/dsync"
|
||||
"github.com/minio/minio/internal/hash"
|
||||
)
|
||||
|
||||
@@ -851,9 +854,6 @@ func accessObjectVersions(ctx context.Context, z *erasureServerPools, src int, b
|
||||
return nil, errAccessTierRemoteVersion
|
||||
}
|
||||
}
|
||||
if len(versions) == 1 && versions[0].Deleted {
|
||||
return nil, errAccessTierNotEligible
|
||||
}
|
||||
versionsSorter(versions).reverse()
|
||||
return versions, nil
|
||||
}
|
||||
@@ -872,24 +872,24 @@ func accessObjectBytes(ctx context.Context, z *erasureServerPools, src int, buck
|
||||
return total, nil
|
||||
}
|
||||
|
||||
func deleteAccessTierPoolObject(ctx context.Context, z *erasureServerPools, pool int, bucket, object string) error {
|
||||
func deleteAccessTierPoolVersion(ctx context.Context, z *erasureServerPools, pool int, bucket, object, versionID string) error {
|
||||
if versionID == "" {
|
||||
versionID = nullVersionID
|
||||
}
|
||||
_, err := z.serverPools[pool].DeleteObject(ctx, bucket, encodeDirObject(object), ObjectOptions{
|
||||
DeletePrefix: true, DeletePrefixObject: true, NoLock: true, NoAuditLog: true,
|
||||
Versioned: true, VersionID: versionID, NoLock: true, NoAuditLog: true,
|
||||
})
|
||||
if isErrObjectNotFound(err) || isErrVersionNotFound(err) {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
// rollbackAccessTierDestination removes only what this move wrote.
|
||||
//
|
||||
// A blanket prefix delete would be wrong: the top-level namespace lock is
|
||||
// taken on pool 0 / set 0 (see erasureServerPools.NewNSLock), while a client
|
||||
// PutObject locks the hashed set of whichever pool it lands in, and lockers
|
||||
// are per set. A concurrent client write can therefore land on the
|
||||
// destination while a move is in flight, and it must survive our rollback.
|
||||
//
|
||||
// Versioned writes are removed by exact version ID, which can never touch a
|
||||
// client's version. An unversioned write has no ID to target, so it is only
|
||||
// removed while it still carries this move's own marker stamp.
|
||||
// Existing destination versions are never included in written. Both commit
|
||||
// locks remain held during rollback; the marker check also refuses cleanup of
|
||||
// an unversioned copy that no longer belongs to this attempt.
|
||||
func rollbackAccessTierDestination(ctx context.Context, z *erasureServerPools, dst int, bucket, object string, written []string, movedAt int64) {
|
||||
stamp := accessTierStamp(dst, movedAt)
|
||||
for _, versionID := range written {
|
||||
@@ -907,27 +907,94 @@ func rollbackAccessTierDestination(ctx context.Context, z *erasureServerPools, d
|
||||
ilmLogIf(ctx, fmt.Errorf("access tier rollback skipped for %s/%s in pool %d: destination was overwritten concurrently", bucket, object, dst))
|
||||
continue
|
||||
}
|
||||
if err := deleteAccessTierPoolObject(ctx, z, dst, bucket, object); err != nil {
|
||||
ilmLogIf(ctx, fmt.Errorf("access tier rollback failed for %s/%s in pool %d: %w", bucket, object, dst, err))
|
||||
}
|
||||
continue
|
||||
}
|
||||
_, err := z.serverPools[dst].DeleteObject(ctx, bucket, encodeDirObject(object), ObjectOptions{
|
||||
Versioned: true, VersionID: versionID, NoLock: true, NoAuditLog: true,
|
||||
})
|
||||
if err != nil && !isErrObjectNotFound(err) && !isErrVersionNotFound(err) {
|
||||
if err := deleteAccessTierPoolVersion(ctx, z, dst, bucket, object, versionID); err != nil {
|
||||
ilmLogIf(ctx, fmt.Errorf("access tier rollback failed for %s/%s version %s in pool %d: %w", bucket, object, versionID, dst, err))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// moveObjectPool copies the complete version stack oldest-first and removes
|
||||
// the source only after every destination write succeeded.
|
||||
//
|
||||
// The namespace lock excludes concurrent deletes, which take the same
|
||||
// top-level lock. It does not exclude a concurrent PutObject, which locks the
|
||||
// hashed set of its target pool; see rollbackAccessTierDestination for how
|
||||
// that is contained.
|
||||
// lockAccessTierObject excludes reads/deletes using the top-level namespace
|
||||
// and writes committing in either hashed set. Distributed sets can share lock
|
||||
// servers, so lock each distinct server once, in a stable order. Background
|
||||
// moves require all involved lock servers; ordinary S3 quorum rules stay intact.
|
||||
func lockAccessTierObject(ctx context.Context, z *erasureServerPools, bucket, object string, src, dst int) (context.Context, func(), error) {
|
||||
sets := []*erasureObjects{z.serverPools[0].sets[0], z.serverPools[src].getHashedSet(object), z.serverPools[dst].getHashedSet(object)}
|
||||
var locks []RWLocker
|
||||
local := make(map[*nsLockMap]bool)
|
||||
remote := make(map[string]dsync.NetLocker)
|
||||
var owner string
|
||||
for _, set := range sets {
|
||||
if !set.nsMutex.isDistErasure {
|
||||
if !local[set.nsMutex] {
|
||||
locks = append(locks, set.NewNSLock(bucket, object))
|
||||
local[set.nsMutex] = true
|
||||
}
|
||||
continue
|
||||
}
|
||||
peers, lockOwner := set.getLockers()
|
||||
if len(peers) == 0 {
|
||||
return nil, nil, errAccessTierNotEligible
|
||||
}
|
||||
owner = lockOwner
|
||||
for _, peer := range peers {
|
||||
if peer == nil {
|
||||
return nil, nil, errAccessTierNotEligible
|
||||
}
|
||||
remote[peer.String()] = peer
|
||||
}
|
||||
}
|
||||
for _, address := range slices.Sorted(maps.Keys(remote)) {
|
||||
peer := remote[address]
|
||||
locks = append(locks, (&nsLockMap{isDistErasure: true}).NewNSLock(func() ([]dsync.NetLocker, string) {
|
||||
return []dsync.NetLocker{peer}, owner
|
||||
}, bucket, object))
|
||||
}
|
||||
var unlockers []func()
|
||||
unlock := func() {
|
||||
for i := len(unlockers) - 1; i >= 0; i-- {
|
||||
unlockers[i]()
|
||||
}
|
||||
}
|
||||
for _, lk := range locks {
|
||||
lkctx, err := lk.GetLock(ctx, globalDeleteOperationTimeout)
|
||||
if err != nil {
|
||||
unlock()
|
||||
return nil, nil, err
|
||||
}
|
||||
ctx = lkctx.Context()
|
||||
unlockers = append(unlockers, func() { lk.Unlock(lkctx) })
|
||||
}
|
||||
return ctx, unlock, nil
|
||||
}
|
||||
|
||||
// A retry may reuse an identical destination version. Conflicting versions
|
||||
// are left intact, including null versions and independently changed tags or
|
||||
// retention settings. Physical layout and move bookkeeping may differ by pool.
|
||||
func sameAccessTierVersion(a, b FileInfo) bool {
|
||||
if a.VersionID != b.VersionID || a.Deleted != b.Deleted || a.Size != b.Size || !a.ModTime.Equal(b.ModTime) || !bytes.Equal(a.Checksum, b.Checksum) || len(a.Parts) != len(b.Parts) {
|
||||
return false
|
||||
}
|
||||
for i, part := range a.Parts {
|
||||
other := b.Parts[i]
|
||||
if part.Number != other.Number || part.Size != other.Size || part.ActualSize != other.ActualSize || part.ETag != other.ETag || !bytes.Equal(part.Index, other.Index) || !maps.Equal(part.Checksums, other.Checksums) {
|
||||
return false
|
||||
}
|
||||
}
|
||||
metadata := func(fi FileInfo) map[string]string {
|
||||
m := maps.Clone(fi.Metadata)
|
||||
delete(m, accessTierMetadataKey)
|
||||
delete(m, xMinIODataMov)
|
||||
delete(m, ReservedMetadataPrefixLower+"inline-data")
|
||||
delete(m, minIOErasureUpgraded)
|
||||
return m
|
||||
}
|
||||
return maps.Equal(metadata(a), metadata(b))
|
||||
}
|
||||
|
||||
// moveObjectPool copies missing versions oldest-first, preserving versions
|
||||
// already at the destination. The source is removed only after every source
|
||||
// version is present there. This also resumes after an uncertain source delete.
|
||||
func moveObjectPool(ctx context.Context, z *erasureServerPools, bucket, object string, src, dst int, beforeCopy func(ObjectInfo, uint64) error) (uint64, error) {
|
||||
if bucket == minioMetaBucket || src == dst || src < 0 || dst < 0 || src >= len(z.serverPools) || dst >= len(z.serverPools) {
|
||||
return 0, errAccessTierNotEligible
|
||||
@@ -936,18 +1003,19 @@ func moveObjectPool(ctx context.Context, z *erasureServerPools, bucket, object s
|
||||
return 0, errAccessTierNotEligible
|
||||
}
|
||||
|
||||
lk := z.NewNSLock(bucket, object)
|
||||
lkctx, err := lk.GetLock(ctx, globalDeleteOperationTimeout)
|
||||
ctx, unlock, err := lockAccessTierObject(ctx, z, bucket, encodeDirObject(object), src, dst)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
ctx = lkctx.Context()
|
||||
defer lk.Unlock(lkctx)
|
||||
defer unlock()
|
||||
|
||||
versions, err := accessObjectVersions(ctx, z, src, bucket, object)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
if len(versions) == 1 && versions[0].Deleted {
|
||||
return 0, errAccessTierNotEligible
|
||||
}
|
||||
versioned := globalBucketVersioningSys.PrefixEnabled(bucket, object)
|
||||
latest := versions[len(versions)-1].ToObjectInfo(bucket, object, versioned)
|
||||
var total uint64
|
||||
@@ -962,20 +1030,20 @@ func moveObjectPool(ctx context.Context, z *erasureServerPools, bucket, object s
|
||||
}
|
||||
}
|
||||
|
||||
// A failed earlier attempt may have copied the full stack but failed to
|
||||
// remove the source. It is safe to discard only a destination carrying
|
||||
// our internal marker; an unrelated split-brain copy is left untouched.
|
||||
dstOI, dstErr := z.serverPools[dst].GetObjectInfo(ctx, bucket, encodeDirObject(object), ObjectOptions{NoLock: true})
|
||||
if dstErr == nil {
|
||||
markedPool, _, marked := parseAccessTierStamp(dstOI.UserDefined)
|
||||
if !marked || markedPool != dst {
|
||||
// Never clear a pre-existing destination: it may hold the only surviving
|
||||
// copy of a version after a partially committed source deletion.
|
||||
dstVersions, dstErr := accessObjectVersions(ctx, z, dst, bucket, object)
|
||||
if dstErr != nil && !isErrObjectNotFound(dstErr) && !isErrVersionNotFound(dstErr) {
|
||||
return 0, dstErr
|
||||
}
|
||||
existing := make(map[string]FileInfo, len(dstVersions))
|
||||
for _, version := range dstVersions {
|
||||
existing[version.VersionID] = version
|
||||
}
|
||||
for _, version := range versions {
|
||||
if prior, ok := existing[version.VersionID]; ok && !sameAccessTierVersion(version, prior) {
|
||||
return 0, errAccessTierNotEligible
|
||||
}
|
||||
if err := deleteAccessTierPoolObject(ctx, z, dst, bucket, object); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
} else if !isErrObjectNotFound(dstErr) && !isErrVersionNotFound(dstErr) {
|
||||
return 0, dstErr
|
||||
}
|
||||
|
||||
movedAt := time.Now().UnixNano()
|
||||
@@ -994,6 +1062,9 @@ func moveObjectPool(ctx context.Context, z *erasureServerPools, bucket, object s
|
||||
|
||||
set := z.serverPools[src].getHashedSet(encodeDirObject(object))
|
||||
for _, version := range versions {
|
||||
if _, ok := existing[version.VersionID]; ok {
|
||||
continue
|
||||
}
|
||||
if err := ctx.Err(); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
@@ -1029,7 +1100,7 @@ func moveObjectPool(ctx context.Context, z *erasureServerPools, bucket, object s
|
||||
|
||||
// Re-read while still holding the namespace lock and refuse the source
|
||||
// delete if the latest version changed despite the lock contract.
|
||||
current, err := z.serverPools[src].GetObjectInfo(ctx, bucket, encodeDirObject(object), ObjectOptions{NoLock: true})
|
||||
current, err := z.serverPools[src].GetObjectInfo(ctx, bucket, encodeDirObject(object), ObjectOptions{NoLock: true, Versioned: versioned})
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
@@ -1040,9 +1111,12 @@ func moveObjectPool(ctx context.Context, z *erasureServerPools, bucket, object s
|
||||
// the complete destination in that case: preserving two copies is safer
|
||||
// than risking zero.
|
||||
rollbackDestination = false
|
||||
err = deleteAccessTierPoolObject(ctx, z, src, bucket, object)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
// A recursive prefix delete could also remove another key such as
|
||||
// "object/child" on this set. Delete only the versions we actually copied.
|
||||
for _, version := range versions {
|
||||
if err := deleteAccessTierPoolVersion(ctx, z, src, bucket, object, version.VersionID); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
}
|
||||
return total, nil
|
||||
}
|
||||
@@ -1060,6 +1134,9 @@ func moveAccessTierVersion(ctx context.Context, z *erasureServerPools, src, dst
|
||||
metadata = make(map[string]string, 1)
|
||||
}
|
||||
metadata[accessTierMetadataKey] = accessTierStamp(dst, movedAt)
|
||||
// Multipart ETags are derived from plaintext part hashes at upload time.
|
||||
// A raw encrypted copy must retain that ETag instead of hashing ciphertext.
|
||||
metadata["etag"] = oi.ETag
|
||||
|
||||
if oi.isMultipart() {
|
||||
res, err := z.NewMultipartUpload(ctx, bucket, oi.Name, ObjectOptions{
|
||||
@@ -1092,7 +1169,8 @@ func moveAccessTierVersion(ctx context.Context, z *erasureServerPools, src, dst
|
||||
parts[i] = CompletePart{
|
||||
ETag: pi.ETag, PartNumber: pi.PartNumber,
|
||||
ChecksumCRC32: pi.ChecksumCRC32, ChecksumCRC32C: pi.ChecksumCRC32C,
|
||||
ChecksumSHA1: pi.ChecksumSHA1, ChecksumSHA256: pi.ChecksumSHA256,
|
||||
ChecksumCRC64NVME: pi.ChecksumCRC64NVME,
|
||||
ChecksumSHA1: pi.ChecksumSHA1, ChecksumSHA256: pi.ChecksumSHA256,
|
||||
}
|
||||
}
|
||||
_, err = z.CompleteMultipartUpload(ctx, bucket, oi.Name, res.UploadID, parts, ObjectOptions{
|
||||
|
||||
@@ -327,6 +327,15 @@ Demotion is not blocked by these caps and is processed before promotion.
|
||||
Access moves pause during rebalance or decommission, never target a suspended
|
||||
pool, skip remotely transitioned objects and objects with excessive version
|
||||
counts, and recheck eligibility while holding the object namespace lock.
|
||||
Moves also lock the source and destination write locations. All lock servers
|
||||
for those locations must be reachable; otherwise the background move is
|
||||
deferred. Ordinary S3 requests keep their existing quorum requirements.
|
||||
|
||||
Each move preserves version IDs, delete markers, ETags, checksums, encryption
|
||||
and user metadata. A retry copies only missing versions and retains versions
|
||||
already at the destination. If the same version ID has conflicting metadata
|
||||
in the two pools, the move is skipped and both copies are left intact. Source
|
||||
data is removed only after the complete version stack exists at the destination.
|
||||
|
||||
The hit counter is intentionally best effort. Only successfully served GET
|
||||
requests count; HEAD requests do not. Counters are merged across nodes and
|
||||
@@ -348,6 +357,17 @@ Access-tier activity is exposed under /minio/metrics/v3/ilm, including move
|
||||
counts, moved bytes, queue depth, hot bytes per bucket, failed moves, dropped
|
||||
GET samples, and separate skip counters for each capacity limit.
|
||||
|
||||
To stop scheduling moves, set `access_tiering=off` (and remove any environment
|
||||
override). Objects already moved remain in their current pools and stay
|
||||
accessible; disabling the feature does not move them back. Before downgrading
|
||||
to a release without access tiering, disable it and let active moves finish.
|
||||
The object storage format is unchanged. The data-usage cache advances from
|
||||
v8 to v9: this release reads both, but an older binary discards v9 caches and
|
||||
rebuilds usage statistics through the scanner. Usage and quota statistics can
|
||||
therefore take time to repopulate after a downgrade. Keep a copy of lifecycle
|
||||
XML containing Silo extensions, since an older binary may omit those fields
|
||||
when rewriting a lifecycle rule.
|
||||
|
||||
## Explore Further
|
||||
|
||||
- [MinIO Go client API reference (S3-compatible SDK)](https://pkg.go.dev/github.com/minio/minio-go/v7)
|
||||
|
||||
Reference in New Issue
Block a user