from uuid import UUID from fastapi import APIRouter, Depends, Query from app.api.deps import ( AsyncSession, CurrentUser, get_db, get_pagination, require_tenant, ) from app.api.permissions import require_permissions from app.permissions.definitions import ( MESSAGES_CANCEL, MESSAGES_SEND, MESSAGES_VIEW, QUEUE_MANAGE, QUEUE_VIEW, ) from app.schemas.common import DeliveryEventOut, MessageOut, MessageSendRequest, QueueItemOut from app.services.message_service import MessageService from app.services.queue_engine import QueueEngine from shared.pagination import PaginationParams router = APIRouter() @router.post("/send", response_model=list[MessageOut], status_code=201) async def send_message( body: MessageSendRequest, tenant_id: UUID = Depends(require_tenant), session: AsyncSession = Depends(get_db), user: CurrentUser = Depends(require_permissions(MESSAGES_SEND)), ): return await MessageService(session).send( tenant_id, body.model_dump(), actor_id=user.user_id ) @router.get("", response_model=list[MessageOut]) async def list_messages( status: str | None = Query(default=None), tenant_id: UUID = Depends(require_tenant), session: AsyncSession = Depends(get_db), pagination: PaginationParams = Depends(get_pagination), _: CurrentUser = Depends(require_permissions(MESSAGES_VIEW)), ): return await MessageService(session).list( tenant_id, status=status, offset=pagination.offset, limit=pagination.limit, ) @router.get("/stats") async def message_stats( tenant_id: UUID = Depends(require_tenant), session: AsyncSession = Depends(get_db), _: CurrentUser = Depends(require_permissions(MESSAGES_VIEW)), ): return await MessageService(session).stats(tenant_id) @router.get("/{message_id}", response_model=MessageOut) async def get_message( message_id: UUID, tenant_id: UUID = Depends(require_tenant), session: AsyncSession = Depends(get_db), _: CurrentUser = Depends(require_permissions(MESSAGES_VIEW)), ): return await MessageService(session).get(tenant_id, message_id) @router.get("/{message_id}/timeline", response_model=list[DeliveryEventOut]) async def message_timeline( message_id: UUID, tenant_id: UUID = Depends(require_tenant), session: AsyncSession = Depends(get_db), _: CurrentUser = Depends(require_permissions(MESSAGES_VIEW)), ): return await MessageService(session).timeline(tenant_id, message_id) @router.post("/{message_id}/cancel", response_model=MessageOut) async def cancel_message( message_id: UUID, tenant_id: UUID = Depends(require_tenant), session: AsyncSession = Depends(get_db), _: CurrentUser = Depends(require_permissions(MESSAGES_CANCEL)), ): return await MessageService(session).cancel(tenant_id, message_id) @router.post("/queue/process") async def process_queue( tenant_id: UUID = Depends(require_tenant), session: AsyncSession = Depends(get_db), _: CurrentUser = Depends(require_permissions(QUEUE_MANAGE)), limit: int = Query(default=20, ge=1, le=200), ): processed = await QueueEngine(session).process_due(tenant_id, limit=limit) return {"processed": processed} @router.get("/queue/dead-letters", response_model=list[QueueItemOut]) async def dead_letters( tenant_id: UUID = Depends(require_tenant), session: AsyncSession = Depends(get_db), _: CurrentUser = Depends(require_permissions(QUEUE_VIEW)), ): return await QueueEngine(session).list_dead_letters(tenant_id) @router.get("/queue/stats") async def queue_stats( tenant_id: UUID = Depends(require_tenant), session: AsyncSession = Depends(get_db), _: CurrentUser = Depends(require_permissions(QUEUE_VIEW)), ): return await QueueEngine(session).stats(tenant_id)