codekingpro/portable-devtools
115k
1import asyncio2from datetime import datetime3import functools4import os5import socket6import time7from urllib.parse import parse_qs8from .admin import EventBuffer9from .exceptions import ConnectionRefusedError10 11HOSTNAME = socket.gethostname()12PID = os.getpid()13 14 15class InstrumentedAsyncServer:16 def __init__(self, sio, auth=None, namespace='/admin', read_only=False,17 server_id=None, mode='development', server_stats_interval=2):18 """Instrument the Socket.IO server for monitoring with the `Socket.IO19 Admin UI <https://socket.io/docs/v4/admin-ui/>`_.20 """21 if auth is None:22 raise ValueError('auth must be specified')23 self.sio = sio24 self.auth = auth25 self.admin_namespace = namespace26 self.read_only = read_only27 self.server_id = server_id or (28 self.sio.manager.host_id if hasattr(self.sio.manager, 'host_id')29 else HOSTNAME30 )31 self.mode = mode32 self.server_stats_interval = server_stats_interval33 self.admin_queue = []34 self.event_buffer = EventBuffer()35 36 # task that emits "server_stats" every 2 seconds37 self.stop_stats_event = None38 self.stats_task = None39 40 # monkey-patch the server to report metrics to the admin UI41 self.instrument()42 43 def instrument(self):44 self.sio.on('connect', self.admin_connect,45 namespace=self.admin_namespace)46 47 if self.mode == 'development':48 if not self.read_only: # pragma: no branch49 self.sio.on('emit', self.admin_emit,50 namespace=self.admin_namespace)51 self.sio.on('join', self.admin_enter_room,52 namespace=self.admin_namespace)53 self.sio.on('leave', self.admin_leave_room,54 namespace=self.admin_namespace)55 self.sio.on('_disconnect', self.admin_disconnect,56 namespace=self.admin_namespace)57 58 # track socket connection times59 self.sio.manager._timestamps = {}60 61 # report socket.io connections62 self.sio.manager.__connect = self.sio.manager.connect63 self.sio.manager.connect = self._connect64 65 # report socket.io disconnection66 self.sio.manager.__disconnect = self.sio.manager.disconnect67 self.sio.manager.disconnect = self._disconnect68 69 # report join rooms70 self.sio.manager.__basic_enter_room = \71 self.sio.manager.basic_enter_room72 self.sio.manager.basic_enter_room = self._basic_enter_room73 74 # report leave rooms75 self.sio.manager.__basic_leave_room = \76 self.sio.manager.basic_leave_room77 self.sio.manager.basic_leave_room = self._basic_leave_room78 79 # report emit events80 self.sio.manager.__emit = self.sio.manager.emit81 self.sio.manager.emit = self._emit82 83 # report receive events84 self.sio.__handle_event_internal = self.sio._handle_event_internal85 self.sio._handle_event_internal = self._handle_event_internal86 87 # report engine.io connections88 self.sio.eio.on('connect', self._handle_eio_connect)89 self.sio.eio.on('disconnect', self._handle_eio_disconnect)90 91 # report polling packets92 from engineio.async_socket import AsyncSocket93 self.sio.eio.__ok = self.sio.eio._ok94 self.sio.eio._ok = self._eio_http_response95 AsyncSocket.__handle_post_request = AsyncSocket.handle_post_request96 AsyncSocket.handle_post_request = functools.partialmethod(97 self.__class__._eio_handle_post_request, self)98 99 # report websocket packets100 AsyncSocket.__websocket_handler = AsyncSocket._websocket_handler101 AsyncSocket._websocket_handler = functools.partialmethod(102 self.__class__._eio_websocket_handler, self)103 104 # report connected sockets with each ping105 if self.mode == 'development':106 AsyncSocket.__send_ping = AsyncSocket._send_ping107 AsyncSocket._send_ping = functools.partialmethod(108 self.__class__._eio_send_ping, self)109 110 def uninstrument(self): # pragma: no cover111 if self.mode == 'development':112 self.sio.manager.connect = self.sio.manager.__connect113 self.sio.manager.disconnect = self.sio.manager.__disconnect114 self.sio.manager.basic_enter_room = \115 self.sio.manager.__basic_enter_room116 self.sio.manager.basic_leave_room = \117 self.sio.manager.__basic_leave_room118 self.sio.manager.emit = self.sio.manager.__emit119 self.sio._handle_event_internal = self.sio.__handle_event_internal120 self.sio.eio._ok = self.sio.eio.__ok121 122 from engineio.async_socket import AsyncSocket123 AsyncSocket.handle_post_request = AsyncSocket.__handle_post_request124 AsyncSocket._websocket_handler = AsyncSocket.__websocket_handler125 if self.mode == 'development':126 AsyncSocket._send_ping = AsyncSocket.__send_ping127 128 async def admin_connect(self, sid, environ, client_auth):129 authenticated = True130 if self.auth:131 authenticated = False132 if isinstance(self.auth, dict):133 authenticated = client_auth == self.auth134 elif isinstance(self.auth, list):135 authenticated = client_auth in self.auth136 else:137 if asyncio.iscoroutinefunction(self.auth):138 authenticated = await self.auth(client_auth)139 else:140 authenticated = self.auth(client_auth)141 if not authenticated:142 raise ConnectionRefusedError('authentication failed')143 144 async def config(sid):145 await self.sio.sleep(0.1)146 147 # supported features148 features = ['AGGREGATED_EVENTS']149 if not self.read_only:150 features += ['EMIT', 'JOIN', 'LEAVE', 'DISCONNECT', 'MJOIN',151 'MLEAVE', 'MDISCONNECT']152 if self.mode == 'development':153 features.append('ALL_EVENTS')154 await self.sio.emit('config', {'supportedFeatures': features},155 to=sid, namespace=self.admin_namespace)156 157 # send current sockets158 if self.mode == 'development':159 all_sockets = []160 for nsp in self.sio.manager.get_namespaces():161 for sid, eio_sid in self.sio.manager.get_participants(162 nsp, None):163 all_sockets.append(164 self.serialize_socket(sid, nsp, eio_sid))165 await self.sio.emit('all_sockets', all_sockets, to=sid,166 namespace=self.admin_namespace)167 168 self.sio.start_background_task(config, sid)169 self.stop_stats_event = self.sio.eio.create_event()170 self.stats_task = self.sio.start_background_task(171 self._emit_server_stats)172 173 async def admin_emit(self, _, namespace, room_filter, event, *data):174 await self.sio.emit(event, data, to=room_filter, namespace=namespace)175 176 async def admin_enter_room(self, _, namespace, room, room_filter=None):177 for sid, _ in self.sio.manager.get_participants(178 namespace, room_filter):179 await self.sio.enter_room(sid, room, namespace=namespace)180 181 async def admin_leave_room(self, _, namespace, room, room_filter=None):182 for sid, _ in self.sio.manager.get_participants(183 namespace, room_filter):184 await self.sio.leave_room(sid, room, namespace=namespace)185 186 async def admin_disconnect(self, _, namespace, close, room_filter=None):187 for sid, _ in self.sio.manager.get_participants(188 namespace, room_filter):189 await self.sio.disconnect(sid, namespace=namespace)190 191 async def shutdown(self):192 if self.stats_task: # pragma: no branch193 self.stop_stats_event.set()194 await asyncio.gather(self.stats_task)195 196 async def _connect(self, eio_sid, namespace):197 sid = await self.sio.manager.__connect(eio_sid, namespace)198 t = time.time()199 self.sio.manager._timestamps[sid] = t200 serialized_socket = self.serialize_socket(sid, namespace, eio_sid)201 await self.sio.emit('socket_connected', (202 serialized_socket,203 datetime.utcfromtimestamp(t).isoformat() + 'Z',204 ), namespace=self.admin_namespace)205 return sid206 207 async def _disconnect(self, sid, namespace, **kwargs):208 del self.sio.manager._timestamps[sid]209 await self.sio.emit('socket_disconnected', (210 namespace,211 sid,212 'N/A',213 datetime.utcnow().isoformat() + 'Z',214 ), namespace=self.admin_namespace)215 return await self.sio.manager.__disconnect(sid, namespace, **kwargs)216 217 async def _check_for_upgrade(self, eio_sid, sid,218 namespace): # pragma: no cover219 for _ in range(5):220 await self.sio.sleep(5)221 try:222 if self.sio.eio._get_socket(eio_sid).upgraded:223 await self.sio.emit('socket_updated', {224 'id': sid,225 'nsp': namespace,226 'transport': 'websocket',227 }, namespace=self.admin_namespace)228 break229 except KeyError:230 pass231 232 def _basic_enter_room(self, sid, namespace, room, eio_sid=None):233 ret = self.sio.manager.__basic_enter_room(sid, namespace, room,234 eio_sid)235 if room:236 self.admin_queue.append(('room_joined', (237 namespace,238 room,239 sid,240 datetime.utcnow().isoformat() + 'Z',241 )))242 return ret243 244 def _basic_leave_room(self, sid, namespace, room):245 if room:246 self.admin_queue.append(('room_left', (247 namespace,248 room,249 sid,250 datetime.utcnow().isoformat() + 'Z',251 )))252 return self.sio.manager.__basic_leave_room(sid, namespace, room)253 254 async def _emit(self, event, data, namespace, room=None, skip_sid=None,255 callback=None, **kwargs):256 ret = await self.sio.manager.__emit(257 event, data, namespace, room=room, skip_sid=skip_sid,258 callback=callback, **kwargs)259 if namespace != self.admin_namespace:260 event_data = [event] + list(data) if isinstance(data, tuple) \261 else [data]262 if not isinstance(skip_sid, list): # pragma: no branch263 skip_sid = [skip_sid]264 for sid, _ in self.sio.manager.get_participants(namespace, room):265 if sid not in skip_sid:266 await self.sio.emit('event_sent', (267 namespace,268 sid,269 event_data,270 datetime.utcnow().isoformat() + 'Z',271 ), namespace=self.admin_namespace)272 return ret273 274 async def _handle_event_internal(self, server, sid, eio_sid, data,275 namespace, id):276 ret = await self.sio.__handle_event_internal(server, sid, eio_sid,277 data, namespace, id)278 await self.sio.emit('event_received', (279 namespace,280 sid,281 data,282 datetime.utcnow().isoformat() + 'Z',283 ), namespace=self.admin_namespace)284 return ret285 286 async def _handle_eio_connect(self, eio_sid, environ):287 if self.stop_stats_event is None:288 self.stop_stats_event = self.sio.eio.create_event()289 self.stats_task = self.sio.start_background_task(290 self._emit_server_stats)291 292 self.event_buffer.push('rawConnection')293 return await self.sio._handle_eio_connect(eio_sid, environ)294 295 async def _handle_eio_disconnect(self, eio_sid):296 self.event_buffer.push('rawDisconnection')297 return await self.sio._handle_eio_disconnect(eio_sid)298 299 def _eio_http_response(self, packets=None, headers=None, jsonp_index=None):300 ret = self.sio.eio.__ok(packets=packets, headers=headers,301 jsonp_index=jsonp_index)302 self.event_buffer.push('packetsOut')303 self.event_buffer.push('bytesOut', len(ret['response']))304 return ret305 306 async def _eio_handle_post_request(socket, self, environ):307 ret = await socket.__handle_post_request(environ)308 self.event_buffer.push('packetsIn')309 self.event_buffer.push(310 'bytesIn', int(environ.get('CONTENT_LENGTH', 0)))311 return ret312 313 async def _eio_websocket_handler(socket, self, ws):314 async def _send(ws, data):315 self.event_buffer.push('packetsOut')316 self.event_buffer.push('bytesOut', len(data))317 return await ws.__send(data)318 319 async def _wait(ws):320 ret = await ws.__wait()321 self.event_buffer.push('packetsIn')322 self.event_buffer.push('bytesIn', len(ret or ''))323 return ret324 325 ws.__send = ws.send326 ws.send = functools.partial(_send, ws)327 ws.__wait = ws.wait328 ws.wait = functools.partial(_wait, ws)329 return await socket.__websocket_handler(ws)330 331 async def _eio_send_ping(socket, self): # pragma: no cover332 eio_sid = socket.sid333 t = time.time()334 for namespace in self.sio.manager.get_namespaces():335 sid = self.sio.manager.sid_from_eio_sid(eio_sid, namespace)336 if sid:337 serialized_socket = self.serialize_socket(sid, namespace,338 eio_sid)339 await self.sio.emit('socket_connected', (340 serialized_socket,341 datetime.utcfromtimestamp(t).isoformat() + 'Z',342 ), namespace=self.admin_namespace)343 return await socket.__send_ping()344 345 async def _emit_server_stats(self):346 start_time = time.time()347 namespaces = list(self.sio.handlers.keys())348 namespaces.sort()349 while not self.stop_stats_event.is_set():350 await self.sio.sleep(self.server_stats_interval)351 await self.sio.emit('server_stats', {352 'serverId': self.server_id,353 'hostname': HOSTNAME,354 'pid': PID,355 'uptime': time.time() - start_time,356 'clientsCount': len(self.sio.eio.sockets),357 'pollingClientsCount': len(358 [s for s in self.sio.eio.sockets.values()359 if not s.upgraded]),360 'aggregatedEvents': self.event_buffer.get_and_clear(),361 'namespaces': [{362 'name': nsp,363 'socketsCount': len(self.sio.manager.rooms.get(364 nsp, {None: []}).get(None, []))365 } for nsp in namespaces],366 }, namespace=self.admin_namespace)367 while self.admin_queue:368 event, args = self.admin_queue.pop(0)369 await self.sio.emit(event, args,370 namespace=self.admin_namespace)371 372 def serialize_socket(self, sid, namespace, eio_sid=None):373 if eio_sid is None: # pragma: no cover374 eio_sid = self.sio.manager.eio_sid_from_sid(sid)375 socket = self.sio.eio._get_socket(eio_sid)376 environ = self.sio.environ.get(eio_sid, {})377 tm = self.sio.manager._timestamps[sid] if sid in \378 self.sio.manager._timestamps else 0379 return {380 'id': sid,381 'clientId': eio_sid,382 'transport': 'websocket' if socket.upgraded else 'polling',383 'nsp': namespace,384 'data': {},385 'handshake': {386 'address': environ.get('REMOTE_ADDR', ''),387 'headers': {k[5:].lower(): v for k, v in environ.items()388 if k.startswith('HTTP_')},389 'query': {k: v[0] if len(v) == 1 else v for k, v in parse_qs(390 environ.get('QUERY_STRING', '')).items()},391 'secure': environ.get('wsgi.url_scheme', '') == 'https',392 'url': environ.get('PATH_INFO', ''),393 'issued': tm * 1000,394 'time': datetime.utcfromtimestamp(tm).isoformat() + 'Z'395 if tm else '',396 },397 'rooms': self.sio.manager.get_rooms(sid, namespace),398 }399 