"""Node backend bridge — runs the existing Next.js Respite backend inside the Space container (userspace Node, no root) and reverse-proxies to it. Mounted under /respite/v2/node/*: /status → is the runtime bootstrapped, is the child alive /version → runs `node --version` as a child process /proxy → reverse-proxy (+ query/body/headers) to the Node backend, streaming both ways (SSE included) """ import os import subprocess import httpx from fastapi import Request from fastapi.responses import JSONResponse import node_runtime def _node_path() -> str | None: try: return node_runtime.ensure_node() except Exception: return None CORS_HEADERS = { "Access-Control-Allow-Origin": "*", "Access-Control-Allow-Methods": "GET, POST, OPTIONS", "Access-Control-Allow-Headers": "Content-Type, Authorization, Accept", "Access-Control-Max-Age": "86400", } def register_routes(fa_app): @fa_app.options("/respite/v2/node/proxy/{path:path}") async def _proxy_preflight(path: str): from starlette.responses import PlainTextResponse return PlainTextResponse("", status_code=204, headers=CORS_HEADERS) @fa_app.get("/respite/v2/node/status") def _status(): return JSONResponse({ "runtime_dir": node_runtime.RUNTIME_DIR, "node_bin_exists": os.path.exists(node_runtime.NODE_BIN), "backend_ready": node_runtime.node_ready(), "backend_dir": node_runtime.BACKEND_DIR, }) @fa_app.get("/respite/v2/node/version") def _version(): node = _node_path() if not node: return JSONResponse({"error": "node bootstrap failed"}, status_code=500) r = subprocess.run([node, "--version"], capture_output=True, text=True, timeout=30) return JSONResponse({"stdout": r.stdout.strip(), "stderr": r.stderr.strip()[:400], "code": r.returncode}) @fa_app.get("/respite/v2/node/proxy/{path:path}") @fa_app.post("/respite/v2/node/proxy/{path:path}") async def _proxy(request: Request, path: str): """Forward method, query string, body and headers to the Node backend on 127.0.0.1:3210. Streaming both ways (SSE included).""" if not node_runtime.node_ready(): start_ok = node_runtime.start_node_backend() if not start_ok: return JSONResponse({"error": "node backend not running"}, status_code=503) url = f"http://127.0.0.1:{node_runtime.NODE_PORT}/{path}" if request.url.query: url += "?" + request.url.query headers = { k: v for k, v in request.headers.items() if k.lower() not in ("host", "connection", "content-length", "accept-encoding") } body = await request.body() try: client = httpx.AsyncClient(timeout=httpx.Timeout(300.0, connect=10.0)) req = client.build_request( request.method, url, headers=headers, content=body or None ) upstream = await client.send(req, stream=True) resp_headers = { k: v for k, v in upstream.headers.items() if k.lower() not in ("content-length", "transfer-encoding", "connection") } from starlette.responses import StreamingResponse async def stream_gen(): try: async for chunk in upstream.aiter_bytes(): yield chunk finally: await upstream.aclose() await client.aclose() resp_headers.update(CORS_HEADERS) return StreamingResponse( stream_gen(), status_code=upstream.status_code, headers=resp_headers ) except Exception as e: return JSONResponse({"error": str(e)}, status_code=502) @fa_app.get("/respite/v2/live-frames") async def _live_frames(q: str = "Write three short paragraphs about rivers."): """TEMP DIAGNOSTIC: open a Live socket, log outputTranscription frame timings/count. Returns only timing metadata — never the API key.""" import asyncio, json, time websockets = None try: import websockets as _ws websockets = _ws except ImportError: return JSONResponse({"error": "websockets not installed"}) key = os.environ.get("GOOGLE_API_KEY") or os.environ.get("GEMINI_API_KEY") if not key: return JSONResponse({"error": "no google key in env"}) url = f"wss://generativelanguage.googleapis.com/ws/google.ai.generativelanguage.v1beta.GenerativeService.BidiGenerateContent?key={key}" frames = [{"t": 0, "note": "route entered"}] t0 = time.monotonic() async def run(): async with websockets.connect(url, max_size=None) as ws: await ws.send(json.dumps({"setup": {"model": "models/gemini-3.8-live"}})) setup = True while True: try: raw = await asyncio.wait_for(ws.recv(), timeout=45) except asyncio.TimeoutError: frames.append({"t": round(time.monotonic() - t0, 2), "err": "timeout"}) break except Exception as e: frames.append({"t": round(time.monotonic() - t0, 2), "err": repr(e)[:200]}) break try: f = json.loads(raw) except Exception: continue if f.get("setupComplete"): await ws.send(json.dumps({"clientContent": {"turns": [{"role": "user", "parts": [{"text": q}]}], "turnComplete": True}})) setup = False continue frames.append({"t": round(time.monotonic() - t0, 2), "keys": sorted(f.keys()), "raw": json.dumps(f)[:400]}) if len(frames) >= 6: break try: await asyncio.wait_for(run(), timeout=50) except Exception as e: return JSONResponse({"error": str(e), "frames": frames}) return JSONResponse({"frames": frames, "count": len(frames), "done": True})