codekingpro/portable-devtools
115k
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 