Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
client.py606 linesDownload Raw Back to engineio
1from base64 import b64encode2from engineio.json import JSONDecodeError3import logging4import queue5import ssl6import threading7import time8import urllib9 10try:11    import requests12except ImportError:  # pragma: no cover13    requests = None14try:15    import websocket16except ImportError:  # pragma: no cover17    websocket = None18from . import base_client19from . import exceptions20from . import packet21from . import payload22 23default_logger = logging.getLogger('engineio.client')24 25 26class Client(base_client.BaseClient):27    """An Engine.IO client.28 29    This class implements a fully compliant Engine.IO web client with support30    for websocket and long-polling transports.31 32    :param logger: To enable logging set to ``True`` or pass a logger object to33                   use. To disable logging set to ``False``. The default is34                   ``False``. Note that fatal errors are logged even when35                   ``logger`` is ``False``.36    :param json: An alternative json module to use for encoding and decoding37                 packets. Custom json modules must have ``dumps`` and ``loads``38                 functions that are compatible with the standard library39                 versions.40    :param request_timeout: A timeout in seconds for requests. The default is41                            5 seconds.42    :param http_session: an initialized ``requests.Session`` object to be used43                         when sending requests to the server. Use it if you44                         need to add special client options such as proxy45                         servers, SSL certificates, custom CA bundle, etc.46    :param ssl_verify: ``True`` to verify SSL certificates, or ``False`` to47                       skip SSL certificate verification, allowing48                       connections to servers with self signed certificates.49                       The default is ``True``.50    :param handle_sigint: Set to ``True`` to automatically handle disconnection51                          when the process is interrupted, or to ``False`` to52                          leave interrupt handling to the calling application.53                          Interrupt handling can only be enabled when the54                          client instance is created in the main thread.55    :param websocket_extra_options: Dictionary containing additional keyword56                                    arguments passed to57                                    ``websocket.create_connection()``.58    """59    def connect(self, url, headers=None, transports=None,60                engineio_path='engine.io'):61        """Connect to an Engine.IO server.62 63        :param url: The URL of the Engine.IO server. It can include custom64                    query string parameters if required by the server.65        :param headers: A dictionary with custom headers to send with the66                        connection request.67        :param transports: The list of allowed transports. Valid transports68                           are ``'polling'`` and ``'websocket'``. If not69                           given, the polling transport is connected first,70                           then an upgrade to websocket is attempted.71        :param engineio_path: The endpoint where the Engine.IO server is72                              installed. The default value is appropriate for73                              most cases.74 75        Example usage::76 77            eio = engineio.Client()78            eio.connect('http://localhost:5000')79        """80        if self.state != 'disconnected':81            raise ValueError('Client is not in a disconnected state')82        valid_transports = ['polling', 'websocket']83        if transports is not None:84            if isinstance(transports, str):85                transports = [transports]86            transports = [transport for transport in transports87                          if transport in valid_transports]88            if not transports:89                raise ValueError('No valid transports provided')90        self.transports = transports or valid_transports91        self.queue = self.create_queue()92        return getattr(self, '_connect_' + self.transports[0])(93            url, headers or {}, engineio_path)94 95    def wait(self):96        """Wait until the connection with the server ends.97 98        Client applications can use this function to block the main thread99        during the life of the connection.100        """101        if self.read_loop_task:102            self.read_loop_task.join()103 104    def send(self, data):105        """Send a message to the server.106 107        :param data: The data to send to the server. Data can be of type108                     ``str``, ``bytes``, ``list`` or ``dict``. If a ``list``109                     or ``dict``, the data will be serialized as JSON.110        """111        self._send_packet(packet.Packet(packet.MESSAGE, data=data))112 113    def disconnect(self, abort=False):114        """Disconnect from the server.115 116        :param abort: If set to ``True``, do not wait for background tasks117                      associated with the connection to end.118        """119        if self.state == 'connected':120            self._send_packet(packet.Packet(packet.CLOSE))121            self.queue.put(None)122            self.state = 'disconnecting'123            self._trigger_event('disconnect', run_async=False)124            if self.current_transport == 'websocket':125                self.ws.close()126            if not abort:127                self.read_loop_task.join()128            self.state = 'disconnected'129            try:130                base_client.connected_clients.remove(self)131            except ValueError:  # pragma: no cover132                pass133        self._reset()134 135    def start_background_task(self, target, *args, **kwargs):136        """Start a background task.137 138        This is a utility function that applications can use to start a139        background task.140 141        :param target: the target function to execute.142        :param args: arguments to pass to the function.143        :param kwargs: keyword arguments to pass to the function.144 145        This function returns an object that represents the background task,146        on which the ``join()`` method can be invoked to wait for the task to147        complete.148        """149        th = threading.Thread(target=target, args=args, kwargs=kwargs,150                              daemon=True)151        th.start()152        return th153 154    def sleep(self, seconds=0):155        """Sleep for the requested amount of time."""156        return time.sleep(seconds)157 158    def create_queue(self, *args, **kwargs):159        """Create a queue object."""160        q = queue.Queue(*args, **kwargs)161        q.Empty = queue.Empty162        return q163 164    def create_event(self, *args, **kwargs):165        """Create an event object."""166        return threading.Event(*args, **kwargs)167 168    def _connect_polling(self, url, headers, engineio_path):169        """Establish a long-polling connection to the Engine.IO server."""170        if requests is None:  # pragma: no cover171            # not installed172            self.logger.error('requests package is not installed -- cannot '173                              'send HTTP requests!')174            return175        self.base_url = self._get_engineio_url(url, engineio_path, 'polling')176        self.logger.info('Attempting polling connection to ' + self.base_url)177        r = self._send_request(178            'GET', self.base_url + self._get_url_timestamp(), headers=headers,179            timeout=self.request_timeout)180        if r is None or isinstance(r, str):181            self._reset()182            raise exceptions.ConnectionError(183                r or 'Connection refused by the server')184        if r.status_code < 200 or r.status_code >= 300:185            self._reset()186            try:187                arg = r.json()188            except JSONDecodeError:189                arg = None190            raise exceptions.ConnectionError(191                'Unexpected status code {} in server response'.format(192                    r.status_code), arg)193        try:194            p = payload.Payload(encoded_payload=r.content.decode('utf-8'))195        except ValueError:196            raise exceptions.ConnectionError(197                'Unexpected response from server') from None198        open_packet = p.packets[0]199        if open_packet.packet_type != packet.OPEN:200            raise exceptions.ConnectionError(201                'OPEN packet not returned by server')202        self.logger.info(203            'Polling connection accepted with ' + str(open_packet.data))204        self.sid = open_packet.data['sid']205        self.upgrades = open_packet.data['upgrades']206        self.ping_interval = int(open_packet.data['pingInterval']) / 1000.0207        self.ping_timeout = int(open_packet.data['pingTimeout']) / 1000.0208        self.current_transport = 'polling'209        self.base_url += '&sid=' + self.sid210 211        self.state = 'connected'212        base_client.connected_clients.append(self)213        self._trigger_event('connect', run_async=False)214 215        for pkt in p.packets[1:]:216            self._receive_packet(pkt)217 218        if 'websocket' in self.upgrades and 'websocket' in self.transports:219            # attempt to upgrade to websocket220            if self._connect_websocket(url, headers, engineio_path):221                # upgrade to websocket succeeded, we're done here222                return223 224        # start background tasks associated with this client225        self.write_loop_task = self.start_background_task(self._write_loop)226        self.read_loop_task = self.start_background_task(227            self._read_loop_polling)228 229    def _connect_websocket(self, url, headers, engineio_path):230        """Establish or upgrade to a WebSocket connection with the server."""231        if websocket is None:  # pragma: no cover232            # not installed233            self.logger.error('websocket-client package not installed, only '234                              'polling transport is available')235            return False236        websocket_url = self._get_engineio_url(url, engineio_path, 'websocket')237        if self.sid:238            self.logger.info(239                'Attempting WebSocket upgrade to ' + websocket_url)240            upgrade = True241            websocket_url += '&sid=' + self.sid242        else:243            upgrade = False244            self.base_url = websocket_url245            self.logger.info(246                'Attempting WebSocket connection to ' + websocket_url)247 248        # get cookies and other settings from the long-polling connection249        # so that they are preserved when connecting to the WebSocket route250        cookies = None251        extra_options = {}252        if self.http:253            # cookies254            cookies = '; '.join(["{}={}".format(cookie.name, cookie.value)255                                 for cookie in self.http.cookies])256            for header, value in headers.items():257                if header.lower() == 'cookie':258                    if cookies:259                        cookies += '; '260                    cookies += value261                    del headers[header]262                    break263 264            # auth265            if 'Authorization' not in headers and self.http.auth is not None:266                if not isinstance(self.http.auth, tuple):  # pragma: no cover267                    raise ValueError('Only basic authentication is supported')268                basic_auth = '{}:{}'.format(269                    self.http.auth[0], self.http.auth[1]).encode('utf-8')270                basic_auth = b64encode(basic_auth).decode('utf-8')271                headers['Authorization'] = 'Basic ' + basic_auth272 273            # cert274            # this can be given as ('certfile', 'keyfile') or just 'certfile'275            if isinstance(self.http.cert, tuple):276                extra_options['sslopt'] = {277                    'certfile': self.http.cert[0],278                    'keyfile': self.http.cert[1]}279            elif self.http.cert:280                extra_options['sslopt'] = {'certfile': self.http.cert}281 282            # proxies283            if self.http.proxies:284                proxy_url = None285                if websocket_url.startswith('ws://'):286                    proxy_url = self.http.proxies.get(287                        'ws', self.http.proxies.get('http'))288                else:  # wss://289                    proxy_url = self.http.proxies.get(290                        'wss', self.http.proxies.get('https'))291                if proxy_url:292                    parsed_url = urllib.parse.urlparse(293                        proxy_url if '://' in proxy_url294                        else 'scheme://' + proxy_url)295                    extra_options['http_proxy_host'] = parsed_url.hostname296                    extra_options['http_proxy_port'] = parsed_url.port297                    extra_options['http_proxy_auth'] = (298                        (parsed_url.username, parsed_url.password)299                        if parsed_url.username or parsed_url.password300                        else None)301 302            # verify303            if isinstance(self.http.verify, str):304                if 'sslopt' in extra_options:305                    extra_options['sslopt']['ca_certs'] = self.http.verify306                else:307                    extra_options['sslopt'] = {'ca_certs': self.http.verify}308            elif not self.http.verify:309                self.ssl_verify = False310 311        if not self.ssl_verify:312            if 'sslopt' in extra_options:313                extra_options['sslopt'].update({"cert_reqs": ssl.CERT_NONE})314            else:315                extra_options['sslopt'] = {"cert_reqs": ssl.CERT_NONE}316 317        # combine internally generated options with the ones supplied by the318        # caller. The caller's options take precedence.319        headers.update(self.websocket_extra_options.pop('header', {}))320        extra_options['header'] = headers321        extra_options['cookie'] = cookies322        extra_options['enable_multithread'] = True323        extra_options['timeout'] = self.request_timeout324        extra_options.update(self.websocket_extra_options)325        try:326            ws = websocket.create_connection(327                websocket_url + self._get_url_timestamp(), **extra_options)328        except (ConnectionError, IOError, websocket.WebSocketException):329            if upgrade:330                self.logger.warning(331                    'WebSocket upgrade failed: connection error')332                return False333            else:334                raise exceptions.ConnectionError('Connection error')335        if upgrade:336            p = packet.Packet(packet.PING, data='probe').encode()337            try:338                ws.send(p)339            except Exception as e:  # pragma: no cover340                self.logger.warning(341                    'WebSocket upgrade failed: unexpected send exception: %s',342                    str(e))343                return False344            try:345                p = ws.recv()346            except Exception as e:  # pragma: no cover347                self.logger.warning(348                    'WebSocket upgrade failed: unexpected recv exception: %s',349                    str(e))350                return False351            pkt = packet.Packet(encoded_packet=p)352            if pkt.packet_type != packet.PONG or pkt.data != 'probe':353                self.logger.warning(354                    'WebSocket upgrade failed: no PONG packet')355                return False356            p = packet.Packet(packet.UPGRADE).encode()357            try:358                ws.send(p)359            except Exception as e:  # pragma: no cover360                self.logger.warning(361                    'WebSocket upgrade failed: unexpected send exception: %s',362                    str(e))363                return False364            self.current_transport = 'websocket'365            self.logger.info('WebSocket upgrade was successful')366        else:367            try:368                p = ws.recv()369            except Exception as e:  # pragma: no cover370                raise exceptions.ConnectionError(371                    'Unexpected recv exception: ' + str(e))372            open_packet = packet.Packet(encoded_packet=p)373            if open_packet.packet_type != packet.OPEN:374                raise exceptions.ConnectionError('no OPEN packet')375            self.logger.info(376                'WebSocket connection accepted with ' + str(open_packet.data))377            self.sid = open_packet.data['sid']378            self.upgrades = open_packet.data['upgrades']379            self.ping_interval = int(open_packet.data['pingInterval']) / 1000.0380            self.ping_timeout = int(open_packet.data['pingTimeout']) / 1000.0381            self.current_transport = 'websocket'382 383            self.state = 'connected'384            base_client.connected_clients.append(self)385            self._trigger_event('connect', run_async=False)386        self.ws = ws387        self.ws.settimeout(self.ping_interval + self.ping_timeout)388 389        # start background tasks associated with this client390        self.write_loop_task = self.start_background_task(self._write_loop)391        self.read_loop_task = self.start_background_task(392            self._read_loop_websocket)393        return True394 395    def _receive_packet(self, pkt):396        """Handle incoming packets from the server."""397        packet_name = packet.packet_names[pkt.packet_type] \398            if pkt.packet_type < len(packet.packet_names) else 'UNKNOWN'399        self.logger.info(400            'Received packet %s data %s', packet_name,401            pkt.data if not isinstance(pkt.data, bytes) else '<binary>')402        if pkt.packet_type == packet.MESSAGE:403            self._trigger_event('message', pkt.data, run_async=True)404        elif pkt.packet_type == packet.PING:405            self._send_packet(packet.Packet(packet.PONG, pkt.data))406        elif pkt.packet_type == packet.CLOSE:407            self.disconnect(abort=True)408        elif pkt.packet_type == packet.NOOP:409            pass410        else:411            self.logger.error('Received unexpected packet of type %s',412                              pkt.packet_type)413 414    def _send_packet(self, pkt):415        """Queue a packet to be sent to the server."""416        if self.state != 'connected':417            return418        self.queue.put(pkt)419        self.logger.info(420            'Sending packet %s data %s',421            packet.packet_names[pkt.packet_type],422            pkt.data if not isinstance(pkt.data, bytes) else '<binary>')423 424    def _send_request(425            self, method, url, headers=None, body=None,426            timeout=None):  # pragma: no cover427        if self.http is None:428            self.http = requests.Session()429        if not self.ssl_verify:430            self.http.verify = False431        try:432            return self.http.request(method, url, headers=headers, data=body,433                                     timeout=timeout)434        except requests.exceptions.RequestException as exc:435            self.logger.info('HTTP %s request to %s failed with error %s.',436                             method, url, exc)437            return str(exc)438 439    def _trigger_event(self, event, *args, **kwargs):440        """Invoke an event handler."""441        run_async = kwargs.pop('run_async', False)442        if event in self.handlers:443            if run_async:444                return self.start_background_task(self.handlers[event], *args)445            else:446                try:447                    return self.handlers[event](*args)448                except:449                    self.logger.exception(event + ' handler error')450 451    def _read_loop_polling(self):452        """Read packets by polling the Engine.IO server."""453        while self.state == 'connected' and self.write_loop_task:454            self.logger.info(455                'Sending polling GET request to ' + self.base_url)456            r = self._send_request(457                'GET', self.base_url + self._get_url_timestamp(),458                timeout=max(self.ping_interval, self.ping_timeout) + 5)459            if r is None or isinstance(r, str):460                self.logger.warning(461                    r or 'Connection refused by the server, aborting')462                self.queue.put(None)463                break464            if r.status_code < 200 or r.status_code >= 300:465                self.logger.warning('Unexpected status code %s in server '466                                    'response, aborting', r.status_code)467                self.queue.put(None)468                break469            try:470                p = payload.Payload(encoded_payload=r.content.decode('utf-8'))471            except ValueError:472                self.logger.warning(473                    'Unexpected packet from server, aborting')474                self.queue.put(None)475                break476            for pkt in p.packets:477                self._receive_packet(pkt)478 479        if self.write_loop_task:  # pragma: no branch480            self.logger.info('Waiting for write loop task to end')481            self.write_loop_task.join()482        if self.state == 'connected':483            self._trigger_event('disconnect', run_async=False)484            try:485                base_client.connected_clients.remove(self)486            except ValueError:  # pragma: no cover487                pass488            self._reset()489        self.logger.info('Exiting read loop task')490 491    def _read_loop_websocket(self):492        """Read packets from the Engine.IO WebSocket connection."""493        while self.state == 'connected':494            p = None495            try:496                p = self.ws.recv()497                if len(p) == 0:  # pragma: no cover498                    # websocket client can return an empty string after close499                    raise websocket.WebSocketConnectionClosedException()500            except websocket.WebSocketTimeoutException:501                self.logger.warning(502                    'Server has stopped communicating, aborting')503                self.queue.put(None)504                break505            except websocket.WebSocketConnectionClosedException:506                self.logger.warning(507                    'WebSocket connection was closed, aborting')508                self.queue.put(None)509                break510            except Exception as e:  # pragma: no cover511                if type(e) is OSError and e.errno == 9:512                    self.logger.info(513                        'WebSocket connection is closing, aborting',514                        str(e))515                else:516                    self.logger.info(517                        'Unexpected error receiving packet: "%s", aborting',518                        str(e))519                self.queue.put(None)520                break521            try:522                pkt = packet.Packet(encoded_packet=p)523            except Exception as e:  # pragma: no cover524                self.logger.info(525                    'Unexpected error decoding packet: "%s", aborting', str(e))526                self.queue.put(None)527                break528            self._receive_packet(pkt)529 530        if self.write_loop_task:  # pragma: no branch531            self.logger.info('Waiting for write loop task to end')532            self.write_loop_task.join()533        if self.state == 'connected':534            self._trigger_event('disconnect', run_async=False)535            try:536                base_client.connected_clients.remove(self)537            except ValueError:  # pragma: no cover538                pass539            self._reset()540        self.logger.info('Exiting read loop task')541 542    def _write_loop(self):543        """This background task sends packages to the server as they are544        pushed to the send queue.545        """546        while self.state == 'connected':547            # to simplify the timeout handling, use the maximum of the548            # ping interval and ping timeout as timeout, with an extra 5549            # seconds grace period550            timeout = max(self.ping_interval, self.ping_timeout) + 5551            packets = None552            try:553                packets = [self.queue.get(timeout=timeout)]554            except self.queue.Empty:555                self.logger.error('packet queue is empty, aborting')556                break557            if packets == [None]:558                self.queue.task_done()559                packets = []560            else:561                while True:562                    try:563                        packets.append(self.queue.get(block=False))564                    except self.queue.Empty:565                        break566                    if packets[-1] is None:567                        packets = packets[:-1]568                        self.queue.task_done()569                        break570            if not packets:571                # empty packet list returned -> connection closed572                break573            if self.current_transport == 'polling':574                p = payload.Payload(packets=packets)575                r = self._send_request(576                    'POST', self.base_url, body=p.encode(),577                    headers={'Content-Type': 'text/plain'},578                    timeout=self.request_timeout)579                for pkt in packets:580                    self.queue.task_done()581                if r is None or isinstance(r, str):582                    self.logger.warning(583                        r or 'Connection refused by the server, aborting')584                    break585                if r.status_code < 200 or r.status_code >= 300:586                    self.logger.warning('Unexpected status code %s in server '587                                        'response, aborting', r.status_code)588                    self.write_loop_task = None589                    break590            else:591                # websocket592                try:593                    for pkt in packets:594                        encoded_packet = pkt.encode()595                        if pkt.binary:596                            self.ws.send_binary(encoded_packet)597                        else:598                            self.ws.send(encoded_packet)599                        self.queue.task_done()600                except (websocket.WebSocketConnectionClosedException,601                        BrokenPipeError, OSError):602                    self.logger.warning(603                        'WebSocket connection was closed, aborting')604                    break605        self.logger.info('Exiting write loop task')606 
codekingpro/portable-devtools · Team Ai