From 9a4df95f572ab243ca3f7ed6cda2b31aeef1b34c Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Tue, 28 Jul 2026 10:13:20 +0800 Subject: [PATCH 01/10] schedule, server: fix replicated state semantics Treat the RuleChecker pending list as a hint and revalidate missing placement-rule peers against the current Store topology. Report PENDING only when no add or swap target exists; otherwise report INPROGRESS. Also avoid reporting uncovered Region ranges as replicated. Signed-off-by: Ryan Leung --- pkg/mcs/scheduling/server/apis/v1/api.go | 2 +- pkg/schedule/checker/rule_checker.go | 118 ++++++++++++++++++++++ pkg/schedule/checker/rule_checker_test.go | 57 +++++++++++ pkg/schedule/handler/handler.go | 57 +++++++++-- server/api/region.go | 2 +- tests/server/api/region_test.go | 49 +++++++++ 6 files changed, 273 insertions(+), 12 deletions(-) diff --git a/pkg/mcs/scheduling/server/apis/v1/api.go b/pkg/mcs/scheduling/server/apis/v1/api.go index 52bcc399509..9ed8a016f46 100644 --- a/pkg/mcs/scheduling/server/apis/v1/api.go +++ b/pkg/mcs/scheduling/server/apis/v1/api.go @@ -1466,7 +1466,7 @@ func splitRegions(c *gin.Context) { } // @Tags region -// @Summary Check if regions in the given key ranges are replicated. Returns 'REPLICATED', 'INPROGRESS', or 'PENDING'. 'PENDING' means that there is at least one region pending for scheduling. Similarly, 'INPROGRESS' means there is at least one region in scheduling. +// @Summary Check if regions in the given key ranges are replicated. Returns 'REPLICATED', 'INPROGRESS', or 'PENDING'. 'PENDING' means that the current Store topology cannot satisfy at least one required placement-rule peer. 'INPROGRESS' means that placement is incomplete but PD can continue or retry scheduling, including incomplete Region metadata in the requested range. // @Param startKey query string true "Regions start key, hex encoded" // @Param endKey query string true "Regions end key, hex encoded" // @Produce plain diff --git a/pkg/schedule/checker/rule_checker.go b/pkg/schedule/checker/rule_checker.go index 24a419a5af9..a48ab0facc1 100644 --- a/pkg/schedule/checker/rule_checker.go +++ b/pkg/schedule/checker/rule_checker.go @@ -34,11 +34,24 @@ import ( "github.com/tikv/pd/pkg/schedule/operator" "github.com/tikv/pd/pkg/schedule/placement" "github.com/tikv/pd/pkg/schedule/types" + "github.com/tikv/pd/pkg/statistics" "github.com/tikv/pd/pkg/utils/syncutil" ) const maxPendingListLen = 100000 +// RegionPlacementState describes the current placement-rule scheduling state. +type RegionPlacementState int + +const ( + // RegionPlacementStateReplicated means no placement work remains. + RegionPlacementStateReplicated RegionPlacementState = iota + // RegionPlacementStateInProgress means placement is incomplete but can progress or retry. + RegionPlacementStateInProgress + // RegionPlacementStatePending means the current topology cannot satisfy the placement rules. + RegionPlacementStatePending +) + // RuleChecker fix/improve region by placement rules. type RuleChecker struct { PauseController @@ -67,6 +80,111 @@ func (*RuleChecker) Name() string { return types.RuleChecker.String() } +// GetRegionPlacementState revalidates the placement state without creating an +// Operator. It is intended for Regions already recorded in RuleChecker's +// pending list. +func (c *RuleChecker) GetRegionPlacementState(region *core.RegionInfo) RegionPlacementState { + if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 { + return RegionPlacementStateInProgress + } + fit := c.ruleManager.FitRegionWithoutCache(c.cluster, region) + if len(fit.RuleFits) == 0 || len(fit.OrphanPeers) > 0 { + return RegionPlacementStateInProgress + } + + hasUnfixablePlacement := false + for _, rf := range fit.RuleFits { + if len(rf.Peers) < rf.Rule.Count { + isWitness := rf.Rule.IsWitness && isWitnessEnabled(c.cluster) + storeID, filteredByTemporaryState := c.strategy(region, rf.Rule, isWitness). + SelectStoreToAdd(getRuleFitStores(c.cluster, rf)) + if storeID != 0 || filteredByTemporaryState || c.hasStoreToSwapForMissingRule(region, fit, rf) { + return RegionPlacementStateInProgress + } + hasUnfixablePlacement = true + continue + } + + ruleUnfixable := false + for _, peer := range rf.Peers { + store := c.cluster.GetStore(peer.GetStoreId()) + if store == nil { + return RegionPlacementStateInProgress + } + + var fastFailover bool + switch { + case c.isDownPeer(region, peer): + if !c.isStoreDownTimeHitMaxDownTime(peer.GetStoreId()) { + return RegionPlacementStateInProgress + } + fastFailover = isWitnessEnabled(c.cluster) && store.IsTiKV() + case c.isOfflinePeer(peer): + fastFailover = isWitnessEnabled(c.cluster) && store.IsTiKV() && rf.Rule.IsWitness + default: + continue + } + + storeID, filteredByTemporaryState := c.strategy(region, rf.Rule, fastFailover). + SelectStoreToFix(getRuleFitStores(c.cluster, rf), peer.GetStoreId()) + if storeID != 0 || filteredByTemporaryState { + return RegionPlacementStateInProgress + } + hasUnfixablePlacement = true + ruleUnfixable = true + break + } + if ruleUnfixable { + continue + } + if len(rf.PeersWithDifferentRole) > 0 { + return RegionPlacementStateInProgress + } + if len(rf.Rule.LocationLabels) == 0 { + continue + } + + isWitness := rf.Rule.IsWitness && isWitnessEnabled(c.cluster) + _, newStoreID, filteredByTemporaryState := c.strategy(region, rf.Rule, isWitness). + getBetterLocation(c.cluster, region, fit, rf) + if newStoreID != 0 || filteredByTemporaryState { + return RegionPlacementStateInProgress + } + if !statistics.IsRegionLabelIsolationSatisfied( + getRuleFitStores(c.cluster, rf), + rf.Rule.LocationLabels, + rf.Rule.IsolationLevel, + ) { + hasUnfixablePlacement = true + } + } + if hasUnfixablePlacement { + return RegionPlacementStatePending + } + return RegionPlacementStateReplicated +} + +func (c *RuleChecker) hasStoreToSwapForMissingRule(region *core.RegionInfo, fit *placement.RegionFit, missingRuleFit *placement.RuleFit) bool { + for _, peer := range region.GetPeers() { + store := c.cluster.GetStore(peer.GetStoreId()) + if store == nil || !placement.MatchLabelConstraints(store, missingRuleFit.Rule.LabelConstraints) { + continue + } + oldRuleFit := fit.GetRuleFit(peer.GetId()) + if oldRuleFit == nil || oldRuleFit == missingRuleFit || !oldRuleFit.IsSatisfied() { + continue + } + + fastFailover := isWitnessEnabled(c.cluster) && store.IsTiKV() && oldRuleFit.Rule.IsWitness + storeID, filteredByTemporaryState := c.strategy(region, oldRuleFit.Rule, fastFailover). + SelectStoreToFix(getRuleFitStores(c.cluster, oldRuleFit), peer.GetStoreId()) + if storeID != 0 || filteredByTemporaryState { + return true + } + } + return false +} + // GetType returns RuleChecker's type. func (*RuleChecker) GetType() types.CheckerSchedulerType { return types.RuleChecker diff --git a/pkg/schedule/checker/rule_checker_test.go b/pkg/schedule/checker/rule_checker_test.go index 43d25b92cf7..67fbc6174e7 100644 --- a/pkg/schedule/checker/rule_checker_test.go +++ b/pkg/schedule/checker/rule_checker_test.go @@ -2217,9 +2217,11 @@ func (suite *ruleCheckerTestSuite) TestPendingList() { re.Nil(op) _, exist := suite.rc.pendingList.Get(1) re.True(exist) + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(suite.cluster.GetRegion(1))) // add more stores suite.cluster.AddLeaderStore(3, 1) + re.Equal(RegionPlacementStateInProgress, suite.rc.GetRegionPlacementState(suite.cluster.GetRegion(1))) op = suite.rc.Check(suite.cluster.GetRegion(1)) re.NotNil(op) re.Equal("add-rule-peer", op.Desc()) @@ -2229,6 +2231,61 @@ func (suite *ruleCheckerTestSuite) TestPendingList() { re.False(exist) } +func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingOfflinePeer() { + re := suite.Require() + suite.cluster.AddLeaderStore(1, 1) + suite.cluster.AddLeaderStore(2, 1) + suite.cluster.AddLeaderStore(3, 1) + suite.cluster.AddLeaderRegion(1, 1, 2, 3) + suite.cluster.SetStoreOffline(3) + + region := suite.cluster.GetRegion(1) + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(region)) + + suite.cluster.AddLeaderStore(4, 1) + re.Equal(RegionPlacementStateInProgress, suite.rc.GetRegionPlacementState(region)) +} + +func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingDownPeer() { + re := suite.Require() + suite.cluster.AddLeaderStore(1, 1) + suite.cluster.AddLeaderStore(2, 1) + suite.cluster.AddLeaderStore(3, 1) + suite.cluster.AddLeaderRegion(1, 1, 2, 3) + suite.cluster.SetStoreDown(3) + + region := suite.cluster.GetRegion(1) + region = region.Clone(core.WithDownPeers([]*pdpb.PeerStats{{Peer: region.GetStorePeer(3)}})) + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(region)) + + suite.cluster.AddLeaderStore(4, 1) + re.Equal(RegionPlacementStateInProgress, suite.rc.GetRegionPlacementState(region)) +} + +func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingUnmetIsolation() { + re := suite.Require() + suite.cluster.AddLabelsStore(1, 1, map[string]string{"zone": "z1"}) + suite.cluster.AddLabelsStore(2, 1, map[string]string{"zone": "z1"}) + suite.cluster.AddLabelsStore(3, 1, map[string]string{"zone": "z1"}) + suite.cluster.AddLeaderRegion(1, 1, 2, 3) + re.NoError(suite.ruleManager.SetRule(&placement.Rule{ + GroupID: placement.DefaultGroupID, + ID: "test", + Index: 100, + Override: true, + Role: placement.Voter, + Count: 3, + LocationLabels: []string{"zone"}, + IsolationLevel: "zone", + })) + + region := suite.cluster.GetRegion(1) + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(region)) + + suite.cluster.AddLabelsStore(4, 1, map[string]string{"zone": "z2"}) + re.Equal(RegionPlacementStateInProgress, suite.rc.GetRegionPlacementState(region)) +} + func (suite *ruleCheckerTestSuite) TestLocationLabels() { re := suite.Require() suite.cluster.AddLabelsStore(1, 1, map[string]string{"zone": "z1", "rack": "r1", "host": "h1"}) diff --git a/pkg/schedule/handler/handler.go b/pkg/schedule/handler/handler.go index bf193c030a7..7254bc126c1 100644 --- a/pkg/schedule/handler/handler.go +++ b/pkg/schedule/handler/handler.go @@ -34,6 +34,7 @@ import ( "github.com/tikv/pd/pkg/core" "github.com/tikv/pd/pkg/errs" "github.com/tikv/pd/pkg/schedule" + "github.com/tikv/pd/pkg/schedule/checker" sche "github.com/tikv/pd/pkg/schedule/core" "github.com/tikv/pd/pkg/schedule/filter" "github.com/tikv/pd/pkg/schedule/labeler" @@ -1290,25 +1291,61 @@ func (h *Handler) CheckRegionsReplicated(startKeyHex, endKeyHex string) (string, return "", errs.ErrNotBootstrapped.GenWithStackByArgs() } regions := c.ScanRegions(startKey, endKey, -1) + if !regionsCoverRange(regions, startKey, endKey) { + return "INPROGRESS", nil + } state := "REPLICATED" + placementRulesEnabled := c.GetSharedConfig().IsPlacementRulesEnabled() for _, region := range regions { - if !filter.IsRegionReplicated(c, region) { + if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 { state = "INPROGRESS" - if co.IsPendingRegion(region.GetID()) { - state = "PENDING" - break + continue + } + if placementRulesEnabled { + pending := co.IsPendingRegion(region.GetID()) + failpoint.Inject("mockPending", func(val failpoint.Value) { + if mockPending, ok := val.(bool); ok { + pending = mockPending + } + }) + if pending { + switch co.GetRuleChecker().GetRegionPlacementState(region) { + case checker.RegionPlacementStatePending: + return "PENDING", nil + case checker.RegionPlacementStateInProgress: + state = "INPROGRESS" + } + continue } } - } - failpoint.Inject("mockPending", func(val failpoint.Value) { - aok, ok := val.(bool) - if ok && aok { - state = "PENDING" + if !filter.IsRegionReplicated(c, region) { + state = "INPROGRESS" } - }) + } return state, nil } +func regionsCoverRange(regions []*core.RegionInfo, startKey, endKey []byte) bool { + if len(regions) == 0 { + return false + } + first := regions[0] + if bytes.Compare(first.GetStartKey(), startKey) > 0 || + (len(first.GetEndKey()) > 0 && bytes.Compare(startKey, first.GetEndKey()) >= 0) { + return false + } + for i := 1; i < len(regions); i++ { + if !bytes.Equal(regions[i-1].GetEndKey(), regions[i].GetStartKey()) { + return false + } + } + lastEndKey := regions[len(regions)-1].GetEndKey() + if len(endKey) == 0 { + return len(lastEndKey) == 0 + } + return len(lastEndKey) == 0 || bytes.Compare(lastEndKey, endKey) >= 0 +} + // GetRuleManager returns the rule manager. func (h *Handler) GetRuleManager() (*placement.RuleManager, error) { c := h.GetCluster() diff --git a/server/api/region.go b/server/api/region.go index 179c4f09c00..04a461ac1ed 100644 --- a/server/api/region.go +++ b/server/api/region.go @@ -136,7 +136,7 @@ func (h *regionHandler) GetRegion(w http.ResponseWriter, r *http.Request) { // CheckRegionsReplicated checks if regions in the given key ranges are replicated. // // @Tags region -// @Summary Check if regions in the given key ranges are replicated. Returns 'REPLICATED', 'INPROGRESS', or 'PENDING'. 'PENDING' means that there is at least one region pending for scheduling. Similarly, 'INPROGRESS' means there is at least one region in scheduling. +// @Summary Check if regions in the given key ranges are replicated. Returns 'REPLICATED', 'INPROGRESS', or 'PENDING'. 'PENDING' means that the current Store topology cannot satisfy at least one required placement-rule peer. 'INPROGRESS' means that placement is incomplete but PD can continue or retry scheduling, including incomplete Region metadata in the requested range. // @Param startKey query string true "Regions start key, hex encoded" // @Param endKey query string true "Regions end key, hex encoded" // @Produce plain diff --git a/tests/server/api/region_test.go b/tests/server/api/region_test.go index 550437ca2a5..6dcbeb6ff44 100644 --- a/tests/server/api/region_test.go +++ b/tests/server/api/region_test.go @@ -202,11 +202,60 @@ func (suite *regionTestSuite) checkRegionsReplicated(cluster *tests.TestCluster) return status == "REPLICATED" }) + re.NoError(failpoint.Enable("github.com/tikv/pd/pkg/schedule/handler/mockPending", "return(true)")) + err = testutil.ReadGetJSON(re, tests.TestDialClient, url, &status) + re.NoError(err) + re.Equal("REPLICATED", status) + re.NoError(failpoint.Disable("github.com/tikv/pd/pkg/schedule/handler/mockPending")) + + // A rule can have enough peers but still be pending because an offline + // peer has no replacement target. + offlineStore := &metapb.Store{ + Id: 1, + State: metapb.StoreState_Offline, + NodeState: metapb.NodeState_Removing, + } + tests.MustPutStore(re, cluster, offlineStore) re.NoError(failpoint.Enable("github.com/tikv/pd/pkg/schedule/handler/mockPending", "return(true)")) err = testutil.ReadGetJSON(re, tests.TestDialClient, url, &status) re.NoError(err) re.Equal("PENDING", status) + targetStore := &metapb.Store{ + Id: 2, + State: metapb.StoreState_Up, + NodeState: metapb.NodeState_Serving, + } + tests.MustPutStore(re, cluster, targetStore) + availableTarget := leader.GetRaftCluster().GetStore(2).Clone(core.SetStoreStats(&pdpb.StoreStats{ + Capacity: uint64(10 * units.GiB), + Available: uint64(10 * units.GiB), + })) + leader.GetRaftCluster().GetBasicCluster().PutStore(availableTarget) + if schedulingServer := cluster.GetSchedulingPrimaryServer(); schedulingServer != nil { + schedulingServer.GetCluster().PutStore(availableTarget) + } + err = testutil.ReadGetJSON(re, tests.TestDialClient, url, &status) + re.NoError(err) + re.Equal("INPROGRESS", status) re.NoError(failpoint.Disable("github.com/tikv/pd/pkg/schedule/handler/mockPending")) + targetStore.State = metapb.StoreState_Offline + targetStore.NodeState = metapb.NodeState_Removing + tests.MustPutStore(re, cluster, targetStore) + tests.MustPutStore(re, cluster, s1) + + // Missing Region metadata must never be reported as complete. + noRegionURL := fmt.Sprintf(`%s/regions/replicated?startKey=%s&endKey=%s`, urlPrefix, + hex.EncodeToString(r1.GetEndKey()), hex.EncodeToString([]byte("c"))) + err = testutil.ReadGetJSON(re, tests.TestDialClient, noRegionURL, &status) + re.NoError(err) + re.Equal("INPROGRESS", status) + + gapURL := fmt.Sprintf(`%s/regions/replicated?startKey=%s&endKey=%s`, urlPrefix, + hex.EncodeToString([]byte("0")), hex.EncodeToString(r1.GetEndKey())) + err = testutil.ReadGetJSON(re, tests.TestDialClient, gapURL, &status) + re.NoError(err) + re.Equal("INPROGRESS", status) + // test multiple rules r1 = core.NewTestRegionInfo(2, 1, []byte("a"), []byte("b")) r1.GetMeta().Peers = append(r1.GetMeta().Peers, &metapb.Peer{Id: 5, StoreId: 1}) From 703ccc15776952fb0e91eee3b3c864714de0a8d7 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Tue, 28 Jul 2026 16:03:22 +0800 Subject: [PATCH 02/10] schedule: reduce replicated-state lookup overhead Signed-off-by: Ryan Leung --- pkg/schedule/checker/checker_controller.go | 2 +- pkg/schedule/checker/replica_strategy.go | 52 +++++++++++----- pkg/schedule/checker/rule_checker.go | 72 ++++++++++++++++++---- pkg/schedule/checker/rule_checker_test.go | 49 +++++++++++++++ pkg/schedule/handler/handler.go | 16 +++-- 5 files changed, 159 insertions(+), 32 deletions(-) diff --git a/pkg/schedule/checker/checker_controller.go b/pkg/schedule/checker/checker_controller.go index 7da5ae7bb49..275612be178 100644 --- a/pkg/schedule/checker/checker_controller.go +++ b/pkg/schedule/checker/checker_controller.go @@ -535,7 +535,7 @@ func (c *Controller) ClearSuspectKeyRanges() { // IsPendingRegion returns true if the given region is in the pending list. func (c *Controller) IsPendingRegion(regionID uint64) bool { - _, exist := c.ruleChecker.pendingList.Get(regionID) + _, exist := c.ruleChecker.pendingList.Peek(regionID) return exist } diff --git a/pkg/schedule/checker/replica_strategy.go b/pkg/schedule/checker/replica_strategy.go index f374e723e92..0d1b7929d21 100644 --- a/pkg/schedule/checker/replica_strategy.go +++ b/pkg/schedule/checker/replica_strategy.go @@ -30,13 +30,15 @@ import ( // ReplicaStrategy collects some utilities to manipulate region peers. It // exists to allow replica_checker and rule_checker to reuse common logics. type ReplicaStrategy struct { - checkerName string // replica-checker / rule-checker - cluster sche.CheckerCluster - locationLabels []string - isolationLevel string - region *core.RegionInfo - extraFilters []filter.Filter - fastFailover bool + checkerName string // replica-checker / rule-checker + cluster sche.CheckerCluster + locationLabels []string + isolationLevel string + region *core.RegionInfo + extraFilters []filter.Filter + fastFailover bool + storeCandidates []*core.StoreInfo + storeCandidatesPrepared bool } // SelectStoreToAdd returns the store to add a replica to a region. @@ -62,11 +64,15 @@ func (s *ReplicaStrategy) SelectStoreToAdd(coLocationStores []*core.StoreInfo, e if s.fastFailover { level = constant.Urgent } - filters := []filter.Filter{ - filter.NewExcludedFilter(s.checkerName, nil, s.region.GetStoreIDs()), - filter.NewStorageThresholdFilter(s.checkerName), - filter.NewSpecialUseFilter(s.checkerName), - &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowTemporaryStates: true, OperatorLevel: level}, + stores := s.storeCandidates + filters := []filter.Filter{filter.NewExcludedFilter(s.checkerName, nil, s.region.GetStoreIDs())} + if !s.storeCandidatesPrepared { + stores = s.cluster.GetStores() + filters = append(filters, + filter.NewStorageThresholdFilter(s.checkerName), + filter.NewSpecialUseFilter(s.checkerName), + &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowTemporaryStates: true, OperatorLevel: level}, + ) } if len(s.locationLabels) > 0 && s.isolationLevel != "" { filters = append(filters, filter.NewIsolationFilter(s.checkerName, s.isolationLevel, s.locationLabels, coLocationStores)) @@ -74,13 +80,13 @@ func (s *ReplicaStrategy) SelectStoreToAdd(coLocationStores []*core.StoreInfo, e if len(extraFilters) > 0 { filters = append(filters, extraFilters...) } - if len(s.extraFilters) > 0 { + if !s.storeCandidatesPrepared && len(s.extraFilters) > 0 { filters = append(filters, s.extraFilters...) } isolationComparer := filter.IsolationComparer(s.locationLabels, coLocationStores) strictStateFilter := &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowFastFailover: s.fastFailover, OperatorLevel: level} - targetCandidate := filter.NewCandidates(s.cluster.GetStores()). + targetCandidate := filter.NewCandidates(stores). FilterTarget(s.cluster.GetCheckerConfig(), nil, nil, filters...). KeepTheTopStores(isolationComparer, false) // greater isolation score is better if targetCandidate.Len() == 0 { @@ -94,6 +100,24 @@ func (s *ReplicaStrategy) SelectStoreToAdd(coLocationStores []*core.StoreInfo, e return target.GetID(), false } +// prepareStoreCandidates applies the Region-independent target filters so the +// result can be reused by placement-state checks using the same Rule. +func (s *ReplicaStrategy) prepareStoreCandidates(stores []*core.StoreInfo) { + level := constant.High + if s.fastFailover { + level = constant.Urgent + } + filters := []filter.Filter{ + filter.NewStorageThresholdFilter(s.checkerName), + filter.NewSpecialUseFilter(s.checkerName), + &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowTemporaryStates: true, OperatorLevel: level}, + } + filters = append(filters, s.extraFilters...) + s.storeCandidates = filter.NewCandidates(stores). + FilterTarget(s.cluster.GetCheckerConfig(), nil, nil, filters...).Stores + s.storeCandidatesPrepared = true +} + // SelectStoreToFix returns a store to replace down/offline old peer. The location // placement after scheduling is allowed to be worse than original. func (s *ReplicaStrategy) SelectStoreToFix(coLocationStores []*core.StoreInfo, old uint64) (uint64, bool) { diff --git a/pkg/schedule/checker/rule_checker.go b/pkg/schedule/checker/rule_checker.go index a48ab0facc1..2899513b77b 100644 --- a/pkg/schedule/checker/rule_checker.go +++ b/pkg/schedule/checker/rule_checker.go @@ -80,10 +80,45 @@ func (*RuleChecker) Name() string { return types.RuleChecker.String() } +type placementStateCandidateKey struct { + rule *placement.Rule + fastFailover bool +} + +type placementStateContext struct { + storesLoaded bool + stores []*core.StoreInfo + candidates map[placementStateCandidateKey][]*core.StoreInfo +} + +func (ctx *placementStateContext) getStrategy(c *RuleChecker, region *core.RegionInfo, rule *placement.Rule, fastFailover bool) *ReplicaStrategy { + strategy := c.strategy(region, rule, fastFailover) + key := placementStateCandidateKey{rule: rule, fastFailover: fastFailover} + if candidates, ok := ctx.candidates[key]; ok { + strategy.storeCandidates = candidates + strategy.storeCandidatesPrepared = true + return strategy + } + if !ctx.storesLoaded { + ctx.stores = c.cluster.GetStores() + ctx.storesLoaded = true + } + if ctx.candidates == nil { + ctx.candidates = make(map[placementStateCandidateKey][]*core.StoreInfo) + } + strategy.prepareStoreCandidates(ctx.stores) + ctx.candidates[key] = strategy.storeCandidates + return strategy +} + // GetRegionPlacementState revalidates the placement state without creating an // Operator. It is intended for Regions already recorded in RuleChecker's // pending list. func (c *RuleChecker) GetRegionPlacementState(region *core.RegionInfo) RegionPlacementState { + return c.getRegionPlacementState(region, &placementStateContext{}) +} + +func (c *RuleChecker) getRegionPlacementState(region *core.RegionInfo, context *placementStateContext) RegionPlacementState { if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 { return RegionPlacementStateInProgress } @@ -96,9 +131,10 @@ func (c *RuleChecker) GetRegionPlacementState(region *core.RegionInfo) RegionPla for _, rf := range fit.RuleFits { if len(rf.Peers) < rf.Rule.Count { isWitness := rf.Rule.IsWitness && isWitnessEnabled(c.cluster) - storeID, filteredByTemporaryState := c.strategy(region, rf.Rule, isWitness). - SelectStoreToAdd(getRuleFitStores(c.cluster, rf)) - if storeID != 0 || filteredByTemporaryState || c.hasStoreToSwapForMissingRule(region, fit, rf) { + strategy := context.getStrategy(c, region, rf.Rule, isWitness) + storeID, filteredByTemporaryState := strategy.SelectStoreToAdd(getRuleFitStores(c.cluster, rf)) + if storeID != 0 || filteredByTemporaryState || + c.hasStoreToSwapForMissingRule(region, fit, rf, context) { return RegionPlacementStateInProgress } hasUnfixablePlacement = true @@ -125,8 +161,8 @@ func (c *RuleChecker) GetRegionPlacementState(region *core.RegionInfo) RegionPla continue } - storeID, filteredByTemporaryState := c.strategy(region, rf.Rule, fastFailover). - SelectStoreToFix(getRuleFitStores(c.cluster, rf), peer.GetStoreId()) + strategy := context.getStrategy(c, region, rf.Rule, fastFailover) + storeID, filteredByTemporaryState := strategy.SelectStoreToFix(getRuleFitStores(c.cluster, rf), peer.GetStoreId()) if storeID != 0 || filteredByTemporaryState { return RegionPlacementStateInProgress } @@ -143,10 +179,9 @@ func (c *RuleChecker) GetRegionPlacementState(region *core.RegionInfo) RegionPla if len(rf.Rule.LocationLabels) == 0 { continue } - isWitness := rf.Rule.IsWitness && isWitnessEnabled(c.cluster) - _, newStoreID, filteredByTemporaryState := c.strategy(region, rf.Rule, isWitness). - getBetterLocation(c.cluster, region, fit, rf) + strategy := context.getStrategy(c, region, rf.Rule, isWitness) + _, newStoreID, filteredByTemporaryState := strategy.getBetterLocation(c.cluster, region, fit, rf) if newStoreID != 0 || filteredByTemporaryState { return RegionPlacementStateInProgress } @@ -164,7 +199,22 @@ func (c *RuleChecker) GetRegionPlacementState(region *core.RegionInfo) RegionPla return RegionPlacementStateReplicated } -func (c *RuleChecker) hasStoreToSwapForMissingRule(region *core.RegionInfo, fit *placement.RegionFit, missingRuleFit *placement.RuleFit) bool { +// GetRegionsPlacementState returns the aggregate placement state of Regions. +func (c *RuleChecker) GetRegionsPlacementState(regions []*core.RegionInfo) RegionPlacementState { + state := RegionPlacementStateReplicated + context := &placementStateContext{} + for _, region := range regions { + switch c.getRegionPlacementState(region, context) { + case RegionPlacementStatePending: + return RegionPlacementStatePending + case RegionPlacementStateInProgress: + state = RegionPlacementStateInProgress + } + } + return state +} + +func (c *RuleChecker) hasStoreToSwapForMissingRule(region *core.RegionInfo, fit *placement.RegionFit, missingRuleFit *placement.RuleFit, context *placementStateContext) bool { for _, peer := range region.GetPeers() { store := c.cluster.GetStore(peer.GetStoreId()) if store == nil || !placement.MatchLabelConstraints(store, missingRuleFit.Rule.LabelConstraints) { @@ -176,8 +226,8 @@ func (c *RuleChecker) hasStoreToSwapForMissingRule(region *core.RegionInfo, fit } fastFailover := isWitnessEnabled(c.cluster) && store.IsTiKV() && oldRuleFit.Rule.IsWitness - storeID, filteredByTemporaryState := c.strategy(region, oldRuleFit.Rule, fastFailover). - SelectStoreToFix(getRuleFitStores(c.cluster, oldRuleFit), peer.GetStoreId()) + strategy := context.getStrategy(c, region, oldRuleFit.Rule, fastFailover) + storeID, filteredByTemporaryState := strategy.SelectStoreToFix(getRuleFitStores(c.cluster, oldRuleFit), peer.GetStoreId()) if storeID != 0 || filteredByTemporaryState { return true } diff --git a/pkg/schedule/checker/rule_checker_test.go b/pkg/schedule/checker/rule_checker_test.go index 67fbc6174e7..dc507279ad9 100644 --- a/pkg/schedule/checker/rule_checker_test.go +++ b/pkg/schedule/checker/rule_checker_test.go @@ -54,6 +54,16 @@ type ruleCheckerTestSuite struct { cancel context.CancelFunc } +type getStoresCountingCluster struct { + *mockcluster.Cluster + getStoresCount int +} + +func (c *getStoresCountingCluster) GetStores() []*core.StoreInfo { + c.getStoresCount++ + return c.Cluster.GetStores() +} + func (suite *ruleCheckerTestSuite) SetupTest() { cfg := mockconfig.NewTestOptions() suite.ctx, suite.cancel = context.WithCancel(context.Background()) @@ -2231,6 +2241,45 @@ func (suite *ruleCheckerTestSuite) TestPendingList() { re.False(exist) } +func (suite *ruleCheckerTestSuite) TestPendingListLookupDoesNotRefreshLRU() { + re := suite.Require() + suite.rc.pendingList = cache.NewDefaultCache(2) + suite.rc.pendingList.Put(1, nil) + suite.rc.pendingList.Put(2, nil) + controller := &Controller{ruleChecker: suite.rc} + + re.True(controller.IsPendingRegion(1)) + suite.rc.pendingList.Put(3, nil) + + _, exists := suite.rc.pendingList.Peek(1) + re.False(exists) + _, exists = suite.rc.pendingList.Peek(2) + re.True(exists) +} + +func (suite *ruleCheckerTestSuite) TestRegionsPlacementStateLoadsStoresOnce() { + re := suite.Require() + suite.cluster.AddLeaderStore(1, 1) + suite.cluster.AddLeaderStore(2, 1) + suite.cluster.AddLeaderStore(3, 1) + regions := make([]*core.RegionInfo, 0, 10) + for i := uint64(1); i <= 10; i++ { + suite.cluster.AddLeaderRegion(i, 1, 2) + regions = append(regions, suite.cluster.GetRegion(i)) + } + + cluster := &getStoresCountingCluster{Cluster: suite.cluster} + checker := NewRuleChecker( + suite.ctx, + cluster, + suite.ruleManager, + cache.NewIDTTL(suite.ctx, time.Minute, 3*time.Minute), + ) + + re.Equal(RegionPlacementStateInProgress, checker.GetRegionsPlacementState(regions)) + re.Equal(1, cluster.getStoresCount) +} + func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingOfflinePeer() { re := suite.Require() suite.cluster.AddLeaderStore(1, 1) diff --git a/pkg/schedule/handler/handler.go b/pkg/schedule/handler/handler.go index 7254bc126c1..840619f7fcd 100644 --- a/pkg/schedule/handler/handler.go +++ b/pkg/schedule/handler/handler.go @@ -1296,6 +1296,7 @@ func (h *Handler) CheckRegionsReplicated(startKeyHex, endKeyHex string) (string, } state := "REPLICATED" placementRulesEnabled := c.GetSharedConfig().IsPlacementRulesEnabled() + var pendingRegions []*core.RegionInfo for _, region := range regions { if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 { state = "INPROGRESS" @@ -1309,12 +1310,7 @@ func (h *Handler) CheckRegionsReplicated(startKeyHex, endKeyHex string) (string, } }) if pending { - switch co.GetRuleChecker().GetRegionPlacementState(region) { - case checker.RegionPlacementStatePending: - return "PENDING", nil - case checker.RegionPlacementStateInProgress: - state = "INPROGRESS" - } + pendingRegions = append(pendingRegions, region) continue } } @@ -1322,6 +1318,14 @@ func (h *Handler) CheckRegionsReplicated(startKeyHex, endKeyHex string) (string, state = "INPROGRESS" } } + if len(pendingRegions) > 0 { + switch co.GetRuleChecker().GetRegionsPlacementState(pendingRegions) { + case checker.RegionPlacementStatePending: + return "PENDING", nil + case checker.RegionPlacementStateInProgress: + state = "INPROGRESS" + } + } return state, nil } From c419ffbe3f0a4724ef9d105c19679ec1e207dd02 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Tue, 28 Jul 2026 17:46:42 +0800 Subject: [PATCH 03/10] schedule: harden replicated-state checks Evaluate placement states without relying on the RuleChecker pending LRU, distinguish retryable work from unschedulable topology, and keep completed-range checks on the fast path. Reuse Store candidates across incomplete Regions and cover active replica operators and missing Region metadata. Signed-off-by: Ryan Leung --- pkg/mcs/scheduling/server/apis/v1/api.go | 2 +- pkg/schedule/checker/checker_controller.go | 2 +- pkg/schedule/checker/rule_checker.go | 44 +++++++++++++--- pkg/schedule/checker/rule_checker_test.go | 60 ++++++++++++++++------ pkg/schedule/handler/handler.go | 55 ++++++++++++-------- server/api/region.go | 2 +- tests/server/api/region_test.go | 11 +--- 7 files changed, 120 insertions(+), 56 deletions(-) diff --git a/pkg/mcs/scheduling/server/apis/v1/api.go b/pkg/mcs/scheduling/server/apis/v1/api.go index 9ed8a016f46..e126c3983d7 100644 --- a/pkg/mcs/scheduling/server/apis/v1/api.go +++ b/pkg/mcs/scheduling/server/apis/v1/api.go @@ -1466,7 +1466,7 @@ func splitRegions(c *gin.Context) { } // @Tags region -// @Summary Check if regions in the given key ranges are replicated. Returns 'REPLICATED', 'INPROGRESS', or 'PENDING'. 'PENDING' means that the current Store topology cannot satisfy at least one required placement-rule peer. 'INPROGRESS' means that placement is incomplete but PD can continue or retry scheduling, including incomplete Region metadata in the requested range. +// @Summary Check if regions in the given key ranges are replicated. Returns 'REPLICATED', 'INPROGRESS', or 'PENDING'. 'PENDING' means that PD cannot currently make scheduling progress because the Store topology has no eligible placement. 'INPROGRESS' means that placement is incomplete but PD can continue or retry scheduling, including incomplete Region metadata in the requested range. // @Param startKey query string true "Regions start key, hex encoded" // @Param endKey query string true "Regions end key, hex encoded" // @Produce plain diff --git a/pkg/schedule/checker/checker_controller.go b/pkg/schedule/checker/checker_controller.go index 275612be178..7da5ae7bb49 100644 --- a/pkg/schedule/checker/checker_controller.go +++ b/pkg/schedule/checker/checker_controller.go @@ -535,7 +535,7 @@ func (c *Controller) ClearSuspectKeyRanges() { // IsPendingRegion returns true if the given region is in the pending list. func (c *Controller) IsPendingRegion(regionID uint64) bool { - _, exist := c.ruleChecker.pendingList.Peek(regionID) + _, exist := c.ruleChecker.pendingList.Get(regionID) return exist } diff --git a/pkg/schedule/checker/rule_checker.go b/pkg/schedule/checker/rule_checker.go index 2899513b77b..97dd4e70862 100644 --- a/pkg/schedule/checker/rule_checker.go +++ b/pkg/schedule/checker/rule_checker.go @@ -111,18 +111,20 @@ func (ctx *placementStateContext) getStrategy(c *RuleChecker, region *core.Regio return strategy } -// GetRegionPlacementState revalidates the placement state without creating an -// Operator. It is intended for Regions already recorded in RuleChecker's -// pending list. +// GetRegionPlacementState evaluates the placement state without creating an +// Operator or updating RuleChecker's caches and metrics. func (c *RuleChecker) GetRegionPlacementState(region *core.RegionInfo) RegionPlacementState { - return c.getRegionPlacementState(region, &placementStateContext{}) + return c.evaluateRegionPlacementState(region, &placementStateContext{}) } -func (c *RuleChecker) getRegionPlacementState(region *core.RegionInfo, context *placementStateContext) RegionPlacementState { - if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 { +func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, context *placementStateContext) RegionPlacementState { + if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 || c.pendingProcessedRegions.Exists(region.GetID()) { return RegionPlacementStateInProgress } fit := c.ruleManager.FitRegionWithoutCache(c.cluster, region) + if isRegionPlacementSatisfied(region, fit) { + return RegionPlacementStateReplicated + } if len(fit.RuleFits) == 0 || len(fit.OrphanPeers) > 0 { return RegionPlacementStateInProgress } @@ -199,12 +201,40 @@ func (c *RuleChecker) getRegionPlacementState(region *core.RegionInfo, context * return RegionPlacementStateReplicated } +func isRegionPlacementSatisfied(region *core.RegionInfo, fit *placement.RegionFit) bool { + if len(fit.RuleFits) == 0 || len(fit.OrphanPeers) > 0 { + return false + } + for _, rf := range fit.RuleFits { + if !rf.IsSatisfied() { + return false + } + if len(rf.Stores) != len(rf.Peers) { + return false + } + for _, peer := range rf.Peers { + if region.GetDownPeer(peer.GetId()) != nil { + return false + } + } + for _, store := range rf.Stores { + if !store.IsPreparing() && !store.IsServing() { + return false + } + } + if !statistics.IsRegionLabelIsolationSatisfied(rf.Stores, rf.Rule.LocationLabels, rf.Rule.IsolationLevel) { + return false + } + } + return true +} + // GetRegionsPlacementState returns the aggregate placement state of Regions. func (c *RuleChecker) GetRegionsPlacementState(regions []*core.RegionInfo) RegionPlacementState { state := RegionPlacementStateReplicated context := &placementStateContext{} for _, region := range regions { - switch c.getRegionPlacementState(region, context) { + switch c.evaluateRegionPlacementState(region, context) { case RegionPlacementStatePending: return RegionPlacementStatePending case RegionPlacementStateInProgress: diff --git a/pkg/schedule/checker/rule_checker_test.go b/pkg/schedule/checker/rule_checker_test.go index dc507279ad9..96e1f358937 100644 --- a/pkg/schedule/checker/rule_checker_test.go +++ b/pkg/schedule/checker/rule_checker_test.go @@ -2241,22 +2241,6 @@ func (suite *ruleCheckerTestSuite) TestPendingList() { re.False(exist) } -func (suite *ruleCheckerTestSuite) TestPendingListLookupDoesNotRefreshLRU() { - re := suite.Require() - suite.rc.pendingList = cache.NewDefaultCache(2) - suite.rc.pendingList.Put(1, nil) - suite.rc.pendingList.Put(2, nil) - controller := &Controller{ruleChecker: suite.rc} - - re.True(controller.IsPendingRegion(1)) - suite.rc.pendingList.Put(3, nil) - - _, exists := suite.rc.pendingList.Peek(1) - re.False(exists) - _, exists = suite.rc.pendingList.Peek(2) - re.True(exists) -} - func (suite *ruleCheckerTestSuite) TestRegionsPlacementStateLoadsStoresOnce() { re := suite.Require() suite.cluster.AddLeaderStore(1, 1) @@ -2280,6 +2264,50 @@ func (suite *ruleCheckerTestSuite) TestRegionsPlacementStateLoadsStoresOnce() { re.Equal(1, cluster.getStoresCount) } +func (suite *ruleCheckerTestSuite) TestRegionsPlacementStateSkipsStoresForCompletedRegions() { + re := suite.Require() + suite.cluster.AddLabelsStore(1, 1, map[string]string{"zone": "z1"}) + suite.cluster.AddLabelsStore(2, 1, map[string]string{"zone": "z2"}) + suite.cluster.AddLabelsStore(3, 1, map[string]string{"zone": "z3"}) + re.NoError(suite.ruleManager.SetRule(&placement.Rule{ + GroupID: placement.DefaultGroupID, + ID: placement.DefaultRuleID, + Role: placement.Voter, + Count: 3, + LocationLabels: []string{"zone"}, + IsolationLevel: "zone", + })) + regions := make([]*core.RegionInfo, 0, 10) + for i := uint64(1); i <= 10; i++ { + suite.cluster.AddLeaderRegion(i, 1, 2, 3) + regions = append(regions, suite.cluster.GetRegion(i)) + } + + cluster := &getStoresCountingCluster{Cluster: suite.cluster} + checker := NewRuleChecker( + suite.ctx, + cluster, + suite.ruleManager, + cache.NewIDTTL(suite.ctx, time.Minute, 3*time.Minute), + ) + + re.Equal(RegionPlacementStateReplicated, checker.GetRegionsPlacementState(regions)) + re.Zero(cluster.getStoresCount) +} + +func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingProcessedIsInProgress() { + re := suite.Require() + suite.cluster.AddLeaderStore(1, 1) + suite.cluster.AddLeaderStore(2, 1) + suite.cluster.AddLeaderStore(3, 1) + suite.cluster.AddLeaderRegion(1, 1, 2, 3) + pendingProcessed := cache.NewIDTTL(suite.ctx, time.Minute, 3*time.Minute) + pendingProcessed.Put(1, nil) + checker := NewRuleChecker(suite.ctx, suite.cluster, suite.ruleManager, pendingProcessed) + + re.Equal(RegionPlacementStateInProgress, checker.GetRegionPlacementState(suite.cluster.GetRegion(1))) +} + func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingOfflinePeer() { re := suite.Require() suite.cluster.AddLeaderStore(1, 1) diff --git a/pkg/schedule/handler/handler.go b/pkg/schedule/handler/handler.go index 840619f7fcd..4efed166681 100644 --- a/pkg/schedule/handler/handler.go +++ b/pkg/schedule/handler/handler.go @@ -1294,39 +1294,52 @@ func (h *Handler) CheckRegionsReplicated(startKeyHex, endKeyHex string) (string, if !regionsCoverRange(regions, startKey, endKey) { return "INPROGRESS", nil } + if c.GetSharedConfig().IsPlacementRulesEnabled() { + switch co.GetRuleChecker().GetRegionsPlacementState(regions) { + case checker.RegionPlacementStatePending: + return "PENDING", nil + case checker.RegionPlacementStateInProgress: + return "INPROGRESS", nil + } + if hasReplicaOperator(co.GetOperatorController(), regions) { + return "INPROGRESS", nil + } + return "REPLICATED", nil + } + state := "REPLICATED" - placementRulesEnabled := c.GetSharedConfig().IsPlacementRulesEnabled() - var pendingRegions []*core.RegionInfo for _, region := range regions { if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 { state = "INPROGRESS" continue } - if placementRulesEnabled { - pending := co.IsPendingRegion(region.GetID()) - failpoint.Inject("mockPending", func(val failpoint.Value) { - if mockPending, ok := val.(bool); ok { - pending = mockPending - } - }) - if pending { - pendingRegions = append(pendingRegions, region) - continue - } - } if !filter.IsRegionReplicated(c, region) { state = "INPROGRESS" } } - if len(pendingRegions) > 0 { - switch co.GetRuleChecker().GetRegionsPlacementState(pendingRegions) { - case checker.RegionPlacementStatePending: - return "PENDING", nil - case checker.RegionPlacementStateInProgress: - state = "INPROGRESS" + return state, nil +} + +func hasReplicaOperator(controller *operator.Controller, regions []*core.RegionInfo) bool { + operators := controller.GetOperatorsOfKind(operator.OpReplica) + for _, op := range controller.GetWaitingOperators() { + if op.Kind()&operator.OpReplica != 0 { + operators = append(operators, op) } } - return state, nil + if len(operators) == 0 { + return false + } + operatorRegions := make(map[uint64]struct{}, len(operators)) + for _, op := range operators { + operatorRegions[op.RegionID()] = struct{}{} + } + for _, region := range regions { + if _, ok := operatorRegions[region.GetID()]; ok { + return true + } + } + return false } func regionsCoverRange(regions []*core.RegionInfo, startKey, endKey []byte) bool { diff --git a/server/api/region.go b/server/api/region.go index 04a461ac1ed..389bb4ba765 100644 --- a/server/api/region.go +++ b/server/api/region.go @@ -136,7 +136,7 @@ func (h *regionHandler) GetRegion(w http.ResponseWriter, r *http.Request) { // CheckRegionsReplicated checks if regions in the given key ranges are replicated. // // @Tags region -// @Summary Check if regions in the given key ranges are replicated. Returns 'REPLICATED', 'INPROGRESS', or 'PENDING'. 'PENDING' means that the current Store topology cannot satisfy at least one required placement-rule peer. 'INPROGRESS' means that placement is incomplete but PD can continue or retry scheduling, including incomplete Region metadata in the requested range. +// @Summary Check if regions in the given key ranges are replicated. Returns 'REPLICATED', 'INPROGRESS', or 'PENDING'. 'PENDING' means that PD cannot currently make scheduling progress because the Store topology has no eligible placement. 'INPROGRESS' means that placement is incomplete but PD can continue or retry scheduling, including incomplete Region metadata in the requested range. // @Param startKey query string true "Regions start key, hex encoded" // @Param endKey query string true "Regions end key, hex encoded" // @Produce plain diff --git a/tests/server/api/region_test.go b/tests/server/api/region_test.go index 6dcbeb6ff44..969bb21f6e1 100644 --- a/tests/server/api/region_test.go +++ b/tests/server/api/region_test.go @@ -202,12 +202,6 @@ func (suite *regionTestSuite) checkRegionsReplicated(cluster *tests.TestCluster) return status == "REPLICATED" }) - re.NoError(failpoint.Enable("github.com/tikv/pd/pkg/schedule/handler/mockPending", "return(true)")) - err = testutil.ReadGetJSON(re, tests.TestDialClient, url, &status) - re.NoError(err) - re.Equal("REPLICATED", status) - re.NoError(failpoint.Disable("github.com/tikv/pd/pkg/schedule/handler/mockPending")) - // A rule can have enough peers but still be pending because an offline // peer has no replacement target. offlineStore := &metapb.Store{ @@ -216,7 +210,6 @@ func (suite *regionTestSuite) checkRegionsReplicated(cluster *tests.TestCluster) NodeState: metapb.NodeState_Removing, } tests.MustPutStore(re, cluster, offlineStore) - re.NoError(failpoint.Enable("github.com/tikv/pd/pkg/schedule/handler/mockPending", "return(true)")) err = testutil.ReadGetJSON(re, tests.TestDialClient, url, &status) re.NoError(err) re.Equal("PENDING", status) @@ -237,7 +230,6 @@ func (suite *regionTestSuite) checkRegionsReplicated(cluster *tests.TestCluster) err = testutil.ReadGetJSON(re, tests.TestDialClient, url, &status) re.NoError(err) re.Equal("INPROGRESS", status) - re.NoError(failpoint.Disable("github.com/tikv/pd/pkg/schedule/handler/mockPending")) targetStore.State = metapb.StoreState_Offline targetStore.NodeState = metapb.NodeState_Removing tests.MustPutStore(re, cluster, targetStore) @@ -311,10 +303,11 @@ func (suite *regionTestSuite) checkRegionsReplicated(cluster *tests.TestCluster) return s1 || s2 }) + // The new rule requires more peers, but there is no available target store. testutil.Eventually(re, func() bool { err = testutil.ReadGetJSON(re, tests.TestDialClient, url, &status) re.NoError(err) - return status == "INPROGRESS" + return status == "PENDING" }) r1 = core.NewTestRegionInfo(2, 1, []byte("a"), []byte("b")) From 5913b77ebeab8117db637868767616a5ac336492 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Wed, 29 Jul 2026 17:05:06 +0800 Subject: [PATCH 04/10] schedule: fix replicated state convergence Signed-off-by: Ryan Leung --- pkg/schedule/checker/checker_controller.go | 6 +++ pkg/schedule/checker/rule_checker.go | 59 ++++++++++++++++++++-- pkg/schedule/checker/rule_checker_test.go | 12 +++++ tests/server/api/region_test.go | 30 ++++++----- 4 files changed, 91 insertions(+), 16 deletions(-) diff --git a/pkg/schedule/checker/checker_controller.go b/pkg/schedule/checker/checker_controller.go index 7da5ae7bb49..e2149846b11 100644 --- a/pkg/schedule/checker/checker_controller.go +++ b/pkg/schedule/checker/checker_controller.go @@ -266,6 +266,12 @@ func (c *Controller) checkPendingProcessedRegions() { pendingProcessedRegionsGauge.Set(float64(len(ids))) for _, id := range ids { region := c.cluster.GetRegion(id) + if region == nil { + continue + } + // Remove the old entry before checking. A checker will add it back if + // temporary conditions still prevent scheduling. + c.RemovePendingProcessedRegion(id) c.tryAddOperators(region) } } diff --git a/pkg/schedule/checker/rule_checker.go b/pkg/schedule/checker/rule_checker.go index 97dd4e70862..0cc5fa40e15 100644 --- a/pkg/schedule/checker/rule_checker.go +++ b/pkg/schedule/checker/rule_checker.go @@ -175,8 +175,16 @@ func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, cont if ruleUnfixable { continue } - if len(rf.PeersWithDifferentRole) > 0 { - return RegionPlacementStateInProgress + for _, peer := range rf.PeersWithDifferentRole { + if c.getPeerRolePlacementState(region, fit, rf, peer) == RegionPlacementStateInProgress { + return RegionPlacementStateInProgress + } + hasUnfixablePlacement = true + ruleUnfixable = true + break + } + if ruleUnfixable { + continue } if len(rf.Rule.LocationLabels) == 0 { continue @@ -244,6 +252,43 @@ func (c *RuleChecker) GetRegionsPlacementState(regions []*core.RegionInfo) Regio return state } +func (c *RuleChecker) getPeerRolePlacementState( + region *core.RegionInfo, + fit *placement.RegionFit, + rf *placement.RuleFit, + peer *metapb.Peer, +) RegionPlacementState { + if core.IsLearner(peer) && rf.Rule.Role != placement.Learner { + return RegionPlacementStateInProgress + } + if region.GetLeader().GetId() != peer.GetId() && rf.Rule.Role == placement.Leader { + if c.allowLeaderWithTemporaryStates(fit, peer, true) { + return RegionPlacementStateInProgress + } + return RegionPlacementStatePending + } + if region.GetLeader().GetId() == peer.GetId() && rf.Rule.Role == placement.Follower { + for _, candidate := range region.GetPeers() { + if c.cluster.GetStore(candidate.GetStoreId()) == nil || + c.allowLeaderWithTemporaryStates(fit, candidate, true) { + return RegionPlacementStateInProgress + } + } + return RegionPlacementStatePending + } + if core.IsVoter(peer) && rf.Rule.Role == placement.Learner { + return RegionPlacementStateInProgress + } + if region.GetLeader().GetId() == peer.GetId() && rf.Rule.IsWitness { + return RegionPlacementStatePending + } + if (!core.IsWitness(peer) && rf.Rule.IsWitness && isWitnessEnabled(c.cluster)) || + (core.IsWitness(peer) && (!rf.Rule.IsWitness || !isWitnessEnabled(c.cluster))) { + return RegionPlacementStateInProgress + } + return RegionPlacementStatePending +} + func (c *RuleChecker) hasStoreToSwapForMissingRule(region *core.RegionInfo, fit *placement.RegionFit, missingRuleFit *placement.RuleFit, context *placementStateContext) bool { for _, peer := range region.GetPeers() { store := c.cluster.GetStore(peer.GetStoreId()) @@ -598,6 +643,10 @@ func (c *RuleChecker) fixLooseMatchPeer(region *core.RegionInfo, fit *placement. } func (c *RuleChecker) allowLeader(fit *placement.RegionFit, peer *metapb.Peer) bool { + return c.allowLeaderWithTemporaryStates(fit, peer, false) +} + +func (c *RuleChecker) allowLeaderWithTemporaryStates(fit *placement.RegionFit, peer *metapb.Peer, allowTemporaryStates bool) bool { if core.IsLearner(peer) || core.IsWitness(peer) { return false } @@ -605,7 +654,11 @@ func (c *RuleChecker) allowLeader(fit *placement.RegionFit, peer *metapb.Peer) b if s == nil { return false } - stateFilter := &filter.StoreStateFilter{ActionScope: "rule-checker", TransferLeader: true} + stateFilter := &filter.StoreStateFilter{ + ActionScope: "rule-checker", + TransferLeader: true, + AllowTemporaryStates: allowTemporaryStates, + } if !stateFilter.Target(c.cluster.GetCheckerConfig(), s).IsOK() { return false } diff --git a/pkg/schedule/checker/rule_checker_test.go b/pkg/schedule/checker/rule_checker_test.go index 96e1f358937..bd9e08191b6 100644 --- a/pkg/schedule/checker/rule_checker_test.go +++ b/pkg/schedule/checker/rule_checker_test.go @@ -34,6 +34,7 @@ import ( "github.com/tikv/pd/pkg/errs" "github.com/tikv/pd/pkg/mock/mockcluster" "github.com/tikv/pd/pkg/mock/mockconfig" + "github.com/tikv/pd/pkg/schedule/config" "github.com/tikv/pd/pkg/schedule/operator" "github.com/tikv/pd/pkg/schedule/placement" "github.com/tikv/pd/pkg/utils/operatorutil" @@ -419,6 +420,17 @@ func (suite *ruleCheckerTestSuite) TestFixRoleLeader() { re.NotNil(op) re.Equal("fix-follower-role", op.Desc()) re.Equal(uint64(3), op.Step(0).(operator.TransferLeader).ToStore) + re.Equal(RegionPlacementStateInProgress, suite.rc.GetRegionPlacementState(suite.cluster.GetRegion(1))) + + // A temporary Store state can recover, so PD can retry the leader transfer. + suite.cluster.SetStoreBusy(3, true) + re.Equal(RegionPlacementStateInProgress, suite.rc.GetRegionPlacementState(suite.cluster.GetRegion(1))) + suite.cluster.SetStoreBusy(3, false) + + // A reject-leader target cannot satisfy the role with the current topology. + suite.cluster.SetLabelProperty(config.RejectLeader, "role", "voter") + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(suite.cluster.GetRegion(1))) + re.Nil(suite.rc.Check(suite.cluster.GetRegion(1))) } func (suite *ruleCheckerTestSuite) TestFixRoleLeaderIssue3130() { diff --git a/tests/server/api/region_test.go b/tests/server/api/region_test.go index 969bb21f6e1..9fcf3589780 100644 --- a/tests/server/api/region_test.go +++ b/tests/server/api/region_test.go @@ -204,12 +204,18 @@ func (suite *regionTestSuite) checkRegionsReplicated(cluster *tests.TestCluster) // A rule can have enough peers but still be pending because an offline // peer has no replacement target. - offlineStore := &metapb.Store{ - Id: 1, - State: metapb.StoreState_Offline, - NodeState: metapb.NodeState_Removing, + putStore := func(store *core.StoreInfo) { + leader.GetRaftCluster().GetBasicCluster().PutStore(store) + if schedulingServer := cluster.GetSchedulingPrimaryServer(); schedulingServer != nil { + schedulingServer.GetCluster().PutStore(store) + } } - tests.MustPutStore(re, cluster, offlineStore) + servingStore := leader.GetRaftCluster().GetStore(1) + offlineStore := servingStore.Clone( + core.SetStoreState(metapb.StoreState_Offline, false), + core.SetNodeState(metapb.NodeState_Removing), + ) + putStore(offlineStore) err = testutil.ReadGetJSON(re, tests.TestDialClient, url, &status) re.NoError(err) re.Equal("PENDING", status) @@ -223,17 +229,15 @@ func (suite *regionTestSuite) checkRegionsReplicated(cluster *tests.TestCluster) Capacity: uint64(10 * units.GiB), Available: uint64(10 * units.GiB), })) - leader.GetRaftCluster().GetBasicCluster().PutStore(availableTarget) - if schedulingServer := cluster.GetSchedulingPrimaryServer(); schedulingServer != nil { - schedulingServer.GetCluster().PutStore(availableTarget) - } + putStore(availableTarget) err = testutil.ReadGetJSON(re, tests.TestDialClient, url, &status) re.NoError(err) re.Equal("INPROGRESS", status) - targetStore.State = metapb.StoreState_Offline - targetStore.NodeState = metapb.NodeState_Removing - tests.MustPutStore(re, cluster, targetStore) - tests.MustPutStore(re, cluster, s1) + putStore(availableTarget.Clone( + core.SetStoreState(metapb.StoreState_Offline, false), + core.SetNodeState(metapb.NodeState_Removing), + )) + putStore(servingStore) // Missing Region metadata must never be reported as complete. noRegionURL := fmt.Sprintf(`%s/regions/replicated?startKey=%s&endKey=%s`, urlPrefix, From d4cd707fe1af2c85fe4ade814a1dd7d6ee04eaa7 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 30 Jul 2026 10:57:01 +0800 Subject: [PATCH 05/10] schedule: align replicated state with schedulable work Signed-off-by: Ryan Leung --- pkg/schedule/checker/rule_checker.go | 151 +++++++++++++++++- pkg/schedule/checker/rule_checker_test.go | 125 +++++++++++++++ pkg/schedule/handler/handler.go | 20 +-- pkg/schedule/operator/operator_controller.go | 11 ++ pkg/schedule/operator/waiting_operator.go | 57 ++++++- .../operator/waiting_operator_test.go | 18 +++ 6 files changed, 352 insertions(+), 30 deletions(-) diff --git a/pkg/schedule/checker/rule_checker.go b/pkg/schedule/checker/rule_checker.go index 0cc5fa40e15..192ec12de75 100644 --- a/pkg/schedule/checker/rule_checker.go +++ b/pkg/schedule/checker/rule_checker.go @@ -125,11 +125,14 @@ func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, cont if isRegionPlacementSatisfied(region, fit) { return RegionPlacementStateReplicated } - if len(fit.RuleFits) == 0 || len(fit.OrphanPeers) > 0 { + if len(fit.RuleFits) == 0 { + return RegionPlacementStateInProgress + } + if c.hasOrphanPeerAction(region, fit) { return RegionPlacementStateInProgress } - hasUnfixablePlacement := false + hasUnfixablePlacement := len(fit.OrphanPeers) > 0 for _, rf := range fit.RuleFits { if len(rf.Peers) < rf.Rule.Count { isWitness := rf.Rule.IsWitness && isWitnessEnabled(c.cluster) @@ -268,15 +271,15 @@ func (c *RuleChecker) getPeerRolePlacementState( return RegionPlacementStatePending } if region.GetLeader().GetId() == peer.GetId() && rf.Rule.Role == placement.Follower { - for _, candidate := range region.GetPeers() { - if c.cluster.GetStore(candidate.GetStoreId()) == nil || - c.allowLeaderWithTemporaryStates(fit, candidate, true) { - return RegionPlacementStateInProgress - } + if c.hasAlternativeLeader(region, fit, peer.GetId()) { + return RegionPlacementStateInProgress } return RegionPlacementStatePending } if core.IsVoter(peer) && rf.Rule.Role == placement.Learner { + if region.GetLeader().GetId() == peer.GetId() && !c.hasAlternativeLeader(region, fit, peer.GetId()) { + return RegionPlacementStatePending + } return RegionPlacementStateInProgress } if region.GetLeader().GetId() == peer.GetId() && rf.Rule.IsWitness { @@ -289,6 +292,134 @@ func (c *RuleChecker) getPeerRolePlacementState( return RegionPlacementStatePending } +func (c *RuleChecker) hasAlternativeLeader(region *core.RegionInfo, fit *placement.RegionFit, excludedPeerID uint64) bool { + for _, candidate := range region.GetPeers() { + if candidate.GetId() == excludedPeerID || core.IsLearner(candidate) || core.IsWitness(candidate) { + continue + } + if region.GetPendingPeer(candidate.GetId()) != nil { + return true + } + store := c.cluster.GetStore(candidate.GetStoreId()) + if region.GetDownPeer(candidate.GetId()) != nil { + if store == nil || store.DownTime() < c.cluster.GetCheckerConfig().GetMaxStoreDownTime() { + return true + } + continue + } + if store == nil || c.allowLeaderWithTemporaryStates(fit, candidate, true) { + return true + } + } + return false +} + +func (c *RuleChecker) hasOrphanPeerAction(region *core.RegionInfo, fit *placement.RegionFit) bool { + if len(fit.OrphanPeers) == 0 { + return false + } + + isPendingPeer := func(peer *metapb.Peer) bool { + return region.GetPendingPeer(peer.GetId()) != nil + } + isDownPeer := func(peer *metapb.Peer) bool { + return region.GetDownPeer(peer.GetId()) != nil + } + isDisconnectedPeer := func(peer *metapb.Peer) bool { + store := c.cluster.GetStore(peer.GetStoreId()) + return store == nil || store.IsDisconnected() + } + canRemovePeer := func(peer *metapb.Peer) bool { + return region.GetLeader().GetId() != peer.GetId() || + c.hasAlternativeLeader(region, fit, peer.GetId()) + } + + allRulesSatisfiedAndHealthy := true + var pinDownPeer *metapb.Peer + for _, rf := range fit.RuleFits { + if !rf.IsSatisfied() { + allRulesSatisfiedAndHealthy = false + break + } + for _, peer := range rf.Peers { + if isPendingPeer(peer) { + return true + } + if isDownPeer(peer) { + if !c.isStoreDownTimeHitMaxDownTime(peer.GetStoreId()) { + return true + } + pinDownPeer = peer + allRulesSatisfiedAndHealthy = false + break + } + if isDisconnectedPeer(peer) { + return true + } + } + if !allRulesSatisfiedAndHealthy { + break + } + } + if allRulesSatisfiedAndHealthy { + // fixOrphanPeers always tries the first orphan peer in this case. + return canRemovePeer(fit.OrphanPeers[0]) + } + + if pinDownPeer != nil { + for _, orphanPeer := range fit.OrphanPeers { + if pinDownPeer.GetId() == orphanPeer.GetId() || isPendingPeer(orphanPeer) || + isDownPeer(orphanPeer) || isDisconnectedPeer(orphanPeer) || + pinDownPeer.GetIsWitness() || orphanPeer.GetIsWitness() { + continue + } + if !isDisconnectedPeer(pinDownPeer) { + continue + } + dstStore := c.cluster.GetStore(orphanPeer.GetStoreId()) + if !fit.Replace(pinDownPeer.GetStoreId(), dstStore) { + continue + } + destRole := pinDownPeer.GetRole() + orphanRole := orphanPeer.GetRole() + switch { + case orphanRole == metapb.PeerRole_Learner && destRole == metapb.PeerRole_Voter: + return true + case orphanRole == metapb.PeerRole_Voter && destRole == metapb.PeerRole_Learner: + return c.cluster.GetSharedConfig().IsUseJointConsensus() + case orphanRole == destRole: + return canRemovePeer(pinDownPeer) + } + } + } + + if len(fit.OrphanPeers) < 2 { + return false + } + var disconnectedPeer *metapb.Peer + for _, orphanPeer := range fit.OrphanPeers { + if isDisconnectedPeer(orphanPeer) { + disconnectedPeer = orphanPeer + break + } + } + hasHealthyPeer := false + for _, orphanPeer := range fit.OrphanPeers { + if isPendingPeer(orphanPeer) || isDownPeer(orphanPeer) { + // fixOrphanPeers returns after trying the first unhealthy orphan. + return canRemovePeer(orphanPeer) + } + if hasHealthyPeer && fit.ExtraCount() > 0 { + if disconnectedPeer != nil { + return canRemovePeer(disconnectedPeer) + } + return canRemovePeer(orphanPeer) + } + hasHealthyPeer = true + } + return false +} + func (c *RuleChecker) hasStoreToSwapForMissingRule(region *core.RegionInfo, fit *placement.RegionFit, missingRuleFit *placement.RuleFit, context *placementStateContext) bool { for _, peer := range region.GetPeers() { store := c.cluster.GetStore(peer.GetStoreId()) @@ -601,6 +732,12 @@ func (c *RuleChecker) fixLooseMatchPeer(region *core.RegionInfo, fit *placement. if region.GetLeader().GetId() == peer.GetId() && rf.Rule.Role == placement.Follower { ruleCheckerFixFollowerRoleCounter.Inc() for _, p := range region.GetPeers() { + if p.GetId() == peer.GetId() { + continue + } + if region.GetPendingPeer(p.GetId()) != nil || region.GetDownPeer(p.GetId()) != nil { + continue + } if c.allowLeader(fit, p) { return operator.CreateTransferLeaderOperator("fix-follower-role", c.cluster, region, p.GetStoreId(), []uint64{}, 0) } diff --git a/pkg/schedule/checker/rule_checker_test.go b/pkg/schedule/checker/rule_checker_test.go index bd9e08191b6..28f726b17e7 100644 --- a/pkg/schedule/checker/rule_checker_test.go +++ b/pkg/schedule/checker/rule_checker_test.go @@ -433,6 +433,44 @@ func (suite *ruleCheckerTestSuite) TestFixRoleLeader() { re.Nil(suite.rc.Check(suite.cluster.GetRegion(1))) } +func (suite *ruleCheckerTestSuite) TestFixFollowerRoleWithOverlappingRules() { + re := suite.Require() + suite.cluster.AddLabelsStore(1, 1, map[string]string{"role": "follower"}) + suite.cluster.AddLabelsStore(2, 1, map[string]string{"role": "voter"}) + suite.cluster.AddLeaderRegion(1, 1, 2) + re.NoError(suite.ruleManager.SetRule(&placement.Rule{ + GroupID: placement.DefaultGroupID, + ID: "follower", + Index: 100, + Role: placement.Follower, + Count: 1, + LabelConstraints: []placement.LabelConstraint{ + {Key: "role", Op: placement.In, Values: []string{"follower"}}, + }, + })) + re.NoError(suite.ruleManager.SetRule(&placement.Rule{ + GroupID: placement.DefaultGroupID, + ID: "voter", + Index: 101, + Role: placement.Voter, + Count: 1, + })) + re.NoError(suite.ruleManager.DeleteRule(placement.DefaultGroupID, placement.DefaultRuleID)) + + region := suite.cluster.GetRegion(1) + re.Equal(RegionPlacementStateInProgress, suite.rc.GetRegionPlacementState(region)) + op := suite.rc.Check(region) + re.NotNil(op) + re.Equal("fix-follower-role", op.Desc()) + re.Equal(uint64(2), op.Step(0).(operator.TransferLeader).ToStore) + + // The current leader also matches the broad voter rule, but it must not be + // selected as its own transfer target when the other peer rejects leaders. + suite.cluster.SetLabelProperty(config.RejectLeader, "role", "voter") + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(region)) + re.Nil(suite.rc.Check(region)) +} + func (suite *ruleCheckerTestSuite) TestFixRoleLeaderIssue3130() { re := suite.Require() suite.cluster.AddLabelsStore(1, 1, map[string]string{"role": "follower"}) @@ -2202,6 +2240,41 @@ func (suite *ruleCheckerTestSuite) TestDemoteVoter() { re.Equal("fix-demote-voter", op.Desc()) } +func (suite *ruleCheckerTestSuite) TestDemoteLeaderRequiresAlternativeLeader() { + re := suite.Require() + suite.cluster.AddLabelsStore(1, 1, map[string]string{"role": "learner"}) + suite.cluster.AddLabelsStore(2, 1, map[string]string{"role": "voter"}) + suite.cluster.AddLeaderRegion(1, 1, 2) + re.NoError(suite.ruleManager.SetRule(&placement.Rule{ + GroupID: placement.DefaultGroupID, + ID: "learner", + Index: 100, + Role: placement.Learner, + Count: 1, + LabelConstraints: []placement.LabelConstraint{ + {Key: "role", Op: placement.In, Values: []string{"learner"}}, + }, + })) + re.NoError(suite.ruleManager.SetRule(&placement.Rule{ + GroupID: placement.DefaultGroupID, + ID: "voter", + Index: 101, + Role: placement.Voter, + Count: 1, + })) + re.NoError(suite.ruleManager.DeleteRule(placement.DefaultGroupID, placement.DefaultRuleID)) + + region := suite.cluster.GetRegion(1) + re.Equal(RegionPlacementStateInProgress, suite.rc.GetRegionPlacementState(region)) + op := suite.rc.Check(region) + re.NotNil(op) + re.Equal("fix-demote-voter", op.Desc()) + + suite.cluster.SetLabelProperty(config.RejectLeader, "role", "voter") + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(region)) + re.Nil(suite.rc.Check(region)) +} + func (suite *ruleCheckerTestSuite) TestOfflineAndDownStore() { re := suite.Require() suite.cluster.AddLabelsStore(1, 1, map[string]string{"zone": "z1"}) @@ -2253,6 +2326,58 @@ func (suite *ruleCheckerTestSuite) TestPendingList() { re.False(exist) } +func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingWithUnfixableOrphan() { + re := suite.Require() + suite.cluster.AddLabelsStore(1, 1, map[string]string{"zone": "z1"}) + suite.cluster.AddLabelsStore(2, 1, map[string]string{"zone": "z2"}) + suite.cluster.AddLeaderRegion(1, 1, 2) + re.NoError(suite.ruleManager.SetRule(&placement.Rule{ + GroupID: placement.DefaultGroupID, + ID: placement.DefaultRuleID, + Role: placement.Voter, + Count: 2, + LabelConstraints: []placement.LabelConstraint{ + {Key: "zone", Op: placement.In, Values: []string{"z1"}}, + }, + })) + + region := suite.cluster.GetRegion(1) + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(region)) + re.Nil(suite.rc.Check(region)) + + suite.cluster.AddLabelsStore(3, 1, map[string]string{"zone": "z1"}) + re.Equal(RegionPlacementStateInProgress, suite.rc.GetRegionPlacementState(region)) + re.NotNil(suite.rc.Check(region)) +} + +func (suite *ruleCheckerTestSuite) TestRegionPlacementFollowsOrphanRemovalOrder() { + re := suite.Require() + suite.cluster.AddLabelsStore(1, 1, map[string]string{"role": "blocked", "zone": "z1"}) + suite.cluster.AddLabelsStore(2, 1, map[string]string{"role": "blocked", "zone": "z2"}) + suite.cluster.AddLabelsStore(3, 1, map[string]string{"role": "blocked", "zone": "z3"}) + suite.cluster.AddLeaderRegion(1, 1, 2, 3) + suite.cluster.SetLabelProperty(config.RejectLeader, "role", "blocked") + re.NoError(suite.ruleManager.SetRule(&placement.Rule{ + GroupID: placement.DefaultGroupID, + ID: placement.DefaultRuleID, + Role: placement.Voter, + Count: 1, + LabelConstraints: []placement.LabelConstraint{ + {Key: "zone", Op: placement.In, Values: []string{"z2"}}, + }, + })) + + region := suite.cluster.GetRegion(1) + fit := suite.ruleManager.FitRegionWithoutCache(suite.cluster, region) + re.Len(fit.OrphanPeers, 2) + re.Equal(uint64(1), fit.OrphanPeers[0].GetStoreId()) + // The checker always tries the first orphan. It cannot remove the leader + // because no other peer is allowed to lead, even though the second orphan + // itself would be removable. + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(region)) + re.Nil(suite.rc.Check(region)) +} + func (suite *ruleCheckerTestSuite) TestRegionsPlacementStateLoadsStoresOnce() { re := suite.Require() suite.cluster.AddLeaderStore(1, 1) diff --git a/pkg/schedule/handler/handler.go b/pkg/schedule/handler/handler.go index 4efed166681..fd4544117e2 100644 --- a/pkg/schedule/handler/handler.go +++ b/pkg/schedule/handler/handler.go @@ -1321,25 +1321,11 @@ func (h *Handler) CheckRegionsReplicated(startKeyHex, endKeyHex string) (string, } func hasReplicaOperator(controller *operator.Controller, regions []*core.RegionInfo) bool { - operators := controller.GetOperatorsOfKind(operator.OpReplica) - for _, op := range controller.GetWaitingOperators() { - if op.Kind()&operator.OpReplica != 0 { - operators = append(operators, op) - } - } - if len(operators) == 0 { - return false - } - operatorRegions := make(map[uint64]struct{}, len(operators)) - for _, op := range operators { - operatorRegions[op.RegionID()] = struct{}{} - } + regionIDs := make([]uint64, 0, len(regions)) for _, region := range regions { - if _, ok := operatorRegions[region.GetID()]; ok { - return true - } + regionIDs = append(regionIDs, region.GetID()) } - return false + return controller.HasOperatorOfKind(regionIDs, operator.OpReplica) } func regionsCoverRange(regions []*core.RegionInfo, startKey, endKey []byte) bool { diff --git a/pkg/schedule/operator/operator_controller.go b/pkg/schedule/operator/operator_controller.go index 80cffcc399b..b940b009125 100644 --- a/pkg/schedule/operator/operator_controller.go +++ b/pkg/schedule/operator/operator_controller.go @@ -808,6 +808,17 @@ func (oc *Controller) GetWaitingOperators() []*Operator { return oc.wop.ListOperator() } +// HasOperatorOfKind checks whether any of the Regions has a running or waiting +// operator of the specified kind. +func (oc *Controller) HasOperatorOfKind(regionIDs []uint64, mask OpKind) bool { + for _, regionID := range regionIDs { + if op := oc.GetOperator(regionID); op != nil && op.Kind()&mask != 0 { + return true + } + } + return oc.wop.HasOperator(regionIDs, mask) +} + // GetOperatorsOfKind returns the running operators of the kind. func (oc *Controller) GetOperatorsOfKind(mask OpKind) []*Operator { operators := make([]*Operator, 0, oc.opNotifierQueue.len()) diff --git a/pkg/schedule/operator/waiting_operator.go b/pkg/schedule/operator/waiting_operator.go index 70013e143c7..93076fb8da4 100644 --- a/pkg/schedule/operator/waiting_operator.go +++ b/pkg/schedule/operator/waiting_operator.go @@ -29,6 +29,7 @@ type WaitingOperator interface { PutMergeOperators(op []*Operator) GetOperator() []*Operator ListOperator() []*Operator + HasOperator(regionIDs []uint64, mask OpKind) bool } // bucket is used to maintain the operators created by a specific scheduler. @@ -39,9 +40,10 @@ type bucket struct { // randBuckets is an implementation of waiting operators type randBuckets struct { - mu syncutil.Mutex - totalWeight float64 - buckets []*bucket + mu syncutil.RWMutex + totalWeight float64 + buckets []*bucket + operatorsByRegion map[uint64][]*Operator } // newRandBuckets creates a random buckets. @@ -52,7 +54,10 @@ func newRandBuckets() *randBuckets { weight: priorityWeight[i], }) } - return &randBuckets{buckets: buckets} + return &randBuckets{ + buckets: buckets, + operatorsByRegion: make(map[uint64][]*Operator), + } } // PutOperator puts an operator into the random buckets. @@ -65,6 +70,7 @@ func (b *randBuckets) PutOperator(op *Operator) { b.totalWeight += bucket.weight } bucket.ops = append(bucket.ops, op) + b.operatorsByRegion[op.RegionID()] = append(b.operatorsByRegion[op.RegionID()], op) } // PutMergeOperators puts two operators into the random buckets. @@ -80,12 +86,15 @@ func (b *randBuckets) PutMergeOperators(ops []*Operator) { b.totalWeight += bucket.weight } bucket.ops = append(bucket.ops, ops...) + for _, op := range ops { + b.operatorsByRegion[op.RegionID()] = append(b.operatorsByRegion[op.RegionID()], op) + } } // ListOperator lists all operator in the random buckets. func (b *randBuckets) ListOperator() []*Operator { - b.mu.Lock() - defer b.mu.Unlock() + b.mu.RLock() + defer b.mu.RUnlock() var ops []*Operator for i := range b.buckets { bucket := b.buckets[i] @@ -94,6 +103,21 @@ func (b *randBuckets) ListOperator() []*Operator { return ops } +// HasOperator checks whether any of the Regions has a waiting operator of the +// specified kind. +func (b *randBuckets) HasOperator(regionIDs []uint64, mask OpKind) bool { + b.mu.RLock() + defer b.mu.RUnlock() + for _, regionID := range regionIDs { + for _, op := range b.operatorsByRegion[regionID] { + if op.Kind()&mask != 0 { + return true + } + } + } + return false +} + // GetOperator gets an operator from the random buckets. func (b *randBuckets) GetOperator() []*Operator { b.mu.Lock() @@ -125,6 +149,9 @@ func (b *randBuckets) GetOperator() []*Operator { if len(bucket.ops) == 0 { b.totalWeight -= bucket.weight } + for _, op := range res { + b.removeOperatorFromRegion(op) + } return res } sum += proportion @@ -132,6 +159,24 @@ func (b *randBuckets) GetOperator() []*Operator { return nil } +func (b *randBuckets) removeOperatorFromRegion(op *Operator) { + regionID := op.RegionID() + ops := b.operatorsByRegion[regionID] + for i, waitingOp := range ops { + if waitingOp != op { + continue + } + ops[i] = ops[len(ops)-1] + ops = ops[:len(ops)-1] + break + } + if len(ops) == 0 { + delete(b.operatorsByRegion, regionID) + return + } + b.operatorsByRegion[regionID] = ops +} + // waitingOperatorStatus is used to limit the count of each kind of operators. type waitingOperatorStatus struct { mu syncutil.Mutex diff --git a/pkg/schedule/operator/waiting_operator_test.go b/pkg/schedule/operator/waiting_operator_test.go index b7b82321670..966480f3d9e 100644 --- a/pkg/schedule/operator/waiting_operator_test.go +++ b/pkg/schedule/operator/waiting_operator_test.go @@ -65,6 +65,24 @@ func TestListOperator(t *testing.T) { re.Len(rb.ListOperator(), len(priorityWeight)) } +func TestHasWaitingOperator(t *testing.T) { + re := require.New(t) + rb := newRandBuckets() + regionOp := NewTestOperator(1, &metapb.RegionEpoch{}, OpRegion, RemovePeer{FromStore: 1}) + replicaOp := NewTestOperator(1, &metapb.RegionEpoch{}, OpRegion|OpReplica, RemovePeer{FromStore: 2}) + rb.PutOperator(regionOp) + rb.PutOperator(replicaOp) + + re.True(rb.HasOperator([]uint64{2, 1}, OpReplica)) + re.False(rb.HasOperator([]uint64{2}, OpReplica)) + re.True(rb.HasOperator([]uint64{1}, OpRegion)) + + re.Equal(regionOp, rb.GetOperator()[0]) + re.True(rb.HasOperator([]uint64{1}, OpReplica)) + re.Equal(replicaOp, rb.GetOperator()[0]) + re.False(rb.HasOperator([]uint64{1}, OpReplica)) +} + func TestRandomBucketsWithMergeRegion(t *testing.T) { re := require.New(t) rb := newRandBuckets() From 720da3011f1a096de7347bb5d29e56d8c9ed6093 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 30 Jul 2026 14:07:17 +0800 Subject: [PATCH 06/10] schedule: reduce replicated-state query overhead Signed-off-by: Ryan Leung --- pkg/schedule/checker/checker_controller.go | 2 +- pkg/schedule/checker/replica_strategy.go | 148 ++++++-- pkg/schedule/checker/rule_checker.go | 324 ++++++++---------- pkg/schedule/handler/handler.go | 20 +- pkg/schedule/operator/operator_controller.go | 11 - pkg/schedule/operator/waiting_operator.go | 53 +-- .../operator/waiting_operator_test.go | 18 - 7 files changed, 278 insertions(+), 298 deletions(-) diff --git a/pkg/schedule/checker/checker_controller.go b/pkg/schedule/checker/checker_controller.go index e2149846b11..767dc40c2a0 100644 --- a/pkg/schedule/checker/checker_controller.go +++ b/pkg/schedule/checker/checker_controller.go @@ -266,12 +266,12 @@ func (c *Controller) checkPendingProcessedRegions() { pendingProcessedRegionsGauge.Set(float64(len(ids))) for _, id := range ids { region := c.cluster.GetRegion(id) + c.RemovePendingProcessedRegion(id) if region == nil { continue } // Remove the old entry before checking. A checker will add it back if // temporary conditions still prevent scheduling. - c.RemovePendingProcessedRegion(id) c.tryAddOperators(region) } } diff --git a/pkg/schedule/checker/replica_strategy.go b/pkg/schedule/checker/replica_strategy.go index 0d1b7929d21..0bd8501e592 100644 --- a/pkg/schedule/checker/replica_strategy.go +++ b/pkg/schedule/checker/replica_strategy.go @@ -30,15 +30,13 @@ import ( // ReplicaStrategy collects some utilities to manipulate region peers. It // exists to allow replica_checker and rule_checker to reuse common logics. type ReplicaStrategy struct { - checkerName string // replica-checker / rule-checker - cluster sche.CheckerCluster - locationLabels []string - isolationLevel string - region *core.RegionInfo - extraFilters []filter.Filter - fastFailover bool - storeCandidates []*core.StoreInfo - storeCandidatesPrepared bool + checkerName string // replica-checker / rule-checker + cluster sche.CheckerCluster + locationLabels []string + isolationLevel string + region *core.RegionInfo + extraFilters []filter.Filter + fastFailover bool } // SelectStoreToAdd returns the store to add a replica to a region. @@ -64,15 +62,11 @@ func (s *ReplicaStrategy) SelectStoreToAdd(coLocationStores []*core.StoreInfo, e if s.fastFailover { level = constant.Urgent } - stores := s.storeCandidates - filters := []filter.Filter{filter.NewExcludedFilter(s.checkerName, nil, s.region.GetStoreIDs())} - if !s.storeCandidatesPrepared { - stores = s.cluster.GetStores() - filters = append(filters, - filter.NewStorageThresholdFilter(s.checkerName), - filter.NewSpecialUseFilter(s.checkerName), - &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowTemporaryStates: true, OperatorLevel: level}, - ) + filters := []filter.Filter{ + filter.NewExcludedFilter(s.checkerName, nil, s.region.GetStoreIDs()), + filter.NewStorageThresholdFilter(s.checkerName), + filter.NewSpecialUseFilter(s.checkerName), + &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowTemporaryStates: true, OperatorLevel: level}, } if len(s.locationLabels) > 0 && s.isolationLevel != "" { filters = append(filters, filter.NewIsolationFilter(s.checkerName, s.isolationLevel, s.locationLabels, coLocationStores)) @@ -80,13 +74,13 @@ func (s *ReplicaStrategy) SelectStoreToAdd(coLocationStores []*core.StoreInfo, e if len(extraFilters) > 0 { filters = append(filters, extraFilters...) } - if !s.storeCandidatesPrepared && len(s.extraFilters) > 0 { + if len(s.extraFilters) > 0 { filters = append(filters, s.extraFilters...) } isolationComparer := filter.IsolationComparer(s.locationLabels, coLocationStores) strictStateFilter := &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowFastFailover: s.fastFailover, OperatorLevel: level} - targetCandidate := filter.NewCandidates(stores). + targetCandidate := filter.NewCandidates(s.cluster.GetStores()). FilterTarget(s.cluster.GetCheckerConfig(), nil, nil, filters...). KeepTheTopStores(isolationComparer, false) // greater isolation score is better if targetCandidate.Len() == 0 { @@ -100,9 +94,17 @@ func (s *ReplicaStrategy) SelectStoreToAdd(coLocationStores []*core.StoreInfo, e return target.GetID(), false } -// prepareStoreCandidates applies the Region-independent target filters so the -// result can be reused by placement-state checks using the same Rule. -func (s *ReplicaStrategy) prepareStoreCandidates(stores []*core.StoreInfo) { +type storeCandidateSet struct { + stores []*core.StoreInfo + filters []filter.Filter + next int + candidates []*core.StoreInfo +} + +// newStoreCandidateSet creates a lazily evaluated candidate set. Static Store +// filters are evaluated at most once per Rule, and scanning stops as soon as a +// Region finds a usable Store. +func (s *ReplicaStrategy) newStoreCandidateSet(stores []*core.StoreInfo) *storeCandidateSet { level := constant.High if s.fastFailover { level = constant.Urgent @@ -113,9 +115,57 @@ func (s *ReplicaStrategy) prepareStoreCandidates(stores []*core.StoreInfo) { &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowTemporaryStates: true, OperatorLevel: level}, } filters = append(filters, s.extraFilters...) - s.storeCandidates = filter.NewCandidates(stores). - FilterTarget(s.cluster.GetCheckerConfig(), nil, nil, filters...).Stores - s.storeCandidatesPrepared = true + return &storeCandidateSet{stores: stores, filters: filters} +} + +// hasStoreToAdd checks whether at least one Store can satisfy the Region- +// specific exclusion and isolation constraints. +func (s *ReplicaStrategy) hasStoreToAdd( + candidateSet *storeCandidateSet, + coLocationStores []*core.StoreInfo, + extraFilters ...filter.Filter, +) bool { + conf := s.cluster.GetCheckerConfig() + var isolationFilter filter.Filter + if len(s.locationLabels) > 0 && s.isolationLevel != "" { + isolationFilter = filter.NewIsolationFilter(s.checkerName, s.isolationLevel, s.locationLabels, coLocationStores) + } + matchesRegion := func(store *core.StoreInfo) bool { + if s.region.GetStorePeer(store.GetID()) != nil || + (isolationFilter != nil && !isolationFilter.Target(conf, store).IsOK()) { + return false + } + for _, extraFilter := range extraFilters { + if !extraFilter.Target(conf, store).IsOK() { + return false + } + } + return true + } + for _, store := range candidateSet.candidates { + if matchesRegion(store) { + return true + } + } + for candidateSet.next < len(candidateSet.stores) { + store := candidateSet.stores[candidateSet.next] + candidateSet.next++ + matched := true + for _, storeFilter := range candidateSet.filters { + if !storeFilter.Target(conf, store).IsOK() { + matched = false + break + } + } + if !matched { + continue + } + candidateSet.candidates = append(candidateSet.candidates, store) + if matchesRegion(store) { + return true + } + } + return false } // SelectStoreToFix returns a store to replace down/offline old peer. The location @@ -134,6 +184,52 @@ func (s *ReplicaStrategy) SelectStoreToFix(coLocationStores []*core.StoreInfo, o return s.SelectStoreToAdd(coLocationStores) } +func (s *ReplicaStrategy) hasStoreToFix(candidateSet *storeCandidateSet, coLocationStores []*core.StoreInfo, old uint64) bool { + if len(coLocationStores) == 0 { + return false + } + swapStoreToFirst(coLocationStores, old) + if len(coLocationStores) > 1 { + coLocationStores = coLocationStores[1:] + } + return s.hasStoreToAdd(candidateSet, coLocationStores) +} + +func (s *ReplicaStrategy) hasBetterLocation( + candidateSet *storeCandidateSet, + cluster sche.SharedCluster, + region *core.RegionInfo, + fit *placement.RegionFit, + ruleFit *placement.RuleFit, +) bool { + oldStoreID, _ := s.selectStoreToRemoveWithTempState(ruleFit.Stores) + if oldStoreID == 0 { + return false + } + oldStore := cluster.GetStore(oldStoreID) + if oldStore == nil { + return false + } + var coLocationStores []*core.StoreInfo + for _, store := range cluster.GetRegionStores(region) { + if store.GetLabelValue(core.EngineKey) != oldStore.GetLabelValue(core.EngineKey) { + continue + } + for _, rule := range fit.GetRules() { + if rule.Role == ruleFit.Rule.Role && placement.MatchLabelConstraints(store, rule.LabelConstraints) { + coLocationStores = append(coLocationStores, store) + break + } + } + } + if len(coLocationStores) == 0 { + return false + } + swapStoreToFirst(coLocationStores, oldStoreID) + locationImprover := filter.NewLocationImprover(s.checkerName, s.locationLabels, coLocationStores, oldStore) + return s.hasStoreToAdd(candidateSet, coLocationStores[1:], locationImprover) +} + // SelectStoreToImprove returns a store to replace oldStore. The location // placement after scheduling should be better than original. func (s *ReplicaStrategy) SelectStoreToImprove(coLocationStores []*core.StoreInfo, old uint64) (uint64, bool) { diff --git a/pkg/schedule/checker/rule_checker.go b/pkg/schedule/checker/rule_checker.go index 192ec12de75..1a61a461d83 100644 --- a/pkg/schedule/checker/rule_checker.go +++ b/pkg/schedule/checker/rule_checker.go @@ -88,27 +88,30 @@ type placementStateCandidateKey struct { type placementStateContext struct { storesLoaded bool stores []*core.StoreInfo - candidates map[placementStateCandidateKey][]*core.StoreInfo + candidates map[placementStateCandidateKey]*storeCandidateSet } -func (ctx *placementStateContext) getStrategy(c *RuleChecker, region *core.RegionInfo, rule *placement.Rule, fastFailover bool) *ReplicaStrategy { +func (ctx *placementStateContext) getStrategy( + c *RuleChecker, + region *core.RegionInfo, + rule *placement.Rule, + fastFailover bool, +) (*ReplicaStrategy, *storeCandidateSet) { strategy := c.strategy(region, rule, fastFailover) key := placementStateCandidateKey{rule: rule, fastFailover: fastFailover} - if candidates, ok := ctx.candidates[key]; ok { - strategy.storeCandidates = candidates - strategy.storeCandidatesPrepared = true - return strategy + if candidateSet, ok := ctx.candidates[key]; ok { + return strategy, candidateSet } if !ctx.storesLoaded { ctx.stores = c.cluster.GetStores() ctx.storesLoaded = true } if ctx.candidates == nil { - ctx.candidates = make(map[placementStateCandidateKey][]*core.StoreInfo) + ctx.candidates = make(map[placementStateCandidateKey]*storeCandidateSet) } - strategy.prepareStoreCandidates(ctx.stores) - ctx.candidates[key] = strategy.storeCandidates - return strategy + candidateSet := strategy.newStoreCandidateSet(ctx.stores) + ctx.candidates[key] = candidateSet + return strategy, candidateSet } // GetRegionPlacementState evaluates the placement state without creating an @@ -128,7 +131,7 @@ func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, cont if len(fit.RuleFits) == 0 { return RegionPlacementStateInProgress } - if c.hasOrphanPeerAction(region, fit) { + if c.planOrphanPeerAction(region, fit).canProgress() { return RegionPlacementStateInProgress } @@ -136,9 +139,8 @@ func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, cont for _, rf := range fit.RuleFits { if len(rf.Peers) < rf.Rule.Count { isWitness := rf.Rule.IsWitness && isWitnessEnabled(c.cluster) - strategy := context.getStrategy(c, region, rf.Rule, isWitness) - storeID, filteredByTemporaryState := strategy.SelectStoreToAdd(getRuleFitStores(c.cluster, rf)) - if storeID != 0 || filteredByTemporaryState || + strategy, candidateSet := context.getStrategy(c, region, rf.Rule, isWitness) + if strategy.hasStoreToAdd(candidateSet, getRuleFitStores(c.cluster, rf)) || c.hasStoreToSwapForMissingRule(region, fit, rf, context) { return RegionPlacementStateInProgress } @@ -166,9 +168,8 @@ func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, cont continue } - strategy := context.getStrategy(c, region, rf.Rule, fastFailover) - storeID, filteredByTemporaryState := strategy.SelectStoreToFix(getRuleFitStores(c.cluster, rf), peer.GetStoreId()) - if storeID != 0 || filteredByTemporaryState { + strategy, candidateSet := context.getStrategy(c, region, rf.Rule, fastFailover) + if strategy.hasStoreToFix(candidateSet, getRuleFitStores(c.cluster, rf), peer.GetStoreId()) { return RegionPlacementStateInProgress } hasUnfixablePlacement = true @@ -193,9 +194,8 @@ func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, cont continue } isWitness := rf.Rule.IsWitness && isWitnessEnabled(c.cluster) - strategy := context.getStrategy(c, region, rf.Rule, isWitness) - _, newStoreID, filteredByTemporaryState := strategy.getBetterLocation(c.cluster, region, fit, rf) - if newStoreID != 0 || filteredByTemporaryState { + strategy, candidateSet := context.getStrategy(c, region, rf.Rule, isWitness) + if strategy.hasBetterLocation(candidateSet, c.cluster, region, fit, rf) { return RegionPlacementStateInProgress } if !statistics.IsRegionLabelIsolationSatisfied( @@ -314,9 +314,34 @@ func (c *RuleChecker) hasAlternativeLeader(region *core.RegionInfo, fit *placeme return false } -func (c *RuleChecker) hasOrphanPeerAction(region *core.RegionInfo, fit *placement.RegionFit) bool { +type orphanPeerActionKind uint8 + +const ( + orphanPeerActionNone orphanPeerActionKind = iota + orphanPeerActionRemove + orphanPeerActionPromoteAndRemove + orphanPeerActionDemoteAndRemove +) + +type orphanPeerAction struct { + kind orphanPeerActionKind + desc string + peer *metapb.Peer + peerToRemove *metapb.Peer + retryable bool + schedulable bool + replacement bool + noFitCount int + replaceCount int +} + +func (a orphanPeerAction) canProgress() bool { + return a.retryable || (a.kind != orphanPeerActionNone && a.schedulable) +} + +func (c *RuleChecker) planOrphanPeerAction(region *core.RegionInfo, fit *placement.RegionFit) orphanPeerAction { if len(fit.OrphanPeers) == 0 { - return false + return orphanPeerAction{} } isPendingPeer := func(peer *metapb.Peer) bool { @@ -333,37 +358,46 @@ func (c *RuleChecker) hasOrphanPeerAction(region *core.RegionInfo, fit *placemen return region.GetLeader().GetId() != peer.GetId() || c.hasAlternativeLeader(region, fit, peer.GetId()) } + checkDownPeer := func(peers []*metapb.Peer) (*metapb.Peer, bool, bool) { + for _, peer := range peers { + down := isDownPeer(peer) + disconnected := isDisconnectedPeer(peer) + if disconnected || down { + retryable := disconnected && !down + if down { + retryable = !c.isStoreDownTimeHitMaxDownTime(peer.GetStoreId()) + } + return peer, true, retryable + } + if isPendingPeer(peer) { + return nil, true, true + } + } + return nil, false, false + } - allRulesSatisfiedAndHealthy := true + action := orphanPeerAction{} + hasUnhealthyFit := false var pinDownPeer *metapb.Peer for _, rf := range fit.RuleFits { if !rf.IsSatisfied() { - allRulesSatisfiedAndHealthy = false + hasUnhealthyFit = true break } - for _, peer := range rf.Peers { - if isPendingPeer(peer) { - return true - } - if isDownPeer(peer) { - if !c.isStoreDownTimeHitMaxDownTime(peer.GetStoreId()) { - return true - } - pinDownPeer = peer - allRulesSatisfiedAndHealthy = false - break - } - if isDisconnectedPeer(peer) { - return true - } - } - if !allRulesSatisfiedAndHealthy { + var retryable bool + pinDownPeer, hasUnhealthyFit, retryable = checkDownPeer(rf.Peers) + action.retryable = retryable + if hasUnhealthyFit { break } } - if allRulesSatisfiedAndHealthy { - // fixOrphanPeers always tries the first orphan peer in this case. - return canRemovePeer(fit.OrphanPeers[0]) + if !hasUnhealthyFit { + peer := fit.OrphanPeers[0] + action.kind = orphanPeerActionRemove + action.desc = "remove-orphan-peer" + action.peer = peer + action.schedulable = canRemovePeer(peer) + return action } if pinDownPeer != nil { @@ -378,23 +412,42 @@ func (c *RuleChecker) hasOrphanPeerAction(region *core.RegionInfo, fit *placemen } dstStore := c.cluster.GetStore(orphanPeer.GetStoreId()) if !fit.Replace(pinDownPeer.GetStoreId(), dstStore) { + action.noFitCount++ continue } + action.replaceCount++ destRole := pinDownPeer.GetRole() orphanRole := orphanPeer.GetRole() switch { case orphanRole == metapb.PeerRole_Learner && destRole == metapb.PeerRole_Voter: - return true + action.kind = orphanPeerActionPromoteAndRemove + action.desc = "replace-down-peer-with-orphan-peer" + action.peer = orphanPeer + action.peerToRemove = pinDownPeer + action.schedulable = true + action.replacement = true + return action case orphanRole == metapb.PeerRole_Voter && destRole == metapb.PeerRole_Learner: - return c.cluster.GetSharedConfig().IsUseJointConsensus() - case orphanRole == destRole: - return canRemovePeer(pinDownPeer) + action.kind = orphanPeerActionDemoteAndRemove + action.desc = "replace-down-peer-with-orphan-peer" + action.peer = orphanPeer + action.peerToRemove = pinDownPeer + action.schedulable = c.cluster.GetSharedConfig().IsUseJointConsensus() + action.replacement = true + return action + case orphanRole == destRole && !dstStore.IsDisconnected(): + action.kind = orphanPeerActionRemove + action.desc = "remove-replaced-orphan-peer" + action.peer = pinDownPeer + action.schedulable = canRemovePeer(pinDownPeer) + action.replacement = true + return action } } } if len(fit.OrphanPeers) < 2 { - return false + return action } var disconnectedPeer *metapb.Peer for _, orphanPeer := range fit.OrphanPeers { @@ -406,18 +459,25 @@ func (c *RuleChecker) hasOrphanPeerAction(region *core.RegionInfo, fit *placemen hasHealthyPeer := false for _, orphanPeer := range fit.OrphanPeers { if isPendingPeer(orphanPeer) || isDownPeer(orphanPeer) { - // fixOrphanPeers returns after trying the first unhealthy orphan. - return canRemovePeer(orphanPeer) + action.kind = orphanPeerActionRemove + action.desc = "remove-unhealthy-orphan-peer" + action.peer = orphanPeer + action.schedulable = canRemovePeer(orphanPeer) + return action } if hasHealthyPeer && fit.ExtraCount() > 0 { if disconnectedPeer != nil { - return canRemovePeer(disconnectedPeer) + orphanPeer = disconnectedPeer } - return canRemovePeer(orphanPeer) + action.kind = orphanPeerActionRemove + action.desc = "remove-orphan-peer" + action.peer = orphanPeer + action.schedulable = canRemovePeer(orphanPeer) + return action } hasHealthyPeer = true } - return false + return action } func (c *RuleChecker) hasStoreToSwapForMissingRule(region *core.RegionInfo, fit *placement.RegionFit, missingRuleFit *placement.RuleFit, context *placementStateContext) bool { @@ -432,9 +492,8 @@ func (c *RuleChecker) hasStoreToSwapForMissingRule(region *core.RegionInfo, fit } fastFailover := isWitnessEnabled(c.cluster) && store.IsTiKV() && oldRuleFit.Rule.IsWitness - strategy := context.getStrategy(c, region, oldRuleFit.Rule, fastFailover) - storeID, filteredByTemporaryState := strategy.SelectStoreToFix(getRuleFitStores(c.cluster, oldRuleFit), peer.GetStoreId()) - if storeID != 0 || filteredByTemporaryState { + strategy, candidateSet := context.getStrategy(c, region, oldRuleFit.Rule, fastFailover) + if strategy.hasStoreToFix(candidateSet, getRuleFitStores(c.cluster, oldRuleFit), peer.GetStoreId()) { return true } } @@ -831,145 +890,30 @@ func (c *RuleChecker) fixOrphanPeers(region *core.RegionInfo, fit *placement.Reg if len(fit.OrphanPeers) == 0 { return nil, nil } - - isPendingPeer := func(id uint64) bool { - for _, pendingPeer := range region.GetPendingPeers() { - if pendingPeer.GetId() == id { - return true - } - } - return false - } - - isDownPeer := func(id uint64) bool { - for _, downPeer := range region.GetDownPeers() { - if downPeer.Peer.GetId() == id { - return true - } - } - return false - } - - isUnhealthyPeer := func(id uint64) bool { - return isPendingPeer(id) || isDownPeer(id) - } - - isInDisconnectedStore := func(p *metapb.Peer) bool { - // avoid to meet down store when fix orphan peers, - // isInDisconnectedStore is usually more strictly than IsUnhealthy. - store := c.cluster.GetStore(p.GetStoreId()) - if store == nil { - return true - } - return store.IsDisconnected() + action := c.planOrphanPeerAction(region, fit) + if action.noFitCount > 0 { + ruleCheckerReplaceOrphanPeerNoFitCounter.Add(float64(action.noFitCount)) } - - checkDownPeer := func(peers []*metapb.Peer) (*metapb.Peer, bool) { - for _, p := range peers { - if isInDisconnectedStore(p) || isDownPeer(p.GetId()) { - return p, true - } - if isPendingPeer(p.GetId()) { - return nil, true - } - } - return nil, false + if action.replaceCount > 0 { + ruleCheckerReplaceOrphanPeerCounter.Add(float64(action.replaceCount)) } - - // remove orphan peers only when all rules are satisfied (count+role) and all peers selected - // by RuleFits is not pending or down. - var pinDownPeer *metapb.Peer - hasUnhealthyFit := false - for _, rf := range fit.RuleFits { - if !rf.IsSatisfied() { - hasUnhealthyFit = true - break - } - pinDownPeer, hasUnhealthyFit = checkDownPeer(rf.Peers) - if hasUnhealthyFit { - break - } + if action.kind == orphanPeerActionNone { + ruleCheckerSkipRemoveOrphanPeerCounter.Inc() + return nil, nil } - - // If hasUnhealthyFit is false, it is safe to delete the OrphanPeer. - if !hasUnhealthyFit { + if !action.replacement { ruleCheckerRemoveOrphanPeerCounter.Inc() - return operator.CreateRemovePeerOperator("remove-orphan-peer", c.cluster, 0, region, fit.OrphanPeers[0].StoreId) } - - // try to use orphan peers to replace unhealthy down peers. - for _, orphanPeer := range fit.OrphanPeers { - if pinDownPeer != nil { - if pinDownPeer.GetId() == orphanPeer.GetId() { - continue - } - // make sure the orphan peer is healthy. - if isUnhealthyPeer(orphanPeer.GetId()) || isInDisconnectedStore(orphanPeer) { - continue - } - // no consider witness in this path. - if pinDownPeer.GetIsWitness() || orphanPeer.GetIsWitness() { - continue - } - // pinDownPeer's store should be disconnected, because we use more strict judge before. - if !isInDisconnectedStore(pinDownPeer) { - continue - } - // check if down peer can replace with orphan peer. - dstStore := c.cluster.GetStore(orphanPeer.GetStoreId()) - if fit.Replace(pinDownPeer.GetStoreId(), dstStore) { - destRole := pinDownPeer.GetRole() - orphanPeerRole := orphanPeer.GetRole() - ruleCheckerReplaceOrphanPeerCounter.Inc() - switch { - case orphanPeerRole == metapb.PeerRole_Learner && destRole == metapb.PeerRole_Voter: - return operator.CreatePromoteLearnerOperatorAndRemovePeer("replace-down-peer-with-orphan-peer", c.cluster, region, orphanPeer, pinDownPeer) - case orphanPeerRole == metapb.PeerRole_Voter && destRole == metapb.PeerRole_Learner: - return operator.CreateDemoteLearnerOperatorAndRemovePeer("replace-down-peer-with-orphan-peer", c.cluster, region, orphanPeer, pinDownPeer) - case orphanPeerRole == destRole && isInDisconnectedStore(pinDownPeer) && !dstStore.IsDisconnected(): - return operator.CreateRemovePeerOperator("remove-replaced-orphan-peer", c.cluster, 0, region, pinDownPeer.GetStoreId()) - default: - // destRole should not same with orphanPeerRole. if role is same, it fit with orphanPeer should be better than now. - // destRole never be leader, so we not consider it. - } - } else { - ruleCheckerReplaceOrphanPeerNoFitCounter.Inc() - } - } - } - - extra := fit.ExtraCount() - // If hasUnhealthyFit is true, try to remove unhealthy orphan peers only if number of OrphanPeers is >= 2. - // Ref https://github.com/tikv/pd/issues/4045 - if len(fit.OrphanPeers) >= 2 { - hasHealthPeer := false - var disconnectedPeer *metapb.Peer - for _, orphanPeer := range fit.OrphanPeers { - if isInDisconnectedStore(orphanPeer) { - disconnectedPeer = orphanPeer - break - } - } - for _, orphanPeer := range fit.OrphanPeers { - if isUnhealthyPeer(orphanPeer.GetId()) { - ruleCheckerRemoveOrphanPeerCounter.Inc() - return operator.CreateRemovePeerOperator("remove-unhealthy-orphan-peer", c.cluster, 0, region, orphanPeer.StoreId) - } - // The healthy orphan peer can be removed to keep the high availability only if the peer count is greater than the rule requirement. - if hasHealthPeer && extra > 0 { - // there already exists a healthy orphan peer, so we can remove other orphan Peers. - ruleCheckerRemoveOrphanPeerCounter.Inc() - // if there exists a disconnected orphan peer, we will pick it to remove firstly. - if disconnectedPeer != nil { - return operator.CreateRemovePeerOperator("remove-orphan-peer", c.cluster, 0, region, disconnectedPeer.StoreId) - } - return operator.CreateRemovePeerOperator("remove-orphan-peer", c.cluster, 0, region, orphanPeer.StoreId) - } - hasHealthPeer = true - } + switch action.kind { + case orphanPeerActionRemove: + return operator.CreateRemovePeerOperator(action.desc, c.cluster, 0, region, action.peer.GetStoreId()) + case orphanPeerActionPromoteAndRemove: + return operator.CreatePromoteLearnerOperatorAndRemovePeer(action.desc, c.cluster, region, action.peer, action.peerToRemove) + case orphanPeerActionDemoteAndRemove: + return operator.CreateDemoteLearnerOperatorAndRemovePeer(action.desc, c.cluster, region, action.peer, action.peerToRemove) + default: + return nil, nil } - ruleCheckerSkipRemoveOrphanPeerCounter.Inc() - return nil, nil } func (c *RuleChecker) isDownPeer(region *core.RegionInfo, peer *metapb.Peer) bool { diff --git a/pkg/schedule/handler/handler.go b/pkg/schedule/handler/handler.go index fd4544117e2..4efed166681 100644 --- a/pkg/schedule/handler/handler.go +++ b/pkg/schedule/handler/handler.go @@ -1321,11 +1321,25 @@ func (h *Handler) CheckRegionsReplicated(startKeyHex, endKeyHex string) (string, } func hasReplicaOperator(controller *operator.Controller, regions []*core.RegionInfo) bool { - regionIDs := make([]uint64, 0, len(regions)) + operators := controller.GetOperatorsOfKind(operator.OpReplica) + for _, op := range controller.GetWaitingOperators() { + if op.Kind()&operator.OpReplica != 0 { + operators = append(operators, op) + } + } + if len(operators) == 0 { + return false + } + operatorRegions := make(map[uint64]struct{}, len(operators)) + for _, op := range operators { + operatorRegions[op.RegionID()] = struct{}{} + } for _, region := range regions { - regionIDs = append(regionIDs, region.GetID()) + if _, ok := operatorRegions[region.GetID()]; ok { + return true + } } - return controller.HasOperatorOfKind(regionIDs, operator.OpReplica) + return false } func regionsCoverRange(regions []*core.RegionInfo, startKey, endKey []byte) bool { diff --git a/pkg/schedule/operator/operator_controller.go b/pkg/schedule/operator/operator_controller.go index b940b009125..80cffcc399b 100644 --- a/pkg/schedule/operator/operator_controller.go +++ b/pkg/schedule/operator/operator_controller.go @@ -808,17 +808,6 @@ func (oc *Controller) GetWaitingOperators() []*Operator { return oc.wop.ListOperator() } -// HasOperatorOfKind checks whether any of the Regions has a running or waiting -// operator of the specified kind. -func (oc *Controller) HasOperatorOfKind(regionIDs []uint64, mask OpKind) bool { - for _, regionID := range regionIDs { - if op := oc.GetOperator(regionID); op != nil && op.Kind()&mask != 0 { - return true - } - } - return oc.wop.HasOperator(regionIDs, mask) -} - // GetOperatorsOfKind returns the running operators of the kind. func (oc *Controller) GetOperatorsOfKind(mask OpKind) []*Operator { operators := make([]*Operator, 0, oc.opNotifierQueue.len()) diff --git a/pkg/schedule/operator/waiting_operator.go b/pkg/schedule/operator/waiting_operator.go index 93076fb8da4..5709e776f6c 100644 --- a/pkg/schedule/operator/waiting_operator.go +++ b/pkg/schedule/operator/waiting_operator.go @@ -29,7 +29,6 @@ type WaitingOperator interface { PutMergeOperators(op []*Operator) GetOperator() []*Operator ListOperator() []*Operator - HasOperator(regionIDs []uint64, mask OpKind) bool } // bucket is used to maintain the operators created by a specific scheduler. @@ -40,10 +39,9 @@ type bucket struct { // randBuckets is an implementation of waiting operators type randBuckets struct { - mu syncutil.RWMutex - totalWeight float64 - buckets []*bucket - operatorsByRegion map[uint64][]*Operator + mu syncutil.RWMutex + totalWeight float64 + buckets []*bucket } // newRandBuckets creates a random buckets. @@ -54,10 +52,7 @@ func newRandBuckets() *randBuckets { weight: priorityWeight[i], }) } - return &randBuckets{ - buckets: buckets, - operatorsByRegion: make(map[uint64][]*Operator), - } + return &randBuckets{buckets: buckets} } // PutOperator puts an operator into the random buckets. @@ -70,7 +65,6 @@ func (b *randBuckets) PutOperator(op *Operator) { b.totalWeight += bucket.weight } bucket.ops = append(bucket.ops, op) - b.operatorsByRegion[op.RegionID()] = append(b.operatorsByRegion[op.RegionID()], op) } // PutMergeOperators puts two operators into the random buckets. @@ -86,9 +80,6 @@ func (b *randBuckets) PutMergeOperators(ops []*Operator) { b.totalWeight += bucket.weight } bucket.ops = append(bucket.ops, ops...) - for _, op := range ops { - b.operatorsByRegion[op.RegionID()] = append(b.operatorsByRegion[op.RegionID()], op) - } } // ListOperator lists all operator in the random buckets. @@ -103,21 +94,6 @@ func (b *randBuckets) ListOperator() []*Operator { return ops } -// HasOperator checks whether any of the Regions has a waiting operator of the -// specified kind. -func (b *randBuckets) HasOperator(regionIDs []uint64, mask OpKind) bool { - b.mu.RLock() - defer b.mu.RUnlock() - for _, regionID := range regionIDs { - for _, op := range b.operatorsByRegion[regionID] { - if op.Kind()&mask != 0 { - return true - } - } - } - return false -} - // GetOperator gets an operator from the random buckets. func (b *randBuckets) GetOperator() []*Operator { b.mu.Lock() @@ -149,9 +125,6 @@ func (b *randBuckets) GetOperator() []*Operator { if len(bucket.ops) == 0 { b.totalWeight -= bucket.weight } - for _, op := range res { - b.removeOperatorFromRegion(op) - } return res } sum += proportion @@ -159,24 +132,6 @@ func (b *randBuckets) GetOperator() []*Operator { return nil } -func (b *randBuckets) removeOperatorFromRegion(op *Operator) { - regionID := op.RegionID() - ops := b.operatorsByRegion[regionID] - for i, waitingOp := range ops { - if waitingOp != op { - continue - } - ops[i] = ops[len(ops)-1] - ops = ops[:len(ops)-1] - break - } - if len(ops) == 0 { - delete(b.operatorsByRegion, regionID) - return - } - b.operatorsByRegion[regionID] = ops -} - // waitingOperatorStatus is used to limit the count of each kind of operators. type waitingOperatorStatus struct { mu syncutil.Mutex diff --git a/pkg/schedule/operator/waiting_operator_test.go b/pkg/schedule/operator/waiting_operator_test.go index 966480f3d9e..b7b82321670 100644 --- a/pkg/schedule/operator/waiting_operator_test.go +++ b/pkg/schedule/operator/waiting_operator_test.go @@ -65,24 +65,6 @@ func TestListOperator(t *testing.T) { re.Len(rb.ListOperator(), len(priorityWeight)) } -func TestHasWaitingOperator(t *testing.T) { - re := require.New(t) - rb := newRandBuckets() - regionOp := NewTestOperator(1, &metapb.RegionEpoch{}, OpRegion, RemovePeer{FromStore: 1}) - replicaOp := NewTestOperator(1, &metapb.RegionEpoch{}, OpRegion|OpReplica, RemovePeer{FromStore: 2}) - rb.PutOperator(regionOp) - rb.PutOperator(replicaOp) - - re.True(rb.HasOperator([]uint64{2, 1}, OpReplica)) - re.False(rb.HasOperator([]uint64{2}, OpReplica)) - re.True(rb.HasOperator([]uint64{1}, OpRegion)) - - re.Equal(regionOp, rb.GetOperator()[0]) - re.True(rb.HasOperator([]uint64{1}, OpReplica)) - re.Equal(replicaOp, rb.GetOperator()[0]) - re.False(rb.HasOperator([]uint64{1}, OpReplica)) -} - func TestRandomBucketsWithMergeRegion(t *testing.T) { re := require.New(t) rb := newRandBuckets() From ddf9aea13f58a88c0e1335b918f7cc1deae3527c Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 30 Jul 2026 15:18:26 +0800 Subject: [PATCH 07/10] schedule: align placement state retries Signed-off-by: Ryan Leung --- pkg/schedule/checker/checker_controller.go | 4 +- pkg/schedule/checker/rule_checker.go | 14 +++- pkg/schedule/checker/rule_checker_test.go | 86 ++++++++++++++++++++++ 3 files changed, 97 insertions(+), 7 deletions(-) diff --git a/pkg/schedule/checker/checker_controller.go b/pkg/schedule/checker/checker_controller.go index 767dc40c2a0..2c328ddcbdc 100644 --- a/pkg/schedule/checker/checker_controller.go +++ b/pkg/schedule/checker/checker_controller.go @@ -266,12 +266,10 @@ func (c *Controller) checkPendingProcessedRegions() { pendingProcessedRegionsGauge.Set(float64(len(ids))) for _, id := range ids { region := c.cluster.GetRegion(id) - c.RemovePendingProcessedRegion(id) if region == nil { + c.RemovePendingProcessedRegion(id) continue } - // Remove the old entry before checking. A checker will add it back if - // temporary conditions still prevent scheduling. c.tryAddOperators(region) } } diff --git a/pkg/schedule/checker/rule_checker.go b/pkg/schedule/checker/rule_checker.go index 1a61a461d83..c72d1ea3493 100644 --- a/pkg/schedule/checker/rule_checker.go +++ b/pkg/schedule/checker/rule_checker.go @@ -121,11 +121,15 @@ func (c *RuleChecker) GetRegionPlacementState(region *core.RegionInfo) RegionPla } func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, context *placementStateContext) RegionPlacementState { - if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 || c.pendingProcessedRegions.Exists(region.GetID()) { + pendingProcessed := c.pendingProcessedRegions.Exists(region.GetID()) + if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 { return RegionPlacementStateInProgress } fit := c.ruleManager.FitRegionWithoutCache(c.cluster, region) if isRegionPlacementSatisfied(region, fit) { + if pendingProcessed { + return RegionPlacementStateInProgress + } return RegionPlacementStateReplicated } if len(fit.RuleFits) == 0 { @@ -209,6 +213,9 @@ func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, cont if hasUnfixablePlacement { return RegionPlacementStatePending } + if pendingProcessed { + return RegionPlacementStateInProgress + } return RegionPlacementStateReplicated } @@ -493,9 +500,8 @@ func (c *RuleChecker) hasStoreToSwapForMissingRule(region *core.RegionInfo, fit fastFailover := isWitnessEnabled(c.cluster) && store.IsTiKV() && oldRuleFit.Rule.IsWitness strategy, candidateSet := context.getStrategy(c, region, oldRuleFit.Rule, fastFailover) - if strategy.hasStoreToFix(candidateSet, getRuleFitStores(c.cluster, oldRuleFit), peer.GetStoreId()) { - return true - } + // addRulePeerWithOptions stops after trying the first matching peer. + return strategy.hasStoreToFix(candidateSet, getRuleFitStores(c.cluster, oldRuleFit), peer.GetStoreId()) } return false } diff --git a/pkg/schedule/checker/rule_checker_test.go b/pkg/schedule/checker/rule_checker_test.go index 28f726b17e7..04f3f7bc9c3 100644 --- a/pkg/schedule/checker/rule_checker_test.go +++ b/pkg/schedule/checker/rule_checker_test.go @@ -35,6 +35,7 @@ import ( "github.com/tikv/pd/pkg/mock/mockcluster" "github.com/tikv/pd/pkg/mock/mockconfig" "github.com/tikv/pd/pkg/schedule/config" + "github.com/tikv/pd/pkg/schedule/hbstream" "github.com/tikv/pd/pkg/schedule/operator" "github.com/tikv/pd/pkg/schedule/placement" "github.com/tikv/pd/pkg/utils/operatorutil" @@ -2378,6 +2379,56 @@ func (suite *ruleCheckerTestSuite) TestRegionPlacementFollowsOrphanRemovalOrder( re.Nil(suite.rc.Check(region)) } +func (suite *ruleCheckerTestSuite) TestRegionPlacementFollowsSwapPeerOrder() { + re := suite.Require() + suite.cluster.AddLabelsStore(1, 1, map[string]string{"rule": "a", "shared": "yes"}) + suite.cluster.AddLabelsStore(2, 1, map[string]string{"rule": "b", "shared": "yes"}) + suite.cluster.AddLabelsStore(3, 1, map[string]string{"rule": "b"}) + suite.cluster.AddLeaderRegion(1, 1, 2) + + for _, rule := range []*placement.Rule{ + { + GroupID: placement.DefaultGroupID, + ID: "rule-a", + Index: 100, + Role: placement.Voter, + Count: 1, + LabelConstraints: []placement.LabelConstraint{ + {Key: "rule", Op: placement.In, Values: []string{"a"}}, + }, + }, + { + GroupID: placement.DefaultGroupID, + ID: "rule-b", + Index: 101, + Role: placement.Voter, + Count: 1, + LabelConstraints: []placement.LabelConstraint{ + {Key: "rule", Op: placement.In, Values: []string{"b"}}, + }, + }, + { + GroupID: placement.DefaultGroupID, + ID: "missing", + Index: 102, + Role: placement.Learner, + Count: 1, + LabelConstraints: []placement.LabelConstraint{ + {Key: "shared", Op: placement.In, Values: []string{"yes"}}, + }, + }, + } { + re.NoError(suite.ruleManager.SetRule(rule)) + } + re.NoError(suite.ruleManager.DeleteRule(placement.DefaultGroupID, placement.DefaultRuleID)) + + region := suite.cluster.GetRegion(1) + // The first matching peer has no replacement target. The checker stops + // there even though a later matching peer could be replaced through Store 3. + re.Equal(RegionPlacementStatePending, suite.rc.GetRegionPlacementState(region)) + re.Nil(suite.rc.Check(region)) +} + func (suite *ruleCheckerTestSuite) TestRegionsPlacementStateLoadsStoresOnce() { re := suite.Require() suite.cluster.AddLeaderStore(1, 1) @@ -2445,6 +2496,41 @@ func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingProcessedIsInProgre re.Equal(RegionPlacementStateInProgress, checker.GetRegionPlacementState(suite.cluster.GetRegion(1))) } +func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingProcessedDoesNotHidePending() { + re := suite.Require() + suite.cluster.AddLeaderStore(1, 1) + suite.cluster.AddLeaderStore(2, 1) + suite.cluster.AddLeaderRegion(1, 1, 2) + pendingProcessed := cache.NewIDTTL(suite.ctx, time.Minute, 3*time.Minute) + pendingProcessed.Put(1, nil) + checker := NewRuleChecker(suite.ctx, suite.cluster, suite.ruleManager, pendingProcessed) + + re.Equal(RegionPlacementStatePending, checker.GetRegionPlacementState(suite.cluster.GetRegion(1))) +} + +func (suite *ruleCheckerTestSuite) TestPendingProcessedRetryRetainedWithoutLeader() { + re := suite.Require() + suite.cluster.AddLeaderStore(1, 1) + suite.cluster.AddLeaderRegion(1, 1) + region := suite.cluster.GetRegion(1).Clone(core.WithLeader(nil)) + suite.cluster.PutRegion(region) + + stream := hbstream.NewTestHeartbeatStreams(suite.ctx, suite.cluster, false) + defer stream.Close() + opController := operator.NewController( + suite.ctx, + suite.cluster.GetBasicCluster(), + suite.cluster.GetSharedConfig(), + stream, + ) + controller := NewController(suite.ctx, suite.cluster, suite.cluster.GetCheckerConfig(), opController) + controller.AddPendingProcessedRegions(false, region.GetID(), 2) + + controller.checkPendingProcessedRegions() + // The retryable Region is retained, while the vanished Region is removed. + re.Equal([]uint64{region.GetID()}, controller.GetPendingProcessedRegions()) +} + func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingOfflinePeer() { re := suite.Require() suite.cluster.AddLeaderStore(1, 1) From 9a9824817c7bcb054c6ac54e0102d660fa9bbc64 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 30 Jul 2026 17:21:57 +0800 Subject: [PATCH 08/10] tests: fix replicated state expectations Signed-off-by: Ryan Leung --- tests/integrations/client/http_client_test.go | 3 ++- tests/server/api/region_test.go | 17 +++++++++++++++++ 2 files changed, 19 insertions(+), 1 deletion(-) diff --git a/tests/integrations/client/http_client_test.go b/tests/integrations/client/http_client_test.go index 03c0399883c..92782196930 100644 --- a/tests/integrations/client/http_client_test.go +++ b/tests/integrations/client/http_client_test.go @@ -180,7 +180,8 @@ func (suite *httpClientTestSuite) TestMeta() { re.Len(regions.Regions, 3) state, err := client.GetRegionsReplicatedStateByKeyRange(ctx, pd.NewKeyRange([]byte("a1"), []byte("a3"))) re.NoError(err) - re.Equal("INPROGRESS", state) + // All candidate Stores are low on space, so no placement action can progress. + re.Equal("PENDING", state) regionStats, err := client.GetRegionStatusByKeyRange(ctx, pd.NewKeyRange([]byte("a1"), []byte("a3")), false) re.NoError(err) re.Positive(regionStats.Count) diff --git a/tests/server/api/region_test.go b/tests/server/api/region_test.go index 9fcf3589780..a0316cee55d 100644 --- a/tests/server/api/region_test.go +++ b/tests/server/api/region_test.go @@ -133,6 +133,11 @@ func (suite *regionTestSuite) checkAccelerateRegionsScheduleInRanges(cluster *te } func (suite *regionTestSuite) TestCheckRegionsReplicated() { + re := suite.Require() + re.NoError(failpoint.Enable("github.com/tikv/pd/pkg/schedule/checker/skipCheckSuspectRanges", "return(true)")) + defer func() { + re.NoError(failpoint.Disable("github.com/tikv/pd/pkg/schedule/checker/skipCheckSuspectRanges")) + }() suite.env.RunTest(suite.checkRegionsReplicated) } @@ -140,6 +145,18 @@ func (suite *regionTestSuite) checkRegionsReplicated(cluster *tests.TestCluster) re := suite.Require() pauseAllCheckers(re, cluster) leader := cluster.GetLeaderServer() + var checkerController *checker.Controller + if sche := cluster.GetSchedulingPrimaryServer(); sche == nil { + checkerController = leader.GetRaftCluster().GetCoordinator().GetCheckerController() + } else { + checkerController = sche.GetCluster().GetCoordinator().GetCheckerController() + } + checkerController.ClearSuspectKeyRanges() + checkerController.ClearPendingProcessedRegions() + defer func() { + checkerController.ClearSuspectKeyRanges() + checkerController.ClearPendingProcessedRegions() + }() urlPrefix := leader.GetAddr() + "/pd/api/v1" // add test region From a7813db3645ba2ebfea78767b3eec92d2eb5c516 Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 30 Jul 2026 17:51:25 +0800 Subject: [PATCH 09/10] schedule: ignore retry cache in placement state Signed-off-by: Ryan Leung --- pkg/schedule/checker/rule_checker.go | 7 ------- pkg/schedule/checker/rule_checker_test.go | 4 ++-- 2 files changed, 2 insertions(+), 9 deletions(-) diff --git a/pkg/schedule/checker/rule_checker.go b/pkg/schedule/checker/rule_checker.go index c72d1ea3493..c3798413ff0 100644 --- a/pkg/schedule/checker/rule_checker.go +++ b/pkg/schedule/checker/rule_checker.go @@ -121,15 +121,11 @@ func (c *RuleChecker) GetRegionPlacementState(region *core.RegionInfo) RegionPla } func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, context *placementStateContext) RegionPlacementState { - pendingProcessed := c.pendingProcessedRegions.Exists(region.GetID()) if region.GetLeader() == nil || len(region.GetPendingPeers()) > 0 { return RegionPlacementStateInProgress } fit := c.ruleManager.FitRegionWithoutCache(c.cluster, region) if isRegionPlacementSatisfied(region, fit) { - if pendingProcessed { - return RegionPlacementStateInProgress - } return RegionPlacementStateReplicated } if len(fit.RuleFits) == 0 { @@ -213,9 +209,6 @@ func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, cont if hasUnfixablePlacement { return RegionPlacementStatePending } - if pendingProcessed { - return RegionPlacementStateInProgress - } return RegionPlacementStateReplicated } diff --git a/pkg/schedule/checker/rule_checker_test.go b/pkg/schedule/checker/rule_checker_test.go index 04f3f7bc9c3..48b616c8180 100644 --- a/pkg/schedule/checker/rule_checker_test.go +++ b/pkg/schedule/checker/rule_checker_test.go @@ -2483,7 +2483,7 @@ func (suite *ruleCheckerTestSuite) TestRegionsPlacementStateSkipsStoresForComple re.Zero(cluster.getStoresCount) } -func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingProcessedIsInProgress() { +func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingProcessedDoesNotHideReplicated() { re := suite.Require() suite.cluster.AddLeaderStore(1, 1) suite.cluster.AddLeaderStore(2, 1) @@ -2493,7 +2493,7 @@ func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingProcessedIsInProgre pendingProcessed.Put(1, nil) checker := NewRuleChecker(suite.ctx, suite.cluster, suite.ruleManager, pendingProcessed) - re.Equal(RegionPlacementStateInProgress, checker.GetRegionPlacementState(suite.cluster.GetRegion(1))) + re.Equal(RegionPlacementStateReplicated, checker.GetRegionPlacementState(suite.cluster.GetRegion(1))) } func (suite *ruleCheckerTestSuite) TestRegionPlacementPendingProcessedDoesNotHidePending() { From 8e9a02f96d16bf0e4dffd75878247ad2c4ab369f Mon Sep 17 00:00:00 2001 From: Ryan Leung Date: Thu, 30 Jul 2026 18:28:56 +0800 Subject: [PATCH 10/10] schedule: deduplicate placement state helpers Signed-off-by: Ryan Leung --- pkg/schedule/checker/replica_strategy.go | 75 +++++++++--------------- pkg/schedule/checker/rule_checker.go | 56 ++++++------------ 2 files changed, 46 insertions(+), 85 deletions(-) diff --git a/pkg/schedule/checker/replica_strategy.go b/pkg/schedule/checker/replica_strategy.go index 0bd8501e592..e627d65a725 100644 --- a/pkg/schedule/checker/replica_strategy.go +++ b/pkg/schedule/checker/replica_strategy.go @@ -39,6 +39,13 @@ type ReplicaStrategy struct { fastFailover bool } +func (s *ReplicaStrategy) operatorLevel() constant.PriorityLevel { + if s.fastFailover { + return constant.Urgent + } + return constant.High +} + // SelectStoreToAdd returns the store to add a replica to a region. // `coLocationStores` are the stores used to compare location with target // store. @@ -58,10 +65,7 @@ func (s *ReplicaStrategy) SelectStoreToAdd(coLocationStores []*core.StoreInfo, e // // The reason for it is to prevent the non-optimal replica placement due // to the short-term state, resulting in redundant scheduling. - level := constant.High - if s.fastFailover { - level = constant.Urgent - } + level := s.operatorLevel() filters := []filter.Filter{ filter.NewExcludedFilter(s.checkerName, nil, s.region.GetStoreIDs()), filter.NewStorageThresholdFilter(s.checkerName), @@ -105,14 +109,10 @@ type storeCandidateSet struct { // filters are evaluated at most once per Rule, and scanning stops as soon as a // Region finds a usable Store. func (s *ReplicaStrategy) newStoreCandidateSet(stores []*core.StoreInfo) *storeCandidateSet { - level := constant.High - if s.fastFailover { - level = constant.Urgent - } filters := []filter.Filter{ filter.NewStorageThresholdFilter(s.checkerName), filter.NewSpecialUseFilter(s.checkerName), - &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowTemporaryStates: true, OperatorLevel: level}, + &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowTemporaryStates: true, OperatorLevel: s.operatorLevel()}, } filters = append(filters, s.extraFilters...) return &storeCandidateSet{stores: stores, filters: filters} @@ -195,6 +195,22 @@ func (s *ReplicaStrategy) hasStoreToFix(candidateSet *storeCandidateSet, coLocat return s.hasStoreToAdd(candidateSet, coLocationStores) } +func collectCoLocationStores(cluster sche.SharedCluster, region *core.RegionInfo, fit *placement.RegionFit, ruleFit *placement.RuleFit, oldStore *core.StoreInfo) []*core.StoreInfo { + var stores []*core.StoreInfo + for _, store := range cluster.GetRegionStores(region) { + if store.GetLabelValue(core.EngineKey) != oldStore.GetLabelValue(core.EngineKey) { + continue + } + for _, rule := range fit.GetRules() { + if rule.Role == ruleFit.Rule.Role && placement.MatchLabelConstraints(store, rule.LabelConstraints) { + stores = append(stores, store) + break + } + } + } + return stores +} + func (s *ReplicaStrategy) hasBetterLocation( candidateSet *storeCandidateSet, cluster sche.SharedCluster, @@ -210,18 +226,7 @@ func (s *ReplicaStrategy) hasBetterLocation( if oldStore == nil { return false } - var coLocationStores []*core.StoreInfo - for _, store := range cluster.GetRegionStores(region) { - if store.GetLabelValue(core.EngineKey) != oldStore.GetLabelValue(core.EngineKey) { - continue - } - for _, rule := range fit.GetRules() { - if rule.Role == ruleFit.Rule.Role && placement.MatchLabelConstraints(store, rule.LabelConstraints) { - coLocationStores = append(coLocationStores, store) - break - } - } - } + coLocationStores := collectCoLocationStores(cluster, region, fit, ruleFit, oldStore) if len(coLocationStores) == 0 { return false } @@ -263,12 +268,8 @@ func swapStoreToFirst(stores []*core.StoreInfo, id uint64) { // SelectStoreToRemove returns the best option to remove from the region. func (s *ReplicaStrategy) SelectStoreToRemove(coLocationStores []*core.StoreInfo) uint64 { isolationComparer := filter.IsolationComparer(s.locationLabels, coLocationStores) - level := constant.High - if s.fastFailover { - level = constant.Urgent - } source := filter.NewCandidates(coLocationStores). - FilterSource(s.cluster.GetCheckerConfig(), nil, nil, &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, OperatorLevel: level}). + FilterSource(s.cluster.GetCheckerConfig(), nil, nil, &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, OperatorLevel: s.operatorLevel()}). KeepTheTopStores(isolationComparer, true). PickTheTopStore(filter.RegionScoreComparer(s.cluster.GetCheckerConfig()), false) if source == nil { @@ -280,10 +281,7 @@ func (s *ReplicaStrategy) SelectStoreToRemove(coLocationStores []*core.StoreInfo func (s *ReplicaStrategy) selectStoreToRemoveWithTempState(coLocationStores []*core.StoreInfo) (uint64, bool) { isolationComparer := filter.IsolationComparer(s.locationLabels, coLocationStores) - level := constant.High - if s.fastFailover { - level = constant.Urgent - } + level := s.operatorLevel() sourceCandidate := filter.NewCandidates(coLocationStores). FilterSource(s.cluster.GetCheckerConfig(), nil, nil, &filter.StoreStateFilter{ActionScope: s.checkerName, MoveRegion: true, AllowTemporaryStates: true, OperatorLevel: level}). KeepTheTopStores(isolationComparer, true) @@ -311,22 +309,7 @@ func (s *ReplicaStrategy) getBetterLocation(cluster sche.SharedCluster, region * if oldStore == nil { return 0, 0, false } - var coLocationStores []*core.StoreInfo - regionStores := cluster.GetRegionStores(region) - for _, store := range regionStores { - if store.GetLabelValue(core.EngineKey) != oldStore.GetLabelValue(core.EngineKey) { - continue - } - for _, r := range fit.GetRules() { - if r.Role != rf.Rule.Role { - continue - } - if placement.MatchLabelConstraints(store, r.LabelConstraints) { - coLocationStores = append(coLocationStores, store) - break - } - } - } + coLocationStores := collectCoLocationStores(cluster, region, fit, rf, oldStore) newStoreID, filterByTempState = s.SelectStoreToImprove(coLocationStores, oldStoreID) if sourceFilterByTempState { if newStoreID != 0 || filterByTempState { diff --git a/pkg/schedule/checker/rule_checker.go b/pkg/schedule/checker/rule_checker.go index c3798413ff0..4fcaa0a8544 100644 --- a/pkg/schedule/checker/rule_checker.go +++ b/pkg/schedule/checker/rule_checker.go @@ -213,13 +213,10 @@ func (c *RuleChecker) evaluateRegionPlacementState(region *core.RegionInfo, cont } func isRegionPlacementSatisfied(region *core.RegionInfo, fit *placement.RegionFit) bool { - if len(fit.RuleFits) == 0 || len(fit.OrphanPeers) > 0 { + if !fit.IsSatisfied() { return false } for _, rf := range fit.RuleFits { - if !rf.IsSatisfied() { - return false - } if len(rf.Stores) != len(rf.Peers) { return false } @@ -339,6 +336,16 @@ func (a orphanPeerAction) canProgress() bool { return a.retryable || (a.kind != orphanPeerActionNone && a.schedulable) } +func (a orphanPeerAction) withAction(kind orphanPeerActionKind, desc string, peer, peerToRemove *metapb.Peer, schedulable, replacement bool) orphanPeerAction { + a.kind = kind + a.desc = desc + a.peer = peer + a.peerToRemove = peerToRemove + a.schedulable = schedulable + a.replacement = replacement + return a +} + func (c *RuleChecker) planOrphanPeerAction(region *core.RegionInfo, fit *placement.RegionFit) orphanPeerAction { if len(fit.OrphanPeers) == 0 { return orphanPeerAction{} @@ -393,11 +400,7 @@ func (c *RuleChecker) planOrphanPeerAction(region *core.RegionInfo, fit *placeme } if !hasUnhealthyFit { peer := fit.OrphanPeers[0] - action.kind = orphanPeerActionRemove - action.desc = "remove-orphan-peer" - action.peer = peer - action.schedulable = canRemovePeer(peer) - return action + return action.withAction(orphanPeerActionRemove, "remove-orphan-peer", peer, nil, canRemovePeer(peer), false) } if pinDownPeer != nil { @@ -420,28 +423,11 @@ func (c *RuleChecker) planOrphanPeerAction(region *core.RegionInfo, fit *placeme orphanRole := orphanPeer.GetRole() switch { case orphanRole == metapb.PeerRole_Learner && destRole == metapb.PeerRole_Voter: - action.kind = orphanPeerActionPromoteAndRemove - action.desc = "replace-down-peer-with-orphan-peer" - action.peer = orphanPeer - action.peerToRemove = pinDownPeer - action.schedulable = true - action.replacement = true - return action + return action.withAction(orphanPeerActionPromoteAndRemove, "replace-down-peer-with-orphan-peer", orphanPeer, pinDownPeer, true, true) case orphanRole == metapb.PeerRole_Voter && destRole == metapb.PeerRole_Learner: - action.kind = orphanPeerActionDemoteAndRemove - action.desc = "replace-down-peer-with-orphan-peer" - action.peer = orphanPeer - action.peerToRemove = pinDownPeer - action.schedulable = c.cluster.GetSharedConfig().IsUseJointConsensus() - action.replacement = true - return action + return action.withAction(orphanPeerActionDemoteAndRemove, "replace-down-peer-with-orphan-peer", orphanPeer, pinDownPeer, c.cluster.GetSharedConfig().IsUseJointConsensus(), true) case orphanRole == destRole && !dstStore.IsDisconnected(): - action.kind = orphanPeerActionRemove - action.desc = "remove-replaced-orphan-peer" - action.peer = pinDownPeer - action.schedulable = canRemovePeer(pinDownPeer) - action.replacement = true - return action + return action.withAction(orphanPeerActionRemove, "remove-replaced-orphan-peer", pinDownPeer, nil, canRemovePeer(pinDownPeer), true) } } } @@ -459,21 +445,13 @@ func (c *RuleChecker) planOrphanPeerAction(region *core.RegionInfo, fit *placeme hasHealthyPeer := false for _, orphanPeer := range fit.OrphanPeers { if isPendingPeer(orphanPeer) || isDownPeer(orphanPeer) { - action.kind = orphanPeerActionRemove - action.desc = "remove-unhealthy-orphan-peer" - action.peer = orphanPeer - action.schedulable = canRemovePeer(orphanPeer) - return action + return action.withAction(orphanPeerActionRemove, "remove-unhealthy-orphan-peer", orphanPeer, nil, canRemovePeer(orphanPeer), false) } if hasHealthyPeer && fit.ExtraCount() > 0 { if disconnectedPeer != nil { orphanPeer = disconnectedPeer } - action.kind = orphanPeerActionRemove - action.desc = "remove-orphan-peer" - action.peer = orphanPeer - action.schedulable = canRemovePeer(orphanPeer) - return action + return action.withAction(orphanPeerActionRemove, "remove-orphan-peer", orphanPeer, nil, canRemovePeer(orphanPeer), false) } hasHealthyPeer = true }