diff --git a/cmd/ilm-access-move_test.go b/cmd/ilm-access-move_test.go new file mode 100644 index 000000000..9d150edef --- /dev/null +++ b/cmd/ilm-access-move_test.go @@ -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") + } +} diff --git a/cmd/ilm-access-tier.go b/cmd/ilm-access-tier.go index 28007b141..cc69974ce 100644 --- a/cmd/ilm-access-tier.go +++ b/cmd/ilm-access-tier.go @@ -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{ diff --git a/docs/bucket/lifecycle/README.md b/docs/bucket/lifecycle/README.md index 3e5f5d5dd..362203b3f 100644 --- a/docs/bucket/lifecycle/README.md +++ b/docs/bucket/lifecycle/README.md @@ -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)