""" backend/api/hf_monitor.py — Monitoraggio centralizzato HF Spaces (ARCH-P5.2-MONITOR) Aggrega salute e stato di tutti gli HF Spaces configurati: - Polling salute via /health endpoint di ogni Space - Stato runtime via HuggingFace API (RUNNING/SLEEPING/BUILDING) - Cache interna con TTL per evitare flooding degli endpoint - Background polling ogni 60s Endpoints: GET /api/hf-monitor/spaces — stato aggregato tutti gli Space (cached) POST /api/hf-monitor/spaces/refresh — forza aggiornamento immediato GET /api/hf-monitor/ping/{space_id} — ping live singolo Space Invarianti: - Zero crash se uno Space non risponde (tutto in try/except) - Token non esposto nei response - Cache aggiornata in background senza bloccare l'API - Idempotente: start_monitor() sicuro da chiamare più volte """ from __future__ import annotations import asyncio import logging import os import time from dataclasses import dataclass, field, asdict from typing import Any import httpx from fastapi import APIRouter, Depends, HTTPException from .auth_guard import AuthRole, require_role _logger = logging.getLogger("api.hf_monitor") router = APIRouter() # ── Configurazione ────────────────────────────────────────────────────────────── _PING_TIMEOUT_S = 8.0 # timeout ping singolo Space _HF_API_TIMEOUT_S = 6.0 # timeout HuggingFace runtime API _CACHE_TTL_S = 60.0 # TTL cache risultati _POLL_INTERVAL_S = 60.0 # intervallo background polling _HF_API_BASE = "https://huggingface.co/api/spaces" # ── Tipi ──────────────────────────────────────────────────────────────────────── @dataclass class SpaceConfig: space_id: str # ID univoco interno (es. "brain", "daemon") name: str # Nome display (es. "Brain / Backend") base_url: str # URL base dello Space hf_repo_id: str = "" # ID repo HF (es. "arjanit98/terminal") — per API runtime role: str = "" # ruolo: brain | daemon | worker | compute @dataclass class SpaceStatus: space_id: str name: str base_url: str role: str # Health health: str = "unknown" # ok | degraded | down | unknown latency_ms: float = 0.0 http_status: int = 0 # HF Runtime stage hf_stage: str = "unknown" # RUNNING | SLEEPING | BUILDING | FAILED | unknown # Meta checked_at: float = field(default_factory=time.time) error: str = "" def _build_spaces_from_env() -> list[SpaceConfig]: """Legge configurazioni degli Space dalle env var. Zero crash su var mancanti.""" spaces: list[SpaceConfig] = [] def _add(space_id: str, name: str, env_url: str, hf_repo_id: str = "", role: str = "") -> None: url = os.getenv(env_url, "").rstrip("/") if url: spaces.append(SpaceConfig( space_id=space_id, name=name, base_url=url, hf_repo_id=hf_repo_id, role=role, )) else: _logger.debug("hf_monitor: %s non configurato (%s vuoto)", space_id, env_url) # Fleet — aggiungere nuove righe al crescere degli Space _add("brain", "Brain / Backend", "HF_SPACE_URL", role="brain") _add("daemon", "Daemon / Telegram", "HF_SPACE_B_URL", role="daemon") _add("worker-a", "Worker A (Collab)", "HF_SPACE_C_URL", role="worker") _add("worker-b", "Worker B", "HF_SPACE_D_URL", role="worker") _add("worker-c", "Worker C", "HF_SPACE_E_URL", role="worker") # ── Ruoli futuri — skippati automaticamente se env var non configurata ── _add("executor", "Brain Executor", "HF_SPACE_EXECUTOR_URL", role="executor") _add("browser-worker", "Browser Worker", "HF_SPACE_BROWSER_URL", role="browser-worker") _add("memory-worker", "Memory Worker", "HF_SPACE_MEMORY_URL", role="memory-worker") _add("staging", "Staging / Collab B", "HF_SPACE_STAGING_URL", role="staging") _add("oracle", "Oracle Cloud VM", "ORACLE_CLOUD_VM_URL", role="compute") return spaces # Singleton — letto all'import del modulo _SPACES: list[SpaceConfig] = _build_spaces_from_env() # ── Cache ──────────────────────────────────────────────────────────────────────── _cache: dict[str, SpaceStatus] = {} _cache_ts: float = 0.0 _poll_task: asyncio.Task | None = None # type: ignore[type-arg] async def _ping_space(client: httpx.AsyncClient, cfg: SpaceConfig) -> SpaceStatus: """Ping /health di uno Space. Non solleva mai eccezioni.""" t0 = time.time() status = SpaceStatus( space_id=cfg.space_id, name=cfg.name, base_url=cfg.base_url, role=cfg.role, ) try: url = f"{cfg.base_url}/health" resp = await client.get(url, timeout=_PING_TIMEOUT_S) ms = (time.time() - t0) * 1000 status.http_status = resp.status_code status.latency_ms = round(ms, 1) if resp.status_code < 400: status.health = "ok" elif resp.status_code < 500: status.health = "degraded" else: status.health = "down" except httpx.TimeoutException: status.health = "down" status.error = "timeout" status.latency_ms = _PING_TIMEOUT_S * 1000 except Exception as exc: status.health = "down" status.error = str(exc)[:120] status.checked_at = time.time() return status async def _get_hf_stage(client: httpx.AsyncClient, repo_id: str, hf_token: str) -> str: """Recupera lo stage runtime da HF API. Ritorna 'unknown' su qualsiasi errore.""" if not repo_id or not hf_token: return "unknown" try: url = f"{_HF_API_BASE}/{repo_id}/runtime" hdrs = {"Authorization": f"Bearer {hf_token}"} resp = await client.get(url, headers=hdrs, timeout=_HF_API_TIMEOUT_S) if resp.status_code == 200: return str(resp.json().get("stage", "unknown")) except Exception as exc: _logger.debug("hf_monitor: HF runtime API %s: %s", repo_id, exc) return "unknown" async def _refresh_all() -> dict[str, SpaceStatus]: """Aggiorna lo stato di tutti gli Space configurati in parallelo.""" global _cache, _cache_ts if not _SPACES: return {} hf_token = os.getenv("HF_TOKEN", "") async with httpx.AsyncClient() as client: ping_results: list[SpaceStatus] = list(await asyncio.gather( *[_ping_space(client, cfg) for cfg in _SPACES], return_exceptions=False, )) stage_results: list[str] = list(await asyncio.gather( *[_get_hf_stage(client, cfg.hf_repo_id, hf_token) for cfg in _SPACES], return_exceptions=False, )) for status, stage in zip(ping_results, stage_results): status.hf_stage = stage new_cache = {s.space_id: s for s in ping_results} _cache = new_cache _cache_ts = time.time() _logger.info( "hf_monitor: refresh OK — %d space: %s", len(ping_results), ", ".join(f"{s.space_id}={s.health}" for s in ping_results), ) return new_cache async def _background_poll() -> None: """Loop di polling in background — mai si ferma, mai solleva.""" while True: try: await _refresh_all() except Exception as exc: _logger.warning("hf_monitor: errore polling: %s", exc) await asyncio.sleep(_POLL_INTERVAL_S) def start_monitor() -> None: """Avvia il background polling. Idempotente: sicuro da chiamare più volte.""" global _poll_task if _poll_task and not _poll_task.done(): return _poll_task = asyncio.ensure_future(_background_poll()) _logger.info("hf_monitor: polling avviato — %d space, ogni %ds", len(_SPACES), int(_POLL_INTERVAL_S)) def _status_to_dict(s: SpaceStatus) -> dict[str, Any]: return asdict(s) # ── Endpoints ─────────────────────────────────────────────────────────────────── @router.get( "/api/hf-monitor/spaces", summary="Stato aggregato di tutti gli HF Spaces monitorati (ARCH-P5.2)", ) async def get_spaces_status( _auth: None = Depends(require_role(AuthRole.MACHINE)), ) -> dict[str, Any]: """ Restituisce lo stato cached di tutti gli HF Spaces. La cache si aggiorna ogni 60s in background; la prima call forza un refresh. """ global _cache, _cache_ts if not _cache: await _refresh_all() return { "spaces": [_status_to_dict(s) for s in _cache.values()], "total": len(_cache), "cache_age_s": round(time.time() - _cache_ts, 1) if _cache_ts else None, "spaces_configured": len(_SPACES), } @router.post( "/api/hf-monitor/spaces/refresh", summary="Forza aggiornamento immediato stato HF Spaces", ) async def force_refresh( _auth: None = Depends(require_role(AuthRole.MACHINE)), ) -> dict[str, Any]: """Forza refresh sincrono di tutti gli Space. Può richiedere fino a 8s.""" await _refresh_all() return { "spaces": [_status_to_dict(s) for s in _cache.values()], "total": len(_cache), "refreshed_at": time.time(), } @router.get( "/api/hf-monitor/ping/{space_id}", summary="Ping live di un singolo HF Space", ) async def ping_single_space( space_id: str, _auth: None = Depends(require_role(AuthRole.MACHINE)), ) -> dict[str, Any]: """Ping diretto e sincrono di un singolo Space (bypassa la cache).""" cfg = next((c for c in _SPACES if c.space_id == space_id), None) if not cfg: raise HTTPException( status_code=404, detail=f"Space '{space_id}' non configurato. " f"Space disponibili: {[c.space_id for c in _SPACES]}", ) hf_token = os.getenv("HF_TOKEN", "") async with httpx.AsyncClient() as client: status = await _ping_space(client, cfg) stage = await _get_hf_stage(client, cfg.hf_repo_id, hf_token) status.hf_stage = stage return _status_to_dict(status)