Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
client.py524 linesDownload Raw Back to socketio
1import random2 3import engineio4 5from . import base_client6from . import exceptions7from . import packet8 9 10class Client(base_client.BaseClient):11    """A Socket.IO client.12 13    This class implements a fully compliant Socket.IO web client with support14    for websocket and long-polling transports.15 16    :param reconnection: ``True`` if the client should automatically attempt to17                         reconnect to the server after an interruption, or18                         ``False`` to not reconnect. The default is ``True``.19    :param reconnection_attempts: How many reconnection attempts to issue20                                  before giving up, or 0 for infinite attempts.21                                  The default is 0.22    :param reconnection_delay: How long to wait in seconds before the first23                               reconnection attempt. Each successive attempt24                               doubles this delay.25    :param reconnection_delay_max: The maximum delay between reconnection26                                   attempts.27    :param randomization_factor: Randomization amount for each delay between28                                 reconnection attempts. The default is 0.5,29                                 which means that each delay is randomly30                                 adjusted by +/- 50%.31    :param logger: To enable logging set to ``True`` or pass a logger object to32                   use. To disable logging set to ``False``. The default is33                   ``False``. Note that fatal errors are logged even when34                   ``logger`` is ``False``.35    :param serializer: The serialization method to use when transmitting36                       packets. Valid values are ``'default'``, ``'pickle'``,37                       ``'msgpack'`` and ``'cbor'``. Alternatively, a subclass38                       of the :class:`Packet` class with custom implementations39                       of the ``encode()`` and ``decode()`` methods can be40                       provided. Client and server must use compatible41                       serializers.42    :param json: An alternative json module to use for encoding and decoding43                 packets. Custom json modules must have ``dumps`` and ``loads``44                 functions that are compatible with the standard library45                 versions.46    :param handle_sigint: Set to ``True`` to automatically handle disconnection47                          when the process is interrupted, or to ``False`` to48                          leave interrupt handling to the calling application.49                          Interrupt handling can only be enabled when the50                          client instance is created in the main thread.51 52    The Engine.IO configuration supports the following settings:53 54    :param request_timeout: A timeout in seconds for requests. The default is55                            5 seconds.56    :param http_session: an initialized ``requests.Session`` object to be used57                         when sending requests to the server. Use it if you58                         need to add special client options such as proxy59                         servers, SSL certificates, etc.60    :param ssl_verify: ``True`` to verify SSL certificates, or ``False`` to61                       skip SSL certificate verification, allowing62                       connections to servers with self signed certificates.63                       The default is ``True``.64    :param engineio_logger: To enable Engine.IO logging set to ``True`` or pass65                            a logger object to use. To disable logging set to66                            ``False``. The default is ``False``. Note that67                            fatal errors are logged even when68                            ``engineio_logger`` is ``False``.69    """70    def connect(self, url, headers={}, auth=None, transports=None,71                namespaces=None, socketio_path='socket.io', wait=True,72                wait_timeout=1, retry=False):73        """Connect to a Socket.IO server.74 75        :param url: The URL of the Socket.IO server. It can include custom76                    query string parameters if required by the server. If a77                    function is provided, the client will invoke it to obtain78                    the URL each time a connection or reconnection is79                    attempted.80        :param headers: A dictionary with custom headers to send with the81                        connection request. If a function is provided, the82                        client will invoke it to obtain the headers dictionary83                        each time a connection or reconnection is attempted.84        :param auth: Authentication data passed to the server with the85                     connection request, normally a dictionary with one or86                     more string key/value pairs. If a function is provided,87                     the client will invoke it to obtain the authentication88                     data each time a connection or reconnection is attempted.89        :param transports: The list of allowed transports. Valid transports90                           are ``'polling'`` and ``'websocket'``. If not91                           given, the polling transport is connected first,92                           then an upgrade to websocket is attempted.93        :param namespaces: The namespaces to connect as a string or list of94                           strings. If not given, the namespaces that have95                           registered event handlers are connected.96        :param socketio_path: The endpoint where the Socket.IO server is97                              installed. The default value is appropriate for98                              most cases.99        :param wait: if set to ``True`` (the default) the call only returns100                     when all the namespaces are connected. If set to101                     ``False``, the call returns as soon as the Engine.IO102                     transport is connected, and the namespaces will connect103                     in the background.104        :param wait_timeout: How long the client should wait for the105                             connection. The default is 1 second. This106                             argument is only considered when ``wait`` is set107                             to ``True``.108        :param retry: Apply the reconnection logic if the initial connection109                      attempt fails. The default is ``False``.110 111        Example usage::112 113            sio = socketio.Client()114            sio.connect('http://localhost:5000')115        """116        if self.connected:117            raise exceptions.ConnectionError('Already connected')118 119        self.connection_url = url120        self.connection_headers = headers121        self.connection_auth = auth122        self.connection_transports = transports123        self.connection_namespaces = namespaces124        self.socketio_path = socketio_path125 126        if namespaces is None:127            namespaces = list(set(self.handlers.keys()).union(128                set(self.namespace_handlers.keys())))129            if len(namespaces) == 0:130                namespaces = ['/']131        elif isinstance(namespaces, str):132            namespaces = [namespaces]133        self.connection_namespaces = namespaces134        self.namespaces = {}135        if self._connect_event is None:136            self._connect_event = self.eio.create_event()137        else:138            self._connect_event.clear()139        real_url = self._get_real_value(self.connection_url)140        real_headers = self._get_real_value(self.connection_headers)141        try:142            self.eio.connect(real_url, headers=real_headers,143                             transports=transports,144                             engineio_path=socketio_path)145        except engineio.exceptions.ConnectionError as exc:146            for n in self.connection_namespaces:147                self._trigger_event(148                    'connect_error', n,149                    exc.args[1] if len(exc.args) > 1 else exc.args[0])150            if retry:  # pragma: no cover151                self._handle_reconnect()152                if self.eio.state == 'connected':153                    return154            raise exceptions.ConnectionError(exc.args[0]) from None155 156        if wait:157            while self._connect_event.wait(timeout=wait_timeout):158                self._connect_event.clear()159                if set(self.namespaces) == set(self.connection_namespaces):160                    break161            if set(self.namespaces) != set(self.connection_namespaces):162                self.disconnect()163                raise exceptions.ConnectionError(164                    'One or more namespaces failed to connect')165 166        self.connected = True167 168    def wait(self):169        """Wait until the connection with the server ends.170 171        Client applications can use this function to block the main thread172        during the life of the connection.173        """174        while True:175            self.eio.wait()176            self.sleep(1)  # give the reconnect task time to start up177            if not self._reconnect_task:178                break179            self._reconnect_task.join()180            if self.eio.state != 'connected':181                break182 183    def emit(self, event, data=None, namespace=None, callback=None):184        """Emit a custom event to the server.185 186        :param event: The event name. It can be any string. The event names187                      ``'connect'``, ``'message'`` and ``'disconnect'`` are188                      reserved and should not be used.189        :param data: The data to send to the server. Data can be of190                     type ``str``, ``bytes``, ``list`` or ``dict``. To send191                     multiple arguments, use a tuple where each element is of192                     one of the types indicated above.193        :param namespace: The Socket.IO namespace for the event. If this194                          argument is omitted the event is emitted to the195                          default namespace.196        :param callback: If given, this function will be called to acknowledge197                         the server has received the message. The arguments198                         that will be passed to the function are those provided199                         by the server.200 201        Note: this method is not thread safe. If multiple threads are emitting202        at the same time on the same client connection, messages composed of203        multiple packets may end up being sent in an incorrect sequence. Use204        standard concurrency solutions (such as a Lock object) to prevent this205        situation.206        """207        namespace = namespace or '/'208        if namespace not in self.namespaces:209            raise exceptions.BadNamespaceError(210                namespace + ' is not a connected namespace.')211        self.logger.info('Emitting event "%s" [%s]', event, namespace)212        if callback is not None:213            id = self._generate_ack_id(namespace, callback)214        else:215            id = None216        # tuples are expanded to multiple arguments, everything else is sent217        # as a single argument218        if isinstance(data, tuple):219            data = list(data)220        elif data is not None:221            data = [data]222        else:223            data = []224        self._send_packet(self.packet_class(packet.EVENT, namespace=namespace,225                                            data=[event] + data, id=id))226 227    def send(self, data, namespace=None, callback=None):228        """Send a message to the server.229 230        This function emits an event with the name ``'message'``. Use231        :func:`emit` to issue custom event names.232 233        :param data: The data to send to the server. Data can be of234                     type ``str``, ``bytes``, ``list`` or ``dict``. To send235                     multiple arguments, use a tuple where each element is of236                     one of the types indicated above.237        :param namespace: The Socket.IO namespace for the event. If this238                          argument is omitted the event is emitted to the239                          default namespace.240        :param callback: If given, this function will be called to acknowledge241                         the server has received the message. The arguments242                         that will be passed to the function are those provided243                         by the server.244        """245        self.emit('message', data=data, namespace=namespace,246                  callback=callback)247 248    def call(self, event, data=None, namespace=None, timeout=60):249        """Emit a custom event to the server and wait for the response.250 251        This method issues an emit with a callback and waits for the callback252        to be invoked before returning. If the callback isn't invoked before253        the timeout, then a ``TimeoutError`` exception is raised. If the254        Socket.IO connection drops during the wait, this method still waits255        until the specified timeout.256 257        :param event: The event name. It can be any string. The event names258                      ``'connect'``, ``'message'`` and ``'disconnect'`` are259                      reserved and should not be used.260        :param data: The data to send to the server. Data can be of261                     type ``str``, ``bytes``, ``list`` or ``dict``. To send262                     multiple arguments, use a tuple where each element is of263                     one of the types indicated above.264        :param namespace: The Socket.IO namespace for the event. If this265                          argument is omitted the event is emitted to the266                          default namespace.267        :param timeout: The waiting timeout. If the timeout is reached before268                        the server acknowledges the event, then a269                        ``TimeoutError`` exception is raised.270 271        Note: this method is not thread safe. If multiple threads are emitting272        at the same time on the same client connection, messages composed of273        multiple packets may end up being sent in an incorrect sequence. Use274        standard concurrency solutions (such as a Lock object) to prevent this275        situation.276        """277        callback_event = self.eio.create_event()278        callback_args = []279 280        def event_callback(*args):281            callback_args.append(args)282            callback_event.set()283 284        self.emit(event, data=data, namespace=namespace,285                  callback=event_callback)286        if not callback_event.wait(timeout=timeout):287            raise exceptions.TimeoutError()288        return callback_args[0] if len(callback_args[0]) > 1 \289            else callback_args[0][0] if len(callback_args[0]) == 1 \290            else None291 292    def disconnect(self):293        """Disconnect from the server."""294        # here we just request the disconnection295        # later in _handle_eio_disconnect we invoke the disconnect handler296        for n in self.namespaces:297            self._send_packet(self.packet_class(298                packet.DISCONNECT, namespace=n))299        self.eio.disconnect(abort=True)300 301    def start_background_task(self, target, *args, **kwargs):302        """Start a background task using the appropriate async model.303 304        This is a utility function that applications can use to start a305        background task using the method that is compatible with the306        selected async mode.307 308        :param target: the target function to execute.309        :param args: arguments to pass to the function.310        :param kwargs: keyword arguments to pass to the function.311 312        This function returns an object that represents the background task,313        on which the ``join()`` methond can be invoked to wait for the task to314        complete.315        """316        return self.eio.start_background_task(target, *args, **kwargs)317 318    def sleep(self, seconds=0):319        """Sleep for the requested amount of time using the appropriate async320        model.321 322        This is a utility function that applications can use to put a task to323        sleep without having to worry about using the correct call for the324        selected async mode.325        """326        return self.eio.sleep(seconds)327 328    def _get_real_value(self, value):329        """Return the actual value, for parameters that can also be given as330        callables."""331        if not callable(value):332            return value333        return value()334 335    def _send_packet(self, pkt):336        """Send a Socket.IO packet to the server."""337        encoded_packet = pkt.encode()338        if isinstance(encoded_packet, list):339            for ep in encoded_packet:340                self.eio.send(ep)341        else:342            self.eio.send(encoded_packet)343 344    def _handle_connect(self, namespace, data):345        namespace = namespace or '/'346        if namespace not in self.namespaces:347            self.logger.info('Namespace {} is connected'.format(namespace))348            self.namespaces[namespace] = (data or {}).get('sid', self.sid)349            self._trigger_event('connect', namespace=namespace)350            self._connect_event.set()351 352    def _handle_disconnect(self, namespace):353        if not self.connected:354            return355        namespace = namespace or '/'356        self._trigger_event('disconnect', namespace=namespace)357        self._trigger_event('__disconnect_final', namespace=namespace)358        if namespace in self.namespaces:359            del self.namespaces[namespace]360        if not self.namespaces:361            self.connected = False362            self.eio.disconnect(abort=True)363 364    def _handle_event(self, namespace, id, data):365        namespace = namespace or '/'366        self.logger.info('Received event "%s" [%s]', data[0], namespace)367        r = self._trigger_event(data[0], namespace, *data[1:])368        if id is not None:369            # send ACK packet with the response returned by the handler370            # tuples are expanded as multiple arguments371            if r is None:372                data = []373            elif isinstance(r, tuple):374                data = list(r)375            else:376                data = [r]377            self._send_packet(self.packet_class(378                packet.ACK, namespace=namespace, id=id, data=data))379 380    def _handle_ack(self, namespace, id, data):381        namespace = namespace or '/'382        self.logger.info('Received ack [%s]', namespace)383        callback = None384        try:385            callback = self.callbacks[namespace][id]386        except KeyError:387            # if we get an unknown callback we just ignore it388            self.logger.warning('Unknown callback received, ignoring.')389        else:390            del self.callbacks[namespace][id]391        if callback is not None:392            callback(*data)393 394    def _handle_error(self, namespace, data):395        namespace = namespace or '/'396        self.logger.info('Connection to namespace {} was rejected'.format(397            namespace))398        if data is None:399            data = tuple()400        elif not isinstance(data, (tuple, list)):401            data = (data,)402        self._trigger_event('connect_error', namespace, *data)403        self._connect_event.set()404        if namespace in self.namespaces:405            del self.namespaces[namespace]406        if namespace == '/':407            self.namespaces = {}408            self.connected = False409 410    def _trigger_event(self, event, namespace, *args):411        """Invoke an application event handler."""412        # first see if we have an explicit handler for the event413        handler, args = self._get_event_handler(event, namespace, args)414        if handler:415            return handler(*args)416 417        # or else, forward the event to a namespace handler if one exists418        handler, args = self._get_namespace_handler(namespace, args)419        if handler:420            return handler.trigger_event(event, *args)421 422    def _handle_reconnect(self):423        if self._reconnect_abort is None:  # pragma: no cover424            self._reconnect_abort = self.eio.create_event()425        self._reconnect_abort.clear()426        base_client.reconnecting_clients.append(self)427        attempt_count = 0428        current_delay = self.reconnection_delay429        while True:430            delay = current_delay431            current_delay *= 2432            if delay > self.reconnection_delay_max:433                delay = self.reconnection_delay_max434            delay += self.randomization_factor * (2 * random.random() - 1)435            self.logger.info(436                'Connection failed, new attempt in {:.02f} seconds'.format(437                    delay))438            if self._reconnect_abort.wait(delay):439                self.logger.info('Reconnect task aborted')440                for n in self.connection_namespaces:441                    self._trigger_event('__disconnect_final', namespace=n)442                break443            attempt_count += 1444            try:445                self.connect(self.connection_url,446                             headers=self.connection_headers,447                             auth=self.connection_auth,448                             transports=self.connection_transports,449                             namespaces=self.connection_namespaces,450                             socketio_path=self.socketio_path,451                             retry=False)452            except (exceptions.ConnectionError, ValueError):453                pass454            else:455                self.logger.info('Reconnection successful')456                self._reconnect_task = None457                break458            if self.reconnection_attempts and \459                    attempt_count >= self.reconnection_attempts:460                self.logger.info(461                    'Maximum reconnection attempts reached, giving up')462                for n in self.connection_namespaces:463                    self._trigger_event('__disconnect_final', namespace=n)464                break465        base_client.reconnecting_clients.remove(self)466 467    def _handle_eio_connect(self):468        """Handle the Engine.IO connection event."""469        self.logger.info('Engine.IO connection established')470        self.sid = self.eio.sid471        real_auth = self._get_real_value(self.connection_auth) or {}472        for n in self.connection_namespaces:473            self._send_packet(self.packet_class(474                packet.CONNECT, data=real_auth, namespace=n))475 476    def _handle_eio_message(self, data):477        """Dispatch Engine.IO messages."""478        if self._binary_packet:479            pkt = self._binary_packet480            if pkt.add_attachment(data):481                self._binary_packet = None482                if pkt.packet_type == packet.BINARY_EVENT:483                    self._handle_event(pkt.namespace, pkt.id, pkt.data)484                else:485                    self._handle_ack(pkt.namespace, pkt.id, pkt.data)486        else:487            pkt = self.packet_class(encoded_packet=data)488            if pkt.packet_type == packet.CONNECT:489                self._handle_connect(pkt.namespace, pkt.data)490            elif pkt.packet_type == packet.DISCONNECT:491                self._handle_disconnect(pkt.namespace)492            elif pkt.packet_type == packet.EVENT:493                self._handle_event(pkt.namespace, pkt.id, pkt.data)494            elif pkt.packet_type == packet.ACK:495                self._handle_ack(pkt.namespace, pkt.id, pkt.data)496            elif pkt.packet_type == packet.BINARY_EVENT or \497                    pkt.packet_type == packet.BINARY_ACK:498                self._binary_packet = pkt499            elif pkt.packet_type == packet.CONNECT_ERROR:500                self._handle_error(pkt.namespace, pkt.data)501            else:502                raise ValueError('Unknown packet type.')503 504    def _handle_eio_disconnect(self):505        """Handle the Engine.IO disconnection event."""506        self.logger.info('Engine.IO connection dropped')507        will_reconnect = self.reconnection and self.eio.state == 'connected'508        if self.connected:509            for n in self.namespaces:510                self._trigger_event('disconnect', namespace=n)511                if not will_reconnect:512                    self._trigger_event('__disconnect_final', namespace=n)513            self.namespaces = {}514            self.connected = False515        self.callbacks = {}516        self._binary_packet = None517        self.sid = None518        if will_reconnect:519            self._reconnect_task = self.start_background_task(520                self._handle_reconnect)521 522    def _engineio_client_class(self):523        return engineio.Client524 
codekingpro/portable-devtools · Team Ai