Spaces:
Paused
Paused
| import asyncio | |
| from fastapi import FastAPI, WebSocket | |
| from fastapi.responses import HTMLResponse | |
| import logging | |
| logging.basicConfig(level=logging.INFO) | |
| logger = logging.getLogger(__name__) | |
| app = FastAPI() | |
| def read_root(): | |
| return HTMLResponse("<h1>Redis WebSocket Proxy Running</h1>") | |
| async def websocket_endpoint(websocket: WebSocket): | |
| await websocket.accept() | |
| logger.info("WebSocket connection accepted.") | |
| try: | |
| # Connect to local Redis TCP port | |
| reader, writer = await asyncio.open_connection("127.0.0.1", 6379) | |
| logger.info("Connected to local Redis.") | |
| except Exception as e: | |
| logger.error(f"Failed to connect to local Redis: {e}") | |
| await websocket.close() | |
| return | |
| async def pipe_ws_to_tcp(): | |
| try: | |
| while True: | |
| data = await websocket.receive_bytes() | |
| writer.write(data) | |
| await writer.drain() | |
| except Exception as e: | |
| logger.info(f"ws_to_tcp closed: {e}") | |
| async def pipe_tcp_to_ws(): | |
| try: | |
| while True: | |
| data = await reader.read(4096) | |
| if not data: | |
| break | |
| await websocket.send_bytes(data) | |
| except Exception as e: | |
| logger.info(f"tcp_to_ws closed: {e}") | |
| try: | |
| await asyncio.gather( | |
| pipe_ws_to_tcp(), | |
| pipe_tcp_to_ws(), | |
| return_exceptions=True | |
| ) | |
| finally: | |
| writer.close() | |
| try: | |
| await websocket.close() | |
| except: | |
| pass | |
| logger.info("Connection closed.") | |