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})