Skip to content
Open
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 35 additions & 29 deletions CubeMaster/pkg/base/localcache/localcache.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,8 @@ type LocalCache struct {
localCacheConfig *LocalCacheConfig
sharedCalls util.SharedCalls
consecutiveFailNum int64
expiredUse atomic.Bool
destroyOnce sync.Once
}

func NewCache(name string, loader LoaderFunc, localCacheConfig *LocalCacheConfig) *LocalCache {
Expand All @@ -77,6 +79,7 @@ func NewCache(name string, loader LoaderFunc, localCacheConfig *LocalCacheConfig
localCache.loadFile(localCacheConfig.LoadFileName)
}
localCache.localCacheConfig = localCache.SetupConfig(localCacheConfig)
localCache.expiredUse.Store(localCache.localCacheConfig.ExpiredUse)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

SetupConfig deliberately guards a nil config (if localCacheConfig == nil { return nil }), so localCache.localCacheConfig can be nil here and this new unconditional localCacheConfig.ExpiredUse dereference turns a nil config into a panic at construction time — previously the panic was deferred to the first Get on an expired entry. No current caller passes nil, so this is latent rather than a live bug, but the new code bypasses the nil guard that SetupConfig provides. Consider localCache.expiredUse.Store(localCacheConfig != nil && localCacheConfig.ExpiredUse) or an explicit nil check.


localCache.chShrinkCache = make(chan bool, 1)
localCache.chCacheExit = make(chan bool)
Expand All @@ -99,36 +102,39 @@ func NewCache(name string, loader LoaderFunc, localCacheConfig *LocalCacheConfig
}

func (localCache *LocalCache) Destroy() {
if localCache != nil {
if localCache == nil {
return
}
localCache.destroyOnce.Do(func() {
CubeLog.Infof("LruCache(%s) Destroy", localCache.name)
if localCache.chCacheExit != nil {
if localCache.localCacheConfig.OpenCacheFile {
localCache.saveFile(localCache.localCacheConfig.LoadFileName)
}
localCache.cache.Flush()
close(localCache.chCacheExit)
localCache.waitGroup.Wait()
localCache.chCacheExit = nil
if localCache.chCacheExit == nil {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Dead code now: after this PR nothing assigns chCacheExit = nil (that was the race being removed), so this guard can only be true for a LocalCache not created via NewCache. In that case the surrounding destroyOnce.Do marks the cache as destroyed and a later, legitimate Destroy is silently skipped. Consider removing the guard, or making the "not initialized via NewCache" no-op explicit.

return
}
}
if localCache.localCacheConfig.OpenCacheFile {
localCache.saveFile(localCache.localCacheConfig.LoadFileName)
}
localCache.cache.Flush()
close(localCache.chCacheExit)
localCache.waitGroup.Wait()
})
}

func (localCache *LocalCache) Get(ctx context.Context, key string) (interface{}, bool, error) {
item, found := localCache.cache.Get(key)
if found {
element := item.(*list.Element)

localCache.Lock()

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The PR description says Get reads the *CacheValue pointer and does its MoveToBack inside one critical section, but the MoveToBack at the end of Get is still a separate lock acquisition. The separation is actually load-bearing — promoting the entry before the refresh decision would make a failed refresh promote the entry, which TestFailingRefreshDoesNotPromoteTheEntry explicitly forbids — so the code is correct, just worth a comment. Note it also means every hit now acquires the cache-wide mutex twice (snapshot read + MoveToBack), doubling exclusive-lock contention on the hot read path of this shared cache.

itm := element.Value.(*util.CacheValue)
localCache.valueList.MoveToBack(element)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Behavior change vs. the PR description. The description says "the LRU effect is the same — exactly one MoveToBack per hit, on access", but on the default config (ExpiredUse: false) the old code returned early from the sync-refresh branch — return r, f, errwithout MoveToBack, both on refresh success and on refresh failure without demotion. Moving MoveToBack to before the refresh decision means every found hit now promotes the element, including entries whose loader keeps failing with demotion disabled. Those failing entries are now pinned at the back of the list, so shrinkCache's front-based eviction won't reclaim them; they live until the CacheExpiredRemove ticker (default 24 h) instead of drifting toward the eviction front. If this is intended, it's worth stating in the description; if not, the promotion should stay on the non-refresh hit path only.

localCache.Unlock()

if time.Now().Add(-itm.Expired).After(time.Unix(itm.LastAccess, 0)) {

if !localCache.localCacheConfig.ExpiredUse {
if !localCache.expiredUse.Load() {
r, f, err := localCache.loadAndRefresh(ctx, key)

if err != nil && localCache.localCacheConfig.DemotionExpiredUse {

localCache.Lock()
localCache.valueList.MoveToBack(element)
localCache.Unlock()
return itm.Value, true, nil
}

Expand All @@ -142,9 +148,6 @@ func (localCache *LocalCache) Get(ctx context.Context, key string) (interface{},
localCache.asyncRefresh(ctx, key)
}

localCache.Lock()
localCache.valueList.MoveToBack(element)
localCache.Unlock()
return itm.Value, true, nil
}

Expand All @@ -159,14 +162,17 @@ func (localCache *LocalCache) put(key string, val interface{}, expired time.Dura
if item, found := localCache.cache.Get(key); found {
element := item.(*list.Element)
localCache.Lock()
prev := element.Value.(*util.CacheValue)
next := &util.CacheValue{
Key: prev.Key,
Value: val,
LastAccess: time.Now().Unix(),
Expired: expired}
element.Value = next

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

put() now writes element.Value = next while holding localCache's mutex, but saveFile() (unchanged in this PR, called from Destroy()) reads the same *list.Element.Value field — value.Value.(*util.CacheValue) at line 328 — without holding that mutex.

This is a newly-introduced data race. Before this PR, element.Value was write-once (assigned at PushBack), so saveFile's unlocked read was safe; now every refresh of an existing key writes the field. The race is reachable: asyncRefresh() spawns goroutines that are not tracked by waitGroup, and concurrent Gets can still be inside loadAndRefresh → put while Destroy() runs saveFile. The new race tests never overlap Destroy with an in-flight put, so the reported green -race run wouldn't catch it.

Fix suggestion: snapshot the *util.CacheValue under localCache.Lock() in saveFile (e.g. read value.Value into a local while holding the lock, then gob-encode after), so all readers of element.Value are serialized.

localCache.valueList.MoveToBack(element)
localCache.Unlock()
itm := element.Value.(*util.CacheValue)
atomic.AddInt64(&localCache.curCacheSize, -itm.Size())
itm.Value = val
itm.Expired = expired
itm.LastAccess = time.Now().Unix()
atomic.AddInt64(&localCache.curCacheSize, itm.Size())
atomic.AddInt64(&localCache.curCacheSize, -prev.Size())
atomic.AddInt64(&localCache.curCacheSize, next.Size())
} else {
itm := &util.CacheValue{
Key: key,
Expand All @@ -180,7 +186,7 @@ func (localCache *LocalCache) put(key string, val interface{}, expired time.Dura
localCache.cache.Set(key, element, -1)
}

if localCache.curCacheSize >= localCache.localCacheConfig.HighCacheSize {
if atomic.LoadInt64(&localCache.curCacheSize) >= localCache.localCacheConfig.HighCacheSize {
localCache.chShrinkCache <- true
}
}
Expand All @@ -198,7 +204,7 @@ func (localCache *LocalCache) loadAndRefresh(ctx context.Context, key string) (i
CubeLog.Errorf("Cache LoadAndRefresh Error:%s, %v, %s", key, found, err)
return nil, err
} else {
localCache.consecutiveFailNum = 0
atomic.StoreInt64(&localCache.consecutiveFailNum, 0)
}

if !found {
Expand Down Expand Up @@ -305,12 +311,12 @@ func (localCache *LocalCache) errStrategy() {
int64(localCache.localCacheConfig.MaxConsecutiveFailNum) {

if localCache.localCacheConfig.DemotionExpiredUse {
localCache.localCacheConfig.ExpiredUse = true
localCache.expiredUse.Store(true)
}
} else {

if localCache.localCacheConfig.DemotionExpiredUse && localCache.localCacheConfig.ExpiredUse {
localCache.localCacheConfig.ExpiredUse = false
if localCache.localCacheConfig.DemotionExpiredUse && localCache.expiredUse.Load() {
localCache.expiredUse.Store(false)
}
}
}
Expand Down
90 changes: 90 additions & 0 deletions CubeMaster/pkg/base/localcache/localcache_race_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
// Copyright (c) 2024 Tencent Inc.
// SPDX-License-Identifier: Apache-2.0
//

package localcache

import (
"context"
"sync"
"sync/atomic"
"testing"
"time"
)

func TestConcurrentGetAndRefreshOnSameKey(t *testing.T) {
var loads int64
localCache := NewCache("race-probe",
func(ctx context.Context, key string) (interface{}, bool, error) {
atomic.AddInt64(&loads, 1)
return RandString(16), true, nil
},
&LocalCacheConfig{
LowCacheSize: 1000000,
HighCacheSize: 2000000,
Expired: time.Millisecond,
AsyncRefreshBefore: time.Millisecond,
MaxAsyncRefreshNum: 100,
ExpiredUse: true,
})
defer localCache.Destroy()

ctx := context.Background()
const readers = 32
const iterations = 300

var wg sync.WaitGroup
wg.Add(readers)
for i := 0; i < readers; i++ {
go func() {
defer wg.Done()
for j := 0; j < iterations; j++ {
v, found, err := localCache.Get(ctx, "hot-key")
if err != nil {
t.Errorf("Get: %v", err)
return
}
if found {
if _, ok := v.(string); !ok {
t.Errorf("Get returned %T, want string; a torn interface read", v)
return
}
}
}
}()
}
wg.Wait()

if atomic.LoadInt64(&loads) == 0 {
t.Fatal("loader never ran; the test did not exercise the refresh path")
}
}

func TestDestroyIsIdempotent(t *testing.T) {
localCache := NewCache("destroy-probe",
func(ctx context.Context, key string) (interface{}, bool, error) {
return "v", true, nil
},
&LocalCacheConfig{LowCacheSize: 1000, HighCacheSize: 2000, Expired: time.Minute})

localCache.Destroy()
localCache.Destroy()
}

func TestConcurrentDestroyDoesNotRaceWithBackgroundLoops(t *testing.T) {
localCache := NewCache("destroy-race-probe",
func(ctx context.Context, key string) (interface{}, bool, error) {
return "v", true, nil
},
&LocalCacheConfig{LowCacheSize: 1000, HighCacheSize: 2000, Expired: time.Minute})

var wg sync.WaitGroup
wg.Add(4)
for i := 0; i < 4; i++ {
go func() {
defer wg.Done()
localCache.Destroy()
}()
}
wg.Wait()
}
8 changes: 4 additions & 4 deletions CubeMaster/pkg/base/localcache/localcache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,13 +45,13 @@ func TestLocalCache(t *testing.T) {
for i := 0; i < 2700; i++ {

wg.Add(1)
ctx = context.WithValue(ctx, ctxKey, i)
go func() {
iterCtx := context.WithValue(ctx, ctxKey, i)
go func(c context.Context) {
defer wg.Done()
start := time.Now()
localCache.Get(ctx, RandString(8))
localCache.Get(c, RandString(8))
fmt.Printf("=====%d====\n", time.Since(start).Milliseconds())
}()
}(iterCtx)

}
wg.Wait()
Expand Down
Loading