From 23a9e249930fbd604df5433f510c2289d533a7ef Mon Sep 17 00:00:00 2001 From: disksing Date: Tue, 16 Sep 2025 16:16:50 +0800 Subject: [PATCH 01/21] resource_group: implement async loading for resource groups to improve startup performance (#411) * feat: implement async loading for resource groups - Add async loading mechanism to reduce startup time - Use atomic operations for loading state management - Implement lazy loading for individual resource groups - Add retry mechanism with infinite retries for reliability - Return error for list requests during loading - Extend storage interface for single group loading - Optimize loading logic to avoid partial loading issues This change significantly improves startup performance by loading resource groups asynchronously while maintaining data integrity. Signed-off-by: disksing * tiny fix Signed-off-by: disksing * test: add comprehensive async loading test with simplified blocking control - Simplify control mechanism from 4 to 2 control points: * blockBeforeLoad: blocks before starting load operation * blockAfterLoad: blocks after loading is completed - Add more test data (test-group-3, test-group-4) for better coverage - Test operations during async loading (read, update, delete, list) - Test operations after async loading completes - Verify syncLoadedGroups mechanism prevents group resurrection - Ensure proper error handling during loading state This test validates the complete async loading workflow with simplified control and comprehensive scenario coverage. Signed-off-by: disksing * fix: resolve testifylint issues in test files - Remove unnecessary fmt.Sprintf calls in assert messages - Use require instead of assert for error assertions - Remove unused fmt imports This fixes all testifylint warnings in the test files. Signed-off-by: disksing * fix: replace assert.NoError with require.NoError in manager_async_test.go - Fix testifylint require-error violations on lines 342, 347, and 353 - Use require.NoError for error assertions to ensure test stops on failure Signed-off-by: disksing * feat: add metrics for resource group loading operations - Add asyncLoadGroupDuration histogram to track async loading performance - Add syncLoadGroupCounter to count synchronous loading operations - Include duration in async loading completion logs for better observability This helps monitor the performance of resource group loading and understand the loading patterns in the system. Signed-off-by: disksing * update error code Signed-off-by: disksing * minor fix Signed-off-by: disksing * fix default group Signed-off-by: disksing * extract addDefaultGroup Signed-off-by: disksing * minor fix Signed-off-by: disksing * fix static check Signed-off-by: disksing * fix lint Signed-off-by: disksing * fix manager reload Signed-off-by: disksing * fix update default group Signed-off-by: disksing * fix test Signed-off-by: disksing * fix when load failed Signed-off-by: disksing --------- Signed-off-by: disksing (cherry picked from commit da1b8ba1e3873401aef0fcbd99c0898a654952e5) --- errors.toml | 5 + pkg/errs/errno.go | 1 + pkg/mcs/resourcemanager/server/manager.go | 282 ++++++++++++++++-- .../server/manager_async_test.go | 150 ++++++++++ .../resourcemanager/server/manager_test.go | 4 +- pkg/mcs/resourcemanager/server/metrics.go | 18 ++ pkg/storage/endpoint/resource_group.go | 12 + .../resourcemanager/resource_manager_test.go | 9 + 8 files changed, 454 insertions(+), 27 deletions(-) create mode 100644 pkg/mcs/resourcemanager/server/manager_async_test.go diff --git a/errors.toml b/errors.toml index 33acad6fbfb..791560c9de1 100644 --- a/errors.toml +++ b/errors.toml @@ -926,6 +926,11 @@ error = ''' empty region ''' +["PD:resourcemanager:ErrResourceGroupsLoading"] +error = ''' +resource groups are still being loaded, please try again later +''' + ["PD:schedule:ErrCreateOperator"] error = ''' unable to create operator, %s diff --git a/pkg/errs/errno.go b/pkg/errs/errno.go index a887d539f69..415f51061a3 100644 --- a/pkg/errs/errno.go +++ b/pkg/errs/errno.go @@ -541,6 +541,7 @@ var ( ErrResourceGroupNotExists = errors.Normalize("the %s resource group does not exist", errors.RFCCodeText("PD:resourcemanager:ErrGroupNotExists")) ErrDeleteReservedGroup = errors.Normalize("cannot delete reserved group", errors.RFCCodeText("PD:resourcemanager:ErrDeleteReservedGroup")) ErrInvalidGroup = errors.Normalize("invalid group settings, please check %s", errors.RFCCodeText("PD:resourcemanager:ErrInvalidGroup")) + ErrResourceGroupsLoading = errors.Normalize("resource groups are still being loaded, please try again later", errors.RFCCodeText("PD:resourcemanager:ErrResourceGroupsLoading")) ) // Microservice errors diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 4f811a0392e..606223a5b95 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -21,6 +21,7 @@ import ( "os" "strings" "sync" + "sync/atomic" "time" "github.com/prometheus/client_golang/prometheus/push" @@ -114,8 +115,22 @@ type Manager struct { metrics *metrics // ruCollector is used to collect the RU metering data. ruCollector *ruCollector + // async loading state management + loadingState int32 // atomic access + // syncLoadedGroups records groups that were loaded synchronously (e.g., by lazy loading) + syncLoadedGroups map[trackerKey]bool } +// LoadingState represents the current loading state of resource groups +const ( + // LoadingStateNotStarted means resource groups haven't started loading + LoadingStateNotStarted int32 = iota + // LoadingStateInProgress means resource groups are being loaded asynchronously + LoadingStateInProgress + // LoadingStateCompleted means all resource groups have been loaded + LoadingStateCompleted +) + // factoryProvider is a factory provider for the manager, which injects some specialized functions // that need to be retrieved from the `bs.Server` instance without interacting with its interface. type factoryProvider interface { @@ -149,6 +164,8 @@ func newManagerBase(controllerConfig *ControllerConfig, writeRole ResourceGroupW keyspaceIDLookup: make(map[string]uint32), metrics: newMetrics(), ruCollector: newRUCollector(), + loadingState: LoadingStateNotStarted, + syncLoadedGroups: make(map[trackerKey]bool), } } @@ -194,7 +211,10 @@ func NewMetadataOnlyManager[T metadataFactoryProvider](srv bs.Server) (*Manager, nil, ) m.srv = srv - if err := m.initMetadata(); err != nil { + if err := m.initControllerConfig(); err != nil { + return nil, err + } + if err := m.loadKeyspaceResourceGroups(); err != nil { return nil, err } return m, nil @@ -325,13 +345,16 @@ func (m *Manager) Init(ctx context.Context) error { m.wg.Wait() return err } + atomic.StoreInt32(&m.loadingState, LoadingStateCompleted) } else { - if err := m.initMetadata(); err != nil { - return err - } // This context is derived from the leader/primary context, it will be canceled // from the outside loop when the leader/primary step down. ctx, m.cancel = context.WithCancel(ctx) + if err := m.initMetadata(ctx); err != nil { + m.cancel() + m.wg.Wait() + return err + } } m.wg.Add(1) // Start the background metrics flusher. @@ -369,34 +392,140 @@ func (m *Manager) initControllerConfig() error { return nil } -func (m *Manager) initMetadata() error { +func (m *Manager) initMetadata(ctx context.Context) error { if err := m.initControllerConfig(); err != nil { return err } - // Load keyspace resource groups from the storage. - return m.loadKeyspaceResourceGroups() + m.Lock() + m.krgms = make(map[uint32]*keyspaceResourceGroupManager) + m.syncLoadedGroups = make(map[trackerKey]bool) + atomic.StoreInt32(&m.loadingState, LoadingStateNotStarted) + m.Unlock() + + m.initReserved() + if err := m.loadServiceLimits(); err != nil { + return err + } + + m.wg.Add(1) + go m.asyncLoadResourceGroups(ctx) + return nil +} + +func (m *Manager) loadServiceLimits() error { + return m.storage.LoadServiceLimits(func(keyspaceID uint32, serviceLimit float64) { + m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false).setServiceLimitFromStorage(serviceLimit) + }) } func (m *Manager) loadKeyspaceResourceGroups() error { - // Empty the keyspace resource group manager map before the loading. + tempKrgms, err := m.loadKeyspaceResourceGroupsFromStorage() + if err != nil { + return err + } m.Lock() - m.krgms = make(map[uint32]*keyspaceResourceGroupManager) + m.krgms = tempKrgms + m.syncLoadedGroups = nil + atomic.StoreInt32(&m.loadingState, LoadingStateCompleted) m.Unlock() - // Load keyspace resource group meta info from the storage. + m.initReserved() + return m.loadServiceLimits() +} + +func (m *Manager) asyncLoadResourceGroups(ctx context.Context) { + defer logutil.LogPanic() + defer m.wg.Done() + + const retryInterval = 10 * time.Second + retry := 0 + for { + select { + case <-ctx.Done(): + log.Info("async loading resource groups cancelled") + return + default: + } + if retry > 0 { + log.Info("retrying async loading resource groups", zap.Int("retry", retry)) + timer := time.NewTimer(retryInterval) + select { + case <-ctx.Done(): + timer.Stop() + log.Info("async loading resource groups cancelled") + return + case <-timer.C: + } + } + + atomic.StoreInt32(&m.loadingState, LoadingStateInProgress) + startTime := time.Now() + tempKrgms, err := m.loadKeyspaceResourceGroupsFromStorage() + if err != nil { + log.Error("failed to load resource groups", zap.Error(err), zap.Int("retry", retry)) + atomic.StoreInt32(&m.loadingState, LoadingStateNotStarted) + retry++ + continue + } + + loaded := 0 + m.Lock() + for keyspaceID, tempKrgm := range tempKrgms { + krgm := m.krgms[keyspaceID] + if krgm == nil { + krgm = newKeyspaceResourceGroupManager(keyspaceID, m.storage, m.writeRole) + m.krgms[keyspaceID] = krgm + } + groupsToSync := make([]*ResourceGroup, 0) + tempKrgm.RLock() + krgm.Lock() + for name, group := range tempKrgm.groups { + key := trackerKey{keyspaceID: keyspaceID, groupName: name} + if !m.syncLoadedGroups[key] { + krgm.groups[name] = group + groupsToSync = append(groupsToSync, group) + loaded++ + } + } + krgm.Unlock() + tempKrgm.RUnlock() + for _, group := range groupsToSync { + krgm.syncBurstabilityWithServiceLimit(group) + } + } + m.syncLoadedGroups = nil + m.Unlock() + + m.initReserved() + atomic.StoreInt32(&m.loadingState, LoadingStateCompleted) + duration := time.Since(startTime) + asyncLoadGroupDuration.Observe(duration.Seconds()) + log.Info("async loading resource groups completed", zap.Int("loaded-groups", loaded), zap.Duration("duration", duration)) + return + } +} + +func (m *Manager) loadKeyspaceResourceGroupsFromStorage() (map[uint32]*keyspaceResourceGroupManager, error) { + tempKrgms := make(map[uint32]*keyspaceResourceGroupManager) + getOrCreateTempKrgm := func(keyspaceID uint32) *keyspaceResourceGroupManager { + krgm, ok := tempKrgms[keyspaceID] + if !ok { + krgm = newKeyspaceResourceGroupManager(keyspaceID, m.storage, m.writeRole) + tempKrgms[keyspaceID] = krgm + } + return krgm + } if err := m.storage.LoadResourceGroupSettings(func(keyspaceID uint32, name string, rawValue string) { - // Since the default resource group might be loaded from the storage, we don't need to initialize it here. - err := m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false).addResourceGroupFromRaw(name, rawValue) + err := getOrCreateTempKrgm(keyspaceID).addResourceGroupFromRaw(name, rawValue) if err != nil { log.Error("failed to add resource group to the keyspace resource group manager", zap.Uint32("keyspace-id", keyspaceID), zap.String("group-name", name), zap.Error(err)) } }); err != nil { - return err + return nil, err } - // Load keyspace resource group states from the storage. if err := m.storage.LoadResourceGroupStates(func(keyspaceID uint32, name string, rawValue string) { - krgm := m.getKeyspaceResourceGroupManager(keyspaceID) + krgm := tempKrgms[keyspaceID] if krgm == nil { log.Warn("failed to get the corresponding keyspace resource group manager", zap.Uint32("keyspace-id", keyspaceID), zap.String("group-name", name)) @@ -408,18 +537,81 @@ func (m *Manager) loadKeyspaceResourceGroups() error { zap.Uint32("keyspace-id", keyspaceID), zap.String("group-name", name), zap.Error(err)) } }); err != nil { + return nil, err + } + return tempKrgms, nil +} + +func (m *Manager) loadResourceGroup(keyspaceID uint32, name string) (*ResourceGroup, error) { + rawValue, err := m.storage.LoadResourceGroupSetting(keyspaceID, name) + if err != nil { + return nil, err + } + if rawValue == "" { + return nil, errs.ErrResourceGroupNotExists.FastGenByArgs(name) + } + krgm := newKeyspaceResourceGroupManager(keyspaceID, m.storage, m.writeRole) + if err := krgm.addResourceGroupFromRaw(name, rawValue); err != nil { + return nil, err + } + state, err := m.storage.LoadResourceGroupState(keyspaceID, name) + if err == nil && state != "" { + if err := krgm.setRawStatesIntoResourceGroup(name, state); err != nil { + return nil, err + } + } + return krgm.getMutableResourceGroup(name), nil +} + +func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) error { + if atomic.LoadInt32(&m.loadingState) == LoadingStateCompleted { + return nil + } + krgm := m.getKeyspaceResourceGroupManager(keyspaceID) + if krgm != nil { + if group := krgm.getMutableResourceGroup(name); group != nil { + return nil + } + } + if name == DefaultResourceGroupName { + m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, true) + return nil + } + group, err := m.loadResourceGroup(keyspaceID, name) + if err != nil { return err } - // Initialize the reserved keyspace resource group manager and default resource groups. - m.initReserved() - // Load service limits from the storage after all resource groups are loaded. - return m.loadServiceLimits() + krgm = m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false) + inserted := false + krgm.Lock() + if _, exists := krgm.groups[name]; !exists { + krgm.groups[name] = group + inserted = true + } + krgm.Unlock() + if inserted { + krgm.syncBurstabilityWithServiceLimit(group) + } + markKey := trackerKey{keyspaceID: keyspaceID, groupName: name} + m.Lock() + if m.syncLoadedGroups != nil { + m.syncLoadedGroups[markKey] = true + } + m.Unlock() + syncLoadGroupCounter.Inc() + return nil } -func (m *Manager) loadServiceLimits() error { - return m.storage.LoadServiceLimits(func(keyspaceID uint32, serviceLimit float64) { - m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false).setServiceLimitFromStorage(serviceLimit) - }) +func (m *Manager) markResourceGroupSyncLoaded(keyspaceID uint32, name string) { + m.Lock() + defer m.Unlock() + if m.syncLoadedGroups != nil { + m.syncLoadedGroups[trackerKey{keyspaceID: keyspaceID, groupName: name}] = true + } +} + +func (m *Manager) isResourceGroupLoadingComplete() bool { + return atomic.LoadInt32(&m.loadingState) == LoadingStateCompleted } func cloneControllerConfig(cfg *ControllerConfig) *ControllerConfig { @@ -456,6 +648,7 @@ func (m *Manager) applyResourceGroupSettingFromRaw(keyspaceID uint32, name, rawV zap.Error(err)) return err } + m.markResourceGroupSyncLoaded(keyspaceID, name) return nil } @@ -492,6 +685,7 @@ func (m *Manager) applyResourceGroupStatesFromRaw(keyspaceID uint32, name, rawVa zap.Error(err)) return err } + m.markResourceGroupSyncLoaded(keyspaceID, name) return nil } @@ -574,7 +768,14 @@ func (m *Manager) AddResourceGroup(grouppb *rmpb.ResourceGroup) error { if krgm == nil { return errs.ErrKeyspaceNotExists.FastGenByArgs(keyspaceID) } - return krgm.addResourceGroup(grouppb) + if err := m.loadResourceGroupIfNeeded(keyspaceID, grouppb.Name); err != nil { + log.Debug("failed to load resource group", zap.Uint32("keyspace-id", keyspaceID), zap.String("name", grouppb.Name), zap.Error(err)) + } + if err := krgm.addResourceGroup(grouppb); err != nil { + return err + } + m.markResourceGroupSyncLoaded(keyspaceID, grouppb.Name) + return nil } // ModifyResourceGroup modifies an existing resource group. @@ -583,11 +784,19 @@ func (m *Manager) ModifyResourceGroup(grouppb *rmpb.ResourceGroup) error { return errMetadataWriteDisabled } keyspaceID := ExtractKeyspaceID(grouppb.GetKeyspaceId()) + if err := m.loadResourceGroupIfNeeded(keyspaceID, grouppb.Name); err != nil { + log.Debug("failed to load resource group", zap.Uint32("keyspace-id", keyspaceID), zap.String("name", grouppb.Name), zap.Error(err)) + return err + } krgm, err := m.accessKeyspaceResourceGroupManager(keyspaceID, grouppb.Name) if err != nil { return err } - return krgm.modifyResourceGroup(grouppb) + if err := krgm.modifyResourceGroup(grouppb); err != nil { + return err + } + m.markResourceGroupSyncLoaded(keyspaceID, grouppb.Name) + return nil } // DeleteResourceGroup deletes a resource group. @@ -595,16 +804,28 @@ func (m *Manager) DeleteResourceGroup(keyspaceID uint32, name string) error { if !m.writeRole.AllowsMetadataWrite() { return errMetadataWriteDisabled } + if err := m.loadResourceGroupIfNeeded(keyspaceID, name); err != nil { + log.Debug("failed to load resource group", zap.Uint32("keyspace-id", keyspaceID), zap.String("name", name), zap.Error(err)) + return err + } // "default" group can't be deleted, so there is not need to call accessKeyspaceResourceGroupManager krgm := m.getKeyspaceResourceGroupManager(keyspaceID) if krgm == nil { return errs.ErrKeyspaceNotExists.FastGenByArgs(keyspaceID) } - return krgm.deleteResourceGroup(name) + if err := krgm.deleteResourceGroup(name); err != nil { + return err + } + m.markResourceGroupSyncLoaded(keyspaceID, name) + return nil } // GetResourceGroup returns a copy of a resource group. func (m *Manager) GetResourceGroup(keyspaceID uint32, name string, withStats bool) (*ResourceGroup, error) { + if err := m.loadResourceGroupIfNeeded(keyspaceID, name); err != nil { + log.Debug("failed to load resource group", zap.Uint32("keyspace-id", keyspaceID), zap.String("name", name), zap.Error(err)) + return nil, err + } krgm, err := m.accessKeyspaceResourceGroupManager(keyspaceID, name) if err != nil { return nil, err @@ -614,6 +835,10 @@ func (m *Manager) GetResourceGroup(keyspaceID uint32, name string, withStats boo // GetMutableResourceGroup returns a mutable resource group. func (m *Manager) GetMutableResourceGroup(keyspaceID uint32, name string) (*ResourceGroup, error) { + if err := m.loadResourceGroupIfNeeded(keyspaceID, name); err != nil { + log.Debug("failed to load resource group", zap.Uint32("keyspace-id", keyspaceID), zap.String("name", name), zap.Error(err)) + return nil, err + } krgm, err := m.accessKeyspaceResourceGroupManager(keyspaceID, name) if err != nil { return nil, err @@ -622,7 +847,12 @@ func (m *Manager) GetMutableResourceGroup(keyspaceID uint32, name string) (*Reso } // GetResourceGroupList returns copies of resource group list. +// Returns error if resource groups are still being loaded asynchronously. func (m *Manager) GetResourceGroupList(keyspaceID uint32, withStats bool) ([]*ResourceGroup, error) { + if !m.isResourceGroupLoadingComplete() { + log.Debug("resource groups are still being loaded, cannot return list") + return nil, errs.ErrResourceGroupsLoading + } krgm, err := m.accessKeyspaceResourceGroupManager(keyspaceID, DefaultResourceGroupName) if err != nil { return nil, err diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go new file mode 100644 index 00000000000..7bf7d51ef1e --- /dev/null +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -0,0 +1,150 @@ +// Copyright 2026 TiKV Project Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package server + +import ( + "context" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/pingcap/kvproto/pkg/resource_manager" + + "github.com/tikv/pd/pkg/errs" + "github.com/tikv/pd/pkg/storage" + "github.com/tikv/pd/pkg/utils/testutil" +) + +type blockingResourceGroupStorage struct { + storage.Storage + + once sync.Once + entered chan struct{} + release chan struct{} +} + +func newBlockingResourceGroupStorage() *blockingResourceGroupStorage { + return &blockingResourceGroupStorage{ + Storage: storage.NewStorageWithMemoryBackend(), + entered: make(chan struct{}), + release: make(chan struct{}), + } +} + +func (s *blockingResourceGroupStorage) LoadResourceGroupSettings(f func(keyspaceID uint32, name, rawValue string)) error { + s.once.Do(func() { + close(s.entered) + <-s.release + }) + return s.Storage.LoadResourceGroupSettings(f) +} + +func (s *blockingResourceGroupStorage) waitEntered(t *testing.T) { + t.Helper() + select { + case <-s.entered: + case <-time.After(time.Second): + t.Fatal("timed out waiting for async resource group loading") + } +} + +func (s *blockingResourceGroupStorage) unblock() { + close(s.release) +} + +func newAsyncTestGroup(name string, fillRate uint64) *resource_manager.ResourceGroup { + return &resource_manager.ResourceGroup{ + Name: name, + Mode: resource_manager.GroupMode_RUMode, + Priority: middlePriority, + RUSettings: &resource_manager.GroupRequestUnitSettings{ + RU: &resource_manager.TokenBucket{ + Settings: &resource_manager.TokenLimitSettings{ + FillRate: fillRate, + BurstLimit: int64(fillRate), + }, + }, + }, + } +} + +func stopAsyncTestManager(m *Manager) { + if m.cancel != nil { + m.cancel() + } + m.wg.Wait() +} + +func TestAsyncLoadResourceGroupsLazyGet(t *testing.T) { + re := require.New(t) + store := newBlockingResourceGroupStorage() + re.NoError(store.SaveResourceGroupSetting(1, "lazy-group", newAsyncTestGroup("lazy-group", 100))) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + re.NoError(m.Init(context.Background())) + defer stopAsyncTestManager(m) + + store.waitEntered(t) + + _, err := m.GetResourceGroupList(1, false) + re.ErrorIs(err, errs.ErrResourceGroupsLoading) + + group, err := m.GetResourceGroup(1, "lazy-group", false) + re.NoError(err) + re.NotNil(group) + re.Equal("lazy-group", group.Name) + re.Equal(float64(100), group.RUSettings.RU.getFillRate()) + + store.unblock() + testutil.Eventually(re, func() bool { + groups, err := m.GetResourceGroupList(1, false) + return err == nil && len(groups) == 2 + }, testutil.WithTickInterval(20*time.Millisecond)) +} + +func TestAsyncLoadResourceGroupsDoesNotRestoreDeletedLazyGroup(t *testing.T) { + re := require.New(t) + store := newBlockingResourceGroupStorage() + re.NoError(store.SaveResourceGroupSetting(1, "deleted-group", newAsyncTestGroup("deleted-group", 100))) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + re.NoError(m.Init(context.Background())) + defer stopAsyncTestManager(m) + + store.waitEntered(t) + + group, err := m.GetResourceGroup(1, "deleted-group", false) + re.NoError(err) + re.NotNil(group) + re.NoError(m.DeleteResourceGroup(1, "deleted-group")) + + store.unblock() + testutil.Eventually(re, func() bool { + groups, err := m.GetResourceGroupList(1, false) + if err != nil { + return false + } + for _, group := range groups { + if group.Name == "deleted-group" { + return false + } + } + return true + }, testutil.WithTickInterval(20*time.Millisecond)) +} diff --git a/pkg/mcs/resourcemanager/server/manager_test.go b/pkg/mcs/resourcemanager/server/manager_test.go index 8d98ffc6dbc..7898fce3ba7 100644 --- a/pkg/mcs/resourcemanager/server/manager_test.go +++ b/pkg/mcs/resourcemanager/server/manager_test.go @@ -408,7 +408,9 @@ func TestInitManager(t *testing.T) { m.storage = storage err = m.Init(ctx) re.NoError(err) - re.Len(m.getKeyspaceResourceGroupManagers(), 2) + testutil.Eventually(re, func() bool { + return len(m.getKeyspaceResourceGroupManagers()) == 2 + }) // Get the default resource group. rg, err := m.GetResourceGroup(1, DefaultResourceGroupName, true) re.NoError(err) diff --git a/pkg/mcs/resourcemanager/server/metrics.go b/pkg/mcs/resourcemanager/server/metrics.go index 5d4907768bd..47b507efc83 100644 --- a/pkg/mcs/resourcemanager/server/metrics.go +++ b/pkg/mcs/resourcemanager/server/metrics.go @@ -220,6 +220,22 @@ var ( Help: "The duration of pushing RU metrics to Prometheus.", Buckets: prometheus.DefBuckets, }) + + syncLoadGroupCounter = prometheus.NewCounter( + prometheus.CounterOpts{ + Namespace: namespace, + Subsystem: serverSubsystem, + Name: "sync_load_group_counter", + Help: "The number of the sync load group.", + }) + + asyncLoadGroupDuration = prometheus.NewHistogram( + prometheus.HistogramOpts{ + Namespace: namespace, + Subsystem: serverSubsystem, + Name: "async_load_group_duration_seconds", + Help: "The duration of the async load group.", + }) ) type metrics struct { @@ -274,6 +290,8 @@ func init() { prometheus.MustRegister(overrideSettings) prometheus.MustRegister(serviceLimit) prometheus.MustRegister(pushRUMetricsDuration) + prometheus.MustRegister(syncLoadGroupCounter) + prometheus.MustRegister(asyncLoadGroupDuration) } func newMetrics() *metrics { diff --git a/pkg/storage/endpoint/resource_group.go b/pkg/storage/endpoint/resource_group.go index 68a62977721..a54eb7a271b 100644 --- a/pkg/storage/endpoint/resource_group.go +++ b/pkg/storage/endpoint/resource_group.go @@ -30,9 +30,11 @@ import ( // ResourceGroupStorage defines the storage operations on the resource group. type ResourceGroupStorage interface { LoadResourceGroupSettings(f func(keyspaceID uint32, name, rawValue string)) error + LoadResourceGroupSetting(keyspaceID uint32, name string) (string, error) SaveResourceGroupSetting(keyspaceID uint32, name string, msg proto.Message) error DeleteResourceGroupSetting(keyspaceID uint32, name string) error LoadResourceGroupStates(f func(keyspaceID uint32, name, rawValue string)) error + LoadResourceGroupState(keyspaceID uint32, name string) (string, error) SaveResourceGroupStates(keyspaceID uint32, name string, obj any) error DeleteResourceGroupStates(keyspaceID uint32, name string) error SaveControllerConfig(config any) error @@ -73,6 +75,11 @@ func (se *StorageEndpoint) LoadResourceGroupSettings(f func(keyspaceID uint32, n }) } +// LoadResourceGroupSetting loads a specific resource group from storage. +func (se *StorageEndpoint) LoadResourceGroupSetting(keyspaceID uint32, name string) (string, error) { + return se.Load(keypath.KeyspaceResourceGroupSettingPath(keyspaceID, name)) +} + // SaveResourceGroupStates stores a resource group to storage. func (se *StorageEndpoint) SaveResourceGroupStates(keyspaceID uint32, name string, obj any) error { return se.saveJSON(keypath.KeyspaceResourceGroupStatePath(keyspaceID, name), obj) @@ -102,6 +109,11 @@ func (se *StorageEndpoint) LoadResourceGroupStates(f func(keyspaceID uint32, nam }) } +// LoadResourceGroupState loads a specific resource group state from storage. +func (se *StorageEndpoint) LoadResourceGroupState(keyspaceID uint32, name string) (string, error) { + return se.Load(keypath.KeyspaceResourceGroupStatePath(keyspaceID, name)) +} + // SaveControllerConfig stores the resource controller config to storage. func (se *StorageEndpoint) SaveControllerConfig(config any) error { return se.saveJSON(keypath.ControllerConfigPath(), config) diff --git a/tests/integrations/mcs/resourcemanager/resource_manager_test.go b/tests/integrations/mcs/resourcemanager/resource_manager_test.go index 1808175f804..79ec16ebd59 100644 --- a/tests/integrations/mcs/resourcemanager/resource_manager_test.go +++ b/tests/integrations/mcs/resourcemanager/resource_manager_test.go @@ -193,6 +193,7 @@ func (suite *resourceManagerClientTestSuite) SetupSuite() { // Ensure RM service discovery has picked up the standalone endpoint before running tests. waitResourceManagerServiceURL(re, suite.client, true) } + waitAsyncLoadResourceGroups(re, suite.client) suite.initGroups = []*rmpb.ResourceGroup{ { @@ -271,6 +272,13 @@ func waitResourceManagerServiceURL(re *require.Assertions, cli pd.Client, wantNo }) } +func waitAsyncLoadResourceGroups(re *require.Assertions, cli pd.Client) { + testutil.Eventually(re, func() bool { + _, err := cli.ListResourceGroups(context.TODO()) + return err == nil + }, testutil.WithTickInterval(100*time.Millisecond)) +} + func TestSwitchModeDuringWorkload(t *testing.T) { for _, tc := range []struct { name string @@ -545,6 +553,7 @@ func (suite *resourceManagerClientTestSuite) resignAndWaitLeader(re *require.Ass newLeader := suite.cluster.GetServer(suite.cluster.WaitLeader()) re.NotNil(newLeader) waitLeaderServingClient(re, suite.client, newLeader.GetAddr()) + waitAsyncLoadResourceGroups(re, suite.client) } func (suite *resourceManagerClientTestSuite) TestWatchResourceGroup() { From 966d279550b3f8e5fb06693a45bea01e3353934b Mon Sep 17 00:00:00 2001 From: bufferflies <1045931706@qq.com> Date: Thu, 11 Jun 2026 04:42:29 +0200 Subject: [PATCH 02/21] tests: wait for resource group reload after restart Signed-off-by: bufferflies <1045931706@qq.com> --- errors.toml | 8 ++++---- pkg/errs/errno.go | 2 +- pkg/mcs/resourcemanager/server/manager.go | 7 ++++++- pkg/mcs/resourcemanager/server/metrics.go | 6 +++--- .../resourcemanager/resource_manager_test.go | 18 ++++++++++++------ 5 files changed, 26 insertions(+), 15 deletions(-) diff --git a/errors.toml b/errors.toml index 791560c9de1..43413f79a05 100644 --- a/errors.toml +++ b/errors.toml @@ -921,14 +921,14 @@ error = ''' keyspace not found with name: %s ''' -["PD:scatter:ErrEmptyRegion"] +["PD:resourcemanager:ErrResourceGroupsLoading"] error = ''' -empty region +resource groups are still being loaded, please try again later ''' -["PD:resourcemanager:ErrResourceGroupsLoading"] +["PD:scatter:ErrEmptyRegion"] error = ''' -resource groups are still being loaded, please try again later +empty region ''' ["PD:schedule:ErrCreateOperator"] diff --git a/pkg/errs/errno.go b/pkg/errs/errno.go index 415f51061a3..9ddb0b7b400 100644 --- a/pkg/errs/errno.go +++ b/pkg/errs/errno.go @@ -541,7 +541,7 @@ var ( ErrResourceGroupNotExists = errors.Normalize("the %s resource group does not exist", errors.RFCCodeText("PD:resourcemanager:ErrGroupNotExists")) ErrDeleteReservedGroup = errors.Normalize("cannot delete reserved group", errors.RFCCodeText("PD:resourcemanager:ErrDeleteReservedGroup")) ErrInvalidGroup = errors.Normalize("invalid group settings, please check %s", errors.RFCCodeText("PD:resourcemanager:ErrInvalidGroup")) - ErrResourceGroupsLoading = errors.Normalize("resource groups are still being loaded, please try again later", errors.RFCCodeText("PD:resourcemanager:ErrResourceGroupsLoading")) + ErrResourceGroupsLoading = errors.Normalize("resource groups are still being loaded, please try again later", errors.RFCCodeText("PD:resourcemanager:ErrResourceGroupsLoading")) ) // Microservice errors diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 606223a5b95..329aea39be8 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -555,7 +555,12 @@ func (m *Manager) loadResourceGroup(keyspaceID uint32, name string) (*ResourceGr return nil, err } state, err := m.storage.LoadResourceGroupState(keyspaceID, name) - if err == nil && state != "" { + if err != nil { + log.Warn("failed to load resource group state, continuing without state", + zap.Uint32("keyspace-id", keyspaceID), + zap.String("group-name", name), + zap.Error(err)) + } else if state != "" { if err := krgm.setRawStatesIntoResourceGroup(name, state); err != nil { return nil, err } diff --git a/pkg/mcs/resourcemanager/server/metrics.go b/pkg/mcs/resourcemanager/server/metrics.go index 47b507efc83..2fa373aa59b 100644 --- a/pkg/mcs/resourcemanager/server/metrics.go +++ b/pkg/mcs/resourcemanager/server/metrics.go @@ -225,8 +225,8 @@ var ( prometheus.CounterOpts{ Namespace: namespace, Subsystem: serverSubsystem, - Name: "sync_load_group_counter", - Help: "The number of the sync load group.", + Name: "sync_load_groups_total", + Help: "Total number of resource groups loaded synchronously.", }) asyncLoadGroupDuration = prometheus.NewHistogram( @@ -234,7 +234,7 @@ var ( Namespace: namespace, Subsystem: serverSubsystem, Name: "async_load_group_duration_seconds", - Help: "The duration of the async load group.", + Help: "Duration of asynchronous resource group loading in seconds.", }) ) diff --git a/tests/integrations/mcs/resourcemanager/resource_manager_test.go b/tests/integrations/mcs/resourcemanager/resource_manager_test.go index 79ec16ebd59..c3492aa50fd 100644 --- a/tests/integrations/mcs/resourcemanager/resource_manager_test.go +++ b/tests/integrations/mcs/resourcemanager/resource_manager_test.go @@ -21,6 +21,7 @@ import ( "io" "math/rand/v2" "net/http" + "reflect" "strconv" "strings" "sync/atomic" @@ -634,16 +635,21 @@ func (suite *resourceManagerClientTestSuite) TestWatchResourceGroup() { re.NoError(err) re.Contains(resp, "Success!") // Make sure the resource group active - meta, err = controller.GetResourceGroup(group.Name) - re.NotNil(meta) - re.NoError(err) + testutil.Eventually(re, func() bool { + meta, err = controller.GetResourceGroup(group.Name) + if err != nil || meta == nil { + return false + } + meta = controller.GetActiveResourceGroup(group.Name) + return meta != nil + }, testutil.WithTickInterval(50*time.Millisecond)) modifySettings(group, 30000) resp, err = cli.ModifyResourceGroup(suite.ctx, group) re.NoError(err) re.Contains(resp, "Success!") testutil.Eventually(re, func() bool { meta = controller.GetActiveResourceGroup(group.Name) - return meta.RUSettings.RU.Settings.FillRate == uint64(30000) + return meta != nil && meta.RUSettings.RU.Settings.FillRate == uint64(30000) }, testutil.WithTickInterval(100*time.Millisecond)) re.NoError(failpoint.Disable("github.com/tikv/pd/client/resource_group/controller/watchStreamError")) @@ -1500,8 +1506,8 @@ func (suite *resourceManagerClientTestSuite) TestBasicResourceGroupCURD() { testutil.Eventually(re, func() bool { var err error newGroups, err = cli.ListResourceGroups(suite.ctx) - return err == nil - }, testutil.WithWaitFor(time.Second)) + return err == nil && reflect.DeepEqual(groups, newGroups) + }) re.Equal(groups, newGroups) } From fae9a3ef251824a5143bdef36818f8798bb6ea06 Mon Sep 17 00:00:00 2001 From: bufferflies <1045931706@qq.com> Date: Thu, 11 Jun 2026 05:21:56 +0200 Subject: [PATCH 03/21] resource_group: avoid overwriting default during async load Signed-off-by: bufferflies <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/manager.go | 11 ++++++- .../resourcemanager/resource_manager_test.go | 31 +++++++++++++++++-- 2 files changed, 39 insertions(+), 3 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 329aea39be8..3a06c6fe912 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -403,7 +403,7 @@ func (m *Manager) initMetadata(ctx context.Context) error { atomic.StoreInt32(&m.loadingState, LoadingStateNotStarted) m.Unlock() - m.initReserved() + m.initReservedInCache() if err := m.loadServiceLimits(); err != nil { return err } @@ -703,6 +703,15 @@ func (m *Manager) initReserved() { } } +func (m *Manager) initReservedInCache() { + // Initialize the reserved default group in memory before async loading + // without overwriting persisted default group settings. + m.getOrCreateKeyspaceResourceGroupManager(constant.NullKeyspaceID, false).ensureReservedDefaultGroupInCache() + for _, krgm := range m.getKeyspaceResourceGroupManagers() { + krgm.ensureReservedDefaultGroupInCache() + } +} + // UpdateControllerConfigItem updates the controller config item. func (m *Manager) UpdateControllerConfigItem(key string, value any) error { if !m.writeRole.AllowsMetadataWrite() { diff --git a/tests/integrations/mcs/resourcemanager/resource_manager_test.go b/tests/integrations/mcs/resourcemanager/resource_manager_test.go index c3492aa50fd..e03dd2bcc23 100644 --- a/tests/integrations/mcs/resourcemanager/resource_manager_test.go +++ b/tests/integrations/mcs/resourcemanager/resource_manager_test.go @@ -22,12 +22,14 @@ import ( "math/rand/v2" "net/http" "reflect" + "sort" "strconv" "strings" "sync/atomic" "testing" "time" + "github.com/gogo/protobuf/proto" "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" "go.uber.org/goleak" @@ -1502,13 +1504,38 @@ func (suite *resourceManagerClientTestSuite) TestBasicResourceGroupCURD() { // re-connect client as well suite.client = suite.setupPDClient(re) cli = suite.client + expectedGroups := normalizeResourceGroupsForSettingsCompare(groups) var newGroups []*rmpb.ResourceGroup testutil.Eventually(re, func() bool { var err error newGroups, err = cli.ListResourceGroups(suite.ctx) - return err == nil && reflect.DeepEqual(groups, newGroups) + return err == nil && reflect.DeepEqual(expectedGroups, normalizeResourceGroupsForSettingsCompare(newGroups)) }) - re.Equal(groups, newGroups) + re.Equal(expectedGroups, normalizeResourceGroupsForSettingsCompare(newGroups)) +} + +func normalizeResourceGroupsForSettingsCompare(groups []*rmpb.ResourceGroup) []*rmpb.ResourceGroup { + normalized := make([]*rmpb.ResourceGroup, 0, len(groups)) + for _, group := range groups { + cloned := proto.Clone(group).(*rmpb.ResourceGroup) + cloned.RUStats = nil + resetTokenBucketRuntimeState(cloned.GetRUSettings().GetRU()) + rawSettings := cloned.GetRawResourceSettings() + resetTokenBucketRuntimeState(rawSettings.GetCpu()) + resetTokenBucketRuntimeState(rawSettings.GetIoRead()) + resetTokenBucketRuntimeState(rawSettings.GetIoWrite()) + normalized = append(normalized, cloned) + } + sort.Slice(normalized, func(i, j int) bool { + return normalized[i].GetName() < normalized[j].GetName() + }) + return normalized +} + +func resetTokenBucketRuntimeState(bucket *rmpb.TokenBucket) { + if bucket != nil { + bucket.Tokens = 0 + } } func (suite *resourceManagerClientTestSuite) TestResourceGroupRUConsumption() { From f7cf2cdcaf13aad4af4b49060f9edbd27a557642 Mon Sep 17 00:00:00 2001 From: bufferflies <1045931706@qq.com> Date: Thu, 11 Jun 2026 06:01:05 +0200 Subject: [PATCH 04/21] tests: avoid direct proto dependency in RM integration Signed-off-by: bufferflies <1045931706@qq.com> --- .../mcs/resourcemanager/resource_manager_test.go | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/tests/integrations/mcs/resourcemanager/resource_manager_test.go b/tests/integrations/mcs/resourcemanager/resource_manager_test.go index e03dd2bcc23..675c82c2a4c 100644 --- a/tests/integrations/mcs/resourcemanager/resource_manager_test.go +++ b/tests/integrations/mcs/resourcemanager/resource_manager_test.go @@ -29,7 +29,6 @@ import ( "testing" "time" - "github.com/gogo/protobuf/proto" "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" "go.uber.org/goleak" @@ -1517,7 +1516,9 @@ func (suite *resourceManagerClientTestSuite) TestBasicResourceGroupCURD() { func normalizeResourceGroupsForSettingsCompare(groups []*rmpb.ResourceGroup) []*rmpb.ResourceGroup { normalized := make([]*rmpb.ResourceGroup, 0, len(groups)) for _, group := range groups { - cloned := proto.Clone(group).(*rmpb.ResourceGroup) + cloned := typeutil.DeepClone(group, func() *rmpb.ResourceGroup { + return &rmpb.ResourceGroup{} + }) cloned.RUStats = nil resetTokenBucketRuntimeState(cloned.GetRUSettings().GetRU()) rawSettings := cloned.GetRawResourceSettings() From a20e813777d94410189363e90fdbd046b8b41777 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Wed, 8 Jul 2026 09:50:36 +0800 Subject: [PATCH 05/21] resource_group: fix lazy-load ordering and default-group review findings Address PR #10873 review comments: - grpc_service: AcquireTokenBuckets checked accessKeyspaceResourceGroupManager before GetMutableResourceGroup, bypassing the lazy-load path for non-default groups whose keyspace manager wasn't in memory yet. Reorder so the lazy load runs first. - manager: loadResourceGroupIfNeeded synthesized the reserved default group unconditionally, which could persist synthetic defaults over a customized one still on disk. Try the storage load first and only fall back to the synthetic default on an explicit not-found error. - manager: AddResourceGroup silently swallowed all loadResourceGroupIfNeeded errors, including transient storage failures. Only ignore the expected not-found case and propagate others. - manager_async_test: blockingResourceGroupStorage.unblock could hang stopAsyncTestManager's wg.Wait() if a test aborted before calling it. Make unblock idempotent and defer it right after Init so teardown can't deadlock. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- .../resourcemanager/server/grpc_service.go | 14 ++++++++------ pkg/mcs/resourcemanager/server/manager.go | 16 ++++++++++------ .../server/manager_async_test.go | 19 +++++++++++++++---- 3 files changed, 33 insertions(+), 16 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/grpc_service.go b/pkg/mcs/resourcemanager/server/grpc_service.go index 875299c6025..bbb73b2340d 100644 --- a/pkg/mcs/resourcemanager/server/grpc_service.go +++ b/pkg/mcs/resourcemanager/server/grpc_service.go @@ -231,18 +231,20 @@ func (s *Service) AcquireTokenBuckets(stream rmpb.ResourceManager_AcquireTokenBu zap.Uint32("keyspace-id", keyspaceID), zap.String("resource-group", resourceGroupName), ) + // Get the resource group from manager to acquire token buckets. This also + // triggers lazy loading of the group if async loading hasn't completed yet, + // so it must happen before accessKeyspaceResourceGroupManager below. + rg, err := s.manager.GetMutableResourceGroup(keyspaceID, resourceGroupName) + if rg == nil { + log.Warn("resource group not found", append(requestFields, zap.Error(err))...) + continue + } // Get keyspace resource group manager to apply service limit later. krgm, err := s.manager.accessKeyspaceResourceGroupManager(keyspaceID, resourceGroupName) if krgm == nil { log.Warn("keyspace resource group manager not found", append(requestFields, zap.Error(err))...) continue } - // Get the resource group from manager to acquire token buckets. - rg, err := s.manager.GetMutableResourceGroup(keyspaceID, resourceGroupName) - if rg == nil { - log.Warn("resource group not found", append(requestFields, zap.Error(err))...) - continue - } // Send the consumption to update the metrics. err = s.manager.dispatchConsumption(req) if err != nil { diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 3a06c6fe912..b83c89902a2 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -578,12 +578,14 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro return nil } } - if name == DefaultResourceGroupName { - m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, true) - return nil - } group, err := m.loadResourceGroup(keyspaceID, name) if err != nil { + if name == DefaultResourceGroupName && errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(name)) { + // No persisted default group settings exist yet (e.g. a brand-new + // keyspace), so it's safe to synthesize the reserved default group. + m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, true) + return nil + } return err } krgm = m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false) @@ -782,8 +784,10 @@ func (m *Manager) AddResourceGroup(grouppb *rmpb.ResourceGroup) error { if krgm == nil { return errs.ErrKeyspaceNotExists.FastGenByArgs(keyspaceID) } - if err := m.loadResourceGroupIfNeeded(keyspaceID, grouppb.Name); err != nil { - log.Debug("failed to load resource group", zap.Uint32("keyspace-id", keyspaceID), zap.String("name", grouppb.Name), zap.Error(err)) + if err := m.loadResourceGroupIfNeeded(keyspaceID, grouppb.Name); err != nil && + !errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(grouppb.Name)) { + log.Warn("failed to load resource group before add", zap.Uint32("keyspace-id", keyspaceID), zap.String("name", grouppb.Name), zap.Error(err)) + return err } if err := krgm.addResourceGroup(grouppb); err != nil { return err diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go index 7bf7d51ef1e..66daf955803 100644 --- a/pkg/mcs/resourcemanager/server/manager_async_test.go +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -32,9 +32,10 @@ import ( type blockingResourceGroupStorage struct { storage.Storage - once sync.Once - entered chan struct{} - release chan struct{} + once sync.Once + releaseOnce sync.Once + entered chan struct{} + release chan struct{} } func newBlockingResourceGroupStorage() *blockingResourceGroupStorage { @@ -63,7 +64,9 @@ func (s *blockingResourceGroupStorage) waitEntered(t *testing.T) { } func (s *blockingResourceGroupStorage) unblock() { - close(s.release) + s.releaseOnce.Do(func() { + close(s.release) + }) } func newAsyncTestGroup(name string, fillRate uint64) *resource_manager.ResourceGroup { @@ -98,6 +101,10 @@ func TestAsyncLoadResourceGroupsLazyGet(t *testing.T) { m.storage = store re.NoError(m.Init(context.Background())) defer stopAsyncTestManager(m) + // Unblock the async loader first (LIFO) so stopAsyncTestManager's wg.Wait() + // cannot hang if a later assertion aborts the test before the explicit + // store.unblock() call below is reached. + defer store.unblock() store.waitEntered(t) @@ -126,6 +133,10 @@ func TestAsyncLoadResourceGroupsDoesNotRestoreDeletedLazyGroup(t *testing.T) { m.storage = store re.NoError(m.Init(context.Background())) defer stopAsyncTestManager(m) + // Unblock the async loader first (LIFO) so stopAsyncTestManager's wg.Wait() + // cannot hang if a later assertion aborts the test before the explicit + // store.unblock() call below is reached. + defer store.unblock() store.waitEntered(t) From cbe73bcf816b860d4240a6819ccccb695cb74e71 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Wed, 8 Jul 2026 15:13:19 +0800 Subject: [PATCH 06/21] resource_group: propagate non-not-found errors in AcquireTokenBuckets GetMutableResourceGroup can return (nil, err) for real failures (e.g. ErrKeyspaceNotExists, storage errors during lazy load), not just a missing group. Treating every rg == nil as "not found" silently dropped token requests on real errors instead of failing the stream. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/grpc_service.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/pkg/mcs/resourcemanager/server/grpc_service.go b/pkg/mcs/resourcemanager/server/grpc_service.go index bbb73b2340d..0945494cdc5 100644 --- a/pkg/mcs/resourcemanager/server/grpc_service.go +++ b/pkg/mcs/resourcemanager/server/grpc_service.go @@ -236,6 +236,9 @@ func (s *Service) AcquireTokenBuckets(stream rmpb.ResourceManager_AcquireTokenBu // so it must happen before accessKeyspaceResourceGroupManager below. rg, err := s.manager.GetMutableResourceGroup(keyspaceID, resourceGroupName) if rg == nil { + if err != nil && !errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(resourceGroupName)) { + return err + } log.Warn("resource group not found", append(requestFields, zap.Error(err))...) continue } From 25ec642094e0a5edd17dd85048a2d5a77ed9109c Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Mon, 13 Jul 2026 18:40:07 +0800 Subject: [PATCH 07/21] resource_group: don't mark lazily-loaded group synced on state load failure When loadResourceGroup's LoadResourceGroupState read fails, it logs a warning and still returns the group using default state. But loadResourceGroupIfNeeded then unconditionally marked the group as sync-loaded, which made asyncLoadResourceGroups's merge permanently skip it, even if the later bulk load successfully read the real persisted state. The group would keep its default/fresh state (e.g. a re-initialized token bucket) for the rest of the manager's lifetime. Only mark the group sync-loaded when its state was actually read, so the async bulk merge remains free to fill in the correct state afterward if nothing else has modified the group in the meantime. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/manager.go | 37 +++++++++++++++-------- 1 file changed, 24 insertions(+), 13 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index b83c89902a2..c2c71e131c7 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -542,17 +542,21 @@ func (m *Manager) loadKeyspaceResourceGroupsFromStorage() (map[uint32]*keyspaceR return tempKrgms, nil } -func (m *Manager) loadResourceGroup(keyspaceID uint32, name string) (*ResourceGroup, error) { +// loadResourceGroup loads a single resource group from storage. The returned +// stateLoaded reports whether the group's persisted state was successfully +// read; the caller must not mark such a group as sync-loaded, so that a +// concurrent or later async bulk load can still fill in its real state. +func (m *Manager) loadResourceGroup(keyspaceID uint32, name string) (group *ResourceGroup, stateLoaded bool, err error) { rawValue, err := m.storage.LoadResourceGroupSetting(keyspaceID, name) if err != nil { - return nil, err + return nil, false, err } if rawValue == "" { - return nil, errs.ErrResourceGroupNotExists.FastGenByArgs(name) + return nil, false, errs.ErrResourceGroupNotExists.FastGenByArgs(name) } krgm := newKeyspaceResourceGroupManager(keyspaceID, m.storage, m.writeRole) if err := krgm.addResourceGroupFromRaw(name, rawValue); err != nil { - return nil, err + return nil, false, err } state, err := m.storage.LoadResourceGroupState(keyspaceID, name) if err != nil { @@ -560,12 +564,14 @@ func (m *Manager) loadResourceGroup(keyspaceID uint32, name string) (*ResourceGr zap.Uint32("keyspace-id", keyspaceID), zap.String("group-name", name), zap.Error(err)) - } else if state != "" { + return krgm.getMutableResourceGroup(name), false, nil + } + if state != "" { if err := krgm.setRawStatesIntoResourceGroup(name, state); err != nil { - return nil, err + return nil, false, err } } - return krgm.getMutableResourceGroup(name), nil + return krgm.getMutableResourceGroup(name), true, nil } func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) error { @@ -578,7 +584,7 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro return nil } } - group, err := m.loadResourceGroup(keyspaceID, name) + group, stateLoaded, err := m.loadResourceGroup(keyspaceID, name) if err != nil { if name == DefaultResourceGroupName && errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(name)) { // No persisted default group settings exist yet (e.g. a brand-new @@ -599,12 +605,17 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro if inserted { krgm.syncBurstabilityWithServiceLimit(group) } - markKey := trackerKey{keyspaceID: keyspaceID, groupName: name} - m.Lock() - if m.syncLoadedGroups != nil { - m.syncLoadedGroups[markKey] = true + // Only mark the group as sync-loaded when its persisted state was actually + // read; otherwise a later async bulk load must remain free to fill in the + // real state instead of being skipped forever. + if stateLoaded { + markKey := trackerKey{keyspaceID: keyspaceID, groupName: name} + m.Lock() + if m.syncLoadedGroups != nil { + m.syncLoadedGroups[markKey] = true + } + m.Unlock() } - m.Unlock() syncLoadGroupCounter.Inc() return nil } From 306fd57c17bf04fbac9881038ecb1f07d85d8c67 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Tue, 14 Jul 2026 10:26:28 +0800 Subject: [PATCH 08/21] resource_group: distinguish reserved default placeholder from confirmed data initReservedInCache installs a synthetic default resource group into the cache before async loading starts. loadResourceGroupIfNeeded's early cache-hit check treated any cached entry as already-loaded, so it never attempted the persisted point load for a keyspace's default group as long as the synthetic placeholder occupied the slot. A customized default group would be served with built-in settings for the whole startup window, and the async bulk merge could be blocked from correcting it if a concurrent write marked it sync-loaded first. Track which cache entries are still just placeholders (keyspaceResourceGroupManager.reservedGroups, guarded by the same lock as groups) and clear the mark on any real write. loadResourceGroupIfNeeded now only treats a cached entry as satisfying the call when it isn't a placeholder, and is allowed to replace a placeholder with the result of a real storage load. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- .../server/keyspace_manager.go | 36 +++++++++++++++++-- pkg/mcs/resourcemanager/server/manager.go | 13 +++++-- 2 files changed, 45 insertions(+), 4 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/keyspace_manager.go b/pkg/mcs/resourcemanager/server/keyspace_manager.go index d3d8a13daec..8aabecedded 100644 --- a/pkg/mcs/resourcemanager/server/keyspace_manager.go +++ b/pkg/mcs/resourcemanager/server/keyspace_manager.go @@ -69,7 +69,13 @@ type consumptionItem struct { type keyspaceResourceGroupManager struct { syncutil.RWMutex - groups map[string]*ResourceGroup + groups map[string]*ResourceGroup + // reservedGroups tracks names whose entry in groups is still just the + // synthetic placeholder inserted by ensureReservedDefaultGroupInCache or + // restoreDefaultResourceGroupFromReserved, not yet confirmed by a storage + // load or a real write. It shares the same lock as groups so a lazy load + // can atomically decide whether it's safe to replace the placeholder. + reservedGroups map[string]struct{} groupRUTrackers map[string]*groupRUTracker serviceLimiter *serviceLimiter @@ -89,6 +95,7 @@ func newKeyspaceResourceGroupManager( } return &keyspaceResourceGroupManager{ groups: make(map[string]*ResourceGroup), + reservedGroups: make(map[string]struct{}), groupRUTrackers: make(map[string]*groupRUTracker), keyspaceID: keyspaceID, storage: storage, @@ -159,6 +166,9 @@ func (krgm *keyspaceResourceGroupManager) upsertResourceGroupFromRaw(name string zap.Uint32("keyspace-id", krgm.keyspaceID), zap.String("name", name), zap.String("raw-value", rawValue), zap.Error(err)) return err } + krgm.Lock() + delete(krgm.reservedGroups, group.Name) + krgm.Unlock() krgm.syncBurstabilityWithServiceLimit(existing) return nil } @@ -166,6 +176,7 @@ func (krgm *keyspaceResourceGroupManager) upsertResourceGroupFromRaw(name string resourceGroup := FromProtoResourceGroup(group) krgm.Lock() krgm.groups[group.Name] = resourceGroup + delete(krgm.reservedGroups, group.Name) krgm.Unlock() krgm.syncBurstabilityWithServiceLimit(resourceGroup) return nil @@ -175,6 +186,7 @@ func (krgm *keyspaceResourceGroupManager) deleteResourceGroupFromCache(name stri krgm.Lock() delete(krgm.groups, name) delete(krgm.groupRUTrackers, name) + delete(krgm.reservedGroups, name) krgm.Unlock() } @@ -218,6 +230,7 @@ func (krgm *keyspaceResourceGroupManager) ensureReservedDefaultGroupInCache() { krgm.Lock() if _, ok := krgm.groups[DefaultResourceGroupName]; !ok { krgm.groups[DefaultResourceGroupName] = defaultGroup + krgm.reservedGroups[DefaultResourceGroupName] = struct{}{} inserted = true } krgm.Unlock() @@ -245,6 +258,7 @@ func (krgm *keyspaceResourceGroupManager) restoreDefaultResourceGroupFromReserve defaultGroup := newDefaultResourceGroup() krgm.Lock() krgm.groups[DefaultResourceGroupName] = defaultGroup + krgm.reservedGroups[DefaultResourceGroupName] = struct{}{} krgm.Unlock() krgm.syncBurstabilityWithServiceLimit(defaultGroup) } @@ -266,6 +280,7 @@ func (krgm *keyspaceResourceGroupManager) addResourceGroup(grouppb *rmpb.Resourc } krgm.Lock() krgm.groups[group.Name] = group + delete(krgm.reservedGroups, group.Name) krgm.Unlock() krgm.syncBurstabilityWithServiceLimit(group) return nil @@ -288,7 +303,13 @@ func (krgm *keyspaceResourceGroupManager) modifyResourceGroup(group *rmpb.Resour if err != nil { return err } - return curGroup.persistSettings(krgm.keyspaceID, krgm.storage) + if err := curGroup.persistSettings(krgm.keyspaceID, krgm.storage); err != nil { + return err + } + krgm.Lock() + delete(krgm.reservedGroups, group.Name) + krgm.Unlock() + return nil } func (krgm *keyspaceResourceGroupManager) deleteResourceGroup(name string) error { @@ -330,6 +351,17 @@ func (krgm *keyspaceResourceGroupManager) deleteResourceGroup(name string) error return nil } +// isReserved reports whether name's cached entry is still just the synthetic +// placeholder inserted by ensureReservedDefaultGroupInCache or +// restoreDefaultResourceGroupFromReserved, not yet confirmed by a storage +// load or a real write. +func (krgm *keyspaceResourceGroupManager) isReserved(name string) bool { + krgm.RLock() + defer krgm.RUnlock() + _, ok := krgm.reservedGroups[name] + return ok +} + func (krgm *keyspaceResourceGroupManager) getResourceGroup(name string, withStats bool) *ResourceGroup { krgm.RLock() defer krgm.RUnlock() diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index c2c71e131c7..8f465019196 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -580,7 +580,10 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro } krgm := m.getKeyspaceResourceGroupManager(keyspaceID) if krgm != nil { - if group := krgm.getMutableResourceGroup(name); group != nil { + // A cached entry only satisfies this call if it's confirmed data, not + // just a synthetic placeholder (e.g. from ensureReservedDefaultGroupInCache) + // installed before async loading had a chance to run. + if group := krgm.getMutableResourceGroup(name); group != nil && !krgm.isReserved(name) { return nil } } @@ -588,7 +591,7 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro if err != nil { if name == DefaultResourceGroupName && errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(name)) { // No persisted default group settings exist yet (e.g. a brand-new - // keyspace), so it's safe to synthesize the reserved default group. + // keyspace), so it's safe to keep serving the reserved placeholder. m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, true) return nil } @@ -600,7 +603,13 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro if _, exists := krgm.groups[name]; !exists { krgm.groups[name] = group inserted = true + } else if _, reserved := krgm.reservedGroups[name]; reserved { + // The existing entry is just an unconfirmed placeholder; the freshly + // loaded group is the real, confirmed data, so replacing it is safe. + krgm.groups[name] = group + inserted = true } + delete(krgm.reservedGroups, name) krgm.Unlock() if inserted { krgm.syncBurstabilityWithServiceLimit(group) From 042d6a77b12d5b6bd6111e6bba28c424d0bf0677 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Tue, 14 Jul 2026 15:24:34 +0800 Subject: [PATCH 09/21] resource_group: keep placeholder state out of the persist loop and reserved marker persistResourceGroupRunningState persisted every entry in krgm.groups regardless of reservedGroups, so the persist loop (running concurrently with async loading) could write a synthetic default's fresh token state back to storage before the real persisted state was ever loaded, permanently overwriting it. Skip reserved entries there. Also fix an oversight in the previous placeholder-tracking commit: loadResourceGroupIfNeeded cleared a group's reserved mark unconditionally after a lazy load, even when LoadResourceGroupState had failed and the group only carries default state. Only clear the mark when the state was actually read, otherwise keep it reserved so a later call or the async bulk merge can still fill in the real state instead of the partial result being treated as final. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/keyspace_manager.go | 8 ++++++++ pkg/mcs/resourcemanager/server/manager.go | 10 +++++++++- 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/pkg/mcs/resourcemanager/server/keyspace_manager.go b/pkg/mcs/resourcemanager/server/keyspace_manager.go index 8aabecedded..0396f556b3b 100644 --- a/pkg/mcs/resourcemanager/server/keyspace_manager.go +++ b/pkg/mcs/resourcemanager/server/keyspace_manager.go @@ -416,6 +416,14 @@ func (krgm *keyspaceResourceGroupManager) persistResourceGroupRunningState() { for idx := range keys { krgm.RLock() group, ok := krgm.groups[keys[idx]] + _, reserved := krgm.reservedGroups[keys[idx]] + if ok && reserved { + // The entry is still just an unconfirmed placeholder (e.g. the + // synthetic default installed before async loading completes); + // persisting its fresh state would permanently overwrite any + // real persisted state still waiting to be loaded. + ok = false + } if ok { if err := group.persistStates(krgm.keyspaceID, krgm.storage); err != nil { log.Error("persist keyspace resource group state failed", diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 8f465019196..940135ea044 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -609,7 +609,15 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro krgm.groups[name] = group inserted = true } - delete(krgm.reservedGroups, name) + // Only clear the placeholder mark once the state was actually read; a + // metadata-only group (state load failed) must stay reserved so a later + // call or the async bulk merge remains free to fill in the real state, + // instead of this partial result being treated as final forever. + if stateLoaded { + delete(krgm.reservedGroups, name) + } else { + krgm.reservedGroups[name] = struct{}{} + } krgm.Unlock() if inserted { krgm.syncBurstabilityWithServiceLimit(group) From d2cddc91a5d342bfe762cfa7494bab29c6662cbf Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Wed, 15 Jul 2026 10:58:51 +0800 Subject: [PATCH 10/21] resource_group: route default-group creation through the safe load path getOrCreateKeyspaceResourceGroupManager(id, true) is used by both AddResourceGroup and SetKeyspaceServiceLimit to make sure a keyspace's default resource group exists. It called initDefaultResourceGroup directly, which persists a synthetic default whenever one isn't cached yet, with no attempt to check storage first. For a keyspace touched for the first time while async loading is still in progress, this could overwrite a customized default group's stored settings before the real data had a chance to load - the same class of bug fixed earlier in loadResourceGroupIfNeeded, reachable via two more entry points. Make initDefault=true go through loadResourceGroupIfNeeded (storage point load first, synthesize only on confirmed not-found) while async loading is in progress. Once loading has completed, any default still missing from the cache is confirmed absent from storage, so it's synthesized directly as before. loadResourceGroupIfNeeded's own not-found fallback now calls initDefaultResourceGroup directly instead of going through getOrCreateKeyspaceResourceGroupManager(id, true), to avoid recursing back into itself. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/manager.go | 24 +++++++++++++++++++---- 1 file changed, 20 insertions(+), 4 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 940135ea044..a2c59bfd21c 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -295,6 +295,13 @@ func (m *Manager) GetRUVersionPolicy() *RUVersionPolicy { return m.controllerConfig.RUVersionPolicy.Clone() } +// getOrCreateKeyspaceResourceGroupManager returns the keyspace resource group +// manager for keyspaceID, creating it if needed. When initDefault is true, it +// also ensures the default resource group is present. While async loading is +// still in progress this goes through loadResourceGroupIfNeeded, which tries +// a storage point load first so a customized default is never clobbered by a +// blindly synthesized one. Once loading has completed, any default missing +// from the cache truly doesn't exist anywhere, so it's synthesized directly. func (m *Manager) getOrCreateKeyspaceResourceGroupManager(keyspaceID uint32, initDefault bool) *keyspaceResourceGroupManager { m.Lock() krgm, ok := m.krgms[keyspaceID] @@ -303,9 +310,15 @@ func (m *Manager) getOrCreateKeyspaceResourceGroupManager(keyspaceID uint32, ini m.krgms[keyspaceID] = krgm } m.Unlock() - // Init the default resource group if needed. if initDefault { - krgm.initDefaultResourceGroup() + if atomic.LoadInt32(&m.loadingState) == LoadingStateCompleted { + // Async loading (if any) has already finished, so if the default + // group isn't cached yet it truly doesn't exist anywhere; it's + // safe to synthesize and persist it directly. + krgm.initDefaultResourceGroup() + } else if err := m.loadResourceGroupIfNeeded(keyspaceID, DefaultResourceGroupName); err != nil { + log.Debug("failed to load default resource group", zap.Uint32("keyspace-id", keyspaceID), zap.Error(err)) + } } return krgm } @@ -591,8 +604,11 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro if err != nil { if name == DefaultResourceGroupName && errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(name)) { // No persisted default group settings exist yet (e.g. a brand-new - // keyspace), so it's safe to keep serving the reserved placeholder. - m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, true) + // keyspace), so it's safe to synthesize the reserved default group. + // This calls initDefaultResourceGroup directly instead of going + // through getOrCreateKeyspaceResourceGroupManager(id, true), which + // now routes back into this same function and would recurse. + m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false).initDefaultResourceGroup() return nil } return err From 8f743d909575deddbd22927a69c690b2eefeb99a Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Wed, 15 Jul 2026 11:02:53 +0800 Subject: [PATCH 11/21] resource_group: don't let Modify confirm a state-unconfirmed group krgm.modifyResourceGroup unconditionally cleared reservedGroups after patching settings, and Manager.ModifyResourceGroup unconditionally called markResourceGroupSyncLoaded afterwards. Modify only patches settings, it never loads or establishes a group's state, so if the preceding lazy load had failed to read the persisted state (stateLoaded=false), a Modify call would incorrectly promote the entry to fully confirmed at both the keyspace-manager level (reservedGroups) and the manager level (syncLoadedGroups). The async bulk merge would then skip it forever, so the real persisted running state would never get applied and the cached token bucket could remain initialized with empty state indefinitely. Leave reservedGroups untouched in modifyResourceGroup, and only mark the group sync-loaded in ModifyResourceGroup when it isn't still reserved. This lets the existing retry-on-next-access path (and the async bulk merge) keep trying to fill in the real state instead of the settings-only patch being treated as a full confirmation. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/keyspace_manager.go | 11 ++++------- pkg/mcs/resourcemanager/server/manager.go | 8 +++++++- 2 files changed, 11 insertions(+), 8 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/keyspace_manager.go b/pkg/mcs/resourcemanager/server/keyspace_manager.go index 0396f556b3b..8c900c54796 100644 --- a/pkg/mcs/resourcemanager/server/keyspace_manager.go +++ b/pkg/mcs/resourcemanager/server/keyspace_manager.go @@ -303,13 +303,10 @@ func (krgm *keyspaceResourceGroupManager) modifyResourceGroup(group *rmpb.Resour if err != nil { return err } - if err := curGroup.persistSettings(krgm.keyspaceID, krgm.storage); err != nil { - return err - } - krgm.Lock() - delete(krgm.reservedGroups, group.Name) - krgm.Unlock() - return nil + // Deliberately not clearing reservedGroups here: modifying only patches + // settings, it never establishes the group's state, so it must not make + // a state-unconfirmed entry look fully confirmed. + return curGroup.persistSettings(krgm.keyspaceID, krgm.storage) } func (krgm *keyspaceResourceGroupManager) deleteResourceGroup(name string) error { diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index a2c59bfd21c..01636ee2deb 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -857,7 +857,13 @@ func (m *Manager) ModifyResourceGroup(grouppb *rmpb.ResourceGroup) error { if err := krgm.modifyResourceGroup(grouppb); err != nil { return err } - m.markResourceGroupSyncLoaded(keyspaceID, grouppb.Name) + // Modifying only patches settings, it never establishes the group's + // state. If the state still hasn't been confirmed (isReserved), marking + // it sync-loaded here would make the async bulk merge skip it forever, + // so the persisted running state would never get applied. + if !krgm.isReserved(grouppb.Name) { + m.markResourceGroupSyncLoaded(keyspaceID, grouppb.Name) + } return nil } From eb3e774035369f01a0f2b7f5d63858338db11141 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Wed, 15 Jul 2026 16:12:50 +0800 Subject: [PATCH 12/21] resource_group: add lazy-load coverage for legacy null-keyspace groups Review feedback questioned whether the point loaders (LoadResourceGroupSetting/LoadResourceGroupState) miss the legacy, pre-keyspace resource_group/* path that the bulk loaders fall back to for constant.NullKeyspaceID. They don't: KeyspaceResourceGroupSettingPath and KeyspaceResourceGroupStatePath already resolve to the legacy path for NullKeyspaceID, the same helper both point and bulk loaders use. Add a regression test exercising the actual lazy-load path (GetResourceGroup during blocked async loading) against a legacy group saved under NullKeyspaceID, to make this guarantee explicit and catch any future divergence between the point and bulk loaders. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- .../server/manager_async_test.go | 32 +++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go index 66daf955803..af67e55112f 100644 --- a/pkg/mcs/resourcemanager/server/manager_async_test.go +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -25,6 +25,7 @@ import ( "github.com/pingcap/kvproto/pkg/resource_manager" "github.com/tikv/pd/pkg/errs" + "github.com/tikv/pd/pkg/keyspace/constant" "github.com/tikv/pd/pkg/storage" "github.com/tikv/pd/pkg/utils/testutil" ) @@ -159,3 +160,34 @@ func TestAsyncLoadResourceGroupsDoesNotRestoreDeletedLazyGroup(t *testing.T) { return true }, testutil.WithTickInterval(20*time.Millisecond)) } + +// TestAsyncLoadResourceGroupsLazyGetLegacyKeyspace guards against the point +// loaders (LoadResourceGroupSetting/LoadResourceGroupState) diverging from +// the bulk loaders on legacy, pre-keyspace resource groups: those are saved +// under constant.NullKeyspaceID, and a lazy Get during async loading must be +// able to find one the same way the bulk scan would once it completes. +func TestAsyncLoadResourceGroupsLazyGetLegacyKeyspace(t *testing.T) { + re := require.New(t) + store := newBlockingResourceGroupStorage() + re.NoError(store.SaveResourceGroupSetting(constant.NullKeyspaceID, "legacy-group", newAsyncTestGroup("legacy-group", 100))) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + re.NoError(m.Init(context.Background())) + defer stopAsyncTestManager(m) + defer store.unblock() + + store.waitEntered(t) + + group, err := m.GetResourceGroup(constant.NullKeyspaceID, "legacy-group", false) + re.NoError(err) + re.NotNil(group) + re.Equal("legacy-group", group.Name) + re.Equal(float64(100), group.RUSettings.RU.getFillRate()) + + store.unblock() + testutil.Eventually(re, func() bool { + group, err := m.GetResourceGroup(constant.NullKeyspaceID, "legacy-group", false) + return err == nil && group != nil + }, testutil.WithTickInterval(20*time.Millisecond)) +} From a32358a51102619087950904da20de7816fcccc0 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Thu, 16 Jul 2026 15:35:11 +0800 Subject: [PATCH 13/21] resource_group: clear reserved marker when async merge installs confirmed data asyncLoadResourceGroups's merge loop installed the fully-loaded (settings and state) confirmed group into krgm.groups but never cleared reservedGroups for it. A group that got marked reserved by an earlier lazy-load state-read failure would stay reserved forever even after the bulk load correctly recovered it, permanently excluding it from the state persist loop and from loadResourceGroupIfNeeded's fast path. Clear the reserved marker in the same merge step that installs the confirmed data. Add a regression test that injects a one-time LoadResourceGroupState failure during a lazy Get, then lets the async bulk load complete, and asserts the reserved marker is cleared once the confirmed data lands; verified the test fails without the fix. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/manager.go | 5 ++ .../server/manager_async_test.go | 54 +++++++++++++++++++ 2 files changed, 59 insertions(+) diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 01636ee2deb..e45be2cad34 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -496,6 +496,11 @@ func (m *Manager) asyncLoadResourceGroups(ctx context.Context) { key := trackerKey{keyspaceID: keyspaceID, groupName: name} if !m.syncLoadedGroups[key] { krgm.groups[name] = group + // This group is now confirmed, fully-loaded data (settings + // and state); it must no longer be treated as an + // unconfirmed placeholder by loadResourceGroupIfNeeded or + // skipped by the state persist loop. + delete(krgm.reservedGroups, name) groupsToSync = append(groupsToSync, group) loaded++ } diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go index af67e55112f..73db769b883 100644 --- a/pkg/mcs/resourcemanager/server/manager_async_test.go +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -16,7 +16,9 @@ package server import ( "context" + "errors" "sync" + "sync/atomic" "testing" "time" @@ -37,6 +39,10 @@ type blockingResourceGroupStorage struct { releaseOnce sync.Once entered chan struct{} release chan struct{} + + // failNextState, when true, makes the very next LoadResourceGroupState + // call fail once, then resets itself. + failNextState atomic.Bool } func newBlockingResourceGroupStorage() *blockingResourceGroupStorage { @@ -55,6 +61,13 @@ func (s *blockingResourceGroupStorage) LoadResourceGroupSettings(f func(keyspace return s.Storage.LoadResourceGroupSettings(f) } +func (s *blockingResourceGroupStorage) LoadResourceGroupState(keyspaceID uint32, name string) (string, error) { + if s.failNextState.CompareAndSwap(true, false) { + return "", errors.New("injected resource group state load failure") + } + return s.Storage.LoadResourceGroupState(keyspaceID, name) +} + func (s *blockingResourceGroupStorage) waitEntered(t *testing.T) { t.Helper() select { @@ -191,3 +204,44 @@ func TestAsyncLoadResourceGroupsLazyGetLegacyKeyspace(t *testing.T) { return err == nil && group != nil }, testutil.WithTickInterval(20*time.Millisecond)) } + +// TestAsyncLoadResourceGroupsRecoversFromStateLoadFailure guards against a +// group getting stuck marked reserved forever after a transient +// LoadResourceGroupState failure during lazy loading: once the async bulk +// load subsequently installs the fully-loaded (settings and state) +// confirmed data for the same group, the reserved marker must be cleared, +// otherwise loadResourceGroupIfNeeded and the state persist loop would keep +// treating already-recovered, correct data as an unconfirmed placeholder. +func TestAsyncLoadResourceGroupsRecoversFromStateLoadFailure(t *testing.T) { + re := require.New(t) + store := newBlockingResourceGroupStorage() + group := newAsyncTestGroup("flaky-group", 100) + re.NoError(store.SaveResourceGroupSetting(1, "flaky-group", group)) + re.NoError(store.SaveResourceGroupStates(1, "flaky-group", FromProtoResourceGroup(group).GetGroupStates())) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + re.NoError(m.Init(context.Background())) + defer stopAsyncTestManager(m) + defer store.unblock() + + store.waitEntered(t) + + // Make the lazy load's own state read fail once, so the group is cached + // as a metadata-only, still-reserved entry. + store.failNextState.Store(true) + fetched, err := m.GetResourceGroup(1, "flaky-group", false) + re.NoError(err) + re.NotNil(fetched) + + krgm := m.getKeyspaceResourceGroupManager(1) + re.NotNil(krgm) + re.True(krgm.isReserved("flaky-group"), "group should still be reserved after a failed state load") + + // Let the async bulk load proceed; its own state read is unaffected + // (failNextState was already consumed) and should install confirmed data. + store.unblock() + testutil.Eventually(re, func() bool { + return !krgm.isReserved("flaky-group") + }, testutil.WithTickInterval(20*time.Millisecond)) +} From 44e1d75943090c9fd7b385456f44feff0370143b Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Mon, 20 Jul 2026 14:13:50 +0800 Subject: [PATCH 14/21] resource_group: fix unparam lint in async loading tests newAsyncTestGroup's fillRate parameter always received 100, which tripped the unparam linter in the statics check. Drop the constant parameter and hoist the value into a named asyncTestGroupFillRate constant shared by the helper and the fill-rate assertions. Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- .../server/manager_async_test.go | 23 +++++++++++-------- 1 file changed, 14 insertions(+), 9 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go index 73db769b883..fdcf0704ff8 100644 --- a/pkg/mcs/resourcemanager/server/manager_async_test.go +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -83,7 +83,12 @@ func (s *blockingResourceGroupStorage) unblock() { }) } -func newAsyncTestGroup(name string, fillRate uint64) *resource_manager.ResourceGroup { +// asyncTestGroupFillRate is the fill rate used by all async-loading test +// groups; kept as a named constant so the setup and the assertions stay in +// sync. +const asyncTestGroupFillRate = 100 + +func newAsyncTestGroup(name string) *resource_manager.ResourceGroup { return &resource_manager.ResourceGroup{ Name: name, Mode: resource_manager.GroupMode_RUMode, @@ -91,8 +96,8 @@ func newAsyncTestGroup(name string, fillRate uint64) *resource_manager.ResourceG RUSettings: &resource_manager.GroupRequestUnitSettings{ RU: &resource_manager.TokenBucket{ Settings: &resource_manager.TokenLimitSettings{ - FillRate: fillRate, - BurstLimit: int64(fillRate), + FillRate: asyncTestGroupFillRate, + BurstLimit: asyncTestGroupFillRate, }, }, }, @@ -109,7 +114,7 @@ func stopAsyncTestManager(m *Manager) { func TestAsyncLoadResourceGroupsLazyGet(t *testing.T) { re := require.New(t) store := newBlockingResourceGroupStorage() - re.NoError(store.SaveResourceGroupSetting(1, "lazy-group", newAsyncTestGroup("lazy-group", 100))) + re.NoError(store.SaveResourceGroupSetting(1, "lazy-group", newAsyncTestGroup("lazy-group"))) m := NewManager[*mockConfigProvider](&mockConfigProvider{}) m.storage = store @@ -129,7 +134,7 @@ func TestAsyncLoadResourceGroupsLazyGet(t *testing.T) { re.NoError(err) re.NotNil(group) re.Equal("lazy-group", group.Name) - re.Equal(float64(100), group.RUSettings.RU.getFillRate()) + re.Equal(float64(asyncTestGroupFillRate), group.RUSettings.RU.getFillRate()) store.unblock() testutil.Eventually(re, func() bool { @@ -141,7 +146,7 @@ func TestAsyncLoadResourceGroupsLazyGet(t *testing.T) { func TestAsyncLoadResourceGroupsDoesNotRestoreDeletedLazyGroup(t *testing.T) { re := require.New(t) store := newBlockingResourceGroupStorage() - re.NoError(store.SaveResourceGroupSetting(1, "deleted-group", newAsyncTestGroup("deleted-group", 100))) + re.NoError(store.SaveResourceGroupSetting(1, "deleted-group", newAsyncTestGroup("deleted-group"))) m := NewManager[*mockConfigProvider](&mockConfigProvider{}) m.storage = store @@ -182,7 +187,7 @@ func TestAsyncLoadResourceGroupsDoesNotRestoreDeletedLazyGroup(t *testing.T) { func TestAsyncLoadResourceGroupsLazyGetLegacyKeyspace(t *testing.T) { re := require.New(t) store := newBlockingResourceGroupStorage() - re.NoError(store.SaveResourceGroupSetting(constant.NullKeyspaceID, "legacy-group", newAsyncTestGroup("legacy-group", 100))) + re.NoError(store.SaveResourceGroupSetting(constant.NullKeyspaceID, "legacy-group", newAsyncTestGroup("legacy-group"))) m := NewManager[*mockConfigProvider](&mockConfigProvider{}) m.storage = store @@ -196,7 +201,7 @@ func TestAsyncLoadResourceGroupsLazyGetLegacyKeyspace(t *testing.T) { re.NoError(err) re.NotNil(group) re.Equal("legacy-group", group.Name) - re.Equal(float64(100), group.RUSettings.RU.getFillRate()) + re.Equal(float64(asyncTestGroupFillRate), group.RUSettings.RU.getFillRate()) store.unblock() testutil.Eventually(re, func() bool { @@ -215,7 +220,7 @@ func TestAsyncLoadResourceGroupsLazyGetLegacyKeyspace(t *testing.T) { func TestAsyncLoadResourceGroupsRecoversFromStateLoadFailure(t *testing.T) { re := require.New(t) store := newBlockingResourceGroupStorage() - group := newAsyncTestGroup("flaky-group", 100) + group := newAsyncTestGroup("flaky-group") re.NoError(store.SaveResourceGroupSetting(1, "flaky-group", group)) re.NoError(store.SaveResourceGroupStates(1, "flaky-group", FromProtoResourceGroup(group).GetGroupStates())) From 5fe46baff337295e3d11a506e259c5fae21f4c3b Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Mon, 20 Jul 2026 15:41:35 +0800 Subject: [PATCH 15/21] resource_group: don't let a racing lazy load resurrect a deleted group loadResourceGroupIfNeeded reads a group from storage without holding the keyspace lock, then inserts it under the lock. A concurrent Delete that completed between the read and the insert would leave the stale insert observing an empty cache and re-adding the just-deleted group. Because the async bulk scan no longer contains that entry, the resurrected group stayed visible for the rest of the manager's lifetime. Add a per-keyspace deleteGen counter, bumped under the write lock on every cache removal. The lazy load snapshots it before its lock-free storage read and re-checks it under the insert lock; if it changed, a Delete raced and the stale result is dropped instead of inserted. A monotonic counter is used rather than a per-name tombstone map so there is no unbounded state to clear and no lifecycle window where a late in-flight reader could still slip through. Add a deterministic regression test that pauses a lazy load right after its storage read, deletes the group, then releases the lazy load and asserts it is not resurrected (verified to fail without the fix). Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- .../server/keyspace_manager.go | 15 +++ pkg/mcs/resourcemanager/server/manager.go | 16 +++- .../server/manager_async_test.go | 91 ++++++++++++++++++- 3 files changed, 117 insertions(+), 5 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/keyspace_manager.go b/pkg/mcs/resourcemanager/server/keyspace_manager.go index 8c900c54796..f9638d850ea 100644 --- a/pkg/mcs/resourcemanager/server/keyspace_manager.go +++ b/pkg/mcs/resourcemanager/server/keyspace_manager.go @@ -78,6 +78,11 @@ type keyspaceResourceGroupManager struct { reservedGroups map[string]struct{} groupRUTrackers map[string]*groupRUTracker serviceLimiter *serviceLimiter + // deleteGen is bumped under the write lock every time a group is removed + // from the cache. A lazy load snapshots it before its lock-free storage + // read and re-checks it under the lock before inserting, so a group + // deleted after the read is never resurrected by the now-stale result. + deleteGen uint64 keyspaceID uint32 storage endpoint.ResourceGroupStorage @@ -187,9 +192,19 @@ func (krgm *keyspaceResourceGroupManager) deleteResourceGroupFromCache(name stri delete(krgm.groups, name) delete(krgm.groupRUTrackers, name) delete(krgm.reservedGroups, name) + // Signal any in-flight lazy load that a delete happened, so it won't + // reinsert a copy read from storage before this deletion. + krgm.deleteGen++ krgm.Unlock() } +// loadDeleteGen returns the current delete generation counter. +func (krgm *keyspaceResourceGroupManager) loadDeleteGen() uint64 { + krgm.RLock() + defer krgm.RUnlock() + return krgm.deleteGen +} + func (krgm *keyspaceResourceGroupManager) setRawStatesIntoResourceGroup(name string, rawValue string) error { tokens := &GroupStates{} if err := json.Unmarshal([]byte(rawValue), tokens); err != nil { diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index e45be2cad34..aaa22d6a5b6 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -605,6 +605,12 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro return nil } } + // Ensure the keyspace manager exists and snapshot its delete generation + // before the lock-free storage read below, so a concurrent Delete that + // lands after the read is detected under the insert lock and can't be + // undone by this now-stale result. + krgm = m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false) + deleteGen := krgm.loadDeleteGen() group, stateLoaded, err := m.loadResourceGroup(keyspaceID, name) if err != nil { if name == DefaultResourceGroupName && errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(name)) { @@ -613,14 +619,20 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro // This calls initDefaultResourceGroup directly instead of going // through getOrCreateKeyspaceResourceGroupManager(id, true), which // now routes back into this same function and would recurse. - m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false).initDefaultResourceGroup() + krgm.initDefaultResourceGroup() return nil } return err } - krgm = m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false) inserted := false krgm.Lock() + if krgm.deleteGen != deleteGen { + // A Delete raced with our storage read; the result may be stale, so + // don't insert it. A later request or the async bulk merge will reload + // the group if it still exists. + krgm.Unlock() + return nil + } if _, exists := krgm.groups[name]; !exists { krgm.groups[name] = group inserted = true diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go index fdcf0704ff8..f29e22e2fd0 100644 --- a/pkg/mcs/resourcemanager/server/manager_async_test.go +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -43,13 +43,22 @@ type blockingResourceGroupStorage struct { // failNextState, when true, makes the very next LoadResourceGroupState // call fail once, then resets itself. failNextState atomic.Bool + + // pausePointState, when true, makes the very next LoadResourceGroupState + // call signal pointReached and then block on pointRelease, so a test can + // hold a lazy load right after its storage read but before it inserts. + pausePointState atomic.Bool + pointReached chan struct{} + pointRelease chan struct{} } func newBlockingResourceGroupStorage() *blockingResourceGroupStorage { return &blockingResourceGroupStorage{ - Storage: storage.NewStorageWithMemoryBackend(), - entered: make(chan struct{}), - release: make(chan struct{}), + Storage: storage.NewStorageWithMemoryBackend(), + entered: make(chan struct{}), + release: make(chan struct{}), + pointReached: make(chan struct{}), + pointRelease: make(chan struct{}), } } @@ -65,6 +74,10 @@ func (s *blockingResourceGroupStorage) LoadResourceGroupState(keyspaceID uint32, if s.failNextState.CompareAndSwap(true, false) { return "", errors.New("injected resource group state load failure") } + if s.pausePointState.CompareAndSwap(true, false) { + close(s.pointReached) + <-s.pointRelease + } return s.Storage.LoadResourceGroupState(keyspaceID, name) } @@ -250,3 +263,75 @@ func TestAsyncLoadResourceGroupsRecoversFromStateLoadFailure(t *testing.T) { return !krgm.isReserved("flaky-group") }, testutil.WithTickInterval(20*time.Millisecond)) } + +// TestAsyncLoadResourceGroupsDeleteRaceDoesNotResurrect reproduces the +// lazy-load vs concurrent Delete race deterministically: a lazy load reads a +// group from storage, then a Delete removes it before the lazy load inserts. +// The stale insert must be rejected (via the delete-generation check) so the +// deleted group is not resurrected for the rest of the manager's lifetime. +func TestAsyncLoadResourceGroupsDeleteRaceDoesNotResurrect(t *testing.T) { + re := require.New(t) + store := newBlockingResourceGroupStorage() + group := newAsyncTestGroup("race-group") + re.NoError(store.SaveResourceGroupSetting(1, "race-group", group)) + re.NoError(store.SaveResourceGroupStates(1, "race-group", FromProtoResourceGroup(group).GetGroupStates())) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + re.NoError(m.Init(context.Background())) + defer stopAsyncTestManager(m) + defer store.unblock() + + // Async bulk load is blocked, so loadingState stays in progress and lazy + // loading is active. + store.waitEntered(t) + + // Start a lazy Get that will pause inside its state read, i.e. after it has + // read the group from storage but before it inserts into the cache. + store.pausePointState.Store(true) + var ( + gotGroup *ResourceGroup + gotErr error + ) + getDone := make(chan struct{}) + go func() { + defer close(getDone) + gotGroup, gotErr = m.GetResourceGroup(1, "race-group", false) + }() + + select { + case <-store.pointReached: + case <-time.After(time.Second): + t.Fatal("timed out waiting for the lazy load to reach its state read") + } + + // While the lazy load is paused, delete the group. Delete does its own + // (unpaused) load-then-delete, removing it from storage and cache and + // bumping the delete generation. + re.NoError(m.DeleteResourceGroup(1, "race-group")) + + // Release the paused lazy load; its now-stale insert must be rejected. + close(store.pointRelease) + <-getDone + re.NoError(gotErr) + re.Nil(gotGroup, "the racing lazy load must observe the group as deleted") + + krgm := m.getKeyspaceResourceGroupManager(1) + re.NotNil(krgm) + re.Nil(krgm.getMutableResourceGroup("race-group"), "deleted group must not be resurrected by the racing lazy load") + + // Finishing async loading must not bring the deleted group back either. + store.unblock() + testutil.Eventually(re, func() bool { + groups, err := m.GetResourceGroupList(1, false) + if err != nil { + return false + } + for _, g := range groups { + if g.Name == "race-group" { + return false + } + } + return true + }, testutil.WithTickInterval(20*time.Millisecond)) +} From 229cd5dd5a08d7d9a58d12d98fef3eccb3a07451 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Wed, 22 Jul 2026 15:27:39 +0800 Subject: [PATCH 16/21] resource_group: prevent a stale async loader from polluting a new term asyncLoadResourceGroups only checked its context at the top of the retry loop. A loader blocked in a storage scan across a leadership change would, after Init ran again for a new term, wake up and merge its stale scan into the new term's maps, clear the new term's syncLoadedGroups, and publish LoadingStateCompleted while the new loader was still running. Add a loadEpoch counter to the manager, bumped by initMetadata under the manager lock and captured by each loader at start. Every shared-state mutation the loader performs (loading-state transitions and the merge) now re-verifies the epoch inside the same critical section, so a stale loader exits instead of touching the newer term's state. Also re-check context cancellation right after the scans return, and make initControllerConfig publish the config via clone-and-swap under the lock, since a re-initialization can race with the previous term's background goroutines still reading it. Add a deterministic regression test that blocks a term-1 loader in its states scan, reinitializes the manager for term 2, deletes the group, then releases the stale loader and asserts it does not resurrect the group or disturb the new term (verified to fail without the fix). Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/manager.go | 70 ++++++++++++-- .../server/manager_async_test.go | 96 ++++++++++++++++++- 2 files changed, 154 insertions(+), 12 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index aaa22d6a5b6..1d84f775feb 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -119,6 +119,12 @@ type Manager struct { loadingState int32 // atomic access // syncLoadedGroups records groups that were loaded synchronously (e.g., by lazy loading) syncLoadedGroups map[trackerKey]bool + // loadEpoch is bumped (under the manager lock) every time initMetadata + // resets the loading state for a new term. The async loader captures it at + // start and re-checks it before every shared-state mutation, so a loader + // from a previous term that was blocked in a storage scan can never merge + // stale data into, or publish completion for, a newer term. + loadEpoch uint64 } // LoadingState represents the current loading state of resource groups @@ -392,13 +398,22 @@ func (m *Manager) initControllerConfig() error { log.Error("resource controller config load failed", zap.Error(err), zap.String("v", v)) return err } - if err = json.Unmarshal([]byte(v), &m.controllerConfig); err != nil { + // Unmarshal into a clone and publish it under the lock: on a + // re-initialization after a leadership change, the previous term's + // background goroutines may still be reading the current config. + m.RLock() + controllerConfig := cloneControllerConfig(m.controllerConfig) + m.RUnlock() + if err = json.Unmarshal([]byte(v), &controllerConfig); err != nil { log.Warn("un-marshall controller config failed, fallback to default", zap.Error(err), zap.String("v", v)) } + m.Lock() + m.controllerConfig = controllerConfig + m.Unlock() // re-save the config to make sure the config has been persisted. if m.writeRole.AllowsMetadataWrite() { - if err := m.storage.SaveControllerConfig(m.controllerConfig); err != nil { + if err := m.storage.SaveControllerConfig(controllerConfig); err != nil { return err } } @@ -413,6 +428,8 @@ func (m *Manager) initMetadata(ctx context.Context) error { m.Lock() m.krgms = make(map[uint32]*keyspaceResourceGroupManager) m.syncLoadedGroups = make(map[trackerKey]bool) + m.loadEpoch++ + epoch := m.loadEpoch atomic.StoreInt32(&m.loadingState, LoadingStateNotStarted) m.Unlock() @@ -422,7 +439,7 @@ func (m *Manager) initMetadata(ctx context.Context) error { } m.wg.Add(1) - go m.asyncLoadResourceGroups(ctx) + go m.asyncLoadResourceGroups(ctx, epoch) return nil } @@ -446,7 +463,21 @@ func (m *Manager) loadKeyspaceResourceGroups() error { return m.loadServiceLimits() } -func (m *Manager) asyncLoadResourceGroups(ctx context.Context) { +// storeLoadingStateIfCurrent stores the loading state only if the manager has +// not been reinitialized since the loader with the given epoch started. It +// returns false when the loader is stale and must exit without touching any +// further shared state. +func (m *Manager) storeLoadingStateIfCurrent(epoch uint64, state int32) bool { + m.Lock() + defer m.Unlock() + if m.loadEpoch != epoch { + return false + } + atomic.StoreInt32(&m.loadingState, state) + return true +} + +func (m *Manager) asyncLoadResourceGroups(ctx context.Context, epoch uint64) { defer logutil.LogPanic() defer m.wg.Done() @@ -471,18 +502,40 @@ func (m *Manager) asyncLoadResourceGroups(ctx context.Context) { } } - atomic.StoreInt32(&m.loadingState, LoadingStateInProgress) + if !m.storeLoadingStateIfCurrent(epoch, LoadingStateInProgress) { + log.Info("async loading resource groups aborted: manager was reinitialized") + return + } startTime := time.Now() tempKrgms, err := m.loadKeyspaceResourceGroupsFromStorage() + // The storage scans above can block for a long time; re-check for + // cancellation before touching any shared state, so a loader whose + // term already ended doesn't pollute a newer term's state. + select { + case <-ctx.Done(): + log.Info("async loading resource groups cancelled") + return + default: + } if err != nil { log.Error("failed to load resource groups", zap.Error(err), zap.Int("retry", retry)) - atomic.StoreInt32(&m.loadingState, LoadingStateNotStarted) + if !m.storeLoadingStateIfCurrent(epoch, LoadingStateNotStarted) { + log.Info("async loading resource groups aborted: manager was reinitialized") + return + } retry++ continue } loaded := 0 m.Lock() + if m.loadEpoch != epoch { + // The manager was reinitialized for a new term while this loader + // was scanning; its result is stale and must not be merged. + m.Unlock() + log.Info("async loading resource groups aborted: manager was reinitialized") + return + } for keyspaceID, tempKrgm := range tempKrgms { krgm := m.krgms[keyspaceID] if krgm == nil { @@ -515,7 +568,10 @@ func (m *Manager) asyncLoadResourceGroups(ctx context.Context) { m.Unlock() m.initReserved() - atomic.StoreInt32(&m.loadingState, LoadingStateCompleted) + if !m.storeLoadingStateIfCurrent(epoch, LoadingStateCompleted) { + log.Info("async loading resource groups aborted: manager was reinitialized") + return + } duration := time.Since(startTime) asyncLoadGroupDuration.Observe(duration.Seconds()) log.Info("async loading resource groups completed", zap.Int("loaded-groups", loaded), zap.Duration("duration", duration)) diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go index f29e22e2fd0..cf1562f66a6 100644 --- a/pkg/mcs/resourcemanager/server/manager_async_test.go +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -50,15 +50,26 @@ type blockingResourceGroupStorage struct { pausePointState atomic.Bool pointReached chan struct{} pointRelease chan struct{} + + // pauseNextStates, when true, makes the very next bulk + // LoadResourceGroupStates call signal statesReached and then block on + // statesRelease, so a test can hold an async loader after it has captured + // the settings scan but before it merges. + pauseNextStates atomic.Bool + statesReached chan struct{} + statesRelease chan struct{} + statesReleaseOnce sync.Once } func newBlockingResourceGroupStorage() *blockingResourceGroupStorage { return &blockingResourceGroupStorage{ - Storage: storage.NewStorageWithMemoryBackend(), - entered: make(chan struct{}), - release: make(chan struct{}), - pointReached: make(chan struct{}), - pointRelease: make(chan struct{}), + Storage: storage.NewStorageWithMemoryBackend(), + entered: make(chan struct{}), + release: make(chan struct{}), + pointReached: make(chan struct{}), + pointRelease: make(chan struct{}), + statesReached: make(chan struct{}), + statesRelease: make(chan struct{}), } } @@ -81,6 +92,14 @@ func (s *blockingResourceGroupStorage) LoadResourceGroupState(keyspaceID uint32, return s.Storage.LoadResourceGroupState(keyspaceID, name) } +func (s *blockingResourceGroupStorage) LoadResourceGroupStates(f func(keyspaceID uint32, name, rawValue string)) error { + if s.pauseNextStates.CompareAndSwap(true, false) { + close(s.statesReached) + <-s.statesRelease + } + return s.Storage.LoadResourceGroupStates(f) +} + func (s *blockingResourceGroupStorage) waitEntered(t *testing.T) { t.Helper() select { @@ -96,6 +115,12 @@ func (s *blockingResourceGroupStorage) unblock() { }) } +func (s *blockingResourceGroupStorage) unblockStates() { + s.statesReleaseOnce.Do(func() { + close(s.statesRelease) + }) +} + // asyncTestGroupFillRate is the fill rate used by all async-loading test // groups; kept as a named constant so the setup and the assertions stay in // sync. @@ -335,3 +360,64 @@ func TestAsyncLoadResourceGroupsDeleteRaceDoesNotResurrect(t *testing.T) { return true }, testutil.WithTickInterval(20*time.Millisecond)) } + +// TestAsyncLoadResourceGroupsStaleLoaderDoesNotPolluteNewTerm reproduces the +// stale-loader race: a loader from an old term is blocked in its storage scan +// while the leadership changes and Init runs again for a new term. When the +// old loader finally wakes up, it must not merge its stale scan into the new +// term's maps, clear the new term's syncLoadedGroups, or publish completion. +func TestAsyncLoadResourceGroupsStaleLoaderDoesNotPolluteNewTerm(t *testing.T) { + re := require.New(t) + store := newBlockingResourceGroupStorage() + group := newAsyncTestGroup("stale-group") + re.NoError(store.SaveResourceGroupSetting(1, "stale-group", group)) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + // Term 1: the loader blocks at the start of its settings scan. + re.NoError(m.Init(context.Background())) + cancelTerm1 := m.cancel + defer stopAsyncTestManager(m) + defer store.unblock() + defer store.unblockStates() + + store.waitEntered(t) + + // Let the term-1 loader run its settings scan (capturing stale-group into + // its temp result) and then block again in the states scan, i.e. after it + // has read storage but before it merges. + store.pauseNextStates.Store(true) + store.unblock() + select { + case <-store.statesReached: + case <-time.After(time.Second): + t.Fatal("timed out waiting for the term-1 loader to reach its states scan") + } + + // Leadership changes: cancel term 1 and reinitialize for term 2. The + // term-2 loader hits neither block (both were consumed) and completes. + cancelTerm1() + re.NoError(m.Init(context.Background())) + testutil.Eventually(re, func() bool { + groups, err := m.GetResourceGroupList(1, false) + return err == nil && len(groups) == 2 + }, testutil.WithTickInterval(20*time.Millisecond)) + + // Delete the group in term 2, after loading completed. + re.NoError(m.DeleteResourceGroup(1, "stale-group")) + + // Release the stale term-1 loader. It must observe its cancelled context / + // stale epoch and exit without resurrecting the deleted group or touching + // the new term's loading state. + store.unblockStates() + time.Sleep(200 * time.Millisecond) + + krgm := m.getKeyspaceResourceGroupManager(1) + re.NotNil(krgm) + re.Nil(krgm.getMutableResourceGroup("stale-group"), "stale loader must not merge into the new term") + groups, err := m.GetResourceGroupList(1, false) + re.NoError(err) + for _, g := range groups { + re.NotEqual("stale-group", g.Name) + } +} From 10b6c2abfc33fb6eb007fbd66e454aabafb174d6 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Wed, 22 Jul 2026 15:34:00 +0800 Subject: [PATCH 17/21] resource_group: don't let the bulk merge clobber modified settings When a lazy load reads a group's settings but fails to read its state, the entry is cached as reserved and deliberately kept out of syncLoadedGroups so the async bulk merge can still recover the real state. But the merge replaced the cache entry wholesale, so a modification persisted after the bulk scan captured its (older) settings was silently lost from the serving cache for the rest of the manager's lifetime, while storage kept the new values. Split the reserved marker into two kinds: reservedPlaceholder (settings and state both synthetic, e.g. the pre-inserted default group), which the merge may still replace wholesale, and reservedStateOnly (settings confirmed from storage, state missing), for which the merge now adopts only the scanned state into the existing entry and keeps its settings. Add a regression test that captures the bulk scan before a Modify, fails the lazy load's state reads, modifies the group, then releases the merge and asserts the modified settings survive while the scanned state is adopted (verified to fail without the fix). Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- .../server/keyspace_manager.go | 35 ++++++--- pkg/mcs/resourcemanager/server/manager.go | 36 +++++++--- .../server/manager_async_test.go | 72 +++++++++++++++++++ 3 files changed, 124 insertions(+), 19 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/keyspace_manager.go b/pkg/mcs/resourcemanager/server/keyspace_manager.go index f9638d850ea..d37983a7872 100644 --- a/pkg/mcs/resourcemanager/server/keyspace_manager.go +++ b/pkg/mcs/resourcemanager/server/keyspace_manager.go @@ -67,15 +67,32 @@ type consumptionItem struct { isTiFlash bool } +// reservedKind describes why a cached resource group entry is still +// considered unconfirmed. +type reservedKind int + +const ( + // reservedPlaceholder marks an entry whose settings and state are both + // synthetic, e.g. the default group pre-inserted by + // ensureReservedDefaultGroupInCache before async loading runs. The async + // bulk merge may replace such an entry wholesale. + reservedPlaceholder reservedKind = iota + // reservedStateOnly marks an entry whose settings were confirmed from + // storage (and may have been modified and persisted since) but whose + // state failed to load. The async bulk merge must only adopt the scanned + // state into it, never replace its settings with the scan's older copy. + reservedStateOnly +) + type keyspaceResourceGroupManager struct { syncutil.RWMutex groups map[string]*ResourceGroup - // reservedGroups tracks names whose entry in groups is still just the - // synthetic placeholder inserted by ensureReservedDefaultGroupInCache or - // restoreDefaultResourceGroupFromReserved, not yet confirmed by a storage - // load or a real write. It shares the same lock as groups so a lazy load - // can atomically decide whether it's safe to replace the placeholder. - reservedGroups map[string]struct{} + // reservedGroups tracks names whose entry in groups is not yet fully + // confirmed by a storage load or a real write, together with how much of + // it is unconfirmed (see reservedKind). It shares the same lock as groups + // so a lazy load can atomically decide whether it's safe to replace the + // entry. + reservedGroups map[string]reservedKind groupRUTrackers map[string]*groupRUTracker serviceLimiter *serviceLimiter // deleteGen is bumped under the write lock every time a group is removed @@ -100,7 +117,7 @@ func newKeyspaceResourceGroupManager( } return &keyspaceResourceGroupManager{ groups: make(map[string]*ResourceGroup), - reservedGroups: make(map[string]struct{}), + reservedGroups: make(map[string]reservedKind), groupRUTrackers: make(map[string]*groupRUTracker), keyspaceID: keyspaceID, storage: storage, @@ -245,7 +262,7 @@ func (krgm *keyspaceResourceGroupManager) ensureReservedDefaultGroupInCache() { krgm.Lock() if _, ok := krgm.groups[DefaultResourceGroupName]; !ok { krgm.groups[DefaultResourceGroupName] = defaultGroup - krgm.reservedGroups[DefaultResourceGroupName] = struct{}{} + krgm.reservedGroups[DefaultResourceGroupName] = reservedPlaceholder inserted = true } krgm.Unlock() @@ -273,7 +290,7 @@ func (krgm *keyspaceResourceGroupManager) restoreDefaultResourceGroupFromReserve defaultGroup := newDefaultResourceGroup() krgm.Lock() krgm.groups[DefaultResourceGroupName] = defaultGroup - krgm.reservedGroups[DefaultResourceGroupName] = struct{}{} + krgm.reservedGroups[DefaultResourceGroupName] = reservedPlaceholder krgm.Unlock() krgm.syncBurstabilityWithServiceLimit(defaultGroup) } diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 1d84f775feb..53bfbc73348 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -547,16 +547,32 @@ func (m *Manager) asyncLoadResourceGroups(ctx context.Context, epoch uint64) { krgm.Lock() for name, group := range tempKrgm.groups { key := trackerKey{keyspaceID: keyspaceID, groupName: name} - if !m.syncLoadedGroups[key] { - krgm.groups[name] = group - // This group is now confirmed, fully-loaded data (settings - // and state); it must no longer be treated as an - // unconfirmed placeholder by loadResourceGroupIfNeeded or - // skipped by the state persist loop. - delete(krgm.reservedGroups, name) - groupsToSync = append(groupsToSync, group) - loaded++ + if m.syncLoadedGroups[key] { + continue + } + if kind, reserved := krgm.reservedGroups[name]; reserved && kind == reservedStateOnly { + if existing, ok := krgm.groups[name]; ok { + // The cached entry's settings are confirmed and may + // carry a modification persisted after this scan + // started; only its state is missing. Adopt the + // scanned state into the existing entry instead of + // replacing it, so the newer settings aren't + // clobbered by the scan's older copy. + existing.SetStatesIntoResourceGroup(group.GetGroupStates()) + delete(krgm.reservedGroups, name) + groupsToSync = append(groupsToSync, existing) + loaded++ + continue + } } + krgm.groups[name] = group + // This group is now confirmed, fully-loaded data (settings + // and state); it must no longer be treated as an + // unconfirmed placeholder by loadResourceGroupIfNeeded or + // skipped by the state persist loop. + delete(krgm.reservedGroups, name) + groupsToSync = append(groupsToSync, group) + loaded++ } krgm.Unlock() tempKrgm.RUnlock() @@ -705,7 +721,7 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro if stateLoaded { delete(krgm.reservedGroups, name) } else { - krgm.reservedGroups[name] = struct{}{} + krgm.reservedGroups[name] = reservedStateOnly } krgm.Unlock() if inserted { diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go index cf1562f66a6..98f8f4d2dd5 100644 --- a/pkg/mcs/resourcemanager/server/manager_async_test.go +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -421,3 +421,75 @@ func TestAsyncLoadResourceGroupsStaleLoaderDoesNotPolluteNewTerm(t *testing.T) { re.NotEqual("stale-group", g.Name) } } + +// TestAsyncLoadResourceGroupsMergeKeepsModifiedSettings reproduces the +// stale-settings clobber: a group's lazy load reads its settings but fails to +// read its state, then the group is modified (and the modification persisted) +// while the async bulk scan still holds the pre-modification settings. The +// bulk merge must adopt only the scanned state into the cached entry, not +// replace it wholesale, so the modified settings survive in the serving cache. +func TestAsyncLoadResourceGroupsMergeKeepsModifiedSettings(t *testing.T) { + re := require.New(t) + store := newBlockingResourceGroupStorage() + group := newAsyncTestGroup("mod-group") + re.NoError(store.SaveResourceGroupSetting(1, "mod-group", group)) + // Persist a state with a recognizable consumption so the test can tell + // that the merge really adopted the scanned state. + states := FromProtoResourceGroup(group).GetGroupStates() + states.RUConsumption.RRU = 123 + re.NoError(store.SaveResourceGroupStates(1, "mod-group", states)) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + re.NoError(m.Init(context.Background())) + defer stopAsyncTestManager(m) + defer store.unblock() + defer store.unblockStates() + + store.waitEntered(t) + + // Let the bulk loader capture the pre-modification settings, then hold it + // right before its states scan (i.e. before it merges). + store.pauseNextStates.Store(true) + store.unblock() + select { + case <-store.statesReached: + case <-time.After(time.Second): + t.Fatal("timed out waiting for the bulk loader to reach its states scan") + } + + // Lazily load the group with a failing state read: it is cached with + // confirmed settings but unconfirmed (fresh) state. + store.failNextState.Store(true) + fetched, err := m.GetResourceGroup(1, "mod-group", false) + re.NoError(err) + re.NotNil(fetched) + krgm := m.getKeyspaceResourceGroupManager(1) + re.NotNil(krgm) + re.True(krgm.isReserved("mod-group")) + + // Modify the group (fill rate 100 -> 200) and persist it. The state read + // of Modify's own lazy load fails again, so the entry stays state-only + // reserved and out of syncLoadedGroups. + store.failNextState.Store(true) + modified := newAsyncTestGroup("mod-group") + modified.RUSettings.RU.Settings.FillRate = 200 + modified.KeyspaceId = &resource_manager.KeyspaceIDValue{Value: 1} + re.NoError(m.ModifyResourceGroup(modified)) + + // Release the bulk loader; its merge must keep the modified settings and + // only adopt the scanned state. + store.unblockStates() + testutil.Eventually(re, func() bool { + _, err := m.GetResourceGroupList(1, false) + return err == nil + }, testutil.WithTickInterval(20*time.Millisecond)) + + got, err := m.GetResourceGroup(1, "mod-group", false) + re.NoError(err) + re.NotNil(got) + re.Equal(float64(200), got.RUSettings.RU.getFillRate(), "modified settings must survive the bulk merge") + re.False(krgm.isReserved("mod-group"), "state adoption must clear the reserved marker") + re.Equal(float64(123), krgm.getMutableResourceGroup("mod-group").GetGroupStates().RUConsumption.RRU, + "the scanned state must be adopted into the cached entry") +} From 64e13136909eaa21d5655d7faa3784118c9fb26c Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Wed, 22 Jul 2026 15:41:45 +0800 Subject: [PATCH 18/21] resource_group: persist fresh-store default and retry lazy loads on delete races Two fixes for the async loading path: initDefaultResourceGroup bailed out whenever any cache entry existed for the default group, including the synthetic placeholder pre-inserted by initReservedInCache. On a fresh store, the confirmed-not-found fallback therefore never created or persisted the default group, and the entry stayed an unconfirmed placeholder for the manager lifetime: its settings were never stored and the persist loop permanently skipped its token and consumption state. Treat a reservedPlaceholder entry as absent (synthesize and persist), while still returning early for confirmed entries and reservedStateOnly ones, whose real settings must not be overwritten with synthetic values. The per-keyspace deleteGen is shared by every group, so deleting group B while group A was being lazily loaded discarded A's valid result and made its request spuriously report the group as missing. Retry the storage read (up to 3 attempts) on a generation mismatch: the re-read observes post-delete storage, so an unrelated deletion just reloads successfully, and a deletion of the group itself now correctly reports not-found instead of silently returning nothing. Add regression tests for both (each verified to fail without its fix). Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- .../server/keyspace_manager.go | 10 +- pkg/mcs/resourcemanager/server/manager.go | 124 ++++++++++-------- .../server/manager_async_test.go | 113 +++++++++++++++- 3 files changed, 187 insertions(+), 60 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/keyspace_manager.go b/pkg/mcs/resourcemanager/server/keyspace_manager.go index d37983a7872..15c9e0af82d 100644 --- a/pkg/mcs/resourcemanager/server/keyspace_manager.go +++ b/pkg/mcs/resourcemanager/server/keyspace_manager.go @@ -240,8 +240,16 @@ func (krgm *keyspaceResourceGroupManager) setRawStatesIntoResourceGroup(name str func (krgm *keyspaceResourceGroupManager) initDefaultResourceGroup() { krgm.RLock() _, ok := krgm.groups[DefaultResourceGroupName] + kind, reserved := krgm.reservedGroups[DefaultResourceGroupName] krgm.RUnlock() - if ok { + // A cached entry only makes initialization unnecessary if it's confirmed + // data, or at least has confirmed settings (reservedStateOnly), which must + // not be overwritten with synthetic ones. A reservedPlaceholder means + // nothing is persisted for the default group (e.g. a fresh store): it must + // still be created and persisted here, otherwise it would stay an + // unconfirmed placeholder forever, with its settings never stored and its + // state persistence permanently skipped. + if ok && (!reserved || kind == reservedStateOnly) { return } defaultGroup := newDefaultResourceGroup() diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 53bfbc73348..f41e7a1d9f3 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -677,69 +677,77 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro return nil } } - // Ensure the keyspace manager exists and snapshot its delete generation - // before the lock-free storage read below, so a concurrent Delete that - // lands after the read is detected under the insert lock and can't be - // undone by this now-stale result. krgm = m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false) - deleteGen := krgm.loadDeleteGen() - group, stateLoaded, err := m.loadResourceGroup(keyspaceID, name) - if err != nil { - if name == DefaultResourceGroupName && errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(name)) { - // No persisted default group settings exist yet (e.g. a brand-new - // keyspace), so it's safe to synthesize the reserved default group. - // This calls initDefaultResourceGroup directly instead of going - // through getOrCreateKeyspaceResourceGroupManager(id, true), which - // now routes back into this same function and would recurse. - krgm.initDefaultResourceGroup() - return nil + // deleteGen is shared by every group in the keyspace, so a concurrent + // Delete of any group invalidates the lock-free storage read below. + // Deletes are rare: retry the read a few times so an unrelated deletion + // doesn't make this request spuriously report the group as missing; if it + // keeps racing, give up without inserting and let a later request or the + // async bulk merge reload the group. + const maxLoadAttempts = 3 + for attempt := 1; ; attempt++ { + // Snapshot the delete generation before the lock-free storage read, + // so a Delete that lands after the read is detected under the insert + // lock and can't be undone by the now-stale result. + deleteGen := krgm.loadDeleteGen() + group, stateLoaded, err := m.loadResourceGroup(keyspaceID, name) + if err != nil { + if name == DefaultResourceGroupName && errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(name)) { + // No persisted default group settings exist yet (e.g. a brand-new + // keyspace), so it's safe to synthesize the reserved default group. + // This calls initDefaultResourceGroup directly instead of going + // through getOrCreateKeyspaceResourceGroupManager(id, true), which + // now routes back into this same function and would recurse. + krgm.initDefaultResourceGroup() + return nil + } + return err + } + inserted := false + krgm.Lock() + if krgm.deleteGen != deleteGen { + krgm.Unlock() + if attempt >= maxLoadAttempts { + return nil + } + continue + } + if _, exists := krgm.groups[name]; !exists { + krgm.groups[name] = group + inserted = true + } else if _, reserved := krgm.reservedGroups[name]; reserved { + // The existing entry is just an unconfirmed placeholder; the freshly + // loaded group is the real, confirmed data, so replacing it is safe. + krgm.groups[name] = group + inserted = true + } + // Only clear the placeholder mark once the state was actually read; a + // metadata-only group (state load failed) must stay reserved so a later + // call or the async bulk merge remains free to fill in the real state, + // instead of this partial result being treated as final forever. + if stateLoaded { + delete(krgm.reservedGroups, name) + } else { + krgm.reservedGroups[name] = reservedStateOnly } - return err - } - inserted := false - krgm.Lock() - if krgm.deleteGen != deleteGen { - // A Delete raced with our storage read; the result may be stale, so - // don't insert it. A later request or the async bulk merge will reload - // the group if it still exists. krgm.Unlock() - return nil - } - if _, exists := krgm.groups[name]; !exists { - krgm.groups[name] = group - inserted = true - } else if _, reserved := krgm.reservedGroups[name]; reserved { - // The existing entry is just an unconfirmed placeholder; the freshly - // loaded group is the real, confirmed data, so replacing it is safe. - krgm.groups[name] = group - inserted = true - } - // Only clear the placeholder mark once the state was actually read; a - // metadata-only group (state load failed) must stay reserved so a later - // call or the async bulk merge remains free to fill in the real state, - // instead of this partial result being treated as final forever. - if stateLoaded { - delete(krgm.reservedGroups, name) - } else { - krgm.reservedGroups[name] = reservedStateOnly - } - krgm.Unlock() - if inserted { - krgm.syncBurstabilityWithServiceLimit(group) - } - // Only mark the group as sync-loaded when its persisted state was actually - // read; otherwise a later async bulk load must remain free to fill in the - // real state instead of being skipped forever. - if stateLoaded { - markKey := trackerKey{keyspaceID: keyspaceID, groupName: name} - m.Lock() - if m.syncLoadedGroups != nil { - m.syncLoadedGroups[markKey] = true + if inserted { + krgm.syncBurstabilityWithServiceLimit(group) } - m.Unlock() + // Only mark the group as sync-loaded when its persisted state was actually + // read; otherwise a later async bulk load must remain free to fill in the + // real state instead of being skipped forever. + if stateLoaded { + markKey := trackerKey{keyspaceID: keyspaceID, groupName: name} + m.Lock() + if m.syncLoadedGroups != nil { + m.syncLoadedGroups[markKey] = true + } + m.Unlock() + } + syncLoadGroupCounter.Inc() + return nil } - syncLoadGroupCounter.Inc() - return nil } func (m *Manager) markResourceGroupSyncLoaded(keyspaceID uint32, name string) { diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go index 98f8f4d2dd5..7dd44f521cf 100644 --- a/pkg/mcs/resourcemanager/server/manager_async_test.go +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -338,7 +338,9 @@ func TestAsyncLoadResourceGroupsDeleteRaceDoesNotResurrect(t *testing.T) { // Release the paused lazy load; its now-stale insert must be rejected. close(store.pointRelease) <-getDone - re.NoError(gotErr) + // The generation mismatch makes the lazy load retry its storage read, + // which now correctly observes the group as deleted. + re.ErrorContains(gotErr, "does not exist") re.Nil(gotGroup, "the racing lazy load must observe the group as deleted") krgm := m.getKeyspaceResourceGroupManager(1) @@ -493,3 +495,112 @@ func TestAsyncLoadResourceGroupsMergeKeepsModifiedSettings(t *testing.T) { re.Equal(float64(123), krgm.getMutableResourceGroup("mod-group").GetGroupStates().RUConsumption.RRU, "the scanned state must be adopted into the cached entry") } + +// TestAsyncLoadResourceGroupsFreshStoreDefaultPersisted guards against the +// fresh-store dead end: initReservedInCache pre-inserts a synthetic default +// placeholder, and on a store with nothing persisted, the confirmed-not-found +// fallback used to bail out on its cache-exists check, leaving the default +// group an unconfirmed placeholder forever — settings never persisted and its +// state persistence permanently skipped. +func TestAsyncLoadResourceGroupsFreshStoreDefaultPersisted(t *testing.T) { + re := require.New(t) + // A completely fresh store: nothing persisted at all. + store := newBlockingResourceGroupStorage() + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + re.NoError(m.Init(context.Background())) + defer stopAsyncTestManager(m) + defer store.unblock() + + store.waitEntered(t) + + // Fetch the default group while async loading is still in progress: the + // point load confirms nothing is persisted, so the placeholder must be + // promoted to a real, persisted default group. + group, err := m.GetResourceGroup(constant.NullKeyspaceID, DefaultResourceGroupName, false) + re.NoError(err) + re.NotNil(group) + krgm := m.getKeyspaceResourceGroupManager(constant.NullKeyspaceID) + re.NotNil(krgm) + re.False(krgm.isReserved(DefaultResourceGroupName), "the default group must be confirmed after synthesis") + raw, err := store.LoadResourceGroupSetting(constant.NullKeyspaceID, DefaultResourceGroupName) + re.NoError(err) + re.NotEmpty(raw, "the synthesized default group settings must be persisted") + + // Loading completion must keep it confirmed. + store.unblock() + testutil.Eventually(re, func() bool { + _, err := m.GetResourceGroupList(constant.NullKeyspaceID, false) + return err == nil + }, testutil.WithTickInterval(20*time.Millisecond)) + re.False(krgm.isReserved(DefaultResourceGroupName)) +} + +// TestAsyncLoadResourceGroupsUnrelatedDeleteDoesNotFailLazyLoad guards against +// the delete-generation check being too coarse: deleting group B while group A +// is being lazily loaded must not make A's request spuriously report the group +// as missing — the lazy load retries its storage read and succeeds. +func TestAsyncLoadResourceGroupsUnrelatedDeleteDoesNotFailLazyLoad(t *testing.T) { + re := require.New(t) + store := newBlockingResourceGroupStorage() + re.NoError(store.SaveResourceGroupSetting(1, "group-a", newAsyncTestGroup("group-a"))) + re.NoError(store.SaveResourceGroupSetting(1, "group-b", newAsyncTestGroup("group-b"))) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + re.NoError(m.Init(context.Background())) + defer stopAsyncTestManager(m) + defer store.unblock() + + store.waitEntered(t) + + // Start a lazy Get of group-a and pause it inside its state read, i.e. + // after it has read the group from storage but before it inserts. + store.pausePointState.Store(true) + var ( + gotGroup *ResourceGroup + gotErr error + ) + getDone := make(chan struct{}) + go func() { + defer close(getDone) + gotGroup, gotErr = m.GetResourceGroup(1, "group-a", false) + }() + select { + case <-store.pointReached: + case <-time.After(time.Second): + t.Fatal("timed out waiting for the lazy load to reach its state read") + } + + // Delete the unrelated group-b while group-a's lazy load is paused; this + // bumps the keyspace's delete generation. + re.NoError(m.DeleteResourceGroup(1, "group-b")) + + // Release group-a's lazy load: the generation mismatch must make it retry + // and succeed, not report group-a as missing. + close(store.pointRelease) + <-getDone + re.NoError(gotErr) + re.NotNil(gotGroup, "an unrelated delete must not fail the lazy load") + re.Equal("group-a", gotGroup.Name) + + // After loading completes, group-a is present and group-b stays deleted. + store.unblock() + testutil.Eventually(re, func() bool { + groups, err := m.GetResourceGroupList(1, false) + if err != nil { + return false + } + foundA := false + for _, g := range groups { + if g.Name == "group-b" { + return false + } + if g.Name == "group-a" { + foundA = true + } + } + return foundA + }, testutil.WithTickInterval(20*time.Millisecond)) +} From 57475ed3e2669ab94d02505472dc15f441d294b9 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Wed, 22 Jul 2026 15:46:17 +0800 Subject: [PATCH 19/21] resource_group: guard token bucket state writes with the group lock SetStatesIntoResourceGroup wrote the token bucket fields (Tokens, LastUpdate, Initialized) via setState with no lock, while RequestRU mutates the same fields under the group lock. Both the async bulk merge's state adoption and the metadata watcher's runtime state sync call SetStatesIntoResourceGroup on groups that are already serving token requests, racing with them. Take the group lock around setState; UpdateRUConsumption already locks internally. The lock order stays keyspace-manager lock -> group lock, consistent with every other path (resource_group.go and token_buckets.go never reference the keyspace manager). Signed-off-by: tongjian Signed-off-by: tongjian <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/resource_group.go | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/pkg/mcs/resourcemanager/server/resource_group.go b/pkg/mcs/resourcemanager/server/resource_group.go index 37fdd8bdb85..f1ab18adfdb 100644 --- a/pkg/mcs/resourcemanager/server/resource_group.go +++ b/pkg/mcs/resourcemanager/server/resource_group.go @@ -379,7 +379,13 @@ func (rg *ResourceGroup) SetStatesIntoResourceGroup(states *GroupStates) { switch rg.Mode { case rmpb.GroupMode_RUMode: if state := states.RU; state != nil { + // The group may already be serving requests: setState writes the + // token bucket fields directly, so guard it with the group lock + // the same way RequestRU does. UpdateRUConsumption below locks + // internally. + rg.Lock() rg.RUSettings.RU.setState(state) + rg.Unlock() log.Debug("update group token bucket state", zap.String("name", rg.Name), zap.Any("state", state)) } if states.RUConsumption != nil { From 78e9e69389534883c96c9048eadef248d10785bf Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Thu, 23 Jul 2026 15:41:12 +0800 Subject: [PATCH 20/21] address comment Signed-off-by: tongjian <1045931706@qq.com> --- .../server/keyspace_manager.go | 45 ++------ pkg/mcs/resourcemanager/server/manager.go | 67 ++++------- .../server/manager_async_test.go | 104 +++++++++--------- 3 files changed, 80 insertions(+), 136 deletions(-) diff --git a/pkg/mcs/resourcemanager/server/keyspace_manager.go b/pkg/mcs/resourcemanager/server/keyspace_manager.go index 15c9e0af82d..5c2d61b4401 100644 --- a/pkg/mcs/resourcemanager/server/keyspace_manager.go +++ b/pkg/mcs/resourcemanager/server/keyspace_manager.go @@ -67,32 +67,12 @@ type consumptionItem struct { isTiFlash bool } -// reservedKind describes why a cached resource group entry is still -// considered unconfirmed. -type reservedKind int - -const ( - // reservedPlaceholder marks an entry whose settings and state are both - // synthetic, e.g. the default group pre-inserted by - // ensureReservedDefaultGroupInCache before async loading runs. The async - // bulk merge may replace such an entry wholesale. - reservedPlaceholder reservedKind = iota - // reservedStateOnly marks an entry whose settings were confirmed from - // storage (and may have been modified and persisted since) but whose - // state failed to load. The async bulk merge must only adopt the scanned - // state into it, never replace its settings with the scan's older copy. - reservedStateOnly -) - type keyspaceResourceGroupManager struct { syncutil.RWMutex groups map[string]*ResourceGroup - // reservedGroups tracks names whose entry in groups is not yet fully - // confirmed by a storage load or a real write, together with how much of - // it is unconfirmed (see reservedKind). It shares the same lock as groups - // so a lazy load can atomically decide whether it's safe to replace the - // entry. - reservedGroups map[string]reservedKind + // reservedGroups tracks names whose entry in groups is still a synthetic + // placeholder and not yet confirmed by a storage load or a real write. + reservedGroups map[string]struct{} groupRUTrackers map[string]*groupRUTracker serviceLimiter *serviceLimiter // deleteGen is bumped under the write lock every time a group is removed @@ -117,7 +97,7 @@ func newKeyspaceResourceGroupManager( } return &keyspaceResourceGroupManager{ groups: make(map[string]*ResourceGroup), - reservedGroups: make(map[string]reservedKind), + reservedGroups: make(map[string]struct{}), groupRUTrackers: make(map[string]*groupRUTracker), keyspaceID: keyspaceID, storage: storage, @@ -240,16 +220,13 @@ func (krgm *keyspaceResourceGroupManager) setRawStatesIntoResourceGroup(name str func (krgm *keyspaceResourceGroupManager) initDefaultResourceGroup() { krgm.RLock() _, ok := krgm.groups[DefaultResourceGroupName] - kind, reserved := krgm.reservedGroups[DefaultResourceGroupName] + _, reserved := krgm.reservedGroups[DefaultResourceGroupName] krgm.RUnlock() // A cached entry only makes initialization unnecessary if it's confirmed - // data, or at least has confirmed settings (reservedStateOnly), which must - // not be overwritten with synthetic ones. A reservedPlaceholder means - // nothing is persisted for the default group (e.g. a fresh store): it must - // still be created and persisted here, otherwise it would stay an - // unconfirmed placeholder forever, with its settings never stored and its - // state persistence permanently skipped. - if ok && (!reserved || kind == reservedStateOnly) { + // data. A reserved placeholder means nothing is persisted for the default + // group (e.g. a fresh store): it must still be created and persisted here, + // otherwise its settings are never stored and state persistence stays skipped. + if ok && !reserved { return } defaultGroup := newDefaultResourceGroup() @@ -270,7 +247,7 @@ func (krgm *keyspaceResourceGroupManager) ensureReservedDefaultGroupInCache() { krgm.Lock() if _, ok := krgm.groups[DefaultResourceGroupName]; !ok { krgm.groups[DefaultResourceGroupName] = defaultGroup - krgm.reservedGroups[DefaultResourceGroupName] = reservedPlaceholder + krgm.reservedGroups[DefaultResourceGroupName] = struct{}{} inserted = true } krgm.Unlock() @@ -298,7 +275,7 @@ func (krgm *keyspaceResourceGroupManager) restoreDefaultResourceGroupFromReserve defaultGroup := newDefaultResourceGroup() krgm.Lock() krgm.groups[DefaultResourceGroupName] = defaultGroup - krgm.reservedGroups[DefaultResourceGroupName] = reservedPlaceholder + krgm.reservedGroups[DefaultResourceGroupName] = struct{}{} krgm.Unlock() krgm.syncBurstabilityWithServiceLimit(defaultGroup) } diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index f41e7a1d9f3..dc574b69f00 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -550,21 +550,6 @@ func (m *Manager) asyncLoadResourceGroups(ctx context.Context, epoch uint64) { if m.syncLoadedGroups[key] { continue } - if kind, reserved := krgm.reservedGroups[name]; reserved && kind == reservedStateOnly { - if existing, ok := krgm.groups[name]; ok { - // The cached entry's settings are confirmed and may - // carry a modification persisted after this scan - // started; only its state is missing. Adopt the - // scanned state into the existing entry instead of - // replacing it, so the newer settings aren't - // clobbered by the scan's older copy. - existing.SetStatesIntoResourceGroup(group.GetGroupStates()) - delete(krgm.reservedGroups, name) - groupsToSync = append(groupsToSync, existing) - loaded++ - continue - } - } krgm.groups[name] = group // This group is now confirmed, fully-loaded data (settings // and state); it must no longer be treated as an @@ -632,36 +617,33 @@ func (m *Manager) loadKeyspaceResourceGroupsFromStorage() (map[uint32]*keyspaceR return tempKrgms, nil } -// loadResourceGroup loads a single resource group from storage. The returned -// stateLoaded reports whether the group's persisted state was successfully -// read; the caller must not mark such a group as sync-loaded, so that a -// concurrent or later async bulk load can still fill in its real state. -func (m *Manager) loadResourceGroup(keyspaceID uint32, name string) (group *ResourceGroup, stateLoaded bool, err error) { +// loadResourceGroup loads a single resource group from storage. +func (m *Manager) loadResourceGroup(keyspaceID uint32, name string) (*ResourceGroup, error) { rawValue, err := m.storage.LoadResourceGroupSetting(keyspaceID, name) if err != nil { - return nil, false, err + return nil, err } if rawValue == "" { - return nil, false, errs.ErrResourceGroupNotExists.FastGenByArgs(name) + return nil, errs.ErrResourceGroupNotExists.FastGenByArgs(name) } krgm := newKeyspaceResourceGroupManager(keyspaceID, m.storage, m.writeRole) if err := krgm.addResourceGroupFromRaw(name, rawValue); err != nil { - return nil, false, err + return nil, err } state, err := m.storage.LoadResourceGroupState(keyspaceID, name) if err != nil { - log.Warn("failed to load resource group state, continuing without state", + log.Warn("failed to load resource group state", zap.Uint32("keyspace-id", keyspaceID), zap.String("group-name", name), zap.Error(err)) - return krgm.getMutableResourceGroup(name), false, nil + return nil, err } if state != "" { if err := krgm.setRawStatesIntoResourceGroup(name, state); err != nil { - return nil, false, err + return nil, err } } - return krgm.getMutableResourceGroup(name), true, nil + return krgm.getMutableResourceGroup(name), nil } func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) error { @@ -690,7 +672,7 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro // so a Delete that lands after the read is detected under the insert // lock and can't be undone by the now-stale result. deleteGen := krgm.loadDeleteGen() - group, stateLoaded, err := m.loadResourceGroup(keyspaceID, name) + group, err := m.loadResourceGroup(keyspaceID, name) if err != nil { if name == DefaultResourceGroupName && errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(name)) { // No persisted default group settings exist yet (e.g. a brand-new @@ -704,9 +686,12 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro return err } inserted := false + markKey := trackerKey{keyspaceID: keyspaceID, groupName: name} + m.Lock() krgm.Lock() if krgm.deleteGen != deleteGen { krgm.Unlock() + m.Unlock() if attempt >= maxLoadAttempts { return nil } @@ -721,30 +706,16 @@ func (m *Manager) loadResourceGroupIfNeeded(keyspaceID uint32, name string) erro krgm.groups[name] = group inserted = true } - // Only clear the placeholder mark once the state was actually read; a - // metadata-only group (state load failed) must stay reserved so a later - // call or the async bulk merge remains free to fill in the real state, - // instead of this partial result being treated as final forever. - if stateLoaded { - delete(krgm.reservedGroups, name) - } else { - krgm.reservedGroups[name] = reservedStateOnly - } + delete(krgm.reservedGroups, name) krgm.Unlock() + if m.syncLoadedGroups != nil { + m.syncLoadedGroups[markKey] = true + } + m.Unlock() + failpoint.Inject("lazyLoadAfterCachePublish", func() {}) if inserted { krgm.syncBurstabilityWithServiceLimit(group) } - // Only mark the group as sync-loaded when its persisted state was actually - // read; otherwise a later async bulk load must remain free to fill in the - // real state instead of being skipped forever. - if stateLoaded { - markKey := trackerKey{keyspaceID: keyspaceID, groupName: name} - m.Lock() - if m.syncLoadedGroups != nil { - m.syncLoadedGroups[markKey] = true - } - m.Unlock() - } syncLoadGroupCounter.Inc() return nil } diff --git a/pkg/mcs/resourcemanager/server/manager_async_test.go b/pkg/mcs/resourcemanager/server/manager_async_test.go index 7dd44f521cf..a59f7da05b7 100644 --- a/pkg/mcs/resourcemanager/server/manager_async_test.go +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -24,6 +24,7 @@ import ( "github.com/stretchr/testify/require" + "github.com/pingcap/failpoint" "github.com/pingcap/kvproto/pkg/resource_manager" "github.com/tikv/pd/pkg/errs" @@ -248,14 +249,11 @@ func TestAsyncLoadResourceGroupsLazyGetLegacyKeyspace(t *testing.T) { }, testutil.WithTickInterval(20*time.Millisecond)) } -// TestAsyncLoadResourceGroupsRecoversFromStateLoadFailure guards against a -// group getting stuck marked reserved forever after a transient -// LoadResourceGroupState failure during lazy loading: once the async bulk -// load subsequently installs the fully-loaded (settings and state) -// confirmed data for the same group, the reserved marker must be cleared, -// otherwise loadResourceGroupIfNeeded and the state persist loop would keep -// treating already-recovered, correct data as an unconfirmed placeholder. -func TestAsyncLoadResourceGroupsRecoversFromStateLoadFailure(t *testing.T) { +// TestAsyncLoadResourceGroupsDoesNotServeStateLoadFailure guards against a +// group with confirmed settings but failed state loading being exposed with a +// fresh token bucket before the async bulk loader can recover its persisted +// state. +func TestAsyncLoadResourceGroupsDoesNotServeStateLoadFailure(t *testing.T) { re := require.New(t) store := newBlockingResourceGroupStorage() group := newAsyncTestGroup("flaky-group") @@ -270,22 +268,22 @@ func TestAsyncLoadResourceGroupsRecoversFromStateLoadFailure(t *testing.T) { store.waitEntered(t) - // Make the lazy load's own state read fail once, so the group is cached - // as a metadata-only, still-reserved entry. + // Make the lazy load's own state read fail once. The group must remain + // unavailable rather than being returned with a fresh token bucket state. store.failNextState.Store(true) fetched, err := m.GetResourceGroup(1, "flaky-group", false) - re.NoError(err) - re.NotNil(fetched) + re.Error(err) + re.Nil(fetched) krgm := m.getKeyspaceResourceGroupManager(1) re.NotNil(krgm) - re.True(krgm.isReserved("flaky-group"), "group should still be reserved after a failed state load") + re.Nil(krgm.getMutableResourceGroup("flaky-group"), "failed state load must not publish the group") // Let the async bulk load proceed; its own state read is unaffected // (failNextState was already consumed) and should install confirmed data. store.unblock() testutil.Eventually(re, func() bool { - return !krgm.isReserved("flaky-group") + return krgm.getMutableResourceGroup("flaky-group") != nil && !krgm.isReserved("flaky-group") }, testutil.WithTickInterval(20*time.Millisecond)) } @@ -424,22 +422,16 @@ func TestAsyncLoadResourceGroupsStaleLoaderDoesNotPolluteNewTerm(t *testing.T) { } } -// TestAsyncLoadResourceGroupsMergeKeepsModifiedSettings reproduces the -// stale-settings clobber: a group's lazy load reads its settings but fails to -// read its state, then the group is modified (and the modification persisted) -// while the async bulk scan still holds the pre-modification settings. The -// bulk merge must adopt only the scanned state into the cached entry, not -// replace it wholesale, so the modified settings survive in the serving cache. -func TestAsyncLoadResourceGroupsMergeKeepsModifiedSettings(t *testing.T) { +// TestAsyncLoadResourceGroupsLazyPublishAndMarkAreAtomic reproduces the race +// where a lazy load publishes a cache entry before recording it in +// syncLoadedGroups. A bulk merge entering that gap can overwrite mutable state +// updated by token-bucket handling with its older scan result. +func TestAsyncLoadResourceGroupsLazyPublishAndMarkAreAtomic(t *testing.T) { re := require.New(t) store := newBlockingResourceGroupStorage() - group := newAsyncTestGroup("mod-group") - re.NoError(store.SaveResourceGroupSetting(1, "mod-group", group)) - // Persist a state with a recognizable consumption so the test can tell - // that the merge really adopted the scanned state. - states := FromProtoResourceGroup(group).GetGroupStates() - states.RUConsumption.RRU = 123 - re.NoError(store.SaveResourceGroupStates(1, "mod-group", states)) + group := newAsyncTestGroup("atomic-group") + re.NoError(store.SaveResourceGroupSetting(1, "atomic-group", group)) + re.NoError(store.SaveResourceGroupStates(1, "atomic-group", FromProtoResourceGroup(group).GetGroupStates())) m := NewManager[*mockConfigProvider](&mockConfigProvider{}) m.storage = store @@ -450,8 +442,7 @@ func TestAsyncLoadResourceGroupsMergeKeepsModifiedSettings(t *testing.T) { store.waitEntered(t) - // Let the bulk loader capture the pre-modification settings, then hold it - // right before its states scan (i.e. before it merges). + // Let the bulk loader capture fill rate 100, then hold it before merge. store.pauseNextStates.Store(true) store.unblock() select { @@ -460,40 +451,45 @@ func TestAsyncLoadResourceGroupsMergeKeepsModifiedSettings(t *testing.T) { t.Fatal("timed out waiting for the bulk loader to reach its states scan") } - // Lazily load the group with a failing state read: it is cached with - // confirmed settings but unconfirmed (fresh) state. - store.failNextState.Store(true) - fetched, err := m.GetResourceGroup(1, "mod-group", false) - re.NoError(err) - re.NotNil(fetched) - krgm := m.getKeyspaceResourceGroupManager(1) - re.NotNil(krgm) - re.True(krgm.isReserved("mod-group")) + re.NoError(failpoint.Enable("github.com/tikv/pd/pkg/mcs/resourcemanager/server/lazyLoadAfterCachePublish", `pause`)) + defer func() { + re.NoError(failpoint.Disable("github.com/tikv/pd/pkg/mcs/resourcemanager/server/lazyLoadAfterCachePublish")) + }() - // Modify the group (fill rate 100 -> 200) and persist it. The state read - // of Modify's own lazy load fails again, so the entry stays state-only - // reserved and out of syncLoadedGroups. - store.failNextState.Store(true) - modified := newAsyncTestGroup("mod-group") - modified.RUSettings.RU.Settings.FillRate = 200 - modified.KeyspaceId = &resource_manager.KeyspaceIDValue{Value: 1} - re.NoError(m.ModifyResourceGroup(modified)) + getDone := make(chan struct{}) + go func() { + defer close(getDone) + _, err := m.GetResourceGroup(1, "atomic-group", false) + re.NoError(err) + }() - // Release the bulk loader; its merge must keep the modified settings and - // only adopt the scanned state. + var krgm *keyspaceResourceGroupManager + testutil.Eventually(re, func() bool { + krgm = m.getKeyspaceResourceGroupManager(1) + if krgm == nil { + return false + } + return krgm.getMutableResourceGroup("atomic-group") != nil + }, testutil.WithTickInterval(20*time.Millisecond)) + + krgm.getMutableResourceGroup("atomic-group").UpdateRUConsumption(&resource_manager.Consumption{RRU: 10}) + + // With the lazy load paused at the cache-publish hook, the bulk merge must + // not be able to overwrite the updated cache entry. store.unblockStates() testutil.Eventually(re, func() bool { _, err := m.GetResourceGroupList(1, false) return err == nil }, testutil.WithTickInterval(20*time.Millisecond)) - got, err := m.GetResourceGroup(1, "mod-group", false) + re.NoError(failpoint.Disable("github.com/tikv/pd/pkg/mcs/resourcemanager/server/lazyLoadAfterCachePublish")) + <-getDone + + got, err := m.GetResourceGroup(1, "atomic-group", false) re.NoError(err) re.NotNil(got) - re.Equal(float64(200), got.RUSettings.RU.getFillRate(), "modified settings must survive the bulk merge") - re.False(krgm.isReserved("mod-group"), "state adoption must clear the reserved marker") - re.Equal(float64(123), krgm.getMutableResourceGroup("mod-group").GetGroupStates().RUConsumption.RRU, - "the scanned state must be adopted into the cached entry") + re.Equal(float64(10), krgm.getMutableResourceGroup("atomic-group").GetGroupStates().RUConsumption.RRU, + "bulk merge must not overwrite the published lazy-loaded group") } // TestAsyncLoadResourceGroupsFreshStoreDefaultPersisted guards against the From 3766698423d00c2f7190572a3d245c18a89d34b0 Mon Sep 17 00:00:00 2001 From: tongjian <1045931706@qq.com> Date: Thu, 23 Jul 2026 17:00:04 +0800 Subject: [PATCH 21/21] resource_group: use warn level for retryable bulk load failures The async loader retries indefinitely until the scan succeeds, so a failed attempt is not a terminal condition and warn level fits better. This also satisfies the error-log-review check on new error-level logs. Signed-off-by: tongjian <1045931706@qq.com> --- pkg/mcs/resourcemanager/server/manager.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index dc574b69f00..aff2700d3fc 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -518,7 +518,8 @@ func (m *Manager) asyncLoadResourceGroups(ctx context.Context, epoch uint64) { default: } if err != nil { - log.Error("failed to load resource groups", zap.Error(err), zap.Int("retry", retry)) + // Use warn level since the loader retries indefinitely until it succeeds. + log.Warn("failed to load resource groups", zap.Error(err), zap.Int("retry", retry)) if !m.storeLoadingStateIfCurrent(epoch, LoadingStateNotStarted) { log.Info("async loading resource groups aborted: manager was reinitialized") return