Spaces:
Running
Running
| """ | |
| backend/api/session_manager.py β Session Manager (Fase 1 ADR-S26-S30) | |
| ResponsabilitΓ : gestione del ciclo di vita delle sessioni utente. | |
| Disaccoppia la logica di sessione dal Brain/Executor. | |
| Livelli di storage (tiered): | |
| Tier 1 β In-memory LRU (hot sessions, max 256, sub-ms) | |
| Tier 2 β Redis TTL (sessioni attive cross-restart, TTL 24h) | |
| Tier 3 β Supabase (storico permanente, resume dopo shutdown) | |
| Cycle di vita di una sessione: | |
| CREATED β ACTIVE β [PAUSED] β ENDED | EXPIRED | |
| Invarianti ADR: | |
| S26: ogni sessione Γ¨ persistente (Tier 3) | |
| S27: ogni sessione ha correlation_id per tracciabilitΓ | |
| S29: ogni componente sostituibile via feature flag (SESS_BACKEND=redis|supabase|memory) | |
| Endpoints: | |
| POST /api/sessions β crea sessione (auth: MACHINE) | |
| GET /api/sessions/{session_id} β recupera sessione (auth: MACHINE) | |
| PATCH /api/sessions/{session_id} β aggiorna metadata/status (auth: MACHINE) | |
| DELETE /api/sessions/{session_id}/end β termina sessione (auth: MACHINE) | |
| GET /api/sessions/status β diagnostica (auth: MACHINE) | |
| """ | |
| import asyncio, json, time, uuid, logging, os | |
| from collections import OrderedDict | |
| from typing import Any | |
| from fastapi import APIRouter, Depends, HTTPException | |
| from pydantic import BaseModel, Field | |
| from .auth_guard import require_role, AuthRole | |
| from .state import _sb | |
| _logger = logging.getLogger("api.session_manager") | |
| router = APIRouter( | |
| prefix="/api/sessions", | |
| tags=["session-manager"], | |
| dependencies=[Depends(require_role(AuthRole.MACHINE))], | |
| ) | |
| _SESS_TTL_S = int(os.getenv("SESSION_TTL_SECONDS", str(24 * 3600))) # 24h default | |
| _LRU_MAX = int(os.getenv("SESSION_LRU_MAX", "256")) | |
| _SESS_TABLE = "sessions" | |
| _REDIS_PREFIX = "sess:" | |
| # ββ SessionStatus ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| class SessionStatus: | |
| CREATED = "created" | |
| ACTIVE = "active" | |
| PAUSED = "paused" | |
| ENDED = "ended" | |
| EXPIRED = "expired" | |
| # ββ In-memory LRU (Tier 1) βββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| _lru: "OrderedDict[str, dict]" = OrderedDict() | |
| _lru_lock = asyncio.Lock() | |
| async def _lru_get(session_id: str) -> dict | None: | |
| async with _lru_lock: | |
| if session_id in _lru: | |
| _lru.move_to_end(session_id) | |
| return _lru[session_id] | |
| return None | |
| async def _lru_set(session_id: str, session: dict) -> None: | |
| async with _lru_lock: | |
| _lru[session_id] = session | |
| _lru.move_to_end(session_id) | |
| while len(_lru) > _LRU_MAX: | |
| _lru.popitem(last=False) | |
| # ββ Redis helpers (Tier 2) βββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| async def _redis_get_session(session_id: str) -> dict | None: | |
| try: | |
| import httpx | |
| redis_url = os.getenv("UPSTASH_REDIS_REST_URL", "") | |
| redis_token = os.getenv("UPSTASH_REDIS_REST_TOKEN", "") | |
| if not redis_url: return None | |
| async with httpx.AsyncClient(timeout=2.0) as c: | |
| r = await c.post(redis_url, json=["GET", _REDIS_PREFIX + session_id], | |
| headers={"Authorization": f"Bearer {redis_token}"}) | |
| data = r.json() | |
| if data.get("result"): | |
| return json.loads(data["result"]) | |
| except Exception as exc: | |
| _logger.debug("[session_mgr] redis get skip: %s", exc) | |
| return None | |
| async def _redis_set_session(session_id: str, session: dict) -> None: | |
| try: | |
| import httpx | |
| redis_url = os.getenv("UPSTASH_REDIS_REST_URL", "") | |
| redis_token = os.getenv("UPSTASH_REDIS_REST_TOKEN", "") | |
| if not redis_url: return | |
| async with httpx.AsyncClient(timeout=2.0) as c: | |
| await c.post(redis_url, | |
| json=["SET", _REDIS_PREFIX + session_id, json.dumps(session), "EX", _SESS_TTL_S], | |
| headers={"Authorization": f"Bearer {redis_token}"}) | |
| except Exception as exc: | |
| _logger.debug("[session_mgr] redis set skip: %s", exc) | |
| async def _redis_del_session(session_id: str) -> None: | |
| try: | |
| import httpx | |
| redis_url = os.getenv("UPSTASH_REDIS_REST_URL", "") | |
| redis_token = os.getenv("UPSTASH_REDIS_REST_TOKEN", "") | |
| if not redis_url: return | |
| async with httpx.AsyncClient(timeout=2.0) as c: | |
| await c.post(redis_url, | |
| json=["DEL", _REDIS_PREFIX + session_id], | |
| headers={"Authorization": f"Bearer {redis_token}"}) | |
| except Exception as exc: | |
| _logger.debug("[session_mgr] redis del skip: %s", exc) | |
| # ββ Supabase helpers (Tier 3) ββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| async def _supa_upsert(session: dict) -> None: | |
| if not _sb: return | |
| try: | |
| _sb.table(_SESS_TABLE).upsert({ | |
| "id": session["session_id"], | |
| "status": session["status"], | |
| "user_id": session.get("user_id"), | |
| "correlation_id": session.get("correlation_id"), | |
| "metadata": session.get("metadata", {}), | |
| "created_at": session.get("created_at"), | |
| "last_active_at": session.get("last_active_at"), | |
| "ended_at": session.get("ended_at"), | |
| }).execute() | |
| except Exception as exc: | |
| _logger.debug("[session_mgr] supabase upsert skip: %s", exc) | |
| async def _supa_get(session_id: str) -> dict | None: | |
| if not _sb: return None | |
| try: | |
| res = _sb.table(_SESS_TABLE).select("*").eq("id", session_id).limit(1).execute() | |
| if res.data: | |
| r = res.data[0] | |
| return { | |
| "session_id": r["id"], | |
| "status": r.get("status", SessionStatus.ACTIVE), | |
| "user_id": r.get("user_id"), | |
| "correlation_id": r.get("correlation_id"), | |
| "metadata": r.get("metadata", {}), | |
| "created_at": r.get("created_at"), | |
| "last_active_at": r.get("last_active_at"), | |
| "ended_at": r.get("ended_at"), | |
| } | |
| except Exception as exc: | |
| _logger.debug("[session_mgr] supabase get skip: %s", exc) | |
| return None | |
| # ββ Core session operations βββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| async def create_session(user_id: str | None = None, metadata: dict | None = None) -> dict: | |
| """Crea una nuova sessione e la persiste su tutti i tier. Chiamabile internamente.""" | |
| now = time.time() | |
| session = { | |
| "session_id": str(uuid.uuid4()), | |
| "status": SessionStatus.CREATED, | |
| "user_id": user_id, | |
| "correlation_id": str(uuid.uuid4()), | |
| "metadata": metadata or {}, | |
| "created_at": now, | |
| "last_active_at": now, | |
| "ended_at": None, | |
| } | |
| await _lru_set(session["session_id"], session) | |
| asyncio.create_task(_redis_set_session(session["session_id"], session)) | |
| asyncio.create_task(_supa_upsert(session)) | |
| _logger.info("[session_mgr] created session=%s user=%s", session["session_id"][:8], user_id) | |
| return session | |
| async def get_session(session_id: str) -> dict | None: | |
| """Recupera una sessione, cercando nei tier in ordine di velocitΓ .""" | |
| s = await _lru_get(session_id) | |
| if s: return s | |
| s = await _redis_get_session(session_id) | |
| if s: | |
| await _lru_set(session_id, s) | |
| return s | |
| s = await _supa_get(session_id) | |
| if s: | |
| await _lru_set(session_id, s) | |
| asyncio.create_task(_redis_set_session(session_id, s)) | |
| return s | |
| async def touch_session(session_id: str) -> None: | |
| """Aggiorna last_active_at e rinnova il TTL Redis.""" | |
| s = await get_session(session_id) | |
| if not s: return | |
| s["last_active_at"] = time.time() | |
| s["status"] = SessionStatus.ACTIVE | |
| await _lru_set(session_id, s) | |
| asyncio.create_task(_redis_set_session(session_id, s)) | |
| asyncio.create_task(_supa_upsert(s)) | |
| # ββ Pydantic models ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| class CreateSessionRequest(BaseModel): | |
| user_id: str | None = None | |
| metadata: dict = Field(default_factory=dict) | |
| class PatchSessionRequest(BaseModel): | |
| status: str | None = None | |
| metadata: dict | None = None | |
| # ββ Endpoints ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| async def create_session_endpoint(req: CreateSessionRequest): | |
| session = await create_session(user_id=req.user_id, metadata=req.metadata) | |
| return session | |
| async def session_manager_status(): | |
| async with _lru_lock: | |
| lru_count = len(_lru) | |
| return { | |
| "status": "ok", | |
| "component": "session_manager", | |
| "lru_sessions": lru_count, | |
| "lru_max": _LRU_MAX, | |
| "session_ttl_s": _SESS_TTL_S, | |
| "tiers": { | |
| "memory": "active", | |
| "redis": "active" if os.getenv("UPSTASH_REDIS_REST_URL") else "not_configured", | |
| "supabase": "active" if _sb else "not_configured", | |
| }, | |
| } | |
| async def get_session_endpoint(session_id: str): | |
| s = await get_session(session_id) | |
| if not s: | |
| raise HTTPException(404, detail=f"Sessione {session_id} non trovata") | |
| return s | |
| async def patch_session_endpoint(session_id: str, req: PatchSessionRequest): | |
| s = await get_session(session_id) | |
| if not s: | |
| raise HTTPException(404, detail=f"Sessione {session_id} non trovata") | |
| if req.status: | |
| s["status"] = req.status | |
| if req.metadata is not None: | |
| s["metadata"].update(req.metadata) | |
| s["last_active_at"] = time.time() | |
| await _lru_set(session_id, s) | |
| asyncio.create_task(_redis_set_session(session_id, s)) | |
| asyncio.create_task(_supa_upsert(s)) | |
| return s | |
| async def end_session_endpoint(session_id: str): | |
| s = await get_session(session_id) | |
| if not s: | |
| raise HTTPException(404, detail=f"Sessione {session_id} non trovata") | |
| s["status"] = SessionStatus.ENDED | |
| s["ended_at"] = time.time() | |
| await _lru_set(session_id, s) | |
| asyncio.create_task(_redis_del_session(session_id)) | |
| asyncio.create_task(_supa_upsert(s)) | |
| _logger.info("[session_mgr] ended session=%s", session_id[:8]) | |
| return {"session_id": session_id, "status": SessionStatus.ENDED} | |