Spaces:
Running
Running
File size: 22,987 Bytes
5ac39cd d66556b 5ac39cd d66556b 5ac39cd d66556b 5ac39cd d66556b 5ac39cd 03e5649 5ac39cd 03e5649 5ac39cd 03e5649 5ac39cd d66556b 5ac39cd 03e5649 5ac39cd 03e5649 5ac39cd d66556b 5ac39cd d66556b 5ac39cd d66556b 5ac39cd 03e5649 5ac39cd 03e5649 5ac39cd 03e5649 5ac39cd 03e5649 d66556b 03e5649 d66556b 5ac39cd d66556b 5ac39cd d66556b 5ac39cd d66556b 03e5649 5ac39cd d66556b | 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 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 | """
backend/main.py β FastAPI entry point (S354: split in APIRouter modules).
Mappa dei router:
api.conversations β /api/conversations/**
api.agent_memory β /api/memory/agent/**
api.files β /api/files/**
api.agent β /api/agent/**, /api/reason/loop, /api/unified/loop, /run_loop
api.exec β /api/exec, /api/execute-shell, /api/pip-install, /api/agent/fix
api.search β /api/search, /api/fetch-page, /api/analyze-image
api.providers β /health, /api/tools, /api/status, /api/ai/health, /api/providers/heartbeat
api.vault β /api/vault/**
api.webhook β /api/webhook/**, /api/public/chat
api.terminal β /ws/terminal, GET /api/terminal/packages
api.browser β /api/browser/** (pre-esistente)
api.coding β /api/coding/** (pre-esistente)
api.web β /web/** (pre-esistente)
api.mcp β /api/mcp (P19-B3: MCP JSON-RPC 2.0 server)
"""
import os, time, logging, secrets as _secrets_mod
from fastapi import FastAPI, Request
from fastapi.staticfiles import StaticFiles
from fastapi.responses import JSONResponse
# Gap 2.2: Structured JSON logging β setup PRIMA di qualsiasi altro import
from api.structured_log import setup_structured_logging as _setup_structured_log, router as _logs_router
from api.integrity_manager import router as _integrity_router # P41
_setup_structured_log()
import logging as _boot_logger; _boot_logger.getLogger('agente_ai').info('BOOT: importing FastAPI...')
app = FastAPI(title='Agente AI', version='3.4.2')
# P19-SEC2 (fix concorrente rimosso): `backend/auth/auth_managed.py` non definiva
# alcun `router` β questo import causava un ImportError al boot (app non avviabile).
# Il modulo introduceva anche un secondo sistema di crittografia/token duplicato
# e insicuro (token di sessione hardcoded "secure-session-token"). Rimosso: l'unico
# modulo auth_managed valido resta `backend/api/auth_managed.py`.
_logger = logging.getLogger('agente_ai')
# S274-SEC3: INTERNAL_TOKEN β genera casuale al boot se non configurato.
_GENERATED_TOKEN = _secrets_mod.token_hex(32)
if "INTERNAL_TOKEN" not in os.environ or os.environ["INTERNAL_TOKEN"] == _GENERATED_TOKEN:
os.environ['INTERNAL_TOKEN'] = _GENERATED_TOKEN
_logger.critical('BOOT: INTERNAL_TOKEN not set β ephemeral token generato per questa sessione.')
_logger.critical('BOOT: Ogni restart cambia il token β CF Worker riceve 401 finchΓ© il secret non Γ¨ aggiornato!')
_logger.critical('BOOT: session token (copia in CF Workers secret INTERNAL_TOKEN): %r', _GENERATED_TOKEN)
_logger.critical('BOOT: Fix β imposta INTERNAL_TOKEN uguale su HF Spaces e CF Workers secrets.')
else:
_logger.info('BOOT: INTERNAL_TOKEN configurato OK')
# ββ CORS β env-driven, Safari-safe βββββββββββββββββββββββββββββββββββββββββββ
# ALLOWED_ORIGINS: comma-separated list, e.g. "https://agente-ai.vercel.app,http://localhost:5173"
# Supports wildcard suffix match (*.vercel.app, *.hf.space) for preview URLs.
_ALLOWED_ORIGINS_ENV = os.getenv('ALLOWED_ORIGINS', '')
_ALWAYS_ALLOWED = [
'http://localhost:5173',
'http://localhost:4173',
'http://localhost:3000',
]
_VERCEL_PATTERNS = ['.vercel.app', '.vercel.sh', '.pages.dev']
_HF_PATTERNS = ['.hf.space', '.huggingface.co']
def _build_origin_list() -> list[str]:
origins = list(_ALWAYS_ALLOWED)
if _ALLOWED_ORIGINS_ENV:
for o in _ALLOWED_ORIGINS_ENV.split(','):
o = o.strip()
if o:
origins.append(o)
return origins
_STATIC_ORIGINS = _build_origin_list()
def _is_allowed_origin(origin: str | None) -> bool:
if not origin:
return False
if origin in _STATIC_ORIGINS:
_logger.info('CORS ALLOW (static): %s', origin)
return True
for pat in _VERCEL_PATTERNS + _HF_PATTERNS:
if origin.endswith(pat):
_logger.info('CORS ALLOW (pattern %s): %s', pat, origin)
return True
_logger.warning('CORS BLOCK: %s', origin)
return False
# CORS NOTE (HF Spaces 2026-05-27): HF proxy auto-injects ACAO on all responses.
# This middleware handles preflight OPTIONS for Vercel frontend + documents allowed patterns.
_CORS_HEADERS = {
'Access-Control-Allow-Methods': 'GET,POST,PUT,DELETE,OPTIONS,PATCH',
'Access-Control-Allow-Headers': '*',
'Access-Control-Allow-Credentials': 'true',
'Access-Control-Max-Age': '3600',
'Vary': 'Origin',
}
@app.middleware('http')
async def _cors_middleware(request: Request, call_next):
origin = request.headers.get('origin')
if request.method == 'OPTIONS':
if _is_allowed_origin(origin):
return JSONResponse(content={}, status_code=204, headers={
'Access-Control-Allow-Origin': origin,
**_CORS_HEADERS,
})
return JSONResponse(content={}, status_code=204)
response = await call_next(request)
if _is_allowed_origin(origin):
response.headers['Access-Control-Allow-Origin'] = origin
for k, v in _CORS_HEADERS.items():
response.headers[k] = v
return response
@app.options('/{path:path}')
async def _preflight_fallback(path: str, request: Request):
origin = request.headers.get('origin', '')
if not _is_allowed_origin(origin):
return JSONResponse(content={}, status_code=204)
return JSONResponse(content={}, status_code=204, headers={
'Access-Control-Allow-Origin': origin,
'Access-Control-Allow-Methods': 'GET,POST,PUT,DELETE,OPTIONS,PATCH',
'Access-Control-Allow-Headers': '*',
'Access-Control-Allow-Credentials': 'true',
'Access-Control-Max-Age': '3600',
'Vary': 'Origin',
})
# ββ WARN-1 fix: body size hard limit β S292 βββββββββββββββββββββββββββββββββ
_MAX_BODY_BYTES = 512_000 # 512 KB β increased for PDF/vision content (V006)
@app.middleware('http')
async def _body_size_middleware(request: Request, call_next):
cl = request.headers.get('content-length')
if cl:
try:
if int(cl) > _MAX_BODY_BYTES:
return JSONResponse(
{'detail': f'Payload troppo grande: max {_MAX_BODY_BYTES // 1024}KB. (R-S292)'},
status_code=413,
)
except ValueError:
pass
return await call_next(request)
# ββ S477-SEC4: Rate limiting in-memory per IP ββββββββββββββββββββββββββββββββ
# VITE_INTERNAL_TOKEN Γ¨ visibile nel bundle JS β rate limiting come mitigazione
# pratica all'abuso (CF Worker proxy Γ¨ il fix definitivo, non ancora implementato).
# 120 req/min globale per IP; OPTIONS exempt (preflight CORS non contano).
# In-memory: si resetta a ogni restart HF Space β accettabile su free tier.
_rl_store: dict[str, list[float]] = {}
_RL_WINDOW = 60.0
_RL_LIMIT = 120 # req/minuto per IP
@app.middleware('http')
async def _rate_limit_middleware(request: Request, call_next):
if request.method == "OPTIONS":
return await call_next(request)
ip = (request.client.host if request.client else None) or "unknown"
now = time.monotonic()
hits = [t for t in _rl_store.get(ip, []) if now - t < _RL_WINDOW]
if len(hits) >= _RL_LIMIT:
return JSONResponse(
{"detail": "Too many requests"},
status_code=429,
headers={
"X-RateLimit-Limit": str(_RL_LIMIT),
"X-RateLimit-Remaining": "0",
"X-RateLimit-Reset": str(int(now + _RL_WINDOW)),
"Retry-After": str(int(_RL_WINDOW)),
},
)
hits.append(now)
_rl_store[ip] = hits
# S572: prune _rl_store ogni ~500 req β evita leak memoria con molti IP unici.
# Rimuove IP con zero hit nella finestra (inattivi da > _RL_WINDOW secondi).
if len(_rl_store) > 500:
_cutoff = now - _RL_WINDOW
_stale = [_k for _k, _v in list(_rl_store.items()) if not _v or _v[-1] < _cutoff]
for _k in _stale:
_rl_store.pop(_k, None)
return await call_next(request)
# ββ Include routers (S354 split) βββββββββββββββββββββββββββββββββββββββββββββ
from api.conversations import router as _conv_router
from api.agent_memory import router as _mem_router
from api.files import router as _files_router
from api.agent import router as _agent_router
from api.exec import router as _exec_router
from api.search import router as _search_router
from api.providers import router as _providers_router
from api.vault import router as _vault_router
from api.webhook import router as _webhook_router
from api.terminal import router as _terminal_router
from api.browser import router as _browser_router
from api.coding import router as _coding_router
from api.web import router as _web_router
from api.vision import router as _vision_router # V001: generate_image / analyze_image / search_images
from api.gemini_vision import router as _gemini_vision_router # P48: Gemini 1.5 Flash Vision direct endpoint
from api.email import router as _email_router # V002: send_email via Resend API
from api.database import router as _db_router # V003: database_query PostgreSQL/SQLite
from api.research import router as _research_router # V004: web_research multi-URL + Groq synthesis
from api.deploy import router as _deploy_router # S750: CI status + deploy trigger
from api.scheduler import router as _scheduler_router, start_scheduler as _start_scheduler # GAP-2.1: server-side persistent scheduler
from api.benchmark import router as _benchmark_router # S-BENCH: self-test endpoint /api/debug/benchmark
from api.telemetry import router as _telemetry_router # BG-3: timing metrics /api/telemetry
from api.agent_telemetry import router as _agent_telemetry_router # Gap N4: verdetti cross-session
from api.telegram_webhook import router as _tg_webhook_router # TG-BOT: riceve comandi bot + setup webhook
# notify_bot rimosso β Telegram gestito dal daemon Node.js
from api.incident_registry import router as _incident_router, start_incident_registry as _start_incident_reg # GAP-A1
from api.decision_memory import router as _decision_router, start_decision_memory as _start_decision_mem # GAP-A2
from api.llm_cache import router as _cache_router # Gap 2.3: /api/cache/stats
from api.daemon_status import router as _daemon_status_router # DAEMON-STATUS: /api/daemon/status
from api.auth_guard import AuthRole, require_role # GAP-A6: importa per uso nei router
from api.blackboard import router as _blackboard_router # S-BB: shared blackboard cross-agent via Upstash
from api.job_queue import router as _jq_router # S-DUAL-2: /api/jq/** Redis coordination
from agents.skill_tracker import skill_router as _skill_tracker_router # P17-B2: POST /skill-record + DELETE /skill-stats
from api.mcp import router as _mcp_router # P19-B3: MCP JSON-RPC 2.0 server
from api.auth_managed import router as _auth_managed_router # P38: OAuth one-click connectors
from api.skills import router as _skills_router # P17-B2
# Doc2-1b-FIX: memory/sync router non era montato β endpoint /api/memory/sync/* non raggiungibili
# NOTA: create_memory_sync_router(memory) Γ¨ una factory β richiede l'istanza MemoryManager.
# GAP-5-FIX: memory/sync router montato in _on_startup() (vedi sotto)
app.include_router(_auth_managed_router) # P38: OAuth one-click
app.include_router(_conv_router)
app.include_router(_mem_router)
app.include_router(_files_router)
app.include_router(_agent_router)
app.include_router(_exec_router)
app.include_router(_search_router)
app.include_router(_providers_router)
app.include_router(_vault_router)
app.include_router(_webhook_router)
app.include_router(_terminal_router)
app.include_router(_browser_router)
app.include_router(_coding_router)
app.include_router(_web_router)
app.include_router(_vision_router)
app.include_router(_gemini_vision_router) # P48: /api/vision/gemini + /api/vision/screenshot_analyze
app.include_router(_email_router)
app.include_router(_db_router)
app.include_router(_research_router)
app.include_router(_deploy_router)
app.include_router(_scheduler_router) # GAP-2.1
app.include_router(_benchmark_router) # S-BENCH: /api/debug/benchmark
app.include_router(_tg_webhook_router) # TG-BOT: /api/telegram/webhook + /api/telegram/config/invalidate
app.include_router(_telemetry_router) # BG-3: /api/telemetry
app.include_router(_agent_telemetry_router) # Gap N4: /api/agent-telemetry/sync
app.include_router(_logs_router) # Gap 2.2: /api/logs + /api/logs/frontend
app.include_router(_incident_router) # GAP-A1: Incident Registry
app.include_router(_decision_router) # GAP-A2: Decision Memory
app.include_router(_cache_router) # Gap 2.3: /api/cache/stats
app.include_router(_daemon_status_router) # DAEMON-STATUS: /api/daemon/status
app.include_router(_blackboard_router) # S-BB: /api/blackboard/**
app.include_router(_jq_router) # S-DUAL-2: /api/jq/**
app.include_router(_integrity_router) # P41: /api/integrity/**
if _skill_tracker_router is not None:
app.include_router(_skill_tracker_router) # P17-B2: /api/agent/skill-record + /api/agent/skill-stats (DELETE)
app.include_router(_mcp_router) # P19-B3: /api/mcp β MCP JSON-RPC 2.0
# (memory/sync router montato in _on_startup)
# ββ Startup: heartbeat + warmup ββββββββββββββββββββββββββββββββββββββββββββββββ
import asyncio as _asyncio_main
def _log_task_exc(task: '_asyncio_main.Task[object]', name: str = '') -> None:
"""Done callback β loga eccezioni non gestite nei background task (P2)."""
try:
if not task.cancelled():
exc = task.exception()
if exc:
_logger.error('BG task %r crashed: %s', name or task.get_name(), exc, exc_info=exc)
except Exception:
pass
@app.on_event('startup')
async def _on_startup():
from api.providers import start_heartbeat
start_heartbeat()
_logger.info('BOOT: heartbeat started')
_start_scheduler()
_start_incident_reg() # GAP-A1: Incident Registry
_start_decision_mem() # GAP-A2: Decision Memory
_logger.info('BOOT: scheduler server-side avviato')
# S388: warmup TCP connection pools β inizializza i client Groq con 1 token
# così la prima vera richiesta utente non paga il costo di handshake HTTP/TLS (~80ms per provider).
# GAP-5-FIX: factory montata qui, DOPO _get_mem_manager_async() che garantisce
# MemoryManager.init() completato prima che le richieste arrivino.
try:
from memory.sync import create_memory_sync_router as _create_sync_router
from api.state import _get_mem_manager_async as _gmm_async
_mem = await _gmm_async()
if _mem is not None:
_sync_router = _create_sync_router(_mem)
app.include_router(_sync_router)
_logger.info('BOOT: memory/sync router OK')
else:
_logger.info('BOOT: memory/sync router skip (manager None)')
except Exception as _sync_err:
_logger.warning('BOOT: memory/sync err β %s', _sync_err)
# NOTA: _skills_router non deve dipendere da _mem nΓ© dall'init del MemoryManager β
# regressione introdotta da un commit concorrente che lo aveva spostato dentro l'if
# sopra, disabilitandolo quando il MemoryManager fallisce l'init. Registrato in un
# try/except indipendente e dedicato cosi un fallimento dell'uno non silenzia l'altro
# (audit GAP #2, 2026-07-08).
try:
app.include_router(_skills_router) # P17-B2
except Exception as _skills_err:
_logger.warning('BOOT: skills_router registration err β %s', _skills_err)
import asyncio as _aio
_t_wm = _aio.create_task(_startup_warmup())
_t_wm.add_done_callback(lambda t: _log_task_exc(t, 'startup_warmup'))
# GAP-STATE: ripristina _agent_tasks da Supabase snapshot + avvia bg persist
try:
from api.state import restore_agent_tasks_from_snap as _restore_snap, persist_state_snapshot as _snap_bg
await _restore_snap()
# GAP-2: crash-recovery β task zombi RUNNING vengono resettati a pending
try:
from api.state import _agent_tasks as _agt_cr
_requeued = 0
for _tid_cr, _td_cr in list(_agt_cr.items()):
if _td_cr.get('_snap_restored') and _td_cr.get('status') in ('running', 'RUNNING'):
_td_cr['status'] = 'pending'
_td_cr['_crash_recovered'] = True
_requeued += 1
if _requeued:
_logger.info('BOOT GAP-2: %d task crash-recovered β status reset a pending', _requeued)
except Exception as _cr_err:
_logger.warning('BOOT GAP-2: crash-recover skip β %s', _cr_err)
_t_snap = _aio.create_task(_snap_bg())
_t_snap.add_done_callback(lambda t: _log_task_exc(t, 'snap_bg'))
_logger.info('BOOT: GAP-STATE snapshot bg avviato')
except Exception as _gstate_err:
_logger.warning('BOOT: GAP-STATE skip β %s', _gstate_err)
# NOTA: heartbeat Telegram rimosso β notifiche gestite dal daemon Node.js
# GAP-NEW-5: telemetry alert loop β campiona ogni 5min, alert Telegram su soglie
try:
from api.telemetry import telemetry_alert_loop as _tel_alert
_t_tel = _aio.create_task(_tel_alert())
_t_tel.add_done_callback(lambda t: _log_task_exc(t, 'telemetry_alert_loop'))
_logger.info('BOOT: telemetry alert loop avviato')
except Exception as _tel_err:
_logger.warning('BOOT: telemetry alert skip β %s', _tel_err)
# S-DUAL-2: job queue consumer + load publisher
try:
from api.job_queue import start_job_queue_consumer as _start_jq
_t_jq = _aio.create_task(_start_jq())
_t_jq.add_done_callback(lambda t: _log_task_exc(t, 'job_queue_consumer'))
_logger.info('BOOT: job queue consumer/publisher avviato (SPACE_ROLE=%s)', os.getenv('SPACE_ROLE', 'unknown'))
except Exception as _jq_err:
_logger.warning('BOOT: job queue consumer skip β %s', _jq_err)
# TG-WEBHOOK-AUTO: Registra il webhook all'avvio se USE_WEBHOOK=true
if os.getenv('USE_WEBHOOK', '').lower() == 'true':
try:
from api.telegram_webhook import setup_telegram_webhook as _setup_tg_wh
_t_tg_wh = _aio.create_task(_setup_tg_wh())
_t_tg_wh.add_done_callback(lambda t: _log_task_exc(t, 'setup_tg_wh'))
_logger.info('BOOT: Telegram webhook auto-setup task creato')
except Exception as _tg_wh_err:
_logger.warning('BOOT: Telegram webhook auto-setup skip β %s', _tg_wh_err)
async def _startup_warmup() -> None:
"""
S388: Warmup dei provider Groq al boot.
Spara 1 token a ogni slot Groq in parallelo β preinizializza i connection pool HTTP.
Non blocca il boot, fallback silenzioso su qualsiasi errore.
Attende 1s per permettere a FastAPI di completare il setup.
"""
import asyncio as _aio
await _aio.sleep(1)
try:
from api.state import _get_ai_client
client = _get_ai_client()
if not client or not client.providers:
return
# Warma solo gli slot Groq (veloci, <500ms) β non Gemini/OpenRouter
groq_providers = [p for p in client.providers if p.name.startswith("groq")]
if not groq_providers:
return
async def _warm_one(provider) -> None:
try:
c = client._client_for(provider)
await _aio.wait_for(
_aio.to_thread(
c.chat.completions.create,
model=provider.default_model,
messages=[{"role": "user", "content": "hi"}],
max_tokens=1,
stream=False,
),
timeout=5.0,
)
_logger.info('BOOT: warmup OK β %s (%s)', provider.name, provider.default_model.split('/')[-1][:24])
except Exception as exc:
_logger.warning('BOOT: warmup skip β %s: %s', provider.name, str(exc)[:60])
await _aio.gather(*[_warm_one(p) for p in groq_providers])
except Exception as exc:
_logger.warning('BOOT: warmup failed: %s', exc)
# P17-B4: pip pre-warm β importa i 20 moduli piΓΉ usati dagli script sandbox
# così la prima exec utente non paga il costo di import (~30-200ms/modulo).
# Silenzioso: se non installato, skip.
import importlib as _imp
_PIP_PREWARM = [
"numpy", "pandas", "matplotlib", "requests", "httpx",
"json", "re", "os", "sys", "math",
"datetime", "pathlib", "itertools", "functools", "collections",
"typing", "dataclasses", "io", "base64", "hashlib",
]
for _pkg in _PIP_PREWARM:
try:
_imp.import_module(_pkg)
except Exception:
pass
_logger.info("BOOT: pip pre-warm %d modules done", len(_PIP_PREWARM))
# ββ P17-B3: shutdown β chiudi exec_http_client (evita fd leak) βββββββββββββββ
@app.on_event('shutdown')
async def _on_shutdown_exec_client():
"""P17-B3: cleanup del persistent client httpx al termine del processo."""
try:
from tools.registry import _exec_http_client as _ehc
if _ehc is not None and not _ehc.is_closed:
await _ehc.aclose()
_logger.info('SHUTDOWN: exec_http_client closed (P17-B3)')
except Exception as _e:
_logger.debug('SHUTDOWN: exec_http_client close skipped: %s', _e)
# ββ Frontend static (SPA) ββββββββββββββββββββββββββββββββββοΏ½οΏ½βββββββββββββββββββ
_STATIC_DIR = os.getenv('FRONTEND_DIST', '/app/backend/static')
if os.path.isdir(_STATIC_DIR):
app.mount('/', StaticFiles(directory=_STATIC_DIR, html=True), name='spa')
_logger.info('BOOT: serving frontend from %s', _STATIC_DIR)
else:
_logger.warning('BOOT: no frontend at %s', _STATIC_DIR)
_logger.info('BOOT: main.py v%s ready β %s routes registered β', app.version, len(app.routes))
# Test comment for synchronization
|