Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
leaderelection.py192 linesDownload Raw Back to leaderelection
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 
codekingpro/portable-devtools · Team Ai