import http from aiohttp import web from aiohttp_apispec import docs, headers_schema, json_schema, querystring_schema from apollo_utils.service.clients.aiohttp.utils.response import dump_response_schema from apollo_utils.service.exceptions import NotFound, UnprocessableEntity from apollo_utils.service.schemas.headers import AppHeader from server import config from server.constants.core import BASE_API_PREFIX from server.db.models.core import Application from server.db.models.messages import Event from server.publishers.clients.sns.core import sns_publish_bulk from server.publishers.schemas.events import EventCreatedEvent from server.scenarios.messages.create_bulk import create_multiple_messages_scenario from server.scenarios.messages.get_bulk import get_multiple_messages_scenario from server.schemas.events import Events router = web.RouteTableDef() @router.get(BASE_API_PREFIX + "/service/events/") @docs(tags=["service", "event"], summary="Get an event by id.") @headers_schema(AppHeader) @querystring_schema(Events.Get.Request) @dump_response_schema(Events.Get.Response, apply=True) async def get_event(request: web.Request) -> web.Response: event_id, app_slug = request["querystring"]["id"], request["headers"]["app"] event = await Event.get(id=event_id) if not event: raise NotFound(f"Event with id={event_id} does not exist.") return event @router.post(BASE_API_PREFIX + "/service/events/") @docs(tags=["service", "event"], summary="Create an event.") @headers_schema(AppHeader) @json_schema(Events.Create.Request) @dump_response_schema(Events.Create.Response, code=http.HTTPStatus.CREATED, apply=True) async def post_event(request: web.Request) -> web.Response: data, app_slug = request["json"], request["headers"]["app"] public = data["public"] and config.ALLOW_EVENTS_PUBLISHING app = await Application.get(slug=app_slug) if not app: raise NotFound(f"An application with the slug {app_slug} is not found.") data["app_slug"] = app_slug event = await Event.create(data) errors = None exc = None if public: try: event_public_data = EventCreatedEvent().dump(event) errors = await sns_publish_bulk([event_public_data]) except Exception as _exc: exc = _exc if errors or exc: await Event.delete(id=event.id) raise (exc or UnprocessableEntity(extra=errors.values())) return event @router.post(BASE_API_PREFIX + "/service/events/list/") @docs( tags=["service", "events", "list"], summary="Get multiple events id.", ) @headers_schema(AppHeader) @json_schema(Events.List.Request) @dump_response_schema(Events.List.Response, apply=True) async def get_events(request: web.Request) -> web.Response: return await get_multiple_messages_scenario(Event, data=request["json"], app_slug=request["headers"]["app"]) @router.post(BASE_API_PREFIX + "/v2/service/events/") @docs( tags=["service", "event", "v2"], summary="Create multiple events.", ) @headers_schema(AppHeader) @json_schema(Events.CreateBulk.Request) @dump_response_schema(Events.CreateBulk.Response, apply=False) async def create_events(request: web.Request) -> web.Response: data = request["json"] public = data["public"] and config.ALLOW_EVENTS_PUBLISHING unique_mode, unique_key = data.pop("unique_mode"), data.pop("unique_key", None) return await create_multiple_messages_scenario( Event, data=data, app_slug=request["headers"]["app"], public=public, publish_schema=EventCreatedEvent, check_unique_mode=unique_mode, check_unique_key=unique_key, )