Skip to content
Open
Show file tree
Hide file tree
Changes from all 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
67 changes: 41 additions & 26 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 != nil && localCache.localCacheConfig.ExpiredUse)

localCache.chShrinkCache = make(chan bool, 1)
localCache.chCacheExit = make(chan bool)
Expand All @@ -99,29 +102,35 @@ 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.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 {
Expand Down Expand Up @@ -159,14 +168,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 +192,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 +210,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 @@ -240,10 +252,10 @@ func (localCache *LocalCache) shrinkCache() {
select {
case <-localCache.chShrinkCache:
curTime := time.Now()
curCacheSize := localCache.curCacheSize
curCacheSize := atomic.LoadInt64(&localCache.curCacheSize)
var shrinkNum, shrinkSize, size int64
for {
if localCache.curCacheSize > localCache.localCacheConfig.LowCacheSize {
if atomic.LoadInt64(&localCache.curCacheSize) > localCache.localCacheConfig.LowCacheSize {
localCache.Lock()
element := localCache.valueList.Front()
if element == nil {
Expand Down Expand Up @@ -305,12 +317,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 All @@ -324,10 +336,13 @@ func (localCache *LocalCache) saveFile(file string) {
var itm *util.CacheValue
switch value := item.Object.(type) {
case *list.Element:
localCache.Lock()
stored := value.Value
localCache.Unlock()
var ok bool
itm, ok = value.Value.(*util.CacheValue)
itm, ok = stored.(*util.CacheValue)
if !ok {
CubeLog.Errorf("Cache(%s) cannot persist key %s: list element contains %T", localCache.name, key, value.Value)
CubeLog.Errorf("Cache(%s) cannot persist key %s: list element contains %T", localCache.name, key, stored)
continue
}
case *util.CacheValue:
Expand Down
220 changes: 220 additions & 0 deletions CubeMaster/pkg/base/localcache/localcache_race_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,220 @@
// Copyright (c) 2024 Tencent Inc.
// SPDX-License-Identifier: Apache-2.0
//

package localcache

import (
"context"
"errors"
"path/filepath"
"strconv"
"sync"
"sync/atomic"
"testing"
"time"

"github.com/tencentcloud/CubeSandbox/CubeMaster/pkg/base/localcache/util"
)

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()
}

func TestDestroySnapshotDoesNotRaceWithInFlightRefresh(t *testing.T) {
localCache := NewCache("destroy-snapshot-probe",
func(ctx context.Context, key string) (interface{}, bool, error) {
return RandString(16), true, nil
},
&LocalCacheConfig{
LowCacheSize: 1000000,
HighCacheSize: 2000000,
Expired: time.Millisecond,
AsyncRefreshBefore: time.Millisecond,
MaxAsyncRefreshNum: 100,
ExpiredUse: true,
OpenCacheFile: true,
LoadFileName: filepath.Join(t.TempDir(), "cache.gob"),
})

ctx := context.Background()
const keys = 64
for i := 0; i < keys; i++ {
if _, _, err := localCache.Get(ctx, "k"+strconv.Itoa(i)); err != nil {
t.Fatalf("seed Get: %v", err)
}
}

var stop atomic.Bool
var wg sync.WaitGroup
wg.Add(8)
for r := 0; r < 8; r++ {
go func(r int) {
defer wg.Done()
for i := 0; !stop.Load(); i++ {
if _, _, err := localCache.Get(ctx, "k"+strconv.Itoa((i+r)%keys)); err != nil {
t.Errorf("Get: %v", err)
return
}
}
}(r)
}

time.Sleep(50 * time.Millisecond)
localCache.Destroy()
stop.Store(true)
wg.Wait()
}

func frontKey(t *testing.T, c *LocalCache) string {
t.Helper()
c.Lock()
defer c.Unlock()
front := c.valueList.Front()
if front == nil {
t.Fatal("value list is empty")
}
return front.Value.(*util.CacheValue).Key
}

func TestFailingRefreshDoesNotPromoteTheEntry(t *testing.T) {
var fail atomic.Bool
localCache := NewCache("lru-demotion-probe",
func(ctx context.Context, key string) (interface{}, bool, error) {
if fail.Load() && key == "a" {
return nil, false, errors.New("loader is down")
}
return "v-" + key, true, nil
},
&LocalCacheConfig{
LowCacheSize: 1000000,
HighCacheSize: 2000000,
Expired: time.Millisecond,
ExpiredUse: false,
DemotionExpiredUse: false,
})
defer localCache.Destroy()

ctx := context.Background()
for _, k := range []string{"a", "b"} {
if _, _, err := localCache.Get(ctx, k); err != nil {
t.Fatalf("seed Get(%s): %v", k, err)
}
}
if got := frontKey(t, localCache); got != "a" {
t.Fatalf("front is %q before the probe, want a", got)
}

fail.Store(true)
time.Sleep(5 * time.Millisecond)

if _, _, err := localCache.Get(ctx, "a"); err == nil {
t.Fatal("Get(a) succeeded while the loader was failing")
}
if got := frontKey(t, localCache); got != "a" {
t.Fatalf("a failing entry was promoted: front is %q, want a to stay evictable", got)
}
}

func TestSuccessfulHitPromotesTheEntry(t *testing.T) {
localCache := NewCache("lru-promote-probe",
func(ctx context.Context, key string) (interface{}, bool, error) {
return "v-" + key, true, nil
},
&LocalCacheConfig{
LowCacheSize: 1000000,
HighCacheSize: 2000000,
Expired: time.Hour,
})
defer localCache.Destroy()

ctx := context.Background()
for _, k := range []string{"a", "b"} {
if _, _, err := localCache.Get(ctx, k); err != nil {
t.Fatalf("seed Get(%s): %v", k, err)
}
}
if got := frontKey(t, localCache); got != "a" {
t.Fatalf("front is %q before the probe, want a", got)
}

if _, _, err := localCache.Get(ctx, "a"); err != nil {
t.Fatalf("Get(a): %v", err)
}
if got := frontKey(t, localCache); got != "b" {
t.Fatalf("a fresh hit did not promote: front is %q, want b", got)
}
}
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