diff options
Diffstat (limited to 'ttun_server/endpoints.py')
| -rw-r--r-- | ttun_server/endpoints.py | 57 |
1 files changed, 5 insertions, 52 deletions
diff --git a/ttun_server/endpoints.py b/ttun_server/endpoints.py index 22dcb6d..51c17aa 100644 --- a/ttun_server/endpoints.py +++ b/ttun_server/endpoints.py | |||
| @@ -1,54 +1,7 @@ | |||
| 1 | import logging | 1 | from fastapi import FastAPI |
| 2 | from base64 import b64decode, b64encode | ||
| 3 | from uuid import uuid4 | ||
| 4 | 2 | ||
| 5 | from starlette.background import BackgroundTask | 3 | base_endpoints = FastAPI() |
| 6 | from starlette.requests import Request | ||
| 7 | from starlette.responses import Response | ||
| 8 | 4 | ||
| 9 | from ttun_server.proxy_queue import ProxyQueue | 5 | @base_endpoints.get('/health/') |
| 10 | from ttun_server.types import HttpRequestData, HttpMessageType, HttpMessage | 6 | async def health(): |
| 11 | 7 | return 'OK' | |
| 12 | logger = logging.getLogger(__name__) | ||
| 13 | |||
| 14 | |||
| 15 | async def proxy(request: Request) -> Response: | ||
| 16 | [subdomain, *_] = request.headers['host'].split('.') | ||
| 17 | identifier = str(uuid4()) | ||
| 18 | response_queue = await ProxyQueue.create_for_identifier(identifier) | ||
| 19 | |||
| 20 | try: | ||
| 21 | request_queue = await ProxyQueue.get_for_identifier(subdomain) | ||
| 22 | |||
| 23 | logger.debug('PROXY %s%s ', subdomain, request.url) | ||
| 24 | await request_queue.enqueue( | ||
| 25 | HttpMessage( | ||
| 26 | type=HttpMessageType.request.value, | ||
| 27 | identifier=identifier, | ||
| 28 | payload=HttpRequestData( | ||
| 29 | method=request.method, | ||
| 30 | path=str(request.url).replace(str(request.base_url), '/'), | ||
| 31 | headers=list(request.headers.items()), | ||
| 32 | body=b64encode(await request.body()).decode() | ||
| 33 | ) | ||
| 34 | ) | ||
| 35 | ) | ||
| 36 | |||
| 37 | _response = await response_queue.dequeue() | ||
| 38 | payload = _response['payload'] | ||
| 39 | return Response( | ||
| 40 | status_code=payload['status'], | ||
| 41 | headers=dict(payload['headers']), | ||
| 42 | content=b64decode(payload['body'].encode()), | ||
| 43 | background=BackgroundTask(response_queue.delete) | ||
| 44 | ) | ||
| 45 | except AssertionError: | ||
| 46 | return Response( | ||
| 47 | content='Not Found', | ||
| 48 | status_code=404, | ||
| 49 | background=BackgroundTask(response_queue.delete) | ||
| 50 | ) | ||
| 51 | |||
| 52 | |||
| 53 | async def health(_: Request) -> Response: | ||
| 54 | return Response(content='OK', status_code=200) | ||
