Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
reusable.py496 linesDownload Raw Back to strategy
1"""2"""3 4# Created on 2014.03.235#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 datetime27from os import linesep28from threading import Thread, Lock29from time import sleep30 31from .. import RESTARTABLE, get_config_parameter, AUTO_BIND_DEFAULT, AUTO_BIND_NONE, AUTO_BIND_NO_TLS, AUTO_BIND_TLS_AFTER_BIND, AUTO_BIND_TLS_BEFORE_BIND32from .base import BaseStrategy33from ..core.usage import ConnectionUsage34from ..core.exceptions import LDAPConnectionPoolNameIsMandatoryError, LDAPConnectionPoolNotStartedError, LDAPOperationResult, LDAPExceptionError, LDAPResponseTimeoutError35from ..utils.log import log, log_enabled, ERROR, BASIC36from ..protocol.rfc4511 import LDAP_MAX_INT37 38TERMINATE_REUSABLE = 'TERMINATE_REUSABLE_CONNECTION'39 40BOGUS_BIND = -141BOGUS_UNBIND = -242BOGUS_EXTENDED = -343BOGUS_ABANDON = -444 45try:46    from queue import Queue, Empty47except ImportError:  # Python 248    # noinspection PyUnresolvedReferences49    from Queue import Queue, Empty50 51 52# noinspection PyProtectedMember53class ReusableStrategy(BaseStrategy):54    """55    A pool of reusable SyncWaitRestartable connections with lazy behaviour and limited lifetime.56    The connection using this strategy presents itself as a normal connection, but internally the strategy has a pool of57    connections that can be used as needed. Each connection lives in its own thread and has a busy/available status.58    The strategy performs the requested operation on the first available connection.59    The pool of connections is instantiated at strategy initialization.60    Strategy has two customizable properties, the total number of connections in the pool and the lifetime of each connection.61    When lifetime is expired the connection is closed and will be open again when needed.62    """63    pools = dict()64 65    def receiving(self):66        raise NotImplementedError67 68    def _start_listen(self):69        raise NotImplementedError70 71    def _get_response(self, message_id, timeout):72        raise NotImplementedError73 74    def get_stream(self):75        raise NotImplementedError76 77    def set_stream(self, value):78        raise NotImplementedError79 80    # noinspection PyProtectedMember81    class ConnectionPool(object):82        """83        Container for the Connection Threads84        """85        def __new__(cls, connection):86            if connection.pool_name in ReusableStrategy.pools:  # returns existing connection pool87                pool = ReusableStrategy.pools[connection.pool_name]88                if not pool.started:  # if pool is not started remove it from the pools singleton and create a new onw89                    del ReusableStrategy.pools[connection.pool_name]90                    return object.__new__(cls)91                if connection.pool_keepalive and pool.keepalive != connection.pool_keepalive:  # change lifetime92                    pool.keepalive = connection.pool_keepalive93                if connection.pool_lifetime and pool.lifetime != connection.pool_lifetime:  # change keepalive94                    pool.lifetime = connection.pool_lifetime95                if connection.pool_size and pool.pool_size != connection.pool_size:  # if pool size has changed terminate and recreate the connections96                    pool.terminate_pool()97                    pool.pool_size = connection.pool_size98                return pool99            else:100                return object.__new__(cls)101 102        def __init__(self, connection):103            if not hasattr(self, 'workers'):104                self.name = connection.pool_name105                self.master_connection = connection106                self.workers = []107                self.pool_size = connection.pool_size or get_config_parameter('REUSABLE_THREADED_POOL_SIZE')108                self.lifetime = connection.pool_lifetime or get_config_parameter('REUSABLE_THREADED_LIFETIME')109                self.keepalive = connection.pool_keepalive110                self.request_queue = Queue()111                self.open_pool = False112                self.bind_pool = False113                self.tls_pool = False114                self._incoming = dict()115                self.counter = 0116                self.terminated_usage = ConnectionUsage() if connection._usage else None117                self.terminated = False118                self.pool_lock = Lock()119                ReusableStrategy.pools[self.name] = self120                self.started = False121                if log_enabled(BASIC):122                    log(BASIC, 'instantiated ConnectionPool: <%r>', self)123 124        def __str__(self):125            s = 'POOL: ' + str(self.name) + ' - status: ' + ('started' if self.started else 'terminated')126            s += ' - responses in queue: ' + str(len(self._incoming))127            s += ' - pool size: ' + str(self.pool_size)128            s += ' - lifetime: ' + str(self.lifetime)129            s += ' - keepalive: ' + str(self.keepalive)130            s += ' - open: ' + str(self.open_pool)131            s += ' - bind: ' + str(self.bind_pool)132            s += ' - tls: ' + str(self.tls_pool) + linesep133            s += 'MASTER CONN: ' + str(self.master_connection) + linesep134            s += 'WORKERS:'135            if self.workers:136                for i, worker in enumerate(self.workers):137                    s += linesep + str(i).rjust(5) + ': ' + str(worker)138            else:139                s += linesep + '    no active workers in pool'140 141            return s142 143        def __repr__(self):144            return self.__str__()145 146        def get_info_from_server(self):147            for worker in self.workers:148                with worker.worker_lock:149                    if not worker.connection.server.schema or not worker.connection.server.info:150                        worker.get_info_from_server = True151                    else:152                        worker.get_info_from_server = False153 154        def rebind_pool(self):155            for worker in self.workers:156                with worker.worker_lock:157                    worker.connection.rebind(self.master_connection.user,158                                             self.master_connection.password,159                                             self.master_connection.authentication,160                                             self.master_connection.sasl_mechanism,161                                             self.master_connection.sasl_credentials)162 163        def start_pool(self):164            if not self.started:165                self.create_pool()166                for worker in self.workers:167                    with worker.worker_lock:168                        worker.thread.start()169                self.started = True170                self.terminated = False171                if log_enabled(BASIC):172                    log(BASIC, 'worker started for pool <%s>', self)173                return True174            return False175 176        def create_pool(self):177            if log_enabled(BASIC):178                log(BASIC, 'created pool <%s>', self)179            self.workers = [ReusableStrategy.PooledConnectionWorker(self.master_connection, self.request_queue) for _ in range(self.pool_size)]180 181        def terminate_pool(self):182            if not self.terminated:183                if log_enabled(BASIC):184                    log(BASIC, 'terminating pool <%s>', self)185                self.started = False186                self.request_queue.join()  # waits for all queue pending operations187                for _ in range(len([worker for worker in self.workers if worker.thread.is_alive()])):  # put a TERMINATE signal on the queue for each active thread188                    self.request_queue.put((TERMINATE_REUSABLE, None, None, None))189                self.request_queue.join()  # waits for all queue terminate operations190                self.terminated = True191                if log_enabled(BASIC):192                    log(BASIC, 'pool terminated for <%s>', self)193 194    class PooledConnectionThread(Thread):195        """196        The thread that holds the Reusable connection and receive operation request via the queue197        Result are sent back in the pool._incoming list when ready198        """199        def __init__(self, worker, master_connection):200            Thread.__init__(self)201            self.daemon = True202            self.worker = worker203            self.master_connection = master_connection204            if log_enabled(BASIC):205                log(BASIC, 'instantiated PooledConnectionThread: <%r>', self)206 207        # noinspection PyProtectedMember208        def run(self):209            self.worker.running = True210            terminate = False211            pool = self.master_connection.strategy.pool212            while not terminate:213                try:214                    counter, message_type, request, controls = pool.request_queue.get(block=True, timeout=self.master_connection.strategy.pool.keepalive)215                except Empty:  # issue an Abandon(0) operation to keep the connection live - Abandon(0) is a harmless operation216                    if not self.worker.connection.closed:217                        self.worker.connection.abandon(0)218                    continue219 220                with self.worker.worker_lock:221                    self.worker.busy = True222                    if counter == TERMINATE_REUSABLE:223                        terminate = True224                        if self.worker.connection.bound:225                            try:226                                self.worker.connection.unbind()227                                if log_enabled(BASIC):228                                    log(BASIC, 'thread terminated')229                            except LDAPExceptionError:230                                pass231                    else:232                        if (datetime.now() - self.worker.creation_time).seconds >= self.master_connection.strategy.pool.lifetime:  # destroy and create a new connection233                            try:234                                self.worker.connection.unbind()235                            except LDAPExceptionError:236                                pass237                            self.worker.new_connection()238                            if log_enabled(BASIC):239                                log(BASIC, 'thread respawn')240                        if message_type not in ['bindRequest', 'unbindRequest']:241                            try:242                                if pool.open_pool and self.worker.connection.closed:243                                    self.worker.connection.open(read_server_info=False)244                                    if pool.tls_pool and not self.worker.connection.tls_started:245                                        self.worker.connection.start_tls(read_server_info=False)246                                    if pool.bind_pool and not self.worker.connection.bound:247                                        self.worker.connection.bind(read_server_info=False)248                                elif pool.open_pool and not self.worker.connection.closed:  # connection already open, issues a start_tls249                                    if pool.tls_pool and not self.worker.connection.tls_started:250                                        self.worker.connection.start_tls(read_server_info=False)251                                if self.worker.get_info_from_server and counter:252                                    self.worker.connection.refresh_server_info()253                                    self.worker.get_info_from_server = False254                                response = None255                                result = None256                                if message_type == 'searchRequest':257                                    response = self.worker.connection.post_send_search(self.worker.connection.send(message_type, request, controls))258                                else:259                                    response = self.worker.connection.post_send_single_response(self.worker.connection.send(message_type, request, controls))260                                result = self.worker.connection.result261                                with pool.pool_lock:262                                    pool._incoming[counter] = (response, result, BaseStrategy.decode_request(message_type, request, controls))263                            except LDAPOperationResult as e:  # raise_exceptions has raised an exception. It must be redirected to the original connection thread264                                with pool.pool_lock:265                                    pool._incoming[counter] = (e, None, None)266                                    # pool._incoming[counter] = (type(e)(str(e)), None, None)267                            # except LDAPOperationResult as e:  # raise_exceptions has raised an exception. It must be redirected to the original connection thread268                            #     exc = e269                            # with pool.pool_lock:270                            #     if exc:271                            #         pool._incoming[counter] = (exc, None, None)272                            #     else:273                            #         pool._incoming[counter] = (response, result, BaseStrategy.decode_request(message_type, request, controls))274 275                    self.worker.busy = False276                    pool.request_queue.task_done()277                    self.worker.task_counter += 1278            if log_enabled(BASIC):279                log(BASIC, 'thread terminated')280            if self.master_connection.usage:281                pool.terminated_usage += self.worker.connection.usage282            self.worker.running = False283 284    class PooledConnectionWorker(object):285        """286        Container for the restartable connection. it includes a thread and a lock to execute the connection in the pool287        """288        def __init__(self, connection, request_queue):289            self.master_connection = connection290            self.request_queue = request_queue291            self.running = False292            self.busy = False293            self.get_info_from_server = False294            self.connection = None295            self.creation_time = None296            self.task_counter = 0297            self.new_connection()298            self.thread = ReusableStrategy.PooledConnectionThread(self, self.master_connection)299            self.worker_lock = Lock()300            if log_enabled(BASIC):301                log(BASIC, 'instantiated PooledConnectionWorker: <%s>', self)302 303        def __str__(self):304            s = 'CONN: ' + str(self.connection) + linesep + '       THREAD: '305            s += 'running' if self.running else 'halted'306            s += ' - ' + ('busy' if self.busy else 'available')307            s += ' - ' + ('created at: ' + self.creation_time.isoformat())308            s += ' - time to live: ' + str(self.master_connection.strategy.pool.lifetime - (datetime.now() - self.creation_time).seconds)309            s += ' - requests served: ' + str(self.task_counter)310 311            return s312 313        def new_connection(self):314            from ..core.connection import Connection315            # noinspection PyProtectedMember316            self.creation_time = datetime.now()317            self.connection = Connection(server=self.master_connection.server_pool if self.master_connection.server_pool else self.master_connection.server,318                                         user=self.master_connection.user,319                                         password=self.master_connection.password,320                                         auto_bind=AUTO_BIND_NONE,  # do not perform auto_bind because it reads again the schema321                                         version=self.master_connection.version,322                                         authentication=self.master_connection.authentication,323                                         client_strategy=RESTARTABLE,324                                         auto_referrals=self.master_connection.auto_referrals,325                                         auto_range=self.master_connection.auto_range,326                                         sasl_mechanism=self.master_connection.sasl_mechanism,327                                         sasl_credentials=self.master_connection.sasl_credentials,328                                         check_names=self.master_connection.check_names,329                                         collect_usage=self.master_connection._usage,330                                         read_only=self.master_connection.read_only,331                                         raise_exceptions=self.master_connection.raise_exceptions,332                                         lazy=False,333                                         fast_decoder=self.master_connection.fast_decoder,334                                         receive_timeout=self.master_connection.receive_timeout,335                                         return_empty_attributes=self.master_connection.empty_attributes)336 337            # simulates auto_bind, always with read_server_info=False338            if self.master_connection.auto_bind and self.master_connection.auto_bind not in [AUTO_BIND_NONE, AUTO_BIND_DEFAULT]:339                if log_enabled(BASIC):340                    log(BASIC, 'performing automatic bind for <%s>', self.connection)341                self.connection.open(read_server_info=False)342                if self.master_connection.auto_bind == AUTO_BIND_NO_TLS:343                    self.connection.bind(read_server_info=False)344                elif self.master_connection.auto_bind == AUTO_BIND_TLS_BEFORE_BIND:345                    self.connection.start_tls(read_server_info=False)346                    self.connection.bind(read_server_info=False)347                elif self.master_connection.auto_bind == AUTO_BIND_TLS_AFTER_BIND:348                    self.connection.bind(read_server_info=False)349                    self.connection.start_tls(read_server_info=False)350 351            if self.master_connection.server_pool:352                self.connection.server_pool = self.master_connection.server_pool353                self.connection.server_pool.initialize(self.connection)354 355    # ReusableStrategy methods356    def __init__(self, ldap_connection):357        BaseStrategy.__init__(self, ldap_connection)358        self.sync = False359        self.no_real_dsa = False360        self.pooled = True361        self.can_stream = False362        if hasattr(ldap_connection, 'pool_name') and ldap_connection.pool_name:363            self.pool = ReusableStrategy.ConnectionPool(ldap_connection)364        else:365            if log_enabled(ERROR):366                log(ERROR, 'reusable connection must have a pool_name')367            raise LDAPConnectionPoolNameIsMandatoryError('reusable connection must have a pool_name')368 369    def open(self, reset_usage=True, read_server_info=True):370        # read_server_info not used371        self.pool.open_pool = True372        self.pool.start_pool()373        self.connection.closed = False374        if self.connection.usage:375            if reset_usage or not self.connection._usage.initial_connection_start_time:376                self.connection._usage.start()377 378    def terminate(self):379        self.pool.terminate_pool()380        self.pool.open_pool = False381        self.connection.bound = False382        self.connection.closed = True383        self.pool.bind_pool = False384        self.pool.tls_pool = False385 386    def _close_socket(self):387        """388        Doesn't really close the socket389        """390        self.connection.closed = True391 392        if self.connection.usage:393            self.connection._usage.closed_sockets += 1394 395    def send(self, message_type, request, controls=None):396        if self.pool.started:397            if message_type == 'bindRequest':398                self.pool.bind_pool = True399                counter = BOGUS_BIND400            elif message_type == 'unbindRequest':401                self.pool.bind_pool = False402                counter = BOGUS_UNBIND403            elif message_type == 'abandonRequest':404                counter = BOGUS_ABANDON405            elif message_type == 'extendedReq' and self.connection.starting_tls:406                self.pool.tls_pool = True407                counter = BOGUS_EXTENDED408            else:409                with self.pool.pool_lock:410                    self.pool.counter += 1411                    if self.pool.counter > LDAP_MAX_INT:412                        self.pool.counter = 1413                    counter = self.pool.counter414                self.pool.request_queue.put((counter, message_type, request, controls))415            return counter416        if log_enabled(ERROR):417            log(ERROR, 'reusable connection pool not started')418        raise LDAPConnectionPoolNotStartedError('reusable connection pool not started')419 420    def validate_bind(self, controls):421        # in case of a new connection or different credentials422        if (self.connection.user != self.pool.master_connection.user or423                self.connection.password != self.pool.master_connection.password or424                self.connection.authentication != self.pool.master_connection.authentication or425                self.connection.sasl_mechanism != self.pool.master_connection.sasl_mechanism or426                self.connection.sasl_credentials != self.pool.master_connection.sasl_credentials):427            self.pool.master_connection.user = self.connection.user428            self.pool.master_connection.password = self.connection.password429            self.pool.master_connection.authentication = self.connection.authentication430            self.pool.master_connection.sasl_mechanism = self.connection.sasl_mechanism431            self.pool.master_connection.sasl_credentials = self.connection.sasl_credentials432            self.pool.rebind_pool()433        temp_connection = self.pool.workers[0].connection434        old_lazy = temp_connection.lazy435        temp_connection.lazy = False436        if not self.connection.server.schema or not self.connection.server.info:437            result = self.pool.workers[0].connection.bind(controls=controls)438        else:439            result = self.pool.workers[0].connection.bind(controls=controls, read_server_info=False)440 441        temp_connection.unbind()442        temp_connection.lazy = old_lazy443        if result:444            self.pool.bind_pool = True  # bind pool if bind is validated445        return result446 447    def get_response(self, counter, timeout=None, get_request=False):448        sleeptime = get_config_parameter('RESPONSE_SLEEPTIME')449        request=None450        if timeout is None:451            timeout = get_config_parameter('RESPONSE_WAITING_TIMEOUT')452        if counter == BOGUS_BIND:  # send a bogus bindResponse453            response = list()454            result = {'description': 'success', 'referrals': None, 'type': 'bindResponse', 'result': 0, 'dn': '', 'message': '<bogus Bind response>', 'saslCreds': None}455        elif counter == BOGUS_UNBIND:  # bogus unbind response456            response = None457            result = None458        elif counter == BOGUS_ABANDON:  # abandon cannot be executed because of multiple connections459            response = list()460            result = {'result': 0, 'referrals': None, 'responseName': '1.3.6.1.4.1.1466.20037', 'type': 'extendedResp', 'description': 'success', 'responseValue': 'None', 'dn': '', 'message': '<bogus StartTls response>'}461        elif counter == BOGUS_EXTENDED:  # bogus startTls extended response462            response = list()463            result = {'result': 0, 'referrals': None, 'responseName': '1.3.6.1.4.1.1466.20037', 'type': 'extendedResp', 'description': 'success', 'responseValue': 'None', 'dn': '', 'message': '<bogus StartTls response>'}464            self.connection.starting_tls = False465        else:466            response = None467            result = None468            while timeout >= 0:  # waiting for completed message to appear in _incoming469                try:470                    with self.connection.strategy.pool.pool_lock:471                        response, result, request = self.connection.strategy.pool._incoming.pop(counter)472                except KeyError:473                    sleep(sleeptime)474                    timeout -= sleeptime475                    continue476                break477 478            if timeout <= 0:479                if log_enabled(ERROR):480                    log(ERROR, 'no response from worker threads in Reusable connection')481                raise LDAPResponseTimeoutError('no response from worker threads in Reusable connection')482 483        if isinstance(response, LDAPOperationResult):484            raise response  # an exception has been raised with raise_exceptions485 486        if get_request:487            return response, result, request488 489        return response, result490 491    def post_send_single_response(self, counter):492        return counter493 494    def post_send_search(self, counter):495        return counter496 
codekingpro/portable-devtools · Team Ai