"""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, "❌ Runner benchmark esteso non disponibile.\n"
"Il deployment non ha incluso benchmark-extended.mjs.")
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,
"🎯 Benchmark Extended v5 mirato avviato\n"
"10 categorie più deboli della baseline 39,1 · seed 1337 · task API moderna.",
)
else:
report_path = _REPORT_V7
flags = ["--full", "--json", f"--output={report_path}", "--gap-analysis"]
await send_reply_fn(
chat_id,
"🚀 Benchmark Extended v5 avviato\n"
"20/20 categorie · seed 1337 · task API moderna · durata variabile fino a ~60 min.",
)
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"❌ Errore benchmark Extended:\n{err}")
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"⏱ Timeout benchmark Extended (>{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"💥 Errore critico benchmark: {exc}")
return
report_exists = await asyncio.to_thread(os.path.exists, report_path)
if not report_exists:
await send_reply_fn(chat_id, "⚠️ Benchmark Extended terminato ma report non trovato.")
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"⚠️ Report Extended non leggibile: {exc}")
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"⚠️ Run incompleta: {len(categories)}/{expected_categories} 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"🏆 Benchmark {ver} completato!\n\n"
f"📊 Score agente: {avg}/100\n"
f"📅 Run: {ts}\n"
+ (f"🎯 Modalità: {run_label}\n" if run_label else "") + "\n"
"📈 Confronto vs riferimenti:\n"
f" • Replit: {repl}/100\n"
f" • Cursor: {curs}/100\n"
f" • Devin: {devi}/100\n"
f" • Manus: {manu}/100\n"
]
if verd:
lines.append(f"\n📝 Verdetto: {verd}\n")
if canary:
lines.append(f"⚠️ Canary leak: {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"🧪 Copertura: {len(attempted_categories)}/{expected_categories} categorie tentate · "
f"{scored} valutabili · {skipped} non valutabili\n"
)
if by_cat:
lines.append("\n📂 Per categoria:\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} {avg_cat:5.1f} {cat}\n")
failures = report.get("taskFailures", [])
if failures:
lines.append("\n⚠️ Categorie non valutabili:\n")
for failure in failures[:3]:
cat = failure.get("cat", "?")
reason = str(failure.get("reason", "errore non specificato"))[:100]
lines.append(f" • {cat} — {reason}\n")
if len(failures) > 3:
lines.append(f" ...e altre {len(failures) - 3}.\n")
# Gap cards (prime 3)
gap_cards = report.get("gapCards", [])
if gap_cards:
lines.append(f"\n💡 Gap ({gaps} totali):\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" • {gid} {gtit}" + (f" [{gsev}]" if gsev else "") + "\n")
if gaps > 3:
lines.append(f" ...e altri {gaps - 3} gap.\n")
lines.append("\n🔍 Usa /riepilogo per analisi approfondita.")
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] = ["📋 Briefing Assistente Proattivo\n\n"]
lines.append("🟢 Stato Sistema: 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"📊 Ultimo Score ({ver}): {avg}/100\n")
lines.append(f"📅 Run: {ts}\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("💡 Aree prioritarie:\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" • {gap['name']} → {slug}\n")
else:
lines.append("✨ Nessun gap critico nell'ultimo run.\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"📊 Ultimo Score (v6.2): {score}/100 [{judge}]\n")
lines.append(f"📅 Run: {ts}\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("💡 Aree prioritarie:\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" • {gap['name']} → {slug}\n")
else:
lines.append("✨ Nessun gap critico nell'ultimo run.\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("📊 Nessun report — esegui /bench per generare dati.\n")
lines.append("\n🚀 Task in corso: controlla /status per stato live.")
return "".join(lines)