patdev/k3-a40-bootstrap
01.3k
1"""Pont Anthropic -> OpenAI, pour brancher Claude Code sur un serveur vLLM.2 3Claude Code parle le protocole Anthropic (`POST /v1/messages`, SSE a evenements4nommes). vLLM ne sert que le protocole OpenAI. Ce module traduit dans les deux5sens, en streaming comme en une passe, avec les appels d'outils -- sans quoi un6agent ne peut rien faire.7 8Lance a cote de vLLM sur la meme machine ; ecoute sur PROXY_PORT et relaie vers9UPSTREAM.10"""11 12from __future__ import annotations13 14import asyncio15import json16import os17import re18import time19import uuid20from typing import Any21 22import httpx23from fastapi import FastAPI, HTTPException, Request24from fastapi.responses import JSONResponse, StreamingResponse25 26UPSTREAM = os.environ.get("VL_UPSTREAM", "http://127.0.0.1:8080")27MODEL = os.environ.get("VL_SERVED_NAME", "qwen")28# ---------------------------------------------------------------- SWAP (v70)29# Deux modeles a tour de role sur la meme carte. Le bootstrap ecrit30# `$VL_SWAP_DIR/modele_courant` (cle nom depot) a chaque lancement de vLLM ; le31# pont y lit le modele charge A CHAQUE REQUETE (il survit au swap). Une requete32# pour un modele non charge ecrit `modele_demande`, puis attend que le33# bootstrap ait relance vLLM (jusqu'a VL_SWAP_TIMEOUT s).34SWAP_DIR = os.environ.get("VL_SWAP_DIR", "/travail")35SWAP = dict(x.split(":", 1) for x in os.environ.get("VL_SWAP", "").split(",") if ":" in x)36SWAP_TIMEOUT = float(os.environ.get("VL_SWAP_TIMEOUT", "900"))37 38 39def _courant() -> tuple[str, str, str, str]:40 """(cle bootstrap, nom servi, depot, speculation) du modele charge."""41 try:42 with open(os.path.join(SWAP_DIR, "modele_courant"), encoding="utf-8") as f:43 champs = f.read().split()44 cle, nom, depot = champs[:3]45 return cle, nom, depot, (champs[3] if len(champs) > 3 else "off")46 except Exception:47 return os.environ.get("VL_MODEL_KEY", ""), MODEL, os.environ.get("VL_REAL_MODEL", MODEL), "off"48 49 50def _cle_demandee(model: object) -> tuple[str, str] | None:51 """`ornith` -> (cle, "") ; `ornith+dspark` -> (cle, "dspark") : la speculation52 voyage dans le nom du modele, pour changer de reglage sans recreer le pod."""53 if not isinstance(model, str):54 return None55 nom = model.removesuffix("[1m]")56 if nom.startswith("claude-"):57 nom = nom[len("claude-"):]58 nom, _, spec = nom.partition("+")59 cle = SWAP.get(nom)60 return (cle, spec) if cle else None61TIMEOUT = float(os.environ.get("VL_TIMEOUT", "1800"))62# Sortie maximale du modele. Claude Code demande couramment 64 k, ce que63# vLLM refuse d'un 400 portant sur max_tokens -- la requete entiere echoue64# alors qu'un plafonnement silencieux suffit.65MAX_OUTPUT = int(os.environ.get("VL_MAX_OUTPUT", "32768"))66 67# Claude Code refuse tout identifiant de modele qui ne commence pas par68# "claude-" : il valide le nom avant d'emettre la requete. On expose donc des69# alias conformes, et on ignore le nom recu pour router vers l'unique modele70# reellement charge -- le client choisit une etiquette, pas un moteur.71# Le PREMIER alias nomme le modele reellement charge : sans cela, un client qui72# voit "claude-kimi-k3" croit legitimement executer du Kimi alors que le moteur73# sert du Qwen. Les suivants sont des etiquettes de compatibilite, et tous74# routent vers l'unique modele charge.75_REAL = {"qwen": "claude-qwen3-coder-30b", "kimi": "claude-kimi-linear-48b"}76ALIASES = [77 # Alias NU, sans prefixe. Indispensable pour une fenetre > 200 k.78 # Claude Code 2.1.239 (fonction JFd du binaire) :79 # let n = CLAUDE_CODE_MAX_CONTEXT_TOKENS;80 # if (n !== undefined && n > 0 && !id.startsWith("claude-")) return n;81 # return <defaut 200 000>82 # Autrement dit il REFUSE toute fenetre personnalisee sur un identifiant83 # commencant par "claude-" : avec `claude-flashnext`, ni84 # CLAUDE_CODE_MAX_CONTEXT_TOKENS ni CLAUDE_CODE_AUTO_COMPACT_WINDOW n'ont85 # le moindre effet, et /context affiche 200k quoi qu'on fasse.86 MODEL,87 _REAL.get(MODEL, f"claude-{MODEL}"),88 "claude-kimi-k3",89 "claude-kimi-k3-linear",90 "claude-qwen3-coder",91 "claude-sonnet-4-5", # alias de compatibilite : certains clients92 "claude-3-5-haiku", # codent en dur un modele "rapide" et un "lent"93]94# Variantes "[1m]". Claude Code deduit la fenetre de contexte du NOM du modele :95# un identifiant qu'il ne connait pas est suppose a 200 k, et l'auto-compactage96# se declenche a 200 k meme si le moteur en accepte 1 000 000. Le suffixe [1m]97# est sa convention pour la fenetre du million ; VERIFIE : avec lui,98# l'avertissement "auto-compact will keep this session within 200k" disparait.99for _n in SWAP: # les modeles interchangeables, tous annonces100 ALIASES += [_n, f"claude-{_n}", f"{_n}+dspark", f"{_n}+off"]101ALIASES = [a for pair in ((x, f"{x}[1m]") for x in ALIASES) for a in pair]102ALIASES = list(dict.fromkeys(ALIASES)) # dedoublonne en gardant l'ordre103 104app = FastAPI(title="anthropic-bridge")105 106_verrou_swap = asyncio.Lock()107 108 109async def _assurer_modele(model: object) -> str:110 """Rend le nom servi pour `model`, en declenchant le swap s'il le faut."""111 dem = _cle_demandee(model)112 courant = _courant()113 if not dem:114 return courant[1]115 cle, spec = dem116 117 def satisfait(c):118 # sans "+spec" dans le nom, n'importe quelle speculation du bon modele convient119 return c[0] == cle and (not spec or c[3] == spec)120 121 if satisfait(courant):122 return courant[1]123 async with _verrou_swap:124 courant = _courant()125 if satisfait(courant):126 return courant[1]127 with open(os.path.join(SWAP_DIR, "modele_demande"), "w", encoding="utf-8") as f:128 f.write(f"{cle} {spec}\n")129 t0 = time.time()130 while time.time() - t0 < SWAP_TIMEOUT:131 await asyncio.sleep(3)132 c = _courant()133 if not satisfait(c):134 continue135 try:136 r = await _client.get("/health", timeout=3.0)137 if r.status_code == 200:138 return c[1]139 except Exception:140 pass141 raise HTTPException(status_code=503, detail=f"swap vers {cle} non termine en {SWAP_TIMEOUT:.0f} s")142 143 144@app.middleware("http")145async def _guerir_bassin(request: Request, call_next):146 generation = _bassin_generation147 try:148 return await call_next(request)149 except httpx.PoolTimeout:150 await _reconstruire_bassin(generation)151 return JSONResponse(152 status_code=503,153 content={"type": "error",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 en159 # PoolTimeout. Un RuntimeError("client has been closed") apres une160 # reconstruction, ou un ReadError sur une connexion recyclee, sortait en161 # 500 brut sans jamais declencher la guerison. Meme remede, et le type162 # REEL est journalise pour qu'on ne rediagnostique plus a l'aveugle.163 import traceback164 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 -- les178# flux abandonnes cote client (retry, coupure proxy) gardaient leur connexion179# amont, vLLM generait pour personne, et a 100 connexions (defaut httpx) le pool180# bloquait TOUT, /health compris : pont muet, processus vivant a 3 % CPU.181# Le timeout de pool transforme une saturation en 503 rapide au lieu d'un blocage.182def _nouveau_client() -> httpx.AsyncClient:183 return httpx.AsyncClient(184 base_url=UPSTREAM,185 timeout=httpx.Timeout(TIMEOUT, connect=10.0, pool=10.0),186 limits=httpx.Limits(max_connections=512, max_keepalive_connections=64))187 188 189_client = _nouveau_client()190 191# GUERISON DU BASSIN, ajoutee le 27/08 apres trois saturations en douze heures,192# a intervalle d'environ une heure et quart sous charge multi-agents.193#194# Ce que la mesure a montre, et qui change le diagnostic : au moment du blocage,195# `ss -tan | grep :8000` renvoyait **zero** socket. Le bassin n'est donc pas196# plein de connexions vivantes -- il est plein de PLACES comptabilisees et197# jamais rendues. Consequence : baisser VL_TIMEOUT de 1800 a 240 s n'a rien198# change, parce que la fuite n'est pas temporelle, elle est definitive par flux199# abandonne. Monter `max_connections` ne fait que retarder l'echeance.200#201# On ne repare donc pas la fuite ici (elle est dans httpx/httpcore, pas dans ce202# fichier) : on rend le pont capable d'en sortir seul. A la premiere PoolTimeout203# on remplace le client par un neuf et on ferme l'ancien en arriere-plan. Le204# compteur de generation evite que dix requetes simultanees reconstruisent dix205# fois : seule celle qui a vu la generation courante agit.206_bassin_verrou = asyncio.Lock()207_bassin_generation = 0208 209 210async def _reconstruire_bassin(generation_vue: int) -> None:211 global _client, _bassin_generation212 async with _bassin_verrou:213 if generation_vue != _bassin_generation:214 return215 ancien, _client = _client, _nouveau_client()216 _bassin_generation += 1217 print("[bassin] PoolTimeout -> client httpx reconstruit "218 "(generation %d)" % _bassin_generation, flush=True)219 try:220 await ancien.aclose()221 except Exception:222 pass223 224 225# ------------------------------------------------------------------ garde-fou226# Un modele de 30 a 48 milliards de parametres quantifie en 4 bits abrege : somme227# de reecrire un long fichier, il repond `// ... [previous content] ...` et228# considere le travail fait. Avec --permission-mode acceptEdits, cet abrege est229# ecrit sur le disque sans qu'un diff soit montre : 1841 lignes de source ont230# ainsi disparu. Le modele ne sait pas qu'il a detruit quelque chose, donc il231# rapporte un succes.232#233# C'est un defaut de capacite, pas de configuration : aucun reglage de vLLM ne234# le corrige. Ce qu'on peut faire, en revanche, c'est refuser l'ecriture. Le235# pont voit passer tous les arguments d'outil ; il est le dernier endroit ou236# l'on puisse transformer une destruction silencieuse en refus visible.237 238GUARD = os.environ.get("VL_GUARD", "on") != "off"239 240# Les outils dont un argument atterrit tel quel dans un fichier.241# Seules les ecritures de fichier ENTIER sont gardees. Les editions ciblees242# (Edit, MultiEdit, str_replace_*, NotebookEdit) portent legitimement des243# "..." dans old_string/new_string -- et le message du garde leur disait244# justement d'utiliser une edition ciblee : boucle de refus, observee le 22/08.245_GUARDED = {"write", "write_file", "create_file"}246 247# Marqueurs d'omission EXPLICITES seulement. Les motifs generiques "[...]" et248# "# ..." bloquaient des contenus legitimes (un .md avec une ligne "# ...", un249# script contenant "[...]") : faux positifs observes le 22/08 sur une vraie250# session Claude Code.251PLACEHOLDER = re.compile(252 r"\.\.\.\s*(\[|#|//|--)?\s*(previous|rest of|remaining|existing|unchanged|original)"253 r"|\[\s*\.\.\.\s*(previous|rest|remaining|existing|unchanged|original|snip)[^\]]*\]"254 r"|<\s*(unchanged|snip|elided)\s*>"255 r"|\(\s*(reste|suite) (du|des) ",256 re.I | re.M,257)258 259# Rappel injecte en tete de systeme. Ne remplace pas le garde-fou -- un modele260# qui abrege le fait souvent malgre la consigne -- mais reduit la frequence, et261# ne coute qu'une constante en tete de prefixe, donc le cache reste partage.262SYSTEM_GUARD = (263 "Quand tu ecris un fichier, l'argument `content` doit contenir le fichier "264 "INTEGRAL. N'ecris jamais de marqueur d'omission du type "265 "\"... [previous content] ...\", \"// ... rest of file ...\" ou "266 "\"(reste du fichier inchange)\" : ces marqueurs sont ecrits litteralement "267 "sur le disque et detruisent le fichier. Si le fichier est trop long pour "268 "etre reemis en entier, dis-le et utilise une edition ciblee."269)270 271 272def guard_violation(name: str, args: Any) -> str | None:273 """Renvoie la description de l'abreviation trouvee, ou None."""274 if not GUARD or (name or "").lower() not in _GUARDED:275 return None276 277 def walk(v: Any, path: str) -> str | None:278 if isinstance(v, str):279 m = PLACEHOLDER.search(v)280 return f"{path} contient {m.group(0).strip()!r}" if m else None281 if isinstance(v, dict):282 for k, sub in v.items():283 hit = walk(sub, f"{path}.{k}")284 if hit:285 return hit286 elif isinstance(v, list):287 for i, sub in enumerate(v):288 hit = walk(sub, f"{path}[{i}]")289 if hit:290 return hit291 return None292 293 return walk(args, name)294 295 296def guard_message(detail: str) -> str:297 return (298 "\n\n[garde-fou du pont] Appel d'outil bloque : " + detail + ".\n"299 "Le contenu propose abrege le fichier par un marqueur d'omission, ce qui "300 "l'aurait ecrase par une version incomplete. L'ecriture n'a PAS eu lieu et "301 "le fichier est intact.\n"302 "Reemets le fichier integral, ou fais une edition ciblee qui ne remplace "303 "que les lignes concernees."304 )305 306 307# --------------------------------------------------------------- Anthropic -> OpenAI308 309def _image_part(block: dict) -> dict | None:310 """Traduit un bloc `image` Anthropic vers la partie `image_url` d'OpenAI.311 312 Anthropic decrit l'image par une `source` typee : `base64` porte les octets313 et le type MIME separement, `url` porte un lien. OpenAI attend dans les deux314 cas UNE chaine dans `image_url.url` -- une URI de donnees pour le premier,315 le lien tel quel pour le second.316 317 Sans cette traduction, `_text_of` ignorait purement et simplement les blocs318 image : une capture collee dans Claude Code arrivait au modele comme un319 message vide, et le modele repondait a cote sans que rien ne signale la320 perte.321 """322 src = block.get("source") or {}323 kind = src.get("type")324 if kind == "base64":325 data = src.get("data")326 if not data:327 return None328 mime = src.get("media_type") or "image/png"329 return {"type": "image_url",330 "image_url": {"url": f"data:{mime};base64,{data}"}}331 if kind == "url" and src.get("url"):332 return {"type": "image_url", "image_url": {"url": src["url"]}}333 return None334 335 336def _text_of(content: Any) -> str:337 """Anthropic autorise une chaine ou une liste de blocs typés."""338 if isinstance(content, str):339 return content340 if not isinstance(content, list):341 return ""342 out = []343 for b in content:344 if isinstance(b, str):345 out.append(b)346 elif isinstance(b, dict) and b.get("type") == "text":347 out.append(b.get("text", ""))348 return "".join(out)349 350 351def to_openai(body: dict) -> dict:352 msgs: list[dict] = []353 354 # `system` est un champ separe chez Anthropic, un message de role chez OpenAI.355 # Le rappel anti-abreviation n'est ajoute qu'en presence d'outils : sans356 # outil, aucun contenu n'atteint le disque et la consigne serait du bruit.357 sys = _text_of(body.get("system") or "")358 if GUARD and body.get("tools"):359 sys = (sys + "\n\n" + SYSTEM_GUARD) if sys else SYSTEM_GUARD360 if sys:361 msgs.append({"role": "system", "content": sys})362 363 for m in body.get("messages", []):364 role = m.get("role", "user")365 content = m.get("content")366 367 # Claude Code place un SECOND message systeme en fin de conversation368 # (rappel de contexte). Le gabarit Qwen n'accepte le role systeme qu'en369 # tete et vLLM refuse toute la requete :370 # 400 "System message must be at the beginning."371 # Le pont ne verifiait pas le statut amont, donc l'echec ressortait en372 # reponse VIDE avec stop_reason=end_turn -- Claude Code affichait373 # "Cogitated for 0s" et rien d'autre.374 #375 # On le convertit en message utilisateur plutot que de le fusionner dans376 # le systeme de tete : ce rappel change a chaque tour, et le remonter en377 # tete invaliderait le prefixe partage, donc le cache -- soit 68 % des378 # blocs, mesures.379 if role == "system" and msgs and msgs[0].get("role") == "system":380 txt = _text_of(content)381 if txt:382 msgs.append({"role": "user", "content": txt})383 continue384 385 if isinstance(content, list):386 # Un tour d'assistant peut melanger du texte et des tool_use ; un tour387 # d'utilisateur porte les tool_result. OpenAI separe les deux en388 # `tool_calls` sur l'assistant et en messages de role `tool`.389 texts, calls, results, images, thinks = [], [], [], [], []390 for b in content:391 if not isinstance(b, dict):392 continue393 t = b.get("type")394 if t == "text":395 texts.append(b.get("text", ""))396 elif t == "thinking":397 # Claude Code renvoie le raisonnement du tour precedent dans398 # l'historique (boucles d'outils). Le gabarit Qwen3.5 sait le399 # garder pour le DERNIER tour d'assistant via400 # `reasoning_content` -- exactement la semantique Anthropic.401 # Avant : bloc ignore en silence, le modele perdait son plan402 # entre deux appels d'outil.403 thinks.append(b.get("thinking", ""))404 elif t == "image":405 part = _image_part(b)406 if part:407 images.append(part)408 elif t == "tool_use":409 calls.append({410 "id": b.get("id") or f"call_{uuid.uuid4().hex[:8]}",411 "type": "function",412 "function": {413 "name": b.get("name", ""),414 "arguments": json.dumps(b.get("input") or {}),415 },416 })417 elif t == "tool_result":418 # Un resultat d'outil peut porter des images (capture rendue419 # par un outil). Le role `tool` d'OpenAI n'accepte que du420 # texte : on extrait les images pour les rattacher au tour421 # utilisateur, sinon elles disparaissent en silence.422 rc = b.get("content")423 if isinstance(rc, list):424 for sub in rc:425 if isinstance(sub, dict) and sub.get("type") == "image":426 part = _image_part(sub)427 if part:428 images.append(part)429 results.append({430 "role": "tool",431 "tool_call_id": b.get("tool_use_id", ""),432 "content": _text_of(rc) or "",433 })434 435 if role == "assistant":436 a: dict[str, Any] = {"role": "assistant", "content": "".join(texts) or None}437 if thinks:438 a["reasoning_content"] = chr(10).join(t for t in thinks if t)439 if calls:440 a["tool_calls"] = calls441 msgs.append(a)442 else:443 if images:444 # Contenu multipart : OpenAI n'accepte les images que dans445 # une LISTE de parties, jamais dans une chaine.446 parts: list[dict] = []447 joined = "".join(texts)448 if joined:449 parts.append({"type": "text", "text": joined})450 parts.extend(images)451 msgs.append({"role": "user", "content": parts})452 elif texts:453 msgs.append({"role": "user", "content": "".join(texts)})454 msgs.extend(results)455 else:456 msgs.append({"role": role, "content": content or ""})457 458 out: dict[str, Any] = {459 "model": _courant()[1],460 "messages": msgs,461 "max_tokens": min(int(body.get("max_tokens") or 4096), MAX_OUTPUT),462 "stream": bool(body.get("stream")),463 }464 for src, dst in (("temperature", "temperature"), ("top_p", "top_p"),465 ("stop_sequences", "stop")):466 if body.get(src) is not None:467 out[dst] = body[src]468 469 # Le raisonnement est ACTIF par defaut sur ce modele et consomme des jetons470 # avant le premier caractere de reponse : une requete a max_tokens=40 revient471 # avec content vide et stop_reason=max_tokens. Le protocole Anthropic exprime472 # la coupure par `thinking: {"type": "disabled"}` ; le gabarit Qwen l'attend473 # sous la forme enable_thinking=false, qui fait prefixer un bloc <think> deja474 # ferme. Sans cette traduction, le champ etait recu puis ignore en silence :475 # la case "raisonnement" de la console ne changeait rien.476 # `thinking` a plusieurs formes. La documentation du protocole de passerelle477 # precise que Claude Code envoie `{"type": "adaptive"}` aux modeles recents478 # ET "traite les noms de modeles qu'il ne reconnait pas, tels les alias de479 # passerelle, comme des modeles actuels qui recoivent le champ" -- donc nous.480 # Seul "disabled" doit couper le raisonnement ; "adaptive" et "enabled" le481 # laissent actif, qui est le defaut du gabarit.482 th = body.get("thinking")483 if isinstance(th, dict) and th.get("type") == "disabled":484 out["chat_template_kwargs"] = {"enable_thinking": False}485 elif isinstance(th, dict) and th.get("budget_tokens"):486 # Le niveau de raisonnement de Claude Code arrive ici, en jetons. vLLM487 # 0.27.1 l'applique via `thinking_token_budget` (sampling params) : il488 # ferme le bloc de raisonnement a ce nombre de jetons. Mesure : sans489 # ce champ, budget 512 et budget 8000 donnaient la meme chose.490 try:491 out["thinking_token_budget"] = max(1, int(th["budget_tokens"]))492 except (TypeError, ValueError):493 pass494 # Niveau d'effort (output_config.effort) -> plafond de jetons de raisonnement,495 # memes seuils que le patch natif (vllm_anthropic_effort_patch.py). Le budget496 # explicite ci-dessus gagne ; `disabled` aussi.497 _eff = (body.get("output_config") or {}).get("effort") if isinstance(body.get("output_config"), dict) else None498 _BUDGET = {"low": 1024, "medium": 4096, "high": 16384, "xhigh": 32768, "max": None}499 if _eff in _BUDGET and "thinking_token_budget" not in out and "chat_template_kwargs" not in out:500 if _BUDGET[_eff] is not None:501 out["thinking_token_budget"] = _BUDGET[_eff]502 503 if body.get("tools"):504 out["tools"] = [{505 "type": "function",506 "function": {507 "name": t["name"],508 "description": t.get("description", ""),509 "parameters": t.get("input_schema") or {"type": "object", "properties": {}},510 },511 } for t in body["tools"] if t.get("name")]512 tc = body.get("tool_choice") or {}513 kind = tc.get("type") if isinstance(tc, dict) else None514 if kind == "none":515 out["tool_choice"] = "none"516 elif kind == "any":517 out["tool_choice"] = "required"518 elif kind == "tool" and tc.get("name"):519 out["tool_choice"] = {"type": "function", "function": {"name": tc["name"]}}520 else:521 out["tool_choice"] = "auto"522 return out523 524 525_STOP = {"stop": "end_turn", "length": "max_tokens", "tool_calls": "tool_use"}526 527 528# --------------------------------------------------------------- OpenAI -> Anthropic529 530def _usage_of(usage: dict) -> dict:531 """vLLM expose les jetons servis par le cache de prefixe dans532 `prompt_tokens_details.cached_tokens` ; Anthropic les appelle533 `cache_read_input_tokens`. Sans cette traduction Claude Code affichait534 0 % de cache alors que vLLM en servait 68 %."""535 u = {"input_tokens": usage.get("prompt_tokens", 0) or 0,536 "output_tokens": usage.get("completion_tokens", 0) or 0}537 det = usage.get("prompt_tokens_details") or {}538 cached = det.get("cached_tokens") if isinstance(det, dict) else None539 if cached:540 u["cache_read_input_tokens"] = int(cached)541 u["cache_creation_input_tokens"] = 0542 return u543 544 545def _stop_of(choice: dict) -> tuple[str, str | None]:546 """finish_reason=stop couvre deux cas Anthropic : end_turn, ou547 stop_sequence quand vLLM rapporte la chaine d'arret dans `stop_reason`."""548 fr = choice.get("finish_reason") or "stop"549 sr = choice.get("stop_reason")550 if fr == "stop" and isinstance(sr, str) and sr:551 return "stop_sequence", sr552 return _STOP.get(fr, "end_turn"), None553 554 555def to_anthropic(oai: dict, req_model: str) -> dict:556 choice = (oai.get("choices") or [{}])[0]557 msg = choice.get("message") or {}558 blocks: list[dict] = []559 560 # Le raisonnement arrive ici en non-flux ; en flux il est deja traduit en561 # bloc thinking. Meme forme dans les deux cas, signature comprise.562 reasoning = msg.get("reasoning_content") or msg.get("reasoning")563 if reasoning:564 blocks.append({"type": "thinking", "thinking": reasoning,565 "signature": "vllm-" + uuid.uuid4().hex[:16]})566 if msg.get("content"):567 blocks.append({"type": "text", "text": msg["content"]})568 elif reasoning and not msg.get("tool_calls"):569 # meme repli qu'en flux : pas de texte, pas d'outil -> le raisonnement570 blocks.append({"type": "text", "text": str(reasoning).strip()})571 572 blocked = False573 for c in msg.get("tool_calls") or []:574 fn = c.get("function") or {}575 try:576 args = json.loads(fn.get("arguments") or "{}")577 except json.JSONDecodeError:578 # Un modele quantifie peut emettre du JSON legerement casse ; mieux579 # vaut transmettre la chaine brute que faire tomber la requete.580 args = {"_raw": fn.get("arguments", "")}581 name = fn.get("name", "")582 detail = guard_violation(name, args)583 if detail:584 blocked = True585 blocks.append({"type": "text", "text": guard_message(detail)})586 continue587 blocks.append({588 "type": "tool_use",589 "id": c.get("id") or f"toolu_{uuid.uuid4().hex[:16]}",590 "name": name,591 "input": args,592 })593 594 usage = oai.get("usage") or {}595 if blocked and not any(b["type"] == "tool_use" for b in blocks):596 # Plus aucun outil a executer : le tour se termine sur l'explication.597 return {598 "id": oai.get("id") or f"msg_{uuid.uuid4().hex[:24]}",599 "type": "message", "role": "assistant", "model": req_model,600 "content": blocks, "stop_reason": "end_turn", "stop_sequence": None,601 "usage": _usage_of(usage),602 }603 stop_reason, stop_seq = _stop_of(choice)604 return {605 "id": oai.get("id") or f"msg_{uuid.uuid4().hex[:24]}",606 "type": "message",607 "role": "assistant",608 "model": req_model,609 "content": blocks,610 "stop_reason": stop_reason,611 "stop_sequence": stop_seq,612 "usage": _usage_of(usage),613 }614 615 616def _err_of(body: str, status: int) -> dict:617 """Renvoie l'objet d'erreur amont intact, ou en fabrique un equivalent.618 619 Le libelle compte : c'est sur lui que le client decide s'il peut retenter.620 """621 try:622 d = json.loads(body)623 e = d.get("error")624 if isinstance(e, dict) and e.get("message"):625 return {"type": e.get("type") or "api_error", "message": e["message"]}626 except (json.JSONDecodeError, AttributeError):627 pass628 return {"type": "api_error", "message": body or f"upstream {status}"}629 630 631def _sse(event: str, data: dict) -> bytes:632 return f"event: {event}\ndata: {json.dumps(data, separators=(',', ':'))}\n\n".encode()633 634 635async def _flux_gueri(gen):636 """v73 : les exceptions nees dans un generateur SSE ne traversent pas le637 middleware _guerir_bassin. On les attrape ici : evenement `error` au client638 (sa logique de reprise lit le libelle) et reconstruction du bassin."""639 generation = _bassin_generation640 try:641 async for morceau in gen:642 yield morceau643 except (httpx.HTTPError, RuntimeError) as e:644 import traceback645 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 655async def stream_anthropic(payload: dict, req_model: str, request: Request | None = None):656 """Traduit le flux OpenAI en evenements Anthropic.657 658 Le point delicat est l'indexation des blocs : Anthropic numerote chaque bloc659 de contenu et exige un content_block_start/stop apparie, alors qu'OpenAI660 emet des deltas plats. On ouvre donc un bloc texte a la volee, et un bloc661 tool_use par appel, en fermant le precedent.662 663 Les arguments d'outil, eux, sont accumules et n'emis qu'une fois complets.664 Il le faut : un argument diffuse fragment par fragment ne peut pas etre665 valide, et le client aurait deja commence a ecrire le fichier quand le666 marqueur d'omission apparait.667 668 Mais accumuler veut dire n'envoyer AUCUN octet pendant toute la generation669 de l'appel, et le proxy de Runpod coupe une connexion inactive vers 125 s.670 MESURE : une reecriture de 1354 lignes par l'outil Write mourait a 126 s671 avec zero token et stop=None, alors que la meme reecriture en prose passait672 (le texte, lui, est diffuse au fil de l'eau). Un test voisin est passe de673 justesse a 123,6 s. On emet donc un `ping` periodique tant qu'on accumule :674 le protocole Anthropic le prevoit, les clients l'ignorent, et la connexion675 reste ouverte sans que le garde-fou perde sa capacite a valider avant676 emission.677 """678 mid = f"msg_{uuid.uuid4().hex[:24]}"679 yield _sse("message_start", {680 "type": "message_start",681 "message": {"id": mid, "type": "message", "role": "assistant",682 "model": req_model, "content": [], "stop_reason": None,683 "stop_sequence": None,684 "usage": {"input_tokens": 0, "output_tokens": 0}},685 })686 yield _sse("ping", {"type": "ping"})687 688 idx = -1689 text_open = False690 think_open = False691 n_text = 0692 n_think = 0693 think_buf: list[str] = [] # raisonnement accumule, pour le repli "reponse vide"694 last_usage_sent = -1695 last_usage_ts = 0.0696 # Horodatage du dernier octet REELLEMENT envoye au client, pour savoir quand697 # le silence devient dangereux.698 last_out = time.monotonic()699 KEEPALIVE_S = 15.0700 tools: dict[int, dict] = {} # index OpenAI -> {id, name, args}701 stop_reason = "end_turn"702 out_tokens = 0703 in_tokens = 0704 cached_tokens = 0705 stop_seq: str | None = None706 707 async with _client.stream("POST", "/v1/chat/completions", json=payload) as r:708 # Un 4xx amont renvoie du JSON d'erreur, pas des lignes "data:". Sans ce709 # controle, chaque ligne est ignoree, le flux se termine vide et le710 # client recoit un message sans contenu avec stop_reason=end_turn : une711 # panne parfaitement silencieuse, indiscernable d'un modele muet.712 if r.status_code >= 400:713 body = (await r.aread()).decode("utf-8", "ignore")[:800]714 roles = "/".join(m.get("role", "?") for m in payload.get("messages", []))715 print(f"[amont {r.status_code}] {body}", flush=True)716 print(f"[amont] roles={roles}", flush=True)717 # Le corps d'erreur repart TEL QUEL, dans un evenement SSE `error`.718 # La documentation du protocole de passerelle est explicite : Claude719 # Code se remet tout seul de certains refus -- champ `thinking`,720 # signatures de raisonnement, et justement les messages systeme en721 # milieu de conversation -- mais "la logique de reprise s'appuie sur722 # le libelle de l'erreur amont", et "une passerelle qui enveloppe723 # les erreurs dans son propre format casse la reprise meme si elle724 # preserve le code de statut". Emballer l'erreur dans un bloc texte,725 # comme je le faisais, empechait donc cette reprise.726 yield _sse("error", {"type": "error", "error": _err_of(body, r.status_code)})727 return728 # Le `ping` periodique ne suffisait pas : il ne se declenchait qu'a729 # l'interieur de cette boucle, donc uniquement quand vLLM emettait deja730 # quelque chose. Or pendant le PREREMPLISSAGE vLLM n'envoie rien -- et731 # un prompt de plusieurs centaines de milliers de jetons met plusieurs732 # minutes a etre calcule. Le flux restait muet, et le proxy Runpod733 # coupait vers 125 s.734 #735 # MESURE : une aiguille a ~700 k jetons echouait a 126,5 s avec zero736 # jeton, alors que la meme a ~358 k passait en 94,4 s. Le plafond737 # n'etait pas le modele mais le silence : aucune requete au-dela de738 # ~500 k ne pouvait aboutir, quel que soit le contexte annonce.739 #740 # On decouple donc la lecture amont de l'emission aval : une tache741 # pompe les lignes dans une file, et l'absence de ligne pendant742 # KEEPALIVE_S produit un `ping` au lieu d'un silence.743 file: asyncio.Queue = asyncio.Queue()744 745 async def _pompe() -> None:746 try:747 async for _l in r.aiter_lines():748 await file.put(_l)749 finally:750 await file.put(None)751 752 _tache = asyncio.create_task(_pompe())753 try:754 while True:755 try:756 line = await asyncio.wait_for(file.get(), timeout=KEEPALIVE_S)757 except asyncio.TimeoutError:758 # Client parti pendant le prefill ? On ferme l'amont (vLLM759 # annule alors la generation) au lieu de pomper dans le vide.760 if request is not None and await request.is_disconnected():761 print("[deconnexion] client parti pendant l'attente, amont ferme", flush=True)762 break763 yield _sse("ping", {"type": "ping"})764 last_out = time.monotonic()765 continue766 if line is None:767 break768 # Le `yield` ne leve pas toujours a la deconnexion (le proxy Runpod769 # garde parfois sa connexion ouverte) : on verifie explicitement.770 if request is not None and (out_tokens % 32 == 0) and await request.is_disconnected():771 print(f"[deconnexion] client parti apres {out_tokens} jetons, amont ferme", flush=True)772 break773 if not line.startswith("data: "):774 continue775 chunk = line[6:].strip()776 if chunk == "[DONE]":777 break778 try:779 d = json.loads(chunk)780 except json.JSONDecodeError:781 continue782 783 if d.get("usage"):784 in_tokens = d["usage"].get("prompt_tokens", in_tokens)785 out_tokens = d["usage"].get("completion_tokens", out_tokens)786 _det = d["usage"].get("prompt_tokens_details") or {}787 if isinstance(_det, dict) and _det.get("cached_tokens"):788 cached_tokens = int(_det["cached_tokens"])789 790 ch = (d.get("choices") or [{}])[0]791 delta = ch.get("delta") or {}792 793 # Comptage en direct facon Anthropic : vLLM (continuous_usage_stats)794 # renvoie le total cumule sur chaque chunk (lu ci-dessus). On le795 # relaie dans un message_delta periodique pour que le compteur du796 # client monte en continu au lieu de sauter au total a la fin.797 if out_tokens and out_tokens != last_usage_sent and (time.monotonic() - last_usage_ts) >= 0.5:798 yield _sse("message_delta", {"type": "message_delta",799 "delta": {"stop_reason": None, "stop_sequence": None},800 "usage": {"input_tokens": in_tokens, "output_tokens": out_tokens}})801 last_usage_sent = out_tokens; last_usage_ts = time.monotonic(); last_out = time.monotonic()802 803 # vLLM range le raisonnement dans un champ separe. Sans ce relais,804 # il est simplement perdu : le client ne peut ni l'afficher ni savoir805 # pourquoi la reponse tarde.806 if (delta.get("reasoning_content") or delta.get("reasoning")):807 if not think_open:808 if text_open:809 yield _sse("content_block_stop",810 {"type": "content_block_stop", "index": idx})811 text_open = False812 idx += 1813 think_open = True814 yield _sse("content_block_start", {815 "type": "content_block_start", "index": idx,816 "content_block": {"type": "thinking", "thinking": ""}})817 yield _sse("content_block_delta", {818 "type": "content_block_delta", "index": idx,819 "delta": {"type": "thinking_delta",820 "thinking": (delta.get("reasoning_content") or delta.get("reasoning"))}})821 n_think += len((delta.get("reasoning_content") or delta.get("reasoning")))822 think_buf.append(delta.get("reasoning_content") or delta.get("reasoning"))823 last_out = time.monotonic()824 825 if delta.get("content"):826 if think_open:827 yield _sse("content_block_delta", {828 "type": "content_block_delta", "index": idx,829 "delta": {"type": "signature_delta",830 "signature": "vllm-" + uuid.uuid4().hex[:16]}})831 yield _sse("content_block_stop",832 {"type": "content_block_stop", "index": idx})833 think_open = False834 if not text_open:835 idx += 1836 text_open = True837 yield _sse("content_block_start", {838 "type": "content_block_start", "index": idx,839 "content_block": {"type": "text", "text": ""}})840 yield _sse("content_block_delta", {841 "type": "content_block_delta", "index": idx,842 "delta": {"type": "text_delta", "text": delta["content"]}})843 n_text += len(delta["content"])844 last_out = time.monotonic()845 846 for tc in delta.get("tool_calls") or []:847 i = tc.get("index", 0)848 fn = tc.get("function") or {}849 t = tools.setdefault(i, {"id": None, "name": "", "args": ""})850 if tc.get("id"):851 t["id"] = tc["id"]852 if fn.get("name"):853 t["name"] = fn["name"]854 if fn.get("arguments"):855 t["args"] += fn["arguments"]856 857 # Le silence de l'accumulation est ce qui tuait la connexion.858 if time.monotonic() - last_out > KEEPALIVE_S:859 yield _sse("ping", {"type": "ping"})860 last_out = time.monotonic()861 862 if ch.get("finish_reason"):863 stop_reason, stop_seq = _stop_of(ch)864 finally:865 _tache.cancel()866 867 if think_open:868 # Structure Anthropic complete : la signature precede la fermeture du869 # bloc. Claude Code la renvoie telle quelle dans l'historique ; le pont870 # ne la verifie pas, mais un client strict l'attend.871 yield _sse("content_block_delta", {872 "type": "content_block_delta", "index": idx,873 "delta": {"type": "signature_delta",874 "signature": "vllm-" + uuid.uuid4().hex[:16]}})875 yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx})876 think_open = False877 if text_open:878 yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx})879 text_open = False880 881 emitted = 0882 for i in sorted(tools):883 t = tools[i]884 try:885 args = json.loads(t["args"] or "{}")886 except json.JSONDecodeError:887 args = {"_raw": t["args"]}888 889 detail = guard_violation(t["name"], args)890 if detail:891 print(f"[garde-fou] appel bloque : {detail}", flush=True)892 idx += 1893 yield _sse("content_block_start", {894 "type": "content_block_start", "index": idx,895 "content_block": {"type": "text", "text": ""}})896 yield _sse("content_block_delta", {897 "type": "content_block_delta", "index": idx,898 "delta": {"type": "text_delta", "text": guard_message(detail)}})899 yield _sse("content_block_stop",900 {"type": "content_block_stop", "index": idx})901 continue902 903 idx += 1904 emitted += 1905 yield _sse("content_block_start", {906 "type": "content_block_start", "index": idx,907 "content_block": {"type": "tool_use",908 "id": t["id"] or f"toolu_{uuid.uuid4().hex[:16]}",909 "name": t["name"], "input": {}}})910 # Un seul fragment : l'argument est deja complet et valide.911 yield _sse("content_block_delta", {912 "type": "content_block_delta", "index": idx,913 "delta": {"type": "input_json_delta",914 "partial_json": json.dumps(args, separators=(",", ":"))}})915 yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx})916 917 # Annoncer `tool_use` sans avoir emis d'outil ferait attendre au client un918 # resultat qui ne viendra jamais.919 if stop_reason == "tool_use" and emitted == 0:920 stop_reason = "end_turn"921 922 # Repli "reponse vide" (22/08 nuit) : la compaction automatique de Claude923 # Code a echoue sur "summarization produced empty response" -- le modele924 # avait tout ecrit dans le raisonnement (1 679 car.) et rien en texte. Un925 # client qui attend du texte ne sait rien faire d'un bloc thinking seul :926 # on rend alors le raisonnement comme texte, clairement marque.927 if n_text == 0 and emitted == 0 and "".join(think_buf).strip():928 idx += 1929 yield _sse("content_block_start", {"type": "content_block_start", "index": idx,930 "content_block": {"type": "text", "text": ""}})931 yield _sse("content_block_delta", {"type": "content_block_delta", "index": idx,932 "delta": {"type": "text_delta", "text": "".join(think_buf).strip()}})933 yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx})934 n_text = len("".join(think_buf).strip())935 print("[repli] reponse sans texte ni outil : raisonnement rendu en texte", flush=True)936 937 # Trace compacte d'une requete. Sans elle, un client qui n'affiche rien est938 # indiscernable d'un modele qui ne repond rien : les deux donnent un 200 OK939 # dans le journal d'acces.940 if os.environ.get("VL_TRACE", "on") != "off":941 print(f"[trace] entree={in_tokens} sortie={out_tokens} "942 f"texte={n_text}c raisonnement={n_think}c outils={emitted} "943 f"stop={stop_reason}", flush=True)944 945 _u = {"input_tokens": in_tokens, "output_tokens": out_tokens}946 if cached_tokens:947 _u["cache_read_input_tokens"] = cached_tokens948 _u["cache_creation_input_tokens"] = 0949 yield _sse("message_delta", {950 "type": "message_delta",951 "delta": {"stop_reason": stop_reason, "stop_sequence": stop_seq},952 "usage": _u})953 yield _sse("message_stop", {"type": "message_stop"})954 955 956# ------------------------------------------------------------------------ routes957 958@app.post("/v1/messages")959async def messages(request: Request):960 body = await request.json()961 req_model = body.get("model", MODEL)962 payload = to_openai(body)963 964 if payload.get("stream"):965 payload["stream_options"] = {"include_usage": True, "continuous_usage_stats": True}966 967 async def flux_avec_swap():968 # Le proxy RunPod coupe a ~120 s toute reponse sans octet ; un swap969 # de modele en dure 2 a 6. On tient la ligne avec des pings SSE970 # pendant l'attente, puis on enchaine sur le flux normal.971 tache = asyncio.ensure_future(_assurer_modele(req_model))972 while not tache.done():973 await asyncio.wait({tache}, timeout=15)974 if not tache.done():975 yield _sse("ping", {"type": "ping"})976 try:977 payload["model"] = tache.result()978 except HTTPException as e:979 yield _sse("error", {"type": "error",980 "error": {"type": "overloaded_error", "message": str(e.detail)}})981 return982 async for morceau in stream_anthropic(payload, req_model, request):983 yield morceau984 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)990 if r.status_code != 200:991 return JSONResponse(status_code=r.status_code,992 content={"type": "error",993 "error": {"type": "api_error", "message": r.text[:800]}})994 return JSONResponse(to_anthropic(r.json(), req_model))995 996 997@app.post("/v1/messages/count_tokens")998async def count_tokens(request: Request):999 """Claude Code interroge ce point avant d'envoyer. Une estimation suffit :1000 il s'en sert pour decider de compacter, pas pour facturer."""1001 body = await request.json()1002 # Compte REEL par vLLM (/tokenize en forme chat : gabarit, systeme et1003 # outils compris). L'estimation len/4 ignorait systeme et outils : 131004 # jetons pour un appel qui en pesait bien plus. Repli sur l'estimation si1005 # le tokenizer amont ne repond pas.1006 try:1007 payload = to_openai(body)1008 req = {"model": payload.get("model"), "messages": payload.get("messages", []),1009 "add_generation_prompt": True}1010 if payload.get("tools"):1011 req["tools"] = payload["tools"]1012 r = await _client.post("/tokenize", json=req, timeout=60)1013 if r.status_code == 200:1014 n = r.json().get("count")1015 if isinstance(n, int) and n > 0:1016 return JSONResponse({"input_tokens": n})1017 except Exception: # noqa: BLE0011018 pass1019 n = len(json.dumps(body.get("messages", []))) + len(json.dumps(body.get("system", ""))) \1020 + len(json.dumps(body.get("tools", [])))1021 return JSONResponse({"input_tokens": max(1, n // 4)})1022 1023 1024# ------------------------------------------------------------------- console1025# Une page servie par le pont lui-meme, donc de meme origine que /v1/messages :1026# un fichier ouvert depuis le disque, ou une page hebergee ailleurs, se ferait1027# refuser la requete. C'est aussi la raison pour laquelle le pod n'a besoin que1028# d'un seul port ouvert.1029 1030@app.get("/")1031async def console():1032 from fastapi.responses import HTMLResponse, PlainTextResponse1033 path = os.path.join(os.path.dirname(os.path.abspath(__file__)), "console.html")1034 try:1035 with open(path, encoding="utf-8") as fh:1036 return HTMLResponse(fh.read())1037 except OSError:1038 return PlainTextResponse(1039 "console.html absente a cote du pont.\n"1040 "Le bootstrap la telecharge depuis le Hub au demarrage.\n"1041 f"Attendue ici : {path}\n", status_code=503)1042 1043 1044@app.get("/health")1045async def health():1046 try:1047 r = await _client.get("/health", timeout=5)1048 ok = r.status_code == 2001049 except Exception:1050 ok = False1051 return JSONResponse({"status": "ok" if ok else "upstream_down",1052 "upstream": UPSTREAM, "model": MODEL, "ts": time.time()},1053 status_code=200 if ok else 503)1054 1055 1056@app.get("/v1/models")1057async def models(request: Request):1058 """Deux protocoles sur la meme route.1059 1060 Un client Anthropic attend `{"data":[{"type":"model","id":...}]}` avec des1061 identifiants en "claude-*". Un client OpenAI ou un banc attend la reponse1062 de vLLM, dont il lit `max_model_len`. On distingue sur l'en-tete1063 `anthropic-version`, et on renvoie a chacun ce qu'il sait lire.1064 """1065 if "anthropic-version" not in request.headers:1066 try:1067 r = await _client.get("/v1/models", timeout=2.0)1068 if r.status_code == 200:1069 return JSONResponse(r.json())1070 except Exception:1071 pass1072 1073 ctx = None1074 try:1075 r = await _client.get("/v1/models", timeout=2.0)1076 if r.status_code == 200:1077 ctx = (r.json().get("data") or [{}])[0].get("max_model_len")1078 except Exception:1079 pass1080 1081 data = [{1082 "type": "model",1083 "id": a,1084 # Le nom affiche porte le depot exact : c'est la seule facon pour un1085 # utilisateur de savoir quel modele repond derriere une etiquette.1086 "display_name": f"{a} -> {_courant()[2]}",1087 "created_at": "2026-01-01T00:00:00Z",1088 **({"context_window": ctx} if ctx else {}),1089 } for a in ALIASES]1090 return JSONResponse({"data": data, "has_more": False,1091 "first_id": data[0]["id"], "last_id": data[-1]["id"]})1092 1093 1094@app.get("/v1/models/{model_id}")1095async def model_detail(model_id: str):1096 return JSONResponse({"type": "model", "id": model_id,1097 "display_name": f"{model_id} ({MODEL})",1098 "created_at": "2026-01-01T00:00:00Z"})1099 1100# ------------------------------------------------------------------ passe-plat1101# Le pod n'expose qu'un port. Le pont le prend (c'est l'API que Claude Code1102# consomme) et relaie ces routes vers vLLM, pour que les outils de mesure et les1103# clients OpenAI restent joignables de l'exterieur.1104 1105def _real_model(body: dict) -> dict:1106 """Traduit un alias claude-* vers le nom servi par vLLM.1107 1108 Sans ca, un harnais de mesure qui reprend l'identifiant vu dans /v1/models1109 (`claude-ornith`) recoit un 404 "model not found" de vLLM, releve comme une1110 route absente alors que seul le nom etait inconnu.1111 """1112 m = body.get("model")1113 if isinstance(m, str) and (m in ALIASES or _cle_demandee(m)):1114 body = {**body, "model": _courant()[1]}1115 return body1116 1117 1118@app.post("/v1/chat/completions")1119async def oai_chat(request: Request):1120 body = await request.json()1121 await _assurer_modele(body.get("model"))1122 body = _real_model(body)1123 if body.get("stream"):1124 async def gen():1125 async with _client.stream("POST", "/v1/chat/completions", json=body) as r:1126 async for chunk in r.aiter_raw():1127 yield chunk1128 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 1133 1134@app.post("/v1/completions")1135async def oai_completions(request: Request):1136 body = await request.json()1137 await _assurer_modele(body.get("model"))1138 body = _real_model(body)1139 r = await _client.post("/v1/completions", json=body)1140 return JSONResponse(status_code=r.status_code, content=r.json())1141 1142 1143@app.get("/metrics")1144async def metrics():1145 from fastapi.responses import PlainTextResponse1146 r = await _client.get("/metrics")1147 return PlainTextResponse(r.text, status_code=r.status_code)1148 