"""Ingestion: приём событий по контракту docs/05-ingestion-api.md.
Приём быстрый: жёсткая валидация конверта + дедупликация + вставка в БД +
бросок в Celery (202 мгновенно). Никакой работы с каналами здесь.
"""
import uuid
from datetime import UTC, datetime, timedelta
from fastapi import APIRouter, Depends, HTTPException, Request, Response, status
from pydantic import ValidationError
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.api.schemas import (
AcceptedEvent,
BatchRejected,
BatchResult,
DeliverySummary,
EventEnvelope,
EventStatus,
)
from app.auth.apikeys import resolve_event_source
from app.database import get_db
from app.models import ChannelTarget, Delivery, Event, NotificationType, Source
from app.settings_store import get_setting
from app.worker.celery_app import celery_app
router = APIRouter(prefix="/api/v1/events", tags=["ingestion"])
class Deduplicated(Exception):
"""Сигнал из accept_event: это дубль события (id и статус первого)."""
def __init__(self, event_id: uuid.UUID, event_status: str) -> None:
self.event_id = event_id
self.event_status = event_status
def accept_event(db: Session, source: Source, envelope: EventEnvelope) -> Event:
"""Общая логика приёма (одиночный и batch, и send_test_event из MCP):
проверки + вставка. Принимает источник (не ключ) — MCP вызывает его
для тест-события без выдачи API-ключа себе.
Диспетчеризацию в Celery выполняет вызывающий код после коммита.
"""
if envelope.source != source.name:
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail=(
f"Поле source='{envelope.source}' не совпадает с источником "
f"ключа '{source.name}'"
),
)
known_type = db.execute(
select(NotificationType.id).where(
NotificationType.source_id == source.id,
NotificationType.subject == envelope.subject,
NotificationType.action == envelope.action,
NotificationType.deleted_at.is_(None),
)
).scalar_one_or_none()
if known_type is None:
raise HTTPException(
status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
detail=(
f"Тип ({envelope.source}, {envelope.subject}, {envelope.action}) "
"не зарегистрирован — обратись к админу Synapse"
),
)
# Дедупликация: тот же dedup_key от источника в окне времени — дубль.
# Точность в рантайме под нагрузкой не гарантирована (нет unique-констрейнта);
# race на «одновременно влетевшие» дублей в MVP не ловим.
if envelope.dedup_key:
window = int(get_setting(db, "dedup_window_seconds"))
prior = db.execute(
select(Event)
.where(
Event.source_id == source.id,
Event.dedup_key == envelope.dedup_key,
Event.created_at > datetime.now(UTC) - timedelta(seconds=window),
)
.order_by(Event.created_at.desc())
).scalars().first()
if prior is not None:
raise Deduplicated(prior.id, prior.status)
event = Event(
source_id=source.id,
subject=envelope.subject,
action=envelope.action,
priority=envelope.priority,
payload=envelope.payload,
dedup_key=envelope.dedup_key,
expires_at=(
datetime.now(UTC) + timedelta(seconds=envelope.ttl_seconds)
if envelope.ttl_seconds
else None
),
scheduled_at=envelope.scheduled_at,
status="queued",
)
db.add(event)
# flush: uuid4-дефолт id срабатывает на INSERT; batch строит ответ до коммита.
db.flush()
return event
def _dispatch(event_id: str) -> None:
celery_app.send_task("synapse.ingest", args=[event_id])
@router.post("", status_code=202, response_model=AcceptedEvent)
def ingest(
envelope: EventEnvelope,
request: Request,
db: Session = Depends(get_db),
) -> AcceptedEvent:
api_key, source = resolve_event_source(db, request)
try:
event = accept_event(db, source, envelope)
except Deduplicated as dup:
return AcceptedEvent(id=str(dup.event_id), status=dup.event_status, deduplicated=True)
api_key.last_used_at = datetime.now(UTC)
db.commit()
_dispatch(str(event.id))
return AcceptedEvent(id=str(event.id), status=event.status, deduplicated=False)
@router.post("/batch", response_model=BatchResult)
def ingest_batch(
raw_envelopes: list[dict],
request: Request,
response: Response,
db: Session = Depends(get_db),
) -> BatchResult:
"""Пачка событий (smart-home, Navi). Частичная валидация честная:
ошибки по индексам, валидные приняты; 202, если принято хоть что-то,
иначе 422 (откат — транзакцию нельзя коммитить частично).
"""
api_key, source = resolve_event_source(db, request)
results: list[AcceptedEvent | BatchRejected] = []
accepted: list[str] = []
for i, raw in enumerate(raw_envelopes):
try:
envelope = EventEnvelope.model_validate(raw)
except ValidationError as err:
results.append(BatchRejected(index=i, detail=err.errors()[0].get("msg", "невалидно")))
continue
try:
event = accept_event(db, source, envelope)
results.append(AcceptedEvent(id=str(event.id), status=event.status))
accepted.append(str(event.id))
except Deduplicated as dup:
results.append(
AcceptedEvent(id=str(dup.event_id), status=dup.event_status, deduplicated=True)
)
except HTTPException as err:
results.append(
BatchRejected(
index=i,
detail=err.detail if isinstance(err.detail, str) else "Ошибка приёма",
)
)
if not accepted:
db.rollback()
response.status_code = status.HTTP_422_UNPROCESSABLE_ENTITY
else:
api_key.last_used_at = datetime.now(UTC)
db.commit()
response.status_code = status.HTTP_202_ACCEPTED
for event_id in accepted:
_dispatch(event_id)
return BatchResult(results=results)
@router.get("/{event_id}", response_model=EventStatus)
def event_status(
event_id: uuid.UUID,
request: Request,
db: Session = Depends(get_db),
) -> EventStatus:
"""Статус события: только его источнику, чужой id — 404."""
api_key, source = resolve_event_source(db, request)
event = db.execute(
select(Event).where(Event.id == event_id, Event.source_id == api_key.source_id)
).scalar_one_or_none()
if event is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Событие не найдено")
deliveries = (
db.execute(select(Delivery).where(Delivery.event_id == event.id).order_by(Delivery.id))
.scalars()
.all()
)
out: list[DeliverySummary] = []
for d in deliveries:
target_name: str | None = None
if d.channel_target_id:
target = db.get(ChannelTarget, d.channel_target_id)
target_name = target.name if target else None
out.append(
DeliverySummary(
channel=d.channel,
target=target_name,
status=d.status,
attempts=d.attempts,
error=d.last_error,
rendered_message=d.rendered_message,
)
)
return EventStatus(
id=str(event.id),
source=source.name,
subject=event.subject,
action=event.action,
priority=event.priority,
status=event.status,
created_at=event.created_at,
expires_at=event.expires_at,
deliveries=out,
)