codekingpro/portable-devtools
114k
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 