| 215 | } |
| 216 | |
| 217 | func (e *Etcd) electLoop(ctx context.Context) { |
| 218 | defer e.wg.Done() |
| 219 | for { |
| 220 | select { |
| 221 | case <-e.quitCh: |
| 222 | return |
| 223 | default: |
| 224 | } |
| 225 | |
| 226 | reset: |
| 227 | session, err := concurrency.NewSession(e.client, concurrency.WithTTL(sessionTTL)) |
| 228 | if err != nil { |
| 229 | logger.Get().With( |
| 230 | zap.Error(err), |
| 231 | ).Error("Failed to create session") |
| 232 | time.Sleep(sessionTTL / 3) |
| 233 | continue |
| 234 | } |
| 235 | election := concurrency.NewElection(session, e.electPath) |
| 236 | e.electionCh <- election |
| 237 | for { |
| 238 | if err := election.Campaign(ctx, e.myID); err != nil { |
| 239 | logger.Get().With( |
| 240 | zap.Error(err), |
| 241 | ).Error("Failed to acquire the leader campaign") |
| 242 | continue |
| 243 | } |
| 244 | select { |
| 245 | case <-session.Done(): |
| 246 | logger.Get().Warn("Leader session is done") |
| 247 | goto reset |
| 248 | case <-e.quitCh: |
| 249 | logger.Get().Info("Exit the leader election loop") |
| 250 | return |
| 251 | } |
| 252 | } |
| 253 | } |
| 254 | } |
| 255 | |
| 256 | func (e *Etcd) observeLeaderEvent(ctx context.Context) { |
| 257 | defer e.wg.Done() |