patdev commited on
Commit
d26237d
·
verified ·
1 Parent(s): e0bb44a

keepalive pendant le preremplissage : le flux ne peut plus rester muet

Browse files
Files changed (1) hide show
  1. anthropic_proxy.py +38 -1
anthropic_proxy.py CHANGED
@@ -11,6 +11,7 @@ UPSTREAM.
11
 
12
  from __future__ import annotations
13
 
 
14
  import json
15
  import os
16
  import re
@@ -473,7 +474,41 @@ async def stream_anthropic(payload: dict, req_model: str):
473
  # comme je le faisais, empechait donc cette reprise.
474
  yield _sse("error", {"type": "error", "error": _err_of(body, r.status_code)})
475
  return
476
- async for line in r.aiter_lines():
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
477
  if not line.startswith("data: "):
478
  continue
479
  chunk = line[6:].strip()
@@ -547,6 +582,8 @@ async def stream_anthropic(payload: dict, req_model: str):
547
 
548
  if ch.get("finish_reason"):
549
  stop_reason = _STOP.get(ch["finish_reason"], "end_turn")
 
 
550
 
551
  if think_open:
552
  yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx})
 
11
 
12
  from __future__ import annotations
13
 
14
+ import asyncio
15
  import json
16
  import os
17
  import re
 
474
  # comme je le faisais, empechait donc cette reprise.
475
  yield _sse("error", {"type": "error", "error": _err_of(body, r.status_code)})
476
  return
477
+ # Le `ping` periodique ne suffisait pas : il ne se declenchait qu'a
478
+ # l'interieur de cette boucle, donc uniquement quand vLLM emettait deja
479
+ # quelque chose. Or pendant le PREREMPLISSAGE vLLM n'envoie rien -- et
480
+ # un prompt de plusieurs centaines de milliers de jetons met plusieurs
481
+ # minutes a etre calcule. Le flux restait muet, et le proxy Runpod
482
+ # coupait vers 125 s.
483
+ #
484
+ # MESURE : une aiguille a ~700 k jetons echouait a 126,5 s avec zero
485
+ # jeton, alors que la meme a ~358 k passait en 94,4 s. Le plafond
486
+ # n'etait pas le modele mais le silence : aucune requete au-dela de
487
+ # ~500 k ne pouvait aboutir, quel que soit le contexte annonce.
488
+ #
489
+ # On decouple donc la lecture amont de l'emission aval : une tache
490
+ # pompe les lignes dans une file, et l'absence de ligne pendant
491
+ # KEEPALIVE_S produit un `ping` au lieu d'un silence.
492
+ file: asyncio.Queue = asyncio.Queue()
493
+
494
+ async def _pompe() -> None:
495
+ try:
496
+ async for _l in r.aiter_lines():
497
+ await file.put(_l)
498
+ finally:
499
+ await file.put(None)
500
+
501
+ _tache = asyncio.create_task(_pompe())
502
+ try:
503
+ while True:
504
+ try:
505
+ line = await asyncio.wait_for(file.get(), timeout=KEEPALIVE_S)
506
+ except asyncio.TimeoutError:
507
+ yield _sse("ping", {"type": "ping"})
508
+ last_out = time.monotonic()
509
+ continue
510
+ if line is None:
511
+ break
512
  if not line.startswith("data: "):
513
  continue
514
  chunk = line[6:].strip()
 
582
 
583
  if ch.get("finish_reason"):
584
  stop_reason = _STOP.get(ch["finish_reason"], "end_turn")
585
+ finally:
586
+ _tache.cancel()
587
 
588
  if think_open:
589
  yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx})