revert: remove access-frequency ILM tiering (#60)

Reverse the first-parent diff of a3df317ae0,
including the feature branch compatibility and mover follow-up fixes.
Retain the independent multi-pool correctness fixes from #178 and migrate
their shared test fixture away from access-tier code.

Tolerate retired ILM keys and XML, read old v9 statistics while writing v8,
and document migration without moving objects or rewriting their metadata.
Include regression coverage using a historical scanner/writer v9 fixture.

Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
Feng Ruohang
2026-09-15 06:51:14 +08:00
parent 89637554d6
commit 9b76a21675
44 changed files with 581 additions and 5670 deletions
-1
View File
@@ -39,7 +39,6 @@ const (
lcEventSrc_s3PutObject
lcEventSrc_s3CopyObject
lcEventSrc_s3CompleteMultipartUpload
lcEventSrc_AccessTier
)
//revive:enable:var-naming
+25
View File
@@ -39,3 +39,28 @@ func TestBucketMetadataCorsRoundTrip(t *testing.T) {
t.Fatalf("CorsConfigUpdatedAt not preserved")
}
}
// A persisted retired extension must not prevent the whole bucket's metadata
// from loading, including unrelated versioning and ordinary lifecycle rules.
func TestBucketMetadataRetiredAccessTiering(t *testing.T) {
meta := newBucketMetadata("retired-access")
meta.LifecycleConfigXML = []byte(`<LifecycleConfiguration><AccessTierQuota>500GiB</AccessTierQuota><Rule><ID>access</ID><Status>Enabled</Status><Filter><Prefix>logs/</Prefix></Filter><AccessTransition><Window>10m</Window><PromoteAfterAccesses>10</PromoteAfterAccesses></AccessTransition></Rule><Rule><ID>ordinary</ID><Status>Enabled</Status><Filter><Prefix>expired/</Prefix></Filter><Expiration><Days>30</Days></Expiration></Rule></LifecycleConfiguration>`)
meta.VersioningConfigXML = []byte(`<VersioningConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Status>Enabled</Status></VersioningConfiguration>`)
data, err := meta.MarshalMsg(nil)
if err != nil {
t.Fatal(err)
}
got := newBucketMetadata(meta.Name)
if _, err := got.UnmarshalMsg(data); err != nil {
t.Fatal(err)
}
if err := got.parseAllConfigs(t.Context(), nil); err != nil {
t.Fatalf("bucket metadata failed to load: %v", err)
}
if got.lifecycleConfig == nil || got.lifecycleConfig.HasActiveRules("logs/") || !got.lifecycleConfig.HasActiveRules("expired/") {
t.Fatal("unexpected lifecycle behavior")
}
if got.versioningConfig == nil || !got.versioningConfig.Enabled() {
t.Fatal("unrelated versioning lost")
}
}
+1 -4
View File
@@ -227,7 +227,7 @@ func initHelp() {
},
config.HelpKV{
Key: config.ILMSubSys,
Description: "manage ILM settings for expiration, transition, and access-tier workers",
Description: "manage ILM settings for expiration and transition workers",
Optional: true,
},
}
@@ -704,9 +704,6 @@ func applyDynamicConfigForSubSys(ctx context.Context, objAPI ObjectLayer, s conf
if globalExpiryState != nil {
globalExpiryState.ResizeWorkers(ilmCfg.ExpirationWorkers)
}
if globalAccessTierState != nil {
globalAccessTierState.UpdateWorkers(ilmCfg.AccessWorkers)
}
globalILMConfig.update(ilmCfg)
}
}
-14
View File
@@ -899,7 +899,6 @@ type scannerItem struct {
objectName string // Only the object name without prefixes.
replication replicationConfig
lifeCycle *lifecycle.Lifecycle
poolIdx int
Typ fs.FileMode
heal struct {
enabled bool
@@ -910,7 +909,6 @@ type scannerItem struct {
type sizeSummary struct {
totalSize int64
hotTierSize int64
versions uint64
deleteMarkers uint64
replicatedSize int64
@@ -1161,14 +1159,6 @@ eventLoop:
globalExpiryState.enqueueNoncurrentVersions(i.bucket, toDel, noncurrentEvents)
}
i.alertExcessiveVersions(remainingVersions, cumulativeSize)
if globalILMConfig.accessTieringEnabled() {
for idx, oi := range objInfos {
if oi.IsLatest && events[idx].Action == lifecycle.NoneAction {
applyAccessTransition(ctx, i, oi)
break
}
}
}
}
func evalActionFromLifecycle(ctx context.Context, lc lifecycle.Lifecycle, lr lock.Retention, rcfg *replication.Config, obj ObjectInfo) lifecycle.Event {
@@ -1484,8 +1474,6 @@ const (
ILMFreeVersionDelete = "ilm:free-version-delete"
// ILMTransition - audit trail for ILM transitioning.
ILMTransition = " ilm:transition"
// ILMAccessTier - audit trail for moving objects between server pools.
ILMAccessTier = "ilm:access-tier"
)
func auditLogLifecycle(ctx context.Context, oi ObjectInfo, event string, tags map[string]string, traceFn func(event string, metadata map[string]string, err error)) {
@@ -1497,8 +1485,6 @@ func auditLogLifecycle(ctx context.Context, oi ObjectInfo, event string, tags ma
apiName = "ILMFreeVersionDelete"
case ILMTransition:
apiName = "ILMTransition"
case ILMAccessTier:
apiName = "ILMAccessTier"
}
auditLogInternal(ctx, AuditLogOptions{
Event: event,
+89
View File
@@ -0,0 +1,89 @@
// 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"
"encoding/json"
"os"
"reflect"
"testing"
)
// Read a real pre-removal scanner cache with nonzero hot-tier bytes. A v8
// payload with a relabeled header would not exercise the ignored hts field.
func TestDataUsageCacheReadV9(t *testing.T) {
raw, err := os.ReadFile("testdata/data-usage-v9/data-usage-v9.bin")
if err != nil {
t.Fatal(err)
}
if len(raw) == 0 || raw[0] != 9 {
t.Fatal("fixture is not v9")
}
expected, err := os.ReadFile("testdata/data-usage-v9/data-usage-v9.json")
if err != nil {
t.Fatal(err)
}
var want, got dataUsageCache
if err := json.Unmarshal(expected, &want); err != nil {
t.Fatal(err)
}
if err := got.deserialize(bytes.NewReader(raw)); err != nil {
t.Fatal(err)
}
if !got.Info.LastUpdate.Equal(want.Info.LastUpdate) {
t.Fatal("cache timestamp changed")
}
// msgp restores local time while JSON preserves the UTC representation.
want.Info.LastUpdate = got.Info.LastUpdate
if !reflect.DeepEqual(got, want) {
t.Fatalf("v9 ordinary cache fields changed:\ngot: %#v\nwant: %#v", got, want)
}
flat := got.flatten(*got.root())
if flat.Size != 74962 || flat.Objects != 3 || flat.Versions != 5 || flat.DeleteMarkers != 2 {
t.Fatalf("statistics changed: %+v", flat)
}
if flat.AllTierStats == nil || flat.AllTierStats.Tiers["COLD"] != (tierStats{TotalSize: 65536, NumVersions: 1, NumObjects: 1}) {
t.Fatalf("remote tier statistics changed: %+v", flat.AllTierStats)
}
buckets := []BucketInfo{{Name: "v9-bucket"}}
if gotInfo, wantInfo := got.dui(dataUsageRoot, buckets), want.dui(dataUsageRoot, buckets); !reflect.DeepEqual(gotInfo, wantInfo) {
t.Fatalf("bucket usage aggregation changed: %+v != %+v", gotInfo, wantInfo)
}
var buf bytes.Buffer
if err := got.serializeTo(&buf); err != nil {
t.Fatal(err)
}
if buf.Bytes()[0] != 8 {
t.Fatalf("wrote cache version %d, want 8", buf.Bytes()[0])
}
var roundtrip dataUsageCache
if err := roundtrip.deserialize(&buf); err != nil {
t.Fatal(err)
}
if !reflect.DeepEqual(got, roundtrip) {
t.Fatal("v8 round trip lost ordinary fields")
}
if err := roundtrip.deserialize(bytes.NewReader(raw[:len(raw)/2])); err == nil {
t.Fatal("truncated v9 cache accepted")
}
raw[0] = 10
if err := roundtrip.deserialize(bytes.NewReader(raw)); err == nil {
t.Fatal("unknown cache version accepted")
}
}
+7 -61
View File
@@ -61,7 +61,6 @@ type dataUsageEntry struct {
Children dataUsageHashMap `msg:"ch"`
// These fields do no include any children.
Size int64 `msg:"sz"`
HotTierSize int64 `msg:"hts"`
Objects uint64 `msg:"os"`
Versions uint64 `msg:"vs"` // Versions that are not delete markers.
DeleteMarkers uint64 `msg:"dms"`
@@ -135,8 +134,8 @@ func (ts tierStats) add(u tierStats) tierStats {
}
}
//msgp:encode ignore dataUsageEntryV2 dataUsageEntryV3 dataUsageEntryV4 dataUsageEntryV5 dataUsageEntryV6 dataUsageEntryV7 dataUsageEntryV8
//msgp:marshal ignore dataUsageEntryV2 dataUsageEntryV3 dataUsageEntryV4 dataUsageEntryV5 dataUsageEntryV6 dataUsageEntryV7 dataUsageEntryV8
//msgp:encode ignore dataUsageEntryV2 dataUsageEntryV3 dataUsageEntryV4 dataUsageEntryV5 dataUsageEntryV6 dataUsageEntryV7
//msgp:marshal ignore dataUsageEntryV2 dataUsageEntryV3 dataUsageEntryV4 dataUsageEntryV5 dataUsageEntryV6 dataUsageEntryV7
//msgp:tuple dataUsageEntryV2
type dataUsageEntryV2 struct {
@@ -200,30 +199,14 @@ type dataUsageEntryV7 struct {
Compacted bool `msg:"c"`
}
// dataUsageEntryV8 is the on-disk shape before access-tier accounting was
// introduced. Keep it so caches written by the previous release decode
// without being discarded.
type dataUsageEntryV8 struct {
Children dataUsageHashMap `msg:"ch"`
// These fields do no include any children.
Size int64 `msg:"sz"`
Objects uint64 `msg:"os"`
Versions uint64 `msg:"vs"`
DeleteMarkers uint64 `msg:"dms"`
ObjSizes sizeHistogram `msg:"szs"`
ObjVersions versionsHistogram `msg:"vh"`
AllTierStats *allTierStats `msg:"ats,omitempty"`
Compacted bool `msg:"c"`
}
// dataUsageCache contains a cache of data usage entries latest version.
type dataUsageCache struct {
Info dataUsageCacheInfo
Cache map[string]dataUsageEntry
}
//msgp:encode ignore dataUsageCacheV2 dataUsageCacheV3 dataUsageCacheV4 dataUsageCacheV5 dataUsageCacheV6 dataUsageCacheV7 dataUsageCacheV8
//msgp:marshal ignore dataUsageCacheV2 dataUsageCacheV3 dataUsageCacheV4 dataUsageCacheV5 dataUsageCacheV6 dataUsageCacheV7 dataUsageCacheV8
//msgp:encode ignore dataUsageCacheV2 dataUsageCacheV3 dataUsageCacheV4 dataUsageCacheV5 dataUsageCacheV6 dataUsageCacheV7
//msgp:marshal ignore dataUsageCacheV2 dataUsageCacheV3 dataUsageCacheV4 dataUsageCacheV5 dataUsageCacheV6 dataUsageCacheV7
// dataUsageCacheV2 contains a cache of data usage entries version 2.
type dataUsageCacheV2 struct {
@@ -261,12 +244,6 @@ type dataUsageCacheV7 struct {
Cache map[string]dataUsageEntryV7
}
// dataUsageCacheV8 contains a cache of data usage entries version 8.
type dataUsageCacheV8 struct {
Info dataUsageCacheInfo
Cache map[string]dataUsageEntryV8
}
//msgp:ignore dataUsageEntryInfo
type dataUsageEntryInfo struct {
Name string
@@ -295,7 +272,6 @@ type dataUsageCacheInfo struct {
func (e *dataUsageEntry) addSizes(summary sizeSummary) {
e.Size += summary.totalSize
e.HotTierSize += summary.hotTierSize
e.Versions += summary.versions
e.DeleteMarkers += summary.deleteMarkers
e.ObjSizes.add(summary.totalSize)
@@ -315,7 +291,6 @@ func (e *dataUsageEntry) merge(other dataUsageEntry) {
e.Versions += other.Versions
e.DeleteMarkers += other.DeleteMarkers
e.Size += other.Size
e.HotTierSize += other.HotTierSize
for i, v := range other.ObjSizes[:] {
e.ObjSizes[i] += v
@@ -456,7 +431,6 @@ func (d *dataUsageCache) dui(path string, buckets []BucketInfo) DataUsageInfo {
flat := d.flatten(*e)
dui := DataUsageInfo{
LastUpdate: d.Info.LastUpdate,
ScannerCycle: d.Info.NextCycle,
ObjectsTotalCount: flat.Objects,
VersionsTotalCount: flat.Versions,
DeleteMarkersTotalCount: flat.DeleteMarkers,
@@ -807,7 +781,6 @@ func (d *dataUsageCache) bucketsUsageInfo(buckets []BucketInfo) map[string]Bucke
flat := d.flatten(*e)
bui := BucketUsageInfo{
Size: uint64(flat.Size),
HotTierSize: uint64(max(flat.HotTierSize, 0)),
VersionsCount: flat.Versions,
ObjectsCount: flat.Objects,
DeleteMarkersCount: flat.DeleteMarkers,
@@ -1007,8 +980,8 @@ func (d *dataUsageCache) save(ctx context.Context, store objectIO, name string)
// Bumping the cache version will drop data from previous versions
// and write new data with the new version.
const (
dataUsageCacheVerCurrent = 9
dataUsageCacheVerV8 = 8
dataUsageCacheVerCurrent = 8
dataUsageCacheVerV9 = 9 // Retired access-tier cache; only adds the ignored "hts" entry key.
dataUsageCacheVerV7 = 7
dataUsageCacheVerV6 = 6
dataUsageCacheVerV5 = 5
@@ -1210,34 +1183,7 @@ func (d *dataUsageCache) deserialize(r io.Reader) error {
}
return nil
case dataUsageCacheVerV8:
// Zstd compressed.
dec, err := zstd.NewReader(r, zstd.WithDecoderConcurrency(2))
if err != nil {
return err
}
defer dec.Close()
dold := &dataUsageCacheV8{}
if err = dold.DecodeMsg(msgp.NewReader(dec)); err != nil {
return err
}
d.Info = dold.Info
d.Cache = make(map[string]dataUsageEntry, len(dold.Cache))
for k, v := range dold.Cache {
d.Cache[k] = dataUsageEntry{
Children: v.Children,
Size: v.Size,
Objects: v.Objects,
Versions: v.Versions,
DeleteMarkers: v.DeleteMarkers,
ObjSizes: v.ObjSizes,
ObjVersions: v.ObjVersions,
AllTierStats: v.AllTierStats,
Compacted: v.Compacted,
}
}
return nil
case dataUsageCacheVerCurrent:
case dataUsageCacheVerCurrent, dataUsageCacheVerV9:
// Zstd compressed.
dec, err := zstd.NewReader(r, zstd.WithDecoderConcurrency(2))
if err != nil {
+9 -605
View File
@@ -1591,145 +1591,6 @@ func (z *dataUsageCacheV7) Msgsize() (s int) {
return
}
// DecodeMsg implements msgp.Decodable
func (z *dataUsageCacheV8) DecodeMsg(dc *msgp.Reader) (err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err)
return
}
for zb0001 > 0 {
zb0001--
field, err = dc.ReadMapKeyPtr()
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "Info":
err = z.Info.DecodeMsg(dc)
if err != nil {
err = msgp.WrapError(err, "Info")
return
}
case "Cache":
var zb0002 uint32
zb0002, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err, "Cache")
return
}
if z.Cache == nil {
z.Cache = make(map[string]dataUsageEntryV8, zb0002)
} else if len(z.Cache) > 0 {
clear(z.Cache)
}
for zb0002 > 0 {
zb0002--
var za0001 string
za0001, err = dc.ReadString()
if err != nil {
err = msgp.WrapError(err, "Cache")
return
}
var za0002 dataUsageEntryV8
err = za0002.DecodeMsg(dc)
if err != nil {
err = msgp.WrapError(err, "Cache", za0001)
return
}
z.Cache[za0001] = za0002
}
default:
err = dc.Skip()
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
return
}
// UnmarshalMsg implements msgp.Unmarshaler
func (z *dataUsageCacheV8) UnmarshalMsg(bts []byte) (o []byte, err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
for zb0001 > 0 {
zb0001--
field, bts, err = msgp.ReadMapKeyZC(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "Info":
bts, err = z.Info.UnmarshalMsg(bts)
if err != nil {
err = msgp.WrapError(err, "Info")
return
}
case "Cache":
var zb0002 uint32
zb0002, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Cache")
return
}
if z.Cache == nil {
z.Cache = make(map[string]dataUsageEntryV8, zb0002)
} else if len(z.Cache) > 0 {
clear(z.Cache)
}
for zb0002 > 0 {
var za0002 dataUsageEntryV8
zb0002--
var za0001 string
za0001, bts, err = msgp.ReadStringBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Cache")
return
}
bts, err = za0002.UnmarshalMsg(bts)
if err != nil {
err = msgp.WrapError(err, "Cache", za0001)
return
}
z.Cache[za0001] = za0002
}
default:
bts, err = msgp.Skip(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
o = bts
return
}
// Msgsize returns an upper bound estimate of the number of bytes occupied by the serialized message
func (z *dataUsageCacheV8) Msgsize() (s int) {
s = 1 + 5 + z.Info.Msgsize() + 6 + msgp.MapHeaderSize
if z.Cache != nil {
for za0001, za0002 := range z.Cache {
_ = za0002
s += msgp.StringPrefixSize + len(za0001) + za0002.Msgsize()
}
}
return
}
// DecodeMsg implements msgp.Decodable
func (z *dataUsageEntry) DecodeMsg(dc *msgp.Reader) (err error) {
var field []byte
@@ -1762,12 +1623,6 @@ func (z *dataUsageEntry) DecodeMsg(dc *msgp.Reader) (err error) {
err = msgp.WrapError(err, "Size")
return
}
case "hts":
z.HotTierSize, err = dc.ReadInt64()
if err != nil {
err = msgp.WrapError(err, "HotTierSize")
return
}
case "os":
z.Objects, err = dc.ReadUint64()
if err != nil {
@@ -1946,12 +1801,12 @@ func (z *dataUsageEntry) DecodeMsg(dc *msgp.Reader) (err error) {
// EncodeMsg implements msgp.Encodable
func (z *dataUsageEntry) EncodeMsg(en *msgp.Writer) (err error) {
// check for omitted fields
zb0001Len := uint32(10)
var zb0001Mask uint16 /* 10 bits */
zb0001Len := uint32(9)
var zb0001Mask uint16 /* 9 bits */
_ = zb0001Mask
if z.AllTierStats == nil {
zb0001Len--
zb0001Mask |= 0x100
zb0001Mask |= 0x80
}
// variable map header, size zb0001Len
err = en.Append(0x80 | uint8(zb0001Len))
@@ -1981,16 +1836,6 @@ func (z *dataUsageEntry) EncodeMsg(en *msgp.Writer) (err error) {
err = msgp.WrapError(err, "Size")
return
}
// write "hts"
err = en.Append(0xa3, 0x68, 0x74, 0x73)
if err != nil {
return
}
err = en.WriteInt64(z.HotTierSize)
if err != nil {
err = msgp.WrapError(err, "HotTierSize")
return
}
// write "os"
err = en.Append(0xa2, 0x6f, 0x73)
if err != nil {
@@ -2055,7 +1900,7 @@ func (z *dataUsageEntry) EncodeMsg(en *msgp.Writer) (err error) {
return
}
}
if (zb0001Mask & 0x100) == 0 { // if not omitted
if (zb0001Mask & 0x80) == 0 { // if not omitted
// write "ats"
err = en.Append(0xa3, 0x61, 0x74, 0x73)
if err != nil {
@@ -2136,12 +1981,12 @@ func (z *dataUsageEntry) EncodeMsg(en *msgp.Writer) (err error) {
func (z *dataUsageEntry) MarshalMsg(b []byte) (o []byte, err error) {
o = msgp.Require(b, z.Msgsize())
// check for omitted fields
zb0001Len := uint32(10)
var zb0001Mask uint16 /* 10 bits */
zb0001Len := uint32(9)
var zb0001Mask uint16 /* 9 bits */
_ = zb0001Mask
if z.AllTierStats == nil {
zb0001Len--
zb0001Mask |= 0x100
zb0001Mask |= 0x80
}
// variable map header, size zb0001Len
o = append(o, 0x80|uint8(zb0001Len))
@@ -2158,9 +2003,6 @@ func (z *dataUsageEntry) MarshalMsg(b []byte) (o []byte, err error) {
// string "sz"
o = append(o, 0xa2, 0x73, 0x7a)
o = msgp.AppendInt64(o, z.Size)
// string "hts"
o = append(o, 0xa3, 0x68, 0x74, 0x73)
o = msgp.AppendInt64(o, z.HotTierSize)
// string "os"
o = append(o, 0xa2, 0x6f, 0x73)
o = msgp.AppendUint64(o, z.Objects)
@@ -2182,7 +2024,7 @@ func (z *dataUsageEntry) MarshalMsg(b []byte) (o []byte, err error) {
for za0002 := range z.ObjVersions {
o = msgp.AppendUint64(o, z.ObjVersions[za0002])
}
if (zb0001Mask & 0x100) == 0 { // if not omitted
if (zb0001Mask & 0x80) == 0 { // if not omitted
// string "ats"
o = append(o, 0xa3, 0x61, 0x74, 0x73)
if z.AllTierStats == nil {
@@ -2246,12 +2088,6 @@ func (z *dataUsageEntry) UnmarshalMsg(bts []byte) (o []byte, err error) {
err = msgp.WrapError(err, "Size")
return
}
case "hts":
z.HotTierSize, bts, err = msgp.ReadInt64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "HotTierSize")
return
}
case "os":
z.Objects, bts, err = msgp.ReadUint64Bytes(bts)
if err != nil {
@@ -2429,7 +2265,7 @@ func (z *dataUsageEntry) UnmarshalMsg(bts []byte) (o []byte, err error) {
// Msgsize returns an upper bound estimate of the number of bytes occupied by the serialized message
func (z *dataUsageEntry) Msgsize() (s int) {
s = 1 + 3 + z.Children.Msgsize() + 3 + msgp.Int64Size + 4 + msgp.Int64Size + 3 + msgp.Uint64Size + 3 + msgp.Uint64Size + 4 + msgp.Uint64Size + 4 + msgp.ArrayHeaderSize + (dataUsageBucketLen * (msgp.Uint64Size)) + 3 + msgp.ArrayHeaderSize + (dataUsageVersionLen * (msgp.Uint64Size)) + 4
s = 1 + 3 + z.Children.Msgsize() + 3 + msgp.Int64Size + 3 + msgp.Uint64Size + 3 + msgp.Uint64Size + 4 + msgp.Uint64Size + 4 + msgp.ArrayHeaderSize + (dataUsageBucketLen * (msgp.Uint64Size)) + 3 + msgp.ArrayHeaderSize + (dataUsageVersionLen * (msgp.Uint64Size)) + 4
if z.AllTierStats == nil {
s += msgp.NilSize
} else {
@@ -3422,438 +3258,6 @@ func (z *dataUsageEntryV7) Msgsize() (s int) {
return
}
// DecodeMsg implements msgp.Decodable
func (z *dataUsageEntryV8) DecodeMsg(dc *msgp.Reader) (err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err)
return
}
var zb0001Mask uint8 /* 1 bits */
_ = zb0001Mask
for zb0001 > 0 {
zb0001--
field, err = dc.ReadMapKeyPtr()
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "ch":
err = z.Children.DecodeMsg(dc)
if err != nil {
err = msgp.WrapError(err, "Children")
return
}
case "sz":
z.Size, err = dc.ReadInt64()
if err != nil {
err = msgp.WrapError(err, "Size")
return
}
case "os":
z.Objects, err = dc.ReadUint64()
if err != nil {
err = msgp.WrapError(err, "Objects")
return
}
case "vs":
z.Versions, err = dc.ReadUint64()
if err != nil {
err = msgp.WrapError(err, "Versions")
return
}
case "dms":
z.DeleteMarkers, err = dc.ReadUint64()
if err != nil {
err = msgp.WrapError(err, "DeleteMarkers")
return
}
case "szs":
var zb0002 uint32
zb0002, err = dc.ReadArrayHeader()
if err != nil {
err = msgp.WrapError(err, "ObjSizes")
return
}
if zb0002 != uint32(dataUsageBucketLen) {
err = msgp.ArrayError{Wanted: uint32(dataUsageBucketLen), Got: zb0002}
return
}
for za0001 := range z.ObjSizes {
z.ObjSizes[za0001], err = dc.ReadUint64()
if err != nil {
err = msgp.WrapError(err, "ObjSizes", za0001)
return
}
}
case "vh":
var zb0003 uint32
zb0003, err = dc.ReadArrayHeader()
if err != nil {
err = msgp.WrapError(err, "ObjVersions")
return
}
if zb0003 != uint32(dataUsageVersionLen) {
err = msgp.ArrayError{Wanted: uint32(dataUsageVersionLen), Got: zb0003}
return
}
for za0002 := range z.ObjVersions {
z.ObjVersions[za0002], err = dc.ReadUint64()
if err != nil {
err = msgp.WrapError(err, "ObjVersions", za0002)
return
}
}
case "ats":
if dc.IsNil() {
err = dc.ReadNil()
if err != nil {
err = msgp.WrapError(err, "AllTierStats")
return
}
z.AllTierStats = nil
} else {
if z.AllTierStats == nil {
z.AllTierStats = new(allTierStats)
}
var zb0004 uint32
zb0004, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err, "AllTierStats")
return
}
for zb0004 > 0 {
zb0004--
field, err = dc.ReadMapKeyPtr()
if err != nil {
err = msgp.WrapError(err, "AllTierStats")
return
}
switch msgp.UnsafeString(field) {
case "ts":
var zb0005 uint32
zb0005, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers")
return
}
if z.AllTierStats.Tiers == nil {
z.AllTierStats.Tiers = make(map[string]tierStats, zb0005)
} else if len(z.AllTierStats.Tiers) > 0 {
clear(z.AllTierStats.Tiers)
}
for zb0005 > 0 {
zb0005--
var za0003 string
za0003, err = dc.ReadString()
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers")
return
}
var za0004 tierStats
var zb0006 uint32
zb0006, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003)
return
}
for zb0006 > 0 {
zb0006--
field, err = dc.ReadMapKeyPtr()
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003)
return
}
switch msgp.UnsafeString(field) {
case "ts":
za0004.TotalSize, err = dc.ReadUint64()
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003, "TotalSize")
return
}
case "nv":
za0004.NumVersions, err = dc.ReadInt()
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003, "NumVersions")
return
}
case "no":
za0004.NumObjects, err = dc.ReadInt()
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003, "NumObjects")
return
}
default:
err = dc.Skip()
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003)
return
}
}
}
z.AllTierStats.Tiers[za0003] = za0004
}
default:
err = dc.Skip()
if err != nil {
err = msgp.WrapError(err, "AllTierStats")
return
}
}
}
}
zb0001Mask |= 0x1
case "c":
z.Compacted, err = dc.ReadBool()
if err != nil {
err = msgp.WrapError(err, "Compacted")
return
}
default:
err = dc.Skip()
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
// Clear omitted fields.
if (zb0001Mask & 0x1) == 0 {
z.AllTierStats = nil
}
return
}
// UnmarshalMsg implements msgp.Unmarshaler
func (z *dataUsageEntryV8) UnmarshalMsg(bts []byte) (o []byte, err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
var zb0001Mask uint8 /* 1 bits */
_ = zb0001Mask
for zb0001 > 0 {
zb0001--
field, bts, err = msgp.ReadMapKeyZC(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "ch":
bts, err = z.Children.UnmarshalMsg(bts)
if err != nil {
err = msgp.WrapError(err, "Children")
return
}
case "sz":
z.Size, bts, err = msgp.ReadInt64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "Size")
return
}
case "os":
z.Objects, bts, err = msgp.ReadUint64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "Objects")
return
}
case "vs":
z.Versions, bts, err = msgp.ReadUint64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "Versions")
return
}
case "dms":
z.DeleteMarkers, bts, err = msgp.ReadUint64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "DeleteMarkers")
return
}
case "szs":
var zb0002 uint32
zb0002, bts, err = msgp.ReadArrayHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "ObjSizes")
return
}
if zb0002 != uint32(dataUsageBucketLen) {
err = msgp.ArrayError{Wanted: uint32(dataUsageBucketLen), Got: zb0002}
return
}
for za0001 := range z.ObjSizes {
z.ObjSizes[za0001], bts, err = msgp.ReadUint64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "ObjSizes", za0001)
return
}
}
case "vh":
var zb0003 uint32
zb0003, bts, err = msgp.ReadArrayHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "ObjVersions")
return
}
if zb0003 != uint32(dataUsageVersionLen) {
err = msgp.ArrayError{Wanted: uint32(dataUsageVersionLen), Got: zb0003}
return
}
for za0002 := range z.ObjVersions {
z.ObjVersions[za0002], bts, err = msgp.ReadUint64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "ObjVersions", za0002)
return
}
}
case "ats":
if msgp.IsNil(bts) {
bts, err = msgp.ReadNilBytes(bts)
if err != nil {
return
}
z.AllTierStats = nil
} else {
if z.AllTierStats == nil {
z.AllTierStats = new(allTierStats)
}
var zb0004 uint32
zb0004, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats")
return
}
for zb0004 > 0 {
zb0004--
field, bts, err = msgp.ReadMapKeyZC(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats")
return
}
switch msgp.UnsafeString(field) {
case "ts":
var zb0005 uint32
zb0005, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers")
return
}
if z.AllTierStats.Tiers == nil {
z.AllTierStats.Tiers = make(map[string]tierStats, zb0005)
} else if len(z.AllTierStats.Tiers) > 0 {
clear(z.AllTierStats.Tiers)
}
for zb0005 > 0 {
var za0004 tierStats
zb0005--
var za0003 string
za0003, bts, err = msgp.ReadStringBytes(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers")
return
}
var zb0006 uint32
zb0006, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003)
return
}
for zb0006 > 0 {
zb0006--
field, bts, err = msgp.ReadMapKeyZC(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003)
return
}
switch msgp.UnsafeString(field) {
case "ts":
za0004.TotalSize, bts, err = msgp.ReadUint64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003, "TotalSize")
return
}
case "nv":
za0004.NumVersions, bts, err = msgp.ReadIntBytes(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003, "NumVersions")
return
}
case "no":
za0004.NumObjects, bts, err = msgp.ReadIntBytes(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003, "NumObjects")
return
}
default:
bts, err = msgp.Skip(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats", "Tiers", za0003)
return
}
}
}
z.AllTierStats.Tiers[za0003] = za0004
}
default:
bts, err = msgp.Skip(bts)
if err != nil {
err = msgp.WrapError(err, "AllTierStats")
return
}
}
}
}
zb0001Mask |= 0x1
case "c":
z.Compacted, bts, err = msgp.ReadBoolBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Compacted")
return
}
default:
bts, err = msgp.Skip(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
// Clear omitted fields.
if (zb0001Mask & 0x1) == 0 {
z.AllTierStats = nil
}
o = bts
return
}
// Msgsize returns an upper bound estimate of the number of bytes occupied by the serialized message
func (z *dataUsageEntryV8) Msgsize() (s int) {
s = 1 + 3 + z.Children.Msgsize() + 3 + msgp.Int64Size + 3 + msgp.Uint64Size + 3 + msgp.Uint64Size + 4 + msgp.Uint64Size + 4 + msgp.ArrayHeaderSize + (dataUsageBucketLen * (msgp.Uint64Size)) + 3 + msgp.ArrayHeaderSize + (dataUsageVersionLen * (msgp.Uint64Size)) + 4
if z.AllTierStats == nil {
s += msgp.NilSize
} else {
s += 1 + 3 + msgp.MapHeaderSize
if z.AllTierStats.Tiers != nil {
for za0003, za0004 := range z.AllTierStats.Tiers {
_ = za0004
s += msgp.StringPrefixSize + len(za0003) + 1 + 3 + msgp.Uint64Size + 3 + msgp.IntSize + 3 + msgp.IntSize
}
}
}
s += 2 + msgp.BoolSize
return
}
// DecodeMsg implements msgp.Decodable
func (z *dataUsageHash) DecodeMsg(dc *msgp.Reader) (err error) {
{
+1 -6
View File
@@ -46,8 +46,7 @@ type BucketTargetUsageInfo struct {
// - total objects in a bucket
// - object size histogram per bucket
type BucketUsageInfo struct {
Size uint64 `json:"size"`
HotTierSize uint64 `json:"hotTierSize,omitempty"`
Size uint64 `json:"size"`
// Following five fields suffixed with V1 are here for backward compatibility
// Total Size for objects that have not yet been replicated
ReplicationPendingSizeV1 uint64 `json:"objectsPendingReplicationTotalSize"`
@@ -79,10 +78,6 @@ type DataUsageInfo struct {
// LastUpdate is the timestamp of when the data usage info was last updated.
// This does not indicate a full scan.
LastUpdate time.Time `json:"lastUpdate"`
// ScannerCycle changes only after a complete scanner pass. Background
// consumers use it to distinguish a partial cache update from a baseline
// that has visited every bucket and server pool.
ScannerCycle uint32 `json:"scannerCycle,omitempty"`
// Objects total count across all buckets
ObjectsTotalCount uint64 `json:"objectsCount"`
+17 -5
View File
@@ -108,7 +108,7 @@ func TestPoolsConditionalDeleteReportsOtherPoolFailure(t *testing.T) {
getDisks := set.getDisks
faulty := append([]StorageAPI(nil), getDisks()...)
for i := range faulty {
faulty[i] = accessMoveDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object}
faulty[i] = consistencyDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object}
}
set.getDisks = func() []StorageAPI { return faulty }
defer func() { set.getDisks = getDisks }()
@@ -135,9 +135,8 @@ func testPoolsConditionalDeleteWriter(t *testing.T, multipart bool) {
ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second)
defer cancel()
reader := mustGetPutObjReader(t, bytes.NewBufferString("after"), 5, "", "")
destination := 1
write := func() error {
_, err := z.PutObject(ctx, bucket, object, reader, ObjectOptions{DataMovement: true, SrcPoolIdx: 0, DstPoolIdx: &destination})
_, err := z.PutObject(ctx, bucket, object, reader, ObjectOptions{})
return err
}
if multipart {
@@ -548,7 +547,7 @@ func TestPoolsReplicaCleanupFailureCanRetry(t *testing.T) {
getDisks := set.getDisks
faulty := append([]StorageAPI(nil), getDisks()...)
for i := range faulty {
faulty[i] = accessMoveDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object, version: oi.VersionID}
faulty[i] = consistencyDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object, version: oi.VersionID}
}
set.getDisks = func() []StorageAPI { return faulty }
defer func() { set.getDisks = getDisks }()
@@ -622,7 +621,7 @@ func TestPoolsRetiringCopyPreservesSharedTierObject(t *testing.T) {
getDisks := set.getDisks
faulty := append([]StorageAPI(nil), getDisks()...)
for i := range faulty {
faulty[i] = accessMoveDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object, version: oi.VersionID}
faulty[i] = consistencyDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object, version: oi.VersionID}
}
set.getDisks = func() []StorageAPI { return faulty }
defer func() { set.getDisks = getDisks }()
@@ -680,3 +679,16 @@ func TestPoolsRetiringCopyPreservesSharedTierObject(t *testing.T) {
})
}
}
// Fault injection shared by general multi-pool regressions.
type consistencyDeleteFaultDisk struct {
StorageAPI
bucket, object, version string
}
func (d consistencyDeleteFaultDisk) 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)
}
+23 -96
View File
@@ -615,17 +615,7 @@ func (z *erasureServerPools) getPoolIdxExistingNoLock(ctx context.Context, bucke
})
}
func (z *erasureServerPools) getPoolIdxNoLock(ctx context.Context, bucket, object string, size int64, dstPoolIdx *int) (idx int, err error) {
if dstPoolIdx != nil {
if *dstPoolIdx < 0 || *dstPoolIdx >= len(z.serverPools) {
return -1, errInvalidArgument
}
if z.IsSuspended(*dstPoolIdx) || z.IsPoolRebalancing(*dstPoolIdx) {
return -1, toObjectErr(errDiskFull)
}
return *dstPoolIdx, nil
}
func (z *erasureServerPools) getPoolIdxNoLock(ctx context.Context, bucket, object string, size int64) (idx int, err error) {
idx, err = z.getPoolIdxExistingNoLock(ctx, bucket, object)
if err != nil && !isErrObjectNotFound(err) {
return idx, err
@@ -644,23 +634,13 @@ func (z *erasureServerPools) getPoolIdxNoLock(ctx context.Context, bucket, objec
// getPoolIdx returns the found previous object and its corresponding pool idx,
// if none are found falls back to most available space pool, this function is
// designed to be only used by PutObject, CopyObject (newObject creation) and NewMultipartUpload.
func (z *erasureServerPools) getPoolIdx(ctx context.Context, bucket, object string, size int64, dstPoolIdx *int) (idx int, err error) {
return z.getWritePoolIdx(ctx, bucket, object, size, dstPoolIdx, false)
func (z *erasureServerPools) getPoolIdx(ctx context.Context, bucket, object string, size int64) (idx int, err error) {
return z.getWritePoolIdx(ctx, bucket, object, size, false)
}
// getWritePoolIdx keeps the write-allocation policy when the caller already
// holds the pools-layer object lock and must not reacquire a set read lock.
func (z *erasureServerPools) getWritePoolIdx(ctx context.Context, bucket, object string, size int64, dstPoolIdx *int, noLock bool) (idx int, err error) {
if dstPoolIdx != nil {
if *dstPoolIdx < 0 || *dstPoolIdx >= len(z.serverPools) {
return -1, errInvalidArgument
}
if z.IsSuspended(*dstPoolIdx) || z.IsPoolRebalancing(*dstPoolIdx) {
return -1, toObjectErr(errDiskFull)
}
return *dstPoolIdx, nil
}
func (z *erasureServerPools) getWritePoolIdx(ctx context.Context, bucket, object string, size int64, noLock bool) (idx int, err error) {
pinfo, _, err := z.getPoolInfoExistingWithOpts(ctx, bucket, object, ObjectOptions{
NoLock: noLock,
SkipDecommissioned: true,
@@ -683,13 +663,6 @@ func (z *erasureServerPools) getWritePoolIdx(ctx context.Context, bucket, object
return idx, nil
}
func dataMovementDstPool(opts ObjectOptions) *int {
if !opts.DataMovement {
return nil
}
return opts.DstPoolIdx
}
func (z *erasureServerPools) Shutdown(ctx context.Context) error {
g := errgroup.WithNErrs(len(z.serverPools))
@@ -1158,16 +1131,10 @@ func (z *erasureServerPools) PutObject(ctx context.Context, bucket string, objec
object = encodeDirObject(object)
if z.SinglePool() {
idx, err := z.getPoolIdx(ctx, bucket, object, data.Size(), dataMovementDstPool(opts))
_, err := z.getPoolIdx(ctx, bucket, object, data.Size())
if err != nil {
return ObjectInfo{}, err
}
if dataMovementDstPool(opts) != nil && idx == opts.SrcPoolIdx {
return ObjectInfo{}, DataMovementOverwriteErr{
Bucket: bucket, Object: object, VersionID: opts.VersionID,
Err: errDataMovementSrcDstPoolSame,
}
}
return z.serverPools[0].PutObject(ctx, bucket, object, data, opts)
}
@@ -1182,7 +1149,7 @@ func (z *erasureServerPools) PutObject(ctx context.Context, bucket string, objec
}
opts.NoLock = true
idx, err := z.getWritePoolIdx(ctx, bucket, object, data.Size(), dataMovementDstPool(opts), true)
idx, err := z.getWritePoolIdx(ctx, bucket, object, data.Size(), true)
if err != nil {
return ObjectInfo{}, err
}
@@ -1238,28 +1205,6 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob
return ObjectInfo{}, z.deletePrefix(ctx, bucket, object)
}
// Access-tier moves must recreate delete markers on the explicitly
// selected destination. The regular data-movement path discovers a pool
// from existing object state, which is ambiguous while both source and
// destination temporarily contain the version stack.
if dstPoolIdx := dataMovementDstPool(opts); dstPoolIdx != nil {
if *dstPoolIdx < 0 || *dstPoolIdx >= len(z.serverPools) {
return ObjectInfo{}, errInvalidArgument
}
if *dstPoolIdx == opts.SrcPoolIdx {
return ObjectInfo{}, DataMovementOverwriteErr{
Bucket: bucket, Object: decodeDirObject(object), VersionID: opts.VersionID,
Err: errDataMovementSrcDstPoolSame,
}
}
if z.IsSuspended(*dstPoolIdx) || z.IsPoolRebalancing(*dstPoolIdx) {
return ObjectInfo{}, toObjectErr(errDiskFull)
}
objInfo, err = z.serverPools[*dstPoolIdx].DeleteObject(ctx, bucket, object, opts)
objInfo.Name = decodeDirObject(object)
return objInfo, err
}
if !z.SinglePool() && opts.CheckPrecondFn != nil {
return z.deleteObjectConditional(ctx, bucket, object, opts)
}
@@ -1508,16 +1453,10 @@ func (z *erasureServerPools) CopyObject(ctx context.Context, srcBucket, srcObjec
return oi, err
}
poolIdx, err := z.getPoolIdxNoLock(ctx, dstBucket, dstObject, srcInfo.Size, dataMovementDstPool(dstOpts))
poolIdx, err := z.getPoolIdxNoLock(ctx, dstBucket, dstObject, srcInfo.Size)
if err != nil {
return objInfo, err
}
if dataMovementDstPool(dstOpts) != nil && poolIdx == dstOpts.SrcPoolIdx {
return ObjectInfo{}, DataMovementOverwriteErr{
Bucket: dstBucket, Object: dstObject, VersionID: dstOpts.VersionID,
Err: errDataMovementSrcDstPoolSame,
}
}
if !z.SinglePool() && dstOpts.ReplicaLockReconcile {
dstOpts.replicaObjectInfo = z.replicaObjectInfo
@@ -1977,41 +1916,29 @@ func (z *erasureServerPools) NewMultipartUpload(ctx context.Context, bucket, obj
}()
if z.SinglePool() {
idx, err := z.getPoolIdx(ctx, bucket, object, -1, dataMovementDstPool(opts))
if err != nil {
return nil, err
}
if dataMovementDstPool(opts) != nil && idx == opts.SrcPoolIdx {
return nil, DataMovementOverwriteErr{
Bucket: bucket, Object: object, VersionID: opts.VersionID,
Err: errDataMovementSrcDstPoolSame,
}
}
return z.serverPools[0].NewMultipartUpload(ctx, bucket, object, opts)
}
if dataMovementDstPool(opts) == nil {
for idx, pool := range z.serverPools {
if z.IsSuspended(idx) || z.IsPoolRebalancing(idx) {
continue
}
for idx, pool := range z.serverPools {
if z.IsSuspended(idx) || z.IsPoolRebalancing(idx) {
continue
}
result, err := pool.ListMultipartUploads(ctx, bucket, object, "", "", "", maxUploadsList)
if err != nil {
return nil, err
}
// If there is a multipart upload with the same bucket/object name,
// create the new multipart in the same pool, this will avoid
// creating two multiparts uploads in two different pools.
if len(result.Uploads) != 0 {
return z.serverPools[idx].NewMultipartUpload(ctx, bucket, object, opts)
}
result, err := pool.ListMultipartUploads(ctx, bucket, object, "", "", "", maxUploadsList)
if err != nil {
return nil, err
}
// If there is a multipart upload with the same bucket/object name,
// create the new multipart in the same pool, this will avoid
// creating two multiparts uploads in two different pools
if len(result.Uploads) != 0 {
return z.serverPools[idx].NewMultipartUpload(ctx, bucket, object, opts)
}
}
// any parallel writes on the object will block for this poolIdx
// to return since this holds a read lock on the namespace.
idx, err := z.getPoolIdx(ctx, bucket, object, -1, dataMovementDstPool(opts))
idx, err := z.getPoolIdx(ctx, bucket, object, -1)
if err != nil {
return nil, err
}
@@ -2045,7 +1972,7 @@ func (z *erasureServerPools) PutObjectPart(ctx context.Context, bucket, object,
}
if z.SinglePool() {
_, err := z.getPoolIdx(ctx, bucket, object, data.Size(), dataMovementDstPool(opts))
_, err := z.getPoolIdx(ctx, bucket, object, data.Size())
if err != nil {
return PartInfo{}, err
}
@@ -3219,7 +3146,7 @@ func (z *erasureServerPools) DecomTieredObject(ctx context.Context, bucket, obje
defer ns.Unlock(lkctx)
opts.NoLock = true
}
idx, err := z.getPoolIdxNoLock(ctx, bucket, object, fi.Size, dataMovementDstPool(opts))
idx, err := z.getPoolIdxNoLock(ctx, bucket, object, fi.Size)
if err != nil {
return err
}
-529
View File
@@ -1,529 +0,0 @@
// 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()...)
// Commit removal of the old version, then reject the latest version on
// every disk. This fixes the retry boundary without relying on how a
// partially deleted erasure quorum is subsequently resolved or healed.
for i := range faulty {
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 := z.serverPools[1].GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: old.VersionID}); !isErrVersionNotFound(err) {
t.Fatalf("old source version should already be removed: %v", err)
}
assertAccessMoveVersion(t, z, bucket, object, 1, 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")
}
}
File diff suppressed because it is too large Load Diff
-263
View File
@@ -1,263 +0,0 @@
// Copyright (c) 2015-2026 MinIO, Inc.
package cmd
import (
"bytes"
"context"
"errors"
"strconv"
"strings"
"testing"
"time"
"github.com/klauspost/compress/zstd"
"github.com/minio/minio/internal/bucket/lifecycle"
"github.com/minio/minio/internal/config/ilm"
"github.com/tinylib/msgp/msgp"
)
func accessLifecycleForTest(t *testing.T) *lifecycle.Lifecycle {
t.Helper()
xml := "<LifecycleConfiguration>" +
"<Rule><ID>access</ID><Status>Enabled</Status>" +
"<AccessTransition><Window>10m</Window><PromoteAfterAccesses>100</PromoteAfterAccesses>" +
"<DemoteAfterAccesses>5</DemoteAfterAccesses><DemoteAfterIdle>1h</DemoteAfterIdle></AccessTransition>" +
"</Rule></LifecycleConfiguration>"
lc, err := lifecycle.ParseLifecycleConfig(strings.NewReader(xml))
if err != nil {
t.Fatal(err)
}
return lc
}
func TestAccessTierStampAndDemoteEligibility(t *testing.T) {
oldTracker := globalAccessTracker
globalAccessTracker = newAccessTracker()
t.Cleanup(func() { globalAccessTracker = oldTracker })
now := time.Now()
cfg := ilm.Config{
AccessTiering: true, AccessPools: []int{0, 1},
AccessBinWidth: time.Minute, AccessBins: 12,
AccessMinResidency: 30 * time.Minute,
}
oi := ObjectInfo{
Bucket: "bucket", Name: "object", Size: 10, IsLatest: true,
UserDefined: map[string]string{
accessTierMetadataKey: "0:" + strconv.FormatInt(now.Add(-2*time.Hour).UnixNano(), 10),
},
}
if !accessDemoteEligible(accessLifecycleForTest(t), oi, 0, cfg, now) {
t.Fatal("cold, resident object should be demotion eligible")
}
oi.UserDefined[accessTierMetadataKey] = "0:" + strconv.FormatInt(now.Add(-time.Minute).UnixNano(), 10)
if accessDemoteEligible(accessLifecycleForTest(t), oi, 0, cfg, now) {
t.Fatal("object inside minimum residency was eligible")
}
oi.UserDefined[accessTierMetadataKey] = "broken"
if accessDemoteEligible(accessLifecycleForTest(t), oi, 0, cfg, now) {
t.Fatal("object with malformed marker was eligible")
}
}
func TestAccessTierReservationsEnforceCaps(t *testing.T) {
state := newAccessTierState(t.Context())
state.usageReady = true
state.baseUsage["a"] = 80
state.baseTotal = 180
cfg := ilm.Config{AccessMaxSize: 200}
task := accessTierTask{ctx: t.Context(), bucket: "a", object: "one", bytes: 30}
if reason, ok := state.reservePromotion(task, cfg, 0); ok || reason != "max-size" {
t.Fatalf("max-size reserve = %q/%v", reason, ok)
}
cfg.AccessMaxSize = 0
if reason, ok := state.reservePromotion(task, cfg, 100); ok || reason != "quota" {
t.Fatalf("quota reserve = %q/%v", reason, ok)
}
task.bytes = 20
if reason, ok := state.reservePromotion(task, cfg, 100); !ok || reason != "" {
t.Fatalf("valid reserve = %q/%v", reason, ok)
}
if got := state.bucketUsageLocked("a"); got != 100 {
t.Fatalf("reserved bucket usage = %d, want 100", got)
}
state.releasePending(task, true)
}
func TestDataMovementDestinationIsGated(t *testing.T) {
dst := 0
if got := dataMovementDstPool(ObjectOptions{DstPoolIdx: &dst}); got != nil {
t.Fatal("destination honored without DataMovement")
}
if got := dataMovementDstPool(ObjectOptions{DataMovement: true, DstPoolIdx: &dst}); got == nil || *got != 0 {
t.Fatalf("destination = %v", got)
}
}
func TestForcedDestinationPoolSelection(t *testing.T) {
z := &erasureServerPools{serverPools: make([]*erasureSets, 2)}
dst := 1
if got, err := z.getPoolIdx(context.Background(), "bucket", "object", 1, &dst); err != nil || got != dst {
t.Fatalf("destination = %d, err = %v", got, err)
}
bad := 2
if _, err := z.getPoolIdx(context.Background(), "bucket", "object", 1, &bad); !errors.Is(err, errInvalidArgument) {
t.Fatalf("out-of-range error = %v", err)
}
z.poolMeta.Pools = make([]PoolStatus, 2)
z.poolMeta.Pools[1].Decommission = &PoolDecommissionInfo{}
if _, err := z.getPoolIdx(context.Background(), "bucket", "object", 1, &dst); err == nil {
t.Fatal("suspended destination was accepted")
}
}
func TestHotTierAccounting(t *testing.T) {
entry := dataUsageEntry{}
entry.addSizes(sizeSummary{totalSize: 100, hotTierSize: 60, versions: 1})
entry.merge(dataUsageEntry{Size: 50, HotTierSize: 25})
if entry.Size != 150 || entry.HotTierSize != 85 {
t.Fatalf("usage = size:%d hot:%d", entry.Size, entry.HotTierSize)
}
}
func TestScannerCyclesApart(t *testing.T) {
tests := []struct {
current, previous uint32
want uint32
}{
{current: 12, previous: 10, want: 2},
{current: 0, previous: ^uint32(0), want: 1},
// A cycle counter reset is not evidence that a complete pass covered
// recent moves, so it must not release conservative deltas.
{current: 1, previous: 100, want: 0},
}
for _, tt := range tests {
if got := scannerCyclesApart(tt.current, tt.previous); got != tt.want {
t.Fatalf("scannerCyclesApart(%d, %d) = %d, want %d", tt.current, tt.previous, got, tt.want)
}
}
}
func TestDataUsageCacheV8Migration(t *testing.T) {
var encoded bytes.Buffer
if err := encoded.WriteByte(dataUsageCacheVerV8); err != nil {
t.Fatal(err)
}
zw, err := zstd.NewWriter(&encoded)
if err != nil {
t.Fatal(err)
}
mw := msgp.NewWriter(zw)
if err = mw.WriteMapHeader(2); err != nil {
t.Fatal(err)
}
if err = mw.WriteString("Info"); err != nil {
t.Fatal(err)
}
info := dataUsageCacheInfo{Name: dataUsageRoot, NextCycle: 17, LastUpdate: time.Now().UTC()}
if err = info.EncodeMsg(mw); err != nil {
t.Fatal(err)
}
if err = mw.WriteString("Cache"); err != nil {
t.Fatal(err)
}
if err = mw.WriteMapHeader(1); err != nil {
t.Fatal(err)
}
if err = mw.WriteString("entry"); err != nil {
t.Fatal(err)
}
// A v8 entry has no hts field. Encoding only populated fields also
// verifies that its map decoder retains normal msgpack compatibility.
if err = mw.WriteMapHeader(2); err != nil {
t.Fatal(err)
}
if err = mw.WriteString("sz"); err != nil {
t.Fatal(err)
}
if err = mw.WriteInt64(123); err != nil {
t.Fatal(err)
}
if err = mw.WriteString("os"); err != nil {
t.Fatal(err)
}
if err = mw.WriteUint64(2); err != nil {
t.Fatal(err)
}
if err = mw.Flush(); err != nil {
t.Fatal(err)
}
if err = zw.Close(); err != nil {
t.Fatal(err)
}
var migrated dataUsageCache
if err = migrated.deserialize(bytes.NewReader(encoded.Bytes())); err != nil {
t.Fatal(err)
}
entry := migrated.Cache["entry"]
if entry.Size != 123 || entry.Objects != 2 || entry.HotTierSize != 0 {
t.Fatalf("migrated entry = size:%d objects:%d hot:%d", entry.Size, entry.Objects, entry.HotTierSize)
}
if migrated.Info.NextCycle != 17 || !migrated.Info.LastUpdate.Equal(info.LastUpdate) {
t.Fatalf("migrated cache info = %+v", migrated.Info)
}
}
// The stamp written on every moved version and the stamp the rollback path
// matches against must agree, otherwise rollback silently skips its own work.
func TestAccessTierStampRoundTrip(t *testing.T) {
movedAt := time.Now().UnixNano()
stamp := accessTierStamp(3, movedAt)
pool, at, ok := parseAccessTierStamp(map[string]string{accessTierMetadataKey: stamp})
if !ok {
t.Fatalf("stamp %q did not parse", stamp)
}
if pool != 3 {
t.Fatalf("pool = %d, want 3", pool)
}
if at.UnixNano() != movedAt {
t.Fatalf("movedAt = %d, want %d", at.UnixNano(), movedAt)
}
// A different move of the same object must not match, so a concurrent
// client overwrite is never mistaken for our own copy.
if stamp == accessTierStamp(3, movedAt+1) {
t.Fatal("stamps from distinct moves collided")
}
if stamp == accessTierStamp(4, movedAt) {
t.Fatal("stamps from distinct pools collided")
}
}
func TestAccessHitsMightPromote(t *testing.T) {
lc := accessLifecycleForTest(t)
hot := accessEntry{Bins: []uint32{100}, HeadAt: 0}
cold := accessEntry{Bins: []uint32{1}, HeadAt: 0}
if !accessHitsMightPromote(lc, "object", hot, 60) {
t.Fatal("object above promote threshold was rejected")
}
if accessHitsMightPromote(lc, "object", cold, 60) {
t.Fatal("object below promote threshold was accepted")
}
xml := "<LifecycleConfiguration><Rule><ID>logs</ID><Status>Enabled</Status>" +
"<Filter><Prefix>logs/</Prefix></Filter>" +
"<AccessTransition><Window>10m</Window><PromoteAfterAccesses>100</PromoteAfterAccesses>" +
"<DemoteAfterAccesses>5</DemoteAfterAccesses><DemoteAfterIdle>1h</DemoteAfterIdle></AccessTransition>" +
"</Rule></LifecycleConfiguration>"
prefixed, err := lifecycle.ParseLifecycleConfig(strings.NewReader(xml))
if err != nil {
t.Fatal(err)
}
if accessHitsMightPromote(prefixed, "data/object", hot, 60) {
t.Fatal("object outside the rule prefix was accepted")
}
if !accessHitsMightPromote(prefixed, "logs/object", hot, 60) {
t.Fatal("object inside the rule prefix was rejected")
}
}
-598
View File
@@ -1,598 +0,0 @@
// Copyright (c) 2015-2026 MinIO, Inc.
//
// This file is part of MinIO 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 (
"cmp"
"context"
"fmt"
"slices"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/minio/minio/internal/config/ilm"
"github.com/zeebo/xxh3"
)
//go:generate msgp -file=$GOFILE -unexported
//msgp:ignore accessTracker mergedAccess
const (
// accessTrackerPrefix is where each node publishes its own counters.
// One object per node, merged by every node on a timer.
accessTrackerPrefix = minioConfigPrefix + "/ilm/access"
// accessQueueSize bounds the GET -> tracker handoff. Overflow drops
// samples rather than slowing down reads.
accessQueueSize = 100000
// accessMaxDemoteCandidates bounds what the scanner may hand over in a
// single flush interval.
accessMaxDemoteCandidates = 100000
// accessShardStaleFactor multiplies the flush interval to decide when a
// peer's counters are too old to trust, e.g. after a node is removed.
accessShardStaleFactor = 5
)
func accessShardFresh(now int64, shard accessShard, stale int64) bool {
if shard.UpdatedAt <= 0 {
return false
}
if stale <= 0 {
return true
}
age := now - shard.UpdatedAt
return age >= -stale && age <= stale
}
// accessEntry is a rolling hit counter for one object. Bins[0] is the current
// bin and each subsequent bin is one bin-width older, so a rule asking for
// "100 hits in 10 minutes" sums the newest ceil(10m/binWidth) bins.
//
// A fixed window is used rather than an exponentially decayed score because
// the rule is stated to operators in exactly those terms.
type accessEntry struct {
Bins []uint32 `msg:"b"`
HeadAt int64 `msg:"h"` // unix seconds at the start of Bins[0]
LastAt int64 `msg:"l"` // unix seconds of the most recent hit
}
// demoteCandidate is an object the scanner found sitting on a fast pool with
// no recent reads. Candidates ride the node's own counter shard so they reach
// the leader without a new peer RPC.
type demoteCandidate struct {
Bucket string `msg:"b"`
Object string `msg:"o"`
Pool int `msg:"p"`
}
// accessShard is what one node publishes. BinWidth is carried so a peer that
// has not yet picked up a configuration change is ignored rather than merged
// with mismatched bins.
type accessShard struct {
UpdatedAt int64 `msg:"u"`
BinWidth int64 `msg:"bw"`
Entries map[string]accessEntry `msg:"e"`
Demote []demoteCandidate `msg:"d"`
}
// binStart truncates a unix timestamp to the start of its bin.
func binStart(now, binWidth int64) int64 {
if binWidth <= 0 {
return now
}
return now - now%binWidth
}
// rollTo advances the counter to now, zeroing the bins that elapsed since the
// last update and resizing if the configured bin count changed.
func (e *accessEntry) rollTo(now, binWidth int64, nbins int) {
if nbins <= 0 || binWidth <= 0 {
return
}
if len(e.Bins) != nbins {
resized := make([]uint32, nbins)
copy(resized, e.Bins)
e.Bins = resized
}
head := binStart(now, binWidth)
if e.HeadAt == 0 {
e.HeadAt = head
return
}
steps := (head - e.HeadAt) / binWidth
if steps <= 0 {
return
}
if steps >= int64(nbins) {
clear(e.Bins)
} else {
copy(e.Bins[steps:], e.Bins[:nbins-int(steps)])
clear(e.Bins[:steps])
}
e.HeadAt = head
}
// hits returns the number of accesses recorded over the newest bins covering
// window. A window longer than the configured history is clamped to it.
//
// The newest bin is partial, so the covered span is between window-binWidth
// and window. Operators tune resolution with ilm access_bin_width.
func (e accessEntry) hits(window time.Duration, binWidth int64) uint64 {
if binWidth <= 0 || len(e.Bins) == 0 {
return 0
}
n := int((int64(window/time.Second) + binWidth - 1) / binWidth)
if n < 1 {
n = 1
}
if n > len(e.Bins) {
n = len(e.Bins)
}
var total uint64
for _, v := range e.Bins[:n] {
total += uint64(v)
}
return total
}
// total is the whole retained history, used to decide what to evict.
func (e accessEntry) total() uint64 {
var t uint64
for _, v := range e.Bins {
t += uint64(v)
}
return t
}
// mergeFrom adds another node's counters for the same object. Both sides must
// already be rolled to the same head.
func (e *accessEntry) mergeFrom(o accessEntry) {
for i := range e.Bins {
if i < len(o.Bins) {
total := uint64(e.Bins[i]) + uint64(o.Bins[i])
if total > uint64(^uint32(0)) {
total = uint64(^uint32(0))
}
e.Bins[i] = uint32(total)
}
}
if o.LastAt > e.LastAt {
e.LastAt = o.LastAt
}
}
// mergedAccess is an immutable cluster-wide snapshot. Readers take it from an
// atomic pointer, so the hot scanner and sweep paths never take a lock.
type mergedAccess struct {
entries map[string]accessEntry
binWidth int64
at int64
}
func (m *mergedAccess) hits(key string, window time.Duration) uint64 {
if m == nil {
return 0
}
e, ok := m.entries[key]
if !ok {
return 0
}
return e.hits(window, m.binWidth)
}
func (m *mergedAccess) lastAccess(key string) int64 {
if m == nil {
return 0
}
return m.entries[key].LastAt
}
// accessTracker records how often each object is read.
//
// Ownership is deliberately narrow: the live counter map is touched only by
// run()'s goroutine, so it needs no lock. Everything read from elsewhere goes
// through the immutable merged snapshot.
type accessTracker struct {
ch chan string
enabled atomic.Bool
merged atomic.Pointer[mergedAccess]
// Demote candidates arrive from scanner goroutines, so this one does
// need a lock. It is small: only objects we previously promoted.
demoteMu sync.Mutex
demote map[string]demoteCandidate
dropped atomic.Uint64
}
var globalAccessTracker = newAccessTracker()
func newAccessTracker() *accessTracker {
return &accessTracker{
ch: make(chan string, accessQueueSize),
demote: make(map[string]demoteCandidate),
}
}
// accessKey is the tracker's map key. Bucket names cannot contain '/', so the
// join is unambiguous.
func accessKey(bucket, object string) string {
return bucket + "/" + object
}
func splitAccessKey(key string) (bucket, object string, ok bool) {
bucket, object, ok = strings.Cut(key, "/")
if !ok || bucket == "" || object == "" {
return "", "", false
}
return bucket, object, true
}
// note records one read. It is called from the GET path and must never block
// or allocate meaningfully: on a full queue the sample is dropped.
func (t *accessTracker) note(bucket, object string) {
if t == nil || !t.enabled.Load() {
return
}
select {
case t.ch <- accessKey(bucket, object):
default:
t.dropped.Add(1)
}
}
// noteDemoteCandidate is called by the scanner for an object it found on a
// fast pool that has gone quiet. The leader picks these up on the next merge.
func (t *accessTracker) noteDemoteCandidate(bucket, object string, pool int) {
if t == nil || !t.enabled.Load() {
return
}
t.demoteMu.Lock()
defer t.demoteMu.Unlock()
if len(t.demote) >= accessMaxDemoteCandidates {
return
}
t.demote[accessKey(bucket, object)] = demoteCandidate{Bucket: bucket, Object: object, Pool: pool}
}
// hits reports cluster-wide accesses to an object over window.
func (t *accessTracker) hits(bucket, object string, window time.Duration) uint64 {
if t == nil {
return 0
}
return t.merged.Load().hits(accessKey(bucket, object), window)
}
// lastAccess reports the cluster-wide time an object was last read. A zero
// time means "no read on record", which for demotion purposes is idle.
func (t *accessTracker) lastAccess(bucket, object string) time.Time {
if t == nil {
return time.Time{}
}
sec := t.merged.Load().lastAccess(accessKey(bucket, object))
if sec == 0 {
return time.Time{}
}
return time.Unix(sec, 0)
}
// snapshot returns the current merged view, or nil if none has been published.
func (t *accessTracker) snapshot() *mergedAccess {
if t == nil {
return nil
}
return t.merged.Load()
}
// takeDemoteCandidates drains and returns the pending candidates.
func (t *accessTracker) takeDemoteCandidates() []demoteCandidate {
t.demoteMu.Lock()
defer t.demoteMu.Unlock()
if len(t.demote) == 0 {
return nil
}
out := make([]demoteCandidate, 0, len(t.demote))
for _, c := range t.demote {
out = append(out, c)
}
t.demote = make(map[string]demoteCandidate)
return out
}
// restoreDemoteCandidates puts candidates back when the publisher is still
// busy. Dropping access samples is acceptable; dropping the only scanner
// discovery of an idle promoted object would delay demotion by a full scan.
func (t *accessTracker) restoreDemoteCandidates(candidates []demoteCandidate) {
if len(candidates) == 0 {
return
}
t.demoteMu.Lock()
defer t.demoteMu.Unlock()
for _, c := range candidates {
if len(t.demote) >= accessMaxDemoteCandidates {
return
}
t.demote[accessKey(c.Bucket, c.Object)] = c
}
}
// shardName is this node's counter object. The node name is hashed so that
// host:port never has to be escaped into an object key.
func (t *accessTracker) shardName() string {
return fmt.Sprintf("%s/%016x.bin", accessTrackerPrefix, xxh3.HashString(globalLocalNodeName))
}
// run owns the live counter map. It drains reads, and on every flush interval
// hands a marshaled shard to a background publisher.
//
// Everything here is best effort: this is accounting for a background data
// movement decision, not a durability path.
func (t *accessTracker) run(ctx context.Context, objAPI ObjectLayer) {
cfg := globalILMConfig.accessCfg()
t.enabled.Store(cfg.AccessTiering)
live := make(map[string]accessEntry)
if cfg.AccessTiering {
if buf, err := readConfig(ctx, objAPI, t.shardName()); err == nil {
var previous accessShard
if _, err = previous.UnmarshalMsg(buf); err == nil && previous.BinWidth == int64(cfg.AccessBinWidth/time.Second) {
now := time.Now().Unix()
for key, entry := range previous.Entries {
entry.rollTo(now, previous.BinWidth, cfg.AccessBins)
if entry.total() != 0 {
live[key] = entry
}
}
}
}
}
ticker := time.NewTicker(cfg.AccessFlush)
defer ticker.Stop()
// One publisher goroutine keeps object-layer I/O off the drain loop.
type pub struct {
shard accessShard
cfg ilm.Config
}
pubCh := make(chan pub, 1)
go func() {
for p := range pubCh {
t.publish(ctx, objAPI, p.shard, p.cfg)
}
}()
defer close(pubCh)
for {
select {
case <-ctx.Done():
return
case key := <-t.ch:
now := time.Now().Unix()
e := live[key]
e.rollTo(now, int64(cfg.AccessBinWidth/time.Second), cfg.AccessBins)
if len(e.Bins) > 0 {
if e.Bins[0] != ^uint32(0) {
e.Bins[0]++
}
}
e.LastAt = now
live[key] = e
case <-ticker.C:
newCfg := globalILMConfig.accessCfg()
t.enabled.Store(newCfg.AccessTiering)
if newCfg.AccessFlush != cfg.AccessFlush && newCfg.AccessFlush > 0 {
ticker.Reset(newCfg.AccessFlush)
}
cfg = newCfg
if !cfg.AccessTiering {
// Feature turned off: release the counters rather than
// holding a stale working set for the process lifetime.
clear(live)
t.merged.Store(nil)
continue
}
now := time.Now().Unix()
binWidth := int64(cfg.AccessBinWidth / time.Second)
t.evict(live, now, binWidth, cfg)
shard := accessShard{
UpdatedAt: now,
BinWidth: binWidth,
Entries: make(map[string]accessEntry, len(live)),
Demote: t.takeDemoteCandidates(),
}
for k, e := range live {
e.Bins = slices.Clone(e.Bins)
shard.Entries[k] = e
}
select {
case pubCh <- pub{shard: shard, cfg: cfg}:
default:
// Previous publish still running; skip this round
// rather than queueing work we cannot keep up with.
t.restoreDemoteCandidates(shard.Demote)
}
}
}
}
// evict rolls every counter forward and drops the ones with no hits left in
// the retained history, then enforces the tracked-object cap.
//
// ponytail: amortized O(n) sweep on the flush tick; a heap would only pay off
// past ~10M tracked keys.
func (t *accessTracker) evict(live map[string]accessEntry, now, binWidth int64, cfg ilm.Config) {
for k, e := range live {
e.rollTo(now, binWidth, cfg.AccessBins)
if e.total() == 0 {
delete(live, k)
continue
}
live[k] = e
}
if len(live) <= cfg.AccessMaxTracked {
return
}
// Over the cap: keep the hottest and drop the rest. Objects that fall
// out are cold by construction, which is exactly what the demotion path
// already assumes about anything missing from the map.
type kt struct {
key string
total uint64
}
all := make([]kt, 0, len(live))
for k, e := range live {
all = append(all, kt{k, e.total()})
}
slices.SortFunc(all, func(a, b kt) int { return cmp.Compare(b.total, a.total) })
for _, x := range all[cfg.AccessMaxTracked:] {
delete(live, x.key)
}
}
// publish writes this node's shard and rebuilds the merged snapshot from every
// node's shard. Both halves are best effort.
func (t *accessTracker) publish(ctx context.Context, objAPI ObjectLayer, shard accessShard, cfg ilm.Config) {
buf, err := shard.MarshalMsg(nil)
if err != nil {
ilmLogIf(ctx, err)
return
}
if err := saveConfig(ctx, objAPI, t.shardName(), buf); err != nil {
ilmLogIf(ctx, err)
// Still rebuild the snapshot below: our own counters are already
// in hand and a stale peer view beats no view.
}
t.merged.Store(t.mergeShards(ctx, objAPI, shard, cfg))
}
// mergeShards sums every live node's counters, including our own in-memory
// shard so this node's most recent reads are never a flush behind.
func (t *accessTracker) mergeShards(ctx context.Context, objAPI ObjectLayer, own accessShard, cfg ilm.Config) *mergedAccess {
now := time.Now().Unix()
binWidth := int64(cfg.AccessBinWidth / time.Second)
out := &mergedAccess{
entries: make(map[string]accessEntry, len(own.Entries)),
binWidth: binWidth,
at: now,
}
add := func(s accessShard) {
if s.BinWidth != binWidth {
// A peer has not yet picked up a bin-width change; merging
// its bins would silently mis-scale the counts.
return
}
for k, e := range s.Entries {
e.rollTo(now, binWidth, cfg.AccessBins)
cur, ok := out.entries[k]
if !ok {
cur = accessEntry{Bins: make([]uint32, cfg.AccessBins)}
}
cur.mergeFrom(e)
out.entries[k] = cur
}
}
add(own)
ownName := t.shardName()
stale := int64(cfg.AccessFlush/time.Second) * accessShardStaleFactor
for _, name := range t.listShards(ctx, objAPI) {
if name == ownName {
continue
}
buf, err := readConfig(ctx, objAPI, name)
if err != nil {
continue
}
var s accessShard
if _, err := s.UnmarshalMsg(buf); err != nil {
continue
}
if !accessShardFresh(now, s, stale) {
continue // node is gone or wedged
}
add(s)
}
return out
}
func (t *accessTracker) listShards(ctx context.Context, objAPI ObjectLayer) []string {
res, err := objAPI.ListObjects(ctx, minioMetaBucket, accessTrackerPrefix+"/", "", "", maxObjectList)
if err != nil {
return nil
}
names := make([]string, 0, len(res.Objects))
for _, o := range res.Objects {
if strings.HasSuffix(o.Name, ".bin") {
names = append(names, o.Name)
}
}
return names
}
// collectDemoteCandidates returns every node's pending demote candidates. Only
// the leader calls this, right before running a demotion pass.
//
// Our own shard is read back rather than skipped: run() drains the local map
// into the published shard, so anything found since the last leader pass lives
// there, not in memory. The local map is still drained here to pick up
// candidates recorded since that publish.
func (t *accessTracker) collectDemoteCandidates(ctx context.Context, objAPI ObjectLayer, cfg ilm.Config) []demoteCandidate {
seen := make(map[string]struct{})
var out []demoteCandidate
for _, c := range t.takeDemoteCandidates() {
key := accessKey(c.Bucket, c.Object)
seen[key] = struct{}{}
out = append(out, c)
}
now := time.Now().Unix()
stale := int64(cfg.AccessFlush/time.Second) * accessShardStaleFactor
for _, name := range t.listShards(ctx, objAPI) {
buf, err := readConfig(ctx, objAPI, name)
if err != nil {
continue
}
var s accessShard
if _, err := s.UnmarshalMsg(buf); err != nil {
continue
}
if !accessShardFresh(now, s, stale) {
continue
}
for _, c := range s.Demote {
key := accessKey(c.Bucket, c.Object)
if _, dup := seen[key]; dup {
continue
}
seen[key] = struct{}{}
out = append(out, c)
}
}
return out
}
-742
View File
@@ -1,742 +0,0 @@
// Code generated by github.com/tinylib/msgp DO NOT EDIT.
package cmd
import (
"github.com/tinylib/msgp/msgp"
)
// DecodeMsg implements msgp.Decodable
func (z *accessEntry) DecodeMsg(dc *msgp.Reader) (err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err)
return
}
for zb0001 > 0 {
zb0001--
field, err = dc.ReadMapKeyPtr()
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "b":
var zb0002 uint32
zb0002, err = dc.ReadArrayHeader()
if err != nil {
err = msgp.WrapError(err, "Bins")
return
}
if cap(z.Bins) >= int(zb0002) {
z.Bins = (z.Bins)[:zb0002]
} else {
z.Bins = make([]uint32, zb0002)
}
for za0001 := range z.Bins {
z.Bins[za0001], err = dc.ReadUint32()
if err != nil {
err = msgp.WrapError(err, "Bins", za0001)
return
}
}
case "h":
z.HeadAt, err = dc.ReadInt64()
if err != nil {
err = msgp.WrapError(err, "HeadAt")
return
}
case "l":
z.LastAt, err = dc.ReadInt64()
if err != nil {
err = msgp.WrapError(err, "LastAt")
return
}
default:
err = dc.Skip()
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
return
}
// EncodeMsg implements msgp.Encodable
func (z *accessEntry) EncodeMsg(en *msgp.Writer) (err error) {
// map header, size 3
// write "b"
err = en.Append(0x83, 0xa1, 0x62)
if err != nil {
return
}
err = en.WriteArrayHeader(uint32(len(z.Bins)))
if err != nil {
err = msgp.WrapError(err, "Bins")
return
}
for za0001 := range z.Bins {
err = en.WriteUint32(z.Bins[za0001])
if err != nil {
err = msgp.WrapError(err, "Bins", za0001)
return
}
}
// write "h"
err = en.Append(0xa1, 0x68)
if err != nil {
return
}
err = en.WriteInt64(z.HeadAt)
if err != nil {
err = msgp.WrapError(err, "HeadAt")
return
}
// write "l"
err = en.Append(0xa1, 0x6c)
if err != nil {
return
}
err = en.WriteInt64(z.LastAt)
if err != nil {
err = msgp.WrapError(err, "LastAt")
return
}
return
}
// MarshalMsg implements msgp.Marshaler
func (z *accessEntry) MarshalMsg(b []byte) (o []byte, err error) {
o = msgp.Require(b, z.Msgsize())
// map header, size 3
// string "b"
o = append(o, 0x83, 0xa1, 0x62)
o = msgp.AppendArrayHeader(o, uint32(len(z.Bins)))
for za0001 := range z.Bins {
o = msgp.AppendUint32(o, z.Bins[za0001])
}
// string "h"
o = append(o, 0xa1, 0x68)
o = msgp.AppendInt64(o, z.HeadAt)
// string "l"
o = append(o, 0xa1, 0x6c)
o = msgp.AppendInt64(o, z.LastAt)
return
}
// UnmarshalMsg implements msgp.Unmarshaler
func (z *accessEntry) UnmarshalMsg(bts []byte) (o []byte, err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
for zb0001 > 0 {
zb0001--
field, bts, err = msgp.ReadMapKeyZC(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "b":
var zb0002 uint32
zb0002, bts, err = msgp.ReadArrayHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Bins")
return
}
if cap(z.Bins) >= int(zb0002) {
z.Bins = (z.Bins)[:zb0002]
} else {
z.Bins = make([]uint32, zb0002)
}
for za0001 := range z.Bins {
z.Bins[za0001], bts, err = msgp.ReadUint32Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "Bins", za0001)
return
}
}
case "h":
z.HeadAt, bts, err = msgp.ReadInt64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "HeadAt")
return
}
case "l":
z.LastAt, bts, err = msgp.ReadInt64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "LastAt")
return
}
default:
bts, err = msgp.Skip(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
o = bts
return
}
// Msgsize returns an upper bound estimate of the number of bytes occupied by the serialized message
func (z *accessEntry) Msgsize() (s int) {
s = 1 + 2 + msgp.ArrayHeaderSize + (len(z.Bins) * (msgp.Uint32Size)) + 2 + msgp.Int64Size + 2 + msgp.Int64Size
return
}
// DecodeMsg implements msgp.Decodable
func (z *accessShard) DecodeMsg(dc *msgp.Reader) (err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err)
return
}
for zb0001 > 0 {
zb0001--
field, err = dc.ReadMapKeyPtr()
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "u":
z.UpdatedAt, err = dc.ReadInt64()
if err != nil {
err = msgp.WrapError(err, "UpdatedAt")
return
}
case "bw":
z.BinWidth, err = dc.ReadInt64()
if err != nil {
err = msgp.WrapError(err, "BinWidth")
return
}
case "e":
var zb0002 uint32
zb0002, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err, "Entries")
return
}
if z.Entries == nil {
z.Entries = make(map[string]accessEntry, zb0002)
} else if len(z.Entries) > 0 {
clear(z.Entries)
}
for zb0002 > 0 {
zb0002--
var za0001 string
za0001, err = dc.ReadString()
if err != nil {
err = msgp.WrapError(err, "Entries")
return
}
var za0002 accessEntry
err = za0002.DecodeMsg(dc)
if err != nil {
err = msgp.WrapError(err, "Entries", za0001)
return
}
z.Entries[za0001] = za0002
}
case "d":
var zb0003 uint32
zb0003, err = dc.ReadArrayHeader()
if err != nil {
err = msgp.WrapError(err, "Demote")
return
}
if cap(z.Demote) >= int(zb0003) {
z.Demote = (z.Demote)[:zb0003]
} else {
z.Demote = make([]demoteCandidate, zb0003)
}
for za0003 := range z.Demote {
var zb0004 uint32
zb0004, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err, "Demote", za0003)
return
}
for zb0004 > 0 {
zb0004--
field, err = dc.ReadMapKeyPtr()
if err != nil {
err = msgp.WrapError(err, "Demote", za0003)
return
}
switch msgp.UnsafeString(field) {
case "b":
z.Demote[za0003].Bucket, err = dc.ReadString()
if err != nil {
err = msgp.WrapError(err, "Demote", za0003, "Bucket")
return
}
case "o":
z.Demote[za0003].Object, err = dc.ReadString()
if err != nil {
err = msgp.WrapError(err, "Demote", za0003, "Object")
return
}
case "p":
z.Demote[za0003].Pool, err = dc.ReadInt()
if err != nil {
err = msgp.WrapError(err, "Demote", za0003, "Pool")
return
}
default:
err = dc.Skip()
if err != nil {
err = msgp.WrapError(err, "Demote", za0003)
return
}
}
}
}
default:
err = dc.Skip()
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
return
}
// EncodeMsg implements msgp.Encodable
func (z *accessShard) EncodeMsg(en *msgp.Writer) (err error) {
// map header, size 4
// write "u"
err = en.Append(0x84, 0xa1, 0x75)
if err != nil {
return
}
err = en.WriteInt64(z.UpdatedAt)
if err != nil {
err = msgp.WrapError(err, "UpdatedAt")
return
}
// write "bw"
err = en.Append(0xa2, 0x62, 0x77)
if err != nil {
return
}
err = en.WriteInt64(z.BinWidth)
if err != nil {
err = msgp.WrapError(err, "BinWidth")
return
}
// write "e"
err = en.Append(0xa1, 0x65)
if err != nil {
return
}
err = en.WriteMapHeader(uint32(len(z.Entries)))
if err != nil {
err = msgp.WrapError(err, "Entries")
return
}
for za0001, za0002 := range z.Entries {
err = en.WriteString(za0001)
if err != nil {
err = msgp.WrapError(err, "Entries")
return
}
err = za0002.EncodeMsg(en)
if err != nil {
err = msgp.WrapError(err, "Entries", za0001)
return
}
}
// write "d"
err = en.Append(0xa1, 0x64)
if err != nil {
return
}
err = en.WriteArrayHeader(uint32(len(z.Demote)))
if err != nil {
err = msgp.WrapError(err, "Demote")
return
}
for za0003 := range z.Demote {
// map header, size 3
// write "b"
err = en.Append(0x83, 0xa1, 0x62)
if err != nil {
return
}
err = en.WriteString(z.Demote[za0003].Bucket)
if err != nil {
err = msgp.WrapError(err, "Demote", za0003, "Bucket")
return
}
// write "o"
err = en.Append(0xa1, 0x6f)
if err != nil {
return
}
err = en.WriteString(z.Demote[za0003].Object)
if err != nil {
err = msgp.WrapError(err, "Demote", za0003, "Object")
return
}
// write "p"
err = en.Append(0xa1, 0x70)
if err != nil {
return
}
err = en.WriteInt(z.Demote[za0003].Pool)
if err != nil {
err = msgp.WrapError(err, "Demote", za0003, "Pool")
return
}
}
return
}
// MarshalMsg implements msgp.Marshaler
func (z *accessShard) MarshalMsg(b []byte) (o []byte, err error) {
o = msgp.Require(b, z.Msgsize())
// map header, size 4
// string "u"
o = append(o, 0x84, 0xa1, 0x75)
o = msgp.AppendInt64(o, z.UpdatedAt)
// string "bw"
o = append(o, 0xa2, 0x62, 0x77)
o = msgp.AppendInt64(o, z.BinWidth)
// string "e"
o = append(o, 0xa1, 0x65)
o = msgp.AppendMapHeader(o, uint32(len(z.Entries)))
for za0001, za0002 := range z.Entries {
o = msgp.AppendString(o, za0001)
o, err = za0002.MarshalMsg(o)
if err != nil {
err = msgp.WrapError(err, "Entries", za0001)
return
}
}
// string "d"
o = append(o, 0xa1, 0x64)
o = msgp.AppendArrayHeader(o, uint32(len(z.Demote)))
for za0003 := range z.Demote {
// map header, size 3
// string "b"
o = append(o, 0x83, 0xa1, 0x62)
o = msgp.AppendString(o, z.Demote[za0003].Bucket)
// string "o"
o = append(o, 0xa1, 0x6f)
o = msgp.AppendString(o, z.Demote[za0003].Object)
// string "p"
o = append(o, 0xa1, 0x70)
o = msgp.AppendInt(o, z.Demote[za0003].Pool)
}
return
}
// UnmarshalMsg implements msgp.Unmarshaler
func (z *accessShard) UnmarshalMsg(bts []byte) (o []byte, err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
for zb0001 > 0 {
zb0001--
field, bts, err = msgp.ReadMapKeyZC(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "u":
z.UpdatedAt, bts, err = msgp.ReadInt64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "UpdatedAt")
return
}
case "bw":
z.BinWidth, bts, err = msgp.ReadInt64Bytes(bts)
if err != nil {
err = msgp.WrapError(err, "BinWidth")
return
}
case "e":
var zb0002 uint32
zb0002, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Entries")
return
}
if z.Entries == nil {
z.Entries = make(map[string]accessEntry, zb0002)
} else if len(z.Entries) > 0 {
clear(z.Entries)
}
for zb0002 > 0 {
var za0002 accessEntry
zb0002--
var za0001 string
za0001, bts, err = msgp.ReadStringBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Entries")
return
}
bts, err = za0002.UnmarshalMsg(bts)
if err != nil {
err = msgp.WrapError(err, "Entries", za0001)
return
}
z.Entries[za0001] = za0002
}
case "d":
var zb0003 uint32
zb0003, bts, err = msgp.ReadArrayHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Demote")
return
}
if cap(z.Demote) >= int(zb0003) {
z.Demote = (z.Demote)[:zb0003]
} else {
z.Demote = make([]demoteCandidate, zb0003)
}
for za0003 := range z.Demote {
var zb0004 uint32
zb0004, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Demote", za0003)
return
}
for zb0004 > 0 {
zb0004--
field, bts, err = msgp.ReadMapKeyZC(bts)
if err != nil {
err = msgp.WrapError(err, "Demote", za0003)
return
}
switch msgp.UnsafeString(field) {
case "b":
z.Demote[za0003].Bucket, bts, err = msgp.ReadStringBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Demote", za0003, "Bucket")
return
}
case "o":
z.Demote[za0003].Object, bts, err = msgp.ReadStringBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Demote", za0003, "Object")
return
}
case "p":
z.Demote[za0003].Pool, bts, err = msgp.ReadIntBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Demote", za0003, "Pool")
return
}
default:
bts, err = msgp.Skip(bts)
if err != nil {
err = msgp.WrapError(err, "Demote", za0003)
return
}
}
}
}
default:
bts, err = msgp.Skip(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
o = bts
return
}
// Msgsize returns an upper bound estimate of the number of bytes occupied by the serialized message
func (z *accessShard) Msgsize() (s int) {
s = 1 + 2 + msgp.Int64Size + 3 + msgp.Int64Size + 2 + msgp.MapHeaderSize
if z.Entries != nil {
for za0001, za0002 := range z.Entries {
_ = za0002
s += msgp.StringPrefixSize + len(za0001) + za0002.Msgsize()
}
}
s += 2 + msgp.ArrayHeaderSize
for za0003 := range z.Demote {
s += 1 + 2 + msgp.StringPrefixSize + len(z.Demote[za0003].Bucket) + 2 + msgp.StringPrefixSize + len(z.Demote[za0003].Object) + 2 + msgp.IntSize
}
return
}
// DecodeMsg implements msgp.Decodable
func (z *demoteCandidate) DecodeMsg(dc *msgp.Reader) (err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, err = dc.ReadMapHeader()
if err != nil {
err = msgp.WrapError(err)
return
}
for zb0001 > 0 {
zb0001--
field, err = dc.ReadMapKeyPtr()
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "b":
z.Bucket, err = dc.ReadString()
if err != nil {
err = msgp.WrapError(err, "Bucket")
return
}
case "o":
z.Object, err = dc.ReadString()
if err != nil {
err = msgp.WrapError(err, "Object")
return
}
case "p":
z.Pool, err = dc.ReadInt()
if err != nil {
err = msgp.WrapError(err, "Pool")
return
}
default:
err = dc.Skip()
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
return
}
// EncodeMsg implements msgp.Encodable
func (z demoteCandidate) EncodeMsg(en *msgp.Writer) (err error) {
// map header, size 3
// write "b"
err = en.Append(0x83, 0xa1, 0x62)
if err != nil {
return
}
err = en.WriteString(z.Bucket)
if err != nil {
err = msgp.WrapError(err, "Bucket")
return
}
// write "o"
err = en.Append(0xa1, 0x6f)
if err != nil {
return
}
err = en.WriteString(z.Object)
if err != nil {
err = msgp.WrapError(err, "Object")
return
}
// write "p"
err = en.Append(0xa1, 0x70)
if err != nil {
return
}
err = en.WriteInt(z.Pool)
if err != nil {
err = msgp.WrapError(err, "Pool")
return
}
return
}
// MarshalMsg implements msgp.Marshaler
func (z demoteCandidate) MarshalMsg(b []byte) (o []byte, err error) {
o = msgp.Require(b, z.Msgsize())
// map header, size 3
// string "b"
o = append(o, 0x83, 0xa1, 0x62)
o = msgp.AppendString(o, z.Bucket)
// string "o"
o = append(o, 0xa1, 0x6f)
o = msgp.AppendString(o, z.Object)
// string "p"
o = append(o, 0xa1, 0x70)
o = msgp.AppendInt(o, z.Pool)
return
}
// UnmarshalMsg implements msgp.Unmarshaler
func (z *demoteCandidate) UnmarshalMsg(bts []byte) (o []byte, err error) {
var field []byte
_ = field
var zb0001 uint32
zb0001, bts, err = msgp.ReadMapHeaderBytes(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
for zb0001 > 0 {
zb0001--
field, bts, err = msgp.ReadMapKeyZC(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
switch msgp.UnsafeString(field) {
case "b":
z.Bucket, bts, err = msgp.ReadStringBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Bucket")
return
}
case "o":
z.Object, bts, err = msgp.ReadStringBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Object")
return
}
case "p":
z.Pool, bts, err = msgp.ReadIntBytes(bts)
if err != nil {
err = msgp.WrapError(err, "Pool")
return
}
default:
bts, err = msgp.Skip(bts)
if err != nil {
err = msgp.WrapError(err)
return
}
}
}
o = bts
return
}
// Msgsize returns an upper bound estimate of the number of bytes occupied by the serialized message
func (z demoteCandidate) Msgsize() (s int) {
s = 1 + 2 + msgp.StringPrefixSize + len(z.Bucket) + 2 + msgp.StringPrefixSize + len(z.Object) + 2 + msgp.IntSize
return
}
-349
View File
@@ -1,349 +0,0 @@
// Code generated by github.com/tinylib/msgp DO NOT EDIT.
package cmd
import (
"bytes"
"testing"
"github.com/tinylib/msgp/msgp"
)
func TestMarshalUnmarshalaccessEntry(t *testing.T) {
v := accessEntry{}
bts, err := v.MarshalMsg(nil)
if err != nil {
t.Fatal(err)
}
left, err := v.UnmarshalMsg(bts)
if err != nil {
t.Fatal(err)
}
if len(left) > 0 {
t.Errorf("%d bytes left over after UnmarshalMsg(): %q", len(left), left)
}
left, err = msgp.Skip(bts)
if err != nil {
t.Fatal(err)
}
if len(left) > 0 {
t.Errorf("%d bytes left over after Skip(): %q", len(left), left)
}
}
func BenchmarkMarshalMsgaccessEntry(b *testing.B) {
v := accessEntry{}
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
v.MarshalMsg(nil)
}
}
func BenchmarkAppendMsgaccessEntry(b *testing.B) {
v := accessEntry{}
bts := make([]byte, 0, v.Msgsize())
bts, _ = v.MarshalMsg(bts[0:0])
b.SetBytes(int64(len(bts)))
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
bts, _ = v.MarshalMsg(bts[0:0])
}
}
func BenchmarkUnmarshalaccessEntry(b *testing.B) {
v := accessEntry{}
bts, _ := v.MarshalMsg(nil)
b.ReportAllocs()
b.SetBytes(int64(len(bts)))
b.ResetTimer()
for i := 0; i < b.N; i++ {
_, err := v.UnmarshalMsg(bts)
if err != nil {
b.Fatal(err)
}
}
}
func TestEncodeDecodeaccessEntry(t *testing.T) {
v := accessEntry{}
var buf bytes.Buffer
msgp.Encode(&buf, &v)
m := v.Msgsize()
if buf.Len() > m {
t.Log("WARNING: TestEncodeDecodeaccessEntry Msgsize() is inaccurate")
}
vn := accessEntry{}
err := msgp.Decode(&buf, &vn)
if err != nil {
t.Error(err)
}
buf.Reset()
msgp.Encode(&buf, &v)
err = msgp.NewReader(&buf).Skip()
if err != nil {
t.Error(err)
}
}
func BenchmarkEncodeaccessEntry(b *testing.B) {
v := accessEntry{}
var buf bytes.Buffer
msgp.Encode(&buf, &v)
b.SetBytes(int64(buf.Len()))
en := msgp.NewWriter(msgp.Nowhere)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
v.EncodeMsg(en)
}
en.Flush()
}
func BenchmarkDecodeaccessEntry(b *testing.B) {
v := accessEntry{}
var buf bytes.Buffer
msgp.Encode(&buf, &v)
b.SetBytes(int64(buf.Len()))
rd := msgp.NewEndlessReader(buf.Bytes(), b)
dc := msgp.NewReader(rd)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
err := v.DecodeMsg(dc)
if err != nil {
b.Fatal(err)
}
}
}
func TestMarshalUnmarshalaccessShard(t *testing.T) {
v := accessShard{}
bts, err := v.MarshalMsg(nil)
if err != nil {
t.Fatal(err)
}
left, err := v.UnmarshalMsg(bts)
if err != nil {
t.Fatal(err)
}
if len(left) > 0 {
t.Errorf("%d bytes left over after UnmarshalMsg(): %q", len(left), left)
}
left, err = msgp.Skip(bts)
if err != nil {
t.Fatal(err)
}
if len(left) > 0 {
t.Errorf("%d bytes left over after Skip(): %q", len(left), left)
}
}
func BenchmarkMarshalMsgaccessShard(b *testing.B) {
v := accessShard{}
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
v.MarshalMsg(nil)
}
}
func BenchmarkAppendMsgaccessShard(b *testing.B) {
v := accessShard{}
bts := make([]byte, 0, v.Msgsize())
bts, _ = v.MarshalMsg(bts[0:0])
b.SetBytes(int64(len(bts)))
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
bts, _ = v.MarshalMsg(bts[0:0])
}
}
func BenchmarkUnmarshalaccessShard(b *testing.B) {
v := accessShard{}
bts, _ := v.MarshalMsg(nil)
b.ReportAllocs()
b.SetBytes(int64(len(bts)))
b.ResetTimer()
for i := 0; i < b.N; i++ {
_, err := v.UnmarshalMsg(bts)
if err != nil {
b.Fatal(err)
}
}
}
func TestEncodeDecodeaccessShard(t *testing.T) {
v := accessShard{}
var buf bytes.Buffer
msgp.Encode(&buf, &v)
m := v.Msgsize()
if buf.Len() > m {
t.Log("WARNING: TestEncodeDecodeaccessShard Msgsize() is inaccurate")
}
vn := accessShard{}
err := msgp.Decode(&buf, &vn)
if err != nil {
t.Error(err)
}
buf.Reset()
msgp.Encode(&buf, &v)
err = msgp.NewReader(&buf).Skip()
if err != nil {
t.Error(err)
}
}
func BenchmarkEncodeaccessShard(b *testing.B) {
v := accessShard{}
var buf bytes.Buffer
msgp.Encode(&buf, &v)
b.SetBytes(int64(buf.Len()))
en := msgp.NewWriter(msgp.Nowhere)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
v.EncodeMsg(en)
}
en.Flush()
}
func BenchmarkDecodeaccessShard(b *testing.B) {
v := accessShard{}
var buf bytes.Buffer
msgp.Encode(&buf, &v)
b.SetBytes(int64(buf.Len()))
rd := msgp.NewEndlessReader(buf.Bytes(), b)
dc := msgp.NewReader(rd)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
err := v.DecodeMsg(dc)
if err != nil {
b.Fatal(err)
}
}
}
func TestMarshalUnmarshaldemoteCandidate(t *testing.T) {
v := demoteCandidate{}
bts, err := v.MarshalMsg(nil)
if err != nil {
t.Fatal(err)
}
left, err := v.UnmarshalMsg(bts)
if err != nil {
t.Fatal(err)
}
if len(left) > 0 {
t.Errorf("%d bytes left over after UnmarshalMsg(): %q", len(left), left)
}
left, err = msgp.Skip(bts)
if err != nil {
t.Fatal(err)
}
if len(left) > 0 {
t.Errorf("%d bytes left over after Skip(): %q", len(left), left)
}
}
func BenchmarkMarshalMsgdemoteCandidate(b *testing.B) {
v := demoteCandidate{}
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
v.MarshalMsg(nil)
}
}
func BenchmarkAppendMsgdemoteCandidate(b *testing.B) {
v := demoteCandidate{}
bts := make([]byte, 0, v.Msgsize())
bts, _ = v.MarshalMsg(bts[0:0])
b.SetBytes(int64(len(bts)))
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
bts, _ = v.MarshalMsg(bts[0:0])
}
}
func BenchmarkUnmarshaldemoteCandidate(b *testing.B) {
v := demoteCandidate{}
bts, _ := v.MarshalMsg(nil)
b.ReportAllocs()
b.SetBytes(int64(len(bts)))
b.ResetTimer()
for i := 0; i < b.N; i++ {
_, err := v.UnmarshalMsg(bts)
if err != nil {
b.Fatal(err)
}
}
}
func TestEncodeDecodedemoteCandidate(t *testing.T) {
v := demoteCandidate{}
var buf bytes.Buffer
msgp.Encode(&buf, &v)
m := v.Msgsize()
if buf.Len() > m {
t.Log("WARNING: TestEncodeDecodedemoteCandidate Msgsize() is inaccurate")
}
vn := demoteCandidate{}
err := msgp.Decode(&buf, &vn)
if err != nil {
t.Error(err)
}
buf.Reset()
msgp.Encode(&buf, &v)
err = msgp.NewReader(&buf).Skip()
if err != nil {
t.Error(err)
}
}
func BenchmarkEncodedemoteCandidate(b *testing.B) {
v := demoteCandidate{}
var buf bytes.Buffer
msgp.Encode(&buf, &v)
b.SetBytes(int64(buf.Len()))
en := msgp.NewWriter(msgp.Nowhere)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
v.EncodeMsg(en)
}
en.Flush()
}
func BenchmarkDecodedemoteCandidate(b *testing.B) {
v := demoteCandidate{}
var buf bytes.Buffer
msgp.Encode(&buf, &v)
b.SetBytes(int64(buf.Len()))
rd := msgp.NewEndlessReader(buf.Bytes(), b)
dc := msgp.NewReader(rd)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
err := v.DecodeMsg(dc)
if err != nil {
b.Fatal(err)
}
}
}
-95
View File
@@ -1,95 +0,0 @@
// Copyright (c) 2015-2026 MinIO, Inc.
package cmd
import (
"math"
"testing"
"time"
"github.com/minio/minio/internal/config/ilm"
)
func TestAccessEntryRollAndHits(t *testing.T) {
entry := accessEntry{Bins: []uint32{3, 2, 1}, HeadAt: 100, LastAt: 107}
entry.rollTo(120, 10, 3)
want := []uint32{0, 0, 3}
for i := range want {
if entry.Bins[i] != want[i] {
t.Fatalf("bins = %v, want %v", entry.Bins, want)
}
}
if entry.HeadAt != 120 || entry.LastAt != 107 {
t.Fatalf("head/last = %d/%d", entry.HeadAt, entry.LastAt)
}
entry = accessEntry{Bins: []uint32{10, 20, 30}, HeadAt: 120}
if got := entry.hits(20*time.Second, 10); got != 30 {
t.Fatalf("20s hits = %d, want 30", got)
}
if got := entry.hits(21*time.Second, 10); got != 60 {
t.Fatalf("21s hits = %d, want 60", got)
}
entry.rollTo(200, 10, 3)
if got := entry.total(); got != 0 {
t.Fatalf("expired total = %d, want 0", got)
}
}
func TestAccessEntryMergeSaturates(t *testing.T) {
entry := accessEntry{Bins: []uint32{math.MaxUint32 - 1}}
entry.mergeFrom(accessEntry{Bins: []uint32{10}, LastAt: 50})
if entry.Bins[0] != math.MaxUint32 || entry.LastAt != 50 {
t.Fatalf("merged entry = %+v", entry)
}
}
func TestAccessTrackerEvictsColdest(t *testing.T) {
tracker := newAccessTracker()
live := map[string]accessEntry{
"a": {Bins: []uint32{1}, HeadAt: 100},
"b": {Bins: []uint32{5}, HeadAt: 100},
"c": {Bins: []uint32{3}, HeadAt: 100},
}
tracker.evict(live, 100, 10, ilm.Config{AccessBins: 1, AccessMaxTracked: 2})
if len(live) != 2 {
t.Fatalf("len = %d, want 2", len(live))
}
if _, ok := live["a"]; ok {
t.Fatal("coldest entry was retained")
}
}
func TestAccessKeyRoundTrip(t *testing.T) {
key := accessKey("bucket", "a/b/c")
bucket, object, ok := splitAccessKey(key)
if !ok || bucket != "bucket" || object != "a/b/c" {
t.Fatalf("split %q = %q/%q/%v", key, bucket, object, ok)
}
}
func TestAccessShardFresh(t *testing.T) {
const now = int64(1000)
if !accessShardFresh(now, accessShard{UpdatedAt: 950}, 100) {
t.Fatal("fresh shard rejected")
}
if accessShardFresh(now, accessShard{UpdatedAt: 899}, 100) {
t.Fatal("stale shard accepted")
}
if accessShardFresh(now, accessShard{UpdatedAt: 1101}, 100) {
t.Fatal("far-future shard accepted")
}
if accessShardFresh(now, accessShard{}, 100) {
t.Fatal("zero timestamp accepted")
}
}
func TestAccessTrackerRestoresDemoteCandidates(t *testing.T) {
tracker := newAccessTracker()
candidate := demoteCandidate{Bucket: "bucket", Object: "object", Pool: 1}
tracker.restoreDemoteCandidates([]demoteCandidate{candidate})
got := tracker.takeDemoteCandidates()
if len(got) != 1 || got[0] != candidate {
t.Fatalf("restored candidates = %+v, want %+v", got, candidate)
}
}
-30
View File
@@ -19,7 +19,6 @@ package cmd
import (
"sync"
"time"
"github.com/minio/minio/internal/config/ilm"
)
@@ -28,16 +27,6 @@ var globalILMConfig = ilmConfig{
cfg: ilm.Config{
ExpirationWorkers: 100,
TransitionWorkers: 100,
// Access tiering stays off until configured, but the counter
// geometry must be sane from the start: the tracker divides by
// AccessBinWidth before any config is loaded.
AccessPromoteWatermark: 85,
AccessBinWidth: time.Minute,
AccessBins: 12,
AccessFlush: time.Minute,
AccessMinResidency: 24 * time.Hour,
AccessWorkers: 10,
AccessMaxTracked: 1000000,
},
}
@@ -60,25 +49,6 @@ func (c *ilmConfig) getTransitionWorkers() int {
return c.cfg.TransitionWorkers
}
// accessCfg returns a copy of the access tiering settings. Callers take the
// whole struct rather than one getter per field because the promotion and
// demotion paths need a consistent view of several knobs at once.
func (c *ilmConfig) accessCfg() ilm.Config {
c.mu.RLock()
defer c.mu.RUnlock()
return c.cfg
}
// accessTieringEnabled is the cheap gate used on the scanner path.
func (c *ilmConfig) accessTieringEnabled() bool {
c.mu.RLock()
defer c.mu.RUnlock()
_, ok := c.cfg.HotPool()
return ok
}
func (c *ilmConfig) update(cfg ilm.Config) {
c.mu.Lock()
defer c.mu.Unlock()
+2 -3
View File
@@ -19,12 +19,11 @@ func _() {
_ = x[lcEventSrc_s3PutObject-8]
_ = x[lcEventSrc_s3CopyObject-9]
_ = x[lcEventSrc_s3CompleteMultipartUpload-10]
_ = x[lcEventSrc_AccessTier-11]
}
const _lcEventSrc_name = "NoneHealScannerDecomRebals3HeadObjects3GetObjects3ListObjectss3PutObjects3CopyObjects3CompleteMultipartUploadAccessTier"
const _lcEventSrc_name = "NoneHealScannerDecomRebals3HeadObjects3GetObjects3ListObjectss3PutObjects3CopyObjects3CompleteMultipartUpload"
var _lcEventSrc_index = [...]uint8{0, 4, 8, 15, 20, 25, 37, 48, 61, 72, 84, 109, 119}
var _lcEventSrc_index = [...]uint8{0, 4, 8, 15, 20, 25, 37, 48, 61, 72, 84, 109}
func (i lcEventSrc) String() string {
idx := int(i) - 0
-37
View File
@@ -26,17 +26,6 @@ const (
transitionActiveTasks = "transition_active_tasks"
transitionPendingTasks = "transition_pending_tasks"
transitionMissedImmediateTasks = "transition_missed_immediate_tasks"
accessTierActiveTasks = "access_tier_active_tasks"
accessTierPendingTasks = "access_tier_pending_tasks"
accessTierPromotionsTotal = "access_tier_promotions_total"
accessTierDemotionsTotal = "access_tier_demotions_total"
accessTierBytesMovedTotal = "access_tier_bytes_moved_total"
accessTierFailuresTotal = "access_tier_failures_total"
accessTierSkippedWatermark = "access_tier_skipped_watermark_total"
accessTierSkippedMaxSize = "access_tier_skipped_max_size_total"
accessTierSkippedQuota = "access_tier_skipped_bucket_quota_total"
accessTierHotBytes = "access_tier_hot_bytes"
accessTierSamplesDropped = "access_tier_samples_dropped_total"
versionsScanned = "versions_scanned"
)
@@ -45,17 +34,6 @@ var (
ilmTransitionActiveTasksMD = NewGaugeMD(transitionActiveTasks, "Number of active ILM transition tasks")
ilmTransitionPendingTasksMD = NewGaugeMD(transitionPendingTasks, "Number of pending ILM transition tasks in the queue")
ilmTransitionMissedImmediateTasksMD = NewCounterMD(transitionMissedImmediateTasks, "Number of missed immediate ILM transition tasks")
ilmAccessTierActiveTasksMD = NewGaugeMD(accessTierActiveTasks, "Number of active access-tier pool moves")
ilmAccessTierPendingTasksMD = NewGaugeMD(accessTierPendingTasks, "Number of pending access-tier pool moves")
ilmAccessTierPromotionsTotalMD = NewCounterMD(accessTierPromotionsTotal, "Total objects promoted by access-tier ILM")
ilmAccessTierDemotionsTotalMD = NewCounterMD(accessTierDemotionsTotal, "Total objects demoted by access-tier ILM")
ilmAccessTierBytesMovedTotalMD = NewCounterMD(accessTierBytesMovedTotal, "Total logical bytes moved by access-tier ILM")
ilmAccessTierFailuresTotalMD = NewCounterMD(accessTierFailuresTotal, "Total failed access-tier ILM moves")
ilmAccessTierSkippedWatermarkMD = NewCounterMD(accessTierSkippedWatermark, "Promotions skipped because the hot pool reached its watermark")
ilmAccessTierSkippedMaxSizeMD = NewCounterMD(accessTierSkippedMaxSize, "Promotions skipped because the cluster hot-tier size cap was reached")
ilmAccessTierSkippedQuotaMD = NewCounterMD(accessTierSkippedQuota, "Promotions skipped because the bucket hot-tier quota was reached")
ilmAccessTierHotBytesMD = NewGaugeMD(accessTierHotBytes, "Logical bytes currently accounted to the hot tier", "bucket")
ilmAccessTierSamplesDroppedMD = NewCounterMD(accessTierSamplesDropped, "GET samples dropped because the access tracker queue was full")
ilmVersionsScannedMD = NewCounterMD(versionsScanned, "Total number of object versions checked for ILM actions since server start")
)
@@ -69,21 +47,6 @@ func loadILMMetrics(_ context.Context, m MetricValues, _ *metricsCache) error {
m.Set(transitionPendingTasks, float64(globalTransitionState.PendingTasks()))
m.Set(transitionMissedImmediateTasks, float64(globalTransitionState.MissedImmediateTasks()))
}
if globalAccessTierState != nil {
m.Set(accessTierActiveTasks, float64(globalAccessTierState.ActiveTasks()))
m.Set(accessTierPendingTasks, float64(globalAccessTierState.PendingTasks()))
m.Set(accessTierPromotionsTotal, float64(globalAccessTierState.promotions.Load()))
m.Set(accessTierDemotionsTotal, float64(globalAccessTierState.demotions.Load()))
m.Set(accessTierBytesMovedTotal, float64(globalAccessTierState.bytesMoved.Load()))
m.Set(accessTierFailuresTotal, float64(globalAccessTierState.failures.Load()))
m.Set(accessTierSkippedWatermark, float64(globalAccessTierState.skippedWatermark.Load()))
m.Set(accessTierSkippedMaxSize, float64(globalAccessTierState.skippedMaxSize.Load()))
m.Set(accessTierSkippedQuota, float64(globalAccessTierState.skippedQuota.Load()))
for bucket, bytes := range globalAccessTierState.hotUsageSnapshot() {
m.Set(accessTierHotBytes, float64(bytes), "bucket", bucket)
}
}
m.Set(accessTierSamplesDropped, float64(globalAccessTracker.dropped.Load()))
m.Set(versionsScanned, float64(globalScannerMetrics.lifetime(scannerMetricILM)))
return nil
-11
View File
@@ -392,17 +392,6 @@ func newMetricGroups(r *prometheus.Registry) *metricsV3Collection {
ilmTransitionActiveTasksMD,
ilmTransitionPendingTasksMD,
ilmTransitionMissedImmediateTasksMD,
ilmAccessTierActiveTasksMD,
ilmAccessTierPendingTasksMD,
ilmAccessTierPromotionsTotalMD,
ilmAccessTierDemotionsTotalMD,
ilmAccessTierBytesMovedTotalMD,
ilmAccessTierFailuresTotalMD,
ilmAccessTierSkippedWatermarkMD,
ilmAccessTierSkippedMaxSizeMD,
ilmAccessTierSkippedQuotaMD,
ilmAccessTierHotBytesMD,
ilmAccessTierSamplesDroppedMD,
ilmVersionsScannedMD,
},
loadILMMetrics,
-4
View File
@@ -122,10 +122,6 @@ type ObjectOptions struct {
SkipRebalancing bool
SrcPoolIdx int // set by PutObject/CompleteMultipart operations due to rebalance; used to prevent rebalance src, dst pools to be the same
// DstPoolIdx forces a data-movement write onto a specific server pool.
// It is ignored unless DataMovement is true; a pointer keeps pool zero
// distinguishable from the unset value.
DstPoolIdx *int
DataMovement bool // indicates an going decommisionning or rebalacing
-2
View File
@@ -576,8 +576,6 @@ func (api objectAPIHandlers) getObjectHandler(ctx context.Context, objectAPI Obj
return
}
globalAccessTracker.note(bucket, object)
// Notify object accessed via a GET request.
sendEvent(eventArgs{
EventName: event.ObjectAccessedGet,
-5
View File
@@ -501,7 +501,6 @@ func initAllSubsystems(ctx context.Context) {
globalTierConfigMgr = NewTierConfigMgr()
globalTransitionState = newTransitionState(GlobalContext)
globalAccessTierState = newAccessTierState(GlobalContext)
globalSiteResyncMetrics = newSiteResyncMetrics(GlobalContext)
}
@@ -1070,10 +1069,6 @@ func serverMain(ctx *cli.Context) {
bootstrapTrace("globalTransitionState.Init", func() {
globalTransitionState.Init(newObject)
})
bootstrapTrace("globalAccessTierState.Init", func() {
globalAccessTierState.Init(newObject)
go globalAccessTracker.run(GlobalContext, newObject)
})
go func() {
// Initialize transition tier configuration manager
+27
View File
@@ -0,0 +1,27 @@
# Data-usage v9 retirement fixture
Generated using unmodified `scanDataFolder` and `dataUsageCache.serializeTo` from
Server commit `89637554d60c27cfc51d2281d0a4fe15e415f06d` (Go 1.27.1, darwin/arm64).
The scanner visits three real local files; its size callback supplies synthetic
version/delete-marker/remote-tier summaries, following `TestDataUsageCacheSerialize`.
It is a scanner/cache compatibility fixture, not a distributed object-store test.
The old scanner counts 74,962 bytes, 3 objects, 5 versions and 2 delete markers.
Its nonzero retired `hts` totals 9,426 bytes. Remote tier `COLD` contains 65,536
bytes, one version and one object. The cache includes nested children, both
histograms, a fixed timestamp and scanner cycle 42. The JSON is the old scanner's
complete expected cache (including `HotTierSize`, which the new reader ignores).
Binary SHA-256: `4c9c7e94cea7758fe9b498b19f639188764b2787582783e6cc950f2c44d52a8b`.
To regenerate, create a detached worktree at that exact source commit, copy
`generate.go.txt` to `cmd/retirement-fixture_test.go`, and run:
```sh
SILO_RETIRE_FIXTURE_DIR=/absolute/output/directory go test ./cmd -run '^TestGenerateAccessRetirementV9Fixture$' -count=1 -v
```
The historical production writer adds the version byte, compresses with zstd
and encodes msgp, including the nonzero `hts` field. Do not regenerate using the
retired implementation or by changing the header of a v8 payload. Map order may
change serialized bytes across regeneration; compare the decoded full cache.
Binary file not shown.
+118
View File
@@ -0,0 +1,118 @@
{
"Info": {
"Name": "/",
"NextCycle": 42,
"LastUpdate": "2026-09-15T00:00:00Z",
"SkipHealing": true
},
"Cache": {
"/": {
"Children": {
"v9-bucket": {}
},
"Size": 0,
"HotTierSize": 0,
"Objects": 0,
"Versions": 0,
"DeleteMarkers": 0,
"ObjSizes": [
0,
0,
0,
0,
0,
0,
0,
0,
0,
0,
0
],
"ObjVersions": [
0,
0,
0,
0,
0,
0,
0
],
"AllTierStats": null,
"Compacted": false
},
"v9-bucket": {
"Children": {
"v9-bucket/nested": {}
},
"Size": 1234,
"HotTierSize": 1234,
"Objects": 1,
"Versions": 1,
"DeleteMarkers": 0,
"ObjSizes": [
0,
1,
0,
0,
0,
0,
0,
0,
0,
0,
0
],
"ObjVersions": [
0,
1,
0,
0,
0,
0,
0
],
"AllTierStats": null,
"Compacted": false
},
"v9-bucket/nested": {
"Children": null,
"Size": 73728,
"HotTierSize": 8192,
"Objects": 2,
"Versions": 4,
"DeleteMarkers": 2,
"ObjSizes": [
0,
1,
1,
0,
0,
0,
0,
0,
0,
0,
0
],
"ObjVersions": [
0,
1,
1,
0,
0,
0,
0
],
"AllTierStats": {
"Tiers": {
"COLD": {
"TotalSize": 65536,
"NumVersions": 1,
"NumObjects": 1
}
}
},
"Compacted": true
}
}
}
+82
View File
@@ -0,0 +1,82 @@
// 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/json"
"os"
"path/filepath"
"testing"
"time"
"github.com/minio/minio/internal/cachevalue"
)
// Run only on 89637554d, before access-tiering is removed. This uses the real
// folder scanner and v9 serializer with synthetic file-size/tier summaries,
// following TestDataUsageCacheSerialize; it does not relabel a v8 payload.
func TestGenerateAccessRetirementV9Fixture(t *testing.T) {
if dataUsageCacheVerCurrent != 9 { t.Fatal("requires the historical v9 writer") }
base := t.TempDir()
const bucket = "v9-bucket"
createUsageTestFiles(t, base, bucket, []usageTestFile{
{name: "root", size: 1234},
{name: "nested/versions", size: 8192},
{name: "nested/remote", size: 65536},
})
getSize := func(item scannerItem) (s sizeSummary, err error) {
if item.Typ&os.ModeDir != 0 { return s, nil }
info, err := os.Stat(item.Path)
if err != nil { return s, err }
s.totalSize, s.hotTierSize, s.versions = info.Size(), info.Size(), 1
if filepath.Base(item.Path) == "versions" { s.versions, s.deleteMarkers = 3, 2 }
if filepath.Base(item.Path) == "remote" {
s.hotTierSize = 0
s.tiers = map[string]tierStats{"COLD": {TotalSize: uint64(info.Size()), NumVersions: 1, NumObjects: 1}}
}
return s, nil
}
xls := xlStorage{drivePath: base, diskInfoCache: cachevalue.New[DiskInfo]()}
xls.diskInfoCache.InitOnce(time.Second, cachevalue.Opts{}, func(context.Context) (DiskInfo,error) {
return DiskInfo{Total: 1<<40, Free: 1<<40}, nil
})
cache, err := scanDataFolder(t.Context(), nil, &xls, dataUsageCache{Info:dataUsageCacheInfo{Name:bucket, SkipHealing:true}}, getSize, 0, func()bool{return false})
if err != nil { t.Fatal(err) }
root := *cache.find(bucket)
cache.replace(dataUsageRoot, "", dataUsageEntry{})
cache.replace(bucket, dataUsageRoot, root)
cache.Info.Name = dataUsageRoot
cache.Info.LastUpdate = time.Date(2026, 9, 15, 0, 0, 0, 0, time.UTC)
cache.Info.NextCycle = 42
flat := cache.flatten(*cache.root())
if flat.Size != 74962 || flat.Objects != 3 || flat.Versions != 5 || flat.DeleteMarkers != 2 || flat.HotTierSize != 9426 {
t.Fatalf("unexpected scanner fixture: %+v", flat)
}
var buf bytes.Buffer
if err := cache.serializeTo(&buf); err != nil { t.Fatal(err) }
dir := os.Getenv("SILO_RETIRE_FIXTURE_DIR")
if dir == "" { t.Fatal("SILO_RETIRE_FIXTURE_DIR required") }
if err := os.MkdirAll(dir, 0755); err != nil { t.Fatal(err) }
if err := os.WriteFile(filepath.Join(dir,"data-usage-v9.bin"),buf.Bytes(),0644);err != nil{t.Fatal(err)}
expected, err := json.MarshalIndent(cache,""," ")
if err != nil { t.Fatal(err) }
if err := os.WriteFile(filepath.Join(dir,"data-usage-v9.json"),append(expected,'\n'),0644);err != nil{t.Fatal(err)}
t.Logf("old scanner/v9 writer: bytes=%d size=%d objects=%d versions=%d markers=%d hts=%d",buf.Len(),flat.Size,flat.Objects,flat.Versions,flat.DeleteMarkers,flat.HotTierSize)
}
-5
View File
@@ -583,7 +583,6 @@ func (s *xlStorage) NSScanner(ctx context.Context, cache dataUsageCache, updates
}
poolIdx, setIdx, _ := s.GetDiskLoc()
hotPool, hotPoolOK := globalILMConfig.accessCfg().HotPool()
disks, err := objAPI.GetDisks(poolIdx, setIdx)
if err != nil {
@@ -593,7 +592,6 @@ func (s *xlStorage) NSScanner(ctx context.Context, cache dataUsageCache, updates
cache.Info.updates = updates
dataUsageInfo, err := scanDataFolder(ctx, disks, s, cache, func(item scannerItem) (sizeSummary, error) {
item.poolIdx = poolIdx
// Look for `xl.meta/xl.json' at the leaf.
if !strings.HasSuffix(item.Path, SlashSeparator+xlStorageFormatFile) &&
!strings.HasSuffix(item.Path, SlashSeparator+xlStorageFormatFileV1) {
@@ -657,9 +655,6 @@ func (s *xlStorage) NSScanner(ctx context.Context, cache dataUsageCache, updates
sizeS.versions++
}
sizeS.totalSize += sz
if hotPoolOK && poolIdx == hotPool {
sizeS.hotTierSize += sz
}
// Skip tier accounting if object version is a delete-marker or a free-version
// tracking deleted transitioned objects