Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
23a9e24
resource_group: implement async loading for resource groups to improv…
disksing Sep 16, 2025
966d279
tests: wait for resource group reload after restart
bufferflies Jun 11, 2026
fae9a3e
resource_group: avoid overwriting default during async load
bufferflies Jun 11, 2026
f7cf2cd
tests: avoid direct proto dependency in RM integration
bufferflies Jun 11, 2026
a20e813
resource_group: fix lazy-load ordering and default-group review findings
bufferflies Jul 8, 2026
cbe73bc
resource_group: propagate non-not-found errors in AcquireTokenBuckets
bufferflies Jul 8, 2026
b23dd8f
Merge remote-tracking branch 'origin/master' into fix-10873-review-co…
bufferflies Jul 13, 2026
25ec642
resource_group: don't mark lazily-loaded group synced on state load f…
bufferflies Jul 13, 2026
306fd57
resource_group: distinguish reserved default placeholder from confirm…
bufferflies Jul 14, 2026
042d6a7
resource_group: keep placeholder state out of the persist loop and re…
bufferflies Jul 14, 2026
d2cddc9
resource_group: route default-group creation through the safe load path
bufferflies Jul 15, 2026
8f743d9
resource_group: don't let Modify confirm a state-unconfirmed group
bufferflies Jul 15, 2026
eb3e774
resource_group: add lazy-load coverage for legacy null-keyspace groups
bufferflies Jul 15, 2026
3e6fb18
resource_group: clear reserved marker when async merge installs confi…
bufferflies Jul 16, 2026
fdb5ef4
resource_group: fix unparam lint in async loading tests
bufferflies Jul 20, 2026
267b077
resource_group: don't let a racing lazy load resurrect a deleted group
bufferflies Jul 20, 2026
a4e1729
resource_group: prevent a stale async loader from polluting a new term
bufferflies Jul 22, 2026
1bd5ac8
resource_group: don't let the bulk merge clobber modified settings
bufferflies Jul 22, 2026
f3a6069
resource_group: persist fresh-store default and retry lazy loads on d…
bufferflies Jul 22, 2026
2bc385f
resource_group: guard token bucket state writes with the group lock
bufferflies Jul 22, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions errors.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions pkg/errs/errno.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
17 changes: 11 additions & 6 deletions pkg/mcs/resourcemanager/server/grpc_service.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
// 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 {
Expand Down
81 changes: 79 additions & 2 deletions pkg/mcs/resourcemanager/server/keyspace_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,11 +67,39 @@ 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
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
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
Expand All @@ -89,6 +117,7 @@ func newKeyspaceResourceGroupManager(
}
return &keyspaceResourceGroupManager{
groups: make(map[string]*ResourceGroup),
reservedGroups: make(map[string]reservedKind),
groupRUTrackers: make(map[string]*groupRUTracker),
keyspaceID: keyspaceID,
storage: storage,
Expand Down Expand Up @@ -159,13 +188,17 @@ 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
}

resourceGroup := FromProtoResourceGroup(group)
krgm.Lock()
krgm.groups[group.Name] = resourceGroup
delete(krgm.reservedGroups, group.Name)
krgm.Unlock()
krgm.syncBurstabilityWithServiceLimit(resourceGroup)
return nil
Expand All @@ -175,9 +208,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 {
Expand All @@ -196,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()
Expand All @@ -218,6 +270,7 @@ func (krgm *keyspaceResourceGroupManager) ensureReservedDefaultGroupInCache() {
krgm.Lock()
if _, ok := krgm.groups[DefaultResourceGroupName]; !ok {
krgm.groups[DefaultResourceGroupName] = defaultGroup
krgm.reservedGroups[DefaultResourceGroupName] = reservedPlaceholder
inserted = true
}
krgm.Unlock()
Expand Down Expand Up @@ -245,6 +298,7 @@ func (krgm *keyspaceResourceGroupManager) restoreDefaultResourceGroupFromReserve
defaultGroup := newDefaultResourceGroup()
krgm.Lock()
krgm.groups[DefaultResourceGroupName] = defaultGroup
krgm.reservedGroups[DefaultResourceGroupName] = reservedPlaceholder
krgm.Unlock()
krgm.syncBurstabilityWithServiceLimit(defaultGroup)
}
Expand All @@ -266,6 +320,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
Expand All @@ -288,6 +343,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)
}

Expand Down Expand Up @@ -330,6 +388,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()
Expand Down Expand Up @@ -384,6 +453,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",
Expand Down
Loading
Loading