Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
8 changes: 7 additions & 1 deletion pkg/mcs/discovery/discover_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,11 @@ func TestServiceRegistryEntry(t *testing.T) {
re := require.New(t)
_, client, clean := etcdutil.NewTestEtcdCluster(t, 1, nil)
defer clean()
entry1 := &ServiceRegistryEntry{ServiceAddr: "127.0.0.1:1"}
entry1 := &ServiceRegistryEntry{
ServiceAddr: "127.0.0.1:1",
GitHash: "new-build",
Features: map[string]string{"feature-a": "new-build"},
}
s1, err := entry1.Serialize()
re.NoError(err)
sr1 := NewServiceRegister(context.Background(), client, "test_service", "127.0.0.1:1", s1, DefaultLeaseInSeconds)
Expand All @@ -80,10 +84,12 @@ func TestServiceRegistryEntry(t *testing.T) {
err = returnedEntry1.Deserialize([]byte(endpoints[0]))
re.NoError(err)
re.Equal("127.0.0.1:1", returnedEntry1.ServiceAddr)
re.Equal(entry1.Features, returnedEntry1.Features)
returnedEntry2 := &ServiceRegistryEntry{}
err = returnedEntry2.Deserialize([]byte(endpoints[1]))
re.NoError(err)
re.Equal("127.0.0.1:2", returnedEntry2.ServiceAddr)
re.Nil(returnedEntry2.Features)

sr1.cancel()
sr2.cancel()
Expand Down
2 changes: 2 additions & 0 deletions pkg/mcs/discovery/registry_entry.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ type ServiceRegistryEntry struct {
GitHash string `json:"git-hash"`
DeployPath string `json:"deploy-path"`
StartTimestamp int64 `json:"start-timestamp"`
// Features maps a supported feature to the build hash that advertised it.
Features map[string]string `json:"features,omitempty"`
}

// Serialize this service registry entry
Expand Down
1 change: 1 addition & 0 deletions pkg/mcs/scheduling/server/apis/v1/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -362,6 +362,7 @@ func getConfig(c *gin.Context) {
c.String(http.StatusInternalServerError, "failed to decode primary's config: "+err.Error())
return
}
primaryCfg.Schedule.MigrateDeprecatedFlags()

// Schedule and Replication are dynamic configs managed by primary, so we need to merge them.
mergedCfg := localCfg
Expand Down
9 changes: 5 additions & 4 deletions pkg/mcs/scheduling/server/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -271,6 +271,7 @@ func (o *PersistConfig) GetScheduleConfig() *sc.ScheduleConfig {
// SetScheduleConfig sets the scheduling configuration dynamically.
func (o *PersistConfig) SetScheduleConfig(cfg *sc.ScheduleConfig) {
old := o.GetScheduleConfig()
cfg.SyncDefaultStoreLimitCompat()
o.schedule.Store(cfg)
// The coordinator is not aware of the underlying scheduler config changes,
// we should notify it to update the schedulers proactively.
Expand All @@ -289,6 +290,7 @@ func AdjustScheduleCfg(scheduleCfg *sc.ScheduleConfig) {
scheduleCfg.Schedulers = append(scheduleCfg.Schedulers, ps)
}
}
scheduleCfg.MigrateDeprecatedFlags()
}

// GetReplicationConfig returns replication configurations.
Expand Down Expand Up @@ -517,10 +519,7 @@ func (o *PersistConfig) GetStoreLimit(storeID uint64) (returnSC sc.StoreLimitCon
return limit
}
cfg := o.GetScheduleConfig().Clone()
limitCfg := sc.StoreLimitConfig{
AddPeer: sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer),
RemovePeer: sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.RemovePeer),
}
limitCfg := cfg.GetDefaultStoreLimit()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Config entries written by pre-upgrade PD do not contain default-store-limit; the MCS watcher unmarshals those entries into a zero-valued ScheduleConfig, so this fallback installs {0,0} for an unlisted store. In store-limit v1, zero is treated as unlimited, so a new store whose per-store entry has not been persisted can bypass add/remove-peer throttling after an upgrade.

@King-Dylan King-Dylan Aug 13, 2026

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 96d4b03. Persisted JSON now tracks default-store-limit field presence, and the MCS watcher initializes and migrates the watched schedule config before installing it. Legacy configs without the field now use the built-in default (or store-balance-rate when present), while explicit zero values remain unchanged. TestAdjustScheduleConfigDefaultStoreLimit covers these upgrade cases and the future-store lookup.

