patdev commited on
Commit
4f8a161
·
verified ·
1 Parent(s): 90d66b7

pont v73 : guerison du bassin elargie (HTTPError/RuntimeError) + flux SSE gueri

Browse files
Files changed (1) hide show
  1. anthropic_proxy.py +42 -2
anthropic_proxy.py CHANGED
@@ -154,6 +154,24 @@ async def _guerir_bassin(request: Request, call_next):
154
  "error": {"type": "overloaded_error",
155
  "message": "upstream pool exhausted, "
156
  "connection pool rebuilt; retry"}})
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
157
 
158
  # Pool EXPLICITE. Mesure 22/08 nuit (5 agents Claude Code, ~150 k de contexte) :
159
  # 189 flux ouverts pour 100 termines, 45 connexions amont pour 2 clients -- les
@@ -614,6 +632,26 @@ def _sse(event: str, data: dict) -> bytes:
614
  return f"event: {event}\ndata: {json.dumps(data, separators=(',', ':'))}\n\n".encode()
615
 
616
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
617
  async def stream_anthropic(payload: dict, req_model: str, request: Request | None = None):
618
  """Traduit le flux OpenAI en evenements Anthropic.
619
 
@@ -944,7 +982,8 @@ async def messages(request: Request):
944
  async for morceau in stream_anthropic(payload, req_model, request):
945
  yield morceau
946
 
947
- return StreamingResponse(flux_avec_swap(), media_type="text/event-stream")
 
948
  payload["model"] = await _assurer_modele(req_model)
949
 
950
  r = await _client.post("/v1/chat/completions", json=payload)
@@ -1086,7 +1125,8 @@ async def oai_chat(request: Request):
1086
  async with _client.stream("POST", "/v1/chat/completions", json=body) as r:
1087
  async for chunk in r.aiter_raw():
1088
  yield chunk
1089
- return StreamingResponse(gen(), media_type="text/event-stream")
 
1090
  r = await _client.post("/v1/chat/completions", json=body)
1091
  return JSONResponse(status_code=r.status_code, content=r.json())
1092
 
 
154
  "error": {"type": "overloaded_error",
155
  "message": "upstream pool exhausted, "
156
  "connection pool rebuilt; retry"}})
157
+ except (httpx.HTTPError, RuntimeError) as e:
158
+ # v73 : mesure du 29/08 -- la saturation ne se presente pas toujours en
159
+ # PoolTimeout. Un RuntimeError("client has been closed") apres une
160
+ # reconstruction, ou un ReadError sur une connexion recyclee, sortait en
161
+ # 500 brut sans jamais declencher la guerison. Meme remede, et le type
162
+ # REEL est journalise pour qu'on ne rediagnostique plus a l'aveugle.
163
+ import traceback
164
+ print(f"[bassin?] {type(e).__name__}: {e} sur {request.url.path}",
165
+ flush=True)
166
+ traceback.print_exc()
167
+ await _reconstruire_bassin(generation)
168
+ return JSONResponse(
169
+ status_code=503,
170
+ content={"type": "error",
171
+ "error": {"type": "overloaded_error",
172
+ "message": "upstream relay failed "
173
+ f"({type(e).__name__}), pool "
174
+ "rebuilt; retry"}})
175
 
176
  # Pool EXPLICITE. Mesure 22/08 nuit (5 agents Claude Code, ~150 k de contexte) :
177
  # 189 flux ouverts pour 100 termines, 45 connexions amont pour 2 clients -- les
 
632
  return f"event: {event}\ndata: {json.dumps(data, separators=(',', ':'))}\n\n".encode()
633
 
634
 
635
+ async def _flux_gueri(gen):
636
+ """v73 : les exceptions nees dans un generateur SSE ne traversent pas le
637
+ middleware _guerir_bassin. On les attrape ici : evenement `error` au client
638
+ (sa logique de reprise lit le libelle) et reconstruction du bassin."""
639
+ generation = _bassin_generation
640
+ try:
641
+ async for morceau in gen:
642
+ yield morceau
643
+ except (httpx.HTTPError, RuntimeError) as e:
644
+ import traceback
645
+ print(f"[bassin?] flux: {type(e).__name__}: {e}", flush=True)
646
+ traceback.print_exc()
647
+ asyncio.ensure_future(_reconstruire_bassin(generation))
648
+ yield _sse("error", {
649
+ "type": "error",
650
+ "error": {"type": "overloaded_error",
651
+ "message": f"upstream stream failed "
652
+ f"({type(e).__name__}); retry"}})
653
+
654
+
655
  async def stream_anthropic(payload: dict, req_model: str, request: Request | None = None):
656
  """Traduit le flux OpenAI en evenements Anthropic.
657
 
 
982
  async for morceau in stream_anthropic(payload, req_model, request):
983
  yield morceau
984
 
985
+ return StreamingResponse(_flux_gueri(flux_avec_swap()),
986
+ media_type="text/event-stream")
987
  payload["model"] = await _assurer_modele(req_model)
988
 
989
  r = await _client.post("/v1/chat/completions", json=payload)
 
1125
  async with _client.stream("POST", "/v1/chat/completions", json=body) as r:
1126
  async for chunk in r.aiter_raw():
1127
  yield chunk
1128
+ return StreamingResponse(_flux_gueri(gen()),
1129
+ media_type="text/event-stream")
1130
  r = await _client.post("/v1/chat/completions", json=body)
1131
  return JSONResponse(status_code=r.status_code, content=r.json())
1132