From 61d6b4b7c1d925dc34b3a8ba68d3169dece699da Mon Sep 17 00:00:00 2001 From: "roman.mazhut" Date: Wed, 8 Oct 2025 13:47:37 +0000 Subject: [PATCH] Placement validation and auto-resignation added to election manager --- src/aggregator/aggregator/election_mgr.go | 67 ++++++++++++++++++- .../aggregator/election_mgr_test.go | 55 +++++++++++++++ 2 files changed, 121 insertions(+), 1 deletion(-) diff --git a/src/aggregator/aggregator/election_mgr.go b/src/aggregator/aggregator/election_mgr.go index ab5ed0a235..8a9f9fa167 100644 --- a/src/aggregator/aggregator/election_mgr.go +++ b/src/aggregator/aggregator/election_mgr.go @@ -348,12 +348,13 @@ func (mgr *electionManager) Open(shardSetID uint32) error { } mgr.state = electionManagerOpen - mgr.Add(5) + mgr.Add(6) go mgr.watchGoalStateChanges(stateChangeWatch) go mgr.verifyPendingFollower(verifyWatch) go mgr.checkCampaignStateLoop() go mgr.campaignLoop(campaignStateWatch) go mgr.reportMetrics() + go mgr.validateLeaderInPlacementLoop() mgr.logger.Info("election manager opened successfully") return nil @@ -890,3 +891,67 @@ func (mgr *electionManager) logError(desc string, err error) { zap.Error(err), ) } + +// validateLeaderInPlacementLoop periodically checks if the current leader is still in placement. +// If the leader is not in placement, it resigns from leadership to trigger re-election. +func (mgr *electionManager) validateLeaderInPlacementLoop() { + defer mgr.Done() + + ticker := time.NewTicker(mgr.campaignStateCheckInterval) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + if err := mgr.validateAndResignIfNeeded(); err != nil { + mgr.logError("error validating leader placement", err) + } + case <-mgr.doneCh: + return + } + } +} + +func (mgr *electionManager) validateAndResignIfNeeded() error { + // Only validate if we are currently the leader + if mgr.ElectionState() != LeaderState { + return nil + } + + // Get current leader from leader service + leader, err := mgr.leaderService.Leader(mgr.electionKey) + if err != nil { + return fmt.Errorf("error getting current leader: %v", err) + } + + // If we're not the leader anymore, no need to validate + if leader != mgr.leaderValue { + return nil + } + + // Get current placement + placement, err := mgr.placementManager.Placement() + if err != nil { + return fmt.Errorf("error getting placement: %v", err) + } + + // Check if leader exists in placement + _, exists := placement.Instance(leader) + if !exists { + mgr.logger.Warn("current leader is not in placement, resigning", + zap.String("leader", leader), + zap.String("electionKey", mgr.electionKey)) + + // Resign leadership since leader is not in placement + ctx, cancel := context.WithTimeout(context.Background(), mgr.electionOpts.ResignTimeout()) + defer cancel() + + if err := mgr.Resign(ctx); err != nil { + return fmt.Errorf("error resigning leadership: %v", err) + } + + mgr.logger.Info("successfully resigned leadership due to leader not in placement") + } + + return nil +} diff --git a/src/aggregator/aggregator/election_mgr_test.go b/src/aggregator/aggregator/election_mgr_test.go index 2a7924a5af..5b455ea693 100644 --- a/src/aggregator/aggregator/election_mgr_test.go +++ b/src/aggregator/aggregator/election_mgr_test.go @@ -1204,3 +1204,58 @@ type enabledRes struct { result bool err error } + +func TestElectionManagerValidateLeaderInPlacementRemovedFromPlacement(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + leaderValue := "test-leader" + shardSetID := uint32(1) + electionKey := "/shardset/1/lock" + + updatedPlacement := placement.NewMockPlacement(ctrl) + updatedPlacement.EXPECT().Instance(leaderValue).Return(nil, false).AnyTimes() + + placementManager := NewMockPlacementManager(ctrl) + placementManager.EXPECT().Placement().Return(updatedPlacement, nil).AnyTimes() + placementManager.EXPECT().Shards().Return(shard.NewShards([]shard.Shard{shard.NewShard(0)}), nil).AnyTimes() + + leaderService := services.NewMockLeaderService(ctrl) + leaderService.EXPECT().Leader(electionKey).Return(leaderValue, nil).AnyTimes() + + resignCalled := make(chan struct{}, 1) + leaderService.EXPECT().Resign(electionKey).DoAndReturn(func(string) error { + select { + case resignCalled <- struct{}{}: + default: + } + return nil + }).AnyTimes() + + leaderService.EXPECT().Campaign(gomock.Any(), gomock.Any()).DoAndReturn(func(string, services.CampaignOptions) (<-chan campaign.Status, error) { + return make(chan campaign.Status), nil + }).AnyTimes() + + opts := testElectionManagerOptions(t, ctrl). + SetPlacementManager(placementManager). + SetLeaderService(leaderService). + SetCampaignStateCheckInterval(50 * time.Millisecond) + + mgr := NewElectionManager(opts).(*electionManager) + mgr.leaderValue = leaderValue + + err := mgr.Open(shardSetID) + require.NoError(t, err) + + mgr.electionStateWatchable.Update(LeaderState) + + select { + case <-resignCalled: + case <-time.After(5 * time.Second): + t.Fatal("Expected Resign() to be called") + } + + mgr.Close() +}