diff --git a/pkg/mcs/discovery/discover_test.go b/pkg/mcs/discovery/discover_test.go index 3c65af0895..a3179baf76 100644 --- a/pkg/mcs/discovery/discover_test.go +++ b/pkg/mcs/discovery/discover_test.go @@ -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) @@ -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() diff --git a/pkg/mcs/discovery/registry_entry.go b/pkg/mcs/discovery/registry_entry.go index 887a8eb7ae..c942ad4d4a 100644 --- a/pkg/mcs/discovery/registry_entry.go +++ b/pkg/mcs/discovery/registry_entry.go @@ -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 diff --git a/pkg/mcs/scheduling/server/apis/v1/api.go b/pkg/mcs/scheduling/server/apis/v1/api.go index 52bcc39950..58c2adda76 100644 --- a/pkg/mcs/scheduling/server/apis/v1/api.go +++ b/pkg/mcs/scheduling/server/apis/v1/api.go @@ -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 diff --git a/pkg/mcs/scheduling/server/config/config.go b/pkg/mcs/scheduling/server/config/config.go index f7aeb6b4bd..139a87ed7b 100644 --- a/pkg/mcs/scheduling/server/config/config.go +++ b/pkg/mcs/scheduling/server/config/config.go @@ -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. @@ -289,6 +290,7 @@ func AdjustScheduleCfg(scheduleCfg *sc.ScheduleConfig) { scheduleCfg.Schedulers = append(scheduleCfg.Schedulers, ps) } } + scheduleCfg.MigrateDeprecatedFlags() } // GetReplicationConfig returns replication configurations. @@ -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() cfg.StoreLimit[storeID] = limitCfg o.SetScheduleConfig(cfg) return o.GetScheduleConfig().StoreLimit[storeID] @@ -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} diff --git a/pkg/mcs/scheduling/server/config/config_test.go b/pkg/mcs/scheduling/server/config/config_test.go index 9353bce326..eef492b8bf 100644 --- a/pkg/mcs/scheduling/server/config/config_test.go +++ b/pkg/mcs/scheduling/server/config/config_test.go @@ -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" @@ -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)) + }) + } +} diff --git a/pkg/mcs/scheduling/server/config/watcher.go b/pkg/mcs/scheduling/server/config/watcher.go index 9213da4f14..00fefb130d 100644 --- a/pkg/mcs/scheduling/server/config/watcher.go +++ b/pkg/mcs/scheduling/server/config/watcher.go @@ -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)) diff --git a/pkg/mcs/scheduling/server/server.go b/pkg/mcs/scheduling/server/server.go index ce621211a4..48fb14e831 100644 --- a/pkg/mcs/scheduling/server/server.go +++ b/pkg/mcs/scheduling/server/server.go @@ -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() diff --git a/pkg/mcs/utils/util.go b/pkg/mcs/utils/util.go index 4058140400..b91208d367 100644 --- a/pkg/mcs/utils/util.go +++ b/pkg/mcs/utils/util.go @@ -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 diff --git a/pkg/member/member.go b/pkg/member/member.go index e8eaa52d29..93dc97601e 100644 --- a/pkg/member/member.go +++ b/pkg/member/member.go @@ -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) @@ -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() diff --git a/pkg/schedule/config/config.go b/pkg/schedule/config/config.go index eacfa1a484..8b4d7aeeac 100644 --- a/pkg/schedule/config/config.go +++ b/pkg/schedule/config/config.go @@ -15,6 +15,8 @@ package config import ( + "encoding/json" + "math" "time" "github.com/pingcap/errors" @@ -227,6 +229,11 @@ type ScheduleConfig struct { // StoreBalanceRate is the maximum of balance rate for each store. // WARN: StoreBalanceRate is deprecated. StoreBalanceRate float64 `toml:"store-balance-rate" json:"store-balance-rate,omitempty"` + // DefaultStoreLimit is the default limit of scheduling for stores. + DefaultStoreLimit StoreLimitConfig `toml:"default-store-limit" json:"default-store-limit"` + // defaultStoreLimitJSONPresence distinguishes omitted fields in persisted legacy JSON + // from explicitly configured zero values. + defaultStoreLimitJSONPresence *storeLimitConfigJSONPresence // StoreLimit is the limit of scheduling for stores. StoreLimit map[uint64]StoreLimitConfig `toml:"store-limit" json:"store-limit"` // TolerantSizeRatio is the ratio of buffer size for balance scheduler. @@ -337,6 +344,39 @@ type ScheduleConfig struct { MaxAffinityMergeRegionSize uint64 `toml:"max-affinity-merge-region-size" json:"max-affinity-merge-region-size"` } +type storeLimitConfigJSONPresence struct { + addPeer bool + removePeer bool +} + +// DefaultStoreLimitCompatStoreID is reserved for persisting the default store +// limit in the legacy per-store map. Store ID 0 is invalid for a real store, so +// pre-feature PD versions preserve this entry when they rewrite the full config. +const DefaultStoreLimitCompatStoreID uint64 = 0 + +// UnmarshalJSON tracks default-store-limit field presence for legacy config migration. +func (c *ScheduleConfig) UnmarshalJSON(data []byte) error { + type scheduleConfig ScheduleConfig + decoded := scheduleConfig(*c) + if err := json.Unmarshal(data, &decoded); err != nil { + return err + } + var fields struct { + DefaultStoreLimit map[string]json.RawMessage `json:"default-store-limit"` + } + if err := json.Unmarshal(data, &fields); err != nil { + return err + } + _, addPeerDefined := fields.DefaultStoreLimit["add-peer"] + _, removePeerDefined := fields.DefaultStoreLimit["remove-peer"] + *c = ScheduleConfig(decoded) + c.defaultStoreLimitJSONPresence = &storeLimitConfigJSONPresence{ + addPeer: addPeerDefined, + removePeer: removePeerDefined, + } + return nil +} + // Clone returns a cloned scheduling configuration. func (c *ScheduleConfig) Clone() *ScheduleConfig { schedulers := append(c.Schedulers[:0:0], c.Schedulers...) @@ -353,6 +393,23 @@ func (c *ScheduleConfig) Clone() *ScheduleConfig { return &cfg } +// SyncDefaultStoreLimitCompat stores the default in a representation that +// pre-feature PD versions already understand and preserve. +func (c *ScheduleConfig) SyncDefaultStoreLimitCompat() { + if c.StoreLimit == nil { + c.StoreLimit = make(map[uint64]StoreLimitConfig) + } + c.StoreLimit[DefaultStoreLimitCompatStoreID] = c.DefaultStoreLimit +} + +// CloneWithoutDefaultStoreLimitCompat returns a public view of the schedule +// config without the internal compatibility entry. +func (c *ScheduleConfig) CloneWithoutDefaultStoreLimitCompat() *ScheduleConfig { + cfg := c.Clone() + delete(cfg.StoreLimit, DefaultStoreLimitCompatStoreID) + return cfg +} + // Adjust adjusts the config. func (c *ScheduleConfig) Adjust(meta *configutil.ConfigMetaData, reloading bool) error { if !meta.IsDefined("max-snapshot-count") { @@ -457,6 +514,9 @@ func (c *ScheduleConfig) Adjust(meta *configutil.ConfigMetaData, reloading bool) } adjustSchedulers(&c.Schedulers, DefaultSchedulers) + defaultStoreLimitMeta := meta.Child("default-store-limit") + c.migrateStoreBalanceRate(defaultStoreLimitMeta) + c.adjustDefaultStoreLimit(defaultStoreLimitMeta) for k, b := range c.migrateConfigurationMap() { v, err := parseDeprecatedFlag(meta, k, *b[0], *b[1]) @@ -466,11 +526,6 @@ func (c *ScheduleConfig) Adjust(meta *configutil.ConfigMetaData, reloading bool) *b[0], *b[1] = false, v // reset old flag false to make it ignored when marshal to JSON } - if c.StoreBalanceRate != 0 { - DefaultStoreLimit = StoreLimit{AddPeer: c.StoreBalanceRate, RemovePeer: c.StoreBalanceRate} - c.StoreBalanceRate = 0 - } - if c.StoreLimit == nil { c.StoreLimit = make(map[uint64]StoreLimitConfig) } @@ -489,6 +544,59 @@ func (c *ScheduleConfig) Adjust(meta *configutil.ConfigMetaData, reloading bool) return c.Validate() } +func (c *ScheduleConfig) adjustDefaultStoreLimit(meta *configutil.ConfigMetaData) { + defaultStoreLimit := DefaultStoreLimitConfig() + if !meta.IsDefined("add-peer") { + configutil.AdjustFloat64(&c.DefaultStoreLimit.AddPeer, defaultStoreLimit.AddPeer) + } + if !meta.IsDefined("remove-peer") { + configutil.AdjustFloat64(&c.DefaultStoreLimit.RemovePeer, defaultStoreLimit.RemovePeer) + } +} + +func (c *ScheduleConfig) migrateStoreBalanceRate(defaultStoreLimitMeta *configutil.ConfigMetaData) { + if c.StoreBalanceRate == 0 { + return + } + defaultStoreLimit := StoreLimitConfig{AddPeer: c.StoreBalanceRate, RemovePeer: c.StoreBalanceRate} + if !defaultStoreLimitMeta.IsDefined("add-peer") { + c.DefaultStoreLimit.AddPeer = defaultStoreLimit.AddPeer + } + if !defaultStoreLimitMeta.IsDefined("remove-peer") { + c.DefaultStoreLimit.RemovePeer = defaultStoreLimit.RemovePeer + } + DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, c.DefaultStoreLimit.AddPeer) + DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, c.DefaultStoreLimit.RemovePeer) + c.StoreBalanceRate = 0 +} + +func (c *ScheduleConfig) migratePersistedStoreLimit() { + addPeerDefined, removePeerDefined := true, true + if c.defaultStoreLimitJSONPresence != nil { + addPeerDefined = c.defaultStoreLimitJSONPresence.addPeer + removePeerDefined = c.defaultStoreLimitJSONPresence.removePeer + } + + defaultStoreLimit := DefaultStoreLimitConfig() + if c.StoreBalanceRate != 0 { + defaultStoreLimit = StoreLimitConfig{AddPeer: c.StoreBalanceRate, RemovePeer: c.StoreBalanceRate} + } + if compatDefault, ok := c.StoreLimit[DefaultStoreLimitCompatStoreID]; ok { + defaultStoreLimit = compatDefault + } + if !addPeerDefined { + c.DefaultStoreLimit.AddPeer = defaultStoreLimit.AddPeer + } + if !removePeerDefined { + c.DefaultStoreLimit.RemovePeer = defaultStoreLimit.RemovePeer + } + DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, c.DefaultStoreLimit.AddPeer) + DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, c.DefaultStoreLimit.RemovePeer) + c.StoreBalanceRate = 0 + c.defaultStoreLimitJSONPresence = nil + c.SyncDefaultStoreLimitCompat() +} + func (c *ScheduleConfig) migrateConfigurationMap() map[string][2]*bool { return map[string][2]*bool{ "remove-down-replica": {&c.DisableRemoveDownReplica, &c.EnableRemoveDownReplica}, @@ -536,10 +644,7 @@ func parseDeprecatedFlag(meta *configutil.ConfigMetaData, name string, old, new // MigrateDeprecatedFlags updates new flags according to deprecated flags. func (c *ScheduleConfig) MigrateDeprecatedFlags() { c.DisableLearner = false - if c.StoreBalanceRate != 0 { - DefaultStoreLimit = StoreLimit{AddPeer: c.StoreBalanceRate, RemovePeer: c.StoreBalanceRate} - c.StoreBalanceRate = 0 - } + c.migratePersistedStoreLimit() for _, b := range c.migrateConfigurationMap() { // If old=false (previously disabled), set both old and new to false. if *b[0] { @@ -550,6 +655,14 @@ func (c *ScheduleConfig) MigrateDeprecatedFlags() { // Validate is used to validate if some scheduling configurations are right. func (c *ScheduleConfig) Validate() error { + if math.IsNaN(c.DefaultStoreLimit.AddPeer) || math.IsInf(c.DefaultStoreLimit.AddPeer, 0) || + c.DefaultStoreLimit.AddPeer < 0 { + return errors.New("default-store-limit.add-peer should be finite and non-negative") + } + if math.IsNaN(c.DefaultStoreLimit.RemovePeer) || math.IsInf(c.DefaultStoreLimit.RemovePeer, 0) || + c.DefaultStoreLimit.RemovePeer < 0 { + return errors.New("default-store-limit.remove-peer should be finite and non-negative") + } if c.TolerantSizeRatio < 0 { return errors.New("tolerant-size-ratio should be non-negative") } @@ -606,6 +719,19 @@ type StoreLimitConfig struct { RemovePeer float64 `toml:"remove-peer" json:"remove-peer"` } +// DefaultStoreLimitConfig returns the current process default store limit config. +func DefaultStoreLimitConfig() StoreLimitConfig { + return StoreLimitConfig{ + AddPeer: DefaultStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer), + RemovePeer: DefaultStoreLimit.GetDefaultStoreLimit(storelimit.RemovePeer), + } +} + +// GetDefaultStoreLimit returns the default store limit config. +func (c *ScheduleConfig) GetDefaultStoreLimit() StoreLimitConfig { + return c.DefaultStoreLimit +} + // SchedulerConfigs is a slice of customized scheduler configuration. type SchedulerConfigs []SchedulerConfig diff --git a/pkg/statistics/store_collection.go b/pkg/statistics/store_collection.go index aef38b70c3..3ca47beb87 100644 --- a/pkg/statistics/store_collection.go +++ b/pkg/statistics/store_collection.go @@ -283,6 +283,9 @@ func (s *storeStatistics) collect() { } for storeID, limit := range s.opt.GetStoresLimit() { + if storeID == config.DefaultStoreLimitCompatStoreID { + continue + } id := strconv.FormatUint(storeID, 10) StoreLimitGauge.WithLabelValues(id, "add-peer").Set(limit.AddPeer) StoreLimitGauge.WithLabelValues(id, "remove-peer").Set(limit.RemovePeer) diff --git a/pkg/utils/keypath/absolute_key_path.go b/pkg/utils/keypath/absolute_key_path.go index 28ca4e8dfa..a6b90a52e9 100644 --- a/pkg/utils/keypath/absolute_key_path.go +++ b/pkg/utils/keypath/absolute_key_path.go @@ -52,6 +52,7 @@ const ( memberBinaryDeployPathFormat = "/pd/%d/member/%d/deploy_path" // "/pd/{cluster_id}/member/{member_id}/deploy_path" memberGitHashPath = "/pd/%d/member/%d/git_hash" // "/pd/{cluster_id}/member/{member_id}/git_hash" memberBinaryVersionPathFormat = "/pd/%d/member/%d/binary_version" // "/pd/{cluster_id}/member/{member_id}/binary_version" + memberFeaturePathFormat = "/pd/%d/member/%d/feature/%s" // "/pd/{cluster_id}/member/{member_id}/feature/{feature}" memberLeaderPriorityPathFormat = "/pd/%d/member/%d/leader_priority" // "/pd/{cluster_id}/member/{member_id}/leader_priority" rulePathFormat = "/pd/%d/rules/%s" // "/pd/{cluster_id}/rules/{rule_id}" diff --git a/pkg/utils/keypath/member.go b/pkg/utils/keypath/member.go index 670b3ff556..607831f7b3 100644 --- a/pkg/utils/keypath/member.go +++ b/pkg/utils/keypath/member.go @@ -30,3 +30,8 @@ func MemberGitHashPath(id uint64) string { func MemberBinaryVersionPath(id uint64) string { return fmt.Sprintf(memberBinaryVersionPathFormat, ClusterID(), id) } + +// MemberFeaturePath returns the path of a feature advertised by a member. +func MemberFeaturePath(id uint64, feature string) string { + return fmt.Sprintf(memberFeaturePathFormat, ClusterID(), id, feature) +} diff --git a/pkg/versioninfo/feature.go b/pkg/versioninfo/feature.go index 24aed22dc9..81b1eb2adc 100644 --- a/pkg/versioninfo/feature.go +++ b/pkg/versioninfo/feature.go @@ -26,6 +26,13 @@ import ( // Feature supported features. type Feature int +const ( + // DefaultStoreLimitPersistence is advertised by PD and Scheduling Service + // members that preserve schedule.default-store-limit when loading and + // persisting the shared scheduling configuration. + DefaultStoreLimitPersistence = "default-store-limit-persistence" +) + // Features list. // The cluster provides corresponding new features if the cluster version // greater than or equal to the required minimum version of the feature. diff --git a/server/api/config.go b/server/api/config.go index 194761443a..b94a7a79e5 100644 --- a/server/api/config.go +++ b/server/api/config.go @@ -79,7 +79,7 @@ func (h *confHandler) GetConfig(w http.ResponseWriter, r *http.Request) { } mergedCfg := localCfg mergedCfg.Replication = leaderCfg.Replication - mergedCfg.Schedule = leaderCfg.Schedule + mergedCfg.Schedule = *leaderCfg.Schedule.CloneWithoutDefaultStoreLimitCompat() h.rd.JSON(w, http.StatusOK, mergedCfg) return } @@ -91,7 +91,7 @@ func (h *confHandler) GetConfig(w http.ResponseWriter, r *http.Request) { h.rd.JSON(w, http.StatusInternalServerError, err.Error()) return } - cfg.Schedule = schedulingServerConfig.Schedule + cfg.Schedule = *schedulingServerConfig.Schedule.CloneWithoutDefaultStoreLimitCompat() cfg.Replication = schedulingServerConfig.Replication } else { cfg.Schedule.MaxMergeRegionKeys = cfg.Schedule.GetMaxMergeRegionKeys() @@ -418,7 +418,7 @@ func (h *confHandler) GetScheduleConfig(w http.ResponseWriter, r *http.Request) h.rd.JSON(w, http.StatusInternalServerError, err.Error()) return } - h.rd.JSON(w, http.StatusOK, cfg.Schedule) + h.rd.JSON(w, http.StatusOK, cfg.Schedule.CloneWithoutDefaultStoreLimitCompat()) return } cfg := h.svr.GetScheduleConfig() @@ -679,9 +679,7 @@ func (h *confHandler) getLeaderConfig() (*config.Config, error) { if err != nil { return nil, err } - var leaderConfig config.Config - err = json.Unmarshal(b, &leaderConfig) - return &leaderConfig, err + return unmarshalRemoteConfig(b) } func (h *confHandler) getSchedulingServerConfig() (*config.Config, error) { @@ -702,12 +700,16 @@ func (h *confHandler) getSchedulingServerConfig() (*config.Config, error) { if err != nil { return nil, err } - var schedulingServerConfig config.Config - err = json.Unmarshal(b, &schedulingServerConfig) - if err != nil { + return unmarshalRemoteConfig(b) +} + +func unmarshalRemoteConfig(data []byte) (*config.Config, error) { + var cfg config.Config + if err := json.Unmarshal(data, &cfg); err != nil { return nil, err } - return &schedulingServerConfig, nil + cfg.Schedule.MigrateDeprecatedFlags() + return &cfg, nil } func (h *confHandler) updateControllerConfig(key string, value any) error { diff --git a/server/api/config_test.go b/server/api/config_test.go new file mode 100644 index 0000000000..6580e77fd5 --- /dev/null +++ b/server/api/config_test.go @@ -0,0 +1,85 @@ +// 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 api + +import ( + "encoding/json" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/tikv/pd/pkg/core/storelimit" + sc "github.com/tikv/pd/pkg/schedule/config" +) + +func TestUnmarshalRemoteConfigMigratesDefaultStoreLimit(t *testing.T) { + oldAddPeer := sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer) + oldRemovePeer := sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.RemovePeer) + t.Cleanup(func() { + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, oldAddPeer) + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, oldRemovePeer) + }) + + testCases := []struct { + name string + data string + expected sc.StoreLimitConfig + }{ + { + name: "missing default uses process default", + data: `{"schedule":{}}`, + expected: sc.StoreLimitConfig{AddPeer: 15, RemovePeer: 15}, + }, + { + name: "legacy rate backfills missing default", + data: `{"schedule":{"store-balance-rate":60}}`, + expected: sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 60}, + }, + { + name: "explicit zero remains unlimited", + data: `{"schedule":{"store-balance-rate":60,"default-store-limit":{"add-peer":0,"remove-peer":0}}}`, + expected: sc.StoreLimitConfig{AddPeer: 0, RemovePeer: 0}, + }, + { + name: "compatibility entry survives old primary rewrite", + data: `{"schedule":{"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) { + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, 15) + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, 15) + + cfg, err := unmarshalRemoteConfig([]byte(testCase.data)) + require.NoError(t, err) + require.Equal(t, testCase.expected, cfg.Schedule.DefaultStoreLimit) + require.Equal(t, testCase.expected, + cfg.Schedule.StoreLimit[sc.DefaultStoreLimitCompatStoreID]) + require.Zero(t, cfg.Schedule.StoreBalanceRate) + _, exposed := cfg.Schedule.CloneWithoutDefaultStoreLimitCompat().StoreLimit[sc.DefaultStoreLimitCompatStoreID] + require.False(t, exposed) + + data, err := json.Marshal(cfg) + require.NoError(t, err) + var roundTrip struct { + Schedule sc.ScheduleConfig `json:"schedule"` + } + require.NoError(t, json.Unmarshal(data, &roundTrip)) + require.Equal(t, testCase.expected, roundTrip.Schedule.DefaultStoreLimit) + }) + } +} diff --git a/server/cluster/cluster.go b/server/cluster/cluster.go index 80a8a4b689..5b09644543 100644 --- a/server/cluster/cluster.go +++ b/server/cluster/cluster.go @@ -2389,10 +2389,7 @@ func (c *RaftCluster) AddStoreLimit(store *metapb.Store) { return } - slc := sc.StoreLimitConfig{ - AddPeer: sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer), - RemovePeer: sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.RemovePeer), - } + slc := cfg.GetDefaultStoreLimit() if core.IsStoreContainLabel(store, core.EngineKey, core.EngineTiFlash) { slc = sc.StoreLimitConfig{ AddPeer: sc.DefaultTiFlashStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer), @@ -2591,14 +2588,13 @@ func (c *RaftCluster) SetStoreLimit(storeID uint64, typ storelimit.Type, ratePer // SetAllStoresLimit sets all store limit for a given type and rate. func (c *RaftCluster) SetAllStoresLimit(typ storelimit.Type, ratePerMin float64) error { old := c.opt.GetScheduleConfig().Clone() - oldAdd := sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer) - oldRemove := sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.RemovePeer) c.opt.SetAllStoresLimit(typ, ratePerMin) if err := c.opt.Persist(c.storage); err != nil { // roll back the store limit c.opt.SetScheduleConfig(old) - sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, oldAdd) - sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, oldRemove) + oldDefaultStoreLimit := old.GetDefaultStoreLimit() + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, oldDefaultStoreLimit.AddPeer) + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, oldDefaultStoreLimit.RemovePeer) log.Error("persist store limit meet error", errs.ZapError(err)) return err } diff --git a/server/cluster/cluster_test.go b/server/cluster/cluster_test.go index 5a30302c4c..31c0a7a64a 100644 --- a/server/cluster/cluster_test.go +++ b/server/cluster/cluster_test.go @@ -3344,6 +3344,41 @@ func TestStoreLimitChangeRefreshLimiter(t *testing.T) { re.True(store.IsAvailable(storelimit.AddPeer, constant.Low)) } +func TestAddStoreLimitUsesPersistedDefaultStoreLimit(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) + }() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + _, opt, err := newTestScheduleConfig() + re.NoError(err) + rc := newTestRaftCluster(ctx, mockid.NewIDAllocator(), opt, storage.NewStorageWithMemoryBackend()) + opt.SetAllStoresLimit(storelimit.AddPeer, 60) + + // Simulate a restarted process whose package-level default goes back to the built-in value. + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, 15) + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, 15) + + rc.AddStoreLimit(&metapb.Store{Id: 1}) + re.Equal(sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 15}, opt.GetScheduleConfig().StoreLimit[1]) + + opt.SetStoreLimit(2, storelimit.RemovePeer, 80) + rc.AddStoreLimit(&metapb.Store{Id: 2}) + re.Equal(sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 80}, opt.GetScheduleConfig().StoreLimit[2]) + + rc.AddStoreLimit(&metapb.Store{ + Id: 3, + Labels: []*metapb.StoreLabel{{Key: core.EngineKey, Value: core.EngineTiFlash}}, + }) + re.Equal(sc.StoreLimitConfig{AddPeer: 30, RemovePeer: 30}, opt.GetScheduleConfig().StoreLimit[3]) +} + func TestPatrolRegionConcurrency(t *testing.T) { re := require.New(t) diff --git a/server/config/config_test.go b/server/config/config_test.go index 6a9411ef88..8f255da609 100644 --- a/server/config/config_test.go +++ b/server/config/config_test.go @@ -28,6 +28,7 @@ import ( "github.com/stretchr/testify/require" "go.uber.org/goleak" + "github.com/tikv/pd/pkg/core/storelimit" "github.com/tikv/pd/pkg/ratelimit" sc "github.com/tikv/pd/pkg/schedule/config" "github.com/tikv/pd/pkg/storage" @@ -81,6 +82,220 @@ func TestReloadConfig(t *testing.T) { re.Equal(int64(512), newOpt.GetMaxMovableHotPeerSize()) } +func TestReloadDefaultStoreLimit(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) + + opt, err := newTestScheduleOption() + re.NoError(err) + opt.SetAllStoresLimit(storelimit.AddPeer, 60) + re.Equal(sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 15}, opt.GetScheduleConfig().DefaultStoreLimit) + re.Equal(sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 15}, + opt.GetScheduleConfig().StoreLimit[sc.DefaultStoreLimitCompatStoreID]) + + storage := storage.NewStorageWithMemoryBackend() + re.NoError(opt.Persist(storage)) + + // Simulate a restarted process whose package-level default goes back to the built-in value. + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, 15) + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, 15) + newOpt, err := newTestScheduleOption() + re.NoError(err) + re.NoError(newOpt.Reload(storage)) + + expected := sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 15} + re.Equal(expected, newOpt.GetScheduleConfig().DefaultStoreLimit) + re.Equal(expected, newOpt.GetStoreLimit(100)) + + newOpt.SetStoreLimit(101, storelimit.RemovePeer, 70) + re.Equal(sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 70}, newOpt.GetStoreLimit(101)) + + cfg := newOpt.GetScheduleConfig().Clone() + cfg.DefaultStoreLimit.AddPeer = 0 + newOpt.SetScheduleConfig(cfg) + re.NoError(newOpt.Persist(storage)) + + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, 15) + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, 15) + reloadedOpt, err := newTestScheduleOption() + re.NoError(err) + re.NoError(reloadedOpt.Reload(storage)) + re.Equal(sc.StoreLimitConfig{AddPeer: 0, RemovePeer: 15}, reloadedOpt.GetScheduleConfig().DefaultStoreLimit) + re.Equal(sc.StoreLimitConfig{AddPeer: 0, RemovePeer: 15}, reloadedOpt.GetStoreLimit(102)) +} + +func TestReloadDefaultStoreLimitAfterPreFeatureRewrite(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) + + opt, err := newTestScheduleOption() + re.NoError(err) + opt.SetAllStoresLimit(storelimit.AddPeer, 60) + opt.SetAllStoresLimit(storelimit.RemovePeer, 0) + expected := sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 0} + re.Equal(expected, opt.GetScheduleConfig().StoreLimit[sc.DefaultStoreLimitCompatStoreID]) + + // Simulate a pre-feature leader loading and rewriting the full config. It + // drops default-store-limit, but preserves the legacy per-store map. + type preFeatureScheduleConfig struct { + MaxSnapshotCount uint64 `json:"max-snapshot-count"` + StoreLimit map[uint64]sc.StoreLimitConfig `json:"store-limit"` + } + type preFeatureConfig struct { + Schedule preFeatureScheduleConfig `json:"schedule"` + } + storage := storage.NewStorageWithMemoryBackend() + re.NoError(storage.SaveConfig(&preFeatureConfig{ + Schedule: preFeatureScheduleConfig{ + MaxSnapshotCount: 10, + StoreLimit: opt.GetScheduleConfig().Clone().StoreLimit, + }, + })) + + // A new leader restores the public field from the compatibility entry, + // including an explicit zero value. + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.AddPeer, 15) + sc.DefaultStoreLimit.SetDefaultStoreLimit(storelimit.RemovePeer, 15) + reloadedOpt, err := newTestScheduleOption() + re.NoError(err) + re.NoError(reloadedOpt.Reload(storage)) + re.Equal(uint64(10), reloadedOpt.GetMaxSnapshotCount()) + re.Equal(expected, reloadedOpt.GetScheduleConfig().DefaultStoreLimit) + re.Equal(expected, reloadedOpt.GetStoreLimit(100)) + re.Equal(expected, reloadedOpt.GetScheduleConfig().StoreLimit[sc.DefaultStoreLimitCompatStoreID]) + + publicCfg := reloadedOpt.GetScheduleConfig().CloneWithoutDefaultStoreLimitCompat() + _, ok := publicCfg.StoreLimit[sc.DefaultStoreLimitCompatStoreID] + re.False(ok) +} + +func TestDefaultStoreLimitAdjust(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) + + cases := []struct { + name string + config string + expect sc.StoreLimitConfig + }{ + { + name: "preserve explicit zero", + config: ` +[schedule.default-store-limit] +add-peer = 0 +remove-peer = 60 +`, + expect: sc.StoreLimitConfig{AddPeer: 0, RemovePeer: 60}, + }, + { + name: "store balance rate backfills undefined field", + config: ` +[schedule] +store-balance-rate = 50 + +[schedule.default-store-limit] +add-peer = 0 +`, + expect: sc.StoreLimitConfig{AddPeer: 0, RemovePeer: 50}, + }, + { + name: "explicit default store limit wins over store balance rate", + config: ` +[schedule] +store-balance-rate = 50 + +[schedule.default-store-limit] +add-peer = 60 +remove-peer = 70 +`, + expect: sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 70}, + }, + } + for _, testCase := range cases { + t.Run(testCase.name, func(t *testing.T) { + cfg := NewConfig() + meta, err := toml.Decode(testCase.config, cfg) + require.NoError(t, err) + require.NoError(t, cfg.Adjust(&meta, false)) + require.Equal(t, testCase.expect, cfg.Schedule.DefaultStoreLimit) + }) + } + + schedule := &sc.ScheduleConfig{} + re.NoError(json.Unmarshal([]byte(`{"store-balance-rate":50}`), schedule)) + schedule.MigrateDeprecatedFlags() + re.Equal(sc.StoreLimitConfig{AddPeer: 50, RemovePeer: 50}, schedule.DefaultStoreLimit) + + schedule = &sc.ScheduleConfig{} + re.NoError(json.Unmarshal([]byte(`{"store-balance-rate":50,"default-store-limit":{"add-peer":0,"remove-peer":60}}`), schedule)) + schedule.MigrateDeprecatedFlags() + re.Equal(sc.StoreLimitConfig{AddPeer: 0, RemovePeer: 60}, schedule.DefaultStoreLimit) + + schedule = &sc.ScheduleConfig{} + re.NoError(json.Unmarshal([]byte(`{"store-limit":{"0":{"add-peer":70,"remove-peer":0}}}`), schedule)) + schedule.MigrateDeprecatedFlags() + re.Equal(sc.StoreLimitConfig{AddPeer: 70, RemovePeer: 0}, schedule.DefaultStoreLimit) + re.Equal(schedule.DefaultStoreLimit, schedule.StoreLimit[sc.DefaultStoreLimitCompatStoreID]) + + schedule = &sc.ScheduleConfig{} + re.NoError(json.Unmarshal([]byte(`{"default-store-limit":{"add-peer":0,"remove-peer":60},"store-limit":{"0":{"add-peer":70,"remove-peer":80}}}`), schedule)) + schedule.MigrateDeprecatedFlags() + re.Equal(sc.StoreLimitConfig{AddPeer: 0, RemovePeer: 60}, schedule.DefaultStoreLimit) + re.Equal(schedule.DefaultStoreLimit, schedule.StoreLimit[sc.DefaultStoreLimitCompatStoreID]) +} + +func TestReloadLegacyStoreBalanceRate(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) + + type legacyScheduleConfig struct { + StoreBalanceRate float64 `json:"store-balance-rate"` + } + type legacyConfig struct { + Schedule legacyScheduleConfig `json:"schedule"` + } + storage := storage.NewStorageWithMemoryBackend() + re.NoError(storage.SaveConfig(&legacyConfig{ + Schedule: legacyScheduleConfig{StoreBalanceRate: 60}, + })) + + opt, err := newTestScheduleOption() + re.NoError(err) + re.NoError(opt.Reload(storage)) + expected := sc.StoreLimitConfig{AddPeer: 60, RemovePeer: 60} + re.Equal(expected, opt.GetScheduleConfig().DefaultStoreLimit) + re.Equal(expected, opt.GetStoreLimit(100)) + re.Zero(opt.GetScheduleConfig().StoreBalanceRate) +} + func TestReloadUpgrade(t *testing.T) { re := require.New(t) opt, err := newTestScheduleOption() @@ -144,6 +359,17 @@ func TestValidation(t *testing.T) { re.Error(cfg.Schedule.Validate()) cfg.Schedule.LowSpaceRatio = 0.8 re.NoError(cfg.Schedule.Validate()) + cfg.Schedule.DefaultStoreLimit.AddPeer = -1 + re.ErrorContains(cfg.Schedule.Validate(), "default-store-limit.add-peer") + cfg.Schedule.DefaultStoreLimit.AddPeer = math.Inf(1) + re.ErrorContains(cfg.Schedule.Validate(), "default-store-limit.add-peer") + cfg.Schedule.DefaultStoreLimit.AddPeer = 15 + cfg.Schedule.DefaultStoreLimit.RemovePeer = -1 + re.ErrorContains(cfg.Schedule.Validate(), "default-store-limit.remove-peer") + cfg.Schedule.DefaultStoreLimit.RemovePeer = math.NaN() + re.ErrorContains(cfg.Schedule.Validate(), "default-store-limit.remove-peer") + cfg.Schedule.DefaultStoreLimit.RemovePeer = 15 + re.NoError(cfg.Schedule.Validate()) cfg.Schedule.TolerantSizeRatio = -0.6 re.Error(cfg.Schedule.Validate()) // check quota diff --git a/server/config/persist_options.go b/server/config/persist_options.go index 15cda11b13..eb40897dc5 100644 --- a/server/config/persist_options.go +++ b/server/config/persist_options.go @@ -83,6 +83,7 @@ func (o *PersistOptions) GetScheduleConfig() *sc.ScheduleConfig { // SetScheduleConfig sets the PD scheduling configuration. func (o *PersistOptions) SetScheduleConfig(cfg *sc.ScheduleConfig) { + cfg.SyncDefaultStoreLimitCompat() o.schedule.Store(cfg) } @@ -390,19 +391,20 @@ func (o *PersistOptions) SetMaxMergeRegionKeys(maxMergeRegionKeys uint64) { // SetStoreLimit sets a store limit for a given type and rate. func (o *PersistOptions) SetStoreLimit(storeID uint64, typ storelimit.Type, ratePerMin float64) { v := o.GetScheduleConfig().Clone() + defaultStoreLimit := v.GetDefaultStoreLimit() var slc sc.StoreLimitConfig var rate float64 switch typ { case storelimit.AddPeer: if _, ok := v.StoreLimit[storeID]; !ok { - rate = sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.RemovePeer) + rate = defaultStoreLimit.RemovePeer } else { rate = v.StoreLimit[storeID].RemovePeer } slc = sc.StoreLimitConfig{AddPeer: ratePerMin, RemovePeer: rate} case storelimit.RemovePeer: if _, ok := v.StoreLimit[storeID]; !ok { - rate = sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer) + rate = defaultStoreLimit.AddPeer } else { rate = v.StoreLimit[storeID].AddPeer } @@ -417,12 +419,14 @@ func (o *PersistOptions) SetAllStoresLimit(typ storelimit.Type, ratePerMin float 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} @@ -494,10 +498,7 @@ func (o *PersistOptions) GetStoreLimit(storeID uint64) (returnSC sc.StoreLimitCo return limit } cfg := o.GetScheduleConfig().Clone() - limitCfg := sc.StoreLimitConfig{ - AddPeer: sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.AddPeer), - RemovePeer: sc.DefaultStoreLimit.GetDefaultStoreLimit(storelimit.RemovePeer), - } + limitCfg := cfg.GetDefaultStoreLimit() cfg.StoreLimit[storeID] = limitCfg o.SetScheduleConfig(cfg) return o.GetScheduleConfig().StoreLimit[storeID] @@ -789,9 +790,11 @@ func (o *PersistOptions) SwitchRaftV2(storage endpoint.ConfigStorage) error { // Persist saves the configuration to the storage. func (o *PersistOptions) Persist(storage endpoint.ConfigStorage) error { + schedule := o.GetScheduleConfig().Clone() + schedule.SyncDefaultStoreLimitCompat() cfg := &persistedConfig{ Config: &Config{ - Schedule: *o.GetScheduleConfig(), + Schedule: *schedule, Replication: *o.GetReplicationConfig(), PDServerCfg: *o.GetPDServerConfig(), ReplicationMode: *o.GetReplicationModeConfig(), diff --git a/server/handler.go b/server/handler.go index 57e09eb943..aae48f1b99 100644 --- a/server/handler.go +++ b/server/handler.go @@ -255,6 +255,9 @@ func (h *Handler) SetAllStoresLimit(ratePerMin float64, limitType storelimit.Typ if err != nil { return err } + if err := h.s.checkDefaultStoreLimitPersistenceSupport(); err != nil { + return err + } return c.SetAllStoresLimit(limitType, ratePerMin) } diff --git a/server/server.go b/server/server.go index 58e3caa31e..80c2121bc7 100644 --- a/server/server.go +++ b/server/server.go @@ -498,6 +498,10 @@ func (s *Server) startServer(ctx context.Context) error { if err := s.member.SetMemberGitHash(s.member.ID(), versioninfo.PDGitHash); err != nil { return err } + if err := s.member.SetMemberFeature( + s.member.ID(), versioninfo.DefaultStoreLimitPersistence, versioninfo.PDGitHash); err != nil { + return err + } s.idAllocator = id.NewAllocator(&id.AllocatorParams{ Client: s.client, Label: id.DefaultLabel, @@ -1196,7 +1200,7 @@ func (s *Server) GetServiceMiddlewareConfig() *config.ServiceMiddlewareConfig { // GetConfig gets the config information. func (s *Server) GetConfig() *config.Config { cfg := s.cfg.Clone() - cfg.Schedule = *s.persistOptions.GetScheduleConfig().Clone() + cfg.Schedule = *s.persistOptions.GetScheduleConfig().CloneWithoutDefaultStoreLimitCompat() cfg.Replication = *s.persistOptions.GetReplicationConfig().Clone() cfg.PDServerCfg = *s.persistOptions.GetPDServerConfig().Clone() cfg.ReplicationMode = *s.persistOptions.GetReplicationModeConfig() @@ -1274,7 +1278,7 @@ func (s *Server) SetMicroserviceConfig(cfg config.MicroserviceConfig) error { // GetScheduleConfig gets the balance config information. func (s *Server) GetScheduleConfig() *sc.ScheduleConfig { - return s.persistOptions.GetScheduleConfig().Clone() + return s.persistOptions.GetScheduleConfig().CloneWithoutDefaultStoreLimitCompat() } // SetScheduleConfig sets the balance config information. @@ -1287,6 +1291,11 @@ func (s *Server) SetScheduleConfig(cfg sc.ScheduleConfig) error { return err } old := s.persistOptions.GetScheduleConfig() + if cfg.DefaultStoreLimit != old.DefaultStoreLimit { + if err := s.checkDefaultStoreLimitPersistenceSupport(); err != nil { + return err + } + } s.persistOptions.SetScheduleConfig(&cfg) if err := s.persistOptions.Persist(s.storage); err != nil { s.persistOptions.SetScheduleConfig(old) diff --git a/server/store_limit_compatibility_test.go b/server/store_limit_compatibility_test.go new file mode 100644 index 0000000000..267b5bab10 --- /dev/null +++ b/server/store_limit_compatibility_test.go @@ -0,0 +1,32 @@ +// 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 ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestCheckDefaultStoreLimitPersistenceFeature(t *testing.T) { + re := require.New(t) + re.NoError(checkDefaultStoreLimitPersistenceFeature("PD", "pd-1", "new-build", "new-build")) + re.ErrorContains( + checkDefaultStoreLimitPersistenceFeature("PD", "pd-1", "new-build", ""), + "PD member pd-1 does not support persisted default store limits") + re.ErrorContains( + checkDefaultStoreLimitPersistenceFeature("Scheduling Service", "scheduling-1", "old-build", "new-build"), + "Scheduling Service member scheduling-1 does not support persisted default store limits") +} diff --git a/server/util.go b/server/util.go index a1fce6a868..5937f4908c 100644 --- a/server/util.go +++ b/server/util.go @@ -29,6 +29,8 @@ import ( "github.com/pingcap/log" "github.com/tikv/pd/pkg/errs" + "github.com/tikv/pd/pkg/mcs/discovery" + mcsconstant "github.com/tikv/pd/pkg/mcs/utils/constant" "github.com/tikv/pd/pkg/utils/apiutil" "github.com/tikv/pd/pkg/utils/keypath" "github.com/tikv/pd/pkg/versioninfo" @@ -58,6 +60,63 @@ func CheckPDVersionWithClusterVersion(opt *config.PersistOptions) { } } +func checkDefaultStoreLimitPersistenceFeature(component, name, gitHash, featureGitHash string) error { + if featureGitHash == "" || featureGitHash != gitHash { + return errors.Errorf( + "cannot update default store limit while %s member %s does not support persisted default store limits", + component, name) + } + return nil +} + +// checkDefaultStoreLimitPersistenceSupport prevents a new persisted schedule +// field from being activated while a currently registered old PD or Scheduling +// Service member can still become leader/primary and drop the unknown field on +// its next persist. It is a rolling-upgrade gate, not a downgrade guard: a +// pre-feature binary does not understand this check or the persisted field. +func (s *Server) checkDefaultStoreLimitPersistenceSupport() error { + members, err := s.ReloadMembers() + if err != nil { + return errors.Annotate(err, "failed to load PD members before updating default store limit") + } + for _, member := range members { + gitHash, err := s.GetMember().GetMemberGitHash(member.GetMemberId()) + if err != nil { + return errors.Annotatef(err, "failed to load git hash for PD member %s", member.GetName()) + } + featureGitHash, err := s.GetMember().GetMemberFeature( + member.GetMemberId(), versioninfo.DefaultStoreLimitPersistence) + if err != nil { + return errors.Errorf( + "cannot update default store limit while PD member %s does not support persisted default store limits", + member.GetName()) + } + if err := checkDefaultStoreLimitPersistenceFeature( + "PD", member.GetName(), gitHash, featureGitHash); err != nil { + return err + } + } + + if !s.IsServiceIndependent(mcsconstant.SchedulingServiceName) { + return nil + } + schedulingMembers, err := discovery.GetMSMembers(mcsconstant.SchedulingServiceName, s.GetClient()) + if err != nil { + return errors.Annotate(err, "failed to load Scheduling Service members before updating default store limit") + } + if len(schedulingMembers) == 0 { + return errors.New("cannot update default store limit without a registered Scheduling Service member") + } + for _, member := range schedulingMembers { + if err := checkDefaultStoreLimitPersistenceFeature( + "Scheduling Service", member.Name, member.GitHash, + member.Features[versioninfo.DefaultStoreLimitPersistence]); err != nil { + return err + } + } + return nil +} + func checkBootstrapRequest(req *pdpb.BootstrapRequest) error { clusterID := keypath.ClusterID() // TODO: do more check for request fields validation. diff --git a/tests/server/api/store_test.go b/tests/server/api/store_test.go index 9538368f8e..3aae983ade 100644 --- a/tests/server/api/store_test.go +++ b/tests/server/api/store_test.go @@ -15,6 +15,7 @@ package api import ( + "context" "encoding/json" "fmt" "io" @@ -31,7 +32,10 @@ import ( "github.com/pingcap/kvproto/pkg/pdpb" "github.com/tikv/pd/pkg/core" + "github.com/tikv/pd/pkg/mcs/discovery" + "github.com/tikv/pd/pkg/mcs/utils/constant" "github.com/tikv/pd/pkg/response" + "github.com/tikv/pd/pkg/utils/keypath" "github.com/tikv/pd/pkg/utils/testutil" "github.com/tikv/pd/pkg/utils/typeutil" "github.com/tikv/pd/pkg/versioninfo" @@ -47,6 +51,39 @@ func TestStoreTestSuite(t *testing.T) { suite.Run(t, new(storeTestSuite)) } +func TestDefaultStoreLimitRequiresAllPDsToSupportPersistence(t *testing.T) { + re := require.New(t) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + cluster, err := tests.NewTestCluster(ctx, 3) + re.NoError(err) + defer cluster.Destroy() + re.NoError(cluster.RunInitialServers()) + re.NotEmpty(cluster.WaitLeader()) + re.NoError(cluster.GetLeaderServer().BootstrapCluster()) + + leader := cluster.GetLeaderServer() + var follower *tests.TestServer + for _, server := range cluster.GetServers() { + if server.GetServerID() != leader.GetServerID() { + follower = server + break + } + } + re.NotNil(follower) + _, err = cluster.GetEtcdClient().Delete(context.Background(), keypath.MemberFeaturePath( + follower.GetServerID(), versioninfo.DefaultStoreLimitPersistence)) + re.NoError(err) + + url := leader.GetAddr() + "/pd/api/v1/stores/limit" + body := []byte(`{"rate":60,"type":"add-peer"}`) + err = testutil.CheckPostJSON(tests.TestDialClient, url, body, + testutil.StatusNotOK(re), + testutil.StringContain(re, "does not support persisted default store limits")) + re.NoError(err) + re.Equal(float64(15), leader.GetPersistOptions().GetScheduleConfig().DefaultStoreLimit.AddPeer) +} + func (suite *storeTestSuite) SetupSuite() { suite.env = tests.NewSchedulingTestEnvironment(suite.T()) } @@ -151,6 +188,57 @@ func (suite *storeTestSuite) TestStores() { suite.env.RunTestInNonMicroserviceEnv(suite.checkStoreLabel) } +func (suite *storeTestSuite) TestStoreLimitPersistenceCompatibility() { + suite.env.RunTest(suite.checkStoreLimitPersistenceCompatibility) +} + +func (suite *storeTestSuite) checkStoreLimitPersistenceCompatibility(cluster *tests.TestCluster) { + re := suite.Require() + leader := cluster.GetLeaderServer() + url := leader.GetAddr() + "/pd/api/v1/stores/limit" + body := []byte(`{"rate":60,"type":"add-peer"}`) + + _, err := cluster.GetEtcdClient().Delete(context.Background(), keypath.MemberFeaturePath( + leader.GetServerID(), versioninfo.DefaultStoreLimitPersistence)) + re.NoError(err) + err = testutil.CheckPostJSON(tests.TestDialClient, url, body, + testutil.StatusNotOK(re), + testutil.StringContain(re, "does not support persisted default store limits")) + re.NoError(err) + + err = leader.GetServer().GetMember().SetMemberFeature( + leader.GetServerID(), versioninfo.DefaultStoreLimitPersistence, versioninfo.PDGitHash) + re.NoError(err) + if schedulingServer := cluster.GetSchedulingPrimaryServer(); schedulingServer != nil { + entry := &discovery.ServiceRegistryEntry{ + Name: schedulingServer.Name(), + ServiceAddr: schedulingServer.GetAdvertiseListenAddr(), + Version: versioninfo.PDReleaseVersion, + GitHash: versioninfo.PDGitHash, + } + serializedEntry, err := entry.Serialize() + re.NoError(err) + registryPath := keypath.RegistryPath(constant.SchedulingServiceName, entry.ServiceAddr) + _, err = cluster.GetEtcdClient().Put(context.Background(), registryPath, serializedEntry) + re.NoError(err) + err = testutil.CheckPostJSON(tests.TestDialClient, url, body, + testutil.StatusNotOK(re), + testutil.StringContain(re, "Scheduling Service member")) + re.NoError(err) + + entry.Features = map[string]string{ + versioninfo.DefaultStoreLimitPersistence: versioninfo.PDGitHash, + } + serializedEntry, err = entry.Serialize() + re.NoError(err) + _, err = cluster.GetEtcdClient().Put(context.Background(), registryPath, serializedEntry) + re.NoError(err) + } + err = testutil.CheckPostJSON(tests.TestDialClient, url, body, testutil.StatusOK(re)) + re.NoError(err) + re.Equal(float64(60), leader.GetPersistOptions().GetScheduleConfig().DefaultStoreLimit.AddPeer) +} + func (suite *storeTestSuite) checkGetAllLimit(cluster *tests.TestCluster) { re := suite.Require() diff --git a/tests/server/config/config_test.go b/tests/server/config/config_test.go index 8beed55466..1155b5af55 100644 --- a/tests/server/config/config_test.go +++ b/tests/server/config/config_test.go @@ -305,6 +305,20 @@ func (suite *configTestSuite) checkConfigSchedule(cluster *tests.TestCluster) { re.NoError(testutil.ReadGetJSON(re, tests.TestDialClient, addr, scheduleConfig1)) return reflect.DeepEqual(*scheduleConfig1, *scheduleConfig) }) + + invalidDefaultStoreLimit := map[string]any{ + "default-store-limit": map[string]any{ + "add-peer": -1, + "remove-peer": scheduleConfig.DefaultStoreLimit.RemovePeer, + }, + } + postData, err = json.Marshal(invalidDefaultStoreLimit) + re.NoError(err) + err = testutil.CheckPostJSON(tests.TestDialClient, addr, postData, + testutil.StatusNotOK(re), + testutil.StringContain(re, "default-store-limit.add-peer should be finite and non-negative")) + re.NoError(err) + re.Equal(scheduleConfig.DefaultStoreLimit, leaderServer.GetPersistOptions().GetScheduleConfig().DefaultStoreLimit) } func (suite *configTestSuite) TestConfigReplication() {