Spaces:
Running on Zero
Running on Zero
File size: 6,365 Bytes
6da16aa ba9ed6e 6da16aa ba9ed6e bb11950 6da16aa ba9ed6e 60e15d5 ba9ed6e 60e15d5 ba9ed6e bb11950 abb3d7d ba9ed6e abb3d7d ba9ed6e abb3d7d 60e15d5 abb3d7d ba9ed6e 79c96e0 8686316 79c96e0 271301b 79c96e0 f9312e5 79c96e0 8686316 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 | """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})
|