Spaces:
Running
Running
| """backend/api/benchmark_handler.py β Gestore benchmark via Telegram. | |
| Versione v7 (GAP-BENCH-2): usa benchmark-extended.mjs (comprehensive, 2319 righe) | |
| con flag --json. Mantiene compatibilitΓ backward con report v6.2-stress per /riepilogo. | |
| Sprint history: | |
| v6.2 β benchmark-ultra-v6.2-stress.mjs (208 righe, stress+Groq judge) | |
| v7 β benchmark-extended.mjs (2319 righe, 10+ categorie, HF datasets, | |
| ref vs Replit/Cursor/Devin/Manus) | |
| Fix invarianti portati da v6.2: | |
| - FIX-1: process.kill() + wait() su TimeoutError β nessun zombie su Railway. | |
| - FIX-2: _load_gap_map() usa importlib isolato β nessuna sys.path pollution. | |
| - FIX-3: I/O su file in asyncio.to_thread β non blocca l'event loop. | |
| """ | |
| from __future__ import annotations | |
| import asyncio, importlib, importlib.util, json, logging, os | |
| from typing import Any | |
| logger = logging.getLogger("agente_ai.benchmark_handler") | |
| # ββ Percorsi server Railway ββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| # Lo Space HF esegue il backend in /app; Railway puΓ² impostare REPO_ROOT. | |
| _REPO_ROOT = os.getenv("REPO_ROOT", "/app") | |
| # Extended v5: 20 categorie. Gli Space possono montare il repository in | |
| # /home/user/app anche quando il Dockerfile dichiara WORKDIR=/app. | |
| _BENCH_SCRIPT_CANDIDATES = ( | |
| os.getenv("BENCHMARK_RUNNER_PATH", "").strip(), | |
| os.path.join(_REPO_ROOT, "benchmark-extended.mjs"), | |
| "/home/user/app/benchmark-extended.mjs", | |
| "/app/benchmark-extended.mjs", | |
| ) | |
| _BENCH_SCRIPT = next( | |
| (candidate for candidate in _BENCH_SCRIPT_CANDIDATES if candidate and os.path.isfile(candidate)), | |
| os.path.join(_REPO_ROOT, "benchmark-extended.mjs"), | |
| ) | |
| _REPORT_V7 = "/tmp/agente-ai/benchmark-v5-latest.json" | |
| _REPORT_V7_WEAK = "/tmp/agente-ai/benchmark-v5-weak-latest.json" | |
| _WEAK_CATEGORIES = ( | |
| "sql", "context_window", "reasoning", "data_analysis", "research_synthesis", | |
| "mmlu", "technical_writing", "code_correct", "feature", "security", | |
| ) | |
| # 20 task seriali possono richiedere piΓΉ di 12 minuti con provider gratuiti. | |
| _BENCH_TIMEOUT = float(os.getenv("BENCH_TIMEOUT_SECS", "3600")) | |
| # v6.2 β usato come fallback in get_smart_summary per compatibilitΓ | |
| _REPORT_V6 = os.path.join(_REPO_ROOT, "benchmark-stress-report.json") | |
| async def run_benchmark_task(chat_id: int, send_reply_fn, mode: str = "full") -> None: | |
| """Esegue il benchmark Extended v5 su tutte le 20 categorie via API task moderna.""" | |
| if not await asyncio.to_thread(os.path.isfile, _BENCH_SCRIPT): | |
| await send_reply_fn(chat_id, "β <b>Runner benchmark esteso non disponibile.</b>\n" | |
| "Il deployment non ha incluso <code>benchmark-extended.mjs</code>.") | |
| return | |
| is_weak_run = mode == "weak" | |
| if is_weak_run: | |
| report_path = _REPORT_V7_WEAK | |
| flags = [ | |
| f"--categories={','.join(_WEAK_CATEGORIES)}", "--json", | |
| f"--output={report_path}", "--gap-analysis", | |
| ] | |
| await send_reply_fn( | |
| chat_id, | |
| "π― <b>Benchmark Extended v5 mirato avviato</b>\n" | |
| "<i>10 categorie piΓΉ deboli della baseline 39,1 Β· seed 1337 Β· task API moderna.</i>", | |
| ) | |
| else: | |
| report_path = _REPORT_V7 | |
| flags = ["--full", "--json", f"--output={report_path}", "--gap-analysis"] | |
| await send_reply_fn( | |
| chat_id, | |
| "π <b>Benchmark Extended v5 avviato</b>\n" | |
| "<i>20/20 categorie Β· seed 1337 Β· task API moderna Β· durata variabile fino a ~60 min.</i>", | |
| ) | |
| env = { | |
| **os.environ, | |
| "INTERNAL_TOKEN": os.getenv("INTERNAL_TOKEN", ""), | |
| "BENCHMARK_BASE_URL": os.getenv("BENCHMARK_BASE_URL", "http://127.0.0.1:7860"), | |
| } | |
| process: asyncio.subprocess.Process | None = None | |
| try: | |
| process = await asyncio.create_subprocess_exec( | |
| "node", _BENCH_SCRIPT, *flags, | |
| stdout=asyncio.subprocess.PIPE, | |
| stderr=asyncio.subprocess.PIPE, | |
| env=env, | |
| ) | |
| _stdout, stderr = await asyncio.wait_for(process.communicate(), timeout=_BENCH_TIMEOUT) | |
| if process.returncode != 0: | |
| err = stderr.decode(errors="replace")[:400] | |
| logger.error("Extended benchmark failed rc=%d: %s", process.returncode, err) | |
| await send_reply_fn(chat_id, f"β <b>Errore benchmark Extended:</b>\n<code>{err}</code>") | |
| return | |
| except asyncio.TimeoutError: | |
| if process is not None: | |
| try: | |
| process.kill() | |
| await process.wait() | |
| except Exception: | |
| pass | |
| logger.warning("Extended benchmark timeout (>%.0fs) β process killed", _BENCH_TIMEOUT) | |
| await send_reply_fn(chat_id, f"β± <b>Timeout benchmark Extended</b> (>{int(_BENCH_TIMEOUT // 60)} min) β processo terminato.") | |
| return | |
| except Exception as exc: | |
| if process is not None: | |
| try: | |
| process.kill() | |
| await process.wait() | |
| except Exception: | |
| pass | |
| logger.exception("run_benchmark_task extended error") | |
| await send_reply_fn(chat_id, f"π₯ <b>Errore critico benchmark:</b> <code>{exc}</code>") | |
| return | |
| report_exists = await asyncio.to_thread(os.path.exists, report_path) | |
| if not report_exists: | |
| await send_reply_fn(chat_id, "β οΈ <b>Benchmark Extended terminato ma report non trovato.</b>") | |
| return | |
| try: | |
| report: dict[str, Any] = await asyncio.to_thread(_read_json, report_path) | |
| except Exception as exc: | |
| await send_reply_fn(chat_id, f"β οΈ <b>Report Extended non leggibile:</b> <code>{exc}</code>") | |
| return | |
| categories = {str(task.get("cat", "")) for task in report.get("tasks", []) if task.get("cat")} | |
| expected_categories = len(_WEAK_CATEGORIES) if is_weak_run else 20 | |
| if len(categories) != expected_categories: | |
| await send_reply_fn(chat_id, f"β οΈ <b>Run incompleta:</b> <code>{len(categories)}/{expected_categories}</code> categorie nel report." | |
| " Nessun risultato incompleto viene presentato come benchmark completo.") | |
| return | |
| await send_reply_fn(chat_id, _format_v7_report(report, expected_categories=expected_categories, run_label="mirato Β· categorie deboli" if is_weak_run else None)) | |
| def _format_v7_report(report: dict[str, Any], *, expected_categories: int = 20, run_label: str | None = None) -> str: | |
| """Formatta il report v7 per Telegram HTML.""" | |
| s = report.get("summary", {}) | |
| ts = (report.get("timestamp") or "")[:16].replace("T", " ") | |
| ver = report.get("version", "extended-v7") | |
| avg = s.get("avgScore", "N/A") | |
| repl = s.get("avgReplit", "N/A") | |
| curs = s.get("avgCursor", "N/A") | |
| devi = s.get("avgDevin", "N/A") | |
| manu = s.get("avgManus", "N/A") | |
| gaps = s.get("gapCount", 0) | |
| verd = s.get("verdict", "") | |
| canary = s.get("canaryLeaks", 0) | |
| lines: list[str] = [ | |
| f"π <b>Benchmark {ver} completato!</b>\n\n" | |
| f"π <b>Score agente:</b> <code>{avg}/100</code>\n" | |
| f"π <b>Run:</b> <code>{ts}</code>\n" | |
| + (f"π― <b>ModalitΓ :</b> <code>{run_label}</code>\n" if run_label else "") + "\n" | |
| "π <b>Confronto vs riferimenti:</b>\n" | |
| f" β’ Replit: <code>{repl}/100</code>\n" | |
| f" β’ Cursor: <code>{curs}/100</code>\n" | |
| f" β’ Devin: <code>{devi}/100</code>\n" | |
| f" β’ Manus: <code>{manu}/100</code>\n" | |
| ] | |
| if verd: | |
| lines.append(f"\nπ <b>Verdetto:</b> {verd}\n") | |
| if canary: | |
| lines.append(f"β οΈ <b>Canary leak:</b> {canary} task\n") | |
| # Score per categoria. Una categoria in timeout resta tentata ma non entra | |
| # nella media: non va trasformata silenziosamente in uno score pari a zero. | |
| tasks = report.get("tasks", []) | |
| if tasks: | |
| attempted_categories = {str(t.get("cat")) for t in tasks if t.get("cat")} | |
| by_cat: dict[str, list[float]] = {} | |
| for t in tasks: | |
| cat = t.get("cat", "?") | |
| sc = t.get("score") | |
| if isinstance(sc, (int, float)): | |
| by_cat.setdefault(cat, []).append(float(sc)) | |
| attempted = s.get("attemptedTaskCount", len(tasks)) | |
| scored = s.get("scoredTaskCount", sum(len(v) for v in by_cat.values())) | |
| skipped = s.get("skippedTaskCount", max(0, attempted - scored)) | |
| lines.append( | |
| f"π§ͺ <b>Copertura:</b> <code>{len(attempted_categories)}/{expected_categories} categorie tentate Β· " | |
| f"{scored} valutabili Β· {skipped} non valutabili</code>\n" | |
| ) | |
| if by_cat: | |
| lines.append("\nπ <b>Per categoria:</b>\n") | |
| for cat, scores in sorted(by_cat.items()): | |
| avg_cat = sum(scores) / len(scores) | |
| icon = "π’" if avg_cat >= 70 else "π‘" if avg_cat >= 50 else "π΄" | |
| lines.append(f" {icon} <code>{avg_cat:5.1f}</code> {cat}\n") | |
| failures = report.get("taskFailures", []) | |
| if failures: | |
| lines.append("\nβ οΈ <b>Categorie non valutabili:</b>\n") | |
| for failure in failures[:3]: | |
| cat = failure.get("cat", "?") | |
| reason = str(failure.get("reason", "errore non specificato"))[:100] | |
| lines.append(f" β’ <code>{cat}</code> β {reason}\n") | |
| if len(failures) > 3: | |
| lines.append(f" <i>...e altre {len(failures) - 3}.</i>\n") | |
| # Gap cards (prime 3) | |
| gap_cards = report.get("gapCards", []) | |
| if gap_cards: | |
| lines.append(f"\nπ‘ <b>Gap ({gaps} totali):</b>\n") | |
| for gc in gap_cards[:3]: | |
| gid = gc.get("id", "?") | |
| gtit = gc.get("title", gc.get("name", "")) | |
| gsev = gc.get("severity", "") | |
| lines.append(f" β’ <code>{gid}</code> {gtit}" + (f" [{gsev}]" if gsev else "") + "\n") | |
| if gaps > 3: | |
| lines.append(f" <i>...e altri {gaps - 3} gap.</i>\n") | |
| lines.append("\nπ <i>Usa /riepilogo per analisi approfondita.</i>") | |
| return "".join(lines) | |
| # ββ Helpers βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def _read_json(path: str) -> dict[str, Any]: | |
| """Lettura JSON sincrona β da eseguire sempre in asyncio.to_thread.""" | |
| with open(path, encoding="utf-8") as f: | |
| return json.load(f) | |
| def _load_gap_map() -> tuple[list, list]: | |
| """Importa GAPS e CATS da scripts/gap_map.py senza inquinare sys.path. | |
| FIX-2: usa importlib.util.spec_from_file_location per un import isolato. | |
| """ | |
| gap_map_path = os.path.join(_REPO_ROOT, "scripts", "gap_map.py") | |
| if not os.path.exists(gap_map_path): | |
| logger.debug("gap_map.py not found at %s", gap_map_path) | |
| return [], [] | |
| try: | |
| spec = importlib.util.spec_from_file_location("_gap_map_isolated", gap_map_path) | |
| if spec is None or spec.loader is None: | |
| return [], [] | |
| mod = importlib.util.module_from_spec(spec) | |
| spec.loader.exec_module(mod) # type: ignore[union-attr] | |
| return getattr(mod, "GAPS", []), getattr(mod, "CATS", []) | |
| except Exception as exc: | |
| logger.warning("_load_gap_map failed: %s", exc) | |
| return [], [] | |
| async def get_smart_summary(chat_id: int) -> str: | |
| """Genera riepilogo strutturato: stato sistema + ultimo score + gap priority. | |
| GAP-BENCH-2: prova prima report v7, fallback a v6.2 per compatibilitΓ . | |
| """ | |
| GAPS, _ = await asyncio.to_thread(_load_gap_map) | |
| lines: list[str] = ["π <b>Briefing Assistente Proattivo</b>\n\n"] | |
| lines.append("π’ <b>Stato Sistema:</b> Railway operativo\n") | |
| # Prova v7 per prima | |
| v7_exists = await asyncio.to_thread(os.path.exists, _REPORT_V7) | |
| v6_exists = await asyncio.to_thread(os.path.exists, _REPORT_V6) | |
| if v7_exists: | |
| try: | |
| report: dict[str, Any] = await asyncio.to_thread(_read_json, _REPORT_V7) | |
| s = report.get("summary", {}) | |
| avg = s.get("avgScore", "N/A") | |
| ts = (report.get("timestamp") or "")[:16].replace("T", " ") | |
| ver = report.get("version", "v7") | |
| lines.append(f"π <b>Ultimo Score ({ver}):</b> <code>{avg}/100</code>\n") | |
| lines.append(f"π <b>Run:</b> <code>{ts}</code>\n\n") | |
| # Categorie con score < 50 β alimenta gap priority | |
| failed_cats: set[str] = set() | |
| for t in report.get("tasks", []): | |
| sc = t.get("score") | |
| if isinstance(sc, (int, float)) and sc < 50: | |
| failed_cats.add(t.get("cat", "").lower()) | |
| if failed_cats and GAPS: | |
| lines.append("π‘ <b>Aree prioritarie:</b>\n") | |
| for gap in GAPS: | |
| if any(c in failed_cats for c in gap.get("categories", [])): | |
| slug = ( | |
| f"feature/{gap['id'].lower()}" | |
| f"-{gap['name'].lower().replace(' ', '-')}" | |
| ) | |
| lines.append(f" β’ <code>{gap['name']}</code> β <code>{slug}</code>\n") | |
| else: | |
| lines.append("β¨ <b>Nessun gap critico nell'ultimo run.</b>\n") | |
| except Exception as exc: | |
| logger.warning("get_smart_summary v7 error: %s", exc) | |
| lines.append("β οΈ Errore lettura report v7 β esegui /bench per aggiornare.\n") | |
| elif v6_exists: | |
| # Fallback v6.2 β compatibilitΓ | |
| try: | |
| report = await asyncio.to_thread(_read_json, _REPORT_V6) | |
| score = report.get("finalScore", "N/A") | |
| ts = (report.get("timestamp") or "")[:16].replace("T", " ") | |
| judge = report.get("judge", "heuristic") | |
| lines.append(f"π <b>Ultimo Score (v6.2):</b> <code>{score}/100</code> [{judge}]\n") | |
| lines.append(f"π <b>Run:</b> <code>{ts}</code>\n\n") | |
| failed_cats: set[str] = set() | |
| for res in report.get("results", []): | |
| if res.get("score", {}).get("total", 100) < 50: | |
| rid = res.get("id", "") | |
| if rid.startswith("STRESS_AMB"): failed_cats.add("recovery") | |
| if rid.startswith("STRESS_REC"): failed_cats.add("recovery") | |
| if rid.startswith("STRESS_MEM"): failed_cats.add("memory_context") | |
| if rid.startswith("WEB"): failed_cats.add("orchestration") | |
| if rid.startswith("CODE"): failed_cats.add("bug_fix") | |
| if rid.startswith("REASON"): failed_cats.add("reasoning") | |
| if failed_cats and GAPS: | |
| lines.append("π‘ <b>Aree prioritarie:</b>\n") | |
| for gap in GAPS: | |
| if any(c in failed_cats for c in gap.get("categories", [])): | |
| slug = ( | |
| f"feature/{gap['id'].lower()}" | |
| f"-{gap['name'].lower().replace(' ', '-')}" | |
| ) | |
| lines.append(f" β’ <code>{gap['name']}</code> β <code>{slug}</code>\n") | |
| else: | |
| lines.append("β¨ <b>Nessun gap critico nell'ultimo run.</b>\n") | |
| except Exception as exc: | |
| logger.warning("get_smart_summary v6 error: %s", exc) | |
| lines.append("β οΈ Errore lettura report v6.2 β esegui /bench per aggiornare.\n") | |
| else: | |
| lines.append("π <b>Nessun report</b> β esegui <code>/bench</code> per generare dati.\n") | |
| lines.append("\nπ <b>Task in corso:</b> controlla /status per stato live.") | |
| return "".join(lines) | |