# retoor import asyncio import logging import httpx import websockets from fastapi import APIRouter, Request, WebSocket from starlette.responses import Response from devplacepy.services.containers import api, store from devplacepy.utils import not_found from devplacepy.services.audit import record as audit logger = logging.getLogger(__name__) router = APIRouter() HOP_HEADERS = { "connection", "keep-alive", "proxy-authenticate", "proxy-authorization", "te", "trailers", "transfer-encoding", "upgrade", "host", "content-length", "content-encoding", } METHODS = ["GET", "POST", "PUT", "DELETE", "PATCH", "OPTIONS", "HEAD"] def _forward_headers(request: Request, prefix: str) -> dict: headers = {k: v for k, v in request.headers.items() if k.lower() not in HOP_HEADERS} headers["X-Forwarded-Prefix"] = prefix headers["X-Script-Name"] = prefix headers["X-Forwarded-Host"] = request.headers.get( "host", request.url.hostname or "" ) headers["X-Forwarded-Proto"] = request.headers.get( "x-forwarded-proto", request.url.scheme ) headers["Accept-Encoding"] = "identity" return headers def _inject_base(body: bytes, prefix: str) -> bytes: lowered = body.lower() if b"'.encode() head = lowered.find(b"", head) if head != -1 else lowered.find(b">", lowered.find(b" str: if value.startswith("/") and not value.startswith("//"): return prefix + value return value def _resolve(slug: str): instance = store.find_instance_by_ingress(slug) if not instance or instance.get("status") != store.ST_RUNNING: return None, None, None host, port = api.proxy_target(instance) return instance, host, port @router.api_route("/{slug}", methods=METHODS) @router.api_route("/{slug}/{path:path}", methods=METHODS) async def proxy_http(request: Request, slug: str, path: str = ""): instance, host, port = _resolve(slug) if instance is None: raise not_found("No running container is published at this address") if not host or not port: return Response( "the published container has no reachable port", status_code=502 ) audit.record( request, "proxy.access", target_type="instance", target_uid=instance["uid"], target_label=instance.get("name"), metadata={"slug": slug, "path": path, "method": request.method}, summary=f"request proxied to instance {instance.get('name')} via ingress {slug}", links=[audit.instance(instance["uid"], instance.get("name"))], ) prefix = f"/p/{slug}" url = f"http://{host}:{port}/{path}" headers = _forward_headers(request, prefix) body = await request.body() try: async with httpx.AsyncClient(timeout=60.0, follow_redirects=False) as client: upstream = await client.request( request.method, url, params=request.query_params, headers=headers, content=body, ) except httpx.HTTPError as exc: return Response(f"upstream error: {exc}", status_code=502) out_headers = { k: v for k, v in upstream.headers.items() if k.lower() not in HOP_HEADERS and k.lower() != "set-cookie" } if "location" in out_headers: out_headers["location"] = _rewrite_location(out_headers["location"], prefix) content_type = upstream.headers.get("content-type", "") content = upstream.content if "text/html" in content_type.lower(): content = _inject_base(content, prefix) response = Response( content=content, status_code=upstream.status_code, headers=out_headers, media_type=content_type or None, ) for cookie in upstream.headers.get_list("set-cookie"): response.headers.append("set-cookie", cookie) return response @router.websocket("/{slug}") @router.websocket("/{slug}/{path:path}") async def proxy_ws(websocket: WebSocket, slug: str, path: str = ""): instance, host, port = _resolve(slug) if instance is None or not host or not port: await websocket.close(code=1011) return upstream_url = f"ws://{host}:{port}/{path}" if websocket.url.query: upstream_url += f"?{websocket.url.query}" await websocket.accept() audit.record( websocket, "proxy.access", user=None, target_type="instance", target_uid=instance["uid"], target_label=instance.get("name"), metadata={"slug": slug, "path": path, "protocol": "websocket"}, summary=f"websocket proxied to instance {instance.get('name')} via ingress {slug}", links=[audit.instance(instance["uid"], instance.get("name"))], ) try: async with websockets.connect( upstream_url, open_timeout=10, max_size=None ) as upstream: await _pump(websocket, upstream) except Exception as exc: logger.debug("ws proxy %s failed: %s", slug, exc) try: await websocket.close(code=1011) except Exception: pass async def _pump(client_ws: WebSocket, upstream) -> None: async def client_to_upstream(): try: while True: message = await client_ws.receive() if message["type"] == "websocket.disconnect": break if message.get("text") is not None: await upstream.send(message["text"]) elif message.get("bytes") is not None: await upstream.send(message["bytes"]) except Exception: pass finally: await upstream.close() async def upstream_to_client(): try: async for message in upstream: if isinstance(message, (bytes, bytearray)): await client_ws.send_bytes(bytes(message)) else: await client_ws.send_text(message) except Exception: pass finally: try: await client_ws.close() except Exception: pass await asyncio.gather(client_to_upstream(), upstream_to_client())