codekingpro/portable-devtools
114k
1# Copyright 2021 The Kubernetes Authors.2#3# Licensed under the Apache License, Version 2.0 (the "License");4# you may not use this file except in compliance with the License.5# You may obtain a copy of the License at6#7# http://www.apache.org/licenses/LICENSE-2.08#9# Unless required by applicable law or agreed to in writing, software10# distributed under the License is distributed on an "AS IS" BASIS,11# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.12# See the License for the specific language governing permissions and13# limitations under the License.14 15import datetime16import sys17import time18import json19import threading20from .leaderelectionrecord import LeaderElectionRecord21import logging22# if condition to be removed when support for python2 will be removed23if sys.version_info > (3, 0):24 from http import HTTPStatus25else:26 import httplib27logger = logging.getLogger("leaderelection")28 29"""30This package implements leader election using an annotation in a Kubernetes object.31The onstarted_leading function is run in a thread and when it returns, if it does 32it might not be safe to run it again in a process.33 34At first all candidates are considered followers. The one to create a lock or update35an existing lock first becomes the leader and remains so until it keeps renewing its36lease.37"""38 39 40class LeaderElection:41 def __init__(self, election_config):42 if election_config is None:43 sys.exit("argument config not passed")44 45 # Latest record observed in the created lock object46 self.observed_record = None47 48 # The configuration set for this candidate49 self.election_config = election_config50 51 # Latest update time of the lock52 self.observed_time_milliseconds = 053 54 # Point of entry to Leader election55 def run(self):56 # Try to create/ acquire a lock57 if self.acquire():58 logger.info("{} successfully acquired lease".format(self.election_config.lock.identity))59 60 # Start leading and call OnStartedLeading()61 threading.daemon = True62 threading.Thread(target=self.election_config.onstarted_leading).start()63 64 self.renew_loop()65 66 # Failed to update lease, run OnStoppedLeading callback67 self.election_config.onstopped_leading()68 69 def acquire(self):70 # Follower71 logger.info("{} is a follower".format(self.election_config.lock.identity))72 retry_period = self.election_config.retry_period73 74 while True:75 succeeded = self.try_acquire_or_renew()76 77 if succeeded:78 return True79 80 time.sleep(retry_period)81 82 def renew_loop(self):83 # Leader84 logger.info("Leader has entered renew loop and will try to update lease continuously")85 86 retry_period = self.election_config.retry_period87 renew_deadline = self.election_config.renew_deadline * 100088 89 while True:90 timeout = int(time.time() * 1000) + renew_deadline91 succeeded = False92 93 while int(time.time() * 1000) < timeout:94 succeeded = self.try_acquire_or_renew()95 96 if succeeded:97 break98 time.sleep(retry_period)99 100 if succeeded:101 time.sleep(retry_period)102 continue103 104 # failed to renew, return105 return106 107 def try_acquire_or_renew(self):108 now_timestamp = time.time()109 now = datetime.datetime.fromtimestamp(now_timestamp)110 111 # Check if lock is created112 lock_status, old_election_record = self.election_config.lock.get(self.election_config.lock.name,113 self.election_config.lock.namespace)114 115 # create a default Election record for this candidate116 leader_election_record = LeaderElectionRecord(self.election_config.lock.identity,117 str(self.election_config.lease_duration), str(now), str(now))118 119 # A lock is not created with that name, try to create one120 if not lock_status:121 # To be removed when support for python2 will be removed122 if sys.version_info > (3, 0):123 if json.loads(old_election_record.body)['code'] != HTTPStatus.NOT_FOUND:124 logger.info("Error retrieving resource lock {} as {}".format(self.election_config.lock.name,125 old_election_record.reason))126 return False127 else:128 if json.loads(old_election_record.body)['code'] != httplib.NOT_FOUND:129 logger.info("Error retrieving resource lock {} as {}".format(self.election_config.lock.name,130 old_election_record.reason))131 return False132 133 logger.info("{} is trying to create a lock".format(leader_election_record.holder_identity))134 create_status = self.election_config.lock.create(name=self.election_config.lock.name,135 namespace=self.election_config.lock.namespace,136 election_record=leader_election_record)137 138 if create_status is False:139 logger.info("{} Failed to create lock".format(leader_election_record.holder_identity))140 return False141 142 self.observed_record = leader_election_record143 self.observed_time_milliseconds = int(time.time() * 1000)144 return True145 146 # A lock exists with that name147 # Validate old_election_record148 if old_election_record is None:149 # try to update lock with proper annotation and election record150 return self.update_lock(leader_election_record)151 152 if (old_election_record.holder_identity is None or old_election_record.lease_duration is None153 or old_election_record.acquire_time is None or old_election_record.renew_time is None):154 # try to update lock with proper annotation and election record155 return self.update_lock(leader_election_record)156 157 # Report transitions158 if self.observed_record and self.observed_record.holder_identity != old_election_record.holder_identity:159 logger.info("Leader has switched to {}".format(old_election_record.holder_identity))160 161 if self.observed_record is None or old_election_record.__dict__ != self.observed_record.__dict__:162 self.observed_record = old_election_record163 self.observed_time_milliseconds = int(time.time() * 1000)164 165 # If This candidate is not the leader and lease duration is yet to finish166 if (self.election_config.lock.identity != self.observed_record.holder_identity167 and self.observed_time_milliseconds + self.election_config.lease_duration * 1000 > int(now_timestamp * 1000)):168 logger.info("yet to finish lease_duration, lease held by {} and has not expired".format(old_election_record.holder_identity))169 return False170 171 # If this candidate is the Leader172 if self.election_config.lock.identity == self.observed_record.holder_identity:173 # Leader updates renewTime, but keeps acquire_time unchanged174 leader_election_record.acquire_time = self.observed_record.acquire_time175 176 return self.update_lock(leader_election_record)177 178 def update_lock(self, leader_election_record):179 # Update object with latest election record180 update_status = self.election_config.lock.update(self.election_config.lock.name,181 self.election_config.lock.namespace,182 leader_election_record)183 184 if update_status is False:185 logger.info("{} failed to acquire lease".format(leader_election_record.holder_identity))186 return False187 188 self.observed_record = leader_election_record189 self.observed_time_milliseconds = int(time.time() * 1000)190 logger.info("leader {} has successfully acquired lease".format(leader_election_record.holder_identity))191 return True192 