Team Ai
Datasetpublic

codekingpro/portable-devtools

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