| import asyncio |
| import sys |
| import logging |
| from aiohttp import web, WSMsgType |
|
|
| logging.basicConfig(level=logging.INFO, stream=sys.stderr, |
| format='%(asctime)s %(levelname)s %(message)s') |
| log = logging.getLogger("proxy") |
|
|
|
|
| async def pipe(reader, writer): |
| try: |
| while True: |
| data = await reader.read(65536) |
| if not data: |
| break |
| writer.write(data) |
| await writer.drain() |
| except Exception: |
| pass |
|
|
|
|
| async def handle_ws(request): |
| ws = web.WebSocketResponse() |
| await ws.prepare(request) |
| log.info("WebSocket opened") |
|
|
| try: |
| msg = await ws.receive() |
| target = msg.data |
| log.info(f"CONNECT {target}") |
| host, port = target.rsplit(":", 1) |
| port = int(port) |
|
|
| reader, writer = await asyncio.wait_for( |
| asyncio.open_connection(host, port), timeout=10 |
| ) |
| log.info(f"CONNECTED to {host}:{port}") |
|
|
| async def ws_to_tcp(): |
| async for msg in ws: |
| if msg.type == WSMsgType.BINARY: |
| writer.write(msg.data) |
| await writer.drain() |
|
|
| async def tcp_to_ws(): |
| while True: |
| data = await reader.read(65536) |
| if not data: |
| break |
| await ws.send_bytes(data) |
|
|
| done, pending = await asyncio.wait( |
| [asyncio.create_task(ws_to_tcp()), asyncio.create_task(tcp_to_ws())], |
| return_when=asyncio.FIRST_COMPLETED, |
| ) |
| for t in pending: |
| t.cancel() |
| writer.close() |
| log.info(f"DISCONNECT from {host}:{port}") |
|
|
| except Exception as e: |
| log.error(f"ERROR: {type(e).__name__}: {e}") |
|
|
| return ws |
|
|
|
|
| async def handle_index(request): |
| return web.Response(text="US Proxy OK") |
|
|
|
|
| app = web.Application() |
| app.router.add_get("/", handle_index) |
| app.router.add_get("/ws", handle_ws) |
| app.router.add_get("/", handle_ws) |
|
|
| if __name__ == "__main__": |
| log.info("Starting on 0.0.0.0:7860") |
| web.run_app(app, host="0.0.0.0", port=7860, print=None) |
|
|