cfg.StoreLimit[storeID] = limitCfg
o.SetScheduleConfig(cfg)
return o.GetScheduleConfig().StoreLimit[storeID]
Expand Down Expand Up @@ -597,12 +596,14 @@ func (o *PersistConfig) SetAllStoresLimit(typ storelimit.Type, ratePerMin float6
v := o.GetScheduleConfig().Clone()
switch typ {
case storelimit.AddPeer:
v.DefaultStoreLimit.AddPeer = ratePerMin
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, ratePerMin)
for storeID := range v.StoreLimit {
sc := sc.StoreLimitConfig{AddPeer: ratePerMin, RemovePeer: v.StoreLimit[storeID].RemovePeer}
v.StoreLimit[storeID] = sc
}
case storelimit.RemovePeer:
v.DefaultStoreLimit.RemovePeer = ratePerMin
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, ratePerMin)
for storeID := range v.StoreLimit {
sc := sc.StoreLimitConfig{AddPeer: v.StoreLimit[storeID].AddPeer, RemovePeer: ratePerMin}
Expand Down
99 changes: 99 additions & 0 deletions pkg/mcs/scheduling/server/config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,14 @@
package config

import (
"encoding/json"
"testing"
"time"

"github.com/stretchr/testify/require"
"go.uber.org/goleak"

"github.com/tikv/pd/pkg/core/storelimit"
sc "github.com/tikv/pd/pkg/schedule/config"
"github.com/tikv/pd/pkg/schedule/types"
"github.com/tikv/pd/pkg/utils/testutil"
Expand Down Expand Up @@ -88,3 +90,100 @@ func newScheduleConfig(schedulerType types.CheckerSchedulerType) *sc.ScheduleCon
},
}
}

func TestPersistConfigDefaultStoreLimit(t *testing.T) {
re := require.New(t)
oldAddPeer := sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer)
oldRemovePeer := sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.RemovePeer)
defer func() {
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, oldAddPeer)
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, oldRemovePeer)
}()
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, 15)
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, 15)

cfg := NewConfig()
re.NoError(cfg.adjust(nil))
persistConfig := NewPersistConfig(cfg, nil)
persistConfig.GetScheduleConfig().StoreLimit[1] = sc.StoreLimitConfig{AddPeer: 10, RemovePeer: 20}

persistConfig.SetAllStoresLimit(storelimit.AddPeer, 60)
re.Equal(sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 15}, persistConfig.GetScheduleConfig().DefaultStoreLimit)
re.Equal(sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 15},
persistConfig.GetScheduleConfig().StoreLimit[sc.DefaultStoreLimitCompatStoreID])
re.Equal(sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 20}, persistConfig.GetStoreLimit(1))

data, err := json.Marshal(persistConfig.GetScheduleConfig())
re.NoError(err)
var reloadedScheduleConfig sc.ScheduleConfig
re.NoError(json.Unmarshal(data, &reloadedScheduleConfig))
restartedConfig := NewConfig()
restartedConfig.Schedule = reloadedScheduleConfig
restartedPersistConfig := NewPersistConfig(restartedConfig, nil)

sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, 15)
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, 25)
re.Equal(sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 15}, restartedPersistConfig.GetStoreLimit(2))
}

func TestAdjustScheduleConfigDefaultStoreLimit(t *testing.T) {
oldAddPeer := sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer)
oldRemovePeer := sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.RemovePeer)
defer func() {
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, oldAddPeer)
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, oldRemovePeer)
}()

testCases := []struct {
name string
config string
expected sc.StoreLimitConfig
}{
{
name: "legacy config without store limit default",
config: `{"store-limit":{}}`,
expected: sc.StoreLimitConfig{AddPeer: 15, RemovePeer: 15},
},
{
name: "legacy store balance rate",
config: `{"store-balance-rate":60,"store-limit":{}}`,
expected: sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 60},
},
{
name: "explicit zero wins over legacy store balance rate",
config: `{"store-balance-rate":60,"default-store-limit":{"add-peer":0,"remove-peer":0},"store-limit":{}}`,
expected: sc.StoreLimitConfig{AddPeer: 0, RemovePeer: 0},
},
{
name: "legacy store balance rate backfills an omitted field",
config: `{"store-balance-rate":60,"default-store-limit":{"add-peer":0},"store-limit":{}}`,
expected: sc.StoreLimitConfig{AddPeer: 0, RemovePeer: 60},
},
{
name: "compatibility entry survives a pre-feature rewrite",
config: `{"store-limit":{"0":{"add-peer":70,"remove-peer":0}}}`,
expected: sc.StoreLimitConfig{AddPeer: 70, RemovePeer: 0},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
re := require.New(t)
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, 15)
sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, 15)
watchedConfig := &persistedConfig{
Schedule: sc.ScheduleConfig{DefaultStoreLimit: sc.DefaultStoreLimitConfig()},
}
re.NoError(json.Unmarshal([]byte(`{"schedule":`+testCase.config+`}`), watchedConfig))
AdjustScheduleCfg(&watchedConfig.Schedule)
re.Equal(testCase.expected, watchedConfig.Schedule.DefaultStoreLimit)
re.Equal(testCase.expected,
watchedConfig.Schedule.StoreLimit[sc.DefaultStoreLimitCompatStoreID])
re.Zero(watchedConfig.Schedule.StoreBalanceRate)

