diff --git a/staging/src/k8s.io/apiserver/pkg/storage/cacher/cacher_testing_utils_test.go b/staging/src/k8s.io/apiserver/pkg/storage/cacher/cacher_testing_utils_test.go index cf83231016b..9090e4c6923 100644 --- a/staging/src/k8s.io/apiserver/pkg/storage/cacher/cacher_testing_utils_test.go +++ b/staging/src/k8s.io/apiserver/pkg/storage/cacher/cacher_testing_utils_test.go @@ -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 } diff --git a/staging/src/k8s.io/apiserver/pkg/storage/etcd3/stats.go b/staging/src/k8s.io/apiserver/pkg/storage/etcd3/stats.go index af7bae69377..e2dc7f22f47 100644 --- a/staging/src/k8s.io/apiserver/pkg/storage/etcd3/stats.go +++ b/staging/src/k8s.io/apiserver/pkg/storage/etcd3/stats.go @@ -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) { diff --git a/staging/src/k8s.io/apiserver/pkg/storage/etcd3/stats_test.go b/staging/src/k8s.io/apiserver/pkg/storage/etcd3/stats_test.go index bc64c9973db..77637fe3336 100644 --- a/staging/src/k8s.io/apiserver/pkg/storage/etcd3/stats_test.go +++ b/staging/src/k8s.io/apiserver/pkg/storage/etcd3/stats_test.go @@ -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) diff --git a/staging/src/k8s.io/apiserver/pkg/storage/etcd3/store.go b/staging/src/k8s.io/apiserver/pkg/storage/etcd3/store.go index dbd7f675d4b..808a87f2c6d 100644 --- a/staging/src/k8s.io/apiserver/pkg/storage/etcd3/store.go +++ b/staging/src/k8s.io/apiserver/pkg/storage/etcd3/store.go @@ -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) diff --git a/staging/src/k8s.io/apiserver/pkg/storage/etcd3/store_test.go b/staging/src/k8s.io/apiserver/pkg/storage/etcd3/store_test.go index b91c2192a65..6d425a0b255 100644 --- a/staging/src/k8s.io/apiserver/pkg/storage/etcd3/store_test.go +++ b/staging/src/k8s.io/apiserver/pkg/storage/etcd3/store_test.go @@ -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 } diff --git a/staging/src/k8s.io/apiserver/pkg/storage/storagebackend/factory/etcd3.go b/staging/src/k8s.io/apiserver/pkg/storage/storagebackend/factory/etcd3.go index 1397cfa4b83..0006afe0e3e 100644 --- a/staging/src/k8s.io/apiserver/pkg/storage/storagebackend/factory/etcd3.go +++ b/staging/src/k8s.io/apiserver/pkg/storage/storagebackend/factory/etcd3.go @@ -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