patdev commited on
Commit
75e5bb9
·
verified ·
1 Parent(s): 45676a7

pont: fermeture amont a la deconnexion client, pool httpx 512 + timeout de pool (fuite de flux -> pont muet)

Browse files
Files changed (1) hide show
  1. anthropic_proxy.py +22 -3
anthropic_proxy.py CHANGED
@@ -57,7 +57,16 @@ ALIASES = [a for pair in ((x, f"{x}[1m]") for x in ALIASES) for a in pair]
57
  ALIASES = list(dict.fromkeys(ALIASES)) # dedoublonne en gardant l'ordre
58
 
59
  app = FastAPI(title="anthropic-bridge")
60
- _client = httpx.AsyncClient(base_url=UPSTREAM, timeout=TIMEOUT)
 
 
 
 
 
 
 
 
 
61
 
62
 
63
  # ------------------------------------------------------------------ garde-fou
@@ -467,7 +476,7 @@ def _sse(event: str, data: dict) -> bytes:
467
  return f"event: {event}\ndata: {json.dumps(data, separators=(',', ':'))}\n\n".encode()
468
 
469
 
470
- async def stream_anthropic(payload: dict, req_model: str):
471
  """Traduit le flux OpenAI en evenements Anthropic.
472
 
473
  Le point delicat est l'indexation des blocs : Anthropic numerote chaque bloc
@@ -567,11 +576,21 @@ async def stream_anthropic(payload: dict, req_model: str):
567
  try:
568
  line = await asyncio.wait_for(file.get(), timeout=KEEPALIVE_S)
569
  except asyncio.TimeoutError:
 
 
 
 
 
570
  yield _sse("ping", {"type": "ping"})
571
  last_out = time.monotonic()
572
  continue
573
  if line is None:
574
  break
 
 
 
 
 
575
  if not line.startswith("data: "):
576
  continue
577
  chunk = line[6:].strip()
@@ -739,7 +758,7 @@ async def messages(request: Request):
739
 
740
  if payload.get("stream"):
741
  payload["stream_options"] = {"include_usage": True}
742
- return StreamingResponse(stream_anthropic(payload, req_model),
743
  media_type="text/event-stream")
744
 
745
  r = await _client.post("/v1/chat/completions", json=payload)
 
57
  ALIASES = list(dict.fromkeys(ALIASES)) # dedoublonne en gardant l'ordre
58
 
59
  app = FastAPI(title="anthropic-bridge")
60
+ # Pool EXPLICITE. Mesure 22/08 nuit (5 agents Claude Code, ~150 k de contexte) :
61
+ # 189 flux ouverts pour 100 termines, 45 connexions amont pour 2 clients -- les
62
+ # flux abandonnes cote client (retry, coupure proxy) gardaient leur connexion
63
+ # amont, vLLM generait pour personne, et a 100 connexions (defaut httpx) le pool
64
+ # bloquait TOUT, /health compris : pont muet, processus vivant a 3 % CPU.
65
+ # Le timeout de pool transforme une saturation en 503 rapide au lieu d'un blocage.
66
+ _client = httpx.AsyncClient(
67
+ base_url=UPSTREAM,
68
+ timeout=httpx.Timeout(TIMEOUT, connect=10.0, pool=10.0),
69
+ limits=httpx.Limits(max_connections=512, max_keepalive_connections=64))
70
 
71
 
72
  # ------------------------------------------------------------------ garde-fou
 
476
  return f"event: {event}\ndata: {json.dumps(data, separators=(',', ':'))}\n\n".encode()
477
 
478
 
479
+ async def stream_anthropic(payload: dict, req_model: str, request: Request | None = None):
480
  """Traduit le flux OpenAI en evenements Anthropic.
481
 
482
  Le point delicat est l'indexation des blocs : Anthropic numerote chaque bloc
 
576
  try:
577
  line = await asyncio.wait_for(file.get(), timeout=KEEPALIVE_S)
578
  except asyncio.TimeoutError:
579
+ # Client parti pendant le prefill ? On ferme l'amont (vLLM
580
+ # annule alors la generation) au lieu de pomper dans le vide.
581
+ if request is not None and await request.is_disconnected():
582
+ print("[deconnexion] client parti pendant l'attente, amont ferme", flush=True)
583
+ break
584
  yield _sse("ping", {"type": "ping"})
585
  last_out = time.monotonic()
586
  continue
587
  if line is None:
588
  break
589
+ # Le `yield` ne leve pas toujours a la deconnexion (le proxy Runpod
590
+ # garde parfois sa connexion ouverte) : on verifie explicitement.
591
+ if request is not None and (out_tokens % 32 == 0) and await request.is_disconnected():
592
+ print(f"[deconnexion] client parti apres {out_tokens} jetons, amont ferme", flush=True)
593
+ break
594
  if not line.startswith("data: "):
595
  continue
596
  chunk = line[6:].strip()
 
758
 
759
  if payload.get("stream"):
760
  payload["stream_options"] = {"include_usage": True}
761
+ return StreamingResponse(stream_anthropic(payload, req_model, request),
762
  media_type="text/event-stream")
763
 
764
  r = await _client.post("/v1/chat/completions", json=payload)