TorbatYar/backend/services/communication/app/api/v1/messages.py
Mortezakoohjani e41ecfad4c Sync platform docs, infra, and module services with Accounting integration.
Include Loyalty/Communication/Sports Center backends and registry updates alongside production nginx and compose wiring.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-25 22:35:23 +03:30

122 lines
3.9 KiB
Python

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)