cfg := NewConfig()
cfg.Schedule = watchedConfig.Schedule
persistConfig := NewPersistConfig(cfg, nil)
re.Equal(testCase.expected, persistConfig.GetStoreLimit(100))
})
}
}
4 changes: 3 additions & 1 deletion pkg/mcs/scheduling/server/config/watcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,9 @@ func (cw *Watcher) getSchedulersController() *schedulers.Controller {

func (cw *Watcher) initializeConfigWatcher() error {
putFn := func(kv *mvccpb.KeyValue) error {
cfg := &persistedConfig{}
cfg := &persistedConfig{
Schedule: sc.ScheduleConfig{DefaultStoreLimit: sc.DefaultStoreLimitConfig()},
}
if err := json.Unmarshal(kv.Value, cfg); err != nil {
log.Warn("failed to unmarshal scheduling config entry",
zap.String("event-kv-key", string(kv.Key)), zap.Error(err))
Expand Down
2 changes: 1 addition & 1 deletion pkg/mcs/scheduling/server/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -627,7 +627,7 @@ func (s *Server) GetPersistConfig() *config.PersistConfig {
// GetConfig gets the config.
func (s *Server) GetConfig() *config.Config {
cfg := s.cfg.Clone()
cfg.Schedule = *s.persistConfig.GetScheduleConfig().Clone()
cfg.Schedule = *s.persistConfig.GetScheduleConfig().CloneWithoutDefaultStoreLimitCompat()
cfg.Replication = *s.persistConfig.GetReplicationConfig().Clone()
cfg.ClusterVersion = *s.persistConfig.GetClusterVersion()
cfg.Schedule.MaxMergeRegionKeys = cfg.Schedule.GetMaxMergeRegionKeys()
Expand Down
5 changes: 5 additions & 0 deletions pkg/mcs/utils/util.go
Original file line number Diff line number Diff line change
Expand Up @@ -291,6 +291,11 @@ func Register(s server, serviceName string) (*discovery.ServiceRegistryEntry, *d
StartTimestamp: s.StartTimestamp(),
Name: s.Name(),
}
if serviceName == constant.SchedulingServiceName {
serviceID.Features = map[string]string{
versioninfo.DefaultStoreLimitPersistence: versioninfo.PDGitHash,
}
}
serializedEntry, err := serviceID.Serialize()
if err != nil {
return nil, nil, err
Expand Down
27 changes: 27 additions & 0 deletions pkg/member/member.go
Original file line number Diff line number Diff line change
Expand Up @@ -482,6 +482,19 @@ func (m *Member) GetMemberGitHash(id uint64) (string, error) {
return string(res.Kvs[0].Value), nil
}

// GetMemberFeature loads the build hash that advertised a member feature.
func (m *Member) GetMemberFeature(id uint64, feature string) (string, error) {
key := keypath.MemberFeaturePath(id, feature)
res, err := etcdutil.EtcdKVGet(m.client, key)
if err != nil {
return "", err
}
if len(res.Kvs) == 0 {
return "", errs.ErrEtcdKVGetResponse.FastGenByArgs("no value")
}
return string(res.Kvs[0].Value), nil
}

// SetMemberBinaryVersion saves a member's binary version.
func (m *Member) SetMemberBinaryVersion(id uint64, releaseVersion string) error {
key := keypath.MemberBinaryVersionPath(id)
Expand Down Expand Up @@ -510,6 +523,20 @@ func (m *Member) SetMemberGitHash(id uint64, gitHash string) error {
return nil
}

// SetMemberFeature advertises a member feature for the current build hash.
func (m *Member) SetMemberFeature(id uint64, feature, gitHash string) error {
key := keypath.MemberFeaturePath(id, feature)
txn := kv.NewSlowLogTxn(m.client)
res, err := txn.Then(clientv3.OpPut(key, gitHash)).Commit()
if err != nil {
return errors.WithStack(err)
}
if !res.Succeeded {
return errors.New("failed to save member feature")
}
return nil
}

// Close gracefully shuts down all servers/listeners.
func (m *Member) Close() {
m.Etcd().Close()
Expand Down
Loading