Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
async_pubsub_manager.py243 linesDownload Raw Back to socketio
1import asyncio2from functools import partial3import uuid4 5from engineio import json6import pickle7 8from .async_manager import AsyncManager9 10 11class AsyncPubSubManager(AsyncManager):12    """Manage a client list attached to a pub/sub backend under asyncio.13 14    This is a base class that enables multiple servers to share the list of15    clients, with the servers communicating events through a pub/sub backend.16    The use of a pub/sub backend also allows any client connected to the17    backend to emit events addressed to Socket.IO clients.18 19    The actual backends must be implemented by subclasses, this class only20    provides a pub/sub generic framework for asyncio applications.21 22    :param channel: The channel name on which the server sends and receives23                    notifications.24    """25    name = 'asyncpubsub'26 27    def __init__(self, channel='socketio', write_only=False, logger=None):28        super().__init__()29        self.channel = channel30        self.write_only = write_only31        self.host_id = uuid.uuid4().hex32        self.logger = logger33 34    def initialize(self):35        super().initialize()36        if not self.write_only:37            self.thread = self.server.start_background_task(self._thread)38        self._get_logger().info(self.name + ' backend initialized.')39 40    async def emit(self, event, data, namespace=None, room=None, skip_sid=None,41                   callback=None, **kwargs):42        """Emit a message to a single client, a room, or all the clients43        connected to the namespace.44 45        This method takes care or propagating the message to all the servers46        that are connected through the message queue.47 48        The parameters are the same as in :meth:`.Server.emit`.49 50        Note: this method is a coroutine.51        """52        if kwargs.get('ignore_queue'):53            return await super().emit(54                event, data, namespace=namespace, room=room, skip_sid=skip_sid,55                callback=callback)56        namespace = namespace or '/'57        if callback is not None:58            if self.server is None:59                raise RuntimeError('Callbacks can only be issued from the '60                                   'context of a server.')61            if room is None:62                raise ValueError('Cannot use callback without a room set.')63            id = self._generate_ack_id(room, callback)64            callback = (room, namespace, id)65        else:66            callback = None67        message = {'method': 'emit', 'event': event, 'data': data,68                   'namespace': namespace, 'room': room,69                   'skip_sid': skip_sid, 'callback': callback,70                   'host_id': self.host_id}71        await self._handle_emit(message)  # handle in this host72        await self._publish(message)  # notify other hosts73 74    async def can_disconnect(self, sid, namespace):75        if self.is_connected(sid, namespace):76            # client is in this server, so we can disconnect directly77            return await super().can_disconnect(sid, namespace)78        else:79            # client is in another server, so we post request to the queue80            await self._publish({'method': 'disconnect', 'sid': sid,81                                 'namespace': namespace or '/',82                                 'host_id': self.host_id})83 84    async def disconnect(self, sid, namespace, **kwargs):85        if kwargs.get('ignore_queue'):86            return await super().disconnect(87                sid, namespace=namespace)88        message = {'method': 'disconnect', 'sid': sid,89                   'namespace': namespace or '/', 'host_id': self.host_id}90        await self._handle_disconnect(message)  # handle in this host91        await self._publish(message)  # notify other hosts92 93    async def enter_room(self, sid, namespace, room, eio_sid=None):94        if self.is_connected(sid, namespace):95            # client is in this server, so we can disconnect directly96            return await super().enter_room(sid, namespace, room,97                                            eio_sid=eio_sid)98        else:99            message = {'method': 'enter_room', 'sid': sid, 'room': room,100                       'namespace': namespace or '/', 'host_id': self.host_id}101            await self._publish(message)  # notify other hosts102 103    async def leave_room(self, sid, namespace, room):104        if self.is_connected(sid, namespace):105            # client is in this server, so we can disconnect directly106            return await super().leave_room(sid, namespace, room)107        else:108            message = {'method': 'leave_room', 'sid': sid, 'room': room,109                       'namespace': namespace or '/', 'host_id': self.host_id}110            await self._publish(message)  # notify other hosts111 112    async def close_room(self, room, namespace=None):113        message = {'method': 'close_room', 'room': room,114                   'namespace': namespace or '/', 'host_id': self.host_id}115        await self._handle_close_room(message)  # handle in this host116        await self._publish(message)  # notify other hosts117 118    async def _publish(self, data):119        """Publish a message on the Socket.IO channel.120 121        This method needs to be implemented by the different subclasses that122        support pub/sub backends.123        """124        raise NotImplementedError('This method must be implemented in a '125                                  'subclass.')  # pragma: no cover126 127    async def _listen(self):128        """Return the next message published on the Socket.IO channel,129        blocking until a message is available.130 131        This method needs to be implemented by the different subclasses that132        support pub/sub backends.133        """134        raise NotImplementedError('This method must be implemented in a '135                                  'subclass.')  # pragma: no cover136 137    async def _handle_emit(self, message):138        # Events with callbacks are very tricky to handle across hosts139        # Here in the receiving end we set up a local callback that preserves140        # the callback host and id from the sender141        remote_callback = message.get('callback')142        remote_host_id = message.get('host_id')143        if remote_callback is not None and len(remote_callback) == 3:144            callback = partial(self._return_callback, remote_host_id,145                               *remote_callback)146        else:147            callback = None148        await super().emit(message['event'], message['data'],149                           namespace=message.get('namespace'),150                           room=message.get('room'),151                           skip_sid=message.get('skip_sid'),152                           callback=callback)153 154    async def _handle_callback(self, message):155        if self.host_id == message.get('host_id'):156            try:157                sid = message['sid']158                id = message['id']159                args = message['args']160            except KeyError:161                return162            await self.trigger_callback(sid, id, args)163 164    async def _return_callback(self, host_id, sid, namespace, callback_id,165                               *args):166        # When an event callback is received, the callback is returned back167        # the sender, which is identified by the host_id168        if host_id == self.host_id:169            await self.trigger_callback(sid, callback_id, args)170        else:171            await self._publish({'method': 'callback', 'host_id': host_id,172                                 'sid': sid, 'namespace': namespace,173                                 'id': callback_id, 'args': args})174 175    async def _handle_disconnect(self, message):176        await self.server.disconnect(sid=message.get('sid'),177                                     namespace=message.get('namespace'),178                                     ignore_queue=True)179 180    async def _handle_enter_room(self, message):181        sid = message.get('sid')182        namespace = message.get('namespace')183        if self.is_connected(sid, namespace):184            await super().enter_room(sid, namespace, message.get('room'))185 186    async def _handle_leave_room(self, message):187        sid = message.get('sid')188        namespace = message.get('namespace')189        if self.is_connected(sid, namespace):190            await super().leave_room(sid, namespace, message.get('room'))191 192    async def _handle_close_room(self, message):193        await super().close_room(room=message.get('room'),194                                 namespace=message.get('namespace'))195 196    async def _thread(self):197        while True:198            try:199                async for message in self._listen():  # pragma: no branch200                    data = None201                    if isinstance(message, dict):202                        data = message203                    else:204                        if isinstance(message, bytes):  # pragma: no cover205                            try:206                                data = pickle.loads(message)207                            except:208                                pass209                        if data is None:210                            try:211                                data = json.loads(message)212                            except:213                                pass214                    if data and 'method' in data:215                        self._get_logger().debug('pubsub message: {}'.format(216                            data['method']))217                        try:218                            if data['method'] == 'callback':219                                await self._handle_callback(data)220                            elif data.get('host_id') != self.host_id:221                                if data['method'] == 'emit':222                                    await self._handle_emit(data)223                                elif data['method'] == 'disconnect':224                                    await self._handle_disconnect(data)225                                elif data['method'] == 'enter_room':226                                    await self._handle_enter_room(data)227                                elif data['method'] == 'leave_room':228                                    await self._handle_leave_room(data)229                                elif data['method'] == 'close_room':230                                    await self._handle_close_room(data)231                        except asyncio.CancelledError:232                            raise  # let the outer try/except handle it233                        except Exception:234                            self.server.logger.exception(235                                'Handler error in pubsub listening thread')236                self.server.logger.error('pubsub listen() exited unexpectedly')237                break  # loop should never exit except in unit tests!238            except asyncio.CancelledError:  # pragma: no cover239                break240            except Exception:  # pragma: no cover241                self.server.logger.exception('Unexpected Error in pubsub '242                                             'listening thread')243 
codekingpro/portable-devtools · Team Ai