Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes15kdownloads
admin.py406 linesDownload Raw Back to socketio
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 
codekingpro/portable-devtools · Team Ai