codekingpro/portable-devtools
114k
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 