Buffy commited on
Commit
abb3d7d
·
1 Parent(s): 71ea0d0

Proxy: forward method/query/body/headers, stream SSE both ways

Browse files
Files changed (1) hide show
  1. node_probe.py +32 -8
node_probe.py CHANGED
@@ -81,19 +81,43 @@ def register_routes(fa_app):
81
 
82
  @fa_app.get("/respite/v2/node/proxy/{path:path}")
83
  @fa_app.post("/respite/v2/node/proxy/{path:path}")
84
- def _proxy(path: str):
85
- """Forward to the Node backend on 127.0.0.1:3210 (streaming)."""
 
86
  if not node_runtime.node_ready():
87
  start_ok = node_runtime.start_node_backend()
88
  if not start_ok:
89
  return JSONResponse({"error": "node backend not running"}, status_code=503)
90
  url = f"http://127.0.0.1:{node_runtime.NODE_PORT}/{path}"
 
 
 
 
 
 
 
91
  try:
92
- with httpx.stream("GET", url, timeout=30.0) as r:
93
- headers = dict(r.headers)
94
- body = b""
95
- for chunk in r.iter_bytes():
96
- body += chunk
97
- return Response(content=body, status_code=r.status_code, headers=headers)
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
98
  except Exception as e:
99
  return JSONResponse({"error": str(e)}, status_code=502)
 
81
 
82
  @fa_app.get("/respite/v2/node/proxy/{path:path}")
83
  @fa_app.post("/respite/v2/node/proxy/{path:path}")
84
+ async def _proxy(request, path: str):
85
+ """Forward method, query string, body and headers to the Node backend
86
+ on 127.0.0.1:3210. Streaming both ways (SSE included)."""
87
  if not node_runtime.node_ready():
88
  start_ok = node_runtime.start_node_backend()
89
  if not start_ok:
90
  return JSONResponse({"error": "node backend not running"}, status_code=503)
91
  url = f"http://127.0.0.1:{node_runtime.NODE_PORT}/{path}"
92
+ if request.url.query:
93
+ url += "?" + request.url.query
94
+ headers = {
95
+ k: v for k, v in request.headers.items()
96
+ if k.lower() not in ("host", "connection", "content-length", "accept-encoding")
97
+ }
98
+ body = await request.body()
99
  try:
100
+ client = httpx.AsyncClient(timeout=httpx.Timeout(300.0, connect=10.0))
101
+ req = client.build_request(
102
+ request.method, url, headers=headers, content=body or None
103
+ )
104
+ upstream = await client.send(req, stream=True)
105
+ resp_headers = {
106
+ k: v for k, v in upstream.headers.items()
107
+ if k.lower() not in ("content-length", "transfer-encoding", "connection")
108
+ }
109
+ from starlette.responses import StreamingResponse
110
+
111
+ async def stream_gen():
112
+ try:
113
+ async for chunk in upstream.aiter_bytes():
114
+ yield chunk
115
+ finally:
116
+ await upstream.aclose()
117
+ await client.aclose()
118
+
119
+ return StreamingResponse(
120
+ stream_gen(), status_code=upstream.status_code, headers=resp_headers
121
+ )
122
  except Exception as e:
123
  return JSONResponse({"error": str(e)}, status_code=502)