import json import logging import uuid from datetime import date from typing import Any, cast import redis from confluent_kafka import KafkaError, Message, Producer from fansifter_common.adapters.db import Database from fansifter_common.constants import QA_ENVIRONMENT from httpx import HTTPStatusError from app.connectors.aws.s3 import S3Client from app.connectors.db.repositories.account_to_fan_response_email_address_mapping import ( AccountToFanResponseEmailAddressMappingRepository, ) from app.connectors.ows_preference_center import OwsPreferenceCenterClient from app.connectors.sendgrid import SendgridClient from app.models import ReplyRequest, UnsubRequest logger = logging.getLogger(__name__) class SendgridWebhooksUnsubHandler: def __init__( self, kafka_producer: Producer, kafka_sendgrid_inbound_topic: str, s3_client: S3Client, failed_events_s3_bucket: str, ows_preference_center: OwsPreferenceCenterClient, ) -> None: self.kafka_producer = kafka_producer self.kafka_sendgrid_inbound_topic = kafka_sendgrid_inbound_topic self.s3_client = s3_client self.failed_events_s3_bucket = failed_events_s3_bucket self.ows_preference_center = ows_preference_center def handle( self, request: UnsubRequest, event_dict: dict[str, Any] ) -> dict[str, int]: try: self.unsubscribe_fan(request) self.send_kafka_message(request) logger.info( "Unsubscribe request has been processed", extra={"email_hash": request.email_hash}, ) except Exception as exc: if isinstance(exc, HTTPStatusError) and exc.response.status_code == 404: logger.warning( "Fan not found in OWS Preference Center", extra={"email_hash": request.email_hash}, ) else: logger.exception("Handler exception") self.s3_client.put_object( bucket=self.failed_events_s3_bucket, key=f"inbound/{date.today().isoformat()}/{uuid.uuid4()}.json", body=json.dumps(event_dict), ) return {"statusCode": 200} def unsubscribe_fan(self, request: UnsubRequest) -> None: self.ows_preference_center.unsubscribe( email=request.email, email_type=request.email_type.value, email_id=request.email_id, automated_email_trigger_id=request.automated_email_trigger_id, ) def send_kafka_message(self, request: UnsubRequest) -> None: self.kafka_producer.produce( self.kafka_sendgrid_inbound_topic, key=str(uuid.uuid4()), value=request.model_dump_json( include={ "email", "email_campaign_id", "email_type", "email_id", "automated_email_trigger_id", "timestamp", "recipient_email", } ), callback=self.kafka_producer_callback, ) self.kafka_producer.flush() @staticmethod def kafka_producer_callback(error: KafkaError | None, message: Message) -> None: if error: logger.error("Kafka producer error %s", error) class SendgridWebhooksReplyHandler: def __init__( self, kafka_producer: Producer, kafka_sendgrid_reply_topic: str, db: Database, repository: AccountToFanResponseEmailAddressMappingRepository, sendgrid_client: SendgridClient, redis_client: redis.Redis, cache_key_template: str, requests_count_limit: int, requests_count_reset_seconds: int, s3_client: S3Client, failed_events_s3_bucket: str, default_forward_email: str, from_email_template: str, accounts_to_forward_on_qa_env: list[int], env: str, ) -> None: self.kafka_producer = kafka_producer self.kafka_sendgrid_reply_topic = kafka_sendgrid_reply_topic self.db = db self.repository = repository self.sendgrid_client = sendgrid_client self.redis_client = redis_client self.cache_key_template = cache_key_template self.requests_count_limit = requests_count_limit self.requests_count_reset_seconds = requests_count_reset_seconds self.s3_client = s3_client self.failed_events_s3_bucket = failed_events_s3_bucket self.default_forward_email = default_forward_email self.from_email_template = from_email_template self.accounts_to_forward_on_qa_env = accounts_to_forward_on_qa_env self.env = env def handle( self, request: ReplyRequest, event_dict: dict[str, Any] ) -> dict[str, int]: try: requests_count = self.check_cache(request) if requests_count > self.requests_count_limit: logger.warning( "Too many requests for the same email, skipping processing", extra={ "email_hash": request.email_hash, "requests_count": requests_count, }, ) return {"statusCode": 200} self.forward_email(request) self.send_kafka_message(request) logger.info( "Reply event has been processed", extra={"email_hash": request.email_hash}, ) except Exception: logger.exception("Handler exception") self.s3_client.put_object( bucket=self.failed_events_s3_bucket, key=f"reply/{date.today().isoformat()}/{uuid.uuid4()}.json", body=json.dumps(event_dict), ) return {"statusCode": 200} def check_cache(self, request: ReplyRequest) -> int: cache_key = self.cache_key_template.format( email_id=request.email_id, sender_email_hash=request.email_hash ) requests_count = cast(int, self.redis_client.incr(cache_key)) if requests_count == 1: self.redis_client.expire(cache_key, self.requests_count_reset_seconds) return requests_count def forward_email(self, request: ReplyRequest) -> None: from_email = self.from_email_template.format(unique_part=str(uuid.uuid4())[:8]) with self.db.session_factory(): forward_email_mapping = self.repository.get_forward_email_mapping( email_id=request.email_id, email_type=request.email_type.value ) # We do not want to forward any emails on QA environment # except test account if forward_email_mapping and self.env == QA_ENVIRONMENT: if ( forward_email_mapping.vendor_id not in self.accounts_to_forward_on_qa_env ): logger.info( "Skipping forwarding email for account in QA environment", extra={"vendor_id": forward_email_mapping.vendor_id}, ) return self.sendgrid_client.send_email( to_email=forward_email_mapping.email_address if forward_email_mapping else self.default_forward_email, subject=request.subject, text_content=request.text, html_content=request.html, from_email=from_email, from_email_display_name=request.email, reply_to=request.email, categories=["Fan Reply Forwarding"], custom_args={"send_type": "fan_reply_forwarding"}, ) def send_kafka_message(self, request: ReplyRequest) -> None: self.kafka_producer.produce( self.kafka_sendgrid_reply_topic, key=str(uuid.uuid4()), value=request.model_dump_json( include={ "email", "email_campaign_id", "email_type", "email_id", "automated_email_trigger_id", "timestamp", "recipient_email", "subject", "text", "html", } ), callback=self.kafka_producer_callback, ) self.kafka_producer.flush() @staticmethod def kafka_producer_callback(error: KafkaError | None, message: Message) -> None: if error: logger.error("Kafka producer error %s", error)