From 803b28a39acb12b587d058a1602a84c25803d4d6 Mon Sep 17 00:00:00 2001 From: R0m1k3 Date: Tue, 7 Jul 2026 15:56:04 +0200 Subject: [PATCH] Stabilise Ollama et les benchmarks longs --- backend/app/bench.py | 13 +++++++++- backend/app/ollama_client.py | 23 ++++++++++++----- frontend/src/api/client.ts | 48 +++++++++++++++++++++++++----------- 3 files changed, 62 insertions(+), 22 deletions(-) diff --git a/backend/app/bench.py b/backend/app/bench.py index 3200341..c0f658d 100644 --- a/backend/app/bench.py +++ b/backend/app/bench.py @@ -6,6 +6,7 @@ respect d'un format. Score /100, stocké en base et affiché dans l'UI. """ from __future__ import annotations +import asyncio import json import re import subprocess @@ -189,14 +190,24 @@ async def run_bench(model: str) -> AsyncIterator[dict]: details = [] for name, fn in TASKS: yield {"type": "task_start", "task": name} + task = asyncio.create_task(fn(model)) try: - score, detail = await fn(model) + while not task.done(): + done, _ = await asyncio.wait({task}, timeout=10) + if not done: + # Empêche OpenResty/Nginx de fermer le SSE pendant une longue + # génération d'un gros modèle. + yield {"type": "heartbeat", "task": name} + score, detail = await task except (OllamaError, httpx.HTTPError, OSError) as exc: score, detail = 0, f"erreur : {str(exc)[:80]}" except Exception as exc: # Une épreuve défaillante ne doit pas couper silencieusement le SSE : # elle vaut zéro et les autres épreuves continuent. score, detail = 0, f"épreuve interrompue : {str(exc)[:80]}" + finally: + if not task.done(): + task.cancel() total += score details.append({"task": name, "score": score, "detail": detail}) yield {"type": "task_done", "task": name, "score": score, "detail": detail} diff --git a/backend/app/ollama_client.py b/backend/app/ollama_client.py index f15cef6..b91a2d7 100644 --- a/backend/app/ollama_client.py +++ b/backend/app/ollama_client.py @@ -5,6 +5,7 @@ total sur le streaming et n'embarquer aucune dépendance superflue. """ from __future__ import annotations +import asyncio import json from typing import AsyncIterator @@ -126,12 +127,22 @@ class OllamaClient: # laisse jusqu'à dix minutes pour les gros modèles ou un stockage lent. timeout = httpx.Timeout(connect=10.0, read=600.0, write=30.0, pool=10.0) async with httpx.AsyncClient(timeout=timeout, follow_redirects=True) as client: - resp = await client.post( - f"{self.host}/api/generate", - json={"model": model, "keep_alive": keep_alive, "stream": False}, - ) - resp.raise_for_status() - return resp.json() + for attempt in range(2): + resp = await client.post( + f"{self.host}/api/generate", + json={"model": model, "keep_alive": keep_alive, "stream": False}, + ) + try: + # Inclut le corps JSON d'Ollama dans l'erreur (OOM, runner…), + # contrairement à raise_for_status qui ne montrait que « 500 ». + await _raise_for_stream_status(resp) + except OllamaError: + if resp.status_code >= 500 and attempt == 0: + await asyncio.sleep(2) + continue + raise + return resp.json() + raise OllamaError("préchargement interrompu sans réponse") async def chat( self, diff --git a/frontend/src/api/client.ts b/frontend/src/api/client.ts index 540710d..cca75e1 100644 --- a/frontend/src/api/client.ts +++ b/frontend/src/api/client.ts @@ -205,11 +205,21 @@ export async function runBench( model: string, onProgress: (task: string, score: number | null, detail?: string) => void ): Promise { - const res = await fetch("/api/bench", { - method: "POST", - headers: { "Content-Type": "application/json" }, - body: JSON.stringify({ model }), - }); + let res: Response; + try { + res = await fetch("/api/bench", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ model }), + }); + } catch (err) { + const raw = err instanceof Error ? err.message : "connexion impossible"; + throw new Error( + /network error|failed to fetch|load failed/i.test(raw) + ? "connexion au benchmark interrompue par le réseau ou le reverse proxy" + : raw + ); + } if (!res.ok) { throw await apiError(res, `benchmark refusé (${res.status})`); } @@ -238,18 +248,26 @@ export async function runBench( final = { score: payload.score, details: payload.details, at: Date.now() / 1000 }; }; - while (true) { - const { done, value } = await reader.read(); - if (done) break; - buffer += decoder.decode(value, { stream: true }).replace(/\r\n/g, "\n"); - const events = buffer.split("\n\n"); - buffer = events.pop() ?? ""; - for (const block of events) { - if (block.trim()) dispatch(block); + try { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + buffer += decoder.decode(value, { stream: true }).replace(/\r\n/g, "\n"); + const events = buffer.split("\n\n"); + buffer = events.pop() ?? ""; + for (const block of events) { + if (block.trim()) dispatch(block); + } } + buffer += decoder.decode(); + if (buffer.trim()) dispatch(buffer); + } catch (err) { + const raw = err instanceof Error ? err.message : "connexion interrompue"; + if (/network error|failed to fetch|load failed/i.test(raw)) { + throw new Error("flux du benchmark interrompu par le reverse proxy"); + } + throw err; } - buffer += decoder.decode(); - if (buffer.trim()) dispatch(buffer); if (!final) throw new Error("le benchmark s'est interrompu avant le résultat"); return final; }