import json import logging import uuid from dataclasses import dataclass from datetime import date from typing import Any from uuid import uuid4 from confluent_kafka import KafkaError, Message, Producer from pydantic import TypeAdapter, ValidationError from app.auth import Auth from app.connectors.aws.s3 import S3Client from app.exceptions import AuthException from app.models import BounceEvent, Event, ReplyForwardingEvent, Request EventsValidator = TypeAdapter(list[Event | BounceEvent | ReplyForwardingEvent]) logger = logging.getLogger(__name__) @dataclass(kw_only=True) class SendgridEventHandler: auth: Auth kafka_producer: Producer kafka_sendgrid_events_topic: str s3_client: S3Client failed_events_s3_bucket: str def handle(self, request_dict: dict[str, Any]) -> dict[str, int]: try: request = Request.model_validate(request_dict) self.auth.validate(request) events = EventsValidator.validate_json(request.body) for event in events: if isinstance(event, ReplyForwardingEvent): logger.info( "Skipping fan reply forwarding event", extra={"event": event.model_dump()}, ) continue self.kafka_producer.produce( self.kafka_sendgrid_events_topic, key=str(uuid4()), value=event.model_dump_json(), callback=self.kafka_producer_callback, ) self.kafka_producer.flush() logger.info( "Events processing success", extra={"events_count": len(events)}, ) status_code = 200 except AuthException: logger.exception("Auth error") status_code = 401 except ValidationError: logger.exception("Events validation error") status_code = 400 except Exception: logger.exception("Handler exception") status_code = 500 if status_code != 200: self.s3_client.put_object( bucket=self.failed_events_s3_bucket, key=f"general/{date.today().isoformat()}/{uuid.uuid4()}.json", body=json.dumps(request_dict), ) return {"statusCode": status_code} @staticmethod def kafka_producer_callback(error: KafkaError | None, _message: Message) -> None: if error: logger.error("Kafka producer error %s", error)