ai-memory-backend / api /benchmark_handler.py
Baida07's picture
sync: 166 file da Baida98/AI@a6ac2424e11e5c320c5ff688e1ce7addac64cdab (local-fallback deploy-all) (#12)
80065ea
Raw
History Blame
15.8 kB
"""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)