codekingpro/portable-devtools
114k
1import asyncio2 3from engineio import packet as eio_packet4from socketio import packet5from .base_manager import BaseManager6 7 8class AsyncManager(BaseManager):9 """Manage a client list for an asyncio server."""10 async def can_disconnect(self, sid, namespace):11 return self.is_connected(sid, namespace)12 13 async def emit(self, event, data, namespace, room=None, skip_sid=None,14 callback=None, **kwargs):15 """Emit a message to a single client, a room, or all the clients16 connected to the namespace.17 18 Note: this method is a coroutine.19 """20 if namespace not in self.rooms:21 return22 if isinstance(data, tuple):23 # tuples are expanded to multiple arguments, everything else is24 # sent as a single argument25 data = list(data)26 elif data is not None:27 data = [data]28 else:29 data = []30 if not isinstance(skip_sid, list):31 skip_sid = [skip_sid]32 tasks = []33 if not callback:34 # when callbacks aren't used the packets sent to each recipient are35 # identical, so they can be generated once and reused36 pkt = self.server.packet_class(37 packet.EVENT, namespace=namespace, data=[event] + data)38 encoded_packet = pkt.encode()39 if not isinstance(encoded_packet, list):40 encoded_packet = [encoded_packet]41 eio_pkt = [eio_packet.Packet(eio_packet.MESSAGE, p)42 for p in encoded_packet]43 for sid, eio_sid in self.get_participants(namespace, room):44 if sid not in skip_sid:45 for p in eio_pkt:46 tasks.append(asyncio.create_task(47 self.server._send_eio_packet(eio_sid, p)))48 else:49 # callbacks are used, so each recipient must be sent a packet that50 # contains a unique callback id51 # note that callbacks when addressing a group of people are52 # implemented but not tested or supported53 for sid, eio_sid in self.get_participants(namespace, room):54 if sid not in skip_sid: # pragma: no branch55 id = self._generate_ack_id(sid, callback)56 pkt = self.server.packet_class(57 packet.EVENT, namespace=namespace, data=[event] + data,58 id=id)59 tasks.append(asyncio.create_task(60 self.server._send_packet(eio_sid, pkt)))61 if tasks == []: # pragma: no cover62 return63 await asyncio.wait(tasks)64 65 async def connect(self, eio_sid, namespace):66 """Register a client connection to a namespace.67 68 Note: this method is a coroutine.69 """70 return super().connect(eio_sid, namespace)71 72 async def disconnect(self, sid, namespace, **kwargs):73 """Disconnect a client.74 75 Note: this method is a coroutine.76 """77 return self.basic_disconnect(sid, namespace, **kwargs)78 79 async def enter_room(self, sid, namespace, room, eio_sid=None):80 """Add a client to a room.81 82 Note: this method is a coroutine.83 """84 return self.basic_enter_room(sid, namespace, room, eio_sid=eio_sid)85 86 async def leave_room(self, sid, namespace, room):87 """Remove a client from a room.88 89 Note: this method is a coroutine.90 """91 return self.basic_leave_room(sid, namespace, room)92 93 async def close_room(self, room, namespace):94 """Remove all participants from a room.95 96 Note: this method is a coroutine.97 """98 return self.basic_close_room(room, namespace)99 100 async def trigger_callback(self, sid, id, data):101 """Invoke an application callback.102 103 Note: this method is a coroutine.104 """105 callback = None106 try:107 callback = self.callbacks[sid][id]108 except KeyError:109 # if we get an unknown callback we just ignore it110 self._get_logger().warning('Unknown callback received, ignoring.')111 else:112 del self.callbacks[sid][id]113 if callback is not None:114 ret = callback(*data)115 if asyncio.iscoroutine(ret):116 try:117 await ret118 except asyncio.CancelledError: # pragma: no cover119 pass120 