""" backend/api/agent_telemetry.py — Sync verdetti agentTelemetry.ts cross-session/device. Gap N4: agentTelemetry.ts usava solo localStorage → dati di calibrazione persi su altri device o dopo clear della cache. Questo endpoint li persiste su /data/ del volume HF Space (stesso usato dalla memoria episodica — mai si perde tra restart). Endpoints: POST /api/agent-telemetry/sync — riceve TelemetryStore dal client, merge server, persiste, ritorna il merged aggiornato. GET /api/agent-telemetry/sync — ritorna store server-side (load al boot del client). DELETE /api/agent-telemetry/sync — reset admin/debug. Auth: nessuna (dati aggregati, zero PII — stesso pattern di /api/telemetry esistente). Rate limit: middleware globale 120 req/min/IP già applicato da main.py. Merge: additive — per ogni (system, verdict) prende max(count) e max(lastSeenMs). I contatori non diminuiscono mai (protezione da client con dati parziali). """ import os, json, logging, time, asyncio from pathlib import Path from fastapi import APIRouter, Depends from .auth_guard import require_role, AuthRole from fastapi.responses import JSONResponse from pydantic import BaseModel router = APIRouter( dependencies=[Depends(require_role(AuthRole.MACHINE))]) # GAP-1-fix: router-level auth _logger = logging.getLogger("agente_ai") # ─── Storage ────────────────────────────────────────────────────────────────── _DATA_DIR = Path(os.getenv("DATA_DIR", "/data")) _TEL_FILE = _DATA_DIR / "agent_telemetry.json" _MAX_STORE = 500 # max entry totali prima del pruning # Serializza write concorrenti: due POST simultanei da device diversi leggerebbero # lo stesso store e si sovrascriverebbero. asyncio.Lock() è safe a livello di modulo # in Python 3.10+ (non richiede event loop attivo all'init del modulo). _store_lock = asyncio.Lock() _RUNTIME_PHASES = ('auth', 'queue', 'provider', 'tool', 'persistence') _RUNTIME_MAX_SAMPLES = 200 _runtime_lock = asyncio.Lock() _runtime_store: dict[str, dict] = { phase: {'samples_ms': [], 'ok': 0, 'errors': {}} for phase in _RUNTIME_PHASES } async def record_runtime_phase(phase: str, duration_ms: float = 0.0, outcome: str = 'ok', error_class: str | None = None) -> None: if phase not in _runtime_store: return async with _runtime_lock: bucket = _runtime_store[phase] samples = bucket['samples_ms'] samples.append(max(0.0, round(float(duration_ms), 2))) if len(samples) > _RUNTIME_MAX_SAMPLES: del samples[:-_RUNTIME_MAX_SAMPLES] if outcome == 'ok': bucket['ok'] += 1 else: key = error_class or outcome or 'unknown' bucket['errors'][key] = bucket['errors'].get(key, 0) + 1 def _percentile(samples: list[float], percentile: float) -> float: if not samples: return 0.0 ordered = sorted(samples) index = min(len(ordered) - 1, int(round((percentile / 100) * (len(ordered) - 1)))) return ordered[index] async def runtime_snapshot() -> dict[str, dict]: async with _runtime_lock: return { phase: { 'count': len(bucket['samples_ms']), 'ok': bucket['ok'], 'errors': dict(bucket['errors']), 'p50_ms': _percentile(bucket['samples_ms'], 50), 'p95_ms': _percentile(bucket['samples_ms'], 95), } for phase, bucket in _runtime_store.items() } # ─── Store helpers ──────────────────────────────────────────────────────────── def _load_store() -> dict: """Carica /data/agent_telemetry.json. Ritorna {} in caso di errore.""" try: if _TEL_FILE.exists(): return json.loads(_TEL_FILE.read_text("utf-8")) except Exception as exc: _logger.warning("agent_telemetry: load error — %s", exc) return {} def _save_store(store: dict) -> None: """Persiste store su disco. Fail-open.""" try: _DATA_DIR.mkdir(parents=True, exist_ok=True) _TEL_FILE.write_text(json.dumps(store, separators=(",", ":")), "utf-8") except Exception as exc: _logger.warning("agent_telemetry: save error — %s", exc) def _merge(server: dict, client: dict) -> dict: """ Merge additive cross-device. Regola: per ogni (system, verdict) prende max(count) e max(lastSeenMs). Un client con dati parziali non può mai ridurre i contatori server. """ merged: dict = {k: dict(v) for k, v in server.items()} for sys_name, verdicts in client.items(): if not isinstance(verdicts, dict): continue if sys_name not in merged: merged[sys_name] = {} srv_sys = merged[sys_name] for verdict, stats in verdicts.items(): if not isinstance(stats, dict): continue c_count = int(stats.get("count", 0)) c_ts = int(stats.get("lastSeenMs", 0)) prev = srv_sys.get(verdict, {"count": 0, "lastSeenMs": 0}) srv_sys[verdict] = { "count": max(int(prev.get("count", 0)), c_count), "lastSeenMs": max(int(prev.get("lastSeenMs", 0)), c_ts), } return merged def _prune(store: dict) -> dict: """Se totale entry > _MAX_STORE, sacrifica i sistemi meno usati.""" total = sum(len(v) for v in store.values()) if total <= _MAX_STORE: return store sorted_sys = sorted( store.items(), key=lambda kv: sum(s.get("count", 0) for s in kv[1].values()), reverse=True, ) pruned: dict = {} kept = 0 for sys_name, verdicts in sorted_sys: n = len(verdicts) if kept + n > _MAX_STORE: break pruned[sys_name] = verdicts kept += n return pruned # ─── Pydantic models ────────────────────────────────────────────────────────── class TelemetrySyncBody(BaseModel): """ Payload POST dal client — struttura identica a TelemetryStore di agentTelemetry.ts: { "crossCritic": { "pass": {"count": 5, "lastSeenMs": 1718000000000} }, ... } """ data: dict # ─── Endpoints ──────────────────────────────────────────────────────────────── @router.get("/api/agent-telemetry/runtime") async def get_runtime_telemetry() -> JSONResponse: return JSONResponse({'ok': True, 'phases': await runtime_snapshot(), 'server_ts': int(time.time() * 1000)}) @router.post("/api/agent-telemetry/sync") async def post_agent_telemetry(body: TelemetrySyncBody) -> JSONResponse: """ POST /api/agent-telemetry/sync Riceve il TelemetryStore locale del client (agentTelemetry.ts localStorage). Lo merge con lo store server-side (additive — max count), lo persiste su /data/agent_telemetry.json, e ritorna il merged. Il client sostituisce il suo localStorage con il merged ricevuto: i dati da altri device sono ora disponibili localmente. """ try: async with _store_lock: server = _load_store() merged = _prune(_merge(server, body.data)) _save_store(merged) return JSONResponse({ "ok": True, "merged": merged, "server_ts": int(time.time() * 1000), }) except Exception as exc: _logger.error("agent_telemetry POST error: %s", exc) return JSONResponse({"ok": False, "error": str(exc)[:120]}, status_code=500) @router.get("/api/agent-telemetry/sync") async def get_agent_telemetry() -> JSONResponse: """ GET /api/agent-telemetry/sync Ritorna lo store server-side. Il client lo carica all'avvio e lo mergia con il localStorage locale (additive — max count) prima di iniziare a registrare nuovi eventi. """ try: store = _load_store() return JSONResponse({ "ok": True, "data": store, "server_ts": int(time.time() * 1000), }) except Exception as exc: _logger.error("agent_telemetry GET error: %s", exc) return JSONResponse({"ok": False, "error": str(exc)[:120]}, status_code=500) @router.delete("/api/agent-telemetry/sync") async def delete_agent_telemetry() -> JSONResponse: """ DELETE /api/agent-telemetry/sync — solo per admin/debug. Cancella lo store server-side. Non tocca il localStorage dei client. """ try: async with _store_lock: if _TEL_FILE.exists(): _TEL_FILE.unlink() return JSONResponse({"ok": True, "deleted": True}) except Exception as exc: _logger.error("agent_telemetry DELETE error: %s", exc) return JSONResponse({"ok": False, "error": str(exc)[:120]}, status_code=500)