Instructions to use patdev/k3-a40-bootstrap with libraries, inference providers, notebooks, and local apps. Follow these links to get started.
- Notebooks
- Google Colab
- Kaggle
- Local Apps Settings
- llama.cpp
How to use patdev/k3-a40-bootstrap with llama.cpp:
Install (macOS, Linux)
curl -LsSf https://llama.app/install.sh | sh # Start a local OpenAI-compatible server with a web UI: llama serve -hf patdev/k3-a40-bootstrap:BF16 # Run inference directly in the terminal: llama cli -hf patdev/k3-a40-bootstrap:BF16
Install from WinGet (Windows)
winget install llama.cpp # Start a local OpenAI-compatible server with a web UI: llama serve -hf patdev/k3-a40-bootstrap:BF16 # Run inference directly in the terminal: llama cli -hf patdev/k3-a40-bootstrap:BF16
Use pre-built binary
# Download pre-built binary from: # https://github.com/ggerganov/llama.cpp/releases # Start a local OpenAI-compatible server with a web UI: ./llama-server -hf patdev/k3-a40-bootstrap:BF16 # Run inference directly in the terminal: ./llama-cli -hf patdev/k3-a40-bootstrap:BF16
Build from source code
git clone https://github.com/ggerganov/llama.cpp.git cd llama.cpp cmake -B build cmake --build build -j --target llama-server llama-cli # Start a local OpenAI-compatible server with a web UI: ./build/bin/llama-server -hf patdev/k3-a40-bootstrap:BF16 # Run inference directly in the terminal: ./build/bin/llama-cli -hf patdev/k3-a40-bootstrap:BF16
Use Docker
docker model run hf.co/patdev/k3-a40-bootstrap:BF16
- LM Studio
- Jan
- Ollama
How to use patdev/k3-a40-bootstrap with Ollama:
ollama run hf.co/patdev/k3-a40-bootstrap:BF16
- Unsloth Desktop
- Docker Model Runner
How to use patdev/k3-a40-bootstrap with Docker Model Runner:
docker model run hf.co/patdev/k3-a40-bootstrap:BF16
- Lemonade
How to use patdev/k3-a40-bootstrap with Lemonade:
Pull the model
# Download Lemonade from https://lemonade-server.ai/ lemonade pull patdev/k3-a40-bootstrap:BF16
Run and chat with the model
lemonade run user.k3-a40-bootstrap-BF16
List all available models
lemonade list
- Atomic Chat
Download anthropic_proxy.py from patdev/k3-a40-bootstrap: direct link, hf CLI and curl.
- Browser
- Download file 52.8 kB
-
https://huggingface.co/patdev/k3-a40-bootstrap/resolve/main/anthropic_proxy.py
- Command line
-
hf download hf://patdev/k3-a40-bootstrap/anthropic_proxy.py
-
curl -L -o anthropic_proxy.py https://huggingface.co/patdev/k3-a40-bootstrap/resolve/main/anthropic_proxy.py
52.8 kB
| """Pont Anthropic -> OpenAI, pour brancher Claude Code sur un serveur vLLM. | |
| Claude Code parle le protocole Anthropic (`POST /v1/messages`, SSE a evenements | |
| nommes). vLLM ne sert que le protocole OpenAI. Ce module traduit dans les deux | |
| sens, en streaming comme en une passe, avec les appels d'outils -- sans quoi un | |
| agent ne peut rien faire. | |
| Lance a cote de vLLM sur la meme machine ; ecoute sur PROXY_PORT et relaie vers | |
| UPSTREAM. | |
| """ | |
| from __future__ import annotations | |
| import asyncio | |
| import json | |
| import os | |
| import re | |
| import time | |
| import uuid | |
| from typing import Any | |
| import httpx | |
| from fastapi import FastAPI, HTTPException, Request | |
| from fastapi.responses import JSONResponse, StreamingResponse | |
| UPSTREAM = os.environ.get("VL_UPSTREAM", "http://127.0.0.1:8080") | |
| MODEL = os.environ.get("VL_SERVED_NAME", "qwen") | |
| # ---------------------------------------------------------------- SWAP (v70) | |
| # Deux modeles a tour de role sur la meme carte. Le bootstrap ecrit | |
| # `$VL_SWAP_DIR/modele_courant` (cle nom depot) a chaque lancement de vLLM ; le | |
| # pont y lit le modele charge A CHAQUE REQUETE (il survit au swap). Une requete | |
| # pour un modele non charge ecrit `modele_demande`, puis attend que le | |
| # bootstrap ait relance vLLM (jusqu'a VL_SWAP_TIMEOUT s). | |
| SWAP_DIR = os.environ.get("VL_SWAP_DIR", "/travail") | |
| SWAP = dict(x.split(":", 1) for x in os.environ.get("VL_SWAP", "").split(",") if ":" in x) | |
| SWAP_TIMEOUT = float(os.environ.get("VL_SWAP_TIMEOUT", "900")) | |
| def _courant() -> tuple[str, str, str, str]: | |
| """(cle bootstrap, nom servi, depot, speculation) du modele charge.""" | |
| try: | |
| with open(os.path.join(SWAP_DIR, "modele_courant"), encoding="utf-8") as f: | |
| champs = f.read().split() | |
| cle, nom, depot = champs[:3] | |
| return cle, nom, depot, (champs[3] if len(champs) > 3 else "off") | |
| except Exception: | |
| return os.environ.get("VL_MODEL_KEY", ""), MODEL, os.environ.get("VL_REAL_MODEL", MODEL), "off" | |
| def _cle_demandee(model: object) -> tuple[str, str] | None: | |
| """`ornith` -> (cle, "") ; `ornith+dspark` -> (cle, "dspark") : la speculation | |
| voyage dans le nom du modele, pour changer de reglage sans recreer le pod.""" | |
| if not isinstance(model, str): | |
| return None | |
| nom = model.removesuffix("[1m]") | |
| if nom.startswith("claude-"): | |
| nom = nom[len("claude-"):] | |
| nom, _, spec = nom.partition("+") | |
| cle = SWAP.get(nom) | |
| return (cle, spec) if cle else None | |
| TIMEOUT = float(os.environ.get("VL_TIMEOUT", "1800")) | |
| # Sortie maximale du modele. Claude Code demande couramment 64 k, ce que | |
| # vLLM refuse d'un 400 portant sur max_tokens -- la requete entiere echoue | |
| # alors qu'un plafonnement silencieux suffit. | |
| MAX_OUTPUT = int(os.environ.get("VL_MAX_OUTPUT", "32768")) | |
| # Claude Code refuse tout identifiant de modele qui ne commence pas par | |
| # "claude-" : il valide le nom avant d'emettre la requete. On expose donc des | |
| # alias conformes, et on ignore le nom recu pour router vers l'unique modele | |
| # reellement charge -- le client choisit une etiquette, pas un moteur. | |
| # Le PREMIER alias nomme le modele reellement charge : sans cela, un client qui | |
| # voit "claude-kimi-k3" croit legitimement executer du Kimi alors que le moteur | |
| # sert du Qwen. Les suivants sont des etiquettes de compatibilite, et tous | |
| # routent vers l'unique modele charge. | |
| _REAL = {"qwen": "claude-qwen3-coder-30b", "kimi": "claude-kimi-linear-48b"} | |
| ALIASES = [ | |
| # Alias NU, sans prefixe. Indispensable pour une fenetre > 200 k. | |
| # Claude Code 2.1.239 (fonction JFd du binaire) : | |
| # let n = CLAUDE_CODE_MAX_CONTEXT_TOKENS; | |
| # if (n !== undefined && n > 0 && !id.startsWith("claude-")) return n; | |
| # return <defaut 200 000> | |
| # Autrement dit il REFUSE toute fenetre personnalisee sur un identifiant | |
| # commencant par "claude-" : avec `claude-flashnext`, ni | |
| # CLAUDE_CODE_MAX_CONTEXT_TOKENS ni CLAUDE_CODE_AUTO_COMPACT_WINDOW n'ont | |
| # le moindre effet, et /context affiche 200k quoi qu'on fasse. | |
| MODEL, | |
| _REAL.get(MODEL, f"claude-{MODEL}"), | |
| "claude-kimi-k3", | |
| "claude-kimi-k3-linear", | |
| "claude-qwen3-coder", | |
| "claude-sonnet-4-5", # alias de compatibilite : certains clients | |
| "claude-3-5-haiku", # codent en dur un modele "rapide" et un "lent" | |
| ] | |
| # Variantes "[1m]". Claude Code deduit la fenetre de contexte du NOM du modele : | |
| # un identifiant qu'il ne connait pas est suppose a 200 k, et l'auto-compactage | |
| # se declenche a 200 k meme si le moteur en accepte 1 000 000. Le suffixe [1m] | |
| # est sa convention pour la fenetre du million ; VERIFIE : avec lui, | |
| # l'avertissement "auto-compact will keep this session within 200k" disparait. | |
| for _n in SWAP: # les modeles interchangeables, tous annonces | |
| ALIASES += [_n, f"claude-{_n}", f"{_n}+dspark", f"{_n}+off"] | |
| ALIASES = [a for pair in ((x, f"{x}[1m]") for x in ALIASES) for a in pair] | |
| ALIASES = list(dict.fromkeys(ALIASES)) # dedoublonne en gardant l'ordre | |
| app = FastAPI(title="anthropic-bridge") | |
| _verrou_swap = asyncio.Lock() | |
| async def _assurer_modele(model: object) -> str: | |
| """Rend le nom servi pour `model`, en declenchant le swap s'il le faut.""" | |
| dem = _cle_demandee(model) | |
| courant = _courant() | |
| if not dem: | |
| return courant[1] | |
| cle, spec = dem | |
| def satisfait(c): | |
| # sans "+spec" dans le nom, n'importe quelle speculation du bon modele convient | |
| return c[0] == cle and (not spec or c[3] == spec) | |
| if satisfait(courant): | |
| return courant[1] | |
| async with _verrou_swap: | |
| courant = _courant() | |
| if satisfait(courant): | |
| return courant[1] | |
| with open(os.path.join(SWAP_DIR, "modele_demande"), "w", encoding="utf-8") as f: | |
| f.write(f"{cle} {spec}\n") | |
| t0 = time.time() | |
| while time.time() - t0 < SWAP_TIMEOUT: | |
| await asyncio.sleep(3) | |
| c = _courant() | |
| if not satisfait(c): | |
| continue | |
| try: | |
| r = await _client.get("/health", timeout=3.0) | |
| if r.status_code == 200: | |
| return c[1] | |
| except Exception: | |
| pass | |
| raise HTTPException(status_code=503, detail=f"swap vers {cle} non termine en {SWAP_TIMEOUT:.0f} s") | |
| async def _guerir_bassin(request: Request, call_next): | |
| generation = _bassin_generation | |
| try: | |
| return await call_next(request) | |
| except httpx.PoolTimeout: | |
| await _reconstruire_bassin(generation) | |
| return JSONResponse( | |
| status_code=503, | |
| content={"type": "error", | |
| "error": {"type": "overloaded_error", | |
| "message": "upstream pool exhausted, " | |
| "connection pool rebuilt; retry"}}) | |
| except (httpx.HTTPError, RuntimeError) as e: | |
| # v73 : mesure du 29/08 -- la saturation ne se presente pas toujours en | |
| # PoolTimeout. Un RuntimeError("client has been closed") apres une | |
| # reconstruction, ou un ReadError sur une connexion recyclee, sortait en | |
| # 500 brut sans jamais declencher la guerison. Meme remede, et le type | |
| # REEL est journalise pour qu'on ne rediagnostique plus a l'aveugle. | |
| import traceback | |
| print(f"[bassin?] {type(e).__name__}: {e} sur {request.url.path}", | |
| flush=True) | |
| traceback.print_exc() | |
| await _reconstruire_bassin(generation) | |
| return JSONResponse( | |
| status_code=503, | |
| content={"type": "error", | |
| "error": {"type": "overloaded_error", | |
| "message": "upstream relay failed " | |
| f"({type(e).__name__}), pool " | |
| "rebuilt; retry"}}) | |
| # Pool EXPLICITE. Mesure 22/08 nuit (5 agents Claude Code, ~150 k de contexte) : | |
| # 189 flux ouverts pour 100 termines, 45 connexions amont pour 2 clients -- les | |
| # flux abandonnes cote client (retry, coupure proxy) gardaient leur connexion | |
| # amont, vLLM generait pour personne, et a 100 connexions (defaut httpx) le pool | |
| # bloquait TOUT, /health compris : pont muet, processus vivant a 3 % CPU. | |
| # Le timeout de pool transforme une saturation en 503 rapide au lieu d'un blocage. | |
| def _nouveau_client() -> httpx.AsyncClient: | |
| return httpx.AsyncClient( | |
| base_url=UPSTREAM, | |
| timeout=httpx.Timeout(TIMEOUT, connect=10.0, pool=10.0), | |
| limits=httpx.Limits(max_connections=512, max_keepalive_connections=64)) | |
| _client = _nouveau_client() | |
| # GUERISON DU BASSIN, ajoutee le 27/08 apres trois saturations en douze heures, | |
| # a intervalle d'environ une heure et quart sous charge multi-agents. | |
| # | |
| # Ce que la mesure a montre, et qui change le diagnostic : au moment du blocage, | |
| # `ss -tan | grep :8000` renvoyait **zero** socket. Le bassin n'est donc pas | |
| # plein de connexions vivantes -- il est plein de PLACES comptabilisees et | |
| # jamais rendues. Consequence : baisser VL_TIMEOUT de 1800 a 240 s n'a rien | |
| # change, parce que la fuite n'est pas temporelle, elle est definitive par flux | |
| # abandonne. Monter `max_connections` ne fait que retarder l'echeance. | |
| # | |
| # On ne repare donc pas la fuite ici (elle est dans httpx/httpcore, pas dans ce | |
| # fichier) : on rend le pont capable d'en sortir seul. A la premiere PoolTimeout | |
| # on remplace le client par un neuf et on ferme l'ancien en arriere-plan. Le | |
| # compteur de generation evite que dix requetes simultanees reconstruisent dix | |
| # fois : seule celle qui a vu la generation courante agit. | |
| _bassin_verrou = asyncio.Lock() | |
| _bassin_generation = 0 | |
| async def _reconstruire_bassin(generation_vue: int) -> None: | |
| global _client, _bassin_generation | |
| async with _bassin_verrou: | |
| if generation_vue != _bassin_generation: | |
| return | |
| ancien, _client = _client, _nouveau_client() | |
| _bassin_generation += 1 | |
| print("[bassin] PoolTimeout -> client httpx reconstruit " | |
| "(generation %d)" % _bassin_generation, flush=True) | |
| try: | |
| await ancien.aclose() | |
| except Exception: | |
| pass | |
| # ------------------------------------------------------------------ garde-fou | |
| # Un modele de 30 a 48 milliards de parametres quantifie en 4 bits abrege : somme | |
| # de reecrire un long fichier, il repond `// ... [previous content] ...` et | |
| # considere le travail fait. Avec --permission-mode acceptEdits, cet abrege est | |
| # ecrit sur le disque sans qu'un diff soit montre : 1841 lignes de source ont | |
| # ainsi disparu. Le modele ne sait pas qu'il a detruit quelque chose, donc il | |
| # rapporte un succes. | |
| # | |
| # C'est un defaut de capacite, pas de configuration : aucun reglage de vLLM ne | |
| # le corrige. Ce qu'on peut faire, en revanche, c'est refuser l'ecriture. Le | |
| # pont voit passer tous les arguments d'outil ; il est le dernier endroit ou | |
| # l'on puisse transformer une destruction silencieuse en refus visible. | |
| GUARD = os.environ.get("VL_GUARD", "on") != "off" | |
| # Les outils dont un argument atterrit tel quel dans un fichier. | |
| # Seules les ecritures de fichier ENTIER sont gardees. Les editions ciblees | |
| # (Edit, MultiEdit, str_replace_*, NotebookEdit) portent legitimement des | |
| # "..." dans old_string/new_string -- et le message du garde leur disait | |
| # justement d'utiliser une edition ciblee : boucle de refus, observee le 22/08. | |
| _GUARDED = {"write", "write_file", "create_file"} | |
| # Marqueurs d'omission EXPLICITES seulement. Les motifs generiques "[...]" et | |
| # "# ..." bloquaient des contenus legitimes (un .md avec une ligne "# ...", un | |
| # script contenant "[...]") : faux positifs observes le 22/08 sur une vraie | |
| # session Claude Code. | |
| PLACEHOLDER = re.compile( | |
| r"\.\.\.\s*(\[|#|//|--)?\s*(previous|rest of|remaining|existing|unchanged|original)" | |
| r"|\[\s*\.\.\.\s*(previous|rest|remaining|existing|unchanged|original|snip)[^\]]*\]" | |
| r"|<\s*(unchanged|snip|elided)\s*>" | |
| r"|\(\s*(reste|suite) (du|des) ", | |
| re.I | re.M, | |
| ) | |
| # Rappel injecte en tete de systeme. Ne remplace pas le garde-fou -- un modele | |
| # qui abrege le fait souvent malgre la consigne -- mais reduit la frequence, et | |
| # ne coute qu'une constante en tete de prefixe, donc le cache reste partage. | |
| SYSTEM_GUARD = ( | |
| "Quand tu ecris un fichier, l'argument `content` doit contenir le fichier " | |
| "INTEGRAL. N'ecris jamais de marqueur d'omission du type " | |
| "\"... [previous content] ...\", \"// ... rest of file ...\" ou " | |
| "\"(reste du fichier inchange)\" : ces marqueurs sont ecrits litteralement " | |
| "sur le disque et detruisent le fichier. Si le fichier est trop long pour " | |
| "etre reemis en entier, dis-le et utilise une edition ciblee." | |
| ) | |
| def guard_violation(name: str, args: Any) -> str | None: | |
| """Renvoie la description de l'abreviation trouvee, ou None.""" | |
| if not GUARD or (name or "").lower() not in _GUARDED: | |
| return None | |
| def walk(v: Any, path: str) -> str | None: | |
| if isinstance(v, str): | |
| m = PLACEHOLDER.search(v) | |
| return f"{path} contient {m.group(0).strip()!r}" if m else None | |
| if isinstance(v, dict): | |
| for k, sub in v.items(): | |
| hit = walk(sub, f"{path}.{k}") | |
| if hit: | |
| return hit | |
| elif isinstance(v, list): | |
| for i, sub in enumerate(v): | |
| hit = walk(sub, f"{path}[{i}]") | |
| if hit: | |
| return hit | |
| return None | |
| return walk(args, name) | |
| def guard_message(detail: str) -> str: | |
| return ( | |
| "\n\n[garde-fou du pont] Appel d'outil bloque : " + detail + ".\n" | |
| "Le contenu propose abrege le fichier par un marqueur d'omission, ce qui " | |
| "l'aurait ecrase par une version incomplete. L'ecriture n'a PAS eu lieu et " | |
| "le fichier est intact.\n" | |
| "Reemets le fichier integral, ou fais une edition ciblee qui ne remplace " | |
| "que les lignes concernees." | |
| ) | |
| # --------------------------------------------------------------- Anthropic -> OpenAI | |
| def _image_part(block: dict) -> dict | None: | |
| """Traduit un bloc `image` Anthropic vers la partie `image_url` d'OpenAI. | |
| Anthropic decrit l'image par une `source` typee : `base64` porte les octets | |
| et le type MIME separement, `url` porte un lien. OpenAI attend dans les deux | |
| cas UNE chaine dans `image_url.url` -- une URI de donnees pour le premier, | |
| le lien tel quel pour le second. | |
| Sans cette traduction, `_text_of` ignorait purement et simplement les blocs | |
| image : une capture collee dans Claude Code arrivait au modele comme un | |
| message vide, et le modele repondait a cote sans que rien ne signale la | |
| perte. | |
| """ | |
| src = block.get("source") or {} | |
| kind = src.get("type") | |
| if kind == "base64": | |
| data = src.get("data") | |
| if not data: | |
| return None | |
| mime = src.get("media_type") or "image/png" | |
| return {"type": "image_url", | |
| "image_url": {"url": f"data:{mime};base64,{data}"}} | |
| if kind == "url" and src.get("url"): | |
| return {"type": "image_url", "image_url": {"url": src["url"]}} | |
| return None | |
| def _text_of(content: Any) -> str: | |
| """Anthropic autorise une chaine ou une liste de blocs typés.""" | |
| if isinstance(content, str): | |
| return content | |
| if not isinstance(content, list): | |
| return "" | |
| out = [] | |
| for b in content: | |
| if isinstance(b, str): | |
| out.append(b) | |
| elif isinstance(b, dict) and b.get("type") == "text": | |
| out.append(b.get("text", "")) | |
| return "".join(out) | |
| def to_openai(body: dict) -> dict: | |
| msgs: list[dict] = [] | |
| # `system` est un champ separe chez Anthropic, un message de role chez OpenAI. | |
| # Le rappel anti-abreviation n'est ajoute qu'en presence d'outils : sans | |
| # outil, aucun contenu n'atteint le disque et la consigne serait du bruit. | |
| sys = _text_of(body.get("system") or "") | |
| if GUARD and body.get("tools"): | |
| sys = (sys + "\n\n" + SYSTEM_GUARD) if sys else SYSTEM_GUARD | |
| if sys: | |
| msgs.append({"role": "system", "content": sys}) | |
| for m in body.get("messages", []): | |
| role = m.get("role", "user") | |
| content = m.get("content") | |
| # Claude Code place un SECOND message systeme en fin de conversation | |
| # (rappel de contexte). Le gabarit Qwen n'accepte le role systeme qu'en | |
| # tete et vLLM refuse toute la requete : | |
| # 400 "System message must be at the beginning." | |
| # Le pont ne verifiait pas le statut amont, donc l'echec ressortait en | |
| # reponse VIDE avec stop_reason=end_turn -- Claude Code affichait | |
| # "Cogitated for 0s" et rien d'autre. | |
| # | |
| # On le convertit en message utilisateur plutot que de le fusionner dans | |
| # le systeme de tete : ce rappel change a chaque tour, et le remonter en | |
| # tete invaliderait le prefixe partage, donc le cache -- soit 68 % des | |
| # blocs, mesures. | |
| if role == "system" and msgs and msgs[0].get("role") == "system": | |
| txt = _text_of(content) | |
| if txt: | |
| msgs.append({"role": "user", "content": txt}) | |
| continue | |
| if isinstance(content, list): | |
| # Un tour d'assistant peut melanger du texte et des tool_use ; un tour | |
| # d'utilisateur porte les tool_result. OpenAI separe les deux en | |
| # `tool_calls` sur l'assistant et en messages de role `tool`. | |
| texts, calls, results, images, thinks = [], [], [], [], [] | |
| for b in content: | |
| if not isinstance(b, dict): | |
| continue | |
| t = b.get("type") | |
| if t == "text": | |
| texts.append(b.get("text", "")) | |
| elif t == "thinking": | |
| # Claude Code renvoie le raisonnement du tour precedent dans | |
| # l'historique (boucles d'outils). Le gabarit Qwen3.5 sait le | |
| # garder pour le DERNIER tour d'assistant via | |
| # `reasoning_content` -- exactement la semantique Anthropic. | |
| # Avant : bloc ignore en silence, le modele perdait son plan | |
| # entre deux appels d'outil. | |
| thinks.append(b.get("thinking", "")) | |
| elif t == "image": | |
| part = _image_part(b) | |
| if part: | |
| images.append(part) | |
| elif t == "tool_use": | |
| calls.append({ | |
| "id": b.get("id") or f"call_{uuid.uuid4().hex[:8]}", | |
| "type": "function", | |
| "function": { | |
| "name": b.get("name", ""), | |
| "arguments": json.dumps(b.get("input") or {}), | |
| }, | |
| }) | |
| elif t == "tool_result": | |
| # Un resultat d'outil peut porter des images (capture rendue | |
| # par un outil). Le role `tool` d'OpenAI n'accepte que du | |
| # texte : on extrait les images pour les rattacher au tour | |
| # utilisateur, sinon elles disparaissent en silence. | |
| rc = b.get("content") | |
| if isinstance(rc, list): | |
| for sub in rc: | |
| if isinstance(sub, dict) and sub.get("type") == "image": | |
| part = _image_part(sub) | |
| if part: | |
| images.append(part) | |
| results.append({ | |
| "role": "tool", | |
| "tool_call_id": b.get("tool_use_id", ""), | |
| "content": _text_of(rc) or "", | |
| }) | |
| if role == "assistant": | |
| a: dict[str, Any] = {"role": "assistant", "content": "".join(texts) or None} | |
| if thinks: | |
| a["reasoning_content"] = chr(10).join(t for t in thinks if t) | |
| if calls: | |
| a["tool_calls"] = calls | |
| msgs.append(a) | |
| else: | |
| if images: | |
| # Contenu multipart : OpenAI n'accepte les images que dans | |
| # une LISTE de parties, jamais dans une chaine. | |
| parts: list[dict] = [] | |
| joined = "".join(texts) | |
| if joined: | |
| parts.append({"type": "text", "text": joined}) | |
| parts.extend(images) | |
| msgs.append({"role": "user", "content": parts}) | |
| elif texts: | |
| msgs.append({"role": "user", "content": "".join(texts)}) | |
| msgs.extend(results) | |
| else: | |
| msgs.append({"role": role, "content": content or ""}) | |
| out: dict[str, Any] = { | |
| "model": _courant()[1], | |
| "messages": msgs, | |
| "max_tokens": min(int(body.get("max_tokens") or 4096), MAX_OUTPUT), | |
| "stream": bool(body.get("stream")), | |
| } | |
| for src, dst in (("temperature", "temperature"), ("top_p", "top_p"), | |
| ("stop_sequences", "stop")): | |
| if body.get(src) is not None: | |
| out[dst] = body[src] | |
| # Le raisonnement est ACTIF par defaut sur ce modele et consomme des jetons | |
| # avant le premier caractere de reponse : une requete a max_tokens=40 revient | |
| # avec content vide et stop_reason=max_tokens. Le protocole Anthropic exprime | |
| # la coupure par `thinking: {"type": "disabled"}` ; le gabarit Qwen l'attend | |
| # sous la forme enable_thinking=false, qui fait prefixer un bloc <think> deja | |
| # ferme. Sans cette traduction, le champ etait recu puis ignore en silence : | |
| # la case "raisonnement" de la console ne changeait rien. | |
| # `thinking` a plusieurs formes. La documentation du protocole de passerelle | |
| # precise que Claude Code envoie `{"type": "adaptive"}` aux modeles recents | |
| # ET "traite les noms de modeles qu'il ne reconnait pas, tels les alias de | |
| # passerelle, comme des modeles actuels qui recoivent le champ" -- donc nous. | |
| # Seul "disabled" doit couper le raisonnement ; "adaptive" et "enabled" le | |
| # laissent actif, qui est le defaut du gabarit. | |
| th = body.get("thinking") | |
| if isinstance(th, dict) and th.get("type") == "disabled": | |
| out["chat_template_kwargs"] = {"enable_thinking": False} | |
| elif isinstance(th, dict) and th.get("budget_tokens"): | |
| # Le niveau de raisonnement de Claude Code arrive ici, en jetons. vLLM | |
| # 0.27.1 l'applique via `thinking_token_budget` (sampling params) : il | |
| # ferme le bloc de raisonnement a ce nombre de jetons. Mesure : sans | |
| # ce champ, budget 512 et budget 8000 donnaient la meme chose. | |
| try: | |
| out["thinking_token_budget"] = max(1, int(th["budget_tokens"])) | |
| except (TypeError, ValueError): | |
| pass | |
| # Niveau d'effort (output_config.effort) -> plafond de jetons de raisonnement, | |
| # memes seuils que le patch natif (vllm_anthropic_effort_patch.py). Le budget | |
| # explicite ci-dessus gagne ; `disabled` aussi. | |
| _eff = (body.get("output_config") or {}).get("effort") if isinstance(body.get("output_config"), dict) else None | |
| _BUDGET = {"low": 1024, "medium": 4096, "high": 16384, "xhigh": 32768, "max": None} | |
| if _eff in _BUDGET and "thinking_token_budget" not in out and "chat_template_kwargs" not in out: | |
| if _BUDGET[_eff] is not None: | |
| out["thinking_token_budget"] = _BUDGET[_eff] | |
| if body.get("tools"): | |
| out["tools"] = [{ | |
| "type": "function", | |
| "function": { | |
| "name": t["name"], | |
| "description": t.get("description", ""), | |
| "parameters": t.get("input_schema") or {"type": "object", "properties": {}}, | |
| }, | |
| } for t in body["tools"] if t.get("name")] | |
| tc = body.get("tool_choice") or {} | |
| kind = tc.get("type") if isinstance(tc, dict) else None | |
| if kind == "none": | |
| out["tool_choice"] = "none" | |
| elif kind == "any": | |
| out["tool_choice"] = "required" | |
| elif kind == "tool" and tc.get("name"): | |
| out["tool_choice"] = {"type": "function", "function": {"name": tc["name"]}} | |
| else: | |
| out["tool_choice"] = "auto" | |
| return out | |
| _STOP = {"stop": "end_turn", "length": "max_tokens", "tool_calls": "tool_use"} | |
| # --------------------------------------------------------------- OpenAI -> Anthropic | |
| def _usage_of(usage: dict) -> dict: | |
| """vLLM expose les jetons servis par le cache de prefixe dans | |
| `prompt_tokens_details.cached_tokens` ; Anthropic les appelle | |
| `cache_read_input_tokens`. Sans cette traduction Claude Code affichait | |
| 0 % de cache alors que vLLM en servait 68 %.""" | |
| u = {"input_tokens": usage.get("prompt_tokens", 0) or 0, | |
| "output_tokens": usage.get("completion_tokens", 0) or 0} | |
| det = usage.get("prompt_tokens_details") or {} | |
| cached = det.get("cached_tokens") if isinstance(det, dict) else None | |
| if cached: | |
| u["cache_read_input_tokens"] = int(cached) | |
| u["cache_creation_input_tokens"] = 0 | |
| return u | |
| def _stop_of(choice: dict) -> tuple[str, str | None]: | |
| """finish_reason=stop couvre deux cas Anthropic : end_turn, ou | |
| stop_sequence quand vLLM rapporte la chaine d'arret dans `stop_reason`.""" | |
| fr = choice.get("finish_reason") or "stop" | |
| sr = choice.get("stop_reason") | |
| if fr == "stop" and isinstance(sr, str) and sr: | |
| return "stop_sequence", sr | |
| return _STOP.get(fr, "end_turn"), None | |
| def to_anthropic(oai: dict, req_model: str) -> dict: | |
| choice = (oai.get("choices") or [{}])[0] | |
| msg = choice.get("message") or {} | |
| blocks: list[dict] = [] | |
| # Le raisonnement arrive ici en non-flux ; en flux il est deja traduit en | |
| # bloc thinking. Meme forme dans les deux cas, signature comprise. | |
| reasoning = msg.get("reasoning_content") or msg.get("reasoning") | |
| if reasoning: | |
| blocks.append({"type": "thinking", "thinking": reasoning, | |
| "signature": "vllm-" + uuid.uuid4().hex[:16]}) | |
| if msg.get("content"): | |
| blocks.append({"type": "text", "text": msg["content"]}) | |
| elif reasoning and not msg.get("tool_calls"): | |
| # meme repli qu'en flux : pas de texte, pas d'outil -> le raisonnement | |
| blocks.append({"type": "text", "text": str(reasoning).strip()}) | |
| blocked = False | |
| for c in msg.get("tool_calls") or []: | |
| fn = c.get("function") or {} | |
| try: | |
| args = json.loads(fn.get("arguments") or "{}") | |
| except json.JSONDecodeError: | |
| # Un modele quantifie peut emettre du JSON legerement casse ; mieux | |
| # vaut transmettre la chaine brute que faire tomber la requete. | |
| args = {"_raw": fn.get("arguments", "")} | |
| name = fn.get("name", "") | |
| detail = guard_violation(name, args) | |
| if detail: | |
| blocked = True | |
| blocks.append({"type": "text", "text": guard_message(detail)}) | |
| continue | |
| blocks.append({ | |
| "type": "tool_use", | |
| "id": c.get("id") or f"toolu_{uuid.uuid4().hex[:16]}", | |
| "name": name, | |
| "input": args, | |
| }) | |
| usage = oai.get("usage") or {} | |
| if blocked and not any(b["type"] == "tool_use" for b in blocks): | |
| # Plus aucun outil a executer : le tour se termine sur l'explication. | |
| return { | |
| "id": oai.get("id") or f"msg_{uuid.uuid4().hex[:24]}", | |
| "type": "message", "role": "assistant", "model": req_model, | |
| "content": blocks, "stop_reason": "end_turn", "stop_sequence": None, | |
| "usage": _usage_of(usage), | |
| } | |
| stop_reason, stop_seq = _stop_of(choice) | |
| return { | |
| "id": oai.get("id") or f"msg_{uuid.uuid4().hex[:24]}", | |
| "type": "message", | |
| "role": "assistant", | |
| "model": req_model, | |
| "content": blocks, | |
| "stop_reason": stop_reason, | |
| "stop_sequence": stop_seq, | |
| "usage": _usage_of(usage), | |
| } | |
| def _err_of(body: str, status: int) -> dict: | |
| """Renvoie l'objet d'erreur amont intact, ou en fabrique un equivalent. | |
| Le libelle compte : c'est sur lui que le client decide s'il peut retenter. | |
| """ | |
| try: | |
| d = json.loads(body) | |
| e = d.get("error") | |
| if isinstance(e, dict) and e.get("message"): | |
| return {"type": e.get("type") or "api_error", "message": e["message"]} | |
| except (json.JSONDecodeError, AttributeError): | |
| pass | |
| return {"type": "api_error", "message": body or f"upstream {status}"} | |
| def _sse(event: str, data: dict) -> bytes: | |
| return f"event: {event}\ndata: {json.dumps(data, separators=(',', ':'))}\n\n".encode() | |
| async def _flux_gueri(gen): | |
| """v73 : les exceptions nees dans un generateur SSE ne traversent pas le | |
| middleware _guerir_bassin. On les attrape ici : evenement `error` au client | |
| (sa logique de reprise lit le libelle) et reconstruction du bassin.""" | |
| generation = _bassin_generation | |
| try: | |
| async for morceau in gen: | |
| yield morceau | |
| except (httpx.HTTPError, RuntimeError) as e: | |
| import traceback | |
| print(f"[bassin?] flux: {type(e).__name__}: {e}", flush=True) | |
| traceback.print_exc() | |
| asyncio.ensure_future(_reconstruire_bassin(generation)) | |
| yield _sse("error", { | |
| "type": "error", | |
| "error": {"type": "overloaded_error", | |
| "message": f"upstream stream failed " | |
| f"({type(e).__name__}); retry"}}) | |
| async def stream_anthropic(payload: dict, req_model: str, request: Request | None = None): | |
| """Traduit le flux OpenAI en evenements Anthropic. | |
| Le point delicat est l'indexation des blocs : Anthropic numerote chaque bloc | |
| de contenu et exige un content_block_start/stop apparie, alors qu'OpenAI | |
| emet des deltas plats. On ouvre donc un bloc texte a la volee, et un bloc | |
| tool_use par appel, en fermant le precedent. | |
| Les arguments d'outil, eux, sont accumules et n'emis qu'une fois complets. | |
| Il le faut : un argument diffuse fragment par fragment ne peut pas etre | |
| valide, et le client aurait deja commence a ecrire le fichier quand le | |
| marqueur d'omission apparait. | |
| Mais accumuler veut dire n'envoyer AUCUN octet pendant toute la generation | |
| de l'appel, et le proxy de Runpod coupe une connexion inactive vers 125 s. | |
| MESURE : une reecriture de 1354 lignes par l'outil Write mourait a 126 s | |
| avec zero token et stop=None, alors que la meme reecriture en prose passait | |
| (le texte, lui, est diffuse au fil de l'eau). Un test voisin est passe de | |
| justesse a 123,6 s. On emet donc un `ping` periodique tant qu'on accumule : | |
| le protocole Anthropic le prevoit, les clients l'ignorent, et la connexion | |
| reste ouverte sans que le garde-fou perde sa capacite a valider avant | |
| emission. | |
| """ | |
| mid = f"msg_{uuid.uuid4().hex[:24]}" | |
| yield _sse("message_start", { | |
| "type": "message_start", | |
| "message": {"id": mid, "type": "message", "role": "assistant", | |
| "model": req_model, "content": [], "stop_reason": None, | |
| "stop_sequence": None, | |
| "usage": {"input_tokens": 0, "output_tokens": 0}}, | |
| }) | |
| yield _sse("ping", {"type": "ping"}) | |
| idx = -1 | |
| text_open = False | |
| think_open = False | |
| n_text = 0 | |
| n_think = 0 | |
| think_buf: list[str] = [] # raisonnement accumule, pour le repli "reponse vide" | |
| last_usage_sent = -1 | |
| last_usage_ts = 0.0 | |
| # Horodatage du dernier octet REELLEMENT envoye au client, pour savoir quand | |
| # le silence devient dangereux. | |
| last_out = time.monotonic() | |
| KEEPALIVE_S = 15.0 | |
| tools: dict[int, dict] = {} # index OpenAI -> {id, name, args} | |
| stop_reason = "end_turn" | |
| out_tokens = 0 | |
| in_tokens = 0 | |
| cached_tokens = 0 | |
| stop_seq: str | None = None | |
| async with _client.stream("POST", "/v1/chat/completions", json=payload) as r: | |
| # Un 4xx amont renvoie du JSON d'erreur, pas des lignes "data:". Sans ce | |
| # controle, chaque ligne est ignoree, le flux se termine vide et le | |
| # client recoit un message sans contenu avec stop_reason=end_turn : une | |
| # panne parfaitement silencieuse, indiscernable d'un modele muet. | |
| if r.status_code >= 400: | |
| body = (await r.aread()).decode("utf-8", "ignore")[:800] | |
| roles = "/".join(m.get("role", "?") for m in payload.get("messages", [])) | |
| print(f"[amont {r.status_code}] {body}", flush=True) | |
| print(f"[amont] roles={roles}", flush=True) | |
| # Le corps d'erreur repart TEL QUEL, dans un evenement SSE `error`. | |
| # La documentation du protocole de passerelle est explicite : Claude | |
| # Code se remet tout seul de certains refus -- champ `thinking`, | |
| # signatures de raisonnement, et justement les messages systeme en | |
| # milieu de conversation -- mais "la logique de reprise s'appuie sur | |
| # le libelle de l'erreur amont", et "une passerelle qui enveloppe | |
| # les erreurs dans son propre format casse la reprise meme si elle | |
| # preserve le code de statut". Emballer l'erreur dans un bloc texte, | |
| # comme je le faisais, empechait donc cette reprise. | |
| yield _sse("error", {"type": "error", "error": _err_of(body, r.status_code)}) | |
| return | |
| # Le `ping` periodique ne suffisait pas : il ne se declenchait qu'a | |
| # l'interieur de cette boucle, donc uniquement quand vLLM emettait deja | |
| # quelque chose. Or pendant le PREREMPLISSAGE vLLM n'envoie rien -- et | |
| # un prompt de plusieurs centaines de milliers de jetons met plusieurs | |
| # minutes a etre calcule. Le flux restait muet, et le proxy Runpod | |
| # coupait vers 125 s. | |
| # | |
| # MESURE : une aiguille a ~700 k jetons echouait a 126,5 s avec zero | |
| # jeton, alors que la meme a ~358 k passait en 94,4 s. Le plafond | |
| # n'etait pas le modele mais le silence : aucune requete au-dela de | |
| # ~500 k ne pouvait aboutir, quel que soit le contexte annonce. | |
| # | |
| # On decouple donc la lecture amont de l'emission aval : une tache | |
| # pompe les lignes dans une file, et l'absence de ligne pendant | |
| # KEEPALIVE_S produit un `ping` au lieu d'un silence. | |
| file: asyncio.Queue = asyncio.Queue() | |
| async def _pompe() -> None: | |
| try: | |
| async for _l in r.aiter_lines(): | |
| await file.put(_l) | |
| finally: | |
| await file.put(None) | |
| _tache = asyncio.create_task(_pompe()) | |
| try: | |
| while True: | |
| try: | |
| line = await asyncio.wait_for(file.get(), timeout=KEEPALIVE_S) | |
| except asyncio.TimeoutError: | |
| # Client parti pendant le prefill ? On ferme l'amont (vLLM | |
| # annule alors la generation) au lieu de pomper dans le vide. | |
| if request is not None and await request.is_disconnected(): | |
| print("[deconnexion] client parti pendant l'attente, amont ferme", flush=True) | |
| break | |
| yield _sse("ping", {"type": "ping"}) | |
| last_out = time.monotonic() | |
| continue | |
| if line is None: | |
| break | |
| # Le `yield` ne leve pas toujours a la deconnexion (le proxy Runpod | |
| # garde parfois sa connexion ouverte) : on verifie explicitement. | |
| if request is not None and (out_tokens % 32 == 0) and await request.is_disconnected(): | |
| print(f"[deconnexion] client parti apres {out_tokens} jetons, amont ferme", flush=True) | |
| break | |
| if not line.startswith("data: "): | |
| continue | |
| chunk = line[6:].strip() | |
| if chunk == "[DONE]": | |
| break | |
| try: | |
| d = json.loads(chunk) | |
| except json.JSONDecodeError: | |
| continue | |
| if d.get("usage"): | |
| in_tokens = d["usage"].get("prompt_tokens", in_tokens) | |
| out_tokens = d["usage"].get("completion_tokens", out_tokens) | |
| _det = d["usage"].get("prompt_tokens_details") or {} | |
| if isinstance(_det, dict) and _det.get("cached_tokens"): | |
| cached_tokens = int(_det["cached_tokens"]) | |
| ch = (d.get("choices") or [{}])[0] | |
| delta = ch.get("delta") or {} | |
| # Comptage en direct facon Anthropic : vLLM (continuous_usage_stats) | |
| # renvoie le total cumule sur chaque chunk (lu ci-dessus). On le | |
| # relaie dans un message_delta periodique pour que le compteur du | |
| # client monte en continu au lieu de sauter au total a la fin. | |
| if out_tokens and out_tokens != last_usage_sent and (time.monotonic() - last_usage_ts) >= 0.5: | |
| yield _sse("message_delta", {"type": "message_delta", | |
| "delta": {"stop_reason": None, "stop_sequence": None}, | |
| "usage": {"input_tokens": in_tokens, "output_tokens": out_tokens}}) | |
| last_usage_sent = out_tokens; last_usage_ts = time.monotonic(); last_out = time.monotonic() | |
| # vLLM range le raisonnement dans un champ separe. Sans ce relais, | |
| # il est simplement perdu : le client ne peut ni l'afficher ni savoir | |
| # pourquoi la reponse tarde. | |
| if (delta.get("reasoning_content") or delta.get("reasoning")): | |
| if not think_open: | |
| if text_open: | |
| yield _sse("content_block_stop", | |
| {"type": "content_block_stop", "index": idx}) | |
| text_open = False | |
| idx += 1 | |
| think_open = True | |
| yield _sse("content_block_start", { | |
| "type": "content_block_start", "index": idx, | |
| "content_block": {"type": "thinking", "thinking": ""}}) | |
| yield _sse("content_block_delta", { | |
| "type": "content_block_delta", "index": idx, | |
| "delta": {"type": "thinking_delta", | |
| "thinking": (delta.get("reasoning_content") or delta.get("reasoning"))}}) | |
| n_think += len((delta.get("reasoning_content") or delta.get("reasoning"))) | |
| think_buf.append(delta.get("reasoning_content") or delta.get("reasoning")) | |
| last_out = time.monotonic() | |
| if delta.get("content"): | |
| if think_open: | |
| yield _sse("content_block_delta", { | |
| "type": "content_block_delta", "index": idx, | |
| "delta": {"type": "signature_delta", | |
| "signature": "vllm-" + uuid.uuid4().hex[:16]}}) | |
| yield _sse("content_block_stop", | |
| {"type": "content_block_stop", "index": idx}) | |
| think_open = False | |
| if not text_open: | |
| idx += 1 | |
| text_open = True | |
| yield _sse("content_block_start", { | |
| "type": "content_block_start", "index": idx, | |
| "content_block": {"type": "text", "text": ""}}) | |
| yield _sse("content_block_delta", { | |
| "type": "content_block_delta", "index": idx, | |
| "delta": {"type": "text_delta", "text": delta["content"]}}) | |
| n_text += len(delta["content"]) | |
| last_out = time.monotonic() | |
| for tc in delta.get("tool_calls") or []: | |
| i = tc.get("index", 0) | |
| fn = tc.get("function") or {} | |
| t = tools.setdefault(i, {"id": None, "name": "", "args": ""}) | |
| if tc.get("id"): | |
| t["id"] = tc["id"] | |
| if fn.get("name"): | |
| t["name"] = fn["name"] | |
| if fn.get("arguments"): | |
| t["args"] += fn["arguments"] | |
| # Le silence de l'accumulation est ce qui tuait la connexion. | |
| if time.monotonic() - last_out > KEEPALIVE_S: | |
| yield _sse("ping", {"type": "ping"}) | |
| last_out = time.monotonic() | |
| if ch.get("finish_reason"): | |
| stop_reason, stop_seq = _stop_of(ch) | |
| finally: | |
| _tache.cancel() | |
| if think_open: | |
| # Structure Anthropic complete : la signature precede la fermeture du | |
| # bloc. Claude Code la renvoie telle quelle dans l'historique ; le pont | |
| # ne la verifie pas, mais un client strict l'attend. | |
| yield _sse("content_block_delta", { | |
| "type": "content_block_delta", "index": idx, | |
| "delta": {"type": "signature_delta", | |
| "signature": "vllm-" + uuid.uuid4().hex[:16]}}) | |
| yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx}) | |
| think_open = False | |
| if text_open: | |
| yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx}) | |
| text_open = False | |
| emitted = 0 | |
| for i in sorted(tools): | |
| t = tools[i] | |
| try: | |
| args = json.loads(t["args"] or "{}") | |
| except json.JSONDecodeError: | |
| args = {"_raw": t["args"]} | |
| detail = guard_violation(t["name"], args) | |
| if detail: | |
| print(f"[garde-fou] appel bloque : {detail}", flush=True) | |
| idx += 1 | |
| yield _sse("content_block_start", { | |
| "type": "content_block_start", "index": idx, | |
| "content_block": {"type": "text", "text": ""}}) | |
| yield _sse("content_block_delta", { | |
| "type": "content_block_delta", "index": idx, | |
| "delta": {"type": "text_delta", "text": guard_message(detail)}}) | |
| yield _sse("content_block_stop", | |
| {"type": "content_block_stop", "index": idx}) | |
| continue | |
| idx += 1 | |
| emitted += 1 | |
| yield _sse("content_block_start", { | |
| "type": "content_block_start", "index": idx, | |
| "content_block": {"type": "tool_use", | |
| "id": t["id"] or f"toolu_{uuid.uuid4().hex[:16]}", | |
| "name": t["name"], "input": {}}}) | |
| # Un seul fragment : l'argument est deja complet et valide. | |
| yield _sse("content_block_delta", { | |
| "type": "content_block_delta", "index": idx, | |
| "delta": {"type": "input_json_delta", | |
| "partial_json": json.dumps(args, separators=(",", ":"))}}) | |
| yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx}) | |
| # Annoncer `tool_use` sans avoir emis d'outil ferait attendre au client un | |
| # resultat qui ne viendra jamais. | |
| if stop_reason == "tool_use" and emitted == 0: | |
| stop_reason = "end_turn" | |
| # Repli "reponse vide" (22/08 nuit) : la compaction automatique de Claude | |
| # Code a echoue sur "summarization produced empty response" -- le modele | |
| # avait tout ecrit dans le raisonnement (1 679 car.) et rien en texte. Un | |
| # client qui attend du texte ne sait rien faire d'un bloc thinking seul : | |
| # on rend alors le raisonnement comme texte, clairement marque. | |
| if n_text == 0 and emitted == 0 and "".join(think_buf).strip(): | |
| idx += 1 | |
| yield _sse("content_block_start", {"type": "content_block_start", "index": idx, | |
| "content_block": {"type": "text", "text": ""}}) | |
| yield _sse("content_block_delta", {"type": "content_block_delta", "index": idx, | |
| "delta": {"type": "text_delta", "text": "".join(think_buf).strip()}}) | |
| yield _sse("content_block_stop", {"type": "content_block_stop", "index": idx}) | |
| n_text = len("".join(think_buf).strip()) | |
| print("[repli] reponse sans texte ni outil : raisonnement rendu en texte", flush=True) | |
| # Trace compacte d'une requete. Sans elle, un client qui n'affiche rien est | |
| # indiscernable d'un modele qui ne repond rien : les deux donnent un 200 OK | |
| # dans le journal d'acces. | |
| if os.environ.get("VL_TRACE", "on") != "off": | |
| print(f"[trace] entree={in_tokens} sortie={out_tokens} " | |
| f"texte={n_text}c raisonnement={n_think}c outils={emitted} " | |
| f"stop={stop_reason}", flush=True) | |
| _u = {"input_tokens": in_tokens, "output_tokens": out_tokens} | |
| if cached_tokens: | |
| _u["cache_read_input_tokens"] = cached_tokens | |
| _u["cache_creation_input_tokens"] = 0 | |
| yield _sse("message_delta", { | |
| "type": "message_delta", | |
| "delta": {"stop_reason": stop_reason, "stop_sequence": stop_seq}, | |
| "usage": _u}) | |
| yield _sse("message_stop", {"type": "message_stop"}) | |
| # ------------------------------------------------------------------------ routes | |
| async def messages(request: Request): | |
| body = await request.json() | |
| req_model = body.get("model", MODEL) | |
| payload = to_openai(body) | |
| if payload.get("stream"): | |
| payload["stream_options"] = {"include_usage": True, "continuous_usage_stats": True} | |
| async def flux_avec_swap(): | |
| # Le proxy RunPod coupe a ~120 s toute reponse sans octet ; un swap | |
| # de modele en dure 2 a 6. On tient la ligne avec des pings SSE | |
| # pendant l'attente, puis on enchaine sur le flux normal. | |
| tache = asyncio.ensure_future(_assurer_modele(req_model)) | |
| while not tache.done(): | |
| await asyncio.wait({tache}, timeout=15) | |
| if not tache.done(): | |
| yield _sse("ping", {"type": "ping"}) | |
| try: | |
| payload["model"] = tache.result() | |
| except HTTPException as e: | |
| yield _sse("error", {"type": "error", | |
| "error": {"type": "overloaded_error", "message": str(e.detail)}}) | |
| return | |
| async for morceau in stream_anthropic(payload, req_model, request): | |
| yield morceau | |
| return StreamingResponse(_flux_gueri(flux_avec_swap()), | |
| media_type="text/event-stream") | |
| payload["model"] = await _assurer_modele(req_model) | |
| r = await _client.post("/v1/chat/completions", json=payload) | |
| if r.status_code != 200: | |
| return JSONResponse(status_code=r.status_code, | |
| content={"type": "error", | |
| "error": {"type": "api_error", "message": r.text[:800]}}) | |
| return JSONResponse(to_anthropic(r.json(), req_model)) | |
| async def count_tokens(request: Request): | |
| """Claude Code interroge ce point avant d'envoyer. Une estimation suffit : | |
| il s'en sert pour decider de compacter, pas pour facturer.""" | |
| body = await request.json() | |
| # Compte REEL par vLLM (/tokenize en forme chat : gabarit, systeme et | |
| # outils compris). L'estimation len/4 ignorait systeme et outils : 13 | |
| # jetons pour un appel qui en pesait bien plus. Repli sur l'estimation si | |
| # le tokenizer amont ne repond pas. | |
| try: | |
| payload = to_openai(body) | |
| req = {"model": payload.get("model"), "messages": payload.get("messages", []), | |
| "add_generation_prompt": True} | |
| if payload.get("tools"): | |
| req["tools"] = payload["tools"] | |
| r = await _client.post("/tokenize", json=req, timeout=60) | |
| if r.status_code == 200: | |
| n = r.json().get("count") | |
| if isinstance(n, int) and n > 0: | |
| return JSONResponse({"input_tokens": n}) | |
| except Exception: # noqa: BLE001 | |
| pass | |
| n = len(json.dumps(body.get("messages", []))) + len(json.dumps(body.get("system", ""))) \ | |
| + len(json.dumps(body.get("tools", []))) | |
| return JSONResponse({"input_tokens": max(1, n // 4)}) | |
| # ------------------------------------------------------------------- console | |
| # Une page servie par le pont lui-meme, donc de meme origine que /v1/messages : | |
| # un fichier ouvert depuis le disque, ou une page hebergee ailleurs, se ferait | |
| # refuser la requete. C'est aussi la raison pour laquelle le pod n'a besoin que | |
| # d'un seul port ouvert. | |
| async def console(): | |
| from fastapi.responses import HTMLResponse, PlainTextResponse | |
| path = os.path.join(os.path.dirname(os.path.abspath(__file__)), "console.html") | |
| try: | |
| with open(path, encoding="utf-8") as fh: | |
| return HTMLResponse(fh.read()) | |
| except OSError: | |
| return PlainTextResponse( | |
| "console.html absente a cote du pont.\n" | |
| "Le bootstrap la telecharge depuis le Hub au demarrage.\n" | |
| f"Attendue ici : {path}\n", status_code=503) | |
| async def health(): | |
| try: | |
| r = await _client.get("/health", timeout=5) | |
| ok = r.status_code == 200 | |
| except Exception: | |
| ok = False | |
| return JSONResponse({"status": "ok" if ok else "upstream_down", | |
| "upstream": UPSTREAM, "model": MODEL, "ts": time.time()}, | |
| status_code=200 if ok else 503) | |
| async def models(request: Request): | |
| """Deux protocoles sur la meme route. | |
| Un client Anthropic attend `{"data":[{"type":"model","id":...}]}` avec des | |
| identifiants en "claude-*". Un client OpenAI ou un banc attend la reponse | |
| de vLLM, dont il lit `max_model_len`. On distingue sur l'en-tete | |
| `anthropic-version`, et on renvoie a chacun ce qu'il sait lire. | |
| """ | |
| if "anthropic-version" not in request.headers: | |
| try: | |
| r = await _client.get("/v1/models", timeout=2.0) | |
| if r.status_code == 200: | |
| return JSONResponse(r.json()) | |
| except Exception: | |
| pass | |
| ctx = None | |
| try: | |
| r = await _client.get("/v1/models", timeout=2.0) | |
| if r.status_code == 200: | |
| ctx = (r.json().get("data") or [{}])[0].get("max_model_len") | |
| except Exception: | |
| pass | |
| data = [{ | |
| "type": "model", | |
| "id": a, | |
| # Le nom affiche porte le depot exact : c'est la seule facon pour un | |
| # utilisateur de savoir quel modele repond derriere une etiquette. | |
| "display_name": f"{a} -> {_courant()[2]}", | |
| "created_at": "2026-01-01T00:00:00Z", | |
| **({"context_window": ctx} if ctx else {}), | |
| } for a in ALIASES] | |
| return JSONResponse({"data": data, "has_more": False, | |
| "first_id": data[0]["id"], "last_id": data[-1]["id"]}) | |
| async def model_detail(model_id: str): | |
| return JSONResponse({"type": "model", "id": model_id, | |
| "display_name": f"{model_id} ({MODEL})", | |
| "created_at": "2026-01-01T00:00:00Z"}) | |
| # ------------------------------------------------------------------ passe-plat | |
| # Le pod n'expose qu'un port. Le pont le prend (c'est l'API que Claude Code | |
| # consomme) et relaie ces routes vers vLLM, pour que les outils de mesure et les | |
| # clients OpenAI restent joignables de l'exterieur. | |
| def _real_model(body: dict) -> dict: | |
| """Traduit un alias claude-* vers le nom servi par vLLM. | |
| Sans ca, un harnais de mesure qui reprend l'identifiant vu dans /v1/models | |
| (`claude-ornith`) recoit un 404 "model not found" de vLLM, releve comme une | |
| route absente alors que seul le nom etait inconnu. | |
| """ | |
| m = body.get("model") | |
| if isinstance(m, str) and (m in ALIASES or _cle_demandee(m)): | |
| body = {**body, "model": _courant()[1]} | |
| return body | |
| async def oai_chat(request: Request): | |
| body = await request.json() | |
| await _assurer_modele(body.get("model")) | |
| body = _real_model(body) | |
| if body.get("stream"): | |
| async def gen(): | |
| async with _client.stream("POST", "/v1/chat/completions", json=body) as r: | |
| async for chunk in r.aiter_raw(): | |
| yield chunk | |
| return StreamingResponse(_flux_gueri(gen()), | |
| media_type="text/event-stream") | |
| r = await _client.post("/v1/chat/completions", json=body) | |
| return JSONResponse(status_code=r.status_code, content=r.json()) | |
| async def oai_completions(request: Request): | |
| body = await request.json() | |
| await _assurer_modele(body.get("model")) | |
| body = _real_model(body) | |
| r = await _client.post("/v1/completions", json=body) | |
| return JSONResponse(status_code=r.status_code, content=r.json()) | |
| async def metrics(): | |
| from fastapi.responses import PlainTextResponse | |
| r = await _client.get("/metrics") | |
| return PlainTextResponse(r.text, status_code=r.status_code) | |