"""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)