Terminal / api /agent_telemetry.py
Baida07's picture
sync: 191 file da Baida98/AI@5943984b (2026-08-29 20:27 UTC) (#146)
72f7ead
Raw
History Blame
9.28 kB
"""
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)