Run background cleanup goroutine

This commit is contained in:
Marek Siarkowicz
2025-06-18 17:24:17 +02:00
parent ec78b8305a
commit 7cb2417999
6 changed files with 76 additions and 17 deletions

View File

@@ -71,6 +71,7 @@ func newEtcdTestStorage(t testing.TB, prefix string) (*etcd3testing.EtcdTestServ
etcd3.NewDefaultLeaseManagerConfig(),
etcd3.NewDefaultDecoder(codec, versioner),
versioner)
t.Cleanup(storage.Close)
return server, storage
}

View File

@@ -20,17 +20,31 @@ import (
"context"
"sync"
"sync/atomic"
"time"
"go.etcd.io/etcd/api/v3/mvccpb"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/apiserver/pkg/storage"
"k8s.io/klog/v2"
)
const sizerRefreshInterval = time.Minute
type keysFunc func(context.Context) ([]string, error)
func newStatsCache(getKeys keysFunc) *statsCache {
sc := &statsCache{
getKeys: getKeys,
stop: make(chan struct{}),
keys: make(map[string]sizeRevision),
}
sc.wg.Add(1)
go func() {
defer sc.wg.Done()
sc.run()
}()
return sc
}
@@ -44,7 +58,10 @@ func newStatsCache(getKeys keysFunc) *statsCache {
// This approach may leak keys if delete events are not observed,
// thus we run a background goroutine to periodically cleanup keys if needed.
type statsCache struct {
getKeys keysFunc
getKeys keysFunc
stop chan struct{}
wg sync.WaitGroup
lastKeyCleanup atomic.Pointer[time.Time]
lock sync.Mutex
keys map[string]sizeRevision
@@ -72,6 +89,36 @@ func (sc *statsCache) Stats(ctx context.Context) (storage.Stats, error) {
return stats, nil
}
func (sc *statsCache) Close() {
close(sc.stop)
sc.wg.Wait()
}
func (sc *statsCache) run() {
err := wait.PollUntilContextCancel(wait.ContextForChannel(sc.stop), sizerRefreshInterval, false, func(ctx context.Context) (done bool, err error) {
sc.cleanKeysIfNeeded(ctx)
return false, nil
})
if err != nil {
klog.InfoS("Sizer exiting")
}
}
func (sc *statsCache) cleanKeysIfNeeded(ctx context.Context) {
lastKeyCleanup := sc.lastKeyCleanup.Load()
if lastKeyCleanup != nil && time.Since(*lastKeyCleanup) < sizerRefreshInterval {
return
}
// Don't execute getKeys under lock.
keys, err := sc.getKeys(ctx)
if err != nil {
klog.InfoS("Error getting keys", "err", err)
}
sc.lock.Lock()
defer sc.lock.Unlock()
sc.cleanKeys(keys)
}
func (sc *statsCache) cleanKeys(keepKeys []string) {
newKeys := make(map[string]sizeRevision, len(keepKeys))
for _, key := range keepKeys {
@@ -82,6 +129,8 @@ func (sc *statsCache) cleanKeys(keepKeys []string) {
newKeys[key] = keySizeRevision
}
sc.keys = newKeys
now := time.Now()
sc.lastKeyCleanup.Store(&now)
}
func (sc *statsCache) keySizes() (totalSize int64) {

View File

@@ -28,6 +28,7 @@ import (
func TestStatsCache(t *testing.T) {
ctx := t.Context()
store := newStatsCache(func(ctx context.Context) ([]string, error) { return []string{}, nil })
defer store.Close()
stats, err := store.Stats(ctx)
require.NoError(t, err)

View File

@@ -202,6 +202,12 @@ func (s *store) Versioner() storage.Versioner {
return s.versioner
}
func (s *store) Close() {
if s.stats != nil {
s.stats.Close()
}
}
// Get implements storage.Interface.Get.
func (s *store) Get(ctx context.Context, key string, opts storage.GetOptions, out runtime.Object) error {
preparedKey, err := s.prepareKey(key)

View File

@@ -600,6 +600,7 @@ func testSetup(t testing.TB, opts ...setupOption) (context.Context, *store, *kub
NewDefaultDecoder(setupOpts.codec, versioner),
versioner,
)
t.Cleanup(store.Close)
ctx := context.Background()
return ctx, store, client
}

View File

@@ -445,17 +445,6 @@ func newETCD3Storage(c storagebackend.ConfigForResource, newFunc, newListFunc fu
return nil, nil, err
}
var once sync.Once
destroyFunc := func() {
// we know that storage destroy funcs are called multiple times (due to reuse in subresources).
// Hence, we only destroy once.
// TODO: fix duplicated storage destroy calls higher level
once.Do(func() {
stopCompactor()
stopDBSizeMonitor()
client.Close()
})
}
transformer := c.Transformer
if transformer == nil {
transformer = identity.NewEncryptCheckTransformer()
@@ -468,12 +457,24 @@ func newETCD3Storage(c storagebackend.ConfigForResource, newFunc, newListFunc fu
transformer = etcd3.WithCorruptObjErrorHandlingTransformer(transformer)
decoder = etcd3.WithCorruptObjErrorHandlingDecoder(decoder)
}
var store storage.Interface
store = etcd3.New(client, c.Codec, newFunc, newListFunc, c.Prefix, resourcePrefix, c.GroupResource, transformer, c.LeaseManagerConfig, decoder, versioner)
if utilfeature.DefaultFeatureGate.Enabled(genericfeatures.AllowUnsafeMalformedObjectDeletion) {
store = etcd3.NewStoreWithUnsafeCorruptObjectDeletion(store, c.GroupResource)
store := etcd3.New(client, c.Codec, newFunc, newListFunc, c.Prefix, resourcePrefix, c.GroupResource, transformer, c.LeaseManagerConfig, decoder, versioner)
var once sync.Once
destroyFunc := func() {
// we know that storage destroy funcs are called multiple times (due to reuse in subresources).
// Hence, we only destroy once.
// TODO: fix duplicated storage destroy calls higher level
once.Do(func() {
stopCompactor()
stopDBSizeMonitor()
store.Close()
_ = client.Close()
})
}
return store, destroyFunc, nil
var storage storage.Interface = store
if utilfeature.DefaultFeatureGate.Enabled(genericfeatures.AllowUnsafeMalformedObjectDeletion) {
storage = etcd3.NewStoreWithUnsafeCorruptObjectDeletion(storage, c.GroupResource)
}
return storage, destroyFunc, nil
}
// startDBSizeMonitorPerEndpoint starts a loop to monitor etcd database size and update the