Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
pooling.py330 linesDownload Raw Back to core
1"""2"""3 4# Created on 2014.03.145#6# Author: Giovanni Cannata7#8# Copyright 2014 - 2020 Giovanni Cannata9#10# This file is part of ldap3.11#12# ldap3 is free software: you can redistribute it and/or modify13# it under the terms of the GNU Lesser General Public License as published14# by the Free Software Foundation, either version 3 of the License, or15# (at your option) any later version.16#17# ldap3 is distributed in the hope that it will be useful,18# but WITHOUT ANY WARRANTY; without even the implied warranty of19# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the20# GNU Lesser General Public License for more details.21#22# You should have received a copy of the GNU Lesser General Public License23# along with ldap3 in the COPYING and COPYING.LESSER files.24# If not, see <http://www.gnu.org/licenses/>.25 26from datetime import datetime, MINYEAR27from os import linesep28from random import randint29from time import sleep30 31from .. import FIRST, ROUND_ROBIN, RANDOM, SEQUENCE_TYPES, STRING_TYPES, get_config_parameter32from .exceptions import LDAPUnknownStrategyError, LDAPServerPoolError, LDAPServerPoolExhaustedError33from .server import Server34from ..utils.log import log, log_enabled, ERROR, BASIC, NETWORK35 36POOLING_STRATEGIES = [FIRST, ROUND_ROBIN, RANDOM]37 38 39class ServerState(object):40    def __init__(self, server, last_checked_time, available):41        self.server = server42        self.last_checked_time = last_checked_time43        self.available = available44 45 46class ServerPoolState(object):47    def __init__(self, server_pool):48        self.server_states = []  # each element is a ServerState49        self.strategy = server_pool.strategy50        self.server_pool = server_pool51        self.last_used_server = 052        self.refresh()53        self.initialize_time = datetime.now()54 55        if log_enabled(BASIC):56            log(BASIC, 'instantiated ServerPoolState: <%r>', self)57 58    def __str__(self):59        s = 'servers: ' + linesep60        if self.server_states:61            for state in self.server_states:62                s += str(state.server) + linesep63        else:64            s += 'None' + linesep65        s += 'Pool strategy: ' + str(self.strategy) + linesep66        s += ' - Last used server: ' + ('None' if self.last_used_server == -1 else str(self.server_states[self.last_used_server].server))67 68        return s69 70    def refresh(self):71        self.server_states = []72        for server in self.server_pool.servers:73            self.server_states.append(ServerState(server, datetime(MINYEAR, 1, 1), True))  # server, smallest date ever, supposed available74        self.last_used_server = randint(0, len(self.server_states) - 1)75 76    def get_current_server(self):77        return self.server_states[self.last_used_server].server78 79    def get_server(self):80        if self.server_states:81            if self.server_pool.strategy == FIRST:82                if self.server_pool.active:83                    # returns the first active server84                    self.last_used_server = self.find_active_server(starting=0)85                else:86                    # returns always the first server - no pooling87                    self.last_used_server = 088            elif self.server_pool.strategy == ROUND_ROBIN:89                if self.server_pool.active:90                    # returns the next active server in a circular range91                    self.last_used_server = self.find_active_server(self.last_used_server + 1)92                else:93                    # returns the next server in a circular range94                    self.last_used_server = self.last_used_server + 1 if (self.last_used_server + 1) < len(self.server_states) else 095            elif self.server_pool.strategy == RANDOM:96                if self.server_pool.active:97                    self.last_used_server = self.find_active_random_server()98                else:99                    # returns a random server in the pool100                    self.last_used_server = randint(0, len(self.server_states) - 1)101            else:102                if log_enabled(ERROR):103                    log(ERROR, 'unknown server pooling strategy <%s>', self.server_pool.strategy)104                raise LDAPUnknownStrategyError('unknown server pooling strategy')105            if log_enabled(BASIC):106                log(BASIC, 'server returned from Server Pool: <%s>', self.last_used_server)107            return self.server_states[self.last_used_server].server108        else:109            if log_enabled(ERROR):110                log(ERROR, 'no servers in Server Pool <%s>', self)111            raise LDAPServerPoolError('no servers in server pool')112 113    def find_active_random_server(self):114        counter = self.server_pool.active  # can be True for "forever" or the number of cycles to try115        while counter:116            if log_enabled(NETWORK):117                log(NETWORK, 'entering loop for finding active server in pool <%s>', self)118            temp_list = self.server_states[:]  # copy119            while temp_list:120                # pops a random server from a temp list and checks its121                # availability, if not available tries another one122                server_state = temp_list.pop(randint(0, len(temp_list) - 1))123                if not server_state.available:  # server is offline124                    if (isinstance(self.server_pool.exhaust, bool) and self.server_pool.exhaust) or (datetime.now() - server_state.last_checked_time).seconds < self.server_pool.exhaust:  # keeps server offline125                        if log_enabled(NETWORK):126                            log(NETWORK, 'server <%s> excluded from checking because it is offline', server_state.server)127                        continue128                    if log_enabled(NETWORK):129                            log(NETWORK, 'server <%s> reinserted in pool', server_state.server)130                server_state.last_checked_time = datetime.now()131                if log_enabled(NETWORK):132                    log(NETWORK, 'checking server <%s> for availability', server_state.server)133                if server_state.server.check_availability():134                    # returns a random active server in the pool135                    server_state.available = True136                    return self.server_states.index(server_state)137                else:138                    server_state.available = False139            if not isinstance(self.server_pool.active, bool):140                counter -= 1141        if log_enabled(ERROR):142            log(ERROR, 'no random active server available in Server Pool <%s> after maximum number of tries', self)143        raise LDAPServerPoolExhaustedError('no random active server available in server pool after maximum number of tries')144 145    def find_active_server(self, starting):146        conf_pool_timeout = get_config_parameter('POOLING_LOOP_TIMEOUT')147        counter = self.server_pool.active  # can be True for "forever" or the number of cycles to try148        if starting >= len(self.server_states):149            starting = 0150 151        while counter:152            if log_enabled(NETWORK):153                log(NETWORK, 'entering loop number <%s> for finding active server in pool <%s>', counter, self)154            index = -1155            pool_size = len(self.server_states)156            while index < pool_size - 1:157                index += 1158                offset = index + starting if index + starting < pool_size else index + starting - pool_size159                server_state = self.server_states[offset]160                if not server_state.available:  # server is offline161                    if (isinstance(self.server_pool.exhaust, bool) and self.server_pool.exhaust) or (datetime.now() - server_state.last_checked_time).seconds < self.server_pool.exhaust:  # keeps server offline162                        if log_enabled(NETWORK):163                            if isinstance(self.server_pool.exhaust, bool):164                                log(NETWORK, 'server <%s> excluded from checking because is offline', server_state.server)165                            else:166                                log(NETWORK, 'server <%s> excluded from checking because is offline for %d seconds', server_state.server, (self.server_pool.exhaust - (datetime.now() - server_state.last_checked_time).seconds))167                        continue168                    if log_enabled(NETWORK):169                            log(NETWORK, 'server <%s> reinserted in pool', server_state.server)170                server_state.last_checked_time = datetime.now()171                if log_enabled(NETWORK):172                    log(NETWORK, 'checking server <%s> for availability', server_state.server)173                if server_state.server.check_availability():174                    server_state.available = True175                    return offset176                else:177                    server_state.available = False  # sets server offline178 179            if not isinstance(self.server_pool.active, bool):180                counter -= 1181            if log_enabled(NETWORK):182                log(NETWORK, 'waiting for %d seconds before retrying pool servers cycle', conf_pool_timeout)183            sleep(conf_pool_timeout)184 185        if log_enabled(ERROR):186            log(ERROR, 'no active server available in Server Pool <%s> after maximum number of tries', self)187        raise LDAPServerPoolExhaustedError('no active server available in server pool after maximum number of tries')188 189    def __len__(self):190        return len(self.server_states)191 192 193class ServerPool(object):194    def __init__(self,195                 servers=None,196                 pool_strategy=ROUND_ROBIN,197                 active=True,198                 exhaust=False,199                 single_state=True):200 201        if pool_strategy not in POOLING_STRATEGIES:202            if log_enabled(ERROR):203                log(ERROR, 'unknown pooling strategy <%s>', pool_strategy)204            raise LDAPUnknownStrategyError('unknown pooling strategy')205        if exhaust and not active:206            if log_enabled(ERROR):207                log(ERROR, 'cannot instantiate pool with exhaust and not active')208            raise LDAPServerPoolError('pools can be exhausted only when checking for active servers')209        self.servers = []210        self.pool_states = dict()211        self.active = active212        self.exhaust = exhaust213        self.single = single_state214        self._pool_state = None # used for storing the global state of the pool215        if isinstance(servers, SEQUENCE_TYPES + (Server, )):216            self.add(servers)217        elif isinstance(servers, STRING_TYPES):218            self.add(Server(servers))219        self.strategy = pool_strategy220 221        if log_enabled(BASIC):222            log(BASIC, 'instantiated ServerPool: <%r>', self)223 224    def __str__(self):225            s = 'servers: ' + linesep226            if self.servers:227                for server in self.servers:228                    s += str(server) + linesep229            else:230                s += 'None' + linesep231            s += 'Pool strategy: ' + str(self.strategy)232            s += ' - ' + 'active: ' + (str(self.active) if self.active else 'False')233            s += ' - ' + 'exhaust pool: ' + (str(self.exhaust) if self.exhaust else 'False')234            return s235 236    def __repr__(self):237        r = 'ServerPool(servers='238        if self.servers:239            r += '['240            for server in self.servers:241                r += server.__repr__() + ', '242            r = r[:-2] + ']'243        else:244            r += 'None'245        r += ', pool_strategy={0.strategy!r}'.format(self)246        r += ', active={0.active!r}'.format(self)247        r += ', exhaust={0.exhaust!r}'.format(self)248        r += ')'249 250        return r251 252    def __len__(self):253        return len(self.servers)254 255    def __getitem__(self, item):256        return self.servers[item]257 258    def __iter__(self):259        return self.servers.__iter__()260 261    def add(self, servers):262        if isinstance(servers, Server):263            if servers not in self.servers:264                self.servers.append(servers)265        elif isinstance(servers, STRING_TYPES):266            self.servers.append(Server(servers))267        elif isinstance(servers, SEQUENCE_TYPES):268            for server in servers:269                if isinstance(server, Server):270                    self.servers.append(server)271                elif isinstance(server, STRING_TYPES):272                    self.servers.append(Server(server))273                else:274                    if log_enabled(ERROR):275                        log(ERROR, 'element must be a server in Server Pool <%s>', self)276                    raise LDAPServerPoolError('server in ServerPool must be a Server')277        else:278            if log_enabled(ERROR):279                log(ERROR, 'server must be a Server of a list of Servers when adding to Server Pool <%s>', self)280            raise LDAPServerPoolError('server must be a Server or a list of Server')281 282        if self.single:283            if self._pool_state:284                self._pool_state.refresh()285        else:286            for connection in self.pool_states:287                # notifies connections using this pool to refresh288                self.pool_states[connection].refresh()289 290    def remove(self, server):291        if server in self.servers:292            self.servers.remove(server)293        else:294            if log_enabled(ERROR):295                log(ERROR, 'server %s to be removed not in Server Pool <%s>', server, self)296            raise LDAPServerPoolError('server not in server pool')297 298        if self.single:299            if self._pool_state:300                self._pool_state.refresh()301        else:302            for connection in self.pool_states:303                # notifies connections using this pool to refresh304                self.pool_states[connection].refresh()305 306    def initialize(self, connection):307        # registers pool_state in ServerPool object308        if self.single:309            if not self._pool_state:310                self._pool_state = ServerPoolState(self)311            self.pool_states[connection] = self._pool_state312        else:313            self.pool_states[connection] = ServerPoolState(self)314 315    def get_server(self, connection):316        if connection in self.pool_states:317            return self.pool_states[connection].get_server()318        else:319            if log_enabled(ERROR):320                log(ERROR, 'connection <%s> not in Server Pool State <%s>', connection, self)321            raise LDAPServerPoolError('connection not in ServerPoolState')322 323    def get_current_server(self, connection):324        if connection in self.pool_states:325            return self.pool_states[connection].get_current_server()326        else:327            if log_enabled(ERROR):328                log(ERROR, 'connection <%s> not in Server Pool State <%s>', connection, self)329            raise LDAPServerPoolError('connection not in ServerPoolState')330 
codekingpro/portable-devtools · Team Ai