""" 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 ────────────────────────────────────────────────────────────────── @router.post("", summary="Crea una nuova sessione") async def create_session_endpoint(req: CreateSessionRequest): session = await create_session(user_id=req.user_id, metadata=req.metadata) return session @router.get("/status", summary="Diagnostica Session Manager") 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", }, } @router.get("/{session_id}", summary="Recupera sessione") 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 @router.patch("/{session_id}", summary="Aggiorna sessione") 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 @router.delete("/{session_id}/end", summary="Termina sessione") 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}