Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
async_client.py562 linesDownload Raw Back to socketio
1import asyncio2import logging3import random4 5import engineio6 7from . import base_client8from . import exceptions9from . import packet10 11default_logger = logging.getLogger('socketio.client')12 13 14class AsyncClient(base_client.BaseClient):15    """A Socket.IO client for asyncio.16 17    This class implements a fully compliant Socket.IO web client with support18    for websocket and long-polling transports.19 20    :param reconnection: ``True`` if the client should automatically attempt to21                         reconnect to the server after an interruption, or22                         ``False`` to not reconnect. The default is ``True``.23    :param reconnection_attempts: How many reconnection attempts to issue24                                  before giving up, or 0 for infinite attempts.25                                  The default is 0.26    :param reconnection_delay: How long to wait in seconds before the first27                               reconnection attempt. Each successive attempt28                               doubles this delay.29    :param reconnection_delay_max: The maximum delay between reconnection30                                   attempts.31    :param randomization_factor: Randomization amount for each delay between32                                 reconnection attempts. The default is 0.5,33                                 which means that each delay is randomly34                                 adjusted by +/- 50%.35    :param logger: To enable logging set to ``True`` or pass a logger object to36                   use. To disable logging set to ``False``. The default is37                   ``False``. Note that fatal errors are logged even when38                   ``logger`` is ``False``.39    :param json: An alternative json module to use for encoding and decoding40                 packets. Custom json modules must have ``dumps`` and ``loads``41                 functions that are compatible with the standard library42                 versions.43    :param handle_sigint: Set to ``True`` to automatically handle disconnection44                          when the process is interrupted, or to ``False`` to45                          leave interrupt handling to the calling application.46                          Interrupt handling can only be enabled when the47                          client instance is created in the main thread.48 49    The Engine.IO configuration supports the following settings:50 51    :param request_timeout: A timeout in seconds for requests. The default is52                            5 seconds.53    :param http_session: an initialized ``aiohttp.ClientSession`` object to be54                         used when sending requests to the server. Use it if55                         you need to add special client options such as proxy56                         servers, SSL certificates, etc.57    :param ssl_verify: ``True`` to verify SSL certificates, or ``False`` to58                       skip SSL certificate verification, allowing59                       connections to servers with self signed certificates.60                       The default is ``True``.61    :param engineio_logger: To enable Engine.IO logging set to ``True`` or pass62                            a logger object to use. To disable logging set to63                            ``False``. The default is ``False``. Note that64                            fatal errors are logged even when65                            ``engineio_logger`` is ``False``.66    """67    def is_asyncio_based(self):68        return True69 70    async 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        Note: this method is a coroutine.112 113        Example usage::114 115            sio = socketio.AsyncClient()116            sio.connect('http://localhost:5000')117        """118        if self.connected:119            raise exceptions.ConnectionError('Already connected')120 121        self.connection_url = url122        self.connection_headers = headers123        self.connection_auth = auth124        self.connection_transports = transports125        self.connection_namespaces = namespaces126        self.socketio_path = socketio_path127 128        if namespaces is None:129            namespaces = list(set(self.handlers.keys()).union(130                set(self.namespace_handlers.keys())))131            if len(namespaces) == 0:132                namespaces = ['/']133        elif isinstance(namespaces, str):134            namespaces = [namespaces]135        self.connection_namespaces = namespaces136        self.namespaces = {}137        if self._connect_event is None:138            self._connect_event = self.eio.create_event()139        else:140            self._connect_event.clear()141        real_url = await self._get_real_value(self.connection_url)142        real_headers = await self._get_real_value(self.connection_headers)143        try:144            await self.eio.connect(real_url, headers=real_headers,145                                   transports=transports,146                                   engineio_path=socketio_path)147        except engineio.exceptions.ConnectionError as exc:148            for n in self.connection_namespaces:149                await self._trigger_event(150                    'connect_error', n,151                    exc.args[1] if len(exc.args) > 1 else exc.args[0])152            if retry:  # pragma: no cover153                await self._handle_reconnect()154                if self.eio.state == 'connected':155                    return156            raise exceptions.ConnectionError(exc.args[0]) from None157 158        if wait:159            try:160                while True:161                    await asyncio.wait_for(self._connect_event.wait(),162                                           wait_timeout)163                    self._connect_event.clear()164                    if set(self.namespaces) == set(self.connection_namespaces):165                        break166            except asyncio.TimeoutError:167                pass168            if set(self.namespaces) != set(self.connection_namespaces):169                await self.disconnect()170                raise exceptions.ConnectionError(171                    'One or more namespaces failed to connect')172 173        self.connected = True174 175    async def wait(self):176        """Wait until the connection with the server ends.177 178        Client applications can use this function to block the main thread179        during the life of the connection.180 181        Note: this method is a coroutine.182        """183        while True:184            await self.eio.wait()185            await self.sleep(1)  # give the reconnect task time to start up186            if not self._reconnect_task:187                break188            await self._reconnect_task189            if self.eio.state != 'connected':190                break191 192    async def emit(self, event, data=None, namespace=None, callback=None):193        """Emit a custom event to the server.194 195        :param event: The event name. It can be any string. The event names196                      ``'connect'``, ``'message'`` and ``'disconnect'`` are197                      reserved and should not be used.198        :param data: The data to send to the server. Data can be of199                     type ``str``, ``bytes``, ``list`` or ``dict``. To send200                     multiple arguments, use a tuple where each element is of201                     one of the types indicated above.202        :param namespace: The Socket.IO namespace for the event. If this203                          argument is omitted the event is emitted to the204                          default namespace.205        :param callback: If given, this function will be called to acknowledge206                         the server has received the message. The arguments207                         that will be passed to the function are those provided208                         by the server.209 210        Note: this method is not designed to be used concurrently. If multiple211        tasks are emitting at the same time on the same client connection, then212        messages composed of multiple packets may end up being sent in an213        incorrect sequence. Use standard concurrency solutions (such as a Lock214        object) to prevent this situation.215 216        Note 2: this method is a coroutine.217        """218        namespace = namespace or '/'219        if namespace not in self.namespaces:220            raise exceptions.BadNamespaceError(221                namespace + ' is not a connected namespace.')222        self.logger.info('Emitting event "%s" [%s]', event, namespace)223        if callback is not None:224            id = self._generate_ack_id(namespace, callback)225        else:226            id = None227        # tuples are expanded to multiple arguments, everything else is sent228        # as a single argument229        if isinstance(data, tuple):230            data = list(data)231        elif data is not None:232            data = [data]233        else:234            data = []235        await self._send_packet(self.packet_class(236            packet.EVENT, namespace=namespace, data=[event] + data, id=id))237 238    async def send(self, data, namespace=None, callback=None):239        """Send a message to the server.240 241        This function emits an event with the name ``'message'``. Use242        :func:`emit` to issue custom event names.243 244        :param data: The data to send to the server. Data can be of245                     type ``str``, ``bytes``, ``list`` or ``dict``. To send246                     multiple arguments, use a tuple where each element is of247                     one of the types indicated above.248        :param namespace: The Socket.IO namespace for the event. If this249                          argument is omitted the event is emitted to the250                          default namespace.251        :param callback: If given, this function will be called to acknowledge252                         the server has received the message. The arguments253                         that will be passed to the function are those provided254                         by the server.255 256        Note: this method is a coroutine.257        """258        await self.emit('message', data=data, namespace=namespace,259                        callback=callback)260 261    async def call(self, event, data=None, namespace=None, timeout=60):262        """Emit a custom event to the server and wait for the response.263 264        This method issues an emit with a callback and waits for the callback265        to be invoked before returning. If the callback isn't invoked before266        the timeout, then a ``TimeoutError`` exception is raised. If the267        Socket.IO connection drops during the wait, this method still waits268        until the specified timeout.269 270        :param event: The event name. It can be any string. The event names271                      ``'connect'``, ``'message'`` and ``'disconnect'`` are272                      reserved and should not be used.273        :param data: The data to send to the server. Data can be of274                     type ``str``, ``bytes``, ``list`` or ``dict``. To send275                     multiple arguments, use a tuple where each element is of276                     one of the types indicated above.277        :param namespace: The Socket.IO namespace for the event. If this278                          argument is omitted the event is emitted to the279                          default namespace.280        :param timeout: The waiting timeout. If the timeout is reached before281                        the server acknowledges the event, then a282                        ``TimeoutError`` exception is raised.283 284        Note: this method is not designed to be used concurrently. If multiple285        tasks are emitting at the same time on the same client connection, then286        messages composed of multiple packets may end up being sent in an287        incorrect sequence. Use standard concurrency solutions (such as a Lock288        object) to prevent this situation.289 290        Note 2: this method is a coroutine.291        """292        callback_event = self.eio.create_event()293        callback_args = []294 295        def event_callback(*args):296            callback_args.append(args)297            callback_event.set()298 299        await self.emit(event, data=data, namespace=namespace,300                        callback=event_callback)301        try:302            await asyncio.wait_for(callback_event.wait(), timeout)303        except asyncio.TimeoutError:304            raise exceptions.TimeoutError() from None305        return callback_args[0] if len(callback_args[0]) > 1 \306            else callback_args[0][0] if len(callback_args[0]) == 1 \307            else None308 309    async def disconnect(self):310        """Disconnect from the server.311 312        Note: this method is a coroutine.313        """314        # here we just request the disconnection315        # later in _handle_eio_disconnect we invoke the disconnect handler316        for n in self.namespaces:317            await self._send_packet(self.packet_class(packet.DISCONNECT,318                                    namespace=n))319        await self.eio.disconnect(abort=True)320 321    def start_background_task(self, target, *args, **kwargs):322        """Start a background task using the appropriate async model.323 324        This is a utility function that applications can use to start a325        background task using the method that is compatible with the326        selected async mode.327 328        :param target: the target function to execute.329        :param args: arguments to pass to the function.330        :param kwargs: keyword arguments to pass to the function.331 332        The return value is a ``asyncio.Task`` object.333        """334        return self.eio.start_background_task(target, *args, **kwargs)335 336    async def sleep(self, seconds=0):337        """Sleep for the requested amount of time using the appropriate async338        model.339 340        This is a utility function that applications can use to put a task to341        sleep without having to worry about using the correct call for the342        selected async mode.343 344        Note: this method is a coroutine.345        """346        return await self.eio.sleep(seconds)347 348    async def _get_real_value(self, value):349        """Return the actual value, for parameters that can also be given as350        callables."""351        if not callable(value):352            return value353        if asyncio.iscoroutinefunction(value):354            return await value()355        return value()356 357    async def _send_packet(self, pkt):358        """Send a Socket.IO packet to the server."""359        encoded_packet = pkt.encode()360        if isinstance(encoded_packet, list):361            for ep in encoded_packet:362                await self.eio.send(ep)363        else:364            await self.eio.send(encoded_packet)365 366    async def _handle_connect(self, namespace, data):367        namespace = namespace or '/'368        if namespace not in self.namespaces:369            self.logger.info('Namespace {} is connected'.format(namespace))370            self.namespaces[namespace] = (data or {}).get('sid', self.sid)371            await self._trigger_event('connect', namespace=namespace)372            self._connect_event.set()373 374    async def _handle_disconnect(self, namespace):375        if not self.connected:376            return377        namespace = namespace or '/'378        await self._trigger_event('disconnect', namespace=namespace)379        await self._trigger_event('__disconnect_final', namespace=namespace)380        if namespace in self.namespaces:381            del self.namespaces[namespace]382        if not self.namespaces:383            self.connected = False384            await self.eio.disconnect(abort=True)385 386    async def _handle_event(self, namespace, id, data):387        namespace = namespace or '/'388        self.logger.info('Received event "%s" [%s]', data[0], namespace)389        r = await self._trigger_event(data[0], namespace, *data[1:])390        if id is not None:391            # send ACK packet with the response returned by the handler392            # tuples are expanded as multiple arguments393            if r is None:394                data = []395            elif isinstance(r, tuple):396                data = list(r)397            else:398                data = [r]399            await self._send_packet(self.packet_class(400                packet.ACK, namespace=namespace, id=id, data=data))401 402    async def _handle_ack(self, namespace, id, data):403        namespace = namespace or '/'404        self.logger.info('Received ack [%s]', namespace)405        callback = None406        try:407            callback = self.callbacks[namespace][id]408        except KeyError:409            # if we get an unknown callback we just ignore it410            self.logger.warning('Unknown callback received, ignoring.')411        else:412            del self.callbacks[namespace][id]413        if callback is not None:414            if asyncio.iscoroutinefunction(callback):415                await callback(*data)416            else:417                callback(*data)418 419    async def _handle_error(self, namespace, data):420        namespace = namespace or '/'421        self.logger.info('Connection to namespace {} was rejected'.format(422            namespace))423        if data is None:424            data = tuple()425        elif not isinstance(data, (tuple, list)):426            data = (data,)427        await self._trigger_event('connect_error', namespace, *data)428        self._connect_event.set()429        if namespace in self.namespaces:430            del self.namespaces[namespace]431        if namespace == '/':432            self.namespaces = {}433            self.connected = False434 435    async def _trigger_event(self, event, namespace, *args):436        """Invoke an application event handler."""437        # first see if we have an explicit handler for the event438        handler, args = self._get_event_handler(event, namespace, args)439        if handler:440            if asyncio.iscoroutinefunction(handler):441                try:442                    ret = await handler(*args)443                except asyncio.CancelledError:  # pragma: no cover444                    ret = None445            else:446                ret = handler(*args)447            return ret448 449        # or else, forward the event to a namepsace handler if one exists450        handler, args = self._get_namespace_handler(namespace, args)451        if handler:452            return await handler.trigger_event(event, *args)453 454    async def _handle_reconnect(self):455        if self._reconnect_abort is None:  # pragma: no cover456            self._reconnect_abort = self.eio.create_event()457        self._reconnect_abort.clear()458        base_client.reconnecting_clients.append(self)459        attempt_count = 0460        current_delay = self.reconnection_delay461        while True:462            delay = current_delay463            current_delay *= 2464            if delay > self.reconnection_delay_max:465                delay = self.reconnection_delay_max466            delay += self.randomization_factor * (2 * random.random() - 1)467            self.logger.info(468                'Connection failed, new attempt in {:.02f} seconds'.format(469                    delay))470            try:471                await asyncio.wait_for(self._reconnect_abort.wait(), delay)472                self.logger.info('Reconnect task aborted')473                for n in self.connection_namespaces:474                    await self._trigger_event('__disconnect_final',475                                              namespace=n)476                break477            except (asyncio.TimeoutError, asyncio.CancelledError):478                pass479            attempt_count += 1480            try:481                await self.connect(self.connection_url,482                                   headers=self.connection_headers,483                                   auth=self.connection_auth,484                                   transports=self.connection_transports,485                                   namespaces=self.connection_namespaces,486                                   socketio_path=self.socketio_path,487                                   retry=False)488            except (exceptions.ConnectionError, ValueError):489                pass490            else:491                self.logger.info('Reconnection successful')492                self._reconnect_task = None493                break494            if self.reconnection_attempts and \495                    attempt_count >= self.reconnection_attempts:496                self.logger.info(497                    'Maximum reconnection attempts reached, giving up')498                for n in self.connection_namespaces:499                    await self._trigger_event('__disconnect_final',500                                              namespace=n)501                break502        base_client.reconnecting_clients.remove(self)503 504    async def _handle_eio_connect(self):505        """Handle the Engine.IO connection event."""506        self.logger.info('Engine.IO connection established')507        self.sid = self.eio.sid508        real_auth = await self._get_real_value(self.connection_auth) or {}509        for n in self.connection_namespaces:510            await self._send_packet(self.packet_class(511                packet.CONNECT, data=real_auth, namespace=n))512 513    async def _handle_eio_message(self, data):514        """Dispatch Engine.IO messages."""515        if self._binary_packet:516            pkt = self._binary_packet517            if pkt.add_attachment(data):518                self._binary_packet = None519                if pkt.packet_type == packet.BINARY_EVENT:520                    await self._handle_event(pkt.namespace, pkt.id, pkt.data)521                else:522                    await self._handle_ack(pkt.namespace, pkt.id, pkt.data)523        else:524            pkt = self.packet_class(encoded_packet=data)525            if pkt.packet_type == packet.CONNECT:526                await self._handle_connect(pkt.namespace, pkt.data)527            elif pkt.packet_type == packet.DISCONNECT:528                await self._handle_disconnect(pkt.namespace)529            elif pkt.packet_type == packet.EVENT:530                await self._handle_event(pkt.namespace, pkt.id, pkt.data)531            elif pkt.packet_type == packet.ACK:532                await self._handle_ack(pkt.namespace, pkt.id, pkt.data)533            elif pkt.packet_type == packet.BINARY_EVENT or \534                    pkt.packet_type == packet.BINARY_ACK:535                self._binary_packet = pkt536            elif pkt.packet_type == packet.CONNECT_ERROR:537                await self._handle_error(pkt.namespace, pkt.data)538            else:539                raise ValueError('Unknown packet type.')540 541    async def _handle_eio_disconnect(self):542        """Handle the Engine.IO disconnection event."""543        self.logger.info('Engine.IO connection dropped')544        will_reconnect = self.reconnection and self.eio.state == 'connected'545        if self.connected:546            for n in self.namespaces:547                await self._trigger_event('disconnect', namespace=n)548                if not will_reconnect:549                    await self._trigger_event('__disconnect_final',550                                              namespace=n)551            self.namespaces = {}552            self.connected = False553        self.callbacks = {}554        self._binary_packet = None555        self.sid = None556        if will_reconnect:557            self._reconnect_task = self.start_background_task(558                self._handle_reconnect)559 560    def _engineio_client_class(self):561        return engineio.AsyncClient562 
codekingpro/portable-devtools · Team Ai