Respite-API / node_probe.py
Buffy
diag: route-entered marker
8686316
Raw
History Blame Contribute Delete
6.37 kB
"""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 <path> (+ 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})