Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
pubsub_manager.py233 linesDownload Raw Back to socketio
1from functools import partial2import uuid3 4from engineio import json5import pickle6 7from .manager import Manager8 9 10class PubSubManager(Manager):11    """Manage a client list attached to a pub/sub backend.12 13    This is a base class that enables multiple servers to share the list of14    clients, with the servers communicating events through a pub/sub backend.15    The use of a pub/sub backend also allows any client connected to the16    backend to emit events addressed to Socket.IO clients.17 18    The actual backends must be implemented by subclasses, this class only19    provides a pub/sub generic framework.20 21    :param channel: The channel name on which the server sends and receives22                    notifications.23    """24    name = 'pubsub'25 26    def __init__(self, channel='socketio', write_only=False, logger=None):27        super().__init__()28        self.channel = channel29        self.write_only = write_only30        self.host_id = uuid.uuid4().hex31        self.logger = logger32 33    def initialize(self):34        super().initialize()35        if not self.write_only:36            self.thread = self.server.start_background_task(self._thread)37        self._get_logger().info(self.name + ' backend initialized.')38 39    def emit(self, event, data, namespace=None, room=None, skip_sid=None,40             callback=None, **kwargs):41        """Emit a message to a single client, a room, or all the clients42        connected to the namespace.43 44        This method takes care or propagating the message to all the servers45        that are connected through the message queue.46 47        The parameters are the same as in :meth:`.Server.emit`.48        """49        if kwargs.get('ignore_queue'):50            return super().emit(51                event, data, namespace=namespace, room=room, skip_sid=skip_sid,52                callback=callback)53        namespace = namespace or '/'54        if callback is not None:55            if self.server is None:56                raise RuntimeError('Callbacks can only be issued from the '57                                   'context of a server.')58            if room is None:59                raise ValueError('Cannot use callback without a room set.')60            id = self._generate_ack_id(room, callback)61            callback = (room, namespace, id)62        else:63            callback = None64        message = {'method': 'emit', 'event': event, 'data': data,65                   'namespace': namespace, 'room': room,66                   'skip_sid': skip_sid, 'callback': callback,67                   'host_id': self.host_id}68        self._handle_emit(message)  # handle in this host69        self._publish(message)  # notify other hosts70 71    def can_disconnect(self, sid, namespace):72        if self.is_connected(sid, namespace):73            # client is in this server, so we can disconnect directly74            return super().can_disconnect(sid, namespace)75        else:76            # client is in another server, so we post request to the queue77            message = {'method': 'disconnect', 'sid': sid,78                       'namespace': namespace or '/', 'host_id': self.host_id}79            self._handle_disconnect(message)  # handle in this host80            self._publish(message)  # notify other hosts81 82    def disconnect(self, sid, namespace=None, **kwargs):83        if kwargs.get('ignore_queue'):84            return super().disconnect(sid, namespace=namespace)85        message = {'method': 'disconnect', 'sid': sid,86                   'namespace': namespace or '/', 'host_id': self.host_id}87        self._handle_disconnect(message)  # handle in this host88        self._publish(message)  # notify other hosts89 90    def enter_room(self, sid, namespace, room, eio_sid=None):91        if self.is_connected(sid, namespace):92            # client is in this server, so we can add to the room directly93            return super().enter_room(sid, namespace, room, eio_sid=eio_sid)94        else:95            message = {'method': 'enter_room', 'sid': sid, 'room': room,96                       'namespace': namespace or '/', 'host_id': self.host_id}97            self._publish(message)  # notify other hosts98 99    def leave_room(self, sid, namespace, room):100        if self.is_connected(sid, namespace):101            # client is in this server, so we can remove from the room directly102            return super().leave_room(sid, namespace, room)103        else:104            message = {'method': 'leave_room', 'sid': sid, 'room': room,105                       'namespace': namespace or '/', 'host_id': self.host_id}106            self._publish(message)  # notify other hosts107 108    def close_room(self, room, namespace=None):109        message = {'method': 'close_room', 'room': room,110                   'namespace': namespace or '/', 'host_id': self.host_id}111        self._handle_close_room(message)  # handle in this host112        self._publish(message)  # notify other hosts113 114    def _publish(self, data):115        """Publish a message on the Socket.IO channel.116 117        This method needs to be implemented by the different subclasses that118        support pub/sub backends.119        """120        raise NotImplementedError('This method must be implemented in a '121                                  'subclass.')  # pragma: no cover122 123    def _listen(self):124        """Return the next message published on the Socket.IO channel,125        blocking until a message is available.126 127        This method needs to be implemented by the different subclasses that128        support pub/sub backends.129        """130        raise NotImplementedError('This method must be implemented in a '131                                  'subclass.')  # pragma: no cover132 133    def _handle_emit(self, message):134        # Events with callbacks are very tricky to handle across hosts135        # Here in the receiving end we set up a local callback that preserves136        # the callback host and id from the sender137        remote_callback = message.get('callback')138        remote_host_id = message.get('host_id')139        if remote_callback is not None and len(remote_callback) == 3:140            callback = partial(self._return_callback, remote_host_id,141                               *remote_callback)142        else:143            callback = None144        super().emit(message['event'], message['data'],145                     namespace=message.get('namespace'),146                     room=message.get('room'),147                     skip_sid=message.get('skip_sid'), callback=callback)148 149    def _handle_callback(self, message):150        if self.host_id == message.get('host_id'):151            try:152                sid = message['sid']153                id = message['id']154                args = message['args']155            except KeyError:156                return157            self.trigger_callback(sid, id, args)158 159    def _return_callback(self, host_id, sid, namespace, callback_id, *args):160        # When an event callback is received, the callback is returned back161        # to the sender, which is identified by the host_id162        if host_id == self.host_id:163            self.trigger_callback(sid, callback_id, args)164        else:165            self._publish({'method': 'callback', 'host_id': host_id,166                           'sid': sid, 'namespace': namespace,167                           'id': callback_id, 'args': args})168 169    def _handle_disconnect(self, message):170        self.server.disconnect(sid=message.get('sid'),171                               namespace=message.get('namespace'),172                               ignore_queue=True)173 174    def _handle_enter_room(self, message):175        sid = message.get('sid')176        namespace = message.get('namespace')177        if self.is_connected(sid, namespace):178            super().enter_room(sid, namespace, message.get('room'))179 180    def _handle_leave_room(self, message):181        sid = message.get('sid')182        namespace = message.get('namespace')183        if self.is_connected(sid, namespace):184            super().leave_room(sid, namespace, message.get('room'))185 186    def _handle_close_room(self, message):187        super().close_room(room=message.get('room'),188                           namespace=message.get('namespace'))189 190    def _thread(self):191        while True:192            try:193                for message in self._listen():194                    data = None195                    if isinstance(message, dict):196                        data = message197                    else:198                        if isinstance(message, bytes):  # pragma: no cover199                            try:200                                data = pickle.loads(message)201                            except:202                                pass203                        if data is None:204                            try:205                                data = json.loads(message)206                            except:207                                pass208                    if data and 'method' in data:209                        self._get_logger().debug('pubsub message: {}'.format(210                            data['method']))211                        try:212                            if data['method'] == 'callback':213                                self._handle_callback(data)214                            elif data.get('host_id') != self.host_id:215                                if data['method'] == 'emit':216                                    self._handle_emit(data)217                                elif data['method'] == 'disconnect':218                                    self._handle_disconnect(data)219                                elif data['method'] == 'enter_room':220                                    self._handle_enter_room(data)221                                elif data['method'] == 'leave_room':222                                    self._handle_leave_room(data)223                                elif data['method'] == 'close_room':224                                    self._handle_close_room(data)225                        except Exception:226                            self.server.logger.exception(227                                'Handler error in pubsub listening thread')228                self.server.logger.error('pubsub listen() exited unexpectedly')229                break  # loop should never exit except in unit tests!230            except Exception:  # pragma: no cover231                self.server.logger.exception('Unexpected Error in pubsub '232                                             'listening thread')233 
codekingpro/portable-devtools · Team Ai