Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 66 additions & 1 deletion src/aggregator/aggregator/election_mgr.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
}
55 changes: 55 additions & 0 deletions src/aggregator/aggregator/election_mgr_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
}