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