codekingpro/portable-devtools
115k
1from datetime import datetime2import functools3import os4import socket5import time6from urllib.parse import parse_qs7from .exceptions import ConnectionRefusedError8 9HOSTNAME = socket.gethostname()10PID = os.getpid()11 12 13class EventBuffer:14 def __init__(self):15 self.buffer = {}16 17 def push(self, type, count=1):18 timestamp = int(time.time()) * 100019 key = '{};{}'.format(timestamp, type)20 if key not in self.buffer:21 self.buffer[key] = {22 'timestamp': timestamp,23 'type': type,24 'count': count,25 }26 else:27 self.buffer[key]['count'] += count28 29 def get_and_clear(self):30 buffer = self.buffer31 self.buffer = {}32 return [value for value in buffer.values()]33 34 35class InstrumentedServer:36 def __init__(self, sio, auth=None, mode='development', read_only=False,37 server_id=None, namespace='/admin', server_stats_interval=2):38 """Instrument the Socket.IO server for monitoring with the `Socket.IO39 Admin UI <https://socket.io/docs/v4/admin-ui/>`_.40 """41 if auth is None:42 raise ValueError('auth must be specified')43 self.sio = sio44 self.auth = auth45 self.admin_namespace = namespace46 self.read_only = read_only47 self.server_id = server_id or (48 self.sio.manager.host_id if hasattr(self.sio.manager, 'host_id')49 else HOSTNAME50 )51 self.mode = mode52 self.server_stats_interval = server_stats_interval53 self.event_buffer = EventBuffer()54 55 # task that emits "server_stats" every 2 seconds56 self.stop_stats_event = None57 self.stats_task = None58 59 # monkey-patch the server to report metrics to the admin UI60 self.instrument()61 62 def instrument(self):63 self.sio.on('connect', self.admin_connect,64 namespace=self.admin_namespace)65 66 if self.mode == 'development':67 if not self.read_only: # pragma: no branch68 self.sio.on('emit', self.admin_emit,69 namespace=self.admin_namespace)70 self.sio.on('join', self.admin_enter_room,71 namespace=self.admin_namespace)72 self.sio.on('leave', self.admin_leave_room,73 namespace=self.admin_namespace)74 self.sio.on('_disconnect', self.admin_disconnect,75 namespace=self.admin_namespace)76 77 # track socket connection times78 self.sio.manager._timestamps = {}79 80 # report socket.io connections81 self.sio.manager.__connect = self.sio.manager.connect82 self.sio.manager.connect = self._connect83 84 # report socket.io disconnection85 self.sio.manager.__disconnect = self.sio.manager.disconnect86 self.sio.manager.disconnect = self._disconnect87 88 # report join rooms89 self.sio.manager.__basic_enter_room = \90 self.sio.manager.basic_enter_room91 self.sio.manager.basic_enter_room = self._basic_enter_room92 93 # report leave rooms94 self.sio.manager.__basic_leave_room = \95 self.sio.manager.basic_leave_room96 self.sio.manager.basic_leave_room = self._basic_leave_room97 98 # report emit events99 self.sio.manager.__emit = self.sio.manager.emit100 self.sio.manager.emit = self._emit101 102 # report receive events103 self.sio.__handle_event_internal = self.sio._handle_event_internal104 self.sio._handle_event_internal = self._handle_event_internal105 106 # report engine.io connections107 self.sio.eio.on('connect', self._handle_eio_connect)108 self.sio.eio.on('disconnect', self._handle_eio_disconnect)109 110 # report polling packets111 from engineio.socket import Socket112 self.sio.eio.__ok = self.sio.eio._ok113 self.sio.eio._ok = self._eio_http_response114 Socket.__handle_post_request = Socket.handle_post_request115 Socket.handle_post_request = functools.partialmethod(116 self.__class__._eio_handle_post_request, self)117 118 # report websocket packets119 Socket.__websocket_handler = Socket._websocket_handler120 Socket._websocket_handler = functools.partialmethod(121 self.__class__._eio_websocket_handler, self)122 123 # report connected sockets with each ping124 if self.mode == 'development':125 Socket.__send_ping = Socket._send_ping126 Socket._send_ping = functools.partialmethod(127 self.__class__._eio_send_ping, self)128 129 def uninstrument(self): # pragma: no cover130 if self.mode == 'development':131 self.sio.manager.connect = self.sio.manager.__connect132 self.sio.manager.disconnect = self.sio.manager.__disconnect133 self.sio.manager.basic_enter_room = \134 self.sio.manager.__basic_enter_room135 self.sio.manager.basic_leave_room = \136 self.sio.manager.__basic_leave_room137 self.sio.manager.emit = self.sio.manager.__emit138 self.sio._handle_event_internal = self.sio.__handle_event_internal139 self.sio.eio._ok = self.sio.eio.__ok140 141 from engineio.socket import Socket142 Socket.handle_post_request = Socket.__handle_post_request143 Socket._websocket_handler = Socket.__websocket_handler144 if self.mode == 'development':145 Socket._send_ping = Socket.__send_ping146 147 def admin_connect(self, sid, environ, client_auth):148 if self.auth:149 authenticated = False150 if isinstance(self.auth, dict):151 authenticated = client_auth == self.auth152 elif isinstance(self.auth, list):153 authenticated = client_auth in self.auth154 else:155 authenticated = self.auth(client_auth)156 if not authenticated:157 raise ConnectionRefusedError('authentication failed')158 159 def config(sid):160 self.sio.sleep(0.1)161 162 # supported features163 features = ['AGGREGATED_EVENTS']164 if not self.read_only:165 features += ['EMIT', 'JOIN', 'LEAVE', 'DISCONNECT', 'MJOIN',166 'MLEAVE', 'MDISCONNECT']167 if self.mode == 'development':168 features.append('ALL_EVENTS')169 self.sio.emit('config', {'supportedFeatures': features},170 to=sid, namespace=self.admin_namespace)171 172 # send current sockets173 if self.mode == 'development':174 all_sockets = []175 for nsp in self.sio.manager.get_namespaces():176 for sid, eio_sid in self.sio.manager.get_participants(177 nsp, None):178 all_sockets.append(179 self.serialize_socket(sid, nsp, eio_sid))180 self.sio.emit('all_sockets', all_sockets, to=sid,181 namespace=self.admin_namespace)182 183 self.sio.start_background_task(config, sid)184 185 def admin_emit(self, _, namespace, room_filter, event, *data):186 self.sio.emit(event, data, to=room_filter, namespace=namespace)187 188 def admin_enter_room(self, _, namespace, room, room_filter=None):189 for sid, _ in self.sio.manager.get_participants(190 namespace, room_filter):191 self.sio.enter_room(sid, room, namespace=namespace)192 193 def admin_leave_room(self, _, namespace, room, room_filter=None):194 for sid, _ in self.sio.manager.get_participants(195 namespace, room_filter):196 self.sio.leave_room(sid, room, namespace=namespace)197 198 def admin_disconnect(self, _, namespace, close, room_filter=None):199 for sid, _ in self.sio.manager.get_participants(200 namespace, room_filter):201 self.sio.disconnect(sid, namespace=namespace)202 203 def shutdown(self):204 if self.stats_task: # pragma: no branch205 self.stop_stats_event.set()206 self.stats_task.join()207 208 def _connect(self, eio_sid, namespace):209 sid = self.sio.manager.__connect(eio_sid, namespace)210 t = time.time()211 self.sio.manager._timestamps[sid] = t212 serialized_socket = self.serialize_socket(sid, namespace, eio_sid)213 self.sio.emit('socket_connected', (214 serialized_socket,215 datetime.utcfromtimestamp(t).isoformat() + 'Z',216 ), namespace=self.admin_namespace)217 return sid218 219 def _disconnect(self, sid, namespace, **kwargs):220 del self.sio.manager._timestamps[sid]221 self.sio.emit('socket_disconnected', (222 namespace,223 sid,224 'N/A',225 datetime.utcnow().isoformat() + 'Z',226 ), namespace=self.admin_namespace)227 return self.sio.manager.__disconnect(sid, namespace, **kwargs)228 229 def _check_for_upgrade(self, eio_sid, sid, namespace): # pragma: no cover230 for _ in range(5):231 self.sio.sleep(5)232 try:233 if self.sio.eio._get_socket(eio_sid).upgraded:234 self.sio.emit('socket_updated', {235 'id': sid,236 'nsp': namespace,237 'transport': 'websocket',238 }, namespace=self.admin_namespace)239 break240 except KeyError:241 pass242 243 def _basic_enter_room(self, sid, namespace, room, eio_sid=None):244 ret = self.sio.manager.__basic_enter_room(sid, namespace, room,245 eio_sid)246 if room:247 self.sio.emit('room_joined', (248 namespace,249 room,250 sid,251 datetime.utcnow().isoformat() + 'Z',252 ), namespace=self.admin_namespace)253 return ret254 255 def _basic_leave_room(self, sid, namespace, room):256 if room:257 self.sio.emit('room_left', (258 namespace,259 room,260 sid,261 datetime.utcnow().isoformat() + 'Z',262 ), namespace=self.admin_namespace)263 return self.sio.manager.__basic_leave_room(sid, namespace, room)264 265 def _emit(self, event, data, namespace, room=None, skip_sid=None,266 callback=None, **kwargs):267 ret = self.sio.manager.__emit(event, data, namespace, room=room,268 skip_sid=skip_sid, callback=callback,269 **kwargs)270 if namespace != self.admin_namespace:271 event_data = [event] + list(data) if isinstance(data, tuple) \272 else [data]273 if not isinstance(skip_sid, list): # pragma: no branch274 skip_sid = [skip_sid]275 for sid, _ in self.sio.manager.get_participants(namespace, room):276 if sid not in skip_sid:277 self.sio.emit('event_sent', (278 namespace,279 sid,280 event_data,281 datetime.utcnow().isoformat() + 'Z',282 ), namespace=self.admin_namespace)283 return ret284 285 def _handle_event_internal(self, server, sid, eio_sid, data, namespace,286 id):287 ret = self.sio.__handle_event_internal(server, sid, eio_sid, data,288 namespace, id)289 self.sio.emit('event_received', (290 namespace,291 sid,292 data,293 datetime.utcnow().isoformat() + 'Z',294 ), namespace=self.admin_namespace)295 return ret296 297 def _handle_eio_connect(self, eio_sid, environ):298 if self.stop_stats_event is None:299 self.stop_stats_event = self.sio.eio.create_event()300 self.stats_task = self.sio.start_background_task(301 self._emit_server_stats)302 303 self.event_buffer.push('rawConnection')304 return self.sio._handle_eio_connect(eio_sid, environ)305 306 def _handle_eio_disconnect(self, eio_sid):307 self.event_buffer.push('rawDisconnection')308 return self.sio._handle_eio_disconnect(eio_sid)309 310 def _eio_http_response(self, packets=None, headers=None, jsonp_index=None):311 ret = self.sio.eio.__ok(packets=packets, headers=headers,312 jsonp_index=jsonp_index)313 self.event_buffer.push('packetsOut')314 self.event_buffer.push('bytesOut', len(ret['response']))315 return ret316 317 def _eio_handle_post_request(socket, self, environ):318 ret = socket.__handle_post_request(environ)319 self.event_buffer.push('packetsIn')320 self.event_buffer.push(321 'bytesIn', int(environ.get('CONTENT_LENGTH', 0)))322 return ret323 324 def _eio_websocket_handler(socket, self, ws):325 def _send(ws, data, *args, **kwargs):326 self.event_buffer.push('packetsOut')327 self.event_buffer.push('bytesOut', len(data))328 return ws.__send(data, *args, **kwargs)329 330 def _wait(ws):331 ret = ws.__wait()332 self.event_buffer.push('packetsIn')333 self.event_buffer.push('bytesIn', len(ret or ''))334 return ret335 336 ws.__send = ws.send337 ws.send = functools.partial(_send, ws)338 ws.__wait = ws.wait339 ws.wait = functools.partial(_wait, ws)340 return socket.__websocket_handler(ws)341 342 def _eio_send_ping(socket, self): # pragma: no cover343 eio_sid = socket.sid344 t = time.time()345 for namespace in self.sio.manager.get_namespaces():346 sid = self.sio.manager.sid_from_eio_sid(eio_sid, namespace)347 if sid:348 serialized_socket = self.serialize_socket(sid, namespace,349 eio_sid)350 self.sio.emit('socket_connected', (351 serialized_socket,352 datetime.utcfromtimestamp(t).isoformat() + 'Z',353 ), namespace=self.admin_namespace)354 return socket.__send_ping()355 356 def _emit_server_stats(self):357 start_time = time.time()358 namespaces = list(self.sio.handlers.keys())359 namespaces.sort()360 while not self.stop_stats_event.is_set():361 self.sio.sleep(self.server_stats_interval)362 self.sio.emit('server_stats', {363 'serverId': self.server_id,364 'hostname': HOSTNAME,365 'pid': PID,366 'uptime': time.time() - start_time,367 'clientsCount': len(self.sio.eio.sockets),368 'pollingClientsCount': len(369 [s for s in self.sio.eio.sockets.values()370 if not s.upgraded]),371 'aggregatedEvents': self.event_buffer.get_and_clear(),372 'namespaces': [{373 'name': nsp,374 'socketsCount': len(self.sio.manager.rooms.get(375 nsp, {None: []}).get(None, []))376 } for nsp in namespaces],377 }, namespace=self.admin_namespace)378 379 def serialize_socket(self, sid, namespace, eio_sid=None):380 if eio_sid is None: # pragma: no cover381 eio_sid = self.sio.manager.eio_sid_from_sid(sid)382 socket = self.sio.eio._get_socket(eio_sid)383 environ = self.sio.environ.get(eio_sid, {})384 tm = self.sio.manager._timestamps[sid] if sid in \385 self.sio.manager._timestamps else 0386 return {387 'id': sid,388 'clientId': eio_sid,389 'transport': 'websocket' if socket.upgraded else 'polling',390 'nsp': namespace,391 'data': {},392 'handshake': {393 'address': environ.get('REMOTE_ADDR', ''),394 'headers': {k[5:].lower(): v for k, v in environ.items()395 if k.startswith('HTTP_')},396 'query': {k: v[0] if len(v) == 1 else v for k, v in parse_qs(397 environ.get('QUERY_STRING', '')).items()},398 'secure': environ.get('wsgi.url_scheme', '') == 'https',399 'url': environ.get('PATH_INFO', ''),400 'issued': tm * 1000,401 'time': datetime.utcfromtimestamp(tm).isoformat() + 'Z'402 if tm else '',403 },404 'rooms': self.sio.manager.get_rooms(sid, namespace),405 }406 