diff --git a/errors.toml b/errors.toml index 33acad6fbf..43413f79a0 100644 --- a/errors.toml +++ b/errors.toml @@ -921,6 +921,11 @@ error = ''' keyspace not found with name: %s ''' +["PD:resourcemanager:ErrResourceGroupsLoading"] +error = ''' +resource groups are still being loaded, please try again later +''' + ["PD:scatter:ErrEmptyRegion"] error = ''' empty region diff --git a/pkg/errs/errno.go b/pkg/errs/errno.go index a887d539f6..9ddb0b7b40 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/grpc_service.go b/pkg/mcs/resourcemanager/server/grpc_service.go index 875299c602..0945494cdc 100644 --- a/pkg/mcs/resourcemanager/server/grpc_service.go +++ b/pkg/mcs/resourcemanager/server/grpc_service.go @@ -231,18 +231,23 @@ 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 { + if err != nil && !errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(resourceGroupName)) { + return err + } + 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/keyspace_manager.go b/pkg/mcs/resourcemanager/server/keyspace_manager.go index d3d8a13dae..5c2d61b440 100644 --- a/pkg/mcs/resourcemanager/server/keyspace_manager.go +++ b/pkg/mcs/resourcemanager/server/keyspace_manager.go @@ -69,9 +69,17 @@ 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 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 + // 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 @@ -89,6 +97,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 +168,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 +178,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,9 +188,20 @@ func (krgm *keyspaceResourceGroupManager) deleteResourceGroupFromCache(name stri krgm.Lock() 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 { @@ -196,8 +220,13 @@ func (krgm *keyspaceResourceGroupManager) setRawStatesIntoResourceGroup(name str func (krgm *keyspaceResourceGroupManager) initDefaultResourceGroup() { krgm.RLock() _, ok := krgm.groups[DefaultResourceGroupName] + _, reserved := krgm.reservedGroups[DefaultResourceGroupName] krgm.RUnlock() - if ok { + // A cached entry only makes initialization unnecessary if it's confirmed + // 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() @@ -218,6 +247,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 +275,7 @@ func (krgm *keyspaceResourceGroupManager) restoreDefaultResourceGroupFromReserve defaultGroup := newDefaultResourceGroup() krgm.Lock() krgm.groups[DefaultResourceGroupName] = defaultGroup + krgm.reservedGroups[DefaultResourceGroupName] = struct{}{} krgm.Unlock() krgm.syncBurstabilityWithServiceLimit(defaultGroup) } @@ -266,6 +297,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,6 +320,9 @@ func (krgm *keyspaceResourceGroupManager) modifyResourceGroup(group *rmpb.Resour if err != nil { return err } + // 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) } @@ -330,6 +365,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() @@ -384,6 +430,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 4f811a0392..a811b6d042 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,28 @@ 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 + // 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 +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 +170,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 +217,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 @@ -275,6 +301,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] @@ -283,9 +316,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 } @@ -325,13 +364,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. @@ -356,47 +398,210 @@ 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 } } 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) + m.loadEpoch++ + epoch := m.loadEpoch + atomic.StoreInt32(&m.loadingState, LoadingStateNotStarted) + m.Unlock() + + m.initReservedInCache() + if err := m.loadServiceLimits(); err != nil { + return err + } + + m.wg.Add(1) + go m.asyncLoadResourceGroups(ctx, epoch) + 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() +} + +// 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() + + 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: + } + } + + 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 { + // 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 + } + 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 { + 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] { + 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() + for _, group := range groupsToSync { + krgm.syncBurstabilityWithServiceLimit(group) + } + } + m.syncLoadedGroups = nil + m.Unlock() + + m.initReserved() + 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)) + 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 +613,164 @@ func (m *Manager) loadKeyspaceResourceGroups() error { zap.Uint32("keyspace-id", keyspaceID), zap.String("group-name", name), zap.Error(err)) } }); err != nil { - return err + return nil, 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() + return tempKrgms, nil } -func (m *Manager) loadServiceLimits() error { - return m.storage.LoadServiceLimits(func(keyspaceID uint32, serviceLimit float64) { - m.getOrCreateKeyspaceResourceGroupManager(keyspaceID, false).setServiceLimitFromStorage(serviceLimit) - }) +// 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, 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 { + log.Warn("failed to load resource group state", + zap.Uint32("keyspace-id", keyspaceID), + zap.String("group-name", name), + zap.Error(err)) + return nil, err + } + if 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 { + // 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 + } + } + // The lock-free storage read below can be invalidated while it runs: a + // concurrent Delete of any group in the keyspace bumps deleteGen, and a + // leadership change replaces m.krgms and m.syncLoadedGroups (bumping + // loadEpoch). Both are rare: retry the read a few times against freshly + // captured state so neither makes this request spuriously fail or publish + // into a detached manager; 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++ { + // Capture the load epoch and the current keyspace manager atomically, + // and re-capture them on every retry: publishing into a previous + // term's detached manager while marking the new term's map would make + // the new bulk merge skip a group its cache doesn't contain. + m.Lock() + epoch := m.loadEpoch + krgm = m.krgms[keyspaceID] + if krgm == nil { + krgm = newKeyspaceResourceGroupManager(keyspaceID, m.storage, m.writeRole) + m.krgms[keyspaceID] = krgm + } + m.Unlock() + // 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, err := m.loadResourceGroup(keyspaceID, name) + if err != nil { + if name == DefaultResourceGroupName && errors.ErrorEqual(err, errs.ErrResourceGroupNotExists.FastGenByArgs(name)) { + m.RLock() + stale := m.loadEpoch != epoch || m.krgms[keyspaceID] != krgm + m.RUnlock() + if stale { + if attempt >= maxLoadAttempts { + return nil + } + continue + } + // 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 + markKey := trackerKey{keyspaceID: keyspaceID, groupName: name} + m.Lock() + if m.loadEpoch != epoch || m.krgms[keyspaceID] != krgm { + // The manager was reinitialized for a new term while the storage + // read was in flight; retry against the new term's state. + m.Unlock() + if attempt >= maxLoadAttempts { + return nil + } + continue + } + krgm.Lock() + if krgm.deleteGen != deleteGen { + krgm.Unlock() + m.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 + } + 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) + } + syncLoadGroupCounter.Inc() + return nil + } +} + +// markResourceGroupSyncLoaded records that the group in krgm was written or +// fully loaded synchronously. The caller passes the keyspace manager it +// actually mutated: if that manager is no longer the live one (the manager was +// reinitialized for a new term while the caller was blocked on storage I/O), +// the marker is skipped, since marking the new term's map for a group its +// cache doesn't contain would make the bulk merge skip loading it. +func (m *Manager) markResourceGroupSyncLoaded(keyspaceID uint32, krgm *keyspaceResourceGroupManager, name string) { + m.Lock() + defer m.Unlock() + if m.krgms[keyspaceID] != krgm { + return + } + 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 +807,7 @@ func (m *Manager) applyResourceGroupSettingFromRaw(keyspaceID uint32, name, rawV zap.Error(err)) return err } + m.markResourceGroupSyncLoaded(keyspaceID, krgm, name) return nil } @@ -492,6 +844,7 @@ func (m *Manager) applyResourceGroupStatesFromRaw(keyspaceID uint32, name, rawVa zap.Error(err)) return err } + m.markResourceGroupSyncLoaded(keyspaceID, krgm, name) return nil } @@ -504,6 +857,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() { @@ -574,7 +936,16 @@ 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 && + !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 + } + m.markResourceGroupSyncLoaded(keyspaceID, krgm, grouppb.Name) + return nil } // ModifyResourceGroup modifies an existing resource group. @@ -583,11 +954,25 @@ 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 + } + // 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, krgm, grouppb.Name) + } + return nil } // DeleteResourceGroup deletes a resource group. @@ -595,16 +980,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, krgm, 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 +1011,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 +1023,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 0000000000..6304e19cc8 --- /dev/null +++ b/pkg/mcs/resourcemanager/server/manager_async_test.go @@ -0,0 +1,664 @@ +// 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" + "errors" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/pingcap/failpoint" + "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" +) + +type blockingResourceGroupStorage struct { + storage.Storage + + once sync.Once + 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 + + // 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{} + + // 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{}), + statesReached: make(chan struct{}), + statesRelease: 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) LoadResourceGroupState(keyspaceID uint32, name string) (string, error) { + 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) +} + +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 { + case <-s.entered: + case <-time.After(time.Second): + t.Fatal("timed out waiting for async resource group loading") + } +} + +func (s *blockingResourceGroupStorage) unblock() { + s.releaseOnce.Do(func() { + close(s.release) + }) +} + +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. +const asyncTestGroupFillRate = 100 + +func newAsyncTestGroup(name string) *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: asyncTestGroupFillRate, + BurstLimit: asyncTestGroupFillRate, + }, + }, + }, + } +} + +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"))) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + 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) + + _, 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(asyncTestGroupFillRate), 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"))) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + 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) + + 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)) +} + +// 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"))) + + 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(asyncTestGroupFillRate), 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)) +} + +// 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") + 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. 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.Error(err) + re.Nil(fetched) + + krgm := m.getKeyspaceResourceGroupManager(1) + re.NotNil(krgm) + 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.getMutableResourceGroup("flaky-group") != nil && !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 + // 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) + 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)) +} + +// 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) + } +} + +// 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("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 + re.NoError(m.Init(context.Background())) + defer stopAsyncTestManager(m) + defer store.unblock() + defer store.unblockStates() + + store.waitEntered(t) + + // Let the bulk loader capture fill rate 100, then hold it before merge. + 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") + } + + 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")) + }() + + getDone := make(chan struct{}) + go func() { + defer close(getDone) + _, err := m.GetResourceGroup(1, "atomic-group", false) + re.NoError(err) + }() + + 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)) + + 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(10), krgm.getMutableResourceGroup("atomic-group").GetGroupStates().RUConsumption.RRU, + "bulk merge must not overwrite the published lazy-loaded group") +} + +// 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)) +} + +// TestAsyncLoadResourceGroupsStaleLazyLoadRetriesNewTerm reproduces the +// cross-term lazy load: the load captures the keyspace manager, then blocks in +// its storage read while the leadership changes and Init replaces m.krgms and +// syncLoadedGroups. On resume it must not publish into the detached old +// manager while marking the group in the new term's map (which would make the +// new bulk merge skip a group its cache doesn't contain); instead it retries +// against the freshly captured state and publishes into the new term. +func TestAsyncLoadResourceGroupsStaleLazyLoadRetriesNewTerm(t *testing.T) { + re := require.New(t) + store := newBlockingResourceGroupStorage() + re.NoError(store.SaveResourceGroupSetting(1, "cross-term", newAsyncTestGroup("cross-term"))) + + m := NewManager[*mockConfigProvider](&mockConfigProvider{}) + m.storage = store + re.NoError(m.Init(context.Background())) + cancelTerm1 := m.cancel + defer stopAsyncTestManager(m) + defer store.unblock() + + store.waitEntered(t) + + // The term-1 lazy load pauses inside its state read, holding the old + // term's keyspace manager. + store.pausePointState.Store(true) + var ( + gotGroup *ResourceGroup + gotErr error + ) + getDone := make(chan struct{}) + go func() { + defer close(getDone) + gotGroup, gotErr = m.GetResourceGroup(1, "cross-term", false) + }() + select { + case <-store.pointReached: + case <-time.After(time.Second): + t.Fatal("timed out waiting for the lazy load to reach its state read") + } + + // Leadership changes: reinitialize the manager for term 2 while the + // term-1 lazy load is still blocked. Term 2's bulk loader parks on the + // same settings-scan block until store.unblock(). + cancelTerm1() + re.NoError(m.Init(context.Background())) + + // Release the stale lazy load: it must detect the term change, retry, and + // publish into the new term's manager, so the request still succeeds. + close(store.pointRelease) + <-getDone + re.NoError(gotErr) + re.NotNil(gotGroup, "the cross-term lazy load must retry and succeed against the new term") + re.Equal("cross-term", gotGroup.Name) + + // Finish loading; the group must remain present after the term-2 bulk + // merge (it was correctly marked in the same term it was published in). + store.unblock() + testutil.Eventually(re, func() bool { + g, err := m.GetResourceGroup(1, "cross-term", false) + return err == nil && g != nil + }, testutil.WithTickInterval(20*time.Millisecond)) +} diff --git a/pkg/mcs/resourcemanager/server/manager_test.go b/pkg/mcs/resourcemanager/server/manager_test.go index 16ecac456d..c045597240 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 2a7fd14482..3bdac0f9fd 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_groups_total", + Help: "Total number of resource groups loaded synchronously.", + }) + + asyncLoadGroupDuration = prometheus.NewHistogram( + prometheus.HistogramOpts{ + Namespace: namespace, + Subsystem: serverSubsystem, + Name: "async_load_group_duration_seconds", + Help: "Duration of asynchronous resource group loading in seconds.", + }) ) 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/mcs/resourcemanager/server/resource_group.go b/pkg/mcs/resourcemanager/server/resource_group.go index 37fdd8bdb8..f1ab18adfd 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 { diff --git a/pkg/storage/endpoint/resource_group.go b/pkg/storage/endpoint/resource_group.go index 68a6297772..a54eb7a271 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 d30a5cb372..3d95b2d7a2 100644 --- a/tests/integrations/mcs/resourcemanager/resource_manager_test.go +++ b/tests/integrations/mcs/resourcemanager/resource_manager_test.go @@ -22,6 +22,7 @@ import ( "math/rand/v2" "net/http" "reflect" + "sort" "strconv" "strings" "sync" @@ -227,6 +228,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{ { @@ -305,6 +307,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 @@ -647,6 +656,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() { @@ -727,16 +737,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")) @@ -1599,13 +1614,40 @@ 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 - }, testutil.WithWaitFor(time.Second)) - re.Equal(groups, newGroups) + return err == nil && reflect.DeepEqual(expectedGroups, normalizeResourceGroupsForSettingsCompare(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 := typeutil.DeepClone(group, func() *rmpb.ResourceGroup { + return &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() {