Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes15kdownloads
socket.py251 linesDownload Raw Back to engineio
1import sys2import time3 4from . import base_socket5from . import exceptions6from . import packet7from . import payload8 9 10class Socket(base_socket.BaseSocket):11    """An Engine.IO socket."""12    def poll(self):13        """Wait for packets to send to the client."""14        queue_empty = self.server.get_queue_empty_exception()15        try:16            packets = [self.queue.get(17                timeout=self.server.ping_interval + self.server.ping_timeout)]18            self.queue.task_done()19        except queue_empty:20            raise exceptions.QueueEmpty()21        if packets == [None]:22            return []23        while True:24            try:25                pkt = self.queue.get(block=False)26                self.queue.task_done()27                if pkt is None:28                    self.queue.put(None)29                    break30                packets.append(pkt)31            except queue_empty:32                break33        return packets34 35    def receive(self, pkt):36        """Receive packet from the client."""37        packet_name = packet.packet_names[pkt.packet_type] \38            if pkt.packet_type < len(packet.packet_names) else 'UNKNOWN'39        self.server.logger.info('%s: Received packet %s data %s',40                                self.sid, packet_name,41                                pkt.data if not isinstance(pkt.data, bytes)42                                else '<binary>')43        if pkt.packet_type == packet.PONG:44            self.schedule_ping()45        elif pkt.packet_type == packet.MESSAGE:46            self.server._trigger_event('message', self.sid, pkt.data,47                                       run_async=self.server.async_handlers)48        elif pkt.packet_type == packet.UPGRADE:49            self.send(packet.Packet(packet.NOOP))50        elif pkt.packet_type == packet.CLOSE:51            self.close(wait=False, abort=True)52        else:53            raise exceptions.UnknownPacketError()54 55    def check_ping_timeout(self):56        """Make sure the client is still responding to pings."""57        if self.closed:58            raise exceptions.SocketIsClosedError()59        if self.last_ping and \60                time.time() - self.last_ping > self.server.ping_timeout:61            self.server.logger.info('%s: Client is gone, closing socket',62                                    self.sid)63            # Passing abort=False here will cause close() to write a64            # CLOSE packet. This has the effect of updating half-open sockets65            # to their correct state of disconnected66            self.close(wait=False, abort=False)67            return False68        return True69 70    def send(self, pkt):71        """Send a packet to the client."""72        if not self.check_ping_timeout():73            return74        else:75            self.queue.put(pkt)76        self.server.logger.info('%s: Sending packet %s data %s',77                                self.sid, packet.packet_names[pkt.packet_type],78                                pkt.data if not isinstance(pkt.data, bytes)79                                else '<binary>')80 81    def handle_get_request(self, environ, start_response):82        """Handle a long-polling GET request from the client."""83        connections = [84            s.strip()85            for s in environ.get('HTTP_CONNECTION', '').lower().split(',')]86        transport = environ.get('HTTP_UPGRADE', '').lower()87        if 'upgrade' in connections and transport in self.upgrade_protocols:88            self.server.logger.info('%s: Received request to upgrade to %s',89                                    self.sid, transport)90            return getattr(self, '_upgrade_' + transport)(environ,91                                                          start_response)92        if self.upgrading or self.upgraded:93            # we are upgrading to WebSocket, do not return any more packets94            # through the polling endpoint95            return [packet.Packet(packet.NOOP)]96        try:97            packets = self.poll()98        except exceptions.QueueEmpty:99            exc = sys.exc_info()100            self.close(wait=False)101            raise exc[1].with_traceback(exc[2])102        return packets103 104    def handle_post_request(self, environ):105        """Handle a long-polling POST request from the client."""106        length = int(environ.get('CONTENT_LENGTH', '0'))107        if length > self.server.max_http_buffer_size:108            raise exceptions.ContentTooLongError()109        else:110            body = environ['wsgi.input'].read(length).decode('utf-8')111            p = payload.Payload(encoded_payload=body)112            for pkt in p.packets:113                self.receive(pkt)114 115    def close(self, wait=True, abort=False):116        """Close the socket connection."""117        if not self.closed and not self.closing:118            self.closing = True119            self.server._trigger_event('disconnect', self.sid, run_async=False)120            if not abort:121                self.send(packet.Packet(packet.CLOSE))122            self.closed = True123            self.queue.put(None)124            if wait:125                self.queue.join()126 127    def schedule_ping(self):128        self.server.start_background_task(self._send_ping)129 130    def _send_ping(self):131        self.last_ping = None132        self.server.sleep(self.server.ping_interval)133        if not self.closing and not self.closed:134            self.last_ping = time.time()135            self.send(packet.Packet(packet.PING))136 137    def _upgrade_websocket(self, environ, start_response):138        """Upgrade the connection from polling to websocket."""139        if self.upgraded:140            raise IOError('Socket has been upgraded already')141        if self.server._async['websocket'] is None:142            # the selected async mode does not support websocket143            return self.server._bad_request()144        ws = self.server._async['websocket'](145            self._websocket_handler, self.server)146        return ws(environ, start_response)147 148    def _websocket_handler(self, ws):149        """Engine.IO handler for websocket transport."""150        def websocket_wait():151            data = ws.wait()152            if data and len(data) > self.server.max_http_buffer_size:153                raise ValueError('packet is too large')154            return data155 156        # try to set a socket timeout matching the configured ping interval157        # and timeout158        for attr in ['_sock', 'socket']:  # pragma: no cover159            if hasattr(ws, attr) and hasattr(getattr(ws, attr), 'settimeout'):160                getattr(ws, attr).settimeout(161                    self.server.ping_interval + self.server.ping_timeout)162 163        if self.connected:164            # the socket was already connected, so this is an upgrade165            self.upgrading = True  # hold packet sends during the upgrade166 167            pkt = websocket_wait()168            decoded_pkt = packet.Packet(encoded_packet=pkt)169            if decoded_pkt.packet_type != packet.PING or \170                    decoded_pkt.data != 'probe':171                self.server.logger.info(172                    '%s: Failed websocket upgrade, no PING packet', self.sid)173                self.upgrading = False174                return []175            ws.send(packet.Packet(packet.PONG, data='probe').encode())176            self.queue.put(packet.Packet(packet.NOOP))  # end poll177 178            pkt = websocket_wait()179            decoded_pkt = packet.Packet(encoded_packet=pkt)180            if decoded_pkt.packet_type != packet.UPGRADE:181                self.upgraded = False182                self.server.logger.info(183                    ('%s: Failed websocket upgrade, expected UPGRADE packet, '184                     'received %s instead.'),185                    self.sid, pkt)186                self.upgrading = False187                return []188            self.upgraded = True189            self.upgrading = False190        else:191            self.connected = True192            self.upgraded = True193 194        # start separate writer thread195        def writer():196            while True:197                packets = None198                try:199                    packets = self.poll()200                except exceptions.QueueEmpty:201                    break202                if not packets:203                    # empty packet list returned -> connection closed204                    break205                try:206                    for pkt in packets:207                        ws.send(pkt.encode())208                except:209                    break210            ws.close()211 212        writer_task = self.server.start_background_task(writer)213 214        self.server.logger.info(215            '%s: Upgrade to websocket successful', self.sid)216 217        while True:218            p = None219            try:220                p = websocket_wait()221            except Exception as e:222                # if the socket is already closed, we can assume this is a223                # downstream error of that224                if not self.closed:  # pragma: no cover225                    self.server.logger.info(226                        '%s: Unexpected error "%s", closing connection',227                        self.sid, str(e))228                break229            if p is None:230                # connection closed by client231                break232            pkt = packet.Packet(encoded_packet=p)233            try:234                self.receive(pkt)235            except exceptions.UnknownPacketError:  # pragma: no cover236                pass237            except exceptions.SocketIsClosedError:  # pragma: no cover238                self.server.logger.info('Receive error -- socket is closed')239                break240            except:  # pragma: no cover241                # if we get an unexpected exception we log the error and exit242                # the connection properly243                self.server.logger.exception('Unknown receive error')244                break245 246        self.queue.put(None)  # unlock the writer task so that it can exit247        writer_task.join()248        self.close(wait=False, abort=True)249 250        return []251 
codekingpro/portable-devtools · Team Ai