codekingpro/portable-devtools
114k
1import asyncio2import sys3from urllib.parse import urlsplit4from .. import exceptions5 6import tornado.web7import tornado.websocket8 9 10def get_tornado_handler(engineio_server):11 class Handler(tornado.websocket.WebSocketHandler): # pragma: no cover12 def __init__(self, *args, **kwargs):13 super().__init__(*args, **kwargs)14 if isinstance(engineio_server.cors_allowed_origins, str):15 if engineio_server.cors_allowed_origins == '*':16 self.allowed_origins = None17 else:18 self.allowed_origins = [19 engineio_server.cors_allowed_origins]20 else:21 self.allowed_origins = engineio_server.cors_allowed_origins22 self.receive_queue = asyncio.Queue()23 24 async def get(self, *args, **kwargs):25 if self.request.headers.get('Upgrade', '').lower() == 'websocket':26 ret = super().get(*args, **kwargs)27 if asyncio.iscoroutine(ret):28 await ret29 else:30 await engineio_server.handle_request(self)31 32 async def open(self, *args, **kwargs):33 # this is the handler for the websocket request34 asyncio.ensure_future(engineio_server.handle_request(self))35 36 async def post(self, *args, **kwargs):37 await engineio_server.handle_request(self)38 39 async def options(self, *args, **kwargs):40 await engineio_server.handle_request(self)41 42 async def on_message(self, message):43 await self.receive_queue.put(message)44 45 async def get_next_message(self):46 return await self.receive_queue.get()47 48 def on_close(self):49 self.receive_queue.put_nowait(None)50 51 def check_origin(self, origin):52 if self.allowed_origins is None or origin in self.allowed_origins:53 return True54 return super().check_origin(origin)55 56 def get_compression_options(self):57 # enable compression58 return {}59 60 return Handler61 62 63def translate_request(handler):64 """This function takes the arguments passed to the request handler and65 uses them to generate a WSGI compatible environ dictionary.66 """67 class AwaitablePayload(object):68 def __init__(self, payload):69 self.payload = payload or b''70 71 async def read(self, length=None):72 if length is None:73 r = self.payload74 self.payload = b''75 else:76 r = self.payload[:length]77 self.payload = self.payload[length:]78 return r79 80 payload = handler.request.body81 82 uri_parts = urlsplit(handler.request.path)83 full_uri = handler.request.path84 if handler.request.query: # pragma: no cover85 full_uri += '?' + handler.request.query86 environ = {87 'wsgi.input': AwaitablePayload(payload),88 'wsgi.errors': sys.stderr,89 'wsgi.version': (1, 0),90 'wsgi.async': True,91 'wsgi.multithread': False,92 'wsgi.multiprocess': False,93 'wsgi.run_once': False,94 'SERVER_SOFTWARE': 'aiohttp',95 'REQUEST_METHOD': handler.request.method,96 'QUERY_STRING': handler.request.query or '',97 'RAW_URI': full_uri,98 'SERVER_PROTOCOL': 'HTTP/%s' % handler.request.version,99 'REMOTE_ADDR': '127.0.0.1',100 'REMOTE_PORT': '0',101 'SERVER_NAME': 'aiohttp',102 'SERVER_PORT': '0',103 'tornado.handler': handler104 }105 106 for hdr_name, hdr_value in handler.request.headers.items():107 hdr_name = hdr_name.upper()108 if hdr_name == 'CONTENT-TYPE':109 environ['CONTENT_TYPE'] = hdr_value110 continue111 elif hdr_name == 'CONTENT-LENGTH':112 environ['CONTENT_LENGTH'] = hdr_value113 continue114 115 key = 'HTTP_%s' % hdr_name.replace('-', '_')116 environ[key] = hdr_value117 118 environ['wsgi.url_scheme'] = environ.get('HTTP_X_FORWARDED_PROTO', 'http')119 120 path_info = uri_parts.path121 122 environ['PATH_INFO'] = path_info123 environ['SCRIPT_NAME'] = ''124 125 return environ126 127 128def make_response(status, headers, payload, environ):129 """This function generates an appropriate response object for this async130 mode.131 """132 tornado_handler = environ['tornado.handler']133 try:134 tornado_handler.set_status(int(status.split()[0]))135 except RuntimeError: # pragma: no cover136 # for websocket connections Tornado does not accept a response, since137 # it already emitted the 101 status code138 return139 for header, value in headers:140 tornado_handler.set_header(header, value)141 tornado_handler.write(payload)142 tornado_handler.finish()143 144 145class WebSocket(object): # pragma: no cover146 """147 This wrapper class provides a tornado WebSocket interface that is148 somewhat compatible with eventlet's implementation.149 """150 def __init__(self, handler, server):151 self.handler = handler152 self.tornado_handler = None153 154 async def __call__(self, environ):155 self.tornado_handler = environ['tornado.handler']156 self.environ = environ157 await self.handler(self)158 159 async def close(self):160 self.tornado_handler.close()161 162 async def send(self, message):163 try:164 self.tornado_handler.write_message(165 message, binary=isinstance(message, bytes))166 except tornado.websocket.WebSocketClosedError:167 raise exceptions.EngineIOError()168 169 async def wait(self):170 msg = await self.tornado_handler.get_next_message()171 if not isinstance(msg, bytes) and \172 not isinstance(msg, str):173 raise IOError()174 return msg175 176 177_async = {178 'asyncio': True,179 'translate_request': translate_request,180 'make_response': make_response,181 'websocket': WebSocket,182}183 