Terminal / api /event_store.py
Baida-A
Initial clean deploy (Reverse Proxy removed)
28a08e7
Raw
History Blame
8.02 kB
"""
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 ──────────────────────────────────────────────────────────────────
@router.post("/store", summary="Salva evento nello store")
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}")
@router.get("/replay", summary="Replay eventi filtrati")
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}")
@router.get("/store/{event_id}", summary="Recupera evento singolo")
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}")
@router.get("/store/status", summary="Diagnostica Event Store")
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)}