File size: 7,421 Bytes
28a08e7
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
"""
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()

# ─── 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.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)