Spaces:
Running
Running
| """ | |
| backend/api/event_store.py β Event Store (persistenza, Fase 1 ADR-S26-S30) | |
| ResponsabilitΓ : SALVARE tutti gli eventi per replayability, debugging, benchmark. | |
| NON instrada β per pub/sub usa event_bus.py. | |
| Schema Supabase (tabella `event_store`, auto-created se non esiste): | |
| id UUID PK default gen_random_uuid() | |
| topic TEXT NOT NULL | |
| payload JSONB NOT NULL default '{}' | |
| correlation_id TEXT | |
| session_id TEXT | |
| source TEXT | |
| created_at TIMESTAMPTZ NOT NULL default now() | |
| Invarianti ADR: | |
| S26: ogni workflow (e ogni evento del workflow) Γ¨ persistente | |
| S27: correlation_id garantisce tracciabilitΓ cross-componente | |
| Endpoints: | |
| POST /api/events/store β salva evento (auth: MACHINE) | |
| GET /api/events/replay β replay eventi filtrati (auth: MACHINE) | |
| GET /api/events/store/{id} β recupera evento singolo (auth: MACHINE) | |
| GET /api/events/store/status β diagnostica store (auth: MACHINE) | |
| """ | |
| import json, time, uuid, logging, os | |
| from typing import Any | |
| from fastapi import APIRouter, Depends, Query, HTTPException | |
| from pydantic import BaseModel, Field | |
| from .auth_guard import require_role, AuthRole | |
| from .state import _sb | |
| _logger = logging.getLogger("api.event_store") | |
| router = APIRouter( | |
| prefix="/api/events", | |
| tags=["event-store"], | |
| dependencies=[Depends(require_role(AuthRole.MACHINE))], | |
| ) | |
| _TABLE = "event_store" | |
| # ββ Auto-create table (best-effort, richiede service role key) βββββββββββββββββ | |
| _TABLE_CREATED = False | |
| async def _ensure_table() -> bool: | |
| """Crea la tabella event_store su Supabase se non esiste. Best-effort.""" | |
| global _TABLE_CREATED | |
| if _TABLE_CREATED: | |
| return True | |
| if not _sb: | |
| return False | |
| try: | |
| # Prova una SELECT β se la tabella non esiste, Supabase ritorna un errore | |
| res = _sb.table(_TABLE).select("id").limit(1).execute() | |
| _TABLE_CREATED = True | |
| return True | |
| except Exception as exc: | |
| _logger.warning("[event_store] tabella '%s' non raggiungibile: %s β " | |
| "crea manualmente con migration Supabase", _TABLE, exc) | |
| return False | |
| # ββ Pydantic models ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| class StoreEventRequest(BaseModel): | |
| topic: str | |
| payload: dict = Field(default_factory=dict) | |
| correlation_id: str | None = None | |
| session_id: str | None = None | |
| source: str | None = None | |
| class StoredEvent(BaseModel): | |
| id: str | |
| topic: str | |
| payload: dict | |
| correlation_id: str | None | |
| session_id: str | None | |
| source: str | None | |
| created_at: str | None | |
| # ββ Endpoints ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| async def store_event(req: StoreEventRequest) -> StoredEvent: | |
| """ | |
| Persiste un evento nel Supabase Event Store. | |
| Chiamato automaticamente dall'event_bus (via hook) o esplicitamente | |
| dai componenti che vogliono garantire persistenza. | |
| """ | |
| await _ensure_table() | |
| if not _sb: | |
| raise HTTPException(503, detail="Event Store non disponibile (Supabase non configurato)") | |
| record = { | |
| "topic": req.topic, | |
| "payload": req.payload, | |
| "correlation_id": req.correlation_id or str(uuid.uuid4()), | |
| "session_id": req.session_id, | |
| "source": req.source, | |
| } | |
| try: | |
| res = _sb.table(_TABLE).insert(record).execute() | |
| row = res.data[0] if res.data else {**record, "id": str(uuid.uuid4()), "created_at": None} | |
| return StoredEvent(**{ | |
| "id": row.get("id", ""), | |
| "topic": row.get("topic", req.topic), | |
| "payload": row.get("payload", req.payload), | |
| "correlation_id": row.get("correlation_id"), | |
| "session_id": row.get("session_id"), | |
| "source": row.get("source"), | |
| "created_at": str(row.get("created_at", "")), | |
| }) | |
| except Exception as exc: | |
| _logger.error("[event_store] insert failed: %s", exc) | |
| raise HTTPException(500, detail=f"Event Store insert error: {exc}") | |
| async def replay_events( | |
| topic: str | None = Query(None, description="Filtra per topic"), | |
| session_id: str | None = Query(None, description="Filtra per session_id"), | |
| correlation_id: str | None = Query(None, description="Filtra per correlation_id"), | |
| from_ts: float | None = Query(None, description="Unix timestamp minimo (created_at >=)"), | |
| limit: int = Query(100, ge=1, le=1000), | |
| ): | |
| """ | |
| Recupera eventi filtrati dall'Event Store. Supporta replay per debugging, | |
| test di regressione e audit trail. | |
| """ | |
| await _ensure_table() | |
| if not _sb: | |
| raise HTTPException(503, detail="Event Store non disponibile") | |
| try: | |
| q = _sb.table(_TABLE).select("*").order("created_at", desc=True).limit(limit) | |
| if topic: q = q.eq("topic", topic) | |
| if session_id: q = q.eq("session_id", session_id) | |
| if correlation_id: q = q.eq("correlation_id", correlation_id) | |
| if from_ts: | |
| import datetime | |
| dt = datetime.datetime.utcfromtimestamp(from_ts).isoformat() + "Z" | |
| q = q.gte("created_at", dt) | |
| res = q.execute() | |
| return { | |
| "events": res.data or [], | |
| "count": len(res.data or []), | |
| "filters": { | |
| "topic": topic, "session_id": session_id, | |
| "correlation_id": correlation_id, "from_ts": from_ts, "limit": limit, | |
| }, | |
| } | |
| except Exception as exc: | |
| _logger.error("[event_store] replay failed: %s", exc) | |
| raise HTTPException(500, detail=f"Event Store query error: {exc}") | |
| async def get_event(event_id: str) -> StoredEvent: | |
| """Recupera un evento specifico per ID.""" | |
| await _ensure_table() | |
| if not _sb: | |
| raise HTTPException(503, detail="Event Store non disponibile") | |
| try: | |
| res = _sb.table(_TABLE).select("*").eq("id", event_id).limit(1).execute() | |
| if not res.data: | |
| raise HTTPException(404, detail=f"Evento {event_id} non trovato") | |
| row = res.data[0] | |
| return StoredEvent(**{ | |
| "id": row.get("id", event_id), | |
| "topic": row.get("topic", ""), | |
| "payload": row.get("payload", {}), | |
| "correlation_id": row.get("correlation_id"), | |
| "session_id": row.get("session_id"), | |
| "source": row.get("source"), | |
| "created_at": str(row.get("created_at", "")), | |
| }) | |
| except HTTPException: | |
| raise | |
| except Exception as exc: | |
| raise HTTPException(500, detail=f"Event Store get error: {exc}") | |
| async def store_status(): | |
| """Verifica connettivitΓ dello store e restituisce statistiche.""" | |
| if not _sb: | |
| return {"status": "unavailable", "reason": "Supabase non configurato"} | |
| try: | |
| res = _sb.table(_TABLE).select("topic", count="exact").execute() | |
| total = res.count if hasattr(res, "count") and res.count else len(res.data or []) | |
| return { | |
| "status": "ok", | |
| "component": "event_store", | |
| "total_events": total, | |
| "table": _TABLE, | |
| } | |
| except Exception as exc: | |
| return {"status": "error", "detail": str(exc)} | |